use std::error::Error;
use std::net::{SocketAddr, TcpStream};
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
use ed25519_dalek::{Signer, SigningKey};
use liminal::protocol::{
CausalContext, Frame, MessageEnvelope, ProtocolError, ProtocolVersion, SchemaId,
WorkerRegistration, decode, encode, encoded_len,
};
use liminal_server::ServerError;
use liminal_server::auth_pass::{PassPrincipal, WirePassV1};
use liminal_server::config::{
AuthConfig, ChannelDef, LimitsConfig, PassConfig, ServerConfig, ServicesConfig, WebSocketConfig,
};
use liminal_server::server::connection::{
ConnectionNotifier, ConnectionServices, ConnectionSupervisor, LiminalConnectionServices,
WebSocketListener,
};
use liminal_server::server::listener::ServerListener;
use liminal_server::server::shutdown::run_shutdown_sequence;
use tungstenite::Message;
use tungstenite::client::IntoClientRequest;
use tungstenite::protocol::WebSocket;
const BEARER: &str = "embedder-bearer";
const INSIDE_CHANNEL: &str = "workspace/acme/events";
const OUTSIDE_CHANNEL: &str = "events";
const PATH: &str = "/liminal";
const DEADLINE: Duration = Duration::from_secs(5);
struct Vector {
seed: [u8; 32],
verifying_key_hex: String,
pass: WirePassV1,
}
fn vector() -> Result<Vector, Box<dyn Error>> {
let value: serde_json::Value =
serde_json::from_str(include_str!("../test-vectors/wire-pass-v1.json"))?;
let field = |name: &str| -> Result<&str, Box<dyn Error>> {
value[name]
.as_str()
.ok_or_else(|| format!("vector field {name} missing").into())
};
let seed: [u8; 32] = hex::decode(field("private_key_seed_hex")?)?
.try_into()
.map_err(|_| "seed length")?;
let verifying_key_hex = field("registry_verifying_key_hex")?.to_owned();
let pass = WirePassV1::parse(&hex::decode(field("pass_hex")?)?)
.map_err(|error| format!("vector pass parse failed: {error:?}"))?;
let signing = SigningKey::from_bytes(&seed);
assert_eq!(
hex::encode(signing.verifying_key().to_bytes()),
verifying_key_hex,
"the vector's seed must derive the vector's registry verifying key"
);
Ok(Vector {
seed,
verifying_key_hex,
pass,
})
}
fn now_secs() -> Result<u64, Box<dyn Error>> {
Ok(SystemTime::now().duration_since(UNIX_EPOCH)?.as_secs())
}
fn mint_pass(vector: &Vector, now: u64) -> Result<Vec<u8>, Box<dyn Error>> {
mint_pass_for(vector, now, &vector.pass.participant)
}
fn mint_pass_for(vector: &Vector, now: u64, participant: &[u8]) -> Result<Vec<u8>, Box<dyn Error>> {
let mut pass = vector.pass.clone();
pass.participant = participant.to_vec();
pass.issued_at = now.saturating_sub(60);
pass.expires_at = now.saturating_add(3_600);
let unsigned = pass
.canonical_unsigned_bytes()
.map_err(|error| format!("canonical encode failed: {error:?}"))?;
pass.signature = SigningKey::from_bytes(&vector.seed)
.sign(&unsigned)
.to_bytes();
pass.canonical_bytes()
.map_err(|error| format!("canonical encode failed: {error:?}").into())
}
fn pass_config(vector: &Vector) -> PassConfig {
PassConfig {
registry_verifying_key: vector.verifying_key_hex.clone(),
maximum_clock_skew_seconds: 0,
}
}
fn auth_config(vector: &Vector) -> AuthConfig {
AuthConfig {
token: BEARER.to_owned(),
pass: Some(pass_config(vector)),
}
}
fn server_config() -> Result<ServerConfig, Box<dyn Error>> {
let health = std::net::TcpListener::bind("127.0.0.1:0")?;
let health_listen_address = health.local_addr()?;
drop(health);
let channel = |name: &str| ChannelDef {
name: name.to_owned(),
schema_ref: None,
durable: false,
loaded_schema: None,
};
Ok(ServerConfig {
listen_address: "127.0.0.1:0".parse()?,
health_listen_address,
drain_timeout_ms: 30_000,
channels: vec![channel(INSIDE_CHANNEL), channel(OUTSIDE_CHANNEL)],
routing_rules: Vec::new(),
persistence_path: None,
cluster: None,
auth: None,
services: ServicesConfig::default(),
limits: LimitsConfig::default(),
websocket: None,
participant: None,
})
}
fn bind_ws(
supervisor: &ConnectionSupervisor,
) -> Result<(WebSocketListener, SocketAddr), Box<dyn Error>> {
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 address = ws.local_addr();
Ok((ws, address))
}
fn encode_frame(frame: &Frame) -> Result<Vec<u8>, Box<dyn Error>> {
let len = encoded_len(frame).map_err(|error| format!("encoded_len: {error}"))?;
let mut bytes = vec![0_u8; len];
let written = encode(frame, &mut bytes).map_err(|error| format!("encode: {error}"))?;
bytes.truncate(written);
Ok(bytes)
}
fn connect_frame(token: &[u8]) -> Frame {
connect_frame_at_version(token, ProtocolVersion::new(1, 0))
}
fn connect_frame_at_version(token: &[u8], version: ProtocolVersion) -> Frame {
Frame::Connect {
flags: 0,
min_version: version,
max_version: version,
auth_token: token.to_vec(),
}
}
fn subscribe_frame(stream_id: u32, channel: &str) -> Frame {
Frame::Subscribe {
flags: 0,
stream_id,
channel: channel.to_owned(),
accepted_schemas: Vec::new(),
max_in_flight: 8,
}
}
const fn envelope(payload: Vec<u8>) -> MessageEnvelope {
MessageEnvelope::new(
SchemaId::new([0_u8; SchemaId::WIRE_LEN]),
CausalContext::independent(),
payload,
)
}
fn ws_connect(address: SocketAddr) -> Result<WebSocket<TcpStream>, Box<dyn Error>> {
let stream = TcpStream::connect(address)?;
stream.set_nodelay(true)?;
stream.set_read_timeout(Some(DEADLINE))?;
let request = format!("ws://{address}{PATH}").into_client_request()?;
let (socket, _response) = tungstenite::client::client(request, stream)
.map_err(|error| format!("websocket client handshake failed: {error}"))?;
socket
.get_ref()
.set_read_timeout(Some(Duration::from_millis(200)))?;
Ok(socket)
}
fn ws_read_frame(socket: &mut WebSocket<TcpStream>) -> Result<Frame, Box<dyn Error>> {
ws_read_frame_within(socket, DEADLINE)
}
fn ws_read_frame_within(
socket: &mut WebSocket<TcpStream>,
bound: Duration,
) -> Result<Frame, Box<dyn Error>> {
let deadline = Instant::now() + bound;
loop {
match socket.read() {
Ok(Message::Binary(bytes)) => {
let (frame, consumed) = decode(&bytes)?;
if consumed != bytes.len() {
return Err("binary message carried trailing bytes".into());
}
return Ok(frame);
}
Ok(Message::Ping(_) | Message::Pong(_)) => {}
Ok(other) => return Err(format!("unexpected websocket message: {other:?}").into()),
Err(tungstenite::Error::Io(error))
if matches!(
error.kind(),
std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
) =>
{
if Instant::now() >= deadline {
return Err("timed out reading a websocket message".into());
}
}
Err(error) => return Err(format!("websocket read failed: {error}").into()),
}
}
}
fn ws_send_frame(socket: &mut WebSocket<TcpStream>, frame: &Frame) -> Result<(), Box<dyn Error>> {
socket.send(Message::Binary(encode_frame(frame)?.into()))?;
Ok(())
}
fn ws_connect_with(
socket: &mut WebSocket<TcpStream>,
token: &[u8],
) -> Result<Frame, Box<dyn Error>> {
ws_send_frame(socket, &connect_frame(token))?;
ws_read_frame(socket)
}
fn expect_connect_ack(frame: &Frame) -> Result<(), Box<dyn Error>> {
match frame {
Frame::ConnectAck { .. } => Ok(()),
other => Err(format!("expected ConnectAck, got {other:?}").into()),
}
}
fn expect_connect_error_message(frame: &Frame, expected: &str) -> Result<(), Box<dyn Error>> {
match frame {
Frame::ConnectError { message, .. } if message.as_deref() == Some(expected) => Ok(()),
other => Err(format!("expected ConnectError {expected:?}, got {other:?}").into()),
}
}
fn ws_subscribe(
socket: &mut WebSocket<TcpStream>,
stream_id: u32,
channel: &str,
) -> Result<bool, Box<dyn Error>> {
ws_send_frame(socket, &subscribe_frame(stream_id, channel))?;
match ws_read_frame(socket)? {
Frame::SubscribeAck { .. } => Ok(true),
Frame::SubscribeError { .. } => Ok(false),
other => Err(format!("expected SubscribeAck or SubscribeError, got {other:?}").into()),
}
}
#[test]
fn embedder_keeps_services_and_verifies_passes_at_connect() -> Result<(), Box<dyn Error>> {
let vector = vector()?;
let config = server_config()?;
let services = Arc::new(LiminalConnectionServices::from_config(&config)?);
let supervisor = ConnectionSupervisor::builder(services.clone())
.auth(&auth_config(&vector))?
.limits(config.limits)
.build()?;
let (_ws, address) = bind_ws(&supervisor)?;
let mut socket = ws_connect(address)?;
expect_connect_ack(&ws_connect_with(
&mut socket,
&mint_pass(&vector, now_secs()?)?,
)?)?;
assert!(
!ws_subscribe(&mut socket, 1, OUTSIDE_CHANNEL)?,
"a pass-stamped connection must be refused outside its live prefix"
);
assert!(
ws_subscribe(&mut socket, 2, INSIDE_CHANNEL)?,
"a pass-stamped connection must be admitted inside its live prefix"
);
let payload = br#""in-process""#.to_vec();
let outcome = services.publish(INSIDE_CHANNEL, &envelope(payload.clone()), None)?;
assert!(
outcome.delivered,
"the in-process publish must be accepted by the pass-stamped subscriber"
);
match ws_read_frame(&mut socket)? {
Frame::Deliver {
envelope: delivered,
..
} => assert_eq!(delivered.payload, payload),
other => return Err(format!("expected Deliver, got {other:?}").into()),
}
Ok(())
}
#[test]
fn builder_bearer_path_is_unchanged_and_a_tampered_pass_is_refused_by_name()
-> Result<(), Box<dyn Error>> {
let vector = vector()?;
let config = server_config()?;
let services = Arc::new(LiminalConnectionServices::from_config(&config)?);
let supervisor = ConnectionSupervisor::builder(services)
.auth(&auth_config(&vector))?
.build()?;
let (_ws, address) = bind_ws(&supervisor)?;
let mut bearer = ws_connect(address)?;
expect_connect_ack(&ws_connect_with(&mut bearer, BEARER.as_bytes())?)?;
assert!(
ws_subscribe(&mut bearer, 1, OUTSIDE_CHANNEL)?,
"a bearer connection carries no principal and is admitted everywhere"
);
let mut tampered_pass = mint_pass(&vector, now_secs()?)?;
let last = tampered_pass.last_mut().ok_or("minted pass is empty")?;
*last ^= 0x01;
let mut tampered = ws_connect(address)?;
expect_connect_error_message(
&ws_connect_with(&mut tampered, &tampered_pass)?,
"connection pass signature check failed",
)?;
Ok(())
}
#[test]
fn existing_services_constructors_still_carry_no_pass_verifier() -> Result<(), Box<dyn Error>> {
let vector = vector()?;
let config = server_config()?;
let services = Arc::new(LiminalConnectionServices::from_config(&config)?);
let supervisor =
ConnectionSupervisor::with_services_and_auth(services, Some(BEARER.as_bytes().to_vec()))?;
let (_ws, address) = bind_ws(&supervisor)?;
let mut socket = ws_connect(address)?;
expect_connect_error_message(
&ws_connect_with(&mut socket, &mint_pass(&vector, now_secs()?)?)?,
"connection authentication token rejected",
)?;
Ok(())
}
#[test]
fn builder_pass_config_is_validated_exactly_as_from_config() -> Result<(), Box<dyn Error>> {
for bad_key in ["not-hex", "00ff"] {
let auth = AuthConfig {
token: BEARER.to_owned(),
pass: Some(PassConfig {
registry_verifying_key: bad_key.to_owned(),
maximum_clock_skew_seconds: 0,
}),
};
let mut config = server_config()?;
config.auth = Some(auth.clone());
let from_config = ConnectionSupervisor::from_config(&config)
.err()
.ok_or("from_config accepted a malformed registry key")?;
let services = Arc::new(LiminalConnectionServices::from_config(&config)?);
let builder = ConnectionSupervisor::builder(services)
.auth(&auth)
.err()
.ok_or("the builder accepted a malformed registry key")?;
assert_eq!(builder.to_string(), from_config.to_string());
}
Ok(())
}
#[derive(Clone, Debug, PartialEq, Eq)]
enum PresenceEvent {
Attached { pid: u64, principal: PassPrincipal },
Detached { pid: u64, principal: PassPrincipal },
}
#[derive(Debug, Default)]
struct PresenceRecorder {
events: Mutex<Vec<PresenceEvent>>,
changed: Condvar,
worker_calls: AtomicUsize,
}
impl PresenceRecorder {
fn record(&self, event: PresenceEvent) {
if let Ok(mut events) = self.events.lock() {
events.push(event);
}
self.changed.notify_all();
}
fn wait_for_len(
&self,
expected: usize,
deadline: Instant,
) -> Result<Vec<PresenceEvent>, Box<dyn Error>> {
let mut events = self
.events
.lock()
.map_err(|error| format!("presence recorder poisoned: {error}"))?;
while events.len() < expected {
let remaining = deadline
.checked_duration_since(Instant::now())
.ok_or_else(|| {
format!("timed out waiting for {expected} presence events; observed {events:?}")
})?;
let (guard, _timeout) = self
.changed
.wait_timeout(events, remaining)
.map_err(|error| format!("presence recorder poisoned: {error}"))?;
events = guard;
}
Ok(events.clone())
}
fn settled_after(&self, dwell: Duration) -> Result<Vec<PresenceEvent>, Box<dyn Error>> {
let deadline = Instant::now() + dwell;
let mut events = self
.events
.lock()
.map_err(|error| format!("presence recorder poisoned: {error}"))?;
while let Some(remaining) = deadline.checked_duration_since(Instant::now()) {
let (guard, _timeout) = self
.changed
.wait_timeout(events, remaining)
.map_err(|error| format!("presence recorder poisoned: {error}"))?;
events = guard;
}
Ok(events.clone())
}
fn worker_calls(&self) -> usize {
self.worker_calls.load(Ordering::SeqCst)
}
}
impl ConnectionNotifier for PresenceRecorder {
fn on_worker_registered(
&self,
_pid: u64,
_registration: &WorkerRegistration,
) -> Result<(), ServerError> {
self.worker_calls.fetch_add(1, Ordering::SeqCst);
Ok(())
}
fn on_worker_unregistered(&self, _pid: u64) {
self.worker_calls.fetch_add(1, Ordering::SeqCst);
}
fn on_pass_attached(&self, pid: u64, principal: &PassPrincipal) {
self.record(PresenceEvent::Attached {
pid,
principal: principal.clone(),
});
}
fn on_pass_detached(&self, pid: u64, principal: &PassPrincipal) {
self.record(PresenceEvent::Detached {
pid,
principal: principal.clone(),
});
}
}
struct Plane {
tcp: ServerListener,
ws: Option<WebSocketListener>,
supervisor: ConnectionSupervisor,
ws_addr: SocketAddr,
}
impl Plane {
fn start(vector: &Vector, recorder: Arc<PresenceRecorder>) -> Result<Self, Box<dyn Error>> {
let config = server_config()?;
let services = Arc::new(LiminalConnectionServices::from_config(&config)?);
let supervisor = ConnectionSupervisor::builder(services)
.auth(&auth_config(vector))?
.notifier(recorder)
.limits(config.limits)
.build()?;
let tcp = ServerListener::bind(&config, supervisor.clone())?;
let (ws, ws_addr) = bind_ws(&supervisor)?;
Ok(Self {
tcp,
ws: Some(ws),
supervisor,
ws_addr,
})
}
fn shutdown(mut self) -> Result<(), Box<dyn Error>> {
let mut ws = self.ws.take().ok_or("websocket listener missing")?;
run_shutdown_sequence(
&mut self.tcp,
Some(&mut ws),
&self.supervisor,
Duration::from_millis(1_500),
)?;
Ok(())
}
}
fn ws_close_cleanly(socket: &mut WebSocket<TcpStream>) -> Result<(), Box<dyn Error>> {
socket.close(None)?;
let deadline = Instant::now() + DEADLINE;
loop {
match socket.read() {
Err(tungstenite::Error::ConnectionClosed) | Ok(Message::Close(_)) => return Ok(()),
Ok(_) => {}
Err(tungstenite::Error::Io(error))
if matches!(
error.kind(),
std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
) =>
{
if Instant::now() >= deadline {
return Err("timed out driving the close handshake".into());
}
}
Err(_) => return Ok(()),
}
}
}
fn expect_attached(event: &PresenceEvent) -> Result<(u64, &PassPrincipal), Box<dyn Error>> {
match event {
PresenceEvent::Attached { pid, principal } => Ok((*pid, principal)),
other @ PresenceEvent::Detached { .. } => {
Err(format!("expected an attach, got {other:?}").into())
}
}
}
#[test]
fn pass_presence_attach_detach_fires_once_per_connection_and_never_on_a_timer()
-> Result<(), Box<dyn Error>> {
let vector = vector()?;
let recorder = Arc::new(PresenceRecorder::default());
let plane = Plane::start(&vector, Arc::clone(&recorder))?;
let now = now_secs()?;
let mut socket_a = ws_connect(plane.ws_addr)?;
expect_connect_ack(&ws_connect_with(&mut socket_a, &mint_pass(&vector, now)?)?)?;
let events = recorder.wait_for_len(1, Instant::now() + DEADLINE)?;
let (pid_a, principal_a) = expect_attached(&events[0])?;
assert_eq!(principal_a.participant, b"participant-42");
assert_eq!(principal_a.public_key, vector.pass.public_key);
assert_eq!(principal_a.live, "workspace/acme/");
assert!(principal_a.may_enroll);
assert_eq!(
principal_a
.conversations
.iter()
.copied()
.collect::<Vec<_>>(),
vec![7, 42, 9_001]
);
assert!(
plane.supervisor.active_connection_pids().contains(&pid_a),
"the attach must name a connection the supervisor tracks"
);
let principal_a = principal_a.clone();
let mut socket_b = ws_connect(plane.ws_addr)?;
expect_connect_ack(&ws_connect_with(
&mut socket_b,
&mint_pass_for(&vector, now, b"participant-43")?,
)?)?;
let events = recorder.wait_for_len(2, Instant::now() + DEADLINE)?;
let (pid_b, principal_b) = expect_attached(&events[1])?;
assert_ne!(pid_a, pid_b, "each connection carries its own id");
assert!(plane.supervisor.active_connection_pids().contains(&pid_b));
assert_eq!(principal_b.participant, b"participant-43");
assert_ne!(
principal_a, *principal_b,
"the two principals must differ for the per-connection assertion to discriminate"
);
assert_eq!(
principal_b.public_key, principal_a.public_key,
"same registry key, same scope: only the participant differs"
);
let principal_b = principal_b.clone();
let settled = recorder.settled_after(Duration::from_millis(400))?;
assert_eq!(
settled.len(),
2,
"no attach or detach may fire while both connections simply stay open: {settled:?}"
);
drop(socket_a);
let events = recorder.wait_for_len(3, Instant::now() + DEADLINE)?;
assert_eq!(
events,
vec![
PresenceEvent::Attached {
pid: pid_a,
principal: principal_a.clone(),
},
PresenceEvent::Attached {
pid: pid_b,
principal: principal_b.clone(),
},
PresenceEvent::Detached {
pid: pid_a,
principal: principal_a.clone(),
},
],
"exactly attach A, attach B, detach A, in order, each with its own principal"
);
ws_close_cleanly(&mut socket_b)?;
let events = recorder.wait_for_len(4, Instant::now() + DEADLINE)?;
assert_eq!(
events[3],
PresenceEvent::Detached {
pid: pid_b,
principal: principal_b,
}
);
let mut socket_c = ws_connect(plane.ws_addr)?;
expect_connect_ack(&ws_connect_with(&mut socket_c, &mint_pass(&vector, now)?)?)?;
let events = recorder.wait_for_len(5, Instant::now() + DEADLINE)?;
let (pid_c, _) = expect_attached(&events[4])?;
plane.shutdown()?;
let events = recorder.wait_for_len(6, Instant::now() + DEADLINE)?;
assert_eq!(
events[5],
PresenceEvent::Detached {
pid: pid_c,
principal: principal_a,
}
);
let settled = recorder.settled_after(Duration::from_millis(200))?;
assert_eq!(settled.len(), 6, "no event may fire after the last detach");
assert_eq!(
recorder.worker_calls(),
0,
"pass presence never rides the worker-registration hooks"
);
Ok(())
}
#[test]
fn bearer_and_refused_connections_fire_no_presence() -> Result<(), Box<dyn Error>> {
let vector = vector()?;
let recorder = Arc::new(PresenceRecorder::default());
let plane = Plane::start(&vector, Arc::clone(&recorder))?;
let mut bearer = ws_connect(plane.ws_addr)?;
expect_connect_ack(&ws_connect_with(&mut bearer, BEARER.as_bytes())?)?;
assert!(ws_subscribe(&mut bearer, 1, OUTSIDE_CHANNEL)?);
let mut tampered_pass = mint_pass(&vector, now_secs()?)?;
let last = tampered_pass.last_mut().ok_or("minted pass is empty")?;
*last ^= 0x01;
let mut refused = ws_connect(plane.ws_addr)?;
expect_connect_error_message(
&ws_connect_with(&mut refused, &tampered_pass)?,
"connection pass signature check failed",
)?;
let mut unsupported = ws_connect(plane.ws_addr)?;
ws_send_frame(
&mut unsupported,
&connect_frame_at_version(
&mint_pass(&vector, now_secs()?)?,
ProtocolVersion::new(99, 0),
),
)?;
match ws_read_frame(&mut unsupported)? {
Frame::ConnectError { .. } => {}
other => return Err(format!("expected a version ConnectError, got {other:?}").into()),
}
assert_eq!(
recorder.settled_after(Duration::from_millis(300))?,
Vec::new()
);
ws_close_cleanly(&mut bearer)?;
drop(refused);
drop(unsupported);
plane.shutdown()?;
assert_eq!(
recorder.settled_after(Duration::from_millis(300))?,
Vec::new(),
"a bearer connection's close and the shutdown fire no detach"
);
assert_eq!(recorder.worker_calls(), 0);
Ok(())
}
#[test]
fn in_process_publisher_reaches_pass_stamped_subscribers_through_the_named_constructor()
-> Result<(), Box<dyn Error>> {
let vector = vector()?;
let config = server_config()?;
let services = Arc::new(LiminalConnectionServices::from_config(&config)?);
let supervisor = ConnectionSupervisor::builder(services.clone())
.auth(&AuthConfig {
token: BEARER.to_owned(),
pass: Some(pass_config(&vector)),
})?
.limits(config.limits)
.build()?;
let (_ws, address) = bind_ws(&supervisor)?;
let mut bad_signature = mint_pass(&vector, now_secs()?)?;
let last = bad_signature.last_mut().ok_or("minted pass is empty")?;
*last ^= 0x01;
let mut refused = ws_connect(address)?;
match ws_connect_with(&mut refused, &bad_signature)? {
Frame::ConnectError {
reason_code,
message,
..
} => {
assert_eq!(
reason_code,
ProtocolError::AuthenticationFailure { message: None }.reason_code(),
"a bad-signature pass is an AuthenticationFailure, by reason code"
);
assert_eq!(
message.as_deref(),
Some("connection pass signature check failed"),
"the refusal names the check that failed"
);
}
other => return Err(format!("expected a named ConnectError, got {other:?}").into()),
}
let mut subscriber = ws_connect(address)?;
expect_connect_ack(&ws_connect_with(
&mut subscriber,
&mint_pass(&vector, now_secs()?)?,
)?)?;
assert!(
ws_subscribe(&mut subscriber, 1, INSIDE_CHANNEL)?,
"the good pass must be admitted inside its live prefix"
);
let quiet = ws_read_frame_within(&mut subscriber, Duration::from_millis(500));
assert!(
quiet.is_err(),
"nothing was published, yet a frame arrived: {quiet:?}"
);
let payload = br#""from the door, in-process""#.to_vec();
let outcome = services.publish(INSIDE_CHANNEL, &envelope(payload.clone()), None)?;
assert!(
outcome.delivered,
"the in-process publish must be accepted by the pass-stamped subscriber"
);
match ws_read_frame_within(&mut subscriber, DEADLINE)? {
Frame::Deliver {
envelope: delivered,
..
} => assert_eq!(delivered.payload, payload),
other => return Err(format!("expected Deliver, got {other:?}").into()),
}
Ok(())
}