THE ORIONIS API
Build with clarity.
Explore the building blocks of an async-first Python framework. Every module, class, and method — connected, searchable, and ready to build with.
class documentation
class RedisQueueDriver(IQueueDriver): (source)
Constructor: RedisQueueDriver(config, client)
Persist delayed jobs and fenced reservations through atomic Lua scripts.
| Method | __init__ |
Configure a lazy Redis transport and queue key namespace. |
| Async Method | _execute |
Execute one atomic Redis transition and translate transport errors. |
| Method | _keys |
Return colocated keys for one logical queue. |
| Async Method | clear |
Remove every job and reservation in one queue atomically. |
| Async Method | close |
Close the transport only when this driver owns its Redis client. |
| Async Method | delete |
Acknowledge a currently owned reservation atomically. |
| Async Method | push |
Store an envelope and its availability atomically. |
| Async Method | release |
Release a currently owned reservation with optional delay. |
| Async Method | reserve |
Claim one due job in queue priority order. |
| Async Method | size |
Count ready, delayed, and reserved jobs in one queue. |
| Class Variable | __slots__ |
Undocumented |
| Instance Variable | _client |
Undocumented |
| Instance Variable | _owns |
Undocumented |
| Instance Variable | _prefix |
Undocumented |
Configure a lazy Redis transport and queue key namespace.
| Parameters | |
config:RedisConfig | None, optional | Validated host, port, database, password, and namespace settings. Use the configuration entity's environment defaults when omitted. |
client:Redis or None, optional | Externally owned async client, primarily for explicit test fakes. |
| Returns | |
None | Initialize a lazy transport without opening network connections. |
| Raises | |
QueueStorageError | If the configuration type or queue namespace is invalid. |
TypeError | If an environment default has an invalid type. |
ValueError | If an environment default fails its configuration validation. |
async def _executeScript(self, script:
str, keys: tuple[ str, ...], arguments: tuple[ str | bytes | float, ...]) -> object:
(source)
¶
Execute one atomic Redis transition and translate transport errors.
| Parameters | |
script:str | Framework-owned Lua transition source. |
keys:tuple of str | Redis keys touched by the transition. |
arguments:tuple | Bound transition arguments. |
| Returns | |
object | Redis reply supplied by the executed script. |
| Raises | |
QueueStorageError | If Redis cannot perform the transition. |
Close the transport only when this driver owns its Redis client.
| Returns | |
None | Release driver-owned pooled Redis connections. |
Acknowledge a currently owned reservation atomically.
| Parameters | |
reserved:ReservedJob | Reservation owned by the acknowledging worker. |
| Returns | |
bool | Whether an owned, unexpired reservation was removed. |
Store an envelope and its availability atomically.
| Parameters | |
envelope:JobEnvelope | Immutable serialized job configuration and state. |
delay:float, optional | Seconds before the job becomes available. |
| Returns | |
str | Persisted job identifier. |
| Raises | |
QueueStorageError | If the identifier already exists or persistence fails. |
Release a currently owned reservation with optional delay.
| Parameters | |
reserved:ReservedJob | Reservation owned by the releasing worker. |
delay:float, optional | Seconds before the job becomes available again. |
| Returns | |
bool | Whether an owned, unexpired reservation was released. |
async def reserve(self, queues:
tuple[ str, ...], retry_after: float) -> ReservedJob | None:
(source)
¶
Claim one due job in queue priority order.
| Parameters | |
queues:tuple of str | Logical queues ordered from highest to lowest priority. |
retryfloat | Lease duration in seconds. |
| Returns | |
ReservedJob or None | Acquired reservation, or no available job. |
| Raises | |
QueueStorageError | If Redis returns invalid reservation data. |