Skip to content

reflexr.sql

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

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

Install reflexr[postgres] (asyncpg) or reflexr[sqlite] (aiosqlite). Upgrade the database with migrate, then build workspaces over a SqlStorage::

engine = create_async_engine("postgresql+asyncpg://localhost/app")
await migrate(engine)
workspaces = Workspaces(SqlStorage(engine), events=[...], rules=[...])

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(
    workspace: WorkspaceRef,
) -> 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.

head_seq async

head_seq(workspace: WorkspaceRef) -> int

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

read async

read(
    workspace: WorkspaceRef,
    *,
    after_seq: int = 0,
    before_seq: int | None = None,
    types: Collection[str] | 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 of types, if given: the first limit or the last last. The filter and both ends of the window are in the query, over ix_reflexr_events_type when it names types, and the last are read from the end of the log.

subscribe async

subscribe(
    workspace: WorkspaceRef, *, 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.

run async

run(workspace: WorkspaceRef, run_id: RunId) -> Run | None

Return a run, or None.

runs async

runs(
    workspace: WorkspaceRef,
    *,
    rule: RuleName | None = None,
    status: RunStatus | None = None,
    scope_key: ScopeKey | None = None,
    limit: int | None = None,
) -> list[Run]

Return runs, newest first, optionally filtered.

dead_letters async

dead_letters(
    workspace: WorkspaceRef, *, rule: RuleName | None = None
) -> list[EvaluationError]

Return the envelopes rules could not evaluate, oldest first.

progress async

progress(
    workspace: WorkspaceRef,
) -> dict[RuleName, RuleProgress]

Return every rule's progress in a workspace.

schedules async

schedules(workspace: WorkspaceRef) -> dict[str, datetime]

Return the time of each schedule's last tick in a workspace.

stored_rules async

stored_rules(workspace: WorkspaceRef) -> list[StoredRule]

Return a workspace's active stored rules, by name, in code-point order.

workspaces async

workspaces() -> list[WorkspaceRef]

Return every workspace with a log.

due_runs async

due_runs(
    *, now: datetime, limit: int, policy: RunPolicy
) -> list[tuple[WorkspaceRef, Run]]

Return the runs to attempt now that can start, across workspaces, oldest first.

A running run is due when its lease is missing or has expired by now. A pending or retrying run of an ordered rule is due only if no earlier run of its scope holds it, as a NOT EXISTS over ix_reflexr_runs_scope_head finds. Runs of the disabled rules are left out.

acquire_lease async

acquire_lease(
    workspace: WorkspaceRef,
    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(
    workspace: WorkspaceRef, key: str, holder: str
) -> None

Release a lease if holder has it.

cursor async

cursor(workspace: WorkspaceRef, name: str) -> int

Return how far a named consumer of the log has got: the seq saved, or 0.

save_cursor async

save_cursor(
    workspace: WorkspaceRef, name: str, seq: int
) -> None

Save how far a named consumer of the log has got; a cursor only moves forward.

The cursor's row is locked while it is compared, so saves from several processes leave the furthest.

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, 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 artifactr's src/artifactr/sql/sqlite.py; change both.

Parameters:

Name Type Description Default
url str | URL

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

required
**kwargs Any

Passed to sqlalchemy.ext.asyncio.create_async_engine. Give poolclass (with its own arguments) to pool connections differently.

{}

Returns:

Type Description
AsyncEngine

The engine. Dispose of it when you are done.

Schema

migrate async

migrate(engine: AsyncEngine) -> None

Upgrade the database to reflexr'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 reflexr'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.