use crate::sync::helpers::setup;
use eidetica::Result;
use eidetica::auth::crypto::{format_public_key, generate_keypair};
use eidetica::sync::peer_types::{Address, PeerStatus};
use std::time::Duration;
use tokio::time::sleep;
#[tokio::test]
async fn test_sync_queue_operations() -> Result<()> {
let (_instance, sync) = setup();
let (_, verifying_key) = generate_keypair();
let peer_pubkey = format_public_key(&verifying_key);
let queue = sync.sync_queue();
assert_eq!(queue.total_entries(), 0, "Queue should start empty");
let entry_id1 = eidetica::Entry::builder("entry1".to_string())
.build()
.id()
.clone();
let entry_id2 = eidetica::Entry::builder("entry2".to_string())
.build()
.id()
.clone();
let tree_id = eidetica::Entry::builder("test_tree".to_string())
.build()
.id()
.clone();
assert!(
queue.queue_entry(peer_pubkey, &entry_id1, &tree_id)?,
"First entry should be added"
);
assert!(
queue.queue_entry(peer_pubkey, &entry_id2, &tree_id)?,
"Second entry should be added"
);
assert!(
!queue.queue_entry(peer_pubkey, &entry_id1, &tree_id)?,
"Duplicate entry should not be added"
);
let pending = queue.get_pending_entries(peer_pubkey)?;
assert_eq!(pending.len(), 2, "Should have 2 pending entries");
Ok(())
}
#[tokio::test]
async fn test_flush_worker_processes_queue() -> Result<()> {
let (_instance, mut sync) = setup();
let (_, verifying_key) = generate_keypair();
let peer_pubkey = format_public_key(&verifying_key);
sync.register_peer(peer_pubkey, Some("Test Peer"))?;
sync.update_peer_status(peer_pubkey, PeerStatus::Active)?;
sync.add_peer_address(peer_pubkey, Address::http("127.0.0.1:8080"))?;
sync.start_flush_worker_async().await?;
assert!(
sync.is_flush_worker_running_async().await,
"Flush worker should be running"
);
let queue = sync.sync_queue();
for i in 0..3 {
let entry_id = eidetica::Entry::builder(format!("entry_{i}"))
.build()
.id()
.clone();
let tree_id_obj = eidetica::Entry::builder("test_tree".to_string())
.build()
.id()
.clone();
queue.queue_entry(peer_pubkey, &entry_id, &tree_id_obj)?;
}
let pending_before = queue.get_pending_entries(peer_pubkey)?;
assert_eq!(pending_before.len(), 3, "Should have 3 pending entries");
sleep(Duration::from_secs(2)).await;
sync.stop_flush_worker_async().await?;
assert!(
!sync.is_flush_worker_running_async().await,
"Flush worker should be stopped"
);
Ok(())
}
#[tokio::test]
async fn test_flush_worker_lifecycle() -> Result<()> {
let (_instance, mut sync) = setup();
assert!(
!sync.is_flush_worker_running_async().await,
"Worker should not be running initially"
);
sync.start_flush_worker_async().await?;
assert!(
sync.is_flush_worker_running_async().await,
"Worker should be running after start"
);
let result = sync.start_flush_worker_async().await;
assert!(result.is_err(), "Starting worker twice should fail");
sync.stop_flush_worker_async().await?;
assert!(
!sync.is_flush_worker_running_async().await,
"Worker should be stopped"
);
sync.stop_flush_worker_async().await?;
assert!(
!sync.is_flush_worker_running_async().await,
"Worker should still be stopped"
);
Ok(())
}
#[tokio::test]
async fn test_queue_with_peer_management() -> Result<()> {
let (_instance, mut sync) = setup();
let (_, verifying_key) = generate_keypair();
let peer_pubkey = format_public_key(&verifying_key);
sync.register_peer(peer_pubkey, Some("Test Peer"))?;
let queue = sync.sync_queue();
let entry_id = eidetica::Entry::builder("entry1".to_string())
.build()
.id()
.clone();
let tree_id = eidetica::Entry::builder("test_tree".to_string())
.build()
.id()
.clone();
queue.queue_entry(peer_pubkey, &entry_id, &tree_id)?;
let pending = queue.get_pending_entries(peer_pubkey)?;
assert_eq!(pending.len(), 1, "Should have 1 pending entry");
queue.cleanup_failed_entries(0)?; let pending_after = queue.get_pending_entries(peer_pubkey)?;
assert_eq!(
pending_after.len(),
0,
"Entry with 0 attempts exceeds max_retries=0"
);
queue.queue_entry(peer_pubkey, &entry_id, &tree_id)?;
queue.cleanup_failed_entries(3)?; let pending_final = queue.get_pending_entries(peer_pubkey)?;
assert_eq!(
pending_final.len(),
1,
"Entry with 0 attempts should remain with max_retries=3"
);
Ok(())
}
#[tokio::test]
async fn test_queue_size_triggers_flush() -> Result<()> {
let (_instance, mut sync) = setup();
let (_, verifying_key) = generate_keypair();
let peer_pubkey = format_public_key(&verifying_key);
sync.register_peer(peer_pubkey, Some("Test Peer"))?;
sync.update_peer_status(peer_pubkey, PeerStatus::Active)?;
let queue = sync.sync_queue();
let config = queue.config();
let max_size = config.max_queue_size;
for i in 0..max_size {
let entry_id = eidetica::Entry::builder(format!("entry_{i}"))
.build()
.id()
.clone();
let tree_id = eidetica::Entry::builder("test_tree".to_string())
.build()
.id()
.clone();
queue.queue_entry(peer_pubkey, &entry_id, &tree_id)?;
}
let peers_needing_flush = queue.get_peers_needing_flush()?;
assert!(
peers_needing_flush.contains(&peer_pubkey.to_string()),
"Peer should need flushing when queue size limit reached"
);
Ok(())
}