#![cfg(feature = "noxu")]
use std::sync::Arc;
use std::time::Duration;
use dyniak::datastore::xa_net::{
serve_xa_peer, DnodeXaTransport, XaPeer, XaTransport, XaTransportError,
};
use dyniak::datastore::xa_wire::{WireXid, XaWriteOp};
use dyniak::datastore::XaParticipant;
use dynomite::io::mbuf::MbufPool;
use dynomite::proto::dnode::{dmsg_write, DmsgType};
use tempfile::TempDir;
use tokio::io::{AsyncReadExt as _, AsyncWriteExt as _};
use tokio::net::{TcpListener, TcpStream};
const FORMAT_ID: i32 = 0x6479_6e6b;
fn scratch_dir() -> TempDir {
let base = std::path::Path::new("/scratch");
if base.is_dir() {
TempDir::new_in(base).expect("tempdir in /scratch")
} else {
TempDir::new().expect("tempdir")
}
}
fn open_participant(dir: &TempDir, name: &[u8]) -> XaParticipant {
XaParticipant::open(dir.path(), name.to_vec()).expect("open participant")
}
fn wire(bqual: &[u8]) -> WireXid {
WireXid {
format_id: FORMAT_ID,
gtrid: 1u64.to_be_bytes().to_vec(),
bqual: bqual.to_vec(),
}
}
fn frame(ty: DmsgType, payload: &[u8]) -> Vec<u8> {
let pool = MbufPool::default();
let mut header = pool.get();
let plen = u32::try_from(payload.len()).unwrap();
dmsg_write(&mut header, 1, ty, 0, true, None, plen).expect("header");
let mut out = header.readable().to_vec();
out.extend_from_slice(payload);
out
}
#[tokio::test]
async fn connect_failure_is_transport_error() {
let addr = "127.0.0.1:1".parse().unwrap();
let transport = DnodeXaTransport::new(addr).with_timeout(Duration::from_millis(200));
let res = transport.commit(&wire(b"west"), b"west").await;
assert!(matches!(res, Err(XaTransportError::Transport(_))));
}
#[tokio::test]
async fn silent_peer_times_out_the_phase() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let _server = tokio::spawn(async move {
let (mut s, _) = listener.accept().await.unwrap();
let mut buf = [0u8; 256];
let _ = s.read(&mut buf).await;
tokio::time::sleep(Duration::from_secs(5)).await;
drop(s);
});
let transport = DnodeXaTransport::new(addr).with_timeout(Duration::from_millis(150));
let res = transport.commit(&wire(b"west"), b"west").await;
assert!(matches!(res, Err(XaTransportError::Timeout)), "{res:?}");
}
#[tokio::test]
async fn prepare_wrong_reply_type_is_transport_error() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let _server = tokio::spawn(async move {
let (mut s, _) = listener.accept().await.unwrap();
let mut buf = [0u8; 512];
let _ = s.read(&mut buf).await;
let reply = frame(DmsgType::XaAck, &[1u8]);
s.write_all(&reply).await.unwrap();
tokio::time::sleep(Duration::from_millis(200)).await;
});
let transport = DnodeXaTransport::new(addr).with_timeout(Duration::from_secs(2));
let writes = vec![XaWriteOp::Put {
bucket: b"u".to_vec(),
key: b"alice".to_vec(),
value: b"a".to_vec(),
indexes: vec![],
}];
let res = transport.prepare(&wire(b"west"), b"west", &writes).await;
match res {
Err(XaTransportError::Transport(msg)) => {
assert!(msg.contains("expected XaVote"), "{msg}");
}
other => panic!("expected wrong-reply transport error, got {other:?}"),
}
}
#[tokio::test]
async fn commit_wrong_reply_type_is_transport_error() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let _server = tokio::spawn(async move {
let (mut s, _) = listener.accept().await.unwrap();
let mut buf = [0u8; 512];
let _ = s.read(&mut buf).await;
let reply = frame(DmsgType::XaVote, &[0u8]);
s.write_all(&reply).await.unwrap();
tokio::time::sleep(Duration::from_millis(200)).await;
});
let transport = DnodeXaTransport::new(addr).with_timeout(Duration::from_secs(2));
let res = transport.commit(&wire(b"west"), b"west").await;
match res {
Err(XaTransportError::Transport(msg)) => {
assert!(msg.contains("expected XaAck"), "{msg}");
}
other => panic!("expected wrong-reply transport error, got {other:?}"),
}
}
#[tokio::test]
async fn commit_ack_false_is_transport_error() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let _server = tokio::spawn(async move {
let (mut s, _) = listener.accept().await.unwrap();
let mut buf = [0u8; 512];
let _ = s.read(&mut buf).await;
let reply = frame(DmsgType::XaAck, &[0u8]);
s.write_all(&reply).await.unwrap();
tokio::time::sleep(Duration::from_millis(200)).await;
});
let transport = DnodeXaTransport::new(addr).with_timeout(Duration::from_secs(2));
let res = transport.commit(&wire(b"west"), b"west").await;
match res {
Err(XaTransportError::Transport(msg)) => {
assert!(msg.contains("unresolved"), "{msg}");
}
other => panic!("expected unresolved transport error, got {other:?}"),
}
}
#[tokio::test]
async fn peer_closing_connection_is_transport_error() {
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let _server = tokio::spawn(async move {
let (mut s, _) = listener.accept().await.unwrap();
let mut buf = [0u8; 512];
let _ = s.read(&mut buf).await;
s.shutdown().await.unwrap();
drop(s);
});
let transport = DnodeXaTransport::new(addr).with_timeout(Duration::from_secs(2));
let res = transport.commit(&wire(b"west"), b"west").await;
assert!(
matches!(res, Err(XaTransportError::Transport(_))),
"{res:?}"
);
}
#[tokio::test]
async fn peer_plane_rejects_unexpected_dnode_type() {
let d = scratch_dir();
let peer = Arc::new(XaPeer::new(vec![(
b"west".to_vec(),
open_participant(&d, b"west"),
)]));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let _server = tokio::spawn(async move {
let _ = serve_xa_peer(listener, peer).await;
});
let mut conn = TcpStream::connect(addr).await.unwrap();
let bad = frame(DmsgType::GossipSyn, b"hello");
conn.write_all(&bad).await.unwrap();
let mut buf = [0u8; 64];
let n = conn.read(&mut buf).await.unwrap_or(0);
assert_eq!(n, 0, "peer closed the connection after an unexpected type");
}
#[tokio::test]
async fn persistent_connection_is_reused_across_phases() {
let d = scratch_dir();
let peer = Arc::new(XaPeer::new(vec![(
b"west".to_vec(),
open_participant(&d, b"west"),
)]));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let _server = tokio::spawn(async move {
let _ = serve_xa_peer(listener, peer).await;
});
let transport = DnodeXaTransport::new(addr).with_timeout(Duration::from_secs(2));
let writes = vec![XaWriteOp::Put {
bucket: b"u".to_vec(),
key: b"alice".to_vec(),
value: b"a".to_vec(),
indexes: vec![],
}];
let vote = transport
.prepare(&wire(b"west"), b"west", &writes)
.await
.expect("prepare");
assert_eq!(vote, dyniak::datastore::xa_wire::XaVote::Ok);
transport
.commit(&wire(b"west"), b"west")
.await
.expect("commit");
}
#[tokio::test]
async fn split_frame_is_reassembled_across_reads() {
let d = scratch_dir();
let peer = Arc::new(XaPeer::new(vec![(
b"west".to_vec(),
open_participant(&d, b"west"),
)]));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let _server = tokio::spawn(async move {
let _ = serve_xa_peer(listener, peer).await;
});
let payload = dyniak::datastore::xa_wire::XaResolveMsg {
xid: wire(b"west"),
env: b"west".to_vec(),
}
.encode();
let full = frame(DmsgType::XaCommit, &payload);
let split_at = full.len() - payload.len().min(full.len() / 2).max(1);
let mut conn = TcpStream::connect(addr).await.unwrap();
conn.write_all(&full[..split_at]).await.unwrap();
conn.flush().await.unwrap();
tokio::time::sleep(Duration::from_millis(50)).await;
conn.write_all(&full[split_at..]).await.unwrap();
let mut buf = [0u8; 64];
let n = conn.read(&mut buf).await.unwrap();
assert!(n > 0, "server replied after reassembling the split frame");
}
#[tokio::test]
async fn garbage_header_is_a_parse_error() {
let d = scratch_dir();
let peer = Arc::new(XaPeer::new(vec![(
b"west".to_vec(),
open_participant(&d, b"west"),
)]));
let listener = TcpListener::bind("127.0.0.1:0").await.unwrap();
let addr = listener.local_addr().unwrap();
let _server = tokio::spawn(async move {
let _ = serve_xa_peer(listener, peer).await;
});
let mut conn = TcpStream::connect(addr).await.unwrap();
conn.write_all(b"GARBAGE!").await.unwrap();
let mut buf = [0u8; 64];
let n = conn.read(&mut buf).await.unwrap_or(0);
assert_eq!(n, 0, "peer closed the connection after a parse error");
}