返回 DeepSeek-Reasonix
metadata.go
根目录 / internal / sessioncatalog / metadata.go
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
282 lines GO