projectors

package
v0.7.1 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 5, 2026 License: Apache-2.0 Imports: 13 Imported by: 0

Documentation

Overview

The SQLite cold-start retry policy. It lives beside the projector database because the only retried operation is that database's first journal_mode pragma.

Projector database bootstrap. Migrations and the global path policy remain owned by the storage layer; this package owns the projector connection settings.

Package projectors applies session, message and part events to the SQLite tables that mirror flat storage. Each Apply call is one transaction.

The session SQLite schema. The project table that two foreign keys reference is owned by the caller.

Index

Constants

View Source
const (
	EventSessionCreated     = "session.created"
	EventSessionUpdated     = "session.updated"
	EventSessionDeleted     = "session.deleted"
	EventMessageUpdated     = "message.updated"
	EventMessageRemoved     = "message.removed"
	EventMessagePartRemoved = "message.part.removed"
	EventMessagePartUpdated = "message.part.updated"
)

Variables

View Source
var SchemaStatements = []string{
	`CREATE TABLE IF NOT EXISTS session (
		id text PRIMARY KEY,
		project_id text NOT NULL,
		workspace_id text,
		parent_id text,
		slug text NOT NULL,
		directory text NOT NULL,
		path text,
		title text NOT NULL,
		version text NOT NULL,
		share_url text,
		summary_additions integer,
		summary_deletions integer,
		summary_files integer,
		summary_diffs text,
		revert text,
		permission text,
		agent text,
		model text,
		time_created integer NOT NULL,
		time_updated integer NOT NULL,
		time_compacting integer,
		time_archived integer,
		CONSTRAINT fk_session_project_id_project_id_fk
			FOREIGN KEY (project_id) REFERENCES project(id) ON DELETE CASCADE
	)`,
	`CREATE INDEX IF NOT EXISTS session_project_idx ON session (project_id)`,
	`CREATE INDEX IF NOT EXISTS session_workspace_idx ON session (workspace_id)`,
	`CREATE INDEX IF NOT EXISTS session_parent_idx ON session (parent_id)`,
	`CREATE TABLE IF NOT EXISTS message (
		id text PRIMARY KEY,
		session_id text NOT NULL,
		time_created integer NOT NULL,
		time_updated integer NOT NULL,
		data text NOT NULL,
		CONSTRAINT fk_message_session_id_session_id_fk
			FOREIGN KEY (session_id) REFERENCES session(id) ON DELETE CASCADE
	)`,
	`CREATE INDEX IF NOT EXISTS message_session_time_created_id_idx
		ON message (session_id, time_created, id)`,
	`CREATE TABLE IF NOT EXISTS part (
		id text PRIMARY KEY,
		message_id text NOT NULL,
		session_id text NOT NULL,
		time_created integer NOT NULL,
		time_updated integer NOT NULL,
		data text NOT NULL,
		CONSTRAINT fk_part_message_id_message_id_fk
			FOREIGN KEY (message_id) REFERENCES message(id) ON DELETE CASCADE
	)`,
	`CREATE INDEX IF NOT EXISTS part_message_id_id_idx ON part (message_id, id)`,
	`CREATE INDEX IF NOT EXISTS part_session_idx ON part (session_id)`,
	`CREATE TABLE IF NOT EXISTS todo (
		session_id text NOT NULL,
		content text NOT NULL,
		status text NOT NULL,
		priority text NOT NULL,
		position integer NOT NULL,
		time_created integer NOT NULL,
		time_updated integer NOT NULL,
		CONSTRAINT todo_pk PRIMARY KEY (session_id, position),
		CONSTRAINT fk_todo_session_id_session_id_fk
			FOREIGN KEY (session_id) REFERENCES session(id) ON DELETE CASCADE
	)`,
	`CREATE INDEX IF NOT EXISTS todo_session_idx ON todo (session_id)`,
	`CREATE TABLE IF NOT EXISTS session_message (
		id text PRIMARY KEY,
		session_id text NOT NULL,
		type text NOT NULL,
		time_created integer NOT NULL,
		time_updated integer NOT NULL,
		data text NOT NULL,
		CONSTRAINT fk_session_message_session_id_session_id_fk
			FOREIGN KEY (session_id) REFERENCES session(id) ON DELETE CASCADE
	)`,
	`CREATE INDEX IF NOT EXISTS session_message_session_idx
		ON session_message (session_id)`,
	`CREATE INDEX IF NOT EXISTS session_message_session_type_idx
		ON session_message (session_id, type)`,
	`CREATE INDEX IF NOT EXISTS session_message_time_created_idx
		ON session_message (time_created)`,
	`CREATE TABLE IF NOT EXISTS permission (
		project_id text PRIMARY KEY,
		time_created integer NOT NULL,
		time_updated integer NOT NULL,
		data text NOT NULL,
		CONSTRAINT fk_permission_project_id_project_id_fk
			FOREIGN KEY (project_id) REFERENCES project(id) ON DELETE CASCADE
	)`,
}

SchemaStatements creates the session-owned tables and indexes.

Functions

func ApplySchema

func ApplySchema(ctx context.Context, db *sql.DB) error

ApplySchema installs the session-owned tables and indexes. The project table referenced by two foreign keys must be installed by the caller.

func Configure

func Configure(ctx context.Context, db SQLExecutor, retry BusyRetryOptions) error

Configure applies the startup PRAGMA sequence. Only journal_mode is retried: it is the first file-touching statement, and so the point where concurrent processes opening the same database collide at cold start.

func IsBusyError

func IsBusyError(err error, depth ...int) bool

IsBusyError reports whether err or one of at most depth wrapped causes is a SQLite BUSY-class error. The default depth is 5, meaning the outer error plus five causes are inspected.

func Open

func Open(ctx context.Context, path string, retry BusyRetryOptions) (*sql.DB, error)

Open creates the cgo-free SQLite connection and configures it for projector use. A single physical connection keeps connection-local PRAGMAs effective.

func WithBusyRetry

func WithBusyRetry[T any](fn func() (T, error), opts BusyRetryOptions) (T, error)

WithBusyRetry runs fn and retries only BUSY-class errors, sleeping a 100–400 ms (inclusive, by default) jitter between attempts.

Types

type BusyRetryError

type BusyRetryError struct {
	Message string
	Cause   error
}

BusyRetryError is returned after all BUSY-class attempts are exhausted. Unwrap preserves the last SQLite failure as the cause.

func (*BusyRetryError) Error

func (e *BusyRetryError) Error() string

func (*BusyRetryError) Unwrap

func (e *BusyRetryError) Unwrap() error

type BusyRetryOptions

type BusyRetryOptions struct {
	MaxAttempts *int
	BaseDelayMS *int
	MaxDelayMS  *int
	DBPath      string
	Random      func() float64
	Sleep       func(time.Duration)
	Log         io.Writer
}

BusyRetryOptions controls WithBusyRetry. Nil numeric fields select the defaults; pointers preserve the distinction between omitted and explicitly zero options.

type Event

type Event struct {
	ID   string          `json:"id"`
	Type string          `json:"type"`
	Data json.RawMessage `json:"data"`
}

Event is the serialized subset consumed by a projector.

type NotFoundError

type NotFoundError struct {
	Message string
}

NotFoundError is returned when session.updated names a session that does not exist. The useful detail is in Message.

func (*NotFoundError) Error

func (e *NotFoundError) Error() string

type PartialRow

type PartialRow map[string]any

PartialRow is a set of column values keyed by column name.

func ToPartialRow

func ToPartialRow(info json.RawMessage) (PartialRow, error)

ToPartialRow maps a JSON session patch onto the snake_case update columns. A nested field whose parent is null yields a null column.

type SQLExecutor

type SQLExecutor interface {
	ExecContext(context.Context, string, ...any) (sql.Result, error)
}

SQLExecutor is the minimal database surface needed by Configure.

type Store

type Store struct {
	// contains filtered or unexported fields
}

Store applies the ordered projector registry to a SQLite database.

func NewStore

func NewStore(db *sql.DB, options StoreOptions) *Store

NewStore binds the projector chain to an initialized database.

func (*Store) Apply

func (s *Store) Apply(ctx context.Context, event Event) error

Apply projects one event inside a SQLite transaction.

func (*Store) ApplyReconcileTx

func (s *Store) ApplyReconcileTx(ctx context.Context, tx *sql.Tx, event Event) error

ApplyReconcileTx upserts every authoritative field from a flat-storage event. Startup reconciliation must repair stale key and timestamp columns as well as JSON payloads, so it does not go through the live projectors.

func (*Store) ApplyTx

func (s *Store) ApplyTx(ctx context.Context, tx *sql.Tx, event Event) error

ApplyTx projects one event into an existing transaction.

type StoreOptions

type StoreOptions struct {
	Now  func() int64
	Warn func(Warning)
}

StoreOptions supplies the clock and the warning sink.

type Warning

type Warning struct {
	Message string
	Fields  PartialRow
}

Warning is emitted for the two deliberately ignored late-write cases.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL