#![cfg(feature = "ipc")]
use lazily::{
Delta, DeltaApplyStatus, DeltaOp, EdgeSnapshot, IpcMessage, NodeId, NodeSnapshot, NodeState,
OpKind, PeerId, PeerPermissions, RemoteOp, SHM_BLOB_HEADER_LEN, ShmBlobArena,
ShmBlobArenaError, Snapshot,
};
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]),
},
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]),
},
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 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
);
}