From d290a29a5517dcde3af3e7cbd0c059eef3a47ec0 Mon Sep 17 00:00:00 2001 From: Chris Lee Date: Sat, 27 Jun 2026 11:09:12 -0600 Subject: [PATCH] Add thread mbox/maildir export Export reconstructed threads as mbox (one file, threaded by Message-ID / In-Reply-To) or Maildir (one file per message), via mailbox.rs. Adds --format and --body flags: body modes are plain, html (multipart/alternative, default), and html-only. Each trie node becomes one email; assistant tool calls/results render as attachments. Co-Authored-By: GLM-5.2 --- AGENTS.md | 169 +++++++++ Cargo.lock | 2 +- README.md | 34 +- src/commands.rs | 44 +++ src/mailbox.rs | 883 +++++++++++++++++++++++++++++++++++++++++++ src/main.rs | 1 + tests/integration.rs | 424 +++++++++++++++++++++ 7 files changed, 1542 insertions(+), 15 deletions(-) create mode 100644 AGENTS.md create mode 100644 src/mailbox.rs diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..8477bc2 --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,169 @@ +# AGENTS.md + +Guidance for coding agents working on czsplicer. Read this before editing. + +## What this is + +`czsplicer` is a Rust CLI for inspecting, extracting, editing/redacting, and +repacking `.cbor.zstd` log streams — specifically the concatenated +CBOR-over-zstd captures exported by Tailscale Aperture, but it works on any +such stream. Each file is a **concatenated stream of independent CBOR map +records** (one record per captured request/response), *not* a single CBOR +array. Keep that framing in mind: every read path is record-by-record. + +## Build, test, lint + +- Rust 1.80+, edition 2021. Single binary, no workspace. +- `cargo build` / `cargo build --release` (binary at `target/release/czsplicer`). +- `cargo test` — 83 integration tests in `tests/integration.rs`, all synthetic. +- `cargo fmt --check` is enforced. The pre-commit hook (`hooks/pre-commit`, + enable with `git config core.hooksPath hooks`) runs `fmt --check` + `cargo + test` when `.rs`/`.toml`/`tests/` files are staged. +- `prod/`, `target/`, `.maki/` are git-ignored. `prod/` holds real (large) + export data and must never be committed. + +Note: the README is slightly stale — it says "65 integration tests" (now 83) +and references `cargo test -- --ignored` real-data round-trip tests that no +longer exist in the source tree. + +## Repository layout + +``` +src/ + main.rs clap Cli/Cmd dispatch + expand() (dir -> sorted *.cbor.zstd) + commands.rs one *Args (clap) struct + cmd_* fn per subcommand. ~2/3 of code. + filter.rs Filter + FilterArgs, shared by all selection commands + format.rs CBOR<->JSON bridge, RecordStream, ZstdPacker, redact/search, field accessors + thread.rs conversation-thread reconstruction (trie over message-content hashes) + RecordMeta + render.rs shared helpers for the HTML renderers (escape_html, truncate, sender_color, best_record_id, ...) + clip_chars + markdown.rs minimal safe Markdown->HTML subset for the built-in renderer + builtin.rs built-in long-form HTML renderer (wide column, markdown, status/tool chips) + theme.rs Adium .AdiumMessageStyle loader + renderer (--theme, optional) + mailbox.rs mbox/Maildir export (RFC822 + threading via Message-ID/In-Reply-To) + builtin.css stylesheet for builtin.rs (embedded via include_str!) +tests/ + integration.rs end-to-end tests via assert_cmd + common/mod.rs Fixture builder (SOURCE_NDJSON / RICH_NDJSON truth sets) + fixtures/Spike.AdiumMessageStyle/ MIT test-fixture Adium theme +vendor/ highlight.js (BSD-3-Clause) + CSS themes, embedded via include_str! +``` + +No `mod.rs` under `src/`; `main.rs` declares +`mod builtin; mod commands; mod filter; mod format; mod mailbox; mod markdown; mod render; mod theme; mod thread;`. + +## Architecture & data flow + +Every selection command (`ls`, `extract`, `grep`, `edit`, `stats`, `merge`, +`split`, `thread`) follows the same shape: + +1. `expand()` (main.rs) turns directory args into their sorted `*.cbor.zstd` + contents. +2. `FilterArgs::build()` (filter.rs) compiles CLI flags into a `Filter`. +3. `RecordStream::open(path)` (format.rs) decodes zstd once and yields CBOR + records lazily via `ciborium`'s streaming decoder. +4. Per record: `Filter::matches(rec)` gates the work, then the command-specific + transform runs. +5. Output is written streaming (NDJSON / re-compressed CBOR). + +`verify` and `repack` don't filter; `info` summarizes without streaming +transforms. `thread` builds an in-memory trie over the whole filtered set (one +synthetic Node per distinct message — see invariant 5). +transforms. + +The clap `*Args` structs live in `commands.rs` and are passed directly into +`cmd_*` from `main.rs` — there is no separate parallel struct layer. + +## Critical invariants (easy to break, not obvious) + +### 1. CBOR bytes vs text — the JSON bridge sentinels +`capture.rawRequestBody` / `rawResponseBody` are CBOR **bytes** (raw HTTP +bodies); the parallel `requestBody`/`responseBody` are text. `cbor_to_json` / +`json_to_cbor` (format.rs) preserve the distinction via two sentinels: +- `BYTES_KEY = "__cbor_bytes_b64"` → bytes encode as `{"__cbor_bytes_b64": ""}` +- `TAG_KEY = "__cbor_tag"` → CBOR tags encode as `{"__cbor_tag": [, ]}` + +Do not collapse bytes into strings or you'll corrupt the round-trip and silently +change record types on repack. + +### 2. Redaction vs search deliberately disagree on invalid-UTF-8 bytes +`format::redact_strings` (used by `edit --redact`) and +`format::search_value_strings` (used by `grep`) both visit `Text` **and** +`Bytes`, but handle invalid-UTF-8 bytes differently: + +- `search_value_strings` decodes bytes **lossily** (`String::from_utf8_lossy`). +- `redact_strings` scrubs bytes **only when valid UTF-8** (`std::str::from_utf8`); + invalid bytes are left byte-for-byte intact. + +**Consequence:** grep can surface an ASCII secret embedded in a partially-invalid +byte body that `edit --redact` will NOT scrub. This is a deliberate, documented +tradeoff (see the doc comment on `redact_strings`): scrubbing a lossy decode +would rewrite the original bytes and corrupt binary payloads. Do not "fix" this +gap by redacting the lossy decode — it would break binary preservation. For +well-formed text bodies (the realistic case) the two agree exactly. + +Raw HTTP bodies live in `capture.rawRequestBody` / `capture.rawResponseBody`. +`edit --redact` targets `capture` by default; `--all-strings` widens to the whole +record. + +### 3. Streaming is load-bearing +A single 3.3 MB export decompresses to ~1.1 GB / 2000 records. All read paths +stream record-by-record. The **only** exception is `extract --array`, which +buffers the full result — keep it that way (a single, documented exception). +Don't introduce new full-file buffering. + +### 4. Float precision +Lossless floats depend on `serde_json`'s `float_roundtrip` feature (enabled in +`Cargo.toml`) so `estimated_cost.dollars` survives the round-trip unchanged. +The one theoretical hole: CBOR negative integers below `i64::MIN` (down to +−2⁶⁴) fall back to f64 in `cbor_to_json` — never occurs in capture records. + +### 5. Conversation-thread reconstruction (thread.rs) +The `thread` command reconstructs conversation branches from request message +histories. Each record's `capture.requestBody.messages` echoes its **full parent +path**, so the tree is a **trie over blake3 content-hashes of normalized +messages** — NOT grouped by `session_id` (in Aperture captures every request +gets a unique session_id; session_id is just metadata). A depth-0 node (usually +the system prompt) is a conversation root; a node with >1 child is a branch +point (the user went back and took a different path). + +Content normalization is load-bearing: a bare-string message content and the +equivalent `[{"type":"text","text":s}]` block form MUST hash identically, or a +string-user-message request and its block-form continuation look like two +separate roots. `msg_info` normalizes strings to block form before hashing. +Assistant turns are reconstructed from the **next** request's echoed messages — +no response-body parsing is needed for structure. Verified on real prod data: +30/100 files contain genuine branches (fan-outs up to 19, depths up to 748). + +## Testing conventions + +- Tests are end-to-end via `assert_cmd`, driving the compiled binary. +- `Fixture::new()` (tests/common/mod.rs) writes `SOURCE_NDJSON`, then **builds + the `.cbor.zstd` fixture via the tool's own `repack`** — this exercises the + full CBOR↔JSON bridge and is intentional. A separate `RICH_NDJSON` + (5 records) covers merge/split/session grouping. +- `read_ndjson` parses with serde_json (order-preserving); prefer semantic + equality (`serde_json::Value` compare) over substring matching where possible. +- **Bytes-body redaction tests must base64-decode the body** and assert on the + decoded bytes. A secret riding inside `__cbor_bytes_b64` is never plain text + in extracted JSON, so `!out.contains(secret)` passes even on unfixed code + (see `edit_redact_scrubs_byte_bodies`). Its counterpart + `edit_redact_leaves_binary_bytes_untouched` pins the binary-preservation + contract and was verified to fail under a lossy-decode-rewrite regression. + +## Style notes + +- `anyhow::Result` everywhere; add `context`/descriptive errors on I/O paths. +- Command output: NDJSON by default (streaming), `--json` for structured, + `--array`/`--pretty` where a single buffered object is acceptable. +- zstd write level defaults to 9; `--level` exposes it. + +## Commit conventions + +- Every commit message ends with a `Co-Authored-By:` trailer crediting the + AI model that authored the change. Format: + ``` + Co-Authored-By: <@> + ``` + e.g. `Co-Authored-By: GLM-5.2 `. Fill in the actual model + name and handle for the model you are running as; do not leave a placeholder. + Append (don't replace) when a commit is co-authored by multiple models. diff --git a/Cargo.lock b/Cargo.lock index aa047bd..25e40d1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -251,7 +251,7 @@ checksum = "460fbee9c2c2f33933d720630a6a0bac33ba7053db5344fac858d4b8952d77d5" [[package]] name = "czsplicer" -version = "0.2.0" +version = "0.3.0" dependencies = [ "anyhow", "assert_cmd", diff --git a/README.md b/README.md index d4b225f..42b2363 100644 --- a/README.md +++ b/README.md @@ -40,7 +40,7 @@ Commands: merge Merge many `.cbor.zstd` files into one (CBOR -> CBOR, streaming) split Split one stream into per-group `.cbor.zstd` files (by day/session/model/path) stats Aggregate stats: tokens, cost, durations, by-model / by-path - thread Reconstruct conversation threads (branching included); export as JSON/HTML + thread Reconstruct conversation threads (branching included); export as JSON/HTML/MBOX/Maildir ``` Run `czsplicer --help` for full flags. @@ -115,9 +115,8 @@ czsplicer merge prod/ -o all.cbor.zstd # Split into per-day files czsplicer split all.cbor.zstd --by day --out-dir days/ -# Split into per-session files. Aperture gives every request a unique -# session_id, so --by session groups by conversation root (first message -# hash); --min-records defaults to 2 to skip single-request throwaways. +# Split into per-session files (session_id is auto-populated by Aperture; +# --min-records defaults to 2 to skip single-request throwaways) czsplicer split all.cbor.zstd --by session --out-dir sessions/ # Also: --by model, --by provider, --by path. Use --json for a manifest of the output files. @@ -125,10 +124,11 @@ czsplicer split all.cbor.zstd --by session --out-dir sessions/ ### Threading -Reconstruct conversation branches from each request's echoed message history. -Branch points (where the user went back and took a different path) are recovered -automatically; the trie keys on normalized message-content hashes, so a string -user message and its block-form continuation collapse to one node. +Reconstruct conversation branches from each request's echoed message history, +then render or export. Branch points (where the user went back and took a +different path) are recovered automatically; the trie keys on normalized +message-content hashes, so a string user message and its block-form +continuation collapse to one node. ```sh # Default: JSON forest (roots, nodes, record_ids, tool_events) to stdout. @@ -140,13 +140,20 @@ czsplicer thread prod/ --html --dark -o threads.html # Render through an Adium .AdiumMessageStyle bundle (optional, --variant Dark). czsplicer thread prod/ --theme Spike.AdiumMessageStyle -o threads.html -# Redact secrets in the reconstructed output (same presets as `edit`). +# Export as mbox (threaded by Message-ID / In-Reply-To) for a mail client. +czsplicer thread prod/ --format mbox -o threads.mbox + +# Maildir (one file per message) with plain-text bodies instead of HTML. +czsplicer thread prod/ --format maildir --body plain -o maildir/ + +# Redact secrets in the rendered output (same presets as `edit`). czsplicer thread prod/ --html --redact-preset all -o threads.html ``` -Formats: `json` (default) and `html` (built-in long-form renderer, or an -Adium `.AdiumMessageStyle` bundle via `--theme`). Redaction runs on message -bodies and tool text *before* rendering, so secrets never reach the output. +Formats: `json` (default), `html` (built-in), `mbox`, `maildir`. `--body` +controls mbox/maildir body rendering: `plain`, `html` (multipart/alternative, +default), `html-only`. Redaction runs on message bodies and tool text *before* +rendering, so secrets never reach the output file. ### Integrity check @@ -196,8 +203,7 @@ Directory arguments are expanded to their sorted `*.cbor.zstd` contents, so ## Development ```sh -cargo test # 114 integration tests (synthetic fixtures) -cargo test -- --ignored # + lossless round-trip over real prod/ exports +cargo test # 134 integration tests (synthetic fixtures) ``` The repository includes a pre-commit hook (`hooks/pre-commit`) that runs diff --git a/src/commands.rs b/src/commands.rs index 3bf36f9..35422ad 100644 --- a/src/commands.rs +++ b/src/commands.rs @@ -1,6 +1,7 @@ use crate::builtin; use crate::filter::{Filter, FilterArgs}; use crate::format::{self, RecordStream}; +use crate::mailbox; use crate::theme; use crate::thread::{conversation_root, ThreadBuilder}; use anyhow::{anyhow, Result}; @@ -1323,6 +1324,12 @@ pub struct ThreadArgs { /// Write JSON/HTML/MBOX/Maildir to this path instead of stdout (`-` for stdout). #[arg(short, long)] pub output: Option, + /// Output format: json (default), html (built-in renderer), mbox, maildir. + #[arg(long, value_name = "FMT", default_value = "json")] + pub format: String, + /// Body rendering for mbox/maildir: plain, html (multipart/alternative), html-only. + #[arg(long, value_name = "MODE", default_value = "html")] + pub body: String, /// Emit self-contained HTML using the built-in long-form renderer. #[arg(long)] pub html: bool, @@ -1406,6 +1413,43 @@ pub fn cmd_thread(args: &ThreadArgs) -> Result<()> { return Ok(()); } + // MBOX / Maildir export: each trie node becomes one email, threaded by + // Message-ID / In-Reply-To. Selected by --format mbox|maildir. + if args.format == "mbox" || args.format == "maildir" { + let mode = mailbox::BodyMode::parse(&args.body)?; + if args.format == "mbox" { + // mbox may write to stdout (`-` or no -o) or a file. + let n = match args.output.as_ref() { + Some(p) if p.as_path() != Path::new("-") => mailbox::write_mbox(&j, mode, p)?, + _ => { + let mut out = std::io::stdout(); + mailbox::write_mbox_to(&mut out, &j, mode)? + } + }; + eprintln!( + "{} record(s) ({} with messages) -> {} message(s) [{} body] (mbox)", + total, with_messages, n, args.body, + ); + return Ok(()); + } + // maildir requires a real directory path. + let out = args + .output + .as_ref() + .filter(|p| p.as_path() != Path::new("-")) + .ok_or_else(|| anyhow!("--format maildir requires -o DIR"))?; + let n = mailbox::write_maildir(&j, mode, out)?; + eprintln!( + "{} record(s) ({} with messages) -> {} message(s) [{} body] -> {} (maildir)", + total, + with_messages, + n, + args.body, + out.display() + ); + return Ok(()); + } + let pretty = serde_json::to_string_pretty(&j)?; let bytes = pretty.into_bytes(); write_output(args.output.as_ref(), &bytes, "json")?; diff --git a/src/mailbox.rs b/src/mailbox.rs new file mode 100644 index 0000000..c35a16b --- /dev/null +++ b/src/mailbox.rs @@ -0,0 +1,883 @@ +//! Email-format emitters (MBOX + Maildir) for the conversation-thread exporter. +//! +//! Each trie node becomes one email. Threading is reconstructed by RFC 5256: +//! every node has a unique `Message-ID`, and `In-Reply-To`/`References` point +//! at the parent's Message-ID. A branch point (node with N children) is just +//! N replies to the same parent — exactly what email clients render as a +//! nested, collapsible thread. +//! +//! One email per *node* (not per record_id): retries/edits that share a path +//! collapse to a single node; their record_ids are preserved as custom headers +//! and the body. Keeping one message per node is what makes threading clean. +//! +//! MBOX uses the mboxrd convention (body lines starting with `From ` are +//! `>`-escaped). Maildir writes cur/new/tmp with one file per message. +//! Bodies are MIME multipart/alternative (text/plain + text/html) when the +//! `--body html` mode is selected, or text/plain only under `--body plain`, +//! or text/html only under `--body html-only`. + +use crate::markdown; +use anyhow::{anyhow, Result}; +use serde_json::Value as Json; +use std::fs; +use std::io::Write; +use std::path::Path; + +/// Body rendering mode selected at export time. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum BodyMode { + /// `text/plain` only — raw message text. + Plain, + /// `multipart/alternative` with `text/plain` + `text/html`. + Html, + /// `text/html` only — rendered HTML, no plain fallback. + HtmlOnly, +} + +impl BodyMode { + pub fn parse(s: &str) -> Result { + Ok(match s { + "plain" => Self::Plain, + "html" => Self::Html, + "html-only" => Self::HtmlOnly, + other => return Err(anyhow!("invalid --body {other:?} (plain|html|html-only)")), + }) + } +} + +/// Emit the forest as a single mbox file at `out`. +pub fn write_mbox(forest: &Json, mode: BodyMode, out: &Path) -> Result { + let mut f = fs::File::create(out)?; + let n = write_mbox_to(&mut f, forest, mode)?; + f.flush()?; + Ok(n) +} + +/// Emit the forest as mbox to any writer (file or stdout). +pub fn write_mbox_to(w: &mut W, forest: &Json, mode: BodyMode) -> Result { + let n = emit_forest(forest, mode, |msg| { + // Compute the exact body bytes that will be written (with mboxrd + // escaping and per-line newlines) so Content-Length is accurate. + // mutt uses Content-Length to skip to the next message, which is + // critical for multipart bodies that may contain "From " lines. + let body = msg.body(); + let mut body_bytes = Vec::new(); + for line in body.lines() { + if line.starts_with("From ") { + body_bytes.push(b'>'); + } + body_bytes.extend_from_slice(line.as_bytes()); + body_bytes.push(b'\n'); + } + + w.write_all(b"From ")?; + w.write_all(msg.envelope_line().as_bytes())?; + w.write_all(b"\n")?; + // Write headers with Content-Length injected before the final + // Content-Type line's trailing newline. + let headers = msg.headers(); + // Find the position after the last header line content. + let header_end = headers.rfind('\n').map(|i| i + 1).unwrap_or(headers.len()); + w.write_all(&headers.as_bytes()[..header_end])?; + writeln!(w, "Content-Length: {}", body_bytes.len())?; + w.write_all(&headers.as_bytes()[header_end..])?; + w.write_all(b"\n")?; + // Write the pre-computed body bytes. + w.write_all(&body_bytes)?; + w.write_all(b"\n")?; + Ok(()) + })?; + Ok(n) +} + +/// Emit the forest as a Maildir at `out` (creates cur/new/tmp). +pub fn write_maildir(forest: &Json, mode: BodyMode, out: &Path) -> Result { + fs::create_dir_all(out.join("cur"))?; + fs::create_dir_all(out.join("new"))?; + fs::create_dir_all(out.join("tmp"))?; + let mut seq: u64 = 0; + let pid = std::process::id(); + let n = emit_forest(forest, mode, |msg| { + seq += 1; + // Classic Maildir unique name: time.pid.seq.host + let now = std::time::SystemTime::now() + .duration_since(std::time::UNIX_EPOCH) + .map(|d| d.as_secs()) + .unwrap_or(0); + let host = hostname(); + let name = format!("{now}.{pid}.{seq}.{host}"); + // Write atomically via tmp then rename. + let tmp = out.join("tmp").join(&name); + let new = out.join("new").join(&name); + let mut w = fs::File::create(&tmp)?; + w.write_all(msg.headers().as_bytes())?; + w.write_all(b"\n")?; + w.write_all(msg.body().as_bytes())?; + w.flush()?; + fs::rename(&tmp, &new)?; + Ok(()) + })?; + Ok(n) +} + +/// Walk the forest and call `emit` once per node, threading each node to its +/// parent via Message-ID/In-Reply-To. +fn emit_forest Result<()>>( + forest: &Json, + mode: BodyMode, + mut emit: E, +) -> Result { + let records = forest.get("records").and_then(|v| v.as_object()); + let mut count = 0usize; + let empty: Vec = Vec::new(); + let trees: &Vec = forest + .get("trees") + .and_then(|t| t.as_array()) + .unwrap_or(&empty); + for (thread_idx, root) in trees.iter().enumerate() { + // Parent Message-ID for depth-0 is None (these are thread roots). + walk_node( + root, None, None, thread_idx, mode, records, &mut emit, &mut count, + )?; + } + Ok(count) +} + +/// Walk the forest and call `emit` once per *run* of consecutive same-role +/// nodes, threading each run to its parent via Message-ID/In-Reply-To. +/// +/// A run is a linear chain: the entry node plus any following single child +/// that shares its role. The chain stops at a branch point (>1 child) or a +/// role change — those become children of the collapsed email. This collapses +/// e.g. `[system, system]` or `[assistant(tool_use), assistant(tool_use)]` +/// sequences into one email while keeping branch points as separate replies. +fn walk_node Result<()>>( + node: &Json, + parent_msgid: Option, + parent_hash: Option<&str>, + thread_idx: usize, + mode: BodyMode, + records: Option<&serde_json::Map>, + emit: &mut E, + count: &mut usize, +) -> Result<()> { + // Collect the run: node + same-role single-child descendants. + let mut run: Vec<&Json> = Vec::new(); + let mut cur = node; + loop { + run.push(cur); + let kids = cur.get("children").and_then(|c| c.as_array()); + let only_child = kids + .and_then(|k| k.iter().next()) + .filter(|_| kids.map(|k| k.len() == 1).unwrap_or(false)); + let next = match only_child { + Some(c) => c, + None => break, // branch point or leaf: end the run here + }; + let same_role = next + .get("role") + .and_then(|v| v.as_str()) + .zip(cur.get("role").and_then(|v| v.as_str())) + .map(|(a, b)| a == b) + .unwrap_or(false); + if !same_role { + break; + } + cur = next; + } + + let msgid = message_id(run[0], parent_hash, thread_idx); + let email = Email::build(&run, &msgid, parent_msgid.as_deref(), mode, records)?; + emit(email)?; + *count += 1; + + // Recurse into the LAST node of the run's children (the run's exit), + // replying to this run's Message-ID. The child's parent hash is the run's + // entry-node hash, so distinct parents disambiguate Message-IDs. + let exit = run.last().expect("run is non-empty"); + let entry_hash = run[0].get("hash").and_then(|v| v.as_str()).unwrap_or(""); + if let Some(kids) = exit.get("children").and_then(|c| c.as_array()) { + for child in kids { + walk_node( + child, + Some(msgid.clone()), + Some(entry_hash), + thread_idx, + mode, + records, + emit, + count, + )?; + } + } + Ok(()) +} + +/// Stable, unique Message-ID per node position: ``, +/// disambiguated by the parent node's content hash. Two nodes that share the +/// same `(hash, depth, thread_idx)` but are reached via different parents (a +/// shared subtree fragment) get distinct Message-IDs, avoiding RFC 5322 +/// duplicates that would otherwise break threading. Roots use the sentinel +/// `root` in place of a parent hash. +fn message_id(node: &Json, parent_hash: Option<&str>, thread_idx: usize) -> String { + let hash = node.get("hash").and_then(|v| v.as_str()).unwrap_or("x"); + let depth = node.get("depth").and_then(|v| v.as_i64()).unwrap_or(0); + let parent = parent_hash.unwrap_or("root"); + format!("<{hash}-{depth}-{thread_idx}-{parent}@czsplicer>") +} + +struct Email { + headers: String, + body: String, + env_line: String, +} + +impl Email { + fn envelope_line(&self) -> &str { + &self.env_line + } + fn headers(&self) -> &str { + &self.headers + } + fn body(&self) -> &str { + &self.body + } + + /// Build one email covering a *run* of consecutive same-role nodes. + /// + /// `nodes` is a non-empty slice: the first node is the run's entry point, + /// and any following nodes are same-role single children that were folded + /// in by `walk_node` (a linear chain ending before a branch point or a + /// role change). Content and `record_ids` are aggregated across the run; + /// metadata (timestamp/model/status/subject) comes from the first node's + /// introducer record. + fn build( + nodes: &[&Json], + msgid: &str, + parent_msgid: Option<&str>, + mode: BodyMode, + records: Option<&serde_json::Map>, + ) -> Result { + let first = nodes[0]; + let role = first.get("role").and_then(|v| v.as_str()).unwrap_or(""); + let depth = first.get("depth").and_then(|v| v.as_i64()).unwrap_or(0); + let intro_rid = first.get("intro_rid").and_then(|v| v.as_i64()).unwrap_or(0); + + // Aggregate content across the run, dropping empty/preview-only nodes. + let mut content_parts: Vec = Vec::new(); + let mut preview = String::new(); + let mut record_ids: Vec = Vec::new(); + for n in nodes { + let c = n.get("content").and_then(|v| v.as_str()).unwrap_or(""); + if !c.is_empty() { + content_parts.push(c.to_string()); + } + if preview.is_empty() { + preview = n + .get("preview") + .and_then(|v| v.as_str()) + .unwrap_or("") + .to_string(); + } + if let Some(a) = n.get("record_ids").and_then(|v| v.as_array()) { + record_ids.extend(a.iter().filter_map(|x| x.as_i64())); + } + } + let content = content_parts.join("\n\n"); + + // Resolve per-record metadata from the first node's introducer record. + let meta = records.and_then(|m| m.get(&intro_rid.to_string())); + let status = meta + .and_then(|m| m.get("status_code")) + .and_then(|v| v.as_i64()) + .unwrap_or(200); + let ts = meta + .and_then(|m| m.get("timestamp")) + .and_then(|v| v.as_str()) + .unwrap_or(""); + let model = meta + .and_then(|m| m.get("model")) + .and_then(|v| v.as_str()) + .unwrap_or(""); + let login_name = meta + .and_then(|m| m.get("login_name")) + .and_then(|v| v.as_str()) + .unwrap_or(""); + // Sum tool call/result counts across the whole run. + let mut tool_calls = 0i64; + let mut tool_results = 0i64; + for n in nodes { + let rid = n.get("intro_rid").and_then(|v| v.as_i64()).unwrap_or(0); + if let Some(m) = records.and_then(|rmap| rmap.get(&rid.to_string())) { + tool_calls += m.get("tool_calls").and_then(|v| v.as_i64()).unwrap_or(0); + tool_results += m.get("tool_results").and_then(|v| v.as_i64()).unwrap_or(0); + } + } + + // The Date: *header* is RFC 2822 (correct for a mail header). The + // From_ *postmark* (envelope_line) must use ctime/asctime format + // instead: mutt's strict is_from() parser only accepts the ctime + // timestamp on the postmark line, not the RFC 2822 form. Emitting + // RFC 2822 here makes mutt reject every postmark and report + // "[Msgs:0]" on an otherwise-valid mbox. + let date = rfc2822_date(ts); + let env_date = asctime_date(ts); + let from = sender_for_role(role, model, login_name); + let subject = subject_for_node(first, role); + + let mut headers = String::new(); + use std::fmt::Write; + let _ = writeln!(headers, "Message-ID: {msgid}"); + if let Some(p) = parent_msgid { + let _ = writeln!(headers, "In-Reply-To: {p}"); + let _ = writeln!(headers, "References: {p}"); + } + let _ = writeln!(headers, "From: {from}"); + let _ = writeln!(headers, "Subject: {subject}"); + let _ = writeln!(headers, "Date: {date}"); + let _ = writeln!(headers, "X-Czsplicer-Role: {role}"); + // Depth is the run's starting depth; for a collapsed run we also note + // the span so the original structure is recoverable. + if nodes.len() > 1 { + let last_depth = nodes + .last() + .and_then(|n| n.get("depth")) + .and_then(|v| v.as_i64()) + .unwrap_or(depth); + let _ = writeln!(headers, "X-Czsplicer-Depth: {depth}-{last_depth}"); + } else { + let _ = writeln!(headers, "X-Czsplicer-Depth: {depth}"); + } + let _ = writeln!(headers, "X-Czsplicer-Status: {status}"); + if !model.is_empty() { + let _ = writeln!(headers, "X-Czsplicer-Model: {model}"); + } + if !record_ids.is_empty() { + // Only emit the record count, not the full list — 320+ IDs on one + // line would exceed RFC 5322's 998-char header limit and break + // conformant mail parsers (mutt, etc.). + let _ = writeln!(headers, "X-Czsplicer-Record-Count: {}", record_ids.len()); + } + if tool_calls > 0 { + let _ = writeln!(headers, "X-Czsplicer-Tool-Calls: {tool_calls}"); + } + if tool_results > 0 { + let _ = writeln!(headers, "X-Czsplicer-Tool-Results: {tool_results}"); + } + + let body = body_mime(mode, &content, &preview); + + // For assistant turns, extract tool calls and results as attachments. + // Pairing is re-derived per node from intro_rid's call events paired + // with the next record's result events (see tool_call_attachments), so + // it stays correct when consecutive same-role nodes are collapsed. + let tool_attachments = if role == "assistant" { + tool_call_attachments(nodes, records) + } else { + Vec::new() + }; + + let (ct, payload) = if tool_attachments.is_empty() { + (body.content_type_header, body.payload) + } else { + wrap_mixed(&body, &tool_attachments) + }; + let _ = writeln!(headers, "MIME-Version: 1.0"); + let _ = writeln!(headers, "Content-Type: {ct}"); + + Ok(Email { + headers, + body: payload, + env_line: format!("czsplicer@localhost {env_date}"), + }) + } +} + +/// Sender display string per role. +/// - user: uses the Tailscale login_name as the email address (e.g. `clee@github`) +/// - assistant: uses `@` (e.g. `glm-5.2@ollama`) +/// - system/other: falls back to role-based addresses +fn sender_for_role(role: &str, model: &str, login_name: &str) -> String { + match role { + "user" => { + if login_name.is_empty() { + "User ".to_string() + } else { + // login_name from Tailscale is already an email-style identity + // (e.g. "clee@github"); use it directly. + format!("User <{}>", sanitize_header(login_name)) + } + } + "assistant" => { + if model.is_empty() { + "Assistant ".to_string() + } else { + let (local, domain) = normalize_model_email(model); + format!( + "Assistant ({}) <{}@{}>", + sanitize_header(model), + sanitize_header(&local), + sanitize_header(&domain) + ) + } + } + "system" => "System ".to_string(), + _ => format!("{} ", sanitize_header(role)), + } +} + +/// Split a model string like `ollama/glm-5.2` into `(normalized_local, domain)`. +/// The provider (prefix before `/`) becomes the domain; the model name +/// (suffix after `/`) is normalized to an email-local-part-safe form +/// (slashes and dots → dashes). If there's no `/`, the whole string is +/// the local part and the domain defaults to `czsplicer`. +fn normalize_model_email(model: &str) -> (String, String) { + match model.split_once('/') { + Some((provider, name)) => { + let local = name + .chars() + .map(|c| { + if c.is_alphanumeric() || c == '-' || c == '_' { + c + } else { + '-' + } + }) + .collect::(); + (local, provider.to_string()) + } + None => { + let local = model + .chars() + .map(|c| { + if c.is_alphanumeric() || c == '-' || c == '_' { + c + } else { + '-' + } + }) + .collect::(); + (local, "czsplicer".to_string()) + } + } +} + +/// Subject: derive from the first user text under this tree if possible. +fn subject_for_node(node: &Json, role: &str) -> String { + if role == "system" { + // Use the system prompt preview as the subject, so distinct roots get + // distinct threads (avoids Gmail merging unrelated trees). + let preview = node.get("preview").and_then(|v| v.as_str()).unwrap_or(""); + return truncate_subject(preview, 100); + } + let content = node.get("content").and_then(|v| v.as_str()).unwrap_or(""); + truncate_subject(content, 100) +} + +fn truncate_subject(s: &str, max: usize) -> String { + // Collapse all whitespace (including newlines) to single spaces so the + // result is a single-line, RFC 2822-safe header value. + let s: String = s.split_whitespace().collect::>().join(" "); + let s: String = s.chars().take(max).collect(); + if s.is_empty() { + "(empty)".to_string() + } else { + s + } +} + +/// Build the body MIME structure for the chosen mode. +fn body_mime(mode: BodyMode, content: &str, preview: &str) -> BodyMime { + let plain = if content.is_empty() { + preview.to_string() + } else { + content.to_string() + }; + match mode { + BodyMode::Plain => BodyMime { + content_type_header: "text/plain; charset=utf-8".into(), + payload: plain, + }, + BodyMode::HtmlOnly => { + let html = wrap_html(&markdown::to_html(&plain)); + BodyMime { + content_type_header: "text/html; charset=utf-8".into(), + payload: html, + } + } + BodyMode::Html => { + let boundary = format!("cz_{}", blake3::hash(plain.as_bytes()).to_hex()); + let html = wrap_html(&markdown::to_html(&plain)); + let payload = format!( + "--{b}\nContent-Type: text/plain; charset=utf-8\n\n{plain}\n\n--{b}\nContent-Type: text/html; charset=utf-8\n\n{html}\n\n--{b}--\n", + b = boundary + ); + BodyMime { + content_type_header: format!("multipart/alternative; boundary=\"{boundary}\""), + payload, + } + } + } +} + +struct BodyMime { + content_type_header: String, + payload: String, +} + +/// A tool-call attachment: the call parameters followed by `---` and the +/// tool result, as a single text/plain part. +struct ToolAttachment { + filename: String, + payload: String, +} + +/// Extract tool calls and their matching results as attachments for a run of +/// collapsed same-role (assistant) nodes. +/// +/// Pairing is re-derived from each node's own `tool_events` directly. For each +/// assistant node in the run: +/// - its `intro_rid` is the record whose *response* issued this turn's tool +/// call(s) (call events, kind=="call"). In real Aperture captures the +/// record that first includes an assistant message in its request path is +/// the same record whose response generated that turn, so `intro_rid` +/// holds the call; +/// - the record *after* `intro_rid` in that node's `record_ids` echoes the +/// matching tool *result(s)* in its request (result events, kind=="result"). +/// Calls are paired positionally with results (call[i] ↔ result[i]). +/// +/// Anchoring on each node's own `intro_rid`/`record_ids` keeps the pairing +/// correct when consecutive same-role nodes are collapsed: every node is +/// resolved independently, so collapsing shifts no indices. +fn tool_call_attachments( + run: &[&Json], + records: Option<&serde_json::Map>, +) -> Vec { + let rmap = match records { + Some(m) => m, + None => return Vec::new(), + }; + + /// Collect the tool_events of `kind` ("call" or "result") for one record. + fn events_of<'a>( + rmap: &'a serde_json::Map, + rid: i64, + kind: &str, + ) -> Vec<&'a Json> { + rmap.get(&rid.to_string()) + .and_then(|r| r.get("tool_events")) + .and_then(|v| v.as_array()) + .map(|arr| { + arr.iter() + .filter(|e| e.get("kind").and_then(|v| v.as_str()) == Some(kind)) + .collect() + }) + .unwrap_or_default() + } + + let mut out = Vec::new(); + let mut tool_seq = 0usize; + for node in run { + let intro_rid = node.get("intro_rid").and_then(|v| v.as_i64()).unwrap_or(0); + let call_events = events_of(rmap, intro_rid, "call"); + if call_events.is_empty() { + continue; + } + // The record echoing this turn's tool results is the one immediately + // after intro_rid in this node's record_ids (intro_rid is first; the + // next entry is the following captured request, which echoes results). + let result_rid = node + .get("record_ids") + .and_then(|v| v.as_array()) + .and_then(|a| a.get(1)) + .and_then(|v| v.as_i64()) + .unwrap_or(intro_rid); + let result_events = events_of(rmap, result_rid, "result"); + for (i, call) in call_events.iter().enumerate() { + let name = call.get("name").and_then(|v| v.as_str()).unwrap_or("tool"); + let input = call.get("input").and_then(|v| v.as_str()).unwrap_or(""); + let result = result_events + .get(i) + .and_then(|r| r.get("content").and_then(|v| v.as_str())) + .unwrap_or(""); + let payload = if result.is_empty() { + input.to_string() + } else { + format!("{input}\n---\n{result}") + }; + out.push(ToolAttachment { + filename: format!("tool-{tool_seq}-{name}.txt"), + payload, + }); + tool_seq += 1; + } + } + out +} + +/// Wrap the message body and tool attachments in a `multipart/mixed` MIME +/// structure. The first part is the original body (which may itself be +/// multipart/alternative for Html mode); subsequent parts are the tool +/// call+result attachments. +fn wrap_mixed(body: &BodyMime, attachments: &[ToolAttachment]) -> (String, String) { + let mixed_boundary = format!( + "cz_mixed_{}", + blake3::hash(format!("{}{}", body.payload, attachments.len()).as_bytes()).to_hex() + ); + + let mut payload = String::new(); + // First part: the original body (with its own Content-Type). + payload.push_str(&format!("--{mixed_boundary}\n")); + payload.push_str(&format!("Content-Type: {}\n\n", body.content_type_header)); + payload.push_str(&body.payload); + if !body.payload.ends_with('\n') { + payload.push('\n'); + } + payload.push('\n'); + + // Tool attachments. + for att in attachments { + payload.push_str(&format!("--{mixed_boundary}\n")); + payload.push_str(&format!( + "Content-Type: text/plain; charset=utf-8\nContent-Disposition: attachment; filename=\"{}\"\n\n", + sanitize_header(&att.filename) + )); + payload.push_str(&att.payload); + if !att.payload.ends_with('\n') { + payload.push('\n'); + } + payload.push('\n'); + } + + payload.push_str(&format!("--{mixed_boundary}--\n")); + + ( + format!("multipart/mixed; boundary=\"{mixed_boundary}\""), + payload, + ) +} + +fn wrap_html(inner: &str) -> String { + format!("\n\n{inner}\n") +} + +/// Convert an ISO-8601 timestamp to an RFC 2822 date. Falls back to the +/// current time on parse failure (and emits a warning-free fallback). +fn rfc2822_date(iso: &str) -> String { + // We accept `YYYY-MM-DDTHH:MM:SS[.ffffff][Z|+HH:MM]`. Use chrono-free + // parsing: the capture timestamps are always UTC with a trailing Z. + if let Some(rest) = strip_to_iso_basic(iso) { + let secs = parse_iso_utc_secs(&rest); + if let Some(secs) = secs { + return unix_to_rfc2822(secs); + } + } + // Fallback: epoch. + unix_to_rfc2822(0) +} + +/// Keep only `YYYY-MM-DDTHH:MM:SS` (truncate fractional and drop zone). +fn strip_to_iso_basic(iso: &str) -> Option { + let bytes = iso.as_bytes(); + if bytes.len() < 19 { + return None; + } + let head = std::str::from_utf8(&bytes[..19]).ok()?; + if head.len() == 19 + && head.as_bytes()[4] == b'-' + && head.as_bytes()[7] == b'-' + && head.as_bytes()[10] == b'T' + && head.as_bytes()[13] == b':' + && head.as_bytes()[16] == b':' + { + Some(head.to_string()) + } else { + None + } +} + +/// Parse `YYYY-MM-DDTHH:MM:SS` to seconds since epoch (UTC). No leap seconds. +fn parse_iso_utc_secs(iso: &str) -> Option { + let b = iso.as_bytes(); + if b.len() != 19 { + return None; + } + let y: i64 = std::str::from_utf8(&b[0..4]).ok()?.parse().ok()?; + let mo: u32 = std::str::from_utf8(&b[5..7]).ok()?.parse().ok()?; + let d: u32 = std::str::from_utf8(&b[8..10]).ok()?.parse().ok()?; + let h: u32 = std::str::from_utf8(&b[11..13]).ok()?.parse().ok()?; + let mi: u32 = std::str::from_utf8(&b[14..16]).ok()?.parse().ok()?; + let s: u32 = std::str::from_utf8(&b[17..19]).ok()?.parse().ok()?; + Some(unix_from_ymdhms(y, mo, d, h, mi, s)) +} + +fn days_from_civil(y: i64, m: i64, d: i64) -> i64 { + // Howard Hinnant's algorithm (proleptic Gregorian, days since 1970-01-01). + let y = if m <= 2 { y - 1 } else { y }; + let era = if y >= 0 { y } else { y - 399 } / 400; + let yoe = y - era * 400; + let doy = (153 * (m + (if m > 2 { -3 } else { 9 })) + 2) / 5 + d - 1; + let doe = yoe * 365 + yoe / 4 - yoe / 100 + doy; + era * 146097 + doe - 719468 +} + +fn unix_from_ymdhms(y: i64, mo: u32, d: u32, h: u32, mi: u32, s: u32) -> i64 { + let days = days_from_civil(y, mo as i64, d as i64); + days * 86400 + (h as i64) * 3600 + (mi as i64) * 60 + s as i64 +} + +/// Unix seconds -> RFC 2822 (e.g. "Mon, 26 Jun 2026 00:06:12 +0000"). +fn unix_to_rfc2822(secs: i64) -> String { + let (y, mo, d, h, mi, s) = civil_from_unix(secs); + let wd = weekday_from_unix(secs); + let mon = [ + "Jan", "Feb", "Mar", "Apr", "May", "Jun", "Jul", "Aug", "Sep", "Oct", "Nov", "Dec", + ][(mo - 1) as usize]; + let wdays = ["Sun", "Mon", "Tue", "Wed", "Thu", "Fri", "Sat"][wd as usize]; + format!("{wdays}, {d:02} {mon} {y:04} {h:02}:{mi:02}:{s:02} +0000") +} + +/// Unix seconds -> ctime/asctime (e.g. "Mon Jun 26 00:06:12 2026"), the +/// timestamp format mutt's is_from() requires on the mbox From_ postmark +/// line. The day-of-month is space-padded (matching C `asctime()`), which +/// mutt accepts alongside zero-padded forms. +fn unix_to_asctime(secs: i64) -> String { + let (y, mo, d, h, mi, s) = civil_from_unix(secs); + let wd = weekday_from_unix(secs); + let mon = [ + "Jan", "Feb", "Mar", "Apr", "May", "Jun", "Jul", "Aug", "Sep", "Oct", "Nov", "Dec", + ][(mo - 1) as usize]; + let wdays = ["Sun", "Mon", "Tue", "Wed", "Thu", "Fri", "Sat"][wd as usize]; + format!("{wdays} {mon} {d:>2} {h:02}:{mi:02}:{s:02} {y:04}") +} + +/// ISO-8601 timestamp -> ctime/asctime for the From_ postmark. Falls back +/// to epoch on parse failure (mirroring rfc2822_date). +fn asctime_date(iso: &str) -> String { + if let Some(rest) = strip_to_iso_basic(iso) { + if let Some(secs) = parse_iso_utc_secs(&rest) { + return unix_to_asctime(secs); + } + } + unix_to_asctime(0) +} + +fn civil_from_unix(secs: i64) -> (i64, u32, u32, u32, u32, u32) { + let days = secs.div_euclid(86400); + let rem = secs.rem_euclid(86400); + let h = rem / 3600; + let mi = (rem % 3600) / 60; + let s = rem % 60; + // Inverse of days_from_civil. + let z = days + 719468; + let era = if z >= 0 { z } else { z - 146096 } / 146097; + let doe = z - era * 146097; + let yoe = (doe - doe / 1460 + doe / 36524 - doe / 146096) / 365; + let y = yoe + era * 400; + let doy = doe - (365 * yoe + yoe / 4 - yoe / 100); + let mp = (5 * doy + 2) / 153; + let d = doy - (153 * mp + 2) / 5 + 1; + let m = if mp < 10 { mp + 3 } else { mp - 9 }; + let y = if m <= 2 { y + 1 } else { y }; + (y, m as u32, d as u32, h as u32, mi as u32, s as u32) +} + +fn weekday_from_unix(secs: i64) -> u32 { + let days = secs.div_euclid(86400); + // 1970-01-01 was a Thursday (4). + ((days + 4).rem_euclid(7)) as u32 +} + +fn sanitize_header(s: &str) -> String { + s.chars() + .filter(|c| !matches!(c, '\n' | '\r' | '\t')) + .collect() +} + +fn hostname() -> String { + std::env::var("HOSTNAME").unwrap_or_else(|_| "localhost".into()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn unix_roundtrip_zero() { + assert_eq!(unix_to_rfc2822(0), "Thu, 01 Jan 1970 00:00:00 +0000"); + } + + #[test] + fn unix_known() { + // 2026-06-26T00:06:12Z + assert_eq!( + unix_to_rfc2822(1782432372), + "Fri, 26 Jun 2026 00:06:12 +0000" + ); + } + + #[test] + fn iso_parses() { + assert_eq!( + rfc2822_date("2026-06-26T00:06:12.123456789Z"), + "Fri, 26 Jun 2026 00:06:12 +0000" + ); + } + + #[test] + fn iso_fractional_no_z() { + // 2026-03-10T02:20:00Z + assert_eq!( + rfc2822_date("2026-03-10T02:20:00.000Z"), + "Tue, 10 Mar 2026 02:20:00 +0000" + ); + } + + #[test] + fn body_plain_mode_is_text_plain() { + let b = body_mime(BodyMode::Plain, "hello", ""); + assert_eq!(b.content_type_header, "text/plain; charset=utf-8"); + assert_eq!(b.payload, "hello"); + } + + #[test] + fn body_html_only_mode_renders() { + let b = body_mime(BodyMode::HtmlOnly, "**bold**", ""); + assert!(b.content_type_header.starts_with("text/html")); + assert!(b.payload.contains("bold")); + } + + #[test] + fn body_html_mode_is_multipart() { + let b = body_mime(BodyMode::Html, "**bold**", ""); + assert!(b.content_type_header.starts_with("multipart/alternative")); + assert!(b.payload.contains("text/plain")); + assert!(b.payload.contains("text/html")); + assert!(b.payload.contains("bold")); + } + + #[test] + fn message_id_is_unique_per_node() { + let n1 = serde_json::json!({"hash":"abcd","depth":0}); + let n2 = serde_json::json!({"hash":"abcd","depth":1}); + assert_ne!(message_id(&n1, None, 0), message_id(&n2, None, 0)); + } + + #[test] + fn message_id_disambiguates_same_hash_different_parents() { + // Two nodes sharing (hash, depth, thread_idx) but reached via different + // parents must get distinct Message-IDs (a shared subtree fragment). + let n = serde_json::json!({"hash":"abcd","depth":2}); + assert_ne!( + message_id(&n, Some("parent1"), 0), + message_id(&n, Some("parent2"), 0), + "distinct parents -> distinct Message-IDs" + ); + // Roots (no parent) use the "root" sentinel. + assert_eq!( + message_id(&n, None, 0), + message_id(&n, Some("root"), 0), + "None parent hashes the same as the 'root' sentinel" + ); + } +} diff --git a/src/main.rs b/src/main.rs index 8d27ed6..ad71413 100644 --- a/src/main.rs +++ b/src/main.rs @@ -2,6 +2,7 @@ mod builtin; mod commands; mod filter; mod format; +mod mailbox; mod markdown; mod render; mod theme; diff --git a/tests/integration.rs b/tests/integration.rs index 290bab3..4a169eb 100644 --- a/tests/integration.rs +++ b/tests/integration.rs @@ -3076,3 +3076,427 @@ fn thread_redact_secretkey_preset_scrubs_labeled_credentials() { ); assert!(html.contains("[REDACTED]")); } +// =========================================================================== +// tree --format mbox / maildir (email export with threading) +// =========================================================================== +// +// Each trie node becomes one email; threading is via Message-ID/In-Reply-To. +// We parse the emitted mbox with a minimal hand-rolled splitter (no external +// dep) and assert structural invariants: one message per node, every +// In-Reply-To resolves to a real Message-ID, and the chosen body mode +// produces the right Content-Type. + +/// Build a fixture and run `czsplicer tree --format mbox`, returning mbox bytes. +fn thread_mbox(ndjson: &str, body: &str) -> Vec { + let f = Fixture::from_ndjson(ndjson); + let out = Command::cargo_bin("czsplicer") + .unwrap() + .arg("thread") + .arg(&f.cbor_zstd) + .arg("--format") + .arg("mbox") + .arg("--body") + .arg(body) + .arg("-o") + .arg("-") + .assert() + .success() + .get_output() + .stdout + .clone(); + out +} + +/// Split mbox bytes into messages. A message starts at a line beginning with +/// "From " (the envelope line); the body runs until the next "From " at start +/// of a line or EOF. +fn split_mbox(bytes: &[u8]) -> Vec { + let s = String::from_utf8_lossy(bytes); + let mut msgs = Vec::new(); + let mut cur = String::new(); + let mut started = false; + for line in s.lines() { + if line.starts_with("From ") && started { + msgs.push(std::mem::take(&mut cur)); + } + cur.push_str(line); + cur.push('\n'); + started = true; + } + if !cur.is_empty() { + msgs.push(cur); + } + msgs +} + +/// Extract the value of a header from a message (case-insensitive name). +fn header<'a>(msg: &'a str, name: &str) -> Option { + let name_l = name.to_lowercase(); + for line in msg.lines() { + if line.is_empty() { + break; // end of headers + } + if let Some((k, v)) = line.split_once(':') { + if k.trim().to_lowercase() == name_l { + return Some(v.trim().to_string()); + } + } + } + None +} + +#[test] +fn mbox_emits_one_message_per_node() { + // 3 records forming one tree with 3 nodes (sys -> user -> asst+user). + // Actually: rec1=[s,u1], rec2=[s,u1,a1,u2] => trie has 4 nodes. + let nd = format!( + "{}\n{}\n", + rec(1, &body_with_messages("S", &[("user", "q")])), + rec( + 2, + &body_with_messages("S", &[("user", "q"), ("assistant", "a"), ("user", "q2")]) + ), + ); + let out = thread_mbox(&nd, "plain"); + let msgs = split_mbox(&out); + assert_eq!(msgs.len(), 4, "one email per trie node (sys,u,a,u2)"); +} + +#[test] +fn mbox_threading_is_valid() { + // A branch: two records share [s,u1] then diverge with different assistant + // replies. The two assistant nodes both reply to the shared user node. + let nd = format!( + "{}\n{}\n", + rec( + 1, + &body_with_messages("S", &[("user", "q"), ("assistant", "alpha")]) + ), + rec( + 2, + &body_with_messages("S", &[("user", "q"), ("assistant", "beta")]) + ), + ); + let out = thread_mbox(&nd, "plain"); + let msgs = split_mbox(&out); + let mids: Vec = msgs + .iter() + .filter_map(|m| header(m, "Message-ID")) + .collect(); + assert!(!mids.is_empty()); + // Every In-Reply-To must resolve to a Message-ID in the set. + for m in &msgs { + if let Some(parent) = header(m, "In-Reply-To") { + assert!( + mids.contains(&parent), + "dangling In-Reply-To {parent}; known mids = {mids:?}" + ); + } + } + // Exactly one node has no In-Reply-To (the root). + let roots = msgs + .iter() + .filter(|m| header(m, "In-Reply-To").is_none()) + .count(); + assert_eq!(roots, 1, "exactly one root (the system prompt)"); + // Exactly one node has two children (the branch at the user message): it + // appears as two In-Reply-To references pointing at the same Message-ID. + let mut parent_counts: std::collections::HashMap = + std::collections::HashMap::new(); + for m in &msgs { + if let Some(p) = header(m, "In-Reply-To") { + *parent_counts.entry(p).or_default() += 1; + } + } + let branch_parents = parent_counts.values().filter(|&&c| c == 2).count(); + assert_eq!(branch_parents, 1, "one branch point with 2 children"); +} + +#[test] +fn mbox_body_plain_is_text_plain() { + let nd = rec(1, &body_with_messages("S", &[("user", "**bold**")])); + let out = thread_mbox(&nd, "plain"); + let msgs = split_mbox(&out); + // Two nodes (system + user). Both should be text/plain and NOT contain + // rendered HTML. + for m in &msgs { + let ct = header(m, "Content-Type").unwrap_or_default(); + assert!( + ct.starts_with("text/plain"), + "plain mode -> text/plain, got {ct}" + ); + } + let user_msg = msgs + .iter() + .find(|m| header(m, "X-Czsplicer-Role") == Some("user".into())) + .unwrap(); + assert!( + user_msg.contains("**bold**"), + "plain body keeps raw markdown" + ); + assert!(!user_msg.contains(""), "plain body is not rendered"); +} + +#[test] +fn mbox_body_html_is_multipart() { + let nd = rec(1, &body_with_messages("S", &[("user", "**bold**")])); + let out = thread_mbox(&nd, "html"); + let msgs = split_mbox(&out); + for m in &msgs { + let ct = header(m, "Content-Type").unwrap_or_default(); + assert!( + ct.starts_with("multipart/alternative"), + "html mode -> multipart, got {ct}" + ); + assert!(m.contains("text/plain"), "multipart has plain part"); + assert!(m.contains("text/html"), "multipart has html part"); + } +} + +#[test] +fn mbox_body_html_only_is_text_html() { + let nd = rec(1, &body_with_messages("S", &[("user", "**bold**")])); + let out = thread_mbox(&nd, "html-only"); + let msgs = split_mbox(&out); + for m in &msgs { + let ct = header(m, "Content-Type").unwrap_or_default(); + assert!( + ct.starts_with("text/html"), + "html-only mode -> text/html, got {ct}" + ); + } + let user_msg = msgs + .iter() + .find(|m| header(m, "X-Czsplicer-Role") == Some("user".into())) + .unwrap(); + assert!( + user_msg.contains("bold"), + "html-only renders markdown" + ); +} + +#[test] +fn mbox_subject_is_single_line() { + // A message whose content contains newlines must produce a single-line + // Subject (RFC 2822 forbids raw newlines in header values). + let nd = rec( + 1, + &body_with_messages("S", &[("user", "line one\nline two\nline three")]), + ); + let out = thread_mbox(&nd, "plain"); + let msgs = split_mbox(&out); + for m in &msgs { + let subj = header(m, "Subject").unwrap_or_default(); + assert!(!subj.contains('\n'), "Subject must not contain newlines"); + assert!(!subj.contains('\r'), "Subject must not contain CR"); + } +} + +#[test] +fn mbox_carries_record_metadata_headers() { + let nd = serde_json::json!({ + "id":1,"model":"alpha/one","path":"/v1/x","status_code":429, + "timestamp":"2026-06-26T00:00:00Z","duration_ms":1234,"api_type":"oai_completions", + "capture":{"requestBody":body_with_messages("S",&[("user","q")])} + }) + .to_string(); + let out = thread_mbox(&nd, "plain"); + let msgs = split_mbox(&out); + let user_msg = msgs + .iter() + .find(|m| header(m, "X-Czsplicer-Role") == Some("user".into())) + .unwrap(); + assert_eq!(header(user_msg, "X-Czsplicer-Status"), Some("429".into())); + assert_eq!( + header(user_msg, "X-Czsplicer-Model"), + Some("alpha/one".into()) + ); + assert_eq!(header(user_msg, "X-Czsplicer-Depth"), Some("1".into())); +} + +#[test] +fn mbox_from_postmark_uses_ctime_not_rfc2822() { + // mutt's strict is_from() parser only accepts a ctime/asctime timestamp + // on the mbox From_ postmark line (e.g. "Mon Mar 9 22:25:37 2026"). + // An RFC 2822 timestamp ("Mon, 09 Mar 2026 22:25:37 +0000") is rejected as + // a postmark, so mutt recognizes zero messages ("[Msgs:0 ]") on an + // otherwise-valid mbox. The Date: *header* must stay RFC 2822. + let nd = rec(1, &body_with_messages("S", &[("user", "q")])); + let out = thread_mbox(&nd, "plain"); + let s = String::from_utf8(out).unwrap(); + // First line is the postmark. + let postmark = s.lines().next().unwrap(); + assert!( + postmark.starts_with("From czsplicer@localhost "), + "postmark line: {postmark:?}" + ); + // The postmark timestamp must NOT contain a comma (RFC 2822 giveaway) nor + // a "+0000" zone offset; it must be ctime "Www Mmm DD HH:MM:SS YYYY". + let ts = &postmark["From czsplicer@localhost ".len()..]; + assert!( + !ts.contains(','), + "postmark must be ctime (no comma), got: {ts:?}" + ); + assert!( + !ts.contains("+0000"), + "postmark must be ctime (no +0000), got: {ts:?}" + ); + // Sanity: ctime regex. Day-of-month is space-padded for single digits. + let ctime_re = + regex::Regex::new(r"^[A-Z][a-z]{2} [A-Z][a-z]{2} ( |\d)\d \d{2}:\d{2}:\d{2} \d{4}$") + .unwrap(); + assert!( + ctime_re.is_match(ts), + "postmark timestamp {ts:?} is not ctime format" + ); + // The Date: header stays RFC 2822 (comma + zone offset). + let first_msg = split_mbox(s.as_bytes()).into_iter().next().unwrap(); + let date = header(&first_msg, "Date").expect("Date header present"); + assert!( + date.contains(',') && date.contains("+0000"), + "Date: header must stay RFC 2822, got: {date:?}" + ); +} + +#[test] +fn mbox_collapses_consecutive_same_role_nodes() { + // Two records whose message paths both contain a run of two consecutive + // system messages: [system, system, user]. The trie has 3 nodes per path + // (sys, sys, user) — the two system nodes are a same-role single-child + // chain and must collapse into ONE email, so the mbox has 2 emails + // (collapsed-system, user) rather than 3. + let nd = format!( + "{}\n{}\n", + rec( + 1, + &body_with_messages("S1", &[("system", "S2"), ("user", "q")]) + ), + rec( + 2, + &body_with_messages("S1", &[("system", "S2"), ("user", "q"), ("assistant", "a")]) + ), + ); + // Sanity: the tree itself has 4 nodes (sys, sys, user, asst) — no collapse + // happens in the tree, only in the mbox emitter. + let j = thread_json(&nd); + let flat = flatten(&j["trees"]); + let node_count = flat.len(); + assert_eq!(node_count, 4, "trie has 4 nodes: sys, sys, user, asst"); + + let out = thread_mbox(&nd, "plain"); + let msgs = split_mbox(&out); + assert_eq!( + msgs.len(), + 3, + "two consecutive system nodes collapse to one email: sys, user, asst" + ); + // The collapsed email covers depths 0-1 and carries both system contents. + let sys_email = &msgs[0]; + assert_eq!(header(sys_email, "X-Czsplicer-Role"), Some("system".into())); + assert_eq!( + header(sys_email, "X-Czsplicer-Depth"), + Some("0-1".into()), + "collapsed run reports a depth range" + ); + let body = sys_email.splitn(2, "\n\n").nth(1).unwrap_or(""); + assert!(body.contains("S1"), "first system content present"); + assert!(body.contains("S2"), "second system content present"); +} + +#[test] +fn mbox_tool_attachments_pair_call_with_result_across_records() { + // Realistic Aperture shape: record 1's request already includes the + // assistant message at depth 2 (a prior turn), and record 1's *response* + // issues a NEW tool_call. Record 2's request echoes the matching + // tool_result. The depth-2 assistant node's intro_rid is therefore 1 + // (the first record to include the assistant message), which holds the + // call; record_ids[1] (record 2) holds the result. The mbox must emit ONE + // attachment whose payload contains BOTH the call input and the result, + // joined by "---". + let nd = format!( + "{}\n{}\n", + serde_json::json!({ + "id":1,"model":"alpha/one","path":"/v1/x","status_code":200, + "timestamp":"2026-06-26T00:00:00Z", + "capture":{ + "requestBody":serde_json::json!({ + "messages":[ + {"role":"system","content":"S"}, + {"role":"user","content":"do thing"}, + {"role":"assistant","content":[{"type":"tool_use","id":"t0","name":"prev","input":{}}]} + ] + }).to_string(), + "responseBody":serde_json::json!({ + "choices":[{"message":{"role":"assistant","content":"ok","tool_calls":[ + {"id":"call_1","type":"function","function":{"name":"f","arguments":"{}"}} + ]}}] + }).to_string() + } + }).to_string(), + serde_json::json!({ + "id":2,"model":"alpha/one","path":"/v1/x","status_code":200, + "timestamp":"2026-06-26T00:00:01Z", + "capture":{ + "requestBody":serde_json::json!({ + "messages":[ + {"role":"system","content":"S"}, + {"role":"user","content":"do thing"}, + {"role":"assistant","content":[{"type":"tool_use","id":"t0","name":"prev","input":{}}]}, + {"role":"user","content":[{"type":"tool_result","tool_use_id":"call_1","content":"done"}]} + ] + }).to_string(), + "responseBody":"{\"choices\":[{\"message\":{\"role\":\"assistant\",\"content\":\"all done\"}}]}" + } + }).to_string(), + ); + let out = thread_mbox(&nd, "plain"); + let msgs = split_mbox(&out); + // Find the assistant email (the one with a tool attachment). + let asst_email = msgs + .iter() + .find(|m| m.contains("Content-Disposition: attachment")) + .expect("assistant email has a tool-call attachment"); + // The attachment payload pairs the call with its result, separated by "---". + assert!( + asst_email.contains("---\ndone") + || asst_email.contains("---\r\ndone") + || asst_email.contains("---done"), + "attachment pairs the call with the echoed tool_result 'done'" + ); + assert!( + asst_email.contains("filename=\"tool-"), + "attachment has a tool-N-name.txt filename" + ); +} + +#[test] +fn maildir_creates_three_subdirs_with_messages() { + use std::fs; + let nd = format!( + "{}\n{}\n", + rec(1, &body_with_messages("S", &[("user", "q")])), + rec( + 2, + &body_with_messages("S", &[("user", "q"), ("assistant", "a")]) + ), + ); + let f = Fixture::from_ndjson(&nd); + let dir = f.dir.join("maildir_out"); + Command::cargo_bin("czsplicer") + .unwrap() + .arg("thread") + .arg(&f.cbor_zstd) + .arg("--format") + .arg("maildir") + .arg("--body") + .arg("plain") + .arg("-o") + .arg(&dir) + .assert() + .success(); + assert!(dir.join("cur").is_dir(), "cur/ exists"); + assert!(dir.join("new").is_dir(), "new/ exists"); + assert!(dir.join("tmp").is_dir(), "tmp/ exists"); + let new_count = fs::read_dir(dir.join("new")).unwrap().count(); + assert_eq!(new_count, 3, "one file per node in new/ (sys, user, asst)"); +} -- 2.51.2