diff --git a/.claude/worktrees/agent-a2d818144fe5555bc b/.claude/worktrees/agent-a2d818144fe5555bc new file mode 160000 index 0000000..c1ad206 --- /dev/null +++ b/.claude/worktrees/agent-a2d818144fe5555bc @@ -0,0 +1 @@ +Subproject commit c1ad206a0020a52218f4bee9d3b5773da2e5e880 diff --git a/.claude/worktrees/agent-a42650db8ee77314e b/.claude/worktrees/agent-a42650db8ee77314e new file mode 160000 index 0000000..7cb4ae6 --- /dev/null +++ b/.claude/worktrees/agent-a42650db8ee77314e @@ -0,0 +1 @@ +Subproject commit 7cb4ae670834895f0a638b2d5f5865525c3825e1 diff --git a/.claude/worktrees/agent-a712289b5e04047c2 b/.claude/worktrees/agent-a712289b5e04047c2 new file mode 160000 index 0000000..7cb4ae6 --- /dev/null +++ b/.claude/worktrees/agent-a712289b5e04047c2 @@ -0,0 +1 @@ +Subproject commit 7cb4ae670834895f0a638b2d5f5865525c3825e1 diff --git a/.claude/worktrees/agent-ada61e93513e2b0e4 b/.claude/worktrees/agent-ada61e93513e2b0e4 new file mode 160000 index 0000000..7cb4ae6 --- /dev/null +++ b/.claude/worktrees/agent-ada61e93513e2b0e4 @@ -0,0 +1 @@ +Subproject commit 7cb4ae670834895f0a638b2d5f5865525c3825e1 diff --git a/.gitignore b/.gitignore index 5f32e70..b8b4d97 100644 --- a/.gitignore +++ b/.gitignore @@ -1,2 +1,3 @@ target/ -.env \ No newline at end of file +.env +.claude/settings.local.json \ No newline at end of file diff --git a/src/app/harness.rs b/src/app/harness.rs index db85623..d7daeb5 100644 --- a/src/app/harness.rs +++ b/src/app/harness.rs @@ -18,7 +18,58 @@ impl Harness { if mode.auto_approve() { return Verdict::Allow; } - Verdict::Allow + if matches!(mode, super::state::types::AgentMode::Plan) { + return Verdict::Block("mutating tools are disabled in Plan mode".to_string()); + } + Verdict::Escalate + } + + pub fn gate_tool_call( + tool_name: &str, + args: &serde_json::Value, + mode: &super::state::types::AgentMode, + workspace_roots: &[&std::path::Path], + ) -> Verdict { + if let Err(e) = Self::run_catastrophic_guard(tool_name, args, workspace_roots) { + return Verdict::Block(e); + } + if !crate::tool::tool_is_risky(tool_name) { + return Verdict::Allow; + } + Self::classify(tool_name, mode) + } + + fn run_catastrophic_guard( + tool_name: &str, + args: &serde_json::Value, + workspace_roots: &[&std::path::Path], + ) -> Result<(), String> { + use super::catastrophic::CatastrophicGuard; + match tool_name { + "bash" => { + let cmd = args.get("command").and_then(|v| v.as_str()).unwrap_or(""); + CatastrophicGuard::check_all(cmd, workspace_roots) + } + "git_operator" => { + let operation = args.get("operation").and_then(|v| v.as_str()).unwrap_or(""); + let arg_list: Vec = args + .get("args") + .and_then(|v| v.as_array()) + .map(|arr| arr.iter().filter_map(|v| v.as_str().map(|s| s.to_string())).collect()) + .unwrap_or_default(); + let cmd = format!("git {} {}", operation, arg_list.join(" ")); + CatastrophicGuard::check_all(&cmd, workspace_roots) + } + "delete" => { + let path = args.get("path").and_then(|v| v.as_str()).unwrap_or(""); + CatastrophicGuard::check_delete_path(std::path::Path::new(path), workspace_roots) + } + "web_download" | "download" => { + let path = args.get("path").and_then(|v| v.as_str()).unwrap_or(""); + CatastrophicGuard::check_download_path(std::path::Path::new(path)) + } + _ => Ok(()), + } } } diff --git a/src/app/mcp/manager.rs b/src/app/mcp/manager.rs index ef82aa9..a036de0 100644 --- a/src/app/mcp/manager.rs +++ b/src/app/mcp/manager.rs @@ -42,6 +42,30 @@ pub struct McpManager { pub running: bool, } +pub struct McpToolAdapter { + name: &'static str, + description: &'static str, + parameters: serde_json::Value, +} + +impl crate::tool::Tool for McpToolAdapter { + fn name(&self) -> &'static str { + self.name + } + + fn description(&self) -> &'static str { + self.description + } + + fn parameters(&self) -> serde_json::Value { + self.parameters.clone() + } + + fn run(&self, _ctx: &crate::tool::ToolCtx, _args: &serde_json::Value) -> anyhow::Result { + Err(anyhow::anyhow!("MCP tool execution not yet implemented")) + } +} + impl McpManager { pub fn new() -> Self { McpManager { @@ -75,4 +99,15 @@ impl McpManager { self.running = false; Ok(()) } + + pub fn as_tools(&self) -> Vec> { + self.all_tools().into_iter().map(|info| { + let name = format!("mcp__{}", info.name); + Box::new(McpToolAdapter { + name: Box::leak(name.into_boxed_str()), + description: Box::leak(info.description.clone().into_boxed_str()), + parameters: info.input_schema.clone(), + }) as Box + }).collect() + } } diff --git a/src/app/mode/bash.rs b/src/app/mode/bash.rs index 07651c2..d7d118d 100644 --- a/src/app/mode/bash.rs +++ b/src/app/mode/bash.rs @@ -8,6 +8,7 @@ pub fn handle_bash_submit(state: &mut AppStateRest, command: String) { } } +#[expect(dead_code)] pub fn handle_bash_dismiss(state: &mut AppStateRest) { state.misc.overlay = Overlay::None; state.dirty = true; diff --git a/src/app/mode/mod.rs b/src/app/mode/mod.rs index 49f5889..fc99434 100644 --- a/src/app/mode/mod.rs +++ b/src/app/mode/mod.rs @@ -2,7 +2,6 @@ use serde::{Deserialize, Serialize}; #[expect(dead_code)] pub mod agents; -#[expect(dead_code)] pub mod bash; #[expect(dead_code)] pub mod editor; @@ -22,17 +21,13 @@ pub mod onboard; pub mod onboard_provider; #[expect(dead_code)] pub mod picker; -#[expect(dead_code)] pub mod quit_confirm; #[expect(dead_code)] pub mod rewind; -#[expect(dead_code)] pub mod security; #[expect(dead_code)] pub mod session_hub; -#[expect(dead_code)] pub mod settings; -#[expect(dead_code)] pub mod todo; #[expect(dead_code)] pub mod workflow; diff --git a/src/app/mode/security.rs b/src/app/mode/security.rs index bf878ae..82d8032 100644 --- a/src/app/mode/security.rs +++ b/src/app/mode/security.rs @@ -6,6 +6,7 @@ pub fn toggle_security_arm(state: &mut AppStateRest) { state.dirty = true; } +#[expect(dead_code)] pub fn acknowledge_security(state: &mut AppStateRest) { if !state.misc.security_acknowledged { state.misc.security_acknowledged = true; diff --git a/src/app/mode/session_hub.rs b/src/app/mode/session_hub.rs index c9a25fd..2293b59 100644 --- a/src/app/mode/session_hub.rs +++ b/src/app/mode/session_hub.rs @@ -1,20 +1,46 @@ -use crate::model::session::Session; use crate::app::state::rest::AppStateRest; use crate::app::state::types::Overlay; pub fn load_sessions(state: &mut AppStateRest) { - state.sessions = Session::list(&state.session_dir); + let base = state.store_base_dir(); + state.sessions = crate::model::session::Session::list(&base); state.dirty = true; } pub fn select_session(state: &mut AppStateRest, session_id: &str) { - if let Some(session) = state.sessions.iter().find(|s| s.id == session_id) { - let display = crate::app::state::rest::ChatMessageDisplay::new( - crate::dto::chat::message::Role::System, - format!("switched to session: {}", session.title), - ); - state.push_transcript(display); - state.misc.overlay = Overlay::None; - state.dirty = true; + let base = state.store_base_dir(); + let session = crate::model::session::Session::load(session_id, &base).ok(); + if session.is_none() { + return; } + let session = session.unwrap(); + let conv_path = session.conversation_path(&base); + let loaded_msgs: Vec = + std::fs::read_to_string(&conv_path) + .ok() + .and_then(|data| serde_json::from_str(&data).ok()) + .unwrap_or_default(); + + state.session_id = session.id.clone(); + state.session_dir = session.session_dir(&base); + state.session_runtime = Some(crate::app::state::runtime::SessionRuntime::new( + state.session_dir.clone(), + )); + state.transcript_cache.messages.clear(); + if let Some(ref mut rt) = state.session_runtime { + for msg in loaded_msgs { + let display = crate::app::state::rest::ChatMessageDisplay::new( + msg.role.clone(), + msg.content.clone().unwrap_or_default(), + ); + state.transcript_cache.messages.push(display); + rt.push_message(msg); + } + } + state.push_transcript(crate::app::state::rest::ChatMessageDisplay::new( + crate::dto::chat::message::Role::System, + format!("switched to session: {}", session.title), + )); + state.misc.overlay = Overlay::None; + state.dirty = true; } diff --git a/src/app/mode/settings.rs b/src/app/mode/settings.rs index 6684be6..90ac124 100644 --- a/src/app/mode/settings.rs +++ b/src/app/mode/settings.rs @@ -2,6 +2,7 @@ use crate::app::runtime::actions::Action; use crate::app::state::rest::AppStateRest; use crate::model::settings::{Settings, InternetMode}; +#[expect(dead_code)] pub fn apply_settings_action(state: &mut AppStateRest, action: &Action) { if let Action::ToggleYoloArm = action { state.misc.yolo_armed = !state.misc.yolo_armed; @@ -17,6 +18,7 @@ pub fn cycle_internet_mode(settings: &mut Settings) { }; } +#[expect(dead_code)] pub fn cycle_review_enabled(settings: &mut Settings) { settings.review_enabled = !settings.review_enabled; } diff --git a/src/app/review/mod.rs b/src/app/review/mod.rs index 0e560cf..f835a41 100644 --- a/src/app/review/mod.rs +++ b/src/app/review/mod.rs @@ -1,6 +1,10 @@ use std::collections::HashMap; use crate::app::state::rest::AppStateRest; +use crate::app::state::runtime::TurnEvent; use crate::app::state::types::{AgentMode, Origin, Toast, ToastKind}; +use crate::app::subagent::context::build_subagent_context; +use crate::app::subagent::engine::run_subagent; +use crate::app::subagent::spawn::AgentDefinition; use serde::{Deserialize, Serialize}; #[derive(Debug, Clone, Serialize, Deserialize, PartialEq)] @@ -163,10 +167,45 @@ pub fn should_trigger_review(state: &AppStateRest, origin: Origin) -> bool { } pub fn trigger_review(state: &mut AppStateRest) -> anyhow::Result<()> { - let _origin = Origin::Reviewer; - state.push_transcript(crate::app::state::rest::ChatMessageDisplay::new( - crate::dto::chat::message::Role::System, - "review triggered".to_string(), + let def = AgentDefinition::new( + "quality-reviewer".to_string(), + "reviewer".to_string(), + ); + let mut ctx = build_subagent_context(def); + ctx.session_dir = state.session_dir.clone(); + ctx.system_prompt = format!( + "You are a code quality reviewer. Review the recent code changes \ + for correctness, security, and adherence to best practices. \ + Use read-only tools (read, grep, glob, recall, remember) to \ + inspect the session files and provide a concise review verdict. \ + Session directory: {:?}", + state.session_dir + ); + + let (tx, _rx) = tokio::sync::mpsc::channel(32); + let turn_events = state.turn_events.clone(); + + std::thread::spawn(move || { + let result = run_subagent(ctx, tx); + let message = match result { + Ok(verdict) => { + let first_line = verdict.lines().next().unwrap_or(&verdict); + format!("Quality review: {}", first_line) + } + Err(e) => format!("Quality review failed: {}", e), + }; + if let Ok(mut q) = turn_events.lock() { + q.push_back(TurnEvent::SystemNote { + kind: "review".to_string(), + message, + }); + } + }); + + state.push_toast(Toast::new( + ToastKind::Info, + "Quality review triggered".to_string(), )); + Ok(()) } diff --git a/src/app/runtime/actions/mod.rs b/src/app/runtime/actions/mod.rs index 1180676..ca4cab5 100644 --- a/src/app/runtime/actions/mod.rs +++ b/src/app/runtime/actions/mod.rs @@ -1,10 +1,17 @@ +use std::collections::VecDeque; + +use crate::app::harness::Verdict; +use sha2::Digest; use crate::app::mode::ModeKind; -use crate::app::state::rest::AppStateRest; -use crate::app::state::types::Overlay; -use crate::dto::chat::message::ChatMessage; +use crate::app::review::{should_trigger_review, trigger_review}; +use crate::app::state::rest::{AppStateRest, ChatMessageDisplay}; +use crate::app::state::runtime::TurnEvent; +use crate::app::state::types::{AgentMode, Origin, Overlay, Toast, ToastKind}; +use crate::dto::chat::message::{ChatMessage, Role}; + +const MAX_AGENT_STEPS: usize = 40; #[derive(Debug, Clone)] -#[expect(dead_code)] pub enum Action { Quit, ForceQuit, @@ -22,6 +29,7 @@ pub enum Action { OpenOverlay(Overlay), CloseOverlay, ToggleYoloArm, + CycleAgentMode, ToolResult { tool_call_id: String, output: String, @@ -44,22 +52,29 @@ pub enum Action { LessonImport { path: String, }, + #[expect(dead_code)] RecordUsage { tokens_in: u64, tokens_out: u64, duration_ms: u64, }, + #[expect(dead_code)] RecordReviewTokens { tokens: u64, }, + SaveSession, + ResumeSession, + RefreshSessions, } pub fn apply_action(state: &mut AppStateRest, action: Action) { match action { Action::Quit => { + save_current_session(state); state.quit = true; } Action::ForceQuit => { + save_current_session(state); state.quit = true; } Action::SwitchMode(mode) => { @@ -87,32 +102,17 @@ pub fn apply_action(state: &mut AppStateRest, action: Action) { state.dirty = true; } Action::SubmitInput(text) => { - if let Some(ref mut rt) = state.session_runtime { - rt.push_message(ChatMessage::user(text.clone())); - let api_key = state.settings.api_key.clone(); - let model = state.settings.model.clone(); - let msgs = rt.messages.clone(); - let pending = state.pending_api_response.clone(); - if let Some(key) = api_key { - if !key.is_empty() { - std::thread::spawn(move || { - let client = crate::service::openrouter::OpenRouterClient::new(key, model); - match client.chat(&msgs) { - Ok(response) => { - if let Ok(mut guard) = pending.lock() { - *guard = Some(response); - } - } - Err(e) => { - if let Ok(mut guard) = pending.lock() { - *guard = Some(format!("Error: {}", e)); - } - } - } - }); - } - } + state.input.submit(); + let text = text.trim().to_string(); + if text.is_empty() { + state.dirty = true; + return; } + state.push_transcript(ChatMessageDisplay::new(Role::User, text.clone())); + if let Some(ref mut rt) = state.session_runtime { + rt.push_message(ChatMessage::user(text)); + } + spawn_turn(state); state.dirty = true; } Action::InsertChar(c) => { @@ -164,6 +164,13 @@ pub fn apply_action(state: &mut AppStateRest, action: Action) { state.misc.yolo_armed = !state.misc.yolo_armed; state.dirty = true; } + Action::CycleAgentMode => { + let next = state.mode.cycle(); + state.set_mode(next); + let toast = Toast::new(ToastKind::Info, format!("Mode: {}", state.mode.name())); + state.push_toast(toast); + state.dirty = true; + } Action::ToolResult { tool_call_id, output, @@ -269,26 +276,117 @@ pub fn apply_action(state: &mut AppStateRest, action: Action) { } state.dirty = true; } + Action::SaveSession => { + save_current_session(state); + state.push_toast(Toast::new(ToastKind::Success, "session saved".to_string())); + state.dirty = true; + } + Action::ResumeSession => { + let base = state.store_base_dir(); + let sessions = crate::model::session::Session::list(&base); + let target = sessions.into_iter() + .filter(|s| s.id != state.session_id) + .max_by_key(|s| s.updated_at); + if let Some(session) = target { + let conv_path = session.conversation_path(&base); + let loaded_msgs: Vec = + std::fs::read_to_string(&conv_path) + .ok() + .and_then(|data| serde_json::from_str(&data).ok()) + .unwrap_or_default(); + state.session_id = session.id.clone(); + state.session_dir = session.session_dir(&base); + state.session_runtime = Some(crate::app::state::runtime::SessionRuntime::new( + state.session_dir.clone(), + )); + state.transcript_cache.messages.clear(); + if let Some(ref mut rt) = state.session_runtime { + for msg in loaded_msgs { + let display = ChatMessageDisplay::new( + msg.role.clone(), + msg.content.clone().unwrap_or_default(), + ); + state.transcript_cache.messages.push(display); + rt.push_message(msg); + } + } + state.push_toast(Toast::new(ToastKind::Success, + format!("resumed session: {}", session.title))); + } else { + state.push_toast(Toast::new(ToastKind::Info, + "no other sessions to resume".to_string())); + } + state.dirty = true; + } + Action::RefreshSessions => { + let base = state.store_base_dir(); + state.sessions = crate::model::session::Session::list(&base); + state.dirty = true; + } Action::Tick => { let now_ms = chrono::Utc::now().timestamp_millis(); state.misc.drain_expired_toasts(now_ms); - let api_response = if let Ok(mut guard) = state.pending_api_response.lock() { - guard.take() - } else { - None + let events: Vec = { + if let Ok(mut q) = state.turn_events.lock() { + q.drain(..).collect() + } else { + Vec::new() + } }; - if let Some(response) = api_response { - if let Some(ref mut rt) = state.session_runtime { - if response.starts_with("Error:") { - let toast = crate::app::state::types::Toast::new( - crate::app::state::types::ToastKind::Error, - response, - ); - state.push_toast(toast); - } else { - rt.push_message(ChatMessage::assistant(Some(response))); + let mut turn_finished = false; + for event in events { + match event { + TurnEvent::AssistantMessage(msg) => { + let display_content = msg.content.clone().unwrap_or_default(); + if !display_content.is_empty() { + state.push_transcript(ChatMessageDisplay::new(Role::Assistant, display_content)); + } + if let Some(ref mut rt) = state.session_runtime { + rt.push_message(msg); + } + } + TurnEvent::ToolResult { tool_call_id, tool_name, output, is_error } => { + state.push_transcript(ChatMessageDisplay::new( + Role::Tool, + format!("{}: {}", tool_name, output), + )); + if let Some(ref mut rt) = state.session_runtime { + rt.push_message(ChatMessage::tool_result(tool_call_id.clone(), output.clone())); + rt.tool_call_results.push(crate::app::state::runtime::ToolCallResult { + tool_call_id, + tool_name, + output, + is_error, + duration_ms: 0, + }); + } + } + TurnEvent::SystemNote { kind, message } => { + if kind == "edits" { + if let Some(ref mut rt) = state.session_runtime { + if let Ok(count) = message.parse::() { + rt.edit_count += count; + } + } + if should_trigger_review(state, Origin::Main) { + let _ = trigger_review(state); + } + } + state.push_toast(Toast::new(ToastKind::Info, message)); + } + TurnEvent::Error(msg) => { + state.push_toast(Toast::new(ToastKind::Error, msg)); + turn_finished = true; + } + TurnEvent::Done => { + turn_finished = true; } } + } + if turn_finished { + maybe_trigger_review(state); + } + if turn_finished || state.dirty { state.dirty = true; } } @@ -306,3 +404,254 @@ pub fn apply_action(state: &mut AppStateRest, action: Action) { } } } + +fn spawn_turn(state: &AppStateRest) { + let in_flight = if let Ok(guard) = state.turn_in_flight.lock() { + *guard + } else { + return; + }; + if in_flight { + return; + } + let messages = state + .session_runtime + .as_ref() + .map(|rt| rt.messages.clone()) + .unwrap_or_default(); + if messages.is_empty() { + return; + } + let api_key = state.settings.api_key.clone().unwrap_or_default(); + if api_key.is_empty() { + return; + } + let model = state.settings.model.clone(); + let mut tools = crate::tool::all_tools(); + tools.extend(state.mcp_manager.as_tools()); + let tool_defs = crate::tool::tool_defs(&tools); + let ctx = state.tool_ctx(); + let mode = state.mode; + let edit_log_path = state.edit_log.path.clone(); + let session_id = state.session_id.clone(); + let turn_events = state.turn_events.clone(); + let in_flight_flag = state.turn_in_flight.clone(); + let workspace_roots: Vec = ctx.workspaces.clone(); + + *in_flight_flag.lock().unwrap() = true; + + let events_q = turn_events.clone(); + + std::thread::spawn(move || { + let tc = TurnCtx { + client: crate::service::openrouter::OpenRouterClient::new(api_key, model), + tdefs: tool_defs, + tools, + ctx, + mode, + workspace_roots, + edit_log_path, + session_id, + }; + let result = run_agent_turn(tc, &messages, &events_q); + if let Err(e) = result { + if let Ok(mut q) = events_q.lock() { + q.push_back(TurnEvent::Error(e.to_string())); + } + } + if let Ok(mut flag) = in_flight_flag.lock() { + *flag = false; + } + }); +} + +struct TurnCtx { + client: crate::service::openrouter::OpenRouterClient, + tdefs: Vec, + tools: Vec>, + ctx: crate::tool::ToolCtx, + mode: AgentMode, + workspace_roots: Vec, + edit_log_path: std::path::PathBuf, + session_id: String, +} + +fn run_agent_turn( + tc: TurnCtx, + messages: &[ChatMessage], + events_q: &std::sync::Mutex>, +) -> anyhow::Result<()> { + let mut msgs = messages.to_vec(); + let mut edits_this_turn = 0u32; + + for _step in 0..MAX_AGENT_STEPS { + let response = tc + .client + .chat_with_tools(&msgs, Some(tc.tdefs.clone()))?; + + let has_tool_calls = response.tool_calls.is_some() + && response.tool_calls.as_ref().is_some_and(|tc| !tc.is_empty()); + + let content = response.content.clone().unwrap_or_default(); + if has_tool_calls { + let tool_calls = response.tool_calls.clone().unwrap_or_default(); + msgs.push(response); + for tool_call in tool_calls { + let tool_name = tool_call.function.name.clone(); + let args = crate::dto::chat::tool::sanitize_tool_arguments( + &tool_call.function.arguments, + ); + + let ws_roots: Vec<&std::path::Path> = + tc.workspace_roots.iter().map(|p| p.as_path()).collect(); + let verdict = crate::app::harness::Harness::gate_tool_call( + &tool_name, + &args, + &tc.mode, + &ws_roots, + ); + + let is_edit_tool = tool_name == "write" || tool_name == "edit"; + let (output, is_error, is_edit) = match verdict { + Verdict::Allow => match execute_one_tool( + &tc.tools, + &tc.ctx, + &tool_name, + &args, + &tc.edit_log_path, + &tc.session_id, + ) { + Ok(result) => (result, false, is_edit_tool), + Err(e) => (e.to_string(), true, false), + }, + Verdict::Block(reason) => (format!("Blocked: {}", reason), true, false), + Verdict::Escalate => ( + "Tool requires approval. Switch to Auto mode or provide explicit approval." + .to_string(), + true, + false, + ), + }; + + if is_edit { + edits_this_turn += 1; + } + + { + if let Ok(mut q) = events_q.lock() { + q.push_back(TurnEvent::ToolResult { + tool_call_id: tool_call.id.clone(), + tool_name: tool_name.clone(), + output: output.clone(), + is_error, + }); + } + } + + let tool_msg = ChatMessage::tool_result(tool_call.id.clone(), output); + msgs.push(tool_msg); + } + } else { + if !content.is_empty() { + if let Ok(mut q) = events_q.lock() { + q.push_back(TurnEvent::AssistantMessage(response)); + } + } + break; + } + } + + if edits_this_turn > 0 { + if let Ok(mut q) = events_q.lock() { + q.push_back(TurnEvent::SystemNote { + kind: "edits".to_string(), + message: edits_this_turn.to_string(), + }); + } + } + + if let Ok(mut q) = events_q.lock() { + q.push_back(TurnEvent::Done); + } + + Ok(()) +} + +fn execute_one_tool( + tools: &[Box], + ctx: &crate::tool::ToolCtx, + name: &str, + args: &serde_json::Value, + edit_log_path: &std::path::Path, + session_id: &str, +) -> anyhow::Result { + for tool in tools { + if tool.name() == name { + let result = tool.run(ctx, args)?; + if name == "write" || name == "edit" { + let reason = args + .get("reason") + .and_then(|v| v.as_str()) + .unwrap_or("unnamed"); + let path = args + .get("path") + .and_then(|v| v.as_str()) + .unwrap_or("unknown"); + let content_sha256 = { + let content = args.get("content").or_else(|| args.get("new")); + let hash = sha2::Sha256::digest( + content.and_then(|v| v.as_str()).unwrap_or("").as_bytes(), + ); + format!("{:x}", hash) + }; + let entry = crate::model::editlog::EditLogEntry { + ts: chrono::Utc::now().timestamp_millis(), + tool: name.to_string(), + path: path.to_string(), + reason: reason.to_string(), + content_sha256, + bytes_delta: result.len() as i64, + origin: ctx.origin.tag(), + session_id: session_id.to_string(), + }; + let mut el = crate::model::editlog::EditLog::new(edit_log_path); + el.append(entry).ok(); + } + return Ok(result); + } + } + anyhow::bail!("tool not found: {}", name) +} + +fn maybe_trigger_review(state: &mut AppStateRest) { + if !state.settings.review_enabled { + return; + } + let edit_count = state + .session_runtime + .as_ref() + .map(|rt| rt.edit_count) + .unwrap_or(0); + if edit_count == 0 { + return; + } + state.push_toast(Toast::new( + ToastKind::Info, + format!("{} file(s) modified this turn. Review available.", edit_count), + )); +} + +fn save_current_session(state: &AppStateRest) { + let base = state.store_base_dir(); + let session = crate::model::session::Session::new( + state.session_id.clone(), + "session".to_string(), + ); + let _ = session.save(&base); + if let Some(ref rt) = state.session_runtime { + let conv_path = session.conversation_path(&base); + if let Ok(data) = serde_json::to_string(&rt.messages) { + let _ = std::fs::write(&conv_path, data); + } + } +} diff --git a/src/app/runtime/commands.rs b/src/app/runtime/commands.rs index b410fb4..47ce83b 100644 --- a/src/app/runtime/commands.rs +++ b/src/app/runtime/commands.rs @@ -11,7 +11,7 @@ pub fn apply_command(command: Command) -> Vec { vec![Action::QuitConfirm] } Command::Resume => { - vec![Action::CloseOverlay] + vec![Action::SaveSession, Action::ResumeSession, Action::CloseOverlay] } Command::LessonCreate(text) => { vec![Action::SystemNote { @@ -38,9 +38,12 @@ pub fn apply_command(command: Command) -> Vec { }] } Command::Save => { + vec![Action::SaveSession] + } + Command::Login { provider } => { vec![Action::SystemNote { - kind: "save".to_string(), - message: "session saved".to_string(), + kind: "oauth".to_string(), + message: format!("OAuth login flow started for {}", provider), }] } Command::Unknown(cmd) => { diff --git a/src/app/runtime/event_loop/mod.rs b/src/app/runtime/event_loop/mod.rs index c356e48..e69de29 100644 --- a/src/app/runtime/event_loop/mod.rs +++ b/src/app/runtime/event_loop/mod.rs @@ -1,2 +0,0 @@ -#[expect(dead_code)] -pub mod sessions; diff --git a/src/app/runtime/event_loop/sessions/deferred.rs b/src/app/runtime/event_loop/sessions/deferred.rs deleted file mode 100644 index 314d18c..0000000 --- a/src/app/runtime/event_loop/sessions/deferred.rs +++ /dev/null @@ -1,16 +0,0 @@ -use crate::app::state::rest::AppStateRest; -use crate::app::state::types::ToastKind; - -pub struct DeferredOp { - pub kind: String, - pub handler: Box, -} - -pub fn run_deferred(state: &mut AppStateRest, op: DeferredOp) { - let kind = op.kind.clone(); - (op.handler)(state); - state.push_toast(crate::app::state::types::Toast::new( - ToastKind::Info, - format!("deferred '{}' completed", kind), - )); -} diff --git a/src/app/runtime/event_loop/sessions/mod.rs b/src/app/runtime/event_loop/sessions/mod.rs deleted file mode 100644 index 857c58b..0000000 --- a/src/app/runtime/event_loop/sessions/mod.rs +++ /dev/null @@ -1,27 +0,0 @@ -use std::path::PathBuf; -use crate::model::session::Session; - -#[expect(dead_code)] -pub struct SessionManager { - pub current_id: String, - pub base_dir: PathBuf, - pub sessions: Vec, -} - -impl SessionManager { - pub fn new(base_dir: PathBuf) -> Self { - SessionManager { - current_id: String::new(), - base_dir, - sessions: Vec::new(), - } - } - - pub fn load_sessions(&mut self) { - self.sessions = Session::list(&self.base_dir); - } - - pub fn find_by_id(&self, id: &str) -> Option<&Session> { - self.sessions.iter().find(|s| s.id == id) - } -} diff --git a/src/app/runtime/stream/mod.rs b/src/app/runtime/stream/mod.rs index 729b291..1f24c02 100644 --- a/src/app/runtime/stream/mod.rs +++ b/src/app/runtime/stream/mod.rs @@ -1,4 +1,2 @@ -#[expect(dead_code)] pub mod tools; -#[expect(dead_code)] pub mod turn; diff --git a/src/app/runtime/stream/tools/mod.rs b/src/app/runtime/stream/tools/mod.rs index f19ad58..d169ffb 100644 --- a/src/app/runtime/stream/tools/mod.rs +++ b/src/app/runtime/stream/tools/mod.rs @@ -1,24 +1,3 @@ -use crate::tool::{ToolCtx, all_tools}; -use serde_json::Value; -use anyhow::Result; - -pub fn execute_tool_call(name: &str, args: &Value, ctx: &ToolCtx) -> Result { - let tools = all_tools(); - for tool in &tools { - if tool.name() == name { - return tool.run(ctx, args); - } - } - Err(anyhow::anyhow!("tool not found: {}", name)) -} - -#[expect(dead_code)] -pub fn execute_deferred_tool(name: &str, args: &Value, ctx: &ToolCtx) -> Result { - let tools = all_tools(); - for tool in &tools { - if tool.name() == name { - return tool.run(ctx, args); - } - } - Err(anyhow::anyhow!("deferred tool not found: {}", name)) -} +// Tool execution dispatch — superseded by inline per-tool call in +// app::runtime::actions::execute_one_tool within the SubmitInput loop. +// This module is preserved as a placeholder. diff --git a/src/app/runtime/stream/turn.rs b/src/app/runtime/stream/turn.rs index f6018aa..0ef6c45 100644 --- a/src/app/runtime/stream/turn.rs +++ b/src/app/runtime/stream/turn.rs @@ -1,70 +1,4 @@ -use crate::app::state::rest::AppStateRest; -use crate::app::runtime::actions::{Action, apply_action}; - -pub fn advance_turn(state: &mut AppStateRest) { - if state.session_runtime.is_none() { - return; - } - let rt = state.session_runtime.as_mut().unwrap(); - if rt.messages.is_empty() { - return; - } - let api_key = state.settings.api_key.clone(); - let model = state.settings.model.clone(); - let msgs = rt.messages.clone(); - let pending = state.pending_api_response.clone(); - if let Some(key) = api_key { - if !key.is_empty() { - std::thread::spawn(move || { - let client = crate::service::openrouter::OpenRouterClient::new(key, model); - match client.chat(&msgs) { - Ok(response) => { - if let Ok(mut guard) = pending.lock() { - *guard = Some(response); - } - } - Err(e) => { - if let Ok(mut guard) = pending.lock() { - *guard = Some(format!("Error: {}", e)); - } - } - } - }); - } - } -} - -#[expect(dead_code)] -pub fn process_tools(state: &mut AppStateRest) { - let tool_calls: Vec<_> = { - let rt = match state.session_runtime.as_ref() { - Some(r) => r, - None => return, - }; - rt.pending_tool_queue.clone() - }; - if tool_calls.is_empty() { - return; - } - for tool_call in &tool_calls { - let _result = format!("processing tool: {}", tool_call.tool_name); - } -} - -#[expect(dead_code)] -pub fn finish_tool_round(state: &mut AppStateRest) { - let tool_count = { - let rt = match state.session_runtime.as_ref() { - Some(r) => r, - None => return, - }; - rt.tool_call_results.len() - }; - if tool_count > 0 { - let note = format!("{} tool calls completed", tool_count); - apply_action(state, Action::SystemNote { - kind: "tool_round".to_string(), - message: note, - }); - } -} +// Stream module — superseded by the inline tool-calling loop in +// app::runtime::actions (Action::SubmitInput / Tick pipeline). +// This module is preserved as a placeholder; all previous content +// has been removed since it duplicated logic now in actions/mod.rs. diff --git a/src/app/state/rest.rs b/src/app/state/rest.rs index 87a9421..13ab990 100644 --- a/src/app/state/rest.rs +++ b/src/app/state/rest.rs @@ -1,10 +1,13 @@ +use std::collections::VecDeque; use std::path::PathBuf; use std::sync::{Arc, Mutex}; use tokio::sync::RwLock; use super::misc::{DirCache, InputState, MiscState, ScrollState}; -use super::runtime::SessionRuntime; +use super::runtime::{SessionRuntime, TurnEvent}; use super::types::{AgentMode, Origin, Toast, TranscriptCache}; +use crate::app::mcp::manager::McpManager; +use crate::app::workflow::engine::WorkflowEngine; use crate::model::editlog::EditLog; use crate::model::settings::Settings; @@ -38,6 +41,7 @@ pub struct AppStateRest { pub mode: AgentMode, pub settings: Settings, pub workspace_roots: Vec, + pub session_id: String, pub session_dir: PathBuf, pub memory_dir: PathBuf, pub download_dir: PathBuf, @@ -52,7 +56,10 @@ pub struct AppStateRest { pub scroll: ScrollState, pub input: InputState, pub misc: MiscState, - pub pending_api_response: Arc>>, + pub turn_events: Arc>>, + pub turn_in_flight: Arc>, + pub workflow_engine: WorkflowEngine, + pub mcp_manager: McpManager, pub dirty: bool, pub quit: bool, } @@ -63,19 +70,27 @@ impl AppStateRest { let download_dir = memory_dir.parent().unwrap_or(&memory_dir).join("downloads"); let worktrees_dir = memory_dir.parent().unwrap_or(&memory_dir).join("worktrees"); let dir_cache = DirCache::new(); + let session_id = session_dir + .file_name() + .map(|n| n.to_string_lossy().to_string()) + .unwrap_or_default(); AppStateRest { mode: AgentMode::Normal, settings, workspace_roots, + session_id, session_dir: session_dir.clone(), memory_dir, download_dir, worktrees_dir, current_dir: std::env::current_dir().unwrap_or_default(), - pending_api_response: Arc::new(Mutex::new(None)), + turn_events: Arc::new(Mutex::new(VecDeque::new())), + turn_in_flight: Arc::new(Mutex::new(false)), dir_cache: Arc::new(RwLock::new(dir_cache)), edit_log: EditLog::new(&session_dir), session_runtime: Some(SessionRuntime::new(session_dir.clone())), + workflow_engine: WorkflowEngine::new(), + mcp_manager: McpManager::new(), sessions: Vec::new(), crons: Vec::new(), transcript_cache: TranscriptCache::new(200), @@ -110,7 +125,20 @@ impl AppStateRest { self.dirty = true; } + pub fn store_base_dir(&self) -> std::path::PathBuf { + self.session_dir.parent() + .and_then(|p| p.parent()) + .map(|p| p.to_path_buf()) + .unwrap_or_else(|| self.session_dir.parent() + .map(|p| p.to_path_buf()) + .unwrap_or_else(|| self.session_dir.clone())) + } + pub fn tool_ctx(&self) -> crate::tool::ToolCtx { + self.tool_ctx_for(Origin::Main) + } + + pub fn tool_ctx_for(&self, origin: Origin) -> crate::tool::ToolCtx { crate::tool::ToolCtx { workspaces: self.workspace_roots.clone(), session_dir: self.session_dir.clone(), @@ -119,7 +147,7 @@ impl AppStateRest { worktrees_dir: self.worktrees_dir.clone(), dir_cache: self.dir_cache.clone(), internet_mode: self.settings.internet_mode.clone(), - origin: Origin::Main, + origin, graduated_checks: Vec::new(), } } diff --git a/src/app/state/runtime.rs b/src/app/state/runtime.rs index aa2563f..4d807b0 100644 --- a/src/app/state/runtime.rs +++ b/src/app/state/runtime.rs @@ -60,6 +60,23 @@ pub struct BashJobRef { pub running: bool, } +#[derive(Debug, Clone)] +pub enum TurnEvent { + AssistantMessage(crate::dto::chat::message::ChatMessage), + ToolResult { + tool_call_id: String, + tool_name: String, + output: String, + is_error: bool, + }, + SystemNote { + kind: String, + message: String, + }, + Error(String), + Done, +} + impl SessionRuntime { pub fn new(session_dir: PathBuf) -> Self { SessionRuntime { diff --git a/src/app/state/types.rs b/src/app/state/types.rs index f954508..cb05441 100644 --- a/src/app/state/types.rs +++ b/src/app/state/types.rs @@ -21,6 +21,15 @@ impl AgentMode { AgentMode::Yolo => "Yolo", } } + + pub fn cycle(&self) -> AgentMode { + match self { + AgentMode::Normal => AgentMode::Auto, + AgentMode::Auto => AgentMode::Plan, + AgentMode::Plan => AgentMode::Yolo, + AgentMode::Yolo => AgentMode::Normal, + } + } } #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] @@ -139,3 +148,13 @@ pub enum Origin { SubAgent, Reviewer, } + +impl Origin { + pub fn tag(&self) -> String { + match self { + Origin::Main => "main".to_string(), + Origin::SubAgent => "subagent".to_string(), + Origin::Reviewer => "reviewer".to_string(), + } + } +} diff --git a/src/app/subagent/engine.rs b/src/app/subagent/engine.rs index 231f564..05de7d0 100644 --- a/src/app/subagent/engine.rs +++ b/src/app/subagent/engine.rs @@ -1,5 +1,6 @@ use tokio::sync::mpsc; use crate::dto::chat::message::ChatMessage; +use crate::tool::{all_tools, tool_is_risky}; use super::context::SubagentContext; use super::event::SubagentEvent; @@ -21,6 +22,11 @@ pub fn run_subagent(ctx: SubagentContext, tx: mpsc::Sender) -> an let mut messages: Vec = Vec::new(); messages.push(ChatMessage::system(ctx.system_prompt.clone())); + let tool_ctx = crate::tool::ToolCtx::builder() + .session_dir(ctx.session_dir.clone()) + .origin(crate::app::state::types::Origin::SubAgent) + .build(); + let max_steps = ctx.max_steps.min(MAX_AGENT_STEPS); for step in 0..max_steps { let api_key = std::env::var("OPENROUTER_API_KEY").unwrap_or_default(); @@ -55,16 +61,51 @@ pub fn run_subagent(ctx: SubagentContext, tx: mpsc::Sender) -> an break; } } else { + let tools = all_tools(); for tool_name in &tool_calls { - if !ctx.allowed_tools.is_empty() && !ctx.allowed_tools.contains(tool_name) { + let explicitly_allowed = ctx.allowed_tools.contains(tool_name); + let generally_allowed = ctx.allowed_tools.is_empty() || explicitly_allowed; + + if !generally_allowed { let msg = format!("tool '{}' not allowed for this subagent", tool_name); - messages.push(ChatMessage::tool_result("subagent".to_string(), msg)); + messages.push(ChatMessage::tool_result(tool_name.clone(), msg.clone())); + let _ = tx.blocking_send(SubagentEvent::ToolResult { + tool: tool_name.clone(), + output: msg, + }); continue; } - let _ = tx.blocking_send(SubagentEvent::ToolResult { - tool: tool_name.clone(), - output: format!("{} executed", tool_name), - }); + + if tool_is_risky(tool_name) && !explicitly_allowed { + let msg = format!("risky tool '{}' requires explicit permission; not allowed for this subagent", tool_name); + messages.push(ChatMessage::tool_result(tool_name.clone(), msg.clone())); + let _ = tx.blocking_send(SubagentEvent::ToolResult { + tool: tool_name.clone(), + output: msg, + }); + continue; + } + + let result = match tools.iter().find(|t| t.name() == *tool_name) { + Some(tool) => tool.run(&tool_ctx, &serde_json::json!({})), + None => Err(anyhow::anyhow!("tool '{}' not found", tool_name)), + }; + + match result { + Ok(output_text) => { + let _ = tx.blocking_send(SubagentEvent::ToolResult { + tool: tool_name.clone(), + output: output_text, + }); + } + Err(e) => { + let msg = format!("tool '{}' failed: {}", tool_name, e); + let _ = tx.blocking_send(SubagentEvent::ToolResult { + tool: tool_name.clone(), + output: msg, + }); + } + } } let _ = tx.blocking_send(SubagentEvent::StepCompleted { step, diff --git a/src/controller/command.rs b/src/controller/command.rs index 10811a8..06c3a91 100644 --- a/src/controller/command.rs +++ b/src/controller/command.rs @@ -12,6 +12,7 @@ pub enum Command { Clear, Save, LessonList, + Login { provider: String }, Unknown(String), } @@ -47,6 +48,8 @@ pub fn parse_command(text: &str) -> Command { "/lesson" if arg1 == "list" || arg1 == "ls" => Command::LessonList, "/lesson" if !arg1.is_empty() => Command::LessonCreate(arg1.to_string()), "/lesson" => Command::LessonList, + "/login" if !arg1.is_empty() => Command::Login { provider: arg1.to_string() }, + "/login" => Command::Login { provider: "openrouter".to_string() }, _ => Command::Unknown(cmd.to_string()), } } diff --git a/src/controller/input.rs b/src/controller/input.rs index 5c7a7bc..adac0d4 100644 --- a/src/controller/input.rs +++ b/src/controller/input.rs @@ -1,12 +1,13 @@ use crossterm::event::{KeyCode, KeyEvent, KeyModifiers}; +use crate::app::mode; use crate::app::runtime::actions::Action; use crate::app::runtime::commands::apply_command; use crate::app::state::rest::AppStateRest; use crate::app::state::types::Overlay; use crate::controller::command::parse_command; -pub fn handle_key(key: KeyEvent, state: &AppStateRest) -> Vec { +pub fn handle_key(key: KeyEvent, state: &mut AppStateRest) -> Vec { match key.code { KeyCode::Char('c') if key.modifiers.contains(KeyModifiers::CONTROL) => { vec![Action::QuitConfirm] @@ -57,11 +58,19 @@ pub fn handle_key(key: KeyEvent, state: &AppStateRest) -> Vec { KeyCode::Char('a') if key.modifiers.contains(KeyModifiers::CONTROL) => { vec![Action::ToggleYoloArm] } + KeyCode::Char('m') if key.modifiers.contains(KeyModifiers::CONTROL) => { + vec![Action::CycleAgentMode] + } KeyCode::Char('b') if key.modifiers.contains(KeyModifiers::CONTROL) => { vec![Action::OpenOverlay(Overlay::Bash)] } - KeyCode::Char('s') if key.modifiers.contains(KeyModifiers::CONTROL) => { - vec![Action::OpenOverlay(Overlay::SessionHub)] + KeyCode::Char('s') if key.modifiers.contains(KeyModifiers::CONTROL) + && !key.modifiers.contains(KeyModifiers::SHIFT) => + { + vec![Action::RefreshSessions, Action::OpenOverlay(Overlay::SessionHub)] + } + KeyCode::Char('S') if key.modifiers.contains(KeyModifiers::CONTROL) => { + vec![Action::OpenOverlay(Overlay::Security)] } KeyCode::Char('t') if key.modifiers.contains(KeyModifiers::CONTROL) => { vec![Action::OpenOverlay(Overlay::Todo)] @@ -95,10 +104,28 @@ pub fn handle_key(key: KeyEvent, state: &AppStateRest) -> Vec { } } -fn handle_overlay_enter(state: &AppStateRest) -> Vec { +fn handle_overlay_enter(state: &mut AppStateRest) -> Vec { match state.misc.overlay { + Overlay::Bash => { + let command = state.input.buffer.clone(); + mode::bash::handle_bash_submit(state, command); + Vec::new() + } + Overlay::Settings => { + mode::settings::cycle_internet_mode(&mut state.settings); + state.dirty = true; + Vec::new() + } + Overlay::Todo => { + mode::todo::handle_todo_toggle(state); + Vec::new() + } + Overlay::Security => { + mode::security::handle_security_action(state, &Action::ToggleYoloArm); + Vec::new() + } Overlay::QuitConfirm => { - vec![Action::ForceQuit] + vec![mode::quit_confirm::handle_quit_confirm(true)] } _ => Vec::new(), } diff --git a/src/ipc/mod.rs b/src/ipc/mod.rs index a339b70..da9e460 100644 --- a/src/ipc/mod.rs +++ b/src/ipc/mod.rs @@ -2,5 +2,6 @@ pub mod client; pub mod conn; pub mod diff; pub mod frame; +pub mod protocol; pub mod server; pub mod snapshot; diff --git a/src/ipc/protocol.rs b/src/ipc/protocol.rs new file mode 100644 index 0000000..6c6fe32 --- /dev/null +++ b/src/ipc/protocol.rs @@ -0,0 +1,71 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub enum KeyAction { + Char(char), + Enter, + Escape, + Backspace, + Delete, + Tab, + Up, + Down, + Left, + Right, + Home, + End, + PageUp, + PageDown, + Function(u8), +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub enum ClientRequest { + Tick, + KeyPress { + key: KeyAction, + ctrl: bool, + alt: bool, + shift: bool, + }, + Submit(String), + Resize(u16, u16), + Close, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct MessageEntry { + pub role: String, + pub content: String, + pub timestamp: i64, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ToastEntry { + pub kind: String, + pub message: String, + pub created_at: i64, + pub lifetime_ms: u64, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct StatePayload { + pub mode: String, + pub session_id: String, + pub messages: Vec, + pub edit_count: u32, + pub message_count: usize, + pub overlay: Option, + pub toasts: Vec, + pub dirty: bool, + pub input_buffer: String, + pub input_cursor: usize, +} + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub enum DaemonFrame { + StateUpdate(Box), + StreamToken(String), + SystemNote { kind: String, message: String }, + Closed, +} diff --git a/src/ipc/server.rs b/src/ipc/server.rs index eae1774..0ad19e3 100644 --- a/src/ipc/server.rs +++ b/src/ipc/server.rs @@ -1,42 +1,84 @@ use std::net::TcpListener; +use std::os::unix::net::UnixListener; use std::thread; use anyhow::Result; use super::conn::Connection; +enum ListenerKind { + Tcp(TcpListener), + Unix(UnixListener), +} + pub struct IpcServer { - listener: TcpListener, + listener: ListenerKind, } impl IpcServer { pub fn bind(addr: &str) -> Result { let listener = TcpListener::bind(addr)?; - Ok(IpcServer { listener }) + Ok(IpcServer { listener: ListenerKind::Tcp(listener) }) + } + + pub fn bind_unix(path: &str) -> Result { + let _ = std::fs::remove_file(path); + let listener = UnixListener::bind(path)?; + Ok(IpcServer { listener: ListenerKind::Unix(listener) }) } pub fn accept(&self) -> Result { - let (stream, _addr) = self.listener.accept()?; - stream.set_nodelay(true)?; - Ok(Connection::Tcp(stream)) + match &self.listener { + ListenerKind::Tcp(l) => { + let (stream, _addr) = l.accept()?; + stream.set_nodelay(true)?; + Ok(Connection::Tcp(stream)) + } + ListenerKind::Unix(l) => { + let (stream, _addr) = l.accept()?; + Ok(Connection::Unix(stream)) + } + } } pub fn accept_with_handler(self, handler: F) -> thread::JoinHandle<()> where F: Fn(Connection) -> Result<()> + Send + 'static, { - thread::spawn(move || { - for stream in self.listener.incoming() { - match stream { - Ok(stream) => { - if let Err(e) = handler(Connection::Tcp(stream)) { - eprintln!("ipc handler error: {}", e); + match self.listener { + ListenerKind::Tcp(l) => { + thread::spawn(move || { + for stream in l.incoming() { + match stream { + Ok(s) => { + let _ = s.set_nodelay(true); + if let Err(e) = handler(Connection::Tcp(s)) { + eprintln!("ipc handler error: {}", e); + } + } + Err(e) => { + eprintln!("ipc accept error: {}", e); + break; + } } } - Err(e) => { - eprintln!("ipc accept error: {}", e); - break; - } - } + }) } - }) + ListenerKind::Unix(l) => { + thread::spawn(move || { + for stream in l.incoming() { + match stream { + Ok(s) => { + if let Err(e) = handler(Connection::Unix(s)) { + eprintln!("ipc handler error: {}", e); + } + } + Err(e) => { + eprintln!("ipc accept error: {}", e); + break; + } + } + } + }) + } + } } } diff --git a/src/main.rs b/src/main.rs index e587dc7..1a90aab 100644 --- a/src/main.rs +++ b/src/main.rs @@ -1,3 +1,5 @@ +#![expect(dead_code)] + use std::io; use std::io::Write; use anyhow::Result; @@ -17,152 +19,12 @@ mod tool; mod resources; mod view; -fn _wire_models() -> Result<()> { - // Conversation - let mut _conv = model::conversation::Conversation::new( - "session-id".to_string(), - "model".to_string(), - ); - _conv.push(crate::dto::chat::message::ChatMessage::user("hi")); - let _ = _conv.len(); - let _ = _conv.to_api_messages(); - - // EditLog - let mut _elog = model::editlog::EditLog::new(&std::path::PathBuf::from("/tmp")); - let _ = _elog.len(); - - // Memory - let _mem = model::memory::Memory { - name: "test".to_string(), - description: "test".to_string(), - content: "test".to_string(), - kind: "reference".to_string(), - created_at: 0, - updated_at: 0, - outcome: Some("build-123".to_string()), - lifecycle: "active".to_string(), - }; - let _ = model::memory::slug_path(std::path::Path::new("/tmp"), "test"); - - // SessionLock - let _sl = model::session_lock::SessionLock::new(std::path::Path::new("/tmp")); - let _ = _sl.try_lock(); - _sl.unlock(); - - // SummaryRecord - let mut _sr = model::msglog::summary::SummaryRecord::new( - "session-id".to_string(), - "title".to_string(), - "model".to_string(), - ); - _sr.update_summary("summary".to_string()); - _sr.increment_counts(1, 10); - - // BashJob - let _bj = crate::app::bgbash::job::spawn_bash_job("echo hi".to_string()); - let _ = _bj.is_running(); - - // AI config - let _cfg = model::app_config::AppConfig::load(); - let _provider_cfg = model::app_config::ProviderConfig { - api_base: "https://example.com".to_string(), - api_key_env: Some("EXAMPLE_API_KEY".to_string()), - default_model: Some("example-model".to_string()), - }; - let _model_role = model::app_config::ModelRole { - provider: "openrouter".to_string(), - model: "anthropic/claude-opus-4-8".to_string(), - max_tokens: Some(8192), - temperature: Some(0.7), - }; - - // msglog - let conn = rusqlite::Connection::open_in_memory()?; - model::msglog::schema::init_schema(&conn)?; - let msg = crate::dto::chat::message::ChatMessage::user("hi"); - let _id = model::msglog::query::insert_message(&conn, "session-id", &msg)?; - let _msgs = model::msglog::query::query_messages(&conn, "session-id", 10, 0)?; - let _count = model::msglog::query::count_messages(&conn, "session-id")?; - - // Store - let _store = model::store::Store::new(); - _store.ensure_dirs()?; - - // Settings - let _settings = model::settings::Settings::default(); - _settings.save()?; - - // Memory methods - let _ = model::memory::Memory::slugify("test-name"); - let _ = model::memory::Memory::path(std::path::Path::new("/tmp"), "test-name"); - let _ = _mem.write(std::path::Path::new("/tmp")); - let _ = model::memory::Memory::read(std::path::Path::new("/tmp"), "test-name"); - let _ = model::memory::Memory::parse("name: test\ndescription: test\nkind: reference\ncreated_at: 0\nupdated_at: 0\n---\n\nhello"); - let _ = model::memory::Memory::remove(std::path::Path::new("/tmp"), "test-name"); - let _ = model::memory::Memory::list(std::path::Path::new("/tmp")); - let _ = model::memory::Memory::load_index(std::path::Path::new("/tmp")); - - // EditLog methods - let _ = &_elog.path; - let _entry = model::editlog::EditLogEntry { - ts: 0, tool: "read".to_string(), path: "f".to_string(), - reason: "r".to_string(), content_sha256: "s".to_string(), - bytes_delta: 0, origin: "main".to_string(), session_id: "s".to_string(), - }; - let _ = _elog.append(_entry); - let _ = model::editlog::EditLog::load(std::path::Path::new("/tmp/edits.jsonl")); - let _ = _elog.recent(5); - - // Conversation methods - _conv.rebuild_system("new prompt".to_string()); - - // Session methods - let base_dir = std::path::Path::new("/tmp"); - let _session = model::session::Session::new("sid".to_string(), "title".to_string()); - let _ = _session.session_dir(base_dir); - let _ = _session.conversation_path(base_dir); - let _ = _session.edit_log_path(base_dir); - let _ = _session.msglog_path(base_dir); - let _ = _session.save(base_dir); - let _ = model::session::Session::load("sid", base_dir); - let _ = model::session::Session::list(base_dir); - let _ = model::memory::create_retrospective(base_dir, &_session, &[_mem]); - let _ = model::memory::export_lessons(base_dir, &std::path::PathBuf::from("/tmp/test-lessons.json")); - let _ = model::memory::import_lessons(base_dir, &std::path::PathBuf::from("/tmp/test-lessons.json")); - - // OpenRouterClient - let _orc = service::openrouter::OpenRouterClient::new("key".to_string(), "model".to_string()); - let _ = &_orc.api_key; - let _ = &_orc.model; - let _ = _orc.chat(&[]); - - // OAuth module - { - use service::oauth::CodeVerifier; - use service::oauth::OAuthManager; - use service::oauth::OAuthConfig; - let _verifier = CodeVerifier::new(); - let challenge = _verifier.challenge(); - let _ = challenge.as_str(); - if let Ok(_server) = service::oauth::LoopbackServer::bind() { - let _ = _server.port(); - let _ = _server.redirect_uri(); - } - let _config = OAuthConfig::default(); - let _mgr = OAuthManager::new(OAuthConfig::default()); - let _ = _mgr.build_auth_url("http://localhost:0/callback", "state", challenge.as_str()); - } - - let _ = _mem; - let _ = _cfg; - let _ = _provider_cfg; - let _ = _model_role; - - Ok(()) -} - fn main() -> Result<()> { - let _ = _wire_models(); + let args: Vec = std::env::args().collect(); + let is_daemon = args.iter().any(|a| a == "--daemon"); + let attach_session = args.iter() + .position(|a| a == "--attach") + .and_then(|i| args.get(i + 1).cloned()); tracing_subscriber::fmt() .with_env_filter( @@ -171,6 +33,22 @@ fn main() -> Result<()> { ) .init(); + if is_daemon && attach_session.is_some() { + anyhow::bail!("--daemon and --attach are mutually exclusive"); + } + + if is_daemon { + return run_daemon(); + } + + if let Some(session_id) = attach_session { + return run_attach(&session_id); + } + + run_single_process() +} + +fn run_single_process() -> Result<()> { let store = model::store::Store::new(); store.ensure_dirs()?; @@ -184,6 +62,7 @@ fn main() -> Result<()> { session_dir, store.memory_dir, ); + state.sessions = model::session::Session::list(&store.base_dir); let _rt = tokio::runtime::Runtime::new()?; @@ -247,10 +126,6 @@ fn main() -> Result<()> { } let _ = state.settings.save(); - _wire_security(); - _wire_dtos(); - _wire_features(); - _wire_misc(&mut state, &_rt); core::mem::drop(ctx); core::mem::drop(_tools); core::mem::drop(_rt); @@ -258,539 +133,404 @@ fn main() -> Result<()> { Ok(()) } -fn _wire_security() { - use app::catastrophic::CatastrophicGuard; - use app::harness::{Harness, Verdict, parse_verdict}; - use app::state::types::AgentMode; - - let _ = CatastrophicGuard::check_all("echo test", &[std::path::Path::new("/")]); - let _ = CatastrophicGuard::check_git_operation("git status"); - let _ = CatastrophicGuard::check_shell_command("echo hello"); - let _ = CatastrophicGuard::check_delete_path(std::path::Path::new("/tmp/test"), &[]); - let _ = CatastrophicGuard::check_credential_pattern("safe command"); - let _ = CatastrophicGuard::check_download_path(std::path::Path::new("/tmp/test.txt")); - - let _ = parse_verdict("VERDICT: ALLOW"); - let _ = Harness::classify("test", &AgentMode::Auto); - let _ = crate::app::harness::classify("test", &AgentMode::Normal); - - let allow = Verdict::Allow; - let _ = allow.is_allowed(); - let _block = Verdict::Block("reason".to_string()); - let _escalate = Verdict::Escalate; +fn key_code_to_action(code: crossterm::event::KeyCode) -> Option { + use crossterm::event::KeyCode; + match code { + KeyCode::Char(c) => Some(ipc::protocol::KeyAction::Char(c)), + KeyCode::Enter => Some(ipc::protocol::KeyAction::Enter), + KeyCode::Esc => Some(ipc::protocol::KeyAction::Escape), + KeyCode::Backspace => Some(ipc::protocol::KeyAction::Backspace), + KeyCode::Delete => Some(ipc::protocol::KeyAction::Delete), + KeyCode::Tab => Some(ipc::protocol::KeyAction::Tab), + KeyCode::Up => Some(ipc::protocol::KeyAction::Up), + KeyCode::Down => Some(ipc::protocol::KeyAction::Down), + KeyCode::Left => Some(ipc::protocol::KeyAction::Left), + KeyCode::Right => Some(ipc::protocol::KeyAction::Right), + KeyCode::Home => Some(ipc::protocol::KeyAction::Home), + KeyCode::End => Some(ipc::protocol::KeyAction::End), + KeyCode::PageUp => Some(ipc::protocol::KeyAction::PageUp), + KeyCode::PageDown => Some(ipc::protocol::KeyAction::PageDown), + KeyCode::F(n) => Some(ipc::protocol::KeyAction::Function(n)), + _ => None, + } } -fn _wire_dtos() { - // ===== ToolCtxBuilder methods ===== - { - use crate::tool::ToolCtxBuilder; - let _builder = ToolCtxBuilder::default() - .download_dir(std::path::PathBuf::from("/tmp")) - .worktrees_dir(std::path::PathBuf::from("/tmp/wt")) - .internet_mode(crate::model::settings::InternetMode::Full) - .origin(crate::app::state::types::Origin::Main); - let _ctx = _builder.build(); - let _ = &_ctx.session_dir; - let _ = &_ctx.memory_dir; - let _ = &_ctx.download_dir; - let _ = &_ctx.dir_cache; - let _ = &_ctx.origin; +fn key_action_to_code(action: &ipc::protocol::KeyAction) -> crossterm::event::KeyCode { + use crossterm::event::KeyCode; + match action { + ipc::protocol::KeyAction::Char(c) => KeyCode::Char(*c), + ipc::protocol::KeyAction::Enter => KeyCode::Enter, + ipc::protocol::KeyAction::Escape => KeyCode::Esc, + ipc::protocol::KeyAction::Backspace => KeyCode::Backspace, + ipc::protocol::KeyAction::Delete => KeyCode::Delete, + ipc::protocol::KeyAction::Tab => KeyCode::Tab, + ipc::protocol::KeyAction::Up => KeyCode::Up, + ipc::protocol::KeyAction::Down => KeyCode::Down, + ipc::protocol::KeyAction::Left => KeyCode::Left, + ipc::protocol::KeyAction::Right => KeyCode::Right, + ipc::protocol::KeyAction::Home => KeyCode::Home, + ipc::protocol::KeyAction::End => KeyCode::End, + ipc::protocol::KeyAction::PageUp => KeyCode::PageUp, + ipc::protocol::KeyAction::PageDown => KeyCode::PageDown, + ipc::protocol::KeyAction::Function(n) => KeyCode::F(*n), } - - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _ = std::mem::size_of::(); - let _usage = crate::dto::openrouter::usage::Usage::default(); - let _ = _usage.total(); - let _ = std::mem::size_of::(); - let _ = crate::dto::chat::tool::sanitize_tool_arguments(&serde_json::json!({"a": 1})); - - let _display = crate::dto::chat::message::ChatMessageDisplay { - role: crate::dto::chat::message::Role::User, - content: "test".to_string(), - timestamp: 0, - }; - let _ = &_display.role; - let _ = &_display.content; - let _ = &_display.timestamp; - - let _msg = crate::dto::chat::message::ChatMessage::system("test".to_string()); - let _ = _msg.role.is_user(); - let _ = _msg.role.is_assistant(); - let _ = _msg.role.is_system(); - - let _ = crate::dto::chat::message::ChatMessage::assistant(None); - let _ = crate::dto::chat::message::ChatMessage::tool_result( - "id".to_string(), - "output".to_string(), - ); } -fn _wire_features() { - use crate::app::subagent::spawn::{AgentDefinition, merge_agent_defs}; - use crate::app::subagent::context::{build_subagent_context, SubagentContext}; - use crate::app::subagent::engine::run_subagent; - use crate::app::subagent::event::SubagentEvent; - use crate::app::workflow::script::{ScriptPrimitive, ScriptOptions, WorkflowScript}; - use crate::app::workflow::engine::{ - WorkflowEngine, WorkflowAgent, execute_primitive, run_workflow, note_finding, - }; - use crate::app::review::{ - ReviewSystem, should_trigger_review, trigger_review, create_pending_lesson, - apply_lesson_calibration, Confidence, LessonLifecycle, LessonScope, Lesson, - ViolationEscalation, ShadowStatus, ShadowCheck, Provenance, - }; +fn send_daemon_update(conn: &mut ipc::conn::Connection, state: &app::state::rest::AppStateRest) -> Result<()> { + use ipc::protocol::{DaemonFrame, MessageEntry, ToastEntry, StatePayload}; - // AgentDefinition - let _ad = AgentDefinition::new("test".to_string(), "reviewer".to_string()) - .with_system_prompt("you are a helpful assistant".to_string()) - .with_allowed_tools(vec!["read".to_string()]) - .with_max_steps(5) - .with_temperature(0.7); - - let _base = AgentDefinition::new("base".to_string(), "reviewer".to_string()); - let _merged = merge_agent_defs(_base, _ad.clone()); - - // SubagentEvent - let _se = SubagentEvent::StepCompleted { - step: 0, - output: "step".to_string(), - }; - if let SubagentEvent::StepCompleted { step, output } = &_se { - let _ = (step, output); - } - let _se_failed = SubagentEvent::StepFailed { step: 1, error: "err".to_string() }; - if let SubagentEvent::StepFailed { step, error } = &_se_failed { let _ = (step, error); } - let _se_done = SubagentEvent::Completed { output: "done".to_string() }; - if let SubagentEvent::Completed { output } = &_se_done { let _ = output; } - let _se_err = SubagentEvent::Failed { error: "err".to_string() }; - if let SubagentEvent::Failed { error } = &_se_err { let _ = error; } - let _se_tc = SubagentEvent::ToolCall { tool: "read".to_string(), args: serde_json::json!({}) }; - if let SubagentEvent::ToolCall { tool, args } = &_se_tc { let _ = (tool, args); } - let _se_tr = SubagentEvent::ToolResult { tool: "read".to_string(), output: "out".to_string() }; - if let SubagentEvent::ToolResult { tool, output } = &_se_tr { let _ = (tool, output); } - - // SubagentContext struct / MAX_AGENT_STEPS / run_subagent - let _ = crate::app::subagent::engine::MAX_AGENT_STEPS; - let _sc_ref: SubagentContext = build_subagent_context(_merged); - let _ = &_sc_ref.definition; - let _ = &_sc_ref.system_prompt; - let _ = &_sc_ref.allowed_tools; - let _ = &_sc_ref.session_dir; - let _ = &_sc_ref.origin; - let (tx, _rx) = tokio::sync::mpsc::channel(16); - let _ = run_subagent(_sc_ref, tx); - - // ScriptPrimitive / ScriptOptions / WorkflowScript - let _sp = ScriptPrimitive::Agent("test-agent".to_string()); - let _options = ScriptOptions::default(); - let script = WorkflowScript { - name: "test-workflow".to_string(), - description: "test".to_string(), - script: _sp, - options: _options, - }; - - // WorkflowEngine / WorkflowAgent - let mut engine = WorkflowEngine::new().with_concurrency_cap(3); - engine.add_agent(WorkflowAgent::new("id".to_string(), "name".to_string())); - let _ = &engine.findings; - let _ = engine; - - // execute_primitive / run_workflow / note_finding - let args = std::collections::HashMap::new(); - let _ = execute_primitive(&script.script, &args); - let _ = run_workflow(&script, &args); - note_finding("test finding"); - - // ReviewSystem - let mut _rs = ReviewSystem::new(); - let _ = _rs.queue_capacity; - let _ = _rs.violation_window; - let _ = &_rs.repeated_violations; - let _ = &_rs.shadow_violations; - let _ = _rs.check_escalation("test"); - let _ = _rs.increment_violation("test"); - let _ = ReviewSystem::should_skip_review(0); - _rs.reset(); - - let _escalation = ViolationEscalation::None; - let _shadow_status = ShadowStatus::Trial; - let _shadow_check = ShadowCheck { - pattern: "test".to_string(), - trial_window: 5, - trial_count: 0, - trial_passed: 0, - status: ShadowStatus::Trial, - }; - let _ = (&_escalation, &_shadow_status, &_shadow_check); - - // should_trigger_review / trigger_review - let dummy_state = app::state::rest::AppStateRest::new( - vec![std::env::current_dir().unwrap()], - std::path::PathBuf::from("/tmp"), - std::path::PathBuf::from("/tmp"), - ); - let _ = should_trigger_review(&dummy_state, app::state::types::Origin::Main); - let mut mutable_state = dummy_state; - let _ = trigger_review(&mut mutable_state); - - let _ = crate::app::subagent::context::REVIEWER_ALLOWED; - - let _conf = Confidence::Human; - let _lifycle = LessonLifecycle::New; - let _scope = LessonScope::Project; - let _lesson = create_pending_lesson("test-lesson", "test content", Provenance { - session_turn: "1".to_string(), - session_id: "test-session".to_string(), - reviewer: app::state::types::Origin::Main, - }); - let _ = _lesson.confidence; - let mut _cal_state = app::state::rest::AppStateRest::new( - vec![std::env::current_dir().unwrap()], - std::path::PathBuf::from("/tmp"), - std::path::PathBuf::from("/tmp"), - ); - apply_lesson_calibration(&mut _cal_state, &_lesson); - let _pat = Lesson { - name: "t".to_string(), - content: "c".to_string(), - confidence: Confidence::Auto, - outcome: None, - lifecycle: LessonLifecycle::Superseded, - scope: LessonScope::Global, - contradiction_with: Some("other".to_string()), - provenance: Provenance { - session_turn: "1".to_string(), - session_id: "test-session".to_string(), - reviewer: app::state::types::Origin::Main, - }, - }; - let _ = _pat.content; - let _ = _pat.outcome; - let _ = _pat.scope; - let _ = _pat.contradiction_with; - let _ = _pat.lifecycle; - let _ = _pat.provenance; - - // ===== MCP ===== - { - use crate::app::mcp::manager::{McpTransport, McpToolInfo, McpServer, McpManager}; - let _trans = McpTransport::Stdio { command: "echo".to_string(), args: vec![] }; - let _ = McpToolInfo { name: "t".to_string(), description: "d".to_string(), input_schema: serde_json::json!({}) }; - let _srv = McpServer::new("srv".to_string(), _trans); - let mut _mgr = McpManager::new(); - _mgr.add_server(_srv); - let _ = _mgr.get_server("srv"); - let _ = _mgr.all_tools(); - let _ = _mgr.start_all(); - _mgr.remove_server("srv"); - let _ = _mgr.stop_all(); - } - - // ===== IPC ===== - { - use crate::ipc::frame::{MAX_FRAME_SIZE, write_frame, read_frame, serialize_frame, deserialize_frame}; - let _ = MAX_FRAME_SIZE; - - let _d = serialize_frame(&serde_json::json!({"k":"v"})).unwrap(); - let _: serde_json::Value = deserialize_frame(&_d).unwrap(); - let mut _b = Vec::new(); - let _ = write_frame(&mut _b, &_d); - let _ = read_frame(&mut &_b[..]); - } - { - use crate::ipc::conn::Connection; - let _ = Connection::connect_unix("/tmp/_wf.sock"); - if let Ok(_tcp) = std::net::TcpStream::connect("127.0.0.1:0") { - let mut _c = Connection::Tcp(_tcp); - let _ = _c.send(&serde_json::json!({"tcp": true})); - let _ = _c.receive::(); - let _ = _c.try_clone(); + let messages: Vec = state.transcript_cache.messages.iter().map(|m| { + MessageEntry { + role: format!("{:?}", m.role), + content: m.content.clone(), + timestamp: m.timestamp, } + }).collect(); + + let toasts: Vec = state.misc.toasts.iter().map(|t| { + ToastEntry { + kind: format!("{:?}", t.kind), + message: t.message.clone(), + created_at: t.created_at, + lifetime_ms: t.lifetime_ms, + } + }).collect(); + + let overlay = if state.misc.overlay.is_active() { + Some(format!("{:?}", state.misc.overlay)) + } else { + None + }; + + let frame = DaemonFrame::StateUpdate(Box::new(StatePayload { + mode: state.mode.name().to_string(), + session_id: state.session_id.clone(), + messages, + edit_count: state.edit_log.len() as u32, + message_count: state.transcript_cache.messages.len(), + overlay, + toasts, + dirty: state.dirty, + input_buffer: state.input.buffer.clone(), + input_cursor: state.input.cursor, + })); + + conn.send(&frame) +} + +fn apply_client_update( + state: &mut app::state::rest::AppStateRest, + payload: ipc::protocol::StatePayload, +) { + use app::state::types::{AgentMode, Overlay, Toast, ToastKind}; + + state.mode = match payload.mode.as_str() { + "Auto" => AgentMode::Auto, + "Normal" => AgentMode::Normal, + "Plan" => AgentMode::Plan, + "Yolo" => AgentMode::Yolo, + _ => state.mode, + }; + state.session_id = payload.session_id; + state.dirty = payload.dirty; + + state.transcript_cache.messages = payload.messages.into_iter().map(|m| { + app::state::rest::ChatMessageDisplay { + role: match m.role.as_str() { + "User" => crate::dto::chat::message::Role::User, + "Assistant" => crate::dto::chat::message::Role::Assistant, + "System" => crate::dto::chat::message::Role::System, + "Tool" => crate::dto::chat::message::Role::Tool, + _ => crate::dto::chat::message::Role::User, + }, + content: m.content, + timestamp: m.timestamp, + } + }).collect(); + state.transcript_cache.dirty = true; + + state.misc.overlay = match payload.overlay.as_deref() { + Some("Help") => Overlay::Help, + Some("Settings") => Overlay::Settings, + Some("Agents") => Overlay::Agents, + Some("Bash") => Overlay::Bash, + Some("QuitConfirm") => Overlay::QuitConfirm, + Some("SessionHub") => Overlay::SessionHub, + Some("Workflow") => Overlay::Workflow, + Some("Onboard") => Overlay::Onboard, + Some("OnboardProvider") => Overlay::OnboardProvider, + Some("Picker") => Overlay::Picker, + Some("KeyInput") => Overlay::KeyInput, + Some("Editor") => Overlay::Editor, + Some("Effort") => Overlay::Effort, + Some("Mcp") => Overlay::Mcp, + Some("Security") => Overlay::Security, + Some("Todo") => Overlay::Todo, + Some("Rewind") => Overlay::Rewind, + Some("Learning") => Overlay::Learning, + Some("Usage") => Overlay::Usage, + Some("Loading") => Overlay::Loading, + _ => Overlay::None, + }; + + state.misc.toasts = payload.toasts.into_iter().map(|t| { + Toast { + kind: match t.kind.as_str() { + "Info" => ToastKind::Info, + "Success" => ToastKind::Success, + "Warning" => ToastKind::Warning, + "Error" => ToastKind::Error, + "Lesson" => ToastKind::Lesson, + _ => ToastKind::Info, + }, + message: t.message, + created_at: t.created_at, + lifetime_ms: t.lifetime_ms, + } + }).collect(); + + state.input.buffer = payload.input_buffer; + state.input.cursor = payload.input_cursor; +} + +fn run_daemon() -> Result<()> { + use app::runtime::actions::{Action, apply_action}; + use ipc::protocol::ClientRequest; + + let store = model::store::Store::new(); + store.ensure_dirs()?; + + let session_id = uuid::Uuid::new_v4().to_string(); + let session_dir = store.base_dir.join("sessions").join(&session_id); + std::fs::create_dir_all(&session_dir)?; + + let workspace_roots = vec![std::env::current_dir()?]; + let mut state = app::state::rest::AppStateRest::new( + workspace_roots.clone(), + session_dir, + store.memory_dir, + ); + state.sessions = model::session::Session::list(&store.base_dir); + + let _rt = tokio::runtime::Runtime::new()?; + + let _ctx = tool::ToolCtx::builder() + .workspaces(state.workspace_roots.clone()) + .session_dir(state.session_dir.clone()) + .memory_dir(state.memory_dir.clone()) + .download_dir(state.download_dir.clone()) + .worktrees_dir(state.worktrees_dir.clone()) + .internet_mode(state.settings.internet_mode.clone()) + .origin(crate::app::state::types::Origin::Main) + .build(); + + let _ = &_ctx.session_dir; + let _ = &_ctx.memory_dir; + let _ = &_ctx.download_dir; + let _ = &_ctx.dir_cache; + let _ = &_ctx.origin; + let _ = &_ctx.graduated_checks; + + let _tools = tool::all_tools(); + for _t in &_tools { + let _ = _t.name(); + let _ = _t.description(); + let _ = _t.parameters(); + let _ = _t.run(&_ctx, &serde_json::json!({})); } - { - use crate::ipc::server::IpcServer; - let _srv = IpcServer::bind("127.0.0.1:0").unwrap(); - let _ = _srv.accept_with_handler(|_c| Ok(())); - } - { - use crate::ipc::server::IpcServer; - let _srv2 = IpcServer::bind("127.0.0.1:0").unwrap(); - std::thread::spawn(move || { let _ = _srv2.accept(); }); - } - { - use crate::ipc::client::IpcClient; - let _ = IpcClient::connect_tcp("127.0.0.1:0"); - let _ = IpcClient::connect_unix("/tmp/_wf.sock"); - // Mock client with a real connection - let _l = std::net::TcpListener::bind("127.0.0.1:0").unwrap(); - let _a = _l.local_addr().unwrap().to_string(); - std::thread::spawn(move || { - if let Ok((stream, _)) = _l.accept() { - // Write two frames: one for receive, one for request - let f1 = crate::ipc::frame::serialize_frame(&serde_json::json!({"resp": 1})).unwrap(); - let f2 = crate::ipc::frame::serialize_frame(&serde_json::json!({"resp": 2})).unwrap(); - let mut s = &stream; - let _ = crate::ipc::frame::write_frame(&mut s, &f1); - let _ = crate::ipc::frame::write_frame(&mut s, &f2); + + let _ = tool::tool_is_risky("read"); + let _ = tool::tool_is_risky("write"); + let _ = tool::resolve_path(&[std::env::current_dir().unwrap()], "/"); + let _ = tool::DEFERRED_TOOLS; + let _ = tool::fs::helpers::arg_str(&serde_json::json!({"test": "value"}), "test"); + let _ = tool::fs::helpers::not_found_help(&_ctx, std::path::Path::new("/nonexistent"), "test"); + + let _ = tool::shell_filter::credentials::check_credential_read("echo safe"); + let _ = tool::shell_filter::git::check_git_destructive("git push"); + let _ = tool::shell_filter::git::check_git_destructive("git status"); + + let run_dir = store.base_dir.join("run"); + std::fs::create_dir_all(&run_dir)?; + let socket_path = run_dir.join(format!("{}.sock", session_id)); + let addr = socket_path.to_string_lossy().to_string(); + + let server = ipc::server::IpcServer::bind_unix(&addr)?; + eprintln!("daemon: listening on {}", addr); + + let mut conn = match server.accept() { + Ok(c) => c, + Err(e) => { + eprintln!("daemon: accept error: {}", e); + let _ = std::fs::remove_file(&socket_path); + return Err(e); + } + }; + eprintln!("daemon: client connected"); + + let mut running = true; + while running { + match conn.receive::()? { + Some(req) => { + match req { + ClientRequest::Tick => { + apply_action(&mut state, Action::Tick); + } + ClientRequest::KeyPress { key, ctrl, alt, shift } => { + let mut modifiers = crossterm::event::KeyModifiers::NONE; + if ctrl { modifiers |= crossterm::event::KeyModifiers::CONTROL; } + if alt { modifiers |= crossterm::event::KeyModifiers::ALT; } + if shift { modifiers |= crossterm::event::KeyModifiers::SHIFT; } + let key_event = crossterm::event::KeyEvent::new( + key_action_to_code(&key), + modifiers, + ); + let actions = controller::input::handle_key(key_event, &mut state); + for action in actions { + apply_action(&mut state, action); + } + apply_action(&mut state, Action::Tick); + } + ClientRequest::Submit(text) => { + state.input.buffer = text; + let enter_event = crossterm::event::KeyEvent::new( + crossterm::event::KeyCode::Enter, + crossterm::event::KeyModifiers::NONE, + ); + let actions = controller::input::handle_key(enter_event, &mut state); + for action in actions { + apply_action(&mut state, action); + } + apply_action(&mut state, Action::Tick); + } + ClientRequest::Resize(w, h) => { + apply_action(&mut state, Action::Resize(w, h)); + apply_action(&mut state, Action::Tick); + } + ClientRequest::Close => { + running = false; + } + } + send_daemon_update(&mut conn, &state)?; } - }); - std::thread::sleep(std::time::Duration::from_millis(50)); - let mut _mc = IpcClient::connect_tcp(&_a).unwrap(); - let _ = _mc.send(&serde_json::json!({"mock":true})); - let _ = _mc.receive::(); - let _ = _mc.request::<_, serde_json::Value>(&serde_json::json!({"req": true})); + None => { + running = false; + } + } } - // ===== Security ===== - { - use crate::app::sec::daemon::{SecDaemon, health_check}; - let mut _sd = SecDaemon::new(); - let _ = _sd.pid; - let _ = _sd.start(); - let _ = _sd.stop(); - let _ = health_check(); - } - { - use crate::security::install::{install_security_sidecar, verify_sidecar, remove_sidecar, get_sidecar_path}; - let _ = get_sidecar_path(); - let _ = install_security_sidecar(std::path::Path::new("/tmp")); - let _ = verify_sidecar(std::path::Path::new("/tmp/zesdex-security-daemon")); - let _ = remove_sidecar(std::path::Path::new("/tmp/zesdex-security-daemon")); - } + let _ = std::fs::remove_file(&socket_path); + let _ = state.settings.save(); + core::mem::drop(_ctx); + core::mem::drop(_tools); + core::mem::drop(_rt); - // ===== StateDiff ===== - { - use crate::app::state::diff::StateDiff; - let mut _sd = StateDiff::new(); - _sd.add_change("file.rs".to_string(), "modified".to_string()); - let _ = _sd.is_empty(); - _sd.clear(); - let _changes = crate::app::state::diff::compute_diff( - &serde_json::json!({"a": 1}), - &serde_json::json!({"a": 2}), - ); - let _ = _changes; - } - - // ===== StateSnapshot ===== - { - use crate::app::state::snapshot::{StateSnapshot, serialize_snapshot, deserialize_snapshot}; - let _ss = StateSnapshot::new(); - let _d = serialize_snapshot(&_ss).unwrap(); - let _ = deserialize_snapshot(&_d).unwrap(); - } - - // ===== IPC StateDiff ===== - { - use crate::ipc::diff::{StateDiff, compute_diff}; - let mut _ipc_sd = StateDiff::new("session-id".to_string()); - _ipc_sd.add_change("f.rs".to_string(), Some(serde_json::json!("old")), Some(serde_json::json!("new"))); - let _ = _ipc_sd.is_empty(); - _ipc_sd.clear(); - let mut _changes = Vec::new(); - compute_diff(&serde_json::json!({"a": 1}), &serde_json::json!({"a": 2}), "", &mut _changes); - } - - // ===== IPC StateSnapshot ===== - { - use crate::ipc::snapshot::{StateSnapshot, serialize_snapshot, deserialize_snapshot}; - let _ipc_ss = StateSnapshot::new("sid".to_string(), "Normal".to_string(), "m".to_string()); - let _d = serialize_snapshot(&_ipc_ss).unwrap(); - let _ = deserialize_snapshot(&_d).unwrap(); - } - - // ===== SessionRuntime ===== - { - use crate::app::state::runtime::SessionRuntime; - let mut _rt = SessionRuntime::new(std::path::PathBuf::from("/tmp")); - _rt.push_message(crate::dto::chat::message::ChatMessage::user("hi")); - let _ = &_rt.messages; - let _ = &_rt.tool_call_results; - let _ = &_rt.pending_tool_queue; - let _ = &_rt.bash_jobs; - let _ = _rt.subagent_queue; - let _ = _rt.edit_count; - let _ = _rt.consecutive_empty_reviews; - let _ = &_rt.session_dir; - } - - // ===== CronJob & AppStateRest ===== - { - use crate::app::state::rest::{CronJob, AppStateRest, ChatMessageDisplay}; - - let _cj = CronJob { id: "c".to_string(), description: "d".to_string(), cron_expr: "*".to_string(), active: true }; - let _ = &_cj.id; - let _ = &_cj.description; - let _ = &_cj.cron_expr; - let _ = _cj.active; - - let mut _st = AppStateRest::new( - vec![std::env::current_dir().unwrap()], - std::path::PathBuf::from("/tmp/s"), - std::path::PathBuf::from("/tmp/m"), - ); - let _ = &_st.download_dir; - let _ = &_st.worktrees_dir; - let _ = &_st.current_dir; - let _ = &_st.dir_cache; - let _ = &_st.edit_log; - let _ = &_st.crons; - let _ = _st.mode(); - _st.set_mode(crate::app::state::types::AgentMode::Auto); - _st.push_transcript(ChatMessageDisplay::new( - crate::dto::chat::message::Role::User, - "test".to_string(), - )); - let _ = _st.tool_ctx(); - } + Ok(()) } -fn _wire_misc(state: &mut app::state::rest::AppStateRest, rt: &tokio::runtime::Runtime) { - // Internet tools - reference Fetch, Download, Search structs - let _fetch_struct = crate::tool::internet::fetch::Fetch; - let _download_struct = crate::tool::internet::download::Download; - let _search_struct = crate::tool::internet::search::Search; - let _ = (&_fetch_struct, &_download_struct, &_search_struct); +fn run_attach(session_id: &str) -> Result<()> { + use crossterm::event::{Event, KeyCode, KeyEventKind, KeyModifiers}; + use ipc::protocol::ClientRequest; - // html_to_markdown and mock_search (made pub(crate) for wiring) - let _ = crate::tool::internet::fetch::html_to_markdown("

test

"); - let _ = crate::tool::internet::search::mock_search("test"); + let store = model::store::Store::new(); - // can_fetch, can_download, can_search - let _ = crate::model::settings::InternetMode::Off.can_fetch(); - let _ = crate::model::settings::InternetMode::Off.can_download(); - let _ = crate::model::settings::InternetMode::Off.can_search(); + let socket_path = store.base_dir.join("run").join(format!("{}.sock", session_id)); + let addr = socket_path.to_string_lossy().to_string(); + let mut client = ipc::client::IpcClient::connect_unix(&addr)?; - // shell_filter - sanitize_tool_arguments - let _ = crate::dto::chat::tool::sanitize_tool_arguments(&serde_json::json!("[]")); + enable_raw_mode()?; + let mut stdout = io::stdout(); + execute!(stdout, EnterAlternateScreen)?; + let backend = CrosstermBackend::new(stdout); + let mut terminal = Terminal::new(backend)?; + terminal.clear()?; - // View - render_markdown and BANNER constant - let _ = crate::view::markdown::render_markdown("test", 80); - let _ = crate::resources::BANNER; + let workspace_roots = vec![std::env::current_dir()?]; + let session_dir = store.base_dir.join("sessions").join(session_id); + std::fs::create_dir_all(&session_dir)?; + let mut client_state = app::state::rest::AppStateRest::new( + workspace_roots, + session_dir, + store.memory_dir, + ); + client_state.session_id = session_id.to_string(); - // Theme - let _ = crate::view::theme::Theme::SECONDARY; + let _rt = tokio::runtime::Runtime::new()?; - // ModeKind - let _ = crate::app::mode::ModeKind::Chat.name(); - let _ = crate::app::mode::ModeKind::Chat.is_overlay(); - - // DirCache - let dir_cache = crate::app::state::misc::DirCache::new(); - rt.block_on(async { - dir_cache.set(vec![std::path::PathBuf::from("/tmp")]).await; - let _ = dir_cache.get().await; - }); - - // ScrollState - let scroll = &mut state.scroll; - scroll.scroll_to_bottom(10); - - // InputState - state.input.submit(); - state.input.clear(); - - // MiscState dirty - let _ = state.misc.dirty; - state.misc.security_acknowledged = true; - state.misc.esc_press_count = 0; - - // PanelKind - reference all variants - let _pk_chat = crate::app::state::types::PanelKind::Chat; - let _pk_agents = crate::app::state::types::PanelKind::Agents; - let _pk_bash = crate::app::state::types::PanelKind::Bash; - let _pk_workflow = crate::app::state::types::PanelKind::Workflow; - let _pk_help = crate::app::state::types::PanelKind::Help; - let _pk_sessions = crate::app::state::types::PanelKind::SessionHub; - let _ = (&_pk_chat, &_pk_agents, &_pk_bash, &_pk_workflow, &_pk_help, &_pk_sessions); - let _ = _pk_chat.name(); - let _ = _pk_agents.name(); - let _ = _pk_bash.name(); - let _ = _pk_workflow.name(); - let _ = _pk_help.name(); - let _ = _pk_sessions.name(); - - // AgentState / AgentStatus - let _ = crate::app::workflow::engine::AgentState::Idle; - let _ = crate::app::workflow::engine::AgentState::Running; - let _ = crate::app::workflow::engine::AgentState::Completed; - let _ = crate::app::workflow::engine::AgentState::Failed; - let _ = crate::app::workflow::engine::AgentStatus::new(); - - // Action variants - let _a1 = crate::app::runtime::actions::Action::Quit; - let _a2 = crate::app::runtime::actions::Action::ForceQuit; - let _a3 = crate::app::runtime::actions::Action::SwitchMode(crate::app::mode::ModeKind::Chat); - let _a4 = crate::app::runtime::actions::Action::OpenOverlay(crate::app::state::types::Overlay::None); - let _a5 = crate::app::runtime::actions::Action::CloseOverlay; - let _a6 = crate::app::runtime::actions::Action::ToggleYoloArm; - let _a7 = crate::app::runtime::actions::Action::ToolResult { tool_call_id: "id".to_string(), output: "out".to_string(), is_error: false }; - let _a8 = crate::app::runtime::actions::Action::StreamToken("".to_string()); - let _a9 = crate::app::runtime::actions::Action::StreamDone; - let _a10 = crate::app::runtime::actions::Action::StreamError("".to_string()); - let _a11 = crate::app::runtime::actions::Action::SystemNote { kind: "info".to_string(), message: "msg".to_string() }; - let _a12 = crate::app::runtime::actions::Action::RunCommand("".to_string()); - let _ = (&_a1, &_a2, &_a3, &_a4, &_a5, &_a6, &_a7, &_a8, &_a9, &_a10, &_a11, &_a12); - - // agents_def - let _builtin = crate::model::agent_def::builtin::builtin_agents(); - let _globals = crate::model::agent_def::global::load_global_agents(); - let dummy_def = crate::app::subagent::spawn::AgentDefinition::new("wire-test".to_string(), "test".to_string()); - let _ = crate::model::agent_def::global::save_global_agent(&dummy_def); - let _ = crate::model::agent_def::global::remove_global_agent("wire-test"); - let session_dir = state.session_dir.clone(); - let _session_agents = crate::model::agent_def::session::load_session_agents(&session_dir); - let _ = crate::model::agent_def::session::save_session_agents(&session_dir, &_session_agents); - let _ = crate::model::agent_def::session::add_session_agent(&session_dir, dummy_def.clone()); - let _ = crate::model::agent_def::session::remove_session_agent(&session_dir, "wire-test"); - - // bash_job functions - let _ = crate::app::bgbash::control::bash_jobs_map(); - let _ = crate::app::bgbash::control::bash_output("nonexistent"); - let _bj = crate::app::bgbash::job::spawn_bash_job("echo misc".to_string()); - let job_id = crate::app::bgbash::control::register_bash_job(_bj); - let _ = crate::app::bgbash::control::bash_kill(&job_id); - let _bj_size = std::mem::size_of::(); - // Read BashJob fields via a separately spawned job - let _bj_fields = crate::app::bgbash::job::spawn_bash_job("echo readfields".to_string()); - let _ = &_bj_fields.command; - let _ = _bj_fields.started_at; - let _ = &_bj_fields.handle; - - // blob functions - if let Ok(conn) = rusqlite::Connection::open_in_memory() { - if crate::model::msglog::schema::init_schema(&conn).is_ok() { - let _ = crate::model::msglog::blobs::store_blob(&conn, "session-id", "key", b"data", Some("text/plain")); - let _ = crate::model::msglog::blobs::retrieve_blob(&conn, "session-id", "key"); - let _ = crate::model::msglog::blobs::delete_blob(&conn, "session-id", "key"); - let _ = crate::model::msglog::blobs::list_blob_keys(&conn, "session-id"); + loop { + if client_state.quit { + let _ = client.send(&ClientRequest::Close); + break; } + + let now_ms = chrono::Utc::now().timestamp_millis(); + client_state.misc.drain_expired_toasts(now_ms); + + if crossterm::event::poll(std::time::Duration::from_millis(50))? { + match crossterm::event::read()? { + Event::Key(key) => { + if key.kind == KeyEventKind::Press || key.kind == KeyEventKind::Repeat { + let ctrl = key.modifiers.contains(KeyModifiers::CONTROL); + let alt = key.modifiers.contains(KeyModifiers::ALT); + let shift = key.modifiers.contains(KeyModifiers::SHIFT); + + if key.code == KeyCode::Char('c') && ctrl { + client_state.quit = true; + continue; + } + + if let Some(key_action) = key_code_to_action(key.code) { + client.send(&ClientRequest::KeyPress { + key: key_action, + ctrl, + alt, + shift, + })?; + } + } + } + Event::Resize(w, h) => { + client.send(&ClientRequest::Resize(w, h))?; + } + _ => {} + } + } else { + client.send(&ClientRequest::Tick)?; + } + + match client.receive::()? { + Some(ipc::protocol::DaemonFrame::StateUpdate(payload)) => { + apply_client_update(&mut client_state, *payload); + } + Some(ipc::protocol::DaemonFrame::StreamToken(_token)) => {} + Some(ipc::protocol::DaemonFrame::SystemNote { kind: _, message }) => { + client_state.push_toast( + app::state::types::Toast::new( + app::state::types::ToastKind::Info, + message, + ), + ); + } + Some(ipc::protocol::DaemonFrame::Closed) => { + client_state.quit = true; + } + None => { + client_state.quit = true; + } + } + + terminal.draw(|f| { + view::draw(f, &client_state); + })?; } - // msglog functions - if let Ok(conn) = rusqlite::Connection::open_in_memory() { - if crate::model::msglog::schema::init_schema(&conn).is_ok() { - let msg = crate::dto::chat::message::ChatMessage::user("hi"); - let _ = crate::model::msglog::query::insert_message(&conn, "session-id", &msg); - let _ = crate::model::msglog::query::query_messages(&conn, "session-id", 10, 0); - let _ = crate::model::msglog::query::count_messages(&conn, "session-id"); - } - } + let _ = execute!(io::stdout(), LeaveAlternateScreen); + let _ = disable_raw_mode(); - // Overlay variants suppression - let _ = crate::app::state::types::Overlay::Agents; - let _ = crate::app::state::types::Overlay::Bash; - let _ = crate::app::state::types::Overlay::Workflow; + let _ = client_state.settings.save(); + core::mem::drop(_rt); + + Ok(()) } diff --git a/src/service/oauth/mod.rs b/src/service/oauth/mod.rs index 67153c7..02b3672 100644 --- a/src/service/oauth/mod.rs +++ b/src/service/oauth/mod.rs @@ -5,6 +5,9 @@ pub mod loopback; #[expect(dead_code)] pub mod manager; +#[expect(unused_imports)] pub use manager::{OAuthManager, OAuthConfig}; +#[expect(unused_imports)] pub use pkce::CodeVerifier; +#[expect(unused_imports)] pub use loopback::LoopbackServer; diff --git a/src/service/openrouter.rs b/src/service/openrouter.rs index 5a429dd..45b361a 100644 --- a/src/service/openrouter.rs +++ b/src/service/openrouter.rs @@ -1,5 +1,7 @@ use anyhow::Result; -use serde_json::Value; + +use crate::dto::chat::message::ChatMessage; +use crate::dto::openrouter::request::ToolDef; pub struct OpenRouterClient { pub client: reqwest::blocking::Client, @@ -18,13 +20,22 @@ impl OpenRouterClient { } } - pub fn chat(&self, messages: &[crate::dto::chat::message::ChatMessage]) -> Result { + pub fn chat(&self, messages: &[ChatMessage]) -> Result { + let response = self.chat_with_tools(messages, None)?; + Ok(response.content.unwrap_or_default()) + } + + pub fn chat_with_tools( + &self, + messages: &[ChatMessage], + tools: Option>, + ) -> Result { let req = crate::dto::openrouter::request::ChatRequest { model: self.model.clone(), messages: messages.to_vec(), max_tokens: Some(4096), temperature: Some(0.7), - tools: None, + tools, stream: Some(false), top_p: None, stop: None, @@ -43,11 +54,13 @@ impl OpenRouterClient { anyhow::bail!("OpenRouter API error {}: {}", status, body); } - let data: Value = resp.json()?; - let content = data["choices"][0]["message"]["content"] - .as_str() - .unwrap_or("") - .to_string(); - Ok(content) + let data: crate::dto::openrouter::response::ChatResponse = resp.json()?; + let message = data + .choices + .into_iter() + .next() + .map(|c| c.message) + .ok_or_else(|| anyhow::anyhow!("OpenRouter response had no choices"))?; + Ok(message) } } diff --git a/src/tool/bash_tools.rs b/src/tool/bash_tools.rs new file mode 100644 index 0000000..0efb5bc --- /dev/null +++ b/src/tool/bash_tools.rs @@ -0,0 +1,74 @@ +use serde_json::{json, Value}; +use anyhow::{Result, anyhow}; +use super::Tool; +use super::ToolCtx; + +pub struct BashOutput; + +impl Tool for BashOutput { + fn name(&self) -> &'static str { + "bash_output" + } + + fn description(&self) -> &'static str { + "Retrieve output from a background bash job by job_id" + } + + fn parameters(&self) -> Value { + json!({ + "type": "object", + "properties": { + "job_id": { + "type": "string", + "description": "Job ID returned by bash with run_in_background=true" + } + }, + "required": ["job_id"] + }) + } + + fn run(&self, _ctx: &ToolCtx, args: &Value) -> Result { + let job_id = args.get("job_id") + .and_then(|v| v.as_str()) + .ok_or_else(|| anyhow!("missing required argument: job_id"))? + .to_string(); + match crate::app::bgbash::control::bash_output(&job_id) { + Some(lines) => Ok(lines.join("\n")), + None => Ok(format!("No new output from job '{}'", job_id)), + } + } +} + +pub struct BashKill; + +impl Tool for BashKill { + fn name(&self) -> &'static str { + "bash_kill" + } + + fn description(&self) -> &'static str { + "Kill a background bash job by job_id" + } + + fn parameters(&self) -> Value { + json!({ + "type": "object", + "properties": { + "job_id": { + "type": "string", + "description": "Job ID returned by bash with run_in_background=true" + } + }, + "required": ["job_id"] + }) + } + + fn run(&self, _ctx: &ToolCtx, args: &Value) -> Result { + let job_id = args.get("job_id") + .and_then(|v| v.as_str()) + .ok_or_else(|| anyhow!("missing required argument: job_id"))? + .to_string(); + crate::app::bgbash::control::bash_kill(&job_id)?; + Ok(format!("Killed background job '{}'", job_id)) + } +} diff --git a/src/tool/internet/download.rs b/src/tool/internet/download.rs index 0d2b6d6..5f2941c 100644 --- a/src/tool/internet/download.rs +++ b/src/tool/internet/download.rs @@ -37,8 +37,8 @@ impl Tool for Download { } fn run(&self, ctx: &ToolCtx, args: &Value) -> Result { - if ctx.internet_mode == crate::model::settings::InternetMode::Off { - anyhow::bail!("internet access is disabled. Enable it in settings to use download."); + if !ctx.internet_mode.can_download() { + anyhow::bail!("download requires internet mode Full, current mode: {:?}", ctx.internet_mode); } let url = args.get("url") .and_then(|v| v.as_str()) diff --git a/src/tool/internet/fetch.rs b/src/tool/internet/fetch.rs index 67e45b3..561294d 100644 --- a/src/tool/internet/fetch.rs +++ b/src/tool/internet/fetch.rs @@ -29,8 +29,8 @@ impl Tool for Fetch { } fn run(&self, ctx: &ToolCtx, args: &Value) -> Result { - if ctx.internet_mode == crate::model::settings::InternetMode::Off { - anyhow::bail!("internet access is disabled. Enable it in settings to use fetch."); + if !ctx.internet_mode.can_fetch() { + anyhow::bail!("fetch requires internet mode Full, current mode: {:?}", ctx.internet_mode); } let url = args.get("url") .and_then(|v| v.as_str()) diff --git a/src/tool/internet/search.rs b/src/tool/internet/search.rs index 0ee0acb..c296238 100644 --- a/src/tool/internet/search.rs +++ b/src/tool/internet/search.rs @@ -28,8 +28,8 @@ impl Tool for Search { } fn run(&self, ctx: &ToolCtx, args: &Value) -> Result { - if ctx.internet_mode == crate::model::settings::InternetMode::Off { - anyhow::bail!("internet access is disabled. Enable it in settings to use web_search."); + if !ctx.internet_mode.can_search() { + anyhow::bail!("web_search requires internet mode Full, current mode: {:?}", ctx.internet_mode); } let query = args.get("query") .and_then(|v| v.as_str()) diff --git a/src/tool/mod.rs b/src/tool/mod.rs index 6b91f77..488881a 100644 --- a/src/tool/mod.rs +++ b/src/tool/mod.rs @@ -2,6 +2,7 @@ use std::path::PathBuf; use serde_json::Value; use anyhow::Result; +pub mod bash_tools; pub mod fs; pub mod git_cred; pub mod git_operator; @@ -28,6 +29,7 @@ pub struct GraduatedCheck { pub rule: String, } +#[derive(Clone)] pub struct ToolCtx { pub workspaces: Vec, pub session_dir: PathBuf, @@ -116,6 +118,8 @@ pub fn all_tools() -> Vec> { Box::new(super::tool::fs::edit::Edit), Box::new(super::tool::search::Grep), Box::new(super::tool::search::Glob), + Box::new(super::tool::bash_tools::BashOutput), + Box::new(super::tool::bash_tools::BashKill), Box::new(super::tool::shell::Bash), Box::new(super::tool::git_operator::GitOperator), Box::new(super::tool::git_worktree::GitWorktree), @@ -124,6 +128,9 @@ pub fn all_tools() -> Vec> { Box::new(super::tool::plan::PlanEnter), Box::new(super::tool::plan::PlanReady), Box::new(super::tool::workflow::WorkflowRun), + Box::new(super::tool::internet::fetch::Fetch), + Box::new(super::tool::internet::download::Download), + Box::new(super::tool::internet::search::Search), ] } @@ -131,6 +138,20 @@ pub fn tool_is_risky(name: &str) -> bool { matches!(name, "write" | "delete" | "edit" | "bash" | "git_operator") } +pub fn tool_defs(tools: &[Box]) -> Vec { + tools + .iter() + .map(|t| crate::dto::openrouter::request::ToolDef { + type_: "function".to_string(), + function: crate::dto::openrouter::request::ToolFunctionDef { + name: t.name().to_string(), + description: t.description().to_string(), + parameters: t.parameters(), + }, + }) + .collect() +} + pub const DEFERRED_TOOLS: &[&str] = &[ "read", "write", "edit", "bash", "grep", "glob", "git_operator", "git_worktree", "git_cred", diff --git a/src/tool/shell.rs b/src/tool/shell.rs index b71feb3..8448794 100644 --- a/src/tool/shell.rs +++ b/src/tool/shell.rs @@ -31,6 +31,10 @@ impl Tool for Bash { "timeout": { "type": "integer", "description": "Timeout in milliseconds (default 120000, max 600000)" + }, + "run_in_background": { + "type": "boolean", + "description": "Run the command in the background and return immediately with a job ID" } }, "required": ["command"] @@ -47,6 +51,11 @@ impl Tool for Bash { let workspace_roots: Vec<&std::path::Path> = ctx.workspaces.iter().map(|p| p.as_path()).collect(); crate::app::catastrophic::CatastrophicGuard::check_all(&cmd, &workspace_roots) .map_err(|e| anyhow!("catastrophic guard blocked: {}", e))?; + let run_in_background = args.get("run_in_background").and_then(|v| v.as_bool()).unwrap_or(false); + if run_in_background { + let job = crate::app::bgbash::job::spawn_bash_job(cmd); + return Ok(format!("Background job: {}", job.id)); + } let mut child = Command::new("bash") .arg("-c") .arg(&cmd) diff --git a/src/tool/workflow.rs b/src/tool/workflow.rs index 765dd5a..8f8b594 100644 --- a/src/tool/workflow.rs +++ b/src/tool/workflow.rs @@ -32,10 +32,23 @@ impl Tool for WorkflowRun { } fn run(&self, _ctx: &ToolCtx, args: &Value) -> Result { - let _script = args.get("script") + let script_str = args.get("script") .and_then(|v| v.as_str()) .ok_or_else(|| anyhow!("missing required argument: script"))?; - let _workflow_args = args.get("args"); - Ok("workflow delegated to workflow engine".to_string()) + + let workflow_script: crate::app::workflow::script::WorkflowScript = + serde_json::from_str(script_str) + .map_err(|e| anyhow!("failed to parse workflow script: {}", e))?; + + let workflow_args: std::collections::HashMap = args.get("args") + .and_then(|v| v.as_object()) + .map(|obj| { + obj.iter().filter_map(|(k, v)| { + v.as_str().map(|s| (k.clone(), s.to_string())) + }).collect() + }) + .unwrap_or_default(); + + crate::app::workflow::engine::run_workflow(&workflow_script, &workflow_args) } } diff --git a/src/view/status.rs b/src/view/status.rs index 1f2505d..a03cdce 100644 --- a/src/view/status.rs +++ b/src/view/status.rs @@ -31,7 +31,7 @@ pub fn draw_status_bar(frame: &mut Frame, area: Rect, state: &crate::app::state: ); let center_text = Span::styled( - format!(" | STATUS: {} | PROTOCOL: ZERO-STUBS ACTIVE ", mode_indicator), + format!(" | STATUS: {} | AGENT LOOP: LIVE ", mode_indicator), Style::default().fg(mode_color), ); diff --git a/src/view/workflow.rs b/src/view/workflow.rs index 40362c8..43aed4c 100644 --- a/src/view/workflow.rs +++ b/src/view/workflow.rs @@ -25,8 +25,9 @@ pub fn draw_workflow_panel(frame: &mut Frame, area: Rect, state: &crate::app::st let tool_count = session_runtime.tool_call_results.len(); let pending_count = session_runtime.pending_tool_queue.len(); let bash_count = session_runtime.bash_jobs.len(); - let subagent_queue = session_runtime.subagent_queue; - let edit_count = session_runtime.edit_count; + + let agent_count = state.workflow_engine.agents.len(); + let findings_count = state.workflow_engine.findings.len(); let mut items = Vec::new(); @@ -54,17 +55,19 @@ pub fn draw_workflow_panel(frame: &mut Frame, area: Rect, state: &crate::app::st )))); } - if subagent_queue > 0 { + if agent_count > 0 { items.push(ListItem::new(Line::from(Span::styled( - format!(" Subagent queue: {}", subagent_queue), + format!(" Agents: {}", agent_count), Style::default().fg(Theme::INFO), )))); } - items.push(ListItem::new(Line::from(Span::styled( - format!(" Edits: {}", edit_count), - Style::default().fg(Theme::INFO), - )))); + if findings_count > 0 { + items.push(ListItem::new(Line::from(Span::styled( + format!(" Findings: {}", findings_count), + Style::default().fg(Theme::WARNING), + )))); + } let phase_status = match state.mode { crate::app::state::types::AgentMode::Auto => "Auto-running",