Storage

Backends and their configuration. Most applications only pass a config dict to CairnDB.configure. See Configuration.

Configuration

cairndb.storage.create_storage(storage_type: str, **kwargs: Any) → BlobStorage[source]

Factory — create a storage backend from its type name.

Parameters:
  • storage_type – One of “filesystem”, “s3”, “azure”, “gcs”

  • **kwargs – Fields of the matching StorageConfig subclass.

Returns:

A configured BlobStorage instance

Raises:
class cairndb.storage.StorageConfig[source]

Base class for storage backend configuration.

Instantiate a backend subclass directly, or dispatch on the type discriminator with from_dict() / from_env().

create_storage() → BlobStorage[source]

Create a storage backend instance from this config.

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

Build the config subclass selected by data["type"].

type defaults to "filesystem"; the remaining keys are the selected subclass’s fields.

Raises:
classmethod from_env() → StorageConfig[source]

Build configuration from CAIRNDB_* environment variables.

Supported variables:

CAIRNDB_STORAGE_TYPE - filesystem | s3 | azure | gcs CAIRNDB_STORAGE_PATH - filesystem root path CAIRNDB_STORAGE_PREFIX - Blob key prefix CAIRNDB_S3_BUCKET - S3 bucket CAIRNDB_GCS_BUCKET - GCS bucket (else CAIRNDB_S3_BUCKET) CAIRNDB_S3_REGION - AWS region CAIRNDB_S3_ENDPOINT_URL - Custom S3 endpoint CAIRNDB_AZURE_CONTAINER - Azure container name CAIRNDB_AZURE_CONNECTION_STRING - Azure connection string CAIRNDB_AZURE_ACCOUNT_URL - Azure account URL CAIRNDB_GCS_PROJECT - GCP project ID CAIRNDB_GCS_CREDENTIALS_PATH - GCP key file

class cairndb.storage.FilesystemStorageConfig(path: str)[source]

Bases: StorageConfig

Local filesystem storage (dev/testing).

class cairndb.storage.S3StorageConfig(bucket: str, prefix: str = '', region: str | None = None, endpoint_url: str | None = None)[source]

Bases: StorageConfig

Amazon S3, and S3-compatible stores with conditional-write support.

class cairndb.storage.GCSStorageConfig(bucket: str, prefix: str = '', project: str | None = None, credentials_path: str | None = None)[source]

Bases: StorageConfig

Google Cloud Storage.

class cairndb.storage.AzureStorageConfig(container: str, prefix: str = '', connection_string: str | None = None, account_url: str | None = None)[source]

Bases: StorageConfig

Azure Blob Storage.

Backend interface

class cairndb.storage.BlobStorage[source]

Abstract interface for blob storage backends.

All backends (filesystem, S3, GCS, Azure) implement this interface. In the log/snapshot API every write is either a put-if-absent of an immutable object or a delete (garbage collection); nothing is ever overwritten. The generic object API additionally supports mutable, etag-guarded documents (compare-and-swap replacement).

abstractmethod async put_commit(number: int, data: bytes) → bool[source]

Write commit object number if and only if it does not exist.

Returns:

True if this call created the object (the writer won the race); False if the object already exists (lost the race).

Raises:

StorageError – On any failure other than losing the race.

abstractmethod async get_commit(number: int) → bytes | None[source]

Read commit object number.

Returns:

The serialized commit, or None if it does not exist.

Raises:

StorageError – On any failure other than absence.

abstractmethod async list_commits(after: int = 0) → list[int][source]

List commit numbers greater than after, ascending.

Used for bootstrap/resync; steady-state polling should use get_commit(last + 1) instead.

abstractmethod async put_snapshot(schema: str, number: int, data: bytes) → bool[source]

Write the snapshot for schema at commit number if absent.

Returns:

True if created; False if it already exists (an identical snapshot — replay is deterministic — so False is success too).

abstractmethod async list_snapshots(schema: str) → list[int][source]

List snapshot commit numbers for schema, ascending.

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

Return the highest snapshot commit number for schema, or None.

abstractmethod async get_snapshot(schema: str, number: int) → bytes[source]

Read the snapshot for schema at commit number.

Raises:

StorageError – If the snapshot does not exist or the read fails.

abstractmethod async delete_commits_before(number: int) → int[source]

Delete commit objects with number < number. Returns count deleted.

abstractmethod async delete_snapshots_before(schema: str, number: int) → int[source]

Delete snapshots for schema with number < number. Returns count deleted.

abstractmethod get_object_sync(key: str) → StoredObject | None[source]

Read object key together with its current etag.

Returns:

A StoredObject, or None if the object does not exist.

Raises:

StorageError – On any failure other than absence.

abstractmethod put_object_sync(key: str, data: bytes, *, if_match: str | None = None, if_absent: bool = False) → str | None[source]

Write object key, optionally guarded by a precondition.

With no precondition the write is unconditional. With if_absent=True it succeeds only if the object does not exist (put-if-absent). With if_match=<etag> it succeeds only if the object still carries that etag (compare-and-swap); an object that has been deleted in the meantime also fails the precondition.

Returns:

The new etag on success; None if the precondition failed.

Raises:
  • ValueError – If both preconditions are given.

  • StorageError – On any failure other than a failed precondition.

append_object_sync(key: str, data: bytes) → bool[source]

Append data to object key, creating the object if absent.

Atomic per call with respect to concurrent appends: each call’s bytes land contiguously and none are lost, though ordering across concurrent appenders is unspecified. Intended for append-only streams (e.g. .jsonl observability logs); do not mix concurrent appends with unconditional put_object_sync rewrites of the same key — a full rewrite may discard a concurrent append.

This default implementation is a compare-and-swap read-modify-write loop (O(object size) per append). Backends with a native append primitive override it with a true O(len(data)) append.

Returns:

True when the append landed; False when the CAS fallback exhausted its retries under contention.

Raises:

StorageError – On any failure other than contention.

abstractmethod delete_object_sync(key: str) → None[source]

Delete object key. Deleting a missing object is a no-op.

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

List keys starting with prefix, ascending, relative to the store root.

async get_object(key: str) → StoredObject | None[source]

Async wrapper around get_object_sync().

async put_object(key: str, data: bytes, *, if_match: str | None = None, if_absent: bool = False) → str | None[source]

Async wrapper around put_object_sync().

async append_object(key: str, data: bytes) → bool[source]

Async wrapper around append_object_sync().

async delete_object(key: str) → None[source]

Async wrapper around delete_object_sync().

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

Async wrapper around list_objects_sync().

class cairndb.storage.StoredObject(data: bytes, etag: str)[source]

An object’s content together with the etag observed at read time.

The etag is an opaque, backend-specific token; it is only meaningful when passed back to the same store as the if_match precondition of BlobStorage.put_object().

Backends

class cairndb.storage.filesystem.FilesystemStorage(root_path: str | Path)[source]

Bases: BlobStorage

Local filesystem storage backend.

Put-if-absent is emulated with write-to-temp + os.link(): the link is atomic and fails with EEXIST if the target exists, and readers can never observe a partially written object. This makes the backend a faithful race simulator for tests, in addition to being the dev backend.

Storage layout mirrors the blob layout:

{root_path}/log/{number:012d}.msgpack {root_path}/snapshots/v{schema}/{number:012d}.sqlite

Initialize filesystem storage.

Parameters:

root_path – Root directory for all storage

async put_commit(number: int, data: bytes) → bool[source]

Write commit object number if and only if it does not exist.

Returns:

True if this call created the object (the writer won the race); False if the object already exists (lost the race).

Raises:

StorageError – On any failure other than losing the race.

async get_commit(number: int) → bytes | None[source]

Read commit object number.

Returns:

The serialized commit, or None if it does not exist.

Raises:

StorageError – On any failure other than absence.

async list_commits(after: int = 0) → list[int][source]

List commit numbers greater than after, ascending.

Used for bootstrap/resync; steady-state polling should use get_commit(last + 1) instead.

async put_snapshot(schema: str, number: int, data: bytes) → bool[source]

Write the snapshot for schema at commit number if absent.

Returns:

True if created; False if it already exists (an identical snapshot — replay is deterministic — so False is success too).

async list_snapshots(schema: str) → list[int][source]

List snapshot commit numbers for schema, ascending.

async get_snapshot(schema: str, number: int) → bytes[source]

Read the snapshot for schema at commit number.

Raises:

StorageError – If the snapshot does not exist or the read fails.

async delete_commits_before(number: int) → int[source]

Delete commit objects with number < number. Returns count deleted.

async delete_snapshots_before(schema: str, number: int) → int[source]

Delete snapshots for schema with number < number. Returns count deleted.

get_object_sync(key: str) → StoredObject | None[source]

Read object key together with its current etag.

Returns:

A StoredObject, or None if the object does not exist.

Raises:

StorageError – On any failure other than absence.

put_object_sync(key: str, data: bytes, *, if_match: str | None = None, if_absent: bool = False) → str | None[source]

Write object key, optionally guarded by a precondition.

With no precondition the write is unconditional. With if_absent=True it succeeds only if the object does not exist (put-if-absent). With if_match=<etag> it succeeds only if the object still carries that etag (compare-and-swap); an object that has been deleted in the meantime also fails the precondition.

Returns:

The new etag on success; None if the precondition failed.

Raises:
  • ValueError – If both preconditions are given.

  • StorageError – On any failure other than a failed precondition.

append_object_sync(key: str, data: bytes) → bool[source]

True append: O(len(data)) under the same lock CAS writers hold.

delete_object_sync(key: str) → None[source]

Delete object key. Deleting a missing object is a no-op.

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

List keys starting with prefix, ascending, relative to the store root.

class cairndb.storage.s3.S3Storage(bucket: str, prefix: str = '', region: str | None = None, endpoint_url: str | None = None)[source]

Bases: BlobStorage

Amazon S3 storage backend.

Put-if-absent uses IfNoneMatch=”*” (native S3 conditional writes; also supported by MinIO and most S3-compatible stores since late 2024).

Initialize S3 storage.

Parameters:
  • bucket – S3 bucket name

  • prefix – Key prefix for all objects

  • region – AWS region (e.g. “us-east-1”)

  • endpoint_url – Custom endpoint for S3-compatible stores (localstack, MinIO)

async put_commit(number: int, data: bytes) → bool[source]

Write commit object number if and only if it does not exist.

Returns:

True if this call created the object (the writer won the race); False if the object already exists (lost the race).

Raises:

StorageError – On any failure other than losing the race.

async get_commit(number: int) → bytes | None[source]

Read commit object number.

Returns:

The serialized commit, or None if it does not exist.

Raises:

StorageError – On any failure other than absence.

async list_commits(after: int = 0) → list[int][source]

List commit numbers greater than after, ascending.

Used for bootstrap/resync; steady-state polling should use get_commit(last + 1) instead.

async put_snapshot(schema: str, number: int, data: bytes) → bool[source]

Write the snapshot for schema at commit number if absent.

Returns:

True if created; False if it already exists (an identical snapshot — replay is deterministic — so False is success too).

async list_snapshots(schema: str) → list[int][source]

List snapshot commit numbers for schema, ascending.

async get_snapshot(schema: str, number: int) → bytes[source]

Read the snapshot for schema at commit number.

Raises:

StorageError – If the snapshot does not exist or the read fails.

async delete_commits_before(number: int) → int[source]

Delete commit objects with number < number. Returns count deleted.

async delete_snapshots_before(schema: str, number: int) → int[source]

Delete snapshots for schema with number < number. Returns count deleted.

get_object_sync(key: str) → StoredObject | None[source]

Read object key together with its current etag.

Returns:

A StoredObject, or None if the object does not exist.

Raises:

StorageError – On any failure other than absence.

put_object_sync(key: str, data: bytes, *, if_match: str | None = None, if_absent: bool = False) → str | None[source]

Write object key, optionally guarded by a precondition.

With no precondition the write is unconditional. With if_absent=True it succeeds only if the object does not exist (put-if-absent). With if_match=<etag> it succeeds only if the object still carries that etag (compare-and-swap); an object that has been deleted in the meantime also fails the precondition.

Returns:

The new etag on success; None if the precondition failed.

Raises:
  • ValueError – If both preconditions are given.

  • StorageError – On any failure other than a failed precondition.

delete_object_sync(key: str) → None[source]

Delete object key. Deleting a missing object is a no-op.

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

List keys starting with prefix, ascending, relative to the store root.

class cairndb.storage.gcs.GCSStorage(bucket: str, prefix: str = '', project: str | None = None, credentials_path: str | None = None)[source]

Bases: BlobStorage

Google Cloud Storage backend.

Put-if-absent uses if_generation_match=0, which succeeds only when the object does not exist (generation 0).

Initialize GCS storage.

Parameters:
  • bucket – GCS bucket name

  • prefix – Object prefix for all blobs

  • project – GCP project ID

  • credentials_path – Path to a service-account JSON key file

async put_commit(number: int, data: bytes) → bool[source]

Write commit object number if and only if it does not exist.

Returns:

True if this call created the object (the writer won the race); False if the object already exists (lost the race).

Raises:

StorageError – On any failure other than losing the race.

async get_commit(number: int) → bytes | None[source]

Read commit object number.

Returns:

The serialized commit, or None if it does not exist.

Raises:

StorageError – On any failure other than absence.

async list_commits(after: int = 0) → list[int][source]

List commit numbers greater than after, ascending.

Used for bootstrap/resync; steady-state polling should use get_commit(last + 1) instead.

async put_snapshot(schema: str, number: int, data: bytes) → bool[source]

Write the snapshot for schema at commit number if absent.

Returns:

True if created; False if it already exists (an identical snapshot — replay is deterministic — so False is success too).

async list_snapshots(schema: str) → list[int][source]

List snapshot commit numbers for schema, ascending.

async get_snapshot(schema: str, number: int) → bytes[source]

Read the snapshot for schema at commit number.

Raises:

StorageError – If the snapshot does not exist or the read fails.

async delete_commits_before(number: int) → int[source]

Delete commit objects with number < number. Returns count deleted.

async delete_snapshots_before(schema: str, number: int) → int[source]

Delete snapshots for schema with number < number. Returns count deleted.

get_object_sync(key: str) → StoredObject | None[source]

Read object key together with its current etag.

Returns:

A StoredObject, or None if the object does not exist.

Raises:

StorageError – On any failure other than absence.

put_object_sync(key: str, data: bytes, *, if_match: str | None = None, if_absent: bool = False) → str | None[source]

Write object key, optionally guarded by a precondition.

With no precondition the write is unconditional. With if_absent=True it succeeds only if the object does not exist (put-if-absent). With if_match=<etag> it succeeds only if the object still carries that etag (compare-and-swap); an object that has been deleted in the meantime also fails the precondition.

Returns:

The new etag on success; None if the precondition failed.

Raises:
  • ValueError – If both preconditions are given.

  • StorageError – On any failure other than a failed precondition.

delete_object_sync(key: str) → None[source]

Delete object key. Deleting a missing object is a no-op.

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

List keys starting with prefix, ascending, relative to the store root.

class cairndb.storage.azure.AzureBlobStorage(container: str, prefix: str = '', connection_string: str | None = None, account_url: str | None = None)[source]

Bases: BlobStorage

Azure Blob Storage backend.

Put-if-absent uses upload_blob(…, overwrite=False), which sets the If-None-Match: * condition and raises ResourceExistsError on conflict.

Provide either connection_string (simplest) or account_url. When only account_url is given, DefaultAzureCredential is used for authentication.

Initialize Azure Blob Storage.

Parameters:
  • container – Azure container name

  • prefix – Blob name prefix

  • connection_string – Full Azure Storage connection string

  • account_url – Storage account URL (uses DefaultAzureCredential)

async put_commit(number: int, data: bytes) → bool[source]

Write commit object number if and only if it does not exist.

Returns:

True if this call created the object (the writer won the race); False if the object already exists (lost the race).

Raises:

StorageError – On any failure other than losing the race.

async get_commit(number: int) → bytes | None[source]

Read commit object number.

Returns:

The serialized commit, or None if it does not exist.

Raises:

StorageError – On any failure other than absence.

async list_commits(after: int = 0) → list[int][source]

List commit numbers greater than after, ascending.

Used for bootstrap/resync; steady-state polling should use get_commit(last + 1) instead.

async put_snapshot(schema: str, number: int, data: bytes) → bool[source]

Write the snapshot for schema at commit number if absent.

Returns:

True if created; False if it already exists (an identical snapshot — replay is deterministic — so False is success too).

async list_snapshots(schema: str) → list[int][source]

List snapshot commit numbers for schema, ascending.

async get_snapshot(schema: str, number: int) → bytes[source]

Read the snapshot for schema at commit number.

Raises:

StorageError – If the snapshot does not exist or the read fails.

async delete_commits_before(number: int) → int[source]

Delete commit objects with number < number. Returns count deleted.

async delete_snapshots_before(schema: str, number: int) → int[source]

Delete snapshots for schema with number < number. Returns count deleted.

get_object_sync(key: str) → StoredObject | None[source]

Read object key together with its current etag.

Returns:

A StoredObject, or None if the object does not exist.

Raises:

StorageError – On any failure other than absence.

put_object_sync(key: str, data: bytes, *, if_match: str | None = None, if_absent: bool = False) → str | None[source]

Write object key, optionally guarded by a precondition.

With no precondition the write is unconditional. With if_absent=True it succeeds only if the object does not exist (put-if-absent). With if_match=<etag> it succeeds only if the object still carries that etag (compare-and-swap); an object that has been deleted in the meantime also fails the precondition.

Returns:

The new etag on success; None if the precondition failed.

Raises:
  • ValueError – If both preconditions are given.

  • StorageError – On any failure other than a failed precondition.

append_object_sync(key: str, data: bytes) → bool[source]

True append via Azure Append Blobs.

The object is created as an append blob on first use; append_block is server-side atomic, so concurrent appenders interleave whole blocks and never lose bytes. A key that already holds a block blob (written before appends existed, or by put_object_sync) cannot change type in place — those fall back to the base compare-and-swap rewrite.

delete_object_sync(key: str) → None[source]

Delete object key. Deleting a missing object is a no-op.

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

List keys starting with prefix, ascending, relative to the store root.

class cairndb.engine.NamespacedStorage(base: BlobStorage, namespace: str)[source]

Bases: BlobStorage

A BlobStorage view whose commit log and snapshots live under a prefix.

Commit and snapshot operations are rerouted through the base store’s generic conditional-object API under {namespace}/; the generic object API itself is delegated unprefixed, so coordination documents stay global.

async put_commit(number: int, data: bytes) → bool[source]

Write commit object number if and only if it does not exist.

Returns:

True if this call created the object (the writer won the race); False if the object already exists (lost the race).

Raises:

StorageError – On any failure other than losing the race.

async get_commit(number: int) → bytes | None[source]

Read commit object number.

Returns:

The serialized commit, or None if it does not exist.

Raises:

StorageError – On any failure other than absence.

async list_commits(after: int = 0) → list[int][source]

List commit numbers greater than after, ascending.

Used for bootstrap/resync; steady-state polling should use get_commit(last + 1) instead.

async delete_commits_before(number: int) → int[source]

Delete commit objects with number < number. Returns count deleted.

async put_snapshot(schema: str, number: int, data: bytes) → bool[source]

Write the snapshot for schema at commit number if absent.

Returns:

True if created; False if it already exists (an identical snapshot — replay is deterministic — so False is success too).

async list_snapshots(schema: str) → list[int][source]

List snapshot commit numbers for schema, ascending.

async get_snapshot(schema: str, number: int) → bytes[source]

Read the snapshot for schema at commit number.

Raises:

StorageError – If the snapshot does not exist or the read fails.

async delete_snapshots_before(schema: str, number: int) → int[source]

Delete snapshots for schema with number < number. Returns count deleted.

get_object_sync(key: str) → StoredObject | None[source]

Read object key together with its current etag.

Returns:

A StoredObject, or None if the object does not exist.

Raises:

StorageError – On any failure other than absence.

put_object_sync(key: str, data: bytes, *, if_match: str | None = None, if_absent: bool = False) → str | None[source]

Write object key, optionally guarded by a precondition.

With no precondition the write is unconditional. With if_absent=True it succeeds only if the object does not exist (put-if-absent). With if_match=<etag> it succeeds only if the object still carries that etag (compare-and-swap); an object that has been deleted in the meantime also fails the precondition.

Returns:

The new etag on success; None if the precondition failed.

Raises:
  • ValueError – If both preconditions are given.

  • StorageError – On any failure other than a failed precondition.

delete_object_sync(key: str) → None[source]

Delete object key. Deleting a missing object is a no-op.

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

List keys starting with prefix, ascending, relative to the store root.

cairndb.engine.log_storage(base: BlobStorage, name: str | None) → BlobStorage[source]

The storage view holding log name’s commits and snapshots.

The root log (name None) is the base storage itself; a named log is the base namespaced under logs/{name}/. Commit/snapshot consumers — the Committer, projections, the snapshot and GC jobs — operate on the log this view selects.

cairndb.storage.base.DEFAULT_SCHEMA_VERSION = '1'

Default projection schema version, shared by every component that names snapshots (Projection, ClientConfig, SnapshotBuilder, the CLI) so that defaults alone always agree on the snapshots/v{version}/ prefix.