File
Blob: cmd/server/kv_provider_test.go
| 1 | package server |
| 2 | |
| 3 | import ( |
| 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 | |
| 16 | func 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 | |
| 55 | func 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 | |
| 114 | func 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 | |
| 125 | func 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 | |
| 154 | func 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 | |
| 167 | func 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 | |
| 183 | func 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 | |
| 202 | func 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 | |
| 222 | func 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 | |
| 237 | func 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 | |
| 266 | func 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 | } |