ORIONIS API REFERENCE

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

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 _handleFailure Choose a retry release or terminal failure from the job's retry policy.
Async Method _releaseCancelled Release a cancelled reservation without replacing its cancellation.
Method _validateReservation 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_resolver Undocumented
Instance Variable _processed Undocumented
Instance Variable _queues Undocumented
Instance Variable _retry_after Undocumented
Instance Variable _running Undocumented
Instance Variable _serializer Undocumented
Instance Variable _sleep Undocumented
Instance Variable _started Undocumented
Instance Variable _stop Undocumented
def __init__(self, app: IApplication, driver: IQueueDriver, serializer: IJobSerializer, options: WorkerOptions, failed: IFailedJobRepository | FailedResolver | None = None): (source)

Initialize worker settings without opening backend connections.

Parameters
app:IApplicationApplication container used to create an isolated scope per job.
driver:IQueueDriverBackend used to reserve, release and acknowledge jobs.
serializer:IJobSerializerSerializer used to restore envelopes and registered jobs.
options:WorkerOptionsValidated connection, queues, concurrency, lease and poll settings.
failed:IFailedJobRepository | FailedResolver | None, optionalFailure repository or async resolver invoked on terminal failure. If None, terminal failures cannot be recorded.
Returns
NoneStore dependencies and initialize counters and stop state.
async def _consume(self, *, stop_when_empty: bool, max_jobs: int | None): (source)

Consume reservations within the shared job budget.

Parameters
stop_when_empty:boolEnd this consumer when the backend has no ready job.
max_jobs:int | NoneShared reservation limit, or None to run without a job limit.
Returns
NoneProcess reserved jobs and update the shared counters.
Raises
ExceptionIf reservation or execution raises an unhandled error.
async def _fail(self, reserved: ReservedJob, exception: Exception) -> bool: (source)

Record a terminal failure before deleting the owned reservation.

Parameters
reserved:ReservedJobReservation whose lease must still be valid before recording failure.
exception:ExceptionOriginal job exception stored in the failure repository.
Returns
boolTrue if deletion succeeds; False if the lease is expired or lost.
Raises
QueueConfigurationErrorIf no failed-job repository can be resolved.
ExceptionIf repository resolution, failure recording or job deletion fails.
async def _finish(self, context: JobContext, *, raise_errors: bool): (source)

Acknowledge unfinished work or propagate an immediate job failure.

Parameters
context:JobContextLifecycle state containing the reservation's explicit transitions.
raise_errors:boolRaise a recorded failure when immediate execution requires propagation.
Returns
NoneDelete unfinished work or preserve an existing lifecycle transition.
Raises
ExceptionIf deletion fails or a recorded failure must be propagated.
async def _handleFailure(self, reserved: 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:ReservedJobCurrent reservation used to release or record the failed attempt.
envelope:JobEnvelope | NoneDecoded retry policy, or None if the envelope is unavailable.
exception:ExceptionOriginal error used to determine retry eligibility and record failure.
raise_errors:boolDisable retries so immediate execution records a terminal failure.
Returns
NoneAttempt a delayed retry release or terminal failure recording.
Raises
QueueConfigurationErrorIf a terminal failure has no available failed-job repository.
ExceptionIf the backend or failure repository cannot complete the transition.
async def _releaseCancelled(self, reserved: ReservedJob): (source)

Release a cancelled reservation without replacing its cancellation.

Parameters
reserved:ReservedJobReservation whose execution was interrupted before completion.
Returns
NoneAttempt lease release and log backend errors without re-raising them.
def _validateReservation(self, reserved: ReservedJob, envelope: JobEnvelope): (source)

Validate reservation routing, timeout and retry eligibility.

Parameters
reserved:ReservedJobBackend reservation containing routing metadata and attempt count.
envelope:JobEnvelopeDecoded dispatch metadata and retry policy to match and enforce.
Returns
NoneAccept a reservation that matches its envelope and retry limits.
Raises
QueuePayloadErrorIf routing metadata differs or the attempt count is invalid.
QueueConfigurationErrorIf the timeout is missing or not shorter than worker retry_after.
QueueRetryErrorIf attempts exceed max_tries or the retry deadline has expired.
async def execute(self, reserved: ReservedJob, *, raise_errors: bool = False): (source)

Execute a reserved job with isolated DI and lease-aware completion.

Parameters
reserved:ReservedJobBackend reservation containing the payload, attempt count and lease.
raise_errors:bool, optionalPropagate execution errors and bypass automatic retries when true.
Returns
NoneAttempt acknowledgement, retry release or terminal failure recording. Skip expired reservations without executing their jobs.
Raises
asyncio.CancelledErrorIf cancelled, after attempting to release unfinished work.
ExceptionIf execution or completion fails and raise_errors is true.
async def run(self, *, stop_when_empty: bool = False, max_jobs: int | None = None) -> int: (source)

Consume jobs concurrently until a stopping condition is reached.

Parameters
stop_when_empty:bool, optionalStop each consumer when no job is ready, even if delayed jobs remain.
max_jobs:int | None, optionalLimit reservations across all consumers; None means no limit.
Returns
intNumber of processed reservations, including retried or failed jobs.
Raises
QueueConfigurationErrorIf already running or a supplied job limit is not a positive integer.
asyncio.CancelledErrorIf cancelled, after in-flight consumers finish.
ExceptionIf a consumer fails, after the remaining consumers finish.
def stop(self): (source)

Request a graceful stop without waiting for active jobs.

Returns
NoneSet the stop event; active consumers complete their current work.

Undocumented

_concurrency = (source)

Undocumented

_connection = (source)

Undocumented

Undocumented

Undocumented

_failed_resolver = (source)

Undocumented

_processed: int = (source)

Undocumented

Undocumented

_retry_after = (source)

Undocumented

_running: bool = (source)

Undocumented

_serializer = (source)

Undocumented

Undocumented

_started: int = (source)

Undocumented

Undocumented