返回 DeepSeek-Reasonix
upsert.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "database/sql"
6 "errors"
7 )
8
9 func (c *Catalog) UpsertSession(ctx context.Context, record SessionRecord) error {
10 record = normalizeSessionRecord(record)
11 record.enqueueSequence = c.mutationSeq.Add(1)
12 return c.upsertSessions(ctx, []SessionRecord{record}, nil, "write")
13 }
14
15 func (c *Catalog) upsertSessions(ctx context.Context, records []SessionRecord, generations map[string]int64, reason string) error {
16 _, err := c.upsertSessionsWithNotification(ctx, records, generations, reason, true, upsertExactSource)
17 return err
18 }
19
20 func (c *Catalog) upsertExactPathSession(ctx context.Context, record SessionRecord) (bool, error) {
21 dirty, err := c.upsertSessionsWithNotification(ctx, []SessionRecord{record}, nil, "write", true, upsertExactSource)
22 return len(dirty) > 0, err
23 }
24
25 func (c *Catalog) upsertSessionsWithNotification(ctx context.Context, records []SessionRecord, generations map[string]int64, reason string, notify bool, mode sessionUpsertMode) (map[string]DirectoryTarget, error) {
26 dirtyDirectories := map[string]DirectoryTarget{}
27 if len(records) == 0 {
28 return dirtyDirectories, nil
29 }
30 c.mutationMu.Lock()
31 defer c.mutationMu.Unlock()
32 var err error
33 records, err = c.prepareUpsertRecords(ctx, records, dirtyDirectories, mode)
34 if err != nil || len(records) == 0 {
35 return dirtyDirectories, err
36 }
37 tx, err := c.db.BeginTx(ctx, nil)
38 if err != nil {
39 return dirtyDirectories, err
40 }
41 affected := map[TopicKey]struct{}{}
42 roots := map[string]struct{}{}
43 directoryGenerations := map[string]int64{}
44 changed := false
45 for _, raw := range records {
46 record := normalizeSessionRecord(raw)
47 pathKey := c.pathKey(record.Path)
48 directoryKey := c.pathKey(record.Directory)
49 remapped, err := removeRemappedSessionIdentity(ctx, tx, record.Path, pathKey)
50 if err != nil {
51 _ = tx.Rollback()
52 return dirtyDirectories, err
53 }
54 for _, key := range remapped {
55 affected[key] = struct{}{}
56 }
57 generation := int64(0)
58 if generations != nil {
59 generation = generations[record.Path]
60 } else if cached, ok := directoryGenerations[directoryKey]; ok {
61 generation = cached
62 } else {
63 _ = tx.QueryRowContext(ctx, `SELECT scan_generation FROM catalog_directories WHERE path_key=?`, directoryKey).Scan(&generation)
64 directoryGenerations[directoryKey] = generation
65 }
66 if c.opts.MetadataOnly && mode == upsertDirectoryProjection && record.metadataUnchanged && len(remapped) == 0 {
67 seen, err := refreshUnchangedMetadata(ctx, tx, record, pathKey, generation)
68 if err != nil {
69 _ = tx.Rollback()
70 return dirtyDirectories, err
71 }
72 if seen {
73 continue
74 }
75 }
76 var previous TopicKey
77 if err := tx.QueryRowContext(ctx, `SELECT scope,workspace_root,workspace_root_key,topic_id FROM catalog_sessions WHERE path_key=?`, pathKey).
78 Scan(&previous.Scope, &previous.WorkspaceRoot, &previous.workspaceKey, &previous.TopicID); err == nil && previous.TopicID != "" {
79 affected[previous] = struct{}{}
80 } else if err != nil && !errors.Is(err, sql.ErrNoRows) {
81 _ = tx.Rollback()
82 return dirtyDirectories, err
83 }
84 if err := c.upsertSessionRow(ctx, tx, record, pathKey, directoryKey, generation, mode); err != nil {
85 _ = tx.Rollback()
86 return dirtyDirectories, err
87 }
88 if record.TopicID != "" {
89 affected[TopicKey{Scope: record.Scope, WorkspaceRoot: record.WorkspaceRoot,
90 workspaceKey: c.workspaceRootKey(record.Scope, record.WorkspaceRoot), TopicID: record.TopicID}] = struct{}{}
91 }
92 if err := c.updateFoldedTopicTombstones(ctx, tx, previous, record, c.opts.Now().UnixMilli()); err != nil {
93 _ = tx.Rollback()
94 return dirtyDirectories, err
95 }
96 roots[record.WorkspaceRoot] = struct{}{}
97 changed = true
98 }
99 if !changed {
100 // Presence belongs to scan completion. An unchanged batch must not
101 // rewrite topic aggregates or make every sidebar refresh its snapshot.
102 return dirtyDirectories, tx.Commit()
103 }
104 for key := range affected {
105 if err := c.recomputeTopic(ctx, tx, key); err != nil {
106 _ = tx.Rollback()
107 return dirtyDirectories, err
108 }
109 }
110 revision, err := bumpRevision(ctx, tx)
111 if err != nil {
112 _ = tx.Rollback()
113 return dirtyDirectories, err
114 }
115 if err := tx.Commit(); err != nil {
116 return dirtyDirectories, err
117 }
118 if notify {
119 c.publishRevision(revision, mapKeys(roots), reason)
120 } else {
121 c.rememberRevision(revision)
122 }
123 c.refreshCounts(ctx)
124 return dirtyDirectories, nil
125 }
126
127 const sessionInsertSQL = `INSERT INTO catalog_sessions(
128 path,path_key,directory,directory_key,scope,workspace_root,workspace_root_key,topic_id,topic_title,custom_title,
129 created_at,last_activity_at,preview,turns,turns_state,recovered,
130 recovery_reason,recovery_digest,parent_id,recovery_copy,recovery_group_id,
131 recovery_role,recovery_canonical,logical_topic_id,ordinary_visible,content_fingerprint,
132 meta_fingerprint,health,missing_since,seen_generation
133 ,repair_state,repair_attempts,repair_retry_at,repair_error_kind,repair_source_fingerprint,repair_engine_version
134 ,log_format,head_count,selected_head_id
135 ) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?)
136 ON CONFLICT(path_key) DO UPDATE SET `
137
138 const repairScheduleUpdateSQL = `
139 repair_state=CASE
140 WHEN excluded.turns_state<>'unknown' THEN 'complete'
141 WHEN catalog_sessions.repair_source_fingerprint<>excluded.repair_source_fingerprint
142 OR catalog_sessions.repair_engine_version<>excluded.repair_engine_version THEN 'pending'
143 ELSE catalog_sessions.repair_state END,
144 repair_attempts=CASE
145 WHEN excluded.turns_state<>'unknown'
146 OR catalog_sessions.repair_source_fingerprint<>excluded.repair_source_fingerprint
147 OR catalog_sessions.repair_engine_version<>excluded.repair_engine_version THEN 0
148 ELSE catalog_sessions.repair_attempts END,
149 repair_retry_at=CASE
150 WHEN excluded.turns_state<>'unknown'
151 OR catalog_sessions.repair_source_fingerprint<>excluded.repair_source_fingerprint
152 OR catalog_sessions.repair_engine_version<>excluded.repair_engine_version THEN 0
153 ELSE catalog_sessions.repair_retry_at END,
154 repair_error_kind=CASE
155 WHEN excluded.turns_state<>'unknown'
156 OR catalog_sessions.repair_source_fingerprint<>excluded.repair_source_fingerprint
157 OR catalog_sessions.repair_engine_version<>excluded.repair_engine_version THEN ''
158 ELSE catalog_sessions.repair_error_kind END,
159 repair_source_fingerprint=excluded.repair_source_fingerprint,
160 repair_engine_version=excluded.repair_engine_version,
161 log_format=excluded.log_format, head_count=excluded.head_count, selected_head_id=excluded.selected_head_id`
162
163 const directoryProjectionUpdateSQL = `
164 path=excluded.path, directory=excluded.directory, directory_key=excluded.directory_key, scope=excluded.scope,
165 workspace_root=excluded.workspace_root, workspace_root_key=excluded.workspace_root_key, topic_id=excluded.topic_id,
166 topic_title=excluded.topic_title, custom_title=excluded.custom_title,
167 created_at=excluded.created_at, last_activity_at=excluded.last_activity_at,
168 preview=excluded.preview, turns=excluded.turns,
169 turns_state=excluded.turns_state, recovered=excluded.recovered,
170 recovery_reason=excluded.recovery_reason,
171 recovery_digest=excluded.recovery_digest, parent_id=excluded.parent_id,
172 recovery_copy=excluded.recovery_copy,
173 recovery_group_id=excluded.recovery_group_id,
174 recovery_role=excluded.recovery_role,
175 recovery_canonical=excluded.recovery_canonical,
176 logical_topic_id=excluded.logical_topic_id,
177 ordinary_visible=excluded.ordinary_visible,
178 content_fingerprint=excluded.content_fingerprint,
179 meta_fingerprint=excluded.meta_fingerprint, health=excluded.health,
180 missing_since=0, seen_generation=MAX(catalog_sessions.seen_generation, excluded.seen_generation),` + repairScheduleUpdateSQL
181
182 const exactSourceUpdateSQL = `
183 path=excluded.path, directory=excluded.directory, directory_key=excluded.directory_key, scope=excluded.scope,
184 workspace_root=excluded.workspace_root, workspace_root_key=excluded.workspace_root_key,
185 topic_id=CASE
186 WHEN catalog_sessions.recovered=1 OR excluded.recovered=1 OR catalog_sessions.recovery_group_id<>''
187 THEN catalog_sessions.topic_id ELSE excluded.topic_id END,
188 topic_title=CASE
189 WHEN catalog_sessions.recovered=1 OR excluded.recovered=1 OR catalog_sessions.recovery_group_id<>''
190 THEN catalog_sessions.topic_title ELSE excluded.topic_title END,
191 custom_title=excluded.custom_title,
192 created_at=excluded.created_at, last_activity_at=excluded.last_activity_at,
193 preview=excluded.preview, turns=excluded.turns,
194 turns_state=excluded.turns_state, recovered=excluded.recovered,
195 recovery_reason=excluded.recovery_reason,
196 recovery_digest=excluded.recovery_digest, parent_id=excluded.parent_id,
197 recovery_copy=catalog_sessions.recovery_copy,
198 recovery_group_id=catalog_sessions.recovery_group_id,
199 recovery_role=catalog_sessions.recovery_role,
200 recovery_canonical=catalog_sessions.recovery_canonical,
201 logical_topic_id=catalog_sessions.logical_topic_id,
202 ordinary_visible=catalog_sessions.ordinary_visible,
203 content_fingerprint=excluded.content_fingerprint,
204 meta_fingerprint=excluded.meta_fingerprint, health=excluded.health,
205 missing_since=0, seen_generation=MAX(catalog_sessions.seen_generation, excluded.seen_generation),` + repairScheduleUpdateSQL
206
207 func (c *Catalog) upsertSessionRow(ctx context.Context, tx *sql.Tx, record SessionRecord, pathKey, directoryKey string, generation int64, mode sessionUpsertMode) error {
208 updateSQL := directoryProjectionUpdateSQL
209 if mode == upsertExactSource {
210 updateSQL = exactSourceUpdateSQL
211 }
212 if _, err := tx.ExecContext(ctx, sessionInsertSQL+updateSQL, c.sessionRowValues(record, pathKey, directoryKey, generation)...); err != nil {
213 return err
214 }
215 return upsertHeadRows(ctx, tx, pathKey, record.heads)
216 }
217
218 func (c *Catalog) sessionRowValues(record SessionRecord, pathKey, directoryKey string, generation int64) []any {
219 repairState := "complete"
220 if record.TurnsState == TurnsUnknown {
221 repairState = "pending"
222 }
223 return []any{
224 record.Path, pathKey, record.Directory, directoryKey, record.Scope, record.WorkspaceRoot,
225 c.workspaceRootKey(record.Scope, record.WorkspaceRoot), record.TopicID, record.TopicTitle, record.CustomTitle, record.CreatedAt,
226 record.LastActivityAt, record.Preview, record.Turns, record.TurnsState,
227 record.Recovered, record.RecoveryReason, record.RecoveryDigest,
228 record.ParentID, boolToInt(record.RecoveryCopy), record.RecoveryGroupID,
229 record.RecoveryRole, boolToInt(record.RecoveryCanonical),
230 record.LogicalTopicID, boolToInt(record.OrdinaryVisible),
231 record.ContentFingerprint, record.MetaFingerprint,
232 record.Health, 0, generation,
233 repairState, 0, 0, "", repairSourceFingerprint(record), repairEngineVersion,
234 max(record.LogFormat, 1), record.HeadCount, record.SelectedHeadID,
235 }
236 }
237
238 func repairSourceFingerprint(record SessionRecord) string {
239 return record.ContentFingerprint + "\x00" + record.MetaFingerprint
240 }
241
241 lines GO