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