返回 DeepSeek-Reasonix
metadata_integrity_test.go
根目录 / internal / sessioncatalog / metadata_integrity_test.go
1 package sessioncatalog
2
3 import (
4 "context"
5 "errors"
6 "os"
7 "path/filepath"
8 "testing"
9 "time"
10 )
11
12 func TestDeferredAuditDoesNotBlockPageAndClosesEveryLeaseOnCorruption(t *testing.T) {
13 entered, release := make(chan struct{}), make(chan struct{})
14 c, err := Open(t.Context(), Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), MetadataOnly: true,
15 DeferredMetadataIntegrity: true, StartPaused: true, RevisionFloor: 90,
16 verifyMetadata: func(ctx context.Context) error {
17 close(entered)
18 select {
19 case <-release:
20 return errors.New("projection integrity check: damaged index")
21 case <-ctx.Done():
22 return ctx.Err()
23 }
24 }})
25 if err != nil {
26 t.Fatal(err)
27 }
28 t.Cleanup(func() { _ = c.Close(context.Background()) })
29 if c.Status().Revision < 90 {
30 t.Fatal("replacement lost the revision fence")
31 }
32 select {
33 case <-entered:
34 t.Fatal("audit ignored discovery admission")
35 default:
36 }
37 c.ResumeDiscovery()
38 <-entered
39 lease, err := c.OpenReadLease(t.Context())
40 if err != nil {
41 t.Fatal(err)
42 }
43 page, err := c.ListOrdinarySessions(lease.Context(t.Context()), OrdinaryPageRequest{Scope: "global", Limit: 50})
44 if err != nil || len(page) != 0 {
45 t.Fatalf("page waited for audit: %v %v", page, err)
46 }
47 close(release)
48 <-c.Invalidated()
49 if _, err := c.OpenReadLease(t.Context()); !errors.Is(err, ErrCatalogInvalidated) {
50 t.Fatalf("accepted invalid generation: %v", err)
51 }
52 if _, err := c.ListOrdinarySessions(lease.Context(t.Context()), OrdinaryPageRequest{Scope: "global"}); !errors.Is(err, ErrCatalogInvalidated) {
53 t.Fatalf("old snapshot remained readable: %v", err)
54 }
55 if err := c.Close(t.Context()); err != nil {
56 t.Fatal(err)
57 }
58 if err := lease.db.Ping(); err == nil {
59 t.Fatal("catalog shutdown left a dedicated lease connection open")
60 }
61 lease.Close()
62 }
63
64 func TestDeferredAuditCancellationAndTransientErrorsDoNotRevokeCatalog(t *testing.T) {
65 for _, failure := range []error{errors.New("database is locked (5) (SQLITE_BUSY)"), context.Canceled} {
66 t.Run(failure.Error(), func(t *testing.T) {
67 attempts := 0
68 var waits []time.Duration
69 c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, DeferredMetadataIntegrity: true,
70 verifyMetadata: func(context.Context) error { attempts++; return failure },
71 waitMetadataRetry: func(_ context.Context, delay time.Duration) error { waits = append(waits, delay); return nil }})
72 if err != nil {
73 t.Fatal(err)
74 }
75 defer c.Close(context.Background())
76 <-c.integrityDone
77 if errors.Is(failure, context.Canceled) {
78 if attempts != 1 || len(waits) != 0 {
79 t.Fatal("cancellation counted as a failed audit")
80 }
81 } else if attempts != 4 || len(waits) != 3 || waits[0] != time.Second || waits[1] != 5*time.Second || waits[2] != 30*time.Second {
82 t.Fatalf("unexpected retry schedule: %d %v", attempts, waits)
83 }
84 select {
85 case <-c.Invalidated():
86 t.Fatal("transient/canceled audit revoked catalog")
87 default:
88 }
89 })
90 }
91 }
92
93 func TestCatalogReadReportsCorruptionToLifecycleOwner(t *testing.T) {
94 revisions := 0
95 c, err := Open(t.Context(), Options{Path: filepath.Join(t.TempDir(), "catalog.sqlite"), MetadataOnly: true, StartPaused: true,
96 OnRevision: func(uint64, []string, string) { revisions++ }})
97 if err != nil {
98 t.Fatal(err)
99 }
100 defer c.Close(context.Background())
101 broken := filepath.Join(t.TempDir(), "broken.sqlite")
102 if err := os.WriteFile(broken, []byte("not a sqlite database"), 0600); err != nil {
103 t.Fatal(err)
104 }
105 rows, err := c.readDB(t.Context()).QueryContext(t.Context(), `ATTACH DATABASE ? AS damaged`, broken)
106 if rows != nil {
107 rows.Close()
108 }
109 if err == nil {
110 t.Fatal("broken database did not fail")
111 }
112 select {
113 case <-c.Invalidated():
114 default:
115 t.Fatal("query corruption did not revoke the generation")
116 }
117 if _, err := os.Stat(broken); err != nil {
118 t.Fatal("reader modified source before lifecycle cleanup")
119 }
120 c.publishRevision(30, nil, "late_publication")
121 if revisions != 0 {
122 t.Fatal("revoked catalog emitted a late revision")
123 }
124 }
125
126 func TestDeferredIntegrityCannotCertifyContentCatalog(t *testing.T) {
127 if c, err := Open(t.Context(), Options{InMemory: true, DeferredMetadataIntegrity: true}); err == nil {
128 c.Close(context.Background())
129 t.Fatal("non-advisory catalog skipped initial integrity validation")
130 }
131 }
132
133 func TestCloseCancelsAnAdmittedIntegrityAudit(t *testing.T) {
134 entered := make(chan struct{})
135 c, err := Open(t.Context(), Options{InMemory: true, MetadataOnly: true, DeferredMetadataIntegrity: true,
136 verifyMetadata: func(ctx context.Context) error {
137 close(entered)
138 <-ctx.Done()
139 return ctx.Err()
140 }})
141 if err != nil {
142 t.Fatal(err)
143 }
144 <-entered
145 if err := c.Close(t.Context()); err != nil {
146 t.Fatal(err)
147 }
148 select {
149 case <-c.integrityDone:
150 default:
151 t.Fatal("close did not join integrity audit")
152 }
153 select {
154 case <-c.Invalidated():
155 t.Fatal("shutdown treated cancellation as corruption")
156 default:
157 }
158 }
159
159 lines GO