feat: bootstrap durable server store

This commit is contained in:
2026-08-31 09:11:58 +00:00
parent 392129c253
commit 8a88a72e62
13 changed files with 822 additions and 12 deletions
+34
View File
@@ -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
}
+33
View File
@@ -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
}
+169
View File
@@ -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;
`
+149
View File
@@ -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
}