返回 DeepSeek-Reasonix
execution_test.go
根目录 / internal / session / execution_test.go
1 package session
2
3 import (
4 "context"
5 "errors"
6 "sync"
7 "testing"
8 )
9
10 type testExecution struct {
11 mu sync.Mutex
12 phase RuntimePhase
13 cancel context.CancelFunc
14 ctx context.Context
15 gen uint64
16 runtime *Runtime
17 }
18
19 type blockingRejectExecution struct {
20 entered chan<- struct{}
21 release <-chan struct{}
22 }
23
24 type advancingCancelExecution struct {
25 *testExecution
26 advance func()
27 }
28
29 func (e *advancingCancelExecution) Cancel() bool { e.advance(); return true }
30
31 func TestMaintenanceCancelAcknowledgementCannotCancelSuccessor(t *testing.T) {
32 _, runtime := reviewRuntime(t)
33 exec := &advancingCancelExecution{testExecution: &testExecution{runtime: runtime}}
34 exec.gen = runtime.BindExecution(exec)
35 if !runtime.BeginExecution(exec.gen, MaintenanceActivity) {
36 t.Fatal("maintenance admission failed")
37 }
38 exec.advance = func() {
39 if !runtime.FinishMaintenanceExecution(exec.gen) || !runtime.BeginExecution(exec.gen, "turn") {
40 t.Fatal("could not hand off to queued turn")
41 }
42 }
43 if !runtime.Cancel() {
44 t.Fatal("cancellation was not acknowledged")
45 }
46 if got := runtime.StateSnapshot(); got.Phase != RuntimeRunning || got.Activity != "turn" {
47 t.Fatalf("late acknowledgement changed successor: %+v", got)
48 }
49 runtime.NoteExecution(exec.gen, RuntimeIdle, "")
50 runtime.UnbindExecution(exec.gen)
51 }
52
53 func (e *blockingRejectExecution) Snapshot() RuntimeSnapshot {
54 return RuntimeSnapshot{Phase: RuntimeIdle}
55 }
56 func (e *blockingRejectExecution) Cancel() bool {
57 e.entered <- struct{}{}
58 <-e.release
59 return false
60 }
61
62 func bindTestExecution(t *testing.T, runtime *Runtime, name string) (context.Context, *testExecution) {
63 t.Helper()
64 ctx, cancel := context.WithCancel(context.Background())
65 exec := &testExecution{phase: RuntimeRunning, cancel: cancel, ctx: ctx, gen: 1, runtime: runtime}
66 exec.gen = runtime.BindExecution(exec)
67 runtime.NoteExecution(exec.gen, RuntimeRunning, name)
68 t.Cleanup(exec.Finish)
69 return ctx, exec
70 }
71
72 func (e *testExecution) Snapshot() RuntimeSnapshot {
73 e.mu.Lock()
74 defer e.mu.Unlock()
75 return RuntimeSnapshot{Phase: e.phase}
76 }
77
78 func (e *testExecution) Cancel() bool {
79 e.mu.Lock()
80 defer e.mu.Unlock()
81 if e.phase != RuntimeRunning && e.phase != RuntimeCancelling {
82 return false
83 }
84 e.phase = RuntimeCancelling
85 if e.cancel != nil {
86 e.cancel()
87 }
88 return true
89 }
90
91 func (e *testExecution) Finish() {
92 e.mu.Lock()
93 if e.phase == RuntimeIdle || e.phase == RuntimeClosed {
94 e.mu.Unlock()
95 return
96 }
97 e.phase = RuntimeIdle
98 e.mu.Unlock()
99 if e.runtime != nil {
100 e.runtime.NoteExecution(e.gen, RuntimeIdle, "")
101 }
102 }
103
104 func TestBindExecutionDoesNotSilentlyReplaceExistingOwner(t *testing.T) {
105 _, runtime := reviewRuntime(t)
106 first := &testExecution{phase: RuntimeRunning, runtime: runtime}
107 first.gen = runtime.BindExecution(first)
108 runtime.NoteExecution(first.gen, RuntimeRunning, "first")
109
110 second := &testExecution{phase: RuntimeRunning, runtime: runtime}
111 if gen := runtime.BindExecution(second); gen != 0 {
112 t.Fatalf("second bind generation = %d, want rejection", gen)
113 }
114 if !runtime.Cancel() {
115 t.Fatal("cancel did not reach the original execution owner")
116 }
117 if got := first.phase; got != RuntimeCancelling {
118 t.Fatalf("original owner phase = %s, want cancelling", got)
119 }
120 // A runtime with a live execution refuses to close, which would strand its
121 // writer lease past the test.
122 first.Finish()
123 }
124
125 func TestUnbindDoesNotClearNewerExecutionGeneration(t *testing.T) {
126 _, runtime := reviewRuntime(t)
127 first := &testExecution{phase: RuntimeRunning, runtime: runtime}
128 second := &testExecution{phase: RuntimeRunning, runtime: runtime}
129 first.gen = runtime.BindExecution(first)
130 second.gen = runtime.ReplaceExecution(first.gen, second)
131 if second.gen == 0 {
132 t.Fatal("replace idle execution")
133 }
134 runtime.NoteExecution(second.gen, RuntimeRunning, "new")
135 runtime.UnbindExecution(first.gen)
136 if !runtime.Cancel() {
137 t.Fatal("new generation lost cancel after old unbind")
138 }
139 if got := runtime.StateSnapshot().Phase; got != RuntimeCancelling {
140 t.Fatalf("phase = %s, want cancelling", got)
141 }
142 runtime.NoteExecution(first.gen, RuntimeIdle, "")
143 if got := runtime.StateSnapshot().Phase; got != RuntimeCancelling {
144 t.Fatalf("old generation cleared new phase: %s", got)
145 }
146 second.Finish()
147 }
148
149 func TestSessionAcceptsCancelHistoryWhileCancelling(t *testing.T) {
150 _, runtime := reviewRuntime(t)
151 _, exec := bindTestExecution(t, runtime, "turn")
152 if !runtime.Cancel() {
153 t.Fatal("cancel")
154 }
155 payload := []byte(`{"messages":[],"reason":"cancel-or-recovery-rewrite"}`)
156 if _, err := runtime.Session().Append(t.Context(), Batch{
157 OperationID: "history-replace",
158 TurnID: "turn-1",
159 Events: []Event{{Kind: "history/replace", Payload: payload}},
160 }); err != nil {
161 t.Fatalf("history/replace during cancel: %v", err)
162 }
163 if _, err := runtime.Session().Append(t.Context(), Batch{
164 OperationID: "turn-end",
165 TurnID: "turn-1",
166 Events: []Event{{Kind: "turn/end", Payload: []byte(`{"status":"interrupted"}`)}},
167 }); err != nil {
168 t.Fatalf("turn/end during cancel: %v", err)
169 }
170 exec.Finish()
171 if got := runtime.StateSnapshot().Phase; got != RuntimeIdle {
172 t.Fatalf("phase after finish = %s", got)
173 }
174 }
175
176 func TestOldControllerCannotIdleNewGeneration(t *testing.T) {
177 _, runtime := reviewRuntime(t)
178 old := &testExecution{phase: RuntimeRunning, runtime: runtime}
179 old.gen = runtime.BindExecution(old)
180 next := &testExecution{phase: RuntimeIdle, runtime: runtime}
181 next.gen = runtime.ReplaceExecution(old.gen, next)
182 if next.gen == 0 {
183 t.Fatal("replace idle execution")
184 }
185 runtime.NoteExecution(next.gen, RuntimeRunning, "new")
186 runtime.NoteExecution(old.gen, RuntimeIdle, "")
187 if got := runtime.StateSnapshot().Phase; got != RuntimeRunning {
188 t.Fatalf("old finish cleared new turn: %s", got)
189 }
190 // next was built idle, so Finish is a no-op for it; idle the generation the
191 // runtime actually observes or the runtime stays busy and cannot close.
192 runtime.NoteExecution(next.gen, RuntimeIdle, "")
193 }
194
195 func TestReplaceExecutionRejectsBusyOwner(t *testing.T) {
196 _, runtime := reviewRuntime(t)
197 old := &testExecution{phase: RuntimeRunning, runtime: runtime}
198 old.gen = runtime.BindExecution(old)
199 runtime.NoteExecution(old.gen, RuntimeRunning, "old")
200 if gen := runtime.ReplaceExecution(old.gen, &testExecution{}); gen != 0 {
201 t.Fatalf("busy replacement generation = %d, want rejection", gen)
202 }
203 if got := runtime.StateSnapshot().Phase; got != RuntimeRunning {
204 t.Fatalf("phase after rejected replacement = %s, want running", got)
205 }
206 old.Finish()
207 }
208
209 func TestStaleExecutionGenerationCannotCommitPreparedBatch(t *testing.T) {
210 _, runtime := reviewRuntime(t)
211 old := &testExecution{phase: RuntimeIdle, runtime: runtime}
212 old.gen = runtime.BindExecution(old)
213 prepared, err := runtime.Session().PrepareBatchContext(t.Context(), "old-config", Batch{
214 Events: []Event{{Kind: "session/config", Payload: []byte(`{"modelRef":"old"}`)}},
215 })
216 if err != nil {
217 t.Fatal(err)
218 }
219 next := &testExecution{phase: RuntimeIdle, runtime: runtime}
220 next.gen = runtime.ReplaceExecution(old.gen, next)
221 if next.gen == 0 {
222 t.Fatal("replace idle execution")
223 }
224 if _, err := runtime.CommitPreparedForExecution(old.gen, prepared); !errors.Is(err, ErrStaleExecution) {
225 t.Fatalf("stale commit error = %v, want %v", err, ErrStaleExecution)
226 }
227 if got := runtime.Session().ExecutionSnapshot().EventSequence; got != 0 {
228 t.Fatalf("stale generation committed sequence %d", got)
229 }
230 current, err := runtime.Session().PrepareBatchContext(t.Context(), "new-config", Batch{
231 Events: []Event{{Kind: "session/config", Payload: []byte(`{"modelRef":"new"}`)}},
232 })
233 if err != nil {
234 t.Fatal(err)
235 }
236 if _, err := runtime.CommitPreparedForExecution(next.gen, current); err != nil {
237 t.Fatalf("current generation commit: %v", err)
238 }
239 }
240
241 func TestCancelRetriesAcrossExecutionCutover(t *testing.T) {
242 _, runtime := reviewRuntime(t)
243 entered := make(chan struct{}, 1)
244 release := make(chan struct{})
245 old := &blockingRejectExecution{entered: entered, release: release}
246 oldGeneration := runtime.BindExecution(old)
247 if oldGeneration == 0 {
248 t.Fatal("bind outgoing execution")
249 }
250 result := make(chan bool, 1)
251 go func() { result <- runtime.Cancel() }()
252 <-entered
253 next := &testExecution{phase: RuntimeRunning, runtime: runtime}
254 next.gen = runtime.ReplaceExecution(oldGeneration, next)
255 if next.gen == 0 {
256 t.Fatal("replace execution while cancel is in flight")
257 }
258 close(release)
259 if accepted := <-result; !accepted {
260 t.Fatal("cancel was lost across execution cutover")
261 }
262 if got := next.phase; got != RuntimeCancelling {
263 t.Fatalf("replacement phase = %s, want cancelling", got)
264 }
265 }
266
267 func TestRuntimeFinalizingIsBusy(t *testing.T) {
268 if !RuntimeFinalizing.Busy() {
269 t.Fatal("finalizing phase must remain busy until terminal commit finishes")
270 }
271 }
272
272 lines GO