返回 DeepSeek-Reasonix
recorder_test.go
根目录 / internal / taskmonitor / recorder_test.go
1 package taskmonitor
2
3 import (
4 "context"
5 "fmt"
6 "sync"
7 "testing"
8 "time"
9
10 "reasonix/internal/jobs"
11 )
12
13 func newRecorderForTest(t *testing.T, projectDir string) (*TaskRecorder, *FileStore) {
14 t.Helper()
15 store := NewFileStore(".reasonix/tasks")
16 r := NewTaskRecorder(store, projectDir, func() string { return "sess-1" })
17 return r, store
18 }
19
20 func TestTaskRecorder_Lifecycle(t *testing.T) {
21 dir := t.TempDir()
22 r, store := newRecorderForTest(t, dir)
23 ctx := context.Background()
24
25 r.RecordStart("task-1", "task", "demo")
26 snap, err := store.GetTask(ctx, dir, monitorTaskID("sess-1", "task-1"))
27 if err != nil || snap == nil {
28 t.Fatalf("GetTask after start: %+v, %v", snap, err)
29 }
30 if snap.State != TaskStateRunning || snap.RuntimeState != RuntimeStateAlive || snap.Version != 1 || snap.SessionID != "sess-1" {
31 t.Fatalf("snapshot after start = %+v", snap)
32 }
33 if snap.JobID != "task-1" {
34 t.Fatalf("snapshot job id = %q, want task-1", snap.JobID)
35 }
36
37 r.RecordDone("task-1", jobs.Done, nil)
38 snap, _ = store.GetTask(ctx, dir, monitorTaskID("sess-1", "task-1"))
39 if snap.State != TaskStateSucceeded || snap.RuntimeState != RuntimeStateExited || snap.Version != 2 {
40 t.Fatalf("snapshot after done = %+v", snap)
41 }
42
43 events, err := store.ListEvents(ctx, dir, monitorTaskID("sess-1", "task-1"), 0)
44 if err != nil || len(events) != 2 {
45 t.Fatalf("events = %+v, %v", events, err)
46 }
47 if events[0].EventType != "state_change" || events[0].State != TaskStateRunning || events[0].RuntimeState != RuntimeStateAlive || events[0].Sequence != 1 {
48 t.Fatalf("event[0] = %+v", events[0])
49 }
50 if events[1].State != TaskStateSucceeded || events[1].RuntimeState != RuntimeStateExited || events[1].Sequence != 2 {
51 t.Fatalf("event[1] = %+v", events[1])
52 }
53 }
54
55 func TestTaskRecorder_HeartbeatRenewsExpiredOwnedLease(t *testing.T) {
56 dir := t.TempDir()
57 r, store := newRecorderForTest(t, dir)
58 ctx := context.Background()
59 monitorID := monitorTaskID("sess-1", "task-1")
60
61 r.RecordStart("task-1", "task", "demo")
62 // Drive the renewal deterministically instead of waiting for the ticker.
63 r.stopHeartbeat(monitorID)
64 raw, err := store.getTaskRaw(ctx, dir, monitorID)
65 if err != nil || raw == nil {
66 t.Fatalf("raw task after start: %+v, %v", raw, err)
67 }
68 raw.Version++
69 raw.RuntimeLeaseUntil = time.Now().Add(-time.Minute)
70 if err := store.SaveTask(ctx, dir, *raw); err != nil {
71 t.Fatal(err)
72 }
73
74 observed, err := store.GetTask(ctx, dir, monitorID)
75 if err != nil || observed == nil || observed.State != TaskStateStale || observed.RuntimeState != RuntimeStateExited {
76 t.Fatalf("expired observed task = %+v, err=%v", observed, err)
77 }
78 if !r.renewHeartbeat(ctx, monitorID) {
79 t.Fatal("live owner failed to renew its expired persisted lease")
80 }
81 observed, err = store.GetTask(ctx, dir, monitorID)
82 if err != nil || observed == nil || observed.State != TaskStateRunning || observed.RuntimeState != RuntimeStateAlive {
83 t.Fatalf("renewed observed task = %+v, err=%v", observed, err)
84 }
85 if !observed.RuntimeLeaseUntil.After(time.Now()) {
86 t.Fatalf("renewed lease = %v, want future deadline", observed.RuntimeLeaseUntil)
87 }
88
89 r.RecordDone("task-1", jobs.Done, nil)
90 }
91
92 func TestTaskRecorder_OldOwnerCannotRenewNewRuntimeGeneration(t *testing.T) {
93 dir := t.TempDir()
94 r, store := newRecorderForTest(t, dir)
95 ctx := context.Background()
96 monitorID := monitorTaskID("sess-1", "task-1")
97
98 r.RecordStart("task-1", "task", "demo")
99 r.stopHeartbeat(monitorID)
100 raw, err := store.getTaskRaw(ctx, dir, monitorID)
101 if err != nil || raw == nil {
102 t.Fatalf("raw task after start: %+v, %v", raw, err)
103 }
104 raw.Version++
105 raw.RuntimeOwnerID = "new-runtime-owner"
106 raw.RuntimeLeaseUntil = time.Now().Add(-time.Minute)
107 if err := store.SaveTask(ctx, dir, *raw); err != nil {
108 t.Fatal(err)
109 }
110 if r.renewHeartbeat(ctx, monitorID) {
111 t.Fatal("older recorder renewed a newer runtime generation")
112 }
113 after, err := store.getTaskRaw(ctx, dir, monitorID)
114 if err != nil || after == nil {
115 t.Fatalf("raw task after rejected renewal: %+v, %v", after, err)
116 }
117 if after.RuntimeOwnerID != "new-runtime-owner" || !after.RuntimeLeaseUntil.Equal(raw.RuntimeLeaseUntil) {
118 t.Fatalf("rejected renewal mutated newer runtime: %+v", after)
119 }
120 }
121
122 func TestTaskRecorder_FailedUsesContentFreeErrorCode(t *testing.T) {
123 dir := t.TempDir()
124 r, store := newRecorderForTest(t, dir)
125 ctx := context.Background()
126
127 r.RecordStart("bash-1", "bash", "")
128 r.RecordDone("bash-1", jobs.Failed, fmt.Errorf(`command "deploy --token secret" failed in /Users/alice/private`))
129 snap, _ := store.GetTask(ctx, dir, monitorTaskID("sess-1", "bash-1"))
130 if snap.State != TaskStateFailed || snap.ErrorCode != "job_failed" || snap.ErrorSummary != "" {
131 t.Fatalf("snapshot = %+v", snap)
132 }
133 events, err := store.ListEvents(ctx, dir, monitorTaskID("sess-1", "bash-1"), 0)
134 if err != nil || len(events) != 2 {
135 t.Fatalf("events = %+v, err=%v", events, err)
136 }
137 if events[1].ErrorCode != "job_failed" || events[1].ErrorSummary != "" {
138 t.Fatalf("terminal event exposed error content: %+v", events[1])
139 }
140 }
141
142 func TestTaskRecorder_KilledAndInterruptedMapToCancelled(t *testing.T) {
143 dir := t.TempDir()
144 r, store := newRecorderForTest(t, dir)
145 ctx := context.Background()
146
147 r.RecordStart("t1", "task", "")
148 r.RecordDone("t1", jobs.Killed, nil)
149 snap, _ := store.GetTask(ctx, dir, monitorTaskID("sess-1", "t1"))
150 if snap.State != TaskStateCancelled {
151 t.Fatalf("killed -> %v, want cancelled", snap.State)
152 }
153
154 r.RecordStart("t2", "task", "")
155 r.RecordDone("t2", jobs.Interrupted, nil)
156 snap, _ = store.GetTask(ctx, dir, monitorTaskID("sess-1", "t2"))
157 if snap.State != TaskStateCancelled {
158 t.Fatalf("interrupted -> %v, want cancelled", snap.State)
159 }
160 }
161
162 func TestTaskRecorder_RestartUsesDistinctMonitorID(t *testing.T) {
163 dir := t.TempDir()
164 r, store := newRecorderForTest(t, dir)
165 ctx := context.Background()
166
167 // First lifecycle.
168 r.RecordStart("task-1", "task", "")
169 r.RecordDone("task-1", jobs.Done, nil)
170 first, _ := store.GetTask(ctx, dir, monitorTaskID("sess-1", "task-1"))
171
172 // A new session with the same local job ID gets a distinct monitor key.
173 r2 := NewTaskRecorder(store, dir, func() string { return "sess-2" })
174 r2.RecordStart("task-1", "task", "")
175 second, _ := store.GetTask(ctx, dir, monitorTaskID("sess-2", "task-1"))
176 if second == nil || second.Version != 1 {
177 t.Fatalf("second snapshot = %+v, want a new version-1 lifecycle", second)
178 }
179 if second.CreatedAt.Equal(first.CreatedAt) {
180 t.Fatalf("second lifecycle reused creation time: %v", second.CreatedAt)
181 }
182 if second.State != TaskStateRunning || second.RuntimeState != RuntimeStateAlive {
183 t.Fatalf("state = %v, want running", second.State)
184 }
185 }
186
187 func TestTaskRecorder_NonTerminalStatusDoesNotUpdate(t *testing.T) {
188 dir := t.TempDir()
189 r, store := newRecorderForTest(t, dir)
190 ctx := context.Background()
191
192 r.RecordStart("t1", "bash", "")
193 r.RecordDone("t1", jobs.Running, nil) // never happens in practice; guard anyway
194 snap, _ := store.GetTask(ctx, dir, monitorTaskID("sess-1", "t1"))
195 if snap.State != TaskStateRunning || snap.Version != 1 {
196 t.Fatalf("snapshot = %+v", snap)
197 }
198 }
199
200 type blockingExitedSaveStore struct {
201 WriteStore
202 once sync.Once
203 blocked chan struct{}
204 release chan struct{}
205 }
206
207 func (s *blockingExitedSaveStore) SaveTask(ctx context.Context, projectDir string, snap TaskSnapshot) error {
208 if snap.RuntimeState == RuntimeStateExited {
209 s.once.Do(func() {
210 close(s.blocked)
211 <-s.release
212 })
213 }
214 return s.WriteStore.SaveTask(ctx, projectDir, snap)
215 }
216
217 func TestTaskRecorder_DoneRetriesAfterConcurrentControlUpdate(t *testing.T) {
218 base := NewInMemoryStore()
219 now := time.Now()
220 if err := base.UpsertTask("/p", TaskSnapshot{
221 SchemaVersion: 1, TaskID: monitorTaskID("session-1", "task-1"), SessionID: "session-1",
222 State: TaskStateRunning, RuntimeState: RuntimeStateAlive, Version: 1,
223 CreatedAt: now, UpdatedAt: now,
224 }); err != nil {
225 t.Fatal(err)
226 }
227 store := &blockingExitedSaveStore{
228 WriteStore: base,
229 blocked: make(chan struct{}),
230 release: make(chan struct{}),
231 }
232 recorder := NewTaskRecorder(store, "/p", func() string { return "session-1" })
233 recorder.rememberMonitorID("task-1", monitorTaskID("session-1", "task-1"))
234 done := make(chan struct{})
235 go func() {
236 recorder.RecordDone("task-1", jobs.Killed, nil)
237 close(done)
238 }()
239 <-store.blocked
240
241 control := NewControlService(base)
242 res, err := control.StopTaskWithKiller(context.Background(), "/p", monitorTaskID("session-1", "task-1"), 1, "", "", &mockKiller{fn: func(string, string) bool { return true }})
243 if err != nil || !res.Accepted {
244 t.Fatalf("concurrent stop: result=%+v err=%v", res, err)
245 }
246 close(store.release)
247 <-done
248
249 snap, err := base.GetTask(context.Background(), "/p", monitorTaskID("session-1", "task-1"))
250 if err != nil || snap == nil {
251 t.Fatalf("GetTask: snap=%+v err=%v", snap, err)
252 }
253 if snap.State != TaskStateCancelled || snap.RuntimeState != RuntimeStateExited || snap.Version != 3 {
254 t.Fatalf("completion evidence was lost after CAS retry: %+v", snap)
255 }
256 }
257
258 func TestControlAcceptsRecorderCompletionThatWinsPostKillCAS(t *testing.T) {
259 store := NewInMemoryStore()
260 now := time.Now()
261 monitorID := monitorTaskID("session-1", "task-1")
262 if err := store.UpsertTask("/p", TaskSnapshot{
263 SchemaVersion: 1, TaskID: monitorID, JobID: "task-1", SessionID: "session-1",
264 State: TaskStateRunning, RuntimeState: RuntimeStateAlive, Version: 1,
265 CreatedAt: now, UpdatedAt: now,
266 }); err != nil {
267 t.Fatal(err)
268 }
269 recorder := NewTaskRecorder(store, "/p", func() string { return "session-1" })
270 recorder.rememberMonitorID("task-1", monitorID)
271 killer := &mockKiller{fn: func(sessionID, jobID string) bool {
272 if sessionID != "session-1" || jobID != "task-1" {
273 t.Fatalf("runtime route = %q/%q, want session-1/task-1", sessionID, jobID)
274 }
275 recorder.RecordDone("task-1", jobs.Killed, nil)
276 return true
277 }}
278
279 res, err := NewControlService(store).StopTaskWithKiller(context.Background(), "/p", monitorID, 1, "", "stop-once", killer)
280 if err != nil || !res.Accepted {
281 t.Fatalf("stop after recorder completion: result=%+v err=%v", res, err)
282 }
283 if res.State != TaskStateCancelled || res.RuntimeState != RuntimeStateExited || res.Version != 2 {
284 t.Fatalf("stop result lost recorder completion: %+v", res)
285 }
286 idem, err := store.CheckIdempotency(context.Background(), "/p", "stop-once")
287 if err != nil || idem == nil || idem.Pending {
288 t.Fatalf("idempotency result = %+v, err=%v; want finalized", idem, err)
289 }
290 }
291
292 func TestTaskRecorder_UnknownTaskDoneIsNoop(t *testing.T) {
293 dir := t.TempDir()
294 r, store := newRecorderForTest(t, dir)
295 ctx := context.Background()
296
297 r.RecordDone("never-started", jobs.Done, nil) // must not panic or write anything
298 tasks, err := store.ListTasks(ctx, dir)
299 if err != nil || len(tasks) != 0 {
300 t.Fatalf("tasks = %+v, %v", tasks, err)
301 }
302 }
303
304 func TestTaskRecorder_EmptySessionIDAllowed(t *testing.T) {
305 dir := t.TempDir()
306 store := NewFileStore(".reasonix/tasks")
307 r := NewTaskRecorder(store, dir, func() string { return "" })
308 ctx := context.Background()
309
310 r.RecordStart("t1", "bash", "")
311 tasks, err := store.ListTasks(ctx, dir)
312 if err != nil || len(tasks) != 1 {
313 t.Fatalf("ListTasks: %+v, %v", tasks, err)
314 }
315 snap := tasks[0]
316 if snap.TaskID == "t1" || snap.SessionID != "" {
317 t.Fatalf("snapshot identity = %+v, want unique sessionless ID", snap)
318 }
319 events, err := store.ListEvents(ctx, dir, snap.TaskID, 0)
320 if err != nil || len(events) != 1 {
321 t.Fatalf("events: %+v, %v", events, err)
322 }
323 }
324
325 func TestTaskRecorder_SameJobIDAcrossSessionsUsesDistinctMonitorIDs(t *testing.T) {
326 store := NewFileStore(".reasonix/tasks")
327 projectDir := t.TempDir()
328 r1 := NewTaskRecorder(store, projectDir, func() string { return "session-a" })
329 r2 := NewTaskRecorder(store, projectDir, func() string { return "session-b" })
330
331 r1.RecordStart("task-1", "task", "first")
332 r2.RecordStart("task-1", "task", "second")
333 r1.RecordDone("task-1", jobs.Done, nil)
334 r2.RecordDone("task-1", jobs.Failed, context.DeadlineExceeded)
335
336 tasks, err := store.ListTasks(context.Background(), projectDir)
337 if err != nil {
338 t.Fatal(err)
339 }
340 if len(tasks) != 2 {
341 t.Fatalf("tasks = %+v, want two independent lifecycles", tasks)
342 }
343 seen := map[string]TaskSnapshot{}
344 for _, task := range tasks {
345 seen[task.TaskID] = task
346 }
347 first, ok := seen["session-a--task-1"]
348 if !ok || first.State != TaskStateSucceeded || first.SessionID != "session-a" {
349 t.Fatalf("session-a task = %+v", first)
350 }
351 second, ok := seen["session-b--task-1"]
352 if !ok || second.State != TaskStateFailed || second.SessionID != "session-b" {
353 t.Fatalf("session-b task = %+v", second)
354 }
355 }
356
356 lines GO