5 Commits
Author SHA1 Message Date
Bendik Aagaard Lynghaug 0ad4b732e5 chore: Release cnats version 0.3.0
Release / build (aarch64, aarch64, /var/lib/gitea-runner/sccache) (push) Failing after 5s
Release / build (x86_64, ubuntu-latest, /sccache) (push) Failing after 2m45s
Release / update-aur (push) Skipped
Release / docker (push) Failing after 10m7s
2026-10-02 18:16:33 +02:00
Bendik Aagaard LynghaugandClaude Fable 5.1 78e9c9ff94 Calls: screen share, mic/cam toggles, status badges, pin + fullscreen, call timer
Screen sharing swaps the outgoing video track on every peer's sender with
replaceTrack (same m-line, no renegotiation); the camera returns when the
share stops, including via the browser's own "Stop sharing" bar. Mic and
camera toggles flip track.enabled. A new Status signal kind broadcasts
mic/cam/screen state so tiles can show badges; newcomers get it directed
when their Join arrives.

Tiles carry a name label, click to pin as a stage (a peer's screen share
pins itself), double-click for fullscreen. The call bar has an elapsed
timer and a participant count. Controls are icon-only inline SVGs on
currentColor: plain with a strike when off, a soft glow and ticking LED
when live. Red is reserved for leave.

Also fixed along the way: the peer map was a non-reactive StoredValue, so
the tile list never re-rendered when someone joined; ICE candidates that
arrive before the remote description are now buffered instead of
rejected; a peer whose connection fails or closes is dropped.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01L3sD29ozDvqA7jozeZXZoB
2026-10-02 17:26:00 +02:00
Bendik Aagaard LynghaugandClaude Opus 4.8 c1c337c2a2 ci: no toolchain installs on the bare aarch64 (klokka) runner; builtin:checkout
The aarch64 release leg runs on the klokka host. Gate the sccache /
cargo-binstall / cargo-leptos installs to the ephemeral x86_64 container; on
the bare runner check-and-fail instead (the sccache tarball was x86_64-only
anyway, so it was both a host mutation and a wrong-arch binary). Switch
checkout to v4 builtin:checkout (native Go, no Node/download) in all jobs.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-09-28 12:37:04 +02:00
Bendik Aagaard LynghaugandClaude Opus 4.8 9de5d7e2d9 CI: sccache-backed compiled-crate cache (persistent, replaces the dead cache service)
x86_64 uses /sccache (host dir bind-mounted into the job container by the
runner config); aarch64 on the klokka host uses a runner-owned dir. Both
persist across releases via RUSTC_WRAPPER=sccache.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01GLUwWE2KmFPzhKaf67tWbx
2026-09-13 12:37:53 +02:00
Bendik Aagaard LynghaugandClaude Opus 4.8 260f2791d0 CI: drop actions/cache — the Gitea cache backend times out (CreateCacheEntry), so it only added retry delay and never saved/restored
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01GLUwWE2KmFPzhKaf67tWbx
2026-09-13 12:21:14 +02:00
9 changed files with 888 additions and 143 deletions
+35 -12
View File
@@ -15,36 +15,59 @@ jobs:
include:
- arch: x86_64
runs-on: ubuntu-latest
sccache_dir: /sccache
- arch: aarch64
runs-on: aarch64
sccache_dir: /var/lib/gitea-runner/sccache
runs-on: ${{ matrix.runs-on }}
steps:
- uses: actions/checkout@v4
- uses: builtin:checkout
- name: Install Rust stable
uses: dtolnay/rust-toolchain@stable
with:
targets: wasm32-unknown-unknown
- name: Cache cargo registry and tools
uses: actions/cache@v3
with:
path: |
~/.cargo/registry
~/.cargo/bin
key: ${{ runner.os }}-${{ matrix.arch }}-cargo-leptos-${{ hashFiles('Cargo.lock') }}
# Compiled-crate cache via sccache. x86_64 uses /sccache (a persistent
# host dir bind-mounted into the job container by the runner config);
# aarch64 runs on the klokka host and uses a runner-owned dir. Both
# persist across releases, unlike the (unreachable) Gitea cache service.
- name: Set up sccache
run: |
echo "RUSTC_WRAPPER=sccache" >> "$GITHUB_ENV"
echo "SCCACHE_DIR=${{ matrix.sccache_dir }}" >> "$GITHUB_ENV"
mkdir -p "${{ matrix.sccache_dir }}"
if command -v sccache >/dev/null 2>&1; then sccache --version; exit 0; fi
# Bare runners (aarch64 = klokka host) are pre-provisioned; CI must not
# install onto them (and this tarball is x86_64-only). Only the
# ephemeral x86_64 container installs; a bare host missing it fails loud.
if [ "${{ matrix.arch }}" != x86_64 ]; then
echo "::error::sccache missing on the bare ${{ matrix.arch }} runner — provision klokka; CI must not install on bare hosts"; exit 1
fi
V=0.8.2
curl -sSL "https://github.com/mozilla/sccache/releases/download/v${V}/sccache-v${V}-x86_64-unknown-linux-musl.tar.gz" | tar -xz
sudo install -m0755 "sccache-v${V}-x86_64-unknown-linux-musl/sccache" /usr/local/bin/sccache
sccache --version
- name: Install cargo-binstall
run: |
command -v cargo-binstall || \
if command -v cargo-binstall >/dev/null 2>&1; then exit 0; fi
if [ "${{ matrix.arch }}" != x86_64 ]; then
echo "::error::cargo-binstall missing on the bare ${{ matrix.arch }} runner — provision klokka"; exit 1
fi
curl -L --proto '=https' --tlsv1.2 -sSf \
https://raw.githubusercontent.com/cargo-bins/cargo-binstall/main/install-from-binstall-release.sh \
| bash
- name: Install cargo-leptos
run: command -v cargo-leptos || cargo binstall cargo-leptos --locked --no-confirm
run: |
if command -v cargo-leptos >/dev/null 2>&1; then exit 0; fi
if [ "${{ matrix.arch }}" != x86_64 ]; then
echo "::error::cargo-leptos missing on the bare ${{ matrix.arch }} runner — provision klokka"; exit 1
fi
cargo binstall cargo-leptos --locked --no-confirm
- name: Build
run: cargo leptos build --release
@@ -101,7 +124,7 @@ jobs:
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
- uses: builtin:checkout
- name: Log in to Docker Hub
run: echo "${{ secrets.DOCKERHUB_TOKEN }}" | docker login -u bendik --password-stdin
@@ -118,7 +141,7 @@ jobs:
runs-on: aarch64
steps:
- uses: actions/checkout@v4
- uses: builtin:checkout
- name: Compute checksums and update PKGBUILD
run: |
Generated
+1 -1
View File
@@ -345,7 +345,7 @@ dependencies = [
[[package]]
name = "cnats"
version = "0.2.6"
version = "0.3.0"
dependencies = [
"anyhow",
"async-nats",
+3 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "cnats"
version = "0.2.6"
version = "0.3.0"
edition = "2021"
[lib]
@@ -63,6 +63,8 @@ web-sys = { version = "0.3", features = [
"RtcIceCandidateInit",
"RtcPeerConnectionIceEvent",
"RtcRtpSender",
"RtcPeerConnectionState",
"DisplayMediaStreamConstraints",
"RtcTrackEvent",
"RtcRtpTransceiver",
"RtcOfferOptions",
+26
View File
@@ -110,6 +110,8 @@ src/
app.rs Leptos UI (login gate + chat console)
auth.rs shared User type + current_user server fn
chat.rs shared ChatMessage/rooms + send_message server fn (publishes to NATS)
call.rs call signaling types + send_signal server fn (call.room.*, never archived)
webrtc.rs browser-only mesh WebRTC: peers, mic/cam toggles, screen share
server/
oidc.rs Kanidm OIDC login/callback/logout handlers
sse.rs NATS → browser SSE bridge (one subscription per client)
@@ -117,6 +119,30 @@ src/
style/main.css the console theme
```
## Calls
The `lobby` room has a mesh WebRTC call (every pair of browsers connects
directly; fine for a handful of people, not more). Signaling rides NATS on
`call.room.<room>` via `/call-sse/<room>` and the `send_signal` server fn, so
it is never archived. Only a public STUN server is configured: calls across
symmetric NATs will not connect without a TURN relay, which this app does not
run.
In a call:
- **MIC / CAM** toggle your mic and camera (the track stays attached and sends
silence/black, so toggling is instant and needs no renegotiation).
- **SCR ▶ share** shares a screen, window, or tab. It swaps the outgoing video
track on every peer with `replaceTrack`; your camera comes back when you
stop, or when the browser's own "Stop sharing" bar is used.
- Tiles show the participant's name and MIC ✕ / CAM ✕ / SCR badges, kept in
sync by a `Status` signal on the same subject.
- **Click** a tile to pin it as the stage (a shared screen pins itself);
**double-click** for fullscreen.
- A peer whose tab closes without leaving is dropped once their connection
fails, and ICE candidates that arrive before the offer are buffered rather
than lost.
## Notes & production hardening
- Sessions are in-memory (`tower-sessions` `MemoryStore`): restart logs
+1 -1
View File
@@ -1,6 +1,6 @@
# Maintainer: Bendik Aagaard Lynghaug <bendik.lynghaug@gmail.com>
pkgname=cnats
pkgver=0.2.6
pkgver=0.3.0
pkgrel=1
pkgdesc="Web chat over NATS subjects with Kanidm SSO (Leptos SSR)"
arch=('x86_64' 'aarch64')
+313 -63
View File
@@ -422,7 +422,6 @@ fn CallPanel(room: Memo<String>, me: String) -> impl IntoView {
let call_state = StoredValue::new_local(CallState::new(room.get_untracked(), me.clone()));
let in_call = call_state.get_value().in_call;
let local_video_ref = NodeRef::<leptos::html::Video>::new();
let es_handle = StoredValue::new_local(None::<web_sys::EventSource>);
Effect::new(move |_| {
@@ -448,54 +447,22 @@ fn CallPanel(room: Memo<String>, me: String) -> impl IntoView {
}
});
// Mirror the local MediaStream into the preview <video> element -
// `srcObject` has no HTML attribute form, has to be set via JS.
//
// Also force `.muted` via the JS property here, not just the
// `muted` attribute on the element below: every video tile in this
// app is built client-side via `document.createElement` (never
// parsed from HTML), and browsers only seed the live `.muted`
// property from the `muted` *attribute* for parser-inserted
// elements. Without this, the local preview plays back the user's
// own mic through their speakers - audible as feedback in a
// solo call.
Effect::new(move |_| {
let stream = call_state.get_value().local_stream().get();
if let Some(el) = local_video_ref.get() {
el.set_muted(true);
el.set_src_object(stream.as_ref());
}
});
let on_join = move |_| {
let cs = call_state.get_value();
leptos::task::spawn_local(async move {
cs.join().await;
});
};
let on_leave = move |_| {
let cs = call_state.get_value();
leptos::task::spawn_local(async move {
cs.leave().await;
});
};
view! {
<div class="call-panel">
{move || {
if in_call.get() {
view! {
<CallActive
local_video_ref=local_video_ref
call_state=call_state.get_value()
on_leave=on_leave
/>
}
.into_any()
view! { <CallActive call_state=call_state.get_value()/> }.into_any()
} else {
view! {
<button class="call-join" on:click=on_join>
"☎ join call"
<button class="call-join" title="join call" aria-label="join call" on:click=on_join>
"☎"
</button>
}
.into_any()
@@ -515,8 +482,8 @@ fn CallPanel(room: Memo<String>, me: String) -> impl IntoView {
let _ = (room, me);
view! {
<div class="call-panel">
<button class="call-join" disabled=true>
"☎ join call"
<button class="call-join" title="join call" aria-label="join call" disabled=true>
"☎"
</button>
</div>
}
@@ -524,57 +491,340 @@ fn CallPanel(room: Memo<String>, me: String) -> impl IntoView {
}
}
/// The in-call subtree (video grid + leave button), split out of
/// The in-call subtree (video grid + control bar), split out of
/// `CallPanel` so its `<For>`-over-peers view doesn't get inlined as a type
/// parameter of `CallPanel`'s own `if`/`else` branch - that inlining is what
/// was blowing the compiler's query recursion limit once mesh calling's
/// nested `Show`/`For` landed inside `ChatShell`'s own `Show`.
///
/// Layout: a flat grid of equal tiles, or - once a tile is pinned (click)
/// - that tile as a big "stage" with the rest in a strip below. A peer
/// starting a screen share gets pinned automatically; nobody shares a
/// screen to have it shown at thumbnail size.
#[cfg(feature = "hydrate")]
#[component]
fn CallActive(
local_video_ref: NodeRef<leptos::html::Video>,
call_state: crate::webrtc::CallState,
on_leave: impl Fn(leptos::ev::MouseEvent) + 'static,
) -> impl IntoView {
fn CallActive(call_state: crate::webrtc::CallState) -> impl IntoView {
let me = call_state.me().to_string();
let status = call_state.status;
let pinned = RwSignal::new(None::<String>);
let cs = StoredValue::new_local(call_state.clone());
// Call timer. The interval is cleared when this subtree unmounts (leave).
let elapsed = RwSignal::new(0u64);
{
let joined_at = call_state.joined_at;
let handle = set_interval_with_handle(
move || {
let secs = ((js_sys::Date::now() - joined_at.get_untracked()) / 1000.0).max(0.0);
elapsed.set(secs as u64);
},
std::time::Duration::from_secs(1),
);
on_cleanup(move || {
if let Ok(h) = handle {
h.clear();
}
});
}
// Auto-pin whoever starts sharing, and release that pin (only that
// one - a manual pin is left alone) when they stop.
let auto_pinned = RwSignal::new(false);
{
let cs2 = call_state.clone();
Effect::new(move |_| {
let sharing = cs2
.peer_views()
.into_iter()
.find(|p| p.status.get().screen)
.map(|p| p.id);
match sharing {
Some(id) => {
if pinned.get_untracked().as_deref() != Some(id.as_str()) {
pinned.set(Some(id));
auto_pinned.set(true);
}
}
None => {
if auto_pinned.get_untracked() {
pinned.set(None);
auto_pinned.set(false);
}
}
}
});
}
// Drop a pin pointing at a peer who has left.
{
let cs2 = call_state.clone();
Effect::new(move |_| {
let ids: Vec<String> = cs2.peer_views().into_iter().map(|p| p.id).collect();
if let Some(p) = pinned.get_untracked() {
if p != me && !ids.contains(&p) {
pinned.set(None);
}
}
});
}
let participants = Memo::new(move |_| cs.get_value().peer_views().len() + 1);
let me_id = call_state.me().to_string();
let me_for_tile = me_id.clone();
view! {
<div class="call-active">
<div class="video-grid">
<video
class="video-tile video-tile-local"
node_ref=local_video_ref
autoplay=true
muted=true
playsinline=true
></video>
<div class="video-grid" class:staged=move || pinned.get().is_some()>
<VideoTile
id=me_for_tile
label="you".to_string()
local=true
stream=Signal::derive(move || cs.get_value().preview_stream())
status=status.into()
pinned=pinned
auto_pinned=auto_pinned
/>
<For
each=move || call_state.peer_streams()
key=|(id, _)| id.clone()
children=move |(_id, stream)| {
view! { <PeerVideoTile stream=stream/> }
each=move || cs.get_value().peer_views()
key=|p| p.id.clone()
children=move |p| {
view! {
<VideoTile
id=p.id.clone()
label=format!("@{}", p.id)
local=false
stream=p.stream.into()
status=p.status.into()
pinned=pinned
auto_pinned=auto_pinned
/>
}
}
/>
</div>
<button class="call-leave" on:click=on_leave>
"⏏ leave call"
<div class="call-bar">
<span class="call-meta">
<span class="led live"></span>
<span class="call-timer">{move || format_elapsed(elapsed.get())}</span>
<span class="call-count" title="participants">
<CallIcon kind=CallIconKind::People off=Signal::derive(|| false)/>
{move || participants.get().to_string()}
</span>
</span>
<span class="call-ctls">
<button
class="call-ctl"
class:live=move || status.get().mic
title="microphone"
aria-label="microphone"
aria-pressed=move || status.get().mic.to_string()
on:click=move |_| cs.get_value().toggle_mic()
>
<span class="led"></span>
<CallIcon kind=CallIconKind::Mic off=Signal::derive(move || !status.get().mic)/>
</button>
<button
class="call-ctl"
class:live=move || status.get().cam
title="camera"
aria-label="camera"
aria-pressed=move || status.get().cam.to_string()
on:click=move |_| cs.get_value().toggle_cam()
>
<span class="led"></span>
<CallIcon kind=CallIconKind::Cam off=Signal::derive(move || !status.get().cam)/>
</button>
<button
class="call-ctl"
class:on=move || status.get().screen
title="share screen"
aria-label="share screen"
aria-pressed=move || status.get().screen.to_string()
on:click=move |_| cs.get_value().toggle_screen_share()
>
<CallIcon kind=CallIconKind::Screen off=Signal::derive(|| false)/>
</button>
<button
class="call-leave"
title="leave call"
aria-label="leave call"
on:click=move |_| {
let cs = cs.get_value();
leptos::task::spawn_local(async move {
cs.leave().await;
});
}
>
"⏏"
</button>
</span>
</div>
</div>
}
}
/// One remote participant's video tile - a plain child component so each
/// tile gets its own `NodeRef`/effect pair instead of trying to juggle a
/// `Vec` of node refs by hand in the parent.
#[cfg(feature = "hydrate")]
fn format_elapsed(secs: u64) -> String {
let (h, m, s) = (secs / 3600, (secs % 3600) / 60, secs % 60);
if h > 0 {
format!("{h}:{m:02}:{s:02}")
} else {
format!("{m:02}:{s:02}")
}
}
#[cfg(feature = "hydrate")]
#[derive(Clone, Copy, PartialEq, Eq)]
enum CallIconKind {
Mic,
Cam,
Screen,
People,
}
/// Inline SVG glyphs for the call controls and tile badges - stroked in
/// `currentColor` like `PulseMark`, so they take whatever colour the
/// surrounding control has. `off` draws a diagonal strike through the
/// glyph (muted / camera off).
#[cfg(feature = "hydrate")]
#[component]
fn PeerVideoTile(stream: RwSignal<Option<web_sys::MediaStream>>) -> impl IntoView {
fn CallIcon(kind: CallIconKind, off: Signal<bool>) -> impl IntoView {
let body = match kind {
CallIconKind::Mic => view! {
<rect x="9" y="3" width="6" height="11" rx="3"></rect>
<path d="M5 11a7 7 0 0 0 14 0M12 18v3M9 21h6"></path>
}
.into_any(),
CallIconKind::Cam => view! {
<rect x="3" y="7" width="13" height="10" rx="2"></rect>
<path d="M16 11l5-3v8l-5-3z"></path>
}
.into_any(),
CallIconKind::Screen => view! {
<rect x="3" y="4" width="18" height="12" rx="2"></rect>
<path d="M8 20h8M12 16v4"></path>
}
.into_any(),
CallIconKind::People => view! {
<circle cx="9" cy="8" r="3.5"></circle>
<path d="M3 20v-1a6 6 0 0 1 12 0v1M16 5a3.5 3.5 0 0 1 0 7M21 20v-1a6 6 0 0 0-4-5.6"></path>
}
.into_any(),
};
view! {
<svg
class="call-icon"
class:off=move || off.get()
viewBox="0 0 24 24"
fill="none"
stroke="currentColor"
stroke-width="2"
stroke-linecap="round"
stroke-linejoin="round"
aria-hidden="true"
>
{body}
<Show when=move || off.get()>
<path class="call-icon-strike-bg" d="M3 3l18 18"></path>
<path class="call-icon-strike" d="M3 3l18 18"></path>
</Show>
</svg>
}
}
/// One participant's tile: the `<video>` plus a name label and mic/cam/
/// screen badges. Click pins it as the stage; double-click goes
/// fullscreen (handy for reading a shared screen). A child component so
/// each tile owns its `NodeRef`/effect pair instead of the parent juggling
/// a `Vec` of node refs by hand.
#[cfg(feature = "hydrate")]
#[component]
fn VideoTile(
id: String,
label: String,
local: bool,
stream: Signal<Option<web_sys::MediaStream>>,
status: Signal<crate::call::MediaStatus>,
pinned: RwSignal<Option<String>>,
auto_pinned: RwSignal<bool>,
) -> impl IntoView {
let video_ref = NodeRef::<leptos::html::Video>::new();
// Mirror the MediaStream into the <video> - `srcObject` has no HTML
// attribute form, has to be set via JS.
//
// For the local tile also force `.muted` via the JS property, not just
// the `muted` attribute: every tile here is built client-side via
// `document.createElement` (never parsed from HTML), and browsers only
// seed the live `.muted` property from the attribute for
// parser-inserted elements. Without this the local preview plays the
// user's own mic back through their speakers.
Effect::new(move |_| {
let s = stream.get();
if let Some(el) = video_ref.get() {
if local {
el.set_muted(true);
}
el.set_src_object(s.as_ref());
}
});
view! { <video class="video-tile" node_ref=video_ref autoplay=true playsinline=true></video> }
let pin_id = id.clone();
let is_pinned = Memo::new(move |_| pinned.get().as_deref() == Some(pin_id.as_str()));
let click_id = id.clone();
let on_click = move |_| {
auto_pinned.set(false);
pinned.update(|p| {
if p.as_deref() == Some(click_id.as_str()) {
*p = None;
} else {
*p = Some(click_id.clone());
}
});
};
let on_dblclick = move |_| {
if let Some(el) = video_ref.get() {
let _ = el.request_fullscreen();
}
};
view! {
<div
class="tile"
class:local=local
class:pinned=move || is_pinned.get()
class:screen=move || status.get().screen
class:cam-off=move || !status.get().cam && !status.get().screen
on:click=on_click
on:dblclick=on_dblclick
title="click to pin · double-click for fullscreen"
>
<video
class="video-tile"
node_ref=video_ref
autoplay=true
muted=local
playsinline=true
></video>
<span class="tile-label">{label}</span>
<span class="tile-flags">
<Show when=move || status.get().screen>
<span class="flag on" title="sharing screen">
<CallIcon kind=CallIconKind::Screen off=Signal::derive(|| false)/>
</span>
</Show>
<Show when=move || !status.get().mic>
<span class="flag" title="muted">
<CallIcon kind=CallIconKind::Mic off=Signal::derive(|| true)/>
</span>
</Show>
<Show when=move || !status.get().cam && !status.get().screen>
<span class="flag" title="camera off">
<CallIcon kind=CallIconKind::Cam off=Signal::derive(|| true)/>
</span>
</Show>
</span>
</div>
}
}
/// Bridges `/call-sse/{room}` into `CallState::handle_signal`. A separate
+27
View File
@@ -36,6 +36,33 @@ pub enum CallSignalKind {
/// A single trickled ICE candidate, JSON-encoded
/// (`RTCIceCandidateInit`, produced client-side).
IceCandidate(String),
/// Mic/camera/screen state, so other tiles can show "muted" or "sharing
/// screen" without sniffing tracks. Broadcast on every change, and sent
/// directly to each newcomer when their `Join` arrives (they don't
/// know anything about us yet).
Status(MediaStatus),
}
/// What a participant is currently sending. `mic`/`cam` are the
/// `MediaStreamTrack.enabled` flags (muted tracks keep flowing as silence /
/// black, which is what lets un-muting be instant and renegotiation-free);
/// `screen` means their outgoing video track is a display capture rather
/// than the camera.
#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct MediaStatus {
pub mic: bool,
pub cam: bool,
pub screen: bool,
}
impl Default for MediaStatus {
fn default() -> Self {
Self {
mic: true,
cam: true,
screen: false,
}
}
}
/// Publishes a call-signaling message to the room's signaling subject.
+335 -50
View File
@@ -5,6 +5,11 @@
//! infrastructure this pass deliberately isn't standing up. Mesh topology
//! (every pair of peers connects directly) is fine at the ~4-person scale
//! this is scoped for; it does not scale further than that.
//!
//! Screen sharing swaps the outgoing video track on every peer's
//! `RTCRtpSender` via `replaceTrack` - same m-line, no renegotiation, so
//! the signaling path never sees it. Peers only learn about it through the
//! `Status` broadcast (for the "sharing" badge); the video just changes.
#![cfg(feature = "hydrate")]
use std::collections::HashMap;
@@ -14,11 +19,12 @@ use leptos::prelude::*;
use wasm_bindgen::{prelude::*, JsCast};
use wasm_bindgen_futures::JsFuture;
use web_sys::{
MediaStream, MediaStreamConstraints, RtcConfiguration, RtcIceCandidateInit,
RtcIceServer, RtcPeerConnection, RtcSdpType, RtcSessionDescriptionInit,
DisplayMediaStreamConstraints, MediaStream, MediaStreamConstraints, MediaStreamTrack,
RtcConfiguration, RtcIceCandidateInit, RtcIceServer, RtcPeerConnection,
RtcPeerConnectionState, RtcRtpSender, RtcSdpType, RtcSessionDescriptionInit,
};
use crate::call::{send_signal, CallSignalKind};
use crate::call::{send_signal, CallSignalKind, MediaStatus};
const STUN_URL: &str = "stun:stun.l.google.com:19302";
@@ -27,6 +33,23 @@ const STUN_URL: &str = "stun:stun.l.google.com:19302";
struct Peer {
conn: RtcPeerConnection,
stream: RwSignal<Option<MediaStream>>,
status: RwSignal<MediaStatus>,
/// ICE candidates that arrived before `setRemoteDescription` resolved.
/// Offer and candidates travel as separate server-fn POSTs, so nothing
/// guarantees their order on arrival; `addIceCandidate` before a remote
/// description throws, and the candidate would be lost - a call that
/// then "just doesn't connect". Buffered here, flushed once the remote
/// description is in.
pending_ice: Vec<RtcIceCandidateInit>,
remote_description_set: bool,
}
/// A remote peer as the UI sees it: id (username), stream, status.
#[derive(Clone)]
pub struct PeerView {
pub id: String,
pub stream: RwSignal<Option<MediaStream>>,
pub status: RwSignal<MediaStatus>,
}
/// Call state for one room. Lives for as long as the user is in the call;
@@ -35,9 +58,20 @@ struct Peer {
pub struct CallState {
room: String,
me: String,
/// Camera + mic from `getUserMedia`. Also the stream every outgoing
/// track is grouped under (`addTrack(track, stream)`), including the
/// screen track, so the remote's `ontrack` always sees one stream.
local_stream: RwSignal<Option<MediaStream>>,
peers: StoredValue<HashMap<String, Peer>, LocalStorage>,
/// Display capture while sharing; `None` otherwise.
screen_stream: RwSignal<Option<MediaStream>>,
/// `RwSignal` (not `StoredValue`) so the tile list re-renders when
/// peers come and go. `LocalStorage` because `RtcPeerConnection` is
/// `!Send`.
peers: RwSignal<HashMap<String, Peer>, LocalStorage>,
pub in_call: RwSignal<bool>,
pub status: RwSignal<MediaStatus>,
/// `Date.now()` at join, for the call timer.
pub joined_at: RwSignal<f64>,
}
fn new_peer_connection() -> Result<RtcPeerConnection, JsValue> {
@@ -61,12 +95,40 @@ async fn get_local_stream() -> Result<MediaStream, JsValue> {
stream.dyn_into::<MediaStream>()
}
fn attach_local_tracks(pc: &RtcPeerConnection, stream: &MediaStream) {
for track in stream.get_tracks().iter() {
if let Ok(track) = track.dyn_into::<web_sys::MediaStreamTrack>() {
pc.add_track_0(&track, stream);
async fn get_display_stream() -> Result<MediaStream, JsValue> {
let window = web_sys::window().ok_or("no window")?;
let media_devices = window.navigator().media_devices()?;
let constraints = DisplayMediaStreamConstraints::new();
constraints.set_video(&JsValue::TRUE);
let promise = media_devices.get_display_media_with_constraints(&constraints)?;
let stream = JsFuture::from(promise).await?;
stream.dyn_into::<MediaStream>()
}
fn tracks(stream: &MediaStream) -> Vec<MediaStreamTrack> {
stream
.get_tracks()
.iter()
.filter_map(|t| t.dyn_into::<MediaStreamTrack>().ok())
.collect()
}
fn first_video_track(stream: &MediaStream) -> Option<MediaStreamTrack> {
tracks(stream).into_iter().find(|t| t.kind() == "video")
}
fn stop_all(stream: &MediaStream) {
for track in tracks(stream) {
track.stop();
}
}
/// The sender carrying our outgoing video on this connection, if any.
fn video_sender(pc: &RtcPeerConnection) -> Option<RtcRtpSender> {
pc.get_senders()
.iter()
.filter_map(|s| s.dyn_into::<RtcRtpSender>().ok())
.find(|s| s.track().map(|t| t.kind() == "video").unwrap_or(false))
}
/// Reads the `sdp` field off whatever `create_offer`/`create_answer`
@@ -94,20 +156,39 @@ impl CallState {
room,
me,
local_stream: RwSignal::new(None),
peers: StoredValue::new_local(HashMap::new()),
screen_stream: RwSignal::new(None),
peers: RwSignal::new_local(HashMap::new()),
in_call: RwSignal::new(false),
status: RwSignal::new(MediaStatus::default()),
joined_at: RwSignal::new(0.0),
}
}
pub fn local_stream(&self) -> ReadSignal<Option<MediaStream>> {
self.local_stream.read_only()
pub fn me(&self) -> &str {
&self.me
}
/// Streams for currently-known peers, keyed by their peer id (username).
/// Recomputed each call - fine at mesh scale (~4 peers).
pub fn peer_streams(&self) -> Vec<(String, RwSignal<Option<MediaStream>>)> {
self.peers
.with_value(|p| p.iter().map(|(id, peer)| (id.clone(), peer.stream)).collect())
/// What the local preview tile should show: the screen while sharing,
/// the camera otherwise.
pub fn preview_stream(&self) -> Option<MediaStream> {
self.screen_stream.get().or_else(|| self.local_stream.get())
}
/// Currently-known peers, keyed by their peer id (username). Tracked:
/// re-runs whoever reads it when a peer joins or leaves.
pub fn peer_views(&self) -> Vec<PeerView> {
let mut views: Vec<PeerView> = self.peers.with(|p| {
p.iter()
.map(|(id, peer)| PeerView {
id: id.clone(),
stream: peer.stream,
status: peer.status,
})
.collect()
});
// HashMap order is arbitrary; keep tiles from shuffling on re-render.
views.sort_by(|a, b| a.id.cmp(&b.id));
views
}
/// getUserMedia, then broadcast Join so existing participants know to
@@ -120,30 +201,149 @@ impl CallState {
return;
}
}
self.status.set(MediaStatus::default());
self.joined_at.set(js_sys::Date::now());
self.in_call.set(true);
let _ = send_signal(self.room.clone(), None, CallSignalKind::Join).await;
}
/// Tears down every peer connection, stops all local tracks (releases
/// the camera/mic), and tells the room we're gone.
/// the camera/mic/screen), and tells the room we're gone.
pub async fn leave(&self) {
self.peers.update_value(|peers| {
self.peers.update(|peers| {
for (_, peer) in peers.drain() {
peer.conn.close();
}
});
if let Some(stream) = self.screen_stream.get_untracked() {
stop_all(&stream);
}
if let Some(stream) = self.local_stream.get_untracked() {
for track in stream.get_tracks().iter() {
if let Ok(track) = track.dyn_into::<web_sys::MediaStreamTrack>() {
track.stop();
}
}
stop_all(&stream);
}
self.screen_stream.set(None);
self.local_stream.set(None);
self.in_call.set(false);
let _ = send_signal(self.room.clone(), None, CallSignalKind::Leave).await;
}
// -- local media controls ------------------------------------------------
fn broadcast_status(&self, to: Option<String>) {
let room = self.room.clone();
let status = self.status.get_untracked();
leptos::task::spawn_local(async move {
let _ = send_signal(room, to, CallSignalKind::Status(status)).await;
});
}
fn set_local_tracks_enabled(&self, kind: &str, enabled: bool) {
if let Some(stream) = self.local_stream.get_untracked() {
for track in tracks(&stream).into_iter().filter(|t| t.kind() == kind) {
track.set_enabled(enabled);
}
}
}
/// Mute/unmute the mic. Flips `enabled` on the audio track (sends
/// silence) rather than removing it - instant, and no renegotiation.
pub fn toggle_mic(&self) {
let mic = !self.status.get_untracked().mic;
self.set_local_tracks_enabled("audio", mic);
self.status.update(|s| s.mic = mic);
self.broadcast_status(None);
}
/// Camera on/off. Only affects the camera track - while screen sharing
/// it's a no-op for peers until sharing stops, which is what you'd
/// expect ("my camera is off" shouldn't black out the slides).
pub fn toggle_cam(&self) {
let cam = !self.status.get_untracked().cam;
self.set_local_tracks_enabled("video", cam);
self.status.update(|s| s.cam = cam);
self.broadcast_status(None);
}
pub fn toggle_screen_share(&self) {
if self.status.get_untracked().screen {
self.stop_screen_share();
} else {
let this = self.clone();
leptos::task::spawn_local(async move {
this.start_screen_share().await;
});
}
}
/// Swap every peer's outgoing video track for a track from `stream`
/// (or back to the camera when `stream` is the camera stream).
fn replace_outgoing_video(&self, track: Option<&MediaStreamTrack>) {
self.peers.with_untracked(|peers| {
for peer in peers.values() {
if let Some(sender) = video_sender(&peer.conn) {
let _ = sender.replace_track(track);
}
}
});
}
async fn start_screen_share(&self) {
let stream = match get_display_stream().await {
Ok(s) => s,
// Most often: the user dismissed the picker. Not an error worth
// surfacing.
Err(e) => {
leptos::logging::log!("getDisplayMedia declined: {e:?}");
return;
}
};
// User may have left while the picker was up.
if !self.in_call.get_untracked() {
stop_all(&stream);
return;
}
let Some(track) = first_video_track(&stream) else {
stop_all(&stream);
return;
};
// Browsers put their own "Stop sharing" affordance on screen
// capture; when the user hits it the track ends under us and we
// have to fall back to the camera ourselves.
{
let this = self.clone();
let onended = Closure::<dyn FnMut()>::new(move || {
if this.status.get_untracked().screen {
this.stop_screen_share();
}
});
track.set_onended(Some(onended.as_ref().unchecked_ref()));
onended.forget();
}
self.replace_outgoing_video(Some(&track));
self.screen_stream.set(Some(stream));
self.status.update(|s| s.screen = true);
self.broadcast_status(None);
}
fn stop_screen_share(&self) {
let camera_track = self
.local_stream
.get_untracked()
.as_ref()
.and_then(first_video_track);
self.replace_outgoing_video(camera_track.as_ref());
if let Some(stream) = self.screen_stream.get_untracked() {
stop_all(&stream);
}
self.screen_stream.set(None);
self.status.update(|s| s.screen = false);
self.broadcast_status(None);
}
// -- signaling -----------------------------------------------------------
/// One incoming signal from `/call-sse/{room}`. Ignores our own
/// broadcasts and anything not addressed to us (directed messages are
/// broadcast NATS-wide and filtered client-side - see call.rs).
@@ -159,43 +359,71 @@ impl CallState {
match kind {
CallSignalKind::Join => {
// A new peer announced themselves - we initiate the offer.
self.start_offer(from);
}
CallSignalKind::Leave => {
self.peers.update_value(|peers| {
if let Some(peer) = peers.remove(&from) {
peer.conn.close();
}
});
// A new peer announced themselves - we initiate the offer,
// and tell them what we're sending (they know nothing yet).
self.start_offer(from.clone());
self.broadcast_status(Some(from));
}
CallSignalKind::Leave => self.remove_peer(&from),
CallSignalKind::Offer(sdp) => self.handle_offer(from, sdp),
CallSignalKind::Answer(sdp) => self.handle_answer(from, sdp),
CallSignalKind::IceCandidate(candidate_json) => {
self.handle_ice_candidate(from, candidate_json)
}
CallSignalKind::Status(status) => {
// May arrive before their offer does - `ensure_peer` is
// idempotent, so just create the slot early.
if self.ensure_peer(&from).is_some() {
self.peers.with_untracked(|peers| {
if let Some(peer) = peers.get(&from) {
peer.status.set(status);
}
});
}
}
}
}
fn remove_peer(&self, peer_id: &str) {
self.peers.update(|peers| {
if let Some(peer) = peers.remove(peer_id) {
peer.conn.close();
}
});
}
fn ensure_peer(&self, peer_id: &str) -> Option<RtcPeerConnection> {
if let Some(pc) = self
.peers
.with_value(|peers| peers.get(peer_id).map(|p| p.conn.clone()))
.with_untracked(|peers| peers.get(peer_id).map(|p| p.conn.clone()))
{
return Some(pc);
}
let pc = new_peer_connection().ok()?;
let Some(local) = self.local_stream.get_untracked() else {
return None;
};
attach_local_tracks(&pc, &local);
let local = self.local_stream.get_untracked()?;
// Audio from the camera stream; video is whatever we're currently
// showing - the screen if a share is in progress, else the camera.
for track in tracks(&local).into_iter().filter(|t| t.kind() == "audio") {
pc.add_track_0(&track, &local);
}
let video = self
.screen_stream
.get_untracked()
.as_ref()
.and_then(first_video_track)
.or_else(|| first_video_track(&local));
if let Some(track) = video {
pc.add_track_0(&track, &local);
}
let remote_stream = RwSignal::new(None::<MediaStream>);
{
let remote_stream = remote_stream;
let ontrack = Closure::<dyn FnMut(web_sys::RtcTrackEvent)>::new(move |ev: web_sys::RtcTrackEvent| {
remote_stream.set(Some(ev.streams().get(0).dyn_into().unwrap()));
if let Ok(stream) = ev.streams().get(0).dyn_into::<MediaStream>() {
remote_stream.set(Some(stream));
}
});
pc.set_ontrack(Some(ontrack.as_ref().unchecked_ref()));
ontrack.forget();
@@ -229,12 +457,33 @@ impl CallState {
onicecandidate.forget();
}
self.peers.update_value(|peers| {
// A peer whose tab just closed never sends `Leave`; drop them when
// the transport gives up so their tile doesn't linger forever.
{
let this = self.clone();
let pc2 = pc.clone();
let peer_id = peer_id.to_string();
let onstate = Closure::<dyn FnMut()>::new(move || {
if matches!(
pc2.connection_state(),
RtcPeerConnectionState::Failed | RtcPeerConnectionState::Closed
) {
this.remove_peer(&peer_id);
}
});
pc.set_onconnectionstatechange(Some(onstate.as_ref().unchecked_ref()));
onstate.forget();
}
self.peers.update(|peers| {
peers.insert(
peer_id.to_string(),
Peer {
conn: pc.clone(),
stream: remote_stream,
status: RwSignal::new(MediaStatus::default()),
pending_ice: Vec::new(),
remote_description_set: false,
},
);
});
@@ -261,11 +510,31 @@ impl CallState {
});
}
/// Marks the peer's remote description as set and replays any ICE
/// candidates that arrived too early.
fn flush_pending_ice(&self, peer_id: &str, pc: &RtcPeerConnection) {
let pending = self.peers.try_update_untracked(|peers| {
peers.get_mut(peer_id).map(|peer| {
peer.remote_description_set = true;
std::mem::take(&mut peer.pending_ice)
})
});
for init in pending.flatten().unwrap_or_default() {
let pc = pc.clone();
leptos::task::spawn_local(async move {
let _ = JsFuture::from(
pc.add_ice_candidate_with_opt_rtc_ice_candidate_init(Some(&init)),
)
.await;
});
}
}
fn handle_offer(&self, from: String, sdp: String) {
let Some(pc) = self.ensure_peer(&from) else {
return;
};
let room = self.room.clone();
let this = self.clone();
leptos::task::spawn_local(async move {
let remote_desc = RtcSessionDescriptionInit::new(RtcSdpType::Offer);
remote_desc.set_sdp(&sdp);
@@ -275,6 +544,7 @@ impl CallState {
{
return;
}
this.flush_pending_ice(&from, &pc);
let Ok(resolved) = JsFuture::from(pc.create_answer()).await else {
return;
};
@@ -289,35 +559,50 @@ impl CallState {
{
return;
}
let _ = send_signal(room, Some(from), CallSignalKind::Answer(answer_sdp)).await;
let _ = send_signal(this.room.clone(), Some(from), CallSignalKind::Answer(answer_sdp)).await;
});
}
fn handle_answer(&self, from: String, sdp: String) {
let Some(pc) = self
.peers
.with_value(|peers| peers.get(&from).map(|p| p.conn.clone()))
.with_untracked(|peers| peers.get(&from).map(|p| p.conn.clone()))
else {
return;
};
let this = self.clone();
leptos::task::spawn_local(async move {
let remote_desc = RtcSessionDescriptionInit::new(RtcSdpType::Answer);
remote_desc.set_sdp(&sdp);
let _ = JsFuture::from(pc.set_remote_description(&remote_desc)).await;
if JsFuture::from(pc.set_remote_description(&remote_desc))
.await
.is_ok()
{
this.flush_pending_ice(&from, &pc);
}
});
}
fn handle_ice_candidate(&self, from: String, candidate_json: String) {
let Some(pc) = self
.peers
.with_value(|peers| peers.get(&from).map(|p| p.conn.clone()))
else {
return;
};
let Ok(parsed) = js_sys::JSON::parse(&candidate_json) else {
return;
};
let init: RtcIceCandidateInit = parsed.unchecked_into();
// Either queue it (remote description not in yet) or hand back the
// connection to add it to right now.
let ready = self.peers.try_update_untracked(|peers| {
let peer = peers.get_mut(&from)?;
if peer.remote_description_set {
Some(peer.conn.clone())
} else {
peer.pending_ice.push(init.clone());
None
}
});
let Some(pc) = ready.flatten() else {
return;
};
leptos::task::spawn_local(async move {
let _ = JsFuture::from(
pc.add_ice_candidate_with_opt_rtc_ice_candidate_init(Some(&init)),
+141 -9
View File
@@ -534,8 +534,8 @@ body::before {
.call-join {
font-family: var(--mono);
font-weight: 600;
font-size: 0.78rem;
letter-spacing: 0.08em;
font-size: 1rem;
line-height: 1;
color: var(--ink-0);
background: var(--signal);
border: none;
@@ -554,26 +554,156 @@ body::before {
gap: 0.6rem;
}
.video-tile {
.tile {
position: relative;
width: 200px;
height: 150px;
background: var(--ink-0);
border: 1px solid var(--line-hot);
object-fit: cover;
cursor: pointer;
user-select: none;
transition: border-color 120ms;
}
.video-tile-local { border-color: var(--signal-dim); }
.tile:hover { border-color: var(--signal-dim); }
.tile.local { border-color: var(--signal-dim); }
.tile.pinned { border-color: var(--signal); box-shadow: 0 0 0 1px var(--signal); }
.video-tile {
display: block;
width: 100%;
height: 100%;
object-fit: cover;
background: var(--ink-0);
}
/* a shared screen is letterboxed, never cropped */
.tile.screen .video-tile { object-fit: contain; }
.tile.cam-off .video-tile { opacity: 0.25; }
/* stage mode: the pinned tile fills the row, everyone else in a strip */
.video-grid.staged .tile.pinned {
order: -1;
flex: 1 0 100%;
width: 100%;
height: auto;
aspect-ratio: 16 / 9;
max-height: 60vh;
}
.video-grid.staged .tile.pinned .video-tile { object-fit: contain; }
.video-grid.staged .tile:not(.pinned) { width: 140px; height: 105px; }
.tile-label,
.tile-flags {
position: absolute;
font-family: var(--mono);
font-size: 0.66rem;
letter-spacing: 0.06em;
line-height: 1;
pointer-events: none;
}
.tile-label {
left: 0;
bottom: 0;
padding: 0.3rem 0.45rem;
color: var(--text);
background: color-mix(in srgb, var(--ink-0) 75%, transparent);
}
.tile-flags {
top: 0.3rem;
right: 0.3rem;
display: flex;
gap: 0.25rem;
}
.flag {
display: inline-flex;
padding: 0.22rem;
color: var(--text-dim);
background: color-mix(in srgb, var(--ink-0) 80%, transparent);
border: 1px solid var(--line-hot);
}
.flag .call-icon { width: 12px; height: 12px; }
.flag.on { color: var(--ink-0); background: var(--signal); border-color: var(--signal); }
/* control bar */
.call-bar {
display: flex;
flex-wrap: wrap;
align-items: center;
justify-content: space-between;
gap: 0.6rem 1rem;
font-family: var(--mono);
font-size: 0.78rem;
}
.call-meta {
display: inline-flex;
align-items: center;
gap: 0.6rem;
color: var(--text-dim);
letter-spacing: 0.06em;
}
.call-timer { color: var(--text); font-variant-numeric: tabular-nums; }
.call-count { display: inline-flex; align-items: center; gap: 0.3rem; font-variant-numeric: tabular-nums; }
.call-count .call-icon { width: 14px; height: 14px; }
.call-ctls { display: inline-flex; flex-wrap: wrap; gap: 0.4rem; }
.call-ctl {
display: inline-flex;
align-items: center;
gap: 0.45rem;
color: var(--text);
background: none;
border: 1px solid var(--line-hot);
padding: 0.5rem 0.7rem;
cursor: pointer;
transition: background 120ms, border-color 120ms, color 120ms;
}
.call-ctl:hover { border-color: var(--signal-dim); }
/* capture controls: plain when off, and when live a soft glow with the
same ticking LED as the bus status - "this is recording", not an alarm */
.call-ctl .led { width: 6px; height: 6px; }
.call-ctl.live {
border-color: var(--signal-dim);
box-shadow: 0 0 10px color-mix(in srgb, var(--signal) 30%, transparent);
}
.call-ctl.live .led {
background: var(--signal);
box-shadow: 0 0 8px var(--signal);
animation: blink 1.4s steps(1) infinite;
}
.call-ctl.on { color: var(--ink-0); background: var(--signal); border-color: var(--signal); }
.call-icon { display: block; width: 18px; height: 18px; }
/* strike: a background stroke in the control's own fill colour first, so
the slash reads as cutting the glyph rather than overlapping it */
.call-icon-strike-bg { stroke: var(--ink-1); stroke-width: 5; }
.call-ctl.on .call-icon-strike-bg { stroke: var(--signal); }
.flag .call-icon-strike-bg { stroke: var(--ink-0); }
.call-leave {
align-self: flex-start;
font-family: var(--mono);
font-weight: 600;
font-size: 0.78rem;
letter-spacing: 0.08em;
font-size: 0.95rem;
line-height: 18px;
color: var(--text);
background: none;
border: 1px solid var(--alarm);
padding: 0.5rem 0.9rem;
padding: 0.5rem 0.8rem;
cursor: pointer;
transition: background 120ms;
}
@@ -618,6 +748,8 @@ body::before {
.room-link { border-left: 0; border-bottom: 2px solid transparent; }
.room-link.active { border-bottom-color: var(--signal); }
.msg { grid-template-columns: auto 1fr; grid-template-rows: auto auto; }
.tile, .video-grid.staged .tile:not(.pinned) { width: calc(50% - 0.3rem); height: auto; aspect-ratio: 4 / 3; }
.video-grid.staged .tile.pinned { width: 100%; }
.msg-time { grid-row: 1; order: 2; justify-self: end; }
.msg-text { grid-column: 1 / -1; }
}