use std::error::Error;
use std::sync::Arc;
use liminal_sdk::remote::{FlushMode, PushClient};
use liminal_server::config::{ChannelDef, LimitsConfig, ServerConfig, ServicesConfig};
use liminal_server::server::connection::LiminalConnectionServices;
use liminal_server::server::{ConnectionSupervisor, ServerListener};
const CHANNEL: &str = "app.events";
const SERVER_ERROR_CODE: u16 = 0xFFFF;
struct RunningServer {
tcp: ServerListener,
addr: String,
}
impl RunningServer {
fn start() -> Result<Self, Box<dyn Error>> {
let config = ServerConfig {
listen_address: "127.0.0.1:0".parse()?,
health_listen_address: "127.0.0.1:0".parse()?,
drain_timeout_ms: 4_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: LimitsConfig::default(),
participant: None,
websocket: None,
};
let services = Arc::new(LiminalConnectionServices::from_config(&config)?);
let supervisor =
ConnectionSupervisor::with_services_auth_and_limits(services, None, config.limits)?;
let tcp = ServerListener::bind(&config, supervisor)?;
let addr = tcp.local_addr().to_string();
Ok(Self { tcp, addr })
}
}
#[test]
fn schema_rejection_surfaces_through_flush_with_raw_reason() -> Result<(), Box<dyn Error>> {
let server = RunningServer::start()?;
let publisher = PushClient::connect(&server.addr)?;
publisher.publish(CHANNEL, b"\"ok-0\"".to_vec())?;
publisher.publish(CHANNEL, b"not-json".to_vec())?;
publisher.publish(CHANNEL, b"\"ok-1\"".to_vec())?;
let outcome = publisher.flush()?;
assert_eq!(
outcome.unresolved(),
0,
"all three verdicts arrive inside the flush budget"
);
assert!(!outcome.is_proven_accepted());
let failures = outcome.failures();
assert_eq!(
failures.len(),
1,
"exactly the invalid publish was rejected"
);
assert_eq!(failures[0].reason_code(), SERVER_ERROR_CODE);
let message = failures[0]
.message()
.ok_or("rejection carried no server message")?;
assert!(
message.contains("invalid JSON payload"),
"raw server schema-mismatch text is surfaced verbatim: {message}"
);
drop(publisher);
server.tcp.shutdown()?;
Ok(())
}
#[test]
fn close_reports_clean_acceptance_and_half_close() -> Result<(), Box<dyn Error>> {
const BURST: usize = 16;
let server = RunningServer::start()?;
let publisher = PushClient::connect(&server.addr)?;
for index in 0..BURST {
publisher.publish(CHANNEL, format!("\"event-{index}\"").into_bytes())?;
}
let outcome = publisher.close()?;
assert!(
outcome.is_proven_accepted(),
"every burst publish must be proven accepted: {outcome:?}"
);
assert_eq!(outcome.mode(), FlushMode::FlushedAndHalfClosed);
server.tcp.shutdown()?;
Ok(())
}