Skip to main content

oxicode/rpc_mode/
handlers.rs

1//! RPC actor and stdio entry point.
2
3use crate::App;
4use crate::app::agent_session::{AgentSession, AgentSessionHandle};
5use crate::store::session::SessionManager;
6use crate::store::settings::Settings;
7use anyhow::{Context, Result};
8use oxicode_agent::{Agent, AgentEvent, AgentHooks, ToolExecutionMode};
9use serde::Serialize;
10use serde_json::Value;
11use std::sync::Arc;
12use std::sync::atomic::Ordering;
13use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
14use tokio::sync::mpsc;
15
16use super::protocol::*;
17
18type OutputSender = mpsc::UnboundedSender<WriterFrame>;
19
20#[derive(Serialize)]
21#[serde(tag = "type", rename_all = "snake_case")]
22pub(crate) enum OutputFrame {
23    Ready,
24}
25
26enum WriterFrame {
27    Response(RpcResponse),
28    Event(RpcEvent),
29    Json(Value),
30}
31
32struct RpcActor {
33    session: AgentSessionHandle,
34    output: mpsc::UnboundedSender<WriterFrame>,
35    active_run: Option<tokio::task::JoinHandle<Result<()>>>,
36    agent: Arc<Agent>,
37    settings: Settings,
38    cwd: String,
39    active_bash: Arc<parking_lot::Mutex<Option<tokio::sync::oneshot::Sender<()>>>>,
40    /// Shared with every swapped-in session so /steer, /follow_up, and
41    /// Ctrl+C continue to work after `/new`, `/resume`, `/fork`, etc.
42    session_state: crate::SessionState,
43}
44
45/// Run the RPC server over JSON Lines on stdin/stdout.
46pub async fn run_rpc_mode(app: App) -> Result<()> {
47    let cwd = std::env::current_dir()
48        .unwrap_or_else(|_| std::path::PathBuf::from("."))
49        .to_string_lossy()
50        .into_owned();
51    let agent = app.agent();
52    let settings = app.settings().clone();
53    let session_state = app.session_state().clone();
54    let session = AgentSession::new(
55        Arc::clone(&agent),
56        settings.clone(),
57        SessionManager::create(&cwd, None),
58        cwd.clone(),
59        session_state.clone(),
60    );
61    let session = session.clone_handle();
62
63    let (command_tx, mut command_rx) = mpsc::unbounded_channel::<RpcCommand>();
64    let (output_tx, output_rx) = mpsc::unbounded_channel::<WriterFrame>();
65
66    let reader = tokio::spawn(read_commands(command_tx, output_tx.clone()));
67    let writer = tokio::spawn(write_frames(output_rx));
68
69    output_tx
70        .send(WriterFrame::Json(
71            serde_json::to_value(OutputFrame::Ready).expect("ready frame is serializable"),
72        ))
73        .map_err(|_| anyhow::anyhow!("RPC stdout writer stopped during startup"))?;
74
75    let mut actor = RpcActor {
76        session,
77        output: output_tx,
78        active_run: None,
79        agent,
80        settings,
81        cwd,
82        active_bash: Arc::new(parking_lot::Mutex::new(None)),
83        session_state,
84    };
85
86    while let Some(command) = command_rx.recv().await {
87        actor.handle(command).await;
88    }
89
90    actor.abort_active_run().await;
91    drop(actor.output);
92    reader.abort();
93    writer.await.context("RPC stdout writer task failed")??;
94    Ok(())
95}
96
97async fn read_commands(command_tx: mpsc::UnboundedSender<RpcCommand>, output: OutputSender) {
98    let mut lines = BufReader::new(tokio::io::stdin()).lines();
99    loop {
100        match lines.next_line().await {
101            Ok(Some(line)) if line.trim().is_empty() => continue,
102            Ok(Some(line)) => match parse_command_line(&line) {
103                Ok(command) => {
104                    if command_tx.send(command).is_err() {
105                        break;
106                    }
107                }
108                Err(response) => {
109                    if output.send(WriterFrame::Response(response)).is_err() {
110                        break;
111                    }
112                }
113            },
114            Ok(None) => break,
115            Err(error) => {
116                let _ = output.send(WriterFrame::Response(error_response(
117                    None,
118                    "read",
119                    format!("Failed to read stdin: {error}"),
120                )));
121                break;
122            }
123        }
124    }
125}
126
127pub(crate) fn parse_command_line(line: &str) -> std::result::Result<RpcCommand, RpcResponse> {
128    let value = parse_json_line(line).map_err(|error| {
129        error_response(None, "parse", format!("Failed to parse command: {error}"))
130    })?;
131    if value.get("jsonrpc").is_some() {
132        return Err(error_response(
133            value.get("id").map(Value::to_string),
134            "jsonrpc",
135            "JSON-RPC framing is not supported by the actor RPC protocol; send native JSONL commands",
136        ));
137    }
138    if value.get("type").and_then(Value::as_str) == Some("extension_ui_response") {
139        return Err(error_response(
140            None,
141            "extension_ui_response",
142            "extension UI responses are not yet supported in RPC mode",
143        ));
144    }
145    // Extract the request ID before `from_value` consumes `value`. Without
146    // this, parse-error responses carry `id: null` and the RPC client's
147    // response-matcher (`send_and_wait`) never matches them — causing the
148    // client to hang until timeout (F-3 test, 2026-07-25).
149    let id = value
150        .get("id")
151        .and_then(Value::as_str)
152        .map(|s| s.to_string());
153    serde_json::from_value(value)
154        .map_err(|error| error_response(id, "parse", format!("Parse error: {error}")))
155}
156
157async fn write_frames(mut rx: mpsc::UnboundedReceiver<WriterFrame>) -> Result<()> {
158    let mut stdout = tokio::io::stdout();
159    while let Some(frame) = rx.recv().await {
160        let line = match frame {
161            WriterFrame::Response(response) => serialize_json_line_obj(&response),
162            WriterFrame::Event(event) => serialize_json_line_obj(&event),
163            WriterFrame::Json(value) => serialize_json_line(&value),
164        };
165        stdout.write_all(line.as_bytes()).await?;
166        stdout.flush().await?;
167    }
168    Ok(())
169}
170
171impl RpcActor {
172    async fn handle(&mut self, command: RpcCommand) {
173        match command {
174            RpcCommand::Prompt {
175                id,
176                message,
177                images,
178                streaming_behavior: _,
179            } => self.prompt(id, message, images),
180            RpcCommand::Steer {
181                id,
182                message,
183                images,
184            } => match build_user_message(message, &images) {
185                Ok(msg) => {
186                    self.session.steer_sync_message(msg);
187                    self.respond(success_response(id, "steer", None));
188                }
189                Err(e) => self.respond(error_response(id, "steer", e)),
190            },
191            RpcCommand::FollowUp {
192                id,
193                message,
194                images,
195            } => match build_user_message(message, &images) {
196                Ok(msg) => {
197                    self.session.follow_up_sync_message(msg);
198                    self.respond(success_response(id, "follow_up", None));
199                }
200                Err(e) => self.respond(error_response(id, "follow_up", e)),
201            },
202            RpcCommand::Abort { id } => {
203                self.session.abort().await;
204                self.session.agent_ref().cancel();
205                self.respond(success_response(id, "abort", None));
206            }
207            RpcCommand::NewSession { id, .. } => {
208                let handle = self
209                    .swap_session(SessionManager::create(&self.cwd, None))
210                    .await;
211                self.respond(success_response(
212                    id,
213                    "new_session",
214                    Some(serde_json::json!({ "session_id": handle.session_id() })),
215                ));
216            }
217            RpcCommand::GetState { id } => self.respond(success_response(
218                id,
219                "get_state",
220                Some(self.session_state_value()),
221            )),
222            RpcCommand::SetModel {
223                id,
224                provider,
225                model_id,
226            } => {
227                let full_model_id = if model_id.contains('/') {
228                    model_id
229                } else {
230                    format!("{provider}/{model_id}")
231                };
232                match self.session.set_model(&full_model_id) {
233                    Ok(()) => self.respond(success_response(
234                        id,
235                        "set_model",
236                        Some(serde_json::json!({ "model": full_model_id })),
237                    )),
238                    Err(error) => self.respond(error_response(id, "set_model", error.to_string())),
239                }
240            }
241            RpcCommand::CycleModel { id } => match self.session.cycle_model() {
242                Some(model) => self.respond(success_response(
243                    id,
244                    "cycle_model",
245                    Some(serde_json::json!({ "model": model })),
246                )),
247                None => self.respond(error_response(
248                    id,
249                    "cycle_model",
250                    "no scoped models configured; use set_model to pick a model",
251                )),
252            },
253            RpcCommand::GetAvailableModels { id } => {
254                let models: Vec<_> = oxicode_sdk::get_all_models()
255                    .map(|entry| {
256                        serde_json::json!({
257                            "provider": entry.provider,
258                            "id": entry.id,
259                        })
260                    })
261                    .collect();
262                self.respond(success_response(
263                    id,
264                    "get_available_models",
265                    Some(serde_json::json!({ "models": models })),
266                ));
267            }
268            RpcCommand::SetThinkingLevel { id, level } => {
269                match crate::store::settings::parse_thinking_level(&level) {
270                    Some(level) => {
271                        self.session.set_thinking_level(level);
272                        self.respond(success_response(id, "set_thinking_level", None));
273                    }
274                    None => self.respond(error_response(
275                        id,
276                        "set_thinking_level",
277                        format!("invalid thinking level: {level}"),
278                    )),
279                }
280            }
281            RpcCommand::CycleThinkingLevel { id } => match self.session.cycle_thinking_level() {
282                Some(level) => self.respond(success_response(
283                    id,
284                    "cycle_thinking_level",
285                    Some(serde_json::json!({ "level": format_thinking_level(level) })),
286                )),
287                None => self.respond(error_response(
288                    id,
289                    "cycle_thinking_level",
290                    "the active model does not support another thinking level",
291                )),
292            },
293            RpcCommand::SetSteeringMode { id, mode } => {
294                self.session.set_steering_mode(mode);
295                self.respond(success_response(id, "set_steering_mode", None));
296            }
297            RpcCommand::SetFollowUpMode { id, mode } => {
298                self.session.set_follow_up_mode(mode);
299                self.respond(success_response(id, "set_follow_up_mode", None));
300            }
301            RpcCommand::Compact {
302                id,
303                custom_instructions,
304            } => match self.session.compact(custom_instructions).await {
305                Ok(result) => self.respond(success_response(
306                    id,
307                    "compact",
308                    Some(serde_json::json!({
309                        "tokens_before": result.tokens_before,
310                        "message_count": self.session.messages().len(),
311                    })),
312                )),
313                Err(error) => self.respond(error_response(id, "compact", error.to_string())),
314            },
315            RpcCommand::SetAutoCompaction { id, enabled } => {
316                self.session.set_auto_compaction(enabled);
317                self.respond(success_response(id, "set_auto_compaction", None));
318            }
319            RpcCommand::SetAutoRetry { id, enabled } => {
320                self.session.set_auto_retry(enabled);
321                self.respond(success_response(id, "set_auto_retry", None));
322            }
323            RpcCommand::AbortRetry { id } => {
324                self.session.cancel_auto_retry();
325                self.respond(success_response(id, "abort_retry", None));
326            }
327            RpcCommand::Bash { id, command } => self.run_bash(id, command),
328            RpcCommand::AbortBash { id } => {
329                let sender = self.active_bash.lock().take();
330                match sender {
331                    Some(tx) => {
332                        let _ = tx.send(());
333                        self.respond(success_response(id, "abort_bash", None));
334                    }
335                    None => self.respond(error_response(
336                        id,
337                        "abort_bash",
338                        "no bash command is running",
339                    )),
340                }
341            }
342            RpcCommand::GetSessionStats { id } => {
343                let state = self.session.state();
344                let stats = self.session.session_stats();
345                self.respond(success_response(
346                    id,
347                    "get_session_stats",
348                    Some(serde_json::json!({
349                        "session_id": stats.session_id,
350                        "message_count": stats.total_messages,
351                        "user_messages": stats.user_messages,
352                        "assistant_messages": stats.assistant_messages,
353                        "tool_calls": stats.tool_calls,
354                        "tool_results": stats.tool_results,
355                        "token_count": state.estimate_tokens(),
356                    })),
357                ));
358            }
359            RpcCommand::GetLastAssistantText { id } => {
360                let text = self
361                    .session
362                    .state()
363                    .messages
364                    .iter()
365                    .rev()
366                    .find_map(|message| {
367                        if let oxicode_sdk::Message::Assistant(message) = message {
368                            Some(message.text_content())
369                        } else {
370                            None
371                        }
372                    });
373                self.respond(success_response(
374                    id,
375                    "get_last_assistant_text",
376                    Some(serde_json::json!({ "text": text })),
377                ));
378            }
379            RpcCommand::SetSessionName { id, name } => {
380                self.session.set_session_name(name);
381                self.respond(success_response(id, "set_session_name", None));
382            }
383            RpcCommand::GetMessages { id } => self.respond(success_response(
384                id,
385                "get_messages",
386                Some(serde_json::json!({ "messages": self.session.messages() })),
387            )),
388            RpcCommand::GetCommands { id } => {
389                let commands: Vec<_> =
390                    crate::tui_vt::slash::registry::SlashRegistry::builtin_commands()
391                        .into_iter()
392                        .map(|(name, description, aliases)| {
393                            serde_json::json!({
394                                "name": name,
395                                "description": description,
396                                "aliases": aliases,
397                            })
398                        })
399                        .collect();
400                self.respond(success_response(
401                    id,
402                    "get_commands",
403                    Some(serde_json::json!({ "commands": commands })),
404                ));
405            }
406            RpcCommand::ExportHtml { id, output_path } => match output_path {
407                Some(path) => match self.session.export_html() {
408                    Ok(html) => match std::fs::write(&path, &html) {
409                        Ok(()) => self.respond(success_response(
410                            id,
411                            "export_html",
412                            Some(serde_json::json!({ "path": path })),
413                        )),
414                        Err(e) => self.respond(error_response(
415                            id,
416                            "export_html",
417                            format!("failed to write {path}: {e}"),
418                        )),
419                    },
420                    Err(e) => self.respond(error_response(id, "export_html", e.to_string())),
421                },
422                None => self.respond(error_response(id, "export_html", "output_path is required")),
423            },
424            RpcCommand::SwitchSession { id, session_path } => {
425                if session_path.is_empty() || !std::path::Path::new(&session_path).exists() {
426                    self.respond(error_response(
427                        id,
428                        "switch_session",
429                        format!("session file not found: {session_path}"),
430                    ));
431                } else {
432                    let sm = SessionManager::open(&session_path, None, Some(&self.cwd));
433                    let handle = self.swap_session(sm).await;
434                    self.respond(success_response(
435                        id,
436                        "switch_session",
437                        Some(serde_json::json!({ "session_id": handle.session_id() })),
438                    ));
439                }
440            }
441            RpcCommand::Fork { id, entry_id } => {
442                // `branch_from_entry` delegates to the live session manager
443                // (with its current entries) and returns the path of a new
444                // session file containing entries up to and including
445                // `entry_id`. Open that file and `swap_session` to it so
446                // future commands operate on the fork.
447                match self.session.branch_from_entry(&entry_id) {
448                    Ok(new_path) => {
449                        let new_sm = SessionManager::open(&new_path, None, Some(&self.cwd));
450                        let handle = self.swap_session(new_sm).await;
451                        self.respond(success_response(
452                            id,
453                            "fork",
454                            Some(serde_json::json!({ "session_id": handle.session_id() })),
455                        ));
456                    }
457                    Err(e) => self.respond(error_response(id, "fork", e)),
458                }
459            }
460            RpcCommand::Clone { id } => match self.session.session_file() {
461                Some(path) => match SessionManager::fork_from(&path, &self.cwd, None) {
462                    Ok(sm) => {
463                        let handle = self.swap_session(sm).await;
464                        self.respond(success_response(
465                            id,
466                            "clone",
467                            Some(serde_json::json!({ "session_id": handle.session_id() })),
468                        ));
469                    }
470                    Err(e) => self.respond(error_response(id, "clone", e)),
471                },
472                None => self.respond(error_response(
473                    id,
474                    "clone",
475                    "no current session file to clone",
476                )),
477            },
478            RpcCommand::GetForkMessages { id } => {
479                self.respond(success_response(
480                    id,
481                    "get_fork_messages",
482                    Some(serde_json::json!({ "messages": self.session.messages() })),
483                ));
484            }
485        }
486    }
487    async fn swap_session(&mut self, session_manager: SessionManager) -> AgentSessionHandle {
488        let new_session = AgentSession::new(
489            Arc::clone(&self.agent),
490            self.settings.clone(),
491            session_manager,
492            self.cwd.clone(),
493            self.session_state.clone(),
494        );
495        let handle = new_session.clone_handle();
496        // Replace the active handle. The previous session is dropped here;
497        // any in-flight prompts continue running on their cloned handles.
498        self.session = handle.clone();
499        handle
500    }
501
502    fn prompt(&mut self, id: Option<String>, message: String, images: Option<Vec<ImageData>>) {
503        if self.session.is_streaming() {
504            self.respond(error_response(
505                id,
506                "prompt",
507                "an agent run is already active; use steer or follow_up",
508            ));
509            return;
510        }
511
512        self.session.reset_should_stop();
513        self.session.agent_ref().reset_cancel();
514        self.session.streaming_flag().store(true, Ordering::SeqCst);
515
516        let session = self.session.clone_handle();
517        let prompt_message = match build_user_message(message.clone(), &images) {
518            Ok(m) => m,
519            Err(e) => {
520                self.respond(error_response(id, "prompt", e));
521                return;
522            }
523        };
524        self.session.persist_user_message(message.clone());
525        let output = self.output.clone();
526        let agent = session.agent_ref();
527        let (event_tx, event_rx) = std::sync::mpsc::channel::<AgentEvent>();
528        let forwarder = tokio::task::spawn_blocking(move || {
529            while let Ok(event) = event_rx.recv() {
530                session.forward_event_to_extensions(&event);
531                if let AgentEvent::MessageEnd { message } = &event {
532                    session.persist_event_message(message);
533                }
534                if let Some(event) = agent_event_to_rpc(&event)
535                    && output.send(WriterFrame::Event(event)).is_err()
536                {
537                    break;
538                }
539            }
540        });
541        let session = self.session.clone_handle();
542
543        let agent_run = tokio::task::spawn_blocking(move || {
544            let runtime = tokio::runtime::Builder::new_current_thread()
545                .enable_all()
546                .build()
547                .context("failed to build RPC agent runtime")?;
548            runtime.block_on(async {
549                let local = tokio::task::LocalSet::new();
550                local
551                    .run_until(agent.run_with_channel_message(prompt_message, event_tx))
552                    .await
553            })
554        });
555        let output = self.output.clone();
556        self.active_run = Some(tokio::spawn(async move {
557            let result = agent_run.await;
558            let _ = forwarder.await;
559            session.persist();
560            session.streaming_flag().store(false, Ordering::SeqCst);
561            match result {
562                Ok(Ok(_)) => {}
563                Ok(Err(error)) => {
564                    let _ = output.send(WriterFrame::Event(RpcEvent::Error {
565                        message: error.to_string(),
566                    }));
567                }
568                Err(error) => {
569                    let _ = output.send(WriterFrame::Event(RpcEvent::Error {
570                        message: format!("RPC agent task failed: {error}"),
571                    }));
572                }
573            }
574            Ok(())
575        }));
576
577        self.respond(success_response(
578            id,
579            "prompt",
580            Some(serde_json::json!({ "accepted": true })),
581        ));
582    }
583
584    fn run_bash(&self, id: Option<String>, command: String) {
585        if is_dangerous_rpc_command(&command) {
586            tracing::warn!("RPC bash command contains dangerous pattern: {:?}", command);
587        }
588
589        // Reject a fresh bash command while another is still in flight — the
590        // RPC protocol only exposes a single active_bash slot for abort.
591        {
592            let mut slot = self.active_bash.lock();
593            if slot.is_some() {
594                self.respond(error_response(
595                    id,
596                    "bash",
597                    "another bash command is already running",
598                ));
599                return;
600            }
601            let (tx, rx) = tokio::sync::oneshot::channel::<()>();
602            *slot = Some(tx);
603            // Drop the lock before spawning to release the mutex.
604            drop(slot);
605            self.spawn_bash(id, command, rx);
606        }
607    }
608
609    fn spawn_bash(
610        &self,
611        id: Option<String>,
612        command: String,
613        abort_rx: tokio::sync::oneshot::Receiver<()>,
614    ) {
615        let output = self.output.clone();
616        let active_bash = Arc::clone(&self.active_bash);
617        tokio::spawn(async move {
618            use std::process::Stdio;
619            use tokio::io::AsyncReadExt;
620            let mut cmd = tokio::process::Command::new("sh");
621            cmd.arg("-c")
622                .arg(&command)
623                .stdin(Stdio::null())
624                .stdout(Stdio::piped())
625                .stderr(Stdio::piped())
626                .kill_on_drop(true);
627            let mut child = match cmd.spawn() {
628                Ok(child) => child,
629                Err(error) => {
630                    // Clear the active slot before responding so a future
631                    // bash command can run.
632                    active_bash.lock().take();
633                    let _ = output.send(WriterFrame::Response(error_response(
634                        id,
635                        "bash",
636                        error.to_string(),
637                    )));
638                    return;
639                }
640            };
641            // Detach piped handles BEFORE the select and drain them concurrently
642            // with `wait()`. Otherwise, a command that writes more than the OS
643            // pipe buffer (~64 KB) would block on write while we wait, deadlocking
644            // forever. On the abort path `start_kill()` + `wait()` closes the
645            // pipes so the read tasks see EOF and finish with whatever is buffered.
646            let mut stdout_pipe = child.stdout.take();
647            let mut stderr_pipe = child.stderr.take();
648            let stdout_task = tokio::spawn(async move {
649                use tokio::io::AsyncReadExt;
650                let mut buf = Vec::new();
651                if let Some(pipe) = stdout_pipe.as_mut() {
652                    let _ = pipe.read_to_end(&mut buf).await;
653                }
654                buf
655            });
656            let stderr_task = tokio::spawn(async move {
657                use tokio::io::AsyncReadExt;
658                let mut buf = Vec::new();
659                if let Some(pipe) = stderr_pipe.as_mut() {
660                    let _ = pipe.read_to_end(&mut buf).await;
661                }
662                buf
663            });
664            let aborted;
665            let exit_status = tokio::select! {
666                biased;
667                _ = abort_rx => {
668                    aborted = true;
669                    // Abort requested: kill the child and reap it so the
670                    // piped streams see EOF and the read tasks finish.
671                    let _ = child.start_kill();
672                    child.wait().await
673                }
674                status = child.wait() => {
675                    aborted = false;
676                    status
677                }
678            };
679            // Collect the concurrent read tasks (they always complete after
680            // the child exits / is killed, which closes the pipes).
681            let stdout_bytes = stdout_task.await.unwrap_or_default();
682            let stderr_bytes = stderr_task.await.unwrap_or_default();
683            // Clear the slot regardless of outcome so the next bash can run.
684            active_bash.lock().take();
685            let response = match exit_status {
686                Ok(status) => success_response(
687                    id,
688                    "bash",
689                    Some(serde_json::json!({
690                        "stdout": String::from_utf8_lossy(&stdout_bytes),
691                        "stderr": String::from_utf8_lossy(&stderr_bytes),
692                        "exit_code": status.code(),
693                        "aborted": aborted,
694                    })),
695                ),
696                Err(error) => error_response(id, "bash", error.to_string()),
697            };
698            let _ = output.send(WriterFrame::Response(response));
699        });
700    }
701
702    fn session_state_value(&self) -> Value {
703        let state = self.session.state();
704        let model_id = self.session.model_id();
705        let (provider, id) = model_id
706            .split_once('/')
707            .map(|(provider, id)| (provider.to_string(), id.to_string()))
708            .unwrap_or_else(|| (String::new(), model_id));
709        serde_json::json!({
710            "model": ModelInfo { provider, id },
711            "thinking_level": format_thinking_level(self.session.thinking_level()),
712            "is_streaming": self.session.is_streaming(),
713            "is_compacting": self.session.is_compacting(),
714            "steering_mode": self.session.steering_mode(),
715            "follow_up_mode": self.session.follow_up_mode(),
716            "session_id": self.session.session_id(),
717            "auto_compaction_enabled": self.session.auto_compaction_enabled(),
718            "message_count": state.messages.len(),
719            "pending_message_count": self.session.pending_message_count(),
720            "iteration": state.iteration,
721            "stop_reason": state.stop_reason,
722        })
723    }
724
725    fn respond(&self, response: RpcResponse) {
726        let _ = self.output.send(WriterFrame::Response(response));
727    }
728
729    async fn abort_active_run(&mut self) {
730        self.session.abort().await;
731        self.session.agent_ref().cancel();
732        if let Some(handle) = self.active_run.take() {
733            let _ = handle.await;
734        }
735    }
736}
737
738pub(crate) fn agent_event_to_rpc(event: &AgentEvent) -> Option<RpcEvent> {
739    match event {
740        AgentEvent::AgentStart { .. } | AgentEvent::Start { .. } => Some(RpcEvent::AgentStart),
741        AgentEvent::AgentEnd { .. } | AgentEvent::Complete { .. } | AgentEvent::Cancelled => {
742            Some(RpcEvent::AgentEnd)
743        }
744        AgentEvent::Thinking | AgentEvent::ThinkingDelta { .. } => Some(RpcEvent::Thinking),
745        AgentEvent::ThinkingEnd => Some(RpcEvent::ThinkingEnd),
746        AgentEvent::TextChunk { text } => Some(RpcEvent::TextChunk { text: text.clone() }),
747        AgentEvent::MessageUpdate { delta, .. }
748            if delta.as_text().is_some_and(|t| !t.is_empty()) =>
749        {
750            Some(RpcEvent::TextChunk {
751                text: delta.as_text().unwrap_or("").to_string(),
752            })
753        }
754        AgentEvent::ToolCallDelta {
755            tool_call_id,
756            args_delta,
757        } => Some(RpcEvent::ToolCallDelta {
758            tool_call_id: tool_call_id.clone(),
759            args_delta: args_delta.clone(),
760        }),
761        AgentEvent::ToolExecutionStart { tool_name, .. }
762        | AgentEvent::ToolStart { tool_name, .. } => Some(RpcEvent::ToolStart {
763            tool: tool_name.clone(),
764        }),
765        AgentEvent::ToolExecutionEnd { tool_name, .. } => Some(RpcEvent::ToolEnd {
766            tool: tool_name.clone(),
767        }),
768        AgentEvent::Error { message, .. } | AgentEvent::ToolError { error: message, .. } => {
769            Some(RpcEvent::Error {
770                message: message.clone(),
771            })
772        }
773        _ => None,
774    }
775}
776
777fn build_user_message(
778    text: String,
779    images: &Option<Vec<ImageData>>,
780) -> Result<oxicode_ai::Message, String> {
781    let mut blocks: Vec<oxicode_ai::ContentBlock> = Vec::new();
782    if !text.is_empty() {
783        blocks.push(oxicode_ai::ContentBlock::Text(
784            oxicode_ai::TextContent::new(text),
785        ));
786    }
787    if let Some(images) = images.as_ref() {
788        for image in images {
789            // Accept either raw base64 or a `data:<mime>;base64,...` prefix.
790            let raw = image.source.as_str();
791            let (mime, payload) = if let Some(rest) = raw.strip_prefix("data:") {
792                match rest.split_once(";base64,") {
793                    Some((mime, b64)) => (mime.to_string(), b64.to_string()),
794                    None => {
795                        return Err(format!(
796                            "image source for {} is not a base64 data URL",
797                            image.media_type
798                        ));
799                    }
800                }
801            } else {
802                let mime = if image.media_type.is_empty() {
803                    "image/png".to_string()
804                } else {
805                    image.media_type.clone()
806                };
807                (mime, raw.to_string())
808            };
809            blocks.push(oxicode_ai::ContentBlock::Image(
810                oxicode_ai::ImageContent::new(payload, mime),
811            ));
812        }
813    }
814    if blocks.is_empty() {
815        return Err("steer/follow_up requires a non-empty message".to_string());
816    }
817    Ok(oxicode_ai::Message::User(oxicode_ai::UserMessage::new(
818        blocks,
819    )))
820}
821
822fn success_response(id: Option<String>, command: &str, data: Option<Value>) -> RpcResponse {
823    RpcResponse::Response {
824        id,
825        command: command.to_string(),
826        success: true,
827        data,
828        error: None,
829    }
830}
831
832fn error_response(id: Option<String>, command: &str, error: impl Into<String>) -> RpcResponse {
833    RpcResponse::Response {
834        id,
835        command: command.to_string(),
836        success: false,
837        data: None,
838        error: Some(error.into()),
839    }
840}
841
842pub(crate) fn unsupported_response(id: Option<String>, command: &str) -> RpcResponse {
843    error_response(
844        id,
845        command,
846        format!("{command} is not yet supported in RPC mode"),
847    )
848}
849
850fn format_thinking_level(level: crate::store::settings::ThinkingLevel) -> &'static str {
851    match level {
852        crate::store::settings::ThinkingLevel::Off => "off",
853        crate::store::settings::ThinkingLevel::Minimal => "minimal",
854        crate::store::settings::ThinkingLevel::Low => "low",
855        crate::store::settings::ThinkingLevel::Medium => "medium",
856        crate::store::settings::ThinkingLevel::High => "high",
857        crate::store::settings::ThinkingLevel::XHigh => "xhigh",
858    }
859}
860
861fn is_dangerous_rpc_command(command: &str) -> bool {
862    let lower = command.to_lowercase();
863    lower.contains("/etc/passwd")
864        || lower.contains("id_rsa")
865        || lower.contains("curl | nc")
866        || lower.contains("/dev/tcp/")
867        || lower.contains("rm -rf /")
868        || lower.contains("> /etc/")
869        || lower.contains("mkfifo")
870}