use crate::cli::App;
use crate::errors::{err, ErrorCode};
use anyhow::{Context, Result};
use chrono::Utc;
use serde_json::{json, Value};
use std::io::Read;
use std::path::PathBuf;
use uuid::Uuid;
pub const EVENT_SCHEMA_VERSION: u32 = 1;
pub fn redact(text: &str, extra_values: &[String]) -> (String, Vec<String>) {
let mut out = text.to_string();
let mut applied: Vec<String> = Vec::new();
for (key, value) in std::env::vars() {
if value.len() >= 8 && out.contains(&value) {
out = out.replace(&value, "[REDACTED:env]");
applied.push(format!("env:{key}"));
}
}
for value in extra_values {
if !value.is_empty() && out.contains(value.as_str()) {
out = out.replace(value.as_str(), "[REDACTED:configured]");
applied.push("configured".into());
}
}
const PREFIXES: &[&str] = &[
"AKIA",
"ASIA",
"sk-",
"sk_live_",
"pk_live_",
"ghp_",
"gho_",
"github_pat_",
"xoxb-",
"xoxp-",
"glpat-",
"AIza",
"ya29.",
"eyJhbGciOi",
];
for prefix in PREFIXES {
let mut result = String::with_capacity(out.len());
let mut rest = out.as_str();
let mut hit = false;
while let Some(pos) = rest.find(prefix) {
let (before, after) = rest.split_at(pos);
result.push_str(before);
let token_len = after
.char_indices()
.take_while(|(_, c)| {
c.is_ascii_alphanumeric() || matches!(c, '-' | '_' | '.' | '+' | '/' | '=')
})
.count();
if token_len >= 12 {
result.push_str("[REDACTED:credential]");
hit = true;
rest = &after[token_len..];
} else {
result.push_str(&after[..token_len.max(prefix.len()).min(after.len())]);
rest = &after[token_len.max(prefix.len()).min(after.len())..];
}
}
result.push_str(rest);
out = result;
if hit {
applied.push(format!("credential:{prefix}"));
}
}
for marker in ["Bearer ", "bearer "] {
let mut search_from = 0usize;
while let Some(rel) = out[search_from..].find(marker) {
let start = search_from + rel + marker.len();
let end = out[start..]
.char_indices()
.take_while(|(_, c)| !c.is_whitespace() && *c != '"' && *c != '\'')
.map(|(i, c)| start + i + c.len_utf8())
.last()
.unwrap_or(start);
if end - start >= 8 && !out[start..end].starts_with("[REDACTED") {
out.replace_range(start..end, "[REDACTED:bearer]");
applied.push("bearer".into());
search_from = start + "[REDACTED:bearer]".len();
} else {
search_from = end.max(start + 1).min(out.len());
}
if search_from >= out.len() {
break;
}
}
}
applied.sort();
applied.dedup();
(out, applied)
}
pub fn ingest(app: &App, agent: &str, event: Option<&str>, session: Option<&str>) -> Result<()> {
if app.config.audit.level == "off" {
return Ok(());
}
let mut raw = Vec::new();
std::io::stdin()
.take(app.config.audit.max_event_bytes)
.read_to_end(&mut raw)?;
let payload: Value = serde_json::from_slice(&raw)
.unwrap_or_else(|_| json!({ "raw": String::from_utf8_lossy(&raw).to_string() }));
let session_id = session
.map(String::from)
.or_else(|| {
payload
.get("session_id")
.and_then(|v| v.as_str())
.map(String::from)
})
.or_else(|| {
payload
.get("conversation_id")
.and_then(|v| v.as_str())
.map(String::from)
})
.unwrap_or_else(|| format!("unattributed-{}", Utc::now().format("%Y%m%d")));
let event_type = event
.map(String::from)
.or_else(|| {
payload
.get("hook_event_name")
.and_then(|v| v.as_str())
.map(String::from)
})
.or_else(|| {
payload
.get("event")
.and_then(|v| v.as_str())
.map(String::from)
})
.unwrap_or_else(|| "unknown".into());
let payload = match app.config.audit.level.as_str() {
"metadata" => json!({}),
_ => payload,
};
let (payload_text, redactions) = if app.config.redaction.enabled {
redact(&payload.to_string(), &app.config.redaction.patterns)
} else {
(payload.to_string(), vec![])
};
let payload: Value = serde_json::from_str(&payload_text).unwrap_or(json!({}));
let event = json!({
"schema_version": EVENT_SCHEMA_VERSION,
"event_id": Uuid::now_v7().to_string(),
"session_id": session_id,
"writer_id": app.writer_id()?,
"agent": agent,
"event_type": event_type,
"timestamp": Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Millis, true),
"repo_root_hash": blake3::hash(app.repo.root.to_string_lossy().as_bytes()).to_hex()[..16].to_string(),
"branch": app.repo.branch(),
"head": app.repo.head_oid(),
"payload": payload,
"redactions": redactions,
});
let spool_dir = app.repo.state_dir().join("spool");
std::fs::create_dir_all(&spool_dir)?;
let safe_session: String = event["session_id"]
.as_str()
.unwrap_or("unknown")
.chars()
.map(|c| {
if c.is_ascii_alphanumeric() || c == '-' {
c
} else {
'_'
}
})
.collect();
let spool_path = spool_dir.join(format!("{safe_session}.ndjson"));
use std::io::Write as _;
let mut line = serde_json::to_vec(&event)?;
line.push(b'\n');
let mut file = std::fs::OpenOptions::new()
.read(true)
.create(true)
.append(true)
.open(&spool_path)
.map_err(|e| err(ErrorCode::AuditWriteFailed, format!("spool: {e}")))?;
file.lock()
.map_err(|e| err(ErrorCode::AuditWriteFailed, format!("spool lock: {e}")))?;
file.write_all(&line)
.map_err(|e| err(ErrorCode::AuditWriteFailed, e.to_string()))?;
Ok(())
}
fn spool_sessions(app: &App) -> Result<Vec<(String, PathBuf, usize)>> {
let spool_dir = app.repo.state_dir().join("spool");
let mut out = Vec::new();
if !spool_dir.exists() {
return Ok(out);
}
for entry in std::fs::read_dir(&spool_dir)? {
let entry = entry?;
let path = entry.path();
if path.extension().is_some_and(|e| e == "ndjson") {
let session = path
.file_stem()
.unwrap_or_default()
.to_string_lossy()
.to_string();
let count = std::fs::read_to_string(&path)
.map(|t| {
t.lines()
.filter(|l| serde_json::from_str::<Value>(l).is_ok())
.count()
})
.unwrap_or(0);
out.push((session, path, count));
}
}
out.sort();
Ok(out)
}
pub fn list(app: &App) -> Result<()> {
let sessions = spool_sessions(app)?;
let staging = app.repo.shared_dir().join("audit-staging");
let mut finalized: Vec<String> = Vec::new();
if staging.exists() {
for entry in std::fs::read_dir(&staging)? {
finalized.push(entry?.file_name().to_string_lossy().to_string());
}
finalized.sort();
}
if app.json {
let open: Vec<Value> = sessions
.iter()
.map(|(s, _, n)| json!({ "session": s, "events": n, "state": "open" }))
.collect();
println!("{}", json!({ "open": open, "finalized": finalized }));
} else {
for (s, _, n) in &sessions {
println!("open {s} ({n} events)");
}
for s in &finalized {
println!("finalized {s}");
}
if sessions.is_empty() && finalized.is_empty() {
println!("No audit sessions.");
}
}
Ok(())
}
pub fn show(app: &App, session: &str, events: bool) -> Result<()> {
let sessions = spool_sessions(app)?;
if let Some((_, path, count)) = sessions.iter().find(|(s, _, _)| s == session) {
let text = std::fs::read_to_string(path)?;
let parsed: Vec<Value> = text
.lines()
.filter_map(|l| serde_json::from_str(l).ok())
.collect();
let types: std::collections::BTreeMap<String, usize> =
parsed.iter().fold(Default::default(), |mut acc, e| {
*acc.entry(e["event_type"].as_str().unwrap_or("unknown").to_string())
.or_default() += 1;
acc
});
if events {
for e in &parsed {
println!("{e}");
}
} else if app.json {
println!(
"{}",
json!({ "session": session, "events": count, "event_types": types })
);
} else {
println!("session {session}: {count} events");
for (t, n) in types {
println!(" {t}: {n}");
}
println!("use --events to expand raw events");
}
return Ok(());
}
Err(err(
ErrorCode::InvalidRecord,
format!("unknown audit session '{session}'"),
))
}
pub fn finalize(app: &App, session: &str) -> Result<()> {
let sessions = spool_sessions(app)?;
let Some((_, spool_path, count)) = sessions.into_iter().find(|(s, _, _)| s == session) else {
return Err(err(
ErrorCode::InvalidRecord,
format!("unknown session '{session}'"),
));
};
let text = std::fs::read_to_string(&spool_path)?;
let (valid_lines, corrupt_lines): (Vec<&str>, Vec<&str>) = text
.lines()
.partition(|l| serde_json::from_str::<Value>(l).is_ok());
let ndjson = valid_lines.join("\n");
let compressed = zstd::encode_all(ndjson.as_bytes(), 9)
.map_err(|e| err(ErrorCode::AuditWriteFailed, format!("zstd: {e}")))?;
let now = Utc::now();
let summary = json!({
"session_id": session,
"finalized_at": now.to_rfc3339_opts(chrono::SecondsFormat::Secs, true),
"writer_id": app.writer_id()?,
"events": count,
"events_bytes_compressed": compressed.len(),
"schema_version": EVENT_SCHEMA_VERSION,
});
let staging = app.repo.shared_dir().join("audit-staging").join(session);
std::fs::create_dir_all(&staging)?;
std::fs::write(
staging.join("summary.json"),
serde_json::to_string_pretty(&summary)?,
)?;
std::fs::write(staging.join("events.jsonl.zst"), &compressed)?;
let mut ref_written = false;
if app.config.audit.write_git_refs {
match write_audit_ref(app, session, &now, &staging) {
Ok(()) => ref_written = true,
Err(e) => {
eprintln!("warning: audit ref update failed ({e}); bundle preserved for retry");
}
}
}
if !corrupt_lines.is_empty() {
let quarantine = spool_path.with_extension("ndjson.corrupt");
use std::io::Write as _;
let mut q = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(&quarantine)
.map_err(|e| err(ErrorCode::AuditWriteFailed, format!("quarantine: {e}")))?;
for line in &corrupt_lines {
writeln!(q, "{line}").map_err(|e| err(ErrorCode::AuditWriteFailed, e.to_string()))?;
}
eprintln!(
"warning: {} unparseable event line(s) quarantined to {}",
corrupt_lines.len(),
quarantine.display()
);
}
std::fs::remove_file(&spool_path).ok();
if app.json {
println!(
"{}",
json!({ "session": session, "events": count, "staging": staging.to_string_lossy(), "audit_ref_written": ref_written })
);
} else {
println!(
"finalized {session}: {count} events -> {}",
staging.display()
);
if ref_written {
println!("written to refs/memlay/audit/{}", app.writer_id()?);
}
}
Ok(())
}
fn write_audit_ref(
app: &App,
session: &str,
now: &chrono::DateTime<Utc>,
staging: &std::path::Path,
) -> Result<()> {
let writer = app.writer_id()?;
let ref_name = format!("refs/memlay/audit/{writer}");
let repo = &app.repo;
let hash_file = |path: &std::path::Path| -> Result<String> {
repo.run(&["hash-object", "-w", &path.to_string_lossy()])
};
let summary_oid = hash_file(&staging.join("summary.json"))?;
let events_oid = hash_file(&staging.join("events.jsonl.zst"))?;
let mktree = |entries: &[(String, String, String)]| -> Result<String> {
let input: String = entries
.iter()
.map(|(mode_type, oid, name)| format!("{mode_type} {oid}\t{name}\n"))
.collect();
let out = std::process::Command::new("git")
.args(["mktree"])
.current_dir(&repo.root)
.stdin(std::process::Stdio::piped())
.stdout(std::process::Stdio::piped())
.stderr(std::process::Stdio::piped())
.spawn()
.and_then(|mut c| {
use std::io::Write as _;
c.stdin.as_mut().unwrap().write_all(input.as_bytes())?;
c.wait_with_output()
})
.context("git mktree")?;
if !out.status.success() {
anyhow::bail!("mktree: {}", String::from_utf8_lossy(&out.stderr));
}
Ok(String::from_utf8_lossy(&out.stdout).trim().to_string())
};
let leaf = mktree(&[
("100644 blob".into(), events_oid, "events.jsonl.zst".into()),
("100644 blob".into(), summary_oid, "summary.json".into()),
])?;
let mut tree = leaf;
for name in [
session.to_string(),
now.format("%d").to_string(),
now.format("%m").to_string(),
now.format("%Y").to_string(),
"sessions".to_string(),
] {
tree = mktree(&[("040000 tree".into(), tree.clone(), name)])?;
}
for _attempt in 0..3 {
let old = repo
.run(&["rev-parse", "--verify", "--quiet", &ref_name])
.ok()
.filter(|s| !s.is_empty());
let msg = format!("memlay audit: session {session}");
let commit = match &old {
Some(parent) => repo.run(&["commit-tree", &tree, "-p", parent, "-m", &msg])?,
None => repo.run(&["commit-tree", &tree, "-m", &msg])?,
};
let result = match &old {
Some(parent) => repo.run(&["update-ref", &ref_name, &commit, parent]),
None => repo.run(&["update-ref", &ref_name, &commit, ""]),
};
if result.is_ok() {
return Ok(());
}
}
Err(err(
ErrorCode::AuditWriteFailed,
"compare-and-swap failed after retries",
))
}
pub fn sync_refs(app: &App) -> Result<()> {
let remote = app.config.audit.remote.clone();
let writer = app.writer_id()?;
let push = app.repo.run(&[
"push",
&remote,
&format!("refs/memlay/audit/{writer}:refs/memlay/audit/{writer}"),
]);
let fetch = app
.repo
.run(&["fetch", &remote, "refs/memlay/audit/*:refs/memlay/audit/*"]);
match (&push, &fetch) {
(Ok(_), Ok(_)) => {
if !app.json {
println!("audit refs synchronized with {remote}");
}
Ok(())
}
_ => Err(err(
ErrorCode::AuditWriteFailed,
format!(
"push: {}; fetch: {}",
push.err()
.map(|e| e.to_string())
.unwrap_or_else(|| "ok".into()),
fetch
.err()
.map(|e| e.to_string())
.unwrap_or_else(|| "ok".into())
),
)),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn redacts_common_credentials() {
let (out, applied) = redact(
"key AKIAIOSFODNN7EXAMPLE and token ghp_abcdefghijklmnopqrstuvwxyz123456 plus Bearer abcdef123456789",
&[],
);
assert!(!out.contains("AKIAIOSFODNN7EXAMPLE"), "{out}");
assert!(!out.contains("ghp_abcdefghijklmnop"), "{out}");
assert!(!out.contains("abcdef123456789"), "{out}");
assert!(!applied.is_empty());
}
#[test]
fn redacts_configured_patterns_and_never_returns_them() {
let secret = "super-secret-value-123".to_string();
let (out, _) = redact(&format!("data {secret} end"), std::slice::from_ref(&secret));
assert!(!out.contains(&secret));
}
#[test]
fn short_prefix_lookalikes_survive() {
let (out, _) = redact("skiing sk-i is fun", &[]);
assert!(out.contains("skiing"), "{out}");
}
}