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