返回 DeepSeek-Reasonix
history_index_test.go
根目录 / internal / session / history_index_test.go
1 package session
2
3 import (
4 "context"
5 "encoding/json"
6 "os"
7 "path/filepath"
8 "strings"
9 "testing"
10
11 "reasonix/internal/projectiondb"
12 "reasonix/internal/provider"
13 )
14
15 func searchHistoryReady(t *testing.T, query *Query, ref SessionRef, text, cursor string, limit int) SearchHistoryPage {
16 t.Helper()
17 for {
18 page, err := query.SearchHistory(t.Context(), ref, text, cursor, limit)
19 if err != nil {
20 t.Fatal(err)
21 }
22 if page.Status == "ready" {
23 return page
24 }
25 if page.Status != "preparing" {
26 t.Fatalf("search preparation = %+v", page)
27 }
28 query.searchMu.Lock()
29 preparation := query.searchBuilds[ref.SessionID]
30 query.searchMu.Unlock()
31 if preparation == nil {
32 t.Fatalf("search preparing without a worker for %s", ref.SessionID)
33 }
34 t.Cleanup(func() {
35 query.Close()
36 <-preparation.done
37 })
38 select {
39 case <-preparation.done:
40 if preparation.err != nil {
41 t.Fatal(preparation.err)
42 }
43 case <-t.Context().Done():
44 t.Fatal(t.Context().Err())
45 }
46 }
47 }
48
49 // A rebuild may scan appends made after its initial stat. Its continuation
50 // offset must describe the same completed commit as its sequence watermark.
51 func TestHistoryIndexRebuildPairsScannedSequenceAndOffset(t *testing.T) {
52 root := filepath.Join(t.TempDir(), "sessions")
53 persistence := NewFilesystemPersistence(root)
54 service, err := NewService("local", persistence)
55 if err != nil {
56 t.Fatal(err)
57 }
58 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
59 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "rebuild-cut"})
60 if err != nil {
61 t.Fatal(err)
62 }
63 appendMessage := func(id string) {
64 t.Helper()
65 payload, err := json.Marshal(map[string]any{"message": provider.Message{ID: id, Role: provider.RoleAssistant, Content: id}})
66 if err != nil {
67 t.Fatal(err)
68 }
69 if _, err = runtime.Session().AppendBatch(t.Context(), id, []Event{{Kind: "message/complete", Payload: payload}}); err != nil {
70 t.Fatal(err)
71 }
72 if _, err = runtime.Session().Flush(t.Context()); err != nil {
73 t.Fatal(err)
74 }
75 }
76 appendMessage("before-stat")
77 dir := filepath.Join(root, runtime.Ref().SessionID)
78 staleRevision, err := revisionOfLog(dir)
79 if err != nil {
80 t.Fatal(err)
81 }
82 appendMessage("after-stat")
83 path := historyIndexPath(root, runtime.Ref().SessionID)
84 if err := rebuildHistoryIndex(t.Context(), dir, path, runtime.Ref().SessionID, staleRevision); err != nil {
85 t.Fatal(err)
86 }
87 appendMessage("after-scan")
88 page := historyPageReady(t, service.Query(), runtime.Ref(), "", 32)
89 if len(page.Messages) != 3 {
90 t.Fatalf("messages = %+v", page.Messages)
91 }
92 for i, id := range []string{"before-stat", "after-stat", "after-scan"} {
93 if page.Messages[i].MessageID != id {
94 t.Fatalf("message %d = %+v", i, page.Messages[i])
95 }
96 }
97 }
98
99 func waitHistoryPage(t *testing.T, query *Query, ref SessionRef, cursor string, limit int) (MessageHistoryPage, error) {
100 t.Helper()
101 for {
102 page, err := query.HistoryPage(t.Context(), ref, cursor, limit)
103 if err != nil || page.Status != "preparing" {
104 return page, err
105 }
106 if err := waitHistoryPreparation(t, query, ref); err != nil {
107 return MessageHistoryPage{}, err
108 }
109 }
110 }
111
112 func waitHistoryPreparation(t *testing.T, query *Query, ref SessionRef) error {
113 t.Helper()
114 // Join the worker on test cleanup even if an assertion interrupts the wait.
115 t.Cleanup(query.Close)
116 query.historyMu.Lock()
117 preparation := query.historyBuilds[ref.SessionID]
118 query.historyMu.Unlock()
119 if preparation == nil {
120 // A synchronous reader can own the lock without a rebuild record.
121 // Join a preparation that waits for that reader and verifies readiness.
122 filesystem, ok := query.persistence.(*FilesystemPersistence)
123 if !ok {
124 t.Fatal("history preparation requires filesystem persistence")
125 }
126 preparation = query.prepareHistoryLocator(filesystem, ref.SessionID, historyIndexPath(filesystem.Root, ref.SessionID))
127 }
128 select {
129 case <-preparation.done:
130 return preparation.err
131 case <-t.Context().Done():
132 return t.Context().Err()
133 }
134 }
135
136 func historyPageReady(t *testing.T, query *Query, ref SessionRef, cursor string, limit int) MessageHistoryPage {
137 t.Helper()
138 page, err := waitHistoryPage(t, query, ref, cursor, limit)
139 if err != nil {
140 t.Fatal(err)
141 }
142 if page.Status != "ready" {
143 t.Fatalf("history page = %+v", page)
144 }
145 return page
146 }
147
148 func TestExternalHistoryColdOpenDefersBodiesBeforeModelReset(t *testing.T) {
149 root := filepath.Join(t.TempDir(), "sessions-v4")
150 service, err := NewService("local", NewFilesystemPersistence(root))
151 if err != nil {
152 t.Fatal(err)
153 }
154 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
155 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "bounded-open"})
156 if err != nil {
157 t.Fatal(err)
158 }
159 oldPayload, err := json.Marshal(map[string]any{"message": provider.Message{ID: "old", Role: provider.RoleUser, Content: strings.Repeat("old", 40<<10)}})
160 if err != nil {
161 t.Fatal(err)
162 }
163 if _, err := runtime.Session().AppendBatch(t.Context(), "old", []Event{{Kind: "message/complete", Payload: oldPayload}}); err != nil {
164 t.Fatal(err)
165 }
166 current := provider.Message{ID: "current", Role: provider.RoleUser, Content: "current workset"}
167 currentPayload, err := json.Marshal(map[string]any{"messages": []provider.Message{current}, "reason": "bounded cold open"})
168 if err != nil {
169 t.Fatal(err)
170 }
171 if _, err := runtime.Session().AppendBatch(t.Context(), "reset", []Event{{Kind: "model/context-replace", Payload: currentPayload}}); err != nil {
172 t.Fatal(err)
173 }
174 if _, err := runtime.Session().Flush(t.Context()); err != nil {
175 t.Fatal(err)
176 }
177 ref := runtime.Ref()
178 if err := service.Close(t.Context(), ref); err != nil {
179 t.Fatal(err)
180 }
181
182 log, err := os.Open(filepath.Join(root, ref.SessionID, "events.frames"))
183 if err != nil {
184 t.Fatal(err)
185 }
186 var historicalDigest string
187 err = scanV4CommitFileRefs(t.Context(), log, 0, 1, contentStoreForSessionDir(filepath.Join(root, ref.SessionID)), nil, func(_ int64, commit Commit) bool {
188 for _, event := range commit.Events {
189 if event.Kind == "message/complete" && event.PayloadRef != nil {
190 historicalDigest = event.PayloadRef.Digest
191 }
192 }
193 return true
194 })
195 _ = log.Close()
196 if err != nil || historicalDigest == "" {
197 t.Fatalf("historical content reference = %q, %v", historicalDigest, err)
198 }
199 object := filepath.Join(root, ".content-v1", "objects", historicalDigest[:2], historicalDigest[2:4], historicalDigest)
200 if err := os.Remove(object); err != nil {
201 t.Fatal(err)
202 }
203
204 reopenedService, err := NewService("local", NewFilesystemPersistence(root))
205 if err != nil {
206 t.Fatal(err)
207 }
208 t.Cleanup(func() { _ = reopenedService.CloseAll(context.Background()) })
209 binding, err := reopenedService.Open(t.Context(), ref)
210 if err != nil {
211 t.Fatalf("cold open resolved retired history body: %v", err)
212 }
213 defer binding.Release(context.Background())
214 model := binding.Runtime().Session().DeriveMessages()
215 if len(model) != 1 || model[0].ID != current.ID || model[0].Content != current.Content {
216 t.Fatalf("cold model projection = %+v", model)
217 }
218 if _, err := waitHistoryPage(t, reopenedService.Query(), ref, "", 100); err == nil {
219 t.Fatal("history query accepted a missing referenced body")
220 }
221 }
222
223 func TestHistoryPageKeepsSnapshotAndAuthorizesReferencedContent(t *testing.T) {
224 root := filepath.Join(t.TempDir(), "sessions-v4")
225 service, err := NewService("local", NewFilesystemPersistence(root))
226 if err != nil {
227 t.Fatal(err)
228 }
229 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
230 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "paged"})
231 if err != nil {
232 t.Fatal(err)
233 }
234 t.Cleanup(func() {
235 if err := service.Close(context.Background(), runtime.Ref()); err != nil {
236 t.Error(err)
237 }
238 })
239 appendMessage := func(id, content string) {
240 payload, err := json.Marshal(map[string]any{"message": provider.Message{ID: id, Role: provider.RoleUser, Content: content}})
241 if err != nil {
242 t.Fatal(err)
243 }
244 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: "message-" + id, Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
245 t.Fatal(err)
246 }
247 if _, err := runtime.Session().Flush(t.Context()); err != nil {
248 t.Fatal(err)
249 }
250 }
251 appendMessage("one", "first")
252 appendMessage("two", strings.Repeat("large", 20000))
253 appendMessage("three", "third")
254 runtime.Session().mu.Lock()
255 residentMessages := len(runtime.Session().projection.Messages)
256 residentModel := len(runtime.Session().projection.ModelMessages)
257 runtime.Session().mu.Unlock()
258 if residentMessages != 0 || residentModel != 3 {
259 t.Fatalf("service runtime retained durable UI bodies: messages=%d model=%d", residentMessages, residentModel)
260 }
261 ref := runtime.Ref()
262 first := historyPageReady(t, service.Query(), ref, "", 1)
263 if len(first.Messages) != 1 || first.Messages[0].MessageID != "three" || !first.HasMore || first.NextCursor == "" {
264 t.Fatalf("first page = %+v", first)
265 }
266 appendMessage("four", "must not enter the fixed snapshot")
267 second, err := waitHistoryPage(t, service.Query(), ref, first.NextCursor, 10)
268 if err != nil {
269 t.Fatal(err)
270 }
271 if len(second.Messages) != 2 || second.Messages[0].MessageID != "one" || second.Messages[1].MessageID != "two" {
272 t.Fatalf("fixed snapshot second page = %+v", second)
273 }
274 large := second.Messages[1]
275 if large.ContentRef == nil || len(large.Inline) != 0 {
276 t.Fatalf("large message was not referenced: %+v", large)
277 }
278 if err := os.RemoveAll(historyIndexPath(root, "paged")); err != nil {
279 t.Fatal(err)
280 }
281 chunk, err := service.Query().ReadContent(t.Context(), ref, *large.ContentRef, 0, min(64, large.ContentRef.Bytes))
282 if err != nil {
283 t.Fatal(err)
284 }
285 var decoded provider.Message
286 full, err := service.Query().ReadContent(t.Context(), ref, *large.ContentRef, 0, large.ContentRef.Bytes)
287 if err != nil {
288 t.Fatal(err)
289 }
290 if len(chunk) == 0 || json.Unmarshal(full, &decoded) != nil || decoded.ID != "two" {
291 t.Fatalf("resolved content prefix=%q id=%q", chunk, decoded.ID)
292 }
293 foreign := *large.ContentRef
294 foreign.Digest = strings.Repeat("0", len(foreign.Digest))
295 if _, err := service.Query().ReadContent(t.Context(), ref, foreign, 0, 1); err == nil {
296 t.Fatal("content hash without a session reference was authorized")
297 }
298 if _, err := service.Query().ReadContent(t.Context(), ref, *large.ContentRef, 0, (1<<20)+1); err == nil {
299 t.Fatal("oversized content range was accepted")
300 }
301 }
302
303 func TestLocateMessageReturnsFixedSnapshotCursorWithoutBody(t *testing.T) {
304 root := filepath.Join(t.TempDir(), "sessions-v4")
305 service, err := NewService("local", NewFilesystemPersistence(root))
306 if err != nil {
307 t.Fatal(err)
308 }
309 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
310 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "locate"})
311 if err != nil {
312 t.Fatal(err)
313 }
314 for _, id := range []string{"one", "two", "three"} {
315 payload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: id, Role: provider.RoleUser, Content: id}})
316 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: "locate-" + id, Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
317 t.Fatal(err)
318 }
319 }
320 if _, err := runtime.Session().Flush(t.Context()); err != nil {
321 t.Fatal(err)
322 }
323 _ = historyPageReady(t, service.Query(), runtime.Ref(), "", 1)
324 location, err := service.Query().LocateMessage(t.Context(), runtime.Ref(), "two", 0)
325 if err != nil || location.Status != "ready" || location.Position == 0 || location.Cursor == "" {
326 t.Fatalf("location = %+v, %v", location, err)
327 }
328 page, err := service.Query().HistoryPage(t.Context(), runtime.Ref(), location.Cursor, 1)
329 if err != nil || len(page.Messages) != 1 || page.Messages[0].MessageID != "two" {
330 t.Fatalf("located page = %+v, %v", page, err)
331 }
332 }
333
334 func TestHistoryLocatorGenerationSurvivesRebuildAndChangesOnReplacement(t *testing.T) {
335 root := filepath.Join(t.TempDir(), "sessions-v4")
336 service, err := NewService("local", NewFilesystemPersistence(root))
337 if err != nil {
338 t.Fatal(err)
339 }
340 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
341 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "locator-generation"})
342 if err != nil {
343 t.Fatal(err)
344 }
345 for _, id := range []string{"one", "two"} {
346 payload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: id, Role: provider.RoleUser, Content: id}})
347 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: "generation-" + id, Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
348 t.Fatal(err)
349 }
350 }
351 if _, err := runtime.Session().Flush(t.Context()); err != nil {
352 t.Fatal(err)
353 }
354 first := historyPageReady(t, service.Query(), runtime.Ref(), "", 1)
355 if first.NextCursor == "" || first.Generation == "" {
356 t.Fatalf("first page = %+v", first)
357 }
358 if err := os.Remove(historyIndexPath(root, runtime.Ref().SessionID)); err != nil {
359 t.Fatal(err)
360 }
361 service.Query().historyMu.Lock()
362 delete(service.Query().historyBuilds, runtime.Ref().SessionID)
363 service.Query().historyMu.Unlock()
364 rebuilt := historyPageReady(t, service.Query(), runtime.Ref(), "", 1)
365 if rebuilt.Generation != first.Generation {
366 t.Fatalf("ordinary rebuild changed generation: %q -> %q", first.Generation, rebuilt.Generation)
367 }
368 if page, err := service.Query().HistoryPage(t.Context(), runtime.Ref(), first.NextCursor, 1); err != nil || page.Status != "ready" || len(page.Messages) != 1 || page.Messages[0].MessageID != "one" {
369 t.Fatalf("cursor after rebuild = %+v, %v", page, err)
370 }
371 replacement := []provider.Message{{ID: "replacement", Role: provider.RoleUser, Content: "replacement"}}
372 payload, _ := json.Marshal(map[string]any{"messages": replacement})
373 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: "replace-history", Events: []Event{{Kind: "history/replace", Payload: payload}}}); err != nil {
374 t.Fatal(err)
375 }
376 if _, err := runtime.Session().Flush(t.Context()); err != nil {
377 t.Fatal(err)
378 }
379 if _, err := service.Query().HistoryPage(t.Context(), runtime.Ref(), "", 1); err != nil {
380 t.Fatal(err)
381 }
382 stale, err := waitHistoryPage(t, service.Query(), runtime.Ref(), first.NextCursor, 1)
383 if err != nil || stale.Status != "stale_cursor" {
384 t.Fatalf("cursor after replacement = %+v, %v", stale, err)
385 }
386 }
387
388 func TestSearchHistoryUsesStableSnapshotAndOpaqueQueryCursor(t *testing.T) {
389 root := filepath.Join(t.TempDir(), "sessions-v4")
390 service, err := NewService("local", NewFilesystemPersistence(root))
391 if err != nil {
392 t.Fatal(err)
393 }
394 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
395 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "search"})
396 if err != nil {
397 t.Fatal(err)
398 }
399 t.Cleanup(func() {
400 if err := service.Close(context.Background(), runtime.Ref()); err != nil {
401 t.Error(err)
402 }
403 })
404 appendMessage := func(id, content string) {
405 payload, _ := json.Marshal(map[string]any{"message": provider.Message{ID: id, Role: provider.RoleUser, Content: content}})
406 if _, err := runtime.Session().Append(t.Context(), Batch{OperationID: id, Events: []Event{{Kind: "message/complete", Payload: payload}}}); err != nil {
407 t.Fatal(err)
408 }
409 if _, err := runtime.Session().Flush(t.Context()); err != nil {
410 t.Fatal(err)
411 }
412 }
413 appendMessage("one", "first needle")
414 appendMessage("two", "second needle")
415 appendMessage("three", "unrelated")
416 preparing, err := service.Query().SearchHistory(t.Context(), runtime.Ref(), "needle", "", 1)
417 if err != nil || preparing.Status != "preparing" {
418 t.Fatalf("first search did not return preparation state: %+v, %v", preparing, err)
419 }
420 first := searchHistoryReady(t, service.Query(), runtime.Ref(), "needle", "", 1)
421 if len(first.Hits) != 1 || first.Hits[0].MessageID != "two" || !first.HasMore || first.NextCursor == "" {
422 t.Fatalf("first search = %+v", first)
423 }
424 appendMessage("four", "new needle outside snapshot")
425 second, err := service.Query().SearchHistory(t.Context(), runtime.Ref(), "needle", first.NextCursor, 10)
426 if err != nil {
427 t.Fatal(err)
428 }
429 if len(second.Hits) != 1 || second.Hits[0].MessageID != "one" {
430 t.Fatalf("second search = %+v", second)
431 }
432 stale, err := service.Query().SearchHistory(t.Context(), runtime.Ref(), "different", first.NextCursor, 10)
433 if err != nil || stale.Status != "stale_cursor" {
434 t.Fatalf("search cursor mismatch = %+v, %v", stale, err)
435 }
436 }
437
438 func TestSearchHistoryCoversInlineFieldsAndReferencedBodies(t *testing.T) {
439 root := filepath.Join(t.TempDir(), "sessions-v4")
440 service, err := NewService("local", NewFilesystemPersistence(root))
441 if err != nil {
442 t.Fatal(err)
443 }
444 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
445 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "search-storage"})
446 if err != nil {
447 t.Fatal(err)
448 }
449 t.Cleanup(func() {
450 if err := service.Close(context.Background(), runtime.Ref()); err != nil {
451 t.Error(err)
452 }
453 })
454 messages := []provider.Message{
455 {ID: "inline", Role: provider.RoleAssistant, RawContent: "raw-field-needle", ReasoningContent: "reasoning-field-needle"},
456 {ID: "referenced", Role: provider.RoleUser, Content: strings.Repeat("large-body-", 7000) + "referenced-field-needle"},
457 }
458 for _, message := range messages {
459 payload, marshalErr := json.Marshal(map[string]any{"message": message})
460 if marshalErr != nil {
461 t.Fatal(marshalErr)
462 }
463 if _, err := runtime.Session().AppendBatch(t.Context(), "append-"+message.ID, []Event{{Kind: "message/complete", Payload: payload}}); err != nil {
464 t.Fatal(err)
465 }
466 }
467 if _, err := runtime.Session().Flush(t.Context()); err != nil {
468 t.Fatal(err)
469 }
470 for query, want := range map[string]string{
471 "raw-field-needle": "inline",
472 "reasoning-field-needle": "inline",
473 "referenced-field-needle": "referenced",
474 } {
475 page := searchHistoryReady(t, service.Query(), runtime.Ref(), query, "", 10)
476 if len(page.Hits) != 1 || page.Hits[0].MessageID != want {
477 t.Fatalf("search %q = %+v, want %q", query, page.Hits, want)
478 }
479 }
480 }
481
482 func TestSearchHistoryKeepsLiteralUnicodeSubstringSemanticsIndependently(t *testing.T) {
483 root := filepath.Join(t.TempDir(), "sessions-v4")
484 service, err := NewService("local", NewFilesystemPersistence(root))
485 if err != nil {
486 t.Fatal(err)
487 }
488 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
489 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "literal-search"})
490 if err != nil {
491 t.Fatal(err)
492 }
493 appendRecoveryTestMessage(t, runtime.Session(), "cn", `大会话恢复包含字面符号 "OR" % _`)
494 appendRecoveryTestMessage(t, runtime.Session(), "other", "unrelated")
495 if _, err := runtime.Session().Flush(t.Context()); err != nil {
496 t.Fatal(err)
497 }
498 for _, query := range []string{"会", "会话", "话恢", `"OR"`, "%", "_"} {
499 page := searchHistoryReady(t, service.Query(), runtime.Ref(), query, "", 10)
500 if page.Status != "ready" || page.CoverageSequence != runtime.Session().EventSequence() || len(page.Hits) != 1 || page.Hits[0].MessageID != "cn" {
501 t.Fatalf("search %q = %+v", query, page)
502 }
503 }
504 searchPath := searchIndexPath(root, "literal-search")
505 before, err := os.Stat(searchPath)
506 if err != nil {
507 t.Fatal(err)
508 }
509 if err := os.RemoveAll(historyIndexPath(root, "literal-search")); err != nil {
510 t.Fatal(err)
511 }
512 appendRecoveryTestMessage(t, runtime.Session(), "new", "新的会话子串")
513 if _, err := runtime.Session().Flush(t.Context()); err != nil {
514 t.Fatal(err)
515 }
516 page := searchHistoryReady(t, service.Query(), runtime.Ref(), "会话", "", 10)
517 after, err := os.Stat(searchPath)
518 if err != nil {
519 t.Fatal(err)
520 }
521 if !os.SameFile(before, after) {
522 t.Fatal("search index was rebuilt instead of incrementally advanced")
523 }
524 if len(page.Hits) != 2 || page.Hits[0].MessageID != "new" || page.Hits[1].MessageID != "cn" {
525 t.Fatalf("independent incremental search = %+v", page)
526 }
527 }
528
529 func TestHistoryIndexHasSnapshotPositionIndex(t *testing.T) {
530 path := filepath.Join(t.TempDir(), "history.sqlite")
531 handle, err := projectiondb.Open(t.Context(), projectiondb.OpenOptions{Path: path, Migrations: historyMigrations, RequireDisk: true, MaxOpenConns: 1})
532 if err != nil {
533 t.Fatal(err)
534 }
535 defer handle.DB.Close()
536 rows, err := handle.DB.QueryContext(context.Background(), `PRAGMA index_list(messages)`)
537 if err != nil {
538 t.Fatal(err)
539 }
540 defer rows.Close()
541 found := false
542 for rows.Next() {
543 var sequence, unique, partial int
544 var name, origin string
545 if err := rows.Scan(&sequence, &name, &unique, &origin, &partial); err != nil {
546 t.Fatal(err)
547 }
548 found = found || name == "messages_snapshot_position"
549 }
550 if err := rows.Err(); err != nil {
551 t.Fatal(err)
552 }
553 if !found {
554 t.Fatal("snapshot-position query index is missing")
555 }
556 }
557
558 func TestHistoryIndexAdvancesInPlaceAfterAppend(t *testing.T) {
559 root := filepath.Join(t.TempDir(), "sessions-v4")
560 service, err := NewService("local", NewFilesystemPersistence(root))
561 if err != nil {
562 t.Fatal(err)
563 }
564 t.Cleanup(func() { _ = service.CloseAll(context.Background()) })
565 runtime, err := service.Create(t.Context(), CreateOptions{SessionID: "incremental"})
566 if err != nil {
567 t.Fatal(err)
568 }
569 appendRecoveryTestMessage(t, runtime.Session(), "first", "first")
570 if _, err := runtime.Session().Flush(t.Context()); err != nil {
571 t.Fatal(err)
572 }
573 _ = historyPageReady(t, service.Query(), runtime.Ref(), "", 100)
574 path := historyIndexPath(root, "incremental")
575 before, err := os.Stat(path)
576 if err != nil {
577 t.Fatal(err)
578 }
579 handle, err := projectiondb.Open(t.Context(), projectiondb.OpenOptions{Path: path, Migrations: historyMigrations, RequireDisk: true, MaxOpenConns: 1})
580 if err != nil {
581 t.Fatal(err)
582 }
583 var inlineBytes, searchBytes int64
584 if err := handle.DB.QueryRowContext(t.Context(), `SELECT COALESCE(SUM(length(inline)),0),COALESCE(SUM(length(search_text)),0) FROM messages`).Scan(&inlineBytes, &searchBytes); err != nil {
585 _ = handle.DB.Close()
586 t.Fatal(err)
587 }
588 _ = handle.DB.Close()
589 if inlineBytes != 0 || searchBytes != 0 {
590 t.Fatalf("locator retained message bodies: inline=%d search=%d", inlineBytes, searchBytes)
591 }
592 appendRecoveryTestMessage(t, runtime.Session(), "second", "second")
593 if _, err := runtime.Session().Flush(t.Context()); err != nil {
594 t.Fatal(err)
595 }
596 page, err := waitHistoryPage(t, service.Query(), runtime.Ref(), "", 100)
597 if err != nil {
598 t.Fatal(err)
599 }
600 after, err := os.Stat(path)
601 if err != nil {
602 t.Fatal(err)
603 }
604 if !os.SameFile(before, after) {
605 t.Fatal("history locator was replaced instead of advanced in place")
606 }
607 if len(page.Messages) != 2 || page.Messages[0].MessageID != "first" || page.Messages[1].MessageID != "second" {
608 t.Fatalf("incremental page = %+v", page)
609 }
610 }
611
611 lines GO