返回 DeepSeek-Reasonix
metadata_integrity.go
根目录 / internal / sessioncatalog / metadata_integrity.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "errors"
6 "time"
7
8 "reasonix/internal/projectiondb"
9 )
10
11 var ErrCatalogInvalidated = errors.New("catalog generation invalidated")
12
13 // Invalidated closes only when corruption is proven. The Desktop owner closes
14 // this generation and reopens with synchronous validation/quarantine enabled.
15 func (c *Catalog) Invalidated() <-chan struct{} { return c.invalidated }
16
17 func (c *Catalog) readable() error {
18 select {
19 case <-c.invalidated:
20 return ErrCatalogInvalidated
21 default:
22 return nil
23 }
24 }
25
26 func (c *Catalog) observeDatabaseError(err error) {
27 if !c.opts.MetadataOnly || !projectiondb.IsCorruptionError(err) {
28 return
29 }
30 c.invalidateOnce.Do(func() {
31 c.invalidReason = err
32 c.statusMu.Lock()
33 c.status.State, c.status.LastError = StateDegraded, err.Error()
34 c.statusMu.Unlock()
35 if c.workerCancel != nil {
36 c.workerCancel()
37 }
38 close(c.invalidated)
39 })
40 }
41
42 func (c *Catalog) verifyMetadataIntegrity() {
43 defer c.workers.Done()
44 defer close(c.integrityDone)
45 select {
46 case <-c.workerCtx.Done():
47 return
48 case <-c.discoveryStart:
49 }
50 wait := c.opts.waitMetadataRetry
51 if wait == nil {
52 wait = waitMetadataAuditRetry
53 }
54 backoff := []time.Duration{time.Second, 5 * time.Second, 30 * time.Second}
55 lastError := ""
56 for attempt := 0; ; attempt++ {
57 err := c.checkMetadataIntegrity()
58 if c.workerCtx.Err() != nil || errors.Is(err, context.Canceled) {
59 return
60 }
61 if err == nil {
62 c.statusMu.Lock()
63 if lastError != "" && c.status.LastError == lastError {
64 c.status.State, c.status.LastError = StateReady, ""
65 }
66 c.statusMu.Unlock()
67 return
68 }
69 lastError = err.Error()
70 c.statusMu.Lock()
71 c.status.State, c.status.LastError = StateDegraded, lastError
72 c.statusMu.Unlock()
73 c.observeDatabaseError(err)
74 if c.readable() != nil || attempt == len(backoff) {
75 return
76 }
77 if err := wait(c.workerCtx, backoff[attempt]); err != nil {
78 return
79 }
80 }
81 }
82
83 func (c *Catalog) checkMetadataIntegrity() error {
84 if c.opts.Maintenance != nil {
85 release, err := c.opts.Maintenance.Background(c.workerCtx)
86 if err != nil {
87 return err
88 }
89 defer release()
90 }
91 verify := c.opts.verifyMetadata
92 if verify == nil {
93 verify = func(ctx context.Context) error { return projectiondb.CheckIntegrity(ctx, c.db) }
94 }
95 return verify(c.workerCtx)
96 }
97
98 func waitMetadataAuditRetry(ctx context.Context, delay time.Duration) error {
99 timer := time.NewTimer(delay)
100 defer timer.Stop()
101 select {
102 case <-ctx.Done():
103 return ctx.Err()
104 case <-timer.C:
105 return nil
106 }
107 }
108
109 func (c *Catalog) registerReadLease(lease *ReadLease) bool {
110 c.readLeasesMu.Lock()
111 defer c.readLeasesMu.Unlock()
112 if c.readLeasesClosed || c.workerCtx.Err() != nil || c.readable() != nil {
113 return false
114 }
115 c.readLeases[lease] = struct{}{}
116 return true
117 }
118
119 func (c *Catalog) closeReadLeases() {
120 c.readLeasesMu.Lock()
121 c.readLeasesClosed = true
122 leases := make([]*ReadLease, 0, len(c.readLeases))
123 for lease := range c.readLeases {
124 leases = append(leases, lease)
125 }
126 c.readLeasesMu.Unlock()
127 for _, lease := range leases {
128 lease.Close()
129 }
130 }
131
131 lines GO