#![allow(dead_code)]
use std::{
collections::HashMap,
sync::{
atomic::{AtomicU64, AtomicUsize, Ordering},
Arc, Mutex, Weak,
},
time::{Duration, Instant},
};
use serde_json::{json, Value};
use tokio::sync::{broadcast, mpsc, oneshot};
const DEFAULT_IDLE_WINDOW: Duration = Duration::from_secs(5 * 60);
const DEFAULT_PRESENCE_TTL: Duration = Duration::from_secs(30);
const DEFAULT_SWEEP_INTERVAL: Duration = Duration::from_secs(10);
const BROADCAST_CAPACITY: usize = 256;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum RealtimeChannel {
Events,
Presence,
DocSync,
}
#[derive(Debug, Clone)]
pub enum Frame {
Event(Value),
Presence(Value),
DocSync(Vec<u8>),
}
impl Frame {
pub fn channel(&self) -> RealtimeChannel {
match self {
Frame::Event(_) => RealtimeChannel::Events,
Frame::Presence(_) => RealtimeChannel::Presence,
Frame::DocSync(_) => RealtimeChannel::DocSync,
}
}
}
const EVENT_NAME_KEY: &str = "__ryu_event";
const EVENT_DATA_KEY: &str = "data";
static NEXT_CONN_ID: AtomicU64 = AtomicU64::new(1);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct ConnId(u64);
impl ConnId {
fn next() -> Self {
Self(NEXT_CONN_ID.fetch_add(1, Ordering::Relaxed))
}
pub fn get(self) -> u64 {
self.0
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Event {
pub name: String,
pub payload: Value,
}
impl Event {
pub fn decode(frame: &Frame) -> Option<Event> {
let Frame::Event(value) = frame else {
return None;
};
let name = value.get(EVENT_NAME_KEY)?.as_str()?.to_string();
let payload = value.get(EVENT_DATA_KEY).cloned().unwrap_or(Value::Null);
Some(Event { name, payload })
}
}
fn encode_event(name: impl Into<String>, payload: Value) -> Value {
let mut map = serde_json::Map::with_capacity(2);
map.insert(EVENT_NAME_KEY.to_string(), Value::String(name.into()));
map.insert(EVENT_DATA_KEY.to_string(), payload);
Value::Object(map)
}
#[derive(Debug, Clone, Copy)]
pub struct RoomConfig {
pub idle_window: Duration,
pub presence_ttl: Duration,
pub sweep_interval: Duration,
}
impl Default for RoomConfig {
fn default() -> Self {
Self {
idle_window: DEFAULT_IDLE_WINDOW,
presence_ttl: DEFAULT_PRESENCE_TTL,
sweep_interval: DEFAULT_SWEEP_INTERVAL,
}
}
}
enum RoomCommand {
Joined,
Left { member_id: String },
Presence { member_id: String, value: Value },
PresenceMembers { reply: oneshot::Sender<Vec<String>> },
OpenConn {
conn_id: ConnId,
tx: mpsc::UnboundedSender<Frame>,
},
CloseConn { conn_id: ConnId },
SendTo { conn_id: ConnId, frame: Frame },
ConnCount { reply: oneshot::Sender<usize> },
}
type RoomMap = HashMap<String, RoomHandle>;
#[derive(Clone)]
pub struct RoomRegistry {
inner: Arc<Mutex<RoomMap>>,
config: RoomConfig,
}
impl RoomRegistry {
pub fn new() -> Self {
Self::with_config(RoomConfig::default())
}
pub fn with_config(config: RoomConfig) -> Self {
Self {
inner: Arc::new(Mutex::new(HashMap::new())),
config,
}
}
pub fn get_or_create(&self, room_id: &str) -> RoomHandle {
let mut map = self.lock();
if let Some(handle) = map.get(room_id) {
return handle.clone();
}
let handle = self.spawn_room(room_id.to_string());
map.insert(room_id.to_string(), handle.clone());
handle
}
pub fn join(&self, room_id: &str, member_id: impl Into<String>) -> RoomMembership {
let mut map = self.lock();
let handle = match map.get(room_id) {
Some(handle) => handle.clone(),
None => {
let handle = self.spawn_room(room_id.to_string());
map.insert(room_id.to_string(), handle.clone());
handle
}
};
handle.members.fetch_add(1, Ordering::SeqCst);
drop(map);
let _ = handle.cmd.send(RoomCommand::Joined);
RoomMembership {
handle,
member_id: member_id.into(),
left: false,
}
}
pub fn publish_event(&self, room_id: &str, value: Value) {
if let Some(handle) = self.lock().get(room_id) {
let _ = handle.broadcast.send(Frame::Event(value));
}
}
pub fn publish_presence(&self, room_id: &str, member_id: &str, value: Value) {
if let Some(handle) = self.lock().get(room_id) {
handle.publish_presence(member_id, value);
}
}
pub fn broadcast_event(&self, room_id: &str, name: impl Into<String>, payload: Value) {
if let Some(handle) = self.lock().get(room_id) {
handle.broadcast_event(name, payload);
}
}
pub fn send_event(&self, room_id: &str, conn: ConnId, name: impl Into<String>, payload: Value) {
if let Some(handle) = self.lock().get(room_id) {
handle.send_event(conn, name, payload);
}
}
pub fn room_count(&self) -> usize {
self.lock().len()
}
fn spawn_room(&self, room_id: String) -> RoomHandle {
let (broadcast_tx, _rx) = broadcast::channel(BROADCAST_CAPACITY);
let (cmd_tx, cmd_rx) = mpsc::unbounded_channel();
let members = Arc::new(AtomicUsize::new(0));
let handle = RoomHandle {
room_id: room_id.clone(),
broadcast: broadcast_tx.clone(),
cmd: cmd_tx,
members: Arc::clone(&members),
};
let registry = Arc::downgrade(&self.inner);
let config = self.config;
tokio::spawn(run_room(
room_id,
members,
broadcast_tx,
cmd_rx,
registry,
config,
));
handle
}
fn lock(&self) -> std::sync::MutexGuard<'_, RoomMap> {
self.inner.lock().unwrap_or_else(|e| e.into_inner())
}
}
impl Default for RoomRegistry {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone)]
pub struct RoomHandle {
room_id: String,
broadcast: broadcast::Sender<Frame>,
cmd: mpsc::UnboundedSender<RoomCommand>,
members: Arc<AtomicUsize>,
}
impl RoomHandle {
pub fn room_id(&self) -> &str {
&self.room_id
}
pub fn subscribe(&self) -> broadcast::Receiver<Frame> {
self.broadcast.subscribe()
}
pub fn member_count(&self) -> usize {
self.members.load(Ordering::SeqCst)
}
pub fn join(&self, member_id: impl Into<String>) -> RoomMembership {
self.members.fetch_add(1, Ordering::SeqCst);
let _ = self.cmd.send(RoomCommand::Joined);
RoomMembership {
handle: self.clone(),
member_id: member_id.into(),
left: false,
}
}
pub fn publish_event(&self, value: Value) {
let _ = self.broadcast.send(Frame::Event(value));
}
pub fn publish_presence(&self, member_id: &str, value: Value) {
let _ = self.cmd.send(RoomCommand::Presence {
member_id: member_id.to_string(),
value,
});
}
pub fn publish_doc_sync(&self, bytes: Vec<u8>) {
let _ = self.broadcast.send(Frame::DocSync(bytes));
}
pub fn broadcast_event(&self, name: impl Into<String>, payload: Value) {
let _ = self
.broadcast
.send(Frame::Event(encode_event(name, payload)));
}
pub fn send_event(&self, conn: ConnId, name: impl Into<String>, payload: Value) {
let _ = self.cmd.send(RoomCommand::SendTo {
conn_id: conn,
frame: Frame::Event(encode_event(name, payload)),
});
}
pub fn open_connection(&self) -> Connection {
let conn_id = ConnId::next();
let (tx, targeted_rx) = mpsc::unbounded_channel();
let _ = self.cmd.send(RoomCommand::OpenConn { conn_id, tx });
Connection {
conn_id,
cmd: self.cmd.clone(),
broadcast_rx: self.broadcast.subscribe(),
targeted_rx,
broadcast_open: true,
targeted_open: true,
}
}
pub async fn conn_count(&self) -> usize {
let (reply, rx) = oneshot::channel();
if self.cmd.send(RoomCommand::ConnCount { reply }).is_err() {
return 0;
}
rx.await.unwrap_or(0)
}
pub async fn presence_members(&self) -> Vec<String> {
let (reply, rx) = oneshot::channel();
if self
.cmd
.send(RoomCommand::PresenceMembers { reply })
.is_err()
{
return Vec::new();
}
rx.await.unwrap_or_default()
}
}
pub struct RoomMembership {
handle: RoomHandle,
member_id: String,
left: bool,
}
impl RoomMembership {
pub fn member_id(&self) -> &str {
&self.member_id
}
pub fn handle(&self) -> &RoomHandle {
&self.handle
}
pub fn subscribe(&self) -> broadcast::Receiver<Frame> {
self.handle.subscribe()
}
pub fn publish_presence(&self, value: Value) {
self.handle.publish_presence(&self.member_id, value);
}
pub fn open_connection(&self) -> Connection {
self.handle.open_connection()
}
pub fn leave(&mut self) {
if self.left {
return;
}
self.left = true;
self.handle.members.fetch_sub(1, Ordering::SeqCst);
let _ = self.handle.cmd.send(RoomCommand::Left {
member_id: self.member_id.clone(),
});
}
}
impl Drop for RoomMembership {
fn drop(&mut self) {
self.leave();
}
}
pub struct Connection {
conn_id: ConnId,
cmd: mpsc::UnboundedSender<RoomCommand>,
broadcast_rx: broadcast::Receiver<Frame>,
targeted_rx: mpsc::UnboundedReceiver<Frame>,
broadcast_open: bool,
targeted_open: bool,
}
impl Connection {
pub fn id(&self) -> ConnId {
self.conn_id
}
pub async fn recv(&mut self) -> Option<Event> {
loop {
if !self.broadcast_open && !self.targeted_open {
return None;
}
let frame = tokio::select! {
biased;
targeted = self.targeted_rx.recv(), if self.targeted_open => match targeted {
Some(frame) => frame,
None => {
self.targeted_open = false;
continue;
}
},
broadcast = self.broadcast_rx.recv(), if self.broadcast_open => match broadcast {
Ok(frame) => frame,
Err(broadcast::error::RecvError::Lagged(_)) => continue,
Err(broadcast::error::RecvError::Closed) => {
self.broadcast_open = false;
continue;
}
},
};
if let Some(event) = Event::decode(&frame) {
return Some(event);
}
}
}
}
impl Drop for Connection {
fn drop(&mut self) {
let _ = self.cmd.send(RoomCommand::CloseConn {
conn_id: self.conn_id,
});
}
}
async fn run_room(
room_id: String,
members: Arc<AtomicUsize>,
broadcast_tx: broadcast::Sender<Frame>,
mut cmd_rx: mpsc::UnboundedReceiver<RoomCommand>,
registry: Weak<Mutex<RoomMap>>,
config: RoomConfig,
) {
let mut presence: HashMap<String, (Value, Instant)> = HashMap::new();
let mut conns: HashMap<ConnId, mpsc::UnboundedSender<Frame>> = HashMap::new();
let mut empty_since: Option<Instant> = Some(Instant::now());
let mut sweep = tokio::time::interval(config.sweep_interval);
sweep.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
sweep.tick().await;
loop {
tokio::select! {
cmd = cmd_rx.recv() => {
match cmd {
None => {
evict(®istry, &room_id, "all handles dropped");
return;
}
Some(RoomCommand::Joined) => {
empty_since = None;
}
Some(RoomCommand::Left { member_id }) => {
if presence.remove(&member_id).is_some() {
let _ = broadcast_tx.send(Frame::Presence(presence_leave(&member_id)));
}
if members.load(Ordering::SeqCst) == 0 {
empty_since = Some(Instant::now());
}
}
Some(RoomCommand::Presence { member_id, value }) => {
presence.insert(member_id, (value.clone(), Instant::now()));
let _ = broadcast_tx.send(Frame::Presence(value));
}
Some(RoomCommand::PresenceMembers { reply }) => {
let mut ids: Vec<String> = presence.keys().cloned().collect();
ids.sort();
let _ = reply.send(ids);
}
Some(RoomCommand::OpenConn { conn_id, tx }) => {
conns.insert(conn_id, tx);
}
Some(RoomCommand::CloseConn { conn_id }) => {
conns.remove(&conn_id);
}
Some(RoomCommand::SendTo { conn_id, frame }) => {
if let Some(tx) = conns.get(&conn_id) {
if tx.send(frame).is_err() {
conns.remove(&conn_id);
}
}
}
Some(RoomCommand::ConnCount { reply }) => {
let _ = reply.send(conns.len());
}
}
}
_ = sweep.tick() => {
let ttl = config.presence_ttl;
let stale: Vec<String> = presence
.iter()
.filter(|(_, (_, seen))| seen.elapsed() >= ttl)
.map(|(id, _)| id.clone())
.collect();
for id in stale {
presence.remove(&id);
let _ = broadcast_tx.send(Frame::Presence(presence_leave(&id)));
}
if let Some(since) = empty_since {
if since.elapsed() >= config.idle_window
&& try_evict(®istry, &room_id, &members)
{
return;
}
}
}
}
}
}
fn presence_leave(member_id: &str) -> Value {
json!({ "type": "presence_leave", "member_id": member_id })
}
fn try_evict(registry: &Weak<Mutex<RoomMap>>, room_id: &str, members: &Arc<AtomicUsize>) -> bool {
let Some(map) = registry.upgrade() else {
return true;
};
let mut map = map.lock().unwrap_or_else(|e| e.into_inner());
if members.load(Ordering::SeqCst) != 0 {
return false;
}
map.remove(room_id);
tracing::info!(room_id, "realtime: hibernating idle room (0 members)");
true
}
fn evict(registry: &Weak<Mutex<RoomMap>>, room_id: &str, reason: &str) {
if let Some(map) = registry.upgrade() {
map.lock()
.unwrap_or_else(|e| e.into_inner())
.remove(room_id);
}
tracing::info!(room_id, reason, "realtime: evicting room");
}
#[cfg(test)]
mod tests {
use super::*;
fn fast_config() -> RoomConfig {
RoomConfig {
idle_window: Duration::from_millis(80),
presence_ttl: Duration::from_millis(100),
sweep_interval: Duration::from_millis(20),
}
}
#[tokio::test]
async fn get_or_create_is_idempotent() {
let reg = RoomRegistry::new();
let a = reg.get_or_create("room-1");
let b = reg.get_or_create("room-1");
let mut rx = b.subscribe();
a.publish_event(json!({"n": 1}));
let frame = rx.recv().await.expect("frame");
match frame {
Frame::Event(v) => assert_eq!(v["n"], 1),
_ => panic!("expected event frame"),
}
assert_eq!(reg.room_count(), 1, "one logical room");
}
#[tokio::test]
async fn join_leave_member_counting() {
let reg = RoomRegistry::new();
let handle = reg.get_or_create("room-2");
assert_eq!(handle.member_count(), 0);
let m1 = handle.join("alice");
assert_eq!(handle.member_count(), 1);
let m2 = handle.join("bob");
assert_eq!(handle.member_count(), 2);
drop(m2);
assert_eq!(handle.member_count(), 1);
let mut m1 = m1;
m1.leave();
assert_eq!(handle.member_count(), 0);
drop(m1);
assert_eq!(handle.member_count(), 0);
}
#[tokio::test]
async fn registry_join_counts_and_recreates() {
let reg = RoomRegistry::new();
let m1 = reg.join("room-j", "alice");
assert_eq!(reg.room_count(), 1);
assert_eq!(m1.handle().member_count(), 1);
let m2 = reg.join("room-j", "bob");
assert_eq!(m2.handle().member_count(), 2);
drop(m1);
drop(m2);
let m3 = reg.join("room-j", "carol");
assert_eq!(m3.handle().member_count(), 1);
assert_eq!(reg.room_count(), 1);
}
#[tokio::test]
async fn published_event_reaches_subscriber() {
let reg = RoomRegistry::new();
let handle = reg.get_or_create("room-3");
let _member = handle.join("alice");
let mut rx = handle.subscribe();
reg.publish_event("room-3", json!({"type": "message", "id": "m1"}));
let frame = rx.recv().await.expect("frame");
match frame {
Frame::Event(v) => {
assert_eq!(v["type"], "message");
assert_eq!(v["id"], "m1");
}
_ => panic!("expected event frame"),
}
assert_eq!(handle.subscribe().len(), 0, "fresh receiver has no backlog");
}
#[tokio::test]
async fn presence_delta_is_broadcast() {
let reg = RoomRegistry::new();
let handle = reg.get_or_create("room-4");
let member = handle.join("alice");
let mut rx = handle.subscribe();
member.publish_presence(json!({"member_id": "alice", "cursor": [1, 2]}));
let frame = rx.recv().await.expect("frame");
match frame {
Frame::Presence(v) => assert_eq!(v["cursor"][0], 1),
_ => panic!("expected presence frame"),
}
}
#[tokio::test]
async fn presence_ttl_is_reaped() {
let reg = RoomRegistry::with_config(fast_config());
let handle = reg.get_or_create("room-5");
let member = handle.join("alice");
let mut rx = handle.subscribe();
member.publish_presence(json!({"member_id": "alice"}));
let _ = rx.recv().await.expect("upsert");
assert_eq!(handle.presence_members().await, vec!["alice".to_string()]);
tokio::time::sleep(Duration::from_millis(220)).await;
assert!(
handle.presence_members().await.is_empty(),
"stale presence should be reaped"
);
let mut saw_leave = false;
while let Ok(frame) = rx.try_recv() {
if let Frame::Presence(v) = frame {
if v["type"] == "presence_leave" {
saw_leave = true;
}
}
}
assert!(saw_leave, "expected a presence_leave delta on reap");
drop(member);
}
#[tokio::test]
async fn idle_room_hibernates() {
let reg = RoomRegistry::with_config(fast_config());
let handle = reg.get_or_create("room-6");
{
let _m = handle.join("alice");
assert_eq!(reg.room_count(), 1);
}
tokio::time::sleep(Duration::from_millis(220)).await;
assert_eq!(reg.room_count(), 0, "idle room should hibernate");
let handle2 = reg.get_or_create("room-6");
let _m2 = handle2.join("bob");
assert_eq!(reg.room_count(), 1, "room rehydrates on next join");
}
#[tokio::test]
async fn publish_event_to_absent_room_is_noop() {
let reg = RoomRegistry::new();
reg.publish_event("ghost", json!({"x": 1}));
assert_eq!(reg.room_count(), 0);
}
#[test]
fn frame_channel_tags() {
assert_eq!(Frame::Event(json!({})).channel(), RealtimeChannel::Events);
assert_eq!(
Frame::Presence(json!({})).channel(),
RealtimeChannel::Presence
);
assert_eq!(
Frame::DocSync(vec![1, 2, 3]).channel(),
RealtimeChannel::DocSync
);
}
#[test]
fn event_decode_only_matches_the_envelope() {
let frame = Frame::Event(encode_event("chat.message", json!({"id": "m1"})));
let ev = Event::decode(&frame).expect("named event");
assert_eq!(ev.name, "chat.message");
assert_eq!(ev.payload["id"], "m1");
assert!(Event::decode(&Frame::Event(json!({"id": "raw"}))).is_none());
assert!(Event::decode(&Frame::Presence(json!({"cursor": [1, 2]}))).is_none());
assert!(Event::decode(&Frame::DocSync(vec![1, 2, 3])).is_none());
}
#[tokio::test]
async fn broadcast_event_reaches_typed_and_raw_subscribers() {
let reg = RoomRegistry::new();
let handle = reg.get_or_create("evt-room");
let mut conn = handle.open_connection();
let mut raw = handle.subscribe();
reg.broadcast_event("evt-room", "counter.tick", json!({"n": 7}));
let ev = conn.recv().await.expect("typed event");
assert_eq!(ev.name, "counter.tick");
assert_eq!(ev.payload["n"], 7);
match raw.recv().await.expect("raw frame") {
Frame::Event(v) => {
let decoded = Event::decode(&Frame::Event(v)).expect("envelope");
assert_eq!(decoded.name, "counter.tick");
}
other => panic!("expected Frame::Event, got {other:?}"),
}
}
#[tokio::test]
async fn plugin_contributions_broadcast_payload_is_self_describing() {
let reg = RoomRegistry::new();
let handle = reg.get_or_create("system:plugins");
let mut raw = handle.subscribe();
reg.broadcast_event(
"system:plugins",
"plugin.contributions.changed",
json!({"type": "contributions_changed"}),
);
match raw.recv().await.expect("raw frame") {
frame @ Frame::Event(_) => {
let ev = Event::decode(&frame).expect("envelope");
assert_eq!(ev.name, "plugin.contributions.changed");
assert_eq!(ev.payload["type"], "contributions_changed");
}
other => panic!("expected Frame::Event, got {other:?}"),
}
}
#[tokio::test]
async fn send_event_is_isolated_to_its_connection() {
let reg = RoomRegistry::new();
let handle = reg.get_or_create("target-room");
let mut conn_a = handle.open_connection();
let mut conn_b = handle.open_connection();
let mut raw = handle.subscribe();
let a_id = conn_a.id();
handle.send_event(a_id, "secret", json!({"for": "a"}));
assert_eq!(handle.conn_count().await, 2);
handle.broadcast_event("marker", json!({}));
let first = conn_a.recv().await.expect("a first");
assert_eq!(first.name, "secret");
assert_eq!(first.payload["for"], "a");
let second = conn_a.recv().await.expect("a second");
assert_eq!(second.name, "marker");
let b_first = conn_b.recv().await.expect("b first");
assert_eq!(b_first.name, "marker");
match raw.recv().await.expect("raw first") {
Frame::Event(v) => {
assert_eq!(v[EVENT_NAME_KEY], "marker");
}
other => panic!("expected marker frame, got {other:?}"),
}
assert!(raw.try_recv().is_err(), "raw saw exactly one frame");
}
#[tokio::test]
async fn dropping_a_connection_prunes_it_from_the_actor() {
let reg = RoomRegistry::new();
let handle = reg.get_or_create("drop-room");
let conn_a = handle.open_connection();
let conn_b = handle.open_connection();
assert_eq!(handle.conn_count().await, 2);
let a_id = conn_a.id();
drop(conn_a);
assert_eq!(handle.conn_count().await, 1);
handle.send_event(a_id, "ghost", json!({}));
handle.broadcast_event("alive", json!({}));
let mut conn_b = conn_b;
assert_eq!(conn_b.recv().await.expect("b").name, "alive");
}
#[tokio::test]
async fn typed_reader_skips_non_event_frames() {
let reg = RoomRegistry::new();
let handle = reg.get_or_create("skip-room");
let mut conn = handle.open_connection();
handle.publish_event(json!({"legacy": true}));
handle.broadcast_event("real", json!({"ok": 1}));
let ev = conn.recv().await.expect("named event");
assert_eq!(ev.name, "real");
assert_eq!(ev.payload["ok"], 1);
}
#[test]
fn conn_ids_are_process_unique_and_monotonic() {
let a = ConnId::next();
let b = ConnId::next();
assert_ne!(a, b);
assert!(b.get() > a.get());
}
}