Skip to main content

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, &current).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}