返回 DeepSeek-Reasonix
session_display_pager_dag_fence_test.go
根目录 / internal / agent / session_display_pager_dag_fence_test.go
1 package agent
2
3 import (
4 "bytes"
5 "context"
6 "database/sql"
7 "errors"
8 "os"
9 "testing"
10
11 "reasonix/internal/historywork"
12 "reasonix/internal/projectiondb"
13 "reasonix/internal/store"
14 )
15
16 func TestDAGPagerResumeFencesSourceAndBranchChanges(t *testing.T) {
17 for _, mode := range []string{"rewrite", "replace", "checkpoint", "branch", "truncated"} {
18 t.Run(mode, func(t *testing.T) {
19 source, cache := dagResumeFixture(t)
20 opts, fingerprint, size := dagResumeOptions(t, source, cache, "")
21 ctx, cancel := context.WithCancel(t.Context())
22 err := projectiondb.Rebuild(ctx, opts, func(ctx context.Context, db *sql.DB) error {
23 return buildDAGDisplayPagerObserved(ctx, db, source, fingerprint, "", size, func(stage string, count int) {
24 if stage == "scan" && count >= historywork.BatchEntries {
25 cancel()
26 }
27 })
28 })
29 cancel()
30 if !errors.Is(err, context.Canceled) {
31 t.Fatal(err)
32 }
33 head := "fork"
34 requested := ""
35 if mode == "branch" {
36 head, requested = SessionMainHead, SessionMainHead
37 } else {
38 mutateDAGPagerSource(t, source, mode)
39 }
40 pager, err := OpenDisplayPager(t.Context(), source, cache, requested)
41 if mode == "truncated" {
42 if err == nil {
43 pager.Close()
44 t.Fatal("truncated source published as complete")
45 }
46 return
47 }
48 if err != nil {
49 t.Fatal(err)
50 }
51 defer pager.Close()
52 assertDAGPagerReplay(t, pager, source, head)
53 })
54 }
55 }
56
57 func TestDAGPagerRejectsSourceChangeBeforePublication(t *testing.T) {
58 for _, mode := range []string{"rewrite", "replace", "checkpoint"} {
59 t.Run(mode, func(t *testing.T) {
60 source, cache := dagResumeFixture(t)
61 opts, fingerprint, size := dagResumeOptions(t, source, cache, "")
62 changed := false
63 err := projectiondb.Rebuild(t.Context(), opts, func(ctx context.Context, db *sql.DB) error {
64 return buildDAGDisplayPagerObserved(ctx, db, source, fingerprint, "", size, func(stage string, count int) {
65 if !changed && stage == "projection" && count >= historywork.BatchEntries {
66 changed = true
67 mutateDAGPagerSource(t, source, mode)
68 }
69 })
70 })
71 if !changed || !errors.Is(err, ErrDisplaySourceChanged) {
72 t.Fatalf("changed source accepted: %v", err)
73 }
74 if _, err := os.Stat(cache); !os.IsNotExist(err) {
75 t.Fatalf("changed projection published: %v", err)
76 }
77 })
78 }
79 }
80
81 func mutateDAGPagerSource(t *testing.T, source, mode string) {
82 t.Helper()
83 path := store.SessionEventLog(source)
84 if mode == "checkpoint" {
85 path = source
86 }
87 info, err := os.Stat(path)
88 if err != nil {
89 t.Fatal(err)
90 }
91 body, err := os.ReadFile(path)
92 if err != nil {
93 t.Fatal(err)
94 }
95 body = bytes.ReplaceAll(body, []byte("question"), []byte("QUESTION"))
96 switch mode {
97 case "checkpoint":
98 body = []byte("changed checkpoint")
99 case "truncated":
100 body = append(body, []byte(`{"schema_version":2,`)...)
101 }
102 writePath := path
103 if mode == "replace" {
104 writePath += ".replacement"
105 }
106 if err := os.WriteFile(writePath, body, 0600); err != nil {
107 t.Fatal(err)
108 }
109 if mode == "replace" {
110 // Windows rejects direct replacement of an open target even with delete
111 // sharing. Moving the old path first still tests fencing the new identity.
112 if err := os.Rename(path, path+".previous"); err != nil {
113 t.Fatal(err)
114 }
115 if err := os.Rename(writePath, path); err != nil {
116 t.Fatal(err)
117 }
118 }
119 if err := os.Chtimes(path, info.ModTime(), info.ModTime()); err != nil {
120 t.Fatal(err)
121 }
122 }
123
123 lines GO