use std::io;
use std::path::{Path, PathBuf};
use crate::ports::{Clock, GitRepo, IdGen, JournalLock, JournalStore};
use crate::protocol::journal::{
EventKind, JournalEvent, Phase, PhaseOutcome, PhaseRecord, RunState, RunStatus,
JOURNAL_SCHEMA_VERSION,
};
#[derive(Debug, Clone)]
pub struct JournalPaths {
releases_dir: PathBuf,
}
impl JournalPaths {
pub fn new(releases_dir: impl Into<PathBuf>) -> Self {
Self {
releases_dir: releases_dir.into(),
}
}
pub fn from_git(git: &dyn GitRepo, override_dir: Option<&Path>) -> io::Result<Self> {
let releases_dir = match override_dir {
Some(dir) => dir.to_path_buf(),
None => git.git_common_dir()?.join("ossctl").join("releases"),
};
Ok(Self { releases_dir })
}
#[must_use]
pub fn releases_dir(&self) -> &Path {
&self.releases_dir
}
#[must_use]
pub fn lock_file(&self) -> PathBuf {
self.releases_dir.join(".lock")
}
#[must_use]
pub fn run_dir(&self, run_id: &str) -> PathBuf {
self.releases_dir.join(run_id)
}
#[must_use]
pub fn journal_file(&self, run_id: &str) -> PathBuf {
self.run_dir(run_id).join("journal.jsonl")
}
#[must_use]
pub fn manifest_file(&self, run_id: &str) -> PathBuf {
self.run_dir(run_id).join("manifest.json")
}
}
#[must_use]
pub fn reduce(events: &[JournalEvent]) -> RunState {
let mut ordered: Vec<&JournalEvent> = events.iter().collect();
ordered.sort_by_key(|e| e.seq);
let mut state = RunState::empty();
for ev in ordered {
apply(&mut state, ev);
}
state
}
pub fn apply(state: &mut RunState, event: &JournalEvent) {
if event.seq <= state.applied_seq {
return;
}
if matches!(state.status, RunStatus::Completed | RunStatus::Abandoned) {
state.applied_seq = event.seq;
return;
}
match &event.kind {
EventKind::RunCreated {
run_id,
plan_id,
version,
targets,
} => {
state.run_id.clone_from(run_id);
state.plan_id.clone_from(plan_id);
state.version.clone_from(version);
state.targets.clone_from(targets);
state.created_ts = event.ts;
state.status = RunStatus::InProgress;
}
EventKind::PhaseEntered { phase } => {
state.current_phase = Some(*phase);
}
EventKind::PhaseCompleted { phase, outcome } => {
upsert_phase(&mut state.phases, *phase, *outcome);
if state.current_phase == Some(*phase) {
state.current_phase = None;
}
if *phase == Phase::Tag && *outcome == PhaseOutcome::Ok {
state.status = RunStatus::Completed;
}
}
EventKind::TargetDryRun { target } => {
state.dry_run.insert(target.clone());
}
EventKind::TargetBuilt { target } => {
state.built.insert(target.clone());
}
EventKind::TargetPublished { target, receipt } => {
state.published.insert(target.clone(), receipt.clone());
}
EventKind::TargetCancelled { target, reason } => {
state.cancelled.insert(target.clone(), reason.clone());
}
EventKind::TagCreatedLocal { tag } => {
state.tags.entry(tag.clone()).or_default().created_local = true;
}
EventKind::TagPushedRemote { tag } => {
state.tags.entry(tag.clone()).or_default().pushed_remote = true;
}
EventKind::GithubReleaseCreated { tag, url } => {
let t = state.tags.entry(tag.clone()).or_default();
t.github_release = true;
t.github_release_url.clone_from(url);
}
EventKind::RunAbandoned { reason } => {
state.status = RunStatus::Abandoned;
state.abandon_reason = Some(reason.clone());
}
}
state.applied_seq = event.seq;
state.updated_ts = event.ts;
}
fn upsert_phase(phases: &mut Vec<PhaseRecord>, phase: Phase, outcome: PhaseOutcome) {
if let Some(rec) = phases.iter_mut().find(|r| r.phase == phase) {
rec.outcome = outcome;
} else {
phases.push(PhaseRecord { phase, outcome });
phases.sort_by_key(|r| r.phase);
}
}
pub fn read_events(store: &dyn JournalStore, path: &Path) -> io::Result<Vec<JournalEvent>> {
#[derive(serde::Deserialize)]
struct Envelope {
schema_version: u32,
}
let lines = store.read_lines(path)?;
let mut events = Vec::with_capacity(lines.len());
for (idx, line) in lines.iter().enumerate() {
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
if let Ok(envelope) = serde_json::from_str::<Envelope>(trimmed) {
if envelope.schema_version > JOURNAL_SCHEMA_VERSION {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
format!(
"release journal {}: line {} has schema_version {} but this \
ossctl understands at most {}; upgrade ossctl to resume this run",
path.display(),
idx + 1,
envelope.schema_version,
JOURNAL_SCHEMA_VERSION
),
));
}
}
let event: JournalEvent = serde_json::from_str(trimmed).map_err(|e| {
io::Error::new(
io::ErrorKind::InvalidData,
format!(
"release journal {}: line {} is not a recognized event \
(corrupt, or written by a newer ossctl): {e}",
path.display(),
idx + 1
),
)
})?;
events.push(event);
}
events.sort_by_key(|e| e.seq);
Ok(events)
}
fn validate_run_id(run_id: &str) -> io::Result<()> {
let bad = run_id.is_empty()
|| run_id == "."
|| run_id == ".."
|| run_id.contains('/')
|| run_id.contains('\\')
|| run_id.contains('\0');
if bad {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
format!("invalid run id {run_id:?}: must be a single path segment"),
));
}
Ok(())
}
pub fn load_state(
store: &dyn JournalStore,
paths: &JournalPaths,
run_id: &str,
) -> io::Result<Option<RunState>> {
validate_run_id(run_id)?;
if let Some(bytes) = store.read(&paths.manifest_file(run_id))? {
if let Ok(state) = serde_json::from_slice::<RunState>(&bytes) {
if state.run_id == run_id && state.schema_version <= JOURNAL_SCHEMA_VERSION {
return Ok(Some(state));
}
}
}
let events = read_events(store, &paths.journal_file(run_id))?;
if events.is_empty() {
return Ok(None);
}
Ok(Some(reduce(&events)))
}
pub fn read_run_state(
store: &dyn JournalStore,
paths: &JournalPaths,
run_id: &str,
) -> io::Result<Option<RunState>> {
validate_run_id(run_id)?;
let events = read_events(store, &paths.journal_file(run_id))?;
if events.is_empty() {
return Ok(None);
}
Ok(Some(reduce(&events)))
}
pub fn read_run(
store: &dyn JournalStore,
paths: &JournalPaths,
run_id: &str,
) -> io::Result<Option<(Vec<JournalEvent>, RunState)>> {
validate_run_id(run_id)?;
let events = read_events(store, &paths.journal_file(run_id))?;
if events.is_empty() {
return Ok(None);
}
let state = reduce(&events);
Ok(Some((events, state)))
}
pub fn list_runs(store: &dyn JournalStore, paths: &JournalPaths) -> io::Result<Vec<String>> {
let mut runs: Vec<String> = store
.list_dir(paths.releases_dir())?
.into_iter()
.filter(|name| validate_run_id(name).is_ok())
.filter(|name| {
store
.read_lines(&paths.journal_file(name))
.is_ok_and(|lines| lines.iter().any(|l| !l.trim().is_empty()))
})
.collect();
runs.sort();
Ok(runs)
}
pub struct Journal<'a> {
store: &'a dyn JournalStore,
clock: &'a dyn Clock,
paths: JournalPaths,
run_id: String,
state: RunState,
_lock: Box<dyn JournalLock>,
}
impl<'a> Journal<'a> {
pub fn create(
store: &'a dyn JournalStore,
clock: &'a dyn Clock,
idgen: &dyn IdGen,
paths: JournalPaths,
plan_id: String,
version: String,
targets: Vec<String>,
) -> io::Result<Self> {
let lock = store.lock_exclusive(&paths.lock_file())?;
let run_id = idgen.new_id();
let mut journal = Self {
store,
clock,
paths,
run_id: run_id.clone(),
state: RunState::empty(),
_lock: lock,
};
journal.append(EventKind::RunCreated {
run_id,
plan_id,
version,
targets,
})?;
Ok(journal)
}
pub fn open(
store: &'a dyn JournalStore,
clock: &'a dyn Clock,
paths: JournalPaths,
run_id: &str,
) -> io::Result<Self> {
validate_run_id(run_id)?;
let lock = store.lock_exclusive(&paths.lock_file())?;
let events = read_events(store, &paths.journal_file(run_id))?;
if events.is_empty() {
return Err(io::Error::new(
io::ErrorKind::NotFound,
format!("no release journal for run {run_id}"),
));
}
let state = reduce(&events);
let journal = Self {
store,
clock,
paths,
run_id: run_id.to_string(),
state,
_lock: lock,
};
let _ = journal.persist_manifest();
Ok(journal)
}
#[must_use]
pub fn run_id(&self) -> &str {
&self.run_id
}
#[must_use]
pub fn state(&self) -> &RunState {
&self.state
}
#[must_use]
pub fn paths(&self) -> &JournalPaths {
&self.paths
}
pub fn append(&mut self, kind: EventKind) -> io::Result<&RunState> {
let event = JournalEvent {
schema_version: JOURNAL_SCHEMA_VERSION,
seq: self.state.applied_seq + 1,
ts: self.clock.now_unix(),
idempotency_key: kind.idempotency_key(),
kind,
};
let line = serde_json::to_string(&event).map_err(|e| {
io::Error::new(io::ErrorKind::InvalidData, format!("serialize event: {e}"))
})?;
self.store
.append_line(&self.paths.journal_file(&self.run_id), &line)?;
apply(&mut self.state, &event);
let _ = self.persist_manifest();
Ok(&self.state)
}
fn persist_manifest(&self) -> io::Result<()> {
let bytes = serde_json::to_vec_pretty(&self.state).map_err(|e| {
io::Error::new(
io::ErrorKind::InvalidData,
format!("serialize manifest: {e}"),
)
})?;
self.store
.write_atomic(&self.paths.manifest_file(&self.run_id), &bytes)
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::protocol::journal::{PublishReceipt, TagState};
use std::cell::RefCell;
use std::collections::{HashMap, HashSet};
use std::rc::Rc;
#[derive(Default)]
struct StoreInner {
files: HashMap<PathBuf, Vec<u8>>,
locked: HashSet<PathBuf>,
fail_next_atomic: bool,
}
#[derive(Clone, Default)]
struct FakeStore {
inner: Rc<RefCell<StoreInner>>,
}
impl FakeStore {
fn journal_lines(&self, path: &Path) -> Vec<String> {
self.inner
.borrow()
.files
.get(path)
.map(|b| {
String::from_utf8_lossy(b)
.lines()
.map(str::to_string)
.collect()
})
.unwrap_or_default()
}
fn arm_atomic_failure(&self) {
self.inner.borrow_mut().fail_next_atomic = true;
}
}
struct FakeLock {
inner: Rc<RefCell<StoreInner>>,
path: PathBuf,
}
impl JournalLock for FakeLock {}
impl Drop for FakeLock {
fn drop(&mut self) {
self.inner.borrow_mut().locked.remove(&self.path);
}
}
impl JournalStore for FakeStore {
fn lock_exclusive(&self, lock_path: &Path) -> io::Result<Box<dyn JournalLock>> {
let mut inner = self.inner.borrow_mut();
if inner.locked.contains(lock_path) {
return Err(io::Error::new(
io::ErrorKind::WouldBlock,
"another release cut holds the lock",
));
}
inner.locked.insert(lock_path.to_path_buf());
Ok(Box::new(FakeLock {
inner: Rc::clone(&self.inner),
path: lock_path.to_path_buf(),
}))
}
fn append_line(&self, path: &Path, line: &str) -> io::Result<()> {
let mut inner = self.inner.borrow_mut();
let buf = inner.files.entry(path.to_path_buf()).or_default();
buf.extend_from_slice(line.as_bytes());
buf.push(b'\n');
Ok(())
}
fn read_lines(&self, path: &Path) -> io::Result<Vec<String>> {
Ok(self.journal_lines(path))
}
fn read(&self, path: &Path) -> io::Result<Option<Vec<u8>>> {
Ok(self.inner.borrow().files.get(path).cloned())
}
fn write_atomic(&self, path: &Path, bytes: &[u8]) -> io::Result<()> {
let mut inner = self.inner.borrow_mut();
if inner.fail_next_atomic {
inner.fail_next_atomic = false;
return Err(io::Error::other("injected atomic-write crash"));
}
inner.files.insert(path.to_path_buf(), bytes.to_vec());
Ok(())
}
fn list_dir(&self, dir: &Path) -> io::Result<Vec<String>> {
let inner = self.inner.borrow();
let mut names: HashSet<String> = HashSet::new();
for path in inner.files.keys() {
if let Ok(rest) = path.strip_prefix(dir) {
if let Some(first) = rest.components().next() {
names.insert(first.as_os_str().to_string_lossy().into_owned());
}
}
}
Ok(names.into_iter().collect())
}
}
struct FakeClock {
t: std::cell::Cell<u64>,
}
impl FakeClock {
fn at(t: u64) -> Self {
Self {
t: std::cell::Cell::new(t),
}
}
}
impl Clock for FakeClock {
fn now_unix(&self) -> u64 {
let now = self.t.get();
self.t.set(now + 1); now
}
}
struct FakeIdGen {
id: String,
}
impl IdGen for FakeIdGen {
fn new_id(&self) -> String {
self.id.clone()
}
}
struct FakeGit {
common_dir: PathBuf,
}
impl GitRepo for FakeGit {
fn head_commit(&self) -> io::Result<String> {
Ok("deadbeef".into())
}
fn is_work_tree(&self) -> bool {
true
}
fn shortlog(&self, _since: Option<&str>) -> io::Result<String> {
Ok(String::new())
}
fn tags(&self) -> io::Result<Vec<String>> {
Ok(Vec::new())
}
fn git_common_dir(&self) -> io::Result<PathBuf> {
Ok(self.common_dir.clone())
}
}
fn paths() -> JournalPaths {
JournalPaths::new("/repo/.git/ossctl/releases")
}
fn receipt(version: &str) -> PublishReceipt {
PublishReceipt {
ecosystem: "cargo".into(),
package: Some("ossctl".into()),
version: version.into(),
registry_url: Some("https://crates.io/crates/ossctl".into()),
digest: Some("sha256:abc".into()),
}
}
fn sample_events() -> Vec<JournalEvent> {
let kinds = vec![
EventKind::RunCreated {
run_id: "RUN01".into(),
plan_id: "plan-abc".into(),
version: "0.1.0".into(),
targets: vec!["cargo".into(), "npm".into()],
},
EventKind::PhaseEntered {
phase: Phase::DryRun,
},
EventKind::TargetDryRun {
target: "cargo".into(),
},
EventKind::PhaseCompleted {
phase: Phase::DryRun,
outcome: PhaseOutcome::Ok,
},
EventKind::PhaseEntered {
phase: Phase::Publish,
},
EventKind::TargetPublished {
target: "cargo".into(),
receipt: receipt("0.1.0"),
},
];
kinds
.into_iter()
.enumerate()
.map(|(i, kind)| JournalEvent {
schema_version: JOURNAL_SCHEMA_VERSION,
seq: (i + 1) as u64,
ts: 1000 + i as u64,
idempotency_key: kind.idempotency_key(),
kind,
})
.collect()
}
#[test]
fn paths_resolve_under_git_common_dir() {
let git = FakeGit {
common_dir: PathBuf::from("/repo/.git"),
};
let p = JournalPaths::from_git(&git, None).unwrap();
assert_eq!(p.releases_dir(), Path::new("/repo/.git/ossctl/releases"));
assert_eq!(
p.journal_file("RUN01"),
Path::new("/repo/.git/ossctl/releases/RUN01/journal.jsonl")
);
assert_eq!(
p.manifest_file("RUN01"),
Path::new("/repo/.git/ossctl/releases/RUN01/manifest.json")
);
assert_eq!(p.lock_file(), Path::new("/repo/.git/ossctl/releases/.lock"));
}
#[test]
fn path_override_wins_over_git() {
let git = FakeGit {
common_dir: PathBuf::from("/repo/.git"),
};
let p = JournalPaths::from_git(&git, Some(Path::new("/ci/journal"))).unwrap();
assert_eq!(p.releases_dir(), Path::new("/ci/journal"));
}
#[test]
fn reduce_is_deterministic() {
let events = sample_events();
let a = reduce(&events);
let b = reduce(&events);
assert_eq!(a, b);
assert_eq!(a.run_id, "RUN01");
assert_eq!(a.plan_id, "plan-abc");
assert_eq!(a.targets, vec!["cargo".to_string(), "npm".to_string()]);
assert!(a.dry_run.contains("cargo"));
assert_eq!(a.published.get("cargo").unwrap().version, "0.1.0");
assert_eq!(a.current_phase, Some(Phase::Publish));
assert_eq!(a.applied_seq, 6);
}
#[test]
fn reduce_ignores_slice_order() {
let mut events = sample_events();
events.reverse();
let out = reduce(&events);
assert_eq!(out, reduce(&sample_events()));
}
#[test]
fn replaying_a_seen_event_is_a_no_op() {
let events = sample_events();
let mut state = reduce(&events);
let before = state.clone();
apply(&mut state, &events[2]);
assert_eq!(state, before);
for ev in &events {
apply(&mut state, ev);
}
assert_eq!(state, before);
}
#[test]
fn structural_idempotency_of_publish_and_tags() {
let mut state = RunState::empty();
let mk = |seq: u64, kind: EventKind| JournalEvent {
schema_version: JOURNAL_SCHEMA_VERSION,
seq,
ts: seq,
idempotency_key: kind.idempotency_key(),
kind,
};
apply(
&mut state,
&mk(
1,
EventKind::RunCreated {
run_id: "R".into(),
plan_id: "p".into(),
version: "0.1.0".into(),
targets: vec!["cargo".into()],
},
),
);
apply(
&mut state,
&mk(
2,
EventKind::TargetPublished {
target: "cargo".into(),
receipt: receipt("0.1.0"),
},
),
);
apply(
&mut state,
&mk(
3,
EventKind::TargetPublished {
target: "cargo".into(),
receipt: receipt("0.1.1"),
},
),
);
assert_eq!(state.published.len(), 1);
assert_eq!(state.published.get("cargo").unwrap().version, "0.1.1");
apply(
&mut state,
&mk(
4,
EventKind::TagCreatedLocal {
tag: "v0.1.1".into(),
},
),
);
apply(
&mut state,
&mk(
5,
EventKind::TagPushedRemote {
tag: "v0.1.1".into(),
},
),
);
assert_eq!(
state.tags.get("v0.1.1"),
Some(&TagState {
created_local: true,
pushed_remote: true,
github_release: false,
github_release_url: None,
})
);
}
#[test]
fn tag_phase_ok_completes_the_run() {
let mut state = reduce(&sample_events());
assert_eq!(state.status, RunStatus::InProgress);
let seq = state.applied_seq + 1;
apply(
&mut state,
&JournalEvent {
schema_version: JOURNAL_SCHEMA_VERSION,
seq,
ts: 9000,
idempotency_key: "phase_completed:tag".into(),
kind: EventKind::PhaseCompleted {
phase: Phase::Tag,
outcome: PhaseOutcome::Ok,
},
},
);
assert_eq!(state.status, RunStatus::Completed);
}
#[test]
fn run_abandoned_is_terminal_with_reason() {
let mut state = reduce(&sample_events());
let seq = state.applied_seq + 1;
apply(
&mut state,
&JournalEvent {
schema_version: JOURNAL_SCHEMA_VERSION,
seq,
ts: 9000,
idempotency_key: "run_abandoned".into(),
kind: EventKind::RunAbandoned {
reason: "OTP timeout".into(),
},
},
);
assert_eq!(state.status, RunStatus::Abandoned);
assert_eq!(state.abandon_reason.as_deref(), Some("OTP timeout"));
}
#[test]
fn create_writes_run_created_and_manifest() {
let store = FakeStore::default();
let clock = FakeClock::at(1000);
let idgen = FakeIdGen { id: "RUN01".into() };
let journal = Journal::create(
&store,
&clock,
&idgen,
paths(),
"plan-abc".into(),
"0.1.0".into(),
vec!["cargo".into()],
)
.unwrap();
assert_eq!(journal.run_id(), "RUN01");
assert_eq!(journal.state().run_id, "RUN01");
assert_eq!(journal.state().applied_seq, 1);
let lines = store.journal_lines(&paths().journal_file("RUN01"));
assert_eq!(lines.len(), 1);
let manifest = store
.inner
.borrow()
.files
.get(&paths().manifest_file("RUN01"))
.cloned()
.unwrap();
let loaded: RunState = serde_json::from_slice(&manifest).unwrap();
assert_eq!(&loaded, journal.state());
}
#[test]
fn append_records_facts_and_a_failed_phase_can_later_complete_ok() {
let store = FakeStore::default();
let clock = FakeClock::at(1000);
let idgen = FakeIdGen { id: "RUN01".into() };
let mut journal = Journal::create(
&store,
&clock,
&idgen,
paths(),
"plan-abc".into(),
"0.1.0".into(),
vec!["cargo".into()],
)
.unwrap();
journal
.append(EventKind::PhaseCompleted {
phase: Phase::Publish,
outcome: PhaseOutcome::Failed,
})
.unwrap();
journal
.append(EventKind::PhaseCompleted {
phase: Phase::Publish,
outcome: PhaseOutcome::Ok,
})
.unwrap();
let lines = store.journal_lines(&paths().journal_file("RUN01"));
assert_eq!(lines.len(), 3);
let publish = journal
.state()
.phases
.iter()
.filter(|r| r.phase == Phase::Publish)
.collect::<Vec<_>>();
assert_eq!(publish.len(), 1);
assert_eq!(publish[0].outcome, PhaseOutcome::Ok);
}
#[test]
fn terminal_state_freezes_further_events() {
let mut state = reduce(&sample_events());
let published_before = state.published.clone();
let mut seq = state.applied_seq;
let mut next = |kind: EventKind| {
seq += 1;
JournalEvent {
schema_version: JOURNAL_SCHEMA_VERSION,
seq,
ts: 9000 + seq,
idempotency_key: kind.idempotency_key(),
kind,
}
};
apply(
&mut state,
&next(EventKind::RunAbandoned {
reason: "aborted".into(),
}),
);
assert_eq!(state.status, RunStatus::Abandoned);
apply(
&mut state,
&next(EventKind::TargetPublished {
target: "npm".into(),
receipt: receipt("9.9.9"),
}),
);
assert_eq!(state.status, RunStatus::Abandoned);
assert_eq!(state.published, published_before);
assert!(!state.published.contains_key("npm"));
}
#[test]
fn event_survives_a_manifest_write_crash() {
let store = FakeStore::default();
let clock = FakeClock::at(1000);
let idgen = FakeIdGen { id: "RUN01".into() };
let mut journal = Journal::create(
&store,
&clock,
&idgen,
paths(),
"plan-abc".into(),
"0.1.0".into(),
vec!["cargo".into()],
)
.unwrap();
store.arm_atomic_failure();
journal
.append(EventKind::TargetPublished {
target: "cargo".into(),
receipt: receipt("0.1.0"),
})
.unwrap();
drop(journal);
let clock2 = FakeClock::at(2000);
let reopened = Journal::open(&store, &clock2, paths(), "RUN01").unwrap();
assert!(reopened.state().published.contains_key("cargo"));
assert_eq!(reopened.state().applied_seq, 2);
let manifest = store
.inner
.borrow()
.files
.get(&paths().manifest_file("RUN01"))
.cloned()
.unwrap();
let loaded: RunState = serde_json::from_slice(&manifest).unwrap();
assert_eq!(&loaded, reopened.state());
}
#[test]
fn open_reduces_from_journal_when_manifest_absent() {
let store = FakeStore::default();
for ev in sample_events() {
let line = serde_json::to_string(&ev).unwrap();
store
.append_line(&paths().journal_file("RUN01"), &line)
.unwrap();
}
let clock = FakeClock::at(1000);
let journal = Journal::open(&store, &clock, paths(), "RUN01").unwrap();
assert_eq!(journal.state(), &reduce(&sample_events()));
}
#[test]
fn second_create_fails_while_lock_is_held() {
let store = FakeStore::default();
let clock = FakeClock::at(1000);
let idgen = FakeIdGen { id: "RUN01".into() };
let held = Journal::create(
&store,
&clock,
&idgen,
paths(),
"plan-abc".into(),
"0.1.0".into(),
vec!["cargo".into()],
)
.unwrap();
let clock2 = FakeClock::at(2000);
let idgen2 = FakeIdGen { id: "RUN02".into() };
let result = Journal::create(
&store,
&clock2,
&idgen2,
paths(),
"plan-def".into(),
"0.1.0".into(),
vec!["cargo".into()],
);
let err = result.err().expect("concurrent create must fail");
assert_eq!(err.kind(), io::ErrorKind::WouldBlock);
drop(held);
let clock3 = FakeClock::at(3000);
let idgen3 = FakeIdGen { id: "RUN02".into() };
assert!(Journal::create(
&store,
&clock3,
&idgen3,
paths(),
"plan-def".into(),
"0.1.0".into(),
vec!["cargo".into()],
)
.is_ok());
}
#[test]
fn load_state_is_read_only_and_takes_no_lock() {
let store = FakeStore::default();
let clock = FakeClock::at(1000);
let idgen = FakeIdGen { id: "RUN01".into() };
let journal = Journal::create(
&store,
&clock,
&idgen,
paths(),
"plan-abc".into(),
"0.1.0".into(),
vec!["cargo".into()],
)
.unwrap();
let loaded = load_state(&store, &paths(), "RUN01").unwrap().unwrap();
assert_eq!(&loaded, journal.state());
assert!(load_state(&store, &paths(), "MISSING").unwrap().is_none());
}
#[test]
fn load_state_prefers_manifest_but_falls_back_to_journal() {
let store = FakeStore::default();
let clock = FakeClock::at(1000);
let idgen = FakeIdGen { id: "RUN01".into() };
let journal = Journal::create(
&store,
&clock,
&idgen,
paths(),
"plan-abc".into(),
"0.1.0".into(),
vec!["cargo".into()],
)
.unwrap();
let expected = journal.state().clone();
drop(journal);
assert_eq!(
load_state(&store, &paths(), "RUN01").unwrap(),
Some(expected.clone())
);
store
.write_atomic(&paths().manifest_file("RUN01"), b"{not json")
.unwrap();
assert_eq!(
load_state(&store, &paths(), "RUN01").unwrap(),
Some(expected)
);
}
#[test]
fn read_run_state_reduces_from_journal_without_writing() {
let store = FakeStore::default();
for ev in sample_events() {
let line = serde_json::to_string(&ev).unwrap();
store
.append_line(&paths().journal_file("RUN01"), &line)
.unwrap();
}
let files_before = store.inner.borrow().files.clone();
let state = read_run_state(&store, &paths(), "RUN01").unwrap().unwrap();
assert_eq!(state, reduce(&sample_events()));
assert_eq!(
store.inner.borrow().files,
files_before,
"read_run_state must not write anything"
);
assert!(
store.inner.borrow().locked.is_empty(),
"read_run_state must not take the lock"
);
assert!(
!store
.inner
.borrow()
.files
.contains_key(&paths().manifest_file("RUN01")),
"read_run_state must not materialize a manifest"
);
assert!(read_run_state(&store, &paths(), "MISSING")
.unwrap()
.is_none());
assert_eq!(
read_run_state(&store, &paths(), "../escape")
.unwrap_err()
.kind(),
io::ErrorKind::InvalidInput
);
}
#[test]
fn read_run_returns_events_and_reduced_state() {
let store = FakeStore::default();
for ev in sample_events() {
let line = serde_json::to_string(&ev).unwrap();
store
.append_line(&paths().journal_file("RUN01"), &line)
.unwrap();
}
let files_before = store.inner.borrow().files.clone();
let (events, state) = read_run(&store, &paths(), "RUN01").unwrap().unwrap();
assert_eq!(events, sample_events());
assert_eq!(state, reduce(&sample_events()));
assert_eq!(store.inner.borrow().files, files_before);
assert!(store.inner.borrow().locked.is_empty());
assert!(read_run(&store, &paths(), "MISSING").unwrap().is_none());
assert_eq!(
read_run(&store, &paths(), "../escape").unwrap_err().kind(),
io::ErrorKind::InvalidInput
);
}
#[test]
fn read_run_state_ignores_a_stale_manifest_fast_path() {
let store = FakeStore::default();
for ev in sample_events() {
let line = serde_json::to_string(&ev).unwrap();
store
.append_line(&paths().journal_file("RUN01"), &line)
.unwrap();
}
let mut stale = RunState::empty();
stale.run_id = "RUN01".into();
stale.plan_id = "STALE".into();
store
.write_atomic(
&paths().manifest_file("RUN01"),
&serde_json::to_vec(&stale).unwrap(),
)
.unwrap();
assert_eq!(
load_state(&store, &paths(), "RUN01")
.unwrap()
.unwrap()
.plan_id,
"STALE"
);
assert_eq!(
read_run_state(&store, &paths(), "RUN01")
.unwrap()
.unwrap()
.plan_id,
"plan-abc",
);
}
#[test]
fn run_id_validation_rejects_path_traversal() {
let store = FakeStore::default();
let clock = FakeClock::at(1000);
for bad in ["..", "a/b", "", ".", "x/../y"] {
assert_eq!(
load_state(&store, &paths(), bad).unwrap_err().kind(),
io::ErrorKind::InvalidInput,
"load_state must reject run id {bad:?}"
);
let err = Journal::open(&store, &clock, paths(), bad).err().unwrap();
assert_eq!(err.kind(), io::ErrorKind::InvalidInput);
}
}
#[test]
fn list_runs_excludes_lock_and_stray_files() {
let store = FakeStore::default();
let clock = FakeClock::at(1000);
for id in ["RUN01", "RUN02"] {
let idgen = FakeIdGen { id: id.into() };
let _j = Journal::create(
&store,
&clock,
&idgen,
paths(),
"plan".into(),
"0.1.0".into(),
vec!["cargo".into()],
)
.unwrap();
}
store
.write_atomic(&paths().releases_dir().join("journal.jsonl.tmp"), b"x")
.unwrap();
store
.write_atomic(&paths().releases_dir().join(".lock"), b"")
.unwrap();
let runs = list_runs(&store, &paths()).unwrap();
assert_eq!(runs, vec!["RUN01".to_string(), "RUN02".to_string()]);
}
#[test]
fn read_events_refuses_a_too_new_schema_version() {
let store = FakeStore::default();
let mut ev = sample_events()[0].clone();
ev.schema_version = JOURNAL_SCHEMA_VERSION + 1;
let line = serde_json::to_string(&ev).unwrap();
store
.append_line(&paths().journal_file("RUN01"), &line)
.unwrap();
let err = read_events(&store, &paths().journal_file("RUN01")).unwrap_err();
assert_eq!(err.kind(), io::ErrorKind::InvalidData);
}
#[test]
fn read_events_tolerates_unknown_additive_fields() {
let store = FakeStore::default();
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}"#;
store
.append_line(&paths().journal_file("RUN01"), line)
.unwrap();
let events = read_events(&store, &paths().journal_file("RUN01")).unwrap();
assert_eq!(events.len(), 1);
assert_eq!(events[0].seq, 1);
}
#[test]
fn run_status_as_str_matches_serde() {
for s in [
RunStatus::InProgress,
RunStatus::Completed,
RunStatus::Abandoned,
] {
assert_eq!(
serde_json::to_value(s).unwrap(),
serde_json::Value::String(s.as_str().to_string()),
"as_str() drifted from serde for {s:?}"
);
}
}
#[test]
fn read_events_skips_blank_lines() {
let store = FakeStore::default();
store
.append_line(&paths().journal_file("RUN01"), "")
.unwrap();
assert!(read_events(&store, &paths().journal_file("RUN01"))
.unwrap()
.is_empty());
}
}