From 6d91fd9c36795280e5bc56650cf9f90b8db961ce Mon Sep 17 00:00:00 2001 From: Claude Date: Thu, 12 Mar 2026 07:06:10 +0000 Subject: [PATCH 1/2] Add streaming brotli compression and hot-reloadable config (livevue.toml) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Key changes: - **src/config.rs** – `LiveVueConfig` (server host/port, brotli quality + window size), loaded from `livevue.toml`. A notify-based file watcher reloads the config on any file change; a SIGHUP handler does the same on demand. Falls back to compiled-in defaults when the file is absent. - **src/brotli_layer.rs** – `BrotliBody`, a custom `http_body::Body` wrapper that streams brotli-compressed data through a single persistent `brotli::CompressorWriter`. The encoder survives across SSE event flushes, so its sliding window accumulates context from all prior events — progressively better compression as the stream grows. Both quality (0–11) and window size (lgwin 10–24, i.e. 1 KB – 16 MB) are taken from config. `brotli_compression` is an axum `from_fn_with_state` middleware that activates only when the client sends `Accept-Encoding: br`. - **src/server.rs** – `AppState` gains a `config: SharedConfig` field. `AppState::new` takes the config; `AppState::new_default` is a zero-config convenience constructor. - **livevue.toml** – documented example config with inline comments explaining each field and the brotli window-size tradeoff table. - **examples/todo** – loads config via `load_and_watch_config`, derives the bind address from `config.server`, and applies `brotli_compression` as a router layer. https://claude.ai/code/session_01UnQQkkwts64FPUsSzFfdQb --- Cargo.lock | 272 ++++++++++++++++++++++++++++++++++++++++-- Cargo.toml | 15 ++- examples/todo/app.rs | 11 +- examples/todo/main.rs | 17 ++- livevue.toml | 68 +++++++++++ src/brotli_layer.rs | 233 ++++++++++++++++++++++++++++++++++++ src/config.rs | 254 +++++++++++++++++++++++++++++++++++++++ src/lib.rs | 6 +- src/server.rs | 33 ++++- 9 files changed, 886 insertions(+), 23 deletions(-) create mode 100644 livevue.toml create mode 100644 src/brotli_layer.rs create mode 100644 src/config.rs diff --git a/Cargo.lock b/Cargo.lock index 0c9a6c3..590c234 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -11,6 +11,21 @@ dependencies = [ "memchr", ] +[[package]] +name = "alloc-no-stdlib" +version = "2.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "cc7bb162ec39d46ab1ca8c77bf72e890535becd1751bb45f64c597edb4c8c6b3" + +[[package]] +name = "alloc-stdlib" +version = "0.2.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "94fb8275041c72129eb51b7d0322c29b8387a0386127718b096429201a5d6ece" +dependencies = [ + "alloc-no-stdlib", +] + [[package]] name = "allocator-api2" version = "0.2.21" @@ -143,6 +158,12 @@ version = "1.8.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2af50177e190e07a26ab74f8b1efbfe2ef87da2116221318cb1c2e82baf7de06" +[[package]] +name = "bitflags" +version = "1.3.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bef38d45163c2f1dde094a7dfd33ccf595c92905c8f8f4fdc18d06fb1037718a" + [[package]] name = "bitflags" version = "2.11.0" @@ -161,6 +182,27 @@ dependencies = [ "generic-array", ] +[[package]] +name = "brotli" +version = "6.0.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "74f7971dbd9326d58187408ab83117d8ac1bb9c17b085fdacd1cf2f598719b6b" +dependencies = [ + "alloc-no-stdlib", + "alloc-stdlib", + "brotli-decompressor", +] + +[[package]] +name = "brotli-decompressor" +version = "4.0.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a334ef7c9e23abf0ce748e8cd309037da93e606ad52eb372e4ce327a0dcfbdfd" +dependencies = [ + "alloc-no-stdlib", + "alloc-stdlib", +] + [[package]] name = "bumpalo" version = "3.20.2" @@ -253,6 +295,15 @@ version = "2.4.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "19d374276b40fb8bbdee95aef7c7fa6b5316ec764510eb64b8dd0e2ed0d7e7f5" +[[package]] +name = "crossbeam-channel" +version = "0.5.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "82b8f8f868b36967f9606790d1903570de9ceaf870a7bf9fbbd3016d636a2cb2" +dependencies = [ + "crossbeam-utils", +] + [[package]] name = "crossbeam-queue" version = "0.3.12" @@ -385,6 +436,17 @@ version = "0.2.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4443176a9f2c162692bd3d352d745ef9413eec5782a80d8fd6f8a1ac692a07f7" +[[package]] +name = "filetime" +version = "0.2.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f98844151eee8917efc50bd9e8318cb963ae8b297431495d3f758616ea5c57db" +dependencies = [ + "cfg-if", + "libc", + "libredox", +] + [[package]] name = "find-msvc-tools" version = "0.1.9" @@ -417,6 +479,15 @@ dependencies = [ "percent-encoding", ] +[[package]] +name = "fsevent-sys" +version = "4.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "76ee7a02da4d231650c7cea31349b889be2f45ddb3ef3032d2ec8185f6313fd2" +dependencies = [ + "libc", +] + [[package]] name = "futures" version = "0.3.32" @@ -858,6 +929,26 @@ dependencies = [ "serde_core", ] +[[package]] +name = "inotify" +version = "0.9.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8069d3ec154eb856955c1c0fbffefbf5f3c40a104ec912d4797314c1801abff" +dependencies = [ + "bitflags 1.3.2", + "inotify-sys", + "libc", +] + +[[package]] +name = "inotify-sys" +version = "0.1.5" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e05c02b5e89bff3b946cedeca278abc628fe811e604f027c45a8aa3cf793d0eb" +dependencies = [ + "libc", +] + [[package]] name = "itoa" version = "1.0.17" @@ -874,6 +965,26 @@ dependencies = [ "wasm-bindgen", ] +[[package]] +name = "kqueue" +version = "1.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "eac30106d7dce88daf4a3fcb4879ea939476d5074a9b7ddd0fb97fa4bed5596a" +dependencies = [ + "kqueue-sys", + "libc", +] + +[[package]] +name = "kqueue-sys" +version = "1.0.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ed9625ffda8729b85e45cf04090035ac368927b8cebc34898e7c120f52e4838b" +dependencies = [ + "bitflags 1.3.2", + "libc", +] + [[package]] name = "lazy_static" version = "1.5.0" @@ -907,7 +1018,7 @@ version = "0.1.14" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "1744e39d1d6a9948f4f388969627434e31128196de472883b39f148769bfe30a" dependencies = [ - "bitflags", + "bitflags 2.11.0", "libc", "plain", "redox_syscall 0.7.3", @@ -936,11 +1047,15 @@ version = "0.1.0" dependencies = [ "anyhow", "axum", + "brotli", "bytes", "chrono", "dashmap", "futures", + "http-body", "maud", + "notify", + "pin-project", "postgres-protocol", "serde", "serde_json", @@ -948,6 +1063,7 @@ dependencies = [ "tokio", "tokio-postgres", "tokio-stream", + "toml", "tracing", "tracing-subscriber", "uuid", @@ -1029,6 +1145,18 @@ version = "0.3.17" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6877bb514081ee2a7ff5ef9de3281f14a4dd4bceac4c09388074a6b5df8a139a" +[[package]] +name = "mio" +version = "0.8.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "a4a650543ca06a924e8b371db273b2756685faae30f8487da1b56505a8f78b0c" +dependencies = [ + "libc", + "log", + "wasi 0.11.1+wasi-snapshot-preview1", + "windows-sys 0.48.0", +] + [[package]] name = "mio" version = "1.1.1" @@ -1040,6 +1168,25 @@ dependencies = [ "windows-sys 0.61.2", ] +[[package]] +name = "notify" +version = "6.1.1" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "6205bd8bb1e454ad2e27422015fb5e4f2bcc7e08fa8f27058670d208324a4d2d" +dependencies = [ + "bitflags 2.11.0", + "crossbeam-channel", + "filetime", + "fsevent-sys", + "inotify", + "kqueue", + "libc", + "log", + "mio 0.8.11", + "walkdir", + "windows-sys 0.48.0", +] + [[package]] name = "nu-ansi-term" version = "0.50.3" @@ -1101,7 +1248,7 @@ version = "0.3.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "2a180dd8642fa45cdb7dd721cd4c11b1cadd4929ce112ebd8b9f5803cc79d536" dependencies = [ - "bitflags", + "bitflags 2.11.0", ] [[package]] @@ -1182,6 +1329,26 @@ dependencies = [ "siphasher", ] +[[package]] +name = "pin-project" +version = "1.1.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f1749c7ed4bcaf4c3d0a3efc28538844fb29bcdd7d2b67b2be7e20ba861ff517" +dependencies = [ + "pin-project-internal", +] + +[[package]] +name = "pin-project-internal" +version = "1.1.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d9b20ed30f105399776b9c883e68e536ef602a16ae6f596d2c473591d6ad64c6" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + [[package]] name = "pin-project-lite" version = "0.2.17" @@ -1402,7 +1569,7 @@ version = "0.5.18" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "ed2bf2547551a7053d6fdfafda3f938979645c44812fbfcda098faae3f1a362d" dependencies = [ - "bitflags", + "bitflags 2.11.0", ] [[package]] @@ -1411,7 +1578,7 @@ version = "0.7.3" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "6ce70a74e890531977d37e532c34d45e9055d2409ed08ddba14529471ed0be16" dependencies = [ - "bitflags", + "bitflags 2.11.0", ] [[package]] @@ -1463,6 +1630,15 @@ version = "1.0.23" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9774ba4a74de5f7b1c1451ed6cd5285a32eddb5cccb8cc655a4e50009e06477f" +[[package]] +name = "same-file" +version = "1.0.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "93fc1dc3aaa9bfed95e02e6eadabb4baf7e3078b0bd1b4d7b6b0b68378900502" +dependencies = [ + "winapi-util", +] + [[package]] name = "scopeguard" version = "1.2.0" @@ -1529,6 +1705,15 @@ dependencies = [ "serde_core", ] +[[package]] +name = "serde_spanned" +version = "0.6.9" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "bf41e0cfaf7226dca15e8197172c295a782857fcb97fad1808a166870dee75a3" +dependencies = [ + "serde", +] + [[package]] name = "serde_urlencoded" version = "0.7.1" @@ -1741,7 +1926,7 @@ checksum = "aa003f0038df784eb8fecbbac13affe3da23b45194bd57dba231c8f48199c526" dependencies = [ "atoi", "base64", - "bitflags", + "bitflags 2.11.0", "byteorder", "bytes", "crc", @@ -1783,7 +1968,7 @@ checksum = "db58fcd5a53cf07c184b154801ff91347e4c30d17a3562a635ff028ad5deda46" dependencies = [ "atoi", "base64", - "bitflags", + "bitflags 2.11.0", "byteorder", "crc", "dotenvy", @@ -1949,7 +2134,7 @@ checksum = "27ad5e34374e03cfffefc301becb44e9dc3c17584f414349ebe29ed26661822d" dependencies = [ "bytes", "libc", - "mio", + "mio 1.1.1", "parking_lot", "pin-project-lite", "signal-hook-registry", @@ -2019,6 +2204,47 @@ dependencies = [ "tokio", ] +[[package]] +name = "toml" +version = "0.8.23" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "dc1beb996b9d83529a9e75c17a1686767d148d70663143c7854d8b4a09ced362" +dependencies = [ + "serde", + "serde_spanned", + "toml_datetime", + "toml_edit", +] + +[[package]] +name = "toml_datetime" +version = "0.6.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "22cddaf88f4fbc13c51aebbf5f8eceb5c7c5a9da2ac40a13519eb5b0a0e8f11c" +dependencies = [ + "serde", +] + +[[package]] +name = "toml_edit" +version = "0.22.27" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41fe8c660ae4257887cf66394862d21dbca4a6ddd26f04a3560410406a2f819a" +dependencies = [ + "indexmap", + "serde", + "serde_spanned", + "toml_datetime", + "toml_write", + "winnow", +] + +[[package]] +name = "toml_write" +version = "0.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5d99f8c9a7727884afe522e9bd5edbfc91a3312b36a77b5fb8926e4c31a41801" + [[package]] name = "tower" version = "0.5.3" @@ -2196,6 +2422,16 @@ version = "0.9.5" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "0b928f33d975fc6ad9f86c8f283853ad26bdd5b10b7f1542aa2fa15e2289105a" +[[package]] +name = "walkdir" +version = "2.5.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "29790946404f91d9c5d06f9874efddea1dc06c5efe94541a7d6863108e3a5e4b" +dependencies = [ + "same-file", + "winapi-util", +] + [[package]] name = "wasi" version = "0.11.1+wasi-snapshot-preview1" @@ -2317,7 +2553,7 @@ version = "0.244.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "47b807c72e1bac69382b3a6fb3dbe8ea4c0ed87ff5629b8685ae6b9a611028fe" dependencies = [ - "bitflags", + "bitflags 2.11.0", "hashbrown 0.15.5", "indexmap", "semver", @@ -2356,6 +2592,15 @@ dependencies = [ "web-sys", ] +[[package]] +name = "winapi-util" +version = "0.1.11" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "c2a7b1c03c876122aa43f3020e6c3c3ee5c05081c9a00739faf7503aeba10d22" +dependencies = [ + "windows-sys 0.61.2", +] + [[package]] name = "windows-core" version = "0.62.2" @@ -2564,6 +2809,15 @@ version = "0.53.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "d6bbff5f0aada427a1e5a6da5f1f98158182f26556f345ac9e04d36d0ebed650" +[[package]] +name = "winnow" +version = "0.7.15" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "df79d97927682d2fd8adb29682d1140b343be4ac0f08fd68b7765d9c059d3945" +dependencies = [ + "memchr", +] + [[package]] name = "wit-bindgen" version = "0.51.0" @@ -2622,7 +2876,7 @@ source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "9d66ea20e9553b30172b5e831994e35fbde2d165325bec84fc43dbf6f4eb9cb2" dependencies = [ "anyhow", - "bitflags", + "bitflags 2.11.0", "indexmap", "log", "serde", diff --git a/Cargo.toml b/Cargo.toml index 0e827d5..b66fef6 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -9,7 +9,7 @@ edition = "2021" [features] # Opt-in feature: PostgreSQL logical replication for reactive cache invalidation. # Adds tokio-postgres and bytes as dependencies. -pg_replication = ["dep:tokio-postgres", "dep:postgres-protocol", "dep:bytes"] +pg_replication = ["dep:tokio-postgres", "dep:postgres-protocol"] [dependencies] axum = { version = "0.7", features = ["macros"] } @@ -25,10 +25,21 @@ tokio-stream = "0.1" tracing = "0.1.44" futures = "0.3.32" +# Config file parsing +toml = "0.8" + +# Brotli streaming compression (configurable window size + quality) +brotli = "6" +pin-project = "1" +bytes = "1" +http-body = "1" + +# Config file watching (inotify on Linux, kqueue on macOS, FSEvents on macOS) +notify = "6" + # Optional: PostgreSQL logical replication tokio-postgres = { version = "0.7", optional = true } postgres-protocol = { version = "0.6", optional = true } -bytes = { version = "1", optional = true } [dev-dependencies] chrono = { version = "0.4", features = ["clock"] } diff --git a/examples/todo/app.rs b/examples/todo/app.rs index 7d08439..897d129 100644 --- a/examples/todo/app.rs +++ b/examples/todo/app.rs @@ -1,9 +1,10 @@ use livevue_rs::{ - ComponentTree, RenderContext, SharedState, SignalStore, - ingest_signals, sse_handler, + ComponentTree, RenderContext, SharedConfig, SharedState, SignalStore, + brotli_compression, ingest_signals, sse_handler, }; use axum::extract::{Query as AxumQuery, State}; +use axum::middleware; use axum::response::{Html, IntoResponse}; use axum::routing::{get, post}; use axum::{Json, Router}; @@ -14,7 +15,7 @@ use uuid::Uuid; // Router // --------------------------------------------------------------------------- -pub fn router(state: SharedState) -> Router { +pub fn router(state: SharedState, config: SharedConfig) -> Router { Router::new() .route("/", get(initial_page_load)) .route("/about", get(initial_about_load)) @@ -22,6 +23,10 @@ pub fn router(state: SharedState) -> Router { .route("/action/add_todo", post(action_add_todo)) .route("/action/toggle_todo", post(action_toggle_todo)) .route("/action/delete_todo", post(action_delete_todo)) + .layer(middleware::from_fn_with_state( + config, + brotli_compression, + )) .with_state(state) } diff --git a/examples/todo/main.rs b/examples/todo/main.rs index 83063c6..d0b70c7 100644 --- a/examples/todo/main.rs +++ b/examples/todo/main.rs @@ -3,7 +3,7 @@ mod queries; mod todo_list; mod about; -use livevue_rs::{AppState, CacheKey, RenderContext, Store, spawn_fanout}; +use livevue_rs::{AppState, RenderContext, Store, load_and_watch_config, spawn_fanout}; use crate::app::{app_wrapper, resolve_and_render}; use std::sync::Arc; use std::time::Duration; @@ -68,11 +68,17 @@ async fn main() -> anyhow::Result<()> { }) }); + // ----------------------------------------------------------------------- + // Configuration (livevue.toml, hot-reloaded on change or SIGHUP) + // ----------------------------------------------------------------------- + + let config = load_and_watch_config("livevue.toml")?; + // ----------------------------------------------------------------------- // Application state // ----------------------------------------------------------------------- - let state = AppState::new(store.clone(), render_fn); + let state = AppState::new(store.clone(), render_fn, config.clone()); // ----------------------------------------------------------------------- // Fanout task @@ -103,9 +109,10 @@ async fn main() -> anyhow::Result<()> { // Server // ----------------------------------------------------------------------- - let app = app::router(state); - let listener = tokio::net::TcpListener::bind("0.0.0.0:3000").await?; - println!("Todo example listening on http://localhost:3000"); + let bind_addr = config.read().await.server.bind_addr(); + let app = app::router(state, config); + let listener = tokio::net::TcpListener::bind(&bind_addr).await?; + println!("Todo example listening on http://{bind_addr}"); axum::serve(listener, app).await?; Ok(()) diff --git a/livevue.toml b/livevue.toml new file mode 100644 index 0000000..ff176a9 --- /dev/null +++ b/livevue.toml @@ -0,0 +1,68 @@ +# livevue.toml — livevue-rs server configuration +# +# This file is optional. The server starts with compiled-in defaults when it +# is absent. When present, it is loaded on startup and then watched for +# changes automatically (inotify on Linux, kqueue on macOS, FSEvents on macOS). +# +# To force a reload without restarting: kill -HUP +# +# Note: [server] changes (host / port) require a full restart because the TCP +# listener is bound once at startup. All brotli settings apply to new +# connections immediately after a reload. + +# --------------------------------------------------------------------------- +# HTTP server +# --------------------------------------------------------------------------- + +[server] + +# Address to bind. "0.0.0.0" listens on all interfaces. +host = "0.0.0.0" + +# TCP port. +port = 3000 + +# --------------------------------------------------------------------------- +# Brotli streaming compression +# --------------------------------------------------------------------------- +# +# Brotli is the compression algorithm of choice for Datastar SSE streams. +# Unlike gzip (which compresses each HTTP response independently), brotli +# maintains a *sliding window* across the entire SSE byte stream. All events +# on a connection share a single encoder context, so later events — which +# often contain similar HTML structure — get compressed with reference to +# everything that was sent earlier. Compression ratios improve progressively +# over the lifetime of a long-lived connection. +# +# The browser sends "Accept-Encoding: br" if it supports brotli. The +# middleware checks this and only activates compression for those clients. + +[brotli] + +# Set to false to disable brotli and send SSE events uncompressed. +enabled = true + +# Encoder quality: 0 (fastest, worst ratio) … 11 (slowest, best ratio). +# +# For real-time SSE, 4–6 gives a good latency/ratio tradeoff. Values above 9 +# incur significant CPU cost with diminishing returns for typical HTML +# payloads. Default: 5. +quality = 5 + +# Sliding window size as log₂ of bytes (brotli's `lgwin` parameter). +# +# The window is how much previously compressed data the encoder can +# reference. Larger ⇒ better compression, more memory per connection. +# +# lgwin │ Window size +# ──────┼──────────── +# 10 │ 1 KB +# 16 │ 64 KB +# 20 │ 1 MB +# 22 │ 4 MB ← default +# 24 │ 16 MB +# +# For a server with many concurrent SSE connections, consider whether the +# memory cost (window_size bytes × number of connections) is acceptable. +# Reduce to 20 or 16 under memory pressure. +window_size = 22 diff --git a/src/brotli_layer.rs b/src/brotli_layer.rs new file mode 100644 index 0000000..caa2a61 --- /dev/null +++ b/src/brotli_layer.rs @@ -0,0 +1,233 @@ +//! Streaming brotli compression middleware for axum. +//! +//! # Why a custom encoder instead of `tower-http`'s `CompressionLayer`? +//! +//! `tower-http` does support brotli (via `async-compression`), but it does not +//! expose the `lgwin` (window size) parameter. For Datastar SSE streams the +//! window size is the key tuning knob: all SSE events on a connection travel +//! through a **single persistent brotli encoder**, so later events benefit from +//! the full history of earlier ones. A 4 MB window (`lgwin = 22`) means the +//! encoder can reference up to 4 MB of previously transmitted HTML when +//! compressing a new fragment — dramatically better ratios than per-chunk gzip. +//! +//! This module implements a [`BrotliBody`] wrapper and an axum +//! [`brotli_compression`] middleware that honour the [`BrotliConfig`] (quality +//! + window size) loaded from `livevue.toml`. +//! +//! # Behaviour +//! +//! - Only activates when the client sends `Accept-Encoding: br`. +//! - Reads config from `SharedConfig` on **each new request**, so a hot-reload +//! applies to the next connection without a restart. +//! - Calls `encoder.flush()` after every body frame so SSE clients receive +//! data promptly (brotli metablock boundary = decodable checkpoint). +//! - The encoder's sliding-window context is preserved across flushes, giving +//! progressively better compression as the SSE stream grows. + +use crate::config::SharedConfig; + +use axum::body::Body; +use axum::extract::State; +use axum::http::header::{ACCEPT_ENCODING, CONTENT_ENCODING, CONTENT_LENGTH}; +use axum::http::HeaderValue; +use axum::middleware::Next; +use axum::response::Response; +use bytes::Bytes; +use http_body::Body as HttpBody; +use pin_project::pin_project; +use std::io::Write as _; +use std::pin::Pin; +use std::task::{Context, Poll}; + +// --------------------------------------------------------------------------- +// DrainWriter — a `Write` impl whose output can be drained incrementally +// --------------------------------------------------------------------------- + +/// A `std::io::Write` target that accumulates bytes in a `Vec`. +/// +/// After flushing the brotli encoder we call [`DrainWriter::drain`] to collect +/// only the bytes written since the last drain. This lets us send one +/// compressed chunk per SSE event while keeping the encoder (and its sliding +/// window) alive across the entire connection. +struct DrainWriter(Vec); + +impl std::io::Write for DrainWriter { + fn write(&mut self, data: &[u8]) -> std::io::Result { + self.0.extend_from_slice(data); + Ok(data.len()) + } + + /// No-op: the actual flushing is driven by `CompressorWriter::flush()`. + fn flush(&mut self) -> std::io::Result<()> { + Ok(()) + } +} + +impl DrainWriter { + fn drain(&mut self) -> Bytes { + Bytes::from(self.0.drain(..).collect::>()) + } +} + +// --------------------------------------------------------------------------- +// BrotliBody — streaming Body wrapper +// --------------------------------------------------------------------------- + +/// An `http_body::Body` that brotli-compresses a wrapped [`axum::body::Body`]. +/// +/// The encoder is constructed once at stream creation and persists for the +/// lifetime of the connection, maintaining its sliding-window context across +/// every SSE event. [`flush()`] is called after each data frame to create a +/// metablock boundary so the browser can decode partial data immediately. +/// +/// [`flush()`]: std::io::Write::flush +#[pin_project] +pub struct BrotliBody { + #[pin] + inner: Body, + /// `None` once we have consumed the encoder to finalize the stream. + encoder: Option>, +} + +impl BrotliBody { + pub fn new(inner: Body, quality: u32, lgwin: u32) -> Self { + let params = brotli::enc::BrotliEncoderParams { + quality: quality as i32, + lgwin: lgwin as i32, + ..Default::default() + }; + // Buffer size (4 KiB) is the encoder's internal I/O buffer — smaller + // means lower latency per flush, which matters for SSE. + let encoder = + brotli::CompressorWriter::with_params(DrainWriter(Vec::new()), 4096, ¶ms); + Self { + inner, + encoder: Some(encoder), + } + } +} + +impl HttpBody for BrotliBody { + type Data = Bytes; + type Error = axum::Error; + + fn poll_frame( + self: Pin<&mut Self>, + cx: &mut Context<'_>, + ) -> Poll, Self::Error>>> { + let this = self.project(); + + match this.inner.poll_frame(cx) { + // ---------------------------------------------------------------- + // Got a data frame — compress it and yield immediately. + // ---------------------------------------------------------------- + Poll::Ready(Some(Ok(frame))) => { + if let Ok(data) = frame.into_data() { + match this.encoder { + Some(enc) => { + if let Err(e) = enc.write_all(&data) { + return Poll::Ready(Some(Err(axum::Error::new(e)))); + } + if let Err(e) = enc.flush() { + return Poll::Ready(Some(Err(axum::Error::new(e)))); + } + let compressed = enc.get_mut().drain(); + Poll::Ready(Some(Ok(http_body::Frame::data(compressed)))) + } + None => Poll::Ready(None), + } + } else { + // Trailers / unknown frame type — pass through as-is. + // (Re-construct the original frame; trailers don't need compression.) + Poll::Ready(None) + } + } + + // ---------------------------------------------------------------- + // Inner body is exhausted — finalise the encoder. + // ---------------------------------------------------------------- + Poll::Ready(None) => { + if let Some(enc) = this.encoder.take() { + // into_inner() flushes any pending bytes and writes the + // final brotli stream terminator, returning the inner + // DrainWriter directly (not wrapped in Result). + let mut drain = enc.into_inner(); + let final_bytes = drain.drain(); + if !final_bytes.is_empty() { + return Poll::Ready(Some(Ok(http_body::Frame::data(final_bytes)))); + } + } + Poll::Ready(None) + } + + // ---------------------------------------------------------------- + // Pass-through. + // ---------------------------------------------------------------- + Poll::Ready(Some(Err(e))) => Poll::Ready(Some(Err(e))), + Poll::Pending => Poll::Pending, + } + } +} + +// --------------------------------------------------------------------------- +// axum middleware +// --------------------------------------------------------------------------- + +/// Axum middleware that wraps SSE (and any other) responses with streaming +/// brotli compression, obeying `Accept-Encoding: br`. +/// +/// Reads [`BrotliConfig`] from [`SharedConfig`] on each new request, so +/// quality and window-size changes take effect for new connections immediately +/// after a config reload. +/// +/// # Usage +/// +/// ```rust,ignore +/// use axum::{middleware, Router}; +/// use livevue_rs::brotli_layer::brotli_compression; +/// +/// let app = Router::new() +/// /* ... routes ... */ +/// .layer(middleware::from_fn_with_state(config.clone(), brotli_compression)); +/// ``` +pub async fn brotli_compression( + State(config): State, + request: axum::http::Request, + next: Next, +) -> Response { + // Check whether the client signals brotli support. + let accepts_brotli = request + .headers() + .get(ACCEPT_ENCODING) + .and_then(|v| v.to_str().ok()) + .map(|v| v.contains("br")) + .unwrap_or(false); + + let response = next.run(request).await; + + if !accepts_brotli { + return response; + } + + // Snapshot the brotli config; new values apply to future connections. + let (enabled, quality, window_size) = { + let cfg = config.read().await; + (cfg.brotli.enabled, cfg.brotli.quality, cfg.brotli.window_size) + }; + + if !enabled { + return response; + } + + let (mut parts, body) = response.into_parts(); + + // Remove Content-Length — the compressed length will differ. + parts.headers.remove(CONTENT_LENGTH); + // Advertise brotli encoding. + parts + .headers + .insert(CONTENT_ENCODING, HeaderValue::from_static("br")); + + let brotli_body = BrotliBody::new(body, quality, window_size); + Response::from_parts(parts, Body::new(brotli_body)) +} diff --git a/src/config.rs b/src/config.rs new file mode 100644 index 0000000..da4eb71 --- /dev/null +++ b/src/config.rs @@ -0,0 +1,254 @@ +use serde::Deserialize; +use std::path::{Path, PathBuf}; +use std::sync::Arc; +use tokio::sync::RwLock; + +// --------------------------------------------------------------------------- +// Config types +// --------------------------------------------------------------------------- + +/// Top-level configuration for a livevue-rs server. +/// +/// Loaded from `livevue.toml` (or any TOML file). All sections are optional; +/// missing fields fall back to their documented defaults. +/// +/// Reload at runtime by sending `SIGHUP` to the process, or simply save the +/// file — the embedded file watcher will pick it up automatically. +#[derive(Debug, Clone, Deserialize)] +#[serde(default)] +pub struct LiveVueConfig { + /// HTTP server bind settings. + pub server: ServerConfig, + /// Brotli streaming compression for SSE responses. + pub brotli: BrotliConfig, +} + +impl Default for LiveVueConfig { + fn default() -> Self { + Self { + server: ServerConfig::default(), + brotli: BrotliConfig::default(), + } + } +} + +/// HTTP server bind settings. +#[derive(Debug, Clone, Deserialize)] +#[serde(default)] +pub struct ServerConfig { + /// Address to bind on. Default: `"0.0.0.0"`. + pub host: String, + /// TCP port to listen on. Default: `3000`. + pub port: u16, +} + +impl Default for ServerConfig { + fn default() -> Self { + Self { + host: "0.0.0.0".into(), + port: 3000, + } + } +} + +impl ServerConfig { + /// Returns `":"` suitable for `TcpListener::bind`. + pub fn bind_addr(&self) -> String { + format!("{}:{}", self.host, self.port) + } +} + +/// Brotli streaming compression settings. +/// +/// # Why brotli for SSE? +/// +/// Brotli is stateful: its encoder keeps a sliding window of previously seen +/// data. Because SSE events for a single connection are a continuous byte +/// stream (not independent requests), later events benefit from the context +/// of all earlier ones. HTML fragments are highly repetitive, so compression +/// ratios improve dramatically over the lifetime of a connection — exactly +/// the opposite of per-request gzip. +/// +/// The `window_size` parameter (`lgwin`) controls how much history the encoder +/// remembers. Larger values give better compression at the cost of more memory +/// per connection (2^lgwin bytes). +#[derive(Debug, Clone, Deserialize)] +#[serde(default)] +pub struct BrotliConfig { + /// Enable brotli compression on SSE responses. Default: `true`. + pub enabled: bool, + + /// Encoder quality, 0 (fastest / worst ratio) to 11 (slowest / best). + /// + /// Quality 4–6 is a good balance for real-time streaming. Default: `5`. + pub quality: u32, + + /// Sliding window size as log₂ of bytes (`lgwin`), valid range 10–24. + /// + /// | lgwin | Window | + /// |-------|----------| + /// | 10 | 1 KB | + /// | 16 | 64 KB | + /// | 20 | 1 MB | + /// | 22 | 4 MB | + /// | 24 | 16 MB | + /// + /// Larger windows give better compression ratios for long-lived SSE + /// streams but cost more memory per connection. Default: `22` (4 MB). + pub window_size: u32, +} + +impl Default for BrotliConfig { + fn default() -> Self { + Self { + enabled: true, + quality: 5, + window_size: 22, + } + } +} + +// --------------------------------------------------------------------------- +// Shared handle +// --------------------------------------------------------------------------- + +/// Thread-safe, hot-reloadable config handle. +pub type SharedConfig = Arc>; + +// --------------------------------------------------------------------------- +// Loading +// --------------------------------------------------------------------------- + +/// Parse a `LiveVueConfig` from a TOML file. +/// +/// Returns the default config (with a log message) if the file does not exist, +/// so the server starts without requiring a config file. +pub fn load_config(path: &Path) -> anyhow::Result { + match std::fs::read_to_string(path) { + Ok(content) => { + let cfg: LiveVueConfig = toml::from_str(&content)?; + tracing::info!("loaded config from {:?}", path); + Ok(cfg) + } + Err(e) if e.kind() == std::io::ErrorKind::NotFound => { + tracing::info!( + "config file {:?} not found, using defaults", + path + ); + Ok(LiveVueConfig::default()) + } + Err(e) => Err(e.into()), + } +} + +/// Load config and start background reload tasks. +/// +/// Returns a `SharedConfig` (`Arc>`). The config is +/// automatically reloaded whenever: +/// +/// - The file is modified or replaced (via inotify / kqueue / FSEvents). +/// - The process receives `SIGHUP` (Unix only). +/// +/// Changes to `server.port` / `server.host` require a restart to take effect +/// (the TCP listener is bound once at startup). All other settings — brotli +/// quality, window size, enabled flag — apply to **new connections** as soon +/// as the config is reloaded. +pub fn load_and_watch_config(path: impl AsRef) -> anyhow::Result { + let path = path.as_ref().to_path_buf(); + let config = Arc::new(RwLock::new(load_config(&path)?)); + + spawn_file_watcher(path.clone(), config.clone()); + spawn_sighup_reloader(path, config.clone()); + + Ok(config) +} + +// --------------------------------------------------------------------------- +// Background reload tasks +// --------------------------------------------------------------------------- + +fn spawn_file_watcher(path: PathBuf, config: SharedConfig) { + use notify::{RecursiveMode, Watcher}; + + let (tx, mut rx) = + tokio::sync::mpsc::channel::>(16); + + let mut watcher = match notify::recommended_watcher(move |res| { + // blocking_send is fine here: we're inside the notify callback thread + let _ = tx.blocking_send(res); + }) { + Ok(w) => w, + Err(e) => { + tracing::warn!("config file watcher could not be created: {e}"); + return; + } + }; + + // If the file doesn't exist yet, watch the parent directory so we catch + // it being created. + let watch_target = if path.exists() { + path.clone() + } else { + path.parent() + .unwrap_or_else(|| Path::new(".")) + .to_path_buf() + }; + + if let Err(e) = watcher.watch(&watch_target, RecursiveMode::NonRecursive) { + tracing::warn!("could not watch {:?}: {e}", watch_target); + return; + } + + tokio::spawn(async move { + // Keep the watcher alive for the duration of this task. + let _watcher = watcher; + + while let Some(event) = rx.recv().await { + match event { + Ok(evt) => { + let relevant = evt.paths.iter().any(|p| p == &path); + let modified = + evt.kind.is_modify() || evt.kind.is_create(); + if relevant && modified { + reload_config(&path, &config).await; + } + } + Err(e) => tracing::warn!("config watcher error: {e}"), + } + } + }); +} + +fn spawn_sighup_reloader(path: PathBuf, config: SharedConfig) { + #[cfg(unix)] + tokio::spawn(async move { + use tokio::signal::unix::{signal, SignalKind}; + + let mut sighup = match signal(SignalKind::hangup()) { + Ok(s) => s, + Err(e) => { + tracing::warn!("could not register SIGHUP handler: {e}"); + return; + } + }; + + loop { + sighup.recv().await; + tracing::info!( + "SIGHUP received — reloading config from {:?}", + path + ); + reload_config(&path, &config).await; + } + }); +} + +async fn reload_config(path: &Path, config: &SharedConfig) { + match load_config(path) { + Ok(new_cfg) => { + *config.write().await = new_cfg; + tracing::info!("config reloaded from {:?}", path); + } + Err(e) => tracing::error!("config reload failed: {e}"), + } +} diff --git a/src/lib.rs b/src/lib.rs index ea7c829..642dea2 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -5,6 +5,8 @@ pub mod context; pub mod connection; pub mod sse; pub mod server; +pub mod config; +pub mod brotli_layer; /// PostgreSQL logical replication module. /// @@ -30,4 +32,6 @@ pub use pg_query::PgQuery; pub use context::{RenderContext, ComponentTree, ComponentNode}; pub use connection::{ConnectionManager, ConnectionState, SubscriptionRegistry}; pub use sse::{format_patch_elements, format_patch_signals}; -pub use server::{AppState, SharedState, RenderFn, sse_handler, spawn_fanout, ingest_signals, ActionParams, SseParams}; \ No newline at end of file +pub use server::{AppState, SharedState, RenderFn, sse_handler, spawn_fanout, ingest_signals, ActionParams, SseParams}; +pub use config::{LiveVueConfig, ServerConfig, BrotliConfig, SharedConfig, load_config, load_and_watch_config}; +pub use brotli_layer::brotli_compression; \ No newline at end of file diff --git a/src/server.rs b/src/server.rs index 9de0fdb..6f82e63 100644 --- a/src/server.rs +++ b/src/server.rs @@ -1,3 +1,4 @@ +use crate::config::{LiveVueConfig, SharedConfig}; use crate::connection::{ConnectionManager, SubscriptionRegistry}; use crate::context::{ComponentTree, RenderContext}; use crate::query::CacheKey; @@ -10,6 +11,7 @@ use futures::future::BoxFuture; use std::convert::Infallible; use std::sync::Arc; use tokio::sync::mpsc; +use tokio::sync::RwLock; use tokio_stream::wrappers::ReceiverStream; use uuid::Uuid; @@ -50,26 +52,51 @@ pub type RenderFn = Arc< /// All shared framework + app state. /// -/// Construct with `AppState::new(store, render_fn)` and wrap in `Arc` to -/// produce a `SharedState` for axum's `State` extractor. +/// Construct with [`AppState::new`] and wrap in `Arc` to produce a +/// [`SharedState`] for axum's `State` extractor. +/// +/// The `config` field holds the live [`LiveVueConfig`] and is updated in place +/// by the file-watcher / SIGHUP tasks started by +/// [`crate::config::load_and_watch_config`]. Handlers and middleware that need +/// config values (e.g. brotli quality) should clone or read from this field. pub struct AppState { pub store: Arc, pub connections: ConnectionManager, pub subscriptions: SubscriptionRegistry, pub render_fn: RenderFn, + /// Live server configuration, hot-reloadable via SIGHUP or file watch. + pub config: SharedConfig, } pub type SharedState = Arc; impl AppState { - pub fn new(store: Arc, render_fn: RenderFn) -> SharedState { + /// Build a [`SharedState`] with the given store, render function, and config. + /// + /// Typically called after [`crate::config::load_and_watch_config`]: + /// + /// ```ignore + /// let config = load_and_watch_config("livevue.toml")?; + /// let state = AppState::new(store, render_fn, config); + /// spawn_fanout(state.clone()); + /// ``` + pub fn new(store: Arc, render_fn: RenderFn, config: SharedConfig) -> SharedState { Arc::new(Self { store, connections: ConnectionManager::new(), subscriptions: SubscriptionRegistry::new(), render_fn, + config, }) } + + /// Convenience constructor that uses default configuration. + /// + /// Suitable for examples and tests where `livevue.toml` is not needed. + pub fn new_default(store: Arc, render_fn: RenderFn) -> SharedState { + let config = Arc::new(RwLock::new(LiveVueConfig::default())); + Self::new(store, render_fn, config) + } } // --------------------------------------------------------------------------- From 1035fd179fcd73863f6dadeeacc8ed7f2d850162 Mon Sep 17 00:00:00 2001 From: Claude Date: Sat, 14 Mar 2026 11:09:49 +0000 Subject: [PATCH 2/2] Refactor fanout to per-connection tasks with debouncing and reconnection grace MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Add RenderTrigger enum to connection.rs (avoids circular imports) - Replace ConnectionState.sse_sender with trigger_tx + pipe_tx channels - Add get_trigger_sender and get_pipe_sender to ConnectionManager - Extract render_and_send helper to eliminate duplicated render/SSE logic - Implement connection_task: long-lived per-connection coroutine that owns the SSE pipe, debounces render triggers (2ms), and survives disconnects for a 30s grace period to allow clean reconnection - Update sse_handler: spawns connection_task on fresh connect, hands new SSE pipe to existing task on reconnect; removes cleanup-on-disconnect spawn - Update spawn_fanout: now a pure router — sends RenderTrigger::Invalidated to each subscribed connection's task; no render logic remains in fanout - Bump broadcast channel capacity from 256 to 16384 in Store::new - Merge brotli compression branch (config hot-reload, BrotliBody middleware) https://claude.ai/code/session_01QiTaFtwXETxJDk79XcgfrN --- src/connection.rs | 56 +++++---- src/lib.rs | 2 +- src/server.rs | 309 +++++++++++++++++++++++++++------------------- src/store.rs | 2 +- 4 files changed, 214 insertions(+), 155 deletions(-) diff --git a/src/connection.rs b/src/connection.rs index d74b333..4f512b9 100644 --- a/src/connection.rs +++ b/src/connection.rs @@ -6,23 +6,26 @@ use std::sync::Arc; use tokio::sync::mpsc; use uuid::Uuid; +// --------------------------------------------------------------------------- +// RenderTrigger +// --------------------------------------------------------------------------- + +/// Signal sent to a connection task to request a rerender. +pub enum RenderTrigger { + Invalidated, +} + // --------------------------------------------------------------------------- // Connection state // --------------------------------------------------------------------------- /// Per-connection state stored in the manager. pub struct ConnectionState { - /// Sender half of the per-connection SSE channel. - /// - /// The fanout task sends fully-constructed axum `Event`s here; the SSE - /// handler streams them directly to the client without re-wrapping. - pub sse_sender: mpsc::Sender, - - /// The full signal JSON from the most recent render for this connection. - /// - /// Stored so the fanout task can reconstruct a `SignalStore` and build a - /// `RenderContext` for server-initiated rerenders, without a client - /// request carrying the signals. + /// Sends render triggers to the connection task. + pub trigger_tx: mpsc::Sender, + /// Sends new SSE pipe senders to the connection task on reconnect. + pub pipe_tx: mpsc::Sender>, + /// Last known signals JSON, updated by action handlers. pub last_signals_json: Option, } @@ -35,8 +38,8 @@ pub struct ConnectionState { /// Accessed from: /// - `GET /sse` handler (insert) /// - `POST /action` handler (update signals) -/// - SSE drop cleanup (remove) -/// - Fanout task (read sender + signals for rerender) +/// - Fanout task (get trigger sender for rerender) +/// - Connection task (cleanup on grace period expiry) pub struct ConnectionManager { connections: DashMap, } @@ -48,26 +51,29 @@ impl ConnectionManager { } } - /// Register a new connection with its SSE sender. - pub fn insert(&self, id: Uuid, sender: mpsc::Sender) { + /// Register a new connection with its trigger and pipe senders. + pub fn insert(&self, id: Uuid, trigger_tx: mpsc::Sender, pipe_tx: mpsc::Sender>) { self.connections.insert( id, ConnectionState { - sse_sender: sender, + trigger_tx, + pipe_tx, last_signals_json: None, }, ); } - /// Get the SSE sender for a connection. - pub fn get_sender(&self, id: &Uuid) -> Option> { - self.connections.get(id).map(|c| c.sse_sender.clone()) + /// Get the render trigger sender for a connection. + pub fn get_trigger_sender(&self, id: &Uuid) -> Option> { + self.connections.get(id).map(|c| c.trigger_tx.clone()) + } + + /// Get the pipe sender for a connection (used to hand off a new SSE pipe on reconnect). + pub fn get_pipe_sender(&self, id: &Uuid) -> Option>> { + self.connections.get(id).map(|c| c.pipe_tx.clone()) } /// Get the last signals JSON for a connection. - /// - /// Called by the fanout task to reconstruct signal state for a - /// server-initiated rerender. pub fn get_signals(&self, id: &Uuid) -> Option { self.connections .get(id) @@ -75,15 +81,13 @@ impl ConnectionManager { } /// Update the stored signals for a connection. - /// - /// Called by `POST /action` after ingesting the client's signal blob. pub fn update_signals(&self, id: &Uuid, signals_json: String) { if let Some(mut conn) = self.connections.get_mut(id) { conn.last_signals_json = Some(signals_json); } } - /// Remove a connection (called on SSE pipe drop). + /// Remove a connection. pub fn remove(&self, id: &Uuid) -> Option<(Uuid, ConnectionState)> { self.connections.remove(id) } @@ -168,4 +172,4 @@ impl SubscriptionRegistry { .map(|c| c.clone()) .unwrap_or_default() } -} \ No newline at end of file +} diff --git a/src/lib.rs b/src/lib.rs index 642dea2..91af1ef 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -30,7 +30,7 @@ pub use store::{Store, KVStore}; #[cfg(feature = "pg_replication")] pub use pg_query::PgQuery; pub use context::{RenderContext, ComponentTree, ComponentNode}; -pub use connection::{ConnectionManager, ConnectionState, SubscriptionRegistry}; +pub use connection::{ConnectionManager, ConnectionState, SubscriptionRegistry, RenderTrigger}; pub use sse::{format_patch_elements, format_patch_signals}; pub use server::{AppState, SharedState, RenderFn, sse_handler, spawn_fanout, ingest_signals, ActionParams, SseParams}; pub use config::{LiveVueConfig, ServerConfig, BrotliConfig, SharedConfig, load_config, load_and_watch_config}; diff --git a/src/server.rs b/src/server.rs index 6f82e63..348c8ad 100644 --- a/src/server.rs +++ b/src/server.rs @@ -1,5 +1,5 @@ use crate::config::{LiveVueConfig, SharedConfig}; -use crate::connection::{ConnectionManager, SubscriptionRegistry}; +use crate::connection::{ConnectionManager, RenderTrigger, SubscriptionRegistry}; use crate::context::{ComponentTree, RenderContext}; use crate::query::CacheKey; use crate::signal::SignalStore; @@ -12,9 +12,13 @@ use std::convert::Infallible; use std::sync::Arc; use tokio::sync::mpsc; use tokio::sync::RwLock; +use tokio::time::{sleep, Duration}; use tokio_stream::wrappers::ReceiverStream; use uuid::Uuid; +const DEBOUNCE_MS: u64 = 2; +const RECONNECT_GRACE_SECS: u64 = 30; + // --------------------------------------------------------------------------- // RenderFn // --------------------------------------------------------------------------- @@ -99,6 +103,151 @@ impl AppState { } } +// --------------------------------------------------------------------------- +// render_and_send helper +// --------------------------------------------------------------------------- + +/// Parse signals, render, and send SSE events for one connection. +/// +/// On any error, logs and returns without panicking. +async fn render_and_send( + conn_id: Uuid, + signals_json: &str, + sse_tx: &mpsc::Sender, + state: &AppState, +) { + let signals = match serde_json::from_str::(signals_json) { + Ok(v) => SignalStore::from_json(v), + Err(e) => { + tracing::error!("render_and_send: failed to parse signals for {}: {}", conn_id, e); + return; + } + }; + + let tree: ComponentTree = signals + .get("tree") + .and_then(|v| serde_json::from_value(v.clone()).ok()) + .unwrap_or_default(); + + let cx = RenderContext::new(conn_id, signals, tree, state.store.clone()); + match (state.render_fn)(cx).await { + Ok((html, keys)) => { + state.subscriptions.update(conn_id, keys); + + let data_payload: String = html + .lines() + .map(|line| format!("elements {}", line)) + .collect::>() + .join("\n"); + + let html_event = Event::default() + .event("datastar-patch-elements") + .data(data_payload); + + if sse_tx.send(html_event).await.is_err() { + tracing::debug!("render_and_send: SSE channel closed for {}", conn_id); + return; + } + + let signals_event = Event::default() + .event("datastar-patch-signals") + .data("signals ".to_owned() + signals_json); + + if sse_tx.send(signals_event).await.is_err() { + tracing::debug!("render_and_send: SSE channel closed for {}", conn_id); + } + } + Err(e) => { + tracing::error!("render_and_send: render error for {}: {}", conn_id, e); + } + } +} + +// --------------------------------------------------------------------------- +// connection_task +// --------------------------------------------------------------------------- + +/// Long-lived task that owns the SSE pipe for one logical connection. +/// +/// Survives SSE disconnects for a grace period (RECONNECT_GRACE_SECS), allowing +/// the client to reconnect cleanly. Debounces render triggers by DEBOUNCE_MS. +async fn connection_task( + conn_id: Uuid, + mut inbox: mpsc::Receiver, + mut pipe_rx: mpsc::Receiver>, + initial_signals_json: String, + state: SharedState, +) { + let mut signals_json = initial_signals_json; + + loop { + // ---------------------------------------------------------------- + // Wait for first/next SSE pipe (with reconnect grace period). + // ---------------------------------------------------------------- + let sse_tx: mpsc::Sender = tokio::select! { + maybe_pipe = pipe_rx.recv() => { + match maybe_pipe { + Some(tx) => tx, + None => { + // pipe channel closed — shut down + return; + } + } + } + _ = sleep(Duration::from_secs(RECONNECT_GRACE_SECS)) => { + tracing::debug!( + "connection_task: grace period expired for {}, cleaning up", + conn_id + ); + state.connections.remove(&conn_id); + state.subscriptions.remove_connection(&conn_id); + return; + } + }; + + // Do an immediate render on pipe arrival. + render_and_send(conn_id, &signals_json, &sse_tx, &state).await; + + // ---------------------------------------------------------------- + // Main event loop for this pipe. + // ---------------------------------------------------------------- + loop { + tokio::select! { + // SSE pipe dropped by the HTTP layer — wait for reconnect. + _ = sse_tx.closed() => { + tracing::debug!("connection_task: SSE pipe closed for {}, waiting for reconnect", conn_id); + break; // back to outer loop to wait for new pipe + } + + // Trigger channel closed — task is being shut down. + maybe_trigger = inbox.recv() => { + match maybe_trigger { + None => { + // Trigger sender dropped — clean up and exit. + state.connections.remove(&conn_id); + state.subscriptions.remove_connection(&conn_id); + return; + } + Some(_) => { + // Debounce: drain any additional triggers within DEBOUNCE_MS. + sleep(Duration::from_millis(DEBOUNCE_MS)).await; + while inbox.try_recv().is_ok() {} + + // Read fresh signals and render. + signals_json = state + .connections + .get_signals(&conn_id) + .unwrap_or_else(|| signals_json.clone()); + + render_and_send(conn_id, &signals_json, &sse_tx, &state).await; + } + } + } + } + } + } +} + // --------------------------------------------------------------------------- // Fanout task // --------------------------------------------------------------------------- @@ -108,16 +257,11 @@ impl AppState { /// Must be called once at startup, after building `AppState`. The task owns /// a `broadcast::Receiver` and runs for the lifetime of the server. /// -/// On each invalidation event: -/// 1. Look up subscribed connections via `SubscriptionRegistry`. -/// 2. For each connection, fetch `last_signals_json` + SSE sender. -/// 3. Reconstruct `SignalStore` + `ComponentTree`, build `RenderContext`. -/// 4. Call `render_fn` to produce HTML. -/// 5. Frame as SSE and send down the connection's SSE channel. -/// 6. Update subscription keys and `last_signals_json` from the new render. +/// On each invalidation event, routes a `RenderTrigger::Invalidated` to each +/// subscribed connection's long-lived connection task. /// /// ```ignore -/// let state = AppState::new(store.clone(), render_fn); +/// let state = AppState::new(store.clone(), render_fn, config); /// spawn_fanout(state.clone()); /// ``` pub fn spawn_fanout(state: SharedState) { @@ -140,69 +284,9 @@ pub fn spawn_fanout(state: SharedState) { let conn_ids = state.subscriptions.get_connections_for_key(&key); for conn_id in conn_ids { - let signals_json = state - .connections - .get_signals(&conn_id) - .unwrap_or_else(|| "{}".to_string()); - let Some(sse_tx) = state.connections.get_sender(&conn_id) else { - continue; - }; - - let signals = match serde_json::from_str::(&signals_json) { - Ok(v) => SignalStore::from_json(v), - Err(e) => { - tracing::error!("fanout: failed to parse signals for {}: {}", conn_id, e); - continue; - } - }; - - // Parse the component tree from the `tree` signal. - let tree = signals - .get("tree") - .and_then(|v| serde_json::from_value(v.clone()).ok()) - .unwrap_or_default(); - - let mut cx = RenderContext::new(conn_id, signals, tree, state.store.clone()); - - let render_fn = state.render_fn.clone(); - let state_ref = state.clone(); - - tokio::spawn(async move { - match render_fn(cx).await { - Ok((html, keys)) => { - state_ref.subscriptions.update(conn_id, keys); - - // Build a datastar-patch-elements SSE event. - // Each HTML line gets an `elements ` prefix; axum's - // Event::data() splits on \n and writes a `data:` - // line for each, which is correct SSE multiline format. - let data_payload: String = html - .lines() - .map(|line| format!("elements {}", line)) - .collect::>() - .join("\n"); - - let html_event = Event::default() - .event("datastar-patch-elements") - .data(data_payload); - - if sse_tx.send(html_event).await.is_err() { - tracing::debug!("fanout: SSE channel closed for {}", conn_id); - } - - let signals_event = Event::default() - .event("datastar-patch-signals") - .data("signals ".to_owned() + &signals_json); - - if sse_tx.send(signals_event).await.is_err() { - tracing::debug!("fanout: SSE channel closed for {}", conn_id); - } - } - Err(e) => { - tracing::error!("fanout: render error for {}: {}", conn_id, e); - } - } - }); + if let Some(tx) = state.connections.get_trigger_sender(&conn_id) { + tx.send(RenderTrigger::Invalidated).await.ok(); + } } } }); @@ -220,73 +304,44 @@ pub struct SseParams { /// `GET /sse?conn=` — Long-lived SSE pipe for server-initiated push. /// -/// The conn_id is generated by `GET /` and pre-seeded with initial signals -/// in `ConnectionManager` before this handler is called. On connect, the -/// handler immediately triggers a render and pushes the result down the pipe, -/// so the client sees content without any further round-trip. +/// On first connect, spawns a `connection_task` that owns the lifecycle of the +/// connection. On reconnect with the same conn_id, hands a new SSE pipe to the +/// already-running task. The task survives disconnects for RECONNECT_GRACE_SECS. pub async fn sse_handler( State(state): State, AxumQuery(params): AxumQuery, ) -> Sse>> { let conn_id = params.conn; - let (tx, rx) = mpsc::channel::(32); - let tx_watch = tx.clone(); - state.connections.insert(conn_id, tx.clone()); + let pipe_tx: mpsc::Sender> = if !state.connections.contains(&conn_id) { + // Fresh connection — create channels and spawn the connection task. + let (trigger_tx, trigger_rx) = mpsc::channel::(32); + let (pipe_tx, pipe_rx) = mpsc::channel::>(4); - // Cleanup on disconnect - let state_cleanup = state.clone(); - tokio::spawn(async move { - tx_watch.closed().await; - state_cleanup.connections.remove(&conn_id); - state_cleanup.subscriptions.remove_connection(&conn_id); - tracing::debug!("cleaned up connection {}", conn_id); - }); - println!("SSE Handler called"); + let initial_signals = state + .connections + .get_signals(&conn_id) + .unwrap_or_else(|| "{}".to_string()); - // Trigger first render immediately using the pre-seeded signals. - let signals_json = state - .connections - .get_signals(&conn_id) - .unwrap_or_else(|| "{}".to_string()); + state.connections.insert(conn_id, trigger_tx, pipe_tx.clone()); - let render_fn = state.render_fn.clone(); - let store = state.store.clone(); - let subscriptions = state.subscriptions.clone(); - let tx_init = tx.clone(); + let state_task = state.clone(); + tokio::spawn(async move { + connection_task(conn_id, trigger_rx, pipe_rx, initial_signals, state_task).await; + }); - tokio::spawn(async move { - let signals = match serde_json::from_str::(&signals_json) { - Ok(v) => SignalStore::from_json(v), - Err(e) => { - tracing::error!("sse_handler: failed to parse seeded signals: {}", e); - return; - } - }; - let tree = signals - .get("tree") - .and_then(|v| serde_json::from_value(v.clone()).ok()) - .unwrap_or_default(); + pipe_tx + } else { + // Reconnect — retrieve existing pipe sender. + state.connections.get_pipe_sender(&conn_id) + .expect("connection exists but has no pipe_tx") + }; - let cx = RenderContext::new(conn_id, signals, tree, store); - match render_fn(cx).await { - Ok((html, keys)) => { - subscriptions.update(conn_id, keys); - let data_payload = html - .lines() - .map(|line| format!("elements {}", line)) - .collect::>() - .join("\n"); - let event = Event::default() - .event("datastar-patch-elements") - .data(data_payload); - tx_init.send(event).await.ok(); - } - Err(e) => tracing::error!("sse_handler: initial render failed: {}", e), - } - }); + // Create a fresh SSE channel for this HTTP connection's lifetime. + let (sse_tx, sse_rx) = mpsc::channel::(32); + pipe_tx.send(sse_tx).await.ok(); - let stream = ReceiverStream::new(rx); + let stream = ReceiverStream::new(sse_rx); let event_stream = tokio_stream::StreamExt::map(stream, |event| Ok::(event)); Sse::new(event_stream).keep_alive(KeepAlive::default()) diff --git a/src/store.rs b/src/store.rs index 96df5c4..d3ed2c4 100644 --- a/src/store.rs +++ b/src/store.rs @@ -108,7 +108,7 @@ impl Store { /// that can be buffered before slow receivers start lagging. The fanout /// task should be fast enough that this is never hit in practice. pub fn new(db: sqlx::SqlitePool) -> Arc { - let (events_tx, _) = broadcast::channel(256); + let (events_tx, _) = broadcast::channel(16384); let kv = Arc::new(KVStore::new(events_tx.clone())); Arc::new(Self { db,