1use std::fs::{self, File, OpenOptions};
31use std::io::Write;
32use std::path::{Path, PathBuf};
33
34#[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#[derive(Debug)]
53pub enum JournalError {
54 Io(std::io::Error),
55 Json(serde_json::Error),
56 BrokenChain {
59 index: u64,
60 expected: String,
61 actual: String,
62 },
63 RecordTampered {
66 index: u64,
67 expected: String,
68 actual: String,
69 },
70 MissingRecord {
72 index: u64,
73 },
74 LockBusy,
76 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
126pub struct Journal {
132 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 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 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 pub fn exists(&self) -> bool {
179 self.dir.is_dir()
180 }
181}
182
183fn safe_name(s: &str) -> String {
187 s.chars()
188 .map(|c| match c {
189 ':' | '/' | '\\' | ' ' | '.' => '_',
190 c => c,
191 })
192 .collect()
193}
194
195#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default)]
200pub struct Head {
201 pub index: u64,
203 pub digest: String,
205 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#[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#[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
266pub 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
292pub 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 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 let prior_count = list_uses_for_grant(j, &rec.grant_id)?.len() as u32;
333 rec.use_number = prior_count.saturating_add(1);
334 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
352pub 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
373pub 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 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
466fn 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
493pub fn rebuild_indexes(j: &Journal) -> Result<u64, JournalError> {
497 let dir = j.indexes_dir();
498 if dir.is_dir() {
499 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; }
520 _ => {}
521 }
522 }
523 Ok(rebuilt)
524}
525
526fn 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 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 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
577pub 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 continue;
649 }
650 }
651 count += 1;
652 }
653 if head.index != 0 && head.digest != prior_digest {
656 return Err(JournalError::MissingRecord { index: head.index });
657 }
658 Ok(count)
659}
660
661pub 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 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 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, };
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
732fn load_use_record(j: &Journal, index: u64) -> Result<Option<ApprovalUse>, JournalError> {
733 let dir = j.records_dir();
734 if !dir.is_dir() {
735 return Ok(None);
736 }
737 let prefix = format!("{:010}.approval-use.", index);
738 for entry in fs::read_dir(&dir)? {
739 let entry = entry?;
740 let name = entry.file_name().to_string_lossy().into_owned();
741 if name.starts_with(&prefix) {
742 let bytes = fs::read(entry.path())?;
743 let rec: ApprovalUse = serde_json::from_slice(&bytes)?;
744 return Ok(Some(rec));
745 }
746 }
747 Ok(None)
748}
749
750pub fn find_use_for_action(
768 j: &Journal,
769 grant_id: &str,
770 nonce_digest: &str,
771 max_uses_hint: Option<u32>,
772) -> Result<Option<(ApprovalUse, ReplayCheck)>, JournalError> {
773 if !j.exists() {
774 return Ok(None);
775 }
776 let index_path = j.by_nonce_path(nonce_digest);
777 if !index_path.exists() {
778 return Ok(None);
779 }
780 let raw = fs::read_to_string(&index_path)?;
781 let mut latest: Option<ApprovalUse> = None;
788 for line in raw.lines() {
789 let idx: u64 = match line.trim().parse() {
790 Ok(n) => n,
791 Err(_) => continue,
792 };
793 if let Some(rec) = load_use_record(j, idx)? {
794 if rec.grant_id == grant_id {
795 latest = Some(rec);
796 }
797 }
798 }
799 let Some(rec) = latest else { return Ok(None) };
800
801 let stored_max = rec.max_uses;
802 let max_uses = max_uses_hint.or(stored_max);
803 let passed = match max_uses {
804 Some(m) => rec.use_number <= m,
805 None => true,
806 };
807 let details = match max_uses {
808 Some(m) => format!(
809 "local Approval Use Journal passed, use {}/{}",
810 rec.use_number, m
811 ),
812 None => format!(
813 "local Approval Use Journal: use {} of unbounded grant",
814 rec.use_number
815 ),
816 };
817 Ok(Some((
818 rec.clone(),
819 ReplayCheck {
820 level: ReplayCheckLevel::LocalJournal,
821 use_number: Some(rec.use_number),
822 max_uses,
823 passed: Some(passed),
824 details: Some(details),
825 },
826 )))
827}
828
829pub fn list_uses_for_grant(j: &Journal, grant_id: &str) -> Result<Vec<ApprovalUse>, JournalError> {
832 if !j.exists() {
833 return Ok(Vec::new());
834 }
835 let index_path = j.by_grant_path(grant_id);
836 if !index_path.exists() {
837 return Ok(Vec::new());
838 }
839 let raw = fs::read_to_string(&index_path)?;
840 let mut out = Vec::new();
841 for line in raw.lines() {
842 let idx: u64 = match line.trim().parse() {
843 Ok(n) => n,
844 Err(_) => continue,
845 };
846 if let Some(rec) = load_use_record(j, idx)? {
847 out.push(rec);
848 }
849 }
850 Ok(out)
851}
852
853#[cfg(test)]
858mod tests {
859 use super::*;
860 use tempfile::tempdir;
861
862 fn sample_use(use_id: &str, grant_id: &str, nonce_digest: &str, n: u32) -> ApprovalUse {
863 ApprovalUse {
864 type_: TYPE_APPROVAL_USE.into(),
865 use_id: use_id.into(),
866 grant_id: grant_id.into(),
867 grant_digest: "sha256:00".into(),
868 nonce_digest: nonce_digest.into(),
869 actor: "agent://deployer".into(),
870 action: "deploy.production".into(),
871 subject: "env://production".into(),
872 session_id: None,
873 action_artifact_id: None,
874 receipt_digest: None,
875 use_number: n,
876 max_uses: Some(2),
877 idempotency_key: None,
878 created_at: "2026-04-30T07:00:00Z".into(),
879 expires_at: None,
880 previous_record_digest: String::new(), record_digest: String::new(), signature: None,
883 signature_alg: None,
884 signing_key_id: None,
885 }
886 }
887
888 #[test]
889 fn first_append_creates_layout_and_head() {
890 let dir = tempdir().unwrap();
891 let j = Journal::new(dir.path());
892 let head = append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
893 assert_eq!(head.index, 1);
894 assert!(j.records_dir().is_dir());
895 assert!(j.heads_dir().is_dir());
896 assert!(j.current_head_path().is_file());
897 assert!(j.meta_path().is_file());
898 assert!(j.by_grant_path("g1").is_file());
900 assert!(j.by_nonce_path("sha256:nn1").is_file());
901 }
902
903 #[test]
904 fn second_append_links_previous_record_digest() {
905 let dir = tempdir().unwrap();
906 let j = Journal::new(dir.path());
907 let h1 = append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
908 let h2 = append_use(&j, sample_use("use_2", "g1", "sha256:nn2", 2)).unwrap();
909 assert_eq!(h2.index, 2);
910 let recs = iter_records(&j).unwrap();
912 assert_eq!(recs.len(), 2);
913 let (_, _, bytes) = &recs[1];
914 let r2: ApprovalUse = serde_json::from_slice(bytes).unwrap();
915 assert_eq!(r2.previous_record_digest, h1.digest);
916 }
917
918 #[test]
919 fn verify_integrity_passes_on_intact_chain() {
920 let dir = tempdir().unwrap();
921 let j = Journal::new(dir.path());
922 for i in 1..=5 {
923 let nd = format!("sha256:nn{i}");
924 append_use(&j, sample_use(&format!("use_{i}"), "g1", &nd, i)).unwrap();
925 }
926 assert_eq!(verify_integrity(&j).unwrap(), 5);
927 }
928
929 #[test]
930 fn editing_a_record_breaks_integrity() {
931 let dir = tempdir().unwrap();
932 let j = Journal::new(dir.path());
933 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
934 let entries: Vec<_> = fs::read_dir(j.records_dir()).unwrap().collect();
936 let entry = entries.into_iter().next().unwrap().unwrap();
937 let mut json: serde_json::Value =
938 serde_json::from_slice(&fs::read(entry.path()).unwrap()).unwrap();
939 json["actor"] = "agent://attacker".into();
940 fs::write(entry.path(), serde_json::to_vec_pretty(&json).unwrap()).unwrap();
941
942 let err = verify_integrity(&j).unwrap_err();
943 assert!(
944 matches!(err, JournalError::RecordTampered { .. }),
945 "expected RecordTampered, got {err:?}"
946 );
947 }
948
949 #[test]
950 fn deleting_a_record_breaks_integrity_or_head_continuity() {
951 let dir = tempdir().unwrap();
952 let j = Journal::new(dir.path());
953 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
954 append_use(&j, sample_use("use_2", "g1", "sha256:nn2", 2)).unwrap();
955 let entries: Vec<_> = fs::read_dir(j.records_dir())
957 .unwrap()
958 .map(|e| e.unwrap().path())
959 .collect();
960 let trailing = entries.iter().max().unwrap();
961 fs::remove_file(trailing).unwrap();
962
963 let err = verify_integrity(&j).unwrap_err();
964 assert!(
965 matches!(err, JournalError::MissingRecord { .. }),
966 "expected MissingRecord, got {err:?}"
967 );
968 }
969
970 #[test]
971 fn indexes_can_be_rebuilt_from_records() {
972 let dir = tempdir().unwrap();
973 let j = Journal::new(dir.path());
974 for i in 1..=3 {
975 let nd = format!("sha256:nn{i}");
976 append_use(&j, sample_use(&format!("use_{i}"), "g1", &nd, i)).unwrap();
977 }
978 fs::remove_dir_all(j.indexes_dir()).unwrap();
980
981 let rebuilt = rebuild_indexes(&j).unwrap();
982 assert_eq!(rebuilt, 3);
983 assert!(j.by_grant_path("g1").is_file());
984 assert!(j.by_nonce_path("sha256:nn1").is_file());
985 }
986
987 #[test]
988 fn check_replay_reports_use_count_and_max() {
989 let dir = tempdir().unwrap();
990 let j = Journal::new(dir.path());
991 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
993 append_use(&j, sample_use("use_2", "g1", "sha256:nn1", 2)).unwrap();
994
995 let r = check_replay(&j, "g1", "sha256:nn1", Some(2)).unwrap();
997 assert_eq!(r.level, ReplayCheckLevel::LocalJournal);
998 assert_eq!(r.use_number, Some(3));
999 assert_eq!(r.max_uses, Some(2));
1000 assert_eq!(r.passed, Some(false));
1001 }
1002
1003 #[test]
1004 fn check_replay_passes_when_under_max() {
1005 let dir = tempdir().unwrap();
1006 let j = Journal::new(dir.path());
1007 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1008 let r = check_replay(&j, "g1", "sha256:nn1", Some(2)).unwrap();
1009 assert_eq!(r.use_number, Some(2));
1010 assert_eq!(r.passed, Some(true));
1011 }
1012
1013 #[test]
1014 fn check_replay_no_journal_returns_not_performed() {
1015 let dir = tempdir().unwrap();
1016 let absent = dir.path().join("nope");
1017 let j = Journal::new(&absent);
1018 let r = check_replay(&j, "g1", "sha256:nn1", Some(1)).unwrap();
1019 assert_eq!(r.level, ReplayCheckLevel::NotPerformed);
1020 assert!(r.use_number.is_none());
1021 }
1022
1023 #[test]
1024 fn check_replay_unbounded_grant_passes_with_count() {
1025 let dir = tempdir().unwrap();
1026 let j = Journal::new(dir.path());
1027 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1028 let mut u = sample_use("use_2", "g2", "sha256:other", 1);
1032 u.max_uses = None;
1033 append_use(&j, u).unwrap();
1034
1035 let r = check_replay(&j, "g2", "sha256:other", None).unwrap();
1036 assert!(r.passed.unwrap());
1037 assert!(r.max_uses.is_none());
1038 }
1039
1040 #[test]
1041 fn list_uses_for_grant_returns_records_in_order() {
1042 let dir = tempdir().unwrap();
1043 let j = Journal::new(dir.path());
1044 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1045 append_use(&j, sample_use("use_2", "g2", "sha256:nn2", 1)).unwrap();
1046 append_use(&j, sample_use("use_3", "g1", "sha256:nn3", 2)).unwrap();
1047 let g1 = list_uses_for_grant(&j, "g1").unwrap();
1048 assert_eq!(g1.len(), 2);
1049 assert_eq!(g1[0].use_id, "use_1");
1050 assert_eq!(g1[1].use_id, "use_3");
1051 }
1052
1053 #[test]
1054 fn lock_keeps_two_appends_serial() {
1055 let dir = tempdir().unwrap();
1058 let j = Journal::new(dir.path());
1059 fs::create_dir_all(j.locks_dir()).unwrap();
1060 let held = OpenOptions::new()
1061 .read(true)
1062 .write(true)
1063 .create(true)
1064 .truncate(false)
1065 .open(j.lock_path())
1066 .unwrap();
1067 held.try_lock_exclusive().unwrap();
1068
1069 let err = append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap_err();
1070 assert!(matches!(err, JournalError::LockBusy));
1071
1072 let _ = fs2::FileExt::unlock(&held);
1073 }
1074
1075 #[test]
1076 fn revocation_appends_into_chain() {
1077 let dir = tempdir().unwrap();
1078 let j = Journal::new(dir.path());
1079 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1080 let rev = ApprovalRevocation {
1081 type_: TYPE_APPROVAL_REVOCATION.into(),
1082 revocation_id: "rev_1".into(),
1083 grant_id: "g1".into(),
1084 grant_digest: "sha256:00".into(),
1085 revoker: "human://alice".into(),
1086 reason: Some("rotated key".into()),
1087 created_at: "2026-04-30T07:01:00Z".into(),
1088 previous_record_digest: String::new(),
1089 record_digest: String::new(),
1090 signature: None,
1091 signature_alg: None,
1092 signing_key_id: None,
1093 };
1094 let h = append_revocation(&j, rev).unwrap();
1095 assert_eq!(h.index, 2);
1096 assert_eq!(verify_integrity(&j).unwrap(), 2);
1097 }
1098
1099 #[test]
1100 fn record_files_contain_no_raw_nonce_or_signature_secrets() {
1101 let dir = tempdir().unwrap();
1106 let j = Journal::new(dir.path());
1107 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1108 let entries: Vec<_> = fs::read_dir(j.records_dir())
1109 .unwrap()
1110 .map(|e| e.unwrap().path())
1111 .collect();
1112 let bytes = fs::read(&entries[0]).unwrap();
1113 let json: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
1114 let obj = json.as_object().unwrap();
1115 for forbidden in [
1116 "nonce",
1117 "command",
1118 "prompt",
1119 "file_content",
1120 "bearer_token",
1121 "api_key",
1122 ] {
1123 assert!(
1124 !obj.contains_key(forbidden),
1125 "journal record must not contain `{forbidden}`",
1126 );
1127 }
1128 assert!(obj.contains_key("nonce_digest"));
1130 }
1131
1132 #[test]
1135 fn reserve_use_first_call_succeeds_and_stamps_use_number() {
1136 let dir = tempdir().unwrap();
1139 let j = Journal::new(dir.path());
1140 let mut rec = sample_use("use_1", "g1", "sha256:nn1", 0);
1141 rec.use_number = 0;
1142 let head = reserve_use(&j, rec, Some(1)).unwrap();
1143 assert_eq!(head.index, 1);
1144 let stored = list_uses_for_grant(&j, "g1").unwrap();
1145 assert_eq!(stored.len(), 1);
1146 assert_eq!(
1147 stored[0].use_number, 1,
1148 "reserve_use must stamp use_number=1 for the first use"
1149 );
1150 }
1151
1152 #[test]
1153 fn reserve_use_max_uses_1_serial_second_call_rejects() {
1154 let dir = tempdir().unwrap();
1157 let j = Journal::new(dir.path());
1158 reserve_use(&j, sample_use("use_1", "g1", "sha256:nn_a", 0), Some(1)).unwrap();
1159
1160 let err = reserve_use(&j, sample_use("use_2", "g1", "sha256:nn_a", 0), Some(1))
1161 .expect_err("second consume of max_uses=1 grant must fail");
1162 match err {
1163 JournalError::MaxUsesExceeded {
1164 grant_id,
1165 max_uses,
1166 current,
1167 } => {
1168 assert_eq!(grant_id, "g1");
1169 assert_eq!(max_uses, 1);
1170 assert_eq!(current, 1);
1171 }
1172 other => panic!("expected MaxUsesExceeded, got {other:?}"),
1173 }
1174 let stored = list_uses_for_grant(&j, "g1").unwrap();
1176 assert_eq!(stored.len(), 1, "rejected reserve must not append");
1177 }
1178
1179 #[test]
1180 fn reserve_use_max_uses_2_two_uses_pass_third_rejects() {
1181 let dir = tempdir().unwrap();
1183 let j = Journal::new(dir.path());
1184 let mut a = sample_use("use_1", "g1", "sha256:nn_a", 0);
1185 a.max_uses = Some(2);
1186 let mut b = sample_use("use_2", "g1", "sha256:nn_b", 0);
1187 b.max_uses = Some(2);
1188 reserve_use(&j, a, Some(2)).unwrap();
1189 reserve_use(&j, b, Some(2)).unwrap();
1190 let mut c = sample_use("use_3", "g1", "sha256:nn_c", 0);
1193 c.max_uses = Some(2);
1194 reserve_use(&j, c, Some(2)).unwrap();
1195 let mut a2 = sample_use("use_1b", "g1", "sha256:nn_a", 0);
1199 a2.max_uses = Some(2);
1200 reserve_use(&j, a2, Some(2)).unwrap();
1201 let mut a3 = sample_use("use_1c", "g1", "sha256:nn_a", 0);
1203 a3.max_uses = Some(2);
1204 let err = reserve_use(&j, a3, Some(2)).expect_err("third use of same nonce must fail");
1205 assert!(matches!(err, JournalError::MaxUsesExceeded { .. }));
1206 }
1207
1208 #[test]
1221 fn reserve_use_retry_after_lock_busy_does_not_bypass_max_uses() {
1222 let dir = tempdir().unwrap();
1223 let j = Journal::new(dir.path());
1224 reserve_use(&j, sample_use("use_1", "g1", "sha256:nn_retry", 0), Some(1)).unwrap();
1226 for i in 0..5 {
1229 let err = reserve_use(
1230 &j,
1231 sample_use(&format!("use_retry_{i}"), "g1", "sha256:nn_retry", 0),
1232 Some(1),
1233 )
1234 .expect_err("retry must fail");
1235 assert!(matches!(err, JournalError::MaxUsesExceeded { .. }));
1236 }
1237 let stored = list_uses_for_grant(&j, "g1").unwrap();
1238 assert_eq!(
1239 stored.len(),
1240 1,
1241 "exactly one record on disk despite 5 retries"
1242 );
1243 }
1244
1245 #[test]
1246 fn reserve_use_concurrent_max_uses_1_only_one_succeeds() {
1247 use std::sync::atomic::{AtomicUsize, Ordering};
1261 use std::sync::Arc;
1262 use std::thread;
1263
1264 let dir = tempdir().unwrap();
1265 let dir_path = Arc::new(dir.path().to_path_buf());
1266 let success = Arc::new(AtomicUsize::new(0));
1267 let lock_busy = Arc::new(AtomicUsize::new(0));
1268 let max_exceeded = Arc::new(AtomicUsize::new(0));
1269
1270 let mut handles = Vec::new();
1271 for i in 0..8 {
1272 let dir_path = Arc::clone(&dir_path);
1273 let success = Arc::clone(&success);
1274 let lock_busy = Arc::clone(&lock_busy);
1275 let max_exceeded = Arc::clone(&max_exceeded);
1276 handles.push(thread::spawn(move || {
1277 let j = Journal::new(dir_path.as_path());
1278 let rec = sample_use(&format!("use_{i}"), "g1", "sha256:race_nonce", 0);
1279 match reserve_use(&j, rec, Some(1)) {
1280 Ok(_) => {
1281 success.fetch_add(1, Ordering::SeqCst);
1282 }
1283 Err(JournalError::LockBusy) => {
1284 lock_busy.fetch_add(1, Ordering::SeqCst);
1285 }
1286 Err(JournalError::MaxUsesExceeded { .. }) => {
1287 max_exceeded.fetch_add(1, Ordering::SeqCst);
1288 }
1289 Err(other) => panic!("unexpected error: {other:?}"),
1290 }
1291 }));
1292 }
1293 for h in handles {
1294 h.join().unwrap();
1295 }
1296
1297 let s = success.load(Ordering::SeqCst);
1298 let lb = lock_busy.load(Ordering::SeqCst);
1299 let me = max_exceeded.load(Ordering::SeqCst);
1300 assert_eq!(s, 1, "exactly one of 8 concurrent reserves must succeed; got {s} (lock_busy={lb}, max_exceeded={me})");
1301 assert_eq!(s + lb + me, 8, "every thread accounted for");
1302
1303 let stored = list_uses_for_grant(&Journal::new(dir.path()), "g1").unwrap();
1305 let same_nonce = stored
1306 .iter()
1307 .filter(|u| u.nonce_digest == "sha256:race_nonce")
1308 .count();
1309 assert_eq!(
1310 same_nonce, 1,
1311 "exactly one record on disk for the contested nonce"
1312 );
1313 }
1314}