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 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 4d44dd5c..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 @@ -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,21 +44,54 @@ func newDedupIndex(capacity int) *dedupIndex { capacity: capacity, ll: list.New(), items: make(map[string]*list.Element, capacity), + inFlight: make(map[string]chan struct{}), } } -// 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() +// 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() - if el, ok := d.items[key]; ok { - d.ll.MoveToFront(el) - return true + // 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 { @@ -67,7 +101,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_test.go b/services/nvpair-workload-manager/dedup_test.go index 6c4e8d68..40cee379 100644 --- a/services/nvpair-workload-manager/dedup_test.go +++ b/services/nvpair-workload-manager/dedup_test.go @@ -3,18 +3,181 @@ package main -import "testing" +import ( + "errors" + "sync/atomic" + "testing" + "time" +) -func TestDedupSeenOrAdd(t *testing.T) { +// 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: +// +// 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 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") } } @@ -22,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") } } @@ -37,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") } } @@ -62,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") } } @@ -82,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") } } @@ -109,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") } } @@ -146,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") } } @@ -171,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") } } diff --git a/services/nvpair-workload-manager/manager.go b/services/nvpair-workload-manager/manager.go index 84ea634d..92c843fa 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,11 +88,23 @@ 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 + // 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 + // 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 + ctx context.Context cancel context.CancelFunc } @@ -124,6 +137,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 +171,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 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 // and to re-assert our workloads immediately: peers that could not receive @@ -331,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) @@ -346,24 +366,59 @@ 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) } -// 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. 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. 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 { 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) 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 { + case <-ctx.Done(): + return + case frame := <-m.broadcastCh: + m.broadcaster.Broadcast(ctx, frame) + } + } } // trackActive records the latest event for a local-origin workload. A @@ -435,8 +490,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() @@ -453,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 new file mode 100644 index 00000000..349e99b0 --- /dev/null +++ b/services/nvpair-workload-manager/manager_test.go @@ -0,0 +1,225 @@ +// 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" + "log/slog" + "net" + "net/http" + "net/http/httptest" + "strconv" + "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 } + +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()) + } + } +} + +// 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/server.go b/services/nvpair-workload-manager/server.go index 742b5607..d6019ebe 100644 --- a/services/nvpair-workload-manager/server.go +++ b/services/nvpair-workload-manager/server.go @@ -170,13 +170,23 @@ 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)) { + // + // 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) + var duplicate bool + if !resync { + duplicate, err = s.dedup.emitOnce(keyLifecycle(wl), func() error { return s.emitUpsert(wl) }) + } else { + err = s.emitUpsert(wl) + } + if duplicate { 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 { + 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. @@ -204,13 +214,17 @@ func (s *Server) handleRemove(w http.ResponseWriter, msg *Message) { return } - if s.dedup.seenOrAdd(keyRemove(nodeID, workloadID)) { + // 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 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..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 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. -- **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. +- **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. +- **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`.