返回 DeepSeek-Reasonix
subagent_progress_test.go
根目录 / internal / agent / subagent_progress_test.go
1 package agent
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "strings"
8 "sync"
9 "testing"
10 "time"
11 "unicode/utf8"
12
13 "reasonix/internal/event"
14 "reasonix/internal/jobs"
15 "reasonix/internal/provider"
16 "reasonix/internal/tool"
17 )
18
19 // fakeProgressClock drives the merger's pacing deterministically: tests advance
20 // the clock instead of sleeping, and the merger's timer fires exactly when the
21 // fake time passes its deadline.
22 type fakeProgressClock struct {
23 mu sync.Mutex
24 now time.Time
25 timers []*fakeProgressTimer
26 }
27
28 func newFakeProgressClock(t0 time.Time) *fakeProgressClock {
29 return &fakeProgressClock{now: t0}
30 }
31
32 func (f *fakeProgressClock) Now() time.Time {
33 f.mu.Lock()
34 defer f.mu.Unlock()
35 return f.now
36 }
37
38 func (f *fakeProgressClock) NewTimer(d time.Duration) progressTimer {
39 f.mu.Lock()
40 defer f.mu.Unlock()
41 t := &fakeProgressTimer{clock: f, ch: make(chan time.Time, 1), deadline: f.now.Add(d)}
42 f.timers = append(f.timers, t)
43 return t
44 }
45
46 // Advance moves the clock forward and fires every due, armed timer. Fires are
47 // delivered to the timer channel only when it is not already holding a value,
48 // so stale fires never block the test.
49 func (f *fakeProgressClock) Advance(d time.Duration) {
50 f.mu.Lock()
51 f.now = f.now.Add(d)
52 timers := make([]*fakeProgressTimer, len(f.timers))
53 copy(timers, f.timers)
54 f.mu.Unlock()
55 now := f.Now()
56 var due []*fakeProgressTimer
57 for _, t := range timers {
58 t.mu.Lock()
59 if !t.stopped && !t.deadline.IsZero() && !t.deadline.After(now) && !t.fired {
60 t.fired = true
61 due = append(due, t)
62 }
63 t.mu.Unlock()
64 }
65 for _, t := range due {
66 select {
67 case t.ch <- time.Time{}:
68 default:
69 }
70 }
71 }
72
73 type fakeProgressTimer struct {
74 clock *fakeProgressClock
75 ch chan time.Time
76 deadline time.Time
77 fired bool
78 stopped bool
79 mu sync.Mutex
80 }
81
82 func (t *fakeProgressTimer) C() <-chan time.Time { return t.ch }
83
84 func (t *fakeProgressTimer) Reset(d time.Duration) bool {
85 t.mu.Lock()
86 defer t.mu.Unlock()
87 t.deadline = t.clock.Now().Add(d)
88 t.fired = false
89 t.stopped = false
90 return true
91 }
92
93 func (t *fakeProgressTimer) Stop() bool {
94 t.mu.Lock()
95 defer t.mu.Unlock()
96 t.stopped = true
97 return true
98 }
99
100 // chanSink delivers emitted events to a channel so tests wait on the pipeline
101 // instead of sleeping on the real clock.
102 type chanSink struct {
103 ch chan event.Event
104 }
105
106 func (s chanSink) Emit(e event.Event) { s.ch <- e }
107
108 func waitEvent(t *testing.T, ch chan event.Event, desc string) event.Event {
109 t.Helper()
110 select {
111 case e := <-ch:
112 return e
113 case <-time.After(2 * time.Second):
114 t.Fatalf("timed out waiting for %s", desc)
115 return event.Event{}
116 }
117 }
118
119 // collectFor drains the sink until it stays quiet for quietFor, bounding the
120 // wait for asynchronous flusher emission without relying on real-time sleeps
121 // for correctness.
122 func collectFor(t *testing.T, ch chan event.Event, quietFor time.Duration) []event.Event {
123 t.Helper()
124 var out []event.Event
125 for {
126 select {
127 case e := <-ch:
128 out = append(out, e)
129 case <-time.After(quietFor):
130 return out
131 }
132 }
133 }
134
135 // newTestTracker joins a tracker to a shared merger so tests control the clock
136 // and the merger is guaranteed closed when the test ends.
137 func newTestTracker(t *testing.T, clock progressClock, sink event.Sink, childID string) *subagentProgressTracker {
138 t.Helper()
139 merger := newSubagentProgressMerger(clock, sink, "group-1")
140 t.Cleanup(merger.Close)
141 ctx := withSubagentProgressMerger(withCallContext(context.Background(), childID, sink, nil, false), merger)
142 return newSubagentProgressTracker(ctx, subSink(ctx))
143 }
144
145 func progressName(e event.Event) string { return e.Tool.Name }
146 func progressOutput(e event.Event) string { return e.Tool.Output }
147
148 func TestSubagentProgressStatusFirstSendImmediateThenMerges(t *testing.T) {
149 clock := newFakeProgressClock(time.Unix(0, 0))
150 ch := make(chan event.Event, 64)
151 trk := newTestTracker(t, clock, chanSink{ch: ch}, "child-1")
152 trk.running()
153
154 first := waitEvent(t, ch, "first status event")
155 if first.Tool.Name != event.SubagentProgressStatusName || first.Tool.Output != string(subagentPhaseRunning) {
156 t.Fatalf("first status event = %+v, want running status", first.Tool)
157 }
158 if first.Tool.ID != "child-1" {
159 t.Fatalf("status ID = %q, want child-1", first.Tool.ID)
160 }
161 if first.Tool.ParentID != "group-1" {
162 t.Fatalf("status ParentID = %q, want group-1", first.Tool.ParentID)
163 }
164
165 // A phase change inside the 250ms window merges into the slot: no second
166 // event is due until the window after the previous send.
167 trk.setPhase(subagentPhaseReasoning)
168 select {
169 case e := <-ch:
170 t.Fatalf("status merged too early: %+v", e.Tool)
171 case <-time.After(50 * time.Millisecond):
172 }
173
174 clock.Advance(subagentProgressMergeWindow)
175 merged := waitEvent(t, ch, "merged status event")
176 if merged.Tool.Output != string(subagentPhaseReasoning) {
177 t.Fatalf("merged status = %q, want reasoning", merged.Tool.Output)
178 }
179 }
180
181 func TestSubagentProgressPreviewMergesWithinWindow(t *testing.T) {
182 clock := newFakeProgressClock(time.Unix(0, 0))
183 ch := make(chan event.Event, 64)
184 trk := newTestTracker(t, clock, chanSink{ch: ch}, "child-1")
185 trk.running()
186 waitEvent(t, ch, "running status")
187 trk.wrap().Emit(event.Event{Kind: event.Reasoning, Text: "first "})
188 // Nothing is due before the 250ms window elapses.
189 select {
190 case e := <-ch:
191 t.Fatalf("preview sent before merge window: %+v", e.Tool)
192 case <-time.After(50 * time.Millisecond):
193 }
194
195 // Deltas arriving inside the window merge into one slot.
196 trk.wrap().Emit(event.Event{Kind: event.Reasoning, Text: "second"})
197 trk.wrap().Emit(event.Event{Kind: event.Reasoning, Text: " third"})
198 clock.Advance(subagentProgressMergeWindow)
199
200 // The reasoning phase transition and the preview are both due at the
201 // window; the status event is emitted first.
202 status := waitEvent(t, ch, "merged status")
203 if status.Tool.Name != event.SubagentProgressStatusName || status.Tool.Output != string(subagentPhaseReasoning) {
204 t.Fatalf("merged status = %+v, want reasoning phase", status.Tool)
205 }
206 merged := waitEvent(t, ch, "merged preview")
207 if merged.Tool.Name != event.SubagentProgressReasoningName {
208 t.Fatalf("preview name = %q, want %q", merged.Tool.Name, event.SubagentProgressReasoningName)
209 }
210 if merged.Tool.Output != "first second third" {
211 t.Fatalf("preview output = %q, want merged deltas", merged.Tool.Output)
212 }
213 if merged.Tool.Truncated {
214 t.Fatal("merged preview must not be marked truncated")
215 }
216 // The window is per (child, channel): a text delta is due on its own timer,
217 // with the responding phase transition emitted first.
218 trk.wrap().Emit(event.Event{Kind: event.Text, Text: "text"})
219 clock.Advance(subagentProgressMergeWindow)
220 resp := waitEvent(t, ch, "responding status")
221 if resp.Tool.Name != event.SubagentProgressStatusName || resp.Tool.Output != string(subagentPhaseResponding) {
222 t.Fatalf("responding status = %+v", resp.Tool)
223 }
224 text := waitEvent(t, ch, "text preview")
225 if text.Tool.Name != event.SubagentProgressTextName || text.Tool.Output != "text" {
226 t.Fatalf("text preview = %+v, want text channel", text.Tool)
227 }
228 }
229
230 func TestSubagentProgressWrapConvertsAndForwards(t *testing.T) {
231 clock := newFakeProgressClock(time.Unix(0, 0))
232 progressCh := make(chan event.Event, 256)
233 merger := newSubagentProgressMerger(clock, chanSink{ch: progressCh}, "group-1")
234 t.Cleanup(merger.Close)
235 parent := &recordSink{}
236 ctx := withSubagentProgressMerger(withCallContext(context.Background(), "task-1", parent, nil, false), merger)
237 trk := newSubagentProgressTracker(ctx, subSink(ctx))
238 trk.running()
239 wrap := trk.wrap()
240
241 // Child reasoning/text/notice/retrying become reserved progress channels.
242 wrap.Emit(event.Event{Kind: event.Reasoning, Text: "think a"})
243 wrap.Emit(event.Event{Kind: event.Reasoning, Text: " think b"})
244 wrap.Emit(event.Event{Kind: event.Text, Text: "answer"})
245 wrap.Emit(event.Event{Kind: event.Notice, Text: "heads up"})
246 wrap.Emit(event.Event{Kind: event.Notice, Detail: "detail only"})
247 wrap.Emit(event.Event{Kind: event.Retrying, RetryAttempt: 2, RetryMax: 3})
248
249 // Message and other parent-visible bodies must never be forwarded.
250 wrap.Emit(event.Event{Kind: event.Message, Text: "parent body", Reasoning: "parent reasoning"})
251 wrap.Emit(event.Event{Kind: event.TurnStarted})
252 wrap.Emit(event.Event{Kind: event.TurnDone})
253
254 // Real tool activity passes through to the parent, namespaced as before.
255 wrap.Emit(event.Event{Kind: event.ToolDispatch, Tool: event.Tool{ID: "bash_1", Name: "bash"}})
256 wrap.Emit(event.Event{Kind: event.ToolProgress, Tool: event.Tool{ID: "bash_1", Output: "chunk"}})
257 wrap.Emit(event.Event{Kind: event.Usage, ModelRef: "m"})
258
259 trk.finish(nil, nil)
260
261 got := map[string]string{}
262 var statuses []string
263 for _, e := range collectFor(t, progressCh, 100*time.Millisecond) {
264 if progressName(e) == event.SubagentProgressStatusName {
265 statuses = append(statuses, progressOutput(e))
266 } else {
267 got[progressName(e)] += progressOutput(e)
268 }
269 }
270 if len(statuses) != 3 || statuses[0] != string(subagentPhaseRunning) || statuses[1] != string(subagentPhaseTool) || statuses[2] != string(subagentPhaseCompleted) {
271 t.Fatalf("statuses = %v, want running → tool → completed", statuses)
272 }
273 if got[event.SubagentProgressReasoningName] != "think a think b" {
274 t.Fatalf("reasoning preview = %q, want both deltas merged", got[event.SubagentProgressReasoningName])
275 }
276 if got[event.SubagentProgressTextName] != "answer" {
277 t.Fatalf("text preview = %q", got[event.SubagentProgressTextName])
278 }
279 notice := got[event.SubagentProgressNoticeName]
280 if !strings.Contains(notice, "heads up") || !strings.Contains(notice, "detail only") {
281 t.Fatalf("notice preview = %q, want both notice texts", notice)
282 }
283
284 // Tool events: forwarded namespaced; no reasoning/text/message leakage.
285 var forwardNames, forwardIDs []string
286 for _, e := range parent.kinds(event.ToolDispatch) {
287 forwardNames = append(forwardNames, e.Tool.Name)
288 forwardIDs = append(forwardIDs, e.Tool.ID)
289 }
290 if len(forwardNames) != 1 || forwardNames[0] != "bash" || forwardIDs[0] != "task-1/bash_1" {
291 t.Fatalf("forwarded dispatch = %v %v, want namespaced bash", forwardNames, forwardIDs)
292 }
293 if tp := parent.kinds(event.ToolProgress); len(tp) != 1 || tp[0].Tool.ID != "task-1/bash_1" {
294 t.Fatalf("forwarded tool progress = %+v, want namespaced chunk", tp)
295 }
296 for _, kind := range []event.Kind{event.Reasoning, event.Text, event.Message, event.Notice, event.Retrying, event.TurnStarted, event.TurnDone} {
297 if n := len(parent.kinds(kind)); n != 0 {
298 t.Fatalf("parent received %d %v events; sub-agent bodies must not be forwarded", n, kind)
299 }
300 }
301 if len(parent.kinds(event.Usage)) != 1 {
302 t.Fatal("usage must still be forwarded")
303 }
304 }
305
306 func TestSubagentProgressTwoChildrenDoNotInterleave(t *testing.T) {
307 clock := newFakeProgressClock(time.Unix(0, 0))
308 ch := make(chan event.Event, 64)
309 merger := newSubagentProgressMerger(clock, chanSink{ch: ch}, "group-1")
310 t.Cleanup(merger.Close)
311
312 merger.statusEvent("a", subagentPhaseRunning)
313 merger.statusEvent("b", subagentPhaseRunning)
314 merger.deltaEvent("a", subagentProgressChanReasoning, "AAA")
315 merger.deltaEvent("b", subagentProgressChanReasoning, "BBB")
316
317 clock.Advance(subagentProgressMergeWindow)
318 got := collectFor(t, ch, 100*time.Millisecond)
319 // Each child's preview carries only its own content, keyed by its own ID.
320 var aGot, bGot []string
321 for _, e := range got {
322 switch {
323 case e.Tool.ID == "a" && progressName(e) == event.SubagentProgressReasoningName:
324 aGot = append(aGot, progressOutput(e))
325 case e.Tool.ID == "b" && progressName(e) == event.SubagentProgressReasoningName:
326 bGot = append(bGot, progressOutput(e))
327 case progressName(e) == event.SubagentProgressReasoningName:
328 t.Fatalf("preview for unknown child: %+v", e.Tool)
329 }
330 }
331 if strings.Join(aGot, "") != "AAA" || strings.Join(bGot, "") != "BBB" {
332 t.Fatalf("children interleaved: a=%v b=%v", aGot, bGot)
333 }
334 }
335
336 func TestSubagentProgressTerminalExactlyOnce(t *testing.T) {
337 cases := []struct {
338 name string
339 ctxErr error
340 runErr error
341 want string
342 }{
343 {"completed", nil, nil, string(subagentPhaseCompleted)},
344 {"cancelled", context.Canceled, errors.New("stop"), string(subagentPhaseCancelled)},
345 {"deadline", context.DeadlineExceeded, nil, string(subagentPhaseCancelled)},
346 {"failed", nil, errors.New("provider exploded"), string(subagentPhaseFailed)},
347 }
348 for _, tc := range cases {
349 t.Run(tc.name, func(t *testing.T) {
350 clock := newFakeProgressClock(time.Unix(0, 0))
351 ch := make(chan event.Event, 64)
352 trk := newTestTracker(t, clock, chanSink{ch: ch}, "task-1")
353 trk.running()
354 trk.finish(tc.ctxErr, tc.runErr)
355 trk.finish(nil, errors.New("second finish must be ignored"))
356
357 var terminals int
358 for _, e := range collectFor(t, ch, 100*time.Millisecond) {
359 if progressName(e) != event.SubagentProgressStatusName {
360 continue
361 }
362 if progressOutput(e) == string(subagentPhaseCompleted) || progressOutput(e) == string(subagentPhaseCancelled) || progressOutput(e) == string(subagentPhaseFailed) {
363 terminals++
364 if progressOutput(e) != tc.want {
365 t.Fatalf("terminal = %q, want %q", progressOutput(e), tc.want)
366 }
367 if e.Tool.DurationMs < 0 {
368 t.Fatalf("terminal DurationMs = %d, want >= 0", e.Tool.DurationMs)
369 }
370 }
371 }
372 if terminals != 1 {
373 t.Fatalf("terminal statuses = %d, want exactly one", terminals)
374 }
375 })
376 }
377 }
378
379 func TestSubagentProgressFlushPrecedesTerminal(t *testing.T) {
380 clock := newFakeProgressClock(time.Unix(0, 0))
381 ch := make(chan event.Event, 64)
382 trk := newTestTracker(t, clock, chanSink{ch: ch}, "task-1")
383 trk.running()
384 trk.wrap().Emit(event.Event{Kind: event.Reasoning, Text: "pending think"})
385 trk.wrap().Emit(event.Event{Kind: event.Text, Text: "pending answer"})
386
387 trk.finish(nil, nil)
388
389 var order, outputs []string
390 for _, e := range collectFor(t, ch, 100*time.Millisecond) {
391 order = append(order, progressName(e))
392 outputs = append(outputs, progressOutput(e))
393 }
394 // Pending phase + previews flush before the terminal status event: running
395 // (direct), the merged responding phase (reasoning→responding overwrote
396 // the slot), the reasoning preview, the text preview, then completed.
397 wantNames := []string{
398 event.SubagentProgressStatusName,
399 event.SubagentProgressStatusName,
400 event.SubagentProgressReasoningName,
401 event.SubagentProgressTextName,
402 event.SubagentProgressStatusName,
403 }
404 wantOutputs := []string{string(subagentPhaseRunning), string(subagentPhaseResponding), "", "", string(subagentPhaseCompleted)}
405 if len(order) != len(wantNames) {
406 t.Fatalf("event order = %v, want %v", order, wantNames)
407 }
408 for i := range wantNames {
409 if order[i] != wantNames[i] {
410 t.Fatalf("event %d name = %s, want %s", i, order[i], wantNames[i])
411 }
412 if order[i] == event.SubagentProgressStatusName && outputs[i] != wantOutputs[i] {
413 t.Fatalf("event %d status = %q, want %q", i, outputs[i], wantOutputs[i])
414 }
415 }
416 if outputs[2] != "pending think" || outputs[3] != "pending answer" {
417 t.Fatalf("flushed previews = %q %q, want pending think / pending answer", outputs[2], outputs[3])
418 }
419 }
420
421 func TestSubagentProgressLateEventsIgnoredAfterTerminal(t *testing.T) {
422 clock := newFakeProgressClock(time.Unix(0, 0))
423 ch := make(chan event.Event, 64)
424 trk := newTestTracker(t, clock, chanSink{ch: ch}, "task-1")
425 trk.running()
426 trk.finish(nil, nil)
427 // Consume the legitimate pre-terminal activity: running + completed.
428 waitEvent(t, ch, "running status")
429 waitEvent(t, ch, "completed terminal")
430
431 wrap := trk.wrap()
432 wrap.Emit(event.Event{Kind: event.Text, Text: "late"})
433 trk.setPhase(subagentPhaseRunning)
434 trk.finish(nil, errors.New("late finish"))
435
436 if got := collectFor(t, ch, 100*time.Millisecond); len(got) != 0 {
437 t.Fatalf("late events emitted after terminal: %+v", got)
438 }
439 }
440
441 func TestSubagentProgressUtf8TailKeepsRuneBoundaries(t *testing.T) {
442 delta := strings.Repeat("世", 4096) // 12 KiB of multi-byte pending text
443 clock := newFakeProgressClock(time.Unix(0, 0))
444 ch := make(chan event.Event, 64)
445 trk := newTestTracker(t, clock, chanSink{ch: ch}, "task-1")
446 trk.running()
447 wrap := trk.wrap()
448 wrap.Emit(event.Event{Kind: event.Reasoning, Text: delta})
449 wrap.Emit(event.Event{Kind: event.Text, Text: delta})
450 wrap.Emit(event.Event{Kind: event.Notice, Text: delta})
451
452 trk.finish(nil, nil)
453
454 pending := 0
455 for _, e := range collectFor(t, ch, 100*time.Millisecond) {
456 if e.Tool.Name == event.SubagentProgressStatusName {
457 continue
458 }
459 pending += len(e.Tool.Output)
460 if !utf8.ValidString(e.Tool.Output) {
461 t.Fatalf("preview split a multi-byte rune: %q", e.Tool.Output)
462 }
463 if !e.Tool.Truncated {
464 t.Fatalf("overflowing preview %q must set Truncated", e.Tool.Name)
465 }
466 }
467 if pending > subagentProgressMaxPendingBytes {
468 t.Fatalf("flushed pending = %d bytes, want <= %d", pending, subagentProgressMaxPendingBytes)
469 }
470 }
471
472 func TestSubagentProgressGroupBudgetBoundsBurstAndServesAll(t *testing.T) {
473 clock := newFakeProgressClock(time.Unix(0, 0))
474 ch := make(chan event.Event, 512)
475 merger := newSubagentProgressMerger(clock, chanSink{ch: ch}, "group-1")
476 t.Cleanup(merger.Close)
477
478 const n = 64
479 for i := 0; i < n; i++ {
480 child := "child-" + string(rune('0'+i/10)) + string(rune('0'+i%10))
481 // Ordinary phase transitions share the budget with previews: 64
482 // status changes alone must not exceed the 32 events/s contract.
483 merger.statusEvent(child, subagentPhaseRunning)
484 merger.deltaEvent(child, subagentProgressChanReasoning, strings.Repeat("x", 256))
485 }
486
487 // The first wave is capped by the group budget (32 events/s) across
488 // statuses and previews together.
489 clock.Advance(subagentProgressMergeWindow)
490 first := collectFor(t, ch, 100*time.Millisecond)
491 firstNonTerminal := 0
492 for _, e := range first {
493 if progressName(e) == event.SubagentProgressStatusName || progressName(e) == event.SubagentProgressReasoningName {
494 firstNonTerminal++
495 }
496 }
497 if firstNonTerminal > subagentProgressGroupBurst+8 { // burst + one refill at the wake
498 t.Fatalf("first-wave non-terminal events = %d, want capped by the 32/sec group budget", firstNonTerminal)
499 }
500
501 // Once the budget refills, every child is served exactly once — no child
502 // starves behind a high-activity sibling.
503 var rest []event.Event
504 for i := 0; i < 16; i++ {
505 clock.Advance(time.Second)
506 rest = append(rest, collectFor(t, ch, 50*time.Millisecond)...)
507 }
508 statuses := 0
509 served := map[string]int{}
510 for _, batch := range [][]event.Event{first, rest} {
511 for _, e := range batch {
512 switch {
513 case progressName(e) == event.SubagentProgressStatusName:
514 statuses++
515 case progressName(e) == event.SubagentProgressReasoningName:
516 served[e.Tool.ID]++
517 }
518 }
519 }
520 if statuses != n {
521 t.Fatalf("status events = %d, want all %d", statuses, n)
522 }
523 if len(served) != n {
524 t.Fatalf("served %d children, want all %d", len(served), n)
525 }
526 for id, count := range served {
527 if count != 1 {
528 t.Fatalf("child %s served %d times, want exactly once", id, count)
529 }
530 }
531 }
532
533 // TestSubagentProgressTrimTruncationPropagates proves a budget trim that drops
534 // buffered content marks the loss on the next actually-emitted channel, so
535 // frontends always learn that some preview content was discarded.
536 func TestSubagentProgressTrimTruncationPropagates(t *testing.T) {
537 clock := newFakeProgressClock(time.Unix(0, 0))
538 ch := make(chan event.Event, 64)
539 merger := newSubagentProgressMerger(clock, chanSink{ch: ch}, "group-1")
540 t.Cleanup(merger.Close)
541
542 merger.statusEvent("child-1", subagentPhaseRunning)
543 waitEvent(t, ch, "running status")
544 delta := strings.Repeat("世", 4096) // 12 KiB per channel; the shared 8 KiB budget trims
545 merger.deltaEvent("child-1", subagentProgressChanReasoning, delta)
546 merger.deltaEvent("child-1", subagentProgressChanText, delta)
547 merger.deltaEvent("child-1", subagentProgressChanNotice, delta)
548 merger.flushChild("child-1", subagentPhaseCompleted, 5)
549
550 var textEvent *event.Event
551 got := collectFor(t, ch, 100*time.Millisecond)
552 for i := range got {
553 if progressName(got[i]) == event.SubagentProgressTextName {
554 textEvent = &got[i]
555 }
556 }
557 if textEvent == nil {
558 t.Fatalf("no text preview emitted: %+v", got)
559 }
560 if !textEvent.Tool.Truncated {
561 t.Fatalf("budget-trimmed preview must carry Truncated: %+v", textEvent.Tool)
562 }
563 if textEvent.Tool.Output == "" || !utf8.ValidString(textEvent.Tool.Output) {
564 t.Fatalf("trimmed preview must keep a UTF-8-safe tail: %+v", textEvent.Tool)
565 }
566 }
567
568 func TestSubagentProgressMergerCloseIdempotentAndQuiet(t *testing.T) {
569 clock := newFakeProgressClock(time.Unix(0, 0))
570 ch := make(chan event.Event, 64)
571 merger := newSubagentProgressMerger(clock, chanSink{ch: ch}, "group-1")
572 merger.statusEvent("child-1", subagentPhaseRunning)
573 waitEvent(t, ch, "status before close")
574
575 merger.Close()
576 merger.Close() // idempotent
577 // Events after close are dropped, never panic.
578 merger.statusEvent("child-1", subagentPhaseReasoning)
579 merger.deltaEvent("child-1", subagentProgressChanReasoning, "dropped")
580 merger.flushChild("child-1", subagentPhaseCompleted, 12)
581 if got := collectFor(t, ch, 50*time.Millisecond); len(got) != 0 {
582 t.Fatalf("events after Close = %+v, want none", got)
583 }
584 clock.Advance(time.Second) // must not panic or deadlock
585 }
586
587 // reasoningTextProvider scripts one reasoning + text turn, so the integration
588 // tests can assert exactly what the progress pipeline forwards.
589 type reasoningTextProvider struct{}
590
591 func (reasoningTextProvider) Name() string { return "reasoning-text" }
592
593 func (reasoningTextProvider) Stream(context.Context, provider.Request) (<-chan provider.Chunk, error) {
594 ch := make(chan provider.Chunk, 3)
595 ch <- provider.Chunk{Type: provider.ChunkReasoning, Text: "thinking hard"}
596 ch <- provider.Chunk{Type: provider.ChunkText, Text: "final answer"}
597 ch <- provider.Chunk{Type: provider.ChunkDone}
598 close(ch)
599 return ch, nil
600 }
601
602 // streamErrorProvider fails the stream with a fixed error.
603 type streamErrorProvider struct{ err error }
604
605 func (p *streamErrorProvider) Name() string { return "stream-error" }
606
607 func (p *streamErrorProvider) Stream(context.Context, provider.Request) (<-chan provider.Chunk, error) {
608 return nil, p.err
609 }
610
611 func TestRunProfileSpecEmitsSubagentProgress(t *testing.T) {
612 rec := &recordSink{}
613 ctx := withCallContext(context.Background(), "task-1", rec, nil, false)
614 task := newTestTaskTool(t, reasoningTextProvider{}, tool.NewRegistry(), "sys", "", "", nil)
615
616 out, err := task.RunProfileSpec(ctx, ProfileExecSpec{
617 Kind: "task", Name: "task", Prompt: "do the thing",
618 SystemPrompt: "sys", AllowNoTools: true,
619 })
620 if err != nil {
621 t.Fatalf("RunProfileSpec: %v", err)
622 }
623 if !strings.Contains(out, "final answer") {
624 t.Fatalf("result = %q, want the child's final answer", out)
625 }
626
627 var order []string
628 for _, e := range rec.kinds(event.ToolProgress) {
629 order = append(order, progressName(e)+":"+progressOutput(e))
630 }
631 want := []string{
632 event.SubagentProgressStatusName + ":running",
633 // The child's reasoning→responding transition merges into the status
634 // slot and is flushed right before the previews.
635 event.SubagentProgressStatusName + ":responding",
636 event.SubagentProgressReasoningName + ":thinking hard",
637 event.SubagentProgressTextName + ":final answer",
638 event.SubagentProgressStatusName + ":completed",
639 }
640 if len(order) != len(want) {
641 t.Fatalf("progress events = %v, want %v", order, want)
642 }
643 for i := range want {
644 if order[i] != want[i] {
645 t.Fatalf("progress events = %v, want %v", order, want)
646 }
647 }
648 for _, e := range rec.kinds(event.ToolProgress) {
649 if e.Tool.ID != "task-1" {
650 t.Fatalf("progress ID = %q, want task-1", e.Tool.ID)
651 }
652 if e.Tool.ParentID != "" {
653 t.Fatalf("single-task progress ParentID = %q, want empty", e.Tool.ParentID)
654 }
655 }
656 // Child bodies never leak into the parent stream.
657 for _, kind := range []event.Kind{event.Reasoning, event.Text, event.Message, event.Notice, event.Retrying, event.TurnStarted, event.TurnDone} {
658 if n := len(rec.kinds(kind)); n != 0 {
659 t.Fatalf("parent received %d %v events from a sub-agent run", n, kind)
660 }
661 }
662 }
663
664 func TestRunProfileSpecProgressCancelledTerminal(t *testing.T) {
665 rec := &recordSink{}
666 ctx, cancel := context.WithCancel(withCallContext(context.Background(), "task-1", rec, nil, false))
667 cancel()
668 task := newTestTaskTool(t, reasoningTextProvider{}, tool.NewRegistry(), "sys", "", "", nil)
669
670 if _, err := task.RunProfileSpec(ctx, ProfileExecSpec{
671 Kind: "task", Name: "task", Prompt: "do the thing",
672 SystemPrompt: "sys", AllowNoTools: true,
673 }); err == nil {
674 t.Fatal("cancelled RunProfileSpec must return an error")
675 }
676
677 var terminals []string
678 for _, e := range rec.kinds(event.ToolProgress) {
679 if progressName(e) == event.SubagentProgressStatusName && progressOutput(e) == string(subagentPhaseCancelled) {
680 terminals = append(terminals, progressOutput(e))
681 }
682 }
683 if len(terminals) != 1 {
684 t.Fatalf("cancelled terminal count = %d, want exactly one", len(terminals))
685 }
686 if len(rec.kinds(event.ToolProgress)) < 2 {
687 t.Fatalf("progress events = %d, want running + cancelled at minimum", len(rec.kinds(event.ToolProgress)))
688 }
689 }
690
691 func TestRunProfileSpecProgressFailedOnProviderError(t *testing.T) {
692 rec := &recordSink{}
693 ctx := withCallContext(context.Background(), "task-1", rec, nil, false)
694 task := newTestTaskTool(t, &streamErrorProvider{err: errors.New("transport cut")}, tool.NewRegistry(), "sys", "", "", nil)
695
696 if _, err := task.RunProfileSpec(ctx, ProfileExecSpec{
697 Kind: "task", Name: "task", Prompt: "do the thing",
698 SystemPrompt: "sys", AllowNoTools: true,
699 }); err == nil {
700 t.Fatal("provider error must propagate")
701 }
702
703 terminals := 0
704 for _, e := range rec.kinds(event.ToolProgress) {
705 if progressName(e) == event.SubagentProgressStatusName && progressOutput(e) == string(subagentPhaseFailed) {
706 terminals++
707 }
708 }
709 if terminals != 1 {
710 t.Fatalf("failed terminal count = %d, want exactly one", terminals)
711 }
712 }
713
714 func TestRunProfileSpecProgressPanicEmitsFailed(t *testing.T) {
715 rec := &recordSink{}
716 ctx := withCallContext(context.Background(), "task-1", rec, nil, false)
717 task := newTestTaskTool(t, panicProvider{name: "boom"}, tool.NewRegistry(), "sys", "", "", nil)
718
719 panicked := false
720 func() {
721 defer func() {
722 if recover() == nil {
723 t.Error("panic must propagate after the failed terminal is emitted")
724 } else {
725 panicked = true
726 }
727 }()
728 task.RunProfileSpec(ctx, ProfileExecSpec{
729 Kind: "task", Name: "task", Prompt: "do the thing",
730 SystemPrompt: "sys", AllowNoTools: true,
731 })
732 }()
733 if !panicked {
734 t.Fatal("provider panic must propagate through RunProfileSpec")
735 }
736
737 terminals := 0
738 for _, e := range rec.kinds(event.ToolProgress) {
739 if progressName(e) == event.SubagentProgressStatusName && progressOutput(e) == string(subagentPhaseFailed) {
740 terminals++
741 }
742 }
743 if terminals != 1 {
744 t.Fatalf("panic failed-terminal count = %d, want exactly one", terminals)
745 }
746 }
747
748 func TestBackgroundTaskEmitsQueuedRunningCompleted(t *testing.T) {
749 rec := &recordSink{}
750 jm := jobs.NewManager(event.Discard)
751 defer jm.Close()
752 ctx := jobs.WithManager(withCallContext(context.Background(), "bg-task", rec, nil, false), jm)
753 ctx = jobs.WithSession(ctx, "sess-bg")
754 ctx = WithParentSession(ctx, "sess-bg")
755
756 sched := NewSubagentScheduler(1, 1)
757 holdRelease, err := sched.Acquire(context.Background(), AcquireRequest{Writer: false})
758 if err != nil {
759 t.Fatal(err)
760 }
761 defer holdRelease()
762
763 started := make(chan struct{})
764 task := newTestTaskTool(t, &blockingProvider{started: started}, tool.NewRegistry(), "sys", "", "", nil).
765 WithScheduler(sched)
766
767 done := make(chan string, 1)
768 go func() {
769 out, err := task.Execute(ctx, json.RawMessage(`{"prompt":"work","run_in_background":true,"description":"bg"}`))
770 if err != nil {
771 done <- "err:" + err.Error()
772 return
773 }
774 done <- out
775 }()
776
777 var jobID string
778 select {
779 case out := <-done:
780 if !strings.Contains(out, "Started background task") {
781 t.Fatalf("background start output = %q", out)
782 }
783 jobID = extractJobID(out)
784 case <-time.After(2 * time.Second):
785 t.Fatal("background task did not return a job id while the slot was held")
786 }
787
788 // Registered but not yet executing: the queued status is emitted
789 // synchronously at registration and must never be merged away.
790 queued := false
791 for _, e := range rec.kinds(event.ToolProgress) {
792 if progressName(e) == event.SubagentProgressStatusName && progressOutput(e) == string(subagentPhaseQueued) && e.Tool.ID == "bg-task" {
793 queued = true
794 }
795 }
796 if !queued {
797 t.Fatal("background task never emitted a queued status at registration")
798 }
799
800 // Free the slot: the job acquires it, runs, and emits its terminal.
801 holdRelease()
802 select {
803 case <-started:
804 case <-time.After(2 * time.Second):
805 t.Fatal("background job never started after slot release")
806 }
807 if jobID != "" {
808 result := jm.WaitForSession(context.Background(), "sess-bg", []string{jobID}, 5)
809 if len(result) != 1 || result[0].Status != jobs.Done {
810 t.Fatalf("background job result = %+v, want one completed job", result)
811 }
812 }
813
814 waitStatus := func(want string) {
815 t.Helper()
816 deadline := time.Now().Add(2 * time.Second)
817 for time.Now().Before(deadline) {
818 for _, e := range rec.kinds(event.ToolProgress) {
819 if progressName(e) == event.SubagentProgressStatusName && progressOutput(e) == want && e.Tool.ID == "bg-task" {
820 return
821 }
822 }
823 time.Sleep(time.Millisecond)
824 }
825 t.Fatalf("never saw %q status", want)
826 }
827 waitStatus(string(subagentPhaseRunning))
828 waitStatus(string(subagentPhaseCompleted))
829
830 // Exactly one terminal for the whole lifecycle.
831 terminals := 0
832 for _, e := range rec.kinds(event.ToolProgress) {
833 if progressName(e) == event.SubagentProgressStatusName &&
834 (progressOutput(e) == string(subagentPhaseCompleted) || progressOutput(e) == string(subagentPhaseFailed) || progressOutput(e) == string(subagentPhaseCancelled)) {
835 terminals++
836 }
837 }
838 if terminals != 1 {
839 t.Fatalf("background terminal statuses = %d, want exactly one", terminals)
840 }
841 }
842
843 // TestParallelTasksGroupLifecycleEvents proves the group card gets an
844 // explicit lifecycle from the tool itself: running when children start and
845 // exactly one terminal after every child settles, keyed by the group call ID
846 // — frontends never need to infer group completion from observed children.
847 func TestParallelTasksGroupLifecycleEvents(t *testing.T) {
848 rec := &recordSink{}
849 task := newTestTaskTool(t, parallelStaticProvider{}, tool.NewRegistry(), "sys", "", "", nil)
850 parallel := NewParallelTasksTool(task, tool.NewRegistry())
851 ctx := withCallContext(context.Background(), "parallel-call", rec, nil, false)
852
853 if _, err := parallel.Execute(ctx, json.RawMessage(`{
854 "tasks": [{"prompt": "first"}, {"prompt": "second"}]
855 }`)); err != nil {
856 t.Fatalf("Execute: %v", err)
857 }
858
859 var groupStatuses []string
860 childStatuses := map[string][]string{}
861 for _, e := range rec.kinds(event.ToolProgress) {
862 if progressName(e) != event.SubagentProgressStatusName {
863 continue
864 }
865 switch {
866 case e.Tool.ID == "parallel-call":
867 groupStatuses = append(groupStatuses, progressOutput(e))
868 case strings.HasPrefix(e.Tool.ID, "parallel-call/"):
869 childStatuses[e.Tool.ID] = append(childStatuses[e.Tool.ID], progressOutput(e))
870 }
871 }
872 if len(groupStatuses) != 2 || groupStatuses[0] != string(subagentPhaseRunning) || groupStatuses[1] != string(subagentPhaseCompleted) {
873 t.Fatalf("group lifecycle = %v, want running → completed", groupStatuses)
874 }
875 for id, st := range childStatuses {
876 if len(st) < 2 || st[0] != string(subagentPhaseRunning) || st[len(st)-1] != string(subagentPhaseCompleted) {
877 t.Fatalf("child %s lifecycle = %v, want running → … → completed", id, st)
878 }
879 terminals := 0
880 for _, out := range st {
881 if isTerminalStatusOutput(out) {
882 terminals++
883 }
884 }
885 if terminals != 1 {
886 t.Fatalf("child %s terminals = %d, want exactly one", id, terminals)
887 }
888 }
889 if len(childStatuses) != 2 {
890 t.Fatalf("child status cards = %d, want 2", len(childStatuses))
891 }
892 }
893
894 // TestParallelTasksGroupLifecycleCancelled proves a cancelled group emits
895 // exactly one cancelled terminal.
896 func TestParallelTasksGroupLifecycleCancelled(t *testing.T) {
897 rec := &recordSink{}
898 started := make(chan struct{})
899 task := newTestTaskTool(t, &cancelBlockingProvider{started: started}, tool.NewRegistry(), "sys", "", "", nil)
900 parallel := NewParallelTasksTool(task, tool.NewRegistry())
901 ctx, cancel := context.WithCancel(withCallContext(context.Background(), "parallel-call", rec, nil, false))
902 defer cancel()
903 go func() {
904 <-started
905 cancel()
906 }()
907
908 if _, err := parallel.Execute(ctx, json.RawMessage(`{
909 "tasks": [{"prompt": "first"}, {"prompt": "second"}]
910 }`)); err == nil {
911 t.Fatal("cancelled Execute must return an error")
912 }
913
914 terminals := 0
915 for _, e := range rec.kinds(event.ToolProgress) {
916 if progressName(e) != event.SubagentProgressStatusName || e.Tool.ID != "parallel-call" {
917 continue
918 }
919 if isTerminalStatusOutput(progressOutput(e)) {
920 terminals++
921 if progressOutput(e) != string(subagentPhaseCancelled) {
922 t.Fatalf("group terminal = %q, want cancelled", progressOutput(e))
923 }
924 }
925 }
926 if terminals != 1 {
927 t.Fatalf("group terminals = %d, want exactly one", terminals)
928 }
929 }
930
931 type cancelBlockingProvider struct {
932 started chan struct{}
933 once sync.Once
934 }
935
936 func (p *cancelBlockingProvider) Name() string { return "cancel-blocking" }
937
938 func (p *cancelBlockingProvider) Stream(ctx context.Context, _ provider.Request) (<-chan provider.Chunk, error) {
939 p.once.Do(func() { close(p.started) })
940 <-ctx.Done()
941 return nil, ctx.Err()
942 }
943
944 // TestParallelTasksGroupLifecycleFailedOnValidation proves validation failures
945 // still emit a failed terminal for the group card.
946 func TestParallelTasksGroupLifecycleFailedOnValidation(t *testing.T) {
947 rec := &recordSink{}
948 parallel := &ParallelTasksTool{} // unconfigured: fails after merger setup
949 ctx := withCallContext(context.Background(), "parallel-call", rec, nil, false)
950
951 if _, err := parallel.Execute(ctx, json.RawMessage(`{"tasks":[{"prompt":"x"}]}`)); err == nil {
952 t.Fatal("unconfigured parallel_tasks must fail")
953 }
954
955 terminals := 0
956 ran := false
957 for _, e := range rec.kinds(event.ToolProgress) {
958 if progressName(e) != event.SubagentProgressStatusName || e.Tool.ID != "parallel-call" {
959 continue
960 }
961 if progressOutput(e) == string(subagentPhaseRunning) {
962 ran = true
963 }
964 if isTerminalStatusOutput(progressOutput(e)) {
965 terminals++
966 if progressOutput(e) != string(subagentPhaseFailed) {
967 t.Fatalf("group terminal = %q, want failed", progressOutput(e))
968 }
969 }
970 }
971 if ran {
972 t.Fatal("validation failure must not emit running")
973 }
974 if terminals != 1 {
975 t.Fatalf("group terminals = %d, want exactly one failed", terminals)
976 }
977 }
978
979 func isTerminalStatusOutput(out string) bool {
980 return out == string(subagentPhaseCompleted) || out == string(subagentPhaseFailed) || out == string(subagentPhaseCancelled)
981 }
982
982 lines GO