返回 DeepSeek-Reasonix
turn_events_test.go
根目录 / internal / control / turn_events_test.go
1 package control
2
3 import (
4 "context"
5 "errors"
6 "os"
7 "path/filepath"
8 "sync/atomic"
9 "testing"
10 "time"
11
12 "reasonix/internal/event"
13 "reasonix/internal/turnevent"
14 )
15
16 type turnEventGateRunner struct {
17 started chan struct{}
18 release chan struct{}
19 calls atomic.Int32
20 }
21
22 func (r *turnEventGateRunner) Run(context.Context, string) error {
23 r.calls.Add(1)
24 close(r.started)
25 <-r.release
26 return nil
27 }
28
29 func TestTurnAdmissionIsDurableBeforeRunnerStarts(t *testing.T) {
30 dir := t.TempDir()
31 runner := &turnEventGateRunner{started: make(chan struct{}), release: make(chan struct{})}
32 done := make(chan event.Event, 1)
33 c := newOwnedTestController(t, Options{
34 Runner: runner,
35 Sink: event.FuncSink(func(e event.Event) {
36 if e.Kind == event.TurnDone {
37 done <- e
38 }
39 }),
40 SessionDir: dir, SessionPath: filepath.Join(dir, "session.jsonl"),
41 })
42 t.Cleanup(c.Close)
43
44 c.Submit("run")
45 select {
46 case <-runner.started:
47 case <-time.After(5 * time.Second):
48 t.Fatal("runner did not start")
49 }
50 records, err := c.TurnEventsAfter(0)
51 if err != nil {
52 t.Fatalf("TurnEventsAfter: %v", err)
53 }
54 if len(records) < 2 || records[0].Status != event.TurnQueued || records[1].Kind != "turn_started" || records[1].Status != event.TurnInProgress {
55 t.Fatalf("admission prefix = %+v, want queued then durable in_progress start", records)
56 }
57 close(runner.release)
58 select {
59 case <-done:
60 case <-time.After(5 * time.Second):
61 t.Fatal("turn did not finish")
62 }
63 records, err = c.TurnEventsAfter(0)
64 if err != nil {
65 t.Fatalf("TurnEventsAfter terminal: %v", err)
66 }
67 started := 0
68 for _, record := range records {
69 if record.Kind == "turn_started" {
70 started++
71 }
72 }
73 if started != 1 {
74 t.Fatalf("turn_started records = %d, want exactly one", started)
75 }
76 }
77
78 func TestTurnAdmissionLedgerFailureDoesNotRunProvider(t *testing.T) {
79 dir := t.TempDir()
80 blockedParent := filepath.Join(dir, "not-a-directory")
81 if err := os.WriteFile(blockedParent, []byte("block"), 0o600); err != nil {
82 t.Fatalf("write blocker: %v", err)
83 }
84 runner := &turnEventGateRunner{started: make(chan struct{}), release: make(chan struct{})}
85 done := make(chan event.Event, 1)
86 c := newOwnedTestController(t, Options{
87 Runner: runner,
88 Sink: event.FuncSink(func(e event.Event) {
89 if e.Kind == event.TurnDone {
90 done <- e
91 }
92 }),
93 SessionDir: dir, SessionPath: filepath.Join(blockedParent, "session.jsonl"),
94 })
95 t.Cleanup(c.Close)
96
97 c.Submit("must not reach provider")
98 select {
99 case terminal := <-done:
100 if !errors.Is(terminal.Err, turnevent.ErrTurnLedgerUnavailable) {
101 t.Fatalf("terminal error = %v, want explicit ledger admission failure", terminal.Err)
102 }
103 case <-time.After(5 * time.Second):
104 t.Fatal("failed admission did not terminate")
105 }
106 if got := runner.calls.Load(); got != 0 {
107 t.Fatalf("runner calls = %d, want provider side effects blocked", got)
108 }
109 }
110
111 func TestAsyncStreamLedgerFailureCancelsTurnWithoutPublishingChunk(t *testing.T) {
112 root := filepath.Join(t.TempDir(), "session-dir")
113 if err := os.MkdirAll(root, 0o700); err != nil {
114 t.Fatal(err)
115 }
116 started := make(chan struct{})
117 release := make(chan struct{})
118 cancelled := make(chan struct{})
119 done := make(chan event.Event, 1)
120 var publishedText atomic.Int32
121 c := newOwnedTestController(t, Options{
122 Sink: event.FuncSink(func(e event.Event) {
123 if e.Kind == event.Text {
124 publishedText.Add(1)
125 }
126 if e.Kind == event.TurnDone {
127 done <- e
128 }
129 }),
130 SessionDir: root, SessionPath: filepath.Join(root, "session.jsonl"),
131 })
132 t.Cleanup(c.Close)
133
134 c.runGuarded(func(ctx context.Context) error {
135 close(started)
136 <-release
137 c.sink.Emit(event.Event{Kind: event.ToolResult, Tool: event.Tool{ID: "failed-store", Name: "probe", Output: "must stay behind the v3 log"}})
138 <-ctx.Done()
139 close(cancelled)
140 return ctx.Err()
141 })
142 <-started
143 v3 := c.sessionEventStore()
144 if v3 == nil {
145 t.Fatal("controller did not open a v3 event store")
146 }
147 if err := v3.Close(context.Background()); err != nil {
148 t.Fatalf("close active v3 handle: %v", err)
149 }
150 close(release)
151
152 select {
153 case <-cancelled:
154 case <-time.After(5 * time.Second):
155 t.Fatal("asynchronous stream persistence failure did not cancel the turn")
156 }
157 select {
158 case terminal := <-done:
159 if terminal.Status != event.TurnFailed || terminal.Err == nil {
160 t.Fatalf("control-plane terminal = %+v, want explicit storage failure", terminal)
161 }
162 case <-time.After(5 * time.Second):
163 t.Fatal("storage failure did not release frontend running state")
164 }
165 if got := publishedText.Load(); got != 0 {
166 t.Fatalf("published text chunks = %d, want none before durable append", got)
167 }
168 if err := c.turnEventLedgerError(); err == nil {
169 t.Fatal("controller accepted new work after the ledger was poisoned")
170 }
171 }
172
172 lines GO