Skip to main content

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}