Skip to main content

Runtime

Struct Runtime 

Source
pub struct Runtime { /* private fields */ }
Expand description

The runtime.

Implementations§

Source§

impl Runtime

Source

pub async fn run_batch( &self, id: BatchId, spec: &BatchSpec, ) -> Result<BatchReport, RuntimeError>

Start a batch, or continue one.

Passing an existing BatchId resumes it: processing picks up after the stored cursor, and any item that was reserved but never finished is replayed under its original run id rather than started again.

Returns when the source is exhausted or max_items is reached. A failing item does not stop the batch — that is the failure isolation a batch exists for, and a caller who wants the opposite should be running one run, not many.

Source

pub async fn batch_report( &self, id: BatchId, ) -> Result<BatchReport, RuntimeError>

The batch’s current tally, without processing anything.

Source§

impl Runtime

Source

pub fn builder_on<B: FullBackend + 'static>(store: Arc<B>) -> RuntimeBuilder

A builder with the whole case layer wired to one backend.

The one-call form of the six-cast litany: journal, cases, tasks, events, timers and memory all backed by store, which is what a deployment on a single RedbStore or PostgresStore means anyway. The à-la-carte methods remain — cases and friends override an individual store afterwards, and builder still starts from the journal alone for a plane that wants nothing else.

Blob storage is deliberately not included: bytes routinely live in a different system than rows, so blobs stays an explicit decision.

Source

pub fn builder(store: Arc<dyn JournalStore>) -> RuntimeBuilder

Source

pub fn cases(&self) -> Option<&Arc<dyn CaseStore>>

The case store, if this runtime has one.

Source

pub fn tasks(&self) -> Option<&Arc<dyn TaskStore>>

The worklist, if this runtime has one.

Source

pub fn events(&self) -> Option<&Arc<dyn EventStore>>

The inbound-event store, if this runtime has one.

Source

pub const fn budget(&self) -> &Budget

The ceilings every run under this runtime starts with.

Readable because “what is this plane allowed to spend” is a question an operator asks of a running system, and answering it by re-reading the config that should have been applied is how a misapplied budget stays invisible.

Source

pub fn blobs(&self) -> Option<&Arc<dyn BlobStore>>

The blob store, if this runtime has one.

Source

pub async fn drill(&self) -> Result<DrillReport, RuntimeError>

Run the live half of the case-layer drill with this plane’s own stores.

The wiring is the point of this method existing: the drill’s value is holding each case’s references against the stores the plane actually runs with, and an embedder assembling drill::Stores by hand could hand it a different blob store than the one their runs write to — a drill that passes over the wrong bucket. Stores this runtime does not have are reported as unchecked by the drill itself, which is a different and better answer than an error: a plane with no blob store has no bytes to lose.

§Errors

If this runtime has no case store — the drill walks cases, so there is nothing to drill — or if the case layer cannot be enumerated.

Source

pub fn timers(&self) -> Option<&Arc<dyn TimerStore>>

The timer store, if this runtime has one.

Source

pub fn batches(&self) -> Option<&Arc<dyn BatchStore>>

Source

pub async fn request_cancel( &self, run: RunId, actor: &str, reason: &str, ) -> Result<bool, RuntimeError>

Ask a run to stop, and drive the stop if nothing else will.

The request is durable before this returns, so an operator who gets an acknowledgement has one whether or not the run was reachable. What happens next depends on where the run is:

  • Suspended — nothing is executing, so this resumes the run itself. It observes the request at its first step boundary and unwinds.
  • Running here or elsewhere — the owner observes the request at its next step boundary. Nothing is interrupted mid-effect, deliberately: stopping between “announced” and “recorded” manufactures the in-doubt case the effect protocol exists to avoid.
  • Already concluded — the request is recorded and does nothing. A sealed run is not reopened by an operator changing their mind.

Returns whether this call recorded the request. A second caller gets false: the first asker stays on the record, because “who intervened” must not be rewritten by a retry.

§Errors

RuntimeError if the store is unreachable, or if resuming a suspended run fails.

Source

pub async fn set_halt(&self, reason: Option<&str>) -> Result<(), RuntimeError>

The stop request standing against a run, if any.

Stop this tenant from starting new work, or let it start again.

The emergency stop. Some(reason) halts, None lifts, and the reason is required because the next person to look will be somebody else, possibly at three in the morning, and why is the whole question.

What it stops, precisely. New admissions, across every instance, because the flag is in the store rather than in this process — a switch that stops only the instance it was thrown on is the in-process-counter failure arriving during an incident. Refusals are their own error, not a ceiling: a ceiling means not right now and invites a retry, which is exactly what somebody pulling this switch is trying to stop.

What it does not stop, deliberately. Runs already executing, and suspended runs resuming. Those are existing work, and refusing to let them continue would strand them mid-saga with reversals unrun — turning an incident into a second one. To stop work in flight, cancel it: that unwinds what it did and records who asked. This is the front door, not a power cut, and saying so is the difference between a control an operator can reason about and one they discover the shape of during an outage.

Requires a quota store; without one there is nowhere durable to keep the flag, and an emergency stop that a restart forgets is not one.

§Errors

If no quota store is wired, or the store is unreachable.

Source

pub async fn halted(&self) -> Result<Option<String>, RuntimeError>

Why this tenant is halted, if it is.

§Errors

If no quota store is wired, or the store is unreachable.

Source

pub async fn cancellation( &self, run: RunId, ) -> Result<Option<Cancellation>, RuntimeError>

§Errors

If the store is unreachable.

Source

pub fn policy(&self) -> Option<&Arc<dyn PolicyEngine>>

The policy engine, if this runtime has one.

Exposed because a surface that faces strangers has to be able to ask whether one exists at all. Inside the process the caller is the embedder’s own code and an absent engine is a choice; on a socket it is a hole, and the HTTP surface refuses to start without one.

Source

pub fn tenant(&self) -> &TenantId

Which tenant this plane runs as.

Source

pub async fn record_break_glass( &self, actor: &str, roles: &[String], reason: &str, ) -> Result<RunId, RuntimeError>

Record an operator deliberately crossing into this tenant, and return the run that holds the record.

Every other tenancy control keeps a cross-tenant read from being reached by accident. This is the designed exception, and the rule for it is that an exception without a record is indistinguishable from the breach it is meant to be. So the access is written into this tenant’s journal — the one whose data is about to be reached — in a sealed run of its own, exactly as a sweep writes its decisions. It therefore inherits the hash chain, the per-record signature and the Merkle inclusion, and the offline audit tool reports it without being taught that break-glass exists.

The record is written before any data is served, and a failure to write it is a failure to access. That direction is the whole control: the alternative — serve first, record best-effort — is a break-glass that works exactly as well when its own evidence is lost.

§Errors

If reason is blank, or if the record cannot be written. An unexplained exception is the thing this record exists to prevent, so it is refused rather than stored empty.

Source

pub fn journal(&self) -> &Arc<dyn JournalStore> ⓘ

The journal, for reads the runtime does not mediate.

This is also the embedder’s resumable output stream: JournalStore::read is a seq-cursored read over a run’s records, so polling from the last sequence seen is a durable, reconnect-safe progress feed that any instance can serve — the exact mechanism the A2A streaming surface is built on. There is deliberately no curated event type between an embedder and the records: a third vocabulary beside the journal’s and the wire’s would be a second truth wearing ergonomics, and the records are already the history the stream must agree with.

Source

pub async fn case_of(&self, run: RunId) -> Result<Option<CaseId>, RuntimeError>

Which case a run belongs to, or None if it belongs to none.

Read from the journal, not from a column beside it. The binding is stamped on the run’s own records at admission, so answering from there is answering from the plan of record — a case-store column saying the same thing would be a second copy of one fact, and the two could disagree about a run the case layer never saw.

It is the first question an operator surface asks, which is why it is a method rather than a documented one-liner: every caller was otherwise going to reach for the first record and read body.case off it, and a caller who reached for the last one instead would still be right today and wrong the moment a run is admitted before its case is known.

§Errors

If the journal cannot be read. An unknown run is Ok(None) rather than an error: no such run and a run in no case are both honest answers to this question, and neither is a fault.

Source

pub fn store(&self) -> &Arc<dyn JournalStore> ⓘ

Source

pub fn owner_id(&self) -> &str

This process instance’s identity, as it appears in run leases.

Not the agent’s name. Several instances of one agent are normal, and a lease is renewed without a fencing bump only when the holder is the same owner — so two instances sharing this string would each renew the other’s lease and both write to one run. See RuntimeBuilder::owner.

Public because “which instance holds this run” is a question an operator asks of a stuck system, and the answer is otherwise only in a store row.

Source

pub async fn run( &self, target: &str, input: Tainted<Value>, ) -> Result<RunOutcome, RuntimeError>

Execute a fresh run with no case attached.

The input arrives labelled, and that is the whole of the decision. Tainted::trusted is what an operator writes for a literal, a constant or a configuration value they vouch for; anything that came from outside — an inbound event, a queue message, a counterparty’s payload — keeps the label it arrived with.

This took a bare Value and admitted it as Trusted by default. A deployment whose runs are started by inbound events passed counterparty data straight in, and three controls went quiet at once: require_trusted protected fields were satisfied by attacker-chosen values, the egress ceiling had nothing untrusted to join with, and the journal recorded no contact with outside data. Nothing failed and the suite stayed green, which is the profile of every other default this crate declines to offer.

A run_trusted/run_tainted pair was tried first and is worse: it doubles every shape, puts run_trusted_in_case one word away from run_tainted_in_case, and still lets run_trusted(cap, payload) compile over data nobody vouched for. A label is a value — it can be computed, threaded through an adapter, or derived from where a message arrived, none of which a method name can do. spawn already worked this way; run was the outlier.

§Errors

RuntimeError::NoProvider when no skill provides target, and whatever admission refuses — policy, quota, a halted tenant, a lease.

Source

pub async fn spawn( self: &Arc<Self>, target: &str, input: Tainted<Value>, ) -> Result<RunId, RuntimeError>

Admit a run and let it proceed in the background, returning its id.

The asynchronous counterpart to run, for callers that want a handle rather than an answer — A2A’s return_immediately, a queue worker, an operator kicking something off.

Admission happens before this returns. The policy gate, the lease and the admission records are all written first, so a refusal is an error here and not a task that never appears, and the id handed back can be read immediately. What continues in the background is the work.

The run is durable, so a process that dies mid-flight leaves a journal another instance resumes; the background task is where the work happens, not where it is kept.

§Panics

Outside a Tokio runtime, as the rest of this crate’s timing does.

§Errors

As run — anything admission itself refuses.

Source

pub async fn spawn_in_case( self: &Arc<Self>, target: &str, input: Tainted<Value>, case: CaseId, ) -> Result<RunId, RuntimeError>

Source

pub async fn spawn_correlated( self: &Arc<Self>, target: &str, input: Tainted<Value>, case_kind: &str, keys: &[CorrelationKey], ) -> Result<RunId, RuntimeError>

Source

pub async fn run_in_case( &self, target: &str, input: Tainted<Value>, case: CaseId, ) -> Result<RunOutcome, RuntimeError>

Start a new immutable run inside a case that already exists.

Source

pub async fn run_correlated( &self, target: &str, input: Tainted<Value>, case_kind: &str, keys: &[CorrelationKey], ) -> Result<RunOutcome, RuntimeError>

Join or open a case by business key, then run.

Correlation happens before planning, because which case a message belongs to is a question of fact, not of judgement: it is a deterministic lookup on business keys, never a model call. If an open case matches any key the run joins it; otherwise a case is opened.

Source

pub async fn run_plan( &self, plan: PlanIR, input: Tainted<Value>, ) -> Result<RunOutcome, RuntimeError>

Execute an explicit multi-step plan.

The plan is validated and frozen before the first step runs: one that would fail at step seven must not begin at step one.

Source

pub async fn run_plan_correlated( &self, plan: PlanIR, input: Tainted<Value>, case_kind: &str, keys: &[CorrelationKey], ) -> Result<RunOutcome, RuntimeError>

Execute an explicit plan inside a long-lived case.

Source

pub async fn replay( &self, run: RunId, mode: Mode, ) -> Result<RunOutcome, RuntimeError>

Re-execute a recorded run from its journal.

  • Mode::Strict verifies determinism: every effect must match, and the run must not want any effect the journal lacks.
  • Mode::Resume recovers a crashed run: history is replayed, then execution continues live from wherever the record ends.

Either way no external effect is performed for anything already in the journal. That is the whole point — a resumed run does not re-issue the invoice it already issued.

Source§

impl Runtime

Inbound event delivery.

Source

pub async fn deliver( &self, event: &InboundEvent, ) -> Result<Delivery, RuntimeError>

Deliver an inbound event, resuming whichever run was waiting for it.

§Ordering

The event is stored before anyone looks for a waiter. That ordering is the whole reason this works: a message can arrive before its run reaches the wait, and one that is matched-then-discarded leaves that run waiting forever for something that already happened.

A Delivery::Buffered result is therefore normal and not an error — it means “held until someone asks”. Only the sweep (EventStore::sweep_unclaimed) decides an event is genuinely unroutable, because that is a claim about the future rather than about this instant.

Source

pub async fn deliver_to( &self, run: RunId, event: &InboundEvent, ) -> Result<Delivery, RuntimeError>

Deliver an inbound event to exactly run.

This is the task-addressed counterpart to Runtime::deliver. It is used by protocols such as A2A where a follow-up carries a concrete task id. Correlation alone is insufficient there: two tasks may wait on the same business key, and resuming the oldest would violate the request.

The event store atomically inserts and claims the event for this run. A run that is not waiting leaves no buffered event behind for another run.

Source

pub async fn sweep_events(&self, grace: Duration) -> Result<usize, RuntimeError>

Retire events that nobody claimed within grace.

A non-empty dead-letter list means a correlation key is wrong somewhere: the message arrived, was held, and no run ever asked for it. That is the failure which otherwise presents as a process silently never completing, so it is worth alerting on rather than logging. grace is a std::time::Duration for the reason StepCtx::deadline’s warn_before is: a negative grace window is meaningless, and the signed type could express it — a cutoff moved forward of now, retiring events that had not yet had their chance. It is also the Duration the caller has.

Source§

impl Runtime

Source

pub async fn sweep( &self, now: Timestamp, event_grace: Duration, ) -> Result<SweepReport, RuntimeError>

Run one sweep.

Idempotent: a transition already applied is not applied twice, so calling this on a timer — or twice by accident, or from two instances — is safe.

now is passed in rather than read so the caller controls the clock. That keeps the sweeper testable at all, and lets a simulation drive it through a year in milliseconds.

Source

pub async fn census(&self, now: Timestamp) -> Result<Census, StoreError>

Read the gauges: how much is open right now.

Queried from the stores rather than accumulated, because a counter incremented on open and decremented on close drifts permanently the first time a process dies between the state change and the emission — and it drifts plausibly, which is worse than drifting obviously.

Stores that were never configured contribute zero rather than an error: a plane with no human tasks has no open tasks, and refusing to report the gauges it does have would be the silent failure this exists to prevent.

Source

pub async fn fire_timers( &self, now: Timestamp, ) -> Result<WokenRuns, RuntimeError>

Wake the runs whose instant has arrived.

The shape mirrors event delivery exactly, and for the same reasons: claim atomically so two sweepers cannot wake one run twice, record the wake-up as the sleeping effect’s result, then let replay do the rest. The resumed run reads the timer back like any other completed effect, so none of the suspension machinery exists twice.

One run’s failure does not block the runs behind it. The batch used to stop at the first error, which handed a single unresumable run a veto over every later wake in the tick — and the failed run heals anyway: its wake was recorded under a lease this pass acquired and never released, so the recovery pass finds it once that lease lapses. The failure is counted rather than returned, because a caller that got a count and an error would have neither.

Source

pub async fn answer_task( &self, id: TaskId, decision: &Decision, ) -> Result<Delivery, RuntimeError>

Record a human decision and resume the run waiting on it.

The decision is delivered as an ordinary inbound event, so it travels the same buffered, deduplicated, single-consumer path as any other message — including the case where the run has not yet reached its wait.

Source

pub async fn decide_task( &self, id: TaskId, decision: &Decision, roles: &[String], ) -> Result<Delivery, RuntimeError>

Complete a task on behalf of a person, enforcing eligibility first.

The claim is checked before the decision is recorded: an approval from somebody who was not permitted to give it is worse than no approval, because it looks like one.

Trait Implementations§

Source§

impl Clone for Runtime

Source§

fn clone(&self) -> Runtime

Returns a duplicate of the value. Read more
1.0.0 (const: unstable) · Source§

fn clone_from(&mut self, source: &Self)

Performs copy-assignment from source. Read more
Source§

impl Debug for Runtime

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<T> CloneToUninit for T
where T: Clone,

Source§

unsafe fn clone_to_uninit(&self, dest: *mut u8)

🔬This is a nightly-only experimental API. (clone_to_uninit)
Performs copy-assignment from self to dest. Read more
Source§

impl<T> DynClone for T
where T: Clone,

Source§

fn __clone_box(&self, _: Private) -> *mut ()

Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> FromRef<T> for T
where T: Clone,

Source§

fn from_ref(input: &T) -> T

Converts to this type from a reference to the input type.
Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> IntoEither for T

Source§

fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ

Converts self into a Left variant of Either<Self, Self> if into_left is true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
where F: FnOnce(&Self) -> bool,

Converts self into a Left variant of Either<Self, Self> if into_left(&self) returns true. Converts self into a Right variant of Either<Self, Self> otherwise. Read more
Source§

impl<Unshared, Shared> IntoShared<Shared> for Unshared
where Shared: FromUnshared<Unshared>,

Source§

fn into_shared(self) -> Shared

Creates a shared type from an unshared type.
Source§

impl<T> MaybeSend for T
where T: Send,

Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> ToOwned for T
where T: Clone,

Source§

type Owned = T

The resulting type after obtaining ownership.
Source§

fn to_owned(&self) -> T

Creates owned data from borrowed data, usually by cloning. Read more
Source§

fn clone_into(&self, target: &mut T)

Uses borrowed data to replace owned data, usually by cloning. Read more
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = Infallible

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, <T as TryFrom<U>>::Error>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more