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 if *phase == Phase::Tag && *outcome == PhaseOutcome::Ok {
195 state.status = RunStatus::Completed;
196 }
197 }
198 EventKind::TargetDryRun { target } => {
199 state.dry_run.insert(target.clone());
200 }
201 EventKind::TargetBuilt { target } => {
202 state.built.insert(target.clone());
203 }
204 EventKind::TargetPublished { target, receipt } => {
205 state.published.insert(target.clone(), receipt.clone());
206 }
207 EventKind::TargetCancelled { target, reason } => {
208 state.cancelled.insert(target.clone(), reason.clone());
209 }
210 EventKind::TagCreatedLocal { tag } => {
211 state.tags.entry(tag.clone()).or_default().created_local = true;
212 }
213 EventKind::TagPushedRemote { tag } => {
214 state.tags.entry(tag.clone()).or_default().pushed_remote = true;
215 }
216 EventKind::GithubReleaseCreated { tag, url } => {
217 let t = state.tags.entry(tag.clone()).or_default();
218 t.github_release = true;
219 t.github_release_url.clone_from(url);
220 }
221 EventKind::RunAbandoned { reason } => {
222 state.status = RunStatus::Abandoned;
223 state.abandon_reason = Some(reason.clone());
224 }
225 }
226 state.applied_seq = event.seq;
227 state.updated_ts = event.ts;
228}
229
230fn upsert_phase(phases: &mut Vec<PhaseRecord>, phase: Phase, outcome: PhaseOutcome) {
233 if let Some(rec) = phases.iter_mut().find(|r| r.phase == phase) {
234 rec.outcome = outcome;
235 } else {
236 phases.push(PhaseRecord { phase, outcome });
237 phases.sort_by_key(|r| r.phase);
238 }
239}
240
241pub fn read_events(store: &dyn JournalStore, path: &Path) -> io::Result<Vec<JournalEvent>> {
254 #[derive(serde::Deserialize)]
258 struct Envelope {
259 schema_version: u32,
260 }
261
262 let lines = store.read_lines(path)?;
263 let mut events = Vec::with_capacity(lines.len());
264 for (idx, line) in lines.iter().enumerate() {
265 let trimmed = line.trim();
266 if trimmed.is_empty() {
267 continue;
268 }
269 if let Ok(envelope) = serde_json::from_str::<Envelope>(trimmed) {
273 if envelope.schema_version > JOURNAL_SCHEMA_VERSION {
274 return Err(io::Error::new(
275 io::ErrorKind::InvalidData,
276 format!(
277 "release journal {}: line {} has schema_version {} but this \
278 ossctl understands at most {}; upgrade ossctl to resume this run",
279 path.display(),
280 idx + 1,
281 envelope.schema_version,
282 JOURNAL_SCHEMA_VERSION
283 ),
284 ));
285 }
286 }
287 let event: JournalEvent = serde_json::from_str(trimmed).map_err(|e| {
288 io::Error::new(
289 io::ErrorKind::InvalidData,
290 format!(
291 "release journal {}: line {} is not a recognized event \
292 (corrupt, or written by a newer ossctl): {e}",
293 path.display(),
294 idx + 1
295 ),
296 )
297 })?;
298 events.push(event);
299 }
300 events.sort_by_key(|e| e.seq);
301 Ok(events)
302}
303
304fn validate_run_id(run_id: &str) -> io::Result<()> {
310 let bad = run_id.is_empty()
311 || run_id == "."
312 || run_id == ".."
313 || run_id.contains('/')
314 || run_id.contains('\\')
315 || run_id.contains('\0');
316 if bad {
317 return Err(io::Error::new(
318 io::ErrorKind::InvalidInput,
319 format!("invalid run id {run_id:?}: must be a single path segment"),
320 ));
321 }
322 Ok(())
323}
324
325pub fn load_state(
341 store: &dyn JournalStore,
342 paths: &JournalPaths,
343 run_id: &str,
344) -> io::Result<Option<RunState>> {
345 validate_run_id(run_id)?;
346 if let Some(bytes) = store.read(&paths.manifest_file(run_id))? {
348 if let Ok(state) = serde_json::from_slice::<RunState>(&bytes) {
349 if state.run_id == run_id && state.schema_version <= JOURNAL_SCHEMA_VERSION {
350 return Ok(Some(state));
351 }
352 }
353 }
355 let events = read_events(store, &paths.journal_file(run_id))?;
356 if events.is_empty() {
357 return Ok(None);
358 }
359 Ok(Some(reduce(&events)))
360}
361
362pub fn read_run_state(
379 store: &dyn JournalStore,
380 paths: &JournalPaths,
381 run_id: &str,
382) -> io::Result<Option<RunState>> {
383 validate_run_id(run_id)?;
384 let events = read_events(store, &paths.journal_file(run_id))?;
385 if events.is_empty() {
386 return Ok(None);
387 }
388 Ok(Some(reduce(&events)))
389}
390
391pub fn read_run(
409 store: &dyn JournalStore,
410 paths: &JournalPaths,
411 run_id: &str,
412) -> io::Result<Option<(Vec<JournalEvent>, RunState)>> {
413 validate_run_id(run_id)?;
414 let events = read_events(store, &paths.journal_file(run_id))?;
415 if events.is_empty() {
416 return Ok(None);
417 }
418 let state = reduce(&events);
419 Ok(Some((events, state)))
420}
421
422pub fn list_runs(store: &dyn JournalStore, paths: &JournalPaths) -> io::Result<Vec<String>> {
432 let mut runs: Vec<String> = store
433 .list_dir(paths.releases_dir())?
434 .into_iter()
435 .filter(|name| validate_run_id(name).is_ok())
436 .filter(|name| {
437 store
440 .read_lines(&paths.journal_file(name))
441 .is_ok_and(|lines| lines.iter().any(|l| !l.trim().is_empty()))
442 })
443 .collect();
444 runs.sort();
445 Ok(runs)
446}
447
448pub struct Journal<'a> {
457 store: &'a dyn JournalStore,
458 clock: &'a dyn Clock,
459 paths: JournalPaths,
460 run_id: String,
461 state: RunState,
462 _lock: Box<dyn JournalLock>,
464}
465
466impl<'a> Journal<'a> {
467 pub fn create(
475 store: &'a dyn JournalStore,
476 clock: &'a dyn Clock,
477 idgen: &dyn IdGen,
478 paths: JournalPaths,
479 plan_id: String,
480 version: String,
481 targets: Vec<String>,
482 ) -> io::Result<Self> {
483 let lock = store.lock_exclusive(&paths.lock_file())?;
484 let run_id = idgen.new_id();
485 let mut journal = Self {
486 store,
487 clock,
488 paths,
489 run_id: run_id.clone(),
490 state: RunState::empty(),
491 _lock: lock,
492 };
493 journal.append(EventKind::RunCreated {
494 run_id,
495 plan_id,
496 version,
497 targets,
498 })?;
499 Ok(journal)
500 }
501
502 pub fn open(
515 store: &'a dyn JournalStore,
516 clock: &'a dyn Clock,
517 paths: JournalPaths,
518 run_id: &str,
519 ) -> io::Result<Self> {
520 validate_run_id(run_id)?;
521 let lock = store.lock_exclusive(&paths.lock_file())?;
522 let events = read_events(store, &paths.journal_file(run_id))?;
523 if events.is_empty() {
524 return Err(io::Error::new(
525 io::ErrorKind::NotFound,
526 format!("no release journal for run {run_id}"),
527 ));
528 }
529 let state = reduce(&events);
530 let journal = Self {
531 store,
532 clock,
533 paths,
534 run_id: run_id.to_string(),
535 state,
536 _lock: lock,
537 };
538 let _ = journal.persist_manifest();
540 Ok(journal)
541 }
542
543 #[must_use]
545 pub fn run_id(&self) -> &str {
546 &self.run_id
547 }
548
549 #[must_use]
551 pub fn state(&self) -> &RunState {
552 &self.state
553 }
554
555 #[must_use]
557 pub fn paths(&self) -> &JournalPaths {
558 &self.paths
559 }
560
561 pub fn append(&mut self, kind: EventKind) -> io::Result<&RunState> {
583 let event = JournalEvent {
584 schema_version: JOURNAL_SCHEMA_VERSION,
585 seq: self.state.applied_seq + 1,
586 ts: self.clock.now_unix(),
587 idempotency_key: kind.idempotency_key(),
588 kind,
589 };
590 let line = serde_json::to_string(&event).map_err(|e| {
591 io::Error::new(io::ErrorKind::InvalidData, format!("serialize event: {e}"))
592 })?;
593 self.store
595 .append_line(&self.paths.journal_file(&self.run_id), &line)?;
596 apply(&mut self.state, &event);
598 let _ = self.persist_manifest();
600 Ok(&self.state)
601 }
602
603 fn persist_manifest(&self) -> io::Result<()> {
605 let bytes = serde_json::to_vec_pretty(&self.state).map_err(|e| {
606 io::Error::new(
607 io::ErrorKind::InvalidData,
608 format!("serialize manifest: {e}"),
609 )
610 })?;
611 self.store
612 .write_atomic(&self.paths.manifest_file(&self.run_id), &bytes)
613 }
614}
615
616#[cfg(test)]
617mod tests {
618 use super::*;
619 use crate::protocol::journal::{PublishReceipt, TagState};
620 use std::cell::RefCell;
621 use std::collections::{HashMap, HashSet};
622 use std::rc::Rc;
623
624 #[derive(Default)]
627 struct StoreInner {
628 files: HashMap<PathBuf, Vec<u8>>,
630 locked: HashSet<PathBuf>,
632 fail_next_atomic: bool,
634 }
635
636 #[derive(Clone, Default)]
637 struct FakeStore {
638 inner: Rc<RefCell<StoreInner>>,
639 }
640
641 impl FakeStore {
642 fn journal_lines(&self, path: &Path) -> Vec<String> {
643 self.inner
644 .borrow()
645 .files
646 .get(path)
647 .map(|b| {
648 String::from_utf8_lossy(b)
649 .lines()
650 .map(str::to_string)
651 .collect()
652 })
653 .unwrap_or_default()
654 }
655
656 fn arm_atomic_failure(&self) {
657 self.inner.borrow_mut().fail_next_atomic = true;
658 }
659 }
660
661 struct FakeLock {
663 inner: Rc<RefCell<StoreInner>>,
664 path: PathBuf,
665 }
666
667 impl JournalLock for FakeLock {}
668
669 impl Drop for FakeLock {
670 fn drop(&mut self) {
671 self.inner.borrow_mut().locked.remove(&self.path);
672 }
673 }
674
675 impl JournalStore for FakeStore {
676 fn lock_exclusive(&self, lock_path: &Path) -> io::Result<Box<dyn JournalLock>> {
677 let mut inner = self.inner.borrow_mut();
678 if inner.locked.contains(lock_path) {
679 return Err(io::Error::new(
680 io::ErrorKind::WouldBlock,
681 "another release cut holds the lock",
682 ));
683 }
684 inner.locked.insert(lock_path.to_path_buf());
685 Ok(Box::new(FakeLock {
686 inner: Rc::clone(&self.inner),
687 path: lock_path.to_path_buf(),
688 }))
689 }
690
691 fn append_line(&self, path: &Path, line: &str) -> io::Result<()> {
692 let mut inner = self.inner.borrow_mut();
693 let buf = inner.files.entry(path.to_path_buf()).or_default();
694 buf.extend_from_slice(line.as_bytes());
695 buf.push(b'\n');
696 Ok(())
697 }
698
699 fn read_lines(&self, path: &Path) -> io::Result<Vec<String>> {
700 Ok(self.journal_lines(path))
701 }
702
703 fn read(&self, path: &Path) -> io::Result<Option<Vec<u8>>> {
704 Ok(self.inner.borrow().files.get(path).cloned())
705 }
706
707 fn write_atomic(&self, path: &Path, bytes: &[u8]) -> io::Result<()> {
708 let mut inner = self.inner.borrow_mut();
709 if inner.fail_next_atomic {
710 inner.fail_next_atomic = false;
711 return Err(io::Error::other("injected atomic-write crash"));
712 }
713 inner.files.insert(path.to_path_buf(), bytes.to_vec());
714 Ok(())
715 }
716
717 fn list_dir(&self, dir: &Path) -> io::Result<Vec<String>> {
718 let inner = self.inner.borrow();
719 let mut names: HashSet<String> = HashSet::new();
720 for path in inner.files.keys() {
721 if let Ok(rest) = path.strip_prefix(dir) {
723 if let Some(first) = rest.components().next() {
724 names.insert(first.as_os_str().to_string_lossy().into_owned());
725 }
726 }
727 }
728 Ok(names.into_iter().collect())
729 }
730 }
731
732 struct FakeClock {
733 t: std::cell::Cell<u64>,
734 }
735 impl FakeClock {
736 fn at(t: u64) -> Self {
737 Self {
738 t: std::cell::Cell::new(t),
739 }
740 }
741 }
742 impl Clock for FakeClock {
743 fn now_unix(&self) -> u64 {
744 let now = self.t.get();
745 self.t.set(now + 1); now
747 }
748 }
749
750 struct FakeIdGen {
751 id: String,
752 }
753 impl IdGen for FakeIdGen {
754 fn new_id(&self) -> String {
755 self.id.clone()
756 }
757 }
758
759 struct FakeGit {
760 common_dir: PathBuf,
761 }
762 impl GitRepo for FakeGit {
763 fn head_commit(&self) -> io::Result<String> {
764 Ok("deadbeef".into())
765 }
766 fn is_work_tree(&self) -> bool {
767 true
768 }
769 fn shortlog(&self, _since: Option<&str>) -> io::Result<String> {
770 Ok(String::new())
771 }
772 fn tags(&self) -> io::Result<Vec<String>> {
773 Ok(Vec::new())
774 }
775 fn git_common_dir(&self) -> io::Result<PathBuf> {
776 Ok(self.common_dir.clone())
777 }
778 }
779
780 fn paths() -> JournalPaths {
781 JournalPaths::new("/repo/.git/ossctl/releases")
782 }
783
784 fn receipt(version: &str) -> PublishReceipt {
785 PublishReceipt {
786 ecosystem: "cargo".into(),
787 package: Some("ossctl".into()),
788 version: version.into(),
789 registry_url: Some("https://crates.io/crates/ossctl".into()),
790 digest: Some("sha256:abc".into()),
791 }
792 }
793
794 fn sample_events() -> Vec<JournalEvent> {
796 let kinds = vec![
797 EventKind::RunCreated {
798 run_id: "RUN01".into(),
799 plan_id: "plan-abc".into(),
800 version: "0.1.0".into(),
801 targets: vec!["cargo".into(), "npm".into()],
802 },
803 EventKind::PhaseEntered {
804 phase: Phase::DryRun,
805 },
806 EventKind::TargetDryRun {
807 target: "cargo".into(),
808 },
809 EventKind::PhaseCompleted {
810 phase: Phase::DryRun,
811 outcome: PhaseOutcome::Ok,
812 },
813 EventKind::PhaseEntered {
814 phase: Phase::Publish,
815 },
816 EventKind::TargetPublished {
817 target: "cargo".into(),
818 receipt: receipt("0.1.0"),
819 },
820 ];
821 kinds
822 .into_iter()
823 .enumerate()
824 .map(|(i, kind)| JournalEvent {
825 schema_version: JOURNAL_SCHEMA_VERSION,
826 seq: (i + 1) as u64,
827 ts: 1000 + i as u64,
828 idempotency_key: kind.idempotency_key(),
829 kind,
830 })
831 .collect()
832 }
833
834 #[test]
837 fn paths_resolve_under_git_common_dir() {
838 let git = FakeGit {
839 common_dir: PathBuf::from("/repo/.git"),
840 };
841 let p = JournalPaths::from_git(&git, None).unwrap();
842 assert_eq!(p.releases_dir(), Path::new("/repo/.git/ossctl/releases"));
843 assert_eq!(
844 p.journal_file("RUN01"),
845 Path::new("/repo/.git/ossctl/releases/RUN01/journal.jsonl")
846 );
847 assert_eq!(
848 p.manifest_file("RUN01"),
849 Path::new("/repo/.git/ossctl/releases/RUN01/manifest.json")
850 );
851 assert_eq!(p.lock_file(), Path::new("/repo/.git/ossctl/releases/.lock"));
852 }
853
854 #[test]
855 fn path_override_wins_over_git() {
856 let git = FakeGit {
857 common_dir: PathBuf::from("/repo/.git"),
858 };
859 let p = JournalPaths::from_git(&git, Some(Path::new("/ci/journal"))).unwrap();
860 assert_eq!(p.releases_dir(), Path::new("/ci/journal"));
861 }
862
863 #[test]
866 fn reduce_is_deterministic() {
867 let events = sample_events();
868 let a = reduce(&events);
869 let b = reduce(&events);
870 assert_eq!(a, b);
871 assert_eq!(a.run_id, "RUN01");
872 assert_eq!(a.plan_id, "plan-abc");
873 assert_eq!(a.targets, vec!["cargo".to_string(), "npm".to_string()]);
874 assert!(a.dry_run.contains("cargo"));
875 assert_eq!(a.published.get("cargo").unwrap().version, "0.1.0");
876 assert_eq!(a.current_phase, Some(Phase::Publish));
878 assert_eq!(a.applied_seq, 6);
879 }
880
881 #[test]
882 fn reduce_ignores_slice_order() {
883 let mut events = sample_events();
884 events.reverse();
885 let out = reduce(&events);
886 assert_eq!(out, reduce(&sample_events()));
888 }
889
890 #[test]
891 fn replaying_a_seen_event_is_a_no_op() {
892 let events = sample_events();
893 let mut state = reduce(&events);
894 let before = state.clone();
895 apply(&mut state, &events[2]);
897 assert_eq!(state, before);
898 for ev in &events {
900 apply(&mut state, ev);
901 }
902 assert_eq!(state, before);
903 }
904
905 #[test]
906 fn structural_idempotency_of_publish_and_tags() {
907 let mut state = RunState::empty();
910 let mk = |seq: u64, kind: EventKind| JournalEvent {
911 schema_version: JOURNAL_SCHEMA_VERSION,
912 seq,
913 ts: seq,
914 idempotency_key: kind.idempotency_key(),
915 kind,
916 };
917 apply(
918 &mut state,
919 &mk(
920 1,
921 EventKind::RunCreated {
922 run_id: "R".into(),
923 plan_id: "p".into(),
924 version: "0.1.0".into(),
925 targets: vec!["cargo".into()],
926 },
927 ),
928 );
929 apply(
930 &mut state,
931 &mk(
932 2,
933 EventKind::TargetPublished {
934 target: "cargo".into(),
935 receipt: receipt("0.1.0"),
936 },
937 ),
938 );
939 apply(
940 &mut state,
941 &mk(
942 3,
943 EventKind::TargetPublished {
944 target: "cargo".into(),
945 receipt: receipt("0.1.1"),
946 },
947 ),
948 );
949 assert_eq!(state.published.len(), 1);
950 assert_eq!(state.published.get("cargo").unwrap().version, "0.1.1");
951
952 apply(
953 &mut state,
954 &mk(
955 4,
956 EventKind::TagCreatedLocal {
957 tag: "v0.1.1".into(),
958 },
959 ),
960 );
961 apply(
962 &mut state,
963 &mk(
964 5,
965 EventKind::TagPushedRemote {
966 tag: "v0.1.1".into(),
967 },
968 ),
969 );
970 assert_eq!(
971 state.tags.get("v0.1.1"),
972 Some(&TagState {
973 created_local: true,
974 pushed_remote: true,
975 github_release: false,
976 github_release_url: None,
977 })
978 );
979 }
980
981 #[test]
982 fn tag_phase_ok_completes_the_run() {
983 let mut state = reduce(&sample_events());
984 assert_eq!(state.status, RunStatus::InProgress);
985 let seq = state.applied_seq + 1;
986 apply(
987 &mut state,
988 &JournalEvent {
989 schema_version: JOURNAL_SCHEMA_VERSION,
990 seq,
991 ts: 9000,
992 idempotency_key: "phase_completed:tag".into(),
993 kind: EventKind::PhaseCompleted {
994 phase: Phase::Tag,
995 outcome: PhaseOutcome::Ok,
996 },
997 },
998 );
999 assert_eq!(state.status, RunStatus::Completed);
1000 }
1001
1002 #[test]
1003 fn run_abandoned_is_terminal_with_reason() {
1004 let mut state = reduce(&sample_events());
1005 let seq = state.applied_seq + 1;
1006 apply(
1007 &mut state,
1008 &JournalEvent {
1009 schema_version: JOURNAL_SCHEMA_VERSION,
1010 seq,
1011 ts: 9000,
1012 idempotency_key: "run_abandoned".into(),
1013 kind: EventKind::RunAbandoned {
1014 reason: "OTP timeout".into(),
1015 },
1016 },
1017 );
1018 assert_eq!(state.status, RunStatus::Abandoned);
1019 assert_eq!(state.abandon_reason.as_deref(), Some("OTP timeout"));
1020 }
1021
1022 #[test]
1025 fn create_writes_run_created_and_manifest() {
1026 let store = FakeStore::default();
1027 let clock = FakeClock::at(1000);
1028 let idgen = FakeIdGen { id: "RUN01".into() };
1029 let journal = Journal::create(
1030 &store,
1031 &clock,
1032 &idgen,
1033 paths(),
1034 "plan-abc".into(),
1035 "0.1.0".into(),
1036 vec!["cargo".into()],
1037 )
1038 .unwrap();
1039 assert_eq!(journal.run_id(), "RUN01");
1040 assert_eq!(journal.state().run_id, "RUN01");
1041 assert_eq!(journal.state().applied_seq, 1);
1042
1043 let lines = store.journal_lines(&paths().journal_file("RUN01"));
1045 assert_eq!(lines.len(), 1);
1046 let manifest = store
1048 .inner
1049 .borrow()
1050 .files
1051 .get(&paths().manifest_file("RUN01"))
1052 .cloned()
1053 .unwrap();
1054 let loaded: RunState = serde_json::from_slice(&manifest).unwrap();
1055 assert_eq!(&loaded, journal.state());
1056 }
1057
1058 #[test]
1059 fn append_records_facts_and_a_failed_phase_can_later_complete_ok() {
1060 let store = FakeStore::default();
1064 let clock = FakeClock::at(1000);
1065 let idgen = FakeIdGen { id: "RUN01".into() };
1066 let mut journal = Journal::create(
1067 &store,
1068 &clock,
1069 &idgen,
1070 paths(),
1071 "plan-abc".into(),
1072 "0.1.0".into(),
1073 vec!["cargo".into()],
1074 )
1075 .unwrap();
1076 journal
1077 .append(EventKind::PhaseCompleted {
1078 phase: Phase::Publish,
1079 outcome: PhaseOutcome::Failed,
1080 })
1081 .unwrap();
1082 journal
1083 .append(EventKind::PhaseCompleted {
1084 phase: Phase::Publish,
1085 outcome: PhaseOutcome::Ok,
1086 })
1087 .unwrap();
1088 let lines = store.journal_lines(&paths().journal_file("RUN01"));
1090 assert_eq!(lines.len(), 3);
1091 let publish = journal
1093 .state()
1094 .phases
1095 .iter()
1096 .filter(|r| r.phase == Phase::Publish)
1097 .collect::<Vec<_>>();
1098 assert_eq!(publish.len(), 1);
1099 assert_eq!(publish[0].outcome, PhaseOutcome::Ok);
1100 }
1101
1102 #[test]
1103 fn terminal_state_freezes_further_events() {
1104 let mut state = reduce(&sample_events());
1106 let published_before = state.published.clone();
1107 let mut seq = state.applied_seq;
1108 let mut next = |kind: EventKind| {
1109 seq += 1;
1110 JournalEvent {
1111 schema_version: JOURNAL_SCHEMA_VERSION,
1112 seq,
1113 ts: 9000 + seq,
1114 idempotency_key: kind.idempotency_key(),
1115 kind,
1116 }
1117 };
1118 apply(
1119 &mut state,
1120 &next(EventKind::RunAbandoned {
1121 reason: "aborted".into(),
1122 }),
1123 );
1124 assert_eq!(state.status, RunStatus::Abandoned);
1125 apply(
1127 &mut state,
1128 &next(EventKind::TargetPublished {
1129 target: "npm".into(),
1130 receipt: receipt("9.9.9"),
1131 }),
1132 );
1133 assert_eq!(state.status, RunStatus::Abandoned);
1134 assert_eq!(state.published, published_before);
1135 assert!(!state.published.contains_key("npm"));
1136 }
1137
1138 #[test]
1141 fn event_survives_a_manifest_write_crash() {
1142 let store = FakeStore::default();
1143 let clock = FakeClock::at(1000);
1144 let idgen = FakeIdGen { id: "RUN01".into() };
1145 let mut journal = Journal::create(
1146 &store,
1147 &clock,
1148 &idgen,
1149 paths(),
1150 "plan-abc".into(),
1151 "0.1.0".into(),
1152 vec!["cargo".into()],
1153 )
1154 .unwrap();
1155 store.arm_atomic_failure();
1159 journal
1160 .append(EventKind::TargetPublished {
1161 target: "cargo".into(),
1162 receipt: receipt("0.1.0"),
1163 })
1164 .unwrap();
1165 drop(journal); let clock2 = FakeClock::at(2000);
1170 let reopened = Journal::open(&store, &clock2, paths(), "RUN01").unwrap();
1171 assert!(reopened.state().published.contains_key("cargo"));
1172 assert_eq!(reopened.state().applied_seq, 2);
1173 let manifest = store
1175 .inner
1176 .borrow()
1177 .files
1178 .get(&paths().manifest_file("RUN01"))
1179 .cloned()
1180 .unwrap();
1181 let loaded: RunState = serde_json::from_slice(&manifest).unwrap();
1182 assert_eq!(&loaded, reopened.state());
1183 }
1184
1185 #[test]
1186 fn open_reduces_from_journal_when_manifest_absent() {
1187 let store = FakeStore::default();
1189 for ev in sample_events() {
1190 let line = serde_json::to_string(&ev).unwrap();
1191 store
1192 .append_line(&paths().journal_file("RUN01"), &line)
1193 .unwrap();
1194 }
1195 let clock = FakeClock::at(1000);
1196 let journal = Journal::open(&store, &clock, paths(), "RUN01").unwrap();
1197 assert_eq!(journal.state(), &reduce(&sample_events()));
1198 }
1199
1200 #[test]
1203 fn second_create_fails_while_lock_is_held() {
1204 let store = FakeStore::default();
1205 let clock = FakeClock::at(1000);
1206 let idgen = FakeIdGen { id: "RUN01".into() };
1207 let held = Journal::create(
1208 &store,
1209 &clock,
1210 &idgen,
1211 paths(),
1212 "plan-abc".into(),
1213 "0.1.0".into(),
1214 vec!["cargo".into()],
1215 )
1216 .unwrap();
1217
1218 let clock2 = FakeClock::at(2000);
1220 let idgen2 = FakeIdGen { id: "RUN02".into() };
1221 let result = Journal::create(
1224 &store,
1225 &clock2,
1226 &idgen2,
1227 paths(),
1228 "plan-def".into(),
1229 "0.1.0".into(),
1230 vec!["cargo".into()],
1231 );
1232 let err = result.err().expect("concurrent create must fail");
1233 assert_eq!(err.kind(), io::ErrorKind::WouldBlock);
1234
1235 drop(held);
1237 let clock3 = FakeClock::at(3000);
1238 let idgen3 = FakeIdGen { id: "RUN02".into() };
1239 assert!(Journal::create(
1240 &store,
1241 &clock3,
1242 &idgen3,
1243 paths(),
1244 "plan-def".into(),
1245 "0.1.0".into(),
1246 vec!["cargo".into()],
1247 )
1248 .is_ok());
1249 }
1250
1251 #[test]
1254 fn load_state_is_read_only_and_takes_no_lock() {
1255 let store = FakeStore::default();
1256 let clock = FakeClock::at(1000);
1257 let idgen = FakeIdGen { id: "RUN01".into() };
1258 let 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 let loaded = load_state(&store, &paths(), "RUN01").unwrap().unwrap();
1270 assert_eq!(&loaded, journal.state());
1271 assert!(load_state(&store, &paths(), "MISSING").unwrap().is_none());
1272 }
1273
1274 #[test]
1275 fn load_state_prefers_manifest_but_falls_back_to_journal() {
1276 let store = FakeStore::default();
1277 let clock = FakeClock::at(1000);
1278 let idgen = FakeIdGen { id: "RUN01".into() };
1279 let journal = Journal::create(
1280 &store,
1281 &clock,
1282 &idgen,
1283 paths(),
1284 "plan-abc".into(),
1285 "0.1.0".into(),
1286 vec!["cargo".into()],
1287 )
1288 .unwrap();
1289 let expected = journal.state().clone();
1290 drop(journal);
1291
1292 assert_eq!(
1294 load_state(&store, &paths(), "RUN01").unwrap(),
1295 Some(expected.clone())
1296 );
1297
1298 store
1300 .write_atomic(&paths().manifest_file("RUN01"), b"{not json")
1301 .unwrap();
1302 assert_eq!(
1303 load_state(&store, &paths(), "RUN01").unwrap(),
1304 Some(expected)
1305 );
1306 }
1307
1308 #[test]
1309 fn read_run_state_reduces_from_journal_without_writing() {
1310 let store = FakeStore::default();
1312 for ev in sample_events() {
1313 let line = serde_json::to_string(&ev).unwrap();
1314 store
1315 .append_line(&paths().journal_file("RUN01"), &line)
1316 .unwrap();
1317 }
1318 let files_before = store.inner.borrow().files.clone();
1320
1321 let state = read_run_state(&store, &paths(), "RUN01").unwrap().unwrap();
1322 assert_eq!(state, reduce(&sample_events()));
1323
1324 assert_eq!(
1326 store.inner.borrow().files,
1327 files_before,
1328 "read_run_state must not write anything"
1329 );
1330 assert!(
1331 store.inner.borrow().locked.is_empty(),
1332 "read_run_state must not take the lock"
1333 );
1334 assert!(
1335 !store
1336 .inner
1337 .borrow()
1338 .files
1339 .contains_key(&paths().manifest_file("RUN01")),
1340 "read_run_state must not materialize a manifest"
1341 );
1342
1343 assert!(read_run_state(&store, &paths(), "MISSING")
1345 .unwrap()
1346 .is_none());
1347 assert_eq!(
1348 read_run_state(&store, &paths(), "../escape")
1349 .unwrap_err()
1350 .kind(),
1351 io::ErrorKind::InvalidInput
1352 );
1353 }
1354
1355 #[test]
1356 fn read_run_returns_events_and_reduced_state() {
1357 let store = FakeStore::default();
1359 for ev in sample_events() {
1360 let line = serde_json::to_string(&ev).unwrap();
1361 store
1362 .append_line(&paths().journal_file("RUN01"), &line)
1363 .unwrap();
1364 }
1365 let files_before = store.inner.borrow().files.clone();
1366
1367 let (events, state) = read_run(&store, &paths(), "RUN01").unwrap().unwrap();
1368 assert_eq!(events, sample_events());
1369 assert_eq!(state, reduce(&sample_events()));
1370
1371 assert_eq!(store.inner.borrow().files, files_before);
1373 assert!(store.inner.borrow().locked.is_empty());
1374
1375 assert!(read_run(&store, &paths(), "MISSING").unwrap().is_none());
1377 assert_eq!(
1378 read_run(&store, &paths(), "../escape").unwrap_err().kind(),
1379 io::ErrorKind::InvalidInput
1380 );
1381 }
1382
1383 #[test]
1384 fn read_run_state_ignores_a_stale_manifest_fast_path() {
1385 let store = FakeStore::default();
1388 for ev in sample_events() {
1389 let line = serde_json::to_string(&ev).unwrap();
1390 store
1391 .append_line(&paths().journal_file("RUN01"), &line)
1392 .unwrap();
1393 }
1394 let mut stale = RunState::empty();
1396 stale.run_id = "RUN01".into();
1397 stale.plan_id = "STALE".into();
1398 store
1399 .write_atomic(
1400 &paths().manifest_file("RUN01"),
1401 &serde_json::to_vec(&stale).unwrap(),
1402 )
1403 .unwrap();
1404
1405 assert_eq!(
1407 load_state(&store, &paths(), "RUN01")
1408 .unwrap()
1409 .unwrap()
1410 .plan_id,
1411 "STALE"
1412 );
1413 assert_eq!(
1414 read_run_state(&store, &paths(), "RUN01")
1415 .unwrap()
1416 .unwrap()
1417 .plan_id,
1418 "plan-abc",
1419 );
1420 }
1421
1422 #[test]
1423 fn run_id_validation_rejects_path_traversal() {
1424 let store = FakeStore::default();
1425 let clock = FakeClock::at(1000);
1426 for bad in ["..", "a/b", "", ".", "x/../y"] {
1427 assert_eq!(
1428 load_state(&store, &paths(), bad).unwrap_err().kind(),
1429 io::ErrorKind::InvalidInput,
1430 "load_state must reject run id {bad:?}"
1431 );
1432 let err = Journal::open(&store, &clock, paths(), bad).err().unwrap();
1433 assert_eq!(err.kind(), io::ErrorKind::InvalidInput);
1434 }
1435 }
1436
1437 #[test]
1438 fn list_runs_excludes_lock_and_stray_files() {
1439 let store = FakeStore::default();
1440 let clock = FakeClock::at(1000);
1441 for id in ["RUN01", "RUN02"] {
1442 let idgen = FakeIdGen { id: id.into() };
1443 let _j = Journal::create(
1444 &store,
1445 &clock,
1446 &idgen,
1447 paths(),
1448 "plan".into(),
1449 "0.1.0".into(),
1450 vec!["cargo".into()],
1451 )
1452 .unwrap();
1453 }
1454 store
1457 .write_atomic(&paths().releases_dir().join("journal.jsonl.tmp"), b"x")
1458 .unwrap();
1459 store
1460 .write_atomic(&paths().releases_dir().join(".lock"), b"")
1461 .unwrap();
1462 let runs = list_runs(&store, &paths()).unwrap();
1463 assert_eq!(runs, vec!["RUN01".to_string(), "RUN02".to_string()]);
1464 }
1465
1466 #[test]
1469 fn read_events_refuses_a_too_new_schema_version() {
1470 let store = FakeStore::default();
1471 let mut ev = sample_events()[0].clone();
1472 ev.schema_version = JOURNAL_SCHEMA_VERSION + 1;
1473 let line = serde_json::to_string(&ev).unwrap();
1474 store
1475 .append_line(&paths().journal_file("RUN01"), &line)
1476 .unwrap();
1477 let err = read_events(&store, &paths().journal_file("RUN01")).unwrap_err();
1478 assert_eq!(err.kind(), io::ErrorKind::InvalidData);
1479 }
1480
1481 #[test]
1482 fn read_events_tolerates_unknown_additive_fields() {
1483 let store = FakeStore::default();
1485 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}"#;
1486 store
1487 .append_line(&paths().journal_file("RUN01"), line)
1488 .unwrap();
1489 let events = read_events(&store, &paths().journal_file("RUN01")).unwrap();
1490 assert_eq!(events.len(), 1);
1491 assert_eq!(events[0].seq, 1);
1492 }
1493
1494 #[test]
1495 fn run_status_as_str_matches_serde() {
1496 for s in [
1497 RunStatus::InProgress,
1498 RunStatus::Completed,
1499 RunStatus::Abandoned,
1500 ] {
1501 assert_eq!(
1502 serde_json::to_value(s).unwrap(),
1503 serde_json::Value::String(s.as_str().to_string()),
1504 "as_str() drifted from serde for {s:?}"
1505 );
1506 }
1507 }
1508
1509 #[test]
1510 fn read_events_skips_blank_lines() {
1511 let store = FakeStore::default();
1512 store
1513 .append_line(&paths().journal_file("RUN01"), "")
1514 .unwrap();
1515 assert!(read_events(&store, &paths().journal_file("RUN01"))
1516 .unwrap()
1517 .is_empty());
1518 }
1519}