From 50d60e24741517221884ccdd77086830919bb7bc Mon Sep 17 00:00:00 2001 From: Sylvester Kaczmarek <16242628+sylvesterkaczmarek@users.noreply.github.com> Date: Fri, 25 Sep 2026 12:06:13 +0100 Subject: [PATCH] Retry failed scheduler priority notifications Fixes #121. Add regression coverage and preserve existing successful behavior. Signed-off-by: Sylvester Kaczmarek <16242628+sylvesterkaczmarek@users.noreply.github.com> --- services/nvpair-job-scheduler/README.md | 8 ++ .../nvpair-job-scheduler/delivery_test.go | 108 ++++++++++++++++++ services/nvpair-job-scheduler/schedule.go | 7 +- 3 files changed, 120 insertions(+), 3 deletions(-) create mode 100644 services/nvpair-job-scheduler/delivery_test.go diff --git a/services/nvpair-job-scheduler/README.md b/services/nvpair-job-scheduler/README.md index 68090b04..63ffc434 100644 --- a/services/nvpair-job-scheduler/README.md +++ b/services/nvpair-job-scheduler/README.md @@ -129,3 +129,11 @@ go test ./... - [`../nvpair-workload-manager/README.md`](../nvpair-workload-manager/README.md) — where workload state originates - [`../VERSIONING.md`](../VERSIONING.md) — SemVer bump rules + +## Delivery failures + +The last-emitted snapshot and timestamp are updated only after the local +JSON-RPC notification write succeeds. A failed update leaves the previous +successful snapshot intact, so ordinary reconciliation retries changed ranks +without requiring another workload event. Each engine is tracked independently; +a successful delivery is not repeated merely because another engine failed. diff --git a/services/nvpair-job-scheduler/delivery_test.go b/services/nvpair-job-scheduler/delivery_test.go new file mode 100644 index 00000000..273580b7 --- /dev/null +++ b/services/nvpair-job-scheduler/delivery_test.go @@ -0,0 +1,108 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "io" + "testing" +) + +// unavailableWriter is an in-memory test fixture; it opens no files or sockets. +type unavailableWriter struct{} + +func (unavailableWriter) Read([]byte) (int, error) { return 0, io.EOF } +func (unavailableWriter) Write([]byte) (int, error) { return 0, io.ErrClosedPipe } + +func TestFailedNotificationDoesNotAdvanceStatus(t *testing.T) { + m := mgrWith(unavailableWriter{}, []string{"test-node-a"}) + m.recomputeAll(false) + for _, engine := range schedulerEngines { + status := m.status().Engines[engine] + if status.LastEmittedAt != 0 || len(status.Emitted) != 0 { + t.Fatalf("failed notification was recorded as emitted: %+v", status) + } + } +} + +// recoveringWriter simulates one failed local codec write, then records frames. +type recoveringWriter struct { + capRW + remaining int + shortWrite bool +} + +func (w *recoveringWriter) Write(data []byte) (int, error) { + if w.remaining > 0 { + w.remaining-- + if w.shortWrite { + return 0, nil + } + return 0, io.ErrClosedPipe + } + return w.capRW.Write(data) +} + +func TestFailedNotificationRetriesWithoutChangingRanks(t *testing.T) { + for _, shortWrite := range []bool{false, true} { + name := "write-error" + if shortWrite { + name = "zero-byte-write" + } + t.Run(name, func(t *testing.T) { + writer := &recoveringWriter{remaining: 1, shortWrite: shortWrite} + m := mgrWith(writer, []string{"test-node-a", "test-node-b"}) + m.recomputeAll(false) + m.recomputeAll(false) + m.recomputeAll(false) + for _, engine := range schedulerEngines { + orders := writer.orders(engine) + if len(orders) != 1 { + t.Fatalf("%s received %d frames, want one successful delivery", engine, len(orders)) + } + assertStrs(t, orders[0], []string{"test-node-a", "test-node-b"}) + if m.status().Engines[engine].LastEmittedAt == 0 { + t.Fatal("successful delivery was not recorded") + } + } + }) + } +} + +func TestFailedChangedNotificationRetainsDeliveredRanks(t *testing.T) { + writer := &recoveringWriter{} + m := mgrWith(writer, []string{"test-node-a", "test-node-b"}) + m.recomputeAll(false) + engine := schedulerEngines[0] + previous := m.status().Engines[engine] + writer.remaining = 1 + m.nodes["test-node-c"] = true + m.recomputeAll(false) + current := m.status().Engines[engine] + if !equalRanks(current.Emitted, previous.Emitted) || current.LastEmittedAt != previous.LastEmittedAt { + t.Fatal("failed update replaced the last successfully delivered snapshot") + } + m.recomputeAll(false) + orders := writer.orders(engine) + if len(orders) != 2 { + t.Fatalf("received %d snapshots, want initial and recovered update", len(orders)) + } + assertStrs(t, orders[1], []string{"test-node-a", "test-node-b", "test-node-c"}) +} + +func TestFailedForcedNotificationRetainsDeliveredTimestamp(t *testing.T) { + writer := &recoveringWriter{} + m := mgrWith(writer, []string{"test-node-a"}) + m.recomputeAll(false) + engine := schedulerEngines[0] + m.emitted[engine] = engineState{ranks: m.emitted[engine].ranks, lastEmittedAt: 1} + writer.remaining = 1 + m.recomputeAll(true) + if got := m.status().Engines[engine].LastEmittedAt; got != 1 { + t.Fatalf("failed forced delivery advanced timestamp to %d", got) + } + m.recomputeAll(false) + if got := len(writer.orders(engine)); got != 1 { + t.Fatalf("previously delivered unchanged ranks were emitted %d times", got) + } +} diff --git a/services/nvpair-job-scheduler/schedule.go b/services/nvpair-job-scheduler/schedule.go index 2bfbf7fc..2e0f3a95 100644 --- a/services/nvpair-job-scheduler/schedule.go +++ b/services/nvpair-job-scheduler/schedule.go @@ -136,9 +136,6 @@ func (m *Manager) emitIfChanged(engine string, order []string, ranks []NodeRank, m.mu.Lock() prev := m.emitted[engine] changed := force || !equalRanks(prev.ranks, ranks) - if changed { - m.emitted[engine] = engineState{ranks: ranks, lastEmittedAt: time.Now().UnixMilli()} - } m.mu.Unlock() if !changed { @@ -149,6 +146,10 @@ func (m *Manager) emitIfChanged(engine string, order []string, ranks []NodeRank, slog.Warn("emit schedule:priority failed", "engine", engine, "err", err) return } + m.mu.Lock() + m.emitted[engine] = engineState{ranks: ranks, lastEmittedAt: time.Now().UnixMilli()} + m.mu.Unlock() + slog.Info("emitted priority", "engine", engine, "nodes", order) }