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
18 changes: 12 additions & 6 deletions internal/controller/exec_sessions.go
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -247,10 +249,14 @@ 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()

if block {
return subscriber.sendLocked(frame)
}

if subscriber.alreadySentLocked(frame) {
return true
}
Expand Down Expand Up @@ -629,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)
}
}
Expand Down
129 changes: 129 additions & 0 deletions internal/controller/exec_sessions_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -266,6 +266,135 @@ 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 := range frameCount {
session.recordFrame(&execstream.Frame{
Type: execstream.FrameTypeStdout,
Data: []byte{byte(i)},
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,
})
}()

require.Eventually(t, func() bool {
return len(subscriber.frames) == cap(subscriber.frames)
}, time.Second, time.Millisecond)

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)
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 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)
Expand Down