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        self.send_input_with_images(prompt, Vec::new()).await
487    }
488
489    async fn send_input_with_images(
490        self: Arc<Self>,
491        prompt: String,
492        image_urls: Vec<String>,
493    ) -> Result<(), SdkError> {
494        let method = self.method(
495            crate::FrontendFacadeMethod::SendInput,
496            "supercode/frontend/send_input",
497            SdkOperation::Input,
498        )?;
499        self.request(
500            method,
501            json!({"sessionId":self.session_id,"prompt":prompt,"image_urls":image_urls}),
502            SdkOperation::Input,
503        )
504        .await
505        .map(|_| ())
506    }
507
508    async fn submit(&self, prompt: String) -> Result<String, SdkError> {
509        let method = self.method(
510            crate::FrontendFacadeMethod::Submit,
511            "supercode/frontend/submit",
512            SdkOperation::Input,
513        )?;
514        let result = self
515            .request(
516                method,
517                json!({"sessionId":self.session_id,"prompt":prompt}),
518                SdkOperation::Input,
519            )
520            .await?;
521        result
522            .get("reply")
523            .and_then(Value::as_str)
524            .map(str::to_owned)
525            .ok_or_else(|| SdkError::Transport("ACP submit omitted reply".into()))
526    }
527
528    async fn submit_with_images(
529        &self,
530        prompt: String,
531        image_urls: Vec<String>,
532    ) -> Result<String, SdkError> {
533        let method = self.method(
534            crate::FrontendFacadeMethod::Submit,
535            "supercode/frontend/submit",
536            SdkOperation::Input,
537        )?;
538        let result = self
539            .request(
540                method,
541                json!({"sessionId":self.session_id,"prompt":prompt,"image_urls":image_urls}),
542                SdkOperation::Input,
543            )
544            .await?;
545        result
546            .get("reply")
547            .and_then(Value::as_str)
548            .map(str::to_owned)
549            .ok_or_else(|| SdkError::Transport("ACP submit omitted reply".into()))
550    }
551
552    async fn interrupt(&self) -> Result<bool, SdkError> {
553        let method = self.method(
554            crate::FrontendFacadeMethod::Interrupt,
555            "supercode/frontend/interrupt",
556            SdkOperation::Interrupt,
557        )?;
558        let result = self
559            .request(
560                method,
561                json!({"sessionId":self.session_id}),
562                SdkOperation::Interrupt,
563            )
564            .await?;
565        Ok(result
566            .get("interrupted")
567            .and_then(Value::as_bool)
568            .unwrap_or(false))
569    }
570
571    async fn steer(&self, prompt: String) -> Result<(), SdkError> {
572        let method = self.method(
573            crate::FrontendFacadeMethod::Steer,
574            "session/steer",
575            SdkOperation::Steer,
576        )?;
577        self.request(
578            method,
579            json!({"sessionId":self.session_id,"text":prompt}),
580            SdkOperation::Steer,
581        )
582        .await
583        .map(|_| ())
584    }
585
586    async fn respond(&self, response: FrontendResponse) -> Result<(), SdkError> {
587        let method = self.method(
588            crate::FrontendFacadeMethod::Respond,
589            "session/respond",
590            SdkOperation::Respond,
591        )?;
592        self.request(
593            method,
594            json!({"sessionId":self.session_id,"response":response}),
595            SdkOperation::Respond,
596        )
597        .await
598        .map(|_| ())
599    }
600
601    async fn invoke(
602        &self,
603        operation: FrontendOperationInvocation,
604    ) -> Result<FrontendOperationResult, SdkError> {
605        let method = self.method(
606            crate::FrontendFacadeMethod::Invoke,
607            "supercode/frontend/invoke",
608            SdkOperation::Input,
609        )?;
610        let result = self
611            .request(
612                method,
613                json!({"sessionId":self.session_id,"operation":operation}),
614                SdkOperation::Input,
615            )
616            .await?;
617        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
618    }
619
620    async fn lease_snapshot(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
621        let method = self.method(
622            crate::FrontendFacadeMethod::Lease,
623            "supercode/frontend/lease",
624            SdkOperation::Events,
625        )?;
626        let result = self
627            .request(
628                method,
629                json!({"sessionId":self.session_id}),
630                SdkOperation::Events,
631            )
632            .await?;
633        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
634    }
635
636    async fn take_control(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
637        let method = self.method(
638            crate::FrontendFacadeMethod::TakeControl,
639            "supercode/frontend/take_control",
640            SdkOperation::Input,
641        )?;
642        let result = self
643            .request(
644                method,
645                json!({"sessionId":self.session_id}),
646                SdkOperation::Input,
647            )
648            .await?;
649        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
650    }
651
652    async fn heartbeat(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
653        let method = self.method(
654            crate::FrontendFacadeMethod::Heartbeat,
655            "supercode/frontend/heartbeat",
656            SdkOperation::Events,
657        )?;
658        let result = self
659            .request(
660                method,
661                json!({"sessionId":self.session_id}),
662                SdkOperation::Events,
663            )
664            .await?;
665        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
666    }
667
668    async fn detach(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
669        let method = self.method(
670            crate::FrontendFacadeMethod::Detach,
671            "supercode/frontend/detach",
672            SdkOperation::Events,
673        )?;
674        let result = self
675            .request(
676                method,
677                json!({"sessionId":self.session_id}),
678                SdkOperation::Events,
679            )
680            .await?;
681        serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
682    }
683
684    async fn close(&self) -> Result<(), SdkError> {
685        let descriptor = self.describe().await?;
686        if !descriptor.actions.close {
687            return Err(SdkError::UnsupportedAction(
688                SdkOperation::Close.action_name(),
689            ));
690        }
691        let method = self.method(
692            crate::FrontendFacadeMethod::Close,
693            "supercode/frontend/close",
694            SdkOperation::Close,
695        )?;
696        self.request(
697            method,
698            json!({"sessionId":self.session_id}),
699            SdkOperation::Close,
700        )
701        .await?;
702        self.client.close().await.map_err(transport)
703    }
704}
705
706fn transport(error: impl std::fmt::Display) -> SdkError {
707    SdkError::Transport(error.to_string())
708}
709
710fn decode_sdk_error(raw: &str, operation: SdkOperation) -> SdkError {
711    let value = serde_json::from_str::<Value>(raw).unwrap_or(Value::Null);
712    let name = value.get("name").and_then(Value::as_str);
713    let message = value
714        .get("message")
715        .and_then(Value::as_str)
716        .unwrap_or(raw)
717        .to_string();
718    match name {
719        Some("unauthenticated") => SdkError::Unauthenticated,
720        Some("unauthorized") => SdkError::Unauthorized {
721            permission: value
722                .get("permission")
723                .and_then(Value::as_str)
724                .unwrap_or("unknown")
725                .to_string(),
726        },
727        Some("controller_required") => SdkError::ControllerRequired {
728            holder: value
729                .get("holder")
730                .and_then(Value::as_str)
731                .map(str::to_owned),
732            expires_at_ms: value.get("expiresAtMs").and_then(Value::as_u64),
733        },
734        Some("lease_expired") => SdkError::LeaseExpired,
735        Some("invalid_argument") => SdkError::InvalidArgument { operation, message },
736        Some("not_found") => SdkError::NotFound { operation, message },
737        Some("busy") => crate::RuntimeSubmitError::Busy.into(),
738        Some("unsupported_action") => SdkError::UnsupportedAction(operation.action_name()),
739        Some("execution") => SdkError::Execution { operation, message },
740        Some("transport") => SdkError::Transport(message),
741        _ => SdkError::Transport(message),
742    }
743}
744
745#[cfg(test)]
746mod tests {
747    use super::*;
748
749    #[test]
750    fn negotiated_methods_only_remove_capabilities() {
751        let methods = [
752            "supercode/frontend/submit",
753            "supercode/frontend/send_input",
754            "supercode/frontend/detach",
755            "session/respond",
756        ]
757        .into_iter()
758        .map(str::to_owned)
759        .collect();
760        let mut actions = crate::FrontendActions {
761            submit: true,
762            interrupt: true,
763            steer: true,
764            respond: true,
765            detach: true,
766            close: true,
767        };
768        mask_actions(&methods, &mut actions);
769        assert!(actions.submit);
770        assert!(actions.respond);
771        assert!(actions.detach);
772        assert!(!actions.interrupt);
773        assert!(!actions.steer);
774        assert!(!actions.close);
775    }
776
777    #[test]
778    fn named_sdk_errors_survive_acp_envelopes() {
779        let error = decode_sdk_error(
780            r#"{"name":"busy","operation":"input","message":"occupied"}"#,
781            SdkOperation::Input,
782        );
783        assert_eq!(error.code(), crate::SdkErrorCode::Busy);
784        assert_eq!(error.operation(), Some(SdkOperation::Input));
785    }
786
787    #[test]
788    fn acknowledged_snapshot_history_is_not_rendered_twice() {
789        let mut history = vec![crate::ChatMessage::user("already rendered")];
790        suppress_acknowledged_history(&mut history, 8, 8);
791        assert!(history.is_empty());
792
793        history = vec![crate::ChatMessage::user("must not be lost")];
794        suppress_acknowledged_history(&mut history, 8, 7);
795        assert_eq!(history.len(), 1);
796    }
797
798    #[test]
799    fn checkpoint_schema_and_session_are_fail_closed() {
800        let valid = AcpFrontendCheckpoint {
801            schema_version: 1,
802            session_id: "runtime-a".into(),
803            acknowledged_sequence: 42,
804        };
805        assert!(validate_checkpoint(&valid, "runtime-a").is_ok());
806
807        let mut wrong_schema = valid.clone();
808        wrong_schema.schema_version = 2;
809        assert!(validate_checkpoint(&wrong_schema, "runtime-a")
810            .unwrap_err()
811            .to_string()
812            .contains("schema 2"));
813
814        let mut wrong_session = valid;
815        wrong_session.session_id = "runtime-b".into();
816        assert!(validate_checkpoint(&wrong_session, "runtime-a")
817            .unwrap_err()
818            .to_string()
819            .contains("runtime-b"));
820    }
821
822    #[test]
823    fn transient_disconnect_never_advances_the_canonical_cursor() {
824        let transient = crate::SdkEvent {
825            sequence: 43,
826            kind: "runtime_disconnected".into(),
827            payload: json!({
828                "type":"runtime_disconnected",
829                "_meta":{"supercode":{"transient":true,"transport":"acp"}}
830            }),
831        };
832        assert!(!crate::frontend::event_advances_acknowledgement(&transient));
833
834        let canonical = crate::SdkEvent {
835            sequence: 43,
836            kind: "text_delta".into(),
837            payload: json!({"type":"text_delta","text":"next owner event"}),
838        };
839        assert!(crate::frontend::event_advances_acknowledgement(&canonical));
840    }
841
842    #[test]
843    fn reconnect_replays_owner_event_after_transient_sequence_collision() {
844        let descriptor = || crate::FrontendRuntimeDescriptor {
845            schema_version: crate::FRONTEND_RUNTIME_SCHEMA_VERSION,
846            session_id: "runtime-a".into(),
847            source_harness: None,
848            emulation_profile: None,
849            active_modules: Vec::new(),
850            commands: Vec::new(),
851            operations: Vec::new(),
852            actions: crate::FrontendActions {
853                submit: false,
854                interrupt: false,
855                steer: false,
856                respond: false,
857                detach: true,
858                close: false,
859            },
860            display: crate::FrontendDisplayCapabilities {
861                event_kinds: Vec::new(),
862                opaque_fallback: true,
863            },
864            model: "test".into(),
865            turn_state: crate::FrontendTurnState::Idle,
866            connection_state: crate::FrontendConnectionState::Connected,
867            extensions: Default::default(),
868        };
869        let acknowledged = Arc::new(AtomicU64::new(42));
870        let transient = crate::SdkEvent {
871            sequence: 43,
872            kind: "runtime_disconnected".into(),
873            payload: json!({
874                "type":"runtime_disconnected",
875                "_meta":{"supercode":{"transient":true,"transport":"acp"}}
876            }),
877        };
878        let (_first_sender, first_live) = broadcast::channel(1);
879        let mut first = FrontendAttachment::from_snapshot_after(
880            FrontendAttachSnapshot {
881                descriptor: descriptor(),
882                history: Vec::new(),
883                history_cursor: 0,
884                replay: std::collections::VecDeque::from([transient]),
885            },
886            first_live,
887            42,
888        )
889        .with_acknowledgement(acknowledged.clone());
890        assert_eq!(first.next_replay_event().unwrap().sequence, 43);
891        assert_eq!(acknowledged.load(Ordering::SeqCst), 42);
892        drop(first);
893
894        let canonical = crate::SdkEvent {
895            sequence: 43,
896            kind: "text_delta".into(),
897            payload: json!({"type":"text_delta","text":"next owner event"}),
898        };
899        let (_second_sender, second_live) = broadcast::channel(1);
900        let mut second = FrontendAttachment::from_snapshot_after(
901            FrontendAttachSnapshot {
902                descriptor: descriptor(),
903                history: Vec::new(),
904                history_cursor: 0,
905                replay: std::collections::VecDeque::from([canonical]),
906            },
907            second_live,
908            acknowledged.load(Ordering::SeqCst),
909        )
910        .with_acknowledgement(acknowledged.clone());
911        assert_eq!(second.next_replay_event().unwrap().sequence, 43);
912        assert_eq!(acknowledged.load(Ordering::SeqCst), 43);
913    }
914}