1use 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
14pub const DEFAULT_MAX_ROWS: usize = 16;
16pub const DEFAULT_SINK_TIMEOUT: Duration = Duration::from_secs(5);
18const MAX_ROWS_CEILING: usize = 23;
22const MAX_CURSOR_BYTES: usize = 4_096;
23const MAX_BACKOFF_MS: u64 = 900_000;
24
25pub struct OutcomeDelivery<O> {
27 pub sink: O,
29 pub audience: String,
31 pub metrics: Arc<dyn Metrics>,
33 pub sleep: Arc<dyn Sleep>,
35 pub sink_timeout: Duration,
37 pub max_rows: usize,
39 pub clock: Option<Arc<dyn Clock>>,
42 pub fire_budget: Option<Duration>,
46}
47
48impl<O> OutcomeDelivery<O> {
49 #[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 #[must_use]
71 pub fn with_clock(mut self, clock: Arc<dyn Clock>) -> Self {
72 self.clock = Some(clock);
73 self
74 }
75
76 #[must_use]
78 pub fn with_fire_budget(mut self, budget: Duration) -> Self {
79 self.fire_budget = Some(budget);
80 self
81 }
82
83 #[must_use]
85 pub fn with_sink_timeout(mut self, timeout: Duration) -> Self {
86 self.sink_timeout = timeout;
87 self
88 }
89
90 #[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)] 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 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 let mut stopped_at: Option<usize> = None;
230 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 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 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 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 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 (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 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 ®istry,
485 clock,
486 now,
487 &TickBudget::default(),
488 )
489 .await
490 .unwrap()
491 }
492
493 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 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 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 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 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 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 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 ®istry,
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 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 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 assert!(timer_due(&store).await >= 1_100);
1021 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 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 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 #[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 #[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 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}