use communitas_core::crdt::*;
use communitas_core::message_sync::MessageSyncService;
use communitas_core::test_harness::TestHarness;
use std::time::Duration;
#[tokio::test]
#[ignore] async fn test_two_peers_sync_stream() {
let harness = TestHarness::new(2).await.expect("harness creation failed");
harness.mesh().await.expect("mesh setup failed");
harness
.wait_until_connected(1, Duration::from_secs(5))
.await
.expect("connection failed");
let node_a = harness.get_node(0).await.expect("node 0 not found");
let node_b = harness.get_node(1).await.expect("node 1 not found");
let peer_id_a = node_a.read().await.four_words.clone();
let peer_id_b = node_b.read().await.four_words.clone();
let sync_a = MessageSyncService::new(peer_id_a.clone());
let sync_b = MessageSyncService::new(peer_id_b.clone());
let entity_id = "contact-alice";
for i in 0..10 {
sync_a
.send_message(
entity_id.to_string(),
EntityType::Person,
MessageContent {
text: format!("Message {}", i),
author: "Peer A".to_string(),
attachments: None,
},
None,
)
.await
.expect("send failed");
}
let messages_a = sync_a
.get_all_messages(entity_id)
.await
.expect("get messages failed");
let result = sync_b
.handle_sync_response(messages_a)
.await
.expect("sync failed");
assert_eq!(result.messages_added, 10, "Should add 10 messages");
let messages_b = sync_b
.get_messages(entity_id)
.await
.expect("get messages failed");
assert_eq!(messages_b.len(), 10, "Peer B should have 10 messages");
for i in 0..9 {
let ord = messages_b[i]
.metadata
.vector_clock
.compare(&messages_b[i + 1].metadata.vector_clock);
assert!(
!matches!(ord, ClockOrdering::After),
"Messages should be in causal order"
);
}
harness.cleanup().await.expect("cleanup failed");
}
#[tokio::test]
#[ignore] async fn test_missing_range_repair() {
let harness = TestHarness::new(2).await.expect("harness creation failed");
harness.set_loss(0, 1, 0.4).await;
harness.mesh().await.expect("mesh setup failed");
harness
.wait_until_connected(1, Duration::from_secs(10))
.await
.expect("connection failed");
let node_a = harness.get_node(0).await.expect("node 0 not found");
let node_b = harness.get_node(1).await.expect("node 1 not found");
let peer_id_a = node_a.read().await.four_words.clone();
let peer_id_b = node_b.read().await.four_words.clone();
let sync_a = MessageSyncService::new(peer_id_a);
let sync_b = MessageSyncService::new(peer_id_b);
let entity_id = "contact-test";
for i in 0..20 {
sync_a
.send_message(
entity_id.to_string(),
EntityType::Person,
MessageContent {
text: format!("Message {}", i),
author: "Peer A".to_string(),
attachments: None,
},
None,
)
.await
.expect("send failed");
}
let messages_a = sync_a
.get_all_messages(entity_id)
.await
.expect("get failed");
let _ = sync_b.handle_sync_response(messages_a.clone()).await;
let state_b = sync_b
.get_sync_state(entity_id)
.await
.expect("state failed");
if state_b.message_count < 20 {
assert!(
!state_b.missing_messages.is_empty(),
"Should detect missing messages"
);
let _repair_result = sync_b
.handle_sync_response(messages_a)
.await
.expect("repair failed");
let final_messages = sync_b.get_messages(entity_id).await.expect("get failed");
assert_eq!(final_messages.len(), 20, "Should repair all messages");
}
harness.cleanup().await.expect("cleanup failed");
}
#[tokio::test]
#[ignore] async fn test_convergence_after_partition() {
let harness = TestHarness::new(3).await.expect("harness creation failed");
harness.mesh().await.expect("mesh setup failed");
harness
.wait_until_connected(3, Duration::from_secs(10))
.await
.expect("initial connection failed");
let peer_ids: Vec<String> = {
let mut ids = Vec::new();
for i in 0..3 {
let node = harness.get_node(i).await.expect("node not found");
ids.push(node.read().await.four_words.clone());
}
ids
};
let sync_services: Vec<MessageSyncService> = peer_ids
.iter()
.map(|id| MessageSyncService::new(id.clone()))
.collect();
let entity_id = "project-communitas";
harness
.partition(&[0], &[1, 2])
.await
.expect("partition failed");
for i in 0..3 {
sync_services[0]
.send_message(
entity_id.to_string(),
EntityType::Project,
MessageContent {
text: format!("Partition A message {}", i),
author: "Node 0".to_string(),
attachments: None,
},
None,
)
.await
.expect("send failed");
}
for i in 0..2 {
sync_services[1]
.send_message(
entity_id.to_string(),
EntityType::Project,
MessageContent {
text: format!("Partition B message {}", i),
author: "Node 1".to_string(),
attachments: None,
},
None,
)
.await
.expect("send failed");
}
let messages_1 = sync_services[1]
.get_all_messages(entity_id)
.await
.expect("get failed");
sync_services[2]
.handle_sync_response(messages_1)
.await
.expect("sync failed");
harness.heal().await.expect("heal failed");
harness
.wait_until_connected(3, Duration::from_secs(10))
.await
.expect("reconnection failed");
for i in 0..3 {
for j in 0..3 {
if i != j {
let messages = sync_services[i]
.get_all_messages(entity_id)
.await
.expect("get failed");
sync_services[j]
.handle_sync_response(messages)
.await
.expect("sync failed");
}
}
}
for (idx, sync) in sync_services.iter().enumerate() {
let messages = sync.get_messages(entity_id).await.expect("get failed");
assert_eq!(
messages.len(),
5,
"Node {} should have 5 messages after convergence",
idx
);
}
harness.cleanup().await.expect("cleanup failed");
}
#[tokio::test]
#[ignore] async fn test_three_peer_bidirectional_sync() {
let harness = TestHarness::new(3).await.expect("harness creation failed");
harness.mesh().await.expect("mesh setup failed");
harness
.wait_until_connected(3, Duration::from_secs(10))
.await
.expect("connection failed");
let peer_ids: Vec<String> = {
let mut ids = Vec::new();
for i in 0..3 {
let node = harness.get_node(i).await.expect("node not found");
ids.push(node.read().await.four_words.clone());
}
ids
};
let sync_services: Vec<MessageSyncService> = peer_ids
.iter()
.map(|id| MessageSyncService::new(id.clone()))
.collect();
let entity_id = "contact-bob";
for (idx, sync) in sync_services.iter().enumerate() {
for i in 0..2 {
sync.send_message(
entity_id.to_string(),
EntityType::Person,
MessageContent {
text: format!("Peer {} message {}", idx, i),
author: format!("Peer {}", idx),
attachments: None,
},
None,
)
.await
.expect("send failed");
}
}
for i in 0..3 {
for j in 0..3 {
if i != j {
let messages = sync_services[i]
.get_all_messages(entity_id)
.await
.expect("get failed");
sync_services[j]
.handle_sync_response(messages)
.await
.expect("sync failed");
}
}
}
for (idx, sync) in sync_services.iter().enumerate() {
let messages = sync.get_messages(entity_id).await.expect("get failed");
assert_eq!(messages.len(), 6, "Peer {} should have 6 messages", idx);
}
harness.cleanup().await.expect("cleanup failed");
}
#[tokio::test]
async fn test_out_of_order_local() {
let sender = MessageSyncService::new("peer-sender".to_string());
let receiver = MessageSyncService::new("peer-receiver".to_string());
let entity_id = "contact-test";
let mut messages = Vec::new();
for i in 0..3 {
let msg = sender
.send_message(
entity_id.to_string(),
EntityType::Person,
MessageContent {
text: format!("Message {}", i),
author: "Sender".to_string(),
attachments: None,
},
None,
)
.await
.expect("send failed");
messages.push(msg);
}
receiver
.receive_message(messages[0].clone())
.await
.expect("receive 0 failed");
let result = receiver
.receive_message(messages[2].clone())
.await
.expect("receive 2 failed");
assert!(
!result.accepted,
"Message 2 should be rejected as out of order"
);
assert!(result.out_of_order, "Should be marked out of order");
let state = receiver
.get_sync_state(entity_id)
.await
.expect("state failed");
assert_eq!(
state.out_of_order_messages.len(),
1,
"Should have 1 pending message"
);
receiver
.receive_message(messages[1].clone())
.await
.expect("receive 1 failed");
let final_messages = receiver.get_messages(entity_id).await.expect("get failed");
assert_eq!(final_messages.len(), 3, "Should have all 3 messages");
}
#[tokio::test]
async fn test_duplicate_detection() {
let sync = MessageSyncService::new("peer-test".to_string());
let entity_id = "contact-alice";
let msg = sync
.send_message(
entity_id.to_string(),
EntityType::Person,
MessageContent {
text: "Test".to_string(),
author: "Test".to_string(),
attachments: None,
},
None,
)
.await
.expect("send failed");
let _ = sync.receive_message(msg.clone()).await;
let messages = sync.get_messages(entity_id).await.expect("get failed");
assert_eq!(messages.len(), 1, "Duplicates should be ignored");
}