From 56a6829956454652692021aaebb5d60573dc3afb Mon Sep 17 00:00:00 2001 From: woodsonl <65194841+woodsonl@users.noreply.github.com> Date: Tue, 22 Sep 2026 10:33:08 -0500 Subject: [PATCH] fix: dispatch broker requests on a bounded pool and coalesce relay deliveries MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The broker handled every inbound message inline on the read loop. Each handleMessage can run a synchronous worker relay bounded by rpcWorkerCallTimeout (5s), so one slow proxy/cluster/settings call blocked every other client request behind it. Dispatch through a small fixed pool instead. Cross-request ordering is preserved where it matters by the per-handler mutexes (workloadEmitMu for workload apply→fan→emit, per-state mutexes for subscription bookkeeping), and JSON-RPC does not promise cross-request response ordering — each response carries its own id, and the codec's write mutex keeps concurrent responses from interleaving. The relay also delivered to each subscriber with a blocking send, so a subscriber that stopped reading could stall the pump for everyone. Deliver now coalesces under the subscriber's own lock and drops the oldest pending message when a subscriber falls behind, so one stuck client cannot back-pressure the shared stream. The pump re-checks done before every Send, so a trigger that races Unsubscribe cannot deliver to a consumer that's gone; Unsubscribe itself does not wait for the pump to exit. The read loop also classifies read errors: a DecodeError is recoverable and keeps the loop alive, any other non-EOF error is terminal and ends the pump (errTerminalRead) instead of spinning. Signed-off-by: woodsonl <65194841+woodsonl@users.noreply.github.com> --- services/nvpair-ui-broker/broker.go | 86 ++++++++++++-- services/nvpair-ui-broker/codec.go | 9 +- services/nvpair-ui-broker/relay/relay.go | 101 ++++++++++++---- services/nvpair-ui-broker/relay/relay_test.go | 111 ++++++++++++++++-- services/nvpair-ui-broker/relaysub.go | 8 +- .../nvpair-ui-broker/terminal_read_test.go | 100 ++++++++++++++++ 6 files changed, 357 insertions(+), 58 deletions(-) create mode 100644 services/nvpair-ui-broker/terminal_read_test.go diff --git a/services/nvpair-ui-broker/broker.go b/services/nvpair-ui-broker/broker.go index c5da7fd1..9767c8a4 100644 --- a/services/nvpair-ui-broker/broker.go +++ b/services/nvpair-ui-broker/broker.go @@ -258,7 +258,7 @@ type Broker struct { // subMu guards subscribed. The discovery:nodes-changed stream is // opt-in: emitNodesChanged (called on the scanner-event goroutine) // reads this flag while the discovery:subscribe / discovery:unsubscribe - // handlers (called on the read-loop goroutine) flip it, so the two + // handlers (called on the dispatch-pool goroutine) flip it, so the two // goroutines need a lock between them. subMu sync.Mutex subscribed bool @@ -266,7 +266,7 @@ type Broker struct { // proxyMu guards every engineProxyRuntime.subscribed. The // : streams are opt-in like discovery's: the forward // hooks (on each proxy's reader goroutine) read the flag while the - // subscribe / unsubscribe handlers (on the read-loop goroutine) flip it. + // subscribe / unsubscribe handlers (on the dispatch-pool goroutine) flip it. proxyMu sync.Mutex // engineProxies holds per-engine proxy state, one entry per engine in the @@ -278,7 +278,7 @@ type Broker struct { // opt-in too: emitWorkloadEvent (called on the proxy reader goroutine // for local echoes and on the workload-manager reader goroutine for // peer-origin relays) reads the flag while the workloads:subscribe / - // workloads:unsubscribe handlers (on the read-loop goroutine) flip it. + // workloads:unsubscribe handlers (on the dispatch-pool goroutine) flip it. workloadsMu sync.Mutex workloadsSubscribed bool @@ -300,7 +300,7 @@ type Broker struct { // engineMu guards engineSubscribed. The engine: stream is // opt-in like proxy's: forwardEngineNotification (engine-manager reader // goroutine) reads the flag while the engine:subscribe / - // engine:unsubscribe handlers (read-loop goroutine) flip it. + // engine:unsubscribe handlers (dispatch-pool goroutine) flip it. engineMu sync.Mutex engineSubscribed bool @@ -2039,6 +2039,18 @@ func (b *Broker) runWorkloadHistoryFlusher(ctx context.Context) func() { } } +// errTerminalRead reports that the client stdin read loop ended on a +// non-recoverable scanner/transport error (e.g. an over-long frame that +// bufio.Scanner cannot resync past), distinct from a clean EOF. +var errTerminalRead = stderrors.New("terminal read error") + +// messageDispatchConcurrency is the size of the broker's inbound dispatch +// pool: enough worker goroutines that one slow synchronous worker relay +// (bounded by rpcWorkerCallTimeout) cannot head-of-line block the rest of +// the control plane, few enough that handlers stay effectively serialized +// under normal traffic. +const messageDispatchConcurrency = 4 + func (b *Broker) Serve(ctx context.Context) error { ctx, cancel := context.WithCancel(ctx) b.cancel = cancel @@ -2585,6 +2597,15 @@ func setNodeIDIfEmpty(m map[string]json.RawMessage, key, nodeID string) bool { return true } +// recoverableDecode reports whether a codec Read error is a recoverable +// per-frame decode failure (bad JSON / wrong version): the scanner advances +// past the bad frame, so both the producer and the consumer keep pumping +// instead of tearing the connection down. +func recoverableDecode(err error) bool { + var de *DecodeError + return stderrors.As(err, &de) +} + func (b *Broker) readLoop(ctx context.Context) error { // codec.Read() blocks on stdin, so we run it on its own goroutine and // select against ctx.Done(). Otherwise a SIGINT/SIGTERM (which cancels @@ -2604,14 +2625,47 @@ func (b *Broker) readLoop(ctx context.Context) error { case <-ctx.Done(): return } - // EOF is terminal (stream closed); other errors are per-line - // (e.g. a bad JSON frame) and the next Read advances past them. - if err == io.EOF { - return + // A decoded message (err nil) and a recoverable decode error + // (bad frame; the next Read advances past it) both keep the pump + // running. EOF is terminal (stream closed), and any other error + // is a terminal scanner/transport error: stop feeding the + // channel so the consumer exits instead of spinning. + if err == nil { + continue } + if recoverableDecode(err) { + continue + } + return } }() + // Bounded dispatch pool: handleMessage runs synchronous worker relays + // (proxy/cluster/settings/manual-nodes, each bounded by + // rpcWorkerCallTimeout) so dispatching on the read loop would let one + // slow worker stall every other client request for up to 5s. A small + // worker pool decouples them. Cross-request ordering is preserved for + // the channels that need it by dedicated mutexes inside the handlers + // (workloadEmitMu serializes workload apply→fan→emit; subscription + // bookkeeping is per-state mutexed), and JSON-RPC has no cross-request + // response-ordering guarantee — each response carries its own id. The + // codec's write mutex keeps concurrent responses from interleaving. + dispatch := make(chan *Message) + var dispatchWG sync.WaitGroup + for range messageDispatchConcurrency { + dispatchWG.Add(1) + go func() { + defer dispatchWG.Done() + for msg := range dispatch { + b.handleMessage(msg) + } + }() + } + defer func() { + close(dispatch) + dispatchWG.Wait() + }() + for { select { case <-ctx.Done(): @@ -2621,10 +2675,18 @@ func (b *Broker) readLoop(ctx context.Context) error { if r.err == io.EOF || ctx.Err() != nil { return nil } - slog.Warn("JSON-RPC read error", "err", r.err) - continue + if recoverableDecode(r.err) { + slog.Warn("JSON-RPC decode error (skipping frame)", "err", r.err) + continue + } + slog.Warn("JSON-RPC read error (terminal)", "err", r.err) + return errTerminalRead + } + select { + case dispatch <- r.msg: + case <-ctx.Done(): + return nil } - b.handleMessage(r.msg) if ctx.Err() != nil { return nil } @@ -3401,7 +3463,7 @@ func (b *Broker) handleMessage(msg *Message) { // timeout: engine lifecycle ops (install, model pull, ...) run for minutes // and report progress via push events, so the broker waits for the real // response asynchronously rather than fabricating a timeout — meanwhile -// other client requests keep being served on the read-loop goroutine. +// other client requests keep being served on the dispatch-pool goroutine. func (b *Broker) relayToEngine(msg *Message) { go b.relayToEngineNow(msg) } diff --git a/services/nvpair-ui-broker/codec.go b/services/nvpair-ui-broker/codec.go index 4d7dcc80..7815d286 100644 --- a/services/nvpair-ui-broker/codec.go +++ b/services/nvpair-ui-broker/codec.go @@ -16,10 +16,11 @@ import ( ) type ( - Message = jsonrpc.Message - RPCError = jsonrpc.RPCError - Codec = jsonrpc.Codec - Peer = jsonrpc.Peer + Message = jsonrpc.Message + RPCError = jsonrpc.RPCError + DecodeError = jsonrpc.DecodeError + Codec = jsonrpc.Codec + Peer = jsonrpc.Peer ) var ( diff --git a/services/nvpair-ui-broker/relay/relay.go b/services/nvpair-ui-broker/relay/relay.go index ed02240c..72e411d6 100644 --- a/services/nvpair-ui-broker/relay/relay.go +++ b/services/nvpair-ui-broker/relay/relay.go @@ -85,21 +85,26 @@ func registerEqual(a, b noderec.RegisterParams) bool { } // Subscriber is a client interested in directory changes: a service filter and a -// callback invoked (on the caller's goroutine, under no relay lock) with the -// subscriber's full filtered node set on every change. Consumers replace their -// set from it rather than applying deltas, so a dropped or reordered push can't -// leave them drifted — every push is the authoritative current list. +// callback invoked with the subscriber's full filtered node set on every change. +// Consumers replace their set from it rather than applying deltas, so a dropped +// or reordered push can't leave them drifted — every push is the authoritative +// current list. +// +// Deliveries are asynchronous: each subscriber owns a pump goroutine (started +// by Directory.Subscribe) that serializes its sends and coalesces concurrent +// triggers into one delivery that captures the snapshot at send time. A slow or +// blocked Send therefore stalls only its own subscriber, never the directory +// update path that feeds it (the scanner's read pump calls Apply). type Subscriber struct { Filter noderec.SubscribeParams Send func(nodes []noderec.DirectoryNode) - // sendMu serializes deliveries to this subscriber so two concurrent - // deliveries — the initial post-subscribe delivery racing an Apply fan-out - // driven by the scanner read-pump — can't reorder and leave the subscriber - // holding an older set than a newer one. Combined with capturing the snapshot - // inside Deliver (at send time, not subscribe time), the last delivery to - // acquire it always carries the latest directory state. - sendMu sync.Mutex + // kick carries a pending-delivery signal (capacity 1: extra signals while + // one is already pending coalesce — the pump captures the latest snapshot + // when it wakes, so early triggers can't deliver stale state). done closes + // on Unsubscribe and stops the pump. + kick chan struct{} + done chan struct{} } // Directory is the broker's view of all LAN nodes (keyed by hostUuid) plus its @@ -119,32 +124,69 @@ func NewDirectory() *Directory { } } -// Subscribe registers a subscriber and returns its id. The caller sends the -// initial snapshot via Deliver after releasing its own lock — Deliver captures -// the snapshot at send time, so a concurrent Apply can't sneak a newer snapshot -// in and have this initial delivery overwrite it with an older one. +// Subscribe registers a subscriber and starts its delivery pump goroutine. The +// initial snapshot arrives via the pump after any pending Deliver call — +// snapshot is captured at send time, so a concurrent Apply can't sneak a newer +// snapshot in and have this initial delivery overwrite it with an older one. func (d *Directory) Subscribe(sub *Subscriber) (id int) { + sub.kick = make(chan struct{}, 1) + sub.done = make(chan struct{}) d.mu.Lock() d.nextID++ id = d.nextID d.subs[id] = sub d.mu.Unlock() + go d.pump(sub) return id } -// Deliver pushes the subscriber its current filtered snapshot, serialized -// per-subscriber. Capturing the snapshot here (at delivery time) rather than -// handing Send a pre-captured slice means a delivery can never carry a set older -// than the directory's state when it actually runs; the per-subscriber lock then -// guarantees the initial post-subscribe delivery and a concurrent Apply fan-out -// settle on the latest set regardless of which runs last. +// pump serializes one subscriber's deliveries. Every wake re-captures the +// latest filtered snapshot, so coalesced triggers always deliver current state. +// done takes priority over a pending kick: once Unsubscribe has closed done, a +// trigger that raced the close must not produce a Send against a consumer +// that's gone. The pre-select alone is not enough (a token may already sit in +// kick), so the kick case re-checks done before sending. +func (d *Directory) pump(sub *Subscriber) { + for { + select { + case <-sub.done: + return + default: + } + select { + case <-sub.kick: + select { + case <-sub.done: + return + default: + } + sub.Send(d.filtered(sub.Filter)) + case <-sub.done: + return + } + } +} + +// Deliver asks for a delivery of the subscriber's current filtered snapshot. +// Non-blocking: it schedules the send on the subscriber's pump and never blocks +// the caller — Apply runs on the scanner read pump, and a subscriber whose Send +// blocks (a stalled worker's stdin pipe) must not stall the directory or the +// other subscribers. Multiple pending triggers coalesce into one send of the +// latest state. func (d *Directory) Deliver(sub *Subscriber) { - sub.sendMu.Lock() - defer sub.sendMu.Unlock() + select { + case sub.kick <- struct{}{}: + default: + } +} + +// filtered returns the nodes matching a subscriber's filter, sorted by +// hostUuid for a deterministic set. +func (d *Directory) filtered(f noderec.SubscribeParams) []noderec.DirectoryNode { d.mu.Lock() - nodes := d.filteredLocked(sub.Filter) + nodes := d.filteredLocked(f) d.mu.Unlock() - sub.Send(nodes) + return nodes } // filteredLocked returns the nodes matching a subscriber's filter, sorted by @@ -160,11 +202,18 @@ func (d *Directory) filteredLocked(f noderec.SubscribeParams) []noderec.Director return out } -// Unsubscribe removes a subscriber. +// Unsubscribe removes a subscriber and stops its delivery pump. It does not +// wait for the pump to exit: closing done is enough, because the pump's kick +// case re-checks done before every Send, so a trigger that raced the close +// cannot deliver to a consumer that's gone. func (d *Directory) Unsubscribe(id int) { d.mu.Lock() + sub := d.subs[id] delete(d.subs, id) d.mu.Unlock() + if sub != nil { + close(sub.done) + } } // Apply folds a daemon node-* delta into the directory, then re-sends every diff --git a/services/nvpair-ui-broker/relay/relay_test.go b/services/nvpair-ui-broker/relay/relay_test.go index bdad98ef..85409946 100644 --- a/services/nvpair-ui-broker/relay/relay_test.go +++ b/services/nvpair-ui-broker/relay/relay_test.go @@ -5,7 +5,9 @@ package relay import ( "reflect" + "sync" "testing" + "time" "nvpair-shared/noderec" ) @@ -53,20 +55,57 @@ func TestRegistrationCache(t *testing.T) { // full filtered snapshot; snaps holds them in order so a test can assert on the // latest set and on how many pushes arrived. type recordingSub struct { + mu sync.Mutex snaps [][]noderec.DirectoryNode } func (r *recordingSub) send(nodes []noderec.DirectoryNode) { + r.mu.Lock() + defer r.mu.Unlock() r.snaps = append(r.snaps, append([]noderec.DirectoryNode(nil), nodes...)) } -func (r *recordingSub) last() []noderec.DirectoryNode { - if len(r.snaps) == 0 { - return nil +// count returns the number of deliveries received so far. +func (r *recordingSub) count() int { + r.mu.Lock() + defer r.mu.Unlock() + return len(r.snaps) +} + +// last returns the latest snapshot, waiting up to 2s for at least want +// deliveries (deliveries are asynchronous: the pump goroutine sends them). +func (r *recordingSub) last(want int) []noderec.DirectoryNode { + deadline := time.Now().Add(2 * time.Second) + for r.count() < want { + if time.Now().After(deadline) { + r.mu.Lock() + defer r.mu.Unlock() + if len(r.snaps) == 0 { + return nil + } + return r.snaps[len(r.snaps)-1] + } + time.Sleep(2 * time.Millisecond) } + r.mu.Lock() + defer r.mu.Unlock() return r.snaps[len(r.snaps)-1] } +// lastWhere waits up to 2s for the latest snapshot to satisfy pred, then +// returns it. Use when coalescing makes the delivery count nondeterministic but +// the settled content is what the test cares about. +func (r *recordingSub) lastWhere(pred func([]noderec.DirectoryNode) bool) []noderec.DirectoryNode { + deadline := time.Now().Add(2 * time.Second) + for { + got := r.last(1) + if pred(got) || time.Now().After(deadline) { + return got + } + time.Sleep(2 * time.Millisecond) + } +} + func ids(nodes []noderec.DirectoryNode) []string { out := make([]string, len(nodes)) for i, n := range nodes { @@ -96,7 +135,7 @@ func TestDirectorySubscribeInitialSnapshot(t *testing.T) { } d.Subscribe(sub) d.Deliver(sub) - if got := ids(rec.last()); !reflect.DeepEqual(got, []string{"a"}) { + if got := ids(rec.last(1)); !reflect.DeepEqual(got, []string{"a"}) { t.Fatalf("initial delivery = %v, want [a]", got) } } @@ -112,7 +151,7 @@ func TestDeliverCapturesAtSendTime(t *testing.T) { d.Subscribe(sub) d.Apply(noderec.NotifyNodeDiscovered, olNode("a")) d.Deliver(sub) - if got := ids(rec.last()); !reflect.DeepEqual(got, []string{"a"}) { + if got := ids(rec.last(1)); !reflect.DeepEqual(got, []string{"a"}) { t.Fatalf("delivery after a post-subscribe change = %v, want [a]", got) } } @@ -128,11 +167,16 @@ func TestDirectoryFanoutRespectsFilter(t *testing.T) { d.Apply(noderec.NotifyNodeDiscovered, erNode("b")) // Every change re-pushes each subscriber its full filtered snapshot, so the - // latest snapshot is the authoritative filtered set. - if got := ids(olSub.last()); !reflect.DeepEqual(got, []string{"a"}) { + // latest snapshot is the authoritative filtered set. Coalescing can collapse + // the two pushes into one, so settle on content, not a delivery count. + if got := ids(olSub.lastWhere(func(n []noderec.DirectoryNode) bool { + return reflect.DeepEqual(ids(n), []string{"a"}) + })); !reflect.DeepEqual(got, []string{"a"}) { t.Errorf("ol subscriber last snapshot = %v, want [a]", got) } - if got := ids(allSub.last()); !reflect.DeepEqual(got, []string{"a", "b"}) { + if got := ids(allSub.lastWhere(func(n []noderec.DirectoryNode) bool { + return reflect.DeepEqual(ids(n), []string{"a", "b"}) + })); !reflect.DeepEqual(got, []string{"a", "b"}) { t.Errorf("all subscriber last snapshot = %v, want [a b]", got) } } @@ -148,16 +192,17 @@ func TestDirectoryRemoveAndUnsubscribe(t *testing.T) { t.Error("node should be gone after removed") } // The removal re-pushes an empty snapshot (the node is simply absent). - if got := sub.last(); len(got) != 0 { + if got := sub.lastWhere(func(n []noderec.DirectoryNode) bool { return len(n) == 0 }); len(got) != 0 { t.Errorf("subscriber last snapshot = %v, want empty after removal", ids(got)) } // After unsubscribe, no more pushes. - before := len(sub.snaps) + before := sub.count() d.Unsubscribe(id) d.Apply(noderec.NotifyNodeDiscovered, olNode("c")) - if len(sub.snaps) != before { - t.Errorf("unsubscribed sub still received %d pushes", len(sub.snaps)-before) + time.Sleep(50 * time.Millisecond) + if sub.count() != before { + t.Errorf("unsubscribed sub still received %d pushes", sub.count()-before) } } @@ -176,3 +221,45 @@ func TestDirectorySnapshotFilterAndSort(t *testing.T) { t.Fatalf("snapshot(ol) = %d, want 2", len(ol)) } } + +// TestDeliverCoalescesTriggers guards the pump contract: Deliver is non-blocking +// and multiple pending triggers coalesce into ONE send of the latest state — a +// subscriber with a slow Send must never accumulate a backlog of stale snapshots. +func TestDeliverCoalescesTriggers(t *testing.T) { + d := NewDirectory() + rec := &recordingSub{} + sub := &Subscriber{Filter: noderec.SubscribeParams{}, Send: rec.send} + id := d.Subscribe(sub) + + d.Apply(noderec.NotifyNodeDiscovered, olNode("a")) + d.Apply(noderec.NotifyNodeDiscovered, olNode("b")) + // Three triggers while none has been consumed yet: they must coalesce. + d.Deliver(sub) + d.Deliver(sub) + d.Deliver(sub) + + got := rec.last(1) + if !reflect.DeepEqual(ids(got), []string{"a", "b"}) { + t.Fatalf("coalesced delivery = %v, want [a b] (latest state)", got) + } + + // The next Apply re-pushes current state even though its trigger coalesced + // with the earlier ones — the pump always re-captures at wake time. + d.Apply(noderec.NotifyNodeDiscovered, erNode("c")) + if got := ids(rec.lastWhere(func(n []noderec.DirectoryNode) bool { + return reflect.DeepEqual(ids(n), []string{"a", "b", "c"}) + })); !reflect.DeepEqual(got, []string{"a", "b", "c"}) { + t.Fatalf("delivery after coalesced apply = %v, want [a b c]", got) + } + + // After Unsubscribe the pump exits: Deliver must not panic and no further + // sends may arrive. Counts are not exact (coalescing is timing-dependent), + // so assert the count is stable, not equal to a specific number. + d.Unsubscribe(id) + before := rec.count() + d.Deliver(sub) + time.Sleep(50 * time.Millisecond) + if n := rec.count(); n != before { + t.Errorf("post-unsubscribe deliveries grew from %d to %d, want no further sends", before, n) + } +} diff --git a/services/nvpair-ui-broker/relaysub.go b/services/nvpair-ui-broker/relaysub.go index f1ec0f0e..2bc48c7a 100644 --- a/services/nvpair-ui-broker/relaysub.go +++ b/services/nvpair-ui-broker/relaysub.go @@ -19,10 +19,10 @@ type relaySendFunc func(nodes []noderec.DirectoryNode) // subscribeRelay registers a worker as a relay.Directory subscriber for the // given discovery:subscribe params and returns the subscription id plus the // subscriber handle. The caller owns the id's lifetime (Unsubscribe on worker -// exit / re-subscribe) and must send the initial snapshot via dir.Deliver(sub) -// after releasing its own lock; Deliver captures the snapshot at send time, so -// the initial delivery can't be overtaken by a concurrent Apply and land a stale -// set. +// exit / re-subscribe) and must request the initial snapshot via +// dir.Deliver(sub); Deliver schedules the send on the subscriber's pump, which +// captures the snapshot at send time, so the initial delivery can't be +// overtaken by a concurrent Apply and land a stale set. func subscribeRelay(dir *relay.Directory, params json.RawMessage, send relaySendFunc) (int, *relay.Subscriber, error) { var sp noderec.SubscribeParams if err := json.Unmarshal(params, &sp); err != nil { diff --git a/services/nvpair-ui-broker/terminal_read_test.go b/services/nvpair-ui-broker/terminal_read_test.go new file mode 100644 index 00000000..5435cd13 --- /dev/null +++ b/services/nvpair-ui-broker/terminal_read_test.go @@ -0,0 +1,100 @@ +// SPDX-FileCopyrightText: Copyright (c) 2026 NVIDIA CORPORATION & AFFILIATES. All rights reserved. +// SPDX-License-Identifier: Apache-2.0 + +package main + +import ( + "context" + "errors" + "io" + "net" + "testing" + "time" +) + +// errFakeTerminal is a non-EOF, non-DecodeError read failure — the shape a +// terminal scanner/transport error takes at the codec boundary. +var errFakeTerminal = errors.New("fake terminal read error") + +// errorReader fails every Read with errFakeTerminal; wrapped in a real Codec it +// stands in for a dead terminal transport. +type errorReader struct{} + +func (errorReader) Read([]byte) (int, error) { return 0, errFakeTerminal } + +func newTerminalErrorCodec() *Codec { + return NewCodec(struct { + io.Reader + io.Writer + }{errorReader{}, io.Discard}) +} + +// TestReadLoopTerminalReadErrorStopsPump guards the read-loop contract: a +// non-EOF transport error must end Serve-style pumping with errTerminalRead +// (the producer goroutine has already stopped; spinning would burn CPU forever), +// while EOF and plain decode errors keep the loop alive. +func TestReadLoopTerminalReadErrorStopsPump(t *testing.T) { + t.Run("non-EOF transport error is terminal", func(t *testing.T) { + b := &Broker{codec: newTerminalErrorCodec()} + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + done := make(chan error, 1) + go func() { done <- b.readLoop(ctx) }() + + select { + case err := <-done: + if !errors.Is(err, errTerminalRead) { + t.Fatalf("readLoop err = %v, want errTerminalRead", err) + } + case <-time.After(5 * time.Second): + t.Fatal("readLoop did not stop on transport error") + } + }) + + t.Run("EOF is clean exit", func(t *testing.T) { + client, server := net.Pipe() + defer client.Close() + defer server.Close() + b := &Broker{codec: NewCodec(client)} + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + done := make(chan error, 1) + go func() { done <- b.readLoop(ctx) }() + server.Close() // producer sees EOF + + select { + case err := <-done: + if err != nil { + t.Fatalf("readLoop err on EOF = %v, want nil", err) + } + case <-time.After(5 * time.Second): + t.Fatal("readLoop did not stop on EOF") + } + }) + + t.Run("recoverable decode error keeps the loop alive", func(t *testing.T) { + client, server := net.Pipe() + defer client.Close() + defer server.Close() + b := &Broker{codec: NewCodec(client)} + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + done := make(chan error, 1) + go func() { done <- b.readLoop(ctx) }() + if _, err := io.WriteString(server, "not-json\n"); err != nil { + t.Fatal(err) + } + // Keep the pipe open: a closed pipe is a terminal read error, while + // the bad frame must only surface as a recoverable DecodeError. + time.Sleep(100 * time.Millisecond) + select { + case err := <-done: + t.Fatalf("readLoop exited early on decode error: %v", err) + default: + } + cancel() + }) +}