taquba_workflow/signal.rs
1//! Durable signal delivery: one buffered signal per correlation key,
2//! held in three key spaces of the caller KV namespace. The waiter
3//! index maps a correlation key to the step job scheduled to wait on
4//! it, the buffer holds a signal that arrived while no waiter was
5//! registered and the delivered record stores a consumed payload under
6//! `(run id, step)` so a redelivered step observes the same signal.
7//! Registration consumes a buffered signal at the previous step's
8//! settlement; delivery to a registered waiter wakes its scheduled
9//! job; a waiter promoted by its timeout consumes the buffer at claim
10//! time.
11
12use std::collections::HashMap;
13use std::time::Duration;
14
15use taquba::{JobRecord, JobStatus, Queue, SettlementEffects, WorkerError};
16use tracing::{debug, warn};
17
18use crate::error::{Result, worker_error};
19use crate::keys::{
20 HEADER_SIGNAL_DELIVERED, HEADER_SIGNAL_WAIT, RunId, signal_buf_kv_key, signal_delivered_kv_key,
21 signal_wait_kv_key,
22};
23use crate::runner::{StepError, StepRunner};
24use crate::runtime::{RuntimeCore, RuntimeInner, StepEnqueueOpts, WorkflowRuntime};
25use crate::terminal::TerminalHook;
26use crate::worker::ClaimedStep;
27
28/// Outcome of [`WorkflowRuntime::signal`].
29#[derive(Debug, Clone, Copy, PartialEq, Eq)]
30pub enum SignalOutcome {
31 /// A waiter was registered for the correlation key and has been woken;
32 /// its step runs next with [`Step::signal`](crate::Step::signal) set to the payload.
33 Delivered,
34 /// No waiter was woken. The signal is buffered durably under the
35 /// correlation key and is consumed by the next waiter registered for
36 /// it, or discarded by [`WorkflowRuntime::clear_signal`].
37 Buffered,
38}
39
40/// Number of waiter-index reads [`WorkflowRuntime::signal`] performs
41/// before concluding no waiter exists, and the pause between them. A
42/// registration settling concurrently becomes visible within one
43/// acknowledgement commit, which these reads cover.
44const SIGNAL_WAIT_READ_ATTEMPTS: u32 = 10;
45const SIGNAL_WAIT_READ_INTERVAL: Duration = Duration::from_millis(25);
46
47/// Remove the signal entry at `key` if it still holds `expected`. A
48/// failed removal is logged and otherwise ignored: the entry is
49/// residue that the next resolution of its key overwrites or drops.
50async fn remove_entry(queue: &Queue, key: &[u8], expected: &[u8]) {
51 if let Err(err) = queue.kv_compare_delete(key, expected).await {
52 debug!(key = %String::from_utf8_lossy(key), "signal entry removal failed: {err}");
53 }
54}
55
56impl<R: StepRunner, H: TerminalHook> WorkflowRuntime<R, H> {
57 /// Deliver a signal for `correlation_key`, waking the run waiting on
58 /// it via [`Trigger::OnSignal`](crate::Trigger::OnSignal).
59 ///
60 /// When a waiter is registered and still waiting, its next step is
61 /// promoted immediately and observes `payload` on [`Step::signal`](crate::Step::signal);
62 /// the call returns [`SignalOutcome::Delivered`]. Otherwise the signal
63 /// is buffered durably under the correlation key and the next waiter
64 /// registered for it consumes the buffered payload at its
65 /// registration; the call returns [`SignalOutcome::Buffered`].
66 ///
67 /// One buffered signal is held per correlation key: a second signal
68 /// before consumption replaces the first. A buffered signal persists
69 /// until a waiter consumes it or [`Self::clear_signal`] discards it.
70 /// The buffer write is durable before the call returns, so a signal
71 /// is never lost once this call returns; delivery to a waiter whose
72 /// registration is settling concurrently falls back to the buffer and
73 /// reaches it no later than its timeout.
74 pub async fn signal(&self, correlation_key: &str, payload: Vec<u8>) -> Result<SignalOutcome> {
75 let queue = &self.inner.core.queue;
76 let buf_key = signal_buf_kv_key(correlation_key);
77 // Buffer first, durably: a waiter registering concurrently reads
78 // the buffer at its settlement, so the signal is never lost even
79 // if the waiter index is not yet visible below.
80 queue.kv_put(&buf_key, &payload).await?;
81
82 let wait_key = signal_wait_kv_key(correlation_key);
83 for attempt in 0..SIGNAL_WAIT_READ_ATTEMPTS {
84 if attempt > 0 {
85 tokio::time::sleep(SIGNAL_WAIT_READ_INTERVAL).await;
86 }
87 let Some(waiter) = queue.view().kv_get(&wait_key).await? else {
88 continue;
89 };
90 let Ok(job_id) = std::str::from_utf8(&waiter).map(str::to_string) else {
91 remove_entry(queue, &wait_key, &waiter).await;
92 return Ok(SignalOutcome::Buffered);
93 };
94 return match queue.wake_scheduled(&job_id, Some(payload.clone())).await? {
95 taquba::WakeOutcome::Woken => {
96 remove_entry(queue, &buf_key, &payload).await;
97 remove_entry(queue, &wait_key, &waiter).await;
98 Ok(SignalOutcome::Delivered)
99 }
100 taquba::WakeOutcome::NotScheduled | taquba::WakeOutcome::NotFound => {
101 // Stale index entry: the waiter was already promoted or
102 // its run was cancelled. Remove the entry; the signal
103 // stays buffered.
104 remove_entry(queue, &wait_key, &waiter).await;
105 Ok(SignalOutcome::Buffered)
106 }
107 };
108 }
109 Ok(SignalOutcome::Buffered)
110 }
111
112 /// Discard the buffered signal for `correlation_key`, if one exists.
113 /// Returns `true` when a buffered signal was removed.
114 pub async fn clear_signal(&self, correlation_key: &str) -> Result<bool> {
115 let queue = &self.inner.core.queue;
116 let buf_key = signal_buf_kv_key(correlation_key);
117 loop {
118 let Some(current) = queue.view().kv_get(&buf_key).await? else {
119 return Ok(false);
120 };
121 if queue.kv_compare_delete(&buf_key, ¤t).await? {
122 return Ok(true);
123 }
124 }
125 }
126}
127
128impl<R: StepRunner, H: TerminalHook> RuntimeInner<R, H> {
129 /// Build the effects that advance the run of `claimed` to a step
130 /// that waits for a signal for `correlation_key`. When a buffered
131 /// signal already exists it is consumed: the next step is enqueued
132 /// immediately, the payload is recorded under the durable delivered
133 /// key and the buffer entry is deleted, all in the acknowledgement
134 /// transaction. Otherwise the next step is scheduled `timeout` from
135 /// now and the waiter index entry joins the same transaction.
136 pub(crate) async fn advance_on_signal(
137 &self,
138 claimed: &ClaimedStep<'_>,
139 payload: Vec<u8>,
140 correlation_key: &str,
141 timeout: Duration,
142 input_hash: [u8; 32],
143 ) -> std::result::Result<SettlementEffects, WorkerError> {
144 let wait_key = signal_wait_kv_key(correlation_key);
145
146 // One waiter per correlation key: a registration while a live
147 // waiter is at the key is rejected. The check reads current
148 // state, and a stale index entry (its job no longer scheduled) is
149 // overwritten. The rejection terminates the run as a permanent
150 // step error does.
151 if let Some(existing) = self
152 .core
153 .queue
154 .view()
155 .kv_get(&wait_key)
156 .await
157 .map_err(worker_error)?
158 && let Ok(existing_id) = std::str::from_utf8(&existing)
159 && let Ok(Some(job)) = self.core.queue.view().job_record(existing_id).await
160 && job.status == JobStatus::Scheduled
161 {
162 let message =
163 format!("a waiter is already registered for correlation key `{correlation_key}`");
164 return Err(self
165 .terminating_failure(claimed, StepError::permanent(message), input_hash)
166 .await);
167 }
168
169 let buf_key = signal_buf_kv_key(correlation_key);
170 match self
171 .core
172 .queue
173 .view()
174 .kv_get(&buf_key)
175 .await
176 .map_err(worker_error)?
177 {
178 Some(buffered) => {
179 let opts = StepEnqueueOpts {
180 reserved_headers: claimed
181 .reserved_headers_with((HEADER_SIGNAL_DELIVERED, "1".to_string())),
182 ..claimed.next_step_opts()
183 };
184 let delivered_key =
185 signal_delivered_kv_key(&claimed.run_id, claimed.step_number + 1);
186 let buffered = buffered.to_vec();
187 let mut effects = self
188 .core
189 .advance_with_kv(claimed, payload, opts, |_| {
190 HashMap::from([(delivered_key, buffered)])
191 })
192 .await;
193 effects.kv_deletes.push(buf_key);
194 Ok(effects)
195 }
196 None => {
197 let opts = StepEnqueueOpts {
198 run_at: Some(self.core.run_at_after(timeout)),
199 reserved_headers: claimed
200 .reserved_headers_with((HEADER_SIGNAL_WAIT, correlation_key.to_string())),
201 ..claimed.next_step_opts()
202 };
203 let effects = self
204 .core
205 .advance_with_kv(claimed, payload, opts, |job_id| {
206 HashMap::from([(wait_key, job_id.as_bytes().to_vec())])
207 })
208 .await;
209 Ok(effects)
210 }
211 }
212 }
213}
214
215impl RuntimeCore {
216 /// Resolve the signal delivery for a claimed step job: the payload to
217 /// expose on [`Step::signal`](crate::Step::signal) and the durable signal entries to delete
218 /// with the step's settlement.
219 pub(crate) async fn resolve_step_signal(
220 &self,
221 job: &JobRecord,
222 run_id: &RunId,
223 step_number: u32,
224 ) -> Result<(Option<Vec<u8>>, Vec<Vec<u8>>)> {
225 if job.headers.contains_key(HEADER_SIGNAL_DELIVERED) {
226 let delivered_key = signal_delivered_kv_key(run_id, step_number);
227 let payload = self
228 .queue
229 .view()
230 .kv_get(&delivered_key)
231 .await?
232 .map(|b| b.to_vec());
233 if payload.is_none() {
234 warn!(run_id = %run_id, step_number, "delivered signal record is missing");
235 }
236 return Ok((payload, vec![delivered_key]));
237 }
238
239 let Some(correlation_key) = job.headers.get(HEADER_SIGNAL_WAIT) else {
240 return Ok((None, Vec::new()));
241 };
242 let wait_key = signal_wait_kv_key(correlation_key);
243 remove_entry(&self.queue, &wait_key, job.id.as_bytes()).await;
244
245 if job.woken_at.is_some() {
246 // A signal promoted this job early. The payload is on the job
247 // record, so redelivery of this step observes it again.
248 return Ok((job.wake_payload.clone(), Vec::new()));
249 }
250
251 // The timeout promoted this job. A prior attempt of this step may
252 // already have consumed the buffer into the delivered record.
253 let delivered_key = signal_delivered_kv_key(run_id, step_number);
254 if let Some(prior) = self.queue.view().kv_get(&delivered_key).await? {
255 return Ok((Some(prior.to_vec()), vec![delivered_key]));
256 }
257 // A signal buffered after this waiter's settlement read of the
258 // buffer, without winning the wake, is consumed here so it is
259 // delivered rather than dropped. The delivery is recorded before
260 // the buffer is consumed, so a retry of this step observes the
261 // same signal.
262 let buf_key = signal_buf_kv_key(correlation_key);
263 if let Some(buffered) = self.queue.view().kv_get(&buf_key).await? {
264 let buffered = buffered.to_vec();
265 self.queue.kv_put(&delivered_key, &buffered).await?;
266 remove_entry(&self.queue, &buf_key, &buffered).await;
267 return Ok((Some(buffered), vec![delivered_key]));
268 }
269 Ok((None, Vec::new()))
270 }
271}