taquba_workflow/jobs/handle.rs
1use std::future::{Future, IntoFuture};
2use std::marker::PhantomData;
3use std::pin::Pin;
4use std::time::Duration;
5
6use crate::{RunId, RunOutcome, RunStatus, RunTermination, StepErrorKind, TerminalStatus};
7use thiserror::Error;
8
9use crate::jobs::job::Job;
10use crate::jobs::runner::JobRuntime;
11use crate::{Error, Result};
12
13/// The logical failure outcome of a job that did not succeed.
14///
15/// Distinct from [`Error`](enum@Error), which is an infrastructure
16/// failure: a `JobError` means the job terminated unsuccessfully. The
17/// concrete `Job::Error` value is not persisted, so this holds its
18/// [`Display`](std::fmt::Display) message and classification.
19#[derive(Debug, Clone, Error)]
20#[error("job failed ({kind:?}): {message}")]
21pub struct JobError {
22 /// Whether the failure was classified transient (the job exhausted its
23 /// attempts) or permanent (dead-lettered on the failing attempt).
24 /// Transient for a job cancelled or terminated outside its handler.
25 pub kind: StepErrorKind,
26 /// The failure message.
27 pub message: String,
28}
29
30/// The error produced by awaiting a [`JobHandle`] directly (via `.await`).
31///
32/// Flattens the two failure modes (infrastructure errors and the job's own
33/// logical failure) into one type so `handle.await?` yields the job's
34/// `Output` directly.
35#[derive(Debug, Error)]
36pub enum JoinError {
37 /// An infrastructure error occurred while submitting, waiting or reading
38 /// the outcome.
39 #[error(transparent)]
40 Infra(#[from] Error),
41 /// The job ran to a terminal state but did not succeed.
42 #[error(transparent)]
43 Job(#[from] JobError),
44}
45
46/// A handle to a submitted job.
47///
48/// Returned by [`JobRunner::submit`](crate::jobs::JobRunner::submit). Await it
49/// directly for the typed result, or use [`join`](Self::join),
50/// [`fetch_result`](Self::fetch_result) and [`status`](Self::status) for
51/// more control.
52///
53/// Awaiting is [`WorkflowRuntime::wait`](crate::WorkflowRuntime::wait)
54/// on the job's run, which relies on Taquba's in-process completion
55/// notification, so a handle is awaited in the same process that runs
56/// the job. The outcome is durable regardless:
57/// [`fetch_result`](Self::fetch_result) reads it back from object
58/// storage after a restart.
59pub struct JobHandle<J: Job> {
60 id: RunId,
61 runtime: JobRuntime,
62 newly_submitted: bool,
63 _marker: PhantomData<fn() -> J>,
64}
65
66impl<J: Job> Clone for JobHandle<J> {
67 fn clone(&self) -> Self {
68 Self {
69 id: self.id.clone(),
70 runtime: self.runtime.clone(),
71 newly_submitted: self.newly_submitted,
72 _marker: PhantomData,
73 }
74 }
75}
76
77impl<J: Job> JobHandle<J> {
78 pub(crate) fn new(id: RunId, runtime: JobRuntime, newly_submitted: bool) -> Self {
79 Self {
80 id,
81 runtime,
82 newly_submitted,
83 _marker: PhantomData,
84 }
85 }
86
87 /// The job's identifier: a ULID, or the digest of the job's
88 /// [`idempotency_key`](Job::idempotency_key) when it has one, so a
89 /// submission that matched an earlier job returns that job's id.
90 pub fn id(&self) -> &RunId {
91 &self.id
92 }
93
94 /// True if the call that produced this handle submitted a new job;
95 /// false if the call matched an in-flight or completed submission with
96 /// the same [`Job::idempotency_key`](crate::jobs::Job::idempotency_key) and
97 /// payload.
98 ///
99 /// For submissions without an `idempotency_key`, the value is always
100 /// `true`. The value reflects the call that returned this handle and
101 /// does not update as the job progresses; clones preserve it.
102 pub fn newly_submitted(&self) -> bool {
103 self.newly_submitted
104 }
105
106 /// The job's status, read from its durable state. A terminated job
107 /// reports [`RunState::Terminated`](crate::RunState::Terminated)
108 /// until [`JobRunnerBuilder::retention`](crate::jobs::JobRunnerBuilder::retention)
109 /// removes its terminal record. Use
110 /// [`fetch_result`](Self::fetch_result) to read a terminal outcome.
111 pub async fn status(&self) -> Result<Option<RunStatus>> {
112 self.runtime.status(&self.id).await
113 }
114
115 /// Read the job's persisted result without waiting.
116 ///
117 /// Returns `None` when no run result record exists for this job: it
118 /// is still pending or in flight, it terminated without a worker
119 /// recording a result (a lease expiry dead-lettered it, or it was
120 /// cancelled while pending), or the record was removed by retention.
121 ///
122 /// Reads from object storage, so it works across process restarts.
123 pub async fn fetch_result(&self) -> Result<Option<std::result::Result<J::Output, JobError>>> {
124 match self
125 .runtime
126 .inner
127 .core
128 .view
129 .recorded_result(&self.id)
130 .await?
131 {
132 None => Ok(None),
133 Some(result) => decode_end::<J>(result.termination, Some(result.outcome)).map(Some),
134 }
135 }
136
137 /// Wait for the job to reach a terminal state and return its outcome.
138 ///
139 /// Waits indefinitely. Use [`join_timeout`](Self::join_timeout) to bound
140 /// the wait.
141 pub async fn join(&self) -> Result<std::result::Result<J::Output, JobError>> {
142 let end = self.runtime.wait(&self.id).await?;
143 decode_end::<J>(end.termination, end.outcome)
144 }
145
146 /// Wait up to `timeout` for the job to reach a terminal state.
147 ///
148 /// Returns `Ok(None)` if the timeout elapses first. On completion
149 /// the run result record is decoded; a job that reached a terminal
150 /// state without one (a lease expiry dead-lettered it, or it was
151 /// cancelled) is reported as a transient [`JobError`] from the
152 /// termination the runtime retains, or with a generic message when
153 /// it retains none.
154 ///
155 /// Returns [`Error::RunNotFound`] if the runtime has no record of
156 /// the job.
157 pub async fn join_timeout(
158 &self,
159 timeout: Duration,
160 ) -> Result<Option<std::result::Result<J::Output, JobError>>> {
161 match self.runtime.wait_timeout(&self.id, timeout).await? {
162 None => Ok(None),
163 Some(end) => decode_end::<J>(end.termination, end.outcome).map(Some),
164 }
165 }
166}
167
168/// The typed result of a terminated job: the decoded output of a
169/// succeeded job, or the [`JobError`] of one that failed, was cancelled
170/// or terminated without recording an outcome.
171pub(crate) fn decode_end<J: Job>(
172 termination: RunTermination,
173 outcome: Option<RunOutcome>,
174) -> Result<std::result::Result<J::Output, JobError>> {
175 if let Some(outcome) = outcome
176 && outcome.status == TerminalStatus::Succeeded
177 {
178 let output = outcome.result.unwrap_or_default();
179 return Ok(Ok(rmp_serde::from_slice(&output)?));
180 }
181 let message = termination.error.unwrap_or_else(|| {
182 match termination.status {
183 TerminalStatus::Cancelled => "job cancelled",
184 _ => "job terminated without recording an outcome",
185 }
186 .to_string()
187 });
188 Ok(Err(JobError {
189 kind: termination.error_kind.unwrap_or(StepErrorKind::Transient),
190 message,
191 }))
192}
193
194impl<J: Job> IntoFuture for JobHandle<J> {
195 type Output = std::result::Result<J::Output, JoinError>;
196 type IntoFuture = Pin<Box<dyn Future<Output = Self::Output> + Send>>;
197
198 fn into_future(self) -> Self::IntoFuture {
199 Box::pin(async move {
200 match self.join().await {
201 Ok(Ok(output)) => Ok(output),
202 Ok(Err(job_error)) => Err(JoinError::Job(job_error)),
203 Err(infra) => Err(JoinError::Infra(infra)),
204 }
205 })
206 }
207}