use std::cell::RefCell;
use std::fs::{self, File, OpenOptions};
use std::io::{BufRead, BufReader, Write};
#[cfg(unix)]
use std::os::fd::AsRawFd;
use std::path::{Path, PathBuf};
use std::process::Command;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::{Arc, Mutex, OnceLock};
use serde::{Deserialize, Serialize};
use time::OffsetDateTime;
use tracing::{debug, warn};
use crate::engine::worktrees::main_repo_root;
use crate::exec::{AgentCaller, Exec, AGENT_CALLER_ENV};
use crate::id::{ExecId, TraceId};
use crate::store::sqlite::SqliteStore;
const JOURNAL_ROOT: &str = ".lf/journal/runs";
const JOURNAL_EXCLUDE_ENTRY: &str = ".lf/journal/";
pub(crate) const EXEC_PROCESS_ROOT: &str = "runtime/exec-processes";
pub const LF_TRACE_ID_ENV: &str = "LF_TRACE_ID";
pub const LF_PROCESS_ID_ENV: &str = "LF_PROCESS_ID";
#[cfg(test)]
pub(crate) fn test_env_lock() -> std::sync::MutexGuard<'static, ()> {
static LOCK: std::sync::OnceLock<std::sync::Mutex<()>> = std::sync::OnceLock::new();
LOCK.get_or_init(|| std::sync::Mutex::new(()))
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
#[cfg(test)]
thread_local! {
static TEST_LEDGER_DB_PATH: RefCell<Option<PathBuf>> = const { RefCell::new(None) };
}
#[cfg(test)]
#[derive(Debug)]
pub(crate) struct TestLedgerGuard {
_lock: std::sync::MutexGuard<'static, ()>,
previous_lf_home: Option<std::ffi::OsString>,
previous_db_path: Option<std::ffi::OsString>,
previous_control_home: Option<std::ffi::OsString>,
previous_control_db_path: Option<std::ffi::OsString>,
previous_test_path: Option<PathBuf>,
home: tempfile::TempDir,
}
#[cfg(test)]
impl TestLedgerGuard {
pub(crate) fn new() -> Self {
let lock = test_env_lock();
let home = tempfile::TempDir::new().expect("test ledger home");
let previous_lf_home = std::env::var_os("LF_HOME");
let previous_db_path = std::env::var_os("LF_DB_PATH");
let previous_control_home = std::env::var_os(crate::store::CONTROL_HOME_ENV);
let previous_control_db_path = std::env::var_os(crate::store::CONTROL_DB_PATH_ENV);
std::env::remove_var("LF_HOME");
std::env::remove_var("LF_DB_PATH");
std::env::remove_var(crate::store::CONTROL_HOME_ENV);
std::env::remove_var(crate::store::CONTROL_DB_PATH_ENV);
std::env::set_var("LF_HOME", home.path());
let previous_test_path =
TEST_LEDGER_DB_PATH.with(|path| path.replace(Some(home.path().join("loopflow.db"))));
Self {
_lock: lock,
previous_lf_home,
previous_db_path,
previous_control_home,
previous_control_db_path,
previous_test_path,
home,
}
}
pub(crate) fn home(&self) -> &Path {
self.home.path()
}
pub(crate) fn set_db_path(&self, path: PathBuf) {
TEST_LEDGER_DB_PATH.with(|current| *current.borrow_mut() = Some(path));
}
}
#[cfg(test)]
impl Drop for TestLedgerGuard {
fn drop(&mut self) {
TEST_LEDGER_DB_PATH.with(|path| *path.borrow_mut() = self.previous_test_path.take());
match &self.previous_lf_home {
Some(value) => std::env::set_var("LF_HOME", value),
None => std::env::remove_var("LF_HOME"),
}
match &self.previous_db_path {
Some(value) => std::env::set_var("LF_DB_PATH", value),
None => std::env::remove_var("LF_DB_PATH"),
}
match &self.previous_control_home {
Some(value) => std::env::set_var(crate::store::CONTROL_HOME_ENV, value),
None => std::env::remove_var(crate::store::CONTROL_HOME_ENV),
}
match &self.previous_control_db_path {
Some(value) => std::env::set_var(crate::store::CONTROL_DB_PATH_ENV, value),
None => std::env::remove_var(crate::store::CONTROL_DB_PATH_ENV),
}
}
}
thread_local! {
static RUN_CONTEXT: RefCell<Option<RunContext>> = const { RefCell::new(None) };
}
static PROCESS_STARTED_AT: OnceLock<i64> = OnceLock::new();
static PROCESS_CONTEXT: Mutex<Option<RunContext>> = Mutex::new(None);
#[derive(Debug, Clone)]
struct RunContext {
run_id: TraceId,
process_id: ExecId,
parent_process_id: Option<ExecId>,
agent_caller: Option<AgentCaller>,
started_at: i64,
process_started_at: Option<i64>,
cwd: PathBuf,
ledger_path: PathBuf,
ledger: Option<SqliteStore>,
command: Option<String>,
run_dir: Option<PathBuf>,
repo: Option<String>,
wave: Option<String>,
finished: Arc<AtomicBool>,
minted_run_id: bool,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub(crate) struct ExecProcessReceipt {
pub schema_version: u32,
pub trace_id: String,
pub exec_id: String,
pub pid: u32,
pub started_at: i64,
}
impl ExecProcessReceipt {
fn process_evidence(&self) -> ProcessIdentityEvidence {
match process_started_at(self.pid) {
Ok(Some(started_at)) if (started_at - self.started_at).abs() <= 3 => {
ProcessIdentityEvidence::Live
}
Ok(Some(_)) | Ok(None) => ProcessIdentityEvidence::Dead,
Err(_) => ProcessIdentityEvidence::Unknown,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum ProcessIdentityEvidence {
Live,
Dead,
Unknown,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum LfNode {
Run,
Flow,
Skill,
}
#[derive(Debug, Clone, Copy, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "snake_case")]
pub enum LfEventType {
Started,
Completed,
Errored,
Escalated,
}
#[derive(Debug, Clone, Default, PartialEq, Eq)]
pub struct LfEventFields {
pub wave_name: Option<String>,
pub worktree: Option<String>,
pub command: Option<Vec<String>>,
pub flow: Option<String>,
pub skill: Option<String>,
pub index: Option<u32>,
pub error: Option<String>,
pub signal: Option<String>,
pub exit_code: Option<i32>,
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
pub struct LfEvent {
pub run_id: TraceId,
#[serde(with = "time::serde::rfc3339")]
pub ts: OffsetDateTime,
pub node: LfNode,
pub event: LfEventType,
#[serde(skip_serializing_if = "Option::is_none")]
pub wave_name: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub worktree: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub command: Option<Vec<String>>,
#[serde(skip_serializing_if = "Option::is_none")]
pub flow: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub skill: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub index: Option<u32>,
#[serde(skip_serializing_if = "Option::is_none")]
pub error: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub signal: Option<String>,
}
pub fn emit(repo_root: &Path, node: LfNode, event: LfEventType, fields: LfEventFields) {
if let Err(err) = try_emit(repo_root, node, event, fields) {
debug!(
error = %err,
repo = %repo_root.display(),
?node,
?event,
"journal append failed"
);
}
}
pub fn with_process(run: impl FnOnce() -> anyhow::Result<()>) -> anyhow::Result<()> {
PROCESS_STARTED_AT
.set(OffsetDateTime::now_utc().unix_timestamp())
.expect("one lf entry point per process");
let result = run();
if current_context().is_none() {
observe_process(&std::env::args().collect::<Vec<_>>());
}
crate::engine::agent::wait_for_interrupt_cleanup();
if let Some(context) = current_context() {
finish_runtime(&context.cwd, &result);
} else {
eprintln!("Exec history unavailable: no compatible process ledger for this process");
}
result
}
pub fn observe_process(command: &[String]) {
if current_context().is_some() {
return;
}
let observe = || -> anyhow::Result<()> {
let path = crate::store::observation_database_path()?;
let ledger = SqliteStore::open_existing_execs(&path)?;
let directory = std::env::current_dir()?;
let fields = LfEventFields {
command: Some(command.to_vec()),
worktree: Some(directory.display().to_string()),
..LfEventFields::default()
};
create_run_context(&directory, &fields, path, Some(ledger))?;
try_emit(&directory, LfNode::Run, LfEventType::Started, fields)?;
Ok(())
};
if let Err(error) = observe() {
debug!(%error, "early Exec observation unavailable");
}
}
pub fn command_exit_code<T>(result: &anyhow::Result<T>) -> u8 {
match result {
Ok(_) => 0,
Err(error) => error
.downcast_ref::<crate::exec::CommandExit>()
.map_or(1, |exit| exit.0),
}
}
pub(crate) fn is_cli_process() -> bool {
PROCESS_STARTED_AT.get().is_some()
}
pub fn with_runtime<T>(
repo_root: &Path,
command: &[String],
run: impl FnOnce() -> anyhow::Result<T>,
) -> anyhow::Result<T> {
if is_cli_process() || current_context().is_some() {
return run();
}
admit_process(repo_root, command);
let result = run();
finish_runtime(repo_root, &result);
result
}
pub fn admit_process(repo_root: &Path, command: &[String]) {
let attribution = crate::work::wave::context::exec_attribution(Some(repo_root));
if let Some(failure) = attribution.failure.as_deref() {
warn!(
error = failure,
"ambient wave identity failed validation; run attributed to no wave \
— pass --wave <name> to recover"
);
}
emit(
repo_root,
LfNode::Run,
LfEventType::Started,
LfEventFields {
wave_name: attribution.wave,
error: attribution.failure,
worktree: Some(repo_root.display().to_string()),
command: Some(command.to_vec()),
..LfEventFields::default()
},
);
}
fn finish_runtime<T>(directory: &Path, result: &anyhow::Result<T>) {
let code = command_exit_code(result);
emit(
directory,
LfNode::Run,
if code == 0 {
LfEventType::Completed
} else {
LfEventType::Errored
},
LfEventFields {
error: result
.as_ref()
.err()
.filter(|_| code != 0)
.map(|error| format!("{error:#}")),
exit_code: Some(i32::from(code)),
..LfEventFields::default()
},
);
}
pub fn runs_root(worktree: &Path) -> PathBuf {
worktree.join(JOURNAL_ROOT)
}
pub fn events_path(run_dir: &Path) -> PathBuf {
run_dir.join("events.jsonl")
}
pub fn read_events(run_dir: &Path) -> Result<Vec<LfEvent>, std::io::Error> {
let path = events_path(run_dir);
if !path.exists() {
return Ok(Vec::new());
}
let file = fs::File::open(path)?;
let mut events = Vec::new();
for line in BufReader::new(file).lines() {
let line = line?;
if line.trim().is_empty() {
continue;
}
let event = serde_json::from_str(&line).map_err(std::io::Error::other)?;
events.push(event);
}
Ok(events)
}
fn try_emit(
repo_root: &Path,
node: LfNode,
event: LfEventType,
fields: LfEventFields,
) -> Result<(), std::io::Error> {
let is_run_started = matches!((node, event), (LfNode::Run, LfEventType::Started));
let maybe_context = if is_run_started {
ensure_run_context(repo_root, &fields)?
} else {
current_context()
};
let Some(context) = maybe_context else {
return Ok(());
};
let terminal = matches!(node, LfNode::Run)
&& matches!(
event,
LfEventType::Completed | LfEventType::Errored | LfEventType::Escalated
);
if terminal && context.finished.swap(true, Ordering::AcqRel) {
return Ok(());
}
let exit_code = fields.exit_code;
let event = LfEvent {
run_id: context.run_id.clone(),
ts: if is_run_started {
OffsetDateTime::from_unix_timestamp(context.started_at)
.map_err(std::io::Error::other)?
} else {
OffsetDateTime::now_utc()
},
node,
event,
wave_name: fields.wave_name,
worktree: fields.worktree,
command: fields.command,
flow: fields.flow,
skill: fields.skill,
index: fields.index,
error: fields.error,
signal: fields.signal,
};
if let Some(run_dir) = &context.run_dir {
if let Err(error) = append_event(run_dir, &event) {
warn!(%error, path = %run_dir.display(), "file journal append failed; recording to ledger");
}
}
ledger_insert(&context, &event, repo_root, exit_code);
if matches!(node, LfNode::Run)
&& matches!(
event.event,
LfEventType::Completed | LfEventType::Errored | LfEventType::Escalated
)
{
remove_exec_process_receipt();
if context.minted_run_id {
std::env::remove_var(LF_TRACE_ID_ENV);
}
std::env::remove_var(LF_PROCESS_ID_ENV);
clear_context();
}
Ok(())
}
fn ledger_insert(context: &RunContext, event: &LfEvent, repo_root: &Path, exit_code: Option<i32>) {
if event.node != LfNode::Run {
return;
}
let outcome = match event.event {
LfEventType::Completed => Some("succeeded"),
LfEventType::Errored => Some("failed"),
LfEventType::Escalated => Some("interrupted"),
_ => None,
};
let record = Exec {
id: context.process_id.clone(),
trace_id: event.run_id.clone(),
parent_exec_id: context.parent_process_id.clone(),
via_agent: Some(context.agent_caller.is_some()),
caller_session_id: context
.agent_caller
.as_ref()
.map(|caller| caller.session_id.clone()),
caller_provider_generation: context
.agent_caller
.as_ref()
.map(|caller| caller.provider_generation),
command: context.command.clone(),
repo: context.repo.clone(),
cwd: Some(repo_root.display().to_string()),
started_at: context.started_at,
completed_at: outcome.map(|_| event.ts.unix_timestamp()),
outcome: outcome.map(str::to_owned),
exit_code: exit_code.or(match event.event {
LfEventType::Completed => Some(0),
LfEventType::Errored => Some(1),
LfEventType::Escalated => Some(130),
_ => None,
}),
signal: event.signal.clone(),
error: event.error.clone(),
};
match context
.ledger
.clone()
.map_or_else(|| SqliteStore::new(&context.ledger_path), Ok)
{
Ok(store) => {
if let Err(err) = store.record_exec(&record) {
if first_ledger_failure() {
warn!(error = %err, run_id = %record.trace_id, "ledger insert failed — this run is not being recorded");
} else {
debug!(error = %err, run_id = %record.trace_id, "ledger insert failed");
}
}
}
Err(err) => {
if first_ledger_failure() {
warn!(error = %err, "ledger unavailable — runs are not being recorded");
} else {
debug!(error = %err, "ledger unavailable");
}
}
}
}
fn first_ledger_failure() -> bool {
static WARNED: AtomicBool = AtomicBool::new(false);
!WARNED.swap(true, Ordering::Relaxed)
}
pub fn open_ledger() -> Result<SqliteStore, crate::store::StoreError> {
SqliteStore::new(&ledger_db_path()?)
}
#[cfg(not(test))]
fn ledger_db_path() -> Result<PathBuf, crate::store::StoreError> {
crate::store::database_path_from_env()
.map_err(|error| crate::store::StoreError::InvalidData(error.to_string()))
}
#[cfg(test)]
fn ledger_db_path() -> Result<PathBuf, crate::store::StoreError> {
if let Some(path) = TEST_LEDGER_DB_PATH.with(|path| path.borrow().clone()) {
return Ok(path);
}
static TEST_HOME: std::sync::OnceLock<tempfile::TempDir> = std::sync::OnceLock::new();
Ok(TEST_HOME
.get_or_init(|| tempfile::TempDir::new().expect("test ledger home"))
.path()
.join("loopflow.db"))
}
fn ensure_run_context(
repo_root: &Path,
fields: &LfEventFields,
) -> Result<Option<RunContext>, std::io::Error> {
if let Some(context) = current_context() {
return Ok(Some(context));
}
create_run_context(
repo_root,
fields,
ledger_db_path().map_err(std::io::Error::other)?,
None,
)
}
fn create_run_context(
repo_root: &Path,
fields: &LfEventFields,
ledger_path: PathBuf,
ledger: Option<SqliteStore>,
) -> Result<Option<RunContext>, std::io::Error> {
let early = ledger.is_some();
let same_store = ledger_db_path()
.is_ok_and(|path| crate::store::same_database_file(&path, &ledger_path).unwrap_or(false));
let main_repo = (!early).then(|| main_repo_root(repo_root).ok()).flatten();
let wave_name = if early {
None
} else {
let attribution = crate::work::wave::context::exec_attribution(main_repo.as_deref());
if let Some(failure) = attribution.failure.as_deref() {
debug!(
error = failure,
"ambient wave identity failed validation; Exec has no wave"
);
}
attribution.wave
};
let agent_caller = same_store
.then(|| std::env::var(AGENT_CALLER_ENV).ok())
.flatten()
.map(|value| serde_json::from_str::<AgentCaller>(&value))
.transpose()
.map_err(std::io::Error::other)?;
std::env::remove_var(AGENT_CALLER_ENV);
let agent_parent = agent_caller.as_ref().and_then(|caller| {
match ledger
.clone()
.map_or_else(open_ledger, Ok)
.and_then(|store| store.agent_parent(caller))
{
Ok(parent) => parent,
Err(error) => {
warn!(%error, "agent caller could not be resolved");
None
}
}
});
let inherited_trace = if agent_caller.is_some() {
agent_parent
.as_ref()
.and_then(|(_, trace)| TraceId::parse(trace).ok())
} else {
same_store.then(|| configured_run_id(repo_root)).flatten()
};
let (run_id, minted_run_id) = match inherited_trace {
Some(run_id) => (run_id, false),
None => {
let run_id = TraceId::default();
std::env::set_var(LF_TRACE_ID_ENV, run_id.as_str());
(run_id, true)
}
};
let parent_process_id = if agent_caller.is_some() {
agent_parent.map(|(parent, _)| parent)
} else {
(!minted_run_id)
.then(|| {
std::env::var(LF_PROCESS_ID_ENV)
.ok()
.and_then(|value| ExecId::parse(&value).ok())
})
.flatten()
.filter(parent_is_recorded)
};
std::env::set_var(LF_TRACE_ID_ENV, run_id.as_str());
let process_id = ExecId::default();
std::env::set_var(LF_PROCESS_ID_ENV, process_id.as_str());
let run_dir = if early {
None
} else {
match crate::repo::discover_repo_root(repo_root)
.map_err(std::io::Error::other)
.and_then(|root| {
root.ok_or_else(|| std::io::Error::other("no repository for file journal"))
})
.and_then(|root| {
ensure_journal_ignored(&root)?;
let dir = runs_root(&root).join(run_id.as_str());
fs::create_dir_all(&dir)?;
Ok(dir)
}) {
Ok(dir) => Some(dir),
Err(err) => {
debug!(
error = %err,
repo = %repo_root.display(),
"file journal unavailable; recording to ledger only"
);
None
}
}
};
let repo = main_repo
.as_deref()
.unwrap_or(repo_root)
.display()
.to_string();
let process_started_at = process_started_at(std::process::id()).unwrap_or_else(|error| {
debug!(%error, "Exec process evidence unavailable; recording command history only");
None
});
let context = RunContext {
run_id,
process_id,
parent_process_id,
agent_caller,
started_at: PROCESS_STARTED_AT
.get()
.copied()
.unwrap_or_else(|| OffsetDateTime::now_utc().unix_timestamp()),
process_started_at,
cwd: repo_root.to_path_buf(),
ledger_path,
ledger,
command: fields
.command
.as_ref()
.and_then(|argv| serde_json::to_string(argv).ok()),
run_dir,
repo: (!early).then_some(repo),
wave: wave_name.clone(),
finished: Arc::new(AtomicBool::new(false)),
minted_run_id,
};
set_context(context.clone());
let interrupted = context.clone();
let directory = repo_root.to_path_buf();
crate::engine::agent::register_interrupt_cleanup(move || {
if interrupted.finished.swap(true, Ordering::AcqRel) {
return;
}
let event = LfEvent {
run_id: interrupted.run_id.clone(),
ts: OffsetDateTime::now_utc(),
node: LfNode::Run,
event: LfEventType::Escalated,
wave_name: interrupted.wave.clone(),
worktree: Some(directory.display().to_string()),
command: None,
flow: None,
skill: None,
index: None,
error: None,
signal: None,
};
ledger_insert(&interrupted, &event, &directory, Some(130));
});
if same_store {
if let Err(error) = write_exec_process_receipt(&context) {
debug!(error = %error, exec_id = %context.process_id, "live Exec receipt unavailable");
} else {
crate::engine::agent::register_interrupt_cleanup(remove_exec_process_receipt);
}
}
if let Some(wave_name) = wave_name {
if fields.wave_name.as_deref() != Some(wave_name.as_str()) {
debug!(
expected_wave = %wave_name,
observed_wave = ?fields.wave_name,
repo = %repo_root.display(),
"journal run start received mismatched wave metadata"
);
}
}
Ok(Some(context))
}
fn parent_is_recorded(parent: &ExecId) -> bool {
let path = match ledger_db_path() {
Ok(path) if path.exists() => path,
Ok(_) => return false,
Err(err) => {
debug!(error = %err, "no ledger path; keeping the inherited parent process id");
return true;
}
};
match SqliteStore::open_execs_read_only(&path)
.and_then(|store| store.process_is_recorded(parent.as_str()))
{
Ok(recorded) => recorded,
Err(err) => {
debug!(
error = %err,
parent = parent.as_str(),
"ledger unreadable; keeping the inherited parent process id"
);
true
}
}
}
fn configured_run_id(repo_root: &Path) -> Option<TraceId> {
let value = std::env::var(LF_TRACE_ID_ENV).ok()?;
let trimmed = value.trim();
if trimmed.is_empty() {
return None;
}
match trimmed.parse() {
Ok(run_id) => Some(run_id),
Err(err) => {
debug!(
env = LF_TRACE_ID_ENV,
value = trimmed,
repo = %repo_root.display(),
error = %err,
"ignoring invalid journal run id override"
);
None
}
}
}
fn append_event(run_dir: &Path, event: &LfEvent) -> Result<(), std::io::Error> {
let mut file = OpenOptions::new()
.create(true)
.append(true)
.open(events_path(run_dir))?;
let _lock = lock_file(&file)?;
let mut line = serde_json::to_vec(event).map_err(std::io::Error::other)?;
line.push(b'\n');
file.write_all(&line)?;
Ok(())
}
#[cfg(unix)]
struct FileLock {
fd: std::os::fd::RawFd,
}
#[cfg(unix)]
fn lock_file(file: &File) -> Result<FileLock, std::io::Error> {
let fd = file.as_raw_fd();
loop {
if unsafe { libc::flock(fd, libc::LOCK_EX) } == 0 {
return Ok(FileLock { fd });
}
let err = std::io::Error::last_os_error();
if err.kind() != std::io::ErrorKind::Interrupted {
return Err(err);
}
}
}
#[cfg(unix)]
impl Drop for FileLock {
fn drop(&mut self) {
let _ = unsafe { libc::flock(self.fd, libc::LOCK_UN) };
}
}
#[cfg(not(unix))]
struct FileLock;
#[cfg(not(unix))]
fn lock_file(_file: &File) -> Result<FileLock, std::io::Error> {
Ok(FileLock)
}
fn current_context() -> Option<RunContext> {
if is_cli_process() {
return PROCESS_CONTEXT
.lock()
.expect("process context mutex poisoned")
.clone();
}
RUN_CONTEXT.with(|cell| cell.borrow().clone())
}
pub(crate) fn current_process_identity() -> Option<crate::durable::TaskWorkerOwner> {
let context = current_context()?;
Some(crate::durable::TaskWorkerOwner {
trace_id: context.run_id,
exec_id: context.process_id,
pid: std::process::id(),
started_at: context.process_started_at?,
})
}
pub(crate) fn current_exec_id() -> Option<ExecId> {
current_context().map(|context| context.process_id)
}
pub(crate) fn task_worker_owner_evidence(
owner: &crate::durable::TaskWorkerOwner,
) -> ProcessIdentityEvidence {
let receipts = match read_exec_process_receipts_at(&crate::store::lf_home_dir()) {
Ok(receipts) => receipts,
Err(_) => return ProcessIdentityEvidence::Unknown,
};
let receipt = receipts.iter().find(|receipt| {
receipt.trace_id == owner.trace_id.as_str()
&& receipt.exec_id == owner.exec_id.as_str()
&& receipt.pid == owner.pid
&& receipt.started_at == owner.started_at
});
if let Some(receipt) = receipt {
return receipt.process_evidence();
}
let Ok(path) = crate::store::observability_database_path() else {
return ProcessIdentityEvidence::Unknown;
};
if !path.exists() {
return ProcessIdentityEvidence::Unknown;
}
let Ok(store) = SqliteStore::open_execs_read_only(&path) else {
return ProcessIdentityEvidence::Unknown;
};
match store.exec(&owner.exec_id) {
Ok(Some(exec)) if exec.trace_id == owner.trace_id && exec.completed_at.is_some() => {
ProcessIdentityEvidence::Dead
}
_ => ProcessIdentityEvidence::Unknown,
}
}
pub(crate) fn exec_process_evidence(store: &SqliteStore, exec: &ExecId) -> ProcessIdentityEvidence {
let Ok(receipts) = read_exec_process_receipts_at(&crate::store::lf_home_dir()) else {
return ProcessIdentityEvidence::Unknown;
};
if let Some(receipt) = receipts
.iter()
.find(|receipt| receipt.exec_id == exec.as_str())
{
return receipt.process_evidence();
}
match store.exec(exec) {
Ok(Some(record)) if record.completed_at.is_some() => ProcessIdentityEvidence::Dead,
_ => ProcessIdentityEvidence::Unknown,
}
}
pub(crate) fn process_started_at(pid: u32) -> Result<Option<i64>, std::io::Error> {
let output = Command::new("ps")
.args(["-p", &pid.to_string(), "-o", "etime="])
.output()?;
if !output.status.success() {
if output.status.code() == Some(1) && output.stdout.is_empty() && output.stderr.is_empty() {
return Ok(None);
}
return Err(std::io::Error::other(format!(
"process start-time query failed: {}",
String::from_utf8_lossy(&output.stderr).trim()
)));
}
let elapsed = String::from_utf8_lossy(&output.stdout);
let seconds = elapsed_seconds(elapsed.trim())
.ok_or_else(|| std::io::Error::other("process start-time query returned invalid age"))?;
Ok(Some(
OffsetDateTime::now_utc()
.unix_timestamp()
.saturating_sub(i64::try_from(seconds).unwrap_or(i64::MAX)),
))
}
fn elapsed_seconds(value: &str) -> Option<u64> {
let (days, clock) = match value.split_once('-') {
Some((days, clock)) => (days.parse().ok()?, clock),
None => (0_u64, value),
};
let parts = clock
.split(':')
.map(str::parse::<u64>)
.collect::<Result<Vec<_>, _>>()
.ok()?;
let clock = match parts.as_slice() {
[minutes, seconds] => minutes.checked_mul(60)?.checked_add(*seconds)?,
[hours, minutes, seconds] => hours
.checked_mul(3_600)?
.checked_add(minutes.checked_mul(60)?)?
.checked_add(*seconds)?,
_ => return None,
};
days.checked_mul(86_400)?.checked_add(clock)
}
fn set_context(context: RunContext) {
if is_cli_process() {
*PROCESS_CONTEXT
.lock()
.expect("process context mutex poisoned") = Some(context);
return;
}
RUN_CONTEXT.with(|cell| {
*cell.borrow_mut() = Some(context);
});
}
fn clear_context() {
if is_cli_process() {
*PROCESS_CONTEXT
.lock()
.expect("process context mutex poisoned") = None;
return;
}
RUN_CONTEXT.with(|cell| {
*cell.borrow_mut() = None;
});
}
pub(crate) fn read_exec_process_receipts_at(
lf_home: &Path,
) -> Result<Vec<ExecProcessReceipt>, std::io::Error> {
let root = lf_home.join(EXEC_PROCESS_ROOT);
let entries = match fs::read_dir(root) {
Ok(entries) => entries,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(error) => return Err(error),
};
let mut receipts = Vec::new();
for entry in entries {
let entry = entry?;
if !entry.file_type()?.is_file() {
continue;
}
let path = entry.path();
if path.extension().and_then(|value| value.to_str()) != Some("json")
|| path
.file_stem()
.and_then(|value| value.to_str())
.and_then(|value| value.parse::<u32>().ok())
.is_none()
{
continue;
}
let Ok(content) = fs::read(path) else {
continue;
};
let Ok(receipt) = serde_json::from_slice::<ExecProcessReceipt>(&content) else {
continue;
};
if receipt.schema_version == 1 {
receipts.push(receipt);
}
}
Ok(receipts)
}
pub(crate) fn remove_exec_process_receipt_at(
lf_home: &Path,
pid: u32,
) -> Result<bool, std::io::Error> {
let path = lf_home.join(EXEC_PROCESS_ROOT).join(format!("{pid}.json"));
match fs::remove_file(path) {
Ok(()) => Ok(true),
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Ok(false),
Err(error) => Err(error),
}
}
fn write_exec_process_receipt(context: &RunContext) -> Result<(), std::io::Error> {
let Some(started_at) = context.process_started_at else {
return Ok(());
};
let root = crate::store::lf_home_dir().join(EXEC_PROCESS_ROOT);
fs::create_dir_all(&root)?;
let pid = std::process::id();
let receipt = ExecProcessReceipt {
schema_version: 1,
trace_id: context.run_id.to_string(),
exec_id: context.process_id.to_string(),
pid,
started_at,
};
let bytes = serde_json::to_vec(&receipt).map_err(std::io::Error::other)?;
let path = root.join(format!("{pid}.json"));
let temporary = root.join(format!(".{pid}.json.tmp"));
fs::write(&temporary, bytes)?;
fs::rename(temporary, path)
}
fn remove_exec_process_receipt() {
let path = crate::store::lf_home_dir()
.join(EXEC_PROCESS_ROOT)
.join(format!("{}.json", std::process::id()));
if let Err(error) = fs::remove_file(path) {
if error.kind() != std::io::ErrorKind::NotFound {
debug!(error = %error, "failed to remove live Exec receipt");
}
}
}
fn ensure_journal_ignored(repo_root: &Path) -> Result<(), std::io::Error> {
let output = Command::new("git")
.arg("-C")
.arg(repo_root)
.args(["rev-parse", "--git-path", "info/exclude"])
.output()?;
if !output.status.success() {
return Err(std::io::Error::other(format!(
"git rev-parse --git-path info/exclude failed: {}",
String::from_utf8_lossy(&output.stderr).trim()
)));
}
let mut exclude_path = PathBuf::from(String::from_utf8_lossy(&output.stdout).trim());
if exclude_path.is_relative() {
exclude_path = repo_root.join(exclude_path);
}
if let Some(parent) = exclude_path.parent() {
fs::create_dir_all(parent)?;
}
let existing = fs::read_to_string(&exclude_path).unwrap_or_default();
if existing
.lines()
.any(|line| line.trim() == JOURNAL_EXCLUDE_ENTRY)
{
return Ok(());
}
let mut updated = existing;
if !updated.is_empty() && !updated.ends_with('\n') {
updated.push('\n');
}
updated.push_str(JOURNAL_EXCLUDE_ENTRY);
updated.push('\n');
fs::write(exclude_path, updated)
}
#[cfg(test)]
mod tests {
use super::{
emit, events_path, read_events, runs_root, LfEvent, LfEventFields, LfEventType, LfNode,
ProcessIdentityEvidence, TestLedgerGuard,
};
use crate::engine::git::is_clean;
use crate::id::{ExecId, TraceId};
use loopflow_test_support::TestRepo;
use std::path::PathBuf;
use std::process::{Command, Stdio};
const CHILD_APPEND_ENV: &str = "LOOPFLOW_JOURNAL_APPEND_CHILD";
const CHILD_EVENT_COUNT_ENV: &str = "LOOPFLOW_JOURNAL_CHILD_EVENT_COUNT";
const CHILD_RUN_DIR_ENV: &str = "LOOPFLOW_JOURNAL_CHILD_RUN_DIR";
const CHILD_WRITER_ENV: &str = "LOOPFLOW_JOURNAL_CHILD_WRITER";
struct AmbientStorage {
_lock: std::sync::MutexGuard<'static, ()>,
previous_lf_home: Option<std::ffi::OsString>,
previous_db_path: Option<std::ffi::OsString>,
}
impl AmbientStorage {
fn seed(home: &std::path::Path, db_path: &std::path::Path) -> Self {
let lock = super::test_env_lock();
let previous_lf_home = std::env::var_os("LF_HOME");
let previous_db_path = std::env::var_os("LF_DB_PATH");
std::env::set_var("LF_HOME", home);
std::env::set_var("LF_DB_PATH", db_path);
Self {
_lock: lock,
previous_lf_home,
previous_db_path,
}
}
}
impl Drop for AmbientStorage {
fn drop(&mut self) {
match &self.previous_lf_home {
Some(value) => std::env::set_var("LF_HOME", value),
None => std::env::remove_var("LF_HOME"),
}
match &self.previous_db_path {
Some(value) => std::env::set_var("LF_DB_PATH", value),
None => std::env::remove_var("LF_DB_PATH"),
}
}
}
fn with_run_id_env<T>(value: Option<&str>, run: impl FnOnce() -> T) -> T {
let _guard = journal_test_guard();
super::clear_context();
let previous = std::env::var(super::LF_TRACE_ID_ENV).ok();
let previous_process = std::env::var(super::LF_PROCESS_ID_ENV).ok();
std::env::remove_var(super::LF_PROCESS_ID_ENV);
match value {
Some(value) => std::env::set_var(super::LF_TRACE_ID_ENV, value),
None => std::env::remove_var(super::LF_TRACE_ID_ENV),
}
let result = run();
super::clear_context();
match previous {
Some(value) => std::env::set_var(super::LF_TRACE_ID_ENV, value),
None => std::env::remove_var(super::LF_TRACE_ID_ENV),
}
match previous_process {
Some(value) => std::env::set_var(super::LF_PROCESS_ID_ENV, value),
None => std::env::remove_var(super::LF_PROCESS_ID_ENV),
}
result
}
fn journal_test_guard() -> TestLedgerGuard {
let guard = TestLedgerGuard::new();
super::clear_context();
std::env::remove_var(super::LF_TRACE_ID_ENV);
std::env::remove_var(super::LF_PROCESS_ID_ENV);
guard
}
#[test]
fn explicit_test_database_path_controls_the_ledger() {
let guard = journal_test_guard();
let path = guard.home().join("explicit.db");
guard.set_db_path(path.clone());
let opened = super::open_ledger();
opened.expect("open explicit ledger");
assert!(path.exists());
}
#[test]
fn exact_process_evidence_distinguishes_a_live_exec_from_its_completion() {
let _guard = journal_test_guard();
let repo = TestRepo::new();
emit(
repo.path(),
LfNode::Run,
LfEventType::Started,
started_fields(
&["lf".to_string(), "task".to_string()],
repo.path(),
"runtime",
),
);
let owner = super::current_process_identity().expect("live Run identity");
assert_eq!(
super::task_worker_owner_evidence(&owner),
ProcessIdentityEvidence::Live
);
emit(
repo.path(),
LfNode::Run,
LfEventType::Completed,
LfEventFields::default(),
);
assert_eq!(
super::task_worker_owner_evidence(&owner),
ProcessIdentityEvidence::Dead
);
}
#[test]
fn a_killed_registered_exec_is_authoritatively_dead() {
let guard = journal_test_guard();
let mut child = Command::new("sleep")
.arg("30")
.spawn()
.expect("spawn owned process");
let owner = crate::durable::TaskWorkerOwner {
trace_id: TraceId::new(),
exec_id: ExecId::new(),
pid: child.id(),
started_at: time::OffsetDateTime::now_utc().unix_timestamp(),
};
let receipt = super::ExecProcessReceipt {
schema_version: 1,
trace_id: owner.trace_id.to_string(),
exec_id: owner.exec_id.to_string(),
pid: owner.pid,
started_at: owner.started_at,
};
let root = guard.home().join(super::EXEC_PROCESS_ROOT);
std::fs::create_dir_all(&root).expect("create receipt directory");
std::fs::write(
root.join(format!("{}.json", owner.pid)),
serde_json::to_vec(&receipt).expect("serialize receipt"),
)
.expect("write receipt");
assert_eq!(
super::task_worker_owner_evidence(&owner),
ProcessIdentityEvidence::Live
);
child.kill().expect("kill owned process");
child.wait().expect("reap owned process");
assert_eq!(
super::task_worker_owner_evidence(&owner),
ProcessIdentityEvidence::Dead
);
}
#[test]
fn unit_test_ledger_ignores_ambient_storage_paths() {
let ambient_home = tempfile::tempdir().expect("ambient home");
let ambient_db_dir = tempfile::tempdir().expect("ambient db dir");
let ambient_db = ambient_db_dir.path().join("production.db");
let _ambient = AmbientStorage::seed(ambient_home.path(), &ambient_db);
let resolved = super::ledger_db_path().expect("test ledger path");
super::open_ledger().expect("open test ledger");
assert_ne!(resolved, ambient_db);
assert_ne!(resolved, ambient_home.path().join("loopflow.db"));
assert!(!ambient_db.exists());
assert!(!ambient_home.path().join("loopflow.db").exists());
}
fn started_fields(
command: &[String],
worktree: &std::path::Path,
wave_name: &str,
) -> LfEventFields {
LfEventFields {
wave_name: Some(wave_name.to_string()),
worktree: Some(worktree.display().to_string()),
command: Some(command.to_vec()),
..LfEventFields::default()
}
}
fn only_run_dir(worktree: &std::path::Path) -> std::path::PathBuf {
let mut entries = std::fs::read_dir(runs_root(worktree))
.expect("read runs")
.map(|entry| entry.expect("run dir entry").path())
.collect::<Vec<_>>();
assert_eq!(entries.len(), 1, "expected a single journal run dir");
entries.pop().expect("run dir")
}
#[test]
fn concurrent_child_process_appends_keep_events_jsonl_parseable() {
let tmp = tempfile::TempDir::new().expect("temp journal");
let run_dir = tmp.path().join("run");
std::fs::create_dir_all(&run_dir).expect("create run dir");
let current_exe = std::env::current_exe().expect("current test binary");
let writers = 8;
let events_per_writer = 20;
let mut children = Vec::new();
for writer in 0..writers {
let child = Command::new(¤t_exe)
.arg("journal_child_process_appends_events_for_concurrency_regression")
.arg("--nocapture")
.env(CHILD_APPEND_ENV, "1")
.env(CHILD_RUN_DIR_ENV, &run_dir)
.env(CHILD_WRITER_ENV, writer.to_string())
.env(CHILD_EVENT_COUNT_ENV, events_per_writer.to_string())
.stdout(Stdio::null())
.stderr(Stdio::null())
.spawn()
.expect("spawn child journal writer");
children.push(child);
}
for mut child in children {
let status = child.wait().expect("wait for child journal writer");
assert!(status.success(), "child journal writer failed: {status}");
}
let raw = std::fs::read_to_string(events_path(&run_dir)).expect("read events.jsonl");
let lines = raw
.lines()
.filter(|line| !line.trim().is_empty())
.collect::<Vec<_>>();
assert_eq!(lines.len(), writers * events_per_writer);
for (line_number, line) in lines.iter().enumerate() {
serde_json::from_str::<LfEvent>(line).unwrap_or_else(|err| {
panic!("line {} is malformed JSONL: {err}: {line}", line_number + 1)
});
}
let parsed = read_events(&run_dir).expect("parse all events");
assert_eq!(parsed.len(), writers * events_per_writer);
}
#[test]
fn journal_child_process_appends_events_for_concurrency_regression() {
if std::env::var(CHILD_APPEND_ENV).ok().as_deref() != Some("1") {
return;
}
let run_dir = PathBuf::from(std::env::var(CHILD_RUN_DIR_ENV).expect("child run dir env"));
let writer = std::env::var(CHILD_WRITER_ENV).expect("child writer env");
let event_count = std::env::var(CHILD_EVENT_COUNT_ENV)
.expect("child event count env")
.parse::<usize>()
.expect("child event count");
let run_id = TraceId::parse("8985c55b-9864-4c2b-860f-b7054a71bbea").expect("run id");
for index in 0..event_count {
let event = LfEvent {
run_id: run_id.clone(),
ts: time::OffsetDateTime::now_utc(),
node: LfNode::Skill,
event: LfEventType::Errored,
wave_name: Some("meta".to_string()),
worktree: Some(run_dir.display().to_string()),
command: None,
flow: Some("garden".to_string()),
skill: Some(format!("writer-{writer}-{index}")),
index: Some(index as u32),
error: Some(format!(
"writer-{writer}-event-{index}:{}",
"x".repeat(16 * 1024)
)),
signal: None,
};
super::append_event(&run_dir, &event).expect("child append event");
}
}
#[test]
fn a_nested_lf_gets_its_own_span_and_names_its_parent() {
let _guard = journal_test_guard();
let repo = TestRepo::new();
let fields = started_fields(&["lf".to_string(), "wave".to_string()], repo.path(), "main");
super::emit(
repo.path(),
LfNode::Run,
LfEventType::Started,
fields.clone(),
);
let parent = super::current_context().expect("parent");
super::clear_context();
let child = super::ensure_run_context(repo.path(), &fields)
.expect("child context")
.expect("child");
assert_ne!(parent.process_id, child.process_id);
assert_eq!(
child.run_id, parent.run_id,
"a nested lf stays in the trace"
);
assert_eq!(child.parent_process_id, Some(parent.process_id));
super::clear_context();
std::env::remove_var(super::LF_PROCESS_ID_ENV);
std::env::remove_var(super::LF_TRACE_ID_ENV);
}
#[test]
fn an_inherited_parent_the_ledger_never_recorded_is_dropped() {
let _guard = journal_test_guard();
let repo = TestRepo::new();
let fields = started_fields(
&["lf".to_string(), "status".to_string()],
repo.path(),
"main",
);
super::emit(
repo.path(),
LfNode::Run,
LfEventType::Started,
fields.clone(),
);
let recorded = super::current_context().expect("recorded run");
super::clear_context();
let ghost = ExecId::new();
std::env::set_var(super::LF_TRACE_ID_ENV, recorded.run_id.as_str());
std::env::set_var(super::LF_PROCESS_ID_ENV, ghost.as_str());
let context = super::ensure_run_context(repo.path(), &fields)
.expect("run context")
.expect("context");
assert!(
super::parent_is_recorded(&recorded.process_id),
"the ledger must really hold a parent, or the assertion below \
passes for the wrong reason"
);
assert_eq!(
context.parent_process_id, None,
"a parent the ledger never recorded is a ghost, not lineage"
);
assert!(
!context.minted_run_id,
"the trace still groups; only the false pointer goes"
);
super::clear_context();
std::env::remove_var(super::LF_PROCESS_ID_ENV);
std::env::remove_var(super::LF_TRACE_ID_ENV);
}
#[test]
fn a_detached_body_inherits_a_recorded_parent_across_its_launcher() {
let _guard = journal_test_guard();
let repo = TestRepo::new();
let fields = started_fields(&["lf".to_string(), "task".to_string()], repo.path(), "main");
super::emit(
repo.path(),
LfNode::Run,
LfEventType::Started,
fields.clone(),
);
let launcher = super::current_context().expect("launcher");
super::emit(
repo.path(),
LfNode::Run,
LfEventType::Completed,
LfEventFields::default(),
);
std::env::set_var(super::LF_TRACE_ID_ENV, launcher.run_id.as_str());
std::env::set_var(super::LF_PROCESS_ID_ENV, launcher.process_id.as_str());
super::clear_context();
let body = super::ensure_run_context(repo.path(), &fields)
.expect("body context")
.expect("body");
assert_eq!(body.parent_process_id, Some(launcher.process_id));
assert_eq!(body.run_id, launcher.run_id);
super::clear_context();
std::env::remove_var(super::LF_PROCESS_ID_ENV);
std::env::remove_var(super::LF_TRACE_ID_ENV);
}
#[test]
fn minting_a_fresh_run_drops_a_stale_cross_trace_parent() {
let _guard = journal_test_guard();
let repo = TestRepo::new();
let fields = started_fields(
&["lf".to_string(), "kickoff".to_string()],
repo.path(),
"main",
);
std::env::set_var(super::LF_PROCESS_ID_ENV, ExecId::new().as_str());
let context = super::ensure_run_context(repo.path(), &fields)
.expect("run context")
.expect("context");
assert!(context.minted_run_id, "no LF_TRACE_ID means a fresh trace");
assert_eq!(
context.parent_process_id, None,
"a fresh trace has no in-trace parent to name"
);
super::clear_context();
std::env::remove_var(super::LF_PROCESS_ID_ENV);
std::env::remove_var(super::LF_TRACE_ID_ENV);
}
#[test]
fn journal_writes_run_flow_and_skill_events_in_wave_worktree() {
let _guard = journal_test_guard();
let repo = TestRepo::new();
let worktree = repo.create_named_worktree("runtime");
let command = vec!["lf".to_string(), "build".to_string()];
emit(
&worktree,
LfNode::Run,
LfEventType::Started,
started_fields(&command, &worktree, "runtime"),
);
emit(
&worktree,
LfNode::Flow,
LfEventType::Started,
LfEventFields {
flow: Some("build".to_string()),
..LfEventFields::default()
},
);
emit(
&worktree,
LfNode::Skill,
LfEventType::Started,
LfEventFields {
skill: Some("implement".to_string()),
index: Some(0),
..LfEventFields::default()
},
);
emit(
&worktree,
LfNode::Skill,
LfEventType::Completed,
LfEventFields {
skill: Some("implement".to_string()),
index: Some(0),
..LfEventFields::default()
},
);
emit(
&worktree,
LfNode::Flow,
LfEventType::Completed,
LfEventFields::default(),
);
emit(
&worktree,
LfNode::Run,
LfEventType::Completed,
LfEventFields::default(),
);
let run_dir = only_run_dir(&worktree);
let events = read_events(&run_dir).expect("read events");
assert_eq!(events.len(), 6);
assert_eq!(events[0].node, LfNode::Run);
assert_eq!(events[0].event, LfEventType::Started);
assert_eq!(events[0].wave_name.as_deref(), Some("runtime"));
assert_eq!(events[1].node, LfNode::Flow);
assert_eq!(events[1].flow.as_deref(), Some("build"));
assert_eq!(events[2].node, LfNode::Skill);
assert_eq!(events[2].skill.as_deref(), Some("implement"));
assert_eq!(events[2].index, Some(0));
assert_eq!(events[3].event, LfEventType::Completed);
assert_eq!(events[4].node, LfNode::Flow);
assert_eq!(events[5].node, LfNode::Run);
assert_eq!(events[5].event, LfEventType::Completed);
}
#[test]
fn journal_keeps_worktree_clean() {
let _guard = journal_test_guard();
let repo = TestRepo::new();
let worktree = repo.create_named_worktree("runtime");
let command = vec!["lf".to_string(), "build".to_string()];
emit(
&worktree,
LfNode::Run,
LfEventType::Started,
started_fields(&command, &worktree, "runtime"),
);
assert!(is_clean(&worktree).expect("worktree should stay clean"));
}
#[test]
fn run_lifecycle_publishes_and_removes_exact_process_ownership() {
let guard = journal_test_guard();
let repo = TestRepo::new();
let worktree = repo.create_named_worktree("runtime");
let command = vec!["lf".to_string(), "build".to_string()];
emit(
&worktree,
LfNode::Run,
LfEventType::Started,
started_fields(&command, &worktree, "runtime"),
);
let receipts =
super::read_exec_process_receipts_at(guard.home()).expect("read live receipt");
assert_eq!(receipts.len(), 1);
assert_eq!(receipts[0].pid, std::process::id());
assert_eq!(
receipts[0].exec_id,
std::env::var(super::LF_PROCESS_ID_ENV).expect("current Exec id")
);
emit(
&worktree,
LfNode::Run,
LfEventType::Completed,
LfEventFields::default(),
);
assert!(super::read_exec_process_receipts_at(guard.home())
.expect("read terminal receipt state")
.is_empty());
}
#[test]
fn terminal_run_events_clear_context_for_the_next_run() {
let _guard = journal_test_guard();
let repo = TestRepo::new();
let worktree = repo.create_named_worktree("runtime");
let command = vec!["lf".to_string(), "build".to_string()];
emit(
&worktree,
LfNode::Run,
LfEventType::Started,
started_fields(&command, &worktree, "runtime"),
);
emit(
&worktree,
LfNode::Run,
LfEventType::Completed,
LfEventFields::default(),
);
emit(
&worktree,
LfNode::Run,
LfEventType::Started,
started_fields(&command, &worktree, "runtime"),
);
emit(
&worktree,
LfNode::Run,
LfEventType::Completed,
LfEventFields::default(),
);
let entries = std::fs::read_dir(runs_root(&worktree))
.expect("read runs")
.count();
assert_eq!(entries, 2);
}
#[test]
fn journal_uses_configured_run_id_when_present() {
with_run_id_env(Some("7c22895f-e4c1-49cc-a95d-2267e2356f16"), || {
let repo = TestRepo::new();
let worktree = repo.create_named_worktree("runtime");
let command = vec!["lf".to_string(), "build".to_string()];
emit(
&worktree,
LfNode::Run,
LfEventType::Started,
started_fields(&command, &worktree, "runtime"),
);
emit(
&worktree,
LfNode::Run,
LfEventType::Completed,
LfEventFields::default(),
);
let run_dir = only_run_dir(&worktree);
let run_id = run_dir
.file_name()
.and_then(|name| name.to_str())
.expect("run dir name");
assert_eq!(run_id, "7c22895f-e4c1-49cc-a95d-2267e2356f16");
});
}
#[test]
fn invalid_configured_run_id_falls_back_to_generated_id() {
with_run_id_env(Some("not-a-uuid"), || {
let repo = TestRepo::new();
let worktree = repo.create_named_worktree("runtime");
let command = vec!["lf".to_string(), "build".to_string()];
emit(
&worktree,
LfNode::Run,
LfEventType::Started,
started_fields(&command, &worktree, "runtime"),
);
emit(
&worktree,
LfNode::Run,
LfEventType::Completed,
LfEventFields::default(),
);
let run_dir = only_run_dir(&worktree);
let run_id = run_dir
.file_name()
.and_then(|name| name.to_str())
.expect("run dir name");
assert!(
TraceId::parse(run_id).is_ok(),
"expected generated UUID run id"
);
assert_ne!(run_id, "not-a-uuid");
});
}
}