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 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}