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            if let Some(e) = crate::org::limits::translate(self, &env.error) {
327                return Err(e);
328            }
329            return Err(Error::Api {
330                method: method.to_string(),
331                path: path.to_string(),
332                status,
333                message: if env.error.is_empty() {
334                    format!("HTTP {status}")
335                } else {
336                    env.error
337                },
338            });
339        }
340        if env.kind == "async" {
341            return Ok((
342                Reply::Async {
343                    operation: env.operation,
344                    metadata: env.metadata,
345                },
346                etag,
347            ));
348        }
349        Ok((Reply::Sync(env.metadata), etag))
350    }
351
352    /// `GET` returning the body and its ETag.
353    #[doc(hidden)]
354    pub fn get_etag(&self, path: &str) -> Result<(Value, Option<String>)> {
355        match self.request_etag("GET", path, None, None, self.timeouts.request)? {
356            (Reply::Sync(v), e) => Ok((v, e)),
357            (Reply::Async { metadata, .. }, e) => Ok((metadata, e)),
358        }
359    }
360
361    /// A mutation guarded by `If-Match`, waited on like [`Client::mutate`].
362    #[doc(hidden)]
363    pub fn mutate_if_match(
364        &self,
365        method: &str,
366        path: &str,
367        body: &Value,
368        etag: Option<&str>,
369        step: &str,
370        deadline: Duration,
371    ) -> Result<Value> {
372        match self.request_etag(method, path, Some(body), etag, self.timeouts.request)? {
373            (Reply::Sync(v), _) => Ok(v),
374            (Reply::Async { operation, .. }, _) => self.wait_operation(&operation, step, deadline),
375        }
376    }
377
378    #[doc(hidden)]
379    pub fn get(&self, path: &str) -> Result<Value> {
380        match self.request("GET", path, None, self.timeouts.request)? {
381            Reply::Sync(v) => Ok(v),
382            Reply::Async { metadata, .. } => Ok(metadata),
383        }
384    }
385
386    /// `GET` returning `None` on 404.
387    #[doc(hidden)]
388    pub fn get_opt(&self, path: &str) -> Result<Option<Value>> {
389        match self.get(path) {
390            Ok(v) => Ok(Some(v)),
391            Err(e) if e.is_not_found() => Ok(None),
392            Err(e) => Err(e),
393        }
394    }
395
396    /// Perform a mutation and, if it is an operation, wait for it under `deadline`.
397    /// On deadline the operation is cancelled (if incus allows it) and
398    /// [`Error::OperationTimeout`] names `step`.
399    #[doc(hidden)]
400    pub fn mutate(
401        &self,
402        method: &str,
403        path: &str,
404        body: Option<&Value>,
405        step: &str,
406        deadline: Duration,
407    ) -> Result<Value> {
408        match self.request(method, path, body, self.timeouts.request)? {
409            Reply::Sync(v) => Ok(v),
410            Reply::Async { operation, .. } => self.wait_operation(&operation, step, deadline),
411        }
412    }
413
414    /// Wait for an operation to finish. Returns its metadata on success.
415    pub fn wait_operation(&self, operation: &str, step: &str, deadline: Duration) -> Result<Value> {
416        let started = Instant::now();
417        loop {
418            let remaining = deadline.saturating_sub(started.elapsed());
419            // incus takes whole seconds (0 = answer now). Poll in slices of at most
420            // 30s so no single request outlives its socket timeout by much, and
421            // finish a sub-second remainder with short client-side sleeps.
422            let secs = remaining.as_secs().min(30);
423            if let Some(v) = self.poll_operation(operation, step, secs)? {
424                return Ok(v);
425            }
426            if started.elapsed() >= deadline {
427                let status = self
428                    .get(operation)
429                    .ok()
430                    .and_then(|op| op.get("status").and_then(Value::as_str).map(String::from))
431                    .unwrap_or_else(|| "running".into())
432                    .to_lowercase();
433                let cancelled = self.cancel_operation(operation).is_ok();
434                return Err(Error::OperationTimeout {
435                    step: step.to_string(),
436                    operation: operation.to_string(),
437                    status,
438                    waited: started.elapsed(),
439                    cancelled,
440                });
441            }
442            if secs == 0 {
443                std::thread::sleep(remaining.min(Duration::from_millis(50)));
444            }
445        }
446    }
447
448    /// Wait up to `secs` whole seconds (0 = just look) for an operation.
449    /// `Ok(Some(metadata))` on success, `Ok(None)` while it is still running.
450    pub(crate) fn poll_operation(
451        &self,
452        operation: &str,
453        step: &str,
454        secs: u64,
455    ) -> Result<Option<Value>> {
456        let op = self.get_with_timeout(
457            &format!("{operation}/wait?timeout={secs}"),
458            Duration::from_secs(secs) + self.timeouts.request,
459        )?;
460        let code = op.get("status_code").and_then(Value::as_i64).unwrap_or(0);
461        match code {
462            200 => Ok(Some(op.get("metadata").cloned().unwrap_or(Value::Null))),
463            400 | 401 => {
464                let err = op
465                    .get("err")
466                    .and_then(Value::as_str)
467                    .filter(|s| !s.is_empty())
468                    .unwrap_or(if code == 401 {
469                        "operation cancelled"
470                    } else {
471                        "operation failed"
472                    });
473                if let Some(e) = crate::org::limits::translate(self, err) {
474                    return Err(e);
475                }
476                Err(Error::OperationFailed {
477                    step: step.to_string(),
478                    message: err.to_string(),
479                })
480            }
481            _ => Ok(None),
482        }
483    }
484
485    /// Start an operation without waiting for it; returns its path
486    /// (`/1.0/operations/<id>`). Low-level: most callers want the sandbox API.
487    pub fn start_operation(
488        &self,
489        method: &str,
490        path: &str,
491        body: Option<&Value>,
492    ) -> Result<String> {
493        match self.request(method, path, body, self.timeouts.request)? {
494            Reply::Async { operation, .. } => Ok(operation),
495            Reply::Sync(_) => Err(Error::Protocol(format!(
496                "{method} {path} did not start an operation"
497            ))),
498        }
499    }
500
501    fn get_with_timeout(&self, path: &str, timeout: Duration) -> Result<Value> {
502        match self.request("GET", path, None, timeout)? {
503            Reply::Sync(v) => Ok(v),
504            Reply::Async { metadata, .. } => Ok(metadata),
505        }
506    }
507
508    /// `DELETE /1.0/operations/<id>`. Fails for operations incus will not cancel.
509    pub fn cancel_operation(&self, operation: &str) -> Result<()> {
510        self.request("DELETE", operation, None, self.timeouts.request)?;
511        Ok(())
512    }
513
514    /// Write a file into an instance (running or stopped), replacing it.
515    /// Parent directories must exist; see [`Client::make_dir`].
516    pub fn push_file(
517        &self,
518        instance: &str,
519        path: &str,
520        data: &[u8],
521        uid: u32,
522        gid: u32,
523        mode: u32,
524    ) -> Result<()> {
525        self.file_request(instance, path, data, uid, gid, mode, "file")
526    }
527
528    /// Create a directory in an instance; an existing one is left as is.
529    pub fn make_dir(
530        &self,
531        instance: &str,
532        path: &str,
533        uid: u32,
534        gid: u32,
535        mode: u32,
536    ) -> Result<()> {
537        match self.file_request(instance, path, &[], uid, gid, mode, "directory") {
538            Err(Error::Api {
539                status,
540                ref message,
541                ..
542            }) if status == 409 || message.to_ascii_lowercase().contains("exist") => Ok(()),
543            r => r,
544        }
545    }
546
547    #[expect(clippy::too_many_arguments)]
548    fn file_request(
549        &self,
550        instance: &str,
551        path: &str,
552        data: &[u8],
553        uid: u32,
554        gid: u32,
555        mode: u32,
556        kind: &str,
557    ) -> Result<()> {
558        let url = format!(
559            "/1.0/instances/{}/files?path={}",
560            encode_segment(instance),
561            encode_query(path)
562        );
563        let headers = [
564            ("Content-Type", "application/octet-stream".to_string()),
565            ("X-Incus-uid", uid.to_string()),
566            ("X-Incus-gid", gid.to_string()),
567            ("X-Incus-mode", format!("{mode:04o}")),
568            ("X-Incus-type", kind.to_string()),
569            ("X-Incus-write", "overwrite".to_string()),
570        ];
571        let r = self.raw_bytes("POST", &url, data, &headers, self.timeouts.request)?;
572        if r.status >= 400 {
573            let message = serde_json::from_slice::<Envelope>(&r.body)
574                .map(|e| e.error)
575                .ok()
576                .filter(|e| !e.is_empty())
577                .unwrap_or_else(|| format!("HTTP {}", r.status));
578            return Err(Error::Api {
579                method: "POST".into(),
580                path: url,
581                status: r.status,
582                message,
583            });
584        }
585        Ok(())
586    }
587
588    /// A `GET` whose answer is not the JSON envelope (`/1.0/metrics`).
589    pub(crate) fn get_raw(&self, path: &str) -> Result<Vec<u8>> {
590        let r = self.raw_bytes("GET", path, &[], &[], self.timeouts.request)?;
591        match r.status {
592            200 => Ok(r.body),
593            status => Err(Error::Api {
594                method: "GET".into(),
595                path: path.into(),
596                status,
597                message: String::from_utf8_lossy(&r.body[..r.body.len().min(200)]).into_owned(),
598            }),
599        }
600    }
601
602    /// An instance's console log as incus keeps it (an OCI app's output).
603    pub fn console_log(&self, instance: &str) -> Result<Vec<u8>> {
604        let url = format!("/1.0/instances/{}/console", encode_segment(instance));
605        let r = self.raw_bytes("GET", &url, &[], &[], self.timeouts.request)?;
606        match r.status {
607            200 => Ok(r.body),
608            404 => Ok(Vec::new()),
609            status => Err(Error::Api {
610                method: "GET".into(),
611                path: url,
612                status,
613                message: serde_json::from_slice::<Envelope>(&r.body)
614                    .map(|e| e.error)
615                    .unwrap_or_default(),
616            }),
617        }
618    }
619
620    /// Read a file from an instance; `None` if it does not exist.
621    pub fn read_file(&self, instance: &str, path: &str) -> Result<Option<Vec<u8>>> {
622        let url = format!(
623            "/1.0/instances/{}/files?path={}",
624            encode_segment(instance),
625            encode_query(path)
626        );
627        let r = self.raw_bytes("GET", &url, &[], &[], self.timeouts.request)?;
628        match r.status {
629            200 => Ok(Some(r.body)),
630            404 => Ok(None),
631            status => Err(Error::Api {
632                method: "GET".into(),
633                path: url,
634                status,
635                message: serde_json::from_slice::<Envelope>(&r.body)
636                    .map(|e| e.error)
637                    .unwrap_or_default(),
638            }),
639        }
640    }
641
642    /// Follow `/1.0/events?QUERY` (all projects when the query says so) as
643    /// a websocket. Reads time out every 5 s so a follower can check
644    /// whether to stop; a timeout is not an error.
645    pub fn events_websocket(&self, query: &str) -> Result<tungstenite::WebSocket<UnixStream>> {
646        let stream = self.connect(self.timeouts.request)?;
647        let url = format!("ws://incus/1.0/events?{query}");
648        let (ws, _resp) = tungstenite::client::client(url.as_str(), stream)
649            .map_err(|e| Error::WebSocket(format!("handshake for /1.0/events: {e}")))?;
650        ws.get_ref()
651            .set_read_timeout(Some(Duration::from_secs(5)))?;
652        ws.get_ref()
653            .set_write_timeout(Some(self.timeouts.request))?;
654        Ok(ws)
655    }
656
657    /// Open one of an operation's websockets (exec stdin/stdout/stderr/control).
658    pub(crate) fn websocket(
659        &self,
660        operation: &str,
661        secret: &str,
662    ) -> Result<tungstenite::WebSocket<UnixStream>> {
663        let path = self.with_project(&format!(
664            "{operation}/websocket?secret={}",
665            encode_query(secret)
666        ));
667        let stream = self.connect(self.timeouts.request)?;
668        let url = format!("ws://incus{path}");
669        let (ws, _resp) = tungstenite::client::client(url.as_str(), stream)
670            .map_err(|e| Error::WebSocket(format!("handshake for {operation}: {e}")))?;
671        // Exec output has no default timeout: a quiet process is not a stuck one.
672        ws.get_ref().set_read_timeout(None)?;
673        ws.get_ref()
674            .set_write_timeout(Some(self.timeouts.request))?;
675        Ok(ws)
676    }
677}
678
679impl Default for Client {
680    fn default() -> Self {
681        Client::new()
682    }
683}
684
685/// Percent-encode a query value (RFC 3986 unreserved set passes through).
686pub(crate) fn encode_query(s: &str) -> String {
687    let mut out = String::with_capacity(s.len());
688    for b in s.bytes() {
689        if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~') {
690            out.push(b as char);
691        } else {
692            out.push_str(&format!("%{b:02X}"));
693        }
694    }
695    out
696}
697
698/// Percent-encode one path segment.
699#[doc(hidden)]
700pub fn encode_segment(s: &str) -> String {
701    encode_query(s)
702}
703
704#[derive(Debug, PartialEq)]
705struct RawResponse {
706    status: u16,
707    body: Vec<u8>,
708    etag: Option<String>,
709}
710
711/// If `buf` holds a complete HTTP response, return it with its body decoded.
712fn complete_response(buf: &[u8]) -> Result<Option<RawResponse>> {
713    let mut headers = [httparse::EMPTY_HEADER; 64];
714    let mut resp = httparse::Response::new(&mut headers);
715    let head_len = match resp
716        .parse(buf)
717        .map_err(|e| Error::Protocol(format!("bad HTTP response: {e}")))?
718    {
719        httparse::Status::Complete(n) => n,
720        httparse::Status::Partial => return Ok(None),
721    };
722    let status = resp.code.unwrap_or(0);
723    let mut content_length: Option<usize> = None;
724    let mut chunked = false;
725    let mut etag = None;
726    for h in resp.headers.iter() {
727        if h.name.eq_ignore_ascii_case("etag") {
728            etag = Some(String::from_utf8_lossy(h.value).trim().to_string());
729        }
730        if h.name.eq_ignore_ascii_case("content-length") {
731            content_length = std::str::from_utf8(h.value)
732                .ok()
733                .and_then(|v| v.trim().parse().ok());
734        } else if h.name.eq_ignore_ascii_case("transfer-encoding")
735            && String::from_utf8_lossy(h.value)
736                .to_ascii_lowercase()
737                .contains("chunked")
738        {
739            chunked = true;
740        }
741    }
742    let body = &buf[head_len..];
743    if chunked {
744        return Ok(decode_chunked(body).map(|b| RawResponse {
745            status,
746            body: b,
747            etag,
748        }));
749    }
750    match content_length {
751        Some(n) if body.len() >= n => Ok(Some(RawResponse {
752            status,
753            body: body[..n].to_vec(),
754            etag,
755        })),
756        Some(_) => Ok(None),
757        // No length and not chunked: body runs to EOF; the caller decides.
758        None => Ok(None),
759    }
760}
761
762/// Decode a chunked body; `None` if it is not complete yet.
763fn decode_chunked(mut body: &[u8]) -> Option<Vec<u8>> {
764    let mut out = Vec::new();
765    loop {
766        let line_end = body.windows(2).position(|w| w == b"\r\n")?;
767        let size_str = std::str::from_utf8(&body[..line_end]).ok()?;
768        let size = usize::from_str_radix(size_str.split(';').next()?.trim(), 16).ok()?;
769        body = &body[line_end + 2..];
770        if size == 0 {
771            return Some(out);
772        }
773        if body.len() < size + 2 {
774            return None;
775        }
776        out.extend_from_slice(&body[..size]);
777        body = &body[size + 2..];
778    }
779}
780
781/// When a response had no length and was cut by EOF, treat what we have as the body.
782fn eof_body(buf: &[u8]) -> Option<RawResponse> {
783    let mut headers = [httparse::EMPTY_HEADER; 64];
784    let mut resp = httparse::Response::new(&mut headers);
785    match resp.parse(buf).ok()? {
786        httparse::Status::Complete(n) => Some(RawResponse {
787            status: resp.code?,
788            body: buf[n..].to_vec(),
789            etag: None,
790        }),
791        httparse::Status::Partial => None,
792    }
793}
794
795#[cfg(test)]
796mod tests {
797    use super::*;
798
799    #[test]
800    fn an_incus_without_oci_images_is_named_with_the_way_out() {
801        let old = serde_json::json!({
802            "api_extensions": ["disk_initial_copy"],
803            "environment": {"server_version": "6.0.4"},
804        });
805        let m = oci_unsupported(&old).unwrap();
806        assert!(m.contains("incus 6.0.4"), "{m}");
807        assert!(m.contains("Zabbly"), "{m}");
808        assert!(m.contains("6.3"), "{m}");
809        let new = serde_json::json!({
810            "api_extensions": ["disk_initial_copy", "instance_oci"],
811            "environment": {"server_version": "7.5.1"},
812        });
813        assert_eq!(oci_unsupported(&new), None);
814        // An answer without the list (an odd proxy) is not a verdict.
815        assert_eq!(oci_unsupported(&serde_json::json!({})), None);
816    }
817
818    #[test]
819    fn parses_content_length_response() {
820        let raw = b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nETag: \"abc\"\r\n\r\nhello";
821        let r = complete_response(raw).unwrap().unwrap();
822        assert_eq!((r.status, r.body.as_slice()), (200, &b"hello"[..]));
823        assert_eq!(r.etag.as_deref(), Some("\"abc\""));
824        let partial = b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nhel";
825        assert_eq!(complete_response(partial).unwrap(), None);
826    }
827
828    #[test]
829    fn parses_chunked_response() {
830        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";
831        let r = complete_response(raw).unwrap().unwrap();
832        assert_eq!(r.body, b"hello world".to_vec());
833        let partial = b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nhel";
834        assert_eq!(complete_response(partial).unwrap(), None);
835    }
836
837    #[test]
838    fn encodes_query_values() {
839        assert_eq!(encode_query("a b/c"), "a%20b%2Fc");
840        assert_eq!(encode_query("abc-_.~"), "abc-_.~");
841    }
842}