Skip to main content

mkit_server/timers/
outcome_delivery.rs

1//! Bounded, at-least-once terminal outcome delivery.
2
3use std::sync::Arc;
4use std::time::Duration;
5
6use crate::pipeline::{Outcome, OutcomeSink};
7use crate::rt::{BoxFuture, Clock, Sleep, with_timeout};
8use crate::store::codec;
9use crate::store::outbox::{guard, plan_ack};
10use crate::store::{Batch, Cursor, NamespaceStore, StoreError, Value, keys};
11use crate::telemetry::Metrics;
12use crate::timers::{DueTimer, Fired, TimerCtx, TimerHandler, TimerKind, registry::kinds};
13
14/// Default rows examined per fire.
15pub const DEFAULT_MAX_ROWS: usize = 16;
16/// Default bound on one sink call.
17pub const DEFAULT_SINK_TIMEOUT: Duration = Duration::from_secs(5);
18// Acknowledgments can cost four operations per row, plus timer completion or
19// reschedule and the shared backlog guard/write. At 23 rows the worst case is
20// 4 * 23 + 2 + 4 = 98 operations, below the shared 100-op transaction cap.
21const MAX_ROWS_CEILING: usize = 23;
22const MAX_CURSOR_BYTES: usize = 4_096;
23const MAX_BACKOFF_MS: u64 = 900_000;
24
25/// Kind-8 delivery driver. The sink deduplicates by reservation id.
26pub struct OutcomeDelivery<O> {
27    /// Deployment sink.
28    pub sink: O,
29    /// Canonical mkit server origin.
30    pub audience: String,
31    /// Rows/bytes gauges and synthetic-row counter.
32    pub metrics: Arc<dyn Metrics>,
33    /// Timer for the per-call sink bound.
34    pub sleep: Arc<dyn Sleep>,
35    /// Bound on one sink call; a timeout is a [`DeliveryError`](crate::pipeline::DeliveryError).
36    pub sink_timeout: Duration,
37    /// Rows examined (and so sink calls attempted) per fire, at least 1.
38    pub max_rows: usize,
39    /// Wall clock for the per-fire bound; without one a fire is bounded only
40    /// by `max_rows` and `sink_timeout`.
41    pub clock: Option<Arc<dyn Clock>>,
42    /// Wall-clock bound on one fire; `None` means twice `sink_timeout`. Once
43    /// exceeded the fire makes no further sink call and reschedules now, so a
44    /// slow-but-successful sink cannot hold a driver.
45    pub fire_budget: Option<Duration>,
46}
47
48impl<O> OutcomeDelivery<O> {
49    /// A driver with the default 5 s sink timeout and 16 rows per fire.
50    #[must_use]
51    pub fn new(
52        sink: O,
53        audience: String,
54        metrics: Arc<dyn Metrics>,
55        sleep: Arc<dyn Sleep>,
56    ) -> Self {
57        Self {
58            sink,
59            audience,
60            metrics,
61            sleep,
62            sink_timeout: DEFAULT_SINK_TIMEOUT,
63            max_rows: DEFAULT_MAX_ROWS,
64            clock: None,
65            fire_budget: None,
66        }
67    }
68
69    /// Bound each fire's wall-clock time (default twice the sink timeout).
70    #[must_use]
71    pub fn with_clock(mut self, clock: Arc<dyn Clock>) -> Self {
72        self.clock = Some(clock);
73        self
74    }
75
76    /// Override the per-fire wall-clock bound; needs [`Self::with_clock`].
77    #[must_use]
78    pub fn with_fire_budget(mut self, budget: Duration) -> Self {
79        self.fire_budget = Some(budget);
80        self
81    }
82
83    /// Override the per-call sink timeout.
84    #[must_use]
85    pub fn with_sink_timeout(mut self, timeout: Duration) -> Self {
86        self.sink_timeout = timeout;
87        self
88    }
89
90    /// Override the rows per fire (clamped to 1..=23 to fit atomic acknowledgments).
91    #[must_use]
92    pub fn with_max_rows(mut self, rows: usize) -> Self {
93        self.max_rows = rows.clamp(1, MAX_ROWS_CEILING);
94        self
95    }
96}
97
98impl<O> core::fmt::Debug for OutcomeDelivery<O> {
99    fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
100        f.debug_struct("OutcomeDelivery")
101            .field("audience", &self.audience)
102            .finish_non_exhaustive()
103    }
104}
105
106impl<O: OutcomeSink, S: NamespaceStore> TimerHandler<S> for OutcomeDelivery<O> {
107    fn kind(&self) -> TimerKind {
108        kinds::OUTCOME_DELIVERY
109    }
110
111    fn fire<'a>(
112        &'a self,
113        ctx: &'a TimerCtx<'a, S>,
114        timer: &'a DueTimer,
115    ) -> BoxFuture<'a, Result<Fired, StoreError>> {
116        Box::pin(self.deliver(ctx, timer))
117    }
118}
119
120fn decode_timer(value: &Value) -> Result<(u32, Option<Cursor>), StoreError> {
121    if value.as_bytes().is_empty() {
122        return Ok((0, None));
123    }
124    let bytes = value.as_bytes();
125    if bytes.len() < 6 {
126        return Err(StoreError::Corrupt("bad outcome timer cursor".into()));
127    }
128    let attempt = u32::from_be_bytes(
129        bytes[..4]
130            .try_into()
131            .map_err(|_| StoreError::Corrupt("bad outcome timer attempt".into()))?,
132    );
133    let length = u16::from_be_bytes(
134        bytes[4..6]
135            .try_into()
136            .map_err(|_| StoreError::Corrupt("bad outcome timer length".into()))?,
137    ) as usize;
138    if length > MAX_CURSOR_BYTES || bytes.len() != 6 + length {
139        return Err(StoreError::Corrupt("bad outcome timer cursor".into()));
140    }
141    Ok((
142        attempt,
143        (length > 0).then(|| Cursor::new(bytes[6..].to_vec())),
144    ))
145}
146
147fn encode_timer(attempt: u32, cursor: Option<&Cursor>) -> Result<Value, StoreError> {
148    let bytes = cursor.map_or(&[][..], Cursor::as_bytes);
149    if bytes.len() > MAX_CURSOR_BYTES {
150        return Err(StoreError::Corrupt("outcome scan cursor too large".into()));
151    }
152    let mut value = Vec::with_capacity(6 + bytes.len());
153    value.extend_from_slice(&attempt.to_be_bytes());
154    let length = u16::try_from(bytes.len())
155        .map_err(|_| StoreError::Corrupt("outcome scan cursor too large".into()))?;
156    value.extend_from_slice(&length.to_be_bytes());
157    value.extend_from_slice(bytes);
158    Ok(Value::new(value))
159}
160
161fn backoff(partition: &crate::store::Partition, attempt: u32) -> Result<u64, StoreError> {
162    let base = 1_000u64
163        .saturating_mul(1u64 << attempt.min(20))
164        .min(MAX_BACKOFF_MS);
165    let mut seed = partition.encode()?.to_vec();
166    seed.extend_from_slice(&attempt.to_be_bytes());
167    let hash = blake3::hash(&seed);
168    let jitter =
169        u64::from(u16::from_be_bytes([hash.as_bytes()[0], hash.as_bytes()[1]])) % (base / 5 + 1);
170    Ok(base.saturating_add(jitter).min(MAX_BACKOFF_MS))
171}
172
173impl<O: OutcomeSink> OutcomeDelivery<O> {
174    #[allow(clippy::cast_precision_loss, clippy::too_many_lines)] // Metrics approximate large counters; one fire plans one guarded delivery batch.
175    async fn deliver<S: NamespaceStore>(
176        &self,
177        ctx: &TimerCtx<'_, S>,
178        timer: &DueTimer,
179    ) -> Result<Fired, StoreError> {
180        let oc_key = keys::outcome_backlog();
181        let oc = ctx.store.get(ctx.partition, &oc_key).await?;
182        let backlog = oc
183            .as_ref()
184            .map(codec::decode_backlog)
185            .transpose()?
186            .unwrap_or_default();
187        let labels_rows = [("shard_kind", ctx.partition.kind()), ("unit", "rows")];
188        let labels_bytes = [("shard_kind", ctx.partition.kind()), ("unit", "bytes")];
189        self.metrics.gauge(
190            "mkit_server_outbox_backlog",
191            &labels_rows,
192            backlog.rows as f64,
193        );
194        self.metrics.gauge(
195            "mkit_server_outbox_backlog",
196            &labels_bytes,
197            backlog.bytes as f64,
198        );
199        if backlog.rows == 0 {
200            return Ok(Fired::Done(
201                Batch::new()
202                    .require(guard(oc_key.clone(), oc.as_ref()))
203                    .delete(oc_key),
204            ));
205        }
206        let (attempt, cursor) = decode_timer(&timer.value)?;
207        let (start, end) = keys::class_range(keys::TAG_OUTCOME_PENDING);
208        // The scan cursor names the last examined row. Scanning past the
209        // delivery budget would strand the unexamined suffix until a full
210        // rotation, so the scan page is the budget.
211        let max_rows = self.max_rows.clamp(1, MAX_ROWS_CEILING);
212        let scan_rows = u32::try_from(max_rows).unwrap_or(u32::MAX);
213        let mut after = cursor.as_ref();
214        let mut page = ctx
215            .store
216            .scan(ctx.partition, &start, &end, after, scan_rows)
217            .await?;
218        if page.entries.is_empty() && after.is_some() && page.next.is_none() {
219            after = None;
220            page = ctx
221                .store
222                .scan(ctx.partition, &start, &end, after, scan_rows)
223                .await?;
224        }
225        let mut batch = Batch::new();
226        let mut delivered = 0u64;
227        let mut sink_retry_ms = 0u64;
228        // Index in `page` of the row whose delivery failed or timed out.
229        let mut stopped_at: Option<usize> = None;
230        // Index of the first row not tried because the fire's wall-clock
231        // budget ran out.
232        let mut cut_at: Option<usize> = None;
233        let started_ms = self.clock.as_ref().map(|clock| clock.now_ms());
234        let budget_ms = i64::try_from(
235            self.fire_budget
236                .unwrap_or_else(|| self.sink_timeout.saturating_mul(2))
237                .as_millis(),
238        )
239        .unwrap_or(i64::MAX);
240        // Acknowledged rows are planned after the sink loop, against a fresh
241        // `oc`, so slow sink awaits sit outside the read-modify-write window.
242        let mut acked: Vec<(String, Value, u64)> = Vec::new();
243        for (index, (key, _)) in page.entries.iter().enumerate().take(max_rows) {
244            let Some(keys::ParsedKey::OutcomePending {
245                seq,
246                reservation_id: rid,
247            }) = keys::parse(key)
248            else {
249                tracing::warn!("invalid outcome index row; retaining");
250                continue;
251            };
252            let Some(value) = ctx
253                .store
254                .get(ctx.partition, &keys::reservation(&rid)?)
255                .await?
256            else {
257                tracing::warn!("missing terminal outcome row; retaining index");
258                continue;
259            };
260            let record = match codec::decode_reservation(&value) {
261                Ok(record) => record,
262                Err(err) => {
263                    tracing::warn!(error = %err, "undecodable terminal outcome; retaining");
264                    continue;
265                }
266            };
267            let outcome =
268                match Outcome::from_reservation(rid.clone(), self.audience.clone(), record) {
269                    Ok(outcome) => outcome,
270                    Err(reason) => {
271                        tracing::warn!(reason, "indexed outcome cannot be delivered; retaining");
272                        continue;
273                    }
274                };
275            let acknowledged = if rid.starts_with("s:") {
276                self.metrics
277                    .incr("mkit_server_synthetic_outcomes_acked", &[], 1);
278                true
279            } else {
280                if let (Some(clock), Some(started)) = (&self.clock, started_ms)
281                    && clock.now_ms().saturating_sub(started) >= budget_ms
282                {
283                    cut_at = Some(index);
284                    break;
285                }
286                match with_timeout(&*self.sleep, self.sink_timeout, self.sink.deliver(&outcome))
287                    .await
288                {
289                    Ok(Ok(())) => true,
290                    Ok(Err(err)) => {
291                        sink_retry_ms = sink_retry_ms.max(err.retry_after.map_or(0, |hint| {
292                            u64::try_from(hint.as_millis()).unwrap_or(u64::MAX)
293                        }));
294                        tracing::warn!(reason = %err.reason, "outcome delivery failed");
295                        stopped_at = Some(index);
296                        false
297                    }
298                    Err(_) => {
299                        tracing::warn!("outcome delivery timed out");
300                        stopped_at = Some(index);
301                        false
302                    }
303                }
304            };
305            if acknowledged {
306                acked.push((rid, value, seq));
307            }
308            // A failing or hung sink is not asked again in this fire.
309            if stopped_at.is_some() {
310                break;
311            }
312        }
313        let fresh_oc = if acked.is_empty() {
314            oc.clone()
315        } else {
316            ctx.store.get(ctx.partition, &oc_key).await?
317        };
318        let fresh_backlog = fresh_oc
319            .as_ref()
320            .map(codec::decode_backlog)
321            .transpose()?
322            .unwrap_or_default();
323        // Completion and acknowledgments must share a snapshot. A later append
324        // fails this guard, retaining the timer for the next fire to re-plan.
325        batch.preconditions.push(guard(oc_key, fresh_oc.as_ref()));
326        for (rid, value, seq) in &acked {
327            match plan_ack(
328                rid,
329                value,
330                *seq,
331                fresh_oc.as_ref(),
332                &mut batch.preconditions,
333                &mut batch.writes,
334            ) {
335                Ok(()) => delivered += 1,
336                Err(err) => {
337                    tracing::warn!(error = %err, "outcome acknowledgment not planned; retaining");
338                }
339            }
340        }
341        if delivered == fresh_backlog.rows {
342            return Ok(Fired::Done(batch));
343        }
344        // The next fire resumes just after the failed row. Resuming after the
345        // page instead would, when the page reached the end of the range, wrap
346        // to the head and retry the same failing row first forever, starving
347        // every row behind it. The page's cursor is opaque, so ask the store
348        // for the one after the failed row.
349        let resume = match (stopped_at, cut_at) {
350            (Some(index), _) => {
351                let limit = u32::try_from(index + 1).unwrap_or(u32::MAX);
352                ctx.store
353                    .scan(ctx.partition, &start, &end, after, limit)
354                    .await?
355                    .next
356            }
357            // Resume at the first untried row: after the one before it.
358            (None, Some(index)) if index > 0 => {
359                let limit = u32::try_from(index).unwrap_or(u32::MAX);
360                ctx.store
361                    .scan(ctx.partition, &start, &end, after, limit)
362                    .await?
363                    .next
364            }
365            (None, _) => page.next,
366        };
367        // Only a failed or hung sink call backs off. A page that delivered
368        // nothing because every row was junk or synthetic is not a sink fault.
369        let next_attempt = if delivered > 0 {
370            0
371        } else if stopped_at.is_some() {
372            attempt.saturating_add(1)
373        } else {
374            attempt
375        };
376        let delay = if delivered > 0 {
377            0
378        } else {
379            backoff(ctx.partition, attempt)?
380                .max(sink_retry_ms)
381                .min(MAX_BACKOFF_MS)
382        };
383        let due_at_ms = ctx
384            .now_ms
385            .saturating_add(delay)
386            .max(timer.due_at_ms.saturating_add(1));
387        Ok(Fired::Reschedule {
388            due_at_ms,
389            value: encode_timer(next_attempt, resume.as_ref())?,
390            batch,
391        })
392    }
393}
394
395#[cfg(test)]
396mod tests {
397    use super::*;
398    use crate::memory::MemoryKv;
399    use crate::pipeline::{DeliveryError, OutcomeKind};
400    use crate::repo::NamespaceKey;
401    use crate::rt::{ManualClock, ManualSleep};
402    use crate::store::codec::{Backlog, ReservationV1};
403    use crate::store::outbox::{OutboxBuilder, Terminal};
404    use crate::store::{BatchOutcome, Key, Partition, Precondition, Write};
405    use crate::telemetry::NoopMetrics;
406    use crate::timers::{TickBudget, TimerRegistry, run_due};
407    use std::sync::Mutex;
408
409    #[derive(Default)]
410    struct Capture {
411        seen: Mutex<Vec<Outcome>>,
412        fail: bool,
413    }
414    impl OutcomeSink for Capture {
415        async fn deliver(&self, outcome: &Outcome) -> Result<(), DeliveryError> {
416            self.seen.lock().unwrap().push(outcome.clone());
417            if self.fail {
418                Err(DeliveryError::new("secret sink failure", None))
419            } else {
420                Ok(())
421            }
422        }
423    }
424
425    fn partition() -> Partition {
426        Partition::Namespace(NamespaceKey::deployment_default())
427    }
428
429    async fn seed(store: &MemoryKv, ids: &[&str]) {
430        let p = partition();
431        let prior = codec::encode_reservation(&ReservationV1::Ticketed { ticket_id: [7; 32] });
432        let mut initial = Batch::new();
433        for rid in ids {
434            initial
435                .writes
436                .push(Write::Put(keys::reservation(rid).unwrap(), prior.clone()));
437        }
438        assert_eq!(
439            store.apply(&p, initial).await.unwrap(),
440            BatchOutcome::Committed
441        );
442        let mut outbox = OutboxBuilder::new(None, None).unwrap();
443        for rid in ids {
444            outbox.outcome(
445                rid,
446                &prior,
447                Terminal::new(ReservationV1::Committed {
448                    repository: "repo".into(),
449                    occurred_at_ms: 100,
450                    bytes_stored: 1,
451                    new_to_repo: 1,
452                    new_to_store: 1,
453                    refs: vec![],
454                })
455                .unwrap(),
456            );
457        }
458        let mut batch = Batch::new();
459        outbox
460            .try_finish(&mut batch.preconditions, &mut batch.writes)
461            .unwrap();
462        assert_eq!(batch.writes.iter().filter(|write| matches!(write, Write::Put(key, _) if matches!(keys::parse(key), Some(keys::ParsedKey::Timer { kind: 8, .. })))).count(), 1);
463        assert_eq!(
464            store.apply(&p, batch).await.unwrap(),
465            BatchOutcome::Committed
466        );
467    }
468
469    async fn fire(
470        store: &MemoryKv,
471        clock: &ManualClock,
472        sink: Arc<impl OutcomeSink>,
473        now: u64,
474    ) -> crate::timers::RunReport {
475        let registry = TimerRegistry::new().register(OutcomeDelivery::new(
476            sink,
477            "https://example.test".into(),
478            Arc::new(NoopMetrics),
479            Arc::new(ManualSleep::new()),
480        ));
481        run_due(
482            store,
483            &partition(),
484            &registry,
485            clock,
486            now,
487            &TickBudget::default(),
488        )
489        .await
490        .unwrap()
491    }
492
493    // A second terminal outcome arrives inside the first sink await.
494    struct AppendDuringAck(Arc<MemoryKv>);
495    impl OutcomeSink for AppendDuringAck {
496        async fn deliver(&self, outcome: &Outcome) -> Result<(), DeliveryError> {
497            assert_eq!(outcome.reservation_id, "first");
498            let prior = codec::encode_reservation(&ReservationV1::Ticketed { ticket_id: [8; 32] });
499            let p = partition();
500            let os = self.0.get(&p, &keys::outbox_sequence()).await.unwrap();
501            let oc = self.0.get(&p, &keys::outcome_backlog()).await.unwrap();
502            let mut outbox = OutboxBuilder::new(os.as_ref(), oc.as_ref()).unwrap();
503            self.0
504                .apply(
505                    &p,
506                    Batch::new().put(keys::reservation("second").unwrap(), prior.clone()),
507                )
508                .await
509                .unwrap();
510            outbox.outcome(
511                "second",
512                &prior,
513                Terminal::new(ReservationV1::Aborted {
514                    repository: "repo".into(),
515                    occurred_at_ms: 100,
516                    reason: codec::AbortReason::RefConflict,
517                    detail: String::new(),
518                })
519                .unwrap(),
520            );
521            let mut batch = Batch::new();
522            outbox
523                .try_finish(&mut batch.preconditions, &mut batch.writes)
524                .unwrap();
525            assert_eq!(
526                self.0.apply(&p, batch).await.unwrap(),
527                BatchOutcome::Committed
528            );
529            Ok(())
530        }
531    }
532
533    #[tokio::test]
534    async fn append_during_ack_keeps_a_timer_and_delivers_on_the_next_fire() {
535        let clock = Arc::new(ManualClock::new(100));
536        let store = Arc::new(MemoryKv::with_clock(clock.clone()));
537        seed(&store, &["first"]).await;
538        let report = fire(
539            &store,
540            &clock,
541            Arc::new(AppendDuringAck(store.clone())),
542            100,
543        )
544        .await;
545        assert_eq!(report.fired, 1);
546        assert_eq!(backlog_rows(&store).await, 1);
547        let (start, end) = keys::class_range(keys::TAG_TIMER);
548        let timers = store
549            .scan(&partition(), &start, &end, None, 10)
550            .await
551            .unwrap();
552        assert_eq!(timers.entries.len(), 1, "one delivery timer remains");
553        let (start, end) = keys::class_range(keys::TAG_OUTCOME_PENDING);
554        let pending = store
555            .scan(&partition(), &start, &end, None, 10)
556            .await
557            .unwrap();
558        assert_eq!(pending.entries.len(), 1);
559        let due = timer_due(&store).await;
560        assert!(due <= 101, "remaining outcome is due now");
561        clock.set(i64::try_from(due).unwrap());
562        let sink = Arc::new(Capture::default());
563        assert_eq!(fire(&store, &clock, sink.clone(), due).await.fired, 1);
564        assert_eq!(sink.seen.lock().unwrap()[0].reservation_id, "second");
565        assert!(matches!(
566            sink.seen.lock().unwrap()[0].kind,
567            OutcomeKind::Aborted { .. }
568        ));
569        assert_eq!(backlog_rows(&store).await, 0);
570        assert!(
571            store
572                .scan(&partition(), &start, &end, None, 10)
573                .await
574                .unwrap()
575                .entries
576                .is_empty()
577        );
578    }
579
580    async fn planned_delivery(store: &MemoryKv, sink: Arc<Capture>) -> (Key, Value, Fired) {
581        let (start, end) = keys::class_range(keys::TAG_TIMER);
582        let (key, value) = store
583            .scan(&partition(), &start, &end, None, 10)
584            .await
585            .unwrap()
586            .entries[0]
587            .clone();
588        let Some(keys::ParsedKey::Timer {
589            due_at_ms,
590            kind,
591            reference,
592        }) = keys::parse(&key)
593        else {
594            panic!("delivery timer");
595        };
596        let timer = DueTimer {
597            due_at_ms,
598            kind: crate::timers::TimerKind::new(kind),
599            reference,
600            value: value.clone(),
601        };
602        let delivery = OutcomeDelivery::new(
603            sink,
604            "https://example.test".into(),
605            Arc::new(NoopMetrics),
606            Arc::new(ManualSleep::new()),
607        );
608        let fired = delivery
609            .fire(
610                &TimerCtx {
611                    store,
612                    partition: &partition(),
613                    now_ms: 100,
614                },
615                &timer,
616            )
617            .await
618            .unwrap();
619        (key, value, fired)
620    }
621
622    #[tokio::test]
623    async fn append_after_fresh_read_fails_the_guard_and_replans_delivery() {
624        let clock = Arc::new(ManualClock::new(100));
625        let store = Arc::new(MemoryKv::with_clock(clock.clone()));
626        seed(&store, &["first"]).await;
627        let sink = Arc::new(Capture::default());
628        let prior_oc = store
629            .get(&partition(), &keys::outcome_backlog())
630            .await
631            .unwrap()
632            .unwrap();
633        let (key, value, fired) = planned_delivery(&store, sink.clone()).await;
634        let Fired::Done(batch) = fired else {
635            panic!("original row acknowledged");
636        };
637        assert!(
638            batch
639                .preconditions
640                .contains(&Precondition::Equals(keys::outcome_backlog(), prior_oc))
641        );
642        // Append after the handler's fresh read and before applying its Done.
643        let first = sink.seen.lock().unwrap()[0].clone();
644        AppendDuringAck(store.clone())
645            .deliver(&first)
646            .await
647            .unwrap();
648        let batch = batch
649            .require(Precondition::Equals(key.clone(), value))
650            .delete(key.clone());
651        assert!(matches!(
652            store.apply(&partition(), batch).await.unwrap(),
653            BatchOutcome::PreconditionFailed { .. }
654        ));
655        assert_eq!(backlog_rows(&store).await, 2);
656        assert!(store.get(&partition(), &key).await.unwrap().is_some());
657        assert_eq!(fire(&store, &clock, sink.clone(), 100).await.fired, 1);
658        assert_eq!(backlog_rows(&store).await, 0);
659        assert!(store.get(&partition(), &key).await.unwrap().is_none());
660        let seen = sink.seen.lock().unwrap();
661        assert_eq!(
662            seen.iter().filter(|o| o.reservation_id == "first").count(),
663            2
664        );
665        assert_eq!(
666            seen.iter().filter(|o| o.reservation_id == "second").count(),
667            1
668        );
669    }
670
671    #[tokio::test]
672    async fn no_acknowledgments_guard_the_initial_backlog() {
673        let store = MemoryKv::with_clock(Arc::new(ManualClock::new(100)));
674        seed(&store, &["first"]).await;
675        let prior = store
676            .get(&partition(), &keys::outcome_backlog())
677            .await
678            .unwrap()
679            .unwrap();
680        let sink = Arc::new(Capture {
681            fail: true,
682            ..Capture::default()
683        });
684        let (_, _, fired) = planned_delivery(&store, sink).await;
685        let Fired::Reschedule { batch, .. } = fired else {
686            panic!("failed sink retries");
687        };
688        assert!(
689            batch
690                .preconditions
691                .contains(&Precondition::Equals(keys::outcome_backlog(), prior))
692        );
693        assert!(batch.writes.is_empty());
694    }
695
696    #[tokio::test]
697    async fn all_acknowledgments_return_done_with_an_empty_backlog() {
698        let store = MemoryKv::with_clock(Arc::new(ManualClock::new(100)));
699        seed(&store, &["first", "second"]).await;
700        let sink = Arc::new(Capture::default());
701        let (key, value, fired) = planned_delivery(&store, sink.clone()).await;
702        let Fired::Done(batch) = fired else {
703            panic!("all rows acknowledged");
704        };
705        assert_eq!(
706            store
707                .apply(
708                    &partition(),
709                    batch
710                        .require(Precondition::Equals(key.clone(), value))
711                        .delete(key.clone())
712                )
713                .await
714                .unwrap(),
715            BatchOutcome::Committed
716        );
717        assert_eq!(sink.seen.lock().unwrap().len(), 2);
718        assert_eq!(backlog_rows(&store).await, 0);
719        assert!(store.get(&partition(), &key).await.unwrap().is_none());
720    }
721
722    #[tokio::test]
723    async fn delivers_in_sequence_and_acks_synthetic_locally() {
724        let clock = Arc::new(ManualClock::new(100));
725        let store = MemoryKv::with_clock(clock.clone());
726        seed(&store, &["one", "s:local", "three"]).await;
727        let sink = Arc::new(Capture::default());
728        let report = fire(&store, &clock, sink.clone(), 100).await;
729        assert_eq!(report.fired, 1);
730        let seen = sink.seen.lock().unwrap().clone();
731        assert_eq!(
732            seen.iter()
733                .map(|outcome| outcome.reservation_id.as_str())
734                .collect::<Vec<_>>(),
735            ["one", "three"]
736        );
737        assert!(
738            seen.iter()
739                .all(|outcome| matches!(outcome.kind, OutcomeKind::Committed { .. }))
740        );
741        assert_eq!(
742            store
743                .get(&partition(), &keys::outcome_backlog())
744                .await
745                .unwrap(),
746            None
747        );
748        for rid in ["one", "s:local", "three"] {
749            assert_eq!(
750                store
751                    .get(&partition(), &keys::reservation(rid).unwrap())
752                    .await
753                    .unwrap(),
754                None
755            );
756        }
757    }
758
759    #[tokio::test]
760    async fn sink_error_and_poison_row_retry_without_deletion() {
761        let clock = Arc::new(ManualClock::new(100));
762        let store = MemoryKv::with_clock(clock.clone());
763        seed(&store, &["poison", "behind"]).await;
764        let poison = keys::reservation("poison").unwrap();
765        assert_eq!(
766            store
767                .apply(
768                    &partition(),
769                    Batch::new().put(poison.clone(), Value::new(b"bad".to_vec()))
770                )
771                .await
772                .unwrap(),
773            BatchOutcome::Committed
774        );
775        let sink = Arc::new(Capture::default());
776        let report = fire(&store, &clock, sink.clone(), 100).await;
777        assert_eq!(report.fired, 1);
778        assert_eq!(sink.seen.lock().unwrap()[0].reservation_id, "behind");
779        assert_eq!(
780            codec::decode_backlog(
781                &store
782                    .get(&partition(), &keys::outcome_backlog())
783                    .await
784                    .unwrap()
785                    .unwrap()
786            )
787            .unwrap()
788            .rows,
789            1
790        );
791        assert!(store.get(&partition(), &poison).await.unwrap().is_some());
792
793        let failing = Arc::new(Capture {
794            fail: true,
795            ..Capture::default()
796        });
797        let second = MemoryKv::with_clock(clock.clone());
798        seed(&second, &["retry"]).await;
799        let report = fire(&second, &clock, failing.clone(), 100).await;
800        assert_eq!(report.fired, 1);
801        assert_eq!(
802            codec::decode_backlog(
803                &second
804                    .get(&partition(), &keys::outcome_backlog())
805                    .await
806                    .unwrap()
807                    .unwrap()
808            )
809            .unwrap()
810            .rows,
811            1
812        );
813        let (start, end) = keys::class_range(keys::TAG_TIMER);
814        let page = second
815            .scan(&partition(), &start, &end, None, 10)
816            .await
817            .unwrap();
818        let due = page
819            .entries
820            .iter()
821            .find_map(|(key, _)| match keys::parse(key) {
822                Some(keys::ParsedKey::Timer {
823                    kind: 8, due_at_ms, ..
824                }) => Some(due_at_ms),
825                _ => None,
826            })
827            .unwrap();
828        assert!(due >= 1_100);
829        let first = failing.seen.lock().unwrap()[0].clone();
830        clock.set(i64::try_from(due).unwrap());
831        assert_eq!(fire(&second, &clock, failing.clone(), due).await.fired, 1);
832        assert_eq!(first, failing.seen.lock().unwrap()[1]);
833        assert_eq!(first.reservation_id, "retry");
834        assert!(matches!(
835            codec::decode_backlog(
836                &second
837                    .get(&partition(), &keys::outcome_backlog())
838                    .await
839                    .unwrap()
840                    .unwrap()
841            )
842            .unwrap(),
843            Backlog { rows: 1, .. }
844        ));
845    }
846
847    /// A sink that never answers.
848    struct Hang(Mutex<u32>);
849    impl OutcomeSink for Hang {
850        async fn deliver(&self, _: &Outcome) -> Result<(), DeliveryError> {
851            *self.0.lock().unwrap() += 1;
852            futures::future::pending().await
853        }
854    }
855
856    async fn backlog_rows(store: &MemoryKv) -> u64 {
857        store
858            .get(&partition(), &keys::outcome_backlog())
859            .await
860            .unwrap()
861            .map_or(0, |v| codec::decode_backlog(&v).unwrap().rows)
862    }
863
864    async fn timer_due(store: &MemoryKv) -> u64 {
865        let (start, end) = keys::class_range(keys::TAG_TIMER);
866        let page = store
867            .scan(&partition(), &start, &end, None, 10)
868            .await
869            .unwrap();
870        page.entries
871            .iter()
872            .find_map(|(key, _)| match keys::parse(key) {
873                Some(keys::ParsedKey::Timer {
874                    kind: 8, due_at_ms, ..
875                }) => Some(due_at_ms),
876                _ => None,
877            })
878            .unwrap()
879    }
880
881    /// A sink that succeeds, but only after the clock advances by `step_ms`.
882    struct Slow {
883        clock: Arc<ManualClock>,
884        step_ms: i64,
885        seen: Mutex<Vec<String>>,
886    }
887    impl OutcomeSink for Slow {
888        async fn deliver(&self, outcome: &Outcome) -> Result<(), DeliveryError> {
889            self.seen
890                .lock()
891                .unwrap()
892                .push(outcome.reservation_id.clone());
893            self.clock.set(self.clock.now_ms() + self.step_ms);
894            Ok(())
895        }
896    }
897
898    #[tokio::test]
899    async fn a_slow_but_successful_sink_is_cut_at_the_fire_budget() {
900        let clock = Arc::new(ManualClock::new(100));
901        let store = MemoryKv::with_clock(clock.clone());
902        seed(&store, &["a", "b", "c", "d"]).await;
903        let sink = Arc::new(Slow {
904            clock: clock.clone(),
905            step_ms: 400,
906            seen: Mutex::new(Vec::new()),
907        });
908        // Budget 2 x 250 ms: two 400 ms calls exhaust it.
909        let delivery = OutcomeDelivery::new(
910            sink.clone(),
911            "https://example.test".into(),
912            Arc::new(NoopMetrics),
913            Arc::new(ManualSleep::new()),
914        )
915        .with_sink_timeout(Duration::from_millis(250))
916        .with_clock(clock.clone());
917        fire_with(&store, &clock, delivery, 100).await;
918        assert_eq!(*sink.seen.lock().unwrap(), ["a", "b"]);
919        assert_eq!(backlog_rows(&store).await, 2);
920        // Rescheduled now (not backed off), resuming at the untried rows.
921        let due = timer_due(&store).await;
922        assert!(due <= 101, "due {due}");
923        clock.set(i64::try_from(due).unwrap());
924        let delivery = OutcomeDelivery::new(
925            sink.clone(),
926            "https://example.test".into(),
927            Arc::new(NoopMetrics),
928            Arc::new(ManualSleep::new()),
929        );
930        fire_with(&store, &clock, delivery, due).await;
931        assert_eq!(*sink.seen.lock().unwrap(), ["a", "b", "c", "d"]);
932        assert_eq!(backlog_rows(&store).await, 0);
933    }
934
935    #[tokio::test]
936    async fn an_undeliverable_page_does_not_raise_the_backoff() {
937        let clock = Arc::new(ManualClock::new(100));
938        let store = MemoryKv::with_clock(clock.clone());
939        seed(&store, &["poison"]).await;
940        store
941            .apply(
942                &partition(),
943                Batch::new().put(
944                    keys::reservation("poison").unwrap(),
945                    Value::new(b"bad".to_vec()),
946                ),
947            )
948            .await
949            .unwrap();
950        let sink = Arc::new(Capture::default());
951        let mut now = 100u64;
952        let mut dues = Vec::new();
953        for _ in 0..3 {
954            clock.set(i64::try_from(now).unwrap());
955            fire(&store, &clock, sink.clone(), now).await;
956            let due = timer_due(&store).await;
957            dues.push(due - now);
958            now = due;
959        }
960        assert!(sink.seen.lock().unwrap().is_empty());
961        // No sink call failed, so the delay stays at the first backoff step.
962        assert!(dues.iter().all(|gap| *gap <= 1_200), "{dues:?}");
963    }
964
965    async fn fire_with(
966        store: &MemoryKv,
967        clock: &ManualClock,
968        delivery: OutcomeDelivery<Arc<impl OutcomeSink>>,
969        now: u64,
970    ) {
971        let registry = TimerRegistry::new().register(delivery);
972        run_due(
973            store,
974            &partition(),
975            &registry,
976            clock,
977            now,
978            &TickBudget::default(),
979        )
980        .await
981        .unwrap();
982    }
983
984    #[tokio::test]
985    async fn a_hanging_sink_is_cut_at_the_timeout_and_stops_the_fire() {
986        let clock = Arc::new(ManualClock::new(100));
987        let store = MemoryKv::with_clock(clock.clone());
988        seed(&store, &["a", "b", "c"]).await;
989        let sink = Arc::new(Hang(Mutex::new(0)));
990        let sleeper = ManualSleep::elapsed();
991        let delivery = OutcomeDelivery::new(
992            sink.clone(),
993            "https://example.test".into(),
994            Arc::new(NoopMetrics),
995            Arc::new(sleeper.clone()),
996        )
997        .with_sink_timeout(Duration::from_millis(250));
998        fire_with(&store, &clock, delivery, 100).await;
999        // One call, cut by the 250 ms timer; the other rows are not tried.
1000        assert_eq!(*sink.0.lock().unwrap(), 1);
1001        assert_eq!(sleeper.requested(), [Duration::from_millis(250)]);
1002        assert_eq!(backlog_rows(&store).await, 3);
1003        // Backed off: the retry is later than now.
1004        assert!(timer_due(&store).await >= 1_100);
1005    }
1006
1007    #[tokio::test]
1008    async fn first_failure_stops_the_fire_and_no_row_is_lost() {
1009        let clock = Arc::new(ManualClock::new(100));
1010        let store = MemoryKv::with_clock(clock.clone());
1011        seed(&store, &["a", "b", "c"]).await;
1012        let failing = Arc::new(Capture {
1013            fail: true,
1014            ..Capture::default()
1015        });
1016        fire(&store, &clock, failing.clone(), 100).await;
1017        assert_eq!(failing.seen.lock().unwrap().len(), 1);
1018        assert_eq!(backlog_rows(&store).await, 3);
1019        // Backed off: the retry is later than now.
1020        assert!(timer_due(&store).await >= 1_100);
1021        // A healthy sink later delivers every row.
1022        let ok = Arc::new(Capture::default());
1023        let (start, end) = keys::class_range(keys::TAG_TIMER);
1024        let page = store
1025            .scan(&partition(), &start, &end, None, 10)
1026            .await
1027            .unwrap();
1028        let due = page
1029            .entries
1030            .iter()
1031            .find_map(|(key, _)| match keys::parse(key) {
1032                Some(keys::ParsedKey::Timer {
1033                    kind: 8, due_at_ms, ..
1034                }) => Some(due_at_ms),
1035                _ => None,
1036            })
1037            .unwrap();
1038        clock.set(i64::try_from(due).unwrap());
1039        // The fire resumes after the failed row `a`, so `b` and `c` go first;
1040        // `a` follows once the cursor wraps.
1041        fire(&store, &clock, ok.clone(), due).await;
1042        assert_eq!(ok.seen.lock().unwrap().len(), 2);
1043        assert_eq!(backlog_rows(&store).await, 1);
1044        let (start, end) = keys::class_range(keys::TAG_TIMER);
1045        let page = store
1046            .scan(&partition(), &start, &end, None, 10)
1047            .await
1048            .unwrap();
1049        let due = page
1050            .entries
1051            .iter()
1052            .find_map(|(key, _)| match keys::parse(key) {
1053                Some(keys::ParsedKey::Timer {
1054                    kind: 8, due_at_ms, ..
1055                }) => Some(due_at_ms),
1056                _ => None,
1057            })
1058            .unwrap();
1059        clock.set(i64::try_from(due).unwrap());
1060        fire(&store, &clock, ok.clone(), due).await;
1061        assert_eq!(ok.seen.lock().unwrap().len(), 3);
1062        assert_eq!(backlog_rows(&store).await, 0);
1063    }
1064
1065    /// A sink that rejects one reservation id and takes the rest.
1066    struct Picky(Mutex<Vec<String>>);
1067    impl OutcomeSink for Picky {
1068        async fn deliver(&self, outcome: &Outcome) -> Result<(), DeliveryError> {
1069            if outcome.reservation_id == "b" {
1070                return Err(DeliveryError::new("rejected", None));
1071            }
1072            self.0.lock().unwrap().push(outcome.reservation_id.clone());
1073            Ok(())
1074        }
1075    }
1076
1077    /// A row the sink keeps refusing does not starve the rows behind it: the
1078    /// fire stops at it, and the next fire resumes after it rather than
1079    /// retrying it first again.
1080    #[tokio::test]
1081    async fn a_refused_row_does_not_starve_the_rows_behind_it() {
1082        let clock = Arc::new(ManualClock::new(100));
1083        let store = MemoryKv::with_clock(clock.clone());
1084        seed(&store, &["a", "b", "c", "d"]).await;
1085        let sink = Arc::new(Picky(Mutex::new(Vec::new())));
1086        let mut now = 100u64;
1087        for _ in 0..3 {
1088            fire_with(
1089                &store,
1090                &clock,
1091                OutcomeDelivery::new(
1092                    sink.clone(),
1093                    "https://example.test".into(),
1094                    Arc::new(NoopMetrics),
1095                    Arc::new(ManualSleep::new()),
1096                ),
1097                now,
1098            )
1099            .await;
1100            now += 1_000_000;
1101            clock.set(i64::try_from(now).unwrap());
1102        }
1103        assert_eq!(*sink.0.lock().unwrap(), ["a", "c", "d"]);
1104        assert_eq!(backlog_rows(&store).await, 1, "only b remains");
1105    }
1106
1107    #[tokio::test]
1108    async fn max_rows_bounds_one_fire_and_the_cursor_rotates() {
1109        let clock = Arc::new(ManualClock::new(100));
1110        let store = MemoryKv::with_clock(clock.clone());
1111        seed(&store, &["a", "b", "c", "d", "e"]).await;
1112        let sink = Arc::new(Capture::default());
1113        let mk = |sink: Arc<Capture>| {
1114            OutcomeDelivery::new(
1115                sink,
1116                "https://example.test".into(),
1117                Arc::new(NoopMetrics),
1118                Arc::new(ManualSleep::new()),
1119            )
1120            .with_max_rows(2)
1121        };
1122        fire_with(&store, &clock, mk(sink.clone()), 100).await;
1123        assert_eq!(sink.seen.lock().unwrap().len(), 2);
1124        assert_eq!(backlog_rows(&store).await, 3);
1125        clock.set(101);
1126        fire_with(&store, &clock, mk(sink.clone()), 101).await;
1127        clock.set(102);
1128        fire_with(&store, &clock, mk(sink.clone()), 102).await;
1129        assert_eq!(sink.seen.lock().unwrap().len(), 5);
1130        assert_eq!(backlog_rows(&store).await, 0);
1131    }
1132
1133    #[tokio::test]
1134    async fn configured_max_rows_24_or_25_drains_without_exceeding_atomic_batch_limit() {
1135        for configured in [24, 25] {
1136            let clock = Arc::new(ManualClock::new(100));
1137            let store = MemoryKv::with_clock(clock.clone());
1138            let ids: Vec<_> = (0..25).map(|i| format!("row-{i}")).collect();
1139            let ids: Vec<_> = ids.iter().map(String::as_str).collect();
1140            seed(&store, &ids).await;
1141            let sink = Arc::new(Capture::default());
1142            for tick in 0..3 {
1143                let now = 100 + tick * 1_000_000;
1144                clock.set(i64::try_from(now).unwrap());
1145                fire_with(
1146                    &store,
1147                    &clock,
1148                    OutcomeDelivery::new(
1149                        sink.clone(),
1150                        "https://example.test".into(),
1151                        Arc::new(NoopMetrics),
1152                        Arc::new(ManualSleep::new()),
1153                    )
1154                    .with_max_rows(configured),
1155                    now,
1156                )
1157                .await;
1158            }
1159            assert_eq!(
1160                sink.seen.lock().unwrap().len(),
1161                25,
1162                "configured row bound {configured}"
1163            );
1164            assert_eq!(
1165                backlog_rows(&store).await,
1166                0,
1167                "configured row bound {configured}"
1168            );
1169        }
1170    }
1171
1172    #[cfg(feature = "remote-hooks")]
1173    async fn rows(store: &MemoryKv) -> u64 {
1174        codec::decode_backlog(
1175            &store
1176                .get(&partition(), &keys::outcome_backlog())
1177                .await
1178                .unwrap()
1179                .unwrap(),
1180        )
1181        .unwrap()
1182        .rows
1183    }
1184
1185    /// A remote sink that cannot reach its hook leaves the row queued, and each
1186    /// retry is signed afresh (new nonce and validity window).
1187    #[cfg(feature = "remote-hooks")]
1188    #[tokio::test]
1189    async fn a_failing_remote_hook_keeps_the_row_and_each_retry_is_signed_afresh() {
1190        use crate::hooks::HookClient;
1191        use crate::hooks::RemoteOutcomes;
1192        use crate::hooks::tests::{MockChannel, Step, channel_of, signer};
1193        use crate::rt::{Clock, ManualSleep};
1194
1195        let clock = Arc::new(ManualClock::new(100));
1196        let store = MemoryKv::with_clock(clock.clone());
1197        seed(&store, &["retry"]).await;
1198        // The client shares the test clock, so a retry's window moves.
1199        let hook = Arc::new(
1200            HookClient::new(
1201                MockChannel::new(Step::Fail(crate::hooks::ChannelError::Timeout)),
1202                "https://example.test",
1203                Some(signer()),
1204                clock.clone() as Arc<dyn Clock>,
1205                Arc::new(ManualSleep::new()),
1206            )
1207            .unwrap(),
1208        );
1209        let sink = Arc::new(RemoteOutcomes::new(hook.clone()));
1210        assert_eq!(fire(&store, &clock, sink.clone(), 100).await.fired, 1);
1211        assert_eq!(rows(&store).await, 1);
1212        let (start, end) = keys::class_range(keys::TAG_TIMER);
1213        let page = store
1214            .scan(&partition(), &start, &end, None, 10)
1215            .await
1216            .unwrap();
1217        let due = page
1218            .entries
1219            .iter()
1220            .find_map(|(key, _)| match keys::parse(key) {
1221                Some(keys::ParsedKey::Timer {
1222                    kind: 8, due_at_ms, ..
1223                }) => Some(due_at_ms),
1224                _ => None,
1225            })
1226            .unwrap();
1227        clock.set(i64::try_from(due).unwrap());
1228        assert_eq!(fire(&store, &clock, sink, due).await.fired, 1);
1229        assert_eq!(rows(&store).await, 1);
1230
1231        let seen = channel_of(&hook).seen.lock().unwrap();
1232        assert_eq!(seen.len(), 2);
1233        let header = |i: usize, name: &str| {
1234            seen[i]
1235                .headers
1236                .iter()
1237                .find(|(n, _)| *n == name)
1238                .unwrap()
1239                .1
1240                .clone()
1241        };
1242        assert_ne!(
1243            header(0, "X-Mkit-Hook-Nonce"),
1244            header(1, "X-Mkit-Hook-Nonce")
1245        );
1246        assert_ne!(
1247            header(0, "X-Mkit-Hook-Created-At"),
1248            header(1, "X-Mkit-Hook-Created-At")
1249        );
1250        assert_eq!(seen[0].body, seen[1].body);
1251    }
1252}