Skip to main content

supercode_harness/
acp_frontend.rs

1//! First-class ACP client implementation of the canonical SDK runtime.
2//!
3//! This adapter joins an already-running SDK-owned runtime through an ACP
4//! bridge process. The bridge is a transport lease only: dropping or
5//! detaching this client never owns or stops the remote runtime.
6
7use std::collections::BTreeSet;
8use std::path::PathBuf;
9use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
10use std::sync::Arc;
11
12use async_trait::async_trait;
13use serde::{Deserialize, Serialize};
14use serde_json::{json, Value};
15use tokio::sync::broadcast;
16
17use crate::frontend::{
18    FrontendAttachSnapshot, FrontendAttachment, FrontendOperationInvocation,
19    FrontendOperationResult, FrontendResponse, FrontendRuntimeDescriptor,
20};
21use crate::runtime::{JsonLineClient, RuntimeEndpoint, RuntimeLaunch};
22use crate::{SdkError, SdkOperation, SdkRuntime};
23
24const FRONTEND_EXTENSION: &str = "/agentCapabilities/_meta/supercode/frontend";
25
26/// Connection parameters for one ACP frontend attachment.
27#[derive(Debug, Clone)]
28pub struct AcpFrontendConnectOptions {
29    /// Command that exposes an already-running runtime as ACP over stdio.
30    /// A typical value is `supercode acp --connect <url>` plus its token.
31    pub launch: RuntimeLaunch,
32    /// Working directory for the bridge process.
33    pub cwd: Option<PathBuf>,
34    /// Stable session to resume. When absent, ACP discovery opens the single
35    /// runtime represented by the endpoint and adopts the returned id.
36    pub session_id: Option<String>,
37    /// Last canonical event fully processed by an earlier client instance.
38    pub after_sequence: Option<u64>,
39}
40
41/// Durable delivery cursor for restarting an ACP frontend without replaying
42/// events it already consumed.
43#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
44pub struct AcpFrontendCheckpoint {
45    /// Checkpoint schema version.
46    pub schema_version: u32,
47    /// Stable SDK runtime/session identity.
48    pub session_id: String,
49    /// Highest canonical event delivered through the frontend attachment.
50    pub acknowledged_sequence: u64,
51}
52
53/// ACP transport implementing the same [`SdkRuntime`] consumed by local and
54/// HTTP frontends.
55pub struct AcpFrontendRuntime {
56    client: Arc<JsonLineClient>,
57    events: broadcast::Sender<crate::SdkEvent>,
58    session_id: String,
59    methods: BTreeSet<String>,
60    event_method: String,
61    extensions: Value,
62    acknowledged_sequence: Arc<AtomicU64>,
63    last_received_sequence: Arc<AtomicU64>,
64    endpoint: RuntimeEndpoint,
65    disconnected: Arc<AtomicBool>,
66}
67
68impl AcpFrontendRuntime {
69    /// Negotiate ACP, discover or resume one session, and start the canonical
70    /// live-event subscription.
71    pub async fn connect(options: AcpFrontendConnectOptions) -> Result<Arc<Self>, SdkError> {
72        let (client, mut receiver, endpoint) = JsonLineClient::spawn(
73            &options.launch,
74            options.cwd.as_deref(),
75            true,
76            "acp-v1-jsonrpc",
77        )
78        .await
79        .map_err(transport)?;
80        let initialized = client
81            .request(
82                "initialize",
83                json!({
84                    "protocolVersion": crate::acp_server::ACP_PROTOCOL_VERSION,
85                    "clientCapabilities": {
86                        "_meta": {"supercode": {"frontend": {"schemaVersion": 2}}}
87                    },
88                    "clientInfo": {
89                        "name": "supercode",
90                        "title": "Supercode frontend",
91                        "version": env!("CARGO_PKG_VERSION")
92                    }
93                }),
94            )
95            .await
96            .map_err(transport)?;
97        if initialized.get("protocolVersion").and_then(Value::as_u64)
98            != Some(crate::acp_server::ACP_PROTOCOL_VERSION)
99        {
100            return Err(SdkError::Transport(
101                "ACP frontend negotiated an unsupported protocol version".into(),
102            ));
103        }
104        let extensions = initialized
105            .pointer(FRONTEND_EXTENSION)
106            .cloned()
107            .ok_or_else(|| {
108                SdkError::Transport(
109                    "ACP peer does not advertise the Supercode frontend extension".into(),
110                )
111            })?;
112        if extensions
113            .get("runtimeOwnedByClient")
114            .and_then(Value::as_bool)
115            != Some(false)
116        {
117            return Err(SdkError::Transport(
118                "ACP peer does not guarantee non-owning frontend detach semantics".into(),
119            ));
120        }
121        let methods = extensions
122            .get("methods")
123            .and_then(Value::as_array)
124            .into_iter()
125            .flatten()
126            .filter_map(Value::as_str)
127            .map(str::to_owned)
128            .collect::<BTreeSet<_>>();
129        let event_method = extensions
130            .get("eventMethod")
131            .and_then(Value::as_str)
132            .unwrap_or("supercode/frontend/event")
133            .to_string();
134        for (required, legacy) in [
135            (
136                crate::FrontendFacadeMethod::Describe,
137                "supercode/frontend/describe",
138            ),
139            (
140                crate::FrontendFacadeMethod::Attach,
141                "supercode/frontend/attach",
142            ),
143        ] {
144            if !supports_method(&methods, required, legacy) {
145                return Err(SdkError::UnsupportedAction(
146                    SdkOperation::Events.action_name(),
147                ));
148            }
149        }
150
151        let requested_session = options.session_id.as_deref();
152        let (open_method, open_params) = if let Some(session_id) = requested_session {
153            let resume = initialized
154                .pointer("/agentCapabilities/sessionCapabilities/resume")
155                .is_some();
156            let load = initialized
157                .pointer("/agentCapabilities/loadSession")
158                .and_then(Value::as_bool)
159                .unwrap_or(false);
160            let method = if resume {
161                "session/resume"
162            } else if load {
163                "session/load"
164            } else {
165                return Err(SdkError::UnsupportedAction(
166                    SdkOperation::Resume.action_name(),
167                ));
168            };
169            (method, json!({"sessionId":session_id}))
170        } else {
171            ("session/new", json!({}))
172        };
173        let opened = client
174            .request(open_method, open_params)
175            .await
176            .map_err(transport)?;
177        let session_id = opened
178            .get("sessionId")
179            .and_then(Value::as_str)
180            .ok_or_else(|| SdkError::Transport("ACP session open omitted sessionId".into()))?
181            .to_string();
182        if requested_session.is_some_and(|requested| requested != session_id) {
183            return Err(SdkError::NotFound {
184                operation: SdkOperation::Resume,
185                message: format!(
186                    "ACP opened `{session_id}` instead of `{}`",
187                    requested_session.unwrap()
188                ),
189            });
190        }
191
192        let (events, _) = broadcast::channel(1024);
193        let disconnected = Arc::new(AtomicBool::new(false));
194        let event_sender = events.clone();
195        let event_disconnected = disconnected.clone();
196        let event_session_id = session_id.clone();
197        let negotiated_event_method = event_method.clone();
198        let last_seen = Arc::new(AtomicU64::new(options.after_sequence.unwrap_or_default()));
199        let last_received_sequence = last_seen.clone();
200        tokio::spawn(async move {
201            while let Some(message) = receiver.recv().await {
202                if message.get("method").and_then(Value::as_str)
203                    == Some(negotiated_event_method.as_str())
204                    && message.pointer("/params/sessionId").and_then(Value::as_str)
205                        == Some(event_session_id.as_str())
206                {
207                    if let Ok(event) = serde_json::from_value::<crate::SdkEvent>(
208                        message
209                            .pointer("/params/event")
210                            .cloned()
211                            .unwrap_or(Value::Null),
212                    ) {
213                        last_seen.fetch_max(event.sequence, Ordering::SeqCst);
214                        let _ = event_sender.send(event);
215                    }
216                    continue;
217                }
218                let kind = message.get("type").and_then(Value::as_str);
219                if matches!(kind, Some("transport_closed" | "transport_error")) {
220                    event_disconnected.store(true, Ordering::SeqCst);
221                    let sequence = last_seen.fetch_add(1, Ordering::SeqCst) + 1;
222                    let _ = event_sender.send(crate::SdkEvent {
223                        sequence,
224                        kind: "runtime_disconnected".into(),
225                        payload: json!({
226                            "type":"runtime_disconnected",
227                            "message": message.get("message").cloned().unwrap_or_else(|| json!("ACP transport closed")),
228                            "_meta":{"supercode":{"transient":true,"transport":"acp"}}
229                        }),
230                    });
231                    break;
232                }
233            }
234        });
235
236        let runtime = Arc::new(Self {
237            client,
238            events,
239            session_id,
240            methods,
241            event_method,
242            extensions,
243            acknowledged_sequence: Arc::new(AtomicU64::new(
244                options.after_sequence.unwrap_or_default(),
245            )),
246            last_received_sequence,
247            endpoint,
248            disconnected,
249        });
250        // Validate the typed descriptor before exposing the runtime.
251        let _ = runtime.describe().await?;
252        Ok(runtime)
253    }
254
255    /// Stable session discovered or resumed during negotiation.
256    pub fn session_id(&self) -> &str {
257        &self.session_id
258    }
259
260    /// Concrete local bridge endpoint, useful for process supervision.
261    pub fn endpoint(&self) -> &RuntimeEndpoint {
262        &self.endpoint
263    }
264
265    /// Complete negotiated extension object. Unknown fields are retained so
266    /// newer peers can be proxied or inspected without lossy decoding.
267    pub fn extensions(&self) -> &Value {
268        &self.extensions
269    }
270
271    /// Negotiated ACP notification carrying canonical facade events.
272    pub fn event_method(&self) -> &str {
273        &self.event_method
274    }
275
276    /// Restore a delivery cursor after negotiation and before attachment.
277    pub fn restore_checkpoint(&self, checkpoint: AcpFrontendCheckpoint) -> Result<(), SdkError> {
278        validate_checkpoint(&checkpoint, &self.session_id)?;
279        self.acknowledged_sequence
280            .fetch_max(checkpoint.acknowledged_sequence, Ordering::SeqCst);
281        self.last_received_sequence
282            .fetch_max(checkpoint.acknowledged_sequence, Ordering::SeqCst);
283        Ok(())
284    }
285
286    /// Snapshot the highest event actually delivered through an attachment.
287    pub fn checkpoint(&self) -> AcpFrontendCheckpoint {
288        AcpFrontendCheckpoint {
289            schema_version: 1,
290            session_id: self.session_id.clone(),
291            acknowledged_sequence: self.acknowledged_sequence.load(Ordering::SeqCst),
292        }
293    }
294
295    /// Detach this transport lease without stopping the SDK-owned runtime.
296    pub async fn detach(&self) -> Result<(), SdkError> {
297        if self.disconnected.load(Ordering::SeqCst) {
298            return self.client.close().await.map_err(transport);
299        }
300        if supports_method(
301            &self.methods,
302            crate::FrontendFacadeMethod::Detach,
303            "supercode/frontend/detach",
304        ) {
305            let _ = SdkRuntime::detach(self).await?;
306        }
307        self.client.close().await.map_err(transport)
308    }
309
310    fn method(
311        &self,
312        method: crate::FrontendFacadeMethod,
313        legacy: &'static str,
314        operation: SdkOperation,
315    ) -> Result<&'static str, SdkError> {
316        if self.methods.contains(method.wire_name()) {
317            Ok(method.wire_name())
318        } else if self.methods.contains(legacy) {
319            Ok(legacy)
320        } else {
321            Err(SdkError::UnsupportedAction(operation.action_name()))
322        }
323    }
324
325    async fn request(
326        &self,
327        method: &str,
328        params: Value,
329        operation: SdkOperation,
330    ) -> Result<Value, SdkError> {
331        let (_id, response) = self
332            .client
333            .begin_request(method, params)
334            .await
335            .map_err(transport)?;
336        match response.await {
337            Ok(Ok(value)) => Ok(value),
338            Ok(Err(error)) => Err(decode_sdk_error(&error, operation)),
339            Err(_) => Err(SdkError::Closed),
340        }
341    }
342
343    fn mask_descriptor(
344        &self,
345        mut descriptor: FrontendRuntimeDescriptor,
346    ) -> FrontendRuntimeDescriptor {
347        mask_actions(&self.methods, &mut descriptor.actions);
348        descriptor
349    }
350}
351
352fn supports_method(
353    methods: &BTreeSet<String>,
354    method: crate::FrontendFacadeMethod,
355    legacy: &str,
356) -> bool {
357    methods.contains(method.wire_name()) || methods.contains(legacy)
358}
359
360fn validate_checkpoint(
361    checkpoint: &AcpFrontendCheckpoint,
362    session_id: &str,
363) -> Result<(), SdkError> {
364    if checkpoint.schema_version != 1 {
365        return Err(SdkError::Transport(format!(
366            "unsupported ACP frontend checkpoint schema {}",
367            checkpoint.schema_version
368        )));
369    }
370    if checkpoint.session_id != session_id {
371        return Err(SdkError::Transport(format!(
372            "ACP frontend checkpoint is for `{}`, not `{session_id}`",
373            checkpoint.session_id
374        )));
375    }
376    Ok(())
377}
378
379fn mask_actions(methods: &BTreeSet<String>, actions: &mut crate::FrontendActions) {
380    actions.submit &= supports_method(
381        methods,
382        crate::FrontendFacadeMethod::Submit,
383        "supercode/frontend/submit",
384    ) && supports_method(
385        methods,
386        crate::FrontendFacadeMethod::SendInput,
387        "supercode/frontend/send_input",
388    );
389    actions.interrupt &= supports_method(
390        methods,
391        crate::FrontendFacadeMethod::Interrupt,
392        "supercode/frontend/interrupt",
393    );
394    actions.steer &= supports_method(methods, crate::FrontendFacadeMethod::Steer, "session/steer");
395    actions.respond &= supports_method(
396        methods,
397        crate::FrontendFacadeMethod::Respond,
398        "session/respond",
399    );
400    actions.detach &= supports_method(
401        methods,
402        crate::FrontendFacadeMethod::Detach,
403        "supercode/frontend/detach",
404    );
405    // A future peer may advertise an explicit close route. This adapter
406    // never guesses one from process ownership or transport EOF.
407    actions.close &= supports_method(
408        methods,
409        crate::FrontendFacadeMethod::Close,
410        "supercode/frontend/close",
411    );
412}
413
414fn suppress_acknowledged_history(
415    history: &mut Vec<crate::ChatMessage>,
416    history_cursor: u64,
417    acknowledged_sequence: u64,
418) {
419    // Canonical history is an atomic projection through `history_cursor`.
420    // Once a reconnecting frontend has acknowledged that boundary, rendering
421    // the projection again would duplicate already displayed messages. A
422    // cursor inside the projection cannot be mapped to individual messages,
423    // so retain the complete history in that case rather than risk loss.
424    if acknowledged_sequence > 0 && acknowledged_sequence >= history_cursor {
425        history.clear();
426    }
427}
428
429#[async_trait]
430impl SdkRuntime for AcpFrontendRuntime {
431    async fn describe(&self) -> Result<FrontendRuntimeDescriptor, SdkError> {
432        let method = self.method(
433            crate::FrontendFacadeMethod::Describe,
434            "supercode/frontend/describe",
435            SdkOperation::Resume,
436        )?;
437        let descriptor = self
438            .request(
439                method,
440                json!({"sessionId":self.session_id}),
441                SdkOperation::Resume,
442            )
443            .await?;
444        serde_json::from_value(descriptor)
445            .map(|descriptor| self.mask_descriptor(descriptor))
446            .map_err(|error| SdkError::Transport(error.to_string()))
447    }
448
449    async fn attach(&self, history_limit: usize) -> Result<FrontendAttachment, SdkError> {
450        if self.disconnected.load(Ordering::SeqCst) {
451            return Err(SdkError::Closed);
452        }
453        let live = self.events.subscribe();
454        let method = self.method(
455            crate::FrontendFacadeMethod::Attach,
456            "supercode/frontend/attach",
457            SdkOperation::Events,
458        )?;
459        let acknowledged_sequence = self.acknowledged_sequence.load(Ordering::SeqCst);
460        let value = self
461            .request(
462                method,
463                json!({
464                    "sessionId":self.session_id,
465                    "limit":history_limit,
466                    "afterSequence":acknowledged_sequence,
467                }),
468                SdkOperation::Events,
469            )
470            .await?;
471        let mut snapshot: FrontendAttachSnapshot = serde_json::from_value(value)
472            .map_err(|error| SdkError::Transport(error.to_string()))?;
473        snapshot.descriptor = self.mask_descriptor(snapshot.descriptor);
474        suppress_acknowledged_history(
475            &mut snapshot.history,
476            snapshot.history_cursor,
477            acknowledged_sequence,
478        );
479        Ok(
480            FrontendAttachment::from_snapshot_after(snapshot, live, acknowledged_sequence)
481                .with_acknowledgement(self.acknowledged_sequence.clone()),
482        )
483    }
484
485    async fn send_input(self: Arc<Self>, prompt: String) -> Result<(), SdkError> {
486        let method = self.method(
487            crate::FrontendFacadeMethod::SendInput,
488            "supercode/frontend/send_input",
489            SdkOperation::Input,
490        )?;
491        self.request(
492            method,
493            json!({"sessionId":self.session_id,"prompt":prompt}),
494            SdkOperation::Input,
495        )
496        .await
497        .map(|_| ())
498    }
499
500    async fn submit(&self, prompt: String) -> Result<String, SdkError> {
501        let method = self.method(
502            crate::FrontendFacadeMethod::Submit,
503            "supercode/frontend/submit",
504            SdkOperation::Input,
505        )?;
506        let result = self
507            .request(
508                method,
509                json!({"sessionId":self.session_id,"prompt":prompt}),
510                SdkOperation::Input,
511            )
512            .await?;
513        result
514            .get("reply")
515            .and_then(Value::as_str)
516            .map(str::to_owned)
517            .ok_or_else(|| SdkError::Transport("ACP submit omitted reply".into()))
518    }
519
520    async fn submit_with_images(
521        &self,
522        prompt: String,
523        image_urls: Vec<String>,
524    ) -> Result<String, SdkError> {
525        let method = self.method(
526            crate::FrontendFacadeMethod::Submit,
527            "supercode/frontend/submit",
528            SdkOperation::Input,
529        )?;
530        let result = self
531            .request(
532                method,
533                json!({"sessionId":self.session_id,"prompt":prompt,"image_urls":image_urls}),
534                SdkOperation::Input,
535            )
536            .await?;
537        result
538            .get("reply")
539            .and_then(Value::as_str)
540            .map(str::to_owned)
541            .ok_or_else(|| SdkError::Transport("ACP submit omitted reply".into()))
542    }
543
544    async fn interrupt(&self) -> Result<bool, SdkError> {
545        let method = self.method(
546            crate::FrontendFacadeMethod::Interrupt,
547            "supercode/frontend/interrupt",
548            SdkOperation::Interrupt,
549        )?;
550        let result = self
551            .request(
552                method,
553                json!({"sessionId":self.session_id}),
554                SdkOperation::Interrupt,
555            )
556            .await?;
557        Ok(result
558            .get("interrupted")
559            .and_then(Value::as_bool)
560            .unwrap_or(false))
561    }
562
563    async fn steer(&self, prompt: String) -> Result<(), SdkError> {
564        let method = self.method(
565            crate::FrontendFacadeMethod::Steer,
566            "session/steer",
567            SdkOperation::Steer,
568        )?;
569        self.request(
570            method,
571            json!({"sessionId":self.session_id,"text":prompt}),
572            SdkOperation::Steer,
573        )
574        .await
575        .map(|_| ())
576    }
577
578    async fn respond(&self, response: FrontendResponse) -> Result<(), SdkError> {
579        let method = self.method(
580            crate::FrontendFacadeMethod::Respond,
581            "session/respond",
582            SdkOperation::Respond,
583        )?;
584        self.request(
585            method,
586            json!({"sessionId":self.session_id,"response":response}),
587            SdkOperation::Respond,
588        )
589        .await
590        .map(|_| ())
591    }
592
593    async fn invoke(
594        &self,
595        operation: FrontendOperationInvocation,
596    ) -> Result<FrontendOperationResult, SdkError> {
597        let method = self.method(
598            crate::FrontendFacadeMethod::Invoke,
599            "supercode/frontend/invoke",
600            SdkOperation::Input,
601        )?;
602        let result = self
603            .request(
604                method,
605                json!({"sessionId":self.session_id,"operation":operation}),
606                SdkOperation::Input,
607            )
608            .await?;
609        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
610    }
611
612    async fn lease_snapshot(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
613        let method = self.method(
614            crate::FrontendFacadeMethod::Lease,
615            "supercode/frontend/lease",
616            SdkOperation::Events,
617        )?;
618        let result = self
619            .request(
620                method,
621                json!({"sessionId":self.session_id}),
622                SdkOperation::Events,
623            )
624            .await?;
625        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
626    }
627
628    async fn take_control(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
629        let method = self.method(
630            crate::FrontendFacadeMethod::TakeControl,
631            "supercode/frontend/take_control",
632            SdkOperation::Input,
633        )?;
634        let result = self
635            .request(
636                method,
637                json!({"sessionId":self.session_id}),
638                SdkOperation::Input,
639            )
640            .await?;
641        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
642    }
643
644    async fn heartbeat(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
645        let method = self.method(
646            crate::FrontendFacadeMethod::Heartbeat,
647            "supercode/frontend/heartbeat",
648            SdkOperation::Events,
649        )?;
650        let result = self
651            .request(
652                method,
653                json!({"sessionId":self.session_id}),
654                SdkOperation::Events,
655            )
656            .await?;
657        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
658    }
659
660    async fn detach(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
661        let method = self.method(
662            crate::FrontendFacadeMethod::Detach,
663            "supercode/frontend/detach",
664            SdkOperation::Events,
665        )?;
666        let result = self
667            .request(
668                method,
669                json!({"sessionId":self.session_id}),
670                SdkOperation::Events,
671            )
672            .await?;
673        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
674    }
675
676    async fn close(&self) -> Result<(), SdkError> {
677        let descriptor = self.describe().await?;
678        if !descriptor.actions.close {
679            return Err(SdkError::UnsupportedAction(
680                SdkOperation::Close.action_name(),
681            ));
682        }
683        let method = self.method(
684            crate::FrontendFacadeMethod::Close,
685            "supercode/frontend/close",
686            SdkOperation::Close,
687        )?;
688        self.request(
689            method,
690            json!({"sessionId":self.session_id}),
691            SdkOperation::Close,
692        )
693        .await?;
694        self.client.close().await.map_err(transport)
695    }
696}
697
698fn transport(error: impl std::fmt::Display) -> SdkError {
699    SdkError::Transport(error.to_string())
700}
701
702fn decode_sdk_error(raw: &str, operation: SdkOperation) -> SdkError {
703    let value = serde_json::from_str::<Value>(raw).unwrap_or(Value::Null);
704    let name = value.get("name").and_then(Value::as_str);
705    let message = value
706        .get("message")
707        .and_then(Value::as_str)
708        .unwrap_or(raw)
709        .to_string();
710    match name {
711        Some("unauthenticated") => SdkError::Unauthenticated,
712        Some("unauthorized") => SdkError::Unauthorized {
713            permission: value
714                .get("permission")
715                .and_then(Value::as_str)
716                .unwrap_or("unknown")
717                .to_string(),
718        },
719        Some("controller_required") => SdkError::ControllerRequired {
720            holder: value
721                .get("holder")
722                .and_then(Value::as_str)
723                .map(str::to_owned),
724            expires_at_ms: value.get("expiresAtMs").and_then(Value::as_u64),
725        },
726        Some("lease_expired") => SdkError::LeaseExpired,
727        Some("invalid_argument") => SdkError::InvalidArgument { operation, message },
728        Some("not_found") => SdkError::NotFound { operation, message },
729        Some("busy") => crate::RuntimeSubmitError::Busy.into(),
730        Some("unsupported_action") => SdkError::UnsupportedAction(operation.action_name()),
731        Some("execution") => SdkError::Execution { operation, message },
732        Some("transport") => SdkError::Transport(message),
733        _ => SdkError::Transport(message),
734    }
735}
736
737#[cfg(test)]
738mod tests {
739    use super::*;
740
741    #[test]
742    fn negotiated_methods_only_remove_capabilities() {
743        let methods = [
744            "supercode/frontend/submit",
745            "supercode/frontend/send_input",
746            "supercode/frontend/detach",
747            "session/respond",
748        ]
749        .into_iter()
750        .map(str::to_owned)
751        .collect();
752        let mut actions = crate::FrontendActions {
753            submit: true,
754            interrupt: true,
755            steer: true,
756            respond: true,
757            detach: true,
758            close: true,
759        };
760        mask_actions(&methods, &mut actions);
761        assert!(actions.submit);
762        assert!(actions.respond);
763        assert!(actions.detach);
764        assert!(!actions.interrupt);
765        assert!(!actions.steer);
766        assert!(!actions.close);
767    }
768
769    #[test]
770    fn named_sdk_errors_survive_acp_envelopes() {
771        let error = decode_sdk_error(
772            r#"{"name":"busy","operation":"input","message":"occupied"}"#,
773            SdkOperation::Input,
774        );
775        assert_eq!(error.code(), crate::SdkErrorCode::Busy);
776        assert_eq!(error.operation(), Some(SdkOperation::Input));
777    }
778
779    #[test]
780    fn acknowledged_snapshot_history_is_not_rendered_twice() {
781        let mut history = vec![crate::ChatMessage::user("already rendered")];
782        suppress_acknowledged_history(&mut history, 8, 8);
783        assert!(history.is_empty());
784
785        history = vec![crate::ChatMessage::user("must not be lost")];
786        suppress_acknowledged_history(&mut history, 8, 7);
787        assert_eq!(history.len(), 1);
788    }
789
790    #[test]
791    fn checkpoint_schema_and_session_are_fail_closed() {
792        let valid = AcpFrontendCheckpoint {
793            schema_version: 1,
794            session_id: "runtime-a".into(),
795            acknowledged_sequence: 42,
796        };
797        assert!(validate_checkpoint(&valid, "runtime-a").is_ok());
798
799        let mut wrong_schema = valid.clone();
800        wrong_schema.schema_version = 2;
801        assert!(validate_checkpoint(&wrong_schema, "runtime-a")
802            .unwrap_err()
803            .to_string()
804            .contains("schema 2"));
805
806        let mut wrong_session = valid;
807        wrong_session.session_id = "runtime-b".into();
808        assert!(validate_checkpoint(&wrong_session, "runtime-a")
809            .unwrap_err()
810            .to_string()
811            .contains("runtime-b"));
812    }
813
814    #[test]
815    fn transient_disconnect_never_advances_the_canonical_cursor() {
816        let transient = crate::SdkEvent {
817            sequence: 43,
818            kind: "runtime_disconnected".into(),
819            payload: json!({
820                "type":"runtime_disconnected",
821                "_meta":{"supercode":{"transient":true,"transport":"acp"}}
822            }),
823        };
824        assert!(!crate::frontend::event_advances_acknowledgement(&transient));
825
826        let canonical = crate::SdkEvent {
827            sequence: 43,
828            kind: "text_delta".into(),
829            payload: json!({"type":"text_delta","text":"next owner event"}),
830        };
831        assert!(crate::frontend::event_advances_acknowledgement(&canonical));
832    }
833
834    #[test]
835    fn reconnect_replays_owner_event_after_transient_sequence_collision() {
836        let descriptor = || crate::FrontendRuntimeDescriptor {
837            schema_version: crate::FRONTEND_RUNTIME_SCHEMA_VERSION,
838            session_id: "runtime-a".into(),
839            source_harness: None,
840            emulation_profile: None,
841            active_modules: Vec::new(),
842            commands: Vec::new(),
843            operations: Vec::new(),
844            actions: crate::FrontendActions {
845                submit: false,
846                interrupt: false,
847                steer: false,
848                respond: false,
849                detach: true,
850                close: false,
851            },
852            display: crate::FrontendDisplayCapabilities {
853                event_kinds: Vec::new(),
854                opaque_fallback: true,
855            },
856            model: "test".into(),
857            turn_state: crate::FrontendTurnState::Idle,
858            connection_state: crate::FrontendConnectionState::Connected,
859            extensions: Default::default(),
860        };
861        let acknowledged = Arc::new(AtomicU64::new(42));
862        let transient = crate::SdkEvent {
863            sequence: 43,
864            kind: "runtime_disconnected".into(),
865            payload: json!({
866                "type":"runtime_disconnected",
867                "_meta":{"supercode":{"transient":true,"transport":"acp"}}
868            }),
869        };
870        let (_first_sender, first_live) = broadcast::channel(1);
871        let mut first = FrontendAttachment::from_snapshot_after(
872            FrontendAttachSnapshot {
873                descriptor: descriptor(),
874                history: Vec::new(),
875                history_cursor: 0,
876                replay: std::collections::VecDeque::from([transient]),
877            },
878            first_live,
879            42,
880        )
881        .with_acknowledgement(acknowledged.clone());
882        assert_eq!(first.next_replay_event().unwrap().sequence, 43);
883        assert_eq!(acknowledged.load(Ordering::SeqCst), 42);
884        drop(first);
885
886        let canonical = crate::SdkEvent {
887            sequence: 43,
888            kind: "text_delta".into(),
889            payload: json!({"type":"text_delta","text":"next owner event"}),
890        };
891        let (_second_sender, second_live) = broadcast::channel(1);
892        let mut second = FrontendAttachment::from_snapshot_after(
893            FrontendAttachSnapshot {
894                descriptor: descriptor(),
895                history: Vec::new(),
896                history_cursor: 0,
897                replay: std::collections::VecDeque::from([canonical]),
898            },
899            second_live,
900            acknowledged.load(Ordering::SeqCst),
901        )
902        .with_acknowledgement(acknowledged.clone());
903        assert_eq!(second.next_replay_event().unwrap().sequence, 43);
904        assert_eq!(acknowledged.load(Ordering::SeqCst), 43);
905    }
906}