1use std::collections::{BTreeMap, VecDeque};
9#[cfg(feature = "adapter-api")]
10use std::sync::atomic::AtomicBool;
11use std::sync::atomic::{AtomicU64, Ordering};
12use std::sync::Arc;
13#[cfg(feature = "adapter-api")]
14use std::sync::Weak;
15
16use async_trait::async_trait;
17#[cfg(feature = "adapter-api")]
18use futures::StreamExt;
19use serde::{Deserialize, Serialize};
20#[cfg(feature = "adapter-api")]
21use serde_json::json;
22use serde_json::Value;
23use tokio::sync::broadcast;
24
25#[cfg(feature = "adapter-api")]
26use crate::sdk::RuntimeSubmitError;
27pub use crate::sdk::SdkError as FrontendRuntimeError;
28pub use crate::sdk::SdkEvent as FrontendEvent;
29pub use crate::sdk::SdkRuntime as FrontendRuntime;
30use crate::server::RpcEngine;
31use crate::ChatMessage;
32
33pub const FRONTEND_RUNTIME_SCHEMA_VERSION: u32 = 2;
35
36pub(crate) const FRONTEND_EVENT_SCHEMA_VERSION: u32 = 1;
41
42pub const FRONTEND_REPLAY_CAPACITY: usize = 4096;
44
45#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
47#[serde(rename_all = "snake_case")]
48pub enum FrontendTurnState {
49 Idle,
51 Busy,
53}
54
55#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
57#[serde(rename_all = "snake_case")]
58pub enum FrontendConnectionState {
59 Connected,
61 ShuttingDown,
63}
64
65#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
67pub struct FrontendActions {
68 pub submit: bool,
70 pub interrupt: bool,
72 pub steer: bool,
74 pub respond: bool,
76 pub detach: bool,
78 pub close: bool,
80}
81
82#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
84pub struct FrontendDisplayCapabilities {
85 pub event_kinds: Vec<String>,
87 pub opaque_fallback: bool,
89}
90
91#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
93pub struct FrontendCommandDescriptor {
94 pub name: String,
96 pub description: Option<String>,
98 #[serde(default, skip_serializing_if = "Option::is_none")]
100 pub argument_hint: Option<String>,
101}
102
103#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
110#[serde(rename_all = "snake_case")]
111pub enum FrontendOperationKind {
112 Prompt,
114 File,
116 Model,
118 Session,
120 Subagent,
122 Image,
124 Reduction,
126 Context,
133}
134
135#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
137pub struct FrontendOperationDescriptor {
138 pub id: String,
140 pub kind: FrontendOperationKind,
142 pub command: Option<FrontendCommandDescriptor>,
144}
145
146#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
148#[serde(tag = "kind", rename_all = "snake_case")]
149pub enum FrontendOperationInvocation {
150 Prompt {
152 operation_id: String,
154 arguments: String,
156 },
157 Context {
160 operation_id: String,
162 },
163 Model {
169 operation_id: String,
171 model: String,
175 },
176}
177
178impl FrontendOperationInvocation {
179 pub fn operation_id(&self) -> &str {
181 match self {
182 Self::Prompt { operation_id, .. }
183 | Self::Context { operation_id }
184 | Self::Model { operation_id, .. } => operation_id,
185 }
186 }
187}
188
189#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
191#[serde(tag = "kind", rename_all = "snake_case")]
192pub enum FrontendOperationResult {
193 Prompt {
195 reply: String,
197 },
198 Context {
200 usage: crate::ContextUsage,
202 },
203 Model {
205 model: String,
207 previous: String,
210 },
211}
212
213#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
215pub struct FrontendRuntimeMetadata {
216 pub source_harness: Option<String>,
218 pub emulation_profile: Option<String>,
220}
221
222#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
224pub struct FrontendRuntimeDescriptor {
225 pub schema_version: u32,
227 pub session_id: String,
229 pub source_harness: Option<String>,
231 pub emulation_profile: Option<String>,
233 pub active_modules: Vec<String>,
235 pub commands: Vec<FrontendCommandDescriptor>,
237 #[serde(default)]
239 pub operations: Vec<FrontendOperationDescriptor>,
240 pub actions: FrontendActions,
242 pub display: FrontendDisplayCapabilities,
244 pub model: String,
246 pub turn_state: FrontendTurnState,
248 pub connection_state: FrontendConnectionState,
250 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
252 pub extensions: BTreeMap<String, Value>,
253}
254
255#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
258pub struct FrontendAttachSnapshot {
259 pub descriptor: FrontendRuntimeDescriptor,
261 pub history: Vec<ChatMessage>,
263 pub history_cursor: u64,
265 pub replay: VecDeque<FrontendEvent>,
267}
268
269#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
271#[serde(rename_all = "snake_case")]
272pub enum FrontendRequestKind {
273 Approval,
275 Elicitation,
277 #[serde(other)]
280 Other,
281}
282
283#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
285pub struct FrontendRequest {
286 pub id: u64,
288 pub kind: FrontendRequestKind,
290 pub payload: Value,
292}
293
294#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
296#[serde(rename_all = "snake_case")]
297pub enum FrontendApprovalDecision {
298 Deny,
300 Allow,
302 AllowForSession,
304}
305
306#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
308#[serde(rename_all = "snake_case")]
309pub enum FrontendElicitationAction {
310 Accept,
312 Decline,
314 Cancel,
316}
317
318#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
320#[serde(tag = "kind", rename_all = "snake_case")]
321pub enum FrontendResponse {
322 Approval {
324 request_id: u64,
326 decision: FrontendApprovalDecision,
328 },
329 Elicitation {
331 request_id: u64,
333 action: FrontendElicitationAction,
335 content: Option<Value>,
337 },
338 Other {
340 request_id: u64,
342 action: FrontendElicitationAction,
344 content: Option<Value>,
346 },
347}
348
349impl FrontendResponse {
350 pub(crate) fn request_id(&self) -> u64 {
351 match self {
352 Self::Approval { request_id, .. }
353 | Self::Elicitation { request_id, .. }
354 | Self::Other { request_id, .. } => *request_id,
355 }
356 }
357}
358
359pub struct FrontendAttachment {
361 pub descriptor: FrontendRuntimeDescriptor,
363 pub history: Vec<ChatMessage>,
365 pub history_cursor: u64,
367 pub(crate) replay: VecDeque<FrontendEvent>,
368 live: broadcast::Receiver<FrontendEvent>,
369 delivered: u64,
370 acknowledged: Option<Arc<AtomicU64>>,
371 _transport_lease: Option<Arc<()>>,
372}
373
374impl FrontendAttachment {
375 pub fn from_snapshot(
379 snapshot: FrontendAttachSnapshot,
380 live: broadcast::Receiver<FrontendEvent>,
381 ) -> Self {
382 Self::from_snapshot_after(snapshot, live, 0)
383 }
384
385 pub fn from_snapshot_after(
390 snapshot: FrontendAttachSnapshot,
391 live: broadcast::Receiver<FrontendEvent>,
392 acknowledged_sequence: u64,
393 ) -> Self {
394 let delivered = snapshot.history_cursor.max(acknowledged_sequence);
395 Self::new_with_delivered(
396 snapshot.descriptor,
397 snapshot.history,
398 snapshot.history_cursor,
399 snapshot.replay,
400 live,
401 None,
402 delivered,
403 )
404 }
405
406 pub(crate) fn new(
407 descriptor: FrontendRuntimeDescriptor,
408 history: Vec<ChatMessage>,
409 history_cursor: u64,
410 replay: VecDeque<FrontendEvent>,
411 live: broadcast::Receiver<FrontendEvent>,
412 transport_lease: Option<Arc<()>>,
413 ) -> Self {
414 let delivered = history_cursor;
415 Self::new_with_delivered(
416 descriptor,
417 history,
418 history_cursor,
419 replay,
420 live,
421 transport_lease,
422 delivered,
423 )
424 }
425
426 fn new_with_delivered(
427 descriptor: FrontendRuntimeDescriptor,
428 history: Vec<ChatMessage>,
429 history_cursor: u64,
430 replay: VecDeque<FrontendEvent>,
431 live: broadcast::Receiver<FrontendEvent>,
432 transport_lease: Option<Arc<()>>,
433 delivered: u64,
434 ) -> Self {
435 Self {
436 descriptor,
437 history,
438 history_cursor,
439 replay,
440 live,
441 delivered,
442 acknowledged: None,
443 _transport_lease: transport_lease,
444 }
445 }
446
447 #[cfg(feature = "adapter-acp")]
448 pub(crate) fn with_acknowledgement(mut self, acknowledged: Arc<AtomicU64>) -> Self {
449 acknowledged.fetch_max(self.history_cursor, Ordering::SeqCst);
450 self.acknowledged = Some(acknowledged);
451 self
452 }
453
454 fn acknowledge(&self, event: &FrontendEvent) {
455 if !event_advances_acknowledgement(event) {
456 return;
457 }
458 if let Some(acknowledged) = &self.acknowledged {
459 acknowledged.fetch_max(event.sequence, Ordering::SeqCst);
460 }
461 }
462
463 pub async fn next_event(&mut self) -> Result<FrontendEvent, FrontendRuntimeError> {
467 loop {
468 let event = match self.next_replay_event() {
469 Some(event) => return Ok(event),
470 None => match self.live.recv().await {
471 Ok(event) => event,
472 Err(broadcast::error::RecvError::Lagged(count)) => {
473 return Err(FrontendRuntimeError::ReplayGap(count));
474 }
475 Err(broadcast::error::RecvError::Closed) => {
476 return Err(FrontendRuntimeError::Closed);
477 }
478 },
479 };
480 if event.sequence <= self.delivered {
481 continue;
482 }
483 self.delivered = event.sequence;
484 self.acknowledge(&event);
485 return Ok(event);
486 }
487 }
488
489 pub fn next_replay_event(&mut self) -> Option<FrontendEvent> {
494 while let Some(event) = self.replay.pop_front() {
495 if event.sequence <= self.delivered {
496 continue;
497 }
498 self.delivered = event.sequence;
499 self.acknowledge(&event);
500 return Some(event);
501 }
502 None
503 }
504}
505
506pub(crate) fn event_advances_acknowledgement(event: &FrontendEvent) -> bool {
507 event
508 .payload
509 .pointer("/_meta/supercode/transient")
510 .and_then(Value::as_bool)
511 != Some(true)
512}
513
514pub(crate) struct FrontendProjectionState {
516 pub(crate) history: Vec<ChatMessage>,
517 pub(crate) history_cursor: u64,
518 pub(crate) next_sequence: u64,
519 pub(crate) replay: VecDeque<FrontendEvent>,
520}
521
522#[async_trait]
523impl FrontendRuntime for RpcEngine {
524 async fn describe(&self) -> Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
525 Ok(self.frontend_descriptor())
526 }
527
528 async fn attach(
529 &self,
530 history_limit: usize,
531 ) -> Result<FrontendAttachment, FrontendRuntimeError> {
532 self.frontend_attach(history_limit)
533 }
534
535 async fn send_input(self: Arc<Self>, prompt: String) -> Result<(), FrontendRuntimeError> {
536 RpcEngine::send_input(&self, prompt)?;
537 Ok(())
538 }
539
540 async fn send_input_with_images(
541 self: Arc<Self>,
542 prompt: String,
543 image_urls: Vec<String>,
544 ) -> Result<(), FrontendRuntimeError> {
545 RpcEngine::send_input_with_images(&self, prompt, image_urls)?;
546 Ok(())
547 }
548
549 async fn submit(&self, prompt: String) -> Result<String, FrontendRuntimeError> {
550 Ok(RpcEngine::submit(self, prompt).await?)
551 }
552
553 async fn submit_with_images(
554 &self,
555 prompt: String,
556 image_urls: Vec<String>,
557 ) -> Result<String, FrontendRuntimeError> {
558 Ok(RpcEngine::submit_with_images(self, prompt, image_urls).await?)
559 }
560
561 async fn interrupt(&self) -> Result<bool, FrontendRuntimeError> {
562 Ok(RpcEngine::interrupt(self).await)
563 }
564
565 async fn steer(&self, prompt: String) -> Result<(), FrontendRuntimeError> {
566 RpcEngine::steer(self, prompt)
567 }
568
569 async fn respond(&self, response: FrontendResponse) -> Result<(), FrontendRuntimeError> {
570 RpcEngine::respond(self, response)
571 }
572
573 async fn invoke(
574 &self,
575 operation: FrontendOperationInvocation,
576 ) -> Result<FrontendOperationResult, FrontendRuntimeError> {
577 RpcEngine::invoke(self, operation).await
578 }
579
580 async fn close(&self) -> Result<(), FrontendRuntimeError> {
581 RpcEngine::shutdown(self).await;
582 Ok(())
583 }
584}
585
586#[cfg(feature = "adapter-api")]
591pub struct HttpFrontendRuntime {
592 base_url: String,
593 token: String,
594 client_id: crate::RuntimeClientId,
595 authorization: crate::RuntimeAuthorization,
596 client: reqwest::Client,
597 events: broadcast::Sender<FrontendEvent>,
598 next_id: AtomicU64,
599 lifecycle: Arc<()>,
600 disconnected: AtomicBool,
601}
602
603#[cfg(feature = "adapter-api")]
604const KEPT_PROBE_IDLE: std::time::Duration = std::time::Duration::from_secs(600);
606
607type KeptProbes =
608 std::collections::HashMap<(String, String), (Arc<HttpFrontendRuntime>, std::time::Instant)>;
609
610fn kept_probes() -> &'static std::sync::Mutex<KeptProbes> {
611 static KEPT: std::sync::OnceLock<std::sync::Mutex<KeptProbes>> = std::sync::OnceLock::new();
612 KEPT.get_or_init(Default::default)
613}
614
615fn kept_probe_runtime() -> &'static tokio::runtime::Runtime {
617 static RUNTIME: std::sync::OnceLock<tokio::runtime::Runtime> = std::sync::OnceLock::new();
618 RUNTIME.get_or_init(|| {
619 tokio::runtime::Builder::new_multi_thread()
620 .worker_threads(1)
621 .thread_name("supercode-runtime-probes")
622 .enable_all()
623 .build()
624 .expect("the runtime-probe runtime starts")
625 })
626}
627
628impl HttpFrontendRuntime {
629 pub async fn connect(
632 base_url: impl Into<String>,
633 token: impl Into<String>,
634 ) -> Result<Arc<Self>, FrontendRuntimeError> {
635 let mut random = [0_u8; 16];
636 getrandom::getrandom(&mut random).map_err(|error| {
637 FrontendRuntimeError::Transport(format!(
638 "cannot generate runtime client identity: {error}"
639 ))
640 })?;
641 let suffix = random
642 .iter()
643 .map(|byte| format!("{byte:02x}"))
644 .collect::<String>();
645 let client_id = crate::RuntimeClientId::parse(format!("http-{suffix}"))
646 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
647 Self::connect_with_client_id(base_url, token, client_id).await
648 }
649
650 pub async fn connect_with_client_id(
653 base_url: impl Into<String>,
654 token: impl Into<String>,
655 client_id: crate::RuntimeClientId,
656 ) -> Result<Arc<Self>, FrontendRuntimeError> {
657 Self::connect_with_authorization(
658 base_url,
659 token,
660 client_id,
661 crate::RuntimeAuthorization::owner(),
662 )
663 .await
664 }
665
666 pub async fn connect_with_authorization(
670 base_url: impl Into<String>,
671 token: impl Into<String>,
672 client_id: crate::RuntimeClientId,
673 authorization: crate::RuntimeAuthorization,
674 ) -> Result<Arc<Self>, FrontendRuntimeError> {
675 Self::connect_inner(base_url, token, client_id, authorization, true)
676 .await
677 .map(|(runtime, _)| runtime)
678 }
679
680 pub(crate) async fn probe_described(
687 base_url: impl Into<String>,
688 token: impl Into<String>,
689 client_id: crate::RuntimeClientId,
690 ) -> Result<(Arc<Self>, FrontendRuntimeDescriptor), FrontendRuntimeError> {
691 Self::connect_inner(
692 base_url,
693 token,
694 client_id,
695 crate::RuntimeAuthorization::observer(),
696 false,
697 )
698 .await
699 }
700
701 pub(crate) async fn probe_described_kept(
708 base_url: String,
709 token: String,
710 client_id: crate::RuntimeClientId,
711 ) -> Result<(Arc<Self>, FrontendRuntimeDescriptor), FrontendRuntimeError> {
712 kept_probe_runtime()
713 .spawn(async move {
714 let key = (base_url.clone(), token.clone());
715 let kept = {
716 let mut kept = kept_probes().lock().unwrap_or_else(|e| e.into_inner());
717 kept.retain(|_, (_, used)| used.elapsed() < KEPT_PROBE_IDLE);
718 kept.get(&key).map(|(runtime, _)| runtime.clone())
719 };
720 if let Some(runtime) = kept {
721 match FrontendRuntime::describe(runtime.as_ref()).await {
722 Ok(descriptor) => {
723 if let Some(entry) = kept_probes()
724 .lock()
725 .unwrap_or_else(|e| e.into_inner())
726 .get_mut(&key)
727 {
728 entry.1 = std::time::Instant::now();
729 }
730 return Ok((runtime, descriptor));
731 }
732 Err(_) => {
733 kept_probes()
734 .lock()
735 .unwrap_or_else(|e| e.into_inner())
736 .remove(&key);
737 }
738 }
739 }
740 let (runtime, descriptor) =
741 Self::probe_described(base_url, token, client_id).await?;
742 kept_probes()
743 .lock()
744 .unwrap_or_else(|e| e.into_inner())
745 .insert(key, (runtime.clone(), std::time::Instant::now()));
746 Ok((runtime, descriptor))
747 })
748 .await
749 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?
750 }
751
752 async fn connect_inner(
753 base_url: impl Into<String>,
754 token: impl Into<String>,
755 client_id: crate::RuntimeClientId,
756 authorization: crate::RuntimeAuthorization,
757 stream_events: bool,
758 ) -> Result<(Arc<Self>, FrontendRuntimeDescriptor), FrontendRuntimeError> {
759 let runtime = Arc::new(Self {
760 base_url: base_url.into().trim_end_matches('/').to_string(),
761 token: token.into(),
762 client_id,
763 authorization,
764 client: reqwest::Client::new(),
765 events: broadcast::channel(1024).0,
766 next_id: AtomicU64::new(1),
767 lifecycle: Arc::new(()),
768 disconnected: AtomicBool::new(false),
769 });
770 let descriptor: FrontendRuntimeDescriptor = runtime
772 .rpc_typed(crate::FrontendFacadeMethod::Describe.wire_name(), json!({}))
773 .await?;
774 if stream_events {
775 Self::start_event_stream(&runtime).await?;
776 }
777 Ok((runtime, descriptor))
778 }
779
780 async fn start_event_stream(runtime: &Arc<Self>) -> Result<(), FrontendRuntimeError> {
781 let (ready_tx, ready_rx) = tokio::sync::oneshot::channel();
782 let weak = Arc::downgrade(runtime);
783 let lifecycle = Arc::downgrade(&runtime.lifecycle);
784 tokio::spawn(async move {
785 Self::run_event_stream(weak, lifecycle, ready_tx).await;
786 });
787 ready_rx.await.map_err(|_| {
788 FrontendRuntimeError::Transport("frontend event stream exited before startup".into())
789 })?
790 }
791
792 async fn run_event_stream(
793 weak: Weak<Self>,
794 lifecycle: Weak<()>,
795 ready: tokio::sync::oneshot::Sender<Result<(), FrontendRuntimeError>>,
796 ) {
797 let Some(runtime) = weak.upgrade() else {
798 let _ = ready.send(Err(FrontendRuntimeError::Closed));
799 return;
800 };
801 let request = runtime
802 .client
803 .get(format!("{}/frontend/events", runtime.base_url))
804 .bearer_auth(&runtime.token)
805 .header("x-supercode-client-id", runtime.client_id.as_str())
806 .header(
807 "x-supercode-permissions",
808 runtime.authorization.header_value(),
809 );
810 let events = runtime.events.clone();
811 drop(runtime);
812 let response = request.send().await;
813 let response = match response {
814 Ok(response) if response.status().is_success() => response,
815 Ok(response) => {
816 let _ = ready.send(Err(FrontendRuntimeError::Transport(format!(
817 "frontend event stream returned {}",
818 response.status()
819 ))));
820 return;
821 }
822 Err(error) => {
823 let _ = ready.send(Err(FrontendRuntimeError::Transport(error.to_string())));
824 return;
825 }
826 };
827 let _ = ready.send(Ok(()));
828 let mut stream = response.bytes_stream();
829 let mut pending = Vec::<u8>::new();
830 let mut liveness = tokio::time::interval(std::time::Duration::from_millis(100));
831 loop {
832 let chunk = tokio::select! {
833 _ = liveness.tick() => {
834 if lifecycle.strong_count() == 0 {
835 break;
836 }
837 if weak
838 .upgrade()
839 .is_some_and(|runtime| runtime.disconnected.load(Ordering::SeqCst))
840 {
841 break;
842 }
843 continue;
844 }
845 chunk = stream.next() => chunk,
846 };
847 let Some(chunk) = chunk else {
848 break;
849 };
850 let Ok(chunk) = chunk else {
851 break;
852 };
853 pending.extend_from_slice(&chunk);
854 while let Some(position) = pending.iter().position(|byte| *byte == b'\n') {
855 let line = pending.drain(..=position).collect::<Vec<_>>();
856 let line = String::from_utf8_lossy(&line);
857 let Some(data) = line.trim_end().strip_prefix("data: ") else {
858 continue;
859 };
860 if let Ok(event) = serde_json::from_str::<FrontendEvent>(data) {
861 let _ = events.send(event);
862 }
863 }
864 }
865 if let Some(runtime) = weak.upgrade() {
866 runtime.disconnected.store(true, Ordering::SeqCst);
867 let _ = runtime.events.send(FrontendEvent::new(
868 u64::MAX,
869 json!({
870 "type": "runtime_disconnected",
871 "schema_version": FRONTEND_EVENT_SCHEMA_VERSION
872 }),
873 ));
874 }
875 }
876
877 async fn rpc(&self, method: &str, params: Value) -> Result<Value, FrontendRuntimeError> {
878 let id = self.next_id.fetch_add(1, Ordering::SeqCst);
879 let requested_operation = params
880 .pointer("/operation/operation_id")
881 .and_then(Value::as_str)
882 .map(str::to_owned);
883 let response = self
884 .client
885 .post(format!("{}/rpc", self.base_url))
886 .bearer_auth(&self.token)
887 .header("x-supercode-client-id", self.client_id.as_str())
888 .header("x-supercode-permissions", self.authorization.header_value())
889 .json(&json!({"id": id, "method": method, "params": params}))
890 .send()
891 .await
892 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
893 if !response.status().is_success() {
894 return Err(FrontendRuntimeError::Transport(format!(
895 "SDK HTTP RPC returned {}",
896 response.status()
897 )));
898 }
899 let value: Value = response
900 .json()
901 .await
902 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))?;
903 if let Some(error) = value.get("error") {
904 let code = error.get("code").and_then(Value::as_i64);
905 let name = error.get("name").and_then(Value::as_str);
906 let operation = error
907 .get("operation")
908 .and_then(Value::as_str)
909 .and_then(crate::SdkOperation::from_action_name);
910 let message = error
911 .get("message")
912 .and_then(Value::as_str)
913 .unwrap_or("SDK runtime request failed")
914 .to_string();
915 return Err(match (name, code) {
916 (Some("unauthenticated"), _) | (_, Some(-32030)) => {
917 FrontendRuntimeError::Unauthenticated
918 }
919 (Some("unauthorized"), _) | (_, Some(-32031)) => {
920 FrontendRuntimeError::Unauthorized {
921 permission: error
922 .get("permission")
923 .and_then(Value::as_str)
924 .unwrap_or("unknown")
925 .to_string(),
926 }
927 }
928 (Some("controller_required"), _) | (_, Some(-32032)) => {
929 FrontendRuntimeError::ControllerRequired {
930 holder: error
931 .get("holder")
932 .and_then(Value::as_str)
933 .map(str::to_owned),
934 expires_at_ms: error.get("expiresAtMs").and_then(Value::as_u64),
935 }
936 }
937 (Some("lease_expired"), _) | (_, Some(-32033)) => {
938 FrontendRuntimeError::LeaseExpired
939 }
940 (_, Some(-32023)) => FrontendRuntimeError::UnsupportedOperation(
941 requested_operation.unwrap_or(message),
942 ),
943 (Some("unsupported_action"), _) => FrontendRuntimeError::UnsupportedAction(
944 operation
945 .unwrap_or_else(|| {
946 crate::SdkOperation::from_action_name(method)
947 .unwrap_or(crate::SdkOperation::Respond)
948 })
949 .action_name(),
950 ),
951 (Some("not_found"), Some(-32021)) => {
952 let request_id = params
953 .pointer("/response/request_id")
954 .and_then(Value::as_u64)
955 .unwrap_or_default();
956 FrontendRuntimeError::UnknownRequest(request_id)
957 }
958 (Some("invalid_argument"), _) => FrontendRuntimeError::InvalidResponse(message),
959 (_, Some(-32000)) => RuntimeSubmitError::Busy.into(),
960 (_, Some(-32001)) => RuntimeSubmitError::Interrupted.into(),
961 (_, Some(-32002)) => RuntimeSubmitError::Agent(message).into(),
962 (_, Some(-32020)) => FrontendRuntimeError::UnsupportedAction(
963 crate::SdkOperation::from_action_name(method)
964 .unwrap_or(crate::SdkOperation::Respond)
965 .action_name(),
966 ),
967 (_, Some(-32021)) => {
968 let request_id = params
969 .pointer("/response/request_id")
970 .and_then(Value::as_u64)
971 .unwrap_or_default();
972 FrontendRuntimeError::UnknownRequest(request_id)
973 }
974 (_, Some(-32022)) => FrontendRuntimeError::InvalidResponse(message),
975 _ => FrontendRuntimeError::Transport(message),
976 });
977 }
978 Ok(value.get("result").cloned().unwrap_or(Value::Null))
979 }
980
981 async fn rpc_typed<T: serde::de::DeserializeOwned>(
982 &self,
983 method: &str,
984 params: Value,
985 ) -> Result<T, FrontendRuntimeError> {
986 serde_json::from_value(self.rpc(method, params).await?)
987 .map_err(|error| FrontendRuntimeError::Transport(error.to_string()))
988 }
989
990 pub fn client_id(&self) -> &crate::RuntimeClientId {
992 &self.client_id
993 }
994
995 pub fn is_disconnected(&self) -> bool {
998 self.disconnected.load(Ordering::SeqCst)
999 }
1000
1001 pub async fn take_control(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1004 self.rpc_typed(
1005 crate::FrontendFacadeMethod::TakeControl.wire_name(),
1006 json!({}),
1007 )
1008 .await
1009 }
1010
1011 pub async fn acquire_control(
1013 &self,
1014 ) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1015 self.rpc_typed(
1016 crate::FrontendFacadeMethod::AcquireControl.wire_name(),
1017 json!({}),
1018 )
1019 .await
1020 }
1021
1022 pub async fn heartbeat(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1024 self.rpc_typed(
1025 crate::FrontendFacadeMethod::Heartbeat.wire_name(),
1026 json!({}),
1027 )
1028 .await
1029 }
1030
1031 pub async fn lease_snapshot(
1033 &self,
1034 ) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1035 self.rpc_typed(crate::FrontendFacadeMethod::Lease.wire_name(), json!({}))
1036 .await
1037 }
1038
1039 pub async fn detach(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1041 let snapshot = self
1042 .rpc_typed(crate::FrontendFacadeMethod::Detach.wire_name(), json!({}))
1043 .await?;
1044 self.disconnected.store(true, Ordering::SeqCst);
1045 Ok(snapshot)
1046 }
1047}
1048
1049#[async_trait]
1050#[cfg(feature = "adapter-api")]
1051impl FrontendRuntime for HttpFrontendRuntime {
1052 async fn describe(&self) -> Result<FrontendRuntimeDescriptor, FrontendRuntimeError> {
1053 self.rpc_typed(crate::FrontendFacadeMethod::Describe.wire_name(), json!({}))
1054 .await
1055 }
1056
1057 async fn attach(
1058 &self,
1059 history_limit: usize,
1060 ) -> Result<FrontendAttachment, FrontendRuntimeError> {
1061 if self.disconnected.load(Ordering::SeqCst) {
1062 return Err(FrontendRuntimeError::Closed);
1063 }
1064 let live = self.events.subscribe();
1068 let snapshot: FrontendAttachSnapshot = self
1069 .rpc_typed(
1070 crate::FrontendFacadeMethod::Attach.wire_name(),
1071 json!({"limit": history_limit}),
1072 )
1073 .await?;
1074 Ok(FrontendAttachment::new(
1075 snapshot.descriptor,
1076 snapshot.history,
1077 snapshot.history_cursor,
1078 snapshot.replay,
1079 live,
1080 Some(self.lifecycle.clone()),
1081 ))
1082 }
1083
1084 async fn send_input(self: Arc<Self>, prompt: String) -> Result<(), FrontendRuntimeError> {
1085 self.send_input_with_images(prompt, Vec::new()).await
1086 }
1087
1088 async fn send_input_with_images(
1089 self: Arc<Self>,
1090 prompt: String,
1091 image_urls: Vec<String>,
1092 ) -> Result<(), FrontendRuntimeError> {
1093 self.rpc(
1094 crate::FrontendFacadeMethod::SendInput.wire_name(),
1095 json!({"prompt": prompt, "image_urls": image_urls}),
1096 )
1097 .await?;
1098 Ok(())
1099 }
1100
1101 async fn submit(&self, prompt: String) -> Result<String, FrontendRuntimeError> {
1102 let result = self
1103 .rpc(
1104 crate::FrontendFacadeMethod::Submit.wire_name(),
1105 json!({"prompt": prompt}),
1106 )
1107 .await?;
1108 Ok(result
1109 .get("reply")
1110 .and_then(Value::as_str)
1111 .unwrap_or_default()
1112 .to_string())
1113 }
1114
1115 async fn submit_with_images(
1116 &self,
1117 prompt: String,
1118 image_urls: Vec<String>,
1119 ) -> Result<String, FrontendRuntimeError> {
1120 let result = self
1121 .rpc(
1122 crate::FrontendFacadeMethod::Submit.wire_name(),
1123 json!({"prompt": prompt, "image_urls": image_urls}),
1124 )
1125 .await?;
1126 Ok(result
1127 .get("reply")
1128 .and_then(Value::as_str)
1129 .unwrap_or_default()
1130 .to_string())
1131 }
1132
1133 async fn interrupt(&self) -> Result<bool, FrontendRuntimeError> {
1134 let result = self
1135 .rpc(
1136 crate::FrontendFacadeMethod::Interrupt.wire_name(),
1137 json!({}),
1138 )
1139 .await?;
1140 Ok(result
1141 .get("interrupted")
1142 .and_then(Value::as_bool)
1143 .unwrap_or(false))
1144 }
1145
1146 async fn steer(&self, prompt: String) -> Result<(), FrontendRuntimeError> {
1147 self.rpc(
1148 crate::FrontendFacadeMethod::Steer.wire_name(),
1149 json!({"prompt": prompt}),
1150 )
1151 .await?;
1152 Ok(())
1153 }
1154
1155 async fn respond(&self, response: FrontendResponse) -> Result<(), FrontendRuntimeError> {
1156 self.rpc(
1157 crate::FrontendFacadeMethod::Respond.wire_name(),
1158 json!({"response": response}),
1159 )
1160 .await?;
1161 Ok(())
1162 }
1163
1164 async fn invoke(
1165 &self,
1166 operation: FrontendOperationInvocation,
1167 ) -> Result<FrontendOperationResult, FrontendRuntimeError> {
1168 self.rpc_typed(
1169 crate::FrontendFacadeMethod::Invoke.wire_name(),
1170 json!({"operation": operation}),
1171 )
1172 .await
1173 }
1174
1175 async fn lease_snapshot(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1176 HttpFrontendRuntime::lease_snapshot(self).await
1177 }
1178
1179 async fn take_control(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1180 HttpFrontendRuntime::take_control(self).await
1181 }
1182
1183 async fn acquire_control(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1184 HttpFrontendRuntime::acquire_control(self).await
1185 }
1186
1187 async fn heartbeat(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1188 HttpFrontendRuntime::heartbeat(self).await
1189 }
1190
1191 async fn detach(&self) -> Result<crate::RuntimeLeaseSnapshot, FrontendRuntimeError> {
1192 HttpFrontendRuntime::detach(self).await
1193 }
1194
1195 async fn close(&self) -> Result<(), FrontendRuntimeError> {
1196 self.rpc(crate::FrontendFacadeMethod::Close.wire_name(), json!({}))
1197 .await?;
1198 self.disconnected.store(true, Ordering::SeqCst);
1199 Ok(())
1200 }
1201}