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
732pub 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
775pub fn find_use_for_action(
793 j: &Journal,
794 grant_id: &str,
795 nonce_digest: &str,
796 max_uses_hint: Option<u32>,
797) -> Result<Option<(ApprovalUse, ReplayCheck)>, JournalError> {
798 if !j.exists() {
799 return Ok(None);
800 }
801 let index_path = j.by_nonce_path(nonce_digest);
802 if !index_path.exists() {
803 return Ok(None);
804 }
805 let raw = fs::read_to_string(&index_path)?;
806 let mut latest: Option<ApprovalUse> = None;
813 for line in raw.lines() {
814 let idx: u64 = match line.trim().parse() {
815 Ok(n) => n,
816 Err(_) => continue,
817 };
818 if let Some(rec) = load_use_record(j, idx)? {
819 if rec.grant_id == grant_id {
820 latest = Some(rec);
821 }
822 }
823 }
824 let Some(rec) = latest else { return Ok(None) };
825
826 let stored_max = rec.max_uses;
827 let max_uses = max_uses_hint.or(stored_max);
828 let passed = match max_uses {
829 Some(m) => rec.use_number <= m,
830 None => true,
831 };
832 let details = match max_uses {
833 Some(m) => format!(
834 "local Approval Use Journal passed, use {}/{}",
835 rec.use_number, m
836 ),
837 None => format!(
838 "local Approval Use Journal: use {} of unbounded grant",
839 rec.use_number
840 ),
841 };
842 Ok(Some((
843 rec.clone(),
844 ReplayCheck {
845 level: ReplayCheckLevel::LocalJournal,
846 use_number: Some(rec.use_number),
847 max_uses,
848 passed: Some(passed),
849 details: Some(details),
850 },
851 )))
852}
853
854pub fn list_uses_for_grant(j: &Journal, grant_id: &str) -> Result<Vec<ApprovalUse>, JournalError> {
857 if !j.exists() {
858 return Ok(Vec::new());
859 }
860 let index_path = j.by_grant_path(grant_id);
861 if !index_path.exists() {
862 return Ok(Vec::new());
863 }
864 let raw = fs::read_to_string(&index_path)?;
865 let mut out = Vec::new();
866 for line in raw.lines() {
867 let idx: u64 = match line.trim().parse() {
868 Ok(n) => n,
869 Err(_) => continue,
870 };
871 if let Some(rec) = load_use_record(j, idx)? {
872 out.push(rec);
873 }
874 }
875 Ok(out)
876}
877
878#[cfg(test)]
883mod tests {
884 use super::*;
885 use tempfile::tempdir;
886
887 fn sample_use(use_id: &str, grant_id: &str, nonce_digest: &str, n: u32) -> ApprovalUse {
888 ApprovalUse {
889 type_: TYPE_APPROVAL_USE.into(),
890 use_id: use_id.into(),
891 grant_id: grant_id.into(),
892 grant_digest: "sha256:00".into(),
893 nonce_digest: nonce_digest.into(),
894 actor: "agent://deployer".into(),
895 action: "deploy.production".into(),
896 subject: "env://production".into(),
897 session_id: None,
898 action_artifact_id: None,
899 receipt_digest: None,
900 use_number: n,
901 max_uses: Some(2),
902 idempotency_key: None,
903 created_at: "2026-04-30T07:00:00Z".into(),
904 expires_at: None,
905 previous_record_digest: String::new(), record_digest: String::new(), signature: None,
908 signature_alg: None,
909 signing_key_id: None,
910 }
911 }
912
913 #[test]
914 fn first_append_creates_layout_and_head() {
915 let dir = tempdir().unwrap();
916 let j = Journal::new(dir.path());
917 let head = append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
918 assert_eq!(head.index, 1);
919 assert!(j.records_dir().is_dir());
920 assert!(j.heads_dir().is_dir());
921 assert!(j.current_head_path().is_file());
922 assert!(j.meta_path().is_file());
923 assert!(j.by_grant_path("g1").is_file());
925 assert!(j.by_nonce_path("sha256:nn1").is_file());
926 }
927
928 #[test]
929 fn second_append_links_previous_record_digest() {
930 let dir = tempdir().unwrap();
931 let j = Journal::new(dir.path());
932 let h1 = append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
933 let h2 = append_use(&j, sample_use("use_2", "g1", "sha256:nn2", 2)).unwrap();
934 assert_eq!(h2.index, 2);
935 let recs = iter_records(&j).unwrap();
937 assert_eq!(recs.len(), 2);
938 let (_, _, bytes) = &recs[1];
939 let r2: ApprovalUse = serde_json::from_slice(bytes).unwrap();
940 assert_eq!(r2.previous_record_digest, h1.digest);
941 }
942
943 #[test]
944 fn verify_integrity_passes_on_intact_chain() {
945 let dir = tempdir().unwrap();
946 let j = Journal::new(dir.path());
947 for i in 1..=5 {
948 let nd = format!("sha256:nn{i}");
949 append_use(&j, sample_use(&format!("use_{i}"), "g1", &nd, i)).unwrap();
950 }
951 assert_eq!(verify_integrity(&j).unwrap(), 5);
952 }
953
954 #[test]
955 fn editing_a_record_breaks_integrity() {
956 let dir = tempdir().unwrap();
957 let j = Journal::new(dir.path());
958 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
959 let entries: Vec<_> = fs::read_dir(j.records_dir()).unwrap().collect();
961 let entry = entries.into_iter().next().unwrap().unwrap();
962 let mut json: serde_json::Value =
963 serde_json::from_slice(&fs::read(entry.path()).unwrap()).unwrap();
964 json["actor"] = "agent://attacker".into();
965 fs::write(entry.path(), serde_json::to_vec_pretty(&json).unwrap()).unwrap();
966
967 let err = verify_integrity(&j).unwrap_err();
968 assert!(
969 matches!(err, JournalError::RecordTampered { .. }),
970 "expected RecordTampered, got {err:?}"
971 );
972 }
973
974 #[test]
975 fn deleting_a_record_breaks_integrity_or_head_continuity() {
976 let dir = tempdir().unwrap();
977 let j = Journal::new(dir.path());
978 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
979 append_use(&j, sample_use("use_2", "g1", "sha256:nn2", 2)).unwrap();
980 let entries: Vec<_> = fs::read_dir(j.records_dir())
982 .unwrap()
983 .map(|e| e.unwrap().path())
984 .collect();
985 let trailing = entries.iter().max().unwrap();
986 fs::remove_file(trailing).unwrap();
987
988 let err = verify_integrity(&j).unwrap_err();
989 assert!(
990 matches!(err, JournalError::MissingRecord { .. }),
991 "expected MissingRecord, got {err:?}"
992 );
993 }
994
995 #[test]
996 fn indexes_can_be_rebuilt_from_records() {
997 let dir = tempdir().unwrap();
998 let j = Journal::new(dir.path());
999 for i in 1..=3 {
1000 let nd = format!("sha256:nn{i}");
1001 append_use(&j, sample_use(&format!("use_{i}"), "g1", &nd, i)).unwrap();
1002 }
1003 fs::remove_dir_all(j.indexes_dir()).unwrap();
1005
1006 let rebuilt = rebuild_indexes(&j).unwrap();
1007 assert_eq!(rebuilt, 3);
1008 assert!(j.by_grant_path("g1").is_file());
1009 assert!(j.by_nonce_path("sha256:nn1").is_file());
1010 }
1011
1012 #[test]
1013 fn check_replay_reports_use_count_and_max() {
1014 let dir = tempdir().unwrap();
1015 let j = Journal::new(dir.path());
1016 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1018 append_use(&j, sample_use("use_2", "g1", "sha256:nn1", 2)).unwrap();
1019
1020 let r = check_replay(&j, "g1", "sha256:nn1", Some(2)).unwrap();
1022 assert_eq!(r.level, ReplayCheckLevel::LocalJournal);
1023 assert_eq!(r.use_number, Some(3));
1024 assert_eq!(r.max_uses, Some(2));
1025 assert_eq!(r.passed, Some(false));
1026 }
1027
1028 #[test]
1029 fn check_replay_passes_when_under_max() {
1030 let dir = tempdir().unwrap();
1031 let j = Journal::new(dir.path());
1032 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1033 let r = check_replay(&j, "g1", "sha256:nn1", Some(2)).unwrap();
1034 assert_eq!(r.use_number, Some(2));
1035 assert_eq!(r.passed, Some(true));
1036 }
1037
1038 #[test]
1039 fn check_replay_no_journal_returns_not_performed() {
1040 let dir = tempdir().unwrap();
1041 let absent = dir.path().join("nope");
1042 let j = Journal::new(&absent);
1043 let r = check_replay(&j, "g1", "sha256:nn1", Some(1)).unwrap();
1044 assert_eq!(r.level, ReplayCheckLevel::NotPerformed);
1045 assert!(r.use_number.is_none());
1046 }
1047
1048 #[test]
1049 fn check_replay_unbounded_grant_passes_with_count() {
1050 let dir = tempdir().unwrap();
1051 let j = Journal::new(dir.path());
1052 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1053 let mut u = sample_use("use_2", "g2", "sha256:other", 1);
1057 u.max_uses = None;
1058 append_use(&j, u).unwrap();
1059
1060 let r = check_replay(&j, "g2", "sha256:other", None).unwrap();
1061 assert!(r.passed.unwrap());
1062 assert!(r.max_uses.is_none());
1063 }
1064
1065 #[test]
1066 fn list_uses_for_grant_returns_records_in_order() {
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 append_use(&j, sample_use("use_2", "g2", "sha256:nn2", 1)).unwrap();
1071 append_use(&j, sample_use("use_3", "g1", "sha256:nn3", 2)).unwrap();
1072 let g1 = list_uses_for_grant(&j, "g1").unwrap();
1073 assert_eq!(g1.len(), 2);
1074 assert_eq!(g1[0].use_id, "use_1");
1075 assert_eq!(g1[1].use_id, "use_3");
1076 }
1077
1078 #[test]
1079 fn lock_keeps_two_appends_serial() {
1080 let dir = tempdir().unwrap();
1083 let j = Journal::new(dir.path());
1084 fs::create_dir_all(j.locks_dir()).unwrap();
1085 let held = OpenOptions::new()
1086 .read(true)
1087 .write(true)
1088 .create(true)
1089 .truncate(false)
1090 .open(j.lock_path())
1091 .unwrap();
1092 held.try_lock_exclusive().unwrap();
1093
1094 let err = append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap_err();
1095 assert!(matches!(err, JournalError::LockBusy));
1096
1097 let _ = fs2::FileExt::unlock(&held);
1098 }
1099
1100 #[test]
1101 fn revocation_appends_into_chain() {
1102 let dir = tempdir().unwrap();
1103 let j = Journal::new(dir.path());
1104 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1105 let rev = ApprovalRevocation {
1106 type_: TYPE_APPROVAL_REVOCATION.into(),
1107 revocation_id: "rev_1".into(),
1108 grant_id: "g1".into(),
1109 grant_digest: "sha256:00".into(),
1110 revoker: "human://alice".into(),
1111 reason: Some("rotated key".into()),
1112 created_at: "2026-04-30T07:01:00Z".into(),
1113 previous_record_digest: String::new(),
1114 record_digest: String::new(),
1115 signature: None,
1116 signature_alg: None,
1117 signing_key_id: None,
1118 };
1119 let h = append_revocation(&j, rev).unwrap();
1120 assert_eq!(h.index, 2);
1121 assert_eq!(verify_integrity(&j).unwrap(), 2);
1122 }
1123
1124 #[test]
1125 fn record_files_contain_no_raw_nonce_or_signature_secrets() {
1126 let dir = tempdir().unwrap();
1131 let j = Journal::new(dir.path());
1132 append_use(&j, sample_use("use_1", "g1", "sha256:nn1", 1)).unwrap();
1133 let entries: Vec<_> = fs::read_dir(j.records_dir())
1134 .unwrap()
1135 .map(|e| e.unwrap().path())
1136 .collect();
1137 let bytes = fs::read(&entries[0]).unwrap();
1138 let json: serde_json::Value = serde_json::from_slice(&bytes).unwrap();
1139 let obj = json.as_object().unwrap();
1140 for forbidden in [
1141 "nonce",
1142 "command",
1143 "prompt",
1144 "file_content",
1145 "bearer_token",
1146 "api_key",
1147 ] {
1148 assert!(
1149 !obj.contains_key(forbidden),
1150 "journal record must not contain `{forbidden}`",
1151 );
1152 }
1153 assert!(obj.contains_key("nonce_digest"));
1155 }
1156
1157 #[test]
1160 fn reserve_use_first_call_succeeds_and_stamps_use_number() {
1161 let dir = tempdir().unwrap();
1164 let j = Journal::new(dir.path());
1165 let mut rec = sample_use("use_1", "g1", "sha256:nn1", 0);
1166 rec.use_number = 0;
1167 let head = reserve_use(&j, rec, Some(1)).unwrap();
1168 assert_eq!(head.index, 1);
1169 let stored = list_uses_for_grant(&j, "g1").unwrap();
1170 assert_eq!(stored.len(), 1);
1171 assert_eq!(
1172 stored[0].use_number, 1,
1173 "reserve_use must stamp use_number=1 for the first use"
1174 );
1175 }
1176
1177 #[test]
1178 fn reserve_use_max_uses_1_serial_second_call_rejects() {
1179 let dir = tempdir().unwrap();
1182 let j = Journal::new(dir.path());
1183 reserve_use(&j, sample_use("use_1", "g1", "sha256:nn_a", 0), Some(1)).unwrap();
1184
1185 let err = reserve_use(&j, sample_use("use_2", "g1", "sha256:nn_a", 0), Some(1))
1186 .expect_err("second consume of max_uses=1 grant must fail");
1187 match err {
1188 JournalError::MaxUsesExceeded {
1189 grant_id,
1190 max_uses,
1191 current,
1192 } => {
1193 assert_eq!(grant_id, "g1");
1194 assert_eq!(max_uses, 1);
1195 assert_eq!(current, 1);
1196 }
1197 other => panic!("expected MaxUsesExceeded, got {other:?}"),
1198 }
1199 let stored = list_uses_for_grant(&j, "g1").unwrap();
1201 assert_eq!(stored.len(), 1, "rejected reserve must not append");
1202 }
1203
1204 #[test]
1205 fn reserve_use_max_uses_2_two_uses_pass_third_rejects() {
1206 let dir = tempdir().unwrap();
1208 let j = Journal::new(dir.path());
1209 let mut a = sample_use("use_1", "g1", "sha256:nn_a", 0);
1210 a.max_uses = Some(2);
1211 let mut b = sample_use("use_2", "g1", "sha256:nn_b", 0);
1212 b.max_uses = Some(2);
1213 reserve_use(&j, a, Some(2)).unwrap();
1214 reserve_use(&j, b, Some(2)).unwrap();
1215 let mut c = sample_use("use_3", "g1", "sha256:nn_c", 0);
1218 c.max_uses = Some(2);
1219 reserve_use(&j, c, Some(2)).unwrap();
1220 let mut a2 = sample_use("use_1b", "g1", "sha256:nn_a", 0);
1224 a2.max_uses = Some(2);
1225 reserve_use(&j, a2, Some(2)).unwrap();
1226 let mut a3 = sample_use("use_1c", "g1", "sha256:nn_a", 0);
1228 a3.max_uses = Some(2);
1229 let err = reserve_use(&j, a3, Some(2)).expect_err("third use of same nonce must fail");
1230 assert!(matches!(err, JournalError::MaxUsesExceeded { .. }));
1231 }
1232
1233 #[test]
1246 fn reserve_use_retry_after_lock_busy_does_not_bypass_max_uses() {
1247 let dir = tempdir().unwrap();
1248 let j = Journal::new(dir.path());
1249 reserve_use(&j, sample_use("use_1", "g1", "sha256:nn_retry", 0), Some(1)).unwrap();
1251 for i in 0..5 {
1254 let err = reserve_use(
1255 &j,
1256 sample_use(&format!("use_retry_{i}"), "g1", "sha256:nn_retry", 0),
1257 Some(1),
1258 )
1259 .expect_err("retry must fail");
1260 assert!(matches!(err, JournalError::MaxUsesExceeded { .. }));
1261 }
1262 let stored = list_uses_for_grant(&j, "g1").unwrap();
1263 assert_eq!(
1264 stored.len(),
1265 1,
1266 "exactly one record on disk despite 5 retries"
1267 );
1268 }
1269
1270 #[test]
1271 fn reserve_use_concurrent_max_uses_1_only_one_succeeds() {
1272 use std::sync::atomic::{AtomicUsize, Ordering};
1286 use std::sync::Arc;
1287 use std::thread;
1288
1289 let dir = tempdir().unwrap();
1290 let dir_path = Arc::new(dir.path().to_path_buf());
1291 let success = Arc::new(AtomicUsize::new(0));
1292 let lock_busy = Arc::new(AtomicUsize::new(0));
1293 let max_exceeded = Arc::new(AtomicUsize::new(0));
1294
1295 let mut handles = Vec::new();
1296 for i in 0..8 {
1297 let dir_path = Arc::clone(&dir_path);
1298 let success = Arc::clone(&success);
1299 let lock_busy = Arc::clone(&lock_busy);
1300 let max_exceeded = Arc::clone(&max_exceeded);
1301 handles.push(thread::spawn(move || {
1302 let j = Journal::new(dir_path.as_path());
1303 let rec = sample_use(&format!("use_{i}"), "g1", "sha256:race_nonce", 0);
1304 match reserve_use(&j, rec, Some(1)) {
1305 Ok(_) => {
1306 success.fetch_add(1, Ordering::SeqCst);
1307 }
1308 Err(JournalError::LockBusy) => {
1309 lock_busy.fetch_add(1, Ordering::SeqCst);
1310 }
1311 Err(JournalError::MaxUsesExceeded { .. }) => {
1312 max_exceeded.fetch_add(1, Ordering::SeqCst);
1313 }
1314 Err(other) => panic!("unexpected error: {other:?}"),
1315 }
1316 }));
1317 }
1318 for h in handles {
1319 h.join().unwrap();
1320 }
1321
1322 let s = success.load(Ordering::SeqCst);
1323 let lb = lock_busy.load(Ordering::SeqCst);
1324 let me = max_exceeded.load(Ordering::SeqCst);
1325 assert_eq!(s, 1, "exactly one of 8 concurrent reserves must succeed; got {s} (lock_busy={lb}, max_exceeded={me})");
1326 assert_eq!(s + lb + me, 8, "every thread accounted for");
1327
1328 let stored = list_uses_for_grant(&Journal::new(dir.path()), "g1").unwrap();
1330 let same_nonce = stored
1331 .iter()
1332 .filter(|u| u.nonce_digest == "sha256:race_nonce")
1333 .count();
1334 assert_eq!(
1335 same_nonce, 1,
1336 "exactly one record on disk for the contested nonce"
1337 );
1338 }
1339}