From 68d2c6c2a1ed6fda6659fe895b7d42bdf017040a Mon Sep 17 00:00:00 2001 From: Alex Kotliarskyi Date: Mon, 17 Aug 2026 16:42:42 -0700 Subject: [PATCH 1/3] Apply backpressure to live exec output subscribers --- internal/controller/exec_sessions.go | 19 +-------- internal/controller/exec_sessions_test.go | 50 +++++++++++++++++++++++ 2 files changed, 51 insertions(+), 18 deletions(-) diff --git a/internal/controller/exec_sessions.go b/internal/controller/exec_sessions.go index 8ef0f31..b467640 100644 --- a/internal/controller/exec_sessions.go +++ b/internal/controller/exec_sessions.go @@ -251,24 +251,7 @@ func (subscriber *execSessionSubscriber) enqueue(frame *execstream.Frame) bool { subscriber.sendMu.Lock() defer subscriber.sendMu.Unlock() - if subscriber.alreadySentLocked(frame) { - return true - } - - select { - case <-subscriber.closed: - return false - default: - } - - select { - case subscriber.frames <- subscriber.markSentLocked(frame): - return true - case <-subscriber.closed: - return false - default: - return false - } + return subscriber.sendLocked(frame) } func (subscriber *execSessionSubscriber) sendHistory(frames []*execstream.Frame) bool { diff --git a/internal/controller/exec_sessions_test.go b/internal/controller/exec_sessions_test.go index e569def..eeece3f 100644 --- a/internal/controller/exec_sessions_test.go +++ b/internal/controller/exec_sessions_test.go @@ -266,6 +266,56 @@ func TestExecSessionHistoryReplayStreamsPastSubscriberBuffer(t *testing.T) { }, time.Second, 10*time.Millisecond) } +func TestExecSessionLiveOutputAppliesBackpressure(t *testing.T) { + session := newManualExecSessionForTest(execSessionKey{vmName: "vm", sessionID: "session"}, nil) + session.policy = legacyExecSessionPolicy + t.Cleanup(session.close) + + subscriber, err := session.attach() + require.NoError(t, err) + + const frameCount = 256 + done := make(chan struct{}) + go func() { + defer close(done) + for i := 0; i < frameCount; i++ { + session.recordFrame(&execstream.Frame{ + Type: execstream.FrameTypeStdout, + Data: []byte{byte(i)}, + }) + } + session.recordFrame(&execstream.Frame{ + Type: execstream.FrameTypeExit, + Exit: &execstream.Exit{Code: 0}, + }) + }() + + require.Eventually(t, func() bool { + return len(subscriber.frames) == cap(subscriber.frames) + }, time.Second, time.Millisecond) + + for i := 0; i < frameCount; i++ { + frame, ok := <-subscriber.frames + require.True(t, ok, "subscriber closed before output frame %d", i) + require.Equal(t, execstream.FrameTypeStdout, frame.Type) + require.Equal(t, []byte{byte(i)}, frame.Data) + } + + exitFrame, ok := <-subscriber.frames + require.True(t, ok, "subscriber closed before the exit frame") + require.Equal(t, execstream.FrameTypeExit, exitFrame.Type) + require.EqualValues(t, 0, exitFrame.Exit.Code) + + require.Eventually(t, func() bool { + select { + case <-done: + return true + default: + return false + } + }, time.Second, time.Millisecond) +} + func TestExecSessionDetachKeepsProcessAlive(t *testing.T) { registry := newExecSessionRegistry() session := newManualExecSessionForTest(execSessionKey{vmName: "vm", sessionID: "session"}, registry) From 2818b340a431306621a3aab34a1c80d3c15ad882 Mon Sep 17 00:00:00 2001 From: Alex Kotliarskyi Date: Mon, 17 Aug 2026 17:09:02 -0700 Subject: [PATCH 2/3] Use integer range loops in exec backpressure test --- internal/controller/exec_sessions_test.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/internal/controller/exec_sessions_test.go b/internal/controller/exec_sessions_test.go index eeece3f..73cb696 100644 --- a/internal/controller/exec_sessions_test.go +++ b/internal/controller/exec_sessions_test.go @@ -278,7 +278,7 @@ func TestExecSessionLiveOutputAppliesBackpressure(t *testing.T) { done := make(chan struct{}) go func() { defer close(done) - for i := 0; i < frameCount; i++ { + for i := range frameCount { session.recordFrame(&execstream.Frame{ Type: execstream.FrameTypeStdout, Data: []byte{byte(i)}, @@ -294,7 +294,7 @@ func TestExecSessionLiveOutputAppliesBackpressure(t *testing.T) { return len(subscriber.frames) == cap(subscriber.frames) }, time.Second, time.Millisecond) - for i := 0; i < frameCount; i++ { + for i := range frameCount { frame, ok := <-subscriber.frames require.True(t, ok, "subscriber closed before output frame %d", i) require.Equal(t, execstream.FrameTypeStdout, frame.Type) From 7ac2673059bc24828f775114e2b717ae90dcd09f Mon Sep 17 00:00:00 2001 From: Yibo Zhuang Date: Tue, 25 Aug 2026 14:47:43 -0700 Subject: [PATCH 3/3] Preserve reconnectable exec subscriber isolation --- internal/controller/exec_sessions.go | 37 ++++++++-- internal/controller/exec_sessions_test.go | 87 +++++++++++++++++++++-- 2 files changed, 113 insertions(+), 11 deletions(-) diff --git a/internal/controller/exec_sessions.go b/internal/controller/exec_sessions.go index b467640..e887fbb 100644 --- a/internal/controller/exec_sessions.go +++ b/internal/controller/exec_sessions.go @@ -15,14 +15,16 @@ import ( const execSessionReplayBufferBytes = 4 * 1024 * 1024 type execSessionPolicy struct { - closeOnDetach bool - retainAfterExit bool - replayEnabled bool + closeOnDetach bool + retainAfterExit bool + replayEnabled bool + blockOnSubscriberBackpressure bool } var ( legacyExecSessionPolicy = execSessionPolicy{ - closeOnDetach: true, + closeOnDetach: true, + blockOnSubscriberBackpressure: true, } reconnectableExecSessionPolicy = execSessionPolicy{ retainAfterExit: true, @@ -247,11 +249,32 @@ func newExecSessionSubscriber() *execSessionSubscriber { } } -func (subscriber *execSessionSubscriber) enqueue(frame *execstream.Frame) bool { +func (subscriber *execSessionSubscriber) enqueue(frame *execstream.Frame, block bool) bool { subscriber.sendMu.Lock() defer subscriber.sendMu.Unlock() - return subscriber.sendLocked(frame) + if block { + return subscriber.sendLocked(frame) + } + + if subscriber.alreadySentLocked(frame) { + return true + } + + select { + case <-subscriber.closed: + return false + default: + } + + select { + case subscriber.frames <- subscriber.markSentLocked(frame): + return true + case <-subscriber.closed: + return false + default: + return false + } } func (subscriber *execSessionSubscriber) sendHistory(frames []*execstream.Frame) bool { @@ -612,7 +635,7 @@ func (session *execSession) recordFrame(frame *execstream.Frame) { session.mu.Unlock() for _, subscriber := range subscribers { - if !subscriber.enqueue(frame) { + if !subscriber.enqueue(frame, session.policy.blockOnSubscriberBackpressure) { session.dropSubscriber(subscriber) } } diff --git a/internal/controller/exec_sessions_test.go b/internal/controller/exec_sessions_test.go index 73cb696..e8bfca9 100644 --- a/internal/controller/exec_sessions_test.go +++ b/internal/controller/exec_sessions_test.go @@ -280,13 +280,21 @@ func TestExecSessionLiveOutputAppliesBackpressure(t *testing.T) { defer close(done) for i := range frameCount { session.recordFrame(&execstream.Frame{ - Type: execstream.FrameTypeStdout, - Data: []byte{byte(i)}, + Type: execstream.FrameTypeStdout, + Data: []byte{byte(i)}, + Terminal: nil, + Exit: nil, + Error: "", + Watermark: 0, }) } session.recordFrame(&execstream.Frame{ - Type: execstream.FrameTypeExit, - Exit: &execstream.Exit{Code: 0}, + Type: execstream.FrameTypeExit, + Data: nil, + Terminal: nil, + Exit: &execstream.Exit{Code: 0}, + Error: "", + Watermark: 0, }) }() @@ -316,6 +324,77 @@ func TestExecSessionLiveOutputAppliesBackpressure(t *testing.T) { }, time.Second, time.Millisecond) } +func TestReconnectableExecSessionDropsStalledSubscriber(t *testing.T) { + session := newManualExecSessionForTest(execSessionKey{vmName: "vm", sessionID: "session"}, nil) + t.Cleanup(session.close) + + stalledSubscriber, err := session.attach() + require.NoError(t, err) + + for i := range cap(stalledSubscriber.frames) { + session.recordFrame(&execstream.Frame{ + Type: execstream.FrameTypeStdout, + Data: []byte{byte(i)}, + Terminal: nil, + Exit: nil, + Error: "", + Watermark: 0, + }) + } + + healthySubscriber, err := session.attach() + require.NoError(t, err) + + done := make(chan struct{}) + go func() { + defer close(done) + session.recordFrame(&execstream.Frame{ + Type: execstream.FrameTypeStdout, + Data: []byte("still running"), + Terminal: nil, + Exit: nil, + Error: "", + Watermark: 0, + }) + session.recordFrame(&execstream.Frame{ + Type: execstream.FrameTypeExit, + Data: nil, + Terminal: nil, + Exit: &execstream.Exit{Code: 0}, + Error: "", + Watermark: 0, + }) + }() + + select { + case <-done: + case <-time.After(time.Second): + t.Fatal("a stalled reconnectable subscriber blocked live output") + } + + outputFrame := <-healthySubscriber.frames + require.Equal(t, execstream.FrameTypeStdout, outputFrame.Type) + require.Equal(t, []byte("still running"), outputFrame.Data) + + exitFrame := <-healthySubscriber.frames + require.Equal(t, execstream.FrameTypeExit, exitFrame.Type) + require.EqualValues(t, 0, exitFrame.Exit.Code) + + select { + case <-stalledSubscriber.closed: + default: + t.Fatal("the stalled reconnectable subscriber was not dropped") + } + + reconnectedSubscriber, err := session.attach() + require.NoError(t, err) + session.sendHistory(reconnectedSubscriber, uint64(cap(stalledSubscriber.frames))) + + require.Equal(t, execstream.FrameTypeStdout, (<-reconnectedSubscriber.frames).Type) + require.Equal(t, execstream.FrameTypeExit, (<-reconnectedSubscriber.frames).Type) + require.Equal(t, execstream.FrameTypeNoMoreHistory, (<-reconnectedSubscriber.frames).Type) +} + func TestExecSessionDetachKeepsProcessAlive(t *testing.T) { registry := newExecSessionRegistry() session := newManualExecSessionForTest(execSessionKey{vmName: "vm", sessionID: "session"}, registry)