1use 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 session_state: crate::SessionState,
43}
44
45pub 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 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 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 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 {
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(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 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 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 let _ = child.start_kill();
672 child.wait().await
673 }
674 status = child.wait() => {
675 aborted = false;
676 status
677 }
678 };
679 let stdout_bytes = stdout_task.await.unwrap_or_default();
682 let stderr_bytes = stderr_task.await.unwrap_or_default();
683 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 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}