返回 DeepSeek-Reasonix
metadata_projection_rows.go
根目录 / internal / sessioncatalog / metadata_projection_rows.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "database/sql"
6 "errors"
7
8 "reasonix/internal/historywork"
9 )
10
11 func (c *Catalog) projectRegistryRow(ctx context.Context, tx *sql.Tx, p ProjectRecord, rootKey string) (bool, error) {
12 result, err := tx.ExecContext(ctx, `INSERT INTO catalog_projects(
13 scope,workspace_root,workspace_root_key,title,color,pinned,sort_order,updated_at)
14 VALUES(?,?,?,?,?,?,?,?) ON CONFLICT(scope,workspace_root_key) DO UPDATE SET
15 workspace_root=excluded.workspace_root,title=excluded.title,color=excluded.color,
16 pinned=excluded.pinned,sort_order=excluded.sort_order,updated_at=excluded.updated_at
17 WHERE workspace_root<>excluded.workspace_root OR title<>excluded.title OR color<>excluded.color
18 OR pinned<>excluded.pinned OR sort_order<>excluded.sort_order`,
19 p.Scope, p.WorkspaceRoot, rootKey, p.Title, p.Color, p.Pinned, p.SortOrder, c.opts.Now().UnixMilli())
20 return registryMutationChanged(result, err)
21 }
22
23 func (c *Catalog) topicRegistryRow(ctx context.Context, tx *sql.Tx, topic TopicMetadata, rootKey string) (bool, error) {
24 // Compare the actual projection, not a remembered input fingerprint: old
25 // writers and source mutations also own these cache rows. Empty titles and
26 // unspecified creation times retain the existing upsert semantics.
27 var same bool
28 err := tx.QueryRowContext(ctx, `SELECT metadata_present=1 AND workspace_root=?
29 AND title=CASE WHEN ?<>'' THEN ? ELSE COALESCE(NULLIF((
30 SELECT s.topic_title FROM catalog_sessions s WHERE s.scope=catalog_topics.scope
31 AND s.workspace_root_key=catalog_topics.workspace_root_key AND s.topic_id=catalog_topics.topic_id
32 ORDER BY s.recovery_copy ASC,s.last_activity_at DESC,s.path ASC LIMIT 1),''),title) END
33 AND title_source=? AND pinned=? AND sort_order=? AND (?=0 OR created_at=?)
34 FROM catalog_topics WHERE scope=? AND workspace_root_key=? AND topic_id=?`,
35 topic.WorkspaceRoot, topic.Title, topic.Title, topic.TitleSource, topic.Pinned, topic.SortOrder,
36 topic.CreatedAt, topic.CreatedAt, topic.Scope, rootKey, topic.TopicID).Scan(&same)
37 if err != nil && !errors.Is(err, sql.ErrNoRows) {
38 return false, err
39 }
40 if same {
41 return false, nil
42 }
43 return true, c.upsertTopicMetadata(ctx, tx, topic)
44 }
45
46 type registryRetirementRow struct {
47 key TopicKey
48 root string
49 }
50
51 func (c *Catalog) retireRegistryRows(ctx context.Context, coordinator *historywork.Coordinator, seen map[TopicKey]bool, topics bool) error {
52 var after TopicKey
53 var pending []registryRetirementRow
54 var exhausted bool
55 return c.metadataSlices(ctx, coordinator, func(ctx context.Context, tx *sql.Tx, budget *metadataSlice) (bool, error) {
56 if len(pending) == 0 {
57 var err error
58 pending, err = registryRetirementPage(ctx, tx, after, topics)
59 if err != nil {
60 return false, err
61 }
62 exhausted = len(pending) < historywork.BatchEntries
63 }
64 for len(pending) > 0 && budget.available() {
65 row := pending[0]
66 changed := false
67 if !seen[row.key] {
68 var err error
69 changed, err = retireRegistryRow(ctx, tx, row.key, topics)
70 if err != nil {
71 return false, err
72 }
73 }
74 budget.record(row.root, changed, len(row.root)+len(row.key.TopicID)+64)
75 after, pending = row.key, pending[1:]
76 }
77 return exhausted && len(pending) == 0, nil
78 })
79 }
80
81 func registryRetirementPage(ctx context.Context, tx *sql.Tx, after TopicKey, topics bool) ([]registryRetirementRow, error) {
82 query := `SELECT scope,workspace_root_key,'',workspace_root FROM catalog_projects
83 WHERE (scope,workspace_root_key)>(?,?) ORDER BY scope,workspace_root_key LIMIT ?`
84 args := []any{after.Scope, after.workspaceKey, historywork.BatchEntries}
85 if topics {
86 query = `SELECT scope,workspace_root_key,topic_id,workspace_root FROM catalog_topics
87 INDEXED BY idx_catalog_topics_registered_metadata WHERE metadata_present=1
88 AND (scope,workspace_root_key,topic_id)>(?,?,?)
89 ORDER BY scope,workspace_root_key,topic_id LIMIT ?`
90 args = []any{after.Scope, after.workspaceKey, after.TopicID, historywork.BatchEntries}
91 }
92 rows, err := tx.QueryContext(ctx, query, args...)
93 if err != nil {
94 return nil, err
95 }
96 defer rows.Close()
97 result := []registryRetirementRow{}
98 for rows.Next() {
99 var row registryRetirementRow
100 if err := rows.Scan(&row.key.Scope, &row.key.workspaceKey, &row.key.TopicID, &row.root); err != nil {
101 return nil, err
102 }
103 result = append(result, row)
104 }
105 return result, rows.Err()
106 }
107
108 func retireRegistryRow(ctx context.Context, tx *sql.Tx, key TopicKey, topics bool) (bool, error) {
109 if !topics {
110 return registryMutationChanged(tx.ExecContext(ctx, `DELETE FROM catalog_projects WHERE scope=? AND workspace_root_key=?`, key.Scope, key.workspaceKey))
111 }
112 changed, err := registryMutationChanged(tx.ExecContext(ctx, `UPDATE catalog_topics SET metadata_present=0
113 WHERE scope=? AND workspace_root_key=? AND topic_id=? AND metadata_present=1`, key.Scope, key.workspaceKey, key.TopicID))
114 if err != nil {
115 return false, err
116 }
117 _, err = tx.ExecContext(ctx, `DELETE FROM catalog_topics WHERE scope=? AND workspace_root_key=? AND topic_id=? AND `+orphanMetadataPredicate,
118 key.Scope, key.workspaceKey, key.TopicID)
119 return changed, err
120 }
121
122 func registryMutationChanged(result sql.Result, err error) (bool, error) {
123 if err != nil {
124 return false, err
125 }
126 rows, err := result.RowsAffected()
127 return rows > 0, err
128 }
129
129 lines GO