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 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 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 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
891pub 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#[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(), record_digest: String::new(), 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 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 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 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 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 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 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 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 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 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 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 assert!(obj.contains_key("nonce_digest"));
1192 }
1193
1194 #[test]
1197 fn reserve_use_first_call_succeeds_and_stamps_use_number() {
1198 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 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 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 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 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 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 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 #[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 reserve_use(&j, sample_use("use_1", "g1", "sha256:nn_retry", 0), Some(1)).unwrap();
1288 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 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 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}