| 1 | package sessioncatalog |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "database/sql" |
| 6 | "strings" |
| 7 | "time" |
| 8 | |
| 9 | "reasonix/internal/historywork" |
| 10 | ) |
| 11 | |
| 12 | // Registry input is already a complete caller-owned observation. Applying it |
| 13 | // is incremental: additions finish before retirement starts, and each slice |
| 14 | // releases the common writer. Cancellation never retires unvisited input. |
| 15 | // No second authoritative registry or persisted format is introduced. |
| 16 | func (c *Catalog) syncMetadataIncremental(ctx context.Context, projects []ProjectRecord, topics []TopicMetadata) error { |
| 17 | coordinator := c.opts.Maintenance |
| 18 | if coordinator == nil { |
| 19 | coordinator = &historywork.Coordinator{} |
| 20 | } |
| 21 | seenProjects := map[TopicKey]bool{} |
| 22 | seenTopics := map[TopicKey]bool{} |
| 23 | projectIndex, topicIndex := 0, 0 |
| 24 | err := c.metadataSlices(ctx, coordinator, func(ctx context.Context, tx *sql.Tx, budget *metadataSlice) (bool, error) { |
| 25 | if _, err := tx.ExecContext(ctx, metadataTopicIndex); err != nil { |
| 26 | return false, err |
| 27 | } |
| 28 | for projectIndex < len(projects) && budget.available() { |
| 29 | project := projects[projectIndex] |
| 30 | project.Scope, project.WorkspaceRoot = normalizeScope(project.Scope, project.WorkspaceRoot) |
| 31 | key := TopicKey{Scope: project.Scope, workspaceKey: c.workspaceRootKey(project.Scope, project.WorkspaceRoot)} |
| 32 | changed, err := c.projectRegistryRow(ctx, tx, project, key.workspaceKey) |
| 33 | if err != nil { |
| 34 | return false, err |
| 35 | } |
| 36 | seenProjects[key] = true |
| 37 | budget.record(project.WorkspaceRoot, changed, len(project.Title)+len(project.WorkspaceRoot)+len(project.Color)+64) |
| 38 | projectIndex++ |
| 39 | } |
| 40 | for topicIndex < len(topics) && budget.available() { |
| 41 | topic := topics[topicIndex] |
| 42 | topic.Scope, topic.WorkspaceRoot = normalizeScope(topic.Scope, topic.WorkspaceRoot) |
| 43 | if strings.TrimSpace(topic.TopicID) != "" { |
| 44 | key := TopicKey{Scope: topic.Scope, workspaceKey: c.workspaceRootKey(topic.Scope, topic.WorkspaceRoot), TopicID: topic.TopicID} |
| 45 | skip, err := c.skipFoldedRecoveryShell(ctx, tx, topic) |
| 46 | if err != nil { |
| 47 | return false, err |
| 48 | } |
| 49 | changed := false |
| 50 | if !skip { |
| 51 | changed, err = c.topicRegistryRow(ctx, tx, topic, key.workspaceKey) |
| 52 | if err != nil { |
| 53 | return false, err |
| 54 | } |
| 55 | seenTopics[key] = true |
| 56 | } |
| 57 | budget.record(topic.WorkspaceRoot, changed, len(topic.Title)+len(topic.TopicID)+len(topic.WorkspaceRoot)+64) |
| 58 | } else { |
| 59 | budget.record("", false, 1) |
| 60 | } |
| 61 | topicIndex++ |
| 62 | } |
| 63 | return projectIndex == len(projects) && topicIndex == len(topics), nil |
| 64 | }) |
| 65 | if err != nil { |
| 66 | return err |
| 67 | } |
| 68 | // Keyset pages cover only registered rows, including rows written by older |
| 69 | // versions. Session-derived topics are owned by source reconciliation. |
| 70 | if err := c.retireRegistryRows(ctx, coordinator, seenProjects, false); err != nil { |
| 71 | return err |
| 72 | } |
| 73 | return c.retireRegistryRows(ctx, coordinator, seenTopics, true) |
| 74 | } |
| 75 | |
| 76 | type metadataSlice struct { |
| 77 | started time.Time |
| 78 | entries int |
| 79 | bytes int64 |
| 80 | roots map[string]struct{} |
| 81 | } |
| 82 | |
| 83 | func (s *metadataSlice) available() bool { |
| 84 | // Always admit one item, even when scheduler latency consumed the slice. |
| 85 | return s.entries == 0 || s.entries < historywork.BatchEntries && s.bytes < historywork.BatchBytes && time.Since(s.started) < historywork.SliceDuration |
| 86 | } |
| 87 | |
| 88 | func (s *metadataSlice) record(root string, changed bool, bytes int) { |
| 89 | s.entries++ |
| 90 | s.bytes += int64(bytes) |
| 91 | if changed { |
| 92 | s.roots[root] = struct{}{} |
| 93 | } |
| 94 | } |
| 95 | |
| 96 | func (c *Catalog) metadataSlices(ctx context.Context, coordinator *historywork.Coordinator, step func(context.Context, *sql.Tx, *metadataSlice) (bool, error)) error { |
| 97 | for { |
| 98 | release, err := coordinator.BackgroundSlice(ctx, false) |
| 99 | if err != nil { |
| 100 | return err |
| 101 | } |
| 102 | budget := &metadataSlice{started: time.Now(), roots: map[string]struct{}{}} |
| 103 | done, err := c.commitMetadataSlice(ctx, budget, step) |
| 104 | release(budget.bytes) |
| 105 | if err != nil { |
| 106 | return err |
| 107 | } |
| 108 | if c.testMetadataSliceHook != nil { |
| 109 | c.testMetadataSliceHook(budget.entries) |
| 110 | } |
| 111 | if done { |
| 112 | return nil |
| 113 | } |
| 114 | } |
| 115 | } |
| 116 | |
| 117 | func (c *Catalog) commitMetadataSlice(ctx context.Context, budget *metadataSlice, step func(context.Context, *sql.Tx, *metadataSlice) (bool, error)) (bool, error) { |
| 118 | ctx, cancel := context.WithTimeout(ctx, 30*time.Second) |
| 119 | defer cancel() |
| 120 | c.mutationMu.Lock() |
| 121 | defer c.mutationMu.Unlock() |
| 122 | tx, err := c.db.BeginTx(ctx, nil) |
| 123 | if err != nil { |
| 124 | return false, err |
| 125 | } |
| 126 | defer func() { _ = tx.Rollback() }() |
| 127 | done, err := step(ctx, tx, budget) |
| 128 | if err != nil { |
| 129 | return false, err |
| 130 | } |
| 131 | var revision uint64 |
| 132 | if len(budget.roots) > 0 { |
| 133 | revision, err = bumpRevision(ctx, tx) |
| 134 | if err != nil { |
| 135 | return false, err |
| 136 | } |
| 137 | } |
| 138 | if err := tx.Commit(); err != nil { |
| 139 | return false, err |
| 140 | } |
| 141 | if revision != 0 { |
| 142 | c.publishRevision(revision, mapKeys(budget.roots), "metadata") |
| 143 | } |
| 144 | return done, nil |
| 145 | } |
| 146 |