返回 DeepSeek-Reasonix
catalog_test.go
根目录 / internal / sessioncatalog / catalog_test.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "fmt"
6 "os"
7 "path/filepath"
8 "runtime"
9 "strings"
10 "testing"
11 "time"
12
13 "reasonix/internal/agent"
14 "reasonix/internal/projectiondb"
15 )
16
17 func TestReconcileMakesUnknownCountsVisibleWithoutReadingTranscript(t *testing.T) {
18 t.Parallel()
19 ctx := context.Background()
20 dir := t.TempDir()
21 path := filepath.Join(dir, "legacy.jsonl")
22 if err := os.WriteFile(path, []byte("not valid jsonl\n"), 0o000); err != nil {
23 t.Fatal(err)
24 }
25 t.Cleanup(func() { _ = os.Chmod(path, 0o600) })
26 if err := agent.SaveBranchMeta(path, agent.BranchMeta{
27 Scope: "project",
28 WorkspaceRoot: "/workspace",
29 TopicID: "topic-1",
30 TopicTitle: "Legacy topic",
31 SchemaVersion: 1,
32 Turns: 0,
33 }); err != nil {
34 t.Fatal(err)
35 }
36
37 catalog, err := Open(ctx, Options{
38 Path: filepath.Join(t.TempDir(), "catalog.sqlite"),
39 DisableRepair: true,
40 })
41 if err != nil {
42 t.Fatal(err)
43 }
44 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
45
46 if err := catalog.ReconcileDirectory(ctx, DirectoryTarget{
47 Path: dir,
48 Scope: "project",
49 WorkspaceRoot: "/workspace",
50 }); err != nil {
51 t.Fatal(err)
52 }
53 page, err := catalog.ListTopics(ctx, TopicPageRequest{
54 Scope: "project",
55 WorkspaceRoot: "/workspace",
56 Limit: 50,
57 })
58 if err != nil {
59 t.Fatal(err)
60 }
61 if len(page.Items) != 1 || len(page.Items[0].Sessions) != 1 {
62 t.Fatalf("page = %#v, want one visible topic/session", page)
63 }
64 if got := page.Items[0].Sessions[0].TurnsState; got != TurnsUnknown {
65 t.Fatalf("turns state = %q, want %q", got, TurnsUnknown)
66 }
67 }
68
69 func TestDirectoryScanReadyOnlyAfterFirstReconcile(t *testing.T) {
70 t.Parallel()
71 ctx := context.Background()
72 dir := t.TempDir()
73 if err := os.WriteFile(filepath.Join(dir, "chat.jsonl"), []byte(`{"role":"user","content":"hi"}`+"\n"), 0o600); err != nil {
74 t.Fatal(err)
75 }
76 catalog, err := Open(ctx, Options{InMemory: true, DisableRepair: true})
77 if err != nil {
78 t.Fatal(err)
79 }
80 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
81 if catalog.DirectoryScanReady(ctx, dir) {
82 t.Fatal("opened catalog must not report the directory ready before the first scan")
83 }
84 if catalog.HasWorkspaceRecords(ctx, "global", "") {
85 t.Fatal("opened catalog must not report workspace records before the first scan")
86 }
87 if err := catalog.ReconcileDirectory(ctx, DirectoryTarget{Path: dir, Scope: "global"}); err != nil {
88 t.Fatal(err)
89 }
90 if !catalog.DirectoryScanReady(ctx, dir) {
91 t.Fatal("directory must be ready after ReconcileDirectory finishes")
92 }
93 if !catalog.HasWorkspaceRecords(ctx, "global", "") {
94 t.Fatal("reconciled directory should report workspace records")
95 }
96 }
97
98 func TestListTopicsUsesStableKeysetCursor(t *testing.T) {
99 t.Parallel()
100 ctx := context.Background()
101 catalog, err := Open(ctx, Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true})
102 if err != nil {
103 t.Fatal(err)
104 }
105 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
106
107 base := time.Date(2026, 8, 10, 10, 0, 0, 0, time.UTC)
108 for i, topicID := range []string{"a", "b", "c"} {
109 if err := catalog.UpsertSession(ctx, SessionRecord{
110 Path: filepath.Join("/sessions", topicID+".jsonl"),
111 Directory: "/sessions",
112 Scope: "global",
113 TopicID: topicID,
114 TopicTitle: topicID,
115 LastActivityAt: base.Add(time.Duration(i) * time.Minute).UnixMilli(),
116 Turns: i + 1,
117 TurnsState: TurnsValid,
118 Health: HealthOK,
119 }); err != nil {
120 t.Fatal(err)
121 }
122 }
123 first, err := catalog.ListTopics(ctx, TopicPageRequest{Scope: "global", Limit: 2})
124 if err != nil {
125 t.Fatal(err)
126 }
127 if len(first.Items) != 2 || first.NextCursor == "" {
128 t.Fatalf("first page = %#v", first)
129 }
130 second, err := catalog.ListTopics(ctx, TopicPageRequest{Scope: "global", Limit: 2, Cursor: first.NextCursor})
131 if err != nil {
132 t.Fatal(err)
133 }
134 if len(second.Items) != 1 || second.NextCursor != "" {
135 t.Fatalf("second page = %#v", second)
136 }
137 if first.Items[0].TopicID != "c" || first.Items[1].TopicID != "b" || second.Items[0].TopicID != "a" {
138 t.Fatalf("keyset order = %q, %q, %q", first.Items[0].TopicID, first.Items[1].TopicID, second.Items[0].TopicID)
139 }
140 }
141
142 func TestListTopicsSortsByCreationTimeAcrossPages(t *testing.T) {
143 t.Parallel()
144 ctx := context.Background()
145 catalog, err := Open(ctx, Options{InMemory: true, DisableRepair: true})
146 if err != nil {
147 t.Fatal(err)
148 }
149 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
150
151 for _, record := range []SessionRecord{
152 {Path: "/sessions/a.jsonl", Directory: "/sessions", Scope: "global", TopicID: "a", TopicTitle: "a", CreatedAt: 300, LastActivityAt: 100, Turns: 1, TurnsState: TurnsValid, Health: HealthOK},
153 {Path: "/sessions/b.jsonl", Directory: "/sessions", Scope: "global", TopicID: "b", TopicTitle: "b", CreatedAt: 200, LastActivityAt: 300, Turns: 1, TurnsState: TurnsValid, Health: HealthOK},
154 {Path: "/sessions/c.jsonl", Directory: "/sessions", Scope: "global", TopicID: "c", TopicTitle: "c", CreatedAt: 100, LastActivityAt: 200, Turns: 1, TurnsState: TurnsValid, Health: HealthOK},
155 } {
156 if err := catalog.UpsertSession(ctx, record); err != nil {
157 t.Fatal(err)
158 }
159 }
160
161 first, err := catalog.ListTopics(ctx, TopicPageRequest{Scope: "global", Limit: 2, SortMode: "created"})
162 if err != nil {
163 t.Fatal(err)
164 }
165 second, err := catalog.ListTopics(ctx, TopicPageRequest{Scope: "global", Limit: 2, SortMode: "created", Cursor: first.NextCursor})
166 if err != nil {
167 t.Fatal(err)
168 }
169 if len(first.Items) != 2 || first.NextCursor == "" || len(second.Items) != 1 || second.NextCursor != "" {
170 t.Fatalf("created pages = first %#v second %#v", first, second)
171 }
172 if first.Items[0].TopicID != "a" || first.Items[1].TopicID != "b" || second.Items[0].TopicID != "c" {
173 t.Fatalf("created order = %q, %q, %q; want a, b, c", first.Items[0].TopicID, first.Items[1].TopicID, second.Items[0].TopicID)
174 }
175 }
176
177 func TestGetTopicIsNotLimitedByPageSize(t *testing.T) {
178 t.Parallel()
179 ctx := context.Background()
180 catalog, err := Open(ctx, Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true})
181 if err != nil {
182 t.Fatal(err)
183 }
184 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
185
186 for i := 0; i <= MaxLimit; i++ {
187 topicID := fmt.Sprintf("topic-%03d", i)
188 if err := catalog.UpsertSession(ctx, SessionRecord{
189 Path: filepath.Join("/sessions", topicID+".jsonl"), Directory: "/sessions",
190 Scope: "global", TopicID: topicID, TopicTitle: topicID,
191 LastActivityAt: int64(MaxLimit - i), TurnsState: TurnsValid, Health: HealthOK,
192 }); err != nil {
193 t.Fatal(err)
194 }
195 }
196
197 topic, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "topic-200"})
198 if err != nil {
199 t.Fatal(err)
200 }
201 if !ok || topic.TopicID != "topic-200" || len(topic.Sessions) != 1 {
202 t.Fatalf("topic = %#v, ok=%v", topic, ok)
203 }
204 }
205
206 func TestUnchangedDirectorySignatureSkipsReconcileRevision(t *testing.T) {
207 t.Parallel()
208 ctx := context.Background()
209 dir := t.TempDir()
210 path := filepath.Join(dir, "session.jsonl")
211 if err := os.WriteFile(path, []byte("{}\n"), 0o600); err != nil {
212 t.Fatal(err)
213 }
214 if err := agent.SaveBranchMeta(path, agent.BranchMeta{
215 Scope: "global", TopicID: "topic", TopicTitle: "Topic",
216 SchemaVersion: agent.BranchMetaCountsVersion, Turns: 1,
217 }); err != nil {
218 t.Fatal(err)
219 }
220 catalog, err := Open(ctx, Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true})
221 if err != nil {
222 t.Fatal(err)
223 }
224 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
225 target := DirectoryTarget{Path: dir, Scope: "global"}
226 if err := catalog.ReconcileDirectory(ctx, target); err != nil {
227 t.Fatal(err)
228 }
229 revision := catalog.Status().Revision
230 if err := catalog.ReconcileDirectory(ctx, target); err != nil {
231 t.Fatal(err)
232 }
233 if got := catalog.Status().Revision; got != revision {
234 t.Fatalf("unchanged scan bumped revision: got %d want %d", got, revision)
235 }
236 }
237
238 func TestSyncMetadataRemovesOnlyMetadataOnlyTopics(t *testing.T) {
239 t.Parallel()
240 ctx := context.Background()
241 catalog, err := Open(ctx, Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true})
242 if err != nil {
243 t.Fatal(err)
244 }
245 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
246 if err := catalog.SyncMetadata(ctx, nil, []TopicMetadata{
247 {Scope: "global", TopicID: "metadata-only", Title: "Metadata"},
248 {Scope: "global", TopicID: "with-session", Title: "Session"},
249 }); err != nil {
250 t.Fatal(err)
251 }
252 if err := catalog.UpsertSession(ctx, SessionRecord{
253 Path: "/sessions/with-session.jsonl", Directory: "/sessions", Scope: "global",
254 TopicID: "with-session", Turns: 1, TurnsState: TurnsValid, Health: HealthOK,
255 }); err != nil {
256 t.Fatal(err)
257 }
258 if err := catalog.SyncMetadata(ctx, nil, nil); err != nil {
259 t.Fatal(err)
260 }
261 if _, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "metadata-only"}); err != nil || ok {
262 t.Fatalf("metadata-only topic survived removal: ok=%v err=%v", ok, err)
263 }
264 if _, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "with-session"}); err != nil || !ok {
265 t.Fatalf("session-derived topic was removed: ok=%v err=%v", ok, err)
266 }
267 }
268
269 func TestSchemaMigrationLedgerRecordsEveryVersion(t *testing.T) {
270 t.Parallel()
271 ctx := context.Background()
272 catalog, err := Open(ctx, Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true})
273 if err != nil {
274 t.Fatal(err)
275 }
276 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
277 rows, err := catalog.db.QueryContext(ctx, `SELECT version FROM schema_migrations ORDER BY version`)
278 if err != nil {
279 t.Fatal(err)
280 }
281 defer rows.Close()
282 versions := []int{}
283 for rows.Next() {
284 var version int
285 if err := rows.Scan(&version); err != nil {
286 t.Fatal(err)
287 }
288 versions = append(versions, version)
289 }
290 if fmt.Sprint(versions) != "[1 2 3 4 5 6 7 8 9 10 11 12 13]" {
291 t.Fatalf("schema migration ledger = %v", versions)
292 }
293 }
294
295 func TestSchemaV13RebuildsFilesystemIdentityProjection(t *testing.T) {
296 ctx := context.Background()
297 path := filepath.Join(t.TempDir(), "catalog.sqlite")
298 legacy, err := projectiondb.Open(ctx, projectiondb.OpenOptions{
299 Path: path, Migrations: sessionMigrations()[:12], Now: time.Now,
300 })
301 if err != nil {
302 t.Fatal(err)
303 }
304 if _, err := legacy.DB.ExecContext(ctx, `INSERT INTO catalog_directories(path,path_key,scope) VALUES('/old','/old','global')`); err != nil {
305 t.Fatal(err)
306 }
307 if err := legacy.DB.Close(); err != nil {
308 t.Fatal(err)
309 }
310 catalog, err := Open(ctx, Options{Path: path, DisableRepair: true})
311 if err != nil {
312 t.Fatal(err)
313 }
314 defer catalog.Close(context.Background())
315 var directories int
316 if err := catalog.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM catalog_directories`).Scan(&directories); err != nil {
317 t.Fatal(err)
318 }
319 if directories != 0 {
320 t.Fatalf("stale identity rows survived v13 migration: %d", directories)
321 }
322 }
323
324 func TestSchemaV11MigratesLegacyUnknownRowsIntoPersistentScheduler(t *testing.T) {
325 ctx := context.Background()
326 path := filepath.Join(t.TempDir(), "catalog.sqlite")
327 legacy, err := projectiondb.Open(ctx, projectiondb.OpenOptions{
328 Path: path, Migrations: sessionMigrations()[:10], Now: time.Now,
329 })
330 if err != nil {
331 t.Fatal(err)
332 }
333 if _, err := legacy.DB.ExecContext(ctx, `INSERT INTO catalog_sessions(path,directory,scope,turns_state)
334 VALUES('/sessions/legacy.jsonl','/sessions','global','unknown')`); err != nil {
335 t.Fatal(err)
336 }
337 if err := legacy.DB.Close(); err != nil {
338 t.Fatal(err)
339 }
340
341 migrated, err := projectiondb.Open(ctx, projectiondb.OpenOptions{Path: path, Migrations: sessionMigrations()[:11], Now: time.Now})
342 if err != nil {
343 t.Fatal(err)
344 }
345 t.Cleanup(func() { _ = migrated.DB.Close() })
346 catalog := &Catalog{db: migrated.DB}
347 var state string
348 var attempts, retryAt, engine int
349 if err := catalog.db.QueryRowContext(ctx, `SELECT repair_state,repair_attempts,repair_retry_at,repair_engine_version
350 FROM catalog_sessions WHERE path='/sessions/legacy.jsonl'`).Scan(&state, &attempts, &retryAt, &engine); err != nil {
351 t.Fatal(err)
352 }
353 if state != "pending" || attempts != 0 || retryAt != 0 || engine != 0 {
354 t.Fatalf("migrated repair schedule = %s/%d/%d/%d", state, attempts, retryAt, engine)
355 }
356 if err := catalog.resetRepairSchedule(ctx); err != nil {
357 t.Fatal(err)
358 }
359 if _, err := catalog.db.ExecContext(ctx, `UPDATE catalog_sessions SET repair_state='blocked',repair_attempts=7,
360 repair_error_kind='unsupported',repair_engine_version=0 WHERE path='/sessions/legacy.jsonl'`); err != nil {
361 t.Fatal(err)
362 }
363 if err := catalog.resetRepairSchedule(ctx); err != nil {
364 t.Fatal(err)
365 }
366 if err := catalog.db.QueryRowContext(ctx, `SELECT repair_state,repair_attempts,repair_retry_at,repair_engine_version
367 FROM catalog_sessions WHERE path='/sessions/legacy.jsonl'`).Scan(&state, &attempts, &retryAt, &engine); err != nil {
368 t.Fatal(err)
369 }
370 if state != "pending" || attempts != 0 || retryAt != 0 || engine != repairEngineVersion {
371 t.Fatalf("repair engine reset = %s/%d/%d/%d", state, attempts, retryAt, engine)
372 }
373 }
374
375 func TestListSessionsUsesRevisionBoundKeysetCursor(t *testing.T) {
376 t.Parallel()
377 ctx := context.Background()
378 catalog, err := Open(ctx, Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true})
379 if err != nil {
380 t.Fatal(err)
381 }
382 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
383 for i, name := range []string{"a", "b", "c"} {
384 if err := catalog.UpsertSession(ctx, SessionRecord{
385 Path: filepath.Join("/sessions", name+".jsonl"), Directory: "/sessions", Scope: "global",
386 TopicID: name, CustomTitle: "title " + name, LastActivityAt: int64(i + 1),
387 TurnsState: TurnsValid, Health: HealthOK,
388 }); err != nil {
389 t.Fatal(err)
390 }
391 }
392 first, err := catalog.ListSessions(ctx, SessionPageRequest{Scope: "all", Limit: 2})
393 if err != nil || len(first.Items) != 2 || first.NextCursor == "" {
394 t.Fatalf("first=%#v err=%v", first, err)
395 }
396 second, err := catalog.ListSessions(ctx, SessionPageRequest{Scope: "all", Limit: 2, Cursor: first.NextCursor})
397 if err != nil || len(second.Items) != 1 || second.Items[0].CustomTitle != "title a" {
398 t.Fatalf("second=%#v err=%v", second, err)
399 }
400 if err := catalog.UpsertSession(ctx, SessionRecord{Path: "/sessions/d.jsonl", Directory: "/sessions", Scope: "global", TopicID: "d", LastActivityAt: 4, TurnsState: TurnsValid, Health: HealthOK}); err != nil {
401 t.Fatal(err)
402 }
403 stale, err := catalog.ListSessions(ctx, SessionPageRequest{Scope: "all", Cursor: first.NextCursor})
404 if err != nil || !stale.StaleCursor || len(stale.Items) != 0 {
405 t.Fatalf("stale=%#v err=%v", stale, err)
406 }
407 }
408
409 func TestDirectWriteDuringScanIsNotMarkedMissing(t *testing.T) {
410 t.Parallel()
411 ctx := context.Background()
412 now := time.Date(2026, 8, 10, 10, 0, 0, 0, time.UTC)
413 dir := t.TempDir()
414 catalog, err := Open(ctx, Options{
415 Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true,
416 Now: func() time.Time { return now },
417 })
418 if err != nil {
419 t.Fatal(err)
420 }
421 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
422 target := DirectoryTarget{Path: dir, Scope: "global"}
423 generation, _, err := catalog.beginDirectoryScan(ctx, target, "test", now.UnixMilli())
424 if err != nil {
425 t.Fatal(err)
426 }
427 path := filepath.Join(dir, "late.jsonl")
428 if err := catalog.UpsertSession(ctx, SessionRecord{
429 Path: path, Directory: dir, Scope: "global", TopicID: "late",
430 TurnsState: TurnsUnknown, Health: HealthOK,
431 }); err != nil {
432 t.Fatal(err)
433 }
434 if err := catalog.finishDirectoryScan(ctx, target, "test", generation, now.UnixMilli(), 0); err != nil {
435 t.Fatal(err)
436 }
437 topic, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "late"})
438 if err != nil || !ok || len(topic.Sessions) != 1 || topic.Sessions[0].Health != HealthOK {
439 t.Fatalf("late write was marked missing: topic=%#v ok=%v err=%v", topic, ok, err)
440 }
441 }
442
443 func TestSessionTopicMoveRecomputesOldTopic(t *testing.T) {
444 t.Parallel()
445 ctx := context.Background()
446 catalog, err := Open(ctx, Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true})
447 if err != nil {
448 t.Fatal(err)
449 }
450 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
451 record := SessionRecord{
452 Path: "/sessions/moved.jsonl", Directory: "/sessions", Scope: "global",
453 TopicID: "old", Turns: 1, TurnsState: TurnsValid, Health: HealthOK,
454 }
455 if err := catalog.UpsertSession(ctx, record); err != nil {
456 t.Fatal(err)
457 }
458 record.TopicID = "new"
459 if err := catalog.UpsertSession(ctx, record); err != nil {
460 t.Fatal(err)
461 }
462 if _, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "old"}); err != nil || ok {
463 t.Fatalf("old topic survived move: ok=%v err=%v", ok, err)
464 }
465 if _, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "new"}); err != nil || !ok {
466 t.Fatalf("new topic missing after move: ok=%v err=%v", ok, err)
467 }
468 }
469
470 func TestWriterQueueCoalescesBySessionPath(t *testing.T) {
471 t.Parallel()
472 ctx := context.Background()
473 catalog, err := Open(ctx, Options{
474 Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true, QueueCapacity: 1,
475 })
476 if err != nil {
477 t.Fatal(err)
478 }
479 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
480 for turns := 1; turns <= 100; turns++ {
481 if ok := catalog.EnqueueSession(SessionRecord{
482 Path: "/sessions/coalesced.jsonl", Directory: "/sessions", Scope: "global",
483 TopicID: "coalesced", Turns: turns, TurnsState: TurnsValid, Health: HealthOK,
484 }); !ok {
485 t.Fatalf("same-path update %d was rejected by a one-slot queue", turns)
486 }
487 }
488 deadline := time.Now().Add(2 * time.Second)
489 for time.Now().Before(deadline) {
490 topic, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "coalesced"})
491 if err != nil {
492 t.Fatal(err)
493 }
494 if ok && topic.Turns == 100 {
495 return
496 }
497 time.Sleep(10 * time.Millisecond)
498 }
499 t.Fatal("coalesced writer did not persist the latest record")
500 }
501
502 func TestRemoveSessionTombstoneHidesTopicBeforeDurableDelete(t *testing.T) {
503 t.Parallel()
504 ctx := context.Background()
505 dir := t.TempDir()
506 path := filepath.Join(dir, "session.jsonl")
507 catalog, err := Open(ctx, Options{
508 Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true,
509 })
510 if err != nil {
511 t.Fatal(err)
512 }
513 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
514 record := SessionRecord{
515 Path: path, Directory: dir, Scope: "global", TopicID: "topic_tombstone",
516 Turns: 1, TurnsState: TurnsValid, Health: HealthOK, LastActivityAt: time.Now().UnixMilli(),
517 }
518 if err := catalog.UpsertSession(ctx, record); err != nil {
519 t.Fatal(err)
520 }
521 // Hold the directory lock so durable DELETE blocks while the short caller
522 // context expires — the read-visible tombstone must still hide the topic.
523 dirLock := catalog.directoryLock(dir)
524 dirLock.Lock()
525 defer dirLock.Unlock()
526
527 short, cancel := context.WithTimeout(ctx, 30*time.Millisecond)
528 defer cancel()
529 // RemoveSession should return promptly (overlay path) without waiting for
530 // the held directory lock forever.
531 done := make(chan error, 1)
532 go func() {
533 done <- catalog.RemoveSession(short, path, "test_tombstone_overlay")
534 }()
535 select {
536 case err := <-done:
537 if err != nil {
538 t.Fatalf("RemoveSession: %v", err)
539 }
540 case <-time.After(2 * time.Second):
541 t.Fatal("RemoveSession blocked on directory lock instead of recording tombstone first")
542 }
543 if page, err := catalog.ListTopics(ctx, TopicPageRequest{Scope: "global", Limit: 50}); err != nil {
544 t.Fatal(err)
545 } else {
546 for _, item := range page.Items {
547 if item.TopicID == "topic_tombstone" {
548 t.Fatalf("tombstoned topic still visible in ListTopics: %+v", page.Items)
549 }
550 }
551 }
552 if _, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "topic_tombstone"}); err != nil || ok {
553 t.Fatalf("GetTopic after tombstone: ok=%v err=%v", ok, err)
554 }
555 if _, ok, err := catalog.GetSession(ctx, path); err != nil || ok {
556 t.Fatalf("GetSession after tombstone: ok=%v err=%v", ok, err)
557 }
558 }
559
560 func TestRemoveSessionWinsOverQueuedStaleWriteAndAllowsLaterRecreation(t *testing.T) {
561 t.Parallel()
562 ctx := context.Background()
563 dir := t.TempDir()
564 path := filepath.Join(dir, "session.jsonl")
565 catalog, err := Open(ctx, Options{
566 Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true,
567 })
568 if err != nil {
569 t.Fatal(err)
570 }
571 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
572 record := SessionRecord{
573 Path: path, Directory: dir, Scope: "global", TopicID: "topic",
574 Turns: 1, TurnsState: TurnsValid, Health: HealthOK,
575 }
576 if err := catalog.UpsertSession(ctx, record); err != nil {
577 t.Fatal(err)
578 }
579 record.Turns = 99
580 if !catalog.EnqueueSession(record) {
581 t.Fatal("queue stale write")
582 }
583 if err := catalog.RemoveSession(ctx, path, "test_remove"); err != nil {
584 t.Fatal(err)
585 }
586 time.Sleep(50 * time.Millisecond)
587 if _, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "topic"}); err != nil || ok {
588 t.Fatalf("queued write resurrected removed row: ok=%v err=%v", ok, err)
589 }
590 if err := os.WriteFile(path, []byte("{}\n"), 0o600); err != nil {
591 t.Fatal(err)
592 }
593 if err := agent.SaveBranchMeta(path, agent.BranchMeta{
594 Scope: "global", TopicID: "topic", SchemaVersion: agent.BranchMetaCountsVersion, Turns: 2,
595 }); err != nil {
596 t.Fatal(err)
597 }
598 if err := catalog.ReconcileDirectory(ctx, DirectoryTarget{Path: dir, Scope: "global"}); err != nil {
599 t.Fatal(err)
600 }
601 if topic, ok, err := catalog.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "topic"}); err != nil || !ok || topic.Turns != 2 {
602 t.Fatalf("new external file did not supersede removal: topic=%#v ok=%v err=%v", topic, ok, err)
603 }
604 }
605
606 func TestCloseCancelsCatalogWorkerContext(t *testing.T) {
607 catalog, err := Open(context.Background(), Options{
608 Path: filepath.Join(t.TempDir(), "catalog.sqlite"), DisableRepair: true,
609 })
610 if err != nil {
611 t.Fatal(err)
612 }
613 workerDone := catalog.workerCtx.Done()
614 timeout := time.Second
615 if runtime.GOOS == "windows" {
616 timeout = 5 * time.Second
617 }
618 ctx, cancel := context.WithTimeout(context.Background(), timeout)
619 defer cancel()
620 if err := catalog.Close(ctx); err != nil {
621 t.Fatal(err)
622 }
623 select {
624 case <-workerDone:
625 default:
626 t.Fatal("catalog close left worker context active")
627 }
628 }
629
630 func BenchmarkListTopicsWarmCatalog10K(b *testing.B) {
631 ctx := context.Background()
632 catalog, err := Open(ctx, Options{
633 Path: filepath.Join(b.TempDir(), "catalog.sqlite"), DisableRepair: true,
634 })
635 if err != nil {
636 b.Fatal(err)
637 }
638 b.Cleanup(func() { _ = catalog.Close(context.Background()) })
639 records := make([]SessionRecord, 10_000)
640 for i := range records {
641 records[i] = SessionRecord{
642 Path: filepath.Join("/sessions", fmt.Sprintf("%05d.jsonl", i)), Directory: "/sessions",
643 Scope: "project", WorkspaceRoot: "/workspace", TopicID: fmt.Sprintf("topic-%05d", i),
644 Preview: fmt.Sprintf("synthetic session %05d", i), Turns: i%20 + 1,
645 TurnsState: TurnsValid, Health: HealthOK, LastActivityAt: int64(i + 1),
646 }
647 }
648 for start := 0; start < len(records); start += 64 {
649 end := min(start+64, len(records))
650 if err := catalog.upsertSessions(ctx, records[start:end], nil, "benchmark-setup"); err != nil {
651 b.Fatal(err)
652 }
653 }
654 b.ResetTimer()
655 for range b.N {
656 page, err := catalog.ListTopics(ctx, TopicPageRequest{
657 Scope: "project", WorkspaceRoot: "/workspace", Limit: 50,
658 })
659 if err != nil || len(page.Items) != 50 {
660 b.Fatalf("page len=%d err=%v", len(page.Items), err)
661 }
662 }
663 }
664
665 func TestMissingSessionRequiresTwoScansAndGraceBeforeRemoval(t *testing.T) {
666 t.Parallel()
667 ctx := context.Background()
668 now := time.Date(2026, 8, 10, 10, 0, 0, 0, time.UTC)
669 dir := t.TempDir()
670 path := filepath.Join(dir, "session.jsonl")
671 if err := os.WriteFile(path, []byte("{}\n"), 0o600); err != nil {
672 t.Fatal(err)
673 }
674 if err := agent.SaveBranchMeta(path, agent.BranchMeta{
675 Scope: "global", TopicID: "topic", TopicTitle: "Topic",
676 SchemaVersion: agent.BranchMetaCountsVersion, Turns: 1,
677 }); err != nil {
678 t.Fatal(err)
679 }
680 catalog, err := Open(ctx, Options{
681 Path: filepath.Join(t.TempDir(), "catalog.sqlite"),
682 DisableRepair: true,
683 MissingGrace: time.Minute,
684 Now: func() time.Time { return now },
685 })
686 if err != nil {
687 t.Fatal(err)
688 }
689 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
690 if err := catalog.ReconcileDirectory(ctx, DirectoryTarget{Path: dir, Scope: "global"}); err != nil {
691 t.Fatal(err)
692 }
693 if err := os.Remove(path); err != nil {
694 t.Fatal(err)
695 }
696 if err := catalog.ReconcileDirectory(ctx, DirectoryTarget{Path: dir, Scope: "global"}); err != nil {
697 t.Fatal(err)
698 }
699 if page, err := catalog.ListTopics(ctx, TopicPageRequest{Scope: "global", Limit: 50}); err != nil || len(page.Items) != 1 {
700 t.Fatalf("first missing scan removed row: page=%#v err=%v", page, err)
701 }
702 now = now.Add(2 * time.Minute)
703 if err := catalog.ReconcileDirectory(ctx, DirectoryTarget{Path: dir, Scope: "global"}); err != nil {
704 t.Fatal(err)
705 }
706 if page, err := catalog.ListTopics(ctx, TopicPageRequest{Scope: "global", Limit: 50}); err != nil || len(page.Items) != 0 {
707 t.Fatalf("stale row survived second scan after grace: page=%#v err=%v", page, err)
708 }
709 }
710
711 func TestOpenQuarantinesCorruptCatalogAndRebuildsProjection(t *testing.T) {
712 t.Parallel()
713 ctx := context.Background()
714 dir := t.TempDir()
715 path := filepath.Join(dir, "catalog.sqlite")
716 if err := os.WriteFile(path, []byte("not sqlite"), 0o600); err != nil {
717 t.Fatal(err)
718 }
719 catalog, err := Open(ctx, Options{Path: path, DisableRepair: true})
720 if err != nil {
721 t.Fatal(err)
722 }
723 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
724 status := catalog.Status()
725 if status.State != StateReady || status.QuarantinedPath == "" {
726 t.Fatalf("status = %#v, want ready catalog with quarantined path", status)
727 }
728 if _, err := os.Stat(status.QuarantinedPath); err != nil {
729 t.Fatalf("quarantined catalog: %v", err)
730 }
731 }
732
733 func TestOpenBlankPathUsesMemoryWithoutWritingCWD(t *testing.T) {
734 // A blank path (CacheDir unavailable or caller override) must use memory
735 // and must not create a relative session-catalog file under cwd.
736 wd := t.TempDir()
737 t.Chdir(wd)
738 catalog, err := Open(context.Background(), Options{Path: " ", DisableRepair: true})
739 if err != nil {
740 t.Fatal(err)
741 }
742 t.Cleanup(func() { _ = catalog.Close(context.Background()) })
743 status := catalog.Status()
744 if status.Mode != ModeMemory {
745 t.Fatalf("status=%#v, want memory mode for blank path", status)
746 }
747 entries, err := os.ReadDir(wd)
748 if err != nil {
749 t.Fatal(err)
750 }
751 for _, entry := range entries {
752 if strings.Contains(entry.Name(), "session-catalog") || strings.HasSuffix(entry.Name(), ".sqlite") {
753 t.Fatalf("blank path wrote projection into cwd: %s", entry.Name())
754 }
755 }
756 }
757
758 func TestRebuildFailureKeepsExistingCatalog(t *testing.T) {
759 t.Parallel()
760 ctx := context.Background()
761 path := filepath.Join(t.TempDir(), "catalog.sqlite")
762 catalog, err := Open(ctx, Options{Path: path, DisableRepair: true})
763 if err != nil {
764 t.Fatal(err)
765 }
766 dir := t.TempDir()
767 session := filepath.Join(dir, "keep.jsonl")
768 if err := os.WriteFile(session, []byte(`{"role":"user","content":"hello"}`+"\n"), 0o600); err != nil {
769 t.Fatal(err)
770 }
771 if err := catalog.UpsertSession(ctx, SessionRecord{
772 Path: session, Directory: dir, Scope: "global", TopicID: "keep",
773 Turns: 1, TurnsState: TurnsValid, Health: HealthOK, Preview: "hello",
774 }); err != nil {
775 t.Fatal(err)
776 }
777 if err := catalog.Close(context.Background()); err != nil {
778 t.Fatal(err)
779 }
780 // ListSessionOrder on a regular file fails; Rebuild must keep the old DB.
781 fileTarget := filepath.Join(t.TempDir(), "not-a-dir")
782 if err := os.WriteFile(fileTarget, []byte("x"), 0o600); err != nil {
783 t.Fatal(err)
784 }
785 if _, err := Rebuild(ctx, path, []DirectoryTarget{{Path: fileTarget, Scope: "global"}}); err == nil {
786 t.Fatal("expected rebuild failure for file path")
787 }
788 restored, err := Open(ctx, Options{Path: path, DisableRepair: true})
789 if err != nil {
790 t.Fatal(err)
791 }
792 t.Cleanup(func() { _ = restored.Close(context.Background()) })
793 topic, ok, err := restored.GetTopic(ctx, TopicKey{Scope: "global", TopicID: "keep"})
794 if err != nil || !ok || len(topic.Sessions) != 1 {
795 t.Fatalf("rebuild failure lost catalog: ok=%v topic=%#v err=%v", ok, topic, err)
796 }
797 }
798
798 lines GO