Skip to content

artifactr.sql

The sql, postgres and sqlite extras. See SQL storage.

SQL storage: workspaces in PostgreSQL or SQLite, through SQLAlchemy 2's asyncio extension.

Install artifactr-ai[postgres] (asyncpg) or artifactr-ai[sqlite] (aiosqlite). Upgrade the database with migrate, then open workspaces over a SqlStorage::

engine = create_async_engine("postgresql+asyncpg://localhost/app")
await migrate(engine)
workspaces = Workspaces(SqlStorage(engine))

For SQLite, create the engine with create_sqlite_engine. create_schema creates the tables without migrations, for tests and prototypes.

Storage

SqlStorage

SqlStorage(
    engine: AsyncEngine,
    *,
    clock: Clock = _utc_now,
    poll_interval: timedelta = timedelta(seconds=0.5),
)

Storage in a SQL database, through SQLAlchemy's asyncio extension.

Create the tables with migrate first. The same code runs on PostgreSQL (with its default READ COMMITTED isolation) and on SQLite, whose engines must come from create_sqlite_engine.

Parameters:

Name Type Description Default
engine AsyncEngine

The database to use. The storage does not dispose of it.

required
clock Clock

Returns the current time, for envelope timestamps and lease expiry. Defaults to the system clock in UTC.

_utc_now
poll_interval timedelta

How often a subscription checks the log for envelopes committed by other processes. Commits made through this storage wake its subscriptions at once.

timedelta(seconds=0.5)

transaction async

transaction(scope: Scope) -> AsyncGenerator[_Transaction]

Begin a transaction; it holds the workspace row's lock until it ends.

artifact async

artifact(
    scope: Scope, artifact_id: ArtifactId
) -> Versioned[Artifact] | None

Return an artifact's current version, or None.

artifacts async

artifacts(
    scope: Scope,
    *,
    kind: str | None = None,
    include_archived: bool = False,
) -> list[Versioned[Artifact]]

Return current artifacts, optionally of one kind, oldest first.

revisions async

revisions(
    scope: Scope, artifact_id: ArtifactId
) -> list[Revision]

Return an artifact's revisions, oldest first.

thread async

thread(scope: Scope, thread_id: ThreadId) -> Thread | None

Return a thread, or None.

threads async

threads(scope: Scope) -> list[Thread]

Return every thread, oldest first.

proposal async

proposal(
    scope: Scope, proposal_id: ProposalId
) -> Proposal | None

Return a proposal, or None.

proposals async

proposals(
    scope: Scope,
    *,
    status: Literal["pending", "accepted", "rejected"]
    | None = None,
) -> list[Proposal]

Return proposals, optionally with one status, oldest first.

run async

run(scope: Scope, run_id: RunId) -> Run | None

Return a run, or None.

runs async

runs(
    scope: Scope,
    *,
    thread_id: ThreadId | None = None,
    status: RunStatus | None = None,
) -> list[Run]

Return runs, optionally of one thread and with one status, oldest first.

head_seq async

head_seq(scope: Scope) -> int

Return the log's latest seq, or 0 if it is empty.

read async

read(
    scope: Scope,
    *,
    after_seq: int = 0,
    limit: int | None = None,
) -> list[Envelope]

Return logged envelopes with seq greater than after_seq, in order.

subscribe async

subscribe(
    scope: Scope, *, after_seq: int = 0
) -> AsyncIterator[Envelope]

Yield stored envelopes after after_seq, then each new one as it commits.

New envelopes arrive at once when they were committed through this storage, and within poll_interval otherwise.

history async

history(
    scope: Scope, thread_id: ThreadId
) -> Sequence[HistoryChunk]

Return a thread's history chunks, in the order they were appended.

acquire_lease async

acquire_lease(
    scope: Scope, key: str, holder: str, ttl: timedelta
) -> bool

Take or renew an exclusive lease. Return False if another holder has it.

release_lease async

release_lease(scope: Scope, key: str, holder: str) -> None

Release a lease if holder has it.

Engines and schema

create_sqlite_engine

create_sqlite_engine(
    url: str | URL, **kwargs: Any
) -> AsyncEngine

Create an async SQLite engine for SqlStorage.

Every transaction on the engine, reads included, holds the database's write lock, so transactions run one at a time. A transaction waits for the lock for up to the timeout connect argument of Python's sqlite3 (5 seconds unless connect_args says otherwise). Use a database file (or a shared-cache memory database): a plain :memory: database is a single connection that concurrent transactions would share.

Parameters:

Name Type Description Default
url str | URL

A SQLite URL for an async driver, such as sqlite+aiosqlite:///app.db (install artifactr-ai[sqlite]).

required
**kwargs Any

Passed to sqlalchemy.ext.asyncio.create_async_engine.

{}

Returns:

Type Description
AsyncEngine

The engine. Dispose of it when you are done.

migrate async

migrate(engine: AsyncEngine) -> None

Upgrade the database to artifactr's latest schema, creating the tables if needed.

Run it when the application deploys or starts. It is safe to run on every start: a database that is already up to date is left alone.

create_schema async

create_schema(engine: AsyncEngine) -> None

Create artifactr's tables straight from the models, for tests and prototypes.

It skips Alembic, so the database records no schema version and cannot be upgraded with migrate later. Use migrate for databases that last.