Session Server (Rust)
Overview
The JellyWatchParty session server is an asynchronous Rust application using Axum for WebSocket handling and Tokio as the async runtime. It manages rooms, clients, and playback synchronization in memory.
Module Structure
src/
├── main.rs # Entry point, listener + graceful shutdown
├── types.rs # Data structures
├── routes.rs # Axum router, origin guard, CORS
├── tasks.rs # Background tasks (zombie cleanup, shutdown)
├── messaging.rs # Message sending functions
├── auth.rs # JWT authentication (optional)
├── utils.rs # Utilities (timestamp, random tokens)
├── password.rs # Room password hashing, constant-time compare
├── admin/ # Admin panel (second listener, own port)
│ ├── mod.rs # Router, session/CSRF middleware, security headers
│ ├── config.rs # ADMIN_* environment variables
│ ├── auth.rs # Login, sessions, failed-login throttle
│ ├── api.rs # JSON API (overview, rooms, members, host)
│ ├── integrations.rs # Bot settings, link codes, activity log
│ ├── ui.rs # Serves the embedded UI
│ └── ui/ # index.html, app.js, app.css (no build step)
├── events.rs # "Rooms changed" counter (wakes the integration long poll)
├── integration/ # Chat integrations (Discord bot sidecar), see integration-api.md
│ ├── mod.rs # Integration state, user cache, guards, audit, reaper
│ ├── config.rs # DATA_DIR, INTEGRATION_*, *_INTEGRATION_TOKEN
│ ├── store.rs # integrations.json: settings, code HMACs, links
│ ├── actions.rs # Who may do what; chat room operations
│ ├── view.rs # Room list for sidecars
│ └── api.rs # Token-protected HTTP API (own listener)
├── jellyfin/ # Admin panel device bridge (Jellyfin REST API)
│ ├── mod.rs # JELLYFIN_* configuration
│ ├── api.rs # /Sessions, Playing/{cmd}, PlayNow (reqwest)
│ ├── logic.rs # Pure host/receiver decisions + position estimate
│ ├── bridge.rs # Poller, one task per bridged device, add/remove
│ └── time.rs # ISO-8601 timestamp parsing
├── ws/
│ ├── mod.rs
│ ├── connection.rs # WebSocket connection lifecycle
│ ├── dispatch.rs # Message dispatching and error sending
│ ├── constants.rs # Protocol constants and limits
│ ├── validation.rs # Message validation
│ ├── pending_play.rs # Pending play logic
│ └── handlers/
│ ├── mod.rs
│ ├── auth.rs # Authentication handler
│ ├── chat.rs # Chat message handler
│ ├── create.rs # Room creation
│ ├── join.rs # Room joining
│ ├── media.rs # set_media (host changes the room's item)
│ ├── misc.rs # Ping, ready, leave_room, client_log, unknown
│ └── playback.rs # player_event, state_update
└── room/
├── mod.rs
├── leave.rs # Client leave / disconnect
├── close.rs # Room closure
├── ops.rs # Shared room operations (add/move, kick, set host, close, groups)
└── reconnect.rs # Grace-period disconnect + reattachment
Module: main.rs
Description
Application entry point. Builds the Axum router and serves it with graceful shutdown.
Main Function
#[tokio::main]
async fn main() {
// Thread-safe shared state
let clients: Clients = Arc::new(RwLock::new(HashMap::new()));
let rooms: Rooms = Arc::new(RwLock::new(HashMap::new()));
// GET /ws (origin-guarded upgrade) + GET /health
let app = routes::build_router(clients, rooms, jwt_config, allowed_origins);
// Listen on 0.0.0.0:3000 (HOST/PORT override the defaults)
let listener = tokio::net::TcpListener::bind(addr).await.unwrap();
axum::serve(listener, app)
.with_graceful_shutdown(async { shutdown_rx.await.ok(); })
.await
.unwrap();
}
The server speaks HTTP/1.1 only. That is deliberate: every client performs an HTTP/1.1
Upgradefor WebSockets, so theaxum/hyperhttp2feature is left off andh2never enters the dependency tree.
Global State
| Variable | Type | Description |
|---|---|---|
clients |
Clients |
HashMap of connected clients |
rooms |
Rooms |
HashMap of active rooms |
Module: types.rs
Description
Defines the data structures used by the server.
Type Aliases
pub type Clients = Arc<RwLock<HashMap<String, Client>>>;
pub type Rooms = Arc<RwLock<HashMap<String, Room>>>;
Struct Client
Represents a connected WebSocket client.
| Field | Type | Description |
|---|---|---|
sender |
mpsc::Sender<...> |
Bounded channel for sending messages to client |
room_id |
Option<String> |
Current room ID (if in a room) |
user_id |
Option<String> |
Jellyfin user ID (set after auth) |
user_name |
Option<String> |
Display name (set after auth) |
authenticated |
bool |
Whether the client has authenticated |
message_count |
u32 |
Messages sent in current rate-limit window |
last_reset |
u64 |
Timestamp of last rate-limit window reset |
last_seen |
u64 |
Timestamp of last activity (for zombie detection) |
Struct Room
Represents a watch party room.
| Field | Type | Serialized | Description |
|---|---|---|---|
room_id |
String |
Yes | Unique identifier (UUID) |
name |
String |
Yes | Display name |
host_id |
String |
Yes | Host client ID |
media_id |
Option<String> |
Yes | Jellyfin media ID |
clients |
Vec<String> |
Yes | Participant client IDs |
ready_clients |
HashSet<String> |
Yes | Clients ready to receive play |
pending_play |
Option<PendingPlay> |
Yes | Pending play action |
state |
PlaybackState |
Yes | Current playback state |
last_state_ts |
u64 |
No | Last accepted state_update timestamp |
last_command_ts |
u64 |
No | Last player_event timestamp (cooldown) |
Struct PlaybackState
Room playback state.
| Field | Type | Description |
|---|---|---|
position |
f64 |
Position in seconds |
play_state |
String |
"playing" or "paused" |
Struct PendingPlay
Pending play action waiting for all clients to be ready.
| Field | Type | Description |
|---|---|---|
position |
f64 |
Position to start at |
created_at |
u64 |
Creation timestamp |
Struct WsMessage
WebSocket message format.
| Field | Type | JSON Key | Description |
|---|---|---|---|
msg_type |
String |
"type" |
Message type |
room |
Option<String> |
"room" |
Room ID |
client |
Option<String> |
"client" |
Sender client ID |
payload |
Option<Value> |
"payload" |
Message data |
ts |
u64 |
"ts" |
Client timestamp |
server_ts |
Option<u64> |
"server_ts" |
Server timestamp |
Module: ws/
Description
Handles WebSocket connections and main business logic. Split into sub-modules: connection.rs (lifecycle), dispatch.rs (message routing), constants.rs (protocol constants), validation.rs (message validation), pending_play.rs (play scheduling), and handlers/ (per-message-type logic).
Constants (ws/constants.rs)
| Constant | Value | Description |
|---|---|---|
PLAY_SCHEDULE_MS |
1000 | Delay before play execution (ms) |
CONTROL_SCHEDULE_MS |
300 | Delay before pause/seek execution (ms) |
MAX_READY_WAIT_MS |
2000 | Max wait time for ready clients (ms) |
MIN_STATE_UPDATE_INTERVAL_MS |
500 | Min interval between state updates (ms) |
POSITION_JITTER_THRESHOLD |
0.5 | Position noise threshold (seconds) |
COMMAND_COOLDOWN_MS |
2000 | Cooldown after player_event (ms) |
MAX_MESSAGE_SIZE |
65536 | Maximum message size (64 KB) |
Function client_connection
Manages a client connection lifecycle.
pub async fn client_connection(ws: WebSocket, clients: Clients, rooms: Rooms) {
// 1. Split WebSocket into sender/receiver
let (client_ws_sender, mut client_ws_rcv) = ws.split();
// 2. Create mpsc channel for async sending
let (client_sender, client_rcv) = mpsc::channel(CHANNEL_BUFFER_SIZE);
// 3. Task to forward messages to WebSocket
tokio::spawn(async move {
client_rcv.forward(client_ws_sender).await;
});
// 4. Generate UUID for client
let client_id = Uuid::new_v4().to_string();
// 5. Register client
clients.lock().unwrap().insert(client_id, Client { sender, room_id: None });
// 6. Send client_hello with ID
send_to_client(&client_id, &WsMessage {
msg_type: "client_hello",
payload: { "client_id": client_id }
});
// 7. Send room list
send_room_list(&client_id, &clients, &rooms);
// 8. Message receive loop
while let Some(msg) = client_ws_rcv.next().await {
client_msg(&client_id, msg, &clients, &rooms).await;
}
// 9. Cleanup on disconnect
handle_disconnect(&client_id, &clients, &rooms);
}
Function all_ready
Checks if all clients in a room are ready.
fn all_ready(room: &Room) -> bool {
room.ready_clients.len() >= room.clients.len()
}
Function broadcast_scheduled_play
Broadcasts a scheduled play event to all participants.
fn broadcast_scheduled_play(room: &mut Room, clients: &Clients, position: f64, target_server_ts: u64) {
// 1. Update room state
room.state.position = position;
room.state.play_state = "playing";
// 2. Create message with target_server_ts
let msg = WsMessage {
msg_type: "player_event",
payload: { "action": "play", "position": position, "target_server_ts": target_server_ts },
server_ts: target_server_ts
};
// 3. Broadcast to all room clients
broadcast_to_room(room, &clients, &msg, None);
}
Message Processing
The client_msg function handles incoming messages based on type:
list_rooms
Sends room list to the client.
create_room
Creates a new room with the sender as host.
join_room
Adds client to an existing room.
ready
Marks client as ready; triggers pending play if all ready.
leave_room
Removes client from room; closes room if host leaves.
player_event
Validates host permissions, applies action, broadcasts to room.
state_update
Applies filtering (cooldown, rate limit, jitter), broadcasts if accepted.
set_media (handlers/media.rs)
Host-only. Changes room.media_id, resets room.state/pending_play/
ready_clients (to just the host), and broadcasts media_changed to
the rest of the room plus a refreshed room_list. A no-op if the id
already matches or the sender isn’t the host. See protocol.md for the full
payload/effects and issue #71 for why this exists — before it,
media_id was write-once at create_room.
ping
Responds with pong for latency measurement.
Module: room/
Description
Manages room lifecycle and client disconnection. Split into leave.rs (client leave/disconnect), close.rs (room closure), reconnect.rs (grace-period disconnect + reattachment) and ops.rs.
ops.rs holds the operations that both the websocket handlers and the
admin API use, all on already-locked maps (rooms first, then clients):
| Function | Used by | What it does |
|---|---|---|
add_member |
join_room, admin “add” |
Moves a client into a room. If it is in another room it leaves that one first (normal leave notifications there). Makes it host if the room has none. Sends room_state (with admin_moved when an admin did it), participants_update and participants. No password check. |
kick_member |
admin “remove” | A normal leave for the room, plus room_closed with a reason to the removed client. |
set_host |
admin “make host” | Hands over the host role, drops a pending play, sends host_changed and participants. |
close_room |
host starting a new room, admin “close” | Removes the room and tells its members why. |
create_group / update_room |
admin | Creates an empty, hostless group (password optional); renames it or sets/clears its password. |
A room may be hostless (host_id empty) only while it is an empty
admin-created group. Host-only messages are ignored then (no client id
is empty), and the first member to arrive becomes host. Groups nobody
joins are removed after ADMIN_EMPTY_GROUP_TTL_SECS
(tasks::spawn_empty_group_reaper).
Function schedule_disconnect (room/reconnect.rs)
Called when a client’s WebSocket connection ends (close, error, or
zombie reap) — instead of tearing the client down immediately, this
captures the client’s current (about-to-be-dead) sender and spawns a
90-second timer (RECONNECT_GRACE_SECS):
const RECONNECT_GRACE_SECS: u64 = 90;
pub async fn schedule_disconnect(client_id: String, clients: Clients, rooms: Rooms) {
let stale_sender = /* clone the client's current sender */;
tokio::spawn(async move {
tokio::time::sleep(Duration::from_secs(RECONNECT_GRACE_SECS)).await;
// If the client reconnected, connection.rs already swapped in a
// new sender for this client_id — comparing channel identity
// tells us whether that happened.
let never_reconnected = /* clients[client_id].sender.same_channel(&stale_sender) */;
if never_reconnected {
handle_disconnect(&client_id, &clients, &rooms).await;
}
// else: reconnected within the grace period, keep room state as-is
});
}
If the client reconnects within the window using the same persistent
client_id (see Persistent Client ID below),
src/server/src/ws/connection.rs reattaches it to its existing entry
and calls resend_room_state (also in reconnect.rs) to bring it back
up to date — no room_closed is ever sent, and other participants never
see a disruption. Only if the grace period elapses with no reconnect
does handle_disconnect actually run.
Function handle_disconnect
Called once a client is confirmed gone (grace period expired with no reconnect).
pub fn handle_disconnect(client_id: &str, clients: &Clients, rooms: &Rooms) {
// 1. Remove client from their room
handle_leave(client_id, &mut clients, &mut rooms);
// 2. Remove client from the list
clients.remove(client_id);
// 3. Update room list for all
broadcast_room_list(clients, rooms);
}
Function handle_leave
Removes a client from a room.
pub fn handle_leave(client_id: &str, clients: &mut HashMap, rooms: &mut HashMap) {
if let Some(room_id) = client.room_id.take() {
if let Some(room) = rooms.get_mut(&room_id) {
// Remove client
room.clients.retain(|id| id != client_id);
room.ready_clients.remove(client_id);
// If host, cancel pending_play
if room.host_id == client_id {
room.pending_play = None;
}
// Close room if empty or host leaves
if room.clients.is_empty() || room.host_id == client_id {
for cid in &room.clients {
send_to_client(cid, { "type": "room_closed" });
}
rooms.remove(&room_id);
} else {
broadcast_to_room(room, { "type": "client_left", "client": client_id });
}
}
}
}
Module: jellyfin/
Description
Lets the admin panel put Jellyfin sessions that can’t run the web client
(TV apps, Fladder, …) into rooms. Only started with the admin panel and
JELLYFIN_URL + JELLYFIN_API_KEY.
- A bridged device is an ordinary client entry with
kind: ClientKind::Bridge, added throughops::add_member. Its outbound channel is read by a task inbridge.rsinstead of a socket. Bridge entries are skipped by the zombie reaper and can never be reattached to over/ws. - Host or receiver is not stored: on every tick the task checks
whether it is
room.host_id. As host it turns the device’s state intoset_media/player_event/state_update; as receiver it sends the devicePlayNow/Pause/Unpause/Seek. Both go throughws::dispatch_internal, i.e. the same handlers (host checks, ready gate, scheduling) as websocket traffic, includingreadyandclient_status. - One poller reads
GET /SessionseveryBRIDGE_POLL_INTERVAL_MS, only while a bridge exists or the panel looked within 30 s, and publishes a snapshot on awatchchannel that wakes every bridge task. - Positions: Jellyfin only updates
PlayState.PositionTickswhen the device reports progress, sologic::device_viewextrapolates fromLastPlaybackCheckIn(capped at three report intervals, 30-120 s). Jellyfin timestamps are shifted by an estimated clock offset (ClockEstimator: the smallest recentfetched_at - check-inover changed check-ins). After sending a seek or pause, the receiver assumes it took effect until a newer check-in arrives, for at most 10 s (FollowerMemory), so stale reports don’t cause repeat seeks. - Scheduled plays: the task notes each
player_eventplay’starget_server_tsand position. Until the room reports a newer state it measures the room position from that target, keeps receivers’ play/pause untouched before it (ahold), and wakes up 300 ms before the target to unpause them. - Remote-control commands run in their own task, so the bridge keeps reading room messages while Jellyfin answers.
- Receivers: a device that was on the room’s item and left it was
stopped on purpose and is left alone until the room changes item;
PlayNowis tried at most 3 times per item. - A bridge that becomes host marks the room started (no start countdown).
- Receiver thresholds: seek beyond 2 s drift (4 s cooldown, 1 s lead while
playing), pause/unpause 2.5 s cooldown,
PlayNow15 s cooldown. /Sessionsis fetched withoutactiveWithinSeconds(an idle TV makes no requests and would drop out); the admin list filters byLastActivityDate(16 min) instead, always keeping bridged devices. Each session is parsed on its own, so one odd entry can’t hide the rest.- Authenticates with
Authorization: MediaBrowser ... DeviceId= "jellywatchparty-session-server", Token="<api key>". Jellyfin treats API-key callers as privileged for remote control (SessionManager.AssertCanControl). Its own session is hidden. - A device missing from
/Sessionsfor 90 s leaves its room; a room closing stops its bridges.addreserves the device under one lock, so concurrent adds can’t bridge it twice;removedetaches in its own task, so a dropped HTTP request can’t leave an orphan.
Jellyfin API used (checked against Jellyfin 12, release-12.z)
| Call | Jellyfin side | Notes |
|---|---|---|
GET /Sessions |
SessionController.GetSessions |
With an API key, every session is returned (isApiKey). Fields read: Id, UserId (GUID, “N” format), UserName, Client, DeviceName, DeviceId, SupportsRemoteControl, NowPlayingItem.{Id,Name,SeriesName}, PlayState.{PositionTicks,IsPaused}, LastPlaybackCheckIn, LastActivityDate (UTC, ends in Z). |
POST /Sessions/{id}/Playing/{Pause,Unpause,Seek}?seekPositionTicks= |
SendPlaystateCommand |
204 on success. |
POST /Sessions/{id}/Playing?playCommand=PlayNow&itemIds=&startPositionTicks= |
Play |
Single episode + “next episode autoplay” makes Jellyfin queue the rest of the series. |
Auth header: Authorization: MediaBrowser Client="...", Device="...",
DeviceId="jellywatchparty-session-server", Version="...", Token="<API key>".
Jellyfin needs Client/Version/DeviceId to build the calling session; for an
API key it replaces Client with the key’s name and treats the caller as
privileged in SessionManager.AssertCanControl, so other users’ sessions
can be controlled. The server’s own session is recognised by its DeviceId
and hidden.
Module: messaging.rs
Description
Message sending utility functions.
Functions
send_room_list(client_id, clients, rooms)- Send room list to specific clientbroadcast_room_list(clients, rooms)- Send room list to all clientssend_to_client(client_id, clients, msg)- Send message to specific clientbroadcast_to_room(room, clients, msg, exclude)- Broadcast to room memberssend_error(client_id, clients, message)- Send error message (also available inws/dispatch.rs)
Module: auth.rs
Description
Optional JWT authentication.
Validation
pub fn validate_token(token: &str, secret: &str) -> Result<Claims, Error> {
let mut validation = Validation::new(Algorithm::HS256);
validation.validate_exp = true; // Enforce expiration
validation.leeway = 60; // 60 seconds tolerance
decode::<Claims>(token, &DecodingKey::from_secret(secret.as_ref()), &validation)
}
Concurrency Model
┌──────────────────────────────────────────────────────────────────┐
│ Tokio Runtime │
├──────────────────────────────────────────────────────────────────┤
│ ┌────────────────┐ ┌────────────────┐ ┌────────────────┐ │
│ │ Task: Client1 │ │ Task: Client2 │ │ Task: Client3 │ │
│ │ WebSocket │ │ WebSocket │ │ WebSocket │ │
│ └───────┬────────┘ └───────┬────────┘ └───────┬────────┘ │
│ │ │ │ │
│ └───────────────────┼───────────────────┘ │
│ │ │
│ ▼ │
│ Arc<RwLock<Clients>> │
│ Arc<RwLock<Rooms>> │
│ │
└──────────────────────────────────────────────────────────────────┘
Design Considerations
- RwLock: Read-heavy workload; multiple readers, exclusive writer
- No deadlock: Handlers that need both maps always take
roomsbeforeclients(including leave/disconnect and every admin action) - Message cloning: one
OutboundMessageis serialized once and cloned per recipient for efficient broadcasting - Bounded channels: Backpressure via bounded
mpsc::Senderper client
Reconnect and Room Lifecycle
Persistent Client ID
Reconnection is matched by identity, not by luck: the client generates a
UUID once per browser tab and stores it in sessionStorage
(getPersistentClientId()/withClientId() in
src/clients/jellyfin-web/ws/connection.js), then sends it as
?client_id=<uuid> on every WebSocket connection attempt (see
Protocol). If the server sees a connection with a client_id
that still has a live (or grace-period) entry, it swaps in the new
transport and keeps the existing room/host state (ws/connection.rs)
instead of registering a new client. A client-supplied ID is only trusted
if it looks like a real UUIDv4 — anything else falls back to a freshly
minted server-side ID.
Client ids are visible to everyone in a room, so reattaching also needs
the entry’s resume secret: a random 64-hex-char value created with the
entry and sent only to its owner in client_hello. The client stores it
next to its id and sends it as &resume=. Without the right secret
(decide_attach in ws/connection.rs, constant-time compare), a
connection asking for an id that is in use gets a fresh id instead of
taking over the entry’s room membership and host role. The secret is
replaced on every reattach.
Each websocket connection also gets a conn_id. When a connection ends,
it marks the entry disconnected and schedules the grace-period check only
if it is still the entry’s current connection; after the grace period the
entry is removed only if that same connection is still attached and
either ended or stayed silent for 60 s (reconnect::should_evict). This
keeps a half-dead old socket, noticed long after the client reconnected,
from evicting the live session.
Reconnection Behavior
| Scenario | Behavior |
|---|---|
| Any client reconnects within 90s | Reattaches to the same client entry; if they were in a room, room_state is resent and their host/guest role is restored |
| Host reconnects within 90s | Room stays open the whole time; participants see no disruption |
| Host does not reconnect within 90s, others remain | Earliest-joined remaining participant is promoted to host; host_changed broadcast; room stays open |
| Host does not reconnect within 90s, no one else remains | Room is torn down, room_closed broadcast (no other participants left to receive it) |
| Server restart | All rooms lost (in-memory only); clients reconnect to an empty server |
Auto-reconnect on the client retries with exponential backoff
(RECONNECT_BASE_MS up to RECONNECT_MAX_MS), not a fixed interval, and
shows “Reconnecting…” in the UI while autoReconnect=true.
Room Capacity and Scaling
Design limits: 20 clients per room, all state in-memory (rooms are
ephemeral by design), single server instance (sufficient for typical use —
see Configuration: Multi-Instance Setup
for why not to run more than one). At capacity, a join attempt gets a
"Room is full" error instead of being added.
| Rooms | Clients/Room | Total Clients | Expected Behavior |
|---|---|---|---|
| 10 | 5 | 50 | Excellent |
| 50 | 10 | 500 | Good |
| 100 | 15 | 1500 | Acceptable (monitor memory) |
| 200+ | 20 | 4000+ | May need resource limits |
Bottlenecks in rough order of likelihood: memory (~2KB/client, ~5KB/room), network (proportional to message rate × clients), CPU (minimal — message relay, no heavy computation).