pub struct WorkflowRuntime<R, H> { /* private fields */ }Expand description
Durable runtime for workflow runs. Cheap to clone (internally Arc).
Implementations§
Source§impl<R: StepRunner, H: TerminalHook> WorkflowRuntime<R, H>
impl<R: StepRunner, H: TerminalHook> WorkflowRuntime<R, H>
Sourcepub fn builder(
queue: Arc<Queue>,
object_store: Arc<dyn ObjectStore>,
runner: R,
terminal_hook: H,
) -> WorkflowRuntimeBuilder<R, H>
pub fn builder( queue: Arc<Queue>, object_store: Arc<dyn ObjectStore>, runner: R, terminal_hook: H, ) -> WorkflowRuntimeBuilder<R, H>
Start configuring a runtime. Takes the four required dependencies
(Taquba queue, object store, StepRunner, TerminalHook); optional
fields are set via WorkflowRuntimeBuilder methods before build.
The object store backs Delivery::memo; it does not need to be
the same store the Queue was opened with, though sharing one store
is the common case (just clone the Arc). Use a distinct
WorkflowRuntimeBuilder::memo_prefix when multiple runtimes share one
store.
Use crate::NoopTerminalHook if you don’t need terminal callbacks.
Sourcepub async fn submit(&self, spec: RunSpec) -> Result<SubmitOutcome>
pub async fn submit(&self, spec: RunSpec) -> Result<SubmitOutcome>
Submit a new run. Enqueues step 0 with payload spec.input.
Idempotent on (run_id, spec.input): if a run with the same id is
already active (its durable run record in Taquba’s user KV
namespace exists, whichever process submitted it) and
spec.input matches the original submission, this
call is a no-op and the returned SubmitOutcome has
newly_submitted = false. A re-submission of an active run_id
with a different input is rejected with Error::InputMismatch;
pick a fresh run_id for a new run.
Sourcepub fn view(&self) -> &WorkflowView
pub fn view(&self) -> &WorkflowView
The runtime’s view of its store: the read-only queries over the queue and the memo store the runtime writes to.
Sourcepub async fn status(&self, run_id: &RunId) -> Result<Option<RunStatus>>
pub async fn status(&self, run_id: &RunId) -> Result<Option<RunStatus>>
WorkflowView::status through the runtime’s view, so the status is
available after a restart and from any runtime over the same queue.
Sourcepub async fn outcome(&self, run_id: &RunId) -> Result<Option<RunOutcome>>
pub async fn outcome(&self, run_id: &RunId) -> Result<Option<RunOutcome>>
WorkflowView::outcome through the runtime’s view.
Sourcepub async fn wait(&self, run_id: &RunId) -> Result<RunEnd>
pub async fn wait(&self, run_id: &RunId) -> Result<RunEnd>
Wait until the run run_id terminates and report its end. The
wait follows the run’s current step across its steps, so it
answers for a run of any length, after a restart and from any
runtime over the same queue; a run already terminated is
reported at once from its records. A step dead-lettered outside
the worker is terminated by the worker’s dead-step
reconciliation, which the wait polls for at the poll interval.
Returns Error::RunNotFound for a run the runtime has no
record of: never submitted, or terminated and swept.
Sourcepub async fn wait_timeout(
&self,
run_id: &RunId,
timeout: Duration,
) -> Result<Option<RunEnd>>
pub async fn wait_timeout( &self, run_id: &RunId, timeout: Duration, ) -> Result<Option<RunEnd>>
Self::wait bounded by timeout; Ok(None) when the timeout
elapses first.
Sourcepub async fn cancel(&self, run_id: &RunId) -> Result<bool>
pub async fn cancel(&self, run_id: &RunId) -> Result<bool>
Request cancellation of an active run.
Returns Ok(true) once the request is recorded on the run’s
durable record, or Ok(false) if the run is unknown or already
terminal, including a run whose current step the queue
dead-lettered outside the worker, which the worker’s dead-step
reconciliation terminates as failed. The request reaches a run
after a restart and from any runtime over the same queue.
The run terminates as TerminalStatus::Cancelled and its
notification job is enqueued for the terminal hook:
- Pending / scheduled step: the queued step job is removed and the notification enqueued in one transaction before this call returns; the hook runs from a worker afterwards.
- Running step: cancellation is delivered to the runner via
Delivery::cancel_token; runners that watch the token short-circuit immediately. Runners that ignore the token are allowed to run to completion (futures cannot be safely aborted mid-step). In both cases the runner’sStepOutcome/StepErroris discarded and the worker settles the run once the step returns, with any pending transient retry suppressed and the step acked rather than nacked. - A step claimed after the request is settled as cancelled without running.
Cancellation is best-effort: a run whose terminal step settles while the request is recorded keeps the outcome it committed. A request applies only to the run it is recorded on, and a later submission of the same run id starts without it.
Sourcepub fn spawn<F>(&self, shutdown: F) -> RunnerHandle
pub fn spawn<F>(&self, shutdown: F) -> RunnerHandle
Spawn Self::run as a Tokio task and return a handle for
graceful shutdown. The worker runs until shutdown resolves or
RunnerHandle::shutdown is called; in-flight steps finish
either way.
Sourcepub async fn run<F>(&self, shutdown: F) -> Result<()>
pub async fn run<F>(&self, shutdown: F) -> Result<()>
Drive the step worker loop until shutdown resolves. Spawns up
to max_concurrent_steps step processors, the dead-step
reconciliation that terminates runs whose step the queue
dead-lettered outside the worker and, when
WorkflowRuntimeBuilder::memo_retention is set, a
memo-retention sweeper, all running in parallel. All halt cleanly
when shutdown resolves or the worker errors.
Source§impl<R: StepRunner, H: TerminalHook> WorkflowRuntime<R, H>
impl<R: StepRunner, H: TerminalHook> WorkflowRuntime<R, H>
Sourcepub async fn signal(
&self,
correlation_key: &str,
payload: Vec<u8>,
) -> Result<SignalOutcome>
pub async fn signal( &self, correlation_key: &str, payload: Vec<u8>, ) -> Result<SignalOutcome>
Deliver a signal for correlation_key, waking the run waiting on
it via Trigger::OnSignal.
When a waiter is registered and still waiting, its next step is
promoted immediately and observes payload on Step::signal;
the call returns SignalOutcome::Delivered. Otherwise the signal
is buffered durably under the correlation key and the next waiter
registered for it consumes the buffered payload at its
registration; the call returns SignalOutcome::Buffered.
One buffered signal is held per correlation key: a second signal
before consumption replaces the first. A buffered signal persists
until a waiter consumes it or Self::clear_signal discards it.
The buffer write is durable before the call returns, so a signal
is never lost once this call returns; delivery to a waiter whose
registration is settling concurrently falls back to the buffer and
reaches it no later than its timeout.
Sourcepub async fn clear_signal(&self, correlation_key: &str) -> Result<bool>
pub async fn clear_signal(&self, correlation_key: &str) -> Result<bool>
Discard the buffered signal for correlation_key, if one exists.
Returns true when a buffered signal was removed.
Trait Implementations§
Auto Trait Implementations§
impl<R, H> !RefUnwindSafe for WorkflowRuntime<R, H>
impl<R, H> !UnwindSafe for WorkflowRuntime<R, H>
impl<R, H> Freeze for WorkflowRuntime<R, H>
impl<R, H> Send for WorkflowRuntime<R, H>
impl<R, H> Sync for WorkflowRuntime<R, H>
impl<R, H> Unpin for WorkflowRuntime<R, H>
impl<R, H> UnsafeUnpin for WorkflowRuntime<R, H>where
Arc<RuntimeInner<R, H>>: UnsafeUnpin,
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
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 moreSource§impl<T> Paint for Twhere
T: ?Sized,
impl<T> Paint for Twhere
T: ?Sized,
Source§fn fg(&self, value: Color) -> Painted<&T>
fn fg(&self, value: Color) -> Painted<&T>
Returns a styled value derived from self with the foreground set to
value.
This method should be used rarely. Instead, prefer to use color-specific
builder methods like red() and
green(), which have the same functionality but are
pithier.
§Example
Set foreground color to white using fg():
use yansi::{Paint, Color};
painted.fg(Color::White);Set foreground color to white using white().
use yansi::Paint;
painted.white();Source§fn bright_black(&self) -> Painted<&T>
fn bright_black(&self) -> Painted<&T>
Source§fn bright_red(&self) -> Painted<&T>
fn bright_red(&self) -> Painted<&T>
Source§fn bright_green(&self) -> Painted<&T>
fn bright_green(&self) -> Painted<&T>
Source§fn bright_yellow(&self) -> Painted<&T>
fn bright_yellow(&self) -> Painted<&T>
Source§fn bright_blue(&self) -> Painted<&T>
fn bright_blue(&self) -> Painted<&T>
Source§fn bright_magenta(&self) -> Painted<&T>
fn bright_magenta(&self) -> Painted<&T>
Source§fn bright_cyan(&self) -> Painted<&T>
fn bright_cyan(&self) -> Painted<&T>
Source§fn bright_white(&self) -> Painted<&T>
fn bright_white(&self) -> Painted<&T>
Source§fn bg(&self, value: Color) -> Painted<&T>
fn bg(&self, value: Color) -> Painted<&T>
Returns a styled value derived from self with the background set to
value.
This method should be used rarely. Instead, prefer to use color-specific
builder methods like on_red() and
on_green(), which have the same functionality but
are pithier.
§Example
Set background color to red using fg():
use yansi::{Paint, Color};
painted.bg(Color::Red);Set background color to red using on_red().
use yansi::Paint;
painted.on_red();Source§fn on_primary(&self) -> Painted<&T>
fn on_primary(&self) -> Painted<&T>
Source§fn on_magenta(&self) -> Painted<&T>
fn on_magenta(&self) -> Painted<&T>
Source§fn on_bright_black(&self) -> Painted<&T>
fn on_bright_black(&self) -> Painted<&T>
Source§fn on_bright_red(&self) -> Painted<&T>
fn on_bright_red(&self) -> Painted<&T>
Source§fn on_bright_green(&self) -> Painted<&T>
fn on_bright_green(&self) -> Painted<&T>
Source§fn on_bright_yellow(&self) -> Painted<&T>
fn on_bright_yellow(&self) -> Painted<&T>
Source§fn on_bright_blue(&self) -> Painted<&T>
fn on_bright_blue(&self) -> Painted<&T>
Source§fn on_bright_magenta(&self) -> Painted<&T>
fn on_bright_magenta(&self) -> Painted<&T>
Source§fn on_bright_cyan(&self) -> Painted<&T>
fn on_bright_cyan(&self) -> Painted<&T>
Source§fn on_bright_white(&self) -> Painted<&T>
fn on_bright_white(&self) -> Painted<&T>
Source§fn attr(&self, value: Attribute) -> Painted<&T>
fn attr(&self, value: Attribute) -> Painted<&T>
Enables the styling Attribute value.
This method should be used rarely. Instead, prefer to use
attribute-specific builder methods like bold() and
underline(), which have the same functionality
but are pithier.
§Example
Make text bold using attr():
use yansi::{Paint, Attribute};
painted.attr(Attribute::Bold);Make text bold using using bold().
use yansi::Paint;
painted.bold();Source§fn rapid_blink(&self) -> Painted<&T>
fn rapid_blink(&self) -> Painted<&T>
Source§fn quirk(&self, value: Quirk) -> Painted<&T>
fn quirk(&self, value: Quirk) -> Painted<&T>
Enables the yansi Quirk value.
This method should be used rarely. Instead, prefer to use quirk-specific
builder methods like mask() and
wrap(), which have the same functionality but are
pithier.
§Example
Enable wrapping using .quirk():
use yansi::{Paint, Quirk};
painted.quirk(Quirk::Wrap);Enable wrapping using wrap().
use yansi::Paint;
painted.wrap();Source§fn clear(&self) -> Painted<&T>
👎Deprecated since 1.0.1: renamed to resetting() due to conflicts with Vec::clear().
The clear() method will be removed in a future release.
fn clear(&self) -> Painted<&T>
renamed to resetting() due to conflicts with Vec::clear().
The clear() method will be removed in a future release.
Source§fn whenever(&self, value: Condition) -> Painted<&T>
fn whenever(&self, value: Condition) -> Painted<&T>
Conditionally enable styling based on whether the Condition value
applies. Replaces any previous condition.
See the crate level docs for more details.
§Example
Enable styling painted only when both stdout and stderr are TTYs:
use yansi::{Paint, Condition};
painted.red().on_yellow().whenever(Condition::STDOUTERR_ARE_TTY);