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.

claim_sync(key: str, value: Any) → ClaimResult[source]

Sync twin of claim().

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 registry to replay with an existing HandlerRegistry — the same one cairndb snapshot --handlers loads — instead of registering handlers with Projection.on().

async wait_for_sequence(projection: Projection, sequence: SequenceNumber | str, timeout: float = 30.0) → bool[source]

Convenience: block until projection has applied sequence.

async close() → None[source]

Stop projection updaters and drain/close every log committer.

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; put with if_match is a compare-and-swap, with if_absent a 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.

delete_sync(key: str) → None[source]

Delete key; deleting a missing object is a no-op.

list_sync(prefix: str = '') → list[str][source]

List keys under prefix, ascending.

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 list(prefix: str = '') → list[str][source]

Async twin of list_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.

value is the claimant’s own value when won is 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 — bumps epoch, 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 of renew + self.state checks 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 update_state(fn: Callable[[Any], Any]) → Any[source]

Async twin of update_state_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.

model may be any class exposing model_dump_json() and model_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.

delete_sync() → None[source]

Delete the document; deleting a missing document is a no-op.

async get() → tuple[Any, str] | None[source]

Async twin of get_sync().

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().

property committer: Committer

The log’s committer, created on first use.

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.

async current_tail() → int[source]

Highest existing commit number (0 = empty log). One LIST.

async close() → None[source]

Drain and close the committer, if one was created.

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.

async get(key: str) → bytes | None[source]

Read key, recording it in the read set.

Returns this transaction’s own staged write when there is one (read-your-writes); the underlying read is still recorded so the transaction conflicts with concurrent writers of key.

put(key: str, data: bytes) → None[source]

Stage a write of key.

delete(key: str) → None[source]

Stage a deletion of key.

note(event_type: str, payload: dict[str, Any]) → None[source]

Attach an audit event to the transaction record.

Notes are embedded in the record itself; projections of the _tx log can fold them. They are not appended to any other log.

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 bootstrap

  • db_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 start() → None[source]

Start the background polling updater.

async stop() → None[source]

Stop the background updater.

async wait_for(sequence: SequenceNumber | str, timeout: float = 30.0) → bool[source]

Block until the projection has applied sequence (read-your-writes).

property path: str

Path of the local SQLite projection file.

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.