pub struct DurableProcessWorker { /* private fields */ }Expand description
Reconstructable background-process worker.
Implementations§
Source§impl DurableProcessWorker
impl DurableProcessWorker
pub fn new(config: DurableProcessWorkerConfig) -> Self
pub fn config(&self) -> &DurableProcessWorkerConfig
pub async fn run_process( &self, registration: ProcessRegistration, execution_context: ProcessExecutionContext, cancellation: CancellationToken, ) -> Result<ProcessAwaitOutput, PluginError>
pub async fn run_process_with_scoped_effect_controller( &self, registration: ProcessRegistration, execution_context: ProcessExecutionContext, scoped_effect_controller: ScopedEffectController<'_>, cancellation: CancellationToken, ) -> Result<ProcessAwaitOutput, PluginError>
Sourcepub async fn run_process_segment_with_scoped_effect_controller(
&self,
registration: ProcessRegistration,
execution_context: ProcessExecutionContext,
scoped_effect_controller: ScopedEffectController<'_>,
cancellation: CancellationToken,
handover: Option<SegmentHandover>,
) -> Result<ProcessRunOutcome, PluginError>
pub async fn run_process_segment_with_scoped_effect_controller( &self, registration: ProcessRegistration, execution_context: ProcessExecutionContext, scoped_effect_controller: ScopedEffectController<'_>, cancellation: CancellationToken, handover: Option<SegmentHandover>, ) -> Result<ProcessRunOutcome, PluginError>
Run exactly one engine segment. Durable substrates use this method so a
non-terminal boundary can end the current substrate invocation; inline
callers continue to use the looping run_process_with_* methods.
Sourcepub async fn drive_pending_processes(&self) -> Result<(), PluginError>
pub async fn drive_pending_processes(&self) -> Result<(), PluginError>
Queue every non-terminal process this worker can claim and execute the runnable rows inline, driving each to a terminal state.
This is the sole inline executor for every process start: live tool and subagent starts, trigger deliveries, admin starts, session-open passes, and crash recovery all enter through this worklist. The drive:
- lists every non-terminal process (
ProcessRegistry::list_non_terminal); - claims the durable single-owner
ProcessLeaseover each — a process already leased live by another owner is skipped (it is being run by that owner right now) unless persisted liveness metadata proves that owner definitely dead, in which case the lease is reclaimed with the fenced CAS discipline ofProcessRegistry::reclaim_process_lease; either way a non-terminal process is re-run by exactly one owner (lease fencing); - queues the full worklist in a worker-scoped scheduler whose shared execution budget spans repeated host-driven passes, then runs claimed processes on this worker’s wired controller while renewing the lease across the long-running execution so a healthy recovery is not swept out from under itself;
- atomically writes the terminal outcome and releases the validated lease.
Idempotent by process_id: terminal processes are never in the worklist,
and a process that became terminal between the list and the claim is
detected after claiming and skipped, so re-running a recovery sweep does
not double-execute completed work.
Sourcepub async fn drain_owner_bound_work(
&self,
) -> Result<ProcessDrainReport, PluginError>
pub async fn drain_owner_bound_work( &self, ) -> Result<ProcessDrainReport, PluginError>
Graceful owner drain: terminalize this host’s own started OwnerBound
work as Abandoned{OwnerDrain} at close (ADR 0019).
This is an explicit host lever on the worker, never an implicit
consequence of closing a session. Processes are global and outlive any
one session ([ADR 0011]), so LashSession::close/park must not touch
them; a host that wants its in-flight owner-bound work terminalized at
shutdown calls this on the worker it is tearing down.
Drain sequence (the operations runbook owns the surrounding steps; this is the terminal-writing step):
- stop admitting new work to this worker;
- cancel or await the worker’s in-flight run tasks so they release their per-run leases — for Rerunnable in-flight work that is the whole story: stopping the local run task without any terminal write leaves the row non-terminal so the next worker re-runs it (its contract);
- call this lever: for every non-terminal OwnerBound row this exact
worker started (
first_started.owner == self.config.lease_owner), claim a fresh drain lease and, being the owner completing its own work, writeAbandoned{OwnerDrain}under it — the ordinary graceful completion path, respecting the single-writer rule.
A row still held by a live foreign lease (an in-flight run under one of
this worker’s own recovery incarnations that step 2 has not yet released)
is deferred rather than reclaimed, so the drain never races a still-live
run; such a row reaches Abandoned on the next drain pass or at a peer’s
recovery sweep. Rows started by a different owner, not-yet-started
OwnerBound rows (still claimable by anyone), Rerunnable rows, and
Externally-Owned rows are all left untouched.
[ADR 0011]: durable process registration is session-independent.
pub async fn request_process_cancel( &self, process_id: &str, reason: Option<String>, ) -> Result<(), PluginError>
Trait Implementations§
Source§impl Clone for DurableProcessWorker
impl Clone for DurableProcessWorker
Source§fn clone(&self) -> DurableProcessWorker
fn clone(&self) -> DurableProcessWorker
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read more