1use std::io;
53use std::path::{Path, PathBuf};
54
55use crate::ports::{Clock, GitRepo, IdGen, JournalLock, JournalStore};
56use crate::protocol::journal::{
57 EventKind, JournalEvent, Phase, PhaseOutcome, PhaseRecord, RunState, RunStatus,
58 JOURNAL_SCHEMA_VERSION,
59};
60
61#[derive(Debug, Clone)]
68pub struct JournalPaths {
69 releases_dir: PathBuf,
70}
71
72impl JournalPaths {
73 pub fn new(releases_dir: impl Into<PathBuf>) -> Self {
76 Self {
77 releases_dir: releases_dir.into(),
78 }
79 }
80
81 pub fn from_git(git: &dyn GitRepo, override_dir: Option<&Path>) -> io::Result<Self> {
88 let releases_dir = match override_dir {
89 Some(dir) => dir.to_path_buf(),
90 None => git.git_common_dir()?.join("ossctl").join("releases"),
91 };
92 Ok(Self { releases_dir })
93 }
94
95 #[must_use]
97 pub fn releases_dir(&self) -> &Path {
98 &self.releases_dir
99 }
100
101 #[must_use]
103 pub fn lock_file(&self) -> PathBuf {
104 self.releases_dir.join(".lock")
105 }
106
107 #[must_use]
109 pub fn run_dir(&self, run_id: &str) -> PathBuf {
110 self.releases_dir.join(run_id)
111 }
112
113 #[must_use]
115 pub fn journal_file(&self, run_id: &str) -> PathBuf {
116 self.run_dir(run_id).join("journal.jsonl")
117 }
118
119 #[must_use]
121 pub fn manifest_file(&self, run_id: &str) -> PathBuf {
122 self.run_dir(run_id).join("manifest.json")
123 }
124}
125
126#[must_use]
135pub fn reduce(events: &[JournalEvent]) -> RunState {
136 let mut ordered: Vec<&JournalEvent> = events.iter().collect();
137 ordered.sort_by_key(|e| e.seq);
138 let mut state = RunState::empty();
139 for ev in ordered {
140 apply(&mut state, ev);
141 }
142 state
143}
144
145pub fn apply(state: &mut RunState, event: &JournalEvent) {
158 if event.seq <= state.applied_seq {
161 return;
162 }
163 if matches!(state.status, RunStatus::Completed | RunStatus::Abandoned) {
167 state.applied_seq = event.seq;
168 return;
169 }
170 match &event.kind {
171 EventKind::RunCreated {
172 run_id,
173 plan_id,
174 version,
175 targets,
176 } => {
177 state.run_id.clone_from(run_id);
178 state.plan_id.clone_from(plan_id);
179 state.version.clone_from(version);
180 state.targets.clone_from(targets);
181 state.created_ts = event.ts;
182 state.status = RunStatus::InProgress;
183 }
184 EventKind::PhaseEntered { phase } => {
185 state.current_phase = Some(*phase);
186 }
187 EventKind::PhaseCompleted { phase, outcome } => {
188 upsert_phase(&mut state.phases, *phase, *outcome);
189 if state.current_phase == Some(*phase) {
190 state.current_phase = None;
191 }
192 let completes = match phase {
207 Phase::Dist => true,
208 Phase::Tag => event.schema_version < 2,
209 _ => false,
210 };
211 if completes && *outcome == PhaseOutcome::Ok {
212 state.status = RunStatus::Completed;
213 }
214 }
215 EventKind::TargetDryRun { target } => {
216 state.dry_run.insert(target.clone());
217 }
218 EventKind::TargetBuilt { target } => {
219 state.built.insert(target.clone());
220 }
221 EventKind::TargetPublished { target, receipt } => {
222 state.published.insert(target.clone(), receipt.clone());
223 }
224 EventKind::TargetCancelled { target, reason } => {
225 state.cancelled.insert(target.clone(), reason.clone());
226 }
227 EventKind::TargetDelegated { target, .. } => {
228 state.delegated.insert(target.clone());
229 }
230 EventKind::TagCreatedLocal { tag } => {
231 state.tags.entry(tag.clone()).or_default().created_local = true;
232 }
233 EventKind::TagPushedRemote { tag } => {
234 state.tags.entry(tag.clone()).or_default().pushed_remote = true;
235 }
236 EventKind::GithubReleaseCreated { tag, url } => {
237 let t = state.tags.entry(tag.clone()).or_default();
238 t.github_release = true;
239 t.github_release_url.clone_from(url);
240 }
241 EventKind::GithubReleaseDelegated { tag, .. } => {
242 state
243 .tags
244 .entry(tag.clone())
245 .or_default()
246 .github_release_delegated = true;
247 }
248 EventKind::RunAbandoned { reason } => {
249 state.status = RunStatus::Abandoned;
250 state.abandon_reason = Some(reason.clone());
251 }
252 }
253 state.applied_seq = event.seq;
254 state.updated_ts = event.ts;
255}
256
257fn upsert_phase(phases: &mut Vec<PhaseRecord>, phase: Phase, outcome: PhaseOutcome) {
260 if let Some(rec) = phases.iter_mut().find(|r| r.phase == phase) {
261 rec.outcome = outcome;
262 } else {
263 phases.push(PhaseRecord { phase, outcome });
264 phases.sort_by_key(|r| r.phase);
265 }
266}
267
268pub fn read_events(store: &dyn JournalStore, path: &Path) -> io::Result<Vec<JournalEvent>> {
281 #[derive(serde::Deserialize)]
285 struct Envelope {
286 schema_version: u32,
287 }
288
289 let lines = store.read_lines(path)?;
290 let mut events = Vec::with_capacity(lines.len());
291 for (idx, line) in lines.iter().enumerate() {
292 let trimmed = line.trim();
293 if trimmed.is_empty() {
294 continue;
295 }
296 if let Ok(envelope) = serde_json::from_str::<Envelope>(trimmed) {
300 if envelope.schema_version > JOURNAL_SCHEMA_VERSION {
301 return Err(io::Error::new(
302 io::ErrorKind::InvalidData,
303 format!(
304 "release journal {}: line {} has schema_version {} but this \
305 ossctl understands at most {}; upgrade ossctl to resume this run",
306 path.display(),
307 idx + 1,
308 envelope.schema_version,
309 JOURNAL_SCHEMA_VERSION
310 ),
311 ));
312 }
313 }
314 let event: JournalEvent = serde_json::from_str(trimmed).map_err(|e| {
315 io::Error::new(
316 io::ErrorKind::InvalidData,
317 format!(
318 "release journal {}: line {} is not a recognized event \
319 (corrupt, or written by a newer ossctl): {e}",
320 path.display(),
321 idx + 1
322 ),
323 )
324 })?;
325 events.push(event);
326 }
327 events.sort_by_key(|e| e.seq);
328 Ok(events)
329}
330
331fn validate_run_id(run_id: &str) -> io::Result<()> {
337 let bad = run_id.is_empty()
338 || run_id == "."
339 || run_id == ".."
340 || run_id.contains('/')
341 || run_id.contains('\\')
342 || run_id.contains('\0');
343 if bad {
344 return Err(io::Error::new(
345 io::ErrorKind::InvalidInput,
346 format!("invalid run id {run_id:?}: must be a single path segment"),
347 ));
348 }
349 Ok(())
350}
351
352pub fn load_state(
368 store: &dyn JournalStore,
369 paths: &JournalPaths,
370 run_id: &str,
371) -> io::Result<Option<RunState>> {
372 validate_run_id(run_id)?;
373 if let Some(bytes) = store.read(&paths.manifest_file(run_id))? {
375 if let Ok(state) = serde_json::from_slice::<RunState>(&bytes) {
376 if state.run_id == run_id && state.schema_version <= JOURNAL_SCHEMA_VERSION {
377 return Ok(Some(state));
378 }
379 }
380 }
382 let events = read_events(store, &paths.journal_file(run_id))?;
383 if events.is_empty() {
384 return Ok(None);
385 }
386 Ok(Some(reduce(&events)))
387}
388
389pub fn read_run_state(
406 store: &dyn JournalStore,
407 paths: &JournalPaths,
408 run_id: &str,
409) -> io::Result<Option<RunState>> {
410 validate_run_id(run_id)?;
411 let events = read_events(store, &paths.journal_file(run_id))?;
412 if events.is_empty() {
413 return Ok(None);
414 }
415 Ok(Some(reduce(&events)))
416}
417
418pub fn read_run(
436 store: &dyn JournalStore,
437 paths: &JournalPaths,
438 run_id: &str,
439) -> io::Result<Option<(Vec<JournalEvent>, RunState)>> {
440 validate_run_id(run_id)?;
441 let events = read_events(store, &paths.journal_file(run_id))?;
442 if events.is_empty() {
443 return Ok(None);
444 }
445 let state = reduce(&events);
446 Ok(Some((events, state)))
447}
448
449pub fn list_runs(store: &dyn JournalStore, paths: &JournalPaths) -> io::Result<Vec<String>> {
459 let mut runs: Vec<String> = store
460 .list_dir(paths.releases_dir())?
461 .into_iter()
462 .filter(|name| validate_run_id(name).is_ok())
463 .filter(|name| {
464 store
467 .read_lines(&paths.journal_file(name))
468 .is_ok_and(|lines| lines.iter().any(|l| !l.trim().is_empty()))
469 })
470 .collect();
471 runs.sort();
472 Ok(runs)
473}
474
475pub struct Journal<'a> {
484 store: &'a dyn JournalStore,
485 clock: &'a dyn Clock,
486 paths: JournalPaths,
487 run_id: String,
488 state: RunState,
489 _lock: Box<dyn JournalLock>,
491}
492
493impl<'a> Journal<'a> {
494 pub fn create(
502 store: &'a dyn JournalStore,
503 clock: &'a dyn Clock,
504 idgen: &dyn IdGen,
505 paths: JournalPaths,
506 plan_id: String,
507 version: String,
508 targets: Vec<String>,
509 ) -> io::Result<Self> {
510 let lock = store.lock_exclusive(&paths.lock_file())?;
511 let run_id = idgen.new_id();
512 let mut journal = Self {
513 store,
514 clock,
515 paths,
516 run_id: run_id.clone(),
517 state: RunState::empty(),
518 _lock: lock,
519 };
520 journal.append(EventKind::RunCreated {
521 run_id,
522 plan_id,
523 version,
524 targets,
525 })?;
526 Ok(journal)
527 }
528
529 pub fn open(
542 store: &'a dyn JournalStore,
543 clock: &'a dyn Clock,
544 paths: JournalPaths,
545 run_id: &str,
546 ) -> io::Result<Self> {
547 validate_run_id(run_id)?;
548 let lock = store.lock_exclusive(&paths.lock_file())?;
549 let events = read_events(store, &paths.journal_file(run_id))?;
550 if events.is_empty() {
551 return Err(io::Error::new(
552 io::ErrorKind::NotFound,
553 format!("no release journal for run {run_id}"),
554 ));
555 }
556 let state = reduce(&events);
557 let journal = Self {
558 store,
559 clock,
560 paths,
561 run_id: run_id.to_string(),
562 state,
563 _lock: lock,
564 };
565 let _ = journal.persist_manifest();
567 Ok(journal)
568 }
569
570 #[must_use]
572 pub fn run_id(&self) -> &str {
573 &self.run_id
574 }
575
576 #[must_use]
578 pub fn state(&self) -> &RunState {
579 &self.state
580 }
581
582 #[must_use]
584 pub fn paths(&self) -> &JournalPaths {
585 &self.paths
586 }
587
588 pub fn append(&mut self, kind: EventKind) -> io::Result<&RunState> {
610 let event = JournalEvent {
611 schema_version: JOURNAL_SCHEMA_VERSION,
612 seq: self.state.applied_seq + 1,
613 ts: self.clock.now_unix(),
614 idempotency_key: kind.idempotency_key(),
615 kind,
616 };
617 let line = serde_json::to_string(&event).map_err(|e| {
618 io::Error::new(io::ErrorKind::InvalidData, format!("serialize event: {e}"))
619 })?;
620 self.store
622 .append_line(&self.paths.journal_file(&self.run_id), &line)?;
623 apply(&mut self.state, &event);
625 let _ = self.persist_manifest();
627 Ok(&self.state)
628 }
629
630 fn persist_manifest(&self) -> io::Result<()> {
632 let bytes = serde_json::to_vec_pretty(&self.state).map_err(|e| {
633 io::Error::new(
634 io::ErrorKind::InvalidData,
635 format!("serialize manifest: {e}"),
636 )
637 })?;
638 self.store
639 .write_atomic(&self.paths.manifest_file(&self.run_id), &bytes)
640 }
641}
642
643#[cfg(test)]
644mod tests {
645 use super::*;
646 use crate::protocol::journal::{PublishReceipt, TagState};
647 use std::cell::RefCell;
648 use std::collections::{HashMap, HashSet};
649 use std::rc::Rc;
650
651 #[derive(Default)]
654 struct StoreInner {
655 files: HashMap<PathBuf, Vec<u8>>,
657 locked: HashSet<PathBuf>,
659 fail_next_atomic: bool,
661 }
662
663 #[derive(Clone, Default)]
664 struct FakeStore {
665 inner: Rc<RefCell<StoreInner>>,
666 }
667
668 impl FakeStore {
669 fn journal_lines(&self, path: &Path) -> Vec<String> {
670 self.inner
671 .borrow()
672 .files
673 .get(path)
674 .map(|b| {
675 String::from_utf8_lossy(b)
676 .lines()
677 .map(str::to_string)
678 .collect()
679 })
680 .unwrap_or_default()
681 }
682
683 fn arm_atomic_failure(&self) {
684 self.inner.borrow_mut().fail_next_atomic = true;
685 }
686 }
687
688 struct FakeLock {
690 inner: Rc<RefCell<StoreInner>>,
691 path: PathBuf,
692 }
693
694 impl JournalLock for FakeLock {}
695
696 impl Drop for FakeLock {
697 fn drop(&mut self) {
698 self.inner.borrow_mut().locked.remove(&self.path);
699 }
700 }
701
702 impl JournalStore for FakeStore {
703 fn lock_exclusive(&self, lock_path: &Path) -> io::Result<Box<dyn JournalLock>> {
704 let mut inner = self.inner.borrow_mut();
705 if inner.locked.contains(lock_path) {
706 return Err(io::Error::new(
707 io::ErrorKind::WouldBlock,
708 "another release cut holds the lock",
709 ));
710 }
711 inner.locked.insert(lock_path.to_path_buf());
712 Ok(Box::new(FakeLock {
713 inner: Rc::clone(&self.inner),
714 path: lock_path.to_path_buf(),
715 }))
716 }
717
718 fn append_line(&self, path: &Path, line: &str) -> io::Result<()> {
719 let mut inner = self.inner.borrow_mut();
720 let buf = inner.files.entry(path.to_path_buf()).or_default();
721 buf.extend_from_slice(line.as_bytes());
722 buf.push(b'\n');
723 Ok(())
724 }
725
726 fn read_lines(&self, path: &Path) -> io::Result<Vec<String>> {
727 Ok(self.journal_lines(path))
728 }
729
730 fn read(&self, path: &Path) -> io::Result<Option<Vec<u8>>> {
731 Ok(self.inner.borrow().files.get(path).cloned())
732 }
733
734 fn write_atomic(&self, path: &Path, bytes: &[u8]) -> io::Result<()> {
735 let mut inner = self.inner.borrow_mut();
736 if inner.fail_next_atomic {
737 inner.fail_next_atomic = false;
738 return Err(io::Error::other("injected atomic-write crash"));
739 }
740 inner.files.insert(path.to_path_buf(), bytes.to_vec());
741 Ok(())
742 }
743
744 fn list_dir(&self, dir: &Path) -> io::Result<Vec<String>> {
745 let inner = self.inner.borrow();
746 let mut names: HashSet<String> = HashSet::new();
747 for path in inner.files.keys() {
748 if let Ok(rest) = path.strip_prefix(dir) {
750 if let Some(first) = rest.components().next() {
751 names.insert(first.as_os_str().to_string_lossy().into_owned());
752 }
753 }
754 }
755 Ok(names.into_iter().collect())
756 }
757 }
758
759 struct FakeClock {
760 t: std::cell::Cell<u64>,
761 }
762 impl FakeClock {
763 fn at(t: u64) -> Self {
764 Self {
765 t: std::cell::Cell::new(t),
766 }
767 }
768 }
769 impl Clock for FakeClock {
770 fn now_unix(&self) -> u64 {
771 let now = self.t.get();
772 self.t.set(now + 1); now
774 }
775 }
776
777 struct FakeIdGen {
778 id: String,
779 }
780 impl IdGen for FakeIdGen {
781 fn new_id(&self) -> String {
782 self.id.clone()
783 }
784 }
785
786 struct FakeGit {
787 common_dir: PathBuf,
788 }
789 impl GitRepo for FakeGit {
790 fn head_commit(&self) -> io::Result<String> {
791 Ok("deadbeef".into())
792 }
793 fn is_work_tree(&self) -> bool {
794 true
795 }
796 fn shortlog(&self, _since: Option<&str>) -> io::Result<String> {
797 Ok(String::new())
798 }
799 fn tags(&self) -> io::Result<Vec<String>> {
800 Ok(Vec::new())
801 }
802 fn git_common_dir(&self) -> io::Result<PathBuf> {
803 Ok(self.common_dir.clone())
804 }
805 }
806
807 fn paths() -> JournalPaths {
808 JournalPaths::new("/repo/.git/ossctl/releases")
809 }
810
811 fn receipt(version: &str) -> PublishReceipt {
812 PublishReceipt {
813 ecosystem: "cargo".into(),
814 package: Some("ossctl".into()),
815 version: version.into(),
816 registry_url: Some("https://crates.io/crates/ossctl".into()),
817 digest: Some("sha256:abc".into()),
818 }
819 }
820
821 fn sample_events() -> Vec<JournalEvent> {
823 let kinds = vec![
824 EventKind::RunCreated {
825 run_id: "RUN01".into(),
826 plan_id: "plan-abc".into(),
827 version: "0.1.0".into(),
828 targets: vec!["cargo".into(), "npm".into()],
829 },
830 EventKind::PhaseEntered {
831 phase: Phase::DryRun,
832 },
833 EventKind::TargetDryRun {
834 target: "cargo".into(),
835 },
836 EventKind::PhaseCompleted {
837 phase: Phase::DryRun,
838 outcome: PhaseOutcome::Ok,
839 },
840 EventKind::PhaseEntered {
841 phase: Phase::Publish,
842 },
843 EventKind::TargetPublished {
844 target: "cargo".into(),
845 receipt: receipt("0.1.0"),
846 },
847 ];
848 kinds
849 .into_iter()
850 .enumerate()
851 .map(|(i, kind)| JournalEvent {
852 schema_version: JOURNAL_SCHEMA_VERSION,
853 seq: (i + 1) as u64,
854 ts: 1000 + i as u64,
855 idempotency_key: kind.idempotency_key(),
856 kind,
857 })
858 .collect()
859 }
860
861 #[test]
864 fn paths_resolve_under_git_common_dir() {
865 let git = FakeGit {
866 common_dir: PathBuf::from("/repo/.git"),
867 };
868 let p = JournalPaths::from_git(&git, None).unwrap();
869 assert_eq!(p.releases_dir(), Path::new("/repo/.git/ossctl/releases"));
870 assert_eq!(
871 p.journal_file("RUN01"),
872 Path::new("/repo/.git/ossctl/releases/RUN01/journal.jsonl")
873 );
874 assert_eq!(
875 p.manifest_file("RUN01"),
876 Path::new("/repo/.git/ossctl/releases/RUN01/manifest.json")
877 );
878 assert_eq!(p.lock_file(), Path::new("/repo/.git/ossctl/releases/.lock"));
879 }
880
881 #[test]
882 fn path_override_wins_over_git() {
883 let git = FakeGit {
884 common_dir: PathBuf::from("/repo/.git"),
885 };
886 let p = JournalPaths::from_git(&git, Some(Path::new("/ci/journal"))).unwrap();
887 assert_eq!(p.releases_dir(), Path::new("/ci/journal"));
888 }
889
890 #[test]
893 fn reduce_is_deterministic() {
894 let events = sample_events();
895 let a = reduce(&events);
896 let b = reduce(&events);
897 assert_eq!(a, b);
898 assert_eq!(a.run_id, "RUN01");
899 assert_eq!(a.plan_id, "plan-abc");
900 assert_eq!(a.targets, vec!["cargo".to_string(), "npm".to_string()]);
901 assert!(a.dry_run.contains("cargo"));
902 assert_eq!(a.published.get("cargo").unwrap().version, "0.1.0");
903 assert_eq!(a.current_phase, Some(Phase::Publish));
905 assert_eq!(a.applied_seq, 6);
906 }
907
908 #[test]
909 fn reduce_ignores_slice_order() {
910 let mut events = sample_events();
911 events.reverse();
912 let out = reduce(&events);
913 assert_eq!(out, reduce(&sample_events()));
915 }
916
917 #[test]
918 fn replaying_a_seen_event_is_a_no_op() {
919 let events = sample_events();
920 let mut state = reduce(&events);
921 let before = state.clone();
922 apply(&mut state, &events[2]);
924 assert_eq!(state, before);
925 for ev in &events {
927 apply(&mut state, ev);
928 }
929 assert_eq!(state, before);
930 }
931
932 #[test]
933 fn structural_idempotency_of_publish_and_tags() {
934 let mut state = RunState::empty();
937 let mk = |seq: u64, kind: EventKind| JournalEvent {
938 schema_version: JOURNAL_SCHEMA_VERSION,
939 seq,
940 ts: seq,
941 idempotency_key: kind.idempotency_key(),
942 kind,
943 };
944 apply(
945 &mut state,
946 &mk(
947 1,
948 EventKind::RunCreated {
949 run_id: "R".into(),
950 plan_id: "p".into(),
951 version: "0.1.0".into(),
952 targets: vec!["cargo".into()],
953 },
954 ),
955 );
956 apply(
957 &mut state,
958 &mk(
959 2,
960 EventKind::TargetPublished {
961 target: "cargo".into(),
962 receipt: receipt("0.1.0"),
963 },
964 ),
965 );
966 apply(
967 &mut state,
968 &mk(
969 3,
970 EventKind::TargetPublished {
971 target: "cargo".into(),
972 receipt: receipt("0.1.1"),
973 },
974 ),
975 );
976 assert_eq!(state.published.len(), 1);
977 assert_eq!(state.published.get("cargo").unwrap().version, "0.1.1");
978
979 apply(
980 &mut state,
981 &mk(
982 4,
983 EventKind::TagCreatedLocal {
984 tag: "v0.1.1".into(),
985 },
986 ),
987 );
988 apply(
989 &mut state,
990 &mk(
991 5,
992 EventKind::TagPushedRemote {
993 tag: "v0.1.1".into(),
994 },
995 ),
996 );
997 assert_eq!(
998 state.tags.get("v0.1.1"),
999 Some(&TagState {
1000 created_local: true,
1001 pushed_remote: true,
1002 github_release: false,
1003 github_release_url: None,
1004 github_release_delegated: false,
1005 })
1006 );
1007 }
1008
1009 #[test]
1010 fn github_release_delegation_reduces_to_the_tag_state_flag() {
1011 let mut state = RunState::empty();
1015 let mk = |seq: u64, kind: EventKind| JournalEvent {
1016 schema_version: JOURNAL_SCHEMA_VERSION,
1017 seq,
1018 ts: seq,
1019 idempotency_key: kind.idempotency_key(),
1020 kind,
1021 };
1022 apply(
1023 &mut state,
1024 &mk(
1025 1,
1026 EventKind::TagCreatedLocal {
1027 tag: "v1.0.0".into(),
1028 },
1029 ),
1030 );
1031 apply(
1032 &mut state,
1033 &mk(
1034 2,
1035 EventKind::TagPushedRemote {
1036 tag: "v1.0.0".into(),
1037 },
1038 ),
1039 );
1040 apply(
1041 &mut state,
1042 &mk(
1043 3,
1044 EventKind::GithubReleaseDelegated {
1045 tag: "v1.0.0".into(),
1046 delegated_to: "cargo-dist".into(),
1047 },
1048 ),
1049 );
1050 assert_eq!(
1051 state.tags.get("v1.0.0"),
1052 Some(&TagState {
1053 created_local: true,
1054 pushed_remote: true,
1055 github_release: false,
1056 github_release_url: None,
1057 github_release_delegated: true,
1058 })
1059 );
1060 }
1061
1062 #[test]
1063 fn dist_phase_ok_completes_the_run_tag_ok_alone_does_not() {
1064 let mut state = reduce(&sample_events());
1065 assert_eq!(state.status, RunStatus::InProgress);
1066 let mut seq = state.applied_seq;
1067 let mut push = |state: &mut RunState, phase: Phase| {
1068 seq += 1;
1069 apply(
1070 state,
1071 &JournalEvent {
1072 schema_version: JOURNAL_SCHEMA_VERSION,
1073 seq,
1074 ts: 9000 + seq,
1075 idempotency_key: format!("phase_completed:{}", phase.as_str()),
1076 kind: EventKind::PhaseCompleted {
1077 phase,
1078 outcome: PhaseOutcome::Ok,
1079 },
1080 },
1081 );
1082 };
1083 push(&mut state, Phase::Tag);
1085 assert_eq!(state.status, RunStatus::InProgress);
1086 push(&mut state, Phase::Dist);
1087 assert_eq!(state.status, RunStatus::Completed);
1088 }
1089
1090 #[test]
1091 fn a_v1_tag_ok_completes_the_run_for_backward_compat() {
1092 let mut state = reduce(&sample_events());
1097 assert_eq!(state.status, RunStatus::InProgress);
1098 let seq = state.applied_seq + 1;
1099 apply(
1100 &mut state,
1101 &JournalEvent {
1102 schema_version: 1,
1103 seq,
1104 ts: 9000,
1105 idempotency_key: "phase_completed:tag".into(),
1106 kind: EventKind::PhaseCompleted {
1107 phase: Phase::Tag,
1108 outcome: PhaseOutcome::Ok,
1109 },
1110 },
1111 );
1112 assert_eq!(state.status, RunStatus::Completed);
1113 }
1114
1115 #[test]
1116 fn run_abandoned_is_terminal_with_reason() {
1117 let mut state = reduce(&sample_events());
1118 let seq = state.applied_seq + 1;
1119 apply(
1120 &mut state,
1121 &JournalEvent {
1122 schema_version: JOURNAL_SCHEMA_VERSION,
1123 seq,
1124 ts: 9000,
1125 idempotency_key: "run_abandoned".into(),
1126 kind: EventKind::RunAbandoned {
1127 reason: "OTP timeout".into(),
1128 },
1129 },
1130 );
1131 assert_eq!(state.status, RunStatus::Abandoned);
1132 assert_eq!(state.abandon_reason.as_deref(), Some("OTP timeout"));
1133 }
1134
1135 #[test]
1138 fn create_writes_run_created_and_manifest() {
1139 let store = FakeStore::default();
1140 let clock = FakeClock::at(1000);
1141 let idgen = FakeIdGen { id: "RUN01".into() };
1142 let journal = Journal::create(
1143 &store,
1144 &clock,
1145 &idgen,
1146 paths(),
1147 "plan-abc".into(),
1148 "0.1.0".into(),
1149 vec!["cargo".into()],
1150 )
1151 .unwrap();
1152 assert_eq!(journal.run_id(), "RUN01");
1153 assert_eq!(journal.state().run_id, "RUN01");
1154 assert_eq!(journal.state().applied_seq, 1);
1155
1156 let lines = store.journal_lines(&paths().journal_file("RUN01"));
1158 assert_eq!(lines.len(), 1);
1159 let manifest = store
1161 .inner
1162 .borrow()
1163 .files
1164 .get(&paths().manifest_file("RUN01"))
1165 .cloned()
1166 .unwrap();
1167 let loaded: RunState = serde_json::from_slice(&manifest).unwrap();
1168 assert_eq!(&loaded, journal.state());
1169 }
1170
1171 #[test]
1172 fn append_records_facts_and_a_failed_phase_can_later_complete_ok() {
1173 let store = FakeStore::default();
1177 let clock = FakeClock::at(1000);
1178 let idgen = FakeIdGen { id: "RUN01".into() };
1179 let mut journal = Journal::create(
1180 &store,
1181 &clock,
1182 &idgen,
1183 paths(),
1184 "plan-abc".into(),
1185 "0.1.0".into(),
1186 vec!["cargo".into()],
1187 )
1188 .unwrap();
1189 journal
1190 .append(EventKind::PhaseCompleted {
1191 phase: Phase::Publish,
1192 outcome: PhaseOutcome::Failed,
1193 })
1194 .unwrap();
1195 journal
1196 .append(EventKind::PhaseCompleted {
1197 phase: Phase::Publish,
1198 outcome: PhaseOutcome::Ok,
1199 })
1200 .unwrap();
1201 let lines = store.journal_lines(&paths().journal_file("RUN01"));
1203 assert_eq!(lines.len(), 3);
1204 let publish = journal
1206 .state()
1207 .phases
1208 .iter()
1209 .filter(|r| r.phase == Phase::Publish)
1210 .collect::<Vec<_>>();
1211 assert_eq!(publish.len(), 1);
1212 assert_eq!(publish[0].outcome, PhaseOutcome::Ok);
1213 }
1214
1215 #[test]
1216 fn terminal_state_freezes_further_events() {
1217 let mut state = reduce(&sample_events());
1219 let published_before = state.published.clone();
1220 let mut seq = state.applied_seq;
1221 let mut next = |kind: EventKind| {
1222 seq += 1;
1223 JournalEvent {
1224 schema_version: JOURNAL_SCHEMA_VERSION,
1225 seq,
1226 ts: 9000 + seq,
1227 idempotency_key: kind.idempotency_key(),
1228 kind,
1229 }
1230 };
1231 apply(
1232 &mut state,
1233 &next(EventKind::RunAbandoned {
1234 reason: "aborted".into(),
1235 }),
1236 );
1237 assert_eq!(state.status, RunStatus::Abandoned);
1238 apply(
1240 &mut state,
1241 &next(EventKind::TargetPublished {
1242 target: "npm".into(),
1243 receipt: receipt("9.9.9"),
1244 }),
1245 );
1246 assert_eq!(state.status, RunStatus::Abandoned);
1247 assert_eq!(state.published, published_before);
1248 assert!(!state.published.contains_key("npm"));
1249 }
1250
1251 #[test]
1254 fn event_survives_a_manifest_write_crash() {
1255 let store = FakeStore::default();
1256 let clock = FakeClock::at(1000);
1257 let idgen = FakeIdGen { id: "RUN01".into() };
1258 let mut journal = Journal::create(
1259 &store,
1260 &clock,
1261 &idgen,
1262 paths(),
1263 "plan-abc".into(),
1264 "0.1.0".into(),
1265 vec!["cargo".into()],
1266 )
1267 .unwrap();
1268 store.arm_atomic_failure();
1272 journal
1273 .append(EventKind::TargetPublished {
1274 target: "cargo".into(),
1275 receipt: receipt("0.1.0"),
1276 })
1277 .unwrap();
1278 drop(journal); let clock2 = FakeClock::at(2000);
1283 let reopened = Journal::open(&store, &clock2, paths(), "RUN01").unwrap();
1284 assert!(reopened.state().published.contains_key("cargo"));
1285 assert_eq!(reopened.state().applied_seq, 2);
1286 let manifest = store
1288 .inner
1289 .borrow()
1290 .files
1291 .get(&paths().manifest_file("RUN01"))
1292 .cloned()
1293 .unwrap();
1294 let loaded: RunState = serde_json::from_slice(&manifest).unwrap();
1295 assert_eq!(&loaded, reopened.state());
1296 }
1297
1298 #[test]
1299 fn open_reduces_from_journal_when_manifest_absent() {
1300 let store = FakeStore::default();
1302 for ev in sample_events() {
1303 let line = serde_json::to_string(&ev).unwrap();
1304 store
1305 .append_line(&paths().journal_file("RUN01"), &line)
1306 .unwrap();
1307 }
1308 let clock = FakeClock::at(1000);
1309 let journal = Journal::open(&store, &clock, paths(), "RUN01").unwrap();
1310 assert_eq!(journal.state(), &reduce(&sample_events()));
1311 }
1312
1313 #[test]
1316 fn second_create_fails_while_lock_is_held() {
1317 let store = FakeStore::default();
1318 let clock = FakeClock::at(1000);
1319 let idgen = FakeIdGen { id: "RUN01".into() };
1320 let held = Journal::create(
1321 &store,
1322 &clock,
1323 &idgen,
1324 paths(),
1325 "plan-abc".into(),
1326 "0.1.0".into(),
1327 vec!["cargo".into()],
1328 )
1329 .unwrap();
1330
1331 let clock2 = FakeClock::at(2000);
1333 let idgen2 = FakeIdGen { id: "RUN02".into() };
1334 let result = Journal::create(
1337 &store,
1338 &clock2,
1339 &idgen2,
1340 paths(),
1341 "plan-def".into(),
1342 "0.1.0".into(),
1343 vec!["cargo".into()],
1344 );
1345 let err = result.err().expect("concurrent create must fail");
1346 assert_eq!(err.kind(), io::ErrorKind::WouldBlock);
1347
1348 drop(held);
1350 let clock3 = FakeClock::at(3000);
1351 let idgen3 = FakeIdGen { id: "RUN02".into() };
1352 assert!(Journal::create(
1353 &store,
1354 &clock3,
1355 &idgen3,
1356 paths(),
1357 "plan-def".into(),
1358 "0.1.0".into(),
1359 vec!["cargo".into()],
1360 )
1361 .is_ok());
1362 }
1363
1364 #[test]
1367 fn load_state_is_read_only_and_takes_no_lock() {
1368 let store = FakeStore::default();
1369 let clock = FakeClock::at(1000);
1370 let idgen = FakeIdGen { id: "RUN01".into() };
1371 let journal = Journal::create(
1372 &store,
1373 &clock,
1374 &idgen,
1375 paths(),
1376 "plan-abc".into(),
1377 "0.1.0".into(),
1378 vec!["cargo".into()],
1379 )
1380 .unwrap();
1381 let loaded = load_state(&store, &paths(), "RUN01").unwrap().unwrap();
1383 assert_eq!(&loaded, journal.state());
1384 assert!(load_state(&store, &paths(), "MISSING").unwrap().is_none());
1385 }
1386
1387 #[test]
1388 fn load_state_prefers_manifest_but_falls_back_to_journal() {
1389 let store = FakeStore::default();
1390 let clock = FakeClock::at(1000);
1391 let idgen = FakeIdGen { id: "RUN01".into() };
1392 let journal = Journal::create(
1393 &store,
1394 &clock,
1395 &idgen,
1396 paths(),
1397 "plan-abc".into(),
1398 "0.1.0".into(),
1399 vec!["cargo".into()],
1400 )
1401 .unwrap();
1402 let expected = journal.state().clone();
1403 drop(journal);
1404
1405 assert_eq!(
1407 load_state(&store, &paths(), "RUN01").unwrap(),
1408 Some(expected.clone())
1409 );
1410
1411 store
1413 .write_atomic(&paths().manifest_file("RUN01"), b"{not json")
1414 .unwrap();
1415 assert_eq!(
1416 load_state(&store, &paths(), "RUN01").unwrap(),
1417 Some(expected)
1418 );
1419 }
1420
1421 #[test]
1422 fn read_run_state_reduces_from_journal_without_writing() {
1423 let store = FakeStore::default();
1425 for ev in sample_events() {
1426 let line = serde_json::to_string(&ev).unwrap();
1427 store
1428 .append_line(&paths().journal_file("RUN01"), &line)
1429 .unwrap();
1430 }
1431 let files_before = store.inner.borrow().files.clone();
1433
1434 let state = read_run_state(&store, &paths(), "RUN01").unwrap().unwrap();
1435 assert_eq!(state, reduce(&sample_events()));
1436
1437 assert_eq!(
1439 store.inner.borrow().files,
1440 files_before,
1441 "read_run_state must not write anything"
1442 );
1443 assert!(
1444 store.inner.borrow().locked.is_empty(),
1445 "read_run_state must not take the lock"
1446 );
1447 assert!(
1448 !store
1449 .inner
1450 .borrow()
1451 .files
1452 .contains_key(&paths().manifest_file("RUN01")),
1453 "read_run_state must not materialize a manifest"
1454 );
1455
1456 assert!(read_run_state(&store, &paths(), "MISSING")
1458 .unwrap()
1459 .is_none());
1460 assert_eq!(
1461 read_run_state(&store, &paths(), "../escape")
1462 .unwrap_err()
1463 .kind(),
1464 io::ErrorKind::InvalidInput
1465 );
1466 }
1467
1468 #[test]
1469 fn read_run_returns_events_and_reduced_state() {
1470 let store = FakeStore::default();
1472 for ev in sample_events() {
1473 let line = serde_json::to_string(&ev).unwrap();
1474 store
1475 .append_line(&paths().journal_file("RUN01"), &line)
1476 .unwrap();
1477 }
1478 let files_before = store.inner.borrow().files.clone();
1479
1480 let (events, state) = read_run(&store, &paths(), "RUN01").unwrap().unwrap();
1481 assert_eq!(events, sample_events());
1482 assert_eq!(state, reduce(&sample_events()));
1483
1484 assert_eq!(store.inner.borrow().files, files_before);
1486 assert!(store.inner.borrow().locked.is_empty());
1487
1488 assert!(read_run(&store, &paths(), "MISSING").unwrap().is_none());
1490 assert_eq!(
1491 read_run(&store, &paths(), "../escape").unwrap_err().kind(),
1492 io::ErrorKind::InvalidInput
1493 );
1494 }
1495
1496 #[test]
1497 fn read_run_state_ignores_a_stale_manifest_fast_path() {
1498 let store = FakeStore::default();
1501 for ev in sample_events() {
1502 let line = serde_json::to_string(&ev).unwrap();
1503 store
1504 .append_line(&paths().journal_file("RUN01"), &line)
1505 .unwrap();
1506 }
1507 let mut stale = RunState::empty();
1509 stale.run_id = "RUN01".into();
1510 stale.plan_id = "STALE".into();
1511 store
1512 .write_atomic(
1513 &paths().manifest_file("RUN01"),
1514 &serde_json::to_vec(&stale).unwrap(),
1515 )
1516 .unwrap();
1517
1518 assert_eq!(
1520 load_state(&store, &paths(), "RUN01")
1521 .unwrap()
1522 .unwrap()
1523 .plan_id,
1524 "STALE"
1525 );
1526 assert_eq!(
1527 read_run_state(&store, &paths(), "RUN01")
1528 .unwrap()
1529 .unwrap()
1530 .plan_id,
1531 "plan-abc",
1532 );
1533 }
1534
1535 #[test]
1536 fn run_id_validation_rejects_path_traversal() {
1537 let store = FakeStore::default();
1538 let clock = FakeClock::at(1000);
1539 for bad in ["..", "a/b", "", ".", "x/../y"] {
1540 assert_eq!(
1541 load_state(&store, &paths(), bad).unwrap_err().kind(),
1542 io::ErrorKind::InvalidInput,
1543 "load_state must reject run id {bad:?}"
1544 );
1545 let err = Journal::open(&store, &clock, paths(), bad).err().unwrap();
1546 assert_eq!(err.kind(), io::ErrorKind::InvalidInput);
1547 }
1548 }
1549
1550 #[test]
1551 fn list_runs_excludes_lock_and_stray_files() {
1552 let store = FakeStore::default();
1553 let clock = FakeClock::at(1000);
1554 for id in ["RUN01", "RUN02"] {
1555 let idgen = FakeIdGen { id: id.into() };
1556 let _j = Journal::create(
1557 &store,
1558 &clock,
1559 &idgen,
1560 paths(),
1561 "plan".into(),
1562 "0.1.0".into(),
1563 vec!["cargo".into()],
1564 )
1565 .unwrap();
1566 }
1567 store
1570 .write_atomic(&paths().releases_dir().join("journal.jsonl.tmp"), b"x")
1571 .unwrap();
1572 store
1573 .write_atomic(&paths().releases_dir().join(".lock"), b"")
1574 .unwrap();
1575 let runs = list_runs(&store, &paths()).unwrap();
1576 assert_eq!(runs, vec!["RUN01".to_string(), "RUN02".to_string()]);
1577 }
1578
1579 #[test]
1582 fn read_events_refuses_a_too_new_schema_version() {
1583 let store = FakeStore::default();
1584 let mut ev = sample_events()[0].clone();
1585 ev.schema_version = JOURNAL_SCHEMA_VERSION + 1;
1586 let line = serde_json::to_string(&ev).unwrap();
1587 store
1588 .append_line(&paths().journal_file("RUN01"), &line)
1589 .unwrap();
1590 let err = read_events(&store, &paths().journal_file("RUN01")).unwrap_err();
1591 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1592 }
1593
1594 #[test]
1595 fn read_events_tolerates_unknown_additive_fields() {
1596 let store = FakeStore::default();
1598 let line = r#"{"schema_version":1,"seq":1,"ts":1000,"idempotency_key":"run_created","kind":"run_created","run_id":"R","plan_id":"p","version":"0.1.0","targets":[],"future_field":42}"#;
1599 store
1600 .append_line(&paths().journal_file("RUN01"), line)
1601 .unwrap();
1602 let events = read_events(&store, &paths().journal_file("RUN01")).unwrap();
1603 assert_eq!(events.len(), 1);
1604 assert_eq!(events[0].seq, 1);
1605 }
1606
1607 #[test]
1608 fn run_status_as_str_matches_serde() {
1609 for s in [
1610 RunStatus::InProgress,
1611 RunStatus::Completed,
1612 RunStatus::Abandoned,
1613 ] {
1614 assert_eq!(
1615 serde_json::to_value(s).unwrap(),
1616 serde_json::Value::String(s.as_str().to_string()),
1617 "as_str() drifted from serde for {s:?}"
1618 );
1619 }
1620 }
1621
1622 #[test]
1623 fn read_events_skips_blank_lines() {
1624 let store = FakeStore::default();
1625 store
1626 .append_line(&paths().journal_file("RUN01"), "")
1627 .unwrap();
1628 assert!(read_events(&store, &paths().journal_file("RUN01"))
1629 .unwrap()
1630 .is_empty());
1631 }
1632}