mkit_server/relay/mod.rs
1//! Source-side, at-least-once outbox delivery. Target watermarks make
2//! redelivery idempotent, including a crash before source cleanup.
3//!
4//! Each source persists one relay scan cycle in `rs`: the `os` snapshot at
5//! cycle start, the last inspected sequence, and up to 32 blocked targets.
6//! During an active cycle (`cursor < cycle_end`), every undelivered row with
7//! `seq <= cursor` has its target in `blocked`. A completed marker
8//! (`cursor == cycle_end`) is exempt: the next fire resets the cursor and
9//! blocked set before scanning from the head. Target batches apply each
10//! target's rows in ascending sequence. If an older row is behind the active
11//! cursor, its target is blocked; if it lies ahead, ascending scanning reaches
12//! it first. Thus no later row for a target is applied while an older one is
13//! undelivered. The cycle end excludes new rows until the next cycle.
14//! Only delivery failures block targets; reaching the target budget pauses
15//! the cursor before the next target. Blocked targets are retried at each
16//! cycle start. While fewer than `MAX_BLOCKED_TARGETS` distinct failing
17//! targets precede it, every healthy target is eventually delivered: a fire
18//! that sees a deliverable row delivers at least one, delivered rows are
19//! deleted in the same guarded checkpoint, and fires without delivery back
20//! off. No closed-form fire bound is claimed; the throughput regressions in
21//! `tests.rs` pin fire counts for representative schedules. A mid-cycle
22//! append first waits for the next cycle. This can exceed
23//! `RELAY_LAG_BOUND_MS` in time; WP-1.23c's `namespace_relay_watermark`
24//! must tolerate that lag.
25
26mod content;
27mod deliver;
28mod enqueue;
29mod hook;
30pub use content::{ContentTakedownV1, HolderRelayHook, TakedownRequestTimer};
31
32pub use deliver::RelayHandler;
33pub use enqueue::{
34 RelayEnqueueSnapshot, commit_relay_rows, enqueue_relay_rows, relay_delivered_through,
35};
36pub(crate) use hook::AUDIT_CAPACITY;
37pub use hook::{NoHook, RelayHook};
38
39use crate::store::{NamespaceStore, Partition, StoreError, codec, keys};
40
41/// P-15: consumers allow a 60-second relay lag window.
42pub const RELAY_LAG_BOUND_MS: u64 = 60_000;
43
44/// Source work per fire. Targets are processed sequentially.
45#[derive(Debug, Clone, Copy)]
46#[non_exhaustive]
47pub struct RelayBudget {
48 /// Base row budget (default 256). A fire inspects at most four times this
49 /// many rows; selected targets may receive every row in that scan window.
50 pub max_rows: u32,
51 /// Maximum distinct targets attempted (default 16). Later targets pause
52 /// the scan without joining the failed-target blocked set.
53 pub max_targets: u32,
54 /// Optional maximum target store calls per target per fire, clamped to
55 /// at least two. Workers use two (one watermark get and one apply), and
56 /// total calls per fire are bounded by this times `max_targets`; native
57 /// defaults to no cap.
58 pub max_target_calls: Option<u32>,
59}
60/// Worker Paid relay targets attempted per fire.
61pub const WORKER_PAID_RELAY_TARGETS: u32 = 32;
62/// Worker Free relay targets attempted per fire.
63pub const WORKER_FREE_RELAY_TARGETS: u32 = 8;
64/// Worker Paid relay fires per alarm.
65pub const WORKER_PAID_RELAY_FIRES: u32 = 8;
66/// Worker Free relay fires per alarm.
67pub const WORKER_FREE_RELAY_FIRES: u32 = 2;
68/// Worker target calls allowed per target per fire.
69pub const WORKER_RELAY_CALLS_PER_TARGET: u32 = 2;
70impl Default for RelayBudget {
71 fn default() -> Self {
72 Self {
73 max_rows: 256,
74 max_targets: 16,
75 max_target_calls: None,
76 }
77 }
78}
79
80/// A commit-time lower bound for this source's undelivered outbox: every
81/// undelivered row committed at or after the returned time (+1).
82///
83/// Rows are in commit order (the `os` guard), but `at_ms` is the writer's
84/// plan-time reading and is **not** monotonic in seq: a writer may read its
85/// clock before retrying on `os`. The first row's `at_ms` still bounds every
86/// later row's commit time from below, because later rows commit after it.
87/// So the value is a valid lower bound, but it can move **backwards** as
88/// rows are delivered. A consumer (WP-1.23c's coordinator) keeps the running
89/// maximum it has seen, which stays safe for the same reason. The bound
90/// assumes the writer's clock is not ahead of the committing store's clock
91/// by more than the lease margin. Saturates at zero; an empty source
92/// reports `now_ms`.
93pub async fn relay_watermark<S: NamespaceStore>(
94 store: &S,
95 p: &Partition,
96 now_ms: u64,
97) -> Result<u64, StoreError> {
98 Ok(source_relay_state(store, p, now_ms).await?.0)
99}
100
101/// Source lower bound and whether its relay outbox is empty, in one scan.
102pub async fn source_relay_state<S: NamespaceStore>(
103 store: &S,
104 p: &Partition,
105 now_ms: u64,
106) -> Result<(u64, bool), StoreError> {
107 let (start, end) = keys::class_range(keys::TAG_RELAY);
108 let page = store.scan(p, &start, &end, None, 1).await?;
109 match page.entries.first() {
110 Some((key, value)) => {
111 if !matches!(keys::parse(key), Some(keys::ParsedKey::Relay(_))) {
112 return Err(StoreError::Corrupt("invalid relay key".into()));
113 }
114 Ok((codec::decode_relay(value)?.at_ms.saturating_sub(1), false))
115 }
116 None if page.next.is_none() => Ok((now_ms, true)),
117 None => Err(StoreError::Corrupt(
118 "relay scan returned an empty nonterminal page".into(),
119 )),
120 }
121}
122
123#[cfg(all(test, feature = "memory"))]
124mod tests;