22 Commits
Author SHA1 Message Date
Bendik Aagaard LynghaugandClaude Opus 4.8 8ab1b50310 chore: Release cnats version 0.3.2
Release / build (x86_64, ubuntu-latest, /sccache) (push) Successful in 1m30s
Release / docker (push) Successful in 13s
Release / build (aarch64, aarch64, /var/lib/gitea-runner/sccache) (push) Successful in 12m52s
Release / update-aur (push) Successful in 28s
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-10-03 10:17:22 +02:00
Bendik Aagaard LynghaugandClaude Opus 4.8 ede1b8497f fix(call): mirror SSR node structure so the lobby call button hydrates
CallPanel's hydrate body wraps the join button in a dynamic child
(`{move || if in_call {..} else {button}}`), which serializes hydration
marker comments; the SSR body rendered a *static* button (plus a
`disabled=true` the hydrate button lacks). That structural mismatch meant
`on:click` never attached on the first lobby load — the call button was dead
until a client-side nav re-rendered the subtree. Mirror the dynamic-child
wrapper and the button's serialized attributes in the SSR branch.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-10-03 10:16:50 +02:00
Bendik Aagaard Lynghaug cc88f3c325 chore: Release cnats version 0.3.1
Release / build (aarch64, aarch64, /var/lib/gitea-runner/sccache) (push) Failing after 1s
Release / build (x86_64, ubuntu-latest, /sccache) (push) Successful in 1m32s
Release / update-aur (push) Skipped
Release / docker (push) Successful in 2m11s
2026-10-02 18:33:24 +02:00
Bendik Aagaard LynghaugandClaude Fable 5.1 e4c0094f34 Fix release build: box the call tiles, icons and in-call subtree
`cargo leptos build --release` died with "queries overflow the depth
limit": every tile's Show/icon views inlined into one enormous
hydrate_async type for the whole call subtree. VideoTile, CallIcon and
CallActive now return AnyView, the same cut CallPanel got in 0.2.1. The
dev profile never tripped it, which is why 0.3.0 shipped broken.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01L3sD29ozDvqA7jozeZXZoB
2026-10-02 18:33:23 +02:00
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
Bendik Aagaard Lynghaug ed0e0229d4 chore: Release cnats version 0.2.6
Release / build (x86_64, ubuntu-latest) (push) Successful in 6m31s
Release / build (aarch64, aarch64) (push) Successful in 15m12s
Release / update-aur (push) Successful in 49s
Release / docker (push) Failing after 20m5s
2026-09-13 12:16:43 +02:00
Bendik Aagaard LynghaugandClaude Opus 4.8 e998860e09 Archive consumer: capped exponential backoff (resets after a healthy run)
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01GLUwWE2KmFPzhKaf67tWbx
2026-09-13 12:16:29 +02:00
Bendik Aagaard LynghaugandClaude Opus 4.8 7361a39618 CI: Arch-registry publish uses scoped REGISTRY_TOKEN, gated + non-fatal
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01GLUwWE2KmFPzhKaf67tWbx
2026-09-13 11:38:09 +02:00
Bendik Aagaard LynghaugandClaude Opus 4.8 0754e9e3e3 CI: fix Arch-registry publish (use makepkg --packagelist; !strip for cross-arch repack)
The runner's PKGEXT is .pkg.tar.xz, so the .zst glob never matched;
--packagelist yields the exact filename, and options=('!strip') lets the
foreign-arch binary be packaged on the aarch64 host without stripping.

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01GLUwWE2KmFPzhKaf67tWbx
2026-09-13 11:05:38 +02:00
Bendik Aagaard LynghaugandClaude Fable 5 903b8ccbb9 Rehome to project.uhhm.no: PKGBUILD URLs, ephemeral CI token, Arch registry publishing
- PKGBUILD url/source now point at this instance's releases.
- Release uploads use the run's own ephemeral token instead of the
  GITEA_TOKEN secret.
- The publish job also builds both architectures' packages (repack
  PKGBUILD, CARCH override) and uploads them to the instance Arch
  package registry (repository name: uhhm).

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_01GLUwWE2KmFPzhKaf67tWbx
2026-09-12 17:30:46 +02:00
Bendik Aagaard Lynghaug 1f47337c40 chore: Release cnats version 0.2.5 2026-08-16 14:42:39 +02:00
Bendik Aagaard LynghaugandClaude Fable 5 e2caee536f Light theme: green-tinted paper print following prefers-color-scheme
Palette override in a light media query; scanlines, vignette, card
shadow and signal glows moved to tokens/color-mix so both prints share
one set of rules. color-scheme on :root brings UA scrollbars and form
controls along.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-08-16 14:40:49 +02:00
Bendik Aagaard Lynghaug 080b9c46d6 chore: Release cnats version 0.2.4 2026-07-30 15:06:52 +02:00
Bendik Aagaard Lynghaug 3807294333 Fix mic feedback in solo calls: force .muted on the local preview tile
Every video tile is built client-side via document.createElement, never
parsed from HTML - browsers only seed the live .muted property from the
muted attribute for parser-inserted elements, so the attribute alone
never actually silenced the local preview. Result: your own mic played
back through your own speakers even with no other peers in the call.
Now sets el.set_muted(true) directly whenever the preview's srcObject
is (re)assigned.
2026-07-30 15:06:14 +02:00
Bendik Aagaard Lynghaug 9c89cfbe9c chore: Release cnats version 0.2.3 2026-07-29 10:50:09 +02:00
Bendik Aagaard Lynghaug 9326f93576 Fix hydration mismatch in CallPanel that crashed room switching
The SSR branch rendered an empty call-panel div while the hydrate
branch expected a child button node. On any full load/refresh of the
lobby room, the mismatched hydration cursor hit tachys's
unreachable!() panic path, trapping the wasm instance and killing all
reactivity (routing, room switching, SSE) for the rest of the page
load, while the already-rendered SSR HTML stayed visually intact.

SSR now renders the same default join-button markup hydrate expects.
2026-07-29 10:50:02 +02:00
Bendik Aagaard Lynghaug 9858460696 chore: Release cnats version 0.2.2 2026-07-28 12:03:23 +02:00
Bendik Aagaard Lynghaug 9550fa1f72 CI: install cargo-leptos via cargo-binstall instead of building from source
Building cargo-leptos from source pulls in swc/lightningcss/rhai and was
OOM-killing the aarch64 (Raspberry Pi) runner. cargo-leptos publishes
prebuilt aarch64-unknown-linux-gnu binaries, so binstall sidesteps the
compile entirely.
2026-07-28 12:02:26 +02:00
10 changed files with 1008 additions and 147 deletions
+72 -16
View File
@@ -15,29 +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: |
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 install cargo-leptos --locked
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
@@ -58,7 +88,7 @@ jobs:
- name: Create release
run: |
curl -sX POST \
-H "Authorization: token ${{ secrets.GITEA_TOKEN }}" \
-H "Authorization: token ${{ secrets.GITHUB_TOKEN }}" \
-H "Content-Type: application/json" \
"${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases" \
-d "{\"tag_name\":\"${{ gitea.ref_name }}\",\"name\":\"${{ gitea.ref_name }}\"}" \
@@ -67,24 +97,24 @@ jobs:
- name: Upload assets
run: |
RELEASE_ID=$(curl -s \
-H "Authorization: token ${{ secrets.GITEA_TOKEN }}" \
-H "Authorization: token ${{ secrets.GITHUB_TOKEN }}" \
"${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases/tags/${{ gitea.ref_name }}" \
| jq -r '.id')
for FILE in "${{ env.TARBALL }}" "${{ env.TARBALL }}.sha256"; do
# Remove any existing asset with the same name so re-runs stay clean
EXISTING=$(curl -s \
-H "Authorization: token ${{ secrets.GITEA_TOKEN }}" \
-H "Authorization: token ${{ secrets.GITHUB_TOKEN }}" \
"${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases/${RELEASE_ID}/assets" \
| jq -r ".[] | select(.name == \"${FILE}\") | .id")
for AID in $EXISTING; do
curl -sX DELETE \
-H "Authorization: token ${{ secrets.GITEA_TOKEN }}" \
-H "Authorization: token ${{ secrets.GITHUB_TOKEN }}" \
"${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases/${RELEASE_ID}/assets/${AID}"
done
curl -sX POST \
-H "Authorization: token ${{ secrets.GITEA_TOKEN }}" \
-H "Authorization: token ${{ secrets.GITHUB_TOKEN }}" \
-H "Content-Type: application/octet-stream" \
"${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases/${RELEASE_ID}/assets?name=${FILE}" \
--data-binary "@${FILE}" --fail-with-body
@@ -94,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
@@ -111,7 +141,7 @@ jobs:
runs-on: aarch64
steps:
- uses: actions/checkout@v4
- uses: builtin:checkout
- name: Compute checksums and update PKGBUILD
run: |
@@ -125,6 +155,32 @@ jobs:
sed -i "s/sha256sums_x86_64=('.*')/sha256sums_x86_64=('${SUM_X86}')/" aur/PKGBUILD
sed -i "s/sha256sums_aarch64=('.*')/sha256sums_aarch64=('${SUM_AARCH}')/" aur/PKGBUILD
# Also publish the built packages to this instance's Arch registry
# (docs.gitea.com/usage/packages/arch). The PKGBUILD only repacks the
# release tarballs, so CARCH can produce both architectures from this
# one host. Consumers: see the infrastructure README.
# Best-effort mirror to the instance Arch registry. The ephemeral
# GITHUB_TOKEN is not accepted as a package-write credential, so this
# uses a dedicated REGISTRY_TOKEN secret (a write:package token for bl);
# if it is unset the step is skipped, and continue-on-error keeps a
# registry hiccup from failing the release or the AUR push.
- name: Publish to the Arch package registry
continue-on-error: true
run: |
set -euo pipefail
if [ -z "${{ secrets.REGISTRY_TOKEN }}" ]; then
echo "::warning::REGISTRY_TOKEN not set — skipping Arch registry publish"
exit 0
fi
cd aur
for carch in aarch64 x86_64; do
pkgfile=$(CARCH="$carch" makepkg --packagelist | tail -1)
CARCH="$carch" makepkg -f --nodeps --noconfirm --skipinteg
curl --fail-with-body --user "bl:${{ secrets.REGISTRY_TOKEN }}" \
--upload-file "$pkgfile" \
"${{ gitea.server_url }}/api/packages/${{ gitea.repository_owner }}/arch/uhhm"
done
- name: Push to AUR
env:
AUR_SSH_KEY: ${{ secrets.AUR_SSH_KEY }}
Generated
+1 -1
View File
@@ -345,7 +345,7 @@ dependencies = [
[[package]]
name = "cnats"
version = "0.2.1"
version = "0.3.2"
dependencies = [
"anyhow",
"async-nats",
+3 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "cnats"
version = "0.2.1"
version = "0.3.2"
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
+5 -4
View File
@@ -1,10 +1,11 @@
# Maintainer: Bendik Aagaard Lynghaug <bendik.lynghaug@gmail.com>
pkgname=cnats
pkgver=0.2.1
pkgver=0.3.1
pkgrel=1
pkgdesc="Web chat over NATS subjects with Kanidm SSO (Leptos SSR)"
arch=('x86_64' 'aarch64')
url="https://prosjekt.klingenbergbygg.no/bl/cnats"
options=('!strip')
url="https://project.uhhm.no/bl/cnats"
license=('MIT')
depends=('glibc' 'gcc-libs')
optdepends=(
@@ -14,8 +15,8 @@ optdepends=(
provides=('cnats')
conflicts=('cnats-git' 'cnats-bin')
backup=('etc/cnats/env')
source_x86_64=("cnats-v${pkgver}-x86_64.tar.gz::https://prosjekt.klingenbergbygg.no/bl/cnats/releases/download/v${pkgver}/cnats-v${pkgver}-x86_64.tar.gz")
source_aarch64=("cnats-v${pkgver}-aarch64.tar.gz::https://prosjekt.klingenbergbygg.no/bl/cnats/releases/download/v${pkgver}/cnats-v${pkgver}-aarch64.tar.gz")
source_x86_64=("cnats-v${pkgver}-x86_64.tar.gz::https://project.uhhm.no/bl/cnats/releases/download/v${pkgver}/cnats-v${pkgver}-x86_64.tar.gz")
source_aarch64=("cnats-v${pkgver}-aarch64.tar.gz::https://project.uhhm.no/bl/cnats/releases/download/v${pkgver}/cnats-v${pkgver}-aarch64.tar.gz")
sha256sums_x86_64=('SKIP')
sha256sums_aarch64=('SKIP')
+342 -52
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,44 +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.
Effect::new(move |_| {
let stream = call_state.get_value().local_stream().get();
if let Some(el) = local_video_ref.get() {
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()
@@ -495,64 +472,377 @@ fn CallPanel(room: Memo<String>, me: String) -> impl IntoView {
}
.into_any()
}
// Must mirror the hydrate branch's default (not-in-call) markup AND its
// node structure exactly. The hydrate branch puts the button inside a
// dynamic child (`{move || ...}`), which serializes hydration-marker
// comments around it; if SSR emits a plain static button instead, the
// hydrate walk can't find the dynamic node, so it never attaches
// `on:click` - the call button is dead on the first lobby load until a
// client-side nav re-renders the subtree. So: same `{move || ...}` wrapper,
// and only the attributes the hydrate button actually serializes to
// (class/title/aria-label; `on:click` is a JS listener, not HTML). No
// `disabled` here - the hydrate button has none and hydration won't clear
// a stray SSR-only attribute.
#[cfg(not(feature = "hydrate"))]
{
let _ = (room, me);
view! { <div class="call-panel"></div> }.into_any()
view! {
<div class="call-panel">
{move || {
view! {
<button class="call-join" title="join call" aria-label="join call">
"☎"
</button>
}
.into_any()
}}
</div>
}
.into_any()
}
}
/// 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>
}
.into_any()
}
#[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}")
}
}
/// 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")]
#[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>
}
.into_any()
}
/// One participant's tile: the `<video>` plus a name label and mic/cam/
/// screen badges. Returns `AnyView` (as do `CallIcon` and `CallActive`):
/// left generic, every tile's `Show`s and icons inline into one enormous
/// `hydrate_async` type for the whole call subtree, and the release build
/// dies with "queries overflow the depth limit" - the same wall
/// `CallPanel` hit in 0.2.1. 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>
}
.into_any()
}
/// 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.
+18 -3
View File
@@ -2,7 +2,7 @@
//! into Postgres, so history survives restarts and includes messages
//! published by any client on the bus (not just this app).
use std::time::Duration;
use std::time::{Duration, Instant};
use async_nats::jetstream;
use futures::StreamExt;
@@ -38,11 +38,26 @@ pub async fn init_schema(pool: &PgPool) -> anyhow::Result<()> {
/// Runs forever; (re)creates the stream/consumer and retries on any failure,
/// so a NATS or Postgres outage never takes the chat server down.
pub async fn run_consumer(nats: async_nats::Client, pool: PgPool) {
const MIN_BACKOFF: Duration = Duration::from_secs(5);
const MAX_BACKOFF: Duration = Duration::from_secs(60);
let mut backoff = MIN_BACKOFF;
loop {
let started = Instant::now();
if let Err(err) = consume(&nats, &pool).await {
tracing::error!("archive consumer failed: {err:#}; retrying in 5s");
// A failure after a long healthy run is a fresh incident, not an
// escalating one - reset the backoff so we retry promptly.
if started.elapsed() >= MAX_BACKOFF {
backoff = MIN_BACKOFF;
}
tracing::error!(
"archive consumer failed after {:?}: {err:#}; retrying in {}s",
started.elapsed(),
backoff.as_secs()
);
tokio::time::sleep(backoff).await;
// Cap the backoff so a persistent outage doesn't hammer NATS/PG.
backoff = (backoff * 2).min(MAX_BACKOFF);
}
tokio::time::sleep(Duration::from_secs(5)).await;
}
}
+336 -51
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,14 +95,42 @@ 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`
/// resolved to, and builds a fresh `RtcSessionDescriptionInit` from it -
/// simpler and more reliable than trying to cast the resolved JsValue
@@ -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)),
+176 -17
View File
@@ -1,8 +1,10 @@
/* ── cnats · message-bus console ─────────────────────────────────────────
dark phosphor terminal: deep green-black ground, mint signal, amber id.
phosphor terminal in two prints: dark (green-black ground, mint signal)
and light (green-tinted paper, forest signal), following the OS scheme.
type: Archivo (UI voice) + IBM Plex Mono (wire voice). */
:root {
color-scheme: light dark;
--ink-0: #060a09;
--ink-1: #0b1210;
--ink-2: #101a17;
@@ -16,10 +18,34 @@
--signal-dim: #2a8f6c;
--amber: #ffb454;
--alarm: #ff6b6b;
--scanline: rgba(255, 255, 255, 0.015);
--vignette: rgba(0, 0, 0, 0.45);
--card-shadow: rgba(0, 0, 0, 0.55);
--mono: "IBM Plex Mono", ui-monospace, monospace;
--sans: "Archivo", system-ui, sans-serif;
}
@media (prefers-color-scheme: light) {
:root {
--ink-0: #f3f6f4;
--ink-1: #eaf0ec;
--ink-2: #dfe8e2;
--ink-3: #d2ded6;
--line: #c3d2c9;
--line-hot: #a3bcae;
--text: #14211c;
--text-dim: #46584f;
--text-faint: #74887e;
--signal: #0b7a52;
--signal-dim: #2f8a66;
--amber: #a85f00;
--alarm: #c73f3f;
--scanline: rgba(6, 10, 9, 0.02);
--vignette: rgba(6, 10, 9, 0.06);
--card-shadow: rgba(20, 33, 28, 0.18);
}
}
* { margin: 0; padding: 0; box-sizing: border-box; }
html, body { height: 100%; }
@@ -39,8 +65,8 @@ body::before {
pointer-events: none;
z-index: 999;
background:
repeating-linear-gradient(0deg, rgba(255, 255, 255, 0.015) 0 1px, transparent 1px 3px),
radial-gradient(ellipse 120% 90% at 50% 40%, transparent 55%, rgba(0, 0, 0, 0.45));
repeating-linear-gradient(0deg, var(--scanline) 0 1px, transparent 1px 3px),
radial-gradient(ellipse 120% 90% at 50% 40%, transparent 55%, var(--vignette));
}
::selection { background: var(--signal); color: var(--ink-0); }
@@ -53,7 +79,7 @@ body::before {
place-items: center;
padding: 2rem;
background:
radial-gradient(ellipse 60% 45% at 50% 0%, rgba(78, 240, 177, 0.07), transparent 70%),
radial-gradient(ellipse 60% 45% at 50% 0%, color-mix(in srgb, var(--signal) 7%, transparent), transparent 70%),
linear-gradient(var(--ink-0), var(--ink-1));
}
@@ -69,7 +95,7 @@ body::before {
background: linear-gradient(160deg, var(--ink-2), var(--ink-1) 60%);
padding: 3rem 2.75rem 2.5rem;
position: relative;
box-shadow: 0 40px 80px rgba(0, 0, 0, 0.55);
box-shadow: 0 40px 80px var(--card-shadow);
animation: rise 0.5s cubic-bezier(0.2, 0.9, 0.3, 1) both;
}
@@ -129,7 +155,8 @@ body::before {
.gate-btn:hover {
background: transparent;
color: var(--signal);
box-shadow: 0 0 24px rgba(78, 240, 177, 0.25), inset 0 0 12px rgba(78, 240, 177, 0.08);
box-shadow: 0 0 24px color-mix(in srgb, var(--signal) 25%, transparent),
inset 0 0 12px color-mix(in srgb, var(--signal) 8%, transparent);
}
.gate-btn-glyph { font-size: 1.05rem; }
@@ -458,7 +485,7 @@ body::before {
.composer:focus-within {
border-color: var(--signal-dim);
box-shadow: 0 0 0 1px var(--signal-dim), 0 0 30px rgba(78, 240, 177, 0.08);
box-shadow: 0 0 0 1px var(--signal-dim), 0 0 30px color-mix(in srgb, var(--signal) 8%, transparent);
}
.composer-prompt {
@@ -507,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;
@@ -527,31 +554,161 @@ 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;
}
.call-leave:hover { background: rgba(255, 90, 90, 0.12); }
.call-leave:hover { background: color-mix(in srgb, var(--alarm) 12%, transparent); }
/* motion */
@@ -591,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; }
}