use std::io::Read;
use std::process::ExitCode;
use std::sync::Arc;
use locode_core::{
CacheHint, EngineConfig, EventSink, FnSink, Host, HostConfig, InstructionsConfig, PackContext,
PathPolicy, ProviderInit, ProviderRegistry, SamplingArgs, Session, SkillsConfig,
};
use crate::cli::{Cli, OutputFormat};
use crate::output;
pub struct PreRunError(pub String);
impl<E: std::fmt::Display> From<E> for PreRunError {
fn from(e: E) -> Self {
PreRunError(e.to_string())
}
}
pub async fn run(cli: Cli, providers: &ProviderRegistry) -> Result<ExitCode, PreRunError> {
#[cfg(unix)]
let cancel_slot = crate::signal::install_sigterm();
let prompt = resolve_prompt(cli.prompt.as_deref())?;
let cwd = match &cli.cwd {
Some(dir) => dir.clone(),
None => std::env::current_dir()?,
};
let cwd = std::fs::canonicalize(&cwd)
.map_err(|e| PreRunError(format!("--cwd {}: {e}", cwd.display())))?;
let add_dirs = canonicalize_add_dirs(&cli.add_dir)?;
let mut host_config = host_config_for(&cwd, &add_dirs);
if !cli.restricted {
host_config.path_policy = PathPolicy::Unrestricted;
}
let host = Arc::new(Host::new(host_config)?);
output::warning_line(if cli.restricted {
output::RESTRICTED_MODE_NOTICE
} else {
output::UNRESTRICTED_MODE_NOTICE
});
let settings_load = load_settings_reporting(&cwd, cli.settings.as_deref());
let settings = settings_load.settings;
let effort = resolve_effort_reporting(cli.effort, settings.effort.as_deref());
let extends_dirs = settings_load.extends_dirs;
let identity = resolve_identity(&cli, &cwd, &settings)?;
let pack = locode_core::resolve(&identity.harness)?;
let registry = pack.build_registry(&host);
let session_id = identity.session_id.clone();
let built = providers
.build(
&identity.api_schema,
&ProviderInit {
session_id: session_id.clone(),
model: identity.model_override.clone(),
},
)
.map_err(|e| PreRunError(e.to_string()))?;
let (provider, model) = (built.provider, built.model);
enforce_wire_requirement(pack, provider.api_schema())?;
let pack_ctx = PackContext {
cwd: cwd.clone(),
os: std::env::consts::OS.to_string(),
shell: std::env::var("SHELL").unwrap_or_else(|_| "/bin/sh".to_string()),
date: chrono::Local::now().format("%Y-%m-%d").to_string(),
headless: true,
is_git_repo: detect_git_repo(&cwd),
model: Some(model.clone()),
os_version: os_version(),
timezone: timezone(),
strip_identity: cli.strip_identity,
};
let preamble = match &identity.resumed {
Some(resumed) => resumed.history.clone(),
None => pack.preamble(&pack_ctx),
};
let user_prompt = pack.shape_user_prompt(&prompt);
let config = EngineConfig {
session_id,
harness: pack.name().to_string(),
api_schema: provider.api_schema().to_string(),
model,
cwd: cwd.clone(),
workspace_root: cwd,
max_turns: cli.max_turns,
sampling_args: SamplingArgs {
reasoning_effort: effort.map(Into::into),
..SamplingArgs::default()
},
cache_hint: CacheHint::Standard,
streaming: cli.stream,
instructions: InstructionsConfig {
enabled: !cli.no_project_instructions,
root_stop_pattern: settings.root_stop_pattern.clone(),
extends_dirs: extends_dirs.clone(),
extra_roots: add_dirs.clone(),
..InstructionsConfig::default()
},
skills: SkillsConfig {
extends_dirs,
extra: settings.skills_extra.clone(),
extra_roots: add_dirs.clone(),
..SkillsConfig::enabled()
},
..EngineConfig::default()
};
let mut trace = build_trace_writer(&cli, &identity, &config.cwd);
let sink = make_sink(cli.output_format, trace.take());
let mut session = Session::new(provider, registry, preamble, config, sink);
#[cfg(unix)]
crate::signal::arm(&cancel_slot, session.cancel_handle());
let report = session.run_text(user_prompt).await;
match cli.output_format {
OutputFormat::Json => output::write_json_line(&report),
OutputFormat::Text => output::write_text(report.final_message.as_deref().unwrap_or("")),
OutputFormat::StreamJson => {} }
Ok(output::exit_code(report.status))
}
struct ResumedSession {
path: std::path::PathBuf,
history: Vec<locode_core::Message>,
}
struct RunIdentity {
harness: String,
api_schema: String,
model_override: Option<String>,
session_id: String,
resumed: Option<ResumedSession>,
}
fn resolve_identity(
cli: &Cli,
cwd: &std::path::Path,
settings: &locode_core::Settings,
) -> Result<RunIdentity, PreRunError> {
let recovered = if cli.continue_session || cli.resume.is_some() {
let home = locode_core::locode_home().map_err(PreRunError)?;
let root = home.join("sessions");
let path = if let Some(id) = &cli.resume {
locode_core::find_rollout_by_id(&root, cwd, id)
.ok_or_else(|| PreRunError(format!("--resume: no session `{id}` found")))?
} else {
locode_core::find_latest_rollout(&root, cwd).ok_or_else(|| {
PreRunError(format!(
"--continue: no session found for {}",
cwd.display()
))
})?
};
let contents = locode_core::read_rollout(&path).map_err(PreRunError)?;
Some((path, contents))
} else {
None
};
if let Some((path, contents)) = recovered {
let meta = contents.meta;
if let Some(flag) = cli.harness
&& flag.as_str() != meta.harness
{
return Err(PreRunError(format!(
"--harness {} conflicts with the resumed session's harness `{}`",
flag.as_str(),
meta.harness
)));
}
if let Some(flag) = &cli.api_schema
&& flag != &meta.api_schema
{
return Err(PreRunError(format!(
"--api-schema {flag} conflicts with the resumed session's wire `{}` \
(a session never crosses wires)",
meta.api_schema
)));
}
return Ok(RunIdentity {
harness: meta.harness.clone(),
api_schema: meta.api_schema.clone(),
model_override: cli.model.clone().or_else(|| settings.model.clone()),
session_id: meta.session_id.clone(),
resumed: Some(ResumedSession {
path,
history: contents.history,
}),
});
}
Ok(RunIdentity {
harness: match cli.harness {
Some(harness) => harness.as_str().to_string(),
None => settings
.harness
.clone()
.unwrap_or_else(|| "claude".to_string()),
},
api_schema: cli
.api_schema
.clone()
.or_else(|| settings.api_schema.clone())
.unwrap_or_else(|| "anthropic".to_string()),
model_override: cli.model.clone().or_else(|| settings.model.clone()),
session_id: new_session_id(),
resumed: None,
})
}
fn build_trace_writer(
cli: &Cli,
identity: &RunIdentity,
cwd: &std::path::Path,
) -> Option<locode_core::TraceWriter> {
let root = locode_core::locode_home()
.ok()
.filter(|_| !cli.no_session_persistence)?
.join("sessions");
match &identity.resumed {
Some(resumed) => locode_core::TraceWriter::resume(resumed.path.clone(), root)
.map_err(|e| output::warning_line(&format!("trace: {e}; tracing disabled")))
.ok(),
None => Some(locode_core::TraceWriter::new(
root,
locode_core::TraceExtras {
cli_version: env!("CARGO_PKG_VERSION").to_string(),
git: git_meta(cwd),
..Default::default()
},
)),
}
}
fn make_sink(
output_format: OutputFormat,
mut trace: Option<locode_core::TraceWriter>,
) -> Box<dyn EventSink> {
let stream = matches!(output_format, OutputFormat::StreamJson);
Box::new(FnSink(move |event| {
if let Some(writer) = trace.as_mut() {
writer.on_event(&event);
if let Some(e) = writer.take_error() {
output::warning_line(&format!("trace: {e}; tracing disabled"));
}
}
if stream && in_whole_message_trace(&event) {
output::write_json_line(&event);
}
}))
}
fn enforce_wire_requirement(pack: &dyn locode_core::Pack, schema: &str) -> Result<(), PreRunError> {
if schema != "mock"
&& let Some(required) = pack.required_api_schemas()
&& !required.contains(&schema)
{
return Err(PreRunError(format!(
"harness `{}` requires one of these wires: {}; got `--api-schema {}`",
pack.name(),
required.join(", "),
schema,
)));
}
Ok(())
}
fn resolve_prompt(arg: Option<&str>) -> Result<String, PreRunError> {
let prompt = match arg {
Some("-") | None => {
let mut buf = String::new();
std::io::stdin().read_to_string(&mut buf)?;
buf
}
Some(text) => text.to_string(),
};
let prompt = prompt.trim().to_string();
if prompt.is_empty() {
return Err(PreRunError(
"no prompt: pass it as the positional argument or on stdin".to_string(),
));
}
Ok(prompt)
}
fn git_meta(cwd: &std::path::Path) -> Option<locode_core::GitMeta> {
if !detect_git_repo(cwd) {
return None;
}
let run = |args: &[&str]| -> Option<String> {
let out = std::process::Command::new("git")
.arg("-C")
.arg(cwd)
.args(args)
.output()
.ok()?;
if !out.status.success() {
return None;
}
let s = String::from_utf8_lossy(&out.stdout).trim().to_string();
(!s.is_empty()).then_some(s)
};
Some(locode_core::GitMeta {
root: run(&["rev-parse", "--show-toplevel"]).map(std::path::PathBuf::from),
branch: run(&["rev-parse", "--abbrev-ref", "HEAD"]),
head: run(&["rev-parse", "HEAD"]),
remote: run(&["remote", "get-url", "origin"]),
})
}
fn detect_git_repo(cwd: &std::path::Path) -> bool {
cwd.ancestors().any(|dir| dir.join(".git").exists())
}
fn os_version() -> Option<String> {
#[cfg(unix)]
{
let out = std::process::Command::new("uname")
.args(["-s", "-r"])
.output()
.ok()?;
if !out.status.success() {
return None;
}
let s = String::from_utf8_lossy(&out.stdout).trim().to_string();
(!s.is_empty()).then_some(s)
}
#[cfg(not(unix))]
{
None
}
}
fn timezone() -> Option<String> {
if let Ok(tz) = std::env::var("TZ") {
let tz = tz.trim();
if !tz.is_empty() {
return Some(tz.to_string());
}
}
#[cfg(unix)]
{
let target = std::fs::read_link("/etc/localtime").ok()?;
let s = target.to_string_lossy();
s.split_once("zoneinfo/")
.map(|(_, name)| name.to_string())
.filter(|name| !name.is_empty())
}
#[cfg(not(unix))]
{
None
}
}
fn new_session_id() -> String {
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map_or(0, |d| d.as_millis());
format!("sess-{now}-{}", std::process::id())
}
fn in_whole_message_trace(event: &locode_core::Event) -> bool {
!matches!(
event,
locode_core::Event::MessageDelta { .. } | locode_core::Event::MessageDeltaReset { .. }
)
}
pub fn canonicalize_add_dirs(
dirs: &[std::path::PathBuf],
) -> Result<Vec<std::path::PathBuf>, PreRunError> {
dirs.iter()
.map(|dir| {
std::fs::canonicalize(dir)
.map_err(|e| PreRunError(format!("--add-dir {}: {e}", dir.display())))
})
.collect()
}
fn load_settings_reporting(
cwd: &std::path::Path,
inline: Option<&str>,
) -> locode_core::SettingsLoad {
let load = locode_core::load_settings(cwd, inline);
for warning in &load.warnings {
output::warning_line(warning);
}
load
}
fn host_config_for(cwd: &std::path::Path, add_dirs: &[std::path::PathBuf]) -> HostConfig {
let mut config = HostConfig::new(cwd);
config.extra_roots = add_dirs.to_vec();
config
}
fn resolve_effort_reporting(
flag: Option<crate::EffortArg>,
setting: Option<&str>,
) -> Option<locode_core::Effort> {
let mut warnings = Vec::new();
let effort = resolve_effort(flag, setting, &mut warnings);
for warning in &warnings {
output::warning_line(warning);
}
effort
}
#[must_use]
pub fn resolve_effort(
flag: Option<crate::EffortArg>,
setting: Option<&str>,
warnings: &mut Vec<String>,
) -> Option<locode_core::Effort> {
if let Some(flag) = flag {
return Some(flag.into());
}
let raw = setting?;
if let Some(effort) = locode_core::Effort::parse(raw) {
return Some(effort);
}
warnings.push(format!(
"settings: unknown effort {raw:?} — using the API default (expected one of {})",
locode_core::Effort::ALL
.iter()
.map(|e| e.as_str())
.collect::<Vec<_>>()
.join(", ")
));
None
}
#[cfg(test)]
mod tests {
use super::{enforce_wire_requirement, in_whole_message_trace};
use locode_core::{Event, Message, Role};
#[test]
fn codex_rejects_a_non_responses_wire() {
let codex = locode_core::resolve("codex").unwrap();
let err = enforce_wire_requirement(codex, "anthropic").expect_err("mismatch");
assert!(err.0.contains("codex"), "{}", err.0);
assert!(err.0.contains("openai-responses"), "{}", err.0);
assert!(err.0.contains("anthropic"), "{}", err.0);
assert!(enforce_wire_requirement(codex, "openai-responses").is_ok());
assert!(enforce_wire_requirement(codex, "mock").is_ok());
}
#[test]
fn wire_agnostic_packs_accept_any_wire() {
let grok = locode_core::resolve("grok").unwrap();
assert!(enforce_wire_requirement(grok, "anthropic").is_ok());
assert!(enforce_wire_requirement(grok, "openai-responses").is_ok());
}
#[test]
fn stream_json_trace_drops_message_deltas_keeps_whole_messages() {
assert!(!in_whole_message_trace(&Event::MessageDelta {
text: "tok".into()
}));
assert!(in_whole_message_trace(&Event::Message {
message: Message {
role: Role::Assistant,
content: vec![],
},
}));
assert!(in_whole_message_trace(&Event::Error {
message: "e".into()
}));
}
}