//! Unit-level coverage for `SyncPeer` and pipe halves over the mmap backend. //! //! These tests exercise variants that don't require socket setup, lifting //! `clawsync-transport/src/peer.rs` coverage off the floor. TCP/QUIC variants //! are already covered by `e2e_sync.rs` and `e2e_quic_sync.rs`. use clawsync_transport::{ MmapChannel, MmapReceiver, MmapSender, PipeReadHalf, PipeWriteHalf, SyncMessage, SyncPeer, }; use tempfile::tempdir; const RING: usize = 64 * 1024; fn build_mmap_pair() -> (MmapSender, MmapReceiver, tempfile::TempDir) { let dir = tempdir().expect("tempdir"); let path = dir.path().join("ch.mmap"); let _ = MmapChannel::create(&path, RING).expect("create channel"); let tx = MmapSender::open(&path).expect("open sender"); let rx = MmapReceiver::open(&path).expect("open receiver"); (tx, rx, dir) } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn mmap_peer_send_recv_roundtrip() { let (send_ch, recv_ch, _dir) = build_mmap_pair(); let mut peer = SyncPeer::Mmap { send_ch, recv_ch }; peer.send(&SyncMessage::Ack { revision: 7 }) .await .expect("send"); match peer.recv().await.expect("recv") { SyncMessage::Ack { revision } => assert_eq!(revision, 7), other => panic!("unexpected message: {other:?}"), } } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn mmap_peer_shutdown_is_noop_ok() { let (send_ch, recv_ch, _dir) = build_mmap_pair(); let peer = SyncPeer::Mmap { send_ch, recv_ch }; peer.shutdown().await.expect("mmap shutdown is Ok(())"); } #[test] fn mmap_peer_quic_helpers_return_none() { let dir = tempdir().unwrap(); let path = dir.path().join("ch.mmap"); let _ = MmapChannel::create(&path, RING).unwrap(); let peer = SyncPeer::Mmap { send_ch: MmapSender::open(&path).unwrap(), recv_ch: MmapReceiver::open(&path).unwrap(), }; assert!(peer.as_stream_peer().is_none()); assert!(peer.quic_conn_clone().is_none()); // close_quic is a no-op for non-QUIC peers — it must not panic. peer.close_quic(); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn mmap_peer_into_pipe_halves_roundtrip() { let (send_ch, recv_ch, _dir) = build_mmap_pair(); let peer = SyncPeer::Mmap { send_ch, recv_ch }; let (mut read_half, mut write_half) = peer.into_pipe_halves(); assert!(matches!(read_half, PipeReadHalf::Mmap(_))); assert!(matches!(write_half, PipeWriteHalf::Mmap(_))); write_half .send(&SyncMessage::Ack { revision: 42 }) .await .expect("pipe send"); match read_half.recv().await.expect("pipe recv") { SyncMessage::Ack { revision } => assert_eq!(revision, 42), other => panic!("unexpected: {other:?}"), } } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn mmap_pipe_write_half_shutdown_is_ok() { let (send_ch, recv_ch, _dir) = build_mmap_pair(); let (_r, w) = SyncPeer::Mmap { send_ch, recv_ch }.into_pipe_halves(); w.shutdown().await.expect("mmap pipe shutdown"); } #[tokio::test(flavor = "multi_thread", worker_threads = 2)] async fn mmap_pipe_write_half_shutdown_push_is_ok() { let (send_ch, recv_ch, _dir) = build_mmap_pair(); let (_r, w) = SyncPeer::Mmap { send_ch, recv_ch }.into_pipe_halves(); // For non-QUIC transports, shutdown_push() delegates to shutdown(). w.shutdown_push().await.expect("mmap pipe shutdown_push"); }