client-core: the app's own log ring, and POST /client-log to get it off a phone
Iris tests iris builds on a phone with no adb, and Android forbids one app reading another's logcat, so a `log::info!` in the app can only reach her if the app carries its own copy and sends it somewhere. `client_core::log_ring` is that copy: a bounded ring (2000 lines / 256 KiB, whichever bites first) behind a `log::Log` backend that forwards to whichever real logger the platform installed, so `logcat` and the desktop terminal see exactly what they saw before. Reading does not consume -- the report and the uploader are two readers of one ring. `client_core::log_upload` drains it into ai-server's new `POST /client-log`, which re-emits each line into the server's own tracing output. Dev Updater already shows that as ai-server's runtime log, so nothing new is built there. A failed batch is retried from the same cursor, and nothing in the upload path calls `log!` -- it would land in the ring it is draining. Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
1 parent
9cd1263080
commit
977bdb9ee0
6 files changed
+1027
No files matched your search
Generated
+1
@@ -47,6 +47,7 @@ name = "client-core"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"event-model",
|
||||
"log",
|
||||
"pulldown-cmark",
|
||||
"serde",
|
||||
"serde_json",
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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<LogLine>,
|
||||
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<Mutex<Inner>>);
|
||||
|
||||
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<R>(&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<LogLine> {
|
||||
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<LogLine>, u64) {
|
||||
self.with(|inner| {
|
||||
let lines: Vec<LogLine> = 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<u64> {
|
||||
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::<Vec<_>>()
|
||||
.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<dyn log::Log>,
|
||||
}
|
||||
|
||||
impl RingLogger {
|
||||
pub fn new(ring: LogRing, inner: Box<dyn log::Log>) -> 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<dyn log::Log>,
|
||||
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<String> = 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<String> = 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<u64> = 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<Mutex<Vec<String>>>, 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<String> = ring.snapshot().into_iter().map(|l| l.message).collect();
|
||||
assert_eq!(held, ["kept", "filtered"], "the ring keeps both");
|
||||
}
|
||||
}
|
||||
@@ -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<String>,
|
||||
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<dyn Transport>,
|
||||
source: String,
|
||||
cursor: u64,
|
||||
status: Arc<Mutex<UploadStatus>>,
|
||||
}
|
||||
|
||||
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<dyn Transport>, source: impl Into<String>) -> 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<Mutex<UploadStatus>> {
|
||||
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::<Vec<_>>(),
|
||||
});
|
||||
|
||||
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<Signal>, Condvar)>,
|
||||
thread: Option<std::thread::JoinHandle<()>>,
|
||||
status: Arc<Mutex<UploadStatus>>,
|
||||
}
|
||||
|
||||
/// 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<dyn Transport>,
|
||||
source: impl Into<String>,
|
||||
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<Vec<serde_json::Value>>,
|
||||
status: Mutex<u16>,
|
||||
}
|
||||
|
||||
impl Fake {
|
||||
fn new() -> Arc<Self> {
|
||||
Arc::new(Self {
|
||||
posted: Mutex::new(Vec::new()),
|
||||
status: Mutex::new(200),
|
||||
})
|
||||
}
|
||||
fn bodies(&self) -> Vec<serde_json::Value> {
|
||||
self.posted.lock().unwrap().clone()
|
||||
}
|
||||
}
|
||||
|
||||
impl Transport for Fake {
|
||||
fn request(
|
||||
&self,
|
||||
method: &str,
|
||||
path: &str,
|
||||
body: Option<Body>,
|
||||
) -> Result<RawResponse, ApiError> {
|
||||
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<Box<dyn std::io::Read + Send>, 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");
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user