use crate::checkpoint::{Checkpoint, CheckpointError, FrontierEntry};
use crate::oplog::{Hlc, OpRecord, WallClock};
use serde::{Deserialize, Serialize};
use std::collections::BTreeMap;
use std::fmt;
use std::fs::{self, File, OpenOptions};
use std::io::Write;
use std::path::{Path, PathBuf};
pub type Frontier = BTreeMap<String, u64>;
pub fn frontier_of(ops: &[OpRecord]) -> Frontier {
let mut frontier = Frontier::new();
for op in ops {
let entry = frontier.entry(op.device_id.clone()).or_insert(op.seq);
if op.seq > *entry {
*entry = op.seq;
}
}
frontier
}
pub fn checkpoint_frontier(checkpoint: &Checkpoint) -> Frontier {
checkpoint
.frontier
.iter()
.map(|(device, entry)| (device.clone(), entry.seq))
.collect()
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum DeviceStatus {
Active,
Evicted,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct RosterEntry {
pub device_id: String,
pub added_at: Hlc,
pub last_seen_ms: u64,
pub acked: Option<Hlc>,
pub status: DeviceStatus,
}
#[derive(Debug, Clone, Default)]
pub struct RelayConfig {
pub eviction_horizon_ms: Option<u64>,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct PushOutcome {
pub accepted: usize,
pub deduped: usize,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct PullResult {
pub ops: Vec<OpRecord>,
pub latest_checkpoint: Option<String>,
}
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct AckOutcome {
pub advanced: bool,
pub reinstated: bool,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
pub struct GcReport {
pub dropped: BTreeMap<String, usize>,
}
impl GcReport {
pub fn total(&self) -> usize {
self.dropped.values().sum()
}
}
#[derive(Debug)]
pub enum RelayError {
Chain {
device_id: String,
detail: String,
},
Fork {
device_id: String,
seq: u64,
},
Gap {
device_id: String,
expected: u64,
found: u64,
},
ForeignOps {
device_id: String,
op_device: String,
},
FrontierTruncated {
device_id: String,
dropped_below: u64,
},
Checkpoint(CheckpointError),
CheckpointFrontierUnverified {
device_id: String,
seq: u64,
detail: String,
},
Io(std::io::Error),
}
impl fmt::Display for RelayError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
RelayError::Chain { device_id, detail } => {
write!(f, "relay push rejected for {device_id}: {detail}")
}
RelayError::Fork { device_id, seq } => write!(
f,
"relay push rejected: a different op already holds {device_id} seq {seq} — \
device chain fork"
),
RelayError::Gap {
device_id,
expected,
found,
} => write!(
f,
"relay push rejected for {device_id}: seq gap (relay expects {expected}, \
got {found}) — push contiguously"
),
RelayError::ForeignOps {
device_id,
op_device,
} => write!(
f,
"relay push rejected: device {device_id} pushed an op emitted by {op_device} — \
a device pushes only its own chain"
),
RelayError::FrontierTruncated {
device_id,
dropped_below,
} => write!(
f,
"pull frontier reaches into GC'd space (device {device_id}: ops below seq \
{dropped_below} were truncated) — cold-bootstrap from the latest checkpoint"
),
RelayError::Checkpoint(e) => write!(f, "relay checkpoint rejected: {e}"),
RelayError::CheckpointFrontierUnverified {
device_id,
seq,
detail,
} => write!(
f,
"relay checkpoint rejected: frontier for device {device_id} at seq {seq} does \
not match the relay-held chain ({detail}) — forged or foreign checkpoint, \
refusing to store it (it would become GC coverage for ops it does not cover)"
),
RelayError::Io(e) => write!(f, "relay io error: {e}"),
}
}
}
impl std::error::Error for RelayError {}
pub trait Relay {
fn register(&mut self, device_id: &str) -> Result<RosterEntry, RelayError>;
fn push(&mut self, device_id: &str, ops: &[OpRecord]) -> Result<PushOutcome, RelayError>;
fn pull(&mut self, device_id: &str, since: &Frontier) -> Result<PullResult, RelayError>;
fn ack(&mut self, device_id: &str, frontier: Hlc) -> Result<AckOutcome, RelayError>;
fn checkpoint_put(
&mut self,
device_id: &str,
checkpoint: &Checkpoint,
) -> Result<bool, RelayError>;
fn checkpoint_get(&mut self) -> Result<Option<Checkpoint>, RelayError>;
fn roster(&mut self) -> Result<Vec<RosterEntry>, RelayError>;
fn stable_frontier(&mut self) -> Result<Option<Hlc>, RelayError>;
fn gc(&mut self) -> Result<GcReport, RelayError>;
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
struct DeviceChain {
ops: BTreeMap<u64, OpRecord>,
dropped_below: u64,
dropped_head: Option<(u64, String, Hlc)>,
}
impl DeviceChain {
fn head(&self) -> Option<(u64, &str, &Hlc)> {
self.ops
.iter()
.next_back()
.map(|(seq, op)| (*seq, op.op_id.as_str(), &op.hlc))
.or_else(|| {
self.dropped_head
.as_ref()
.map(|(seq, id, hlc)| (*seq, id.as_str(), hlc))
})
}
fn next_seq(&self) -> u64 {
self.head().map(|(seq, _, _)| seq + 1).unwrap_or(0)
}
}
#[derive(Debug, Default, Serialize, Deserialize)]
struct RelayState {
roster: BTreeMap<String, RosterEntry>,
chains: BTreeMap<String, DeviceChain>,
#[serde(skip)]
checkpoints: BTreeMap<String, Checkpoint>,
latest_checkpoint: Option<String>,
}
pub(crate) fn frontier_dominates(a: &Checkpoint, b: &Checkpoint) -> bool {
b.frontier.iter().all(|(device, entry)| {
a.frontier
.get(device)
.is_some_and(|ae: &FrontierEntry| ae.seq >= entry.seq)
})
}
impl RelayState {
fn touch_and_sweep(&mut self, device_id: &str, now_ms: u64, config: &RelayConfig) {
let entry = self
.roster
.entry(device_id.to_string())
.or_insert_with(|| RosterEntry {
device_id: device_id.to_string(),
added_at: Hlc {
wall_ms: now_ms,
counter: 0,
device_id: device_id.to_string(),
},
last_seen_ms: now_ms,
acked: None,
status: DeviceStatus::Active,
});
entry.last_seen_ms = now_ms;
if let Some(horizon) = config.eviction_horizon_ms {
for entry in self.roster.values_mut() {
if entry.status == DeviceStatus::Active
&& now_ms.saturating_sub(entry.last_seen_ms) > horizon
{
entry.status = DeviceStatus::Evicted;
}
}
}
}
fn stable_frontier(&self) -> Option<Hlc> {
let active: Vec<&RosterEntry> = self
.roster
.values()
.filter(|e| e.status == DeviceStatus::Active)
.collect();
if active.is_empty() || active.iter().any(|e| e.acked.is_none()) {
return None;
}
active.iter().filter_map(|e| e.acked.clone()).min()
}
fn push(&mut self, device_id: &str, ops: &[OpRecord]) -> Result<PushOutcome, RelayError> {
let mut outcome = PushOutcome::default();
let mut sorted: Vec<&OpRecord> = ops.iter().collect();
sorted.sort_by_key(|op| op.seq);
for op in sorted {
if op.device_id != device_id {
return Err(RelayError::ForeignOps {
device_id: device_id.to_string(),
op_device: op.device_id.clone(),
});
}
if !op.id_valid() {
return Err(RelayError::Chain {
device_id: device_id.to_string(),
detail: format!("op {}: stored op_id does not match content", op.op_id),
});
}
if op.hlc.device_id != op.device_id {
return Err(RelayError::Chain {
device_id: device_id.to_string(),
detail: format!("op {}: hlc.device_id != device_id", op.op_id),
});
}
let chain = self.chains.entry(device_id.to_string()).or_default();
if op.seq < chain.dropped_below {
outcome.deduped += 1;
continue;
}
if let Some(existing) = chain.ops.get(&op.seq) {
if existing.op_id == op.op_id {
outcome.deduped += 1;
continue;
}
return Err(RelayError::Fork {
device_id: device_id.to_string(),
seq: op.seq,
});
}
let expected = chain.next_seq();
if op.seq != expected {
return Err(RelayError::Gap {
device_id: device_id.to_string(),
expected,
found: op.seq,
});
}
match chain.head() {
Some((_, head_id, head_hlc)) => {
if op.prev.as_deref() != Some(head_id) {
return Err(RelayError::Chain {
device_id: device_id.to_string(),
detail: format!(
"op {}: prev does not link the relay-held head {head_id}",
op.op_id
),
});
}
if op.hlc <= *head_hlc {
return Err(RelayError::Chain {
device_id: device_id.to_string(),
detail: format!(
"op {}: hlc does not advance past the relay-held head",
op.op_id
),
});
}
}
None => {
if op.prev.is_some() {
return Err(RelayError::Chain {
device_id: device_id.to_string(),
detail: format!("op {}: seq 0 must have no prev", op.op_id),
});
}
}
}
chain.ops.insert(op.seq, op.clone());
outcome.accepted += 1;
}
Ok(outcome)
}
fn pull(&self, since: &Frontier) -> Result<PullResult, RelayError> {
let mut ops = Vec::new();
for (device_id, chain) in &self.chains {
let start = since.get(device_id).map(|held| held + 1).unwrap_or(0);
if start < chain.dropped_below {
return Err(RelayError::FrontierTruncated {
device_id: device_id.clone(),
dropped_below: chain.dropped_below,
});
}
ops.extend(chain.ops.range(start..).map(|(_, op)| op.clone()));
}
ops.sort_by(|a, b| (&a.hlc, &a.op_id).cmp(&(&b.hlc, &b.op_id)));
Ok(PullResult {
ops,
latest_checkpoint: self.latest_checkpoint.clone(),
})
}
fn ack(&mut self, device_id: &str, frontier: Hlc) -> AckOutcome {
let others_frontier = {
let others: Vec<&RosterEntry> = self
.roster
.values()
.filter(|e| e.status == DeviceStatus::Active && e.device_id != device_id)
.collect();
if others.is_empty() || others.iter().any(|e| e.acked.is_none()) {
None
} else {
others.iter().filter_map(|e| e.acked.clone()).min()
}
};
let entry = self.roster.get_mut(device_id).expect("touched before ack");
let advanced = match &entry.acked {
Some(current) if frontier <= *current => false,
_ => {
entry.acked = Some(frontier);
true
}
};
let mut reinstated = false;
if entry.status == DeviceStatus::Evicted {
let caught_up = match (&entry.acked, &others_frontier) {
(Some(acked), Some(frontier)) => acked >= frontier,
(Some(_), None) => true,
(None, _) => false,
};
if caught_up {
entry.status = DeviceStatus::Active;
reinstated = true;
}
}
AckOutcome {
advanced,
reinstated,
}
}
fn validate_frontier(&self, checkpoint: &Checkpoint) -> Result<(), RelayError> {
for (device_id, entry) in &checkpoint.frontier {
let unverified = |detail: &str| RelayError::CheckpointFrontierUnverified {
device_id: device_id.clone(),
seq: entry.seq,
detail: detail.to_string(),
};
let Some(chain) = self.chains.get(device_id) else {
return Err(unverified("relay holds no chain for this device"));
};
if entry.seq >= chain.dropped_below {
match chain.ops.get(&entry.seq) {
Some(op) if op.op_id == entry.head => {}
Some(_) => {
return Err(unverified(
"frontier head does not match the relay-held op at this seq",
))
}
None => {
return Err(unverified(
"relay holds no op at the claimed frontier seq (claims coverage \
beyond its chain head)",
))
}
}
} else {
match &chain.dropped_head {
Some((seq, id, _)) if *seq == entry.seq => {
if id != &entry.head {
return Err(unverified(
"frontier head does not match the relay's GC'd dropped head",
));
}
}
Some((seq, _, _)) if entry.seq < *seq => {}
_ => {
return Err(unverified(
"claimed frontier seq is below the relay's GC floor with no \
matching record",
))
}
}
}
}
Ok(())
}
fn checkpoint_put(&mut self, checkpoint: &Checkpoint) -> Result<bool, RelayError> {
checkpoint.verify().map_err(RelayError::Checkpoint)?;
self.validate_frontier(checkpoint)?;
let stored = if self.checkpoints.contains_key(&checkpoint.checkpoint_hash) {
false
} else {
self.checkpoints
.insert(checkpoint.checkpoint_hash.clone(), checkpoint.clone());
true
};
let advance = match self
.latest_checkpoint
.as_ref()
.and_then(|hash| self.checkpoints.get(hash))
{
Some(current) => {
checkpoint.checkpoint_hash != current.checkpoint_hash
&& frontier_dominates(checkpoint, current)
}
None => true,
};
if advance {
self.latest_checkpoint = Some(checkpoint.checkpoint_hash.clone());
}
Ok(stored)
}
fn checkpoint_get(&self) -> Option<Checkpoint> {
self.latest_checkpoint
.as_ref()
.and_then(|hash| self.checkpoints.get(hash))
.cloned()
}
fn gc(&mut self) -> GcReport {
let mut report = GcReport::default();
let Some(frontier) = self.stable_frontier() else {
return report;
};
let mut max_covered: BTreeMap<String, u64> = BTreeMap::new();
for ckpt in self.checkpoints.values() {
for (device, entry) in &ckpt.frontier {
let slot = max_covered.entry(device.clone()).or_insert(entry.seq);
if entry.seq > *slot {
*slot = entry.seq;
}
}
}
for (device_id, chain) in &mut self.chains {
let covered_through = max_covered.get(device_id).copied();
let mut droppable: Vec<u64> = Vec::new();
for (seq, op) in &chain.ops {
let below_frontier = op.hlc <= frontier;
let covered = covered_through.is_some_and(|through| through >= *seq);
if below_frontier && covered {
droppable.push(*seq);
} else {
break;
}
}
for seq in &droppable {
let op = chain.ops.remove(seq).expect("collected from the map");
chain.dropped_below = seq + 1;
chain.dropped_head = Some((*seq, op.op_id, op.hlc));
}
if !droppable.is_empty() {
report.dropped.insert(device_id.clone(), droppable.len());
}
}
report
}
}
pub struct InMemoryRelay {
state: RelayState,
config: RelayConfig,
wall: WallClock,
}
impl fmt::Debug for InMemoryRelay {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("InMemoryRelay")
.field("state", &self.state)
.field("config", &self.config)
.finish_non_exhaustive()
}
}
impl InMemoryRelay {
pub fn new(config: RelayConfig, wall: WallClock) -> Self {
Self {
state: RelayState::default(),
config,
wall,
}
}
}
impl Relay for InMemoryRelay {
fn register(&mut self, device_id: &str) -> Result<RosterEntry, RelayError> {
let now = (self.wall)();
self.state.touch_and_sweep(device_id, now, &self.config);
Ok(self.state.roster[device_id].clone())
}
fn push(&mut self, device_id: &str, ops: &[OpRecord]) -> Result<PushOutcome, RelayError> {
let now = (self.wall)();
self.state.touch_and_sweep(device_id, now, &self.config);
self.state.push(device_id, ops)
}
fn pull(&mut self, device_id: &str, since: &Frontier) -> Result<PullResult, RelayError> {
let now = (self.wall)();
self.state.touch_and_sweep(device_id, now, &self.config);
self.state.pull(since)
}
fn ack(&mut self, device_id: &str, frontier: Hlc) -> Result<AckOutcome, RelayError> {
let now = (self.wall)();
self.state.touch_and_sweep(device_id, now, &self.config);
Ok(self.state.ack(device_id, frontier))
}
fn checkpoint_put(
&mut self,
device_id: &str,
checkpoint: &Checkpoint,
) -> Result<bool, RelayError> {
let now = (self.wall)();
self.state.touch_and_sweep(device_id, now, &self.config);
self.state.checkpoint_put(checkpoint)
}
fn checkpoint_get(&mut self) -> Result<Option<Checkpoint>, RelayError> {
Ok(self.state.checkpoint_get())
}
fn roster(&mut self) -> Result<Vec<RosterEntry>, RelayError> {
Ok(self.state.roster.values().cloned().collect())
}
fn stable_frontier(&mut self) -> Result<Option<Hlc>, RelayError> {
Ok(self.state.stable_frontier())
}
fn gc(&mut self) -> Result<GcReport, RelayError> {
Ok(self.state.gc())
}
}
pub struct FsRelay {
dir: PathBuf,
config: RelayConfig,
wall: WallClock,
checkpoint_cache: BTreeMap<String, Checkpoint>,
}
impl fmt::Debug for FsRelay {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("FsRelay")
.field("dir", &self.dir)
.field("config", &self.config)
.finish_non_exhaustive()
}
}
impl FsRelay {
pub fn open(dir: &Path, config: RelayConfig, wall: WallClock) -> std::io::Result<Self> {
fs::create_dir_all(dir.join("checkpoints"))?;
Ok(Self {
dir: dir.to_path_buf(),
config,
wall,
checkpoint_cache: BTreeMap::new(),
})
}
fn state_path(&self) -> PathBuf {
self.dir.join("relay-state.json")
}
fn checkpoints_dir(&self) -> PathBuf {
self.dir.join("checkpoints")
}
fn with_state<T>(
&mut self,
f: impl FnOnce(&mut RelayState, u64, &RelayConfig) -> Result<T, RelayError>,
) -> Result<T, RelayError> {
let lock = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(self.dir.join("relay.lock"))
.map_err(RelayError::Io)?;
lock.lock().map_err(RelayError::Io)?;
let mut state: RelayState = match fs::read_to_string(self.state_path()) {
Ok(raw) => {
serde_json::from_str(&raw).map_err(|e| RelayError::Io(std::io::Error::other(e)))?
}
Err(e) if e.kind() == std::io::ErrorKind::NotFound => RelayState::default(),
Err(e) => return Err(RelayError::Io(e)),
};
for entry in fs::read_dir(self.checkpoints_dir()).map_err(RelayError::Io)? {
let path = entry.map_err(RelayError::Io)?.path();
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
let Some(hash) = name.strip_suffix(".checkpoint.json") else {
continue;
};
let ckpt = match self.checkpoint_cache.get(hash) {
Some(cached) => cached.clone(),
None => {
let ckpt = Checkpoint::load(&path).map_err(RelayError::Checkpoint)?;
self.checkpoint_cache
.insert(ckpt.checkpoint_hash.clone(), ckpt.clone());
ckpt
}
};
state.checkpoints.insert(ckpt.checkpoint_hash.clone(), ckpt);
}
let now = (self.wall)();
let result = f(&mut state, now, &self.config);
for ckpt in state.checkpoints.values() {
let path = self.checkpoints_dir().join(ckpt.file_name());
if !path.exists() {
ckpt.save(&self.checkpoints_dir()).map_err(RelayError::Io)?;
}
self.checkpoint_cache
.entry(ckpt.checkpoint_hash.clone())
.or_insert_with(|| ckpt.clone());
}
let tmp = self.dir.join("relay-state.json.tmp");
{
let mut file = File::create(&tmp).map_err(RelayError::Io)?;
file.write_all(
serde_json::to_string(&state)
.map_err(|e| RelayError::Io(std::io::Error::other(e)))?
.as_bytes(),
)
.map_err(RelayError::Io)?;
file.sync_all().map_err(RelayError::Io)?;
}
fs::rename(&tmp, self.state_path()).map_err(RelayError::Io)?;
result
}
}
impl Relay for FsRelay {
fn register(&mut self, device_id: &str) -> Result<RosterEntry, RelayError> {
self.with_state(|state, now, config| {
state.touch_and_sweep(device_id, now, config);
Ok(state.roster[device_id].clone())
})
}
fn push(&mut self, device_id: &str, ops: &[OpRecord]) -> Result<PushOutcome, RelayError> {
self.with_state(|state, now, config| {
state.touch_and_sweep(device_id, now, config);
state.push(device_id, ops)
})
}
fn pull(&mut self, device_id: &str, since: &Frontier) -> Result<PullResult, RelayError> {
self.with_state(|state, now, config| {
state.touch_and_sweep(device_id, now, config);
state.pull(since)
})
}
fn ack(&mut self, device_id: &str, frontier: Hlc) -> Result<AckOutcome, RelayError> {
self.with_state(|state, now, config| {
state.touch_and_sweep(device_id, now, config);
Ok(state.ack(device_id, frontier))
})
}
fn checkpoint_put(
&mut self,
device_id: &str,
checkpoint: &Checkpoint,
) -> Result<bool, RelayError> {
self.with_state(|state, now, config| {
state.touch_and_sweep(device_id, now, config);
state.checkpoint_put(checkpoint)
})
}
fn checkpoint_get(&mut self) -> Result<Option<Checkpoint>, RelayError> {
self.with_state(|state, _, _| Ok(state.checkpoint_get()))
}
fn roster(&mut self) -> Result<Vec<RosterEntry>, RelayError> {
self.with_state(|state, _, _| Ok(state.roster.values().cloned().collect()))
}
fn stable_frontier(&mut self) -> Result<Option<Hlc>, RelayError> {
self.with_state(|state, _, _| Ok(state.stable_frontier()))
}
fn gc(&mut self) -> Result<GcReport, RelayError> {
self.with_state(|state, _, _| Ok(state.gc()))
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::oplog::{DeviceLog, Scope, Surface};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
fn manual_clock() -> (Arc<AtomicU64>, WallClock) {
let t = Arc::new(AtomicU64::new(0));
let reader = t.clone();
(t, Arc::new(move || reader.load(Ordering::SeqCst)))
}
fn mem_relay() -> InMemoryRelay {
InMemoryRelay::new(RelayConfig::default(), Arc::new(|| 0))
}
fn ops_for(device: &str, n: usize) -> (DeviceLog, Vec<OpRecord>) {
let mut log = DeviceLog::new(device);
let ops = (0..n)
.map(|i| {
log.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": format!("{device}-f{i}")}),
)
})
.collect();
(log, ops)
}
#[test]
fn push_validates_the_chain_and_dedups_retransmission() {
let mut relay = mem_relay();
let (mut log, ops) = ops_for("a", 3);
let outcome = relay.push("a", &ops).unwrap();
assert_eq!(
outcome,
PushOutcome {
accepted: 3,
deduped: 0
}
);
let again = relay.push("a", &ops).unwrap();
assert_eq!(
again,
PushOutcome {
accepted: 0,
deduped: 3
}
);
let next = log.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": "x"}),
);
assert_eq!(relay.push("a", &[next]).unwrap().accepted, 1);
log.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": "skipped"}),
);
let ahead = log.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": "y"}),
);
assert!(matches!(
relay.push("a", &[ahead]),
Err(RelayError::Gap {
expected: 4,
found: 5,
..
})
));
let mut forked = DeviceLog::new("a");
let f0 = forked.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": "evil"}),
);
assert!(matches!(
relay.push("a", &[f0]),
Err(RelayError::Fork { seq: 0, .. })
));
let (_, b_ops) = ops_for("b", 1);
assert!(matches!(
relay.push("a", &b_ops),
Err(RelayError::ForeignOps { .. })
));
let mut tampered = ops[0].clone();
tampered.payload = serde_json::json!({"forged": true});
assert!(matches!(
relay.push("b", &[tampered]),
Err(RelayError::ForeignOps { .. })
));
let mut own_tampered = ops[0].clone();
own_tampered.payload = serde_json::json!({"id": "a-f0", "forged": true});
assert!(matches!(
relay.push("a", &[own_tampered]),
Err(RelayError::Chain { .. })
));
}
#[test]
fn pull_is_a_seq_cursor_and_serves_canonical_order() {
let mut relay = mem_relay();
let (_, a_ops) = ops_for("a", 3);
let (_, b_ops) = ops_for("b", 2);
relay.push("a", &a_ops).unwrap();
relay.push("b", &b_ops).unwrap();
let all = relay.pull("c", &Frontier::new()).unwrap();
assert_eq!(all.ops.len(), 5);
assert!(all.latest_checkpoint.is_none());
let mut since = Frontier::new();
since.insert("a".into(), 1);
since.insert("b".into(), 1);
let tail = relay.pull("c", &since).unwrap();
assert_eq!(tail.ops.len(), 1);
assert_eq!(tail.ops[0].seq, 2);
assert_eq!(tail.ops[0].device_id, "a");
}
#[test]
fn stable_frontier_requires_every_active_device_acked() {
let mut relay = mem_relay();
let (_, a_ops) = ops_for("a", 2);
relay.push("a", &a_ops).unwrap();
relay.register("b").unwrap();
relay.ack("a", a_ops[1].hlc.clone()).unwrap();
assert_eq!(relay.stable_frontier().unwrap(), None);
relay.ack("b", a_ops[0].hlc.clone()).unwrap();
assert_eq!(
relay.stable_frontier().unwrap(),
Some(a_ops[0].hlc.clone()),
"min(acked)"
);
let outcome = relay.ack("b", a_ops[0].hlc.clone()).unwrap();
assert!(!outcome.advanced);
let outcome = relay.ack("b", a_ops[1].hlc.clone()).unwrap();
assert!(outcome.advanced);
assert_eq!(relay.stable_frontier().unwrap(), Some(a_ops[1].hlc.clone()));
}
#[test]
fn gc_requires_both_frontier_and_covering_checkpoint() {
let mut relay = mem_relay();
let (_, ops) = ops_for("a", 4);
relay.push("a", &ops).unwrap();
relay.ack("a", ops[3].hlc.clone()).unwrap();
assert_eq!(relay.gc().unwrap().total(), 0);
let ckpt = Checkpoint::from_ops(&ops[..2]).unwrap();
assert!(relay.checkpoint_put("a", &ckpt).unwrap());
let report = relay.gc().unwrap();
assert_eq!(report.dropped["a"], 2);
let mut log = DeviceLog::resume("a", &ops).unwrap();
let newer = log.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": "n"}),
);
relay.push("a", std::slice::from_ref(&newer)).unwrap();
let full_ckpt = {
let mut all = ops.clone();
all.push(newer);
Checkpoint::from_ops(&all).unwrap()
};
relay.checkpoint_put("a", &full_ckpt).unwrap();
let report = relay.gc().unwrap();
assert_eq!(report.dropped["a"], 2);
let survivors = relay.pull("b", &Frontier::new());
assert!(matches!(
survivors,
Err(RelayError::FrontierTruncated {
dropped_below: 4,
..
})
));
let mut since = Frontier::new();
since.insert("a".into(), 3);
assert_eq!(relay.pull("b", &since).unwrap().ops.len(), 1);
}
#[test]
fn gc_preserves_chain_continuity_for_later_pushes() {
let mut relay = mem_relay();
let (mut log, ops) = ops_for("a", 3);
relay.push("a", &ops).unwrap();
relay.ack("a", ops[2].hlc.clone()).unwrap();
let ckpt = Checkpoint::from_ops(&ops).unwrap();
relay.checkpoint_put("a", &ckpt).unwrap();
assert_eq!(relay.gc().unwrap().dropped["a"], 3);
let next = log.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": "n"}),
);
assert_eq!(relay.push("a", &[next]).unwrap().accepted, 1);
assert_eq!(
relay.push("a", &ops).unwrap(),
PushOutcome {
accepted: 0,
deduped: 3
}
);
}
#[test]
fn forged_checkpoint_frontier_is_rejected_and_never_becomes_gc_coverage() {
let mut relay = mem_relay();
let mut a = DeviceLog::new("a");
let real: Vec<OpRecord> = (0..3)
.map(|i| {
a.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": format!("real-{i}")}),
)
})
.collect();
relay.push("a", &real).unwrap();
relay.ack("a", real[2].hlc.clone()).unwrap();
let mut fake = DeviceLog::new("a");
let other: Vec<OpRecord> = (0..3)
.map(|i| {
fake.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": format!("other-{i}")}),
)
})
.collect();
let forged = Checkpoint::from_ops(&other).unwrap();
assert!(matches!(
relay.checkpoint_put("a", &forged),
Err(RelayError::CheckpointFrontierUnverified { device_id, .. }) if device_id == "a"
));
assert_eq!(
relay.gc().unwrap().total(),
0,
"no forged coverage, no data loss"
);
assert!(relay.checkpoint_get().unwrap().is_none());
let served = relay.pull("a", &Frontier::new()).unwrap();
assert_eq!(served.ops.len(), 3);
assert!(served
.ops
.iter()
.all(|op| op.payload["id"].as_str().unwrap().starts_with("real-")));
let mut a2 = DeviceLog::resume("a", &real).unwrap();
let mut ahead = real.clone();
ahead.push(a2.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": "real-3"}),
));
let claims_beyond = Checkpoint::from_ops(&ahead).unwrap(); assert!(matches!(
relay.checkpoint_put("a", &claims_beyond),
Err(RelayError::CheckpointFrontierUnverified { seq: 3, .. })
));
let mut c = DeviceLog::new("c");
let c_ops: Vec<OpRecord> = (0..2)
.map(|i| {
c.append(
Scope::Personal,
Surface::Knowledge,
serde_json::json!({"id": format!("c{i}")}),
)
})
.collect();
let unknown_dev = Checkpoint::from_ops(&c_ops).unwrap();
assert!(matches!(
relay.checkpoint_put("a", &unknown_dev),
Err(RelayError::CheckpointFrontierUnverified { device_id, .. }) if device_id == "c"
));
let honest = Checkpoint::from_ops(&real[..2]).unwrap();
assert!(relay.checkpoint_put("a", &honest).unwrap());
assert_eq!(relay.gc().unwrap().dropped["a"], 2);
}
#[test]
fn checkpoint_dedup_keys_on_whole_record_address_not_state_hash() {
let mut a = DeviceLog::new("a");
let mut b = DeviceLog::new("b");
let fact = serde_json::json!({"id": "f1", "body": "hi"});
let oa = a.append(Scope::Personal, Surface::Knowledge, fact.clone());
let ob = b.append(Scope::Personal, Surface::Knowledge, fact);
let just_a = Checkpoint::from_ops(std::slice::from_ref(&oa)).unwrap();
let both = Checkpoint::from_ops(&[oa.clone(), ob.clone()]).unwrap();
assert_eq!(
just_a.state_hash, both.state_hash,
"the cross-device dedup collision"
);
let mut relay = mem_relay();
relay.push("a", std::slice::from_ref(&oa)).unwrap();
relay.push("b", std::slice::from_ref(&ob)).unwrap();
assert!(relay.checkpoint_put("a", &just_a).unwrap());
assert!(
relay.checkpoint_put("b", &both).unwrap(),
"same state_hash is NOT a dedup"
);
assert!(
!relay.checkpoint_put("a", &just_a).unwrap(),
"same checkpoint_hash IS"
);
assert_eq!(
relay.checkpoint_get().unwrap().unwrap().checkpoint_hash,
both.checkpoint_hash
);
relay.checkpoint_put("a", &just_a).unwrap();
assert_eq!(
relay.checkpoint_get().unwrap().unwrap().checkpoint_hash,
both.checkpoint_hash,
"dominance-monotone pointer"
);
let mut forged = both.clone();
forged.state.logs.clear();
assert!(matches!(
relay.checkpoint_put("b", &forged),
Err(RelayError::Checkpoint(CheckpointError::HashMismatch { .. }))
));
}
#[test]
fn eviction_unpins_the_frontier_and_ack_reinstates() {
let (t, wall) = manual_clock();
let mut relay = InMemoryRelay::new(
RelayConfig {
eviction_horizon_ms: Some(1_000),
},
wall,
);
let (_, a_ops) = ops_for("a", 2);
relay.push("a", &a_ops).unwrap();
t.store(100, Ordering::SeqCst);
relay.ack("a", a_ops[1].hlc.clone()).unwrap();
relay.ack("c", a_ops[0].hlc.clone()).unwrap(); assert_eq!(relay.stable_frontier().unwrap(), Some(a_ops[0].hlc.clone()));
t.store(2_000, Ordering::SeqCst);
relay.register("a").unwrap();
let roster: BTreeMap<String, RosterEntry> = relay
.roster()
.unwrap()
.into_iter()
.map(|e| (e.device_id.clone(), e))
.collect();
assert_eq!(roster["c"].status, DeviceStatus::Evicted);
assert_eq!(roster["a"].status, DeviceStatus::Active);
assert_eq!(
relay.stable_frontier().unwrap(),
Some(a_ops[1].hlc.clone()),
"the evicted device's ack no longer holds the frontier"
);
let outcome = relay.ack("c", a_ops[0].hlc.clone()).unwrap();
assert!(!outcome.advanced && !outcome.reinstated);
assert_eq!(relay.stable_frontier().unwrap(), Some(a_ops[1].hlc.clone()));
let outcome = relay.ack("c", a_ops[1].hlc.clone()).unwrap();
assert!(outcome.advanced && outcome.reinstated);
let roster: BTreeMap<String, RosterEntry> = relay
.roster()
.unwrap()
.into_iter()
.map(|e| (e.device_id.clone(), e))
.collect();
assert_eq!(roster["c"].status, DeviceStatus::Active);
}
#[test]
fn fs_relay_matches_in_memory_semantics_and_persists() {
let dir = tempfile::tempdir().unwrap();
let (_, ops) = ops_for("a", 3);
{
let mut relay =
FsRelay::open(dir.path(), RelayConfig::default(), Arc::new(|| 7)).unwrap();
assert_eq!(relay.push("a", &ops).unwrap().accepted, 3);
relay.ack("a", ops[2].hlc.clone()).unwrap();
let ckpt = Checkpoint::from_ops(&ops[..2]).unwrap();
assert!(relay.checkpoint_put("a", &ckpt).unwrap());
}
let mut relay = FsRelay::open(dir.path(), RelayConfig::default(), Arc::new(|| 8)).unwrap();
assert_eq!(relay.stable_frontier().unwrap(), Some(ops[2].hlc.clone()));
assert_eq!(relay.pull("b", &Frontier::new()).unwrap().ops, ops);
assert_eq!(relay.stable_frontier().unwrap(), None);
relay.ack("b", ops[2].hlc.clone()).unwrap();
assert_eq!(
relay.push("a", &ops).unwrap(),
PushOutcome {
accepted: 0,
deduped: 3
}
);
let report = relay.gc().unwrap();
assert_eq!(report.dropped["a"], 2);
let roster = relay.roster().unwrap();
assert_eq!(roster.len(), 2);
assert_eq!(
roster[0].added_at.wall_ms, 7,
"roster added_at survives restart"
);
let ckpt = relay.checkpoint_get().unwrap().unwrap();
assert_eq!(ckpt.frontier["a"].seq, 1);
let path = dir
.path()
.join("checkpoints")
.join(format!("{}.checkpoint.json", ckpt.checkpoint_hash));
let raw = fs::read_to_string(&path).unwrap();
fs::write(&path, raw.replace("\"a-f0\"", "\"a-f0-forged\"")).unwrap();
let mut fresh = FsRelay::open(dir.path(), RelayConfig::default(), Arc::new(|| 9)).unwrap();
assert!(matches!(
fresh.checkpoint_get(),
Err(RelayError::Checkpoint(_))
));
}
}