sqlite

package
v0.1.45 Latest Latest
Warning

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

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

Documentation

Overview

AppendBatch: a batch applied in one writer transaction, a savepoint per entry, and recorded in the ledger by id in that same transaction.

Compaction in a sqlite file: dropping the rows a kind's rules select from the middle of its streams, mirroring that into an index, and merging streams into a new one.

Dynamic kinds in a sqlite file: a column added for every key a row brings, typed by its first value, and every column the catalog records adopted.

Kind tables for a batch: a kind the batch declares itself is reconciled additively with the file, and any other kind is resolved as for an append.

A kind table's indexes: the time index a kind's newest-first pages read, and the indexes it declares, recorded in its catalog entry and never dropped.

The batch ledger: every batch the file applied, by id, with its outcome, written in the batch's own transaction so a batch is applied exactly once.

A backend opened read-only: it reads a file another process writes, hands every mutation to that process as a batch, and can be promoted to write.

Replacing stored rows: an append to a kind that replaces them removes the rows its keys name before storing the appended ones after the high seq.

Package sqlite stores record streams in a SQLite file: a table per kind, keyed (stream_id, seq), with the kind's columns typed from its schema, and a record_streams table describing every stream.

The same file is the query index a `sql` profile reads through ReadDSN, which is what makes a stream pageable, filterable and exportable by the native profile engine. Used as a durable backend it is the index already; behind another backend an Indexer keeps it caught up.

Index

Constants

View Source
const CatalogVersion = catalogVersion

CatalogVersion is the catalog version this build reads and writes, for a process telling another which files it can open.

Variables

This section is empty.

Functions

func VersionedPath

func VersionedPath(configured string) string

VersionedPath is the file this build opens for the configured path <dir>/<file>: <dir>/v<CatalogVersion>/<file>, made absolute.

Types

type Backend

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

Backend is a recordstore.Backend over a SQLite file.

func Open

func Open(options Options) (*Backend, error)

Open opens (creating when absent) the file at options.Path.

func (*Backend) Append

func (b *Backend) Append(ctx context.Context, stream, kind string, rows []recordstore.Row) (recordstore.AppendResult, error)

Append numbers the rows it keeps after the stream's high seq.

func (*Backend) AppendBatch

func (b *Backend) AppendBatch(ctx context.Context, batch recordstore.Batch) (recordstore.BatchResult, error)

AppendBatch applies batch's entries in order, each one committed or rolled back alone, and records the outcome under batch.ID; a read-only backend submits it instead. A batch id the ledger already holds is answered with the outcome it recorded, applying nothing. Rows are stamped as appended now, when the batch is applied, so a kind retaining rows keeps them from then rather than from when they were made.

func (*Backend) BatchOutcome

func (b *Backend) BatchOutcome(ctx context.Context, id string) (recordstore.BatchResult, bool, error)

BatchOutcome reads batch id's recorded outcome.

func (*Backend) Close

func (b *Backend) Close() error

Close stops the sweeper, waiting out a sweep in progress, and closes the file.

func (*Backend) Compact

func (b *Backend) Compact(ctx context.Context) (int, error)

Compact drops, from every stream of every kind with compact rules, the rows a rule selects, evaluating at most compactBudget rows, and reports how many it dropped. A stream's last row always stays. A derived index and a read-only backend compact nothing: the one mirrors its source, the other leaves the file to its writer.

func (*Backend) Delete

func (b *Backend) Delete(ctx context.Context, stream string) error

func (*Backend) DeleteSeqs

func (b *Backend) DeleteSeqs(ctx context.Context, stream string, seqs []int64, compactions int64) (recordstore.Meta, error)

DeleteSeqs mirrors a source compaction into the index: it drops the indexed rows at seqs and records the source's compaction count.

func (*Backend) Derived

func (b *Backend) Derived() bool

Derived reports whether this file can discard indexed rows and rebuild them from a separate source.

func (*Backend) Expire

func (b *Backend) Expire(ctx context.Context, stream string, ttl time.Duration) error

Expire moves stream's expiry to ttl from now.

func (*Backend) Import

Import stores rows under source's generation and seqs. The requested first seq must follow the indexed high seq exactly, and the seqs must be contiguous, unless the kind replaces stored rows: its source leaves replaced rows' seqs behind, so its seqs only have to increase, and each imported row replaces the indexed row under its key. A stream the index does not hold yet starts at the first seq, since its source may have trimmed the seqs below.

func (*Backend) Lease

func (b *Backend) Lease() func()

Lease holds off every row and schema mutation until the returned function is called, so an external reader can issue multiple stable paging statements.

func (*Backend) Merge

func (b *Backend) Merge(ctx context.Context, target string, sources []string, options recordstore.MergeOptions) (recordstore.AppendResult, error)

Merge creates target from the rows of sources; see recordstore.Merger.

func (*Backend) Meta

func (b *Backend) Meta(ctx context.Context, stream string) (recordstore.Meta, error)

Meta describes stream.

func (*Backend) Path

func (b *Backend) Path() string

Path is the database file.

func (*Backend) Prepare

func (b *Backend) Prepare(ctx context.Context, source recordstore.Meta) (recordstore.Meta, bool, error)

Prepare reconciles source's table before comparing the indexed stream. A derived index removes an older incarnation so its first import starts at 1.

func (*Backend) ProducerSeq

func (b *Backend) ProducerSeq(ctx context.Context, instance string) (int64, error)

ProducerSeq is the highest seq of the batches the ledger holds from producer instance, or zero for a producer it holds none of.

func (*Backend) Promote

func (b *Backend) Promote(ctx context.Context) error

Promote makes a read-only backend write the file itself, for the process that has become its writer: it opens the writer, checks the catalog, and starts sweeping. Promoting a backend that writes already does nothing.

func (*Backend) ReadDSN

func (b *Backend) ReadDSN() string

ReadDSN is the read-only DSN a profile connection reads the file through.

func (*Backend) Reopen

func (b *Backend) Reopen(ctx context.Context, stream, generation string) error

func (*Backend) Scan

func (b *Backend) Scan(ctx context.Context, stream string, afterSeq int64, fn func(int64, recordstore.Row) error) error

Scan reads stream's rows after afterSeq in seq order.

func (*Backend) Seal

func (b *Backend) Seal(ctx context.Context, stream string) error

Seal marks stream complete on its record_streams row. An index mirrors a sealed source through it once it holds every row the source does.

func (*Backend) SetExpiry

func (b *Backend) SetExpiry(ctx context.Context, stream string, expiresAt *time.Time) error

SetExpiry makes an index expire at the source's exact deadline. A nil deadline keeps it for as long as its source exists.

func (*Backend) Sweep

func (b *Backend) Sweep(ctx context.Context) (int, error)

Sweep removes every stream whose expiry has passed, rows included, and reports how many it removed.

func (*Backend) SweepBatches

func (b *Backend) SweepBatches(ctx context.Context, before time.Time, keep func(id string) bool) (int, error)

SweepBatches removes the ledger rows of batches applied before before, but for those keep reports: a batch whose directory still exists must stay, or ingesting it again would apply it twice.

func (*Backend) Table

func (b *Backend) Table(kind string) (sqlitetable.Table, error)

Table is kind's table, created on first use. It is exported for the typed result registry, which reads a kind through the table's own columns.

func (*Backend) Trim

func (b *Backend) Trim(ctx context.Context, stream string, before time.Time) (recordstore.Meta, error)

Trim removes the rows appended before before.

func (*Backend) TrimBelow

func (b *Backend) TrimBelow(ctx context.Context, stream string, lowSeq int64) (recordstore.Meta, error)

TrimBelow drops the indexed rows below lowSeq, mirroring a trim of the source. An index that had not reached lowSeq is moved up to it as an empty stream, so the next import starts at lowSeq.

func (*Backend) WatchChanges

func (b *Backend) WatchChanges(ctx context.Context, changed func()) error

WatchChanges calls changed whenever a commit to the file lands, until ctx ends: the writer's, when this backend reads, and its own once it writes — which covers the batches it ingests on other processes' behalf, not only the appends a Notifier sees made through it.

type Options

type Options struct {
	// Path is the database file as configured, <dir>/<file>. The backend opens
	// <dir>/v<catalog version>/<file>, so a build reading another catalog
	// version keeps its own file; see versionedFile. It must be a file rather
	// than memory: a profile reads it through a connection of its own.
	Path string

	// Schema resolves a kind to its columns, key and retention. A kind it
	// refuses cannot be written.
	Schema recordstore.SchemaResolver

	// TTL is how long a stream lives from its first append unless Expire moves
	// it, or how long a row of a kind retaining rows lives from its own append.
	// Zero keeps a stream until something expires it, and refuses a kind
	// retaining rows.
	TTL time.Duration

	// Derived marks the file as an index rebuilt from a source. A kind whose
	// columns changed since its table was created is then dropped with every
	// stream it held, for an Indexer to refill; in a durable file the same
	// change is an error, because those rows exist nowhere else.
	Derived bool

	// Now is the clock streams are stamped and expired by. Nil is time.Now.
	Now func() time.Time

	// SweepInterval is how often the file removes the streams whose expiry
	// has passed, rows included, until it is closed. It is required: a stream
	// expires from a ttl or from the source an index mirrors, and a file that
	// never swept would keep every row it was ever given. A read-only backend
	// sweeps once promoted.
	SweepInterval time.Duration

	// ReadOnly opens the file another process writes, which must exist at
	// this build's catalog version, without writing it: no catalog creation,
	// no sweeper. Every mutation is handed to Submit as a batch, and a kind's
	// table is adopted from the catalog, or declared through Submit when the
	// file lacks it. A derived index cannot be read-only.
	ReadOnly bool

	// Submit hands a read-only backend's batches to the process writing the
	// file. Without it every mutation fails with sqlite.ErrReadOnly.
	Submit recordstore.Submitter
}

Options configure a sqlite backend.

Jump to

Keyboard shortcuts

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