// Package spool owns the durable, command-local state retained by a client // daemon between network sessions and process restarts. package spool import ( "context" "crypto/sha256" "database/sql" "errors" "fmt" "net/url" "path/filepath" "strconv" "strings" "sync" "time" "github.com/rvbox/rvbox/internal/domain" _ "modernc.org/sqlite" ) var ( ErrUnsafeDataDirectory = errors.New("unsafe client spool directory") ErrAlreadyOpen = errors.New("client spool directory is already locked") ErrIdentityCorrupt = errors.New("client instance identity is corrupt") ErrCommandConflict = errors.New("command UUID has different immutable content") ErrAlreadyExecuted = errors.New("command UUID was already executed") ErrUnknownCommand = errors.New("unknown client command") ErrInvalidEventAck = errors.New("event acknowledgement is beyond assigned sequence") ErrEventsPending = errors.New("command has events pending server acknowledgement") ) const DefaultTombstoneLimit uint64 = 1_000_000 type Options struct { DataDir string BusyTimeout time.Duration TombstoneLimit uint64 QuotaLimits QuotaLimits // MaxExecutionSpecBytes bounds the decompressed protobuf retained for an // accepted command. It is separate from MaxScriptBytes because a script // body is uploaded in a different table and lifecycle. MaxExecutionSpecBytes uint64 MaxScriptBytes uint64 } type Store struct { db *sql.DB unlock func() error dataDir string identity domain.UUID tombstoneLimit uint64 quotaLimits QuotaLimits maxScriptBytes uint64 maxExecutionSpecBytes uint64 mu sync.Mutex } // Open recovers an existing spool or creates an empty one. The caller must // surface ErrIdentityCorrupt as dirty health; it is deliberately never healed // by assigning a new client identity. 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, errors.New("busy timeout must be positive") } if options.TombstoneLimit == 0 { options.TombstoneLimit = DefaultTombstoneLimit } if options.QuotaLimits == (QuotaLimits{}) { options.QuotaLimits = DefaultQuotaLimits() } if err := options.QuotaLimits.Validate(); err != nil { return nil, err } if options.MaxScriptBytes == 0 { options.MaxScriptBytes = 10 << 20 if options.MaxScriptBytes > options.QuotaLimits.CommandTotalBytes { options.MaxScriptBytes = options.QuotaLimits.CommandTotalBytes } } if options.MaxExecutionSpecBytes == 0 { options.MaxExecutionSpecBytes = 768 << 10 if options.MaxExecutionSpecBytes > options.QuotaLimits.CommandTotalBytes { options.MaxExecutionSpecBytes = options.QuotaLimits.CommandTotalBytes } } if options.MaxExecutionSpecBytes > options.QuotaLimits.CommandTotalBytes { return nil, errors.New("execution specification maximum exceeds client command quota") } if options.MaxScriptBytes > options.QuotaLimits.CommandTotalBytes { return nil, errors.New("script maximum exceeds client command quota") } if err := ensurePrivateDirectory(options.DataDir); err != nil { return nil, err } unlock, err := acquireInstanceLock(filepath.Join(options.DataDir, "spool.lock")) if err != nil { return nil, err } identity, err := loadOrCreateIdentity(filepath.Join(options.DataDir, "client-instance-id"), domain.NewUUIDv7) if err != nil { _ = unlock() return nil, err } databasePath := filepath.Join(options.DataDir, "spool.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)+")") db, err := sql.Open("sqlite", sqliteFileURI(databasePath, query)) if err != nil { _ = unlock() return nil, err } db.SetMaxOpenConns(1) db.SetMaxIdleConns(1) store := &Store{db: db, unlock: unlock, dataDir: options.DataDir, identity: identity, tombstoneLimit: options.TombstoneLimit, quotaLimits: options.QuotaLimits, maxScriptBytes: options.MaxScriptBytes, maxExecutionSpecBytes: options.MaxExecutionSpecBytes} if err := db.PingContext(ctx); err != nil { _ = store.Close() return nil, fmt.Errorf("open client spool SQLite: %w", err) } if err := applyMigrations(ctx, db); err != nil { _ = store.Close() return nil, err } return store, nil } // sqliteFileURI creates a file: URI with no authority. SQLite requires a // Windows drive path to be /C:/... in the URI path; without the leading slash // net/url renders C: as an authority (file://C:/...), which SQLite rejects. func sqliteFileURI(path string, query url.Values) string { path = strings.ReplaceAll(path, `\`, "/") if len(path) >= 2 && path[1] == ':' { path = "/" + path } return (&url.URL{Scheme: "file", Path: path, RawQuery: query.Encode()}).String() } func (store *Store) ClientInstanceID() domain.UUID { return store.identity } 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 immutableDigest(payload []byte) [32]byte { return sha256.Sum256(payload) }