use std::io::{Read, Write};
use std::net::{TcpListener, TcpStream};
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::mpsc;
use std::thread;
use std::time::{Duration, Instant};
use crate::ErrorCode;
use crate::ingress::sender::fail_next_catch_up_allocation_for_test;
use crate::ingress::sender::has_any_sfa_file as slot_has_sfa_file;
use crate::ingress::sender::qwp_ws::fail_next_recovered_dict_copy_for_test;
use crate::ingress::{
Buffer, ColumnName, Protocol, QwpWsEncodeScratch, QwpWsErrorCategory, QwpWsErrorPolicy,
QwpWsProgress, SenderBuilder, SymbolGlobalDict, TableName, TimestampNanos,
};
#[cfg(feature = "sync-sender-http")]
use crate::ingress::ProtocolVersion;
use crate::ws::frame::OPCODE_CLOSE;
pub(crate) const WS_GUID: &str = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11";
const FIRST_WIRE_SEQUENCE: u64 = 0;
const QWP_STATUS_OK: u8 = 0x00;
const QWP_STATUS_DURABLE_ACK: u8 = 0x02;
const QWP_STATUS_SCHEMA_MISMATCH: u8 = 0x03;
const QWP_STATUS_PARSE_ERROR: u8 = 0x05;
const QWP_STATUS_WRITE_ERROR: u8 = 0x09;
const QWP_STATUS_NOT_WRITABLE: u8 = 0x0C;
const QWP_WS_PUBLIC_BENCH_DEFAULT_ROWS: usize = 20_000_000;
const QWP_WS_PUBLIC_BENCH_DEFAULT_BATCH_SIZE: usize = 1000;
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
enum ProgressCase {
Background,
Manual,
}
impl ProgressCase {
fn name(self) -> &'static str {
match self {
Self::Background => "background",
Self::Manual => "manual",
}
}
}
fn build_qwp_ws_sender_from_builder(
progress: ProgressCase,
builder: SenderBuilder,
) -> crate::ingress::Sender {
match progress {
ProgressCase::Background => builder.build().unwrap(),
ProgressCase::Manual => builder
.qwp_ws_progress(QwpWsProgress::Manual)
.unwrap()
.build()
.unwrap(),
}
}
fn build_qwp_ws_sender(progress: ProgressCase, port: u16) -> crate::ingress::Sender {
build_qwp_ws_sender_from_builder(
progress,
SenderBuilder::new(Protocol::Ws, "127.0.0.1", port),
)
}
struct MockResult {
request_lines: Vec<String>,
received_frames: Vec<Vec<u8>>,
}
pub(crate) fn read_request_until_blank<R: Read>(stream: &mut R) -> std::io::Result<Vec<u8>> {
let mut buf = Vec::new();
let mut tmp = [0u8; 256];
loop {
let n = stream.read(&mut tmp)?;
if n == 0 {
break;
}
buf.extend_from_slice(&tmp[..n]);
if buf.windows(4).any(|w| w == b"\r\n\r\n") {
break;
}
}
Ok(buf)
}
pub(crate) fn parse_header(req: &str, name: &str) -> Option<String> {
for line in req.split("\r\n").skip(1) {
if let Some((k, v)) = line.split_once(':')
&& k.trim().eq_ignore_ascii_case(name)
{
return Some(v.trim().to_string());
}
}
None
}
pub(crate) fn read_frame(stream: &mut TcpStream) -> std::io::Result<(bool, u8, Vec<u8>)> {
let mut hdr = [0u8; 2];
stream.read_exact(&mut hdr)?;
let fin = (hdr[0] & 0x80) != 0;
let opcode = hdr[0] & 0x0F;
let masked = (hdr[1] & 0x80) != 0;
let len_short = hdr[1] & 0x7F;
let payload_len = match len_short {
126 => {
let mut b = [0u8; 2];
stream.read_exact(&mut b)?;
u16::from_be_bytes(b) as usize
}
127 => {
let mut b = [0u8; 8];
stream.read_exact(&mut b)?;
u64::from_be_bytes(b) as usize
}
n => n as usize,
};
let mut mask = [0u8; 4];
if masked {
stream.read_exact(&mut mask)?;
}
let mut payload = vec![0u8; payload_len];
stream.read_exact(&mut payload)?;
if masked {
for (i, b) in payload.iter_mut().enumerate() {
*b ^= mask[i & 3];
}
}
Ok((fin, opcode, payload))
}
pub(crate) fn write_server_binary_frame(
stream: &mut TcpStream,
payload: &[u8],
) -> std::io::Result<()> {
let mut frame = vec![0x82];
let plen = payload.len();
if plen <= 125 {
frame.push(plen as u8);
} else if plen <= 0xFFFF {
frame.push(126);
frame.extend_from_slice(&(plen as u16).to_be_bytes());
} else {
frame.push(127);
frame.extend_from_slice(&(plen as u64).to_be_bytes());
}
frame.extend_from_slice(payload);
stream.write_all(&frame)
}
pub(crate) fn perform_server_upgrade(stream: &mut TcpStream) -> std::io::Result<Vec<String>> {
perform_server_upgrade_with_version(stream, 1)
}
pub(crate) fn perform_server_upgrade_durable(
stream: &mut TcpStream,
) -> std::io::Result<Vec<String>> {
stream.set_read_timeout(Some(Duration::from_secs(5)))?;
stream.set_write_timeout(Some(Duration::from_secs(5)))?;
Ok(upgrade_mock_stream_with_durable_ack(stream, true))
}
pub(crate) fn perform_server_upgrade_with_version(
stream: &mut TcpStream,
version: u8,
) -> std::io::Result<Vec<String>> {
stream.set_read_timeout(Some(Duration::from_secs(5)))?;
stream.set_write_timeout(Some(Duration::from_secs(5)))?;
let req_bytes = read_request_until_blank(stream)?;
let req = String::from_utf8_lossy(&req_bytes).to_string();
let request_lines: Vec<String> = req
.split("\r\n")
.take_while(|l| !l.is_empty())
.map(String::from)
.collect();
let key = parse_header(&req, "Sec-WebSocket-Key").expect("missing Sec-WebSocket-Key");
let accept = compute_accept(&key);
let response = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: {version}\r\n\
\r\n"
);
stream.write_all(response.as_bytes())?;
Ok(request_lines)
}
#[cfg(feature = "sync-reader-qwp-ws")]
pub(crate) fn write_server_info_frame(stream: &mut TcpStream) -> std::io::Result<()> {
let mut payload = Vec::new();
payload.push(0x18); payload.push(0x00); payload.extend_from_slice(&1u64.to_le_bytes()); payload.extend_from_slice(&0u32.to_le_bytes()); payload.extend_from_slice(&0i64.to_le_bytes()); payload.extend_from_slice(&0u16.to_le_bytes()); payload.extend_from_slice(&0u16.to_le_bytes());
let mut frame = Vec::with_capacity(12 + payload.len());
frame.extend_from_slice(b"QWP1");
frame.push(1); frame.push(0); frame.extend_from_slice(&0u16.to_le_bytes()); frame.extend_from_slice(&(payload.len() as u32).to_le_bytes());
frame.extend_from_slice(&payload);
write_server_frame(stream, 0x2, &frame, false)
}
pub(crate) fn write_server_frame(
stream: &mut TcpStream,
opcode: u8,
payload: &[u8],
masked: bool,
) -> std::io::Result<()> {
let mut frame = vec![0x80 | (opcode & 0x0F)];
let mask_bit = if masked { 0x80 } else { 0 };
let plen = payload.len();
if plen <= 125 {
frame.push(mask_bit | plen as u8);
} else if plen <= 0xFFFF {
frame.push(mask_bit | 126);
frame.extend_from_slice(&(plen as u16).to_be_bytes());
} else {
frame.push(mask_bit | 127);
frame.extend_from_slice(&(plen as u64).to_be_bytes());
}
if masked {
let mask = [0u8; 4];
frame.extend_from_slice(&mask);
for (index, byte) in payload.iter().enumerate() {
frame.push(*byte ^ mask[index & 3]);
}
} else {
frame.extend_from_slice(payload);
}
stream.write_all(&frame)
}
fn write_server_close_frame(
stream: &mut TcpStream,
code: u16,
reason: &str,
) -> std::io::Result<()> {
let mut payload = Vec::new();
payload.extend_from_slice(&code.to_be_bytes());
payload.extend_from_slice(reason.as_bytes());
let mut frame = vec![0x88];
let plen = payload.len();
if plen <= 125 {
frame.push(plen as u8);
} else {
frame.push(126);
frame.extend_from_slice(&(plen as u16).to_be_bytes());
}
frame.extend_from_slice(&payload);
stream.write_all(&frame)
}
fn write_raw_ws_frame(stream: &mut TcpStream, byte0: u8, payload: &[u8]) -> std::io::Result<()> {
let mut frame = vec![byte0];
let plen = payload.len();
if plen <= 125 {
frame.push(plen as u8);
} else if plen <= 0xFFFF {
frame.push(126);
frame.extend_from_slice(&(plen as u16).to_be_bytes());
} else {
frame.push(127);
frame.extend_from_slice(&(plen as u64).to_be_bytes());
}
frame.extend_from_slice(payload);
stream.write_all(&frame)
}
pub(crate) fn write_qwp_ok_response(stream: &mut TcpStream, wire_seq: u64) -> std::io::Result<()> {
let mut ok = Vec::new();
ok.push(QWP_STATUS_OK);
ok.extend_from_slice(&wire_seq.to_le_bytes());
append_table_seq_txns(&mut ok, &[]);
write_server_binary_frame(stream, &ok)
}
pub(crate) fn write_qwp_ok_response_with_table_entries(
stream: &mut TcpStream,
wire_seq: u64,
entries: &[(&str, i64)],
) -> std::io::Result<()> {
let mut ok = Vec::new();
ok.push(QWP_STATUS_OK);
ok.extend_from_slice(&wire_seq.to_le_bytes());
append_table_seq_txns(&mut ok, entries);
write_server_binary_frame(stream, &ok)
}
pub(crate) fn write_qwp_durable_ack_response(
stream: &mut TcpStream,
entries: &[(&str, i64)],
) -> std::io::Result<()> {
let mut ack = Vec::new();
ack.push(QWP_STATUS_DURABLE_ACK);
append_table_seq_txns(&mut ack, entries);
write_server_binary_frame(stream, &ack)
}
fn append_table_seq_txns(payload: &mut Vec<u8>, entries: &[(&str, i64)]) {
payload.extend_from_slice(&(entries.len() as u16).to_le_bytes());
for (table, seq_txn) in entries {
payload.extend_from_slice(&(table.len() as u16).to_le_bytes());
payload.extend_from_slice(table.as_bytes());
payload.extend_from_slice(&seq_txn.to_le_bytes());
}
}
pub(crate) fn write_qwp_error_response(
stream: &mut TcpStream,
status: u8,
wire_seq: u64,
msg: &[u8],
) -> std::io::Result<()> {
let mut err = Vec::new();
err.push(status);
err.extend_from_slice(&wire_seq.to_le_bytes());
err.extend_from_slice(&(msg.len() as u16).to_le_bytes());
err.extend_from_slice(msg);
write_server_binary_frame(stream, &err)
}
pub(crate) fn compute_accept(key_b64: &str) -> String {
use base64ct::{Base64, Encoding};
let combined = format!("{key_b64}{WS_GUID}");
let digest = sha1(combined.as_bytes());
Base64::encode_string(&digest)
}
fn upgrade_mock_stream(stream: &mut TcpStream) -> Vec<String> {
upgrade_mock_stream_with_durable_ack(stream, false)
}
fn upgrade_mock_stream_with_durable_ack(
stream: &mut TcpStream,
durable_ack_enabled: bool,
) -> Vec<String> {
let req_bytes = read_request_until_blank(stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let request_lines: Vec<String> = req
.split("\r\n")
.take_while(|l| !l.is_empty())
.map(String::from)
.collect();
let key = parse_header(&req, "Sec-WebSocket-Key").expect("missing Sec-WebSocket-Key");
let accept = compute_accept(&key);
let durable_ack_header = if durable_ack_enabled {
"X-QWP-Durable-Ack: enabled\r\n"
} else {
""
};
let response = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
{durable_ack_header}\
\r\n"
);
stream.write_all(response.as_bytes()).unwrap();
request_lines
}
fn upgrade_mock_stream_with_version(stream: &mut TcpStream, version: u8) -> Vec<String> {
let req_bytes = read_request_until_blank(stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let request_lines: Vec<String> = req
.split("\r\n")
.take_while(|l| !l.is_empty())
.map(String::from)
.collect();
let key = parse_header(&req, "Sec-WebSocket-Key").expect("missing Sec-WebSocket-Key");
let accept = compute_accept(&key);
let response = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: {version}\r\n\
\r\n"
);
stream.write_all(response.as_bytes()).unwrap();
request_lines
}
fn upgrade_mock_stream_without_upgrade_header(stream: &mut TcpStream) {
let req_bytes = read_request_until_blank(stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").expect("missing Sec-WebSocket-Key");
let accept = compute_accept(&key);
let response = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(response.as_bytes()).unwrap();
}
pub(crate) fn sha1(input: &[u8]) -> [u8; 20] {
let (mut h0, mut h1, mut h2, mut h3, mut h4) = (
0x67452301u32,
0xEFCDAB89,
0x98BADCFE,
0x10325476,
0xC3D2E1F0,
);
let bit_len = (input.len() as u64).wrapping_mul(8);
let mut p = Vec::with_capacity(input.len() + 64);
p.extend_from_slice(input);
p.push(0x80);
while p.len() % 64 != 56 {
p.push(0);
}
p.extend_from_slice(&bit_len.to_be_bytes());
let mut w = [0u32; 80];
for chunk in p.chunks_exact(64) {
for (i, word) in chunk.chunks_exact(4).enumerate() {
w[i] = u32::from_be_bytes([word[0], word[1], word[2], word[3]]);
}
for i in 16..80 {
w[i] = (w[i - 3] ^ w[i - 8] ^ w[i - 14] ^ w[i - 16]).rotate_left(1);
}
let (mut a, mut b, mut c, mut d, mut e) = (h0, h1, h2, h3, h4);
for (i, &wi) in w.iter().enumerate() {
let (f, k) = match i {
0..=19 => ((b & c) | ((!b) & d), 0x5A827999u32),
20..=39 => (b ^ c ^ d, 0x6ED9EBA1),
40..=59 => ((b & c) | (b & d) | (c & d), 0x8F1BBCDC),
_ => (b ^ c ^ d, 0xCA62C1D6),
};
let t = a
.rotate_left(5)
.wrapping_add(f)
.wrapping_add(e)
.wrapping_add(k)
.wrapping_add(wi);
e = d;
d = c;
c = b.rotate_left(30);
b = a;
a = t;
}
h0 = h0.wrapping_add(a);
h1 = h1.wrapping_add(b);
h2 = h2.wrapping_add(c);
h3 = h3.wrapping_add(d);
h4 = h4.wrapping_add(e);
}
let mut out = [0u8; 20];
for (i, h) in [h0, h1, h2, h3, h4].iter().enumerate() {
out[i * 4..i * 4 + 4].copy_from_slice(&h.to_be_bytes());
}
out
}
fn spawn_mock_server() -> (u16, mpsc::Receiver<MockResult>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let request_lines = perform_server_upgrade(&mut stream).unwrap();
let mut received_frames = Vec::new();
if let Ok((_fin, _opcode, payload)) = read_frame(&mut stream) {
received_frames.push(payload);
let _ = write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE);
}
let _ = tx.send(MockResult {
request_lines,
received_frames,
});
thread::sleep(Duration::from_millis(50));
});
(port, rx)
}
fn spawn_recovery_mock_server() -> (u16, mpsc::Receiver<MockResult>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let request_lines = perform_server_upgrade(&mut stream).unwrap();
let mut received_frames = Vec::new();
if read_frame(&mut stream).is_ok()
&& let Ok((_fin, _opcode, payload)) = read_frame(&mut stream)
{
received_frames.push(payload);
let _ = write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE + 1);
}
let _ = tx.send(MockResult {
request_lines,
received_frames,
});
thread::sleep(Duration::from_millis(50));
});
(port, rx)
}
fn spawn_one_response_server(response: MockQwpResponse) -> (u16, mpsc::Receiver<MockResult>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let request_lines = perform_server_upgrade(&mut stream).unwrap();
let mut received_frames = Vec::new();
if let Ok((_fin, _opcode, payload)) = read_frame(&mut stream) {
received_frames.push(payload);
let _ = write_mock_qwp_response(&mut stream, response);
}
let _ = tx.send(MockResult {
request_lines,
received_frames,
});
thread::sleep(Duration::from_millis(50));
});
(port, rx)
}
#[derive(Clone, Copy)]
enum MockQwpResponse {
Error {
status: u8,
wire_seq: u64,
message: &'static [u8],
},
}
#[derive(Clone, Copy)]
enum RecycleServerAction {
WriteError,
NonOrderlyClose,
}
fn write_mock_qwp_response(
stream: &mut TcpStream,
response: MockQwpResponse,
) -> std::io::Result<()> {
match response {
MockQwpResponse::Error {
status,
wire_seq,
message,
} => write_qwp_error_response(stream, status, wire_seq, message),
}
}
fn spawn_recycling_server(
action: RecycleServerAction,
run_for: Duration,
) -> (u16, thread::JoinHandle<usize>, Arc<AtomicUsize>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let port = listener.local_addr().unwrap().port();
let connection_count = Arc::new(AtomicUsize::new(0));
let shared_count = Arc::clone(&connection_count);
let handle = thread::spawn(move || {
let deadline = Instant::now() + run_for;
let hard_deadline = Instant::now() + Duration::from_secs(10);
let mut connections = 0usize;
while Instant::now() < deadline || (connections < 2 && Instant::now() < hard_deadline) {
match listener.accept() {
Ok((mut stream, _)) => {
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream(&mut stream);
match read_frame(&mut stream) {
Ok((_fin, 0x2, _payload)) => {
connections += 1;
match action {
RecycleServerAction::WriteError => {
let _ = write_qwp_error_response(
&mut stream,
QWP_STATUS_WRITE_ERROR,
FIRST_WIRE_SEQUENCE,
b"retry later",
);
}
RecycleServerAction::NonOrderlyClose => {
let _ =
write_server_close_frame(&mut stream, 1002, "retry later");
}
}
shared_count.store(connections, Ordering::Release);
}
Ok((_fin, _opcode, _payload)) => {}
Err(err)
if matches!(
err.kind(),
std::io::ErrorKind::UnexpectedEof
| std::io::ErrorKind::ConnectionReset
| std::io::ErrorKind::ConnectionAborted
| std::io::ErrorKind::BrokenPipe
) => {}
Err(err) => panic!("recycling server failed to read frame: {err}"),
}
}
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(5));
}
Err(err) => panic!("recycling server accept failed: {err}"),
}
}
connections
});
(port, handle, connection_count)
}
fn spawn_stalled_after_first_frame_server() -> (u16, mpsc::Receiver<Vec<u8>>, mpsc::Sender<()>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frame_tx, frame_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _opcode, payload) = read_frame(&mut stream).unwrap();
frame_tx.send(payload).unwrap();
let _ = release_rx.recv_timeout(Duration::from_secs(5));
});
(port, frame_rx, release_tx)
}
fn spawn_silent_never_acking_server() -> (u16, mpsc::Receiver<Vec<u8>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frame_tx, frame_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
while let Ok((_fin, _opcode, payload)) = read_frame(&mut stream) {
if frame_tx.send(payload).is_err() {
break;
}
}
});
(port, frame_rx)
}
fn spawn_delayed_durable_ack_server() -> (
u16,
mpsc::Receiver<Vec<u8>>,
mpsc::Sender<()>,
mpsc::Sender<()>,
) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frame_tx, frame_rx) = mpsc::channel();
let (ok_tx, ok_rx) = mpsc::channel();
let (durable_tx, durable_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
perform_server_upgrade_durable(&mut stream).unwrap();
let (_fin, _opcode, payload) = read_frame(&mut stream).unwrap();
frame_tx.send(payload).unwrap();
ok_rx.recv_timeout(Duration::from_secs(5)).unwrap();
write_qwp_ok_response_with_table_entries(
&mut stream,
FIRST_WIRE_SEQUENCE,
&[("trades", 10)],
)
.unwrap();
durable_rx.recv_timeout(Duration::from_secs(5)).unwrap();
write_qwp_durable_ack_response(&mut stream, &[("trades", 10)]).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(10)))
.unwrap();
let mut sink = [0u8; 256];
while matches!(stream.read(&mut sink), Ok(n) if n > 0) {}
});
(port, frame_rx, ok_tx, durable_tx)
}
struct DurableBacklogServer {
port: u16,
initial_ok_count: Arc<AtomicUsize>,
disconnect_tx: mpsc::Sender<usize>,
replayed_rx: mpsc::Receiver<usize>,
release_durable_tx: mpsc::Sender<()>,
resumed_rx: mpsc::Receiver<()>,
done_tx: mpsc::Sender<()>,
handle: thread::JoinHandle<()>,
}
fn qwp_frame_has_tables(payload: &[u8]) -> bool {
assert!(
payload.len() >= 8,
"short QWP frame: {} bytes",
payload.len()
);
assert_eq!(&payload[0..4], b"QWP1");
u16::from_le_bytes([payload[6], payload[7]]) != 0
}
fn spawn_durable_backlog_reconnect_server() -> DurableBacklogServer {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let initial_ok_count = Arc::new(AtomicUsize::new(0));
let server_initial_ok_count = Arc::clone(&initial_ok_count);
let (disconnect_tx, disconnect_rx) = mpsc::channel();
let (replayed_tx, replayed_rx) = mpsc::channel();
let (release_durable_tx, release_durable_rx) = mpsc::channel();
let (resumed_tx, resumed_rx) = mpsc::channel();
let (done_tx, done_rx) = mpsc::channel();
let handle = thread::spawn(move || {
let (mut initial, _) = listener.accept().unwrap();
perform_server_upgrade_durable(&mut initial).unwrap();
initial
.set_read_timeout(Some(Duration::from_millis(100)))
.unwrap();
let mut initial_frames = 0usize;
let mut wire_seq = FIRST_WIRE_SEQUENCE;
let expected_replay = loop {
match disconnect_rx.try_recv() {
Ok(expected_replay) => break expected_replay,
Err(mpsc::TryRecvError::Empty) => {}
Err(mpsc::TryRecvError::Disconnected) => {
panic!("durable backlog test dropped the disconnect signal")
}
}
match read_frame(&mut initial) {
Ok((_fin, 0x2, payload)) => {
let response_wire_seq = wire_seq;
wire_seq += 1;
if !qwp_frame_has_tables(&payload) {
continue;
}
initial_frames += 1;
write_qwp_ok_response_with_table_entries(
&mut initial,
response_wire_seq,
&[("trades", initial_frames as i64)],
)
.unwrap();
server_initial_ok_count.store(initial_frames, Ordering::Release);
}
Ok((_fin, 0x9, payload)) => {
write_server_frame(&mut initial, 0xA, &payload, false).unwrap();
}
Ok((_fin, 0x8, _payload)) => {
panic!("client closed before the durable backlog was released")
}
Ok((_fin, _opcode, _payload)) => {}
Err(err)
if matches!(
err.kind(),
std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
) => {}
Err(err) => panic!("initial durable backlog connection failed: {err}"),
}
};
assert_eq!(
initial_frames, expected_replay,
"every published frame must receive an ordinary OK before reconnect"
);
drop(initial);
let (mut replay, _) = listener.accept().unwrap();
perform_server_upgrade_durable(&mut replay).unwrap();
let mut replayed = 0usize;
let mut wire_seq = FIRST_WIRE_SEQUENCE;
while replayed < expected_replay {
match read_frame(&mut replay) {
Ok((_fin, 0x2, payload)) => {
let response_wire_seq = wire_seq;
wire_seq += 1;
if !qwp_frame_has_tables(&payload) {
continue;
}
replayed += 1;
write_qwp_ok_response_with_table_entries(
&mut replay,
response_wire_seq,
&[("trades", replayed as i64)],
)
.unwrap();
}
Ok((_fin, 0x9, payload)) => {
write_server_frame(&mut replay, 0xA, &payload, false).unwrap();
}
Ok((_fin, 0x8, _payload)) => {
panic!("client closed before replaying the durable backlog")
}
Ok((_fin, _opcode, _payload)) => {}
Err(err) => panic!("durable backlog replay failed: {err}"),
}
}
replayed_tx.send(replayed).unwrap();
release_durable_rx
.recv_timeout(Duration::from_secs(5))
.unwrap();
write_qwp_durable_ack_response(&mut replay, &[("trades", replayed as i64)]).unwrap();
loop {
match read_frame(&mut replay) {
Ok((_fin, 0x2, payload)) => {
let response_wire_seq = wire_seq;
wire_seq += 1;
if !qwp_frame_has_tables(&payload) {
continue;
}
let resumed_txn = replayed as i64 + 1;
write_qwp_ok_response_with_table_entries(
&mut replay,
response_wire_seq,
&[("trades", resumed_txn)],
)
.unwrap();
write_qwp_durable_ack_response(&mut replay, &[("trades", resumed_txn)])
.unwrap();
resumed_tx.send(()).unwrap();
break;
}
Ok((_fin, 0x9, payload)) => {
write_server_frame(&mut replay, 0xA, &payload, false).unwrap();
}
Ok((_fin, 0x8, _payload)) => {
panic!("client closed before producer capacity recovered")
}
Ok((_fin, _opcode, _payload)) => {}
Err(err) => panic!("post-durable publication failed: {err}"),
}
}
done_rx.recv_timeout(Duration::from_secs(5)).unwrap();
});
DurableBacklogServer {
port,
initial_ok_count,
disconnect_tx,
replayed_rx,
release_durable_tx,
resumed_rx,
done_tx,
handle,
}
}
fn spawn_ack_each_frame_server() -> (u16, thread::JoinHandle<usize>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let handle = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
let mut next_wire_seq = FIRST_WIRE_SEQUENCE;
let mut binary_frames = 0usize;
loop {
match read_frame(&mut stream) {
Ok((_fin, 0x2, _payload)) => {
write_qwp_ok_response(&mut stream, next_wire_seq).unwrap();
next_wire_seq += 1;
binary_frames += 1;
}
Ok((_fin, 0x8, _payload)) => break,
Ok((_fin, 0x9, payload)) => {
write_server_frame(&mut stream, 0xA, &payload, false).unwrap();
}
Ok((_fin, _opcode, _payload)) => {}
Err(err)
if matches!(
err.kind(),
std::io::ErrorKind::UnexpectedEof
| std::io::ErrorKind::ConnectionReset
| std::io::ErrorKind::ConnectionAborted
| std::io::ErrorKind::BrokenPipe
) =>
{
break;
}
Err(err) => panic!("benchmark ACK server failed to read frame: {err}"),
}
}
binary_frames
});
(port, handle)
}
fn wait_until<F: FnMut() -> bool>(timeout: Duration, mut predicate: F) -> bool {
let deadline = Instant::now() + timeout;
while Instant::now() < deadline {
if predicate() {
return true;
}
thread::sleep(Duration::from_millis(10));
}
predicate()
}
fn spawn_upgrade_only_server() -> u16 {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(10)))
.unwrap();
let mut sink = [0u8; 256];
while matches!(stream.read(&mut sink), Ok(n) if n > 0) {}
});
port
}
fn qwp_ws_replay_encoded_len(buf: &Buffer) -> usize {
let mut scratch = QwpWsEncodeScratch::new();
let mut global_dict = SymbolGlobalDict::new();
buf.as_qwp_ws()
.unwrap()
.encode_ws_replay_message(&mut scratch, &mut global_dict, 1)
.unwrap();
scratch.message.len()
}
fn qwp_ws_public_bench_env_usize(name: &str, default: usize) -> usize {
std::env::var(name)
.ok()
.and_then(|value| value.parse().ok())
.unwrap_or(default)
}
fn qwp_ws_public_bench_env_bool(name: &str) -> bool {
matches!(
std::env::var(name).as_deref(),
Ok("1" | "on" | "true" | "yes")
)
}
#[derive(Clone, Copy)]
enum QwpWsPublicBenchWorkload {
Base,
Numeric,
Symbol,
String,
Full,
}
impl QwpWsPublicBenchWorkload {
fn from_env() -> Self {
match std::env::var("QWP_WS_PUBLIC_BENCH_WORKLOAD")
.unwrap_or_else(|_| "full".to_string())
.as_str()
{
"base" => Self::Base,
"numeric" => Self::Numeric,
"symbol" => Self::Symbol,
"string" => Self::String,
"full" => Self::Full,
other => panic!("unknown QWP_WS_PUBLIC_BENCH_WORKLOAD: {other}"),
}
}
fn as_str(self) -> &'static str {
match self {
Self::Base => "base",
Self::Numeric => "numeric",
Self::Symbol => "symbol",
Self::String => "string",
Self::Full => "full",
}
}
}
fn fill_qwp_ws_public_benchmark_batch(
buffer: &mut Buffer,
workload: QwpWsPublicBenchWorkload,
prevalidated_names: bool,
batch_idx: usize,
batch_size: usize,
rows_in_batch: usize,
) {
const SYMBOLS: [&str; 8] = [
"SYM000", "SYM001", "SYM002", "SYM003", "SYM004", "SYM005", "SYM006", "SYM007",
];
const VENUES: [&str; 8] = ["ldn", "nyc", "ams", "fra", "sin", "hkg", "tyo", "sfo"];
for row_idx in 0..rows_in_batch {
let seq = (batch_idx * batch_size + row_idx) as i64;
match workload {
QwpWsPublicBenchWorkload::Base => {
bench_table(buffer, prevalidated_names).unwrap();
bench_qty(buffer, prevalidated_names, seq).unwrap();
buffer.at(TimestampNanos::new(seq)).unwrap();
}
QwpWsPublicBenchWorkload::Numeric => {
bench_table(buffer, prevalidated_names).unwrap();
bench_qty(buffer, prevalidated_names, seq).unwrap();
bench_px(buffer, prevalidated_names, 100.0 + (seq & 1023) as f64).unwrap();
buffer.at(TimestampNanos::new(seq)).unwrap();
}
QwpWsPublicBenchWorkload::Symbol => {
bench_table(buffer, prevalidated_names).unwrap();
bench_symbol(buffer, prevalidated_names, SYMBOLS[row_idx & 7]).unwrap();
bench_qty(buffer, prevalidated_names, seq).unwrap();
buffer.at(TimestampNanos::new(seq)).unwrap();
}
QwpWsPublicBenchWorkload::String => {
bench_table(buffer, prevalidated_names).unwrap();
bench_qty(buffer, prevalidated_names, seq).unwrap();
bench_venue(buffer, prevalidated_names, VENUES[row_idx & 7]).unwrap();
buffer.at(TimestampNanos::new(seq)).unwrap();
}
QwpWsPublicBenchWorkload::Full => {
bench_table(buffer, prevalidated_names).unwrap();
bench_symbol(buffer, prevalidated_names, SYMBOLS[row_idx & 7]).unwrap();
bench_qty(buffer, prevalidated_names, seq).unwrap();
bench_px(buffer, prevalidated_names, 100.0 + (seq & 1023) as f64).unwrap();
bench_venue(buffer, prevalidated_names, VENUES[row_idx & 7]).unwrap();
bench_event_ts(buffer, prevalidated_names, TimestampNanos::new(seq)).unwrap();
buffer.at(TimestampNanos::new(seq)).unwrap();
}
}
}
}
fn bench_table(buffer: &mut Buffer, prevalidated_names: bool) -> crate::Result<&mut Buffer> {
if prevalidated_names {
buffer.table(TableName::new_unchecked("trades"))
} else {
buffer.table("trades")
}
}
fn bench_symbol<'a>(
buffer: &'a mut Buffer,
prevalidated_names: bool,
value: &str,
) -> crate::Result<&'a mut Buffer> {
if prevalidated_names {
buffer.symbol(ColumnName::new_unchecked("sym"), value)
} else {
buffer.symbol("sym", value)
}
}
fn bench_qty(
buffer: &mut Buffer,
prevalidated_names: bool,
value: i64,
) -> crate::Result<&mut Buffer> {
if prevalidated_names {
buffer.column_i64(ColumnName::new_unchecked("qty"), value)
} else {
buffer.column_i64("qty", value)
}
}
fn bench_px(
buffer: &mut Buffer,
prevalidated_names: bool,
value: f64,
) -> crate::Result<&mut Buffer> {
if prevalidated_names {
buffer.column_f64(ColumnName::new_unchecked("px"), value)
} else {
buffer.column_f64("px", value)
}
}
fn bench_venue<'a>(
buffer: &'a mut Buffer,
prevalidated_names: bool,
value: &str,
) -> crate::Result<&'a mut Buffer> {
if prevalidated_names {
buffer.column_str(ColumnName::new_unchecked("venue"), value)
} else {
buffer.column_str("venue", value)
}
}
fn bench_event_ts(
buffer: &mut Buffer,
prevalidated_names: bool,
value: TimestampNanos,
) -> crate::Result<&mut Buffer> {
if prevalidated_names {
buffer.column_ts(ColumnName::new_unchecked("event_ts"), value)
} else {
buffer.column_ts("event_ts", value)
}
}
fn no_symbol_frame_at_local_hint_overcount_boundary(max: usize) -> Buffer {
for len in 0..max {
let mut buf = Buffer::qwp_ws_with_max_name_len(127);
let note = "x".repeat(len);
buf.table("trades")
.unwrap()
.column_str("note", note.as_str())
.unwrap()
.at_now()
.unwrap();
let encoded_len = qwp_ws_replay_encoded_len(&buf);
if buf.len() > max && encoded_len <= max {
return buf;
}
}
panic!("no QWP/WS size-boundary frame found for max={max}");
}
fn spawn_manual_orphan_drain_server() -> (u16, mpsc::Receiver<Vec<u8>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut orphan, _) = listener.accept().unwrap();
perform_server_upgrade(&mut orphan).unwrap();
let (_fin, _opcode, _catch_up) = read_frame(&mut orphan).unwrap();
let (_fin, _opcode, payload) = read_frame(&mut orphan).unwrap();
write_qwp_ok_response(&mut orphan, FIRST_WIRE_SEQUENCE + 1).unwrap();
tx.send(payload).unwrap();
thread::sleep(Duration::from_millis(50));
});
(port, rx)
}
fn spawn_two_frame_orphan_drain_server() -> (u16, mpsc::Receiver<Vec<Vec<u8>>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut orphan, _) = listener.accept().unwrap();
orphan
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
perform_server_upgrade(&mut orphan).unwrap();
let (_fin, _opcode, _catch_up) = read_frame(&mut orphan).unwrap();
let (_fin, _opcode, first) = read_frame(&mut orphan).unwrap();
write_qwp_ok_response(&mut orphan, FIRST_WIRE_SEQUENCE + 1).unwrap();
let (_fin, _opcode, second) = read_frame(&mut orphan).unwrap();
write_qwp_ok_response(&mut orphan, FIRST_WIRE_SEQUENCE + 2).unwrap();
tx.send(vec![first, second]).unwrap();
});
(port, rx)
}
fn spawn_manual_orphan_reject_server(status: u8) -> (u16, mpsc::Receiver<Vec<u8>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut orphan, _) = listener.accept().unwrap();
perform_server_upgrade(&mut orphan).unwrap();
let (_fin, _opcode, _catch_up) = read_frame(&mut orphan).unwrap();
let (_fin, _opcode, payload) = read_frame(&mut orphan).unwrap();
write_qwp_error_response(&mut orphan, status, FIRST_WIRE_SEQUENCE + 1, b"bad orphan")
.unwrap();
tx.send(payload).unwrap();
thread::sleep(Duration::from_millis(50));
});
(port, rx)
}
struct TerminalThenDrainOrphanServer {
port: u16,
terminal_rx: mpsc::Receiver<Vec<u8>>,
drained_rx: mpsc::Receiver<Vec<u8>>,
handle: thread::JoinHandle<()>,
}
fn spawn_terminal_then_drain_orphan_server() -> TerminalThenDrainOrphanServer {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (terminal_tx, terminal_rx) = mpsc::channel();
let (drained_tx, drained_rx) = mpsc::channel();
let server = thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut terminal_orphan, _) = listener.accept().unwrap();
perform_server_upgrade(&mut terminal_orphan).unwrap();
let (_fin, _opcode, _catch_up) = read_frame(&mut terminal_orphan).unwrap();
let (_fin, _opcode, terminal_payload) = read_frame(&mut terminal_orphan).unwrap();
write_qwp_error_response(
&mut terminal_orphan,
QWP_STATUS_PARSE_ERROR,
FIRST_WIRE_SEQUENCE + 1,
b"bad orphan",
)
.unwrap();
terminal_tx.send(terminal_payload).unwrap();
let (mut drainable_orphan, _) = listener.accept().unwrap();
perform_server_upgrade(&mut drainable_orphan).unwrap();
let (_fin, _opcode, _catch_up) = read_frame(&mut drainable_orphan).unwrap();
let (_fin, _opcode, drained_payload) = read_frame(&mut drainable_orphan).unwrap();
write_qwp_ok_response(&mut drainable_orphan, FIRST_WIRE_SEQUENCE + 1).unwrap();
drained_tx.send(drained_payload).unwrap();
thread::sleep(Duration::from_millis(50));
});
TerminalThenDrainOrphanServer {
port,
terminal_rx,
drained_rx,
handle: server,
}
}
struct RoleRejectThenDrainOrphanServer {
port: u16,
rejected_rx: mpsc::Receiver<Vec<u8>>,
drained_rx: mpsc::Receiver<Vec<u8>>,
handle: thread::JoinHandle<()>,
}
struct CatchUpFailureThenDrainOrphanServer {
port: u16,
retried_rx: mpsc::Receiver<()>,
drained_rx: mpsc::Receiver<Vec<u8>>,
handle: thread::JoinHandle<()>,
}
fn spawn_catch_up_failure_then_drain_orphan_server(
first_max_batch_size: Option<usize>,
) -> CatchUpFailureThenDrainOrphanServer {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (retried_tx, retried_rx) = mpsc::channel();
let (drained_tx, drained_rx) = mpsc::channel();
let handle = thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut failed_orphan, _) = listener.accept().unwrap();
upgrade_mock_stream_with_max_batch_size(&mut failed_orphan, first_max_batch_size);
let (mut retry_orphan, _) = listener.accept().unwrap();
retried_tx.send(()).unwrap();
drop(failed_orphan);
retry_orphan
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
perform_server_upgrade(&mut retry_orphan).unwrap();
let (_fin, _opcode, _catch_up) = read_frame(&mut retry_orphan).unwrap();
let (_fin, _opcode, drained_payload) = read_frame(&mut retry_orphan).unwrap();
write_qwp_ok_response(&mut retry_orphan, FIRST_WIRE_SEQUENCE + 1).unwrap();
drained_tx.send(drained_payload).unwrap();
});
CatchUpFailureThenDrainOrphanServer {
port,
retried_rx,
drained_rx,
handle,
}
}
fn spawn_role_reject_then_drain_orphan_server() -> RoleRejectThenDrainOrphanServer {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (rejected_tx, rejected_rx) = mpsc::channel();
let (drained_tx, drained_rx) = mpsc::channel();
let handle = thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut replica, _) = listener.accept().unwrap();
perform_server_upgrade(&mut replica).unwrap();
let (_fin, _opcode, _catch_up) = read_frame(&mut replica).unwrap();
let (_fin, _opcode, rejected_payload) = read_frame(&mut replica).unwrap();
write_qwp_error_response(
&mut replica,
QWP_STATUS_NOT_WRITABLE,
FIRST_WIRE_SEQUENCE + 1,
b"replica",
)
.unwrap();
rejected_tx.send(rejected_payload).unwrap();
let (mut primary, _) = listener.accept().unwrap();
perform_server_upgrade(&mut primary).unwrap();
let (_fin, _opcode, _catch_up) = read_frame(&mut primary).unwrap();
let (_fin, _opcode, drained_payload) = read_frame(&mut primary).unwrap();
write_qwp_ok_response(&mut primary, FIRST_WIRE_SEQUENCE + 1).unwrap();
drained_tx.send(drained_payload).unwrap();
thread::sleep(Duration::from_millis(50));
});
RoleRejectThenDrainOrphanServer {
port,
rejected_rx,
drained_rx,
handle,
}
}
struct ReplicaWindowThenPromoteServer {
port: u16,
promoted: Arc<AtomicBool>,
reject_count: Arc<AtomicUsize>,
handle: thread::JoinHandle<Vec<Vec<u8>>>,
}
fn spawn_replica_window_then_promote_server() -> ReplicaWindowThenPromoteServer {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let port = listener.local_addr().unwrap().port();
let promoted = Arc::new(AtomicBool::new(false));
let reject_count = Arc::new(AtomicUsize::new(0));
let thread_promoted = Arc::clone(&promoted);
let thread_reject_count = Arc::clone(&reject_count);
let handle = thread::spawn(move || {
let started = Instant::now();
while started.elapsed() < Duration::from_secs(15) {
let (mut stream, _) = match listener.accept() {
Ok(accepted) => accepted,
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(2));
continue;
}
Err(err) => panic!("replica-window listener failed: {err}"),
};
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
if !thread_promoted.load(Ordering::Acquire) {
let _ = read_request_until_blank(&mut stream).unwrap();
stream
.write_all(
b"HTTP/1.1 421 Misdirected Request\r\n\
Connection: close\r\n\
Content-Length: 0\r\n\
X-QuestDB-Role: REPLICA\r\n\
\r\n",
)
.unwrap();
thread_reject_count.fetch_add(1, Ordering::AcqRel);
continue;
}
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _opcode, payload) = read_frame(&mut stream).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
thread::sleep(Duration::from_millis(100));
return vec![payload];
}
Vec::new()
});
ReplicaWindowThenPromoteServer {
port,
promoted,
reject_count,
handle,
}
}
fn spawn_stalled_background_orphan_drain_server() -> (u16, mpsc::Receiver<Vec<u8>>, mpsc::Sender<()>)
{
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut orphan, _) = listener.accept().unwrap();
perform_server_upgrade(&mut orphan).unwrap();
let (_fin, _opcode, payload) = read_frame(&mut orphan).unwrap();
tx.send(payload).unwrap();
let _ = release_rx.recv_timeout(Duration::from_secs(6));
});
(port, rx, release_tx)
}
fn spawn_stalled_background_orphan_connect_server() -> (u16, mpsc::Receiver<()>, mpsc::Sender<()>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (accepted_tx, accepted_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (_orphan, _) = listener.accept().unwrap();
accepted_tx.send(()).unwrap();
let _ = release_rx.recv_timeout(Duration::from_secs(10));
});
(port, accepted_rx, release_tx)
}
fn spawn_stalled_first_background_orphan_connect_server() -> (
u16,
mpsc::Receiver<TcpListener>,
mpsc::Sender<()>,
thread::JoinHandle<()>,
) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (accepted_tx, accepted_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
let server = thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (_orphan, _) = listener.accept().unwrap();
accepted_tx.send(listener.try_clone().unwrap()).unwrap();
let _ = release_rx.recv_timeout(Duration::from_secs(10));
});
(port, accepted_rx, release_tx, server)
}
fn spawn_blocked_background_orphan_send_server() -> (
u16,
mpsc::Receiver<()>,
mpsc::Sender<()>,
thread::JoinHandle<()>,
) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (send_started_tx, send_started_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
let server = thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut orphan, _) = listener.accept().unwrap();
socket2::SockRef::from(&orphan)
.set_recv_buffer_size(4096)
.unwrap();
perform_server_upgrade(&mut orphan).unwrap();
orphan
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
let mut peek_buf = [0u8; 4096];
loop {
if orphan.peek(&mut peek_buf).unwrap() == peek_buf.len() {
break;
}
thread::yield_now();
}
send_started_tx.send(()).unwrap();
let _ = release_rx.recv_timeout(Duration::from_secs(10));
});
(port, send_started_rx, release_tx, server)
}
fn spawn_role_reject_upgrade_server(
expected_attempts: usize,
role: &'static str,
) -> (u16, thread::JoinHandle<usize>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let port = listener.local_addr().unwrap().port();
let handle = thread::spawn(move || {
let started = Instant::now();
let mut attempts = 0usize;
while attempts < expected_attempts && started.elapsed() < Duration::from_secs(5) {
match listener.accept() {
Ok((mut stream, _)) => {
attempts += 1;
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = read_request_until_blank(&mut stream).unwrap();
let response = format!(
"HTTP/1.1 421 Misdirected Request\r\n\
Connection: close\r\n\
Content-Length: 0\r\n\
X-QuestDB-Role: {role}\r\n\
\r\n"
);
stream.write_all(response.as_bytes()).unwrap();
}
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(5));
}
Err(err) => panic!("role-reject listener failed: {err}"),
}
}
attempts
});
(port, handle)
}
fn spawn_no_durable_ack_upgrade_server(done: Arc<AtomicBool>) -> (u16, thread::JoinHandle<usize>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let port = listener.local_addr().unwrap().port();
let handle = thread::spawn(move || {
let mut attempts = 0usize;
while !done.load(Ordering::Acquire) {
match listener.accept() {
Ok((mut stream, _)) => {
attempts += 1;
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = upgrade_mock_stream(&mut stream);
}
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(5));
}
Err(err) => panic!("no-durable-ack listener failed: {err}"),
}
}
attempts
});
(port, handle)
}
fn seed_orphan_slot(sf_dir: &Path) {
seed_orphan_slot_named(sf_dir, "orphan");
}
fn seed_orphan_slot_named(sf_dir: &Path, sender_id: &str) {
seed_orphan_slot_named_with_symbol(sf_dir, sender_id, "old");
}
fn seed_orphan_slot_named_with_symbol(sf_dir: &Path, sender_id: &str, symbol: &str) {
let seed_port = spawn_upgrade_only_server();
let seed_conf = format!(
"ws::addr=127.0.0.1:{seed_port};qwp_ws_progress=manual;\
sf_dir={};sender_id={sender_id};sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.display(),
);
let mut seed_sender = SenderBuilder::from_conf(&seed_conf)
.unwrap()
.build()
.unwrap();
let mut seed_buf = seed_sender.new_buffer();
seed_buf
.table("orphaned")
.unwrap()
.symbol("src", symbol)
.unwrap()
.column_i64("value", 42)
.unwrap()
.at_now()
.unwrap();
seed_sender.flush(&mut seed_buf).unwrap();
drop(seed_sender);
}
fn seed_orphan_slot_with_two_delta_frames(sf_dir: &Path) {
let seed_port = spawn_upgrade_only_server();
let seed_conf = format!(
"ws::addr=127.0.0.1:{seed_port};qwp_ws_progress=manual;\
sf_dir={};sender_id=orphan;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.display()
);
let mut seed_sender = SenderBuilder::from_conf(&seed_conf)
.unwrap()
.build()
.unwrap();
for symbol in ["alpha", "bravo"] {
let mut seed_buf = seed_sender.new_buffer();
seed_buf
.table("orphaned")
.unwrap()
.symbol("src", symbol)
.unwrap()
.column_i64("value", 42)
.unwrap()
.at_now()
.unwrap();
seed_sender.flush(&mut seed_buf).unwrap();
}
drop(seed_sender);
}
fn seed_large_orphan_slot(sf_dir: &Path) {
let seed_port = spawn_upgrade_only_server();
let seed_conf = format!(
"ws::addr=127.0.0.1:{seed_port};qwp_ws_progress=manual;\
sf_dir={};sender_id=orphan;sf_max_segment_bytes=16777216;",
sf_dir.display()
);
let mut seed_sender = SenderBuilder::from_conf(&seed_conf)
.unwrap()
.build()
.unwrap();
let value = "x".repeat(8 * 1024 * 1024);
let mut seed_buf = seed_sender.new_buffer();
seed_buf
.table("orphaned")
.unwrap()
.column_str("payload", value.as_str())
.unwrap()
.at_now()
.unwrap();
seed_sender.flush(&mut seed_buf).unwrap();
drop(seed_sender);
}
fn spawn_orphan_drain_all_frames_server() -> (u16, mpsc::Receiver<Vec<u8>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut orphan, _) = listener.accept().unwrap();
orphan
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
orphan
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
perform_server_upgrade(&mut orphan).unwrap();
let mut wire_seq = FIRST_WIRE_SEQUENCE;
while let Ok((_fin, opcode, payload)) = read_frame(&mut orphan) {
if opcode == OPCODE_CLOSE {
break;
}
if write_qwp_ok_response(&mut orphan, wire_seq).is_err() {
break;
}
wire_seq += 1;
if tx.send(payload).is_err() {
break;
}
}
});
(port, rx)
}
#[test]
fn qwp_ws_round_trip_minimal_message() {
let (port, rx) = spawn_mock_server();
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let result = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(
result
.request_lines
.first()
.unwrap()
.contains("/api/v4/write"),
"request line: {:?}",
result.request_lines.first()
);
let has_max_version = result
.request_lines
.iter()
.any(|l| l.eq_ignore_ascii_case("X-QWP-Max-Version: 1"));
assert!(has_max_version, "expected X-QWP-Max-Version header");
let frame = result.received_frames.first().expect("frame received");
assert!(frame.len() >= 12, "frame too small: {}", frame.len());
assert_eq!(&frame[0..4], b"QWP1");
assert_eq!(frame[4], 1, "version");
assert_eq!(frame[5] & 0x08, 0x08, "FLAG_DELTA_SYMBOL_DICT must be set");
let table_count = u16::from_le_bytes([frame[6], frame[7]]);
assert_eq!(table_count, 1);
let payload_len = u32::from_le_bytes([frame[8], frame[9], frame[10], frame[11]]) as usize;
assert_eq!(12 + payload_len, frame.len());
let payload = &frame[12..];
assert_eq!(payload[0], 0x00, "delta_start = 0 (varint)");
assert_eq!(payload[1], 0x01, "delta_count = 1 (varint)");
assert_eq!(payload[2], 0x07);
assert_eq!(&payload[3..10], b"ETH-USD");
}
#[test]
fn qwp_ws_max_buf_size_allows_frame_when_encoded_replay_len_fits_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let max = 1024;
let (port, rx) = spawn_mock_server();
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_buf_size(max)
.unwrap();
let mut sender = match progress {
ProgressCase::Background => builder.build().unwrap(),
ProgressCase::Manual => builder
.qwp_ws_progress(QwpWsProgress::Manual)
.unwrap()
.build()
.unwrap(),
};
let mut buf = no_symbol_frame_at_local_hint_overcount_boundary(max);
let local_hint = buf.len();
let encoded_len = qwp_ws_replay_encoded_len(&buf);
assert!(local_hint > max, "local_hint={local_hint}, max={max}");
assert!(encoded_len <= max, "encoded_len={encoded_len}, max={max}");
sender.flush(&mut buf).unwrap();
if progress == ProgressCase::Manual {
assert!(
sender.drive_once().unwrap(),
"{} mode did not report foreground progress",
progress.name()
);
}
let result = rx.recv_timeout(Duration::from_secs(5)).unwrap();
let frame = result.received_frames.first().expect("frame received");
assert_eq!(frame.len(), encoded_len, "mode={}", progress.name());
assert!(frame.len() <= max, "mode={}", progress.name());
}
}
#[test]
fn qwp_ws_max_buf_size_rejects_oversized_replay_frame_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let mut buf = Buffer::qwp_ws_with_max_name_len(127);
for idx in 0..131 {
buf.table(format!("t{idx}").as_str())
.unwrap()
.column_i64(format!("c{idx}").as_str(), idx)
.unwrap()
.at_now()
.unwrap();
}
let encoded_len = qwp_ws_replay_encoded_len(&buf);
assert!(encoded_len > 1024, "encoded_len={encoded_len}");
let max = encoded_len - 1;
let port = spawn_upgrade_only_server();
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_buf_size(max)
.unwrap();
let mut sender = match progress {
ProgressCase::Background => builder.build().unwrap(),
ProgressCase::Manual => builder
.qwp_ws_progress(QwpWsProgress::Manual)
.unwrap()
.build()
.unwrap(),
};
let err = sender.flush(&mut buf).unwrap_err();
assert_eq!(
err.code(),
crate::ErrorCode::InvalidApiCall,
"mode={}",
progress.name()
);
assert_eq!(
err.msg(),
format!(
"Could not flush buffer: QWP/WebSocket encoded message size of {encoded_len} exceeds maximum configured allowed size of {max} bytes."
),
"mode={}",
progress.name()
);
assert!(!buf.is_empty(), "mode={}", progress.name());
}
}
fn assert_durable_ack_without_opt_in(err: crate::Error, mode: ProgressCase) {
assert_eq!(
err.code(),
ErrorCode::InvalidApiCall,
"mode={}",
mode.name()
);
assert_eq!(
err.msg(),
"AckLevel::Durable requires the pool to be opened with \
`request_durable_ack=on` in the connect string.",
"mode={}",
mode.name()
);
}
#[test]
fn qwp_ws_wait_durable_without_opt_in_fails_with_no_published_frames_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let port = spawn_upgrade_only_server();
let mut sender = build_qwp_ws_sender(progress, port);
let err = sender
.wait(crate::ingress::AckLevel::Durable, Duration::from_secs(5))
.expect_err("durable wait without opt-in must fail before the empty-stream shortcut");
assert_durable_ack_without_opt_in(err, progress);
assert_eq!(sender.published_fsn().unwrap(), None);
}
}
#[test]
fn qwp_ws_wait_durable_without_opt_in_fails_after_ok_ack_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let (port, rx) = spawn_mock_server();
let mut sender = build_qwp_ws_sender(progress, port);
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
sender
.wait(crate::ingress::AckLevel::Ok, Duration::from_secs(5))
.unwrap_or_else(|e| panic!("mode={}: {e}", progress.name()));
assert_eq!(sender.acked_fsn().unwrap(), Some(fsn));
let err = sender
.wait(crate::ingress::AckLevel::Durable, Duration::from_secs(5))
.expect_err("durable wait without opt-in must fail after OK coverage too");
assert_durable_ack_without_opt_in(err, progress);
let result = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(result.received_frames.len(), 1, "mode={}", progress.name());
}
}
#[test]
fn qwp_ws_publish_ack_completes_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let (port, rx) = spawn_mock_server();
let mut sender = build_qwp_ws_sender(progress, port);
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
assert!(!sender.must_close(), "mode={}", progress.name());
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert_eq!(fsn, 0, "mode={}", progress.name());
assert!(buf.is_empty(), "mode={}", progress.name());
sender
.wait(crate::ingress::AckLevel::Ok, Duration::from_secs(5))
.unwrap_or_else(|e| panic!("mode={}: {e}", progress.name()));
assert_eq!(sender.published_fsn().unwrap(), Some(fsn));
assert_eq!(sender.acked_fsn().unwrap(), Some(fsn));
assert!(!sender.must_close(), "mode={}", progress.name());
let result = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(result.received_frames.len(), 1, "mode={}", progress.name());
assert_eq!(&result.received_frames[0][0..4], b"QWP1");
}
}
#[test]
fn qwp_ws_schema_reject_terminalizes_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let (port, rx) = spawn_one_response_server(MockQwpResponse::Error {
status: QWP_STATUS_SCHEMA_MISMATCH,
wire_seq: FIRST_WIRE_SEQUENCE,
message: b"bad schema",
});
let mut sender = build_qwp_ws_sender(progress, port);
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert_eq!(fsn, 0, "mode={}", progress.name());
let err = sender
.wait(crate::ingress::AckLevel::Ok, Duration::from_secs(5))
.unwrap_err();
assert_eq!(err.code(), ErrorCode::ServerRejection);
assert!(
err.msg().contains("bad schema"),
"mode={}, got: {}",
progress.name(),
err.msg()
);
assert_eq!(
err.qwp_ws_rejection().map(|error| error.category),
Some(QwpWsErrorCategory::SchemaMismatch),
"mode={}",
progress.name()
);
let received = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(
received.received_frames.len(),
1,
"mode={}",
progress.name()
);
let qwp_error = sender.poll_qwp_ws_error().unwrap().unwrap();
assert_eq!(qwp_error.category, QwpWsErrorCategory::SchemaMismatch);
assert_eq!(qwp_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(qwp_error.status, Some(QWP_STATUS_SCHEMA_MISMATCH));
assert_eq!(qwp_error.message.as_deref(), Some("bad schema"));
assert_eq!(qwp_error.message_sequence, Some(FIRST_WIRE_SEQUENCE));
assert_eq!(qwp_error.from_fsn, fsn);
assert_eq!(qwp_error.to_fsn, fsn);
assert_eq!(sender.poll_qwp_ws_error().unwrap(), None);
}
}
#[test]
fn qwp_ws_terminal_reject_terminalizes_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let (port, rx) = spawn_one_response_server(MockQwpResponse::Error {
status: QWP_STATUS_PARSE_ERROR,
wire_seq: FIRST_WIRE_SEQUENCE,
message: b"bad column",
});
let mut sender = build_qwp_ws_sender(progress, port);
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert_eq!(fsn, 0, "mode={}", progress.name());
let err = sender
.wait(crate::ingress::AckLevel::Ok, Duration::from_secs(5))
.unwrap_err();
assert_eq!(err.code(), ErrorCode::ServerRejection);
assert!(
err.msg().contains("bad column"),
"mode={}, got: {}",
progress.name(),
err.msg()
);
assert_eq!(
err.qwp_ws_rejection().map(|error| error.category),
Some(QwpWsErrorCategory::ParseError),
"mode={}",
progress.name()
);
let result = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(result.received_frames.len(), 1, "mode={}", progress.name());
let qwp_error = sender.poll_qwp_ws_error().unwrap().unwrap();
assert_eq!(qwp_error.category, QwpWsErrorCategory::ParseError);
assert_eq!(qwp_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(qwp_error.status, Some(QWP_STATUS_PARSE_ERROR));
assert_eq!(qwp_error.message.as_deref(), Some("bad column"));
assert_eq!(qwp_error.message_sequence, Some(FIRST_WIRE_SEQUENCE));
assert_eq!(qwp_error.from_fsn, fsn);
assert_eq!(qwp_error.to_fsn, fsn);
let mut later = sender.new_buffer();
later
.table("trades")
.unwrap()
.column_i64("qty", 2)
.unwrap()
.at_now()
.unwrap();
let later_err = sender.flush(&mut later).unwrap_err();
assert_eq!(later_err.code(), ErrorCode::ServerRejection);
assert!(
later_err.msg().contains("bad column"),
"mode={}, got: {}",
progress.name(),
later_err.msg()
);
assert!(!later.is_empty(), "mode={}", progress.name());
assert!(sender.must_close(), "mode={}", progress.name());
}
}
#[test]
fn qwp_ws_wire_has_no_frame_count_cap() {
const FRAMES: usize = 200;
let (port, frames) = spawn_silent_never_acking_server();
let conf = format!("ws::addr=127.0.0.1:{port};");
let mut sender = SenderBuilder::from_conf(conf).unwrap().build().unwrap();
for qty in 0..FRAMES as i64 {
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", qty)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
}
for i in 0..FRAMES {
let payload = frames
.recv_timeout(Duration::from_secs(5))
.unwrap_or_else(|_| {
panic!("only {i} of {FRAMES} frames reached the server without any ack")
});
assert_eq!(&payload[0..4], b"QWP1");
}
}
#[test]
fn qwp_ws_backpressure_timeout_matches_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let (port, frame_rx, release_tx) = spawn_stalled_after_first_frame_server();
let conf = format!(
"ws::addr=127.0.0.1:{port};\
qwp_ws_progress={};\
sf_max_segment_bytes=512;\
sf_max_total_bytes=1024;\
sf_append_deadline_millis=20;",
progress.name()
);
let mut sender = SenderBuilder::from_conf(conf).unwrap().build().unwrap();
let mut first = sender.new_buffer();
first
.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
let first_fsn = sender.flush_and_get_fsn(&mut first).unwrap().unwrap();
assert_eq!(first_fsn, 0, "mode={}", progress.name());
if progress == ProgressCase::Manual {
assert!(sender.drive_once().unwrap(), "mode={}", progress.name());
}
assert_eq!(
&frame_rx.recv_timeout(Duration::from_secs(5)).unwrap()[0..4],
b"QWP1"
);
let mut backpressured = None;
for _ in 0..64 {
let mut next = sender.new_buffer();
next.table("trades")
.unwrap()
.column_i64("qty", 2)
.unwrap()
.at_now()
.unwrap();
if let Err(err) = sender.flush(&mut next) {
assert!(!next.is_empty(), "mode={}", progress.name());
backpressured = Some(err);
break;
}
}
let err = backpressured.unwrap_or_else(|| {
panic!(
"the full segment ring must reject a flush, mode={}",
progress.name()
)
});
assert_eq!(err.code(), ErrorCode::SocketError);
assert!(
err.msg()
.starts_with("QWP/WebSocket Store-and-Forward append timed out"),
"mode={}, got: {}",
progress.name(),
err.msg()
);
let _ = release_tx.send(());
}
}
#[test]
fn qwp_ws_durable_ack_requires_upgrade_echo() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (request_tx, request_rx) = mpsc::channel();
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let request_lines = upgrade_mock_stream(&mut stream);
request_tx.send(request_lines).unwrap();
});
let conf = format!("ws::addr=127.0.0.1:{port};request_durable_ack=on;");
let err = SenderBuilder::from_conf(conf).unwrap().build().unwrap_err();
assert!(
err.msg().contains("server did not enable durable ACK"),
"got: {}",
err.msg()
);
let request_lines = request_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(
request_lines
.iter()
.any(|line| line.eq_ignore_ascii_case("X-QWP-Request-Durable-Ack: true"))
);
server.join().unwrap();
}
#[test]
fn qwp_ws_durable_ack_completion_waits_for_durable_confirmation_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (allow_ack_tx, allow_ack_rx) = mpsc::channel();
let (done_tx, done_rx) = mpsc::channel();
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let request_lines = upgrade_mock_stream_with_durable_ack(&mut stream, true);
assert!(
request_lines
.iter()
.any(|line| line.eq_ignore_ascii_case("X-QWP-Request-Durable-Ack: true"))
);
let (_, opcode, payload) = read_frame(&mut stream).unwrap();
assert_eq!(opcode, 0x2);
assert_eq!(&payload[0..4], b"QWP1");
write_qwp_ok_response_with_table_entries(
&mut stream,
FIRST_WIRE_SEQUENCE,
&[("trades", 10)],
)
.unwrap();
allow_ack_rx.recv_timeout(Duration::from_secs(5)).unwrap();
write_qwp_durable_ack_response(&mut stream, &[("trades", 10)]).unwrap();
let _ = done_rx.recv_timeout(Duration::from_secs(5));
});
let conf = format!(
"ws::addr=127.0.0.1:{port};\
qwp_ws_progress={};\
request_durable_ack=on;",
progress.name()
);
let mut sender = SenderBuilder::from_conf(conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
let err = sender
.wait(crate::ingress::AckLevel::Durable, Duration::from_millis(50))
.expect_err("gated durable ack must time out the bounded wait");
assert_eq!(
err.code(),
ErrorCode::FailoverRetry,
"mode={}: {}",
progress.name(),
err.msg()
);
assert!(
err.msg().contains("timed out"),
"mode={}: {}",
progress.name(),
err.msg()
);
assert_eq!(
sender.acked_fsn().unwrap(),
None,
"mode={}",
progress.name()
);
allow_ack_tx.send(()).unwrap();
sender
.wait(crate::ingress::AckLevel::Durable, Duration::from_secs(5))
.unwrap_or_else(|e| panic!("mode={}: {e}", progress.name()));
assert_eq!(sender.acked_fsn().unwrap(), Some(fsn));
done_tx.send(()).unwrap();
server.join().unwrap();
}
}
#[test]
fn qwp_ws_sender_fsn_watermarks_and_close_drain_work_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let (port, rx) = spawn_mock_server();
let mut sender = build_qwp_ws_sender(progress, port);
assert_eq!(sender.published_fsn().unwrap(), None);
assert_eq!(sender.acked_fsn().unwrap(), None);
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert_eq!(fsn, 0, "mode={}", progress.name());
assert!(buf.is_empty(), "mode={}", progress.name());
assert_eq!(sender.published_fsn().unwrap(), Some(fsn));
sender.close_drain().unwrap();
assert_eq!(sender.acked_fsn().unwrap(), Some(fsn));
let result = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(result.received_frames.len(), 1, "mode={}", progress.name());
}
}
#[test]
fn sender_sfa_fully_delivered_tracks_ok_and_durable_watermarks() {
let (port, frame_rx, ok_tx, durable_tx) = spawn_delayed_durable_ack_server();
let conf = format!("ws::addr=127.0.0.1:{port};request_durable_ack=on;");
let mut sender = SenderBuilder::from_conf(conf).unwrap().build().unwrap();
assert!(sender.sfa_fully_delivered(false));
assert!(sender.sfa_fully_delivered(true));
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert_eq!(fsn, FIRST_WIRE_SEQUENCE);
let payload = frame_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&payload[0..4], b"QWP1");
assert!(!sender.sfa_fully_delivered(false));
assert!(!sender.sfa_fully_delivered(true));
ok_tx.send(()).unwrap();
assert!(
wait_until(Duration::from_secs(5), || sender.sfa_fully_delivered(false)),
"OK watermark should cover the published frame"
);
assert!(
!sender.sfa_fully_delivered(true),
"durable watermark must wait for durable ACK coverage"
);
durable_tx.send(()).unwrap();
sender
.wait(crate::ingress::AckLevel::Durable, Duration::from_secs(5))
.unwrap();
assert!(sender.sfa_fully_delivered(false));
assert!(sender.sfa_fully_delivered(true));
#[cfg(feature = "sync-sender-http")]
{
let http_sender = SenderBuilder::new(Protocol::Http, "127.0.0.1", 1)
.protocol_version(ProtocolVersion::V1)
.unwrap()
.build()
.unwrap();
assert!(http_sender.sfa_fully_delivered(false));
assert!(http_sender.sfa_fully_delivered(true));
}
}
#[test]
fn qwp_ws_deep_durable_backlog_fills_byte_ring_replays_and_recovers() {
const MIN_DEEP_BACKLOG: usize = 64;
let DurableBacklogServer {
port,
initial_ok_count,
disconnect_tx,
replayed_rx,
release_durable_tx,
resumed_rx,
done_tx,
handle,
} = spawn_durable_backlog_reconnect_server();
let conf = format!(
"ws::addr=127.0.0.1:{port};\
request_durable_ack=on;\
sf_max_segment_bytes=512;\
sf_max_total_bytes=8192;\
sf_append_deadline_millis=100;\
durable_ack_keepalive_interval_millis=10;\
reconnect_initial_backoff_millis=1;\
reconnect_max_backoff_millis=1;\
reconnect_max_duration_millis=5000;"
);
let mut sender = SenderBuilder::from_conf(conf).unwrap().build().unwrap();
let mut published = 0usize;
let mut last_fsn = None;
let (mut blocked, backpressure) = loop {
let mut next = sender.new_buffer();
next.table("trades")
.unwrap()
.column_i64("qty", published as i64)
.unwrap()
.at_now()
.unwrap();
match sender.flush_and_get_fsn(&mut next) {
Ok(Some(fsn)) => {
published += 1;
last_fsn = Some(fsn);
}
Ok(None) => panic!("non-empty durable frame did not publish an FSN"),
Err(err) => break (next, err),
}
};
assert!(!blocked.is_empty());
assert_eq!(backpressure.code(), ErrorCode::SocketError);
assert!(
backpressure
.msg()
.contains("timed out waiting for ACK-driven segment trim"),
"expected byte-ring backpressure, got: {}",
backpressure.msg()
);
assert!(
published >= MIN_DEEP_BACKLOG,
"byte ring held only {published} frames; test did not build a deep backlog"
);
assert!(
wait_until(Duration::from_secs(5), || {
initial_ok_count.load(Ordering::Acquire) == published
}),
"server ordinary-OKed only {} of {published} frames",
initial_ok_count.load(Ordering::Acquire)
);
assert!(
wait_until(Duration::from_secs(5), || sender.sfa_fully_delivered(false)),
"ordinary OK watermark did not cover the deep backlog"
);
assert!(
!sender.sfa_fully_delivered(true),
"durable watermark advanced while durable ACKs were withheld"
);
assert_eq!(sender.acked_fsn().unwrap(), None);
disconnect_tx.send(published).unwrap();
let replayed = replayed_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(replayed, published);
assert!(
wait_until(Duration::from_secs(5), || sender.sfa_fully_delivered(false)),
"ordinary OK watermark did not recover after replay"
);
assert!(!sender.sfa_fully_delivered(true));
assert!(
wait_until(Duration::from_secs(5), || {
let totals = sender.qwp_ws_totals().unwrap();
totals.reconnects_succeeded >= 1 && totals.frames_replayed >= published as u64
}),
"reconnect/replay counters did not cover the deep backlog: {:?}",
sender.qwp_ws_totals().unwrap()
);
let totals = sender.qwp_ws_totals().unwrap();
assert_eq!(totals.reconnects_succeeded, 1);
assert_eq!(totals.frames_replayed, published as u64);
release_durable_tx.send(()).unwrap();
sender
.wait(crate::ingress::AckLevel::Durable, Duration::from_secs(5))
.unwrap();
let last_fsn = last_fsn.unwrap();
assert_eq!(sender.acked_fsn().unwrap(), Some(last_fsn));
assert!(sender.sfa_fully_delivered(true));
let resumed_fsn = sender
.flush_and_get_fsn(&mut blocked)
.unwrap()
.expect("the blocked frame must publish after durable segment trim");
assert_eq!(resumed_fsn, last_fsn + 1);
sender
.wait(crate::ingress::AckLevel::Durable, Duration::from_secs(5))
.unwrap();
resumed_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(sender.acked_fsn().unwrap(), Some(resumed_fsn));
done_tx.send(()).unwrap();
drop(sender);
handle.join().unwrap();
}
#[test]
fn qwp_ws_close_flush_timeout_minus_one_skips_close_drain_wait() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frame_tx, frame_rx) = mpsc::channel();
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream(&mut stream);
let (_fin, _opcode, payload) = read_frame(&mut stream).unwrap();
frame_tx.send(payload).unwrap();
let mut sink = [0u8; 256];
while matches!(stream.read(&mut sink), Ok(n) if n > 0) {}
});
let conf = format!("ws::addr=127.0.0.1:{port};close_flush_timeout_millis=-1;");
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert_eq!(fsn, 0);
let frame = frame_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
let started = Instant::now();
sender.close_drain().unwrap();
assert!(
started.elapsed() < Duration::from_secs(2),
"close_drain waited despite close_flush_timeout_millis=-1"
);
drop(sender);
server.join().unwrap();
}
#[test]
fn qwp_ws_drop_interrupts_blocked_background_send() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (send_started_tx, send_started_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
socket2::SockRef::from(&stream)
.set_recv_buffer_size(4096)
.unwrap();
upgrade_mock_stream(&mut stream);
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
let mut byte = [0u8; 1];
stream.peek(&mut byte).unwrap();
send_started_tx.send(()).unwrap();
let _ = release_rx.recv_timeout(Duration::from_secs(10));
});
let conf = format!(
"ws::addr=127.0.0.1:{port};\
close_flush_timeout_millis=-1;\
sf_max_segment_bytes=16777216;"
);
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let value = "x".repeat(8 * 1024 * 1024);
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_str("payload", value.as_str())
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
send_started_rx
.recv_timeout(Duration::from_secs(5))
.unwrap();
let started = Instant::now();
drop(sender);
let elapsed = started.elapsed();
let _ = release_tx.send(());
server.join().unwrap();
assert!(
elapsed < Duration::from_secs(1),
"drop waited for the socket write timeout: {elapsed:?}"
);
}
fn assert_qwp_ws_drop_interrupts_stalled_connect(scheme: &str, tls_options: &str) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let wait_for_client_hello = scheme == "wss";
let (connect_stalled_tx, connect_stalled_rx) = mpsc::channel();
let (release_tx, release_rx) = mpsc::channel();
let server = thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
if wait_for_client_hello {
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
let mut record_header = [0u8; 5];
stream.read_exact(&mut record_header).unwrap();
assert_eq!(record_header[0], 0x16, "expected a TLS handshake record");
let record_len = u16::from_be_bytes([record_header[3], record_header[4]]) as usize;
let mut record = vec![0u8; record_len];
stream.read_exact(&mut record).unwrap();
assert_eq!(record.first(), Some(&0x01), "expected a TLS ClientHello");
thread::sleep(Duration::from_millis(200));
}
connect_stalled_tx.send(()).unwrap();
let _ = release_rx.recv_timeout(Duration::from_secs(10));
});
let conf = format!(
"{scheme}::addr=127.0.0.1:{port};\
initial_connect_retry=async;\
{tls_options}"
);
let sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
connect_stalled_rx
.recv_timeout(Duration::from_secs(5))
.unwrap();
let started = Instant::now();
drop(sender);
let elapsed = started.elapsed();
let _ = release_tx.send(());
server.join().unwrap();
assert!(
elapsed < Duration::from_secs(1),
"{scheme} drop did not interrupt the stalled connect phase: {elapsed:?}"
);
}
#[test]
fn qwp_ws_drop_interrupts_stalled_websocket_upgrade() {
assert_qwp_ws_drop_interrupts_stalled_connect("ws", "");
}
#[test]
fn qwp_ws_drop_interrupts_stalled_tls_handshake() {
let mut cert = std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR"));
cert.pop();
cert.push("tls_certs/server_rootCA.pem");
let tls_options = format!("tls_roots={};", cert.display());
assert_qwp_ws_drop_interrupts_stalled_connect("wss", &tls_options);
}
#[test]
#[ignore = "performance benchmark"]
fn qwp_ws_public_sender_batch_throughput_benchmark() {
let rows =
qwp_ws_public_bench_env_usize("QWP_WS_PUBLIC_BENCH_ROWS", QWP_WS_PUBLIC_BENCH_DEFAULT_ROWS);
let batch_size = qwp_ws_public_bench_env_usize(
"QWP_WS_PUBLIC_BENCH_BATCH_SIZE",
QWP_WS_PUBLIC_BENCH_DEFAULT_BATCH_SIZE,
);
let workload = QwpWsPublicBenchWorkload::from_env();
let prevalidated_names = qwp_ws_public_bench_env_bool("QWP_WS_PUBLIC_BENCH_PREVALIDATED_NAMES");
assert!(rows > 0);
assert!(batch_size > 0);
let (port, server) = spawn_ack_each_frame_server();
let conf = format!("ws::addr=127.0.0.1:{port};");
let mut sender = SenderBuilder::from_conf(conf).unwrap().build().unwrap();
let mut buffer = sender.new_buffer();
fill_qwp_ws_public_benchmark_batch(
&mut buffer,
workload,
prevalidated_names,
0,
batch_size,
batch_size,
);
sender.flush(&mut buffer).unwrap();
let started = Instant::now();
let mut build_elapsed = Duration::ZERO;
let mut flush_elapsed = Duration::ZERO;
let mut published_rows = 0usize;
let mut batch_idx = 0usize;
while published_rows < rows {
let rows_in_batch = (rows - published_rows).min(batch_size);
let build_started = Instant::now();
fill_qwp_ws_public_benchmark_batch(
&mut buffer,
workload,
prevalidated_names,
batch_idx,
batch_size,
rows_in_batch,
);
build_elapsed += build_started.elapsed();
let flush_started = Instant::now();
sender.flush(&mut buffer).unwrap();
flush_elapsed += flush_started.elapsed();
published_rows += rows_in_batch;
batch_idx += 1;
}
let close_started = Instant::now();
sender.close_drain().unwrap();
let close_elapsed = close_started.elapsed();
let elapsed = started.elapsed();
drop(sender);
let binary_frames = server.join().unwrap();
let expected_frames = rows.div_ceil(batch_size) + 1;
assert_eq!(binary_frames, expected_frames);
eprintln!(
"qwp_ws_public_sender_batch_throughput workload={} prevalidated_names={} rows={} batch_size={} frames={} total_ms={} build_ms={} flush_ms={} close_ms={} rows_per_sec={:.2}",
workload.as_str(),
prevalidated_names,
rows,
batch_size,
binary_frames,
elapsed.as_millis(),
build_elapsed.as_millis(),
flush_elapsed.as_millis(),
close_elapsed.as_millis(),
rows as f64 / elapsed.as_secs_f64()
);
eprintln!(
"qwp_ws_public_sender_batch_build workload={} prevalidated_names={} rows={} batch_size={} elapsed_ms={} rows_per_sec={:.2}",
workload.as_str(),
prevalidated_names,
rows,
batch_size,
build_elapsed.as_millis(),
rows as f64 / build_elapsed.as_secs_f64()
);
eprintln!(
"qwp_ws_public_sender_batch_flush workload={} prevalidated_names={} rows={} batch_size={} elapsed_ms={} rows_per_sec={:.2}",
workload.as_str(),
prevalidated_names,
rows,
batch_size,
flush_elapsed.as_millis(),
rows as f64 / flush_elapsed.as_secs_f64()
);
}
#[test]
fn qwp_ws_manual_sender_can_pipeline_before_waiting() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frames_tx, frames_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let req_bytes = read_request_until_blank(&mut stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let response = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(response.as_bytes()).unwrap();
let mut received = Vec::new();
let (_fin, _opcode, first) = read_frame(&mut stream).unwrap();
received.push(first);
let (_fin, _opcode, second) = read_frame(&mut stream).unwrap();
received.push(second);
frames_tx.send(received).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE + 1).unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.qwp_ws_progress(QwpWsProgress::Manual)
.unwrap()
.build()
.unwrap();
let mut first = sender.new_buffer();
first
.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
let first_fsn = sender.flush_and_get_fsn(&mut first).unwrap().unwrap();
assert!(first.is_empty());
assert_eq!(first_fsn, 0);
let mut second = sender.new_buffer();
second
.table("trades")
.unwrap()
.symbol("sym", "BTC-USD")
.unwrap()
.column_i64("qty", 11)
.unwrap()
.at_now()
.unwrap();
let second_fsn = sender.flush_and_get_fsn(&mut second).unwrap().unwrap();
assert!(second.is_empty());
assert_eq!(second_fsn, 1);
assert_eq!(sender.published_fsn().unwrap(), Some(second_fsn));
assert_eq!(sender.acked_fsn().unwrap(), None);
assert!(sender.drive_once().unwrap());
assert!(sender.drive_once().unwrap());
let frames = frames_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(frames.len(), 2);
assert_eq!(&frames[0][0..4], b"QWP1");
assert_eq!(&frames[1][0..4], b"QWP1");
sender
.wait(crate::ingress::AckLevel::Ok, Duration::from_secs(5))
.unwrap();
assert_eq!(sender.acked_fsn().unwrap(), Some(second_fsn));
}
#[test]
fn qwp_ws_manual_sender_schema_rejection_terminalizes_without_ack_advance() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let req_bytes = read_request_until_blank(&mut stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let response = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(response.as_bytes()).unwrap();
let (_fin, _opcode, _first) = read_frame(&mut stream).unwrap();
write_qwp_error_response(
&mut stream,
QWP_STATUS_SCHEMA_MISMATCH,
FIRST_WIRE_SEQUENCE,
b"first bad",
)
.unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.qwp_ws_progress(QwpWsProgress::Manual)
.unwrap()
.build()
.unwrap();
let mut first = sender.new_buffer();
first
.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
let first_fsn = sender.flush_and_get_fsn(&mut first).unwrap().unwrap();
assert_eq!(first_fsn, 0);
assert_eq!(sender.acked_fsn().unwrap(), None);
let deadline = Instant::now() + Duration::from_secs(5);
let err = loop {
match sender.drive_once() {
Ok(_) if Instant::now() < deadline => {
assert_eq!(sender.acked_fsn().ok().flatten(), None);
thread::sleep(Duration::from_millis(10));
}
Ok(progressed) => {
panic!("schema rejection did not terminalize sender; last progressed={progressed}")
}
Err(err) => break err,
}
};
assert_eq!(err.code(), ErrorCode::ServerRejection);
assert_eq!(
err.qwp_ws_rejection().map(|error| error.category),
Some(QwpWsErrorCategory::SchemaMismatch)
);
let first_error = sender.poll_qwp_ws_error().unwrap().unwrap();
assert_eq!(first_error.category, QwpWsErrorCategory::SchemaMismatch);
assert_eq!(first_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(first_error.status, Some(QWP_STATUS_SCHEMA_MISMATCH));
assert_eq!(first_error.message.as_deref(), Some("first bad"));
assert_eq!(first_error.message_sequence, Some(FIRST_WIRE_SEQUENCE));
assert_eq!(first_error.from_fsn, first_fsn);
assert_eq!(first_error.to_fsn, first_fsn);
assert_eq!(sender.poll_qwp_ws_error().unwrap(), None);
assert_eq!(sender.qwp_ws_errors_dropped().unwrap(), 0);
}
#[test]
fn qwp_ws_store_and_forward_config_opens_java_slot_layout() {
let (port, rx) = spawn_mock_server();
let sf_dir = tempfile::TempDir::new().unwrap();
let conf = format!(
"ws::addr=127.0.0.1:{port};sf_dir={};sender_id=primary;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let _ = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(sf_dir.path().join("primary").join(".lock").exists());
}
#[test]
fn qwp_ws_store_and_forward_rejects_one_segment_total_capacity() {
let sf_dir = tempfile::TempDir::new().unwrap();
let conf = format!(
"ws::addr=127.0.0.1:1;sf_dir={};sender_id=primary;\
sf_max_segment_bytes=256;sf_max_total_bytes=256;",
sf_dir.path().display()
);
let err = SenderBuilder::from_conf(conf).unwrap().build().unwrap_err();
assert_eq!(err.code(), crate::ErrorCode::SocketError);
assert!(
err.msg().contains("Store-and-Forward queue") && err.msg().contains("InvalidCapacity"),
"got: {}",
err.msg()
);
assert!(
!sf_dir
.path()
.join("primary")
.join("sf-initial.sfa")
.exists()
);
}
#[test]
fn qwp_ws_manual_orphan_drainer_replays_sibling_slot() {
let seed_port = spawn_upgrade_only_server();
let sf_dir = tempfile::TempDir::new().unwrap();
let seed_conf = format!(
"ws::addr=127.0.0.1:{seed_port};qwp_ws_progress=manual;\
sf_dir={};sender_id=orphan;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut seed_sender = SenderBuilder::from_conf(&seed_conf)
.unwrap()
.build()
.unwrap();
let mut seed_buf = seed_sender.new_buffer();
seed_buf
.table("orphaned")
.unwrap()
.symbol("src", "old")
.unwrap()
.column_i64("value", 42)
.unwrap()
.at_now()
.unwrap();
seed_sender.flush(&mut seed_buf).unwrap();
drop(seed_sender);
let (port, rx) = spawn_manual_orphan_drain_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let mut orphan_payload = None;
for _ in 0..20 {
let _ = sender.drive_once().unwrap();
if let Ok(payload) = rx.try_recv() {
orphan_payload = Some(payload);
break;
}
thread::sleep(Duration::from_millis(10));
}
let payload = orphan_payload.expect("orphan payload was not replayed");
assert!(!payload.is_empty());
let orphan_slot = sf_dir.path().join("orphan");
for _ in 0..20 {
let _ = sender.drive_once().unwrap();
if !slot_has_sfa_file(&orphan_slot) {
break;
}
thread::sleep(Duration::from_millis(10));
}
assert!(!slot_has_sfa_file(&orphan_slot));
assert!(!sf_dir.path().join("primary").join(".failed").exists());
assert!(!orphan_slot.join(".failed").exists());
}
#[test]
fn qwp_ws_orphan_dict_copy_oom_retries_without_failed_sentinel() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot_with_two_delta_frames(sf_dir.path());
let orphan_slot = sf_dir.path().join("orphan");
let (port, drained_rx) = spawn_two_frame_orphan_drain_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
fail_next_recovered_dict_copy_for_test();
assert!(sender.drive_once().unwrap());
let last_error = orphan_slot.join(".last_error");
assert!(
!orphan_slot.join(".failed").exists(),
"transient dictionary allocation failure must not quarantine the slot"
);
let retry_reason = std::fs::read_to_string(&last_error)
.expect("transient dictionary allocation failure must record its retry reason");
assert!(
retry_reason.contains("injected recovered symbol dictionary allocation failure"),
"unexpected orphan retry reason: {retry_reason}"
);
assert!(
slot_has_sfa_file(&orphan_slot),
"a retryable allocation failure must retain the durable slot"
);
let deadline = Instant::now() + Duration::from_secs(5);
let mut drained_frames = None;
while Instant::now() < deadline {
if let Ok(frames) = drained_rx.try_recv() {
drained_frames = Some(frames);
break;
}
let _ = sender.drive_once().unwrap();
thread::sleep(Duration::from_millis(10));
}
let drained_frames = drained_frames.expect("the retried orphan drain did not complete");
assert_eq!(drained_frames.len(), 2);
let mut delta_pos = 12;
assert_eq!(
read_varint(&drained_frames[1], &mut delta_pos),
1,
"the second persisted frame must depend on recovered dictionary id 0"
);
assert!(
!orphan_slot.join(".failed").exists(),
"a retried allocation failure must never write the permanent sentinel"
);
let fully_drained = wait_until(Duration::from_secs(5), || {
let _ = sender.drive_once().unwrap();
!slot_has_sfa_file(&orphan_slot) && !last_error.exists()
});
assert!(
fully_drained,
"the intact delta slot did not drain cleanly after the allocation failure cleared; \
has_sfa={}, last_error={}",
slot_has_sfa_file(&orphan_slot),
std::fs::read_to_string(&last_error).unwrap_or_default()
);
assert!(!orphan_slot.join(".failed").exists());
}
fn spawn_orphan_capture_first_frame_server() -> (u16, mpsc::Receiver<Vec<u8>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut orphan, _) = listener.accept().unwrap();
orphan
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
orphan
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
perform_server_upgrade(&mut orphan).unwrap();
if let Ok((_fin, _op, first)) = read_frame(&mut orphan) {
let _ = tx.send(first);
let _ = write_qwp_ok_response(&mut orphan, FIRST_WIRE_SEQUENCE);
}
thread::sleep(Duration::from_millis(50));
});
(port, rx)
}
#[test]
fn qwp_ws_orphan_drain_heals_a_zero_extended_side_file_and_replays_via_delta() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot(sf_dir.path());
let side_file = sf_dir.path().join("orphan").join(".symbol-dict");
{
let mut f = std::fs::OpenOptions::new()
.append(true)
.open(&side_file)
.expect("seed must have written a delta-mode side-file");
f.write_all(&[0u8; 4]).unwrap();
}
let (port, rx) = spawn_orphan_capture_first_frame_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let mut first_frame = None;
for _ in 0..60 {
let _ = sender.drive_once().unwrap();
if let Ok(frame) = rx.try_recv() {
first_frame = Some(frame);
break;
}
thread::sleep(Duration::from_millis(10));
}
let frame = first_frame.expect("orphan drainer sent no frame");
assert!(
frame.len() >= 12 && &frame[..4] == b"QWP1",
"not a QWP frame: {:?}",
&frame[..frame.len().min(12)]
);
let table_count = u16::from_le_bytes([frame[6], frame[7]]);
assert_eq!(
table_count, 0,
"the zero tail is healed at open, so the orphan arms delta on the clean \
recovered dictionary and re-registers the real symbol via a table-less \
catch-up first (table_count == 0); a DATA frame first (table_count >= 1) \
would mean it fell back to dense. table_count = {table_count}"
);
assert!(
!sf_dir.path().join("orphan").join(".failed").exists(),
"healing a zero tail must keep the slot recoverable, not fail it"
);
}
#[test]
fn qwp_ws_orphan_drain_discards_a_duplicate_entry_dict_and_arms_the_mirror_empty() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot(sf_dir.path());
let side_file = sf_dir.path().join("orphan").join(".symbol-dict");
{
let existing =
std::fs::read(&side_file).expect("seed must have written a delta-mode side-file");
let entries_region = existing[8..].to_vec();
assert!(
!entries_region.is_empty(),
"seed must have persisted at least one symbol entry to duplicate"
);
let mut f = std::fs::OpenOptions::new()
.append(true)
.open(&side_file)
.unwrap();
f.write_all(&entries_region).unwrap();
}
let (port, rx) = spawn_orphan_capture_first_frame_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let mut first_frame = None;
for _ in 0..60 {
let _ = sender.drive_once().unwrap();
if let Ok(frame) = rx.try_recv() {
first_frame = Some(frame);
break;
}
thread::sleep(Duration::from_millis(10));
}
let frame = first_frame.expect("orphan drainer sent no frame");
assert!(
frame.len() >= 12 && &frame[..4] == b"QWP1",
"not a QWP frame: {:?}",
&frame[..frame.len().min(12)]
);
let table_count = u16::from_le_bytes([frame[6], frame[7]]);
assert!(
table_count >= 1,
"a corrupt (duplicate-entry) recovered dictionary must be DISCARDED, so the \
mirror arms empty, emits no catch-up, and the DATA frame goes first \
(table_count >= 1); a table-less catch-up first (table_count == 0) would mean \
it wrongly seeded the mirror from the corrupt dictionary. \
table_count = {table_count}"
);
assert!(
!sf_dir.path().join("orphan").join(".failed").exists(),
"a corrupt recovered dictionary must leave the slot recoverable, not fail it"
);
}
#[test]
fn qwp_ws_orphan_drain_arms_an_empty_mirror_when_the_recovered_dict_exceeds_the_cap() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot_with_two_delta_frames(sf_dir.path());
let _cap = crate::ingress::TestDictCapGuard::new(1);
let (port, rx) = spawn_orphan_drain_all_frames_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let deadline = Instant::now() + Duration::from_secs(10);
let mut frames: Vec<Vec<u8>> = Vec::new();
let mut data_frames = 0usize;
while Instant::now() < deadline && data_frames < 2 {
let _ = sender.drive_once().unwrap();
while let Ok(frame) = rx.try_recv() {
if frame.len() >= 12 && &frame[..4] == b"QWP1" {
if u16::from_le_bytes([frame[6], frame[7]]) >= 1 {
data_frames += 1;
}
frames.push(frame);
}
}
thread::sleep(Duration::from_millis(10));
}
let first = frames.first().expect(
"the drainer must still connect and replay -- an over-cap recovered \
dictionary must not abort the drain",
);
let first_table_count = u16::from_le_bytes([first[6], first[7]]);
assert!(
first_table_count >= 1,
"the rejected entries must be DISCARDED, so the empty mirror emits no \
catch-up and the DATA frame goes first (table_count >= 1). A table-less \
catch-up first (table_count == 0) would mean the drain re-registered a \
dictionary it had just refused to seed. table_count = {first_table_count}"
);
assert!(
data_frames >= 2,
"both queued frames must replay -- got {data_frames}. One means the drain \
fell back to dense, terminally rejected its `delta_start == 1` frame, and \
abandoned the slot behind a `.failed` sentinel"
);
assert!(
!sf_dir.path().join("orphan").join(".failed").exists(),
"an over-cap side-file is not proven-local unrecoverable: the frames carry \
the dictionary they need, so the slot must drain rather than be quarantined"
);
}
#[test]
fn qwp_ws_orphan_drain_replays_both_frames_when_the_first_dict_chunk_is_corrupt() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot_with_two_delta_frames(sf_dir.path());
let side_file = sf_dir.path().join("orphan").join(".symbol-dict");
{
let mut bytes =
std::fs::read(&side_file).expect("seed must have written a delta-mode side-file");
let idx = bytes
.windows(5)
.position(|w| w == b"alpha")
.expect("alpha payload present");
bytes[idx] = b'X'; std::fs::write(&side_file, &bytes).unwrap();
}
let (port, rx) = spawn_orphan_drain_all_frames_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let deadline = Instant::now() + Duration::from_secs(10);
let mut data_frames = 0usize;
while Instant::now() < deadline && data_frames < 2 {
let _ = sender.drive_once().unwrap();
while let Ok(frame) = rx.try_recv() {
if frame.len() >= 12
&& &frame[..4] == b"QWP1"
&& u16::from_le_bytes([frame[6], frame[7]]) >= 1
{
data_frames += 1;
}
}
thread::sleep(Duration::from_millis(10));
}
assert!(
data_frames >= 2,
"both queued frames must replay -- got {data_frames}. One means the orphan \
slot fell back to dense, terminally rejected its `delta_start == 1` frame, \
and was re-queued by `RetryLater` to re-open, re-decide dense and re-strand: \
a live-lock that never drains and never fails"
);
assert!(
!sf_dir.path().join("orphan").join(".failed").exists(),
"draining a torn-first-chunk slot must keep it recoverable, not fail it"
);
}
#[test]
fn qwp_ws_manual_orphan_drainer_walks_endpoint_list() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot(sf_dir.path());
let bad_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let bad_port = bad_listener.local_addr().unwrap().port();
let (port, rx) = spawn_manual_orphan_drain_server();
drop(bad_listener);
let drain_conf = format!(
"ws::addr=127.0.0.1:{bad_port},127.0.0.1:{port};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let mut orphan_payload = None;
for _ in 0..20 {
let _ = sender.drive_once().unwrap();
if let Ok(payload) = rx.try_recv() {
orphan_payload = Some(payload);
break;
}
thread::sleep(Duration::from_millis(10));
}
assert!(
!orphan_payload
.expect("orphan payload was not replayed")
.is_empty()
);
assert!(!sf_dir.path().join("orphan").join(".failed").exists());
}
#[test]
fn qwp_ws_manual_orphan_drainer_role_reject_tries_next_endpoint() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot(sf_dir.path());
let (reject_port, reject_handle) = spawn_role_reject_upgrade_server(2, "REPLICA");
let (port, rx) = spawn_manual_orphan_drain_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{reject_port},127.0.0.1:{port};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let mut orphan_payload = None;
for _ in 0..20 {
let _ = sender.drive_once().unwrap();
if let Ok(payload) = rx.try_recv() {
orphan_payload = Some(payload);
break;
}
thread::sleep(Duration::from_millis(10));
}
assert!(
!orphan_payload
.expect("orphan payload was not replayed")
.is_empty()
);
assert_eq!(reject_handle.join().unwrap(), 2);
assert!(!sf_dir.path().join("orphan").join(".failed").exists());
}
#[test]
fn qwp_ws_manual_orphan_drainer_terminal_reject_leaves_slot_recoverable() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot(sf_dir.path());
let (port, rx) = spawn_manual_orphan_reject_server(QWP_STATUS_PARSE_ERROR);
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let orphan_slot = sf_dir.path().join("orphan");
let mut orphan_payload = None;
for _ in 0..20 {
let _ = sender.drive_once().unwrap();
if orphan_payload.is_none()
&& let Ok(payload) = rx.try_recv()
{
orphan_payload = Some(payload);
}
if orphan_payload.is_some() && orphan_slot.join(".last_error").exists() {
break;
}
thread::sleep(Duration::from_millis(10));
}
assert!(
!orphan_payload
.expect("orphan payload was not replayed")
.is_empty()
);
assert!(!orphan_slot.join(".failed").exists());
assert!(orphan_slot.join(".last_error").exists());
assert!(slot_has_sfa_file(&orphan_slot));
}
#[test]
fn qwp_ws_background_terminal_orphan_releases_worker_for_next_slot() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot_named(sf_dir.path(), "orphan-a");
seed_orphan_slot_named(sf_dir.path(), "orphan-b");
let server = spawn_terminal_then_drain_orphan_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{};\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
server.port,
sf_dir.path().display()
);
let sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
assert!(
!server
.terminal_rx
.recv_timeout(Duration::from_secs(5))
.expect("first orphan was not replayed and terminally rejected")
.is_empty()
);
assert!(
!server
.drained_rx
.recv_timeout(Duration::from_secs(2))
.expect("terminal orphan retained the only worker; next slot was never drained")
.is_empty()
);
let orphan_a = sf_dir.path().join("orphan-a");
let orphan_b = sf_dir.path().join("orphan-b");
wait_until(Duration::from_secs(5), || {
[&orphan_a, &orphan_b]
.into_iter()
.filter(|slot| slot_has_sfa_file(slot))
.count()
== 1
&& [&orphan_a, &orphan_b]
.into_iter()
.filter(|slot| slot.join(".last_error").exists())
.count()
== 1
});
let terminal_slots = [&orphan_a, &orphan_b]
.into_iter()
.filter(|slot| slot.join(".last_error").exists())
.count();
assert_eq!(
terminal_slots, 1,
"exactly one orphan must record the terminal error"
);
let retained_slots = [&orphan_a, &orphan_b]
.into_iter()
.filter(|slot| slot_has_sfa_file(slot))
.count();
assert_eq!(
retained_slots, 1,
"the terminal slot must retain its data while the next slot fully drains"
);
assert_eq!(
orphan_a.join(".last_error").exists(),
slot_has_sfa_file(&orphan_a)
);
assert_eq!(
orphan_b.join(".last_error").exists(),
slot_has_sfa_file(&orphan_b)
);
assert!(!orphan_a.join(".failed").exists());
assert!(!orphan_b.join(".failed").exists());
drop(sender);
server.handle.join().unwrap();
}
#[test]
fn qwp_ws_background_role_reject_recycles_wire_and_drains_orphan() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot(sf_dir.path());
let server = spawn_role_reject_then_drain_orphan_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{};\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;\
reconnect_initial_backoff_millis=10;reconnect_max_backoff_millis=20;\
sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
server.port,
sf_dir.path().display()
);
let sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let rejected = server
.rejected_rx
.recv_timeout(Duration::from_secs(5))
.expect("orphan was not replayed to the replica");
let drained = server
.drained_rx
.recv_timeout(Duration::from_secs(2))
.expect("role-rejected orphan was not retried on a recycled connection");
assert_eq!(
rejected, drained,
"the recycled connection must replay the same frame"
);
let orphan_slot = sf_dir.path().join("orphan");
let last_error = orphan_slot.join(".last_error");
wait_until(Duration::from_secs(5), || {
!slot_has_sfa_file(&orphan_slot) && !last_error.exists()
});
assert!(!slot_has_sfa_file(&orphan_slot));
assert!(!orphan_slot.join(".failed").exists());
assert!(
!last_error.exists(),
"stale .last_error after drain: {}",
std::fs::read_to_string(&last_error).unwrap_or_default()
);
drop(sender);
server.handle.join().unwrap();
}
#[test]
fn qwp_ws_close_drain_blocks_across_all_replica_window_until_promotion() {
let server = spawn_replica_window_then_promote_server();
let conf = format!(
"ws::addr=127.0.0.1:{};initial_connect_retry=async;\
reconnect_initial_backoff_millis=10;reconnect_max_backoff_millis=50;\
close_flush_timeout_millis=10000;",
server.port
);
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let closer = thread::spawn(move || sender.close_drain());
assert!(
wait_until(Duration::from_secs(5), || {
server.reject_count.load(Ordering::Acquire) >= 2
}),
"server never produced the all-replica window"
);
assert!(
!closer.is_finished(),
"close_drain must stay pending for the whole all-replica window"
);
server.promoted.store(true, Ordering::Release);
closer
.join()
.unwrap()
.expect("close_drain must complete cleanly after promotion");
let delivered = server.handle.join().unwrap();
assert_eq!(
delivered.len(),
1,
"the queued frame must be delivered exactly once"
);
assert_eq!(&delivered[0][0..4], b"QWP1");
}
#[test]
fn qwp_ws_manual_orphan_terminal_retires_slot_and_lets_caller_park() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot(sf_dir.path());
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (done_tx, done_rx) = mpsc::channel::<()>();
let server = thread::spawn(move || {
let (mut foreground, _) = listener.accept().unwrap();
perform_server_upgrade(&mut foreground).unwrap();
let (mut orphan, _) = listener.accept().unwrap();
perform_server_upgrade(&mut orphan).unwrap();
let (_fin, _opcode, _catch_up) = read_frame(&mut orphan).unwrap();
let (_fin, _opcode, payload) = read_frame(&mut orphan).unwrap();
write_qwp_error_response(
&mut orphan,
QWP_STATUS_PARSE_ERROR,
FIRST_WIRE_SEQUENCE + 1,
b"bad orphan",
)
.unwrap();
let _ = done_rx.recv_timeout(Duration::from_secs(10));
payload
});
let conf = format!(
"ws::addr=127.0.0.1:{};qwp_ws_progress=manual;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_frame_rejections=1;poison_min_escalation_window_millis=0;\
reconnect_initial_backoff_millis=10;reconnect_max_backoff_millis=20;\
sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
port,
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let orphan_slot = sf_dir.path().join("orphan");
let last_error = orphan_slot.join(".last_error");
let started = Instant::now();
while started.elapsed() < Duration::from_secs(5) && !last_error.exists() {
if !sender.drive_once().unwrap() {
thread::sleep(Duration::from_millis(1));
}
}
assert!(
last_error.exists(),
"session-sticky terminal must leave a breadcrumb"
);
let quiesce_started = Instant::now();
while sender.drive_once().unwrap() {
assert!(
quiesce_started.elapsed() < Duration::from_secs(2),
"drive_once must stop reporting progress once the orphan is terminal"
);
}
for _ in 0..5 {
assert!(
!sender.drive_once().unwrap(),
"a terminal orphan must let the caller park, not spin"
);
}
assert!(
!orphan_slot.join(".failed").exists(),
"a session-sticky terminal must not poison the slot"
);
assert!(
slot_has_sfa_file(&orphan_slot),
"the undrained frame must stay on disk for the next session"
);
done_tx.send(()).unwrap();
let payload = server.join().unwrap();
assert!(!payload.is_empty());
drop(sender);
}
#[test]
fn qwp_ws_background_catch_up_allocation_failure_reconnects_and_drains_orphan() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot_named_with_symbol(sf_dir.path(), "orphan", "qdb-test-catch-up-allocation");
let server = spawn_catch_up_failure_then_drain_orphan_server(None);
fail_next_catch_up_allocation_for_test();
let drain_conf = format!(
"ws::addr=127.0.0.1:{};\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;\
reconnect_initial_backoff_millis=10;reconnect_max_backoff_millis=20;\
sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
server.port,
sf_dir.path().display()
);
let sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
server
.retried_rx
.recv_timeout(Duration::from_secs(5))
.expect("catch-up allocation failure did not trigger a reconnect");
assert!(
!server
.drained_rx
.recv_timeout(Duration::from_secs(2))
.expect("catch-up allocation failure abandoned the orphan slot")
.is_empty()
);
let orphan_slot = sf_dir.path().join("orphan");
let last_error = orphan_slot.join(".last_error");
assert!(wait_until(Duration::from_secs(5), || {
!slot_has_sfa_file(&orphan_slot) && !last_error.exists()
}));
assert!(!orphan_slot.join(".failed").exists());
drop(sender);
server.handle.join().unwrap();
}
#[test]
fn qwp_ws_background_catch_up_batch_cap_reconnects_and_drains_orphan() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot(sf_dir.path());
let server = spawn_catch_up_failure_then_drain_orphan_server(Some(28));
let drain_conf = format!(
"ws::addr=127.0.0.1:{};\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;\
reconnect_initial_backoff_millis=10;reconnect_max_backoff_millis=20;\
sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
server.port,
sf_dir.path().display()
);
let sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
server
.retried_rx
.recv_timeout(Duration::from_secs(5))
.expect("catch-up batch cap did not trigger a reconnect");
assert!(
!server
.drained_rx
.recv_timeout(Duration::from_secs(2))
.expect("catch-up BatchTooLarge abandoned the orphan slot")
.is_empty()
);
let orphan_slot = sf_dir.path().join("orphan");
let last_error = orphan_slot.join(".last_error");
assert!(wait_until(Duration::from_secs(5), || {
!slot_has_sfa_file(&orphan_slot) && !last_error.exists()
}));
assert!(!orphan_slot.join(".failed").exists());
drop(sender);
server.handle.join().unwrap();
}
#[test]
fn qwp_ws_background_orphan_close_is_bounded_and_leaves_orphan_recoverable() {
let seed_port = spawn_upgrade_only_server();
let sf_dir = tempfile::TempDir::new().unwrap();
let seed_conf = format!(
"ws::addr=127.0.0.1:{seed_port};qwp_ws_progress=manual;\
sf_dir={};sender_id=orphan;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut seed_sender = SenderBuilder::from_conf(&seed_conf)
.unwrap()
.build()
.unwrap();
let mut seed_buf = seed_sender.new_buffer();
seed_buf
.table("orphaned")
.unwrap()
.symbol("src", "old")
.unwrap()
.column_i64("value", 42)
.unwrap()
.at_now()
.unwrap();
seed_sender.flush(&mut seed_buf).unwrap();
drop(seed_sender);
let orphan_slot = sf_dir.path().join("orphan");
assert!(slot_has_sfa_file(&orphan_slot));
let (port, rx, release_stalled_orphan) = spawn_stalled_background_orphan_drain_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let orphan_payload = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(!orphan_payload.is_empty());
let started = Instant::now();
sender.close_drain().unwrap();
let elapsed = started.elapsed();
assert!(
elapsed < Duration::from_secs(5),
"background orphan shutdown took {elapsed:?}"
);
assert!(!orphan_slot.join(".failed").exists());
release_stalled_orphan.send(()).unwrap();
let (recover_port, recover_rx) = spawn_recovery_mock_server();
let recover_conf = format!(
"ws::addr=127.0.0.1:{recover_port};\
sf_dir={};sender_id=orphan;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let retry_deadline = Instant::now() + Duration::from_secs(5);
let mut recover_sender = loop {
match SenderBuilder::from_conf(&recover_conf).unwrap().build() {
Ok(sender) => break sender,
Err(err) if Instant::now() < retry_deadline => {
let _ = err;
thread::sleep(Duration::from_millis(10));
}
Err(err) => panic!("orphan slot was not reusable after close: {err}"),
}
};
recover_sender.close_drain().unwrap();
let recovered = recover_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(recovered.received_frames.len(), 1);
assert!(!recovered.received_frames[0].is_empty());
}
#[test]
fn qwp_ws_background_orphan_close_interrupts_stalled_connect() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot(sf_dir.path());
let (port, orphan_accepted, release_stalled_orphan) =
spawn_stalled_background_orphan_connect_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
orphan_accepted
.recv_timeout(Duration::from_secs(5))
.unwrap();
let started = Instant::now();
sender.close_drain().unwrap();
let elapsed = started.elapsed();
assert!(
elapsed < Duration::from_secs(5),
"background orphan shutdown took {elapsed:?}"
);
let (recover_port, recover_rx) = spawn_recovery_mock_server();
let recover_conf = format!(
"ws::addr=127.0.0.1:{recover_port};\
sf_dir={};sender_id=orphan;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut recover_sender = SenderBuilder::from_conf(&recover_conf)
.unwrap()
.build()
.expect("stalled orphan worker retained the slot lock after close");
recover_sender.close_drain().unwrap();
let recovered = recover_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(recovered.received_frames.len(), 1);
release_stalled_orphan.send(()).unwrap();
}
#[test]
fn qwp_ws_background_orphan_close_does_not_dial_next_slot() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_orphan_slot_named(sf_dir.path(), "orphan-a");
seed_orphan_slot_named(sf_dir.path(), "orphan-b");
let (port, listener_rx, release_stalled_orphan, server) =
spawn_stalled_first_background_orphan_connect_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=256;sf_max_total_bytes=1024;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
let listener = listener_rx.recv_timeout(Duration::from_secs(5)).unwrap();
sender.close_drain().unwrap();
listener.set_nonblocking(true).unwrap();
assert!(matches!(
listener.accept(),
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock
));
assert!(slot_has_sfa_file(&sf_dir.path().join("orphan-a")));
assert!(slot_has_sfa_file(&sf_dir.path().join("orphan-b")));
release_stalled_orphan.send(()).unwrap();
server.join().unwrap();
}
#[test]
fn qwp_ws_background_orphan_close_interrupts_blocked_send() {
let sf_dir = tempfile::TempDir::new().unwrap();
seed_large_orphan_slot(sf_dir.path());
let (port, send_started, release_stalled_orphan, server) =
spawn_blocked_background_orphan_send_server();
let drain_conf = format!(
"ws::addr=127.0.0.1:{port};close_flush_timeout_millis=-1;\
sf_dir={};sender_id=primary;drain_orphans=on;\
max_background_drainers=1;sf_max_segment_bytes=16777216;",
sf_dir.path().display()
);
let mut sender = SenderBuilder::from_conf(&drain_conf)
.unwrap()
.build()
.unwrap();
send_started.recv_timeout(Duration::from_secs(5)).unwrap();
let started = Instant::now();
sender.close_drain().unwrap();
let elapsed = started.elapsed();
assert!(
elapsed < Duration::from_secs(5),
"background orphan shutdown waited for the socket write timeout: {elapsed:?}"
);
let reopen_port = spawn_upgrade_only_server();
let reopen_conf = format!(
"ws::addr=127.0.0.1:{reopen_port};qwp_ws_progress=manual;\
sf_dir={};sender_id=orphan;sf_max_segment_bytes=16777216;",
sf_dir.path().display()
);
let reopened = SenderBuilder::from_conf(&reopen_conf)
.unwrap()
.build()
.expect("blocked orphan worker retained the slot lock after close");
drop(reopened);
assert!(slot_has_sfa_file(&sf_dir.path().join("orphan")));
release_stalled_orphan.send(()).unwrap();
server.join().unwrap();
}
#[test]
fn qwp_ws_subsequent_message_delta_encodes_dictionary_and_reemits_full_schema() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let req_bytes = read_request_until_blank(&mut stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let resp = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(resp.as_bytes()).unwrap();
for seq in 0u64..2 {
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
tx.send(payload).unwrap();
let mut ok = vec![0u8];
ok.extend_from_slice(&seq.to_le_bytes());
ok.extend_from_slice(&0u16.to_le_bytes());
write_server_binary_frame(&mut stream, &ok).unwrap();
}
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "BTC-USD")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
buf.table("trades")
.unwrap()
.symbol("sym", "BTC-USD")
.unwrap()
.column_i64("qty", 2)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let first = rx.recv_timeout(Duration::from_secs(5)).unwrap();
let second = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert!(second.len() >= 12);
let payload = &second[12..];
assert_eq!(payload[0], 0x01, "delta_start = 1");
assert_eq!(payload[1], 0x00, "delta_count = 0");
assert_eq!(first_table_column_count(&first), 2);
assert_eq!(first_table_column_count(&second), 2);
}
#[test]
fn qwp_ws_replay_full_schema_used_when_columns_match() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let req_bytes = read_request_until_blank(&mut stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let resp = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(resp.as_bytes()).unwrap();
for seq in 0u64..3 {
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
tx.send(payload).unwrap();
let mut ok = vec![0u8];
ok.extend_from_slice(&seq.to_le_bytes());
ok.extend_from_slice(&0u16.to_le_bytes());
write_server_binary_frame(&mut stream, &ok).unwrap();
}
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.build()
.unwrap();
for qty in 1..=3 {
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", qty)
.unwrap()
.column_f64("price", qty as f64 + 0.5)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
}
let m1 = rx.recv_timeout(Duration::from_secs(5)).unwrap();
let m2 = rx.recv_timeout(Duration::from_secs(5)).unwrap();
let m3 = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(
first_table_column_count(&m1),
2,
"first message: full schema"
);
assert_eq!(
first_table_column_count(&m2),
2,
"second message: full schema"
);
assert_eq!(
first_table_column_count(&m3),
2,
"third message: full schema"
);
}
#[test]
fn qwp_ws_full_schema_re_emitted_when_columns_change() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let req_bytes = read_request_until_blank(&mut stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let resp = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(resp.as_bytes()).unwrap();
for seq in 0u64..2 {
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
tx.send(payload).unwrap();
let mut ok = vec![0u8];
ok.extend_from_slice(&seq.to_le_bytes());
ok.extend_from_slice(&0u16.to_le_bytes());
write_server_binary_frame(&mut stream, &ok).unwrap();
}
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 2)
.unwrap()
.column_f64("price", 99.9)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let m1 = rx.recv_timeout(Duration::from_secs(5)).unwrap();
let m2 = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(first_table_column_count(&m1), 1);
assert_eq!(
first_table_column_count(&m2),
2,
"different column set must re-emit the full inline schema"
);
}
fn read_varint(buf: &[u8], pos: &mut usize) -> u64 {
let mut shift = 0;
let mut result: u64 = 0;
loop {
let b = buf[*pos];
*pos += 1;
result |= ((b & 0x7F) as u64) << shift;
if b & 0x80 == 0 {
return result;
}
shift += 7;
}
}
fn first_table_column_count(frame: &[u8]) -> u64 {
let mut pos = 12; let _delta_start = read_varint(frame, &mut pos);
let delta_count = read_varint(frame, &mut pos);
for _ in 0..delta_count {
let name_len = read_varint(frame, &mut pos) as usize;
pos += name_len;
}
let name_len = read_varint(frame, &mut pos) as usize;
pos += name_len;
let _row_count = read_varint(frame, &mut pos);
read_varint(frame, &mut pos)
}
#[test]
fn qwp_ws_server_error_response_is_surfaced() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (error_tx, error_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let req_bytes = read_request_until_blank(&mut stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let resp = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(resp.as_bytes()).unwrap();
let _ = read_frame(&mut stream).unwrap();
write_qwp_error_response(
&mut stream,
QWP_STATUS_PARSE_ERROR,
FIRST_WIRE_SEQUENCE,
b"bad column",
)
.unwrap();
let mut sink = [0u8; 256];
while matches!(stream.read(&mut sink), Ok(n) if n > 0) {}
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.qwp_ws_error_handler(move |error| {
error_tx.send(error.clone()).unwrap();
})
.unwrap()
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
let first_fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert!(buf.is_empty());
assert_eq!(first_fsn, 0);
let deadline = std::time::Instant::now() + Duration::from_secs(5);
while sender.qwp_ws_terminal_error().unwrap().is_none() {
assert!(
std::time::Instant::now() < deadline,
"server rejection was not applied within 5s"
);
thread::sleep(Duration::from_millis(1));
}
buf.table("trades")
.unwrap()
.column_i64("qty", 2)
.unwrap()
.at_now()
.unwrap();
let err = sender.flush(&mut buf).unwrap_err();
assert_eq!(err.code(), ErrorCode::ServerRejection);
assert!(
err.msg().contains("bad column"),
"expected server error in message, got: {}",
err.msg()
);
assert_eq!(
err.qwp_ws_rejection().map(|error| error.category),
Some(QwpWsErrorCategory::ParseError)
);
assert!(
!buf.is_empty(),
"terminal async error must not clear a newly prepared buffer"
);
let callback_error = error_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(callback_error.category, QwpWsErrorCategory::ParseError);
assert_eq!(callback_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(callback_error.from_fsn, first_fsn);
let qwp_error = sender.poll_qwp_ws_error().unwrap().unwrap();
assert_eq!(qwp_error.category, QwpWsErrorCategory::ParseError);
assert_eq!(qwp_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(qwp_error.status, Some(QWP_STATUS_PARSE_ERROR));
assert_eq!(qwp_error.message.as_deref(), Some("bad column"));
assert_eq!(qwp_error.message_sequence, Some(FIRST_WIRE_SEQUENCE));
assert_eq!(qwp_error.from_fsn, first_fsn);
assert_eq!(qwp_error.to_fsn, first_fsn);
assert_eq!(sender.poll_qwp_ws_error().unwrap(), None);
assert_eq!(sender.qwp_ws_errors_dropped().unwrap(), 0);
}
#[test]
fn qwp_ws_schema_rejection_terminalizes_and_notifies_handler() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
let (error_tx, error_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let req_bytes = read_request_until_blank(&mut stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let resp = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(resp.as_bytes()).unwrap();
let mut received_frames = Vec::new();
let (_fin, _opcode, first) = read_frame(&mut stream).unwrap();
received_frames.push(first);
write_qwp_error_response(
&mut stream,
QWP_STATUS_SCHEMA_MISMATCH,
FIRST_WIRE_SEQUENCE,
b"bad schema",
)
.unwrap();
tx.send(received_frames).unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.qwp_ws_error_handler(move |error| {
error_tx.send(error.clone()).unwrap();
})
.unwrap()
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
let first_fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert!(buf.is_empty());
assert_eq!(first_fsn, 0);
let received_frames = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(received_frames.len(), 1);
let err = sender
.wait(crate::ingress::AckLevel::Ok, Duration::from_secs(5))
.unwrap_err();
assert_eq!(err.code(), ErrorCode::ServerRejection);
assert_eq!(
err.qwp_ws_rejection().map(|error| error.category),
Some(QwpWsErrorCategory::SchemaMismatch)
);
let qwp_error = sender.poll_qwp_ws_error().unwrap().unwrap();
assert_eq!(qwp_error.category, QwpWsErrorCategory::SchemaMismatch);
assert_eq!(qwp_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(qwp_error.status, Some(QWP_STATUS_SCHEMA_MISMATCH));
assert_eq!(qwp_error.message.as_deref(), Some("bad schema"));
assert_eq!(qwp_error.message_sequence, Some(FIRST_WIRE_SEQUENCE));
assert_eq!(qwp_error.from_fsn, first_fsn);
assert_eq!(qwp_error.to_fsn, first_fsn);
assert_eq!(sender.poll_qwp_ws_error().unwrap(), None);
assert_eq!(sender.qwp_ws_errors_dropped().unwrap(), 0);
let _ = sender.flush(&mut buf);
let callback_error = error_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(callback_error.category, QwpWsErrorCategory::SchemaMismatch);
assert_eq!(callback_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(callback_error.from_fsn, first_fsn);
}
#[test]
fn qwp_ws_repeated_head_close_poison_is_pollable_as_protocol_violation() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frame_tx, frame_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream(&mut stream);
let (_fin, _opcode, payload) = read_frame(&mut stream).unwrap();
frame_tx.send(payload).unwrap();
write_server_close_frame(&mut stream, 1002, "bad frame").unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_frame_rejections(1)
.unwrap()
.poison_min_escalation_window(Duration::ZERO)
.unwrap()
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert!(buf.is_empty());
assert_eq!(fsn, 0);
let frame = frame_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
let mut observed = None;
for _ in 0..500 {
if let Some(error) = sender.poll_qwp_ws_error().unwrap() {
observed = Some(error);
break;
}
thread::sleep(Duration::from_millis(10));
}
let close_error = observed.expect("expected close poison diagnostic");
assert_eq!(close_error.category, QwpWsErrorCategory::ProtocolViolation);
assert_eq!(close_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(close_error.status, None);
assert_eq!(close_error.message_sequence, None);
assert!(close_error.message.as_deref().is_some_and(|message| {
message.contains("QWP/WebSocket frame fsn 0 was closed 1 times without ACK progress")
}));
assert_eq!(close_error.from_fsn, fsn);
assert_eq!(close_error.to_fsn, fsn);
assert_eq!(sender.poll_qwp_ws_error().unwrap(), None);
assert_eq!(sender.qwp_ws_errors_dropped().unwrap(), 0);
let err = sender.flush(&mut buf).unwrap_err();
assert!(
err.msg().contains("closed 1 times without ACK progress"),
"expected close poison in empty flush message, got: {}",
err.msg()
);
assert!(
buf.is_empty(),
"empty terminal flush must leave buffer empty"
);
buf.table("trades").unwrap().column_i64("qty", 2).unwrap();
let err = sender.flush(&mut buf).unwrap_err();
assert!(
err.msg().contains("closed 1 times without ACK progress"),
"expected close poison to dominate incomplete-row validation, got: {}",
err.msg()
);
assert!(
!err.msg().contains("Bad call to `flush`"),
"local buffer validation must not mask close poison: {}",
err.msg()
);
assert!(
!buf.is_empty(),
"terminal async error must not clear a newly prepared buffer"
);
buf.at_now().unwrap();
let err = sender.flush(&mut buf).unwrap_err();
assert!(
err.msg().contains("closed 1 times without ACK progress"),
"expected close poison in message, got: {}",
err.msg()
);
assert!(
!buf.is_empty(),
"terminal async error must not clear a newly prepared buffer"
);
assert!(sender.must_close());
}
#[test]
fn qwp_ws_orderly_close_reconnects_without_poison_strike() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (accept_tx, accept_rx) = mpsc::channel();
thread::spawn(move || {
for round in 0..2 {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream(&mut stream);
let (_fin, opcode, payload) = read_frame(&mut stream).unwrap();
assert_eq!(opcode, 0x2);
accept_tx.send((round, payload)).unwrap();
if round == 0 {
write_server_close_frame(&mut stream, 1000, "role change").unwrap();
} else {
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
}
thread::sleep(Duration::from_millis(50));
}
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_frame_rejections(1)
.unwrap()
.poison_min_escalation_window(Duration::ZERO)
.unwrap()
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
assert_eq!(fsn, 0);
let first = accept_rx.recv_timeout(Duration::from_secs(5)).unwrap();
let second = accept_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(first.0, 0);
assert_eq!(second.0, 1);
assert_eq!(&first.1[0..4], b"QWP1");
assert_eq!(first.1, second.1);
for _ in 0..100 {
if let Some(error) = sender.poll_qwp_ws_error().unwrap() {
assert_ne!(error.category, QwpWsErrorCategory::ProtocolViolation);
}
thread::sleep(Duration::from_millis(10));
}
}
fn run_paced_recycle_scenario(action: RecycleServerAction) -> usize {
let run_for = Duration::from_millis(900);
let (port, handle, connection_count) = spawn_recycling_server(action, run_for);
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_frame_rejections(1000)
.unwrap()
.poison_min_escalation_window(Duration::from_secs(60))
.unwrap()
.reconnect_initial_backoff(Duration::from_millis(150))
.unwrap()
.reconnect_max_backoff(Duration::from_secs(1))
.unwrap()
.reconnect_max_duration(Duration::from_secs(5))
.unwrap()
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
thread::sleep(run_for);
let wait_deadline = Instant::now() + Duration::from_secs(10);
while connection_count.load(Ordering::Acquire) < 2 && Instant::now() < wait_deadline {
thread::sleep(Duration::from_millis(10));
}
drop(sender);
handle.join().unwrap()
}
#[test]
fn qwp_ws_nack_recycles_are_paced_against_healthy_server() {
let connections = run_paced_recycle_scenario(RecycleServerAction::WriteError);
assert!(
(2..=8).contains(&connections),
"expected paced NACK recycles to make a small number of connections, got {connections}"
);
}
#[test]
fn qwp_ws_non_orderly_close_recycles_are_paced() {
let connections = run_paced_recycle_scenario(RecycleServerAction::NonOrderlyClose);
assert!(
(2..=8).contains(&connections),
"expected paced close recycles to make a small number of connections, got {connections}"
);
}
fn assert_server_protocol_violation<F>(write_bad_response: F, expected_message: &'static str)
where
F: FnOnce(&mut TcpStream) -> std::io::Result<()> + Send + 'static,
{
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frame_tx, frame_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream(&mut stream);
let (_fin, _opcode, payload) = read_frame(&mut stream).unwrap();
frame_tx.send(payload).unwrap();
write_bad_response(&mut stream).unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
let fsn = sender.flush_and_get_fsn(&mut buf).unwrap().unwrap();
let frame = frame_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
let mut observed = None;
for _ in 0..500 {
if let Some(error) = sender.poll_qwp_ws_error().unwrap() {
observed = Some(error);
break;
}
thread::sleep(Duration::from_millis(10));
}
let protocol_error = observed.expect("expected protocol violation diagnostic");
assert_eq!(
protocol_error.category,
QwpWsErrorCategory::ProtocolViolation
);
assert_eq!(protocol_error.applied_policy, QwpWsErrorPolicy::Terminal);
assert_eq!(protocol_error.status, None);
assert_eq!(protocol_error.message_sequence, None);
assert_eq!(protocol_error.message.as_deref(), Some(expected_message));
assert_eq!(protocol_error.from_fsn, fsn);
assert_eq!(protocol_error.to_fsn, fsn);
let err = sender.flush(&mut buf).unwrap_err();
assert!(
err.msg().contains(expected_message),
"expected protocol violation in terminal message, got: {}",
err.msg()
);
assert!(sender.must_close());
}
#[test]
fn qwp_ws_masked_server_frame_is_pollable_as_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_frame(stream, 0x2, b"masked", true),
"WebSocket server frame must not be masked",
);
}
#[test]
fn qwp_ws_unknown_opcode_is_pollable_as_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_frame(stream, 0x0B, b"", false),
"Unknown WebSocket opcode: 0xb",
);
}
#[test]
fn qwp_ws_text_response_is_pollable_as_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_frame(stream, 0x1, b"not-qwp", false),
"QWP/WebSocket server response was not a binary frame",
);
}
#[test]
fn qwp_ws_rsv_bits_set_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_raw_ws_frame(stream, 0xC2, b"hello"),
"WebSocket RSV bits are set but no extensions were negotiated",
);
}
#[test]
fn qwp_ws_continuation_without_prior_data_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_raw_ws_frame(stream, 0x80, b""),
"Continuation frame without prior data frame",
);
}
#[test]
fn qwp_ws_new_data_frame_mid_fragment_is_protocol_violation() {
assert_server_protocol_violation(
|stream| {
write_raw_ws_frame(stream, 0x02, b"frag1")?;
write_raw_ws_frame(stream, 0x82, b"frag2")
},
"Unexpected new data frame mid-message",
);
}
#[test]
fn qwp_ws_fragmented_control_frame_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_raw_ws_frame(stream, 0x08, &1000u16.to_be_bytes()),
"WebSocket control frame must not be fragmented",
);
}
#[test]
fn qwp_ws_close_frame_with_one_byte_payload_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_raw_ws_frame(stream, 0x88, &[0x00]),
"WebSocket close frame payload length must be 0 or at least 2 bytes",
);
}
#[test]
fn qwp_ws_close_code_1004_reserved_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_close_frame(stream, 1004, ""),
"WebSocket close frame uses reserved or out-of-range close code: 1004",
);
}
#[test]
fn qwp_ws_close_code_1005_sentinel_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_close_frame(stream, 1005, ""),
"WebSocket close frame uses reserved or out-of-range close code: 1005",
);
}
#[test]
fn qwp_ws_close_code_1006_sentinel_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_close_frame(stream, 1006, ""),
"WebSocket close frame uses reserved or out-of-range close code: 1006",
);
}
#[test]
fn qwp_ws_close_code_1015_sentinel_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_close_frame(stream, 1015, ""),
"WebSocket close frame uses reserved or out-of-range close code: 1015",
);
}
#[test]
fn qwp_ws_close_code_below_1000_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_close_frame(stream, 999, ""),
"WebSocket close frame uses reserved or out-of-range close code: 999",
);
}
#[test]
fn qwp_ws_close_code_reserved_future_range_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_close_frame(stream, 1016, ""),
"WebSocket close frame uses reserved or out-of-range close code: 1016",
);
}
#[test]
fn qwp_ws_close_code_above_4999_is_protocol_violation() {
assert_server_protocol_violation(
|stream| write_server_close_frame(stream, 5000, ""),
"WebSocket close frame uses reserved or out-of-range close code: 5000",
);
}
#[test]
fn qwp_ws_close_frame_with_invalid_utf8_reason_is_protocol_violation() {
assert_server_protocol_violation(
|stream| {
let mut payload = Vec::new();
payload.extend_from_slice(&1000u16.to_be_bytes());
payload.extend_from_slice(&[0xC3, 0x28]);
write_raw_ws_frame(stream, 0x88, &payload)
},
"WebSocket close frame reason is not valid UTF-8",
);
}
#[test]
fn qwp_ws_high_level_flush_returns_before_ack() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frame_tx, frame_rx) = mpsc::channel();
let (ack_tx, ack_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream(&mut stream);
let (_fin, _opcode, payload) = read_frame(&mut stream).unwrap();
frame_tx.send(payload).unwrap();
ack_rx.recv_timeout(Duration::from_secs(5)).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
assert!(buf.is_empty());
let frame = frame_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
ack_tx.send(()).unwrap();
}
#[test]
fn qwp_ws_high_level_flush_and_keep_returns_before_ack_and_preserves_buffer() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frame_tx, frame_rx) = mpsc::channel();
let (ack_tx, ack_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream(&mut stream);
let (_fin, _opcode, payload) = read_frame(&mut stream).unwrap();
frame_tx.send(payload).unwrap();
ack_rx.recv_timeout(Duration::from_secs(5)).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
sender.flush_and_keep(&buf).unwrap();
assert!(!buf.is_empty());
let frame = frame_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
ack_tx.send(()).unwrap();
}
#[test]
fn qwp_ws_high_level_flushes_pipeline_before_ack() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (frames_tx, frames_rx) = mpsc::channel();
let (ack_tx, ack_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream(&mut stream);
let mut frames = Vec::new();
let (_fin, _opcode, first) = read_frame(&mut stream).unwrap();
frames.push(first);
let (_fin, _opcode, second) = read_frame(&mut stream).unwrap();
frames.push(second);
frames_tx.send(frames).unwrap();
ack_rx.recv_timeout(Duration::from_secs(5)).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE + 1).unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.column_i64("qty", 1)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
buf.table("trades")
.unwrap()
.column_i64("qty", 2)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frames = frames_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(frames.len(), 2);
assert!(frames.iter().all(|frame| &frame[0..4] == b"QWP1"));
ack_tx.send(()).unwrap();
}
fn spawn_dropping_then_recovering_server() -> (u16, std::sync::mpsc::Receiver<Vec<u8>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = std::sync::mpsc::channel();
thread::spawn(move || {
let do_upgrade = |stream: &mut TcpStream| {
let req_bytes = read_request_until_blank(stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let resp = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(resp.as_bytes()).unwrap();
};
let (mut s1, _) = listener.accept().unwrap();
s1.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
s1.set_write_timeout(Some(Duration::from_secs(5))).unwrap();
do_upgrade(&mut s1);
let (_fin, _op, payload) = read_frame(&mut s1).unwrap();
tx.send(payload).unwrap();
drop(s1);
let (mut s2, _) = listener.accept().unwrap();
s2.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
s2.set_write_timeout(Some(Duration::from_secs(5))).unwrap();
do_upgrade(&mut s2);
let (_fin, _op, _catch_up) = read_frame(&mut s2).unwrap();
let (_fin, _op, payload) = read_frame(&mut s2).unwrap();
tx.send(payload).unwrap();
let mut ok = vec![0u8];
ok.extend_from_slice(&(FIRST_WIRE_SEQUENCE + 1).to_le_bytes());
ok.extend_from_slice(&0u16.to_le_bytes());
write_server_binary_frame(&mut s2, &ok).unwrap();
thread::sleep(Duration::from_millis(50));
});
(port, rx)
}
#[test]
fn qwp_ws_reconnects_and_replays_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let (port, rx) = spawn_dropping_then_recovering_server();
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.reconnect_initial_backoff(Duration::from_millis(20))
.unwrap()
.reconnect_max_backoff(Duration::from_millis(50))
.unwrap();
let mut sender = build_qwp_ws_sender_from_builder(progress, builder);
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush_and_get_fsn(&mut buf).unwrap();
sender
.wait(crate::ingress::AckLevel::Ok, Duration::from_secs(5))
.unwrap_or_else(|e| panic!("mode={}: {e}", progress.name()));
let frame1 = rx.recv_timeout(Duration::from_secs(5)).unwrap();
let frame2 = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame1[0..4], b"QWP1");
assert_eq!(frame1, frame2, "mode={}", progress.name());
let totals = sender.qwp_ws_totals().unwrap();
assert_eq!(totals.frames_sent, 2, "mode={}", progress.name());
assert_eq!(totals.frames_replayed, 1, "mode={}", progress.name());
}
}
fn catch_up_symbols(frame: &[u8]) -> (u64, Vec<Vec<u8>>) {
assert_eq!(&frame[0..4], b"QWP1", "catch-up frame magic");
assert_eq!(frame[5] & 0x08, 0x08, "catch-up carries the delta flag");
assert_eq!(
u16::from_le_bytes([frame[6], frame[7]]),
0,
"catch-up frame is table-less"
);
let mut pos = 12usize;
let delta_start = read_varint(frame, &mut pos);
let count = read_varint(frame, &mut pos);
let mut symbols = Vec::new();
for _ in 0..count {
let len = read_varint(frame, &mut pos) as usize;
symbols.push(frame[pos..pos + len].to_vec());
pos += len;
}
(delta_start, symbols)
}
fn spawn_reconnect_forwarding_catch_up_server() -> (u16, mpsc::Receiver<Vec<u8>>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut s1, _) = listener.accept().unwrap();
s1.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
perform_server_upgrade(&mut s1).unwrap();
let _ = read_frame(&mut s1);
drop(s1);
let (mut s2, _) = listener.accept().unwrap();
s2.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
s2.set_write_timeout(Some(Duration::from_secs(5))).unwrap();
perform_server_upgrade(&mut s2).unwrap();
let (_fin, _op, catch_up) = read_frame(&mut s2).unwrap();
tx.send(catch_up).unwrap();
let (_fin, _op, _data) = read_frame(&mut s2).unwrap();
write_qwp_ok_response(&mut s2, FIRST_WIRE_SEQUENCE + 1).unwrap();
thread::sleep(Duration::from_millis(50));
});
(port, rx)
}
#[test]
fn qwp_ws_reconnect_catch_up_re_registers_the_real_multi_symbol_dictionary() {
let (port, rx) = spawn_reconnect_forwarding_catch_up_server();
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.reconnect_initial_backoff(Duration::from_millis(20))
.unwrap()
.reconnect_max_backoff(Duration::from_millis(50))
.unwrap();
let mut sender = build_qwp_ws_sender_from_builder(ProgressCase::Background, builder);
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
buf.table("trades")
.unwrap()
.symbol("sym", "BTC-USD")
.unwrap()
.column_i64("qty", 8)
.unwrap()
.at_now()
.unwrap();
sender.flush_and_get_fsn(&mut buf).unwrap();
sender
.wait(crate::ingress::AckLevel::Ok, Duration::from_secs(5))
.unwrap();
let catch_up = rx.recv_timeout(Duration::from_secs(5)).unwrap();
let (delta_start, symbols) = catch_up_symbols(&catch_up);
assert_eq!(
delta_start, 0,
"catch-up re-registers the dictionary from id 0"
);
assert_eq!(
symbols,
vec![b"ETH-USD".to_vec(), b"BTC-USD".to_vec()],
"catch-up re-registers the real encoder-produced dictionary in id order"
);
}
#[test]
fn qwp_ws_midstream_failure_reconnects_to_next_endpoint() {
let first_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let first_port = first_listener.local_addr().unwrap().port();
let second_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let second_port = second_listener.local_addr().unwrap().port();
let (tx, rx) = std::sync::mpsc::channel();
let first_tx = tx.clone();
thread::spawn(move || {
let (mut stream, _) = first_listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
first_tx.send(payload).unwrap();
drop(stream);
});
thread::spawn(move || {
let (mut stream, _) = second_listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, _catch_up) = read_frame(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE + 1).unwrap();
thread::sleep(Duration::from_millis(50));
});
let conf = format!(
"ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};\
reconnect_initial_backoff_millis=1;\
reconnect_max_backoff_millis=1;\
reconnect_max_duration_millis=5000;"
);
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frame1 = rx.recv_timeout(Duration::from_secs(5)).unwrap();
let frame2 = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame1[0..4], b"QWP1");
assert_eq!(frame1, frame2);
}
#[test]
fn qwp_ws_midstream_failure_version_error_on_other_endpoint_stays_retryable() {
let first_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let first_port = first_listener.local_addr().unwrap().port();
let done = Arc::new(AtomicBool::new(false));
let (second_port, second_handle) = spawn_no_durable_ack_upgrade_server(Arc::clone(&done));
let (frame_tx, frame_rx) = mpsc::channel();
let first_thread = thread::spawn(move || {
let (mut stream, _) = first_listener.accept().unwrap();
perform_server_upgrade_durable(&mut stream).unwrap();
let (_fin, _op, _payload) = read_frame(&mut stream).unwrap();
drop(stream);
let (mut stream, _) = first_listener.accept().unwrap();
perform_server_upgrade_durable(&mut stream).unwrap();
let (_fin, _op, _catch_up) = read_frame(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
write_qwp_ok_response_with_table_entries(
&mut stream,
FIRST_WIRE_SEQUENCE + 1,
&[("trades", 10)],
)
.unwrap();
write_qwp_durable_ack_response(&mut stream, &[("trades", 10)]).unwrap();
frame_tx.send(payload).unwrap();
thread::sleep(Duration::from_millis(50));
});
let conf = format!(
"ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};\
request_durable_ack=on;\
reconnect_initial_backoff_millis=1;\
reconnect_max_backoff_millis=1;\
reconnect_max_duration_millis=5000;"
);
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender
.flush(&mut buf)
.expect("mid-stream reconnect must survive a version error on the other endpoint");
let replayed = frame_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&replayed[0..4], b"QWP1");
done.store(true, Ordering::Release);
first_thread.join().unwrap();
assert!(
second_handle.join().unwrap() >= 1,
"the reconnect round should have dialed the version-erroring endpoint"
);
}
#[test]
fn qwp_ws_sync_reconnect_retries_failed_attempt() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (payload_tx, payload_rx) = mpsc::channel();
let (event_tx, event_rx) = mpsc::channel();
thread::spawn(move || {
let do_upgrade = |stream: &mut TcpStream| {
let req_bytes = read_request_until_blank(stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let resp = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
stream.write_all(resp.as_bytes()).unwrap();
};
let (mut s1, _) = listener.accept().unwrap();
s1.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
s1.set_write_timeout(Some(Duration::from_secs(5))).unwrap();
do_upgrade(&mut s1);
let (_fin, _op, payload) = read_frame(&mut s1).unwrap();
payload_tx.send(payload).unwrap();
drop(s1);
let (s2, _) = listener.accept().unwrap();
event_tx.send("failed_reconnect").unwrap();
drop(s2);
let (mut s3, _) = listener.accept().unwrap();
s3.set_read_timeout(Some(Duration::from_secs(5))).unwrap();
s3.set_write_timeout(Some(Duration::from_secs(5))).unwrap();
do_upgrade(&mut s3);
let (_fin, _op, _catch_up) = read_frame(&mut s3).unwrap();
let (_fin, _op, payload) = read_frame(&mut s3).unwrap();
payload_tx.send(payload).unwrap();
let mut ok = vec![0u8];
ok.extend_from_slice(&(FIRST_WIRE_SEQUENCE + 1).to_le_bytes());
ok.extend_from_slice(&0u16.to_le_bytes());
write_server_binary_frame(&mut s3, &ok).unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.reconnect_initial_backoff(Duration::from_millis(1))
.unwrap()
.reconnect_max_backoff(Duration::from_millis(1))
.unwrap()
.reconnect_max_duration(Duration::from_secs(5))
.unwrap()
.build()
.unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
assert_eq!(
event_rx.recv_timeout(Duration::from_secs(5)).unwrap(),
"failed_reconnect"
);
let frame1 = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
let frame2 = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame1[0..4], b"QWP1");
assert_eq!(frame1, frame2);
}
#[test]
fn qwp_ws_sync_initial_connect_retry_survives_dropped_upgrade() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (payload_tx, payload_rx) = mpsc::channel();
let (event_tx, event_rx) = mpsc::channel();
thread::spawn(move || {
let (mut first, _) = listener.accept().unwrap();
first
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
first
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = read_request_until_blank(&mut first).unwrap();
event_tx.send("dropped_initial_upgrade").unwrap();
drop(first);
let (mut second, _) = listener.accept().unwrap();
second
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
second
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let req_bytes = read_request_until_blank(&mut second).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let key = parse_header(&req, "Sec-WebSocket-Key").unwrap();
let accept = compute_accept(&key);
let resp = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
\r\n"
);
second.write_all(resp.as_bytes()).unwrap();
let (_fin, _op, payload) = read_frame(&mut second).unwrap();
payload_tx.send(payload).unwrap();
let mut ok = vec![0u8];
ok.extend_from_slice(&FIRST_WIRE_SEQUENCE.to_le_bytes());
ok.extend_from_slice(&0u16.to_le_bytes());
write_server_binary_frame(&mut second, &ok).unwrap();
thread::sleep(Duration::from_millis(50));
});
let mut sender = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.initial_connect_retry(true)
.unwrap()
.reconnect_initial_backoff(Duration::from_millis(1))
.unwrap()
.reconnect_max_backoff(Duration::from_millis(1))
.unwrap()
.reconnect_max_duration(Duration::from_secs(5))
.unwrap()
.build()
.unwrap();
assert_eq!(
event_rx.recv_timeout(Duration::from_secs(5)).unwrap(),
"dropped_initial_upgrade"
);
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
}
#[test]
fn qwp_ws_initial_connect_walks_endpoint_list_in_off_mode() {
let bad_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let bad_port = bad_listener.local_addr().unwrap().port();
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let good_port = listener.local_addr().unwrap().port();
drop(bad_listener);
let (payload_tx, payload_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
payload_tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
thread::sleep(Duration::from_millis(50));
});
let conf = format!("ws::addr=127.0.0.1:{bad_port},127.0.0.1:{good_port};");
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
}
#[test]
fn qwp_ws_initial_connect_role_reject_tries_next_endpoint() {
let first_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let first_port = first_listener.local_addr().unwrap().port();
let second_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let second_port = second_listener.local_addr().unwrap().port();
let (payload_tx, payload_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = first_listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = read_request_until_blank(&mut stream).unwrap();
stream
.write_all(
b"HTTP/1.1 421 Misdirected Request\r\n\
Connection: close\r\n\
Content-Length: 0\r\n\
X-QuestDB-Role: PRIMARY_CATCHUP\r\n\
\r\n",
)
.unwrap();
});
thread::spawn(move || {
let (mut stream, _) = second_listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
payload_tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
thread::sleep(Duration::from_millis(50));
});
let conf = format!("ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};");
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
}
#[test]
fn qwp_ws_initial_connect_retryable_status_tries_next_endpoint() {
let first_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let first_port = first_listener.local_addr().unwrap().port();
let second_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let second_port = second_listener.local_addr().unwrap().port();
let (payload_tx, payload_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = first_listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = read_request_until_blank(&mut stream).unwrap();
stream
.write_all(
b"HTTP/1.1 500 Internal Server Error\r\n\
Connection: close\r\n\
Content-Length: 0\r\n\
\r\n",
)
.unwrap();
});
thread::spawn(move || {
let (mut stream, _) = second_listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
payload_tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
thread::sleep(Duration::from_millis(50));
});
let conf = format!("ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};");
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
}
#[test]
fn qwp_ws_initial_connect_mixed_role_and_transport_prefers_role_mismatch() {
let (role_port, role_handle) = spawn_role_reject_upgrade_server(1, "PRIMARY_CATCHUP");
let status_listener = TcpListener::bind("127.0.0.1:0").unwrap();
status_listener.set_nonblocking(true).unwrap();
let status_port = status_listener.local_addr().unwrap().port();
let status_handle = thread::spawn(move || {
let started = Instant::now();
let mut attempts = 0usize;
while attempts < 1 && started.elapsed() < Duration::from_secs(5) {
match status_listener.accept() {
Ok((mut stream, _)) => {
attempts += 1;
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = read_request_until_blank(&mut stream).unwrap();
stream
.write_all(
b"HTTP/1.1 500 Internal Server Error\r\n\
Connection: close\r\n\
Content-Length: 0\r\n\
\r\n",
)
.unwrap();
}
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(5));
}
Err(err) => panic!("status listener failed: {err}"),
}
}
attempts
});
let conf = format!(
"ws::addr=127.0.0.1:{role_port},127.0.0.1:{status_port};\
initial_connect_retry=off;"
);
let result = SenderBuilder::from_conf(&conf).unwrap().build();
assert_eq!(role_handle.join().unwrap(), 1);
assert_eq!(status_handle.join().unwrap(), 1);
let err = result.unwrap_err();
assert_eq!(err.code(), ErrorCode::RoleMismatch);
assert!(err.msg().contains("role mismatch"), "got: {}", err.msg());
assert!(err.msg().contains("PRIMARY_CATCHUP"), "got: {}", err.msg());
assert!(!err.msg().contains("HTTP status 500"), "got: {}", err.msg());
}
#[test]
fn qwp_ws_initial_connect_unsupported_version_tries_next_endpoint() {
let first_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let first_port = first_listener.local_addr().unwrap().port();
let second_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let second_port = second_listener.local_addr().unwrap().port();
let (payload_tx, payload_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = first_listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = upgrade_mock_stream_with_version(&mut stream, 2);
});
thread::spawn(move || {
let (mut stream, _) = second_listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
payload_tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
thread::sleep(Duration::from_millis(50));
});
let conf = format!("ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};");
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
}
#[test]
fn qwp_ws_sync_initial_retry_unsupported_version_retries_next_round() {
let first_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let first_port = first_listener.local_addr().unwrap().port();
let second_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let second_port = second_listener.local_addr().unwrap().port();
let first_attempts = Arc::new(AtomicUsize::new(0));
let second_attempts = Arc::new(AtomicUsize::new(0));
let (payload_tx, payload_rx) = mpsc::channel();
let (done_tx, done_rx) = mpsc::channel();
let first_handle = {
let first_attempts = Arc::clone(&first_attempts);
thread::spawn(move || {
let (mut stream, _) = first_listener.accept().unwrap();
first_attempts.fetch_add(1, Ordering::AcqRel);
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = upgrade_mock_stream_with_version(&mut stream, 2);
let (mut stream, _) = first_listener.accept().unwrap();
first_attempts.fetch_add(1, Ordering::AcqRel);
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
payload_tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
let _ = done_rx.recv_timeout(Duration::from_secs(5));
})
};
let second_handle = {
let second_attempts = Arc::clone(&second_attempts);
thread::spawn(move || {
let (mut stream, _) = second_listener.accept().unwrap();
second_attempts.fetch_add(1, Ordering::AcqRel);
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = upgrade_mock_stream_with_version(&mut stream, 2);
})
};
let conf = format!(
"ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};\
initial_connect_retry=sync;\
reconnect_initial_backoff_millis=1;\
reconnect_max_backoff_millis=1;\
reconnect_max_duration_millis=2000;"
);
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
done_tx.send(()).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
first_handle.join().unwrap();
second_handle.join().unwrap();
assert_eq!(&frame[0..4], b"QWP1");
assert_eq!(first_attempts.load(Ordering::Acquire), 2);
assert_eq!(second_attempts.load(Ordering::Acquire), 1);
}
#[test]
fn qwp_ws_sync_initial_retry_malformed_101_retries_after_round_exhaustion() {
let first_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let first_port = first_listener.local_addr().unwrap().port();
let second_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let second_port = second_listener.local_addr().unwrap().port();
let first_attempts = Arc::new(AtomicUsize::new(0));
let second_attempts = Arc::new(AtomicUsize::new(0));
let (payload_tx, payload_rx) = mpsc::channel();
let (done_tx, done_rx) = mpsc::channel();
let first_handle = {
let first_attempts = Arc::clone(&first_attempts);
thread::spawn(move || {
let (mut stream, _) = first_listener.accept().unwrap();
first_attempts.fetch_add(1, Ordering::AcqRel);
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream_without_upgrade_header(&mut stream);
let (mut stream, _) = first_listener.accept().unwrap();
first_attempts.fetch_add(1, Ordering::AcqRel);
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
payload_tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
let _ = done_rx.recv_timeout(Duration::from_secs(5));
})
};
let second_handle = {
let second_attempts = Arc::clone(&second_attempts);
thread::spawn(move || {
let (mut stream, _) = second_listener.accept().unwrap();
second_attempts.fetch_add(1, Ordering::AcqRel);
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
upgrade_mock_stream_without_upgrade_header(&mut stream);
})
};
let conf = format!(
"ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};\
initial_connect_retry=sync;\
reconnect_initial_backoff_millis=1;\
reconnect_max_backoff_millis=1;\
reconnect_max_duration_millis=2000;"
);
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
done_tx.send(()).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
first_handle.join().unwrap();
second_handle.join().unwrap();
assert_eq!(&frame[0..4], b"QWP1");
assert_eq!(first_attempts.load(Ordering::Acquire), 2);
assert_eq!(second_attempts.load(Ordering::Acquire), 1);
}
#[test]
fn qwp_ws_sync_initial_retry_mixed_role_and_transport_prefers_role_mismatch() {
let (role_port, role_handle) = spawn_role_reject_upgrade_server(1, "PRIMARY_CATCHUP");
let status_listener = TcpListener::bind("127.0.0.1:0").unwrap();
status_listener.set_nonblocking(true).unwrap();
let status_port = status_listener.local_addr().unwrap().port();
let status_handle = thread::spawn(move || {
let started = Instant::now();
let mut attempts = 0usize;
while attempts < 1 && started.elapsed() < Duration::from_secs(5) {
match status_listener.accept() {
Ok((mut stream, _)) => {
attempts += 1;
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = read_request_until_blank(&mut stream).unwrap();
stream
.write_all(
b"HTTP/1.1 500 Internal Server Error\r\n\
Connection: close\r\n\
Content-Length: 0\r\n\
\r\n",
)
.unwrap();
}
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(5));
}
Err(err) => panic!("status listener failed: {err}"),
}
}
attempts
});
let conf = format!(
"ws::addr=127.0.0.1:{role_port},127.0.0.1:{status_port};\
initial_connect_retry=sync;\
reconnect_initial_backoff_millis=10;\
reconnect_max_backoff_millis=10;\
reconnect_max_duration_millis=1;"
);
let result = SenderBuilder::from_conf(&conf).unwrap().build();
assert_eq!(role_handle.join().unwrap(), 1);
assert_eq!(status_handle.join().unwrap(), 1);
let err = result.unwrap_err();
assert_eq!(err.code(), ErrorCode::RoleMismatch);
assert!(
err.msg()
.contains("QWP/WebSocket initial connect retry budget exhausted"),
"got: {}",
err.msg()
);
assert!(err.msg().contains("role mismatch"), "got: {}", err.msg());
assert!(err.msg().contains("PRIMARY_CATCHUP"), "got: {}", err.msg());
assert!(!err.msg().contains("HTTP status 500"), "got: {}", err.msg());
}
#[test]
fn qwp_ws_initial_connect_durable_ack_mismatch_tries_next_endpoint() {
let first_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let first_port = first_listener.local_addr().unwrap().port();
let second_listener = TcpListener::bind("127.0.0.1:0").unwrap();
let second_port = second_listener.local_addr().unwrap().port();
let (payload_tx, payload_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = first_listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = upgrade_mock_stream_with_durable_ack(&mut stream, false);
});
thread::spawn(move || {
let (mut stream, _) = second_listener.accept().unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = upgrade_mock_stream_with_durable_ack(&mut stream, true);
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
payload_tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
write_qwp_durable_ack_response(&mut stream, &[]).unwrap();
thread::sleep(Duration::from_millis(50));
});
let conf = format!(
"ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};\
request_durable_ack=on;"
);
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
}
#[test]
fn qwp_ws_sync_initial_retry_durable_ack_all_role_rejects_retry_until_budget() {
let (first_port, first_handle) = spawn_role_reject_upgrade_server(1, "REPLICA");
let (second_port, second_handle) = spawn_role_reject_upgrade_server(1, "REPLICA");
let conf = format!(
"ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};\
request_durable_ack=on;\
initial_connect_retry=sync;\
reconnect_initial_backoff_millis=100;\
reconnect_max_backoff_millis=100;\
reconnect_max_duration_millis=500;"
);
let err = SenderBuilder::from_conf(&conf)
.unwrap()
.build()
.unwrap_err();
assert_eq!(err.code(), ErrorCode::SocketError);
assert!(
err.msg()
.contains("QWP/WebSocket initial connect retry budget exhausted"),
"got: {}",
err.msg()
);
assert_eq!(first_handle.join().unwrap(), 1);
assert_eq!(second_handle.join().unwrap(), 1);
}
fn spawn_accept_then_drop_server(done: Arc<AtomicBool>) -> (u16, thread::JoinHandle<usize>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
listener.set_nonblocking(true).unwrap();
let port = listener.local_addr().unwrap().port();
let handle = thread::spawn(move || {
let mut attempts = 0usize;
while !done.load(Ordering::Acquire) {
match listener.accept() {
Ok((stream, _)) => {
attempts += 1;
drop(stream);
}
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(5));
}
Err(err) => panic!("accept-then-drop listener failed: {err}"),
}
}
attempts
});
(port, handle)
}
#[test]
fn qwp_ws_sync_initial_retry_version_error_does_not_mask_transient_endpoint_failure() {
let done = Arc::new(AtomicBool::new(false));
let (first_port, first_handle) = spawn_no_durable_ack_upgrade_server(Arc::clone(&done));
let (second_port, second_handle) = spawn_accept_then_drop_server(Arc::clone(&done));
let conf = format!(
"ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};\
request_durable_ack=on;\
initial_connect_retry=sync;\
reconnect_initial_backoff_millis=50;\
reconnect_max_backoff_millis=50;\
reconnect_max_duration_millis=300;"
);
let err = SenderBuilder::from_conf(&conf)
.unwrap()
.build()
.unwrap_err();
done.store(true, Ordering::Release);
assert_eq!(err.code(), ErrorCode::SocketError, "got: {}", err.msg());
assert!(
err.msg()
.contains("QWP/WebSocket initial connect retry budget exhausted"),
"got: {}",
err.msg()
);
let first_attempts = first_handle.join().unwrap();
let _ = second_handle.join().unwrap();
assert!(
first_attempts >= 2,
"expected repeated retries against the version-erroring endpoint, got {first_attempts}"
);
}
#[test]
fn qwp_ws_sync_initial_retry_budget_exhaustion_reports_context() {
let probe = TcpListener::bind("127.0.0.1:0").unwrap();
let port = probe.local_addr().unwrap().port();
drop(probe);
let conf = format!(
"ws::addr=127.0.0.1:{port};\
initial_connect_retry=sync;\
reconnect_initial_backoff_millis=10;\
reconnect_max_backoff_millis=10;\
reconnect_max_duration_millis=1;"
);
let err = SenderBuilder::from_conf(&conf)
.unwrap()
.build()
.unwrap_err();
assert_eq!(err.code(), ErrorCode::SocketError);
assert!(
err.msg()
.contains("QWP/WebSocket initial connect retry budget exhausted"),
"got: {}",
err.msg()
);
assert!(err.msg().contains("attempts=1"), "got: {}", err.msg());
assert!(err.msg().contains("elapsed_ms="), "got: {}", err.msg());
assert!(
err.msg()
.contains("last_error=QWP/WebSocket all endpoints unreachable"),
"got: {}",
err.msg()
);
}
#[test]
fn qwp_ws_sync_initial_retry_resets_non_healthy_between_rounds() {
let first_listener = TcpListener::bind("127.0.0.1:0").unwrap();
first_listener.set_nonblocking(true).unwrap();
let first_port = first_listener.local_addr().unwrap().port();
let second_listener = TcpListener::bind("127.0.0.1:0").unwrap();
second_listener.set_nonblocking(true).unwrap();
let second_port = second_listener.local_addr().unwrap().port();
let first_attempts = Arc::new(AtomicUsize::new(0));
let second_attempts = Arc::new(AtomicUsize::new(0));
let done = Arc::new(AtomicBool::new(false));
let (payload_tx, payload_rx) = mpsc::channel();
let first_attempts_thread = Arc::clone(&first_attempts);
let done_first_thread = Arc::clone(&done);
let first_thread = thread::spawn(move || {
while !done_first_thread.load(Ordering::Acquire) {
match first_listener.accept() {
Ok((mut stream, _)) => {
let attempt = first_attempts_thread.fetch_add(1, Ordering::AcqRel) + 1;
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
if attempt == 1 {
let _ = read_request_until_blank(&mut stream).unwrap();
stream
.write_all(
b"HTTP/1.1 421 Misdirected Request\r\n\
Connection: close\r\n\
Content-Length: 0\r\n\
X-QuestDB-Role: REPLICA\r\n\
\r\n",
)
.unwrap();
} else {
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
payload_tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
thread::sleep(Duration::from_millis(50));
done_first_thread.store(true, Ordering::Release);
}
}
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(5));
}
Err(err) => panic!("first listener failed: {err}"),
}
}
});
let second_attempts_thread = Arc::clone(&second_attempts);
let done_second_thread = Arc::clone(&done);
let second_thread = thread::spawn(move || {
while !done_second_thread.load(Ordering::Acquire) {
match second_listener.accept() {
Ok((mut stream, _)) => {
second_attempts_thread.fetch_add(1, Ordering::AcqRel);
stream.set_nonblocking(false).unwrap();
stream
.set_read_timeout(Some(Duration::from_secs(5)))
.unwrap();
stream
.set_write_timeout(Some(Duration::from_secs(5)))
.unwrap();
let _ = read_request_until_blank(&mut stream).unwrap();
stream
.write_all(
b"HTTP/1.1 500 Internal Server Error\r\n\
Connection: close\r\n\
Content-Length: 0\r\n\
\r\n",
)
.unwrap();
}
Err(err) if err.kind() == std::io::ErrorKind::WouldBlock => {
thread::sleep(Duration::from_millis(5));
}
Err(err) => panic!("second listener failed: {err}"),
}
}
});
let conf = format!(
"ws::addr=127.0.0.1:{first_port},127.0.0.1:{second_port};\
initial_connect_retry=sync;\
reconnect_initial_backoff_millis=100;\
reconnect_max_backoff_millis=100;\
reconnect_max_duration_millis=5000;"
);
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
done.store(true, Ordering::Release);
first_thread.join().unwrap();
second_thread.join().unwrap();
assert_eq!(first_attempts.load(Ordering::Acquire), 2);
assert_eq!(second_attempts.load(Ordering::Acquire), 1);
}
#[test]
fn qwp_ws_async_initial_connect_background_can_flush() {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (payload_tx, payload_rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
perform_server_upgrade(&mut stream).unwrap();
let (_fin, _op, payload) = read_frame(&mut stream).unwrap();
payload_tx.send(payload).unwrap();
write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE).unwrap();
thread::sleep(Duration::from_millis(50));
});
let conf = format!("ws::addr=127.0.0.1:{port};initial_connect_retry=async;");
let mut sender = SenderBuilder::from_conf(&conf).unwrap().build().unwrap();
let mut buf = sender.new_buffer();
buf.table("trades")
.unwrap()
.symbol("sym", "ETH-USD")
.unwrap()
.column_i64("qty", 7)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut buf).unwrap();
let frame = payload_rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(&frame[0..4], b"QWP1");
}
#[test]
fn qwp_ws_async_initial_connect_rejects_manual_progress() {
let conf = "ws::addr=127.0.0.1:1;\
initial_connect_retry=async;\
qwp_ws_progress=manual;";
let err = SenderBuilder::from_conf(conf).unwrap().build().unwrap_err();
assert!(
err.msg().contains("initial_connect_retry=async")
&& err.msg().contains("background progress"),
"got: {}",
err.msg()
);
}
#[test]
fn qwp_ws_from_conf_parses_java_reconnect_keys() {
let conf = "ws::addr=localhost:9000;\
reconnect_max_duration_millis=20000;\
reconnect_initial_backoff_millis=200;\
reconnect_max_backoff_millis=2000;\
initial_connect_retry=on;\
close_flush_timeout_millis=120000;\
request_durable_ack=on;\
durable_ack_keepalive_interval_millis=250;";
SenderBuilder::from_conf(conf).unwrap();
let disabled_keepalive = "ws::addr=localhost:9000;durable_ack_keepalive_interval_millis=0;";
SenderBuilder::from_conf(disabled_keepalive).unwrap();
let conf_sync = "ws::addr=localhost:9000;initial_connect_retry=sync;";
SenderBuilder::from_conf(conf_sync).unwrap();
let conf_false = "ws::addr=localhost:9000;initial_connect_retry=false;";
SenderBuilder::from_conf(conf_false).unwrap();
let conf_multi = "ws::addr=localhost:9000, localhost:9001;addr=localhost:9002;";
SenderBuilder::from_conf(conf_multi).unwrap();
let zone_ignored = "ws::addr=localhost:9000;zone=dc-amsterdam;";
SenderBuilder::from_conf(zone_ignored).unwrap();
#[cfg(feature = "sync-sender-tcp")]
{
let tcp_zone = "tcp::addr=localhost:9009;zone=dc-amsterdam;";
SenderBuilder::from_conf(tcp_zone).unwrap();
}
let target_ignored = "ws::addr=localhost:9000;target=primary;";
SenderBuilder::from_conf(target_ignored).unwrap();
let duplicate = "ws::addr=localhost:9000,localhost:9000;";
let err = SenderBuilder::from_conf(duplicate).unwrap_err();
assert!(err.msg().contains("duplicate"), "got: {}", err.msg());
let duplicate_effective_port = "ws::addr=localhost:9000,localhost:09000;";
let err = SenderBuilder::from_conf(duplicate_effective_port).unwrap_err();
assert!(
err.msg().contains("duplicate") && err.msg().contains("localhost:9000"),
"got: {}",
err.msg()
);
let empty_entry = "ws::addr=localhost:9000,,localhost:9001;";
let err = SenderBuilder::from_conf(empty_entry).unwrap_err();
assert!(err.msg().contains("empty entry"), "got: {}", err.msg());
let empty_host = "ws::addr=:9000;";
let err = SenderBuilder::from_conf(empty_host).unwrap_err();
assert!(err.msg().contains("empty host"), "got: {}", err.msg());
let ipv6_multi = "ws::addr=[::1]:9000,[2001:db8::1]:9001, localhost:9002;";
SenderBuilder::from_conf(ipv6_multi).unwrap();
let unbracketed_ipv6 = "ws::addr=::1:9000;";
let err = SenderBuilder::from_conf(unbracketed_ipv6).unwrap_err();
assert!(err.msg().contains("bracket IPv6"), "got: {}", err.msg());
let unterminated_bracket = "ws::addr=[::1:9000;";
let err = SenderBuilder::from_conf(unterminated_bracket).unwrap_err();
assert!(err.msg().contains("missing ']'"), "got: {}", err.msg());
let bracket_junk = "ws::addr=[::1]9000;";
let err = SenderBuilder::from_conf(bracket_junk).unwrap_err();
assert!(
err.msg().contains("expected ':port' after ']'"),
"got: {}",
err.msg()
);
let bracket_empty_host = "ws::addr=[]:9000;";
let err = SenderBuilder::from_conf(bracket_empty_host).unwrap_err();
assert!(err.msg().contains("empty host"), "got: {}", err.msg());
let invalid_port = "ws::addr=localhost:notaport;";
let err = SenderBuilder::from_conf(invalid_port).unwrap_err();
assert!(err.msg().contains("invalid port"), "got: {}", err.msg());
let zero_port = "ws::addr=localhost:0;";
let err = SenderBuilder::from_conf(zero_port).unwrap_err();
assert!(err.msg().contains("invalid port"), "got: {}", err.msg());
#[cfg(feature = "sync-sender-tcp")]
{
let repeated_tcp_addr = "tcp::addr=localhost:9009;addr=localhost:9010;";
let err = SenderBuilder::from_conf(repeated_tcp_addr).unwrap_err();
assert!(
err.msg().contains("DuplicateKey") || err.msg().contains("duplicate"),
"got: {}",
err.msg()
);
}
let conf_async = "ws::addr=localhost:9000;initial_connect_retry=async;";
SenderBuilder::from_conf(conf_async).unwrap();
let bad = "ws::addr=localhost:9000;initial_connect_retry=maybe;";
let err = SenderBuilder::from_conf(bad).unwrap_err();
assert!(
err.msg().contains("initial_connect_retry"),
"got: {}",
err.msg()
);
}
pub(crate) fn upgrade_mock_stream_with_max_batch_size(
stream: &mut TcpStream,
max_batch_size: Option<usize>,
) -> Vec<String> {
let req_bytes = read_request_until_blank(stream).unwrap();
let req = String::from_utf8_lossy(&req_bytes).to_string();
let request_lines: Vec<String> = req
.split("\r\n")
.take_while(|l| !l.is_empty())
.map(String::from)
.collect();
let key = parse_header(&req, "Sec-WebSocket-Key").expect("missing Sec-WebSocket-Key");
let accept = compute_accept(&key);
let batch_size_header = match max_batch_size {
Some(n) => format!("X-QWP-Max-Batch-Size: {n}\r\n"),
None => String::new(),
};
let response = format!(
"HTTP/1.1 101 Switching Protocols\r\n\
Upgrade: websocket\r\n\
Connection: Upgrade\r\n\
Sec-WebSocket-Accept: {accept}\r\n\
X-QWP-Version: 1\r\n\
{batch_size_header}\
\r\n"
);
stream.write_all(response.as_bytes()).unwrap();
request_lines
}
fn spawn_max_batch_size_server(max_batch_size: Option<usize>) -> (u16, mpsc::Receiver<MockResult>) {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
let request_lines = upgrade_mock_stream_with_max_batch_size(&mut stream, max_batch_size);
let mut received_frames = Vec::new();
if let Ok((_fin, _opcode, payload)) = read_frame(&mut stream) {
received_frames.push(payload);
let _ = write_qwp_ok_response(&mut stream, FIRST_WIRE_SEQUENCE);
}
let _ = tx.send(MockResult {
request_lines,
received_frames,
});
thread::sleep(Duration::from_millis(50));
});
(port, rx)
}
fn spawn_max_batch_size_upgrade_only_server(max_batch_size: Option<usize>) -> u16 {
let listener = TcpListener::bind("127.0.0.1:0").unwrap();
let port = listener.local_addr().unwrap().port();
thread::spawn(move || {
let (mut stream, _) = listener.accept().unwrap();
upgrade_mock_stream_with_max_batch_size(&mut stream, max_batch_size);
stream
.set_read_timeout(Some(Duration::from_secs(10)))
.unwrap();
let mut sink = [0u8; 256];
while matches!(stream.read(&mut sink), Ok(n) if n > 0) {}
});
port
}
fn build_buffer_with_payload_bytes(target_encoded_len: usize) -> Buffer {
let mut buf = Buffer::qwp_ws_with_max_name_len(127);
buf.table("t").unwrap();
let col_name = "v";
let value_len = target_encoded_len;
let big_value: String = "x".repeat(value_len);
buf.column_str(col_name, big_value.as_str())
.unwrap()
.at_now()
.unwrap();
buf
}
#[test]
fn server_max_batch_size_clamps_flush_below_configured_max_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let server_cap: usize = 512;
let port = spawn_max_batch_size_upgrade_only_server(Some(server_cap));
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_buf_size(100 * 1024 * 1024)
.unwrap();
let mut sender = build_qwp_ws_sender_from_builder(progress, builder);
let mut buf = build_buffer_with_payload_bytes(server_cap + 200);
let encoded_len = qwp_ws_replay_encoded_len(&buf);
assert!(
encoded_len > server_cap,
"mode={}: encoded_len={encoded_len} must exceed server_cap={server_cap}",
progress.name()
);
let err = sender.flush(&mut buf).unwrap_err();
assert_eq!(
err.code(),
crate::ErrorCode::InvalidApiCall,
"mode={}",
progress.name()
);
assert!(
err.msg()
.contains("exceeds maximum configured allowed size"),
"mode={}: got: {}",
progress.name(),
err.msg()
);
assert!(
err.msg().contains(&server_cap.to_string()),
"mode={}: error should reference server cap {server_cap}, got: {}",
progress.name(),
err.msg()
);
}
}
#[test]
fn server_max_batch_size_allows_flush_when_encoded_fits_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let server_cap: usize = 100 * 1024;
let (port, rx) = spawn_max_batch_size_server(Some(server_cap));
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_buf_size(100 * 1024 * 1024)
.unwrap();
let mut sender = build_qwp_ws_sender_from_builder(progress, builder);
let mut buf = Buffer::qwp_ws_with_max_name_len(127);
buf.table("t")
.unwrap()
.column_i64("v", 42)
.unwrap()
.at_now()
.unwrap();
let encoded_len = qwp_ws_replay_encoded_len(&buf);
assert!(
encoded_len < server_cap,
"mode={}: encoded_len={encoded_len} must be under server_cap={server_cap}",
progress.name()
);
sender.flush(&mut buf).unwrap();
if progress == ProgressCase::Manual {
assert!(sender.drive_once().unwrap());
}
let result = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(result.received_frames.len(), 1, "mode={}", progress.name());
}
}
#[test]
fn absent_server_max_batch_size_falls_back_to_configured_max_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let configured_max: usize = 1024;
let port = spawn_max_batch_size_upgrade_only_server(None);
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_buf_size(configured_max)
.unwrap();
let mut sender = build_qwp_ws_sender_from_builder(progress, builder);
let mut buf = build_buffer_with_payload_bytes(configured_max + 200);
let encoded_len = qwp_ws_replay_encoded_len(&buf);
assert!(
encoded_len > configured_max,
"mode={}: encoded_len={encoded_len} must exceed configured_max={configured_max}",
progress.name()
);
let err = sender.flush(&mut buf).unwrap_err();
assert_eq!(
err.code(),
crate::ErrorCode::InvalidApiCall,
"mode={}",
progress.name()
);
assert!(
err.msg().contains(&configured_max.to_string()),
"mode={}: error should reference configured max {configured_max}, got: {}",
progress.name(),
err.msg()
);
}
}
#[test]
fn configured_max_wins_when_smaller_than_server_cap_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let configured_max: usize = 1024;
let server_cap: usize = 16 * 1024 * 1024;
let port = spawn_max_batch_size_upgrade_only_server(Some(server_cap));
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_buf_size(configured_max)
.unwrap();
let mut sender = build_qwp_ws_sender_from_builder(progress, builder);
let mut buf = build_buffer_with_payload_bytes(configured_max + 200);
let encoded_len = qwp_ws_replay_encoded_len(&buf);
assert!(
encoded_len > configured_max,
"mode={}: encoded_len={encoded_len} must exceed configured_max={configured_max}",
progress.name()
);
assert!(
encoded_len < server_cap,
"mode={}: encoded_len={encoded_len} must be under server_cap={server_cap}",
progress.name()
);
let err = sender.flush(&mut buf).unwrap_err();
assert_eq!(
err.code(),
crate::ErrorCode::InvalidApiCall,
"mode={}",
progress.name()
);
assert!(
err.msg().contains(&configured_max.to_string()),
"mode={}: error should reference configured max {configured_max}, got: {}",
progress.name(),
err.msg()
);
}
}
#[test]
fn server_cap_prevents_message_too_big_by_rejecting_locally_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let server_cap: usize = 2048;
let (port, _rx) = spawn_max_batch_size_server(Some(server_cap));
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_buf_size(100 * 1024 * 1024)
.unwrap();
let mut sender = build_qwp_ws_sender_from_builder(progress, builder);
let mut small_buf = Buffer::qwp_ws_with_max_name_len(127);
small_buf
.table("t")
.unwrap()
.column_i64("v", 1)
.unwrap()
.at_now()
.unwrap();
sender.flush(&mut small_buf).unwrap();
if progress == ProgressCase::Manual {
assert!(sender.drive_once().unwrap());
}
let mut big_buf = build_buffer_with_payload_bytes(server_cap + 500);
let encoded_len = qwp_ws_replay_encoded_len(&big_buf);
assert!(
encoded_len > server_cap,
"mode={}: encoded_len={encoded_len} must exceed server_cap={server_cap}",
progress.name()
);
let err = sender.flush(&mut big_buf).unwrap_err();
assert_eq!(
err.code(),
crate::ErrorCode::InvalidApiCall,
"mode={}: must get local InvalidApiCall, not a server error",
progress.name()
);
}
}
#[test]
fn server_cap_at_exact_boundary_allows_flush_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let mut buf = Buffer::qwp_ws_with_max_name_len(127);
buf.table("t")
.unwrap()
.column_i64("v", 42)
.unwrap()
.at_now()
.unwrap();
let encoded_len = qwp_ws_replay_encoded_len(&buf);
let server_cap = encoded_len;
let (port, rx) = spawn_max_batch_size_server(Some(server_cap));
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_buf_size(100 * 1024 * 1024)
.unwrap();
let mut sender = build_qwp_ws_sender_from_builder(progress, builder);
sender.flush(&mut buf).unwrap();
if progress == ProgressCase::Manual {
assert!(sender.drive_once().unwrap());
}
let result = rx.recv_timeout(Duration::from_secs(5)).unwrap();
assert_eq!(result.received_frames.len(), 1, "mode={}", progress.name());
assert_eq!(
result.received_frames[0].len(),
encoded_len,
"mode={}",
progress.name()
);
}
}
#[test]
fn server_cap_one_byte_below_encoded_len_rejects_flush_in_all_progress_modes() {
for progress in [ProgressCase::Background, ProgressCase::Manual] {
let mut buf = Buffer::qwp_ws_with_max_name_len(127);
buf.table("t")
.unwrap()
.column_i64("v", 42)
.unwrap()
.at_now()
.unwrap();
let encoded_len = qwp_ws_replay_encoded_len(&buf);
let server_cap = encoded_len - 1;
let port = spawn_max_batch_size_upgrade_only_server(Some(server_cap));
let builder = SenderBuilder::new(Protocol::Ws, "127.0.0.1", port)
.max_buf_size(100 * 1024 * 1024)
.unwrap();
let mut sender = build_qwp_ws_sender_from_builder(progress, builder);
let err = sender.flush(&mut buf).unwrap_err();
assert_eq!(
err.code(),
crate::ErrorCode::InvalidApiCall,
"mode={}",
progress.name()
);
assert!(
err.msg().contains(&server_cap.to_string()),
"mode={}: error should reference server cap {server_cap}, got: {}",
progress.name(),
err.msg()
);
}
}