Client

The lower-level read path: a polling projection with SQLAlchemy sessions, and the replay machinery that Projection and the jobs are built on. See Lower-level API.

Public API

class cairndb.client.CairnDBClient(config: ClientConfig, registry: HandlerRegistry, init_schema=None)[source]

Main entry point for CairnDB client applications.

Responsibilities: - Starts/stops the background polling updater. - Provides read-only SQLAlchemy sessions via get_session(). - Transparently reconnects when the projection file is swapped. - Exposes wait_for_sequence() for read-your-writes consistency.

Typical usage:

registry = HandlerRegistry()

@registry.handler("user.created")
async def handle_user_created(db, entry):
    await db.execute("INSERT INTO users ...")

config = ClientConfig(
    storage=FilesystemStorageConfig(path="./data"),
    db_path="./projection.db",
)

client = CairnDBClient(config, registry)
await client.start()

async with client.get_session() as session:
    result = await session.execute(select(User))

await client.stop()
async start() → None[source]

Start the background projection updater.

async stop() → None[source]

Stop the updater and dispose the SQLAlchemy engine.

get_session() → AsyncIterator[AsyncSession][source]

Yield a read-only AsyncSession against the current projection.

If the projection file has been atomically swapped since the last session (detected by comparing st_mtime), the engine is silently disposed and recreated before the session is opened.

Raises:

FileNotFoundError – When the projection database does not exist.

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

Block until sequence is visible in the projection.

Internally triggers polling updates while waiting. Useful immediately after writing an event when you need to query the result right away.

Parameters:
  • sequence – The sequence number string returned by the write API.

  • timeout – Maximum seconds to wait.

Returns:

True once visible; False if timeout elapsed first.

async trigger_update() → tuple[bool, str | None][source]

Force an immediate projection update outside the normal poll cycle.

Returns:

(applied, new_sequence) – whether new data was applied, and the resulting sequence string (or None).

property is_running: bool

Whether the background updater is currently active.

class cairndb.client.ClientConfig(storage: StorageConfig | None = None, db_path: str = './projection.db', poll_interval_seconds: float = 5.0, use_reflink: bool = True, schema_version: str = '1')[source]

Client configuration for CairnDB.

Configures: - Storage backend (where the commit log and snapshots live) - Projection database path - Updater behavior (polling interval) - Copy-on-write settings

storage may be left None when the storage backend is created separately and passed alongside the config (as BackgroundUpdater and Projector accept).

storage

Storage backend config (see cairndb.storage.config).

Type:

cairndb.storage.config.StorageConfig | None

db_path

Path to the SQLite projection database.

Type:

str

poll_interval_seconds

Polling interval for checking new commits (one GET when idle; 0 exclusive to 3600).

Type:

float

Use copy-on-write (reflink) when available.

Type:

bool

schema_version

Projection schema version; selects the snapshots/v{version}/ prefix used for bootstrap.

Type:

str

create_storage() → BlobStorage[source]

Create the storage backend from the storage config.

property db_path_obj: Path

Get db_path as Path object.

property new_db_path: str

Get new database path for atomic swap.

classmethod from_dict(data: dict[str, Any]) → ClientConfig[source]

Create config from a dictionary.

data["storage"] may be a StorageConfig or a dict with a type discriminator (see StorageConfig.from_dict()).

classmethod from_env() → ClientConfig[source]

Build configuration from CAIRNDB_* environment variables.

Storage variables are handled by StorageConfig.from_env(); client variables:

CAIRNDB_DB_PATH        - SQLite projection path
CAIRNDB_POLL_INTERVAL  - Poll interval seconds
CAIRNDB_SCHEMA_VERSION - Projection schema version
class cairndb.client.HandlerRegistry[source]

Registry for event handlers.

Applications register handlers for specific event types. Handlers are async functions that take a database connection and a SequencedEvent (event fields plus its global sequence number).

Example:

registry = HandlerRegistry()

@registry.handler("user.created")
async def handle_user_created(db, entry):
    await db.execute(
        "INSERT INTO users (id, name) VALUES (?, ?)",
        (entry.payload["id"], entry.payload["name"])
    )

Initialize the handler registry.

handler(event_type: str) → Callable[[Callable[[Any, SequencedEvent], Awaitable[None]]], Callable[[Any, SequencedEvent], Awaitable[None]]][source]

Decorator to register an event handler.

Parameters:

event_type – The event type to handle (e.g., “user.created”)

Returns:

Decorator function

Example:

@registry.handler("user.created")
async def handle_user_created(db, entry):
    ...
register(event_type: str, handler: Callable[[Any, SequencedEvent], Awaitable[None]]) → None[source]

Register an event handler.

Parameters:
  • event_type – The event type to handle

  • handler – Async function that processes the event

Raises:

ValueError – If handler is already registered for this event type

get_handler(event_type: str) → Callable[[Any, SequencedEvent], Awaitable[None]][source]

Get the handler for an event type.

Parameters:

event_type – The event type

Returns:

The registered handler function

Raises:

HandlerNotFoundError – If no handler is registered

has_handler(event_type: str) → bool[source]

Check if a handler is registered for an event type.

Parameters:

event_type – The event type to check

Returns:

True if handler exists, False otherwise

unregister(event_type: str) → None[source]

Unregister a handler.

Parameters:

event_type – The event type to unregister

Raises:

HandlerNotFoundError – If no handler is registered

list_handlers() → list[str][source]

List all registered event types.

Returns:

List of event types with handlers

clear() → None[source]

Clear all registered handlers.

cairndb.client.EventHandler

alias of Callable[[Any, SequencedEvent], Awaitable[None]]

Internals

These classes are stable enough to use, but most applications never need them directly.

class cairndb.client.projector.Projector(config: ClientConfig, storage: BlobStorage, registry: HandlerRegistry, init_schema: Callable[[str], Awaitable[None]] | None = None)[source]

Manages SQLite projection updates with atomic swaps.

Update workflow:

  1. Build the new database:

    • fresh start with a snapshot available: download it as the base (the snapshot carries schema and data)

    • otherwise: copy projection.db -> projection.db.new (COW if available), or create an empty base (metadata + init_schema) if nothing exists

  2. Replay new commits onto the new database (one transaction each)

  3. Atomically rename: projection.db.new -> projection.db

This ensures readers never see partial updates. Updates are serialized with an internal lock, so a background poll and a manual trigger can never double-apply commits.

Initialize the projector.

Parameters:
  • config – Client configuration

  • storage – Blob storage backend

  • registry – Event handler registry

  • init_schema – Optional callback creating application tables on a fresh, empty projection (not needed when snapshots exist or handlers create their own tables)

async initialize_projection() → None[source]

Initialize a new, empty projection database.

Creates the database file, the metadata table, and — when an init_schema callback was provided — the application tables.

Raises:

FileExistsError – If projection already exists

async apply_updates() → tuple[bool, SequenceNumber | None][source]

Bring the projection up to the log tail, atomically.

Steady state costs a single GET (404 = already caught up).

Returns:

Tuple of (updates_applied, current_sequence)

Raises:

ReplayError – If update fails

async rebuild_from_scratch() → SequenceNumber | None[source]

Rebuild the projection from the ledger, discarding local state.

Uses the latest snapshot as the base when available, then replays the log tail. Useful for corruption recovery, schema migrations, and testing.

Returns:

Final sequence number

Raises:

ReplayError – If rebuild fails

async get_current_sequence() → SequenceNumber | None[source]

Get the current projection sequence.

Returns:

Current sequence number, or None if projection is empty

async has_updates_available() → bool[source]

Check if new commits are available (single GET).

Returns:

True if new commits can be applied

class cairndb.client.updater.BackgroundUpdater(config: ClientConfig, storage: BlobStorage, registry: HandlerRegistry, init_schema=None)[source]

Background service keeping the projection eventually consistent.

Polls the log tail every poll_interval_seconds (a single GET when idle) and applies new commits via the projector’s atomic-swap update. Polling is the correctness mechanism; there is no push channel to miss.

Initialize the background updater.

Parameters:
  • config – Client configuration

  • storage – Blob storage backend

  • registry – Event handler registry

  • init_schema – Optional async callback creating application tables on a fresh projection (see Projector)

async start() → None[source]

Start the background update loop.

async stop() → None[source]

Stop the background update loop and wait for completion.

async trigger_update() → tuple[bool, str | None][source]

Apply updates immediately, outside the normal poll interval.

Useful for tests and read-your-writes flows.

Returns:

Tuple of (updated, new_sequence_str)

property is_running: bool

Check if the updater is currently running.

property update_count: int

Get the total number of updates performed.

async wait_for_sequence(target_sequence: str, timeout: float = 30.0) → bool[source]

Wait until the projection reaches a target sequence.

Read-your-writes flow: 1. seq = await committer.append(event) 2. await updater.wait_for_sequence(str(seq)) 3. Query the projection

Parameters:
  • target_sequence – Sequence string as returned by the committer

  • timeout – Maximum time to wait in seconds

Returns:

True if sequence reached, False if timeout

Raises:

ValueError – If invalid sequence format

class cairndb.client.replay.ReplayEngine(storage: BlobStorage, registry: HandlerRegistry)[source]

Engine for replaying commits to a SQLite database.

Responsibilities: - Stream commits from storage (sequential GETs; the log is dense) - Apply events in order via registered handlers - One SQLite transaction per commit (the natural atomicity boundary) - Track last applied sequence in the _cairndb_metadata table

Initialize the replay engine.

Parameters:
  • storage – Blob storage backend

  • registry – Event handler registry

async iter_commits(after: int, end_at: int | None = None) → AsyncIterator[Commit][source]

Yield commits with numbers > after, in order, until a gap (404).

Because the log is dense, a missing number means the tail is reached.

Parameters:
  • after – Yield commits strictly after this number

  • end_at – Stop after this commit number (inclusive), if given

async replay(db_path: str, after: int, end_at: int | None = None) → SequenceNumber | None[source]

Replay commits (after, end_at] onto a database.

Each commit is applied in its own transaction, so readers of a crashed replay never see a partially applied commit.

Parameters:
  • db_path – Path to SQLite database (opened in write mode)

  • after – Replay commits strictly after this number

  • end_at – Stop after this commit number (inclusive), if given

Returns:

Last applied sequence number, or None if nothing was applied

Raises:

ReplayError – If a handler or the database fails

async initialize_metadata_table(db_path: str) → None[source]

Initialize the metadata table for tracking replay state.

Creates a _cairndb_metadata table if it doesn’t exist.

Parameters:

db_path – Path to SQLite database

async get_last_applied_sequence(db_path: str) → SequenceNumber | None[source]

Get the last applied sequence number from the database.

Parameters:

db_path – Path to SQLite database

Returns:

Last applied sequence, or None if no commits have been applied

async update_last_applied_sequence(db_path: str, sequence: SequenceNumber) → None[source]

Update the last applied sequence in the database.

Parameters:
  • db_path – Path to SQLite database

  • sequence – Last applied sequence number

class cairndb.client.DiscoveryService(storage: BlobStorage)[source]

Helpers for clients to locate their starting point and poll the tail.

Initialize the discovery service.

Parameters:

storage – Blob storage backend

async find_latest_snapshot(schema: str) → int | None[source]

Find the latest snapshot for a projection schema version.

Returns:

The snapshot’s commit number, or None if no snapshot exists.

async has_new_commits(current: SequenceNumber | None) → bool[source]

Check whether the log extends past the current position.

This is the steady-state poll: a single GET on the next commit number. 404 means fully caught up.

Parameters:

current – Current projection position, or None if empty

Returns:

True if at least one new commit exists