| 1 | package sessioncatalog |
| 2 | |
| 3 | import ( |
| 4 | "context" |
| 5 | "crypto/sha256" |
| 6 | "database/sql" |
| 7 | "encoding/hex" |
| 8 | "errors" |
| 9 | "fmt" |
| 10 | "net/url" |
| 11 | "os" |
| 12 | "path/filepath" |
| 13 | "runtime" |
| 14 | "strings" |
| 15 | "sync" |
| 16 | "time" |
| 17 | |
| 18 | "reasonix/internal/agent" |
| 19 | "reasonix/internal/projectiondb" |
| 20 | "reasonix/internal/sqliteuri" |
| 21 | ) |
| 22 | |
| 23 | func (c *Catalog) ReconcileDirectory(ctx context.Context, target DirectoryTarget) error { |
| 24 | if c == nil { |
| 25 | return nil |
| 26 | } |
| 27 | return c.reconcileDirectory(ctx, target, c.mutationSeq.Add(1)) |
| 28 | } |
| 29 | |
| 30 | func (c *Catalog) reconcileDirectory(ctx context.Context, target DirectoryTarget, sequence uint64) error { |
| 31 | if c == nil || c.db == nil { |
| 32 | return nil |
| 33 | } |
| 34 | if c.opts.MetadataOnly { |
| 35 | return c.reconcileMetadata(ctx, target, sequence) |
| 36 | } |
| 37 | target.Path = cleanCatalogAccessPath(target.Path) |
| 38 | if target.Path == "" { |
| 39 | return nil |
| 40 | } |
| 41 | lock := c.directoryLock(target.Path) |
| 42 | lock.Lock() |
| 43 | defer lock.Unlock() |
| 44 | target.Scope, target.WorkspaceRoot = normalizeScope(target.Scope, target.WorkspaceRoot) |
| 45 | signature, err := directorySignature(target.Path) |
| 46 | if err != nil { |
| 47 | c.failDirectoryScan(ctx, target.Path, err) |
| 48 | return err |
| 49 | } |
| 50 | if unchanged, err := c.directoryScanCanSkip(ctx, target, signature); err != nil { |
| 51 | return err |
| 52 | } else if unchanged { |
| 53 | return nil |
| 54 | } |
| 55 | now := c.opts.Now().UnixMilli() |
| 56 | generation, _, err := c.beginDirectoryScan(ctx, target, signature, now) |
| 57 | if err != nil { |
| 58 | return err |
| 59 | } |
| 60 | content := newStrictRecoveryContentCache(c.testSessionContentLoadHook) |
| 61 | ordered, err := listSessionOrderWithContent(target.Path, content) |
| 62 | if err != nil { |
| 63 | c.failDirectoryScan(ctx, target.Path, err) |
| 64 | return err |
| 65 | } |
| 66 | records := make([]SessionRecord, 0, len(ordered)) |
| 67 | for start := 0; start < len(ordered); start += 64 { |
| 68 | if err := ctx.Err(); err != nil { |
| 69 | c.failDirectoryScan(context.Background(), target.Path, err) |
| 70 | return err |
| 71 | } |
| 72 | end := min(start+64, len(ordered)) |
| 73 | for _, info := range ordered[start:end] { |
| 74 | records = append(records, recordFromOrder(target, info)) |
| 75 | } |
| 76 | runtime.Gosched() |
| 77 | } |
| 78 | records = c.filterPathMutations(records, sequence) |
| 79 | records, err = c.preserveKnownSourceStates(ctx, target.Path, records) |
| 80 | if err != nil { |
| 81 | c.failDirectoryScan(context.Background(), target.Path, err) |
| 82 | return err |
| 83 | } |
| 84 | for i := range records { |
| 85 | records[i] = classifyRecoveryLineageWithContent(normalizeSessionRecord(records[i]), content) |
| 86 | } |
| 87 | records = promoteCanonicalLeavesWithContent(records, content) |
| 88 | if err := c.commitDirectoryProjection(ctx, target, signature, generation, now, records); err != nil { |
| 89 | c.failDirectoryScan(context.Background(), target.Path, err) |
| 90 | return err |
| 91 | } |
| 92 | c.markDirectoryVerifiedIfStable(ctx, target, signature) |
| 93 | for _, record := range records { |
| 94 | if record.TurnsState == TurnsUnknown { |
| 95 | c.enqueueRepair(record.Path) |
| 96 | } |
| 97 | } |
| 98 | return nil |
| 99 | } |
| 100 | |
| 101 | func directorySignature(dir string) (string, error) { |
| 102 | // os.ReadDir of a plain file returns the file itself on Windows but |
| 103 | // ENOTDIR on POSIX; stat first so both platforms reject non-directories. |
| 104 | info, err := os.Stat(dir) |
| 105 | if err != nil { |
| 106 | if os.IsNotExist(err) { |
| 107 | return "missing", nil |
| 108 | } |
| 109 | return "", err |
| 110 | } |
| 111 | if !info.IsDir() { |
| 112 | return "", fmt.Errorf("not a directory: %s", dir) |
| 113 | } |
| 114 | entries, err := os.ReadDir(dir) |
| 115 | if err != nil { |
| 116 | return "", err |
| 117 | } |
| 118 | hash := sha256.New() |
| 119 | for _, entry := range entries { |
| 120 | name := entry.Name() |
| 121 | if entry.IsDir() || (!strings.HasSuffix(name, ".jsonl") && !strings.HasSuffix(name, ".meta")) { |
| 122 | continue |
| 123 | } |
| 124 | info, err := entry.Info() |
| 125 | if err != nil { |
| 126 | return "", err |
| 127 | } |
| 128 | _, _ = fmt.Fprintf(hash, "%s\x00%d\x00%d\x00%d\n", name, info.Size(), info.ModTime().UnixNano(), info.Mode()) |
| 129 | } |
| 130 | return hex.EncodeToString(hash.Sum(nil)), nil |
| 131 | } |
| 132 | |
| 133 | func (c *Catalog) directoryLock(path string) *sync.Mutex { |
| 134 | path = c.pathKey(path) |
| 135 | c.directoryLocksMu.Lock() |
| 136 | defer c.directoryLocksMu.Unlock() |
| 137 | lock := c.directoryLocks[path] |
| 138 | if lock == nil { |
| 139 | lock = &sync.Mutex{} |
| 140 | c.directoryLocks[path] = lock |
| 141 | } |
| 142 | return lock |
| 143 | } |
| 144 | |
| 145 | // IndexSessionPath indexes one session without walking its directory. |
| 146 | func (c *Catalog) IndexSessionPath(ctx context.Context, target DirectoryTarget, path string) error { |
| 147 | if c == nil { |
| 148 | return nil |
| 149 | } |
| 150 | return c.indexSessionPath(ctx, target, path, c.mutationSeq.Add(1)) |
| 151 | } |
| 152 | |
| 153 | func (c *Catalog) indexSessionPath(ctx context.Context, target DirectoryTarget, path string, sequence uint64) error { |
| 154 | if c.opts.MetadataOnly { |
| 155 | return c.indexMetadataPath(ctx, target, path, sequence) |
| 156 | } |
| 157 | path = cleanCatalogAccessPath(path) |
| 158 | if path == "" { |
| 159 | return nil |
| 160 | } |
| 161 | target.Path = cleanCatalogAccessPath(target.Path) |
| 162 | if target.Path == "" { |
| 163 | target.Path = filepath.Dir(path) |
| 164 | } |
| 165 | // Hold the directory lock so a concurrent scan cannot mark this row missing. |
| 166 | lock := c.directoryLock(target.Path) |
| 167 | lock.Lock() |
| 168 | defer lock.Unlock() |
| 169 | info, err := os.Stat(path) |
| 170 | if err != nil { |
| 171 | if os.IsNotExist(err) { |
| 172 | return nil |
| 173 | } |
| 174 | return err |
| 175 | } |
| 176 | meta, ok, err := agent.LoadBranchMeta(path) |
| 177 | if err != nil { |
| 178 | return err |
| 179 | } |
| 180 | order := agent.SessionOrderInfo{ |
| 181 | Path: path, |
| 182 | CreatedAt: info.ModTime(), |
| 183 | LastActivityAt: info.ModTime(), |
| 184 | ModTime: info.ModTime(), |
| 185 | Scope: target.Scope, |
| 186 | WorkspaceRoot: target.WorkspaceRoot, |
| 187 | } |
| 188 | if ok { |
| 189 | order.CreatedAt = meta.CreatedAt |
| 190 | order.LastActivityAt = meta.UpdatedAt |
| 191 | order.ModTime = meta.UpdatedAt |
| 192 | order.Scope = meta.DefaultScope() |
| 193 | order.WorkspaceRoot = meta.WorkspaceRoot |
| 194 | order.TopicID = meta.TopicID |
| 195 | order.TopicTitle = meta.TopicTitle |
| 196 | order.CustomTitle = meta.CustomTitle |
| 197 | order.Recovered = meta.Recovered |
| 198 | order.RecoveryReason = meta.RecoveryReason |
| 199 | order.RecoveryDigest = meta.RecoveryDigest |
| 200 | order.ParentID = meta.ParentID |
| 201 | order.RecoveryPreferred = agent.RecoveryPreferenceCurrent(path, meta) |
| 202 | order.Turns = meta.Turns |
| 203 | order.Preview = meta.Preview |
| 204 | order.SchemaVersion = meta.SchemaVersion |
| 205 | order.Revision = meta.Revision |
| 206 | order.ContentDigest = meta.ContentDigest |
| 207 | order.ListingRevision = meta.ListingRevision |
| 208 | order.ListingContentDigest = meta.ListingContentDigest |
| 209 | order.HeadID, order.HeadCount, order.LogSchema = meta.HeadID, meta.HeadCount, meta.LogSchema |
| 210 | } |
| 211 | if order.CreatedAt.IsZero() { |
| 212 | order.CreatedAt = info.ModTime() |
| 213 | } |
| 214 | if order.LastActivityAt.IsZero() { |
| 215 | order.LastActivityAt = info.ModTime() |
| 216 | } |
| 217 | record := recordFromOrder(target, order) |
| 218 | record.enqueueSequence = sequence |
| 219 | projectionDirty, err := c.upsertExactPathSession(ctx, record) |
| 220 | if err != nil { |
| 221 | return err |
| 222 | } |
| 223 | if projectionDirty { |
| 224 | // Queue after the exact source row is durable. The non-blocking worker |
| 225 | // will acquire this directory lock after IndexSessionPath returns and |
| 226 | // publish the full sibling-aware projection in one transaction. |
| 227 | c.RequestReconcile(target) |
| 228 | } |
| 229 | if record.TurnsState == TurnsUnknown { |
| 230 | c.enqueueRepair(record.Path) |
| 231 | } |
| 232 | return nil |
| 233 | } |
| 234 | |
| 235 | func recordFromOrder(target DirectoryTarget, info agent.SessionOrderInfo) SessionRecord { |
| 236 | scope, root := normalizeScope(info.Scope, info.WorkspaceRoot) |
| 237 | if info.TopicID == "" { |
| 238 | scope, root = target.Scope, target.WorkspaceRoot |
| 239 | } |
| 240 | // A stale projection is never certified, but its last-known preview and |
| 241 | // count stay as display hints so the row does not vanish during repair. |
| 242 | turnsState := TurnsValid |
| 243 | if !info.ListingProjectionFresh() { |
| 244 | turnsState = TurnsUnknown |
| 245 | } |
| 246 | heads := projectSessionHeads(info) |
| 247 | contentFingerprint := sessionContentFingerprint(info.Path) + heads.fingerprint |
| 248 | metaFingerprint := fileFingerprint(agent.BranchMetaPath(info.Path)) |
| 249 | if heads.stale { |
| 250 | turnsState = TurnsUnknown |
| 251 | } |
| 252 | createdAt := unixMilli(info.CreatedAt) |
| 253 | lastActivityAt := unixMilli(info.LastActivityAt) |
| 254 | // File mtime fills a missing clock. Do not raise a known sidecar UpdatedAt: |
| 255 | // repair and other metadata writes bump mtime without new user turns. |
| 256 | if st, err := os.Stat(info.Path); err == nil { |
| 257 | fileMS := st.ModTime().UnixMilli() |
| 258 | if createdAt <= 0 { |
| 259 | createdAt = fileMS |
| 260 | } |
| 261 | if lastActivityAt <= 0 { |
| 262 | lastActivityAt = fileMS |
| 263 | } |
| 264 | } |
| 265 | return normalizeSessionRecord(SessionRecord{ |
| 266 | Path: info.Path, |
| 267 | Directory: target.Path, |
| 268 | Scope: scope, |
| 269 | WorkspaceRoot: root, |
| 270 | TopicID: info.TopicID, |
| 271 | TopicTitle: info.TopicTitle, |
| 272 | CustomTitle: info.CustomTitle, |
| 273 | CreatedAt: createdAt, |
| 274 | LastActivityAt: lastActivityAt, |
| 275 | Preview: info.Preview, |
| 276 | Turns: info.Turns, |
| 277 | TurnsState: turnsState, |
| 278 | Recovered: info.Recovered, |
| 279 | RecoveryReason: info.RecoveryReason, |
| 280 | RecoveryDigest: info.RecoveryDigest, |
| 281 | ParentID: info.ParentID, |
| 282 | RecoveryPreferred: info.RecoveryPreferred, |
| 283 | RecoveryCopy: false, |
| 284 | LogFormat: heads.logFormat, |
| 285 | HeadCount: heads.headCount, |
| 286 | SelectedHeadID: heads.selected, |
| 287 | heads: heads.heads, |
| 288 | ContentFingerprint: contentFingerprint, |
| 289 | MetaFingerprint: metaFingerprint, |
| 290 | Health: HealthOK, |
| 291 | }) |
| 292 | } |
| 293 | |
| 294 | func unixMilli(value time.Time) int64 { |
| 295 | if value.IsZero() { |
| 296 | return 0 |
| 297 | } |
| 298 | return value.UnixMilli() |
| 299 | } |
| 300 | |
| 301 | func fileFingerprint(path string) string { |
| 302 | info, err := os.Stat(path) |
| 303 | if err != nil { |
| 304 | return "" |
| 305 | } |
| 306 | return fmt.Sprintf("%d:%d", info.Size(), info.ModTime().UnixNano()) |
| 307 | } |
| 308 | |
| 309 | func sessionContentFingerprint(path string) string { |
| 310 | return fileFingerprint(path) + "|" + fileFingerprint(agent.SessionEventLogPath(path)) |
| 311 | } |
| 312 | |
| 313 | // beginDirectoryScan starts or resumes a directory scan. When the previous |
| 314 | // scan for the same signature was interrupted mid-way, the stored scan_cursor |
| 315 | // is returned so ReconcileDirectory continues instead of restarting from 0. |
| 316 | func (c *Catalog) beginDirectoryScan(ctx context.Context, target DirectoryTarget, signature string, now int64) (int64, string, error) { |
| 317 | c.mutationMu.Lock() |
| 318 | defer c.mutationMu.Unlock() |
| 319 | tx, err := c.db.BeginTx(ctx, nil) |
| 320 | if err != nil { |
| 321 | return 0, "", err |
| 322 | } |
| 323 | var previousSig, previousState, previousCursor string |
| 324 | var previousGeneration int64 |
| 325 | pathKey := c.pathKey(target.Path) |
| 326 | if err := c.removeRemappedDirectoryIdentity(ctx, tx, target.Path, pathKey); err != nil { |
| 327 | _ = tx.Rollback() |
| 328 | return 0, "", err |
| 329 | } |
| 330 | err = tx.QueryRowContext(ctx, `SELECT signature,state,scan_cursor,scan_generation FROM catalog_directories WHERE path_key=?`, |
| 331 | pathKey).Scan(&previousSig, &previousState, &previousCursor, &previousGeneration) |
| 332 | resume := err == nil && previousState == "scanning" && previousSig == signature && strings.TrimSpace(previousCursor) != "" |
| 333 | if errors.Is(err, sql.ErrNoRows) { |
| 334 | err = nil |
| 335 | } |
| 336 | if err != nil { |
| 337 | _ = tx.Rollback() |
| 338 | return 0, "", err |
| 339 | } |
| 340 | if resume { |
| 341 | if _, err := tx.ExecContext(ctx, `UPDATE catalog_directories SET path=?,scope=?,workspace_root=?,state='scanning',error='',signature=? WHERE path_key=?`, |
| 342 | target.Path, target.Scope, target.WorkspaceRoot, signature, pathKey); err != nil { |
| 343 | _ = tx.Rollback() |
| 344 | return 0, "", err |
| 345 | } |
| 346 | if previousGeneration == 0 { |
| 347 | previousGeneration = 1 |
| 348 | } |
| 349 | return previousGeneration, previousCursor, tx.Commit() |
| 350 | } |
| 351 | if _, err := tx.ExecContext(ctx, `INSERT INTO catalog_directories(path,path_key,scope,workspace_root,state,error,signature) |
| 352 | VALUES(?,?,?,?,'scanning','',?) ON CONFLICT(path_key) DO UPDATE SET |
| 353 | path=excluded.path,scope=excluded.scope,workspace_root=excluded.workspace_root,state='scanning',error='', |
| 354 | signature=excluded.signature,scan_generation=catalog_directories.scan_generation+1,scan_cursor='',indexed=0`, |
| 355 | target.Path, pathKey, target.Scope, target.WorkspaceRoot, signature); err != nil { |
| 356 | _ = tx.Rollback() |
| 357 | return 0, "", err |
| 358 | } |
| 359 | var generation int64 |
| 360 | if err := tx.QueryRowContext(ctx, `SELECT scan_generation FROM catalog_directories WHERE path_key=?`, pathKey).Scan(&generation); err != nil { |
| 361 | _ = tx.Rollback() |
| 362 | return 0, "", err |
| 363 | } |
| 364 | if generation == 0 { |
| 365 | generation = 1 |
| 366 | if _, err := tx.ExecContext(ctx, `UPDATE catalog_directories SET scan_generation=1 WHERE path_key=?`, pathKey); err != nil { |
| 367 | _ = tx.Rollback() |
| 368 | return 0, "", err |
| 369 | } |
| 370 | } |
| 371 | return generation, "", tx.Commit() |
| 372 | } |
| 373 | |
| 374 | // commitDirectoryProjection publishes a complete sibling-aware directory |
| 375 | // snapshot. Parsing and lineage classification happen before this function; |
| 376 | // readers therefore observe either the previous committed projection or every |
| 377 | // row, tombstone, missing marker, topic aggregate, and readiness update from |
| 378 | // this transaction together. |
| 379 | func (c *Catalog) commitDirectoryProjection(ctx context.Context, target DirectoryTarget, signature string, generation, now int64, records []SessionRecord) error { |
| 380 | c.mutationMu.Lock() |
| 381 | tx, err := c.db.BeginTx(ctx, nil) |
| 382 | if err != nil { |
| 383 | c.mutationMu.Unlock() |
| 384 | return err |
| 385 | } |
| 386 | rollback := func(commitErr error) error { |
| 387 | _ = tx.Rollback() |
| 388 | c.mutationMu.Unlock() |
| 389 | return commitErr |
| 390 | } |
| 391 | stmt, err := tx.PrepareContext(ctx, sessionInsertSQL+directoryProjectionUpdateSQL) |
| 392 | if err != nil { |
| 393 | return rollback(err) |
| 394 | } |
| 395 | affected := map[TopicKey]struct{}{} |
| 396 | directoryKey := c.pathKey(target.Path) |
| 397 | for start := 0; start < len(records); start += 64 { |
| 398 | if err := ctx.Err(); err != nil { |
| 399 | _ = stmt.Close() |
| 400 | return rollback(err) |
| 401 | } |
| 402 | end := min(start+64, len(records)) |
| 403 | for _, record := range records[start:end] { |
| 404 | pathKey := c.pathKey(record.Path) |
| 405 | remapped, err := removeRemappedSessionIdentity(ctx, tx, record.Path, pathKey) |
| 406 | if err != nil { |
| 407 | _ = stmt.Close() |
| 408 | return rollback(err) |
| 409 | } |
| 410 | for _, key := range remapped { |
| 411 | affected[key] = struct{}{} |
| 412 | } |
| 413 | var previous TopicKey |
| 414 | if err := tx.QueryRowContext(ctx, `SELECT scope,workspace_root,workspace_root_key,topic_id FROM catalog_sessions WHERE path_key=?`, pathKey). |
| 415 | Scan(&previous.Scope, &previous.WorkspaceRoot, &previous.workspaceKey, &previous.TopicID); err == nil && previous.TopicID != "" { |
| 416 | affected[previous] = struct{}{} |
| 417 | } else if err != nil && !errors.Is(err, sql.ErrNoRows) { |
| 418 | _ = stmt.Close() |
| 419 | return rollback(err) |
| 420 | } |
| 421 | if err := c.writeDirectoryRow(ctx, tx, stmt, record, pathKey, directoryKey, generation); err != nil { |
| 422 | _ = stmt.Close() |
| 423 | return rollback(err) |
| 424 | } |
| 425 | if record.TopicID != "" { |
| 426 | affected[TopicKey{Scope: record.Scope, WorkspaceRoot: record.WorkspaceRoot, |
| 427 | workspaceKey: c.workspaceRootKey(record.Scope, record.WorkspaceRoot), TopicID: record.TopicID}] = struct{}{} |
| 428 | } |
| 429 | if err := c.updateFoldedTopicTombstones(ctx, tx, previous, record, now); err != nil { |
| 430 | _ = stmt.Close() |
| 431 | return rollback(err) |
| 432 | } |
| 433 | } |
| 434 | if c.testReconcileBatchHook != nil { |
| 435 | c.testReconcileBatchHook(end) |
| 436 | } |
| 437 | runtime.Gosched() |
| 438 | } |
| 439 | if err := stmt.Close(); err != nil { |
| 440 | return rollback(err) |
| 441 | } |
| 442 | |
| 443 | rows, err := tx.QueryContext(ctx, `SELECT scope,workspace_root,workspace_root_key,topic_id FROM catalog_sessions |
| 444 | WHERE directory_key=? AND seen_generation<? AND topic_id<>''`, directoryKey, generation) |
| 445 | if err != nil { |
| 446 | return rollback(err) |
| 447 | } |
| 448 | for rows.Next() { |
| 449 | var key TopicKey |
| 450 | if err := rows.Scan(&key.Scope, &key.WorkspaceRoot, &key.workspaceKey, &key.TopicID); err != nil { |
| 451 | _ = rows.Close() |
| 452 | return rollback(err) |
| 453 | } |
| 454 | affected[key] = struct{}{} |
| 455 | } |
| 456 | if err := rows.Close(); err != nil { |
| 457 | return rollback(err) |
| 458 | } |
| 459 | if _, err := tx.ExecContext(ctx, `UPDATE catalog_sessions SET |
| 460 | missing_since=CASE WHEN missing_since=0 THEN ? ELSE missing_since END, |
| 461 | health='missing' |
| 462 | WHERE directory_key=? AND seen_generation<?`, now, directoryKey, generation); err != nil { |
| 463 | return rollback(err) |
| 464 | } |
| 465 | cutoff := now - c.opts.MissingGrace.Milliseconds() |
| 466 | if _, err := tx.ExecContext(ctx, `DELETE FROM catalog_sessions |
| 467 | WHERE directory_key=? AND seen_generation<? AND missing_since>0 AND missing_since<=?`, directoryKey, generation, cutoff); err != nil { |
| 468 | return rollback(err) |
| 469 | } |
| 470 | for key := range affected { |
| 471 | if err := c.recomputeTopic(ctx, tx, key); err != nil { |
| 472 | return rollback(err) |
| 473 | } |
| 474 | } |
| 475 | if _, err := tx.ExecContext(ctx, `UPDATE catalog_directories SET state='ready',error='',signature=?, |
| 476 | scan_cursor='',indexed=?,total=?,completed_at=? WHERE path_key=?`, signature, len(records), len(records), now, directoryKey); err != nil { |
| 477 | return rollback(err) |
| 478 | } |
| 479 | revision, err := bumpRevision(ctx, tx) |
| 480 | if err != nil { |
| 481 | return rollback(err) |
| 482 | } |
| 483 | if err := tx.Commit(); err != nil { |
| 484 | c.mutationMu.Unlock() |
| 485 | return err |
| 486 | } |
| 487 | c.mutationMu.Unlock() |
| 488 | |
| 489 | c.publishRevision(revision, []string{target.WorkspaceRoot}, "reconcile_complete") |
| 490 | c.refreshCounts(ctx) |
| 491 | c.statusMu.Lock() |
| 492 | c.status.State = StateReady |
| 493 | c.status.LastError = "" |
| 494 | c.statusMu.Unlock() |
| 495 | return nil |
| 496 | } |
| 497 | |
| 498 | func (c *Catalog) finishDirectoryScan(ctx context.Context, target DirectoryTarget, signature string, generation, now int64, total int) error { |
| 499 | c.mutationMu.Lock() |
| 500 | defer c.mutationMu.Unlock() |
| 501 | tx, err := c.db.BeginTx(ctx, nil) |
| 502 | if err != nil { |
| 503 | return err |
| 504 | } |
| 505 | directoryKey := c.pathKey(target.Path) |
| 506 | rows, err := tx.QueryContext(ctx, `SELECT scope,workspace_root,workspace_root_key,topic_id FROM catalog_sessions |
| 507 | WHERE directory_key=? AND seen_generation<? AND topic_id<>''`, directoryKey, generation) |
| 508 | if err != nil { |
| 509 | _ = tx.Rollback() |
| 510 | return err |
| 511 | } |
| 512 | affected := map[TopicKey]struct{}{} |
| 513 | for rows.Next() { |
| 514 | var key TopicKey |
| 515 | if err := rows.Scan(&key.Scope, &key.WorkspaceRoot, &key.workspaceKey, &key.TopicID); err != nil { |
| 516 | _ = rows.Close() |
| 517 | _ = tx.Rollback() |
| 518 | return err |
| 519 | } |
| 520 | affected[key] = struct{}{} |
| 521 | } |
| 522 | if err := rows.Close(); err != nil { |
| 523 | _ = tx.Rollback() |
| 524 | return err |
| 525 | } |
| 526 | if _, err := tx.ExecContext(ctx, `UPDATE catalog_sessions SET |
| 527 | missing_since=CASE WHEN missing_since=0 THEN ? ELSE missing_since END, |
| 528 | health='missing' |
| 529 | WHERE directory_key=? AND seen_generation<?`, now, directoryKey, generation); err != nil { |
| 530 | _ = tx.Rollback() |
| 531 | return err |
| 532 | } |
| 533 | cutoff := now - c.opts.MissingGrace.Milliseconds() |
| 534 | if _, err := tx.ExecContext(ctx, `DELETE FROM catalog_sessions |
| 535 | WHERE directory_key=? AND seen_generation<? AND missing_since>0 AND missing_since<=?`, directoryKey, generation, cutoff); err != nil { |
| 536 | _ = tx.Rollback() |
| 537 | return err |
| 538 | } |
| 539 | for key := range affected { |
| 540 | if err := c.recomputeTopic(ctx, tx, key); err != nil { |
| 541 | _ = tx.Rollback() |
| 542 | return err |
| 543 | } |
| 544 | } |
| 545 | if _, err := tx.ExecContext(ctx, `UPDATE catalog_directories SET state='ready',error='',signature=?, |
| 546 | scan_cursor='',indexed=?,total=?,completed_at=? WHERE path_key=?`, signature, total, total, now, directoryKey); err != nil { |
| 547 | _ = tx.Rollback() |
| 548 | return err |
| 549 | } |
| 550 | revision, err := bumpRevision(ctx, tx) |
| 551 | if err != nil { |
| 552 | _ = tx.Rollback() |
| 553 | return err |
| 554 | } |
| 555 | if err := tx.Commit(); err != nil { |
| 556 | return err |
| 557 | } |
| 558 | c.publishRevision(revision, []string{target.WorkspaceRoot}, "reconcile_complete") |
| 559 | c.refreshCounts(ctx) |
| 560 | c.statusMu.Lock() |
| 561 | c.status.State = StateReady |
| 562 | c.status.LastError = "" |
| 563 | c.statusMu.Unlock() |
| 564 | return nil |
| 565 | } |
| 566 | |
| 567 | // Rebuild replaces only the disposable catalog. Authoritative sessions and |
| 568 | // sidecars are never changed or removed by this operation. The live database |
| 569 | // stays in place until a fully-populated replacement is validated and swapped. |
| 570 | func Rebuild(ctx context.Context, path string, targets []DirectoryTarget) (Status, error) { |
| 571 | return RebuildWithRevisionFloor(ctx, path, targets, 0) |
| 572 | } |
| 573 | |
| 574 | // RebuildWithRevisionFloor preserves the caller's revision epoch while |
| 575 | // atomically replacing the disposable projection. Desktop clients retain |
| 576 | // revision fences across the rebuild, so a replacement must never publish a |
| 577 | // lower revision than the catalog they already rendered. |
| 578 | func RebuildWithRevisionFloor(ctx context.Context, path string, targets []DirectoryTarget, revisionFloor uint64) (Status, error) { |
| 579 | targets = UniqueDirectoryTargets(targets) |
| 580 | if strings.TrimSpace(path) == "" { |
| 581 | path = DefaultPath() |
| 582 | } |
| 583 | if strings.TrimSpace(path) == "" { |
| 584 | catalog, err := Open(ctx, Options{InMemory: true, DisableRepair: true}) |
| 585 | if err != nil { |
| 586 | return Status{}, err |
| 587 | } |
| 588 | if err := setCatalogRevisionFloor(ctx, catalog.db, revisionFloor); err != nil { |
| 589 | closeCtx, cancel := context.WithTimeout(context.Background(), time.Second) |
| 590 | _ = catalog.Close(closeCtx) |
| 591 | cancel() |
| 592 | return Status{}, err |
| 593 | } |
| 594 | catalog.rememberRevision(revisionFloor) |
| 595 | for _, target := range targets { |
| 596 | if err := catalog.ReconcileDirectory(ctx, target); err != nil { |
| 597 | closeCtx, cancel := context.WithTimeout(context.Background(), time.Second) |
| 598 | _ = catalog.Close(closeCtx) |
| 599 | cancel() |
| 600 | return catalog.Status(), err |
| 601 | } |
| 602 | } |
| 603 | status := catalog.Status() |
| 604 | closeCtx, cancel := context.WithTimeout(context.Background(), time.Second) |
| 605 | _ = catalog.Close(closeCtx) |
| 606 | cancel() |
| 607 | return status, nil |
| 608 | } |
| 609 | err := projectiondb.Rebuild(ctx, projectiondb.OpenOptions{ |
| 610 | Path: path, |
| 611 | MemoryName: "session-catalog-rebuild", |
| 612 | Migrations: sessionMigrations(), |
| 613 | RetainBackup: true, |
| 614 | }, func(ctx context.Context, db *sql.DB) error { |
| 615 | if err := setCatalogRevisionFloor(ctx, db, revisionFloor); err != nil { |
| 616 | return err |
| 617 | } |
| 618 | // Populate through a catalog that owns this temporary database handle |
| 619 | // without starting background repair workers. |
| 620 | temp := &Catalog{ |
| 621 | db: db, |
| 622 | opts: Options{Path: path, DisableRepair: true, Now: time.Now, MissingGrace: defaultMissingGrace}, |
| 623 | pathIdentity: PathIdentityKey, |
| 624 | writeQueued: map[string]SessionRecord{}, |
| 625 | directoryLocks: map[string]*sync.Mutex{}, |
| 626 | stop: make(chan struct{}), |
| 627 | status: Status{State: StateReady, Mode: ModeDisk, Path: path, Revision: revisionFloor}, |
| 628 | } |
| 629 | temp.revision.Store(revisionFloor) |
| 630 | temp.workerCtx, temp.workerCancel = context.WithCancel(ctx) |
| 631 | defer temp.workerCancel() |
| 632 | for _, target := range targets { |
| 633 | if err := temp.ReconcileDirectory(ctx, target); err != nil { |
| 634 | return err |
| 635 | } |
| 636 | } |
| 637 | return nil |
| 638 | }) |
| 639 | if err != nil { |
| 640 | return Status{}, err |
| 641 | } |
| 642 | // Open the published replacement briefly for a status snapshot, then close. |
| 643 | catalog, err := Open(ctx, Options{Path: path, DisableRepair: true}) |
| 644 | if err != nil { |
| 645 | return Status{State: StateReady, Mode: ModeDisk, Path: path, Revision: revisionFloor}, nil |
| 646 | } |
| 647 | status := catalog.Status() |
| 648 | closeCtx, cancel := context.WithTimeout(context.Background(), time.Second) |
| 649 | _ = catalog.Close(closeCtx) |
| 650 | cancel() |
| 651 | return status, nil |
| 652 | } |
| 653 | |
| 654 | func setCatalogRevisionFloor(ctx context.Context, db *sql.DB, revisionFloor uint64) error { |
| 655 | if db == nil || revisionFloor == 0 { |
| 656 | return nil |
| 657 | } |
| 658 | _, err := db.ExecContext(ctx, `UPDATE catalog_state SET revision=? WHERE id=1 AND revision<?`, revisionFloor, revisionFloor) |
| 659 | return err |
| 660 | } |
| 661 | |
| 662 | // Inspect is read-only. It never migrates, repairs, quarantines, or rewrites a |
| 663 | // catalog, making it suitable for `reasonix doctor sessions`. |
| 664 | func Inspect(ctx context.Context, path string) (Status, error) { |
| 665 | if strings.TrimSpace(path) == "" { |
| 666 | path = DefaultPath() |
| 667 | } |
| 668 | status := Status{State: StateDegraded, Mode: ModeDisk, Path: path} |
| 669 | if strings.TrimSpace(path) == "" { |
| 670 | status.LastError = "catalog path unavailable" |
| 671 | return status, nil |
| 672 | } |
| 673 | inspection := projectiondb.Inspect(ctx, path) |
| 674 | if err := ctx.Err(); err != nil { |
| 675 | return status, err |
| 676 | } |
| 677 | if inspection.Error != "" && !inspection.Exists { |
| 678 | status.LastError = inspection.Error |
| 679 | return status, nil |
| 680 | } |
| 681 | if !inspection.Exists { |
| 682 | status.LastError = "catalog does not exist" |
| 683 | return status, nil |
| 684 | } |
| 685 | if inspection.Integrity != "" && inspection.Integrity != "ok" { |
| 686 | status.LastError = inspection.Integrity |
| 687 | return status, nil |
| 688 | } |
| 689 | if inspection.Error != "" { |
| 690 | status.LastError = inspection.Error |
| 691 | return status, nil |
| 692 | } |
| 693 | dsn, err := sqliteuri.Disk(path, url.Values{ |
| 694 | "mode": {"ro"}, |
| 695 | "_pragma": {"busy_timeout(150)"}, |
| 696 | }) |
| 697 | if err != nil { |
| 698 | status.LastError = "build read-only database URI: " + err.Error() |
| 699 | return status, nil |
| 700 | } |
| 701 | db, err := sql.Open("sqlite", dsn) |
| 702 | if err != nil { |
| 703 | status.LastError = "open read-only database: " + err.Error() |
| 704 | return status, nil |
| 705 | } |
| 706 | defer db.Close() |
| 707 | result := status |
| 708 | fail := func(stage string, err error) (Status, error) { |
| 709 | if ctxErr := ctx.Err(); ctxErr != nil { |
| 710 | return status, ctxErr |
| 711 | } |
| 712 | if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) { |
| 713 | return status, err |
| 714 | } |
| 715 | status.LastError = stage + ": " + err.Error() |
| 716 | return status, nil |
| 717 | } |
| 718 | queries := []struct { |
| 719 | stage string |
| 720 | query string |
| 721 | dest any |
| 722 | }{ |
| 723 | {"read catalog revision", `SELECT revision FROM catalog_state WHERE id=1`, &result.Revision}, |
| 724 | {"count indexed sessions", `SELECT COUNT(*) FROM catalog_sessions`, &result.Indexed}, |
| 725 | {"count pending repairs", `SELECT COUNT(*) FROM catalog_sessions WHERE turns_state='unknown'`, &result.RepairPending}, |
| 726 | {"count active repairs", `SELECT COUNT(*) FROM catalog_sessions WHERE turns_state='unknown' AND repair_state IN ('pending','active')`, &result.RepairActive}, |
| 727 | {"count deferred repairs", `SELECT COUNT(*) FROM catalog_sessions WHERE turns_state='unknown' AND repair_state='deferred'`, &result.RepairDeferred}, |
| 728 | {"count blocked repairs", `SELECT COUNT(*) FROM catalog_sessions WHERE turns_state='unknown' AND repair_state='blocked'`, &result.RepairBlocked}, |
| 729 | {"read next repair time", `SELECT COALESCE(MIN(repair_retry_at),0) FROM catalog_sessions WHERE turns_state='unknown' AND repair_state='deferred'`, &result.NextRepairAt}, |
| 730 | {"count physical sessions", `SELECT COUNT(*) FROM catalog_sessions WHERE missing_since=0`, &result.PhysicalSessions}, |
| 731 | {"count logical sessions", `SELECT COUNT(*) FROM catalog_topics`, &result.LogicalSessions}, |
| 732 | {"count recovery groups", `SELECT COUNT(DISTINCT recovery_group_id) FROM catalog_sessions WHERE recovered=1 AND recovery_group_id<>'' AND missing_since=0`, &result.RecoveryGroups}, |
| 733 | {"count recovery branches", `SELECT COUNT(*) FROM catalog_sessions WHERE recovered=1 AND missing_since=0`, &result.RecoveryBranches}, |
| 734 | {"count diverged recoveries", `SELECT COUNT(*) FROM catalog_sessions WHERE recovered=1 AND recovery_role='diverged' AND missing_since=0`, &result.RecoveryDiverged}, |
| 735 | {"count cleanup eligible recoveries", `SELECT COUNT(*) FROM catalog_sessions WHERE recovered=1 AND recovery_role='covered_copy' AND missing_since=0`, &result.CleanupEligible}, |
| 736 | } |
| 737 | for _, query := range queries { |
| 738 | if err := db.QueryRowContext(ctx, query.query).Scan(query.dest); err != nil { |
| 739 | return fail(query.stage, err) |
| 740 | } |
| 741 | } |
| 742 | rows, err := db.QueryContext(ctx, `SELECT repair_error_kind,COUNT(*) FROM catalog_sessions |
| 743 | WHERE turns_state='unknown' AND repair_error_kind<>'' GROUP BY repair_error_kind`) |
| 744 | if err != nil { |
| 745 | return fail("count repair error kinds", err) |
| 746 | } |
| 747 | result.RepairErrorKinds = map[string]int64{} |
| 748 | for rows.Next() { |
| 749 | var kind string |
| 750 | var count int64 |
| 751 | if err := rows.Scan(&kind, &count); err != nil { |
| 752 | _ = rows.Close() |
| 753 | return fail("scan repair error kinds", err) |
| 754 | } |
| 755 | result.RepairErrorKinds[kind] = count |
| 756 | } |
| 757 | if err := rows.Err(); err != nil { |
| 758 | _ = rows.Close() |
| 759 | return fail("iterate repair error kinds", err) |
| 760 | } |
| 761 | if err := rows.Close(); err != nil { |
| 762 | return fail("close repair error kinds", err) |
| 763 | } |
| 764 | result.State = StateReady |
| 765 | result.LastError = "" |
| 766 | return result, nil |
| 767 | } |
| 768 |