use std::path::Path;
use anyhow::{anyhow, Context, Result};
use super::{
capture_is_prepared, conversation_background_name, conversation_exec_is_running,
lock_session_exec, publish_prepared_input, session_not_found, start_durable_session, surface,
NativeSession, SessionRecord,
};
use crate::session::{AgentSession, PrimaryScope, SessionKind, TitleSource, WorkSource};
use crate::store::SharedStore;
pub(crate) async fn ensure(
store: &SharedStore,
repo: &Path,
wave: Option<&str>,
) -> Result<SessionRecord> {
let (scope, session) = match wave {
Some(wave) => {
let binding =
crate::ops::resolve_work_binding(store, repo, &format!("wave:{wave}")).await?;
(
PrimaryScope::Wave(binding.wave_id.clone()),
wave_session(&binding),
)
}
None => {
let repo = crate::repository::CanonicalRepo::discover(repo)?;
let session = repository_session(&repo);
(PrimaryScope::Repository(repo), session)
}
};
let _lock = lock_scope(&scope).await?;
let session = store.ensure_primary_session(&scope, None, session).await?;
start(store, session).await
}
pub(crate) async fn replace(store: &SharedStore, id: &str) -> Result<SessionRecord> {
let scope = store
.sqlite
.primary_scope(id)?
.ok_or_else(|| anyhow!("Session {id} is not a primary Session"))?;
let _lock = lock_scope(&scope).await?;
let previous = store
.session(id)
.await?
.ok_or_else(|| session_not_found(id))?;
if previous.completed_at.is_none() {
let lock_id = previous.id.clone();
let _launch = tokio::task::spawn_blocking(move || lock_session_exec(&lock_id)).await??;
stop_client(&previous)
.with_context(|| format!("stop primary Session {id}; it remains primary"))?;
}
let successor = match &scope {
PrimaryScope::Repository(repo) => repository_session(repo),
PrimaryScope::Wave(wave) => wave_session(
&crate::ops::resolve_work_binding(store, &previous.cwd, &format!("wave:{wave}"))
.await?,
),
};
let session = store
.ensure_primary_session(&scope, Some(id), successor)
.await?;
start(store, session).await
}
fn wave_session(binding: &crate::ops::WorkBinding) -> AgentSession {
AgentSession {
wave_id: Some(binding.wave_id.clone()),
work_source: Some(WorkSource::Declared),
..conversation(
&binding.cwd,
binding.agent.as_deref(),
"wave/session",
binding.wave_name.clone(),
)
}
}
fn repository_session(repo: &crate::repository::CanonicalRepo) -> AgentSession {
let title = repo
.as_path()
.file_name()
.map(|name| name.to_string_lossy().into_owned())
.unwrap_or_else(|| repo.to_string());
AgentSession {
repo: Some(repo.to_string()),
..conversation(repo.as_path(), None, "repo/session", title)
}
}
fn conversation(cwd: &Path, agent: Option<&str>, skill: &str, title: String) -> AgentSession {
let agent = crate::ops::task::resolve_task_agent(cwd, agent, None);
let (provider, model) = crate::engine::config::parse_agent(&agent);
AgentSession {
captured: None,
id: uuid::Uuid::new_v4().simple().to_string(),
artifact_key: crate::session_record::new_artifact_key(),
caller_artifact_key: None,
input_published: false,
cwd: cwd.to_path_buf(),
skill: Some(skill.into()),
provider: Some(provider),
model,
node: None,
iterations: None,
task_id: None,
wave_id: None,
flow_session_id: None,
work_source: None,
bound_at: None,
kind: SessionKind::Conversation,
interactive: true,
repo: None,
title,
title_source: TitleSource::Generated,
request: None,
ready_summary: None,
completed_at: None,
created_at: crate::store::rows::now_unix(),
}
}
async fn start(store: &SharedStore, session: AgentSession) -> Result<SessionRecord> {
let session = if session.input_published {
session
} else {
publish_prepared_input(
store,
&session,
crate::session_record::SessionFlowMembership::Independent,
)?
};
if capture_is_prepared(&session.artifact_key)?
&& !conversation_exec_is_running(&session.id).await?
{
let lf = crate::engine::process::resolve_pinned_lf_binary()?;
let argv = vec![
lf.to_string_lossy().to_string(),
"session".to_string(),
"serve-conversation".to_string(),
session.artifact_key.to_string(),
];
start_durable_session(
&conversation_background_name(&session.id),
&session.cwd,
&argv,
)
.await
.with_context(|| {
format!(
"launch primary Session {}; run the same command to retry",
session.id
)
})?;
}
surface(store, &session).await
}
fn stop_client(session: &AgentSession) -> Result<()> {
#[cfg(test)]
if super::action_test::stop(&session.artifact_key) {
return Ok(());
}
let native = NativeSession::of(session)?;
if !native.dir.is_dir() {
return Ok(());
}
crate::lf::commands::util::stop_provider_session(&native.dir, native.provider)
}
async fn lock_scope(scope: &PrimaryScope) -> Result<std::fs::File> {
let key = match scope {
PrimaryScope::Repository(repo) => format!("primary:repository:{repo}"),
PrimaryScope::Wave(wave) => format!("primary:wave:{wave}"),
};
tokio::task::spawn_blocking(move || lock_session_exec(&key))
.await
.context("lock primary Session scope")?
}
#[cfg(test)]
mod tests {
use super::{ensure, replace};
use crate::ops::human_session::tests::{
SessionHome, CONVERSATION_LAUNCHERS, FAILED_CONVERSATION_LAUNCHERS,
};
use crate::ops::human_session::{action_test::NativeClients, conversation_background_name};
use crate::session::SessionKind;
async fn wave(
store: &crate::store::SharedStore,
repo: &loopflow_test_support::TestRepo,
) -> crate::work::wave::Wave {
crate::work::wave::ensure_wave_row(store, repo.path(), "infrastructure")
.await
.unwrap()
}
#[test]
fn ensure_finds_the_same_wave_conversation_and_launches_once() {
let _lock = crate::journal::test_env_lock();
let home = SessionHome::new();
tokio::runtime::Runtime::new().unwrap().block_on(async {
let store = home.store().await;
let repo = loopflow_test_support::TestRepo::new();
let wave = wave(&store, &repo).await;
let first = ensure(&store, repo.path(), Some("infrastructure"))
.await
.unwrap();
let second = ensure(&store, repo.path(), Some("infrastructure"))
.await
.unwrap();
assert_eq!(first.id, second.id);
assert!(CONVERSATION_LAUNCHERS
.lock()
.unwrap()
.contains(&conversation_background_name(&first.id)));
let session = store.session(&first.id).await.unwrap().unwrap();
assert_eq!(session.kind, SessionKind::Conversation);
assert!(session.interactive && session.input_published);
assert_eq!(session.skill.as_deref(), Some("wave/session"));
assert_eq!(session.wave_id.as_ref(), Some(wave.id()));
assert_eq!(session.task_id, None);
});
}
#[test]
fn opening_waits_for_completion_under_the_launch_lock() {
let _lock = crate::journal::test_env_lock();
let home = SessionHome::new();
tokio::runtime::Runtime::new().unwrap().block_on(async {
let store = home.store().await;
let repo = loopflow_test_support::TestRepo::new();
let record = ensure(&store, repo.path(), None).await.unwrap();
let session = store.session(&record.id).await.unwrap().unwrap();
let lock = super::lock_session_exec(&session.id).unwrap();
let opening = super::super::open_waiting(&store, &session.id);
tokio::pin!(opening);
tokio::select! {
biased;
result = &mut opening => panic!("opened while locked: {result:?}"),
() = async {
store.complete_session(&session.id, session.captured).await.unwrap();
drop(lock);
} => {}
}
assert!(opening
.await
.unwrap_err()
.to_string()
.contains("already complete"));
assert_eq!(store.session_inputs(&session.id).await.unwrap().len(), 1);
});
}
#[test]
fn failed_start_keeps_the_session_for_the_next_ensure() {
let _lock = crate::journal::test_env_lock();
let home = SessionHome::new();
tokio::runtime::Runtime::new().unwrap().block_on(async {
let store = home.store().await;
let repo = loopflow_test_support::TestRepo::new();
wave(&store, &repo).await;
let admitted = ensure(&store, repo.path(), Some("infrastructure"))
.await
.unwrap();
let name = conversation_background_name(&admitted.id);
CONVERSATION_LAUNCHERS.lock().unwrap().remove(&name);
FAILED_CONVERSATION_LAUNCHERS
.lock()
.unwrap()
.insert(name.clone());
assert!(ensure(&store, repo.path(), Some("infrastructure"))
.await
.is_err());
FAILED_CONVERSATION_LAUNCHERS.lock().unwrap().clear();
let retried = ensure(&store, repo.path(), Some("infrastructure"))
.await
.unwrap();
assert_eq!(retried.id, admitted.id);
assert!(CONVERSATION_LAUNCHERS.lock().unwrap().contains(&name));
});
}
#[test]
fn replace_stops_the_predecessor_and_repeats_to_the_same_successor() {
let _lock = crate::journal::test_env_lock();
let home = SessionHome::new();
tokio::runtime::Runtime::new().unwrap().block_on(async {
let store = home.store().await;
let repo = loopflow_test_support::TestRepo::new();
wave(&store, &repo).await;
let first = ensure(&store, repo.path(), Some("infrastructure"))
.await
.unwrap();
let previous = store.session(&first.id).await.unwrap().unwrap();
let clients =
NativeClients::new(&first.id, std::slice::from_ref(&previous.artifact_key));
let successor = replace(&store, &first.id).await.unwrap();
assert_ne!(successor.id, first.id);
assert!(clients.active().is_empty());
assert!(store
.session(&first.id)
.await
.unwrap()
.unwrap()
.completed_at
.is_some());
assert_eq!(replace(&store, &first.id).await.unwrap().id, successor.id);
assert_eq!(
ensure(&store, repo.path(), Some("infrastructure"))
.await
.unwrap()
.id,
successor.id
);
});
}
#[test]
fn a_repository_without_waves_has_its_own_conversation() {
let _lock = crate::journal::test_env_lock();
let home = SessionHome::new();
tokio::runtime::Runtime::new().unwrap().block_on(async {
let store = home.store().await;
let repo = loopflow_test_support::TestRepo::new();
let first = ensure(&store, repo.path(), None).await.unwrap();
let second = ensure(&store, repo.path(), None).await.unwrap();
assert_eq!(first.id, second.id);
let session = store.session(&first.id).await.unwrap().unwrap();
assert_eq!(session.skill.as_deref(), Some("repo/session"));
assert_eq!((session.wave_id, session.task_id), (None, None));
let ambient_wave = wave(&store, &repo).await;
std::env::set_var("LF_WAVE_ID", ambient_wave.id().as_str());
assert_eq!(
ensure(&store, repo.path(), None).await.unwrap().id,
first.id
);
std::env::set_var("LF_WAVE_ID", "unknown-wave");
let wave = ensure(&store, repo.path(), Some("infrastructure"))
.await
.unwrap();
assert_ne!(wave.id, first.id);
let successor = replace(&store, &first.id).await.unwrap();
assert_ne!(successor.id, first.id);
assert_eq!(
ensure(&store, repo.path(), None).await.unwrap().id,
successor.id
);
});
}
#[test]
fn replace_refuses_an_ordinary_conversation() {
let _lock = crate::journal::test_env_lock();
let home = SessionHome::new();
tokio::runtime::Runtime::new().unwrap().block_on(async {
let store = home.store().await;
let error = replace(&store, "not-a-primary").await.unwrap_err();
assert!(error.to_string().contains("not a primary Session"));
});
}
}