Engine¶
The CairnDB facade and the objects its methods return. Import them from
the top-level package (from cairndb import CairnDB, Lease, ...). See
Concepts for the semantics behind each layer.
CairnDB¶
- class cairndb.CairnDB(storage: BlobStorage)[source]¶
Serverless database engine on blob storage.
All state lives in the bucket; instances are cheap, stateless views of it and any number of them — across any number of processes — may operate on the same bucket concurrently.
- classmethod configure(config: dict[str, Any] | StorageConfig) CairnDB[source]¶
Build an engine from configuration.
Accepts either a StorageConfig, or a dict with a
"storage"key holding StorageConfig fields:CairnDB.configure({"storage": {"type": "s3", "bucket": "myapp"}})
- async claim(key: str, value: Any) ClaimResult[source]¶
Claim key put-if-absent; losers converge on the winner’s value.
- async lease(key: str, *, ttl: float, holder: str | None = None, steal_if_expired: bool = True, state_fn: Callable[[Any], Any] | None = None) Lease | None[source]¶
Acquire the lease on key, or None if it is actively held.
state_fn(current state → new state, pure, None on fresh creation) is applied atomically with the acquisition; an exception it raises aborts the acquisition with nothing written.
- lease_sync(key: str, *, ttl: float, holder: str | None = None, steal_if_expired: bool = True, state_fn: Callable[[Any], Any] | None = None) Lease | None[source]¶
Sync twin of
lease().
- async attach_lease(key: str, *, holder: str, ttl: float) Lease | None[source]¶
Re-attach to the lease on key that holder owns now, or None.
Reads the lease document and rebuilds the holder’s handle: same epoch, same deadline, nothing written. None when the lease is absent, released, held by someone else, or expired — an expired lease is resumed with
lease(), which takes a fresh epoch. The holder string is a credential: whoever knows it can write.
- attach_lease_sync(key: str, *, holder: str, ttl: float) Lease | None[source]¶
Sync twin of
attach_lease().
- async cooperative_write(key: str, state_fn: Callable[[Any], Any]) Any | None[source]¶
Write into the lease on key from outside the lease, without fencing the holder (e.g. a cancel flag); the holder sees the new state on its next
renew/update_state. Returns the state written, or None when no lease document exists.
- cooperative_write_sync(key: str, state_fn: Callable[[Any], Any]) Any | None[source]¶
Sync twin of
cooperative_write().
- doc(key: str, model: type | None = None) Document[source]¶
A typed, etag-guarded document with a read-modify-write loop.
- log(name: str | None = None) Log[source]¶
The named commit log name, or the root log when omitted.
Logs are cached per name; the same Log (and its committer) is returned for repeated calls.
- transact() Transaction[source]¶
Begin a multi-key transaction (use as
async with).
- async recover_transactions() int[source]¶
Re-apply committed transactions that never reached the object store (crash between commit point and apply). Idempotent.
- projection(name: str, *, version: str = DEFAULT_SCHEMA_VERSION, db_path: str | None = None, log: str | None = None, poll_interval: float = 5.0, init_schema: Callable[[str], Awaitable[None]] | None = None, registry: HandlerRegistry | None = None) Projection[source]¶
A declarative SQLite projection of one log (root by default).
Pass
registryto replay with an existing HandlerRegistry — the same onecairndb snapshot --handlersloads — instead of registering handlers withProjection.on().
- async wait_for_sequence(projection: Projection, sequence: SequenceNumber | str, timeout: float = 30.0) bool[source]¶
Convenience: block until projection has applied sequence.
Layer 0 — Objects¶
- class cairndb.Objects(storage: BlobStorage)[source]¶
Conditional key-addressed objects: get / put / delete / list / wait_for.
Semantics are exactly those of the underlying store (see docs/concepts/storage-model.md): etags are opaque and backend-native;
putwithif_matchis a compare-and-swap, withif_absenta put-if-absent, and returns None when the precondition failed.- get_sync(key: str) StoredObject | None[source]¶
Read key with its current etag, or None if absent.
- put_sync(key: str, data: bytes, *, if_match: str | None = None, if_absent: bool = False) str | None[source]¶
Write key; returns the new etag, or None if the precondition failed.
- async get(key: str) StoredObject | None[source]¶
Async twin of
get_sync().
- async put(key: str, data: bytes, *, if_match: str | None = None, if_absent: bool = False) str | None[source]¶
Async twin of
put_sync().
- async delete(key: str) None[source]¶
Async twin of
delete_sync().
- async wait_for(key: str, *, timeout: float = 60.0, poll_interval: float = 1.0, changed_from: str | None = None) StoredObject[source]¶
Poll until key exists (default) or no longer carries changed_from.
One GET per interval; polling is the correctness mechanism, there is no push channel to miss.
- Parameters:
key – Object key to watch
timeout – Maximum seconds to wait
poll_interval – Seconds between polls
changed_from – If given, wait until the object’s etag differs from this value (an object that has been deleted also counts as changed and returns on its next reappearance)
- Returns:
The observed object.
- Raises:
TimeoutError – If the condition is not met within timeout.
- cairndb.engine.objects.RESERVED_PREFIXES¶
Prefixes owned by the engine; application objects must live elsewhere.
Layer 1 — Coordination¶
- class cairndb.ClaimResult(won: bool, key: str, value: Any, etag: str)[source]¶
Outcome of a claim: exactly one caller wins, everyone converges.
valueis the claimant’s own value whenwonis True, and the winner’s value otherwise — a duplicate or recovered caller reads back the winner’s result and continues on the identical path.
- class cairndb.Lease(storage: BlobStorage, key: str, *, ttl: float, epoch: int, holder: str | None, deadline_at: Timestamp, state: Any, etag: str)[source]¶
An acquired lease on a key.
The lease document is
{epoch, holder, deadline_at, state}and is only ever replaced by CAS. Acquiring — fresh, after a release, or by stealing an expired lease — bumpsepoch, the monotonic fence token.Every write through the lease is guarded by the etag observed at the previous write; when the guard fails and the document’s epoch has advanced, the holder has been fenced and gets
LeaseLost.A guard failure on our own document (epoch and holder unchanged) is a cooperative write made by
cooperative_write_sync()— the fresh state is absorbed and the write retried, so cooperative writes are observed rather than clobbered.Obtain instances via
acquire_sync()/acquire()— or, for an ownership period that is already yours,attach_sync()/attach()— not the constructor.- renew_sync() None[source]¶
Extend the deadline by ttl. Raises LeaseLost if fenced.
State is preserved — including any cooperative write made out-of-band since our last write, which becomes visible on
self.state. A heartbeat loop ofrenew+self.statechecks is the holder side of a cooperative cancel protocol.
- write_sync(state: Any) None[source]¶
Replace the lease’s state payload under the ownership guard.
The replacement is unconditional: a cooperative write made out-of-band since our last write is discarded unread. Holders participating in a cooperative-write protocol should use
update_state_sync().
- update_state_sync(fn: Callable[[Any], Any]) Any[source]¶
Apply fn to the freshest state and write the result, guarded.
fn must be pure — a lost guard re-reads and applies it again, so a concurrent cooperative write ends up inside its input rather than clobbered. Returns the state that was written.
- release_sync(state: Any = None) None[source]¶
Release the lease, optionally recording a final state.
Without an explicit state the freshest state is preserved. The document remains (holder None) so the epoch stays monotonic across re-acquisitions.
- async renew() None[source]¶
Async twin of
renew_sync().
- async write(state: Any) None[source]¶
Async twin of
write_sync().
- async release(state: Any = None) None[source]¶
Async twin of
release_sync().
- class cairndb.Document(storage: BlobStorage, key: str, model: type | None = None)[source]¶
A mutable, etag-guarded document with a read-modify-write loop.
modelmay be any class exposingmodel_dump_json()andmodel_validate_json()(e.g. a pydantic BaseModel), in which case values are validated on read and serialized on write; without it, values are plain JSON-serializable objects.- get_sync() tuple[Any, str] | None[source]¶
Read the document. Returns (value, etag), or None if absent.
- update_sync(fn: Callable[[Any], Any], *, create: Any = None, max_attempts: int = _UPDATE_ATTEMPTS) Any[source]¶
Apply fn to the current value and CAS-write the result, retrying on conflict.
fn must be pure — it may run several times. When the document is absent, fn is applied to create; with no create, an absent document is an error.
- Returns:
The value that was written.
- Raises:
CairnDBError – If the document is absent and no create was given, or the write kept conflicting for max_attempts.
- async update(fn: Callable[[Any], Any], *, create: Any = None, max_attempts: int = _UPDATE_ATTEMPTS) Any[source]¶
Async twin of
update_sync().
- async delete() None[source]¶
Async twin of
delete_sync().
Layer 2 — Logs and transactions¶
- class cairndb.Log(storage: BlobStorage, name: str | None = None, *, committer_config: CommitterConfig | None = None, revalidate: Callable[[list[Event], list[Commit]], Awaitable[list[Event | None]]] | None = None)[source]¶
A named (or the root) commit log: append, read, tail.
The write path is the standard Committer (group commit, durable ack); it is created lazily on first append and closed by
close().- async append(event: Event) SequenceNumber[source]¶
Append one event; durable once returned.
- async append_many(events: list[Event]) list[SequenceNumber][source]¶
Append several events, preserving relative order.
- async read_commits(after: int = 0, end_at: int | None = None) AsyncIterator[Commit][source]¶
Yield commits > after in order until a gap (the tail).
- async read(after: int = 0, end_at: int | None = None) AsyncIterator[SequencedEvent][source]¶
Yield sequenced events from commits > after until the tail.
- async tail(after: int = 0, *, poll_interval: float = 1.0) AsyncIterator[SequencedEvent][source]¶
Follow the log forever: yield events as commits land.
Steady state costs one GET per interval (404 = idle). The caller stops by breaking out of the iteration.
- class cairndb.Transaction(manager: TransactionManager)[source]¶
A staged multi-key transaction. Obtain via
CairnDB.transact().Used as an async context manager: exiting without an exception commits; an exception (or never entering) discards the staged operations.
Layer 3 — Projections¶
- class cairndb.Projection(storage: BlobStorage, name: str, *, version: str = DEFAULT_SCHEMA_VERSION, db_path: str | None = None, poll_interval: float = 5.0, init_schema: Callable[[str], Awaitable[None]] | None = None, registry: HandlerRegistry | None = None)[source]¶
A derived, read-only SQLite view of one log.
- Parameters:
storage – The log’s storage view (already namespaced for a named log — pass
Log.storage)name – Projection name; also names the default db file
version – Projection schema version — selects the
snapshots/v{version}/prefix used for bootstrapdb_path – Local SQLite path (default
./{name}.v{version}.sqlite)poll_interval – Background poll interval in seconds
init_schema – Optional async callback creating application tables on a fresh, empty projection
registry – Handlers to replay with, e.g. the module-level registry the snapshot job also loads, so the app and the job cannot drift apart. Handlers added with
on()register into it. Default: a fresh, empty registry.
- property registry: HandlerRegistry¶
The handlers this projection replays with.
Read-only: replay is wired to this registry at construction, so a replacement could never take effect. Pass
registry=to the constructor instead.
- on(event_type: str) Callable[[Callable[[Any, SequencedEvent], Awaitable[None]]], Callable[[Any, SequencedEvent], Awaitable[None]]][source]¶
Decorator registering the handler for event_type.
Handlers are
async def handler(conn, entry)receiving an aiosqlite connection and a SequencedEvent; one SQLite transaction wraps each commit.
- async refresh() SequenceNumber | None[source]¶
Catch the projection up to the log tail now (one GET when idle).
- Returns:
The projection’s current sequence, or None while empty.
- async wait_for(sequence: SequenceNumber | str, timeout: float = 30.0) bool[source]¶
Block until the projection has applied sequence (read-your-writes).
- connect() Connection[source]¶
Open a read-only connection to the projection.
The projection is derived state: application code must never write to it, and mode=ro enforces that.
- async as_of(commit: int, dest_path: str | None = None) str[source]¶
Build a throwaway projection of the log through commit commit.
Bootstraps from the latest snapshot at or before the target when one exists (subject to GC retention: if the log head before the target was pruned and no usable snapshot remains, the build fails), then replays up to and including commit.
- Returns:
Path of the newly built SQLite file.