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.
Closing the session rolls back whatever it did not commit, and returns its connection.
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.
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.
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.
read
async
¶
read(
scope: Scope,
*,
after_seq: int = 0,
before_seq: int | None = None,
threads: Collection[ThreadId] | None = None,
limit: int | None = None,
last: int | None = None,
) -> list[Envelope]
Return logged envelopes in the window after_seq < seq < before_seq, in order.
Of those delivered to threads, if given: the first limit or the last last.
The filter and both ends of the window are in the query, and the last are read
backwards from the end of the log.
subscribe
async
¶
subscribe(
scope: Scope, *, after_seq: int = 0
) -> AsyncGenerator[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
¶
Take or renew an exclusive lease. Return False if another holder has it.
release_lease
async
¶
Release a lease if holder has it.
cursor
async
¶
Return how far a named consumer of the log has got: the seq saved, or 0.
Engines and schema¶
create_sqlite_engine
¶
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, on the engine's one connection. A transaction waits for the
connection for up to the engine's pool_timeout (30 seconds), and for another process's
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).
Shared verbatim with reflexr's src/reflexr/sql/sqlite.py; change both.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
url
|
str | URL
|
A SQLite URL for an async driver, such as |
required |
**kwargs
|
Any
|
Passed to |
{}
|
Returns:
| Type | Description |
|---|---|
AsyncEngine
|
The engine. Dispose of it when you are done. |
migrate
async
¶
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.