Skip to main content

taquba_workflow/
view.rs

1//! The read-only queries of a workflow store: the status and the outcome of a
2//! run, from the queue's KV namespace and the memo store the runtime writes to.
3
4use taquba::{JobRecord, JobStatus, QueueView};
5
6use crate::durable::{
7    self, DurableCurrentStep, DurableRunRecord, DurableRunResult, DurableTermination,
8};
9use crate::error::{Error, Result};
10use crate::keys::{RunId, outcome_kv_key, run_kv_key, step_kv_key};
11use crate::memo::{MemoStore, RUN_RESULT_MEMO_KEY};
12use crate::runtime::{RunResult, RunState, RunStatus, RunTermination};
13use crate::terminal::RunOutcome;
14
15/// The read-only queries of a workflow store, over a [`QueueView`] and the
16/// [`MemoStore`] the runtime writes to.
17/// [`WorkflowRuntime::view`](crate::WorkflowRuntime::view) returns the
18/// runtime's own view, and a process without a runtime builds a view with
19/// [`WorkflowView::new`] from a [`taquba::QueueReader::view`] and a memo store
20/// at the runtime's prefix.
21#[derive(Clone)]
22pub struct WorkflowView {
23    queue: QueueView,
24    memos: MemoStore,
25}
26
27impl WorkflowView {
28    /// A view over `queue` and `memos`. Through a reader's view the queries are
29    /// of the flushed state the reader last observed.
30    pub fn new(queue: QueueView, memos: MemoStore) -> Self {
31        Self { queue, memos }
32    }
33
34    /// The status of `run_id` from its durable state. An active run is read
35    /// from the run record, the current-step pointer and the step's queue job,
36    /// and a terminated run from its terminal record. A terminated run reports
37    /// [`RunState::Terminated`] until the memo sweep removes its terminal
38    /// record, and a run that is unknown or swept is `None`. A run with a
39    /// pending cancellation request reports [`RunState::Cancelling`] at every
40    /// lifecycle position of its step, until the run terminates.
41    ///
42    /// When a second read returns the current-step pointer with its job still
43    /// absent, the call fails with [`Error::InconsistentRunState`]. The runtime
44    /// does not write a pointer without its job.
45    pub async fn status(&self, run_id: &RunId) -> Result<Option<RunStatus>> {
46        let Some(record) = self.run_record(run_id).await? else {
47            return self.terminated_status(run_id).await;
48        };
49        // The pointer is deleted with the record, so its absence here means
50        // that the run terminated between the two reads.
51        let Some((current, job)) = self.current_job(run_id).await? else {
52            return self.terminated_status(run_id).await;
53        };
54        let state = if record.cancel_requested {
55            RunState::Cancelling
56        } else if job.status == JobStatus::Claimed {
57            RunState::Running
58        } else {
59            RunState::Pending
60        };
61        Ok(Some(RunStatus {
62            run_id: run_id.clone(),
63            state,
64            current_step: current.step_number,
65        }))
66    }
67
68    /// The committed outcome of a terminated run, read from its run result
69    /// record: the result the runner returned or the error that ended the run,
70    /// with the submitter's headers and the final step. The record is written
71    /// by the worker that terminates the run before the terminating settlement
72    /// and is removed with the run's memo entries. It is `None` for a run that
73    /// is unknown, still active or was terminated without a worker (a
74    /// cancellation of a pending step, a dead-letter outside the worker), whose
75    /// status [`Self::status`] reports. A record belongs to the termination its
76    /// terminal record describes: a re-submission of a terminated run id leaves
77    /// the earlier run's record in place until its own termination overwrites
78    /// it, and such a record is not reported.
79    pub async fn outcome(&self, run_id: &RunId) -> Result<Option<RunOutcome>> {
80        if self.current_step_if_active(run_id).await?.is_some() {
81            return Ok(None);
82        }
83        Ok(self
84            .recorded_result(run_id)
85            .await?
86            .map(|result| result.outcome))
87    }
88
89    /// The durable record of `run_id`, when the run is active.
90    pub(crate) async fn run_record(&self, run_id: &RunId) -> Result<Option<DurableRunRecord>> {
91        durable::kv_record(&self.queue, &run_kv_key(run_id)).await
92    }
93
94    /// The current-step pointer of `run_id`, or `None` when the run is not
95    /// active.
96    pub(crate) async fn current_step_if_active(
97        &self,
98        run_id: &RunId,
99    ) -> Result<Option<DurableCurrentStep>> {
100        durable::kv_record(&self.queue, &step_kv_key(run_id)).await
101    }
102
103    /// The current step of `run_id` with its queue job in the stored form of
104    /// [`QueueView::job_record`], or `None` when the run is not active. The
105    /// pointer and the job change in one transaction, so a pointer that moved
106    /// between the two reads is followed. When a second read returns the
107    /// pointer unchanged and its job is still absent, the call fails with
108    /// [`Error::InconsistentRunState`]. The runtime does not write a pointer
109    /// without its job.
110    pub(crate) async fn current_job(
111        &self,
112        run_id: &RunId,
113    ) -> Result<Option<(DurableCurrentStep, JobRecord)>> {
114        let mut absent: Option<String> = None;
115        loop {
116            let Some(current) = self.current_step_if_active(run_id).await? else {
117                return Ok(None);
118            };
119            if let Some(job) = self.queue.job_record(&current.job_id).await? {
120                return Ok(Some((current, job)));
121            }
122            if absent.as_deref() == Some(current.job_id.as_str()) {
123                return Err(Error::InconsistentRunState(run_id.clone()));
124            }
125            absent = Some(current.job_id);
126        }
127    }
128
129    /// The terminal record of `run_id`, or `None` when no record exists.
130    pub(crate) async fn terminal_record(
131        &self,
132        run_id: &RunId,
133    ) -> Result<Option<DurableTermination>> {
134        durable::kv_record(&self.queue, &outcome_kv_key(run_id)).await
135    }
136
137    /// The status of a terminated run from its terminal record, or `None` when
138    /// no record exists.
139    async fn terminated_status(&self, run_id: &RunId) -> Result<Option<RunStatus>> {
140        Ok(self.terminal_record(run_id).await?.map(|record| RunStatus {
141            run_id: run_id.clone(),
142            current_step: record.final_step,
143            state: RunState::Terminated(record.into()),
144        }))
145    }
146
147    /// The run result record of the termination `run_id`'s terminal record
148    /// describes, or `None` when no terminal record remains or the worker that
149    /// terminated the run wrote no record.
150    pub(crate) async fn recorded_result(&self, run_id: &RunId) -> Result<Option<RunResult>> {
151        match self.terminal_record(run_id).await? {
152            Some(termination) => self.run_result_of(run_id, &termination.into()).await,
153            None => Ok(None),
154        }
155    }
156
157    /// The run result record of `run_id` when it belongs to `termination`. A
158    /// record outlives a re-submission of the run id until the next termination
159    /// overwrites it, and a record written before a settlement that did not
160    /// commit outlives the termination that followed, so a record of another
161    /// termination is not reported.
162    pub(crate) async fn run_result_of(
163        &self,
164        run_id: &RunId,
165        termination: &RunTermination,
166    ) -> Result<Option<RunResult>> {
167        Ok(self
168            .run_result(run_id)
169            .await?
170            .filter(|result| result.termination == *termination))
171    }
172
173    /// The run result record of `run_id`, whichever termination it belongs to.
174    /// A record that fails to decode is treated as absent.
175    async fn run_result(&self, run_id: &RunId) -> Result<Option<RunResult>> {
176        let Some(bytes) = self
177            .memos
178            .new_run_memo(run_id)
179            .get(RUN_RESULT_MEMO_KEY)
180            .await?
181        else {
182            return Ok(None);
183        };
184        Ok(
185            durable::decode_or_absent::<DurableRunResult>(&bytes, "run result record", run_id).map(
186                |record| RunResult {
187                    termination: record.termination.into(),
188                    outcome: record.outcome.into(),
189                },
190            ),
191        )
192    }
193}