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