The dev VM is treated as untrusted, and the repo is a read-write virtiofs mount shared with the backend host -- so a CA private key sitting in it is a key that machine can sign with, and a leaf signed by this CA is one the phone's pinned app accepts without question. Pinning against a CA the attacker holds is no pinning at all. So certificates are now generated on the machine that serves them, into $XDG_CONFIG_HOME/ai-app/certs at 0700 with 0600 keys (AI_APP_CERTS overrides), and config.json and session transcripts move to the XDG config and data directories. Transcripts move for a plainer reason than the keys: they are whole conversations, and they were world-readable at 0644. Two smaller things fall out. The host and VM stop sharing one config, which had already put a test token on the production backend. And state stops living where `git clean -xdf` would take the enrollment and every transcript with it. State that predates the move is still read from the repo, with a warning naming where to move it, so an existing install keeps working rather than silently coming up on an empty config -- the precedence is covered by a test, since picking the wrong file would otherwise be silent. Verified: 31 tests, clippy clean; the certificate script writing 0700/0600 into an overridden directory; and the server logging the fallback and serving from it. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com> Claude-Session: https://claude.ai/code/session_017xn8nHw1tw1R6PtiY1eEtw
839 lines
36 KiB
Rust
839 lines
36 KiB
Rust
//! The Claude Code driver: `claude -p` speaking stream-json on stdio,
|
|
//! translated into the common event model.
|
|
//!
|
|
//! Wire format pinned against CLI 2.1.237 by probing (2026-08-24; scripts
|
|
//! summarized here since they live outside the repo):
|
|
//!
|
|
//! - Outbound: `system/init` (carries `session_id`, the `--resume` token),
|
|
//! `stream_event` (raw API deltas; `text_delta` is the streaming text),
|
|
//! consolidated `assistant` messages (their `tool_use` blocks have the
|
|
//! complete input), `user` messages with `tool_result` blocks, a `result`
|
|
//! per turn (usage + cost), `control_request` for anything needing a
|
|
//! human, `control_response` answering ours.
|
|
//! - Permission prompts require the hidden `--permission-prompt-tool stdio`
|
|
//! flag; they arrive as `control_request{subtype:can_use_tool}` and are
|
|
//! answered with `{behavior:"allow",updatedInput}` or
|
|
//! `{behavior:"deny",message}`. `AskUserQuestion` uses the same shape,
|
|
//! with the chosen labels added to `updatedInput` as
|
|
//! `answers:{<question text>:<label>}`.
|
|
//! - Inbound `user` messages sent mid-turn are queued and injected at the
|
|
//! next tool boundary (verified live: the model acknowledged a steer
|
|
//! between two Bash calls) -- the behavior this app exists for.
|
|
//! - `control_request{subtype:set_model}` answers success;
|
|
//! `{subtype:interrupt}` stops the turn.
|
|
|
|
use std::collections::HashMap;
|
|
use std::path::{Path, PathBuf};
|
|
use std::sync::{Arc, Mutex};
|
|
|
|
use anyhow::{Context, Result};
|
|
use serde_json::{Value, json};
|
|
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
|
use tokio::sync::{mpsc, oneshot};
|
|
|
|
use super::driver::{Driver, Event, EventSink, ImageRef, SessionStatus};
|
|
use crate::config::{HostConfig, ProviderConfig, SessionConfig};
|
|
|
|
/// Where the driver remembers its CLI session id between backend runs --
|
|
/// the whole crash-recovery story: respawning with `--resume <id>` picks
|
|
/// the conversation back up from Claude's own session files. Kept in the
|
|
/// session directory rather than config.json so the shared schema stays
|
|
/// free of per-driver state.
|
|
const RESUME_FILE: &str = "claude-session.json";
|
|
|
|
/// Grace period between closing stdin (the polite exit) and SIGKILL.
|
|
const SHUTDOWN_GRACE: std::time::Duration = std::time::Duration::from_secs(5);
|
|
|
|
pub struct ClaudeDriver {
|
|
sink: EventSink,
|
|
/// Lines for the child's stdin; `None` after shutdown started (taking
|
|
/// it closes stdin, which is the CLI's graceful exit signal).
|
|
to_child: Mutex<Option<mpsc::UnboundedSender<String>>>,
|
|
/// Fires SIGKILL if the process outlives the shutdown grace period.
|
|
kill: Mutex<Option<oneshot::Sender<()>>>,
|
|
state: Arc<Mutex<Translator>>,
|
|
session_dir: PathBuf,
|
|
}
|
|
|
|
impl ClaudeDriver {
|
|
pub fn spawn(
|
|
meta: &SessionConfig,
|
|
provider: &ProviderConfig,
|
|
host: Option<&HostConfig>,
|
|
session_dir: &Path,
|
|
sink: EventSink,
|
|
) -> Result<Self> {
|
|
let mut args: Vec<String> = ["-p", "--verbose"].iter().map(|a| a.to_string()).collect();
|
|
let mut push = |flag: &str, value: &str| {
|
|
args.push(flag.to_string());
|
|
args.push(value.to_string());
|
|
};
|
|
push("--input-format", "stream-json");
|
|
push("--output-format", "stream-json");
|
|
// Hidden but load-bearing: without it the CLI resolves permissions
|
|
// itself and nothing ever reaches the phone.
|
|
push("--permission-prompt-tool", "stdio");
|
|
if let Some(model) = &meta.model {
|
|
push("--model", model);
|
|
}
|
|
if let Some(mode) = &meta.permission_mode {
|
|
push("--permission-mode", mode);
|
|
}
|
|
if let Some(resume) = read_resume_token(session_dir) {
|
|
push("--resume", &resume);
|
|
}
|
|
args.push("--include-partial-messages".to_string());
|
|
|
|
let program = provider.command.as_deref().unwrap_or("claude");
|
|
let cwd = meta.cwd.as_deref();
|
|
let where_it_runs = match host {
|
|
Some(host) => format!("on {} ({})", host.name, host.address),
|
|
None => "on this machine".to_string(),
|
|
};
|
|
let mut child = crate::ssh::command(host, program, &args, cwd)
|
|
.spawn()
|
|
.with_context(|| match host {
|
|
Some(host) => format!(
|
|
"couldn't start ssh to run \"{program}\" on {} -- is the ssh client \
|
|
installed here?",
|
|
host.name
|
|
),
|
|
None => format!(
|
|
"couldn't run \"{program}\" on this machine -- is it installed and on \
|
|
PATH? If it lives on another machine, give the session a host to run on.",
|
|
),
|
|
})?;
|
|
tracing::info!("session {} running {program} {where_it_runs}", meta.id);
|
|
|
|
let stdin = child.stdin.take().expect("piped stdin");
|
|
let stdout = child.stdout.take().expect("piped stdout");
|
|
let stderr = child.stderr.take().expect("piped stderr");
|
|
let state = Arc::new(Mutex::new(Translator::new(session_dir.to_path_buf())));
|
|
|
|
// Writer: everything for the child funnels through one channel so
|
|
// driver methods stay sync and writes can't interleave.
|
|
let (to_child, mut from_driver) = mpsc::unbounded_channel::<String>();
|
|
tokio::spawn(async move {
|
|
let mut stdin = stdin;
|
|
while let Some(line) = from_driver.recv().await {
|
|
if stdin.write_all(line.as_bytes()).await.is_err()
|
|
|| stdin.write_all(b"\n").await.is_err()
|
|
|| stdin.flush().await.is_err()
|
|
{
|
|
break;
|
|
}
|
|
}
|
|
// Sender dropped/taken: stdin drops here, closing it -- the
|
|
// CLI's signal to finish up and exit.
|
|
});
|
|
|
|
tokio::spawn(read_stdout(
|
|
stdout,
|
|
Arc::clone(&state),
|
|
sink.clone(),
|
|
session_dir.to_path_buf(),
|
|
));
|
|
|
|
// stderr is diagnostics only; surface it in the log, and keep the
|
|
// last line for the exit report below. For a remote provider this
|
|
// is also where ssh's own failures arrive ("Permission denied",
|
|
// "Could not resolve hostname"), which are the ones a person
|
|
// actually needs to see.
|
|
let last_stderr = Arc::new(Mutex::new(String::new()));
|
|
{
|
|
let last_stderr = Arc::clone(&last_stderr);
|
|
let label = provider.name.clone();
|
|
tokio::spawn(async move {
|
|
let mut lines = BufReader::new(stderr).lines();
|
|
while let Ok(Some(line)) = lines.next_line().await {
|
|
tracing::warn!("{label} stderr: {line}");
|
|
*last_stderr.lock().unwrap() = line;
|
|
}
|
|
});
|
|
}
|
|
|
|
// Monitor: reports process death as an event (with stderr context
|
|
// when it died complaining), and carries the SIGKILL escape hatch.
|
|
let (kill_tx, kill_rx) = oneshot::channel::<()>();
|
|
{
|
|
let sink = sink.clone();
|
|
let label = format!("{} {where_it_runs}", provider.name);
|
|
tokio::spawn(async move {
|
|
let status = tokio::select! {
|
|
status = child.wait() => status.ok(),
|
|
_ = kill_rx => {
|
|
let _ = child.kill().await;
|
|
None
|
|
}
|
|
};
|
|
if let Some(status) = status
|
|
&& !status.success()
|
|
{
|
|
let detail = last_stderr.lock().unwrap().clone();
|
|
let _ = sink.send(Event::Error {
|
|
message: format!(
|
|
"{label} exited with {status}{}",
|
|
if detail.is_empty() { String::new() } else { format!(": {detail}") }
|
|
),
|
|
});
|
|
}
|
|
let _ = sink.send(Event::Status { state: SessionStatus::Exited });
|
|
});
|
|
}
|
|
|
|
Ok(Self {
|
|
sink,
|
|
to_child: Mutex::new(Some(to_child)),
|
|
kill: Mutex::new(Some(kill_tx)),
|
|
state,
|
|
session_dir: session_dir.to_path_buf(),
|
|
})
|
|
}
|
|
|
|
fn send_line(&self, line: String) {
|
|
if let Some(sender) = self.to_child.lock().unwrap().as_ref() {
|
|
let _ = sender.send(line);
|
|
}
|
|
}
|
|
|
|
fn send_control(&self, request: Value) {
|
|
let id = format!("req-{}", super::now() as u64);
|
|
self.send_line(
|
|
json!({"type": "control_request", "request_id": id, "request": request}).to_string(),
|
|
);
|
|
}
|
|
}
|
|
|
|
impl Driver for ClaudeDriver {
|
|
fn send_user_message(&self, text: String, images: Vec<ImageRef>) {
|
|
let mut content = Vec::new();
|
|
for id in &images {
|
|
match attachment_block(&self.session_dir, id) {
|
|
Ok(block) => content.push(block),
|
|
Err(err) => {
|
|
let _ = self.sink.send(Event::Error {
|
|
message: format!("attachment {id} couldn't be sent: {err:#}"),
|
|
});
|
|
}
|
|
}
|
|
}
|
|
if !text.is_empty() {
|
|
content.push(json!({"type": "text", "text": text}));
|
|
}
|
|
// Sent mid-turn this queues for injection at the next tool
|
|
// boundary; sent while idle it starts a turn.
|
|
let _ = self.sink.send(Event::Status { state: SessionStatus::Running });
|
|
self.send_line(
|
|
json!({"type": "user", "message": {"role": "user", "content": content}}).to_string(),
|
|
);
|
|
}
|
|
|
|
fn answer_question(&self, id: &str, answer: &str) {
|
|
let response = {
|
|
let mut state = self.state.lock().unwrap();
|
|
state.answer(id, answer)
|
|
};
|
|
match response {
|
|
AnswerOutcome::Respond(control_response) => {
|
|
let _ = self.sink.send(Event::Status { state: SessionStatus::Running });
|
|
self.send_line(control_response.to_string());
|
|
}
|
|
// A multi-question AskUserQuestion still waiting on the rest.
|
|
AnswerOutcome::Pending => {}
|
|
AnswerOutcome::Unknown => {
|
|
let _ = self.sink.send(Event::Error {
|
|
message: format!("no question {id} is awaiting an answer"),
|
|
});
|
|
}
|
|
}
|
|
}
|
|
|
|
fn interrupt(&self) {
|
|
self.send_control(json!({"subtype": "interrupt"}));
|
|
}
|
|
|
|
fn set_model(&self, model: &str) {
|
|
self.send_control(json!({"subtype": "set_model", "model": model}));
|
|
}
|
|
|
|
fn compact(&self) {
|
|
// Slash commands ride the normal user-message channel.
|
|
self.send_line(
|
|
json!({"type": "user", "message": {"role": "user", "content": [
|
|
{"type": "text", "text": "/compact"}
|
|
]}})
|
|
.to_string(),
|
|
);
|
|
}
|
|
|
|
fn shutdown(&self) {
|
|
// Closing stdin is the polite exit; the kill timer is the escape
|
|
// hatch for a CLI that doesn't oblige.
|
|
self.to_child.lock().unwrap().take();
|
|
if let Some(kill) = self.kill.lock().unwrap().take() {
|
|
tokio::spawn(async move {
|
|
tokio::time::sleep(SHUTDOWN_GRACE).await;
|
|
let _ = kill.send(());
|
|
});
|
|
}
|
|
}
|
|
}
|
|
|
|
async fn read_stdout(
|
|
stdout: tokio::process::ChildStdout,
|
|
state: Arc<Mutex<Translator>>,
|
|
sink: EventSink,
|
|
session_dir: PathBuf,
|
|
) {
|
|
let mut lines = BufReader::new(stdout).lines();
|
|
while let Ok(Some(line)) = lines.next_line().await {
|
|
let Ok(message) = serde_json::from_str::<Value>(&line) else {
|
|
tracing::warn!("unparseable claude output line: {}", &line[..line.len().min(200)]);
|
|
continue;
|
|
};
|
|
let (events, new_session_id) = {
|
|
let mut state = state.lock().unwrap();
|
|
let before = state.session_id.clone();
|
|
let events = state.translate(&message);
|
|
let after = state.session_id.clone();
|
|
(events, if before != after { after } else { None })
|
|
};
|
|
if let Some(session_id) = new_session_id {
|
|
write_resume_token(&session_dir, &session_id);
|
|
}
|
|
for event in events {
|
|
if sink.send(event).is_err() {
|
|
return; // session torn down
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
fn read_resume_token(session_dir: &Path) -> Option<String> {
|
|
let text = std::fs::read_to_string(session_dir.join(RESUME_FILE)).ok()?;
|
|
serde_json::from_str::<Value>(&text)
|
|
.ok()?
|
|
.get("sessionId")?
|
|
.as_str()
|
|
.map(String::from)
|
|
}
|
|
|
|
fn write_resume_token(session_dir: &Path, session_id: &str) {
|
|
let path = session_dir.join(RESUME_FILE);
|
|
if let Err(err) = std::fs::write(&path, json!({"sessionId": session_id}).to_string()) {
|
|
tracing::error!("couldn't persist resume token to {}: {err}", path.display());
|
|
}
|
|
}
|
|
|
|
/// Reads an uploaded attachment into an API image content block.
|
|
fn attachment_block(session_dir: &Path, id: &str) -> Result<Value> {
|
|
// Ids are server-generated hex (see routes::upload_attachment); the
|
|
// check keeps a crafted "id" from naming an arbitrary file.
|
|
if !id.chars().all(|c| c.is_ascii_alphanumeric() || c == '.' || c == '-') {
|
|
anyhow::bail!("invalid attachment id");
|
|
}
|
|
let path = session_dir.join("attachments").join(id);
|
|
let bytes = std::fs::read(&path).with_context(|| format!("read {}", path.display()))?;
|
|
use base64::Engine;
|
|
Ok(json!({
|
|
"type": "image",
|
|
"source": {
|
|
"type": "base64",
|
|
"media_type": media_type_of(id),
|
|
"data": base64::engine::general_purpose::STANDARD.encode(bytes),
|
|
}
|
|
}))
|
|
}
|
|
|
|
fn media_type_of(name: &str) -> &'static str {
|
|
match name.rsplit('.').next() {
|
|
Some("png") => "image/png",
|
|
Some("gif") => "image/gif",
|
|
Some("webp") => "image/webp",
|
|
_ => "image/jpeg",
|
|
}
|
|
}
|
|
|
|
/// What answering a question produced.
|
|
enum AnswerOutcome {
|
|
/// Send this control_response line to the CLI.
|
|
Respond(Value),
|
|
/// Part of a multi-question request; more answers still needed.
|
|
Pending,
|
|
Unknown,
|
|
}
|
|
|
|
/// A `can_use_tool` request we've surfaced to the phone and not yet
|
|
/// answered. For plain permissions there is one implicit question
|
|
/// (Allow/Deny); for AskUserQuestion, one per entry in `questions`.
|
|
struct PendingRequest {
|
|
request_id: String,
|
|
input: Value,
|
|
/// Question text per sub-question, in order -- the keys the answers
|
|
/// map uses. Empty for a plain permission request.
|
|
questions: Vec<String>,
|
|
answers: HashMap<String, String>,
|
|
}
|
|
|
|
/// Translation state: stream-json lines in, common events out. The one
|
|
/// side effect is saving images a tool result carries into the session
|
|
/// dir (they'd bloat the transcript as base64); everything else is pure,
|
|
/// so the dialect mapping is unit-testable from recorded lines.
|
|
struct Translator {
|
|
session_id: Option<String>,
|
|
pending: HashMap<String, PendingRequest>,
|
|
session_dir: PathBuf,
|
|
}
|
|
|
|
impl Translator {
|
|
fn new(session_dir: PathBuf) -> Self {
|
|
Self { session_id: None, pending: HashMap::new(), session_dir }
|
|
}
|
|
fn translate(&mut self, message: &Value) -> Vec<Event> {
|
|
// Events from subagents (Task tool internals) carry a
|
|
// parent_tool_use_id; the transcript shows the Task tool's own
|
|
// start/end instead of every nested step.
|
|
if message.get("parent_tool_use_id").is_some_and(|id| !id.is_null()) {
|
|
return Vec::new();
|
|
}
|
|
match message.get("type").and_then(Value::as_str) {
|
|
Some("system") => {
|
|
if message.get("subtype").and_then(Value::as_str) == Some("init")
|
|
&& let Some(id) = message.get("session_id").and_then(Value::as_str)
|
|
{
|
|
self.session_id = Some(id.to_string());
|
|
}
|
|
Vec::new()
|
|
}
|
|
Some("stream_event") => self.translate_stream_event(&message["event"]),
|
|
Some("assistant") => self.translate_assistant(&message["message"]),
|
|
Some("user") => self.translate_user(message),
|
|
Some("control_request") => self.translate_control_request(message),
|
|
Some("control_response") => {
|
|
let response = &message["response"];
|
|
if response.get("subtype").and_then(Value::as_str) == Some("error") {
|
|
let error = response.get("error").and_then(Value::as_str).unwrap_or("unknown");
|
|
vec![Event::Error { message: format!("claude rejected a request: {error}") }]
|
|
} else {
|
|
Vec::new()
|
|
}
|
|
}
|
|
Some("result") => {
|
|
let usage = &message["usage"];
|
|
let tokens = usage.get("input_tokens").and_then(Value::as_u64).unwrap_or(0)
|
|
+ usage.get("output_tokens").and_then(Value::as_u64).unwrap_or(0);
|
|
let mut events = Vec::new();
|
|
if message.get("is_error").and_then(Value::as_bool).unwrap_or(false) {
|
|
events.push(Event::Error {
|
|
message: message
|
|
.get("result")
|
|
.and_then(Value::as_str)
|
|
.unwrap_or("the turn ended with an error")
|
|
.to_string(),
|
|
});
|
|
}
|
|
if tokens > 0 {
|
|
events.push(Event::UsageDelta { tokens });
|
|
}
|
|
events.push(Event::Status { state: SessionStatus::Idle });
|
|
events
|
|
}
|
|
_ => Vec::new(),
|
|
}
|
|
}
|
|
|
|
/// Raw API streaming: only text deltas become events. Consolidated
|
|
/// blocks arriving later re-carry the same text, so those are skipped
|
|
/// in `translate_assistant` -- one source per fact.
|
|
fn translate_stream_event(&mut self, event: &Value) -> Vec<Event> {
|
|
if event.get("type").and_then(Value::as_str) == Some("content_block_delta")
|
|
&& let Some(delta) = event["delta"].get("text")
|
|
&& event["delta"].get("type").and_then(Value::as_str) == Some("text_delta")
|
|
&& let Some(text) = delta.as_str()
|
|
{
|
|
return vec![Event::AssistantText { delta: text.to_string() }];
|
|
}
|
|
Vec::new()
|
|
}
|
|
|
|
fn translate_assistant(&mut self, message: &Value) -> Vec<Event> {
|
|
let Some(content) = message.get("content").and_then(Value::as_array) else {
|
|
return Vec::new();
|
|
};
|
|
content
|
|
.iter()
|
|
.filter(|block| block.get("type").and_then(Value::as_str) == Some("tool_use"))
|
|
.map(|block| Event::ToolStart {
|
|
id: block.get("id").and_then(Value::as_str).unwrap_or_default().to_string(),
|
|
tool: block.get("name").and_then(Value::as_str).unwrap_or_default().to_string(),
|
|
input: block.get("input").cloned().unwrap_or(Value::Null),
|
|
})
|
|
.collect()
|
|
}
|
|
|
|
fn translate_control_request(&mut self, message: &Value) -> Vec<Event> {
|
|
let request = &message["request"];
|
|
if request.get("subtype").and_then(Value::as_str) != Some("can_use_tool") {
|
|
return Vec::new();
|
|
}
|
|
let request_id =
|
|
message.get("request_id").and_then(Value::as_str).unwrap_or_default().to_string();
|
|
let tool_name = request.get("tool_name").and_then(Value::as_str).unwrap_or("a tool");
|
|
let input = request.get("input").cloned().unwrap_or(Value::Null);
|
|
|
|
let mut events = Vec::new();
|
|
let mut questions = Vec::new();
|
|
if tool_name == "AskUserQuestion" {
|
|
for (i, question) in input
|
|
.get("questions")
|
|
.and_then(Value::as_array)
|
|
.into_iter()
|
|
.flatten()
|
|
.enumerate()
|
|
{
|
|
let text = question
|
|
.get("question")
|
|
.and_then(Value::as_str)
|
|
.unwrap_or("(question)")
|
|
.to_string();
|
|
let options = question
|
|
.get("options")
|
|
.and_then(Value::as_array)
|
|
.into_iter()
|
|
.flatten()
|
|
.filter_map(|option| option.get("label").and_then(Value::as_str))
|
|
.map(String::from)
|
|
.collect();
|
|
events.push(Event::Question {
|
|
id: format!("{request_id}#{i}"),
|
|
prompt: text.clone(),
|
|
options,
|
|
});
|
|
questions.push(text);
|
|
}
|
|
} else {
|
|
let summary = serde_json::to_string_pretty(&input).unwrap_or_default();
|
|
let summary: String = summary.chars().take(600).collect();
|
|
events.push(Event::Question {
|
|
id: request_id.clone(),
|
|
prompt: format!("Allow {tool_name}?\n{summary}"),
|
|
options: vec!["Allow".to_string(), "Deny".to_string()],
|
|
});
|
|
}
|
|
self.pending.insert(
|
|
request_id.clone(),
|
|
PendingRequest { request_id, input, questions, answers: HashMap::new() },
|
|
);
|
|
events.push(Event::Status { state: SessionStatus::AwaitingInput });
|
|
events
|
|
}
|
|
|
|
/// Applies one answer from the phone. Question ids are the control
|
|
/// request id, suffixed `#i` for AskUserQuestion sub-questions.
|
|
fn answer(&mut self, question_id: &str, answer: &str) -> AnswerOutcome {
|
|
let (request_id, sub) = match question_id.split_once('#') {
|
|
Some((request_id, index)) => (request_id, index.parse::<usize>().ok()),
|
|
None => (question_id, None),
|
|
};
|
|
let Some(pending) = self.pending.get_mut(request_id) else {
|
|
return AnswerOutcome::Unknown;
|
|
};
|
|
|
|
let response = if let Some(index) = sub {
|
|
let Some(question) = pending.questions.get(index) else {
|
|
return AnswerOutcome::Unknown;
|
|
};
|
|
pending.answers.insert(question.clone(), answer.to_string());
|
|
if pending.answers.len() < pending.questions.len() {
|
|
return AnswerOutcome::Pending;
|
|
}
|
|
let mut updated = pending.input.clone();
|
|
updated["answers"] = serde_json::to_value(&pending.answers).expect("string map");
|
|
json!({"behavior": "allow", "updatedInput": updated})
|
|
} else if answer.eq_ignore_ascii_case("deny") {
|
|
json!({"behavior": "deny", "message": "The user denied this from the phone."})
|
|
} else {
|
|
json!({"behavior": "allow", "updatedInput": pending.input})
|
|
};
|
|
|
|
let request_id = pending.request_id.clone();
|
|
self.pending.remove(&request_id);
|
|
AnswerOutcome::Respond(json!({
|
|
"type": "control_response",
|
|
"response": {"subtype": "success", "request_id": request_id, "response": response},
|
|
}))
|
|
}
|
|
}
|
|
|
|
impl Translator {
|
|
/// `user` messages: tool results become ToolEnd, with any image parts
|
|
/// saved into the session dir and referenced by an Image event (the
|
|
/// phone fetches them from `/sessions/{id}/files/{ref}`). Replayed and
|
|
/// synthetic user text is skipped -- the manager already recorded the
|
|
/// user's side.
|
|
fn translate_user(&self, message: &Value) -> Vec<Event> {
|
|
let Some(content) = message["message"].get("content").and_then(Value::as_array) else {
|
|
return Vec::new();
|
|
};
|
|
let mut events = Vec::new();
|
|
for block in content {
|
|
if block.get("type").and_then(Value::as_str) != Some("tool_result") {
|
|
continue;
|
|
}
|
|
let mut texts = Vec::new();
|
|
match block.get("content") {
|
|
Some(Value::String(text)) => texts.push(text.clone()),
|
|
Some(Value::Array(parts)) => {
|
|
for part in parts {
|
|
match part.get("type").and_then(Value::as_str) {
|
|
Some("text") => {
|
|
if let Some(text) = part.get("text").and_then(Value::as_str) {
|
|
texts.push(text.to_string());
|
|
}
|
|
}
|
|
Some("image") => {
|
|
if let Some(name) = self.save_image(part) {
|
|
events.push(Event::Image { image: name });
|
|
}
|
|
}
|
|
_ => {}
|
|
}
|
|
}
|
|
}
|
|
_ => {}
|
|
}
|
|
events.push(Event::ToolEnd {
|
|
id: block
|
|
.get("tool_use_id")
|
|
.and_then(Value::as_str)
|
|
.unwrap_or_default()
|
|
.to_string(),
|
|
output: texts.join("\n"),
|
|
});
|
|
}
|
|
events
|
|
}
|
|
|
|
/// Decodes one base64 image block into `files/` and returns its ref.
|
|
fn save_image(&self, part: &Value) -> Option<String> {
|
|
let source = part.get("source")?;
|
|
let data = source.get("data")?.as_str()?;
|
|
use base64::Engine;
|
|
let bytes = base64::engine::general_purpose::STANDARD.decode(data).ok()?;
|
|
let extension = match source.get("media_type").and_then(Value::as_str) {
|
|
Some("image/jpeg") => "jpg",
|
|
Some("image/gif") => "gif",
|
|
Some("image/webp") => "webp",
|
|
_ => "png",
|
|
};
|
|
let name = format!("{}.{extension}", super::random_hex());
|
|
let dir = self.session_dir.join("files");
|
|
if let Err(err) =
|
|
crate::config::create_private_dir(&dir)
|
|
.map_err(std::io::Error::other)
|
|
.and_then(|()| std::fs::write(dir.join(&name), bytes))
|
|
{
|
|
tracing::error!("couldn't save produced image: {err}");
|
|
return None;
|
|
}
|
|
Some(name)
|
|
}
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use super::*;
|
|
|
|
fn translate_lines(translator: &mut Translator, lines: &[&str]) -> Vec<Event> {
|
|
lines
|
|
.iter()
|
|
.flat_map(|line| translator.translate(&serde_json::from_str(line).expect("json")))
|
|
.collect()
|
|
}
|
|
|
|
#[test]
|
|
fn captures_the_resume_token_from_init() {
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
let events = translate_lines(
|
|
&mut translator,
|
|
&[r#"{"type":"system","subtype":"init","cwd":"/x","session_id":"5ecf21da-d53f","tools":[],"model":"claude-haiku-4-5-20251001"}"#],
|
|
);
|
|
assert!(events.is_empty());
|
|
assert_eq!(translator.session_id.as_deref(), Some("5ecf21da-d53f"));
|
|
}
|
|
|
|
#[test]
|
|
fn streams_text_deltas_and_skips_the_consolidated_copy() {
|
|
// Real lines (trimmed) from the 2.1.237 probe.
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
let events = translate_lines(&mut translator, &[
|
|
r#"{"type":"stream_event","event":{"type":"content_block_delta","index":1,"delta":{"type":"text_delta","text":"Done."}},"session_id":"s","parent_tool_use_id":null}"#,
|
|
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"text","text":"Done."}]},"parent_tool_use_id":null,"session_id":"s"}"#,
|
|
r#"{"type":"stream_event","event":{"type":"content_block_delta","index":0,"delta":{"type":"thinking_delta","thinking":"hmm"}},"session_id":"s","parent_tool_use_id":null}"#,
|
|
]);
|
|
assert_eq!(events, vec![Event::AssistantText { delta: "Done.".to_string() }]);
|
|
}
|
|
|
|
#[test]
|
|
fn tool_use_and_result_become_tool_events() {
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
let events = translate_lines(&mut translator, &[
|
|
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"tool_use","id":"toolu_01","name":"Bash","input":{"command":"echo probe-ok"}}]},"parent_tool_use_id":null}"#,
|
|
r#"{"type":"user","message":{"role":"user","content":[{"tool_use_id":"toolu_01","type":"tool_result","content":"probe-ok","is_error":false}]},"parent_tool_use_id":null}"#,
|
|
]);
|
|
assert_eq!(events, vec![
|
|
Event::ToolStart {
|
|
id: "toolu_01".to_string(),
|
|
tool: "Bash".to_string(),
|
|
input: serde_json::json!({"command": "echo probe-ok"}),
|
|
},
|
|
Event::ToolEnd { id: "toolu_01".to_string(), output: "probe-ok".to_string() },
|
|
]);
|
|
}
|
|
|
|
#[test]
|
|
fn subagent_events_are_not_duplicated_into_the_transcript() {
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
let events = translate_lines(&mut translator, &[
|
|
r#"{"type":"assistant","message":{"role":"assistant","content":[{"type":"tool_use","id":"toolu_02","name":"Bash","input":{}}]},"parent_tool_use_id":"toolu_parent"}"#,
|
|
]);
|
|
assert!(events.is_empty());
|
|
}
|
|
|
|
#[test]
|
|
fn a_permission_request_becomes_an_allow_deny_question() {
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
let events = translate_lines(&mut translator, &[
|
|
r#"{"type":"control_request","request_id":"req-1","request":{"subtype":"can_use_tool","tool_name":"Bash","input":{"command":"rm -rf /tmp/x"},"tool_use_id":"toolu_03"}}"#,
|
|
]);
|
|
let Event::Question { id, prompt, options } = &events[0] else {
|
|
panic!("expected a question, got {events:?}");
|
|
};
|
|
assert_eq!(id, "req-1");
|
|
assert!(prompt.contains("Bash") && prompt.contains("rm -rf /tmp/x"));
|
|
assert_eq!(options, &["Allow", "Deny"]);
|
|
assert_eq!(events[1], Event::Status { state: SessionStatus::AwaitingInput });
|
|
|
|
// Allowing echoes the input back; the request is then gone.
|
|
let AnswerOutcome::Respond(response) = translator.answer("req-1", "Allow") else {
|
|
panic!("expected a control response");
|
|
};
|
|
assert_eq!(response["response"]["request_id"], "req-1");
|
|
assert_eq!(response["response"]["response"]["behavior"], "allow");
|
|
assert_eq!(
|
|
response["response"]["response"]["updatedInput"]["command"],
|
|
"rm -rf /tmp/x"
|
|
);
|
|
assert!(matches!(translator.answer("req-1", "Allow"), AnswerOutcome::Unknown));
|
|
}
|
|
|
|
#[test]
|
|
fn denying_a_permission_sends_deny() {
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
translate_lines(&mut translator, &[
|
|
r#"{"type":"control_request","request_id":"req-2","request":{"subtype":"can_use_tool","tool_name":"Write","input":{"file_path":"/etc/passwd"}}}"#,
|
|
]);
|
|
let AnswerOutcome::Respond(response) = translator.answer("req-2", "Deny") else {
|
|
panic!("expected a control response");
|
|
};
|
|
assert_eq!(response["response"]["response"]["behavior"], "deny");
|
|
}
|
|
|
|
#[test]
|
|
fn ask_user_question_rides_the_same_flow_with_answers_keyed_by_question() {
|
|
// The real 2.1.237 shape, verified live: answers go back inside
|
|
// updatedInput, keyed by the question text.
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
let events = translate_lines(&mut translator, &[
|
|
r#"{"type":"control_request","request_id":"req-3","request":{"subtype":"can_use_tool","tool_name":"AskUserQuestion","input":{"questions":[{"question":"Which color?","header":"Color","options":[{"label":"Red"},{"label":"Blue"}],"multiSelect":false},{"question":"Which size?","header":"Size","options":[{"label":"S"},{"label":"L"}],"multiSelect":false}]},"tool_use_id":"toolu_04","requires_user_interaction":true}}"#,
|
|
]);
|
|
let questions: Vec<_> = events
|
|
.iter()
|
|
.filter_map(|event| match event {
|
|
Event::Question { id, prompt, options } => Some((id.clone(), prompt.clone(), options.clone())),
|
|
_ => None,
|
|
})
|
|
.collect();
|
|
assert_eq!(questions.len(), 2);
|
|
assert_eq!(questions[0].0, "req-3#0");
|
|
assert_eq!(questions[0].1, "Which color?");
|
|
assert_eq!(questions[0].2, vec!["Red", "Blue"]);
|
|
|
|
// First answer alone isn't enough; the response goes out when the
|
|
// last sub-question is answered, with all answers aboard.
|
|
assert!(matches!(translator.answer("req-3#0", "Blue"), AnswerOutcome::Pending));
|
|
let AnswerOutcome::Respond(response) = translator.answer("req-3#1", "L") else {
|
|
panic!("expected a control response");
|
|
};
|
|
let updated = &response["response"]["response"]["updatedInput"];
|
|
assert_eq!(updated["answers"]["Which color?"], "Blue");
|
|
assert_eq!(updated["answers"]["Which size?"], "L");
|
|
assert_eq!(updated["questions"][0]["question"], "Which color?");
|
|
}
|
|
|
|
#[test]
|
|
fn images_in_tool_results_are_saved_and_referenced() {
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
// A 1x1 PNG, the smallest real payload worth round-tripping.
|
|
let png = "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAYAAAAfFcSJAAAADUlEQVR42mNk+M9QDwADhgGAWjR9awAAAABJRU5ErkJggg==";
|
|
let line = format!(
|
|
r#"{{"type":"user","message":{{"role":"user","content":[{{"type":"tool_result","tool_use_id":"toolu_05","content":[{{"type":"text","text":"took a screenshot"}},{{"type":"image","source":{{"type":"base64","media_type":"image/png","data":"{png}"}}}}]}}]}},"parent_tool_use_id":null}}"#
|
|
);
|
|
let events = translator.translate(&serde_json::from_str(&line).expect("json"));
|
|
|
|
let Event::Image { image } = &events[0] else {
|
|
panic!("expected an image event, got {events:?}");
|
|
};
|
|
assert!(image.ends_with(".png"));
|
|
let saved = dir.path().join("files").join(image);
|
|
assert!(saved.is_file(), "image not saved at {}", saved.display());
|
|
assert_eq!(events[1], Event::ToolEnd {
|
|
id: "toolu_05".to_string(),
|
|
output: "took a screenshot".to_string(),
|
|
});
|
|
}
|
|
|
|
#[test]
|
|
fn a_turn_result_reports_usage_and_returns_to_idle() {
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
let events = translate_lines(&mut translator, &[
|
|
r#"{"type":"result","subtype":"success","is_error":false,"num_turns":2,"stop_reason":"end_turn","session_id":"s","total_cost_usd":0.0149,"usage":{"input_tokens":18,"output_tokens":164}}"#,
|
|
]);
|
|
assert_eq!(events, vec![
|
|
Event::UsageDelta { tokens: 182 },
|
|
Event::Status { state: SessionStatus::Idle },
|
|
]);
|
|
}
|
|
|
|
#[test]
|
|
fn an_error_result_surfaces_the_message() {
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
let events = translate_lines(&mut translator, &[
|
|
r#"{"type":"result","subtype":"error_during_execution","is_error":true,"result":"something broke","usage":{}}"#,
|
|
]);
|
|
assert_eq!(events[0], Event::Error { message: "something broke".to_string() });
|
|
assert_eq!(*events.last().unwrap(), Event::Status { state: SessionStatus::Idle });
|
|
}
|
|
|
|
#[test]
|
|
fn replayed_and_synthetic_user_text_is_skipped() {
|
|
let dir = tempfile::tempdir().expect("tempdir");
|
|
let mut translator = Translator::new(dir.path().to_path_buf());
|
|
let events = translate_lines(&mut translator, &[
|
|
r#"{"type":"user","message":{"role":"user","content":[{"type":"text","text":"hi"}]},"isReplay":true,"parent_tool_use_id":null}"#,
|
|
r#"{"type":"user","message":{"role":"user","content":[{"type":"text","text":"[continue]"}]},"isSynthetic":true,"parent_tool_use_id":null}"#,
|
|
]);
|
|
assert!(events.is_empty());
|
|
}
|
|
}
|