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.
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.
due_runs
async
¶
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 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 |
required |
**kwargs
|
Any
|
Passed to |
{}
|
Returns:
| Type | Description |
|---|---|
AsyncEngine
|
The engine. Dispose of it when you are done. |
Schema¶
migrate
async
¶
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.