1use crate::*;
2use tokio::io::{AsyncReadExt, AsyncWriteExt};
3
4static SLOTS: tokio::sync::Semaphore = tokio::sync::Semaphore::const_new(4);
5
6pub async fn sample(app: Arc<App>) -> Result<Arc<Document>> {
7 app.cache
8 .get("host-sample".into(), Duration::from_secs(1), || async {
9 call(json!({"operation":"host.sample"})).await
10 })
11 .await
12}
13
14pub async fn call(request: Value) -> Result<Value> {
15 tokio::time::timeout(Duration::from_secs(70), async {
16 let _slot = SLOTS.acquire().await?;
17 let mut socket = tokio::net::UnixStream::connect(env(
18 "STUDIO_HOST_SOCKET",
19 "/run/studio-host/host.sock",
20 ))
21 .await?;
22 #[cfg(target_os = "linux")]
23 if socket.peer_cred()?.uid() != 0 {
24 return Err(Error::new(
25 502,
26 "The host service identity couldn't be verified.",
27 ));
28 }
29 let mut message = serde_json::to_vec(&request)?;
30 message.push(b'\n');
31 if message.len() > 65536 {
32 return Err(Error::new(
33 400,
34 "The host request is too large. Narrow the selection.",
35 ));
36 }
37 socket.write_all(&message).await?;
38 let length = socket.read_u32().await? as usize;
39 if length > 16 * 1024 * 1024 {
40 return Err(Error::new(
41 502,
42 "The host response is too large. Narrow the selection.",
43 ));
44 }
45 let mut bytes = vec![0; length];
46 socket.read_exact(&mut bytes).await?;
47 let mut response: Value = serde_json::from_slice(&bytes)?;
48 if let Some(message) = response["error"].as_str() {
49 return Err(Error::new(
50 response["status"]
51 .as_u64()
52 .filter(|s| (400..=599).contains(s))
53 .unwrap_or(502) as u16,
54 message,
55 ));
56 }
57 response
58 .as_object_mut()
59 .and_then(|v| v.remove("value"))
60 .ok_or_else(|| {
61 Error::new(
62 502,
63 "The host response is incomplete. Check its logs, then retry.",
64 )
65 })
66 })
67 .await
68 .map_err(|_| {
69 Error::new(
70 504,
71 "The host operation is taking too long. Check its logs, then retry.",
72 )
73 })?
74}