Phase 2 complete: images both ways
Inbound: POST /sessions/{id}/attachments stores a picked photo under the
session; message attachmentIds become base64 image blocks in the
stream-json user message (verified live: an uploaded red PNG answered
"Red."). Outbound: image parts in tool results are decoded into the
session's files/ dir and referenced by Image events -- the transcript
stays lean -- and GET /sessions/{id}/files/{ref} serves them (verified
via the Read tool round-tripping the same PNG). The app grows an attach
button (system photo picker, upload-on-pick) and renders Image events
inline with an authenticated pinned fetch. Sent attachments are echoed
into the transcript as Image events so every device shows them.
Attachments and files are addressed under their session (a deviation
from PLAN.md's original bare /attachments -- recorded there) so their
lifecycle is the session directory's: deleting the session is still the
complete path out.
Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017xn8nHw1tw1R6PtiY1eEtw
This commit is contained in:
1 parent
95d389e2b8
commit
f2430671a2
7 files changed
+389
-55
No files matched your search
+61
-2
@@ -11,11 +11,13 @@
|
||||
//! POST /sessions/{id}/interrupt
|
||||
//! POST /sessions/{id}/model {model}
|
||||
//! POST /sessions/{id}/compact
|
||||
//! POST /sessions/{id}/attachments multipart image upload -> {id}, referenced by /message
|
||||
//! GET /sessions/{id}/files/{name} images the session produced or was sent
|
||||
//! DELETE /sessions/{id} kill process, delete transcript + files
|
||||
//! ```
|
||||
//!
|
||||
//! Later phases add: `POST /attachments`, `GET /files/{session}/{id}`,
|
||||
//! `GET /usage`, `GET|PUT /hosts` and `/models` -- see PLAN.md's table.
|
||||
//! Later phases add: `GET /usage`, `GET|PUT /hosts` and `/models` -- see
|
||||
//! PLAN.md's table.
|
||||
//!
|
||||
//! Everything here works purely in the common event model; nothing may
|
||||
//! branch on the session kind (that's what drivers are for).
|
||||
@@ -48,6 +50,10 @@ pub fn router(manager: Arc<SessionManager>) -> Router {
|
||||
.route("/sessions/{id}/interrupt", post(interrupt))
|
||||
.route("/sessions/{id}/model", post(set_model))
|
||||
.route("/sessions/{id}/compact", post(compact))
|
||||
.route("/sessions/{id}/attachments", post(upload_attachment))
|
||||
.route("/sessions/{id}/files/{name}", get(serve_file))
|
||||
// Phone photos overflow axum's 2 MB default body cap.
|
||||
.layer(axum::extract::DefaultBodyLimit::max(32 * 1024 * 1024))
|
||||
// An explicit fallback so the auth middleware (layered around the
|
||||
// whole router in main.rs) also covers unknown paths -- a scanner
|
||||
// gets the same 401 everywhere, never a route map.
|
||||
@@ -202,6 +208,59 @@ async fn compact(
|
||||
Ok(StatusCode::NO_CONTENT)
|
||||
}
|
||||
|
||||
/// Accepts one image (any multipart field) and stores it under the
|
||||
/// session; the returned id goes into a later `/message`'s attachmentIds.
|
||||
async fn upload_attachment(
|
||||
State(manager): State<Arc<SessionManager>>,
|
||||
UrlPath(id): UrlPath<String>,
|
||||
mut multipart: axum::extract::Multipart,
|
||||
) -> Result<axum::Json<serde_json::Value>, ApiError> {
|
||||
let session = lookup(&manager, &id)?;
|
||||
let field = multipart
|
||||
.next_field()
|
||||
.await
|
||||
.map_err(|err| ApiError::BadRequest(format!("bad upload: {err}")))?
|
||||
.ok_or_else(|| ApiError::BadRequest("no file in the upload".to_string()))?;
|
||||
let content_type = field.content_type().unwrap_or("image/jpeg").to_string();
|
||||
let bytes = field
|
||||
.bytes()
|
||||
.await
|
||||
.map_err(|err| ApiError::BadRequest(format!("upload read failed: {err}")))?;
|
||||
let name = session.save_attachment(&bytes, &content_type).map_err(bad_request)?;
|
||||
Ok(axum::Json(serde_json::json!({ "id": name })))
|
||||
}
|
||||
|
||||
/// Serves a session's stored images -- both `files/` (produced by tools)
|
||||
/// and `attachments/` (uploaded from the phone), by the id events and
|
||||
/// uploads reference.
|
||||
async fn serve_file(
|
||||
State(manager): State<Arc<SessionManager>>,
|
||||
UrlPath((id, name)): UrlPath<(String, String)>,
|
||||
) -> Result<Response, ApiError> {
|
||||
// Ids are server-generated hex + extension; anything else (and any
|
||||
// path separator in particular) is refused, not resolved.
|
||||
if !name.chars().all(|c| c.is_ascii_alphanumeric() || c == '.') || name.contains("..") {
|
||||
return Err(ApiError::BadRequest("invalid file id".to_string()));
|
||||
}
|
||||
let session = lookup(&manager, &id)?;
|
||||
let candidates =
|
||||
[session.dir().join("files").join(&name), session.dir().join("attachments").join(&name)];
|
||||
let Some(path) = candidates.iter().find(|path| path.is_file()) else {
|
||||
return Err(ApiError::BadRequest(format!("no file {name} in session {id}")));
|
||||
};
|
||||
let bytes = std::fs::read(path).map_err(|err| {
|
||||
tracing::error!("read {} failed: {err}", path.display());
|
||||
ApiError::BadRequest("file unreadable".to_string())
|
||||
})?;
|
||||
let content_type = match name.rsplit('.').next() {
|
||||
Some("png") => "image/png",
|
||||
Some("gif") => "image/gif",
|
||||
Some("webp") => "image/webp",
|
||||
_ => "image/jpeg",
|
||||
};
|
||||
Ok(([(axum::http::header::CONTENT_TYPE, content_type)], bytes).into_response())
|
||||
}
|
||||
|
||||
#[derive(Deserialize)]
|
||||
struct EventsQuery {
|
||||
#[serde(default)]
|
||||
|
||||
+119
-40
@@ -92,7 +92,7 @@ impl ClaudeDriver {
|
||||
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::default()));
|
||||
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.
|
||||
@@ -354,15 +354,20 @@ struct PendingRequest {
|
||||
answers: HashMap<String, String>,
|
||||
}
|
||||
|
||||
/// Pure translation state: stream-json lines in, common events out. No
|
||||
/// I/O, so the whole dialect mapping is unit-testable from recorded lines.
|
||||
#[derive(Default)]
|
||||
/// 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
|
||||
@@ -381,7 +386,7 @@ impl Translator {
|
||||
}
|
||||
Some("stream_event") => self.translate_stream_event(&message["event"]),
|
||||
Some("assistant") => self.translate_assistant(&message["message"]),
|
||||
Some("user") => translate_user(message),
|
||||
Some("user") => self.translate_user(message),
|
||||
Some("control_request") => self.translate_control_request(message),
|
||||
Some("control_response") => {
|
||||
let response = &message["response"];
|
||||
@@ -539,36 +544,77 @@ impl Translator {
|
||||
}
|
||||
}
|
||||
|
||||
/// `user` messages: tool results become ToolEnd (with any images saved
|
||||
/// out-of-band by the caller -- phase 2b); replayed/synthetic user text is
|
||||
/// skipped, since the manager already recorded the user's side.
|
||||
fn translate_user(message: &Value) -> Vec<Event> {
|
||||
let Some(content) = message["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_result"))
|
||||
.map(|block| {
|
||||
let output = match block.get("content") {
|
||||
Some(Value::String(text)) => text.clone(),
|
||||
Some(Value::Array(parts)) => parts
|
||||
.iter()
|
||||
.filter_map(|part| part.get("text").and_then(Value::as_str))
|
||||
.collect::<Vec<_>>()
|
||||
.join("\n"),
|
||||
_ => String::new(),
|
||||
};
|
||||
Event::ToolEnd {
|
||||
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,
|
||||
}
|
||||
})
|
||||
.collect()
|
||||
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) =
|
||||
std::fs::create_dir_all(&dir).and_then(|_| std::fs::write(dir.join(&name), bytes))
|
||||
{
|
||||
tracing::error!("couldn't save produced image: {err}");
|
||||
return None;
|
||||
}
|
||||
Some(name)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -584,7 +630,8 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn captures_the_resume_token_from_init() {
|
||||
let mut translator = Translator::default();
|
||||
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"}"#],
|
||||
@@ -596,7 +643,8 @@ mod tests {
|
||||
#[test]
|
||||
fn streams_text_deltas_and_skips_the_consolidated_copy() {
|
||||
// Real lines (trimmed) from the 2.1.237 probe.
|
||||
let mut translator = Translator::default();
|
||||
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"}"#,
|
||||
@@ -607,7 +655,8 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn tool_use_and_result_become_tool_events() {
|
||||
let mut translator = Translator::default();
|
||||
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}"#,
|
||||
@@ -624,7 +673,8 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn subagent_events_are_not_duplicated_into_the_transcript() {
|
||||
let mut translator = Translator::default();
|
||||
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"}"#,
|
||||
]);
|
||||
@@ -633,7 +683,8 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn a_permission_request_becomes_an_allow_deny_question() {
|
||||
let mut translator = Translator::default();
|
||||
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"}}"#,
|
||||
]);
|
||||
@@ -660,7 +711,8 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn denying_a_permission_sends_deny() {
|
||||
let mut translator = Translator::default();
|
||||
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"}}}"#,
|
||||
]);
|
||||
@@ -674,7 +726,8 @@ mod tests {
|
||||
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 mut translator = Translator::default();
|
||||
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}}"#,
|
||||
]);
|
||||
@@ -702,9 +755,33 @@ mod tests {
|
||||
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 mut translator = Translator::default();
|
||||
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}}"#,
|
||||
]);
|
||||
@@ -716,7 +793,8 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn an_error_result_surfaces_the_message() {
|
||||
let mut translator = Translator::default();
|
||||
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":{}}"#,
|
||||
]);
|
||||
@@ -726,7 +804,8 @@ mod tests {
|
||||
|
||||
#[test]
|
||||
fn replayed_and_synthetic_user_text_is_skipped() {
|
||||
let mut translator = Translator::default();
|
||||
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}"#,
|
||||
|
||||
@@ -94,6 +94,11 @@ impl LiveSession {
|
||||
/// driver -- which queues it for injection mid-run rather than at the
|
||||
/// end of the turn (the point of the whole app).
|
||||
pub fn send_message(&self, text: String, images: Vec<ImageRef>) {
|
||||
// Attachments render in the transcript like any produced image --
|
||||
// the files route serves uploads by the same ref.
|
||||
for image in &images {
|
||||
let _ = self.sink.send(Event::Image { image: image.clone() });
|
||||
}
|
||||
let _ = self.sink.send(Event::UserMessage { text: text.clone() });
|
||||
self.driver.send_user_message(text, images);
|
||||
}
|
||||
@@ -122,6 +127,30 @@ impl LiveSession {
|
||||
&self.transcript_path
|
||||
}
|
||||
|
||||
/// The session's directory (attachments in, produced files out live in
|
||||
/// `attachments/` and `files/` under it).
|
||||
pub fn dir(&self) -> &Path {
|
||||
self.transcript_path.parent().expect("transcript lives in the session dir")
|
||||
}
|
||||
|
||||
/// Stores one uploaded attachment, returning the id `POST /message`
|
||||
/// references it by. Removed with the session directory on delete --
|
||||
/// the same path out as everything else in it.
|
||||
pub fn save_attachment(&self, bytes: &[u8], content_type: &str) -> Result<String> {
|
||||
let extension = match content_type {
|
||||
"image/png" => "png",
|
||||
"image/gif" => "gif",
|
||||
"image/webp" => "webp",
|
||||
_ => "jpg",
|
||||
};
|
||||
let name = format!("{}.{extension}", random_hex());
|
||||
let dir = self.dir().join("attachments");
|
||||
std::fs::create_dir_all(&dir).with_context(|| format!("create {}", dir.display()))?;
|
||||
std::fs::write(dir.join(&name), bytes)
|
||||
.with_context(|| format!("write attachment {name}"))?;
|
||||
Ok(name)
|
||||
}
|
||||
|
||||
fn info(&self) -> SessionInfo {
|
||||
SessionInfo {
|
||||
id: self.meta.id.clone(),
|
||||
@@ -313,13 +342,18 @@ fn default_title(kind: SessionKind) -> String {
|
||||
}
|
||||
|
||||
/// 8 random bytes, hex -- short enough for a URL, unique enough forever at
|
||||
/// this scale. Still checked against the existing list out of caution.
|
||||
fn unique_id(config: &Config) -> String {
|
||||
/// this scale.
|
||||
pub fn random_hex() -> String {
|
||||
use rand::Rng;
|
||||
let mut bytes = [0u8; 8];
|
||||
rand::rng().fill_bytes(&mut bytes);
|
||||
bytes.iter().map(|b| format!("{b:02x}")).collect()
|
||||
}
|
||||
|
||||
/// A [`random_hex`] id not already taken -- checked out of caution.
|
||||
fn unique_id(config: &Config) -> String {
|
||||
loop {
|
||||
let mut bytes = [0u8; 8];
|
||||
rand::rng().fill_bytes(&mut bytes);
|
||||
let id: String = bytes.iter().map(|b| format!("{b:02x}")).collect();
|
||||
let id = random_hex();
|
||||
if !config.sessions.iter().any(|meta| meta.id == id) {
|
||||
return id;
|
||||
}
|
||||
|
||||
Reference in new issue
Block a user