Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions docs/agent-instructions/architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.<id>` 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.<id>")`, 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.
Expand Down
6 changes: 4 additions & 2 deletions docs/design/fd4-jetstream-durable-delivery.md
Original file line number Diff line number Diff line change
@@ -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.<id>` 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.
Expand All @@ -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.
Expand Down
5 changes: 3 additions & 2 deletions go.mod
Original file line number Diff line number Diff line change
Expand Up @@ -7,11 +7,12 @@ 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
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
Expand Down Expand Up @@ -108,7 +109,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
Expand Down
8 changes: 4 additions & 4 deletions go.sum
Original file line number Diff line number Diff line change
Expand Up @@ -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=
Expand Down
79 changes: 60 additions & 19 deletions meshsync/exec.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -73,8 +81,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.<id> subject; the
// 10s streamChannelPool ticker republishes the active-session list.
execCleanup(h, id)
continue
}
Expand Down Expand Up @@ -139,16 +148,27 @@ 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)
// 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()

// 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.
Expand All @@ -157,6 +177,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()
Expand All @@ -168,19 +195,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.<id> 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) {}),
Expand Down Expand Up @@ -270,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:
Expand All @@ -289,6 +319,17 @@ func (h *Handler) streamSession(id string, req model.ExecRequest, cfg config.Lis
}
}

// unsubscribeSessionInput tears down an exec session's stdin subscription
// (input.<id>) 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 {
Expand Down
130 changes: 130 additions & 0 deletions meshsync/exec_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,130 @@
package meshsync

import (
"context"
"strings"
"sync"
"testing"
"time"

"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
// 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.<id> 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)}

// 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 buffered channel subscribed to input.<id>.
subCh := make(chan *broker.Message, execInputChannelBuffer)

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 (verified by the
// deferred goleak check), so nothing leaks.
}

// 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)
}
}
Loading