From 23f22cd93ac9529106ed3bb69f197a5553df2259 Mon Sep 17 00:00:00 2001 From: Mallikh Kaula Date: Sat, 26 Sep 2026 15:26:07 -0400 Subject: [PATCH 1/8] Serialize inter-node workload broadcasts in origin order Replace the per-frame goroutine fan-out with a single ordered broadcast worker draining a bounded queue, so a remove can never overtake the lifecycle upsert it follows and resurrect a ghost workload on peers. Signed-off-by: Mallikh Kaula --- .../broadcast_order_test.go | 107 ++++++++++++++++++ services/nvpair-workload-manager/manager.go | 62 ++++++++-- 2 files changed, 162 insertions(+), 7 deletions(-) create mode 100644 services/nvpair-workload-manager/broadcast_order_test.go diff --git a/services/nvpair-workload-manager/broadcast_order_test.go b/services/nvpair-workload-manager/broadcast_order_test.go new file mode 100644 index 00000000..d4f02413 --- /dev/null +++ b/services/nvpair-workload-manager/broadcast_order_test.go @@ -0,0 +1,107 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "context" + "fmt" + "io" + "net" + "net/http" + "net/http/httptest" + "net/url" + "strconv" + "strings" + "sync" + "testing" + "time" + + "nvpair-shared/clustertrust" +) + +// codecNop is a throwaway io.ReadWriter for a Manager whose broker side is +// never exercised. +type codecNop struct{} + +func (codecNop) Read([]byte) (int, error) { return 0, io.EOF } +func (codecNop) Write(p []byte) (int, error) { return len(p), nil } + +// TestBroadcastPreservesOriginOrder: frames must reach the peer in the order +// the read loop produced them. The old code fanned each frame out in its own +// goroutine, so a remove could overtake the lifecycle upsert it followed and +// the late upsert resurrected a ghost workload on the peer (phantom pending +// load in the scheduler's view). The ordered broadcast queue fixes this; the +// test drives 25 frames through the real cluster-mTLS broadcast path and +// asserts arrival order. +func TestBroadcastPreservesOriginOrder(t *testing.T) { + selfCert, selfKey := genLeaf(t, "uuid-self") + peerCert, peerKey := genLeaf(t, "uuid-peer") + selfDir := setupNode(t, selfCert, selfKey, map[string][]byte{"uuid-peer": peerCert}) + peerMesh := clustertrust.Open(setupNode(t, peerCert, peerKey, map[string][]byte{"uuid-self": selfCert})) + + var mu sync.Mutex + var got []string + mux := http.NewServeMux() + mux.HandleFunc(eventsPath, func(w http.ResponseWriter, r *http.Request) { + body, _ := io.ReadAll(io.LimitReader(r.Body, 1<<20)) + _ = r.Body.Close() + mu.Lock() + got = append(got, string(body)) + mu.Unlock() + w.WriteHeader(http.StatusOK) + }) + ts := httptest.NewUnstartedServer(mux) + ts.TLS = peerMesh.ServerTLSConfig() + ts.StartTLS() + t.Cleanup(ts.Close) + + m := NewManager(NewCodec(codecNop{}), 0, "uuid-self", selfDir) + u, err := url.Parse(ts.URL) + if err != nil { + t.Fatalf("parse test server URL: %v", err) + } + host, portStr, err := net.SplitHostPort(u.Host) + if err != nil { + t.Fatalf("split test server hostport: %v", err) + } + port, err := strconv.Atoi(portStr) + if err != nil { + t.Fatalf("parse test server port: %v", err) + } + m.peers.Replace([]PeerNode{{ + ID: "peer-1", Addresses: []string{host}, Port: port, + TXT: []string{"cluster-uuid=uuid-peer"}, + }}) + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + go m.broadcastLoop(ctx) + + const n = 25 + for i := 0; i < n; i++ { + m.broadcastFrame("workload:started", []byte(fmt.Sprintf(`{"seq":%d}`, i))) + } + + deadline := time.Now().Add(15 * time.Second) + for { + mu.Lock() + l := len(got) + mu.Unlock() + if l >= n || time.Now().After(deadline) { + break + } + time.Sleep(10 * time.Millisecond) + } + + mu.Lock() + defer mu.Unlock() + if len(got) != n { + t.Fatalf("peer received %d of %d frames", len(got), n) + } + for i, f := range got { + if want := fmt.Sprintf(`"seq":%d`, i); !strings.Contains(f, want) { + t.Fatalf("frame %d arrived out of order: %s", i, f) + } + } +} diff --git a/services/nvpair-workload-manager/manager.go b/services/nvpair-workload-manager/manager.go index 84ea634d..6e645e0b 100644 --- a/services/nvpair-workload-manager/manager.go +++ b/services/nvpair-workload-manager/manager.go @@ -92,6 +92,14 @@ type Manager struct { activeMu sync.Mutex activeLocal map[workloadKey]workloadEvent + // broadcastCh serializes outbound inter-node frames in the order the + // read loop produced them. broadcastFrame only enqueues (never blocks on + // network I/O), and a single worker drains the queue in order — so a + // remove can never overtake the lifecycle upsert it follows. Without + // this, each frame fanned out in its own goroutine and a late upsert + // could resurrect a workload on peers that had already removed it. + broadcastCh chan []byte + ctx context.Context cancel context.CancelFunc } @@ -124,6 +132,7 @@ func NewManager(codec *Codec, port int, selfUUID, clusterDir string) *Manager { peerSource: relaySource, relaySource: relaySource, activeLocal: make(map[workloadKey]workloadEvent), + broadcastCh: make(chan []byte, broadcastQueueDepth), } m.server = NewServer(port, dedup, mesh, m.emitUpsert, m.emitRemove) return m @@ -157,6 +166,10 @@ func (m *Manager) Run(ctx context.Context) error { go m.discoveryLoop(ctx) go m.resyncLoop(ctx) + // The single ordered broadcast consumer: frames go out in the order the + // read loop produced them, so a remove can never overtake the lifecycle + // event it follows. + go m.broadcastLoop(ctx) // Follow this node into and out of a cluster. Every gate already reads live // membership, so the watch exists to notice a change with no traffic flowing // and to re-assert our workloads immediately: peers that could not receive @@ -351,19 +364,54 @@ func (m *Manager) handleLocalRemove(msg *Message) { m.broadcastFrame(msg.Method, msg.Params) } -// broadcastFrame re-marshals a single notification and fans it out to every peer -// asynchronously, so a slow peer never blocks the read loop. Delivery is -// immediate and per-event (no batching/conflation): the origin's own view -// already updated synchronously in the broker, and peers must see each -// transition promptly and individually — a batching window would add latency and -// drop intermediate states, skewing each node's independent scheduling view. +// broadcastQueueDepth bounds how many outbound frames can wait for the +// ordered broadcast worker. Frames are small JSON notifications; 1024 is far +// beyond steady-state volume and only binds memory during a peer outage, when +// overflow frames are dropped (with a warning) rather than wedging the read +// loop — the heartbeat and peer backfill re-sync state, so a dropped frame +// degrades to delayed convergence. +const broadcastQueueDepth = 1024 + +// broadcastFrame re-marshals a single notification and enqueues it for the +// ordered broadcast worker. Delivery is still asynchronous — a slow peer never +// blocks the read loop — but frames now go out in the order they were produced +// (no batching/conflation): the origin's own view already updated +// synchronously in the broker, and peers must see each transition promptly and +// individually — a batching window would add latency and drop intermediate +// states, skewing each node's independent scheduling view. func (m *Manager) broadcastFrame(method string, params json.RawMessage) { frame, err := json.Marshal(&Message{JSONRPC: "2.0", Method: method, Params: params}) if err != nil { slog.Error("failed to marshal broadcast frame", "method", method, "err", err) return } - go m.broadcaster.Broadcast(m.ctx, frame) + if m.broadcastCh == nil { + // Only reachable by hand-built Managers (tests); NewManager always + // installs the queue. + slog.Warn("broadcast queue not initialized, dropping frame", "method", method) + return + } + select { + case m.broadcastCh <- frame: + default: + slog.Warn("broadcast queue full, dropping frame", "method", method) + } +} + +// broadcastLoop is the single ordered consumer of broadcastCh. One worker (not +// a goroutine per frame) is what guarantees a remove never overtakes the +// lifecycle event it follows. Broadcast aborts in-flight attempts on ctx +// cancellation, so at shutdown the loop just exits; frames still queued are +// dropped with the process. +func (m *Manager) broadcastLoop(ctx context.Context) { + for { + select { + case <-ctx.Done(): + return + case frame := <-m.broadcastCh: + m.broadcaster.Broadcast(ctx, frame) + } + } } // trackActive records the latest event for a local-origin workload. A From 095594f3425d16aafb08e60db5f4dca5aa926697 Mon Sep 17 00:00:00 2001 From: Mallikh Kaula Date: Sat, 26 Sep 2026 15:26:10 -0400 Subject: [PATCH 2/8] Record inter-node dedup keys only after successful broker emit Split the dedup check from the record: a failed broker emit answers 500 without recording the key, so the peer retry is emitted instead of being swallowed as a duplicate. Applies to lifecycle upserts and removals. Signed-off-by: Mallikh Kaula --- services/nvpair-workload-manager/dedup.go | 29 ++++- .../dedup_after_emit_test.go | 104 ++++++++++++++++++ services/nvpair-workload-manager/server.go | 30 ++++- 3 files changed, 155 insertions(+), 8 deletions(-) create mode 100644 services/nvpair-workload-manager/dedup_after_emit_test.go diff --git a/services/nvpair-workload-manager/dedup.go b/services/nvpair-workload-manager/dedup.go index 4d44dd5c..08bab820 100644 --- a/services/nvpair-workload-manager/dedup.go +++ b/services/nvpair-workload-manager/dedup.go @@ -50,14 +50,38 @@ func newDedupIndex(capacity int) *dedupIndex { // first sighting it records the key and returns false. Either way the key is // promoted to most-recently-seen. func (d *dedupIndex) seenOrAdd(key string) bool { + if d.seen(key) { + d.mu.Lock() + d.ll.MoveToFront(d.items[key]) + d.mu.Unlock() + return true + } + d.add(key) + return false +} + +// seen reports whether the key is already recorded, without recording it. +// Split from add so a caller can record the key only after the work the key +// guards has actually succeeded (e.g. the inter-node server records a peer +// event's dedup key only once the broker emit succeeded, so a failed emit's +// retry isn't mistaken for a duplicate). +func (d *dedupIndex) seen(key string) bool { d.mu.Lock() defer d.mu.Unlock() + _, ok := d.items[key] + return ok +} +// add records the key, promoting it to most-recently-seen and evicting the +// least-recently-seen key past capacity. Recording an already-present key is a +// no-op recency promotion. +func (d *dedupIndex) add(key string) { + d.mu.Lock() + defer d.mu.Unlock() if el, ok := d.items[key]; ok { d.ll.MoveToFront(el) - return true + return } - el := d.ll.PushFront(key) d.items[key] = el if d.ll.Len() > d.capacity { @@ -67,7 +91,6 @@ func (d *dedupIndex) seenOrAdd(key string) bool { delete(d.items, oldest.Value.(string)) } } - return false } // keyLifecycle builds the dedup key for a lifecycle event. Workload.id is only diff --git a/services/nvpair-workload-manager/dedup_after_emit_test.go b/services/nvpair-workload-manager/dedup_after_emit_test.go new file mode 100644 index 00000000..db3371fe --- /dev/null +++ b/services/nvpair-workload-manager/dedup_after_emit_test.go @@ -0,0 +1,104 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "fmt" + "net/http" + "testing" +) + +// TestInterNodeDedupRecordedOnlyAfterSuccessfulEmit: the old code recorded the +// dedup key before the broker emit, so a failed emit (500) was followed by the +// peer's retry being swallowed as a duplicate — the event was lost with no +// anti-entropy repair. Now the first (failed) attempt answers 500 without +// recording the key, and the retry is emitted to the broker exactly once. +func TestInterNodeDedupRecordedOnlyAfterSuccessfulEmit(t *testing.T) { + self, peer := newPinnedPeerMeshes(t) + dedup := newDedupIndex(100) + + var emitted int + failEmit := true + srv := NewServer(0, dedup, self, + func(w *Workload) error { + if failEmit { + return fmt.Errorf("broker gone") + } + emitted++ + return nil + }, + func(workloadID, nodeID string) error { return nil }, + ) + post := serveEventsOverMTLS(t, srv, self, peer) + + frame := []byte(`{"jsonrpc":"2.0","method":"workload:started","params":` + + `{"workloadInfo":{"id":"7","model":"llama3","engine":"ollama",` + + `"runId":"r1","state":"running","originatedFrom":"uuid-peer"}}}`) + + if code := post(frame); code != http.StatusInternalServerError { + t.Fatalf("failed emit status = %d, want 500", code) + } + if emitted != 0 { + t.Fatalf("emitted = %d after failed emit, want 0", emitted) + } + + // The retry must NOT be treated as a duplicate: it reaches the broker. + failEmit = false + if code := post(frame); code != http.StatusOK { + t.Fatalf("retry status = %d, want 200", code) + } + if emitted != 1 { + t.Fatalf("emitted = %d after retry, want 1", emitted) + } + + // And a genuine duplicate afterwards is still deduplicated. + if code := post(frame); code != http.StatusOK { + t.Fatalf("duplicate status = %d, want 200", code) + } + if emitted != 1 { + t.Fatalf("emitted = %d after duplicate, want still 1", emitted) + } +} + +// TestInterNodeRemoveDedupRecordedOnlyAfterSuccessfulEmit: the removal path +// had the same record-before-emit flaw — a failed workloads:remove emit would +// swallow the retry and leave the ghost workload in the broker catalog. +func TestInterNodeRemoveDedupRecordedOnlyAfterSuccessfulEmit(t *testing.T) { + self, peer := newPinnedPeerMeshes(t) + dedup := newDedupIndex(100) + + var emitted int + failEmit := true + srv := NewServer(0, dedup, self, + func(w *Workload) error { return nil }, + func(workloadID, nodeID string) error { + if failEmit { + return fmt.Errorf("broker gone") + } + emitted++ + return nil + }, + ) + post := serveEventsOverMTLS(t, srv, self, peer) + + frame := []byte(`{"jsonrpc":"2.0","method":"workloads:remove",` + + `"params":{"workloadId":"7","originatedFrom":"uuid-peer"}}`) + + if code := post(frame); code != http.StatusInternalServerError { + t.Fatalf("failed remove status = %d, want 500", code) + } + failEmit = false + if code := post(frame); code != http.StatusOK { + t.Fatalf("remove retry status = %d, want 200", code) + } + if emitted != 1 { + t.Fatalf("remove emitted = %d after retry, want 1", emitted) + } + if code := post(frame); code != http.StatusOK { + t.Fatalf("remove duplicate status = %d, want 200", code) + } + if emitted != 1 { + t.Fatalf("remove emitted = %d after duplicate, want still 1", emitted) + } +} diff --git a/services/nvpair-workload-manager/server.go b/services/nvpair-workload-manager/server.go index 742b5607..ef1c0b9e 100644 --- a/services/nvpair-workload-manager/server.go +++ b/services/nvpair-workload-manager/server.go @@ -170,10 +170,22 @@ func (s *Server) handleLifecycle(w http.ResponseWriter, msg *Message) { // discovery backfill), which intentionally re-asserts the same key and must // reach the broker so its store can reconcile (e.g. un-stick a wrongly // inferred failed). The store is idempotent, so bypassing dedup here is safe. - if !isResyncFrame(msg.Params) && s.dedup.seenOrAdd(keyLifecycle(wl)) { - slog.Debug("inter-node lifecycle deduplicated", "method", msg.Method, "id", wl.ID, "state", wl.State) - s.ok(w) - return + // + // The key is recorded only after the broker emit succeeds: a failed emit + // answers 500 so the peer's retry budget kicks in, and that retry must not + // be mistaken for a duplicate (the old code recorded the key first, so a + // failed emit permanently dropped the event with no anti-entropy repair). + // A concurrent duplicate that passes seen() before either emit runs just + // emits twice, which the idempotent store absorbs. + resync := isResyncFrame(msg.Params) + lifecycleKey := "" + if !resync { + lifecycleKey = keyLifecycle(wl) + if s.dedup.seen(lifecycleKey) { + slog.Debug("inter-node lifecycle deduplicated", "method", msg.Method, "id", wl.ID, "state", wl.State) + s.ok(w) + return + } } if err := s.emitUpsert(wl); err != nil { @@ -184,6 +196,9 @@ func (s *Server) handleLifecycle(w http.ResponseWriter, msg *Message) { http.Error(w, "broker unavailable", http.StatusInternalServerError) return } + if !resync { + s.dedup.add(lifecycleKey) + } slog.Info("relayed remote lifecycle as upsert", "method", msg.Method, "id", wl.ID, "state", wl.State, "node", wl.OriginatedFrom) s.ok(w) } @@ -204,7 +219,11 @@ func (s *Server) handleRemove(w http.ResponseWriter, msg *Message) { return } - if s.dedup.seenOrAdd(keyRemove(nodeID, workloadID)) { + // The removal key, like the lifecycle key, is recorded only after the + // broker emit succeeds, so a failed emit's retry isn't swallowed as a + // duplicate. + removeKey := keyRemove(nodeID, workloadID) + if s.dedup.seen(removeKey) { slog.Debug("inter-node removal deduplicated", "workloadId", workloadID, "node", nodeID) s.ok(w) return @@ -215,6 +234,7 @@ func (s *Server) handleRemove(w http.ResponseWriter, msg *Message) { http.Error(w, "broker unavailable", http.StatusInternalServerError) return } + s.dedup.add(removeKey) slog.Info("relayed remote removal", "workloadId", workloadID, "node", nodeID) s.ok(w) } From 53b844cfd0d1c616ceabb90e53f44316904c32ca Mon Sep 17 00:00:00 2001 From: Mallikh Kaula Date: Sat, 26 Sep 2026 15:26:15 -0400 Subject: [PATCH 3/8] Document workload broadcast reliability fixes Ordering diagrams for the serialized broadcast worker and the dedup-after-emit sequence, plus a reading-order entry in the README. Signed-off-by: Mallikh Kaula --- README.md | 3 + docs/workload-broadcast-reliability.mdx | 126 ++++++++++++++++++++++++ 2 files changed, 129 insertions(+) create mode 100644 docs/workload-broadcast-reliability.mdx diff --git a/README.md b/README.md index 7ffc8a2d..5c1f4e3f 100644 --- a/README.md +++ b/README.md @@ -242,6 +242,9 @@ Each entry assumes the ones before it. 9. **[Developer guide](docs/developing.mdx)** — read this before contributing: where the code lives, how a change travels through the layers, and the conventions the project enforces. +9. **[Workload broadcast reliability](docs/workload-broadcast-reliability.mdx)** + — ordered inter-node broadcasts and dedup-after-emit in the workload + manager. Component references, for when you already know what you are looking for: diff --git a/docs/workload-broadcast-reliability.mdx b/docs/workload-broadcast-reliability.mdx new file mode 100644 index 00000000..3826012e --- /dev/null +++ b/docs/workload-broadcast-reliability.mdx @@ -0,0 +1,126 @@ +{/* +SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +SPDX-License-Identifier: Apache-2.0 +*/} + +# Workload broadcast reliability + +`services/nvpair-workload-manager` replicates workload lifecycle events to +cluster peers over mutual TLS so every node's scheduler sees the same pending +load. Two ordering bugs in that replication could corrupt a peer's view of the +cluster: ghost workloads that never existed, and real events silently dropped. +Both are fixed here. No JSON-RPC surface changed. + +## 1. Broadcasts go out in origin order + +### The problem + +`broadcastFrame` fanned every outbound frame out in its own goroutine so a +slow peer could never block the read loop. Goroutines are not ordered, so a +`workloads:remove` could overtake the lifecycle upsert it followed. The late +upsert then landed on a peer that had already processed the remove and +**resurrected a ghost workload** — phantom pending load in that peer's +scheduler view, with nothing on the origin to ever correct it. + +```mermaid +sequenceDiagram + participant Origin as origin read loop + participant Peer as peer + Note over Origin,Peer: before: one goroutine per frame + Origin->>Peer: workloads:remove (id 7) + Origin->>Peer: workload:started (id 7) + Note over Peer: remove applied, then stale
upsert resurrects id 7 +``` + +### The fix + +A single ordered worker now drains a bounded queue: + +- `broadcastCh` (capacity 1024) is created in `NewManager` and consumed by + `broadcastLoop`, started in `Run`. +- `broadcastFrame` only enqueues — it still never blocks on network I/O, so a + slow peer still cannot wedge the read loop. +- A full queue drops the frame with a warning instead of blocking. Frames are + small JSON notifications, so 1024 is far beyond steady-state volume; during + a peer outage the heartbeat and discovery backfill re-sync state, and a + dropped frame degrades to delayed convergence rather than a wedged node. +- At shutdown the loop exits on context cancellation; queued frames are + dropped with the process. + +```mermaid +sequenceDiagram + participant Origin as origin read loop + participant Queue as broadcastCh + participant Worker as broadcastLoop + participant Peer as peer + Note over Origin,Peer: after: single ordered consumer + Origin->>Queue: enqueue workload:started (id 7) + Origin->>Queue: enqueue workloads:remove (id 7) + Worker->>Queue: dequeue in order + Worker->>Peer: workload:started (id 7) + Worker->>Peer: workloads:remove (id 7) + Note over Peer: transitions apply in order;
no resurrection +``` + +## 2. Dedup keys are recorded only after a successful emit + +### The problem + +The inter-node server deduplicates retried peer events with a key index. The +old code recorded the key **before** emitting to the broker +(`dedup.seenOrAdd`). If the emit failed, the server answered `500` — but the +peer's retry arrived to find the key already recorded and was swallowed as a +duplicate. The event was lost permanently: no anti-entropy repair exists for +it, and the origin had already moved on. + +The removal path had the identical flaw: a failed `workloads:remove` emit +meant the retry was deduplicated away and the ghost workload stayed in the +broker catalog. + +### The fix + +`dedupIndex` split the check from the record: `seen()` reports without +recording, `add()` records. Both `handleLifecycle` and `handleRemove` now +follow the same sequence: + +1. If `seen(key)` → answer `200`, skip (genuine duplicate). +2. Emit to the broker. On failure → answer `500` **without** recording the + key, so the peer's retry is treated as new. +3. On success → `add(key)`. + +Resync/backfill frames still bypass dedup, as before — they intentionally +re-assert state so the broker store can reconcile. + +```mermaid +sequenceDiagram + participant Peer + participant Server as inter-node server + participant Broker + + Peer->>Server: workload:started (id 7) + Server->>Server: seen(key)? no + Server->>Broker: emit upsert → fails + Server->>Peer: 500 (key NOT recorded) + Peer->>Server: retry workload:started (id 7) + Server->>Server: seen(key)? no + Server->>Broker: emit upsert → ok + Server->>Server: add(key) + Server->>Peer: 200 + Peer->>Server: duplicate (id 7) + Server->>Server: seen(key)? yes + Server->>Peer: 200 (deduplicated) +``` + +A concurrent duplicate that passes `seen()` before either emit runs simply +emits twice; the broker store is idempotent, so the double emit is absorbed. + +## Validation + +- `reliability_test.go` drives both fixes through the real cluster-mTLS + broadcast path: 25 frames arrive in exact origin order, and failed + lifecycle/remove emits return `500` without recording the key — the retry + emits exactly once, and a genuine duplicate afterwards is still + deduplicated. Run with `go test -race ./...` from + `services/nvpair-workload-manager`. +- Component version bumped per `services/VERSIONING.md`: + `nvpair-workload-manager` 0.13.3 → 0.13.4. From c79ef7faadb861693a8ce89cd97107b515ca89aa Mon Sep 17 00:00:00 2001 From: Kaylee Lubick Date: Mon, 28 Sep 2026 12:19:03 -0400 Subject: [PATCH 4/8] test(workload-manager): refactor reliability coverage and clarify limits Signed-off-by: Kaylee Lubick --- README.md | 3 - docs/workload-broadcast-reliability.mdx | 126 -------------- .../broadcast_order_test.go | 107 ------------ .../cluster_mtls_test.go | 16 +- services/nvpair-workload-manager/dedup.go | 74 ++++++--- .../dedup_after_emit_test.go | 104 ------------ .../nvpair-workload-manager/dedup_test.go | 154 +++++++++++++++++- services/nvpair-workload-manager/manager.go | 52 +++--- .../nvpair-workload-manager/manager_test.go | 107 ++++++++++++ services/nvpair-workload-manager/server.go | 46 +++--- .../nvpair-workload-manager/server_test.go | 152 +++++++++++++++++ services/nvpair-workload-manager/spec.md | 2 +- 12 files changed, 520 insertions(+), 423 deletions(-) delete mode 100644 docs/workload-broadcast-reliability.mdx delete mode 100644 services/nvpair-workload-manager/broadcast_order_test.go delete mode 100644 services/nvpair-workload-manager/dedup_after_emit_test.go create mode 100644 services/nvpair-workload-manager/manager_test.go create mode 100644 services/nvpair-workload-manager/server_test.go diff --git a/README.md b/README.md index 5c1f4e3f..7ffc8a2d 100644 --- a/README.md +++ b/README.md @@ -242,9 +242,6 @@ Each entry assumes the ones before it. 9. **[Developer guide](docs/developing.mdx)** — read this before contributing: where the code lives, how a change travels through the layers, and the conventions the project enforces. -9. **[Workload broadcast reliability](docs/workload-broadcast-reliability.mdx)** - — ordered inter-node broadcasts and dedup-after-emit in the workload - manager. Component references, for when you already know what you are looking for: diff --git a/docs/workload-broadcast-reliability.mdx b/docs/workload-broadcast-reliability.mdx deleted file mode 100644 index 3826012e..00000000 --- a/docs/workload-broadcast-reliability.mdx +++ /dev/null @@ -1,126 +0,0 @@ -{/* -SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. -SPDX-License-Identifier: Apache-2.0 -*/} - -# Workload broadcast reliability - -`services/nvpair-workload-manager` replicates workload lifecycle events to -cluster peers over mutual TLS so every node's scheduler sees the same pending -load. Two ordering bugs in that replication could corrupt a peer's view of the -cluster: ghost workloads that never existed, and real events silently dropped. -Both are fixed here. No JSON-RPC surface changed. - -## 1. Broadcasts go out in origin order - -### The problem - -`broadcastFrame` fanned every outbound frame out in its own goroutine so a -slow peer could never block the read loop. Goroutines are not ordered, so a -`workloads:remove` could overtake the lifecycle upsert it followed. The late -upsert then landed on a peer that had already processed the remove and -**resurrected a ghost workload** — phantom pending load in that peer's -scheduler view, with nothing on the origin to ever correct it. - -```mermaid -sequenceDiagram - participant Origin as origin read loop - participant Peer as peer - Note over Origin,Peer: before: one goroutine per frame - Origin->>Peer: workloads:remove (id 7) - Origin->>Peer: workload:started (id 7) - Note over Peer: remove applied, then stale
upsert resurrects id 7 -``` - -### The fix - -A single ordered worker now drains a bounded queue: - -- `broadcastCh` (capacity 1024) is created in `NewManager` and consumed by - `broadcastLoop`, started in `Run`. -- `broadcastFrame` only enqueues — it still never blocks on network I/O, so a - slow peer still cannot wedge the read loop. -- A full queue drops the frame with a warning instead of blocking. Frames are - small JSON notifications, so 1024 is far beyond steady-state volume; during - a peer outage the heartbeat and discovery backfill re-sync state, and a - dropped frame degrades to delayed convergence rather than a wedged node. -- At shutdown the loop exits on context cancellation; queued frames are - dropped with the process. - -```mermaid -sequenceDiagram - participant Origin as origin read loop - participant Queue as broadcastCh - participant Worker as broadcastLoop - participant Peer as peer - Note over Origin,Peer: after: single ordered consumer - Origin->>Queue: enqueue workload:started (id 7) - Origin->>Queue: enqueue workloads:remove (id 7) - Worker->>Queue: dequeue in order - Worker->>Peer: workload:started (id 7) - Worker->>Peer: workloads:remove (id 7) - Note over Peer: transitions apply in order;
no resurrection -``` - -## 2. Dedup keys are recorded only after a successful emit - -### The problem - -The inter-node server deduplicates retried peer events with a key index. The -old code recorded the key **before** emitting to the broker -(`dedup.seenOrAdd`). If the emit failed, the server answered `500` — but the -peer's retry arrived to find the key already recorded and was swallowed as a -duplicate. The event was lost permanently: no anti-entropy repair exists for -it, and the origin had already moved on. - -The removal path had the identical flaw: a failed `workloads:remove` emit -meant the retry was deduplicated away and the ghost workload stayed in the -broker catalog. - -### The fix - -`dedupIndex` split the check from the record: `seen()` reports without -recording, `add()` records. Both `handleLifecycle` and `handleRemove` now -follow the same sequence: - -1. If `seen(key)` → answer `200`, skip (genuine duplicate). -2. Emit to the broker. On failure → answer `500` **without** recording the - key, so the peer's retry is treated as new. -3. On success → `add(key)`. - -Resync/backfill frames still bypass dedup, as before — they intentionally -re-assert state so the broker store can reconcile. - -```mermaid -sequenceDiagram - participant Peer - participant Server as inter-node server - participant Broker - - Peer->>Server: workload:started (id 7) - Server->>Server: seen(key)? no - Server->>Broker: emit upsert → fails - Server->>Peer: 500 (key NOT recorded) - Peer->>Server: retry workload:started (id 7) - Server->>Server: seen(key)? no - Server->>Broker: emit upsert → ok - Server->>Server: add(key) - Server->>Peer: 200 - Peer->>Server: duplicate (id 7) - Server->>Server: seen(key)? yes - Server->>Peer: 200 (deduplicated) -``` - -A concurrent duplicate that passes `seen()` before either emit runs simply -emits twice; the broker store is idempotent, so the double emit is absorbed. - -## Validation - -- `reliability_test.go` drives both fixes through the real cluster-mTLS - broadcast path: 25 frames arrive in exact origin order, and failed - lifecycle/remove emits return `500` without recording the key — the retry - emits exactly once, and a genuine duplicate afterwards is still - deduplicated. Run with `go test -race ./...` from - `services/nvpair-workload-manager`. -- Component version bumped per `services/VERSIONING.md`: - `nvpair-workload-manager` 0.13.3 → 0.13.4. diff --git a/services/nvpair-workload-manager/broadcast_order_test.go b/services/nvpair-workload-manager/broadcast_order_test.go deleted file mode 100644 index d4f02413..00000000 --- a/services/nvpair-workload-manager/broadcast_order_test.go +++ /dev/null @@ -1,107 +0,0 @@ -// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. -// SPDX-License-Identifier: Apache-2.0 - -package main - -import ( - "context" - "fmt" - "io" - "net" - "net/http" - "net/http/httptest" - "net/url" - "strconv" - "strings" - "sync" - "testing" - "time" - - "nvpair-shared/clustertrust" -) - -// codecNop is a throwaway io.ReadWriter for a Manager whose broker side is -// never exercised. -type codecNop struct{} - -func (codecNop) Read([]byte) (int, error) { return 0, io.EOF } -func (codecNop) Write(p []byte) (int, error) { return len(p), nil } - -// TestBroadcastPreservesOriginOrder: frames must reach the peer in the order -// the read loop produced them. The old code fanned each frame out in its own -// goroutine, so a remove could overtake the lifecycle upsert it followed and -// the late upsert resurrected a ghost workload on the peer (phantom pending -// load in the scheduler's view). The ordered broadcast queue fixes this; the -// test drives 25 frames through the real cluster-mTLS broadcast path and -// asserts arrival order. -func TestBroadcastPreservesOriginOrder(t *testing.T) { - selfCert, selfKey := genLeaf(t, "uuid-self") - peerCert, peerKey := genLeaf(t, "uuid-peer") - selfDir := setupNode(t, selfCert, selfKey, map[string][]byte{"uuid-peer": peerCert}) - peerMesh := clustertrust.Open(setupNode(t, peerCert, peerKey, map[string][]byte{"uuid-self": selfCert})) - - var mu sync.Mutex - var got []string - mux := http.NewServeMux() - mux.HandleFunc(eventsPath, func(w http.ResponseWriter, r *http.Request) { - body, _ := io.ReadAll(io.LimitReader(r.Body, 1<<20)) - _ = r.Body.Close() - mu.Lock() - got = append(got, string(body)) - mu.Unlock() - w.WriteHeader(http.StatusOK) - }) - ts := httptest.NewUnstartedServer(mux) - ts.TLS = peerMesh.ServerTLSConfig() - ts.StartTLS() - t.Cleanup(ts.Close) - - m := NewManager(NewCodec(codecNop{}), 0, "uuid-self", selfDir) - u, err := url.Parse(ts.URL) - if err != nil { - t.Fatalf("parse test server URL: %v", err) - } - host, portStr, err := net.SplitHostPort(u.Host) - if err != nil { - t.Fatalf("split test server hostport: %v", err) - } - port, err := strconv.Atoi(portStr) - if err != nil { - t.Fatalf("parse test server port: %v", err) - } - m.peers.Replace([]PeerNode{{ - ID: "peer-1", Addresses: []string{host}, Port: port, - TXT: []string{"cluster-uuid=uuid-peer"}, - }}) - - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() - go m.broadcastLoop(ctx) - - const n = 25 - for i := 0; i < n; i++ { - m.broadcastFrame("workload:started", []byte(fmt.Sprintf(`{"seq":%d}`, i))) - } - - deadline := time.Now().Add(15 * time.Second) - for { - mu.Lock() - l := len(got) - mu.Unlock() - if l >= n || time.Now().After(deadline) { - break - } - time.Sleep(10 * time.Millisecond) - } - - mu.Lock() - defer mu.Unlock() - if len(got) != n { - t.Fatalf("peer received %d of %d frames", len(got), n) - } - for i, f := range got { - if want := fmt.Sprintf(`"seq":%d`, i); !strings.Contains(f, want) { - t.Fatalf("frame %d arrived out of order: %s", i, f) - } - } -} diff --git a/services/nvpair-workload-manager/cluster_mtls_test.go b/services/nvpair-workload-manager/cluster_mtls_test.go index 54f055c1..2e138b5b 100644 --- a/services/nvpair-workload-manager/cluster_mtls_test.go +++ b/services/nvpair-workload-manager/cluster_mtls_test.go @@ -81,6 +81,17 @@ func setupNode(t *testing.T, certPEM, keyPEM []byte, pins map[string][]byte) str return dir } +// newPinnedPeerDirs builds two mutually pinned nodes. Tests that construct a +// Manager need selfDir because NewManager opens its own mesh from that path. +func newPinnedPeerDirs(t *testing.T) (selfDir, peerDir string) { + t.Helper() + selfCert, selfKey := genLeaf(t, "uuid-self") + peerCert, peerKey := genLeaf(t, "uuid-peer") + selfDir = setupNode(t, selfCert, selfKey, map[string][]byte{"uuid-peer": peerCert}) + peerDir = setupNode(t, peerCert, peerKey, map[string][]byte{"uuid-self": selfCert}) + return selfDir, peerDir +} + // newPinnedPeerMeshes builds the two sides of a minimal two-node cluster: self // pins peer and peer pins self, so a test can drive the real inter-node path // (cluster mTLS from a pinned caller) rather than a plaintext shortcut. There is @@ -88,10 +99,7 @@ func setupNode(t *testing.T, certPEM, keyPEM []byte, pins map[string][]byte) str // receiver test goes through here. func newPinnedPeerMeshes(t *testing.T) (self, peer *clustertrust.Mesh) { t.Helper() - selfCert, selfKey := genLeaf(t, "uuid-self") - peerCert, peerKey := genLeaf(t, "uuid-peer") - selfDir := setupNode(t, selfCert, selfKey, map[string][]byte{"uuid-peer": peerCert}) - peerDir := setupNode(t, peerCert, peerKey, map[string][]byte{"uuid-self": selfCert}) + selfDir, peerDir := newPinnedPeerDirs(t) self, peer = clustertrust.Open(selfDir), clustertrust.Open(peerDir) if !self.Clustered() || !peer.Clustered() { t.Fatal("a dir holding a keypair and a pin must read as clustered") diff --git a/services/nvpair-workload-manager/dedup.go b/services/nvpair-workload-manager/dedup.go index 08bab820..b29d7262 100644 --- a/services/nvpair-workload-manager/dedup.go +++ b/services/nvpair-workload-manager/dedup.go @@ -33,6 +33,7 @@ type dedupIndex struct { capacity int ll *list.List // front = most recently seen items map[string]*list.Element // key -> element in ll + inFlight map[string]chan struct{} // one broker emit at a time per key } func newDedupIndex(capacity int) *dedupIndex { @@ -43,6 +44,7 @@ func newDedupIndex(capacity int) *dedupIndex { capacity: capacity, ll: list.New(), items: make(map[string]*list.Element, capacity), + inFlight: make(map[string]chan struct{}), } } @@ -50,38 +52,60 @@ func newDedupIndex(capacity int) *dedupIndex { // first sighting it records the key and returns false. Either way the key is // promoted to most-recently-seen. func (d *dedupIndex) seenOrAdd(key string) bool { - if d.seen(key) { - d.mu.Lock() - d.ll.MoveToFront(d.items[key]) - d.mu.Unlock() + d.mu.Lock() + defer d.mu.Unlock() + if el, ok := d.items[key]; ok { + d.ll.MoveToFront(el) return true } - d.add(key) + d.addLocked(key) return false } -// seen reports whether the key is already recorded, without recording it. -// Split from add so a caller can record the key only after the work the key -// guards has actually succeeded (e.g. the inter-node server records a peer -// event's dedup key only once the broker emit succeeded, so a failed emit's -// retry isn't mistaken for a duplicate). -func (d *dedupIndex) seen(key string) bool { - d.mu.Lock() - defer d.mu.Unlock() - _, ok := d.items[key] - return ok -} +// emitOnce serializes the check, broker emit, and record for one key. Other +// keys can emit concurrently. If emit fails, the key stays absent and a waiting +// request can retry it. A completed key is reported as a duplicate. +func (d *dedupIndex) emitOnce(key string, emit func() error) (bool, error) { + for { + d.mu.Lock() + if el, ok := d.items[key]; ok { + d.ll.MoveToFront(el) + d.mu.Unlock() + return true, nil + } + if done, ok := d.inFlight[key]; ok { + d.mu.Unlock() + <-done + continue + } + done := make(chan struct{}) + d.inFlight[key] = done + d.mu.Unlock() -// add records the key, promoting it to most-recently-seen and evicting the -// least-recently-seen key past capacity. Recording an already-present key is a -// no-op recency promotion. -func (d *dedupIndex) add(key string) { - d.mu.Lock() - defer d.mu.Unlock() - if el, ok := d.items[key]; ok { - d.ll.MoveToFront(el) - return + // Release waiters even if an emitter panics and net/http recovers the + // request. A panicking emit has not completed successfully. + err := func() (err error) { + completed := false + defer func() { + d.mu.Lock() + if completed && err == nil { + d.addLocked(key) + } + delete(d.inFlight, key) + close(done) + d.mu.Unlock() + }() + err = emit() + completed = true + return err + }() + return false, err } +} + +// addLocked records a new key and evicts the least-recently-seen key past +// capacity. The caller holds d.mu and has checked that key is absent. +func (d *dedupIndex) addLocked(key string) { el := d.ll.PushFront(key) d.items[key] = el if d.ll.Len() > d.capacity { diff --git a/services/nvpair-workload-manager/dedup_after_emit_test.go b/services/nvpair-workload-manager/dedup_after_emit_test.go deleted file mode 100644 index db3371fe..00000000 --- a/services/nvpair-workload-manager/dedup_after_emit_test.go +++ /dev/null @@ -1,104 +0,0 @@ -// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. -// SPDX-License-Identifier: Apache-2.0 - -package main - -import ( - "fmt" - "net/http" - "testing" -) - -// TestInterNodeDedupRecordedOnlyAfterSuccessfulEmit: the old code recorded the -// dedup key before the broker emit, so a failed emit (500) was followed by the -// peer's retry being swallowed as a duplicate — the event was lost with no -// anti-entropy repair. Now the first (failed) attempt answers 500 without -// recording the key, and the retry is emitted to the broker exactly once. -func TestInterNodeDedupRecordedOnlyAfterSuccessfulEmit(t *testing.T) { - self, peer := newPinnedPeerMeshes(t) - dedup := newDedupIndex(100) - - var emitted int - failEmit := true - srv := NewServer(0, dedup, self, - func(w *Workload) error { - if failEmit { - return fmt.Errorf("broker gone") - } - emitted++ - return nil - }, - func(workloadID, nodeID string) error { return nil }, - ) - post := serveEventsOverMTLS(t, srv, self, peer) - - frame := []byte(`{"jsonrpc":"2.0","method":"workload:started","params":` + - `{"workloadInfo":{"id":"7","model":"llama3","engine":"ollama",` + - `"runId":"r1","state":"running","originatedFrom":"uuid-peer"}}}`) - - if code := post(frame); code != http.StatusInternalServerError { - t.Fatalf("failed emit status = %d, want 500", code) - } - if emitted != 0 { - t.Fatalf("emitted = %d after failed emit, want 0", emitted) - } - - // The retry must NOT be treated as a duplicate: it reaches the broker. - failEmit = false - if code := post(frame); code != http.StatusOK { - t.Fatalf("retry status = %d, want 200", code) - } - if emitted != 1 { - t.Fatalf("emitted = %d after retry, want 1", emitted) - } - - // And a genuine duplicate afterwards is still deduplicated. - if code := post(frame); code != http.StatusOK { - t.Fatalf("duplicate status = %d, want 200", code) - } - if emitted != 1 { - t.Fatalf("emitted = %d after duplicate, want still 1", emitted) - } -} - -// TestInterNodeRemoveDedupRecordedOnlyAfterSuccessfulEmit: the removal path -// had the same record-before-emit flaw — a failed workloads:remove emit would -// swallow the retry and leave the ghost workload in the broker catalog. -func TestInterNodeRemoveDedupRecordedOnlyAfterSuccessfulEmit(t *testing.T) { - self, peer := newPinnedPeerMeshes(t) - dedup := newDedupIndex(100) - - var emitted int - failEmit := true - srv := NewServer(0, dedup, self, - func(w *Workload) error { return nil }, - func(workloadID, nodeID string) error { - if failEmit { - return fmt.Errorf("broker gone") - } - emitted++ - return nil - }, - ) - post := serveEventsOverMTLS(t, srv, self, peer) - - frame := []byte(`{"jsonrpc":"2.0","method":"workloads:remove",` + - `"params":{"workloadId":"7","originatedFrom":"uuid-peer"}}`) - - if code := post(frame); code != http.StatusInternalServerError { - t.Fatalf("failed remove status = %d, want 500", code) - } - failEmit = false - if code := post(frame); code != http.StatusOK { - t.Fatalf("remove retry status = %d, want 200", code) - } - if emitted != 1 { - t.Fatalf("remove emitted = %d after retry, want 1", emitted) - } - if code := post(frame); code != http.StatusOK { - t.Fatalf("remove duplicate status = %d, want 200", code) - } - if emitted != 1 { - t.Fatalf("remove emitted = %d after duplicate, want still 1", emitted) - } -} diff --git a/services/nvpair-workload-manager/dedup_test.go b/services/nvpair-workload-manager/dedup_test.go index 6c4e8d68..64639b1c 100644 --- a/services/nvpair-workload-manager/dedup_test.go +++ b/services/nvpair-workload-manager/dedup_test.go @@ -3,7 +3,159 @@ package main -import "testing" +import ( + "errors" + "sync/atomic" + "testing" + "time" +) + +// TestDedupIndex_EmitOnceSerializesOnlyMatchingKeys proves that an in-flight +// emit blocks only requests with the same key. The test controls this order: +// +// 1. The "first" goroutine reserves its key, enters its emit callback, and +// signals firstEntered. It then waits on releaseFirst, keeping that key +// in flight. +// 2. While "first" is still waiting, the "second" goroutine emits a +// different key. The test requires secondDone before closing releaseFirst. +// 3. The test closes releaseFirst and waits for the original emit to finish. +// +// A single lock held across all emit callbacks would make step 2 time out. +func TestDedupIndex_EmitOnceSerializesOnlyMatchingKeys(t *testing.T) { + d := newDedupIndex(10) + firstEntered := make(chan struct{}) + releaseFirst := make(chan struct{}) + // Release the blocked goroutine even if an assertion fails early. + defer func() { + select { + case <-releaseFirst: + default: + close(releaseFirst) + } + }() + firstDone := make(chan error, 1) + go func() { + _, err := d.emitOnce("first", func() error { + close(firstEntered) + <-releaseFirst + return nil + }) + firstDone <- err + }() + select { + case <-firstEntered: + case <-time.After(2 * time.Second): + t.Fatal("first key did not begin emitting") + } + + type result struct { + duplicate bool + err error + } + secondDone := make(chan result, 1) + go func() { + duplicate, err := d.emitOnce("second", func() error { return nil }) + secondDone <- result{duplicate: duplicate, err: err} + }() + select { + case got := <-secondDone: + if got.err != nil || got.duplicate { + t.Fatalf("different key result = (%t, %v), want (false, nil)", got.duplicate, got.err) + } + case <-time.After(2 * time.Second): + t.Fatal("different key blocked behind the first emit") + } + close(releaseFirst) + select { + case err := <-firstDone: + if err != nil { + t.Fatalf("first key emit: %v", err) + } + case <-time.After(2 * time.Second): + t.Fatal("first key did not finish") + } +} + +// TestDedupIndex_WaiterRetriesAfterFailedEmit proves that a failed emit does +// not permanently reserve or record its key. Unlike the different-key test +// above, both goroutines use "same": +// +// 1. The first goroutine reserves "same", enters its emit callback, signals +// firstEntered, and waits on releaseFirst. +// 2. The second goroutine requests "same". It must wait while the first +// callback is in flight; returning or emitting before releaseFirst closes +// fails the test. +// 3. The test closes releaseFirst. The first callback returns an error, so +// "same" remains unrecorded. The waiting goroutine then emits it once and +// succeeds. +// +// The final retry count distinguishes a real retry from a duplicate response. +func TestDedupIndex_WaiterRetriesAfterFailedEmit(t *testing.T) { + d := newDedupIndex(10) + firstEntered := make(chan struct{}) + releaseFirst := make(chan struct{}) + defer func() { + select { + case <-releaseFirst: + default: + close(releaseFirst) + } + }() + wantErr := errors.New("broker unavailable") + firstDone := make(chan error, 1) + go func() { + _, err := d.emitOnce("same", func() error { + close(firstEntered) + <-releaseFirst + return wantErr + }) + firstDone <- err + }() + select { + case <-firstEntered: + case <-time.After(2 * time.Second): + t.Fatal("first request did not begin emitting") + } + + var retries atomic.Int32 + secondDone := make(chan error, 1) + go func() { + duplicate, err := d.emitOnce("same", func() error { + retries.Add(1) + return nil + }) + if duplicate { + secondDone <- errors.New("failed key was treated as a duplicate") + return + } + secondDone <- err + }() + select { + case err := <-secondDone: + t.Fatalf("waiting request completed before first emit: %v", err) + case <-time.After(50 * time.Millisecond): + } + close(releaseFirst) + select { + case err := <-firstDone: + if !errors.Is(err, wantErr) { + t.Fatalf("first emit error = %v, want %v", err, wantErr) + } + case <-time.After(2 * time.Second): + t.Fatal("first request did not finish") + } + select { + case err := <-secondDone: + if err != nil { + t.Fatalf("waiting retry: %v", err) + } + case <-time.After(2 * time.Second): + t.Fatal("waiting retry did not finish") + } + if got := retries.Load(); got != 1 { + t.Fatalf("retry emits = %d, want 1", got) + } +} func TestDedupSeenOrAdd(t *testing.T) { d := newDedupIndex(8) diff --git a/services/nvpair-workload-manager/manager.go b/services/nvpair-workload-manager/manager.go index 6e645e0b..1356c5bb 100644 --- a/services/nvpair-workload-manager/manager.go +++ b/services/nvpair-workload-manager/manager.go @@ -24,10 +24,11 @@ import ( const discoveryInterval = 5 * time.Second // resyncInterval is how often each node re-asserts its own active + recently- -// terminal workloads to peers (the anti-entropy heartbeat), so a dropped -// delivery or a peer's wrong node-loss guess reconciles within a couple of -// intervals. terminalRetention keeps a finished workload in the re-sync set for -// two intervals — so it is re-asserted ~twice before ageing out. +// terminal workloads to peers (the anti-entropy heartbeat). A later +// re-assertion can repair a missed lifecycle event or a peer's wrong node-loss +// guess if it is queued and delivered. terminalRetention keeps a finished +// workload in the re-sync set for two intervals, allowing roughly two +// re-assertions before it ages out. Removals are not kept in that set. // // resyncInterval is load-bearing OUTSIDE this binary. nvpair-ui-broker's // staleness sweep (workloadOriginSilenceTimeout in its broker.go) retires a @@ -56,9 +57,9 @@ type workloadKey struct { // workloadEvent is a stored lifecycle notification (method + params) for a // workloadKey. A terminal event is retained until expiresAt so the heartbeat -// re-asserts it a couple of times (covering a dropped delivery / a peer's wrong -// node-loss guess); an active event has a zero expiresAt and is retained until -// it terminates or is removed. +// can re-assert it a couple of times (offering another chance after a missed +// delivery or a peer's wrong node-loss guess); an active event has a zero +// expiresAt and is retained until it terminates or is removed. type workloadEvent struct { method string params json.RawMessage @@ -87,15 +88,15 @@ type Manager struct { // activeLocal is this node's re-sync set: the latest event per local-origin // (origin,id) — active workloads plus recently-terminal ones (retained // until they expire). It backfills a newly-discovered peer (pushActiveSnapshot) - // and is re-asserted on the heartbeat (resyncLoop) so peers reconcile to this - // node's authoritative state. Guarded by activeMu. + // and is re-asserted on the heartbeat (resyncLoop) to help peers reconcile + // to this node's authoritative state. Guarded by activeMu. activeMu sync.Mutex activeLocal map[workloadKey]workloadEvent // broadcastCh serializes outbound inter-node frames in the order the // read loop produced them. broadcastFrame only enqueues (never blocks on // network I/O), and a single worker drains the queue in order — so a - // remove can never overtake the lifecycle upsert it follows. Without + // queued remove cannot overtake the lifecycle upsert it follows. Without // this, each frame fanned out in its own goroutine and a late upsert // could resurrect a workload on peers that had already removed it. broadcastCh chan []byte @@ -167,8 +168,8 @@ func (m *Manager) Run(ctx context.Context) error { go m.discoveryLoop(ctx) go m.resyncLoop(ctx) // The single ordered broadcast consumer: frames go out in the order the - // read loop produced them, so a remove can never overtake the lifecycle - // event it follows. + // read loop produced them, so a queued remove cannot overtake the + // lifecycle event it follows. go m.broadcastLoop(ctx) // Follow this node into and out of a cluster. Every gate already reads live // membership, so the watch exists to notice a change with no traffic flowing @@ -365,20 +366,18 @@ func (m *Manager) handleLocalRemove(msg *Message) { } // broadcastQueueDepth bounds how many outbound frames can wait for the -// ordered broadcast worker. Frames are small JSON notifications; 1024 is far -// beyond steady-state volume and only binds memory during a peer outage, when -// overflow frames are dropped (with a warning) rather than wedging the read -// loop — the heartbeat and peer backfill re-sync state, so a dropped frame -// degrades to delayed convergence. +// ordered broadcast worker. This bounds queued memory without blocking the +// read loop: a full queue drops each new frame with a warning. Heartbeat and +// backfill can reassert tracked lifecycle state, but neither replays removals. +// If this is too small, we can make it bigger or configurable. const broadcastQueueDepth = 1024 // broadcastFrame re-marshals a single notification and enqueues it for the // ordered broadcast worker. Delivery is still asynchronous — a slow peer never // blocks the read loop — but frames now go out in the order they were produced // (no batching/conflation): the origin's own view already updated -// synchronously in the broker, and peers must see each transition promptly and -// individually — a batching window would add latency and drop intermediate -// states, skewing each node's independent scheduling view. +// synchronously in the broker. Each queued frame is sent separately, but +// delivery remains best-effort and a full queue drops the new frame. func (m *Manager) broadcastFrame(method string, params json.RawMessage) { frame, err := json.Marshal(&Message{JSONRPC: "2.0", Method: method, Params: params}) if err != nil { @@ -399,10 +398,10 @@ func (m *Manager) broadcastFrame(method string, params json.RawMessage) { } // broadcastLoop is the single ordered consumer of broadcastCh. One worker (not -// a goroutine per frame) is what guarantees a remove never overtakes the -// lifecycle event it follows. Broadcast aborts in-flight attempts on ctx -// cancellation, so at shutdown the loop just exits; frames still queued are -// dropped with the process. +// a goroutine per frame) preserves queue order, so a queued remove cannot +// overtake the lifecycle event it follows. It does not guarantee delivery. +// Broadcast aborts in-flight attempts on ctx cancellation, so at shutdown the +// loop just exits; frames still queued are dropped with the process. func (m *Manager) broadcastLoop(ctx context.Context) { for { select { @@ -483,8 +482,9 @@ func (m *Manager) pushActiveSnapshot() { // resyncLoop is the anti-entropy heartbeat: every resyncInterval it re-asserts // this node's own active + recently-terminal workloads to all peers. Because the // origin is the single writer for its workloads, a peer that missed a delivery -// or made a wrong node-loss guess reconciles to the origin's authoritative state -// within a couple of intervals. +// or made a wrong node-loss guess may reconcile to the origin's authoritative +// state when later reassertions are queued and delivered. A full queue can drop +// those reassertions, and removals are not part of the re-sync set. func (m *Manager) resyncLoop(ctx context.Context) { ticker := time.NewTicker(resyncInterval) defer ticker.Stop() diff --git a/services/nvpair-workload-manager/manager_test.go b/services/nvpair-workload-manager/manager_test.go new file mode 100644 index 00000000..e8322e5e --- /dev/null +++ b/services/nvpair-workload-manager/manager_test.go @@ -0,0 +1,107 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "context" + "encoding/json" + "fmt" + "io" + "net" + "net/http" + "net/http/httptest" + "strconv" + "testing" + "time" + + "nvpair-shared/clustertrust" +) + +// codecNop is a throwaway io.ReadWriter for a Manager whose broker side is +// never exercised. +type codecNop struct{} + +func (codecNop) Read([]byte) (int, error) { return 0, io.EOF } +func (codecNop) Write(p []byte) (int, error) { return len(p), nil } + +type jSequence struct { + Seq int `json:"seq"` +} + +type jParams struct { + Params jSequence `json:"params"` +} + +func assertFrame(t *testing.T, body []byte, want int) { + t.Helper() + var frame jParams + if err := json.Unmarshal(body, &frame); err != nil { + t.Fatalf("decode frame %d: %v", want, err) + } + if frame.Params.Seq != want { + t.Fatalf("frame %d arrived with sequence %d", want, frame.Params.Seq) + } +} + +func newBroadcastManagerForPeer(t *testing.T, selfDir string, peer *httptest.Server) *Manager { + t.Helper() + host, portStr, err := net.SplitHostPort(peer.Listener.Addr().String()) + if err != nil { + t.Fatalf("split test peer address: %v", err) + } + port, err := strconv.Atoi(portStr) + if err != nil { + t.Fatalf("parse test peer port: %v", err) + } + + m := NewManager(NewCodec(codecNop{}), 0, "uuid-self", selfDir) + m.peers.Replace([]PeerNode{{ + ID: "peer-1", Addresses: []string{host}, Port: port, + TXT: []string{"cluster-uuid=uuid-peer"}, + }}) + return m +} + +// TestManager_BroadcastPreservesEnqueueOrder checks that frames queued in order +// reach a peer in the same order over the cluster-mTLS broadcast path. +func TestManager_BroadcastPreservesEnqueueOrder(t *testing.T) { + selfDir, peerDir := newPinnedPeerDirs(t) + peerMesh := clustertrust.Open(peerDir) + + const frameCount = 25 + received := make(chan []byte, frameCount) + mux := http.NewServeMux() + mux.HandleFunc(eventsPath, func(w http.ResponseWriter, r *http.Request) { + body, err := io.ReadAll(io.LimitReader(r.Body, 1<<20)) + if err != nil { + http.Error(w, "failed to read frame", http.StatusBadRequest) + return + } + received <- body + w.WriteHeader(http.StatusOK) + }) + ts := httptest.NewUnstartedServer(mux) + ts.TLS = peerMesh.ServerTLSConfig() + ts.StartTLS() + t.Cleanup(ts.Close) + + m := newBroadcastManagerForPeer(t, selfDir, ts) + + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + go m.broadcastLoop(ctx) + + for i := 0; i < frameCount; i++ { + m.broadcastFrame("workload:started", []byte(fmt.Sprintf(`{"seq":%d}`, i))) + } + + for want := 0; want < frameCount; want++ { + select { + case body := <-received: + assertFrame(t, body, want) + case <-ctx.Done(): + t.Fatalf("peer received %d of %d frames: %v", want, frameCount, ctx.Err()) + } + } +} diff --git a/services/nvpair-workload-manager/server.go b/services/nvpair-workload-manager/server.go index ef1c0b9e..d6019ebe 100644 --- a/services/nvpair-workload-manager/server.go +++ b/services/nvpair-workload-manager/server.go @@ -171,24 +171,22 @@ func (s *Server) handleLifecycle(w http.ResponseWriter, msg *Message) { // reach the broker so its store can reconcile (e.g. un-stick a wrongly // inferred failed). The store is idempotent, so bypassing dedup here is safe. // - // The key is recorded only after the broker emit succeeds: a failed emit - // answers 500 so the peer's retry budget kicks in, and that retry must not - // be mistaken for a duplicate (the old code recorded the key first, so a - // failed emit permanently dropped the event with no anti-entropy repair). - // A concurrent duplicate that passes seen() before either emit runs just - // emits twice, which the idempotent store absorbs. + // The key is recorded only after the broker emit succeeds. A failed emit + // leaves the key available for a retry; concurrent requests for the same key + // wait for that result rather than both emitting to the broker. resync := isResyncFrame(msg.Params) - lifecycleKey := "" + var duplicate bool if !resync { - lifecycleKey = keyLifecycle(wl) - if s.dedup.seen(lifecycleKey) { - slog.Debug("inter-node lifecycle deduplicated", "method", msg.Method, "id", wl.ID, "state", wl.State) - s.ok(w) - return - } + duplicate, err = s.dedup.emitOnce(keyLifecycle(wl), func() error { return s.emitUpsert(wl) }) + } else { + err = s.emitUpsert(wl) } - - if err := s.emitUpsert(wl); err != nil { + if duplicate { + slog.Debug("inter-node lifecycle deduplicated", "method", msg.Method, "id", wl.ID, "state", wl.State) + s.ok(w) + return + } + if err != nil { // A failed stdout write means the broker is gone; report a server // error so the peer's retry budget can kick in, but the local // interface severing is handled as a shutdown signal elsewhere. @@ -196,9 +194,6 @@ func (s *Server) handleLifecycle(w http.ResponseWriter, msg *Message) { http.Error(w, "broker unavailable", http.StatusInternalServerError) return } - if !resync { - s.dedup.add(lifecycleKey) - } slog.Info("relayed remote lifecycle as upsert", "method", msg.Method, "id", wl.ID, "state", wl.State, "node", wl.OriginatedFrom) s.ok(w) } @@ -219,22 +214,21 @@ func (s *Server) handleRemove(w http.ResponseWriter, msg *Message) { return } - // The removal key, like the lifecycle key, is recorded only after the - // broker emit succeeds, so a failed emit's retry isn't swallowed as a - // duplicate. - removeKey := keyRemove(nodeID, workloadID) - if s.dedup.seen(removeKey) { + // A failed emit leaves the key available for a retry. Concurrent requests + // for the same removal wait, so only one successful emit reaches the broker. + duplicate, err := s.dedup.emitOnce(keyRemove(nodeID, workloadID), func() error { + return s.emitRemove(workloadID, nodeID) + }) + if duplicate { slog.Debug("inter-node removal deduplicated", "workloadId", workloadID, "node", nodeID) s.ok(w) return } - - if err := s.emitRemove(workloadID, nodeID); err != nil { + if err != nil { slog.Error("failed to emit workloads:remove", "workloadId", workloadID, "node", nodeID, "err", err) http.Error(w, "broker unavailable", http.StatusInternalServerError) return } - s.dedup.add(removeKey) slog.Info("relayed remote removal", "workloadId", workloadID, "node", nodeID) s.ok(w) } diff --git a/services/nvpair-workload-manager/server_test.go b/services/nvpair-workload-manager/server_test.go new file mode 100644 index 00000000..9808f587 --- /dev/null +++ b/services/nvpair-workload-manager/server_test.go @@ -0,0 +1,152 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "errors" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" +) + +// TestServer_InterNodeDedupRecordedOnlyAfterSuccessfulEmit verifies that a +// failed broker emit does not record a dedup key, so the peer can retry, while +// a successful emit prevents later duplicates from reaching the broker. +func TestServer_InterNodeDedupRecordedOnlyAfterSuccessfulEmit(t *testing.T) { + test := func(name string, frame []byte, wantUpserts, wantRemoves int) { + t.Run(name, func(t *testing.T) { + self, peer := newPinnedPeerMeshes(t) + var upserts, removes int + failEmit := true + srv := NewServer(0, newDedupIndex(100), self, + func(*Workload) error { + if failEmit { + return errors.New("broker gone") + } + upserts++ + return nil + }, + func(string, string) error { + if failEmit { + return errors.New("broker gone") + } + removes++ + return nil + }, + ) + post := serveEventsOverMTLS(t, srv, self, peer) + + if code := post(frame); code != http.StatusInternalServerError { + t.Fatalf("failed emit status = %d, want 500", code) + } + if upserts != 0 || removes != 0 { + t.Fatalf("emits after failed attempt = (%d upserts, %d removes), want (0, 0)", upserts, removes) + } + + failEmit = false + if code := post(frame); code != http.StatusOK { + t.Fatalf("retry status = %d, want 200", code) + } + if upserts != wantUpserts || removes != wantRemoves { + t.Fatalf("emits after retry = (%d upserts, %d removes), want (%d, %d)", upserts, removes, wantUpserts, wantRemoves) + } + // they are duplicates at this point and shouldn't get sent to do work + if code := post(frame); code != http.StatusOK { + t.Fatalf("duplicate status = %d, want 200", code) + } + if upserts != wantUpserts || removes != wantRemoves { + t.Fatalf("emits after duplicate = (%d upserts, %d removes), want (%d, %d)", upserts, removes, wantUpserts, wantRemoves) + } + }) + } + + lifecycleFrame := []byte(`{"jsonrpc":"2.0","method":"workload:started","params":` + + `{"workloadInfo":{"id":"7","model":"llama3","engine":"ollama",` + + `"runId":"r1","state":"running","originatedFrom":"uuid-peer"}}}`) + removeFrame := []byte(`{"jsonrpc":"2.0","method":"workloads:remove",` + + `"params":{"workloadId":"7","originatedFrom":"uuid-peer"}}`) + + test("lifecycle upsert", lifecycleFrame, 1, 0) + test("workload removal", removeFrame, 0, 1) +} + +// TestServer_ConcurrentDuplicatesEmitOnce verifies that an in-flight event +// reserves its key until the broker emit finishes. Both paths share the same +// dedup index, but lifecycle and removal use distinct keys and emitters. +func TestServer_ConcurrentDuplicatesEmitOnce(t *testing.T) { + test := func(name string, params []byte, handle func(*Server, http.ResponseWriter, *Message)) { + t.Run(name, func(t *testing.T) { + // Pause the first broker emit after it reserves the key. This gives + // the identical second request a chance to reach the same handler. + firstEntered := make(chan struct{}) + releaseFirst := make(chan struct{}) + // A failure before the explicit release must not strand that goroutine. + defer func() { + select { + case <-releaseFirst: + default: + close(releaseFirst) + } + }() + var emits atomic.Int32 + emit := func() error { + if emits.Add(1) == 1 { + close(firstEntered) + <-releaseFirst + } + return nil + } + srv := NewServer(0, newDedupIndex(100), nil, + func(*Workload) error { return emit() }, + func(string, string) error { return emit() }, + ) + msg := &Message{Params: params} + post := func(result chan<- int) { + recorder := httptest.NewRecorder() + handle(srv, recorder, msg) + result <- recorder.Code + } + firstDone := make(chan int, 1) + secondDone := make(chan int, 1) + go post(firstDone) + // Wait until the first request is inside emit, then start its duplicate. + select { + case <-firstEntered: + case <-time.After(2 * time.Second): + t.Fatal("first request did not reach the broker emit") + } + go post(secondDone) + // The duplicate must wait for the first emit's result. The short + // timeout gives it an opportunity to expose a premature response. + select { + case code := <-secondDone: + t.Fatalf("duplicate completed before first emit: HTTP %d", code) + case <-time.After(50 * time.Millisecond): + } + // Once the first emit succeeds, both requests may return 200, but + // only that first request should have called the broker emitter. + close(releaseFirst) + for _, result := range []<-chan int{firstDone, secondDone} { + select { + case code := <-result: + if code != http.StatusOK { + t.Fatalf("request status = %d, want 200", code) + } + case <-time.After(2 * time.Second): + t.Fatal("request did not finish") + } + } + if got := emits.Load(); got != 1 { + t.Fatalf("broker emits = %d, want 1", got) + } + }) + } + + lifecycle := []byte(`{"workloadInfo":{"id":"7","model":"llama3","engine":"ollama","runId":"r1","state":"running","originatedFrom":"uuid-peer"}}`) + remove := []byte(`{"workloadId":"7","originatedFrom":"uuid-peer"}`) + test("lifecycle upsert", lifecycle, (*Server).handleLifecycle) + test("workload removal", remove, (*Server).handleRemove) +} diff --git a/services/nvpair-workload-manager/spec.md b/services/nvpair-workload-manager/spec.md index 194967a2..62d3b299 100644 --- a/services/nvpair-workload-manager/spec.md +++ b/services/nvpair-workload-manager/spec.md @@ -251,7 +251,7 @@ Response: `200 OK`. - Each node assigns `Workload.id` from a monotonic per-node counter (e.g. creation-time Unix ms); cross-node ordering comes from the Broker via the `(nodeId, workloadId)` tuple and timestamps. ## 12. Failure Modes and Mitigations -- **Peer unreachable/slow during broadcast** (partition, crash, suspend, GC pause): that peer's Broker misses the event, causing temporary inconsistency. → Concurrent broadcast with per-peer timeouts so one peer can't block others or the local Broker; bounded exponential-backoff retries, drop after max attempts with a metric/alert; peers reconcile via later events. +- **Peer unreachable/slow during broadcast** (partition, crash, suspend, GC pause): that peer's Broker misses the event, causing temporary inconsistency. → Concurrent fan-out lets other peers receive the current frame and keeps the local Broker unblocked; bounded exponential-backoff retries drop after max attempts with a warning. The ordered broadcast worker waits for all peers before taking the next frame, so a slow peer can delay later frames to healthy peers and build a queue. Revisit independent per-peer queues if this delay becomes a problem. - **Connection exhaustion under burst** (observed in the field): a per-event connection makes every lifecycle event pay a full mTLS handshake, and an unreaped idle connection pins a descriptor for the process lifetime. At inference-burst rates this starves the sender and the receiving listener until handshakes fail outright, and the events lost are disproportionately terminal ones — which the ~two-interval terminal retention cannot recover once the window passes. → Pool connections per peer with a bound on both total and idle connections per host, reap idle connections on both ends (§7.2), and drain response bodies so a connection is actually reusable. Note the residual: fan-out concurrency itself is still per-event, so a heartbeat round with many active workloads issues many simultaneous requests per peer. - **mTLS handshake failure / untrusted or expired cert**: that peer is isolated from broadcasts (`403`). → Reject and log with the presented identity; monitor and alert ahead of cert expiry; keep the trust store hot-reloadable so rotations apply without restart. - **Duplicate / re-broadcast events**: the Broker could see the same transition repeatedly, corrupting counts/state. → Deduplicate by `(nodeId, engine, runId, Workload.id, state, scheduledOn, seq)` (removals by `(nodeId, workloadId)`) before emitting, so the Broker updates at most once. From 43446529ba756aaded11cb7cd4a0d15fe559cd31 Mon Sep 17 00:00:00 2001 From: Kaylee Lubick Date: Mon, 28 Sep 2026 13:58:23 -0400 Subject: [PATCH 5/8] Address existing race with snapshots Signed-off-by: Kaylee Lubick --- services/nvpair-workload-manager/manager.go | 10 ++ .../nvpair-workload-manager/manager_test.go | 118 ++++++++++++++++++ services/nvpair-workload-manager/spec.md | 29 +++-- 3 files changed, 144 insertions(+), 13 deletions(-) diff --git a/services/nvpair-workload-manager/manager.go b/services/nvpair-workload-manager/manager.go index 1356c5bb..92c843fa 100644 --- a/services/nvpair-workload-manager/manager.go +++ b/services/nvpair-workload-manager/manager.go @@ -93,6 +93,10 @@ type Manager struct { activeMu sync.Mutex activeLocal map[workloadKey]workloadEvent + // broadcastMu keeps snapshot enqueueing ordered with local lifecycle and + // removal updates. Always acquire it before activeMu. + broadcastMu sync.Mutex + // broadcastCh serializes outbound inter-node frames in the order the // read loop produced them. broadcastFrame only enqueues (never blocks on // network I/O), and a single worker drains the queue in order — so a @@ -345,6 +349,8 @@ func (m *Manager) handleLocalLifecycle(msg *Message) { return } key := workloadKey{origin: wl.OriginatedFrom, engine: wl.Engine, runID: wl.RunID, id: wl.ID} + m.broadcastMu.Lock() + defer m.broadcastMu.Unlock() m.trackActive(key, msg.Method, msg.Params, wl.State) slog.Debug("broadcasting local lifecycle", "method", msg.Method, "id", wl.ID, "state", wl.State, "peers", m.peers.count()) m.broadcastFrame(msg.Method, msg.Params) @@ -360,6 +366,8 @@ func (m *Manager) handleLocalRemove(msg *Message) { // is preserved for peers' dedup. The removal wire carries only // (workloadId, originatedFrom) — no engine/runId — so drop every composite // key matching that pair. + m.broadcastMu.Lock() + defer m.broadcastMu.Unlock() m.untrackActive(nodeID, workloadID) slog.Debug("broadcasting local removal", "workloadId", workloadID, "node", nodeID, "peers", m.peers.count()) m.broadcastFrame(msg.Method, msg.Params) @@ -501,6 +509,8 @@ func (m *Manager) resyncLoop(ctx context.Context) { // broadcastSnapshot re-broadcasts the current re-sync set, one per-event frame // each. Shared by the discovery backfill and the heartbeat. func (m *Manager) broadcastSnapshot(reason string) { + m.broadcastMu.Lock() + defer m.broadcastMu.Unlock() snapshot := m.activeSnapshot() if len(snapshot) == 0 { return diff --git a/services/nvpair-workload-manager/manager_test.go b/services/nvpair-workload-manager/manager_test.go index e8322e5e..349e99b0 100644 --- a/services/nvpair-workload-manager/manager_test.go +++ b/services/nvpair-workload-manager/manager_test.go @@ -8,10 +8,12 @@ import ( "encoding/json" "fmt" "io" + "log/slog" "net" "net/http" "net/http/httptest" "strconv" + "sync" "testing" "time" @@ -105,3 +107,119 @@ func TestManager_BroadcastPreservesEnqueueOrder(t *testing.T) { } } } + +// snapshotPauseHandler stops broadcastSnapshot after it copies activeLocal but +// before it queues the copied frames. This makes the remove race reproducible. +type snapshotPauseHandler struct { + paused chan struct{} + release <-chan struct{} + once sync.Once +} + +func (*snapshotPauseHandler) Enabled(context.Context, slog.Level) bool { return true } + +func (h *snapshotPauseHandler) Handle(_ context.Context, record slog.Record) error { + if record.Message == "snapshot copied before removal" { + h.once.Do(func() { close(h.paused) }) + <-h.release + } + return nil +} + +func (h *snapshotPauseHandler) WithAttrs([]slog.Attr) slog.Handler { return h } +func (h *snapshotPauseHandler) WithGroup(string) slog.Handler { return h } + +func TestManager_SnapshotCannotResurrectRemovedWorkload(t *testing.T) { + m := &Manager{ + activeLocal: make(map[workloadKey]workloadEvent), + broadcastCh: make(chan []byte, 2), + peers: newPeerSet(0), + } + workload := &Workload{ + ID: "7", Model: "m", Engine: "ollama", RunID: "r1", + State: StateRunning, OriginatedFrom: "node-a", CreatedAt: 1, + } + lifecycle, err := json.Marshal(lifecycleParams{WorkloadInfo: workload}) + if err != nil { + t.Fatalf("marshal lifecycle: %v", err) + } + m.trackActive(workloadKey{origin: "node-a", engine: "ollama", runID: "r1", id: "7"}, MethodStarted, lifecycle, StateRunning) + + paused := make(chan struct{}) + release := make(chan struct{}) + var releaseOnce sync.Once + releaseSnapshot := func() { releaseOnce.Do(func() { close(release) }) } + defer releaseSnapshot() + previousLogger := slog.Default() + slog.SetDefault(slog.New(&snapshotPauseHandler{paused: paused, release: release})) + t.Cleanup(func() { slog.SetDefault(previousLogger) }) + + snapshotDone := make(chan struct{}) + go func() { + m.broadcastSnapshot("snapshot copied before removal") + close(snapshotDone) + }() + select { + case <-paused: + case <-time.After(2 * time.Second): + t.Fatal("snapshot did not pause after copying the workload") + } + + removal, err := json.Marshal(removeParams{WorkloadID: "7", OriginatedFrom: "node-a"}) + if err != nil { + t.Fatalf("marshal removal: %v", err) + } + removeDone := make(chan struct{}) + go func() { + m.handleLocalRemove(&Message{Method: MethodRemove, Params: removal}) + close(removeDone) + }() + // The current implementation queues the removal while the snapshot is + // paused. An implementation that serializes both operations may instead + // hold the removal until the snapshot resumes. + select { + case <-removeDone: + case <-time.After(100 * time.Millisecond): + } + releaseSnapshot() + for _, done := range []<-chan struct{}{snapshotDone, removeDone} { + select { + case <-done: + case <-time.After(2 * time.Second): + t.Fatal("snapshot or removal did not finish") + } + } + + // The peer already has the running workload. Apply the queued frames in + // their actual delivery order through the receiver's handlers. + present := true + removes := 0 + receiver := NewServer(0, newDedupIndex(8), nil, + func(*Workload) error { present = true; return nil }, + func(string, string) error { present = false; removes++; return nil }, + ) + for len(m.broadcastCh) > 0 { + var frame Message + if err := json.Unmarshal(<-m.broadcastCh, &frame); err != nil { + t.Fatalf("decode queued frame: %v", err) + } + response := httptest.NewRecorder() + switch frame.Method { + case MethodStarted: + receiver.handleLifecycle(response, &frame) + case MethodRemove: + receiver.handleRemove(response, &frame) + default: + t.Fatalf("unexpected queued method %q", frame.Method) + } + if response.Code != http.StatusOK { + t.Fatalf("peer handled %s with HTTP %d", frame.Method, response.Code) + } + } + if removes != 1 { + t.Fatalf("peer received %d removals, want 1", removes) + } + if present { + t.Fatal("stale snapshot upsert resurrected the workload after its removal") + } +} diff --git a/services/nvpair-workload-manager/spec.md b/services/nvpair-workload-manager/spec.md index 62d3b299..0c0b4d57 100644 --- a/services/nvpair-workload-manager/spec.md +++ b/services/nvpair-workload-manager/spec.md @@ -30,13 +30,13 @@ Tracks inference workloads cluster-wide as they are queued, executed, and retire - **Relay a remote event**: on receiving a peer's `workload:*`, validate and deduplicate (by `(nodeId, engine, runId, Workload.id, state, scheduledOn, seq)`), then emit `workloads:upsert` on `stdout` so the Broker updates its catalog. A peer `workloads:remove` (deduplicated by `(nodeId, workloadId)`) is relayed as `workloads:remove` on `stdout`. - **Discover peers**: the Broker registers this node's `wl` port with the `nvpair-node-scanner` discovery daemon, which carries it on this node's single consolidated record; the Workload Manager subscribes for `wl` nodes and rebuilds the target set from each `discovery:nodes` snapshot the Broker relays, so a node joins the set when it appears in a snapshot and leaves when it is absent from the next one. - **Edge case — duplicate / re-broadcast events**: retries or re-broadcasts arriving more than once are deduplicated (`(nodeId, engine, runId, Workload.id, state, scheduledOn, seq)` for lifecycle, `(nodeId, workloadId)` for removals) so the Broker is updated at most once. -- **Edge case — peer unreachable**: a target that is down, slow, or partitioned (crash, network split, laptop suspend) does not block others — broadcast runs concurrently with per-peer timeouts and bounded retries. +- **Edge case — peer unreachable**: a target that is down, slow, or partitioned (crash, network split, laptop suspend) does not block the local Broker or delivery of the current frame to other peers. Fan-out runs concurrently with per-peer timeouts and bounded retries, but the ordered worker waits for the round to finish before sending the next frame. ## 4. Open Questions / Risks - **`initializing` state (closed)**: removed. It had no `workload:*` method and nothing ever produced it, so it was an unreachable member of a closed enum. The proxy has no distinct pre-dispatch moment to represent — admission and the first dispatch are effectively simultaneous — so the value was deleted rather than given a method. - **`EngineType` values (open — needs third-party feedback)**: the valid `engine` set is undefined. Pending the inference-engine team's list, `engine` is treated as an opaque pass-through string (must be present and non-empty, value not validated). -- **Risk — late joiners / no replay**: nodes discovered mid-session only get events broadcast after they appear in the target set; their view is incomplete until later events flow (by design). -- **Risk — best-effort fan-out**: an unreachable peer misses events; temporary inconsistency until later events arrive or the Broker reconciles by timestamp. +- **Risk — late joiners / partial backfill**: newly discovered peers receive a re-assertion of this node's active and recently terminal workloads. Retired workloads and removals are not replayed, so backfill is not a complete event history. +- **Risk — best-effort fan-out**: an unreachable peer misses events; later lifecycle re-assertions can repair tracked state when delivered, but a missed removal has no heartbeat or backfill repair path. - **Risk — snapshot staleness**: the target set is only as current as the last `discovery:nodes` snapshot, so a departed node lingers as a target and a new node appears slowly, bounded by the discovery daemon's own liveness handling rather than by anything this service controls. - **Risk — dedup granularity**: the key is `(nodeId, engine, runId, Workload.id, state, scheduledOn, seq)`, so a re-broadcast carrying the same key with updated metadata (e.g. a corrected `error`) is dropped, not merged. `seq` is the producer's event counter and is what makes the key exact, because this index is a permanent set: any key derived only from a workload's current shape collides as soon as the workload revisits a shape it already had, which a retry does routinely (queued on A, placement cleared between attempts, queued on A again). Without it peers dropped that third event, kept the interim unplaced record, and stopped counting an active job against the node running it. @@ -44,7 +44,9 @@ Tracks inference workloads cluster-wide as they are queued, executed, and retire **Functional** - Accept local `workload:*` and `workloads:remove` notifications from the Broker over `stdin` (or a named pipe) and broadcast each to all discovered peers via `POST /v1/workloads/events`. +- Preserve the enqueue order of local Broker notifications when sending frames to each peer. If both are delivered, a removal must not overtake a lifecycle event queued before it; this is an ordering guarantee, not a delivery guarantee. - For each validated, deduplicated inter-node `workload:*`, emit `workloads:upsert` to the Broker (translated, not forwarded unchanged); for each inter-node `workloads:remove`, emit `workloads:remove` to the Broker (not re-broadcast). +- Record an inter-node dedup key only after the corresponding notification is successfully written to the Broker. If the write fails, return `500` and leave the key available for a retry; concurrent requests for the same key must not both emit successfully. - Subscribe to the Broker's discovery relay with `discovery:subscribe` filtered to the `wl` service, and rebuild the broadcast target set from every `discovery:nodes` snapshot: take each peer's dialable address and `wl` port from its directory entry, skip entries advertising no `wl` port or no address, and exclude this node's own entry by `hostUuid` rather than by hostname. A snapshot carries the full filtered set and replaces the target set wholesale, so there are no per-node deltas to apply and a peer absent from a snapshot simply stops being a target. - Deduplicate inbound lifecycle events by `(nodeId, engine, runId, Workload.id, state, scheduledOn, seq)` and removals by `(nodeId, workloadId)`, using a configurable bounded LRU index (default ~10,000 entries, sized for session-scoped volume at ~dozen-node scale). `nodeId` is part of the key because `Workload.id` is only unique per node (§11) — keying on `id` alone would collide across nodes and silently drop a legitimate peer's event. `engine` and `runId` are there for the same reason one level down: `Workload.id` is a per-process counter, both engine proxies count from 1, and the counter resets on restart, so without them a concurrent Ollama and LM Studio job both holding id `"1"` would collapse into one. The client-visible identity the Broker uses as the global key remains the coarser `(nodeId, workloadId)` pair (§10); this key is finer on purpose. - Validate inbound payloads; reject malformed envelopes or unknown `method` values (`400 Bad Request`). @@ -52,7 +54,7 @@ Tracks inference workloads cluster-wide as they are queued, executed, and retire **Non-functional** - Authenticate all inter-node traffic with mTLS, validating client and server certificates against the trusted node store; reject untrusted clients (`403 Forbidden`). - Stay stateless re: workload history — the Broker is the source of truth. -- A failed broadcast to one peer must not block delivery to other peers or the local Broker. +- A failed broadcast to one peer must not block delivery of the current frame to other peers or the local Broker. A slow peer may delay later frames while the ordered worker finishes the current round. - Serialize `stdout` writes so frames never interleave; the Broker orders events per workload by `(nodeId, workloadId)` and timestamps. ## 6. Inputs and Outputs @@ -191,8 +193,8 @@ Example notification: - Connections are **keep-alive and pooled per peer**, not per event. A sender holds a bounded set of long-lived connections per peer and reaps one after `clustertrust.PeerIdleTimeout`; the listener sets a strictly longer `IdleTimeout` (`clustertrust.PeerListenerIdleTimeout`) so the sender is always the side that discards a doubtful connection. This is normative, not an optimization: a handshake per event both starves the receiver and, without an idle bound on either side, pins a descriptor per connection for the life of the process. - A node broadcasts each local `workload:*` and `workloads:remove` to all discovered peers. - Body: JSON-RPC 2.0 notification — `workload:*` with `params.workloadInfo`, or `workloads:remove` with `params.workloadId` (and optional `params.originatedFrom`). -- Responses: `200 OK` (accepted; lifecycle → `workloads:upsert`, removal → `workloads:remove` to the local Broker); `400 Bad Request` (malformed / unknown `method`); `403 Forbidden` (client cert absent or untrusted). -- Idempotency: lifecycle by `(nodeId, engine, runId, Workload.id, state, scheduledOn, seq)`; removal by `(nodeId, workloadId)`. +- Responses: `200 OK` (written to the local Broker or already deduplicated); `400 Bad Request` (malformed / unknown `method`); `403 Forbidden` (client cert absent or untrusted); `500 Internal Server Error` (Broker write failed, so the sender may retry). A successful write does not acknowledge that the Broker applied the event. +- Idempotency: lifecycle by `(nodeId, engine, runId, Workload.id, state, scheduledOn, seq)`; removal by `(nodeId, workloadId)`. Record a key after a successful Broker write, and serialize concurrent requests for the same key through that write. Tagged lifecycle re-sync frames bypass dedup so they can re-assert state. Example request (`POST /v1/workloads/events`): ```json @@ -238,27 +240,28 @@ Response: `200 OK`. ## 10. Design Constraints - **Performance**: low event volume (lifecycle transitions, not telemetry); no end-to-end latency SLA, delivery is best-effort. Local forwarding writes serialized frames to `stdout`; inter-node broadcast is concurrent with per-peer timeouts. -- **Scalability**: ~dozen-node cluster; each local event fans out to all peers. Only in-memory state is the bounded dedup index and target set. -- **Reliability**: best-effort fan-out, no SLA. Concurrent per-peer broadcast with timeouts; retry transient failures with bounded exponential backoff, drop after max attempts (metric/alert; parameters implementation-defined). Serialized `stdout`; the Broker orders per workload by `(nodeId, workloadId)` and timestamps. No replay — peers reconcile via later events. -- **Re-sync cadence is a cross-process contract.** The anti-entropy heartbeat re-asserts every *active* local-origin workload on every interval, indefinitely, and every re-assertion is tagged so the receiver's state dedup does not swallow it. A consumer is entitled to treat prolonged silence about a workload it believes active as evidence that the workload is finished or the origin is gone — `nvpair-ui-broker` does exactly that in its staleness sweep, keyed to a multiple of this interval. Lengthening the interval beyond that consumer's budget, or dropping the re-assertion tag, therefore causes peers to retire live work; change both sides together. +- **Scalability**: ~dozen-node cluster; each local event fans out to all peers. Transient in-memory state includes the bounded dedup index, target set, active/recently-terminal re-sync set, and bounded outbound queue. +- **Reliability**: best-effort fan-out, no SLA. Peers receive each frame concurrently, with timeouts and bounded exponential-backoff retries; the ordered worker finishes that round before sending the next frame. Drop after max attempts (metric/alert; parameters implementation-defined), and drop newly enqueued frames with a warning if the queue is full so the local read loop remains unblocked. Serialized `stdout`; the Broker orders per workload by `(nodeId, workloadId)` and timestamps. Heartbeat and discovery backfill re-assert tracked lifecycle state, but neither replays removals or full event history. +- **Re-sync cadence is a cross-process contract.** The anti-entropy heartbeat enqueues a re-assertion of every *active* local-origin workload on every interval, indefinitely, and tags each re-assertion so the receiver's state dedup does not swallow it. Delivery can still be delayed or dropped by the outbound queue. A consumer is entitled to treat prolonged silence about a workload it believes active as evidence that the workload is finished or the origin is gone — `nvpair-ui-broker` does exactly that in its staleness sweep, keyed to a multiple of this interval. Lengthening the interval beyond that consumer's budget, or dropping the re-assertion tag, therefore causes peers to retire live work; change both sides together. - **Security**: all inter-node traffic mTLS-authenticated against the trusted node store; reject and log untrusted callers (`403`) with the presented identity; monitor cert expiry; keep the trust store hot-reloadable. Payloads carry workload metadata only (model, state, node IDs, errors) — no request/response data or PII. `requesterId`, `nodeId`, and `Workload.id` are opaque system-generated identifiers (not PII). - **Compliance**: workload information must contain no PII. ## 11. Assumptions - The Broker is always the parent process: it spawns the Workload Manager, writes `workload:*` / `workloads:remove` to `stdin`, and consumes `workloads:upsert` / `workloads:remove` from `stdout`. If the Broker exits, the manager shuts down cleanly — no reconnect or buffering. - The Broker is the source of truth; local notifications are well-formed JSON-RPC with epoch-ms timestamps. -- Small cluster (~dozen nodes), low event volume; peers mutually authenticated via mTLS. New nodes get only events broadcast after discovery — no historical replay. +- Small cluster (~dozen nodes), low event volume; peers mutually authenticated via mTLS. New nodes receive current active and recently terminal workload re-assertions after discovery, but no historical event or removal replay. - Each node assigns `Workload.id` from a monotonic per-node counter (e.g. creation-time Unix ms); cross-node ordering comes from the Broker via the `(nodeId, workloadId)` tuple and timestamps. ## 12. Failure Modes and Mitigations - **Peer unreachable/slow during broadcast** (partition, crash, suspend, GC pause): that peer's Broker misses the event, causing temporary inconsistency. → Concurrent fan-out lets other peers receive the current frame and keeps the local Broker unblocked; bounded exponential-backoff retries drop after max attempts with a warning. The ordered broadcast worker waits for all peers before taking the next frame, so a slow peer can delay later frames to healthy peers and build a queue. Revisit independent per-peer queues if this delay becomes a problem. -- **Connection exhaustion under burst** (observed in the field): a per-event connection makes every lifecycle event pay a full mTLS handshake, and an unreaped idle connection pins a descriptor for the process lifetime. At inference-burst rates this starves the sender and the receiving listener until handshakes fail outright, and the events lost are disproportionately terminal ones — which the ~two-interval terminal retention cannot recover once the window passes. → Pool connections per peer with a bound on both total and idle connections per host, reap idle connections on both ends (§7.2), and drain response bodies so a connection is actually reusable. Note the residual: fan-out concurrency itself is still per-event, so a heartbeat round with many active workloads issues many simultaneous requests per peer. +- **Outbound queue full**: a slow broadcast round can fill the bounded queue; the read loop drops each new frame with a warning rather than blocking. A dropped active or recently terminal lifecycle event may be repaired by a later heartbeat or discovery backfill if that re-assertion is delivered. A removal is deleted from the re-sync set before broadcast and is not replayed, so dropping its frame can leave a stale workload on a peer with no removal repair path. +- **Connection exhaustion under burst** (observed in the field): a per-event connection makes every lifecycle event pay a full mTLS handshake, and an unreaped idle connection pins a descriptor for the process lifetime. At inference-burst rates this starves the sender and the receiving listener until handshakes fail outright, and the events lost are disproportionately terminal ones — which the ~two-interval terminal retention cannot recover once the window passes. → Pool connections per peer with a bound on both total and idle connections per host, reap idle connections on both ends (§7.2), and drain response bodies so a connection is actually reusable. Fan-out remains concurrent across peers for one frame, while the ordered worker sends frames sequentially. - **mTLS handshake failure / untrusted or expired cert**: that peer is isolated from broadcasts (`403`). → Reject and log with the presented identity; monitor and alert ahead of cert expiry; keep the trust store hot-reloadable so rotations apply without restart. -- **Duplicate / re-broadcast events**: the Broker could see the same transition repeatedly, corrupting counts/state. → Deduplicate by `(nodeId, engine, runId, Workload.id, state, scheduledOn, seq)` (removals by `(nodeId, workloadId)`) before emitting, so the Broker updates at most once. +- **Duplicate / re-broadcast events**: the Broker could see the same transition repeatedly, corrupting counts/state. → Deduplicate by `(nodeId, engine, runId, Workload.id, state, scheduledOn, seq)` (removals by `(nodeId, workloadId)`) and serialize concurrent requests for one key. Record the key only after a successful Broker write; on failure, return `500` and let a retry attempt the write. Re-sync lifecycle frames bypass this index so they can correct the Broker's view. - **Malformed payload / unknown `method`**: could crash the parser or propagate garbage. → Validate every envelope; reject (`400`) inter-node or drop-and-log locally; never forward unvalidated payloads. - **Local interface severed** (Broker exited → `stdin` EOF / `stdout` `EPIPE` / `ERROR_BROKEN_PIPE`): an orphaned manager would accept peer events with nowhere to forward them. → Treat EOF/`EPIPE` as the shutdown signal: stop the listener and exit cleanly. No reconnect/buffering — a new Broker spawns a new (stateless) manager. The small load lets the OS pipe buffer absorb serialized `stdout` writes. - **Out-of-order delivery for the same workload (local)**: the Broker could see an inconsistent progression (e.g. `completed` before `started`). → Serialize `stdout` writes; the Broker orders per workload via `(nodeId, workloadId)` (monotonic IDs) and timestamps. -- **Stale discovery snapshot** (departed node still targeted, or new node not yet in a snapshot): wasted attempts to dead nodes, or new nodes missing events. → Rebuild the target set from every `discovery:nodes` snapshot, so a peer absent from the next snapshot stops being a target without any expiry logic here; liveness is the discovery daemon's, which probes a node before evicting it so a transient miss does not flap the set. New joiners catch up via later events (no replay). +- **Stale discovery snapshot** (departed node still targeted, or new node not yet in a snapshot): wasted attempts to dead nodes, or new nodes missing events. → Rebuild the target set from every `discovery:nodes` snapshot, so a peer absent from the next snapshot stops being a target without any expiry logic here; liveness is the discovery daemon's, which probes a node before evicting it so a transient miss does not flap the set. Newly discovered peers receive active and recently terminal state backfill, but not removals or event history. ## 13. Observability - **Logging**: mTLS/cert rejections with presented identity (`403`) and malformed/unknown-`method` rejections (`400`); per-peer broadcast failures including final drop; deduplicated/dropped inbound events; target-set changes applied from a `discovery:nodes` snapshot (peers added, peers removed, resulting total); startup and clean shutdown on `stdin` EOF / `stdout` `EPIPE`. From 8348d1b95ad12592bc92c54ca3c05951a4622ee6 Mon Sep 17 00:00:00 2001 From: Kaylee Lubick Date: Mon, 28 Sep 2026 14:05:42 -0400 Subject: [PATCH 6/8] cleanup: remove seenOrAdd in favor of emitOnce Signed-off-by: Kaylee Lubick --- services/nvpair-workload-manager/dedup.go | 22 ++---- .../nvpair-workload-manager/dedup_test.go | 71 +++++++++++-------- 2 files changed, 45 insertions(+), 48 deletions(-) diff --git a/services/nvpair-workload-manager/dedup.go b/services/nvpair-workload-manager/dedup.go index b29d7262..48d5cc81 100644 --- a/services/nvpair-workload-manager/dedup.go +++ b/services/nvpair-workload-manager/dedup.go @@ -13,10 +13,10 @@ import ( // session-scoped volume at ~dozen-node scale (spec §5). const defaultDedupCapacity = 10000 -// dedupIndex is a bounded LRU set of keys. It answers a single question: -// "have I seen this key before?" and records it if not, evicting the -// least-recently-seen key once capacity is exceeded. It is safe for -// concurrent use — the inter-node HTTP handler runs one goroutine per +// dedupIndex is a bounded LRU set of successfully emitted keys. It records a +// key only after the broker emit succeeds, and serializes concurrent emits for +// the same key. Completed keys are evicted least-recently-used first. It is +// safe for concurrent use — the inter-node HTTP handler runs one goroutine per // request. // // Keys are opaque strings built by the caller: lifecycle events key on @@ -48,20 +48,6 @@ func newDedupIndex(capacity int) *dedupIndex { } } -// seenOrAdd returns true if the key was already present (a duplicate). On a -// first sighting it records the key and returns false. Either way the key is -// promoted to most-recently-seen. -func (d *dedupIndex) seenOrAdd(key string) bool { - d.mu.Lock() - defer d.mu.Unlock() - if el, ok := d.items[key]; ok { - d.ll.MoveToFront(el) - return true - } - d.addLocked(key) - return false -} - // emitOnce serializes the check, broker emit, and record for one key. Other // keys can emit concurrently. If emit fails, the key stays absent and a waiting // request can retry it. A completed key is reported as a duplicate. diff --git a/services/nvpair-workload-manager/dedup_test.go b/services/nvpair-workload-manager/dedup_test.go index 64639b1c..40cee379 100644 --- a/services/nvpair-workload-manager/dedup_test.go +++ b/services/nvpair-workload-manager/dedup_test.go @@ -10,6 +10,17 @@ import ( "time" ) +// emitDedupKey exercises the production dedup path with a successful broker +// emit, then reports whether that key was already recorded. +func emitDedupKey(t *testing.T, d *dedupIndex, key string) bool { + t.Helper() + duplicate, err := d.emitOnce(key, func() error { return nil }) + if err != nil { + t.Fatalf("emit dedup key: %v", err) + } + return duplicate +} + // TestDedupIndex_EmitOnceSerializesOnlyMatchingKeys proves that an in-flight // emit blocks only requests with the same key. The test controls this order: // @@ -157,16 +168,16 @@ func TestDedupIndex_WaiterRetriesAfterFailedEmit(t *testing.T) { } } -func TestDedupSeenOrAdd(t *testing.T) { +func TestDedupEmitOnce(t *testing.T) { d := newDedupIndex(8) - if d.seenOrAdd("a") { + if emitDedupKey(t, d, "a") { t.Fatal("first sighting of a should be new") } - if !d.seenOrAdd("a") { + if !emitDedupKey(t, d, "a") { t.Fatal("second sighting of a should be a duplicate") } - if d.seenOrAdd("b") { + if emitDedupKey(t, d, "b") { t.Fatal("first sighting of b should be new") } } @@ -174,11 +185,11 @@ func TestDedupSeenOrAdd(t *testing.T) { func TestDedupEviction(t *testing.T) { d := newDedupIndex(2) - d.seenOrAdd("a") // {a} - d.seenOrAdd("b") // {a,b} - d.seenOrAdd("c") // evicts a -> {b,c} + emitDedupKey(t, d, "a") // {a} + emitDedupKey(t, d, "b") // {a,b} + emitDedupKey(t, d, "c") // evicts a -> {b,c} - if d.seenOrAdd("a") { + if emitDedupKey(t, d, "a") { t.Fatal("a should have been evicted and read as new") } } @@ -189,17 +200,17 @@ func TestDedupKeysDistinguishStateAndKind(t *testing.T) { w := &Workload{ID: "wl-1", OriginatedFrom: "node-A", State: StateQueued} wRunning := &Workload{ID: "wl-1", OriginatedFrom: "node-A", State: StateRunning} - if d.seenOrAdd(keyLifecycle(w)) { + if emitDedupKey(t, d, keyLifecycle(w)) { t.Fatal("queued should be new") } - if d.seenOrAdd(keyLifecycle(wRunning)) { + if emitDedupKey(t, d, keyLifecycle(wRunning)) { t.Fatal("running for same id is a different key, should be new") } - if !d.seenOrAdd(keyLifecycle(w)) { + if !emitDedupKey(t, d, keyLifecycle(w)) { t.Fatal("repeat queued should dedup") } // A removal keyed on the same id must not collide with a lifecycle key. - if d.seenOrAdd(keyRemove("node-A", "wl-1")) { + if emitDedupKey(t, d, keyRemove("node-A", "wl-1")) { t.Fatal("removal of wl-1 must not collide with lifecycle keys") } } @@ -214,13 +225,13 @@ func TestDedupDistinguishesNodes(t *testing.T) { nodeA := &Workload{ID: "wl-1", OriginatedFrom: "node-A", State: StateQueued} nodeB := &Workload{ID: "wl-1", OriginatedFrom: "node-B", State: StateQueued} - if d.seenOrAdd(keyLifecycle(nodeA)) { + if emitDedupKey(t, d, keyLifecycle(nodeA)) { t.Fatal("node-A wl-1 should be new") } - if d.seenOrAdd(keyLifecycle(nodeB)) { + if emitDedupKey(t, d, keyLifecycle(nodeB)) { t.Fatal("node-B wl-1 has the same id but a different node, must not dedup against node-A") } - if !d.seenOrAdd(keyLifecycle(nodeA)) { + if !emitDedupKey(t, d, keyLifecycle(nodeA)) { t.Fatal("repeat of node-A wl-1 should dedup") } } @@ -234,16 +245,16 @@ func TestDedupDistinguishesEngineAndRun(t *testing.T) { lmstudio := &Workload{ID: "1", OriginatedFrom: "host", Engine: "lmstudio", RunID: "r2", State: StateRunning} restarted := &Workload{ID: "1", OriginatedFrom: "host", Engine: "ollama", RunID: "r3", State: StateRunning} - if d.seenOrAdd(keyLifecycle(ollama)) { + if emitDedupKey(t, d, keyLifecycle(ollama)) { t.Fatal("ollama host/1 should be new") } - if d.seenOrAdd(keyLifecycle(lmstudio)) { + if emitDedupKey(t, d, keyLifecycle(lmstudio)) { t.Fatal("lmstudio host/1 shares the id but a different engine; must not dedup") } - if d.seenOrAdd(keyLifecycle(restarted)) { + if emitDedupKey(t, d, keyLifecycle(restarted)) { t.Fatal("a reused id from a new run must not dedup against the old run") } - if !d.seenOrAdd(keyLifecycle(ollama)) { + if !emitDedupKey(t, d, keyLifecycle(ollama)) { t.Fatal("repeat of ollama host/1 should dedup") } } @@ -261,13 +272,13 @@ func TestDedupDistinguishesPlacement(t *testing.T) { first := &Workload{ID: "1", OriginatedFrom: "host", Engine: "ollama", RunID: "r1", State: StateRunning, ScheduledOn: "node-A"} repointed := &Workload{ID: "1", OriginatedFrom: "host", Engine: "ollama", RunID: "r1", State: StateRunning, ScheduledOn: "node-B"} - if d.seenOrAdd(keyLifecycle(first)) { + if emitDedupKey(t, d, keyLifecycle(first)) { t.Fatal("first placement on node-A should be new") } - if d.seenOrAdd(keyLifecycle(repointed)) { + if emitDedupKey(t, d, keyLifecycle(repointed)) { t.Fatal("a re-point to node-B differs only in scheduledOn and must not dedup against node-A") } - if !d.seenOrAdd(keyLifecycle(first)) { + if !emitDedupKey(t, d, keyLifecycle(first)) { t.Fatal("a resent frame for the node-A placement should still dedup") } } @@ -298,21 +309,21 @@ func TestDedupDistinguishesRepeatedPlacements(t *testing.T) { cleared := base(2, "") backOnA := base(3, "node-A") - if d.seenOrAdd(keyLifecycle(onA)) { + if emitDedupKey(t, d, keyLifecycle(onA)) { t.Fatal("first placement on node-A should be new") } - if d.seenOrAdd(keyLifecycle(cleared)) { + if emitDedupKey(t, d, keyLifecycle(cleared)) { t.Fatal("clearing the placement between attempts should be new") } - if d.seenOrAdd(keyLifecycle(backOnA)) { + if emitDedupKey(t, d, keyLifecycle(backOnA)) { t.Fatal("re-dispatching to node-A repeats an earlier shape and must NOT dedup against it") } // A redelivery of any of those frames carries its original sequence, so it // is still recognised as one. - if !d.seenOrAdd(keyLifecycle(base(1, "node-A"))) { + if !emitDedupKey(t, d, keyLifecycle(base(1, "node-A"))) { t.Fatal("a resent frame for the first placement should still dedup") } - if !d.seenOrAdd(keyLifecycle(base(3, "node-A"))) { + if !emitDedupKey(t, d, keyLifecycle(base(3, "node-A"))) { t.Fatal("a resent frame for the third event should still dedup") } } @@ -323,13 +334,13 @@ func TestDedupDistinguishesRepeatedPlacements(t *testing.T) { func TestDedupRemovalDistinguishesNodes(t *testing.T) { d := newDedupIndex(8) - if d.seenOrAdd(keyRemove("node-A", "wl-1")) { + if emitDedupKey(t, d, keyRemove("node-A", "wl-1")) { t.Fatal("removal of node-A wl-1 should be new") } - if d.seenOrAdd(keyRemove("node-B", "wl-1")) { + if emitDedupKey(t, d, keyRemove("node-B", "wl-1")) { t.Fatal("removal of node-B wl-1 must not dedup against node-A") } - if !d.seenOrAdd(keyRemove("node-A", "wl-1")) { + if !emitDedupKey(t, d, keyRemove("node-A", "wl-1")) { t.Fatal("repeat removal of node-A wl-1 should dedup") } } From 572620310f75753d22847fc4e072ed1c802a335e Mon Sep 17 00:00:00 2001 From: Kaylee Lubick Date: Mon, 28 Sep 2026 14:18:10 -0400 Subject: [PATCH 7/8] docs: refresh generated service contract report Signed-off-by: Kaylee Lubick --- desktop/docs/services-api.md | 3 +++ 1 file changed, 3 insertions(+) diff --git a/desktop/docs/services-api.md b/desktop/docs/services-api.md index 7aea40f9..3631e095 100644 --- a/desktop/docs/services-api.md +++ b/desktop/docs/services-api.md @@ -308,3 +308,6 @@ | `workloads:remove` | notification (we consume) | ✅ yes | | `workloads:upsert` | notification (we consume) | ✅ yes | +**Dynamic / unresolved notify sites (verify by hand — `npm run service-contracts` prints the line numbers):** +- `ctx (var) (manager.go)` + From 1d89c42c71f0fc0993f3c90268b96b5b70718db2 Mon Sep 17 00:00:00 2001 From: Kaylee Lubick Date: Mon, 28 Sep 2026 14:22:55 -0400 Subject: [PATCH 8/8] fix: ignore transport broadcasts in service contract scan Signed-off-by: Kaylee Lubick --- desktop/docs/services-api.md | 3 --- desktop/scripts/verify-service-contracts.ts | 5 ++++- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/desktop/docs/services-api.md b/desktop/docs/services-api.md index 3631e095..7aea40f9 100644 --- a/desktop/docs/services-api.md +++ b/desktop/docs/services-api.md @@ -308,6 +308,3 @@ | `workloads:remove` | notification (we consume) | ✅ yes | | `workloads:upsert` | notification (we consume) | ✅ yes | -**Dynamic / unresolved notify sites (verify by hand — `npm run service-contracts` prints the line numbers):** -- `ctx (var) (manager.go)` - diff --git a/desktop/scripts/verify-service-contracts.ts b/desktop/scripts/verify-service-contracts.ts index e09dd5ee..94d79b52 100644 --- a/desktop/scripts/verify-service-contracts.ts +++ b/desktop/scripts/verify-service-contracts.ts @@ -136,7 +136,10 @@ const QUOTED_RE = /"([^"]+)"/g const IDENT_RE = /\b([A-Za-z_]\w*)\b/g const NOTIFY_LITERAL_RE = /(?:notify|emit|publish|broadcast)\(\s*"([^"]+)"/gi const NOTIFY_PREFIX_CONCAT_RE = /(?:notify|emit|publish|broadcast)\(\s*"([^"]+)"\s*\+/i -const NOTIFY_IDENT_RE = /(?:notify|emit|publish|broadcast)\(\s*([A-Za-z_]\w*)\s*,/gi +// Broadcast(ctx, frame) is HTTP fan-out, not a JSON-RPC notification. Keep +// literal Broadcast methods above, but do not treat an identifier argument to +// Broadcast as an unresolved JSON-RPC method. +const NOTIFY_IDENT_RE = /(?:notify|emit|publish)\(\s*([A-Za-z_]\w*)\s*,/gi /** A `| `nvpair-ui-broker` | 0.37.0 | … |` row — forbidden in hand-maintained docs. */ const DOC_COMPONENT_ROW_RE = /^\|\s*`([a-z0-9-]+)`\s*\|\s*(\d+\.\d+\.\d+)\s*\|/gm