Skip to main content

WorkflowContext

Struct WorkflowContext 

Source
pub struct WorkflowContext<'a> { /* private fields */ }
Expand description

Replay helper for Rust workflow runtimes.

WorkflowContext is a read-only view over a workflow invocation. It provides deterministic helpers for inspecting persisted history and returning the next command to the engine.

Implementations§

Source§

impl<'a> WorkflowContext<'a>

Source

pub fn new(invocation: &'a WorkflowInvocation) -> Self

Creates a replay context over one immutable runtime invocation.

Source

pub fn run_id(&self) -> &str

Returns the stable run identifier.

Source

pub fn input(&self) -> &JsonValue

Returns the workflow’s initial JSON input.

Source

pub fn spec(&self) -> &WorkflowSpec

Return the immutable workflow definition pinned by run_created.

Source

pub fn has_patch_marker(&self, patch_id: &str) -> bool

Return whether this run was created with a replay-safe patch marker.

Marker presence never changes for an existing run. A compatible runtime can therefore keep both code paths and deterministically replay old unmarked histories alongside new marked histories.

Source

pub fn input_as<T>(&self) -> Result<T>

Decode the workflow input into a host-defined serde type.

Source

pub fn history(&self) -> &[FlowEventEnvelope]

Returns committed history in ascending event-sequence order.

Source

pub fn cancellation_request(&self) -> Option<&CancellationRequest>

Return the durable cleanup-aware cancellation request, when present.

Source

pub fn progress(&self, progress_id: &str) -> Option<&WorkflowProgress>

Return a durable progress update by its idempotency identity.

Source

pub fn child_operation( &self, reference_id: &str, ) -> Option<&ChildOperationReference>

Return a durable child-operation reference by its parent-local id.

Source

pub fn child_workflow_run_id(&self, child_id: &str) -> Option<&str>

Return the engine-generated root run ID for a durable child request.

Source

pub fn child_workflow_outcome( &self, child_id: &str, ) -> Option<&WorkflowTerminalOutcome>

Return the terminal outcome durably observed for a child workflow.

Source

pub fn signal(&self, signal_id: &str) -> Option<&WorkflowSignal>

Return a received signal by its caller-owned delivery identity.

Source

pub fn signal_payload(&self, wait_id: &str) -> Option<&JsonValue>

Return the payload paired with a completed deterministic signal wait.

Source

pub fn signal_payload_as<T>(&self, wait_id: &str) -> Result<Option<T>>

Decode the payload paired with a completed deterministic signal wait.

Source

pub fn step_output(&self, step_id: &str) -> Option<&JsonValue>

Returns the durable JSON output of a completed step.

Source

pub fn step_output_as<T>(&self, step_id: &str) -> Result<Option<T>>

Decodes a completed step output into a host-defined serde type.

Source

pub fn step_completed(&self, step_id: &str) -> bool

Returns whether the step has a durable successful output.

Source

pub fn step_failed(&self, step_id: &str) -> Option<&str>

Returns the terminal error of a step that exhausted its retries.

Source

pub fn wait_completed(&self, wait_id: &str) -> bool

Returns whether a durable timer wait has completed.

Source

pub fn hook_payload(&self, hook_id: &str) -> Option<&JsonValue>

Returns the durable JSON payload received by a hook.

Source

pub fn hook_payload_as<T>(&self, hook_id: &str) -> Result<Option<T>>

Decodes a received hook payload into a host-defined serde type.

Source

pub fn hook_disposed(&self, hook_id: &str) -> bool

Returns whether a hook was explicitly closed without a payload.

Source

pub fn complete(&self, output: JsonValue) -> RuntimeCommand

Returns a command that completes the workflow successfully.

Source

pub fn fail(&self, error: impl Into<String>) -> RuntimeCommand

Returns a command that fails the workflow.

Source

pub fn cancel(&self) -> RuntimeCommand

Finish a previously requested cancellation after cleanup is durable.

Source

pub fn timeout( &self, deadline: DateTime<Utc>, reason: Option<String>, ) -> RuntimeCommand

Finish a run with a typed timeout outcome.

Source

pub fn continue_as_new(&self, input: JsonValue) -> RuntimeCommand

Close this history segment and continue with fresh history and input.

The engine persists the successor identity before creating it and carries the exact current WorkflowSpec into the new run.

Source

pub fn record_progress(&self, progress: WorkflowProgress) -> RuntimeCommand

Persist an idempotently identified progress update and replay.

Persist a child-operation reference and replay.

Source

pub fn start_child_workflow( &self, child_id: impl Into<String>, spec: WorkflowSpec, input: JsonValue, ) -> RuntimeCommand

Start or await a first-class child workflow.

The child ID is stable within this parent history. By default, a parent cancellation request is propagated to an open child and the parent waits for the child’s terminal outcome.

Source

pub fn start_child_workflow_with_policy( &self, child_id: impl Into<String>, spec: WorkflowSpec, input: JsonValue, cancellation_policy: ChildWorkflowCancellationPolicy, ) -> RuntimeCommand

Start or await a child with an explicit cancellation policy.

Source

pub fn child_workflow( &self, child_id: impl Into<String>, spec: WorkflowSpec, input: JsonValue, ) -> ChildWorkflowCommand

Create a child definition for a bounded durable batch.

Source

pub fn child_workflow_with_policy( &self, child_id: impl Into<String>, spec: WorkflowSpec, input: JsonValue, cancellation_policy: ChildWorkflowCancellationPolicy, ) -> ChildWorkflowCommand

Create a batch child definition with an explicit cancellation policy.

Source

pub fn start_child_workflows( &self, children: Vec<ChildWorkflowCommand>, ) -> RuntimeCommand

Durably request a deterministic batch before any child starts.

Source

pub fn schedule_step( &self, step_id: impl Into<String>, step_name: impl Into<String>, input: JsonValue, ) -> RuntimeCommand

Schedules one durable step with the default retry policy.

Source

pub fn schedule_step_with_retry( &self, step_id: impl Into<String>, step_name: impl Into<String>, input: JsonValue, retry: RetryPolicy, ) -> RuntimeCommand

Schedules one durable step with an explicit retry policy.

Source

pub fn step( &self, step_id: impl Into<String>, step_name: impl Into<String>, input: JsonValue, ) -> StepCommand

Creates a step definition with the default retry policy.

Source

pub fn step_with_retry( &self, step_id: impl Into<String>, step_name: impl Into<String>, input: JsonValue, retry: RetryPolicy, ) -> StepCommand

Creates a step definition with an explicit retry policy.

Source

pub fn schedule_steps(&self, steps: Vec<StepCommand>) -> RuntimeCommand

Atomically schedules a deterministic batch of durable steps.

Source

pub fn wait_until( &self, wait_id: impl Into<String>, resume_at: DateTime<Utc>, ) -> RuntimeCommand

Suspends replay until the given UTC deadline becomes ready.

Source

pub fn create_hook( &self, hook_id: impl Into<String>, token: impl Into<String>, metadata: JsonValue, ) -> RuntimeCommand

Creates an externally completable hook with JSON metadata.

Source

pub fn create_hook_with_metadata( &self, hook_id: impl Into<String>, token: impl Into<String>, metadata: HookMetadata, ) -> Result<RuntimeCommand>

Creates an externally completable hook with typed metadata.

Source

pub fn wait_for_signal( &self, wait_id: impl Into<String>, signal_name: impl Into<String>, ) -> RuntimeCommand

Suspend until the next queued signal with signal_name is paired with the stable wait_id.

Auto Trait Implementations§

§

impl<'a> Freeze for WorkflowContext<'a>

§

impl<'a> RefUnwindSafe for WorkflowContext<'a>

§

impl<'a> Send for WorkflowContext<'a>

§

impl<'a> Sync for WorkflowContext<'a>

§

impl<'a> Unpin for WorkflowContext<'a>

§

impl<'a> UnsafeUnpin for WorkflowContext<'a>

§

impl<'a> UnwindSafe for WorkflowContext<'a>

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<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

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> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T> SqlComparable<Option<T>> for T

Source§

impl<T> SqlComparable<T> for T

Source§

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

Source§

type Error = !

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<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