Extract exec replay buffer

This commit is contained in:
Fedor Korotkov 2026-04-28 16:43:32 -04:00
parent 82f827cffb
commit 593d3569a9
2 changed files with 99 additions and 67 deletions

View File

@ -144,6 +144,70 @@ type execReplayFrame struct {
size int size int
} }
type execReplayBuffer struct {
frames []execReplayFrame
bufferBytes int
nextWatermark uint64
ackedWatermark uint64
}
func (buffer *execReplayBuffer) append(frame *execstream.Frame) *execstream.Frame {
frame = cloneExecFrame(frame)
buffer.nextWatermark++
frame.Watermark = buffer.nextWatermark
frameSize := execFrameSize(frame)
buffer.frames = append(buffer.frames, execReplayFrame{
frame: frame,
size: frameSize,
})
buffer.bufferBytes += frameSize
buffer.trimAcknowledged()
buffer.trimToLimit()
return frame
}
func (buffer *execReplayBuffer) ack(watermark uint64) {
if watermark <= buffer.ackedWatermark {
return
}
buffer.ackedWatermark = watermark
buffer.trimAcknowledged()
}
func (buffer *execReplayBuffer) replayAfter(
watermark uint64,
enqueue func(*execstream.Frame) bool,
) bool {
for _, record := range buffer.frames {
if record.frame.Watermark <= watermark {
continue
}
if !enqueue(record.frame) {
return false
}
}
return true
}
func (buffer *execReplayBuffer) trimAcknowledged() {
for len(buffer.frames) > 0 && buffer.frames[0].frame.Watermark <= buffer.ackedWatermark {
buffer.bufferBytes -= buffer.frames[0].size
buffer.frames = buffer.frames[1:]
}
}
func (buffer *execReplayBuffer) trimToLimit() {
for buffer.bufferBytes > execSessionReplayBufferBytes && len(buffer.frames) > 0 {
buffer.bufferBytes -= buffer.frames[0].size
buffer.frames = buffer.frames[1:]
}
}
type execSessionSubscriber struct { type execSessionSubscriber struct {
frames chan *execstream.Frame frames chan *execstream.Frame
} }
@ -179,10 +243,7 @@ type execSession struct {
stdin io.WriteCloser stdin io.WriteCloser
stdinClosed bool stdinClosed bool
subscribers map[*execSessionSubscriber]struct{} subscribers map[*execSessionSubscriber]struct{}
frames []execReplayFrame replay execReplayBuffer
bufferBytes int
nextWatermark uint64
ackedWatermark uint64
started bool started bool
finished bool finished bool
closed bool closed bool
@ -318,12 +379,7 @@ func (session *execSession) ack(watermark uint64) {
session.mu.Lock() session.mu.Lock()
defer session.mu.Unlock() defer session.mu.Unlock()
if watermark <= session.ackedWatermark { session.replay.ack(watermark)
return
}
session.ackedWatermark = watermark
session.trimAcknowledgedLocked()
} }
func (session *execSession) sendHistory( func (session *execSession) sendHistory(
@ -341,21 +397,15 @@ func (session *execSession) sendHistory(
return return
} }
for _, record := range session.frames { if !session.replay.replayAfter(watermark, subscriber.enqueue) {
if record.frame.Watermark <= watermark {
continue
}
if !subscriber.enqueue(record.frame) {
session.detachLocked(subscriber) session.detachLocked(subscriber)
return return
} }
}
if !subscriber.enqueue(&execstream.Frame{ if !subscriber.enqueue(&execstream.Frame{
Type: execstream.FrameTypeNoMoreHistory, Type: execstream.FrameTypeNoMoreHistory,
Watermark: session.nextWatermark, Watermark: session.replay.nextWatermark,
}) { }) {
session.detachLocked(subscriber) session.detachLocked(subscriber)
} }
@ -375,16 +425,10 @@ func (session *execSession) close() {
session.expiryTimer = nil session.expiryTimer = nil
} }
subscribers := make([]*execSessionSubscriber, 0, len(session.subscribers)) subscribers := session.takeSubscribersLocked()
for subscriber := range session.subscribers {
subscribers = append(subscribers, subscriber)
}
session.subscribers = map[*execSessionSubscriber]struct{}{}
session.mu.Unlock() session.mu.Unlock()
for _, subscriber := range subscribers { closeSubscribers(subscribers)
close(subscriber.frames)
}
session.cancel() session.cancel()
_ = session.exec.Close() _ = session.exec.Close()
@ -428,18 +472,10 @@ func (session *execSession) recordFrame(frame *execstream.Frame) {
return return
} }
frame = cloneExecFrame(frame)
if session.policy.replayEnabled { if session.policy.replayEnabled {
session.nextWatermark++ frame = session.replay.append(frame)
frame.Watermark = session.nextWatermark } else {
frame = cloneExecFrame(frame)
session.frames = append(session.frames, execReplayFrame{
frame: frame,
size: execFrameSize(frame),
})
session.bufferBytes += execFrameSize(frame)
session.trimAcknowledgedLocked()
session.trimToLimitLocked()
} }
for subscriber := range session.subscribers { for subscriber := range session.subscribers {
@ -465,16 +501,10 @@ func (session *execSession) markFinished() {
session.expiryTimer = time.AfterFunc(session.exitTTL, session.expire) session.expiryTimer = time.AfterFunc(session.exitTTL, session.expire)
} }
subscribers := make([]*execSessionSubscriber, 0, len(session.subscribers)) subscribers := session.takeSubscribersLocked()
for subscriber := range session.subscribers {
subscribers = append(subscribers, subscriber)
}
session.subscribers = map[*execSessionSubscriber]struct{}{}
session.mu.Unlock() session.mu.Unlock()
for _, subscriber := range subscribers { closeSubscribers(subscribers)
close(subscriber.frames)
}
session.doneOnce.Do(func() { session.doneOnce.Do(func() {
close(session.done) close(session.done)
@ -489,17 +519,19 @@ func (session *execSession) expire() {
session.close() session.close()
} }
func (session *execSession) trimAcknowledgedLocked() { func (session *execSession) takeSubscribersLocked() []*execSessionSubscriber {
for len(session.frames) > 0 && session.frames[0].frame.Watermark <= session.ackedWatermark { subscribers := make([]*execSessionSubscriber, 0, len(session.subscribers))
session.bufferBytes -= session.frames[0].size for subscriber := range session.subscribers {
session.frames = session.frames[1:] subscribers = append(subscribers, subscriber)
} }
session.subscribers = map[*execSessionSubscriber]struct{}{}
return subscribers
} }
func (session *execSession) trimToLimitLocked() { func closeSubscribers(subscribers []*execSessionSubscriber) {
for session.bufferBytes > execSessionReplayBufferBytes && len(session.frames) > 0 { for _, subscriber := range subscribers {
session.bufferBytes -= session.frames[0].size close(subscriber.frames)
session.frames = session.frames[1:]
} }
} }

View File

@ -153,8 +153,8 @@ func TestExecSessionHistoryReplayAndAck(t *testing.T) {
require.EqualValues(t, 3, noMoreHistory.Watermark) require.EqualValues(t, 3, noMoreHistory.Watermark)
session.ack(2) session.ack(2)
require.Len(t, session.frames, 1) require.Len(t, session.replay.frames, 1)
require.EqualValues(t, 3, session.frames[0].frame.Watermark) require.EqualValues(t, 3, session.replay.frames[0].frame.Watermark)
} }
func TestExecSessionDetachKeepsProcessAlive(t *testing.T) { func TestExecSessionDetachKeepsProcessAlive(t *testing.T) {
@ -193,8 +193,8 @@ func TestLegacyExecSessionDoesNotRetainReplayHistory(t *testing.T) {
session.recordFrame(&execstream.Frame{Type: execstream.FrameTypeStdout, Data: []byte("out")}) session.recordFrame(&execstream.Frame{Type: execstream.FrameTypeStdout, Data: []byte("out")})
require.Empty(t, session.frames) require.Empty(t, session.replay.frames)
require.Zero(t, session.nextWatermark) require.Zero(t, session.replay.nextWatermark)
} }
func TestExecSessionCloseIfUnusedClosesIdleSession(t *testing.T) { func TestExecSessionCloseIfUnusedClosesIdleSession(t *testing.T) {