Skip to main content

treeship_core/journal/
mod.rs

1//! Local Approval Use Journal -- v0.9.9 PR 2.
2//!
3//! Per-workspace append-only memory of consumed Approval Grants. The
4//! journal turns the v0.9.6 "package-local only" replay finding into a
5//! local-journal replay finding: with this module wired through, verify
6//! can say "use 1/1 -- local Approval Use Journal passed" instead of
7//! "no global ledger consulted."
8//!
9//! Scope of THIS PR:
10//!   * journal storage (records/, heads/, indexes/, locks/)
11//!   * append-only writes with file lock + atomic temp+rename
12//!   * hash chain via `previous_record_digest`
13//!   * read-only `check_replay` lookup
14//!   * `verify_integrity` chain walk
15//!   * `rebuild_indexes` from records (records are truth)
16//!
17//! Out of scope (later PRs):
18//!   * consume-before-action wiring inside `treeship attest action` (PR 3)
19//!   * package export of journal records (PR 4)
20//!   * Hub checkpoint signing (PR 6 scaffold)
21//!
22//! Privacy rules baked into the layout:
23//!   * `nonce_digest`, never raw nonce
24//!   * no commands, prompts, file contents, bearer tokens, or API keys
25//!     are stored. The journal answers the single question "has this
26//!     (grant_id, nonce_digest) been consumed before, and if so how
27//!     many times?" -- everything else stays in the signed grant +
28//!     receipt where it already is.
29
30use std::fs::{self, File, OpenOptions};
31use std::io::Write;
32use std::path::{Path, PathBuf};
33
34// fs2 is gated to non-wasm targets at the workspace Cargo.toml; the WASM
35// build has no concurrent writers and no real filesystem, so journal
36// operations fall back to a deterministic "no-op write" mode that still
37// keeps the public API building. Same pattern session::event_log uses.
38#[cfg(not(target_family = "wasm"))]
39use fs2::FileExt;
40
41use crate::statements::{
42    approval_revocation_record_digest, approval_use_record_digest,
43    journal_checkpoint_record_digest, ApprovalRevocation, ApprovalUse, JournalCheckpoint,
44    ReplayCheck, ReplayCheckLevel, TYPE_APPROVAL_REVOCATION, TYPE_APPROVAL_USE,
45    TYPE_JOURNAL_CHECKPOINT,
46};
47
48// ---------------------------------------------------------------------------
49// Errors
50// ---------------------------------------------------------------------------
51
52#[derive(Debug)]
53pub enum JournalError {
54    Io(std::io::Error),
55    Json(serde_json::Error),
56    /// `previous_record_digest` on a record didn't match the prior
57    /// record's `record_digest`. The chain is broken.
58    BrokenChain {
59        index: u64,
60        expected: String,
61        actual: String,
62    },
63    /// A record's stored `record_digest` didn't match the recomputed
64    /// digest. The record was tampered after write.
65    RecordTampered {
66        index: u64,
67        expected: String,
68        actual: String,
69    },
70    /// A record file referenced by the head no longer exists.
71    MissingRecord {
72        index: u64,
73    },
74    /// The journal's append lock could not be acquired.
75    LockBusy,
76    /// The append exceeds `max_uses` recorded on prior uses for this
77    /// grant. Surfaced as an error so callers (PR 3) refuse to sign
78    /// the action; PR 2 itself only writes uses passed in by callers,
79    /// so this only fires from `append_use` when the caller didn't
80    /// preflight via `check_replay`.
81    MaxUsesExceeded {
82        grant_id: String,
83        max_uses: u32,
84        current: u32,
85    },
86}
87
88impl std::fmt::Display for JournalError {
89    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
90        match self {
91            Self::Io(e)            => write!(f, "journal io: {e}"),
92            Self::Json(e)          => write!(f, "journal json: {e}"),
93            Self::BrokenChain { index, expected, actual } => write!(
94                f,
95                "journal broken at record {index}: previous_record_digest = {actual}, expected {expected}",
96            ),
97            Self::RecordTampered { index, expected, actual } => write!(
98                f,
99                "journal record {index} tampered: stored digest {expected}, recomputed {actual}",
100            ),
101            Self::MissingRecord { index } => write!(
102                f,
103                "journal record {index} referenced by head but missing on disk",
104            ),
105            Self::LockBusy => write!(f, "journal append lock busy; another process holds it"),
106            Self::MaxUsesExceeded { grant_id, max_uses, current } => write!(
107                f,
108                "approval grant {grant_id} would exceed max_uses ({current}/{max_uses})",
109            ),
110        }
111    }
112}
113
114impl std::error::Error for JournalError {}
115impl From<std::io::Error> for JournalError {
116    fn from(e: std::io::Error) -> Self {
117        Self::Io(e)
118    }
119}
120impl From<serde_json::Error> for JournalError {
121    fn from(e: serde_json::Error) -> Self {
122        Self::Json(e)
123    }
124}
125
126// ---------------------------------------------------------------------------
127// Layout
128// ---------------------------------------------------------------------------
129
130/// Directory layout under `.treeship/journals/approval-use/`.
131pub struct Journal {
132    /// Root directory.
133    pub dir: PathBuf,
134}
135
136impl Journal {
137    pub fn new(dir: impl Into<PathBuf>) -> Self {
138        Self { dir: dir.into() }
139    }
140
141    pub fn records_dir(&self) -> PathBuf {
142        self.dir.join("records")
143    }
144    pub fn heads_dir(&self) -> PathBuf {
145        self.dir.join("heads")
146    }
147    pub fn indexes_dir(&self) -> PathBuf {
148        self.dir.join("indexes")
149    }
150    pub fn locks_dir(&self) -> PathBuf {
151        self.dir.join("locks")
152    }
153    pub fn current_head_path(&self) -> PathBuf {
154        self.heads_dir().join("current.json")
155    }
156    pub fn lock_path(&self) -> PathBuf {
157        self.locks_dir().join("journal.lock")
158    }
159    pub fn meta_path(&self) -> PathBuf {
160        self.dir.join("journal.json")
161    }
162
163    /// Index file for a given grant. Each line is one `record_index`.
164    pub fn by_grant_path(&self, grant_id: &str) -> PathBuf {
165        self.indexes_dir()
166            .join("by-grant")
167            .join(format!("{}.txt", safe_name(grant_id)))
168    }
169
170    /// Index file for a nonce_digest.
171    pub fn by_nonce_path(&self, nonce_digest: &str) -> PathBuf {
172        self.indexes_dir()
173            .join("by-nonce")
174            .join(format!("{}.txt", safe_name(nonce_digest)))
175    }
176
177    /// Returns true iff the journal directory exists.
178    pub fn exists(&self) -> bool {
179        self.dir.is_dir()
180    }
181}
182
183/// Make a filesystem-safe name by replacing path-unsafe chars. Used for
184/// index file names; not a security boundary -- the journal's actual
185/// integrity check is the hash chain.
186fn safe_name(s: &str) -> String {
187    s.chars()
188        .map(|c| match c {
189            ':' | '/' | '\\' | ' ' | '.' => '_',
190            c => c,
191        })
192        .collect()
193}
194
195// ---------------------------------------------------------------------------
196// Head file
197// ---------------------------------------------------------------------------
198
199#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
200pub struct Head {
201    /// 1-indexed; 0 means "no records yet."
202    pub index: u64,
203    /// `record_digest` of the most recent record. Empty when index=0.
204    pub digest: String,
205    /// Updated on every append.
206    pub updated_at: String,
207}
208
209fn read_head(j: &Journal) -> Result<Head, JournalError> {
210    let path = j.current_head_path();
211    if !path.exists() {
212        return Ok(Head::default());
213    }
214    let bytes = fs::read(&path)?;
215    Ok(serde_json::from_slice(&bytes)?)
216}
217
218fn write_head(j: &Journal, head: &Head) -> Result<(), JournalError> {
219    fs::create_dir_all(j.heads_dir())?;
220    let path = j.current_head_path();
221    let tmp = path.with_extension("json.tmp");
222    let json = serde_json::to_vec_pretty(head)?;
223    fs::write(&tmp, json)?;
224    fs::rename(&tmp, &path)?;
225    Ok(())
226}
227
228// ---------------------------------------------------------------------------
229// Append
230// ---------------------------------------------------------------------------
231
232/// Acquire the journal append lock for the duration of the closure. Uses
233/// fs2::FileExt::try_lock_exclusive (the same primitive `session::event_log`
234/// uses) so behavior matches what the rest of the codebase already
235/// trusts.
236#[cfg(not(target_family = "wasm"))]
237fn with_lock<F, T>(j: &Journal, body: F) -> Result<T, JournalError>
238where
239    F: FnOnce() -> Result<T, JournalError>,
240{
241    fs::create_dir_all(j.locks_dir())?;
242    let lock = OpenOptions::new()
243        .read(true)
244        .write(true)
245        .create(true)
246        .truncate(false)
247        .open(j.lock_path())?;
248    if lock.try_lock_exclusive().is_err() {
249        return Err(JournalError::LockBusy);
250    }
251    let result = body();
252    let _ = fs2::FileExt::unlock(&lock);
253    result
254}
255
256/// WASM build: no concurrent writers, no advisory locks. Run the body
257/// directly. Matches `session::event_log`'s wasm fallback.
258#[cfg(target_family = "wasm")]
259fn with_lock<F, T>(_j: &Journal, body: F) -> Result<T, JournalError>
260where
261    F: FnOnce() -> Result<T, JournalError>,
262{
263    body()
264}
265
266/// Append an ApprovalUse to the journal. The caller MUST set
267/// `previous_record_digest` to the current head's digest on the
268/// incoming record; we re-validate before write. `record_digest` is
269/// computed from the canonical form and stamped on the stored record.
270///
271/// Returns the new head's index and digest.
272pub fn append_use(j: &Journal, mut rec: ApprovalUse) -> Result<Head, JournalError> {
273    rec.type_ = TYPE_APPROVAL_USE.into();
274    with_lock(j, || {
275        let head = read_head(j)?;
276        rec.previous_record_digest = head.digest.clone();
277        rec.record_digest = approval_use_record_digest(&rec);
278        let next_index = head.index + 1;
279        write_record_use(j, next_index, &rec)?;
280        update_indexes_for_use(j, next_index, &rec)?;
281        let new_head = Head {
282            index: next_index,
283            digest: rec.record_digest.clone(),
284            updated_at: rec.created_at.clone(),
285        };
286        write_head(j, &new_head)?;
287        ensure_meta(j)?;
288        Ok(new_head)
289    })
290}
291
292/// Atomic check-and-append for the consume path. Combines `check_replay` +
293/// append under a single journal lock so concurrent consume paths cannot
294/// bypass `max_uses` via TOCTOU race.
295///
296/// v0.9.9 PR 3 (`reserve_in_journal` in attest.rs) ran `check_replay` and
297/// derived `use_number` *outside* `with_lock`, then called `append_use`
298/// (which takes the lock only for the write). Two parallel attests could
299/// both pass the pre-lock replay check, then queue serially through the
300/// lock, and both would write — exceeding `max_uses=1`. v0.9.10 closes
301/// that race by doing the check *inside* the lock.
302///
303/// The function also stamps `use_number` from the grant-wide count
304/// observed at lock-acquire time. Callers should pass the record with
305/// `use_number = 0` (or any value); it will be overwritten.
306///
307/// Returns the new head on success. On replay violation, returns
308/// `JournalError::MaxUsesExceeded` and writes nothing — the lock is
309/// released without state change.
310pub fn reserve_use(
311    j: &Journal,
312    mut rec: ApprovalUse,
313    max_uses: Option<u32>,
314) -> Result<Head, JournalError> {
315    rec.type_ = TYPE_APPROVAL_USE.into();
316    with_lock(j, || {
317        // Replay check inside the lock. `check_replay` reads the
318        // by-nonce index; while we hold the exclusive lock, no other
319        // writer can mutate that index, so the count is correct.
320        let replay = check_replay(j, &rec.grant_id, &rec.nonce_digest, max_uses)?;
321        if let Some(false) = replay.passed {
322            let current = replay.use_number.map(|n| n.saturating_sub(1)).unwrap_or(0);
323            return Err(JournalError::MaxUsesExceeded {
324                grant_id: rec.grant_id.clone(),
325                max_uses: replay.max_uses.unwrap_or(0),
326                current,
327            });
328        }
329        // Stamp use_number from grant-wide count, also inside the lock,
330        // so two parallel reservations on the same grant cannot both
331        // claim the same use_number.
332        let prior_count = list_uses_for_grant(j, &rec.grant_id)?.len() as u32;
333        rec.use_number = prior_count.saturating_add(1);
334        // Append.
335        let head = read_head(j)?;
336        rec.previous_record_digest = head.digest.clone();
337        rec.record_digest = approval_use_record_digest(&rec);
338        let next_index = head.index + 1;
339        write_record_use(j, next_index, &rec)?;
340        update_indexes_for_use(j, next_index, &rec)?;
341        let new_head = Head {
342            index: next_index,
343            digest: rec.record_digest.clone(),
344            updated_at: rec.created_at.clone(),
345        };
346        write_head(j, &new_head)?;
347        ensure_meta(j)?;
348        Ok(new_head)
349    })
350}
351
352/// Append an ApprovalRevocation. Sibling of `append_use`.
353pub fn append_revocation(j: &Journal, mut rec: ApprovalRevocation) -> Result<Head, JournalError> {
354    rec.type_ = TYPE_APPROVAL_REVOCATION.into();
355    with_lock(j, || {
356        let head = read_head(j)?;
357        rec.previous_record_digest = head.digest.clone();
358        rec.record_digest = approval_revocation_record_digest(&rec);
359        let next_index = head.index + 1;
360        write_record_revocation(j, next_index, &rec)?;
361        index_grant(j, next_index, &rec.grant_id)?;
362        let new_head = Head {
363            index: next_index,
364            digest: rec.record_digest.clone(),
365            updated_at: rec.created_at.clone(),
366        };
367        write_head(j, &new_head)?;
368        ensure_meta(j)?;
369        Ok(new_head)
370    })
371}
372
373/// Append a JournalCheckpoint over a contiguous range of prior records.
374pub fn append_checkpoint(j: &Journal, mut rec: JournalCheckpoint) -> Result<Head, JournalError> {
375    rec.type_ = TYPE_JOURNAL_CHECKPOINT.into();
376    with_lock(j, || {
377        let head = read_head(j)?;
378        rec.previous_record_digest = head.digest.clone();
379        rec.record_digest = journal_checkpoint_record_digest(&rec);
380        let next_index = head.index + 1;
381        write_record_checkpoint(j, next_index, &rec)?;
382        let new_head = Head {
383            index: next_index,
384            digest: rec.record_digest.clone(),
385            updated_at: rec.created_at.clone(),
386        };
387        write_head(j, &new_head)?;
388        ensure_meta(j)?;
389        Ok(new_head)
390    })
391}
392
393fn record_filename(index: u64, type_: &str, digest: &str) -> String {
394    // Use the digest's hex tail (after "sha256:") so the filename is
395    // bounded length and contains no separators.
396    let tail = digest.strip_prefix("sha256:").unwrap_or(digest);
397    let short = &tail[..tail.len().min(16)];
398    format!("{:010}.{type_}.{short}.json", index)
399}
400
401fn write_record_use(j: &Journal, index: u64, rec: &ApprovalUse) -> Result<(), JournalError> {
402    fs::create_dir_all(j.records_dir())?;
403    let name = record_filename(index, "approval-use", &rec.record_digest);
404    let path = j.records_dir().join(&name);
405    let tmp = path.with_extension("json.tmp");
406    let mut f = File::create(&tmp)?;
407    f.write_all(&serde_json::to_vec_pretty(rec)?)?;
408    f.sync_all()?;
409    fs::rename(&tmp, &path)?;
410    Ok(())
411}
412
413fn write_record_revocation(
414    j: &Journal,
415    index: u64,
416    rec: &ApprovalRevocation,
417) -> Result<(), JournalError> {
418    fs::create_dir_all(j.records_dir())?;
419    let name = record_filename(index, "approval-revocation", &rec.record_digest);
420    let path = j.records_dir().join(&name);
421    let tmp = path.with_extension("json.tmp");
422    let mut f = File::create(&tmp)?;
423    f.write_all(&serde_json::to_vec_pretty(rec)?)?;
424    f.sync_all()?;
425    fs::rename(&tmp, &path)?;
426    Ok(())
427}
428
429fn write_record_checkpoint(
430    j: &Journal,
431    index: u64,
432    rec: &JournalCheckpoint,
433) -> Result<(), JournalError> {
434    fs::create_dir_all(j.records_dir())?;
435    let name = record_filename(index, "journal-checkpoint", &rec.record_digest);
436    let path = j.records_dir().join(&name);
437    let tmp = path.with_extension("json.tmp");
438    let mut f = File::create(&tmp)?;
439    f.write_all(&serde_json::to_vec_pretty(rec)?)?;
440    f.sync_all()?;
441    fs::rename(&tmp, &path)?;
442    Ok(())
443}
444
445fn ensure_meta(j: &Journal) -> Result<(), JournalError> {
446    let path = j.meta_path();
447    if path.exists() {
448        return Ok(());
449    }
450    #[derive(serde::Serialize)]
451    struct Meta<'a> {
452        kind: &'a str,
453        version: &'a str,
454        format: &'a str,
455    }
456    let meta = Meta {
457        kind: "approval-use-journal",
458        version: "v1",
459        format: "json-records",
460    };
461    let bytes = serde_json::to_vec_pretty(&meta)?;
462    fs::write(&path, bytes)?;
463    Ok(())
464}
465
466// ---------------------------------------------------------------------------
467// Indexes (rebuildable cache)
468// ---------------------------------------------------------------------------
469
470fn append_index(path: &Path, line: &str) -> Result<(), JournalError> {
471    if let Some(parent) = path.parent() {
472        fs::create_dir_all(parent)?;
473    }
474    let mut f = OpenOptions::new().append(true).create(true).open(path)?;
475    writeln!(f, "{line}")?;
476    Ok(())
477}
478
479fn index_grant(j: &Journal, index: u64, grant_id: &str) -> Result<(), JournalError> {
480    append_index(&j.by_grant_path(grant_id), &index.to_string())
481}
482
483fn index_nonce(j: &Journal, index: u64, nonce_digest: &str) -> Result<(), JournalError> {
484    append_index(&j.by_nonce_path(nonce_digest), &index.to_string())
485}
486
487fn update_indexes_for_use(j: &Journal, index: u64, rec: &ApprovalUse) -> Result<(), JournalError> {
488    index_grant(j, index, &rec.grant_id)?;
489    index_nonce(j, index, &rec.nonce_digest)?;
490    Ok(())
491}
492
493/// Delete and rebuild every index from the records directory. Records are
494/// truth; indexes are cache. Useful as a recovery tool when an index file
495/// is corrupt or out of sync.
496pub fn rebuild_indexes(j: &Journal) -> Result<u64, JournalError> {
497    let dir = j.indexes_dir();
498    if dir.is_dir() {
499        // Wipe by recursive remove. Atomic enough; the worst-case is a
500        // partially-rebuilt index, which the next call to this function
501        // also recovers from.
502        fs::remove_dir_all(&dir)?;
503    }
504    let mut rebuilt = 0u64;
505    for (idx, kind, bytes) in iter_records(j)? {
506        match kind.as_str() {
507            "approval-use" => {
508                let rec: ApprovalUse = serde_json::from_slice(&bytes)?;
509                update_indexes_for_use(j, idx, &rec)?;
510                rebuilt += 1;
511            }
512            "approval-revocation" => {
513                let rec: ApprovalRevocation = serde_json::from_slice(&bytes)?;
514                index_grant(j, idx, &rec.grant_id)?;
515                rebuilt += 1;
516            }
517            "journal-checkpoint" => {
518                rebuilt += 1; // checkpoints aren't indexed by grant/nonce
519            }
520            _ => {}
521        }
522    }
523    Ok(rebuilt)
524}
525
526// ---------------------------------------------------------------------------
527// Iteration + integrity
528// ---------------------------------------------------------------------------
529
530/// Walk records/ in index order. Returns `(index, kind, bytes)`. Kind is
531/// derived from the filename ("approval-use" / "approval-revocation" /
532/// "journal-checkpoint"). Filenames Treeship doesn't recognize are
533/// skipped silently rather than failing the whole walk -- a future record
534/// type added by a newer version shouldn't break older readers.
535fn iter_records(j: &Journal) -> Result<Vec<(u64, String, Vec<u8>)>, JournalError> {
536    let dir = j.records_dir();
537    if !dir.is_dir() {
538        return Ok(Vec::new());
539    }
540    let mut entries: Vec<(u64, String, PathBuf)> = Vec::new();
541    for entry in fs::read_dir(&dir)? {
542        let entry = entry?;
543        let path = entry.path();
544        if path.extension().and_then(|s| s.to_str()) != Some("json") {
545            continue;
546        }
547        let name = match path.file_name().and_then(|n| n.to_str()) {
548            Some(n) => n,
549            None => continue,
550        };
551        // Filename shape: "<10-digit-index>.<kind>.<short-digest>.json"
552        let mut parts = name.splitn(4, '.');
553        let idx_str = match parts.next() {
554            Some(s) => s,
555            None => continue,
556        };
557        let kind = match parts.next() {
558            Some(s) => s,
559            None => continue,
560        };
561        // index parses as u64
562        let idx = match idx_str.parse::<u64>() {
563            Ok(n) => n,
564            Err(_) => continue,
565        };
566        entries.push((idx, kind.to_string(), path));
567    }
568    entries.sort_by_key(|(idx, _, _)| *idx);
569    let mut out = Vec::with_capacity(entries.len());
570    for (idx, kind, path) in entries {
571        let bytes = fs::read(&path)?;
572        out.push((idx, kind, bytes));
573    }
574    Ok(out)
575}
576
577/// Walk every record in order, recompute each `record_digest`, and check
578/// that each record's `previous_record_digest` matches the prior
579/// record's stored `record_digest`. Returns the number of records walked
580/// or an error pinpointing the first integrity failure.
581pub fn verify_integrity(j: &Journal) -> Result<u64, JournalError> {
582    let mut prior_digest = String::new();
583    let mut count = 0u64;
584    let head = read_head(j)?;
585    for (idx, kind, bytes) in iter_records(j)? {
586        match kind.as_str() {
587            "approval-use" => {
588                let rec: ApprovalUse = serde_json::from_slice(&bytes)?;
589                if rec.previous_record_digest != prior_digest {
590                    return Err(JournalError::BrokenChain {
591                        index: idx,
592                        expected: prior_digest,
593                        actual: rec.previous_record_digest,
594                    });
595                }
596                let recomputed = approval_use_record_digest(&rec);
597                if recomputed != rec.record_digest {
598                    return Err(JournalError::RecordTampered {
599                        index: idx,
600                        expected: rec.record_digest,
601                        actual: recomputed,
602                    });
603                }
604                prior_digest = rec.record_digest;
605            }
606            "approval-revocation" => {
607                let rec: ApprovalRevocation = serde_json::from_slice(&bytes)?;
608                if rec.previous_record_digest != prior_digest {
609                    return Err(JournalError::BrokenChain {
610                        index: idx,
611                        expected: prior_digest,
612                        actual: rec.previous_record_digest,
613                    });
614                }
615                let recomputed = approval_revocation_record_digest(&rec);
616                if recomputed != rec.record_digest {
617                    return Err(JournalError::RecordTampered {
618                        index: idx,
619                        expected: rec.record_digest,
620                        actual: recomputed,
621                    });
622                }
623                prior_digest = rec.record_digest;
624            }
625            "journal-checkpoint" => {
626                let rec: JournalCheckpoint = serde_json::from_slice(&bytes)?;
627                if rec.previous_record_digest != prior_digest {
628                    return Err(JournalError::BrokenChain {
629                        index: idx,
630                        expected: prior_digest,
631                        actual: rec.previous_record_digest,
632                    });
633                }
634                let recomputed = journal_checkpoint_record_digest(&rec);
635                if recomputed != rec.record_digest {
636                    return Err(JournalError::RecordTampered {
637                        index: idx,
638                        expected: rec.record_digest,
639                        actual: recomputed,
640                    });
641                }
642                prior_digest = rec.record_digest;
643            }
644            _ => {
645                // Unknown record kind. Stop the chain check rather than
646                // skip silently -- a newer record type would still need
647                // to participate in the chain.
648                continue;
649            }
650        }
651        count += 1;
652    }
653    // Tail must match the head if records exist; if records were
654    // deleted off the end the head will be stale.
655    if head.index != 0 && head.digest != prior_digest {
656        return Err(JournalError::MissingRecord { index: head.index });
657    }
658    Ok(count)
659}
660
661// ---------------------------------------------------------------------------
662// check_replay
663// ---------------------------------------------------------------------------
664
665/// Check whether (`grant_id`, `nonce_digest`) has already been consumed,
666/// and how many times. Returns a `ReplayCheck` carrying the strongest
667/// level the journal can speak to:
668///
669///   - `NotPerformed` when the journal directory does not exist on disk.
670///     The caller (verify) should fall back to its package-local check.
671///   - `LocalJournal` otherwise. `passed: true` means the use count is
672///     within `max_uses_hint`; `false` means it would exceed.
673///
674/// `max_uses_hint` is what the caller knows from the signed grant's
675/// `ApprovalScope.max_actions`. We accept it as a hint rather than
676/// reading it back from a stored record because the stored uses already
677/// carry their own `max_uses` snapshot, and disagreement between the
678/// hint and the stored value should be visible in `details`.
679pub fn check_replay(
680    j: &Journal,
681    grant_id: &str,
682    nonce_digest: &str,
683    max_uses_hint: Option<u32>,
684) -> Result<ReplayCheck, JournalError> {
685    if !j.exists() {
686        return Ok(ReplayCheck::not_performed());
687    }
688    // Use the by-nonce index: every prior use of the same approval
689    // shares the same nonce_digest, so the index gives us the exact
690    // record list.
691    let index_path = j.by_nonce_path(nonce_digest);
692    let mut current = 0u32;
693    let mut last_max: Option<u32> = None;
694    if index_path.exists() {
695        let raw = fs::read_to_string(&index_path)?;
696        for line in raw.lines() {
697            let idx: u64 = match line.trim().parse() {
698                Ok(n) => n,
699                Err(_) => continue,
700            };
701            if let Some(rec) = load_use_record(j, idx)? {
702                // Only count uses that bind to the same grant_id; the
703                // by-nonce index can in theory share a digest across
704                // grants, though in practice nonces are random.
705                if rec.grant_id == grant_id {
706                    current = current.saturating_add(1);
707                    last_max = rec.max_uses.or(last_max);
708                }
709            }
710        }
711    }
712    let max_uses = max_uses_hint.or(last_max);
713    let passed = match max_uses {
714        Some(m) => current < m,
715        None => true, // unbounded grant; PR 5 reports this honestly
716    };
717    let details = match max_uses {
718        Some(m) => format!("local Approval Use Journal: use {current}/{m}"),
719        None => {
720            format!("local Approval Use Journal: {current} prior use(s); grant has no max_uses")
721        }
722    };
723    Ok(ReplayCheck {
724        level: ReplayCheckLevel::LocalJournal,
725        use_number: Some(current.saturating_add(1)),
726        max_uses,
727        passed: Some(passed),
728        details: Some(details),
729    })
730}
731
732/// Find a use record by its `use_id`, scanning the records directory. The
733/// journal indexes by grant and by nonce; a caller holding only the id an
734/// action's `meta.approval_use_id` names (the VI attestation builder, an
735/// auditor with a receipt in hand) needs this lookup.
736pub fn find_use_by_id(j: &Journal, use_id: &str) -> Result<Option<ApprovalUse>, JournalError> {
737    let dir = j.records_dir();
738    if !dir.is_dir() {
739        return Ok(None);
740    }
741    for entry in fs::read_dir(&dir)? {
742        let entry = entry?;
743        let name = entry.file_name().to_string_lossy().into_owned();
744        if !name.contains(".approval-use.") {
745            continue;
746        }
747        let bytes = fs::read(entry.path())?;
748        if let Ok(rec) = serde_json::from_slice::<ApprovalUse>(&bytes) {
749            if rec.use_id == use_id {
750                return Ok(Some(rec));
751            }
752        }
753    }
754    Ok(None)
755}
756
757fn load_use_record(j: &Journal, index: u64) -> Result<Option<ApprovalUse>, JournalError> {
758    let dir = j.records_dir();
759    if !dir.is_dir() {
760        return Ok(None);
761    }
762    let prefix = format!("{:010}.approval-use.", index);
763    for entry in fs::read_dir(&dir)? {
764        let entry = entry?;
765        let name = entry.file_name().to_string_lossy().into_owned();
766        if name.starts_with(&prefix) {
767            let bytes = fs::read(entry.path())?;
768            let rec: ApprovalUse = serde_json::from_slice(&bytes)?;
769            return Ok(Some(rec));
770        }
771    }
772    Ok(None)
773}
774
775// ---------------------------------------------------------------------------
776// Public read helpers (CLI)
777// ---------------------------------------------------------------------------
778
779/// Find the recorded ApprovalUse for an already-signed action.
780/// Returns the matching use record plus a `ReplayCheck` that answers
781/// the *verify-time* question -- "is the recorded use within max_uses?"
782/// -- as opposed to `check_replay`'s consume-time question -- "would
783/// the next use exceed?". The two questions look the same but have
784/// different boundary semantics:
785///
786///   consume-time: passed = use_number_that_would_be_allocated <= max_uses
787///                 (i.e. current_count < max_uses, since next = current + 1)
788///   verify-time:  passed = recorded_use_number <= max_uses
789///
790/// Verify should call THIS, not check_replay, when reporting on an
791/// action that already has a journal record.
792pub fn find_use_for_action(
793    j: &Journal,
794    grant_id: &str,
795    nonce_digest: &str,
796    max_uses_hint: Option<u32>,
797    action_artifact_id: Option<&str>,
798) -> Result<Option<(ApprovalUse, ReplayCheck)>, JournalError> {
799    if !j.exists() {
800        return Ok(None);
801    }
802    let index_path = j.by_nonce_path(nonce_digest);
803    if !index_path.exists() {
804        return Ok(None);
805    }
806    let raw = fs::read_to_string(&index_path)?;
807    // The use record for the action under verification is the one that
808    // names it: consume-time backfills `action_artifact_id` onto the
809    // record once the action is signed. Matching on (grant_id,
810    // nonce_digest) alone let a second directory's journal, holding its own
811    // `use 1/1` for the same grant, vouch for an action it never recorded
812    // (film findings 2026-09-22, #11). Records without an action id
813    // (written before the backfill existed, or whose backfill failed) fall
814    // back to the most recent match, as before.
815    let mut named_this: Option<ApprovalUse> = None;
816    let mut latest: Option<ApprovalUse> = None;
817    let mut named_other = false;
818    let mut any_unnamed = false;
819    for line in raw.lines() {
820        let idx: u64 = match line.trim().parse() {
821            Ok(n) => n,
822            Err(_) => continue,
823        };
824        if let Some(rec) = load_use_record(j, idx)? {
825            if rec.grant_id == grant_id {
826                match (action_artifact_id, rec.action_artifact_id.as_deref()) {
827                    (Some(want), Some(have)) if want == have => named_this = Some(rec.clone()),
828                    (Some(_), Some(_)) => named_other = true,
829                    _ => any_unnamed = true,
830                }
831                latest = Some(rec);
832            }
833        }
834    }
835    let rec = match (named_this, latest) {
836        (Some(rec), _) => rec,
837        (None, Some(rec)) if named_other && !any_unnamed => {
838            // Every record for this grant and nonce names a different
839            // action. This action's consumption is not in this journal, so
840            // the journal cannot say its use was within max_uses.
841            let m = max_uses_hint.or(rec.max_uses);
842            let details = format!(
843                "local Approval Use Journal has no use record for this action; use {}{} of this grant and nonce was recorded for {}",
844                rec.use_number,
845                m.map(|m| format!("/{m}")).unwrap_or_default(),
846                rec.action_artifact_id.as_deref().unwrap_or("another action")
847            );
848            return Ok(Some((
849                rec.clone(),
850                ReplayCheck {
851                    level: ReplayCheckLevel::LocalJournal,
852                    use_number: Some(rec.use_number),
853                    max_uses: m,
854                    passed: Some(false),
855                    details: Some(details),
856                },
857            )));
858        }
859        (None, Some(rec)) => rec,
860        (None, None) => return Ok(None),
861    };
862
863    let stored_max = rec.max_uses;
864    let max_uses = max_uses_hint.or(stored_max);
865    let passed = match max_uses {
866        Some(m) => rec.use_number <= m,
867        None => true,
868    };
869    let details = match max_uses {
870        Some(m) => format!(
871            "local Approval Use Journal passed, use {}/{}",
872            rec.use_number, m
873        ),
874        None => format!(
875            "local Approval Use Journal: use {} of unbounded grant",
876            rec.use_number
877        ),
878    };
879    Ok(Some((
880        rec.clone(),
881        ReplayCheck {
882            level: ReplayCheckLevel::LocalJournal,
883            use_number: Some(rec.use_number),
884            max_uses,
885            passed: Some(passed),
886            details: Some(details),
887        },
888    )))
889}
890
891/// Every ApprovalUse for `grant_id`. Reads the by-grant index, then
892/// loads each record. Quiet on missing journal.
893pub fn list_uses_for_grant(j: &Journal, grant_id: &str) -> Result<Vec<ApprovalUse>, JournalError> {
894    if !j.exists() {
895        return Ok(Vec::new());
896    }
897    let index_path = j.by_grant_path(grant_id);
898    if !index_path.exists() {
899        return Ok(Vec::new());
900    }
901    let raw = fs::read_to_string(&index_path)?;
902    let mut out = Vec::new();
903    for line in raw.lines() {
904        let idx: u64 = match line.trim().parse() {
905            Ok(n) => n,
906            Err(_) => continue,
907        };
908        if let Some(rec) = load_use_record(j, idx)? {
909            out.push(rec);
910        }
911    }
912    Ok(out)
913}
914
915// ---------------------------------------------------------------------------
916// Tests
917// ---------------------------------------------------------------------------
918
919#[cfg(test)]
920mod tests {
921    use super::*;
922    use tempfile::tempdir;
923
924    fn sample_use(use_id: &str, grant_id: &str, nonce_digest: &str, n: u32) -> ApprovalUse {
925        ApprovalUse {
926            type_: TYPE_APPROVAL_USE.into(),
927            use_id: use_id.into(),
928            grant_id: grant_id.into(),
929            grant_digest: "sha256:00".into(),
930            nonce_digest: nonce_digest.into(),
931            actor: "agent://deployer".into(),
932            action: "deploy.production".into(),
933            subject: "env://production".into(),
934            session_id: None,
935            action_artifact_id: None,
936            receipt_digest: None,
937            use_number: n,
938            max_uses: Some(2),
939            idempotency_key: None,
940            created_at: "2026-04-30T07:00:00Z".into(),
941            expires_at: None,
942            previous_record_digest: String::new(), // append_use rewrites this
943            record_digest: String::new(),          // append_use rewrites this
944            signature: None,
945            signature_alg: None,
946            signing_key_id: None,
947        }
948    }
949
950    #[test]
951    fn first_append_creates_layout_and_head() {
952        let dir = tempdir().unwrap();
953        let j = Journal::new(dir.path());
954        let head = append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
955        assert_eq!(head.index, 1);
956        assert!(j.records_dir().is_dir());
957        assert!(j.heads_dir().is_dir());
958        assert!(j.current_head_path().is_file());
959        assert!(j.meta_path().is_file());
960        // by-grant + by-nonce indexes populated
961        assert!(j.by_grant_path("g1").is_file());
962        assert!(j.by_nonce_path("sha256:nn1").is_file());
963    }
964
965    #[test]
966    fn second_append_links_previous_record_digest() {
967        let dir = tempdir().unwrap();
968        let j = Journal::new(dir.path());
969        let h1 = append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
970        let h2 = append_use(&j, sample_use("use_2", "g1", "sha256:nn2", 2)).unwrap();
971        assert_eq!(h2.index, 2);
972        // Reading record 2 should show previous_record_digest == h1.digest
973        let recs = iter_records(&j).unwrap();
974        assert_eq!(recs.len(), 2);
975        let (_, _, bytes) = &recs[1];
976        let r2: ApprovalUse = serde_json::from_slice(bytes).unwrap();
977        assert_eq!(r2.previous_record_digest, h1.digest);
978    }
979
980    #[test]
981    fn verify_integrity_passes_on_intact_chain() {
982        let dir = tempdir().unwrap();
983        let j = Journal::new(dir.path());
984        for i in 1..=5 {
985            let nd = format!("sha256:nn{i}");
986            append_use(&j, sample_use(&format!("use_{i}"), "g1", &nd, i)).unwrap();
987        }
988        assert_eq!(verify_integrity(&j).unwrap(), 5);
989    }
990
991    #[test]
992    fn editing_a_record_breaks_integrity() {
993        let dir = tempdir().unwrap();
994        let j = Journal::new(dir.path());
995        append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
996        // Find the on-disk record file and corrupt it.
997        let entries: Vec<_> = fs::read_dir(j.records_dir()).unwrap().collect();
998        let entry = entries.into_iter().next().unwrap().unwrap();
999        let mut json: serde_json::Value =
1000            serde_json::from_slice(&fs::read(entry.path()).unwrap()).unwrap();
1001        json["actor"] = "agent://attacker".into();
1002        fs::write(entry.path(), serde_json::to_vec_pretty(&json).unwrap()).unwrap();
1003
1004        let err = verify_integrity(&j).unwrap_err();
1005        assert!(
1006            matches!(err, JournalError::RecordTampered { .. }),
1007            "expected RecordTampered, got {err:?}"
1008        );
1009    }
1010
1011    #[test]
1012    fn deleting_a_record_breaks_integrity_or_head_continuity() {
1013        let dir = tempdir().unwrap();
1014        let j = Journal::new(dir.path());
1015        append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1016        append_use(&j, sample_use("use_2", "g1", "sha256:nn2", 2)).unwrap();
1017        // Remove the trailing record. Head still points at index 2.
1018        let entries: Vec<_> = fs::read_dir(j.records_dir())
1019            .unwrap()
1020            .map(|e| e.unwrap().path())
1021            .collect();
1022        let trailing = entries.iter().max().unwrap();
1023        fs::remove_file(trailing).unwrap();
1024
1025        let err = verify_integrity(&j).unwrap_err();
1026        assert!(
1027            matches!(err, JournalError::MissingRecord { .. }),
1028            "expected MissingRecord, got {err:?}"
1029        );
1030    }
1031
1032    #[test]
1033    fn indexes_can_be_rebuilt_from_records() {
1034        let dir = tempdir().unwrap();
1035        let j = Journal::new(dir.path());
1036        for i in 1..=3 {
1037            let nd = format!("sha256:nn{i}");
1038            append_use(&j, sample_use(&format!("use_{i}"), "g1", &nd, i)).unwrap();
1039        }
1040        // Wipe indexes; check_replay (or rebuild_indexes) should still work.
1041        fs::remove_dir_all(j.indexes_dir()).unwrap();
1042
1043        let rebuilt = rebuild_indexes(&j).unwrap();
1044        assert_eq!(rebuilt, 3);
1045        assert!(j.by_grant_path("g1").is_file());
1046        assert!(j.by_nonce_path("sha256:nn1").is_file());
1047    }
1048
1049    #[test]
1050    fn check_replay_reports_use_count_and_max() {
1051        let dir = tempdir().unwrap();
1052        let j = Journal::new(dir.path());
1053        // Two prior uses of grant g1 with the same nonce_digest.
1054        append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1055        append_use(&j, sample_use("use_2", "g1", "sha256:nn1", 2)).unwrap();
1056
1057        // max_uses_hint = 2: the next use would be 3/2 -> not passed.
1058        let r = check_replay(&j, "g1", "sha256:nn1", Some(2)).unwrap();
1059        assert_eq!(r.level, ReplayCheckLevel::LocalJournal);
1060        assert_eq!(r.use_number, Some(3));
1061        assert_eq!(r.max_uses, Some(2));
1062        assert_eq!(r.passed, Some(false));
1063    }
1064
1065    #[test]
1066    fn check_replay_passes_when_under_max() {
1067        let dir = tempdir().unwrap();
1068        let j = Journal::new(dir.path());
1069        append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1070        let r = check_replay(&j, "g1", "sha256:nn1", Some(2)).unwrap();
1071        assert_eq!(r.use_number, Some(2));
1072        assert_eq!(r.passed, Some(true));
1073    }
1074
1075    #[test]
1076    fn check_replay_no_journal_returns_not_performed() {
1077        let dir = tempdir().unwrap();
1078        let absent = dir.path().join("nope");
1079        let j = Journal::new(&absent);
1080        let r = check_replay(&j, "g1", "sha256:nn1", Some(1)).unwrap();
1081        assert_eq!(r.level, ReplayCheckLevel::NotPerformed);
1082        assert!(r.use_number.is_none());
1083    }
1084
1085    #[test]
1086    fn check_replay_unbounded_grant_passes_with_count() {
1087        let dir = tempdir().unwrap();
1088        let j = Journal::new(dir.path());
1089        append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1090        // No max_uses_hint and stored record's max_uses is Some(2) too,
1091        // so we explicitly set None on a fresh record to test the
1092        // unbounded path.
1093        let mut u = sample_use("use_2", "g2", "sha256:other", 1);
1094        u.max_uses = None;
1095        append_use(&j, u).unwrap();
1096
1097        let r = check_replay(&j, "g2", "sha256:other", None).unwrap();
1098        assert!(r.passed.unwrap());
1099        assert!(r.max_uses.is_none());
1100    }
1101
1102    #[test]
1103    fn list_uses_for_grant_returns_records_in_order() {
1104        let dir = tempdir().unwrap();
1105        let j = Journal::new(dir.path());
1106        append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1107        append_use(&j, sample_use("use_2", "g2", "sha256:nn2", 1)).unwrap();
1108        append_use(&j, sample_use("use_3", "g1", "sha256:nn3", 2)).unwrap();
1109        let g1 = list_uses_for_grant(&j, "g1").unwrap();
1110        assert_eq!(g1.len(), 2);
1111        assert_eq!(g1[0].use_id, "use_1");
1112        assert_eq!(g1[1].use_id, "use_3");
1113    }
1114
1115    #[test]
1116    fn lock_keeps_two_appends_serial() {
1117        // Hold the lock externally; an append should fail with LockBusy
1118        // rather than racing or silently overwriting.
1119        let dir = tempdir().unwrap();
1120        let j = Journal::new(dir.path());
1121        fs::create_dir_all(j.locks_dir()).unwrap();
1122        let held = OpenOptions::new()
1123            .read(true)
1124            .write(true)
1125            .create(true)
1126            .truncate(false)
1127            .open(j.lock_path())
1128            .unwrap();
1129        held.try_lock_exclusive().unwrap();
1130
1131        let err = append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap_err();
1132        assert!(matches!(err, JournalError::LockBusy));
1133
1134        let _ = fs2::FileExt::unlock(&held);
1135    }
1136
1137    #[test]
1138    fn revocation_appends_into_chain() {
1139        let dir = tempdir().unwrap();
1140        let j = Journal::new(dir.path());
1141        append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1142        let rev = ApprovalRevocation {
1143            type_: TYPE_APPROVAL_REVOCATION.into(),
1144            revocation_id: "rev_1".into(),
1145            grant_id: "g1".into(),
1146            grant_digest: "sha256:00".into(),
1147            revoker: "human://alice".into(),
1148            reason: Some("rotated key".into()),
1149            created_at: "2026-04-30T07:01:00Z".into(),
1150            previous_record_digest: String::new(),
1151            record_digest: String::new(),
1152            signature: None,
1153            signature_alg: None,
1154            signing_key_id: None,
1155        };
1156        let h = append_revocation(&j, rev).unwrap();
1157        assert_eq!(h.index, 2);
1158        assert_eq!(verify_integrity(&j).unwrap(), 2);
1159    }
1160
1161    #[test]
1162    fn record_files_contain_no_raw_nonce_or_signature_secrets() {
1163        // Privacy invariant: ApprovalUse has no `nonce` field on the
1164        // struct, so by construction the stored JSON only contains
1165        // `nonce_digest`. This test pins the on-disk shape so a future
1166        // schema change can't sneak in a raw-nonce field.
1167        let dir = tempdir().unwrap();
1168        let j = Journal::new(dir.path());
1169        append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1170        let entries: Vec<_> = fs::read_dir(j.records_dir())
1171            .unwrap()
1172            .map(|e| e.unwrap().path())
1173            .collect();
1174        let bytes = fs::read(&entries[0]).unwrap();
1175        let json: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
1176        let obj = json.as_object().unwrap();
1177        for forbidden in [
1178            "nonce",
1179            "command",
1180            "prompt",
1181            "file_content",
1182            "bearer_token",
1183            "api_key",
1184        ] {
1185            assert!(
1186                !obj.contains_key(forbidden),
1187                "journal record must not contain `{forbidden}`",
1188            );
1189        }
1190        // The digest IS allowed.
1191        assert!(obj.contains_key("nonce_digest"));
1192    }
1193
1194    // -- v0.9.10 PR A: reserve_use atomic check+append regression tests --
1195
1196    #[test]
1197    fn reserve_use_first_call_succeeds_and_stamps_use_number() {
1198        // Caller passes use_number=0; reserve_use stamps it from the
1199        // grant-wide count observed inside the lock.
1200        let dir = tempdir().unwrap();
1201        let j = Journal::new(dir.path());
1202        let mut rec = sample_use("use_1", "g1", "sha256:nn1", 0);
1203        rec.use_number = 0;
1204        let head = reserve_use(&j, rec, Some(1)).unwrap();
1205        assert_eq!(head.index, 1);
1206        let stored = list_uses_for_grant(&j, "g1").unwrap();
1207        assert_eq!(stored.len(), 1);
1208        assert_eq!(
1209            stored[0].use_number, 1,
1210            "reserve_use must stamp use_number=1 for the first use"
1211        );
1212    }
1213
1214    #[test]
1215    fn reserve_use_max_uses_1_serial_second_call_rejects() {
1216        // Sequential second call with the same nonce against
1217        // max_uses=1 must error with MaxUsesExceeded BEFORE writing.
1218        let dir = tempdir().unwrap();
1219        let j = Journal::new(dir.path());
1220        reserve_use(&j, sample_use("use_1", "g1", "sha256:nn_a", 0), Some(1)).unwrap();
1221
1222        let err = reserve_use(&j, sample_use("use_2", "g1", "sha256:nn_a", 0), Some(1))
1223            .expect_err("second consume of max_uses=1 grant must fail");
1224        match err {
1225            JournalError::MaxUsesExceeded {
1226                grant_id,
1227                max_uses,
1228                current,
1229            } => {
1230                assert_eq!(grant_id, "g1");
1231                assert_eq!(max_uses, 1);
1232                assert_eq!(current, 1);
1233            }
1234            other => panic!("expected MaxUsesExceeded, got {other:?}"),
1235        }
1236        // Crucially, the second record is NOT written.
1237        let stored = list_uses_for_grant(&j, "g1").unwrap();
1238        assert_eq!(stored.len(), 1, "rejected reserve must not append");
1239    }
1240
1241    #[test]
1242    fn reserve_use_max_uses_2_two_uses_pass_third_rejects() {
1243        // Legitimate multi-use grant: two distinct nonces, same grant.
1244        let dir = tempdir().unwrap();
1245        let j = Journal::new(dir.path());
1246        let mut a = sample_use("use_1", "g1", "sha256:nn_a", 0);
1247        a.max_uses = Some(2);
1248        let mut b = sample_use("use_2", "g1", "sha256:nn_b", 0);
1249        b.max_uses = Some(2);
1250        reserve_use(&j, a, Some(2)).unwrap();
1251        reserve_use(&j, b, Some(2)).unwrap();
1252        // A third nonce with max_uses=2 is fine (per-nonce check, not
1253        // per-grant); the journal's invariant is single-use-per-nonce.
1254        let mut c = sample_use("use_3", "g1", "sha256:nn_c", 0);
1255        c.max_uses = Some(2);
1256        reserve_use(&j, c, Some(2)).unwrap();
1257        // But a SECOND consume of nn_a violates max_uses=2 because
1258        // that nonce already has 1 use; 1+1 = 2 is within bound, so
1259        // this second use of nn_a is actually allowed — pin that.
1260        let mut a2 = sample_use("use_1b", "g1", "sha256:nn_a", 0);
1261        a2.max_uses = Some(2);
1262        reserve_use(&j, a2, Some(2)).unwrap();
1263        // A THIRD consume of nn_a exceeds max_uses=2.
1264        let mut a3 = sample_use("use_1c", "g1", "sha256:nn_a", 0);
1265        a3.max_uses = Some(2);
1266        let err = reserve_use(&j, a3, Some(2)).expect_err("third use of same nonce must fail");
1267        assert!(matches!(err, JournalError::MaxUsesExceeded { .. }));
1268    }
1269
1270    /// Round-2 hardening: idempotency-key retries through the
1271    /// CLI's `reserve_in_journal` should NOT bypass `reserve_use`'s
1272    /// max_uses gate. The CLI checks the idempotency key against
1273    /// existing uses *before* calling `reserve_use`; this test
1274    /// confirms that path doesn't sneak a second reservation past
1275    /// max_uses=1 just because a flaky retry uses the same key.
1276    ///
1277    /// The journal-level invariant we pin here: even if a caller
1278    /// repeatedly invokes `reserve_use` with the same record after a
1279    /// `LockBusy`, the second call sees the first record (now
1280    /// committed) and rejects with `MaxUsesExceeded`. There is no
1281    /// "free retry" loophole.
1282    #[test]
1283    fn reserve_use_retry_after_lock_busy_does_not_bypass_max_uses() {
1284        let dir = tempdir().unwrap();
1285        let j = Journal::new(dir.path());
1286        // First reserve commits use_1.
1287        reserve_use(&j, sample_use("use_1", "g1", "sha256:nn_retry", 0), Some(1)).unwrap();
1288        // Subsequent reserves with the SAME nonce all fail with
1289        // MaxUsesExceeded -- no retry-bypass window.
1290        for i in 0..5 {
1291            let err = reserve_use(
1292                &j,
1293                sample_use(&format!("use_retry_{i}"), "g1", "sha256:nn_retry", 0),
1294                Some(1),
1295            )
1296            .expect_err("retry must fail");
1297            assert!(matches!(err, JournalError::MaxUsesExceeded { .. }));
1298        }
1299        let stored = list_uses_for_grant(&j, "g1").unwrap();
1300        assert_eq!(
1301            stored.len(),
1302            1,
1303            "exactly one record on disk despite 5 retries"
1304        );
1305    }
1306
1307    #[test]
1308    fn reserve_use_concurrent_max_uses_1_only_one_succeeds() {
1309        // The headline regression test for the v0.9.9 TOCTOU race.
1310        //
1311        // Eight threads race to reserve the same (grant_id, nonce_digest)
1312        // against max_uses=1. With v0.9.9's split check_replay/append_use
1313        // pattern, two threads could both pass the pre-lock replay check
1314        // and write — exceeding max_uses. With v0.9.10's reserve_use,
1315        // the check happens INSIDE the lock; exactly one thread wins.
1316        //
1317        // Outcomes for the 7 losers are a mix of:
1318        //   - LockBusy: the lock was held when they tried try_lock
1319        //   - MaxUsesExceeded: they got the lock after the winner
1320        //     released, saw the winner's record, declined to write
1321        // Both are correct — neither is a bypass.
1322        use std::sync::atomic::{AtomicUsize, Ordering};
1323        use std::sync::Arc;
1324        use std::thread;
1325
1326        let dir = tempdir().unwrap();
1327        let dir_path = Arc::new(dir.path().to_path_buf());
1328        let success = Arc::new(AtomicUsize::new(0));
1329        let lock_busy = Arc::new(AtomicUsize::new(0));
1330        let max_exceeded = Arc::new(AtomicUsize::new(0));
1331
1332        let mut handles = Vec::new();
1333        for i in 0..8 {
1334            let dir_path = Arc::clone(&dir_path);
1335            let success = Arc::clone(&success);
1336            let lock_busy = Arc::clone(&lock_busy);
1337            let max_exceeded = Arc::clone(&max_exceeded);
1338            handles.push(thread::spawn(move || {
1339                let j = Journal::new(dir_path.as_path());
1340                let rec = sample_use(&format!("use_{i}"), "g1", "sha256:race_nonce", 0);
1341                match reserve_use(&j, rec, Some(1)) {
1342                    Ok(_) => {
1343                        success.fetch_add(1, Ordering::SeqCst);
1344                    }
1345                    Err(JournalError::LockBusy) => {
1346                        lock_busy.fetch_add(1, Ordering::SeqCst);
1347                    }
1348                    Err(JournalError::MaxUsesExceeded { .. }) => {
1349                        max_exceeded.fetch_add(1, Ordering::SeqCst);
1350                    }
1351                    Err(other) => panic!("unexpected error: {other:?}"),
1352                }
1353            }));
1354        }
1355        for h in handles {
1356            h.join().unwrap();
1357        }
1358
1359        let s = success.load(Ordering::SeqCst);
1360        let lb = lock_busy.load(Ordering::SeqCst);
1361        let me = max_exceeded.load(Ordering::SeqCst);
1362        assert_eq!(s, 1, "exactly one of 8 concurrent reserves must succeed; got {s} (lock_busy={lb}, max_exceeded={me})");
1363        assert_eq!(s + lb + me, 8, "every thread accounted for");
1364
1365        // Belt-and-braces: only one record actually on disk for this nonce.
1366        let stored = list_uses_for_grant(&Journal::new(dir.path()), "g1").unwrap();
1367        let same_nonce = stored
1368            .iter()
1369            .filter(|u| u.nonce_digest == "sha256:race_nonce")
1370            .count();
1371        assert_eq!(
1372            same_nonce, 1,
1373            "exactly one record on disk for the contested nonce"
1374        );
1375    }
1376}