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 acquire_control(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
653 let method = self.method(
654 crate::FrontendFacadeMethod::AcquireControl,
655 "supercode/frontend/acquire_control",
656 SdkOperation::Input,
657 )?;
658 let result = self
659 .request(
660 method,
661 json!({"sessionId":self.session_id}),
662 SdkOperation::Input,
663 )
664 .await?;
665 serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
666 }
667
668 async fn heartbeat(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
669 let method = self.method(
670 crate::FrontendFacadeMethod::Heartbeat,
671 "supercode/frontend/heartbeat",
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 detach(&self) -> Result<crate::RuntimeLeaseSnapshot, SdkError> {
685 let method = self.method(
686 crate::FrontendFacadeMethod::Detach,
687 "supercode/frontend/detach",
688 SdkOperation::Events,
689 )?;
690 let result = self
691 .request(
692 method,
693 json!({"sessionId":self.session_id}),
694 SdkOperation::Events,
695 )
696 .await?;
697 serde_json::from_value(result).map_err(|error| SdkError::Transport(error.to_string()))
698 }
699
700 async fn close(&self) -> Result<(), SdkError> {
701 let descriptor = self.describe().await?;
702 if !descriptor.actions.close {
703 return Err(SdkError::UnsupportedAction(
704 SdkOperation::Close.action_name(),
705 ));
706 }
707 let method = self.method(
708 crate::FrontendFacadeMethod::Close,
709 "supercode/frontend/close",
710 SdkOperation::Close,
711 )?;
712 self.request(
713 method,
714 json!({"sessionId":self.session_id}),
715 SdkOperation::Close,
716 )
717 .await?;
718 self.client.close().await.map_err(transport)
719 }
720}
721
722fn transport(error: impl std::fmt::Display) -> SdkError {
723 SdkError::Transport(error.to_string())
724}
725
726fn decode_sdk_error(raw: &str, operation: SdkOperation) -> SdkError {
727 let value = serde_json::from_str::<Value>(raw).unwrap_or(Value::Null);
728 let name = value.get("name").and_then(Value::as_str);
729 let message = value
730 .get("message")
731 .and_then(Value::as_str)
732 .unwrap_or(raw)
733 .to_string();
734 match name {
735 Some("unauthenticated") => SdkError::Unauthenticated,
736 Some("unauthorized") => SdkError::Unauthorized {
737 permission: value
738 .get("permission")
739 .and_then(Value::as_str)
740 .unwrap_or("unknown")
741 .to_string(),
742 },
743 Some("controller_required") => SdkError::ControllerRequired {
744 holder: value
745 .get("holder")
746 .and_then(Value::as_str)
747 .map(str::to_owned),
748 expires_at_ms: value.get("expiresAtMs").and_then(Value::as_u64),
749 },
750 Some("lease_expired") => SdkError::LeaseExpired,
751 Some("invalid_argument") => SdkError::InvalidArgument { operation, message },
752 Some("not_found") => SdkError::NotFound { operation, message },
753 Some("busy") => crate::RuntimeSubmitError::Busy.into(),
754 Some("unsupported_action") => SdkError::UnsupportedAction(operation.action_name()),
755 Some("execution") => SdkError::Execution { operation, message },
756 Some("transport") => SdkError::Transport(message),
757 _ => SdkError::Transport(message),
758 }
759}
760
761#[cfg(test)]
762mod tests {
763 use super::*;
764
765 #[test]
766 fn negotiated_methods_only_remove_capabilities() {
767 let methods = [
768 "supercode/frontend/submit",
769 "supercode/frontend/send_input",
770 "supercode/frontend/detach",
771 "session/respond",
772 ]
773 .into_iter()
774 .map(str::to_owned)
775 .collect();
776 let mut actions = crate::FrontendActions {
777 submit: true,
778 interrupt: true,
779 steer: true,
780 respond: true,
781 detach: true,
782 close: true,
783 };
784 mask_actions(&methods, &mut actions);
785 assert!(actions.submit);
786 assert!(actions.respond);
787 assert!(actions.detach);
788 assert!(!actions.interrupt);
789 assert!(!actions.steer);
790 assert!(!actions.close);
791 }
792
793 #[test]
794 fn named_sdk_errors_survive_acp_envelopes() {
795 let error = decode_sdk_error(
796 r#"{"name":"busy","operation":"input","message":"occupied"}"#,
797 SdkOperation::Input,
798 );
799 assert_eq!(error.code(), crate::SdkErrorCode::Busy);
800 assert_eq!(error.operation(), Some(SdkOperation::Input));
801 }
802
803 #[test]
804 fn acknowledged_snapshot_history_is_not_rendered_twice() {
805 let mut history = vec![crate::ChatMessage::user("already rendered")];
806 suppress_acknowledged_history(&mut history, 8, 8);
807 assert!(history.is_empty());
808
809 history = vec![crate::ChatMessage::user("must not be lost")];
810 suppress_acknowledged_history(&mut history, 8, 7);
811 assert_eq!(history.len(), 1);
812 }
813
814 #[test]
815 fn checkpoint_schema_and_session_are_fail_closed() {
816 let valid = AcpFrontendCheckpoint {
817 schema_version: 1,
818 session_id: "runtime-a".into(),
819 acknowledged_sequence: 42,
820 };
821 assert!(validate_checkpoint(&valid, "runtime-a").is_ok());
822
823 let mut wrong_schema = valid.clone();
824 wrong_schema.schema_version = 2;
825 assert!(validate_checkpoint(&wrong_schema, "runtime-a")
826 .unwrap_err()
827 .to_string()
828 .contains("schema 2"));
829
830 let mut wrong_session = valid;
831 wrong_session.session_id = "runtime-b".into();
832 assert!(validate_checkpoint(&wrong_session, "runtime-a")
833 .unwrap_err()
834 .to_string()
835 .contains("runtime-b"));
836 }
837
838 #[test]
839 fn transient_disconnect_never_advances_the_canonical_cursor() {
840 let transient = crate::SdkEvent {
841 sequence: 43,
842 kind: "runtime_disconnected".into(),
843 payload: json!({
844 "type":"runtime_disconnected",
845 "_meta":{"supercode":{"transient":true,"transport":"acp"}}
846 }),
847 };
848 assert!(!crate::frontend::event_advances_acknowledgement(&transient));
849
850 let canonical = crate::SdkEvent {
851 sequence: 43,
852 kind: "text_delta".into(),
853 payload: json!({"type":"text_delta","text":"next owner event"}),
854 };
855 assert!(crate::frontend::event_advances_acknowledgement(&canonical));
856 }
857
858 #[test]
859 fn reconnect_replays_owner_event_after_transient_sequence_collision() {
860 let descriptor = || crate::FrontendRuntimeDescriptor {
861 schema_version: crate::FRONTEND_RUNTIME_SCHEMA_VERSION,
862 session_id: "runtime-a".into(),
863 source_harness: None,
864 emulation_profile: None,
865 active_modules: Vec::new(),
866 commands: Vec::new(),
867 operations: Vec::new(),
868 actions: crate::FrontendActions {
869 submit: false,
870 interrupt: false,
871 steer: false,
872 respond: false,
873 detach: true,
874 close: false,
875 },
876 display: crate::FrontendDisplayCapabilities {
877 event_kinds: Vec::new(),
878 opaque_fallback: true,
879 },
880 model: "test".into(),
881 turn_state: crate::FrontendTurnState::Idle,
882 connection_state: crate::FrontendConnectionState::Connected,
883 extensions: Default::default(),
884 };
885 let acknowledged = Arc::new(AtomicU64::new(42));
886 let transient = crate::SdkEvent {
887 sequence: 43,
888 kind: "runtime_disconnected".into(),
889 payload: json!({
890 "type":"runtime_disconnected",
891 "_meta":{"supercode":{"transient":true,"transport":"acp"}}
892 }),
893 };
894 let (_first_sender, first_live) = broadcast::channel(1);
895 let mut first = FrontendAttachment::from_snapshot_after(
896 FrontendAttachSnapshot {
897 descriptor: descriptor(),
898 history: Vec::new(),
899 history_cursor: 0,
900 replay: std::collections::VecDeque::from([transient]),
901 },
902 first_live,
903 42,
904 )
905 .with_acknowledgement(acknowledged.clone());
906 assert_eq!(first.next_replay_event().unwrap().sequence, 43);
907 assert_eq!(acknowledged.load(Ordering::SeqCst), 42);
908 drop(first);
909
910 let canonical = crate::SdkEvent {
911 sequence: 43,
912 kind: "text_delta".into(),
913 payload: json!({"type":"text_delta","text":"next owner event"}),
914 };
915 let (_second_sender, second_live) = broadcast::channel(1);
916 let mut second = FrontendAttachment::from_snapshot_after(
917 FrontendAttachSnapshot {
918 descriptor: descriptor(),
919 history: Vec::new(),
920 history_cursor: 0,
921 replay: std::collections::VecDeque::from([canonical]),
922 },
923 second_live,
924 acknowledged.load(Ordering::SeqCst),
925 )
926 .with_acknowledgement(acknowledged.clone());
927 assert_eq!(second.next_replay_event().unwrap().sequence, 43);
928 assert_eq!(acknowledged.load(Ordering::SeqCst), 43);
929 }
930}