From bc94cac8373f2138435754ba138b5c6b23ec9c06 Mon Sep 17 00:00:00 2001 From: benthecarman Date: Sat, 19 Sep 2026 01:28:39 -0500 Subject: [PATCH 1/3] Add native Claude Code delegation Let Maple tasks delegate to an installed Claude Code CLI through a Rust transport adapted from Goose. Share Codex's activity, approval, question, and cancellation controls while keeping process ownership in Maple. Cover streaming, resumption, provider isolation, permissions, failures, and process cleanup with native CLI fixtures. --- .../app/assets/icons/claude-mark.svg | 1 + apps/maple-agent/app/src/assets.rs | 2 + apps/maple-agent/app/src/ui/chat/mod.rs | 80 ++- .../maple-agent/app/src/ui/chat/navigation.rs | 24 +- apps/maple-agent/app/src/ui/chat/tests.rs | 191 ++++- .../maple-agent/app/src/ui/chat/transcript.rs | 28 +- apps/maple-agent/app/src/ui/settings.rs | 25 +- .../resources/skills/advisor/SKILL.md | 2 +- .../resources/skills/committee/SKILL.md | 2 +- .../resources/skills/handoff/SKILL.md | 6 +- .../crates/maple-agent/src/agent.rs | 1 + .../maple-agent/src/agent/developer_tools.rs | 24 +- .../src/agent/external_agents/app_server.rs | 76 +- .../src/agent/external_agents/claude.rs | 659 ++++++++++++++++++ .../src/agent/external_agents/codex.rs | 3 +- .../src/agent/external_agents/mod.rs | 290 ++++++-- .../src/agent/external_agents/tests.rs | 543 +++++++++++++-- .../external_agents/tests/claude_fixture.rs | 267 +++++++ .../maple-agent/src/agent/integrations.rs | 164 ++++- .../crates/maple-agent/src/agent/types.rs | 1 + 20 files changed, 2206 insertions(+), 183 deletions(-) create mode 100644 apps/maple-agent/app/assets/icons/claude-mark.svg create mode 100644 apps/maple-agent/crates/maple-agent/src/agent/external_agents/claude.rs create mode 100644 apps/maple-agent/crates/maple-agent/src/agent/external_agents/tests/claude_fixture.rs diff --git a/apps/maple-agent/app/assets/icons/claude-mark.svg b/apps/maple-agent/app/assets/icons/claude-mark.svg new file mode 100644 index 000000000..1beee8612 --- /dev/null +++ b/apps/maple-agent/app/assets/icons/claude-mark.svg @@ -0,0 +1 @@ +Claude \ No newline at end of file diff --git a/apps/maple-agent/app/src/assets.rs b/apps/maple-agent/app/src/assets.rs index c68107426..996d4507c 100644 --- a/apps/maple-agent/app/src/assets.rs +++ b/apps/maple-agent/app/src/assets.rs @@ -20,6 +20,8 @@ assets!( "icons/check.svg", "icons/chevron-down.svg", "icons/chevron-right.svg", + // Claude mark (simple-icons, CC0) for the Claude Code integration card. + "icons/claude-mark.svg", "icons/copy.svg", // Contrast-safe partner marks published at https://cua.ai/branding. "icons/cua-mark-black.svg", diff --git a/apps/maple-agent/app/src/ui/chat/mod.rs b/apps/maple-agent/app/src/ui/chat/mod.rs index 2688ab4fa..78dcf7dbc 100644 --- a/apps/maple-agent/app/src/ui/chat/mod.rs +++ b/apps/maple-agent/app/src/ui/chat/mod.rs @@ -3,7 +3,7 @@ //! facade + event stream. use crate::settings::PermissionMode; -use std::collections::{HashMap, HashSet}; +use std::collections::{BTreeSet, HashMap, HashSet}; use std::sync::Arc; use gpui::{ @@ -168,6 +168,12 @@ pub(crate) struct SpeechState { pub playing: bool, } +#[derive(Default)] +pub(super) struct QuestionSelection { + cursor: Option, + picked: BTreeSet, +} + pub struct ChatScreen { backend: Arc, user_id: String, @@ -388,7 +394,7 @@ pub struct ChatScreen { /// inverts the `tool_details` default for that item. toggled_tools: HashSet, /// Ticked options of a multi-select question. - question_selected: HashMap, + question_selected: HashMap, /// Index of the question being shown within the current card; a batch /// iterates one question at a time instead of listing them all. question_step: usize, @@ -2762,8 +2768,7 @@ impl ChatScreen { self.sidebar.read(cx).step_target(delta) } - /// Ctrl-1 to Ctrl-9: answer the question card with the numbered - /// option, the same as picking it and pressing Answer. + /// Ctrl-1 to Ctrl-9 submits a single choice or toggles a multi-select option. fn pick_and_submit_question_option(&mut self, index: usize, cx: &mut Context) { let Some(question) = self.current_question() else { return; @@ -2779,8 +2784,12 @@ impl ChatScreen { if index >= options { return; } - self.select_question_option(step, index, cx); - self.submit_question(cx); + if question.questions[step].multi_select { + self.toggle_question_option(step, index, cx); + } else { + self.select_question_option(step, index, cx); + self.submit_question(cx); + } } fn escape(&mut self, cx: &mut Context) { @@ -3569,10 +3578,15 @@ impl ChatScreen { .map(|input| input.read(cx).text()) .unwrap_or_default(); let typed = typed.trim(); - if let Some(option_index) = self.question_selected.get(&step) - && let Some(option) = entry.options.get(*option_index) - { - let mut answer = vec![option.label.clone()]; + let mut answer: Vec<_> = self + .question_selected + .get(&step) + .into_iter() + .flat_map(|selection| &selection.picked) + .filter_map(|index| entry.options.get(*index)) + .map(|option| option.label.clone()) + .collect(); + if !answer.is_empty() { if !typed.is_empty() { // The prefix keeps the note from reading as a second // picked option in the echoed answers array. @@ -3669,28 +3683,62 @@ impl ChatScreen { } } - /// Pick one option of one question (single-select, codex shape). + /// Pick an option, replacing earlier choices only for single-select questions. fn select_question_option( &mut self, question_index: usize, option_index: usize, cx: &mut Context, ) { - self.question_selected.insert(question_index, option_index); + let Some(question) = self + .current_question() + .and_then(|q| q.questions.get(question_index)) + else { + return; + }; + if option_index >= question.options.len() { + return; + } + let multi_select = question.multi_select; + let selection = self.question_selected.entry(question_index).or_default(); + selection.cursor = Some(option_index); + if !multi_select { + selection.picked.clear(); + } + selection.picked.insert(option_index); cx.notify(); } + /// Vim navigation moves the cursor without checking multi-select options. + fn focus_question_option(&mut self, step: usize, index: usize, cx: &mut Context) { + if self + .current_question() + .and_then(|q| q.questions.get(step)) + .is_some_and(|q| q.multi_select) + { + self.question_selected.entry(step).or_default().cursor = Some(index); + cx.notify(); + } else { + self.select_question_option(step, index, cx); + } + } + /// Click handler for an option row: clicking the picked option again - /// clears it so a typed answer can stand alone. Ctrl-N keeps plain - /// select semantics because it submits immediately. + /// clears it so a typed answer can stand alone. fn toggle_question_option( &mut self, question_index: usize, option_index: usize, cx: &mut Context, ) { - if self.question_selected.get(&question_index) == Some(&option_index) { - self.question_selected.remove(&question_index); + if let Some(selection) = self.question_selected.get_mut(&question_index) + && selection.picked.remove(&option_index) + { + if selection.picked.is_empty() { + self.question_selected.remove(&question_index); + } else { + selection.cursor = Some(option_index); + } cx.notify(); } else { self.select_question_option(question_index, option_index, cx); diff --git a/apps/maple-agent/app/src/ui/chat/navigation.rs b/apps/maple-agent/app/src/ui/chat/navigation.rs index dfbbfe4f7..a7041940d 100644 --- a/apps/maple-agent/app/src/ui/chat/navigation.rs +++ b/apps/maple-agent/app/src/ui/chat/navigation.rs @@ -336,10 +336,9 @@ impl ChatScreen { .map(|question| question.options.len()) .unwrap_or(0); if len > 0 { - let current = self.question_selected.get(&step).copied(); + let current = self.question_selected.get(&step).and_then(|s| s.cursor); let next = stepped_index(current, len, direction, count); - self.question_selected.insert(step, next); - cx.notify(); + self.focus_question_option(step, next, cx); } return; } @@ -383,9 +382,7 @@ impl ChatScreen { .map(|question| question.options.len()) .unwrap_or(0); if len > 0 { - self.question_selected - .insert(step, if first { 0 } else { len - 1 }); - cx.notify(); + self.focus_question_option(step, if first { 0 } else { len - 1 }, cx); } return; } @@ -558,13 +555,24 @@ impl ChatScreen { .min(question.questions.len().saturating_sub(1)) }) .unwrap_or_default(); - if self.question_selected.contains_key(&step) { + if self + .current_question() + .and_then(|q| q.questions.get(step)) + .is_some_and(|q| q.multi_select) + && let Some(index) = self.question_selected.get(&step).and_then(|s| s.cursor) + { + self.toggle_question_option(step, index, cx); + } else if self + .question_selected + .get(&step) + .is_some_and(|s| !s.picked.is_empty()) + { self.submit_question(cx); } else if let Some(input) = self.pending_question_input.clone() { // Enter before an option is picked is the semantic route into // the card's free-form answer. Enter inside the input still // submits through TextInput's existing callback; j/k followed - // by Enter retains the direct option-submit path. + // by Enter retains the direct submit path for single choices. let handle = input.read(cx).focus_handle(cx); window.focus(&handle, cx); cx.notify(); diff --git a/apps/maple-agent/app/src/ui/chat/tests.rs b/apps/maple-agent/app/src/ui/chat/tests.rs index b9aea2c75..325415c45 100644 --- a/apps/maple-agent/app/src/ui/chat/tests.rs +++ b/apps/maple-agent/app/src/ui/chat/tests.rs @@ -594,6 +594,7 @@ mod state_tests { session_id: "s1".to_string(), request_id: format!("req-{id}"), questions: vec![maple_agent::agent::AgentQuestion { + multi_select: false, id: id.to_string(), header: "Question".to_string(), question: text.to_string(), @@ -624,6 +625,7 @@ mod state_tests { session_id: "s1".to_string(), request_id: "q2".to_string(), questions: vec![maple_agent::agent::AgentQuestion { + multi_select: false, id: "pick".to_string(), header: "Pick".to_string(), question: "Pick one".to_string(), @@ -650,6 +652,172 @@ mod state_tests { }); } + #[gpui::test] + fn test_multi_select_preserves_choices_notes_and_question_steps(cx: &mut TestAppContext) { + let screen = screen(cx); + screen.update(cx, |this, cx| { + let question = |id: &str, multi_select| maple_agent::agent::AgentQuestion { + id: id.into(), + header: "Choose".into(), + question: "Which options?".into(), + multi_select, + options: ["A", "B", "C"] + .into_iter() + .map(|label| maple_agent::agent::AgentQuestionOption { + label: label.into(), + description: String::new(), + }) + .collect(), + }; + this.handle_service_event( + AgentServiceEvent::Question { + session_id: "s1".into(), + request_id: "multi".into(), + questions: vec![question("many", true), question("one", false)], + }, + cx, + ); + this.toggle_question_option(0, 1, cx); + this.toggle_question_option(0, 0, cx); + assert_eq!(this.question_selected[&0].picked, BTreeSet::from([0, 1])); + // Number shortcuts toggle without prematurely sending the answer. + this.pick_and_submit_question_option(0, cx); + assert_eq!(this.question_step, 0); + assert!(this.question_step_answers.is_empty()); + assert_eq!(this.question_selected[&0].picked, BTreeSet::from([1])); + // Moving the Vim cursor must not check the options it passes. + this.focus_question_option(0, 2, cx); + this.focus_question_option(0, 0, cx); + assert_eq!(this.question_selected[&0].picked, BTreeSet::from([1])); + this.pick_and_submit_question_option(0, cx); + this.pending_question_input + .clone() + .unwrap() + .update(cx, |input, cx| input.set_text(" on Linux ", cx)); + this.submit_question(cx); + assert_eq!(this.question_step, 1); + assert!(this.question_selected.is_empty()); + assert_eq!( + this.question_step_answers[&0], + ["A", "B", "Additional note: on Linux"] + ); + this.toggle_question_option(1, 0, cx); + this.toggle_question_option(1, 2, cx); + let answer: serde_json::Value = + serde_json::from_str(&this.composed_question_answer(cx)).unwrap(); + assert_eq!( + answer["answers"]["many"]["answers"], + serde_json::json!(["A", "B", "Additional note: on Linux"]) + ); + assert_eq!( + answer["answers"]["one"]["answers"], + serde_json::json!(["C"]) + ); + }); + } + + /// Run this ignored fixture alone on a private Linux desktop. It mounts + /// the real chat screen without credentials or a live inference request. + #[cfg(target_os = "linux")] + #[test] + #[ignore = "interactive native question-card fixture; requires a private display"] + fn native_question_card_fixture() { + gpui_platform::application() + .with_assets(crate::assets::Assets) + .run(|cx| { + cx.text_system() + .add_fonts( + crate::assets::FONTS + .iter() + .map(|font| std::borrow::Cow::Borrowed(*font)) + .collect(), + ) + .unwrap(); + let _shortcuts = + crate::shortcuts::ShortcutRuntime::bootstrap(&Default::default(), cx); + theme::apply_preference( + theme::Preference::parse(&crate::settings::load_settings().theme), + cx, + ); + let backend = Arc::new( + AgentBackend::new("http://127.0.0.1:9".into(), String::new()).unwrap(), + ); + cx.open_window( + gpui::WindowOptions { + window_bounds: Some(gpui::WindowBounds::Windowed(gpui::Bounds::centered( + None, + gpui::size(px(1100.), px(800.)), + cx, + ))), + titlebar: Some(gpui::TitlebarOptions { + title: Some("Question card fixture".into()), + ..Default::default() + }), + ..Default::default() + }, + |window, cx| { + theme::resolve(window.appearance()); + cx.new(|cx| { + let mut chat = + ChatScreen::new_mounted(backend, "fixture-user".into(), cx); + chat.booting = false; + chat.selected_session = Some("s1".into()); + chat.replace_timeline(vec![user_item( + "prompt", + "Choose which checks to run.", + )]); + chat.handle_service_event( + AgentServiceEvent::Question { + session_id: "s1".into(), + request_id: "fixture".into(), + questions: [true, false] + .into_iter() + .enumerate() + .map(|(index, multi_select)| { + maple_agent::agent::AgentQuestion { + id: format!("q{index}"), + header: "Validation".into(), + multi_select, + question: if multi_select { + "Which checks should run?" + } else { + "Which check should run first?" + } + .into(), + options: [ + "Unit tests", + "Integration tests", + "Native UI checks", + ] + .into_iter() + .map(|label| { + maple_agent::agent::AgentQuestionOption { + label: label.into(), + description: String::new(), + } + }) + .collect(), + } + }) + .collect(), + }, + cx, + ); + chat + }) + }, + ) + .unwrap(); + cx.on_window_closed(|cx, _| { + if cx.windows().is_empty() { + cx.quit(); + } + }) + .detach(); + cx.activate(true); + }); + } + #[gpui::test] fn test_question_typed_text_rides_along_with_picked_option(cx: &mut TestAppContext) { cx.executor().allow_parking(); @@ -658,6 +826,7 @@ mod state_tests { session_id: "s1".to_string(), request_id: "q3".to_string(), questions: vec![maple_agent::agent::AgentQuestion { + multi_select: false, id: "pick".to_string(), header: "Pick".to_string(), question: "Pick one".to_string(), @@ -689,6 +858,7 @@ mod state_tests { session_id: "s1".to_string(), request_id: "q4".to_string(), questions: vec![maple_agent::agent::AgentQuestion { + multi_select: false, id: "pick".to_string(), header: "Pick".to_string(), question: "Pick one".to_string(), @@ -703,7 +873,7 @@ mod state_tests { // Clicking the picked option again clears it; the typed text // then stands alone. this.toggle_question_option(0, 0, cx); - assert_eq!(this.question_selected.get(&0), Some(&0)); + assert_eq!(this.question_selected[&0].picked, BTreeSet::from([0])); this.toggle_question_option(0, 0, cx); assert!(this.question_selected.is_empty()); let input = this.pending_question_input.clone().expect("input exists"); @@ -853,6 +1023,7 @@ mod state_tests { request_id: "batch".to_string(), questions: vec![ maple_agent::agent::AgentQuestion { + multi_select: false, id: "first".to_string(), header: "One".to_string(), question: "First?".to_string(), @@ -862,6 +1033,7 @@ mod state_tests { }], }, maple_agent::agent::AgentQuestion { + multi_select: false, id: "second".to_string(), header: "Two".to_string(), question: "Second?".to_string(), @@ -897,6 +1069,7 @@ mod state_tests { session_id: "s1".to_string(), request_id: id.to_string(), questions: vec![maple_agent::agent::AgentQuestion { + multi_select: false, id: id.to_string(), header: "Question".to_string(), question: format!("Question {id}"), @@ -962,6 +1135,7 @@ mod state_tests { session_id: "s1".to_string(), request_id: "free-form".to_string(), questions: vec![maple_agent::agent::AgentQuestion { + multi_select: false, id: "answer".to_string(), header: "Question".to_string(), question: "What should Maple do?".to_string(), @@ -1014,6 +1188,7 @@ mod state_tests { session_id: "s1".to_string(), request_id: "q9".to_string(), questions: vec![maple_agent::agent::AgentQuestion { + multi_select: false, id: "paused".to_string(), header: "Question".to_string(), question: "Paused?".to_string(), @@ -1291,6 +1466,7 @@ mod state_tests { session_id: "s2".to_string(), request_id: "req-other".to_string(), questions: vec![maple_agent::agent::AgentQuestion { + multi_select: false, id: "other".to_string(), header: "Question".to_string(), question: "From another task".to_string(), @@ -1316,10 +1492,19 @@ mod state_tests { assert_eq!(this.pending_questions.len(), 1); // A second question for the shown session keeps the pick made // on the first card. - this.handle_service_event(one_question("a", "First?"), cx); + let mut event = one_question("a", "First?"); + if let AgentServiceEvent::Question { questions, .. } = &mut event { + questions[0] + .options + .push(maple_agent::agent::AgentQuestionOption { + label: "A".into(), + description: String::new(), + }); + } + this.handle_service_event(event, cx); this.select_question_option(0, 0, cx); this.handle_service_event(one_question("b", "Second?"), cx); - assert_eq!(this.question_selected.get(&0), Some(&0)); + assert_eq!(this.question_selected[&0].picked, BTreeSet::from([0])); // Switching to the other task shows its card. this.set_active_session(summary("s2", "B"), Vec::new(), HashMap::new(), cx); assert_eq!( diff --git a/apps/maple-agent/app/src/ui/chat/transcript.rs b/apps/maple-agent/app/src/ui/chat/transcript.rs index f81b719be..0cb171279 100644 --- a/apps/maple-agent/app/src/ui/chat/transcript.rs +++ b/apps/maple-agent/app/src/ui/chat/transcript.rs @@ -14,7 +14,7 @@ use maple_agent::agent::{ use super::cache::{MAX_DIFF_LINES, MarkdownKind}; use super::commands::ChatCommand; use super::speech::speak_message_button; -use super::{CONTENT_WIDTH, ChatScreen, TranscriptCtx}; +use super::{CONTENT_WIDTH, ChatScreen, QuestionSelection, TranscriptCtx}; use crate::backend::PendingPermission; use crate::ui::icons::{icon, spinner, spinner_with_id}; @@ -1294,7 +1294,7 @@ pub(super) fn render_question_card( question: &crate::backend::PendingQuestion, step: usize, input: Option>, - selected: &HashMap, + selected: &HashMap, cx: &mut Context, ) -> Div { let mut card = div() @@ -1351,11 +1351,25 @@ pub(super) fn render_question_card( .text_color(gpui::rgb(theme::text_secondary())) .child(entry.question.clone()), ); + if entry.multi_select { + block = block.child( + div() + .text_xs() + .text_color(gpui::rgb(theme::text_muted())) + .child(if has_more { + "Select all that apply, then choose Next." + } else { + "Select all that apply, then choose Answer." + }), + ); + } for (option_index, option) in entry.options.iter().enumerate() { - let is_picked = selected.get(&question_index) == Some(&option_index); + let selection = selected.get(&question_index); + let is_picked = selection.is_some_and(|s| s.picked.contains(&option_index)); let marker = div() .size_3() - .rounded_full() + .when(entry.multi_select, |marker| marker.rounded_sm()) + .when(!entry.multi_select, |marker| marker.rounded_full()) .border_1() .border_color(gpui::rgb(if is_picked { theme::accent() @@ -1402,6 +1416,11 @@ pub(super) fn render_question_card( .px_2() .py_1p5() .rounded(theme::RADIUS_SM) + .when( + entry.multi_select + && selection.is_some_and(|s| s.cursor == Some(option_index)), + |row| row.bg(gpui::rgb(theme::bg_sidebar_row_hover())), + ) .hover(|style| { style .bg(gpui::rgb(theme::bg_sidebar_row_hover())) @@ -1501,6 +1520,7 @@ pub(super) fn render_waiting_indicator() -> gpui::Stateful
{ /// external agent whose request Maple relays. pub(super) fn permission_card_heading(tool_name: &str) -> &'static str { match tool_name { + "claude_tool" => "Claude Code wants to use a tool", "codex_command" => "Codex wants to run a command", "codex_file_change" => "Codex wants to change files", _ => "Permission required", diff --git a/apps/maple-agent/app/src/ui/settings.rs b/apps/maple-agent/app/src/ui/settings.rs index cd712a7dc..5b59f7e86 100644 --- a/apps/maple-agent/app/src/ui/settings.rs +++ b/apps/maple-agent/app/src/ui/settings.rs @@ -2559,9 +2559,13 @@ impl SettingsScreen { .size(px(28.)) .text_color(gpui::rgb(theme::text_primary())) .into_any_element() - } else if is_codex(&integration.id) { + } else if matches!(integration.id.as_str(), "codex" | "claude") { gpui::svg() - .path("icons/openai-mark.svg") + .path(if integration.id == "claude" { + "icons/claude-mark.svg" + } else { + "icons/openai-mark.svg" + }) .size(px(26.)) .text_color(gpui::rgb(theme::text_primary())) .into_any_element() @@ -2787,9 +2791,9 @@ fn plan_card(plan: &crate::billing::PlanUsage) -> Div { } fn integration_is_visible(integration: &AgentIntegration) -> bool { - // Codex is worth a row even before it is installed: the card says how + // External agents get a row even before installation: the card says how // to get it, where a hidden row would leave the feature undiscoverable. - if is_codex(&integration.id) { + if integration.is_external_agent() { return true; } if is_cua_driver(&integration.id) @@ -2944,10 +2948,6 @@ fn is_cua_driver(id: &str) -> bool { id == "cua-driver" } -fn is_codex(id: &str) -> bool { - id == "codex" -} - /// Small on/off pill used in list rows. fn pill_button( id: String, @@ -3430,6 +3430,15 @@ mod tests { codex_ready.id = "codex".to_string(); assert!(integration_can_enable(&codex_ready)); assert!(!integration_can_setup(&codex_ready)); + let mut claude_missing = integration(AgentIntegrationAvailability::NotDetected, false); + claude_missing.id = "claude".to_string(); + assert!(integration_is_visible(&claude_missing)); + assert!(!integration_can_enable(&claude_missing)); + let mut claude_needs_setup = + integration(AgentIntegrationAvailability::SetupRequired, false); + claude_needs_setup.id = "claude".to_string(); + assert!(integration_is_visible(&claude_needs_setup)); + assert!(!integration_can_enable(&claude_needs_setup)); let mut external = integration(AgentIntegrationAvailability::Available, true); external.backend = Some(AgentIntegrationBackend::External); diff --git a/apps/maple-agent/crates/maple-agent/resources/skills/advisor/SKILL.md b/apps/maple-agent/crates/maple-agent/resources/skills/advisor/SKILL.md index 55fd98d39..899892e52 100644 --- a/apps/maple-agent/crates/maple-agent/resources/skills/advisor/SKILL.md +++ b/apps/maple-agent/crates/maple-agent/resources/skills/advisor/SKILL.md @@ -1,6 +1,6 @@ --- name: advisor -description: Ask an external agent (Codex) for read-only analysis or review of code, a plan, or a problem, without letting it change anything. Use when the user wants a review, an audit, or advice from another model. +description: Ask an external agent (Codex or Claude Code) for read-only analysis or review of code, a plan, or a problem, without letting it change anything. Use when the user wants a review, an audit, or advice from another model. metadata: maple: external-agents argument-hint: "" diff --git a/apps/maple-agent/crates/maple-agent/resources/skills/committee/SKILL.md b/apps/maple-agent/crates/maple-agent/resources/skills/committee/SKILL.md index 3680818cc..1122e0b16 100644 --- a/apps/maple-agent/crates/maple-agent/resources/skills/committee/SKILL.md +++ b/apps/maple-agent/crates/maple-agent/resources/skills/committee/SKILL.md @@ -1,6 +1,6 @@ --- name: committee -description: Get several independent opinions on a question or design by asking external agents (Codex) and comparing them with your own analysis. Use when the user wants a second opinion, a review from another model, or a comparison of approaches. +description: Get several independent opinions on a question or design by asking external agents (Codex or Claude Code) and comparing them with your own analysis. Use when the user wants a second opinion, a review from another model, or a comparison of approaches. metadata: maple: external-agents argument-hint: "" diff --git a/apps/maple-agent/crates/maple-agent/resources/skills/handoff/SKILL.md b/apps/maple-agent/crates/maple-agent/resources/skills/handoff/SKILL.md index b673d1a80..75821ef6a 100644 --- a/apps/maple-agent/crates/maple-agent/resources/skills/handoff/SKILL.md +++ b/apps/maple-agent/crates/maple-agent/resources/skills/handoff/SKILL.md @@ -1,6 +1,6 @@ --- name: handoff -description: Hand a self-contained piece of work to an external coding agent (Codex) that runs in this project with its own context. Use when the user asks to delegate, hand off, or have Codex implement something. +description: Hand a self-contained piece of work to an external coding agent (Codex or Claude Code) that runs in this project with its own context. Use when the user asks to delegate, hand off, or have Codex or Claude Code implement something. metadata: maple: external-agents argument-hint: "" @@ -8,7 +8,7 @@ metadata: # Hand work to an external agent -Maple can start an external coding agent (Codex today) inside this project. +Maple can start an external coding agent (Codex or Claude Code) inside this project. The agent runs with its own context and its own account. It does not see this conversation. It runs under its own sandbox and approval settings; whatever it asks approval for comes to the user through Maple. @@ -16,7 +16,7 @@ whatever it asks approval for comes to the user through Maple. ## Steps 1. Call `list_agent_providers` first. If no provider is usable, tell the - user what is missing (install, PATH, or `codex login`) and stop. + user what is missing (install, PATH, or the provider’s sign-in command) and stop. 2. Write a self-contained briefing. The agent has zero context, so the briefing must carry everything: - **Task**: what to do, in one or two sentences. diff --git a/apps/maple-agent/crates/maple-agent/src/agent.rs b/apps/maple-agent/crates/maple-agent/src/agent.rs index fa3483a97..d1b2a4e0b 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent.rs @@ -18576,6 +18576,7 @@ mod tests { )); let broker = state.question_broker(); let one_question = |id: &str, text: &str| AgentQuestion { + multi_select: false, id: id.to_string(), header: "Question".to_string(), question: text.to_string(), diff --git a/apps/maple-agent/crates/maple-agent/src/agent/developer_tools.rs b/apps/maple-agent/crates/maple-agent/src/agent/developer_tools.rs index 8ace1e25d..659fe84dd 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent/developer_tools.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent/developer_tools.rs @@ -250,7 +250,7 @@ impl MapleDeveloperClient { Tool::new( AGENT_START_TOOL.to_string(), format!( - "Hand a self-contained piece of work to an external coding agent (an installed harness such as Codex) that runs in the project with its own context and its own account. \ + "Hand a self-contained piece of work to an external coding agent (an installed harness such as Codex or Claude Code) that runs in the project with its own context and its own account. \ The new agent knows nothing about this conversation: write a complete briefing with the task, relevant files, current state, what was tried, decisions made, acceptance criteria, and constraints. \ It runs under its own sandbox and approval settings; whatever it asks approval for comes to the user through Maple, and in Allow all Maple grants it. \ Blocking by default: the call returns the agent's result. With background=true the call returns at once and Maple tells you when the agent finishes; do not poll. \ @@ -261,7 +261,7 @@ Call {LIST_AGENT_PROVIDERS_TOOL} first when unsure what is installed." "properties": { "provider": { "type": "string", - "description": "Which external agent to use, from list_agent_providers (for example \"codex\")" + "description": "Which external agent to use, from list_agent_providers (for example \"codex\" or \"claude\")" }, "prompt": { "type": "string", @@ -2824,6 +2824,10 @@ pub(super) fn parse_user_questions( .take(5) .collect(); questions.push(crate::agent::AgentQuestion { + multi_select: entry + .get("multiSelect") + .and_then(serde_json::Value::as_bool) + .unwrap_or(false), id, header, question, @@ -3110,6 +3114,22 @@ mod tests { assert_eq!(questions[1].question, "Second?"); } + #[test] + fn parse_user_questions_preserves_multi_select_without_changing_the_default() { + let questions = parse_user_questions(&[ + serde_json::json!({"question": "Pick several", "multiSelect": true}), + serde_json::json!({"question": "Pick one"}), + serde_json::json!({"question": "Invalid flag", "multiSelect": "true"}), + ]); + assert!(questions[0].multi_select); + assert!(!questions[1].multi_select); + assert!(!questions[2].multi_select); + assert_eq!( + serde_json::to_value(&questions[0]).unwrap()["multiSelect"], + true + ); + } + #[test] fn parse_user_questions_skips_blank_entries() { let entries: Vec = diff --git a/apps/maple-agent/crates/maple-agent/src/agent/external_agents/app_server.rs b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/app_server.rs index 1888c4ad2..70ac547b5 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent/external_agents/app_server.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/app_server.rs @@ -1,4 +1,4 @@ -//! JSON-RPC 2.0 over newline-delimited stdio, as `codex app-server` speaks it. +//! Codex app-server transport and the shared external-agent client interface. //! //! One reader task owns the child's stdout. It resolves responses to the //! requests this client sent, and hands notifications and server-initiated @@ -21,6 +21,30 @@ const MAX_LINE_BYTES: usize = 4 * 1024 * 1024; /// message for long: approvals run on their own tasks. const SERVER_MESSAGE_CAPACITY: usize = 256; +/// Operations supported by Maple's external-agent host. +#[derive(Clone, Copy, Debug)] +pub(super) enum RequestMethod { + Initialize, + ThreadStart, + ThreadResume, + TurnStart, + TurnSteer, + TurnInterrupt, +} + +impl RequestMethod { + fn codex_method(self) -> &'static str { + match self { + Self::Initialize => "initialize", + Self::ThreadStart => "thread/start", + Self::ThreadResume => "thread/resume", + Self::TurnStart => "turn/start", + Self::TurnSteer => "turn/steer", + Self::TurnInterrupt => "turn/interrupt", + } + } +} + /// A message the server initiated. #[derive(Debug)] pub(super) enum ServerMessage { @@ -36,6 +60,56 @@ pub(super) enum ServerMessage { }, } +/// Both transports expose Maple's existing activity and permission messages. +/// Claude translates these in Rust; its child speaks Claude's native protocol. +pub(super) enum AgentClient { + Codex(Arc), + Claude(Arc), +} + +impl AgentClient { + pub(super) fn closed(&self) -> &CancellationToken { + match self { + Self::Codex(client) => client.closed(), + Self::Claude(client) => client.closed(), + } + } + + pub(super) async fn request( + &self, + method: RequestMethod, + params: Value, + ) -> Result { + match self { + Self::Codex(client) => client.request(method.codex_method(), params).await, + Self::Claude(client) => client.request(method, params).await, + } + } + + pub(super) async fn initialized(&self) -> Result<(), String> { + match self { + Self::Codex(client) => client.notify("initialized", json!({})).await, + Self::Claude(_) => Ok(()), + } + } + + pub(super) async fn respond(&self, id: Value, result: Value) -> Result<(), String> { + match self { + Self::Codex(client) => client.respond(id, result).await, + Self::Claude(client) => client.respond(id, result).await, + } + } + + pub(super) async fn respond_error(&self, id: Value, message: &str) { + match self { + Self::Codex(client) => client.respond_error(id, message).await, + Self::Claude(client) => { + let _ = client.respond(id, json!({"decision": "cancel"})).await; + } + } + } +} + type PendingResponses = Arc>>>>; pub(super) struct AppServerClient { diff --git a/apps/maple-agent/crates/maple-agent/src/agent/external_agents/claude.rs b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/claude.rs new file mode 100644 index 000000000..bc1bd08ab --- /dev/null +++ b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/claude.rs @@ -0,0 +1,659 @@ +//! Claude Code's native stream-json transport, adapted from Goose's +//! `crates/goose/src/providers/claude_code.rs` at 785d655d110746147117d23690e09cc7023aa9dc. +//! Source: https://github.com/AnthonyRonning/goose. +//! +//! The control request/response types and permission exchange originate in +//! Goose's Rust SDK protocol implementation. Maple adapts process ownership, +//! bounded concurrent reads, question answers, and activity projection here; +//! it does not instantiate Goose's provider, whose subprocess is private. +//! Unlike Goose's Auto mode, we never set --dangerously-skip-permissions. + +use super::super::developer_tools::{ + executable_in_search_path, executable_on_path, spawn_contained, +}; +use super::app_server::{RequestMethod, ServerMessage}; +use super::codex; +use futures_util::StreamExt; +use serde::{Deserialize, Serialize}; +use serde_json::{Value, json}; +use std::collections::HashMap; +use std::path::{Path, PathBuf}; +use std::process::Stdio; +use std::sync::atomic::{AtomicBool, AtomicU64, Ordering}; +use std::sync::{Arc, Mutex as StdMutex}; +use std::time::Duration; +use tokio::io::AsyncWriteExt; +use tokio::process::{ChildStdin, ChildStdout}; +use tokio::sync::{Mutex, mpsc, oneshot}; +use tokio_util::codec::{FramedRead, LinesCodec}; +use tokio_util::sync::CancellationToken; + +pub(crate) const PROVIDER_ID: &str = "claude"; +pub(crate) const PROVIDER_NAME: &str = "Claude Code"; +const MAX_LINE_BYTES: usize = 4 * 1024 * 1024; +const AUTH_PROBE_TIMEOUT: Duration = Duration::from_secs(3); +const MAX_AUTH_BYTES: usize = 16 * 1024; +const FAILURE: &str = + "Claude Code could not complete this request. Check its sign-in and configuration."; + +#[derive(Debug, Clone, Default)] +pub(crate) struct ClaudeDetection { + pub(crate) executable: Option, + pub(crate) version: Option, + /// `None` when the installed CLI cannot report its authentication state. + pub(crate) signed_in: Option, + pub(crate) problem: Option, +} + +pub(super) fn find_executable(search_path: Option<&str>) -> Option { + match search_path { + Some(path) => executable_in_search_path("claude", path), + None => executable_on_path("claude"), + } +} + +pub(crate) async fn detect(search_path: Option<&str>) -> ClaudeDetection { + let Some(executable) = find_executable(search_path) else { + return ClaudeDetection::default(); + }; + let mut detection = ClaudeDetection { + executable: Some(executable.clone()), + ..Default::default() + }; + match codex::probe_version(&executable).await { + Ok(version) => { + detection.version = Some(version); + detection.signed_in = probe_auth_status(&executable, search_path).await; + } + Err(_) => { + detection.problem = Some( + "Maple could not run `claude --version`. Check the Claude Code installation." + .into(), + ) + } + } + detection +} + +pub(crate) fn sign_in_hint() -> &'static str { + "Claude Code is not signed in. Run `claude auth login` in a terminal, then try again." +} + +async fn probe_auth_status(executable: &Path, search_path: Option<&str>) -> Option { + let mut command = tokio::process::Command::new(executable); + command + .args(["auth", "status", "--json"]) + .stdin(Stdio::null()) + .stdout(Stdio::piped()) + .stderr(Stdio::null()) + .kill_on_drop(true); + if let Some(path) = search_path { + command.env("PATH", path); + } + let mut child = spawn_contained(command).ok()?; + let stdout = child.as_mut().stdout().take()?; + let result = tokio::time::timeout(AUTH_PROBE_TIMEOUT, async { + tokio::join!( + child.as_mut().wait(), + super::super::bounded_process::read_bounded_stdout( + stdout, + MAX_AUTH_BYTES, + "Claude authentication status", + ) + ) + }) + .await; + child.kill_and_wait().await; + let (status, output) = result.ok()?; + // Read only the boolean; never retain or log account details from the CLI. + #[derive(Deserialize)] + struct AuthStatus { + #[serde(rename = "loggedIn")] + logged_in: bool, + } + let auth: AuthStatus = serde_json::from_slice(&output.ok()?).ok()?; + match (status.ok()?.code(), auth.logged_in) { + (Some(0), true) => Some(true), + (Some(1), false) => Some(false), + _ => None, + } +} + +pub(super) fn new_session_id() -> String { + // Claude requires a UUID for --session-id. Generate RFC 4122 version 4 + // using the runtime's existing random source. + let value = (rand::random::() & !(0xf_u128 << 76 | 0x3_u128 << 62)) + | 0x4_u128 << 76 + | 0x2_u128 << 62; + let hex = format!("{value:032x}"); + format!( + "{}-{}-{}-{}-{}", + &hex[..8], + &hex[8..12], + &hex[12..16], + &hex[16..20], + &hex[20..] + ) +} + +pub(super) fn command_args( + session: &str, + resume: bool, + model: Option<&str>, + effort: Option<&str>, +) -> Vec { + let mut args: Vec = [ + "--input-format", + "stream-json", + "--output-format", + "stream-json", + "--verbose", + "--include-partial-messages", + "--permission-prompt-tool", + "stdio", + "--permission-mode", + "default", + if resume { "--resume" } else { "--session-id" }, + session, + ] + .into_iter() + .map(str::to_string) + .collect(); + if let Some(model) = model { + args.extend(["--model".into(), model.into()]); + } + if let Some(effort) = effort { + args.extend(["--effort".into(), effort.into()]); + } + args +} + +// Adapted from Goose's control protocol types for Claude's SDK wire format. +#[derive(Serialize)] +struct ControlResponse { + #[serde(rename = "type")] + msg_type: &'static str, + response: ControlResponseBody, +} + +#[derive(Serialize)] +struct ControlResponseBody { + subtype: &'static str, + request_id: String, + response: T, +} + +#[derive(Serialize)] +#[serde(tag = "behavior")] +enum PermissionResponse { + #[serde(rename = "allow")] + Allow { + #[serde(rename = "updatedInput")] + updated_input: serde_json::Map, + #[serde(rename = "toolUseID")] + tool_use_id: String, + }, + #[serde(rename = "deny")] + Deny { message: String }, +} + +#[derive(Serialize)] +struct ControlRequest { + #[serde(rename = "type")] + msg_type: &'static str, + request_id: String, + request: ControlRequestBody, +} + +#[derive(Serialize)] +#[serde(tag = "subtype", rename_all = "snake_case")] +enum ControlRequestBody { + Initialize, + Interrupt, +} + +#[derive(Deserialize)] +struct IncomingControlRequest { + request_id: String, + request: IncomingRequestBody, +} + +#[derive(Deserialize)] +#[serde(tag = "subtype")] +enum IncomingRequestBody { + #[serde(rename = "can_use_tool")] + CanUseTool { + tool_name: String, + #[serde(default)] + input: serde_json::Map, + #[serde(default)] + tool_use_id: String, + }, +} + +impl ControlResponse { + fn success(request_id: String, response: T) -> Self { + Self { + msg_type: "control_response", + response: ControlResponseBody { + subtype: "success", + request_id, + response, + }, + } + } +} + +type Pending = StdMutex>>>; + +pub(super) struct Client { + writer: Mutex>, + pending: Pending, + permissions: StdMutex>, + next_id: AtomicU64, + session: String, + interrupted: AtomicBool, + closed: CancellationToken, + sender: mpsc::Sender, +} + +impl Client { + pub(super) fn new( + stdin: ChildStdin, + stdout: ChildStdout, + session: String, + ) -> ( + Arc, + mpsc::Receiver, + tokio::task::JoinHandle<()>, + ) { + let (sender, receiver) = mpsc::channel(256); + let client = Arc::new(Self { + writer: Mutex::new(Some(stdin)), + pending: StdMutex::new(HashMap::new()), + permissions: StdMutex::new(HashMap::new()), + next_id: AtomicU64::new(1), + session, + interrupted: AtomicBool::new(false), + closed: CancellationToken::new(), + sender, + }); + let reader_client = Arc::clone(&client); + let reader = tokio::spawn(async move { + reader_client.read(stdout).await; + }); + (client, receiver, reader) + } + + pub(super) fn closed(&self) -> &CancellationToken { + &self.closed + } + + async fn write(&self, message: impl Serialize) -> Result<(), String> { + let mut bytes = serde_json::to_vec(&message).map_err(|_| FAILURE.to_string())?; + bytes.push(b'\n'); + let mut writer = self.writer.lock().await; + let writer = writer.as_mut().ok_or(FAILURE)?; + writer + .write_all(&bytes) + .await + .map_err(|_| FAILURE.to_string())?; + writer.flush().await.map_err(|_| FAILURE.to_string()) + } + + async fn control(&self, request: ControlRequestBody) -> Result { + let request_id = format!("req_{}", self.next_id.fetch_add(1, Ordering::Relaxed)); + let (tx, rx) = oneshot::channel(); + self.pending.lock().unwrap().insert(request_id.clone(), tx); + let result = async { + self.write(ControlRequest { + msg_type: "control_request", + request_id: request_id.clone(), + request, + }) + .await?; + tokio::select! { + result = rx => result.unwrap_or_else(|_| Err(FAILURE.into())), + _ = self.closed.cancelled() => Err(FAILURE.into()), + } + } + .await; + self.pending.lock().unwrap().remove(&request_id); + result + } + + pub(super) async fn request( + &self, + method: RequestMethod, + params: Value, + ) -> Result { + if self.closed.is_cancelled() { + return Err(FAILURE.into()); + } + match method { + RequestMethod::Initialize => self.control(ControlRequestBody::Initialize).await, + RequestMethod::ThreadStart | RequestMethod::ThreadResume => { + Ok(json!({"thread": {"id": self.session}})) + } + RequestMethod::TurnStart => { + self.notify("turn/started", json!({"turn": {"id": new_session_id()}})) + .await; + self.write(json!({ + "type": "user", "session_id": self.session, + "message": {"role": "user", "content": [{"type": "text", "text": params["input"][0]["text"]}]}, + })).await?; + Ok(json!({})) + } + RequestMethod::TurnInterrupt => { + self.interrupted.store(true, Ordering::Relaxed); + self.control(ControlRequestBody::Interrupt).await + } + RequestMethod::TurnSteer => { + Err("Claude does not support steering an active turn".into()) + } + } + } + + pub(super) async fn respond(&self, id: Value, result: Value) -> Result<(), String> { + let id = id.as_str().ok_or(FAILURE)?; + let request = self.permissions.lock().unwrap().remove(id).ok_or(FAILURE)?; + let IncomingRequestBody::CanUseTool { + tool_name, + mut input, + tool_use_id, + } = request; + let allow = if tool_name == "AskUserQuestion" { + let mut answers = serde_json::Map::new(); + for (i, question) in input + .get("questions") + .and_then(Value::as_array) + .into_iter() + .flatten() + .enumerate() + { + let Some(text) = question["question"].as_str() else { + continue; + }; + let answer = result["answers"][format!("q{i}")]["answers"] + .as_array() + .map(|values| { + values + .iter() + .filter_map(Value::as_str) + .collect::>() + .join(", ") + }) + .unwrap_or_default(); + answers.insert(text.into(), answer.into()); + } + let answered = answers + .values() + .any(|answer| answer.as_str().is_some_and(|text| !text.is_empty())); + input.insert("answers".into(), answers.into()); + answered + } else { + result["decision"] == "accept" + }; + let response = + if allow && !self.interrupted.load(Ordering::Relaxed) && !self.closed.is_cancelled() { + PermissionResponse::Allow { + updated_input: input, + tool_use_id, + } + } else { + PermissionResponse::Deny { + message: "Maple declined this action".into(), + } + }; + self.write(ControlResponse::success(id.into(), response)) + .await + } + + async fn notify(&self, method: &str, params: Value) { + let _ = self + .sender + .send(ServerMessage::Notification { + method: method.into(), + params, + }) + .await; + } + + async fn permission(&self, message: Value) -> Result<(), String> { + let request: IncomingControlRequest = + serde_json::from_value(message).map_err(|_| FAILURE)?; + let IncomingRequestBody::CanUseTool { + ref tool_name, + ref input, + ref tool_use_id, + } = request.request; + let (method, params) = if tool_name == "AskUserQuestion" { + let questions: Vec<_> = input + .get("questions") + .and_then(Value::as_array) + .into_iter() + .flatten() + .enumerate() + .filter(|(_, question)| question.is_object()) + .map(|(i, question)| { + let mut question = question.clone(); + question["id"] = format!("q{i}").into(); + question + }) + .collect(); + ( + "item/tool/requestUserInput", + json!({"questions": questions}), + ) + } else { + ( + "claude/tool/requestApproval", + json!({"tool": tool_name, "input": input, "itemId": tool_use_id}), + ) + }; + let id = request.request_id; + { + let mut pending = self.permissions.lock().unwrap(); + if pending.len() >= 64 || pending.contains_key(&id) { + return Err(FAILURE.into()); + } + pending.insert(id.clone(), request.request); + } + self.sender + .send(ServerMessage::Request { + id: id.into(), + method: method.into(), + params, + }) + .await + .map_err(|_| FAILURE.into()) + } + + async fn read(self: Arc, stdout: ChildStdout) { + // Even task abortion closes requests. No pending permission can carry + // into a replacement process or a subsequent turn. + struct Close(Arc); + impl Drop for Close { + fn drop(&mut self) { + self.0.closed.cancel(); + self.0.pending.lock().unwrap().clear(); + self.0.permissions.lock().unwrap().clear(); + } + } + let _close = Close(Arc::clone(&self)); + let mut lines = FramedRead::new(stdout, LinesCodec::new_with_max_length(MAX_LINE_BYTES)); + let mut activity = Activity::default(); + let mut completed = false; + let mut session_confirmed = false; + while let Some(line) = lines.next().await { + let Ok(line) = line else { + break; + }; + if line.trim().is_empty() { + continue; + } + let Ok(message) = serde_json::from_str::(&line) else { + break; + }; + match message["type"].as_str() { + Some("control_response") => { + let response = &message["response"]; + if let Some(id) = response["request_id"].as_str() + && let Some(tx) = self.pending.lock().unwrap().remove(id) + { + let result = if response["subtype"] == "success" { + Ok(response["response"].clone()) + } else { + Err(FAILURE.into()) + }; + let _ = tx.send(result); + } + } + Some("control_request") => { + if self.permission(message).await.is_err() { + break; + } + } + _ => { + if !session_confirmed + && message["session_id"].as_str() == Some(self.session.as_str()) + && ((message["type"] == "system" && message["subtype"] == "init") + || (message["type"] == "result" + && message["subtype"] == "success" + && message["is_error"] != true)) + { + session_confirmed = true; + self.notify("thread/started", json!({"thread": {"id": self.session}})) + .await; + } + for (method, mut params) in activity.events(&message) { + if method == "turn/completed" { + completed = true; + // EOF lets Claude flush its persisted session and + // exit normally before the next turn resumes it. + self.writer.lock().await.take(); + self.permissions.lock().unwrap().clear(); + } + if method == "turn/completed" && self.interrupted.load(Ordering::Relaxed) { + params = json!({"turn": {"status": "interrupted"}}); + } + self.notify(method, params).await; + } + } + } + } + // The client retains a sender, so premature EOF must explicitly finish + // the turn. EOF after a result must not emit a second completion. + if completed { + return; + } + let turn = if self.interrupted.load(Ordering::Relaxed) { + json!({"status": "interrupted"}) + } else { + json!({"status": "failed", "error": {"message": FAILURE}}) + }; + self.notify("turn/completed", json!({"turn": turn})).await; + } +} + +#[derive(Default)] +struct Activity { + message_id: String, + tools: HashMap, +} + +impl Activity { + fn events(&mut self, message: &Value) -> Vec<(&'static str, Value)> { + let mut events = Vec::new(); + // Nested Claude agents keep their own transcript. Only project the + // delegated agent's top-level messages into Maple's activity row. + if !message["parent_tool_use_id"].is_null() { + return events; + } + match message["type"].as_str() { + Some("stream_event") => { + let event = &message["event"]; + if event["type"] == "message_start" { + self.message_id = event["message"]["id"].as_str().unwrap_or("message").into(); + } else if event["type"] == "content_block_delta" + && event["delta"]["type"] == "text_delta" + { + events.push(( + "item/agentMessage/delta", + json!({"itemId": self.message_id, "delta": event["delta"]["text"]}), + )); + } + } + Some("assistant") => { + let message = &message["message"]; + if let Some(id) = message["id"].as_str() { + self.message_id = id.into(); + } + let blocks = message["content"].as_array().into_iter().flatten(); + let mut text = Vec::new(); + for block in blocks { + if block["type"] == "text" { + if let Some(value) = block["text"].as_str() { + text.push(value); + } + } else if block["type"] == "tool_use" { + let Some(id) = block["id"].as_str() else { + continue; + }; + let name = block["name"].as_str().unwrap_or_default(); + if matches!(name, "Bash" | "Edit" | "Write" | "NotebookEdit") + && self.tools.len() < 1024 + { + self.tools + .insert(id.into(), (name.into(), block["input"].clone())); + } + match name { + "Bash" => events.push(("item/started", json!({"item": {"id": id, "type": "commandExecution", "command": block["input"]["command"]}}))), + "TodoWrite" => events.push(("item/started", json!({"item": {"id": id, "type": "todoList", "items": block["input"]["todos"]}}))), + _ => {} + } + } + } + if !text.is_empty() { + events.push(("item/completed", json!({"item": {"id": self.message_id, "type": "agentMessage", "text": text.join("\n")}}))); + } + } + Some("user") => { + for block in message["message"]["content"] + .as_array() + .into_iter() + .flatten() + { + if block["type"] != "tool_result" { + continue; + } + let Some(id) = block["tool_use_id"].as_str() else { + continue; + }; + let Some((name, args)) = self.tools.remove(id) else { + continue; + }; + let failed = block["is_error"] == true; + if name == "Bash" { + events.push(("item/completed", json!({"item": {"id": id, "type": "commandExecution", "command": args["command"], "status": if failed { "failed" } else { "completed" }}}))); + } else if matches!(name.as_str(), "Edit" | "Write" | "NotebookEdit") + && !failed + && let Some(path) = args["file_path"] + .as_str() + .or(args["notebook_path"].as_str()) + { + events.push(("item/completed", json!({"item": {"id": id, "type": "fileChange", "changes": [{"path": path, "kind": "update"}]}}))); + } + } + } + Some("result") | Some("error") => { + let success = message["type"] == "result" + && message["subtype"] == "success" + && message["is_error"] != true; + events.push(("turn/completed", json!({"turn": {"status": if success { "completed" } else { "failed" }, "error": if success { Value::Null } else { json!({"message": FAILURE}) }}}))); + } + _ => {} + } + events + } +} diff --git a/apps/maple-agent/crates/maple-agent/src/agent/external_agents/codex.rs b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/codex.rs index e9e347386..cc8e67aae 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent/external_agents/codex.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/codex.rs @@ -78,7 +78,7 @@ pub(crate) fn find_executable(search_path: Option<&str>) -> Option { } } -async fn probe_version(executable: &Path) -> Result { +pub(super) async fn probe_version(executable: &Path) -> Result { let mut command = tokio::process::Command::new(executable); command .arg("--version") @@ -546,6 +546,7 @@ pub(super) fn async_question_prompts(questions: &[AsyncQuestion]) -> Vec Result<(), String> { - if provider.trim() == codex::PROVIDER_ID { + if matches!(provider.trim(), codex::PROVIDER_ID | claude::PROVIDER_ID) { Ok(()) } else { Err(format!( @@ -448,37 +449,60 @@ impl ExternalAgentRegistry { call: &ExternalAgentCall, providers: &[String], ) -> CallToolResult { - if !providers - .iter() - .any(|provider| provider == codex::PROVIDER_ID) - { + if providers.is_empty() { return text_result("No external agent providers are enabled for this task."); } - let detection = codex::detect(call.login_path.as_deref()).await; - let mut out = String::new(); - let _ = writeln!(out, "Agent providers available to this task:"); - match (&detection.executable, &detection.problem) { - (None, _) => { - let _ = writeln!( - out, - "- codex: not installed. Ask the user to install the Codex CLI and make sure `codex` is on PATH." - ); - } - (Some(_), Some(problem)) => { - let _ = writeln!(out, "- codex: unusable. {problem}"); - } - (Some(_), None) => { - let version = detection.version.as_deref().unwrap_or("unknown version"); + let mut out = String::from("Agent providers available to this task:\n"); + if providers + .iter() + .any(|provider| provider == claude::PROVIDER_ID) + { + let detection = claude::detect(call.login_path.as_deref()).await; + let detail = if detection.executable.is_none() { + "Install Claude Code and make sure `claude` is on PATH.".to_string() + } else if let Some(problem) = detection.problem { + problem + } else { let sign_in = match detection.signed_in { - Some(true) => "signed in".to_string(), - Some(false) => codex::sign_in_hint().to_string(), - None => "sign-in state unknown".to_string(), + Some(true) => "Signed in.", + Some(false) => claude::sign_in_hint(), + None => "Sign-in state unknown.", }; - let _ = writeln!( - out, - "- codex: {} {version}, {sign_in}. Runs `codex app-server` with the user's own Codex account and configuration; optional `model` and `effort` arguments override its defaults.", - codex::PROVIDER_NAME - ); + format!( + "{}; {sign_in} Supports model and effort overrides.", + detection.version.unwrap_or_default() + ) + }; + let _ = writeln!(out, "- claude: {detail}"); + } + if providers + .iter() + .any(|provider| provider == codex::PROVIDER_ID) + { + let detection = codex::detect(call.login_path.as_deref()).await; + match (&detection.executable, &detection.problem) { + (None, _) => { + let _ = writeln!( + out, + "- codex: not installed. Ask the user to install the Codex CLI and make sure `codex` is on PATH." + ); + } + (Some(_), Some(problem)) => { + let _ = writeln!(out, "- codex: unusable. {problem}"); + } + (Some(_), None) => { + let version = detection.version.as_deref().unwrap_or("unknown version"); + let sign_in = match detection.signed_in { + Some(true) => "signed in".to_string(), + Some(false) => codex::sign_in_hint().to_string(), + None => "sign-in state unknown".to_string(), + }; + let _ = writeln!( + out, + "- codex: {} {version}, {sign_in}. Runs `codex app-server` with the user's own Codex account and configuration; optional `model` and `effort` arguments override its defaults.", + codex::PROVIDER_NAME + ); + } } } let _ = writeln!( @@ -514,10 +538,11 @@ impl ExternalAgentRegistry { } let agent_id = format!( "{}-{}", - codex::PROVIDER_ID, + params.provider.trim(), self.next_agent.fetch_add(1, Ordering::Relaxed) ); let agent = Arc::new(ExternalAgent::new( + params.provider.trim().to_string(), agent_id.clone(), call.session_id.clone(), first_line_label(&prompt), @@ -556,6 +581,9 @@ impl ExternalAgentRegistry { let Some(agent) = self.agent(&call.session_id, ¶ms.agent_id).await else { return error_result(unknown_agent(¶ms.agent_id)); }; + if agent.provider != params.provider.trim() { + return error_result("This agent belongs to a different provider."); + } agent .run_turn( &call, @@ -580,6 +608,9 @@ impl ExternalAgentRegistry { let Some(agent) = self.agent(&call.session_id, ¶ms.agent_id).await else { return error_result(unknown_agent(¶ms.agent_id)); }; + if agent.provider != params.provider.trim() { + return error_result("This agent belongs to a different provider."); + } let activity = agent.activity().await; let guidance = if activity.status == "running" { background_guidance() @@ -599,6 +630,12 @@ impl ExternalAgentRegistry { if let Err(error) = Self::require_provider(¶ms.provider) { return error_result(error); } + let Some(agent) = self.agent(&call.session_id, ¶ms.agent_id).await else { + return error_result(unknown_agent(¶ms.agent_id)); + }; + if agent.provider != params.provider.trim() { + return error_result("This agent belongs to a different provider."); + } match self.cancel(&call.session_id, ¶ms.agent_id).await { Ok(activity) => { let mut result = @@ -730,8 +767,9 @@ struct StoredCall { } struct AgentProcess { + thread_id: String, child: ArmedShellChild, - client: Arc, + client: Arc, reader: tokio::task::JoinHandle<()>, events: tokio::task::JoinHandle<()>, } @@ -755,6 +793,8 @@ struct AgentState { } struct ExternalAgent { + provider: String, + launch: Mutex<()>, agent_id: String, session_id: String, task: String, @@ -769,7 +809,15 @@ struct ExternalAgent { } impl ExternalAgent { + fn provider_name(&self) -> &str { + if self.provider == claude::PROVIDER_ID { + claude::PROVIDER_NAME + } else { + codex::PROVIDER_NAME + } + } fn new( + provider: String, agent_id: String, session_id: String, task: String, @@ -784,7 +832,7 @@ impl ExternalAgent { thread_id: None, turn: None, activity: ExternalAgentActivity { - provider: codex::PROVIDER_ID.to_string(), + provider: provider.clone(), agent_id: agent_id.clone(), status: "idle".to_string(), ..Default::default() @@ -795,6 +843,8 @@ impl ExternalAgent { last_row_emit: None, last_row_id: None, }), + provider, + launch: Mutex::new(()), agent_id, session_id, task, @@ -826,7 +876,7 @@ impl ExternalAgent { elapsed_ms: self.started.elapsed().as_millis().min(u64::MAX as u128) as u64, activity: latest_activity_label(&state.activity), external: Some(ExternalAgentRef { - provider: codex::PROVIDER_ID.to_string(), + provider: self.provider.clone(), agent_id: self.agent_id.clone(), }), }) @@ -839,10 +889,25 @@ impl ExternalAgent { call: &ExternalAgentCall, input: TurnInput, ) -> CallToolResult { + let launch = self.launch.lock().await; if self.cancel.is_cancelled() { return error_result("This external agent has been shut down."); } - if let Err(error) = self.ensure_process(call, input.model.as_deref()).await { + if self.state.lock().await.turn.is_some() { + return error_result(format!( + "Agent {} is still working on its previous turn. Wait for Maple's notice, or check with {AGENT_STATUS_TOOL}.", + self.agent_id + )); + } + let ready = tokio::select! { + biased; + _ = self.cancel.cancelled() => Err("This external agent has been shut down.".to_string()), + _ = call.cancel_token.cancelled() => Err("The external agent launch was cancelled.".to_string()), + _ = call.tool_context.revoked.cancelled() => Err("The external agent context was revoked.".to_string()), + result = tokio::time::timeout(Duration::from_secs(30), self.ensure_process(call, input.model.as_deref(), input.effort.as_deref())) => + result.unwrap_or_else(|_| Err("The external agent did not initialize in time.".to_string())), + }; + if let Err(error) = ready { return error_result(error); } let (client, thread_id, done_rx) = { @@ -853,16 +918,11 @@ impl ExternalAgent { self.agent_id )); } - let Some(client) = state - .process - .as_ref() - .map(|process| Arc::clone(&process.client)) - else { + let Some(process) = state.process.as_ref() else { return error_result("The agent process is not running."); }; - let Some(thread_id) = state.thread_id.clone() else { - return error_result("The agent has no thread."); - }; + let client = Arc::clone(&process.client); + let thread_id = process.thread_id.clone(); let (done_tx, done_rx) = oneshot::channel(); let synthetic = call.row_id.is_none(); let row_id = call.row_id.clone().unwrap_or_else(|| { @@ -900,12 +960,13 @@ impl ExternalAgent { model: input.model.as_deref(), effort: input.effort.as_deref(), }); - if let Err(error) = client.request("turn/start", params).await { + if let Err(error) = client.request(RequestMethod::TurnStart, params).await { self.finish_turn(TurnOutcome::Failed, Some(error.clone())) .await; return error_result(error); } + drop(launch); if input.background { let agent = Arc::clone(self); tokio::spawn(async move { @@ -961,6 +1022,7 @@ impl ExternalAgent { self: &Arc, call: &ExternalAgentCall, model: Option<&str>, + effort: Option<&str>, ) -> Result<(), String> { { let state = self.state.lock().await; @@ -970,17 +1032,40 @@ impl ExternalAgent { return Ok(()); } } - let executable = codex::find_executable(call.login_path.as_deref()).ok_or_else(|| { - "Codex is not installed, or `codex` is not on PATH. Ask the user to install the Codex CLI.".to_string() - })?; + let existing_thread = self.state.lock().await.thread_id.clone(); + let claude_thread = existing_thread + .clone() + .unwrap_or_else(claude::new_session_id); + let (executable, args) = if self.provider == claude::PROVIDER_ID { + let executable = claude::find_executable(call.login_path.as_deref()) + .ok_or("Install Claude Code and make sure `claude` is on PATH.")?; + ( + executable, + claude::command_args(&claude_thread, existing_thread.is_some(), model, effort), + ) + } else { + let executable = codex::find_executable(call.login_path.as_deref()).ok_or_else(|| { + "Codex is not installed, or `codex` is not on PATH. Ask the user to install the Codex CLI.".to_string() + })?; + ( + executable, + codex::app_server_args() + .iter() + .map(|arg| arg.to_string()) + .collect(), + ) + }; let mut command = build_external_agent_command( &executable, - &codex::app_server_args(), - &self.host.project_root, + &args.iter().map(String::as_str).collect::>(), + &self.cwd, call.login_path.as_deref(), Some(&self.session_id), &call.tool_context, )?; + if self.provider == claude::PROVIDER_ID { + command.env_remove("CLAUDECODE"); + } command .stdin(Stdio::piped()) .stdout(Stdio::piped()) @@ -988,41 +1073,54 @@ impl ExternalAgent { .kill_on_drop(true); let mut child = { let _launch = call.tool_context.begin_process_launch(&call.cancel_token)?; - spawn_contained(command).map_err(|error| format!("Failed to start Codex: {error}"))? + spawn_contained(command) + .map_err(|error| format!("Failed to start {}: {error}", self.provider_name()))? }; let stdin = child .as_mut() .stdin() .take() - .ok_or_else(|| "Failed to open Codex's stdin".to_string())?; + .ok_or_else(|| "Failed to open the agent stdin".to_string())?; let stdout = child .as_mut() .stdout() .take() - .ok_or_else(|| "Failed to open Codex's stdout".to_string())?; - let (client, receiver, reader) = AppServerClient::new(stdin, stdout); + .ok_or_else(|| "Failed to open the agent stdout".to_string())?; + let (client, receiver, reader) = if self.provider == claude::PROVIDER_ID { + let (client, receiver, reader) = claude::Client::new(stdin, stdout, claude_thread); + (Arc::new(AgentClient::Claude(client)), receiver, reader) + } else { + let (client, receiver, reader) = AppServerClient::new(stdin, stdout); + (Arc::new(AgentClient::Codex(client)), receiver, reader) + }; client - .request("initialize", codex::initialize_params()) + .request(RequestMethod::Initialize, codex::initialize_params()) .await?; - client.notify("initialized", json!({})).await?; + client.initialized().await?; let thread_id = { let existing = self.state.lock().await.thread_id.clone(); let response = match &existing { Some(thread_id) => { client - .request("thread/resume", codex::thread_resume_params(thread_id)) + .request( + RequestMethod::ThreadResume, + codex::thread_resume_params(thread_id), + ) .await? } None => { client - .request("thread/start", codex::thread_start_params(&self.cwd, model)) + .request( + RequestMethod::ThreadStart, + codex::thread_start_params(&self.cwd, model), + ) .await? } }; match existing { Some(thread_id) => thread_id, None => codex::thread_id_from_response(&response) - .ok_or_else(|| "Codex did not report a thread ID".to_string())?, + .ok_or_else(|| "The agent did not report a thread ID".to_string())?, } }; let events = tokio::spawn(Arc::clone(self).consume_server_messages(receiver)); @@ -1031,9 +1129,14 @@ impl ExternalAgent { previous.reader.abort(); previous.events.abort(); } - state.thread_id = Some(thread_id.clone()); - state.activity.thread_id = Some(thread_id); + // Claude's initial UUID is provisional until the CLI confirms it. + // A launch failure before system/init must start fresh on retry. + if self.provider != claude::PROVIDER_ID { + state.thread_id = Some(thread_id.clone()); + state.activity.thread_id = Some(thread_id.clone()); + } state.process = Some(AgentProcess { + thread_id, child, client, reader, @@ -1046,8 +1149,34 @@ impl ExternalAgent { while let Some(message) = receiver.recv().await { match message { ServerMessage::Notification { method, params } => { - self.handle_event(codex::parse_notification(&method, ¶ms)) - .await; + let event = codex::parse_notification(&method, ¶ms); + if self.provider == claude::PROVIDER_ID + && let CodexEvent::TurnCompleted { ref status, .. } = event + { + // A Claude process serves one turn. Reclaim it before + // reporting completion, including on protocol failure. + // On success stdin is closed: give session writes time + // to flush, then clean up any remaining descendants. + // Hold the lifecycle lock through cleanup so shutdown + // cannot return while this task still owns a child. + let mut state = self.state.lock().await; + if let Some(mut process) = state.process.take() { + if status == "completed" { + let _ = tokio::time::timeout( + Duration::from_secs(2), + process.child.as_mut().wait(), + ) + .await; + } + process.child.kill_and_wait().await; + process.reader.abort(); + // Do not abort process.events: it is this task. + } + drop(state); + self.handle_event(event).await; + return; + } + self.handle_event(event).await; } ServerMessage::Request { id, method, params } => { // An approval waits on the user. It must not stall the @@ -1225,6 +1354,21 @@ impl ExternalAgent { None => return, } }; + if self.provider == claude::PROVIDER_ID && method == "claude/tool/requestApproval" { + let tool = params["tool"].as_str().unwrap_or("tool"); + let arguments = params["input"].as_object().cloned().unwrap_or_default(); + let request = AgentPermissionRequest { + request_id: format!("{}-{}", self.agent_id, id.as_str().unwrap_or("request")), + tool_name: "claude_tool".into(), + arguments, + prompt: Some(format!("Claude Code wants to use {tool}")), + }; + let decision = self + .request_permission(request, format!("use {tool}")) + .await; + let _ = client.respond(id, codex::approval_response(decision)).await; + return; + } let response = match codex::parse_server_request(method, ¶ms) { CodexServerRequest::CommandApproval { item_id, @@ -1366,7 +1510,7 @@ impl ExternalAgent { if let Some((client, thread_id, turn_id)) = steer { match client .request( - "turn/steer", + RequestMethod::TurnSteer, codex::turn_steer_params(&thread_id, &turn_id, &prompt), ) .await @@ -1435,6 +1579,9 @@ impl ExternalAgent { request: AgentPermissionRequest, summary: String, ) -> AgentPermissionDecision { + if self.cancel.is_cancelled() || self.turn_ended().await.is_cancelled() { + return AgentPermissionDecision::Cancel; + } { let modes = self.host.permission_modes.lock().await; if modes @@ -1521,17 +1668,17 @@ impl ExternalAgent { }; ( Arc::clone(&process.client), - state.thread_id.clone(), + process.thread_id.clone(), turn.turn_id.clone(), ) }; - let (Some(thread_id), Some(turn_id)) = (thread_id, turn_id) else { + let Some(turn_id) = turn_id else { // The turn has not been identified yet; the agent will report // it and the caller's settle timeout ends the wait. return; }; let request = client.request( - "turn/interrupt", + RequestMethod::TurnInterrupt, codex::turn_interrupt_params(&thread_id, &turn_id), ); if let Ok(Err(error)) = tokio::time::timeout(INTERRUPT_REQUEST_TIMEOUT, request).await { @@ -1649,13 +1796,12 @@ impl ExternalAgent { let result_text = render_activity(&activity, &completion_guidance(&activity)); let row_id = self.last_row_id().await; let for_model = background_result_message( - &format!("external agent {} ({})", self.agent_id, codex::PROVIDER_ID), + &format!("external agent {} ({})", self.agent_id, self.provider), outcome.status(), &result_text, &format!( "Use {AGENT_STATUS_TOOL}(provider: \"{}\", agent_id: \"{}\") only if you need to inspect its current state again.", - codex::PROVIDER_ID, - self.agent_id + self.provider, self.agent_id ), ); let delivered = self @@ -1713,7 +1859,7 @@ impl ExternalAgent { format!( "External agent {} ({}) {}. {}", self.agent_id, - codex::PROVIDER_NAME, + self.provider_name(), match outcome { TurnOutcome::Completed => "finished", TurnOutcome::Failed => "failed", @@ -1779,7 +1925,7 @@ impl ExternalAgent { task: self.task.clone(), background, external: Some(ExternalAgentRef { - provider: codex::PROVIDER_ID.to_string(), + provider: self.provider.clone(), agent_id: self.agent_id.clone(), }), }, diff --git a/apps/maple-agent/crates/maple-agent/src/agent/external_agents/tests.rs b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/tests.rs index d4f48c511..0d937425e 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent/external_agents/tests.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/tests.rs @@ -1,7 +1,7 @@ -//! Driver tests against a fake `codex app-server`. +//! Driver tests against fake Codex and Claude Code CLIs. //! -//! The fixture is this test binary re-executed as an ignored test. A shell -//! shim named `codex` on a private PATH forwards to it, so the driver +//! Each fixture is this test binary re-executed as an ignored test. Shell +//! shims on a private PATH forward to them, so the driver //! resolves and spawns it exactly as it would the real CLI. Unix only until //! a `.cmd` shim exists for Windows. @@ -11,14 +11,16 @@ use super::*; use crate::agent::tool_context::default_tool_context_spec; use crate::agent::{AgentEventSink, AgentPathLayout, MapleAgentHostResources}; use std::io::{BufRead, Write}; +use std::os::fd::FromRawFd; use std::os::unix::fs::PermissionsExt; -const FIXTURE_TEST: &str = "agent::external_agents::tests::fake_codex_app_server"; -const FIXTURE_MARKER: &str = "MAPLE_FAKE_CODEX"; -const FIXTURE_ARGS: &str = "MAPLE_FAKE_CODEX_ARGS"; -const FIXTURE_MODE: &str = "MAPLE_FAKE_CODEX_MODE"; -const FIXTURE_PID_FILE: &str = "MAPLE_FAKE_CODEX_PID_FILE"; -const FIXTURE_LOG: &str = "MAPLE_FAKE_CODEX_LOG"; +mod claude_fixture; + +const FIXTURE_MARKER: &str = "MAPLE_FAKE_AGENT"; +const FIXTURE_ARGS: &str = "MAPLE_FAKE_AGENT_ARGS"; +const FIXTURE_MODE: &str = "MAPLE_FAKE_AGENT_MODE"; +const FIXTURE_PID_FILE: &str = "MAPLE_FAKE_AGENT_PID_FILE"; +const FIXTURE_LOG: &str = "MAPLE_FAKE_AGENT_LOG"; const WAIT: Duration = Duration::from_secs(20); #[derive(Default)] @@ -66,17 +68,8 @@ impl Harness { fs::create_dir_all(&history).unwrap(); let shim_dir = root.join("bin"); fs::create_dir_all(&shim_dir).unwrap(); - let pid_file = root.join("codex.pid"); - let log_file = root.join("codex.log"); - let shim = format!( - "#!/bin/sh\nexport {FIXTURE_MARKER}=1\nexport {FIXTURE_MODE}='{mode}'\nexport {FIXTURE_PID_FILE}='{}'\nexport {FIXTURE_LOG}='{}'\nexport {FIXTURE_ARGS}=\"$*\"\nexec '{}' '{FIXTURE_TEST}' --exact --ignored --nocapture --test-threads=1\n", - pid_file.display(), - log_file.display(), - std::env::current_exe().unwrap().display(), - ); - let shim_path = shim_dir.join("codex"); - fs::write(&shim_path, shim).unwrap(); - fs::set_permissions(&shim_path, fs::Permissions::from_mode(0o700)).unwrap(); + let pid_file = root.join("agent.pid"); + let log_file = root.join("agent.log"); let paths = AgentPathLayout::from_app_roots(root.join("config"), root.join("data")); let sink = Arc::new(RecordingSink::default()); @@ -100,7 +93,7 @@ impl Harness { lifetime: CancellationToken::new(), }; let registry = Arc::new(ExternalAgentRegistry::new(host.clone())); - Self { + let harness = Self { _temp: temp, project, shim_dir, @@ -110,7 +103,29 @@ impl Harness { registry, pid_file, log_file, - } + }; + harness.install_fixture("codex", mode); + harness.install_fixture("claude", mode); + harness + } + + fn install_fixture(&self, provider: &str, mode: &str) { + let test = match provider { + "codex" => "fake_codex_app_server", + "claude" => "claude_fixture::run", + _ => panic!("unknown fixture provider"), + }; + // Keep libtest's status output off both protocol pipes. Its output + // can race the fixture, so inserting a newline is not sufficient. + let shim = format!( + "#!/bin/sh\nexport {FIXTURE_MARKER}=1\nexport {FIXTURE_MODE}='{mode}'\nexport {FIXTURE_PID_FILE}='{}'\nexport {FIXTURE_LOG}='{}'\nexport {FIXTURE_ARGS}=\"$*\"\nexec '{}' 'agent::external_agents::tests::{test}' --exact --ignored --nocapture --test-threads=1 3>&1 1>/dev/null\n", + self.pid_file.display(), + self.log_file.display(), + std::env::current_exe().unwrap().display(), + ); + let path = self.shim_dir.join(provider); + fs::write(&path, shim).unwrap(); + fs::set_permissions(path, fs::Permissions::from_mode(0o700)).unwrap(); } fn call(&self, session_id: &str, row_id: &str) -> ExternalAgentCall { @@ -147,14 +162,18 @@ impl Harness { } } -async fn wait_for(mut probe: impl FnMut() -> Option) -> T { - let deadline = Instant::now() + WAIT; - loop { - if let Some(value) = probe() { - return value; +#[track_caller] +fn wait_for(mut probe: impl FnMut() -> Option) -> impl Future { + let caller = std::panic::Location::caller(); + async move { + let deadline = Instant::now() + WAIT; + loop { + if let Some(value) = probe() { + return value; + } + assert!(Instant::now() < deadline, "timed out waiting at {caller}"); + tokio::time::sleep(Duration::from_millis(25)).await; } - assert!(Instant::now() < deadline, "timed out waiting"); - tokio::time::sleep(Duration::from_millis(25)).await; } } @@ -171,6 +190,13 @@ fn process_alive(pid: i32) -> bool { unsafe { libc::kill(pid, 0) == 0 } } +fn fixture_output() -> fs::File { + // SAFETY: install_fixture's shim duplicates the protocol pipe to fd 3 + // before redirecting libtest stdout. Each fixture calls this once and + // this File is the sole owner of that descriptor in the child process. + unsafe { fs::File::from_raw_fd(3) } +} + /// The fake app-server. It answers the handshake, starts a thread, and /// on `turn/start` plays a short turn that asks for one command approval /// and reports what decision it got in its final message. In `slow` mode @@ -181,11 +207,7 @@ fn fake_codex_app_server() { if std::env::var_os(FIXTURE_MARKER).is_none() { return; } - let stdout = std::io::stdout(); - let mut out = stdout.lock(); - // libtest prints "test ... " with no newline before the test - // runs; end that line so the first protocol line stands alone. - writeln!(out).unwrap(); + let mut out = fixture_output(); let args = std::env::var(FIXTURE_ARGS).unwrap_or_default(); if args.contains("--version") { writeln!(out, "codex-cli 0.150.0").unwrap(); @@ -200,6 +222,10 @@ fn fake_codex_app_server() { let stdin = std::io::stdin(); let mut lines = stdin.lock().lines(); let mut send = |value: Value| { + // Reproduce libtest status text arriving after fixture startup, + // without a newline. It must not corrupt a protocol response. + print!("fixture status"); + std::io::stdout().flush().unwrap(); writeln!(out, "{value}").unwrap(); out.flush().unwrap(); }; @@ -527,7 +553,7 @@ async fn cancelling_the_run_interrupts_the_turn_and_shutdown_kills_the_process() let registry = Arc::clone(&harness.registry); let call = harness.call("session-3", "row-3"); let cancel = call.cancel_token.clone(); - let turn = tokio::spawn(async move { + let mut turn = tokio::spawn(async move { registry .start( call, @@ -542,19 +568,25 @@ async fn cancelling_the_run_interrupts_the_turn_and_shutdown_kills_the_process() ) .await }); - let pid = harness.fixture_pid().await; - // Wait until the turn is identified, so the interrupt can name it. - { - let sink = Arc::clone(&harness.sink); - wait_for(|| { - sink.events() - .iter() - .any(|event| matches!(event, AgentServiceEvent::TimelineItem { item, .. } if item.id == "row-3")) - .then_some(()) - }) - .await; - } - tokio::time::sleep(Duration::from_millis(200)).await; + let pid = tokio::select! { + pid = harness.fixture_pid() => pid, + result = &mut turn => panic!("Codex fixture stopped before startup: {}", result_text(&result.unwrap())), + }; + // The initial timeline row precedes turn/start. Wait for the actual + // turn ID instead of guessing when the notification has been consumed. + let agent = harness + .registry + .agent("session-3", "codex-1") + .await + .unwrap(); + wait_for(|| { + agent + .state + .try_lock() + .ok() + .and_then(|state| state.turn.as_ref()?.turn_id.as_ref().map(|_| ())) + }) + .await; cancel.cancel(); let result = turn.await.unwrap(); let text = result_text(&result); @@ -734,7 +766,7 @@ async fn registry_rejects_unknown_providers_bad_cwd_and_too_many_agents() { }; let unknown = harness .registry - .start(harness.call("session-5", "r"), start("claude", None)) + .start(harness.call("session-5", "r"), start("unknown", None)) .await; assert_eq!(unknown.is_error, Some(true)); assert!(result_text(&unknown).contains("Unknown agent provider")); @@ -769,6 +801,7 @@ async fn registry_rejects_unknown_providers_bad_cwd_and_too_many_agents() { session.agents.insert( agent_id.clone(), Arc::new(ExternalAgent::new( + "codex".into(), agent_id, "session-5".into(), "idle".into(), @@ -952,3 +985,415 @@ async fn an_async_question_is_answered_in_a_turn_of_maples_own() { ); harness.registry.shutdown_all(Duration::from_secs(5)).await; } + +fn claude_start(background: bool) -> AgentStartParams { + AgentStartParams { + provider: "claude".into(), + prompt: "Fix the fixture".into(), + background, + model: Some("sonnet".into()), + effort: Some("high".into()), + cwd: Some("sub".into()), + } +} + +#[tokio::test] +async fn claude_detection_distinguishes_sign_in_from_probe_failures() { + for (mode, expected) in [ + ("auth-in", Some(true)), + ("auth-out", Some(false)), + ("auth-error", None), + ("auth-inconsistent", None), + ("auth-missing", None), + ("auth-oversized", None), + ("auth-timeout", None), + ] { + let harness = Harness::new(mode); + let detection = claude::detect(harness.shim_dir.to_str()).await; + assert_eq!(detection.signed_in, expected, "{mode}"); + assert!(detection.version.is_some(), "{mode}"); + assert!(detection.problem.is_none(), "{mode}"); + assert!(!format!("{detection:?}").contains("secret-canary")); + let pid = harness.fixture_pid().await; + assert!( + !process_alive(pid), + "{mode}: authentication probe left running" + ); + } +} + +/// Exercises the native Rust transport against a deterministic CLI, with no inference. +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn claude_native_streams_resumes_and_binds_the_provider() { + let harness = Harness::new("approve"); + fs::create_dir(harness.project.join("sub")).unwrap(); + harness.set_mode("claude-task", GooseMode::Auto).await; + let result = harness + .registry + .start(harness.call("claude-task", "r1"), claude_start(false)) + .await; + let activity = result + .structured_content + .as_ref() + .unwrap_or_else(|| panic!("{}", result_text(&result)))[ACTIVITY_KEY] + .clone(); + assert_eq!(activity["status"], "completed", "{activity}"); + assert_eq!(activity["provider"], "claude"); + assert_eq!(activity["text"], "Allowed"); + assert_eq!(activity["commands"][0]["status"], "completed"); + assert_eq!(activity["fileChanges"][0]["path"], "src/lib.rs"); + assert_eq!(activity["todos"][0]["completed"], true); + let agent_id = activity["agentId"].as_str().unwrap().to_string(); + let thread_id = activity["threadId"].as_str().unwrap().to_string(); + let result = harness + .registry + .send( + harness.call("claude-task", "r2"), + AgentSendParams { + provider: "claude".into(), + agent_id: agent_id.clone(), + prompt: "Continue".into(), + background: false, + model: None, + effort: None, + }, + ) + .await; + assert!( + result_text(&result).contains("Status: completed"), + "{}", + result_text(&result) + ); + let log = harness.log(); + assert!(log.contains("--resume")); + assert!( + log.contains("stdin_closed"), + "completed CLI should exit before resumption" + ); + assert!(log.contains(&thread_id)); + assert!(log.contains("--effort")); + assert!(log.contains(&harness.project.join("sub").to_string_lossy().into_owned())); + assert!(!log.contains("dangerously-skip-permissions")); + for provider in ["codex", "unknown"] { + let result = harness + .registry + .cancel_tool( + &harness.call("claude-task", "r3"), + AgentRefParams { + provider: provider.into(), + agent_id: agent_id.clone(), + }, + ) + .await; + assert_eq!(result.is_error, Some(true)); + } + let result = harness + .registry + .status( + &harness.call("another-task", "r4"), + AgentRefParams { + provider: "claude".into(), + agent_id, + }, + ) + .await; + assert_eq!(result.is_error, Some(true)); + let providers = harness + .registry + .list_providers(&harness.call("claude-task", "list"), &["claude".into()]) + .await; + assert!(result_text(&providers).contains("- claude:")); + assert!(!result_text(&providers).contains("- codex:")); + harness.registry.shutdown_all(Duration::from_secs(5)).await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn claude_native_retries_only_resume_confirmed_sessions() { + for (mode, confirmed) in [ + ("pre-init-error", false), + ("pre-init-eof", false), + ("error", true), + ("no-init-success", true), + ] { + let harness = Harness::new(mode); + fs::create_dir(harness.project.join("sub")).unwrap(); + harness.set_mode("retry", GooseMode::Auto).await; + let result = harness + .registry + .start(harness.call("retry", "r1"), claude_start(false)) + .await; + let activity = &result.structured_content.as_ref().unwrap()[ACTIVITY_KEY]; + assert_eq!( + activity["threadId"].is_string(), + confirmed, + "{mode}: {activity}" + ); + assert_eq!( + activity["status"], + if mode == "no-init-success" { + "completed" + } else { + "failed" + } + ); + let agent_id = activity["agentId"].as_str().unwrap().to_string(); + harness.install_fixture("claude", "approve"); + for row in ["r2", "r3"] { + let result = harness + .registry + .send( + harness.call("retry", row), + AgentSendParams { + provider: "claude".into(), + agent_id: agent_id.clone(), + prompt: "Try again".into(), + background: false, + model: None, + effort: None, + }, + ) + .await; + assert_eq!( + result.structured_content.unwrap()[ACTIVITY_KEY]["status"], + "completed" + ); + } + let args: Vec = harness + .log() + .lines() + .map(|line| { + serde_json::from_str::(line) + .unwrap_or_else(|error| panic!("{mode}: invalid fixture log record: {error}")) + }) + .filter_map(|value| value.get("args").cloned()) + .collect(); + assert_eq!(args.len(), 3); + let session = |args: &Value| { + args.as_array() + .unwrap() + .windows(2) + .find_map(|pair| { + (pair[0] == "--session-id" || pair[0] == "--resume").then(|| pair[1].clone()) + }) + .unwrap() + }; + assert!(args[0].as_array().unwrap().contains(&json!("--session-id"))); + assert!(args[1].as_array().unwrap().contains(&json!(if confirmed { + "--resume" + } else { + "--session-id" + }))); + assert_eq!(session(&args[0]) == session(&args[1]), confirmed); + assert!(args[2].as_array().unwrap().contains(&json!("--resume"))); + assert_eq!(session(&args[1]), session(&args[2])); + harness.registry.shutdown_all(Duration::from_secs(5)).await; + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn claude_native_cancellation_kills_cli_and_descendants() { + let harness = Harness::new("slow"); + fs::create_dir(harness.project.join("sub")).unwrap(); + let result = harness + .registry + .start(harness.call("claude-stop", "r1"), claude_start(true)) + .await; + assert!( + result_text(&result).contains("claude-1"), + "{}", + result_text(&result) + ); + let pid = harness.fixture_pid().await; + let child_pid = wait_for(|| { + fs::read_to_string(harness.pid_file.with_extension("pid.child")) + .ok() + .and_then(|pid| pid.parse::().ok()) + }) + .await; + let activity = harness + .registry + .cancel("claude-stop", "claude-1") + .await + .unwrap(); + assert_eq!(activity.status, "cancelled"); + wait_for(|| (!process_alive(pid)).then_some(())).await; + // A killed orphan can briefly remain a zombie before PID 1 reaps it. + wait_for(|| { + (!process_alive(child_pid) + || fs::read_to_string(format!("/proc/{child_pid}/stat")) + .is_ok_and(|stat| stat.contains(") Z "))) + .then_some(()) + }) + .await; + // Reconnect through a new CLI process and resume the session saved before Stop. + harness.install_fixture("claude", "approve"); + harness.set_mode("claude-stop", GooseMode::Auto).await; + let result = harness + .registry + .send( + harness.call("claude-stop", "r2"), + AgentSendParams { + provider: "claude".into(), + agent_id: "claude-1".into(), + prompt: "Resume".into(), + background: false, + model: None, + effort: None, + }, + ) + .await; + assert!( + result_text(&result).contains("Status: completed"), + "{}", + result_text(&result) + ); + assert_eq!( + result.structured_content.unwrap()[ACTIVITY_KEY]["threadId"].as_str(), + activity.thread_id.as_deref() + ); + harness.registry.shutdown_all(Duration::from_secs(5)).await; +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn claude_native_permission_denial_and_questions_use_maple_brokers() { + for mode in ["approve", "question", "multi-question"] { + let harness = Harness::new(mode); + fs::create_dir(harness.project.join("sub")).unwrap(); + let registry = harness.registry.clone(); + let call = harness.call("claude-input", "r1"); + let turn = tokio::spawn(async move { registry.start(call, claude_start(false)).await }); + if mode == "approve" { + let pending = wait_for(|| { + harness + .service + .pending_permissions + .try_lock() + .ok() + .and_then(|entries| { + entries + .iter() + .next() + .map(|(id, entry)| (id.clone(), entry.clone())) + }) + }) + .await; + let (key, entry) = pending; + assert_eq!(entry.request.tool_name, "claude_tool"); + assert_eq!(entry.request.arguments["command"], "cargo test"); + harness + .service + .pending_permissions + .lock() + .await + .remove(&key); + let PendingPermissionOrigin::ExternalAgent(responder) = entry.origin else { + panic!("external responder expected") + }; + assert!(responder.resolve(AgentPermissionDecision::DenyOnce)); + assert!(!responder.resolve(AgentPermissionDecision::AllowOnce)); + } else { + let (request_id, questions) = wait_for(|| { + harness.sink.events().iter().find_map(|event| match event { + AgentServiceEvent::Question { + session_id, + request_id, + questions, + } if session_id == "claude-input" => { + Some((request_id.clone(), questions.clone())) + } + _ => None, + }) + }) + .await; + assert_eq!(questions[0].question, "Tabs or spaces?"); + assert_eq!(questions[0].multi_select, mode == "multi-question"); + assert!( + harness + .service + .answer_question( + &request_id, + if mode == "multi-question" { + r#"{"answers":{"q0":{"answers":["Spaces","Tabs"]}}}"# + } else { + r#"{"answers":{"q0":{"answers":["Spaces"]}}}"# + } + .into() + ) + .await + ); + } + let result = tokio::time::timeout(WAIT, turn).await.unwrap().unwrap(); + let text = result_text(&result); + assert!( + text.contains(if mode == "approve" { + "Denied" + } else if mode == "multi-question" { + "Spaces, Tabs" + } else { + "Spaces" + }), + "{text}" + ); + harness.registry.shutdown_all(Duration::from_secs(5)).await; + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn claude_native_failures_are_not_success_or_unsanitized_output() { + for mode in ["eof", "error", "init-error", "malformed", "oversized"] { + let harness = Harness::new(mode); + fs::create_dir(harness.project.join("sub")).unwrap(); + let result = harness + .registry + .start(harness.call("claude-fail", "r1"), claude_start(false)) + .await; + let text = result_text(&result); + assert!( + text.contains("Status: failed") || result.is_error == Some(true), + "{mode}: {text}" + ); + assert!(!text.contains("secret-canary")); + let pid = harness.fixture_pid().await; + wait_for(|| (!process_alive(pid)).then_some(())).await; + harness.registry.shutdown_all(Duration::from_secs(5)).await; + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn claude_native_stop_withdraws_pending_permission() { + let harness = Harness::new("approve"); + fs::create_dir(harness.project.join("sub")).unwrap(); + let result = harness + .registry + .start(harness.call("claude-pending", "r1"), claude_start(true)) + .await; + assert!(result_text(&result).contains("claude-1")); + let key = wait_for(|| { + harness + .service + .pending_permissions + .try_lock() + .ok() + .and_then(|pending| pending.keys().next().cloned()) + }) + .await; + let pid = harness.fixture_pid().await; + let activity = harness + .registry + .cancel("claude-pending", "claude-1") + .await + .unwrap(); + assert_eq!(activity.status, "cancelled"); + wait_for(|| (!process_alive(pid)).then_some(())).await; + wait_for(|| { + harness + .service + .pending_permissions + .try_lock() + .ok() + .and_then(|pending| (!pending.contains_key(&key)).then_some(())) + }) + .await; + assert!(!harness.log().contains("\"behavior\":\"allow\"")); + harness.registry.shutdown_all(Duration::from_secs(5)).await; +} diff --git a/apps/maple-agent/crates/maple-agent/src/agent/external_agents/tests/claude_fixture.rs b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/tests/claude_fixture.rs new file mode 100644 index 000000000..b5d3fd752 --- /dev/null +++ b/apps/maple-agent/crates/maple-agent/src/agent/external_agents/tests/claude_fixture.rs @@ -0,0 +1,267 @@ +//! Native CLI fixture for the Claude transport. Reuses the driver test binary. + +use super::{ + FIXTURE_ARGS, FIXTURE_LOG, FIXTURE_MARKER, FIXTURE_MODE, FIXTURE_PID_FILE, fixture_output, +}; +use serde_json::{Value, json}; +use std::fs::{self, File, OpenOptions}; +use std::io::{BufRead, Write}; +use std::process::Command; + +fn log_record(log: &mut File, message: &Value) { + // Maple may kill the fixture at any time. Append the JSON and newline + // together so a partial formatted record cannot swallow the next launch. + let mut line = serde_json::to_vec(message).unwrap(); + line.push(b'\n'); + log.write_all(&line).unwrap(); +} + +fn send(out: &mut File, message: Value) { + writeln!(out, "{message}").unwrap(); + out.flush().unwrap(); +} + +fn finish(out: &mut File, session: &str, text: &str) { + send( + out, + json!({"type": "assistant", "message": { + "id": "message-1", "role": "assistant", "model": "fixture", + "content": [{"type": "text", "text": text}], + }}), + ); + send( + out, + json!({"type": "result", "subtype": "success", "is_error": false, + "session_id": session, "result": text, + }), + ); +} + +#[test] +#[ignore = "fake Claude CLI run by the driver tests"] +fn run() { + if std::env::var_os(FIXTURE_MARKER).is_none() { + return; + } + let mut out = fixture_output(); + let args = std::env::var(FIXTURE_ARGS).unwrap(); + let args: Vec<_> = args.split_whitespace().collect(); + if args.contains(&"--version") { + writeln!(out, "2.1.270 (Claude Code)").unwrap(); + return; + } + let mode = std::env::var(FIXTURE_MODE).unwrap(); + if mode == "sleeper" { + // A descendant the "slow" CLI leaves behind: this same binary, + // so the test assumes no `sleep` on the host's paths. + std::thread::sleep(std::time::Duration::from_secs(1000)); + return; + } + if args == ["auth", "status", "--json"] { + fs::write( + std::env::var(FIXTURE_PID_FILE).unwrap(), + std::process::id().to_string(), + ) + .unwrap(); + let exit = match mode.as_str() { + "auth-out" => { + send(&mut out, json!({"loggedIn": false})); + 1 + } + "auth-error" => { + writeln!(out, "secret-canary: unsupported command").unwrap(); + 1 + } + "auth-inconsistent" => { + send(&mut out, json!({"loggedIn": true})); + 1 + } + "auth-missing" => { + send(&mut out, json!({"email": "secret-canary"})); + 0 + } + "auth-oversized" => { + send( + &mut out, + json!({"loggedIn": true, "extra": "x".repeat(20 * 1024)}), + ); + 0 + } + "auth-timeout" => { + std::thread::sleep(std::time::Duration::from_secs(30)); + 0 + } + _ => { + send( + &mut out, + json!({"loggedIn": true, "email": "secret-canary"}), + ); + 0 + } + }; + std::process::exit(exit); + } + let session = args + .windows(2) + .find_map(|pair| matches!(pair[0], "--session-id" | "--resume").then_some(pair[1])) + .expect("a new or resumed Claude session"); + let pid_file = std::env::var(FIXTURE_PID_FILE).unwrap(); + fs::write(&pid_file, std::process::id().to_string()).unwrap(); + let mut log = OpenOptions::new() + .create(true) + .append(true) + .open(std::env::var(FIXTURE_LOG).unwrap()) + .unwrap(); + log_record( + &mut log, + &json!({"args": args, "cwd": std::env::current_dir().unwrap()}), + ); + + for line in std::io::stdin().lock().lines() { + let message: Value = serde_json::from_str(&line.unwrap()).unwrap(); + log_record(&mut log, &message); + match message["type"].as_str() { + Some("control_request") => { + if mode == "init-error" { + send( + &mut out, + json!({"type": "control_response", "response": { + "subtype": "error", "request_id": message["request_id"], "error": "secret-canary", + }}), + ); + continue; + } + send( + &mut out, + json!({"type": "control_response", "response": { + "subtype": "success", "request_id": message["request_id"], "response": {}, + }}), + ); + if message["request"]["subtype"] == "interrupt" { + break; + } + } + Some("user") => { + if mode == "pre-init-eof" { + return; + } + if mode == "pre-init-error" { + send( + &mut out, + json!({"type": "result", "subtype": "error_during_execution", + "is_error": true, "session_id": session, "errors": ["secret-canary"], + }), + ); + continue; + } + if mode != "no-init-success" { + send( + &mut out, + json!({"type": "system", "subtype": "init", "session_id": session}), + ); + } + match mode.as_str() { + "eof" => return, + "malformed" | "oversized" => { + let line = if mode == "malformed" { + "not-json".into() + } else { + "x".repeat(4 * 1024 * 1024 + 1) + }; + writeln!(out, "{line}").unwrap(); + out.flush().unwrap(); + continue; + } + "error" => { + send( + &mut out, + json!({"type": "result", "subtype": "error_during_execution", + "is_error": true, "session_id": session, "errors": ["secret-canary"], + }), + ); + continue; + } + "slow" => { + // Deliberately outlives the CLI to test Maple's process + // group cleanup, including descendants after interrupt. + #[allow(clippy::zombie_processes)] + let child = Command::new(std::env::current_exe().unwrap()) + .args(std::env::args().skip(1)) + .env(FIXTURE_MODE, "sleeper") + .spawn() + .unwrap(); + fs::write(format!("{pid_file}.child"), child.id().to_string()).unwrap(); + continue; + } + _ => {} + } + send( + &mut out, + json!({"type": "stream_event", "event": { + "type": "message_start", "message": {"id": "message-1"}, + }}), + ); + send( + &mut out, + json!({"type": "stream_event", "event": { + "type": "content_block_delta", "delta": {"type": "text_delta", "text": "Working"}, + }}), + ); + let (tool, input) = if matches!(mode.as_str(), "question" | "multi-question") { + ( + "AskUserQuestion", + json!({"questions": [{"header": "Style", "question": "Tabs or spaces?", + "multiSelect": mode == "multi-question", + "options": [{"label": "Spaces", "description": "Use spaces"}, + {"label": "Tabs", "description": "Use tabs"}], + }]}), + ) + } else { + ("Bash", json!({"command": "cargo test"})) + }; + send( + &mut out, + json!({"type": "control_request", "request_id": "permission-1", + "request": {"subtype": "can_use_tool", "tool_name": tool, "input": input, "tool_use_id": "tool-1"}, + }), + ); + } + Some("control_response") => { + let answer = &message["response"]["response"]; + if matches!(mode.as_str(), "question" | "multi-question") { + finish( + &mut out, + session, + &answer["updatedInput"]["answers"].to_string(), + ); + continue; + } + let allowed = answer["behavior"] == "allow"; + if allowed { + send( + &mut out, + json!({"type": "assistant", "message": {"id": "tools", "content": [ + {"type": "tool_use", "id": "c1", "name": "Bash", "input": {"command": "cargo test"}}, + {"type": "tool_use", "id": "f1", "name": "Edit", "input": {"file_path": "src/lib.rs"}}, + {"type": "tool_use", "id": "todo", "name": "TodoWrite", "input": {"todos": [{"content": "Test", "status": "completed"}]}}, + ]}}), + ); + send( + &mut out, + json!({"type": "user", "message": {"content": [ + {"type": "tool_result", "tool_use_id": "c1", "content": "ok"}, + {"type": "tool_result", "tool_use_id": "f1", "content": "ok"}, + ]}}), + ); + } + finish( + &mut out, + session, + if allowed { "Allowed" } else { "Denied" }, + ); + } + _ => panic!("unexpected Claude fixture input"), + } + } + log_record(&mut log, &json!({"event": "stdin_closed"})); +} diff --git a/apps/maple-agent/crates/maple-agent/src/agent/integrations.rs b/apps/maple-agent/crates/maple-agent/src/agent/integrations.rs index 7b87464ff..ac31fc94d 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent/integrations.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent/integrations.rs @@ -5,6 +5,7 @@ //! Maple-hosted implementation. A task freezes that backend choice when it is //! created; account defaults never rewrite existing tasks. +use super::external_agents::claude::{self, ClaudeDetection}; use super::external_agents::codex::{self, CodexDetection}; use super::*; use std::collections::HashSet; @@ -44,13 +45,27 @@ pub(super) struct ExternalAgentIntegration { project: fn(&IntegrationDetections, Option<&StoredIntegration>) -> AgentIntegration, } -pub(super) const EXTERNAL_AGENT_INTEGRATIONS: &[ExternalAgentIntegration] = - &[ExternalAgentIntegration { +pub(super) const EXTERNAL_AGENT_INTEGRATIONS: &[ExternalAgentIntegration] = &[ + ExternalAgentIntegration { id: CODEX_INTEGRATION_ID, name: CODEX_CARD_NAME, description: CODEX_CARD_DESCRIPTION, project: |detections, stored| codex_public(&detections.codex, stored), - }]; + }, + ExternalAgentIntegration { + id: claude::PROVIDER_ID, + name: claude::PROVIDER_NAME, + description: "Let a task hand work to the Claude Code CLI installed on this computer, with its own account.", + project: |detections, stored| claude_public(&detections.claude, stored), + }, +]; + +impl AgentIntegration { + /// Whether this card is an external coding agent in the runtime catalog. + pub fn is_external_agent(&self) -> bool { + external_agent_selection(&self.id).is_some() + } +} pub(super) fn external_agent_selection(id: &str) -> Option<&'static ExternalAgentIntegration> { EXTERNAL_AGENT_INTEGRATIONS @@ -194,6 +209,7 @@ impl StoredIntegrationRegistry { pub(super) struct IntegrationDetections { pub(super) cua: CuaDetection, pub(super) codex: CodexDetection, + pub(super) claude: ClaudeDetection, } /// The Integrations card for Codex. Availability comes from the @@ -237,6 +253,40 @@ fn codex_public( } } +fn claude_public( + detection: &ClaudeDetection, + stored: Option<&StoredIntegration>, +) -> AgentIntegration { + let descriptor = external_agent_selection(claude::PROVIDER_ID).expect("Claude catalog entry"); + AgentIntegration { + id: descriptor.id.into(), + name: descriptor.name.into(), + description: descriptor.description.into(), + availability: match (&detection.executable, &detection.problem) { + (None, _) => AgentIntegrationAvailability::NotDetected, + (_, Some(_)) => AgentIntegrationAvailability::SetupRequired, + _ => AgentIntegrationAvailability::Available, + }, + backend: stored.map(|entry| entry.backend), + version: detection.version.clone(), + standalone_version: None, + permissions: None, + setup_available: false, + enabled_for_new_tasks: stored.is_some_and(|entry| entry.enabled), + detail: Some(if detection.executable.is_none() { + "Install Claude Code and make sure `claude` is on PATH, then reopen this page.".into() + } else if let Some(problem) = &detection.problem { + problem.clone() + } else { + match detection.signed_in { + Some(true) => "Signed in.".into(), + Some(false) => claude::sign_in_hint().into(), + None => "Could not check sign-in. Run `claude auth status` in a terminal.".into(), + } + }), + } +} + /// Whether tasks of this account may delegate to external agents. Read /// on a path that must keep working, so an unusable file reads as off. pub(super) fn external_agents_enabled(paths: &AgentPathLayout, user_id: &str) -> bool { @@ -438,13 +488,14 @@ fn embedded_cua_version() -> Option { None } -/// Discover every curated integration. `codex_search_path` is the PATH -/// to look on for the Codex CLI, when the host knows a fuller one than +/// Discover every curated integration. `search_path` is the PATH +/// to look on for external agents, when the host knows a fuller one than /// the process environment (a macOS GUI launch). -pub(super) async fn detect_integrations(codex_search_path: Option<&str>) -> IntegrationDetections { +pub(super) async fn detect_integrations(search_path: Option<&str>) -> IntegrationDetections { IntegrationDetections { cua: detect_cua_driver().await, - codex: codex::detect(codex_search_path).await, + codex: codex::detect(search_path).await, + claude: claude::detect(search_path).await, } } @@ -982,20 +1033,18 @@ fn validate_stored_registry( } let mut ids = HashSet::new(); for entry in ®istry.integrations { - if !matches!( - entry.id.as_str(), - CUA_DRIVER_INTEGRATION_ID | CODEX_INTEGRATION_ID - ) || !ids.insert(entry.id.as_str()) + if (entry.id != CUA_DRIVER_INTEGRATION_ID && external_agent_selection(&entry.id).is_none()) + || !ids.insert(entry.id.as_str()) { return Err( "Device-local integration settings contain an unknown or duplicate integration" .to_string(), ); } - if entry.id == CODEX_INTEGRATION_ID { + if external_agent_selection(&entry.id).is_some() { if entry.backend != AgentIntegrationBackend::Embedded || entry.external_server.is_some() { - return Err("Device-local Codex settings are invalid".to_string()); + return Err("Device-local external agent settings are invalid".to_string()); } continue; } @@ -1272,6 +1321,7 @@ mod tests { fn cua_only(cua: CuaDetection) -> IntegrationDetections { IntegrationDetections { + claude: ClaudeDetection::default(), cua, codex: CodexDetection::default(), } @@ -1307,10 +1357,14 @@ mod tests { } fn codex_registry(enabled: bool) -> StoredIntegrationRegistry { + provider_registry("codex", enabled) + } + + fn provider_registry(provider: &str, enabled: bool) -> StoredIntegrationRegistry { StoredIntegrationRegistry { version: INTEGRATIONS_FILE_VERSION, integrations: vec![StoredIntegration { - id: CODEX_INTEGRATION_ID.to_string(), + id: provider.to_string(), enabled, backend: AgentIntegrationBackend::Embedded, external_server: None, @@ -2005,6 +2059,7 @@ mod tests { ); let user = "codex-user"; let detections = IntegrationDetections { + claude: ClaudeDetection::default(), cua: CuaDetection::not_detected(), codex: CodexDetection { executable: Some(temporary.path().join("codex")), @@ -2080,6 +2135,7 @@ mod tests { enabled: false, }, &IntegrationDetections { + claude: ClaudeDetection::default(), cua: CuaDetection::not_detected(), codex: CodexDetection::default(), }, @@ -2114,6 +2170,7 @@ mod tests { enabled: true, }; let missing = IntegrationDetections { + claude: ClaudeDetection::default(), cua: CuaDetection::not_detected(), codex: CodexDetection::default(), }; @@ -2126,6 +2183,7 @@ mod tests { ); let old = IntegrationDetections { + claude: ClaudeDetection::default(), cua: CuaDetection::not_detected(), codex: CodexDetection { executable: Some(temporary.path().join("codex")), @@ -2145,4 +2203,82 @@ mod tests { ); assert!(!external_agents_enabled(&paths, user)); } + #[test] + fn claude_card_reports_authentication_without_gating_availability() { + for (signed_in, detail) in [ + (Some(true), "Signed in."), + (Some(false), claude::sign_in_hint()), + ( + None, + "Could not check sign-in. Run `claude auth status` in a terminal.", + ), + ] { + let card = claude_public( + &ClaudeDetection { + executable: Some(PathBuf::from("claude")), + signed_in, + ..Default::default() + }, + None, + ); + assert_eq!(card.detail.as_deref(), Some(detail)); + assert_eq!(card.availability, AgentIntegrationAvailability::Available); + } + } + + #[test] + fn claude_setup_and_defaults_keep_shared_skills_until_both_are_off() { + let temp = tempfile::tempdir().unwrap(); + let paths = + AgentPathLayout::from_app_roots(temp.path().join("config"), temp.path().join("data")); + let mut detections = cua_only(CuaDetection::not_detected()); + let request = AgentSetIntegrationEnabledRequest { + id: "claude".into(), + enabled: true, + }; + assert!( + set_integration_default(&paths, "user", &request, &detections) + .unwrap_err() + .contains("Install Claude Code") + ); + detections.claude = ClaudeDetection { + executable: Some(temp.path().join("claude")), + version: Some("2.1.270".into()), + signed_in: None, + problem: Some("Claude Code could not be started".into()), + }; + assert!( + set_integration_default(&paths, "user", &request, &detections) + .unwrap_err() + .contains("could not be started") + ); + detections.claude.problem = None; + save_stored_integrations(&paths, "user", &codex_registry(true)).unwrap(); + set_integration_default(&paths, "user", &request, &detections).unwrap(); + set_integration_default( + &paths, + "user", + &AgentSetIntegrationEnabledRequest { + id: "codex".into(), + enabled: false, + }, + &detections, + ) + .unwrap(); + assert!(external_agents_enabled(&paths, "user")); + let skills = external_agent_skills_dir(&paths, "user").unwrap(); + assert!(skills.join("handoff/SKILL.md").is_file()); + set_integration_default( + &paths, + "user", + &AgentSetIntegrationEnabledRequest { + id: "claude".into(), + enabled: false, + }, + &detections, + ) + .unwrap(); + assert!(!external_agents_enabled(&paths, "user")); + assert!(!skills.join("handoff/SKILL.md").exists()); + } } diff --git a/apps/maple-agent/crates/maple-agent/src/agent/types.rs b/apps/maple-agent/crates/maple-agent/src/agent/types.rs index 63a33b085..ea5a077f0 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent/types.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent/types.rs @@ -316,6 +316,7 @@ pub struct AgentQuestionOption { #[derive(Debug, Clone, PartialEq, Eq, Serialize)] #[serde(rename_all = "camelCase")] pub struct AgentQuestion { + pub multi_select: bool, pub id: String, pub header: String, pub question: String, From 6292b25cc23bdac7fce155594d531ee4519a34a5 Mon Sep 17 00:00:00 2001 From: benthecarman Date: Sat, 19 Sep 2026 01:28:59 -0500 Subject: [PATCH 2/3] Gate session agents on integration settings Show Codex and Claude in the composer only when enabled in Settings. Require each task to opt in, and block saved selections while Settings is disabled. Preserve task choices across restarts and re-enablement. Exercise admission, composer visibility, and cached tool refresh for both providers, including rejection of direct requests when disabled. --- apps/maple-agent/app/src/ui/settings.rs | 2 +- .../crates/maple-agent/src/agent.rs | 11 +- .../maple-agent/src/agent/integrations.rs | 142 +++++++++++++----- .../crates/maple-agent/src/agent/types.rs | 2 + 4 files changed, 117 insertions(+), 40 deletions(-) diff --git a/apps/maple-agent/app/src/ui/settings.rs b/apps/maple-agent/app/src/ui/settings.rs index 5b59f7e86..16562f048 100644 --- a/apps/maple-agent/app/src/ui/settings.rs +++ b/apps/maple-agent/app/src/ui/settings.rs @@ -2236,7 +2236,7 @@ impl SettingsScreen { div() .text_sm() .text_color(gpui::rgb(theme::text_muted())) - .child("Set integration defaults here. Choose integrations for each task in the composer."), + .child("Manage integrations here. Enable coding agents to make them available in the composer, then select them for each task."), ); if let Some(notice) = &self.integration_notice { diff --git a/apps/maple-agent/crates/maple-agent/src/agent.rs b/apps/maple-agent/crates/maple-agent/src/agent.rs index d1b2a4e0b..5983a3a06 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent.rs @@ -4209,6 +4209,12 @@ impl AgentRuntimeHandle { return Err("External agents are available only in desktop tasks".to_string()); } if request.enabled { + if !external_agent_enabled(&stored_integrations, provider.id) { + return Err( + "Enable this integration in Settings before selecting it for this task" + .to_string(), + ); + } let integrations = project_integrations(&state.host.paths, user_id, &detected)?; let integration = integrations .iter() @@ -8814,8 +8820,7 @@ async fn finish_session_agent( session, allow_embedded_cua, ); - // A task override can keep delegation enabled after the device default - // was disabled (and removed Maple's skills), including after a restart. + // Restore bundled skills for admitted providers, including after a restart. if !selected_external_providers.is_empty() && let Err(error) = sync_external_agent_skills(skills_scope.paths, skills_scope.user_id, true) @@ -8865,7 +8870,7 @@ async fn finish_session_agent( .with_attachment_store(attachment_store) .with_web_enabled(session_web_enabled(session)) .with_desktop_ui_tools(session.session_type != SessionType::Acp) - // Re-read inherited defaults and explicit task choices at every run, + // Re-read Settings gates and explicit task choices at every run, // including cold restores. The driving surface remains the authority: // leasing a desktop task to ACP must remove desktop-only capabilities. .with_external_agents(external_agents.cloned(), selected_external_providers); diff --git a/apps/maple-agent/crates/maple-agent/src/agent/integrations.rs b/apps/maple-agent/crates/maple-agent/src/agent/integrations.rs index ac31fc94d..a35ebc5ef 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent/integrations.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent/integrations.rs @@ -83,21 +83,24 @@ impl ExtensionState for TaskIntegrationOverrides { const VERSION: &'static str = "1"; } +pub(super) fn external_agent_enabled(stored: &StoredIntegrationRegistry, id: &str) -> bool { + stored + .integrations + .iter() + .any(|entry| entry.id == id && entry.enabled) +} + fn session_external_agent_enabled( stored: &StoredIntegrationRegistry, session: &Session, id: &str, ) -> bool { - // Missing metadata is inheritance, including tasks created before this - // integration existed. Never snapshot an inherited false into old tasks. - TaskIntegrationOverrides::from_extension_data(&session.extension_data) - .and_then(|state| state.enabled.get(id).copied()) - .unwrap_or_else(|| { - stored - .integrations - .iter() - .any(|entry| entry.id == id && entry.enabled) - }) + // Settings admits the provider; each task must also explicitly select it. + // A saved selection cannot bypass a disabled account/device integration. + external_agent_enabled(stored, id) + && TaskIntegrationOverrides::from_extension_data(&session.extension_data) + .and_then(|state| state.enabled.get(id).copied()) + .unwrap_or(false) } pub(super) fn session_external_agent_providers( @@ -856,6 +859,7 @@ pub(super) fn project_session_mcp_servers( servers.extend( EXTERNAL_AGENT_INTEGRATIONS .iter() + .filter(|entry| external_agent_enabled(stored, entry.id)) .map(|entry| AgentSessionMcpServer { name: entry.id.to_string(), kind: AgentSessionIntegrationKind::ExternalAgent, @@ -1373,7 +1377,7 @@ mod tests { } #[tokio::test] - async fn old_task_inherits_providers_until_explicitly_overridden() { + async fn tasks_require_both_settings_and_persisted_provider_selection() { let temp = tempfile::tempdir().unwrap(); let manager = SessionManager::new(temp.path().join("sessions")); let session = manager @@ -1388,10 +1392,7 @@ mod tests { assert!( session_external_agent_providers(&codex_registry(false), &session, true).is_empty() ); - assert_eq!( - session_external_agent_providers(&codex_registry(true), &session, true), - ["codex"] - ); + assert!(session_external_agent_providers(&codex_registry(true), &session, true).is_empty()); assert!( session_external_agent_providers(&codex_registry(true), &session, false).is_empty() ); @@ -1409,21 +1410,37 @@ mod tests { let enabled = persist_task_integration_override(&manager, &session.id, "codex", true) .await .unwrap(); + assert!( + session_external_agent_providers(&codex_registry(false), &enabled, true).is_empty() + ); assert_eq!( - session_external_agent_providers(&codex_registry(false), &enabled, true), + session_external_agent_providers(&codex_registry(true), &enabled, true), ["codex"] ); + let other_session = manager + .create_session( + temp.path().to_path_buf(), + "Other task".into(), + SessionType::User, + GooseMode::SmartApprove, + ) + .await + .unwrap(); + assert!( + session_external_agent_providers(&codex_registry(true), &other_session, true) + .is_empty() + ); drop(manager); let manager = SessionManager::new(temp.path().join("sessions")); let restored = manager.get_session(&session.id, false).await.unwrap(); assert_eq!( - session_external_agent_providers(&codex_registry(false), &restored, true), + session_external_agent_providers(&codex_registry(true), &restored, true), ["codex"] ); } #[test] - fn composer_lists_native_integrations_without_prior_setup() { + fn composer_lists_only_settings_enabled_providers_without_selecting_them() { let session = Session { session_type: SessionType::User, ..Session::default() @@ -1431,10 +1448,23 @@ mod tests { let rows = project_session_mcp_servers(&StoredIntegrationRegistry::default(), &[], &session) .unwrap(); - assert!(rows.iter().any(|row| row.name == "codex" - && row.kind == AgentSessionIntegrationKind::ExternalAgent - && row.display_name == "Codex" - && !row.enabled)); + assert!( + !rows + .iter() + .any(|row| row.kind == AgentSessionIntegrationKind::ExternalAgent) + ); + for provider in ["codex", "claude"] { + let rows = + project_session_mcp_servers(&provider_registry(provider, true), &[], &session) + .unwrap(); + let native: Vec<_> = rows + .iter() + .filter(|row| row.kind == AgentSessionIntegrationKind::ExternalAgent) + .collect(); + assert_eq!(native.len(), 1); + assert_eq!(native[0].name, provider); + assert!(!native[0].enabled); + } assert!( rows.iter() .any(|row| row.name == CUA_DRIVER_MCP_NAME && !row.enabled && !row.available) @@ -1456,7 +1486,7 @@ mod tests { assert!( codex_rows .iter() - .any(|row| row.kind == AgentSessionIntegrationKind::ExternalAgent && row.enabled) + .any(|row| row.kind == AgentSessionIntegrationKind::ExternalAgent && !row.enabled) ); let legacy_request: AgentSetSessionMcpServerRequest = serde_json::from_value(json!({ "sessionId": "s1", "name": "codex", "enabled": true, @@ -1472,6 +1502,11 @@ mod tests { /// external process is needed to inspect the model-facing tool catalog. #[tokio::test] async fn resumed_task_refreshes_external_tools_on_every_run() { + assert_resumed_task_refreshes_provider("codex").await; + assert_resumed_task_refreshes_provider("claude").await; + } + + async fn assert_resumed_task_refreshes_provider(provider: &str) { let fixture = super::super::test_support::started_agent_runtime("integration-refresh").await; let service = &fixture.handle.service; @@ -1513,20 +1548,24 @@ mod tests { ); let mut agent = Arc::new(Agent::with_config(config.clone())); let context = SharedAgentToolContext::new(AgentToolContextSpec::default()); - // Old task before enable, cached task after enable, cold restore after - // enable, explicit off, explicit on despite default off, and ACP lease. + // Settings alone never selects a provider. Revocation removes tools + // from warm and cold agents despite a saved selection; re-enabling + // Settings preserves the task's choice. ACP never receives the tools. for (default, override_value, cold, desktop, expected) in [ (false, None, false, true, false), - (true, None, true, true, true), + (true, None, true, true, false), (false, None, false, true, false), - (true, None, false, true, true), + (true, None, false, true, false), + (true, Some(true), false, true, true), + (false, Some(true), false, true, false), + (true, Some(true), true, true, true), + (false, Some(true), true, true, false), (true, Some(false), true, true, false), - (false, Some(true), true, true, true), - (true, None, false, false, false), + (true, Some(true), false, false, false), ] { - save_stored_integrations(paths, user, &codex_registry(default)).unwrap(); + save_stored_integrations(paths, user, &provider_registry(provider, default)).unwrap(); if let Some(enabled) = override_value { - persist_task_integration_override(&manager, &session.id, "codex", enabled) + persist_task_integration_override(&manager, &session.id, provider, enabled) .await .unwrap(); } @@ -1582,23 +1621,54 @@ mod tests { ); } } + save_stored_integrations(paths, user, &provider_registry(provider, false)).unwrap(); + let error = fixture + .handle + .set_session_mcp_server_enabled(AgentSetSessionMcpServerRequest { + session_id: session.id.clone(), + name: provider.to_string(), + kind: AgentSessionIntegrationKind::ExternalAgent, + enabled: true, + }) + .await + .unwrap_err(); + assert!( + error.contains("Enable this integration in Settings"), + "{error}" + ); + let hidden = fixture + .handle + .list_session_mcp_servers(session.id.clone()) + .await + .unwrap(); + assert!( + !hidden + .iter() + .any(|row| row.kind == AgentSessionIntegrationKind::ExternalAgent) + ); let rows = fixture .handle .set_session_mcp_server_enabled(AgentSetSessionMcpServerRequest { session_id: session.id.clone(), - name: "codex".to_string(), + name: provider.to_string(), kind: AgentSessionIntegrationKind::ExternalAgent, enabled: false, }) .await .unwrap(); assert!( - rows.iter() - .any(|row| row.kind == AgentSessionIntegrationKind::ExternalAgent && !row.enabled) + !rows + .iter() + .any(|row| row.kind == AgentSessionIntegrationKind::ExternalAgent) ); let stored_task = manager.get_session(&session.id, false).await.unwrap(); assert!( - session_external_agent_providers(&codex_registry(true), &stored_task, true).is_empty() + session_external_agent_providers( + &provider_registry(provider, true), + &stored_task, + true + ) + .is_empty() ); // A future provider selected alongside an installed but unselected // Codex must not grant Codex access through forged tool arguments. @@ -1622,7 +1692,7 @@ mod tests { .call_tool( &goose::agents::ToolCallContext::new(session.id.clone(), None, None), name, - Some(rmcp::object!({"provider": "codex"})), + Some(rmcp::object!({"provider": provider})), CancellationToken::new(), ) .await diff --git a/apps/maple-agent/crates/maple-agent/src/agent/types.rs b/apps/maple-agent/crates/maple-agent/src/agent/types.rs index ea5a077f0..5c5a31684 100644 --- a/apps/maple-agent/crates/maple-agent/src/agent/types.rs +++ b/apps/maple-agent/crates/maple-agent/src/agent/types.rs @@ -144,6 +144,8 @@ pub struct AgentIntegration { /// the desktop session, so the interface does not offer a button that /// repeats work the user already did. pub setup_available: bool, + /// Settings switch: external agents require this gate plus a per-task + /// selection. For CUA this remains the default for newly created tasks. pub enabled_for_new_tasks: bool, pub detail: Option, } From 7a9b23c7fb4fec84a7eaadc901e86e7c89dd4747 Mon Sep 17 00:00:00 2001 From: benthecarman Date: Sat, 19 Sep 2026 01:28:59 -0500 Subject: [PATCH 3/3] Document Claude setup and session selection Explain Claude CLI setup, native transport provenance, and the separate Settings and task controls. Record the fixture validation workflow so future provider changes preserve the shared lifecycle contract. --- .agents/skills/develop-maple-agent/SKILL.md | 13 ++- apps/maple-agent/README.md | 29 ++++-- apps/maple-agent/docs/external-agents.md | 103 +++++++++++++++++--- 3 files changed, 118 insertions(+), 27 deletions(-) diff --git a/.agents/skills/develop-maple-agent/SKILL.md b/.agents/skills/develop-maple-agent/SKILL.md index 6ba6ed51e..f18cb53d2 100644 --- a/.agents/skills/develop-maple-agent/SKILL.md +++ b/.agents/skills/develop-maple-agent/SKILL.md @@ -75,10 +75,17 @@ cleanup. Do not run raw `cargo clean` against the shared cache. For composer integrations and external providers, read `apps/maple-agent/docs/external-agents.md`. Keep provider metadata in the runtime catalog and pass typed selection kinds through the UI bridge so -user-controlled MCP names cannot shadow provider IDs. Preserve inherited -defaults for tasks without overrides and CUA's existing backend metadata. +user-controlled MCP names cannot shadow provider IDs. External providers must +be enabled in Settings before appearing in the composer, and each task must +explicitly select them. A saved task choice cannot bypass Settings. Preserve +CUA's existing backend metadata. Exercise warm and cold session tool catalogs and ACP exclusion when changing -run-boundary admission. +run-boundary admission. Claude Code uses a native Rust transport adapted from +the pinned Goose provider in `external_agents/claude.rs`. Preserve the source +attribution when changing that adapted code. Keep process ownership in Maple's +contained host. The `claude_native_*` tests re-execute the +Rust test binary as a CLI fixture and need no Claude account or inference +request. Keep fixture launch and environment setup shared with the Codex tests. ## Security and publication diff --git a/apps/maple-agent/README.md b/apps/maple-agent/README.md index 4999ad193..8121f334b 100644 --- a/apps/maple-agent/README.md +++ b/apps/maple-agent/README.md @@ -74,10 +74,10 @@ Cargo manifests and lockfile; Research has an independent dependency graph. keeps its row after the turn ends, and Maple tells the task when it finishes, with a bounded result in the running turn or a new turn Maple starts automatically. The task can use `load` to retrieve any truncated output. -- External agents: a task can hand work to the Codex CLI installed on +- External agents: a task can hand work to Codex or Claude Code installed on this computer with the `agent_start`, `agent_send`, `agent_status`, - `agent_cancel`, and `list_agent_providers` tools, once Codex is enabled - under Settings > Integrations. Codex runs in the project with its own + `agent_cancel`, and `list_agent_providers` tools, once the provider is enabled + under Settings > Integrations. Each agent runs in the project with its own account, context, and sandbox settings; whatever it asks approval for comes to you through Maple's permission card, and Allow all grants it. Its progress streams into the tool call's row and its row above the composer has a Stop @@ -158,16 +158,27 @@ account configuration that may roam between devices. The embedded design, migration rules, privacy boundary, and preview limits are documented in [`docs/embedded-cua.md`](docs/embedded-cua.md). +#### Claude Code + +Settings > Integrations lists Claude Code (`claude`) alongside Codex, with the +same per-task selection, streamed activity, permission cards, and Stop control. +Install the Claude Code CLI on the app's PATH and sign in using +`claude auth login`. Maple uses a Rust transport adapted from Goose's Claude +Code provider. The CLI is the only external runtime dependency. The integration +is off by default. Enable it in Settings to show it in the composer, then +select it for the tasks that should use it. See +[external agents](docs/external-agents.md#how-claude-code-is-driven). + #### Codex Settings > Integrations also lists the Codex CLI when `codex` is on the PATH (the login shell's PATH on macOS). The card shows the installed version and whether Codex is signed in; Maple never runs Codex's sign-in itself. The -toggle is off by default. The composer lists Codex alongside CUA and custom -MCP servers, with an independent choice for each task. Tasks without an -explicit Codex choice inherit the Settings default on every run, including -older tasks; composer overrides survive relaunches. Enabling it gives runs -the external-agent tools and installs the `handoff`, `committee`, and `advisor` skills into the +toggle is off by default. Enabling it makes Codex available in the composer +alongside CUA and custom MCP servers. Each task must select Codex explicitly; +that choice survives relaunches but only applies while Settings enables Codex. +Selecting it gives that task the external-agent tools. Enabling it in Settings +installs the `handoff`, `committee`, and `advisor` skills into the account's Goose skills directory; disabling removes only the files Maple wrote. Codex needs version 0.143 or newer. See [`docs/external-agents.md`](docs/external-agents.md). @@ -433,7 +444,7 @@ The roots follow the platform, the same way the Tauri app's | `/settings.json` | App settings. | | `/agent/accounts//config.json` | Per-account agent configuration (default root, model, custom MCP servers, project trust). May roam between machines. | | `/agent/accounts//goose/config/` | Goose permission file for the account. | -| `/agent/accounts//goose/config/skills/` | Skills the account's tasks can load, including the delegation skills Maple installs while Codex is enabled. | +| `/agent/accounts//goose/config/skills/` | Skills the account's tasks can load, including the delegation skills Maple installs while any external agent is enabled. | | `/agent/goose-runtime/` | Goose process configuration. | | `/auth.json` | Sign-in credentials (mode 0600). Device-local; never in a roaming profile. | | `/agent/accounts//integrations.json` | Per-account defaults and validated launch details for integrations detected on this device. | diff --git a/apps/maple-agent/docs/external-agents.md b/apps/maple-agent/docs/external-agents.md index f056b7e62..345414123 100644 --- a/apps/maple-agent/docs/external-agents.md +++ b/apps/maple-agent/docs/external-agents.md @@ -1,36 +1,39 @@ # External agents A Maple task can hand work to an external coding agent that is installed on -the same computer. Codex is the first provider. The tool contract is generic -so another harness can be added without changing what the task sees. +the same computer. Supported providers are Codex (`codex`) and Claude Code +(`claude`). Both use the same delegation tools and task controls. ## What the task sees Five tools appear when at least one external agent is selected in the -composer's Integrations menu. Settings > Integrations supplies the device -and account default. Tasks without an explicit choice inherit that default -on every run, including tasks created before the integration existed. -Choosing on or off in the composer persists an override for that task. +composer's Integrations menu. Settings > Integrations controls which providers +are available to select for this account on this device. Disabled providers +are hidden from the composer. Enabling a provider in Settings does not select +it for any task: each task starts with external agents off. +Choosing on or off in the composer persists that choice for that task. +Disabling a provider in Settings blocks saved task selections on the next run; +re-enabling it restores their availability without erasing those choices. Changes take effect on its next run; stop an active run before changing its selection. CUA and custom MCP servers have independent switches. -An unavailable integration remains visible and links to Settings for setup. +An integration enabled in Settings but unavailable on the device remains +visible in the composer and links to Settings for setup. The runtime checks installation before enabling a provider and again when launching it. External agents are desktop capabilities: an ACP caller cannot acquire them by resuming a desktop task. The tools are: - | Tool | Purpose | | --- | --- | -| `list_agent_providers` | Which providers are installed, their version, and whether they are signed in. | +| `list_agent_providers` | Selected providers, their installation status and version, and setup or sign-in guidance. | | `agent_start` | Start an agent on a self-contained briefing. Blocking by default; `background: true` returns at once. | | `agent_send` | Give a started agent more instructions in the same thread. | | `agent_status` | Read an agent's status, last message, changed files, and commands. | | `agent_cancel` | Stop an agent's current turn and its process. The thread stays on disk; the next `agent_send` resumes it in a fresh process. | -All tools except `list_agent_providers` take `provider` (`"codex"`). `agent_start` also takes optional +All tools except `list_agent_providers` take `provider` (`"codex"` or `"claude"`). `agent_start` also takes optional `model`, `effort`, and `cwd`. `cwd` must be inside the project root. A result has Paseo's shape: a status line, the agent ID and thread ID, the @@ -78,7 +81,7 @@ default mode, opens Maple's question card and the answer goes back in Codex's own shape; an asynchronous answer that arrives after the turn ended starts a follow-up turn on the same thread. -The child runs in the project root, on the user's login PATH, with the same +The child runs in the requested project directory, on the user's login PATH, with the same environment scrubbing as the shell tool, in its own process group or job so teardown reaches every descendant. It is killed when the runtime stops, on logout, and when its task is deleted. Threads are not ephemeral, so @@ -93,6 +96,67 @@ PATH or the user's shell configuration. The handshake reports the reserved client name `codex_app_server_daemon`, the same non-originating name Paseo uses. +## How Claude Code is driven + +Maple's native Rust transport is adapted from +[Goose's `ClaudeCodeProvider`](https://github.com/AnthonyRonning/goose/blob/785d655d110746147117d23690e09cc7023aa9dc/crates/goose/src/providers/claude_code.rs). +The control protocol types and permission exchange come from that implementation +of Claude's SDK protocol. The adapted source lives in `external_agents/claude.rs` +with its provenance. + +Goose keeps its transport private inside its provider. Maple adapts that code +so its existing host retains control of process launch, the working directory, +scoped environment, cancellation, and descendant cleanup. It also bounds +protocol lines, sanitizes errors, projects activity, and supplies question +answers through `updatedInput`. + +Install the Claude Code CLI and sign in with `claude auth login`. Detection +runs `claude --version` and `claude auth status --json`. The card shows the +CLI's sign-in status or login instructions; an unavailable, malformed, or +timed-out status check is reported as unknown. Maple reads only the sign-in +boolean, without retaining account details or reading credential files. +As with Codex, sign-in status does not gate enabling the integration. +The CLI is the only additional runtime dependency. + +Claude keeps its normal system prompt and configuration. Maple passes +`--permission-mode default` and `--permission-prompt-tool stdio`, never a +bypass-permissions flag. Claude's rules decide which actions need approval; +`can_use_tool` requests go to Maple's current permission mode. Allow all answers +those requests automatically, and Read only shows a one-shot permission card. +`AskUserQuestion` uses Maple's question card, including multiple choices when +Claude sets `multiSelect`. Click options or use numbered shortcuts to toggle +them, then choose Answer (or Next in a batch). Vim navigation moves the cursor; +Enter toggles the current choice. Answers change only that call's input; Maple +never persists an allow rule in Claude's settings. + +Each turn starts a contained Claude CLI process using `--session-id` initially +and `--resume` after the CLI confirms the session with `system/init` or a +successful result. A failure before confirmation leaves the agent retryable +with a fresh session. This allows optional `model` and `effort` launch +arguments to change on each turn. The process runs in the requested project +directory with the shell tool's scoped environment and process containment. +Stop, task deletion, logout, and runtime shutdown reclaim Claude and its +descendants. A subsequent `agent_send` resumes the saved session. You can also +use `claude --resume` from a terminal. + +Claude's text, Bash calls, successful Edit/Write/NotebookEdit calls, and +TodoWrite items feed the existing activity row. Other tools continue to run +under Claude's policy but do not yet have specialized activity summaries. +Command completion shows success or failure; tool-result messages do not +guarantee numeric exit codes. Protocol failures produce a generic error without +forwarding potentially sensitive exception text or CLI stderr. + +The `claude_native_*` runtime tests use a Rust CLI fixture, including a Rust +sleeping child for descendant cleanup; they need no installed `sleep` program +or Claude account. On Linux, the ignored `native_question_card_fixture` app +test mounts the real chat screen with multi-select and single-select questions +for interactive checks on a private desktop. Build it with `cargo test -p +maple-gpui native_question_card_fixture --no-run` through the component build +environment, then launch the reported test binary on the private display with +`native_question_card_fixture --ignored --nocapture --test-threads=1` and +isolated XDG directories. It supplies fixture questions without authenticating +or starting inference; runtime broker tests cover delivery of the answer. + ## Transcript The tool call's row streams the agent's text, its commands with exit @@ -123,7 +187,8 @@ approval therefore still needs `load` before it can finish; completion delivery alone cannot unblock it. Stopping an agent, from its row or with `agent_cancel`, sends -`turn/interrupt`, waits briefly for Codex to confirm, then kills the +`turn/interrupt` (translated to Claude’s native `interrupt` control request), waits +briefly for confirmation, then kills the process group. Codex does not always end a sandboxed command on interrupt, so the kill is what guarantees nothing keeps running. @@ -144,17 +209,18 @@ so the kill is what guarantees nothing keeps running. Register its stable ID, label, description, and Settings projection in `EXTERNAL_AGENT_INTEGRATIONS` in `agent/integrations.rs`, then add discovery to `IntegrationDetections` and a transport adapter under `agent/external_agents/`. -Settings defaults, composer rows, and task overrides use this catalog. +Settings gates, composer rows, and task choices use this catalog. No provider-specific composer branch or database migration is needed. Task overrides live in the versioned `maple_integrations` extension data, -keyed by provider ID. Missing entries mean inheritance, not disabled. CUA +keyed by provider ID. Missing entries mean disabled. A true entry grants access +only while Settings also enables that provider. CUA keeps its existing `maple_cua` backend metadata so an old external driver task cannot silently switch to the embedded backend. The selector carries a typed `kind` alongside `name` and `displayName`. MCP names and external provider IDs are separate domains; a custom MCP -server named `codex` cannot toggle the Codex provider. Older MCP requests +server named `codex` or `claude` cannot toggle either provider. Older MCP requests without `kind` continue to mean MCP. Provider transport adapters still own their protocol, discovery, progress, @@ -167,3 +233,10 @@ Keep listing filtered by that same selection when more adapters are added. run configuration and Goose tool cache with old persisted extension state, warm reuse, cold restore, explicit overrides, and a desktop task leased to ACP. Extend that regression with each adapter's admission behavior. + +The `claude_native_*` tests exercise the Rust transport against a deterministic +fake CLI, without network requests. Both providers' fixtures re-execute the +Rust test binary through CLI shims on a private search PATH. They cover +streamed activity, permissions, questions, resumption, +provider/session isolation, errors, and process-group cancellation. Live +inference and platform packaging remain separate checks.