#![cfg(feature = "ipc")]
use lazily::{
CapabilityHandshake, CrdtOp, CrdtSync, Delta, DeltaApplyStatus, DeltaOp, DeltaSinceRequest,
EdgeSnapshot, IpcMessage, KeyIndex, NODE_KEY_MAX_SEGMENTS, NodeId, NodeKey, NodeKeyError,
NodeSnapshot, NodeState, OpKind, PeerId, PeerPermissions, RemoteOp, SHM_BLOB_HEADER_LEN,
ShmBlobArena, ShmBlobArenaError, Snapshot, WireStamp,
};
const PEER_A: PeerId = PeerId(1);
const PEER_B: PeerId = PeerId(2);
#[test]
fn snapshot_round_trips_through_serde() {
let snapshot = Snapshot::new(
7,
vec![
NodeSnapshot::payload(NodeId(1), "i32", vec![1, 2, 3]),
NodeSnapshot::opaque(NodeId(2), "opaque-type"),
],
vec![EdgeSnapshot::new(NodeId(2), NodeId(1))],
vec![NodeId(1), NodeId(2)],
);
let json = serde_json::to_string(&IpcMessage::Snapshot(snapshot.clone())).unwrap();
let back: IpcMessage = serde_json::from_str(&json).unwrap();
assert_eq!(back, IpcMessage::Snapshot(snapshot));
}
#[test]
fn delta_round_trips_through_serde() {
let delta = Delta::next(
41,
vec![
DeltaOp::cell_set(NodeId(1), vec![10]),
DeltaOp::slot_value(NodeId(2), vec![20]),
DeltaOp::invalidate(NodeId(3)),
DeltaOp::NodeAdd {
node: NodeId(4),
type_tag: "u64".into(),
state: NodeState::Payload(vec![64]),
key: None,
},
DeltaOp::NodeRemove { node: NodeId(5) },
DeltaOp::EdgeAdd {
dependent: NodeId(2),
dependency: NodeId(1),
},
DeltaOp::EdgeRemove {
dependent: NodeId(3),
dependency: NodeId(1),
},
],
);
let json = serde_json::to_string(&IpcMessage::Delta(delta.clone())).unwrap();
let back: IpcMessage = serde_json::from_str(&json).unwrap();
assert_eq!(back, IpcMessage::Delta(delta));
}
#[test]
fn delta_status_accepts_only_sequential_epochs() {
let next = Delta::next(10, vec![]);
assert_eq!(next.apply_status(10), DeltaApplyStatus::Apply);
assert!(next.is_next_after(10));
let gap = Delta::new(12, 13, vec![]);
assert_eq!(
gap.apply_status(10),
DeltaApplyStatus::ResyncRequired {
last_epoch: 10,
base_epoch: 12,
epoch: 13,
}
);
let non_sequential = Delta::new(10, 12, vec![]);
assert_eq!(
non_sequential.apply_status(10),
DeltaApplyStatus::ResyncRequired {
last_epoch: 10,
base_epoch: 10,
epoch: 12,
}
);
}
#[test]
fn snapshot_filter_omits_non_readable_nodes_edges_and_roots() {
let snapshot = Snapshot::new(
5,
vec![
NodeSnapshot::payload(NodeId(1), "i32", vec![1]),
NodeSnapshot::payload(NodeId(2), "i32", vec![2]),
NodeSnapshot::payload(NodeId(3), "i32", vec![3]),
],
vec![
EdgeSnapshot::new(NodeId(2), NodeId(1)),
EdgeSnapshot::new(NodeId(3), NodeId(1)),
],
vec![NodeId(1), NodeId(2), NodeId(3)],
);
let mut permissions = PeerPermissions::new();
permissions.allow_many(PEER_A, OpKind::Read, [NodeId(1), NodeId(2)]);
permissions.allow(PEER_A, RemoteOp::write(NodeId(3)));
let filtered = snapshot.filter_readable(&permissions, PEER_A);
assert_eq!(
filtered.nodes,
vec![
NodeSnapshot::payload(NodeId(1), "i32", vec![1]),
NodeSnapshot::payload(NodeId(2), "i32", vec![2]),
]
);
assert_eq!(
filtered.edges,
vec![EdgeSnapshot::new(NodeId(2), NodeId(1))]
);
assert_eq!(filtered.roots, vec![NodeId(1), NodeId(2)]);
let empty = snapshot.filter_readable(&permissions, PEER_B);
assert!(empty.nodes.is_empty());
assert!(empty.edges.is_empty());
assert!(empty.roots.is_empty());
}
#[test]
fn delta_filter_omits_non_readable_ops_without_redaction() {
let delta = Delta::next(
8,
vec![
DeltaOp::cell_set(NodeId(1), vec![1]),
DeltaOp::slot_value(NodeId(2), vec![2]),
DeltaOp::invalidate(NodeId(3)),
DeltaOp::NodeAdd {
node: NodeId(4),
type_tag: "u8".into(),
state: NodeState::Payload(vec![4]),
key: None,
},
DeltaOp::NodeRemove { node: NodeId(5) },
DeltaOp::EdgeAdd {
dependent: NodeId(2),
dependency: NodeId(1),
},
DeltaOp::EdgeRemove {
dependent: NodeId(3),
dependency: NodeId(1),
},
],
);
let mut permissions = PeerPermissions::new();
permissions.allow_many(PEER_A, OpKind::Read, [NodeId(1), NodeId(2), NodeId(5)]);
let filtered = delta.filter_readable(&permissions, PEER_A);
assert_eq!(
filtered.ops,
vec![
DeltaOp::cell_set(NodeId(1), vec![1]),
DeltaOp::slot_value(NodeId(2), vec![2]),
DeltaOp::NodeRemove { node: NodeId(5) },
DeltaOp::EdgeAdd {
dependent: NodeId(2),
dependency: NodeId(1),
},
]
);
}
#[test]
fn crdt_sync_round_trips_through_serde() {
let sync = CrdtSync::new(
vec![
(
1,
WireStamp {
wall_time: 200,
logical: 0,
peer: 1,
},
),
(
2,
WireStamp {
wall_time: 180,
logical: 3,
peer: 2,
},
),
],
vec![
CrdtOp::new(
NodeId(1),
WireStamp {
wall_time: 200,
logical: 0,
peer: 1,
},
vec![10, 20],
),
CrdtOp::keyed(
NodeId(2),
NodeKey::new("scores/alice").unwrap(),
WireStamp {
wall_time: 180,
logical: 3,
peer: 2,
},
vec![30],
),
],
);
let json = serde_json::to_string(&IpcMessage::CrdtSync(sync.clone())).unwrap();
let back: IpcMessage = serde_json::from_str(&json).unwrap();
assert_eq!(back, IpcMessage::CrdtSync(sync));
}
#[test]
fn crdt_sync_filter_omits_non_readable_ops_but_keeps_frontier() {
let frontier = vec![
(
1,
WireStamp {
wall_time: 200,
logical: 0,
peer: 1,
},
),
(
2,
WireStamp {
wall_time: 200,
logical: 0,
peer: 2,
},
),
];
let sync = CrdtSync::new(
frontier.clone(),
vec![
CrdtOp::new(
NodeId(1),
WireStamp {
wall_time: 1,
logical: 0,
peer: 1,
},
vec![1],
),
CrdtOp::new(
NodeId(2),
WireStamp {
wall_time: 2,
logical: 0,
peer: 1,
},
vec![2],
),
CrdtOp::new(
NodeId(3),
WireStamp {
wall_time: 3,
logical: 0,
peer: 1,
},
vec![3],
),
],
);
let mut permissions = PeerPermissions::new();
permissions.allow_many(PEER_A, OpKind::Read, [NodeId(1), NodeId(3)]);
let filtered = sync.filter_readable(&permissions, PEER_A);
assert_eq!(
filtered.ops.iter().map(|op| op.node).collect::<Vec<_>>(),
vec![NodeId(1), NodeId(3)],
);
assert_eq!(filtered.frontier, frontier);
}
#[test]
fn crdt_sync_frontier_suppress_omits_empty_frontier() {
let ops = vec![CrdtOp::new(
NodeId(7),
WireStamp {
wall_time: 13,
logical: 0,
peer: 1,
},
vec![7],
)];
let suppressed = CrdtSync::ops_only(ops.clone());
assert!(suppressed.is_frontier_suppressed());
let json = serde_json::to_string(&IpcMessage::CrdtSync(suppressed.clone())).unwrap();
assert!(
!json.contains("\"frontier\""),
"frontier should be omitted, got: {json}"
);
let back: IpcMessage = serde_json::from_str(&json).unwrap();
assert_eq!(back, IpcMessage::CrdtSync(suppressed));
}
#[test]
fn crdt_sync_with_frontier_still_serializes_frontier() {
let sync = CrdtSync::new(
vec![(
1,
WireStamp {
wall_time: 5,
logical: 0,
peer: 1,
},
)],
vec![CrdtOp::new(
NodeId(1),
WireStamp {
wall_time: 5,
logical: 0,
peer: 1,
},
vec![1],
)],
);
let json = serde_json::to_string(&IpcMessage::CrdtSync(sync.clone())).unwrap();
assert!(json.contains("\"frontier\""), "frontier should be present");
let back: IpcMessage = serde_json::from_str(&json).unwrap();
assert_eq!(back, IpcMessage::CrdtSync(sync));
}
#[cfg(feature = "json-base64")]
#[test]
fn json_base64_round_trips_and_shrinks_payload() {
let payload = vec![0xACu8; 256];
let snapshot = Snapshot::new(
1,
vec![NodeSnapshot::payload(NodeId(1), "bytes", payload.clone())],
vec![EdgeSnapshot::new(NodeId(1), NodeId(1))],
vec![NodeId(1)],
);
let msg = IpcMessage::Snapshot(snapshot);
let canonical = msg.encode_json().unwrap();
let base64_encoded = msg.encode_json_base64().unwrap();
assert!(
base64_encoded.len() < canonical.len() / 2,
"base64 {} should be < half of canonical {}",
base64_encoded.len(),
canonical.len()
);
let back = IpcMessage::decode_json_base64(&base64_encoded).unwrap();
assert_eq!(back, msg);
}
#[test]
fn json_intern_round_trips_and_dedups_type_tags() {
let nodes: Vec<NodeSnapshot> = (0..64)
.map(|i| {
NodeSnapshot::payload(
NodeId(i + 1),
if i % 2 == 0 { "alpha" } else { "beta" },
vec![i as u8],
)
})
.collect();
let roots: Vec<NodeId> = nodes.iter().map(|n| n.node).collect();
let snapshot = Snapshot::new(1, nodes, vec![], roots);
let msg = IpcMessage::Snapshot(snapshot);
let canonical = serde_json::to_vec(&msg).unwrap();
let interned = msg.encode_json_intern().unwrap();
let interned_str = String::from_utf8(interned.clone()).unwrap();
assert!(
interned_str.contains("\"intern\""),
"intern table should be present"
);
assert!(
interned.len() < canonical.len(),
"interned {} should be < canonical {}",
interned.len(),
canonical.len()
);
let back = IpcMessage::decode_json_intern(&interned).unwrap();
assert_eq!(back, msg);
}
#[test]
fn shm_blob_arena_round_trips_payload_by_descriptor() {
let mut arena = ShmBlobArena::with_capacity(SHM_BLOB_HEADER_LEN + 128).unwrap();
let blob = arena.write_blob(12, b"large context pack").unwrap();
assert_eq!(arena.read_blob(blob).unwrap(), b"large context pack");
assert_eq!(blob.epoch, 12);
assert_eq!(blob.len, "large context pack".len() as u64);
}
#[test]
fn shm_blob_arena_rejects_oversized_payload() {
let mut arena = ShmBlobArena::with_capacity(SHM_BLOB_HEADER_LEN + 4).unwrap();
let err = arena.write_blob(1, b"12345").unwrap_err();
assert_eq!(err, ShmBlobArenaError::BlobTooLarge { len: 5, max_len: 4 });
}
#[test]
fn shm_blob_arena_wrap_rejects_stale_descriptor() {
let mut arena = ShmBlobArena::with_capacity((SHM_BLOB_HEADER_LEN * 2) + 8).unwrap();
let old = arena.write_blob(1, b"old").unwrap();
let _middle = arena.write_blob(2, b"abcd").unwrap();
let _new = arena.write_blob(3, b"new").unwrap();
let err = arena.read_blob(old).unwrap_err();
assert!(matches!(
err,
ShmBlobArenaError::DescriptorMismatch {
field: "generation"
} | ShmBlobArenaError::DescriptorMismatch { field: "checksum" }
));
}
#[test]
fn shm_blob_arena_rejects_torn_payload() {
let mut arena = ShmBlobArena::with_capacity(SHM_BLOB_HEADER_LEN + 32).unwrap();
let blob = arena.write_blob(4, b"payload").unwrap();
let payload_offset = blob.offset as usize + SHM_BLOB_HEADER_LEN;
arena.bytes_mut()[payload_offset] ^= 0xff;
let err = arena.read_blob(blob).unwrap_err();
assert!(matches!(err, ShmBlobArenaError::ChecksumMismatch { .. }));
}
#[test]
fn ipc_messages_can_reference_shared_blobs() {
let mut arena = ShmBlobArena::with_capacity(SHM_BLOB_HEADER_LEN + 128).unwrap();
let blob = arena.write_blob(9, b"large slot value").unwrap();
let snapshot = Snapshot::new(
9,
vec![NodeSnapshot::shared_blob(NodeId(7), "text/plain", blob)],
vec![],
vec![NodeId(7)],
);
let delta = Delta::next(9, vec![DeltaOp::slot_value_blob(NodeId(7), blob)]);
let snapshot_json = serde_json::to_string(&IpcMessage::Snapshot(snapshot.clone())).unwrap();
let delta_json = serde_json::to_string(&IpcMessage::Delta(delta.clone())).unwrap();
assert_eq!(
serde_json::from_str::<IpcMessage>(&snapshot_json).unwrap(),
IpcMessage::Snapshot(snapshot)
);
assert_eq!(
serde_json::from_str::<IpcMessage>(&delta_json).unwrap(),
IpcMessage::Delta(delta)
);
assert_eq!(arena.read_blob(blob).unwrap(), b"large slot value");
}
#[test]
fn ipc_message_bytes_are_channel_agnostic_payloads() {
let message = IpcMessage::Delta(Delta::next(
15,
vec![
DeltaOp::cell_set(NodeId(1), b"cell".to_vec()),
DeltaOp::slot_value(NodeId(2), b"slot".to_vec()),
],
));
let websocket_text_frame = serde_json::to_string(&message).unwrap();
let webrtc_data_frame = websocket_text_frame.as_bytes().to_vec();
let ffi_owned_buffer = webrtc_data_frame.clone();
assert_eq!(
serde_json::from_str::<IpcMessage>(&websocket_text_frame).unwrap(),
message
);
assert_eq!(
serde_json::from_slice::<IpcMessage>(&webrtc_data_frame).unwrap(),
message
);
assert_eq!(
serde_json::from_slice::<IpcMessage>(&ffi_owned_buffer).unwrap(),
message
);
}
#[test]
fn node_key_validates_path_bounds() {
assert!(NodeKey::new("scores/alice").is_ok());
assert_eq!(NodeKey::new("").unwrap_err(), NodeKeyError::Empty);
assert_eq!(
NodeKey::new("a//b").unwrap_err(),
NodeKeyError::EmptySegment
);
assert_eq!(
NodeKey::new("/leading").unwrap_err(),
NodeKeyError::EmptySegment
);
let too_many = vec!["s"; NODE_KEY_MAX_SEGMENTS + 1].join("/");
assert!(matches!(
NodeKey::new(too_many).unwrap_err(),
NodeKeyError::TooManySegments { .. }
));
let too_long = "x".repeat(2000);
assert!(matches!(
NodeKey::new(too_long).unwrap_err(),
NodeKeyError::TooLong { .. }
));
}
#[test]
fn node_key_segments_round_trip() {
let key = NodeKey::from_segments(["outer", "k1", "inner", "k2"]).unwrap();
assert_eq!(key.as_str(), "outer/k1/inner/k2");
assert_eq!(
key.segments().collect::<Vec<_>>(),
vec!["outer", "k1", "inner", "k2"]
);
}
#[test]
fn keyed_node_round_trips_through_json() {
let key = NodeKey::new("scores/alice").unwrap();
let snapshot = Snapshot::new(
1,
vec![NodeSnapshot::payload(NodeId(1), "i32", vec![1]).with_key(key.clone())],
vec![],
vec![NodeId(1)],
);
let message = IpcMessage::Snapshot(snapshot);
let json = serde_json::to_string(&message).unwrap();
assert!(json.contains("scores/alice"));
assert_eq!(
serde_json::from_str::<IpcMessage>(&json).unwrap(),
message,
"keyed snapshot must round-trip through JSON"
);
}
#[test]
fn unkeyed_node_omits_key_in_json() {
let snapshot = Snapshot::new(
1,
vec![NodeSnapshot::payload(NodeId(1), "i32", vec![1])],
vec![],
vec![NodeId(1)],
);
let message = IpcMessage::Snapshot(snapshot);
let json = serde_json::to_string(&message).unwrap();
assert!(
!json.contains("\"key\""),
"unkeyed node must omit the key field in JSON: {json}"
);
let delta = Delta::next(
1,
vec![DeltaOp::NodeAdd {
node: NodeId(2),
type_tag: "i32".into(),
state: NodeState::Payload(vec![2]),
key: None,
}],
);
let delta_json = serde_json::to_string(&IpcMessage::Delta(delta)).unwrap();
assert!(
!delta_json.contains("\"key\""),
"unkeyed NodeAdd must omit the key field in JSON: {delta_json}"
);
}
#[test]
fn node_with_absent_key_decodes_to_none() {
let wire = r#"{"Snapshot":{"epoch":1,"nodes":[{"node":1,"type_tag":"i32","state":{"Payload":[1]}}],"edges":[],"roots":[1]}}"#;
let IpcMessage::Snapshot(snapshot) = serde_json::from_str::<IpcMessage>(wire).unwrap() else {
panic!("expected snapshot");
};
assert_eq!(snapshot.nodes[0].key, None);
}
#[test]
fn key_index_survives_nodeid_churn() {
let key = NodeKey::new("scores/alice").unwrap();
let mut index = KeyIndex::new();
let snapshot = Snapshot::new(
1,
vec![NodeSnapshot::payload(NodeId(1), "i32", vec![1]).with_key(key.clone())],
vec![],
vec![NodeId(1)],
);
index.ingest_snapshot(&snapshot);
assert_eq!(index.node_for_key(&key), Some(NodeId(1)));
assert_eq!(index.key_for_node(NodeId(1)), Some(&key));
let delta = Delta::next(
1,
vec![
DeltaOp::NodeRemove { node: NodeId(1) },
DeltaOp::NodeAdd {
node: NodeId(2),
type_tag: "i32".into(),
state: NodeState::Payload(vec![2]),
key: Some(key.clone()),
},
],
);
index.apply_delta(&delta);
assert_eq!(index.node_for_key(&key), Some(NodeId(2)));
assert_eq!(index.key_for_node(NodeId(1)), None);
assert_eq!(index.key_for_node(NodeId(2)), Some(&key));
assert_eq!(index.len(), 1);
}
#[cfg(feature = "ipc-binary")]
mod binary {
use lazily::{
CrdtOp, CrdtSync, DecodeError, Delta, DeltaOp, EdgeSnapshot, IpcMessage, NodeId, NodeKey,
NodeSnapshot, Snapshot, WireStamp,
};
#[test]
fn ipc_message_binary_round_trip_snapshot() {
let snapshot = Snapshot::new(
7,
vec![
NodeSnapshot::payload(NodeId(1), "i32", vec![1, 2, 3]),
NodeSnapshot::opaque(NodeId(2), "opaque-type"),
],
vec![EdgeSnapshot::new(NodeId(2), NodeId(1))],
vec![NodeId(1), NodeId(2)],
);
let message = IpcMessage::Snapshot(snapshot.clone());
let encoded = message.encode_binary().unwrap();
let decoded = IpcMessage::decode_binary(&encoded).unwrap();
assert_eq!(decoded, message);
}
#[test]
fn ipc_message_binary_round_trip_delta() {
let delta = Delta::next(
3,
vec![
DeltaOp::cell_set(NodeId(1), vec![10, 20]),
DeltaOp::slot_value(NodeId(2), vec![30, 40]),
DeltaOp::invalidate(NodeId(3)),
],
);
let message = IpcMessage::Delta(delta.clone());
let encoded = message.encode_binary().unwrap();
let decoded = IpcMessage::decode_binary(&encoded).unwrap();
assert_eq!(decoded, message);
}
#[test]
fn ipc_message_binary_round_trips_keyed_and_unkeyed_nodes() {
let key = NodeKey::new("scores/alice").unwrap();
let snapshot = Snapshot::new(
7,
vec![
NodeSnapshot::payload(NodeId(1), "i32", vec![1]).with_key(key),
NodeSnapshot::opaque(NodeId(2), "opaque-type"),
],
vec![],
vec![NodeId(1), NodeId(2)],
);
let message = IpcMessage::Snapshot(snapshot);
let encoded = message.encode_binary().unwrap();
let decoded = IpcMessage::decode_binary(&encoded).unwrap();
assert_eq!(decoded, message);
}
#[test]
fn ipc_message_binary_round_trips_crdt_sync() {
let sync = CrdtSync::new(
vec![(
1,
WireStamp {
wall_time: 200,
logical: 0,
peer: 1,
},
)],
vec![
CrdtOp::new(
NodeId(1),
WireStamp {
wall_time: 200,
logical: 0,
peer: 1,
},
vec![9],
),
CrdtOp::keyed(
NodeId(2),
NodeKey::new("scores/alice").unwrap(),
WireStamp {
wall_time: 180,
logical: 1,
peer: 2,
},
vec![8, 7],
),
],
);
let message = IpcMessage::CrdtSync(sync);
let encoded = message.encode_binary().unwrap();
let decoded = IpcMessage::decode_binary(&encoded).unwrap();
assert_eq!(decoded, message);
}
#[test]
fn ipc_message_binary_rejects_invalid_bytes() {
let result = IpcMessage::decode_binary(b"garbage");
assert!(matches!(result, Err(DecodeError::Binary(_))));
}
#[test]
fn ipc_message_binary_is_smaller_than_json() {
let snapshot = Snapshot::new(
42,
vec![NodeSnapshot::payload(NodeId(1), "i32", vec![1, 2, 3, 4])],
vec![EdgeSnapshot::new(NodeId(1), NodeId(2))],
vec![NodeId(1)],
);
let message = IpcMessage::Snapshot(snapshot);
let json_len = serde_json::to_vec(&message).unwrap().len();
let binary_len = message.encode_binary().unwrap().len();
assert!(
binary_len < json_len,
"binary ({binary_len}) should be smaller than json ({json_len})"
);
}
}
#[cfg(feature = "ipc-msgpack")]
mod msgpack {
use lazily::{
CrdtOp, CrdtSync, DecodeError, Delta, DeltaOp, EdgeSnapshot, EncodeError, IpcCodec,
IpcMessage, NodeId, NodeKey, NodeSnapshot, Snapshot, WireStamp,
};
#[test]
fn ipc_message_msgpack_round_trips_snapshot() {
let snapshot = Snapshot::new(
7,
vec![
NodeSnapshot::payload(NodeId(1), "i32", vec![1, 2, 3]),
NodeSnapshot::opaque(NodeId(2), "opaque-type"),
],
vec![EdgeSnapshot::new(NodeId(2), NodeId(1))],
vec![NodeId(1), NodeId(2)],
);
let message = IpcMessage::Snapshot(snapshot);
let encoded = message.encode_msgpack().unwrap();
let decoded = IpcMessage::decode_msgpack(&encoded).unwrap();
assert_eq!(decoded, message);
assert_eq!(IpcCodec::MessagePack.name(), "msgpack");
assert_eq!(IpcCodec::MessagePack.decode(&encoded).unwrap(), message);
assert!(serde_json::from_slice::<IpcMessage>(&encoded).is_err());
}
#[test]
fn ipc_message_msgpack_round_trips_delta() {
let delta = Delta::next(
3,
vec![
DeltaOp::cell_set(NodeId(1), vec![10, 20]),
DeltaOp::slot_value(NodeId(2), vec![30, 40]),
DeltaOp::invalidate(NodeId(3)),
],
);
let message = IpcMessage::Delta(delta);
let encoded = message.encode_msgpack().unwrap();
let decoded = IpcMessage::decode_msgpack(&encoded).unwrap();
assert_eq!(decoded, message);
}
#[test]
fn ipc_message_msgpack_round_trips_keyed_and_unkeyed_nodes() {
let key = NodeKey::new("scores/alice").unwrap();
let snapshot = Snapshot::new(
7,
vec![
NodeSnapshot::payload(NodeId(1), "i32", vec![1]).with_key(key),
NodeSnapshot::opaque(NodeId(2), "opaque-type"),
],
vec![],
vec![NodeId(1), NodeId(2)],
);
let message = IpcMessage::Snapshot(snapshot);
let encoded = message.encode_msgpack().unwrap();
let decoded = IpcMessage::decode_msgpack(&encoded).unwrap();
assert_eq!(decoded, message);
}
#[test]
fn ipc_message_msgpack_round_trips_crdt_sync() {
let sync = CrdtSync::new(
vec![(
1,
WireStamp {
wall_time: 200,
logical: 0,
peer: 1,
},
)],
vec![
CrdtOp::new(
NodeId(1),
WireStamp {
wall_time: 200,
logical: 0,
peer: 1,
},
vec![9],
),
CrdtOp::keyed(
NodeId(2),
NodeKey::new("scores/alice").unwrap(),
WireStamp {
wall_time: 180,
logical: 1,
peer: 2,
},
vec![8, 7],
),
],
);
let message = IpcMessage::CrdtSync(sync);
let encoded = message.encode_msgpack().unwrap();
let decoded = IpcMessage::decode_msgpack(&encoded).unwrap();
assert_eq!(decoded, message);
}
#[test]
fn ipc_message_msgpack_is_smaller_than_json() {
let snapshot = Snapshot::new(
42,
vec![NodeSnapshot::payload(NodeId(1), "i32", vec![1, 2, 3, 4])],
vec![EdgeSnapshot::new(NodeId(1), NodeId(2))],
vec![NodeId(1)],
);
let message = IpcMessage::Snapshot(snapshot);
let json_len = serde_json::to_vec(&message).unwrap().len();
let msgpack_len = message.encode_msgpack().unwrap().len();
assert!(
msgpack_len < json_len,
"msgpack ({msgpack_len}) should be smaller than json ({json_len})"
);
}
#[test]
fn ipc_message_msgpack_rejects_invalid_bytes() {
let result = IpcMessage::decode_msgpack(b"garbage");
assert!(matches!(result, Err(DecodeError::Msgpack(_))));
}
#[test]
fn encode_decode_error_implement_display() {
let decode_err = IpcMessage::decode_msgpack(b"garbage").unwrap_err();
let _ = std::format!("{}", decode_err);
let encode_err =
EncodeError::Msgpack(rmp_serde::to_vec_named(&failing_serialize()).unwrap_err());
let _ = std::format!("{}", encode_err);
}
fn failing_serialize() -> impl serde::Serialize {
struct Failing;
impl serde::Serialize for Failing {
fn serialize<S>(&self, _serializer: S) -> Result<S::Ok, S::Error>
where
S: serde::Serializer,
{
Err(serde::ser::Error::custom("expected failure"))
}
}
Failing
}
}
#[cfg(feature = "ffi")]
mod json_codec {
use lazily::{
DecodeError, EdgeSnapshot, EncodeError, IpcMessage, NodeId, NodeSnapshot, Snapshot,
};
#[test]
fn ipc_message_json_round_trip_snapshot() {
let snapshot = Snapshot::new(
7,
vec![
NodeSnapshot::payload(NodeId(1), "i32", vec![1, 2, 3]),
NodeSnapshot::opaque(NodeId(2), "opaque-type"),
],
vec![EdgeSnapshot::new(NodeId(2), NodeId(1))],
vec![NodeId(1), NodeId(2)],
);
let message = IpcMessage::Snapshot(snapshot);
let encoded = message.encode_json().unwrap();
let decoded = IpcMessage::decode_json(&encoded).unwrap();
assert_eq!(decoded, message);
}
#[test]
fn ipc_message_json_rejects_invalid_bytes() {
let result = IpcMessage::decode_json(b"not json");
assert!(matches!(result, Err(DecodeError::Json(_))));
}
#[test]
fn encode_decode_error_implement_display() {
let decode_err = IpcMessage::decode_json(b"not json").unwrap_err();
let _ = std::format!("{}", decode_err);
let encode_err = EncodeError::Json(serde_json::from_str::<()>("bad").unwrap_err());
let _ = std::format!("{}", encode_err);
}
}
mod capability_handshake {
use super::*;
#[test]
fn new_sets_protocol_defaults() {
let hs = CapabilityHandshake::new(PEER_A, "session-1");
assert_eq!(hs.protocol_id, "lazily-ipc");
assert_eq!(hs.protocol_major_version, 1);
assert_eq!(hs.codec, "json");
assert_eq!(hs.max_frame_size, 1_048_576);
assert!(!hs.fragmentation_supported);
assert!(hs.ordered_reliable);
assert_eq!(hs.peer_id, PEER_A);
assert_eq!(hs.session_id, "session-1");
assert!(hs.features.is_empty());
}
#[test]
fn builders_configure_fields() {
let hs = CapabilityHandshake::new(PEER_A, "s")
.with_codec("msgpack")
.with_max_frame_size(2_097_152)
.with_fragmentation(true)
.with_features(["shared-blob", "signaling-relay"]);
assert_eq!(hs.codec, "msgpack");
assert_eq!(hs.max_frame_size, 2_097_152);
assert!(hs.fragmentation_supported);
assert_eq!(hs.features, ["shared-blob", "signaling-relay"]);
assert!(hs.has_feature("shared-blob"));
assert!(!hs.has_feature("crdt-cell-plane"));
}
#[test]
fn round_trips_through_serde_json() {
let hs = CapabilityHandshake::new(PEER_B, "abc-123")
.with_features(["shared-blob", "signaling-relay"]);
let json = serde_json::to_string(&hs).unwrap();
let back: CapabilityHandshake = serde_json::from_str(&json).unwrap();
assert_eq!(back, hs);
}
#[test]
fn serde_matches_protocol_wire_shape() {
let hs = CapabilityHandshake::new(PeerId(1), "abc-123")
.with_max_frame_size(1_048_576)
.with_features(["shared-blob", "signaling-relay"]);
let value: serde_json::Value =
serde_json::from_str(&serde_json::to_string(&hs).unwrap()).unwrap();
assert_eq!(value["protocol_id"], "lazily-ipc");
assert_eq!(value["protocol_major_version"], 1);
assert_eq!(value["codec"], "json");
assert_eq!(value["max_frame_size"], 1_048_576);
assert_eq!(value["fragmentation_supported"], false);
assert_eq!(value["ordered_reliable"], true);
assert_eq!(value["peer_id"], 1);
assert_eq!(value["session_id"], "abc-123");
assert_eq!(value["features"][0], "shared-blob");
}
#[test]
fn ordered_reliable_defaults_to_true_on_decode() {
let json = r#"{
"protocol_id": "lazily-ipc",
"protocol_major_version": 1,
"codec": "json",
"max_frame_size": 1024,
"peer_id": 5,
"session_id": "s"
}"#;
let hs: CapabilityHandshake = serde_json::from_str(json).unwrap();
assert!(hs.ordered_reliable);
assert!(!hs.fragmentation_supported);
assert!(hs.features.is_empty());
}
#[test]
fn compatible_handshakes_pass() {
let a = CapabilityHandshake::new(PEER_A, "s");
let b = CapabilityHandshake::new(PEER_B, "s");
assert!(a.is_compatible_with(&b));
}
#[test]
fn wrong_protocol_id_fails_closed() {
let mut a = CapabilityHandshake::new(PEER_A, "s");
a.protocol_id = "other".to_owned();
let b = CapabilityHandshake::new(PEER_B, "s");
assert!(!a.is_compatible_with(&b));
}
#[test]
fn major_version_mismatch_fails_closed() {
let mut a = CapabilityHandshake::new(PEER_A, "s");
a.protocol_major_version = 2;
let b = CapabilityHandshake::new(PEER_B, "s");
assert!(!a.is_compatible_with(&b));
assert!(!b.is_compatible_with(&a));
}
#[test]
fn codec_mismatch_fails_closed() {
let a = CapabilityHandshake::new(PEER_A, "s").with_codec("json");
let b = CapabilityHandshake::new(PEER_B, "s").with_codec("postcard");
assert!(!a.is_compatible_with(&b));
}
#[test]
fn unordered_reliable_fails_closed() {
let mut a = CapabilityHandshake::new(PEER_A, "s");
a.ordered_reliable = false;
let b = CapabilityHandshake::new(PEER_B, "s");
assert!(!a.is_compatible_with(&b));
assert!(!b.is_compatible_with(&a));
}
#[test]
fn fragmentation_and_features_do_not_block_compatibility() {
let a = CapabilityHandshake::new(PEER_A, "s")
.with_fragmentation(true)
.with_features(["shared-blob"]);
let b = CapabilityHandshake::new(PEER_B, "s")
.with_fragmentation(false)
.with_features(["signaling-relay"]);
assert!(a.is_compatible_with(&b));
}
#[test]
fn delta_since_request_round_trips() {
let req = DeltaSinceRequest::new(vec![
(
1,
WireStamp {
wall_time: 100,
logical: 2,
peer: 1,
},
),
(
2,
WireStamp {
wall_time: 90,
logical: 5,
peer: 2,
},
),
]);
let json = serde_json::to_string(&IpcMessage::DeltaSinceRequest(req.clone())).unwrap();
assert!(json.contains("DeltaSinceRequest"));
let back: IpcMessage = serde_json::from_str(&json).unwrap();
assert_eq!(back, IpcMessage::DeltaSinceRequest(req));
}
#[test]
fn delta_since_request_is_control() {
let req = DeltaSinceRequest::new(vec![]);
let msg = IpcMessage::DeltaSinceRequest(req);
assert!(msg.is_control());
}
}