1use 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#[derive(Debug, Clone)]
28pub struct AcpFrontendConnectOptions {
29 pub launch: RuntimeLaunch,
32 pub cwd: Option<PathBuf>,
34 pub session_id: Option<String>,
37 pub after_sequence: Option<u64>,
39}
40
41#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
44pub struct AcpFrontendCheckpoint {
45 pub schema_version: u32,
47 pub session_id: String,
49 pub acknowledged_sequence: u64,
51}
52
53pub 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 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 let _ = runtime.describe().await?;
252 Ok(runtime)
253 }
254
255 pub fn session_id(&self) -> &str {
257 &self.session_id
258 }
259
260 pub fn endpoint(&self) -> &RuntimeEndpoint {
262 &self.endpoint
263 }
264
265 pub fn extensions(&self) -> &Value {
268 &self.extensions
269 }
270
271 pub fn event_method(&self) -> &str {
273 &self.event_method
274 }
275
276 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 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 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 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 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}