use crate::id::WaveId;
use crate::work::wave::{Wave, WaveLocator};
use std::collections::{HashMap, HashSet};
use std::path::{Path, PathBuf};
use std::sync::{Mutex, OnceLock};
pub const WAVE_ID_ENV: &str = "LF_WAVE_ID";
pub fn resolve_ambient_wave_name() -> Option<String> {
let repo = crate::repo::find_repo_root().ok();
resolve_managed_wave_sync(repo.as_deref(), None)
.ok()
.map(|wave| wave.name().to_string())
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct RunAttribution {
pub wave: Option<String>,
pub failure: Option<String>,
}
pub fn run_attribution(repo: Option<&Path>) -> RunAttribution {
match resolve_managed_wave_sync(repo, None) {
Ok(wave) => RunAttribution {
wave: Some(wave.name().to_string()),
failure: None,
},
Err(WaveResolveError::NoContext) => RunAttribution {
wave: None,
failure: None,
},
Err(error) => RunAttribution {
wave: None,
failure: Some(attribution_failure_text(&error)),
},
}
}
fn attribution_failure_text(error: &WaveResolveError) -> String {
match error {
WaveResolveError::Registry(_) => format!("{error}; pass --wave <name> to recover"),
_ => error.to_string(),
}
}
pub fn resolve_explicit_wave(name: &str) -> anyhow::Result<Wave> {
let repo = crate::repo::find_repo_root().ok();
resolve_managed_wave_sync(repo.as_deref(), Some(name)).map_err(anyhow::Error::from)
}
#[derive(Debug, thiserror::Error, PartialEq, Eq)]
pub enum WaveResolveError {
#[error("no wave in context; pass --wave <name>")]
NoContext,
#[error(
"ambient wave '{0}' (LF_WAVE_ID) is not in this machine's registry; \
the context is stale — pass --wave <name>"
)]
StaleIdentity(String),
#[error("--wave requires a non-empty wave name")]
EmptyExplicit,
#[error(
"wave '{0}' is not registered on this machine; \
run `lf ls` to list known waves, or pass --wave <known-name>"
)]
UnknownExplicit(String),
#[error("wave '{slug}' is ambiguous; it belongs to: {repositories}")]
AmbiguousWave { slug: String, repositories: String },
#[error("wave {wave_id} belongs to {actual}, not invoking repository {expected}")]
RepositoryMismatch {
wave_id: WaveId,
expected: String,
actual: String,
},
#[error("failed to read wave registry: {0}")]
Registry(String),
}
async fn resolve_slug(
store: &crate::store::Store,
repo: Option<&Path>,
slug: &str,
) -> Result<Wave, WaveResolveError> {
if let Some(repo) = repo {
let locator = WaveLocator::discover(repo, slug)
.map_err(|error| WaveResolveError::Registry(error.to_string()))?;
return store
.get_wave_at(&locator)
.await
.map_err(|error| WaveResolveError::Registry(error.to_string()))?
.ok_or_else(|| WaveResolveError::UnknownExplicit(slug.to_string()));
}
let waves = store
.find_waves_by_slug(slug)
.await
.map_err(|error| WaveResolveError::Registry(error.to_string()))?;
match waves.as_slice() {
[wave] => Ok(wave.clone()),
[] => Err(WaveResolveError::UnknownExplicit(slug.to_string())),
_ => Err(WaveResolveError::AmbiguousWave {
slug: slug.to_string(),
repositories: waves
.iter()
.map(|wave| wave.repo())
.collect::<Vec<_>>()
.join(", "),
}),
}
}
pub async fn resolve_managed_wave(
store: Option<&crate::store::Store>,
repo: Option<&Path>,
explicit: Option<&str>,
env_wave_id: Option<&str>,
) -> Result<Wave, WaveResolveError> {
if let Some(raw) = explicit {
if let Ok(id) = raw.trim().parse::<WaveId>() {
let store = store.ok_or_else(|| {
WaveResolveError::Registry("no wave registry on this machine".to_string())
})?;
let wave = store
.get_wave(&id)
.await
.map_err(|error| WaveResolveError::Registry(error.to_string()))?
.ok_or_else(|| WaveResolveError::UnknownExplicit(raw.to_string()))?;
if wave.is_retired() {
return Ok(wave);
}
if let Some(repo) = repo {
let locator = WaveLocator::discover(repo, wave.name())
.map_err(|error| WaveResolveError::Registry(error.to_string()))?;
let scoped = store
.get_wave_at(&locator)
.await
.map_err(|error| WaveResolveError::Registry(error.to_string()))?;
if scoped.as_ref().map(Wave::id) != Some(&id) {
return Err(WaveResolveError::RepositoryMismatch {
wave_id: id,
expected: locator.repo().to_string(),
actual: wave.repo().to_string(),
});
}
}
return Ok(wave);
}
let slug =
crate::ops::util::normalize_wave_name(raw).ok_or(WaveResolveError::EmptyExplicit)?;
let store = store.ok_or_else(|| {
WaveResolveError::Registry("no wave registry on this machine".to_string())
})?;
return resolve_slug(store, repo, &slug).await;
}
let raw = env_wave_id
.map(str::trim)
.filter(|value| !value.is_empty())
.ok_or(WaveResolveError::NoContext)?;
let store = store.ok_or_else(|| WaveResolveError::StaleIdentity(raw.to_string()))?;
if let Ok(id) = raw.parse::<WaveId>() {
let wave = store
.get_wave(&id)
.await
.map_err(|error| WaveResolveError::Registry(error.to_string()))?
.ok_or_else(|| WaveResolveError::StaleIdentity(raw.to_string()))?;
if wave.is_retired() {
return Ok(wave);
}
if let Some(repo) = repo {
let locator = WaveLocator::discover(repo, wave.name())
.map_err(|error| WaveResolveError::Registry(error.to_string()))?;
let scoped = store
.get_wave_at(&locator)
.await
.map_err(|error| WaveResolveError::Registry(error.to_string()))?;
if let Some(scoped) = scoped {
if scoped.id() == &id {
return Ok(scoped);
}
}
return Err(WaveResolveError::RepositoryMismatch {
wave_id: id,
expected: locator.repo().to_string(),
actual: wave.repo().to_string(),
});
}
return Ok(wave);
}
let slug = crate::ops::util::normalize_wave_name(raw).ok_or(WaveResolveError::NoContext)?;
resolve_slug(store, repo, &slug)
.await
.map_err(|error| match error {
WaveResolveError::UnknownExplicit(_) => {
WaveResolveError::StaleIdentity(raw.to_string())
}
other => other,
})
}
pub fn resolve_managed_wave_sync(
repo: Option<&Path>,
explicit: Option<&str>,
) -> Result<Wave, WaveResolveError> {
let repo = repo.map(Path::to_path_buf);
let explicit = explicit.map(str::to_string);
let env_wave_id = std::env::var(WAVE_ID_ENV).ok();
std::thread::spawn(move || {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.expect("current-thread runtime always builds");
runtime.block_on(async move {
let store = crate::store::open_existing_store().await;
resolve_managed_wave(
store.as_ref(),
repo.as_deref(),
explicit.as_deref(),
env_wave_id.as_deref(),
)
.await
})
})
.join()
.unwrap_or_else(|_| {
Err(WaveResolveError::Registry(
"resolver thread panicked".to_string(),
))
})
}
fn repo_origin(repo_root: &Path) -> PathBuf {
static CACHE: OnceLock<Mutex<HashMap<PathBuf, PathBuf>>> = OnceLock::new();
let cache = CACHE.get_or_init(Default::default);
let mut cache = cache
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
cache
.entry(repo_root.to_path_buf())
.or_insert_with(|| query_repo_origin(repo_root))
.clone()
}
fn query_repo_origin(repo_root: &Path) -> PathBuf {
let not_a_root = || repo_root.to_path_buf();
let Ok(output) = std::process::Command::new("git")
.arg("-C")
.arg(repo_root)
.args([
"rev-parse",
"--path-format=absolute",
"--show-toplevel",
"--git-common-dir",
])
.output()
else {
return not_a_root();
};
if !output.status.success() {
return not_a_root();
}
let stdout = String::from_utf8_lossy(&output.stdout);
let mut lines = stdout.lines().map(str::trim);
let toplevel = PathBuf::from(lines.next().unwrap_or_default());
let common_dir = PathBuf::from(lines.next().unwrap_or_default());
let toplevel = toplevel.canonicalize().unwrap_or(toplevel);
let root = repo_root
.canonicalize()
.unwrap_or_else(|_| repo_root.to_path_buf());
if toplevel != root {
return not_a_root();
}
common_dir
.parent()
.map(Path::to_path_buf)
.unwrap_or_else(|| repo_root.to_path_buf())
}
pub fn wave_origin(repo_root: &Path) -> PathBuf {
repo_origin(repo_root)
}
pub fn gather_wave_memory(repo_root: &Path, wave: &str) -> Option<String> {
let origin = wave_origin(repo_root);
gather_wave_memory_from(&origin, &origin, wave)
}
pub(crate) fn gather_wave_memory_from(
origin: &Path,
content_repo: &Path,
wave: &str,
) -> Option<String> {
let chain = memory_wave_chain(origin, wave).unwrap_or_else(|| vec![wave.to_string()]);
gather_memory_chain(content_repo, &chain)
}
fn memory_wave_chain(origin: &Path, wave: &str) -> Option<Vec<String>> {
let origin = origin.to_path_buf();
let wave = wave.to_string();
std::thread::spawn(move || {
let rt = tokio::runtime::Runtime::new().ok()?;
rt.block_on(async {
let store = crate::store::open_existing_store().await?;
memory_wave_chain_from_store(&store, &origin, &wave).await
})
})
.join()
.ok()
.flatten()
}
async fn memory_wave_chain_from_store(
store: &crate::store::Store,
origin: &Path,
wave: &str,
) -> Option<Vec<String>> {
let locator = WaveLocator::discover(origin, wave).ok()?;
let mut current = store.get_wave_at(&locator).await.ok().flatten()?;
let mut seen = HashSet::new();
let mut chain = Vec::new();
loop {
if !seen.insert(current.id().clone()) {
tracing::warn!(
wave,
"cycle in parent_wave_id; using the acyclic memory prefix"
);
break;
}
chain.push(current.name().to_string());
let Some(parent) = current.parent_wave_id() else {
break;
};
current = match store.get_wave(parent).await.ok().flatten() {
Some(parent) => parent,
None => {
tracing::warn!(wave, parent = %parent, "missing parent wave in memory scope");
break;
}
};
}
chain.reverse();
Some(chain)
}
fn gather_memory_chain(origin: &Path, chain: &[String]) -> Option<String> {
let leaf = chain.last()?;
let scoped = chain
.iter()
.filter_map(|wave| {
let base = super::memory::Memory::for_wave(origin, wave).read();
let memory = render_wave_memory(&base)?;
if chain.len() == 1 {
return Some(memory);
}
let ownership = if wave == leaf {
"owned by"
} else {
"inherited from"
};
Some(format!("## Memory {ownership} {wave}\n\n{memory}"))
})
.collect::<Vec<_>>();
(!scoped.is_empty()).then(|| scoped.join("\n\n"))
}
fn render_wave_memory(base: &str) -> Option<String> {
let base = base.trim();
(!base.is_empty()).then(|| base.to_string())
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn run_attribution_classifies_absent_context_and_hand_set_names() {
let ledger = crate::journal::TestLedgerGuard::new();
let previous = std::env::var(WAVE_ID_ENV).ok();
let repo = crate::repo::find_repo_root().unwrap();
let runtime = tokio::runtime::Runtime::new().unwrap();
runtime.block_on(async {
let store = crate::store::open_ephemeral_store(&crate::store::StorageConfig::sqlite(
ledger.home().join("loopflow.db"),
))
.await
.unwrap();
let locator = WaveLocator::discover(&repo, "product").unwrap();
let wave = Wave::new(
WaveId::new(),
"product".to_string(),
locator.repo().to_string(),
);
store.create_wave(&wave).await.unwrap();
});
std::env::remove_var(WAVE_ID_ENV);
let absent = run_attribution(Some(&repo));
assert_eq!(absent.wave, None);
assert_eq!(absent.failure, None);
std::env::set_var(WAVE_ID_ENV, "product");
let named = run_attribution(Some(&repo));
assert_eq!(named.wave.as_deref(), Some("product"));
assert_eq!(named.failure, None);
std::env::set_var(WAVE_ID_ENV, "ghost");
let unregistered = run_attribution(Some(&repo));
assert_eq!(unregistered.wave, None);
assert!(unregistered
.failure
.is_some_and(|failure| failure.contains("context is stale")));
match previous {
Some(value) => std::env::set_var(WAVE_ID_ENV, value),
None => std::env::remove_var(WAVE_ID_ENV),
}
}
#[test]
fn wave_origin_of_a_worktree_is_the_main_checkout() {
let repo = loopflow_test_support::TestRepo::new();
let worktree = repo.create_named_worktree("origin-check");
let origin = wave_origin(&worktree);
assert_eq!(
origin.canonicalize().unwrap(),
repo.path().canonicalize().unwrap()
);
let tmp = tempfile::tempdir().expect("tempdir");
assert_eq!(wave_origin(tmp.path()), tmp.path());
}
#[test]
fn render_wave_memory_uses_only_the_file() {
assert_eq!(
render_wave_memory("# Memory\n\ncompiled base\n").as_deref(),
Some("# Memory\n\ncompiled base")
);
assert!(render_wave_memory(" \n").is_none());
}
#[test]
fn gather_wave_memory_reads_the_file_without_a_server() {
let tmp = tempfile::tempdir().expect("tempdir");
std::fs::create_dir_all(tmp.path().join("wave/goals")).unwrap();
std::fs::write(
tmp.path().join("wave/goals/MEMORY.md"),
"# Goals\n\ncompiled\n",
)
.unwrap();
let memory = gather_wave_memory(tmp.path(), "goals").expect("memory");
assert_eq!(memory, "# Goals\n\ncompiled");
}
#[tokio::test]
async fn child_memory_walks_parent_scope() {
let tmp = tempfile::tempdir().expect("tempdir");
let store = crate::store::open_ephemeral_store(&crate::store::StorageConfig::sqlite(
tmp.path().join("loopflow.db"),
))
.await
.unwrap();
let parent = crate::work::wave::Wave::new(
WaveId::new(),
"platform".into(),
tmp.path().display().to_string(),
);
let child = crate::work::wave::Wave::new(
WaveId::new(),
"release".into(),
tmp.path().display().to_string(),
)
.with_parent(parent.id().clone());
store.create_wave(&parent).await.unwrap();
store.create_wave(&child).await.unwrap();
std::fs::create_dir_all(tmp.path().join("wave/platform")).unwrap();
std::fs::create_dir_all(tmp.path().join("wave/release")).unwrap();
std::fs::write(
tmp.path().join("wave/platform/MEMORY.md"),
"Parent constraint.",
)
.unwrap();
std::fs::write(tmp.path().join("wave/release/MEMORY.md"), "Child decision.").unwrap();
let chain = memory_wave_chain_from_store(&store, tmp.path(), "release")
.await
.expect("scope resolves");
assert_eq!(chain, ["platform", "release"]);
let memory = gather_memory_chain(tmp.path(), &chain).expect("memory renders");
assert!(memory.contains("## Memory inherited from platform\n\nParent constraint."));
assert!(memory.contains("## Memory owned by release\n\nChild decision."));
assert!(
memory.find("Parent constraint.").unwrap() < memory.find("Child decision.").unwrap()
);
}
}