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 Worker(IWorker): (source)
Constructor: Worker(app, driver, serializer, options, ...)
Consume leased jobs through isolated Orionis dependency scopes.
| Method | __init__ |
Initialize worker settings without opening backend connections. |
| Async Method | _consume |
Consume reservations within the shared job budget. |
| Async Method | _fail |
Record a terminal failure before deleting the owned reservation. |
| Async Method | _finish |
Acknowledge unfinished work or propagate an immediate job failure. |
| Async Method | _handle |
Choose a retry release or terminal failure from the job's retry policy. |
| Async Method | _release |
Release a cancelled reservation without replacing its cancellation. |
| Method | _validate |
Validate reservation routing, timeout and retry eligibility. |
| Async Method | execute |
Execute a reserved job with isolated DI and lease-aware completion. |
| Async Method | run |
Consume jobs concurrently until a stopping condition is reached. |
| Method | stop |
Request a graceful stop without waiting for active jobs. |
| Class Variable | __slots__ |
Undocumented |
| Instance Variable | _app |
Undocumented |
| Instance Variable | _concurrency |
Undocumented |
| Instance Variable | _connection |
Undocumented |
| Instance Variable | _driver |
Undocumented |
| Instance Variable | _failed |
Undocumented |
| Instance Variable | _failed |
Undocumented |
| Instance Variable | _processed |
Undocumented |
| Instance Variable | _queues |
Undocumented |
| Instance Variable | _retry |
Undocumented |
| Instance Variable | _running |
Undocumented |
| Instance Variable | _serializer |
Undocumented |
| Instance Variable | _sleep |
Undocumented |
| Instance Variable | _started |
Undocumented |
| Instance Variable | _stop |
Undocumented |
IApplication, driver: IQueueDriver, serializer: IJobSerializer, options: WorkerOptions, failed: IFailedJobRepository | FailedResolver | None = None):
(source)
¶
Initialize worker settings without opening backend connections.
| Parameters | |
app:IApplication | Application container used to create an isolated scope per job. |
driver:IQueueDriver | Backend used to reserve, release and acknowledge jobs. |
serializer:IJobSerializer | Serializer used to restore envelopes and registered jobs. |
options:WorkerOptions | Validated connection, queues, concurrency, lease and poll settings. |
failed:IFailedJobRepository | FailedResolver | None, optional | Failure repository or async resolver invoked on terminal failure. If None, terminal failures cannot be recorded. |
| Returns | |
None | Store dependencies and initialize counters and stop state. |
Consume reservations within the shared job budget.
| Parameters | |
stopbool | End this consumer when the backend has no ready job. |
maxint | None | Shared reservation limit, or None to run without a job limit. |
| Returns | |
None | Process reserved jobs and update the shared counters. |
| Raises | |
Exception | If reservation or execution raises an unhandled error. |
Record a terminal failure before deleting the owned reservation.
| Parameters | |
reserved:ReservedJob | Reservation whose lease must still be valid before recording failure. |
exception:Exception | Original job exception stored in the failure repository. |
| Returns | |
bool | True if deletion succeeds; False if the lease is expired or lost. |
| Raises | |
QueueConfigurationError | If no failed-job repository can be resolved. |
Exception | If repository resolution, failure recording or job deletion fails. |
Acknowledge unfinished work or propagate an immediate job failure.
| Parameters | |
context:JobContext | Lifecycle state containing the reservation's explicit transitions. |
raisebool | Raise a recorded failure when immediate execution requires propagation. |
| Returns | |
None | Delete unfinished work or preserve an existing lifecycle transition. |
| Raises | |
Exception | If deletion fails or a recorded failure must be propagated. |
ReservedJob, envelope: JobEnvelope | None, exception: Exception, *, raise_errors: bool):
(source)
¶
Choose a retry release or terminal failure from the job's retry policy.
| Parameters | |
reserved:ReservedJob | Current reservation used to release or record the failed attempt. |
envelope:JobEnvelope | None | Decoded retry policy, or None if the envelope is unavailable. |
exception:Exception | Original error used to determine retry eligibility and record failure. |
raisebool | Disable retries so immediate execution records a terminal failure. |
| Returns | |
None | Attempt a delayed retry release or terminal failure recording. |
| Raises | |
QueueConfigurationError | If a terminal failure has no available failed-job repository. |
Exception | If the backend or failure repository cannot complete the transition. |
Release a cancelled reservation without replacing its cancellation.
| Parameters | |
reserved:ReservedJob | Reservation whose execution was interrupted before completion. |
| Returns | |
None | Attempt lease release and log backend errors without re-raising them. |
Validate reservation routing, timeout and retry eligibility.
| Parameters | |
reserved:ReservedJob | Backend reservation containing routing metadata and attempt count. |
envelope:JobEnvelope | Decoded dispatch metadata and retry policy to match and enforce. |
| Returns | |
None | Accept a reservation that matches its envelope and retry limits. |
| Raises | |
QueuePayloadError | If routing metadata differs or the attempt count is invalid. |
QueueConfigurationError | If the timeout is missing or not shorter than worker retry_after. |
QueueRetryError | If attempts exceed max_tries or the retry deadline has expired. |
Execute a reserved job with isolated DI and lease-aware completion.
| Parameters | |
reserved:ReservedJob | Backend reservation containing the payload, attempt count and lease. |
raisebool, optional | Propagate execution errors and bypass automatic retries when true. |
| Returns | |
None | Attempt acknowledgement, retry release or terminal failure recording. Skip expired reservations without executing their jobs. |
| Raises | |
asyncio.CancelledError | If cancelled, after attempting to release unfinished work. |
Exception | If execution or completion fails and raise_errors is true. |
bool = False, max_jobs: int | None = None) -> int:
(source)
¶
Consume jobs concurrently until a stopping condition is reached.
| Parameters | |
stopbool, optional | Stop each consumer when no job is ready, even if delayed jobs remain. |
maxint | None, optional | Limit reservations across all consumers; None means no limit. |
| Returns | |
int | Number of processed reservations, including retried or failed jobs. |
| Raises | |
QueueConfigurationError | If already running or a supplied job limit is not a positive integer. |
asyncio.CancelledError | If cancelled, after in-flight consumers finish. |
Exception | If a consumer fails, after the remaining consumers finish. |
Request a graceful stop without waiting for active jobs.
| Returns | |
None | Set the stop event; active consumers complete their current work. |