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. - Exposeswait_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()
- get_session() AsyncIterator[AsyncSession][source]¶
Yield a read-only
AsyncSessionagainst 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:
Trueonce visible;Falseif timeout elapsed first.
- 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
storagemay 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:
- poll_interval_seconds¶
Polling interval for checking new commits (one GET when idle; 0 exclusive to 3600).
- Type:
- schema_version¶
Projection schema version; selects the snapshots/v{version}/ prefix used for bootstrap.
- Type:
- create_storage() BlobStorage[source]¶
Create the storage backend from the
storageconfig.
- classmethod from_dict(data: dict[str, Any]) ClientConfig[source]¶
Create config from a dictionary.
data["storage"]may be a StorageConfig or a dict with atypediscriminator (seeStorageConfig.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
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:
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
Replay new commits onto the new database (one transaction each)
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
- 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 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)
- 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