use alloc::sync::Arc;
use std::collections::{HashMap, HashSet};
use std::sync::{Mutex, MutexGuard};
use serde::{Deserialize, Serialize};
use thiserror::Error;
use routers_network::Entry;
use routers_transition::matcher::Trip;
use crate::event::VehicleId;
use crate::protocol::ids::{
GraphVersion, ObservationId, OutputId, RegionId, Revision, SchemaVersion, SegmentId,
};
#[derive(Clone, Debug, Serialize, Deserialize)]
#[serde(bound(serialize = "E: Serialize", deserialize = "E: Deserialize<'de>"))]
pub struct VehicleCheckpoint<E: Entry> {
pub trip: Trip<E>,
pub last_input: ObservationId,
pub revision: Revision,
pub segment: SegmentId,
pub finalized_through: Option<i64>,
pub graph: GraphVersion,
pub schema: SchemaVersion,
pub region: RegionId,
pub routing_version: u64,
}
impl<E: Entry> VehicleCheckpoint<E> {
pub fn encode(&self) -> Result<Vec<u8>, postcard::Error> {
postcard::to_allocvec(self)
}
pub fn decode(bytes: &[u8]) -> Result<Self, postcard::Error>
where
E: serde::de::DeserializeOwned,
{
postcard::from_bytes(bytes)
}
pub fn to_stored(&self) -> Result<StoredCheckpoint, postcard::Error> {
Ok(StoredCheckpoint {
revision: self.revision,
segment: self.segment,
bytes: self.encode()?,
})
}
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct StoredCheckpoint {
pub revision: Revision,
pub segment: SegmentId,
pub bytes: Vec<u8>,
}
#[derive(Clone, Debug, PartialEq, Eq)]
pub enum StoredCheckpointState {
NeverSeen,
Present(StoredCheckpoint),
Expired {
revision: Revision,
},
}
impl StoredCheckpointState {
#[must_use]
pub fn revision(&self) -> Option<Revision> {
match self {
Self::NeverSeen => None,
Self::Present(checkpoint) => Some(checkpoint.revision),
Self::Expired { revision } => Some(*revision),
}
}
#[must_use]
pub fn present(&self) -> Option<&StoredCheckpoint> {
match self {
Self::Present(checkpoint) => Some(checkpoint),
Self::NeverSeen | Self::Expired { .. } => None,
}
}
#[must_use]
pub fn is_some(&self) -> bool {
matches!(self, Self::Present(_))
}
#[must_use]
pub fn is_none(&self) -> bool {
!self.is_some()
}
pub fn expect(self, message: &str) -> StoredCheckpoint {
match self {
Self::Present(checkpoint) => checkpoint,
Self::NeverSeen | Self::Expired { .. } => panic!("{message}"),
}
}
pub fn unwrap(self) -> StoredCheckpoint {
self.expect("called StoredCheckpointState::unwrap without a present checkpoint")
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub enum CommitPhase {
Prepared,
Published,
}
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct PreparedCommit {
pub output: OutputId,
pub output_subject: String,
pub output_bytes: Vec<u8>,
pub next_checkpoint: Vec<u8>,
pub next_revision: Revision,
pub next_segment: SegmentId,
pub expected_base: Option<Revision>,
pub phase: CommitPhase,
pub raw: ObservationId,
}
impl PreparedCommit {
pub fn is_published(&self) -> bool {
matches!(self.phase, CommitPhase::Published)
}
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash, Serialize, Deserialize)]
pub struct PartitionFrontier {
pub partition: u16,
pub sequence: u64,
}
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub enum PrepareOutcome {
Prepared,
AlreadyPrepared,
Conflict { actual: Option<Revision> },
Busy { pending: OutputId },
}
#[allow(async_fn_in_trait)]
pub trait CheckpointStore: Clone + Send + Sync + 'static {
type Error: core::error::Error + Send + Sync + 'static;
async fn load(
&self,
vehicle: VehicleId,
) -> Result<(StoredCheckpointState, Option<PreparedCommit>), Self::Error>;
async fn prepare(
&self,
vehicle: VehicleId,
partition: u16,
prepared: PreparedCommit,
) -> Result<PrepareOutcome, Self::Error>;
async fn mark_published(&self, vehicle: VehicleId, output: OutputId)
-> Result<(), Self::Error>;
async fn promote(
&self,
vehicle: VehicleId,
partition: u16,
output: OutputId,
) -> Result<(), Self::Error>;
async fn list_prepared(
&self,
partition: u16,
) -> Result<Vec<(VehicleId, PreparedCommit)>, Self::Error>;
async fn frontier(&self, partition: u16) -> Result<Option<PartitionFrontier>, Self::Error>;
async fn set_frontier(&self, frontier: PartitionFrontier) -> Result<(), Self::Error>;
async fn expire_idle(&self, vehicle: VehicleId) -> Result<(), Self::Error>;
}
#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
pub enum Op {
Load,
Prepare,
MarkPublished,
Promote,
ListPrepared,
Frontier,
SetFrontier,
Expire,
}
#[derive(Debug, Error)]
pub enum MemoryError {
#[error("injected memory-store fault")]
Injected,
#[error("prepared output mismatch: staged {staged}, got {got}")]
OutputMismatch { staged: OutputId, got: OutputId },
#[error("prepared output {output} has not been published")]
NotPublished { output: OutputId },
}
#[derive(Default)]
struct Inner {
checkpoints: HashMap<VehicleId, StoredCheckpoint>,
committed_revisions: HashMap<VehicleId, Revision>,
prepared: HashMap<VehicleId, PreparedCommit>,
index: HashMap<u16, HashSet<VehicleId>>,
frontiers: HashMap<u16, PartitionFrontier>,
faults: HashSet<Op>,
}
#[derive(Clone, Default)]
pub struct MemoryCheckpointStore {
inner: Arc<Mutex<Inner>>,
}
#[derive(Clone, Debug, Default, PartialEq, Eq)]
pub struct MemorySnapshot {
pub checkpoints: HashMap<VehicleId, StoredCheckpoint>,
pub committed_revisions: HashMap<VehicleId, Revision>,
pub prepared: HashMap<VehicleId, PreparedCommit>,
pub index: HashMap<u16, HashSet<VehicleId>>,
pub frontiers: HashMap<u16, PartitionFrontier>,
}
impl MemoryCheckpointStore {
pub fn new() -> Self {
Self::default()
}
pub fn fail_next(&self, op: Op) {
self.lock().faults.insert(op);
}
pub fn snapshot(&self) -> MemorySnapshot {
let inner = self.lock();
MemorySnapshot {
checkpoints: inner.checkpoints.clone(),
committed_revisions: inner.committed_revisions.clone(),
prepared: inner.prepared.clone(),
index: inner.index.clone(),
frontiers: inner.frontiers.clone(),
}
}
pub fn drop_all_checkpoints(&self) {
self.lock().checkpoints.clear();
}
pub fn drop_prepared(&self, vehicle: VehicleId) {
let mut inner = self.lock();
inner.prepared.remove(&vehicle);
for set in inner.index.values_mut() {
set.remove(&vehicle);
}
}
fn lock(&self) -> MutexGuard<'_, Inner> {
self.inner.lock().expect("checkpoint store mutex poisoned")
}
}
impl CheckpointStore for MemoryCheckpointStore {
type Error = MemoryError;
async fn load(
&self,
vehicle: VehicleId,
) -> Result<(StoredCheckpointState, Option<PreparedCommit>), Self::Error> {
let mut inner = self.lock();
if inner.faults.remove(&Op::Load) {
return Err(MemoryError::Injected);
}
let checkpoint = match inner.checkpoints.get(&vehicle).cloned() {
Some(checkpoint) => StoredCheckpointState::Present(checkpoint),
None => inner
.committed_revisions
.get(&vehicle)
.copied()
.map_or(StoredCheckpointState::NeverSeen, |revision| {
StoredCheckpointState::Expired { revision }
}),
};
let prepared = inner.prepared.get(&vehicle).cloned();
Ok((checkpoint, prepared))
}
async fn prepare(
&self,
vehicle: VehicleId,
partition: u16,
prepared: PreparedCommit,
) -> Result<PrepareOutcome, Self::Error> {
let mut inner = self.lock();
if inner.faults.remove(&Op::Prepare) {
return Err(MemoryError::Injected);
}
if let Some(existing) = inner.prepared.get(&vehicle) {
return Ok(if existing.output == prepared.output {
PrepareOutcome::AlreadyPrepared
} else {
PrepareOutcome::Busy {
pending: existing.output,
}
});
}
let actual = inner.committed_revisions.get(&vehicle).copied();
let matches_base = match prepared.expected_base {
None => actual.is_none(),
Some(base) => actual == Some(base),
};
if !matches_base {
return Ok(PrepareOutcome::Conflict { actual });
}
inner.index.entry(partition).or_default().insert(vehicle);
inner.prepared.insert(vehicle, prepared);
Ok(PrepareOutcome::Prepared)
}
async fn mark_published(
&self,
vehicle: VehicleId,
output: OutputId,
) -> Result<(), Self::Error> {
let mut inner = self.lock();
if inner.faults.remove(&Op::MarkPublished) {
return Err(MemoryError::Injected);
}
match inner.prepared.get_mut(&vehicle) {
Some(prepared) if prepared.output == output => {
prepared.phase = CommitPhase::Published;
Ok(())
}
Some(prepared) => Err(MemoryError::OutputMismatch {
staged: prepared.output,
got: output,
}),
None => Ok(()),
}
}
async fn promote(
&self,
vehicle: VehicleId,
partition: u16,
output: OutputId,
) -> Result<(), Self::Error> {
let mut inner = self.lock();
if inner.faults.remove(&Op::Promote) {
return Err(MemoryError::Injected);
}
let stored = match inner.prepared.get(&vehicle) {
Some(prepared) if prepared.output == output && prepared.is_published() => {
StoredCheckpoint {
revision: prepared.next_revision,
segment: prepared.next_segment,
bytes: prepared.next_checkpoint.clone(),
}
}
Some(prepared) if prepared.output == output => {
return Err(MemoryError::NotPublished { output });
}
Some(prepared) => {
return Err(MemoryError::OutputMismatch {
staged: prepared.output,
got: output,
});
}
None => return Ok(()),
};
inner.prepared.remove(&vehicle);
if let Some(set) = inner.index.get_mut(&partition) {
set.remove(&vehicle);
}
inner.committed_revisions.insert(vehicle, stored.revision);
inner.checkpoints.insert(vehicle, stored);
Ok(())
}
async fn list_prepared(
&self,
partition: u16,
) -> Result<Vec<(VehicleId, PreparedCommit)>, Self::Error> {
let mut inner = self.lock();
if inner.faults.remove(&Op::ListPrepared) {
return Err(MemoryError::Injected);
}
let mut out: Vec<(VehicleId, PreparedCommit)> = inner
.index
.get(&partition)
.into_iter()
.flatten()
.filter_map(|vehicle| inner.prepared.get(vehicle).map(|p| (*vehicle, p.clone())))
.collect();
out.sort_by_key(|(vehicle, _)| vehicle.0);
Ok(out)
}
async fn frontier(&self, partition: u16) -> Result<Option<PartitionFrontier>, Self::Error> {
let mut inner = self.lock();
if inner.faults.remove(&Op::Frontier) {
return Err(MemoryError::Injected);
}
Ok(inner.frontiers.get(&partition).copied())
}
async fn set_frontier(&self, frontier: PartitionFrontier) -> Result<(), Self::Error> {
let mut inner = self.lock();
if inner.faults.remove(&Op::SetFrontier) {
return Err(MemoryError::Injected);
}
inner.frontiers.insert(frontier.partition, frontier);
Ok(())
}
async fn expire_idle(&self, vehicle: VehicleId) -> Result<(), Self::Error> {
let mut inner = self.lock();
if inner.faults.remove(&Op::Expire) {
return Err(MemoryError::Injected);
}
inner.checkpoints.remove(&vehicle);
Ok(())
}
}
#[cfg(test)]
mod tests {
use routers_network::mock::MockEntryId;
use crate::protocol::ids::SCHEMA_VERSION;
use super::*;
fn store() -> MemoryCheckpointStore {
MemoryCheckpointStore::new()
}
fn obs(partition: u16, sequence: u64) -> ObservationId {
ObservationId {
partition,
sequence,
}
}
fn prepared(output: u128, expected_base: Option<Revision>) -> PreparedCommit {
PreparedCommit {
output: OutputId(output),
output_subject: "events.matched.v1.p.0".to_owned(),
output_bytes: vec![0xaa, 0xbb],
next_checkpoint: vec![1, 2, 3, 4],
next_revision: Revision(expected_base.map_or(1, |r| r.0 + 1)),
next_segment: SegmentId(7),
expected_base,
phase: CommitPhase::Prepared,
raw: obs(0, 100),
}
}
async fn seed_checkpoint(
s: &MemoryCheckpointStore,
vehicle: VehicleId,
revision: Revision,
segment: SegmentId,
) {
let mut p = prepared(0x1000 + u128::from(revision.0), None);
p.next_revision = revision;
p.next_segment = segment;
assert_eq!(
s.prepare(vehicle, 0, p.clone()).await.unwrap(),
PrepareOutcome::Prepared
);
s.mark_published(vehicle, p.output).await.unwrap();
s.promote(vehicle, 0, p.output).await.unwrap();
}
#[test]
fn is_published_reflects_phase() {
let mut p = prepared(0x1, None);
assert!(!p.is_published());
p.phase = CommitPhase::Published;
assert!(p.is_published());
}
#[tokio::test]
async fn prepare_stages_a_fresh_vehicle() {
let s = store();
let v = VehicleId(1);
let p = prepared(0xa1, None);
assert_eq!(
s.prepare(v, 3, p.clone()).await.unwrap(),
PrepareOutcome::Prepared
);
let (checkpoint, pending) = s.load(v).await.unwrap();
assert!(checkpoint.is_none());
assert_eq!(pending, Some(p.clone()));
assert_eq!(s.list_prepared(3).await.unwrap(), vec![(v, p)]);
}
#[tokio::test]
async fn prepare_same_output_is_idempotent_and_different_is_busy() {
let s = store();
let v = VehicleId(1);
let first = prepared(0xa1, None);
assert_eq!(
s.prepare(v, 0, first.clone()).await.unwrap(),
PrepareOutcome::Prepared
);
assert_eq!(
s.prepare(v, 0, first).await.unwrap(),
PrepareOutcome::AlreadyPrepared
);
assert_eq!(
s.prepare(v, 0, prepared(0xb2, None)).await.unwrap(),
PrepareOutcome::Busy {
pending: OutputId(0xa1),
}
);
}
#[tokio::test]
async fn prepare_conflicts_when_base_disagrees() {
let s = store();
let v = VehicleId(1);
seed_checkpoint(&s, v, Revision(5), SegmentId(5)).await;
assert_eq!(
s.prepare(v, 0, prepared(0xc3, None)).await.unwrap(),
PrepareOutcome::Conflict {
actual: Some(Revision(5)),
}
);
assert_eq!(
s.prepare(v, 0, prepared(0xc3, Some(Revision(4))))
.await
.unwrap(),
PrepareOutcome::Conflict {
actual: Some(Revision(5)),
}
);
assert_eq!(
s.prepare(v, 0, prepared(0xc3, Some(Revision(5))))
.await
.unwrap(),
PrepareOutcome::Prepared
);
}
#[tokio::test]
async fn prepare_conflicts_when_expecting_a_missing_checkpoint() {
let s = store();
let v = VehicleId(2);
assert_eq!(
s.prepare(v, 0, prepared(0xd4, Some(Revision(9))))
.await
.unwrap(),
PrepareOutcome::Conflict { actual: None }
);
}
#[tokio::test]
async fn promote_installs_checkpoint_and_clears_prepared() {
let s = store();
let v = VehicleId(1);
let mut p = prepared(0xa1, None);
p.next_revision = Revision(42);
p.next_segment = SegmentId(7);
p.next_checkpoint = vec![9, 8, 7];
assert_eq!(
s.prepare(v, 4, p.clone()).await.unwrap(),
PrepareOutcome::Prepared
);
s.mark_published(v, p.output).await.unwrap();
s.promote(v, 4, p.output).await.unwrap();
let (checkpoint, pending) = s.load(v).await.unwrap();
assert_eq!(pending, None);
assert_eq!(
checkpoint,
StoredCheckpointState::Present(StoredCheckpoint {
revision: Revision(42),
segment: SegmentId(7),
bytes: vec![9, 8, 7],
})
);
assert_eq!(s.list_prepared(4).await.unwrap(), vec![]);
s.promote(v, 4, p.output).await.unwrap();
assert_eq!(s.load(v).await.unwrap().0.unwrap().revision, Revision(42));
}
#[tokio::test]
async fn promote_rejects_a_different_output() {
let s = store();
let v = VehicleId(1);
s.prepare(v, 0, prepared(0xa1, None)).await.unwrap();
let err = s.promote(v, 0, OutputId(0xffff)).await.unwrap_err();
assert!(matches!(err, MemoryError::OutputMismatch { .. }));
assert!(s.load(v).await.unwrap().1.is_some());
}
#[tokio::test]
async fn promote_rejects_an_unpublished_record() {
let s = store();
let v = VehicleId(1);
let p = prepared(0xa1, None);
s.prepare(v, 0, p.clone()).await.unwrap();
let err = s.promote(v, 0, p.output).await.unwrap_err();
assert!(matches!(err, MemoryError::NotPublished { output } if output == p.output));
assert!(!s.load(v).await.unwrap().1.unwrap().is_published());
}
#[tokio::test]
async fn mark_published_sets_phase_and_rejects_other_output() {
let s = store();
let v = VehicleId(1);
let p = prepared(0xa1, None);
s.prepare(v, 0, p.clone()).await.unwrap();
let err = s.mark_published(v, OutputId(0xffff)).await.unwrap_err();
assert!(matches!(err, MemoryError::OutputMismatch { .. }));
assert!(!s.load(v).await.unwrap().1.unwrap().is_published());
s.mark_published(v, p.output).await.unwrap();
assert!(s.load(v).await.unwrap().1.unwrap().is_published());
s.mark_published(v, p.output).await.unwrap();
assert!(s.load(v).await.unwrap().1.unwrap().is_published());
}
#[tokio::test]
async fn list_prepared_is_scoped_to_a_partition() {
let s = store();
let (a, b, c) = (VehicleId(1), VehicleId(2), VehicleId(3));
s.prepare(a, 0, prepared(0xa, None)).await.unwrap();
s.prepare(b, 0, prepared(0xb, None)).await.unwrap();
s.prepare(c, 1, prepared(0xc, None)).await.unwrap();
let p0: Vec<VehicleId> = s
.list_prepared(0)
.await
.unwrap()
.into_iter()
.map(|(v, _)| v)
.collect();
assert_eq!(p0, vec![a, b]);
let p1: Vec<VehicleId> = s
.list_prepared(1)
.await
.unwrap()
.into_iter()
.map(|(v, _)| v)
.collect();
assert_eq!(p1, vec![c]);
assert_eq!(s.list_prepared(99).await.unwrap(), vec![]);
}
#[tokio::test]
async fn frontier_round_trips_per_partition() {
let s = store();
assert_eq!(s.frontier(7).await.unwrap(), None);
s.set_frontier(PartitionFrontier {
partition: 7,
sequence: 1234,
})
.await
.unwrap();
assert_eq!(
s.frontier(7).await.unwrap(),
Some(PartitionFrontier {
partition: 7,
sequence: 1234,
})
);
s.set_frontier(PartitionFrontier {
partition: 7,
sequence: 2000,
})
.await
.unwrap();
assert_eq!(s.frontier(7).await.unwrap().unwrap().sequence, 2000);
assert_eq!(s.frontier(8).await.unwrap(), None);
}
#[tokio::test]
async fn expire_idle_drops_checkpoint_but_keeps_prepared() {
let s = store();
let v = VehicleId(1);
seed_checkpoint(&s, v, Revision(5), SegmentId(5)).await;
let p = prepared(0xe5, Some(Revision(5)));
s.prepare(v, 0, p.clone()).await.unwrap();
s.expire_idle(v).await.unwrap();
let (checkpoint, pending) = s.load(v).await.unwrap();
assert_eq!(
checkpoint,
StoredCheckpointState::Expired {
revision: Revision(5)
},
"the payload is dropped but continuity is retained"
);
assert_eq!(pending, Some(p), "the prepared record is untouched");
}
async fn invoke(s: &MemoryCheckpointStore, op: Op) -> Result<(), MemoryError> {
let v = VehicleId(7);
match op {
Op::Load => s.load(v).await.map(|_| ()),
Op::Prepare => s.prepare(v, 0, prepared(0x1, None)).await.map(|_| ()),
Op::MarkPublished => s.mark_published(v, OutputId(0x1)).await,
Op::Promote => s.promote(v, 0, OutputId(0x1)).await,
Op::ListPrepared => s.list_prepared(0).await.map(|_| ()),
Op::Frontier => s.frontier(0).await.map(|_| ()),
Op::SetFrontier => {
s.set_frontier(PartitionFrontier {
partition: 0,
sequence: 0,
})
.await
}
Op::Expire => s.expire_idle(v).await,
}
}
#[tokio::test]
async fn every_op_honours_its_injected_fault_once() {
let ops = [
Op::Load,
Op::Prepare,
Op::MarkPublished,
Op::Promote,
Op::ListPrepared,
Op::Frontier,
Op::SetFrontier,
Op::Expire,
];
for op in ops {
let s = store();
s.fail_next(op);
assert!(
matches!(invoke(&s, op).await, Err(MemoryError::Injected)),
"{op:?} should fail once"
);
assert!(
invoke(&s, op).await.is_ok(),
"{op:?} should succeed after the fault clears"
);
}
}
#[tokio::test]
async fn snapshot_is_a_deep_copy_and_state_loss_helpers_work() {
let s = store();
let v = VehicleId(1);
seed_checkpoint(&s, v, Revision(5), SegmentId(5)).await;
s.prepare(v, 2, prepared(0xf6, Some(Revision(5))))
.await
.unwrap();
s.set_frontier(PartitionFrontier {
partition: 2,
sequence: 10,
})
.await
.unwrap();
let snap = s.snapshot();
assert_eq!(snap.checkpoints.len(), 1);
assert_eq!(snap.prepared.len(), 1);
assert_eq!(snap.frontiers.get(&2).unwrap().sequence, 10);
assert!(snap.index.get(&2).unwrap().contains(&v));
s.drop_all_checkpoints();
let (checkpoint, pending) = s.load(v).await.unwrap();
assert!(checkpoint.is_none());
assert!(pending.is_some());
s.drop_prepared(v);
assert!(s.load(v).await.unwrap().1.is_none());
assert!(s.list_prepared(2).await.unwrap().is_empty());
assert_eq!(snap.prepared.len(), 1);
assert_eq!(snap.checkpoints.len(), 1);
}
#[test]
fn vehicle_checkpoint_round_trips() {
let checkpoint = VehicleCheckpoint::<MockEntryId> {
trip: Trip::new(),
last_input: obs(485, 42),
revision: Revision(42),
segment: SegmentId(7),
finalized_through: Some(1_700_000_000_000_000),
graph: GraphVersion::new("g1").unwrap(),
schema: SCHEMA_VERSION,
region: RegionId::new("r1").unwrap(),
routing_version: 1,
};
let bytes = checkpoint.encode().expect("encode");
let decoded = VehicleCheckpoint::<MockEntryId>::decode(&bytes).expect("decode");
assert_eq!(decoded.encode().expect("re-encode"), bytes);
assert_eq!(decoded.last_input, checkpoint.last_input);
assert_eq!(decoded.revision, Revision(42));
assert_eq!(decoded.segment, SegmentId(7));
assert_eq!(decoded.finalized_through, Some(1_700_000_000_000_000));
assert_eq!(decoded.graph, checkpoint.graph);
assert_eq!(decoded.schema, SCHEMA_VERSION);
assert_eq!(decoded.region, checkpoint.region);
}
#[test]
fn to_stored_lifts_revision_and_segment_out_of_the_bytes() {
let checkpoint = VehicleCheckpoint::<MockEntryId> {
trip: Trip::new(),
last_input: obs(0, 1),
revision: Revision(9),
segment: SegmentId(3),
finalized_through: None,
graph: GraphVersion::new("g").unwrap(),
schema: SCHEMA_VERSION,
region: RegionId::new("r").unwrap(),
routing_version: 1,
};
let stored = checkpoint.to_stored().expect("to_stored");
assert_eq!(stored.revision, Revision(9));
assert_eq!(stored.segment, SegmentId(3));
assert_eq!(stored.bytes, checkpoint.encode().unwrap());
}
}