Skip to main content

doiget_core/
provenance.rs

1//! JSON Lines + SHA-256 hash-chained provenance log.
2//!
3//! Binding spec: `docs/PROVENANCE_LOG.md` (NORMATIVE, §3 row schema, §4 hash
4//! chain). Failure semantics: **fail-closed** — callers MUST abort the fetch
5//! if a log write returns `Err`. See `docs/SECURITY.md` §1.8 and ADR-0006.
6//!
7//! # On-disk format
8//!
9//! - JSON Lines (`.jsonl`): one JSON object per line, terminated by `\n` (LF).
10//! - UTF-8. Timestamps are RFC3339 in UTC.
11//! - Each row is appended via a single `write_all` whose payload always ends
12//!   in `\n`, so a partially-written row is detectable as a missing trailing
13//!   newline rather than a torn JSON record.
14//! - In audit-grade mode (the only mode shipped here), the writer flushes the
15//!   `BufWriter` and `fsync`s the file after every row.
16//!
17//! # Hash chain (PROVENANCE_LOG.md §4)
18//!
19//! Each row carries a `prev_hash` and a `this_hash`. The first row's
20//! `prev_hash` is the literal string `"GENESIS"`. Every subsequent row's
21//! `prev_hash` MUST equal the previous row's `this_hash`.
22//!
23//! When a log file rotates (§6 — not yet implemented in this crate; see TODO
24//! below), the first row of the NEW log file also uses `prev_hash =
25//! "GENESIS"`, restarting the chain.
26//!
27//! `this_hash` is computed as:
28//!
29//! ```text
30//! this_hash = lower_hex(SHA-256(canonical_json(row \ {this_hash})))
31//! ```
32//!
33//! where `canonical_json` is **compact JSON (no whitespace) with object keys
34//! sorted lexicographically** (PROVENANCE_LOG.md §4). For a row with fields
35//! `{ts: "...", ts_seq: 1, event: "fetch", ...}`, the canonical bytes begin
36//! with `{"capability":...` because `capability` is the lex-first top-level
37//! key. Downstream `doiget audit-log --verify` (Phase 1+) relies on this
38//! exact rule — do not change the canonicalization without bumping the spec.
39//!
40//! # In-process serialization
41//!
42//! `ProvenanceLog` holds a `Mutex<LogState>`. All `append` calls within the
43//! same process serialize on this mutex, satisfying the "process-local mutex
44//! on log appender" requirement of `docs/SECURITY.md` §1.8. Cross-process
45//! coordination (multiple `doiget` invocations) is out of scope here and
46//! handled by the higher-level `flock`-based store layer.
47//!
48//! # Session id
49//!
50//! `session_id` (PROVENANCE_LOG.md §3) is a 26-char ULID generated **once per
51//! process invocation** by the caller and stamped into every row written
52//! through the resulting [`ProvenanceLog`]. This crate does not generate the
53//! ULID itself — see [`ProvenanceLog::open`] for the contract.
54//!
55//! # Log rotation and retention (§6)
56//!
57//! Implemented (PROVENANCE_LOG.md §6): when `access.log` exceeds
58//! `ROTATE_BYTES` (100 MiB) a subsequent [`ProvenanceLog::append`]
59//! gzip-compresses the full file to `access.log.<YYYY-MM-DD-HHMMSS>.gz`,
60//! removes the old `access.log`, and writes the incoming row as the
61//! first row of a fresh file with `prev_hash = "GENESIS"` (the hash
62//! chain **restarts** per segment — segments are NOT linked). Rotation
63//! is fail-closed: any gzip / rename / unlink failure aborts the
64//! `append` (the caller's fetch aborts) so the chain never silently
65//! skips. At [`ProvenanceLog::open`], rotated `.gz` segments older than
66//! the retention window (`DOIGET_LOG_RETENTION_DAYS`, default 90; `0`
67//! disables) are deleted **best-effort** (a prune failure is logged,
68//! not fatal — pruning is housekeeping, not integrity).
69//! [`verify_all`] verifies the current file plus every rotated `.gz`
70//! segment (each its own GENESIS-rooted chain).
71
72use std::collections::BTreeMap;
73use std::fs::{File, OpenOptions};
74use std::io::{BufRead, BufReader, BufWriter, Write};
75use std::sync::Mutex;
76
77use flate2::read::GzDecoder;
78use flate2::write::GzEncoder;
79use flate2::Compression;
80
81use camino::{Utf8Path, Utf8PathBuf};
82use chrono::{DateTime, Utc};
83use serde::{Deserialize, Serialize};
84use sha2::{Digest, Sha256};
85
86/// One row of the provenance log (PROVENANCE_LOG.md §3).
87///
88/// The on-disk wire field names match the spec table; struct-field order is
89/// **not** load-bearing for the hash because canonicalization sorts keys
90/// lexicographically (see PROVENANCE_LOG.md §4).
91///
92/// **Schema version**: this struct is the **v2** row shape (ADR-0024).
93/// Every v2 row carries `schema_version = "v2"` literally; the
94/// `canonical_digest` field carries the ADR-0021 §1 audit identity of
95/// the fetch on rows where one applies (`Fetch` / `Resolve` /
96/// `StoreWrite`) and is `None` on session bookend rows
97/// (`SessionStart` / `SessionEnd` / `CapabilityResolved`) that have no
98/// ref. v1 rows (pre-Slice-4) lack both fields and MUST be migrated via
99/// [`migrate_v1_to_v2`] before the v2 binary can read them — the
100/// `deny_unknown_fields` + non-defaulted `schema_version` shape ensures
101/// v1 rows fail to parse loudly rather than producing silent hash-chain
102/// mismatches.
103#[derive(Debug, Clone, Serialize, Deserialize)]
104#[serde(deny_unknown_fields)]
105pub struct LogRow {
106    /// RFC3339 UTC timestamp of the append (millisecond precision).
107    pub ts: DateTime<Utc>,
108    /// Per-session monotonic sequence number, starting at 1.
109    pub ts_seq: u64,
110    /// Event class (see [`LogEvent`]).
111    pub event: LogEvent,
112    /// Optional reference (DOI / arXiv id). Wire field name is `ref`.
113    #[serde(rename = "ref")]
114    pub ref_: Option<String>,
115    /// Optional source name (e.g. `unpaywall`).
116    pub source: Option<String>,
117    /// Result (see [`LogResult`]).
118    pub result: LogResult,
119    /// OA license string (`event=fetch`, `result=ok`); `None` otherwise.
120    pub license: Option<String>,
121    /// Bytes written / fetched, on success rows.
122    pub size_bytes: Option<u64>,
123    /// Path to the stored payload, relative to the store root
124    /// (`event=fetch`, `result=ok`); `None` otherwise.
125    pub store_path: Option<String>,
126    /// Capability under which the row was written (REQUIRED, every row).
127    pub capability: Capability,
128    /// 26-char ULID identifying the process invocation (REQUIRED).
129    pub session_id: String,
130    /// Stable error code on failure rows.
131    pub error_code: Option<String>,
132    /// Row schema version. Always [`LOG_SCHEMA_VERSION`] (`"v2"`) for
133    /// new rows written by this build (ADR-0024). v1 rows lack this
134    /// field; they MUST be migrated via [`migrate_v1_to_v2`] first.
135    pub schema_version: String,
136    /// Canonical-digest of the fetch's audit identity (ADR-0021 §1) as
137    /// 64 lowercase hex chars. Present on rows with a `ref` (`Fetch`,
138    /// `Resolve`, `StoreWrite`); `None` on session bookend rows. The
139    /// digest is computed from a [`crate::CanonicalRef`] whose
140    /// `resolver_profile` matches this row's `source` field for
141    /// migrated v1 rows; new v2 rows MAY pass an explicit
142    /// `resolver_profile` distinct from `source`.
143    pub canonical_digest: Option<String>,
144    /// 64 lowercase hex chars, OR the literal string `"GENESIS"` for the
145    /// first row of a fresh log file.
146    pub prev_hash: String,
147    /// 64 lowercase hex chars. SHA-256 of canonical JSON of THIS row with
148    /// the `this_hash` field removed. See module docs.
149    pub this_hash: String,
150}
151
152/// Provenance-log row schema version this build writes
153/// (`docs/PROVENANCE_LOG.md` §3, ADR-0024).
154///
155/// Bumped from `"v1"` (implicit; pre-Slice-4 rows had no
156/// `schema_version` field) to `"v2"` when the `canonical_digest` column
157/// landed. The v1→v2 migration is one-shot, idempotent, and dry-runnable
158/// via [`migrate_v1_to_v2`].
159pub const LOG_SCHEMA_VERSION: &str = "v2";
160
161/// Event class for a log row (PROVENANCE_LOG.md §3).
162///
163/// Note: result-status (`ok`/`err`/`denied`) lives in [`LogResult`], NOT in
164/// the event variant. So `Fetch` covers both successful and failed fetch
165/// attempts; the row's `result` distinguishes them.
166///
167/// `non_exhaustive` so adding new variants is non-breaking.
168#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
169#[serde(rename_all = "snake_case")]
170#[non_exhaustive]
171pub enum LogEvent {
172    /// Process started; first row of a new session.
173    SessionStart,
174    /// Capability resolution finished (allowed / denied / which env var).
175    CapabilityResolved,
176    /// Reference resolved to a fetch URL.
177    Resolve,
178    /// Fetch attempt (success or failure determined by `result`).
179    Fetch,
180    /// Store write attempt (success or failure determined by `result`).
181    StoreWrite,
182    /// Process ended cleanly.
183    SessionEnd,
184    /// A caller asked again, with `force`, about a ref this session had
185    /// already been answered on -- a request repeat suppression would
186    /// otherwise have replayed (#507, ADR-0057). The row is the record that
187    /// the override was used.
188    RepeatForced,
189}
190
191/// Per-row outcome (PROVENANCE_LOG.md §3). `non_exhaustive` for forward
192/// compatibility.
193#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
194#[serde(rename_all = "snake_case")]
195#[non_exhaustive]
196pub enum LogResult {
197    /// The operation succeeded.
198    Ok,
199    /// The operation failed with an error.
200    Err,
201    /// The operation was denied (e.g. capability gate).
202    Denied,
203}
204
205/// Capability under which a row was written (PROVENANCE_LOG.md §3).
206///
207/// `kebab-case` serde rename emits `oa`, `metadata`, `tdm-elsevier`,
208/// `tdm-aps`, `tdm-springer` exactly as the spec requires. `non_exhaustive`
209/// for forward compatibility.
210#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
211#[serde(rename_all = "kebab-case")]
212#[non_exhaustive]
213pub enum Capability {
214    /// Open access tier.
215    Oa,
216    /// Metadata-only access.
217    Metadata,
218    /// Elsevier TDM (Tier 3, opt-in build).
219    TdmElsevier,
220    /// APS TDM (Tier 3, opt-in build).
221    TdmAps,
222    /// Springer TDM (Tier 3, opt-in build).
223    TdmSpringer,
224    /// IEEE TDM (Tier 3, opt-in build).
225    TdmIeee,
226    /// A PDF the user supplied with `doiget add` (#606): not fetched under
227    /// any capability, and recorded as such rather than as `oa`.
228    UserSupplied,
229}
230
231/// Errors emitted by the provenance log writer. Callers MUST treat any
232/// variant as a fail-closed signal and abort the surrounding fetch.
233#[derive(Debug, thiserror::Error)]
234#[non_exhaustive]
235pub enum LogError {
236    /// I/O error opening, reading, writing, or syncing the log file. Includes
237    /// recovery-time corruption detection where the synthetic message is
238    /// `"corrupted log at line N: …"`.
239    #[error("provenance log io error: {0}")]
240    Io(#[from] std::io::Error),
241    /// Serialization of a row to canonical JSON failed.
242    #[error("provenance log serialization error: {0}")]
243    Serialize(#[from] serde_json::Error),
244    /// Path supplied to [`ProvenanceLog::open`] exists but is not a regular
245    /// file (e.g. a directory or symlink).
246    #[error("provenance log path is not a regular file: {0}")]
247    NotARegularFile(Utf8PathBuf),
248}
249
250/// Append-only writer with in-process serialization.
251#[derive(Debug)]
252pub struct ProvenanceLog {
253    path: Utf8PathBuf,
254    state: Mutex<LogState>,
255    session_id: String,
256    /// §6 rotation threshold, resolved ONCE at [`ProvenanceLog::open`]
257    /// (not per-`append`). Reading `DOIGET_LOG_ROTATE_BYTES` once at
258    /// open — rather than on every append — means a log opened without
259    /// the env set keeps the real 100 MiB threshold for its whole life
260    /// even if another (test) thread later mutates that process-global
261    /// env var; this removes a parallel-test race without serializing
262    /// every multi-append test. `0` = rotation disabled.
263    rotate_threshold: u64,
264    /// What this session has told its callers, per ref, fed from the
265    /// `session_end` rows as they are written (#507, ADR-0057).
266    repeat: crate::repeat::RepeatIndex,
267}
268
269/// Mutable internal state, guarded by [`ProvenanceLog::state`].
270#[derive(Debug)]
271struct LogState {
272    /// `ts_seq` of the **next** row to be appended.
273    next_seq: u64,
274    /// 64 lowercase hex chars; [`GENESIS_HASH`] if the log is empty.
275    last_hash: String,
276}
277
278/// The genesis sentinel used as `prev_hash` for the first row of a log file
279/// (PROVENANCE_LOG.md §3, §6). Also written verbatim as the prev-hash of the
280/// first row after a log rotation (the chain restarts per segment).
281const GENESIS_HASH: &str = "GENESIS";
282
283/// Rotate `access.log` once it reaches this size (PROVENANCE_LOG.md §6:
284/// "100 MB"). 100 MiB. Overridable via the `DOIGET_LOG_ROTATE_BYTES`
285/// env var — an internal ops/testing knob (NOT a documented public
286/// surface): tests set it tiny to exercise rotation without writing
287/// 100 MiB; a value of `0` disables rotation.
288const ROTATE_BYTES: u64 = 100 * 1024 * 1024;
289
290/// Default rotated-segment retention (PROVENANCE_LOG.md §6: "90 days").
291/// Overridable via `DOIGET_LOG_RETENTION_DAYS`; `0` disables pruning.
292const DEFAULT_RETENTION_DAYS: i64 = 90;
293
294/// Resolve the rotation threshold: `DOIGET_LOG_ROTATE_BYTES` if set and
295/// parseable, else [`ROTATE_BYTES`]. `0` (or unparsable) → returns the
296/// value as-is (`0` means "never rotate").
297fn rotate_threshold_bytes() -> u64 {
298    match std::env::var("DOIGET_LOG_ROTATE_BYTES") {
299        Ok(s) => s.trim().parse::<u64>().unwrap_or(ROTATE_BYTES),
300        Err(_) => ROTATE_BYTES,
301    }
302}
303
304/// Resolve retention days from `DOIGET_LOG_RETENTION_DAYS`
305/// (default [`DEFAULT_RETENTION_DAYS`]). `0` disables pruning. A
306/// negative / unparsable value falls back to the default with a warn.
307fn retention_days() -> i64 {
308    match std::env::var("DOIGET_LOG_RETENTION_DAYS") {
309        Ok(s) => match s.trim().parse::<i64>() {
310            Ok(n) if n >= 0 => n,
311            _ => {
312                tracing::warn!(
313                    value = %s,
314                    "DOIGET_LOG_RETENTION_DAYS is not a non-negative integer; \
315                     using the {DEFAULT_RETENTION_DAYS}-day default"
316                );
317                DEFAULT_RETENTION_DAYS
318            }
319        },
320        Err(_) => DEFAULT_RETENTION_DAYS,
321    }
322}
323
324/// gzip-compress `path` to `<file_name>.<YYYY-MM-DD-HHMMSS>.gz` (in the
325/// same directory) and unlink `path` (PROVENANCE_LOG.md §6).
326///
327/// Atomic & fail-closed: the gzip is written to a `.tmp`, fsynced, then
328/// `rename`d into place (so a partial `.gz` is never observable), and
329/// only then is the original removed. Every step propagates its error
330/// to the caller (`ProvenanceLog::append`), which is fail-closed — a
331/// rotation failure aborts the surrounding fetch. Crash safety: a crash
332/// after the rename but before the unlink leaves both the full `.gz`
333/// and the (over-size) `access.log`; the next `append` simply rotates
334/// again, producing a second independently-valid segment — wasteful but
335/// never lossy or corrupt.
336fn rotate_log(path: &Utf8Path) -> Result<(), LogError> {
337    let file_name = path.file_name().ok_or_else(|| {
338        LogError::Io(std::io::Error::other(
339            "provenance log path has no file name; cannot rotate",
340        ))
341    })?;
342    let ts = Utc::now().format("%Y-%m-%d-%H%M%S");
343    let gz_name = format!("{file_name}.{ts}.gz");
344    let dir = path.parent().unwrap_or_else(|| Utf8Path::new("."));
345    let gz_path = dir.join(&gz_name);
346    let tmp_path = dir.join(format!("{gz_name}.tmp"));
347
348    {
349        let mut src = File::open(path)?;
350        let tmp = File::create(&tmp_path)?;
351        let mut enc = GzEncoder::new(BufWriter::new(tmp), Compression::default());
352        std::io::copy(&mut src, &mut enc)?;
353        let bufw = enc.finish()?;
354        let tmp = bufw.into_inner().map_err(|e| {
355            LogError::Io(std::io::Error::other(format!(
356                "gz tmp buf flush failed: {}",
357                e.error()
358            )))
359        })?;
360        tmp.sync_all()?;
361    }
362    std::fs::rename(&tmp_path, &gz_path)?;
363    std::fs::remove_file(path)?;
364    Ok(())
365}
366
367/// Rotated `.gz` segments siblings of `current`, sorted ascending. The
368/// embedded `YYYY-MM-DD-HHMMSS` timestamp makes lexicographic order ==
369/// chronological order.
370fn rotated_segments(current: &Utf8Path) -> Vec<Utf8PathBuf> {
371    let Some(file_name) = current.file_name() else {
372        return Vec::new();
373    };
374    let dir = current.parent().unwrap_or_else(|| Utf8Path::new("."));
375    let prefix = format!("{file_name}.");
376    let mut segs: Vec<Utf8PathBuf> = match std::fs::read_dir(dir.as_std_path()) {
377        Ok(rd) => rd
378            .filter_map(|e| e.ok())
379            .filter_map(|e| Utf8PathBuf::from_path_buf(e.path()).ok())
380            .filter(|p| {
381                p.file_name()
382                    .map(|n| n.starts_with(&prefix) && n.ends_with(".gz"))
383                    .unwrap_or(false)
384            })
385            .collect(),
386        Err(_) => Vec::new(),
387    };
388    segs.sort();
389    segs
390}
391
392/// Delete rotated `.gz` segments older than `days` (PROVENANCE_LOG.md
393/// §6 retention). `days <= 0` is a no-op (disabled). **Best-effort**:
394/// pruning is housekeeping, not integrity, so any failure is logged and
395/// skipped — `ProvenanceLog::open` still succeeds.
396fn prune_rotated_segments(current: &Utf8Path, days: i64) {
397    if days <= 0 {
398        return;
399    }
400    let Some(cutoff) = std::time::SystemTime::now()
401        .checked_sub(std::time::Duration::from_secs(days as u64 * 86_400))
402    else {
403        return;
404    };
405    for seg in rotated_segments(current) {
406        let aged = std::fs::metadata(seg.as_std_path())
407            .and_then(|m| m.modified())
408            .map(|mt| mt < cutoff)
409            .unwrap_or(false);
410        if !aged {
411            continue;
412        }
413        match std::fs::remove_file(seg.as_std_path()) {
414            Ok(()) => tracing::info!(
415                segment = %seg,
416                "provenance: pruned rotated segment past retention"
417            ),
418            Err(e) => tracing::warn!(
419                segment = %seg, error = %e,
420                "provenance: failed to prune rotated segment (best-effort; continuing)"
421            ),
422        }
423    }
424}
425
426/// Verify the full provenance history: every rotated `.gz` segment
427/// (oldest→newest) followed by the current `access.log`. Each segment
428/// is its own GENESIS-rooted hash chain (segments are deliberately NOT
429/// linked across a rotation, PROVENANCE_LOG.md §6), so they are
430/// verified independently and reported per-segment.
431///
432/// The audited [`verify`] function itself is unchanged; this only
433/// orchestrates it over the segment set (gunzipping each `.gz` to a
434/// tempfile first).
435///
436/// # Errors
437///
438/// [`LogError::Io`] on a gunzip / tempfile failure. A missing current
439/// `access.log` is not an error ([`verify`] reports it empty).
440pub fn verify_all(current: &Utf8Path) -> Result<Vec<(Utf8PathBuf, VerifyReport)>, LogError> {
441    let mut out = Vec::new();
442    for seg in rotated_segments(current) {
443        let gz = File::open(seg.as_std_path())?;
444        let mut dec = GzDecoder::new(gz);
445        let tmp = tempfile::NamedTempFile::new().map_err(|e| {
446            LogError::Io(std::io::Error::other(format!(
447                "verify_all: tempfile for {seg}: {e}"
448            )))
449        })?;
450        {
451            let mut w = File::create(tmp.path())?;
452            std::io::copy(&mut dec, &mut w)?;
453            w.sync_all()?;
454        }
455        let tmp_utf8 = Utf8Path::from_path(tmp.path()).ok_or_else(|| {
456            LogError::Io(std::io::Error::other("verify_all: non-utf8 tempfile path"))
457        })?;
458        let report = verify(tmp_utf8)?;
459        out.push((seg, report));
460        // `tmp` (and the gunzipped file) drop here, after verify.
461    }
462    let report = verify(current)?;
463    out.push((current.to_path_buf(), report));
464    Ok(out)
465}
466
467/// Caller-supplied fields for a row. The writer fills in `ts`, `ts_seq`,
468/// `session_id`, `prev_hash`, `this_hash`, and the literal
469/// `schema_version = "v2"` (`LOG_SCHEMA_VERSION`).
470///
471/// Callers SHOULD populate [`Self::canonical_digest`] on rows that have
472/// a meaningful audit identity (`Fetch` / `Resolve` / `StoreWrite` rows
473/// with a `ref`), leaving it `None` on session bookend rows. The digest
474/// is produced by [`crate::CanonicalRef::digest_hex`] from a
475/// `(source_type, source_id, resolver_profile, version)` tuple — see
476/// ADR-0021 §1 for the algorithm and ADR-0024 for the implementation
477/// surface.
478#[derive(Debug, Clone)]
479pub struct RowInput<'a> {
480    /// Event class.
481    pub event: LogEvent,
482    /// Result.
483    pub result: LogResult,
484    /// Capability under which the row is written (REQUIRED for every row).
485    pub capability: Capability,
486    /// Optional DOI / arXiv id -- or, for a software citation's GitHub
487    /// requests (`source` `github` / `github-raw`, #614), the repository or
488    /// release URL cited. Those rows carry no `canonical_digest`: the URL is
489    /// not a `Ref`, and nothing is stored under it.
490    pub ref_: Option<&'a str>,
491    /// Optional source name.
492    pub source: Option<&'a str>,
493    /// Optional error code on failure rows.
494    pub error_code: Option<&'a str>,
495    /// Optional payload size in bytes.
496    pub size_bytes: Option<u64>,
497    /// Optional OA license string (set on `event=fetch`, `result=ok`).
498    pub license: Option<&'a str>,
499    /// Optional store path relative to the store root (set on `event=fetch`,
500    /// `result=ok`).
501    pub store_path: Option<&'a str>,
502    /// Optional canonical-digest (ADR-0021 §1) as 64 lowercase hex
503    /// chars. `None` for session bookend / capability-resolution rows;
504    /// SHOULD be `Some` for `Fetch` / `Resolve` / `StoreWrite` rows
505    /// whose `source` field names the resolver. Build via
506    /// [`crate::Ref::promote`] + [`crate::CanonicalRef::digest_hex`].
507    pub canonical_digest: Option<&'a str>,
508}
509
510// ---------------------------------------------------------------------------
511// Canonical-JSON helper (PROVENANCE_LOG.md §4)
512//
513// Hashing rule (CRITICAL — this is the spec contract for `audit-log --verify`):
514//
515//   this_hash = lower_hex(SHA-256(canonical_json(row \ {this_hash})))
516//
517// Canonical JSON = **compact (no whitespace), keys sorted lexicographically,
518// no trailing whitespace** (§4). Struct field order is deliberately NOT
519// load-bearing here; the canonicalizer sorts the resulting object keys via
520// `BTreeMap<String, Value>`, which serializes in lex-sorted key order.
521//
522// Worked example: for the row fragment `{ts_seq: 1, ts: "..."}` (input order),
523// the canonical bytes after lex sort are `{"ts":"...","ts_seq":1}` because
524// `"ts"` < `"ts_seq"` lexicographically. In v2 (ADR-0024) the lex-first
525// top-level key is `"canonical_digest"` — `"canonical_digest"` < `"capability"`
526// because 'n'(110) < 'p'(112) at byte index 2 (both share the `"ca"`
527// prefix). The pre-v2 lex-first key was `"capability"`.
528// ---------------------------------------------------------------------------
529
530/// Serializable shadow of [`LogRow`] **without** `this_hash`. Used solely as
531/// an intermediate to compute the canonical bytes that `this_hash` is the
532/// SHA-256 of. The wire key names match [`LogRow`]'s `serde` attributes.
533///
534/// v2 shape (ADR-0024): includes `schema_version` and
535/// `canonical_digest`. Both fields participate in the hash chain — a
536/// tampered `canonical_digest` is detected by `audit-log --verify`
537/// exactly like a tampered `ref` or `source` would be.
538#[derive(Serialize)]
539struct RowForHash<'a> {
540    ts: DateTime<Utc>,
541    ts_seq: u64,
542    event: LogEvent,
543    #[serde(rename = "ref")]
544    ref_: Option<&'a str>,
545    source: Option<&'a str>,
546    result: LogResult,
547    license: Option<&'a str>,
548    size_bytes: Option<u64>,
549    store_path: Option<&'a str>,
550    capability: Capability,
551    session_id: &'a str,
552    error_code: Option<&'a str>,
553    schema_version: &'a str,
554    canonical_digest: Option<&'a str>,
555    prev_hash: &'a str,
556}
557
558/// Produce canonical-JSON bytes for a row-without-hash, with object keys
559/// sorted lexicographically per PROVENANCE_LOG.md §4.
560///
561/// Implementation: serialize via `serde_json::to_value` to get a `Value`,
562/// require it be an object, then move its entries into a
563/// `BTreeMap<String, Value>` (which serializes with lex-sorted keys) and
564/// re-serialize compactly. No new dependency required.
565fn canonical_json_for_hash(rfh: &RowForHash<'_>) -> Result<Vec<u8>, LogError> {
566    let value = serde_json::to_value(rfh)?;
567    let map = match value {
568        serde_json::Value::Object(m) => m,
569        // RowForHash is always a struct, so this branch is unreachable in
570        // practice; surface as a serde error if it ever changes.
571        _ => {
572            return Err(LogError::Serialize(serde::de::Error::custom(
573                "RowForHash did not serialize to a JSON object",
574            )));
575        }
576    };
577    let sorted: BTreeMap<String, serde_json::Value> = map.into_iter().collect();
578    Ok(serde_json::to_vec(&sorted)?)
579}
580
581/// Compute `this_hash` for the given row-without-hash. Returns 64 lowercase
582/// hex chars.
583fn compute_this_hash(rfh: &RowForHash<'_>) -> Result<String, LogError> {
584    let bytes = canonical_json_for_hash(rfh)?;
585    let digest = Sha256::digest(&bytes);
586    Ok(hex::encode(digest))
587}
588
589impl ProvenanceLog {
590    /// Open or create the log at `path`, stamping every row with
591    /// `session_id`.
592    ///
593    /// `session_id` MUST be a 26-char ULID generated **once per process**
594    /// invocation by the caller. Re-opening the log within the same process
595    /// reuses the same `session_id`; re-opening in a new process gets a new
596    /// one. This crate intentionally does NOT generate the ULID itself —
597    /// callers are responsible for creating one (e.g. via the `ulid` crate
598    /// already present in the workspace) and threading it through.
599    ///
600    /// If the file exists, scan it once to recover the last `ts_seq` and
601    /// `this_hash`. If the file is missing or empty, the first row will use
602    /// `prev_hash = "GENESIS"` and `ts_seq = 1`.
603    ///
604    /// # Errors
605    ///
606    /// Returns [`LogError::Io`] for I/O failures or if any line fails to
607    /// parse as a [`LogRow`] (synthetic message: `"corrupted log at line N: …"`).
608    /// The writer never silently truncates a corrupt log.
609    ///
610    /// Returns [`LogError::NotARegularFile`] if `path` exists but is not a
611    /// regular file (e.g. a directory).
612    pub fn open(path: impl Into<Utf8PathBuf>, session_id: String) -> Result<Self, LogError> {
613        // Production path: the §6 threshold comes from
614        // `DOIGET_LOG_ROTATE_BYTES` (default 100 MiB), resolved ONCE here.
615        Self::open_with_rotate_threshold(path, session_id, rotate_threshold_bytes())
616    }
617
618    /// [`open`](Self::open) with an explicit rotation threshold instead
619    /// of reading `DOIGET_LOG_ROTATE_BYTES`.
620    ///
621    /// This exists so the rotation tests inject a tiny threshold WITHOUT
622    /// mutating the process-global env var: a global env knob raced
623    /// non-`#[serial]` tests (a concurrent test's `open` would cache the
624    /// tiny threshold and spuriously rotate). `#[serial]` only
625    /// serializes `#[serial]` tests, so injection — not serialization —
626    /// is the robust fix. `0` disables rotation.
627    pub(crate) fn open_with_rotate_threshold(
628        path: impl Into<Utf8PathBuf>,
629        session_id: String,
630        rotate_threshold: u64,
631    ) -> Result<Self, LogError> {
632        let path: Utf8PathBuf = path.into();
633
634        // Ensure the parent directory exists. The provenance log defaults to
635        // `<config>/doiget/access.jsonl`, and on a fresh machine (e.g. a CI
636        // runner where `~/.config/doiget` was never created) neither the
637        // recover-state read nor the first append can open the file — the
638        // append fails with ENOENT and `verify` fail-closes on the LogError.
639        // `create_dir_all` is idempotent; a genuine permission failure still
640        // surfaces as a `LogError` (the correct fail-closed signal).
641        if let Some(parent) = path.parent() {
642            if !parent.as_str().is_empty() {
643                std::fs::create_dir_all(parent.as_std_path())?;
644            }
645        }
646
647        // Reject obvious non-files up front so later `OpenOptions::append`
648        // doesn't produce a confusing platform-dependent error.
649        if path.exists() {
650            let md = std::fs::metadata(&path)?;
651            if !md.is_file() {
652                return Err(LogError::NotARegularFile(path));
653            }
654        }
655
656        let (next_seq, last_hash) = recover_state(&path)?;
657
658        // §6 retention: prune rotated `.gz` segments older than the
659        // window. Best-effort — pruning is housekeeping, not integrity,
660        // so a failure is logged and `open` still succeeds (unlike
661        // rotation, which is fail-closed).
662        prune_rotated_segments(&path, retention_days());
663
664        Ok(Self {
665            path,
666            state: Mutex::new(LogState {
667                next_seq,
668                last_hash,
669            }),
670            session_id,
671            rotate_threshold,
672            repeat: crate::repeat::RepeatIndex::default(),
673        })
674    }
675
676    /// Append a row. Computes `prev_hash`, `ts_seq`, `ts`, `session_id`, and
677    /// `this_hash`; the caller only supplies the semantic fields via
678    /// [`RowInput`].
679    ///
680    /// Returns the assigned `ts_seq` on success.
681    ///
682    /// # Errors
683    ///
684    /// Returns [`LogError`] on serialization, I/O, or fsync failure. Callers
685    /// MUST treat this as fail-closed and abort the surrounding fetch.
686    pub fn append(&self, input: RowInput<'_>) -> Result<u64, LogError> {
687        // Hold the mutex for the entire append: serialize + write + flush +
688        // fsync + state update. This is the in-process serialization point
689        // promised by `docs/SECURITY.md` §1.8.
690        //
691        // A poisoned mutex only happens if a previous `append` panicked
692        // mid-write. Surface that as an I/O error rather than propagating
693        // a panic.
694        let mut state = self
695            .state
696            .lock()
697            .map_err(|_| LogError::Io(std::io::Error::other("provenance log mutex poisoned")))?;
698
699        // §6 rotation, BEFORE this row is written. If `access.log` has
700        // reached the threshold, gzip+rename it and reset the in-memory
701        // chain state so this row becomes the GENESIS-rooted first row of
702        // a fresh file. Fail-closed: a rotation error aborts the append
703        // (the `?`), so the caller's fetch aborts and the chain never
704        // silently continues in an over-size or half-rotated file. The
705        // `state` mutex is held, so rotation is serialized with appends.
706        let threshold = self.rotate_threshold;
707        if threshold > 0 {
708            let size = match std::fs::metadata(&self.path) {
709                Ok(m) => m.len(),
710                Err(e) if e.kind() == std::io::ErrorKind::NotFound => 0,
711                Err(e) => return Err(LogError::Io(e)),
712            };
713            if size >= threshold {
714                rotate_log(&self.path)?;
715                state.next_seq = 1;
716                state.last_hash = GENESIS_HASH.to_string();
717            }
718        }
719
720        let ts_seq = state.next_seq;
721        let prev_hash = state.last_hash.clone();
722        let ts = Utc::now();
723
724        let rfh = RowForHash {
725            ts,
726            ts_seq,
727            event: input.event,
728            ref_: input.ref_,
729            source: input.source,
730            result: input.result,
731            license: input.license,
732            size_bytes: input.size_bytes,
733            store_path: input.store_path,
734            capability: input.capability,
735            session_id: &self.session_id,
736            error_code: input.error_code,
737            schema_version: LOG_SCHEMA_VERSION,
738            canonical_digest: input.canonical_digest,
739            prev_hash: &prev_hash,
740        };
741
742        let this_hash = compute_this_hash(&rfh)?;
743
744        // Build the on-disk row. Owned strings here because `LogRow` does
745        // not borrow.
746        let row = LogRow {
747            ts,
748            ts_seq,
749            event: input.event,
750            ref_: input.ref_.map(str::to_string),
751            source: input.source.map(str::to_string),
752            result: input.result,
753            license: input.license.map(str::to_string),
754            size_bytes: input.size_bytes,
755            store_path: input.store_path.map(str::to_string),
756            capability: input.capability,
757            session_id: self.session_id.clone(),
758            error_code: input.error_code.map(str::to_string),
759            schema_version: LOG_SCHEMA_VERSION.to_string(),
760            canonical_digest: input.canonical_digest.map(str::to_string),
761            prev_hash,
762            this_hash: this_hash.clone(),
763        };
764
765        // Serialize, append `\n`, write_all in one syscall, flush BufWriter,
766        // fsync the underlying file. `\n` is part of the same buffer, so a
767        // crash mid-write leaves at most a partial line (no trailing `\n`),
768        // which is detectable on recovery as a corrupted final line.
769        let mut bytes = serde_json::to_vec(&row)?;
770        bytes.push(b'\n');
771
772        let file = OpenOptions::new()
773            .create(true)
774            .append(true)
775            .open(&self.path)?;
776        let mut writer = BufWriter::new(file);
777        writer.write_all(&bytes)?;
778        writer.flush()?;
779        // `into_inner` to recover the underlying File for `sync_all`.
780        let file = writer.into_inner().map_err(|e| {
781            LogError::Io(std::io::Error::other(format!(
782                "buf writer flush failed: {}",
783                e.error()
784            )))
785        })?;
786        file.sync_all()?;
787
788        // Only after a successful fsync do we advance the in-memory state.
789        // If any of the above fails, the next `append` retries from the
790        // same `(ts_seq, prev_hash)` — at most a torn last line on disk.
791        state.next_seq = ts_seq + 1;
792        state.last_hash = this_hash;
793
794        // #507: only a row that is durably on disk feeds repeat
795        // suppression -- the log is the record of what callers were told.
796        if input.event == LogEvent::SessionEnd {
797            if let Some(r) = input.ref_ {
798                let code = input.error_code.and_then(crate::ErrorCode::from_wire);
799                if code.is_some() || input.result == LogResult::Ok {
800                    self.repeat.observe(r, code);
801                }
802            }
803        }
804
805        Ok(ts_seq)
806    }
807
808    /// What this session has already told callers about each ref (#507).
809    #[must_use]
810    pub fn repeat(&self) -> &crate::repeat::RepeatIndex {
811        &self.repeat
812    }
813
814    /// Returns the path the log was opened at. Useful for tests and audit tooling.
815    pub fn path(&self) -> &Utf8Path {
816        &self.path
817    }
818
819    /// Returns the session id stamped into every row written through this
820    /// writer.
821    pub fn session_id(&self) -> &str {
822        &self.session_id
823    }
824}
825
826/// Scan an existing log to recover `(next_seq, last_hash)`.
827///
828/// Walk every line, parse as [`LogRow`], track the last successfully parsed
829/// row. If parsing fails, return [`LogError::Io`] with a synthetic
830/// `"corrupted log at line N: …"` message — never silently truncate.
831fn recover_state(path: &Utf8Path) -> Result<(u64, String), LogError> {
832    let file = match File::open(path) {
833        Ok(f) => f,
834        Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
835            return Ok((1, GENESIS_HASH.to_string()));
836        }
837        Err(e) => return Err(LogError::Io(e)),
838    };
839
840    let reader = BufReader::new(file);
841    let mut last_seq: u64 = 0;
842    let mut last_hash: String = GENESIS_HASH.to_string();
843
844    for (idx, line_res) in reader.lines().enumerate() {
845        let line_no = idx + 1;
846        let line = line_res?;
847        if line.is_empty() {
848            // Tolerate trailing/empty lines silently — they are not data.
849            continue;
850        }
851        let row: LogRow = serde_json::from_str(&line).map_err(|e| {
852            LogError::Io(std::io::Error::new(
853                std::io::ErrorKind::InvalidData,
854                format!("corrupted log at line {}: {}", line_no, e),
855            ))
856        })?;
857        last_seq = row.ts_seq;
858        last_hash = row.this_hash;
859    }
860
861    if last_seq == 0 {
862        Ok((1, GENESIS_HASH.to_string()))
863    } else {
864        Ok((last_seq + 1, last_hash))
865    }
866}
867
868// ---------------------------------------------------------------------------
869// Verification (`doiget audit-log --verify`)
870//
871// The provenance log is a JSON Lines file with a SHA-256 hash chain
872// (PROVENANCE_LOG.md §4). Tampering is detected by recomputing every row's
873// `this_hash` and validating the chain. This module provides the offline
874// verifier; the CLI wrapper lives in `doiget-cli::commands::audit_log`.
875//
876// Failure model: returning `Err` is reserved for I/O failures opening / reading
877// the file. Per-row issues (parse failures, hash/chain mismatches, sequence
878// regressions) are accumulated into [`VerifyReport::errors`] so callers can
879// report them all in one pass — this is the contract Phase 1 ships.
880// ---------------------------------------------------------------------------
881
882/// Outcome of [`verify`]: per-row chain status across the entire log.
883#[derive(Debug, Clone)]
884#[non_exhaustive]
885pub struct VerifyReport {
886    /// Total non-empty lines processed (1-based count).
887    pub total_rows: usize,
888    /// Rows whose hash, chain link, and `ts_seq` all validated.
889    pub ok_rows: usize,
890    /// Issues encountered, in encounter order. Line numbers are 1-based.
891    pub errors: Vec<VerifyIssue>,
892}
893
894impl VerifyReport {
895    /// An empty, all-clear report — used when the log file is absent.
896    fn empty() -> Self {
897        Self {
898            total_rows: 0,
899            ok_rows: 0,
900            errors: Vec::new(),
901        }
902    }
903}
904
905/// A single issue discovered by [`verify`].
906#[derive(Debug, Clone)]
907#[non_exhaustive]
908pub struct VerifyIssue {
909    /// 1-based line number where the issue was detected.
910    pub line: usize,
911    /// Classification of the issue (see [`VerifyIssueKind`]).
912    pub kind: VerifyIssueKind,
913    /// Human-readable description (caller may format for stderr/stdout).
914    pub message: String,
915}
916
917/// Classification of a [`VerifyIssue`]. `non_exhaustive` for forward
918/// compatibility — future kinds may include `SessionIdChange`, etc.
919#[derive(Debug, Clone, Copy, PartialEq, Eq)]
920#[non_exhaustive]
921pub enum VerifyIssueKind {
922    /// Row failed to parse as [`LogRow`] (corrupted JSON or unknown field).
923    ParseError,
924    /// `prev_hash` did not match the previous row's `this_hash` (or the
925    /// genesis sentinel on row 1).
926    PrevHashMismatch,
927    /// Row's stored `this_hash` did not match the recomputed canonical-JSON
928    /// SHA-256.
929    ThisHashMismatch,
930    /// `ts_seq` did not increase strictly monotonically (within a session;
931    /// see PROVENANCE_LOG.md §3 + §6 — chain restarts after rotation are
932    /// permitted to reset `ts_seq` and are detected via the genesis sentinel).
933    SequenceJump,
934}
935
936/// Verify the entire log file at `path`.
937///
938/// Returns `Ok(VerifyReport)` regardless of whether the chain validates;
939/// callers inspect `report.errors.is_empty()` to determine pass/fail.
940/// Returns `Err` only when the file itself cannot be opened or read at the
941/// I/O level.
942///
943/// Behavior:
944///
945/// - A missing file is treated as a clean, empty log (no tampering possible
946///   on bytes that don't exist) and returns an empty report after a `warn!`.
947/// - Empty / blank lines are skipped — they are not data per the writer's
948///   on-disk format (PROVENANCE_LOG.md §2).
949/// - On a row that fails to parse as [`LogRow`], a `ParseError` is recorded
950///   and verification continues on the next line. The chain anchor does NOT
951///   advance through an unparsable row, so the next valid row's `prev_hash`
952///   is checked against the last successfully parsed row (or against
953///   `"GENESIS"` if no valid row has been seen yet).
954/// - A `prev_hash == "GENESIS"` sentinel marks a chain restart (first row of
955///   a fresh / rotated log per §6) and resets the `ts_seq` monotonicity
956///   anchor — `ts_seq` is NOT compared to the prior row across a restart.
957pub fn verify(path: &Utf8Path) -> Result<VerifyReport, LogError> {
958    let file = match File::open(path) {
959        Ok(f) => f,
960        Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
961            tracing::warn!(
962                path = %path,
963                "audit-log verify: log file does not exist; reporting empty"
964            );
965            return Ok(VerifyReport::empty());
966        }
967        Err(e) => return Err(LogError::Io(e)),
968    };
969
970    let reader = BufReader::new(file);
971    let mut report = VerifyReport::empty();
972
973    // Anchor for the chain check: the LAST SUCCESSFULLY PARSED row. The chain
974    // is anchored to the bytes on disk, not to a hypothetical "should have
975    // been". This matches the spec — tampering at row N must surface both as
976    // a hash mismatch on N and as a chain break on N+1.
977    let mut prev_row: Option<LogRow> = None;
978
979    for (idx, line_res) in reader.lines().enumerate() {
980        let line_no = idx + 1;
981        let line = line_res?;
982        if line.is_empty() {
983            continue;
984        }
985
986        report.total_rows += 1;
987
988        let row: LogRow = match serde_json::from_str(&line) {
989            Ok(r) => r,
990            Err(e) => {
991                report.errors.push(VerifyIssue {
992                    line: line_no,
993                    kind: VerifyIssueKind::ParseError,
994                    message: format!("failed to parse row as LogRow: {e}"),
995                });
996                // Chain anchor cannot advance through an unparsable row;
997                // leave `prev_row` untouched so the next valid row's
998                // `prev_hash` is checked against the last-known anchor (or
999                // GENESIS if we never had one).
1000                continue;
1001            }
1002        };
1003
1004        let mut row_ok = true;
1005
1006        // 1. Recompute `this_hash` from canonical JSON (row \ {this_hash}).
1007        let rfh = RowForHash {
1008            ts: row.ts,
1009            ts_seq: row.ts_seq,
1010            event: row.event,
1011            ref_: row.ref_.as_deref(),
1012            source: row.source.as_deref(),
1013            result: row.result,
1014            license: row.license.as_deref(),
1015            size_bytes: row.size_bytes,
1016            store_path: row.store_path.as_deref(),
1017            capability: row.capability,
1018            session_id: &row.session_id,
1019            error_code: row.error_code.as_deref(),
1020            schema_version: &row.schema_version,
1021            canonical_digest: row.canonical_digest.as_deref(),
1022            prev_hash: &row.prev_hash,
1023        };
1024        match compute_this_hash(&rfh) {
1025            Ok(recomputed) => {
1026                if recomputed != row.this_hash {
1027                    report.errors.push(VerifyIssue {
1028                        line: line_no,
1029                        kind: VerifyIssueKind::ThisHashMismatch,
1030                        message: format!(
1031                            "this_hash mismatch: stored={}, recomputed={}",
1032                            row.this_hash, recomputed
1033                        ),
1034                    });
1035                    row_ok = false;
1036                }
1037            }
1038            Err(e) => {
1039                // Canonicalization itself failed — surface as a hash
1040                // mismatch with the underlying error in the message.
1041                report.errors.push(VerifyIssue {
1042                    line: line_no,
1043                    kind: VerifyIssueKind::ThisHashMismatch,
1044                    message: format!("failed to recompute this_hash: {e}"),
1045                });
1046                row_ok = false;
1047            }
1048        }
1049
1050        // 2. Chain link: `prev_hash` matches anchor (GENESIS on row 1 / after
1051        //    a chain restart, prior row's `this_hash` otherwise).
1052        let is_genesis = row.prev_hash == GENESIS_HASH;
1053        match &prev_row {
1054            None => {
1055                // First non-empty row in the file: must declare GENESIS.
1056                if !is_genesis {
1057                    report.errors.push(VerifyIssue {
1058                        line: line_no,
1059                        kind: VerifyIssueKind::PrevHashMismatch,
1060                        message: format!(
1061                            "first row must have prev_hash=\"GENESIS\", got {:?}",
1062                            row.prev_hash
1063                        ),
1064                    });
1065                    row_ok = false;
1066                }
1067            }
1068            Some(prev) => {
1069                if is_genesis {
1070                    // Chain restart (rotation per §6) — accepted, no link
1071                    // check, and the `ts_seq` monotonicity anchor resets
1072                    // (handled below via `is_genesis`).
1073                } else if row.prev_hash != prev.this_hash {
1074                    report.errors.push(VerifyIssue {
1075                        line: line_no,
1076                        kind: VerifyIssueKind::PrevHashMismatch,
1077                        message: format!(
1078                            "prev_hash mismatch: row stores {}, previous row's this_hash is {}",
1079                            row.prev_hash, prev.this_hash
1080                        ),
1081                    });
1082                    row_ok = false;
1083                }
1084            }
1085        }
1086
1087        // 3. ts_seq monotonicity — strictly greater than the previous row's
1088        //    `ts_seq`, EXCEPT across a chain restart (where `ts_seq` resets).
1089        if let Some(prev) = &prev_row {
1090            if !is_genesis && row.ts_seq <= prev.ts_seq {
1091                report.errors.push(VerifyIssue {
1092                    line: line_no,
1093                    kind: VerifyIssueKind::SequenceJump,
1094                    message: format!(
1095                        "ts_seq did not increase strictly: previous={}, current={}",
1096                        prev.ts_seq, row.ts_seq
1097                    ),
1098                });
1099                row_ok = false;
1100            }
1101        }
1102
1103        if row_ok {
1104            report.ok_rows += 1;
1105        }
1106
1107        // Advance the anchor to the just-parsed row (whether or not it had
1108        // issues — the on-disk bytes ARE the chain).
1109        prev_row = Some(row);
1110    }
1111
1112    Ok(report)
1113}
1114
1115// ---------------------------------------------------------------------------
1116// v1 → v2 migration (ADR-0024, `docs/PROVENANCE_LOG.md` §"Schema migration").
1117//
1118// v1 rows lack `schema_version` and `canonical_digest`; the v2 binary
1119// fails loudly when asked to read them (see `recover_state` /
1120// `verify`). The migration recovers a v2 log from a v1 file by:
1121//
1122//   1. Parsing every v1 row via the [`V1LogRow`] shadow struct.
1123//   2. Deriving a [`crate::CanonicalRef`] from the v1 `(ref, source)`
1124//      pair — `source` becomes `resolver_profile`, `version` is `None`
1125//      (ADR-0021 §1 → ADR-0024 migration recipe).
1126//   3. Re-computing the SHA-256 hash chain across the new row
1127//      payloads. The v1 chain is invalidated by the schema change; the
1128//      v2 chain restarts at the first row's stored `prev_hash` (which
1129//      is `"GENESIS"` on a fresh log).
1130//   4. Writing the new rows to `<log_path>.v2-migrated`, then
1131//      atomically renaming it onto `<log_path>` after backing up the
1132//      original to `<log_path>.v1-backup`.
1133//
1134// The migration is **idempotent**: running it on an already-v2 log
1135// re-parses every row as v2, recomputes the same hash chain, and
1136// produces a byte-equivalent output.
1137//
1138// The migration is **dry-runnable**: `dry_run = true` returns a
1139// [`MigrationReport`] summarizing what would change without touching
1140// disk.
1141// ---------------------------------------------------------------------------
1142
1143/// Summary of a [`migrate_v1_to_v2`] run.
1144///
1145/// Marked `#[non_exhaustive]` so future fields (e.g. a per-row error
1146/// list, an aborted-row count) can be added without breaking callers
1147/// that pattern-match.
1148///
1149/// `Serialize` enables `provenance migrate --mode json` (#204) — the
1150/// wire form is `{"rows_rewritten": N, "dry_run": bool,
1151/// "first_row_v1_chain_hash": "...", "first_row_v2_chain_hash": "..."}`.
1152///
1153/// # Wire-format stability (post-#208 self-review §1)
1154///
1155/// Once a release ships with the [`Serialize`] derive, the field
1156/// **names** below become part of the public API. Renaming a field is
1157/// then a semver minor bump and warrants a CHANGELOG \[BREAKING\] note;
1158/// new fields are still safe (per `#[non_exhaustive]`).
1159#[derive(Debug, Clone, Serialize)]
1160#[non_exhaustive]
1161pub struct MigrationReport {
1162    /// Number of rows rewritten (or that WOULD be rewritten under
1163    /// `dry_run`).
1164    pub rows_rewritten: u64,
1165    /// Whether this was a dry-run preview (`true`) or a live rewrite
1166    /// (`false`).
1167    pub dry_run: bool,
1168    /// Stored `this_hash` of the first input row (the v1 chain anchor).
1169    /// `"GENESIS"` is reported as the literal `"GENESIS"` when the log
1170    /// was empty.
1171    pub first_row_v1_chain_hash: String,
1172    /// Recomputed `this_hash` of the first migrated row under the v2
1173    /// canonicalization. Equal to [`Self::first_row_v1_chain_hash`]
1174    /// only if the input was already v2 (idempotent case).
1175    pub first_row_v2_chain_hash: String,
1176}
1177
1178/// v1 row shadow struct used ONLY by [`migrate_v1_to_v2`]. The
1179/// non-defaulted v2 fields (`schema_version`, `canonical_digest`) are
1180/// absent here; `deny_unknown_fields` rejects unexpected v2 fields so a
1181/// v2 row on disk fails to parse as v1, letting the migrator detect
1182/// already-v2 input via fallback to the v2 parser.
1183#[derive(Debug, Clone, Deserialize, Serialize)]
1184#[serde(deny_unknown_fields)]
1185struct V1LogRow {
1186    ts: DateTime<Utc>,
1187    ts_seq: u64,
1188    event: LogEvent,
1189    #[serde(rename = "ref")]
1190    ref_: Option<String>,
1191    source: Option<String>,
1192    result: LogResult,
1193    license: Option<String>,
1194    size_bytes: Option<u64>,
1195    store_path: Option<String>,
1196    capability: Capability,
1197    session_id: String,
1198    error_code: Option<String>,
1199    prev_hash: String,
1200    this_hash: String,
1201}
1202
1203/// Minimal in-memory representation a v1 OR v2 row can be promoted to
1204/// before re-hashing.
1205#[derive(Debug, Clone)]
1206struct MigrationRowSeed {
1207    ts: DateTime<Utc>,
1208    ts_seq: u64,
1209    event: LogEvent,
1210    ref_: Option<String>,
1211    source: Option<String>,
1212    result: LogResult,
1213    license: Option<String>,
1214    size_bytes: Option<u64>,
1215    store_path: Option<String>,
1216    capability: Capability,
1217    session_id: String,
1218    error_code: Option<String>,
1219    /// `None` for v1 inputs (the digest is computed during migration);
1220    /// `Some(...)` for already-v2 inputs (carried through verbatim for
1221    /// idempotency).
1222    canonical_digest_in: Option<String>,
1223    /// As stored on disk in the input. Used only for the
1224    /// `first_row_v1_chain_hash` field of [`MigrationReport`].
1225    stored_this_hash: String,
1226}
1227
1228/// Migrate a v1 provenance log to v2 (ADR-0024).
1229///
1230/// Returns a [`MigrationReport`] describing how many rows were (or
1231/// would be) rewritten and the first-row chain-anchor delta. The
1232/// migration is idempotent: running it twice produces byte-equivalent
1233/// output the second time.
1234///
1235/// On a missing log file, returns a no-op report (`rows_rewritten = 0`,
1236/// `first_row_v1_chain_hash = "GENESIS"`, `first_row_v2_chain_hash =
1237/// "GENESIS"`) — there is nothing to migrate.
1238///
1239/// # Errors
1240///
1241/// Returns [`LogError::Io`] on I/O failures and on rows that fail to
1242/// parse as either v1 or v2 (the synthetic message names the line
1243/// number). Returns [`LogError::Serialize`] on canonicalization
1244/// failures.
1245pub fn migrate_v1_to_v2(log_path: &Utf8Path, dry_run: bool) -> Result<MigrationReport, LogError> {
1246    use std::io::BufRead;
1247
1248    // -- 1. Read the input log, parsing each line as v1 OR (idempotent
1249    //       fallback) v2. --------------------------------------------------
1250    let file = match File::open(log_path) {
1251        Ok(f) => f,
1252        Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
1253            return Ok(MigrationReport {
1254                rows_rewritten: 0,
1255                dry_run,
1256                first_row_v1_chain_hash: GENESIS_HASH.to_string(),
1257                first_row_v2_chain_hash: GENESIS_HASH.to_string(),
1258            });
1259        }
1260        Err(e) => return Err(LogError::Io(e)),
1261    };
1262    let reader = BufReader::new(file);
1263    let mut seeds: Vec<MigrationRowSeed> = Vec::new();
1264
1265    for (idx, line_res) in reader.lines().enumerate() {
1266        let line_no = idx + 1;
1267        let line = line_res?;
1268        if line.is_empty() {
1269            continue;
1270        }
1271        // Try v1 first. If it fails, try v2 (idempotency: re-migrating
1272        // a v2 log MUST succeed and produce equivalent output).
1273        let seed = if let Ok(v1) = serde_json::from_str::<V1LogRow>(&line) {
1274            MigrationRowSeed {
1275                ts: v1.ts,
1276                ts_seq: v1.ts_seq,
1277                event: v1.event,
1278                ref_: v1.ref_,
1279                source: v1.source,
1280                result: v1.result,
1281                license: v1.license,
1282                size_bytes: v1.size_bytes,
1283                store_path: v1.store_path,
1284                capability: v1.capability,
1285                session_id: v1.session_id,
1286                error_code: v1.error_code,
1287                canonical_digest_in: None,
1288                stored_this_hash: v1.this_hash,
1289            }
1290        } else {
1291            match serde_json::from_str::<LogRow>(&line) {
1292                Ok(v2) => MigrationRowSeed {
1293                    ts: v2.ts,
1294                    ts_seq: v2.ts_seq,
1295                    event: v2.event,
1296                    ref_: v2.ref_,
1297                    source: v2.source,
1298                    result: v2.result,
1299                    license: v2.license,
1300                    size_bytes: v2.size_bytes,
1301                    store_path: v2.store_path,
1302                    capability: v2.capability,
1303                    session_id: v2.session_id,
1304                    error_code: v2.error_code,
1305                    canonical_digest_in: v2.canonical_digest,
1306                    stored_this_hash: v2.this_hash,
1307                },
1308                Err(e) => {
1309                    return Err(LogError::Io(std::io::Error::new(
1310                        std::io::ErrorKind::InvalidData,
1311                        format!("migration: line {line_no} is neither v1 nor v2: {e}"),
1312                    )));
1313                }
1314            }
1315        };
1316        seeds.push(seed);
1317    }
1318
1319    // -- 2. Derive `canonical_digest` for each seed that lacks one. ------
1320    //
1321    // For v1 rows: build a CanonicalRef from
1322    //   - source_type from `event`/`ref` shape (DOI prefix `10.` vs
1323    //     arXiv) — we use a heuristic that matches `Ref::parse`'s rule
1324    //     (`starts_with "10."` ⇒ DOI; else arXiv).
1325    //   - source_id = ref value (verbatim).
1326    //   - resolver_profile = source value (verbatim, ADR-0021 §3
1327    //     migration recipe).
1328    //   - version = None.
1329    //
1330    // Rows without a `ref` (session bookend) keep `canonical_digest =
1331    // None` per the v2 row contract.
1332
1333    fn derive_digest(seed: &MigrationRowSeed) -> Option<String> {
1334        let ref_str = seed.ref_.as_deref()?;
1335        let source_key = seed.source.as_deref().unwrap_or("");
1336        // Heuristic: bare DOIs always start `10.`; everything else is
1337        // treated as an arXiv id. Mirrors `Ref::parse` rule 3/4.
1338        let source_type = if ref_str.starts_with("10.") {
1339            crate::SourceType::Doi
1340        } else {
1341            crate::SourceType::Arxiv
1342        };
1343        let c = crate::CanonicalRef::new(source_type, ref_str, source_key, None);
1344        Some(c.digest_hex())
1345    }
1346
1347    let digests: Vec<Option<String>> = seeds
1348        .iter()
1349        .map(|s| s.canonical_digest_in.clone().or_else(|| derive_digest(s)))
1350        .collect();
1351
1352    // -- 3. Rebuild the hash chain across the v2 payloads. ----------------
1353    let mut out_rows: Vec<LogRow> = Vec::with_capacity(seeds.len());
1354    let mut prev_hash: String = GENESIS_HASH.to_string();
1355
1356    for (seed, digest) in seeds.iter().zip(digests.iter()) {
1357        let rfh = RowForHash {
1358            ts: seed.ts,
1359            ts_seq: seed.ts_seq,
1360            event: seed.event,
1361            ref_: seed.ref_.as_deref(),
1362            source: seed.source.as_deref(),
1363            result: seed.result,
1364            license: seed.license.as_deref(),
1365            size_bytes: seed.size_bytes,
1366            store_path: seed.store_path.as_deref(),
1367            capability: seed.capability,
1368            session_id: &seed.session_id,
1369            error_code: seed.error_code.as_deref(),
1370            schema_version: LOG_SCHEMA_VERSION,
1371            canonical_digest: digest.as_deref(),
1372            prev_hash: &prev_hash,
1373        };
1374        let this_hash = compute_this_hash(&rfh)?;
1375        let row = LogRow {
1376            ts: seed.ts,
1377            ts_seq: seed.ts_seq,
1378            event: seed.event,
1379            ref_: seed.ref_.clone(),
1380            source: seed.source.clone(),
1381            result: seed.result,
1382            license: seed.license.clone(),
1383            size_bytes: seed.size_bytes,
1384            store_path: seed.store_path.clone(),
1385            capability: seed.capability,
1386            session_id: seed.session_id.clone(),
1387            error_code: seed.error_code.clone(),
1388            schema_version: LOG_SCHEMA_VERSION.to_string(),
1389            canonical_digest: digest.clone(),
1390            prev_hash: prev_hash.clone(),
1391            this_hash: this_hash.clone(),
1392        };
1393        prev_hash = this_hash;
1394        out_rows.push(row);
1395    }
1396
1397    // -- 4. Build the report. --------------------------------------------
1398    let first_v1_hash = seeds
1399        .first()
1400        .map(|s| s.stored_this_hash.clone())
1401        .unwrap_or_else(|| GENESIS_HASH.to_string());
1402    let first_v2_hash = out_rows
1403        .first()
1404        .map(|r| r.this_hash.clone())
1405        .unwrap_or_else(|| GENESIS_HASH.to_string());
1406    let report = MigrationReport {
1407        rows_rewritten: out_rows.len() as u64,
1408        dry_run,
1409        first_row_v1_chain_hash: first_v1_hash,
1410        first_row_v2_chain_hash: first_v2_hash,
1411    };
1412
1413    if dry_run {
1414        return Ok(report);
1415    }
1416
1417    // -- 5. Live write: stage to `<log_path>.v2-migrated`, back up the
1418    //       v1, then atomically rename. -----------------------------------
1419    let staged_path = with_suffix(log_path, ".v2-migrated");
1420    let backup_path = with_suffix(log_path, ".v1-backup");
1421
1422    {
1423        let staged_file = OpenOptions::new()
1424            .create(true)
1425            .write(true)
1426            .truncate(true)
1427            .open(&staged_path)?;
1428        let mut writer = BufWriter::new(staged_file);
1429        for row in &out_rows {
1430            let mut bytes = serde_json::to_vec(row)?;
1431            bytes.push(b'\n');
1432            writer.write_all(&bytes)?;
1433        }
1434        writer.flush()?;
1435        let file = writer.into_inner().map_err(|e| {
1436            LogError::Io(std::io::Error::other(format!(
1437                "migration buf writer flush failed: {}",
1438                e.error()
1439            )))
1440        })?;
1441        file.sync_all()?;
1442    }
1443
1444    // Sanity-check: the staged file MUST verify clean before we
1445    // commit the swap. If it doesn't, the migration is buggy — abort
1446    // without touching the live log.
1447    let verify_report = verify(&staged_path)?;
1448    if !verify_report.errors.is_empty() {
1449        return Err(LogError::Io(std::io::Error::other(format!(
1450            "migration: staged v2 log failed verify; first issue: {:?}",
1451            verify_report.errors.first()
1452        ))));
1453    }
1454
1455    // Move the original aside as `<log_path>.v1-backup`. Overwriting
1456    // any prior backup is intentional — the user re-running migrate
1457    // expects the most recent original preserved.
1458    if log_path.exists() {
1459        if backup_path.exists() {
1460            std::fs::remove_file(&backup_path)?;
1461        }
1462        std::fs::rename(log_path, &backup_path)?;
1463    }
1464    // Atomically promote the staged file to the live path.
1465    std::fs::rename(&staged_path, log_path)?;
1466
1467    Ok(report)
1468}
1469
1470/// Append a literal suffix to a [`Utf8Path`], producing a sibling path
1471/// in the same directory. Avoids `std::path::PathBuf` per the workspace
1472/// posture rule (`docs/SECURITY.md` §3 — camino-only file paths in
1473/// production code).
1474fn with_suffix(path: &Utf8Path, suffix: &str) -> Utf8PathBuf {
1475    let s = format!("{path}{suffix}");
1476    Utf8PathBuf::from(s)
1477}
1478
1479// ---------------------------------------------------------------------------
1480// Tests
1481// ---------------------------------------------------------------------------
1482
1483#[cfg(test)]
1484#[allow(clippy::expect_used, clippy::unwrap_used, clippy::panic)]
1485mod tests {
1486    use super::*;
1487    use std::fs;
1488    use std::sync::Arc;
1489    use std::thread;
1490
1491    use tempfile::TempDir;
1492
1493    /// Convert a `TempDir`'s `&std::path::Path` to a `Utf8PathBuf`. Tests
1494    /// always run on UTF-8 temp paths in CI; if the OS returns a non-UTF-8
1495    /// path we panic, which is acceptable for a unit test.
1496    fn tmp_dir_utf8(dir: &TempDir) -> Utf8PathBuf {
1497        Utf8PathBuf::from_path_buf(dir.path().to_path_buf()).expect("temp dir path must be UTF-8")
1498    }
1499
1500    /// A fixed 26-char ULID-shaped string used in tests. Real callers use
1501    /// the `ulid` crate; tests pin a constant so output is reproducible.
1502    const TEST_SESSION_ID: &str = "01JCKZ7Q0000000000000000AB";
1503
1504    fn open_log(path: &Utf8Path) -> ProvenanceLog {
1505        ProvenanceLog::open(path, TEST_SESSION_ID.to_string()).expect("open")
1506    }
1507
1508    #[test]
1509    fn open_creates_missing_parent_dir() {
1510        // Regression: opening a log whose parent dir does not yet exist must
1511        // create the dir and succeed (then a row appends cleanly), not abort
1512        // with ENOENT. This is the `doiget verify` failure on a fresh CI
1513        // runner where `<config>/doiget/` was never created.
1514        let dir = TempDir::new().expect("tempdir");
1515        let path = tmp_dir_utf8(&dir)
1516            .join("nested")
1517            .join("doiget")
1518            .join("access.jsonl");
1519        assert!(
1520            !path.parent().expect("has parent").exists(),
1521            "parent dir must not pre-exist for this test to be meaningful"
1522        );
1523        let log = ProvenanceLog::open(&path, TEST_SESSION_ID.to_string())
1524            .expect("open must create the parent dir and succeed");
1525        log.append(empty_input())
1526            .expect("append after auto-created dir");
1527        assert!(path.exists(), "log file written under the auto-created dir");
1528        // End-to-end: the row the verify path would write is actually
1529        // readable back (exercises the full OpenOptions/flush/sync write,
1530        // not just that the file exists).
1531        let rows = read_rows(&path);
1532        assert_eq!(rows.len(), 1, "exactly one row in the auto-created log");
1533    }
1534
1535    fn empty_input() -> RowInput<'static> {
1536        RowInput {
1537            event: LogEvent::Fetch,
1538            result: LogResult::Ok,
1539            capability: Capability::Oa,
1540            ref_: None,
1541            source: None,
1542            error_code: None,
1543            size_bytes: None,
1544            license: None,
1545            store_path: None,
1546            canonical_digest: None,
1547        }
1548    }
1549
1550    /// Read the on-disk log and parse every line into a `LogRow`.
1551    fn read_rows(path: &Utf8Path) -> Vec<LogRow> {
1552        let raw = fs::read_to_string(path).expect("read log");
1553        raw.lines()
1554            .filter(|l| !l.is_empty())
1555            .map(|l| serde_json::from_str::<LogRow>(l).expect("valid LogRow"))
1556            .collect()
1557    }
1558
1559    /// Recompute `this_hash` for a stored row and assert it matches the
1560    /// stored value. Walks the same canonicalization rule as
1561    /// [`compute_this_hash`].
1562    fn verify_this_hash(row: &LogRow) {
1563        let rfh = RowForHash {
1564            ts: row.ts,
1565            ts_seq: row.ts_seq,
1566            event: row.event,
1567            ref_: row.ref_.as_deref(),
1568            source: row.source.as_deref(),
1569            result: row.result,
1570            license: row.license.as_deref(),
1571            size_bytes: row.size_bytes,
1572            store_path: row.store_path.as_deref(),
1573            capability: row.capability,
1574            session_id: &row.session_id,
1575            error_code: row.error_code.as_deref(),
1576            schema_version: &row.schema_version,
1577            canonical_digest: row.canonical_digest.as_deref(),
1578            prev_hash: &row.prev_hash,
1579        };
1580        let recomputed = compute_this_hash(&rfh).expect("hash");
1581        assert_eq!(
1582            recomputed, row.this_hash,
1583            "this_hash mismatch on ts_seq {}",
1584            row.ts_seq
1585        );
1586    }
1587
1588    #[test]
1589    fn first_row_uses_genesis_prev_hash() {
1590        let dir = TempDir::new().expect("tmp");
1591        let path = tmp_dir_utf8(&dir).join("log.jsonl");
1592        let log = open_log(&path);
1593        let seq = log.append(empty_input()).expect("append");
1594        assert_eq!(seq, 1);
1595
1596        let rows = read_rows(&path);
1597        assert_eq!(rows.len(), 1);
1598        assert_eq!(rows[0].ts_seq, 1);
1599        assert_eq!(rows[0].prev_hash, GENESIS_HASH);
1600        assert_eq!(rows[0].this_hash.len(), 64);
1601        assert_eq!(rows[0].session_id, TEST_SESSION_ID);
1602        verify_this_hash(&rows[0]);
1603    }
1604
1605    #[test]
1606    fn subsequent_rows_chain_correctly() {
1607        let dir = TempDir::new().expect("tmp");
1608        let path = tmp_dir_utf8(&dir).join("log.jsonl");
1609        let log = open_log(&path);
1610
1611        for _ in 0..3 {
1612            log.append(empty_input()).expect("append");
1613        }
1614
1615        let rows = read_rows(&path);
1616        assert_eq!(rows.len(), 3);
1617        assert_eq!(rows[0].prev_hash, GENESIS_HASH);
1618        assert_eq!(rows[1].prev_hash, rows[0].this_hash);
1619        assert_eq!(rows[2].prev_hash, rows[1].this_hash);
1620        for r in &rows {
1621            verify_this_hash(r);
1622        }
1623        assert_eq!(rows[0].ts_seq, 1);
1624        assert_eq!(rows[1].ts_seq, 2);
1625        assert_eq!(rows[2].ts_seq, 3);
1626    }
1627
1628    #[test]
1629    fn recovery_after_reopen() {
1630        let dir = TempDir::new().expect("tmp");
1631        let path = tmp_dir_utf8(&dir).join("log.jsonl");
1632
1633        {
1634            let log = open_log(&path);
1635            for _ in 0..3 {
1636                log.append(empty_input()).expect("append");
1637            }
1638        } // drop writer
1639
1640        let log2 = open_log(&path);
1641        let seq = log2.append(empty_input()).expect("append after reopen");
1642        assert_eq!(seq, 4);
1643
1644        let rows = read_rows(&path);
1645        assert_eq!(rows.len(), 4);
1646        assert_eq!(rows[0].prev_hash, GENESIS_HASH);
1647        for i in 1..rows.len() {
1648            assert_eq!(
1649                rows[i].prev_hash,
1650                rows[i - 1].this_hash,
1651                "chain break at row {}",
1652                i + 1
1653            );
1654        }
1655        for (i, r) in rows.iter().enumerate() {
1656            assert_eq!(r.ts_seq, (i + 1) as u64);
1657            verify_this_hash(r);
1658        }
1659    }
1660
1661    #[test]
1662    fn concurrent_writers_in_same_process_serialize() {
1663        let dir = TempDir::new().expect("tmp");
1664        let path = tmp_dir_utf8(&dir).join("log.jsonl");
1665        let log = Arc::new(open_log(&path));
1666
1667        let mut handles = Vec::with_capacity(8);
1668        for _ in 0..8 {
1669            let log = Arc::clone(&log);
1670            handles.push(thread::spawn(move || {
1671                log.append(empty_input()).expect("append")
1672            }));
1673        }
1674        let mut returned: Vec<u64> = handles
1675            .into_iter()
1676            .map(|h| h.join().expect("join"))
1677            .collect();
1678        returned.sort_unstable();
1679        assert_eq!(returned, vec![1, 2, 3, 4, 5, 6, 7, 8]);
1680
1681        let rows = read_rows(&path);
1682        assert_eq!(rows.len(), 8);
1683
1684        // The in-process mutex serializes appends, so file order MUST equal
1685        // ts_seq order: row N (0-indexed) on disk has ts_seq = N+1.
1686        for (i, r) in rows.iter().enumerate() {
1687            assert_eq!(r.ts_seq, (i + 1) as u64, "ts_seq gap at file row {}", i + 1);
1688        }
1689        // Hash chain follows file order.
1690        assert_eq!(rows[0].prev_hash, GENESIS_HASH);
1691        for i in 1..rows.len() {
1692            assert_eq!(
1693                rows[i].prev_hash,
1694                rows[i - 1].this_hash,
1695                "chain break at file row {}",
1696                i + 1
1697            );
1698        }
1699        for r in &rows {
1700            verify_this_hash(r);
1701        }
1702    }
1703
1704    #[test]
1705    fn corrupted_existing_log_fails_open() {
1706        let dir = TempDir::new().expect("tmp");
1707        let path = tmp_dir_utf8(&dir).join("log.jsonl");
1708
1709        // JSON but not a valid LogRow: missing required fields, has unknown
1710        // field. `deny_unknown_fields` ensures the parser refuses.
1711        fs::write(&path, "{\"ts_seq\": 1, \"garbage\": true}\n").expect("write");
1712
1713        let err =
1714            ProvenanceLog::open(&path, TEST_SESSION_ID.to_string()).expect_err("must fail open");
1715        match err {
1716            LogError::Io(io) => {
1717                let msg = io.to_string();
1718                assert!(
1719                    msg.contains("corrupted log at line 1"),
1720                    "expected synthetic corruption message, got: {}",
1721                    msg
1722                );
1723            }
1724            other => panic!("expected LogError::Io, got {:?}", other),
1725        }
1726    }
1727
1728    #[test]
1729    fn rejects_non_regular_file() {
1730        // Pointing the log at a directory must fail with NotARegularFile.
1731        let dir = TempDir::new().expect("tmp");
1732        let err = ProvenanceLog::open(tmp_dir_utf8(&dir), TEST_SESSION_ID.to_string())
1733            .expect_err("must fail");
1734        match err {
1735            LogError::NotARegularFile(_) => {}
1736            other => panic!("expected NotARegularFile, got {:?}", other),
1737        }
1738    }
1739
1740    #[test]
1741    fn canonical_json_excludes_this_hash_field() {
1742        // Spec contract: the hashed bytes do not include `this_hash`. If
1743        // this ever regresses, every previously-written log becomes
1744        // unverifiable.
1745        let rfh = RowForHash {
1746            ts: Utc::now(),
1747            ts_seq: 1,
1748            event: LogEvent::Fetch,
1749            ref_: None,
1750            source: None,
1751            result: LogResult::Ok,
1752            license: None,
1753            size_bytes: None,
1754            store_path: None,
1755            capability: Capability::Oa,
1756            session_id: TEST_SESSION_ID,
1757            error_code: None,
1758            schema_version: LOG_SCHEMA_VERSION,
1759            canonical_digest: None,
1760            prev_hash: GENESIS_HASH,
1761        };
1762        let bytes = canonical_json_for_hash(&rfh).expect("canonicalize");
1763        let s = std::str::from_utf8(&bytes).expect("utf8");
1764        assert!(!s.contains("this_hash"), "this_hash leaked into hash input");
1765        assert!(s.contains("\"prev_hash\":"));
1766    }
1767
1768    #[test]
1769    fn canonical_json_keys_are_lexicographically_sorted() {
1770        // PROVENANCE_LOG.md §4: canonical JSON uses keys sorted
1771        // lexicographically. The lex-first top-level key of a row is
1772        // `capability` ("c..." < "e..." < ...). Build a row and assert the
1773        // canonical bytes start with that key.
1774        let rfh = RowForHash {
1775            ts: Utc::now(),
1776            ts_seq: 1,
1777            event: LogEvent::Fetch,
1778            ref_: Some("10.1234/example"),
1779            source: Some("unpaywall"),
1780            result: LogResult::Ok,
1781            license: Some("CC-BY-4.0"),
1782            size_bytes: Some(1234),
1783            store_path: Some("papers/x.pdf"),
1784            capability: Capability::Oa,
1785            session_id: TEST_SESSION_ID,
1786            error_code: None,
1787            schema_version: LOG_SCHEMA_VERSION,
1788            canonical_digest: Some(
1789                "0000000000000000000000000000000000000000000000000000000000000000",
1790            ),
1791            prev_hash: GENESIS_HASH,
1792        };
1793        let bytes = canonical_json_for_hash(&rfh).expect("canonicalize");
1794        let s = std::str::from_utf8(&bytes).expect("utf8");
1795        // v2: lex-first key is `canonical_digest` (< `capability` because
1796        // 'n' < 'p' at byte index 2). Pre-v2 it was `capability`.
1797        assert!(
1798            s.starts_with("{\"canonical_digest\":"),
1799            "canonical bytes must start with lex-first v2 key, got: {}",
1800            s
1801        );
1802        // Spot-check ordering: `prev_hash` (p) must come before `ref` (r),
1803        // which must come before `result` (re...) — wait, "ref" < "result"
1804        // lexicographically because 'f' < 's' in ascii at index 2 vs 'e' at
1805        // index 2 of "result"... let me just check a couple of unambiguous
1806        // pairs: `event` < `prev_hash`, and `ts` < `ts_seq`.
1807        let event_idx = s.find("\"event\":").expect("event key present");
1808        let prev_idx = s.find("\"prev_hash\":").expect("prev_hash key present");
1809        assert!(event_idx < prev_idx, "event must precede prev_hash");
1810        let ts_idx = s.find("\"ts\":").expect("ts key present");
1811        let tsseq_idx = s.find("\"ts_seq\":").expect("ts_seq key present");
1812        assert!(ts_idx < tsseq_idx, "ts must precede ts_seq");
1813    }
1814
1815    // -----------------------------------------------------------------
1816    // verify() tests — Phase 1 surface for `doiget audit-log --verify`.
1817    // -----------------------------------------------------------------
1818
1819    /// Rewrite a single field's quoted-string value on a specific 1-based
1820    /// line of `path`. Used to simulate tampering. Panics on malformed input
1821    /// — only valid inputs are produced by the test harness.
1822    ///
1823    /// `field_key` is matched as `"field_key":"...old..."` (quoted string
1824    /// JSON value). The new value is the literal string `new_value` (no
1825    /// JSON escaping needed for the test fixtures we use).
1826    fn tamper_string_field(
1827        path: &Utf8Path,
1828        line_no_1based: usize,
1829        field_key: &str,
1830        new_value: &str,
1831    ) {
1832        let raw = fs::read_to_string(path).expect("read log");
1833        let mut lines: Vec<String> = raw.lines().map(str::to_string).collect();
1834        let target = &lines[line_no_1based - 1];
1835        let needle = format!("\"{field_key}\":\"");
1836        let start = target
1837            .find(&needle)
1838            .unwrap_or_else(|| panic!("field {field_key} not found on line {line_no_1based}"))
1839            + needle.len();
1840        let end_rel = target[start..]
1841            .find('"')
1842            .unwrap_or_else(|| panic!("unterminated string for field {field_key}"));
1843        let end = start + end_rel;
1844        let mut new_line = String::with_capacity(target.len());
1845        new_line.push_str(&target[..start]);
1846        new_line.push_str(new_value);
1847        new_line.push_str(&target[end..]);
1848        lines[line_no_1based - 1] = new_line;
1849        let mut out = lines.join("\n");
1850        out.push('\n');
1851        fs::write(path, out).expect("write tampered log");
1852    }
1853
1854    #[test]
1855    fn verify_empty_log_is_ok() {
1856        // Missing file is a clean log — no tampering possible on bytes that
1857        // don't exist. `verify` returns an empty report, not an error.
1858        let dir = TempDir::new().expect("tmp");
1859        let path = tmp_dir_utf8(&dir).join("nonexistent.jsonl");
1860        assert!(!path.exists(), "precondition: file must not exist");
1861
1862        let report = verify(&path).expect("verify must not error on missing file");
1863        assert_eq!(report.total_rows, 0);
1864        assert_eq!(report.ok_rows, 0);
1865        assert!(report.errors.is_empty(), "errors: {:?}", report.errors);
1866    }
1867
1868    #[test]
1869    fn verify_well_formed_chain_passes() {
1870        // Three rows written via the real writer must verify clean.
1871        let dir = TempDir::new().expect("tmp");
1872        let path = tmp_dir_utf8(&dir).join("log.jsonl");
1873        let log = open_log(&path);
1874        for _ in 0..3 {
1875            log.append(empty_input()).expect("append");
1876        }
1877
1878        let report = verify(&path).expect("verify must succeed");
1879        assert_eq!(report.total_rows, 3);
1880        assert_eq!(report.ok_rows, 3);
1881        assert!(
1882            report.errors.is_empty(),
1883            "expected no issues on a well-formed log; got: {:?}",
1884            report.errors
1885        );
1886    }
1887
1888    #[test]
1889    fn verify_detects_tampered_row_hash() {
1890        // Mutate the SECOND row's `this_hash` to a syntactically-valid but
1891        // wrong hash. The recomputed canonical-JSON SHA-256 will not match.
1892        let dir = TempDir::new().expect("tmp");
1893        let path = tmp_dir_utf8(&dir).join("log.jsonl");
1894        let log = open_log(&path);
1895        log.append(empty_input()).expect("append 1");
1896        log.append(empty_input()).expect("append 2");
1897        drop(log);
1898
1899        // 64 lowercase hex chars, all zeros — passes `LogRow` parse, fails hash check.
1900        tamper_string_field(
1901            &path,
1902            2,
1903            "this_hash",
1904            "0000000000000000000000000000000000000000000000000000000000000000",
1905        );
1906
1907        let report = verify(&path).expect("verify must succeed");
1908        assert_eq!(report.total_rows, 2);
1909        // Row 2's hash mismatch breaks both the hash check on row 2 AND the
1910        // chain link from row 2's stored `prev_hash` (still correct) into the
1911        // forward direction. There's no row 3 to fail forward, so we expect
1912        // exactly one issue: the this-hash mismatch on line 2.
1913        let hash_issues: Vec<_> = report
1914            .errors
1915            .iter()
1916            .filter(|e| e.kind == VerifyIssueKind::ThisHashMismatch)
1917            .collect();
1918        assert_eq!(
1919            hash_issues.len(),
1920            1,
1921            "expected exactly one ThisHashMismatch, got {:?}",
1922            report.errors
1923        );
1924        assert_eq!(hash_issues[0].line, 2);
1925    }
1926
1927    #[test]
1928    fn verify_detects_tampered_prev_hash() {
1929        // Mutate the SECOND row's `prev_hash` to a wrong value. This
1930        // invalidates the chain link but the row's own `this_hash` was
1931        // computed with the original `prev_hash`, so the this-hash check
1932        // ALSO fails (hash input changed). We assert at least the prev-hash
1933        // issue is reported on line 2.
1934        let dir = TempDir::new().expect("tmp");
1935        let path = tmp_dir_utf8(&dir).join("log.jsonl");
1936        let log = open_log(&path);
1937        log.append(empty_input()).expect("append 1");
1938        log.append(empty_input()).expect("append 2");
1939        drop(log);
1940
1941        tamper_string_field(
1942            &path,
1943            2,
1944            "prev_hash",
1945            "ffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffffff",
1946        );
1947
1948        let report = verify(&path).expect("verify must succeed");
1949        assert_eq!(report.total_rows, 2);
1950        let prev_issues: Vec<_> = report
1951            .errors
1952            .iter()
1953            .filter(|e| e.kind == VerifyIssueKind::PrevHashMismatch)
1954            .collect();
1955        assert_eq!(
1956            prev_issues.len(),
1957            1,
1958            "expected exactly one PrevHashMismatch, got {:?}",
1959            report.errors
1960        );
1961        assert_eq!(prev_issues[0].line, 2);
1962    }
1963
1964    #[test]
1965    fn verify_detects_corrupted_json() {
1966        // One valid row plus a literal `{"garbage":true}` line. The garbage
1967        // line fails `serde_json::from_str::<LogRow>` (missing fields +
1968        // `deny_unknown_fields`) and surfaces as a `ParseError` on line 2.
1969        let dir = TempDir::new().expect("tmp");
1970        let path = tmp_dir_utf8(&dir).join("log.jsonl");
1971        let log = open_log(&path);
1972        log.append(empty_input()).expect("append 1");
1973        drop(log);
1974
1975        // Append a garbage line directly.
1976        let mut existing = fs::read_to_string(&path).expect("read");
1977        if !existing.ends_with('\n') {
1978            existing.push('\n');
1979        }
1980        existing.push_str("{\"garbage\":true}\n");
1981        fs::write(&path, existing).expect("write");
1982
1983        let report = verify(&path).expect("verify must succeed");
1984        // total_rows counts non-empty lines, so both lines are counted.
1985        assert_eq!(report.total_rows, 2);
1986        let parse_issues: Vec<_> = report
1987            .errors
1988            .iter()
1989            .filter(|e| e.kind == VerifyIssueKind::ParseError)
1990            .collect();
1991        assert_eq!(
1992            parse_issues.len(),
1993            1,
1994            "expected exactly one ParseError, got {:?}",
1995            report.errors
1996        );
1997        assert_eq!(parse_issues[0].line, 2);
1998    }
1999
2000    #[test]
2001    fn capability_serializes_kebab_case() {
2002        // PROVENANCE_LOG.md §3 requires `oa`, `metadata`, `tdm-elsevier`,
2003        // `tdm-aps`, `tdm-springer` on the wire (kebab-case).
2004        let cases = [
2005            (Capability::Oa, "\"oa\""),
2006            (Capability::Metadata, "\"metadata\""),
2007            (Capability::TdmElsevier, "\"tdm-elsevier\""),
2008            (Capability::TdmAps, "\"tdm-aps\""),
2009            (Capability::TdmSpringer, "\"tdm-springer\""),
2010            (Capability::TdmIeee, "\"tdm-ieee\""),
2011        ];
2012        for (cap, expected) in cases {
2013            let got = serde_json::to_string(&cap).expect("serialize");
2014            assert_eq!(
2015                got, expected,
2016                "capability wire format mismatch for {:?}",
2017                cap
2018            );
2019        }
2020    }
2021
2022    // -----------------------------------------------------------------
2023    // #140 — §6 rotation, retention, multi-segment verify.
2024    // -----------------------------------------------------------------
2025
2026    fn gunzip_to_string(gz: &Utf8Path) -> String {
2027        use std::io::Read;
2028        let f = std::fs::File::open(gz.as_std_path()).expect("open gz");
2029        let mut dec = GzDecoder::new(f);
2030        let mut s = String::new();
2031        dec.read_to_string(&mut s).expect("gunzip");
2032        s
2033    }
2034
2035    #[test]
2036    fn rotation_archives_to_gz_and_restarts_genesis_chain() {
2037        let dir = TempDir::new().expect("tmp");
2038        let path = tmp_dir_utf8(&dir).join("access.log");
2039        // Inject a tiny threshold (NOT a global env var — that raced
2040        // non-#[serial] tests): row 1 fits, so the SECOND append
2041        // (size>=50) rotates before it writes. A freshly rotated `.gz`
2042        // is not retention-aged, so the default prune at open is a no-op.
2043        let log = ProvenanceLog::open_with_rotate_threshold(&path, TEST_SESSION_ID.to_string(), 50)
2044            .expect("open");
2045        log.append(empty_input()).expect("append 1");
2046        let row1 = read_rows(&path);
2047        assert_eq!(row1.len(), 1);
2048        assert_eq!(row1[0].prev_hash, GENESIS_HASH);
2049
2050        log.append(empty_input()).expect("append 2 (rotates first)");
2051
2052        // Exactly one rotated segment; it gunzips to the original row 1.
2053        let segs = rotated_segments(&path);
2054        assert_eq!(segs.len(), 1, "one .gz segment expected; got {segs:?}");
2055        let archived: Vec<LogRow> = gunzip_to_string(&segs[0])
2056            .lines()
2057            .filter(|l| !l.is_empty())
2058            .map(|l| serde_json::from_str(l).expect("row"))
2059            .collect();
2060        assert_eq!(archived.len(), 1);
2061        assert_eq!(archived[0].this_hash, row1[0].this_hash);
2062
2063        // The fresh access.log restarts the chain at GENESIS, ts_seq 1.
2064        let cur = read_rows(&path);
2065        assert_eq!(cur.len(), 1, "fresh segment holds only the post-rotate row");
2066        assert_eq!(cur[0].prev_hash, GENESIS_HASH);
2067        assert_eq!(cur[0].ts_seq, 1);
2068
2069        // verify_all sees both segments, each its own clean chain.
2070        let reports = verify_all(&path).expect("verify_all");
2071        assert_eq!(reports.len(), 2, "rotated .gz + current");
2072        for (p, r) in &reports {
2073            assert!(r.errors.is_empty(), "segment {p} must verify clean: {r:?}");
2074        }
2075    }
2076
2077    #[test]
2078    fn rotate_log_is_fail_closed_on_missing_source() {
2079        // The append path propagates this via `?`, so a rotation failure
2080        // aborts the fetch (fail-closed) rather than silently continuing.
2081        let dir = TempDir::new().expect("tmp");
2082        let missing = tmp_dir_utf8(&dir).join("nope.log");
2083        let err = rotate_log(&missing).expect_err("missing source must error");
2084        assert!(matches!(err, LogError::Io(_)), "got {err:?}");
2085    }
2086
2087    #[test]
2088    #[serial_test::serial]
2089    fn prune_respects_retention_window_and_disable() {
2090        let dir = TempDir::new().expect("tmp");
2091        let base = tmp_dir_utf8(&dir);
2092        let path = base.join("access.log");
2093        let old_gz = base.join("access.log.2020-01-01-000000.gz");
2094        let new_gz = base.join("access.log.2999-01-01-000000.gz");
2095
2096        let mk = |p: &Utf8Path, aged: bool| {
2097            let f = std::fs::File::create(p.as_std_path()).expect("create gz");
2098            if aged {
2099                // 100 days ago — older than the 90-day default & a 1-day window.
2100                let when =
2101                    std::time::SystemTime::now() - std::time::Duration::from_secs(100 * 86_400);
2102                f.set_modified(when).expect("set mtime");
2103            }
2104        };
2105
2106        // (a) days=0 disables pruning entirely.
2107        mk(&old_gz, true);
2108        std::env::set_var("DOIGET_LOG_RETENTION_DAYS", "0");
2109        let _ = open_log(&path);
2110        assert!(old_gz.exists(), "days=0 must NOT prune");
2111
2112        // (b) days=1 prunes the aged segment, keeps a fresh one.
2113        mk(&new_gz, false);
2114        std::env::set_var("DOIGET_LOG_RETENTION_DAYS", "1");
2115        let _ = open_log(&path);
2116        assert!(!old_gz.exists(), "aged segment must be pruned at days=1");
2117        assert!(new_gz.exists(), "fresh segment must survive");
2118
2119        std::env::remove_var("DOIGET_LOG_RETENTION_DAYS");
2120    }
2121
2122    #[test]
2123    #[serial_test::serial]
2124    fn retention_days_env_parsing() {
2125        std::env::set_var("DOIGET_LOG_RETENTION_DAYS", "0");
2126        assert_eq!(retention_days(), 0);
2127        std::env::set_var("DOIGET_LOG_RETENTION_DAYS", "30");
2128        assert_eq!(retention_days(), 30);
2129        std::env::set_var("DOIGET_LOG_RETENTION_DAYS", "garbage");
2130        assert_eq!(retention_days(), DEFAULT_RETENTION_DAYS);
2131        std::env::set_var("DOIGET_LOG_RETENTION_DAYS", "-5");
2132        assert_eq!(retention_days(), DEFAULT_RETENTION_DAYS);
2133        std::env::remove_var("DOIGET_LOG_RETENTION_DAYS");
2134        assert_eq!(retention_days(), DEFAULT_RETENTION_DAYS);
2135    }
2136
2137    #[test]
2138    fn verify_all_flags_tampered_segment_independently() {
2139        let dir = TempDir::new().expect("tmp");
2140        let path = tmp_dir_utf8(&dir).join("access.log");
2141        // Inject the tiny threshold (no global env → no cross-test race).
2142        let log = ProvenanceLog::open_with_rotate_threshold(&path, TEST_SESSION_ID.to_string(), 50)
2143            .expect("open");
2144        log.append(empty_input()).expect("append 1");
2145        log.append(empty_input()).expect("append 2 (rotates)");
2146        drop(log);
2147
2148        // Tamper the CURRENT segment's row: set this_hash to a
2149        // syntactically-valid 64-hex string that cannot be the SHA-256
2150        // of any row (all zeros). NOTE: the previous "flip the last char
2151        // to '0'" was a no-op ~1/16 of runs when the real hash already
2152        // ended in '0' (this_hash depends on `Utc::now()`), which is the
2153        // flake this fixes — mirrors `verify_detects_tampered_row_hash`.
2154        let mut cur = read_rows(&path);
2155        let mut bad = cur.remove(0);
2156        bad.this_hash =
2157            "0000000000000000000000000000000000000000000000000000000000000000".to_string();
2158        std::fs::write(
2159            path.as_std_path(),
2160            format!("{}\n", serde_json::to_string(&bad).expect("ser")),
2161        )
2162        .expect("rewrite tampered current");
2163
2164        let reports = verify_all(&path).expect("verify_all");
2165        assert_eq!(reports.len(), 2);
2166        // Oldest first = the rotated .gz (clean); current last (tampered).
2167        let (gz_path, gz_rep) = &reports[0];
2168        let (cur_path, cur_rep) = &reports[1];
2169        assert!(
2170            gz_path.as_str().ends_with(".gz") && gz_rep.errors.is_empty(),
2171            "rotated segment must stay clean: {gz_path} {gz_rep:?}"
2172        );
2173        assert!(
2174            cur_path.file_name() == Some("access.log") && !cur_rep.errors.is_empty(),
2175            "tampered current segment must report issues: {cur_path} {cur_rep:?}"
2176        );
2177    }
2178}