Documentation
¶
Overview ¶
Batches: several stream writes a producer hands the store at once, applied in one transaction with each entry isolated and the outcome recorded by id.
Compaction: rules that drop rows from the middle of a kind's streams, and merging several streams into one.
Dynamic kinds: a column for every key a kind's rows bring, typed by the JSON value class of the first value seen.
Package recordstore keeps append-only record streams: the rows a long capture produces, written as they arrive and read back by position instead of being carried inside whatever reported the capture.
A stream is named by a stream id and holds rows of one kind. Every row gets a per-stream seq, increasing from 1, which is the only position a reader resumes from — never a timestamp or a key, which neither order nor identify a row's position reliably. A kind may still declare a key, and then a stream holds each key once, which is what makes re-ingesting an overlapping source window idempotent. Rows leave a stream from its low end — the whole stream expires, or Trim removes the oldest appends — or, for a kind that replaces stored rows, when an append stores a new row under their key at a higher seq; that is the only way a stream's seqs skip.
Backends live in subpackages: kv (a clicky cache.Store: in-process memory or valkey/redis), sqlite (a file that is also the query index a profile reads), and ndjson (a file per stream). An Indexer mirrors any of them into a sqlite index incrementally, which is what makes a stream written by one process pageable through another's query engine.
A store has one writer process at a time. Backends serialize the writes made within one process; several processes share a store only through recordstore/owner, which elects one to write it while the others hand their writes over the spool.
Index ¶
- Variables
- func InferColumnType(value any) (query.ColumnType, bool)
- func RefuseSealed(meta Meta) error
- func ValidateAppend(stream, kind string) error
- func ValidateKind(kind string) error
- func ValidateStream(stream string) error
- func ValidateTTL(ttl time.Duration) error
- type AppendResult
- type Backend
- type BackendKind
- type Batch
- type BatchAppender
- type BatchEntry
- type BatchError
- type BatchErrorCode
- type BatchOp
- type BatchResult
- type ChangeSource
- type CompactRule
- type EntryResult
- type ImportRequest
- type Index
- type IndexDef
- type Indexer
- type KindOptions
- type KindSchema
- type MergeOptions
- type Merger
- type Meta
- type Notifier
- func (n *Notifier) Append(ctx context.Context, stream, kind string, rows []Row) (AppendResult, error)
- func (n *Notifier) AppendBatch(ctx context.Context, batch Batch) (BatchResult, error)
- func (n *Notifier) BatchOutcome(ctx context.Context, id string) (BatchResult, bool, error)
- func (n *Notifier) Close() error
- func (n *Notifier) Delete(ctx context.Context, stream string) error
- func (n *Notifier) Expire(ctx context.Context, stream string, ttl time.Duration) error
- func (n *Notifier) Meta(ctx context.Context, stream string) (Meta, error)
- func (n *Notifier) Reopen(ctx context.Context, stream, generation string) error
- func (n *Notifier) Scan(ctx context.Context, stream string, afterSeq int64, ...) error
- func (n *Notifier) Seal(ctx context.Context, stream string) error
- func (n *Notifier) Tail(ctx context.Context, stream string, afterSeq int64, ...) error
- func (n *Notifier) Trim(ctx context.Context, stream string, before time.Time) (Meta, error)
- func (n *Notifier) Unwrap() Backend
- func (n *Notifier) Wait(ctx context.Context, stream string, afterSeq int64, generation string) (Meta, error)
- type NotifierOptions
- type OnConflict
- type Producer
- type Reopener
- type Retention
- type Router
- func (r *Router) Append(ctx context.Context, stream, kind string, rows []Row) (AppendResult, error)
- func (r *Router) Backend(ctx context.Context, route string) (Backend, error)
- func (r *Router) Close() error
- func (r *Router) Delete(ctx context.Context, stream string) error
- func (r *Router) Expire(ctx context.Context, stream string, ttl time.Duration) error
- func (r *Router) Forget(route string, backend Backend) error
- func (r *Router) Meta(ctx context.Context, stream string) (Meta, error)
- func (r *Router) Reopen(ctx context.Context, stream, generation string) error
- func (r *Router) Resolve(ctx context.Context) (Backend, error)
- func (r *Router) Route(ctx context.Context) (string, error)
- func (r *Router) Scan(ctx context.Context, stream string, afterSeq int64, fn func(int64, Row) error) error
- func (r *Router) Seal(ctx context.Context, stream string) error
- func (r *Router) Trim(ctx context.Context, stream string, before time.Time) (Meta, error)
- type RouterOptions
- type Row
- type SchemaResolver
- type Schemas
- type Settings
- type StreamLocks
- type Submitter
- type Window
Constants ¶
This section is empty.
Variables ¶
var ( // ErrNotFound reports a stream that does not exist, or no longer does. ErrNotFound = errors.New("record stream not found") // ErrCapacity reports an append refused because the stream, or one row of // it, would exceed what the backend was configured to hold. The rows of a // refused append are not written — none of them. ErrCapacity = errors.New("record stream capacity exceeded") // ErrSealed reports an append refused because its stream was sealed: the // writer declared it complete, and a reader has already taken it as such. ErrSealed = errors.New("record stream sealed") // ErrUnsupported reports a kind the backend cannot store as declared, such // as one that replaces stored rows in a backend whose seqs cannot skip. ErrUnsupported = errors.New("record store operation unsupported") )
var ErrSchemaConflict = errors.New("record kind schema conflict")
ErrSchemaConflict reports rows whose kind is stored with another key, conflict policy or column storage than the rows' own schema declares: storing them would make the rows already written mean something else.
var ErrSchemaMismatch = errors.New("record value mismatches its column type")
ErrSchemaMismatch reports an append refused because a value's type differs from its column's: an inferred column keeps the type its first value gave it, and one append may not give a new column two. None of the refused append's rows are written.
Functions ¶
func InferColumnType ¶
func InferColumnType(value any) (query.ColumnType, bool)
InferColumnType is the column type a dynamic kind gives a key first seen holding value, by the class of value's JSON: string, number, boolean, or json for an object or an array. A null infers nothing. A time is its text: a datetime is never inferred, because a row that reached the store through JSON would hold the same instant as a string and infer another type.
func RefuseSealed ¶
RefuseSealed is the error an append to meta's stream fails with once it is sealed, or nil while it is not.
func ValidateAppend ¶
ValidateAppend checks an append's stream id and kind before a backend writes anything.
func ValidateKind ¶
ValidateKind rejects a kind that could not name a table, a directory and a profile.
func ValidateStream ¶
ValidateStream rejects a stream id a backend could not key by. The alphabet is narrow on purpose: an id becomes a Redis key segment, a file name and a query parameter, and every character outside it is a character one of those would have to escape.
func ValidateTTL ¶
ValidateTTL rejects an expiry that is not in the future.
Types ¶
type AppendResult ¶
type AppendResult struct {
Window Window `json:"window"`
Skipped int64 `json:"skipped"`
Replaced int64 `json:"replaced,omitempty"`
}
AppendResult is what one append stored: the window its rows were numbered into, how many rows it skipped because their key was already stored, and how many stored rows it replaced with the rows it appended.
func AppendTyped ¶
func AppendTyped[T any](ctx context.Context, backend Backend, stream, kind string, items []T) (AppendResult, error)
AppendTyped appends items as rows through their JSON encoding, which is the shape every backend stores and every reader sees.
type Backend ¶
type Backend interface {
// Append adds rows to stream, creating it under kind when it does not
// exist, and returns the window the rows were numbered into. Appending to
// an existing stream under a different kind is an error.
//
// A keyed kind (KindOptions.Key) skips every row whose key the stream
// already holds, atomically with the append, and numbers only the rows it
// keeps; a batch naming one key twice is refused whole. A kind that
// replaces stored rows (OnConflictReplace) instead removes each stored row
// the append names by key and numbers every appended row; a backend that
// cannot refuses such a kind with ErrUnsupported. A kind retaining
// rows (RetainRows) slides the stream's expiry to the backend ttl from now
// and trims the rows appended longer than that ago, in the same write.
Append(ctx context.Context, stream, kind string, rows []Row) (AppendResult, error)
// Meta describes stream, or returns ErrNotFound.
Meta(ctx context.Context, stream string) (Meta, error)
// Scan calls fn with every row after afterSeq, in seq order, stopping at
// the first error fn returns. A scan from below the stream's low seq starts
// at the low seq. An unknown stream is ErrNotFound.
Scan(ctx context.Context, stream string, afterSeq int64, fn func(seq int64, row Row) error) error
// Trim removes the rows appended before before — every append up to the
// last one made before it — and returns the stream's metadata after. The
// rows kept keep their seqs, the low seq moves to the first of them, and
// the keys of the rows removed may be appended again. An unknown stream is
// ErrNotFound.
Trim(ctx context.Context, stream string, before time.Time) (Meta, error)
// Expire removes stream ttl from now, rows appended later included. The
// ttl must be positive; an unknown stream is ErrNotFound.
Expire(ctx context.Context, stream string, ttl time.Duration) error
// Seal marks stream complete (Meta.Sealed): every later Append to it fails
// with ErrSealed, while Scan, Trim and Expire still apply. Sealing a sealed
// stream does nothing; an unknown stream is ErrNotFound. The seal ends with
// the stream, so an id reused after expiry starts unsealed.
Seal(ctx context.Context, stream string) error
// Delete removes one stream and its rows immediately. An unknown stream is ErrNotFound.
Delete(ctx context.Context, stream string) error
Close() error
}
Backend stores record streams.
type BackendKind ¶
type BackendKind string
BackendKind names a backend a store's streams can live in.
const ( BackendKV BackendKind = "kv" BackendSQLite BackendKind = "sqlite" BackendNDJSON BackendKind = "ndjson" )
type Batch ¶
type Batch struct {
ID string `json:"id"`
Producer Producer `json:"producer"`
Schemas []KindSchema `json:"-"`
Entries []BatchEntry `json:"entries"`
}
Batch is the entries a producer writes at once. Schemas, when set, are the kinds as the producer declares them, which the store reconciles additively with the kinds it holds; a kind without one is resolved as for any append.
type BatchAppender ¶
type BatchAppender interface {
AppendBatch(ctx context.Context, batch Batch) (BatchResult, error)
BatchOutcome(ctx context.Context, id string) (BatchResult, bool, error)
}
BatchAppender is a backend that applies batches. AppendBatch applies a batch id once: a repeat returns the outcome recorded the first time, so a producer may hand the same batch over again after a crash. An entry that fails rolls back alone and reports its error in its result; an error AppendBatch returns means nothing was applied. BatchOutcome reads a recorded outcome, reporting false for a batch id never applied.
type BatchEntry ¶
type BatchEntry struct {
Op BatchOp `json:"op"`
Stream string `json:"stream"`
// Kind and Rows are an append's. Seal seals the stream after the append,
// in the same entry.
Kind string `json:"kind,omitempty"`
Rows []Row `json:"-"`
Seal bool `json:"seal,omitempty"`
// Generation, when set, fences the entry to that incarnation of the
// stream: any other one, or none, fails the entry as not found. A reopen
// requires it.
Generation string `json:"generation,omitempty"`
// TTL is an expire's, and Before a trim's.
TTL time.Duration `json:"ttl,omitempty"`
Before time.Time `json:"before,omitzero"`
}
BatchEntry is one write of a batch, doing what the backend method of the same name does.
type BatchError ¶
type BatchError struct {
Code BatchErrorCode `json:"code"`
Message string `json:"message"`
}
BatchError is an entry's error as its batch recorded it.
func NewBatchError ¶
func NewBatchError(err error) *BatchError
NewBatchError codes err by the first sentinel it wraps, in the order the codes are declared; an error wrapping none means the entry itself was invalid.
func (*BatchError) Error ¶
func (e *BatchError) Error() string
func (*BatchError) Unwrap ¶
func (e *BatchError) Unwrap() error
type BatchErrorCode ¶
type BatchErrorCode string
BatchErrorCode classifies an entry's error so it survives being recorded and read back by another process.
const ( BatchErrorSealed BatchErrorCode = "sealed" BatchErrorNotFound BatchErrorCode = "not_found" BatchErrorCapacity BatchErrorCode = "capacity" BatchErrorConflict BatchErrorCode = "conflict" BatchErrorMismatch BatchErrorCode = "mismatch" BatchErrorUnsupported BatchErrorCode = "unsupported" BatchErrorInvalid BatchErrorCode = "invalid" )
type BatchResult ¶
type BatchResult struct {
ID string `json:"id"`
Entries []EntryResult `json:"entries"`
}
BatchResult is what a batch did, entry by entry, in entry order.
type ChangeSource ¶
ChangeSource is a backend another process may write: WatchChanges calls changed whenever that process may have committed, until ctx ends, without knowing which streams it changed.
type CompactRule ¶
CompactRule selects rows a store may drop when it compacts: those Where selects, a CEL expression over the row bound as `row`, and appended longer than OlderThan ago by the kind's time column. A rule sets either or both. A stream's last row is never dropped: it is what keeps its high seq.
type EntryResult ¶
type EntryResult struct {
Append *AppendResult `json:"append,omitempty"`
Meta *Meta `json:"meta,omitempty"`
Error *BatchError `json:"error,omitempty"`
}
EntryResult is what one entry did: an append's result, a trim's metadata, or the error that rolled the entry back.
type ImportRequest ¶
type ImportRequest struct {
Source Meta
First int64
Rows []Row
// Seqs is the source seq of each row, increasing from First. Nil numbers
// Rows contiguously from First. Only a kind that replaces stored rows may
// skip seqs, since only there did the source leave them behind.
Seqs []int64
}
ImportRequest is one source window copied into an index.
type Index ¶
type Index interface {
Backend
// Prepare reconciles the storage for source's kind and returns the indexed
// incarnation of its stream. A derived index removes an older generation
// under the same stream id and reports it absent.
Prepare(ctx context.Context, source Meta) (indexed Meta, found bool, err error)
// Import stores request.Rows under their source seqs. First must be the seq
// after the indexed stream's high seq, or for a stream the index does not
// hold yet the source's low seq; a gap or different generation fails.
Import(ctx context.Context, request ImportRequest) (Window, error)
// TrimBelow mirrors a source trim: it drops the indexed rows below lowSeq,
// and moves an indexed stream that had not reached lowSeq up to it empty.
TrimBelow(ctx context.Context, stream string, lowSeq int64) (Meta, error)
// Derived reports whether the index can discard rows and rebuild them from
// a separate source.
Derived() bool
// SetExpiry mirrors the source's absolute expiry. Nil keeps the indexed
// stream while its source exists.
SetExpiry(ctx context.Context, stream string, expiresAt *time.Time) error
// DeleteSeqs mirrors a source compaction: it drops the indexed rows at
// seqs and records the source's compaction count.
DeleteSeqs(ctx context.Context, stream string, seqs []int64, compactions int64) (Meta, error)
}
Index is a backend that takes rows under the seqs another backend gave them: what an Indexer mirrors a source into. The sqlite backend is one.
type IndexDef ¶
type IndexDef struct {
Columns []string
}
IndexDef is one index of a kind's rows: its declared columns, in order.
type Indexer ¶
type Indexer struct {
// contains filtered or unexported fields
}
Indexer keeps an Index caught up with a source backend, one stream at a time and incrementally by seq: each Ensure copies only what the source gained since the last one.
func NewIndexer ¶
NewIndexer mirrors source into index.
func (*Indexer) Ensure ¶
Ensure brings stream's index up to the source's high seq. A stream the source does not have is ErrNotFound. When the source is the index there is nothing to copy, and Ensure only confirms the stream exists.
type KindOptions ¶
type KindOptions struct {
// Key names a string column whose value identifies a row within a stream.
// A keyed stream holds each key once; OnConflict says what an append of a
// stored key does. Empty leaves the kind unkeyed, where every row appended
// is stored.
Key string
// Retention says how long a stream keeps its rows.
Retention Retention
// OnConflict says what an append of a key the stream holds does. Replace
// needs a Key.
OnConflict OnConflict
// TimeColumn names the datetime column a stream's rows are read newest
// first by. A backend that indexes rows indexes (stream, time desc, seq),
// the order every profile of the kind pages in.
TimeColumn string
// Indexes are the column lists a backend that indexes rows indexes, each
// after the stream id, since every read is of one stream or a few. A
// backend keeps every index a kind ever declared.
Indexes []IndexDef
// Dynamic makes a backend that stores rows by column add a column for
// every key a row brings that the kind does not declare, typed by
// InferColumnType; a value of another type later is ErrSchemaMismatch.
// Other backends store every key anyway.
Dynamic bool
// MaxDynamicColumns caps the columns a dynamic kind may infer; an append
// needing more is ErrCapacity. Zero is 256.
MaxDynamicColumns int
// Compact are the rules a store whose seqs can skip drops rows by when it
// compacts; a store whose seqs cannot refuses the kind with
// ErrUnsupported.
Compact []CompactRule
// Compressed names columns a backend storing rows by column keeps
// compressed. SQL cannot compare what is inside one, so each must be a
// text or structured column that is not the key, the time column or
// indexed, and that offers no filter — a string column's filter must be
// switched off explicitly. For large payloads read whole.
Compressed []string
}
KindOptions say how a kind's streams store its rows.
func (KindOptions) DynamicColumnLimit ¶
func (o KindOptions) DynamicColumnLimit() int
DynamicColumnLimit is how many columns a dynamic kind may infer beyond the ones it declares.
type KindSchema ¶
type KindSchema struct {
Kind string
Columns []query.ColumnDef
Options KindOptions
}
KindSchema is everything a kind declares: its columns and how its streams store them.
func ResolveKind ¶
func ResolveKind(resolver SchemaResolver, kind string) (KindSchema, error)
ResolveKind resolves kind through resolver and checks the schema it returns describes kind, so a backend never stores rows under a schema meant for another kind.
func (KindSchema) RefuseSkippedSeqs ¶
func (k KindSchema) RefuseSkippedSeqs() error
RefuseSkippedSeqs is the error a backend whose seqs cannot skip refuses schema's kind with, or nil for a kind whose seqs never skip.
func (KindSchema) RetentionTTL ¶
RetentionTTL is the ttl a stream of schema's kind keeps each row for, or zero for a kind that keeps its stream whole. A kind retaining rows needs a backend with a ttl: without one, nothing would ever leave the stream.
func (KindSchema) RowKeys ¶
func (k KindSchema) RowKeys(rows []Row) ([]string, error)
RowKeys reads the key of every row of an append to a stream of schema's kind, in row order, or nil for an unkeyed kind. A row without a non-empty string key, or two rows with the same key, refuse the whole append.
func (KindSchema) SkipsSeqs ¶
func (k KindSchema) SkipsSeqs() bool
SkipsSeqs reports whether a stream of the kind may skip seqs: one whose rows are replaced or compacted away leaves their seqs behind.
func (KindSchema) Validate ¶
func (k KindSchema) Validate() error
Validate refuses a schema no backend could store as declared.
type MergeOptions ¶
MergeOptions say what a merge copies. Where, a CEL expression over the row bound as `row`, keeps only the rows it selects; DeleteSources removes the sources in the same write as the merged stream is created.
type Merger ¶
type Merger interface {
Merge(ctx context.Context, target string, sources []string, options MergeOptions) (AppendResult, error)
}
Merger is a backend that merges streams: it creates target, which must not exist, holding the rows of sources — all of one kind — ordered by the kind's time column, then by the order of sources, then by seq. A key held by several sources keeps its first row, or, for a kind that replaces stored rows, its last. The rows are appended now, so a kind retaining rows keeps them from the merge rather than from their first append.
type Meta ¶
type Meta struct {
Stream string `json:"stream"`
Kind string `json:"kind"`
Generation string `json:"generation"`
// Total is how many rows the stream holds and HighSeq the seq of the last
// one; LowSeq is a lower bound of the first one's seq. Seqs only grow, but
// a kind that replaces stored rows leaves the replaced rows' seqs behind,
// so its stream may skip seqs and hold fewer than HighSeq-LowSeq+1 rows.
// All three are reported because an index mirroring the stream may lag. An
// empty stream has LowSeq HighSeq+1.
Total int64 `json:"total"`
LowSeq int64 `json:"lowSeq"`
HighSeq int64 `json:"highSeq"`
UpdatedAt time.Time `json:"updatedAt"`
// ExpiresAt is when the stream is removed, or nil when it is kept until
// something expires it.
ExpiresAt *time.Time `json:"expiresAt,omitempty"`
// Capped reports that an append was refused for capacity: the stream is
// complete up to HighSeq and missing whatever that append carried.
Capped bool `json:"capped,omitempty"`
// Compactions counts the compactions that dropped rows from the middle of
// the stream; an index mirroring it drops the same rows when the count
// changes.
Compactions int64 `json:"compactions,omitempty"`
// Sealed reports that the writer declared the stream complete for now.
// A reader that has read through HighSeq is done until an explicit Reopen;
// readers revisiting a stream must fetch Meta again.
Sealed bool `json:"sealed,omitempty"`
}
Meta describes a stream.
func NewStreamMeta ¶
NewStreamMeta starts one incarnation of stream. Generation distinguishes a stream id reused after expiry from the rows an index previously held for it.
type Notifier ¶
type Notifier struct {
// contains filtered or unexported fields
}
Notifier is a Backend that wakes the readers waiting on a stream as soon as an append to it commits, which is what lets a reader follow a stream rather than poll it. It delegates every call to the backend it wraps.
Only appends made through the Notifier wake a waiter at once: a stream has one writer process (see the package documentation), so that writer appending through the Notifier is the whole of the contract. Anything else is seen at the next recheck — but for a backend that is a ChangeSource, whose changes wake every waiter, which then re-reads its stream.
func NewNotifier ¶
func NewNotifier(backend Backend, options NotifierOptions) (*Notifier, error)
NewNotifier wraps backend.
func (*Notifier) Append ¶
func (n *Notifier) Append(ctx context.Context, stream, kind string, rows []Row) (AppendResult, error)
Append appends through the wrapped backend and, once the append committed, wakes every waiter on stream.
func (*Notifier) AppendBatch ¶
AppendBatch applies batch through the wrapped backend, which must be a BatchAppender, and wakes the waiters on every stream the batch names.
func (*Notifier) BatchOutcome ¶
func (*Notifier) Expire ¶
Expire sets the expiry through the wrapped backend and wakes the stream's waiters.
func (*Notifier) Seal ¶
Seal seals through the wrapped backend and wakes the stream's waiters, which finish once they have read what the stream holds.
func (*Notifier) Tail ¶
func (n *Notifier) Tail(ctx context.Context, stream string, afterSeq int64, fn func(seq int64, row Row) error) error
Tail calls fn with every row of stream after afterSeq, in seq order, and then with every row appended after that as it is appended. It returns nil once ctx ends or it has read through a sealed stream's high seq, fn's first error, and an ErrNotFound error when the stream does not exist or stops existing — expired, removed, or recreated as a new generation.
func (*Notifier) Trim ¶
Trim trims through the wrapped backend and wakes the stream's waiters, which re-read what is left.
func (*Notifier) Wait ¶
func (n *Notifier) Wait(ctx context.Context, stream string, afterSeq int64, generation string) (Meta, error)
Wait blocks until stream holds a row after afterSeq, or is sealed, and returns its metadata: a sealed stream whose HighSeq is at or below afterSeq has nothing more to wait for. generation is the incarnation of the stream the caller has read: a stream that is gone, or that was recreated under another generation since, is ErrNotFound, because the seqs the caller holds no longer name its rows. A cancelled ctx returns ctx.Err().
The metadata is read after subscribing to the next append, so an append that commits between the read and the wait still wakes it.
type NotifierOptions ¶
type NotifierOptions struct {
// RecheckInterval is how often a waiter re-reads its stream's metadata
// while no append wakes it. Appends made through the Notifier wake waiters
// at once; the recheck is what notices what no append announces — a stream
// expired, trimmed away or removed by the backend itself. Required.
RecheckInterval time.Duration
}
NotifierOptions configure NewNotifier.
type OnConflict ¶
type OnConflict int
OnConflict says what an append to a keyed stream does with a row whose key the stream already holds.
const ( // OnConflictSkip keeps the stored row and skips the appended one. OnConflictSkip OnConflict = iota // OnConflictReplace removes the stored row and stores the appended one // after the stream's high seq, in the same write, so a reader following // the stream sees the new row. The replaced row's seq is left behind, so // the stream's seqs skip it. OnConflictReplace )
func (OnConflict) String ¶
func (o OnConflict) String() string
type Producer ¶
type Producer struct {
Instance string `json:"instance"`
Seq int64 `json:"seq"`
PID int `json:"pid,omitempty"`
Host string `json:"host,omitempty"`
Build string `json:"build,omitempty"`
}
Producer identifies the process that made a batch: Instance is unique per process lifetime, and Seq counts that instance's batches from 1.
type Reopener ¶
Reopener permits a writer to resume a sealed stream after checking that its generation is still the one the caller recorded. Readers must take a fresh Meta after a resume; an earlier sealed snapshot remains final for its window.
type Retention ¶
type Retention int
Retention says how long a stream keeps its rows.
const ( // RetainStream keeps a stream whole until it expires, the backend ttl after // its first append unless Expire moves it. RetainStream Retention = iota // RetainRows keeps each row the backend ttl after its own append: every // append slides the stream's expiry to the ttl from now and trims the rows // appended longer than the ttl ago. It is for a stream that accumulates // for as long as something writes to it. RetainRows )
type Router ¶
type Router struct {
// contains filtered or unexported fields
}
Router is a Backend that resolves, per call, the backend of the route the call's context names, opening it on first use and keeping it — so a backend's per-stream append serialization holds across calls.
A stream id is not qualified by its route: routes are kept apart by each owning its backend, so a stream written on one route is ErrNotFound on every other.
func NewRouter ¶
func NewRouter(options RouterOptions) (*Router, error)
NewRouter routes calls by options.Route to backends options.Open opens.
func (*Router) Backend ¶
Backend is route's backend, opened on first use. It serves a caller that has already decided which routes a read may span — a read across routes the caller authorized — rather than the one route a context names.
func (*Router) Forget ¶
Forget drops and closes route's backend, so the next call on the route opens it afresh — for an owner whose store behind the route went away. It does nothing unless the route still holds backend (compared by identity), so an old owner releasing late cannot evict the backend its replacement opened. A call already holding the dropped backend finishes against it.
func (*Router) Resolve ¶
Resolve is the backend of the route ctx names, opened on first use: every call delegates to it, and a caller asking where a stream lives asks it.
type RouterOptions ¶
type RouterOptions struct {
// Route names the route a call's context belongs to — the tenant, the
// environment — and fails for a context that carries none.
Route func(ctx context.Context) (string, error)
// Open opens one route's backend. It is called once per route, and again
// only after Forget drops it. The backend it returns must belong to that
// route alone: never shared with another route, and never the index a
// registry reads, or one route's streams become readable through another.
Open func(ctx context.Context, route string) (Backend, error)
}
RouterOptions configure NewRouter.
type Row ¶
Row is one record: a JSON-shaped map keyed by column name.
func EncodeRow ¶
EncodeRow is item's JSON object as a row. Numbers are kept as json.Number so an int64 larger than a float64 can hold survives the trip.
func Unstored ¶
func Unstored(rows []Row, keys []string, stored func(key string) bool) (kept []Row, keptKeys []string, skipped int64)
Unstored keeps the rows of an append whose key stored does not report as already in the stream, with their keys, and counts the rows it skipped. With no keys (an unkeyed kind) every row is kept.
type SchemaResolver ¶
type SchemaResolver func(kind string) (KindSchema, error)
SchemaResolver resolves a kind to its schema, and refuses a kind nothing declared. Schemas.Kind is one.
type Schemas ¶
type Schemas struct {
// contains filtered or unexported fields
}
Schemas is the catalog of kinds and their columns. A backend resolves a kind through it — to store rows by column (sqlite), to find a kind's key and retention (every backend) — and a typed result registry fills it, so the two agree without either one owning the other.
func (*Schemas) Kind ¶
func (s *Schemas) Kind(kind string) (KindSchema, error)
Kind resolves kind, and refuses a kind nothing declared.
type Settings ¶
type Settings struct {
// Prefix is the property prefix the settings were read under; errors about
// a setting name the key under it.
Prefix string
// Backend is the backend asked for, or empty for the caller's default.
Backend BackendKind
Dir string
TTL time.Duration
NDJSONMaxBytes int64
NDJSONKeepStreams int
}
Settings say where a store's streams live and for how long. A store reads them from commons/properties (-P, env, a properties file) under a prefix of its own:
<prefix>.backend kv | sqlite | ndjson; unset lets Resolve pick <prefix>.dir where sqlite files and ndjson streams live <prefix>.ttl how long a stream is kept, e.g. 30d or 36h <prefix>.ndjson.maxBytes one ndjson stream's cap, e.g. 256MiB <prefix>.ndjson.keepStreams ndjson streams kept per kind
func ReadSettings ¶
ReadSettings reads prefix's properties over defaults. A key that is unset keeps its default; a key that is set but unusable is an error naming it, never a silent fallback to the default.
func (Settings) Resolve ¶
func (s Settings) Resolve(hasKV bool, fallback BackendKind) (BackendKind, error)
Resolve is the backend a store writes to. An explicit backend is used as asked; without one, kv wherever the caller has a kv store — the store every process can share — and fallback, a local file, where it has none.
hasKV is a fact about one caller, which may differ per tenant, so it is resolved per call. An explicit kv where there is no kv store is an error: writing somewhere else would lose the one property kv was asked for.
type StreamLocks ¶
type StreamLocks struct {
// contains filtered or unexported fields
}
StreamLocks serializes work on one stream while leaving other streams free, which is how a backend keeps the single-writer contract inside one process. The zero value is ready to use. A lock nobody holds or waits on is dropped, so the set does not grow with every stream ever written.
func (*StreamLocks) Lock ¶
func (l *StreamLocks) Lock(stream string) (unlock func())
Lock blocks until stream is free and returns the function that frees it.
func (*StreamLocks) TryLock ¶
func (l *StreamLocks) TryLock(stream string) (unlock func(), ok bool)
TryLock takes stream only when it is free, for work that should leave a stream in use alone rather than wait for it.
type Submitter ¶
type Submitter func(ctx context.Context, batch Batch) (BatchResult, error)
Submitter hands a batch to the process that writes a store and returns what the batch did there, for a process that only reads the store. It sets the batch's producer.
Source Files
¶
Directories
¶
| Path | Synopsis |
|---|---|
|
Package kv stores record streams in a clicky cache.Store, so one implementation runs in process (cache.NewMemory) and against valkey/redis (clicky/valkey.NewStore) — the backend a CLI writes to and a server reads.
|
Package kv stores record streams in a clicky cache.Store, so one implementation runs in process (cache.NewMemory) and against valkey/redis (clicky/valkey.NewStore) — the backend a CLI writes to and a server reads. |
|
Sealed streams kept gzip-compressed: a sealed stream's file is rewritten as <file>.gz, read linearly, and written plain again when the stream reopens.
|
Sealed streams kept gzip-compressed: a sealed stream's file is rewritten as <file>.gz, read linearly, and written plain again when the stream reopens. |
|
A reader's writes: published to the spool, the owner nudged through its socket, and the outcome awaited in the ledger the owner commits it to.
|
A reader's writes: published to the spool, the owner nudged through its socket, and the outcome awaited in the ledger the owner commits it to. |
|
Package probe coordinates cursor-based external sources with durable record streams.
|
Package probe coordinates cursor-based external sources with durable record streams. |
|
Read-time enrichment: a result type's hook adding computed columns to the rows its profile serves, checked to leave the rows it was given intact.
|
Read-time enrichment: a result type's hook adding computed columns to the rows its profile serves, checked to leave the rows it was given intact. |
|
recordresultstest
Package recordresultstest holds the result types, stores and tenants the recordresults specs and the HTTP specs that serve them share.
|
Package recordresultstest holds the result types, stores and tenants the recordresults specs and the HTTP specs that serve them share. |
|
Package recordstoretest is the conformance suite every recordstore.Backend runs, so the kv, sqlite and ndjson backends can never disagree about what appending, scanning, expiring and describing a stream mean.
|
Package recordstoretest is the conformance suite every recordstore.Backend runs, so the kv, sqlite and ndjson backends can never disagree about what appending, scanning, expiring and describing a stream mean. |
|
Codecs encode a batch entry's rows into its data file: ndjson and gzipped ndjson here, and others registered by the packages that implement them.
|
Codecs encode a batch entry's rows into its data file: ndjson and gzipped ndjson here, and others registered by the packages that implement them. |
|
parquet
Package parquet is the spool's parquet codec: importing it registers the "parquet" format, for producers exporting large batches in columnar form.
|
Package parquet is the spool's parquet codec: importing it registers the "parquet" format, for producers exporting large batches in columnar form. |
|
AppendBatch: a batch applied in one writer transaction, a savepoint per entry, and recorded in the ledger by id in that same transaction.
|
AppendBatch: a batch applied in one writer transaction, a savepoint per entry, and recorded in the ledger by id in that same transaction. |