From 8a88a72e6299d77211ded4d717f9cfc94cdf43c6 Mon Sep 17 00:00:00 2001 From: cabbage Date: Mon, 31 Aug 2026 09:11:58 +0000 Subject: [PATCH] feat: bootstrap durable server store --- deploy/compose.yaml | 1 + docs/dependency-decisions.md | 1 + docs/testing.md | 19 +- go.mod | 10 +- go.sum | 46 ++++ internal/server/store/lock_unix.go | 34 +++ internal/server/store/lock_windows.go | 33 +++ internal/server/store/migrations.go | 169 +++++++++++++ internal/server/store/store.go | 149 +++++++++++ test/coverage.toml | 36 +++ test/harness/harness.go | 72 +++++- test/harness/harness_test.go | 28 +++ .../store/store_integration_test.go | 236 ++++++++++++++++++ 13 files changed, 822 insertions(+), 12 deletions(-) create mode 100644 internal/server/store/lock_unix.go create mode 100644 internal/server/store/lock_windows.go create mode 100644 internal/server/store/migrations.go create mode 100644 internal/server/store/store.go create mode 100644 test/integration/store/store_integration_test.go diff --git a/deploy/compose.yaml b/deploy/compose.yaml index 7694c1b..2dddcb2 100644 --- a/deploy/compose.yaml +++ b/deploy/compose.yaml @@ -12,6 +12,7 @@ services: GOCACHE: /cache/go-build GOMODCACHE: /cache/go/pkg/mod GOPATH: /cache/go + GOFLAGS: -p=2 HOME: /tmp XDG_CACHE_HOME: /cache/xdg volumes: diff --git a/docs/dependency-decisions.md b/docs/dependency-decisions.md index 850343d..11ff156 100644 --- a/docs/dependency-decisions.md +++ b/docs/dependency-decisions.md @@ -12,3 +12,4 @@ All versions are exact in `go.mod`, generated code, or the toolchain image. | google/uuid 1.6.0 | Parse canonical UUIDs and verify RFC variant/version bits | Stable maintained package; RVBox owns the monotonic UUIDv7 generator so clock and ordering behavior remain directly testable. | | go-toml/v2 2.3.1 | Strict configuration decoding | Last maintained release line before TOML 1.1 parsing was enabled; RVBox v1 intentionally accepts TOML 1.0 only. | | klauspost/compress 1.19.0 | Zstandard command-output compression | Maintained pure-Go codec with decoder memory controls; RVBox additionally limits streamed decoded output before allocation. | +| modernc.org/sqlite 1.57.0 | Durable server/client metadata | Maintained CGO-free SQLite driver with Linux/Windows support and defensive-mode DSN support; its exact generated-code-matched libc version is pinned by `go.mod`. | diff --git a/docs/testing.md b/docs/testing.md index b5b1101..425111e 100644 --- a/docs/testing.md +++ b/docs/testing.md @@ -15,9 +15,10 @@ Run focused unit tests with: scripts/test-unit --package ./internal/domain --run UUIDv7 --race ``` -The integration harness currently provides the Phase 0 `sample` suite. It -proves run isolation and the durable lifecycle without starting an RVBox -service that has not been implemented yet: +The integration harness provides the Phase 0 `sample` suite and the incremental +Phase 2 `store` suite. The latter uses a real temporary SQLite database in WAL +mode and a real segment/audit filesystem; it does not mock the persistence +boundary: ```sh scripts/test-env doctor @@ -29,8 +30,19 @@ scripts/test-env reuse --run-id my-sample scripts/test-integration --suite sample --run-id my-sample --resume scripts/test-env reset --run-id my-sample scripts/test-env purge --run-id my-sample + +scripts/test-integration --suite store --run-id store-smoke +scripts/test-env logs --run-id store-smoke +scripts/test-env collect --run-id store-smoke +scripts/test-env reset --run-id store-smoke +scripts/test-env purge --run-id store-smoke ``` +Suite output is capped at 1 MiB and stored as `artifacts/suite.log`. A failed +run remains inspectable and can be moved back to `ready` with `recover`, then +resumed with the same run ID and deterministic shuffle seed. Test-run cleanup +never removes the shared Go module or build-cache volumes. + Each run owns only `.test-runs/` and resources explicitly recorded in that run's versioned manifest. The journal is append-only and fsynced. `purge` validates the run ID and manifest identity, refuses symlink targets or manifests @@ -43,4 +55,3 @@ cleanup. test reference against source. Native Windows integration/E2E entries remain explicitly blocked until the resettable Windows host is available; Wine or a protocol stub is not treated as equivalent coverage. - diff --git a/go.mod b/go.mod index e1e4267..8e78ff5 100644 --- a/go.mod +++ b/go.mod @@ -6,13 +6,21 @@ require ( github.com/google/uuid v1.6.0 github.com/klauspost/compress v1.19.0 github.com/pelletier/go-toml/v2 v2.3.1 + golang.org/x/sys v0.47.0 google.golang.org/grpc v1.83.2 google.golang.org/protobuf v1.36.12 + modernc.org/sqlite v1.57.0 ) require ( + github.com/dustin/go-humanize v1.0.1 // indirect + github.com/mattn/go-isatty v0.0.24 // indirect + github.com/ncruces/go-strftime v1.0.0 // indirect + github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec // indirect golang.org/x/net v0.58.0 // indirect - golang.org/x/sys v0.47.0 // indirect golang.org/x/text v0.41.0 // indirect google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa // indirect + modernc.org/libc v1.74.4 // indirect + modernc.org/mathutil v1.7.1 // indirect + modernc.org/memory v1.11.0 // indirect ) diff --git a/go.sum b/go.sum index 9e78a30..78b969e 100644 --- a/go.sum +++ b/go.sum @@ -1,5 +1,7 @@ github.com/cespare/xxhash/v2 v2.3.0 h1:UL815xU9SqsFlibzuggzjXhog7bL6oX9BbNZnL2UFvs= github.com/cespare/xxhash/v2 v2.3.0/go.mod h1:VGX0DQ3Q6kWi7AoAeZDth3/j3BFtOZR5XLFGgcrjCOs= +github.com/dustin/go-humanize v1.0.1 h1:GzkhY7T5VNhEkwH0PVJgjz+fX1rhBrR7pRT3mDkpeCY= +github.com/dustin/go-humanize v1.0.1/go.mod h1:Mu1zIs6XwVuF/gI1OepvI0qD18qycQx+mFykh5fBlto= github.com/go-logr/logr v1.4.3 h1:CjnDlHq8ikf6E492q6eKboGOC0T8CDaOvkHCIg8idEI= github.com/go-logr/logr v1.4.3/go.mod h1:9T104GzyrTigFIr8wt5mBrctHMim0Nb2HLGrmQ40KvY= github.com/go-logr/stdr v1.2.2 h1:hSWxHoqTgW2S2qGc0LTAI563KZ5YKYRhT3MFKZMbjag= @@ -8,12 +10,22 @@ github.com/golang/protobuf v1.5.4 h1:i7eJL8qZTpSEXOPTxNKhASYpMn+8e5Q6AdndVa1dWek github.com/golang/protobuf v1.5.4/go.mod h1:lnTiLA8Wa4RWRcIUkrtSVa5nRhsEGBg48fD6rSs7xps= github.com/google/go-cmp v0.7.0 h1:wk8382ETsv4JYUZwIsn6YpYiWiBsYLSJiTsyBybVuN8= github.com/google/go-cmp v0.7.0/go.mod h1:pXiqmnSA92OHEEa9HXL2W4E7lf9JzCmGVUdgjX3N/iU= +github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3 h1:LMLX+LgTNWpfvCBdFebv6EsYotImrt/Ppc5cXIriCSo= +github.com/google/pprof v0.0.0-20260802141513-ef3492d7dac3/go.mod h1:jl5iWTm0/hd5PjEYEOuwAJ57L/CibdZfrqZ5XA5GrCk= github.com/google/uuid v1.6.0 h1:NIvaJDMOsjHA8n1jAhLSgzrAzy1Hgr+hNrb57e+94F0= github.com/google/uuid v1.6.0/go.mod h1:TIyPZe4MgqvfeYDBFedMoGGpEw/LqOeaOT+nhxU+yHo= +github.com/hashicorp/golang-lru/v2 v2.0.7 h1:a+bsQ5rvGLjzHuww6tVxozPZFVghXaHOwFs4luLUK2k= +github.com/hashicorp/golang-lru/v2 v2.0.7/go.mod h1:QeFd9opnmA6QUJc5vARoKUSoFhyfM2/ZepoAG6RGpeM= github.com/klauspost/compress v1.19.0 h1:sXLILfc9jV2QYWkzFOPWStmcUVH2RHEB1JCdY2oVvCQ= github.com/klauspost/compress v1.19.0/go.mod h1:cwPg85FWrGar70rWktvGQj8/hthj3wpl0PGDogxkrSQ= +github.com/mattn/go-isatty v0.0.24 h1:tGZZoVgT/KiqK1c8ocVLeDS8BSWMRd47J3Lbz7vsReI= +github.com/mattn/go-isatty v0.0.24/go.mod h1:nMCL3Zebbrt45jsMDgnfIwz6ydEQApk5oEI3HqDio6A= +github.com/ncruces/go-strftime v1.0.0 h1:HMFp8mLCTPp341M/ZnA4qaf7ZlsbTc+miZjCLOFAw7w= +github.com/ncruces/go-strftime v1.0.0/go.mod h1:Fwc5htZGVVkseilnfgOVb9mKy6w1naJmn9CehxcKcls= github.com/pelletier/go-toml/v2 v2.3.1 h1:MYEvvGnQjeNkRF1qUuGolNtNExTDwct51yp7olPtrEc= github.com/pelletier/go-toml/v2 v2.3.1/go.mod h1:2gIqNv+qfxSVS7cM2xJQKtLSTLUE9V8t9Stt+h56mCY= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec h1:W09IVJc94icq4NjY3clb7Lk8O1qJ8BdBEF8z0ibU0rE= +github.com/remyoudompheng/bigfft v0.0.0-20230129092748-24d4a6f8daec/go.mod h1:qqbHyh8v60DhA7CoWK5oRCqLrMHRGoxYCSS9EjAz6Eo= go.opentelemetry.io/auto/sdk v1.2.1 h1:jXsnJ4Lmnqd11kwkBV2LgLoFMZKizbCi5fNZ/ipaZ64= go.opentelemetry.io/auto/sdk v1.2.1/go.mod h1:KRTj+aOaElaLi+wW1kO/DZRXwkF4C5xPbEe3ZiIhN7Y= go.opentelemetry.io/otel v1.44.0 h1:JjwHmHpA4iZ3wBxluu2fbbE7j4kqlE8jXyAyPXH7HqU= @@ -26,12 +38,18 @@ go.opentelemetry.io/otel/sdk/metric v1.44.0 h1:3LlKgI+VjbVsjNRFZJZAJ30WjXC5VkNRk go.opentelemetry.io/otel/sdk/metric v1.44.0/go.mod h1:5B5pMARnXxKhltooO4xUuCBorl65a4EpnTalObqOigA= go.opentelemetry.io/otel/trace v1.44.0 h1:jxF5CsGYCe74MCRx2X4g7WsY/VBKRqqpNvXlX/6gtIk= go.opentelemetry.io/otel/trace v1.44.0/go.mod h1:oLl1jrMQAVo6v3GAggN+1VH9VIz9iUSvW53sW1Q8PIE= +golang.org/x/mod v0.38.0 h1:MECBjubtXD7yj4HrhIUcywNaGeNVUdfVnxmPajOk4yk= +golang.org/x/mod v0.38.0/go.mod h1:V6Xz0pq8TQ3dGqVQ1FVHuelZpAL0uNhSkk9ogYP3c40= golang.org/x/net v0.58.0 h1:ynWG7rqYi4ccpTEuPZ2QGWHktVEM9DMCj9yzDE0Q7To= golang.org/x/net v0.58.0/go.mod h1:YwCddHnFlT7eLQqVprV19OnhLGtc5xOKgE0RyqgfWAU= +golang.org/x/sync v0.22.0 h1:SZjpbeLmrCk4xhRSZFNZW5gFUeCeFgjekvI/+gfScek= +golang.org/x/sync v0.22.0/go.mod h1:9xrNwdLfx4jkKbNva9FpL6vEN7evnE43NNNJQ2LF3+0= golang.org/x/sys v0.47.0 h1:o7XGOvZQCADBQQ4Y7VNq2dRWQR7JmOUW8Kxx4ZsNgWs= golang.org/x/sys v0.47.0/go.mod h1:4GL1E5IUh+htKOUEOaiffhrAeqysfVGipDYzABqnCmw= golang.org/x/text v0.41.0 h1:vz/seA0lnX87Othu2f/0L24RcgrXD9/YFTSuGjj3rH8= golang.org/x/text v0.41.0/go.mod h1:jvf1O8ajNzZqhSrQBPbutR/EB83Cc0CFrezNQIwbb5M= +golang.org/x/tools v0.48.0 h1:3+hClM1aLL5mjMKm5ovokw9epgRXPuu2tILgismM6RE= +golang.org/x/tools v0.48.0/go.mod h1:08xX0orndb/F7jJxGDicx061tyd5pcMto75YMAXr6lk= gonum.org/v1/gonum v0.17.0 h1:VbpOemQlsSMrYmn7T2OUvQ4dqxQXU+ouZFQsZOx50z4= gonum.org/v1/gonum v0.17.0/go.mod h1:El3tOrEuMpv2UdMrbNlKEh9vd86bmQ6vqIcDwxEOc1E= google.golang.org/genproto/googleapis/rpc v0.0.0-20260526163538-3dc84a4a5aaa h1:mZHHdPZl0dbGHCflZgAq/Q468DWVFcU2whhB2KAo8fk= @@ -40,3 +58,31 @@ google.golang.org/grpc v1.83.2 h1:EManeRomTObA0BU7I8vXgg/78uE5MJ9M8B39EX2WscU= google.golang.org/grpc v1.83.2/go.mod h1:YPI1hK3kDked6iHvgX3tR0y+nX/qpMFKhPgFsokw1S8= google.golang.org/protobuf v1.36.12 h1:pJOKDDOyeXErUroCihFAd5LQuwXBSpVnKGrj5o/fwxc= google.golang.org/protobuf v1.36.12/go.mod h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco= +modernc.org/cc/v4 v4.29.1 h1:MKgdCV3WykTSPqpVrnxdEDS0HEd2FHpKZDzxzU5LyeI= +modernc.org/cc/v4 v4.29.1/go.mod h1:OnovgIhbbMXMu1aISnJ0wvVD1KnW+cAUJkIrAWh+kVI= +modernc.org/ccgo/v4 v4.34.6 h1:sBgfIwyN0TQ9C5hwIeuqyeAKyMWnbvj2fvpF4L11uzU= +modernc.org/ccgo/v4 v4.34.6/go.mod h1:SZ8YcN9NG7XVsQYdm6jYBvi8PQP1qi+kqB6OhjqI3Fk= +modernc.org/fileutil v1.4.0 h1:j6ZzNTftVS054gi281TyLjHPp6CPHr2KCxEXjEbD6SM= +modernc.org/fileutil v1.4.0/go.mod h1:EqdKFDxiByqxLk8ozOxObDSfcVOv/54xDs/DUHdvCUU= +modernc.org/gc/v2 v2.6.5 h1:nyqdV8q46KvTpZlsw66kWqwXRHdjIlJOhG6kxiV/9xI= +modernc.org/gc/v2 v2.6.5/go.mod h1:YgIahr1ypgfe7chRuJi2gD7DBQiKSLMPgBQe9oIiito= +modernc.org/gc/v3 v3.1.4 h1:2g65LGVSmFQrXeITAw97x7hCRvZFcyE1uDP+7Vng7JI= +modernc.org/gc/v3 v3.1.4/go.mod h1:HFK/6AGESC7Ex+EZJhJ2Gni6cTaYpSMmU/cT9RmlfYY= +modernc.org/goabi0 v0.2.0 h1:HvEowk7LxcPd0eq6mVOAEMai46V+i7Jrj13t4AzuNks= +modernc.org/goabi0 v0.2.0/go.mod h1:CEFRnnJhKvWT1c1JTI3Avm+tgOWbkOu5oPA8eH8LnMI= +modernc.org/libc v1.74.4 h1:fX1Omw4o2/1C2iRkkIsrQTasJQldLhRmuPreXLoWs9k= +modernc.org/libc v1.74.4/go.mod h1:eeQAS9W3sZeKYMFubydxJpII9ybHWshk+7or7bLG9co= +modernc.org/mathutil v1.7.1 h1:GCZVGXdaN8gTqB1Mf/usp1Y/hSqgI2vAGGP4jZMCxOU= +modernc.org/mathutil v1.7.1/go.mod h1:4p5IwJITfppl0G4sUEDtCr4DthTaT47/N3aT6MhfgJg= +modernc.org/memory v1.11.0 h1:o4QC8aMQzmcwCK3t3Ux/ZHmwFPzE6hf2Y5LbkRs+hbI= +modernc.org/memory v1.11.0/go.mod h1:/JP4VbVC+K5sU2wZi9bHoq2MAkCnrt2r98UGeSK7Mjw= +modernc.org/opt v0.2.0 h1:tGyef5ApycA7FSEOMraay9SaTk5zmbx7Tu+cJs4QKZg= +modernc.org/opt v0.2.0/go.mod h1:03fq9lsNfvkYSfxrfUhZCWPk1lm4cq4N+Bh//bEtgns= +modernc.org/sortutil v1.2.1 h1:+xyoGf15mM3NMlPDnFqrteY07klSFxLElE2PVuWIJ7w= +modernc.org/sortutil v1.2.1/go.mod h1:7ZI3a3REbai7gzCLcotuw9AC4VZVpYMjDzETGsSMqJE= +modernc.org/sqlite v1.57.0 h1:qNQP6xnx5M0ISNtlnxoOX0+cD5bJ0/gr9aMmndFczzg= +modernc.org/sqlite v1.57.0/go.mod h1:yCJ2cmAaIkHQ25oXWrF8H4O1lIfPYPR26yCEDj2P3pQ= +modernc.org/strutil v1.2.1 h1:UneZBkQA+DX2Rp35KcM69cSsNES9ly8mQWD71HKlOA0= +modernc.org/strutil v1.2.1/go.mod h1:EHkiggD70koQxjVdSBM3JKM7k6L0FbGE5eymy9i3B9A= +modernc.org/token v1.1.0 h1:Xl7Ap9dKaEs5kLoOQeQmPWevfnk/DM5qcLcYlA8ys6Y= +modernc.org/token v1.1.0/go.mod h1:UGzOrNV1mAFSEB63lOFHIpNRUVMvYTc6yu1SMY/XTDM= diff --git a/internal/server/store/lock_unix.go b/internal/server/store/lock_unix.go new file mode 100644 index 0000000..d4d8cc3 --- /dev/null +++ b/internal/server/store/lock_unix.go @@ -0,0 +1,34 @@ +//go:build !windows + +package store + +import ( + "errors" + "os" + "syscall" +) + +func acquireInstanceLock(path string) (func() error, error) { + if err := ensurePrivateFile(path); err != nil { + return nil, err + } + file, err := os.OpenFile(path, os.O_RDWR, 0) + if err != nil { + return nil, err + } + if err := syscall.Flock(int(file.Fd()), syscall.LOCK_EX|syscall.LOCK_NB); err != nil { + _ = file.Close() + if errors.Is(err, syscall.EWOULDBLOCK) { + return nil, ErrAlreadyOpen + } + return nil, err + } + return func() error { + unlockErr := syscall.Flock(int(file.Fd()), syscall.LOCK_UN) + closeErr := file.Close() + if unlockErr != nil { + return unlockErr + } + return closeErr + }, nil +} diff --git a/internal/server/store/lock_windows.go b/internal/server/store/lock_windows.go new file mode 100644 index 0000000..d0dbc33 --- /dev/null +++ b/internal/server/store/lock_windows.go @@ -0,0 +1,33 @@ +//go:build windows + +package store + +import ( + "os" + + "golang.org/x/sys/windows" +) + +func acquireInstanceLock(path string) (func() error, error) { + if err := ensurePrivateFile(path); err != nil { + return nil, err + } + file, err := os.OpenFile(path, os.O_RDWR, 0) + if err != nil { + return nil, err + } + var overlapped windows.Overlapped + err = windows.LockFileEx(windows.Handle(file.Fd()), windows.LOCKFILE_EXCLUSIVE_LOCK|windows.LOCKFILE_FAIL_IMMEDIATELY, 0, 1, 0, &overlapped) + if err != nil { + _ = file.Close() + return nil, ErrAlreadyOpen + } + return func() error { + unlockErr := windows.UnlockFileEx(windows.Handle(file.Fd()), 0, 1, 0, &overlapped) + closeErr := file.Close() + if unlockErr != nil { + return unlockErr + } + return closeErr + }, nil +} diff --git a/internal/server/store/migrations.go b/internal/server/store/migrations.go new file mode 100644 index 0000000..dd11371 --- /dev/null +++ b/internal/server/store/migrations.go @@ -0,0 +1,169 @@ +package store + +import ( + "context" + "crypto/sha256" + "database/sql" + "encoding/hex" + "fmt" +) + +type migration struct { + version uint32 + sql string +} + +var migrations = []migration{{version: 1, sql: schemaV1}} + +func applyMigrations(ctx context.Context, db *sql.DB) error { + if _, err := db.ExecContext(ctx, `CREATE TABLE IF NOT EXISTS schema_migrations ( +version INTEGER PRIMARY KEY CHECK(version > 0), checksum TEXT NOT NULL, applied_at INTEGER NOT NULL +) STRICT`); err != nil { + return fmt.Errorf("create migration table: %w", err) + } + for _, current := range migrations { + checksumBytes := sha256.Sum256([]byte(current.sql)) + checksum := hex.EncodeToString(checksumBytes[:]) + var stored string + err := db.QueryRowContext(ctx, `SELECT checksum FROM schema_migrations WHERE version = ?`, current.version).Scan(&stored) + if err == nil { + if stored != checksum { + return fmt.Errorf("migration %d checksum mismatch", current.version) + } + continue + } + if err != sql.ErrNoRows { + return err + } + tx, err := db.BeginTx(ctx, nil) + if err != nil { + return err + } + if _, err = tx.ExecContext(ctx, current.sql); err == nil { + _, err = tx.ExecContext(ctx, `INSERT INTO schema_migrations(version, checksum, applied_at) VALUES (?, ?, unixepoch())`, current.version, checksum) + } + if err != nil { + _ = tx.Rollback() + return fmt.Errorf("apply migration %d: %w", current.version, err) + } + if err := tx.Commit(); err != nil { + return fmt.Errorf("commit migration %d: %w", current.version, err) + } + } + return nil +} + +const schemaV1 = ` +CREATE TABLE clients ( + client_id TEXT PRIMARY KEY CHECK(length(client_id) BETWEEN 1 AND 128), + platform INTEGER NOT NULL, architecture TEXT NOT NULL, daemon_version TEXT NOT NULL, + daemon_cwd TEXT NOT NULL, supported_shells BLOB NOT NULL, capabilities BLOB NOT NULL, + client_instance_id BLOB NOT NULL CHECK(length(client_instance_id) = 16), + generation INTEGER NOT NULL DEFAULT 0 CHECK(generation >= 0), + connected_at INTEGER, last_seen_at INTEGER, + pending_instance_id BLOB CHECK(pending_instance_id IS NULL OR length(pending_instance_id) = 16), pending_instance_seen_at INTEGER, + charged_bytes INTEGER NOT NULL DEFAULT 0 CHECK(charged_bytes >= 0) +) STRICT; +CREATE TABLE sessions ( + session_id BLOB PRIMARY KEY CHECK(length(session_id) = 16), client_id TEXT NOT NULL REFERENCES clients(client_id) ON DELETE CASCADE, + client_instance_id BLOB NOT NULL CHECK(length(client_instance_id) = 16), generation INTEGER NOT NULL CHECK(generation > 0), + opened_at INTEGER NOT NULL, fenced_at INTEGER, closed_at INTEGER, close_reason TEXT +) STRICT; +CREATE UNIQUE INDEX one_live_session_per_client ON sessions(client_id) WHERE closed_at IS NULL AND fenced_at IS NULL; +CREATE TABLE commands ( + issue_uuid BLOB PRIMARY KEY CHECK(length(issue_uuid) = 16), + client_id TEXT NOT NULL REFERENCES clients(client_id) ON DELETE RESTRICT, + issue_time INTEGER NOT NULL, server_receipt_time INTEGER NOT NULL, queue_expiry_time INTEGER, terminal_time INTEGER, + lifecycle INTEGER NOT NULL CHECK(lifecycle BETWEEN 1 AND 11), + revision INTEGER NOT NULL CHECK(revision > 0), exit_code INTEGER, + retention_status INTEGER NOT NULL DEFAULT 1 CHECK(retention_status BETWEEN 1 AND 3), + last_event_seq INTEGER NOT NULL DEFAULT 0 CHECK(last_event_seq >= 0), + retained_compressed_bytes INTEGER NOT NULL DEFAULT 0 CHECK(retained_compressed_bytes >= 0), + output_truncated INTEGER NOT NULL DEFAULT 0 CHECK(output_truncated IN (0,1)), + output_incomplete INTEGER NOT NULL DEFAULT 0 CHECK(output_incomplete IN (0,1)), + immutable_request_sha256 BLOB NOT NULL CHECK(length(immutable_request_sha256) = 32), + execution_spec BLOB NOT NULL, execution_spec_raw_bytes INTEGER NOT NULL CHECK(execution_spec_raw_bytes >= 0), + execution_spec_stored_bytes INTEGER NOT NULL CHECK(execution_spec_stored_bytes >= 0), + execution_spec_compression INTEGER NOT NULL CHECK(execution_spec_compression IN (1,2)), + windows_execution_identity BLOB, + CHECK((lifecycle BETWEEN 5 AND 11 AND terminal_time IS NOT NULL) OR (lifecycle BETWEEN 1 AND 4 AND terminal_time IS NULL)) +) STRICT; +CREATE INDEX commands_by_client_time ON commands(client_id, issue_time DESC, issue_uuid DESC); +CREATE INDEX commands_dispatch ON commands(client_id, lifecycle, issue_time, issue_uuid); +CREATE TABLE command_payloads ( + issue_uuid BLOB NOT NULL REFERENCES commands(issue_uuid) ON DELETE CASCADE, + kind TEXT NOT NULL, raw_bytes INTEGER NOT NULL CHECK(raw_bytes >= 0), stored_bytes INTEGER NOT NULL CHECK(stored_bytes >= 0), + compression INTEGER NOT NULL CHECK(compression IN (1,2)), sha256 BLOB NOT NULL CHECK(length(sha256) = 32), + inline_data BLOB, segment_path TEXT, + PRIMARY KEY(issue_uuid, kind), CHECK((inline_data IS NULL) != (segment_path IS NULL)) +) STRICT, WITHOUT ROWID; +CREATE TABLE command_events ( + issue_uuid BLOB NOT NULL REFERENCES commands(issue_uuid) ON DELETE CASCADE, + event_seq INTEGER NOT NULL CHECK(event_seq > 0), observed_at INTEGER NOT NULL, server_receipt_time INTEGER NOT NULL, + event_type INTEGER NOT NULL, compression INTEGER NOT NULL CHECK(compression IN (1,2)), + raw_bytes INTEGER NOT NULL CHECK(raw_bytes >= 0), stored_bytes INTEGER NOT NULL CHECK(stored_bytes >= 0), + payload BLOB NOT NULL, immutable_sha256 BLOB NOT NULL CHECK(length(immutable_sha256) = 32), + PRIMARY KEY(issue_uuid, event_seq) +) STRICT, WITHOUT ROWID; +CREATE TABLE output_segments ( + issue_uuid BLOB NOT NULL REFERENCES commands(issue_uuid) ON DELETE CASCADE, + ordinal INTEGER NOT NULL CHECK(ordinal >= 0), path TEXT NOT NULL, committed_end_offset INTEGER NOT NULL CHECK(committed_end_offset >= 0), + min_event_seq INTEGER NOT NULL, max_event_seq INTEGER NOT NULL CHECK(max_event_seq >= min_event_seq), + stream_mix INTEGER NOT NULL CHECK(stream_mix >= 0), + compressed_bytes INTEGER NOT NULL CHECK(compressed_bytes >= 0), raw_bytes INTEGER NOT NULL CHECK(raw_bytes >= 0), + checksum BLOB NOT NULL CHECK(length(checksum) = 32), created_at INTEGER NOT NULL, + PRIMARY KEY(issue_uuid, ordinal), UNIQUE(path) +) STRICT, WITHOUT ROWID; +CREATE TABLE output_truncations ( + issue_uuid BLOB NOT NULL REFERENCES commands(issue_uuid) ON DELETE CASCADE, + first_event_seq INTEGER, last_event_seq INTEGER, removed_compressed_bytes INTEGER, removed_raw_bytes INTEGER NOT NULL, + source INTEGER NOT NULL, reason TEXT NOT NULL, recorded_at INTEGER NOT NULL, + CHECK(removed_compressed_bytes IS NULL OR removed_compressed_bytes >= 0), CHECK(removed_raw_bytes >= 0), + CHECK((first_event_seq IS NULL AND last_event_seq IS NULL) OR (first_event_seq > 0 AND last_event_seq >= first_event_seq)) +) STRICT; +CREATE TABLE stdin_writes ( + issue_uuid BLOB NOT NULL REFERENCES commands(issue_uuid) ON DELETE CASCADE, + write_seq INTEGER NOT NULL CHECK(write_seq > 0), payload BLOB, raw_bytes INTEGER NOT NULL CHECK(raw_bytes >= 0), + stored_bytes INTEGER NOT NULL CHECK(stored_bytes >= 0), compression INTEGER NOT NULL CHECK(compression IN (1,2)), + sha256 BLOB NOT NULL CHECK(length(sha256) = 32), append_newline INTEGER NOT NULL CHECK(append_newline IN (0,1)), + close_intent INTEGER NOT NULL CHECK(close_intent IN (0,1)), acknowledged INTEGER NOT NULL CHECK(acknowledged IN (0,1)), + PRIMARY KEY(issue_uuid, write_seq) +) STRICT, WITHOUT ROWID; +CREATE TABLE control_mutations ( + request_uuid BLOB PRIMARY KEY CHECK(length(request_uuid) = 16), method TEXT NOT NULL, + owner_kind TEXT NOT NULL, owner_id TEXT NOT NULL, target TEXT NOT NULL, + immutable_sha256 BLOB NOT NULL CHECK(length(immutable_sha256) = 32), assigned_write_seq INTEGER, + assigned_revision INTEGER, result BLOB NOT NULL, created_at INTEGER NOT NULL, + CHECK(assigned_write_seq IS NULL OR assigned_write_seq > 0), CHECK(assigned_revision IS NULL OR assigned_revision > 0) +) STRICT; +CREATE TABLE takeover_authorizations ( + client_id TEXT PRIMARY KEY REFERENCES clients(client_id) ON DELETE CASCADE, + pending_instance_id BLOB NOT NULL CHECK(length(pending_instance_id) = 16), + request_uuid BLOB NOT NULL UNIQUE CHECK(length(request_uuid) = 16), created_at INTEGER NOT NULL, expires_at INTEGER NOT NULL, + consumed_at INTEGER, CHECK(expires_at > created_at) +) STRICT; +CREATE TABLE command_tombstones ( + issue_uuid BLOB PRIMARY KEY CHECK(length(issue_uuid) = 16), immutable_sha256 BLOB NOT NULL CHECK(length(immutable_sha256) = 32), + client_id TEXT NOT NULL, terminal_lifecycle INTEGER NOT NULL CHECK(terminal_lifecycle BETWEEN 5 AND 11), + terminal_time INTEGER NOT NULL, acknowledged_at INTEGER NOT NULL +) STRICT; +CREATE INDEX tombstones_fifo ON command_tombstones(acknowledged_at, issue_uuid); +CREATE TABLE audit_events ( + audit_id INTEGER PRIMARY KEY, occurred_at INTEGER NOT NULL, source TEXT NOT NULL, principal TEXT, client_id TEXT, + issue_uuid BLOB CHECK(issue_uuid IS NULL OR length(issue_uuid) = 16), action TEXT NOT NULL, outcome TEXT NOT NULL, + error_code INTEGER, compression INTEGER NOT NULL CHECK(compression IN (1,2)), payload BLOB NOT NULL, + raw_bytes INTEGER NOT NULL CHECK(raw_bytes >= 0), stored_bytes INTEGER NOT NULL CHECK(stored_bytes >= 0), + sha256 BLOB NOT NULL CHECK(length(sha256) = 32) +) STRICT; +CREATE INDEX audit_fifo ON audit_events(occurred_at, audit_id); +CREATE TABLE storage_incidents ( + incident_uuid BLOB PRIMARY KEY CHECK(length(incident_uuid) = 16), detected_at INTEGER NOT NULL, resolved_at INTEGER, + state INTEGER NOT NULL CHECK(state BETWEEN 1 AND 3), kind INTEGER NOT NULL, scope TEXT NOT NULL, scope_key TEXT NOT NULL, + client_id TEXT, issue_uuid BLOB CHECK(issue_uuid IS NULL OR length(issue_uuid) = 16), summary TEXT NOT NULL, + evidence BLOB NOT NULL, resolution_note TEXT CHECK(resolution_note IS NULL OR length(resolution_note) <= 4096), + data_loss INTEGER NOT NULL CHECK(data_loss IN (0,1)), + automatically_repairable INTEGER NOT NULL CHECK(automatically_repairable IN (0,1)), + CHECK((state = 1 AND resolved_at IS NULL) OR (state IN (2,3) AND resolved_at IS NOT NULL)) +) STRICT; +CREATE UNIQUE INDEX one_open_incident_per_scope_kind ON storage_incidents(scope, scope_key, kind) WHERE state = 1; +` diff --git a/internal/server/store/store.go b/internal/server/store/store.go new file mode 100644 index 0000000..2099083 --- /dev/null +++ b/internal/server/store/store.go @@ -0,0 +1,149 @@ +// Package store owns RVBox server persistence and migrations. +package store + +import ( + "context" + "database/sql" + "errors" + "fmt" + "net/url" + "os" + "path/filepath" + "strconv" + "sync" + "time" + + _ "modernc.org/sqlite" +) + +var ( + ErrUnsafeDataDirectory = errors.New("unsafe server data directory") + ErrAlreadyOpen = errors.New("server data directory is already locked") +) + +type Options struct { + DataDir string + BusyTimeout time.Duration +} + +type Store struct { + db *sql.DB + unlock func() error + mu sync.Mutex +} + +func Open(ctx context.Context, options Options) (*Store, error) { + if !filepath.IsAbs(options.DataDir) || filepath.Clean(options.DataDir) == string(filepath.Separator) { + return nil, ErrUnsafeDataDirectory + } + if options.BusyTimeout <= 0 { + return nil, fmt.Errorf("busy timeout must be positive") + } + if err := ensurePrivateDirectory(options.DataDir); err != nil { + return nil, err + } + for _, child := range []string{"segments", "audit"} { + if err := ensurePrivateDirectory(filepath.Join(options.DataDir, child)); err != nil { + return nil, err + } + } + unlock, err := acquireInstanceLock(filepath.Join(options.DataDir, "server.lock")) + if err != nil { + return nil, err + } + + databasePath := filepath.Join(options.DataDir, "rvbox.db") + if err := ensurePrivateFile(databasePath); err != nil { + _ = unlock() + return nil, err + } + query := url.Values{} + query.Add("_defensive", "1") + query.Add("_pragma", "journal_mode(WAL)") + query.Add("_pragma", "foreign_keys(ON)") + query.Add("_pragma", "synchronous(FULL)") + query.Add("_pragma", "busy_timeout("+strconv.FormatInt(options.BusyTimeout.Milliseconds(), 10)+")") + databaseURL := &url.URL{Scheme: "file", Path: filepath.ToSlash(databasePath)} + databaseURL.RawQuery = query.Encode() + dsn := databaseURL.String() + db, err := sql.Open("sqlite", dsn) + if err != nil { + _ = unlock() + return nil, err + } + db.SetMaxOpenConns(1) + db.SetMaxIdleConns(1) + store := &Store{db: db, unlock: unlock} + if err := db.PingContext(ctx); err != nil { + _ = store.Close() + return nil, fmt.Errorf("open SQLite: %w", err) + } + if err := applyMigrations(ctx, db); err != nil { + _ = store.Close() + return nil, err + } + return store, nil +} + +func (store *Store) DB() *sql.DB { return store.db } + +func (store *Store) Close() error { + if store == nil { + return nil + } + store.mu.Lock() + defer store.mu.Unlock() + var result error + if store.db != nil { + result = store.db.Close() + store.db = nil + } + if store.unlock != nil { + if err := store.unlock(); result == nil { + result = err + } + store.unlock = nil + } + return result +} + +func ensurePrivateFile(path string) error { + file, err := os.OpenFile(path, os.O_CREATE|os.O_EXCL|os.O_RDWR, 0o600) + if err == nil { + return file.Close() + } + if !errors.Is(err, os.ErrExist) { + return err + } + info, err := os.Lstat(path) + if err != nil { + return err + } + if !info.Mode().IsRegular() || info.Mode()&os.ModeSymlink != 0 { + return fmt.Errorf("%w: %s is not a regular file", ErrUnsafeDataDirectory, path) + } + if info.Mode().Perm()&0o077 != 0 { + return fmt.Errorf("%w: %s permissions %04o expose private state", ErrUnsafeDataDirectory, path, info.Mode().Perm()) + } + return nil +} + +func ensurePrivateDirectory(path string) error { + info, err := os.Lstat(path) + if os.IsNotExist(err) { + if err := os.MkdirAll(path, 0o700); err != nil { + return err + } + info, err = os.Lstat(path) + } + if err != nil { + return err + } + if info.Mode()&os.ModeSymlink != 0 || !info.IsDir() { + return fmt.Errorf("%w: %s is not a real directory", ErrUnsafeDataDirectory, path) + } + if info.Mode().Perm()&0o077 != 0 { + return fmt.Errorf("%w: %s permissions %04o expose private state", ErrUnsafeDataDirectory, path, info.Mode().Perm()) + } + return nil +} diff --git a/test/coverage.toml b/test/coverage.toml index af09a5c..d91365b 100644 --- a/test/coverage.toml +++ b/test/coverage.toml @@ -165,3 +165,39 @@ id = "HP-WINCTX-02" layer = "integration" status = "blocked_native_windows" tests = [] + +[[requirements]] +id = "HP-STORE-01" +layer = "integration" +status = "implemented" +tests = ["test/integration/store/store_integration_test.go:TestRealSQLiteInitializationAndRestart_HP_STORE_01"] + +[[requirements]] +id = "BH-STORE-01" +layer = "integration" +status = "implemented" +tests = ["test/integration/store/store_integration_test.go:TestSchemaRejectsInvalidIDsAndForeignKeys_BH_STORE_01"] + +[[requirements]] +id = "BH-STORE-02" +layer = "integration" +status = "implemented" +tests = ["test/integration/store/store_integration_test.go:TestUnsafePathsAndCancelledOpen_BH_STORE_02"] + +[[requirements]] +id = "BH-STORE-03" +layer = "unit" +status = "implemented" +tests = ["test/harness/harness_test.go:TestBoundedSuiteLogAndFailedRunRecovery_BH_STORE_03"] + +[[requirements]] +id = "RACE-STORE-01" +layer = "integration" +status = "implemented" +tests = ["test/integration/store/store_integration_test.go:TestSingleInstanceLockAndConcurrentClose_RACE_STORE_01"] + +[[requirements]] +id = "REC-STORE-01" +layer = "integration" +status = "implemented" +tests = ["test/integration/store/store_integration_test.go:TestMigrationChecksumMismatchPreventsOpen_REC_STORE_01"] diff --git a/test/harness/harness.go b/test/harness/harness.go index 650d043..da19743 100644 --- a/test/harness/harness.go +++ b/test/harness/harness.go @@ -16,6 +16,7 @@ import ( "os/exec" "path/filepath" "regexp" + "strconv" "strings" "time" @@ -26,6 +27,7 @@ import ( const ( manifestVersion = 1 repositoryID = "rvbox" + maxSuiteLogSize = 1 << 20 ) var runIDPattern = regexp.MustCompile(`^[a-z0-9][a-z0-9-]{0,63}$`) @@ -124,8 +126,8 @@ func (h *harness) integration(ctx context.Context, args []string) error { if err := flags.Parse(args); err != nil { return err } - if *suite != "sample" { - return fmt.Errorf("suite %q is not implemented yet; available: sample", *suite) + if *suite != "sample" && *suite != "store" { + return fmt.Errorf("suite %q is not implemented yet; available: sample, store", *suite) } var current *manifest @@ -141,7 +143,7 @@ func (h *harness) integration(ctx context.Context, args []string) error { if current.Layer != "integration" || current.Suite != *suite { return errors.New("run layer/suite does not match resume request") } - if current.Phase != "ready" && current.Phase != "interrupted" && current.Phase != "stopped" && current.Phase != "running" { + if current.Phase != "ready" && current.Phase != "interrupted" && current.Phase != "stopped" && current.Phase != "running" && current.Phase != "failed" { return fmt.Errorf("run in phase %q is not resumable; recover or reuse it first", current.Phase) } } else { @@ -151,7 +153,7 @@ func (h *harness) integration(ctx context.Context, args []string) error { } } fmt.Fprintln(h.out, current.RunID) - if err := h.transition(current, "running", "sample-start", "sample integration run started"); err != nil { + if err := h.transition(current, "running", *suite+"-start", *suite+" integration run started"); err != nil { return err } select { @@ -160,10 +162,38 @@ func (h *harness) integration(ctx context.Context, args []string) error { return ctx.Err() default: } - if err := h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "invariant", Status: "passed", Detail: "manifest ownership and journal durability verified"}); err != nil { + if *suite == "sample" { + if err := h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "invariant", Status: "passed", Detail: "manifest ownership and journal durability verified"}); err != nil { + return err + } + } else if err := h.runStoreSuite(ctx, current); err != nil { + _ = h.transition(current, "failed", "store-failed", err.Error()) return err } - return h.transition(current, "completed", "sample-complete", "sample integration run completed") + return h.transition(current, "completed", *suite+"-complete", *suite+" integration run completed") +} + +func (h *harness) runStoreSuite(ctx context.Context, current *manifest) error { + if err := h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "store-real-sqlite", Status: "running", Detail: "running real SQLite/WAL and filesystem cases"}); err != nil { + return err + } + artifactDir := filepath.Join(h.runDir(current.RunID), "artifacts") + if err := os.MkdirAll(artifactDir, 0o700); err != nil { + return err + } + capture := &limitedCapture{limit: maxSuiteLogSize} + command := exec.CommandContext(ctx, "go", "test", "-count=1", "-tags=integration", "-shuffle="+strconv.FormatInt(current.Seed, 10), "-timeout=2m", "./test/integration/store") + command.Stdout = capture + command.Stderr = capture + err := command.Run() + logPath := filepath.Join(artifactDir, "suite.log") + if writeErr := atomicWrite(logPath, capture.Bytes(), 0o600); writeErr != nil { + return writeErr + } + if err != nil { + return fmt.Errorf("store suite failed (bounded log %s): %w", logPath, err) + } + return h.appendJournal(current.RunID, journalEntry{At: h.now(), Step: "store-real-sqlite", Status: "passed", Detail: "real SQLite/WAL and filesystem cases passed"}) } func (h *harness) environmentCommand(command string, args []string) error { @@ -195,7 +225,7 @@ func (h *harness) environmentCommand(command string, args []string) error { case "collect": return h.collect(current) case "recover": - if current.Phase != "interrupted" && current.Phase != "stopped" && current.Phase != "running" { + if current.Phase != "interrupted" && current.Phase != "stopped" && current.Phase != "running" && current.Phase != "failed" { return fmt.Errorf("run in phase %q does not need recovery", current.Phase) } return h.transition(current, "ready", "recover", "run recovered and ready to resume") @@ -215,6 +245,34 @@ func (h *harness) environmentCommand(command string, args []string) error { } } +type limitedCapture struct { + data []byte + limit int + truncated bool +} + +func (capture *limitedCapture) Write(data []byte) (int, error) { + written := len(data) + remaining := capture.limit - len(capture.data) + if remaining > 0 { + if len(data) > remaining { + data = data[:remaining] + } + capture.data = append(capture.data, data...) + } + if written > remaining { + capture.truncated = true + } + return written, nil +} + +func (capture *limitedCapture) Bytes() []byte { + if !capture.truncated { + return capture.data + } + return append(append([]byte(nil), capture.data...), []byte("\n[output truncated by RVBox test harness]\n")...) +} + func (h *harness) create(requestedID, layer, suite string) (*manifest, error) { if requestedID == "" { id, err := domain.NewUUIDv7() diff --git a/test/harness/harness_test.go b/test/harness/harness_test.go index f4bada0..00a728a 100644 --- a/test/harness/harness_test.go +++ b/test/harness/harness_test.go @@ -112,6 +112,34 @@ func TestIntegrationResumeValidation_HP_CFG_01(t *testing.T) { } } +func TestBoundedSuiteLogAndFailedRunRecovery_BH_STORE_03(t *testing.T) { + t.Parallel() + + capture := &limitedCapture{limit: 4} + if written, err := capture.Write([]byte("123456")); err != nil || written != 6 { + t.Fatalf("Write = (%d, %v)", written, err) + } + if got := string(capture.Bytes()); !strings.HasPrefix(got, "1234\n") || !strings.Contains(got, "truncated") { + t.Fatalf("bounded output = %q", got) + } + + h := &harness{root: t.TempDir(), now: time.Now, out: &bytes.Buffer{}} + current, err := h.create("failed-store", "integration", "store") + if err != nil { + t.Fatal(err) + } + if err := h.transition(current, "failed", "store-failed", "injected"); err != nil { + t.Fatal(err) + } + if err := h.environmentCommand("recover", []string{"--run-id", current.RunID}); err != nil { + t.Fatalf("recover failed run: %v", err) + } + loaded, err := h.load(current.RunID) + if err != nil || loaded.Phase != "ready" { + t.Fatalf("recovered run = (%+v, %v)", loaded, err) + } +} + func TestCoverageInventoryReferencesExistingTests_HP_CFG_01(t *testing.T) { t.Parallel() diff --git a/test/integration/store/store_integration_test.go b/test/integration/store/store_integration_test.go new file mode 100644 index 0000000..9fad671 --- /dev/null +++ b/test/integration/store/store_integration_test.go @@ -0,0 +1,236 @@ +//go:build integration + +package store_test + +import ( + "context" + "database/sql" + "errors" + "os" + "path/filepath" + "runtime" + "sort" + "strings" + "sync" + "testing" + "time" + + "github.com/rvbox/rvbox/internal/server/store" +) + +const busyTimeout = 2 * time.Second + +func TestRealSQLiteInitializationAndRestart_HP_STORE_01(t *testing.T) { + t.Parallel() + + dataDir := filepath.Join(t.TempDir(), "state") + opened := openStore(t, dataDir) + + for _, directory := range []string{dataDir, filepath.Join(dataDir, "segments"), filepath.Join(dataDir, "audit")} { + info, err := os.Stat(directory) + if err != nil { + t.Fatal(err) + } + if runtime.GOOS != "windows" && info.Mode().Perm() != 0o700 { + t.Fatalf("%s mode = %04o, want 0700", directory, info.Mode().Perm()) + } + } + for _, file := range []string{"rvbox.db", "server.lock"} { + info, err := os.Stat(filepath.Join(dataDir, file)) + if err != nil { + t.Fatal(err) + } + if runtime.GOOS != "windows" && info.Mode().Perm() != 0o600 { + t.Fatalf("%s mode = %04o, want 0600", file, info.Mode().Perm()) + } + } + + assertPragma(t, opened.DB(), "journal_mode", "wal") + assertPragma(t, opened.DB(), "foreign_keys", "1") + assertPragma(t, opened.DB(), "synchronous", "2") + assertPragma(t, opened.DB(), "busy_timeout", "2000") + + rows, err := opened.DB().Query(`SELECT name FROM sqlite_schema WHERE type = 'table' AND name NOT LIKE 'sqlite_%' ORDER BY name`) + if err != nil { + t.Fatal(err) + } + var names []string + for rows.Next() { + var name string + if err := rows.Scan(&name); err != nil { + t.Fatal(err) + } + names = append(names, name) + } + if err := rows.Close(); err != nil { + t.Fatal(err) + } + want := []string{"audit_events", "clients", "command_events", "command_payloads", "command_tombstones", "commands", "control_mutations", "output_segments", "output_truncations", "schema_migrations", "sessions", "stdin_writes", "storage_incidents", "takeover_authorizations"} + sort.Strings(want) + if strings.Join(names, ",") != strings.Join(want, ",") { + t.Fatalf("tables = %v, want %v", names, want) + } + var migrationCount int + if err := opened.DB().QueryRow(`SELECT count(*) FROM schema_migrations`).Scan(&migrationCount); err != nil || migrationCount != 1 { + t.Fatalf("migration count = %d, err = %v", migrationCount, err) + } + if err := opened.Close(); err != nil { + t.Fatal(err) + } + + reopened := openStore(t, dataDir) + if err := reopened.DB().QueryRow(`SELECT count(*) FROM schema_migrations`).Scan(&migrationCount); err != nil || migrationCount != 1 { + t.Fatalf("reopened migration count = %d, err = %v", migrationCount, err) + } + if err := reopened.Close(); err != nil { + t.Fatal(err) + } +} + +func TestSingleInstanceLockAndConcurrentClose_RACE_STORE_01(t *testing.T) { + t.Parallel() + + dataDir := filepath.Join(t.TempDir(), "state") + first := openStore(t, dataDir) + if _, err := store.Open(context.Background(), store.Options{DataDir: dataDir, BusyTimeout: busyTimeout}); !errors.Is(err, store.ErrAlreadyOpen) { + t.Fatalf("second Open error = %v, want ErrAlreadyOpen", err) + } + var wait sync.WaitGroup + for range 8 { + wait.Add(1) + go func() { + defer wait.Done() + if err := first.Close(); err != nil && !errors.Is(err, sql.ErrConnDone) { + t.Errorf("Close: %v", err) + } + }() + } + wait.Wait() + reopened := openStore(t, dataDir) + if err := reopened.Close(); err != nil { + t.Fatal(err) + } +} + +func TestSchemaRejectsInvalidIDsAndForeignKeys_BH_STORE_01(t *testing.T) { + t.Parallel() + + opened := openStore(t, filepath.Join(t.TempDir(), "state")) + defer opened.Close() + _, err := opened.DB().Exec(`INSERT INTO clients ( +client_id, platform, architecture, daemon_version, daemon_cwd, supported_shells, capabilities, client_instance_id +) VALUES ('client-a', 3, 'amd64', 'test', 'C:\\work', x'', x'', x'01')`) + if err == nil { + t.Fatal("invalid client instance UUID was accepted") + } + _, err = opened.DB().Exec(`INSERT INTO sessions ( +session_id, client_id, client_instance_id, generation, opened_at +) VALUES (?, 'missing-client', ?, 1, 1)`, bytesOf(16, 1), bytesOf(16, 2)) + if err == nil || !strings.Contains(strings.ToLower(err.Error()), "foreign key") { + t.Fatalf("foreign-key insert error = %v", err) + } +} + +func TestUnsafePathsAndCancelledOpen_BH_STORE_02(t *testing.T) { + t.Parallel() + + if _, err := store.Open(context.Background(), store.Options{DataDir: "relative", BusyTimeout: busyTimeout}); !errors.Is(err, store.ErrUnsafeDataDirectory) { + t.Fatalf("relative path error = %v", err) + } + if _, err := store.Open(context.Background(), store.Options{DataDir: string(filepath.Separator), BusyTimeout: busyTimeout}); !errors.Is(err, store.ErrUnsafeDataDirectory) { + t.Fatalf("root path error = %v", err) + } + permissive := filepath.Join(t.TempDir(), "permissive") + if err := os.Mkdir(permissive, 0o755); err != nil { + t.Fatal(err) + } + if _, err := store.Open(context.Background(), store.Options{DataDir: permissive, BusyTimeout: busyTimeout}); !errors.Is(err, store.ErrUnsafeDataDirectory) { + t.Fatalf("permissive path error = %v", err) + } + if runtime.GOOS != "windows" { + target := filepath.Join(t.TempDir(), "target") + if err := os.Mkdir(target, 0o700); err != nil { + t.Fatal(err) + } + link := filepath.Join(t.TempDir(), "linked") + if err := os.Symlink(target, link); err != nil { + t.Fatal(err) + } + if _, err := store.Open(context.Background(), store.Options{DataDir: link, BusyTimeout: busyTimeout}); !errors.Is(err, store.ErrUnsafeDataDirectory) { + t.Fatalf("symlink path error = %v", err) + } + } + + cancelled, cancel := context.WithCancel(context.Background()) + cancel() + dataDir := filepath.Join(t.TempDir(), "cancelled") + if _, err := store.Open(cancelled, store.Options{DataDir: dataDir, BusyTimeout: busyTimeout}); !errors.Is(err, context.Canceled) { + t.Fatalf("cancelled Open error = %v", err) + } + opened := openStore(t, dataDir) + if err := opened.Close(); err != nil { + t.Fatal(err) + } +} + +func TestMigrationChecksumMismatchPreventsOpen_REC_STORE_01(t *testing.T) { + t.Parallel() + + dataDir := filepath.Join(t.TempDir(), "state") + opened := openStore(t, dataDir) + if err := opened.Close(); err != nil { + t.Fatal(err) + } + raw, err := sql.Open("sqlite", filepath.Join(dataDir, "rvbox.db")) + if err != nil { + t.Fatal(err) + } + if _, err := raw.Exec(`UPDATE schema_migrations SET checksum = 'tampered' WHERE version = 1`); err != nil { + t.Fatal(err) + } + if err := raw.Close(); err != nil { + t.Fatal(err) + } + if _, err := store.Open(context.Background(), store.Options{DataDir: dataDir, BusyTimeout: busyTimeout}); err == nil || !strings.Contains(err.Error(), "checksum mismatch") { + t.Fatalf("Open error = %v, want checksum mismatch", err) + } + // A failed open must release the process lock for offline inspection/repair. + raw, err = sql.Open("sqlite", filepath.Join(dataDir, "rvbox.db")) + if err != nil { + t.Fatal(err) + } + if err := raw.Ping(); err != nil { + t.Fatal(err) + } + if err := raw.Close(); err != nil { + t.Fatal(err) + } +} + +func openStore(t *testing.T, dataDir string) *store.Store { + t.Helper() + opened, err := store.Open(context.Background(), store.Options{DataDir: dataDir, BusyTimeout: busyTimeout}) + if err != nil { + t.Fatalf("Open(%s): %v", dataDir, err) + } + return opened +} + +func assertPragma(t *testing.T, database *sql.DB, name, want string) { + t.Helper() + var got string + if err := database.QueryRow(`PRAGMA ` + name).Scan(&got); err != nil { + t.Fatalf("PRAGMA %s: %v", name, err) + } + if strings.ToLower(got) != want { + t.Fatalf("PRAGMA %s = %q, want %q", name, got, want) + } +} + +func bytesOf(length int, value byte) []byte { + result := make([]byte, length) + for index := range result { + result[index] = value + } + return result +}