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 Upgrade for WebSockets, so the axum/hyper http2 feature is left off and h2 never 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 through ops::add_member. Its outbound channel is read by a task in bridge.rs instead 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 into set_media / player_event / state_update; as receiver it sends the device PlayNow / Pause / Unpause / Seek. Both go through ws::dispatch_internal, i.e. the same handlers (host checks, ready gate, scheduling) as websocket traffic, including ready and client_status.
  • One poller reads GET /Sessions every BRIDGE_POLL_INTERVAL_MS, only while a bridge exists or the panel looked within 30 s, and publishes a snapshot on a watch channel that wakes every bridge task.
  • Positions: Jellyfin only updates PlayState.PositionTicks when the device reports progress, so logic::device_view extrapolates from LastPlaybackCheckIn (capped at three report intervals, 30-120 s). Jellyfin timestamps are shifted by an estimated clock offset (ClockEstimator: the smallest recent fetched_at - check-in over 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_event play’s target_server_ts and position. Until the room reports a newer state it measures the room position from that target, keeps receivers’ play/pause untouched before it (a hold), 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; PlayNow is 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, PlayNow 15 s cooldown.
  • /Sessions is fetched without activeWithinSeconds (an idle TV makes no requests and would drop out); the admin list filters by LastActivityDate (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 /Sessions for 90 s leaves its room; a room closing stops its bridges. add reserves the device under one lock, so concurrent adds can’t bridge it twice; remove detaches 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 client
  • broadcast_room_list(clients, rooms) - Send room list to all clients
  • send_to_client(client_id, clients, msg) - Send message to specific client
  • broadcast_to_room(room, clients, msg, exclude) - Broadcast to room members
  • send_error(client_id, clients, message) - Send error message (also available in ws/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

  1. RwLock: Read-heavy workload; multiple readers, exclusive writer
  2. No deadlock: Handlers that need both maps always take rooms before clients (including leave/disconnect and every admin action)
  3. Message cloning: one OutboundMessage is serialized once and cloned per recipient for efficient broadcasting
  4. Bounded channels: Backpressure via bounded mpsc::Sender per 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).


Back to top

JellyWatchParty - Synchronized watch parties for Jellyfin