Documentation
¶
Overview ¶
Package persis defines the storage backend interface for Dagu's control plane.
All control-plane data (dag runs, users, secrets, queues, heartbeats, etc.) flows through Backend → Collection → Record. Domain model changes live inside Record.Data (an opaque blob); the physical schema never changes.
To add a new database backend, implement Backend and Collection only. All store adapters (dagrun, queue, proc, user, secret, …) work immediately.
Index ¶
Constants ¶
This section is empty.
Variables ¶
var ( // ErrNotFound is returned by Get and Claim when no matching record exists. ErrNotFound = errors.New("persis: record not found") // ErrConflict is returned by CompareAndSwap when the current Data does not // match the expected value. ErrConflict = errors.New("persis: compare-and-swap conflict") // ErrCorrupt is returned when a record exists but its stored representation // cannot be decoded by the backend. ErrCorrupt = errors.New("persis: corrupt record") )
Sentinel errors returned by Collection methods. Use errors.Is for matching; backends may wrap these with additional context.
Functions ¶
func Decode ¶
Decode unmarshals rec.Data into v as JSON. Store adapters use this after receiving a Record from Collection.Get or Collection.List.
Types ¶
type Backend ¶
type Backend interface {
// Collection returns the collection with the given name.
// Multiple calls with the same name return the same logical namespace.
Collection(name string) Collection
// Close releases all resources held by the backend (file handles,
// connection pools, etc.). No Collection method may be called after Close.
Close() error
}
Backend is the factory for storage [Collection]s and the sole interface a new database driver must implement.
Collections are created lazily; calling Collection("foo") never returns an error — failures surface on the first operation against the collection.
type Collection ¶
type Collection interface {
// Get returns the record identified by id.
// Returns [ErrNotFound] if no record with that id exists.
Get(ctx context.Context, id string) (*Record, error)
// Put creates or replaces a record.
Put(ctx context.Context, rec *Record) error
// Create atomically inserts rec. Returns [ErrConflict] when a record with
// rec.ID already exists. Backends must guarantee the check-and-insert is
// atomic with respect to concurrent Put, CompareAndSwap, and Create calls.
Create(ctx context.Context, rec *Record) error
// Delete removes the record with the given id.
// Returns nil if the record does not exist.
Delete(ctx context.Context, id string) error
// CompareAndDelete atomically removes expected.ID only when the current
// record still matches expected. Returns [ErrConflict] when it does not.
CompareAndDelete(ctx context.Context, expected *Record) error
// List returns a page of records matching q, ordered by CreatedAt ascending.
List(ctx context.Context, q ListQuery) (*Page, error)
// CompareAndSwap atomically replaces record id only when its current Data
// bytes equal expected. Returns [ErrConflict] when they do not match.
// Used for optimistic concurrency on DAGRunStatus updates.
CompareAndSwap(ctx context.Context, id string, expected, next []byte) error
}
Collection is a named, isolated namespace of [Record]s. All methods must be safe for concurrent use.
Collection names map to distinct physical namespaces — directories in [file.Backend], rows in a single SQL table keyed by (collection, id), or key prefixes in etcd / Cassandra.
type ListQuery ¶
type ListQuery struct {
// Prefix filters records whose ID starts with this string.
// An empty Prefix returns all records in the collection.
Prefix string
// Since and Until bound results by Record.CreatedAt.
Since *time.Time
Until *time.Time
// Cursor resumes a previous [Page] iteration.
// Pass [Page.NextCursor] from the prior call; empty starts from the beginning.
Cursor string
// Limit caps the number of records returned. 0 = backend default.
Limit int
}
ListQuery controls what Collection.List returns.
type Page ¶
Page is the result of a Collection.List call.
type Record ¶
Record is the universal storage primitive for all control-plane data.
ID uses "/" as a hierarchy separator so that a ListQuery.Prefix of "mydag/" returns all records whose IDs start with that prefix — enabling efficient tree traversal without backend-specific query syntax.
Example ID formats:
dag_runs → "mydag/run-abc123/attempt-0" secrets → "default/db-password" sessions → "user-1/sess-xyz" queue_items → "high_9f3a..." (priority + uuid, sort order = dequeue order)
Directories
¶
| Path | Synopsis |
|---|---|
|
Package file implements persis.Backend on the local filesystem.
|
Package file implements persis.Backend on the local filesystem. |
|
audit
Package audit provides a file-based implementation of the audit Store interface.
|
Package audit provides a file-based implementation of the audit Store interface. |
|
eventstore
Package eventstore provides a file-based implementation of the event store.
|
Package eventstore provides a file-based implementation of the event store. |
|
tokensecret
Package tokensecret provides a file-based implementation of auth.TokenSecretProvider.
|
Package tokensecret provides a file-based implementation of auth.TokenSecretProvider. |
|
Package store consolidates small persistence stores that each wrap a persis.Collection.
|
Package store consolidates small persistence stores that each wrap a persis.Collection. |
|
Package testutil provides test helpers for the persistence layer.
|
Package testutil provides test helpers for the persistence layer. |