返回 DeepSeek-Reasonix
reconcile_journal.go
根目录 / internal / sessioncatalog / reconcile_journal.go
1 package sessioncatalog
2
3 import (
4 "context"
5 )
6
7 func (c *Catalog) persistReconcileTarget(target DirectoryTarget) error {
8 c.mutationMu.Lock()
9 defer c.mutationMu.Unlock()
10 _, err := c.db.ExecContext(c.workerCtx, `INSERT INTO catalog_pending_roots(path_key,path,scope,workspace_root,sequence)
11 VALUES(?,?,?,?,?) ON CONFLICT(path_key) DO UPDATE SET path=excluded.path,scope=excluded.scope,
12 workspace_root=excluded.workspace_root,sequence=excluded.sequence
13 WHERE excluded.sequence>catalog_pending_roots.sequence`,
14 queuePathKey(target.Path), target.Path, target.Scope, target.WorkspaceRoot, target.mutationSeq)
15 return err
16 }
17
18 func (c *Catalog) settleReconcileTarget(target DirectoryTarget) {
19 c.mutationMu.Lock()
20 defer c.mutationMu.Unlock()
21 _, _ = c.db.ExecContext(c.workerCtx, `DELETE FROM catalog_pending_roots WHERE path_key=? AND sequence<=?`, queuePathKey(target.Path), target.mutationSeq)
22 }
23
24 // Startup reads only dirty root identities. No directory is enumerated and no
25 // transcript is opened here. Failed and interrupted scans keep their journal.
26 func (c *Catalog) loadReconcileJournal(ctx context.Context) error {
27 rows, err := c.db.QueryContext(ctx, `SELECT path,scope,workspace_root,sequence FROM catalog_pending_roots`)
28 if err != nil {
29 return err
30 }
31 defer rows.Close()
32 for rows.Next() {
33 var target DirectoryTarget
34 if err := rows.Scan(&target.Path, &target.Scope, &target.WorkspaceRoot, &target.mutationSeq); err != nil {
35 return err
36 }
37 if target.mutationSeq > c.mutationSeq.Load() {
38 c.mutationSeq.Store(target.mutationSeq)
39 }
40 key := queuePathKey(target.Path)
41 c.reconcileQueued.Store(key, target)
42 c.reconcileDirty[key] = target
43 }
44 return rows.Err()
45 }
46
46 lines GO