diff --git a/crates/config/src/schema.rs b/crates/config/src/schema.rs index 828d1fb957..3019af6c8d 100644 --- a/crates/config/src/schema.rs +++ b/crates/config/src/schema.rs @@ -374,6 +374,10 @@ pub struct ExternalAgentConfig { #[serde(default)] pub args: Vec, #[serde(default)] + pub models: Vec, + #[serde(default)] + pub efforts: Vec, + #[serde(default)] pub env: HashMap, pub working_dir: Option, pub timeout_secs: Option, diff --git a/crates/config/src/schema/tests.rs b/crates/config/src/schema/tests.rs index 7a8501c663..43ac740bcd 100644 --- a/crates/config/src/schema/tests.rs +++ b/crates/config/src/schema/tests.rs @@ -971,6 +971,23 @@ enabled = true assert_eq!(entry.wire_api, WireApi::ChatCompletions); } +#[test] +fn external_agent_models_from_toml() { + let toml_str = r#" +[external_agents] +enabled = true + +[external_agents.agents.claude-code] +binary = "claude" +models = ["claude-opus-4-8", "claude-sonnet-4-6"] +efforts = ["high", "xhigh"] +"#; + let config: MoltisConfig = toml::from_str(toml_str).unwrap(); + let entry = config.external_agents.agents.get("claude-code").unwrap(); + assert_eq!(entry.models, vec!["claude-opus-4-8", "claude-sonnet-4-6"]); + assert_eq!(entry.efforts, vec!["high", "xhigh"]); +} + #[test] fn provider_entry_wire_api_skip_serializing_default() { let entry = ProviderEntry::default(); diff --git a/crates/config/src/template.rs b/crates/config/src/template.rs index 8320fd737f..4d71973b8b 100644 --- a/crates/config/src/template.rs +++ b/crates/config/src/template.rs @@ -693,6 +693,8 @@ port = {port} # Port number (auto-generated for this i # [external_agents.agents.claude-code] # binary = "claude" # Override binary path (default: look up on $PATH) # args = ["-p", "--output-format", "json"] +# models = ["claude-opus-4-8", "claude-sonnet-4-6"] # Optional model choices shown in /model +# efforts = ["high", "xhigh"] # Optional effort choices shown in /model # working_dir = "." # Override working directory # timeout_secs = 300 # Session timeout # use_tmux = false # Force tmux backend (vs direct PTY) @@ -702,6 +704,8 @@ port = {port} # Port number (auto-generated for this i # [external_agents.agents.codex] # binary = "codex" # args = ["app-server"] +# models = ["gpt-5.5", "gpt-5.4"] +# efforts = ["medium", "high", "xhigh"] # [external_agents.agents.acp] # binary = "/path/to/acp-agent" diff --git a/crates/config/src/validate/schema_map.rs b/crates/config/src/validate/schema_map.rs index b5db4d5f99..85f7a963b1 100644 --- a/crates/config/src/validate/schema_map.rs +++ b/crates/config/src/validate/schema_map.rs @@ -57,6 +57,8 @@ pub(super) fn build_schema_map() -> KnownKeys { Struct(HashMap::from([ ("binary", Leaf), ("args", Array(Box::new(Leaf))), + ("models", Array(Box::new(Leaf))), + ("efforts", Array(Box::new(Leaf))), ("env", Map(Box::new(Leaf))), ("working_dir", Leaf), ("timeout_secs", Leaf), diff --git a/crates/config/src/validate/tests/agents.rs b/crates/config/src/validate/tests/agents.rs index 3575625c3b..74b47304e2 100644 --- a/crates/config/src/validate/tests/agents.rs +++ b/crates/config/src/validate/tests/agents.rs @@ -347,9 +347,13 @@ enabled = true [external_agents.agents.claude-code] binary = "claude" +models = ["claude-opus-4-8", "claude-sonnet-4-6"] +efforts = ["high", "xhigh"] [external_agents.agents.codex] binary = "codex" +models = ["gpt-5.5", "gpt-5.4"] +efforts = ["medium", "high", "xhigh"] "#; let result = validate_toml_str(toml); let warning = result diff --git a/crates/external-agents/src/runtimes/claude_code.rs b/crates/external-agents/src/runtimes/claude_code.rs index c8bd4f3e32..5f9ffd9b55 100644 --- a/crates/external-agents/src/runtimes/claude_code.rs +++ b/crates/external-agents/src/runtimes/claude_code.rs @@ -74,6 +74,8 @@ impl ExternalAgentTransport for ClaudeCodeTransport { spec.working_dir.clone(), spec.timeout_secs, spec.external_session_id.clone(), + spec.model.clone(), + spec.effort.clone(), ))) } } @@ -85,6 +87,8 @@ struct ClaudeCodeSession { working_dir: Option, timeout: Duration, session_id: Option, + model: Option, + effort: Option, status: ExternalAgentStatus, } @@ -96,6 +100,8 @@ impl ClaudeCodeSession { working_dir: Option, timeout_secs: Option, session_id: Option, + model: Option, + effort: Option, ) -> Self { Self { binary, @@ -104,12 +110,26 @@ impl ClaudeCodeSession { working_dir, timeout: Duration::from_secs(timeout_secs.unwrap_or(300)), session_id, + model, + effort, status: ExternalAgentStatus::Idle, } } fn args_for_turn(&self) -> Vec { let mut args = self.base_args.clone(); + if let Some(model) = &self.model + && !has_model_arg(&args) + { + args.push("--model".to_string()); + args.push(model.clone()); + } + if let Some(effort) = &self.effort + && !has_effort_arg(&args) + { + args.push("--effort".to_string()); + args.push(effort.clone()); + } if let Some(session_id) = &self.session_id && !has_resume_arg(&args) { @@ -221,6 +241,16 @@ fn has_resume_arg(args: &[String]) -> bool { .any(|arg| matches!(arg.as_str(), "--resume" | "-r" | "--continue" | "-c")) } +fn has_model_arg(args: &[String]) -> bool { + args.iter() + .any(|arg| matches!(arg.as_str(), "--model" | "-m") || arg.starts_with("--model=")) +} + +fn has_effort_arg(args: &[String]) -> bool { + args.iter() + .any(|arg| arg == "--effort" || arg.starts_with("--effort=")) +} + #[cfg(test)] mod tests { use { @@ -256,6 +286,8 @@ mod tests { None, None, None, + None, + None, ); assert!(!session.args_for_turn().iter().any(|arg| arg == "--resume")); session.session_id = Some("sid".to_string()); @@ -277,6 +309,8 @@ mod tests { None, None, Some("persisted".to_string()), + None, + None, ); assert_eq!(session.external_session_id(), Some("persisted")); @@ -314,6 +348,8 @@ printf '%s\n' '{"result":"ok","session_id":"sid-1"}' None, Some(5), None, + None, + None, ); let first = session @@ -338,4 +374,76 @@ printf '%s\n' '{"result":"ok","session_id":"sid-1"}' fs::remove_dir_all(dir)?; Ok(()) } + + #[test] + fn adds_model_arg_when_model_is_selected() { + let session = ClaudeCodeSession::new( + "claude".to_string(), + vec![ + "-p".to_string(), + "--output-format".to_string(), + "json".to_string(), + ], + HashMap::new(), + None, + None, + Some("sid".to_string()), + Some("opus".to_string()), + Some("xhigh".to_string()), + ); + + assert_eq!(session.args_for_turn(), vec![ + "-p", + "--output-format", + "json", + "--model", + "opus", + "--effort", + "xhigh", + "--resume", + "sid" + ]); + } + + #[test] + fn has_model_arg_detects_equals_form() { + let args = vec!["--model=opus".to_string()]; + assert!(has_model_arg(&args)); + } + + #[test] + fn has_model_arg_detects_flag_form() { + let args = vec!["--model".to_string(), "opus".to_string()]; + assert!(has_model_arg(&args)); + } + + #[test] + fn has_model_arg_detects_short_form() { + let args = vec!["-m".to_string(), "opus".to_string()]; + assert!(has_model_arg(&args)); + } + + #[test] + fn has_model_arg_false_when_absent() { + let args = vec!["-p".to_string(), "--output-format".to_string()]; + assert!(!has_model_arg(&args)); + } + + #[test] + fn has_effort_arg_detects_equals_form() { + let args = vec!["--effort=high".to_string()]; + assert!(has_effort_arg(&args)); + } + + #[test] + fn has_effort_arg_detects_flag_form() { + let args = vec!["--effort".to_string(), "high".to_string()]; + assert!(has_effort_arg(&args)); + } + + #[test] + fn has_effort_arg_false_when_absent() { + let args = vec!["-p".to_string()]; + assert!(!has_effort_arg(&args)); + } } diff --git a/crates/external-agents/src/runtimes/codex.rs b/crates/external-agents/src/runtimes/codex.rs index 3ad7d6e492..24f7973288 100644 --- a/crates/external-agents/src/runtimes/codex.rs +++ b/crates/external-agents/src/runtimes/codex.rs @@ -70,6 +70,7 @@ impl ExternalAgentTransport for CodexTransport { } else { spec.args.clone() }; + let args = args_with_model_and_effort(args, spec.model.as_deref(), spec.effort.as_deref()); Ok(Box::new( CodexAppServerSession::start( binary, @@ -382,6 +383,75 @@ fn extract_usage(value: &Value) -> Option { }) } +fn args_with_model_and_effort( + mut args: Vec, + model: Option<&str>, + effort: Option<&str>, +) -> Vec { + let insert_at = args.iter().position(|arg| arg == "app-server").unwrap_or(0); + let mut inserts = Vec::new(); + if let Some(model) = model + && !has_model_arg(&args) + { + inserts.extend(["--model".to_string(), model.to_string()]); + } + if let Some(effort) = effort + && !has_effort_arg(&args) + { + inserts.extend([ + "-c".to_string(), + format!("model_reasoning_effort=\"{effort}\""), + ]); + } + for (offset, arg) in inserts.into_iter().enumerate() { + args.insert(insert_at + offset, arg); + } + args +} + +fn has_model_arg(args: &[String]) -> bool { + let mut iter = args.iter().peekable(); + while let Some(arg) = iter.next() { + if matches!(arg.as_str(), "--model" | "-m") || arg.starts_with("--model=") { + return true; + } + if matches!(arg.as_str(), "--config" | "-c") + && iter + .peek() + .is_some_and(|next| next.trim_start().starts_with("model=")) + { + return true; + } + if arg + .strip_prefix("--config=") + .is_some_and(|value| value.trim_start().starts_with("model=")) + { + return true; + } + } + false +} + +fn has_effort_arg(args: &[String]) -> bool { + let mut iter = args.iter().peekable(); + while let Some(arg) = iter.next() { + if matches!(arg.as_str(), "--config" | "-c") + && iter + .peek() + .is_some_and(|next| next.trim_start().starts_with("model_reasoning_effort=")) + { + return true; + } + if arg + .strip_prefix("--config=") + .is_some_and(|value| value.trim_start().starts_with("model_reasoning_effort=")) + { + return true; + } + } + false +} + fn token_count_field(value: &Value, fields: &[&str]) -> Option { fields.iter().find_map(|field| { value @@ -443,6 +513,44 @@ mod tests { assert_eq!(camel.output_tokens, 55); } + #[test] + fn inserts_model_and_effort_before_app_server_command() { + assert_eq!( + args_with_model_and_effort( + vec!["app-server".to_string()], + Some("gpt-5.5"), + Some("xhigh") + ), + vec![ + "--model", + "gpt-5.5", + "-c", + "model_reasoning_effort=\"xhigh\"", + "app-server" + ] + ); + assert_eq!( + args_with_model_and_effort( + vec![ + "--model".to_string(), + "configured".to_string(), + "-c".to_string(), + "model_reasoning_effort=\"high\"".to_string(), + "app-server".to_string() + ], + Some("ignored"), + Some("ignored"), + ), + vec![ + "--model", + "configured", + "-c", + "model_reasoning_effort=\"high\"", + "app-server" + ] + ); + } + #[tokio::test] async fn codex_session_reuses_thread_for_multiple_prompts() -> anyhow::Result<()> { let unique = SystemTime::now() @@ -512,4 +620,46 @@ done fs::remove_dir_all(dir)?; Ok(()) } + + #[test] + fn has_model_arg_detects_equals_form() { + let args = vec!["--model=gpt-4".to_string()]; + assert!(has_model_arg(&args)); + } + + #[test] + fn has_model_arg_detects_flag_form() { + let args = vec!["--model".to_string(), "gpt-4".to_string()]; + assert!(has_model_arg(&args)); + } + + #[test] + fn has_model_arg_detects_config_form() { + let args = vec!["-c".to_string(), "model=gpt-4".to_string()]; + assert!(has_model_arg(&args)); + } + + #[test] + fn has_model_arg_false_when_absent() { + let args = vec!["--some-flag".to_string()]; + assert!(!has_model_arg(&args)); + } + + #[test] + fn has_effort_arg_detects_config_form() { + let args = vec!["-c".to_string(), "model_reasoning_effort=high".to_string()]; + assert!(has_effort_arg(&args)); + } + + #[test] + fn has_effort_arg_detects_config_equals_form() { + let args = vec!["--config=model_reasoning_effort=high".to_string()]; + assert!(has_effort_arg(&args)); + } + + #[test] + fn has_effort_arg_false_when_absent() { + let args = vec!["--model".to_string()]; + assert!(!has_effort_arg(&args)); + } } diff --git a/crates/external-agents/src/types.rs b/crates/external-agents/src/types.rs index 088108c2d2..2bd8b64956 100644 --- a/crates/external-agents/src/types.rs +++ b/crates/external-agents/src/types.rs @@ -74,6 +74,8 @@ pub struct ExternalAgentSpec { pub kind: AgentTransportKind, pub session_key: Option, pub external_session_id: Option, + pub model: Option, + pub effort: Option, pub binary: Option, pub args: Vec, pub env: HashMap, @@ -90,6 +92,8 @@ impl ExternalAgentSpec { kind, session_key: None, external_session_id: None, + model: None, + effort: None, binary: None, args: Vec::new(), env: HashMap::new(), diff --git a/crates/gateway/src/channel_events/commands/control_handlers.rs b/crates/gateway/src/channel_events/commands/control_handlers.rs index bbb6005539..57abac6215 100644 --- a/crates/gateway/src/channel_events/commands/control_handlers.rs +++ b/crates/gateway/src/channel_events/commands/control_handlers.rs @@ -3,12 +3,16 @@ use std::sync::Arc; use { moltis_channels::{ChannelReplyTarget, Error as ChannelError, Result as ChannelResult}, moltis_config::ModePreset, + moltis_external_agents::AgentTransportKind, moltis_sessions::metadata::SqliteSessionMetadata, moltis_tools::image_cache::ImageBuilder, }; use crate::{ broadcast::{BroadcastOpts, broadcast}, + external_agents::{ + ExternalAgentModelSelection, external_agent_model_id, parse_external_agent_model_id, + }, state::GatewayState, }; @@ -20,6 +24,8 @@ use super::{ formatting::{format_model_list, unique_providers}, }; +const EXTERNAL_AGENT_PROVIDER: &str = "external-agent"; + // ── Control command handlers ───────────────────────────────────── fn sorted_mode_presets(config: &moltis_config::MoltisConfig) -> Vec<(String, ModePreset)> { @@ -357,25 +363,41 @@ pub(in crate::channel_events) async fn handle_model( let models = models_val .as_array() .ok_or_else(|| ChannelError::invalid_input("bad model list"))?; + let mut choices = models.clone(); + choices.extend(external_agent_model_entries(state).await?); - let current_model = { + let (current_model, current_external_agent) = { let entry = session_metadata.get(session_key).await; - entry.and_then(|e| e.model.clone()) + ( + entry.as_ref().and_then(|e| e.model.clone()), + entry.and_then(|e| e.external_agent_kind), + ) + }; + let current_choice = if let Some(kind) = current_external_agent { + current_model + .as_deref() + .and_then(|model| { + parse_external_agent_model_id(model) + .and_then(|choice| (choice.kind == kind.as_str()).then(|| model.to_string())) + }) + .or_else(|| Some(external_agent_model_id(kind.as_str(), None, None))) + } else { + current_model }; if args.is_empty() { // List unique providers (sorted, deduplicated). - let providers = unique_providers(models); + let providers = unique_providers(&choices); if providers.len() <= 1 { // Single provider -- list models directly. - return Ok(format_model_list(models, current_model.as_deref(), None)); + return Ok(format_model_list(&choices, current_choice.as_deref(), None)); } // Multiple providers -- list them for selection. // Prefix with "providers:" so Telegram handler knows. - let current_provider = current_model.as_deref().and_then(|cm| { - models.iter().find_map(|m| { + let current_provider = current_choice.as_deref().and_then(|cm| { + choices.iter().find_map(|m| { let id = m.get("id").and_then(|v| v.as_str())?; if id == cm { m.get("provider").and_then(|v| v.as_str()).map(String::from) @@ -386,7 +408,7 @@ pub(in crate::channel_events) async fn handle_model( }); let mut lines = vec!["providers:".to_string()]; for (i, p) in providers.iter().enumerate() { - let count = models + let count = choices .iter() .filter(|m| m.get("provider").and_then(|v| v.as_str()) == Some(p)) .count(); @@ -401,8 +423,8 @@ pub(in crate::channel_events) async fn handle_model( } else if let Some(provider) = args.strip_prefix("provider:") { // List models for a specific provider. Ok(format_model_list( - models, - current_model.as_deref(), + &choices, + current_choice.as_deref(), Some(provider), )) } else { @@ -410,13 +432,13 @@ pub(in crate::channel_events) async fn handle_model( let n: usize = args .parse() .map_err(|_| ChannelError::invalid_input("usage: /model [number]"))?; - if n == 0 || n > models.len() { + if n == 0 || n > choices.len() { return Err(ChannelError::invalid_input(format!( "invalid model number. Use 1\u{2013}{}.", - models.len() + choices.len() ))); } - let chosen = &models[n - 1]; + let chosen = &choices[n - 1]; let model_id = chosen .get("id") .and_then(|v| v.as_str()) @@ -425,6 +447,48 @@ pub(in crate::channel_events) async fn handle_model( .get("displayName") .and_then(|v| v.as_str()) .unwrap_or(model_id); + if let Some(selection) = chosen + .get("externalAgentKind") + .and_then(|v| v.as_str()) + .map(|kind| ExternalAgentModelSelection { + kind, + model: chosen.get("externalAgentModel").and_then(|v| v.as_str()), + effort: chosen.get("externalAgentEffort").and_then(|v| v.as_str()), + }) + .or_else(|| parse_external_agent_model_id(model_id)) + { + let kind = selection + .kind + .parse::() + .map_err(ChannelError::invalid_input)?; + state + .services + .external_agent + .bind(serde_json::json!({ + "sessionKey": session_key, + "kind": kind.as_str(), + "model": selection.model, + "effort": selection.effort, + })) + .await + .map_err(ChannelError::unavailable)?; + + broadcast( + state, + "session", + serde_json::json!({ + "kind": "patched", + "sessionKey": session_key, + }), + BroadcastOpts { + drop_if_slow: true, + ..Default::default() + }, + ) + .await; + + return Ok(format!("Backend switched to: {display}")); + } let patch_res = state .services @@ -439,6 +503,7 @@ pub(in crate::channel_events) async fn handle_model( .get("version") .and_then(|v| v.as_u64()) .unwrap_or(0); + unbind_external_agent_if_bound(state, session_key).await?; broadcast( state, @@ -459,6 +524,125 @@ pub(in crate::channel_events) async fn handle_model( } } +async fn external_agent_model_entries( + state: &GatewayState, +) -> ChannelResult> { + let agents = state + .services + .external_agent + .list() + .await + .map_err(ChannelError::unavailable)?; + let Some(agents) = agents.as_array() else { + return Ok(Vec::new()); + }; + Ok(agents + .iter() + .filter(|agent| { + agent + .get("installed") + .and_then(|value| value.as_bool()) + .unwrap_or(false) + }) + .filter_map(|agent| { + let kind = agent + .get("kind") + .and_then(|value| value.as_str()) + .or_else(|| agent.get("name").and_then(|value| value.as_str()))?; + let mut entries = vec![serde_json::json!({ + "id": external_agent_model_id(kind, None, None), + "provider": external_agent_provider(kind), + "displayName": format!("{} (default)", external_agent_display_name(kind)), + "externalAgentKind": kind, + })]; + let models = agent + .get("models") + .and_then(|value| value.as_array()) + .into_iter() + .flatten() + .filter_map(|value| value.as_str()) + .filter(|model| !model.trim().is_empty()) + .collect::>(); + let efforts = agent + .get("efforts") + .and_then(|value| value.as_array()) + .into_iter() + .flatten() + .filter_map(|value| value.as_str()) + .filter(|effort| !effort.trim().is_empty()) + .collect::>(); + for model in models { + if efforts.is_empty() { + entries.push(serde_json::json!({ + "id": external_agent_model_id(kind, Some(model), None), + "provider": external_agent_provider(kind), + "displayName": format!("{}: {model}", external_agent_display_name(kind)), + "externalAgentKind": kind, + "externalAgentModel": model, + })); + } else { + for effort in &efforts { + entries.push(serde_json::json!({ + "id": external_agent_model_id(kind, Some(model), Some(effort)), + "provider": external_agent_provider(kind), + "displayName": format!( + "{}: {model} ({effort})", + external_agent_display_name(kind) + ), + "externalAgentKind": kind, + "externalAgentModel": model, + "externalAgentEffort": effort, + })); + } + } + } + Some(entries) + }) + .flatten() + .collect()) +} + +async fn unbind_external_agent_if_bound( + state: &GatewayState, + session_key: &str, +) -> ChannelResult<()> { + let status = state + .services + .external_agent + .status(serde_json::json!({ "sessionKey": session_key })) + .await + .map_err(ChannelError::unavailable)?; + if status + .get("bound") + .and_then(|value| value.as_bool()) + .unwrap_or(false) + { + state + .services + .external_agent + .unbind(serde_json::json!({ "sessionKey": session_key })) + .await + .map_err(ChannelError::unavailable)?; + } + Ok(()) +} + + +fn external_agent_provider(kind: &str) -> String { + format!("{EXTERNAL_AGENT_PROVIDER}/{kind}") +} + +fn external_agent_display_name(kind: &str) -> &'static str { + match kind { + "claude-code" => "Claude Code CLI", + "codex" => "Codex CLI", + "acp" => "ACP Agent", + "opencode" => "OpenCode CLI", + "pi-agent" => "Pi Agent", + _ => "External Agent", + } +} + pub(in crate::channel_events) async fn handle_sandbox( state: &Arc, session_metadata: &SqliteSessionMetadata, diff --git a/crates/gateway/src/external_agents.rs b/crates/gateway/src/external_agents.rs index eb11f0d6e2..11f1c8b992 100644 --- a/crates/gateway/src/external_agents.rs +++ b/crates/gateway/src/external_agents.rs @@ -15,7 +15,7 @@ use { moltis_sessions::{MessageContent, PersistedMessage}, serde_json::Value, tokio::sync::Mutex, - tracing::warn, + tracing::{info, warn}, }; use moltis_tools::approval::{ApprovalDecision, ApprovalManager}; @@ -166,14 +166,14 @@ impl GatewayExternalAgentService { return Ok(Arc::clone(&entry.session)); } } - let spec = self.spec_for_kind(kind)?; + let entry = self.session_metadata.get(session_key).await; + let selected = entry + .as_ref() + .and_then(|entry| selected_external_agent(entry.model.as_deref(), kind)); + let spec = self.spec_for_kind(kind, selected)?; let mut spec = spec; spec.session_key = Some(session_key.to_string()); - spec.external_session_id = self - .session_metadata - .get(session_key) - .await - .and_then(|entry| entry.external_session_id); + spec.external_session_id = entry.and_then(|entry| entry.external_session_id); let session = Arc::new(Mutex::new(self.registry.start_session(&spec).await?)); live_sessions.insert(key, LiveSessionEntry { session: Arc::clone(&session), @@ -223,11 +223,19 @@ impl GatewayExternalAgentService { } } - fn spec_for_kind(&self, kind: AgentTransportKind) -> anyhow::Result { + fn spec_for_kind( + &self, + kind: AgentTransportKind, + selected: Option, + ) -> anyhow::Result { if !self.config.enabled { anyhow::bail!("external agents are disabled") } let mut spec = ExternalAgentSpec::new(kind); + if let Some(selected) = selected { + spec.model = selected.model; + spec.effort = selected.effort; + } if let Some(agent_config) = self.config.agents.get(kind.as_str()) { spec.binary = agent_config.binary.clone(); spec.args = agent_config.args.clone(); @@ -343,8 +351,30 @@ impl ExternalAgentService for GatewayExternalAgentService { if !self.config.enabled { return Ok(serde_json::json!([])); } - Ok(serde_json::to_value(self.registry.list_agents().await) - .unwrap_or_else(|_| serde_json::json!([]))) + let mut agents = serde_json::to_value(self.registry.list_agents().await) + .unwrap_or_else(|_| serde_json::json!([])); + if let Some(agents) = agents.as_array_mut() { + for agent in agents { + let Some(kind) = agent.get("kind").and_then(Value::as_str) else { + continue; + }; + let models = self + .config + .agents + .get(kind) + .map(|config| config.models.clone()) + .unwrap_or_default(); + let efforts = self + .config + .agents + .get(kind) + .map(|config| config.efforts.clone()) + .unwrap_or_default(); + agent["models"] = serde_json::json!(models); + agent["efforts"] = serde_json::json!(efforts); + } + } + Ok(agents) } async fn bind(&self, params: Value) -> ServiceResult { @@ -362,6 +392,16 @@ impl ExternalAgentService for GatewayExternalAgentService { .ok_or_else(|| "missing kind".to_string())? .parse::() .map_err(|error| error.to_string())?; + let model = params + .get("model") + .and_then(|value| value.as_str()) + .filter(|value| !value.trim().is_empty()) + .map(ToOwned::to_owned); + let effort = params + .get("effort") + .and_then(|value| value.as_str()) + .filter(|value| !value.trim().is_empty()) + .map(ToOwned::to_owned); if !self.registry.has_kind(kind) { return Err(format!("external agent kind is not registered: {kind}").into()); } @@ -370,7 +410,19 @@ impl ExternalAgentService for GatewayExternalAgentService { self.session_metadata .set_external_agent(session_key, Some(kind), None) .await; - Ok(serde_json::json!({ "ok": true, "sessionKey": session_key, "kind": kind.as_str() })) + let selected_model_id = + external_agent_model_id(kind.as_str(), model.as_deref(), effort.as_deref()); + self.session_metadata + .set_model(session_key, Some(selected_model_id.clone())) + .await; + Ok(serde_json::json!({ + "ok": true, + "sessionKey": session_key, + "kind": kind.as_str(), + "model": model, + "effort": effort, + "modelId": selected_model_id, + })) } async fn unbind(&self, params: Value) -> ServiceResult { @@ -398,11 +450,77 @@ impl ExternalAgentService for GatewayExternalAgentService { "bound": kind.is_some(), "sessionKey": session_key, "kind": kind.map(|kind| kind.as_str()), + "model": entry.as_ref().and_then(|entry| { + let kind = entry.external_agent_kind?; + selected_external_agent(entry.model.as_deref(), kind).and_then(|selected| selected.model) + }), + "effort": entry.as_ref().and_then(|entry| { + let kind = entry.external_agent_kind?; + selected_external_agent(entry.model.as_deref(), kind).and_then(|selected| selected.effort) + }), "externalSessionId": entry.and_then(|entry| entry.external_session_id), })) } } +pub(crate) const EXTERNAL_AGENT_MODEL_PREFIX: &str = "external-agent::"; + +#[derive(Debug, Clone, PartialEq, Eq)] +struct SelectedExternalAgent { + model: Option, + effort: Option, +} + +pub(crate) fn external_agent_model_id( + kind: &str, + model: Option<&str>, + effort: Option<&str>, +) -> String { + match (model, effort) { + (Some(model), Some(effort)) => { + format!("{EXTERNAL_AGENT_MODEL_PREFIX}{kind}::{model}::{effort}") + }, + (Some(model), None) => format!("{EXTERNAL_AGENT_MODEL_PREFIX}{kind}::{model}"), + (None, Some(effort)) => format!("{EXTERNAL_AGENT_MODEL_PREFIX}{kind}::default::{effort}"), + (None, None) => format!("{EXTERNAL_AGENT_MODEL_PREFIX}{kind}"), + } +} + +pub(crate) struct ExternalAgentModelSelection<'a> { + pub kind: &'a str, + pub model: Option<&'a str>, + pub effort: Option<&'a str>, +} + +pub(crate) fn parse_external_agent_model_id( + model_id: &str, +) -> Option> { + let suffix = model_id.strip_prefix(EXTERNAL_AGENT_MODEL_PREFIX)?; + let mut parts = suffix.split("::"); + let kind = parts.next()?; + let model = parts.next(); + let effort = parts.next(); + Some(ExternalAgentModelSelection { + kind, + model: model.filter(|m| *m != "default"), + effort, + }) +} + +fn selected_external_agent( + selection_id: Option<&str>, + kind: AgentTransportKind, +) -> Option { + let sel = parse_external_agent_model_id(selection_id?)?; + if sel.kind != kind.as_str() { + return None; + } + Some(SelectedExternalAgent { + model: sel.model.map(ToOwned::to_owned), + effort: sel.effort.map(ToOwned::to_owned), + }) +} + pub struct ExternalAgentChatService { inner: Arc, external_agents: Arc, @@ -450,6 +568,22 @@ impl ExternalAgentChatService { .and_then(|value| value.as_str()) .ok_or_else(|| "external agents currently require text input".to_string())? .to_string(); + let channel_reply_target = params + .get("_channel_reply_target") + .cloned() + .and_then(|value| { + match serde_json::from_value::(value) { + Ok(target) => Some(target), + Err(error) => { + warn!( + session = %session_key, + %error, + "ignoring invalid external-agent channel reply target" + ); + None + }, + } + }); let seq = params.get("_seq").and_then(|value| value.as_u64()); let run_id = uuid::Uuid::new_v4().to_string(); let created_at = now_ms(); @@ -475,6 +609,18 @@ impl ExternalAgentChatService { self.session_metadata .touch(&session_key, history.len() as u32) .await; + let selected = self + .session_metadata + .get(&session_key) + .await + .and_then(|entry| selected_external_agent(entry.model.as_deref(), kind)); + let selected_model = selected + .as_ref() + .and_then(|selected| selected.model.as_deref()); + let selected_effort = selected + .as_ref() + .and_then(|selected| selected.effort.as_deref()); + let message_model = selected_model.unwrap_or_else(|| kind.as_str()); crate::broadcast::broadcast( &self.state, @@ -483,7 +629,7 @@ impl ExternalAgentChatService { "runId": run_id, "sessionKey": session_key, "state": "running", - "model": kind.as_str(), + "model": message_model, "provider": "external-agent", "seq": seq, }), @@ -493,6 +639,16 @@ impl ExternalAgentChatService { let context = context_from_history(&history); let start = std::time::Instant::now(); + info!( + session = %session_key, + kind = kind.as_str(), + model = message_model, + effort = selected_effort.unwrap_or("default"), + run_id, + text_len = text.len(), + channel_reply = channel_reply_target.is_some(), + "external-agent turn starting" + ); let live_session = self .external_agents .session_for_binding(&session_key, kind) @@ -575,7 +731,7 @@ impl ExternalAgentChatService { let assistant_msg = PersistedMessage::Assistant { content: assistant_text.clone(), created_at: Some(now_ms()), - model: Some(kind.as_str().to_string()), + model: Some(message_model.to_string()), provider: Some("external-agent".to_string()), input_tokens: token_usage.as_ref().map(|usage| usage.input_tokens), output_tokens: token_usage.as_ref().map(|usage| usage.output_tokens), @@ -609,7 +765,7 @@ impl ExternalAgentChatService { "sessionKey": session_key, "state": "final", "text": assistant_text, - "model": kind.as_str(), + "model": message_model, "provider": "external-agent", "inputTokens": token_usage.as_ref().map(|usage| usage.input_tokens).unwrap_or(0), "outputTokens": token_usage.as_ref().map(|usage| usage.output_tokens).unwrap_or(0), @@ -621,10 +777,70 @@ impl ExternalAgentChatService { BroadcastOpts::default(), ) .await; + deliver_external_agent_channel_reply( + &self.state, + channel_reply_target, + &assistant_text, + &session_key, + ) + .await; + info!( + session = %session_key, + kind = kind.as_str(), + model = message_model, + effort = selected_effort.unwrap_or("default"), + run_id, + text_len = assistant_text.len(), + duration_ms, + "external-agent turn completed" + ); Ok(serde_json::json!({ "ok": true, "runId": run_id })) } } +async fn deliver_external_agent_channel_reply( + state: &GatewayState, + target: Option, + text: &str, + session_key: &str, +) { + let Some(target) = target else { + return; + }; + if text.trim().is_empty() { + info!( + session_key, + account_id = target.account_id, + chat_id = target.chat_id, + "external-agent channel reply skipped: empty response text" + ); + return; + } + let Some(outbound) = state.services.channel_outbound_arc() else { + warn!( + session_key, + account_id = target.account_id, + chat_id = target.chat_id, + "external-agent channel reply skipped: outbound unavailable" + ); + return; + }; + let to = target.outbound_to().into_owned(); + if let Err(error) = outbound + .send_text(&target.account_id, &to, text, target.message_id.as_deref()) + .await + { + warn!( + session_key, + account_id = target.account_id, + chat_id = target.chat_id, + thread_id = target.thread_id.as_deref().unwrap_or("-"), + %error, + "external-agent channel reply failed" + ); + } +} + #[async_trait] impl ChatService for ExternalAgentChatService { async fn send(&self, params: Value) -> ServiceResult { @@ -816,6 +1032,8 @@ mod tests { struct FakeAgentState { starts: std::sync::atomic::AtomicUsize, prompts: std::sync::Mutex>, + models: std::sync::Mutex>>, + efforts: std::sync::Mutex>>, shutdowns: std::sync::atomic::AtomicUsize, } @@ -839,9 +1057,19 @@ mod tests { async fn start_session( &self, - _spec: &ExternalAgentSpec, + spec: &ExternalAgentSpec, ) -> anyhow::Result> { let start_index = self.state.starts.fetch_add(1, Ordering::SeqCst) + 1; + self.state + .models + .lock() + .unwrap_or_else(|error| error.into_inner()) + .push(spec.model.clone()); + self.state + .efforts + .lock() + .unwrap_or_else(|error| error.into_inner()) + .push(spec.effort.clone()); Ok(Box::new(FakeSession { state: Arc::clone(&self.state), external_session_id: format!("fake-session-{start_index}"), @@ -946,13 +1174,17 @@ mod tests { } fn test_gateway_state() -> Arc { + test_gateway_state_with_services(GatewayServices::noop()) + } + + fn test_gateway_state_with_services(services: GatewayServices) -> Arc { GatewayState::new( ResolvedAuth { mode: AuthMode::Token, token: None, password: None, }, - GatewayServices::noop(), + services, ) } @@ -970,6 +1202,55 @@ mod tests { ) } + #[derive(Default)] + struct RecordingOutbound { + messages: Mutex)>>, + } + + #[async_trait] + impl moltis_channels::ChannelOutbound for RecordingOutbound { + async fn send_text( + &self, + account_id: &str, + to: &str, + text: &str, + reply_to: Option<&str>, + ) -> moltis_channels::Result<()> { + self.messages.lock().await.push(( + account_id.to_string(), + to.to_string(), + text.to_string(), + reply_to.map(ToOwned::to_owned), + )); + Ok(()) + } + + async fn send_media( + &self, + _account_id: &str, + _to: &str, + _payload: &moltis_common::types::ReplyPayload, + _reply_to: Option<&str>, + ) -> moltis_channels::Result<()> { + Ok(()) + } + } + + async fn test_chat_service_with_state( + external_agents: Arc, + metadata: Arc, + session_store: Arc, + state: Arc, + ) -> ExternalAgentChatService { + ExternalAgentChatService::new( + Arc::new(NoopChatService), + external_agents, + state, + session_store, + metadata, + ) + } + #[tokio::test] async fn bind_unbind_and_status_update_metadata() { let metadata = Arc::new(SqliteSessionMetadata::new(sqlite_pool().await)); @@ -977,10 +1258,21 @@ mod tests { let service = fake_external_agents(Arc::clone(&metadata), agent_state); let bound = service - .bind(serde_json::json!({ "sessionKey": "main", "kind": "codex" })) + .bind(serde_json::json!({ + "sessionKey": "main", + "kind": "codex", + "model": "gpt-5.2-codex", + "effort": "xhigh", + })) .await .expect("bind external agent"); assert_eq!(bound["kind"], "codex"); + assert_eq!(bound["model"], "gpt-5.2-codex"); + assert_eq!(bound["effort"], "xhigh"); + assert_eq!( + bound["modelId"], + "external-agent::codex::gpt-5.2-codex::xhigh" + ); let status = service .status(serde_json::json!({ "sessionKey": "main" })) @@ -988,6 +1280,13 @@ mod tests { .expect("status"); assert_eq!(status["bound"], true); assert_eq!(status["kind"], "codex"); + assert_eq!(status["model"], "gpt-5.2-codex"); + assert_eq!(status["effort"], "xhigh"); + let entry = metadata.get("main").await.expect("session entry"); + assert_eq!( + entry.model.as_deref(), + Some("external-agent::codex::gpt-5.2-codex::xhigh") + ); service .unbind(serde_json::json!({ "sessionKey": "main" })) @@ -1001,6 +1300,47 @@ mod tests { assert!(status["kind"].is_null()); } + #[tokio::test] + async fn selected_external_agent_model_is_passed_to_runtime() { + let dir = tempfile::tempdir().unwrap(); + let session_store = Arc::new(SessionStore::new(dir.path().to_path_buf())); + let metadata = Arc::new(SqliteSessionMetadata::new(sqlite_pool().await)); + let agent_state = Arc::new(FakeAgentState::default()); + let external_agents = fake_external_agents(Arc::clone(&metadata), Arc::clone(&agent_state)); + external_agents + .bind(serde_json::json!({ + "sessionKey": "main", + "kind": "codex", + "model": "gpt-5.2-codex", + "effort": "xhigh", + })) + .await + .expect("bind external agent"); + let chat = test_chat_service( + Arc::clone(&external_agents), + Arc::clone(&metadata), + Arc::clone(&session_store), + ) + .await; + + chat.send(serde_json::json!({ "sessionKey": "main", "text": "hello" })) + .await + .expect("send external prompt"); + + let models = agent_state + .models + .lock() + .unwrap_or_else(|error| error.into_inner()) + .clone(); + let efforts = agent_state + .efforts + .lock() + .unwrap_or_else(|error| error.into_inner()) + .clone(); + assert_eq!(models, vec![Some("gpt-5.2-codex".to_string())]); + assert_eq!(efforts, vec![Some("xhigh".to_string())]); + } + #[tokio::test] async fn list_returns_empty_when_external_agents_disabled() { let metadata = Arc::new(SqliteSessionMetadata::new(sqlite_pool().await)); @@ -1131,6 +1471,51 @@ mod tests { ); } + #[tokio::test] + async fn bound_chat_send_delivers_external_reply_to_channel_target() { + let dir = tempfile::tempdir().unwrap(); + let session_store = Arc::new(SessionStore::new(dir.path().to_path_buf())); + let metadata = Arc::new(SqliteSessionMetadata::new(sqlite_pool().await)); + let agent_state = Arc::new(FakeAgentState::default()); + let external_agents = fake_external_agents(Arc::clone(&metadata), Arc::clone(&agent_state)); + external_agents + .bind(serde_json::json!({ "sessionKey": "telegram:bot:123", "kind": "codex" })) + .await + .expect("bind external agent"); + let outbound = Arc::new(RecordingOutbound::default()); + let state = test_gateway_state_with_services( + GatewayServices::noop().with_channel_outbound(outbound.clone()), + ); + let chat = test_chat_service_with_state( + Arc::clone(&external_agents), + Arc::clone(&metadata), + Arc::clone(&session_store), + state, + ) + .await; + + chat.send(serde_json::json!({ + "sessionKey": "telegram:bot:123", + "text": "hello", + "_channel_reply_target": { + "channel_type": "telegram", + "account_id": "bot", + "chat_id": "123", + "message_id": "456", + "thread_id": null, + }, + })) + .await + .expect("send external agent turn"); + + assert_eq!(*outbound.messages.lock().await, vec![( + "bot".to_string(), + "123".to_string(), + "reply to hello".to_string(), + Some("456".to_string()), + )]); + } + #[tokio::test] async fn idle_live_external_sessions_are_evicted_before_reuse() { let metadata = Arc::new(SqliteSessionMetadata::new(sqlite_pool().await)); @@ -1284,4 +1669,57 @@ mod tests { assert_eq!(agent_state.shutdowns.load(Ordering::SeqCst), 1); assert!(external_agents.live_sessions.lock().await.is_empty()); } + + #[test] + fn model_id_round_trips_through_parse() { + let id = external_agent_model_id("claude-code", Some("opus"), Some("high")); + let parsed = parse_external_agent_model_id(&id).unwrap(); + assert_eq!(parsed.kind, "claude-code"); + assert_eq!(parsed.model, Some("opus")); + assert_eq!(parsed.effort, Some("high")); + } + + #[test] + fn model_id_kind_only() { + let id = external_agent_model_id("codex", None, None); + assert_eq!(id, "external-agent::codex"); + let parsed = parse_external_agent_model_id(&id).unwrap(); + assert_eq!(parsed.kind, "codex"); + assert!(parsed.model.is_none()); + assert!(parsed.effort.is_none()); + } + + #[test] + fn model_id_with_model_only() { + let id = external_agent_model_id("claude-code", Some("sonnet"), None); + let parsed = parse_external_agent_model_id(&id).unwrap(); + assert_eq!(parsed.kind, "claude-code"); + assert_eq!(parsed.model, Some("sonnet")); + assert!(parsed.effort.is_none()); + } + + #[test] + fn model_id_effort_only_uses_default_placeholder() { + let id = external_agent_model_id("codex", None, Some("low")); + assert!(id.contains("::default::")); + let parsed = parse_external_agent_model_id(&id).unwrap(); + assert_eq!(parsed.kind, "codex"); + assert!(parsed.model.is_none()); + assert_eq!(parsed.effort, Some("low")); + } + + #[test] + fn parse_rejects_non_external_agent_prefix() { + assert!(parse_external_agent_model_id("openai::gpt-4").is_none()); + } + + #[test] + fn selected_external_agent_filters_by_kind() { + let id = external_agent_model_id("claude-code", Some("opus"), None); + let sel = + selected_external_agent(Some(&id), AgentTransportKind::ClaudeCode).unwrap(); + assert_eq!(sel.model.as_deref(), Some("opus")); + + assert!(selected_external_agent(Some(&id), AgentTransportKind::Codex).is_none()); + } } diff --git a/docs/src/configuration-reference.md b/docs/src/configuration-reference.md index e3b4cb0146..b5e099726a 100644 --- a/docs/src/configuration-reference.md +++ b/docs/src/configuration-reference.md @@ -77,6 +77,8 @@ - [`mcp`](#mcp) - [`mcp.servers.`](#mcpserversname) - [`mcp.servers..oauth`](#mcpserversnameoauth) + - [`external_agents`](#external-agents) + - [`external_agents.agents.`](#external-agentsagentskind) - **Memory** - [`memory`](#memory) - [`memory.qmd`](#memoryqmd) @@ -701,6 +703,34 @@ Each channel account (`channels..`) is an arbitrary --- +--- + +### `external_agents` + +**Struct:** `ExternalAgentsConfig` + +| Key | Type | Default | Description | +|-----|------|---------|-------------| +| `enabled` | bool | `false` | Enable the external CLI agent bridge. | +| `agents` | map | `{}` | Per-agent configuration keyed by agent kind, such as `claude-code` or `codex`. | + + +### `external_agents.agents.` + +**Struct:** `ExternalAgentConfig` + +| Key | Type | Default | Description | +|-----|------|---------|-------------| +| `binary` | optional string | `null` | Override the CLI binary path. | +| `args` | array of strings | `[]` | Override CLI arguments for the external agent runtime. | +| `models` | array of strings | `[]` | Optional model choices shown in `/model`; selected values are passed to the external CLI runtime. | +| `efforts` | array of strings | `[]` | Optional thinking/reasoning effort choices shown in `/model`; selected values are passed to runtimes that support effort. | +| `env` | map | `{}` | Environment variables for this external agent runtime. | +| `working_dir` | optional string | `null` | Working directory for this external agent runtime. | +| `timeout_secs` | optional integer | `null` | Timeout for external agent startup or turns. | +| `use_tmux` | optional bool | `null` | Force tmux-backed runtime when supported. | + + --- ## Memory