返回 DeepSeek-Reasonix
recovery_store_test.go
根目录 / internal / session / recovery_store_test.go
1 package session
2
3 import (
4 "context"
5 "encoding/json"
6 "os"
7 "path/filepath"
8 "strconv"
9 "strings"
10 "testing"
11 "time"
12
13 bolt "go.etcd.io/bbolt"
14 "reasonix/internal/provider"
15 )
16
17 func TestRecoveryCheckpointRestoresTailAndIdempotencyWithoutPrefixReplay(t *testing.T) {
18 dir := filepath.Join(t.TempDir(), "sessions-v4", "fast")
19 s, err := CreateWithOptions(dir, "fast", OpenOptions{ExternalHistory: true})
20 if err != nil {
21 t.Fatal(err)
22 }
23
24 for i := range 80 {
25 message := provider.Message{ID: "message-" + strings.Repeat("x", 8) + strconv.Itoa(i), Role: provider.RoleUser, Content: strings.Repeat("history", 256)}
26 payload, marshalErr := json.Marshal(map[string]any{"message": message})
27 if marshalErr != nil {
28 t.Fatal(marshalErr)
29 }
30 if _, err := s.Append(context.Background(), Batch{OperationID: "operation-" + string(rune(0x100+i)), Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
31 t.Fatal(err)
32 }
33 }
34 if _, err := s.Flush(t.Context()); err != nil {
35 t.Fatal(err)
36 }
37 wantModel, err := json.Marshal(s.DeriveMessages())
38 if err != nil {
39 t.Fatal(err)
40 }
41 if err := s.Close(t.Context()); err != nil {
42 t.Fatal(err)
43 }
44
45 // Every query projection is disposable. Removing it must not disable the
46 // recovery fast path or recent-message access.
47 if err := os.RemoveAll(filepath.Join(filepath.Dir(dir), ".query-cache", "fast")); err != nil {
48 t.Fatal(err)
49 }
50 var stats RecoveryOpenStats
51 reopened, err := OpenWithOptions(dir, "fast", OpenOptions{
52 ExternalHistory: true,
53 ObserveRecovery: func(got RecoveryOpenStats) { stats = got },
54 })
55 if err != nil {
56 t.Fatal(err)
57 }
58 t.Cleanup(func() { _ = reopened.Close(context.Background()) })
59 if !stats.UsedCheckpoint {
60 t.Fatalf("recovery stats = %+v, want checkpoint", stats)
61 }
62 if stats.LogBytesRead >= stats.LogBytesTotal {
63 t.Fatalf("recovery read %d of %d bytes, want bounded tail", stats.LogBytesRead, stats.LogBytesTotal)
64 }
65 gotModel, err := json.Marshal(reopened.DeriveMessages())
66 if err != nil {
67 t.Fatal(err)
68 }
69 if string(gotModel) != string(wantModel) {
70 t.Fatal("provider-visible model projection changed across checkpoint recovery")
71 }
72 recent := reopened.RecentSnapshot()
73 if recent.StorageGeneration == "" || recent.DurableSequence != reopened.EventSequence() || len(recent.Entries) == 0 || len(recent.Entries) > RecentMessageLimit {
74 t.Fatalf("recent snapshot = %+v", recent)
75 }
76
77 // A durable operation remains globally idempotent even though old operation
78 // records are no longer resident in Session.operations.
79 firstPayload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: "message-" + strings.Repeat("x", 8) + "0", Role: provider.RoleUser, Content: strings.Repeat("history", 256)}})
80 first, err := reopened.Append(t.Context(), Batch{OperationID: "operation-" + string(rune(0x100)), Events: []Event{{Kind: "message/complete", Payload: firstPayload}}})
81 if err != nil {
82 t.Fatal(err)
83 }
84 if first.FirstSequence != 1 {
85 t.Fatalf("idempotent commit starts at %d, want 1", first.FirstSequence)
86 }
87 conflictPayload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: "conflict", Role: provider.RoleUser, Content: "different"}})
88 if _, err := reopened.Append(t.Context(), Batch{OperationID: "operation-" + string(rune(0x100)), Events: []Event{{Kind: "message/complete", Payload: conflictPayload}}}); err == nil {
89 t.Fatal("same operation key with different content unexpectedly succeeded")
90 }
91 }
92
93 func TestRecentEntriesKeepAbsoluteVisibleTurns(t *testing.T) {
94 dir := filepath.Join(t.TempDir(), "session")
95 if err := os.MkdirAll(dir, 0o700); err != nil {
96 t.Fatal(err)
97 }
98 messages := []provider.Message{
99 {ID: "eight", Role: provider.RoleUser, Content: "eight"},
100 {ID: "assistant", Role: provider.RoleAssistant, Content: "answer"},
101 {ID: "nine", Role: provider.RoleUser, Content: "nine"},
102 {ID: "ten", Role: provider.RoleUser, Content: "ten"},
103 }
104 entries, err := buildRecentEntries(t.Context(), dir, messages, 10, 10)
105 if err != nil {
106 t.Fatal(err)
107 }
108 want := []int{8, 8, 9, 10}
109 for index := range entries {
110 if entries[index].VisibleTurn != want[index] {
111 t.Fatalf("entry %d visible turn=%d, want %d", index, entries[index].VisibleTurn, want[index])
112 }
113 }
114 }
115
116 func TestRecoveryCheckpointReplaysOnlyDurableTail(t *testing.T) {
117 root := filepath.Join(t.TempDir(), "sessions-v4")
118 dir := filepath.Join(root, "tail")
119 first, err := CreateWithOptions(dir, "tail", OpenOptions{ExternalHistory: true})
120 if err != nil {
121 t.Fatal(err)
122 }
123 appendRecoveryTestMessage(t, first, "one", "first")
124 if _, err := first.Flush(t.Context()); err != nil {
125 t.Fatal(err)
126 }
127 if err := first.Close(t.Context()); err != nil {
128 t.Fatal(err)
129 }
130
131 // Suppress checkpoint publication for one append to simulate a log that is
132 // ahead of its last valid recovery point.
133 second, err := OpenWithOptions(dir, "tail", OpenOptions{ExternalHistory: true, DisableRecoveryPublish: true})
134 if err != nil {
135 t.Fatal(err)
136 }
137 appendRecoveryTestMessage(t, second, "two", "second")
138 if _, err := second.Flush(t.Context()); err != nil {
139 t.Fatal(err)
140 }
141 if err := second.Close(t.Context()); err != nil {
142 t.Fatal(err)
143 }
144
145 var stats RecoveryOpenStats
146 third, err := OpenWithOptions(dir, "tail", OpenOptions{ExternalHistory: true, ObserveRecovery: func(got RecoveryOpenStats) { stats = got }})
147 if err != nil {
148 t.Fatal(err)
149 }
150 defer third.Close(context.Background())
151 if !stats.UsedCheckpoint || stats.TailCommits != 1 {
152 t.Fatalf("recovery stats = %+v, want one tail commit", stats)
153 }
154 if got := third.DeriveMessages(); len(got) != 2 || got[0].Content != "first" || got[1].Content != "second" {
155 t.Fatalf("tail projection = %#v", got)
156 }
157 }
158
159 func TestRecoveryCheckpointFallsBackToPreviousGeneration(t *testing.T) {
160 root := filepath.Join(t.TempDir(), "sessions-v4")
161 dir := filepath.Join(root, "previous")
162 store, err := CreateWithOptions(dir, "previous", OpenOptions{ExternalHistory: true})
163 if err != nil {
164 t.Fatal(err)
165 }
166 appendRecoveryTestMessage(t, store, "one", "first")
167 if _, err := store.Flush(t.Context()); err != nil {
168 t.Fatal(err)
169 }
170 appendRecoveryTestMessage(t, store, "two", "second")
171 if _, err := store.Flush(t.Context()); err != nil {
172 t.Fatal(err)
173 }
174 if err := store.Close(t.Context()); err != nil {
175 t.Fatal(err)
176 }
177
178 path := filepath.Join(recoveryCacheDir(dir), recoveryDBName)
179 db, err := bolt.Open(path, 0o600, &bolt.Options{Timeout: time.Second})
180 if err != nil {
181 t.Fatal(err)
182 }
183 if err := db.Update(func(tx *bolt.Tx) error {
184 return tx.Bucket(recoveryCheckpointBucket).Put(recoveryCurrentKey, []byte("damaged"))
185 }); err != nil {
186 _ = db.Close()
187 t.Fatal(err)
188 }
189 if err := db.Close(); err != nil {
190 t.Fatal(err)
191 }
192
193 var stats RecoveryOpenStats
194 reopened, err := OpenWithOptions(dir, "previous", OpenOptions{ExternalHistory: true, ObserveRecovery: func(got RecoveryOpenStats) { stats = got }})
195 if err != nil {
196 t.Fatal(err)
197 }
198 defer reopened.Close(context.Background())
199 if !stats.UsedCheckpoint {
200 t.Fatalf("fallback recovery stats = %+v", stats)
201 }
202 if got := reopened.DeriveMessages(); len(got) != 2 || got[0].ID != "one" || got[1].ID != "two" {
203 t.Fatalf("fallback messages = %+v", got)
204 }
205 }
206
207 func appendRecoveryTestMessage(t *testing.T, s *Session, id, content string) {
208 t.Helper()
209 payload, err := json.Marshal(map[string]any{"message": provider.Message{ID: id, Role: provider.RoleUser, Content: content}})
210 if err != nil {
211 t.Fatal(err)
212 }
213 if _, err := s.Append(t.Context(), Batch{OperationID: "append-" + id, Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
214 t.Fatal(err)
215 }
216 }
217
217 lines GO