#![allow(clippy::needless_return)]
use beam::Node;
use beam::ack::{AckPolicy, ReplicationStatus};
use beam::actor::Actor;
use beam::adapters::MemoryStorage;
use std::time::Duration;
#[tokio::test]
async fn e2e_put_quorum_succeeds_with_local_ack() {
let storage: Vec<Box<dyn Actor>> = vec![Box::new(MemoryStorage::new()) as Box<dyn Actor>];
let mut node = Node::new_with_config(Default::default(), storage, vec![]);
let result: Result<ReplicationStatus, String> = node
.get("e2e_quorum_local_key")
.put_quorum("e2e_quorum_local_value".into(), AckPolicy::any())
.await;
let status = result.expect("single-node local put_quorum with Any should succeed");
assert_eq!(status.acked_by, 1, "local ack counts as 1 peer");
assert!(
status.quorum_met,
"local ack satisfies AckPolicy::any() (quorum=1)"
);
assert_eq!(status.put_id.len(), 8, "put_id is an 8-char random string");
assert!(
status.elapsed <= Duration::from_secs(9),
"elapsed ({:?}) should be within AckPolicy::any() timeout (9s)",
status.elapsed
);
let got = node
.get("e2e_quorum_local_key")
.once(Some(Duration::from_secs(2)))
.await;
assert_eq!(
got,
Some(beam::types::Value::Text(
"e2e_quorum_local_value".to_string()
)),
"value should be readable after put_quorum resolves"
);
}
#[tokio::test]
async fn e2e_put_quorum_majority_policy_completes_single_node() {
let storage: Vec<Box<dyn Actor>> = vec![Box::new(MemoryStorage::new()) as Box<dyn Actor>];
let mut node = Node::new_with_config(Default::default(), storage, vec![]);
let policy = AckPolicy::for_peer_count(5);
let start = std::time::Instant::now();
let result: Result<ReplicationStatus, String> = node
.get("e2e_quorum_majority_key")
.put_quorum("e2e_quorum_majority_value".into(), policy)
.await;
let elapsed = start.elapsed();
let status =
result.expect("single-node put_quorum with Majority policy completes via local ack");
assert_eq!(status.acked_by, 1, "local ack counts as 1 peer");
assert!(
status.quorum_met,
"drain completes with quorum_met=true (local ack path)"
);
assert!(
elapsed < Duration::from_secs(2),
"single-node drain should complete quickly, got {elapsed:?}"
);
let got = node
.get("e2e_quorum_majority_key")
.once(Some(Duration::from_secs(2)))
.await;
assert_eq!(
got,
Some(beam::types::Value::Text(
"e2e_quorum_majority_value".to_string()
)),
"value should be readable after put_quorum resolves"
);
}
#[tokio::test]
async fn e2e_put_quorum_all_policy_completes_single_node() {
let storage: Vec<Box<dyn Actor>> = vec![Box::new(MemoryStorage::new()) as Box<dyn Actor>];
let mut node = Node::new_with_config(Default::default(), storage, vec![]);
let policy = AckPolicy::all();
let start = std::time::Instant::now();
let result: Result<ReplicationStatus, String> = node
.get("e2e_quorum_all_key")
.put_quorum("e2e_quorum_all_value".into(), policy)
.await;
let elapsed = start.elapsed();
let status = result.expect("single-node put_quorum with All policy completes via local ack");
assert_eq!(status.acked_by, 1, "local ack counts as 1 peer");
assert!(
status.quorum_met,
"drain completes with quorum_met=true (local ack path)"
);
assert!(
elapsed < Duration::from_secs(2),
"single-node drain should complete quickly, got {elapsed:?}"
);
}