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