use std::collections::{HashMap, HashSet};
use std::future::Future;
use std::path::Path;
use std::path::PathBuf;
use std::pin::Pin;
use std::time::Duration;
use chimera_core::*;
use crate::config::OpenCodeConfig;
const READY_POLL_INTERVAL: Duration = Duration::from_millis(500);
const READY_PROBE_TIMEOUT: Duration = Duration::from_secs(1);
const STARTUP_STDERR_TAIL_BYTES: usize = 4096;
pub struct OpenCodeAgentSession {
process: tokio::process::Child,
client: reqwest::Client,
port: u16,
opencode_session_id: String,
directory: Option<String>,
model: Option<String>,
system_prompt: Option<String>,
variant: Option<String>,
_isolated_config_root: Option<tempfile::TempDir>,
}
impl OpenCodeAgentSession {
pub(crate) async fn new(
binary_path: PathBuf,
config: SessionConfig<OpenCodeConfig>,
) -> Result<Self> {
Self::start(binary_path, config, None).await
}
pub(crate) async fn resume(
binary_path: PathBuf,
config: SessionConfig<OpenCodeConfig>,
session_id: &str,
) -> Result<Self> {
Self::start(binary_path, config, Some(session_id.to_owned())).await
}
async fn start(
binary_path: PathBuf,
config: SessionConfig<OpenCodeConfig>,
existing_session_id: Option<String>,
) -> Result<Self> {
let port = find_free_port()?;
let startup_timeout = config.backend.startup_timeout();
validate_launch_environment(&config)?;
let config_json = build_config_json(&config.backend, config.model.as_deref())?;
let (launch_env, isolated_config_root) = build_launch_env(&config.backend, &config_json)?;
let startup_log_root = tempfile::tempdir().map_err(|e| AgentError::Other {
message: "failed to create OpenCode startup log root".into(),
source: Some(Box::new(e)),
})?;
let stderr_log_path = startup_log_root.path().join("opencode-stderr.log");
let stderr_log =
std::fs::File::create(&stderr_log_path).map_err(|e| AgentError::Other {
message: format!(
"failed to create OpenCode startup stderr log {}",
stderr_log_path.display()
),
source: Some(Box::new(e)),
})?;
let mut cmd = if binary_path.join("src/index.ts").exists() {
let mut c = tokio::process::Command::new("bun");
c.arg("run")
.arg("--conditions=browser")
.arg("./src/index.ts")
.current_dir(&binary_path);
c
} else {
tokio::process::Command::new(&binary_path)
};
cmd.arg("serve").arg("--port").arg(port.to_string());
for key in &config.env_remove {
cmd.env_remove(key);
}
for (key, value) in &config.env {
cmd.env(key, value);
}
for (key, value) in &launch_env {
cmd.env(key, value);
}
let process = cmd
.stdout(std::process::Stdio::null())
.stderr(std::process::Stdio::from(stderr_log))
.spawn()
.map_err(|e| AgentError::SpawnFailed { source: e })?;
let mut process = process;
let client = reqwest::Client::builder()
.timeout(Duration::from_secs(300))
.build()
.map_err(|e| AgentError::Other {
message: format!("failed to build HTTP client: {e}"),
source: Some(Box::new(e)),
})?;
let startup_context = OpenCodeStartupContext {
binary_path: binary_path.clone(),
port,
mcp_server_count: config.backend.mcp_servers.len(),
mcp_config_mode: config.backend.mcp_config_mode,
stderr_log_path,
};
wait_for_ready(&client, &mut process, startup_timeout, &startup_context).await?;
let cwd = config
.cwd
.as_ref()
.and_then(|p| p.to_str())
.map(str::to_owned);
let opencode_session_id = match existing_session_id {
Some(session_id) => session_id,
None => {
create_session(
&client,
port,
cwd.as_deref(),
config.backend.session_title.as_deref(),
)
.await?
}
};
let model = config.model.clone();
let system_prompt = config
.system_prompt
.clone()
.or_else(|| config.backend.system_prompt.clone());
let variant = config.backend.variant.clone();
Ok(Self {
process,
client,
port,
opencode_session_id,
directory: cwd,
model,
system_prompt,
variant,
_isolated_config_root: isolated_config_root,
})
}
}
impl Drop for OpenCodeAgentSession {
fn drop(&mut self) {
let _ = self.process.start_kill();
}
}
impl Session for OpenCodeAgentSession {
fn turn_stream(
&mut self,
input: Input,
options: TurnOptions,
) -> Pin<Box<dyn Future<Output = Result<EventStream>> + Send + '_>> {
let text = input_to_text(&input);
let options_validation = validate_turn_options(&options);
let model = self.model.clone();
let directory = self.directory.clone();
let session_id = self.opencode_session_id.clone();
let port = self.port;
let system_prompt = self.system_prompt.clone();
let variant = self.variant.clone();
let client = self.client.clone();
Box::pin(async move {
let text = text?;
options_validation?;
let body = build_prompt_body(
&text,
model.as_deref(),
system_prompt.as_deref(),
variant.as_deref(),
);
let sse_client = reqwest::Client::builder()
.timeout(Duration::from_secs(600))
.build()
.map_err(|e| AgentError::Other {
message: format!("failed to build SSE client: {e}"),
source: Some(Box::new(e)),
})?;
let mut sse_request = sse_client.get(format!("http://127.0.0.1:{port}/event"));
if let Some(directory) = directory.as_deref() {
sse_request = sse_request.query(&[("directory", directory)]);
}
let sse_resp = sse_request.send().await.map_err(|e| AgentError::Other {
message: format!("failed to connect to opencode event stream: {e}"),
source: Some(Box::new(e)),
})?;
let url = format!("http://127.0.0.1:{port}/session/{session_id}/prompt_async");
let mut prompt_request = client.post(&url).json(&body);
if let Some(directory) = directory.as_deref() {
prompt_request = prompt_request.query(&[("directory", directory)]);
}
let resp = prompt_request.send().await.map_err(|e| AgentError::Other {
message: format!("opencode prompt_async request failed: {e}"),
source: Some(Box::new(e)),
})?;
let status = resp.status();
if !status.is_success() && status.as_u16() != 204 {
let err_body = resp.text().await.unwrap_or_default();
return Err(AgentError::TurnFailed {
message: format!("opencode HTTP {status}: {err_body}"),
});
}
let (tx, rx) = tokio::sync::mpsc::channel(64);
let _ = tx.send(Ok(AgentEvent::TurnStarted)).await;
tokio::spawn(read_sse_events(sse_resp, session_id, tx));
Ok(EventStream::from_receiver(rx))
})
}
fn session_id(&self) -> Option<&str> {
Some(&self.opencode_session_id)
}
fn interrupt(&mut self) -> Pin<Box<dyn Future<Output = Result<()>> + Send + '_>> {
Box::pin(async move {
Err(AgentError::ControlError {
message: "opencode agent mode does not support interrupt yet".into(),
})
})
}
}
async fn read_sse_events(
mut resp: reqwest::Response,
session_id: String,
tx: tokio::sync::mpsc::Sender<Result<AgentEvent>>,
) {
let mut raw_buf: Vec<u8> = Vec::new();
let mut state = OpenCodeTurnState::default();
loop {
match resp.chunk().await {
Err(e) => {
let _ = tx
.send(Err(AgentError::Other {
message: format!("SSE stream error: {e}"),
source: Some(Box::new(e)),
}))
.await;
return;
}
Ok(None) => break,
Ok(Some(bytes)) => {
raw_buf.extend_from_slice(&bytes);
}
}
while let Some((end, delimiter_len)) = find_event_end(&raw_buf) {
let event_bytes = raw_buf[..end].to_vec();
raw_buf.drain(..end + delimiter_len);
let event_str = match std::str::from_utf8(&event_bytes) {
Ok(s) => s.to_string(),
Err(_) => continue,
};
let json_str: String = event_str
.lines()
.filter_map(|line| line.strip_prefix("data: "))
.collect::<Vec<_>>()
.join("");
if json_str.is_empty() {
continue;
}
let frame: SseFrame = match serde_json::from_str(&json_str) {
Ok(f) => f,
Err(_) => continue,
};
let done = handle_sse_event(frame, &session_id, &tx, &mut state).await;
if done {
return;
}
}
}
let _ = tx
.send(Err(AgentError::TurnFailed {
message: "opencode event stream closed before session became idle".into(),
}))
.await;
}
async fn handle_sse_event(
frame: SseFrame,
session_id: &str,
tx: &tokio::sync::mpsc::Sender<Result<AgentEvent>>,
state: &mut OpenCodeTurnState,
) -> bool {
let props = &frame.properties;
match frame.event_type.as_str() {
"message.updated" => {
let sid = props["sessionID"].as_str().unwrap_or("");
if sid != session_id {
return false;
}
if let Err(error) = state.observe_message(&props["info"]) {
let _ = tx.send(Err(error)).await;
return true;
}
}
"message.part.delta" => {
let sid = props["sessionID"].as_str().unwrap_or("");
if sid != session_id {
return false;
}
let field = props["field"].as_str().unwrap_or("");
let delta = props["delta"].as_str().unwrap_or("");
if field == "text" && !delta.is_empty() {
state.full_text.push_str(delta);
let part_id = props["partID"].as_str().unwrap_or("").to_string();
let _ = tx
.send(Ok(AgentEvent::TextDelta {
id: part_id,
delta: delta.to_string(),
}))
.await;
}
}
"message.part.updated" => {
let part = &props["part"];
let sid = part["sessionID"].as_str().unwrap_or("");
if sid != session_id {
return false;
}
if part["type"].as_str() != Some("tool") {
return false;
}
let call_id = part["callID"].as_str().unwrap_or("").to_string();
let tool_name = part["tool"].as_str().unwrap_or("").to_string();
let tool_state = &part["state"];
let status = tool_state["status"].as_str().unwrap_or("pending");
let input = {
let v = tool_state["input"].clone();
if v.is_null() {
serde_json::json!({})
} else {
v
}
};
if !state.emitted_tool_starts.contains(&call_id) && status != "pending" {
state.emitted_tool_starts.insert(call_id.clone());
let _ = tx
.send(Ok(AgentEvent::ToolCall {
id: call_id.clone(),
name: tool_name.clone(),
arguments: input.to_string(),
}))
.await;
let _ = tx
.send(Ok(AgentEvent::ActionStarted {
id: call_id.clone(),
action: Action::ToolUse {
name: tool_name.clone(),
input: input.clone(),
result: None,
is_error: false,
},
}))
.await;
}
if status == "completed" || status == "error" {
let is_error = status == "error";
let result = if is_error {
tool_state["error"]
.as_str()
.map(|s| serde_json::Value::String(s.to_string()))
} else {
tool_state["output"]
.as_str()
.map(|s| serde_json::Value::String(s.to_string()))
};
let _ = tx
.send(Ok(AgentEvent::ActionCompleted {
id: call_id,
action: Action::ToolUse {
name: tool_name,
input,
result,
is_error,
},
}))
.await;
}
}
"session.status" => {
let sid = props["sessionID"].as_str().unwrap_or("");
if sid != session_id {
return false;
}
let status_type = props["status"]["type"].as_str().unwrap_or("");
if status_type == "idle" && state.has_turn_activity() {
let completion = match state.completion(session_id) {
Ok(completion) => completion,
Err(error) => {
let _ = tx.send(Err(error)).await;
return true;
}
};
let _ = tx.send(Ok(completion)).await;
return true;
}
}
"session.error" => {
let sid = props["sessionID"].as_str().unwrap_or("");
if !sid.is_empty() && sid != session_id {
return false;
}
let error_msg = props["error"]["data"]["message"]
.as_str()
.or_else(|| props["error"]["message"].as_str())
.or_else(|| props["error"]["properties"]["message"].as_str())
.unwrap_or("unknown error");
let _ = tx
.send(Err(AgentError::TurnFailed {
message: format!("opencode session error: {error_msg}"),
}))
.await;
return true;
}
_ => {}
}
false
}
#[derive(Default)]
struct OpenCodeTurnState {
full_text: String,
emitted_tool_starts: HashSet<String>,
user_message_id: Option<String>,
assistants: HashMap<String, AssistantUsageSnapshot>,
last_completed_assistant_id: Option<String>,
}
impl OpenCodeTurnState {
fn observe_message(&mut self, value: &serde_json::Value) -> Result<()> {
let message: MessageSnapshot =
serde_json::from_value(value.clone()).map_err(|error| AgentError::TurnFailed {
message: format!("invalid opencode message.updated event: {error}"),
})?;
match message.role.as_str() {
"user" => {
if self.user_message_id.is_none() {
self.user_message_id = Some(message.id);
}
}
"assistant" => {
let Some(user_message_id) = self.user_message_id.as_deref() else {
return Ok(());
};
if message.parent_id.as_deref() != Some(user_message_id) {
return Ok(());
}
let snapshot = AssistantUsageSnapshot::try_from(message)?;
if snapshot.finish_reason.is_some() {
self.last_completed_assistant_id = Some(snapshot.id.clone());
}
self.assistants.insert(snapshot.id.clone(), snapshot);
}
_ => {}
}
Ok(())
}
fn has_turn_activity(&self) -> bool {
self.user_message_id.is_some() || !self.assistants.is_empty()
}
fn completion(&self, session_id: &str) -> Result<AgentEvent> {
if self.assistants.is_empty() {
return Err(AgentError::TurnFailed {
message: "opencode session became idle without an assistant response".into(),
});
}
if self
.assistants
.values()
.any(|snapshot| snapshot.finish_reason.is_none())
{
return Err(AgentError::TurnFailed {
message: "opencode session became idle with an incomplete assistant response"
.into(),
});
}
let finish_reason = self
.last_completed_assistant_id
.as_ref()
.and_then(|id| self.assistants.get(id))
.and_then(|snapshot| snapshot.finish_reason.clone());
Ok(AgentEvent::TurnCompleted {
usage: self.aggregate_usage()?,
result: Some(TurnResult {
text: (!self.full_text.is_empty()).then(|| self.full_text.clone()),
session_id: Some(session_id.to_string()),
duration_ms: None,
num_turns: None,
turn_status: Some(TurnStatus::Completed),
provider_meta: finish_reason.map(|finish_reason| TurnProviderMeta::OpenCode {
finish_reason: Some(finish_reason),
}),
}),
})
}
fn aggregate_usage(&self) -> Result<Option<Usage>> {
let snapshots = self.assistants.values().collect::<Vec<_>>();
let input_tokens = checked_sum_optional(snapshots.iter().map(|item| item.input_tokens))?;
let output_tokens = checked_sum_optional(
snapshots
.iter()
.map(|item| checked_add_optional(item.output_tokens, item.reasoning_tokens))
.collect::<Result<Vec<_>>>()?
.into_iter(),
)?;
let cache_read_tokens =
checked_sum_optional(snapshots.iter().map(|item| item.cache_read_tokens))?;
let cache_creation_tokens =
checked_sum_optional(snapshots.iter().map(|item| item.cache_creation_tokens))?;
let total_cost_usd = checked_sum_cost_optional(snapshots.iter().map(|item| item.cost_usd))?;
if input_tokens.is_none()
&& output_tokens.is_none()
&& cache_read_tokens.is_none()
&& cache_creation_tokens.is_none()
&& total_cost_usd.is_none()
{
return Ok(None);
}
Ok(Some(Usage {
input_tokens,
output_tokens,
cache_read_tokens,
cache_creation_tokens,
total_cost_usd,
}))
}
}
fn checked_add_optional(left: Option<u64>, right: Option<u64>) -> Result<Option<u64>> {
match (left, right) {
(Some(left), Some(right)) => {
left.checked_add(right)
.map(Some)
.ok_or_else(|| AgentError::TurnFailed {
message: "opencode usage token count overflowed".into(),
})
}
_ => Ok(None),
}
}
fn checked_sum_optional(mut values: impl Iterator<Item = Option<u64>>) -> Result<Option<u64>> {
values.try_fold(Some(0_u64), |total, value| match (total, value) {
(Some(total), Some(value)) => {
total
.checked_add(value)
.map(Some)
.ok_or_else(|| AgentError::TurnFailed {
message: "opencode usage token count overflowed".into(),
})
}
_ => Ok(None),
})
}
fn checked_sum_cost_optional(mut values: impl Iterator<Item = Option<f64>>) -> Result<Option<f64>> {
values.try_fold(Some(0.0_f64), |total, value| match (total, value) {
(Some(total), Some(value)) => {
let sum = total + value;
if sum.is_finite() {
Ok(Some(sum))
} else {
Err(AgentError::TurnFailed {
message: "opencode usage cost overflowed".into(),
})
}
}
_ => Ok(None),
})
}
#[derive(serde::Deserialize)]
struct MessageSnapshot {
id: String,
role: String,
#[serde(rename = "parentID")]
parent_id: Option<String>,
finish: Option<String>,
cost: Option<f64>,
tokens: Option<MessageTokens>,
}
#[derive(serde::Deserialize)]
struct MessageTokens {
input: Option<u64>,
output: Option<u64>,
reasoning: Option<u64>,
cache: Option<MessageCacheTokens>,
}
#[derive(serde::Deserialize)]
struct MessageCacheTokens {
read: Option<u64>,
write: Option<u64>,
}
struct AssistantUsageSnapshot {
id: String,
finish_reason: Option<String>,
cost_usd: Option<f64>,
input_tokens: Option<u64>,
output_tokens: Option<u64>,
reasoning_tokens: Option<u64>,
cache_read_tokens: Option<u64>,
cache_creation_tokens: Option<u64>,
}
impl TryFrom<MessageSnapshot> for AssistantUsageSnapshot {
type Error = AgentError;
fn try_from(message: MessageSnapshot) -> Result<Self> {
if message
.cost
.is_some_and(|cost| !cost.is_finite() || cost < 0.0)
{
return Err(AgentError::TurnFailed {
message: "opencode reported an invalid usage cost".into(),
});
}
let tokens = message.tokens;
let cache = tokens.as_ref().and_then(|tokens| tokens.cache.as_ref());
Ok(Self {
id: message.id,
finish_reason: message.finish,
cost_usd: message.cost,
input_tokens: tokens.as_ref().and_then(|tokens| tokens.input),
output_tokens: tokens.as_ref().and_then(|tokens| tokens.output),
reasoning_tokens: tokens.as_ref().and_then(|tokens| tokens.reasoning),
cache_read_tokens: cache.and_then(|cache| cache.read),
cache_creation_tokens: cache.and_then(|cache| cache.write),
})
}
}
fn find_event_end(buf: &[u8]) -> Option<(usize, usize)> {
let lf = buf.windows(2).position(|window| window == b"\n\n");
let crlf = buf.windows(4).position(|window| window == b"\r\n\r\n");
match (lf, crlf) {
(Some(lf), Some(crlf)) if lf < crlf => Some((lf, 2)),
(Some(_), Some(crlf)) => Some((crlf, 4)),
(Some(lf), None) => Some((lf, 2)),
(None, Some(crlf)) => Some((crlf, 4)),
(None, None) => None,
}
}
fn find_free_port() -> Result<u16> {
let listener = std::net::TcpListener::bind("127.0.0.1:0").map_err(|e| AgentError::Other {
message: format!("failed to find free port: {e}"),
source: Some(Box::new(e)),
})?;
Ok(listener.local_addr().unwrap().port())
}
fn build_config_json(backend: &OpenCodeConfig, model: Option<&str>) -> Result<String> {
let mcp: serde_json::Map<String, serde_json::Value> = backend
.mcp_servers
.iter()
.map(|s| {
let entry = serde_json::json!({
"type": "local",
"command": s.command,
"environment": s.env,
});
(s.name.clone(), entry)
})
.collect();
let mut config = serde_json::json!({ "mcp": mcp });
if let Some(base_url) = backend.agent_base_url.as_deref() {
let provider_id = model
.and_then(|model| model.split_once('/').map(|(provider, _)| provider))
.filter(|provider| !provider.is_empty())
.ok_or_else(|| AgentError::Other {
message: "opencode agent_base_url requires a provider/model selection".into(),
source: None,
})?;
config["provider"][provider_id]["options"]["baseURL"] =
serde_json::Value::String(base_url.to_owned());
config["small_model"] = serde_json::Value::String(model.expect("validated model").into());
}
serde_json::to_string(&config).map_err(|e| AgentError::Other {
message: format!("failed to serialize opencode config: {e}"),
source: Some(Box::new(e)),
})
}
const RESERVED_ENV: [&str; 5] = [
"OPENCODE_CONFIG_CONTENT",
"OPENCODE_CONFIG",
"OPENCODE_CONFIG_DIR",
"OPENCODE_DISABLE_PROJECT_CONFIG",
"OPENCODE_DISABLE_DEFAULT_PLUGINS",
];
fn validate_launch_environment(config: &SessionConfig<OpenCodeConfig>) -> Result<()> {
if let Some(key) = RESERVED_ENV
.iter()
.find(|key| config.env.contains_key(**key) || config.env_remove.contains(**key))
{
return Err(AgentError::Other {
message: format!("opencode agent session reserves environment variable {key}"),
source: None,
});
}
Ok(())
}
fn build_launch_env(
backend: &OpenCodeConfig,
config_json: &str,
) -> Result<(
std::collections::HashMap<String, String>,
Option<tempfile::TempDir>,
)> {
let mut env = std::collections::HashMap::new();
env.insert("OPENCODE_CONFIG_CONTENT".into(), config_json.to_owned());
if !matches!(backend.mcp_config_mode, McpConfigMode::ExplicitOnly) {
return Ok((env, None));
}
let root = tempfile::tempdir().map_err(|e| AgentError::Other {
message: "failed to create isolated OpenCode config root".into(),
source: Some(Box::new(e)),
})?;
let xdg_config = root.path().join("xdg-config");
let opencode_config_dir = root.path().join("opencode-config");
let explicit_config = root.path().join("opencode.explicit.json");
for dir in [&xdg_config, &opencode_config_dir] {
std::fs::create_dir_all(dir).map_err(|e| AgentError::Other {
message: format!("failed to prepare OpenCode isolated path {}", dir.display()),
source: Some(Box::new(e)),
})?;
}
std::fs::write(&explicit_config, "{}").map_err(|e| AgentError::Other {
message: format!(
"failed to write OpenCode isolated config file {}",
explicit_config.display()
),
source: Some(Box::new(e)),
})?;
env.insert(
"OPENCODE_CONFIG".into(),
explicit_config.display().to_string(),
);
env.insert("XDG_CONFIG_HOME".into(), xdg_config.display().to_string());
env.insert(
"OPENCODE_CONFIG_DIR".into(),
opencode_config_dir.display().to_string(),
);
env.insert("OPENCODE_DISABLE_PROJECT_CONFIG".into(), "true".into());
env.insert("OPENCODE_DISABLE_DEFAULT_PLUGINS".into(), "true".into());
Ok((env, Some(root)))
}
#[derive(Debug)]
struct OpenCodeStartupContext {
binary_path: PathBuf,
port: u16,
mcp_server_count: usize,
mcp_config_mode: McpConfigMode,
stderr_log_path: PathBuf,
}
async fn wait_for_ready(
client: &reqwest::Client,
process: &mut tokio::process::Child,
startup_timeout: Duration,
context: &OpenCodeStartupContext,
) -> Result<()> {
let url = format!("http://127.0.0.1:{}/session", context.port);
let deadline = std::time::Instant::now() + startup_timeout;
let mut last_probe = None;
loop {
if let Some(status) = process.try_wait().map_err(|e| AgentError::Other {
message: format!("failed to inspect opencode child process state: {e}"),
source: Some(Box::new(e)),
})? {
return Err(startup_process_exit_error(status, context));
}
if std::time::Instant::now() > deadline {
return Err(startup_timeout_error(
startup_timeout,
context,
last_probe.as_deref(),
));
}
match client.get(&url).timeout(READY_PROBE_TIMEOUT).send().await {
Ok(r) if r.status().as_u16() < 500 => return Ok(()),
Ok(r) => {
last_probe = Some(format!("HTTP {}", r.status()));
tokio::time::sleep(READY_POLL_INTERVAL).await;
}
Err(e) => {
last_probe = Some(e.to_string());
tokio::time::sleep(READY_POLL_INTERVAL).await;
}
}
}
}
fn startup_timeout_error(
startup_timeout: Duration,
context: &OpenCodeStartupContext,
last_probe: Option<&str>,
) -> AgentError {
let mut message = format!(
"opencode agent startup timed out after {:?} waiting for /session on port {} (binary: {}, mcp servers: {}, mcp mode: {:?})",
startup_timeout,
context.port,
context.binary_path.display(),
context.mcp_server_count,
context.mcp_config_mode,
);
if let Some(last_probe) = last_probe {
message.push_str(&format!(", last probe: {last_probe}"));
}
if let Some(stderr_tail) = read_stderr_tail(&context.stderr_log_path) {
message.push_str(&format!(", stderr tail: {stderr_tail}"));
}
message
.push_str(". Increase OpenCodeConfig.startup_timeout_secs for slower MCP-heavy startup.");
AgentError::Other {
message,
source: Some(Box::new(AgentError::Timeout {
duration: startup_timeout,
})),
}
}
fn startup_process_exit_error(
status: std::process::ExitStatus,
context: &OpenCodeStartupContext,
) -> AgentError {
let stderr_tail = read_stderr_tail(&context.stderr_log_path).unwrap_or_default();
if let Some(code) = status.code() {
AgentError::ProcessFailed {
code,
stderr: format_startup_stderr(stderr_tail, context),
}
} else {
AgentError::ProcessKilled {
stderr: format_startup_stderr(stderr_tail, context),
}
}
}
fn format_startup_stderr(stderr_tail: String, context: &OpenCodeStartupContext) -> String {
let mut message = format!(
"opencode agent exited before readiness on port {} (binary: {}, mcp servers: {}, mcp mode: {:?})",
context.port,
context.binary_path.display(),
context.mcp_server_count,
context.mcp_config_mode,
);
if !stderr_tail.is_empty() {
message.push_str(&format!("; stderr tail: {stderr_tail}"));
}
message
}
fn read_stderr_tail(path: &Path) -> Option<String> {
let bytes = std::fs::read(path).ok()?;
if bytes.is_empty() {
return None;
}
let start = bytes.len().saturating_sub(STARTUP_STDERR_TAIL_BYTES);
let tail = String::from_utf8_lossy(&bytes[start..]).trim().to_string();
if tail.is_empty() { None } else { Some(tail) }
}
async fn create_session(
client: &reqwest::Client,
port: u16,
cwd: Option<&str>,
title: Option<&str>,
) -> Result<String> {
let body = create_session_body(title);
let mut req = client
.post(format!("http://127.0.0.1:{port}/session"))
.json(&body);
if let Some(dir) = cwd {
req = client
.post(format!("http://127.0.0.1:{port}/session"))
.query(&[("directory", dir)])
.json(&body);
}
let resp = req.send().await.map_err(|e| AgentError::Other {
message: format!("failed to create opencode session: {e}"),
source: Some(Box::new(e)),
})?;
if !resp.status().is_success() {
let status = resp.status();
let body = resp.text().await.unwrap_or_default();
return Err(AgentError::Other {
message: format!("failed to create opencode session: HTTP {status}: {body}"),
source: None,
});
}
let session: SessionInfo = resp.json().await.map_err(|e| AgentError::Other {
message: format!("failed to parse session response: {e}"),
source: Some(Box::new(e)),
})?;
Ok(session.id)
}
fn create_session_body(title: Option<&str>) -> serde_json::Value {
let mut body = serde_json::json!({ "permission": [] });
if let Some(title) = title {
body["title"] = serde_json::Value::String(title.to_owned());
}
body
}
fn build_prompt_body(
text: &str,
model: Option<&str>,
system_prompt: Option<&str>,
configured_variant: Option<&str>,
) -> serde_json::Value {
let mut body = serde_json::json!({
"parts": [{ "type": "text", "text": text }],
});
let mut variant = configured_variant.map(str::to_owned);
if let Some(model) = model {
let (provider_id, model_id, parsed_variant) = parse_model_selection(model);
if variant.is_none() {
variant = parsed_variant.map(str::to_owned);
}
body["model"] = serde_json::json!({
"providerID": provider_id,
"modelID": model_id,
});
}
if let Some(system_prompt) = system_prompt {
body["system"] = serde_json::Value::String(system_prompt.to_owned());
}
if let Some(variant) = variant {
body["variant"] = serde_json::Value::String(variant);
}
body
}
fn validate_turn_options(options: &TurnOptions) -> Result<()> {
if options.output_schema.is_some() {
return Err(AgentError::Other {
message: "opencode agent mode does not support output_schema".into(),
source: None,
});
}
Ok(())
}
fn parse_model_selection(model: &str) -> (&str, &str, Option<&str>) {
let mut parts = model.splitn(3, '/');
match (parts.next(), parts.next(), parts.next()) {
(Some(provider), Some(model_id), Some(variant)) => (provider, model_id, Some(variant)),
(Some(provider), Some(model_id), None) => (provider, model_id, None),
(Some(model_id), None, None) => ("opencode", model_id, None),
_ => ("opencode", model, None),
}
}
#[cfg(test)]
mod tests {
use super::*;
#[cfg(unix)]
use std::os::unix::process::ExitStatusExt;
fn user_message(id: &str) -> serde_json::Value {
serde_json::json!({ "id": id, "role": "user" })
}
#[test]
fn finds_lf_and_crlf_sse_event_delimiters() {
assert_eq!(find_event_end(b"data: one\n\nnext"), Some((9, 2)));
assert_eq!(find_event_end(b"data: one\r\n\r\nnext"), Some((9, 4)));
}
#[tokio::test]
async fn translates_message_part_deltas_as_text_deltas() {
let (tx, mut rx) = tokio::sync::mpsc::channel(2);
let mut state = OpenCodeTurnState::default();
let frame = SseFrame {
event_type: "message.part.delta".into(),
properties: serde_json::json!({
"sessionID": "session-1",
"partID": "part-1",
"field": "text",
"delta": "hello"
}),
};
assert!(!handle_sse_event(frame, "session-1", &tx, &mut state).await);
let event = rx.recv().await.unwrap().unwrap();
assert!(matches!(
event,
AgentEvent::TextDelta { id, delta }
if id == "part-1" && delta == "hello"
));
assert_eq!(state.full_text, "hello");
}
fn assistant_message(
id: &str,
parent_id: &str,
finish: Option<&str>,
tokens: [u64; 5],
cost: f64,
) -> serde_json::Value {
let [input, output, reasoning, cache_read, cache_write] = tokens;
serde_json::json!({
"id": id,
"role": "assistant",
"parentID": parent_id,
"finish": finish,
"cost": cost,
"tokens": {
"input": input,
"output": output,
"reasoning": reasoning,
"cache": { "read": cache_read, "write": cache_write }
}
})
}
#[test]
fn turn_state_replaces_snapshots_and_combines_reasoning_output() {
let mut state = OpenCodeTurnState::default();
state.observe_message(&user_message("user-1")).unwrap();
state
.observe_message(&assistant_message(
"assistant-1",
"user-1",
None,
[10, 1, 0, 0, 0],
0.01,
))
.unwrap();
state
.observe_message(&assistant_message(
"assistant-1",
"user-1",
Some("stop"),
[20, 2, 1, 3, 4],
0.02,
))
.unwrap();
let usage = state.aggregate_usage().unwrap().unwrap();
assert_eq!(usage.input_tokens, Some(20));
assert_eq!(usage.output_tokens, Some(3));
assert_eq!(usage.cache_read_tokens, Some(3));
assert_eq!(usage.cache_creation_tokens, Some(4));
assert_eq!(usage.total_cost_usd, Some(0.02));
}
#[test]
fn turn_state_aggregates_tool_loop_messages() {
let mut state = OpenCodeTurnState::default();
state.observe_message(&user_message("user-1")).unwrap();
state
.observe_message(&assistant_message(
"assistant-1",
"user-1",
Some("tool-calls"),
[10, 2, 1, 3, 4],
0.01,
))
.unwrap();
state
.observe_message(&assistant_message(
"assistant-2",
"user-1",
Some("stop"),
[20, 5, 2, 6, 8],
0.02,
))
.unwrap();
let AgentEvent::TurnCompleted { usage, result } = state.completion("session-1").unwrap()
else {
panic!("expected TurnCompleted");
};
let usage = usage.unwrap();
assert_eq!(usage.input_tokens, Some(30));
assert_eq!(usage.output_tokens, Some(10));
assert_eq!(usage.cache_read_tokens, Some(9));
assert_eq!(usage.cache_creation_tokens, Some(12));
assert_eq!(usage.total_cost_usd, Some(0.03));
assert!(matches!(
result.unwrap().provider_meta,
Some(TurnProviderMeta::OpenCode {
finish_reason: Some(reason)
}) if reason == "stop"
));
}
#[test]
fn turn_state_rejects_invalid_or_incomplete_usage() {
let mut incomplete = OpenCodeTurnState::default();
incomplete.observe_message(&user_message("user-1")).unwrap();
incomplete
.observe_message(&assistant_message(
"assistant-1",
"user-1",
None,
[1, 2, 3, 4, 5],
0.01,
))
.unwrap();
assert!(
incomplete
.completion("session-1")
.unwrap_err()
.to_string()
.contains("incomplete assistant response")
);
let mut invalid = OpenCodeTurnState::default();
invalid.observe_message(&user_message("user-1")).unwrap();
assert!(
invalid
.observe_message(&assistant_message(
"assistant-1",
"user-1",
Some("stop"),
[1, 2, 3, 4, 5],
-0.01,
))
.unwrap_err()
.to_string()
.contains("invalid usage cost")
);
}
#[test]
fn build_launch_env_defaults_to_merge_mode() {
let backend = OpenCodeConfig::default();
let (env, isolated_root) = build_launch_env(&backend, "{\"mcp\":{}}").unwrap();
assert_eq!(
env.get("OPENCODE_CONFIG_CONTENT").map(String::as_str),
Some("{\"mcp\":{}}")
);
assert!(isolated_root.is_none());
assert!(!env.contains_key("OPENCODE_DISABLE_PROJECT_CONFIG"));
}
#[test]
fn build_config_json_overrides_selected_provider() {
let backend = OpenCodeConfig::builder()
.agent_base_url("http://127.0.0.1:8080/v1")
.build();
let config = build_config_json(&backend, Some("opencode-go/glm-5.1")).unwrap();
let config: serde_json::Value = serde_json::from_str(&config).unwrap();
assert_eq!(
config["provider"]["opencode-go"]["options"]["baseURL"],
"http://127.0.0.1:8080/v1"
);
assert_eq!(config["small_model"], "opencode-go/glm-5.1");
}
#[test]
fn build_config_json_requires_provider_qualified_model() {
let backend = OpenCodeConfig::builder()
.agent_base_url("http://127.0.0.1:8080/v1")
.build();
let error = build_config_json(&backend, Some("glm-5.1")).unwrap_err();
assert!(error.to_string().contains("provider/model"));
}
#[test]
fn launch_environment_rejects_reserved_overrides() {
let config = SessionConfig::builder()
.env(HashMap::from([(
"OPENCODE_CONFIG_CONTENT".into(),
"ambient".into(),
)]))
.backend(OpenCodeConfig::default())
.build();
let error = validate_launch_environment(&config).unwrap_err();
assert!(error.to_string().contains("OPENCODE_CONFIG_CONTENT"));
}
#[test]
fn build_launch_env_explicit_only_isolates_config_roots() {
let backend = OpenCodeConfig::builder()
.mcp_config_mode(McpConfigMode::ExplicitOnly)
.build();
let (env, isolated_root) = build_launch_env(&backend, "{\"mcp\":{}}").unwrap();
let isolated_root = isolated_root.expect("isolated root");
assert_eq!(
env.get("OPENCODE_DISABLE_PROJECT_CONFIG")
.map(String::as_str),
Some("true")
);
assert_eq!(
env.get("OPENCODE_DISABLE_DEFAULT_PLUGINS")
.map(String::as_str),
Some("true")
);
assert!(!env.contains_key("HOME"));
assert_eq!(
env.get("XDG_CONFIG_HOME").map(String::as_str),
Some(
isolated_root
.path()
.join("xdg-config")
.to_string_lossy()
.as_ref()
)
);
assert!(!env.contains_key("XDG_DATA_HOME"));
assert!(!env.contains_key("XDG_STATE_HOME"));
assert!(!env.contains_key("XDG_CACHE_HOME"));
assert_eq!(
env.get("OPENCODE_CONFIG_DIR").map(String::as_str),
Some(
isolated_root
.path()
.join("opencode-config")
.to_string_lossy()
.as_ref()
)
);
assert_eq!(
env.get("OPENCODE_CONFIG").map(String::as_str),
Some(
isolated_root
.path()
.join("opencode.explicit.json")
.to_string_lossy()
.as_ref()
)
);
}
#[test]
fn build_prompt_body_omits_model_when_unset() {
let body = build_prompt_body("hello", None, None, None);
assert!(body.get("model").is_none());
assert_eq!(body["parts"][0]["text"], "hello");
}
#[test]
fn create_session_body_sets_fixed_title() {
assert_eq!(
create_session_body(Some("chimera-test")),
serde_json::json!({ "permission": [], "title": "chimera-test" })
);
}
#[test]
fn build_prompt_body_parses_variant_from_model_selection() {
let body = build_prompt_body("hello", Some("opencode-go/glm-5/high"), None, None);
assert_eq!(body["model"]["providerID"], "opencode-go");
assert_eq!(body["model"]["modelID"], "glm-5");
assert_eq!(body["variant"], "high");
}
#[test]
fn build_prompt_body_prefers_explicit_variant() {
let body = build_prompt_body(
"hello",
Some("opencode-go/glm-5/high"),
None,
Some("medium"),
);
assert_eq!(body["variant"], "medium");
}
#[test]
fn validate_turn_options_rejects_output_schema() {
let err = validate_turn_options(&TurnOptions {
output_schema: Some(serde_json::json!({"type": "object"})),
..Default::default()
})
.unwrap_err();
assert!(
err.to_string()
.contains("opencode agent mode does not support output_schema")
);
}
#[test]
fn startup_timeout_error_mentions_context_and_override_hint() {
let temp = tempfile::tempdir().unwrap();
let stderr_path = temp.path().join("stderr.log");
std::fs::write(&stderr_path, "booting mcp server one\nstill waiting\n").unwrap();
let context = OpenCodeStartupContext {
binary_path: PathBuf::from("/tmp/opencode"),
port: 43123,
mcp_server_count: 3,
mcp_config_mode: McpConfigMode::ExplicitOnly,
stderr_log_path: stderr_path,
};
let err = startup_timeout_error(
Duration::from_secs(45),
&context,
Some("error sending request for url (http://127.0.0.1:43123/session)"),
);
let rendered = err.to_string();
assert!(rendered.contains("timed out after 45s"));
assert!(rendered.contains("mcp servers: 3"));
assert!(rendered.contains("ExplicitOnly"));
assert!(rendered.contains("last probe: error sending request"));
assert!(rendered.contains("stderr tail: booting mcp server one"));
assert!(rendered.contains("startup_timeout_secs"));
}
#[cfg(unix)]
#[test]
fn startup_process_exit_error_includes_context_and_stderr_tail() {
let temp = tempfile::tempdir().unwrap();
let stderr_path = temp.path().join("stderr.log");
std::fs::write(&stderr_path, "fatal: port already in use\n").unwrap();
let context = OpenCodeStartupContext {
binary_path: PathBuf::from("/tmp/opencode"),
port: 43123,
mcp_server_count: 2,
mcp_config_mode: McpConfigMode::Merge,
stderr_log_path: stderr_path,
};
let err = startup_process_exit_error(std::process::ExitStatus::from_raw(1 << 8), &context);
match err {
AgentError::ProcessFailed { code, stderr } => {
assert_eq!(code, 1);
assert!(stderr.contains("exited before readiness"));
assert!(stderr.contains("mcp servers: 2"));
assert!(stderr.contains("fatal: port already in use"));
}
other => panic!("expected ProcessFailed, got {other:?}"),
}
}
#[test]
fn input_to_text_rejects_images() {
let err = input_to_text(&Input::Structured(vec![InputPart::Image(PathBuf::from(
"/tmp/example.png",
))]))
.unwrap_err();
assert!(
err.to_string()
.contains("opencode agent mode does not support image input yet")
);
}
#[tokio::test]
async fn interrupt_returns_control_error() {
let temp = tempfile::tempdir().unwrap();
let binary_path = temp.path().join("unused-opencode");
std::fs::write(&binary_path, "#!/bin/sh\nexit 1\n").unwrap();
let mut session = OpenCodeAgentSession {
process: tokio::process::Command::new("true").spawn().unwrap(),
client: reqwest::Client::new(),
port: 0,
opencode_session_id: "sess-test".into(),
directory: None,
model: None,
system_prompt: None,
variant: None,
_isolated_config_root: None,
};
let err = session.interrupt().await.unwrap_err();
match err {
AgentError::ControlError { message } => {
assert!(message.contains("does not support interrupt"))
}
other => panic!("expected ControlError, got {other:?}"),
}
}
#[tokio::test]
async fn live_smoke_real_binary() {
if std::env::var("APPS_LLM_RUN_LIVE").as_deref() != Ok("1") {
return;
}
if std::env::var("CHIMERA_OPENCODE_AGENT_LIVE").as_deref() != Ok("1") {
return;
}
if !live_auth_available() {
return;
}
let binary_path = match live_binary_path() {
Some(path) => path,
None => return,
};
let config = SessionConfig::builder()
.model(live_model())
.backend(OpenCodeConfig::default())
.build();
let mut session = OpenCodeAgentSession::new(binary_path, config)
.await
.expect("session");
let output = session
.turn(
Input::Text("Reply with exactly: live-opencode-agent-ok".into()),
TurnOptions::default(),
)
.await
.expect("turn");
let response = output.response.unwrap_or_default();
assert!(
response.to_lowercase().contains("live-opencode-agent-ok"),
"unexpected response: {response}"
);
}
#[tokio::test]
async fn live_resume_real_binary() {
if std::env::var("APPS_LLM_RUN_LIVE").as_deref() != Ok("1") {
return;
}
if std::env::var("CHIMERA_OPENCODE_AGENT_LIVE").as_deref() != Ok("1") {
return;
}
if !live_auth_available() {
return;
}
let binary_path = match live_binary_path() {
Some(path) => path,
None => return,
};
let config = SessionConfig::builder()
.model(live_model())
.backend(OpenCodeConfig::default())
.build();
let mut session = OpenCodeAgentSession::new(binary_path.clone(), config.clone())
.await
.expect("session");
session
.turn(
Input::Text("Reply with exactly: live-opencode-resume-seed".into()),
TurnOptions::default(),
)
.await
.expect("seed turn");
let session_id = session.session_id().unwrap().to_owned();
drop(session);
let mut resumed = OpenCodeAgentSession::resume(binary_path, config, &session_id)
.await
.expect("resume session");
assert_eq!(resumed.session_id(), Some(session_id.as_str()));
let output = resumed
.turn(
Input::Text("Reply with exactly: live-opencode-resume-ok".into()),
TurnOptions::default(),
)
.await
.expect("resumed turn");
let response = output.response.unwrap_or_default();
assert!(
response.to_lowercase().contains("live-opencode-resume-ok"),
"unexpected response: {response}"
);
}
fn live_binary_path() -> Option<PathBuf> {
crate::backend::OpenCodeBackend::from_path()
.ok()
.and_then(|backend| backend.binary_path)
}
fn live_auth_available() -> bool {
for key in [
"OPENCODE_ZEN_KEY",
"OPENCODE_API_KEY",
"OPENCODE_GO_API_KEY",
"OPENROUTER_API_KEY",
] {
if std::env::var(key)
.ok()
.is_some_and(|value| !value.is_empty())
{
return true;
}
}
std::env::var("HOME")
.ok()
.map(PathBuf::from)
.map(|home| home.join(".local/share/opencode/auth.json").exists())
.unwrap_or(false)
}
fn live_model() -> String {
std::env::var("OPENCODE_AGENT_MODEL")
.ok()
.filter(|value| !value.is_empty())
.unwrap_or_else(|| "anthropic/claude-sonnet-4-6".into())
}
}
fn input_to_text(input: &Input) -> Result<String> {
match input {
Input::Text(s) => Ok(s.clone()),
Input::Structured(parts) => {
let parts = parts
.iter()
.map(|p| match p {
InputPart::Text(t) => Ok(t.as_str()),
InputPart::Image(path) => Err(AgentError::Other {
message: format!(
"opencode agent mode does not support image input yet: {}",
path.display()
),
source: None,
}),
_ => Err(AgentError::Other {
message:
"opencode agent mode received an unsupported structured input part"
.into(),
source: None,
}),
})
.collect::<Result<Vec<_>>>()?;
Ok(parts.join("\n\n"))
}
_ => Err(AgentError::Other {
message: "opencode agent mode received an unsupported input shape".into(),
source: None,
}),
}
}
#[derive(serde::Deserialize)]
struct SessionInfo {
id: String,
}
#[derive(serde::Deserialize)]
struct SseFrame {
#[serde(rename = "type")]
event_type: String,
properties: serde_json::Value,
}