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