返回 DeepSeek-Reasonix
catalog.go
根目录 / internal / sessioncatalog / catalog.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "database/sql"
6 "encoding/base64"
7 "encoding/json"
8 "errors"
9 "fmt"
10 "os"
11 "path/filepath"
12 "reasonix/internal/agent"
13 "strings"
14 "sync"
15 "sync/atomic"
16 "time"
17 )
18
19 const defaultMissingGrace = 30 * time.Second
20
21 type Catalog struct {
22 db *sql.DB
23 opts Options
24 pathIdentity func(string) string
25 mutationSeq atomic.Uint64
26 discoveryIDs sync.Map
27 discoverySeq atomic.Uint64
28 revision atomic.Uint64
29 statusMu sync.RWMutex
30 status Status
31 writeCh chan string
32 writeMu sync.Mutex
33 writeQueued map[string]SessionRecord
34 // mutationMu is the process-local SQLite single-writer boundary. WAL permits
35 // concurrent readers, but repair, metadata, and reconcile mutations must not
36 // race into avoidable SQLITE_BUSY failures.
37 mutationMu sync.Mutex
38 metadataSyncOnce sync.Once
39 metadataSyncGate chan struct{} // serializes observations without holding the database writer
40 removedPaths sync.Map
41 repairCh chan string
42 repairQueued sync.Map
43 reconcileCh chan DirectoryTarget
44 reconcileQueued sync.Map
45 reconcileDirtyMu sync.Mutex
46 reconcileDirty map[string]DirectoryTarget
47 reconcileDone map[string]chan struct{}
48 verifiedDirsMu sync.RWMutex
49 verifiedDirs map[string]string
50 pathCh chan sessionPathRequest
51 pathQueueMu sync.Mutex
52 pathQueued sync.Map
53 directoryLocksMu sync.Mutex
54 directoryLocks map[string]*sync.Mutex
55 metadataScans sync.Map
56 discoveryStart chan struct{}
57 discoveryOnce sync.Once
58 priorityWorkspace atomic.Value
59 workerCtx context.Context
60 workerCancel context.CancelFunc
61 stop chan struct{}
62 stopOnce sync.Once
63 workers sync.WaitGroup
64 closeDone chan struct{}
65 closeErr error
66 invalidated chan struct{}
67 invalidReason error
68 invalidateOnce sync.Once
69 integrityDone chan struct{}
70 readLeasesMu sync.Mutex
71 readLeases map[*ReadLease]struct{}
72 readLeasesClosed bool
73 // testReconcileBatchHook deterministically pauses an uncommitted directory
74 // projection. Production catalogs leave it nil.
75 testReconcileBatchHook func(int)
76 // testReconcileStartHook observes queued reconcile waves. Direct explicit
77 // ReconcileDirectory calls do not invoke it.
78 testReconcileStartHook func(DirectoryTarget)
79 // testRepairSessionHook replaces the filesystem repair in scheduler tests.
80 testRepairSessionHook func(context.Context, string) (agent.SessionListingRepairResult, error)
81 // testRepairBatchError injects publication failures by transaction stage.
82 testRepairBatchError func(string) error
83 // testSessionContentLoadHook counts strict lineage snapshot loads.
84 testSessionContentLoadHook func(string)
85 // testPathMutationLoadedHook pauses after reading a removal generation.
86 // Production catalogs leave it nil.
87 testPathMutationLoadedHook func(string)
88 // testScanProgressWriteHook runs after acquiring the shared writer boundary.
89 testScanProgressWriteHook func()
90 // Runs between committed metadata slices, after releasing the writer.
91 testMetadataSliceHook func(int)
92 }
93
94 type sessionPathRequest struct {
95 target DirectoryTarget
96 path string
97 queueKey string
98 sequence uint64
99 }
100
101 type pageCursor struct {
102 Pinned int `json:"p"`
103 ManualOrder bool `json:"m,omitempty"`
104 SortOrder int64 `json:"o,omitempty"`
105 Activity int64 `json:"a"`
106 TopicID string `json:"t"`
107 Binding string `json:"b,omitempty"`
108 }
109
110 func Open(ctx context.Context, opts Options) (*Catalog, error) {
111 if opts.DeferredMetadataIntegrity && !opts.MetadataOnly {
112 return nil, errors.New("deferred integrity requires an advisory metadata catalog")
113 }
114 if opts.MetadataOnly {
115 opts.DisableRepair = true
116 }
117 if opts.Path == "" {
118 opts.Path = DefaultPath()
119 }
120 if opts.Now == nil {
121 opts.Now = time.Now
122 }
123 if opts.MissingGrace <= 0 {
124 opts.MissingGrace = defaultMissingGrace
125 }
126 if opts.QueueCapacity <= 0 {
127 opts.QueueCapacity = 1024
128 }
129 // An empty path (no cache dir) or explicit memory flag must never write a
130 // relative session-catalog file into the current project directory.
131 if strings.TrimSpace(opts.Path) == "" {
132 opts.Path = ""
133 opts.InMemory = true
134 }
135 if !opts.InMemory {
136 if env := strings.TrimSpace(os.Getenv("REASONIX_SESSION_CATALOG_MEMORY")); env == "1" {
137 opts.InMemory = true
138 }
139 }
140
141 c := &Catalog{
142 opts: opts,
143 pathIdentity: PathIdentityKey,
144 writeCh: make(chan string, opts.QueueCapacity),
145 writeQueued: map[string]SessionRecord{},
146 repairCh: make(chan string, opts.QueueCapacity),
147 reconcileCh: make(chan DirectoryTarget, 64),
148 reconcileDirty: map[string]DirectoryTarget{},
149 verifiedDirs: map[string]string{},
150 pathCh: make(chan sessionPathRequest, opts.QueueCapacity),
151 directoryLocks: map[string]*sync.Mutex{},
152 stop: make(chan struct{}),
153 closeDone: make(chan struct{}),
154 invalidated: make(chan struct{}),
155 integrityDone: make(chan struct{}),
156 readLeases: map[*ReadLease]struct{}{},
157 discoveryStart: make(chan struct{}),
158 status: Status{State: StateOpening, Path: opts.Path},
159 }
160 handle, err := openCatalogProjection(ctx, opts)
161 if err != nil {
162 return nil, err
163 }
164 c.db = handle.DB
165 if err := setCatalogRevisionFloor(ctx, c.db, opts.RevisionFloor); err != nil {
166 _ = c.db.Close()
167 return nil, err
168 }
169 c.status.Mode = Mode(handle.Status.Mode)
170 c.status.State = State(handle.Status.State)
171 if c.status.State == "" {
172 c.status.State = StateReady
173 }
174 if c.status.Mode == ModeMemory {
175 c.status.Path = ""
176 } else {
177 c.status.Path = handle.Status.Path
178 }
179 c.status.LastError = handle.Status.LastError
180 c.status.QuarantinedPath = handle.Status.QuarantinedPath
181 if err := c.loadStatus(ctx); err != nil {
182 _ = c.db.Close()
183 return nil, err
184 }
185 if c.readable() != nil {
186 _ = c.db.Close()
187 return nil, c.invalidReason
188 }
189 if !opts.DisableRepair {
190 if err := c.resetRepairSchedule(ctx); err != nil {
191 _ = c.db.Close()
192 return nil, err
193 }
194 c.refreshCounts(ctx)
195 }
196 c.testRepairSessionHook = opts.repairSession
197 c.workerCtx, c.workerCancel = context.WithCancel(context.Background())
198 if !opts.StartPaused {
199 c.ResumeDiscovery()
200 }
201 if opts.MetadataOnly {
202 if err := c.loadReconcileJournal(ctx); err != nil {
203 c.workerCancel()
204 _ = c.db.Close()
205 return nil, err
206 }
207 }
208 c.workers.Add(1)
209 go c.writerLoop()
210 c.workers.Add(1)
211 go c.reconcileLoop()
212 c.workers.Add(1)
213 go c.sessionPathLoop()
214 if !opts.DisableRepair {
215 c.workers.Add(1)
216 go c.repairLoop()
217 c.enqueuePersistedRepairs(ctx)
218 }
219 if opts.DeferredMetadataIntegrity {
220 c.workers.Add(1)
221 go c.verifyMetadataIntegrity()
222 } else {
223 close(c.integrityDone)
224 }
225 return c, nil
226 }
227
228 func (c *Catalog) refreshCounts(ctx context.Context) {
229 if c == nil || c.db == nil {
230 return
231 }
232 if c.opts.MetadataOnly {
233 // Startup and each discovery batch may read root progress, but must
234 // not aggregate every session merely to draw the loading indicator.
235 var indexed, total int64
236 if err := c.readDB(ctx).QueryRowContext(ctx, `SELECT COALESCE(SUM(indexed),0),COALESCE(SUM(total),0) FROM catalog_directories`).Scan(&indexed, &total); err == nil {
237 c.statusMu.Lock()
238 c.status.Indexed, c.status.Total = indexed, max(total, indexed)
239 c.status.SourceCount = max(total, indexed)
240 c.status.Revision = c.revision.Load()
241 c.statusMu.Unlock()
242 }
243 return
244 }
245 var indexed, pending, total, physical, logical, groups, branches, diverged, cleanup int64
246 var active, deferred, blocked int64
247 var nextRepair sql.NullInt64
248 err := c.readDB(ctx).QueryRowContext(ctx, `SELECT
249 COUNT(*),
250 COALESCE(SUM(CASE WHEN turns_state='unknown' THEN 1 ELSE 0 END),0),
251 (SELECT COALESCE(SUM(total),0) FROM catalog_directories),
252 COALESCE(SUM(CASE WHEN missing_since=0 THEN 1 ELSE 0 END),0),
253 (SELECT COUNT(*) FROM catalog_topics),
254 COUNT(DISTINCT CASE WHEN recovered=1 AND recovery_group_id<>'' AND missing_since=0 THEN recovery_group_id END),
255 COALESCE(SUM(CASE WHEN recovered=1 AND missing_since=0 THEN 1 ELSE 0 END),0),
256 COALESCE(SUM(CASE WHEN recovered=1 AND recovery_role='diverged' AND missing_since=0 THEN 1 ELSE 0 END),0),
257 COALESCE(SUM(CASE WHEN recovered=1 AND recovery_role='covered_copy' AND missing_since=0 THEN 1 ELSE 0 END),0),
258 COALESCE(SUM(CASE WHEN turns_state='unknown' AND repair_state IN ('pending','active') THEN 1 ELSE 0 END),0),
259 COALESCE(SUM(CASE WHEN turns_state='unknown' AND repair_state='deferred' THEN 1 ELSE 0 END),0),
260 COALESCE(SUM(CASE WHEN turns_state='unknown' AND repair_state='blocked' THEN 1 ELSE 0 END),0),
261 MIN(CASE WHEN turns_state='unknown' AND repair_state='deferred' THEN repair_retry_at END)
262 FROM catalog_sessions`).Scan(&indexed, &pending, &total, &physical, &logical, &groups, &branches,
263 &diverged, &cleanup, &active, &deferred, &blocked, &nextRepair)
264 if err != nil {
265 return
266 }
267 errorKinds := map[string]int64{}
268 if rows, queryErr := c.readDB(ctx).QueryContext(ctx, `SELECT repair_error_kind,COUNT(*) FROM catalog_sessions
269 WHERE turns_state='unknown' AND repair_error_kind<>'' GROUP BY repair_error_kind`); queryErr == nil {
270 for rows.Next() {
271 var kind string
272 var count int64
273 if rows.Scan(&kind, &count) == nil {
274 errorKinds[kind] = count
275 }
276 }
277 _ = rows.Close()
278 }
279 c.statusMu.Lock()
280 c.status.Indexed = indexed
281 c.status.Total = total
282 c.status.RepairPending = pending
283 c.status.RepairActive = active
284 c.status.RepairDeferred = deferred
285 c.status.RepairBlocked = blocked
286 c.status.NextRepairAt = 0
287 if nextRepair.Valid {
288 c.status.NextRepairAt = nextRepair.Int64
289 }
290 c.status.RepairErrorKinds = errorKinds
291 c.status.PhysicalSessions = physical
292 c.status.LogicalSessions = logical
293 c.status.RecoveryGroups = groups
294 c.status.RecoveryBranches = branches
295 c.status.RecoveryDiverged = diverged
296 c.status.CleanupEligible = cleanup
297 c.status.SourceCount = total
298 c.status.Revision = c.revision.Load()
299 c.statusMu.Unlock()
300 }
301
302 func (c *Catalog) markRepair(reason string, at int64) {
303 if c == nil || strings.TrimSpace(reason) == "" {
304 return
305 }
306 if at <= 0 {
307 at = time.Now().UnixMilli()
308 }
309 c.statusMu.Lock()
310 c.status.RepairReason = strings.TrimSpace(reason)
311 c.status.LastRepairAt = at
312 c.statusMu.Unlock()
313 }
314
315 // MarkRepairReason records a lifecycle-level repair cause (for example, a
316 // clean index-generation cutover) without touching the authoritative session
317 // files. Integrity checks use the internal helper so they can attach their
318 // timestamp at the point of detection.
319 func (c *Catalog) MarkRepairReason(reason string) {
320 if c == nil {
321 return
322 }
323 c.markRepair(reason, c.opts.Now().UnixMilli())
324 }
325
326 func normalizeScope(scope, root string) (string, string) {
327 if strings.TrimSpace(scope) != "project" {
328 return "global", ""
329 }
330 return "project", strings.TrimSpace(root)
331 }
332
333 func normalizeSessionRecord(record SessionRecord) SessionRecord {
334 record.Path = cleanCatalogAccessPath(record.Path)
335 if record.Directory == "" {
336 record.Directory = filepath.Dir(record.Path)
337 }
338 record.Directory = cleanCatalogAccessPath(record.Directory)
339 record.Scope, record.WorkspaceRoot = normalizeScope(record.Scope, record.WorkspaceRoot)
340 if record.TurnsState == "" {
341 record.TurnsState = TurnsUnknown
342 }
343 if record.Health == "" {
344 record.Health = HealthOK
345 }
346 return record
347 }
348
349 func (c *Catalog) pathKey(path string) string {
350 if c != nil && c.pathIdentity != nil {
351 return c.pathIdentity(path)
352 }
353 return PathIdentityKey(path)
354 }
355
356 func (c *Catalog) workspaceRootKey(scope, root string) string {
357 scope, root = normalizeScope(scope, root)
358 if scope != "project" || root == "" {
359 return ""
360 }
361 return c.pathKey(root)
362 }
363
364 // queuePathKey is intentionally lexical. Save observers call the enqueue APIs
365 // synchronously, so filesystem probes (EvalSymlinks/platform case detection)
366 // belong to background workers and the SQLite uniqueness boundary.
367 func queuePathKey(path string) string {
368 return cleanCatalogAccessPath(path)
369 }
370
371 func (c *Catalog) EnqueueSession(record SessionRecord) bool {
372 if c == nil {
373 return false
374 }
375 record = normalizeSessionRecord(record)
376 record.enqueueSequence = c.mutationSeq.Add(1)
377 key := queuePathKey(record.Path)
378 if key == "" {
379 return false
380 }
381 c.writeMu.Lock()
382 if _, loaded := c.writeQueued[key]; loaded {
383 c.writeQueued[key] = record
384 c.writeMu.Unlock()
385 return true
386 }
387 c.writeQueued[key] = record
388 select {
389 case <-c.stop:
390 delete(c.writeQueued, key)
391 c.writeMu.Unlock()
392 return false
393 case c.writeCh <- key:
394 c.writeMu.Unlock()
395 return true
396 default:
397 delete(c.writeQueued, key)
398 c.writeMu.Unlock()
399 return false
400 }
401 }
402
403 func (c *Catalog) takeQueuedWrite(path string) (SessionRecord, bool) {
404 c.writeMu.Lock()
405 defer c.writeMu.Unlock()
406 record, ok := c.writeQueued[path]
407 if ok {
408 delete(c.writeQueued, path)
409 }
410 return record, ok
411 }
412
413 func (c *Catalog) writerLoop() {
414 defer c.workers.Done()
415 ticker := time.NewTicker(20 * time.Millisecond)
416 defer ticker.Stop()
417 pending := map[string]SessionRecord{}
418 flush := func() {
419 if len(pending) == 0 {
420 return
421 }
422 records := make([]SessionRecord, 0, len(pending))
423 for _, record := range pending {
424 records = append(records, record)
425 }
426 pending = map[string]SessionRecord{}
427 ctx, cancel := context.WithTimeout(c.workerCtx, time.Second)
428 _ = c.upsertSessions(ctx, records, nil, "write")
429 cancel()
430 }
431 for {
432 select {
433 case path := <-c.writeCh:
434 if record, ok := c.takeQueuedWrite(path); ok {
435 pending[path] = record
436 }
437 if len(pending) >= 64 {
438 flush()
439 }
440 case <-ticker.C:
441 flush()
442 case <-c.stop:
443 for {
444 select {
445 case path := <-c.writeCh:
446 if record, ok := c.takeQueuedWrite(path); ok {
447 pending[path] = record
448 }
449 default:
450 flush()
451 return
452 }
453 }
454 }
455 }
456 }
457
458 func (c *Catalog) recomputeTopic(ctx context.Context, tx *sql.Tx, key TopicKey) error {
459 key.Scope, key.WorkspaceRoot = normalizeScope(key.Scope, key.WorkspaceRoot)
460 rootKey := key.workspaceKey
461 if key.Scope == "project" && rootKey == "" {
462 rootKey = c.workspaceRootKey(key.Scope, key.WorkspaceRoot)
463 }
464 var count int
465 if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM catalog_sessions WHERE scope=? AND workspace_root_key=? AND topic_id=?`, key.Scope, rootKey, key.TopicID).Scan(&count); err != nil {
466 return err
467 }
468 if count == 0 {
469 _, err := tx.ExecContext(ctx, `DELETE FROM catalog_topics WHERE scope=? AND workspace_root_key=? AND topic_id=?`, key.Scope, rootKey, key.TopicID)
470 return err
471 }
472 if err := removeRemappedTopicIdentity(ctx, tx, key, rootKey); err != nil {
473 return err
474 }
475 // Covered copies skip turn/health totals but still update recency. Adopted
476 // branches are alternate continuations, so preserve the pre-catalog contract:
477 // max(sum(normal turns), max(adopted recovery turns)).
478 _, err := tx.ExecContext(ctx, `INSERT INTO catalog_topics(
479 scope,workspace_root,workspace_root_key,topic_id,title,turns,turns_state,created_at,
480 last_activity_at,recovery_state,recovery_branch_count,
481 recovery_unresolved_count,recovery_cleanup_eligible_count,health
482 ) SELECT ?,?,?,?,
483 COALESCE(NULLIF((SELECT COALESCE(NULLIF(topic_title,''), preview, '')
484 FROM catalog_sessions WHERE scope=? AND workspace_root_key=? AND topic_id=?
485 ORDER BY recovery_copy ASC, last_activity_at DESC, path ASC LIMIT 1),''), ?),
486 MAX(
487 COALESCE(SUM(CASE WHEN recovery_copy=0 AND recovered=0 AND turns_state='valid' THEN turns ELSE 0 END),0),
488 COALESCE(MAX(CASE WHEN recovery_copy=0 AND recovered=1 AND turns_state='valid' THEN turns ELSE 0 END),0)
489 ),
490 CASE WHEN SUM(CASE WHEN recovery_copy=0 AND turns_state='corrupt' THEN 1 ELSE 0 END)>0 THEN 'corrupt'
491 WHEN SUM(CASE WHEN recovery_copy=0 AND turns_state='unknown' THEN 1 ELSE 0 END)>0 THEN 'unknown'
492 WHEN SUM(CASE WHEN recovery_copy=0 THEN 1 ELSE 0 END)=0 THEN 'valid'
493 ELSE 'valid' END,
494 COALESCE(MIN(NULLIF(created_at,0)),0), COALESCE(MAX(last_activity_at),0),
495 CASE WHEN SUM(CASE WHEN recovered=1 AND recovery_role='preferred' THEN 1 ELSE 0 END)>0 THEN 'preferred'
496 WHEN SUM(CASE WHEN recovered=1 AND recovery_role='diverged' THEN 1 ELSE 0 END)>0 THEN 'diverged'
497 WHEN SUM(CASE WHEN recovered=1 AND recovery_role='adopted' THEN 1 ELSE 0 END)>0 THEN 'adopted'
498 WHEN SUM(CASE WHEN recovery_copy=0 THEN 1 ELSE 0 END)=0 THEN 'recovery_only' ELSE '' END,
499 SUM(CASE WHEN recovered=1 THEN 1 ELSE 0 END),
500 CASE WHEN SUM(CASE WHEN recovered=1 AND recovery_role='preferred' THEN 1 ELSE 0 END)>0 THEN 0
501 ELSE SUM(CASE WHEN recovered=1 AND recovery_role='diverged' THEN 1 ELSE 0 END) END,
502 SUM(CASE WHEN recovered=1 AND recovery_role='covered_copy' THEN 1 ELSE 0 END),
503 CASE WHEN SUM(CASE WHEN recovery_copy=0 AND health='corrupt' THEN 1 ELSE 0 END)>0 THEN 'corrupt'
504 WHEN SUM(CASE WHEN recovery_copy=0 AND health='missing' THEN 1 ELSE 0 END)>0 THEN 'missing'
505 ELSE 'ok' END
506 FROM catalog_sessions WHERE scope=? AND workspace_root_key=? AND topic_id=?
507 ON CONFLICT(scope,workspace_root_key,topic_id) DO UPDATE SET
508 title=excluded.title, turns=excluded.turns, turns_state=excluded.turns_state,
509 created_at=excluded.created_at, last_activity_at=excluded.last_activity_at,
510 recovery_state=excluded.recovery_state,
511 recovery_branch_count=excluded.recovery_branch_count,
512 recovery_unresolved_count=excluded.recovery_unresolved_count,
513 recovery_cleanup_eligible_count=excluded.recovery_cleanup_eligible_count,
514 health=excluded.health`,
515 key.Scope, key.WorkspaceRoot, rootKey, key.TopicID,
516 key.Scope, rootKey, key.TopicID, key.TopicID,
517 key.Scope, rootKey, key.TopicID)
518 return err
519 }
520
521 func boolToInt(value bool) int {
522 if value {
523 return 1
524 }
525 return 0
526 }
527
528 func bumpRevision(ctx context.Context, tx *sql.Tx) (uint64, error) {
529 if _, err := tx.ExecContext(ctx, `UPDATE catalog_state SET revision=revision+1 WHERE id=1`); err != nil {
530 return 0, err
531 }
532 var revision uint64
533 if err := tx.QueryRowContext(ctx, `SELECT revision FROM catalog_state WHERE id=1`).Scan(&revision); err != nil {
534 return 0, err
535 }
536 return revision, nil
537 }
538
539 func (c *Catalog) publishRevision(revision uint64, roots []string, reason string) {
540 c.rememberRevision(revision)
541 if c.readable() != nil {
542 return
543 }
544 if c.opts.OnRevision != nil {
545 c.opts.OnRevision(revision, c.registeredRevisionRoots(roots), reason)
546 }
547 }
548
549 func (c *Catalog) registeredRevisionRoots(roots []string) []string {
550 out := make([]string, 0, len(roots))
551 seen := make(map[string]struct{}, len(roots))
552 for _, root := range roots {
553 _, root = normalizeScope("project", root)
554 rootKey := c.workspaceRootKey("project", root)
555 if _, ok := seen[rootKey]; ok {
556 continue
557 }
558 seen[rootKey] = struct{}{}
559 registered := root
560 if rootKey != "" && c.db != nil {
561 var candidate string
562 if err := c.db.QueryRowContext(context.Background(), `SELECT workspace_root FROM catalog_projects
563 WHERE scope='project' AND workspace_root_key=?`, rootKey).Scan(&candidate); err == nil && candidate != "" {
564 registered = candidate
565 }
566 }
567 out = append(out, registered)
568 }
569 return out
570 }
571
572 func (c *Catalog) rememberRevision(revision uint64) {
573 c.revision.Store(revision)
574 c.statusMu.Lock()
575 c.status.Revision = revision
576 c.statusMu.Unlock()
577 }
578
579 func mapKeys(values map[string]struct{}) []string {
580 out := make([]string, 0, len(values))
581 for value := range values {
582 out = append(out, value)
583 }
584 return out
585 }
586
587 func (c *Catalog) listTopicSessionsByRootKey(ctx context.Context, key TopicKey, rootKey string) ([]SessionRecord, error) {
588 out := []SessionRecord{}
589 var cursor *sessionPageCursor
590 maxRows := MaxLimit
591 if view, _ := ctx.Value(readViewKey{}).(*readView); view != nil && view.owner == c {
592 maxRows = 2147483647
593 }
594 var capturedBytes int
595 for len(out) < maxRows {
596 where := `scope=? AND workspace_root_key=? AND topic_id=?`
597 args := []any{key.Scope, rootKey, key.TopicID}
598 if cursor != nil {
599 where += ` AND (last_activity_at<? OR (last_activity_at=? AND path>?))`
600 args = append(args, cursor.Activity, cursor.Activity, cursor.Path)
601 }
602 args = append(args, MaxLimit)
603 rows, err := c.readDB(ctx).QueryContext(ctx, `SELECT `+sessionSelectColumns+` FROM catalog_sessions
604 WHERE `+where+` ORDER BY last_activity_at DESC,path ASC LIMIT ?`, args...)
605 if err != nil {
606 return nil, err
607 }
608 rawCount := 0
609 var lastScanned SessionRecord
610 for rows.Next() {
611 record, err := scanSession(rows)
612 if err != nil {
613 _ = rows.Close()
614 return nil, err
615 }
616 rawCount++
617 lastScanned = record
618 if c.pathRemovedKey(record.pathKey, record.Path) {
619 continue
620 }
621 out = append(out, record)
622 if maxRows != MaxLimit {
623 encoded, err := json.Marshal(record)
624 if err != nil {
625 _ = rows.Close()
626 return nil, err
627 }
628 capturedBytes += len(encoded) + 64
629 if capturedBytes > 64<<20 {
630 _ = rows.Close()
631 return nil, fmt.Errorf("read snapshot resource limit: topic exceeds 64 MiB")
632 }
633 }
634 if len(out) == maxRows {
635 break
636 }
637 }
638 rowsErr := rows.Err()
639 _ = rows.Close()
640 if rowsErr != nil {
641 return nil, rowsErr
642 }
643 if len(out) == maxRows || rawCount < MaxLimit || rawCount == 0 {
644 break
645 }
646 cursor = &sessionPageCursor{Activity: lastScanned.LastActivityAt, Path: lastScanned.Path}
647 }
648 return out, nil
649 }
650
651 func (c *Catalog) GetTopic(ctx context.Context, key TopicKey) (TopicRecord, bool, error) {
652 key.Scope, key.WorkspaceRoot = normalizeScope(key.Scope, key.WorkspaceRoot)
653 key.TopicID = strings.TrimSpace(key.TopicID)
654 rootKey := c.workspaceRootKey(key.Scope, key.WorkspaceRoot)
655 item := TopicRecord{Sessions: []SessionRecord{}}
656 err := c.readDB(ctx).QueryRowContext(ctx, `SELECT scope,workspace_root,topic_id,title,title_source,pinned,
657 CASE WHEN metadata_present=1 THEN sort_order ELSE -1 END,
658 turns,turns_state,created_at,last_activity_at,recovery_state,recovery_branch_count,
659 recovery_unresolved_count,recovery_cleanup_eligible_count,health
660 FROM catalog_topics WHERE scope=? AND workspace_root_key=? AND topic_id=?`,
661 key.Scope, rootKey, key.TopicID).Scan(
662 &item.Scope, &item.WorkspaceRoot, &item.TopicID, &item.Title, &item.TitleSource,
663 &item.Pinned, &item.SortOrder, &item.Turns, &item.TurnsState,
664 &item.CreatedAt, &item.LastActivityAt, &item.RecoveryState, &item.RecoveryBranchCount,
665 &item.RecoveryUnresolvedCount, &item.RecoveryCleanupEligibleCount, &item.Health)
666 if errors.Is(err, sql.ErrNoRows) {
667 return item, false, nil
668 }
669 if err != nil {
670 return item, false, err
671 }
672 item.Sessions, err = c.listTopicSessionsByRootKey(ctx, key, rootKey)
673 if err != nil {
674 return TopicRecord{Sessions: []SessionRecord{}}, false, err
675 }
676 // Tombstone overlay: topic rows may lag behind RemoveSession while the
677 // durable DELETE waits on locks or a short caller context.
678 if len(item.Sessions) == 0 {
679 return TopicRecord{Sessions: []SessionRecord{}}, false, nil
680 }
681 hydrateTopicDisplay(&item)
682 return item, true, nil
683 }
684
685 func topicRepresentativePath(sessions []SessionRecord) string {
686 if path := OrdinaryContinuePath(sessions, ""); path != "" {
687 return path
688 }
689 preferred := PreferredOrdinarySessionPaths(sessions)
690 best := SessionRecord{}
691 found := false
692 for _, session := range sessions {
693 path := strings.TrimSpace(session.Path)
694 _, isPreferred := preferred[path]
695 if !session.OrdinaryVisible && !isPreferred && (session.Recovered || session.RecoveryCopy) {
696 continue
697 }
698 if !found || recoveryRank(session) > recoveryRank(best) ||
699 (recoveryRank(session) == recoveryRank(best) && session.LastActivityAt > best.LastActivityAt) {
700 best = session
701 found = true
702 }
703 }
704 if found {
705 return best.Path
706 }
707 if len(sessions) > 0 {
708 return sessions[0].Path
709 }
710 return ""
711 }
712
713 // EncodeTopicCursor builds an exclusive ListTopics keyset cursor after the
714 // given topic position. Desktop post-filters recovery-only rows and needs the
715 // same cursor shape catalog.ListTopics emits.
716 func EncodeTopicCursor(pinned int, lastActivityAt int64, topicID string) string {
717 return encodeCursor(pageCursor{Pinned: pinned, Activity: lastActivityAt, TopicID: topicID})
718 }
719
720 func EncodeTopicCursorBound(pinned int, lastActivityAt int64, topicID, binding string) string {
721 return encodeCursor(pageCursor{Pinned: pinned, Activity: lastActivityAt, TopicID: topicID, Binding: binding})
722 }
723
724 // EncodeOrderedTopicCursor builds a cursor for a workspace with explicit
725 // manual topic ordering. A negative sortOrder places metadata-free/runtime
726 // topics after every explicitly ranked topic in the same pinned bucket.
727 func EncodeOrderedTopicCursor(pinned, sortOrder int, lastActivityAt int64, topicID string) string {
728 manualSortOrder := int64(sortOrder)
729 if sortOrder < 0 {
730 manualSortOrder = unrankedTopicSortOrder
731 }
732 return encodeCursor(pageCursor{
733 Pinned: pinned, ManualOrder: true, SortOrder: manualSortOrder,
734 Activity: lastActivityAt, TopicID: topicID,
735 })
736 }
737
738 func EncodeOrderedTopicCursorBound(pinned, sortOrder int, lastActivityAt int64, topicID, binding string) string {
739 manualSortOrder := int64(sortOrder)
740 if sortOrder < 0 {
741 manualSortOrder = unrankedTopicSortOrder
742 }
743 return encodeCursor(pageCursor{
744 Pinned: pinned, ManualOrder: true, SortOrder: manualSortOrder,
745 Activity: lastActivityAt, TopicID: topicID, Binding: binding,
746 })
747 }
748
749 func encodeCursor(cursor pageCursor) string {
750 b, _ := json.Marshal(cursor)
751 return base64.RawURLEncoding.EncodeToString(b)
752 }
753
754 func decodeCursor(encoded string) (*pageCursor, error) {
755 if strings.TrimSpace(encoded) == "" {
756 return nil, nil
757 }
758 b, err := base64.RawURLEncoding.DecodeString(encoded)
759 if err != nil {
760 return nil, fmt.Errorf("invalid session catalog cursor: %w", err)
761 }
762 var cursor pageCursor
763 if err := json.Unmarshal(b, &cursor); err != nil || cursor.TopicID == "" {
764 return nil, errors.New("invalid session catalog cursor")
765 }
766 return &cursor, nil
767 }
768
769 func timeFilterCutoff(filter string, now time.Time) int64 {
770 var duration time.Duration
771 value := strings.TrimSpace(strings.ToLower(filter))
772 switch value {
773 case "day", "24h":
774 duration = 24 * time.Hour
775 case "week", "7d":
776 duration = 7 * 24 * time.Hour
777 case "month", "30d":
778 duration = 30 * 24 * time.Hour
779 default:
780 parsed, err := time.ParseDuration(value)
781 if err != nil || parsed <= 0 {
782 return 0
783 }
784 duration = parsed
785 }
786 return now.Add(-duration).UnixMilli()
787 }
788
788 lines GO