diff --git a/client-core/Cargo.lock b/client-core/Cargo.lock index 627138e..08c65d7 100644 --- a/client-core/Cargo.lock +++ b/client-core/Cargo.lock @@ -47,6 +47,7 @@ name = "client-core" version = "0.1.0" dependencies = [ "event-model", + "log", "pulldown-cmark", "serde", "serde_json", diff --git a/client-core/Cargo.toml b/client-core/Cargo.toml index ab8e9ce..79ec4b2 100644 --- a/client-core/Cargo.toml +++ b/client-core/Cargo.toml @@ -37,6 +37,11 @@ ureq = { version = "3", features = ["json"] } # the same parser at the same version, rather than a hand-written splitter # that would drift from it. pulldown-cmark = "0.13.4" +# The logging facade only -- `log_ring` implements a `log::Log` backend and +# wraps whichever real one the platform installed (`android_logger` on the +# phone, `env_logger` on the desktop), which is why neither of those is a +# dependency here. See `log_ring`'s module doc. +log = { version = "0.4.28", features = ["std"] } [dev-dependencies] diff --git a/client-core/src/lib.rs b/client-core/src/lib.rs index c7a91d8..310160e 100644 --- a/client-core/src/lib.rs +++ b/client-core/src/lib.rs @@ -8,6 +8,8 @@ pub mod config; pub mod durations; pub mod event_stream; pub mod highlight; +pub mod log_ring; +pub mod log_upload; pub mod markdown_blocks; pub mod notifications; pub mod sse; diff --git a/client-core/src/log_ring.rs b/client-core/src/log_ring.rs new file mode 100644 index 0000000..ec975ad --- /dev/null +++ b/client-core/src/log_ring.rs @@ -0,0 +1,483 @@ +//! The app's own recent log, held in memory so it can be read back +//! without `logcat`. +//! +//! **Why this exists**: Iris tests iris builds on a GrapheneOS phone with +//! no `adb`, and Android forbids one app reading another's logcat, so +//! nothing outside the process can recover what it wrote. The only way a +//! line reaches her is for the app to carry its own copy. This is that +//! copy: a bounded ring every `log::info!` in the process lands in, on top +//! of whichever platform logger was already installed (`android_logger`, +//! `env_logger`) rather than instead of it -- see [`RingLogger`]. +//! +//! Two consumers, both reading the same ring rather than each keeping +//! their own: the bench app's `Copy report`/`Diagnostics` (which reads +//! [`LogRing::to_text`] and [`LogRing::summary`]) and the uploader in +//! [`crate::log_upload`] (which reads [`LogRing::since`]). That is why +//! reading does not consume: a line the uploader has sent must still be in +//! the report, and a report taken twice must say the same thing. + +use std::collections::VecDeque; +use std::sync::{Arc, Mutex}; +use std::time::{SystemTime, UNIX_EPOCH}; + +/// How many lines a default ring holds, and how many bytes of message. +/// +/// Both bounds apply -- whichever bites first -- because the two failure +/// modes are different: a flood of short lines exhausts the count, and one +/// pathological line (a stack trace, a pretty-printed JSON body) exhausts +/// the bytes. A ring bounded only by lines can hold megabytes; one bounded +/// only by bytes can be emptied by a single line. +pub const DEFAULT_MAX_LINES: usize = 2000; +pub const DEFAULT_MAX_BYTES: usize = 256 * 1024; + +/// One recorded line. `seq` is assigned by the ring and only ever +/// increases, so a reader that remembers where it got to can ask for what +/// came after -- and a gap in the sequence is exactly the lines the bound +/// dropped. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct LogLine { + pub seq: u64, + /// Milliseconds since the unix epoch, from the app's own clock. The + /// app's rather than the receiver's: a line is timestamped when it + /// happened, and an upload can be minutes later or never. + pub at_ms: u64, + pub level: log::Level, + pub target: String, + pub message: String, +} + +impl LogLine { + /// Roughly what the line costs the ring. The two `String`s dominate; + /// the fixed fields are counted as a flat overhead so a ring of empty + /// messages still has a bound. + fn weight(&self) -> usize { + self.target.len() + self.message.len() + 32 + } + + /// `12:34:56.789 INFO iris::android: the message`, the shape a + /// person skims. Time of day only -- the date is in the report's own + /// header, and a ring never spans one. + pub fn format(&self) -> String { + format!( + "{} {:<5} {}: {}", + clock_time(self.at_ms), + self.level, + self.target, + self.message + ) + } +} + +/// `HH:MM:SS.mmm` in UTC from a unix millisecond count, without a date +/// library: the only field this needs is the time of day, and dividing out +/// the day is the whole calculation. Deliberately not local time -- the +/// phone's offset is not knowable here, and a report that says UTC is +/// comparable with the server's log, which is what it gets read against. +fn clock_time(at_ms: u64) -> String { + let ms = at_ms % 1000; + let secs_of_day = (at_ms / 1000) % 86_400; + format!( + "{:02}:{:02}:{:02}.{:03}", + secs_of_day / 3600, + (secs_of_day % 3600) / 60, + secs_of_day % 60, + ms + ) +} + +/// Now, in unix milliseconds. Saturating rather than panicking on a clock +/// before the epoch: a wrong timestamp in a diagnostic is not worth taking +/// the app down for. +pub fn now_ms() -> u64 { + SystemTime::now() + .duration_since(UNIX_EPOCH) + .map(|d| d.as_millis() as u64) + .unwrap_or(0) +} + +#[derive(Debug)] +struct Inner { + lines: VecDeque, + bytes: usize, + max_lines: usize, + max_bytes: usize, + next_seq: u64, + /// How many lines the bounds have discarded since the ring was made. + /// Reported rather than inferred, so "the log starts here" and "the + /// log was cut off here" are distinguishable -- the unknown state the + /// UI rules ask for. + dropped: u64, +} + +/// A bounded, shareable ring of recent log lines. Cloning shares the ring; +/// there is one per process and every holder sees the same lines. +#[derive(Debug, Clone)] +pub struct LogRing(Arc>); + +impl LogRing { + pub fn new(max_lines: usize, max_bytes: usize) -> Self { + assert!( + max_lines > 0 && max_bytes > 0, + "a ring with no room holds nothing" + ); + Self(Arc::new(Mutex::new(Inner { + lines: VecDeque::new(), + bytes: 0, + max_lines, + max_bytes, + next_seq: 0, + dropped: 0, + }))) + } + + /// The bounds this project ships with: [`DEFAULT_MAX_LINES`] and + /// [`DEFAULT_MAX_BYTES`]. + pub fn with_defaults() -> Self { + Self::new(DEFAULT_MAX_LINES, DEFAULT_MAX_BYTES) + } + + /// A poisoned lock is a bug in a panicking logger, not a reason to + /// take the app down a second time -- the ring is a diagnostic, and + /// losing it must not be worse than the fault it was recording. + fn with(&self, f: impl FnOnce(&mut Inner) -> R) -> R { + let mut guard = match self.0.lock() { + Ok(guard) => guard, + Err(poisoned) => poisoned.into_inner(), + }; + f(&mut guard) + } + + /// Records a line, evicting the oldest until both bounds hold again. + pub fn push(&self, level: log::Level, target: &str, message: String) { + self.with(|inner| { + let line = LogLine { + seq: inner.next_seq, + at_ms: now_ms(), + level, + target: target.to_string(), + message, + }; + inner.next_seq += 1; + inner.bytes += line.weight(); + inner.lines.push_back(line); + // `!is_empty()` rather than `len() > 1`: one line larger than + // the whole byte bound is kept, because dropping it would + // leave the ring silently empty while lines were arriving. + while inner.lines.len() > inner.max_lines + || (inner.bytes > inner.max_bytes && inner.lines.len() > 1) + { + if let Some(evicted) = inner.lines.pop_front() { + inner.bytes -= evicted.weight(); + inner.dropped += 1; + } + } + }) + } + + /// Every line held, oldest first. + pub fn snapshot(&self) -> Vec { + self.with(|inner| inner.lines.iter().cloned().collect()) + } + + /// The lines with a sequence number at or after `seq`, oldest first, + /// and the sequence to ask from next time. Does not consume: see this + /// module's doc for why. + pub fn since(&self, seq: u64) -> (Vec, u64) { + self.with(|inner| { + let lines: Vec = inner + .lines + .iter() + .filter(|line| line.seq >= seq) + .cloned() + .collect(); + let next = lines.last().map(|line| line.seq + 1).unwrap_or(seq); + (lines, next) + }) + } + + pub fn len(&self) -> usize { + self.with(|inner| inner.lines.len()) + } + + pub fn is_empty(&self) -> bool { + self.len() == 0 + } + + pub fn dropped(&self) -> u64 { + self.with(|inner| inner.dropped) + } + + /// When the newest line was written, in unix milliseconds, or `None` + /// for a ring nothing has been written to. + pub fn last_at_ms(&self) -> Option { + self.with(|inner| inner.lines.back().map(|line| line.at_ms)) + } + + /// Every line held, formatted one per line -- what `Copy report` + /// appends. + pub fn to_text(&self) -> String { + self.snapshot() + .iter() + .map(LogLine::format) + .collect::>() + .join("\n") + } + + /// One line for a diagnostics pane: how much is held, how much was + /// dropped, and when the last line arrived. "no lines yet" is its own + /// wording rather than a count of zero with a made-up time, because + /// "nothing has been logged" and "logging is not running" would + /// otherwise look the same. + pub fn summary(&self) -> String { + let (len, dropped, last) = self.with(|inner| { + ( + inner.lines.len(), + inner.dropped, + inner.lines.back().map(|line| line.at_ms), + ) + }); + match last { + None => "app log: no lines yet".to_string(), + Some(at) => { + let dropped = if dropped > 0 { + format!(", {dropped} dropped") + } else { + String::new() + }; + format!( + "app log: {len} lines held{dropped}, last {}", + clock_time(at) + ) + } + } + } +} + +/// A `log` backend that records into a [`LogRing`] **and** forwards to the +/// logger the platform already installs, so nothing that reads the +/// platform's log (`logcat`, a terminal) changes. +/// +/// The inner logger is passed in rather than chosen here: `client-core` +/// has no business depending on `android_logger` or `env_logger`, and +/// which one is right is exactly what differs between the two platforms +/// (the sharing rule in AGENTS.md). +pub struct RingLogger { + ring: LogRing, + inner: Box, +} + +impl RingLogger { + pub fn new(ring: LogRing, inner: Box) -> Self { + Self { ring, inner } + } +} + +impl log::Log for RingLogger { + /// True for anything `log`'s own max level lets through: the ring + /// wants everything, even where the platform logger would filter it + /// out. The filter is applied per-logger in [`Self::log`] instead. + fn enabled(&self, _metadata: &log::Metadata) -> bool { + true + } + + fn log(&self, record: &log::Record) { + self.ring + .push(record.level(), record.target(), record.args().to_string()); + if self.inner.enabled(record.metadata()) { + self.inner.log(record); + } + } + + fn flush(&self) { + self.inner.flush(); + } +} + +/// Installs a [`RingLogger`] as the process logger and answers the ring it +/// records into. +/// +/// Fails only if a logger is already installed, which is a programmer +/// error (two initialisation paths) rather than a recoverable condition -- +/// the caller is named in the error so it is findable. +pub fn install( + ring: LogRing, + inner: Box, + max_level: log::LevelFilter, +) -> Result<(), log::SetLoggerError> { + log::set_boxed_logger(Box::new(RingLogger::new(ring, inner)))?; + log::set_max_level(max_level); + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use log::Level; + + fn fill(ring: &LogRing, count: usize) { + for n in 0..count { + ring.push(Level::Info, "test", format!("line {n}")); + } + } + + #[test] + fn lines_come_back_oldest_first() { + let ring = LogRing::new(10, 1 << 20); + fill(&ring, 3); + let text: Vec = ring.snapshot().into_iter().map(|l| l.message).collect(); + assert_eq!(text, ["line 0", "line 1", "line 2"]); + } + + #[test] + fn the_line_bound_drops_the_oldest_and_says_how_many() { + let ring = LogRing::new(3, 1 << 20); + fill(&ring, 5); + let text: Vec = ring.snapshot().into_iter().map(|l| l.message).collect(); + assert_eq!(text, ["line 2", "line 3", "line 4"], "the newest survive"); + assert_eq!(ring.len(), 3); + assert_eq!(ring.dropped(), 2, "and the loss is reported, not silent"); + } + + #[test] + fn the_byte_bound_bites_before_the_line_bound_when_lines_are_large() { + // Room for 1000 lines but only a few hundred bytes. + let ring = LogRing::new(1000, 300); + for n in 0..10 { + ring.push(Level::Info, "t", format!("{n}{}", "x".repeat(100))); + } + assert!( + ring.len() < 10, + "the byte bound evicted: {} held", + ring.len() + ); + assert!(ring.dropped() > 0); + assert!( + ring.snapshot().last().unwrap().message.starts_with('9'), + "and it evicted from the old end" + ); + } + + /// The case the `len() > 1` guard exists for: one line larger than the + /// whole bound must still be readable, or a ring that is over budget + /// reads as a ring nothing was written to. + #[test] + fn one_oversized_line_is_kept_rather_than_leaving_the_ring_empty() { + let ring = LogRing::new(100, 64); + ring.push(Level::Error, "t", "y".repeat(5000)); + assert_eq!(ring.len(), 1); + assert_eq!(ring.dropped(), 0); + } + + #[test] + fn sequence_numbers_only_increase_and_survive_eviction() { + let ring = LogRing::new(2, 1 << 20); + fill(&ring, 5); + let seqs: Vec = ring.snapshot().into_iter().map(|l| l.seq).collect(); + assert_eq!(seqs, [3, 4], "a gap is exactly what was dropped"); + } + + #[test] + fn since_returns_only_what_is_new_and_the_next_cursor() { + let ring = LogRing::new(100, 1 << 20); + fill(&ring, 3); + let (first, cursor) = ring.since(0); + assert_eq!(first.len(), 3); + assert_eq!(cursor, 3); + + let (none, cursor) = ring.since(cursor); + assert!(none.is_empty(), "nothing new yet"); + assert_eq!(cursor, 3, "and the cursor does not move"); + + ring.push(Level::Warn, "test", "later".into()); + let (more, cursor) = ring.since(cursor); + assert_eq!(more.len(), 1); + assert_eq!(more[0].message, "later"); + assert_eq!(cursor, 4); + } + + #[test] + fn reading_does_not_consume() { + let ring = LogRing::new(100, 1 << 20); + fill(&ring, 2); + let (sent, _) = ring.since(0); + assert_eq!(sent.len(), 2); + assert_eq!(ring.len(), 2, "the report still has them after an upload"); + assert_eq!(ring.to_text().lines().count(), 2); + } + + #[test] + fn an_empty_ring_says_so_rather_than_reporting_a_time() { + let ring = LogRing::with_defaults(); + assert_eq!(ring.summary(), "app log: no lines yet"); + assert_eq!(ring.last_at_ms(), None); + assert!(ring.is_empty()); + } + + #[test] + fn the_summary_names_dropped_lines_only_when_there_are_some() { + let ring = LogRing::new(2, 1 << 20); + fill(&ring, 2); + assert!(!ring.summary().contains("dropped"), "{}", ring.summary()); + fill(&ring, 2); + assert!(ring.summary().contains("2 dropped"), "{}", ring.summary()); + } + + #[test] + fn a_line_formats_as_time_level_target_message() { + let line = LogLine { + seq: 0, + // 1970-01-01T12:34:56.789Z, so the arithmetic is checkable by + // hand rather than against another clock. + at_ms: (12 * 3600 + 34 * 60 + 56) * 1000 + 789, + level: Level::Info, + target: "iris::android".into(), + message: "surface created".into(), + } + .format(); + assert_eq!(line, "12:34:56.789 INFO iris::android: surface created"); + } + + /// The forwarding half: a line reaches the ring *and* the logger the + /// platform already had, and one the inner logger filters out is still + /// in the ring. + #[test] + fn the_ring_logger_forwards_to_the_inner_logger() { + use log::Log; + struct Collect(Arc>>, log::Level); + impl Log for Collect { + fn enabled(&self, metadata: &log::Metadata) -> bool { + metadata.level() <= self.1 + } + fn log(&self, record: &log::Record) { + self.0.lock().unwrap().push(record.args().to_string()); + } + fn flush(&self) {} + } + + let seen = Arc::new(Mutex::new(Vec::new())); + let ring = LogRing::with_defaults(); + let logger = RingLogger::new(ring.clone(), Box::new(Collect(seen.clone(), Level::Info))); + logger.log( + &log::Record::builder() + .args(format_args!("kept")) + .level(Level::Info) + .target("t") + .build(), + ); + logger.log( + &log::Record::builder() + .args(format_args!("filtered")) + .level(Level::Debug) + .target("t") + .build(), + ); + + assert_eq!( + *seen.lock().unwrap(), + ["kept"], + "the inner logger's own filter still applies" + ); + let held: Vec = ring.snapshot().into_iter().map(|l| l.message).collect(); + assert_eq!(held, ["kept", "filtered"], "the ring keeps both"); + } +} diff --git a/client-core/src/log_upload.rs b/client-core/src/log_upload.rs new file mode 100644 index 0000000..408b432 --- /dev/null +++ b/client-core/src/log_upload.rs @@ -0,0 +1,427 @@ +//! Sending [`crate::log_ring`]'s lines to `ai-server`, so a phone with no +//! `logcat` still has a way for a `log::info!` to reach a person. +//! +//! **Where they end up**: `POST /client-log` re-emits each line into +//! `ai-server`'s own `tracing` output, which Dev Updater already shows as +//! that component's *runtime log* (it runs `ai-server` as a `Managed` +//! service, and a managed service's stdout is redirected to a file its +//! service script reports). So this needs no new route, storage or viewer +//! in Dev Updater at all -- see `docs/DECISIONS.md`, 2026-09-07. +//! +//! **Nothing here calls `log!`.** Every line this module logged would land +//! in the ring it is draining and be uploaded, so a server that is down +//! would produce a growing conversation with itself. Failures are recorded +//! in [`UploadStatus`] instead and shown in the app's diagnostics pane, +//! which is where somebody looking for "why is nothing arriving" is +//! already looking (UI_RULES.md: a failure is reported where it happened). + +use std::sync::{Arc, Condvar, Mutex}; +use std::time::Duration; + +use crate::api::{ApiError, Body, Transport}; +use crate::log_ring::LogRing; + +/// The most lines one request carries. A phone that has been offline for +/// an hour has thousands waiting, and one request holding all of them is a +/// body the server has to buffer whole; the rest go in the next batch, +/// which the loop takes immediately rather than after the next interval. +pub const MAX_LINES_PER_BATCH: usize = 500; + +/// How much of one message is sent. Long enough for a stack trace line, +/// short enough that one pathological message cannot dominate a batch. +/// Truncation is marked, because a silently shortened line reads as a line +/// that ended there. +pub const MAX_MESSAGE_BYTES: usize = 4096; + +/// The route this posts to, on `server/src/routes.rs`'s surface. +pub const CLIENT_LOG_PATH: &str = "/client-log"; + +/// What the last upload attempt did, for a diagnostics pane. `None` for +/// "nothing has been tried yet", which is deliberately distinct from a +/// success that sent nothing. +#[derive(Debug, Clone, Default)] +pub struct UploadStatus { + pub sent: u64, + pub last_error: Option, + pub attempted: bool, +} + +impl UploadStatus { + /// One line for the diagnostics pane, in the same voice as + /// [`LogRing::summary`]. + pub fn summary(&self) -> String { + match (&self.last_error, self.attempted) { + (Some(err), _) => format!("log upload: failing -- {err} ({} sent so far)", self.sent), + (None, false) => "log upload: not tried yet".to_string(), + (None, true) => format!("log upload: {} lines sent", self.sent), + } + } +} + +/// Drains a [`LogRing`] into `POST /client-log`, remembering how far it +/// got so a line is sent once and stays in the ring for the report. +pub struct LogUploader { + ring: LogRing, + transport: Arc, + source: String, + cursor: u64, + status: Arc>, +} + +impl LogUploader { + /// `source` names the build these lines came from -- it is what + /// distinguishes them in `ai-server`'s log from the server's own + /// lines and from another device's. + pub fn new(ring: LogRing, transport: Arc, source: impl Into) -> Self { + Self { + ring, + transport, + source: source.into(), + cursor: 0, + status: Arc::new(Mutex::new(UploadStatus::default())), + } + } + + /// A handle on what the last attempt did, shareable with the UI. + pub fn status(&self) -> Arc> { + Arc::clone(&self.status) + } + + /// Sends up to [`MAX_LINES_PER_BATCH`] waiting lines. Answers how many + /// went, and whether more are waiting -- the loop uses the second to + /// decide whether to go round again at once. + pub fn flush_once(&mut self) -> Result<(usize, bool), ApiError> { + let (mut lines, mut next) = self.ring.since(self.cursor); + let more = lines.len() > MAX_LINES_PER_BATCH; + if more { + lines.truncate(MAX_LINES_PER_BATCH); + next = lines.last().map(|line| line.seq + 1).unwrap_or(next); + } + if lines.is_empty() { + return Ok((0, false)); + } + + let body = serde_json::json!({ + "source": self.source, + "lines": lines.iter().map(|line| serde_json::json!({ + "seq": line.seq, + "at": line.at_ms, + "level": line.level.as_str(), + "target": line.target, + "message": truncate(&line.message), + })).collect::>(), + }); + + let result = self + .transport + .request("POST", CLIENT_LOG_PATH, Some(Body::Json(body))); + let mut status = self.status.lock().unwrap_or_else(|e| e.into_inner()); + status.attempted = true; + match result { + Ok(response) if (200..300).contains(&response.status) => { + // Only on success: a failed batch is retried from the same + // cursor next time, which is what makes a dropped tunnel + // cost nothing but a delay. + self.cursor = next; + status.sent += lines.len() as u64; + status.last_error = None; + Ok((lines.len(), more)) + } + Ok(response) => { + let message = format!("{} from {CLIENT_LOG_PATH}", response.status); + status.last_error = Some(message.clone()); + Err(ApiError { + message, + status: Some(response.status), + }) + } + Err(err) => { + status.last_error = Some(err.message.clone()); + Err(err) + } + } + } +} + +/// Cuts a message to [`MAX_MESSAGE_BYTES`] on a character boundary, saying +/// so, rather than letting one line dominate a batch. +fn truncate(message: &str) -> String { + if message.len() <= MAX_MESSAGE_BYTES { + return message.to_string(); + } + let mut end = MAX_MESSAGE_BYTES; + while end > 0 && !message.is_char_boundary(end) { + end -= 1; + } + format!("{}… [{} bytes cut]", &message[..end], message.len() - end) +} + +/// The background half: a thread that flushes on a timer and on demand. +/// +/// Its path out is [`Drop`] -- dropping the handle stops the thread and +/// waits for it, so an app that tears the uploader down does not leave one +/// posting behind it. +pub struct LogUpload { + signal: Arc<(Mutex, Condvar)>, + thread: Option>, + status: Arc>, +} + +/// Only "stop": a nudge from [`LogUpload::flush_now`] needs no flag, +/// because the thread's reaction to waking is to flush, and flushing an +/// empty ring costs nothing -- so a spurious wakeup is already correct. +#[derive(Default)] +struct Signal { + stop: bool, +} + +impl LogUpload { + /// Starts the loop. `every` is how long it waits between flushes when + /// nobody nudges it -- a compromise between a line arriving promptly + /// and a radio the app woke for one line. + pub fn spawn( + ring: LogRing, + transport: Arc, + source: impl Into, + every: Duration, + ) -> Self { + let mut uploader = LogUploader::new(ring, transport, source); + let status = uploader.status(); + let signal = Arc::new((Mutex::new(Signal::default()), Condvar::new())); + let thread = { + let signal = Arc::clone(&signal); + std::thread::Builder::new() + .name("client-log-upload".into()) + .spawn(move || { + loop { + // Keep going while a batch was capped, so a + // backlog drains at once rather than one batch per + // interval. + while let Ok((_, more)) = uploader.flush_once() { + if !more { + break; + } + } + let (lock, condvar) = &*signal; + let state = lock.lock().unwrap_or_else(|e| e.into_inner()); + if state.stop { + return; + } + let (state, _) = condvar + .wait_timeout(state, every) + .unwrap_or_else(|e| e.into_inner()); + if state.stop { + return; + } + } + }) + .expect("spawning the log upload thread") + }; + Self { + signal, + thread: Some(thread), + status, + } + } + + /// Sends what is waiting now -- what `Copy report` calls, so the lines + /// a person is about to describe are already on the server. + pub fn flush_now(&self) { + let (lock, condvar) = &*self.signal; + let _state = lock.lock().unwrap_or_else(|e| e.into_inner()); + condvar.notify_all(); + } + + pub fn status(&self) -> UploadStatus { + self.status + .lock() + .unwrap_or_else(|e| e.into_inner()) + .clone() + } +} + +impl Drop for LogUpload { + fn drop(&mut self) { + { + let (lock, condvar) = &*self.signal; + let mut state = lock.lock().unwrap_or_else(|e| e.into_inner()); + state.stop = true; + condvar.notify_all(); + } + if let Some(thread) = self.thread.take() { + let _ = thread.join(); + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::api::RawResponse; + use log::Level; + + /// Records every body posted, and answers whatever status the test set. + struct Fake { + posted: Mutex>, + status: Mutex, + } + + impl Fake { + fn new() -> Arc { + Arc::new(Self { + posted: Mutex::new(Vec::new()), + status: Mutex::new(200), + }) + } + fn bodies(&self) -> Vec { + self.posted.lock().unwrap().clone() + } + } + + impl Transport for Fake { + fn request( + &self, + method: &str, + path: &str, + body: Option, + ) -> Result { + assert_eq!(method, "POST"); + assert_eq!(path, CLIENT_LOG_PATH); + if let Some(Body::Json(value)) = body { + self.posted.lock().unwrap().push(value); + } else { + panic!("the client log is posted as JSON"); + } + Ok(RawResponse { + status: *self.status.lock().unwrap(), + body: Vec::new(), + }) + } + fn stream(&self, _path: &str) -> Result, ApiError> { + unreachable!("the log uploader never streams") + } + } + + fn ring_with(count: usize) -> LogRing { + let ring = LogRing::with_defaults(); + for n in 0..count { + ring.push(Level::Info, "t", format!("line {n}")); + } + ring + } + + #[test] + fn an_empty_ring_posts_nothing() { + let fake = Fake::new(); + let mut uploader = LogUploader::new(LogRing::with_defaults(), fake.clone(), "test"); + assert_eq!(uploader.flush_once().unwrap(), (0, false)); + assert!( + fake.bodies().is_empty(), + "no request at all, not an empty one" + ); + } + + #[test] + fn a_line_is_sent_once() { + let fake = Fake::new(); + let ring = ring_with(3); + let mut uploader = LogUploader::new(ring.clone(), fake.clone(), "test"); + assert_eq!(uploader.flush_once().unwrap().0, 3); + assert_eq!( + uploader.flush_once().unwrap(), + (0, false), + "nothing repeats" + ); + + ring.push(Level::Warn, "t", "later".into()); + assert_eq!(uploader.flush_once().unwrap().0, 1); + assert_eq!(fake.bodies().len(), 2); + assert_eq!(ring.len(), 4, "and the report still holds all of them"); + } + + /// The half the change had no reason to touch: a server that refuses + /// must not lose the lines. + #[test] + fn a_failed_batch_is_retried_from_the_same_place() { + let fake = Fake::new(); + *fake.status.lock().unwrap() = 503; + let mut uploader = LogUploader::new(ring_with(2), fake.clone(), "test"); + assert!(uploader.flush_once().is_err()); + assert!( + uploader.status().lock().unwrap().last_error.is_some(), + "and it says why, where somebody can see it" + ); + + *fake.status.lock().unwrap() = 200; + assert_eq!(uploader.flush_once().unwrap().0, 2, "the same two lines"); + assert!(uploader.status().lock().unwrap().last_error.is_none()); + } + + #[test] + fn a_backlog_is_capped_per_batch_and_says_there_is_more() { + let fake = Fake::new(); + let ring = LogRing::new(MAX_LINES_PER_BATCH * 3, 1 << 30); + for n in 0..(MAX_LINES_PER_BATCH + 7) { + ring.push(Level::Info, "t", format!("{n}")); + } + let mut uploader = LogUploader::new(ring, fake.clone(), "test"); + assert_eq!(uploader.flush_once().unwrap(), (MAX_LINES_PER_BATCH, true)); + assert_eq!(uploader.flush_once().unwrap(), (7, false)); + } + + #[test] + fn the_body_carries_the_source_and_each_line_whole() { + let fake = Fake::new(); + let ring = LogRing::with_defaults(); + ring.push(Level::Error, "iris::android", "surface lost".into()); + LogUploader::new(ring, fake.clone(), "iris-bench 1.2") + .flush_once() + .unwrap(); + let body = &fake.bodies()[0]; + assert_eq!(body["source"], "iris-bench 1.2"); + let line = &body["lines"][0]; + assert_eq!(line["level"], "ERROR"); + assert_eq!(line["target"], "iris::android"); + assert_eq!(line["message"], "surface lost"); + assert!(line["at"].as_u64().is_some(), "the app's own clock"); + } + + #[test] + fn an_enormous_message_is_cut_and_says_so() { + let cut = truncate(&"x".repeat(MAX_MESSAGE_BYTES + 100)); + assert!(cut.starts_with("xxxx")); + assert!(cut.contains("bytes cut"), "{cut}"); + assert!(cut.len() < MAX_MESSAGE_BYTES + 64); + let short = truncate("fine"); + assert_eq!(short, "fine", "a short message is untouched"); + } + + #[test] + fn the_status_line_distinguishes_untried_from_sent_nothing() { + let untried = UploadStatus::default(); + assert_eq!(untried.summary(), "log upload: not tried yet"); + let sent_none = UploadStatus { + attempted: true, + ..Default::default() + }; + assert_eq!(sent_none.summary(), "log upload: 0 lines sent"); + } + + /// The path out: dropping the handle must stop the thread, not leave + /// it posting. + #[test] + fn dropping_the_handle_stops_the_thread() { + let fake = Fake::new(); + let upload = LogUpload::spawn( + ring_with(1), + fake.clone(), + "test", + Duration::from_millis(10), + ); + upload.flush_now(); + drop(upload); + let after = fake.bodies().len(); + std::thread::sleep(Duration::from_millis(60)); + assert_eq!(fake.bodies().len(), after, "nothing posted after the drop"); + } +} diff --git a/server/src/routes.rs b/server/src/routes.rs index aac707e..1a85d68 100644 --- a/server/src/routes.rs +++ b/server/src/routes.rs @@ -62,6 +62,9 @@ //! once the account's usage limit lifts //! GET /notifications SSE: every session's attention-wanting //! moments, live only (see `notifications`) +//! POST /client-log {source, lines} -- a client's own recent log, +//! re-emitted into this server's log (see +//! `client_log`; the phone has no logcat) //! GET /defaults {effort} -- what a new session starts at //! POST /defaults {effort} -- null for the CLI's own default //! GET /usage cached usage windows per provider @@ -158,6 +161,7 @@ pub fn router(manager: Arc) -> Router { .route("/sessions/{id}/permission-mode", post(set_permission_mode)) .route("/sessions/{id}/effort", post(set_effort)) .route("/defaults", get(defaults).post(set_defaults)) + .route("/client-log", post(client_log)) .route("/sessions/{id}/notify", post(set_notify)) .route("/sessions/{id}/auto-resume", post(set_auto_resume)) .route("/notifications", get(notifications)) @@ -1407,6 +1411,111 @@ async fn defaults(State(manager): State>) -> axum::Json, +} + +/// Takes a client's own recent log lines and re-emits them into this +/// server's `tracing` output. +/// +/// **Why a route rather than something on the phone**: Android forbids one +/// app reading another's `logcat`, and the phone this project is tested on +/// has no `adb` at all, so a `log::info!` in the app can only reach a +/// person if the app carries its own copy and sends it somewhere. This +/// server is the somewhere it already has a tunnel, a pinned CA and a +/// bearer token for -- and Dev Updater already shows this server's log as +/// its runtime log, so the line arrives where its reader is already +/// looking with nothing new built there. `docs/DECISIONS.md`, 2026-09-07. +/// +/// Each line is emitted separately, at the level the client recorded it +/// at, with the client's own timestamp in the text -- the tracing +/// subscriber stamps the moment of *arrival*, which can be minutes later +/// or on the other side of a tunnel outage, and presenting that as when it +/// happened would be an inferred value shown as a measured one. +async fn client_log(axum::Json(body): axum::Json) -> Result { + if body.lines.len() > CLIENT_LOG_MAX_LINES { + return Err(ApiError::BadRequest(format!( + "{} lines in one batch; the limit is {CLIENT_LOG_MAX_LINES}", + body.lines.len() + ))); + } + for line in &body.lines { + let at = client_log_time(line.at); + let source = &body.source; + let target = &line.target; + let seq = line.seq; + let message = &line.message; + // The level is chosen here rather than passed, because a tracing + // macro's level is part of the callsite. Anything unrecognised is + // reported at INFO with the word it sent kept, so a client using a + // level this server has not heard of loses the level rather than + // the line. + match line.level.to_ascii_uppercase().as_str() { + "ERROR" => { + tracing::error!(target: "client_log", "[{source} {at} #{seq}] {target}: {message}") + } + "WARN" => { + tracing::warn!(target: "client_log", "[{source} {at} #{seq}] {target}: {message}") + } + "DEBUG" => { + tracing::debug!(target: "client_log", "[{source} {at} #{seq}] {target}: {message}") + } + "TRACE" => { + tracing::trace!(target: "client_log", "[{source} {at} #{seq}] {target}: {message}") + } + "INFO" => { + tracing::info!(target: "client_log", "[{source} {at} #{seq}] {target}: {message}") + } + other => { + tracing::info!(target: "client_log", "[{source} {at} #{seq}] {target}: <{other}> {message}") + } + } + } + Ok(StatusCode::NO_CONTENT) +} + +/// `HH:MM:SS.mmm` UTC from the client's unix milliseconds -- the same +/// formatting `client_core::log_ring` uses, so a line read here and the +/// same line in the app's own copied report say the same time. +fn client_log_time(at_ms: u64) -> String { + let secs_of_day = (at_ms / 1000) % 86_400; + format!( + "{:02}:{:02}:{:02}.{:03}", + secs_of_day / 3600, + (secs_of_day % 3600) / 60, + secs_of_day % 60, + at_ms % 1000 + ) +} + /// Sets what a new session's thinking level is. Applied when a session is /// spawned, so nothing already running changes underneath anybody. async fn set_defaults(