返回 DeepSeek-Reasonix
catalog.go
根目录 / internal / historycatalog / catalog.go
1 package historycatalog
2
3 import (
4 "context"
5 "crypto/sha256"
6 "database/sql"
7 "encoding/hex"
8 "errors"
9 "fmt"
10 "os"
11 "path/filepath"
12 "runtime"
13 "sort"
14 "strings"
15 "sync"
16 "sync/atomic"
17 "time"
18
19 "reasonix/internal/agent"
20 "reasonix/internal/projectiondb"
21 "reasonix/internal/store"
22 )
23
24 const defaultMissingGrace = 30 * time.Second
25
26 type Catalog struct {
27 db *sql.DB
28 opts Options
29 revision atomic.Uint64
30 statusMu sync.RWMutex
31 status Status
32 ctx context.Context
33 cancel context.CancelFunc
34 queue chan string
35 rootCh chan string
36 flushCh chan chan struct{}
37 mu sync.Mutex
38 paths map[string]queuedPath
39 roots map[string]Root
40 dirtyRoots map[string]bool
41 wg sync.WaitGroup
42 closeOnce sync.Once
43 closeDone chan struct{}
44 closeErr error
45 }
46
47 type queuedPath struct {
48 root Root
49 appendFrom int
50 }
51
52 func Open(ctx context.Context, opts Options) (*Catalog, error) {
53 if opts.Path == "" {
54 opts.Path = DefaultPath()
55 }
56 if strings.TrimSpace(opts.Path) == "" {
57 opts.Path = ""
58 opts.InMemory = true
59 }
60 if opts.Now == nil {
61 opts.Now = time.Now
62 }
63 if opts.QueueCapacity <= 0 {
64 opts.QueueCapacity = 1024
65 }
66 if opts.MissingGrace <= 0 {
67 opts.MissingGrace = defaultMissingGrace
68 }
69 if opts.ReconcileInterval <= 0 {
70 // Periodic root rescans are fingerprint-cheap when nothing changed; keep
71 // the interval longer so large installs are not re-walked every minute.
72 opts.ReconcileInterval = 5 * time.Minute
73 }
74 opts.MaxBytes = resolveMaxBytes(opts.MaxBytes, configuredMaxMB())
75 handle, err := projectiondb.Open(ctx, projectiondb.OpenOptions{
76 Path: opts.Path, MemoryName: "history-search", Migrations: migrations(), InMemory: opts.InMemory,
77 MaxOpenConns: 4, Now: opts.Now, SecureDelete: true, AutoVacuum: true,
78 })
79 if err != nil {
80 return nil, err
81 }
82 workerCtx, cancel := context.WithCancel(context.Background())
83 c := &Catalog{db: handle.DB, opts: opts, ctx: workerCtx, cancel: cancel,
84 queue: make(chan string, opts.QueueCapacity), rootCh: make(chan string, 64),
85 flushCh: make(chan chan struct{}, 1),
86 paths: map[string]queuedPath{}, roots: map[string]Root{}, dirtyRoots: map[string]bool{},
87 closeDone: make(chan struct{}),
88 status: Status{State: string(handle.Status.State), Mode: handle.Status.Mode, Path: handle.Status.Path,
89 LastError: handle.Status.LastError, QuarantinedPath: handle.Status.QuarantinedPath}}
90 if err := c.db.QueryRowContext(ctx, `SELECT revision FROM history_state WHERE id=1`).Scan(new(uint64)); err != nil {
91 _ = c.db.Close()
92 cancel()
93 return nil, err
94 }
95 if err := c.ensureTokenizerVersion(ctx); err != nil {
96 _ = c.db.Close()
97 cancel()
98 return nil, err
99 }
100 var revision uint64
101 _ = c.db.QueryRowContext(ctx, `SELECT revision FROM history_state WHERE id=1`).Scan(&revision)
102 c.revision.Store(revision)
103 c.refreshStatus(ctx)
104 c.wg.Add(1)
105 go c.worker()
106 // A far-over-cap index from before the cap existed (#8717) is cheaper to
107 // rebuild than to evict session-by-session; wipe async so startup never
108 // blocks. Registered roots rescan afterwards and re-index truncated.
109 if !opts.InMemory && strings.TrimSpace(opts.Path) != "" &&
110 historyDBFileSize(opts.Path) > rebuildOversizeFactor*opts.MaxBytes {
111 c.wg.Go(func() {
112 c.wipeForRebuild(c.ctx)
113 })
114 }
115 return c, nil
116 }
117
118 func (c *Catalog) ensureTokenizerVersion(ctx context.Context) error {
119 var version int
120 if err := c.db.QueryRowContext(ctx, `SELECT tokenizer_version FROM history_state WHERE id=1`).Scan(&version); err != nil {
121 return err
122 }
123 if version == TokenizerVersion {
124 return nil
125 }
126 tx, err := c.db.BeginTx(ctx, nil)
127 if err != nil {
128 return err
129 }
130 if err := wipeProjectionRows(ctx, tx); err != nil {
131 _ = tx.Rollback()
132 return err
133 }
134 if _, err := tx.ExecContext(ctx, `UPDATE history_state SET tokenizer_version=?,revision=revision+1 WHERE id=1`, TokenizerVersion); err != nil {
135 _ = tx.Rollback()
136 return err
137 }
138 return tx.Commit()
139 }
140
141 func (c *Catalog) RegisterRoot(root Root) bool {
142 if c == nil || strings.TrimSpace(root.Path) == "" {
143 return false
144 }
145 root.Path = filepath.Clean(root.Path)
146 if root.Scope != "project" {
147 root.Scope = "global"
148 root.WorkspaceRoot = ""
149 }
150 c.mu.Lock()
151 _, alreadyRegistered := c.roots[root.Path]
152 c.roots[root.Path] = root
153 c.mu.Unlock()
154 if !alreadyRegistered {
155 c.statusMu.Lock()
156 c.status.Pending++
157 c.statusMu.Unlock()
158 }
159 select {
160 case c.rootCh <- root.Path:
161 return true
162 default:
163 c.markRootDirty(root.Path)
164 return false
165 }
166 }
167
168 // ReconcileRoot performs one deterministic scan. Production callers normally
169 // use RegisterRoot; the synchronous form exists for doctor/reindex and tests.
170 func (c *Catalog) ReconcileRoot(ctx context.Context, root Root) error {
171 root.Path = filepath.Clean(root.Path)
172 if root.Scope != "project" {
173 root.Scope = "global"
174 root.WorkspaceRoot = ""
175 }
176 return c.reconcileRoot(ctx, root)
177 }
178
179 func (c *Catalog) EnqueuePath(root Root, path string) bool {
180 return c.enqueuePath(root, path, -1)
181 }
182
183 func (c *Catalog) EnqueuePersist(root Root, event agent.SessionPersistEvent) bool {
184 appendFrom := event.AppendFrom
185 if event.Rewrite {
186 appendFrom = -1
187 }
188 return c.enqueuePath(root, event.Path, appendFrom)
189 }
190
191 func (c *Catalog) enqueuePath(root Root, path string, appendFrom int) bool {
192 if c == nil || strings.TrimSpace(path) == "" {
193 return false
194 }
195 path = filepath.Clean(path)
196 c.mu.Lock()
197 if queued, exists := c.paths[path]; exists {
198 queued.root = root
199 if queued.appendFrom < 0 || appendFrom < 0 {
200 queued.appendFrom = -1
201 } else if appendFrom < queued.appendFrom {
202 queued.appendFrom = appendFrom
203 }
204 c.paths[path] = queued
205 c.mu.Unlock()
206 return true
207 }
208 c.paths[path] = queuedPath{root: root, appendFrom: appendFrom}
209 c.mu.Unlock()
210 select {
211 case c.queue <- path:
212 return true
213 default:
214 c.mu.Lock()
215 delete(c.paths, path)
216 c.dirtyRoots[root.Path] = true
217 c.mu.Unlock()
218 return false
219 }
220 }
221
222 // EnqueueExisting prioritizes a source already known to the catalog without
223 // requiring the caller to retain its root metadata.
224 func (c *Catalog) EnqueueExisting(ctx context.Context, path string) bool {
225 var root Root
226 err := c.db.QueryRowContext(ctx, `SELECT root,source,scope,workspace_root FROM history_sources WHERE path=?`, filepath.Clean(path)).Scan(
227 &root.Path, &root.Source, &root.Scope, &root.WorkspaceRoot)
228 if err != nil {
229 return false
230 }
231 return c.EnqueuePath(root, path)
232 }
233
234 func (c *Catalog) worker() {
235 defer c.wg.Done()
236 ticker := time.NewTicker(c.opts.ReconcileInterval)
237 defer ticker.Stop()
238 for {
239 if root, ok := c.takeDirtyRoot(); ok {
240 _ = c.reconcileRoot(c.ctx, root)
241 continue
242 }
243 select {
244 case <-c.ctx.Done():
245 return
246 case <-ticker.C:
247 c.markAllRootsDirty()
248 c.governSize(c.ctx)
249 case done := <-c.flushCh:
250 c.drainPending(c.ctx)
251 close(done)
252 case path := <-c.queue:
253 c.mu.Lock()
254 queued := c.paths[path]
255 delete(c.paths, path)
256 c.mu.Unlock()
257 _ = c.indexPath(c.ctx, queued.root, path, 0, queued.appendFrom)
258 case path := <-c.rootCh:
259 c.mu.Lock()
260 root, ok := c.roots[path]
261 c.mu.Unlock()
262 if ok {
263 _ = c.reconcileRoot(c.ctx, root)
264 }
265 }
266 }
267 }
268
269 // Flush drains dirty roots and the path queue until empty or ctx cancels, then
270 // waits for the worker to acknowledge. Callers use this on shutdown so pending
271 // index work is not silently abandoned.
272 func (c *Catalog) Flush(ctx context.Context) error {
273 if c == nil {
274 return nil
275 }
276 done := make(chan struct{})
277 select {
278 case c.flushCh <- done:
279 case <-ctx.Done():
280 return ctx.Err()
281 }
282 select {
283 case <-done:
284 return nil
285 case <-ctx.Done():
286 return ctx.Err()
287 }
288 }
289
290 func (c *Catalog) drainPending(ctx context.Context) {
291 for {
292 if err := ctx.Err(); err != nil {
293 return
294 }
295 if root, ok := c.takeDirtyRoot(); ok {
296 _ = c.reconcileRoot(ctx, root)
297 continue
298 }
299 select {
300 case path := <-c.queue:
301 c.mu.Lock()
302 queued := c.paths[path]
303 delete(c.paths, path)
304 c.mu.Unlock()
305 _ = c.indexPath(ctx, queued.root, path, 0, queued.appendFrom)
306 case path := <-c.rootCh:
307 c.mu.Lock()
308 root, ok := c.roots[path]
309 c.mu.Unlock()
310 if ok {
311 _ = c.reconcileRoot(ctx, root)
312 }
313 default:
314 c.mu.Lock()
315 empty := len(c.paths) == 0 && len(c.dirtyRoots) == 0
316 c.mu.Unlock()
317 if empty && len(c.queue) == 0 && len(c.rootCh) == 0 {
318 return
319 }
320 // Another goroutine may have enqueued between checks; yield once.
321 runtime.Gosched()
322 c.mu.Lock()
323 empty = len(c.paths) == 0 && len(c.dirtyRoots) == 0
324 c.mu.Unlock()
325 if empty && len(c.queue) == 0 && len(c.rootCh) == 0 {
326 return
327 }
328 }
329 }
330 }
331
332 func (c *Catalog) markRootDirty(path string) {
333 c.mu.Lock()
334 c.dirtyRoots[path] = true
335 c.mu.Unlock()
336 }
337
338 func (c *Catalog) markAllRootsDirty() {
339 c.mu.Lock()
340 for path := range c.roots {
341 c.dirtyRoots[path] = true
342 }
343 c.mu.Unlock()
344 }
345
346 func (c *Catalog) takeDirtyRoot() (Root, bool) {
347 c.mu.Lock()
348 defer c.mu.Unlock()
349 for path := range c.dirtyRoots {
350 delete(c.dirtyRoots, path)
351 root, ok := c.roots[path]
352 return root, ok
353 }
354 return Root{}, false
355 }
356
357 func historyRootSignature(paths []string) string {
358 hash := sha256.New()
359 for _, path := range paths {
360 for _, candidate := range []string{path, agent.BranchMetaPath(path)} {
361 info, err := os.Stat(candidate)
362 if err != nil {
363 _, _ = fmt.Fprintf(hash, "%s\x00missing\n", candidate)
364 continue
365 }
366 _, _ = fmt.Fprintf(hash, "%s\x00%d\x00%d\n", candidate, info.Size(), info.ModTime().UnixNano())
367 }
368 }
369 return hex.EncodeToString(hash.Sum(nil))
370 }
371
372 func (c *Catalog) reconcileRoot(ctx context.Context, root Root) error {
373 entries, err := os.ReadDir(root.Path)
374 if err != nil && !errors.Is(err, os.ErrNotExist) {
375 c.setError(err)
376 return err
377 }
378 paths := make([]string, 0, len(entries))
379 for _, entry := range entries {
380 if entry.IsDir() || !store.IsSessionTranscriptName(entry.Name()) {
381 continue
382 }
383 path := filepath.Join(root.Path, entry.Name())
384 if !root.Archive && !agent.IsVisibleSession(path) {
385 continue
386 }
387 paths = append(paths, path)
388 }
389 sort.Strings(paths)
390 signature := historyRootSignature(paths)
391 var previousSig, previousState string
392 _ = c.db.QueryRowContext(ctx, `SELECT signature,state FROM history_roots WHERE path=?`, root.Path).Scan(&previousSig, &previousState)
393 if previousState == "ready" && previousSig == signature && signature != "" {
394 return nil
395 }
396 now := c.opts.Now().UnixMilli()
397 tx, err := c.db.BeginTx(ctx, nil)
398 if err != nil {
399 return err
400 }
401 var generation int64
402 if err := tx.QueryRowContext(ctx, `INSERT INTO history_roots(path,source,scope,workspace_root,signature,scan_generation,state,total)
403 VALUES(?,?,?,?,?,1,'scanning',?) ON CONFLICT(path) DO UPDATE SET source=excluded.source,scope=excluded.scope,
404 workspace_root=excluded.workspace_root,signature=excluded.signature,scan_generation=history_roots.scan_generation+1,state='scanning',error='',total=excluded.total
405 RETURNING scan_generation`, root.Path, root.Source, root.Scope, root.WorkspaceRoot, signature, len(paths)).Scan(&generation); err != nil {
406 _ = tx.Rollback()
407 return err
408 }
409 if err := tx.Commit(); err != nil {
410 return err
411 }
412 indexed := 0
413 for i, path := range paths {
414 if err := ctx.Err(); err != nil {
415 return err
416 }
417 if err := c.indexPath(ctx, root, path, generation, -1); err == nil {
418 indexed++
419 }
420 if (i+1)%32 == 0 {
421 runtime.Gosched()
422 }
423 }
424 tx, err = c.db.BeginTx(ctx, nil)
425 if err != nil {
426 return err
427 }
428 if _, err := tx.ExecContext(ctx, `UPDATE history_sources SET missing_since=CASE WHEN missing_since=0 THEN ? ELSE missing_since END,
429 health='missing' WHERE root=? AND seen_generation<>?`, now, root.Path, generation); err != nil {
430 _ = tx.Rollback()
431 return err
432 }
433 cutoff := now - c.opts.MissingGrace.Milliseconds()
434 if _, err := tx.ExecContext(ctx, `DELETE FROM history_fts WHERE rowid IN (
435 SELECT d.id FROM history_documents d JOIN history_sources s ON s.path=d.source_path
436 WHERE s.root=? AND s.seen_generation<>? AND s.missing_since>0 AND s.missing_since<=?
437 )`, root.Path, generation, cutoff); err != nil {
438 _ = tx.Rollback()
439 return err
440 }
441 if _, err := tx.ExecContext(ctx, `DELETE FROM history_documents WHERE source_path IN (
442 SELECT path FROM history_sources WHERE root=? AND seen_generation<>? AND missing_since>0 AND missing_since<=?
443 )`, root.Path, generation, cutoff); err != nil {
444 _ = tx.Rollback()
445 return err
446 }
447 if _, err := tx.ExecContext(ctx, `DELETE FROM history_sources
448 WHERE root=? AND seen_generation<>? AND missing_since>0 AND missing_since<=?`, root.Path, generation, cutoff); err != nil {
449 _ = tx.Rollback()
450 return err
451 }
452 if _, err := tx.ExecContext(ctx, `UPDATE history_roots SET state='ready',signature=?,indexed=?,completed_at=?,scan_cursor='' WHERE path=?`, signature, indexed, now, root.Path); err != nil {
453 _ = tx.Rollback()
454 return err
455 }
456 revision, err := bump(ctx, tx)
457 if err != nil {
458 _ = tx.Rollback()
459 return err
460 }
461 if err := tx.Commit(); err != nil {
462 return err
463 }
464 c.publish(revision, []string{root.Path}, "reconcile")
465 return nil
466 }
467
468 func fileFingerprint(path string) string {
469 info, err := os.Stat(path)
470 if err != nil {
471 return ""
472 }
473 return fmt.Sprintf("%d:%d", info.Size(), info.ModTime().UnixNano())
474 }
475
476 func (c *Catalog) indexPath(ctx context.Context, root Root, path string, generation int64, appendFrom int) error {
477 if !root.Archive && !agent.IsVisibleSession(path) {
478 return c.Purge(ctx, path)
479 }
480 contentFingerprint := fileFingerprint(path)
481 metaFingerprint := fileFingerprint(agent.BranchMetaPath(path))
482 state, known, identityErr := agent.SessionContentIdentity(path)
483 digest := ""
484 revision := int64(0)
485 if identityErr == nil && known {
486 digest, revision = state.DigestHex, state.Revision
487 }
488 var oldFingerprint, oldMetaFingerprint, oldDigest, oldHealth string
489 var oldGeneration, oldRevision int64
490 var oldMessageCount int
491 err := c.db.QueryRowContext(ctx, `SELECT content_fingerprint,meta_fingerprint,content_digest,seen_generation,
492 content_revision,indexed_message_count,health FROM history_sources WHERE path=?`, path).Scan(
493 &oldFingerprint, &oldMetaFingerprint, &oldDigest, &oldGeneration, &oldRevision, &oldMessageCount, &oldHealth)
494 if sourceProjectionUnchanged(err, oldFingerprint, contentFingerprint, oldMetaFingerprint, metaFingerprint, oldDigest, digest) {
495 if generation != 0 && oldGeneration != generation {
496 // Keep evicted rows evicted: an unchanged file must not re-enter the index.
497 _, _ = c.db.ExecContext(ctx, `UPDATE history_sources SET seen_generation=?,missing_since=0,
498 health=CASE WHEN health='evicted' THEN 'evicted' ELSE 'ok' END WHERE path=?`, generation, path)
499 }
500 return nil
501 }
502 if err != nil && !errors.Is(err, sql.ErrNoRows) {
503 return err
504 }
505 // An evicted projection has no prefix to append onto; fall through to a full reload.
506 if err == nil && known && appendFrom >= 0 && oldHealth != "evicted" {
507 handled, appendErr := c.tryAppendPath(ctx, root, path, generation, appendFrom, oldMessageCount, oldRevision,
508 revision, digest, contentFingerprint, metaFingerprint)
509 if appendErr != nil {
510 return appendErr
511 }
512 if handled {
513 return nil
514 }
515 }
516 session, err := agent.LoadSession(path)
517 if err != nil {
518 _, _ = c.db.ExecContext(ctx, `INSERT INTO history_sources(path,root,source,scope,workspace_root,content_fingerprint,meta_fingerprint,health,last_error,seen_generation)
519 VALUES(?,?,?,?,?,?,?,'corrupt',?,?) ON CONFLICT(path) DO UPDATE SET health='corrupt',last_error=excluded.last_error,
520 content_fingerprint=excluded.content_fingerprint,meta_fingerprint=excluded.meta_fingerprint,seen_generation=excluded.seen_generation`,
521 path, root.Path, root.Source, root.Scope, root.WorkspaceRoot, contentFingerprint, metaFingerprint, err.Error(), generation)
522 c.setError(err)
523 return err
524 }
525 messages := session.Snapshot()
526 if digest == "" {
527 h := sha256.New()
528 for _, doc := range documents(messages) {
529 _, _ = h.Write([]byte(doc.terms))
530 _, _ = h.Write([]byte{0})
531 }
532 digest = hex.EncodeToString(h.Sum(nil))
533 }
534 meta, _, _ := agent.LoadBranchMeta(path)
535 lastActivity := max(int64(0), agent.SessionContentModTime(path).UnixMilli())
536 // Hide stale terms as soon as the authoritative fingerprint changes. Rows
537 // remain available for retry and are atomically replaced below.
538 if _, err := c.db.ExecContext(ctx, `UPDATE history_sources SET health='stale',last_error='' WHERE path=?`, path); err != nil {
539 return err
540 }
541 tx, err := c.db.BeginTx(ctx, nil)
542 if err != nil {
543 return err
544 }
545 if _, err := tx.ExecContext(ctx, `DELETE FROM history_fts WHERE rowid IN (SELECT id FROM history_documents WHERE source_path=?)`, path); err != nil {
546 _ = tx.Rollback()
547 return err
548 }
549 if _, err := tx.ExecContext(ctx, `DELETE FROM history_documents WHERE source_path=?`, path); err != nil {
550 _ = tx.Rollback()
551 return err
552 }
553 _, err = tx.ExecContext(ctx, `INSERT INTO history_sources(path,root,source,scope,workspace_root,content_revision,content_digest,
554 content_fingerprint,meta_fingerprint,message_count,indexed_message_count,custom_title,topic_id,topic_title,preview,created_at,
555 last_activity_at,health,missing_since,seen_generation,last_error) VALUES(?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,?,'ok',0,?,'')
556 ON CONFLICT(path) DO UPDATE SET root=excluded.root,source=excluded.source,scope=excluded.scope,workspace_root=excluded.workspace_root,
557 content_revision=excluded.content_revision,content_digest=excluded.content_digest,content_fingerprint=excluded.content_fingerprint,
558 meta_fingerprint=excluded.meta_fingerprint,message_count=excluded.message_count,indexed_message_count=excluded.indexed_message_count,
559 custom_title=excluded.custom_title,topic_id=excluded.topic_id,topic_title=excluded.topic_title,preview=excluded.preview,
560 created_at=excluded.created_at,last_activity_at=excluded.last_activity_at,health='ok',missing_since=0,
561 seen_generation=excluded.seen_generation,last_error=''`, path, root.Path, root.Source, root.Scope, root.WorkspaceRoot, revision, digest,
562 contentFingerprint, metaFingerprint, len(messages), len(messages), meta.CustomTitle, meta.TopicID, meta.TopicTitle, meta.Preview,
563 meta.CreatedAt.UnixMilli(), lastActivity, generation)
564 if err != nil {
565 _ = tx.Rollback()
566 return err
567 }
568 for _, doc := range documents(messages) {
569 result, err := tx.ExecContext(ctx, `INSERT INTO history_documents(source_path,message_index,part_index,role,kind,tool_name,token_count)
570 VALUES(?,?,?,?,?,?,?)`, path, doc.message, doc.part, doc.role, doc.kind, doc.tool, doc.count)
571 if err != nil {
572 _ = tx.Rollback()
573 return err
574 }
575 rowID, err := result.LastInsertId()
576 if err != nil {
577 _ = tx.Rollback()
578 return err
579 }
580 if _, err := tx.ExecContext(ctx, `INSERT INTO history_fts(rowid,terms) VALUES(?,?)`, rowID, doc.terms); err != nil {
581 _ = tx.Rollback()
582 return err
583 }
584 }
585 newRevision, err := bump(ctx, tx)
586 if err != nil {
587 _ = tx.Rollback()
588 return err
589 }
590 if err := tx.Commit(); err != nil {
591 return err
592 }
593 c.publish(newRevision, []string{root.Path}, "source-indexed")
594 return nil
595 }
596
597 func bump(ctx context.Context, tx *sql.Tx) (uint64, error) {
598 if _, err := tx.ExecContext(ctx, `UPDATE history_state SET revision=revision+1 WHERE id=1`); err != nil {
599 return 0, err
600 }
601 var revision uint64
602 err := tx.QueryRowContext(ctx, `SELECT revision FROM history_state WHERE id=1`).Scan(&revision)
603 return revision, err
604 }
605
606 func (c *Catalog) Purge(ctx context.Context, path string) error {
607 if c == nil {
608 return nil
609 }
610 tx, err := c.db.BeginTx(ctx, nil)
611 if err != nil {
612 return err
613 }
614 if _, err := tx.ExecContext(ctx, `DELETE FROM history_fts WHERE rowid IN (SELECT id FROM history_documents WHERE source_path=?)`, path); err != nil {
615 _ = tx.Rollback()
616 return err
617 }
618 if _, err := tx.ExecContext(ctx, `DELETE FROM history_sources WHERE path=?`, path); err != nil {
619 _ = tx.Rollback()
620 return err
621 }
622 revision, err := bump(ctx, tx)
623 if err != nil {
624 _ = tx.Rollback()
625 return err
626 }
627 if err := tx.Commit(); err != nil {
628 return err
629 }
630 _, _ = c.db.ExecContext(ctx, `PRAGMA wal_checkpoint(TRUNCATE)`)
631 _, _ = c.db.ExecContext(ctx, `PRAGMA incremental_vacuum(64)`)
632 c.publish(revision, nil, "purge")
633 return nil
634 }
635
636 func (c *Catalog) refreshStatus(ctx context.Context) {
637 var indexed, total, pending, failed int64
638 _ = c.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_sources WHERE health='ok'`).Scan(&indexed)
639 _ = c.db.QueryRowContext(ctx, `SELECT COALESCE(SUM(total),0) FROM history_roots`).Scan(&total)
640 _ = c.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_roots WHERE state<>'ready'`).Scan(&pending)
641 _ = c.db.QueryRowContext(ctx, `SELECT COUNT(*) FROM history_sources WHERE health='corrupt'`).Scan(&failed)
642 c.statusMu.Lock()
643 c.status.Indexed, c.status.Total, c.status.Pending, c.status.Failed = indexed, total, pending, failed
644 c.status.Revision = c.revision.Load()
645 c.statusMu.Unlock()
646 }
647 func (c *Catalog) publish(revision uint64, roots []string, reason string) {
648 c.revision.Store(revision)
649 statusCtx := c.ctx
650 if statusCtx == nil {
651 statusCtx = context.Background()
652 }
653 c.refreshStatus(statusCtx)
654 if c.opts.OnRevision != nil {
655 c.opts.OnRevision(c.Status(), roots, reason)
656 }
657 }
658
659 func (c *Catalog) setError(err error) {
660 c.statusMu.Lock()
661 c.status.LastError = err.Error()
662 c.status.Failed++
663 c.statusMu.Unlock()
664 }
665
666 func (c *Catalog) Status() Status {
667 if c == nil {
668 return Status{State: "degraded", Mode: projectiondb.ModeMemory, LastError: "history catalog unavailable"}
669 }
670 c.statusMu.RLock()
671 defer c.statusMu.RUnlock()
672 return c.status
673 }
674
675 func (c *Catalog) Close(ctx context.Context) error {
676 if c == nil {
677 return nil
678 }
679 c.closeOnce.Do(func() {
680 c.cancel()
681 go func() {
682 c.wg.Wait()
683 c.closeErr = c.db.Close()
684 close(c.closeDone)
685 }()
686 })
687 select {
688 case <-c.closeDone:
689 return c.closeErr
690 case <-ctx.Done():
691 return ctx.Err()
692 }
693 }
694
694 lines GO