1use crate::shell_api::remoting;
3pub use crate::shell_api::{RemotingInputEvent, RemotingSignal, RemotingStartOptions};
4use crate::{ShellApiHost, ShellApiRequest};
5use base64::Engine as _;
6use serde::{Deserialize, Serialize};
7use tokio::sync::oneshot;
8
9#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
10#[serde(rename_all = "camelCase")]
11pub struct RemotingSession {
12 pub id: u64,
13 pub state: String,
14}
15
16#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
17#[serde(rename_all = "camelCase")]
18pub struct RemotingSnapshot {
19 pub width: u32,
20 pub height: u32,
21 pub mime_type: String,
22 pub data_base64: String,
23}
24
25#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
26#[serde(rename_all = "camelCase")]
27pub struct RemotingNetworkStats {
28 pub state: String,
29 pub route: String,
30 pub transport: String,
31 pub latency_ms: f64,
32 pub jitter_ms: f64,
33 pub bandwidth_kbps: f64,
34 pub packet_loss_percent: f64,
35}
36
37#[derive(Clone)]
38pub struct Remoting {
39 shell_api: Option<ShellApiHost>,
40 supported: std::sync::Arc<std::sync::atomic::AtomicBool>,
41}
42impl Remoting {
43 #[must_use]
44 pub fn new(shell_api: Option<ShellApiHost>) -> Self {
45 Self {
46 shell_api,
47 supported: std::sync::Arc::new(std::sync::atomic::AtomicBool::new(false)),
48 }
49 }
50 pub fn is_supported(&self) -> bool {
51 self.shell_api
52 .as_ref()
53 .is_some_and(ShellApiHost::is_available)
54 }
55 #[must_use]
56 pub fn supports_control(&self) -> bool {
57 self.supported.load(std::sync::atomic::Ordering::Acquire)
58 }
59 pub async fn refresh_capabilities(&self) {
60 let supported = match self.invoke(remoting::Capabilities {}).await {
61 Ok(capabilities) => {
62 tracing::debug!(available = capabilities.available, protocol = %capabilities.protocol, "native desktop capabilities");
63 if !capabilities.available || capabilities.protocol != "crd-webrtc" {
64 tracing::warn!(available = capabilities.available, protocol = %capabilities.protocol, "native desktop protocol is unavailable");
65 }
66 capabilities.available && capabilities.protocol == "crd-webrtc"
67 }
68 Err(error) => {
69 tracing::warn!(%error, "native desktop capability discovery failed");
70 false
71 }
72 };
73 self.supported
74 .store(supported, std::sync::atomic::Ordering::Release);
75 }
76 pub async fn start(&self, options: RemotingStartOptions) -> Result<RemotingSession, String> {
77 let shell_api = self
78 .shell_api
79 .as_ref()
80 .ok_or("native Shell API is unavailable")?;
81 let cleanup_api = shell_api.clone();
82 let (tx, rx) = oneshot::channel();
83 shell_api
84 .invoke(&remoting::Start { options }, move |result| {
85 if let Err(Ok(session)) = tx.send(result) {
86 let _ = cleanup_api.invoke(&remoting::Stop { id: session.id }, |_| {});
88 }
89 })
90 .map_err(|error| error.to_string())?;
91 let session = tokio::time::timeout(std::time::Duration::from_secs(5), rx)
92 .await
93 .map_err(|_| "native remoting startup timed out")?
94 .map_err(|_| "native remoting startup was cancelled")?
95 .map_err(|error| error.to_string())?;
96 let id = u64::try_from(session.id)
97 .ok()
98 .filter(|id| *id > 0)
99 .ok_or_else(|| "native CRD returned an invalid session ID".to_owned())?;
100 Ok(RemotingSession {
101 id,
102 state: session.state,
103 })
104 }
105 pub async fn process_signal(&self, id: u64, signal: RemotingSignal) -> Result<(), String> {
106 self.invoke(remoting::ProcessSignal {
107 id: session_id(id)?,
108 signal,
109 })
110 .await
111 }
112 pub async fn take_signals(
113 &self,
114 id: u64,
115 after_sequence: u64,
116 ) -> Result<Vec<RemotingSignal>, String> {
117 self.invoke(remoting::TakeSignals {
118 id: session_id(id)?,
119 after_sequence: Some(
120 i64::try_from(after_sequence).map_err(|_| "invalid signal sequence")?,
121 ),
122 })
123 .await
124 }
125 pub async fn input(&self, id: u64, event: RemotingInputEvent) -> Result<(), String> {
126 self.invoke(remoting::SendInput {
127 id: session_id(id)?,
128 event,
129 })
130 .await
131 }
132 pub async fn stop(&self, id: u64) -> Result<(), String> {
133 self.invoke(remoting::Stop {
134 id: session_id(id)?,
135 })
136 .await
137 }
138 pub async fn stats(&self, id: u64) -> Result<RemotingNetworkStats, String> {
139 let stats = self
140 .invoke(remoting::Stats {
141 id: session_id(id)?,
142 })
143 .await?;
144 Ok(RemotingNetworkStats {
145 state: stats.state,
146 route: stats.route,
147 transport: stats.transport,
148 latency_ms: stats.latency_ms,
149 jitter_ms: stats.jitter_ms,
150 bandwidth_kbps: stats.available_bandwidth_kbps,
151 packet_loss_percent: stats.packet_loss_percent,
152 })
153 }
154 pub async fn snapshot(&self, screen_id: Option<i64>) -> Result<RemotingSnapshot, String> {
155 let snapshot = self.invoke(remoting::Capture { screen_id }).await?;
156 Ok(RemotingSnapshot {
157 width: dimension(snapshot.width)?,
158 height: dimension(snapshot.height)?,
159 mime_type: snapshot.mime_type,
160 data_base64: base64::engine::general_purpose::STANDARD.encode(snapshot.data.0),
161 })
162 }
163 pub async fn inject_input(&self, event: RemotingInputEvent) -> Result<(), String> {
164 self.invoke(remoting::InjectInput { event }).await
165 }
166 async fn invoke<R: ShellApiRequest>(&self, request: R) -> Result<R::Response, String> {
167 let shell_api = self
168 .shell_api
169 .as_ref()
170 .ok_or_else(|| "native Shell API is unavailable".to_owned())?;
171 let (tx, rx) = oneshot::channel();
172 shell_api
173 .invoke(&request, move |result| {
174 let _ = tx.send(result);
175 })
176 .map_err(|error| format!("native CRD operation failed: {error}"))?;
177 tokio::time::timeout(std::time::Duration::from_secs(5), rx)
178 .await
179 .map_err(|_| "native remoting deadline exceeded".to_owned())?
180 .map_err(|_| "native CRD reply channel closed".to_owned())?
181 .map_err(|error| format!("native CRD operation failed: {error}"))
182 }
183}
184
185fn session_id(id: u64) -> Result<i64, String> {
186 i64::try_from(id)
187 .ok()
188 .filter(|id| *id > 0)
189 .ok_or_else(|| "invalid session ID".to_owned())
190}
191#[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)]
192fn dimension(value: f64) -> Result<u32, String> {
193 if !value.is_finite() || value.fract() != 0.0 || value <= 0.0 || value > f64::from(u32::MAX) {
194 return Err("native CRD returned an invalid snapshot dimension".to_owned());
195 }
196 Ok(value as u32)
197}