orangeqs.juice.task_manager._scheduler#
ScheduledTask scheduling for the task manager.
Module Contents#
Classes#
Serializes work per target service, so a service never runs two Tasks at once. |
|
A submitted, tracked instance of a Task with lifecycle state. |
|
A ScheduledTask state-transition record for publishing. |
|
Fires ScheduledTasks onto targets, honoring the shared serial line. |
Data#
Lifecycle state of a one-off ScheduledTask. |
|
Endstates of Scheduled Task. |
|
|
API#
- orangeqs.juice.task_manager._scheduler.ScheduledTaskStatus#
None
Lifecycle state of a one-off ScheduledTask.
- orangeqs.juice.task_manager._scheduler.TERMINAL_STATUSES: frozenset[str]#
‘frozenset(…)’
Endstates of Scheduled Task.
- orangeqs.juice.task_manager._scheduler.DispatchFn#
None
(target, task_id, message) -> result: send an occurrence, await its result.
- class orangeqs.juice.task_manager._scheduler.ServiceTaskLock#
Serializes work per target service, so a service never runs two Tasks at once.
Each target has one slot. Each Task for the same service tries to acquire the slot to be executed. This guarantees that at most a single Task is running on the target service at any given time.
- async acquire(target: str) None#
Take
target’s slot, waiting if held (returns at once on a free slot).
- async hold(target: str) collections.abc.AsyncGenerator[None]#
Hold
target’s slot for the block (waits if busy).
- class orangeqs.juice.task_manager._scheduler.ScheduledTask(/, **data: Any)#
Bases:
pydantic.BaseModelA submitted, tracked instance of a Task with lifecycle state.
Task describes what to run; ScheduledTask is the queued/scheduled record of running it.
- task_id: str | None#
None
The fresh task id of the running occurrence, if any (assigned at dispatch).
- task_data: dict[str, Any]#
None
The serialized task envelope, re-sent with a fresh id per occurrence.
- status: orangeqs.juice.task_manager._scheduler.ScheduledTaskStatus#
‘queued’
Current lifecycle state.
- run_at: datetime.datetime | None#
None
Earliest time this ScheduledTask should fire.
- submitted_at: datetime.datetime#
None
When the ScheduledTask was submitted (dispatch ordering tiebreak).
- is_eligible(now: datetime.datetime) bool#
Whether this ScheduledTask is a dispatch candidate at
now.
- class orangeqs.juice.task_manager._scheduler.ScheduledTaskEvent(/, **data: Any)#
Bases:
orangeqs.juice.messaging.EventA ScheduledTask state-transition record for publishing.
- class orangeqs.juice.task_manager._scheduler.Scheduler(dispatch: orangeqs.juice.task_manager._scheduler.DispatchFn, *, service_task_lock: orangeqs.juice.task_manager._scheduler.ServiceTaskLock | None = None, publish: collections.abc.Callable[[orangeqs.juice.task_manager._scheduler.ScheduledTaskEvent], None] | None = None, now_fn: collections.abc.Callable[[], datetime.datetime] = _utcnow, tick_interval: float = 1.0)#
Fires ScheduledTasks onto targets, honoring the shared serial line.
Driven by :meth:
tick, which takes an injectednowso dispatch is deterministic to test. Occurrence completion runs on background tasks; :meth:wait_idleawaits them in tests.- submit(scheduled_task: orangeqs.juice.task_manager._scheduler.ScheduledTask) orangeqs.juice.task_manager._scheduler.ScheduledTask#
Store a ScheduledTask and announce it as queued.
- get(scheduled_task_id: str) orangeqs.juice.task_manager._scheduler.ScheduledTask | None#
Return the Scheduled Task with
scheduled_task_id.
- scheduled_tasks() list[orangeqs.juice.task_manager._scheduler.ScheduledTask]#
Return every tracked ScheduledTask, for the dashboard snapshot query.
- cancel(scheduled_task_id: str) bool#
Cancel a Scheduled Task by id.
A queued Schedule stops firing. A running occurrence is not interrupted — the Scheduled Task is marked cancelled and its running Task is left to finish without recording a result.
Returns
Falseif the Scheduled Task is unknown or already terminal.
- async tick(now: datetime.datetime | None = None) None#
Run one pass: dispatch at most one ScheduledTask per free serial slot.