Skip to main content

isb_core/
client.rs

1//! A small, synchronous incusd REST client over the local unix socket.
2//!
3//! Deliberately not the `incus` binary: a CLI client that hangs (one `incus init`
4//! sat ~11 minutes with no matching server-side operation) cannot be diagnosed or
5//! bounded from outside. Here every request has a socket timeout, and every
6//! mutation is a server operation waited on with a deadline, so a stall surfaces
7//! as the step that stalled.
8
9use std::io::{Read, Write};
10use std::os::unix::net::UnixStream;
11use std::path::{Path, PathBuf};
12use std::time::{Duration, Instant};
13
14use serde::Deserialize;
15use serde_json::Value;
16
17use crate::error::{Error, Result};
18
19#[cfg(test)]
20pub(crate) mod fake;
21mod stream;
22
23/// The sentence for an incus (`GET /1.0` metadata) without the
24/// `instance_oci` API extension, which every OCI image needs.
25pub fn oci_unsupported(info: &Value) -> Option<String> {
26    let ext = info["api_extensions"].as_array()?;
27    if ext.iter().any(|e| e == "instance_oci") {
28        return None;
29    }
30    let version = info["environment"]["server_version"]
31        .as_str()
32        .unwrap_or("this version");
33    Some(format!(
34        "incus {version} is too old for OCI images (`docker:`, `ghcr:`, `registry:`), which need incus 6.3 or later; \
35on Ubuntu 24.04 the distro package is 6.0, so install incus from the Zabbly stable repository instead \
36(the install guide, step 1: incus)"
37    ))
38}
39
40/// Deadlines used by the client. Every request has one; there is no unbounded wait
41/// anywhere except the output of `exec`, which by design has no default timeout.
42#[derive(Debug, Clone)]
43pub struct Timeouts {
44    /// Socket timeout for a single request/response exchange.
45    pub request: Duration,
46    /// Deadline for the instance-create operation (includes unpacking the image).
47    pub create: Duration,
48    /// Deadline for start/stop/restart operations.
49    pub state: Duration,
50    /// Deadline for any other operation (config updates, deletes, volume creation).
51    pub other: Duration,
52    /// After a create stalls past its deadline: how long to let the (usually
53    /// non-cancellable) operation settle before cleaning up and retrying.
54    pub settle: Duration,
55}
56
57impl Default for Timeouts {
58    fn default() -> Self {
59        Timeouts {
60            request: Duration::from_secs(30),
61            create: Duration::from_secs(600),
62            state: Duration::from_secs(120),
63            other: Duration::from_secs(120),
64            settle: Duration::from_secs(30),
65        }
66    }
67}
68
69/// Connection to incusd.
70#[derive(Debug, Clone)]
71pub struct Client {
72    socket: PathBuf,
73    project: Option<String>,
74    #[doc(hidden)]
75    pub timeouts: Timeouts,
76}
77
78/// The standard incusd response envelope.
79#[derive(Debug, Deserialize)]
80struct Envelope {
81    #[serde(rename = "type", default)]
82    kind: String,
83    #[serde(default)]
84    error: String,
85    #[serde(default)]
86    metadata: Value,
87    #[serde(default)]
88    operation: String,
89}
90
91/// What a request returned: sync metadata, or an operation to wait on.
92#[derive(Debug)]
93#[doc(hidden)]
94pub enum Reply {
95    Sync(Value),
96    Async { operation: String, metadata: Value },
97}
98
99impl Client {
100    /// Locate the socket the way the incus tools do: `$INCUS_SOCKET`, else
101    /// `$INCUS_DIR/unix.socket`, else `/var/lib/incus/unix.socket`. On macOS
102    /// the last fallback is the default `isb machine`'s forwarded socket,
103    /// `~/.isb/machine/isb/incus.sock`.
104    pub fn default_socket() -> PathBuf {
105        if let Some(s) = std::env::var_os("INCUS_SOCKET").filter(|s| !s.is_empty()) {
106            return PathBuf::from(s);
107        }
108        if let Some(d) = std::env::var_os("INCUS_DIR").filter(|s| !s.is_empty()) {
109            return PathBuf::from(d).join("unix.socket");
110        }
111        #[cfg(target_os = "macos")]
112        if let Ok(s) = crate::machine::incus_socket(crate::machine::DEFAULT_NAME) {
113            return s;
114        }
115        PathBuf::from("/var/lib/incus/unix.socket")
116    }
117
118    /// Client for the default socket and the default project.
119    pub fn new() -> Self {
120        Client::with_socket(Client::default_socket())
121    }
122
123    pub fn with_socket(socket: impl Into<PathBuf>) -> Self {
124        Client {
125            socket: socket.into(),
126            project: None,
127            timeouts: Timeouts::default(),
128        }
129    }
130
131    /// Use an incus project other than `default`.
132    pub fn project(mut self, project: impl Into<String>) -> Self {
133        let p = project.into();
134        self.project = if p.is_empty() || p == "default" {
135            None
136        } else {
137            Some(p)
138        };
139        self
140    }
141
142    pub fn timeouts(mut self, timeouts: Timeouts) -> Self {
143        self.timeouts = timeouts;
144        self
145    }
146
147    pub fn socket(&self) -> &Path {
148        &self.socket
149    }
150
151    pub fn project_name(&self) -> &str {
152        self.project.as_deref().unwrap_or("default")
153    }
154
155    pub fn get_timeouts(&self) -> &Timeouts {
156        &self.timeouts
157    }
158
159    /// `GET /1.0`: server info. Also a cheap reachability check.
160    pub fn server_info(&self) -> Result<Value> {
161        self.get("/1.0")
162    }
163
164    /// Whether the server lists the API extension `name`.
165    pub fn has_extension(&self, name: &str) -> Result<bool> {
166        Ok(self.server_info()?["api_extensions"]
167            .as_array()
168            .is_some_and(|a| a.iter().any(|e| e == name)))
169    }
170
171    /// The server's version (`environment.server_version`), if it says.
172    pub fn server_version(&self) -> Result<Option<String>> {
173        Ok(self.server_info()?["environment"]["server_version"]
174            .as_str()
175            .map(str::to_string))
176    }
177
178    /// A sentence saying so when this incus cannot run OCI images
179    /// (`docker:`, `ghcr:`, `registry:`), else `None`. Also `None` when
180    /// incusd cannot be asked.
181    pub fn oci_unsupported(&self) -> Option<String> {
182        oci_unsupported(&self.server_info().ok()?)
183    }
184
185    fn with_project(&self, path: &str) -> String {
186        match &self.project {
187            None => path.to_string(),
188            Some(p) => {
189                let sep = if path.contains('?') { '&' } else { '?' };
190                format!("{path}{sep}project={}", encode_query(p))
191            }
192        }
193    }
194
195    fn connect(&self, timeout: Duration) -> Result<UnixStream> {
196        let stream = UnixStream::connect(&self.socket).map_err(|source| Error::Connect {
197            socket: self.socket.display().to_string(),
198            source,
199        })?;
200        stream.set_read_timeout(Some(timeout))?;
201        stream.set_write_timeout(Some(timeout))?;
202        Ok(stream)
203    }
204
205    /// One HTTP/1.1 exchange on a fresh connection, bounded by `timeout` overall.
206    fn raw(
207        &self,
208        method: &str,
209        path: &str,
210        body: Option<&Value>,
211        if_match: Option<&str>,
212        timeout: Duration,
213    ) -> Result<RawResponse> {
214        let payload = match body {
215            Some(v) => serde_json::to_vec(v)?,
216            None => Vec::new(),
217        };
218        let mut headers: Vec<(&str, String)> = Vec::new();
219        if body.is_some() {
220            headers.push(("Content-Type", "application/json".into()));
221        }
222        if let Some(etag) = if_match {
223            headers.push(("If-Match", etag.to_string()));
224        }
225        self.raw_bytes(method, path, &payload, &headers, timeout)
226    }
227
228    /// [`Client::raw`] with an arbitrary body and headers (the file API).
229    fn raw_bytes(
230        &self,
231        method: &str,
232        path: &str,
233        payload: &[u8],
234        extra_headers: &[(&str, String)],
235        timeout: Duration,
236    ) -> Result<RawResponse> {
237        let path = self.with_project(path);
238        let started = Instant::now();
239        let to_err = |e: std::io::Error| -> Error {
240            if matches!(
241                e.kind(),
242                std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
243            ) {
244                Error::RequestTimeout {
245                    method: method.to_string(),
246                    path: path.clone(),
247                    timeout,
248                }
249            } else {
250                Error::Io(e)
251            }
252        };
253        let mut stream = self.connect(timeout)?;
254        let mut head = format!(
255            "{method} {path} HTTP/1.1\r\nHost: incus\r\nUser-Agent: isb/{}\r\nConnection: close\r\n",
256            env!("CARGO_PKG_VERSION")
257        );
258        for (k, v) in extra_headers {
259            head.push_str(&format!("{k}: {v}\r\n"));
260        }
261        head.push_str(&format!("Content-Length: {}\r\n\r\n", payload.len()));
262        stream.write_all(head.as_bytes()).map_err(to_err)?;
263        stream.write_all(payload).map_err(to_err)?;
264        stream.flush().map_err(to_err)?;
265
266        let mut buf = Vec::with_capacity(8192);
267        let mut chunk = [0u8; 16384];
268        loop {
269            let remaining = timeout.saturating_sub(started.elapsed());
270            if remaining.is_zero() {
271                return Err(to_err(std::io::ErrorKind::TimedOut.into()));
272            }
273            // macOS refuses setsockopt (EINVAL) once incusd has closed; the
274            // read cannot block then, so the previous timeout is as good.
275            let _ = stream.set_read_timeout(Some(remaining));
276            let n = stream.read(&mut chunk).map_err(to_err)?;
277            if n == 0 {
278                break;
279            }
280            buf.extend_from_slice(&chunk[..n]);
281            // Stop as soon as a complete response is in hand; incusd honours
282            // Connection: close, but there is no need to depend on it.
283            if let Some(done) = complete_response(&buf)? {
284                return Ok(done);
285            }
286        }
287        complete_response(&buf)?
288            .or_else(|| eof_body(&buf))
289            .ok_or_else(|| Error::Protocol(format!("truncated response to {method} {path}")))
290    }
291
292    /// Send a request and decode the incusd envelope.
293    pub(crate) fn request(
294        &self,
295        method: &str,
296        path: &str,
297        body: Option<&Value>,
298        timeout: Duration,
299    ) -> Result<Reply> {
300        self.request_etag(method, path, body, None, timeout)
301            .map(|(r, _)| r)
302    }
303
304    /// Like [`Client::request`], optionally sending `If-Match`, and returning the
305    /// response ETag.
306    pub(crate) fn request_etag(
307        &self,
308        method: &str,
309        path: &str,
310        body: Option<&Value>,
311        if_match: Option<&str>,
312        timeout: Duration,
313    ) -> Result<(Reply, Option<String>)> {
314        let RawResponse {
315            status,
316            body: bytes,
317            etag,
318        } = self.raw(method, path, body, if_match, timeout)?;
319        let env: Envelope = serde_json::from_slice(&bytes).map_err(|e| {
320            Error::Protocol(format!(
321                "{method} {path}: HTTP {status}, undecodable body ({e}): {}",
322                String::from_utf8_lossy(&bytes[..bytes.len().min(200)])
323            ))
324        })?;
325        if env.kind == "error" || status >= 400 {
326            let request = body.map(|b| (path, b));
327            if let Some(e) = crate::org::limits::translate(self, &env.error, request) {
328                return Err(e);
329            }
330            return Err(Error::Api {
331                method: method.to_string(),
332                path: path.to_string(),
333                status,
334                message: if env.error.is_empty() {
335                    format!("HTTP {status}")
336                } else {
337                    env.error
338                },
339            });
340        }
341        if env.kind == "async" {
342            return Ok((
343                Reply::Async {
344                    operation: env.operation,
345                    metadata: env.metadata,
346                },
347                etag,
348            ));
349        }
350        Ok((Reply::Sync(env.metadata), etag))
351    }
352
353    /// `GET` returning the body and its ETag.
354    #[doc(hidden)]
355    pub fn get_etag(&self, path: &str) -> Result<(Value, Option<String>)> {
356        match self.request_etag("GET", path, None, None, self.timeouts.request)? {
357            (Reply::Sync(v), e) => Ok((v, e)),
358            (Reply::Async { metadata, .. }, e) => Ok((metadata, e)),
359        }
360    }
361
362    /// A mutation guarded by `If-Match`, waited on like [`Client::mutate`].
363    #[doc(hidden)]
364    pub fn mutate_if_match(
365        &self,
366        method: &str,
367        path: &str,
368        body: &Value,
369        etag: Option<&str>,
370        step: &str,
371        deadline: Duration,
372    ) -> Result<Value> {
373        match self.request_etag(method, path, Some(body), etag, self.timeouts.request)? {
374            (Reply::Sync(v), _) => Ok(v),
375            (Reply::Async { operation, .. }, _) => self.wait_operation(&operation, step, deadline),
376        }
377    }
378
379    #[doc(hidden)]
380    pub fn get(&self, path: &str) -> Result<Value> {
381        match self.request("GET", path, None, self.timeouts.request)? {
382            Reply::Sync(v) => Ok(v),
383            Reply::Async { metadata, .. } => Ok(metadata),
384        }
385    }
386
387    /// `GET` returning `None` on 404.
388    #[doc(hidden)]
389    pub fn get_opt(&self, path: &str) -> Result<Option<Value>> {
390        match self.get(path) {
391            Ok(v) => Ok(Some(v)),
392            Err(e) if e.is_not_found() => Ok(None),
393            Err(e) => Err(e),
394        }
395    }
396
397    /// Perform a mutation and, if it is an operation, wait for it under `deadline`.
398    /// On deadline the operation is cancelled (if incus allows it) and
399    /// [`Error::OperationTimeout`] names `step`.
400    #[doc(hidden)]
401    pub fn mutate(
402        &self,
403        method: &str,
404        path: &str,
405        body: Option<&Value>,
406        step: &str,
407        deadline: Duration,
408    ) -> Result<Value> {
409        match self.request(method, path, body, self.timeouts.request)? {
410            Reply::Sync(v) => Ok(v),
411            Reply::Async { operation, .. } => self.wait_operation(&operation, step, deadline),
412        }
413    }
414
415    /// Wait for an operation to finish. Returns its metadata on success.
416    pub fn wait_operation(&self, operation: &str, step: &str, deadline: Duration) -> Result<Value> {
417        let started = Instant::now();
418        loop {
419            let remaining = deadline.saturating_sub(started.elapsed());
420            // incus takes whole seconds (0 = answer now). Poll in slices of at most
421            // 30s so no single request outlives its socket timeout by much, and
422            // finish a sub-second remainder with short client-side sleeps.
423            let secs = remaining.as_secs().min(30);
424            if let Some(v) = self.poll_operation(operation, step, secs)? {
425                return Ok(v);
426            }
427            if started.elapsed() >= deadline {
428                let status = self
429                    .get(operation)
430                    .ok()
431                    .and_then(|op| op.get("status").and_then(Value::as_str).map(String::from))
432                    .unwrap_or_else(|| "running".into())
433                    .to_lowercase();
434                let cancelled = self.cancel_operation(operation).is_ok();
435                return Err(Error::OperationTimeout {
436                    step: step.to_string(),
437                    operation: operation.to_string(),
438                    status,
439                    waited: started.elapsed(),
440                    cancelled,
441                });
442            }
443            if secs == 0 {
444                std::thread::sleep(remaining.min(Duration::from_millis(50)));
445            }
446        }
447    }
448
449    /// Wait up to `secs` whole seconds (0 = just look) for an operation.
450    /// `Ok(Some(metadata))` on success, `Ok(None)` while it is still running.
451    pub(crate) fn poll_operation(
452        &self,
453        operation: &str,
454        step: &str,
455        secs: u64,
456    ) -> Result<Option<Value>> {
457        let op = self.get_with_timeout(
458            &format!("{operation}/wait?timeout={secs}"),
459            Duration::from_secs(secs) + self.timeouts.request,
460        )?;
461        let code = op.get("status_code").and_then(Value::as_i64).unwrap_or(0);
462        match code {
463            200 => Ok(Some(op.get("metadata").cloned().unwrap_or(Value::Null))),
464            400 | 401 => {
465                let err = op
466                    .get("err")
467                    .and_then(Value::as_str)
468                    .filter(|s| !s.is_empty())
469                    .unwrap_or(if code == 401 {
470                        "operation cancelled"
471                    } else {
472                        "operation failed"
473                    });
474                if let Some(e) = crate::org::limits::translate(self, err, None) {
475                    return Err(e);
476                }
477                Err(Error::OperationFailed {
478                    step: step.to_string(),
479                    message: err.to_string(),
480                })
481            }
482            _ => Ok(None),
483        }
484    }
485
486    /// Start an operation without waiting for it; returns its path
487    /// (`/1.0/operations/<id>`). Low-level: most callers want the sandbox API.
488    pub fn start_operation(
489        &self,
490        method: &str,
491        path: &str,
492        body: Option<&Value>,
493    ) -> Result<String> {
494        match self.request(method, path, body, self.timeouts.request)? {
495            Reply::Async { operation, .. } => Ok(operation),
496            Reply::Sync(_) => Err(Error::Protocol(format!(
497                "{method} {path} did not start an operation"
498            ))),
499        }
500    }
501
502    fn get_with_timeout(&self, path: &str, timeout: Duration) -> Result<Value> {
503        match self.request("GET", path, None, timeout)? {
504            Reply::Sync(v) => Ok(v),
505            Reply::Async { metadata, .. } => Ok(metadata),
506        }
507    }
508
509    /// `DELETE /1.0/operations/<id>`. Fails for operations incus will not cancel.
510    pub fn cancel_operation(&self, operation: &str) -> Result<()> {
511        self.request("DELETE", operation, None, self.timeouts.request)?;
512        Ok(())
513    }
514
515    /// Write a file into an instance (running or stopped), replacing it.
516    /// Parent directories must exist; see [`Client::make_dir`].
517    pub fn push_file(
518        &self,
519        instance: &str,
520        path: &str,
521        data: &[u8],
522        uid: u32,
523        gid: u32,
524        mode: u32,
525    ) -> Result<()> {
526        self.file_request(instance, path, data, uid, gid, mode, "file")
527    }
528
529    /// Create a directory in an instance; an existing one is left as is.
530    pub fn make_dir(
531        &self,
532        instance: &str,
533        path: &str,
534        uid: u32,
535        gid: u32,
536        mode: u32,
537    ) -> Result<()> {
538        match self.file_request(instance, path, &[], uid, gid, mode, "directory") {
539            Err(Error::Api {
540                status,
541                ref message,
542                ..
543            }) if status == 409 || message.to_ascii_lowercase().contains("exist") => Ok(()),
544            r => r,
545        }
546    }
547
548    #[expect(clippy::too_many_arguments)]
549    fn file_request(
550        &self,
551        instance: &str,
552        path: &str,
553        data: &[u8],
554        uid: u32,
555        gid: u32,
556        mode: u32,
557        kind: &str,
558    ) -> Result<()> {
559        let url = format!(
560            "/1.0/instances/{}/files?path={}",
561            encode_segment(instance),
562            encode_query(path)
563        );
564        let headers = [
565            ("Content-Type", "application/octet-stream".to_string()),
566            ("X-Incus-uid", uid.to_string()),
567            ("X-Incus-gid", gid.to_string()),
568            ("X-Incus-mode", format!("{mode:04o}")),
569            ("X-Incus-type", kind.to_string()),
570            ("X-Incus-write", "overwrite".to_string()),
571        ];
572        let r = self.raw_bytes("POST", &url, data, &headers, self.timeouts.request)?;
573        if r.status >= 400 {
574            let message = serde_json::from_slice::<Envelope>(&r.body)
575                .map(|e| e.error)
576                .ok()
577                .filter(|e| !e.is_empty())
578                .unwrap_or_else(|| format!("HTTP {}", r.status));
579            return Err(Error::Api {
580                method: "POST".into(),
581                path: url,
582                status: r.status,
583                message,
584            });
585        }
586        Ok(())
587    }
588
589    /// `GET path` upgraded to `protocol` (incus' `/sftp`): the raw stream
590    /// once incusd answers 101, with `timeout` on every read and write.
591    pub(crate) fn upgrade(
592        &self,
593        path: &str,
594        protocol: &str,
595        timeout: Duration,
596    ) -> Result<UnixStream> {
597        let path = self.with_project(path);
598        let mut stream = self.connect(timeout)?;
599        let head = format!(
600            "GET {path} HTTP/1.1\r\nHost: incus\r\nUser-Agent: isb/{}\r\nUpgrade: {protocol}\r\nConnection: Upgrade\r\n\r\n",
601            env!("CARGO_PKG_VERSION")
602        );
603        stream.write_all(head.as_bytes())?;
604        // Byte by byte, so nothing past the headers (the protocol's own
605        // first bytes) is consumed here.
606        let mut buf = Vec::new();
607        let mut b = [0u8; 1];
608        while !buf.ends_with(b"\r\n\r\n") {
609            if buf.len() > 16384 || stream.read(&mut b)? == 0 {
610                return Err(Error::Protocol(format!(
611                    "no answer to the upgrade of {path}"
612                )));
613            }
614            buf.push(b[0]);
615        }
616        let head = String::from_utf8_lossy(&buf);
617        let status = head.split_whitespace().nth(1).unwrap_or_default();
618        if status != "101" {
619            return Err(Error::Api {
620                method: "GET".into(),
621                path,
622                status: status.parse().unwrap_or(0),
623                message: head.lines().next().unwrap_or_default().to_string(),
624            });
625        }
626        Ok(stream)
627    }
628
629    /// A `GET` whose answer is not the JSON envelope (`/1.0/metrics`).
630    pub(crate) fn get_raw(&self, path: &str) -> Result<Vec<u8>> {
631        let r = self.raw_bytes("GET", path, &[], &[], self.timeouts.request)?;
632        match r.status {
633            200 => Ok(r.body),
634            status => Err(Error::Api {
635                method: "GET".into(),
636                path: path.into(),
637                status,
638                message: String::from_utf8_lossy(&r.body[..r.body.len().min(200)]).into_owned(),
639            }),
640        }
641    }
642
643    /// An instance's console log as incus keeps it (an OCI app's output).
644    pub fn console_log(&self, instance: &str) -> Result<Vec<u8>> {
645        let url = format!("/1.0/instances/{}/console", encode_segment(instance));
646        let r = self.raw_bytes("GET", &url, &[], &[], self.timeouts.request)?;
647        match r.status {
648            200 => Ok(r.body),
649            404 => Ok(Vec::new()),
650            status => Err(Error::Api {
651                method: "GET".into(),
652                path: url,
653                status,
654                message: serde_json::from_slice::<Envelope>(&r.body)
655                    .map(|e| e.error)
656                    .unwrap_or_default(),
657            }),
658        }
659    }
660
661    /// Read a file from an instance; `None` if it does not exist.
662    pub fn read_file(&self, instance: &str, path: &str) -> Result<Option<Vec<u8>>> {
663        let url = format!(
664            "/1.0/instances/{}/files?path={}",
665            encode_segment(instance),
666            encode_query(path)
667        );
668        let r = self.raw_bytes("GET", &url, &[], &[], self.timeouts.request)?;
669        match r.status {
670            200 => Ok(Some(r.body)),
671            404 => Ok(None),
672            status => Err(Error::Api {
673                method: "GET".into(),
674                path: url,
675                status,
676                message: serde_json::from_slice::<Envelope>(&r.body)
677                    .map(|e| e.error)
678                    .unwrap_or_default(),
679            }),
680        }
681    }
682
683    /// Follow `/1.0/events?QUERY` (all projects when the query says so) as
684    /// a websocket. Reads time out every 5 s so a follower can check
685    /// whether to stop; a timeout is not an error.
686    pub fn events_websocket(&self, query: &str) -> Result<tungstenite::WebSocket<UnixStream>> {
687        let stream = self.connect(self.timeouts.request)?;
688        let url = format!("ws://incus/1.0/events?{query}");
689        let (ws, _resp) = tungstenite::client::client(url.as_str(), stream)
690            .map_err(|e| Error::WebSocket(format!("handshake for /1.0/events: {e}")))?;
691        ws.get_ref()
692            .set_read_timeout(Some(Duration::from_secs(5)))?;
693        ws.get_ref()
694            .set_write_timeout(Some(self.timeouts.request))?;
695        Ok(ws)
696    }
697
698    /// Open one of an operation's websockets (exec stdin/stdout/stderr/control).
699    pub(crate) fn websocket(
700        &self,
701        operation: &str,
702        secret: &str,
703    ) -> Result<tungstenite::WebSocket<UnixStream>> {
704        let path = self.with_project(&format!(
705            "{operation}/websocket?secret={}",
706            encode_query(secret)
707        ));
708        let stream = self.connect(self.timeouts.request)?;
709        let url = format!("ws://incus{path}");
710        let (ws, _resp) = tungstenite::client::client(url.as_str(), stream)
711            .map_err(|e| Error::WebSocket(format!("handshake for {operation}: {e}")))?;
712        // Exec output has no default timeout: a quiet process is not a stuck one.
713        ws.get_ref().set_read_timeout(None)?;
714        ws.get_ref()
715            .set_write_timeout(Some(self.timeouts.request))?;
716        Ok(ws)
717    }
718}
719
720impl Default for Client {
721    fn default() -> Self {
722        Client::new()
723    }
724}
725
726/// Percent-encode a query value (RFC 3986 unreserved set passes through).
727pub(crate) fn encode_query(s: &str) -> String {
728    let mut out = String::with_capacity(s.len());
729    for b in s.bytes() {
730        if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~') {
731            out.push(b as char);
732        } else {
733            out.push_str(&format!("%{b:02X}"));
734        }
735    }
736    out
737}
738
739/// Percent-encode one path segment.
740#[doc(hidden)]
741pub fn encode_segment(s: &str) -> String {
742    encode_query(s)
743}
744
745#[derive(Debug, PartialEq)]
746struct RawResponse {
747    status: u16,
748    body: Vec<u8>,
749    etag: Option<String>,
750}
751
752/// If `buf` holds a complete HTTP response, return it with its body decoded.
753fn complete_response(buf: &[u8]) -> Result<Option<RawResponse>> {
754    let mut headers = [httparse::EMPTY_HEADER; 64];
755    let mut resp = httparse::Response::new(&mut headers);
756    let head_len = match resp
757        .parse(buf)
758        .map_err(|e| Error::Protocol(format!("bad HTTP response: {e}")))?
759    {
760        httparse::Status::Complete(n) => n,
761        httparse::Status::Partial => return Ok(None),
762    };
763    let status = resp.code.unwrap_or(0);
764    let mut content_length: Option<usize> = None;
765    let mut chunked = false;
766    let mut etag = None;
767    for h in resp.headers.iter() {
768        if h.name.eq_ignore_ascii_case("etag") {
769            etag = Some(String::from_utf8_lossy(h.value).trim().to_string());
770        }
771        if h.name.eq_ignore_ascii_case("content-length") {
772            content_length = std::str::from_utf8(h.value)
773                .ok()
774                .and_then(|v| v.trim().parse().ok());
775        } else if h.name.eq_ignore_ascii_case("transfer-encoding")
776            && String::from_utf8_lossy(h.value)
777                .to_ascii_lowercase()
778                .contains("chunked")
779        {
780            chunked = true;
781        }
782    }
783    let body = &buf[head_len..];
784    if chunked {
785        return Ok(decode_chunked(body).map(|b| RawResponse {
786            status,
787            body: b,
788            etag,
789        }));
790    }
791    match content_length {
792        Some(n) if body.len() >= n => Ok(Some(RawResponse {
793            status,
794            body: body[..n].to_vec(),
795            etag,
796        })),
797        Some(_) => Ok(None),
798        // No length and not chunked: body runs to EOF; the caller decides.
799        None => Ok(None),
800    }
801}
802
803/// Decode a chunked body; `None` if it is not complete yet.
804fn decode_chunked(mut body: &[u8]) -> Option<Vec<u8>> {
805    let mut out = Vec::new();
806    loop {
807        let line_end = body.windows(2).position(|w| w == b"\r\n")?;
808        let size_str = std::str::from_utf8(&body[..line_end]).ok()?;
809        let size = usize::from_str_radix(size_str.split(';').next()?.trim(), 16).ok()?;
810        body = &body[line_end + 2..];
811        if size == 0 {
812            return Some(out);
813        }
814        if body.len() < size + 2 {
815            return None;
816        }
817        out.extend_from_slice(&body[..size]);
818        body = &body[size + 2..];
819    }
820}
821
822/// When a response had no length and was cut by EOF, treat what we have as the body.
823fn eof_body(buf: &[u8]) -> Option<RawResponse> {
824    let mut headers = [httparse::EMPTY_HEADER; 64];
825    let mut resp = httparse::Response::new(&mut headers);
826    match resp.parse(buf).ok()? {
827        httparse::Status::Complete(n) => Some(RawResponse {
828            status: resp.code?,
829            body: buf[n..].to_vec(),
830            etag: None,
831        }),
832        httparse::Status::Partial => None,
833    }
834}
835
836#[cfg(test)]
837mod tests {
838    use super::*;
839
840    #[test]
841    fn an_incus_without_oci_images_is_named_with_the_way_out() {
842        let old = serde_json::json!({
843            "api_extensions": ["disk_initial_copy"],
844            "environment": {"server_version": "6.0.4"},
845        });
846        let m = oci_unsupported(&old).unwrap();
847        assert!(m.contains("incus 6.0.4"), "{m}");
848        assert!(m.contains("Zabbly"), "{m}");
849        assert!(m.contains("6.3"), "{m}");
850        let new = serde_json::json!({
851            "api_extensions": ["disk_initial_copy", "instance_oci"],
852            "environment": {"server_version": "7.5.1"},
853        });
854        assert_eq!(oci_unsupported(&new), None);
855        // An answer without the list (an odd proxy) is not a verdict.
856        assert_eq!(oci_unsupported(&serde_json::json!({})), None);
857    }
858
859    #[test]
860    fn parses_content_length_response() {
861        let raw = b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nETag: \"abc\"\r\n\r\nhello";
862        let r = complete_response(raw).unwrap().unwrap();
863        assert_eq!((r.status, r.body.as_slice()), (200, &b"hello"[..]));
864        assert_eq!(r.etag.as_deref(), Some("\"abc\""));
865        let partial = b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nhel";
866        assert_eq!(complete_response(partial).unwrap(), None);
867    }
868
869    #[test]
870    fn parses_chunked_response() {
871        let raw = b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nhello\r\n6\r\n world\r\n0\r\n\r\n";
872        let r = complete_response(raw).unwrap().unwrap();
873        assert_eq!(r.body, b"hello world".to_vec());
874        let partial = b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nhel";
875        assert_eq!(complete_response(partial).unwrap(), None);
876    }
877
878    #[test]
879    fn encodes_query_values() {
880        assert_eq!(encode_query("a b/c"), "a%20b%2Fc");
881        assert_eq!(encode_query("abc-_.~"), "abc-_.~");
882    }
883}