1use crate::store::{
4 Batch, BatchOutcome, MAX_BATCH_BYTES, MAX_BATCH_OPS, NamespaceStore, Partition, Precondition,
5 StoreCapabilities, StoreError, Value, Write,
6 codec::{self, RelayV1},
7 keys,
8 outbox::MAX_RELAY_PUTS,
9};
10
11#[derive(Debug, Clone)]
14pub struct RelayEnqueueSnapshot {
15 pub sequence: Option<Value>,
17 pub source_lease: Option<Value>,
19 pub deadline_ms: u64,
21}
22
23fn guard(key: crate::Key, value: Option<&Value>) -> Precondition {
24 match value {
25 Some(value) => Precondition::Equals(key, value.clone()),
26 None => Precondition::Absent(key),
27 }
28}
29
30fn lease_lost() -> StoreError {
31 StoreError::unavailable(std::io::Error::other("source epoch lease lost; retry"))
32}
33
34fn batch_for(
35 prior: Option<&Value>,
36 lease: Option<&Value>,
37 seq: u64,
38 row: &RelayV1,
39 now_ms: u64,
40 deadline_ms: u64,
41) -> Result<Batch, StoreError> {
42 if row.puts.len().saturating_add(row.deletes.len()) > MAX_RELAY_PUTS {
43 return Err(StoreError::Invalid(
44 "relay row exceeds MAX_RELAY_PUTS".into(),
45 ));
46 }
47 let value = codec::encode_relay(row)?;
48 let mut batch = Batch::new()
49 .require(Precondition::NotAfter(deadline_ms))
50 .require(guard(keys::outbox_sequence(), prior));
51 if let Some(lease) = lease {
52 batch = batch.require(Precondition::Equals(keys::epoch_lease(), lease.clone()));
53 }
54 Ok(batch
55 .put(keys::outbox_sequence(), codec::encode_u64(seq))
56 .put(
57 keys::timer(now_ms, crate::timers::registry::kinds::RELAY.get(), b""),
58 Value::default(),
59 )
60 .put(keys::relay(seq), value))
61}
62
63pub fn enqueue_relay_rows(
68 os: &RelayEnqueueSnapshot,
69 source: &Partition,
70 rows: &[RelayV1],
71 now_ms: u64,
72) -> Result<Vec<Batch>, StoreError> {
73 if matches!(source, Partition::Ref { .. }) && os.source_lease.is_none() {
74 return Err(StoreError::Corrupt(
75 "D34 relay batch lacks source epoch lease".into(),
76 ));
77 }
78 if rows.is_empty() {
79 return Ok(Vec::new());
80 }
81 if let Some(raw) = &os.source_lease {
82 let lease = codec::decode_epoch_lease(raw)?;
83 if now_ms >= lease.expires_at_ms || os.deadline_ms >= lease.expires_at_ms {
84 return Err(lease_lost());
85 }
86 }
87 let mut seq = os
88 .sequence
89 .as_ref()
90 .map(codec::decode_u64)
91 .transpose()?
92 .unwrap_or(0);
93 let mut prior = os.sequence.clone();
94 let mut batches = Vec::new();
95 let mut index = 0;
96 while index < rows.len() {
97 seq = seq
98 .checked_add(1)
99 .ok_or_else(|| StoreError::Corrupt("outbox sequence overflow".into()))?;
100 let mut batch = batch_for(
101 prior.as_ref(),
102 os.source_lease.as_ref(),
103 seq,
104 &rows[index],
105 now_ms,
106 os.deadline_ms,
107 )?;
108 batch.validate(&StoreCapabilities::full())?;
109 index += 1;
110 while index < rows.len() && batch.preconditions.len() + batch.writes.len() < MAX_BATCH_OPS {
111 let next = seq
112 .checked_add(1)
113 .ok_or_else(|| StoreError::Corrupt("outbox sequence overflow".into()))?;
114 if rows[index]
115 .puts
116 .len()
117 .saturating_add(rows[index].deletes.len())
118 > MAX_RELAY_PUTS
119 {
120 return Err(StoreError::Invalid(
121 "relay row exceeds MAX_RELAY_PUTS".into(),
122 ));
123 }
124 let value = codec::encode_relay(&rows[index])?;
125 let key = keys::relay(next);
126 let mut candidate = batch.clone();
127 candidate.writes.push(Write::Put(key, value));
128 if candidate.validate(&StoreCapabilities::full()).is_err() {
129 break;
130 }
131 batch = candidate;
132 seq = next;
133 index += 1;
134 }
135 let next_os = codec::encode_u64(seq);
136 batch.writes[0] = Write::Put(keys::outbox_sequence(), next_os.clone());
138 debug_assert!(batch.preconditions.len() + batch.writes.len() <= MAX_BATCH_OPS);
139 debug_assert!(
140 batch
141 .writes
142 .iter()
143 .map(|w| match w {
144 Write::Put(k, v) => k.as_bytes().len() + v.as_bytes().len(),
145 Write::Delete(k) => k.as_bytes().len(),
146 })
147 .sum::<usize>()
148 <= MAX_BATCH_BYTES
149 );
150 batches.push(batch);
151 prior = Some(next_os);
152 }
153 Ok(batches)
154}
155
156pub async fn commit_relay_rows<S: NamespaceStore>(
160 store: &S,
161 source: &Partition,
162 rows: &[RelayV1],
163 now_ms: u64,
164 deadline_ms: u64,
165 source_lease: Option<&Value>,
166) -> Result<(), StoreError> {
167 if matches!(source, Partition::Ref { .. }) && source_lease.is_none() {
168 return Err(StoreError::Corrupt(
169 "D34 relay batch lacks source epoch lease".into(),
170 ));
171 }
172 let mut done = 0;
173 let mut losses = 0;
174 while done < rows.len() {
175 let snapshot = RelayEnqueueSnapshot {
176 sequence: store.get(source, &keys::outbox_sequence()).await?,
177 source_lease: source_lease.cloned(),
178 deadline_ms,
179 };
180 let batches = enqueue_relay_rows(&snapshot, source, &rows[done..], now_ms)?;
181 let mut lost = false;
182 for batch in batches {
183 let count = batch.writes.iter().filter(|w| matches!(w, Write::Put(key, _) if matches!(keys::parse(key), Some(keys::ParsedKey::Relay(_))))).count();
184 match store.apply(source, batch).await? {
185 BatchOutcome::Committed => done += count,
186 BatchOutcome::PreconditionFailed { index, .. }
187 if source_lease.is_some() && index == 2 =>
188 {
189 return Err(lease_lost());
190 }
191 BatchOutcome::PreconditionFailed { .. } => {
192 lost = true;
193 break;
194 }
195 BatchOutcome::DeadlinePassed { .. } => {
196 return Err(StoreError::unavailable(std::io::Error::other(
197 "relay enqueue deadline passed",
198 )));
199 }
200 }
201 }
202 if lost {
203 losses += 1;
204 if losses == 8 {
205 return Err(StoreError::unavailable(std::io::Error::other(
206 "relay enqueue contention",
207 )));
208 }
209 }
210 }
211 Ok(())
212}
213
214pub async fn relay_delivered_through<S: NamespaceStore>(
217 store: &S,
218 source: &Partition,
219 seq: u64,
220) -> Result<bool, StoreError> {
221 let allocated = store
222 .get(source, &keys::outbox_sequence())
223 .await?
224 .as_ref()
225 .map(codec::decode_u64)
226 .transpose()?
227 .unwrap_or(0);
228 if allocated < seq {
229 return Ok(false);
230 }
231 let (start, end) = keys::class_range(keys::TAG_RELAY);
232 let page = store.scan(source, &start, &end, None, 1).await?;
233 match page.entries.first().and_then(|(key, _)| keys::parse(key)) {
234 Some(keys::ParsedKey::Relay(first)) => Ok(first > seq),
235 None if page.entries.is_empty() => Ok(true),
236 _ => Err(StoreError::Corrupt("invalid relay head".into())),
237 }
238}
239
240#[cfg(all(test, feature = "memory"))]
241mod tests {
242 use super::*;
243 use crate::memory::MemoryKv;
244 use crate::repo::{NamespaceKey, RepoName};
245 use futures_executor::block_on;
246
247 fn source() -> Partition {
248 Partition::Namespace(NamespaceKey::deployment_default())
249 }
250 fn ref_source() -> Partition {
251 Partition::Ref {
252 ns: NamespaceKey::deployment_default(),
253 repo: RepoName::new("one").expect("valid repository name"),
254 shard_ref: "refs/heads/main".into(),
255 }
256 }
257 fn row(n: u16) -> RelayV1 {
258 RelayV1 {
259 at_ms: 1_000,
260 target: source(),
261 puts: vec![(
262 crate::Key::new(n.to_be_bytes().to_vec()),
263 Value::new(vec![u8::try_from(n % 256).unwrap_or(0)]),
264 )],
265 deletes: Vec::new(),
266 }
267 }
268
269 #[test]
270 fn chains_batches_and_detects_delivery() {
271 let store = MemoryKv::default();
272 let source = source();
273 let rows: Vec<_> = (0..120).map(row).collect();
274 let snapshot = RelayEnqueueSnapshot {
275 sequence: None,
276 source_lease: None,
277 deadline_ms: u64::MAX,
278 };
279 let batches = enqueue_relay_rows(&snapshot, &source, &rows, 1_000).unwrap();
280 assert!(batches.len() >= 2);
281 for batch in batches {
282 assert_eq!(
283 block_on(store.apply(&source, batch)).unwrap(),
284 BatchOutcome::Committed
285 );
286 }
287 assert!(!block_on(relay_delivered_through(&store, &source, 120)).unwrap());
288 let deletes = (1..=120).fold(Batch::new(), |batch, seq| batch.delete(keys::relay(seq)));
289 for chunk in deletes.writes.chunks(80) {
291 let batch = Batch {
292 preconditions: Vec::new(),
293 writes: chunk.to_vec(),
294 };
295 assert_eq!(
296 block_on(store.apply(&source, batch)).unwrap(),
297 BatchOutcome::Committed
298 );
299 }
300 assert!(block_on(relay_delivered_through(&store, &source, 120)).unwrap());
301 assert!(!block_on(relay_delivered_through(&store, &source, 121)).unwrap());
302 }
303
304 #[test]
305 fn stale_sequence_replans_uncommitted_rows() {
306 let store = MemoryKv::default();
307 let source = source();
308 let stale = RelayEnqueueSnapshot {
309 sequence: None,
310 source_lease: None,
311 deadline_ms: u64::MAX,
312 };
313 let batch = enqueue_relay_rows(&stale, &source, &[row(1)], 1_000)
314 .unwrap()
315 .remove(0);
316 assert_eq!(
317 block_on(store.apply(
318 &source,
319 Batch::new().put(keys::outbox_sequence(), codec::encode_u64(7))
320 ))
321 .unwrap(),
322 BatchOutcome::Committed
323 );
324 assert!(matches!(
325 block_on(store.apply(&source, batch)).unwrap(),
326 BatchOutcome::PreconditionFailed { .. }
327 ));
328 block_on(commit_relay_rows(
329 &store,
330 &source,
331 &[row(1)],
332 1_000,
333 u64::MAX,
334 None,
335 ))
336 .unwrap();
337 assert!(
338 block_on(store.get(&source, &keys::relay(8)))
339 .unwrap()
340 .is_some()
341 );
342 }
343
344 #[test]
345 fn expired_source_lease_cannot_enqueue() {
346 let lease = codec::encode_epoch_lease(&codec::EpochLease {
347 authority_ready: None,
348 authority_generation: None,
349 epoch: 1,
350 expires_at_ms: 1_010,
351 config_version: 1,
352 });
353 let snapshot = RelayEnqueueSnapshot {
354 sequence: None,
355 source_lease: Some(lease),
356 deadline_ms: 1_010,
357 };
358 assert!(enqueue_relay_rows(&snapshot, &source(), &[row(1)], 1_000).is_err());
359 }
360
361 #[test]
362 fn ref_source_requires_lease_and_a_lost_lease_is_not_contention() {
363 let source = ref_source();
364 let snapshot = RelayEnqueueSnapshot {
365 sequence: None,
366 source_lease: None,
367 deadline_ms: 2_000,
368 };
369 assert!(matches!(
370 enqueue_relay_rows(&snapshot, &source, &[row(1)], 1_000),
371 Err(StoreError::Corrupt(_))
372 ));
373 let stale = codec::encode_epoch_lease(&codec::EpochLease {
374 authority_ready: None,
375 authority_generation: None,
376 epoch: 1,
377 expires_at_ms: 10_000,
378 config_version: 1,
379 });
380 let current = codec::encode_epoch_lease(&codec::EpochLease {
381 authority_ready: None,
382 authority_generation: None,
383 epoch: 2,
384 expires_at_ms: 10_000,
385 config_version: 1,
386 });
387 let store = MemoryKv::with_clock(std::sync::Arc::new(crate::ManualClock::new(1_000)));
388 block_on(store.apply(&source, Batch::new().put(keys::epoch_lease(), current))).unwrap();
389 let error = block_on(commit_relay_rows(
390 &store,
391 &source,
392 &[row(1)],
393 1_000,
394 2_000,
395 Some(&stale),
396 ))
397 .unwrap_err();
398 assert!(error.to_string().contains("source epoch lease lost"));
399 }
400}