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(¤t.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}