| 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 |