返回 DeepSeek-Reasonix
jsonstore_test.go
根目录 / internal / taskmonitor / jsonstore_test.go
1 package taskmonitor
2
3 import (
4 "context"
5 "encoding/json"
6 "errors"
7 "os"
8 "path/filepath"
9 "runtime"
10 "strings"
11 "sync"
12 "testing"
13 "time"
14 )
15
16 func TestFileStore_ListTasks_EmptyDir(t *testing.T) {
17 dir := t.TempDir()
18 store := NewFileStore(".reasonix/tasks")
19
20 tasks, err := store.ListTasks(context.Background(), dir)
21 if err != nil {
22 t.Fatalf("ListTasks: %v", err)
23 }
24 if len(tasks) != 0 {
25 t.Errorf("expected empty, got %d", len(tasks))
26 }
27 }
28
29 func TestFileStore_RoundTrip(t *testing.T) {
30 dir := t.TempDir()
31 store := NewFileStore(".reasonix/tasks")
32
33 // Write a task snapshot
34 taskDir := filepath.Join(dir, ".reasonix", "tasks", "task-1")
35 if err := os.MkdirAll(taskDir, 0o755); err != nil {
36 t.Fatal(err)
37 }
38 now := time.Now().Truncate(time.Second)
39 snap := TaskSnapshot{
40 SchemaVersion: 1, TaskID: "task-1", SessionID: "s1",
41 State: TaskStateFailed, CreatedAt: now.Add(-time.Hour), UpdatedAt: now,
42 ErrorCode: "TIMEOUT", ErrorSummary: "deadline exceeded",
43 }
44 data, _ := json.Marshal(snap)
45 if err := os.WriteFile(filepath.Join(taskDir, "snapshot.json"), data, 0o644); err != nil {
46 t.Fatal(err)
47 }
48
49 // Read back
50 got, err := store.GetTask(context.Background(), dir, "task-1")
51 if err != nil {
52 t.Fatalf("GetTask: %v", err)
53 }
54 if got == nil {
55 t.Fatal("expected snapshot, got nil")
56 }
57 if got.TaskID != "task-1" || got.ErrorCode != "TIMEOUT" {
58 t.Errorf("mismatch: %+v", got)
59 }
60 }
61
62 func TestFileStore_ListEvents_RoundTrip(t *testing.T) {
63 dir := t.TempDir()
64 store := NewFileStore(".reasonix/tasks")
65
66 taskDir := filepath.Join(dir, ".reasonix", "tasks", "t1")
67 if err := os.MkdirAll(taskDir, 0o755); err != nil {
68 t.Fatal(err)
69 }
70 // Write events as JSONL
71 events := `{"sequence":1,"timestamp":"2025-01-01T00:00:01Z","event_type":"state_change","task_id":"t1","session_id":"s","state":"queued"}
72 {"sequence":2,"timestamp":"2025-01-01T00:00:02Z","event_type":"state_change","task_id":"t1","session_id":"s","state":"running"}
73 {"sequence":3,"timestamp":"2025-01-01T00:00:03Z","event_type":"error","task_id":"t1","session_id":"s","state":"failed","error_code":"E1"}
74 `
75 if err := os.WriteFile(filepath.Join(taskDir, "events.jsonl"), []byte(events), 0o644); err != nil {
76 t.Fatal(err)
77 }
78
79 got, err := store.ListEvents(context.Background(), dir, "t1", 0)
80 if err != nil {
81 t.Fatalf("ListEvents: %v", err)
82 }
83 if len(got) != 3 {
84 t.Fatalf("expected 3 events, got %d", len(got))
85 }
86 if got[2].ErrorCode != "E1" {
87 t.Errorf("expected E1, got %q", got[2].ErrorCode)
88 }
89 }
90
91 func TestFileStore_ListEvents_AfterCursor(t *testing.T) {
92 dir := t.TempDir()
93 store := NewFileStore(".reasonix/tasks")
94 taskDir := filepath.Join(dir, ".reasonix", "tasks", "t1")
95 os.MkdirAll(taskDir, 0o755)
96 events := `{"sequence":1,"timestamp":"2025-01-01T00:00:01Z","event_type":"e","task_id":"t1","session_id":"s","state":"queued"}
97 {"sequence":2,"timestamp":"2025-01-01T00:00:02Z","event_type":"e","task_id":"t1","session_id":"s","state":"running"}
98 `
99 os.WriteFile(filepath.Join(taskDir, "events.jsonl"), []byte(events), 0o644)
100
101 got, _ := store.ListEvents(context.Background(), dir, "t1", 1)
102 if len(got) != 1 || got[0].Sequence != 2 {
103 t.Errorf("expected [seq=2], got %d events, seq=%d", len(got), got[0].Sequence)
104 }
105 }
106
107 func TestFileStore_RejectsPathTraversal_TaskID(t *testing.T) {
108 dir := t.TempDir()
109 store := NewFileStore(".reasonix/tasks")
110
111 _, err := store.GetTask(context.Background(), dir, "../escape")
112 if err == nil || !strings.Contains(err.Error(), "path separator") {
113 t.Fatalf("expected path traversal rejection, got %v", err)
114 }
115 }
116
117 func TestFileStore_AcceptsCleanableProjectDir(t *testing.T) {
118 parent := t.TempDir()
119 store := NewFileStore(".reasonix/tasks")
120 now := time.Now()
121
122 for _, projectDir := range []string{
123 filepath.Join(parent, "nested", "..", "project"),
124 filepath.Join(parent, "project..archive"),
125 } {
126 if err := os.MkdirAll(filepath.Clean(projectDir), 0o755); err != nil {
127 t.Fatal(err)
128 }
129 snap := TaskSnapshot{
130 SchemaVersion: 1, TaskID: "task-1", SessionID: "session-1",
131 State: TaskStateRunning, Version: 1, CreatedAt: now, UpdatedAt: now,
132 }
133 if err := store.SaveTask(context.Background(), projectDir, snap); err != nil {
134 t.Fatalf("SaveTask(%q): %v", projectDir, err)
135 }
136 got, err := store.GetTask(context.Background(), projectDir, snap.TaskID)
137 if err != nil || got == nil || got.TaskID != snap.TaskID {
138 t.Fatalf("GetTask(%q) = %+v, %v", projectDir, got, err)
139 }
140 }
141 }
142
143 func TestFileStore_RejectsEmptyTaskID(t *testing.T) {
144 dir := t.TempDir()
145 store := NewFileStore(".reasonix/tasks")
146
147 _, err := store.GetTask(context.Background(), dir, "")
148 if err == nil || !strings.Contains(err.Error(), "must not be empty") {
149 t.Fatalf("expected empty rejection, got %v", err)
150 }
151 }
152
153 func TestFileStore_RejectsDotTaskID(t *testing.T) {
154 dir := t.TempDir()
155 store := NewFileStore(".reasonix/tasks")
156
157 _, err := store.GetTask(context.Background(), dir, ".")
158 if err == nil || !strings.Contains(err.Error(), "invalid") {
159 t.Fatalf("expected rejection for '.', got %v", err)
160 }
161 }
162
163 func TestFileStore_RejectsDotDotTaskID(t *testing.T) {
164 dir := t.TempDir()
165 store := NewFileStore(".reasonix/tasks")
166
167 _, err := store.GetTask(context.Background(), dir, "..")
168 if err == nil || !strings.Contains(err.Error(), "invalid") {
169 t.Fatalf("expected rejection for '..', got %v", err)
170 }
171 }
172
173 func TestFileStore_RejectsSymlinkTaskDirectory(t *testing.T) {
174 project := t.TempDir()
175 outside := t.TempDir()
176 root := filepath.Join(project, ".reasonix", "tasks")
177 if err := os.MkdirAll(root, 0o700); err != nil {
178 t.Fatal(err)
179 }
180 if err := os.Symlink(outside, filepath.Join(root, "evil")); err != nil {
181 t.Skipf("symlink unavailable: %v", err)
182 }
183 snap := TaskSnapshot{SchemaVersion: 1, TaskID: "evil", SessionID: "s", Version: 1, State: TaskStateRunning, CreatedAt: time.Now(), UpdatedAt: time.Now()}
184 if err := NewFileStore(".reasonix/tasks").SaveTask(context.Background(), project, snap); err == nil {
185 t.Fatal("expected symlink task directory to be rejected")
186 }
187 if _, err := os.Stat(filepath.Join(outside, "snapshot.json")); !os.IsNotExist(err) {
188 t.Fatalf("write escaped through symlink: stat err=%v", err)
189 }
190 }
191
192 func TestFileStore_RejectsSymlinkStoreParent(t *testing.T) {
193 project := t.TempDir()
194 outside := t.TempDir()
195 if err := os.Symlink(outside, filepath.Join(project, ".reasonix")); err != nil {
196 t.Skipf("symlink unavailable: %v", err)
197 }
198 snap := TaskSnapshot{SchemaVersion: 1, TaskID: "t1", SessionID: "s", Version: 1, State: TaskStateRunning, CreatedAt: time.Now(), UpdatedAt: time.Now()}
199 if err := NewFileStore(".reasonix/tasks").SaveTask(context.Background(), project, snap); err == nil {
200 t.Fatal("expected symlink store parent to be rejected")
201 }
202 if _, err := os.Stat(filepath.Join(outside, "tasks", "t1", "snapshot.json")); !os.IsNotExist(err) {
203 t.Fatalf("write escaped through parent symlink: stat err=%v", err)
204 }
205 }
206
207 func TestFileStore_DefaultProjectRejectsSymlinkStoreParent(t *testing.T) {
208 project := t.TempDir()
209 outside := t.TempDir()
210 oldWorkingDir, err := os.Getwd()
211 if err != nil {
212 t.Fatal(err)
213 }
214 if err := os.Chdir(project); err != nil {
215 t.Fatal(err)
216 }
217 t.Cleanup(func() {
218 if err := os.Chdir(oldWorkingDir); err != nil {
219 t.Errorf("restore working directory: %v", err)
220 }
221 })
222 if err := os.Symlink(outside, ".reasonix"); err != nil {
223 t.Skipf("symlink unavailable: %v", err)
224 }
225
226 snap := TaskSnapshot{SchemaVersion: 1, TaskID: "t1", SessionID: "s", Version: 1, State: TaskStateRunning, CreatedAt: time.Now(), UpdatedAt: time.Now()}
227 if err := NewFileStore(".reasonix/tasks").SaveTask(context.Background(), "", snap); err == nil {
228 t.Fatal("expected default project scope to reject symlink store parent")
229 }
230 if _, err := os.Stat(filepath.Join(outside, "tasks", "t1", "snapshot.json")); !os.IsNotExist(err) {
231 t.Fatalf("default-scope write escaped through parent symlink: stat err=%v", err)
232 }
233 }
234
235 func TestFileStore_RejectsSymlinkSnapshotAndEvents(t *testing.T) {
236 project := t.TempDir()
237 outside := t.TempDir()
238 root := filepath.Join(project, ".reasonix", "tasks", "t1")
239 if err := os.MkdirAll(root, 0o700); err != nil {
240 t.Fatal(err)
241 }
242 for _, name := range []string{"snapshot.json", "events.jsonl"} {
243 if err := os.Symlink(filepath.Join(outside, name), filepath.Join(root, name)); err != nil {
244 t.Skipf("symlink unavailable: %v", err)
245 }
246 }
247 store := NewFileStore(".reasonix/tasks")
248 if _, err := store.GetTask(context.Background(), project, "t1"); err == nil {
249 t.Fatal("expected snapshot symlink to be rejected")
250 }
251 if _, err := store.ListEvents(context.Background(), project, "t1", 0); err == nil {
252 t.Fatal("expected events symlink to be rejected")
253 }
254 }
255
256 func TestFileStore_WritablePathsUsePrivateModes(t *testing.T) {
257 project := t.TempDir()
258 store := NewFileStore(".reasonix/tasks")
259 now := time.Now()
260 snap := TaskSnapshot{SchemaVersion: 1, TaskID: "t1", SessionID: "s", Version: 1, State: TaskStateRunning, CreatedAt: now, UpdatedAt: now}
261 if err := store.SaveTask(context.Background(), project, snap); err != nil {
262 t.Fatal(err)
263 }
264 if err := store.AppendAuditEvent(context.Background(), project, TaskEvent{TaskID: "t1", SessionID: "s", EventType: "state_change", State: TaskStateRunning, Timestamp: now}); err != nil {
265 t.Fatal(err)
266 }
267 if runtime.GOOS == "windows" {
268 return // Windows does not expose POSIX permission bits through os.FileMode.
269 }
270 checks := map[string]os.FileMode{
271 filepath.Join(project, ".reasonix", "tasks"): 0o700,
272 filepath.Join(project, ".reasonix", "tasks", "t1"): 0o700,
273 filepath.Join(project, ".reasonix", "tasks", "t1", "snapshot.json"): 0o600,
274 filepath.Join(project, ".reasonix", "tasks", "t1", "events.jsonl"): 0o600,
275 filepath.Join(project, ".reasonix", "tasks", "t1", "task.lock"): 0o600,
276 }
277 for path, want := range checks {
278 info, err := os.Stat(path)
279 if err != nil {
280 t.Fatalf("stat %s: %v", path, err)
281 }
282 if got := info.Mode().Perm(); got != want {
283 t.Errorf("mode %s = %o, want %o", path, got, want)
284 }
285 }
286 }
287
288 func TestFileStore_SaveTask_VersionConflict(t *testing.T) {
289 dir := t.TempDir()
290 store := NewFileStore(".reasonix/tasks")
291 ctx := context.Background()
292 now := time.Now().Truncate(time.Second)
293 v1 := TaskSnapshot{
294 SchemaVersion: 1, TaskID: "t1", SessionID: "s1",
295 State: TaskStateRunning, Version: 1, CreatedAt: now, UpdatedAt: now,
296 }
297 if err := store.SaveTask(ctx, dir, v1); err != nil {
298 t.Fatalf("SaveTask v1: %v", err)
299 }
300 // Same version must conflict.
301 if err := store.SaveTask(ctx, dir, v1); err == nil || !errors.Is(err, ErrStoreVersionConflict) {
302 t.Fatalf("expected version conflict, got %v", err)
303 }
304 // Higher version wins.
305 v2 := v1
306 v2.Version = 2
307 v2.State = TaskStateSucceeded
308 if err := store.SaveTask(ctx, dir, v2); err != nil {
309 t.Fatalf("SaveTask v2: %v", err)
310 }
311 got, err := store.GetTask(ctx, dir, "t1")
312 if err != nil || got == nil || got.Version != 2 {
313 t.Fatalf("read back: %+v, %v", got, err)
314 }
315 }
316
317 // TestFileStore_SaveTask_ConcurrentCAS races two independent FileStore
318 // instances (production: CLI and Desktop processes) advancing the same task
319 // from version 1 to version 2. The per-task lock must guarantee exactly one
320 // winner; the loser observes the version conflict instead of silently
321 // overwriting (the pre-fix TOCTOU).
322 func TestFileStore_SaveTask_ConcurrentCAS(t *testing.T) {
323 dir := t.TempDir()
324 ctx := context.Background()
325 now := time.Now().Truncate(time.Second)
326 seed := NewFileStore(".reasonix/tasks")
327 v1 := TaskSnapshot{
328 SchemaVersion: 1, TaskID: "t1", SessionID: "s1",
329 State: TaskStateRunning, Version: 1, CreatedAt: now, UpdatedAt: now,
330 }
331 if err := seed.SaveTask(ctx, dir, v1); err != nil {
332 t.Fatalf("seed: %v", err)
333 }
334
335 write := func(v uint64) error {
336 snap := v1
337 snap.Version = v
338 snap.State = TaskStateSucceeded
339 return NewFileStore(".reasonix/tasks").SaveTask(ctx, dir, snap)
340 }
341
342 start := make(chan struct{})
343 errs := make([]error, 2)
344 var wg sync.WaitGroup
345 for i := range 2 {
346 wg.Add(1)
347 go func(i int) {
348 defer wg.Done()
349 <-start
350 errs[i] = write(2)
351 }(i)
352 }
353 close(start)
354 wg.Wait()
355
356 ok, conflict := 0, 0
357 for _, err := range errs {
358 switch {
359 case err == nil:
360 ok++
361 case strings.Contains(err.Error(), "version conflict"):
362 conflict++
363 default:
364 t.Fatalf("unexpected error: %v", err)
365 }
366 }
367 if ok != 1 || conflict != 1 {
368 t.Fatalf("want exactly one winner and one conflict, got ok=%d conflict=%d (%v)", ok, conflict, errs)
369 }
370 got, err := seed.GetTask(ctx, dir, "t1")
371 if err != nil || got == nil || got.Version != 2 {
372 t.Fatalf("final state: %+v, %v", got, err)
373 }
374 }
375
376 // TestFileStore_SaveTask_CorruptSnapshotRejected guards the CAS gate: a
377 // corrupt snapshot.json must fail loudly instead of silently bypassing the
378 // version check and being overwritten.
379 func TestFileStore_SaveTask_CorruptSnapshotRejected(t *testing.T) {
380 dir := t.TempDir()
381 store := NewFileStore(".reasonix/tasks")
382 ctx := context.Background()
383 taskDir := filepath.Join(dir, ".reasonix", "tasks", "t1")
384 if err := os.MkdirAll(taskDir, 0o755); err != nil {
385 t.Fatal(err)
386 }
387 if err := os.WriteFile(filepath.Join(taskDir, "snapshot.json"), []byte("{not json"), 0o644); err != nil {
388 t.Fatal(err)
389 }
390 snap := TaskSnapshot{
391 SchemaVersion: 1, TaskID: "t1", SessionID: "s1",
392 State: TaskStateRunning, Version: 1, CreatedAt: time.Now(), UpdatedAt: time.Now(),
393 }
394 err := store.SaveTask(ctx, dir, snap)
395 if err == nil || !strings.Contains(err.Error(), "read current snapshot") {
396 t.Fatalf("expected corrupt-snapshot rejection, got %v", err)
397 }
398 }
399
399 lines GO