#![cfg(feature = "std")]
use std::error::Error;
use std::net::{TcpListener, TcpStream};
use std::sync::mpsc::{Receiver, Sender, channel};
use std::thread::JoinHandle;
use std::time::Duration;
use liminal::protocol::{Frame, ProtocolError, ProtocolVersion, decode, encode, encoded_len};
use liminal_sdk::SdkError;
use liminal_sdk::remote::{FlushMode, OBSERVABILITY_CHANNEL, PushClient};
const CHANNEL: &str = "app.events";
const DEADLINE: Duration = Duration::from_secs(5);
const SERVER_ERROR_CODE: u16 = 0xFFFF;
const STREAM_ID: u32 = 1;
#[derive(Debug, PartialEq, Eq)]
enum ServerEvent {
Publish { channel: String, payload: Vec<u8> },
Eof,
}
type Responder = Box<dyn FnMut(u64, &str, &[u8]) -> Vec<Frame> + Send>;
struct FakeServer {
addr: String,
events: Receiver<ServerEvent>,
handle: Option<JoinHandle<()>>,
}
impl FakeServer {
fn spawn(responder: Responder) -> Result<Self, Box<dyn Error>> {
let listener = TcpListener::bind("127.0.0.1:0")?;
let addr = listener.local_addr()?.to_string();
let (events_tx, events) = channel();
let mut responder = responder;
let handle = std::thread::spawn(move || {
if let Ok((stream, _)) = listener.accept() {
serve_connection(stream, &mut responder, &events_tx);
}
});
Ok(Self {
addr,
events,
handle: Some(handle),
})
}
fn recv_event(&self) -> Result<ServerEvent, Box<dyn Error>> {
Ok(self.events.recv_timeout(DEADLINE)?)
}
}
impl Drop for FakeServer {
fn drop(&mut self) {
if let Some(handle) = self.handle.take() {
handle.join().ok();
}
}
}
fn serve_connection(
mut stream: TcpStream,
responder: &mut Responder,
events: &Sender<ServerEvent>,
) {
let mut buffer: Vec<u8> = Vec::new();
let mut publish_ordinal: u64 = 0;
loop {
match decode(&buffer) {
Ok((frame, consumed)) => {
buffer.drain(..consumed);
match frame {
Frame::Connect { .. } => {
let ack = Frame::ConnectAck {
flags: 0,
selected_version: ProtocolVersion::new(1, 0),
capabilities: 0,
};
if write_frame(&mut stream, &ack).is_err() {
return;
}
}
Frame::Publish {
channel, envelope, ..
} => {
if events
.send(ServerEvent::Publish {
channel: channel.clone(),
payload: envelope.payload.clone(),
})
.is_err()
{
return;
}
for response in responder(publish_ordinal, &channel, &envelope.payload) {
if write_frame(&mut stream, &response).is_err() {
return;
}
}
publish_ordinal += 1;
}
Frame::Disconnect { .. } => return,
_ => {}
}
}
Err(
ProtocolError::IncompleteHeader { .. } | ProtocolError::TruncatedPayload { .. },
) => {
let mut chunk = [0_u8; 4096];
match std::io::Read::read(&mut stream, &mut chunk) {
Ok(0) => {
events.send(ServerEvent::Eof).ok();
return;
}
Ok(read) => {
if let Some(received) = chunk.get(..read) {
buffer.extend_from_slice(received);
} else {
return;
}
}
Err(_) => return,
}
}
Err(_) => return,
}
}
}
fn write_frame(stream: &mut TcpStream, frame: &Frame) -> Result<(), Box<dyn Error>> {
let len = encoded_len(frame)?;
let mut bytes = vec![0_u8; len];
let written = encode(frame, &mut bytes)?;
let encoded = bytes.get(..written).ok_or("encoder byte count invalid")?;
std::io::Write::write_all(stream, encoded)?;
Ok(())
}
const fn ack(ordinal: u64) -> Frame {
Frame::PublishAck {
flags: 0,
stream_id: STREAM_ID,
message_id: ordinal,
}
}
fn rejection_quoting(payload: &[u8]) -> Frame {
Frame::PublishError {
flags: 0,
stream_id: STREAM_ID,
reason_code: SERVER_ERROR_CODE,
message: Some(format!("rejected:{}", String::from_utf8_lossy(payload))),
}
}
fn selective_responder() -> Responder {
Box::new(|ordinal, channel, payload| {
if channel == OBSERVABILITY_CHANNEL {
Vec::new()
} else if payload.windows(6).any(|window| window == b"reject") {
vec![rejection_quoting(payload)]
} else {
vec![ack(ordinal)]
}
})
}
#[test]
fn flush_pairs_verdicts_fifo_and_excludes_observability() -> Result<(), Box<dyn Error>> {
let server = FakeServer::spawn(selective_responder())?;
let client = PushClient::connect(&server.addr)?;
client.publish(CHANNEL, b"accept-0".to_vec())?;
client.publish(OBSERVABILITY_CHANNEL, b"obs-0".to_vec())?;
client.publish(CHANNEL, b"reject-1".to_vec())?;
client.publish(OBSERVABILITY_CHANNEL, b"obs-1".to_vec())?;
client.publish(CHANNEL, b"accept-2".to_vec())?;
let outcome = client.flush()?;
assert_eq!(outcome.unresolved(), 0, "all verdicts arrive inside budget");
assert_eq!(outcome.mode(), FlushMode::VerdictOnly);
assert!(!outcome.is_proven_accepted());
let failures = outcome.failures();
assert_eq!(failures.len(), 1, "exactly one publish was rejected");
assert_eq!(failures[0].reason_code(), SERVER_ERROR_CODE);
assert_eq!(
failures[0].message(),
Some("rejected:reject-1"),
"the rejection paired to the SECOND ordinary publish, in FIFO order"
);
let repeat = client.flush()?;
assert!(repeat.is_proven_accepted());
assert_eq!(repeat.mode(), FlushMode::VerdictOnly);
Ok(())
}
#[test]
fn unsolicited_publish_response_is_a_typed_mechanism_error() -> Result<(), Box<dyn Error>> {
let responder: Responder = Box::new(|_, channel, _| {
if channel == OBSERVABILITY_CHANNEL {
vec![
ack(0),
Frame::Push {
flags: 0,
stream_id: STREAM_ID,
correlation_id: 7,
payload: b"sync".to_vec(),
},
]
} else {
Vec::new()
}
});
let server = FakeServer::spawn(responder)?;
let client = PushClient::connect(&server.addr)?;
client.publish(OBSERVABILITY_CHANNEL, b"obs".to_vec())?;
let push = client.recv_timeout(DEADLINE)?;
assert_eq!(push.payload(), b"sync");
match client.flush() {
Err(SdkError::Protocol { description }) => {
assert!(
description.contains("response-count mismatch"),
"mechanism error names the broken invariant: {description}"
);
}
other => return Err(format!("expected a typed mechanism error, got {other:?}").into()),
}
Ok(())
}
#[test]
fn budget_expiry_reports_unresolved_not_an_error() -> Result<(), Box<dyn Error>> {
let responder: Responder = Box::new(|_, _, _| Vec::new());
let server = FakeServer::spawn(responder)?;
let client = PushClient::connect(&server.addr)?;
client.publish(CHANNEL, b"withheld-0".to_vec())?;
client.publish(CHANNEL, b"withheld-1".to_vec())?;
let outcome = client.flush()?;
assert!(outcome.failures().is_empty());
assert_eq!(outcome.unresolved(), 2);
assert!(
!outcome.is_proven_accepted(),
"unresolved publishes must never read as proven-accepted"
);
Ok(())
}
#[test]
fn close_sole_owner_half_closes_and_discloses_mode() -> Result<(), Box<dyn Error>> {
let server = FakeServer::spawn(selective_responder())?;
let client = PushClient::connect(&server.addr)?;
client.publish(CHANNEL, b"accept-a".to_vec())?;
client.publish(CHANNEL, b"accept-b".to_vec())?;
let outcome = client.close()?;
assert!(outcome.is_proven_accepted());
assert_eq!(outcome.mode(), FlushMode::FlushedAndHalfClosed);
assert!(matches!(server.recv_event()?, ServerEvent::Publish { .. }));
assert!(matches!(server.recv_event()?, ServerEvent::Publish { .. }));
assert_eq!(server.recv_event()?, ServerEvent::Eof);
Ok(())
}
#[test]
fn close_with_live_clone_is_verdict_only_and_keeps_socket_open() -> Result<(), Box<dyn Error>> {
let server = FakeServer::spawn(selective_responder())?;
let client = PushClient::connect(&server.addr)?;
let writer = client.writer_handle();
client.publish(CHANNEL, b"accept-pre".to_vec())?;
let outcome = client.close()?;
assert!(outcome.is_proven_accepted());
assert_eq!(
outcome.mode(),
FlushMode::VerdictOnly,
"a live clone forbids the FIN; the degradation is disclosed, not silent"
);
assert!(matches!(server.recv_event()?, ServerEvent::Publish { .. }));
writer.publish(CHANNEL, b"after-close".to_vec())?;
match server.recv_event()? {
ServerEvent::Publish { channel, payload } => {
assert_eq!(channel, CHANNEL);
assert_eq!(payload, b"after-close".to_vec());
}
other @ ServerEvent::Eof => {
return Err(format!("expected the clone's post-close publish, got {other:?}").into());
}
}
drop(writer);
Ok(())
}