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