tuffite/
remoting.rs

1//! Asynchronous native desktop remoting client. No application policy or transport.
2use 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                    // A timed-out/cancelled creator must not strand a native session.
87                    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}