| 1 | package sessioncatalog |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "database/sql" |
| 6 | "path/filepath" |
| 7 | "strings" |
| 8 | ) |
| 9 | |
| 10 | // SyncMetadata projects desktop project/topic registries. Advisory metadata |
| 11 | // catalogs publish bounded slices; cancellation retains unvisited membership. |
| 12 | // Other modes retain atomic refreshes. Neither removes session-derived topics: |
| 13 | // older writers may have saved sidecars not yet reflected in the registry. |
| 14 | func (c *Catalog) SyncMetadata(ctx context.Context, projects []ProjectRecord, topics []TopicMetadata) (result error) { |
| 15 | if c == nil || c.db == nil { |
| 16 | return nil |
| 17 | } |
| 18 | defer func() { c.observeDatabaseError(result) }() |
| 19 | ctx, cancel := context.WithCancel(ctx) |
| 20 | defer cancel() |
| 21 | if c.workerCtx != nil { |
| 22 | stop := context.AfterFunc(c.workerCtx, cancel) |
| 23 | defer stop() |
| 24 | } |
| 25 | c.metadataSyncOnce.Do(func() { c.metadataSyncGate = make(chan struct{}, 1) }) |
| 26 | select { |
| 27 | case c.metadataSyncGate <- struct{}{}: |
| 28 | defer func() { <-c.metadataSyncGate }() |
| 29 | case <-ctx.Done(): |
| 30 | return ctx.Err() |
| 31 | } |
| 32 | if c.opts.MetadataOnly { |
| 33 | return c.syncMetadataIncremental(ctx, projects, topics) |
| 34 | } |
| 35 | c.mutationMu.Lock() |
| 36 | defer c.mutationMu.Unlock() |
| 37 | tx, err := c.db.BeginTx(ctx, nil) |
| 38 | if err != nil { |
| 39 | return err |
| 40 | } |
| 41 | roots := map[string]struct{}{} |
| 42 | if _, err := tx.ExecContext(ctx, `DELETE FROM catalog_projects`); err != nil { |
| 43 | _ = tx.Rollback() |
| 44 | return err |
| 45 | } |
| 46 | if _, err := tx.ExecContext(ctx, `UPDATE catalog_topics SET metadata_present=0`); err != nil { |
| 47 | _ = tx.Rollback() |
| 48 | return err |
| 49 | } |
| 50 | for _, project := range projects { |
| 51 | project.Scope, project.WorkspaceRoot = normalizeScope(project.Scope, project.WorkspaceRoot) |
| 52 | if _, err := tx.ExecContext(ctx, `INSERT INTO catalog_projects( |
| 53 | scope,workspace_root,workspace_root_key,title,color,pinned,sort_order,updated_at |
| 54 | ) VALUES(?,?,?,?,?,?,?,?) ON CONFLICT(scope,workspace_root_key) DO UPDATE SET |
| 55 | workspace_root=excluded.workspace_root,title=excluded.title,color=excluded.color,pinned=excluded.pinned, |
| 56 | sort_order=excluded.sort_order,updated_at=excluded.updated_at`, |
| 57 | project.Scope, project.WorkspaceRoot, c.workspaceRootKey(project.Scope, project.WorkspaceRoot), project.Title, project.Color, |
| 58 | project.Pinned, project.SortOrder, c.opts.Now().UnixMilli()); err != nil { |
| 59 | _ = tx.Rollback() |
| 60 | return err |
| 61 | } |
| 62 | roots[project.WorkspaceRoot] = struct{}{} |
| 63 | } |
| 64 | for _, topic := range topics { |
| 65 | topic.Scope, topic.WorkspaceRoot = normalizeScope(topic.Scope, topic.WorkspaceRoot) |
| 66 | if strings.TrimSpace(topic.TopicID) == "" { |
| 67 | continue |
| 68 | } |
| 69 | skip, err := c.skipFoldedRecoveryShell(ctx, tx, topic) |
| 70 | if err != nil { |
| 71 | _ = tx.Rollback() |
| 72 | return err |
| 73 | } |
| 74 | if skip { |
| 75 | continue |
| 76 | } |
| 77 | if err := c.upsertTopicMetadata(ctx, tx, topic); err != nil { |
| 78 | _ = tx.Rollback() |
| 79 | return err |
| 80 | } |
| 81 | roots[topic.WorkspaceRoot] = struct{}{} |
| 82 | } |
| 83 | if _, err := tx.ExecContext(ctx, `DELETE FROM catalog_topics WHERE `+orphanMetadataPredicate); err != nil { |
| 84 | _ = tx.Rollback() |
| 85 | return err |
| 86 | } |
| 87 | revision, err := bumpRevision(ctx, tx) |
| 88 | if err != nil { |
| 89 | _ = tx.Rollback() |
| 90 | return err |
| 91 | } |
| 92 | if err := tx.Commit(); err != nil { |
| 93 | return err |
| 94 | } |
| 95 | c.publishRevision(revision, mapKeys(roots), "metadata") |
| 96 | return nil |
| 97 | } |
| 98 | |
| 99 | // skipFoldedRecoveryShell reports whether SyncMetadata must not (re)create a |
| 100 | // metadata topic shell for a folded recovery copy. While a directory scan is |
| 101 | // pending, the copy's rows may still sit under their pre-reanchor topic; once |
| 102 | // lineage projection re-anchors them onto the canonical row, re-creating this |
| 103 | // shell from the registry would re-list the copy as a separate sidebar session |
| 104 | // (#8525/#8551). Explicitly pinned topics survive: the user asked for that row. |
| 105 | func (c *Catalog) skipFoldedRecoveryShell(ctx context.Context, tx *sql.Tx, topic TopicMetadata) (bool, error) { |
| 106 | if topic.Pinned { |
| 107 | return false, nil |
| 108 | } |
| 109 | return c.foldedRecoveryShellHasCanonical(ctx, tx, topic.Scope, topic.WorkspaceRoot, topic.TopicID) |
| 110 | } |
| 111 | |
| 112 | // upsertTopicMetadata applies one registry topic. It inherits live session |
| 113 | // aggregates when present so a metadata-only insert does not publish |
| 114 | // last_activity_at=0 / turns_state=valid and reorder the sidebar ahead of (or |
| 115 | // instead of) the authoritative session rows. |
| 116 | func (c *Catalog) upsertTopicMetadata(ctx context.Context, tx *sql.Tx, topic TopicMetadata) error { |
| 117 | rootKey := c.workspaceRootKey(topic.Scope, topic.WorkspaceRoot) |
| 118 | if err := removeRemappedTopicIdentity(ctx, tx, TopicKey{ |
| 119 | Scope: topic.Scope, WorkspaceRoot: topic.WorkspaceRoot, TopicID: topic.TopicID, |
| 120 | }, rootKey); err != nil { |
| 121 | return err |
| 122 | } |
| 123 | _, err := tx.ExecContext(ctx, `INSERT INTO catalog_topics( |
| 124 | scope,workspace_root,workspace_root_key,topic_id,title,title_source,pinned,sort_order, |
| 125 | turns,turns_state,created_at,last_activity_at,recovery_state,health,metadata_present |
| 126 | ) |
| 127 | SELECT ?,?,?,?,?,?,?,?, |
| 128 | COALESCE((SELECT MAX( |
| 129 | COALESCE(SUM(CASE WHEN recovery_copy=0 AND recovered=0 AND turns_state='valid' THEN turns ELSE 0 END),0), |
| 130 | COALESCE(MAX(CASE WHEN recovery_copy=0 AND recovered=1 AND turns_state='valid' THEN turns ELSE 0 END),0) |
| 131 | ) FROM catalog_sessions WHERE scope=? AND workspace_root_key=? AND topic_id=?),0), |
| 132 | COALESCE((SELECT CASE |
| 133 | WHEN SUM(CASE WHEN recovery_copy=0 AND turns_state='corrupt' THEN 1 ELSE 0 END)>0 THEN 'corrupt' |
| 134 | WHEN SUM(CASE WHEN recovery_copy=0 AND turns_state='unknown' THEN 1 ELSE 0 END)>0 THEN 'unknown' |
| 135 | WHEN SUM(CASE WHEN recovery_copy=0 THEN 1 ELSE 0 END)=0 AND COUNT(*)>0 THEN 'valid' |
| 136 | WHEN COUNT(*)=0 THEN 'unknown' |
| 137 | ELSE 'valid' END |
| 138 | FROM catalog_sessions WHERE scope=? AND workspace_root_key=? AND topic_id=?),'unknown'), |
| 139 | COALESCE(NULLIF(?,0),(SELECT MIN(NULLIF(created_at,0)) FROM catalog_sessions |
| 140 | WHERE scope=? AND workspace_root_key=? AND topic_id=?),0), |
| 141 | COALESCE((SELECT MAX(last_activity_at) FROM catalog_sessions |
| 142 | WHERE scope=? AND workspace_root_key=? AND topic_id=?),0), |
| 143 | COALESCE((SELECT CASE |
| 144 | WHEN COUNT(*)>0 AND SUM(CASE WHEN recovery_copy=0 THEN 1 ELSE 0 END)=0 THEN 'recovery_only' |
| 145 | ELSE '' END |
| 146 | FROM catalog_sessions WHERE scope=? AND workspace_root_key=? AND topic_id=?),''), |
| 147 | COALESCE((SELECT CASE |
| 148 | WHEN SUM(CASE WHEN recovery_copy=0 AND health='corrupt' THEN 1 ELSE 0 END)>0 THEN 'corrupt' |
| 149 | WHEN SUM(CASE WHEN recovery_copy=0 AND health='missing' THEN 1 ELSE 0 END)>0 THEN 'missing' |
| 150 | ELSE 'ok' END |
| 151 | FROM catalog_sessions WHERE scope=? AND workspace_root_key=? AND topic_id=?),'ok'), |
| 152 | 1 |
| 153 | ON CONFLICT(scope,workspace_root_key,topic_id) DO UPDATE SET |
| 154 | workspace_root=excluded.workspace_root, |
| 155 | title=COALESCE(NULLIF(excluded.title,''), |
| 156 | NULLIF((SELECT s.topic_title FROM catalog_sessions s |
| 157 | WHERE s.scope=excluded.scope AND s.workspace_root_key=excluded.workspace_root_key |
| 158 | AND s.topic_id=excluded.topic_id |
| 159 | ORDER BY s.recovery_copy ASC,s.last_activity_at DESC,s.path ASC LIMIT 1),''), |
| 160 | catalog_topics.title), |
| 161 | title_source=excluded.title_source,pinned=excluded.pinned, |
| 162 | sort_order=excluded.sort_order,metadata_present=1, |
| 163 | created_at=CASE WHEN excluded.created_at>0 THEN excluded.created_at ELSE catalog_topics.created_at END, |
| 164 | last_activity_at=CASE WHEN excluded.last_activity_at>catalog_topics.last_activity_at |
| 165 | THEN excluded.last_activity_at ELSE catalog_topics.last_activity_at END, |
| 166 | turns=CASE WHEN excluded.turns>0 THEN excluded.turns ELSE catalog_topics.turns END, |
| 167 | turns_state=CASE WHEN excluded.turns_state<>'' AND excluded.turns_state<>'unknown' |
| 168 | THEN excluded.turns_state ELSE catalog_topics.turns_state END, |
| 169 | recovery_state=excluded.recovery_state`, |
| 170 | topic.Scope, topic.WorkspaceRoot, rootKey, topic.TopicID, topic.Title, |
| 171 | topic.TitleSource, topic.Pinned, topic.SortOrder, |
| 172 | topic.Scope, rootKey, topic.TopicID, |
| 173 | topic.Scope, rootKey, topic.TopicID, |
| 174 | topic.CreatedAt, topic.Scope, rootKey, topic.TopicID, |
| 175 | topic.Scope, rootKey, topic.TopicID, |
| 176 | topic.Scope, rootKey, topic.TopicID, |
| 177 | topic.Scope, rootKey, topic.TopicID) |
| 178 | return err |
| 179 | } |
| 180 | |
| 181 | // foldedRecoveryShellHasCanonical reports whether topicID currently projects |
| 182 | // only recovery sessions whose lineage already has an ordinary/canonical |
| 183 | // representative in the catalog, or was tombstoned by a lineage re-anchor. |
| 184 | // Such a topic is a folded recovery copy's shell: its conversation is already |
| 185 | // listed under the canonical row, so SyncMetadata must not (re)create a |
| 186 | // standalone topic for it. |
| 187 | // |
| 188 | // A canonical representative is either a group member flagged |
| 189 | // ordinary_visible/recovery_canonical, or the non-recovered group root (which |
| 190 | // carries no recovery_group_id of its own, so it is matched by path). |
| 191 | // Lineages with no canonical yet (unresolved, still scanning) are left alone. |
| 192 | func (c *Catalog) foldedRecoveryShellHasCanonical(ctx context.Context, tx *sql.Tx, scope, workspaceRoot, topicID string) (bool, error) { |
| 193 | rootKey := c.workspaceRootKey(scope, workspaceRoot) |
| 194 | var ordinary int |
| 195 | if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM catalog_sessions |
| 196 | WHERE scope=? AND workspace_root_key=? AND topic_id=? AND recovered=0 AND recovery_copy=0`, |
| 197 | scope, rootKey, topicID).Scan(&ordinary); err != nil { |
| 198 | return false, err |
| 199 | } |
| 200 | if ordinary > 0 { |
| 201 | return false, nil |
| 202 | } |
| 203 | var folded int |
| 204 | if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM catalog_folded_topics |
| 205 | WHERE scope=? AND workspace_root_key=? AND topic_id=?`, |
| 206 | scope, rootKey, topicID).Scan(&folded); err != nil { |
| 207 | return false, err |
| 208 | } |
| 209 | if folded > 0 { |
| 210 | return true, nil |
| 211 | } |
| 212 | rows, err := tx.QueryContext(ctx, `SELECT DISTINCT directory, recovery_group_id FROM catalog_sessions |
| 213 | WHERE scope=? AND workspace_root_key=? AND topic_id=? AND recovered=1 AND recovery_group_id<>''`, |
| 214 | scope, rootKey, topicID) |
| 215 | if err != nil { |
| 216 | return false, err |
| 217 | } |
| 218 | type groupRef struct { |
| 219 | directory string |
| 220 | id string |
| 221 | } |
| 222 | groups := []groupRef{} |
| 223 | for rows.Next() { |
| 224 | var group groupRef |
| 225 | if err := rows.Scan(&group.directory, &group.id); err != nil { |
| 226 | rows.Close() |
| 227 | return false, err |
| 228 | } |
| 229 | groups = append(groups, group) |
| 230 | } |
| 231 | if err := rows.Err(); err != nil { |
| 232 | rows.Close() |
| 233 | return false, err |
| 234 | } |
| 235 | rows.Close() |
| 236 | if len(groups) == 0 { |
| 237 | return false, nil |
| 238 | } |
| 239 | for _, group := range groups { |
| 240 | var canonical int |
| 241 | if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM catalog_sessions |
| 242 | WHERE scope=? AND workspace_root_key=? AND recovery_group_id=? AND (ordinary_visible=1 OR recovery_canonical=1)`, |
| 243 | scope, rootKey, group.id).Scan(&canonical); err != nil { |
| 244 | return false, err |
| 245 | } |
| 246 | if canonical > 0 { |
| 247 | return true, nil |
| 248 | } |
| 249 | rootPath := filepath.Join(group.directory, group.id+".jsonl") |
| 250 | var roots int |
| 251 | if err := tx.QueryRowContext(ctx, `SELECT COUNT(*) FROM catalog_sessions |
| 252 | WHERE path_key=? AND recovered=0 AND recovery_copy=0`, c.pathKey(rootPath)).Scan(&roots); err != nil { |
| 253 | return false, err |
| 254 | } |
| 255 | if roots > 0 { |
| 256 | return true, nil |
| 257 | } |
| 258 | } |
| 259 | return false, nil |
| 260 | } |
| 261 | |
| 262 | // rememberFoldedTopic tombstones a topic that lineage projection folded into a |
| 263 | // recovery lineage's canonical row. The tombstone is cleared automatically if |
| 264 | // a session is ever indexed under that topic id again. |
| 265 | func (c *Catalog) rememberFoldedTopic(ctx context.Context, tx *sql.Tx, key TopicKey, foldedAt int64) error { |
| 266 | if strings.TrimSpace(key.TopicID) == "" { |
| 267 | return nil |
| 268 | } |
| 269 | rootKey := c.workspaceRootKey(key.Scope, key.WorkspaceRoot) |
| 270 | if err := removeRemappedFoldedTopicIdentity(ctx, tx, key, rootKey); err != nil { |
| 271 | return err |
| 272 | } |
| 273 | _, err := tx.ExecContext(ctx, `INSERT INTO catalog_folded_topics(scope,workspace_root,workspace_root_key,topic_id,folded_at) |
| 274 | VALUES(?,?,?,?,?) ON CONFLICT(scope,workspace_root_key,topic_id) DO UPDATE SET |
| 275 | workspace_root=excluded.workspace_root,folded_at=excluded.folded_at`, |
| 276 | key.Scope, key.WorkspaceRoot, rootKey, key.TopicID, foldedAt) |
| 277 | return err |
| 278 | } |
| 279 | |
| 280 | // updateFoldedTopicTombstones maintains folded-topic tombstones around a |
| 281 | // session upsert: a session claiming a folded topic id makes it real again, |
| 282 | // and a recovered row moving topics tombstones the shell it left behind. |
| 283 | func (c *Catalog) updateFoldedTopicTombstones(ctx context.Context, tx *sql.Tx, previous TopicKey, record SessionRecord, now int64) error { |
| 284 | if record.TopicID != "" { |
| 285 | if _, err := tx.ExecContext(ctx, `DELETE FROM catalog_folded_topics WHERE scope=? AND workspace_root_key=? AND topic_id=?`, |
| 286 | record.Scope, c.workspaceRootKey(record.Scope, record.WorkspaceRoot), record.TopicID); err != nil { |
| 287 | return err |
| 288 | } |
| 289 | } |
| 290 | if record.Recovered && previous.TopicID != "" && previous.TopicID != record.TopicID { |
| 291 | return c.rememberFoldedTopic(ctx, tx, previous, now) |
| 292 | } |
| 293 | return nil |
| 294 | } |
| 295 |