1use crate::store::{
8 Batch, BatchOutcome, Cursor, Key, NamespaceStore, Partition, Precondition, ScanPage,
9 StoreError, codec, keys,
10};
11
12const PAGE_SIZE: u32 = 128;
13
14#[derive(Debug, thiserror::Error)]
16pub enum WatermarkError {
17 #[error("coordinator lease table is recovering")]
19 Recovering,
20 #[error(transparent)]
22 Store(#[from] StoreError),
23}
24
25#[derive(Debug, Clone, PartialEq, Eq)]
28pub struct WatermarkCheckpoint {
29 coordinator: Partition,
30 cursor: Cursor,
31 ceiling_ms: u64,
32 minimum_ms: u64,
33 recovery_generation: RecoveryGeneration,
34}
35
36pub type RecoveryGeneration = (Option<crate::store::Value>, Option<crate::store::Value>);
38
39fn encode_marker(bytes: &mut Vec<u8>, marker: Option<&crate::store::Value>) {
40 match marker {
41 Some(value) => {
42 bytes.extend_from_slice(
43 &u32::try_from(value.as_bytes().len())
44 .expect("store value length fits u32")
45 .to_be_bytes(),
46 );
47 bytes.extend_from_slice(value.as_bytes());
48 }
49 None => bytes.extend_from_slice(&u32::MAX.to_be_bytes()),
50 }
51}
52
53fn decode_marker(bytes: &[u8], at: &mut usize) -> Result<Option<crate::store::Value>, StoreError> {
54 let len = u32::from_be_bytes(
55 bytes
56 .get(*at..*at + 4)
57 .and_then(|slice| slice.try_into().ok())
58 .ok_or_else(|| StoreError::Invalid("truncated watermark generation".into()))?,
59 );
60 *at += 4;
61 if len == u32::MAX {
62 return Ok(None);
63 }
64 let end = at
65 .checked_add(
66 usize::try_from(len)
67 .map_err(|_| StoreError::Invalid("invalid watermark generation length".into()))?,
68 )
69 .ok_or_else(|| StoreError::Invalid("invalid watermark generation length".into()))?;
70 let value = bytes
71 .get(*at..end)
72 .ok_or_else(|| StoreError::Invalid("truncated watermark generation".into()))?;
73 *at = end;
74 Ok(Some(crate::store::Value::new(value.to_vec())))
75}
76
77impl WatermarkCheckpoint {
78 #[must_use]
84 pub fn encode(&self) -> Vec<u8> {
85 let partition = self
86 .coordinator
87 .encode()
88 .expect("coordinator identity encodes");
89 let mut bytes = Vec::with_capacity(27 + partition.len() + self.cursor.as_bytes().len());
90 bytes.push(2);
91 bytes.extend_from_slice(
92 &u16::try_from(partition.len())
93 .expect("partition length fits")
94 .to_be_bytes(),
95 );
96 bytes.extend_from_slice(&partition);
97 bytes.extend_from_slice(&self.ceiling_ms.to_be_bytes());
98 bytes.extend_from_slice(&self.minimum_ms.to_be_bytes());
99 encode_marker(&mut bytes, self.recovery_generation.0.as_ref());
100 encode_marker(&mut bytes, self.recovery_generation.1.as_ref());
101 bytes.extend_from_slice(self.cursor.as_bytes());
102 bytes
103 }
104
105 pub fn decode(bytes: &[u8]) -> Result<Self, StoreError> {
107 if bytes.len() < 28 || bytes[0] != 2 {
108 return Err(StoreError::Invalid("invalid watermark checkpoint".into()));
109 }
110 let part_len = usize::from(u16::from_be_bytes([bytes[1], bytes[2]]));
111 let end = 3usize
112 .checked_add(part_len)
113 .and_then(|n| n.checked_add(16))
114 .ok_or_else(|| StoreError::Invalid("invalid watermark checkpoint length".into()))?;
115 if bytes.len() < end + 9 {
116 return Err(StoreError::Invalid("truncated watermark checkpoint".into()));
117 }
118 let partition = bytes
119 .get(3..3 + part_len)
120 .ok_or_else(|| StoreError::Invalid("truncated coordinator identity".into()))?;
121 let coordinator = Partition::decode(partition)?;
122 let ceiling_ms = u64::from_be_bytes(
123 bytes
124 .get(end - 16..end - 8)
125 .and_then(|slice| slice.try_into().ok())
126 .ok_or_else(|| StoreError::Invalid("truncated watermark ceiling".into()))?,
127 );
128 let minimum_ms = u64::from_be_bytes(
129 bytes
130 .get(end - 8..end)
131 .and_then(|slice| slice.try_into().ok())
132 .ok_or_else(|| StoreError::Invalid("truncated watermark minimum".into()))?,
133 );
134 if minimum_ms > ceiling_ms {
135 return Err(StoreError::Invalid("invalid watermark minimum".into()));
136 }
137 let mut at = end;
138 let recovery = decode_marker(bytes, &mut at)?;
139 let reconcile = decode_marker(bytes, &mut at)?;
140 if bytes.len() <= at {
141 return Err(StoreError::Invalid("truncated watermark cursor".into()));
142 }
143 recovery
144 .as_ref()
145 .map(codec::decode_lease_recovery)
146 .transpose()?;
147 reconcile.as_ref().map(codec::decode_u64).transpose()?;
148 Ok(Self {
149 coordinator,
150 cursor: Cursor::new(bytes[at..].to_vec()),
151 ceiling_ms,
152 minimum_ms,
153 recovery_generation: (recovery, reconcile),
154 })
155 }
156}
157
158#[derive(Debug, Clone, PartialEq, Eq)]
160pub enum WatermarkStep {
161 Pending(WatermarkCheckpoint),
163 Complete(u64),
165}
166
167pub async fn check_recovery<S: NamespaceStore>(
169 store: &S,
170 coordinator: &Partition,
171 expected: Option<&RecoveryGeneration>,
172) -> Result<RecoveryGeneration, WatermarkError> {
173 if !matches!(
174 coordinator,
175 Partition::Coordinator(_) | Partition::Namespace(_)
176 ) {
177 return Err(WatermarkError::Store(StoreError::Invalid(
178 "watermark requires namespace coordinator".into(),
179 )));
180 }
181 let rows = store
182 .get_many(
183 coordinator,
184 &[keys::lease_recovery(), keys::lease_reconcile()],
185 )
186 .await?;
187 let [recovery, reconcile] = rows.as_slice() else {
188 return Err(WatermarkError::Store(StoreError::Corrupt(
189 "recovery get_many length".into(),
190 )));
191 };
192 let generation = (recovery.clone(), reconcile.clone());
193 if expected.is_some_and(|prior| *prior != generation) {
194 return Err(WatermarkError::Store(StoreError::Invalid(
195 "watermark checkpoint recovery generation changed".into(),
196 )));
197 }
198 let recovered = recovery
199 .as_ref()
200 .map(codec::decode_lease_recovery)
201 .transpose()?;
202 let reconciled = reconcile.as_ref().map(codec::decode_u64).transpose()?;
203 if recovered
204 .and_then(codec::LeaseRecovery::recovery_time)
205 .is_some_and(|resumed| reconciled.is_none_or(|at| at <= resumed))
206 {
207 return Err(WatermarkError::Recovering);
208 }
209 Ok(generation)
210}
211
212pub async fn mark_lease_table_reconciled<S: NamespaceStore>(
216 store: &S,
217 coordinator: &Partition,
218 at_ms: u64,
219) -> Result<(), StoreError> {
220 if !matches!(
221 coordinator,
222 Partition::Coordinator(_) | Partition::Namespace(_)
223 ) {
224 return Err(StoreError::Invalid(
225 "reconciliation requires namespace coordinator".into(),
226 ));
227 }
228 let key = keys::lease_recovery();
229 let value = store
230 .get(coordinator, &key)
231 .await?
232 .ok_or_else(|| StoreError::Invalid("no lease recovery to reconcile".into()))?;
233 let recovery = codec::decode_lease_recovery(&value)?;
234 let later = recovery
235 .recovery_time()
236 .ok_or_else(|| StoreError::Invalid("no actual lease recovery to reconcile".into()))?
237 .checked_add(1)
238 .ok_or_else(|| StoreError::Invalid("lease recovery time has no successor".into()))?;
239 let timestamp = at_ms.max(later);
240 let batch = Batch::new()
241 .require(Precondition::Equals(key, value))
242 .put(keys::lease_reconcile(), codec::encode_u64(timestamp));
243 match store.apply(coordinator, batch).await? {
244 BatchOutcome::Committed => Ok(()),
245 _ => Err(StoreError::Corrupt(
246 "lease reconciliation raced recovery".into(),
247 )),
248 }
249}
250
251fn decode_shard(
252 key: &Key,
253 value: &crate::store::Value,
254 coordinator: &Partition,
255) -> Result<(Partition, u64), StoreError> {
256 let Some(keys::ParsedKey::LeasedShard { repo, shard_ref }) = keys::parse(key) else {
257 return Err(StoreError::Corrupt("invalid leased shard key".into()));
258 };
259 let Partition::Coordinator(ns) = coordinator else {
260 return Err(StoreError::Invalid(
261 "watermark requires coordinator partition".into(),
262 ));
263 };
264 let row = codec::decode_leased_shard(value)?;
265 Ok((
266 Partition::Ref {
267 ns: ns.clone(),
268 repo,
269 shard_ref,
270 },
271 row.relay_watermark_ms,
272 ))
273}
274
275pub async fn namespace_relay_watermark_step<S: NamespaceStore>(
280 store: &S,
281 coordinator: &Partition,
282 now_ms: u64,
283 checkpoint: Option<WatermarkCheckpoint>,
284 limit: u32,
285) -> Result<WatermarkStep, WatermarkError> {
286 if !matches!(coordinator, Partition::Coordinator(_)) {
287 return Err(WatermarkError::Store(StoreError::Invalid(
288 "watermark scan requires coordinator partition".into(),
289 )));
290 }
291 let generation = check_recovery(
292 store,
293 coordinator,
294 checkpoint.as_ref().map(|c| &c.recovery_generation),
295 )
296 .await?;
297 let (start, end) = keys::class_range(keys::TAG_LEASED_SHARD);
298 let (cursor, ceiling_ms, mut minimum_ms) = match checkpoint {
299 Some(c) if c.coordinator == *coordinator => (Some(c.cursor), c.ceiling_ms, c.minimum_ms),
300 Some(_) => {
301 return Err(WatermarkError::Store(StoreError::Invalid(
302 "watermark checkpoint belongs to another coordinator".into(),
303 )));
304 }
305 None => (None, now_ms, now_ms),
306 };
307 let page = store
308 .scan(coordinator, &start, &end, cursor.as_ref(), limit.max(1))
309 .await?;
310 for (key, value) in &page.entries {
311 minimum_ms = minimum_ms.min(decode_shard(key, value, coordinator)?.1);
312 }
313 Ok(if let Some(cursor) = page.next {
314 WatermarkStep::Pending(WatermarkCheckpoint {
315 coordinator: coordinator.clone(),
316 cursor,
317 ceiling_ms,
318 minimum_ms,
319 recovery_generation: generation,
320 })
321 } else {
322 check_recovery(store, coordinator, Some(&generation)).await?;
323 WatermarkStep::Complete(minimum_ms)
324 })
325}
326
327pub async fn namespace_relay_watermark<S: NamespaceStore>(
332 store: &S,
333 coordinator: &Partition,
334 now_ms: u64,
335) -> Result<u64, WatermarkError> {
336 let mut checkpoint = None;
337 loop {
338 match namespace_relay_watermark_step(store, coordinator, now_ms, checkpoint, PAGE_SIZE)
339 .await?
340 {
341 WatermarkStep::Pending(next) => checkpoint = Some(next),
342 WatermarkStep::Complete(value) => return Ok(value),
343 }
344 }
345}
346
347#[derive(Debug, Clone, PartialEq, Eq)]
350pub struct ActiveShardsPage {
351 pub shards: Vec<Partition>,
353 pub next: Option<Cursor>,
355}
356
357pub async fn active_shards<S: NamespaceStore>(
359 store: &S,
360 coordinator: &Partition,
361 cursor: Option<&Cursor>,
362 limit: u32,
363) -> Result<ActiveShardsPage, WatermarkError> {
364 if !matches!(coordinator, Partition::Coordinator(_)) {
365 return Err(WatermarkError::Store(StoreError::Invalid(
366 "shard scan requires coordinator partition".into(),
367 )));
368 }
369 check_recovery(store, coordinator, None).await?;
370 let (start, end) = keys::class_range(keys::TAG_LEASED_SHARD);
371 let ScanPage { entries, next } = store
372 .scan(coordinator, &start, &end, cursor, limit.max(1))
373 .await?;
374 let shards = entries
375 .iter()
376 .map(|(key, value)| decode_shard(key, value, coordinator).map(|(p, _)| p))
377 .collect::<Result<_, _>>()?;
378 Ok(ActiveShardsPage { shards, next })
379}
380
381#[cfg(all(test, feature = "memory"))]
382mod tests {
383 use super::*;
384 use crate::memory::MemoryKv;
385 use crate::repo::{NamespaceKey, RepoName};
386 use crate::store::{Batch, BatchOutcome, PartitionStats, StoreCapabilities, Value};
387 use std::sync::{
388 Arc,
389 atomic::{AtomicBool, Ordering},
390 };
391
392 struct RecoverDuringScan {
393 inner: Arc<MemoryKv>,
394 once: AtomicBool,
395 }
396
397 impl NamespaceStore for RecoverDuringScan {
398 fn capabilities(&self) -> StoreCapabilities {
399 self.inner.capabilities()
400 }
401 async fn get(&self, p: &Partition, k: &Key) -> Result<Option<Value>, StoreError> {
402 self.inner.get(p, k).await
403 }
404 async fn get_many(
405 &self,
406 p: &Partition,
407 keys: &[Key],
408 ) -> Result<Vec<Option<Value>>, StoreError> {
409 self.inner.get_many(p, keys).await
410 }
411 async fn scan(
412 &self,
413 p: &Partition,
414 start: &Key,
415 end: &Key,
416 after: Option<&Cursor>,
417 limit: u32,
418 ) -> Result<ScanPage, StoreError> {
419 let page = self.inner.scan(p, start, end, after, limit).await?;
420 if self.once.swap(false, Ordering::SeqCst) {
421 self.inner
422 .apply(
423 p,
424 Batch::new().put(
425 keys::lease_recovery(),
426 codec::encode_lease_recovery(&codec::LeaseRecovery {
427 authority_fence: None,
428 authority_ready: None,
429 activation_only: None,
430 resumed_at_ms: 110,
431 }),
432 ),
433 )
434 .await?;
435 }
436 Ok(page)
437 }
438 async fn apply(&self, p: &Partition, batch: Batch) -> Result<BatchOutcome, StoreError> {
439 self.inner.apply(p, batch).await
440 }
441 async fn stats(&self, p: &Partition) -> Result<PartitionStats, StoreError> {
442 self.inner.stats(p).await
443 }
444 async fn probe(&self) -> Result<(), StoreError> {
445 self.inner.probe().await
446 }
447 }
448
449 fn coordinator() -> Partition {
450 Partition::Coordinator(NamespaceKey::deployment_default())
451 }
452 fn repo(n: u8) -> RepoName {
453 RepoName::new(format!("r{n}")).expect("test repository name")
454 }
455 fn source(n: u8) -> Partition {
456 Partition::Ref {
457 ns: NamespaceKey::deployment_default(),
458 repo: repo(n),
459 shard_ref: "refs/heads/main".into(),
460 }
461 }
462 fn lease(watermark: u64, expiry: u64) -> codec::LeasedShard {
463 codec::LeasedShard {
464 authority_generation: None,
465 acked_authority_generation: None,
466 epoch: 1,
467 expires_at_ms: expiry,
468 acked_epoch: 1,
469 relay_watermark_ms: watermark,
470 sweep_due_ms: expiry,
471 }
472 }
473 async fn put_lease(store: &MemoryKv, n: u8, watermark: u64, expiry: u64) {
474 store
475 .apply(
476 &coordinator(),
477 Batch::new().put(
478 keys::leased_shard(&repo(n), "refs/heads/main"),
479 codec::encode_leased_shard(&lease(watermark, expiry)),
480 ),
481 )
482 .await
483 .expect("insert test lease");
484 }
485 #[tokio::test]
486 async fn pages_and_inserts_preserve_scan_start_ceiling() {
487 let store = MemoryKv::default();
488 put_lease(&store, 0, 40, 200).await;
489 put_lease(&store, 2, 90, 200).await;
490 let WatermarkStep::Pending(checkpoint) =
491 namespace_relay_watermark_step(&store, &coordinator(), 100, None, 1)
492 .await
493 .unwrap()
494 else {
495 panic!("first page");
496 };
497 let encoded = checkpoint.encode();
498 let checkpoint = WatermarkCheckpoint::decode(&encoded).unwrap();
499 let other = Partition::Coordinator(NamespaceKey::from_stored("another".into()));
500 assert!(matches!(
501 namespace_relay_watermark_step(&store, &other, 150, Some(checkpoint.clone()), 1).await,
502 Err(WatermarkError::Store(StoreError::Invalid(_)))
503 ));
504 put_lease(&store, 1, 70, 200).await;
505 let WatermarkStep::Pending(checkpoint) =
506 namespace_relay_watermark_step(&store, &coordinator(), 150, Some(checkpoint), 1)
507 .await
508 .unwrap()
509 else {
510 panic!("second page");
511 };
512 let WatermarkStep::Complete(value) =
513 namespace_relay_watermark_step(&store, &coordinator(), 150, Some(checkpoint), 1)
514 .await
515 .unwrap()
516 else {
517 panic!("third page");
518 };
519 assert_eq!(value, 40);
520 let page = active_shards(&store, &coordinator(), None, 1)
521 .await
522 .unwrap();
523 assert_eq!(page.shards, vec![source(0)]);
524 assert!(page.next.is_some());
525 }
526
527 #[tokio::test]
528 async fn resume_rejects_a_new_recovery_generation() {
529 let store = MemoryKv::default();
530 put_lease(&store, 0, 40, 200).await;
531 put_lease(&store, 1, 90, 200).await;
532 store
533 .apply(
534 &coordinator(),
535 Batch::new()
536 .put(
537 keys::lease_recovery(),
538 codec::encode_lease_recovery(&codec::LeaseRecovery {
539 authority_fence: None,
540 authority_ready: None,
541 activation_only: None,
542 resumed_at_ms: 90,
543 }),
544 )
545 .put(keys::lease_reconcile(), codec::encode_u64(91)),
546 )
547 .await
548 .unwrap();
549 let WatermarkStep::Pending(checkpoint) =
550 namespace_relay_watermark_step(&store, &coordinator(), 100, None, 1)
551 .await
552 .unwrap()
553 else {
554 panic!("first page");
555 };
556 let checkpoint = WatermarkCheckpoint::decode(&checkpoint.encode()).unwrap();
557 store
558 .apply(
559 &coordinator(),
560 Batch::new()
561 .put(
562 keys::lease_recovery(),
563 codec::encode_lease_recovery(&codec::LeaseRecovery {
564 authority_fence: None,
565 authority_ready: None,
566 activation_only: None,
567 resumed_at_ms: 110,
568 }),
569 )
570 .put(keys::lease_reconcile(), codec::encode_u64(111)),
571 )
572 .await
573 .unwrap();
574 assert!(matches!(
575 namespace_relay_watermark_step(
576 &store,
577 &coordinator(),
578 120,
579 Some(checkpoint.clone()),
580 1
581 )
582 .await,
583 Err(WatermarkError::Store(StoreError::Invalid(_)))
584 ));
585 store
586 .apply(
587 &coordinator(),
588 Batch::new()
589 .put(
590 keys::lease_recovery(),
591 codec::encode_lease_recovery(&codec::LeaseRecovery {
592 authority_fence: None,
593 authority_ready: None,
594 activation_only: None,
595 resumed_at_ms: 130,
596 }),
597 )
598 .delete(keys::lease_reconcile()),
599 )
600 .await
601 .unwrap();
602 assert!(matches!(
603 namespace_relay_watermark_step(
604 &store,
605 &coordinator(),
606 140,
607 Some(checkpoint.clone()),
608 1
609 )
610 .await,
611 Err(WatermarkError::Store(StoreError::Invalid(_)))
612 ));
613 store
614 .apply(
615 &coordinator(),
616 Batch::new().put(keys::lease_recovery(), Value::new(&b"corrupt marker"[..])),
617 )
618 .await
619 .unwrap();
620 assert!(matches!(
621 namespace_relay_watermark_step(&store, &coordinator(), 150, Some(checkpoint), 1).await,
622 Err(WatermarkError::Store(StoreError::Invalid(_)))
623 ));
624 }
625
626 #[tokio::test]
627 async fn final_page_rechecks_recovery_generation() {
628 let inner = Arc::new(MemoryKv::default());
629 put_lease(inner.as_ref(), 0, 40, 200).await;
630 let store = RecoverDuringScan {
631 inner,
632 once: AtomicBool::new(true),
633 };
634 assert!(matches!(
635 namespace_relay_watermark_step(&store, &coordinator(), 100, None, 10).await,
636 Err(WatermarkError::Store(StoreError::Invalid(_)))
637 ));
638 }
639
640 #[tokio::test]
641 async fn a_new_row_behind_the_cursor_is_excluded_by_the_scan_ceiling() {
642 let store = MemoryKv::default();
643 put_lease(&store, 1, 80, 200).await;
644 put_lease(&store, 2, 90, 200).await;
645 let WatermarkStep::Pending(checkpoint) =
646 namespace_relay_watermark_step(&store, &coordinator(), 100, None, 1)
647 .await
648 .unwrap()
649 else {
650 panic!("first page");
651 };
652 put_lease(&store, 0, 0, 200).await;
653 let WatermarkStep::Complete(value) =
654 namespace_relay_watermark_step(&store, &coordinator(), 120, Some(checkpoint), 1)
655 .await
656 .unwrap()
657 else {
658 panic!("last page");
659 };
660 assert_eq!(value, 80);
661 assert!(value <= 100, "the scan-start ceiling remains binding");
662 }
663
664 #[tokio::test]
665 async fn corruption_and_recovery_fail_closed() {
666 let store = MemoryKv::default();
667 store
668 .apply(
669 &coordinator(),
670 Batch::new().put(
671 keys::leased_shard(&repo(0), "refs/heads/main"),
672 Value::new(&b"bad"[..]),
673 ),
674 )
675 .await
676 .unwrap();
677 assert!(matches!(
678 namespace_relay_watermark(&store, &coordinator(), 100).await,
679 Err(WatermarkError::Store(StoreError::Corrupt(_)))
680 ));
681 assert!(
682 active_shards(&store, &coordinator(), None, 10)
683 .await
684 .is_err()
685 );
686 store
687 .apply(
688 &coordinator(),
689 Batch::new().put(
690 keys::lease_recovery(),
691 codec::encode_lease_recovery(&codec::LeaseRecovery {
692 authority_fence: None,
693 authority_ready: None,
694 activation_only: None,
695 resumed_at_ms: 100,
696 }),
697 ),
698 )
699 .await
700 .unwrap();
701 assert!(matches!(
702 namespace_relay_watermark(&store, &coordinator(), 200).await,
703 Err(WatermarkError::Recovering)
704 ));
705 mark_lease_table_reconciled(&store, &coordinator(), 101)
706 .await
707 .unwrap();
708 assert!(matches!(
709 namespace_relay_watermark(&store, &coordinator(), 200).await,
710 Err(WatermarkError::Store(StoreError::Corrupt(_)))
711 ));
712 }
713}