ROADMAP gains a v0.8 entry: the smarm 415effb lifetime change as root cause, the description/instantiation split, the tree shape, the accepted per-call resolution cost, and what stays open (RegistryName newtype, PubSub's unprotectable const, dynamic session actors, hammer.sh's missing feature matrix). The v0.7 "unreproduced test failure" open item is closed out and pointed at it — it was this hang, hidden because hammer.sh builds with default features while the failure needs --all-features load. README and the two module headers still taught the v0.5 rules: in-runtime-only construction, the non-static Arc<OnceLock<..>> cell, and "a relay must never hold a PubSub clone". None of those are true any more; they are replaced by what actually holds now, with a note on what changed for anyone who learnt the old shape. session.rs also gains an honest note that its session actors are the last lifetime in urus implied by a drop rather than stated — they are dynamic, so a fixed ChildSpec list cannot hold them, and the trigger is at least a command now rather than a refcount.
781 lines
30 KiB
Rust
781 lines
30 KiB
Rust
//! Opt-in channel session persistence (the ratified v0.6 design).
|
|
//!
|
|
//! Default channel lifetime is transport-bound: cold start per join,
|
|
//! actor dies with the socket. A [`ChannelSession`] impl registered via
|
|
//! [`PrefixRouter::channel_session`](super::PrefixRouter::channel_session)
|
|
//! changes that: the channel actor outlives its transport, buffering
|
|
//! outbound broadcasts (`VecDeque<Arc<Broadcast<P>>>` — the same `Arc`s
|
|
//! the pubsub relay path carries, zero re-allocation) until the client
|
|
//! rejoins, the buffer fills, or the TTL expires.
|
|
//!
|
|
//! Shape (one registry gen_server per `channel_session` registration,
|
|
//! monomorphic over the session key — no type erasure):
|
|
//!
|
|
//! - `deploy` on the conn actor computes the session key and casts the
|
|
//! whole join handshake to the registry.
|
|
//! - The registry owns `key -> (pid, control sender)`. Existing entry:
|
|
//! the handshake is forwarded as an `Attach` (reattach). Absent or
|
|
//! dead (send failure raced the monitor `Down`): spawn fresh,
|
|
//! monitor, insert. Down prunes, guarded by pid match against
|
|
//! replaced entries.
|
|
//! - The session actor is **not linked** to any connection — outliving
|
|
//! the transport is the point. Its exits: explicit leave, rejected
|
|
//! (re)join, TTL expiry, buffer cap, pubsub relay death, or the
|
|
//! registry's control sender dropping — which is exactly the shutdown
|
|
//! chain (the supervisor shuts the registry down after the endpoint
|
|
//! has drained -> the registry's state drops -> control senders drop
|
|
//! -> every parked session wakes on the closed control arm,
|
|
//! terminates, and exits; `AllDone` composes without links).
|
|
//!
|
|
//! Note this is the last lifetime in urus implied by a drop rather
|
|
//! than stated: sessions are dynamic (one per key), so a fixed
|
|
//! `ChildSpec` list cannot hold them. The trigger is now a command
|
|
//! rather than a refcount, which is what the v0.8 cycle was about, but
|
|
//! a dynamic supervisor (OTP's `simple_one_for_one`) is the honest
|
|
//! shape if smarm grows one.
|
|
//!
|
|
//! As-landed decisions (veto by diff):
|
|
//! - **Every attach calls `ch.join()` again** on the same instance —
|
|
//! the client's `phx_join` needs an ack payload and the channel gets
|
|
//! to re-auth. Session channels must treat `join` as re-entrant.
|
|
//! - **A rejected rejoin ends the session** (error reply, `terminate`,
|
|
//! exit): `Err` means "this transport may not have this channel", and
|
|
//! a zombie session held open for an unauthorized client is a
|
|
//! liability. Next join is a cold start.
|
|
//! - **A second transport evicts the first** (best-effort `phx_close`
|
|
//! on the old socket): a session is single-transport by definition;
|
|
//! two tabs wanting independent channels should not share a session
|
|
//! key.
|
|
//! - **TTL is per detach episode** (reset on every disconnect), and the
|
|
//! pubsub subscription is made once, on the first successful join, to
|
|
//! the first join's topic.
|
|
//! - **`ChannelSession::Key: Clone`** beyond the ratified
|
|
//! `Eq + Hash + Send` — the registry keeps a `pid -> key` reverse
|
|
//! index for monitor-`Down` pruning.
|
|
|
|
use super::*;
|
|
|
|
use smarm::gen_server::{self, GenServer, GenServerBuilder, GenServerCtx, GenServerName};
|
|
use smarm::{Down, Restart, Watcher};
|
|
|
|
use std::collections::VecDeque;
|
|
use std::hash::Hash;
|
|
use std::time::{Duration, Instant};
|
|
|
|
/// Opt-in session persistence config for a channel type. All methods
|
|
/// are static: the config is consulted before any channel instance
|
|
/// exists.
|
|
pub trait ChannelSession<P>: Send + 'static {
|
|
/// What identifies "the same session" across reconnects.
|
|
type Key: Eq + Hash + Send + 'static;
|
|
|
|
/// Derive the session key from a join. Same key, same (live)
|
|
/// channel actor; the join payload is the natural carrier for a
|
|
/// client-held session token.
|
|
fn session_key(topic: &str, join_payload: &P) -> Self::Key;
|
|
|
|
/// Detached broadcasts buffered before the session tears itself
|
|
/// down (next join is a cold start).
|
|
fn buffer_cap() -> usize {
|
|
128
|
|
}
|
|
|
|
/// How long a detached session waits for a rejoin before tearing
|
|
/// itself down. Reset on every detach.
|
|
fn ttl() -> Duration {
|
|
Duration::from_secs(30)
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// The session-aware factory (what channel_session registers)
|
|
// ---------------------------------------------------------------------------
|
|
|
|
pub(super) struct SessionFactory<P: Encode + Decode + Send + Sync + 'static, K: SessionKey> {
|
|
inner: Arc<dyn ChannelFactory<P>>,
|
|
keyfn: fn(&str, &P) -> K,
|
|
/// The registry's registered name. Resolved per deploy, so a
|
|
/// registry restarted by the supervisor is reached transparently —
|
|
/// with its session map empty, which is the honest outcome: the
|
|
/// session actors it tracked died with their control senders.
|
|
registry: GenServerName<Registry<P, K>>,
|
|
/// Everything the registry's `ChildSpec` needs to build it again.
|
|
cap: usize,
|
|
ttl: Duration,
|
|
}
|
|
|
|
/// The registry's working bounds for a session key.
|
|
pub(super) trait SessionKey: Eq + Hash + Clone + Send + 'static {}
|
|
impl<K: Eq + Hash + Clone + Send + 'static> SessionKey for K {}
|
|
|
|
impl<P: Encode + Decode + Send + Sync + 'static, K: SessionKey> SessionFactory<P, K> {
|
|
/// Describe the session route. Spawns nothing; the registry actor is
|
|
/// started from [`children`](ChannelFactory::children).
|
|
pub(super) fn new<S: ChannelSession<P, Key = K>>(
|
|
registry: &'static str,
|
|
inner: Arc<dyn ChannelFactory<P>>,
|
|
) -> Self {
|
|
SessionFactory {
|
|
inner,
|
|
keyfn: S::session_key,
|
|
registry: GenServerName::new(registry),
|
|
cap: S::buffer_cap(),
|
|
ttl: S::ttl(),
|
|
}
|
|
}
|
|
}
|
|
|
|
impl<P: Encode + Decode + Send + Sync + 'static, K: SessionKey> ChannelFactory<P>
|
|
for SessionFactory<P, K>
|
|
{
|
|
fn create(&self, topic: &str) -> Box<dyn Channel<P>> {
|
|
self.inner.create(topic)
|
|
}
|
|
|
|
fn children(&self) -> Vec<ChildSpec> {
|
|
let (name, factory, cap, ttl) = (self.registry, self.inner.clone(), self.cap, self.ttl);
|
|
vec![ChildSpec::new(Restart::Permanent, move || {
|
|
let state = Registry::<P, K> {
|
|
factory: factory.clone(),
|
|
cap,
|
|
ttl,
|
|
sessions: HashMap::new(),
|
|
pid_key: HashMap::new(),
|
|
watcher: None,
|
|
};
|
|
if let Err(e) = GenServerBuilder::new(state).named(name).run() {
|
|
panic!("urus session registry '{}': name already taken: {e:?}", name.as_str());
|
|
}
|
|
})]
|
|
}
|
|
|
|
fn deploy(&self, hs: JoinHandshake<P>) -> ChannelInbox<P>
|
|
where
|
|
P: Encode + Decode,
|
|
{
|
|
let key = (self.keyfn)(&hs.topic, &hs.payload);
|
|
let (tx, rx) = smarm::channel::<In<P>>();
|
|
// Cast: the actor acks the join straight to the socket; the
|
|
// conn actor has nothing to wait for. A dead registry can only
|
|
// mean shutdown — the failed join is moot.
|
|
let _ = gen_server::cast(self.registry, Join { key, hs, inbound: rx });
|
|
ChannelInbox(tx)
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// The registry gen_server: key -> live session actor
|
|
// ---------------------------------------------------------------------------
|
|
|
|
pub(super) struct Join<P: Send + Sync + 'static, K> {
|
|
key: K,
|
|
hs: JoinHandshake<P>,
|
|
inbound: Receiver<In<P>>,
|
|
}
|
|
|
|
struct Registry<P: Encode + Decode + Send + Sync + 'static, K: SessionKey> {
|
|
factory: Arc<dyn ChannelFactory<P>>,
|
|
cap: usize,
|
|
ttl: Duration,
|
|
sessions: HashMap<K, (Pid, Sender<Ctl<P>>)>,
|
|
/// Reverse index for `handle_down` — the reason `Key: Clone`.
|
|
pid_key: HashMap<Pid, K>,
|
|
watcher: Option<Watcher<Registry<P, K>>>,
|
|
}
|
|
|
|
impl<P: Encode + Decode + Send + Sync + 'static, K: SessionKey> GenServer for Registry<P, K> {
|
|
type Call = ();
|
|
type Reply = ();
|
|
type Cast = Join<P, K>;
|
|
type Info = ();
|
|
type Timer = ();
|
|
|
|
fn init(&mut self, ctx: &GenServerCtx<Self>) {
|
|
self.watcher = Some(ctx.watcher());
|
|
}
|
|
|
|
fn handle_call(&mut self, _request: ()) {}
|
|
|
|
fn handle_cast(&mut self, Join { key, hs, inbound }: Join<P, K>) {
|
|
let attach = Ctl::Attach { hs, inbound };
|
|
let attach = match self.sessions.get(&key) {
|
|
Some((_, ctl)) => match ctl.send(attach) {
|
|
Ok(()) => return, // reattached
|
|
// The actor died and its Down hasn't landed yet:
|
|
// recover the handshake, replace below.
|
|
Err(smarm::channel::SendError(a)) => a,
|
|
},
|
|
None => attach,
|
|
};
|
|
let (ctl_tx, ctl_rx) = smarm::channel::<Ctl<P>>();
|
|
let (factory, cap, ttl) = (self.factory.clone(), self.cap, self.ttl);
|
|
let pid = smarm::spawn(move || run_session(factory, cap, ttl, ctl_rx)).pid();
|
|
// Cannot fail: the receiver is owned by the just-spawned actor.
|
|
let _ = ctl_tx.send(attach);
|
|
self.watcher
|
|
.as_ref()
|
|
.expect("watcher set in init")
|
|
.watch(smarm::monitor(pid));
|
|
self.pid_key.insert(pid, key.clone());
|
|
self.sessions.insert(key, (pid, ctl_tx));
|
|
}
|
|
|
|
fn handle_down(&mut self, down: Down) {
|
|
if let Some(key) = self.pid_key.remove(&down.pid) {
|
|
// Pid guard: a stale Down for an entry handle_cast already
|
|
// replaced must not evict the replacement.
|
|
if self.sessions.get(&key).is_some_and(|(p, _)| *p == down.pid) {
|
|
self.sessions.remove(&key);
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// The session actor
|
|
// ---------------------------------------------------------------------------
|
|
|
|
enum Ctl<P: Send + Sync + 'static> {
|
|
Attach { hs: JoinHandshake<P>, inbound: Receiver<In<P>> },
|
|
}
|
|
|
|
enum Step<P: Send + Sync + 'static> {
|
|
Attach { hs: JoinHandshake<P>, inbound: Receiver<In<P>> },
|
|
Attached { inbound: Receiver<In<P>> },
|
|
Detached,
|
|
Exit,
|
|
}
|
|
|
|
fn run_session<P: Encode + Decode + Send + Sync + 'static>(
|
|
factory: Arc<dyn ChannelFactory<P>>,
|
|
cap: usize,
|
|
ttl: Duration,
|
|
ctl: Receiver<Ctl<P>>,
|
|
) {
|
|
// The initial Attach is cast by the registry in the same handler
|
|
// that spawned us. Closed instead = registry died first (shutdown
|
|
// raced the spawn): nothing exists yet, nothing to clean.
|
|
let Ok(Ctl::Attach { hs, inbound }) = ctl.recv() else { return };
|
|
|
|
// Partial moves: ws/bus/topic/join_ref go into the socket,
|
|
// reference/payload stay behind for the first attach below.
|
|
let mut socket = ChannelSocket {
|
|
ws: hs.ws,
|
|
bus: hs.bus,
|
|
topic: hs.topic,
|
|
join_ref: hs.join_ref,
|
|
cur_ref: RefCell::new(None),
|
|
};
|
|
let mut ch: Option<Box<dyn Channel<P>>> = None;
|
|
let mut joined = false;
|
|
let mut bus_rx: Option<Receiver<Arc<Broadcast<P>>>> = None;
|
|
let mut buffer: VecDeque<Arc<Broadcast<P>>> = VecDeque::new();
|
|
|
|
let mut step = attach(
|
|
&factory,
|
|
&mut ch,
|
|
&mut joined,
|
|
&mut bus_rx,
|
|
&mut buffer,
|
|
&mut socket,
|
|
hs.reference,
|
|
hs.payload,
|
|
inbound,
|
|
);
|
|
loop {
|
|
step = match step {
|
|
Step::Attach { hs, inbound } => {
|
|
socket.ws = hs.ws;
|
|
socket.join_ref = hs.join_ref;
|
|
socket.topic = hs.topic;
|
|
attach(
|
|
&factory,
|
|
&mut ch,
|
|
&mut joined,
|
|
&mut bus_rx,
|
|
&mut buffer,
|
|
&mut socket,
|
|
hs.reference,
|
|
hs.payload,
|
|
inbound,
|
|
)
|
|
}
|
|
Step::Attached { inbound } => attached_loop(
|
|
ch.as_deref_mut().expect("attached implies created"),
|
|
&socket,
|
|
&ctl,
|
|
&inbound,
|
|
bus_rx.as_ref().expect("attached implies subscribed"),
|
|
&mut buffer,
|
|
),
|
|
Step::Detached => detached_loop(
|
|
&ctl,
|
|
bus_rx.as_ref().expect("detached only after first subscribe"),
|
|
&mut buffer,
|
|
cap,
|
|
ttl,
|
|
),
|
|
Step::Exit => {
|
|
if joined {
|
|
if let Some(c) = ch.as_deref_mut() {
|
|
c.terminate();
|
|
}
|
|
}
|
|
return;
|
|
}
|
|
};
|
|
}
|
|
}
|
|
|
|
/// One attach handshake: create-once, (re)join, subscribe-once, ack,
|
|
/// drain. On a drain failure the undelivered tail stays in `buffer`
|
|
/// for the next attach.
|
|
#[allow(clippy::too_many_arguments)]
|
|
fn attach<P: Encode + Decode + Send + Sync + 'static>(
|
|
factory: &Arc<dyn ChannelFactory<P>>,
|
|
ch: &mut Option<Box<dyn Channel<P>>>,
|
|
joined: &mut bool,
|
|
bus_rx: &mut Option<Receiver<Arc<Broadcast<P>>>>,
|
|
buffer: &mut VecDeque<Arc<Broadcast<P>>>,
|
|
socket: &mut ChannelSocket<P>,
|
|
reference: Option<String>,
|
|
payload: P,
|
|
inbound: Receiver<In<P>>,
|
|
) -> Step<P> {
|
|
if ch.is_none() {
|
|
*ch = Some(factory.create(&socket.topic));
|
|
}
|
|
let c = ch.as_deref_mut().expect("just created");
|
|
|
|
socket.set_ref(reference);
|
|
let topic = socket.topic.clone();
|
|
match c.join(&topic, payload, socket) {
|
|
Ok(reply) => {
|
|
if bus_rx.is_none() {
|
|
// First successful join: subscribe BEFORE the ok reply
|
|
// (subscribe is a call — once the client sees its ack,
|
|
// membership is a fact, not a race). Once per session:
|
|
// the subscription is actor-scoped, survives detach,
|
|
// and that survival IS the buffering path.
|
|
match socket.bus.subscribe(&topic) {
|
|
Ok(rx) => *bus_rx = Some(rx),
|
|
Err(_) => {
|
|
socket.send_reply(Status::Error, None);
|
|
socket.set_ref(None);
|
|
return Step::Exit;
|
|
}
|
|
}
|
|
}
|
|
*joined = true;
|
|
socket.send_reply(Status::Ok, Some(&reply));
|
|
socket.set_ref(None);
|
|
// Drain with the NEW join generation's ref; buffered pushes
|
|
// have no request to correlate to (reference: None).
|
|
while let Some(b) = buffer.front() {
|
|
let out = P::encode(FrameRef {
|
|
join_ref: socket.join_ref.as_deref(),
|
|
reference: None,
|
|
topic: &socket.topic,
|
|
message: MessageRef::Event { event: &b.event, payload: Some(&b.payload) },
|
|
});
|
|
if socket.ws.send(out).is_ok() {
|
|
buffer.pop_front();
|
|
} else {
|
|
return Step::Detached;
|
|
}
|
|
}
|
|
Step::Attached { inbound }
|
|
}
|
|
Err(reply) => {
|
|
// Rejected (re)join = session over, by decision: error
|
|
// reply, terminate (in run_session, iff it ever joined),
|
|
// exit. Next join is a cold start.
|
|
socket.send_reply(Status::Error, Some(&reply));
|
|
socket.set_ref(None);
|
|
Step::Exit
|
|
}
|
|
}
|
|
}
|
|
|
|
fn attached_loop<P: Encode + Decode + Send + Sync + 'static>(
|
|
ch: &mut dyn Channel<P>,
|
|
socket: &ChannelSocket<P>,
|
|
ctl: &Receiver<Ctl<P>>,
|
|
inbound: &Receiver<In<P>>,
|
|
bus_rx: &Receiver<Arc<Broadcast<P>>>,
|
|
buffer: &mut VecDeque<Arc<Broadcast<P>>>,
|
|
) -> Step<P> {
|
|
// ctl sits at lowest priority: it only ever carries an eviction by
|
|
// a second transport (rare) or the closed-arm shutdown signal. A
|
|
// closed arm stays ready forever, but every closed observation
|
|
// exits the loop immediately — no spin, no starvation window.
|
|
loop {
|
|
match smarm::select(&[inbound, bus_rx, ctl]) {
|
|
0 => match inbound.try_recv() {
|
|
Ok(Some(In::Event { reference, event, payload })) => {
|
|
socket.set_ref(reference);
|
|
ch.handle_in(&event, payload, socket);
|
|
socket.set_ref(None);
|
|
}
|
|
Ok(Some(In::Leave { reference })) => {
|
|
// Explicit leave ends the SESSION, not just the
|
|
// attachment: ok reply, close event, terminate.
|
|
socket.set_ref(reference);
|
|
socket.send_reply(Status::Ok, None);
|
|
socket.set_ref(None);
|
|
close_event(socket);
|
|
return Step::Exit;
|
|
}
|
|
Ok(None) => {}
|
|
// Conn actor finished or replaced this join generation:
|
|
// the transport is gone, the session is not.
|
|
Err(_) => return Step::Detached,
|
|
},
|
|
1 => match bus_rx.try_recv() {
|
|
Ok(Some(b)) => {
|
|
let out = P::encode(FrameRef {
|
|
join_ref: socket.join_ref.as_deref(),
|
|
reference: None,
|
|
topic: &socket.topic,
|
|
message: MessageRef::Event { event: &b.event, payload: Some(&b.payload) },
|
|
});
|
|
if socket.ws.send(out).is_err() {
|
|
// Transport died under us: this broadcast is
|
|
// the buffer's first entry, nothing was lost.
|
|
buffer.push_back(b);
|
|
return Step::Detached;
|
|
}
|
|
}
|
|
Ok(None) => {}
|
|
// The topic relay died: pubsub is gone, session over.
|
|
Err(_) => return Step::Exit,
|
|
},
|
|
2 => match ctl.try_recv() {
|
|
Ok(Some(Ctl::Attach { hs, inbound })) => {
|
|
// A second transport claims the session key: evict
|
|
// this one (best effort — it may already be dead)
|
|
// and hand the session over.
|
|
close_event(socket);
|
|
return Step::Attach { hs, inbound };
|
|
}
|
|
Ok(None) => {}
|
|
// Registry gone = shutdown chain: terminate and exit.
|
|
Err(_) => return Step::Exit,
|
|
},
|
|
_ => unreachable!("three arms"),
|
|
}
|
|
}
|
|
}
|
|
|
|
fn detached_loop<P: Encode + Decode + Send + Sync + 'static>(
|
|
ctl: &Receiver<Ctl<P>>,
|
|
bus_rx: &Receiver<Arc<Broadcast<P>>>,
|
|
buffer: &mut VecDeque<Arc<Broadcast<P>>>,
|
|
cap: usize,
|
|
ttl: Duration,
|
|
) -> Step<P> {
|
|
// TTL is per detach episode: the clock starts now, every time.
|
|
let deadline = Instant::now() + ttl;
|
|
loop {
|
|
// cap.max(1): a cap of 0 means "no buffering tolerated" — the
|
|
// first buffered broadcast tears the session down; an empty
|
|
// buffer still waits out the TTL.
|
|
if buffer.len() >= cap.max(1) {
|
|
return Step::Exit;
|
|
}
|
|
let remaining = deadline.saturating_duration_since(Instant::now());
|
|
match smarm::select_timeout(&[ctl, bus_rx], remaining) {
|
|
// TTL expired with no rejoin: session over.
|
|
None => return Step::Exit,
|
|
Some(0) => match ctl.try_recv() {
|
|
Ok(Some(Ctl::Attach { hs, inbound })) => return Step::Attach { hs, inbound },
|
|
Ok(None) => {}
|
|
Err(_) => return Step::Exit,
|
|
},
|
|
Some(1) => match bus_rx.try_recv() {
|
|
Ok(Some(b)) => buffer.push_back(b),
|
|
Ok(None) => {}
|
|
Err(_) => return Step::Exit,
|
|
},
|
|
Some(_) => unreachable!("two arms"),
|
|
}
|
|
}
|
|
}
|
|
|
|
/// Best-effort `phx_close` push for the socket's current join
|
|
/// generation.
|
|
fn close_event<P: Encode + Decode + Send + Sync + 'static>(socket: &ChannelSocket<P>) {
|
|
let _ = socket.ws.send(P::encode(FrameRef {
|
|
join_ref: socket.join_ref.as_deref(),
|
|
reference: None,
|
|
topic: &socket.topic,
|
|
message: MessageRef::Event { event: EV_CLOSE, payload: None },
|
|
}));
|
|
}
|
|
|
|
// ---------------------------------------------------------------------------
|
|
// Tests — runtime-backed, on the toy pipe codec from testkit. The
|
|
// session-config types (`Cfg*`) are deliberately separate from the
|
|
// channel they configure: `channel_session::<S>` only needs `S` for
|
|
// keying/cap/ttl, the factory builds the channel.
|
|
// ---------------------------------------------------------------------------
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::super::testkit::{eventually, Harness, TP};
|
|
use super::super::*;
|
|
use std::sync::atomic::{AtomicBool, Ordering};
|
|
use std::sync::Mutex;
|
|
use std::time::Duration;
|
|
|
|
/// Counts its joins (cold start replies `joined:1`); flags
|
|
/// terminate; broadcasts on "shout".
|
|
struct Counter {
|
|
joins: u32,
|
|
term: Arc<AtomicBool>,
|
|
reject_rejoin: bool,
|
|
}
|
|
|
|
impl Channel<TP> for Counter {
|
|
fn join(&mut self, _topic: &str, _payload: TP, _s: &ChannelSocket<TP>) -> Result<TP, TP> {
|
|
self.joins += 1;
|
|
if self.reject_rejoin && self.joins > 1 {
|
|
return Err(TP("not again".into()));
|
|
}
|
|
Ok(TP(format!("joined:{}", self.joins)))
|
|
}
|
|
fn handle_in(&mut self, event: &str, payload: TP, s: &ChannelSocket<TP>) {
|
|
match event {
|
|
"shout" => s.broadcast("news", payload),
|
|
_ => s.reply(Status::Ok, TP(format!("echo:{}", payload.0))),
|
|
}
|
|
}
|
|
fn terminate(&mut self) {
|
|
self.term.store(true, Ordering::SeqCst);
|
|
}
|
|
}
|
|
|
|
fn counter_factory(term: &Arc<AtomicBool>, reject_rejoin: bool) -> impl ChannelFactory<TP> {
|
|
let term = term.clone();
|
|
move |_t: &str| -> Box<dyn Channel<TP>> {
|
|
Box::new(Counter { joins: 0, term: term.clone(), reject_rejoin })
|
|
}
|
|
}
|
|
|
|
/// Session config: keyed by topic alone, defaults otherwise.
|
|
struct Cfg;
|
|
impl ChannelSession<TP> for Cfg {
|
|
type Key = String;
|
|
fn session_key(topic: &str, _p: &TP) -> String {
|
|
topic.to_owned()
|
|
}
|
|
}
|
|
|
|
/// Tiny TTL (50ms) for the expiry test.
|
|
struct CfgTtl;
|
|
impl ChannelSession<TP> for CfgTtl {
|
|
type Key = String;
|
|
fn session_key(topic: &str, _p: &TP) -> String {
|
|
topic.to_owned()
|
|
}
|
|
fn ttl() -> Duration {
|
|
Duration::from_millis(50)
|
|
}
|
|
}
|
|
|
|
/// Cap of 2 (and a long TTL so only the cap can fire).
|
|
struct CfgCap;
|
|
impl ChannelSession<TP> for CfgCap {
|
|
type Key = String;
|
|
fn session_key(topic: &str, _p: &TP) -> String {
|
|
topic.to_owned()
|
|
}
|
|
fn buffer_cap() -> usize {
|
|
2
|
|
}
|
|
fn ttl() -> Duration {
|
|
Duration::from_secs(30)
|
|
}
|
|
}
|
|
|
|
fn session_hub<S: ChannelSession<TP, Key = String>>(
|
|
term: &Arc<AtomicBool>,
|
|
reject_rejoin: bool,
|
|
) -> (ChannelHub<TP>, Vec<smarm::ChildSpec>) {
|
|
ChannelHub::new(
|
|
"sess-unit-bus",
|
|
PrefixRouter::new().channel_session::<S>(
|
|
"room:*",
|
|
"sess-unit-registry",
|
|
counter_factory(term, reject_rejoin),
|
|
),
|
|
)
|
|
}
|
|
|
|
#[test]
|
|
fn session_buffers_across_reconnect_and_drains_in_order() {
|
|
let out = Arc::new(Mutex::new(Vec::<String>::new()));
|
|
let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false)));
|
|
smarm::run(move || {
|
|
let (hub, children) = session_hub::<Cfg>(&term, false);
|
|
let _sup = crate::channels::tests::start_hub(&hub, children);
|
|
let mut h1 = Harness::new(&hub);
|
|
h1.send("j1|r1|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h1.recv());
|
|
|
|
// Transport dies. Whether the actor has observed the detach
|
|
// yet or not, both paths (closed inbound arm, failed ws
|
|
// send) land the broadcasts in the buffer.
|
|
drop(h1);
|
|
hub.broadcast("room:a", "news", TP("one".into()));
|
|
hub.broadcast("room:a", "news", TP("two".into()));
|
|
|
|
let mut h2 = Harness::new(&hub);
|
|
h2.send("j2|r2|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h2.recv()); // warm rejoin ack
|
|
out2.lock().unwrap().push(h2.recv()); // drained, in order
|
|
out2.lock().unwrap().push(h2.recv());
|
|
});
|
|
let out = out.lock().unwrap();
|
|
assert_eq!(out[0], "j1|r1|room:a|reply|ok|joined:1");
|
|
assert_eq!(out[1], "j2|r2|room:a|reply|ok|joined:2");
|
|
assert_eq!(out[2], "j2|-|room:a|event|news|one");
|
|
assert_eq!(out[3], "j2|-|room:a|event|news|two");
|
|
}
|
|
|
|
#[test]
|
|
fn session_ttl_expiry_tears_down_then_cold_start() {
|
|
let out = Arc::new(Mutex::new(Vec::<String>::new()));
|
|
let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false)));
|
|
let term2 = term.clone();
|
|
smarm::run(move || {
|
|
let (hub, children) = session_hub::<CfgTtl>(&term2, false);
|
|
let _sup = crate::channels::tests::start_hub(&hub, children);
|
|
let mut h1 = Harness::new(&hub);
|
|
h1.send("j1|r1|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h1.recv());
|
|
drop(h1);
|
|
assert!(eventually(|| term2.load(Ordering::SeqCst)), "ttl never fired");
|
|
|
|
let mut h2 = Harness::new(&hub);
|
|
h2.send("j2|r2|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h2.recv());
|
|
});
|
|
let out = out.lock().unwrap();
|
|
assert_eq!(out[0], "j1|r1|room:a|reply|ok|joined:1");
|
|
assert_eq!(out[1], "j2|r2|room:a|reply|ok|joined:1"); // cold
|
|
}
|
|
|
|
#[test]
|
|
fn session_buffer_cap_tears_down_then_cold_start() {
|
|
let out = Arc::new(Mutex::new(Vec::<String>::new()));
|
|
let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false)));
|
|
smarm::run(move || {
|
|
let (hub, children) = session_hub::<CfgCap>(&term, false);
|
|
let _sup = crate::channels::tests::start_hub(&hub, children);
|
|
let mut h1 = Harness::new(&hub);
|
|
h1.send("j1|r1|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h1.recv());
|
|
drop(h1);
|
|
// Cap is 2: the second buffered broadcast tears it down
|
|
// (ttl is 30s — only the cap can fire here).
|
|
hub.broadcast("room:a", "news", TP("one".into()));
|
|
hub.broadcast("room:a", "news", TP("two".into()));
|
|
assert!(eventually(|| term.load(Ordering::SeqCst)), "cap never fired");
|
|
|
|
let mut h2 = Harness::new(&hub);
|
|
h2.send("j2|r2|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h2.recv());
|
|
});
|
|
let out = out.lock().unwrap();
|
|
assert_eq!(out[0], "j1|r1|room:a|reply|ok|joined:1");
|
|
assert_eq!(out[1], "j2|r2|room:a|reply|ok|joined:1"); // cold
|
|
}
|
|
|
|
#[test]
|
|
fn leave_ends_the_session_not_just_the_attachment() {
|
|
let out = Arc::new(Mutex::new(Vec::<String>::new()));
|
|
let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false)));
|
|
smarm::run(move || {
|
|
let (hub, children) = session_hub::<Cfg>(&term, false);
|
|
let _sup = crate::channels::tests::start_hub(&hub, children);
|
|
let mut h = Harness::new(&hub);
|
|
h.send("j1|r1|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h.recv());
|
|
h.send("j1|r2|room:a|event|phx_leave|");
|
|
out2.lock().unwrap().push(h.recv()); // ok reply
|
|
out2.lock().unwrap().push(h.recv()); // phx_close
|
|
assert!(eventually(|| term.load(Ordering::SeqCst)), "leave never terminated");
|
|
|
|
h.send("j2|r3|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h.recv());
|
|
});
|
|
let out = out.lock().unwrap();
|
|
assert_eq!(out[0], "j1|r1|room:a|reply|ok|joined:1");
|
|
assert_eq!(out[1], "j1|r2|room:a|reply|ok|");
|
|
assert_eq!(out[2], "j1|-|room:a|event|phx_close|");
|
|
assert_eq!(out[3], "j2|r3|room:a|reply|ok|joined:1"); // cold
|
|
}
|
|
|
|
#[test]
|
|
fn second_transport_evicts_first() {
|
|
let out = Arc::new(Mutex::new(Vec::<String>::new()));
|
|
let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false)));
|
|
smarm::run(move || {
|
|
let (hub, children) = session_hub::<Cfg>(&term, false);
|
|
let _sup = crate::channels::tests::start_hub(&hub, children);
|
|
let mut h1 = Harness::new(&hub);
|
|
h1.send("j1|r1|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h1.recv());
|
|
|
|
let mut h2 = Harness::new(&hub);
|
|
h2.send("j2|r2|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h2.recv()); // warm handover ack
|
|
out2.lock().unwrap().push(h1.recv()); // eviction phx_close
|
|
|
|
hub.broadcast("room:a", "news", TP("x".into()));
|
|
out2.lock().unwrap().push(h2.recv()); // only h2 is attached
|
|
h1.assert_silence();
|
|
|
|
// h1's conn-side entry is stale: its inbound receiver was
|
|
// dropped on eviction, so the next event self-prunes with
|
|
// an error reply.
|
|
h1.send("j1|r3|room:a|event|ping|x");
|
|
out2.lock().unwrap().push(h1.recv());
|
|
});
|
|
let out = out.lock().unwrap();
|
|
assert_eq!(out[0], "j1|r1|room:a|reply|ok|joined:1");
|
|
assert_eq!(out[1], "j2|r2|room:a|reply|ok|joined:2");
|
|
assert_eq!(out[2], "j1|-|room:a|event|phx_close|");
|
|
assert_eq!(out[3], "j2|-|room:a|event|news|x");
|
|
assert_eq!(out[4], "j1|r3|room:a|reply|error|");
|
|
}
|
|
|
|
#[test]
|
|
fn rejected_rejoin_ends_the_session() {
|
|
let out = Arc::new(Mutex::new(Vec::<String>::new()));
|
|
let (out2, term) = (out.clone(), Arc::new(AtomicBool::new(false)));
|
|
smarm::run(move || {
|
|
let (hub, children) = session_hub::<Cfg>(&term, true);
|
|
let _sup = crate::channels::tests::start_hub(&hub, children);
|
|
let mut h1 = Harness::new(&hub);
|
|
h1.send("j1|r1|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h1.recv());
|
|
drop(h1);
|
|
|
|
let mut h2 = Harness::new(&hub);
|
|
h2.send("j2|r2|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h2.recv()); // rejoin rejected
|
|
assert!(eventually(|| term.load(Ordering::SeqCst)), "reject never terminated");
|
|
|
|
let mut h3 = Harness::new(&hub);
|
|
h3.send("j3|r3|room:a|event|phx_join|u");
|
|
out2.lock().unwrap().push(h3.recv()); // cold start
|
|
});
|
|
let out = out.lock().unwrap();
|
|
assert_eq!(out[0], "j1|r1|room:a|reply|ok|joined:1");
|
|
assert_eq!(out[1], "j2|r2|room:a|reply|error|not again");
|
|
assert_eq!(out[2], "j3|r3|room:a|reply|ok|joined:1");
|
|
}
|
|
}
|