#![cfg(feature = "noxu")]
use std::collections::HashSet;
use std::sync::{Arc, Mutex};
use dyniak::bucket_props::{BucketProps, BucketPropsRegistry};
use dyniak::proto::pb::{MessageCode, RpbContent, RpbGetReq, RpbPutReq};
use dyniak::quorum::{QUORUM_ALL, QUORUM_ONE};
use dyniak::replication::{RingPoint, RingView};
use dyniak::router::{
BucketRouter, PeerOp, PeerOutbound, RoutingHooks, ACK_STORED, ACK_STORED_DURABLE,
};
use dynomite::embed::hooks::BoxFuture;
use dynomite::embed::Datastore;
use dynomite::hashkit::HashType;
use prost::Message;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::net::{TcpListener, TcpStream};
type StoredMap = std::collections::HashMap<(u32, Vec<u8>), Vec<u8>>;
struct ControllableOutbound {
acking: HashSet<u32>,
durable: HashSet<u32>,
stored: Mutex<StoredMap>,
}
impl std::fmt::Debug for ControllableOutbound {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ControllableOutbound")
.finish_non_exhaustive()
}
}
impl PeerOutbound for ControllableOutbound {
fn dispatch(&self, _peer_idx: u32, _op: PeerOp) -> BoxFuture<'_, ()> {
Box::pin(async {})
}
fn request(&self, peer_idx: u32, op: PeerOp) -> BoxFuture<'_, Option<Vec<u8>>> {
Box::pin(async move {
if !self.acking.contains(&peer_idx) {
return None;
}
match op {
PeerOp::RepairPut { key, storage, .. } => {
self.stored
.lock()
.expect("lock")
.insert((peer_idx, key), storage);
let ack = if self.durable.contains(&peer_idx) {
ACK_STORED_DURABLE
} else {
ACK_STORED
};
Some(vec![ack])
}
PeerOp::Get { key, .. } => Some(
self.stored
.lock()
.expect("lock")
.get(&(peer_idx, key))
.cloned()
.unwrap_or_default(),
),
_ => None,
}
})
}
}
async fn send_frame(stream: &mut TcpStream, code: u8, body: &[u8]) {
let len = u32::try_from(body.len() + 1).expect("len");
stream.write_all(&len.to_be_bytes()).await.expect("len");
stream.write_all(&[code]).await.expect("code");
stream.write_all(body).await.expect("body");
}
async fn recv_frame(stream: &mut TcpStream) -> (u8, Vec<u8>) {
let mut len_buf = [0u8; 4];
stream.read_exact(&mut len_buf).await.expect("len");
let len = u32::from_be_bytes(len_buf) as usize;
let mut buf = vec![0u8; len];
stream.read_exact(&mut buf).await.expect("frame");
(buf[0], buf[1..].to_vec())
}
async fn spawn(
acking: Vec<u32>,
props: BucketProps,
) -> (std::net::SocketAddr, tokio::task::JoinHandle<()>) {
let durable = acking.clone();
spawn_with_durability(acking, durable, props).await
}
async fn spawn_with_durability(
acking: Vec<u32>,
durable: Vec<u32>,
props: BucketProps,
) -> (std::net::SocketAddr, tokio::task::JoinHandle<()>) {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.keep();
let ds: Arc<dyn Datastore> =
Arc::new(dyniak::datastore::NoxuDatastore::open_transactional(&path).expect("noxu"));
let registry = Arc::new(BucketPropsRegistry::new_riak_defaults());
registry.set(b"", b"b", props);
let span = u64::from(u32::MAX);
let pts: Vec<RingPoint> = (0..3u32)
.map(|i| RingPoint::new(u64::from(i) * span / 3, i, "dc1", "r1"))
.collect();
let router = Arc::new(BucketRouter::new(
registry,
Arc::new(RingView::new(pts)),
HashType::Murmur,
));
let hooks = RoutingHooks {
router,
outbound: Arc::new(ControllableOutbound {
acking: acking.into_iter().collect(),
durable: durable.into_iter().collect(),
stored: Mutex::new(StoredMap::new()),
}) as Arc<dyn PeerOutbound>,
local_actor: dyniak::datatypes::ActorId::new("dc1", "n0"),
local_peer_idx: 0,
precommit: None,
postcommit: None,
};
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let admin = Arc::new(dynomite::cluster::admin_rpc::NoopClusterAdmin);
let server = tokio::spawn(async move {
let _ = dyniak::server::serve_pbc_with_routing(listener, ds, admin, hooks).await;
});
(addr, server)
}
struct DownPeers(HashSet<u32>);
impl dyniak::replication::ReplicaLiveness for DownPeers {
fn is_up(&self, peer_idx: u32) -> bool {
!self.0.contains(&peer_idx)
}
}
impl std::fmt::Debug for DownPeers {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("DownPeers").finish_non_exhaustive()
}
}
async fn spawn_with_liveness(
down: Vec<u32>,
acking: Vec<u32>,
durable: Vec<u32>,
props: BucketProps,
) -> (std::net::SocketAddr, tokio::task::JoinHandle<()>) {
let dir = tempfile::tempdir().expect("tempdir");
let path = dir.keep();
let ds: Arc<dyn Datastore> =
Arc::new(dyniak::datastore::NoxuDatastore::open_transactional(&path).expect("noxu"));
let registry = Arc::new(BucketPropsRegistry::new_riak_defaults());
registry.set(b"", b"b", props);
let span = u64::from(u32::MAX);
let pts: Vec<RingPoint> = (0..4u32)
.map(|i| RingPoint::new(u64::from(i) * span / 4, i, "dc1", "r1"))
.collect();
let router = Arc::new(
BucketRouter::new(registry, Arc::new(RingView::new(pts)), HashType::Murmur)
.with_liveness(Arc::new(DownPeers(down.into_iter().collect()))),
);
let hooks = RoutingHooks {
router,
outbound: Arc::new(ControllableOutbound {
acking: acking.into_iter().collect(),
durable: durable.into_iter().collect(),
stored: Mutex::new(StoredMap::new()),
}) as Arc<dyn PeerOutbound>,
local_actor: dyniak::datatypes::ActorId::new("dc1", "n0"),
local_peer_idx: 0,
precommit: None,
postcommit: None,
};
let listener = TcpListener::bind("127.0.0.1:0").await.expect("bind");
let addr = listener.local_addr().expect("addr");
let admin = Arc::new(dynomite::cluster::admin_rpc::NoopClusterAdmin);
let server = tokio::spawn(async move {
let _ = dyniak::server::serve_pbc_with_routing(listener, ds, admin, hooks).await;
});
(addr, server)
}
async fn put(c: &mut TcpStream, w: Option<u32>) -> u8 {
put_with(c, w, None, None).await
}
async fn put_with(c: &mut TcpStream, w: Option<u32>, pw: Option<u32>, dw: Option<u32>) -> u8 {
put_key_with(c, b"k", w, pw, dw).await
}
async fn put_key_with(
c: &mut TcpStream,
key: &[u8],
w: Option<u32>,
pw: Option<u32>,
dw: Option<u32>,
) -> u8 {
let req = RpbPutReq {
bucket: b"b".to_vec(),
key: Some(key.to_vec()),
w,
pw,
dw,
content: Some(RpbContent {
value: b"v".to_vec(),
..RpbContent::default()
}),
..RpbPutReq::default()
};
send_frame(c, MessageCode::PutReq.as_u8(), &req.encode_to_vec()).await;
recv_frame(c).await.0
}
#[tokio::test]
async fn write_quorum_all_fails_when_a_replica_does_not_ack() {
let (addr, server) = spawn(
vec![1],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let code = put(&mut c, Some(QUORUM_ALL)).await;
assert_eq!(
code,
MessageCode::ErrorResp.as_u8(),
"W=all must fail when a replica does not ack"
);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn write_quorum_is_met_when_enough_replicas_ack() {
let (addr, server) = spawn(
vec![1, 2],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let code = put(&mut c, Some(QUORUM_ALL)).await;
assert_eq!(
code,
MessageCode::PutResp.as_u8(),
"W=all succeeds when every replica acks"
);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn write_quorum_one_succeeds_on_local_alone() {
let (addr, server) = spawn(
vec![],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let code = put(&mut c, Some(QUORUM_ONE)).await;
assert_eq!(
code,
MessageCode::PutResp.as_u8(),
"W=one is satisfied by the local write"
);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn read_quorum_all_fails_below_quorum() {
let (addr, server) = spawn(
vec![1],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let get = RpbGetReq {
bucket: b"b".to_vec(),
key: b"k".to_vec(),
r: Some(QUORUM_ALL),
..RpbGetReq::default()
};
send_frame(&mut c, MessageCode::GetReq.as_u8(), &get.encode_to_vec()).await;
let (code, _) = recv_frame(&mut c).await;
assert_eq!(
code,
MessageCode::ErrorResp.as_u8(),
"R=all must fail when fewer than N replicas respond"
);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn durable_write_quorum_fails_when_acks_are_not_durable() {
let (addr, server) = spawn_with_durability(
vec![1, 2],
vec![1],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let code = put_with(&mut c, Some(QUORUM_ALL), None, Some(QUORUM_ALL)).await;
assert_eq!(
code,
MessageCode::ErrorResp.as_u8(),
"DW=all must fail when a replica's ack is not confirmed durable"
);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn durable_write_quorum_succeeds_when_every_ack_is_durable() {
let (addr, server) = spawn_with_durability(
vec![1, 2],
vec![1, 2],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let code = put_with(&mut c, Some(QUORUM_ALL), None, Some(QUORUM_ALL)).await;
assert_eq!(
code,
MessageCode::PutResp.as_u8(),
"DW=all succeeds when every ack is confirmed durable"
);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn primary_write_quorum_fails_when_only_a_fallback_acks() {
let (addr, server) = spawn_with_liveness(
vec![1],
vec![3],
vec![3],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let code = put_key_with(&mut c, b"key8", Some(QUORUM_ONE), Some(2), None).await;
assert_eq!(
code,
MessageCode::ErrorResp.as_u8(),
"PW=2 must fail when only a fallback (not a second primary) acks"
);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn primary_write_quorum_succeeds_when_enough_primaries_ack() {
let (addr, server) = spawn_with_liveness(
vec![1],
vec![2, 3],
vec![2, 3],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let code = put_key_with(&mut c, b"key8", Some(QUORUM_ONE), Some(2), None).await;
assert_eq!(
code,
MessageCode::PutResp.as_u8(),
"PW=2 succeeds once two primaries ack, even with a fallback in play"
);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn primary_read_quorum_fails_when_a_primary_is_silent_even_if_a_fallback_answers() {
let (addr, server) = spawn_with_liveness(
vec![1],
vec![3],
vec![3],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let code = put_key_with(&mut c, b"key8", Some(QUORUM_ONE), None, None).await;
assert_eq!(code, MessageCode::PutResp.as_u8(), "seed put must succeed");
let get = RpbGetReq {
bucket: b"b".to_vec(),
key: b"key8".to_vec(),
pr: Some(2),
..RpbGetReq::default()
};
send_frame(&mut c, MessageCode::GetReq.as_u8(), &get.encode_to_vec()).await;
let (code, _) = recv_frame(&mut c).await;
assert_eq!(
code,
MessageCode::ErrorResp.as_u8(),
"PR=2 must fail when only the fallback (peer 3), not a second \
primary, answers the read"
);
server.abort();
let _ = server.await;
}
#[tokio::test]
async fn primary_read_quorum_succeeds_when_enough_primaries_answer() {
let (addr, server) = spawn_with_liveness(
vec![1],
vec![2, 3],
vec![2, 3],
BucketProps {
n_val: Some(3),
..BucketProps::default()
},
)
.await;
let mut c = TcpStream::connect(addr).await.expect("connect");
let code = put_key_with(&mut c, b"key8", Some(QUORUM_ALL), None, None).await;
assert_eq!(code, MessageCode::PutResp.as_u8(), "seed put must succeed");
let get = RpbGetReq {
bucket: b"b".to_vec(),
key: b"key8".to_vec(),
pr: Some(2),
..RpbGetReq::default()
};
send_frame(&mut c, MessageCode::GetReq.as_u8(), &get.encode_to_vec()).await;
let (code, _) = recv_frame(&mut c).await;
assert_eq!(
code,
MessageCode::GetResp.as_u8(),
"PR=2 succeeds once two primaries answer, even with a fallback in play"
);
server.abort();
let _ = server.await;
}