1use 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
23pub 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#[derive(Debug, Clone)]
43pub struct Timeouts {
44 pub request: Duration,
46 pub create: Duration,
48 pub state: Duration,
50 pub other: Duration,
52 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#[derive(Debug, Clone)]
71pub struct Client {
72 socket: PathBuf,
73 project: Option<String>,
74 #[doc(hidden)]
75 pub timeouts: Timeouts,
76}
77
78#[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#[derive(Debug)]
93#[doc(hidden)]
94pub enum Reply {
95 Sync(Value),
96 Async { operation: String, metadata: Value },
97}
98
99impl Client {
100 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 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 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 pub fn server_info(&self) -> Result<Value> {
161 self.get("/1.0")
162 }
163
164 pub fn oci_unsupported(&self) -> Option<String> {
168 oci_unsupported(&self.server_info().ok()?)
169 }
170
171 fn with_project(&self, path: &str) -> String {
172 match &self.project {
173 None => path.to_string(),
174 Some(p) => {
175 let sep = if path.contains('?') { '&' } else { '?' };
176 format!("{path}{sep}project={}", encode_query(p))
177 }
178 }
179 }
180
181 fn connect(&self, timeout: Duration) -> Result<UnixStream> {
182 let stream = UnixStream::connect(&self.socket).map_err(|source| Error::Connect {
183 socket: self.socket.display().to_string(),
184 source,
185 })?;
186 stream.set_read_timeout(Some(timeout))?;
187 stream.set_write_timeout(Some(timeout))?;
188 Ok(stream)
189 }
190
191 fn raw(
193 &self,
194 method: &str,
195 path: &str,
196 body: Option<&Value>,
197 if_match: Option<&str>,
198 timeout: Duration,
199 ) -> Result<RawResponse> {
200 let payload = match body {
201 Some(v) => serde_json::to_vec(v)?,
202 None => Vec::new(),
203 };
204 let mut headers: Vec<(&str, String)> = Vec::new();
205 if body.is_some() {
206 headers.push(("Content-Type", "application/json".into()));
207 }
208 if let Some(etag) = if_match {
209 headers.push(("If-Match", etag.to_string()));
210 }
211 self.raw_bytes(method, path, &payload, &headers, timeout)
212 }
213
214 fn raw_bytes(
216 &self,
217 method: &str,
218 path: &str,
219 payload: &[u8],
220 extra_headers: &[(&str, String)],
221 timeout: Duration,
222 ) -> Result<RawResponse> {
223 let path = self.with_project(path);
224 let started = Instant::now();
225 let to_err = |e: std::io::Error| -> Error {
226 if matches!(
227 e.kind(),
228 std::io::ErrorKind::WouldBlock | std::io::ErrorKind::TimedOut
229 ) {
230 Error::RequestTimeout {
231 method: method.to_string(),
232 path: path.clone(),
233 timeout,
234 }
235 } else {
236 Error::Io(e)
237 }
238 };
239 let mut stream = self.connect(timeout)?;
240 let mut head = format!(
241 "{method} {path} HTTP/1.1\r\nHost: incus\r\nUser-Agent: isb/{}\r\nConnection: close\r\n",
242 env!("CARGO_PKG_VERSION")
243 );
244 for (k, v) in extra_headers {
245 head.push_str(&format!("{k}: {v}\r\n"));
246 }
247 head.push_str(&format!("Content-Length: {}\r\n\r\n", payload.len()));
248 stream.write_all(head.as_bytes()).map_err(to_err)?;
249 stream.write_all(payload).map_err(to_err)?;
250 stream.flush().map_err(to_err)?;
251
252 let mut buf = Vec::with_capacity(8192);
253 let mut chunk = [0u8; 16384];
254 loop {
255 let remaining = timeout.saturating_sub(started.elapsed());
256 if remaining.is_zero() {
257 return Err(to_err(std::io::ErrorKind::TimedOut.into()));
258 }
259 let _ = stream.set_read_timeout(Some(remaining));
262 let n = stream.read(&mut chunk).map_err(to_err)?;
263 if n == 0 {
264 break;
265 }
266 buf.extend_from_slice(&chunk[..n]);
267 if let Some(done) = complete_response(&buf)? {
270 return Ok(done);
271 }
272 }
273 complete_response(&buf)?
274 .or_else(|| eof_body(&buf))
275 .ok_or_else(|| Error::Protocol(format!("truncated response to {method} {path}")))
276 }
277
278 pub(crate) fn request(
280 &self,
281 method: &str,
282 path: &str,
283 body: Option<&Value>,
284 timeout: Duration,
285 ) -> Result<Reply> {
286 self.request_etag(method, path, body, None, timeout)
287 .map(|(r, _)| r)
288 }
289
290 pub(crate) fn request_etag(
293 &self,
294 method: &str,
295 path: &str,
296 body: Option<&Value>,
297 if_match: Option<&str>,
298 timeout: Duration,
299 ) -> Result<(Reply, Option<String>)> {
300 let RawResponse {
301 status,
302 body: bytes,
303 etag,
304 } = self.raw(method, path, body, if_match, timeout)?;
305 let env: Envelope = serde_json::from_slice(&bytes).map_err(|e| {
306 Error::Protocol(format!(
307 "{method} {path}: HTTP {status}, undecodable body ({e}): {}",
308 String::from_utf8_lossy(&bytes[..bytes.len().min(200)])
309 ))
310 })?;
311 if env.kind == "error" || status >= 400 {
312 if let Some(e) = crate::org::limits::translate(self, &env.error) {
313 return Err(e);
314 }
315 return Err(Error::Api {
316 method: method.to_string(),
317 path: path.to_string(),
318 status,
319 message: if env.error.is_empty() {
320 format!("HTTP {status}")
321 } else {
322 env.error
323 },
324 });
325 }
326 if env.kind == "async" {
327 return Ok((
328 Reply::Async {
329 operation: env.operation,
330 metadata: env.metadata,
331 },
332 etag,
333 ));
334 }
335 Ok((Reply::Sync(env.metadata), etag))
336 }
337
338 #[doc(hidden)]
340 pub fn get_etag(&self, path: &str) -> Result<(Value, Option<String>)> {
341 match self.request_etag("GET", path, None, None, self.timeouts.request)? {
342 (Reply::Sync(v), e) => Ok((v, e)),
343 (Reply::Async { metadata, .. }, e) => Ok((metadata, e)),
344 }
345 }
346
347 #[doc(hidden)]
349 pub fn mutate_if_match(
350 &self,
351 method: &str,
352 path: &str,
353 body: &Value,
354 etag: Option<&str>,
355 step: &str,
356 deadline: Duration,
357 ) -> Result<Value> {
358 match self.request_etag(method, path, Some(body), etag, self.timeouts.request)? {
359 (Reply::Sync(v), _) => Ok(v),
360 (Reply::Async { operation, .. }, _) => self.wait_operation(&operation, step, deadline),
361 }
362 }
363
364 #[doc(hidden)]
365 pub fn get(&self, path: &str) -> Result<Value> {
366 match self.request("GET", path, None, self.timeouts.request)? {
367 Reply::Sync(v) => Ok(v),
368 Reply::Async { metadata, .. } => Ok(metadata),
369 }
370 }
371
372 #[doc(hidden)]
374 pub fn get_opt(&self, path: &str) -> Result<Option<Value>> {
375 match self.get(path) {
376 Ok(v) => Ok(Some(v)),
377 Err(e) if e.is_not_found() => Ok(None),
378 Err(e) => Err(e),
379 }
380 }
381
382 #[doc(hidden)]
386 pub fn mutate(
387 &self,
388 method: &str,
389 path: &str,
390 body: Option<&Value>,
391 step: &str,
392 deadline: Duration,
393 ) -> Result<Value> {
394 match self.request(method, path, body, self.timeouts.request)? {
395 Reply::Sync(v) => Ok(v),
396 Reply::Async { operation, .. } => self.wait_operation(&operation, step, deadline),
397 }
398 }
399
400 pub fn wait_operation(&self, operation: &str, step: &str, deadline: Duration) -> Result<Value> {
402 let started = Instant::now();
403 loop {
404 let remaining = deadline.saturating_sub(started.elapsed());
405 let secs = remaining.as_secs().min(30);
409 if let Some(v) = self.poll_operation(operation, step, secs)? {
410 return Ok(v);
411 }
412 if started.elapsed() >= deadline {
413 let status = self
414 .get(operation)
415 .ok()
416 .and_then(|op| op.get("status").and_then(Value::as_str).map(String::from))
417 .unwrap_or_else(|| "running".into())
418 .to_lowercase();
419 let cancelled = self.cancel_operation(operation).is_ok();
420 return Err(Error::OperationTimeout {
421 step: step.to_string(),
422 operation: operation.to_string(),
423 status,
424 waited: started.elapsed(),
425 cancelled,
426 });
427 }
428 if secs == 0 {
429 std::thread::sleep(remaining.min(Duration::from_millis(50)));
430 }
431 }
432 }
433
434 pub(crate) fn poll_operation(
437 &self,
438 operation: &str,
439 step: &str,
440 secs: u64,
441 ) -> Result<Option<Value>> {
442 let op = self.get_with_timeout(
443 &format!("{operation}/wait?timeout={secs}"),
444 Duration::from_secs(secs) + self.timeouts.request,
445 )?;
446 let code = op.get("status_code").and_then(Value::as_i64).unwrap_or(0);
447 match code {
448 200 => Ok(Some(op.get("metadata").cloned().unwrap_or(Value::Null))),
449 400 | 401 => {
450 let err = op
451 .get("err")
452 .and_then(Value::as_str)
453 .filter(|s| !s.is_empty())
454 .unwrap_or(if code == 401 {
455 "operation cancelled"
456 } else {
457 "operation failed"
458 });
459 if let Some(e) = crate::org::limits::translate(self, err) {
460 return Err(e);
461 }
462 Err(Error::OperationFailed {
463 step: step.to_string(),
464 message: err.to_string(),
465 })
466 }
467 _ => Ok(None),
468 }
469 }
470
471 pub fn start_operation(
474 &self,
475 method: &str,
476 path: &str,
477 body: Option<&Value>,
478 ) -> Result<String> {
479 match self.request(method, path, body, self.timeouts.request)? {
480 Reply::Async { operation, .. } => Ok(operation),
481 Reply::Sync(_) => Err(Error::Protocol(format!(
482 "{method} {path} did not start an operation"
483 ))),
484 }
485 }
486
487 fn get_with_timeout(&self, path: &str, timeout: Duration) -> Result<Value> {
488 match self.request("GET", path, None, timeout)? {
489 Reply::Sync(v) => Ok(v),
490 Reply::Async { metadata, .. } => Ok(metadata),
491 }
492 }
493
494 pub fn cancel_operation(&self, operation: &str) -> Result<()> {
496 self.request("DELETE", operation, None, self.timeouts.request)?;
497 Ok(())
498 }
499
500 pub fn push_file(
503 &self,
504 instance: &str,
505 path: &str,
506 data: &[u8],
507 uid: u32,
508 gid: u32,
509 mode: u32,
510 ) -> Result<()> {
511 self.file_request(instance, path, data, uid, gid, mode, "file")
512 }
513
514 pub fn make_dir(
516 &self,
517 instance: &str,
518 path: &str,
519 uid: u32,
520 gid: u32,
521 mode: u32,
522 ) -> Result<()> {
523 match self.file_request(instance, path, &[], uid, gid, mode, "directory") {
524 Err(Error::Api {
525 status,
526 ref message,
527 ..
528 }) if status == 409 || message.to_ascii_lowercase().contains("exist") => Ok(()),
529 r => r,
530 }
531 }
532
533 #[expect(clippy::too_many_arguments)]
534 fn file_request(
535 &self,
536 instance: &str,
537 path: &str,
538 data: &[u8],
539 uid: u32,
540 gid: u32,
541 mode: u32,
542 kind: &str,
543 ) -> Result<()> {
544 let url = format!(
545 "/1.0/instances/{}/files?path={}",
546 encode_segment(instance),
547 encode_query(path)
548 );
549 let headers = [
550 ("Content-Type", "application/octet-stream".to_string()),
551 ("X-Incus-uid", uid.to_string()),
552 ("X-Incus-gid", gid.to_string()),
553 ("X-Incus-mode", format!("{mode:04o}")),
554 ("X-Incus-type", kind.to_string()),
555 ("X-Incus-write", "overwrite".to_string()),
556 ];
557 let r = self.raw_bytes("POST", &url, data, &headers, self.timeouts.request)?;
558 if r.status >= 400 {
559 let message = serde_json::from_slice::<Envelope>(&r.body)
560 .map(|e| e.error)
561 .ok()
562 .filter(|e| !e.is_empty())
563 .unwrap_or_else(|| format!("HTTP {}", r.status));
564 return Err(Error::Api {
565 method: "POST".into(),
566 path: url,
567 status: r.status,
568 message,
569 });
570 }
571 Ok(())
572 }
573
574 pub(crate) fn get_raw(&self, path: &str) -> Result<Vec<u8>> {
576 let r = self.raw_bytes("GET", path, &[], &[], self.timeouts.request)?;
577 match r.status {
578 200 => Ok(r.body),
579 status => Err(Error::Api {
580 method: "GET".into(),
581 path: path.into(),
582 status,
583 message: String::from_utf8_lossy(&r.body[..r.body.len().min(200)]).into_owned(),
584 }),
585 }
586 }
587
588 pub fn console_log(&self, instance: &str) -> Result<Vec<u8>> {
590 let url = format!("/1.0/instances/{}/console", encode_segment(instance));
591 let r = self.raw_bytes("GET", &url, &[], &[], self.timeouts.request)?;
592 match r.status {
593 200 => Ok(r.body),
594 404 => Ok(Vec::new()),
595 status => Err(Error::Api {
596 method: "GET".into(),
597 path: url,
598 status,
599 message: serde_json::from_slice::<Envelope>(&r.body)
600 .map(|e| e.error)
601 .unwrap_or_default(),
602 }),
603 }
604 }
605
606 pub fn read_file(&self, instance: &str, path: &str) -> Result<Option<Vec<u8>>> {
608 let url = format!(
609 "/1.0/instances/{}/files?path={}",
610 encode_segment(instance),
611 encode_query(path)
612 );
613 let r = self.raw_bytes("GET", &url, &[], &[], self.timeouts.request)?;
614 match r.status {
615 200 => Ok(Some(r.body)),
616 404 => Ok(None),
617 status => Err(Error::Api {
618 method: "GET".into(),
619 path: url,
620 status,
621 message: serde_json::from_slice::<Envelope>(&r.body)
622 .map(|e| e.error)
623 .unwrap_or_default(),
624 }),
625 }
626 }
627
628 pub fn events_websocket(&self, query: &str) -> Result<tungstenite::WebSocket<UnixStream>> {
632 let stream = self.connect(self.timeouts.request)?;
633 let url = format!("ws://incus/1.0/events?{query}");
634 let (ws, _resp) = tungstenite::client::client(url.as_str(), stream)
635 .map_err(|e| Error::WebSocket(format!("handshake for /1.0/events: {e}")))?;
636 ws.get_ref()
637 .set_read_timeout(Some(Duration::from_secs(5)))?;
638 ws.get_ref()
639 .set_write_timeout(Some(self.timeouts.request))?;
640 Ok(ws)
641 }
642
643 pub(crate) fn websocket(
645 &self,
646 operation: &str,
647 secret: &str,
648 ) -> Result<tungstenite::WebSocket<UnixStream>> {
649 let path = self.with_project(&format!(
650 "{operation}/websocket?secret={}",
651 encode_query(secret)
652 ));
653 let stream = self.connect(self.timeouts.request)?;
654 let url = format!("ws://incus{path}");
655 let (ws, _resp) = tungstenite::client::client(url.as_str(), stream)
656 .map_err(|e| Error::WebSocket(format!("handshake for {operation}: {e}")))?;
657 ws.get_ref().set_read_timeout(None)?;
659 ws.get_ref()
660 .set_write_timeout(Some(self.timeouts.request))?;
661 Ok(ws)
662 }
663}
664
665impl Default for Client {
666 fn default() -> Self {
667 Client::new()
668 }
669}
670
671pub(crate) fn encode_query(s: &str) -> String {
673 let mut out = String::with_capacity(s.len());
674 for b in s.bytes() {
675 if b.is_ascii_alphanumeric() || matches!(b, b'-' | b'_' | b'.' | b'~') {
676 out.push(b as char);
677 } else {
678 out.push_str(&format!("%{b:02X}"));
679 }
680 }
681 out
682}
683
684#[doc(hidden)]
686pub fn encode_segment(s: &str) -> String {
687 encode_query(s)
688}
689
690#[derive(Debug, PartialEq)]
691struct RawResponse {
692 status: u16,
693 body: Vec<u8>,
694 etag: Option<String>,
695}
696
697fn complete_response(buf: &[u8]) -> Result<Option<RawResponse>> {
699 let mut headers = [httparse::EMPTY_HEADER; 64];
700 let mut resp = httparse::Response::new(&mut headers);
701 let head_len = match resp
702 .parse(buf)
703 .map_err(|e| Error::Protocol(format!("bad HTTP response: {e}")))?
704 {
705 httparse::Status::Complete(n) => n,
706 httparse::Status::Partial => return Ok(None),
707 };
708 let status = resp.code.unwrap_or(0);
709 let mut content_length: Option<usize> = None;
710 let mut chunked = false;
711 let mut etag = None;
712 for h in resp.headers.iter() {
713 if h.name.eq_ignore_ascii_case("etag") {
714 etag = Some(String::from_utf8_lossy(h.value).trim().to_string());
715 }
716 if h.name.eq_ignore_ascii_case("content-length") {
717 content_length = std::str::from_utf8(h.value)
718 .ok()
719 .and_then(|v| v.trim().parse().ok());
720 } else if h.name.eq_ignore_ascii_case("transfer-encoding")
721 && String::from_utf8_lossy(h.value)
722 .to_ascii_lowercase()
723 .contains("chunked")
724 {
725 chunked = true;
726 }
727 }
728 let body = &buf[head_len..];
729 if chunked {
730 return Ok(decode_chunked(body).map(|b| RawResponse {
731 status,
732 body: b,
733 etag,
734 }));
735 }
736 match content_length {
737 Some(n) if body.len() >= n => Ok(Some(RawResponse {
738 status,
739 body: body[..n].to_vec(),
740 etag,
741 })),
742 Some(_) => Ok(None),
743 None => Ok(None),
745 }
746}
747
748fn decode_chunked(mut body: &[u8]) -> Option<Vec<u8>> {
750 let mut out = Vec::new();
751 loop {
752 let line_end = body.windows(2).position(|w| w == b"\r\n")?;
753 let size_str = std::str::from_utf8(&body[..line_end]).ok()?;
754 let size = usize::from_str_radix(size_str.split(';').next()?.trim(), 16).ok()?;
755 body = &body[line_end + 2..];
756 if size == 0 {
757 return Some(out);
758 }
759 if body.len() < size + 2 {
760 return None;
761 }
762 out.extend_from_slice(&body[..size]);
763 body = &body[size + 2..];
764 }
765}
766
767fn eof_body(buf: &[u8]) -> Option<RawResponse> {
769 let mut headers = [httparse::EMPTY_HEADER; 64];
770 let mut resp = httparse::Response::new(&mut headers);
771 match resp.parse(buf).ok()? {
772 httparse::Status::Complete(n) => Some(RawResponse {
773 status: resp.code?,
774 body: buf[n..].to_vec(),
775 etag: None,
776 }),
777 httparse::Status::Partial => None,
778 }
779}
780
781#[cfg(test)]
782mod tests {
783 use super::*;
784
785 #[test]
786 fn an_incus_without_oci_images_is_named_with_the_way_out() {
787 let old = serde_json::json!({
788 "api_extensions": ["disk_initial_copy"],
789 "environment": {"server_version": "6.0.4"},
790 });
791 let m = oci_unsupported(&old).unwrap();
792 assert!(m.contains("incus 6.0.4"), "{m}");
793 assert!(m.contains("Zabbly"), "{m}");
794 assert!(m.contains("6.3"), "{m}");
795 let new = serde_json::json!({
796 "api_extensions": ["disk_initial_copy", "instance_oci"],
797 "environment": {"server_version": "7.5.1"},
798 });
799 assert_eq!(oci_unsupported(&new), None);
800 assert_eq!(oci_unsupported(&serde_json::json!({})), None);
802 }
803
804 #[test]
805 fn parses_content_length_response() {
806 let raw = b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nETag: \"abc\"\r\n\r\nhello";
807 let r = complete_response(raw).unwrap().unwrap();
808 assert_eq!((r.status, r.body.as_slice()), (200, &b"hello"[..]));
809 assert_eq!(r.etag.as_deref(), Some("\"abc\""));
810 let partial = b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\n\r\nhel";
811 assert_eq!(complete_response(partial).unwrap(), None);
812 }
813
814 #[test]
815 fn parses_chunked_response() {
816 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";
817 let r = complete_response(raw).unwrap().unwrap();
818 assert_eq!(r.body, b"hello world".to_vec());
819 let partial = b"HTTP/1.1 200 OK\r\nTransfer-Encoding: chunked\r\n\r\n5\r\nhel";
820 assert_eq!(complete_response(partial).unwrap(), None);
821 }
822
823 #[test]
824 fn encodes_query_values() {
825 assert_eq!(encode_query("a b/c"), "a%20b%2Fc");
826 assert_eq!(encode_query("abc-_.~"), "abc-_.~");
827 }
828}