pub mod claude;
mod claude_history;
mod claude_mapping;
pub mod codex;
#[cfg(unix)]
pub mod codex_connection;
mod codex_history;
mod codex_mapping;
mod common;
#[cfg(test)]
mod conformance_tests;
mod lf_tag;
pub mod opencode;
pub(crate) mod opencode_history;
mod opencode_mapping;
pub mod opencode_runtime;
pub(crate) use claude_mapping::rate_limit_signal as claude_rate_limit_signal;
pub(crate) use codex_mapping::rate_limit_signal as codex_rate_limit_signal;
pub(crate) use codex_mapping::window_name as codex_window_name;
use anyhow::Result;
use async_trait::async_trait;
use tokio::sync::mpsc;
use crate::chat::types::ConversationEvent;
use crate::engine::agent::AgentConfig;
pub(crate) fn configure_vendor_std_env(command: &mut std::process::Command) -> Result<()> {
let context = crate::engine::process::pinned_execution_context()?;
set_vendor_std_env(command, &context.lf_bin, &context.lf_home, &context.db_path)
}
pub(crate) fn configure_agent_env(command: &mut tokio::process::Command, config: &AgentConfig) {
for name in crate::engine::agent::EXECUTION_IDENTITY_ENV {
command.env_remove(name);
}
command
.envs(&config.env)
.env_remove(crate::engine::process::DISCORD_TOKEN_ENV)
.env_remove(crate::ops::git_operation::LEGACY_WORKTREE_WRITER_ID_ENV)
.env_remove("LOOPFLOW_DIRECTIVE_FILE");
if let Some(path) = &config.directive_relay {
command.env("LOOPFLOW_DIRECTIVE_FILE", path);
}
let program = command
.as_std()
.get_program()
.to_string_lossy()
.into_owned();
crate::provider_auth::apply_provider_env_to_command(&program, command.as_std_mut());
}
pub(crate) fn conversation_environment(
command: &std::process::Command,
config: &AgentConfig,
) -> std::collections::BTreeMap<String, String> {
command
.get_envs()
.filter_map(|(key, value)| {
let key = key.to_string_lossy();
let intended = config.env.contains_key(key.as_ref())
|| matches!(
key.as_ref(),
"PATH"
| "LF_BIN"
| "LF_HOME"
| "LF_DB_PATH"
| "LF_CONTROL_BIN"
| "LF_CONTROL_HOME"
| "LF_CONTROL_DB_PATH"
| "LOOPFLOW_DIRECTIVE_FILE"
);
intended
.then_some(value)
.flatten()
.map(|value| (key.into_owned(), value.to_string_lossy().into_owned()))
})
.collect()
}
fn set_vendor_std_env(
command: &mut std::process::Command,
control_bin: &std::path::Path,
control_home: &std::path::Path,
control_db: &std::path::Path,
) -> Result<()> {
command
.env(crate::store::CONTROL_BIN_ENV, control_bin)
.env(crate::store::CONTROL_HOME_ENV, control_home)
.env(crate::store::CONTROL_DB_PATH_ENV, control_db)
.env_remove("LF_BIN")
.env_remove("LF_HOME")
.env_remove("LF_DB_PATH");
command.env_remove(crate::engine::process::DISCORD_TOKEN_ENV);
if !crate::build_info::provenance().is_release() {
command
.env("LF_HOME", control_home)
.env("LF_DB_PATH", control_db);
if crate::machine_install::selection_for_current_executable()?.is_none() {
command.env("LF_BIN", control_bin);
let mut paths = vec![control_bin
.parent()
.expect("absolute lf has a parent")
.to_path_buf()];
paths.extend(std::env::split_paths(
&std::env::var_os("PATH").unwrap_or_default(),
));
command.env("PATH", std::env::join_paths(paths)?);
}
}
Ok(())
}
#[derive(Debug, Clone)]
pub struct RawProviderEvent {
pub stream: &'static str,
pub line: String,
}
#[cfg(test)]
mod environment_tests {
use std::ffi::OsString;
use std::path::Path;
use super::{configure_agent_env, set_vendor_std_env};
use crate::engine::agent::AgentConfig;
#[test]
fn conversation_tools_observe_sanitized_overrides_and_removals() {
let mut config = AgentConfig::default();
for key in [
"LF_DISCORD_TOKEN",
"LF_WORKTREE_WRITER_ID",
"LOOPFLOW_DIRECTIVE_FILE",
"LF_BIN",
"LF_HOME",
"LF_DB_PATH",
] {
config.env.insert(key.into(), "stale-fixture".into());
}
config
.env
.insert("LF_AGENT_CALLER".into(), "current-fixture".into());
let mut engine = tokio::process::Command::new("vendor");
configure_agent_env(&mut engine, &config);
set_vendor_std_env(
engine.as_std_mut(),
Path::new("/control/lf"),
Path::new("/private"),
Path::new("/private/loopflow.db"),
)
.unwrap();
engine.env("PROVIDER_ACCOUNT_FIXTURE", "not-for-tools");
let tools = super::conversation_environment(engine.as_std(), &config);
let context_check = if crate::build_info::provenance().is_release() {
"test -z \"${LF_BIN+x}${LF_HOME+x}${LF_DB_PATH+x}\""
} else {
"test \"$LF_HOME\" = /private && test \"$LF_DB_PATH\" = /private/loopflow.db"
};
let script = format!("test -z \"${{LF_DISCORD_TOKEN+x}}${{LF_WORKTREE_WRITER_ID+x}}${{LOOPFLOW_DIRECTIVE_FILE+x}}${{PROVIDER_ACCOUNT_FIXTURE+x}}\" && test \"$LF_AGENT_CALLER\" = current-fixture && test \"$LF_CONTROL_HOME\" = /private && {context_check}");
assert!(std::process::Command::new("/bin/sh")
.env_clear()
.envs(tools)
.args(["-c", &script])
.status()
.unwrap()
.success());
for key in ["LF_BIN", "LF_HOME", "LF_DB_PATH"] {
engine.env_remove(key);
}
let tools = super::conversation_environment(engine.as_std(), &config);
assert!(std::process::Command::new("/bin/sh").env_clear().envs(tools)
.args(["-c", "test -z \"${LF_BIN+x}${LF_HOME+x}${LF_DB_PATH+x}\" && test \"$LF_CONTROL_BIN\" = /control/lf"])
.status().unwrap().success());
}
#[tokio::test]
async fn provider_child_cannot_read_the_bridge_token() {
let mut command = tokio::process::Command::new("/bin/sh");
command.args(["-c", "test -z \"${LF_DISCORD_TOKEN+x}\""]);
command.env(crate::engine::process::DISCORD_TOKEN_ENV, "fixture-token");
let mut config = AgentConfig::default();
config.env.insert(
crate::engine::process::DISCORD_TOKEN_ENV.into(),
"fixture-override".into(),
);
configure_agent_env(&mut command, &config);
assert!(command.status().await.unwrap().success());
}
#[test]
fn vendor_environment_preserves_development_and_release_contexts() {
let mut command = std::process::Command::new("vendor");
command
.env("LF_BIN", "/ambient/lf")
.env("LF_HOME", "/production")
.env("LF_DB_PATH", "/production/loopflow.db")
.env("LF_CONTROL_HOME", "/old-control");
set_vendor_std_env(
&mut command,
Path::new("/control/lf"),
Path::new("/custom"),
Path::new("/custom/loopflow.db"),
)
.unwrap();
let environment = command
.get_envs()
.map(|(key, value)| (key.to_string_lossy().to_string(), value.map(OsString::from)))
.collect::<std::collections::HashMap<_, _>>();
let development = !crate::build_info::provenance().is_release();
assert_eq!(
environment["LF_HOME"],
development.then(|| OsString::from("/custom"))
);
assert_eq!(
environment["LF_DB_PATH"],
development.then(|| OsString::from("/custom/loopflow.db"))
);
assert_eq!(
environment["LF_BIN"],
development.then(|| OsString::from("/control/lf"))
);
assert_eq!(
environment["LF_CONTROL_BIN"],
Some(OsString::from("/control/lf"))
);
assert_eq!(
environment["LF_CONTROL_HOME"],
Some(OsString::from("/custom"))
);
assert_eq!(
environment["LF_CONTROL_DB_PATH"],
Some(OsString::from("/custom/loopflow.db"))
);
}
#[test]
fn agent_receives_only_fresh_generic_run_identity() {
let mut command = tokio::process::Command::new("vendor");
command
.env(crate::durable::RUN_ID_ENV, "run_stale")
.env(crate::session_record::RUN_DIR_ENV, "/stale/run");
let mut config = crate::engine::agent::AgentConfig::default();
config.env.insert(
crate::durable::RUN_ID_ENV.to_string(),
"run_fresh".to_string(),
);
config.env.insert(
crate::session_record::RUN_DIR_ENV.to_string(),
"/fresh/run".to_string(),
);
configure_agent_env(&mut command, &config);
let environment = command
.as_std()
.get_envs()
.map(|(key, value)| (key.to_string_lossy().to_string(), value.map(OsString::from)))
.collect::<std::collections::HashMap<_, _>>();
assert_eq!(
environment[crate::durable::RUN_ID_ENV],
Some(OsString::from("run_fresh"))
);
assert_eq!(
environment[crate::session_record::RUN_DIR_ENV],
Some(OsString::from("/fresh/run"))
);
}
#[test]
fn provider_drops_legacy_worktree_writer_authority() {
let mut command = tokio::process::Command::new("vendor");
command.env("LF_WORKTREE_WRITER_ID", "writer_stale");
let mut config = AgentConfig::default();
config.env.insert(
"LF_WORKTREE_WRITER_ID".to_string(),
"writer_explicit".to_string(),
);
configure_agent_env(&mut command, &config);
let environment = command
.as_std()
.get_envs()
.map(|(key, value)| (key.to_string_lossy().to_string(), value.map(OsString::from)))
.collect::<std::collections::HashMap<_, _>>();
assert_eq!(environment["LF_WORKTREE_WRITER_ID"], None);
}
}
#[derive(Debug, thiserror::Error)]
pub enum HarnessError {
#[error("turn already in progress")]
TurnAlreadyInProgress,
}
pub fn is_turn_in_progress(err: &anyhow::Error) -> bool {
matches!(
err.downcast_ref::<HarnessError>(),
Some(HarnessError::TurnAlreadyInProgress)
)
}
pub fn is_terminal_harness_error(code: &str) -> bool {
matches!(code, "codex_disconnected" | "opencode_disconnected")
}
#[derive(Debug, Clone, PartialEq, Eq)]
#[non_exhaustive]
pub enum SendCurrentOutcome {
Sent {
provider_turn_id: String,
},
NotSteerable,
Failed {
error: String,
},
Unknown {
provider_turn_id: Option<String>,
error: String,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum ApprovalPolicy {
AutoApprove,
}
#[async_trait]
pub trait Harness: Send + Sync {
async fn start(&mut self, config: &AgentConfig) -> Result<()>;
async fn send_input(&mut self, content: &str) -> Result<()>;
async fn send_current(&mut self, _content: &str) -> SendCurrentOutcome {
SendCurrentOutcome::NotSteerable
}
async fn interrupt(&mut self) -> Result<()>;
async fn stop(&mut self) -> Result<()>;
fn provider_session_id(&self) -> Option<String>;
fn process_id(&self) -> Option<u32> {
None
}
fn process_group_id(&self) -> Option<u32> {
None
}
fn set_raw_provider_sender(
&mut self,
_raw_provider: Option<mpsc::UnboundedSender<RawProviderEvent>>,
) {
}
fn set_provider_session_id(&mut self, _provider_session_id: Option<String>) {}
fn set_provider_account_id(&mut self, _account_id: Option<crate::store::ProviderAccountId>) {}
fn provider_account_id(&self) -> Option<crate::store::ProviderAccountId> {
None
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum HarnessKind {
Codex,
Claude,
OpenCode,
}
impl HarnessKind {
pub fn parse(name: &str) -> Option<Self> {
match name.trim().to_ascii_lowercase().as_str() {
"codex" => Some(Self::Codex),
"claude" => Some(Self::Claude),
"opencode" => Some(Self::OpenCode),
_ => None,
}
}
pub fn as_str(self) -> &'static str {
match self {
Self::Codex => "codex",
Self::Claude => "claude",
Self::OpenCode => "opencode",
}
}
fn create(
self,
approval: ApprovalPolicy,
event_tx: mpsc::UnboundedSender<ConversationEvent>,
) -> Box<dyn Harness> {
match self {
Self::Codex => Box::new(codex::CodexHarness::new(event_tx, approval)),
Self::Claude => Box::new(claude::ClaudeHarness::new(event_tx)),
Self::OpenCode => Box::new(opencode::OpenCodeHarness::new(event_tx, approval)),
}
}
}
pub fn canonical_harness(name: &str) -> Option<&'static str> {
HarnessKind::parse(name).map(HarnessKind::as_str)
}
pub fn default_create_harness(
name: &str,
approval: ApprovalPolicy,
event_tx: mpsc::UnboundedSender<ConversationEvent>,
) -> Result<Box<dyn Harness>> {
if let Some(kind) = HarnessKind::parse(name) {
return Ok(kind.create(approval, event_tx));
}
anyhow::bail!(
"unsupported session harness: {}",
name.trim().to_lowercase()
)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn canonical_harness_is_case_insensitive_and_trimmed() {
assert_eq!(canonical_harness(" claUDe "), Some("claude"));
assert_eq!(canonical_harness(" CODEX"), Some("codex"));
assert_eq!(canonical_harness("OpenCode"), Some("opencode"));
assert_eq!(canonical_harness("lfharness"), None);
}
#[test]
fn default_create_harness_rejects_unknown() {
let (tx, _rx) = mpsc::unbounded_channel();
match default_create_harness("lfharness", ApprovalPolicy::AutoApprove, tx) {
Ok(_) => panic!("should reject unknown harness"),
Err(err) => assert!(err.to_string().contains("unsupported session harness")),
}
}
#[test]
fn terminal_harness_error_recognizes_disconnects_only() {
assert!(is_terminal_harness_error("opencode_disconnected"));
assert!(is_terminal_harness_error("codex_disconnected"));
assert!(!is_terminal_harness_error("opencode_error"));
assert!(!is_terminal_harness_error("claude_harness_crashed"));
}
#[test]
fn hollow_and_decode_gap_are_turn_terminal_not_session_terminal() {
assert!(!is_terminal_harness_error("opencode_hollow_body"));
assert!(!is_terminal_harness_error("opencode_decode_gap"));
}
#[tokio::test]
async fn current_send_is_decided_from_the_active_turn() {
let (tx, _rx) = mpsc::unbounded_channel();
for name in ["codex", "claude", "opencode"] {
let mut harness = default_create_harness(name, ApprovalPolicy::AutoApprove, tx.clone())
.expect("known harness");
assert_eq!(
harness.send_current("direction").await,
SendCurrentOutcome::NotSteerable,
"an inactive {name} harness has no exact Turn to steer"
);
}
}
}