#![cfg(feature = "noxu")]
use std::path::Path;
use std::sync::Arc;
use tempfile::TempDir;
use dynomite::embed::Datastore;
use dynomite::hashkit::HashType;
use dynomite::msg::ConsistencyLevel;
use dynomite::net::ReplicaApplySink;
use dyniak::bucket_props::{BucketProps, BucketPropsRegistry};
use dyniak::datastore::NoxuDatastore;
use dyniak::datatypes::keyfun::KeyFun;
use dyniak::proto::http::object::HttpObject;
use dyniak::proto::replica_wire::{decode_peer_op, encode_peer_op};
use dyniak::replication::{
plan_replicas, ReplicationPlan, ReplicationStrategy, RingPoint, RingView,
};
use dyniak::router::{BucketRouter, PeerOp};
use dyniak::ReplicaApplier;
use dyniak::router::{PeerOutbound, RoutingHooks};
use dynomite::net::client::BoxFuture;
use std::sync::Mutex as StdMutex;
#[derive(Debug, Default)]
struct CapturingOutbound {
ops: StdMutex<Vec<(u32, PeerOp)>>,
}
impl PeerOutbound for CapturingOutbound {
fn dispatch(&self, peer_idx: u32, op: PeerOp) -> BoxFuture<'_, ()> {
self.ops.lock().expect("lock").push((peer_idx, op));
Box::pin(async {})
}
}
fn scratch_dir() -> TempDir {
let base = Path::new("/scratch");
if base.is_dir() {
TempDir::new_in(base).expect("tempdir in /scratch")
} else {
TempDir::new().expect("tempdir")
}
}
fn five_peer_ring() -> Arc<RingView> {
let span = u64::from(u32::MAX);
let pts: Vec<RingPoint> = (0..5u32)
.map(|i| RingPoint::new(u64::from(i) * span / 5, i, "dc1", "r1"))
.collect();
Arc::new(RingView::new(pts))
}
#[test]
fn router_fan_out_matches_plan_replicas() {
let ring = five_peer_ring();
let registry = Arc::new(BucketPropsRegistry::new_riak_defaults());
registry.set(
b"default",
b"users",
BucketProps {
keyfun: Some(KeyFun::Std),
strategy: Some(ReplicationStrategy::Successors),
n_val: Some(3),
..BucketProps::default()
},
);
let router = BucketRouter::new(registry, Arc::clone(&ring), HashType::Murmur);
for i in 0..64u32 {
let key = format!("key-{i}");
let decision = router.route(b"default", b"users", key.as_bytes());
let expected = plan_replicas(
ring.as_ref(),
decision.key_hash,
3,
ReplicationStrategy::Successors,
ConsistencyLevel::DcOne,
);
let expected_list = match expected {
ReplicationPlan::Successors { .. } => expected.into_replica_list(),
other @ ReplicationPlan::Topology(_) => {
panic!("expected successors plan, got {other:?}")
}
};
let actual: Vec<u32> = decision.replica_list().iter().map(|t| t.peer_idx).collect();
let want: Vec<u32> = expected_list.iter().map(|t| t.peer_idx).collect();
assert_eq!(
actual, want,
"router replica set for {key} diverges from plan_replicas"
);
assert_eq!(actual.len(), 3, "n_val=3 yields 3 replicas for {key}");
let mut sorted = actual.clone();
sorted.sort_unstable();
sorted.dedup();
assert_eq!(sorted.len(), 3, "replicas are distinct peers for {key}");
}
}
#[tokio::test]
async fn receive_apply_put_lands_in_local_noxu() {
let dir = scratch_dir();
let noxu = NoxuDatastore::open_transactional(dir.path()).expect("open noxu");
let ds: Arc<dyn Datastore> = Arc::new(noxu);
let applier = ReplicaApplier::new(Arc::clone(&ds));
let op = PeerOp::Put {
bucket_type: b"default".to_vec(),
bucket: b"users".to_vec(),
key: b"alice".to_vec(),
value: b"alice-replicated".to_vec(),
};
let wire = encode_peer_op(&op);
assert_eq!(decode_peer_op(&wire).expect("decode"), op);
applier.apply(&wire).await;
let stored = ds
.riak_get(
&dyniak::router::composite_storage_bucket(b"default", b"users"),
b"alice",
)
.await
.expect("riak_get")
.expect("object present after replica apply");
let obj = HttpObject::from_storage_bytes(&stored).expect("decode envelope");
assert_eq!(obj.value, b"alice-replicated");
}
#[tokio::test]
async fn receive_apply_del_removes_from_local_noxu() {
let dir = scratch_dir();
let noxu = NoxuDatastore::open_transactional(dir.path()).expect("open noxu");
let ds: Arc<dyn Datastore> = Arc::new(noxu);
let applier = ReplicaApplier::new(Arc::clone(&ds));
applier
.apply(&encode_peer_op(&PeerOp::Put {
bucket_type: b"default".to_vec(),
bucket: b"users".to_vec(),
key: b"bob".to_vec(),
value: b"bob-value".to_vec(),
}))
.await;
assert!(ds
.riak_get(
&dyniak::router::composite_storage_bucket(b"default", b"users"),
b"bob"
)
.await
.expect("get")
.is_some());
applier
.apply(&encode_peer_op(&PeerOp::Del {
bucket_type: b"default".to_vec(),
bucket: b"users".to_vec(),
key: b"bob".to_vec(),
}))
.await;
assert!(
ds.riak_get(
&dyniak::router::composite_storage_bucket(b"default", b"users"),
b"bob"
)
.await
.expect("get")
.is_none(),
"replicated delete removed the object locally"
);
}
#[tokio::test]
async fn outbound_receive_pairing_delivers_write_to_node_b() {
let dir_b = scratch_dir();
let noxu_b = NoxuDatastore::open_transactional(dir_b.path()).expect("open noxu b");
let ds_b: Arc<dyn Datastore> = Arc::new(noxu_b);
let applier_b = ReplicaApplier::new(Arc::clone(&ds_b));
let ring = five_peer_ring();
let registry = Arc::new(BucketPropsRegistry::new_riak_defaults());
registry.set(
b"default",
b"carts",
BucketProps {
keyfun: Some(KeyFun::Std),
strategy: Some(ReplicationStrategy::Successors),
n_val: Some(3),
..BucketProps::default()
},
);
let router = BucketRouter::new(registry, ring, HashType::Murmur);
let decision = router.route(b"default", b"carts", b"cart-9");
assert!(
!decision.replica_list().is_empty(),
"successors plan yields at least the primary"
);
let op = PeerOp::Put {
bucket_type: decision.bucket_type.clone(),
bucket: b"carts".to_vec(),
key: b"cart-9".to_vec(),
value: b"two-items".to_vec(),
};
let wire = encode_peer_op(&op);
applier_b.apply(&wire).await;
let stored = ds_b
.riak_get(
&dyniak::router::composite_storage_bucket(b"default", b"carts"),
b"cart-9",
)
.await
.expect("riak_get")
.expect("replica landed on node B");
let obj = HttpObject::from_storage_bytes(&stored).expect("decode envelope");
assert_eq!(obj.value, b"two-items");
}
#[tokio::test]
async fn http_put_fans_out_to_replicas() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};
let dir = scratch_dir();
let noxu = NoxuDatastore::open_transactional(dir.path()).expect("open noxu");
let ds: Arc<dyn Datastore> = Arc::new(noxu);
let ring = five_peer_ring();
let registry = Arc::new(BucketPropsRegistry::new_riak_defaults());
registry.set(
b"default",
b"users",
BucketProps {
keyfun: Some(KeyFun::Std),
strategy: Some(ReplicationStrategy::Successors),
n_val: Some(3),
..BucketProps::default()
},
);
let router = Arc::new(BucketRouter::new(registry, ring, HashType::Murmur));
let outbound = Arc::new(CapturingOutbound::default());
let hooks = RoutingHooks {
router,
outbound: outbound.clone(),
local_actor: dyniak::datatypes::ActorId::new("dc1", "local"),
local_peer_idx: 0,
precommit: None,
postcommit: None,
};
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind");
let addr = listener.local_addr().expect("addr");
let serve = tokio::spawn(dyniak::serve_http_with_routing(listener, ds, hooks));
let body = br#"{"value":[104,105],"content_type":"text/plain","indexes":[],"links":[]}"#;
let req = format!(
"PUT /buckets/users/keys/alice HTTP/1.1\r\nHost: t\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
body.len()
);
let mut sock = tokio::net::TcpStream::connect(addr).await.expect("connect");
sock.write_all(req.as_bytes()).await.expect("write head");
sock.write_all(body).await.expect("write body");
sock.flush().await.expect("flush");
let mut resp = Vec::new();
let _ = sock.read_to_end(&mut resp).await;
let resp = String::from_utf8_lossy(&resp);
assert!(
resp.contains("204"),
"PUT should be 204 No Content, got: {resp}"
);
serve.abort();
let ops = outbound.ops.lock().expect("lock");
assert!(
(2..=3).contains(&ops.len()),
"n_val=3 fans to the replica set minus self, got {}",
ops.len()
);
for (peer, op) in ops.iter() {
assert_ne!(*peer, 0, "the coordinator does not fan to itself");
match op {
PeerOp::RepairPut {
bucket,
key,
storage,
..
} => {
assert_eq!(bucket, b"users");
assert_eq!(key, b"alice");
let set = dyniak::proto::http::object::SiblingSet::from_storage_bytes(storage)
.expect("decode fanned storage");
assert!(
set.siblings.iter().any(|o| o.value == vec![104u8, 105u8]),
"fanned storage carries the written value"
);
}
other => panic!("expected PeerOp::RepairPut, got {other:?}"),
}
}
let mut peers: Vec<u32> = ops.iter().map(|(p, _)| *p).collect();
peers.sort_unstable();
peers.dedup();
assert!(
(2..=3).contains(&peers.len()),
"replicas are distinct peers"
);
}