19 Commits
Author SHA1 Message Date
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
Bendik Aagaard Lynghaug e0061af9ce chore: Release cnats version 0.2.1 2026-07-28 10:53:51 +02:00
Bendik Aagaard Lynghaug 64bc5222fc Fix release build: box CallPanel's in-call subtree to avoid recursion limit
Mesh calling's nested Show/For inside ChatShell's own Show inlined as a
type parameter, blowing rustc's query recursion limit during the wasm
release build (CI: "queries overflow the depth limit!"). Splitting the
in-call view into its own component and type-erasing both CallPanel
branches with .into_any() keeps the type shallow instead of raising the
crate-wide limit.
2026-07-28 10:53:32 +02:00
Bendik Aagaard Lynghaug 4cc0103448 chore: Release cnats version 0.2.0 2026-07-27 23:55:00 +02:00
Bendik Aagaard Lynghaug 650ed50c21 Kanidm group-gated rooms and minimal mesh calling
Rooms can now require a Kanidm group (via the `groups` OIDC claim,
mapped by `oauth2 update-claim-map` server-side) - dev/ops require
`developers`, enforced at every message path (send, history, SSE).

Adds a minimal WebRTC mesh call feature scoped to the lobby room,
signaled over a separate `call.room.*` NATS subject kept out of the
chat archive: public STUN only, no TURN, no SFU - small groups on
friendly networks, by design.
2026-07-27 23:52:27 +02:00
Bendik Aagaard LynghaugandClaude Fable 5 46bdad1629 Pass NATS credentials explicitly: async-nats ignores userinfo in the URL
async_nats::connect() silently drops user:pass embedded in NATS_URL,
so the server rejected every connection with an authorization violation.
Parse the URL and feed credentials through ConnectOptions instead.

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-18 22:33:20 +02:00
Bendik Aagaard LynghaugandClaude Fable 5 b4dfe0c5ea Add cargo-release config matching gpupaper release flow
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
2026-07-09 14:44:03 +02:00
16 changed files with 1046 additions and 49 deletions
+39 -6
View File
@@ -36,8 +36,15 @@ jobs:
~/.cargo/bin ~/.cargo/bin
key: ${{ runner.os }}-${{ matrix.arch }}-cargo-leptos-${{ hashFiles('Cargo.lock') }} key: ${{ runner.os }}-${{ matrix.arch }}-cargo-leptos-${{ hashFiles('Cargo.lock') }}
- name: Install cargo-binstall
run: |
command -v cargo-binstall || \
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 - name: Install cargo-leptos
run: command -v cargo-leptos || cargo install cargo-leptos --locked run: command -v cargo-leptos || cargo binstall cargo-leptos --locked --no-confirm
- name: Build - name: Build
run: cargo leptos build --release run: cargo leptos build --release
@@ -58,7 +65,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 }}\"}" \
@@ -67,24 +74,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
@@ -125,6 +132,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
+5 -1
View File
@@ -345,15 +345,17 @@ dependencies = [
[[package]] [[package]]
name = "cnats" name = "cnats"
version = "0.1.0" version = "0.2.6"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"async-nats", "async-nats",
"axum", "axum",
"base64 0.22.1",
"chrono", "chrono",
"console_error_panic_hook", "console_error_panic_hook",
"dotenvy", "dotenvy",
"futures", "futures",
"js-sys",
"leptos", "leptos",
"leptos_axum", "leptos_axum",
"leptos_meta", "leptos_meta",
@@ -368,8 +370,10 @@ dependencies = [
"tower-sessions", "tower-sessions",
"tracing", "tracing",
"tracing-subscriber", "tracing-subscriber",
"url",
"uuid", "uuid",
"wasm-bindgen", "wasm-bindgen",
"wasm-bindgen-futures",
"web-sys", "web-sys",
] ]
+33 -1
View File
@@ -1,6 +1,6 @@
[package] [package]
name = "cnats" name = "cnats"
version = "0.1.0" version = "0.2.6"
edition = "2021" edition = "2021"
[lib] [lib]
@@ -21,6 +21,7 @@ tower = { version = "0.5", optional = true }
tower-http = { version = "0.6", features = ["fs", "trace"], optional = true } tower-http = { version = "0.6", features = ["fs", "trace"], optional = true }
tower-sessions = { version = "0.14", optional = true } tower-sessions = { version = "0.14", optional = true }
async-nats = { version = "0.38", optional = true } async-nats = { version = "0.38", optional = true }
url = { version = "2", optional = true }
sqlx = { version = "0.8", default-features = false, features = [ sqlx = { version = "0.8", default-features = false, features = [
"runtime-tokio", "runtime-tokio",
"tls-rustls", "tls-rustls",
@@ -28,6 +29,12 @@ sqlx = { version = "0.8", default-features = false, features = [
"macros", "macros",
], optional = true } ], optional = true }
openidconnect = { version = "4", optional = true } openidconnect = { version = "4", optional = true }
# For pulling the `groups` custom claim out of the already-verified ID
# token's raw JWT payload - openidconnect's Core* type aliases default to
# EmptyAdditionalClaims, and reworking that generic stack for one extra
# field isn't worth it. The token's signature is already checked by
# id_token.claims(...) before this ever runs.
base64 = { version = "0.22", optional = true }
futures = { version = "0.3", optional = true } futures = { version = "0.3", optional = true }
chrono = { version = "0.4", features = ["serde"], optional = true } chrono = { version = "0.4", features = ["serde"], optional = true }
uuid = { version = "1", features = ["v4"], optional = true } uuid = { version = "1", features = ["v4"], optional = true }
@@ -38,12 +45,33 @@ tracing-subscriber = { version = "0.3", features = ["env-filter"], optional = tr
# --- browser only --- # --- browser only ---
wasm-bindgen = { version = "0.2", optional = true } wasm-bindgen = { version = "0.2", optional = true }
wasm-bindgen-futures = { version = "0.4", optional = true }
js-sys = { version = "0.3", optional = true }
console_error_panic_hook = { version = "0.1", optional = true } console_error_panic_hook = { version = "0.1", optional = true }
web-sys = { version = "0.3", features = [ web-sys = { version = "0.3", features = [
"EventSource", "EventSource",
"MessageEvent", "MessageEvent",
"HtmlElement", "HtmlElement",
"Element", "Element",
# --- WebRTC mesh calling ---
"RtcPeerConnection",
"RtcConfiguration",
"RtcIceServer",
"RtcSdpType",
"RtcSessionDescriptionInit",
"RtcIceCandidate",
"RtcIceCandidateInit",
"RtcPeerConnectionIceEvent",
"RtcRtpSender",
"RtcTrackEvent",
"RtcRtpTransceiver",
"RtcOfferOptions",
"MediaStream",
"MediaStreamConstraints",
"MediaStreamTrack",
"MediaDevices",
"Navigator",
"HtmlVideoElement",
], optional = true } ], optional = true }
[features] [features]
@@ -51,6 +79,8 @@ default = []
hydrate = [ hydrate = [
"leptos/hydrate", "leptos/hydrate",
"dep:wasm-bindgen", "dep:wasm-bindgen",
"dep:wasm-bindgen-futures",
"dep:js-sys",
"dep:console_error_panic_hook", "dep:console_error_panic_hook",
"dep:web-sys", "dep:web-sys",
] ]
@@ -65,8 +95,10 @@ ssr = [
"dep:tower-http", "dep:tower-http",
"dep:tower-sessions", "dep:tower-sessions",
"dep:async-nats", "dep:async-nats",
"dep:url",
"dep:sqlx", "dep:sqlx",
"dep:openidconnect", "dep:openidconnect",
"dep:base64",
"dep:futures", "dep:futures",
"dep:chrono", "dep:chrono",
"dep:uuid", "dep:uuid",
+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.1.0 pkgver=0.2.6
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')
+12
View File
@@ -0,0 +1,12 @@
publish = false
allow-branch = ["main"]
[[pre-release-replacements]]
file = "aur/PKGBUILD"
search = "pkgver=.*"
replace = "pkgver={{version}}"
[[pre-release-replacements]]
file = "aur/PKGBUILD"
search = "pkgrel=.*"
replace = "pkgrel=1"
+253 -3
View File
@@ -7,7 +7,10 @@ use leptos_router::{
}; };
use crate::auth::{current_user, User}; use crate::auth::{current_user, User};
use crate::chat::{is_valid_room, room_subject, ChatMessage, SendMessage, DEFAULT_ROOM, ROOMS}; use crate::chat::{
is_authorized_for_room, is_valid_room, room_subject, ChatMessage, SendMessage, DEFAULT_ROOM,
ROOMS,
};
#[cfg(feature = "hydrate")] #[cfg(feature = "hydrate")]
use crate::chat::room_history; use crate::chat::room_history;
@@ -75,7 +78,23 @@ fn ChatPage() -> impl IntoView {
{move || { {move || {
user.get() user.get()
.map(|res| match res { .map(|res| match res {
Ok(Some(u)) => view! { <ChatShell room user=u/> }.into_any(), Ok(Some(u)) => {
// `room`'s own Memo only knows about is_valid_room (no
// user context available that early) - fall back to
// DEFAULT_ROOM here too if this user isn't authorized
// for the room the URL asked for, same as an unknown
// room name already does.
let u2 = u.clone();
let effective_room = Memo::new(move |_| {
let r = room.get();
if is_authorized_for_room(&u2, &r) {
r
} else {
DEFAULT_ROOM.to_string()
}
});
view! { <ChatShell room=effective_room user=u/> }.into_any()
}
_ => view! { <LoginGate/> }.into_any(), _ => view! { <LoginGate/> }.into_any(),
}) })
}} }}
@@ -83,6 +102,28 @@ fn ChatPage() -> impl IntoView {
} }
} }
// ---------------------------------------------------------------------------
// brand mark - the "voice pulse" ornament, currentColor so it always
// matches whatever text color surrounds it (sidebar wordmark vs. the much
// larger gate title)
// ---------------------------------------------------------------------------
#[component]
fn PulseMark() -> impl IntoView {
view! {
<svg
class="pulse-mark"
viewBox="-15 -70 430 140"
fill="none"
stroke="currentColor"
stroke-width="18"
stroke-linecap="round"
>
<path d="M0,0 H60 Q80,-24 100,0 Q120,24 140,0 Q162,-52 184,0 Q206,52 228,0 Q250,-24 270,0 Q290,24 310,0 H400"></path>
</svg>
}
}
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
// unauthenticated: the gate // unauthenticated: the gate
// --------------------------------------------------------------------------- // ---------------------------------------------------------------------------
@@ -94,8 +135,10 @@ fn LoginGate() -> impl IntoView {
<div class="gate-card"> <div class="gate-card">
<div class="gate-badge">"MESSAGE BUS · AUTH REQUIRED"</div> <div class="gate-badge">"MESSAGE BUS · AUTH REQUIRED"</div>
<h1 class="gate-title"> <h1 class="gate-title">
<PulseMark/>
"CN" <span class="gate-title-accent">"ATS"</span> "CN" <span class="gate-title-accent">"ATS"</span>
</h1> </h1>
<p class="wordmark-sub gate-tagline">"chat over the bus"</p>
<p class="gate-sub"> <p class="gate-sub">
"Realtime chat carried on NATS subjects. Identity issued by your Kanidm realm — no separate passwords, no local accounts." "Realtime chat carried on NATS subjects. Identity issued by your Kanidm realm — no separate passwords, no local accounts."
</p> </p>
@@ -210,11 +253,13 @@ fn ChatShell(room: Memo<String>, user: User) -> impl IntoView {
}; };
let me = user.username.clone(); let me = user.username.clone();
let call_me = me.clone();
view! { view! {
<div class="console"> <div class="console">
<nav class="rail"> <nav class="rail">
<div class="wordmark"> <div class="wordmark">
<PulseMark/>
"CN" <span class="wordmark-accent">"ATS"</span> "CN" <span class="wordmark-accent">"ATS"</span>
<span class="wordmark-sub">"chat over the bus"</span> <span class="wordmark-sub">"chat over the bus"</span>
</div> </div>
@@ -223,7 +268,8 @@ fn ChatShell(room: Memo<String>, user: User) -> impl IntoView {
<div class="rooms"> <div class="rooms">
{ROOMS {ROOMS
.iter() .iter()
.map(|(name, desc)| { .filter(|(name, _, _)| is_authorized_for_room(&user, name))
.map(|(name, desc, _)| {
let name = *name; let name = *name;
let desc = *desc; let desc = *desc;
view! { view! {
@@ -276,6 +322,10 @@ fn ChatShell(room: Memo<String>, user: User) -> impl IntoView {
</div> </div>
</header> </header>
<Show when=move || room.get() == "lobby">
<CallPanel room=room me=call_me.clone()/>
</Show>
<div class="stream" node_ref=list_ref> <div class="stream" node_ref=list_ref>
<Show when=move || messages.with(|m| m.is_empty())> <Show when=move || messages.with(|m| m.is_empty())>
<div class="stream-empty"> <div class="stream-empty">
@@ -357,6 +407,206 @@ fn open_event_source(
Some(es) Some(es)
} }
/// Minimal mesh call entry point, scoped to the `lobby` room only (see
/// `ChatShell`'s `<Show when=move || room.get() == "lobby">`) - the
/// signaling itself (`call.rs`, `/call-sse/{room}`) is already generic
/// per-room, so widening this later is a one-line UI change, not an
/// architectural one. `CallState` (browser-only: it holds `web_sys`
/// types) can't exist in the `ssr` build at all, so the two targets get
/// entirely separate bodies rather than sharing signals across the gate.
#[component]
fn CallPanel(room: Memo<String>, me: String) -> impl IntoView {
#[cfg(feature = "hydrate")]
{
use crate::webrtc::CallState;
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 |_| {
let room_name = room.get();
es_handle.update_value(|es| {
if let Some(es) = es.take() {
es.close();
}
});
es_handle.set_value(open_call_event_source(&room_name, call_state.get_value()));
});
on_cleanup(move || {
es_handle.update_value(|es| {
if let Some(es) = es.take() {
es.close();
}
});
let cs = call_state.get_value();
if cs.in_call.get_untracked() {
leptos::task::spawn_local(async move {
cs.leave().await;
});
}
});
// Mirror the local MediaStream into the preview <video> element -
// `srcObject` has no HTML attribute form, has to be set via JS.
//
// Also force `.muted` via the JS property here, not just the
// `muted` attribute on the element below: every video tile in this
// app is built client-side via `document.createElement` (never
// parsed from HTML), and browsers only seed the live `.muted`
// property from the `muted` *attribute* for parser-inserted
// elements. Without this, the local preview plays back the user's
// own mic through their speakers - audible as feedback in a
// solo call.
Effect::new(move |_| {
let stream = call_state.get_value().local_stream().get();
if let Some(el) = local_video_ref.get() {
el.set_muted(true);
el.set_src_object(stream.as_ref());
}
});
let on_join = move |_| {
let cs = call_state.get_value();
leptos::task::spawn_local(async move {
cs.join().await;
});
};
let on_leave = move |_| {
let cs = call_state.get_value();
leptos::task::spawn_local(async move {
cs.leave().await;
});
};
view! {
<div class="call-panel">
{move || {
if in_call.get() {
view! {
<CallActive
local_video_ref=local_video_ref
call_state=call_state.get_value()
on_leave=on_leave
/>
}
.into_any()
} else {
view! {
<button class="call-join" on:click=on_join>
"☎ join call"
</button>
}
.into_any()
}
}}
</div>
}
.into_any()
}
// Must mirror the hydrate branch's default (not-in-call) markup exactly -
// hydration reconciles this SSR output against what the hydrate branch
// above expects to find, and an empty div here (vs. the button hydrate
// wants) is a hydration mismatch that panics and traps the whole wasm
// instance, killing all reactivity on the page.
#[cfg(not(feature = "hydrate"))]
{
let _ = (room, me);
view! {
<div class="call-panel">
<button class="call-join" disabled=true>
"☎ join call"
</button>
</div>
}
.into_any()
}
}
/// The in-call subtree (video grid + leave button), 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`.
#[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 {
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>
<For
each=move || call_state.peer_streams()
key=|(id, _)| id.clone()
children=move |(_id, stream)| {
view! { <PeerVideoTile stream=stream/> }
}
/>
</div>
<button class="call-leave" on:click=on_leave>
"⏏ leave call"
</button>
</div>
}
}
/// One remote participant's video tile - a plain child component so each
/// tile gets its own `NodeRef`/effect pair instead of trying to juggle a
/// `Vec` of node refs by hand in the parent.
#[cfg(feature = "hydrate")]
#[component]
fn PeerVideoTile(stream: RwSignal<Option<web_sys::MediaStream>>) -> impl IntoView {
let video_ref = NodeRef::<leptos::html::Video>::new();
Effect::new(move |_| {
let s = stream.get();
if let Some(el) = video_ref.get() {
el.set_src_object(s.as_ref());
}
});
view! { <video class="video-tile" node_ref=video_ref autoplay=true playsinline=true></video> }
}
/// Bridges `/call-sse/{room}` into `CallState::handle_signal`. A separate
/// function from `open_event_source` (rather than a shared generic) since
/// the event name differs: `sse::call_events` emits a custom `"signal"`
/// SSE event, not the default unnamed one, so this needs
/// `add_event_listener_with_callback` instead of `set_onmessage` (which
/// only fires for the default event type).
#[cfg(feature = "hydrate")]
fn open_call_event_source(
room: &str,
call_state: crate::webrtc::CallState,
) -> Option<web_sys::EventSource> {
use wasm_bindgen::{prelude::Closure, JsCast};
use web_sys::{EventSource, MessageEvent};
let es = EventSource::new(&format!("/call-sse/{room}")).ok()?;
let on_signal = Closure::<dyn FnMut(MessageEvent)>::new(move |ev: MessageEvent| {
if let Some(data) = ev.data().as_string() {
if let Ok(signal) = serde_json::from_str::<crate::call::CallSignal>(&data) {
call_state.handle_signal(signal.from, signal.to, signal.kind);
}
}
});
es.add_event_listener_with_callback("signal", on_signal.as_ref().unchecked_ref())
.ok()?;
on_signal.forget();
Some(es)
}
#[component] #[component]
fn NotFound() -> impl IntoView { fn NotFound() -> impl IntoView {
view! { view! {
+5
View File
@@ -8,6 +8,11 @@ pub struct User {
pub sub: String, pub sub: String,
pub username: String, pub username: String,
pub display_name: String, pub display_name: String,
/// Kanidm group membership, from the `groups` OIDC claim (see
/// `oauth2 update-claim-map`). Fixed at login time - not re-checked
/// live, so a group change only takes effect on the next login.
#[serde(default)]
pub groups: Vec<String>,
} }
pub const SESSION_USER_KEY: &str = "user"; pub const SESSION_USER_KEY: &str = "user";
+86
View File
@@ -0,0 +1,86 @@
use leptos::prelude::*;
use serde::{Deserialize, Serialize};
/// Call signaling for a room, kept entirely separate from `chat::ChatMessage`
/// - deliberately a different NATS subject namespace (`call.room.<room>`,
/// not `chat.room.<room>`) so `server::store`'s JetStream/Postgres archive
/// (scoped to `chat.room.*` only) never sees it. Ephemeral SDP/ICE has no
/// business being durably stored.
pub fn call_subject(room: &str) -> String {
format!("call.room.{room}")
}
/// `from`/`to` are peer ids - currently just the signed-in username (same
/// identity chat messages use). `to: None` is a room-wide broadcast (only
/// `Join`/`Leave` use this); everything else is directed at one peer, with
/// every other browser in the room ignoring it client-side. Mesh calls at
/// this scale (~4 people) don't need per-peer NATS subjects - broadcast +
/// client-side filter is the simplest thing that works.
#[derive(Clone, Debug, Serialize, Deserialize)]
pub struct CallSignal {
pub room: String,
pub from: String,
pub to: Option<String>,
pub kind: CallSignalKind,
}
#[derive(Clone, Debug, Serialize, Deserialize)]
#[serde(tag = "kind", content = "data")]
pub enum CallSignalKind {
/// Announces presence to the room; existing participants respond by
/// initiating an offer to the new peer.
Join,
Leave,
Offer(String),
Answer(String),
/// A single trickled ICE candidate, JSON-encoded
/// (`RTCIceCandidateInit`, produced client-side).
IceCandidate(String),
}
/// Publishes a call-signaling message to the room's signaling subject.
/// Requires a signed-in, room-authorized session - same two checks as
/// `chat::send_message`, and deliberately not shared via a helper since
/// the existing chat functions already establish that each enforcement
/// point re-does this small check inline rather than factoring it out.
#[server]
pub async fn send_signal(
room: String,
to: Option<String>,
kind: CallSignalKind,
) -> Result<(), ServerFnError> {
use crate::auth::{User, SESSION_USER_KEY};
use crate::chat::is_authorized_for_room;
use crate::server::AppState;
if !crate::chat::is_valid_room(&room) {
return Err(ServerFnError::new("unknown room"));
}
let session: tower_sessions::Session = leptos_axum::extract().await?;
let Some(user) = session
.get::<User>(SESSION_USER_KEY)
.await
.map_err(|e| ServerFnError::new(e.to_string()))?
else {
return Err(ServerFnError::new("not signed in"));
};
if !is_authorized_for_room(&user, &room) {
return Err(ServerFnError::new("not authorized for this room"));
}
let state = expect_context::<AppState>();
let signal = CallSignal {
room: room.clone(),
from: user.username,
to,
kind,
};
let payload = serde_json::to_vec(&signal).map_err(|e| ServerFnError::new(e.to_string()))?;
state
.nats
.publish(call_subject(&room), payload.into())
.await
.map_err(|e| ServerFnError::new(format!("nats publish failed: {e}")))?;
Ok(())
}
+33 -9
View File
@@ -3,23 +3,42 @@ use serde::{Deserialize, Serialize};
/// Rooms available in the UI. Each maps to the NATS subject /// Rooms available in the UI. Each maps to the NATS subject
/// `chat.room.<name>`, so any other NATS client on the bus can join in. /// `chat.room.<name>`, so any other NATS client on the bus can join in.
pub const ROOMS: &[(&str, &str)] = &[ /// The third field is the Kanidm group (via the `groups` OIDC claim,
("lobby", "general traffic"), /// see `oauth2 update-claim-map`) required to read/post in that room -
("dev", "build & ship"), /// `None` means open to anyone in `cnats_users`.
("ops", "incidents & infra"), pub const ROOMS: &[(&str, &str, Option<&str>)] = &[
("random", "off the record"), ("lobby", "general traffic", None),
("dev", "build & ship", Some("developers")),
("ops", "incidents & infra", Some("developers")),
("random", "off the record", None),
]; ];
pub const DEFAULT_ROOM: &str = "lobby"; pub const DEFAULT_ROOM: &str = "lobby";
pub fn is_valid_room(room: &str) -> bool { pub fn is_valid_room(room: &str) -> bool {
ROOMS.iter().any(|(name, _)| *name == room) ROOMS.iter().any(|(name, _, _)| *name == room)
} }
pub fn room_subject(room: &str) -> String { pub fn room_subject(room: &str) -> String {
format!("chat.room.{room}") format!("chat.room.{room}")
} }
/// Whether `user` may read/post in `room`. `false` for an unknown room -
/// callers should check `is_valid_room` separately if they need to tell
/// "unknown room" and "not authorized" apart in the error they return.
/// Synchronous and I/O-free: the user's groups are already baked into
/// their session (from the `groups` OIDC claim at login), so this never
/// needs a live Kanidm round-trip - and never gets more current than
/// that login until they sign in again.
pub fn is_authorized_for_room(user: &crate::auth::User, room: &str) -> bool {
ROOMS
.iter()
.find(|(name, _, _)| *name == room)
.is_some_and(|(_, _, required_group)| {
required_group.is_none_or(|g| user.groups.iter().any(|ug| ug == g))
})
}
/// A single chat message as it travels over NATS (JSON-encoded payload). /// A single chat message as it travels over NATS (JSON-encoded payload).
#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)] #[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
pub struct ChatMessage { pub struct ChatMessage {
@@ -62,6 +81,9 @@ pub async fn send_message(room: String, text: String) -> Result<(), ServerFnErro
else { else {
return Err(ServerFnError::new("not signed in")); return Err(ServerFnError::new("not signed in"));
}; };
if !is_authorized_for_room(&user, &room) {
return Err(ServerFnError::new("not authorized for this room"));
}
let state = expect_context::<AppState>(); let state = expect_context::<AppState>();
let now = chrono::Utc::now(); let now = chrono::Utc::now();
@@ -94,13 +116,15 @@ pub async fn room_history(room: String) -> Result<Vec<ChatMessage>, ServerFnErro
return Err(ServerFnError::new("unknown room")); return Err(ServerFnError::new("unknown room"));
} }
let session: tower_sessions::Session = leptos_axum::extract().await?; let session: tower_sessions::Session = leptos_axum::extract().await?;
if session let Some(user) = session
.get::<User>(SESSION_USER_KEY) .get::<User>(SESSION_USER_KEY)
.await .await
.map_err(|e| ServerFnError::new(e.to_string()))? .map_err(|e| ServerFnError::new(e.to_string()))?
.is_none() else {
{
return Err(ServerFnError::new("not signed in")); return Err(ServerFnError::new("not signed in"));
};
if !is_authorized_for_room(&user, &room) {
return Err(ServerFnError::new("not authorized for this room"));
} }
let state = expect_context::<AppState>(); let state = expect_context::<AppState>();
+4
View File
@@ -1,10 +1,14 @@
pub mod app; pub mod app;
pub mod auth; pub mod auth;
pub mod call;
pub mod chat; pub mod chat;
#[cfg(feature = "ssr")] #[cfg(feature = "ssr")]
pub mod server; pub mod server;
#[cfg(feature = "hydrate")]
pub mod webrtc;
#[cfg(feature = "hydrate")] #[cfg(feature = "hydrate")]
#[wasm_bindgen::prelude::wasm_bindgen] #[wasm_bindgen::prelude::wasm_bindgen]
pub fn hydrate() { pub fn hydrate() {
+12 -1
View File
@@ -32,7 +32,17 @@ async fn main() -> anyhow::Result<()> {
let nats_url = let nats_url =
std::env::var("NATS_URL").unwrap_or_else(|_| "nats://127.0.0.1:4222".to_string()); std::env::var("NATS_URL").unwrap_or_else(|_| "nats://127.0.0.1:4222".to_string());
tracing::info!(%nats_url, "connecting to NATS"); tracing::info!(%nats_url, "connecting to NATS");
let nats = async_nats::connect(&nats_url).await?; // async-nats does not honor userinfo embedded in the URL, so pass any
// credentials explicitly via ConnectOptions.
let parsed = url::Url::parse(&nats_url)?;
let mut nats_opts = async_nats::ConnectOptions::new();
if !parsed.username().is_empty() {
nats_opts = nats_opts.user_and_password(
parsed.username().to_string(),
parsed.password().unwrap_or_default().to_string(),
);
}
let nats = nats_opts.connect(&nats_url).await?;
let database_url = std::env::var("DATABASE_URL") let database_url = std::env::var("DATABASE_URL")
.unwrap_or_else(|_| "postgres://cnats:cnats@127.0.0.1:5432/cnats".to_string()); .unwrap_or_else(|_| "postgres://cnats:cnats@127.0.0.1:5432/cnats".to_string());
@@ -80,6 +90,7 @@ async fn main() -> anyhow::Result<()> {
.route("/auth/callback", get(oidc::callback)) .route("/auth/callback", get(oidc::callback))
.route("/auth/logout", get(oidc::logout)) .route("/auth/logout", get(oidc::logout))
.route("/sse/{room}", get(sse::room_events)) .route("/sse/{room}", get(sse::room_events))
.route("/call-sse/{room}", get(sse::call_events))
.route("/api/{*fn_name}", any(server_fn_handler)) .route("/api/{*fn_name}", any(server_fn_handler))
.leptos_routes_with_context( .leptos_routes_with_context(
&state, &state,
+41
View File
@@ -185,10 +185,20 @@ pub async fn callback(
.map(|n| n.as_str().to_string()) .map(|n| n.as_str().to_string())
.unwrap_or_else(|| username.clone()); .unwrap_or_else(|| username.clone());
// `groups` is a custom claim (Kanidm `oauth2 update-claim-map`), not
// something the Core* typed claims struct above knows about. The
// signature is already verified by `id_token.claims(...)` above, so
// re-reading the same payload's raw JSON for one more field is safe -
// just a plain field extraction, not a second verification step.
// IdToken's Serialize impl (not Display - it has none) produces the
// raw compact JWT string "header.payload.signature".
let groups = extract_groups_claim(&id_token);
let user = User { let user = User {
sub: claims.subject().as_str().to_string(), sub: claims.subject().as_str().to_string(),
username, username,
display_name, display_name,
groups,
}; };
// Rotate the session id on privilege change, then store the user. // Rotate the session id on privilege change, then store the user.
@@ -207,3 +217,34 @@ pub async fn logout(session: Session) -> Result<Redirect, HandlerError> {
session.flush().await.map_err(internal)?; session.flush().await.map_err(internal)?;
Ok(Redirect::to("/")) Ok(Redirect::to("/"))
} }
/// Pulls the `groups` custom claim (Kanidm `oauth2 update-claim-map`) out
/// of an ID token's raw JWT payload. `IdToken`'s `Serialize` impl (it has
/// no `Display`) produces the compact "header.payload.signature" string,
/// which is where this reads from - the signature itself is never
/// re-checked here, that already happened via `id_token.claims(...)`
/// before this is called. Defensive by design: any parse failure (no
/// claim, wrong shape) just yields no groups rather than failing login.
fn extract_groups_claim<T: serde::Serialize>(id_token: &T) -> Vec<String> {
use base64::Engine;
let Ok(serde_json::Value::String(compact)) = serde_json::to_value(id_token) else {
return Vec::new();
};
let Some(payload_b64) = compact.split('.').nth(1) else {
return Vec::new();
};
let Ok(payload_bytes) = base64::engine::general_purpose::URL_SAFE_NO_PAD.decode(payload_b64)
else {
return Vec::new();
};
let Ok(payload) = serde_json::from_slice::<serde_json::Value>(&payload_bytes) else {
return Vec::new();
};
payload
.get("groups")
.and_then(|g| g.as_array())
.map(|arr| arr.iter().filter_map(|v| v.as_str().map(String::from)).collect())
.unwrap_or_default()
}
+65 -14
View File
@@ -12,28 +12,45 @@ use futures::{Stream, StreamExt};
use tower_sessions::Session; use tower_sessions::Session;
use crate::auth::{User, SESSION_USER_KEY}; use crate::auth::{User, SESSION_USER_KEY};
use crate::chat::{is_valid_room, room_subject}; use crate::call::call_subject;
use crate::chat::{is_authorized_for_room, is_valid_room, room_subject};
use super::AppState; use super::AppState;
/// GET /sse/{room} — stream the room's NATS subject to the browser. /// Shared by `room_events` and `call_events`: signed in, valid room, and
/// authorized for it (Kanidm group gate, checked against the session's
/// own `groups` - see `chat::is_authorized_for_room`). Note this is only
/// checked once, at connect time - a long-lived SSE stream doesn't get
/// re-checked if the user's groups change mid-connection (same kind of
/// staleness the "still signed in at all" check already has).
async fn authorize_room_stream(
room: &str,
session: &Session,
) -> Result<(), (StatusCode, &'static str)> {
let user = session
.get::<User>(SESSION_USER_KEY)
.await
.ok()
.flatten();
let Some(user) = user else {
return Err((StatusCode::UNAUTHORIZED, "sign in first"));
};
if !is_valid_room(room) {
return Err((StatusCode::NOT_FOUND, "unknown room"));
}
if !is_authorized_for_room(&user, room) {
return Err((StatusCode::FORBIDDEN, "not authorized for this room"));
}
Ok(())
}
/// GET /sse/{room} — stream the room's chat NATS subject to the browser.
pub async fn room_events( pub async fn room_events(
Path(room): Path<String>, Path(room): Path<String>,
State(state): State<AppState>, State(state): State<AppState>,
session: Session, session: Session,
) -> Result<Sse<impl Stream<Item = Result<Event, Infallible>>>, (StatusCode, &'static str)> { ) -> Result<Sse<impl Stream<Item = Result<Event, Infallible>>>, (StatusCode, &'static str)> {
let signed_in = session authorize_room_stream(&room, &session).await?;
.get::<User>(SESSION_USER_KEY)
.await
.ok()
.flatten()
.is_some();
if !signed_in {
return Err((StatusCode::UNAUTHORIZED, "sign in first"));
}
if !is_valid_room(&room) {
return Err((StatusCode::NOT_FOUND, "unknown room"));
}
let subscriber = state let subscriber = state
.nats .nats
@@ -56,3 +73,37 @@ pub async fn room_events(
.text("ping"), .text("ping"),
)) ))
} }
/// GET /call-sse/{room} — stream the room's call-signaling NATS subject
/// (SDP offers/answers, ICE candidates). Deliberately a separate subject
/// namespace (`call.room.*`, not `chat.room.*`) so this never touches the
/// JetStream/Postgres chat archive - ephemeral signaling has no business
/// being durably stored.
pub async fn call_events(
Path(room): Path<String>,
State(state): State<AppState>,
session: Session,
) -> Result<Sse<impl Stream<Item = Result<Event, Infallible>>>, (StatusCode, &'static str)> {
authorize_room_stream(&room, &session).await?;
let subscriber = state
.nats
.subscribe(call_subject(&room))
.await
.map_err(|e| {
tracing::error!("nats subscribe failed: {e}");
(StatusCode::BAD_GATEWAY, "message bus unavailable")
})?;
let stream = subscriber.map(|msg| {
Ok(Event::default()
.event("signal")
.data(String::from_utf8_lossy(&msg.payload).into_owned()))
});
Ok(Sse::new(stream).keep_alive(
KeepAlive::new()
.interval(Duration::from_secs(15))
.text("ping"),
))
}
+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;
} }
} }
+328
View File
@@ -0,0 +1,328 @@
//! Minimal mesh WebRTC calling, browser-only. Public STUN, no TURN - calls
//! across hostile NATs (symmetric NAT, restrictive corporate networks)
//! simply won't connect. That's a known, accepted limitation, not a bug to
//! fix later: real NAT traversal needs a TURN relay, which is real
//! 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.
#![cfg(feature = "hydrate")]
use std::collections::HashMap;
use js_sys::{Array, Reflect};
use leptos::prelude::*;
use wasm_bindgen::{prelude::*, JsCast};
use wasm_bindgen_futures::JsFuture;
use web_sys::{
MediaStream, MediaStreamConstraints, RtcConfiguration, RtcIceCandidateInit,
RtcIceServer, RtcPeerConnection, RtcSdpType, RtcSessionDescriptionInit,
};
use crate::call::{send_signal, CallSignalKind};
const STUN_URL: &str = "stun:stun.l.google.com:19302";
/// One remote participant: their peer connection plus the remote stream
/// their video tile renders once `ontrack` fires.
struct Peer {
conn: RtcPeerConnection,
stream: RwSignal<Option<MediaStream>>,
}
/// Call state for one room. Lives for as long as the user is in the call;
/// dropped (and everything torn down) on "leave".
#[derive(Clone)]
pub struct CallState {
room: String,
me: String,
local_stream: RwSignal<Option<MediaStream>>,
peers: StoredValue<HashMap<String, Peer>, LocalStorage>,
pub in_call: RwSignal<bool>,
}
fn new_peer_connection() -> Result<RtcPeerConnection, JsValue> {
let config = RtcConfiguration::new();
let ice_server = RtcIceServer::new();
ice_server.set_urls(&JsValue::from_str(STUN_URL));
let servers = Array::new();
servers.push(&ice_server);
config.set_ice_servers(&servers);
RtcPeerConnection::new_with_configuration(&config)
}
async fn get_local_stream() -> Result<MediaStream, JsValue> {
let window = web_sys::window().ok_or("no window")?;
let media_devices = window.navigator().media_devices()?;
let constraints = MediaStreamConstraints::new();
constraints.set_video(&JsValue::TRUE);
constraints.set_audio(&JsValue::TRUE);
let promise = media_devices.get_user_media_with_constraints(&constraints)?;
let stream = JsFuture::from(promise).await?;
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);
}
}
}
/// 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
/// directly, since its concrete type varies by browser. Returns the sdp
/// string alongside the desc (not `desc.get_sdp()` afterwards - that
/// returns `Option<String>`, and we already have it as a plain `String`
/// right here).
fn session_description_from_resolved(
resolved: &JsValue,
sdp_type: RtcSdpType,
) -> Result<(RtcSessionDescriptionInit, String), JsValue> {
let sdp = Reflect::get(resolved, &JsValue::from_str("sdp"))?
.as_string()
.ok_or("resolved session description had no sdp field")?;
let desc = RtcSessionDescriptionInit::new(sdp_type);
desc.set_sdp(&sdp);
Ok((desc, sdp))
}
impl CallState {
pub fn new(room: String, me: String) -> Self {
Self {
room,
me,
local_stream: RwSignal::new(None),
peers: StoredValue::new_local(HashMap::new()),
in_call: RwSignal::new(false),
}
}
pub fn local_stream(&self) -> ReadSignal<Option<MediaStream>> {
self.local_stream.read_only()
}
/// 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())
}
/// getUserMedia, then broadcast Join so existing participants know to
/// offer us a connection.
pub async fn join(&self) {
match get_local_stream().await {
Ok(stream) => self.local_stream.set(Some(stream)),
Err(e) => {
leptos::logging::error!("getUserMedia failed: {e:?}");
return;
}
}
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.
pub async fn leave(&self) {
self.peers.update_value(|peers| {
for (_, peer) in peers.drain() {
peer.conn.close();
}
});
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();
}
}
}
self.local_stream.set(None);
self.in_call.set(false);
let _ = send_signal(self.room.clone(), None, CallSignalKind::Leave).await;
}
/// 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).
pub fn handle_signal(&self, from: String, to: Option<String>, kind: CallSignalKind) {
if from == self.me || !self.in_call.get_untracked() {
return;
}
if let Some(to) = &to {
if *to != self.me {
return;
}
}
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();
}
});
}
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)
}
}
}
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()))
{
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 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()));
});
pc.set_ontrack(Some(ontrack.as_ref().unchecked_ref()));
ontrack.forget();
}
{
let room = self.room.clone();
let peer_id = peer_id.to_string();
let onicecandidate =
Closure::<dyn FnMut(web_sys::RtcPeerConnectionIceEvent)>::new(move |ev: web_sys::RtcPeerConnectionIceEvent| {
let Some(candidate) = ev.candidate() else {
return;
};
let Ok(candidate_json) = js_sys::JSON::stringify(&candidate.to_json())
.map(|s| s.as_string().unwrap_or_default())
else {
return;
};
let room = room.clone();
let peer_id = peer_id.clone();
leptos::task::spawn_local(async move {
let _ = send_signal(
room,
Some(peer_id),
CallSignalKind::IceCandidate(candidate_json),
)
.await;
});
});
pc.set_onicecandidate(Some(onicecandidate.as_ref().unchecked_ref()));
onicecandidate.forget();
}
self.peers.update_value(|peers| {
peers.insert(
peer_id.to_string(),
Peer {
conn: pc.clone(),
stream: remote_stream,
},
);
});
Some(pc)
}
fn start_offer(&self, peer_id: String) {
let Some(pc) = self.ensure_peer(&peer_id) else {
return;
};
let room = self.room.clone();
leptos::task::spawn_local(async move {
let Ok(resolved) = JsFuture::from(pc.create_offer()).await else {
return;
};
let Ok((desc, sdp)) = session_description_from_resolved(&resolved, RtcSdpType::Offer)
else {
return;
};
if JsFuture::from(pc.set_local_description(&desc)).await.is_err() {
return;
}
let _ = send_signal(room, Some(peer_id), CallSignalKind::Offer(sdp)).await;
});
}
fn handle_offer(&self, from: String, sdp: String) {
let Some(pc) = self.ensure_peer(&from) else {
return;
};
let room = self.room.clone();
leptos::task::spawn_local(async move {
let remote_desc = RtcSessionDescriptionInit::new(RtcSdpType::Offer);
remote_desc.set_sdp(&sdp);
if JsFuture::from(pc.set_remote_description(&remote_desc))
.await
.is_err()
{
return;
}
let Ok(resolved) = JsFuture::from(pc.create_answer()).await else {
return;
};
let Ok((answer_desc, answer_sdp)) =
session_description_from_resolved(&resolved, RtcSdpType::Answer)
else {
return;
};
if JsFuture::from(pc.set_local_description(&answer_desc))
.await
.is_err()
{
return;
}
let _ = send_signal(room, 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()))
else {
return;
};
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;
});
}
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();
leptos::task::spawn_local(async move {
let _ = JsFuture::from(
pc.add_ice_candidate_with_opt_rtc_ice_candidate_init(Some(&init)),
)
.await;
});
}
}
+107 -7
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; }
@@ -197,6 +224,22 @@ body::before {
text-transform: uppercase; text-transform: uppercase;
} }
/* brand mark - height:1em scales it to whatever font-size the surrounding
heading uses, so the same element/class works at sidebar (1.6rem) and
gate-title (up to 4.5rem) scale with no per-placement sizing rules. */
.pulse-mark {
display: block;
height: 1em;
width: auto;
color: var(--signal);
margin-bottom: 0.3em;
}
.gate-tagline {
font-size: 0.75rem;
margin-top: 0.6rem;
}
.rail-label { .rail-label {
font-family: var(--mono); font-family: var(--mono);
font-size: 0.6rem; font-size: 0.6rem;
@@ -442,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 {
@@ -480,6 +523,63 @@ body::before {
.composer-send:hover { filter: brightness(1.15); } .composer-send:hover { filter: brightness(1.15); }
.composer-send:disabled { filter: grayscale(0.6) brightness(0.7); cursor: wait; } .composer-send:disabled { filter: grayscale(0.6) brightness(0.7); cursor: wait; }
/* call */
.call-panel {
padding: 0.9rem 1.6rem;
border-bottom: 1px solid var(--line);
background: var(--ink-1);
}
.call-join {
font-family: var(--mono);
font-weight: 600;
font-size: 0.78rem;
letter-spacing: 0.08em;
color: var(--ink-0);
background: var(--signal);
border: none;
padding: 0.55rem 1rem;
cursor: pointer;
transition: filter 120ms;
}
.call-join:hover { filter: brightness(1.15); }
.call-active { display: flex; flex-direction: column; gap: 0.8rem; }
.video-grid {
display: flex;
flex-wrap: wrap;
gap: 0.6rem;
}
.video-tile {
width: 200px;
height: 150px;
background: var(--ink-0);
border: 1px solid var(--line-hot);
object-fit: cover;
}
.video-tile-local { border-color: var(--signal-dim); }
.call-leave {
align-self: flex-start;
font-family: var(--mono);
font-weight: 600;
font-size: 0.78rem;
letter-spacing: 0.08em;
color: var(--text);
background: none;
border: 1px solid var(--alarm);
padding: 0.5rem 0.9rem;
cursor: pointer;
transition: background 120ms;
}
.call-leave:hover { background: color-mix(in srgb, var(--alarm) 12%, transparent); }
/* motion */ /* motion */
@keyframes rise { @keyframes rise {