12 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
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
10 changed files with 976 additions and 162 deletions
+66 -17
View File
@@ -15,36 +15,59 @@ jobs:
include: include:
- arch: x86_64 - arch: x86_64
runs-on: ubuntu-latest runs-on: ubuntu-latest
sccache_dir: /sccache
- arch: aarch64 - arch: aarch64
runs-on: aarch64 runs-on: aarch64
sccache_dir: /var/lib/gitea-runner/sccache
runs-on: ${{ matrix.runs-on }} runs-on: ${{ matrix.runs-on }}
steps: steps:
- uses: actions/checkout@v4 - uses: builtin:checkout
- name: Install Rust stable - name: Install Rust stable
uses: dtolnay/rust-toolchain@stable uses: dtolnay/rust-toolchain@stable
with: with:
targets: wasm32-unknown-unknown targets: wasm32-unknown-unknown
- name: Cache cargo registry and tools # Compiled-crate cache via sccache. x86_64 uses /sccache (a persistent
uses: actions/cache@v3 # host dir bind-mounted into the job container by the runner config);
with: # aarch64 runs on the klokka host and uses a runner-owned dir. Both
path: | # persist across releases, unlike the (unreachable) Gitea cache service.
~/.cargo/registry - name: Set up sccache
~/.cargo/bin run: |
key: ${{ runner.os }}-${{ matrix.arch }}-cargo-leptos-${{ hashFiles('Cargo.lock') }} 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 - name: Install cargo-binstall
run: | 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 \ curl -L --proto '=https' --tlsv1.2 -sSf \
https://raw.githubusercontent.com/cargo-bins/cargo-binstall/main/install-from-binstall-release.sh \ https://raw.githubusercontent.com/cargo-bins/cargo-binstall/main/install-from-binstall-release.sh \
| bash | bash
- name: Install cargo-leptos - 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 - name: Build
run: cargo leptos build --release run: cargo leptos build --release
@@ -65,7 +88,7 @@ jobs:
- name: Create release - name: Create release
run: | run: |
curl -sX POST \ curl -sX POST \
-H "Authorization: token ${{ secrets.GITEA_TOKEN }}" \ -H "Authorization: token ${{ secrets.GITHUB_TOKEN }}" \
-H "Content-Type: application/json" \ -H "Content-Type: application/json" \
"${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases" \ "${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases" \
-d "{\"tag_name\":\"${{ gitea.ref_name }}\",\"name\":\"${{ gitea.ref_name }}\"}" \ -d "{\"tag_name\":\"${{ gitea.ref_name }}\",\"name\":\"${{ gitea.ref_name }}\"}" \
@@ -74,24 +97,24 @@ jobs:
- name: Upload assets - name: Upload assets
run: | run: |
RELEASE_ID=$(curl -s \ 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 }}" \ "${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases/tags/${{ gitea.ref_name }}" \
| jq -r '.id') | jq -r '.id')
for FILE in "${{ env.TARBALL }}" "${{ env.TARBALL }}.sha256"; do for FILE in "${{ env.TARBALL }}" "${{ env.TARBALL }}.sha256"; do
# Remove any existing asset with the same name so re-runs stay clean # Remove any existing asset with the same name so re-runs stay clean
EXISTING=$(curl -s \ 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" \ "${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases/${RELEASE_ID}/assets" \
| jq -r ".[] | select(.name == \"${FILE}\") | .id") | jq -r ".[] | select(.name == \"${FILE}\") | .id")
for AID in $EXISTING; do for AID in $EXISTING; do
curl -sX DELETE \ 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}" "${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases/${RELEASE_ID}/assets/${AID}"
done done
curl -sX POST \ curl -sX POST \
-H "Authorization: token ${{ secrets.GITEA_TOKEN }}" \ -H "Authorization: token ${{ secrets.GITHUB_TOKEN }}" \
-H "Content-Type: application/octet-stream" \ -H "Content-Type: application/octet-stream" \
"${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases/${RELEASE_ID}/assets?name=${FILE}" \ "${{ gitea.server_url }}/api/v1/repos/${{ gitea.repository }}/releases/${RELEASE_ID}/assets?name=${FILE}" \
--data-binary "@${FILE}" --fail-with-body --data-binary "@${FILE}" --fail-with-body
@@ -101,7 +124,7 @@ jobs:
runs-on: ubuntu-latest runs-on: ubuntu-latest
steps: steps:
- uses: actions/checkout@v4 - uses: builtin:checkout
- name: Log in to Docker Hub - name: Log in to Docker Hub
run: echo "${{ secrets.DOCKERHUB_TOKEN }}" | docker login -u bendik --password-stdin run: echo "${{ secrets.DOCKERHUB_TOKEN }}" | docker login -u bendik --password-stdin
@@ -118,7 +141,7 @@ jobs:
runs-on: aarch64 runs-on: aarch64
steps: steps:
- uses: actions/checkout@v4 - uses: builtin:checkout
- name: Compute checksums and update PKGBUILD - name: Compute checksums and update PKGBUILD
run: | run: |
@@ -132,6 +155,32 @@ jobs:
sed -i "s/sha256sums_x86_64=('.*')/sha256sums_x86_64=('${SUM_X86}')/" aur/PKGBUILD sed -i "s/sha256sums_x86_64=('.*')/sha256sums_x86_64=('${SUM_X86}')/" aur/PKGBUILD
sed -i "s/sha256sums_aarch64=('.*')/sha256sums_aarch64=('${SUM_AARCH}')/" 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 - name: Push to AUR
env: env:
AUR_SSH_KEY: ${{ secrets.AUR_SSH_KEY }} AUR_SSH_KEY: ${{ secrets.AUR_SSH_KEY }}
Generated
+1 -1
View File
@@ -345,7 +345,7 @@ dependencies = [
[[package]] [[package]]
name = "cnats" name = "cnats"
version = "0.2.4" version = "0.3.0"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"async-nats", "async-nats",
+3 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "cnats" name = "cnats"
version = "0.2.4" version = "0.3.0"
edition = "2021" edition = "2021"
[lib] [lib]
@@ -63,6 +63,8 @@ web-sys = { version = "0.3", features = [
"RtcIceCandidateInit", "RtcIceCandidateInit",
"RtcPeerConnectionIceEvent", "RtcPeerConnectionIceEvent",
"RtcRtpSender", "RtcRtpSender",
"RtcPeerConnectionState",
"DisplayMediaStreamConstraints",
"RtcTrackEvent", "RtcTrackEvent",
"RtcRtpTransceiver", "RtcRtpTransceiver",
"RtcOfferOptions", "RtcOfferOptions",
+26
View File
@@ -110,6 +110,8 @@ src/
app.rs Leptos UI (login gate + chat console) app.rs Leptos UI (login gate + chat console)
auth.rs shared User type + current_user server fn auth.rs shared User type + current_user server fn
chat.rs shared ChatMessage/rooms + send_message server fn (publishes to NATS) 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/ server/
oidc.rs Kanidm OIDC login/callback/logout handlers oidc.rs Kanidm OIDC login/callback/logout handlers
sse.rs NATS → browser SSE bridge (one subscription per client) sse.rs NATS → browser SSE bridge (one subscription per client)
@@ -117,6 +119,30 @@ src/
style/main.css the console theme 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 ## Notes & production hardening
- Sessions are in-memory (`tower-sessions` `MemoryStore`): restart logs - 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> # Maintainer: Bendik Aagaard Lynghaug <bendik.lynghaug@gmail.com>
pkgname=cnats pkgname=cnats
pkgver=0.2.4 pkgver=0.3.0
pkgrel=1 pkgrel=1
pkgdesc="Web chat over NATS subjects with Kanidm SSO (Leptos SSR)" pkgdesc="Web chat over NATS subjects with Kanidm SSO (Leptos SSR)"
arch=('x86_64' 'aarch64') arch=('x86_64' 'aarch64')
url="https://prosjekt.klingenbergbygg.no/bl/cnats" options=('!strip')
url="https://project.uhhm.no/bl/cnats"
license=('MIT') license=('MIT')
depends=('glibc' 'gcc-libs') depends=('glibc' 'gcc-libs')
optdepends=( optdepends=(
@@ -14,8 +15,8 @@ optdepends=(
provides=('cnats') provides=('cnats')
conflicts=('cnats-git' 'cnats-bin') conflicts=('cnats-git' 'cnats-bin')
backup=('etc/cnats/env') 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_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://prosjekt.klingenbergbygg.no/bl/cnats/releases/download/v${pkgver}/cnats-v${pkgver}-aarch64.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_x86_64=('SKIP')
sha256sums_aarch64=('SKIP') sha256sums_aarch64=('SKIP')
+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 call_state = StoredValue::new_local(CallState::new(room.get_untracked(), me.clone()));
let in_call = call_state.get_value().in_call; 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>); let es_handle = StoredValue::new_local(None::<web_sys::EventSource>);
Effect::new(move |_| { 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 on_join = move |_| {
let cs = call_state.get_value(); let cs = call_state.get_value();
leptos::task::spawn_local(async move { leptos::task::spawn_local(async move {
cs.join().await; cs.join().await;
}); });
}; };
let on_leave = move |_| {
let cs = call_state.get_value();
leptos::task::spawn_local(async move {
cs.leave().await;
});
};
view! { view! {
<div class="call-panel"> <div class="call-panel">
{move || { {move || {
if in_call.get() { if in_call.get() {
view! { view! { <CallActive call_state=call_state.get_value()/> }.into_any()
<CallActive
local_video_ref=local_video_ref
call_state=call_state.get_value()
on_leave=on_leave
/>
}
.into_any()
} else { } else {
view! { view! {
<button class="call-join" on:click=on_join> <button class="call-join" title="join call" aria-label="join call" on:click=on_join>
"☎ join call" "☎"
</button> </button>
} }
.into_any() .into_any()
@@ -515,8 +482,8 @@ fn CallPanel(room: Memo<String>, me: String) -> impl IntoView {
let _ = (room, me); let _ = (room, me);
view! { view! {
<div class="call-panel"> <div class="call-panel">
<button class="call-join" disabled=true> <button class="call-join" title="join call" aria-label="join call" disabled=true>
"☎ join call" "☎"
</button> </button>
</div> </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 /// `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 /// parameter of `CallPanel`'s own `if`/`else` branch - that inlining is what
/// was blowing the compiler's query recursion limit once mesh calling's /// was blowing the compiler's query recursion limit once mesh calling's
/// nested `Show`/`For` landed inside `ChatShell`'s own `Show`. /// 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")] #[cfg(feature = "hydrate")]
#[component] #[component]
fn CallActive( fn CallActive(call_state: crate::webrtc::CallState) -> impl IntoView {
local_video_ref: NodeRef<leptos::html::Video>, let me = call_state.me().to_string();
call_state: crate::webrtc::CallState, let status = call_state.status;
on_leave: impl Fn(leptos::ev::MouseEvent) + 'static, let pinned = RwSignal::new(None::<String>);
) -> impl IntoView { 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! { view! {
<div class="call-active"> <div class="call-active">
<div class="video-grid"> <div class="video-grid" class:staged=move || pinned.get().is_some()>
<video <VideoTile
class="video-tile video-tile-local" id=me_for_tile
node_ref=local_video_ref label="you".to_string()
autoplay=true local=true
muted=true stream=Signal::derive(move || cs.get_value().preview_stream())
playsinline=true status=status.into()
></video> pinned=pinned
auto_pinned=auto_pinned
/>
<For <For
each=move || call_state.peer_streams() each=move || cs.get_value().peer_views()
key=|(id, _)| id.clone() key=|p| p.id.clone()
children=move |(_id, stream)| { children=move |p| {
view! { <PeerVideoTile stream=stream/> } 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> </div>
<button class="call-leave" on:click=on_leave> <div class="call-bar">
"⏏ leave call" <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>
<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> </div>
} }
} }
/// One remote participant's video tile - a plain child component so each #[cfg(feature = "hydrate")]
/// tile gets its own `NodeRef`/effect pair instead of trying to juggle a fn format_elapsed(secs: u64) -> String {
/// `Vec` of node refs by hand in the parent. 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")] #[cfg(feature = "hydrate")]
#[component] #[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(); 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 |_| { Effect::new(move |_| {
let s = stream.get(); let s = stream.get();
if let Some(el) = video_ref.get() { if let Some(el) = video_ref.get() {
if local {
el.set_muted(true);
}
el.set_src_object(s.as_ref()); 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 /// 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 /// A single trickled ICE candidate, JSON-encoded
/// (`RTCIceCandidateInit`, produced client-side). /// (`RTCIceCandidateInit`, produced client-side).
IceCandidate(String), 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. /// 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 //! into Postgres, so history survives restarts and includes messages
//! published by any client on the bus (not just this app). //! 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 async_nats::jetstream;
use futures::StreamExt; 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, /// Runs forever; (re)creates the stream/consumer and retries on any failure,
/// so a NATS or Postgres outage never takes the chat server down. /// so a NATS or Postgres outage never takes the chat server down.
pub async fn run_consumer(nats: async_nats::Client, pool: PgPool) { 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 { loop {
let started = Instant::now();
if let Err(err) = consume(&nats, &pool).await { 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 //! infrastructure this pass deliberately isn't standing up. Mesh topology
//! (every pair of peers connects directly) is fine at the ~4-person scale //! (every pair of peers connects directly) is fine at the ~4-person scale
//! this is scoped for; it does not scale further than that. //! 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")] #![cfg(feature = "hydrate")]
use std::collections::HashMap; use std::collections::HashMap;
@@ -14,11 +19,12 @@ use leptos::prelude::*;
use wasm_bindgen::{prelude::*, JsCast}; use wasm_bindgen::{prelude::*, JsCast};
use wasm_bindgen_futures::JsFuture; use wasm_bindgen_futures::JsFuture;
use web_sys::{ use web_sys::{
MediaStream, MediaStreamConstraints, RtcConfiguration, RtcIceCandidateInit, DisplayMediaStreamConstraints, MediaStream, MediaStreamConstraints, MediaStreamTrack,
RtcIceServer, RtcPeerConnection, RtcSdpType, RtcSessionDescriptionInit, 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"; 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 { struct Peer {
conn: RtcPeerConnection, conn: RtcPeerConnection,
stream: RwSignal<Option<MediaStream>>, 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; /// 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 { pub struct CallState {
room: String, room: String,
me: 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>>, 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 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> { fn new_peer_connection() -> Result<RtcPeerConnection, JsValue> {
@@ -61,14 +95,42 @@ async fn get_local_stream() -> Result<MediaStream, JsValue> {
stream.dyn_into::<MediaStream>() stream.dyn_into::<MediaStream>()
} }
fn attach_local_tracks(pc: &RtcPeerConnection, stream: &MediaStream) { async fn get_display_stream() -> Result<MediaStream, JsValue> {
for track in stream.get_tracks().iter() { let window = web_sys::window().ok_or("no window")?;
if let Ok(track) = track.dyn_into::<web_sys::MediaStreamTrack>() { let media_devices = window.navigator().media_devices()?;
pc.add_track_0(&track, stream); 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` /// Reads the `sdp` field off whatever `create_offer`/`create_answer`
/// resolved to, and builds a fresh `RtcSessionDescriptionInit` from it - /// resolved to, and builds a fresh `RtcSessionDescriptionInit` from it -
/// simpler and more reliable than trying to cast the resolved JsValue /// simpler and more reliable than trying to cast the resolved JsValue
@@ -94,20 +156,39 @@ impl CallState {
room, room,
me, me,
local_stream: RwSignal::new(None), 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), in_call: RwSignal::new(false),
status: RwSignal::new(MediaStatus::default()),
joined_at: RwSignal::new(0.0),
} }
} }
pub fn local_stream(&self) -> ReadSignal<Option<MediaStream>> { pub fn me(&self) -> &str {
self.local_stream.read_only() &self.me
} }
/// Streams for currently-known peers, keyed by their peer id (username). /// What the local preview tile should show: the screen while sharing,
/// Recomputed each call - fine at mesh scale (~4 peers). /// the camera otherwise.
pub fn peer_streams(&self) -> Vec<(String, RwSignal<Option<MediaStream>>)> { pub fn preview_stream(&self) -> Option<MediaStream> {
self.peers self.screen_stream.get().or_else(|| self.local_stream.get())
.with_value(|p| p.iter().map(|(id, peer)| (id.clone(), peer.stream)).collect()) }
/// 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 /// getUserMedia, then broadcast Join so existing participants know to
@@ -120,30 +201,149 @@ impl CallState {
return; return;
} }
} }
self.status.set(MediaStatus::default());
self.joined_at.set(js_sys::Date::now());
self.in_call.set(true); self.in_call.set(true);
let _ = send_signal(self.room.clone(), None, CallSignalKind::Join).await; let _ = send_signal(self.room.clone(), None, CallSignalKind::Join).await;
} }
/// Tears down every peer connection, stops all local tracks (releases /// 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) { pub async fn leave(&self) {
self.peers.update_value(|peers| { self.peers.update(|peers| {
for (_, peer) in peers.drain() { for (_, peer) in peers.drain() {
peer.conn.close(); peer.conn.close();
} }
}); });
if let Some(stream) = self.screen_stream.get_untracked() {
stop_all(&stream);
}
if let Some(stream) = self.local_stream.get_untracked() { if let Some(stream) = self.local_stream.get_untracked() {
for track in stream.get_tracks().iter() { stop_all(&stream);
if let Ok(track) = track.dyn_into::<web_sys::MediaStreamTrack>() {
track.stop();
}
}
} }
self.screen_stream.set(None);
self.local_stream.set(None); self.local_stream.set(None);
self.in_call.set(false); self.in_call.set(false);
let _ = send_signal(self.room.clone(), None, CallSignalKind::Leave).await; 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 /// One incoming signal from `/call-sse/{room}`. Ignores our own
/// broadcasts and anything not addressed to us (directed messages are /// broadcasts and anything not addressed to us (directed messages are
/// broadcast NATS-wide and filtered client-side - see call.rs). /// broadcast NATS-wide and filtered client-side - see call.rs).
@@ -159,43 +359,71 @@ impl CallState {
match kind { match kind {
CallSignalKind::Join => { CallSignalKind::Join => {
// A new peer announced themselves - we initiate the offer. // A new peer announced themselves - we initiate the offer,
self.start_offer(from); // and tell them what we're sending (they know nothing yet).
} self.start_offer(from.clone());
CallSignalKind::Leave => { self.broadcast_status(Some(from));
self.peers.update_value(|peers| {
if let Some(peer) = peers.remove(&from) {
peer.conn.close();
}
});
} }
CallSignalKind::Leave => self.remove_peer(&from),
CallSignalKind::Offer(sdp) => self.handle_offer(from, sdp), CallSignalKind::Offer(sdp) => self.handle_offer(from, sdp),
CallSignalKind::Answer(sdp) => self.handle_answer(from, sdp), CallSignalKind::Answer(sdp) => self.handle_answer(from, sdp),
CallSignalKind::IceCandidate(candidate_json) => { CallSignalKind::IceCandidate(candidate_json) => {
self.handle_ice_candidate(from, 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> { fn ensure_peer(&self, peer_id: &str) -> Option<RtcPeerConnection> {
if let Some(pc) = self if let Some(pc) = self
.peers .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); return Some(pc);
} }
let pc = new_peer_connection().ok()?; let pc = new_peer_connection().ok()?;
let Some(local) = self.local_stream.get_untracked() else { let local = self.local_stream.get_untracked()?;
return None;
}; // Audio from the camera stream; video is whatever we're currently
attach_local_tracks(&pc, &local); // 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 = RwSignal::new(None::<MediaStream>);
{ {
let remote_stream = remote_stream;
let ontrack = Closure::<dyn FnMut(web_sys::RtcTrackEvent)>::new(move |ev: web_sys::RtcTrackEvent| { 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())); pc.set_ontrack(Some(ontrack.as_ref().unchecked_ref()));
ontrack.forget(); ontrack.forget();
@@ -229,12 +457,33 @@ impl CallState {
onicecandidate.forget(); 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( peers.insert(
peer_id.to_string(), peer_id.to_string(),
Peer { Peer {
conn: pc.clone(), conn: pc.clone(),
stream: remote_stream, 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) { fn handle_offer(&self, from: String, sdp: String) {
let Some(pc) = self.ensure_peer(&from) else { let Some(pc) = self.ensure_peer(&from) else {
return; return;
}; };
let room = self.room.clone(); let this = self.clone();
leptos::task::spawn_local(async move { leptos::task::spawn_local(async move {
let remote_desc = RtcSessionDescriptionInit::new(RtcSdpType::Offer); let remote_desc = RtcSessionDescriptionInit::new(RtcSdpType::Offer);
remote_desc.set_sdp(&sdp); remote_desc.set_sdp(&sdp);
@@ -275,6 +544,7 @@ impl CallState {
{ {
return; return;
} }
this.flush_pending_ice(&from, &pc);
let Ok(resolved) = JsFuture::from(pc.create_answer()).await else { let Ok(resolved) = JsFuture::from(pc.create_answer()).await else {
return; return;
}; };
@@ -289,35 +559,50 @@ impl CallState {
{ {
return; 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) { fn handle_answer(&self, from: String, sdp: String) {
let Some(pc) = self let Some(pc) = self
.peers .peers
.with_value(|peers| peers.get(&from).map(|p| p.conn.clone())) .with_untracked(|peers| peers.get(&from).map(|p| p.conn.clone()))
else { else {
return; return;
}; };
let this = self.clone();
leptos::task::spawn_local(async move { leptos::task::spawn_local(async move {
let remote_desc = RtcSessionDescriptionInit::new(RtcSdpType::Answer); let remote_desc = RtcSessionDescriptionInit::new(RtcSdpType::Answer);
remote_desc.set_sdp(&sdp); 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) { 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 { let Ok(parsed) = js_sys::JSON::parse(&candidate_json) else {
return; return;
}; };
let init: RtcIceCandidateInit = parsed.unchecked_into(); 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 { leptos::task::spawn_local(async move {
let _ = JsFuture::from( let _ = JsFuture::from(
pc.add_ice_candidate_with_opt_rtc_ice_candidate_init(Some(&init)), pc.add_ice_candidate_with_opt_rtc_ice_candidate_init(Some(&init)),
+176 -17
View File
@@ -1,8 +1,10 @@
/* ── cnats · message-bus console ───────────────────────────────────────── /* ── 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). */ type: Archivo (UI voice) + IBM Plex Mono (wire voice). */
:root { :root {
color-scheme: light dark;
--ink-0: #060a09; --ink-0: #060a09;
--ink-1: #0b1210; --ink-1: #0b1210;
--ink-2: #101a17; --ink-2: #101a17;
@@ -16,10 +18,34 @@
--signal-dim: #2a8f6c; --signal-dim: #2a8f6c;
--amber: #ffb454; --amber: #ffb454;
--alarm: #ff6b6b; --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; --mono: "IBM Plex Mono", ui-monospace, monospace;
--sans: "Archivo", system-ui, sans-serif; --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; } * { margin: 0; padding: 0; box-sizing: border-box; }
html, body { height: 100%; } html, body { height: 100%; }
@@ -39,8 +65,8 @@ body::before {
pointer-events: none; pointer-events: none;
z-index: 999; z-index: 999;
background: background:
repeating-linear-gradient(0deg, rgba(255, 255, 255, 0.015) 0 1px, transparent 1px 3px), repeating-linear-gradient(0deg, var(--scanline) 0 1px, transparent 1px 3px),
radial-gradient(ellipse 120% 90% at 50% 40%, transparent 55%, rgba(0, 0, 0, 0.45)); radial-gradient(ellipse 120% 90% at 50% 40%, transparent 55%, var(--vignette));
} }
::selection { background: var(--signal); color: var(--ink-0); } ::selection { background: var(--signal); color: var(--ink-0); }
@@ -53,7 +79,7 @@ body::before {
place-items: center; place-items: center;
padding: 2rem; padding: 2rem;
background: 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)); 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%); background: linear-gradient(160deg, var(--ink-2), var(--ink-1) 60%);
padding: 3rem 2.75rem 2.5rem; padding: 3rem 2.75rem 2.5rem;
position: relative; 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; animation: rise 0.5s cubic-bezier(0.2, 0.9, 0.3, 1) both;
} }
@@ -129,7 +155,8 @@ body::before {
.gate-btn:hover { .gate-btn:hover {
background: transparent; background: transparent;
color: var(--signal); 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; } .gate-btn-glyph { font-size: 1.05rem; }
@@ -458,7 +485,7 @@ body::before {
.composer:focus-within { .composer:focus-within {
border-color: var(--signal-dim); 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 { .composer-prompt {
@@ -507,8 +534,8 @@ body::before {
.call-join { .call-join {
font-family: var(--mono); font-family: var(--mono);
font-weight: 600; font-weight: 600;
font-size: 0.78rem; font-size: 1rem;
letter-spacing: 0.08em; line-height: 1;
color: var(--ink-0); color: var(--ink-0);
background: var(--signal); background: var(--signal);
border: none; border: none;
@@ -527,31 +554,161 @@ body::before {
gap: 0.6rem; gap: 0.6rem;
} }
.video-tile { .tile {
position: relative;
width: 200px; width: 200px;
height: 150px; height: 150px;
background: var(--ink-0); background: var(--ink-0);
border: 1px solid var(--line-hot); 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 { .call-leave {
align-self: flex-start;
font-family: var(--mono); font-family: var(--mono);
font-weight: 600; font-weight: 600;
font-size: 0.78rem; font-size: 0.95rem;
letter-spacing: 0.08em; line-height: 18px;
color: var(--text); color: var(--text);
background: none; background: none;
border: 1px solid var(--alarm); border: 1px solid var(--alarm);
padding: 0.5rem 0.9rem; padding: 0.5rem 0.8rem;
cursor: pointer; cursor: pointer;
transition: background 120ms; 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 */ /* motion */
@@ -591,6 +748,8 @@ body::before {
.room-link { border-left: 0; border-bottom: 2px solid transparent; } .room-link { border-left: 0; border-bottom: 2px solid transparent; }
.room-link.active { border-bottom-color: var(--signal); } .room-link.active { border-bottom-color: var(--signal); }
.msg { grid-template-columns: auto 1fr; grid-template-rows: auto auto; } .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-time { grid-row: 1; order: 2; justify-self: end; }
.msg-text { grid-column: 1 / -1; } .msg-text { grid-column: 1 / -1; }
} }