From 8c67ddaf9ec88589ec4f926a34da082a9b7bfdf8 Mon Sep 17 00:00:00 2001 From: Lee Calcote Date: Mon, 6 Jul 2026 02:16:37 -0500 Subject: [PATCH 1/2] [MeshSync] Unsubscribe exec input on teardown; drop leaked drain goroutine streamSession subscribed to input. for a session's stdin but could not tear that subscription down, because broker.Handler had no Unsubscribe. It worked around that by parking a `<-done; for range subCh {}` drain goroutine - which never exited (subCh is never closed) and never released the subscription, leaking a goroutine and a NATS subscription per exec session for the process lifetime. MeshKit v1.0.22 (meshery/meshkit#1056) adds broker.Handler.Unsubscribe. terminate() now calls Unsubscribe("input."), which releases the subscription and the broker's delivery goroutine, and the drain goroutine is removed. subCh gets a 1-slot buffer so a delivery already in flight at teardown cannot block the broker's delivery goroutine in the window before Unsubscribe takes effect. Fixes part 2 of #585 (part 1, the channelPool data race, landed in #587). - go.mod: meshkit v1.0.20 -> v1.0.22 - meshsync/exec.go: Unsubscribe on teardown via unsubscribeSessionInput; remove drain goroutine; buffer subCh(1); replace stale cyclop TODO with rationale - meshsync/exec_test.go: teardown/leak regression + error-path tests (-race) - docs: architecture.md interactive-sessions section; fd4 status note (base Unsubscribe shipped as 1-arg in v1.0.22, exec leak fixed) Signed-off-by: Lee Calcote --- docs/agent-instructions/architecture.md | 6 + docs/design/fd4-jetstream-durable-delivery.md | 6 +- go.mod | 4 +- go.sum | 8 +- meshsync/exec.go | 50 ++++-- meshsync/exec_test.go | 145 ++++++++++++++++++ 6 files changed, 196 insertions(+), 23 deletions(-) create mode 100644 meshsync/exec_test.go diff --git a/docs/agent-instructions/architecture.md b/docs/agent-instructions/architecture.md index 891111c0..ddac6fc3 100644 --- a/docs/agent-instructions/architecture.md +++ b/docs/agent-instructions/architecture.md @@ -57,6 +57,12 @@ main.go --parses CLI flags--> pkg/lib/meshsync.Run(...) - `channel.go` / `generic.go` / `system.go` / `broker.go` define small typed channels (e.g. `StructChannel`) used for coordination (stop signals, broker request/response) between the handler and the pipeline - not a general pub/sub system. +## Interactive Sessions (exec / log stream) + +- Meshery Server routes interactive `kubectl exec` and pod-log requests to MeshSync over the broker; `meshsync/exec.go` (`processExecRequest`) and `meshsync/logstream.go` (`processLogRequest`) start one long-lived goroutine per request, keyed by a request id. +- These per-session channels live in a **`sync.Mutex`-guarded `sessions` map** on the Handler (`meshsync/sessions.go`), deliberately separate from `channelPool`, which holds only the fixed system channels (`Stop`/`OS`/`ReSync`) and is read-only after construction. Keeping them apart avoids the concurrent map read/write panic that occurred when session goroutines mutated the same map other goroutines ranged. +- An exec session subscribes to its own `input.` subject (client keystrokes) via `SubscribeWithChannel`. On teardown - stream EOF/error, an explicit stop request, or the global `Stop` - `terminate()` runs once (guarded by a `sync.Once`) and calls `broker.Handler.Unsubscribe("input.")`, which releases the subscription and the broker's delivery goroutine. Before MeshKit exposed `Unsubscribe` (v1.0.22), the subscription could not be torn down and each session parked a drain goroutine that never exited, leaking a goroutine and a subscription per session. + ## Deployment Topology - Meshery Operator's `MeshSync` controller (`meshery/meshery-operator`, `pkg/meshsync/meshsync.go`) renders this binary as a `Deployment` and injects `BROKER_URL` pointing at the Broker's derived NATS endpoint. See that repo's architecture doc for the reconcile side. diff --git a/docs/design/fd4-jetstream-durable-delivery.md b/docs/design/fd4-jetstream-durable-delivery.md index c800f8a3..63e58768 100644 --- a/docs/design/fd4-jetstream-durable-delivery.md +++ b/docs/design/fd4-jetstream-durable-delivery.md @@ -1,5 +1,7 @@ # Blueprint: Durable Delivery via NATS JetStream +> **Status (2026-07): base `Unsubscribe` shipped; exec leak fixed.** The adjacent interface gap this blueprint calls out - `broker.Handler` having no `Unsubscribe` - has since been closed independently of the JetStream work. MeshKit v1.0.22 ([meshery/meshkit#1056](https://github.com/meshery/meshkit/pull/1056)) added **`Unsubscribe(subject string) error`** to `broker.Handler`, implemented for both `Nats` (per-subject `*nats.Subscription` tracking) and `ChannelBrokerHandler` (closes the per-queue delivery channels). Note the signature is a single `subject` - it tears down every queue subscription for that subject - not the `(subject, queue)` pair sketched in sections 4-6 below; read those references as the shipped one-argument form. MeshSync now `Unsubscribe`s its per-session `input.` exec subject on teardown ([meshery/meshsync#585](https://github.com/meshery/meshsync/issues/585)), so the exec-subscription/goroutine leak noted below is fixed. What remains forward work here is durable JetStream delivery itself (the `DurablePublisher`/`DurableSubscriber` capability interfaces); the base-`Unsubscribe` prerequisite is done. + ## 1. Goal & Gap **Goal.** Make MeshSync -> Meshery Server resource-event delivery at-least-once and durable: an event published while Meshery Server is disconnected must not be silently dropped; it must be redelivered once the consumer (or a fresh consumer identity, post-restart) reconnects. @@ -9,12 +11,12 @@ - `broker/nats/nats.go:143-167` (MeshKit `Nats.Publish`) is a straight `nc.Publish(subject, data)` - core NATS semantics: if no subscriber is connected to receive the message at publish time, the message is gone forever. - `server/models/meshsync_events.go:126-131` subscribes with `SubscribeWithChannel("meshery.meshsync.core", "", out)` - again core NATS, no ack, no consumer state, no replay. - Recovery today is **full re-list only**: `broker.ReSyncDiscoveryEntity` (`meshsync/handlers.go:203-204,256-258`) tears down and rebuilds every informer, republishing everything. There is no periodic/automatic trigger for this today - only an explicit request (`server/models/meshsync_events.go:378-392` `Resync()`) or a CRD-detected discovery-config change (`meshsync/handlers.go:444-457`). A future periodic-reconcile feature would call this on a timer; durability must be designed to *complement*, not replace, that belt-and-suspenders path. -- Adjacent interface gap: `broker.Handler` (`meshkit/broker/broker.go:12-15`) has no `Unsubscribe`. `ListenToRequests` (`meshsync/handlers.go:142-185`) registers one permanent `SubscribeWithChannel` for the process lifetime and never tears it down - not a per-call leak today, but it means neither MeshSync nor Server can ever cleanly release a subscription (durable JetStream consumers need exactly this to drain/rebind on config change or shutdown). +- Adjacent interface gap (**since closed - see status note above**): `broker.Handler` had no `Unsubscribe` when this was written; MeshKit v1.0.22 added `Unsubscribe(subject string)`. `ListenToRequests` (`meshsync/handlers.go:142-185`) still registers one permanent `SubscribeWithChannel` for the process lifetime and never tears it down - not a per-call leak today, but durable JetStream consumers will use `Unsubscribe` to drain/rebind on config change or shutdown. ## 2. Current State Per Repo (grounded) **MeshKit** (`meshkit`) -- `broker/broker.go:7-26` - `Handler` interface: `Publish`, `PublishWithChannel`, `Subscribe`, `SubscribeWithChannel`, `Info`, `DeepCopyObject`, `DeepCopyInto`, `IsEmpty`, `CloseConnection`, `ConnectedEndpoints`. No ack/nak, no `Unsubscribe`. +- `broker/broker.go` - `Handler` interface: `Publish`, `PublishWithChannel`, `Subscribe`, `SubscribeWithChannel`, `Unsubscribe` (added in v1.0.22), `Info`, `DeepCopyObject`, `DeepCopyInto`, `IsEmpty`, `CloseConnection`, `ConnectedEndpoints`, `IsConnected`. No ack/nak. - `broker/nats/nats.go` - `Nats` struct wraps a single `*nats.Conn` (core NATS). `New(opts Options)` (line 41) sets `ReconnectWait`/`MaxReconnects`/handlers but never touches JetStream. - `broker/channel/channel.go` - a second, in-process `Handler` implementation (`ChannelBrokerHandler`) used in tests/library mode; any interface change must keep this satisfying the interface too (it already stubs `Subscribe` as a no-op, line 154-163). - `go.mod:30` already pins `github.com/nats-io/nats.go v1.47.0`, which ships the modern `github.com/nats-io/nats.go/jetstream` package - **no new external dependency needed** for JetStream client support. diff --git a/go.mod b/go.mod index 9001df3e..6daa2d8c 100644 --- a/go.mod +++ b/go.mod @@ -7,7 +7,7 @@ go 1.26.4 require ( github.com/buger/jsonparser v1.2.0 github.com/google/uuid v1.6.0 - github.com/meshery/meshkit v1.0.20 + github.com/meshery/meshkit v1.0.22 github.com/myntra/pipeline v0.0.0-20180618182531-2babf4864ce8 github.com/sirupsen/logrus v1.9.4 github.com/spf13/viper v1.21.0 @@ -108,7 +108,7 @@ require ( github.com/mattn/go-isatty v0.0.20 // indirect github.com/mattn/go-runewidth v0.0.19 // indirect github.com/mattn/go-sqlite3 v1.14.32 // indirect - github.com/meshery/schemas v1.3.23 // indirect + github.com/meshery/schemas v1.3.25 // indirect github.com/mitchellh/copystructure v1.2.0 // indirect github.com/mitchellh/go-wordwrap v1.0.1 // indirect github.com/mitchellh/reflectwalk v1.0.2 // indirect diff --git a/go.sum b/go.sum index d3df943f..fc647082 100644 --- a/go.sum +++ b/go.sum @@ -255,10 +255,10 @@ github.com/mattn/go-runewidth v0.0.19/go.mod h1:XBkDxAl56ILZc9knddidhrOlY5R/pDhg github.com/mattn/go-sqlite3 v1.14.22/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y= github.com/mattn/go-sqlite3 v1.14.32 h1:JD12Ag3oLy1zQA+BNn74xRgaBbdhbNIDYvQUEuuErjs= github.com/mattn/go-sqlite3 v1.14.32/go.mod h1:Uh1q+B4BYcTPb+yiD3kU8Ct7aC0hY9fxUwlHK0RXw+Y= -github.com/meshery/meshkit v1.0.20 h1:5hr8HzdF+WyeURl87DsWufOMfC7+seIWyuIIg6Paa8Y= -github.com/meshery/meshkit v1.0.20/go.mod h1:bKbYWZHFpc86Paq5Bs5IT8E7xeyG5DC4baKkgAiLm+Y= -github.com/meshery/schemas v1.3.23 h1:bei/pHj4jG+9WeZHQF7sS6CmfWJA8bE9F+3VMzq+Sas= -github.com/meshery/schemas v1.3.23/go.mod h1:3/0nDQZur8swpo+6+3b3pvYmkzUcdWs36xTQBGRU2Y8= +github.com/meshery/meshkit v1.0.22 h1:c8CgO2+mdF6krvAdA66FF5ScmIHBkTEVt3MtkUmOJuU= +github.com/meshery/meshkit v1.0.22/go.mod h1:2SzXjDsxZOraA9YbnqFzHS+1Tk14zsLuePJzHuZ1SRU= +github.com/meshery/schemas v1.3.25 h1:aUZSol0+kfVpFs4WRsKAD4OaqW0MQkAPQJy+JwkNe+Q= +github.com/meshery/schemas v1.3.25/go.mod h1:3/0nDQZur8swpo+6+3b3pvYmkzUcdWs36xTQBGRU2Y8= github.com/miekg/dns v1.1.57 h1:Jzi7ApEIzwEPLHWRcafCN9LZSBbqQpxjt/wpgvg7wcM= github.com/miekg/dns v1.1.57/go.mod h1:uqRjCRUuEAA6qsOiJvDd+CFo/vW+y5WR6SNmHE55hZk= github.com/mitchellh/copystructure v1.2.0 h1:vpKXTN4ewci03Vljg/q9QvCGUDttBOGBIa15WveJJGw= diff --git a/meshsync/exec.go b/meshsync/exec.go index 37cd64cf..d5e52d86 100644 --- a/meshsync/exec.go +++ b/meshsync/exec.go @@ -73,8 +73,9 @@ func (h *Handler) processExecRequest(obj interface{}, cfg config.ListenerConfig) for _, req := range reqs { id := fmt.Sprintf("exec.%s.%s.%s.%s", req.Namespace, req.Name, req.Container, req.ID) if bool(req.Stop) { - // Stop request: tear down the session if one is running (no-op otherwise). - // TODO: once we have unsubscribe functionality, publish to active sessions subject. + // Stop request: tear down the running session if any (no-op otherwise). + // The session's terminate() unsubscribes its input. subject; the + // 10s streamChannelPool ticker republishes the active-session list. execCleanup(h, id) continue } @@ -139,16 +140,26 @@ func (h *Handler) streamChannelPool() { }() } -// TODO fix cyclop error -// Error: meshsync/exec.go:113:1: calculated cyclomatic complexity for function streamSession is 15, max is 10 (cyclop) +// streamSession wires up one exec session: its input subscription and teardown, +// the TTY exec stream, stdout publishing, and the input loop. Its branch count +// (the four-case select plus the teardown/streaming closures) exceeds cyclop's +// default; the body is a linear setup-then-loop sequence, so splitting it would +// scatter the shared pipes/channels and hurt readability more than it helps. // //nolint:cyclop func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.ListenerConfig) { - subCh := make(chan *broker.Message) + // Buffer one message so a delivery already in flight at teardown lands in the + // buffer instead of blocking the broker's delivery goroutine on an unread + // channel in the window before Unsubscribe (in terminate) takes effect. + subCh := make(chan *broker.Message, 1) tstdin, putStdin := io.Pipe() stdin := io.NopCloser(tstdin) getStdout, stdout := io.Pipe() + // inputSubject is the per-session subject the client publishes stdin to; it is + // subscribed below and torn down in terminate(). + inputSubject := fmt.Sprintf("input.%s", id) + // done is closed exactly once when the session ends (stream EOF/error, // explicit Stop, or the global channels.Stop). Closing the pipes unblocks the // stdout reader and any in-flight stdin write, so no goroutine is left blocked. @@ -157,6 +168,13 @@ func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.Lis terminate := func() { once.Do(func() { close(done) + // Tear down the input subscription so the broker stops delivering to + // subCh and releases the delivery goroutine it started for it. The main + // loop has already stopped reading subCh (it returns on done), so + // without this the subscription and its goroutine would leak for the + // process lifetime. Unsubscribe is idempotent, so the repeated + // terminate() calls (defer + TTY goroutine) are safe. + h.unsubscribeSessionInput(inputSubject) // Closing both ends of both pipes unblocks the TTY streamer, the // stdout reader (Read returns io.ErrClosedPipe) and any stdin writer. _ = putStdin.Close() @@ -168,19 +186,10 @@ func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.Lis } defer terminate() - if err := h.Broker.SubscribeWithChannel(fmt.Sprintf("input.%s", id), generateID(), subCh); err != nil { + if err := h.Broker.SubscribeWithChannel(inputSubject, generateID(), subCh); err != nil { h.Log.Error(ErrExecTerminal(err)) } - // The broker interface exposes no Unsubscribe, so the input. subscription - // cannot be torn down here. Once the session ends, keep draining subCh so the - // broker's delivery goroutine never blocks on an unread channel after we stop. - go func() { - <-done - for range subCh { - } - }() - // Put the terminal into raw mode to prevent it echoing characters twice. t := term.TTY{ Parent: interrupt.New(func(s os.Signal) {}), @@ -289,6 +298,17 @@ func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.Lis } } +// unsubscribeSessionInput tears down an exec session's stdin subscription +// (input.) so the broker stops delivering to its channel and releases the +// delivery goroutine it started for it. It runs during session teardown, so the +// error is logged rather than returned; Unsubscribe is a no-op for a subject +// with no active subscription and is safe to call more than once. +func (h *Handler) unsubscribeSessionInput(subject string) { + if err := h.Broker.Unsubscribe(subject); err != nil { + h.Log.Error(ErrExecTerminal(err)) + } +} + func execCleanup(h *Handler, id string) { structChan, ok := h.getSession(id) if !ok { diff --git a/meshsync/exec_test.go b/meshsync/exec_test.go new file mode 100644 index 00000000..492b115b --- /dev/null +++ b/meshsync/exec_test.go @@ -0,0 +1,145 @@ +package meshsync + +import ( + "context" + "runtime" + "strings" + "sync" + "testing" + "time" + + "github.com/meshery/meshkit/broker" + "github.com/meshery/meshkit/broker/channel" + "github.com/meshery/meshkit/logger" +) + +// recordingBroker wraps a real broker.Handler, recording the subjects passed to +// Unsubscribe (and optionally forcing it to fail) while delegating every other +// method - including the actual teardown - to the embedded handler. +type recordingBroker struct { + broker.Handler + mu sync.Mutex + unsubscribed []string + failWith error +} + +// Compile-time proof that the fake satisfies the (extended) broker.Handler +// interface; if a new method is added to the interface this fails to build. +var _ broker.Handler = (*recordingBroker)(nil) + +func (r *recordingBroker) Unsubscribe(subject string) error { + r.mu.Lock() + r.unsubscribed = append(r.unsubscribed, subject) + r.mu.Unlock() + if r.failWith != nil { + return r.failWith + } + return r.Handler.Unsubscribe(subject) +} + +func (r *recordingBroker) subjects() []string { + r.mu.Lock() + defer r.mu.Unlock() + return append([]string(nil), r.unsubscribed...) +} + +func newTestLogger(t *testing.T) logger.Handler { + t.Helper() + log, err := logger.New("meshsync-test", logger.Options{}) + if err != nil { + t.Fatalf("logger.New: %v", err) + } + return log +} + +// TestUnsubscribeSessionInputTearsDownSubscription is the regression test for +// the exec input-subscription/goroutine leak (meshery/meshsync#585): before +// broker.Handler exposed Unsubscribe, streamSession could not tear down its +// input. subscription and parked a drain goroutine that never exited, +// leaking a subscription and a goroutine per exec session. Teardown must now +// unsubscribe the per-session subject, which releases the broker's delivery +// goroutine. +func TestUnsubscribeSessionInputTearsDownSubscription(t *testing.T) { + rec := &recordingBroker{Handler: channel.NewChannelBrokerHandler()} + h := &Handler{Broker: rec, Log: newTestLogger(t)} + + const id = "exec.ns.pod.ctr.req-1" + subject := "input." + id + // Mirror streamSession: a 1-buffered channel subscribed to input.. + subCh := make(chan *broker.Message, 1) + + baseline := runtime.NumGoroutine() + + if err := rec.SubscribeWithChannel(subject, generateID(), subCh); err != nil { + t.Fatalf("SubscribeWithChannel: %v", err) + } + + // Sanity: a published input message reaches the session channel while the + // subscription is live. + if err := rec.Publish(subject, &broker.Message{ObjectType: broker.ExecInputObject, Object: "hi"}); err != nil { + t.Fatalf("Publish: %v", err) + } + select { + case msg := <-subCh: + if got, _ := msg.Object.(string); got != "hi" { + t.Fatalf("delivered %q, want %q", got, "hi") + } + case <-time.After(2 * time.Second): + t.Fatal("input message was not delivered before Unsubscribe") + } + + // Behavior under test: session teardown unsubscribes the per-session subject. + h.unsubscribeSessionInput(subject) + + if got := rec.subjects(); len(got) != 1 || got[0] != subject { + t.Fatalf("Unsubscribe called with %v, want [%s]", got, subject) + } + + // The real broker actually tore the subscription down: nothing is left + // registered on the subject. + for _, ep := range rec.ConnectedEndpoints() { + if strings.HasPrefix(ep, subject+"::") { + t.Fatalf("subject %s still registered after Unsubscribe: %v", subject, rec.ConnectedEndpoints()) + } + } + + // ...and the delivery goroutine it started has exited, so nothing leaks. + waitForGoroutines(t, baseline, 3*time.Second) +} + +// TestUnsubscribeSessionInputLogsErrorWithoutPanic covers the teardown error +// path: unsubscribeSessionInput runs during session cleanup, so an Unsubscribe +// error must be logged rather than propagated or panicked on. +func TestUnsubscribeSessionInputLogsErrorWithoutPanic(t *testing.T) { + rec := &recordingBroker{ + Handler: channel.NewChannelBrokerHandler(), + failWith: context.Canceled, // stand-in broker failure + } + h := &Handler{Broker: rec, Log: newTestLogger(t)} + + // Must not panic even though Unsubscribe fails. + h.unsubscribeSessionInput("input.exec.ns.pod.ctr.req-err") + + if got := rec.subjects(); len(got) != 1 { + t.Fatalf("expected exactly one Unsubscribe call, got %v", got) + } +} + +// waitForGoroutines fails if the goroutine count has not returned to at most +// baseline within timeout. It guards against the per-session delivery-goroutine +// leak that the old drain-goroutine approach caused. +func waitForGoroutines(t *testing.T, baseline int, timeout time.Duration) { + t.Helper() + deadline := time.Now().Add(timeout) + for { + runtime.Gosched() + if runtime.NumGoroutine() <= baseline { + return + } + if time.Now().After(deadline) { + t.Fatalf("goroutine count did not return to baseline %d (now %d): delivery goroutine leaked", + baseline, runtime.NumGoroutine()) + } + time.Sleep(10 * time.Millisecond) + } +} From 246ed6c9ce94e889037e9e37178d56ae1b4dde93 Mon Sep 17 00:00:00 2001 From: Lee Calcote Date: Mon, 6 Jul 2026 02:31:17 -0500 Subject: [PATCH 2/2] [MeshSync] exec teardown: harden input receive, buffer bursts, goleak test Review follow-ups on #591: - Main input loop reads subCh with the comma-ok form and guards the payload assertion, so a broker that closes its delivery channel (returns nil,false) or a malformed non-string ExecInput message ends/skips the session instead of panicking or spinning. (subCh is not closed by the current NATS/channel brokers, but the receive is now robust if one ever does.) - Widen the input channel buffer from 1 to execInputChannelBuffer (256) so a burst of stdin - or a delivery already in flight at teardown - cannot backpressure the broker's delivery goroutine in the window before Unsubscribe stops delivery. - Replace the flaky runtime.NumGoroutine leak assertion with go.uber.org/goleak (VerifyNone + IgnoreCurrent), which retries with backoff and ignores runtime background goroutines; verified it fails when teardown is skipped. Signed-off-by: Lee Calcote --- go.mod | 1 + meshsync/exec.go | 37 +++++++++++++++++++++++++++++-------- meshsync/exec_test.go | 39 ++++++++++++--------------------------- 3 files changed, 42 insertions(+), 35 deletions(-) diff --git a/go.mod b/go.mod index 6daa2d8c..da471dcf 100644 --- a/go.mod +++ b/go.mod @@ -12,6 +12,7 @@ require ( github.com/sirupsen/logrus v1.9.4 github.com/spf13/viper v1.21.0 github.com/stretchr/testify v1.11.1 + go.uber.org/goleak v1.3.0 golang.org/x/exp v0.0.0-20250106191152-7588d65b2ba8 gorm.io/gorm v1.31.2 gotest.tools/v3 v3.5.2 diff --git a/meshsync/exec.go b/meshsync/exec.go index d5e52d86..8b57d7a6 100644 --- a/meshsync/exec.go +++ b/meshsync/exec.go @@ -41,6 +41,14 @@ import ( // KB stands for KiloByte const KB = 1024 +// execInputChannelBuffer bounds how many stdin messages can queue for a session +// before the broker's delivery goroutine backpressures. A cushion (rather than +// an unbuffered channel) keeps a burst of input - or a delivery already in +// flight when the session tears down - from blocking that delivery goroutine in +// the window before Unsubscribe stops further delivery. It is a cushion, not a +// correctness dependency: teardown unsubscribes the subject regardless. +const execInputChannelBuffer = 256 + // terminalSizeQueueAdapter adapts kubectl's term.TerminalSizeQueue to client-go's remotecommand.TerminalSizeQueue type terminalSizeQueueAdapter struct { queue term.TerminalSizeQueue @@ -148,10 +156,11 @@ func (h *Handler) streamChannelPool() { // //nolint:cyclop func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.ListenerConfig) { - // Buffer one message so a delivery already in flight at teardown lands in the - // buffer instead of blocking the broker's delivery goroutine on an unread - // channel in the window before Unsubscribe (in terminate) takes effect. - subCh := make(chan *broker.Message, 1) + // Buffered (see execInputChannelBuffer) so a burst of stdin - or a delivery + // already in flight at teardown - lands in the buffer instead of blocking the + // broker's delivery goroutine on an unread channel in the window before + // Unsubscribe (in terminate) takes effect. + subCh := make(chan *broker.Message, execInputChannelBuffer) tstdin, putStdin := io.Pipe() stdin := io.NopCloser(tstdin) getStdout, stdout := io.Pipe() @@ -279,10 +288,22 @@ func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.Lis } select { - case msg := <-subCh: - if msg.ObjectType == broker.ExecInputObject { - if _, err := io.CopyBuffer(putStdin, strings.NewReader(msg.Object.(string)+"\n"), nil); err != nil { - h.Log.Error(ErrExecTerminal(err)) + case msg, ok := <-subCh: + if !ok { + // A broker implementation that closes the delivery channel on + // Unsubscribe would make this receive return (nil, false); end the + // session rather than spin on the closed channel or dereference a + // nil message. + h.Log.Debugf("Input channel closed for session %s", id) + return + } + // Guard the payload assertion too: a malformed ExecInput message with a + // non-string object must not panic the session loop. + if msg != nil && msg.ObjectType == broker.ExecInputObject { + if input, isStr := msg.Object.(string); isStr { + if _, err := io.CopyBuffer(putStdin, strings.NewReader(input+"\n"), nil); err != nil { + h.Log.Error(ErrExecTerminal(err)) + } } } case <-sessionCh: diff --git a/meshsync/exec_test.go b/meshsync/exec_test.go index 492b115b..53b9f8d3 100644 --- a/meshsync/exec_test.go +++ b/meshsync/exec_test.go @@ -2,7 +2,6 @@ package meshsync import ( "context" - "runtime" "strings" "sync" "testing" @@ -11,6 +10,7 @@ import ( "github.com/meshery/meshkit/broker" "github.com/meshery/meshkit/broker/channel" "github.com/meshery/meshkit/logger" + "go.uber.org/goleak" ) // recordingBroker wraps a real broker.Handler, recording the subjects passed to @@ -63,12 +63,17 @@ func TestUnsubscribeSessionInputTearsDownSubscription(t *testing.T) { rec := &recordingBroker{Handler: channel.NewChannelBrokerHandler()} h := &Handler{Broker: rec, Log: newTestLogger(t)} + // Snapshot the goroutines already running (broker/logger infra) now; the + // forwarding goroutine SubscribeWithChannel starts below is NOT in this set, + // so goleak fails if it is still alive after teardown. goleak retries with + // backoff, so it is robust to normal runtime scheduling (unlike a raw + // NumGoroutine comparison). + defer goleak.VerifyNone(t, goleak.IgnoreCurrent()) + const id = "exec.ns.pod.ctr.req-1" subject := "input." + id - // Mirror streamSession: a 1-buffered channel subscribed to input.. - subCh := make(chan *broker.Message, 1) - - baseline := runtime.NumGoroutine() + // Mirror streamSession: a buffered channel subscribed to input.. + subCh := make(chan *broker.Message, execInputChannelBuffer) if err := rec.SubscribeWithChannel(subject, generateID(), subCh); err != nil { t.Fatalf("SubscribeWithChannel: %v", err) @@ -102,9 +107,8 @@ func TestUnsubscribeSessionInputTearsDownSubscription(t *testing.T) { t.Fatalf("subject %s still registered after Unsubscribe: %v", subject, rec.ConnectedEndpoints()) } } - - // ...and the delivery goroutine it started has exited, so nothing leaks. - waitForGoroutines(t, baseline, 3*time.Second) + // ...and the delivery goroutine it started has exited (verified by the + // deferred goleak check), so nothing leaks. } // TestUnsubscribeSessionInputLogsErrorWithoutPanic covers the teardown error @@ -124,22 +128,3 @@ func TestUnsubscribeSessionInputLogsErrorWithoutPanic(t *testing.T) { t.Fatalf("expected exactly one Unsubscribe call, got %v", got) } } - -// waitForGoroutines fails if the goroutine count has not returned to at most -// baseline within timeout. It guards against the per-session delivery-goroutine -// leak that the old drain-goroutine approach caused. -func waitForGoroutines(t *testing.T, baseline int, timeout time.Duration) { - t.Helper() - deadline := time.Now().Add(timeout) - for { - runtime.Gosched() - if runtime.NumGoroutine() <= baseline { - return - } - if time.Now().After(deadline) { - t.Fatalf("goroutine count did not return to baseline %d (now %d): delivery goroutine leaked", - baseline, runtime.NumGoroutine()) - } - time.Sleep(10 * time.Millisecond) - } -}