Merge pull request #1677 from entireio/feat/review-nonblocking-tui-sink · Entire
Merge pull request #1677 from entireio/feat/review-nonblocking-tui-sink
372bc14→main·· peyton-alt·3d ago·2 files·+363 added/-31 removed
fix(review): TUI sink must never backpressure the orchestrator
Changes
2
cmd/entire/cli/review
Mtui_sink.go+130/-31
Mtui_sink_test.go+233
8 unmodified lines
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
23
24
25
26
43
44
45
46
47
48
49
3 unmodified lines
53
54
55
36
56
57
58
59
60
61
62
42
63
64
65
66
67
68
27 unmodified lines
96
97
98
99
100
101
102
103
104
105
77
78
106
107
108
109
110
111
112
36 unmodified lines
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
123
124
172
173
174
175
176
177
178
179
180
181
182
1 unmodified line
184
185
186
132
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
135
136
137
138
139
140
203
204
205
206
207
208
209
1 unmodified line
211
212
213
148
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
10 unmodified lines
258
259
260
165
261
262
263
264
5 unmodified lines
270
271
272
177
273
274
275
276
8 unmodified lines
285
286
287
192
288
289
290
291
6 unmodified lines
298
299
300
205
206
207
208
209
210
211
212
301
302
303
304
305
306
307
308
216
217
309
310
311
3 unmodified lines
315
316
317
318
319
320
321
322
323
324
325
326
8 unmodified lines
import (
"context"
"io"
"log/slog"
"sync"
"time"
tea "charm.land/bubbletea/v2"
"golang.org/x/term"
"github.com/entireio/cli/cmd/entire/cli/logging"
reviewtypes "github.com/entireio/cli/cmd/entire/cli/review/types"
)
// teaRunner is the slice of *tea.Program the sink depends on, extracted so
// tests can substitute a program with a deterministically stalled event loop.
type teaRunner interface {
Run() (tea.Model, error)
Send(msg tea.Msg)
Kill()
}
// tuiSinkQueueCap bounds the sink's internal dispatch queue. Program.Send is
// an unbuffered BLOCKING send: if the Bubble Tea Update/render pipeline ever
// stalls, a direct Send from the orchestrator's dispatch goroutine parks
// forever — freezing sink dispatch, the fanIn drain loop, the parsers, and
// reviewer-timeout handling with it (observed live: the 2026-07-07 run-6
// wedge, where the TUI froze mid-run and a 20m --timeout never surfaced).
// The queue absorbs bursts; overflow beyond the cap is dropped and counted —
// a display that can lag must never backpressure the data plane.
const tuiSinkQueueCap = 4096
// TUISink is a Sink that renders a Bubble Tea dashboard. The orchestrator
// calls AgentEvent/RunFinished from a single goroutine (CU4 serial-dispatch
// contract); the sink translates each event into a tea.Msg and sends it via
// Program.Send. Bubble Tea's Send is thread-safe, but we never rely on that
// property — the serial-dispatch promise means Send is only called from the
// orchestrator's dispatch goroutine.
// contract); the sink translates each event into a tea.Msg, enqueues it on a
// bounded internal queue, and a pump goroutine forwards it via Program.Send —
// so only the pump can ever block on a stalled Bubble Tea loop, never the
// orchestrator.
//
// Cancellation: cancel is the same context.CancelFunc that controls the
// orchestrator's run context. The first KeyCtrlC in the dashboard fires this
3 unmodified lines
// root's context, which cancels the same function — no parallel signal.Notify
// goroutine is needed here.
type TUISink struct {
program *tea.Program
program teaRunner
mu sync.Mutex
started bool
finished bool
dropped int
done chan struct{} // closed when the tea.Program exits
msgs chan tea.Msg // bounded dispatch queue drained by the pump
done chan struct{} // closed when the tea.Program exits
pumpDone chan struct{} // closed when the pump goroutine exits
}
// Compile-time interface check.
27 unmodified lines
tea.WithInput(input),
tea.WithoutSignalHandler(), // SIGINT handled by cobra root; KeyCtrlC calls cancel directly
)
return newTUISinkWithProgram(prog)
}
// newTUISinkWithProgram wires a TUISink around any teaRunner; tests inject
// fakes with stalled or recording Send implementations.
func newTUISinkWithProgram(prog teaRunner) *TUISink {
return &TUISink{
program: prog,
done: make(chan struct{}),
program: prog,
msgs: make(chan tea.Msg, tuiSinkQueueCap),
done: make(chan struct{}),
pumpDone: make(chan struct{}),
}
}
36 unmodified lines
_ = err
}
// Pump: the only goroutine allowed to block on Program.Send. When the
// program exits (done closes), a blocked Send unblocks via the program's
// context and the pump drains out. A Send that races program exit (done
// closes while a queued msg is in hand) is equally safe: Bubble Tea's
// Send is a context-guarded select and the msgs channel is never closed,
// so a post-exit Send is an immediate no-op — not a panic, not a block.
go func() {
defer close(s.pumpDone)
for {
select {
case <-s.done:
return
case msg := <-s.msgs:
s.program.Send(msg)
}
}
}()
}
// Wait blocks until the Bubble Tea program exits. Safe to call after Start.
// If Start was never called, Wait returns immediately.
// Wait blocks until the Bubble Tea program exits, with a bounded escalation
// so teardown can never hang: in the normal flow PostRunComplete has already
// quit the program and Wait returns immediately; otherwise (early-error
// return paths, or a wedged loop that survived Kill) Wait gives the program
// one grace period, Kills it, gives it one more, and then abandons the
// goroutine — a stuck display must not hold command exit hostage. Joins the
// pump goroutine whenever the program actually exited. Safe to call after
// Start; if Start was never called, returns immediately.
func (s *TUISink) Wait() {
s.mu.Lock()
started := s.started
1 unmodified line
if !started {
return
}
<-s.done
select {
case <-s.done:
<-s.pumpDone
return
case <-time.After(tuiPostRunCompleteGrace):
}
s.program.Kill()
select {
case <-s.done:
<-s.pumpDone
case <-time.After(tuiPostRunCompleteGrace):
// Bubble Tea never returned from Run despite Kill. Abandon the
// program and pump goroutines rather than hanging teardown.
}
}
// AgentEvent (Sink interface): translate ev into a tea.Msg and Send it to the
// Bubble Tea program. Implements the serial-dispatch contract: the orchestrator
// calls this from a single goroutine.
// Note: Send is safe to call from goroutines other than the TUI's update loop;
// Bubble Tea's implementation queues the message internally.
// AgentEvent (Sink interface): translate ev into a tea.Msg and enqueue it for the
// pump. NEVER blocks: display events beyond the queue cap are dropped and
// counted rather than backpressuring the orchestrator's dispatch goroutine —
// see tuiSinkQueueCap for the incident this guards against.
func (s *TUISink) AgentEvent(agent string, ev reviewtypes.Event) {
s.mu.Lock()
ok := s.started && !s.finished
1 unmodified line
if !ok {
return
}
s.program.Send(agentEventMsg{agent: agent, ev: ev})
select {
case s.msgs <- agentEventMsg{agent: agent, ev: ev}:
default:
s.mu.Lock()
s.dropped++
s.mu.Unlock()
}
}
// enqueueControl enqueues a rare, must-not-be-lost-lightly message (run
// summary, phase transitions, quit) with a bounded wait: worth briefly
// waiting out a transient jam, but a wedged TUI must not hold the run
// hostage — callers all have degradation paths (PostRunComplete falls back
// to Kill; a lost summary leaves the footer stale until quit).
func (s *TUISink) enqueueControl(msg tea.Msg) {
select {
case s.msgs <- msg:
case <-s.done:
case <-time.After(tuiPostRunCompleteGrace):
s.mu.Lock()
s.dropped++
s.mu.Unlock()
}
}
// droppedCount reports how many messages were discarded due to a jammed
// queue. Zero in any healthy run.
func (s *TUISink) droppedCount() int {
s.mu.Lock()
defer s.mu.Unlock()
return s.dropped
}
// RunFinished (Sink interface): mark reviewer execution complete and send the
10 unmodified lines
s.finished = true
s.mu.Unlock()
s.program.Send(runFinishedMsg{summary: summary})
s.enqueueControl(runFinishedMsg{summary: summary})
}
// FinalPhaseStarted updates the TUI with a visible post-run phase such as the
5 unmodified lines
if !ok {
return
}
s.program.Send(finalPhaseStartedMsg{name: name})
s.enqueueControl(finalPhaseStartedMsg{name: name})
}
// FinalPhaseFinished marks the visible post-run phase complete.
8 unmodified lines
if err != nil {
msg.err = err.Error()
}
s.program.Send(msg)
s.enqueueControl(msg)
}
// PostRunComplete exits the TUI and waits for the Bubble Tea program to finish.
6 unmodified lines
return
}
// Program.Send can block if Bubble Tea has not entered its event loop yet.
// Send from a goroutine and fall back to Kill so a lost post-run quit cannot
// leave the CLI stuck on "Finalizing output..." forever.
sent := make(chan struct{})
go func() {
s.program.Send(postRunCompleteMsg{})
close(sent)
}
// enqueueControl is bounded, so this cannot park forever even when the
// Bubble Tea loop is stalled or never entered; the Kill fallback below
// guarantees a lost post-run quit cannot leave the CLI stuck on
// "Finalizing output..." forever.
s.enqueueControl(postRunCompleteMsg{})
select {
case <-s.done:
return
case <-sent:
case <-time.After(tuiPostRunCompleteGrace):
s.program.Kill()
}
3 unmodified lines
case <-time.After(tuiPostRunCompleteGrace):
s.program.Kill()
}
// Surface silent loss: a healthy run never drops. A non-zero count means
// the TUI loop stalled or lagged badly enough to jam the queue — exactly
// the diagnostic a future wedge investigation needs first.
if n := s.droppedCount(); n > 0 {
logging.Debug(context.Background(), "tui sink dropped messages under backpressure",
slog.Int("dropped", n))
}
}
Mcmd/entire/cli/review/tui_sink.go+130/-31
1 unmodified line
2
3
4
5
6
7
8
225 unmodified lines
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
1 unmodified line
import (
"bytes"
"sync"
"testing"
"time"
225 unmodified lines
t.Errorf("invalid fd should yield zero dims, got width=%d height=%d", width, height)
}
}
// --- Non-blocking dispatch (wedge hardening) ---
// wedgedProgram is a teaRunner whose event loop never consumes messages:
// Send blocks until Kill, modeling a Bubble Tea program whose Update/render
// pipeline has stalled (the 2026-07-07 run-6 incident shape). Run blocks
// until Kill so the sink's done channel behaves like a live program's.
type wedgedProgram struct {
killed chan struct{}
}
func newWedgedProgram() *wedgedProgram {
return &wedgedProgram{killed: make(chan struct{})}
}
func (w *wedgedProgram) Run() (tea.Model, error) {
<-w.killed
return nil, nil //nolint:nilnil // mirrors tea.Program.Run's exit shape; callers ignore both values
}
func (w *wedgedProgram) Send(tea.Msg) { <-w.killed }
func (w *wedgedProgram) Kill() {
select {
case <-w.killed:
default:
close(w.killed)
}
}
// recordingProgram is a teaRunner that records every message it receives.
type recordingProgram struct {
killed chan struct{}
mu sync.Mutex
msgs []tea.Msg
}
func newRecordingProgram() *recordingProgram {
return &recordingProgram{killed: make(chan struct{})}
}
func (r *recordingProgram) Run() (tea.Model, error) {
<-r.killed
return nil, nil //nolint:nilnil // mirrors tea.Program.Run's exit shape; callers ignore both values
}
func (r *recordingProgram) Send(msg tea.Msg) {
r.mu.Lock()
r.msgs = append(r.msgs, msg)
r.mu.Unlock()
}
func (r *recordingProgram) Kill() {
select {
case <-r.killed:
default:
close(r.killed)
}
}
func (r *recordingProgram) recorded() []tea.Msg {
r.mu.Lock()
defer r.mu.Unlock()
return append([]tea.Msg(nil), r.msgs...)
}
// TestTUISink_AgentEventNeverBlocksWhenProgramLoopIsWedged pins the wedge
// hardening: a stalled Bubble Tea loop must never backpressure the
// orchestrator. Before the fix, the first AgentEvent after the stall parked
// forever inside Program.Send, freezing sink dispatch, the fanIn drain loop,
// the parsers, and reviewer-timeout handling with them.
func TestTUISink_AgentEventNeverBlocksWhenProgramLoopIsWedged(t *testing.T) {
t.Parallel()
prog := newWedgedProgram()
sink := newTUISinkWithProgram(prog)
sink.Start()
defer func() {
prog.Kill()
sink.Wait()
}()
finished := make(chan struct{})
go func() {
for range 3 * tuiSinkQueueCap {
sink.AgentEvent("agent-a", reviewtypes.AssistantText{Text: "x"})
}
close(finished)
}()
select {
case <-finished:
case <-time.After(5 * time.Second):
t.Fatal("AgentEvent blocked on a wedged TUI loop — orchestrator freeze")
}
if got := sink.droppedCount(); got == 0 {
t.Error("expected overflow drops to be counted when the queue jams")
}
}
// TestTUISink_EventsReachProgramInOrder pins that the async pump preserves
// dispatch order for a healthy program.
func TestTUISink_EventsReachProgramInOrder(t *testing.T) {
t.Parallel()
prog := newRecordingProgram()
sink := newTUISinkWithProgram(prog)
sink.Start()
defer func() {
prog.Kill()
sink.Wait()
}()
for i := range 50 {
sink.AgentEvent("agent-a", reviewtypes.AssistantText{Text: string(rune('a' + i%26))})
}
sink.RunFinished(reviewtypes.RunSummary{})
deadline := time.After(5 * time.Second)
for {
msgs := prog.recorded()
if len(msgs) >= 51 {
for i := range 50 {
if _, ok := msgs[i].(agentEventMsg); !ok {
t.Fatalf("msgs[%%d] = %%T, want agentEventMsg", i, msgs[i])
}
}
if _, ok := msgs[50].(runFinishedMsg); !ok {
t.Fatalf("msgs[50] = %%T, want runFinishedMsg (order violated)", msgs[50])
}
return
}
select {
case <-deadline:
t.Fatalf("only %%d/51 messages reached the program", len(msgs))
case <-time.After(10 * time.Millisecond):
}
}
}
// TestTUISink_RunFinishedBoundedWhenWedged pins that control messages use a
// bounded wait rather than blocking forever when the queue is jammed.
func TestTUISink_RunFinishedBoundedWhenWedged(t *testing.T) {
t.Parallel()
prog := newWedgedProgram()
sink := newTUISinkWithProgram(prog)
sink.Start()
defer func() {
prog.Kill()
sink.Wait()
}()
// Jam the queue.
for range 2 * tuiSinkQueueCap {
sink.AgentEvent("agent-a", reviewtypes.AssistantText{Text: "x"})
}
finished := make(chan struct{})
go func() {
sink.RunFinished(reviewtypes.RunSummary{})
close(finished)
}()
select {
case <-finished:
case <-time.After(tuiPostRunCompleteGrace + 3*time.Second):
t.Fatal("RunFinished blocked past its bounded wait on a wedged TUI")
}
}
// stubbornProgram is a teaRunner whose Run NEVER returns, even after Kill —
// modeling a Bubble Tea teardown stuck restoring a blocked terminal. Send
// unblocks on Kill so the pump can drain, but done never closes.
type stubbornProgram struct {
killed chan struct{}
block chan struct{}
}
func newStubbornProgram() *stubbornProgram {
return &stubbornProgram{killed: make(chan struct{}), block: make(chan struct{})}
}
func (p *stubbornProgram) Run() (tea.Model, error) {
<-p.block // never closed — Run never returns
return nil, nil //nolint:nilnil // unreachable; mirrors tea.Program.Run's shape
}
func (p *stubbornProgram) Send(tea.Msg) { <-p.killed }
func (p *stubbornProgram) Kill() {
select {
case <-p.killed:
default:
close(p.killed)
}
}
// TestTUISink_WaitIsBoundedWhenProgramNeverExits pins the teardown guarantee:
// `defer tuiSink.Wait()` must not hang the command forever when the Bubble
// Tea program never returns from Run, even after Kill. Wait escalates
// (grace → Kill → grace) and then abandons the goroutine.
func TestTUISink_WaitIsBoundedWhenProgramNeverExits(t *testing.T) {
t.Parallel()
prog := newStubbornProgram()
sink := newTUISinkWithProgram(prog)
sink.Start()
finished := make(chan struct{})
go func() {
sink.Wait()
close(finished)
}()
select {
case <-finished:
case <-time.After(2*tuiPostRunCompleteGrace + 3*time.Second):
t.Fatal("Wait hung on a program that never exits — teardown wedge")
}
}
// TestTUISink_WaitJoinsPump pins that a normal Wait joins the pump goroutine
// (no leak between done closing and the pump observing it).
func TestTUISink_WaitJoinsPump(t *testing.T) {
t.Parallel()
prog := newRecordingProgram()
sink := newTUISinkWithProgram(prog)
sink.Start()
prog.Kill()
sink.Wait()
select {
case <-sink.pumpDone:
case <-time.After(2 * time.Second):
t.Fatal("Wait returned before the pump goroutine exited")
}
}