1use std::collections::{HashSet, VecDeque};
11use std::sync::atomic::{AtomicBool, Ordering};
12use std::sync::{Arc, Mutex};
13
14use serde_json::{json, Value};
15use tokio::io::{AsyncBufRead, AsyncBufReadExt, AsyncWrite, AsyncWriteExt};
16use tokio::sync::{broadcast, mpsc};
17
18use crate::server::{RuntimeSubmitError, SERVER_MAX_LINE_BYTES};
19use crate::{
20 ChatMessage, FrontendApprovalDecision, FrontendEvent, FrontendRequest, FrontendRequestKind,
21 FrontendResponse, FrontendRuntimeError, Role, SdkRuntime, FRONTEND_REPLAY_CAPACITY,
22};
23
24const ACP_APPROVAL_REQUEST_ID_PREFIX: &str = "supercode-approval-";
25
26pub const ACP_PROTOCOL_VERSION: u64 = 1;
28
29#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
36pub enum AcpCompatibilityProfile {
37 #[default]
39 Standard,
40 GooseAcd3c135,
42}
43
44#[cfg(feature = "adapter-api")]
47pub type HttpAcpRuntime = crate::HttpFrontendRuntime;
48
49struct AcpRuntimeBridge {
50 runtime: Arc<dyn SdkRuntime>,
51 session_id: String,
52 model: String,
53 history: Vec<ChatMessage>,
54 history_cursor: u64,
55 initial_replay_cursor: u64,
56 events: broadcast::Sender<FrontendEvent>,
57 routing: Mutex<EventRoutingState>,
58}
59
60#[derive(Default)]
61struct EventRoutingState {
62 active: bool,
63 pending: VecDeque<FrontendEvent>,
64 failure: Option<String>,
65}
66
67impl EventRoutingState {
68 fn buffer(&mut self, event: FrontendEvent) -> Result<(), ()> {
69 if self.pending.len() >= FRONTEND_REPLAY_CAPACITY {
70 self.pending.clear();
71 self.failure = Some(format!(
72 "ACP attachment received more than {FRONTEND_REPLAY_CAPACITY} events before session activation; reconnect required to preserve a gap-free stream"
73 ));
74 return Err(());
75 }
76 self.pending.push_back(event);
77 Ok(())
78 }
79
80 fn activate(&mut self, acknowledged_cursor: u64) -> Result<Vec<FrontendEvent>, String> {
81 if let Some(error) = &self.failure {
82 return Err(error.clone());
83 }
84 self.active = true;
85 Ok(self
86 .pending
87 .drain(..)
88 .filter(|event| event.sequence > acknowledged_cursor)
89 .collect())
90 }
91}
92
93impl AcpRuntimeBridge {
94 async fn connect(runtime: Arc<dyn SdkRuntime>) -> Result<Arc<Self>, FrontendRuntimeError> {
95 let mut attachment = runtime.attach(200).await?;
96 let session_id = attachment.descriptor.session_id.clone();
97 let model = attachment.descriptor.model.clone();
98 let history = std::mem::take(&mut attachment.history);
99 let history_cursor = attachment.history_cursor;
100 let initial_replay_cursor = attachment
101 .replay
102 .iter()
103 .map(|event| event.sequence)
104 .max()
105 .unwrap_or(history_cursor);
106 let (events, _) = broadcast::channel(1024);
107 let bridge = Arc::new(Self {
108 runtime,
109 session_id,
110 model,
111 history,
112 history_cursor,
113 initial_replay_cursor,
114 events,
115 routing: Mutex::new(EventRoutingState::default()),
116 });
117 let weak = Arc::downgrade(&bridge);
118 tokio::spawn(async move {
119 loop {
120 let event = match attachment.next_event().await {
121 Ok(event) => event,
122 Err(error) => FrontendEvent {
123 sequence: u64::MAX,
124 kind: "runtime_disconnected".into(),
125 payload: json!({
126 "type": "runtime_disconnected",
127 "message": error.to_string(),
128 }),
129 },
130 };
131 let terminal = disconnect_message(&event).is_some();
132 let Some(bridge) = weak.upgrade() else {
133 return;
134 };
135 let mut routing = bridge
136 .routing
137 .lock()
138 .unwrap_or_else(std::sync::PoisonError::into_inner);
139 if routing.active {
140 drop(routing);
141 let _ = bridge.events.send(event);
142 } else if routing.buffer(event).is_err() {
143 return;
144 }
145 if terminal {
146 return;
147 }
148 }
149 });
150 Ok(bridge)
151 }
152
153 fn subscribe(&self) -> broadcast::Receiver<FrontendEvent> {
154 self.events.subscribe()
155 }
156
157 fn session_id(&self) -> &str {
158 &self.session_id
159 }
160
161 fn activate_and_route(
166 &self,
167 history: Vec<Value>,
168 acknowledged_cursor: u64,
169 tx: &mpsc::UnboundedSender<Value>,
170 ) -> Result<(), String> {
171 let mut routing = self
172 .routing
173 .lock()
174 .unwrap_or_else(std::sync::PoisonError::into_inner);
175 if routing.active {
176 return Ok(());
177 }
178 let pending = routing.activate(acknowledged_cursor)?;
179 for update in history {
180 let _ = tx.send(update);
181 }
182 let mut terminal = None;
183 for event in pending {
184 let _ = route_event_projection(tx, self.session_id(), &event);
185 if disconnect_message(&event).is_some() {
186 terminal = Some(event);
187 break;
188 }
189 }
190 drop(routing);
191 if let Some(event) = terminal {
194 let _ = self.events.send(event);
195 }
196 Ok(())
197 }
198
199 async fn submit(
200 &self,
201 prompt: String,
202 image_urls: Vec<String>,
203 ) -> Result<String, FrontendRuntimeError> {
204 let mut events = self.subscribe();
205 let reply = self.runtime.submit_with_images(prompt, image_urls).await?;
206 let deadline = tokio::time::Instant::now() + std::time::Duration::from_secs(5);
207 loop {
208 match tokio::time::timeout_at(deadline, events.recv()).await {
209 Ok(Ok(event)) if is_turn_boundary(&event) => return Ok(reply),
210 Ok(Ok(event)) => {
211 if let Some(message) = disconnect_message(&event) {
212 return Err(FrontendRuntimeError::Transport(message.to_string()));
213 }
214 }
215 Ok(Err(broadcast::error::RecvError::Lagged(skipped))) => {
216 return Err(FrontendRuntimeError::Transport(format!(
217 "SDK event stream lagged by {skipped} event(s)"
218 )));
219 }
220 Ok(Err(broadcast::error::RecvError::Closed)) => {
221 return Err(FrontendRuntimeError::Closed);
222 }
223 Err(_) => {
224 return Err(FrontendRuntimeError::Transport(
225 "SDK event stream did not confirm turn completion".into(),
226 ));
227 }
228 }
229 }
230 }
231
232 async fn interrupt(&self) -> bool {
233 self.runtime.interrupt().await.unwrap_or(false)
234 }
235}
236
237fn is_turn_boundary(event: &FrontendEvent) -> bool {
238 matches!(
239 event.payload.get("type").and_then(Value::as_str),
240 Some("turn_succeeded" | "turn_failed" | "turn_interrupted")
241 )
242}
243
244pub struct AcpServer {
246 runtime: Arc<AcpRuntimeBridge>,
247 compatibility: AcpCompatibilityProfile,
248 session_open: AtomicBool,
249 history_replayed: AtomicBool,
250 prompt_active: AtomicBool,
251}
252
253enum RouterCommand {
254 Prompt {
255 id: Value,
256 prompt: String,
257 image_urls: Vec<String>,
258 frontend_reply: bool,
259 },
260}
261
262async fn project_prompt(
263 server: &AcpServer,
264 events: &mut tokio::sync::broadcast::Receiver<FrontendEvent>,
265 prompt: String,
266 image_urls: Vec<String>,
267 frontend_reply: bool,
268 tx: &mpsc::UnboundedSender<Value>,
269) -> (Result<Value, FrontendRuntimeError>, bool) {
270 let runtime = server.runtime.clone();
271 let mut submit = tokio::spawn(async move { runtime.submit(prompt, image_urls).await });
275 let mut streamed = String::new();
276 let mut terminal = false;
277 let submit_result = loop {
278 tokio::select! {
279 result = &mut submit => break match result {
280 Ok(result) => result,
281 Err(error) => Err(FrontendRuntimeError::Transport(format!(
282 "ACP runtime submit task failed: {error}"
283 ))),
284 },
285 event = events.recv() => match event {
286 Ok(event) => {
287 if let Some(message) = disconnect_message(&event) {
288 submit.abort();
289 terminal = true;
290 break Err(FrontendRuntimeError::Transport(message.to_string()));
291 }
292 if let Some(text) = assistant_event_text(&event) {
293 streamed.push_str(text);
294 }
295 if !route_event_projection(tx, server.runtime.session_id(), &event) {
296 submit.abort();
297 terminal = true;
298 break Err(FrontendRuntimeError::Closed);
299 }
300 }
301 Err(tokio::sync::broadcast::error::RecvError::Lagged(skipped)) => {
302 submit.abort();
303 terminal = true;
304 break Err(FrontendRuntimeError::Transport(format!(
305 "ACP event stream lagged by {skipped} events"
306 )));
307 }
308 Err(tokio::sync::broadcast::error::RecvError::Closed) => {
309 submit.abort();
310 terminal = true;
311 break Err(FrontendRuntimeError::Closed);
312 }
313 }
314 }
315 };
316
317 while !terminal {
322 match events.try_recv() {
323 Ok(event) => {
324 if let Some(message) = disconnect_message(&event) {
325 terminal = true;
326 return (
327 Err(FrontendRuntimeError::Transport(message.to_string())),
328 terminal,
329 );
330 }
331 if let Some(text) = assistant_event_text(&event) {
332 streamed.push_str(text);
333 }
334 if !route_event_projection(tx, server.runtime.session_id(), &event) {
335 terminal = true;
336 return (Err(FrontendRuntimeError::Closed), terminal);
337 }
338 }
339 Err(tokio::sync::broadcast::error::TryRecvError::Empty) => break,
340 Err(tokio::sync::broadcast::error::TryRecvError::Lagged(skipped)) => {
341 terminal = true;
342 return (
343 Err(FrontendRuntimeError::Transport(format!(
344 "ACP event stream lagged by {skipped} events"
345 ))),
346 terminal,
347 );
348 }
349 Err(tokio::sync::broadcast::error::TryRecvError::Closed) => {
350 terminal = true;
351 return (Err(FrontendRuntimeError::Closed), terminal);
352 }
353 }
354 }
355
356 let result = match submit_result {
357 Ok(reply) if reply.starts_with(&streamed) => {
358 let missing = &reply[streamed.len()..];
359 if !missing.is_empty() {
360 let _ = tx.send(session_update(
361 server.runtime.session_id(),
362 json!({
363 "sessionUpdate": "agent_message_chunk",
364 "content": {"type": "text", "text": missing}
365 }),
366 ));
367 streamed.push_str(missing);
368 }
369 if streamed.is_empty() {
370 Err(FrontendRuntimeError::Execution {
371 operation: crate::SdkOperation::Input,
372 message: "runtime completed without assistant output".into(),
373 })
374 } else {
375 if frontend_reply {
376 Ok(json!({"reply": reply}))
377 } else {
378 Ok(json!({"stopReason": "end_turn"}))
379 }
380 }
381 }
382 Ok(reply) => Err(FrontendRuntimeError::Execution {
383 operation: crate::SdkOperation::Input,
384 message: format!(
385 "runtime reply did not match streamed assistant output (reply {} bytes, stream {} bytes)",
386 reply.len(),
387 streamed.len()
388 ),
389 }),
390 Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Interrupted)) => {
391 Ok(json!({"stopReason": "cancelled"}))
392 }
393 Err(error) => Err(error),
394 };
395 (result, terminal)
396}
397
398impl AcpServer {
399 pub async fn new(runtime: Arc<dyn SdkRuntime>) -> Result<Arc<Self>, FrontendRuntimeError> {
402 Self::new_with_profile(runtime, AcpCompatibilityProfile::Standard).await
403 }
404
405 pub async fn new_with_profile(
409 runtime: Arc<dyn SdkRuntime>,
410 compatibility: AcpCompatibilityProfile,
411 ) -> Result<Arc<Self>, FrontendRuntimeError> {
412 Ok(Arc::new(Self {
413 runtime: AcpRuntimeBridge::connect(runtime).await?,
414 compatibility,
415 session_open: AtomicBool::new(false),
416 history_replayed: AtomicBool::new(false),
417 prompt_active: AtomicBool::new(false),
418 }))
419 }
420
421 pub fn runtime(&self) -> &Arc<dyn SdkRuntime> {
423 &self.runtime.runtime
424 }
425
426 fn initialize(&self) -> Value {
427 let mut frontend_methods = crate::FrontendFacadeMethod::ALL
428 .into_iter()
429 .map(crate::FrontendFacadeMethod::wire_name)
430 .collect::<Vec<_>>();
431 frontend_methods.extend(["session/cancel", "session/steer", "session/respond"]);
432 json!({
433 "protocolVersion": ACP_PROTOCOL_VERSION,
434 "agentCapabilities": {
435 "loadSession": true,
436 "sessionCapabilities": {"resume": {}},
437 "promptCapabilities": {"image": false, "embeddedContext": false},
438 "_meta": {
439 "supercode": {
440 "frontend": {
441 "schemaVersion": crate::frontend::FRONTEND_RUNTIME_SCHEMA_VERSION,
442 "contract": "supercode.frontend.contract.v2",
443 "eventMethod": "frontend.v2.event",
444 "runtimeOwnedByClient": false,
445 "methods": frontend_methods
446 }
447 }
448 }
449 },
450 "authMethods": [],
451 "agentInfo": {
452 "name": "supercode",
453 "title": "Supercode",
454 "version": env!("CARGO_PKG_VERSION")
455 }
456 })
457 }
458
459 fn open_session(&self, requested: Option<&str>) -> Result<Value, String> {
460 if let Some(requested) = requested {
461 if requested != self.runtime.session_id() {
462 return Err(format!(
463 "runtime `{}` is not session `{requested}`",
464 self.runtime.session_id()
465 ));
466 }
467 }
468 self.session_open.store(true, Ordering::SeqCst);
469 Ok(json!({"sessionId": self.runtime.session_id()}))
470 }
471
472 fn goose_defaults(&self) -> Result<Value, String> {
473 if self.compatibility != AcpCompatibilityProfile::GooseAcd3c135 {
474 return Err("Goose defaults are unavailable on the standard ACP profile".into());
475 }
476 Ok(json!({
477 "providerId": "supercode",
478 "modelId": self.runtime.model,
479 }))
480 }
481
482 fn validate_new_session(&self, params: &Value) -> Result<(), String> {
483 if self.compatibility != AcpCompatibilityProfile::GooseAcd3c135 {
484 return Ok(());
485 }
486 if params
487 .get("mcpServers")
488 .and_then(Value::as_array)
489 .is_some_and(|servers| !servers.is_empty())
490 {
491 return Err("Goose frontend cannot mutate the selected runtime's MCP servers".into());
492 }
493 Ok(())
494 }
495
496 fn history_updates_once(&self) -> Vec<Value> {
497 if self.compatibility != AcpCompatibilityProfile::GooseAcd3c135
498 || self.history_replayed.swap(true, Ordering::SeqCst)
499 {
500 return Vec::new();
501 }
502 project_history(
503 self.runtime.session_id(),
504 &self.runtime.history,
505 self.runtime.history_cursor,
506 )
507 }
508
509 fn prompt_text(params: &Value) -> Result<String, String> {
510 let prompt = params
511 .get("prompt")
512 .and_then(Value::as_array)
513 .ok_or_else(|| "session/prompt requires a `prompt` content array".to_string())?;
514 let text = prompt
515 .iter()
516 .filter(|part| part.get("type").and_then(Value::as_str) == Some("text"))
517 .filter_map(|part| part.get("text").and_then(Value::as_str))
518 .collect::<Vec<_>>()
519 .join("\n");
520 if text.is_empty() {
521 Err("session/prompt contains no text content".into())
522 } else {
523 Ok(text)
524 }
525 }
526
527 fn validate_session(&self, params: &Value) -> Result<(), String> {
528 if !self.session_open.load(Ordering::SeqCst) {
529 return Err("open a session before prompting".into());
530 }
531 let requested = params
532 .get("sessionId")
533 .and_then(Value::as_str)
534 .ok_or_else(|| "request omitted `sessionId`".to_string())?;
535 if requested != self.runtime.session_id() {
536 return Err(format!(
537 "runtime `{}` is not session `{requested}`",
538 self.runtime.session_id()
539 ));
540 }
541 Ok(())
542 }
543}
544
545fn response(id: Value, result: Result<Value, String>) -> Value {
546 match result {
547 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
548 Err(message) => json!({
549 "jsonrpc": "2.0",
550 "id": id,
551 "error": {"code": -32000, "message": message}
552 }),
553 }
554}
555
556fn sdk_response(id: Value, result: Result<Value, FrontendRuntimeError>) -> Value {
557 match result {
558 Ok(result) => json!({"jsonrpc": "2.0", "id": id, "result": result}),
559 Err(error) => {
560 let code = match error.code() {
561 crate::SdkErrorCode::Busy => -32000,
562 crate::SdkErrorCode::Unauthenticated => -32030,
563 crate::SdkErrorCode::Unauthorized => -32031,
564 crate::SdkErrorCode::ControllerRequired => -32032,
565 crate::SdkErrorCode::LeaseExpired => -32033,
566 _ => -32002,
567 };
568 let mut envelope = json!({
569 "jsonrpc": "2.0",
570 "id": id,
571 "error": {
572 "code": code,
573 "name": error.code(),
574 "operation": error.operation(),
575 "message": error.to_string(),
576 }
577 });
578 if let Some(detail) = envelope.get_mut("error").and_then(Value::as_object_mut) {
579 match error {
580 FrontendRuntimeError::Unauthorized { permission } => {
581 detail.insert("permission".into(), Value::String(permission));
582 }
583 FrontendRuntimeError::ControllerRequired {
584 holder,
585 expires_at_ms,
586 } => {
587 if let Some(holder) = holder {
588 detail.insert("holder".into(), Value::String(holder));
589 }
590 if let Some(expires_at_ms) = expires_at_ms {
591 detail.insert("expiresAtMs".into(), json!(expires_at_ms));
592 }
593 }
594 _ => {}
595 }
596 }
597 envelope
598 }
599 }
600}
601
602fn session_update(session_id: &str, update: Value) -> Value {
603 json!({
604 "jsonrpc": "2.0",
605 "method": "session/update",
606 "params": {"sessionId": session_id, "update": update}
607 })
608}
609
610fn projected_session_update(session_id: &str, event: &FrontendEvent, mut update: Value) -> Value {
611 if let Some(update) = update.as_object_mut() {
612 update.insert(
613 "_meta".into(),
614 json!({
615 "supercode": {
616 "sdkSequence": event.sequence,
617 "sdkKind": event.kind,
618 }
619 }),
620 );
621 }
622 session_update(session_id, update)
623}
624
625fn frontend_event_notification(session_id: &str, event: &FrontendEvent) -> Value {
626 json!({
627 "jsonrpc": "2.0",
628 "method": "frontend.v2.event",
629 "params": {"sessionId": session_id, "event": event}
630 })
631}
632
633fn approval_request_id(request_id: u64) -> String {
634 format!("{ACP_APPROVAL_REQUEST_ID_PREFIX}{request_id}")
635}
636
637fn approval_request_id_from_wire(message: &Value) -> Option<u64> {
638 let wire_id = message.get("id")?.as_str()?;
639 let request_id = wire_id
640 .strip_prefix(ACP_APPROVAL_REQUEST_ID_PREFIX)?
641 .parse::<u64>()
642 .ok()?;
643 (approval_request_id(request_id) == wire_id).then_some(request_id)
644}
645
646fn project_approval_request(session_id: &str, event: &FrontendEvent) -> Option<Value> {
647 if event.payload.get("type").and_then(Value::as_str) != Some("request") {
648 return None;
649 }
650 let request: FrontendRequest =
651 serde_json::from_value(event.payload.get("request")?.clone()).ok()?;
652 if request.kind != FrontendRequestKind::Approval {
653 return None;
654 }
655 let tool = request
656 .payload
657 .get("tool")
658 .and_then(Value::as_str)
659 .unwrap_or("tool");
660 let title = request
661 .payload
662 .get("subject")
663 .and_then(Value::as_str)
664 .filter(|subject| !subject.is_empty())
665 .unwrap_or(tool);
666 let raw_input = request
667 .payload
668 .get("raw_args")
669 .cloned()
670 .unwrap_or(Value::Null);
671 Some(json!({
672 "jsonrpc":"2.0",
673 "id":approval_request_id(request.id),
674 "method":"session/request_permission",
675 "params":{
676 "sessionId":session_id,
677 "toolCall":{
678 "toolCallId":format!("supercode-request-{}", request.id),
679 "title":title,
680 "kind":"other",
681 "status":"pending",
682 "rawInput":raw_input,
683 "_meta":{"supercode":{"requestId":request.id, "tool":tool}},
684 },
685 "options":[
686 {"optionId":"allow_once", "name":"Allow once", "kind":"allow_once"},
687 {"optionId":"allow_for_session", "name":"Allow for session", "kind":"allow_always"},
688 {"optionId":"deny", "name":"Deny", "kind":"reject_once"},
689 ],
690 },
691 }))
692}
693
694fn decode_approval_response(message: &Value) -> Option<FrontendResponse> {
695 let request_id = approval_request_id_from_wire(message)?;
696 let result = message.get("result");
697 let error = message.get("error");
698 let exact_top_level = message.as_object().is_some_and(|object| {
699 object.len() == 3
700 && object.contains_key("jsonrpc")
701 && object.contains_key("id")
702 && (object.contains_key("result") ^ object.contains_key("error"))
703 });
704 let valid_envelope = message.get("jsonrpc").and_then(Value::as_str) == Some("2.0")
705 && message.get("method").is_none()
706 && exact_top_level
707 && (result.is_some() ^ error.is_some());
708 let decision =
709 if valid_envelope && error.is_none() {
710 match message.pointer("/result/outcome") {
711 Some(outcome)
712 if result.and_then(Value::as_object).is_some_and(|object| {
713 object.len() == 1 && object.contains_key("outcome")
714 }) && outcome.as_object().is_some_and(|object| {
715 object.len() == 2
716 && object.contains_key("outcome")
717 && object.contains_key("optionId")
718 }) && outcome.get("outcome").and_then(Value::as_str) == Some("selected") =>
719 {
720 match outcome.get("optionId").and_then(Value::as_str) {
721 Some("allow_once") => FrontendApprovalDecision::Allow,
722 Some("allow_for_session") => FrontendApprovalDecision::AllowForSession,
723 _ => FrontendApprovalDecision::Deny,
724 }
725 }
726 _ => FrontendApprovalDecision::Deny,
727 }
728 } else {
729 FrontendApprovalDecision::Deny
730 };
731 Some(FrontendResponse::Approval {
732 request_id,
733 decision,
734 })
735}
736
737fn route_event_projection(
738 tx: &mpsc::UnboundedSender<Value>,
739 session_id: &str,
740 event: &FrontendEvent,
741) -> bool {
742 if tx
743 .send(frontend_event_notification(session_id, event))
744 .is_err()
745 {
746 return false;
747 }
748 if let Some(request) = project_approval_request(session_id, event) {
749 if tx.send(request).is_err() {
750 return false;
751 }
752 }
753 if let Some(update) = project_event(session_id, event) {
754 if tx.send(update).is_err() {
755 return false;
756 }
757 }
758 true
759}
760
761fn project_event(session_id: &str, event: &FrontendEvent) -> Option<Value> {
762 let payload = &event.payload;
763 match payload.get("type").and_then(Value::as_str)? {
764 "text_delta" => Some(projected_session_update(
765 session_id,
766 event,
767 json!({
768 "sessionUpdate": "agent_message_chunk",
769 "content": {"type": "text", "text": payload.get("text")?}
770 }),
771 )),
772 "tool_call_started" => Some(projected_session_update(
773 session_id,
774 event,
775 json!({
776 "sessionUpdate": "tool_call",
777 "toolCallId": payload.get("id")?,
778 "title": payload.get("name")?,
779 "kind": "other",
780 "status": "in_progress",
781 "rawInput": payload.get("arguments").cloned().unwrap_or(Value::Null)
782 }),
783 )),
784 "tool_call_completed" => Some(projected_session_update(
785 session_id,
786 event,
787 json!({
788 "sessionUpdate": "tool_call_update",
789 "toolCallId": payload.get("id")?,
790 "status": if payload.get("is_error").and_then(Value::as_bool).unwrap_or(false) {
791 "failed"
792 } else {
793 "completed"
794 },
795 "content": [{
796 "type": "content",
797 "content": {
798 "type": "text",
799 "text": payload.get("output").cloned().unwrap_or(Value::String(String::new()))
800 }
801 }]
802 }),
803 )),
804 "background_output" => Some(projected_session_update(
805 session_id,
806 event,
807 json!({
808 "sessionUpdate": "agent_message_chunk",
809 "content": {"type": "text", "text": payload.get("chunk")?},
810 "nativeEvent": payload
811 }),
812 )),
813 _ => Some(projected_session_update(
816 session_id,
817 event,
818 json!({
819 "sessionUpdate": "agent_thought_chunk",
820 "content": {"type": "text", "text": ""},
821 "nativeEvent": payload,
822 }),
823 )),
824 }
825}
826
827fn canonical_content_blocks(message: &ChatMessage) -> Vec<Value> {
828 if let Some(content) = &message.content {
829 return if content.is_empty() {
830 Vec::new()
831 } else {
832 vec![json!({"type":"text", "text":content})]
833 };
834 }
835 message
836 .content_parts
837 .as_deref()
838 .unwrap_or_default()
839 .iter()
840 .filter_map(|part| match part.get("type").and_then(Value::as_str) {
841 Some("text") => part
842 .get("text")
843 .and_then(Value::as_str)
844 .map(|text| json!({"type":"text", "text":text})),
845 Some("image_url") => {
846 let url = part.pointer("/image_url/url").and_then(Value::as_str)?;
847 if let Some(data) = url.strip_prefix("data:") {
848 let (mime_type, data) = data.split_once(";base64,")?;
849 Some(json!({
850 "type":"image",
851 "data":data,
852 "mimeType":mime_type,
853 }))
854 } else {
855 Some(json!({
856 "type":"resource_link",
857 "name":"image attachment",
858 "uri":url,
859 }))
860 }
861 }
862 _ => None,
863 })
864 .collect()
865}
866
867fn history_update(
868 session_id: &str,
869 history_cursor: u64,
870 history_index: usize,
871 role: Role,
872 update: Value,
873) -> Value {
874 let mut notification = session_update(session_id, update);
875 notification["params"]["update"]["_meta"] = json!({
876 "supercode": {
877 "historyCursor": history_cursor,
878 "historyIndex": history_index,
879 "canonicalRole": match role {
880 Role::System => "system",
881 Role::User => "user",
882 Role::Assistant => "assistant",
883 Role::Tool => "tool",
884 },
885 }
886 });
887 notification
888}
889
890fn project_history(session_id: &str, history: &[ChatMessage], history_cursor: u64) -> Vec<Value> {
894 let mut projected = Vec::new();
895 for (index, message) in history.iter().enumerate() {
896 match message.role {
897 Role::System => {}
899 Role::User | Role::Assistant => {
900 let session_update_kind = if message.role == Role::User {
901 "user_message_chunk"
902 } else {
903 "agent_message_chunk"
904 };
905 for content in canonical_content_blocks(message) {
906 projected.push(history_update(
907 session_id,
908 history_cursor,
909 index,
910 message.role,
911 json!({
912 "sessionUpdate":session_update_kind,
913 "content":content,
914 }),
915 ));
916 }
917 if message.role == Role::Assistant {
918 for call in message.tool_calls() {
919 let raw_input = serde_json::from_str::<Value>(&call.function.arguments)
920 .unwrap_or_else(|_| Value::String(call.function.arguments.clone()));
921 projected.push(history_update(
922 session_id,
923 history_cursor,
924 index,
925 message.role,
926 json!({
927 "sessionUpdate":"tool_call",
928 "toolCallId":call.id,
929 "title":call.function.name,
930 "kind":"other",
931 "status":"in_progress",
932 "rawInput":raw_input,
933 }),
934 ));
935 }
936 }
937 }
938 Role::Tool => {
939 let Some(tool_call_id) = message.tool_call_id.as_deref() else {
940 continue;
941 };
942 let content = canonical_content_blocks(message)
943 .into_iter()
944 .map(|content| json!({"type":"content", "content":content}))
945 .collect::<Vec<_>>();
946 projected.push(history_update(
947 session_id,
948 history_cursor,
949 index,
950 message.role,
951 json!({
952 "sessionUpdate":"tool_call_update",
953 "toolCallId":tool_call_id,
954 "status":"completed",
955 "content":content,
956 }),
957 ));
958 }
959 }
960 }
961 projected
962}
963
964fn assistant_event_text(event: &FrontendEvent) -> Option<&str> {
965 match event.payload.get("type").and_then(Value::as_str)? {
966 "text_delta" => event.payload.get("text").and_then(Value::as_str),
967 _ => None,
968 }
969}
970
971fn disconnect_message(event: &FrontendEvent) -> Option<&str> {
972 (event.payload.get("type").and_then(Value::as_str) == Some("runtime_disconnected"))
973 .then(|| event.payload.get("message").and_then(Value::as_str))
974 .flatten()
975}
976
977pub async fn run_stdio<R, W>(
984 server: Arc<AcpServer>,
985 mut reader: R,
986 writer: W,
987) -> std::io::Result<()>
988where
989 R: AsyncBufRead + Unpin + Send + 'static,
990 W: AsyncWrite + Unpin + Send + 'static,
991{
992 let (out_tx, mut out_rx) = mpsc::unbounded_channel::<Value>();
993 let pending_approvals = Arc::new(Mutex::new(HashSet::<u64>::new()));
994 let writer_pending_approvals = pending_approvals.clone();
995 let mut writer_task = tokio::spawn(async move {
996 let mut writer = writer;
997 while let Some(value) = out_rx.recv().await {
998 if value.get("method").and_then(Value::as_str) == Some("session/request_permission") {
999 if let Some(request_id) = approval_request_id_from_wire(&value) {
1000 writer_pending_approvals
1001 .lock()
1002 .unwrap_or_else(std::sync::PoisonError::into_inner)
1003 .insert(request_id);
1004 }
1005 }
1006 writer.write_all(format!("{value}\n").as_bytes()).await?;
1007 writer.flush().await?;
1008 }
1009 Ok::<(), std::io::Error>(())
1010 });
1011
1012 let mut events = server.runtime.subscribe();
1013 let event_server = server.clone();
1014 let event_tx = out_tx.clone();
1015 let (router_tx, mut router_rx) = mpsc::unbounded_channel::<RouterCommand>();
1016 let (disconnect_tx, mut disconnect_rx) = mpsc::unbounded_channel::<()>();
1017 let event_task = tokio::spawn(async move {
1018 loop {
1019 tokio::select! {
1020 command = router_rx.recv() => match command {
1021 Some(RouterCommand::Prompt { id, prompt, image_urls, frontend_reply }) => {
1022 let (result, terminal) = project_prompt(
1023 &event_server,
1024 &mut events,
1025 prompt,
1026 image_urls,
1027 frontend_reply,
1028 &event_tx,
1029 ).await;
1030 event_server.prompt_active.store(false, Ordering::SeqCst);
1031 let _ = event_tx.send(sdk_response(id, result));
1032 if terminal {
1033 let _ = disconnect_tx.send(());
1034 break;
1035 }
1036 }
1037 None => break,
1038 },
1039 event = events.recv() => match event {
1040 Ok(event) if disconnect_message(&event).is_some() => {
1041 let _ = disconnect_tx.send(());
1042 break;
1043 }
1044 Ok(event) if event_server.session_open.load(Ordering::SeqCst) => {
1045 if !route_event_projection(
1046 &event_tx,
1047 event_server.runtime.session_id(),
1048 &event,
1049 ) {
1050 break;
1051 }
1052 }
1053 Ok(_) => {}
1054 Err(tokio::sync::broadcast::error::RecvError::Lagged(_))
1055 | Err(tokio::sync::broadcast::error::RecvError::Closed) => {
1056 let _ = disconnect_tx.send(());
1057 break;
1058 }
1059 },
1060 }
1061 }
1062 });
1063
1064 let mut writer_finished = false;
1065 let mut terminal_error = None;
1066 loop {
1067 let mut line = String::new();
1068 let bytes = tokio::select! {
1069 bytes = reader.read_line(&mut line) => bytes?,
1070 _ = disconnect_rx.recv() => break,
1071 result = &mut writer_task => {
1072 writer_finished = true;
1073 terminal_error = match result {
1074 Ok(Ok(())) => None,
1075 Ok(Err(error)) => Some(error),
1076 Err(error) => Some(std::io::Error::other(format!(
1077 "ACP writer task failed: {error}"
1078 ))),
1079 };
1080 break;
1081 }
1082 };
1083 if bytes == 0 {
1084 break;
1085 }
1086 if line.len() > SERVER_MAX_LINE_BYTES {
1087 let _ = out_tx.send(response(
1088 Value::Null,
1089 Err("ACP request exceeded size limit".into()),
1090 ));
1091 continue;
1092 }
1093 let request: Value = match serde_json::from_str(line.trim()) {
1094 Ok(request) => request,
1095 Err(error) => {
1096 let _ = out_tx.send(json!({
1097 "jsonrpc": "2.0", "id": null,
1098 "error": {"code": -32700, "message": error.to_string()}
1099 }));
1100 continue;
1101 }
1102 };
1103 if request.get("method").is_none() {
1104 if let Some(response) = decode_approval_response(&request) {
1105 pending_approvals
1106 .lock()
1107 .unwrap_or_else(std::sync::PoisonError::into_inner)
1108 .remove(&response.request_id());
1109 let _ = server.runtime.runtime.respond(response).await;
1110 }
1111 continue;
1113 }
1114 let method = request.get("method").and_then(Value::as_str).unwrap_or("");
1115 let id = request.get("id").cloned();
1116 let params = request.get("params").cloned().unwrap_or_else(|| json!({}));
1117
1118 if method == "session/cancel" {
1119 let server = server.clone();
1120 tokio::spawn(async move {
1121 let _ = server.runtime.interrupt().await;
1122 });
1123 continue;
1124 }
1125
1126 let Some(id) = id else {
1127 continue;
1128 };
1129 match method {
1130 "initialize" => {
1131 let version = params.get("protocolVersion").and_then(Value::as_u64);
1132 let result = if version == Some(ACP_PROTOCOL_VERSION) {
1133 Ok(server.initialize())
1134 } else {
1135 Err(format!("unsupported ACP protocol version {version:?}"))
1136 };
1137 let _ = out_tx.send(response(id, result));
1138 }
1139 "_goose/unstable/defaults/read"
1140 if server.compatibility == AcpCompatibilityProfile::GooseAcd3c135 =>
1141 {
1142 let _ = out_tx.send(response(id, server.goose_defaults()));
1143 }
1144 "session/new" => {
1145 let mut opened = server
1146 .validate_new_session(¶ms)
1147 .and_then(|()| server.open_session(None));
1148 if opened.is_ok() {
1149 let goose = server.compatibility == AcpCompatibilityProfile::GooseAcd3c135;
1154 let acknowledged_cursor = if goose {
1155 server.runtime.history_cursor
1156 } else {
1157 server.runtime.initial_replay_cursor
1158 };
1159 if let Err(error) = server.runtime.activate_and_route(
1160 server.history_updates_once(),
1161 acknowledged_cursor,
1162 &out_tx,
1163 ) {
1164 server.session_open.store(false, Ordering::SeqCst);
1165 opened = Err(error);
1166 }
1167 }
1168 let _ = out_tx.send(response(id, opened));
1169 }
1170 "session/load" | "session/resume" => {
1171 let requested = params.get("sessionId").and_then(Value::as_str);
1172 let mut opened = server.open_session(requested);
1173 if opened.is_ok() {
1174 if let Err(error) = server.runtime.activate_and_route(
1175 Vec::new(),
1176 server.runtime.initial_replay_cursor,
1177 &out_tx,
1178 ) {
1179 server.session_open.store(false, Ordering::SeqCst);
1180 opened = Err(error);
1181 }
1182 }
1183 let _ = out_tx.send(response(id, opened));
1184 }
1185 "session/prompt" => {
1186 let validation = server
1187 .validate_session(¶ms)
1188 .and_then(|_| AcpServer::prompt_text(¶ms));
1189 let prompt = match validation {
1190 Ok(prompt) => prompt,
1191 Err(error) => {
1192 let _ = out_tx.send(response(id, Err(error)));
1193 continue;
1194 }
1195 };
1196 if server
1197 .prompt_active
1198 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
1199 .is_err()
1200 {
1201 let _ = out_tx.send(sdk_response(
1202 id,
1203 Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)),
1204 ));
1205 continue;
1206 }
1207 if let Err(error) = router_tx.send(RouterCommand::Prompt {
1208 id,
1209 prompt,
1210 image_urls: Vec::new(),
1211 frontend_reply: false,
1212 }) {
1213 server.prompt_active.store(false, Ordering::SeqCst);
1214 let RouterCommand::Prompt { id, .. } = error.0;
1215 let _ = out_tx.send(response(
1216 id,
1217 Err("ACP runtime event router is unavailable".into()),
1218 ));
1219 }
1220 }
1221 "session/steer" | "frontend.v2.steer" => {
1222 let result = match server.validate_session(¶ms) {
1223 Ok(()) => match params.get("text").and_then(Value::as_str) {
1224 Some(text) => server.runtime.runtime.steer(text.to_string()).await,
1225 None => Err(FrontendRuntimeError::InvalidArgument {
1226 operation: crate::SdkOperation::Steer,
1227 message: "session/steer requires string `text`".into(),
1228 }),
1229 },
1230 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1231 operation: crate::SdkOperation::Steer,
1232 message,
1233 }),
1234 };
1235 let _ = out_tx.send(sdk_response(id, result.map(|()| json!({}))));
1236 }
1237 "session/respond" | "frontend.v2.respond" => {
1238 let result = match server.validate_session(¶ms) {
1239 Ok(()) => serde_json::from_value::<FrontendResponse>(
1240 params.get("response").cloned().unwrap_or(Value::Null),
1241 )
1242 .map_err(|error| FrontendRuntimeError::InvalidArgument {
1243 operation: crate::SdkOperation::Respond,
1244 message: error.to_string(),
1245 }),
1246 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1247 operation: crate::SdkOperation::Respond,
1248 message,
1249 }),
1250 };
1251 let result = match result {
1252 Ok(response) => server.runtime.runtime.respond(response).await,
1253 Err(error) => Err(error),
1254 };
1255 let _ = out_tx.send(sdk_response(id, result.map(|()| json!({}))));
1256 }
1257 "frontend.v2.describe" | "supercode/frontend/describe" => {
1258 let result = match server.validate_session(¶ms) {
1259 Ok(()) => server
1260 .runtime
1261 .runtime
1262 .describe()
1263 .await
1264 .and_then(|descriptor| {
1265 serde_json::to_value(descriptor)
1266 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
1267 }),
1268 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1269 operation: crate::SdkOperation::Resume,
1270 message,
1271 }),
1272 };
1273 let _ = out_tx.send(sdk_response(id, result));
1274 }
1275 "frontend.v2.attach" | "supercode/frontend/attach" => {
1276 let limit = params
1277 .get("limit")
1278 .and_then(Value::as_u64)
1279 .unwrap_or(50)
1280 .clamp(1, crate::server::SERVER_HISTORY_CAPACITY as u64)
1281 as usize;
1282 let after = params
1283 .get("after_sequence")
1284 .or_else(|| params.get("afterSequence"))
1285 .and_then(Value::as_u64)
1286 .unwrap_or_default();
1287 let result = match server.validate_session(¶ms) {
1288 Ok(()) => match server.runtime.runtime.attach(limit).await {
1289 Ok(attachment) => {
1290 let mut snapshot = crate::FrontendAttachSnapshot {
1291 descriptor: attachment.descriptor,
1292 history: attachment.history,
1293 history_cursor: attachment.history_cursor,
1294 replay: attachment.replay,
1295 };
1296 snapshot.replay.retain(|event| event.sequence > after);
1297 serde_json::to_value(snapshot)
1298 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
1299 }
1300 Err(error) => Err(error),
1301 },
1302 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1303 operation: crate::SdkOperation::Events,
1304 message,
1305 }),
1306 };
1307 let _ = out_tx.send(sdk_response(id, result));
1308 }
1309 "frontend.v2.send_input" | "supercode/frontend/send_input" => {
1310 let result = match server.validate_session(¶ms) {
1311 Ok(()) => match params.get("prompt").and_then(Value::as_str) {
1312 Some(prompt) => {
1313 match acp_image_urls(¶ms, "supercode/frontend/send_input") {
1314 Ok(image_urls) => {
1315 server
1316 .runtime
1317 .runtime
1318 .clone()
1319 .send_input_with_images(prompt.to_string(), image_urls)
1320 .await
1321 }
1322 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1323 operation: crate::SdkOperation::Input,
1324 message,
1325 }),
1326 }
1327 }
1328 None => Err(FrontendRuntimeError::InvalidArgument {
1329 operation: crate::SdkOperation::Input,
1330 message: "supercode/frontend/send_input requires string `prompt`"
1331 .into(),
1332 }),
1333 },
1334 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1335 operation: crate::SdkOperation::Input,
1336 message,
1337 }),
1338 };
1339 let _ = out_tx.send(sdk_response(id, result.map(|()| json!({"accepted":true}))));
1340 }
1341 "frontend.v2.submit" | "supercode/frontend/submit" => {
1342 let prompt = match server.validate_session(¶ms).and_then(|_| {
1343 params
1344 .get("prompt")
1345 .and_then(Value::as_str)
1346 .map(str::to_owned)
1347 .ok_or_else(|| {
1348 "supercode/frontend/submit requires string `prompt`".to_string()
1349 })
1350 }) {
1351 Ok(prompt) => prompt,
1352 Err(error) => {
1353 let _ = out_tx.send(sdk_response(
1354 id,
1355 Err(FrontendRuntimeError::InvalidArgument {
1356 operation: crate::SdkOperation::Input,
1357 message: error,
1358 }),
1359 ));
1360 continue;
1361 }
1362 };
1363 let image_urls = match params.get("image_urls") {
1364 None => Vec::new(),
1365 Some(Value::Array(values)) => {
1366 match values.iter().map(Value::as_str).collect::<Option<Vec<_>>>() {
1367 Some(values) => values.into_iter().map(str::to_owned).collect(),
1368 None => {
1369 let _ = out_tx.send(sdk_response(
1370 id,
1371 Err(FrontendRuntimeError::InvalidArgument {
1372 operation: crate::SdkOperation::Input,
1373 message: "supercode/frontend/submit requires string entries in `image_urls`".into(),
1374 }),
1375 ));
1376 continue;
1377 }
1378 }
1379 }
1380 Some(_) => {
1381 let _ = out_tx.send(sdk_response(
1382 id,
1383 Err(FrontendRuntimeError::InvalidArgument {
1384 operation: crate::SdkOperation::Input,
1385 message: "supercode/frontend/submit requires array `image_urls`"
1386 .into(),
1387 }),
1388 ));
1389 continue;
1390 }
1391 };
1392 if server
1393 .prompt_active
1394 .compare_exchange(false, true, Ordering::SeqCst, Ordering::SeqCst)
1395 .is_err()
1396 {
1397 let _ = out_tx.send(sdk_response(
1398 id,
1399 Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)),
1400 ));
1401 continue;
1402 }
1403 if let Err(error) = router_tx.send(RouterCommand::Prompt {
1404 id,
1405 prompt,
1406 image_urls,
1407 frontend_reply: true,
1408 }) {
1409 server.prompt_active.store(false, Ordering::SeqCst);
1410 let RouterCommand::Prompt { id, .. } = error.0;
1411 let _ = out_tx.send(sdk_response(
1412 id,
1413 Err(FrontendRuntimeError::Transport(
1414 "ACP runtime event router is unavailable".into(),
1415 )),
1416 ));
1417 }
1418 }
1419 "frontend.v2.interrupt" | "supercode/frontend/interrupt" => {
1420 let result = match server.validate_session(¶ms) {
1421 Ok(()) => server
1422 .runtime
1423 .runtime
1424 .interrupt()
1425 .await
1426 .map(|interrupted| json!({"interrupted":interrupted})),
1427 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1428 operation: crate::SdkOperation::Interrupt,
1429 message,
1430 }),
1431 };
1432 let _ = out_tx.send(sdk_response(id, result));
1433 }
1434 "frontend.v2.invoke" | "supercode/frontend/invoke" => {
1435 let operation = params
1436 .get("operation")
1437 .cloned()
1438 .ok_or_else(|| FrontendRuntimeError::InvalidArgument {
1439 operation: crate::SdkOperation::Input,
1440 message: "supercode/frontend/invoke requires `operation`".into(),
1441 })
1442 .and_then(|value| {
1443 serde_json::from_value(value).map_err(|error| {
1444 FrontendRuntimeError::InvalidArgument {
1445 operation: crate::SdkOperation::Input,
1446 message: error.to_string(),
1447 }
1448 })
1449 });
1450 let result = match (server.validate_session(¶ms), operation) {
1451 (Ok(()), Ok(operation)) => server.runtime.runtime.invoke(operation).await,
1452 (Err(message), _) => Err(FrontendRuntimeError::InvalidArgument {
1453 operation: crate::SdkOperation::Input,
1454 message,
1455 }),
1456 (_, Err(error)) => Err(error),
1457 };
1458 let result = result.and_then(|result| {
1459 serde_json::to_value(result)
1460 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
1461 });
1462 let _ = out_tx.send(sdk_response(id, result));
1463 }
1464 "frontend.v2.lease" | "supercode/frontend/lease" => {
1465 let result = match server.validate_session(¶ms) {
1466 Ok(()) => server.runtime.runtime.lease_snapshot().await,
1467 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1468 operation: crate::SdkOperation::Events,
1469 message,
1470 }),
1471 }
1472 .and_then(|snapshot| {
1473 serde_json::to_value(snapshot)
1474 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
1475 });
1476 let _ = out_tx.send(sdk_response(id, result));
1477 }
1478 "frontend.v2.take_control" | "supercode/frontend/take_control" => {
1479 let result = match server.validate_session(¶ms) {
1480 Ok(()) => server.runtime.runtime.take_control().await,
1481 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1482 operation: crate::SdkOperation::Input,
1483 message,
1484 }),
1485 }
1486 .and_then(|snapshot| {
1487 serde_json::to_value(snapshot)
1488 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
1489 });
1490 let _ = out_tx.send(sdk_response(id, result));
1491 }
1492 "frontend.v2.heartbeat" | "supercode/frontend/heartbeat" => {
1493 let result = match server.validate_session(¶ms) {
1494 Ok(()) => server.runtime.runtime.heartbeat().await,
1495 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1496 operation: crate::SdkOperation::Events,
1497 message,
1498 }),
1499 }
1500 .and_then(|snapshot| {
1501 serde_json::to_value(snapshot)
1502 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
1503 });
1504 let _ = out_tx.send(sdk_response(id, result));
1505 }
1506 "frontend.v2.detach" | "supercode/frontend/detach" => {
1507 let result = match server.validate_session(¶ms) {
1508 Ok(()) => server.runtime.runtime.detach().await,
1509 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1510 operation: crate::SdkOperation::Events,
1511 message,
1512 }),
1513 }
1514 .and_then(|snapshot| {
1515 serde_json::to_value(snapshot)
1516 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
1517 });
1518 if result.is_ok() {
1519 server.session_open.store(false, Ordering::SeqCst);
1520 }
1521 let _ = out_tx.send(sdk_response(id, result));
1522 }
1523 "frontend.v2.close" | "supercode/frontend/close" => {
1524 let result = match server.validate_session(¶ms) {
1525 Ok(()) => server
1526 .runtime
1527 .runtime
1528 .close()
1529 .await
1530 .map(|()| json!({"closed":true})),
1531 Err(message) => Err(FrontendRuntimeError::InvalidArgument {
1532 operation: crate::SdkOperation::Close,
1533 message,
1534 }),
1535 };
1536 let _ = out_tx.send(sdk_response(id, result));
1537 }
1538 other => {
1539 let _ = out_tx.send(json!({
1540 "jsonrpc": "2.0", "id": id,
1541 "error": {"code": -32601, "message": format!("unknown ACP method `{other}`")}
1542 }));
1543 }
1544 }
1545 }
1546
1547 drop(out_tx);
1548 event_task.abort();
1549 let _ = event_task.await;
1550 if !writer_finished {
1551 terminal_error = match writer_task.await {
1552 Ok(Ok(())) => terminal_error,
1553 Ok(Err(error)) => Some(error),
1554 Err(error) => Some(std::io::Error::other(format!(
1555 "ACP writer task failed: {error}"
1556 ))),
1557 };
1558 }
1559 let unresolved = pending_approvals
1563 .lock()
1564 .unwrap_or_else(std::sync::PoisonError::into_inner)
1565 .drain()
1566 .collect::<Vec<_>>();
1567 for request_id in unresolved {
1568 let _ = server
1569 .runtime
1570 .runtime
1571 .respond(FrontendResponse::Approval {
1572 request_id,
1573 decision: FrontendApprovalDecision::Deny,
1574 })
1575 .await;
1576 }
1577 match terminal_error {
1578 Some(error) => Err(error),
1579 None => Ok(()),
1580 }
1581}
1582
1583fn acp_image_urls(params: &Value, operation: &str) -> Result<Vec<String>, String> {
1584 let urls = match params.get("image_urls") {
1585 None => Ok(Vec::new()),
1586 Some(Value::Array(values)) => values
1587 .iter()
1588 .map(|value| {
1589 value
1590 .as_str()
1591 .map(str::to_owned)
1592 .ok_or_else(|| format!("{operation} requires string entries in `image_urls`"))
1593 })
1594 .collect(),
1595 Some(_) => Err(format!("{operation} requires array `image_urls`")),
1596 }?;
1597 if urls.len() > 4 {
1598 return Err(format!("{operation} accepts at most 4 images"));
1599 }
1600 let mut total = 0usize;
1601 for url in &urls {
1602 if !(url.starts_with("data:image/")
1603 || url.starts_with("https://")
1604 || url.starts_with("http://"))
1605 {
1606 return Err(format!(
1607 "{operation} images must be image data URLs or HTTP(S) URLs"
1608 ));
1609 }
1610 if url.len() > 12 * 1024 * 1024 {
1611 return Err(format!("{operation} image exceeds the encoded size limit"));
1612 }
1613 total = total.saturating_add(url.len());
1614 }
1615 if total > 32 * 1024 * 1024 {
1616 return Err(format!(
1617 "{operation} images exceed the encoded total size limit"
1618 ));
1619 }
1620 Ok(urls)
1621}
1622
1623#[cfg(test)]
1624mod tests {
1625 use super::*;
1626
1627 #[test]
1628 fn projects_text_and_tool_events_without_transport_state() {
1629 let sdk_event = FrontendEvent {
1630 sequence: 41,
1631 kind: "text_delta".into(),
1632 payload: json!({"type": "text_delta", "text": "hello"}),
1633 };
1634 let text = project_event("s1", &sdk_event).unwrap();
1635 assert_eq!(
1636 text.pointer("/params/update/sessionUpdate")
1637 .and_then(Value::as_str),
1638 Some("agent_message_chunk")
1639 );
1640 assert_eq!(
1641 text.pointer("/params/update/content/text")
1642 .and_then(Value::as_str),
1643 Some("hello")
1644 );
1645 assert_eq!(
1646 text.pointer("/params/update/_meta/supercode/sdkSequence")
1647 .and_then(Value::as_u64),
1648 Some(41)
1649 );
1650 assert_eq!(
1651 text.pointer("/params/update/_meta/supercode/sdkKind")
1652 .and_then(Value::as_str),
1653 Some("text_delta")
1654 );
1655
1656 let tool = project_event(
1657 "s1",
1658 &FrontendEvent {
1659 sequence: 42,
1660 kind: "tool_call_started".into(),
1661 payload: json!({
1662 "type": "tool_call_started", "id": "call-1",
1663 "name": "read_file", "arguments": "{}"
1664 }),
1665 },
1666 )
1667 .unwrap();
1668 assert_eq!(
1669 tool.pointer("/params/update/toolCallId")
1670 .and_then(Value::as_str),
1671 Some("call-1")
1672 );
1673 }
1674
1675 #[test]
1676 fn preserves_named_sdk_errors_in_acp_responses() {
1677 let response = sdk_response(
1678 json!(7),
1679 Err(FrontendRuntimeError::Submit(RuntimeSubmitError::Busy)),
1680 );
1681 assert_eq!(response["error"]["name"], "busy");
1682 assert_eq!(response["error"]["operation"], "input");
1683 }
1684
1685 #[test]
1686 fn approval_requests_and_responses_preserve_runtime_correlation() {
1687 let event = FrontendEvent {
1688 sequence: 43,
1689 kind: "request".into(),
1690 payload: json!({
1691 "type":"request",
1692 "request":{
1693 "id":17,
1694 "kind":"approval",
1695 "payload":{
1696 "tool":"bash",
1697 "subject":"echo reviewed",
1698 "raw_args":{"command":"echo reviewed"},
1699 },
1700 },
1701 }),
1702 };
1703 let request = project_approval_request("s1", &event).unwrap();
1704 assert_eq!(request["id"], "supercode-approval-17");
1705 assert_eq!(request["method"], "session/request_permission");
1706 assert_eq!(request["params"]["sessionId"], "s1");
1707 assert_eq!(request["params"]["toolCall"]["title"], "echo reviewed");
1708 assert_eq!(request["params"]["options"][0]["optionId"], "allow_once");
1709
1710 assert_eq!(
1711 decode_approval_response(&json!({
1712 "jsonrpc":"2.0",
1713 "id":"supercode-approval-17",
1714 "result":{"outcome":{"outcome":"selected", "optionId":"allow_for_session"}},
1715 })),
1716 Some(FrontendResponse::Approval {
1717 request_id: 17,
1718 decision: FrontendApprovalDecision::AllowForSession,
1719 })
1720 );
1721 assert_eq!(
1722 decode_approval_response(&json!({
1723 "jsonrpc":"2.0",
1724 "id":"supercode-approval-17",
1725 "result":{"outcome":{"outcome":"cancelled"}},
1726 })),
1727 Some(FrontendResponse::Approval {
1728 request_id: 17,
1729 decision: FrontendApprovalDecision::Deny,
1730 })
1731 );
1732 for malformed in [
1733 json!({
1734 "id":"supercode-approval-17",
1735 "result":{"outcome":{"outcome":"selected", "optionId":"allow_once"}},
1736 }),
1737 json!({
1738 "jsonrpc":"2.0",
1739 "id":"supercode-approval-17",
1740 "result":{"outcome":{"outcome":"selected", "optionId":"allow_once"}},
1741 "error":{"code":-32603, "message":"invalid"},
1742 }),
1743 json!({
1744 "jsonrpc":"2.0",
1745 "id":"supercode-approval-17",
1746 "result":{"outcome":{"outcome":"selected", "optionId":"unknown"}},
1747 }),
1748 json!({
1749 "jsonrpc":"2.0",
1750 "id":"supercode-approval-17",
1751 "result":{"outcome":{"outcome":"selected", "optionId":"allow_once"}},
1752 "params":{},
1753 }),
1754 json!({
1755 "jsonrpc":"2.0",
1756 "id":"supercode-approval-17",
1757 "result":{
1758 "outcome":{"outcome":"selected", "optionId":"allow_once"},
1759 "extra":true,
1760 },
1761 }),
1762 ] {
1763 assert_eq!(
1764 decode_approval_response(&malformed),
1765 Some(FrontendResponse::Approval {
1766 request_id: 17,
1767 decision: FrontendApprovalDecision::Deny,
1768 })
1769 );
1770 }
1771 assert!(decode_approval_response(&json!({
1772 "jsonrpc":"2.0",
1773 "id":"supercode-approval-017",
1774 "result":{"outcome":{"outcome":"selected", "optionId":"allow_once"}},
1775 }))
1776 .is_none());
1777 assert!(decode_approval_response(&json!({"jsonrpc":"2.0", "id":99})).is_none());
1778 }
1779
1780 #[test]
1781 fn inactive_event_routing_is_bounded_and_fails_instead_of_truncating() {
1782 let mut routing = EventRoutingState::default();
1783 for sequence in 1..=FRONTEND_REPLAY_CAPACITY as u64 {
1784 routing
1785 .buffer(FrontendEvent {
1786 sequence,
1787 kind: "loop_tick".into(),
1788 payload: json!({"type":"loop_tick", "sequence":sequence}),
1789 })
1790 .unwrap();
1791 }
1792 assert_eq!(routing.pending.len(), FRONTEND_REPLAY_CAPACITY);
1793 assert!(routing
1794 .buffer(FrontendEvent {
1795 sequence: FRONTEND_REPLAY_CAPACITY as u64 + 1,
1796 kind: "loop_tick".into(),
1797 payload: json!({"type":"loop_tick"}),
1798 })
1799 .is_err());
1800 assert!(routing.pending.is_empty(), "failed replay is never partial");
1801
1802 let error = routing.activate(0).unwrap_err();
1803 assert!(error.contains(&FRONTEND_REPLAY_CAPACITY.to_string()));
1804 assert!(error.contains("reconnect required"));
1805 assert!(!routing.active, "a gapped attachment must never activate");
1806 }
1807}