//! Integration tests for the TCP transport layer. //! //! Tests that two TidalDb instances can replicate signals over TCP. #![allow(clippy::unwrap_used)] use std::collections::HashMap; use std::net::{SocketAddr, TcpListener}; use std::sync::Arc; use std::time::Duration; use tidaldb::TidalDb; use tidaldb::db::config::{NodeConfig, NodeRole}; use tidaldb::replication::shard::ShardRouter; use tidaldb::replication::{RegionId, ShardId, TcpTransport, Transport}; use tidaldb::schema::{DecaySpec, EntityId, SchemaBuilder, Timestamp, Window}; fn make_schema() -> tidaldb::schema::Schema { let mut builder = SchemaBuilder::new(); let _ = builder .signal( "view", tidaldb::schema::EntityKind::Item, DecaySpec::Exponential { half_life: Duration::from_secs(7 * 24 * 3600), }, ) .windows(&[Window::AllTime]) .velocity(false) .add(); builder.build().unwrap() } /// Pick an ephemeral port by binding to :0 and returning the address. fn ephemeral_addr() -> SocketAddr { let listener = TcpListener::bind("127.0.0.1:0").unwrap(); listener.local_addr().unwrap() } #[test] fn tcp_transport_send_and_receive() { let addr_a = ephemeral_addr(); let addr_b = ephemeral_addr(); let mut peers_a = HashMap::new(); peers_a.insert(ShardId(1), addr_b); let mut peers_b = HashMap::new(); peers_b.insert(ShardId(0), addr_a); let transport_a = Arc::new(TcpTransport::new(ShardId(0), addr_a, peers_a).unwrap()); let transport_b = Arc::new(TcpTransport::new(ShardId(1), addr_b, peers_b).unwrap()); // Send from A to B. let payload = tidaldb::replication::WalSegmentPayload { id: tidaldb::replication::WalSegmentId::new(RegionId::SINGLE, ShardId(0), 1), bytes: vec![0xDE; 512], event_count: 3, }; transport_a .send_segment(ShardId(1), payload) .expect("send should succeed"); // Receive on B. let received = transport_b.recv_segment().expect("should receive segment"); assert_eq!(received.id.seqno, 1); assert_eq!(received.event_count, 3); assert_eq!(received.bytes.len(), 512); } #[test] fn tcp_transport_multiple_segments_fifo() { let addr_a = ephemeral_addr(); let addr_b = ephemeral_addr(); let mut peers_a = HashMap::new(); peers_a.insert(ShardId(1), addr_b); let transport_a = Arc::new(TcpTransport::new(ShardId(0), addr_a, peers_a).unwrap()); let transport_b = Arc::new(TcpTransport::new(ShardId(1), addr_b, HashMap::new()).unwrap()); for seq in 1..=5u64 { let payload = tidaldb::replication::WalSegmentPayload { id: tidaldb::replication::WalSegmentId::new(RegionId::SINGLE, ShardId(0), seq), bytes: vec![seq as u8; 64], event_count: seq, }; transport_a.send_segment(ShardId(1), payload).unwrap(); } for expected_seq in 1..=5u64 { let received = transport_b.recv_segment().unwrap(); assert_eq!(received.id.seqno, expected_seq); assert_eq!(received.event_count, expected_seq); } } #[test] fn leader_follower_signal_replication_over_tcp() { let schema = make_schema(); let leader_addr = ephemeral_addr(); let follower_addr = ephemeral_addr(); // Create TCP transports. let mut leader_peers = HashMap::new(); leader_peers.insert(ShardId(1), follower_addr); let leader_transport = Arc::new(TcpTransport::new(ShardId(0), leader_addr, leader_peers).unwrap()); let follower_transport = Arc::new(TcpTransport::new(ShardId(1), follower_addr, HashMap::new()).unwrap()); // Open leader DB with data dir. let leader_dir = tempfile::tempdir().unwrap(); let leader_db = TidalDb::builder() .with_schema(schema.clone()) .with_data_dir(leader_dir.path()) .with_cluster(NodeConfig { role: NodeRole::Leader, shard_id: ShardId(0), region_id: RegionId::SINGLE, peer_shards: vec![ShardId(1)], router: ShardRouter::single(), }) .with_transport(leader_transport as Arc) .open() .unwrap(); // Open follower DB with data dir. let follower_dir = tempfile::tempdir().unwrap(); let follower_db = TidalDb::builder() .with_schema(schema) .with_data_dir(follower_dir.path()) .with_cluster(NodeConfig { role: NodeRole::Follower, shard_id: ShardId(1), region_id: RegionId::SINGLE, peer_shards: vec![ShardId(0)], router: ShardRouter::single(), }) .with_transport(follower_transport as Arc) .open() .unwrap(); // Write signals on the leader. for i in 1..=10u64 { leader_db .signal("view", EntityId::new(i), 1.0, Timestamp::now()) .unwrap(); } // The leader ships sealed WAL segments on a 2s poll interval. // Write enough signals to force a segment seal (segments seal when a new // segment starts), then wait for the shipper to poll. // // Note: In practice, WAL segment replication requires the WAL to produce // sealed segments. The shipper only ships segments that have a newer segment // after them. This test validates the transport layer works end-to-end; // full replication integration depends on the WAL producing multiple segments. // Verify both databases are healthy. leader_db.health_check().unwrap(); follower_db.health_check().unwrap(); // Verify the leader has the signals we wrote. for i in 1..=10u64 { let score = leader_db.read_decay_score(EntityId::new(i), "view", 0); assert!(score.is_ok(), "leader should have signal for entity {i}"); } leader_db.close().unwrap(); follower_db.close().unwrap(); }