use std::error::Error;
use std::net::SocketAddr;
use std::sync::{Mutex, PoisonError};
use std::time::Duration;
use liminal::protocol::{
CausalContext, Frame, MessageEnvelope, ProtocolVersion, SchemaId, decode, encode, encoded_len,
};
use liminal_sdk::SubscriptionStream;
use liminal_sdk::remote::websocket::WebSocketSubscriptionStream;
use liminal_server::config::{
ChannelDef, LimitsConfig, ServerConfig, ServicesConfig, WebSocketConfig,
};
use liminal_server::server::connection::{ConnectionSupervisor, WebSocketListener};
use liminal_server::server::listener::ServerListener;
const CHANNEL: &str = "events";
const PATH: &str = "/liminal";
const BURST: usize = 400;
const ITERATIONS: usize = 6;
const ITERATIONS_ENV: &str = "LIMINAL_STARVATION_ITERS";
const RECV_WINDOW: Duration = Duration::from_secs(2);
const CONNECT_TIMEOUT: Duration = Duration::from_secs(5);
const BURST_PUBLISH_STREAM: u32 = 7;
static PIN_GATE: Mutex<()> = Mutex::new(());
struct RunningServer {
_tcp: ServerListener,
_ws: WebSocketListener,
supervisor: ConnectionSupervisor,
tcp_addr: SocketAddr,
ws_addr: SocketAddr,
}
impl RunningServer {
fn start() -> Result<Self, Box<dyn Error>> {
Self::start_with(LimitsConfig::default())
}
fn start_with(limits: LimitsConfig) -> Result<Self, Box<dyn Error>> {
let health = std::net::TcpListener::bind("127.0.0.1:0")?;
let health_listen_address = health.local_addr()?;
drop(health);
let config = ServerConfig {
listen_address: "127.0.0.1:0".parse()?,
health_listen_address,
drain_timeout_ms: 30_000,
channels: vec![ChannelDef {
name: CHANNEL.to_owned(),
schema_ref: None,
durable: false,
loaded_schema: None,
}],
routing_rules: Vec::new(),
persistence_path: None,
cluster: None,
auth: None,
services: ServicesConfig::default(),
limits,
websocket: None,
participant: None,
};
let supervisor = ConnectionSupervisor::from_config(&config)?;
let tcp = ServerListener::bind(&config, supervisor.clone())?;
let ws_config = WebSocketConfig {
listen_address: "127.0.0.1:0".parse()?,
path: PATH.to_owned(),
allowed_origins: Vec::new(),
ping_interval_ms: None,
};
let ws = WebSocketListener::bind(&ws_config, supervisor.clone())?;
let tcp_addr = tcp.local_addr();
let ws_addr = ws.local_addr();
Ok(Self {
_tcp: tcp,
_ws: ws,
supervisor,
tcp_addr,
ws_addr,
})
}
fn tcp_address(&self) -> String {
self.tcp_addr.to_string()
}
fn ws_url(&self) -> String {
format!("ws://{}{PATH}", self.ws_addr)
}
}
impl Drop for RunningServer {
fn drop(&mut self) {
self.supervisor.shutdown();
}
}
fn open_tcp_subscription(server: &RunningServer) -> Result<SubscriptionStream, Box<dyn Error>> {
let deadline = std::time::Instant::now() + CONNECT_TIMEOUT;
let mut last_error = None;
while std::time::Instant::now() < deadline {
match SubscriptionStream::open(&server.tcp_address(), CHANNEL, Vec::new()) {
Ok(stream) => return Ok(stream),
Err(error) => {
last_error = Some(error);
std::thread::sleep(Duration::from_millis(20));
}
}
}
Err(format!("tcp subscription never opened: {last_error:?}").into())
}
fn open_ws_subscription(
server: &RunningServer,
) -> Result<WebSocketSubscriptionStream, Box<dyn Error>> {
let deadline = std::time::Instant::now() + CONNECT_TIMEOUT;
let mut last_error = None;
while std::time::Instant::now() < deadline {
match WebSocketSubscriptionStream::open(&server.ws_url(), CHANNEL, Vec::new()) {
Ok(stream) => return Ok(stream),
Err(error) => {
last_error = Some(error);
std::thread::sleep(Duration::from_millis(20));
}
}
}
Err(format!("websocket subscription never opened: {last_error:?}").into())
}
fn write_frame(stream: &mut std::net::TcpStream, frame: &Frame) -> Result<(), Box<dyn Error>> {
use std::io::Write;
let mut bytes = vec![0_u8; encoded_len(frame)?];
let written = encode(frame, &mut bytes)?;
bytes.truncate(written);
stream.write_all(&bytes)?;
Ok(())
}
fn publish_burst_raw_tcp(
server: &RunningServer,
count: usize,
) -> Result<(usize, Vec<String>), Box<dyn Error>> {
use std::io::Read;
use std::sync::mpsc;
let mut stream = std::net::TcpStream::connect(server.tcp_addr)?;
stream.set_read_timeout(Some(Duration::from_secs(10)))?;
write_frame(
&mut stream,
&Frame::Connect {
flags: 0,
min_version: ProtocolVersion::new(1, 0),
max_version: ProtocolVersion::new(1, 0),
auth_token: Vec::new(),
},
)?;
let mut buffer = vec![0_u8; 4096];
let read = stream.read(&mut buffer)?;
decode(&buffer[..read])?;
let mut drain = stream.try_clone()?;
let (done, all_acked) = mpsc::channel();
let reader = std::thread::spawn(move || {
let mut pending: Vec<u8> = Vec::new();
let mut chunk = vec![0_u8; 16384];
let mut accepted = 0_usize;
let mut others: Vec<String> = Vec::new();
loop {
let read = match drain.read(&mut chunk) {
Ok(0) | Err(_) => break,
Ok(read) => read,
};
pending.extend_from_slice(&chunk[..read]);
while let Ok((frame, consumed)) = decode(&pending) {
pending.drain(..consumed);
match frame {
Frame::PublishAck { .. } => accepted += 1,
other => others.push(format!("{other:?}")),
}
}
if accepted >= count {
done.send(()).ok();
}
}
(accepted, others)
});
for index in 0..count {
write_frame(
&mut stream,
&Frame::Publish {
flags: 0,
stream_id: BURST_PUBLISH_STREAM,
channel: CHANNEL.to_owned(),
envelope: MessageEnvelope::new(
SchemaId::new([7_u8; SchemaId::WIRE_LEN]),
CausalContext::independent(),
format!("{{\"id\":{index}}}").into_bytes(),
),
idempotency_key: None,
},
)?;
}
all_acked.recv_timeout(Duration::from_secs(20)).ok();
stream.shutdown(std::net::Shutdown::Both).ok();
let (accepted, others) = reader.join().map_err(|_| "ack reader panicked")?;
Ok((accepted, others))
}
fn drain_ws(stream: &WebSocketSubscriptionStream, wanted: usize) -> (usize, Option<String>) {
let mut received = 0_usize;
while received < wanted {
match stream.recv_timeout(RECV_WINDOW) {
Ok(_) => received += 1,
Err(error) => return (received, Some(format!("{error}"))),
}
}
(received, None)
}
fn iterations() -> usize {
std::env::var(ITERATIONS_ENV)
.ok()
.and_then(|value| value.parse::<usize>().ok())
.filter(|value| *value > 0)
.unwrap_or(ITERATIONS)
}
fn drain_tcp(stream: &SubscriptionStream, wanted: usize) -> (usize, Option<String>) {
let mut received = 0_usize;
while received < wanted {
match stream.recv_timeout(RECV_WINDOW) {
Ok(_) => received += 1,
Err(error) => return (received, Some(format!("{error}"))),
}
}
(received, None)
}
fn burst_once() -> Result<(bool, String), Box<dyn Error>> {
let server = RunningServer::start()?;
let tcp_stream = open_tcp_subscription(&server)?;
let ws_stream = open_ws_subscription(&server)?;
let (accepted, publisher_frames) = publish_burst_raw_tcp(&server, BURST)?;
let (tcp_received, tcp_error) = drain_tcp(&tcp_stream, accepted);
let (ws_received, ws_error) = drain_ws(&ws_stream, accepted);
let detail = format!(
"burst={BURST} accepted={accepted} tcp_received={tcp_received} tcp_error={tcp_error:?} \
ws_received={ws_received} ws_error={ws_error:?} publisher_frames={publisher_frames:?}"
);
Ok((
accepted == BURST && tcp_received == accepted && ws_received == accepted,
detail,
))
}
#[test]
fn a_burst_larger_than_the_delivery_slice_never_starves_a_subscriber()
-> Result<(), Box<dyn Error>> {
let _gate = PIN_GATE.lock().unwrap_or_else(PoisonError::into_inner);
let boots = iterations();
let mut failures = Vec::new();
for index in 0..boots {
let (ok, detail) = burst_once()?;
eprintln!("BURST PIN iteration {index}: ok={ok} :: {detail}");
if !ok {
failures.push(format!("iteration {index}: {detail}"));
}
}
assert!(
failures.is_empty(),
"the burst pin lost a subscriber on {}/{boots} fresh boots:\n{}",
failures.len(),
failures.join("\n")
);
Ok(())
}
#[test]
fn c_the_smallest_configurable_slice_budget_still_delivers_everything()
-> Result<(), Box<dyn Error>> {
const SMALL_BURST: usize = 120;
let _gate = PIN_GATE.lock().unwrap_or_else(PoisonError::into_inner);
let limits = LimitsConfig {
delivery_slice_budget: 1,
..LimitsConfig::default()
};
let server = RunningServer::start_with(limits)?;
let tcp_stream = open_tcp_subscription(&server)?;
let ws_stream = open_ws_subscription(&server)?;
let (accepted, publisher_frames) = publish_burst_raw_tcp(&server, SMALL_BURST)?;
assert_eq!(
accepted, SMALL_BURST,
"the publisher must be acked for every record before delivery is judged; \
otherwise a dead publisher scores 0 == 0 and this pin passes vacuously \
(publisher frames: {publisher_frames:?})"
);
let (tcp_received, tcp_error) = drain_tcp(&tcp_stream, accepted);
let (ws_received, ws_error) = drain_ws(&ws_stream, accepted);
assert_eq!(
tcp_received, accepted,
"a budget of 1 must still drain the TCP subscriber: {tcp_error:?}"
);
assert_eq!(
ws_received, accepted,
"a budget of 1 must still drain the WebSocket subscriber: {ws_error:?}"
);
Ok(())
}
fn mixed_fate_once() -> Result<(bool, String), Box<dyn Error>> {
let server = RunningServer::start()?;
let first = open_ws_subscription(&server)?;
let second = open_ws_subscription(&server)?;
let (accepted, publisher_frames) = publish_burst_raw_tcp(&server, BURST)?;
let (first_received, first_error) = drain_ws(&first, accepted);
let (second_received, second_error) = drain_ws(&second, accepted);
let both_ok = first_received == accepted && second_received == accepted;
let mixed = (first_received == accepted) != (second_received == accepted);
let detail = format!(
"burst={BURST} accepted={accepted} first_received={first_received} \
first_error={first_error:?} second_received={second_received} \
second_error={second_error:?} mixed_fate={mixed} \
publisher_frames={publisher_frames:?}"
);
Ok((accepted == BURST && both_ok, detail))
}
#[test]
fn b_two_websocket_subscribers_on_one_boot_share_the_same_fate() -> Result<(), Box<dyn Error>> {
let _gate = PIN_GATE.lock().unwrap_or_else(PoisonError::into_inner);
let boots = iterations();
let mut failures = Vec::new();
for index in 0..boots {
let (ok, detail) = mixed_fate_once()?;
eprintln!("MIXED-FATE PIN iteration {index}: ok={ok} :: {detail}");
if !ok {
failures.push(format!("iteration {index}: {detail}"));
}
}
assert!(
failures.is_empty(),
"the mixed-fate pin lost a subscriber on {}/{boots} fresh boots:\n{}",
failures.len(),
failures.join("\n")
);
Ok(())
}