use std::sync::Arc;
use anyhow::{bail, Context};
use crate::lf::SessionCommand;
use crate::ops::human_session::{OpenMode, SessionKind, SessionState};
use crate::session_record::SessionTitleSource;
use crate::store::{open_store, storage_config_from_env, Store};
pub fn run(command: &SessionCommand) -> anyhow::Result<()> {
let runtime = tokio::runtime::Runtime::new()?;
let worktree = match command {
SessionCommand::Open { json: false, .. }
| SessionCommand::ServeAsk { .. }
| SessionCommand::ServeFlow { .. } => Some(crate::repo::working_directory()?),
SessionCommand::Complete { id } => {
let store = runtime.block_on(open_shared_store())?;
runtime.block_on(crate::ops::human_session::completion_worktree(&store, id))?
}
_ => None,
};
let Some(worktree) = worktree else {
return runtime.block_on(run_async(command));
};
let argv = std::env::args().collect::<Vec<_>>();
crate::journal::with_runtime(&worktree, &argv, || runtime.block_on(run_async(command)))
}
async fn run_async(command: &SessionCommand) -> anyhow::Result<()> {
match command {
SessionCommand::History {
id,
json,
after,
limit,
} => {
let store = open_shared_store().await?;
if store.session(id).await?.is_none() {
bail!("Session {id} was not found");
}
let events = store.sqlite.session_history(id, *after, *limit)?;
if *json {
println!("{}", serde_json::to_string(&events)?);
} else {
for event in events {
println!(
"{} {} {} {}",
event.seq,
event.provider_turn.as_deref().unwrap_or("unknown-turn"),
event.kind.as_str(),
event.payload
);
}
}
Ok(())
}
SessionCommand::List {
json,
all,
interactive,
history,
limit,
offset,
page,
after,
task,
search,
} => {
let store = open_shared_store().await?;
list(
&store,
*json,
&crate::session::SessionFilter {
repo: if *all {
None
} else {
crate::repository::CanonicalRepo::current()?.map(|repo| repo.to_string())
},
task: task.clone(),
search: search.clone(),
interactive: interactive.interactive(),
history: *history,
limit: *limit,
offset: *offset,
after: page.then(|| after.clone().unwrap_or_default()),
},
)
.await
}
SessionCommand::Open {
id,
json,
replace,
try_open,
} => {
let mode = if *replace {
OpenMode::Replace
} else if *try_open {
OpenMode::Try
} else {
OpenMode::Refuse
};
open(id, *json, mode).await
}
SessionCommand::Complete { id } => complete(id).await,
SessionCommand::Rename {
id,
name,
suggest,
json,
} => rename(id, name, *suggest, *json).await,
SessionCommand::Bind {
id,
task,
dry_run,
json,
} => {
let store = open_shared_store().await?;
if *dry_run {
let preview = crate::ops::human_session::preview_binding(&store, id, task).await?;
if *json {
println!("{}", serde_json::to_string_pretty(&preview)?);
} else {
println!(
"{} → {} ({}) [{}]. Not assigned.",
preview.session_id, preview.identifier, preview.title, preview.task_id
);
}
return Ok(());
}
let session = crate::ops::human_session::bind(&store, id, task).await?;
if *json {
println!("{}", serde_json::to_string_pretty(&session)?);
} else {
println!(
"Session {} belongs to {}.",
session.id,
session.work_path.as_deref().unwrap_or(task)
);
}
Ok(())
}
SessionCommand::Ready { summary } => {
let text = required_text(summary, "ready summary")?;
let store = open_shared_store().await?;
crate::ops::human_session::mark_ready(&store, &text).await?;
println!("Session is ready for your review.");
Ok(())
}
SessionCommand::ServeFlow {
task_id,
invocation_id,
flow,
node_id,
skill,
iteration,
} => {
let store = open_shared_store().await?;
crate::ops::human_session::serve_flow(
store,
task_id.clone(),
invocation_id.clone(),
flow.clone(),
node_id.clone(),
skill.clone(),
*iteration,
)
.await
}
SessionCommand::ServeAsk { run_id } => {
let store = open_shared_store().await?;
crate::ops::human_session::serve_ask(&store, run_id).await
}
SessionCommand::StopRun { run_id } => {
crate::ops::human_session::stop_session_client(run_id)
}
}
}
async fn list(
store: &Arc<Store>,
json: bool,
filter: &crate::session::SessionFilter,
) -> anyhow::Result<()> {
if filter.after.is_some() {
anyhow::ensure!(
filter.limit > 0,
"paged Session inventory requires a positive limit"
);
let mut selection = filter.clone();
selection.limit = filter
.limit
.checked_add(1)
.context("Session page limit is too large")?;
let mut entries = crate::ops::human_session::list(store, &selection).await?;
let next = if entries.len() > filter.limit {
entries.pop();
entries.last().map(|session| session.id.clone())
} else {
None
};
println!(
"{}",
serde_json::to_string_pretty(&crate::ops::human_session::SessionPage {
entries,
next
})?
);
return Ok(());
}
let sessions = crate::ops::human_session::list(store, filter).await?;
if json {
println!("{}", serde_json::to_string_pretty(&sessions)?);
} else if sessions.is_empty() {
println!("No Sessions.");
} else {
for session in sessions {
println!(
"{} {:<7} {} {}",
session.id,
match session.state {
SessionState::Unknown => "unknown",
SessionState::Waiting => "waiting",
SessionState::Active => "active",
SessionState::Ready => "ready",
SessionState::Closed => "closed",
},
session.work_path.as_deref().unwrap_or("Repository"),
session.title
);
for action in session.actions {
println!(
" {} — {}",
action.label,
action.unavailable_reason.as_deref().unwrap_or(&action.help)
);
}
}
}
Ok(())
}
async fn open(id: &str, json: bool, mode: OpenMode) -> anyhow::Result<()> {
let store = open_shared_store().await?;
let session = crate::ops::human_session::open(&store, id, mode, !json).await?;
if json {
println!("{}", serde_json::to_string_pretty(&session)?);
}
Ok(())
}
async fn complete(id: &str) -> anyhow::Result<()> {
let store = open_shared_store().await?;
let session = crate::ops::human_session::complete(&store, id).await?;
match session.kind {
SessionKind::Conversation => println!(
"Session {} completed; its provider history remains resumable.",
session.id
),
SessionKind::Ask => println!(
"Ask session completed: {}",
session
.ready_summary
.expect("completed Ask Session has a ready summary")
),
SessionKind::Flow => println!("Review completed; feedback returned to the Flow."),
}
Ok(())
}
async fn rename(id: &str, name: &[String], suggest: bool, json: bool) -> anyhow::Result<()> {
let requested = required_text(name, "Session name")?;
let source = if suggest {
SessionTitleSource::Generated
} else {
SessionTitleSource::Human
};
let store = open_shared_store().await?;
let session = crate::ops::human_session::rename(&store, id, &requested, source).await?;
if json {
println!("{}", serde_json::to_string_pretty(&session)?);
} else if suggest && session.title_source == SessionTitleSource::Human {
println!(
"Session {} keeps its human-assigned name {:?}.",
session.id, session.title
);
} else {
println!("Session {} is named {:?}.", session.id, session.title);
}
Ok(())
}
fn required_text(args: &[String], label: &str) -> anyhow::Result<String> {
let text = args.join(" ").trim().to_string();
if text.is_empty() {
bail!("{label} cannot be empty");
}
Ok(text)
}
async fn open_shared_store() -> anyhow::Result<Arc<Store>> {
let config = storage_config_from_env().context("resolve the shared Loopflow store")?;
Ok(Arc::new(
open_store(&config)
.await
.context("open the shared Loopflow store")?,
))
}