1use crate::checkpoint::{Checkpoint, CheckpointError, FrontierEntry};
97use crate::oplog::{Hlc, OpRecord, WallClock};
98use serde::{Deserialize, Serialize};
99use std::collections::BTreeMap;
100use std::fmt;
101use std::fs::{self, File, OpenOptions};
102use std::io::Write;
103use std::path::{Path, PathBuf};
104
105pub type Frontier = BTreeMap<String, u64>;
110
111pub fn frontier_of(ops: &[OpRecord]) -> Frontier {
113 let mut frontier = Frontier::new();
114 for op in ops {
115 let entry = frontier.entry(op.device_id.clone()).or_insert(op.seq);
116 if op.seq > *entry {
117 *entry = op.seq;
118 }
119 }
120 frontier
121}
122
123pub fn checkpoint_frontier(checkpoint: &Checkpoint) -> Frontier {
126 checkpoint
127 .frontier
128 .iter()
129 .map(|(device, entry)| (device.clone(), entry.seq))
130 .collect()
131}
132
133#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
136#[serde(rename_all = "snake_case")]
137pub enum DeviceStatus {
138 Active,
139 Evicted,
142}
143
144#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
146pub struct RosterEntry {
147 pub device_id: String,
148 pub added_at: Hlc,
150 pub last_seen_ms: u64,
152 pub acked: Option<Hlc>,
155 pub status: DeviceStatus,
156}
157
158#[derive(Debug, Clone, Default)]
160pub struct RelayConfig {
161 pub eviction_horizon_ms: Option<u64>,
165}
166
167#[derive(Debug, Clone, Default, PartialEq, Eq)]
170pub struct PushOutcome {
171 pub accepted: usize,
173 pub deduped: usize,
177}
178
179#[derive(Debug, Clone, PartialEq)]
181pub struct PullResult {
182 pub ops: Vec<OpRecord>,
185 pub latest_checkpoint: Option<String>,
188}
189
190#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
192pub struct AckOutcome {
193 pub advanced: bool,
195 pub reinstated: bool,
197}
198
199#[derive(Debug, Clone, Default, PartialEq, Eq)]
201pub struct GcReport {
202 pub dropped: BTreeMap<String, usize>,
203}
204
205impl GcReport {
206 pub fn total(&self) -> usize {
207 self.dropped.values().sum()
208 }
209}
210
211#[derive(Debug)]
213pub enum RelayError {
214 Chain { device_id: String, detail: String },
217 Fork { device_id: String, seq: u64 },
220 Gap { device_id: String, expected: u64, found: u64 },
223 ForeignOps { device_id: String, op_device: String },
225 FrontierTruncated { device_id: String, dropped_below: u64 },
230 Checkpoint(CheckpointError),
232 CheckpointFrontierUnverified { device_id: String, seq: u64, detail: String },
242 Io(std::io::Error),
243}
244
245impl fmt::Display for RelayError {
246 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
247 match self {
248 RelayError::Chain { device_id, detail } => {
249 write!(f, "relay push rejected for {device_id}: {detail}")
250 }
251 RelayError::Fork { device_id, seq } => write!(
252 f,
253 "relay push rejected: a different op already holds {device_id} seq {seq} — \
254 device chain fork"
255 ),
256 RelayError::Gap { device_id, expected, found } => write!(
257 f,
258 "relay push rejected for {device_id}: seq gap (relay expects {expected}, \
259 got {found}) — push contiguously"
260 ),
261 RelayError::ForeignOps { device_id, op_device } => write!(
262 f,
263 "relay push rejected: device {device_id} pushed an op emitted by {op_device} — \
264 a device pushes only its own chain"
265 ),
266 RelayError::FrontierTruncated { device_id, dropped_below } => write!(
267 f,
268 "pull frontier reaches into GC'd space (device {device_id}: ops below seq \
269 {dropped_below} were truncated) — cold-bootstrap from the latest checkpoint"
270 ),
271 RelayError::Checkpoint(e) => write!(f, "relay checkpoint rejected: {e}"),
272 RelayError::CheckpointFrontierUnverified { device_id, seq, detail } => write!(
273 f,
274 "relay checkpoint rejected: frontier for device {device_id} at seq {seq} does \
275 not match the relay-held chain ({detail}) — forged or foreign checkpoint, \
276 refusing to store it (it would become GC coverage for ops it does not cover)"
277 ),
278 RelayError::Io(e) => write!(f, "relay io error: {e}"),
279 }
280 }
281}
282
283impl std::error::Error for RelayError {}
284
285pub trait Relay {
289 fn register(&mut self, device_id: &str) -> Result<RosterEntry, RelayError>;
292 fn push(&mut self, device_id: &str, ops: &[OpRecord]) -> Result<PushOutcome, RelayError>;
296 fn pull(&mut self, device_id: &str, since: &Frontier) -> Result<PullResult, RelayError>;
299 fn ack(&mut self, device_id: &str, frontier: Hlc) -> Result<AckOutcome, RelayError>;
303 fn checkpoint_put(
307 &mut self,
308 device_id: &str,
309 checkpoint: &Checkpoint,
310 ) -> Result<bool, RelayError>;
311 fn checkpoint_get(&mut self) -> Result<Option<Checkpoint>, RelayError>;
313 fn roster(&mut self) -> Result<Vec<RosterEntry>, RelayError>;
315 fn stable_frontier(&mut self) -> Result<Option<Hlc>, RelayError>;
318 fn gc(&mut self) -> Result<GcReport, RelayError>;
321}
322
323#[derive(Debug, Clone, Default, Serialize, Deserialize)]
325struct DeviceChain {
326 ops: BTreeMap<u64, OpRecord>,
328 dropped_below: u64,
330 dropped_head: Option<(u64, String, Hlc)>,
333}
334
335impl DeviceChain {
336 fn head(&self) -> Option<(u64, &str, &Hlc)> {
339 self.ops
340 .iter()
341 .next_back()
342 .map(|(seq, op)| (*seq, op.op_id.as_str(), &op.hlc))
343 .or_else(|| {
344 self.dropped_head
345 .as_ref()
346 .map(|(seq, id, hlc)| (*seq, id.as_str(), hlc))
347 })
348 }
349
350 fn next_seq(&self) -> u64 {
351 self.head().map(|(seq, _, _)| seq + 1).unwrap_or(0)
352 }
353}
354
355#[derive(Debug, Default, Serialize, Deserialize)]
358struct RelayState {
359 roster: BTreeMap<String, RosterEntry>,
360 chains: BTreeMap<String, DeviceChain>,
361 #[serde(skip)]
365 checkpoints: BTreeMap<String, Checkpoint>,
366 latest_checkpoint: Option<String>,
368}
369
370pub(crate) fn frontier_dominates(a: &Checkpoint, b: &Checkpoint) -> bool {
374 b.frontier.iter().all(|(device, entry)| {
375 a.frontier
376 .get(device)
377 .is_some_and(|ae: &FrontierEntry| ae.seq >= entry.seq)
378 })
379}
380
381impl RelayState {
382 fn touch_and_sweep(&mut self, device_id: &str, now_ms: u64, config: &RelayConfig) {
386 let entry = self
387 .roster
388 .entry(device_id.to_string())
389 .or_insert_with(|| RosterEntry {
390 device_id: device_id.to_string(),
391 added_at: Hlc { wall_ms: now_ms, counter: 0, device_id: device_id.to_string() },
392 last_seen_ms: now_ms,
393 acked: None,
394 status: DeviceStatus::Active,
395 });
396 entry.last_seen_ms = now_ms;
397 if let Some(horizon) = config.eviction_horizon_ms {
398 for entry in self.roster.values_mut() {
399 if entry.status == DeviceStatus::Active
400 && now_ms.saturating_sub(entry.last_seen_ms) > horizon
401 {
402 entry.status = DeviceStatus::Evicted;
403 }
404 }
405 }
406 }
407
408 fn stable_frontier(&self) -> Option<Hlc> {
409 let active: Vec<&RosterEntry> = self
410 .roster
411 .values()
412 .filter(|e| e.status == DeviceStatus::Active)
413 .collect();
414 if active.is_empty() || active.iter().any(|e| e.acked.is_none()) {
415 return None;
416 }
417 active.iter().filter_map(|e| e.acked.clone()).min()
418 }
419
420 fn push(&mut self, device_id: &str, ops: &[OpRecord]) -> Result<PushOutcome, RelayError> {
421 let mut outcome = PushOutcome::default();
422 let mut sorted: Vec<&OpRecord> = ops.iter().collect();
424 sorted.sort_by_key(|op| op.seq);
425 for op in sorted {
426 if op.device_id != device_id {
427 return Err(RelayError::ForeignOps {
428 device_id: device_id.to_string(),
429 op_device: op.device_id.clone(),
430 });
431 }
432 if !op.id_valid() {
433 return Err(RelayError::Chain {
434 device_id: device_id.to_string(),
435 detail: format!("op {}: stored op_id does not match content", op.op_id),
436 });
437 }
438 if op.hlc.device_id != op.device_id {
439 return Err(RelayError::Chain {
440 device_id: device_id.to_string(),
441 detail: format!("op {}: hlc.device_id != device_id", op.op_id),
442 });
443 }
444 let chain = self.chains.entry(device_id.to_string()).or_default();
445 if op.seq < chain.dropped_below {
446 outcome.deduped += 1;
449 continue;
450 }
451 if let Some(existing) = chain.ops.get(&op.seq) {
452 if existing.op_id == op.op_id {
453 outcome.deduped += 1;
454 continue;
455 }
456 return Err(RelayError::Fork { device_id: device_id.to_string(), seq: op.seq });
457 }
458 let expected = chain.next_seq();
459 if op.seq != expected {
460 return Err(RelayError::Gap {
461 device_id: device_id.to_string(),
462 expected,
463 found: op.seq,
464 });
465 }
466 match chain.head() {
467 Some((_, head_id, head_hlc)) => {
468 if op.prev.as_deref() != Some(head_id) {
469 return Err(RelayError::Chain {
470 device_id: device_id.to_string(),
471 detail: format!(
472 "op {}: prev does not link the relay-held head {head_id}",
473 op.op_id
474 ),
475 });
476 }
477 if op.hlc <= *head_hlc {
478 return Err(RelayError::Chain {
479 device_id: device_id.to_string(),
480 detail: format!(
481 "op {}: hlc does not advance past the relay-held head",
482 op.op_id
483 ),
484 });
485 }
486 }
487 None => {
488 if op.prev.is_some() {
489 return Err(RelayError::Chain {
490 device_id: device_id.to_string(),
491 detail: format!("op {}: seq 0 must have no prev", op.op_id),
492 });
493 }
494 }
495 }
496 chain.ops.insert(op.seq, op.clone());
497 outcome.accepted += 1;
498 }
499 Ok(outcome)
500 }
501
502 fn pull(&self, since: &Frontier) -> Result<PullResult, RelayError> {
503 let mut ops = Vec::new();
504 for (device_id, chain) in &self.chains {
505 let start = since.get(device_id).map(|held| held + 1).unwrap_or(0);
506 if start < chain.dropped_below {
507 return Err(RelayError::FrontierTruncated {
508 device_id: device_id.clone(),
509 dropped_below: chain.dropped_below,
510 });
511 }
512 ops.extend(chain.ops.range(start..).map(|(_, op)| op.clone()));
513 }
514 ops.sort_by(|a, b| (&a.hlc, &a.op_id).cmp(&(&b.hlc, &b.op_id)));
515 Ok(PullResult { ops, latest_checkpoint: self.latest_checkpoint.clone() })
516 }
517
518 fn ack(&mut self, device_id: &str, frontier: Hlc) -> AckOutcome {
519 let others_frontier = {
522 let others: Vec<&RosterEntry> = self
523 .roster
524 .values()
525 .filter(|e| e.status == DeviceStatus::Active && e.device_id != device_id)
526 .collect();
527 if others.is_empty() || others.iter().any(|e| e.acked.is_none()) {
528 None
529 } else {
530 others.iter().filter_map(|e| e.acked.clone()).min()
531 }
532 };
533 let entry = self.roster.get_mut(device_id).expect("touched before ack");
534 let advanced = match &entry.acked {
535 Some(current) if frontier <= *current => false,
536 _ => {
537 entry.acked = Some(frontier);
538 true
539 }
540 };
541 let mut reinstated = false;
542 if entry.status == DeviceStatus::Evicted {
543 let caught_up = match (&entry.acked, &others_frontier) {
547 (Some(acked), Some(frontier)) => acked >= frontier,
548 (Some(_), None) => true,
549 (None, _) => false,
550 };
551 if caught_up {
552 entry.status = DeviceStatus::Active;
553 reinstated = true;
554 }
555 }
556 AckOutcome { advanced, reinstated }
557 }
558
559 fn validate_frontier(&self, checkpoint: &Checkpoint) -> Result<(), RelayError> {
567 for (device_id, entry) in &checkpoint.frontier {
568 let unverified = |detail: &str| RelayError::CheckpointFrontierUnverified {
569 device_id: device_id.clone(),
570 seq: entry.seq,
571 detail: detail.to_string(),
572 };
573 let Some(chain) = self.chains.get(device_id) else {
574 return Err(unverified("relay holds no chain for this device"));
575 };
576 if entry.seq >= chain.dropped_below {
577 match chain.ops.get(&entry.seq) {
579 Some(op) if op.op_id == entry.head => {}
580 Some(_) => {
581 return Err(unverified(
582 "frontier head does not match the relay-held op at this seq",
583 ))
584 }
585 None => {
586 return Err(unverified(
587 "relay holds no op at the claimed frontier seq (claims coverage \
588 beyond its chain head)",
589 ))
590 }
591 }
592 } else {
593 match &chain.dropped_head {
598 Some((seq, id, _)) if *seq == entry.seq => {
599 if id != &entry.head {
600 return Err(unverified(
601 "frontier head does not match the relay's GC'd dropped head",
602 ));
603 }
604 }
605 Some((seq, _, _)) if entry.seq < *seq => {}
606 _ => {
607 return Err(unverified(
608 "claimed frontier seq is below the relay's GC floor with no \
609 matching record",
610 ))
611 }
612 }
613 }
614 }
615 Ok(())
616 }
617
618 fn checkpoint_put(&mut self, checkpoint: &Checkpoint) -> Result<bool, RelayError> {
619 checkpoint.verify().map_err(RelayError::Checkpoint)?;
620 self.validate_frontier(checkpoint)?;
626 let stored = if self.checkpoints.contains_key(&checkpoint.checkpoint_hash) {
629 false
630 } else {
631 self.checkpoints
632 .insert(checkpoint.checkpoint_hash.clone(), checkpoint.clone());
633 true
634 };
635 let advance = match self
636 .latest_checkpoint
637 .as_ref()
638 .and_then(|hash| self.checkpoints.get(hash))
639 {
640 Some(current) => {
641 checkpoint.checkpoint_hash != current.checkpoint_hash
642 && frontier_dominates(checkpoint, current)
643 }
644 None => true,
645 };
646 if advance {
647 self.latest_checkpoint = Some(checkpoint.checkpoint_hash.clone());
648 }
649 Ok(stored)
650 }
651
652 fn checkpoint_get(&self) -> Option<Checkpoint> {
653 self.latest_checkpoint
654 .as_ref()
655 .and_then(|hash| self.checkpoints.get(hash))
656 .cloned()
657 }
658
659 fn gc(&mut self) -> GcReport {
660 let mut report = GcReport::default();
661 let Some(frontier) = self.stable_frontier() else {
662 return report;
663 };
664 let mut max_covered: BTreeMap<String, u64> = BTreeMap::new();
668 for ckpt in self.checkpoints.values() {
669 for (device, entry) in &ckpt.frontier {
670 let slot = max_covered.entry(device.clone()).or_insert(entry.seq);
671 if entry.seq > *slot {
672 *slot = entry.seq;
673 }
674 }
675 }
676 for (device_id, chain) in &mut self.chains {
677 let covered_through = max_covered.get(device_id).copied();
680 let mut droppable: Vec<u64> = Vec::new();
681 for (seq, op) in &chain.ops {
682 let below_frontier = op.hlc <= frontier;
683 let covered = covered_through.is_some_and(|through| through >= *seq);
684 if below_frontier && covered {
685 droppable.push(*seq);
686 } else {
687 break;
688 }
689 }
690 for seq in &droppable {
691 let op = chain.ops.remove(seq).expect("collected from the map");
692 chain.dropped_below = seq + 1;
693 chain.dropped_head = Some((*seq, op.op_id, op.hlc));
694 }
695 if !droppable.is_empty() {
696 report.dropped.insert(device_id.clone(), droppable.len());
697 }
698 }
699 report
700 }
701}
702
703pub struct InMemoryRelay {
705 state: RelayState,
706 config: RelayConfig,
707 wall: WallClock,
708}
709
710impl fmt::Debug for InMemoryRelay {
711 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
712 f.debug_struct("InMemoryRelay")
713 .field("state", &self.state)
714 .field("config", &self.config)
715 .finish_non_exhaustive()
716 }
717}
718
719impl InMemoryRelay {
720 pub fn new(config: RelayConfig, wall: WallClock) -> Self {
721 Self { state: RelayState::default(), config, wall }
722 }
723}
724
725impl Relay for InMemoryRelay {
726 fn register(&mut self, device_id: &str) -> Result<RosterEntry, RelayError> {
727 let now = (self.wall)();
728 self.state.touch_and_sweep(device_id, now, &self.config);
729 Ok(self.state.roster[device_id].clone())
730 }
731
732 fn push(&mut self, device_id: &str, ops: &[OpRecord]) -> Result<PushOutcome, RelayError> {
733 let now = (self.wall)();
734 self.state.touch_and_sweep(device_id, now, &self.config);
735 self.state.push(device_id, ops)
736 }
737
738 fn pull(&mut self, device_id: &str, since: &Frontier) -> Result<PullResult, RelayError> {
739 let now = (self.wall)();
740 self.state.touch_and_sweep(device_id, now, &self.config);
741 self.state.pull(since)
742 }
743
744 fn ack(&mut self, device_id: &str, frontier: Hlc) -> Result<AckOutcome, RelayError> {
745 let now = (self.wall)();
746 self.state.touch_and_sweep(device_id, now, &self.config);
747 Ok(self.state.ack(device_id, frontier))
748 }
749
750 fn checkpoint_put(
751 &mut self,
752 device_id: &str,
753 checkpoint: &Checkpoint,
754 ) -> Result<bool, RelayError> {
755 let now = (self.wall)();
756 self.state.touch_and_sweep(device_id, now, &self.config);
757 self.state.checkpoint_put(checkpoint)
758 }
759
760 fn checkpoint_get(&mut self) -> Result<Option<Checkpoint>, RelayError> {
761 Ok(self.state.checkpoint_get())
762 }
763
764 fn roster(&mut self) -> Result<Vec<RosterEntry>, RelayError> {
765 Ok(self.state.roster.values().cloned().collect())
766 }
767
768 fn stable_frontier(&mut self) -> Result<Option<Hlc>, RelayError> {
769 Ok(self.state.stable_frontier())
770 }
771
772 fn gc(&mut self) -> Result<GcReport, RelayError> {
773 Ok(self.state.gc())
774 }
775}
776
777pub struct FsRelay {
802 dir: PathBuf,
803 config: RelayConfig,
804 wall: WallClock,
805 checkpoint_cache: BTreeMap<String, Checkpoint>,
809}
810
811impl fmt::Debug for FsRelay {
812 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
813 f.debug_struct("FsRelay")
814 .field("dir", &self.dir)
815 .field("config", &self.config)
816 .finish_non_exhaustive()
817 }
818}
819
820impl FsRelay {
821 pub fn open(dir: &Path, config: RelayConfig, wall: WallClock) -> std::io::Result<Self> {
822 fs::create_dir_all(dir.join("checkpoints"))?;
823 Ok(Self {
824 dir: dir.to_path_buf(),
825 config,
826 wall,
827 checkpoint_cache: BTreeMap::new(),
828 })
829 }
830
831 fn state_path(&self) -> PathBuf {
832 self.dir.join("relay-state.json")
833 }
834
835 fn checkpoints_dir(&self) -> PathBuf {
836 self.dir.join("checkpoints")
837 }
838
839 fn with_state<T>(
843 &mut self,
844 f: impl FnOnce(&mut RelayState, u64, &RelayConfig) -> Result<T, RelayError>,
845 ) -> Result<T, RelayError> {
846 let lock = OpenOptions::new()
847 .read(true)
848 .write(true)
849 .create(true)
850 .truncate(false)
851 .open(self.dir.join("relay.lock"))
852 .map_err(RelayError::Io)?;
853 lock.lock().map_err(RelayError::Io)?; let mut state: RelayState = match fs::read_to_string(self.state_path()) {
856 Ok(raw) => serde_json::from_str(&raw)
857 .map_err(|e| RelayError::Io(std::io::Error::other(e)))?,
858 Err(e) if e.kind() == std::io::ErrorKind::NotFound => RelayState::default(),
859 Err(e) => return Err(RelayError::Io(e)),
860 };
861 for entry in fs::read_dir(self.checkpoints_dir()).map_err(RelayError::Io)? {
866 let path = entry.map_err(RelayError::Io)?.path();
867 let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
868 continue;
869 };
870 let Some(hash) = name.strip_suffix(".checkpoint.json") else {
871 continue;
872 };
873 let ckpt = match self.checkpoint_cache.get(hash) {
874 Some(cached) => cached.clone(),
875 None => {
876 let ckpt = Checkpoint::load(&path).map_err(RelayError::Checkpoint)?;
877 self.checkpoint_cache
878 .insert(ckpt.checkpoint_hash.clone(), ckpt.clone());
879 ckpt
880 }
881 };
882 state.checkpoints.insert(ckpt.checkpoint_hash.clone(), ckpt);
883 }
884
885 let now = (self.wall)();
886 let result = f(&mut state, now, &self.config);
887
888 for ckpt in state.checkpoints.values() {
896 let path = self.checkpoints_dir().join(ckpt.file_name());
897 if !path.exists() {
898 ckpt.save(&self.checkpoints_dir()).map_err(RelayError::Io)?;
899 }
900 self.checkpoint_cache
901 .entry(ckpt.checkpoint_hash.clone())
902 .or_insert_with(|| ckpt.clone());
903 }
904 let tmp = self.dir.join("relay-state.json.tmp");
905 {
906 let mut file = File::create(&tmp).map_err(RelayError::Io)?;
907 file.write_all(
908 serde_json::to_string(&state)
909 .map_err(|e| RelayError::Io(std::io::Error::other(e)))?
910 .as_bytes(),
911 )
912 .map_err(RelayError::Io)?;
913 file.sync_all().map_err(RelayError::Io)?;
914 }
915 fs::rename(&tmp, self.state_path()).map_err(RelayError::Io)?;
916 result
917 }
918}
919
920impl Relay for FsRelay {
921 fn register(&mut self, device_id: &str) -> Result<RosterEntry, RelayError> {
922 self.with_state(|state, now, config| {
923 state.touch_and_sweep(device_id, now, config);
924 Ok(state.roster[device_id].clone())
925 })
926 }
927
928 fn push(&mut self, device_id: &str, ops: &[OpRecord]) -> Result<PushOutcome, RelayError> {
929 self.with_state(|state, now, config| {
930 state.touch_and_sweep(device_id, now, config);
931 state.push(device_id, ops)
932 })
933 }
934
935 fn pull(&mut self, device_id: &str, since: &Frontier) -> Result<PullResult, RelayError> {
936 self.with_state(|state, now, config| {
937 state.touch_and_sweep(device_id, now, config);
938 state.pull(since)
939 })
940 }
941
942 fn ack(&mut self, device_id: &str, frontier: Hlc) -> Result<AckOutcome, RelayError> {
943 self.with_state(|state, now, config| {
944 state.touch_and_sweep(device_id, now, config);
945 Ok(state.ack(device_id, frontier))
946 })
947 }
948
949 fn checkpoint_put(
950 &mut self,
951 device_id: &str,
952 checkpoint: &Checkpoint,
953 ) -> Result<bool, RelayError> {
954 self.with_state(|state, now, config| {
955 state.touch_and_sweep(device_id, now, config);
956 state.checkpoint_put(checkpoint)
957 })
958 }
959
960 fn checkpoint_get(&mut self) -> Result<Option<Checkpoint>, RelayError> {
961 self.with_state(|state, _, _| Ok(state.checkpoint_get()))
962 }
963
964 fn roster(&mut self) -> Result<Vec<RosterEntry>, RelayError> {
965 self.with_state(|state, _, _| Ok(state.roster.values().cloned().collect()))
966 }
967
968 fn stable_frontier(&mut self) -> Result<Option<Hlc>, RelayError> {
969 self.with_state(|state, _, _| Ok(state.stable_frontier()))
970 }
971
972 fn gc(&mut self) -> Result<GcReport, RelayError> {
973 self.with_state(|state, _, _| Ok(state.gc()))
974 }
975}
976
977#[cfg(test)]
978mod tests {
979 use super::*;
980 use crate::oplog::{DeviceLog, Scope, Surface};
981 use std::sync::atomic::{AtomicU64, Ordering};
982 use std::sync::Arc;
983
984 fn manual_clock() -> (Arc<AtomicU64>, WallClock) {
985 let t = Arc::new(AtomicU64::new(0));
986 let reader = t.clone();
987 (t, Arc::new(move || reader.load(Ordering::SeqCst)))
988 }
989
990 fn mem_relay() -> InMemoryRelay {
991 InMemoryRelay::new(RelayConfig::default(), Arc::new(|| 0))
992 }
993
994 fn ops_for(device: &str, n: usize) -> (DeviceLog, Vec<OpRecord>) {
995 let mut log = DeviceLog::new(device);
996 let ops = (0..n)
997 .map(|i| {
998 log.append(
999 Scope::Personal,
1000 Surface::Knowledge,
1001 serde_json::json!({"id": format!("{device}-f{i}")}),
1002 )
1003 })
1004 .collect();
1005 (log, ops)
1006 }
1007
1008 #[test]
1009 fn push_validates_the_chain_and_dedups_retransmission() {
1010 let mut relay = mem_relay();
1011 let (mut log, ops) = ops_for("a", 3);
1012
1013 let outcome = relay.push("a", &ops).unwrap();
1014 assert_eq!(outcome, PushOutcome { accepted: 3, deduped: 0 });
1015
1016 let again = relay.push("a", &ops).unwrap();
1018 assert_eq!(again, PushOutcome { accepted: 0, deduped: 3 });
1019
1020 let next = log.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": "x"}));
1022 assert_eq!(relay.push("a", &[next]).unwrap().accepted, 1);
1023
1024 log.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": "skipped"}));
1026 let ahead = log.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": "y"}));
1027 assert!(matches!(
1028 relay.push("a", &[ahead]),
1029 Err(RelayError::Gap { expected: 4, found: 5, .. })
1030 ));
1031
1032 let mut forked = DeviceLog::new("a");
1034 let f0 = forked.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": "evil"}));
1035 assert!(matches!(relay.push("a", &[f0]), Err(RelayError::Fork { seq: 0, .. })));
1036
1037 let (_, b_ops) = ops_for("b", 1);
1039 assert!(matches!(
1040 relay.push("a", &b_ops),
1041 Err(RelayError::ForeignOps { .. })
1042 ));
1043
1044 let mut tampered = ops[0].clone();
1046 tampered.payload = serde_json::json!({"forged": true});
1047 assert!(matches!(relay.push("b", &[tampered]), Err(RelayError::ForeignOps { .. })));
1048 let mut own_tampered = ops[0].clone();
1049 own_tampered.payload = serde_json::json!({"id": "a-f0", "forged": true});
1050 assert!(matches!(relay.push("a", &[own_tampered]), Err(RelayError::Chain { .. })));
1051 }
1052
1053 #[test]
1054 fn pull_is_a_seq_cursor_and_serves_canonical_order() {
1055 let mut relay = mem_relay();
1056 let (_, a_ops) = ops_for("a", 3);
1057 let (_, b_ops) = ops_for("b", 2);
1058 relay.push("a", &a_ops).unwrap();
1059 relay.push("b", &b_ops).unwrap();
1060
1061 let all = relay.pull("c", &Frontier::new()).unwrap();
1063 assert_eq!(all.ops.len(), 5);
1064 assert!(all.latest_checkpoint.is_none());
1065
1066 let mut since = Frontier::new();
1068 since.insert("a".into(), 1);
1069 since.insert("b".into(), 1);
1070 let tail = relay.pull("c", &since).unwrap();
1071 assert_eq!(tail.ops.len(), 1);
1072 assert_eq!(tail.ops[0].seq, 2);
1073 assert_eq!(tail.ops[0].device_id, "a");
1074 }
1075
1076 #[test]
1077 fn stable_frontier_requires_every_active_device_acked() {
1078 let mut relay = mem_relay();
1079 let (_, a_ops) = ops_for("a", 2);
1080 relay.push("a", &a_ops).unwrap();
1081 relay.register("b").unwrap();
1082
1083 relay.ack("a", a_ops[1].hlc.clone()).unwrap();
1086 assert_eq!(relay.stable_frontier().unwrap(), None);
1087
1088 relay.ack("b", a_ops[0].hlc.clone()).unwrap();
1089 assert_eq!(relay.stable_frontier().unwrap(), Some(a_ops[0].hlc.clone()), "min(acked)");
1090
1091 let outcome = relay.ack("b", a_ops[0].hlc.clone()).unwrap();
1093 assert!(!outcome.advanced);
1094 let outcome = relay.ack("b", a_ops[1].hlc.clone()).unwrap();
1095 assert!(outcome.advanced);
1096 assert_eq!(relay.stable_frontier().unwrap(), Some(a_ops[1].hlc.clone()));
1097 }
1098
1099 #[test]
1100 fn gc_requires_both_frontier_and_covering_checkpoint() {
1101 let mut relay = mem_relay();
1102 let (_, ops) = ops_for("a", 4);
1103 relay.push("a", &ops).unwrap();
1104 relay.ack("a", ops[3].hlc.clone()).unwrap();
1105
1106 assert_eq!(relay.gc().unwrap().total(), 0);
1109
1110 let ckpt = Checkpoint::from_ops(&ops[..2]).unwrap();
1112 assert!(relay.checkpoint_put("a", &ckpt).unwrap());
1113 let report = relay.gc().unwrap();
1114 assert_eq!(report.dropped["a"], 2);
1115
1116 let mut log = DeviceLog::resume("a", &ops).unwrap();
1118 let newer = log.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": "n"}));
1119 relay.push("a", std::slice::from_ref(&newer)).unwrap();
1120 let full_ckpt = {
1121 let mut all = ops.clone();
1122 all.push(newer);
1123 Checkpoint::from_ops(&all).unwrap()
1124 };
1125 relay.checkpoint_put("a", &full_ckpt).unwrap();
1126 let report = relay.gc().unwrap();
1129 assert_eq!(report.dropped["a"], 2);
1130 let survivors = relay.pull("b", &Frontier::new());
1131 assert!(matches!(
1133 survivors,
1134 Err(RelayError::FrontierTruncated { dropped_below: 4, .. })
1135 ));
1136 let mut since = Frontier::new();
1138 since.insert("a".into(), 3);
1139 assert_eq!(relay.pull("b", &since).unwrap().ops.len(), 1);
1140 }
1141
1142 #[test]
1143 fn gc_preserves_chain_continuity_for_later_pushes() {
1144 let mut relay = mem_relay();
1145 let (mut log, ops) = ops_for("a", 3);
1146 relay.push("a", &ops).unwrap();
1147 relay.ack("a", ops[2].hlc.clone()).unwrap();
1148 let ckpt = Checkpoint::from_ops(&ops).unwrap();
1149 relay.checkpoint_put("a", &ckpt).unwrap();
1150 assert_eq!(relay.gc().unwrap().dropped["a"], 3);
1151
1152 let next = log.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": "n"}));
1155 assert_eq!(relay.push("a", &[next]).unwrap().accepted, 1);
1156 assert_eq!(relay.push("a", &ops).unwrap(), PushOutcome { accepted: 0, deduped: 3 });
1158 }
1159
1160 #[test]
1161 fn forged_checkpoint_frontier_is_rejected_and_never_becomes_gc_coverage() {
1162 let mut relay = mem_relay();
1168
1169 let mut a = DeviceLog::new("a");
1171 let real: Vec<OpRecord> = (0..3)
1172 .map(|i| {
1173 a.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": format!("real-{i}")}))
1174 })
1175 .collect();
1176 relay.push("a", &real).unwrap();
1177 relay.ack("a", real[2].hlc.clone()).unwrap();
1178
1179 let mut fake = DeviceLog::new("a");
1182 let other: Vec<OpRecord> = (0..3)
1183 .map(|i| {
1184 fake.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": format!("other-{i}")}))
1185 })
1186 .collect();
1187 let forged = Checkpoint::from_ops(&other).unwrap();
1188
1189 assert!(matches!(
1192 relay.checkpoint_put("a", &forged),
1193 Err(RelayError::CheckpointFrontierUnverified { device_id, .. }) if device_id == "a"
1194 ));
1195
1196 assert_eq!(relay.gc().unwrap().total(), 0, "no forged coverage, no data loss");
1201 assert!(relay.checkpoint_get().unwrap().is_none());
1202 let served = relay.pull("a", &Frontier::new()).unwrap();
1204 assert_eq!(served.ops.len(), 3);
1205 assert!(served.ops.iter().all(|op| op.payload["id"].as_str().unwrap().starts_with("real-")));
1206
1207 let mut a2 = DeviceLog::resume("a", &real).unwrap();
1210 let mut ahead = real.clone();
1211 ahead.push(a2.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": "real-3"})));
1212 let claims_beyond = Checkpoint::from_ops(&ahead).unwrap(); assert!(matches!(
1214 relay.checkpoint_put("a", &claims_beyond),
1215 Err(RelayError::CheckpointFrontierUnverified { seq: 3, .. })
1216 ));
1217
1218 let mut c = DeviceLog::new("c");
1221 let c_ops: Vec<OpRecord> = (0..2)
1222 .map(|i| c.append(Scope::Personal, Surface::Knowledge, serde_json::json!({"id": format!("c{i}")})))
1223 .collect();
1224 let unknown_dev = Checkpoint::from_ops(&c_ops).unwrap();
1225 assert!(matches!(
1226 relay.checkpoint_put("a", &unknown_dev),
1227 Err(RelayError::CheckpointFrontierUnverified { device_id, .. }) if device_id == "c"
1228 ));
1229
1230 let honest = Checkpoint::from_ops(&real[..2]).unwrap();
1233 assert!(relay.checkpoint_put("a", &honest).unwrap());
1234 assert_eq!(relay.gc().unwrap().dropped["a"], 2);
1235 }
1236
1237 #[test]
1238 fn checkpoint_dedup_keys_on_whole_record_address_not_state_hash() {
1239 let mut a = DeviceLog::new("a");
1245 let mut b = DeviceLog::new("b");
1246 let fact = serde_json::json!({"id": "f1", "body": "hi"});
1247 let oa = a.append(Scope::Personal, Surface::Knowledge, fact.clone());
1248 let ob = b.append(Scope::Personal, Surface::Knowledge, fact);
1249 let just_a = Checkpoint::from_ops(std::slice::from_ref(&oa)).unwrap();
1250 let both = Checkpoint::from_ops(&[oa.clone(), ob.clone()]).unwrap();
1251 assert_eq!(just_a.state_hash, both.state_hash, "the cross-device dedup collision");
1252
1253 let mut relay = mem_relay();
1254 relay.push("a", std::slice::from_ref(&oa)).unwrap();
1257 relay.push("b", std::slice::from_ref(&ob)).unwrap();
1258 assert!(relay.checkpoint_put("a", &just_a).unwrap());
1259 assert!(relay.checkpoint_put("b", &both).unwrap(), "same state_hash is NOT a dedup");
1260 assert!(!relay.checkpoint_put("a", &just_a).unwrap(), "same checkpoint_hash IS");
1261
1262 assert_eq!(
1265 relay.checkpoint_get().unwrap().unwrap().checkpoint_hash,
1266 both.checkpoint_hash
1267 );
1268 relay.checkpoint_put("a", &just_a).unwrap();
1269 assert_eq!(
1270 relay.checkpoint_get().unwrap().unwrap().checkpoint_hash,
1271 both.checkpoint_hash,
1272 "dominance-monotone pointer"
1273 );
1274
1275 let mut forged = both.clone();
1277 forged.state.logs.clear();
1278 assert!(matches!(
1279 relay.checkpoint_put("b", &forged),
1280 Err(RelayError::Checkpoint(CheckpointError::HashMismatch { .. }))
1281 ));
1282 }
1283
1284 #[test]
1285 fn eviction_unpins_the_frontier_and_ack_reinstates() {
1286 let (t, wall) = manual_clock();
1287 let mut relay = InMemoryRelay::new(
1288 RelayConfig { eviction_horizon_ms: Some(1_000) },
1289 wall,
1290 );
1291 let (_, a_ops) = ops_for("a", 2);
1292 relay.push("a", &a_ops).unwrap();
1293
1294 t.store(100, Ordering::SeqCst);
1295 relay.ack("a", a_ops[1].hlc.clone()).unwrap();
1296 relay.ack("c", a_ops[0].hlc.clone()).unwrap(); assert_eq!(relay.stable_frontier().unwrap(), Some(a_ops[0].hlc.clone()));
1298
1299 t.store(2_000, Ordering::SeqCst);
1301 relay.register("a").unwrap();
1302 let roster: BTreeMap<String, RosterEntry> = relay
1303 .roster()
1304 .unwrap()
1305 .into_iter()
1306 .map(|e| (e.device_id.clone(), e))
1307 .collect();
1308 assert_eq!(roster["c"].status, DeviceStatus::Evicted);
1309 assert_eq!(roster["a"].status, DeviceStatus::Active);
1310 assert_eq!(
1311 relay.stable_frontier().unwrap(),
1312 Some(a_ops[1].hlc.clone()),
1313 "the evicted device's ack no longer holds the frontier"
1314 );
1315
1316 let outcome = relay.ack("c", a_ops[0].hlc.clone()).unwrap();
1319 assert!(!outcome.advanced && !outcome.reinstated);
1320 assert_eq!(relay.stable_frontier().unwrap(), Some(a_ops[1].hlc.clone()));
1321
1322 let outcome = relay.ack("c", a_ops[1].hlc.clone()).unwrap();
1324 assert!(outcome.advanced && outcome.reinstated);
1325 let roster: BTreeMap<String, RosterEntry> = relay
1326 .roster()
1327 .unwrap()
1328 .into_iter()
1329 .map(|e| (e.device_id.clone(), e))
1330 .collect();
1331 assert_eq!(roster["c"].status, DeviceStatus::Active);
1332 }
1333
1334 #[test]
1335 fn fs_relay_matches_in_memory_semantics_and_persists() {
1336 let dir = tempfile::tempdir().unwrap();
1337 let (_, ops) = ops_for("a", 3);
1338 {
1339 let mut relay =
1340 FsRelay::open(dir.path(), RelayConfig::default(), Arc::new(|| 7)).unwrap();
1341 assert_eq!(relay.push("a", &ops).unwrap().accepted, 3);
1342 relay.ack("a", ops[2].hlc.clone()).unwrap();
1343 let ckpt = Checkpoint::from_ops(&ops[..2]).unwrap();
1344 assert!(relay.checkpoint_put("a", &ckpt).unwrap());
1345 }
1346 let mut relay =
1348 FsRelay::open(dir.path(), RelayConfig::default(), Arc::new(|| 8)).unwrap();
1349 assert_eq!(relay.stable_frontier().unwrap(), Some(ops[2].hlc.clone()));
1350 assert_eq!(relay.pull("b", &Frontier::new()).unwrap().ops, ops);
1353 assert_eq!(relay.stable_frontier().unwrap(), None);
1354 relay.ack("b", ops[2].hlc.clone()).unwrap();
1355 assert_eq!(relay.push("a", &ops).unwrap(), PushOutcome { accepted: 0, deduped: 3 });
1356 let report = relay.gc().unwrap();
1357 assert_eq!(report.dropped["a"], 2);
1358 let roster = relay.roster().unwrap();
1359 assert_eq!(roster.len(), 2);
1360 assert_eq!(roster[0].added_at.wall_ms, 7, "roster added_at survives restart");
1361
1362 let ckpt = relay.checkpoint_get().unwrap().unwrap();
1365 assert_eq!(ckpt.frontier["a"].seq, 1);
1366
1367 let path = dir
1374 .path()
1375 .join("checkpoints")
1376 .join(format!("{}.checkpoint.json", ckpt.checkpoint_hash));
1377 let raw = fs::read_to_string(&path).unwrap();
1378 fs::write(&path, raw.replace("\"a-f0\"", "\"a-f0-forged\"")).unwrap();
1379 let mut fresh =
1380 FsRelay::open(dir.path(), RelayConfig::default(), Arc::new(|| 9)).unwrap();
1381 assert!(matches!(
1382 fresh.checkpoint_get(),
1383 Err(RelayError::Checkpoint(_))
1384 ));
1385 }
1386}