Skip to main content

supercode_harness/
acp_server.rs

1//! Inbound Agent Client Protocol adapter for the SDK-owned Supercode runtime.
2//!
3//! [`crate::runtime::AcpRuntimeBackend`] is the outbound half: it makes an ACP
4//! agent look like a protocol-neutral SDK [`crate::RuntimeConnection`].  This
5//! module is the deliberately symmetric inbound half: it makes one canonical
6//! [`crate::server::RpcEngine`] look like an ACP agent.  The ACP adapter owns no
7//! model loop, transcript, scheduler, or persistence state; it only translates
8//! wire messages and subscribes to the runtime's canonical event stream.
9
10use 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
26/// ACP protocol version implemented by this adapter.
27pub const ACP_PROTOCOL_VERSION: u64 = 1;
28
29/// Versioned stock-client compatibility behavior layered over standard ACP.
30///
31/// The default remains donor-neutral ACP. A pinned profile may answer only
32/// the client-local extension calls and ordering quirks proven by its stock
33/// traffic corpus; it never changes the SDK runtime's model, tools, provider,
34/// persistence, or session identity.
35#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
36pub enum AcpCompatibilityProfile {
37    /// Standard ACP with Supercode's typed frontend extensions.
38    #[default]
39    Standard,
40    /// Stock Goose text UI at source commit `acd3c135cff7...`.
41    GooseAcd3c135,
42}
43
44/// Compatibility name for the HTTP implementation of the canonical SDK
45/// runtime contract. ACP no longer owns a separate HTTP client.
46#[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    /// Atomically publish the history/replay/live boundary. The attachment
162    /// reader buffers every event newer than `history_cursor` until this
163    /// method has queued the history prefix and finite replay. Producers
164    /// cannot overtake the boundary because they share `routing`.
165    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        // Wake the connection lifecycle without projecting the terminal a
192        // second time; the event task treats disconnect as control only.
193        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
244/// Inbound ACP projection of one SDK runtime.
245pub 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    // A runtime approval handler may synchronously wait for the frontend's
272    // answer. Keep that wait off this event-routing task so the outbound ACP
273    // permission request can be delivered and its response read concurrently.
274    let mut submit = tokio::spawn(async move { runtime.submit(prompt, image_urls).await });
275    let mut streamed = String::new();
276    // the assistant text since the last tool call: a runtime's reply for a
277    // turn that used tools is its last message, not everything it said
278    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    // The SDK runtime publishes every event for the reply before submit
325    // before submit resolves. This task is the sole subscriber, so draining
326    // its queue is an ordering barrier: no assistant chunk can arrive after
327    // the end_turn response.
328    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    /// Wrap an existing SDK runtime.  Multiple frontends may wrap/subscribe to
419    /// the same runtime; the runtime, not this adapter, remains authoritative.
420    pub async fn new(runtime: Arc<dyn SdkRuntime>) -> Result<Arc<Self>, FrontendRuntimeError> {
421        Self::new_with_profile(runtime, AcpCompatibilityProfile::Standard).await
422    }
423
424    /// Wrap an existing SDK runtime with one version-pinned stock-client
425    /// compatibility profile. The profile affects only ACP projection; the
426    /// supplied runtime remains the sole owner of continuation and state.
427    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    /// Canonical runtime shared by this ACP projection and other frontends.
441    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        // Usage/cache/round-trip completion remain available on the SDK event
833        // bus but have no portable ACP v1 update shape.
834        _ => 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
909/// Project the bounded canonical attachment snapshot into standard ACP
910/// history updates. These messages are a display replay only: the canonical
911/// SDK history is neither copied into nor reconstructed from the client.
912fn 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            // System instructions are runtime authority, not transcript UI.
917            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
1000/// Serve ACP v1 JSON-RPC over a newline-delimited stream.
1001///
1002/// Requests are dispatched concurrently so `session/cancel` can interrupt a
1003/// running `session/prompt`.  A single writer task serializes responses and
1004/// notifications, preventing byte interleaving between event and request
1005/// tasks.
1006pub 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            // JSON-RPC responses never receive a response of their own.
1135            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(&params)
1170                    .and_then(|()| server.open_session(None));
1171                if opened.is_ok() {
1172                    // The pinned stock Goose server sends session updates
1173                    // before resolving session/new, and its client accepts
1174                    // that ordering. Replaying at this exact boundary avoids
1175                    // inventing a client-side session-selection protocol.
1176                    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(&params)
1211                    .and_then(|_| AcpServer::prompt_text(&params));
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(&params) {
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(&params) {
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(&params) {
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(&params) {
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(&params) {
1334                    Ok(()) => match params.get("prompt").and_then(Value::as_str) {
1335                        Some(prompt) => {
1336                            match acp_image_urls(&params, "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(&params).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(&params) {
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(&params), 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(&params) {
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(&params) {
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(&params) {
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(&params) {
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(&params) {
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(&params) {
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    // The writer is now quiescent, so its projected-request inventory cannot
1597    // race this drain. Any unanswered approval belonged to this dead stdio
1598    // attachment and must resolve as Deny rather than block or execute.
1599    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}