diff --git a/CHANGELOG.md b/CHANGELOG.md index 710918696..e1f8e9519 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,26 @@ # Changelog +## 1.0.10-beta.0 — 2026-09-30 + +### Features + +- Uploads, `fp` calls and evaluator calls carry a request id; the daemon also sends batch and machine ids, and `fp` errors show a `ref` (#872) + +### Fixes + +- The daemon logs each failed upload attempt at warning level with its request id and batch id, so every attempt of a batch can be found; before, only the last one was visible (#872) +- Parking a failed batch works when the spool and state directories are on different filesystems (e.g. separate Docker volumes); before, the batch stayed in the spool and was re-sent on every sweep. The move is crash-safe: the original is deleted only after the copy's directory entry is on disk (#872) + +### Docs + +- Troubleshooting, HTTP API and Cloud CLI pages explain the `ref` / `request_id` to quote to support, and where the daemon logs them (#872) + +### Dependencies + +- Python SDK dev lockfile: oauthlib 3.3.1 → 4.0.0 and pyjwt 2.13.0 → 2.15.1, clearing three OSV advisories (#872) +- Pin brace-expansion 5.0.9 → 5.0.12 (package.json override), clearing three OSV advisories (#872) +- Python SDK and `fp` CLI lockfiles: urllib3 2.7.0 → 2.8.0, clearing three OSV advisories (#872) + ## 1.0.9 — 2026-09-29 Action needed if you use Jev: its log-only mode is now `observe`, a `jev.json` still set to `shadow` is refused (Jev stays off until `failproofai jev setup` is run again), and Jev's checks now come only from `failproofai policies add FailproofAI/jev-policies`. For Hermes, `failproofai update` moves every profile from the old shell hooks (never run for cron jobs) to the native plugin, and every agent config failproofai edits is written crash-safely with a `.failproofai-backup`. Collects 1.0.9-beta.0 to beta.2 below. diff --git a/Cargo.lock b/Cargo.lock index ca33eae3c..6e055c6ff 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -283,12 +283,14 @@ dependencies = [ name = "fpai-collect" version = "1.0.10-beta.0" dependencies = [ + "getrandom 0.3.4", "notify", "reqwest", "rusqlite", "ruzstd", "serde", "serde_json", + "sha2", "time", "tokio", "tokio-rustls", diff --git a/bun.lock b/bun.lock index 5debbfb90..af3ba23e6 100644 --- a/bun.lock +++ b/bun.lock @@ -38,7 +38,7 @@ }, }, "overrides": { - "brace-expansion": "5.0.9", + "brace-expansion": "5.0.12", "browserslist": "4.28.8", "eslint-plugin-react-hooks": "7.0.1", "nanoid": "3.3.18", @@ -486,7 +486,7 @@ "bidi-js": ["bidi-js@1.1.0", "", { "dependencies": { "require-from-string": "^2.0.2" } }, "sha512-fX1Onk0tdVPC7obPWB5EbJ1z7NVhLq4m2xZLq2YXBkxzMXIGRpNMU88n0EPgWseKl12J7zXs7qrDxPK4sRs2fg=="], - "brace-expansion": ["brace-expansion@5.0.9", "", { "dependencies": { "balanced-match": "^4.0.2" } }, "sha512-ScQ4IuvIEF1TMlP7Zt+vjJ//9zlPb2SDcxWxM3bk8s6t6GGdJ7KO1dCcTidOPJKePW30LE/2cT7wCyPho9/Wxg=="], + "brace-expansion": ["brace-expansion@5.0.12", "", { "dependencies": { "balanced-match": "^4.0.2" } }, "sha512-YovQ3rzhaLMIrDjNDMkNS01tea93qhEhG5xy8f6+R0l+dw3Ki+5sCoIoI942iuLZTHWogWktgwVDhU09iNEimQ=="], "braces": ["braces@3.0.3", "", { "dependencies": { "fill-range": "^7.1.1" } }, "sha512-yQbXgO/OSZVD2IsiLlro+7Hf6Q18EJrKSEsdoMzKePKXct3gvD8oLcOQdIzGupr5Fj+EDe8gO/lxc1BzfMpxvA=="], diff --git a/crates/failproofaid/src/main.rs b/crates/failproofaid/src/main.rs index e8efe08b8..10e386723 100644 --- a/crates/failproofaid/src/main.rs +++ b/crates/failproofaid/src/main.rs @@ -790,8 +790,10 @@ fn collector_tasks() -> Vec { ingest.key.clone(), cfg.failed_dir.clone(), ) - .map(|u| u.with_redact(cfg.settings.redact)) - { + .map(|u| { + u.with_redact(cfg.settings.redact) + .with_machine_id(cfg.settings.machine_id.as_deref()) + }) { Ok(u) => std::sync::Arc::new(u), Err(err) => { eprintln!("[failproofaid] collector disabled: {err}"); diff --git a/crates/fpai-collect/Cargo.toml b/crates/fpai-collect/Cargo.toml index ff4c9339b..6edb91135 100644 --- a/crates/fpai-collect/Cargo.toml +++ b/crates/fpai-collect/Cargo.toml @@ -49,6 +49,11 @@ rusqlite = { version = "0.40", features = ["bundled"] } # reason rustls is chosen above: no C toolchain on the four cross-compiled legs. ruzstd = "0.9" time = { version = "0.3", default-features = false, features = ["std", "formatting", "macros"] } +# Request ids and batch ids on each upload (see uploader.rs). Both were already +# in the lockfile — sha2 via failproofaid, getrandom via the TLS stack — so +# neither adds a crate to the four cross-compiled legs; `uuid` would have. +getrandom = "0.3" +sha2 = "0.11" [dev-dependencies] tokio = { version = "1", features = ["rt-multi-thread", "macros", "time", "sync", "test-util", "net"] } diff --git a/crates/fpai-collect/src/delivery.rs b/crates/fpai-collect/src/delivery.rs index 8027cb54e..56106c446 100644 --- a/crates/fpai-collect/src/delivery.rs +++ b/crates/fpai-collect/src/delivery.rs @@ -37,7 +37,7 @@ use tokio::sync::Semaphore; use crate::spool::is_batch_file; use crate::supervisor::{Shutdown, TaskError}; -use crate::uploader::{ParkedName, Uploader}; +use crate::uploader::{ParkedName, Uploader, batch_id}; /// Concurrent uploads across the watcher and sweeper combined. /// @@ -141,7 +141,14 @@ impl Delivery { match self.uploader.upload_file(&path).await { Ok(()) => tracing::debug!(file = %path.display(), "uploaded"), - Err(err) => tracing::warn!(file = %path.display(), %err, "upload failed"), + // The batch id ties this to the attempt and park lines, which + // carry the request ids. + Err(err) => tracing::warn!( + file = %path.display(), + batch_id = %batch_id(&path), + %err, + "upload failed" + ), } } } @@ -324,6 +331,12 @@ async fn stale_batches(dir: &Path, min_age: Duration, max: usize) -> Vec bool { + // A dotfile is a park still being copied across filesystems. + !name.starts_with('.') && ParkedName::parse(name).is_auto_retryable() +} + async fn retry_parked(delivery: &Delivery, failed_dir: &Path, sd: &Shutdown) { let Ok(mut rd) = tokio::fs::read_dir(failed_dir).await else { return; @@ -336,7 +349,7 @@ async fn retry_parked(delivery: &Delivery, failed_dir: &Path, sd: &Shutdown) { let Some(name) = path.file_name().and_then(|n| n.to_str()) else { continue; }; - if !ParkedName::parse(name).is_auto_retryable() { + if !sweepable(name) { continue; } let Ok(meta) = entry.metadata().await else { @@ -414,6 +427,14 @@ mod tests { std::fs::remove_dir_all(&dir).ok(); } + #[test] + fn the_parked_sweep_skips_rejected_batches_and_half_copied_parks() { + assert!(sweepable("hooks-a-1-0.a1.jsonl")); + assert!(!sweepable("hooks-a-1-0.a1.c401.jsonl")); + assert!(!sweepable("hooks-a-1-0.a3.jsonl.poison")); + assert!(!sweepable(".hooks-a-1-0.a1.jsonl.partial")); + } + #[tokio::test] async fn the_sweeper_only_claims_batch_files() { let dir = tmpdir("filter"); diff --git a/crates/fpai-collect/src/uploader.rs b/crates/fpai-collect/src/uploader.rs index 310b59909..ee527d9a8 100644 --- a/crates/fpai-collect/src/uploader.rs +++ b/crates/fpai-collect/src/uploader.rs @@ -24,6 +24,14 @@ //! with the attempt count encoded in the filename, and only ever parked as //! `.poison` — never deleted. //! +//! **Every attempt says who it is.** Three headers go with each POST so a +//! failure can be found in the server's logs: `x-request-id` (new per attempt, +//! so the one that 503'd is distinguishable from the one that landed), +//! `x-failproofai-batch-id` (the same for every attempt of one batch, across +//! retries and the `failed/` re-drive) and `x-failproofai-machine-id`. All are +//! headers, never payload fields, for the same reason as the version header: +//! the server's dedup key hashes the payload. +//! //! **Oversized batches are split in memory.** Writing the chunks to disk //! beside the original would create files the watcher had never seen, so it //! would pick them up and post them concurrently with this function — the same @@ -53,6 +61,16 @@ pub const DEFAULT_RETRY_BASE: Duration = Duration::from_millis(1000); pub const DEFAULT_READ_TIMEOUT: Duration = Duration::from_secs(120); /// Attempts a parked batch gets before it is marked poison. pub const DEFAULT_FAILED_RETRIES_MAX: u32 = 3; +/// Header names. The server validates the request id as a W3C trace id (32 +/// lowercase hex) and mints its own when it is not one. +pub const REQUEST_ID_HEADER: &str = "x-request-id"; +pub const BATCH_ID_HEADER: &str = "x-failproofai-batch-id"; +pub const MACHINE_ID_HEADER: &str = "x-failproofai-machine-id"; +/// Sent when `collector.machine_id` is unset or not header-safe. Deliberately +/// NOT the telemetry id: that one identifies a person to analytics, and has no +/// business in a customer's server logs. +pub const UNKNOWN_MACHINE_ID: &str = "unknown"; +const MAX_MACHINE_ID_LEN: usize = 128; /// Ceiling applied to a server-supplied `Retry-After`, so a misconfigured /// header cannot park the uploader for hours. const MAX_RETRY_AFTER: Duration = Duration::from_secs(300); @@ -168,6 +186,7 @@ pub struct Uploader { retry_base: Duration, failed_retries_max: u32, redact: Redact, + machine_id: String, metrics: Arc, } @@ -201,6 +220,7 @@ impl Uploader { retry_base: DEFAULT_RETRY_BASE, failed_retries_max: DEFAULT_FAILED_RETRIES_MAX, redact: Redact::default(), + machine_id: UNKNOWN_MACHINE_ID.to_string(), metrics: Arc::new(UploadMetrics::default()), }) } @@ -210,6 +230,13 @@ impl Uploader { self } + /// `collector.machine_id`, sent on every upload. Unset, empty or not safe + /// to put in a header (see `header_safe_machine_id`) becomes `unknown`. + pub fn with_machine_id(mut self, machine_id: Option<&str>) -> Self { + self.machine_id = header_safe_machine_id(machine_id); + self + } + /// Shorten every delay. Tests only — without it each retry test would wait /// out a real multi-second backoff. #[doc(hidden)] @@ -257,8 +284,11 @@ impl Uploader { /// log context — this never deletes it. async fn post_batch(&self, path: &Path, bytes: Vec) -> Result<(), UploadError> { let mut attempt: u32 = 0; + let batch_id = batch_id(path); loop { + // New per attempt: one attempt is one server request. + let request_id = new_request_id(); let result = self .client .post(&self.url) @@ -268,6 +298,9 @@ impl Uploader { // server's content-hash dedup key — otherwise every upgrade // would make previously-shipped events look new. .header("X-Failproofai-Collector-Version", env!("CARGO_PKG_VERSION")) + .header(REQUEST_ID_HEADER, &request_id) + .header(BATCH_ID_HEADER, &batch_id) + .header(MACHINE_ID_HEADER, &self.machine_id) .body(bytes.clone()) .send() .await; @@ -277,6 +310,9 @@ impl Uploader { match result { Ok(resp) => { let status = resp.status(); + // The server's id for this request: its logs are keyed on + // it. Normally ours echoed back; ours if it sent none. + let request_id = response_request_id(&resp).unwrap_or(request_id); if status.is_success() { // A 2xx is necessary but NOT sufficient. Require the body // to actually be an ingest ack (a numeric `accepted`); @@ -301,8 +337,8 @@ impl Uploader { // attempt is encoded in the filename and becomes // `.poison` at `failed_retries_max`, after which it // is kept forever and never retried again. - if self.record_ack(path, &ack, attempt) { - self.park(path, None, attempt).await; + if self.record_ack(path, &ack, attempt, &request_id) { + self.park(path, None, &request_id).await; return Err(if ack.accepted == 0 { UploadError::StoredNothing { skipped: ack.skipped, @@ -318,7 +354,7 @@ impl Uploader { } Err(_) => { let code = status.as_u16(); - self.park(path, Some(code), attempt).await; + self.park(path, Some(code), &request_id).await; return Err(UploadError::Client { status: code }); } } @@ -330,27 +366,47 @@ impl Uploader { // park a batch the server explicitly asked us to resend. let retryable = status.is_server_error() || code == 408 || code == 429; if !retryable { - self.park(path, Some(code), attempt).await; + self.park(path, Some(code), &request_id).await; return Err(UploadError::Client { status: code }); } if attempt >= self.max_retries { - self.park(path, None, attempt).await; + self.park(path, None, &request_id).await; return Err(UploadError::Server { status: code, attempts: attempt, }); } let wait = retry_after(&resp).unwrap_or_else(|| self.backoff(attempt)); + // WARN, not DEBUG: each attempt has its own request id, and + // this line is the only place that ties it to the batch. + // At DEBUG a customer's log showed the last attempt (on the + // park line) and none of the ones before it. + tracing::warn!( + file = %path.display(), + request_id = %request_id, + batch_id = %batch_id, + attempt, + status = code, + "upload attempt failed; retrying" + ); tokio::time::sleep(wait).await; } Err(err) => { if attempt >= self.max_retries { - self.park(path, None, attempt).await; + self.park(path, None, &request_id).await; return Err(UploadError::Network { attempts: attempt, detail: error_chain(&err), }); } + tracing::warn!( + file = %path.display(), + request_id = %request_id, + batch_id = %batch_id, + attempt, + error = %error_chain(&err), + "upload attempt failed; retrying" + ); tokio::time::sleep(self.backoff(attempt)).await; } } @@ -366,7 +422,7 @@ impl Uploader { /// had it, and nothing anywhere recorded which event it was. That is the /// same permanent, silent loss the fully-skipped branch was added to close, /// reached through an ack that merely looks healthier. - fn record_ack(&self, path: &Path, ack: &IngestAck, attempt: u32) -> bool { + fn record_ack(&self, path: &Path, ack: &IngestAck, attempt: u32, request_id: &str) -> bool { let stored_nothing = ack.accepted == 0 && ack.skipped > 0; self.metrics @@ -421,6 +477,7 @@ impl Uploader { file = %path.display(), skipped = ack.skipped, attempt, + request_id, "the server accepted the request but stored NONE of its events; \ every line was rejected as malformed" ); @@ -429,6 +486,7 @@ impl Uploader { file = %path.display(), accepted = ack.accepted, skipped = ack.skipped, + request_id, "the server skipped some events in this batch" ); } @@ -459,10 +517,18 @@ impl Uploader { /// Never deletes and never overwrites: a collision gets a numeric suffix, /// because the file being moved is the only copy of data the server does /// not have. - async fn park(&self, path: &Path, client_status: Option, _attempt: u32) { - if let Err(err) = self.park_inner(path, client_status).await { + /// + /// `request_id` is the last attempt's — the server's echo when it answered — + /// and goes on the park line, so a `.poison` file on a customer laptop can + /// be matched to the exact server log lines that refused it. It is logged, + /// not written beside the batch: the filename is this module's only + /// metadata store (see `ParkedName`), and a sidecar is exactly what that + /// design rejects. + async fn park(&self, path: &Path, client_status: Option, request_id: &str) { + if let Err(err) = self.park_inner(path, client_status, request_id).await { tracing::error!( file = %path.display(), + request_id, %err, "could not park a failed batch; it stays in the spool and will be retried" ); @@ -473,6 +539,7 @@ impl Uploader { &self, path: &Path, client_status: Option, + request_id: &str, ) -> Result<(), std::io::Error> { tokio::fs::create_dir_all(&self.failed_dir).await?; let name = path @@ -493,10 +560,142 @@ impl Uploader { dest = %dest.display(), attempt = parked.attempt, poison = parked.poison, + request_id, + batch_id = %batch_id(path), "parking an undelivered batch" ); - tokio::fs::rename(path, &dest).await + move_file(path, &dest).await + } +} + +/// `rename`, or copy-then-remove when the two paths are on different +/// filesystems — an SDK spool on its own Docker volume, say. A plain rename +/// fails there with EXDEV, parking failed, and a batch the server rejected +/// (401, poison) stayed in the spool to be re-sent on every sweep. +async fn move_file(from: &Path, to: &Path) -> std::io::Result<()> { + match tokio::fs::rename(from, to).await { + Err(err) if err.kind() == std::io::ErrorKind::CrossesDevices => { + copy_then_remove(from, to).await + } + other => other, + } +} + +/// The cross-filesystem move. The copy is written under a hidden name and +/// renamed into place, so the retry sweep (which skips dotfiles) never sees a +/// half-written batch; the source goes only after that, so a crash midway +/// leaves a duplicate at worst — never a loss, and the server dedups events. +/// +/// "Only after that" has to mean after the rename is on disk, not just done: +/// the new directory entry is made durable (the directory is fsynced) before +/// the source is deleted. Otherwise a power cut can keep the deletion and lose +/// the entry, and the source was the last copy of the batch. +async fn copy_then_remove(from: &Path, to: &Path) -> std::io::Result<()> { + let name = to.file_name().and_then(|n| n.to_str()).unwrap_or("batch"); + let tmp = to.with_file_name(format!(".{name}.partial")); + let copied = async { + tokio::fs::copy(from, &tmp).await?; + tokio::fs::OpenOptions::new() + .write(true) + .open(&tmp) + .await? + .sync_all() + .await + } + .await; + if let Err(err) = copied { + let _ = tokio::fs::remove_file(&tmp).await; + return Err(err); + } + tokio::fs::rename(&tmp, to).await?; + if let Some(dir) = to.parent() { + sync_dir(dir).await?; + } + tokio::fs::remove_file(from).await +} + +/// fsync a directory, so the entries just created in it survive a crash. +/// Unix only: Windows cannot open a directory for flushing, and NTFS journals +/// the rename itself. +async fn sync_dir(dir: &Path) -> std::io::Result<()> { + #[cfg(unix)] + { + tokio::fs::File::open(dir).await?.sync_all().await?; + } + #[cfg(not(unix))] + let _ = dir; + Ok(()) +} + +/// A fresh request id: 32 lowercase hex, i.e. a W3C trace id — the shape the +/// server accepts as-is. Falls back to a clock-derived id if the OS RNG is +/// unavailable: an id only has to be distinct, and failing an upload over one +/// would trade data for a log label. +pub fn new_request_id() -> String { + let mut bytes = [0u8; 16]; + if getrandom::fill(&mut bytes).is_err() { + bytes = nanos_now().to_le_bytes(); + } + // Version-4 / variant bits, so it is also a valid UUID without dashes. + bytes[6] = (bytes[6] & 0x0f) | 0x40; + bytes[8] = (bytes[8] & 0x3f) | 0x80; + lower_hex(&bytes) +} + +/// The batch id for a spool or parked file: SHA-256 of its base name, as 32 +/// hex chars. +/// +/// The BASE name, not the file name: parking renames `x.jsonl` to +/// `x.a1.jsonl`, `x.a2.c503.jsonl`, `x.a3.jsonl.poison` (see `ParkedName`), and +/// the id must stay the same across all of them so "this batch was tried 7 +/// times over 3 hours" is one query. Hashed rather than sent raw because the +/// name carries a session tag. Chunks of one oversized file share the id — +/// they are one batch. +pub fn batch_id(path: &Path) -> String { + use sha2::{Digest, Sha256}; + let name = path + .file_name() + .and_then(|n| n.to_str()) + .unwrap_or("batch.jsonl"); + let base = ParkedName::parse(name).base; + let digest = Sha256::digest(base.as_bytes()); + lower_hex(&digest[..16]) +} + +/// `collector.machine_id` if it is safe to send as a header value, else +/// `unknown`: visible ASCII only (no spaces, controls or non-ASCII) and at most +/// 128 chars. An unsendable header would fail every upload, so it is replaced, +/// never rejected. +pub fn header_safe_machine_id(machine_id: Option<&str>) -> String { + match machine_id.map(str::trim) { + Some(id) + if !id.is_empty() + && id.len() <= MAX_MACHINE_ID_LEN + && id.bytes().all(|b| b.is_ascii_graphic()) => + { + id.to_string() + } + _ => UNKNOWN_MACHINE_ID.to_string(), + } +} + +/// The server's `x-request-id`, when it sent a sane one. +fn response_request_id(resp: &reqwest::Response) -> Option { + let raw = resp.headers().get(REQUEST_ID_HEADER)?.to_str().ok()?.trim(); + (!raw.is_empty() + && raw.len() <= 64 + && raw.bytes().all(|b| b.is_ascii_alphanumeric() || b == b'-')) + .then(|| raw.to_string()) +} + +fn lower_hex(bytes: &[u8]) -> String { + const HEX: &[u8; 16] = b"0123456789abcdef"; + let mut out = String::with_capacity(bytes.len() * 2); + for b in bytes { + out.push(HEX[(b >> 4) as usize] as char); + out.push(HEX[(b & 0x0f) as usize] as char); } + out } fn redact_batch(bytes: &[u8], mode: Redact) -> Vec { @@ -742,6 +941,68 @@ mod tests { assert_eq!(rejoined, body, "splitting must lose nothing"); } + #[test] + fn request_ids_are_distinct_trace_ids() { + let a = new_request_id(); + let b = new_request_id(); + assert_ne!(a, b); + for id in [&a, &b] { + assert_eq!(id.len(), 32); + assert!( + id.bytes().all(|c| matches!(c, b'0'..=b'9' | b'a'..=b'f')), + "{id}" + ); + } + } + + #[test] + fn the_batch_id_survives_every_rename_parking_makes() { + let id = batch_id(Path::new( + "/spool/claude-3ee9c788-1785741108180149712-0.jsonl", + )); + assert_eq!(id.len(), 32); + for parked in [ + "/failed/claude-3ee9c788-1785741108180149712-0.a1.jsonl", + "/failed/claude-3ee9c788-1785741108180149712-0.a2.c503.jsonl", + "/failed/claude-3ee9c788-1785741108180149712-0.a3.jsonl.poison", + ] { + assert_eq!(batch_id(Path::new(parked)), id, "{parked}"); + } + assert_ne!( + batch_id(Path::new( + "/spool/claude-3ee9c788-1785741108180149712-1.jsonl" + )), + id, + "the next batch of the same run is a different batch" + ); + // The raw name carries a session tag; the id must not. + assert!(!id.contains("3ee9c788")); + } + + #[test] + fn an_unsendable_machine_id_becomes_unknown() { + assert_eq!(header_safe_machine_id(Some("m-123")), "m-123"); + assert_eq!(header_safe_machine_id(Some(" m-123 ")), "m-123"); + for bad in [ + None, + Some(""), + Some(" "), + Some("has space"), + Some("new\nline"), + Some("naïve"), + ] { + assert_eq!(header_safe_machine_id(bad), UNKNOWN_MACHINE_ID, "{bad:?}"); + } + assert_eq!( + header_safe_machine_id(Some(&"m".repeat(129))), + UNKNOWN_MACHINE_ID + ); + assert_eq!( + header_safe_machine_id(Some(&"m".repeat(128))), + "m".repeat(128) + ); + } + #[test] fn a_body_under_the_cap_is_one_chunk() { let body = b"one\ntwo\n"; @@ -755,4 +1016,49 @@ mod tests { assert_eq!(chunks.len(), 1); assert_eq!(chunks[0], body); } + + #[cfg(unix)] + #[tokio::test] + async fn a_directory_can_be_flushed_and_a_missing_one_is_an_error() { + let dir = std::env::temp_dir().join(format!("fpai-sync-{}", new_request_id())); + std::fs::create_dir_all(&dir).unwrap(); + sync_dir(&dir).await.unwrap(); + std::fs::remove_dir_all(&dir).ok(); + // The source is only deleted after this succeeds, so a failure here + // must surface, not be swallowed. + assert!(sync_dir(&dir).await.is_err()); + } + + #[tokio::test] + async fn a_cross_filesystem_park_copies_then_removes_and_leaves_no_partial() { + let root = std::env::temp_dir().join(format!("fpai-park-{}", new_request_id())); + let (spool, failed) = (root.join("spool"), root.join("failed")); + std::fs::create_dir_all(&spool).unwrap(); + std::fs::create_dir_all(&failed).unwrap(); + let from = spool.join("hooks-a-1-0.jsonl"); + std::fs::write(&from, "{\"a\":1}\n{\"b\":2}\n").unwrap(); + let to = failed.join("hooks-a-1-0.a1.c401.jsonl"); + + // The EXDEV branch, driven directly: a test cannot count on two mounts. + copy_then_remove(&from, &to).await.unwrap(); + + assert!( + !from.exists(), + "the spool copy must go once the park is in place" + ); + assert_eq!( + std::fs::read_to_string(&to).unwrap(), + "{\"a\":1}\n{\"b\":2}\n" + ); + let left: Vec<_> = std::fs::read_dir(&failed) + .unwrap() + .map(|e| e.unwrap().file_name()) + .collect(); + assert_eq!( + left, + vec![to.file_name().unwrap().to_owned()], + "no .partial may remain" + ); + std::fs::remove_dir_all(&root).ok(); + } } diff --git a/crates/fpai-collect/tests/uploader.rs b/crates/fpai-collect/tests/uploader.rs index 1191143bf..68cb1cb6c 100644 --- a/crates/fpai-collect/tests/uploader.rs +++ b/crates/fpai-collect/tests/uploader.rs @@ -757,3 +757,107 @@ async fn a_partially_skipped_batch_is_parked_not_deleted() { fs::remove_dir_all(&spool).ok(); fs::remove_dir_all(&failed).ok(); } + +/// Every attempt carries the three identity headers: a request id that is new +/// per attempt, and a batch id that is the same across the retries of one batch +/// AND across the `failed/` re-drive, which renames the file. Without the +/// second property "this batch was tried N times" is not one server-log query. +#[tokio::test] +async fn every_attempt_carries_request_batch_and_machine_ids() { + let server = MockServer::start().await; + // Two 503s, then success: three attempts of one batch. + Mock::given(method("POST")) + .and(path("/events")) + .respond_with(ResponseTemplate::new(503)) + .up_to_n_times(2) + .mount(&server) + .await; + Mock::given(method("POST")) + .and(path("/events")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "accepted": 1, "skipped": 0 + }))) + .mount(&server) + .await; + + let spool = tmpdir("ids-spool"); + let failed = tmpdir("ids-failed"); + let batch = write_batch(&spool, "hooks-s-1-0.jsonl", 1); + let up = uploader(&server, &failed).with_machine_id(Some("m-123")); + up.upload_file(&batch).await.unwrap(); + + let requests = server.received_requests().await.unwrap(); + assert_eq!(requests.len(), 3); + let get = |r: &wiremock::Request, h: &str| { + r.headers + .get(h) + .unwrap_or_else(|| panic!("{h} missing")) + .to_str() + .unwrap() + .to_string() + }; + let request_ids: Vec = requests.iter().map(|r| get(r, "x-request-id")).collect(); + let batch_ids: Vec = requests + .iter() + .map(|r| get(r, "x-failproofai-batch-id")) + .collect(); + + for id in &request_ids { + assert_eq!(id.len(), 32, "{id}"); + assert!( + id.bytes().all(|c| matches!(c, b'0'..=b'9' | b'a'..=b'f')), + "{id}" + ); + } + let mut distinct = request_ids.clone(); + distinct.sort(); + distinct.dedup(); + assert_eq!( + distinct.len(), + 3, + "each attempt needs its own request id: {request_ids:?}" + ); + assert!( + batch_ids.iter().all(|b| b == &batch_ids[0]), + "{batch_ids:?}" + ); + assert!( + requests + .iter() + .all(|r| get(r, "x-failproofai-machine-id") == "m-123") + ); + + // The same batch re-driven from `failed/` under its parked name keeps its id. + assert_eq!( + fpai_collect::uploader::batch_id(Path::new("/failed/hooks-s-1-0.a1.jsonl")), + batch_ids[0] + ); + + fs::remove_dir_all(&spool).ok(); + fs::remove_dir_all(&failed).ok(); +} + +#[tokio::test] +async fn an_unset_machine_id_is_sent_as_unknown() { + let server = MockServer::start().await; + Mock::given(method("POST")) + .and(path("/events")) + .and(header("x-failproofai-machine-id", "unknown")) + .respond_with(ResponseTemplate::new(200).set_body_json(serde_json::json!({ + "accepted": 1, "skipped": 0 + }))) + .expect(1) + .mount(&server) + .await; + + let spool = tmpdir("mid-spool"); + let failed = tmpdir("mid-failed"); + let batch = write_batch(&spool, "hooks-s-1-0.jsonl", 1); + uploader(&server, &failed) + .upload_file(&batch) + .await + .unwrap(); + + fs::remove_dir_all(&spool).ok(); + fs::remove_dir_all(&failed).ok(); +} diff --git a/docs/reference/cloud-cli.mdx b/docs/reference/cloud-cli.mdx index 9eca539b5..3bab0efb5 100644 --- a/docs/reference/cloud-cli.mdx +++ b/docs/reference/cloud-cli.mdx @@ -365,7 +365,7 @@ What enforcement actually did. **Session-only**, same reason as above. | Flag | Description | | --- | --- | -| `--json` | Emit machine-readable JSON. | +| `--json` | Emit machine-readable JSON. Errors include the failed request's `request_id`. | | `--base-url ` | Use a self-hosted or development dashboard. | | `--org ` | Select an organization for this invocation. | | `--token ` | Override the saved user-session token. | diff --git a/docs/reference/http-api.mdx b/docs/reference/http-api.mdx index ea88e20eb..106bc94f6 100644 --- a/docs/reference/http-api.mdx +++ b/docs/reference/http-api.mdx @@ -69,6 +69,12 @@ The current specification has complete route, method, parameter, permission, and Use `Content-Type: application/json` for JSON writes. Treat `401` as missing or invalid authentication, `403` as a valid identity without the required permission, `404` as a missing or organization-inaccessible resource, `409` as a state conflict, and `422` as an invalid field or permission value. Error responses include a human-readable message; permission failures also name the required grant. +## Request IDs + +Every response carries an `X-Request-Id` header, and every JSON error body includes the same value as `request_id`. Quote it when you contact support: it identifies that one request. + +You can send your own `X-Request-Id` to correlate a request with your own logs. Use 32 lowercase hexadecimal characters, such as a UUID v4 with the dashes removed. Any other value is replaced with a new ID, which is returned in the response. + Policy enforcement deployment is intentionally managed outside the ordinary public `/v1` surface. Use the supported Cloud deployment workflow. diff --git a/docs/reference/troubleshooting.mdx b/docs/reference/troubleshooting.mdx index 535d47010..06f127817 100644 --- a/docs/reference/troubleshooting.mdx +++ b/docs/reference/troubleshooting.mdx @@ -186,6 +186,25 @@ icon: "wrench" + + + + Errors in the dashboard end with a short reference, for example `ref 4bf92f35`. It identifies that one request, and support can use it to find exactly what happened on the server. Copy it into your report as it appears. + + If a whole page fails to load, the error page shows a `digest` instead. Include that. + + + Human-readable `fp` errors end with the same `ref`. With `--json`, the error object carries the full `request_id`: + + ```bash + fp --json sessions --since 24h + ``` + + + When an upload fails, the daemon's log names a `request_id` and a `batch_id`: on Linux, `sudo journalctl -u failproofaid@$USER | grep batch_id`. Every attempt gets its own `request_id`; the `batch_id` stays the same across retries, so it ties the attempts of one batch together. Include both. + + + -When contacting support, include the CLI version, harness, environment, relevant session or deployment ID, and the output of `failproofai config --status` with secrets removed. +When contacting support, include the CLI version, harness, environment, relevant session or deployment ID, any `ref` or `request_id` from the error, and the output of `failproofai config --status` with secrets removed. diff --git a/fp-cloud-cli/fp_cli/app.py b/fp-cloud-cli/fp_cli/app.py index 56bf6bc7f..96ad7b745 100644 --- a/fp-cloud-cli/fp_cli/app.py +++ b/fp-cloud-cli/fp_cli/app.py @@ -163,6 +163,10 @@ def _format_error_as_line(error: click.ClickException) -> None: # exceptions (usage etc.) use their normal rendering. if isinstance(error, _FpCliError): message = error.message + # The ref is what a person pastes into a support ticket; the full id + # rides in the JSON envelope's `request_id` below. + if error.ref and not _wants_json(): + message = f"{message} · ref {error.ref}" else: message = error.format_message() if hasattr(error, "format_message") else str(error) # Click/Typer sometimes doubles the suggestion ("Did you mean 'x'? Did you diff --git a/fp-cloud-cli/fp_cli/auth.py b/fp-cloud-cli/fp_cli/auth.py index 3b443c5aa..41d78828f 100644 --- a/fp-cloud-cli/fp_cli/auth.py +++ b/fp-cloud-cli/fp_cli/auth.py @@ -10,12 +10,20 @@ from datetime import datetime, timedelta, timezone from typing import Any, Dict, Optional, Tuple +import uuid + import httpx from .config import CliConfig, save_config +from .client import _request_id_of from .errors import ApiError, AuthError, NetworkError +def _request_id_header() -> Dict[str, str]: + """A fresh request id for one login-flow call — same shape as the client's.""" + return {"x-request-id": uuid.uuid4().hex} + + def _iso(dt: datetime) -> str: return dt.astimezone(timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ") @@ -36,7 +44,11 @@ def request_otp( valid request (the server returns 200 even for unknown emails).""" try: with httpx.Client( - base_url=base_url.rstrip("/"), timeout=timeout, transport=transport, verify=verify + base_url=base_url.rstrip("/"), + timeout=timeout, + transport=transport, + verify=verify, + headers=_request_id_header(), ) as client: # Mark this as a CLI login so the server emails the paste-into-terminal # OTP template (with no "open the dashboard" button) instead of the @@ -52,6 +64,7 @@ def request_otp( raise ApiError( f"Failed to request a login code (HTTP {response.status_code}).", status=response.status_code, + request_id=_request_id_of(response), ) @@ -67,7 +80,11 @@ def verify_otp( """Exchange the code for a session token. Returns ``(token, expires_in_secs, user)``.""" try: with httpx.Client( - base_url=base_url.rstrip("/"), timeout=timeout, transport=transport, verify=verify + base_url=base_url.rstrip("/"), + timeout=timeout, + transport=transport, + verify=verify, + headers=_request_id_header(), ) as client: response = client.post( "/api/auth/otp/verify", json={"email": email, "code": code} @@ -76,18 +93,25 @@ def verify_otp( raise _network_error(base_url, exc) if response.status_code == 401: - raise AuthError("That code didn't match or has expired. Sign in again with fp login.") + raise AuthError( + "That code didn't match or has expired. Sign in again with fp login.", + request_id=_request_id_of(response), + ) if response.status_code == 429: raise ApiError( "Too many attempts — wait a bit, then run fp login again.", status=429, + request_id=_request_id_of(response), ) if response.status_code >= 400: # The dashboard proxy collapses the server's wrong/expired-code 401 into a 500 (its # `await res.json()` throws on the server's empty 401 body). So a non-401 4xx/5xx at the # verify step is, in practice, a bad/expired code — surface it as a clean auth failure, # not a raw "HTTP 500". (Real unreachability is a NetworkError, handled above.) - raise AuthError("That code didn't match or has expired. Sign in again with fp login.") + raise AuthError( + "That code didn't match or has expired. Sign in again with fp login.", + request_id=_request_id_of(response), + ) token = response.cookies.get("ae_session") if not token: @@ -151,6 +175,7 @@ def logout( with httpx.Client( base_url=base_url.rstrip("/"), cookies={"ae_session": token}, + headers=_request_id_header(), timeout=timeout, transport=transport, verify=verify, diff --git a/fp-cloud-cli/fp_cli/client.py b/fp-cloud-cli/fp_cli/client.py index c9d62afed..7368e2f9d 100644 --- a/fp-cloud-cli/fp_cli/client.py +++ b/fp-cloud-cli/fp_cli/client.py @@ -19,6 +19,7 @@ from __future__ import annotations import json as _json +import re import uuid from dataclasses import dataclass from enum import Enum @@ -264,8 +265,43 @@ def _path(ctx: ClientContext, path: str) -> str: return path +# A request id as the server may echo it: our own 32 hex, or a dashed UUID from +# an older server. Anything else is not put on an error a person will paste. +_SANE_REQUEST_ID = re.compile(r"^[A-Za-z0-9-]{1,64}$") + + +def _new_request_id() -> str: + """One id per HTTP request — a W3C trace id, the shape the server keeps as-is.""" + return uuid.uuid4().hex + + +def _request_id_of(response: httpx.Response) -> Optional[str]: + """The failed request's id: the server's echo, else the one this CLI sent. + + The server echoes ours back, so the two normally agree. The fallback matters + when something in front of the API answered instead — a proxy, a front door — + and a ref is still worth printing: our id is in the dashboard's access logs. + """ + echoed = (response.headers.get("x-request-id") or "").strip() + if _SANE_REQUEST_ID.match(echoed): + return echoed + try: + return response.request.headers.get("x-request-id") + except RuntimeError: # a Response built without a request (tests) + return None + + +def _sent_request_id(exc: httpx.RequestError) -> Optional[str]: + """The id of a request that never got an answer.""" + try: + return exc.request.headers.get("x-request-id") + except RuntimeError: + return None + + def _client(ctx: ClientContext, *, timeout: Any = None) -> httpx.Client: - headers = {"x-request-id": uuid.uuid4().hex} + # One id per client, and every call site opens a fresh client per request. + headers = {"x-request-id": _new_request_id()} cookies = None # Bearer XOR cookie — an `else`, never two independent `if`s. Sending both # would hand a human's `ae_session` to `/v1` alongside the key, and every @@ -355,7 +391,7 @@ def _raise_for_status(response: httpx.Response, ctx: ClientContext) -> None: "the dashboard's login page, so it landed on the web app instead of " "the API.", status=response.status_code, - request_id=response.headers.get("x-request-id"), + request_id=_request_id_of(response), hint="point --base-url at the server itself, e.g. http://localhost:8080", ) # Session mode is deliberately left alone for the /login case: that 3xx is @@ -375,22 +411,23 @@ def _raise_for_status(response: httpx.Response, ctx: ClientContext) -> None: "The request was redirected and did not reach the API, so it had no " "effect. Nothing was changed.", status=response.status_code, - request_id=response.headers.get("x-request-id"), + request_id=_request_id_of(response), hint=( "check --base-url: a redirect here usually means http:// where the " "server wants https://, or a front door in front of the API" ), ) return - request_id = response.headers.get("x-request-id") + request_id = _request_id_of(response) message = _extract_error(response) if response.status_code == 401: if key_mode: raise AuthError( "The API key was rejected. It may be revoked, mistyped, or issued by a " - "different deployment than --base-url points at." + "different deployment than --base-url points at.", + request_id=request_id, ) - raise AuthError("Session expired or not logged in. Run fp login.") + raise AuthError("Session expired or not logged in. Run fp login.", request_id=request_id) if response.status_code == 403: needed = _required_permission(response) if key_mode: @@ -406,10 +443,13 @@ def _raise_for_status(response: httpx.Response, ctx: ClientContext) -> None: raise ForbiddenError( f"{what}, or it cannot act for this org — the server answers 403 for both.", hint="check the key's grants, and the org you targeted with --org / FP_ORG", + request_id=request_id, ) if needed: - raise ForbiddenError(f"you don't have the {needed} permission") - raise ForbiddenError(message or "you don't have permission for this action") + raise ForbiddenError(f"you don't have the {needed} permission", request_id=request_id) + raise ForbiddenError( + message or "you don't have permission for this action", request_id=request_id + ) if response.status_code == 404: # In key mode a 404 has TWO very different causes, and the wrong reading # sends people hunting for a server bug that isn't there: @@ -434,7 +474,7 @@ def _raise_for_status(response: httpx.Response, ctx: ClientContext) -> None: request_id=request_id, hint="point --base-url at the server itself, e.g. http://localhost:8080", ) - raise NotFoundError(message or "Not found.") + raise NotFoundError(message or "Not found.", request_id=request_id) if response.status_code == 429: retry_after = response.headers.get("retry-after") wait = f" Retry after {retry_after}s." if retry_after else " Please wait a moment and try again." @@ -458,7 +498,8 @@ def _get_json(ctx: ClientContext, path: str, params: Optional[Dict[str, Any]] = response = client.get(url, params=clean) except httpx.RequestError as exc: raise NetworkError( - f"Cannot reach FailproofAI Cloud at {ctx.base_url}: {exc}" + f"Cannot reach FailproofAI Cloud at {ctx.base_url}: {exc}", + request_id=_sent_request_id(exc), ) _raise_for_status(response, ctx) # A 2xx with an empty or non-JSON body is anomalous for a read (e.g. a proxy or @@ -470,7 +511,7 @@ def _get_json(ctx: ClientContext, path: str, params: Optional[Dict[str, Any]] = raise ApiError( "The dashboard returned a malformed (non-JSON) response.", status=response.status_code, - request_id=response.headers.get("x-request-id"), + request_id=_request_id_of(response), ) @@ -495,7 +536,8 @@ def _request_json( response = client.request(method, url, json=json_body, params=clean) except httpx.RequestError as exc: raise NetworkError( - f"Cannot reach FailproofAI Cloud at {ctx.base_url}: {exc}" + f"Cannot reach FailproofAI Cloud at {ctx.base_url}: {exc}", + request_id=_sent_request_id(exc), ) _raise_for_status(response, ctx) # A genuinely empty body (204, or a 200 with no content) is a legitimate @@ -514,7 +556,7 @@ def _request_json( "The dashboard returned a malformed (non-JSON) response, so the request " "may not have been applied.", status=response.status_code, - request_id=response.headers.get("x-request-id"), + request_id=_request_id_of(response), ) @@ -575,7 +617,8 @@ def org_is_accessible(ctx: ClientContext, slug: str) -> bool: response = client.get(_path(probe, "/api/access-granters")) except httpx.RequestError as exc: raise NetworkError( - f"Cannot reach FailproofAI Cloud at {ctx.base_url}: {exc}" + f"Cannot reach FailproofAI Cloud at {ctx.base_url}: {exc}", + request_id=_sent_request_id(exc), ) if response.status_code == 200: return True diff --git a/fp-cloud-cli/fp_cli/errors.py b/fp-cloud-cli/fp_cli/errors.py index ca42a8cbe..7162fedad 100644 --- a/fp-cloud-cli/fp_cli/errors.py +++ b/fp-cloud-cli/fp_cli/errors.py @@ -26,9 +26,24 @@ class FpCliError(click.ClickException): exit_code = 1 - def __init__(self, message: str, *, hint: Optional[str] = None) -> None: + def __init__( + self, + message: str, + *, + hint: Optional[str] = None, + request_id: Optional[str] = None, + ) -> None: super().__init__(message) self.hint = hint + # The id of the HTTP request that failed — the server's echo, else the + # one this CLI sent. On every typed error, not only ApiError, so a 403 or + # a 404 is as traceable in the server's logs as a 500. + self.request_id = request_id + + @property + def ref(self) -> Optional[str]: + """What a person reads and pastes: the first 8 chars of the request id.""" + return self.request_id[:8] if self.request_id else None class KeyModeUnsupportedError(FpCliError): @@ -85,9 +100,8 @@ def __init__( request_id: Optional[str] = None, hint: Optional[str] = None, ) -> None: - super().__init__(message, hint=hint) + super().__init__(message, hint=hint, request_id=request_id) self.status = status - self.request_id = request_id def format_message(self) -> str: parts = [self.message] diff --git a/fp-cloud-cli/tests/test_request_id.py b/fp-cloud-cli/tests/test_request_id.py new file mode 100644 index 000000000..0136a5f1a --- /dev/null +++ b/fp-cloud-cli/tests/test_request_id.py @@ -0,0 +1,111 @@ +"""Every CLI request carries its own id, and every failure names it. + +A customer pastes `ref 4bf92f35`; support finds every server log line of that +request with one query. That only works if (a) each HTTP request has a distinct +id, (b) every typed error — not just ApiError — carries it, and (c) the id is +still there when something in front of the API answered without echoing one. +""" + +from __future__ import annotations + +import json +import re + +import httpx +import pytest +import respx + +from fp_cli import client as api +from fp_cli.app import app +from fp_cli.client import ClientContext +from fp_cli.errors import ApiError, AuthError, ForbiddenError, NetworkError, NotFoundError + +BASE = "http://dash.test" +TRACE_ID = re.compile(r"^[0-9a-f]{32}$") +SERVER_ID = "4bf92f3577b34da6a3ce929d0e0e4736" + + +def ctx() -> ClientContext: + return ClientContext(base_url=BASE, token="tok") + + +@respx.mock +def test_each_request_gets_its_own_trace_id(): + route = respx.get(f"{BASE}/api/sessions").mock( + return_value=httpx.Response(200, json={"sessions": [], "next_cursor": None}) + ) + api.list_sessions(ctx()) + api.list_sessions(ctx()) + ids = [call.request.headers["x-request-id"] for call in route.calls] + assert len(ids) == 2 + assert all(TRACE_ID.match(i) for i in ids), ids + assert ids[0] != ids[1] + + +@pytest.mark.parametrize( + ("status", "error_type"), + [(401, AuthError), (403, ForbiddenError), (404, NotFoundError), (500, ApiError)], +) +@respx.mock +def test_every_typed_error_carries_the_servers_request_id(status, error_type): + respx.get(f"{BASE}/api/sessions").mock( + return_value=httpx.Response(status, json={"error": "x"}, headers={"x-request-id": SERVER_ID}) + ) + with pytest.raises(error_type) as excinfo: + api.list_sessions(ctx()) + assert excinfo.value.request_id == SERVER_ID + assert excinfo.value.ref == SERVER_ID[:8] + + +@respx.mock +def test_without_an_echo_the_error_names_the_id_we_sent(): + # A proxy or front door answered: no x-request-id on the response. The id we + # sent is still the one in the dashboard's access logs. + route = respx.get(f"{BASE}/api/sessions").mock(return_value=httpx.Response(502, text="bad gateway")) + with pytest.raises(ApiError) as excinfo: + api.list_sessions(ctx()) + sent = route.calls.last.request.headers["x-request-id"] + assert excinfo.value.request_id == sent + + +@respx.mock +def test_an_unsane_echo_is_not_put_on_the_error(): + route = respx.get(f"{BASE}/api/sessions").mock( + return_value=httpx.Response(500, json={"error": "x"}, headers={"x-request-id": "a b c"}) + ) + with pytest.raises(ApiError) as excinfo: + api.list_sessions(ctx()) + assert excinfo.value.request_id == route.calls.last.request.headers["x-request-id"] + + +@respx.mock +def test_a_network_error_names_the_id_it_sent(): + route = respx.get(f"{BASE}/api/sessions").mock(side_effect=httpx.ConnectError("refused")) + with pytest.raises(NetworkError) as excinfo: + api.list_sessions(ctx()) + assert excinfo.value.request_id == route.calls.last.request.headers["x-request-id"] + # Exit code contract unchanged. + assert excinfo.value.exit_code == 3 + + +@respx.mock +def test_the_human_error_line_shows_the_ref(logged_in, runner): + respx.get(f"{BASE}/api/sessions").mock( + return_value=httpx.Response(403, json={"error": "forbidden"}, headers={"x-request-id": SERVER_ID}) + ) + result = runner.invoke(app, ["sessions"]) + assert result.exit_code == 5 + assert f"ref {SERVER_ID[:8]}" in result.stderr + + +@respx.mock +def test_the_json_envelope_carries_the_full_id_for_non_api_errors(logged_in, runner): + respx.get(f"{BASE}/api/sessions").mock( + return_value=httpx.Response(403, json={"error": "forbidden"}, headers={"x-request-id": SERVER_ID}) + ) + result = runner.invoke(app, ["--json", "sessions"]) + assert result.exit_code == 5 + data = json.loads(result.stdout) + assert data["request_id"] == SERVER_ID + # The ref is for humans; the envelope's message stays clean. + assert "ref " not in data["error"] diff --git a/fp-cloud-cli/uv.lock b/fp-cloud-cli/uv.lock index 424ec5b17..c21e70fd1 100644 --- a/fp-cloud-cli/uv.lock +++ b/fp-cloud-cli/uv.lock @@ -541,9 +541,9 @@ wheels = [ [[package]] name = "urllib3" -version = "2.7.0" +version = "2.8.0" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" } +sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" }, + { url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" }, ] diff --git a/package.json b/package.json index 2ac14d279..db11a30b1 100644 --- a/package.json +++ b/package.json @@ -117,7 +117,7 @@ "eslint-plugin-react-hooks": "7.0.1", "vite": "8.0.16", "undici": "7.29.1", - "brace-expansion": "5.0.9", + "brace-expansion": "5.0.12", "sharp": "0.35.4", "browserslist": "4.28.8" } diff --git a/sdk/python/failproofai_sdk/evaluator/__init__.py b/sdk/python/failproofai_sdk/evaluator/__init__.py index d735bc929..a70735113 100644 --- a/sdk/python/failproofai_sdk/evaluator/__init__.py +++ b/sdk/python/failproofai_sdk/evaluator/__init__.py @@ -14,7 +14,12 @@ Metric, Score, ) -from failproofai_sdk.evaluator.client import EvaluatorAPIError, EvaluatorClient +from failproofai_sdk.evaluator.client import ( + EvaluatorAPIError, + EvaluatorClient, + current_request_id, + request_id_scope, +) from failproofai_sdk.evaluator.protocol import ( Assignment, AssignmentDefinition, @@ -70,6 +75,8 @@ "Evaluator", "ManagedCompiler", "EvaluatorAPIError", + "current_request_id", + "request_id_scope", "EvaluatorClient", "EvaluatorKind", "HeartbeatRequest", diff --git a/sdk/python/failproofai_sdk/evaluator/client.py b/sdk/python/failproofai_sdk/evaluator/client.py index fc3c28ce7..e2a3c7d79 100644 --- a/sdk/python/failproofai_sdk/evaluator/client.py +++ b/sdk/python/failproofai_sdk/evaluator/client.py @@ -2,11 +2,15 @@ from __future__ import annotations +import contextvars import ipaddress import json import random +import re import time -from collections.abc import Callable, Mapping +import uuid +from collections.abc import Callable, Iterator, Mapping +from contextlib import contextmanager from typing import Any from urllib.error import HTTPError, URLError from urllib.parse import urljoin, urlsplit @@ -42,6 +46,56 @@ _DEFAULT_RESPONSE_LIMIT = 2 * 1024 * 1024 _RETRYABLE_HTTP_STATUSES = frozenset({429, 502, 503, 504}) +REQUEST_ID_HEADER = "X-Request-Id" +# A W3C trace id: the shape the server keeps as-is (it mints its own for +# anything else, which would break the join between this worker's log line and +# the server's). +_TRACE_ID = re.compile(r"^[0-9a-f]{32}$") +# What we accept back from a server before putting it on an error or a log line. +_SANE_ECHO = re.compile(r"^[A-Za-z0-9-]{1,64}$") + +#: The request id the next API call should carry. Unset, every call mints its +#: own. A caller that already has one — a job it is working on behalf of, or a +#: host that correlates its own logs — sets it with :func:`request_id_scope`, and +#: every call inside the scope sends it. ``asyncio.to_thread`` copies context, so +#: a scope opened in the runtime reaches the client's worker thread. +current_request_id: contextvars.ContextVar[str | None] = contextvars.ContextVar( + "failproofai_evaluator_request_id", default=None +) + + +def new_request_id() -> str: + return uuid.uuid4().hex + + +@contextmanager +def request_id_scope(request_id: str | None = None) -> Iterator[str]: + """Send ``request_id`` (or a fresh one) on every API call inside the block. + + An id that is not a trace id (32 lowercase hex) is replaced with a fresh one + rather than sent, so a caller can never put free text on the wire. + """ + value = request_id if request_id and _TRACE_ID.match(request_id) else new_request_id() + token = current_request_id.set(value) + try: + yield value + finally: + current_request_id.reset(token) + + +def _request_id_for_call() -> str: + current = current_request_id.get() + return current if current and _TRACE_ID.match(current) else new_request_id() + + +def _echoed_request_id(headers: Any) -> str | None: + """The server's ``x-request-id`` from a response's headers, when sane.""" + if headers is None: + return None + value = headers.get(REQUEST_ID_HEADER) or headers.get(REQUEST_ID_HEADER.lower()) + value = (value or "").strip() + return value if _SANE_ECHO.match(value) else None + class _RejectRedirects(HTTPRedirectHandler): def redirect_request(self, request, file_pointer, code, message, headers, new_url): @@ -222,6 +276,11 @@ def _json( separators=(",", ":"), ).encode("utf-8") request_headers["Content-Type"] = "application/json" + # One id per call, shared by its retries: a retried call is one logical + # request, and the server's log lines for every attempt then group under + # one id. The caller's scope id wins (see `current_request_id`). + request_id = _request_id_for_call() + request_headers[REQUEST_ID_HEADER] = request_id if headers: request_headers.update(headers) @@ -232,11 +291,18 @@ def _json( ) try: with self._opener(request, timeout=self._timeout_seconds) as response: - return self._decode( - response.read(response_limit + 1), response_limit - ) + try: + return self._decode( + response.read(response_limit + 1), response_limit + ) + except EvaluatorAPIError as error: + error.request_id = ( + _echoed_request_id(getattr(response, "headers", None)) + or request_id + ) + raise except HTTPError as error: - api_error = self._http_error(error, response_limit) + api_error = self._http_error(error, response_limit, request_id) if attempt + 1 == attempts or not api_error.retryable: raise api_error from error except (URLError, TimeoutError, OSError) as error: @@ -246,6 +312,8 @@ def _json( code="transport_error", message=str(error), retryable=True, + # Never answered: ours is the only id there is. + request_id=request_id, ) from error # Jitter is scheduling noise, not a security decision. self._sleeper(random.uniform(0, min(0.25 * (2**attempt), 2.0))) # nosec B311 @@ -279,7 +347,13 @@ def _decode(raw: bytes, limit: int) -> dict[str, Any]: return value @classmethod - def _http_error(cls, error: HTTPError, limit: int) -> EvaluatorAPIError: + def _http_error( + cls, error: HTTPError, limit: int, sent_request_id: str | None = None + ) -> EvaluatorAPIError: + # The server's id for this request: its logs are keyed on it. The + # envelope's own field first, then the response header, then ours — a + # proxy that answered in the server's place echoes nothing. + echoed = _echoed_request_id(error.headers) or sent_request_id raw = error.read(limit + 1) try: response = ErrorResponse.from_wire(cls._decode(raw, limit)) @@ -289,11 +363,12 @@ def _http_error(cls, error: HTTPError, limit: int) -> EvaluatorAPIError: code="http_error", message=f"server returned HTTP {error.code}", retryable=error.code in _RETRYABLE_HTTP_STATUSES, + request_id=echoed, ) return EvaluatorAPIError( status=error.code, code=response.error.code, message=response.error.message, retryable=response.error.retryable, - request_id=response.error.request_id, + request_id=response.error.request_id or echoed, ) diff --git a/sdk/python/failproofai_sdk/evaluator/runtime.py b/sdk/python/failproofai_sdk/evaluator/runtime.py index bfbc3a2e2..451403961 100644 --- a/sdk/python/failproofai_sdk/evaluator/runtime.py +++ b/sdk/python/failproofai_sdk/evaluator/runtime.py @@ -289,9 +289,17 @@ async def run_forever(self) -> None: ) except EvaluatorAPIError as error: self._increment("claim_failures") + # The id rides in the message as well as `extra`: a plain + # `logging.basicConfig` formatter drops extras, and this is + # the id that finds the server's side of the failure. logger.warning( - "evaluator claim failed", - extra={"code": error.code, "retryable": error.retryable}, + "evaluator claim failed (request_id=%s)", + error.request_id, + extra={ + "code": error.code, + "retryable": error.retryable, + "request_id": error.request_id, + }, ) if not error.retryable: raise @@ -372,8 +380,11 @@ async def process_assignment(self, assignment: Assignment) -> None: # transcript that can never shrink. Log and return; the server # terminalizes the assignment as `too_large`. logger.warning( - "assignment %s transcript is too large to evaluate; skipping", + "assignment %s transcript is too large to evaluate; skipping " + "(request_id=%s)", assignment.assignment_id, + error.request_id, + extra={"request_id": error.request_id}, ) self._increment("transcripts_too_large") return @@ -796,10 +807,12 @@ async def _heartbeat( task.cancel() return logger.warning( - "evaluator heartbeat failed", + "evaluator heartbeat failed (request_id=%s)", + error.request_id, extra={ "assignment_id": assignment.assignment_id, "code": error.code, + "request_id": error.request_id, }, ) self._increment("heartbeat_failures") @@ -924,6 +937,12 @@ def _consume_task(task: asyncio.Task[None]) -> None: task.result() except asyncio.CancelledError: pass + except EvaluatorAPIError as error: + logger.exception( + "evaluator assignment failed (request_id=%s)", + error.request_id, + extra={"code": error.code, "request_id": error.request_id}, + ) except Exception: logger.exception("evaluator assignment failed") diff --git a/sdk/python/tests/test_evaluator_request_id.py b/sdk/python/tests/test_evaluator_request_id.py new file mode 100644 index 000000000..d57d19fef --- /dev/null +++ b/sdk/python/tests/test_evaluator_request_id.py @@ -0,0 +1,197 @@ +"""Every evaluator API call carries a request id, and every failure names one. + +The server logs each worker call under its `x-request-id`. A worker log line +that says `request_id=4bf92f35…` is then one query away from the server's side +of the same failure — but only if the worker sent an id the server keeps (a +32-hex trace id), and only if the error carries the id the server logged. +""" + +from __future__ import annotations + +import asyncio +import io +import json +import logging +import re +from email.message import Message +from urllib.error import HTTPError, URLError + +import pytest + +from failproofai_sdk.evaluator import ( + ClaimRequest, + EvaluatorAPIError, + EvaluatorClient, + current_request_id, + request_id_scope, +) + +TRACE_ID = re.compile(r"^[0-9a-f]{32}$") +SERVER_ID = "4bf92f3577b34da6a3ce929d0e0e4736" + + +class Response: + def __init__(self, body, headers=None): + self.body = json.dumps(body).encode() if not isinstance(body, bytes) else body + self.headers = headers or {} + + def __enter__(self): + return self + + def __exit__(self, *args): + return None + + def read(self, amount): + return self.body[:amount] + + +def _client(opener): + return EvaluatorClient( + base_url="https://cloud.example/", + credential="secret", + opener=opener, + sleeper=lambda _: None, + ) + + +def _claim(client): + return client.claim(ClaimRequest("worker", "sha256:x", 1, 20)) + + +def _headers(**values): + message = Message() + for key, value in values.items(): + message[key.replace("_", "-")] = value + return message + + +def test_every_call_sends_a_fresh_trace_id(): + sent = [] + + def opener(request, timeout): + sent.append(request.get_header("X-request-id")) + return Response({"assignments": [], "lease_duration_seconds": 30}) + + client = _client(opener) + for _ in range(2): + try: + _claim(client) + except Exception: # noqa: BLE001 - the response shape is not under test here + pass + assert len(sent) == 2 + assert all(TRACE_ID.match(i or "") for i in sent), sent + assert sent[0] != sent[1] + + +def test_retries_of_one_call_share_its_id(): + sent = [] + + def opener(request, timeout): + sent.append(request.get_header("X-request-id")) + raise URLError("reset") + + with pytest.raises(EvaluatorAPIError) as caught: + # heartbeat retries; claim deliberately does not + _client(opener)._json("POST", "/v1/evaluator/runs/heartbeat", None, retry=True) + assert len(sent) == 4 + assert len(set(sent)) == 1 + # A call that never got an answer names the id it sent. + assert caught.value.request_id == sent[0] + + +def test_a_caller_scope_id_is_honoured(): + sent = [] + + def opener(request, timeout): + sent.append(request.get_header("X-request-id")) + raise URLError("offline") + + with request_id_scope(SERVER_ID) as scoped: + assert scoped == SERVER_ID + with pytest.raises(EvaluatorAPIError): + _claim(_client(opener)) + assert sent == [SERVER_ID] + assert current_request_id.get() is None, "the scope must not leak" + + +def test_a_scope_reaches_the_thread_the_runtime_calls_the_client_on(): + # The runtime calls the client via asyncio.to_thread, which copies context. + sent = [] + + def opener(request, timeout): + sent.append(request.get_header("X-request-id")) + raise URLError("offline") + + async def main(): + with request_id_scope(SERVER_ID): + with pytest.raises(EvaluatorAPIError): + await asyncio.to_thread(_claim, _client(opener)) + + asyncio.run(main()) + assert sent == [SERVER_ID] + + +def test_a_scope_with_free_text_is_replaced_not_sent(): + with request_id_scope("evil\nheader") as scoped: + assert TRACE_ID.match(scoped) + + +def test_an_http_error_carries_the_servers_echoed_id(): + def opener(request, timeout): + raise HTTPError( + request.full_url, 503, "Unavailable", _headers(x_request_id=SERVER_ID), io.BytesIO(b"") + ) + + with pytest.raises(EvaluatorAPIError) as caught: + _claim(_client(opener)) + assert caught.value.status == 503 + assert caught.value.request_id == SERVER_ID + + +def test_without_an_echo_the_error_names_the_id_it_sent(): + sent = [] + + def opener(request, timeout): + sent.append(request.get_header("X-request-id")) + raise HTTPError(request.full_url, 502, "Bad Gateway", _headers(), io.BytesIO(b"")) + + with pytest.raises(EvaluatorAPIError) as caught: + _claim(_client(opener)) + assert caught.value.request_id == sent[0] + + +def test_a_malformed_success_body_still_names_an_id(): + def opener(request, timeout): + return Response(b"", headers=_headers(x_request_id=SERVER_ID)) + + with pytest.raises(EvaluatorAPIError) as caught: + _claim(_client(opener)) + assert caught.value.code == "invalid_response" + assert caught.value.request_id == SERVER_ID + + +def test_the_runtimes_claim_failure_line_names_the_request_id(caplog): + # Drive the real runtime path: a claim the server refuses must leave a warn + # line whose MESSAGE carries the id (a basicConfig formatter drops extras). + from tests.test_evaluator_runtime import FakeClient, _runtime + + from failproofai_sdk.evaluator import Evaluator + + class RejectedClaimClient(FakeClient): + def claim(self, request): + raise EvaluatorAPIError( + status=409, + code="catalog_mismatch", + message="register again", + retryable=False, + request_id=SERVER_ID, + ) + + runtime = _runtime(Evaluator(name="test", version="1"), RejectedClaimClient()) + with caplog.at_level(logging.WARNING, logger="failproofai_sdk.evaluator"): + with pytest.raises(EvaluatorAPIError): + asyncio.run(runtime.run_forever()) + lines = [r for r in caplog.records if "claim failed" in r.getMessage()] + assert lines, [r.getMessage() for r in caplog.records] + assert SERVER_ID in lines[0].getMessage() + assert lines[0].request_id == SERVER_ID diff --git a/sdk/python/uv.lock b/sdk/python/uv.lock index 4ab359166..106a304d4 100644 --- a/sdk/python/uv.lock +++ b/sdk/python/uv.lock @@ -2934,11 +2934,11 @@ wheels = [ [[package]] name = "oauthlib" -version = "3.3.1" +version = "4.0.0" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/0b/5f/19930f824ffeb0ad4372da4812c50edbd1434f678c90c2733e1188edfc63/oauthlib-3.3.1.tar.gz", hash = "sha256:0f0f8aa759826a193cf66c12ea1af1637f87b9b4622d46e866952bb022e538c9", size = 185918, upload-time = "2025-06-19T22:48:08.269Z" } +sdist = { url = "https://files.pythonhosted.org/packages/7a/d8/a1bcc8ba112a627f8ffbdc212a78ce18d3ac07e91a5ca65d27918eee25a1/oauthlib-4.0.0.tar.gz", hash = "sha256:efb274799819440f95b4ab3b818869f1ce9ae26c5beacba0201d1a1b76b54f86", size = 187232, upload-time = "2026-09-28T06:01:18.77Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/be/9c/92789c596b8df838baa98fa71844d84283302f7604ed565dafe5a6b5041a/oauthlib-3.3.1-py3-none-any.whl", hash = "sha256:88119c938d2b8fb88561af5f6ee0eec8cc8d552b7bb1f712743136eb7523b7a1", size = 160065, upload-time = "2025-06-19T22:48:06.508Z" }, + { url = "https://files.pythonhosted.org/packages/d9/f4/78229a1066068ca14fc60fb26cf7381cabe4382261392b90e5f9552722d4/oauthlib-4.0.0-py3-none-any.whl", hash = "sha256:624c28c13a0a59cabf9747dfa52af63be3e512a7f2714df16e91b5b3a145e6cd", size = 159715, upload-time = "2026-09-28T06:01:17.008Z" }, ] [[package]] @@ -4082,14 +4082,14 @@ wheels = [ [[package]] name = "pyjwt" -version = "2.13.0" +version = "2.15.1" source = { registry = "https://pypi.org/simple" } dependencies = [ { name = "typing-extensions", marker = "python_full_version < '3.11'" }, ] -sdist = { url = "https://files.pythonhosted.org/packages/3b/81/58d0ac84e1ef3a3843791d6954d94c0b33d526c75eeb1efbce9d0a4c4077/pyjwt-2.13.0.tar.gz", hash = "sha256:41571c89ca91598c79e8ef18a2d07367d4810fbbd6f637794879baf1b7703423", size = 107515, upload-time = "2026-05-21T19:54:36.618Z" } +sdist = { url = "https://files.pythonhosted.org/packages/43/ea/5194e52748b0da83d71e082d75496eaec6e58f419f5e184786ded517e6a9/pyjwt-2.15.1.tar.gz", hash = "sha256:4f259e80cdfb6b3fc18a7de51fd1ef9ec79652f25019bae68975ca2468a34df8", size = 121252, upload-time = "2026-09-28T18:40:42.598Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/a3/5e/ecf12fdb62546d64385c158514e9b2b671f7832108ef2ecd2020ce0af2d1/pyjwt-2.13.0-py3-none-any.whl", hash = "sha256:66adcc2aff09b3f1bbd95fc1e1577df8ac8723c978552fd43304c8a290ac5728", size = 31274, upload-time = "2026-05-21T19:54:35.362Z" }, + { url = "https://files.pythonhosted.org/packages/50/ca/44de4e75f8aadc457f0634be3b542815078ded46dca30efb960edeecad6e/pyjwt-2.15.1-py3-none-any.whl", hash = "sha256:42d59d631f7768a1028a64c7ff581a9bf7519804daf91fc5b6c56e30eec5e193", size = 33860, upload-time = "2026-09-28T18:40:41.429Z" }, ] [package.optional-dependencies] @@ -5120,11 +5120,11 @@ wheels = [ [[package]] name = "urllib3" -version = "2.7.0" +version = "2.8.0" source = { registry = "https://pypi.org/simple" } -sdist = { url = "https://files.pythonhosted.org/packages/53/0c/06f8b233b8fd13b9e5ee11424ef85419ba0d8ba0b3138bf360be2ff56953/urllib3-2.7.0.tar.gz", hash = "sha256:231e0ec3b63ceb14667c67be60f2f2c40a518cb38b03af60abc813da26505f4c", size = 433602, upload-time = "2026-05-07T16:13:18.596Z" } +sdist = { url = "https://files.pythonhosted.org/packages/e3/05/b17359e1cefb4f909b5e40b1b90a496d987258916dbbf88e842c729f510e/urllib3-2.8.0.tar.gz", hash = "sha256:63bf2ead4c879426ebf22ef2a781eeb4aa3b4ae798a0435506f8687fd5bb9b63", size = 458972, upload-time = "2026-09-15T19:29:36.253Z" } wheels = [ - { url = "https://files.pythonhosted.org/packages/7f/3e/5db95bcf282c52709639744ca2a8b149baccf648e39c8cc87553df9eae0c/urllib3-2.7.0-py3-none-any.whl", hash = "sha256:9fb4c81ebbb1ce9531cce37674bbc6f1360472bc18ca9a553ede278ef7276897", size = 131087, upload-time = "2026-05-07T16:13:17.151Z" }, + { url = "https://files.pythonhosted.org/packages/92/9d/c4e665119135114480843e7ab388fa94d8480650450e6f8e26b70d323a4c/urllib3-2.8.0-py3-none-any.whl", hash = "sha256:0cf3cae568d36aa9576b28dfb35f11328f1cb974ca7647d9475ebb86c75ac6e3", size = 135717, upload-time = "2026-09-15T19:29:34.577Z" }, ] [[package]]