1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
//! Source-side, at-least-once outbox delivery. Target watermarks make
//! redelivery idempotent, including a crash before source cleanup.
//!
//! Each source persists one relay scan cycle in `rs`: the `os` snapshot at
//! cycle start, the last inspected sequence, and up to 32 blocked targets.
//! During an active cycle (`cursor < cycle_end`), every undelivered row with
//! `seq <= cursor` has its target in `blocked`. A completed marker
//! (`cursor == cycle_end`) is exempt: the next fire resets the cursor and
//! blocked set before scanning from the head. Target batches apply each
//! target's rows in ascending sequence. If an older row is behind the active
//! cursor, its target is blocked; if it lies ahead, ascending scanning reaches
//! it first. Thus no later row for a target is applied while an older one is
//! undelivered. The cycle end excludes new rows until the next cycle.
//! Only delivery failures block targets; reaching the target budget pauses
//! the cursor before the next target. Blocked targets are retried at each
//! cycle start. While fewer than `MAX_BLOCKED_TARGETS` distinct failing
//! targets precede it, every healthy target is eventually delivered: a fire
//! that sees a deliverable row delivers at least one, delivered rows are
//! deleted in the same guarded checkpoint, and fires without delivery back
//! off. No closed-form fire bound is claimed; the throughput regressions in
//! `tests.rs` pin fire counts for representative schedules. A mid-cycle
//! append first waits for the next cycle. This can exceed
//! `RELAY_LAG_BOUND_MS` in time; WP-1.23c's `namespace_relay_watermark`
//! must tolerate that lag.
pub use ;
pub use RelayHandler;
pub use ;
pub use AUDIT_CAPACITY;
pub use ;
use crate;
/// P-15: consumers allow a 60-second relay lag window.
pub const RELAY_LAG_BOUND_MS: u64 = 60_000;
/// Source work per fire. Targets are processed sequentially.
/// Worker Paid relay targets attempted per fire.
pub const WORKER_PAID_RELAY_TARGETS: u32 = 32;
/// Worker Free relay targets attempted per fire.
pub const WORKER_FREE_RELAY_TARGETS: u32 = 8;
/// Worker Paid relay fires per alarm.
pub const WORKER_PAID_RELAY_FIRES: u32 = 8;
/// Worker Free relay fires per alarm.
pub const WORKER_FREE_RELAY_FIRES: u32 = 2;
/// Worker target calls allowed per target per fire.
pub const WORKER_RELAY_CALLS_PER_TARGET: u32 = 2;
/// A commit-time lower bound for this source's undelivered outbox: every
/// undelivered row committed at or after the returned time (+1).
///
/// Rows are in commit order (the `os` guard), but `at_ms` is the writer's
/// plan-time reading and is **not** monotonic in seq: a writer may read its
/// clock before retrying on `os`. The first row's `at_ms` still bounds every
/// later row's commit time from below, because later rows commit after it.
/// So the value is a valid lower bound, but it can move **backwards** as
/// rows are delivered. A consumer (WP-1.23c's coordinator) keeps the running
/// maximum it has seen, which stays safe for the same reason. The bound
/// assumes the writer's clock is not ahead of the committing store's clock
/// by more than the lease margin. Saturates at zero; an empty source
/// reports `now_ms`.
pub async
/// Source lower bound and whether its relay outbox is empty, in one scan.
pub async