Files

173 lines
5.6 KiB
Go

// 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) }