diff --git a/src/openhuman/memory/tree/jobs/mod.rs b/src/openhuman/memory/tree/jobs/mod.rs index cdb71ccb0..77e8287f2 100644 --- a/src/openhuman/memory/tree/jobs/mod.rs +++ b/src/openhuman/memory/tree/jobs/mod.rs @@ -26,6 +26,7 @@ //! See [`store::enqueue_tx`] for the in-tx producer entry point. mod handlers; +mod redact; pub mod scheduler; pub mod store; pub mod testing; diff --git a/src/openhuman/memory/tree/jobs/redact.rs b/src/openhuman/memory/tree/jobs/redact.rs new file mode 100644 index 000000000..6c0099a68 --- /dev/null +++ b/src/openhuman/memory/tree/jobs/redact.rs @@ -0,0 +1,226 @@ +//! Log-side scrubber for free-form `reason` / `last_error` strings emitted +//! by [`worker::run_once`] (defer/fail branches) and the matching +//! [`store::mark_failed`] / [`store::mark_deferred`] log lines. +//! +//! Persisted DB state (`mem_tree_jobs.last_error`) keeps the full +//! original string for diagnostics — handlers may attach +//! upstream-provider responses or full anyhow chains there. Logs are a +//! lower-trust sink (often forwarded to remote log aggregators or +//! shared in bug reports), so we apply a uniform scrub policy to: +//! +//! 1. Mask credential-shaped tokens (Bearer, OpenAI `sk-…`, GitHub +//! `ghp_…`, Slack `xox?-…`, generic `api_key=…` / `password=…` / +//! `token=…` assignments). +//! 2. Strip URL userinfo (`https://user:pass@host` → `https://***@host`). +//! 3. Mask bare email addresses. +//! 4. Cap the logged string at [`MAX_LEN`] bytes, suffixed with +//! `…(truncated, N more bytes)` so a reader knows the original +//! string was longer. +//! +//! [`worker::run_once`]: super::worker::run_once + +use std::sync::OnceLock; + +use regex::Regex; + +/// Upper bound on the byte length of a scrubbed string. Long enough to +/// keep the head of an anyhow `{:#}` chain (which usually includes the +/// most useful context) but short enough that one bad job can't flood +/// the log with megabytes of provider response body. +pub(crate) const MAX_LEN: usize = 1024; + +/// Scrub a free-form error / reason string for emission to logs. +/// Returns an owned `String` because every regex pass may rewrite the +/// input; callers are emitting this through `format!` / `log::*!` +/// anyway, so the allocation isn't on a hot path. +pub(crate) fn scrub_for_log(input: &str) -> String { + let mut out = input.to_owned(); + for (re, replacement) in patterns() { + // `Cow::into_owned` is cheap when the regex didn't match. + out = re.replace_all(&out, *replacement).into_owned(); + } + truncate(out) +} + +fn truncate(mut s: String) -> String { + if s.len() <= MAX_LEN { + return s; + } + // Round down to a char boundary to avoid splitting a multi-byte + // UTF-8 sequence — `truncate` itself panics on a non-boundary. + let mut cut = MAX_LEN; + while cut > 0 && !s.is_char_boundary(cut) { + cut -= 1; + } + let dropped = s.len() - cut; + s.truncate(cut); + s.push_str(&format!("…(truncated, {dropped} more bytes)")); + s +} + +fn patterns() -> &'static [(Regex, &'static str)] { + static PATTERNS: OnceLock> = OnceLock::new(); + PATTERNS.get_or_init(|| { + vec![ + // URL userinfo: capture the scheme + ://, then the userinfo + // up to '@'. Replace with ***@. Anchored on `://` to avoid + // eating bare `user@host` strings — those are handled by + // the email rule below. + ( + Regex::new(r"(?P[a-zA-Z][a-zA-Z0-9+.\-]*://)[^\s/@]+@").unwrap(), + "$scheme***@", + ), + // Bearer tokens (with or without an `Authorization:` prefix). + ( + Regex::new(r"(?i)bearer\s+[A-Za-z0-9._\-+/=]+").unwrap(), + "Bearer ***", + ), + // Provider-prefixed credentials with stable, well-known + // shapes. Listed individually so a future reader can see + // exactly which providers are covered. + ( + Regex::new(r"sk-[A-Za-z0-9_\-]{16,}").unwrap(), + "sk-***", + ), + ( + Regex::new(r"ghp_[A-Za-z0-9]{20,}").unwrap(), + "ghp_***", + ), + ( + Regex::new(r"ghs_[A-Za-z0-9]{20,}").unwrap(), + "ghs_***", + ), + ( + Regex::new(r"gho_[A-Za-z0-9]{20,}").unwrap(), + "gho_***", + ), + ( + Regex::new(r"xox[abprs]-[A-Za-z0-9\-]{8,}").unwrap(), + "xox-***", + ), + // Generic `key=value` assignments where the key name implies + // a secret. Matches `api_key`, `apiKey`, `api-key`, + // `password`, `passwd`, `pwd`, `token`, `secret`. Accepts + // a quoted, single-quoted, or bare value; bare values stop + // at the first whitespace / comma / closing bracket so we + // don't eat the rest of the message. + ( + Regex::new( + r#"(?i)(?Papi[_\-]?key|password|passwd|pwd|secret|token|auth)\s*[:=]\s*(?:"[^"]*"|'[^']*'|[^\s,}\)\]]+)"#, + ) + .unwrap(), + "$k=***", + ), + // Bare email addresses. Conservative pattern (no UTF-8 + // local parts, no quoted local parts) — sufficient to mask + // the common cases without eating unrelated `@` symbols. + ( + Regex::new(r"\b[A-Za-z0-9._%+\-]+@[A-Za-z0-9.\-]+\.[A-Za-z]{2,}\b").unwrap(), + "***@***", + ), + ] + }) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn passthrough_for_safe_string() { + let input = "rate_limited: provider returned 429, retry in 30s"; + assert_eq!(scrub_for_log(input), input); + } + + #[test] + fn masks_bearer_token() { + let s = scrub_for_log("Authorization: Bearer abc123.def-456_xyz"); + assert!(s.contains("Bearer ***"), "got {s:?}"); + assert!(!s.contains("abc123")); + } + + #[test] + fn masks_openai_key() { + let s = scrub_for_log("upstream returned 401: invalid sk-abcDEF1234567890ZZZZ key"); + assert!(s.contains("sk-***")); + assert!(!s.contains("sk-abcDEF1234567890ZZZZ")); + } + + #[test] + fn masks_github_token_variants() { + for raw in [ + "ghp_abcdefghij1234567890ABCD", + "ghs_abcdefghij1234567890ABCD", + "gho_abcdefghij1234567890ABCD", + ] { + let s = scrub_for_log(&format!("error: token {raw} rejected")); + assert!(s.contains("***"), "input={raw} out={s}"); + assert!(!s.contains(raw), "input={raw} out={s}"); + } + } + + #[test] + fn masks_slack_token() { + let s = scrub_for_log("posting to slack failed: xoxb-1234567890-abcdEFG"); + assert!(s.contains("xox-***")); + assert!(!s.contains("xoxb-1234567890")); + } + + #[test] + fn masks_generic_secret_assignments() { + let inputs = [ + ("api_key=hunter2 trailing", "api_key=***"), + ("password: hunter2", "password=***"), + ("Token = hunter2,more", "Token=***"), + ("apiKey=\"hunter2\"", "apiKey=***"), + ]; + for (raw, expect) in inputs { + let s = scrub_for_log(raw); + assert!(s.contains(expect), "input={raw:?} out={s:?}"); + assert!(!s.contains("hunter2"), "input={raw:?} out={s:?}"); + } + } + + #[test] + fn strips_url_userinfo() { + let s = scrub_for_log("connect failed https://alice:s3cret@db.internal/x"); + assert!(s.contains("https://***@db.internal/x"), "got {s:?}"); + assert!(!s.contains("alice")); + assert!(!s.contains("s3cret")); + } + + #[test] + fn masks_email() { + let s = scrub_for_log("user alice@example.com triggered job"); + assert!(s.contains("***@***"), "got {s:?}"); + assert!(!s.contains("alice")); + assert!(!s.contains("example.com")); + } + + #[test] + fn truncates_oversized_input() { + let big = "x".repeat(MAX_LEN * 2); + let s = scrub_for_log(&big); + assert!(s.len() < big.len()); + assert!(s.contains("(truncated,")); + } + + #[test] + fn truncate_handles_multibyte_boundary() { + // `é` is 2 bytes in UTF-8; build a string whose naïve cut at + // MAX_LEN would land mid-codepoint. + let mut big = "a".repeat(MAX_LEN - 1); + big.push('é'); + big.push_str(&"b".repeat(64)); + let s = scrub_for_log(&big); + // Must not panic; must produce valid UTF-8 (String guarantees). + assert!(s.contains("(truncated,")); + } + + #[test] + fn idempotent_on_already_scrubbed_string() { + let once = scrub_for_log("Bearer abcdef api_key=hunter2"); + let twice = scrub_for_log(&once); + assert_eq!(once, twice); + } +} diff --git a/src/openhuman/memory/tree/jobs/store.rs b/src/openhuman/memory/tree/jobs/store.rs index d733dea59..c6641d7c7 100644 --- a/src/openhuman/memory/tree/jobs/store.rs +++ b/src/openhuman/memory/tree/jobs/store.rs @@ -22,6 +22,7 @@ use rusqlite::{params, Connection, OptionalExtension, Transaction}; use uuid::Uuid; use crate::openhuman::config::Config; +use crate::openhuman::memory::tree::jobs::redact::scrub_for_log; use crate::openhuman::memory::tree::jobs::types::{Job, JobKind, JobStatus, NewJob}; use crate::openhuman::memory::tree::store::with_connection; @@ -202,10 +203,14 @@ pub fn mark_failed(config: &Config, job: &Job, error: &str) -> Result<()> { with_connection(config, |conn| { let now_ms = Utc::now().timestamp_millis(); + // `error` is free-form (anyhow chain or handler-supplied text) and + // may carry credential-shaped substrings; scrub before logging, + // but keep the original in the DB column for diagnostics. + let error_for_log = scrub_for_log(error); if attempts >= max_attempts { log::warn!( "[memory_tree::jobs] terminal failure id={job_id} \ - attempts={attempts}/{max_attempts} err={error}" + attempts={attempts}/{max_attempts} err={error_for_log}" ); let n = conn.execute( "UPDATE mem_tree_jobs @@ -229,7 +234,7 @@ pub fn mark_failed(config: &Config, job: &Job, error: &str) -> Result<()> { let next_at = now_ms.saturating_add(backoff); log::info!( "[memory_tree::jobs] retry id={job_id} attempt={attempts}/{max_attempts} \ - next_at_ms={next_at} err={error}" + next_at_ms={next_at} err={error_for_log}" ); let n = conn.execute( "UPDATE mem_tree_jobs @@ -531,6 +536,61 @@ mod tests { assert!(row.completed_at_ms.is_some()); } + /// `mark_failed` scrubs only the log emission, not the persisted + /// `last_error` column. A reader of `mem_tree_jobs` should still see + /// the full anyhow / handler-supplied chain so they can root-cause + /// the failure; the scrub is a defense-in-depth for the log sink only. + #[test] + fn mark_failed_persists_full_error_unredacted() { + let (_tmp, cfg) = test_config(); + let payload = AppendBufferPayload { + node: NodeRef::Leaf { + chunk_id: "c1".into(), + }, + target: AppendTarget::Source { + source_id: "slack:#x".into(), + }, + }; + let mut nj = NewJob::append_buffer(&payload).unwrap(); + nj.max_attempts = Some(1); + let id = enqueue(&cfg, &nj).unwrap().unwrap(); + + let claim = claim_next(&cfg, DEFAULT_LOCK_DURATION_MS).unwrap().unwrap(); + let raw = "upstream returned 401: Bearer abc123.def-456 token rejected"; + mark_failed(&cfg, &claim, raw).unwrap(); + + let row = get_job(&cfg, &id).unwrap().unwrap(); + assert_eq!(row.status, JobStatus::Failed); + // The persisted column keeps the full original — the scrub is + // applied at log emission only. + assert_eq!(row.last_error.as_deref(), Some(raw)); + } + + /// Same contract for `mark_deferred`: the log line is scrubbed in + /// `worker::run_once`, but the persisted `last_error` keeps the + /// full handler-supplied reason for diagnostics. + #[test] + fn mark_deferred_persists_full_reason_unredacted() { + let (_tmp, cfg) = test_config(); + let payload = AppendBufferPayload { + node: NodeRef::Leaf { + chunk_id: "c2".into(), + }, + target: AppendTarget::Source { + source_id: "slack:#y".into(), + }, + }; + let nj = NewJob::append_buffer(&payload).unwrap(); + let id = enqueue(&cfg, &nj).unwrap().unwrap(); + + let claim = claim_next(&cfg, DEFAULT_LOCK_DURATION_MS).unwrap().unwrap(); + let raw = "rate_limited by api_key=hunter2; retry after 30s"; + mark_deferred(&cfg, &claim, 0, raw).unwrap(); + + let row = get_job(&cfg, &id).unwrap().unwrap(); + assert_eq!(row.last_error.as_deref(), Some(raw)); + } + #[test] fn recover_stale_locks_resets_running_rows() { let (_tmp, cfg) = test_config(); diff --git a/src/openhuman/memory/tree/jobs/worker.rs b/src/openhuman/memory/tree/jobs/worker.rs index 79464f93b..da302ce4c 100644 --- a/src/openhuman/memory/tree/jobs/worker.rs +++ b/src/openhuman/memory/tree/jobs/worker.rs @@ -16,6 +16,7 @@ use tokio::sync::Notify; use crate::openhuman::config::Config; use crate::openhuman::memory::tree::jobs::handlers; +use crate::openhuman::memory::tree::jobs::redact::scrub_for_log; use crate::openhuman::memory::tree::jobs::store::{ claim_next, mark_deferred, mark_done, mark_failed, recover_stale_locks, DEFAULT_LOCK_DURATION_MS, @@ -164,26 +165,31 @@ pub async fn run_once(config: &Config) -> Result { // claim toward the failure-attempt budget. `mark_deferred` // reverts the bump applied by `claim_next` so the row's // attempts counter stays where it was before this claim. + // + // `reason` is handler-supplied free-form text and may + // include upstream provider responses; scrub for log + // emission while keeping the original in DB state. log::info!( "[memory_tree::jobs] deferred id={} kind={} until_ms={} reason={}", job.id, job.kind.as_str(), until_ms, - reason + scrub_for_log(&reason) ); mark_deferred(config, &job, until_ms, &reason)?; } Err(err) => { // Preserve the full anyhow cause chain in the persisted // last_error so a reader of mem_tree_jobs can see the root - // cause, not just the top-level message. Mirrors the {:#} - // log format used right above. + // cause, not just the top-level message. The log line gets + // the same chain after `scrub_for_log`, since anyhow chains + // commonly embed upstream HTTP bodies / auth headers. let message = format!("{err:#}"); log::warn!( - "[memory_tree::jobs] job failed id={} kind={} err={:#}", + "[memory_tree::jobs] job failed id={} kind={} err={}", job.id, job.kind.as_str(), - err + scrub_for_log(&message) ); mark_failed(config, &job, &message)?; }