On this page
Scheduler, the task executor: runs workflows via execute_workflow, drives topological task groups, handles crash recovery, PG LISTEN/NOTIFY, pause/resume/abort, and Overseer/transport event sinks.
#scheduler.src.orxtra.scheduler._executor
#scheduler.src.orxtra.scheduler._executor
#classify_error
def classify_error(error: Exception) -> ErrorCategoryClassify an exception into an error category.
#Scheduler
#run_consult
async def run_consult(self, agent: str, question: str, variable_values: dict[str, str] | None=None) -> str#run_workflow_check
async def run_workflow_check(self, execution: Execution) -> CheckResult#execute_workflow
async def execute_workflow(self, config: WorkflowConfig) -> None#_crash_recovery
async def _crash_recovery(self) -> NoneThree-pass idempotent crash recovery startup.
Acquires the run lock FIRST so that clean_orphaned does not mark our own run as orphaned.
#_setup_pg_listener
async def _setup_pg_listener(self) -> tuple[asyncio.Task[None] | None, asyncpg.pool.PoolConnectionProxy[Any] | None, object]Set up cross-process signal delivery via PG LISTEN.
Returns (listener_task, listener_conn, notification_callback).
#_cleanup_pg_listener
async def _cleanup_pg_listener(self, pg_listener_task: asyncio.Task[None] | None, pg_listener_conn: asyncpg.pool.PoolConnectionProxy[Any] | None, on_notification: Any) -> NoneClean up PG listener resources.
#_register_workflow_tasks
async def _register_workflow_tasks(self, config: WorkflowConfig) -> dict[str, UUID]Validate config and register all workflow tasks.
#_execute_task_groups
async def _execute_task_groups(self, config: WorkflowConfig, task_id_map: dict[str, UUID]) -> NoneExecute task groups in topological order.
#execute_task
async def execute_task(self, task: TaskSpec, parent_task_id: UUID | None, *, task_id: UUID | None=None, variables: dict[str, Any] | None=None) -> TaskResult#_stop_services
async def _stop_services(self) -> NoneStop all running service instances.
#abort
async def abort(self) -> None#pause
async def pause(self) -> None#resume
async def resume(self) -> None#is_paused
def is_paused(self) -> bool#add_overseer_sink
def add_overseer_sink(self, sink: EventSink[OverseerEvent]) -> NoneAdd an overseer event sink.
#remove_overseer_sink
def remove_overseer_sink(self, sink: EventSink[OverseerEvent]) -> NoneRemove an overseer event sink. No-op if not present.
#add_transport_sink
def add_transport_sink(self, sink: EventSink[TransportEvent]) -> NoneAdd a transport event sink.
Also registers the sink on all currently active sessions so that events from in-progress tasks are delivered immediately.
#remove_transport_sink
def remove_transport_sink(self, sink: EventSink[TransportEvent]) -> NoneRemove a transport event sink. No-op if not present.
Also removes the sink from all currently active sessions.
#_handle_control_signal
async def _handle_control_signal(self, run_id: UUID, new_status: str) -> NoneHandle a control signal from the trace layer.
#_handle_headless_event
async def _handle_headless_event(self, event: OverseerEvent) -> NoneDispatch a headless fallback for an event when no Overseer is configured.
Uses the same FALLBACK_BEHAVIORS/FALLBACK_HANDLERS tables as degraded mode.
#_send_overseer_event
async def _send_overseer_event(self, event: OverseerEvent) -> NoneSend an event to the Overseer with a verify-then-accept retry loop.
#_dispatch_to_overseer_sinks
async def _dispatch_to_overseer_sinks(self, event: OverseerEvent) -> NoneDispatch an OverseerEvent to all registered sinks.
Errors in individual sinks are logged but do not crash the scheduler. Uses a snapshot copy to allow safe concurrent mutation of the sinks list.
#_write_coherence_summary
async def _write_coherence_summary(self) -> NoneAsk the Overseer for a coherence summary and persist it.
Requires an Overseer session -- skipped in headless mode (no meaningful fallback).
#_check_session_handoff
async def _check_session_handoff(self) -> NoneCheck if the Overseer session needs handoff due to token usage.
Requires an Overseer session -- no-op in headless mode.
#_send_pending_advisories
async def _send_pending_advisories(self) -> NoneSend stored structural advisories to the Overseer (or headless fallback handler).