use std::collections::BTreeMap;
use std::future::Future;
use std::path::Path;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use time::OffsetDateTime;
use crate::engine::wave_context::{resolve_ambient_wave, AmbientWaveRef};
use crate::lfd::id::LfdId;
use crate::lfd::types::{Session, SessionStatus, SessionUse, LF_CLI_SOURCE};
use crate::lfdb::{open_existing_store, SharedStore};
pub const WAVE_ID_ENV: &str = "LFD_WAVE_ID";
pub const CHANNEL_ENV: &str = "LFD_CHANNEL";
pub const SESSION_ID_ENV: &str = "LFD_SESSION_ID";
pub const SESSION_INHERITED_ENV: &str = "LFD_SESSION_INHERITED";
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RunContext {
Outside,
OwnSession,
NeedsRegistration {
wave: AmbientWaveRef,
parent_session_id: Option<String>,
},
}
pub fn classify_run_context(
ambient: Option<AmbientWaveRef>,
session_id: Option<&str>,
session_inherited: bool,
) -> RunContext {
let Some(wave) = ambient else {
return RunContext::Outside;
};
let session_id = session_id.filter(|value| !value.is_empty());
if session_id.is_some() && !session_inherited {
return RunContext::OwnSession;
}
RunContext::NeedsRegistration {
wave,
parent_session_id: session_id.map(str::to_string),
}
}
pub fn register_run(step: &str, agent: &str, argv: &[String]) -> Option<RunSession> {
let repo_root = crate::lf::commands::util::find_repo_root().ok();
register_run_in(repo_root.as_deref(), step, agent, argv)
}
fn register_run_in(
repo_root: Option<&Path>,
step: &str,
agent: &str,
argv: &[String],
) -> Option<RunSession> {
let context = classify_run_context(
resolve_ambient_wave(env_var(WAVE_ID_ENV).as_deref(), repo_root),
env_var(SESSION_ID_ENV).as_deref(),
env_var(SESSION_INHERITED_ENV).is_some(),
);
let (wave, parent_session_id) = match context {
RunContext::Outside => return None,
RunContext::OwnSession => {
mark_child_sessions_inherited();
return adopt_own_session();
}
RunContext::NeedsRegistration {
wave,
parent_session_id,
} => (wave, parent_session_id),
};
let parent_session_id: Option<LfdId> = match parent_session_id {
Some(value) => Some(value.parse().ok()?),
None => None,
};
let session = block_on(register_session(
wave,
parent_session_id,
step.to_string(),
agent.to_string(),
argv.to_vec(),
))??;
std::env::set_var(SESSION_ID_ENV, session.session_id());
mark_child_sessions_inherited();
let interrupted = Arc::clone(&session.inner);
crate::engine::agent::register_interrupt_cleanup(move || interrupted.complete(130));
Some(session)
}
fn adopt_own_session() -> Option<RunSession> {
let session_id: LfdId = env_var(SESSION_ID_ENV)?.parse().ok()?;
let (store, session) = block_on(async move {
let store: SharedStore = Arc::new(open_existing_store().await?);
let session = store.get_control_session(&session_id).await.ok()??;
Some((store, session))
})??;
if session.status.is_terminal() {
return None;
}
let session = RunSession {
inner: Arc::new(SessionHandle {
store,
session,
completed: AtomicBool::new(false),
}),
};
let interrupted = Arc::clone(&session.inner);
crate::engine::agent::register_interrupt_cleanup(move || interrupted.complete(130));
Some(session)
}
pub fn mark_child_sessions_inherited() {
if env_var(SESSION_ID_ENV).is_some() {
std::env::set_var(SESSION_INHERITED_ENV, "1");
}
}
#[derive(Debug)]
pub struct RunSession {
inner: Arc<SessionHandle>,
}
impl RunSession {
pub fn complete(&self, exit_code: i32) {
self.inner.complete(exit_code);
}
pub fn session_id(&self) -> &str {
self.inner.session.id.as_str()
}
}
impl Drop for RunSession {
fn drop(&mut self) {
self.inner.complete(1);
}
}
#[derive(Debug)]
struct SessionHandle {
store: SharedStore,
session: Session,
completed: AtomicBool,
}
impl SessionHandle {
fn complete(&self, exit_code: i32) {
if self.completed.swap(true, Ordering::SeqCst) {
return;
}
let store = self.store.clone();
let id = self.session.id.clone();
let _ = block_on(async move {
let Some(mut row) = store.get_control_session(&id).await.ok().flatten() else {
return;
};
if !row.complete(exit_code) {
return;
}
let _ = store.update_control_session(&row).await;
});
}
}
async fn register_session(
wave: AmbientWaveRef,
mut parent_session_id: Option<LfdId>,
step: String,
agent: String,
argv: Vec<String>,
) -> Option<RunSession> {
let store: SharedStore = Arc::new(open_existing_store().await?);
let wave_id = match wave {
AmbientWaveRef::Id(id) => {
let id: LfdId = id.parse().ok()?;
store.get_wave(&id).await.ok()??;
id
}
AmbientWaveRef::Name(name) => store.get_wave_by_name(&name).await.ok()??.id().clone(),
};
if let Some(parent) = &parent_session_id {
if store.get_control_session(parent).await.ok()?.is_none() {
parent_session_id = None;
}
}
let now = OffsetDateTime::now_utc();
let session = Session {
id: LfdId::new(),
wave_id,
run_id: None,
parent_session_id,
session_use: SessionUse::Worker,
step,
agent,
cwd: std::env::current_dir()
.map(|cwd| cwd.display().to_string())
.unwrap_or_default(),
argv,
env: BTreeMap::new(),
source: LF_CLI_SOURCE.to_string(),
tmux_name: current_tmux_session_name().unwrap_or_default(),
status: SessionStatus::Running,
attached_at: Some(now),
started_at: Some(now),
completed_at: None,
created_at: now,
completion_token: None,
};
store.register_session(&session).await.ok()?;
Some(RunSession {
inner: Arc::new(SessionHandle {
store,
session,
completed: AtomicBool::new(false),
}),
})
}
fn block_on<T: Send + 'static>(future: impl Future<Output = T> + Send + 'static) -> Option<T> {
if tokio::runtime::Handle::try_current().is_ok() {
return std::thread::spawn(move || block_on_new_runtime(future))
.join()
.ok();
}
Some(block_on_new_runtime(future))
}
fn block_on_new_runtime<T>(future: impl Future<Output = T>) -> T {
tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime always builds")
.block_on(future)
}
fn env_var(key: &str) -> Option<String> {
std::env::var(key).ok().filter(|value| !value.is_empty())
}
#[cfg(test)]
pub(crate) fn test_env_lock() -> std::sync::MutexGuard<'static, ()> {
static LOCK: std::sync::Mutex<()> = std::sync::Mutex::new(());
LOCK.lock().unwrap_or_else(|poisoned| poisoned.into_inner())
}
fn current_tmux_session_name() -> Option<String> {
std::env::var_os("TMUX")?;
let output = std::process::Command::new("tmux")
.args(["display-message", "-p", "#S"])
.output()
.ok()?;
if !output.status.success() {
return None;
}
let name = String::from_utf8_lossy(&output.stdout).trim().to_string();
(!name.is_empty()).then_some(name)
}
#[cfg(test)]
mod tests {
use std::path::Path;
use time::OffsetDateTime;
use crate::lfd::id::LfdId;
use crate::lfd::types::{
RepoWork, Session, SessionStatus, Wave, WaveStatus, TMUX_TERMINAL_SOURCE,
};
use crate::lfdb::{open_store, StorageConfig};
use crate::engine::wave_context::AmbientWaveRef;
use super::{block_on_new_runtime, classify_run_context, register_run_in, RunContext};
fn env_wave(id: &str) -> Option<AmbientWaveRef> {
Some(AmbientWaveRef::Id(id.to_string()))
}
#[test]
fn no_ambient_wave_means_no_registration() {
assert_eq!(
classify_run_context(None, Some("sess-1"), true),
RunContext::Outside
);
assert_eq!(classify_run_context(None, None, false), RunContext::Outside);
}
#[test]
fn executor_launched_runs_do_not_register_again() {
assert_eq!(
classify_run_context(env_wave("wave-1"), Some("sess-1"), false),
RunContext::OwnSession
);
}
#[test]
fn inherited_session_id_becomes_the_parent() {
assert_eq!(
classify_run_context(env_wave("wave-1"), Some("sess-1"), true),
RunContext::NeedsRegistration {
wave: AmbientWaveRef::Id("wave-1".to_string()),
parent_session_id: Some("sess-1".to_string()),
}
);
}
#[test]
fn wave_context_without_session_registers_a_root_child() {
assert_eq!(
classify_run_context(Some(AmbientWaveRef::Name("goals".to_string())), None, false),
RunContext::NeedsRegistration {
wave: AmbientWaveRef::Name("goals".to_string()),
parent_session_id: None,
}
);
}
fn clear_session_env() {
for key in [
super::WAVE_ID_ENV,
super::SESSION_ID_ENV,
super::SESSION_INHERITED_ENV,
"LFD_DB_PATH",
] {
std::env::remove_var(key);
}
}
fn make_wave(repo: &str) -> Wave {
Wave {
id: LfdId::new(),
name: "registry-wave".to_string(),
primary_flow: "ship-roadmap".to_string(),
goal: "ship-roadmap".to_string(),
metrics: Vec::new(),
repos: vec![RepoWork {
repo: repo.to_string(),
worktree: String::new(),
branch: String::new(),
status: WaveStatus::Idle,
iteration: 0,
cycle_start_iteration: 0,
position: 0,
}],
direction: Vec::new(),
area: Vec::new(),
paused: false,
created_at: Some(OffsetDateTime::now_utc()),
workers: 1,
parent_wave_id: None,
}
}
fn seed_registry(path: &Path) -> Wave {
let wave = make_wave("/tmp/repo");
let seeded = wave.clone();
let path_buf = path.to_path_buf();
block_on_new_runtime(async move {
let store = open_store(&StorageConfig::sqlite(path_buf))
.await
.expect("open registry store");
store.create_wave(&seeded).await.expect("seed wave");
});
std::env::set_var("LFD_DB_PATH", path);
wave
}
fn seed_parent_session(path: &Path, wave: &Wave) -> LfdId {
let now = OffsetDateTime::now_utc();
let parent = Session {
id: LfdId::new(),
wave_id: wave.id().clone(),
run_id: None,
parent_session_id: None,
session_use: crate::lfd::types::SessionUse::WaveAgent,
step: "mind".to_string(),
agent: "lf".to_string(),
cwd: "/tmp/repo".to_string(),
argv: Vec::new(),
env: Default::default(),
source: "wave_server".to_string(),
tmux_name: String::new(),
status: SessionStatus::Running,
attached_at: None,
started_at: Some(now),
completed_at: None,
created_at: now,
completion_token: None,
};
let id = parent.id.clone();
let path = path.to_path_buf();
block_on_new_runtime(async move {
let store = open_store(&StorageConfig::sqlite(path))
.await
.expect("open registry store");
store
.register_session(&parent)
.await
.expect("seed parent session");
});
id
}
fn seed_own_session(path: &Path, wave: &Wave) -> Session {
let now = OffsetDateTime::now_utc();
let session = Session {
id: LfdId::new(),
wave_id: wave.id().clone(),
run_id: None,
parent_session_id: None,
session_use: crate::lfd::types::SessionUse::Worker,
step: "dispatch:implement".to_string(),
agent: "lf".to_string(),
cwd: "/tmp/repo".to_string(),
argv: Vec::new(),
env: Default::default(),
source: TMUX_TERMINAL_SOURCE.to_string(),
tmux_name: "lf-worker".to_string(),
status: SessionStatus::Running,
attached_at: Some(now),
started_at: Some(now),
completed_at: None,
created_at: now,
completion_token: None,
};
let stored = session.clone();
let path = path.to_path_buf();
block_on_new_runtime(async move {
let store = open_store(&StorageConfig::sqlite(path))
.await
.expect("open registry store");
store
.register_session(&stored)
.await
.expect("seed own session");
});
session
}
fn store_session(path: &Path, session: &Session) {
let session = session.clone();
let path = path.to_path_buf();
block_on_new_runtime(async move {
let store = open_store(&StorageConfig::sqlite(path))
.await
.expect("open registry store");
store
.update_control_session(&session)
.await
.expect("update session");
});
}
fn count_sessions(path: &Path, wave: &Wave) -> usize {
let wave_id = wave.id().clone();
let path = path.to_path_buf();
block_on_new_runtime(async move {
let store = open_store(&StorageConfig::sqlite(path))
.await
.expect("open registry store");
store
.list_control_sessions(Some(&wave_id), None)
.await
.expect("session list")
.len()
})
}
fn load_session(path: &Path, id: &str) -> Session {
let id: LfdId = id.parse().expect("session id");
let path = path.to_path_buf();
block_on_new_runtime(async move {
let store = open_store(&StorageConfig::sqlite(path))
.await
.expect("open registry store");
store
.get_control_session(&id)
.await
.expect("session lookup")
.expect("session stored")
})
}
#[test]
fn register_run_without_env_is_a_no_op() {
let _guard = super::test_env_lock();
clear_session_env();
let tmp = tempfile::tempdir().expect("tempdir");
assert!(register_run_in(Some(tmp.path()), "design", "lf", &["lf".to_string()]).is_none());
assert!(register_run_in(None, "design", "lf", &["lf".to_string()]).is_none());
assert!(std::env::var(super::SESSION_INHERITED_ENV).is_err());
}
#[test]
fn bare_run_in_wave_worktree_registers_by_worktree_name() {
let _guard = super::test_env_lock();
clear_session_env();
let repo = loopflow_test_support::TestRepo::new();
let worktree = repo.create_wave_worktree("registry-wave");
let tmp = tempfile::tempdir().expect("tempdir");
let db = tmp.path().join("lfd.db");
let wave = seed_registry(&db);
let argv = vec!["lf".to_string(), "design".to_string()];
let registration =
register_run_in(Some(&worktree), "design", "lf", &argv).expect("registered");
let stored = load_session(&db, registration.session_id());
assert_eq!(stored.status, SessionStatus::Running);
assert_eq!(stored.wave_id, *wave.id());
assert_eq!(stored.parent_session_id, None, "root-parented");
assert_eq!(stored.step, "design");
assert_eq!(
std::env::var(super::SESSION_ID_ENV).as_deref(),
Ok(registration.session_id())
);
registration.complete(0);
assert_eq!(
load_session(&db, registration.session_id()).status,
SessionStatus::Succeeded
);
clear_session_env();
}
#[test]
fn executor_launched_run_marks_children_without_registering() {
let _guard = super::test_env_lock();
clear_session_env();
std::env::set_var(super::WAVE_ID_ENV, "wave-1");
std::env::set_var(super::SESSION_ID_ENV, "sess-1");
let registration = register_run_in(None, "design", "lf", &["lf".to_string()]);
assert!(registration.is_none());
assert_eq!(
std::env::var(super::SESSION_INHERITED_ENV).as_deref(),
Ok("1")
);
assert_eq!(
std::env::var(super::SESSION_ID_ENV).as_deref(),
Ok("sess-1")
);
clear_session_env();
}
#[test]
fn register_run_writes_running_row_and_completes_by_exit_code() {
let _guard = super::test_env_lock();
clear_session_env();
let tmp = tempfile::tempdir().expect("tempdir");
let db = tmp.path().join("lfd.db");
let wave = seed_registry(&db);
let parent = seed_parent_session(&db, &wave);
std::env::set_var(super::WAVE_ID_ENV, wave.id().to_string());
std::env::set_var(super::SESSION_ID_ENV, parent.to_string());
std::env::set_var(super::SESSION_INHERITED_ENV, "1");
let argv = vec!["lf".to_string(), "design".to_string()];
let registration = register_run_in(None, "design", "lf", &argv).expect("registered");
let stored = load_session(&db, registration.session_id());
assert_eq!(stored.status, SessionStatus::Running);
assert_eq!(stored.wave_id, *wave.id());
assert_eq!(stored.parent_session_id, Some(parent));
assert_eq!(stored.argv, argv);
assert_eq!(stored.step, "design");
assert_eq!(stored.source, "lf_cli");
assert_eq!(
std::env::var(super::SESSION_ID_ENV).as_deref(),
Ok(registration.session_id())
);
assert_eq!(
std::env::var(super::SESSION_INHERITED_ENV).as_deref(),
Ok("1")
);
registration.complete(0);
let stored = load_session(&db, registration.session_id());
assert_eq!(stored.status, SessionStatus::Succeeded);
let id = registration.session_id().to_string();
drop(registration);
assert_eq!(load_session(&db, &id).status, SessionStatus::Succeeded);
clear_session_env();
}
#[test]
fn nonzero_exit_and_drop_record_failure() {
let _guard = super::test_env_lock();
clear_session_env();
let tmp = tempfile::tempdir().expect("tempdir");
let db = tmp.path().join("lfd.db");
let wave = seed_registry(&db);
std::env::set_var(super::WAVE_ID_ENV, wave.id().to_string());
let failed = register_run_in(None, "debug", "lf", &["lf".to_string()]).expect("registered");
failed.complete(3);
assert_eq!(
load_session(&db, failed.session_id()).status,
SessionStatus::Failed
);
std::env::remove_var(super::SESSION_ID_ENV);
std::env::remove_var(super::SESSION_INHERITED_ENV);
let dropped =
register_run_in(None, "debug", "lf", &["lf".to_string()]).expect("registered");
let id = dropped.session_id().to_string();
drop(dropped);
assert_eq!(load_session(&db, &id).status, SessionStatus::Failed);
clear_session_env();
}
#[test]
fn register_run_is_silent_when_no_registry_exists() {
let _guard = super::test_env_lock();
clear_session_env();
let tmp = tempfile::tempdir().expect("tempdir");
let db = tmp.path().join("never-created.db");
std::env::set_var("LFD_DB_PATH", &db);
std::env::set_var(super::WAVE_ID_ENV, LfdId::new().to_string());
assert!(register_run_in(None, "design", "lf", &["lf".to_string()]).is_none());
assert!(!db.exists(), "a missing registry must not be conjured");
clear_session_env();
}
#[test]
fn own_session_run_completes_its_existing_row_by_exit_code() {
let _guard = super::test_env_lock();
clear_session_env();
let tmp = tempfile::tempdir().expect("tempdir");
let db = tmp.path().join("lfd.db");
let wave = seed_registry(&db);
let own = seed_own_session(&db, &wave);
std::env::set_var(super::WAVE_ID_ENV, wave.id().to_string());
std::env::set_var(super::SESSION_ID_ENV, own.id.to_string());
let adoption =
register_run_in(None, "implement", "lf", &["lf".to_string()]).expect("adopted");
assert_eq!(adoption.session_id(), own.id.as_str());
assert_eq!(count_sessions(&db, &wave), 1);
assert_eq!(
std::env::var(super::SESSION_ID_ENV).as_deref(),
Ok(own.id.as_str())
);
assert_eq!(
std::env::var(super::SESSION_INHERITED_ENV).as_deref(),
Ok("1")
);
adoption.complete(0);
assert_eq!(
load_session(&db, own.id.as_str()).status,
SessionStatus::Succeeded
);
drop(adoption);
assert_eq!(
load_session(&db, own.id.as_str()).status,
SessionStatus::Succeeded
);
clear_session_env();
}
#[test]
fn own_session_nonzero_exit_and_drop_record_failure() {
let _guard = super::test_env_lock();
clear_session_env();
let tmp = tempfile::tempdir().expect("tempdir");
let db = tmp.path().join("lfd.db");
let wave = seed_registry(&db);
std::env::set_var(super::WAVE_ID_ENV, wave.id().to_string());
let failed_row = seed_own_session(&db, &wave);
std::env::set_var(super::SESSION_ID_ENV, failed_row.id.to_string());
std::env::remove_var(super::SESSION_INHERITED_ENV);
let failed =
register_run_in(None, "implement", "lf", &["lf".to_string()]).expect("adopted");
failed.complete(3);
assert_eq!(
load_session(&db, failed_row.id.as_str()).status,
SessionStatus::Failed
);
let dropped_row = seed_own_session(&db, &wave);
std::env::set_var(super::SESSION_ID_ENV, dropped_row.id.to_string());
std::env::remove_var(super::SESSION_INHERITED_ENV);
let dropped =
register_run_in(None, "implement", "lf", &["lf".to_string()]).expect("adopted");
drop(dropped);
assert_eq!(
load_session(&db, dropped_row.id.as_str()).status,
SessionStatus::Failed
);
clear_session_env();
}
#[test]
fn own_session_with_missing_row_or_store_is_silent() {
let _guard = super::test_env_lock();
clear_session_env();
let tmp = tempfile::tempdir().expect("tempdir");
let db = tmp.path().join("lfd.db");
let wave = seed_registry(&db);
std::env::set_var(super::WAVE_ID_ENV, wave.id().to_string());
std::env::set_var(super::SESSION_ID_ENV, LfdId::new().to_string());
assert!(register_run_in(None, "implement", "lf", &["lf".to_string()]).is_none());
assert_eq!(count_sessions(&db, &wave), 0);
let missing = tmp.path().join("never-created.db");
std::env::set_var("LFD_DB_PATH", &missing);
std::env::remove_var(super::SESSION_INHERITED_ENV);
assert!(register_run_in(None, "implement", "lf", &["lf".to_string()]).is_none());
assert!(!missing.exists(), "a missing registry must not be conjured");
clear_session_env();
}
#[test]
fn own_session_never_overwrites_a_terminal_row() {
let _guard = super::test_env_lock();
clear_session_env();
let tmp = tempfile::tempdir().expect("tempdir");
let db = tmp.path().join("lfd.db");
let wave = seed_registry(&db);
std::env::set_var(super::WAVE_ID_ENV, wave.id().to_string());
let mut canceled = seed_own_session(&db, &wave);
assert!(canceled.cancel());
store_session(&db, &canceled);
std::env::set_var(super::SESSION_ID_ENV, canceled.id.to_string());
assert!(register_run_in(None, "implement", "lf", &["lf".to_string()]).is_none());
assert_eq!(
load_session(&db, canceled.id.as_str()).status,
SessionStatus::Canceled
);
let raced = seed_own_session(&db, &wave);
std::env::set_var(super::SESSION_ID_ENV, raced.id.to_string());
std::env::remove_var(super::SESSION_INHERITED_ENV);
let adoption =
register_run_in(None, "implement", "lf", &["lf".to_string()]).expect("adopted");
let mut closed = raced.clone();
assert!(closed.cancel());
store_session(&db, &closed);
adoption.complete(0);
assert_eq!(
load_session(&db, raced.id.as_str()).status,
SessionStatus::Canceled
);
clear_session_env();
}
#[test]
fn register_run_is_silent_when_wave_is_unknown() {
let _guard = super::test_env_lock();
clear_session_env();
let tmp = tempfile::tempdir().expect("tempdir");
let db = tmp.path().join("lfd.db");
seed_registry(&db);
std::env::set_var(super::WAVE_ID_ENV, LfdId::new().to_string());
assert!(register_run_in(None, "design", "lf", &["lf".to_string()]).is_none());
clear_session_env();
}
}