返回 DeepSeek-Reasonix
history_window_test.go
根目录 / internal / session / history_window_test.go
1 package session
2
3 import (
4 "context"
5 "encoding/json"
6 "path/filepath"
7 "strings"
8 "testing"
9 "time"
10
11 "reasonix/internal/provider"
12 )
13
14 func windowReady(t *testing.T, query *Query, ref SessionRef, req HistoryWindowRequest) HistoryWindowPage {
15 t.Helper()
16 page, err := waitHistoryWindow(t, query, ref, req)
17 if err != nil {
18 t.Fatal(err)
19 }
20 return page
21 }
22
23 func waitHistoryWindow(t *testing.T, query *Query, ref SessionRef, req HistoryWindowRequest) (HistoryWindowPage, error) {
24 t.Helper()
25 deadline := time.Now().Add(5 * time.Second)
26 for {
27 page, err := query.ReadHistoryWindow(t.Context(), ref, req)
28 if err != nil || page.Status != "preparing" {
29 return page, err
30 }
31 if time.Now().After(deadline) {
32 t.Fatalf("window stayed preparing: %+v", page)
33 }
34 time.Sleep(time.Millisecond)
35 }
36 }
37
38 func appendWindowMessages(t *testing.T, runtime *Runtime, ids ...string) {
39 t.Helper()
40 for _, id := range ids {
41 payload, err := json.Marshal(map[string]any{"message": provider.Message{ID: id, Role: provider.RoleUser, Content: "body-" + id}})
42 if err != nil {
43 t.Fatal(err)
44 }
45 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: "message-" + id, Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
46 t.Fatal(err)
47 }
48 }
49 if _, err := runtime.Session().Flush(t.Context()); err != nil {
50 t.Fatal(err)
51 }
52 }
53
54 func windowIDs(t *testing.T, page HistoryWindowPage) []string {
55 t.Helper()
56 ids := make([]string, 0, len(page.Messages))
57 for _, m := range page.Messages {
58 ids = append(ids, m.MessageID)
59 }
60 return ids
61 }
62
63 func idsEqual(a, b []string) bool {
64 if len(a) != len(b) {
65 return false
66 }
67 for i := range a {
68 if a[i] != b[i] {
69 return false
70 }
71 }
72 return true
73 }
74
75 func TestHistoryWindowNewestAndOlderPaging(t *testing.T) {
76 root := filepath.Join(t.TempDir(), "sessions-v4")
77 service, err := NewService("local", NewFilesystemPersistence(root))
78 if err != nil {
79 t.Fatal(err)
80 }
81 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
82 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "windowed"})
83 if err != nil {
84 t.Fatal(err)
85 }
86 t.Cleanup(func() { _ = service.Close(context.Background(), runtime.Ref()) })
87 appendWindowMessages(t, runtime, "m1", "m2", "m3", "m4", "m5")
88 ref := runtime.Ref()
89 query := service.Query()
90
91 newest := windowReady(t, query, ref, HistoryWindowRequest{Anchor: "newest", Limit: 2})
92 if newest.Status != "ready" || !idsEqual(windowIDs(t, newest), []string{"m4", "m5"}) {
93 t.Fatalf("newest page = %v", windowIDs(t, newest))
94 }
95 if !newest.HasOlder || newest.OlderCursor == "" || newest.HasNewer {
96 t.Fatalf("newest page cursors = %+v", newest)
97 }
98 older := windowReady(t, query, ref, HistoryWindowRequest{Anchor: "cursor", Cursor: newest.OlderCursor, Limit: 2})
99 if !idsEqual(windowIDs(t, older), []string{"m2", "m3"}) {
100 t.Fatalf("older page = %v", windowIDs(t, older))
101 }
102 if !older.HasNewer || older.NewerCursor == "" {
103 t.Fatalf("older page must expose a newer cursor: %+v", older)
104 }
105 newer := windowReady(t, query, ref, HistoryWindowRequest{Anchor: "cursor", Cursor: older.NewerCursor, Limit: 5})
106 if !idsEqual(windowIDs(t, newer), []string{"m4", "m5"}) {
107 t.Fatalf("newer continuation = %v", windowIDs(t, newer))
108 }
109 oldest := windowReady(t, query, ref, HistoryWindowRequest{Anchor: "cursor", Cursor: older.OlderCursor, Limit: 5})
110 if !idsEqual(windowIDs(t, oldest), []string{"m1"}) || oldest.HasOlder {
111 t.Fatalf("oldest page = %v hasOlder=%v", windowIDs(t, oldest), oldest.HasOlder)
112 }
113 }
114
115 func TestHistoryWindowMessageAndTurnAnchors(t *testing.T) {
116 root := filepath.Join(t.TempDir(), "sessions-v4")
117 service, err := NewService("local", NewFilesystemPersistence(root))
118 if err != nil {
119 t.Fatal(err)
120 }
121 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
122 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "anchored"})
123 if err != nil {
124 t.Fatal(err)
125 }
126 t.Cleanup(func() { _ = service.Close(context.Background(), runtime.Ref()) })
127 appendWindowMessages(t, runtime, "a1", "a2", "a3", "a4")
128 ref := runtime.Ref()
129 query := service.Query()
130
131 // A search-hit-style jump lands a window around the target without
132 // walking from the newest page.
133 byMessage := windowReady(t, query, ref, HistoryWindowRequest{Anchor: "message", MessageID: "a2", Limit: 2})
134 if !idsEqual(windowIDs(t, byMessage), []string{"a1", "a2"}) {
135 t.Fatalf("message anchor page = %v", windowIDs(t, byMessage))
136 }
137 if !byMessage.HasNewer || byMessage.NewerCursor == "" {
138 t.Fatalf("message anchor must continue newer: %+v", byMessage)
139 }
140 afterTarget := windowReady(t, query, ref, HistoryWindowRequest{Anchor: "cursor", Cursor: byMessage.NewerCursor, Limit: 2})
141 if !idsEqual(windowIDs(t, afterTarget), []string{"a3", "a4"}) {
142 t.Fatalf("newer continuation = %v", windowIDs(t, afterTarget))
143 }
144
145 byTurn := windowReady(t, query, ref, HistoryWindowRequest{Anchor: "turn", Turn: 4, Limit: 1})
146 if len(byTurn.Messages) != 1 || byTurn.Messages[0].MessageID != "a4" {
147 t.Fatalf("turn anchor page = %+v", byTurn)
148 }
149 missing := windowReady(t, query, ref, HistoryWindowRequest{Anchor: "message", MessageID: "absent", Limit: 2})
150 if missing.Status != "not_found" {
151 t.Fatalf("missing anchor status = %s", missing.Status)
152 }
153 }
154
155 func mustEncodeWindowCursor(t *testing.T, c historyWindowCursor) string {
156 t.Helper()
157 cursor, err := encodeHistoryWindowCursor(c)
158 if err != nil {
159 t.Fatal(err)
160 }
161 return cursor
162 }
163
164 func TestHistoryWindowRejectsForeignCursor(t *testing.T) {
165 root := filepath.Join(t.TempDir(), "sessions-v4")
166 service, err := NewService("local", NewFilesystemPersistence(root))
167 if err != nil {
168 t.Fatal(err)
169 }
170 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
171 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "fenced"})
172 if err != nil {
173 t.Fatal(err)
174 }
175 t.Cleanup(func() { _ = service.Close(context.Background(), runtime.Ref()) })
176 appendWindowMessages(t, runtime, "f1", "f2")
177 ref := runtime.Ref()
178 page := windowReady(t, service.Query(), ref, HistoryWindowRequest{Anchor: "newest", Limit: 1})
179 if page.OlderCursor == "" {
180 t.Fatal("expected an older cursor")
181 }
182 parsed, err := decodeHistoryWindowCursor(page.OlderCursor)
183 if err != nil {
184 t.Fatal(err)
185 }
186 // A cursor whose storage generation no longer matches — as after a
187 // storage replacement or projection rebuild — answers stale_cursor.
188 foreign := parsed
189 foreign.Generation = "g-none"
190 stale, err := service.Query().ReadHistoryWindow(t.Context(), ref, HistoryWindowRequest{Anchor: "cursor", Cursor: mustEncodeWindowCursor(t, foreign), Limit: 1})
191 if err != nil {
192 t.Fatal(err)
193 }
194 if stale.Status != "stale_cursor" {
195 t.Fatalf("foreign generation cursor = %s, want stale_cursor", stale.Status)
196 }
197 foreign = parsed
198 foreign.SnapshotSequence += 1 << 40
199 overrun, err := service.Query().ReadHistoryWindow(t.Context(), ref, HistoryWindowRequest{Anchor: "cursor", Cursor: mustEncodeWindowCursor(t, foreign), Limit: 1})
200 if err != nil {
201 t.Fatal(err)
202 }
203 if overrun.Status != "stale_cursor" {
204 t.Fatalf("future snapshot cursor = %s, want stale_cursor", overrun.Status)
205 }
206 }
207
208 func TestReadMessageFieldReturnsBoundedAlignedFragments(t *testing.T) {
209 root := filepath.Join(t.TempDir(), "sessions-v4")
210 service, err := NewService("local", NewFilesystemPersistence(root))
211 if err != nil {
212 t.Fatal(err)
213 }
214 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
215 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "fields"})
216 if err != nil {
217 t.Fatal(err)
218 }
219 t.Cleanup(func() { _ = service.Close(context.Background(), runtime.Ref()) })
220 content := strings.Repeat("汉", 30000) + `quote " and \backslash` + strings.Repeat("字", 3000)
221 payload, err := json.Marshal(map[string]any{"message": provider.Message{ID: "big", Role: provider.RoleUser, Content: content}})
222 if err != nil {
223 t.Fatal(err)
224 }
225 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: "message-big", Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
226 t.Fatal(err)
227 }
228 if _, err := runtime.Session().Flush(t.Context()); err != nil {
229 t.Fatal(err)
230 }
231 ref := runtime.Ref()
232 query := service.Query()
233 // Display the page once so the referenced content gains its range grant.
234 page := historyPageReady(t, query, ref, "", 10)
235 if len(page.Messages) != 1 || page.Messages[0].ContentRef == nil {
236 t.Fatalf("expected a referenced large message, got %+v", page)
237 }
238
239 first, err := query.ReadMessageField(t.Context(), ref, "big", 0, "content", 0, 64)
240 if err != nil {
241 t.Fatal(err)
242 }
243 if first.Status != "ready" || first.TotalBytes == 0 || len(first.Data) == 0 || first.NextOffset == 0 {
244 t.Fatalf("first fragment = status:%s total:%d data:%d next:%d", first.Status, first.TotalBytes, len(first.Data), first.NextOffset)
245 }
246 var assembled []byte
247 assembled = append(assembled, first.Data...)
248 offset := first.NextOffset
249 for offset != 0 {
250 frag, err := query.ReadMessageField(t.Context(), ref, "big", 0, "content", offset, 256<<10)
251 if err != nil {
252 t.Fatal(err)
253 }
254 assembled = append(assembled, frag.Data...)
255 offset = frag.NextOffset
256 }
257 var decoded string
258 if err := json.Unmarshal(assembled, &decoded); err != nil {
259 t.Fatalf("reassembled field is not the original JSON string: %v", err)
260 }
261 if decoded != content {
262 t.Fatal("reassembled content mismatch")
263 }
264 whole, err := query.ReadMessageField(t.Context(), ref, "big", 0, "canonicalMessage", 0, 0)
265 if err != nil {
266 t.Fatal(err)
267 }
268 var message provider.Message
269 if err := json.Unmarshal(whole.Data, &message); err != nil {
270 t.Fatalf("canonicalMessage fragment did not parse: %v", err)
271 }
272 absent, err := query.ReadMessageField(t.Context(), ref, "absent", 0, "content", 0, 64)
273 if err != nil {
274 t.Fatal(err)
275 }
276 if absent.Status != "not_found" {
277 t.Fatalf("unknown message status = %s, want not_found", absent.Status)
278 }
279 if _, err := query.ReadMessageField(t.Context(), ref, "big", 0, "", 0, 64); err == nil {
280 t.Fatal("field read without a field name succeeded")
281 }
282 if _, err := query.ReadMessageField(t.Context(), ref, "big", 0, "content", -1, 64); err == nil {
283 t.Fatal("negative offset accepted")
284 }
285 }
286
286 lines GO