Skip to content
File

Blob: cmd/server/kv_provider_test.go

go291 lines
1package server
2 
3import (
4 "context"
5 "io/fs"
6 "os"
7 "path/filepath"
8 "runtime"
9 "testing"
10 
11 "github.com/stretchr/testify/require"
12 "github.com/urfave/cli/v3"
13 "go.uber.org/zap"
14)
15 
16func TestKVProviderPreflightPreservesExistingStores(t *testing.T) {
17 for _, provider := range []string{"aof", "sqlite"} {
18 t.Run(provider, func(t *testing.T) {
19 dir := t.TempDir()
20 vnodeDir := filepath.Join(dir, "12")
21 cacheDir := filepath.Join(dir, "cache")
22 kv, stop, err := getKVProvider(zap.NewNop(), provider, vnodeDir, cacheDir)
23 require.NoError(t, err)
24 err = kv.Put(context.Background(), []byte("existing"), []byte("preserve me"))
25 stop()
26 require.NoError(t, err)
27 before := snapshotStorage(t, dir)
28 
29 other := "sqlite"
30 if provider == "sqlite" {
31 other = "aof"
32 }
33 // An unmarked store from an older server must not be silently replaced.
34 require.ErrorContains(t, preflightKVProvider(other, dir), "conflicts")
35 require.ErrorContains(t, preflightKVProvider("memory", dir), "conflicts")
36 require.Equal(t, before, snapshotStorage(t, dir))
37 
38 require.NoError(t, preflightKVProvider(provider, dir))
39 before[kvProviderMarker] = provider + "\n"
40 require.Equal(t, before, snapshotStorage(t, dir), "adoption must only add the provider marker")
41 require.NoError(t, preflightKVProvider(provider, dir))
42 require.ErrorContains(t, preflightKVProvider(other, dir), "conflicts")
43 require.Equal(t, before, snapshotStorage(t, dir))
44 
45 kv, stop, err = getKVProvider(zap.NewNop(), provider, vnodeDir, cacheDir)
46 require.NoError(t, err)
47 defer stop()
48 value, err := kv.Get(context.Background(), []byte("existing"))
49 require.NoError(t, err)
50 require.Equal(t, []byte("preserve me"), value)
51 })
52 }
53}
54 
55func TestKVProviderPreflightLayouts(t *testing.T) {
56 tests := []struct {
57 name string
58 provider string
59 dirs []string
60 files map[string]string
61 wantErr string
62 }{
63 {name: "fresh aof", provider: "aof"},
64 {name: "fresh sqlite", provider: "sqlite"},
65 {name: "fresh memory", provider: "memory"},
66 {name: "empty vnode", provider: "aof", dirs: []string{"0", "8"}},
67 {name: "cache alone is not storage", provider: "memory", files: map[string]string{"cache/sqlite3/runtime": "cache"}},
68 {name: "empty aof directory identifies storage", provider: "sqlite", dirs: []string{"0/wal"}, wantErr: "conflicts"},
69 {name: "empty sqlite directory identifies storage", provider: "aof", dirs: []string{"0/sqlite3"}, wantErr: "conflicts"},
70 {name: "sqlite sidecars identify storage", provider: "aof", files: map[string]string{"0/sqlite3/db-wal": "journal", "0/sqlite3/db-shm": "shared memory"}, wantErr: "conflicts"},
71 {name: "same vnode mixture", provider: "aof", dirs: []string{"0/wal", "0/sqlite3"}, wantErr: "mixed aof and sqlite"},
72 {name: "cross vnode mixture", provider: "sqlite", dirs: []string{"0/sqlite3", "99/wal"}, wantErr: "mixed aof and sqlite"},
73 {name: "marker conflict without store", provider: "sqlite", files: map[string]string{kvProviderMarker: "aof\n"}, wantErr: "conflicts"},
74 {name: "marker cannot hide conflicting layout", provider: "aof", dirs: []string{"0/sqlite3"}, files: map[string]string{kvProviderMarker: "aof\n"}, wantErr: "marker identifies aof"},
75 {name: "memory rejects marker", provider: "memory", files: map[string]string{kvProviderMarker: "sqlite\n"}, wantErr: "conflicts"},
76 {name: "memory rejects layout", provider: "memory", dirs: []string{"0/wal"}, wantErr: "conflicts"},
77 {name: "empty marker", provider: "aof", files: map[string]string{kvProviderMarker: ""}, wantErr: "invalid storage provider marker"},
78 {name: "unknown marker", provider: "aof", files: map[string]string{kvProviderMarker: "unknown\n"}, wantErr: "invalid storage provider marker"},
79 {name: "memory marker is invalid", provider: "memory", files: map[string]string{kvProviderMarker: "memory\n"}, wantErr: "invalid storage provider marker"},
80 {name: "marker lacks newline", provider: "aof", files: map[string]string{kvProviderMarker: "aof"}, wantErr: "invalid storage provider marker"},
81 {name: "oversized marker", provider: "aof", files: map[string]string{kvProviderMarker: "aof\naof\naof\naof\naof\n"}, wantErr: "invalid storage provider marker"},
82 {name: "marker is directory", provider: "aof", dirs: []string{kvProviderMarker}, wantErr: "invalid storage provider marker"},
83 {name: "vnode is file", provider: "aof", files: map[string]string{"0": "data"}, wantErr: "reading vnode storage directory"},
84 {name: "provider path is file", provider: "aof", files: map[string]string{"0/wal": "data"}, wantErr: "expected aof storage directory"},
85 {name: "ancient direct aof segments", provider: "aof", files: map[string]string{"0/00000000000000000001": "legacy log"}, wantErr: "unrecognized storage layout"},
86 {name: "unknown provider", provider: "typo", wantErr: "unknown kv provider"},
87 }
88 for _, tt := range tests {
89 t.Run(tt.name, func(t *testing.T) {
90 dir := t.TempDir()
91 for _, path := range tt.dirs {
92 require.NoError(t, os.MkdirAll(filepath.Join(dir, path), 0750))
93 }
94 for path, data := range tt.files {
95 path = filepath.Join(dir, path)
96 require.NoError(t, os.MkdirAll(filepath.Dir(path), 0750))
97 require.NoError(t, os.WriteFile(path, []byte(data), 0600))
98 }
99 before := snapshotStorage(t, dir)
100 err := preflightKVProvider(tt.provider, dir)
101 if tt.wantErr != "" {
102 require.ErrorContains(t, err, tt.wantErr)
103 } else {
104 require.NoError(t, err)
105 if tt.provider != "memory" {
106 before[kvProviderMarker] = tt.provider + "\n"
107 }
108 }
109 require.Equal(t, before, snapshotStorage(t, dir))
110 })
111 }
112}
113 
114func TestKVProviderPreflightNewDirectory(t *testing.T) {
115 dir := filepath.Join(t.TempDir(), "new")
116 require.NoError(t, preflightKVProvider("memory", dir))
117 _, err := os.Stat(dir)
118 require.ErrorIs(t, err, os.ErrNotExist)
119 require.NoError(t, preflightKVProvider("sqlite", dir))
120 data, err := os.ReadFile(filepath.Join(dir, kvProviderMarker))
121 require.NoError(t, err)
122 require.Equal(t, "sqlite\n", string(data))
123}
124 
125func TestKVProviderConcurrentIdentityClaim(t *testing.T) {
126 dir := t.TempDir()
127 start := make(chan struct{})
128 type result struct {
129 provider string
130 err error
131 }
132 results := make(chan result, 2)
133 for _, provider := range []string{"aof", "sqlite"} {
134 go func() {
135 <-start
136 results <- result{provider, preflightKVProvider(provider, dir)}
137 }()
138 }
139 close(start)
140 first, second := <-results, <-results
141 if first.err != nil {
142 first, second = second, first
143 }
144 require.NoError(t, first.err)
145 require.Error(t, second.err)
146 data, err := os.ReadFile(filepath.Join(dir, kvProviderMarker))
147 require.NoError(t, err)
148 require.Equal(t, first.provider+"\n", string(data))
149 entries, err := os.ReadDir(dir)
150 require.NoError(t, err)
151 require.Len(t, entries, 1, "temporary identity files must be cleaned up")
152}
153 
154func TestKVProviderPreflightMarkerWriteFailure(t *testing.T) {
155 if runtime.GOOS == "windows" || os.Geteuid() == 0 {
156 t.Skip("requires Unix directory permissions without root bypass")
157 }
158 dir := t.TempDir()
159 require.NoError(t, os.MkdirAll(filepath.Join(dir, "0/wal"), 0750))
160 before := snapshotStorage(t, dir)
161 require.NoError(t, os.Chmod(dir, 0550))
162 t.Cleanup(func() { _ = os.Chmod(dir, 0750) })
163 require.ErrorContains(t, preflightKVProvider("aof", dir), "creating storage provider marker")
164 require.Equal(t, before, snapshotStorage(t, dir))
165}
166 
167func TestKVProviderPreflightMarkerReadFailure(t *testing.T) {
168 if runtime.GOOS == "windows" || os.Geteuid() == 0 {
169 t.Skip("requires Unix file permissions without root bypass")
170 }
171 dir := t.TempDir()
172 path := filepath.Join(dir, kvProviderMarker)
173 require.NoError(t, os.WriteFile(path, []byte("aof\n"), 0600))
174 require.NoError(t, os.Chmod(path, 0000))
175 t.Cleanup(func() { _ = os.Chmod(path, 0600) })
176 require.ErrorContains(t, preflightKVProvider("aof", dir), "reading storage provider marker")
177 require.NoError(t, os.Chmod(path, 0600))
178 data, err := os.ReadFile(path)
179 require.NoError(t, err)
180 require.Equal(t, "aof\n", string(data))
181}
182 
183func TestKVProviderPreflightAllowsTraversalOnlyAncestor(t *testing.T) {
184 if runtime.GOOS == "windows" || os.Geteuid() == 0 {
185 t.Skip("requires Unix directory permissions without root bypass")
186 }
187 parent := t.TempDir()
188 dir := filepath.Join(parent, "existing", "new", "data")
189 require.NoError(t, os.Mkdir(filepath.Join(parent, "existing"), 0750))
190 // The service can traverse this existing ancestor and read its own storage
191 // directory, but cannot enumerate the ancestor's other contents.
192 require.NoError(t, os.Chmod(parent, 0110))
193 t.Cleanup(func() { _ = os.Chmod(parent, 0750) })
194 require.NoError(t, preflightKVProvider("aof", dir))
195 data, err := os.ReadFile(filepath.Join(dir, kvProviderMarker))
196 require.NoError(t, err)
197 require.Equal(t, "aof\n", string(data))
198 // A normal restart should only need the marker's containing directory.
199 require.NoError(t, preflightKVProvider("aof", dir))
200}
201 
202func TestKVProviderPreflightSyncsNewDirectoryBeforeMarker(t *testing.T) {
203 if runtime.GOOS == "windows" || os.Geteuid() == 0 {
204 t.Skip("requires Unix directory permissions without root bypass")
205 }
206 parent := t.TempDir()
207 dir := filepath.Join(parent, "data")
208 // Creating an entry is allowed, but opening its parent to make that entry
209 // durable is not. Failure must happen before publishing a provider marker.
210 require.NoError(t, os.Chmod(parent, 0330))
211 t.Cleanup(func() { _ = os.Chmod(parent, 0750) })
212 require.ErrorContains(t, preflightKVProvider("aof", dir), "opening directory to sync storage identity")
213 _, err := os.Stat(dir)
214 require.ErrorIs(t, err, os.ErrNotExist, "remove only the new empty directory after its creation cannot be synced")
215 require.NoError(t, os.Chmod(parent, 0750))
216 require.NoError(t, preflightKVProvider("aof", dir))
217 data, err := os.ReadFile(filepath.Join(dir, kvProviderMarker))
218 require.NoError(t, err)
219 require.Equal(t, "aof\n", string(data))
220}
221 
222func TestKVProviderPreflightMarkerSymlink(t *testing.T) {
223 if runtime.GOOS == "windows" {
224 t.Skip("symlink creation may require elevated privileges on Windows")
225 }
226 dir := t.TempDir()
227 target := filepath.Join(t.TempDir(), "marker")
228 require.NoError(t, os.WriteFile(target, []byte("aof\n"), 0600))
229 path := filepath.Join(dir, kvProviderMarker)
230 require.NoError(t, os.Symlink(target, path))
231 require.ErrorContains(t, preflightKVProvider("aof", dir), "invalid storage provider marker")
232 link, err := os.Readlink(path)
233 require.NoError(t, err)
234 require.Equal(t, target, link)
235}
236 
237func TestServerRejectsConflictingStorageBeforeSetup(t *testing.T) {
238 dir := t.TempDir()
239 kv, stop, err := getKVProvider(zap.NewNop(), "aof", filepath.Join(dir, "9"), filepath.Join(dir, "cache"))
240 require.NoError(t, err)
241 err = kv.Put(context.Background(), []byte("existing"), []byte("preserve me"))
242 stop()
243 require.NoError(t, err)
244 before := snapshotStorage(t, dir)
245 
246 // Deliberately unusable later configuration proves preflight runs before
247 // further startup work, even for stored vnodes outside the configured count.
248 cmd := &cli.Command{
249 Name: "specter-test",
250 Metadata: map[string]any{"logger": zap.NewNop()},
251 Flags: []cli.Flag{
252 &cli.StringFlag{Name: "data-dir", Value: dir},
253 &cli.StringFlag{Name: "kv-provider", Value: "sqlite"},
254 &cli.IntFlag{Name: "virtual", Value: 1},
255 &cli.StringFlag{Name: "proxy-buffer", Value: "invalid"},
256 },
257 Action: cmdServer,
258 }
259 err = cmd.Run(t.Context(), []string{cmd.Name})
260 require.ErrorContains(t, err, "storage preflight")
261 require.ErrorContains(t, err, "aof")
262 require.ErrorContains(t, err, "sqlite")
263 require.Equal(t, before, snapshotStorage(t, dir))
264}
265 
266func snapshotStorage(t *testing.T, dir string) map[string]string {
267 t.Helper()
268 files := make(map[string]string)
269 err := filepath.WalkDir(dir, func(path string, entry fs.DirEntry, err error) error {
270 if err != nil {
271 return err
272 }
273 rel, err := filepath.Rel(dir, path)
274 if err != nil {
275 return err
276 }
277 if entry.IsDir() {
278 files[rel+"/"] = ""
279 return nil
280 }
281 data, err := os.ReadFile(path)
282 if err != nil {
283 return err
284 }
285 files[rel] = string(data)
286 return nil
287 })
288 require.NoError(t, err)
289 return files
290}