pub struct StepCtx<'a> { /* private fields */ }Expand description
Per-step execution context.
Implementations§
Source§impl<'a> StepCtx<'a>
impl<'a> StepCtx<'a>
Sourcepub async fn commission(
&mut self,
capability: &str,
input: Tainted<Value>,
) -> Result<Tainted<Value>, StepError>
pub async fn commission( &mut self, capability: &str, input: Tainted<Value>, ) -> Result<Tainted<Value>, StepError>
Commission another agent on this plane, and journal that you did.
The hand-off, done properly. A skill cannot hold an Arc<Runtime> —
the runtime needs the skill before the skill can have the runtime — so
commissioning belongs to the runtime and is reached through here.
Three properties, none of them optional:
- Journaled, so a strict replay reads the answer back instead of commissioning the work a second time. A skill that called another runtime inline would be doing non-deterministic work outside the journal, and replay would re-run the whole room.
- The label travels. A specialist’s answer is untrusted — it came from a model — and the next agent is commissioned with that label intact, so the receiving run’s taint gates judge what they were actually given.
- The cost comes back, and is billed to the commissioning run, so an orchestrator’s ceiling bounds the work it ordered rather than its own idling.
The answer is untrusted whatever the org chart says: another agent’s output is somebody else’s data.
§Errors
StepError if the sub-run fails. Reported as in doubt rather than
did not happen: the commissioned agent may have performed effects
before failing, and its own journal is where that is answered.
Sourcepub fn manifest(&self) -> Option<&Manifest>
pub fn manifest(&self) -> Option<&Manifest>
The declaration this agent runs under.
None when the runtime was wired by builder calls instead. A skill asks
its context which agent it is part of — it does not hold a manifest of
its own, because an agent has skills, and a skill carrying a separate
copy of the agent’s declaration could disagree with the agent about what
the agent is.
What a skill typically wants from it: the system prompt
(Identity::system_prompt),
the model role to call, and
output_schema.
Sourcepub async fn complete(
&mut self,
prompt: &Tainted<Value>,
) -> Result<Tainted<Completion>, StepError>
pub async fn complete( &mut self, prompt: &Tainted<Value>, ) -> Result<Tainted<Completion>, StepError>
Complete a prompt on the manifest’s own model, through the plane’s own driver.
The model-call counterpart to call_tool, and it
closes the same gap. A declarative agent resolves its model from the
declaration and its driver from the plane’s registry; a hand-written
skill had to carry an Arc<dyn ModelProvider> field and name a model
in code — so the skill held wiring its manifest never described, and
the file’s models.privileged governed the declarative tier while the
coded tier read it or did not. This is the path where it cannot be
ignored: the privileged role supplies the model and its reviewed
ceilings, the plane supplies the driver registered under the role’s
provider name, and the manifest’s egress ceiling rides the call.
let completion = cx.complete(&prompt).await?;An explicit (provider, model) remains one construction away —
cx.sink_with(&prompt, |value| ModelCall::new(provider, model, value))
— which is the honest spelling for a call the manifest does not govern.
§Errors
StepError when this skill runs under no manifest, when the manifest
declares no privileged model, when no driver is registered under the
role’s provider name, or whatever the dispatch itself refuses.
Sourcepub async fn complete_with<F>(
&mut self,
prompt: &Tainted<Value>,
tune: F,
) -> Result<Tainted<Completion>, StepError>
pub async fn complete_with<F>( &mut self, prompt: &Tainted<Value>, tune: F, ) -> Result<Tainted<Completion>, StepError>
complete, with the call adjusted before dispatch.
The closure receives the fully-resolved call — model, role ceilings and egress ceiling already applied — and may add what only the skill knows:
let completion = cx
.complete_with(&prompt, |call| call.expecting(schema.clone()))
.await?;It runs after the manifest’s declarations are applied, so a skill can
tighten or reshape the call; what it cannot do is dodge the declared
gate, which still refuses a model the manifest never named.
§Errors
As complete.
Sourcepub async fn call_tool(
&mut self,
tool: ToolId,
arguments: Tainted<Value>,
) -> Result<Tainted<Value>, StepError>
pub async fn call_tool( &mut self, tool: ToolId, arguments: Tainted<Value>, ) -> Result<Tainted<Value>, StepError>
Call a tool through the plane’s own catalogue.
§Why this exists, and why the obvious alternative is a hole
A declarative agent gets its ToolCatalog from the runtime. A
hand-written skill had to construct and carry one:
ToolCall::prepare(&self.catalog, Arc::clone(&self.client), id, args)?and nothing bound self.catalog to the manifest governing that skill.
ToolCatalog::from_manifest is the right primitive and it is one call
away — but the obvious thing, hand-building a catalogue with the tools
you know you call, compiles, runs, and grants the skill reach its
declaration never described. Worse, it can be laxer: a
ToolSafety::read_only entry for a tool the manifest calls mutating
exempts it from the whole-value taint gate and carries
Recovery::Retry, so a timed-out
money-moving call is sent a second time.
RuntimeBuilder::try_build
refuses exactly that divergence — for the plane’s catalogue. A
catalogue built inside a skill never passed under that check. So this is
the same dispatch a declarative agent performs, over the same checked
catalogue, and the drift is unrepresentable rather than merely
discouraged.
§Everything else is unchanged
The manifest gate still refuses a tool this agent’s declaration does not
grant, the protected-field rules still have to match, the egress ceiling
still applies, and the result still comes back
Tainted and untrusted. This narrows what a skill can reach; it grants
nothing.
let overdue = cx
.call_tool(ToolId::new("obsd", "list_overdue_processes"), args)
.await?;§Errors
StepError when this plane has no tool catalogue, when the tool is not
in it, when this agent’s manifest does not grant it, or when the
arguments’ label is refused at the sink.
pub fn run_id(&self) -> RunId
pub fn step_id(&self) -> StepId
Sourcepub fn rng(&mut self) -> &mut impl Rng
pub fn rng(&mut self) -> &mut impl Rng
A deterministic random source.
Seeded from (run_id, step) rather than journaled per draw: the sequence
is reproducible by construction, so replay reproduces it for free and the
journal carries no entropy records at all. Cheaper and stronger than
recording each value — there is no way for the recorded and recomputed
streams to disagree.
Sourcepub async fn now(&mut self) -> Result<Timestamp, StepError>
pub async fn now(&mut self) -> Result<Timestamp, StepError>
The current instant, as a journaled effect.
On replay this returns the instant the original run saw, not the instant now — which is why a replayed run makes the same time-dependent decisions as the run it is reproducing.
Sourcepub async fn note(&mut self, text: impl Into<String>) -> Result<(), StepError>
pub async fn note(&mut self, text: impl Into<String>) -> Result<(), StepError>
Record structured reasoning in the journal, adjacent to the effects it explains.
Adjacency is the point: a note next to the action it claims to justify makes reasoning-versus-action mismatch detectable after the fact and testable under replay. A summary written at the end of a run cannot do that, because by then the ordering evidence is gone.
Sourcepub async fn effect<E: Effect>(
&mut self,
effect: E,
) -> Result<Tainted<E::Output>, StepError>
pub async fn effect<E: Effect>( &mut self, effect: E, ) -> Result<Tainted<E::Output>, StepError>
Perform (or replay) an effect, repeating it if it fails and repeating is safe.
The whole determinism boundary is this function.
§Why one loop covers both replay and live execution
A retry sequence is history like any other. Each attempt has its own effect key (the attempt number is hashed in), so the journal holds attempts 1..N as ordinary consecutive effects, and replay walks them the same way it walks anything else. There is no separate “replay the retries” path to drift out of sync with the live one.
§History outranks policy
While history has attempts left, they are consumed regardless of what
the current RetryPolicy says. A run that made four attempts under
yesterday’s policy still made four attempts, and a replay under a
two-attempt policy that stopped early would leave unconsumed records and
report divergence for a run that did nothing wrong. The policy governs
only what happens after history runs out.
§The result is labelled
An effect is how the deterministic zone reaches the outside world, so
what comes back is the outside world’s data. It arrives as
Tainted, labelled from the effect’s own Effect::trust
declaration — which defaults to untrusted.
That is what makes the architecture hold rather than merely be described. A tool result flowing into a downstream step’s input is labelled automatically, so the replan refusal and the taint gate see it without the skill author having to remember; and a skill that wants to treat a tool response as trusted has to say so, in a call that leaves a record.
Sourcepub async fn sleep_until(&mut self, until: Timestamp) -> Result<(), StepError>
pub async fn sleep_until(&mut self, until: Timestamp) -> Result<(), StepError>
Suspend until an instant, durably.
The run’s frame is persisted and the task is dropped: a sleeping run costs a row, not a thread. A sweep wakes it when the instant arrives, so a plane can hold as many sleeping runs as it has disk and a restart loses none of them.
The instant is recorded under an effect key, so replay reads it back rather than sleeping again — and a run that slept until Tuesday still says Tuesday when it is audited next year.
Sourcepub async fn sleep(&mut self, how_long: Duration) -> Result<(), StepError>
pub async fn sleep(&mut self, how_long: Duration) -> Result<(), StepError>
Suspend for a duration.
The duration is resolved to an instant through the journaled clock, so the wake time is a recorded fact rather than a formula re-evaluated on every replay.
Sourcepub async fn sink_with<E, B, F>(
&mut self,
args: &Tainted<Value>,
build: F,
) -> Result<Tainted<E::Output>, StepError>
pub async fn sink_with<E, B, F>( &mut self, args: &Tainted<Value>, build: F, ) -> Result<Tainted<E::Output>, StepError>
Send a labeled value into a sink, handing the value to the effect and the gate in one motion.
The closure receives the inner value and builds the effect from it, so the bytes the gates check and the bytes the effect sends are one argument rather than two the caller must keep in agreement:
let completion = cx
.sink_with(&prompt, |value| ModelCall::new(provider, model, value))
.await?;This replaces the two-pass spelling — ModelCall::new(.., prompt.peek().clone()) beside cx.sink(call, &prompt) — where the
same data was written twice and a runtime check caught the versions
drifting apart. Passing it once removes the drift at the API instead of
detecting it afterwards; the byte-for-byte binding check still runs
underneath, because a custom effect could bind something other than
what its constructor was handed, and that is a driver bug worth a loud
refusal.
Fallible construction composes: the closure may return
Result<E, impl Into<StepError>> — see BuildsEffect — which is
what ToolCall::prepare needs.
§Errors
Whatever the closure refuses with, and everything sink
refuses: the egress ceiling, the journal ceiling, the whole-value taint
gate, and the per-field provenance rules.
Sourcepub async fn sink<E: Effect>(
&mut self,
effect: E,
args: &Tainted<Value>,
) -> Result<Tainted<E::Output>, StepError>
pub async fn sink<E: Effect>( &mut self, effect: E, args: &Tainted<Value>, ) -> Result<Tainted<E::Output>, StepError>
Send a labeled value into a sink, enforcing the information-flow gates.
Prefer sink_with, which hands the value to the
effect and the gate in one motion. This form remains for effects that
bind their outbound value internally — a governed media fetch derives
its bound arguments from the URL it was constructed over — and for
callers holding an effect built elsewhere.
Two checks, both of which are the reason labels exist at all:
- Egress ceiling — a value’s sensitivity may not exceed what the sink is allowed to receive. This is the exfiltration path that actually matters: not the network, but a legitimate-looking call carrying a secret that was read three steps ago.
- Authority-bearing fields — a mutating sink either refuses all untrusted arguments or declares protected JSON fields with explicit trust, source, and sensitivity constraints.
Both are judged over the effective label at this sink: the base label
improved by exactly the release marks whose destination names this
sink’s identity. A release for tool://ledger/transfer moves nothing
at tool://mail/send.
Sourcepub async fn release(
&mut self,
value: Tainted<Value>,
release: Release,
) -> Result<Tainted<Value>, StepError>
pub async fn release( &mut self, value: Tainted<Value>, release: Release, ) -> Result<Tainted<Value>, StepError>
Grant a destination-scoped release over a whole value or selected structured fields.
The release is policy-authorized and permanently records the releaser,
basis, field scope, destination, evidence, and prior label. It does
not relabel the value: it attaches release marks, and only the
sink whose identity equals the release’s destination — the
provenance-style name, tool://server/name, model:provider/model —
computes an improved effective label from them. Everywhere else the
value keeps its base label: joined into other values, written to
memory, read as label().trust, it is still what it was. A selected
field release is accepted only when the value was assembled with
Tainted::object or
Tainted::array, so precision can never
be invented after provenance was flattened.
Source§impl StepCtx<'_>
Case-scoped operations.
impl StepCtx<'_>
Case-scoped operations.
Available only when the runtime was built with a case store and the run was admitted with correlation keys. A run without a case is a perfectly ordinary run — it simply has no long-lived state to reach.
Sourcepub fn correlation(&self) -> &[CorrelationKey]
pub fn correlation(&self) -> &[CorrelationKey]
The business keys this run’s case is identified by.
Empty when the run has no case. These are the keys as recorded when the run bound to the case, not as the case stands now: a case accumulates keys over months, and reading the store here would make a resumed run see a set the live run never did.
The intended use is scoping durable state to the party a run is about —
Recall::about(cx.correlation_value("meter")?) reads back exactly what a
declarative agent’s subject: "$correlation/meter" wrote.
Sourcepub fn correlation_value(&self, namespace: &str) -> Option<&str>
pub fn correlation_value(&self, namespace: &str) -> Option<&str>
One correlation value by namespace.
None for a run with no case, and for a namespace the case is not keyed
by. Two keys sharing a namespace is a correlation the deployment set up,
not something to arbitrate here, so the first in canonical order wins and
the choice is stable across runs rather than dependent on store order.
Sourcepub async fn case_state(
&mut self,
) -> Result<(Tainted<Value>, CaseVersion), StepError>
pub async fn case_state( &mut self, ) -> Result<(Tainted<Value>, CaseVersion), StepError>
Read the case’s opaque state, and the revision it was read at.
A journaled effect, so a replay reads back what the live run saw rather than whatever the case holds now. Case state is mutable storage shared by every run on the case; reading it is as non-deterministic as reading a clock, and treating it as free was a hole in exactly the property this crate exists to provide.
The version comes back with the value because put_case_state needs
it. Returning the value alone is what makes a lost update easy to write.
§Errors
StepError if this run has no case, or the read fails.
Sourcepub async fn draw(
&mut self,
id: &AuthorityId,
amount: Spend,
) -> Result<Drawn, StepError>
pub async fn draw( &mut self, id: &AuthorityId, amount: Spend, ) -> Result<Drawn, StepError>
Draw on a standing authority, or be refused.
The ceiling that outlives a run: a customer’s approved spend, a purchase
order, a subscription mandate. A Budget bounds
this run and a TenantQuota bounds a billing
period; neither can express an authorization somebody granted once and may
take back.
Journaled, so a replay reads the receipt rather than consuming again, and
idempotent across retries of the same call — see
DrawOnAuthority for why the
deduplication key is the dispatch rather than the effect.
Expiry is evaluated against this run’s journaled clock, so a replay reaches the verdict the live run did rather than today’s.
§Errors
StepError::Store when no authority store is wired, and the effect’s
own error carrying whichever of the five refusals applies — unknown,
exhausted, out of draws, revoked, or expired.
Sourcepub async fn recall(
&mut self,
query: Recall,
) -> Result<Vec<Tainted<MemoryItem>>, StepError>
pub async fn recall( &mut self, query: Recall, ) -> Result<Vec<Tainted<MemoryItem>>, StepError>
Recall what this agent remembers about a subject.
§Every item comes back labelled from its provenance
Never from its content. Text asserting its own reliability is the cheapest thing an attacker can write, so a memory derived from a model, a peer or an inbound message stays untrusted however many times it is re-read — and reaching a mutating sink with it takes the same journaled release as any other untrusted value.
That is the defence against the attack this whole module is shaped by: a poisoned write becomes a standing instruction only if something later treats it as one.
§Journaled, and replayed by version
The selection is recorded — ids, versions, content digests — and a replay re-materialises exactly those versions rather than re-running the search. So a run replayed after the corpus changed reads what it read, not what a fresh ranking would return now.
§Errors
StepError if this plane has no memory store, if the recall fails, or
if a version this run read can no longer be reproduced — which is a
deliberate loud failure, not an empty result: a memory that was forgotten
makes the history that used it unreplayable, and saying so beats
replaying a different memory.
Sourcepub async fn embed(
&mut self,
embedder: Arc<dyn Embedder>,
text: Tainted<String>,
) -> Result<Tainted<Vec<f32>>, StepError>
pub async fn embed( &mut self, embedder: Arc<dyn Embedder>, text: Tainted<String>, ) -> Result<Tainted<Vec<f32>>, StepError>
Turn text into a vector, on the record.
The vector this returns is what SemanticQuery::embedding wants, and
going through here rather than calling an embedding client directly is
what makes semantic retrieval replayable at all: the query vector is in
the retrieval effect’s key, and an embedding service is under no
obligation to return the same floats twice. Journaled, so a strict replay
reads the vector back instead of asking again — and so the call is
metered and the model revision that produced it is on the record beside
the numbers.
The text carries its own label, and the returned vector carries it too: a vector derived from an untrusted document is untrusted, and sending confidential text to an embedding service is an egress like any other.
§Errors
Whatever the effect protocol reports — a refused sink, an exhausted budget, or the embedder’s own failure.
Sourcepub async fn semantic_recall(
&mut self,
retriever: Arc<dyn SemanticRetriever>,
query: Tainted<SemanticQuery>,
) -> Result<Vec<(Tainted<MemoryItem>, f32)>, StepError>
pub async fn semantic_recall( &mut self, retriever: Arc<dyn SemanticRetriever>, query: Tainted<SemanticQuery>, ) -> Result<Vec<(Tainted<MemoryItem>, f32)>, StepError>
Rank governed memories through a derived semantic index.
The retriever returns only immutable (id, version, digest) commitments
and scores. The selection is journaled; live execution and replay then
materialize exact versions from the authoritative memory store and
verify scope and digest before exposing content.
Sourcepub async fn remember(
&mut self,
write: MemoryWrite,
content: Tainted<Value>,
) -> Result<u64, StepError>
pub async fn remember( &mut self, write: MemoryWrite, content: Tainted<Value>, ) -> Result<u64, StepError>
Remember something, as a new version.
Journaled: a replay that wrote again would append a second version of a memory this run wrote once, and the version number the run went on to use would be wrong.
Trust, provenance and sensitivity are derived from content. They are
not fields the caller can declare: allowing a skill to store untrusted
model output with trust: Trusted would be an unjournaled release and a
cross-session laundering primitive.
§Errors
StepError if this plane has no memory store, or the write fails.
Sourcepub async fn sweep_expired_memories(&mut self) -> Result<usize, StepError>
pub async fn sweep_expired_memories(&mut self) -> Result<usize, StepError>
Atomically erase memories expired at the run’s journaled clock.
Legal holds remain authoritative in the backend. The cutoff and removed count are journaled, so strict replay reports the historical decision without mutating memory a second time.
Sourcepub async fn compact(
&mut self,
into: Compaction,
sources: &[Tainted<MemoryItem>],
provider: Arc<dyn ModelProvider>,
model: ModelId,
) -> Result<u64, StepError>
pub async fn compact( &mut self, into: Compaction, sources: &[Tainted<MemoryItem>], provider: Arc<dyn ModelProvider>, model: ModelId, ) -> Result<u64, StepError>
Summarise memories into a new, derived memory.
§The label is derived, never declared
This is the difference between compact and
remember. A writer declares where ordinary content
came from; a summary’s provenance is not a matter of opinion — it is the
join of what was summarised, plus the model that wrote it. Letting a
caller declare it would make compaction the laundering step: read three
untrusted memories, summarise, call the result trusted, and every gate
downstream has nothing to act on.
So the summary is untrusted whenever any input is, carries every input’s sources, and takes the highest sensitivity of any of them.
§It records what it was made from
Sources are recorded with the exact versions read. That is what makes a
summary repairable: a poisoned memory does not stop being a problem
when it is forgotten, because its content keeps arriving in every summary
that absorbed it. MemoryStore::derivatives walks that edge, and
forget_cascading is the form an erasure request needs.
§Compaction is an egress decision
It sends the memories to a model. So Compaction::max_sensitivity
bounds what that model may be shown, and it defaults to Public —
summarising is otherwise the way to move confidential content past a
ceiling that stops every other path, while looking like housekeeping.
§The originals stay
Compaction adds; it does not delete. What a summary is for — fitting a context window — is a reason to stop reading the originals, not a reason to destroy the only record of what the summary claims to represent.
§Errors
StepError if this plane has no memory store, if the model call fails,
or if the write fails.
Sourcepub async fn form_memories(
&mut self,
formation: Formation,
source: Tainted<Value>,
provider: Arc<dyn ModelProvider>,
role: ModelRole,
) -> Result<Vec<(String, u64)>, StepError>
pub async fn form_memories( &mut self, formation: Formation, source: Tainted<Value>, provider: Arc<dyn ModelProvider>, role: ModelRole, ) -> Result<Vec<(String, u64)>, StepError>
Extract a bounded set of durable facts from labelled source material.
Formation is not an ambient hook. The reviewed declaration supplies the
destination and instruction; the model proposes only stable keys and
content. Every proposal remains labelled from the model and source and
is written through remember.
§It takes the whole role, not a model id
Formation is untrusted contact — the source material derives from
whatever the run handled — so a manifest that declares a quarantined
role declares max_tokens and reasoning_effort beside the model for
exactly this call. Taking the id alone was how those two ceilings got
parsed into the digest and then dropped at this seam: a declared
control the runtime silently did not apply. The role’s ceilings now
ride the formation call itself.
Sourcepub fn blobs(&self) -> Result<Arc<dyn BlobStore>, StepError>
pub fn blobs(&self) -> Result<Arc<dyn BlobStore>, StepError>
The blob store for this run, sealed to its case.
Use this rather than a store held from the builder. With a key ring
configured, bytes written here are encrypted under the case’s data key,
and a store obtained any other way writes them in the clear — the two
would disagree about what erasing the case actually erased. It is also
what a skill passes to
ModelCall::with_media, so
materialization reads through the same envelope that sealed the bytes.
§Errors
If no blob store is configured, or the run belongs to no case while a key ring is — there would be no erasure unit to scope the key to, and falling back to storing in the clear would silently drop the guarantee.
Sourcepub async fn store_blob(&mut self, bytes: &[u8]) -> Result<Digest, StepError>
pub async fn store_blob(&mut self, bytes: &[u8]) -> Result<Digest, StepError>
Store bytes in the blob store and record that this case produced them.
The reason this lives on the context rather than on the blob store: the runtime knows which case is running and the blob store deliberately does not — it is content-addressed, and a digest cannot be reversed to find the matter it belonged to. Writing through here means the association is made at the only moment it is knowable, so an erasure request can later be answered by case, which is the only unit anybody actually asks about. The association is made before the blob write: a crash can leave a harmless dangling link, repaired by retry, but never durable bytes that case erasure cannot discover.
Deliberately not a journaled effect. The digest is a pure function of the bytes, so a replay that re-derives it gets the same answer without re-performing anything, and writing content-addressed bytes twice is the same write. What is journaled is whatever the skill does with the digest next — a tool call carrying it, a case-state write recording it.
§Errors
If no blob store is configured, if the write fails, or if this step is not running inside a case.
Sourcepub async fn fetch_media(
&mut self,
fetcher: &GovernedMedia,
url: Tainted<String>,
) -> Result<Tainted<FetchedMedia>, StepError>
pub async fn fetch_media( &mut self, fetcher: &GovernedMedia, url: Tainted<String>, ) -> Result<Tainted<FetchedMedia>, StepError>
Fetch remote media through the governed, replayable ingestion boundary.
The URL stays labelled and is bound byte-for-byte to the fetch effect.
The fetcher checks and pins DNS, validates every redirect, caps time and
bytes, refuses content coding and ungranted media types, runs configured
validators, and writes the bytes to content-addressed blob storage. The
journal receives only FetchedMedia.
When the run belongs to a case, the digest is linked before blob storage
so erase_case can enforce retention even
across crashes. Strict replay consumes the same clock record but never
rewrites that link.
§Errors
If no blob store is configured, any fetch control refuses the URL or response, validation fails, or blob/case storage fails.
Sourcepub async fn put_case_state(
&mut self,
at: CaseVersion,
state: Value,
) -> Result<CaseVersion, StepError>
pub async fn put_case_state( &mut self, at: CaseVersion, state: Value, ) -> Result<CaseVersion, StepError>
Replace the case’s opaque state, if it is still at at.
A journaled effect, so a replay does not write again.
§Why you have to pass the version
A case is shared by every run correlated to it, and the window between reading its state and writing it back contains a model call — which is unbounded. Two runs on one case overlap as a matter of course, and a blind write in that window silently discards whichever one lost, with nothing in the record to show it happened.
Passing the version you read makes that unexpressible: the store rejects a write against a revision the case has moved past. The remedy is to re-read and decide again — not to retry the same write, which is the lost update this exists to prevent.
§Errors
StepError if this run has no case, or if the case has moved on since
at — see StoreError::CaseConflict.
Sourcepub async fn set_case_status(
&mut self,
status: CaseStatus,
) -> Result<(), StepError>
pub async fn set_case_status( &mut self, status: CaseStatus, ) -> Result<(), StepError>
Move the case to a new status.
Sourcepub async fn deadline(
&mut self,
name: impl Into<String>,
spec: &DeadlineSpec,
warn_before: Option<Duration>,
) -> Result<Deadline, StepError>
pub async fn deadline( &mut self, name: impl Into<String>, spec: &DeadlineSpec, warn_before: Option<Duration>, ) -> Result<Deadline, StepError>
Register a durable obligation on the case.
Resolution goes through the configured Calendar as a journaled
effect, so replay reads back the instant the original run registered
rather than recomputing it against whatever the calendar says today.
That is what keeps a corrected holiday table from retroactively moving a
deadline that has already been relied upon.
§Why warn_before is a std::time::Duration
Two reasons, and the second is the one that bites. It used to be
time::Duration, which is signed — so a negative warning offset
parsed, compiled, and put warn_at after the instant it warns about:
a warning that can only fire once the obligation is already breached.
A quantity that only makes sense non-negative is an unsigned type here,
as it is for Spend.
And it is the Duration a caller already has.
sleep takes the standard one, so the public surface had
two types spelled Duration, only one of which came from a crate this
one re-exports — a reader with the obvious use std::time::Duration
met a type error naming a dependency the guides never mentioned.
Source§impl StepCtx<'_>
Durable waits.
impl StepCtx<'_>
Durable waits.
Sourcepub async fn await_event(
&mut self,
spec: &AwaitSpec,
) -> Result<Tainted<Value>, StepError>
pub async fn await_event( &mut self, spec: &AwaitSpec, ) -> Result<Tainted<Value>, StepError>
Wait for an inbound event correlated by business key.
On replay this returns the event that was recorded; on first execution it either finds one already buffered, or suspends the run.
§The ordering that makes this safe
An event can arrive before the run reaches this call — a fast counterparty, a slow earlier step, a retry that overtakes. So this looks in the durable buffer first, and only registers a subscription and suspends if nothing is there. Delivery and waiting meet in the store rather than in time, which is the only way to close the race.
§Errors
Returns StepError::Suspended when the event has not arrived. That is
not a failure — propagate it with ?. Catching it turns a durable
wait into a silent hang: the subscription stays live, the event arrives
later, and it resumes a run that already decided it was finished.
Sourcepub async fn task(&mut self, spec: &TaskSpec) -> Result<Decision, StepError>
pub async fn task(&mut self, spec: &TaskSpec) -> Result<Decision, StepError>
Ask a human, and wait for the answer.
The task is created, the run suspends, and a decision resumes it. Because the task id is derived from the awaiting effect rather than minted, a resumed run addresses the same task instead of opening a second one for the same decision.
§Errors
Returns StepError::Suspended until somebody decides. Propagate it —
see Self::await_event.
Sourcepub async fn open_task(&mut self, spec: &TaskSpec) -> Result<TaskId, StepError>
pub async fn open_task(&mut self, spec: &TaskSpec) -> Result<TaskId, StepError>
Put something in front of a person without waiting for them.
§The control an advisory agent needs
task asks and blocks. That is the right shape when the
answer decides what happens next, and the wrong one when nothing does: an
agent that has finished, whose finding a compliance desk must see, does
not need its run suspended — it needs a row in a worklist. Gating the
answer to achieve that is a worklist that blocks, and it costs one
suspended run per finding at whatever rate the world produces them.
So this opens the row and returns its id. The run continues, and nothing resumes on the decision because nothing is waiting on it.
§Journaled, and the id is derived
It is an ordinary mutating effect: replay reads the id back rather than
opening a second row, and the id is derived from the effect key so a
resume addresses the row it already opened. TaskStore::open is
idempotent on that id, which is what makes an interrupted attempt safe to
repeat.
§The justification is untrusted, deliberately
What a reviewer is shown usually came from a model, and this does not
route it through the sink gate — the same arrangement task
has always had. Refusing untrusted content at a worklist would mean a
task could only ever carry content nobody needs to review. See
OpenTask for the whole argument.
§Errors
StepError if this run has no case, if no task store is wired, or if
the named obligation is not registered on the case.
Source§impl<'a> StepCtx<'a>
impl<'a> StepCtx<'a>
Sourcepub async fn group<'g, R, S>(
&'g mut self,
name: impl Into<String>,
resources: R,
) -> Result<EffectGroup<'g, 'a>, StepError>
pub async fn group<'g, R, S>( &'g mut self, name: impl Into<String>, resources: R, ) -> Result<EffectGroup<'g, 'a>, StepError>
Open a group over a declared set of resources.
Every member must name a resource from this set. Declaring the footprint up front is what makes the frontier mean something: a group that could touch anything has committed to nothing.
§Errors
If a group is already open — groups do not nest, because a nested abort would have to decide whether it takes the outer group with it, and either answer is wrong half the time.
Trait Implementations§
Auto Trait Implementations§
impl<'a> !RefUnwindSafe for StepCtx<'a>
impl<'a> !UnwindSafe for StepCtx<'a>
impl<'a> Freeze for StepCtx<'a>
impl<'a> Send for StepCtx<'a>
impl<'a> Sync for StepCtx<'a>
impl<'a> Unpin for StepCtx<'a>
impl<'a> UnsafeUnpin for StepCtx<'a>
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> 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