pub struct Runtime { /* private fields */ }Expand description
The runtime.
Implementations§
Source§impl Runtime
impl Runtime
Sourcepub async fn run_batch(
&self,
id: BatchId,
spec: &BatchSpec,
) -> Result<BatchReport, RuntimeError>
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.
Sourcepub async fn batch_report(
&self,
id: BatchId,
) -> Result<BatchReport, RuntimeError>
pub async fn batch_report( &self, id: BatchId, ) -> Result<BatchReport, RuntimeError>
The batch’s current tally, without processing anything.
Source§impl Runtime
impl Runtime
Sourcepub fn builder_on<B: FullBackend + 'static>(store: Arc<B>) -> RuntimeBuilder
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.
pub fn builder(store: Arc<dyn JournalStore>) -> RuntimeBuilder
Sourcepub fn events(&self) -> Option<&Arc<dyn EventStore>>
pub fn events(&self) -> Option<&Arc<dyn EventStore>>
The inbound-event store, if this runtime has one.
Sourcepub const fn budget(&self) -> &Budget
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.
Sourcepub async fn drill(&self) -> Result<DrillReport, RuntimeError>
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.
Sourcepub fn timers(&self) -> Option<&Arc<dyn TimerStore>>
pub fn timers(&self) -> Option<&Arc<dyn TimerStore>>
The timer store, if this runtime has one.
pub fn batches(&self) -> Option<&Arc<dyn BatchStore>>
Sourcepub async fn request_cancel(
&self,
run: RunId,
actor: &str,
reason: &str,
) -> Result<bool, RuntimeError>
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.
Sourcepub async fn set_halt(&self, reason: Option<&str>) -> Result<(), RuntimeError>
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.
Sourcepub async fn halted(&self) -> Result<Option<String>, RuntimeError>
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.
Sourcepub async fn cancellation(
&self,
run: RunId,
) -> Result<Option<Cancellation>, RuntimeError>
pub async fn cancellation( &self, run: RunId, ) -> Result<Option<Cancellation>, RuntimeError>
§Errors
If the store is unreachable.
Sourcepub fn policy(&self) -> Option<&Arc<dyn PolicyEngine>>
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.
Sourcepub async fn record_break_glass(
&self,
actor: &str,
roles: &[String],
reason: &str,
) -> Result<RunId, RuntimeError>
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.
Sourcepub fn journal(&self) -> &Arc<dyn JournalStore> ⓘ
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.
Sourcepub async fn case_of(&self, run: RunId) -> Result<Option<CaseId>, RuntimeError>
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.
pub fn store(&self) -> &Arc<dyn JournalStore> ⓘ
Sourcepub fn owner_id(&self) -> &str
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.
Sourcepub async fn run(
&self,
target: &str,
input: Tainted<Value>,
) -> Result<RunOutcome, RuntimeError>
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.
Sourcepub async fn spawn(
self: &Arc<Self>,
target: &str,
input: Tainted<Value>,
) -> Result<RunId, RuntimeError>
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.
pub async fn spawn_in_case( self: &Arc<Self>, target: &str, input: Tainted<Value>, case: CaseId, ) -> Result<RunId, RuntimeError>
Sourcepub async fn run_in_case(
&self,
target: &str,
input: Tainted<Value>,
case: CaseId,
) -> Result<RunOutcome, RuntimeError>
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.
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.
Sourcepub async fn run_plan(
&self,
plan: PlanIR,
input: Tainted<Value>,
) -> Result<RunOutcome, RuntimeError>
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.
Execute an explicit plan inside a long-lived case.
Sourcepub async fn replay(
&self,
run: RunId,
mode: Mode,
) -> Result<RunOutcome, RuntimeError>
pub async fn replay( &self, run: RunId, mode: Mode, ) -> Result<RunOutcome, RuntimeError>
Re-execute a recorded run from its journal.
Mode::Strictverifies determinism: every effect must match, and the run must not want any effect the journal lacks.Mode::Resumerecovers 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.
impl Runtime
Inbound event delivery.
Sourcepub async fn deliver(
&self,
event: &InboundEvent,
) -> Result<Delivery, RuntimeError>
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.
Sourcepub async fn deliver_to(
&self,
run: RunId,
event: &InboundEvent,
) -> Result<Delivery, RuntimeError>
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.
Sourcepub async fn sweep_events(&self, grace: Duration) -> Result<usize, RuntimeError>
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
impl Runtime
Sourcepub async fn sweep(
&self,
now: Timestamp,
event_grace: Duration,
) -> Result<SweepReport, RuntimeError>
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.
Sourcepub async fn census(&self, now: Timestamp) -> Result<Census, StoreError>
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.
Sourcepub async fn fire_timers(
&self,
now: Timestamp,
) -> Result<WokenRuns, RuntimeError>
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.
Sourcepub async fn answer_task(
&self,
id: TaskId,
decision: &Decision,
) -> Result<Delivery, RuntimeError>
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.
Sourcepub async fn decide_task(
&self,
id: TaskId,
decision: &Decision,
roles: &[String],
) -> Result<Delivery, RuntimeError>
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§
Auto Trait Implementations§
impl !RefUnwindSafe for Runtime
impl !UnwindSafe for Runtime
impl Freeze for Runtime
impl Send for Runtime
impl Sync for Runtime
impl Unpin for Runtime
impl UnsafeUnpin for Runtime
Blanket Implementations§
Source§impl<T> BorrowMut<T> for Twhere
T: ?Sized,
impl<T> BorrowMut<T> for Twhere
T: ?Sized,
Source§fn borrow_mut(&mut self) -> &mut T
fn borrow_mut(&mut self) -> &mut T
impl<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
Source§impl<T> CloneToUninit for Twhere
T: Clone,
impl<T> CloneToUninit for Twhere
T: Clone,
Source§impl<T> Instrument for T
impl<T> Instrument for T
Source§fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
fn instrument(self, span: Span) -> Instrumented<Self> ⓘ
Source§fn in_current_span(self) -> Instrumented<Self> ⓘ
fn in_current_span(self) -> Instrumented<Self> ⓘ
Source§impl<T> IntoEither for T
impl<T> IntoEither for T
Source§fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
fn into_either(self, into_left: bool) -> Either<Self, Self> ⓘ
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 moreSource§fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
fn into_either_with<F>(self, into_left: F) -> Either<Self, Self> ⓘ
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