返回 DeepSeek-Reasonix
turn_loop_test.go
根目录 / internal / control / turn_loop_test.go
1 package control
2
3 import (
4 "context"
5 "os"
6 "path/filepath"
7 "sync"
8 "sync/atomic"
9 "testing"
10 "time"
11
12 "reasonix/internal/agent"
13 "reasonix/internal/agent/testutil"
14 "reasonix/internal/event"
15 "reasonix/internal/provider"
16 "reasonix/internal/session"
17 "reasonix/internal/tool"
18 )
19
20 func exclusiveTestController(t *testing.T, sink event.Sink) (*Controller, *session.Service, *session.Runtime) {
21 t.Helper()
22 if sink == nil {
23 sink = event.Discard
24 }
25 service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4")))
26 if err != nil {
27 t.Fatal(err)
28 }
29 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
30 runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "loop"})
31 if err != nil {
32 t.Fatal(err)
33 }
34 exec := agent.New(testutil.NewMock("test"), tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, sink)
35 c := newOwnedTestController(t, Options{
36 Runner: exec, Executor: exec, Sink: sink,
37 SessionService: service, SessionRuntime: runtime, ExclusiveSession: true,
38 })
39 return c, service, runtime
40 }
41
42 func TestIdleCancelThenSendProducesNewTurnID(t *testing.T) {
43 done := make(chan event.Event, 4)
44 c, _, runtime := exclusiveTestController(t, event.FuncSink(func(e event.Event) {
45 if e.Kind == event.TurnDone {
46 done <- e
47 }
48 }))
49 receipt := c.CancelSession()
50 if !receipt.Accepted || !receipt.AlreadyIdle {
51 t.Fatalf("idle cancel = %+v", receipt)
52 }
53 started := make(chan struct{})
54 if got := c.runGuarded(func(context.Context) error {
55 close(started)
56 return nil
57 }); got != turnStarted {
58 t.Fatalf("send after idle cancel = %v", got)
59 }
60 <-started
61 terminal := waitTurnDoneEvent(t, done)
62 if terminal.TurnID == "" {
63 t.Fatal("first turn after idle cancel has no durable turn id")
64 }
65 waitIdleAdmission(t, c)
66 if runtime.StateSnapshot().Phase != session.RuntimeIdle {
67 t.Fatalf("runtime phase = %s", runtime.StateSnapshot().Phase)
68 }
69 }
70
71 func TestCancelDuringStreamingAndToolAndApproval(t *testing.T) {
72 for _, name := range []string{"stream", "tool", "approval"} {
73 t.Run(name, func(t *testing.T) {
74 started := make(chan struct{})
75 c := newOwnedTestController(t, Options{Sink: event.Discard})
76 t.Cleanup(c.Close)
77 c.runGuarded(func(ctx context.Context) error {
78 close(started)
79 <-ctx.Done()
80 return ctx.Err()
81 })
82 <-started
83 c.CancelSession()
84 waitIdleAdmission(t, c)
85 })
86 }
87 }
88
89 func TestInputAfterCancelWakesOnceFIFO(t *testing.T) {
90 var ran atomic.Int32
91 order := make(chan int, 2)
92 started := make(chan struct{})
93 release := make(chan struct{})
94 c := newOwnedTestController(t, Options{Sink: event.Discard})
95 t.Cleanup(c.Close)
96 c.runGuarded(func(ctx context.Context) error {
97 close(started)
98 <-ctx.Done()
99 <-release
100 return ctx.Err()
101 })
102 <-started
103 c.CancelSession()
104 if got := c.runGuarded(func(context.Context) error {
105 ran.Add(1)
106 order <- 1
107 return nil
108 }); got != turnParked {
109 t.Fatalf("first queued = %v, want parked", got)
110 }
111 if got := c.runGuarded(func(context.Context) error {
112 ran.Add(1)
113 order <- 2
114 return nil
115 }); got != turnParked {
116 t.Fatalf("second queued = %v, want parked", got)
117 }
118 close(release)
119 waitIdleAdmission(t, c)
120 if ran.Load() != 2 {
121 t.Fatalf("queued bodies ran %d times, want 2", ran.Load())
122 }
123 if first, second := <-order, <-order; first != 1 || second != 2 {
124 t.Fatalf("fifo order = %d,%d", first, second)
125 }
126 }
127
128 func TestSlowExitDoesNotStartNextTurn(t *testing.T) {
129 started := make(chan struct{})
130 exit := make(chan struct{})
131 nextStarted := make(chan struct{})
132 c := newOwnedTestController(t, Options{Sink: event.Discard})
133 t.Cleanup(c.Close)
134 c.runGuarded(func(ctx context.Context) error {
135 close(started)
136 <-ctx.Done()
137 <-exit
138 return ctx.Err()
139 })
140 <-started
141 c.CancelSession()
142 if got := c.runGuarded(func(context.Context) error {
143 close(nextStarted)
144 return nil
145 }); got != turnParked {
146 t.Fatalf("next admission = %v, want parked until slow exit", got)
147 }
148 select {
149 case <-nextStarted:
150 t.Fatal("next turn started before the cancelled body exited")
151 case <-time.After(50 * time.Millisecond):
152 }
153 close(exit)
154 select {
155 case <-nextStarted:
156 case <-time.After(5 * time.Second):
157 t.Fatal("queued turn did not start after slow exit")
158 }
159 waitIdleAdmission(t, c)
160 }
161
162 func TestDuplicateStopAndConcurrentFinish(t *testing.T) {
163 started := make(chan struct{})
164 c := newOwnedTestController(t, Options{Sink: event.Discard})
165 t.Cleanup(c.Close)
166 c.runGuarded(func(ctx context.Context) error {
167 close(started)
168 <-ctx.Done()
169 return ctx.Err()
170 })
171 <-started
172 var wg sync.WaitGroup
173 for range 8 {
174 wg.Go(func() {
175 c.CancelSession()
176 c.Cancel()
177 })
178 }
179 wg.Wait()
180 waitIdleAdmission(t, c)
181 }
182
183 func TestEachStartedTurnHasOneTerminalEvent(t *testing.T) {
184 var terminals atomic.Int32
185 c := newOwnedTestController(t, Options{Sink: event.FuncSink(func(e event.Event) {
186 if e.Kind == event.TurnDone {
187 terminals.Add(1)
188 }
189 })})
190 t.Cleanup(c.Close)
191 for range 3 {
192 c.runGuarded(func(context.Context) error { return nil })
193 waitIdleAdmission(t, c)
194 }
195 if got := terminals.Load(); got != 3 {
196 t.Fatalf("terminal events = %d, want 3", got)
197 }
198 }
199
200 func TestStopHistoryReplaceThenSendGetsNewTurnID(t *testing.T) {
201 done := make(chan event.Event, 8)
202 c, service, runtime := exclusiveTestController(t, event.FuncSink(func(e event.Event) {
203 if e.Kind == event.TurnDone {
204 done <- e
205 }
206 }))
207 if err := c.RecordSessionMessages(t.Context(), "before-stop", []provider.Message{
208 {ID: "retained-question", Role: provider.RoleUser, Content: "keep this question"},
209 }); err != nil {
210 t.Fatal(err)
211 }
212 c.restoreExecutorFromSessionEvents()
213 if _, err := runtime.Session().Flush(t.Context()); err != nil {
214 t.Fatal(err)
215 }
216 if _, err := service.Query().HistoryShape(t.Context(), runtime.Ref()); err != nil {
217 t.Fatal(err)
218 }
219 started := make(chan struct{})
220 c.runGuarded(func(ctx context.Context) error {
221 close(started)
222 <-ctx.Done()
223 c.replaceSessionAfterCancel(c.executor.Session().Snapshot())
224 return ctx.Err()
225 })
226 <-started
227 c.CancelSession()
228 first := waitTurnDoneEvent(t, done)
229 if first.TurnID == "" {
230 t.Fatal("cancelled turn has no durable id")
231 }
232 waitIdleAdmission(t, c)
233 if err := c.turnEventLedgerError(); err != nil {
234 t.Fatalf("cancel poisoned admission: %v", err)
235 }
236 // Switching after Stop reads the same durable rewrite through the history
237 // index. Admission alone used to pass while these reads failed with a
238 // duplicate (message_id, version) key.
239 other, err := service.Create(t.Context(), session.CreateOptions{SessionID: "other"})
240 if err != nil {
241 t.Fatal(err)
242 }
243 if _, err := c.OpenSession(t.Context(), other.Ref()); err != nil {
244 t.Fatalf("switch away after stop: %v", err)
245 }
246 if _, err := service.Query().HistoryShape(t.Context(), runtime.Ref()); err != nil {
247 t.Fatalf("read stopped session after switching away: %v", err)
248 }
249 if _, err := c.OpenSession(t.Context(), runtime.Ref()); err != nil {
250 t.Fatalf("switch back after stop: %v", err)
251 }
252 page, err := service.Query().ReadHistoryWindow(t.Context(), runtime.Ref(), session.HistoryWindowRequest{Anchor: "newest"})
253 if err != nil || page.Status != "ready" {
254 t.Fatalf("reopened stopped history: %+v, %v", page, err)
255 }
256 found := false
257 for _, message := range page.Messages {
258 found = found || message.MessageID == "retained-question"
259 }
260 if !found {
261 t.Fatal("stopped history lost the retained question")
262 }
263 secondStarted := make(chan struct{})
264 if got := c.runGuarded(func(context.Context) error {
265 close(secondStarted)
266 return nil
267 }); got != turnStarted {
268 t.Fatalf("resend after cancel = %v", got)
269 }
270 <-secondStarted
271 second := waitTurnDoneEvent(t, done)
272 if second.TurnID == "" || second.TurnID == first.TurnID {
273 t.Fatalf("second turn id = %q, first = %q", second.TurnID, first.TurnID)
274 }
275 waitIdleAdmission(t, c)
276 if runtime.StateSnapshot().Phase != session.RuntimeIdle {
277 t.Fatalf("runtime phase = %s", runtime.StateSnapshot().Phase)
278 }
279 }
280
281 func TestPlanStateTailWriteDuringStopDoesNotPoison(t *testing.T) {
282 c, _, runtime := exclusiveTestController(t, event.Discard)
283 started := make(chan struct{})
284 c.runGuarded(func(ctx context.Context) error {
285 close(started)
286 <-ctx.Done()
287 if err := c.appendDomainState("plan/state", []byte(`{"enabled":true}`), "stop-race"); err != nil {
288 t.Errorf("plan/state during stop: %v", err)
289 }
290 return ctx.Err()
291 })
292 <-started
293 c.CancelSession()
294 waitIdleAdmission(t, c)
295 if err := c.turnEventLedgerError(); err != nil {
296 t.Fatalf("plan/state vs stop poisoned ledger: %v", err)
297 }
298 if got := c.runGuarded(func(context.Context) error { return nil }); got != turnStarted {
299 t.Fatalf("admission after plan/state race = %v", got)
300 }
301 waitIdleAdmission(t, c)
302 if runtime.StateSnapshot().Phase != session.RuntimeIdle {
303 t.Fatalf("runtime phase = %s", runtime.StateSnapshot().Phase)
304 }
305 }
306
307 func TestTimeoutRecoveryDropsLateEventsFromNewTurn(t *testing.T) {
308 c := newOwnedTestController(t, Options{Sink: event.Discard, SessionDir: t.TempDir(), SessionPath: filepath.Join(t.TempDir(), "session.jsonl")})
309 t.Cleanup(c.Close)
310 c.testCancelGrace = 10 * time.Millisecond
311 started := make(chan struct{})
312 hold := make(chan struct{})
313 exited := make(chan struct{})
314 c.runGuarded(func(ctx context.Context) error {
315 defer close(exited)
316 close(started)
317 <-ctx.Done()
318 <-hold
319 return ctx.Err()
320 })
321 <-started
322 c.CancelSession()
323 deadline := time.Now().Add(5 * time.Second)
324 for time.Now().Before(deadline) {
325 c.mu.Lock()
326 phase := c.turns.phase
327 c.mu.Unlock()
328 if phase == session.RuntimeRecoveryRequired {
329 break
330 }
331 time.Sleep(time.Millisecond)
332 }
333 late := event.Event{Kind: event.ToolResult, TurnID: "not-the-next-turn", Tool: event.Tool{ID: "late", Name: "bash"}}
334 if !c.discardLateTurnEvent(late) {
335 t.Fatal("late tool result must not apply to a later turn")
336 }
337 if got := c.runGuarded(func(context.Context) error { return nil }); got != turnDroppedWriteAuthority {
338 t.Fatalf("admission during recovery = %v, want blocked", got)
339 }
340 close(hold)
341 select {
342 case <-exited:
343 case <-time.After(5 * time.Second):
344 t.Fatal("timed-out turn did not exit after release")
345 }
346 c.autosaveWG.Wait()
347 }
348
349 func TestOldControllerUnbindDoesNotClearNewGeneration(t *testing.T) {
350 service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4")))
351 if err != nil {
352 t.Fatal(err)
353 }
354 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
355 runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "handoff"})
356 if err != nil {
357 t.Fatal(err)
358 }
359 exec := agent.New(testutil.NewMock("test"), tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard)
360 first := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionService: service, SessionRuntime: runtime, ExclusiveSession: true})
361 first.mu.Lock()
362 oldGen := first.turns.generation
363 first.mu.Unlock()
364 second := newOwnedTestController(t, Options{Executor: exec, Sink: event.Discard, SessionService: service, SessionRuntime: runtime, ExclusiveSession: true})
365 t.Cleanup(func() { first.Close(); second.Close() })
366 if err := second.ActivateSessionExecution(oldGen); err != nil {
367 t.Fatalf("activate replacement: %v", err)
368 }
369 runtime.UnbindExecution(oldGen)
370 started := make(chan struct{})
371 second.runGuarded(func(ctx context.Context) error {
372 close(started)
373 <-ctx.Done()
374 return ctx.Err()
375 })
376 <-started
377 if !runtime.Cancel() {
378 t.Fatal("new generation lost Stop after old unbind")
379 }
380 waitIdleAdmission(t, second)
381 }
382
383 func TestOneRuntimeCannotRunTwoControllerLoops(t *testing.T) {
384 service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4")))
385 if err != nil {
386 t.Fatal(err)
387 }
388 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
389 runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "single-loop"})
390 if err != nil {
391 t.Fatal(err)
392 }
393 newController := func() *Controller {
394 exec := agent.New(testutil.NewMock("test"), tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard)
395 return newOwnedTestController(t, Options{
396 Runner: exec, Executor: exec, Sink: event.Discard,
397 SessionService: service, SessionRuntime: runtime, ExclusiveSession: true,
398 })
399 }
400 first := newController()
401 second := newController()
402 t.Cleanup(func() { first.Close(); second.Close() })
403
404 firstStarted := make(chan struct{})
405 releaseFirst := make(chan struct{})
406 if got := first.runGuarded(func(context.Context) error {
407 close(firstStarted)
408 <-releaseFirst
409 return nil
410 }); got != turnStarted {
411 t.Fatalf("first admission = %v, want started", got)
412 }
413 <-firstStarted
414
415 secondStarted := make(chan struct{})
416 if got := second.runGuarded(func(context.Context) error {
417 close(secondStarted)
418 return nil
419 }); got == turnStarted {
420 t.Fatal("same session runtime admitted a second controller loop concurrently")
421 }
422 select {
423 case <-secondStarted:
424 t.Fatal("second controller body ran")
425 default:
426 }
427 close(releaseFirst)
428 waitIdleAdmission(t, first)
429 }
430
431 func TestReplacementModelContextCommitsOnlyWithExecutionCutover(t *testing.T) {
432 service, err := session.NewService("desktop", session.NewFilesystemPersistence(filepath.Join(t.TempDir(), "sessions-v4")))
433 if err != nil {
434 t.Fatal(err)
435 }
436 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
437 runtime, err := service.Create(t.Context(), session.CreateOptions{SessionID: "atomic-cutover"})
438 if err != nil {
439 t.Fatal(err)
440 }
441 newController := func() *Controller {
442 exec := agent.New(testutil.NewMock("test"), tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard)
443 return newOwnedTestController(t, Options{
444 Runner: exec, Executor: exec, Sink: event.Discard,
445 SessionService: service, SessionRuntime: runtime, ExclusiveSession: true,
446 })
447 }
448 first := newController()
449 second := newController()
450 t.Cleanup(func() { first.Close(); second.Close() })
451 oldGeneration := first.ExecutionGeneration()
452 before := runtime.Session().ExecutionSnapshot().EventSequence
453 messages := []provider.Message{
454 {Role: provider.RoleSystem, Content: "replacement system"},
455 {Role: provider.RoleUser, Content: "preserve me"},
456 }
457 if err := second.AdoptRebuiltModelContext(messages); err != nil {
458 t.Fatalf("stage replacement context: %v", err)
459 }
460 if got := runtime.Session().ExecutionSnapshot().EventSequence; got != before {
461 t.Fatalf("unpublished candidate changed sequence from %d to %d", before, got)
462 }
463 if !runtime.OwnsExecution(oldGeneration) {
464 t.Fatal("staging replacement context stole outgoing execution ownership")
465 }
466 if err := ActivateControllerReplacement(first, second); err != nil {
467 t.Fatalf("activate replacement: %v", err)
468 }
469 after := runtime.Session().ExecutionSnapshot()
470 if after.EventSequence <= before {
471 t.Fatalf("activation did not commit staged model context: before=%d after=%d", before, after.EventSequence)
472 }
473 if !runtime.OwnsExecution(second.ExecutionGeneration()) || runtime.OwnsExecution(oldGeneration) {
474 t.Fatal("execution ownership did not transfer atomically")
475 }
476 if got := after.Projection.ModelMessages; len(got) != len(messages) || got[len(got)-1].Content != "preserve me" {
477 t.Fatalf("committed model context = %+v", got)
478 }
479 }
480
481 func TestDiscardedReplacementLeavesOutgoingExecutionOwner(t *testing.T) {
482 first, service, runtime := exclusiveTestController(t, event.Discard)
483 exec := agent.New(testutil.NewMock("test"), tool.NewRegistry(), agent.NewSession("system"), agent.Options{}, event.Discard)
484 candidate := newOwnedTestController(t, Options{
485 Runner: exec, Executor: exec, Sink: event.Discard,
486 SessionService: service, SessionRuntime: runtime, ExclusiveSession: true,
487 })
488 oldGeneration := first.ExecutionGeneration()
489 if err := candidate.AdoptRebuiltModelContext([]provider.Message{{Role: provider.RoleSystem, Content: "discarded"}}); err != nil {
490 t.Fatal(err)
491 }
492 candidate.ReleaseResources()
493 if !runtime.OwnsExecution(oldGeneration) {
494 t.Fatal("discarded candidate cleared outgoing execution owner")
495 }
496 started := make(chan struct{})
497 if got := first.runGuarded(func(context.Context) error { close(started); return nil }); got != turnStarted {
498 t.Fatalf("outgoing admission after candidate discard = %v", got)
499 }
500 <-started
501 waitIdleAdmission(t, first)
502 }
503
504 func TestCloseDropsLatchedWake(t *testing.T) {
505 done := make(chan event.Event, 2)
506 started := make(chan struct{})
507 exit := make(chan struct{})
508 nextStarted := make(chan struct{})
509 c := newOwnedTestController(t, Options{Sink: event.FuncSink(func(e event.Event) {
510 if e.Kind == event.TurnDone {
511 done <- e
512 }
513 })})
514 c.runGuarded(func(ctx context.Context) error {
515 close(started)
516 <-ctx.Done()
517 <-exit
518 return ctx.Err()
519 })
520 <-started
521 c.CancelSession()
522 if got := c.runGuarded(func(context.Context) error {
523 close(nextStarted)
524 return nil
525 }); got != turnParked {
526 t.Fatalf("wake during abort = %v, want parked", got)
527 }
528 c.Close()
529 close(exit)
530 waitTurnDoneEvent(t, done)
531 select {
532 case <-nextStarted:
533 t.Fatal("close started a latched wake")
534 default:
535 }
536 }
537
538 func TestCloseKeepsExecutionBoundUntilTerminalPublicationCompletes(t *testing.T) {
539 entered := make(chan struct{}, 1)
540 release := make(chan struct{})
541 c, service, runtime := exclusiveTestController(t, holdFinishingWindow(release, entered, nil))
542 generation := c.ExecutionGeneration()
543 started := make(chan struct{})
544 if got := c.runGuarded(func(ctx context.Context) error {
545 close(started)
546 <-ctx.Done()
547 return ctx.Err()
548 }); got != turnStarted {
549 t.Fatalf("admission = %v, want started", got)
550 }
551 <-started
552 c.Close()
553 <-entered
554
555 if got := runtime.StateSnapshot().Phase; got != session.RuntimeFinalizing {
556 t.Fatalf("runtime phase while terminal publication is blocked = %s, want finalizing", got)
557 }
558 if !runtime.OwnsExecution(generation) {
559 t.Fatal("close released execution ownership before terminal publication")
560 }
561 if current, ok := service.Runtime(runtime.Ref()); !ok || current != runtime {
562 t.Fatal("close retired the session runtime before terminal publication")
563 }
564
565 close(release)
566 deadline := time.Now().Add(5 * time.Second)
567 for time.Now().Before(deadline) {
568 if !runtime.OwnsExecution(generation) {
569 return
570 }
571 time.Sleep(time.Millisecond)
572 }
573 t.Fatal("execution ownership was not released after terminal publication")
574 }
575
576 func TestSessionOpenFailureStillFailsClosed(t *testing.T) {
577 blocked := filepath.Join(t.TempDir(), "not-a-directory")
578 if err := os.WriteFile(blocked, []byte("block"), 0o600); err != nil {
579 t.Fatal(err)
580 }
581 done := make(chan event.Event, 1)
582 c := newOwnedTestController(t, Options{
583 Sink: event.FuncSink(func(e event.Event) {
584 if e.Kind == event.TurnDone {
585 done <- e
586 }
587 }),
588 SessionDir: t.TempDir(), SessionPath: filepath.Join(blocked, "session.jsonl"),
589 })
590 t.Cleanup(c.Close)
591 c.Submit("must fail closed")
592 select {
593 case <-done:
594 case <-time.After(5 * time.Second):
595 t.Fatal("open failure did not complete the admission attempt")
596 }
597 if err := c.turnEventLedgerError(); err == nil {
598 t.Fatal("open failure did not fail closed")
599 }
600 }
601
601 lines GO