diff --git a/docs/control-api.md b/docs/control-api.md index 774d993..1093839 100644 --- a/docs/control-api.md +++ b/docs/control-api.md @@ -100,6 +100,22 @@ mode = "raw" # "raw" or "macros"; FLY_MACRO_MODE overr - Refusals are counted as `fly_chat_rejected_total{reason}`, one series per rule: `control`, `charset`, `empty`, `too_long`, `url`, `name`, `deny_list`, `rate_limited`, `malformed`. Acceptances are `fly_chat_accepted_total`, and the ring depth is `fly_chat_ring_lines`. +- **The ring survives a restart (2026-09-22).** Every accepted line rewrites a sidecar, + `/chat-ring.json` (`[paths] hot_dir`, the tmpfs the hot checkpoints use), by the same + atomic sequence a checkpoint commit uses: tmp file, fsync, rename over. At startup, before the + first publish, the file is read back; lines older than 24 hours are dropped, only the newest + `ring` of them are kept, and a missing file is silence. An unreadable, unparseable or + unknown-version file is ignored with a logged warning and an empty panel — which is what a + restart gave before this existed — never a startup failure. Writing it is best-effort too: a + failure is a warning, and the line is still accepted and still on screen. + + The sidecar is **not** part of the checkpoint: it is session state, it adds no chunk to the + `FLYSIM01` envelope and nothing about it enters the compatibility string, so `--print-compatibility` + is unchanged and a build that refuses every checkpoint in a directory still restores the panel. + It lives beside the hot checkpoints because it has their lifetime — a reboot clears the tmpfs — + and `FLY_RESET_STATE=1` clears it along with them (`infra/05-deploy.sh`). The bridge resends + nothing on reconnect: the lines the page shows after a restart are the ones the service already + accepted, with their original event ids and timestamps. `POST /chat` status codes: `202 { eventId }` accepted, `400` malformed body, `403` chat disabled, `422 { error }` a rule refused the line (the error names the rule), `429 { retryAfterMs }` a rate diff --git a/infra/05-deploy.sh b/infra/05-deploy.sh index 892bf3e..3b50552 100755 --- a/infra/05-deploy.sh +++ b/infra/05-deploy.sh @@ -258,9 +258,10 @@ if [ -n "$RELEASE_TARBALL" ]; then # BEFORE the symlink moves. Cost: one dataset load, a second or two. # # FLY_RESET_STATE=1 is the deliberate override: it archives the durable - # checkpoints (kept, never deleted) and clears the tmpfs hot ring, so the - # new build warms up fresh. Everything learned so far is thrown away, which - # is why it is not the default. + # checkpoints (kept, never deleted) and clears the tmpfs hot ring — the hot + # checkpoints and the on-screen chat ring's sidecar — so the new build warms + # up fresh. Everything learned so far is thrown away, which is why it is not + # the default. # ----------------------------------------------------------------------- state_dir="${FLY_STATE_DIR:-/srv/fly/state}" hot_dir="${FLY_STATE_HOT_DIR:-/run/fly/state}" @@ -294,7 +295,7 @@ if [ -n "$RELEASE_TARBALL" ]; then # state_dir is its own mountpoint, so the directory itself cannot be # renamed; its contents move instead. ct_exec "$CTID" -- sh -c "mkdir -p '$archive' && mv '${state_dir}'/*.checkpoint '${state_dir}/manifest.json' '$archive'/ 2>/dev/null; chown -R fly:fly '$archive'" - ct_exec "$CTID" -- sh -c "rm -f '${hot_dir}'/*.checkpoint '${hot_dir}/manifest.json' 2>/dev/null; true" + ct_exec "$CTID" -- sh -c "rm -f '${hot_dir}'/*.checkpoint '${hot_dir}/manifest.json' '${hot_dir}/chat-ring.json' 2>/dev/null; true" else die "05-deploy: REFUSING to deploy release ${version}: its checkpoint compatibility string does not match the live state in ${state_dir}, so flysim would refuse every checkpoint there and then refuse to start at all — a black stream. live state: ${live_compat} diff --git a/services/flysim/crates/flysim/src/chat.rs b/services/flysim/crates/flysim/src/chat.rs index 68625df..7bce9aa 100644 --- a/services/flysim/crates/flysim/src/chat.rs +++ b/services/flysim/crates/flysim/src/chat.rs @@ -30,6 +30,7 @@ use std::collections::{HashMap, VecDeque}; use std::path::{Path, PathBuf}; use std::time::{Duration, Instant}; +use serde::{Deserialize, Serialize}; use unicode_normalization::UnicodeNormalization; use crate::ratelimit::RateLimiter; @@ -426,6 +427,41 @@ impl ChatLimiter { // -- the ring --------------------------------------------------------------------------------- +/// The ring's sidecar file, inside `[paths] hot_dir` (`docs/control-api.md`, `[chat]`). +/// +/// The ring is session state, not simulation state, so it deliberately does **not** travel in the +/// `FLYSIM01` checkpoint envelope and is not in the compatibility string: a build that refuses +/// every checkpoint in a directory still reads this file, and a checkpoint written by any build +/// is byte-for-byte what it always was. It sits beside the hot checkpoints because it has their +/// lifetime — the tmpfs a reboot clears — and because the deliberate reset already clears that +/// directory (`infra/05-deploy.sh`, `FLY_RESET_STATE=1`). +/// +/// Sharing that directory with the hot checkpoints means sharing its mtime, which the watchdog +/// reads as flysim's liveness (`infra/bin/fly-watchdog`, check 1: hot-state mtime younger than +/// 30 s). That is safe here only because this file is written from the sim thread, on the same +/// command path as the line itself: a wedged loop accepts no chat, so it can never refresh the +/// directory behind the watchdog's back. Nothing else may ever write here from another thread. +pub const SIDECAR_FILE: &str = "chat-ring.json"; + +/// Lines older than this are dropped when the sidecar is read: a panel coming back after a long +/// outage should be empty rather than show a day-old conversation as if it were live. +pub const SIDECAR_MAX_AGE_MS: u64 = 24 * 60 * 60 * 1_000; + +/// The sidecar's own format version. Nothing else versions with it, which is the point. +const SIDECAR_VERSION: u32 = 1; + +/// `/chat-ring.json`. +pub fn sidecar_path(hot_dir: &Path) -> PathBuf { + hot_dir.join(SIDECAR_FILE) +} + +/// What the sidecar holds: a version and the ring, oldest first. +#[derive(Debug, Serialize, Deserialize)] +struct Sidecar { + version: u32, + lines: Vec, +} + /// The last `capacity` accepted lines, oldest first, as every snapshot header carries them. #[derive(Debug, Clone)] pub struct ChatRing { @@ -461,6 +497,78 @@ impl ChatRing { pub fn capacity(&self) -> usize { self.capacity } + + /// Write the ring to `/chat-ring.json`: tmp file, fsync, rename over, directory + /// fsync — the same atomic sequence a checkpoint commit uses, so a reader never sees a + /// half-written ring and a crash mid-write leaves the previous one. + /// + /// Called on every accepted line, which the admission limits cap at five a second, onto + /// tmpfs. + pub fn save_sidecar(&self, hot_dir: &Path) -> anyhow::Result<()> { + let sidecar = Sidecar { version: SIDECAR_VERSION, lines: self.lines() }; + let bytes = serde_json::to_vec(&sidecar)?; + crate::store::write_atomic(&sidecar_path(hot_dir), &bytes) + } + + /// Read `/chat-ring.json` into the ring, and answer how many lines it restored. + /// + /// Absent is silence and zero lines — the first run on a fresh box. Unreadable, unparseable + /// or a version this build does not know is zero lines and a logged warning: an empty panel + /// is exactly what a restart gives today, so nothing on this path may ever be fatal. Lines + /// older than [`SIDECAR_MAX_AGE_MS`] are dropped, and only the newest `capacity` survive, + /// whatever the file holds. + pub fn load_sidecar(&mut self, hot_dir: &Path, now_ms: u64) -> usize { + let path = sidecar_path(hot_dir); + let text = match std::fs::read_to_string(&path) { + Ok(text) => text, + Err(error) => { + if error.kind() != std::io::ErrorKind::NotFound { + tracing::warn!( + %error, + path = %path.display(), + "could not read the chat ring sidecar; the panel starts empty" + ); + } + return 0; + } + }; + let sidecar: Sidecar = match serde_json::from_str(&text) { + Ok(sidecar) => sidecar, + Err(error) => { + tracing::warn!( + %error, + path = %path.display(), + "the chat ring sidecar is not readable; ignoring it" + ); + return 0; + } + }; + if sidecar.version != SIDECAR_VERSION { + tracing::warn!( + version = sidecar.version, + path = %path.display(), + "the chat ring sidecar is a version this build does not read; ignoring it" + ); + return 0; + } + + let before = sidecar.lines.len(); + self.lines.clear(); + for line in sidecar.lines { + if now_ms.saturating_sub(line.wall_ms) >= SIDECAR_MAX_AGE_MS { + continue; + } + self.push(line); + } + let restored = self.lines.len(); + if restored < before { + tracing::info!( + dropped = before - restored, + "dropped chat lines older than a day from the sidecar" + ); + } + restored + } } #[cfg(test)] @@ -611,6 +719,106 @@ mod tests { assert_eq!(ChatRing::new(999).capacity(), RING_MAX); } + /// A line `wall_ms` milliseconds into the wall clock, for the sidecar tests. + fn line(id: u64, wall_ms: u64) -> ChatLine { + ChatLine { + id, + wall_ms, + by: format!("viewer_{id}"), + text: format!("line {id}"), + bot: None, + } + } + + #[test] + fn the_ring_round_trips_through_its_sidecar() { + let dir = tempfile::tempdir().unwrap(); + let now_ms = 1_757_000_000_000; + + // Nothing written yet: a fresh box is an empty ring and no complaint. + let mut cold = ChatRing::new(12); + assert_eq!(cold.load_sidecar(dir.path(), now_ms), 0); + assert!(cold.is_empty()); + + let mut ring = ChatRing::new(12); + ring.push(line(1, now_ms - 3_000)); + ring.push(line(2, now_ms - 2_000)); + ring.push(ChatLine { bot: Some(true), ..line(3, now_ms - 1_000) }); + ring.save_sidecar(dir.path()).unwrap(); + + // Beside the hot checkpoints, under the documented name, and nothing else is written. + assert!(sidecar_path(dir.path()).is_file()); + let written: Vec = std::fs::read_dir(dir.path()) + .unwrap() + .map(|entry| entry.unwrap().file_name().to_string_lossy().to_string()) + .collect(); + assert_eq!(written, [SIDECAR_FILE]); + + let mut restored = ChatRing::new(12); + assert_eq!(restored.load_sidecar(dir.path(), now_ms), 3); + assert_eq!(restored.lines(), ring.lines(), "oldest first, bot flag and all"); + + // A smaller ring than the file keeps the newest lines, not the first three it reads. + let mut small = ChatRing::new(2); + assert_eq!(small.load_sidecar(dir.path(), now_ms), 2); + assert_eq!( + small.lines().iter().map(|line| line.id).collect::>(), + [2, 3] + ); + } + + #[test] + fn a_sidecar_that_will_not_parse_is_ignored_rather_than_fatal() { + let dir = tempfile::tempdir().unwrap(); + let now_ms = 1_757_000_000_000; + + for content in [ + "", + "{ not json at all", + r#"{"version":1,"lines":[{"id":"not a number"}]}"#, + r#"{"version":1}"#, + // A format from some future build: readable JSON, unreadable meaning. + r#"{"version":99,"lines":[{"id":1,"wallMs":1757000000000,"by":"a","text":"b"}]}"#, + ] { + std::fs::write(sidecar_path(dir.path()), content).unwrap(); + let mut ring = ChatRing::new(12); + assert_eq!(ring.load_sidecar(dir.path(), now_ms), 0, "{content}"); + assert!(ring.is_empty(), "{content}"); + } + + // And the next accepted line simply writes a good one over it. + let mut ring = ChatRing::new(12); + ring.push(line(7, now_ms)); + ring.save_sidecar(dir.path()).unwrap(); + let mut back = ChatRing::new(12); + assert_eq!(back.load_sidecar(dir.path(), now_ms), 1); + } + + #[test] + fn sidecar_lines_older_than_a_day_are_dropped_on_load() { + let dir = tempfile::tempdir().unwrap(); + let now_ms = 1_757_000_000_000; + + let mut ring = ChatRing::new(12); + ring.push(line(1, now_ms - SIDECAR_MAX_AGE_MS - 1)); + ring.push(line(2, now_ms - SIDECAR_MAX_AGE_MS)); + ring.push(line(3, now_ms - SIDECAR_MAX_AGE_MS + 1)); + ring.push(line(4, now_ms - 1_000)); + ring.save_sidecar(dir.path()).unwrap(); + + let mut restored = ChatRing::new(12); + assert_eq!(restored.load_sidecar(dir.path(), now_ms), 2, "24 h exactly is too old"); + assert_eq!( + restored.lines().iter().map(|line| line.id).collect::>(), + [3, 4] + ); + + // A day later still, the whole file is stale and the panel starts empty. + let mut later = ChatRing::new(12); + assert_eq!(later.load_sidecar(dir.path(), now_ms + SIDECAR_MAX_AGE_MS), 0); + assert!(later.is_empty()); + } + #[test] fn every_reason_has_a_stable_spelling_and_a_unique_index() { let mut seen = std::collections::HashSet::new(); diff --git a/services/flysim/crates/flysim/src/simloop.rs b/services/flysim/crates/flysim/src/simloop.rs index ecf3440..2ab4bc5 100644 --- a/services/flysim/crates/flysim/src/simloop.rs +++ b/services/flysim/crates/flysim/src/simloop.rs @@ -619,6 +619,16 @@ impl Sim { sim.next_generation = sim.durable.highest_generation().max(sim.hot.highest_generation()) + 1; sim.start_writer(); sim.restore_or_warm_up()?; + // The on-screen chat ring, from its sidecar beside the hot checkpoints, before the first + // publish (`docs/control-api.md`, `[chat]`). It is session state and not part of the + // checkpoint envelope, so it is restored whatever the checkpoints did — including on a + // fresh start, where the brain is new but the panel's last dozen lines are not stale. + if config.chat.enabled { + let restored = sim.chat_ring.load_sidecar(&config.paths.hot_dir, now_wall_ms()); + if restored > 0 { + tracing::info!(lines = restored, "restored the on-screen chat ring"); + } + } // A dealt mode, seeded from the network as it now stands. No macro has a random // component since `GO FRONTIER` replaced `WANDER` (`docs/design/macros.md` section 9), so // the seed changes nothing about a run today; taking it here rather than before the @@ -1328,6 +1338,12 @@ impl Sim { text, bot: if bot { Some(true) } else { None }, }); + // The ring survives a restart because it is written here, not because it is in a + // checkpoint: one atomic rename onto tmpfs per accepted line, and a failure is a warning + // rather than a refusal — the line is already on screen. + if let Err(error) = self.chat_ring.save_sidecar(&self.shared.config.paths.hot_dir) { + tracing::warn!(%error, "could not persist the chat ring; it will not survive a restart"); + } Metrics::incr(&self.shared.metrics.chat_accepted_total); Ok(event.id) } diff --git a/services/flysim/crates/flysim/tests/integration.rs b/services/flysim/crates/flysim/tests/integration.rs index a227109..0c9f408 100644 --- a/services/flysim/crates/flysim/tests/integration.rs +++ b/services/flysim/crates/flysim/tests/integration.rs @@ -633,6 +633,38 @@ async fn the_service_streams_takes_sugar_checkpoints_and_resumes_after_being_kil "only the first instance may warm up; the second must restore:\n{log}" ); + // The CHAT panel came back with the service. The ring is session state in a sidecar beside + // the hot checkpoints, never a chunk in the envelope: the durable store holds no copy of it, + // and the checkpoint this restore just read carries no chat text. + let resumed_chat = resumed + .chat + .as_ref() + .expect("chat is enabled, so the header carries a ring"); + let restored_line = resumed_chat + .iter() + .find(|line| line.id == chat_event_id) + .unwrap_or_else(|| panic!("the chat ring did not survive the restart: {resumed_chat:?}")); + assert_eq!(restored_line.by, "integration_test"); + assert_eq!(restored_line.text, "go LEFT!"); + assert!( + resumed_chat.iter().any(|line| line.bot == Some(true)), + "the bot's line came back too: {resumed_chat:?}" + ); + assert!( + dir.path().join("hot/chat-ring.json").is_file(), + "the sidecar lives beside the hot checkpoints" + ); + assert!( + !dir.path().join("state/chat-ring.json").exists(), + "the durable store carries no chat" + ); + let envelope = + std::fs::read(dir.path().join(format!("state/{generation}.checkpoint"))).unwrap(); + assert!( + !envelope.windows(8).any(|window| window == b"go LEFT!"), + "chat text must never enter the checkpoint envelope" + ); + // -- 6. SIGTERM writes a final checkpoint -------------------------------------------- let (_, before) = service.get("/status"); let before = before["checkpoint"]["generation"].as_u64().unwrap();