1use std::collections::BTreeMap;
4use std::path::Path;
5use std::time::{Duration, Instant};
6
7use serde::Serialize;
8use serde_json::{Value, json};
9
10use crate::client::{Client, encode_segment};
11use crate::error::{Error, Result};
12use crate::exec::{self, ExecOptions, ExecOutput, ExecStream, Stdin};
13use crate::idmap::SubIds;
14use crate::lock::NameLock;
15
16pub mod images;
17use crate::plan::{
18 self, Action, Actual, Desired, DesiredDevice, DiffOptions, HostFacts, Props, SandboxPlan,
19 VolumeDefs, split_addr,
20};
21use crate::spec::{ExecDefaults, PortSpec, ReadyCheck, SandboxSpec};
22
23pub const CREATE_TOKEN_KEY: &str = "user.isb.create-token";
26
27pub type Reporter<'a> = &'a mut dyn FnMut(&str);
29
30#[derive(Debug, Clone, Serialize)]
32pub struct SandboxInfo {
33 pub name: String,
34 pub status: String,
35 #[serde(rename = "type")]
36 pub instance_type: String,
37 pub labels: BTreeMap<String, String>,
39 pub config: BTreeMap<String, String>,
40 pub devices: BTreeMap<String, Props>,
42 pub profiles: Vec<String>,
43 pub created_at: String,
44 pub description: String,
45}
46
47impl SandboxInfo {
48 pub fn from_api(v: &Value) -> SandboxInfo {
49 let a = Actual::from_api(v);
50 let labels = a
51 .config
52 .iter()
53 .filter_map(|(k, v)| k.strip_prefix("user.").map(|k| (k.to_string(), v.clone())))
54 .filter(|(k, _): &(String, String)| !k.starts_with("isb."))
56 .collect();
57 SandboxInfo {
58 name: v.get("name").and_then(Value::as_str).unwrap_or("").into(),
59 status: a.status.clone(),
60 instance_type: a.instance_type.clone(),
61 labels,
62 config: a.config,
63 devices: a.devices,
64 profiles: a.profiles,
65 created_at: v
66 .get("created_at")
67 .and_then(Value::as_str)
68 .unwrap_or("")
69 .into(),
70 description: v
71 .get("description")
72 .and_then(Value::as_str)
73 .unwrap_or("")
74 .into(),
75 }
76 }
77
78 pub fn matches(&self, filters: &[LabelFilter]) -> bool {
80 filters
81 .iter()
82 .all(|f| match (&f.value, self.labels.get(&f.key)) {
83 (_, None) => false,
84 (None, Some(_)) => true,
85 (Some(want), Some(have)) => want == have,
86 })
87 }
88}
89
90#[derive(Debug, Clone, PartialEq, Eq)]
92pub struct LabelFilter {
93 pub key: String,
94 pub value: Option<String>,
95}
96
97impl LabelFilter {
98 pub fn parse(s: &str) -> LabelFilter {
99 match s.split_once('=') {
100 Some((k, v)) => LabelFilter {
101 key: k.into(),
102 value: Some(v.into()),
103 },
104 None => LabelFilter {
105 key: s.into(),
106 value: None,
107 },
108 }
109 }
110}
111
112#[derive(Debug, Clone, Default, Serialize)]
114pub struct ApplyReport {
115 pub name: String,
116 pub created: bool,
117 pub applied: Vec<Action>,
118 pub ports: BTreeMap<String, String>,
120 pub restart_needed: Vec<String>,
122}
123
124#[derive(Debug, Clone, Copy)]
126pub struct EnsureOptions {
127 pub diff: DiffOptions,
128 pub wait_ready: bool,
130 pub lock_wait: Duration,
132}
133
134impl Default for EnsureOptions {
135 fn default() -> Self {
136 EnsureOptions {
137 diff: DiffOptions::default(),
138 wait_ready: true,
139 lock_wait: Duration::from_secs(900),
140 }
141 }
142}
143
144fn inst_path(name: &str) -> String {
145 format!("/1.0/instances/{}", encode_segment(name))
146}
147
148pub fn host_facts(client: &Client) -> Result<HostFacts> {
150 let pools = client
151 .get("/1.0/storage-pools")?
152 .as_array()
153 .map(|a| {
154 a.iter()
155 .filter_map(|u| u.as_str())
156 .filter_map(|u| u.rsplit('/').next())
157 .map(|s| s.split('?').next().unwrap_or(s).to_string())
158 .collect()
159 })
160 .unwrap_or_default();
161 let initial_copy = client.has_extension("disk_initial_copy")?;
162 Ok(HostFacts {
163 subids: SubIds::read_host(),
164 pools,
165 path_map: HostFacts::detect_path_map(),
166 initial_copy,
167 incus_version: client.server_version()?,
168 invoking_ids: crate::idmap::invoking_ids(),
169 shared_root: shared_root(),
170 org: crate::org::OrgId::from_incus_project(client.project_name()),
171 registry: crate::registry::info(client)?.map(|i| i.addr),
172 project: client.project_name().to_string(),
173 })
174}
175
176fn shared_root() -> Option<String> {
178 if !cfg!(target_os = "macos") {
179 return None;
180 }
181 let home = std::env::var_os("HOME").filter(|h| !h.is_empty())?;
182 let p = std::path::PathBuf::from(home);
183 Some(p.canonicalize().unwrap_or(p).to_string_lossy().into_owned())
184}
185
186pub fn resolve(
188 client: &Client,
189 spec: &SandboxSpec,
190 defs: &VolumeDefs,
191 base: &Path,
192) -> Result<Desired> {
193 let host = host_facts(client)?;
194 plan::resolve(spec, defs, &host, base)
195}
196
197fn get_actual(client: &Client, name: &str) -> Result<Option<Actual>> {
198 Ok(client
199 .get_opt(&inst_path(name))?
200 .map(|v| Actual::from_api(&v)))
201}
202
203fn volume_exists(client: &Client, pool: &str, name: &str) -> Result<bool> {
204 Ok(client
205 .get_opt(&format!(
206 "/1.0/storage-pools/{}/volumes/custom/{}",
207 encode_segment(pool),
208 encode_segment(name)
209 ))?
210 .is_some())
211}
212
213fn local_image(client: &Client, alias: &str) -> Result<Option<String>> {
215 if let Some(a) = client.get_opt(&format!("/1.0/images/aliases/{}", encode_segment(alias)))? {
216 return Ok(a.get("target").and_then(Value::as_str).map(String::from));
217 }
218 if alias.len() >= 12 && alias.chars().all(|c| c.is_ascii_hexdigit()) {
219 if let Some(i) = client.get_opt(&format!("/1.0/images/{}", encode_segment(alias)))? {
220 return Ok(i
221 .get("fingerprint")
222 .and_then(Value::as_str)
223 .map(String::from));
224 }
225 }
226 Ok(None)
227}
228
229pub fn plan_desired(client: &Client, desired: &Desired, opts: DiffOptions) -> Result<SandboxPlan> {
231 let actual = get_actual(client, &desired.name)?;
232 if actual.is_none()
233 && desired.image.server.is_none()
234 && local_image(client, &desired.image.alias)?.is_none()
235 {
236 return Err(images::missing_on(client, &desired.image.alias));
237 }
238 let mut missing = Vec::new();
239 for v in &desired.volumes {
240 if !volume_exists(client, &v.pool, &v.name)? {
241 missing.push((v.pool.clone(), v.name.clone()));
242 }
243 }
244 let mut plan = plan::diff(desired, actual.as_ref(), &missing, opts)?;
245 if let (Some(a), None) = (&actual, &desired.image.server) {
248 if let (Some(built), Some(now)) = (
249 a.config.get("volatile.base_image"),
250 local_image(client, &desired.image.alias)?,
251 ) {
252 if *built != now {
253 plan.actions.push(Action::Note {
254 message: format!(
255 "image {} is now {} but this instance was built from {}; recreate to pick it up",
256 desired.image.alias,
257 &now[..12.min(now.len())],
258 &built[..12.min(built.len())]
259 ),
260 });
261 }
262 }
263 }
264 Ok(plan)
265}
266
267fn random_token() -> String {
268 let mut b = [0u8; 12];
269 if let Ok(mut f) = std::fs::File::open("/dev/urandom") {
270 use std::io::Read;
271 let _ = f.read_exact(&mut b);
272 }
273 let t = std::time::SystemTime::now()
274 .duration_since(std::time::UNIX_EPOCH)
275 .map(|d| d.as_nanos())
276 .unwrap_or(0);
277 format!(
278 "{}{:x}",
279 b.iter().map(|x| format!("{x:02x}")).collect::<String>(),
280 t & 0xffff
281 )
282}
283
284fn retry_once<T>(report: &mut dyn FnMut(&str), mut f: impl FnMut() -> Result<T>) -> Result<T> {
286 match f() {
287 Err(e) if e.is_timeout() => {
288 report(&format!("{e}; retrying once"));
289 f()
290 }
291 r => r,
292 }
293}
294
295pub fn apply(
297 client: &Client,
298 desired: &Desired,
299 plan: &SandboxPlan,
300 report: Reporter<'_>,
301) -> Result<ApplyReport> {
302 let name = &desired.name;
303 let mut out = ApplyReport {
304 name: name.clone(),
305 ..Default::default()
306 };
307 let mut pending: Vec<&Action> = Vec::new();
308 if let Some(p) = &desired.egress {
309 crate::egress::before_apply(client, p, report)?;
310 }
311 for action in &plan.actions {
312 match action {
313 Action::SetConfig { .. }
314 | Action::AddDevice { .. }
315 | Action::ReplaceDevice { .. }
316 | Action::RemoveDevice { .. } => {
317 pending.push(action);
318 continue;
319 }
320 _ => {}
321 }
322 flush_updates(client, desired, &mut pending, &mut out, report)?;
325 match action {
326 Action::Note { message } => report(&format!("{name}: note: {message}")),
327 Action::CreateVolume {
328 pool,
329 volume,
330 config,
331 } => {
332 report(&format!("{name}: creating volume {volume} on {pool}"));
333 retry_once(report, || {
334 crate::volume::ensure(client, pool, volume, config).map(|_| ())
335 })?;
336 }
337 Action::CreateInstance { .. } => {
338 report(&format!("{name}: creating from {}", desired.image.spec));
339 create_instance(client, desired, report)?;
340 out.created = true;
341 }
342 Action::StartInstance => {
343 report(&format!("{name}: starting"));
344 retry_once(report, || start_instance(client, name))?;
345 }
346 Action::AddPort {
347 device,
348 props,
349 search,
350 } => {
351 let listen = add_port_searching(client, name, device, props, *search)?;
352 report(&format!("{name}: port {device} listening on {listen}"));
353 out.ports.insert(device.clone(), listen);
354 }
355 Action::FixOwner { path, owner } => {
356 if desired.instance_type == crate::spec::InstanceType::VirtualMachine {
357 wait_ready(
359 client,
360 name,
361 &[ReadyCheck::Agent],
362 desired.ready_timeout,
363 &desired.exec,
364 )?;
365 }
366 report(&format!("{name}: chown {owner} {path}"));
367 fix_owner(client, name, path, owner)?;
368 }
369 _ => unreachable!("batched above"),
370 }
371 out.applied.push(action.clone());
372 }
373 flush_updates(client, desired, &mut pending, &mut out, report)?;
374 Ok(out)
375}
376
377fn flush_updates(
378 client: &Client,
379 desired: &Desired,
380 pending: &mut Vec<&Action>,
381 out: &mut ApplyReport,
382 report: &mut dyn FnMut(&str),
383) -> Result<()> {
384 if pending.is_empty() {
385 return Ok(());
386 }
387 let name = desired.name.as_str();
388 for a in pending.iter() {
389 report(&format!("{name}: {a}"));
390 }
391 let actions: Vec<Action> = pending.iter().map(|a| (*a).clone()).collect();
392 retry_once(report, || {
393 update_instance(
394 client,
395 name,
396 &format!("update {name}"),
397 &mut |config, devices| {
398 for a in &actions {
399 match a {
400 Action::SetConfig {
401 key, to, secret, ..
402 } => {
403 let v = match (secret, desired.config.get(key)) {
405 (true, Some(v)) => v,
406 _ => to,
407 };
408 config.insert(key.clone(), json!(v));
409 }
410 Action::AddDevice { device, props } => {
411 devices.insert(device.clone(), json!(props));
412 }
413 Action::ReplaceDevice {
414 device,
415 replaces,
416 to,
417 ..
418 } => {
419 devices.remove(replaces);
420 devices.insert(device.clone(), json!(to));
421 }
422 Action::RemoveDevice { device, .. } => {
423 devices.remove(device);
424 }
425 _ => {}
426 }
427 }
428 Ok(())
429 },
430 )
431 })?;
432 for a in pending.drain(..) {
433 if let Action::SetConfig {
434 key, restart: true, ..
435 } = a
436 {
437 out.restart_needed.push(key.clone());
438 }
439 out.applied.push(a.clone());
440 }
441 Ok(())
442}
443
444type Obj = serde_json::Map<String, Value>;
445
446#[doc(hidden)]
449pub fn update_instance(
450 client: &Client,
451 name: &str,
452 step: &str,
453 modify: &mut dyn FnMut(&mut Obj, &mut Obj) -> Result<()>,
454) -> Result<()> {
455 let path = inst_path(name);
456 for attempt in 0..5 {
457 let (inst, etag) = client.get_etag(&path)?;
458 let mut config = inst
459 .get("config")
460 .and_then(Value::as_object)
461 .cloned()
462 .unwrap_or_default();
463 let mut devices = inst
464 .get("devices")
465 .and_then(Value::as_object)
466 .cloned()
467 .unwrap_or_default();
468 modify(&mut config, &mut devices)?;
469 let body = json!({
470 "architecture": inst.get("architecture"),
471 "config": config,
472 "devices": devices,
473 "ephemeral": inst.get("ephemeral"),
474 "profiles": inst.get("profiles"),
475 "stateful": inst.get("stateful"),
476 "description": inst.get("description"),
477 });
478 match client.mutate_if_match(
479 "PUT",
480 &path,
481 &body,
482 etag.as_deref(),
483 step,
484 client.timeouts.other,
485 ) {
486 Err(Error::Api { status: 412, .. }) if attempt < 4 => continue,
487 r => return r.map(|_| ()),
488 }
489 }
490 unreachable!()
491}
492
493fn create_instance(client: &Client, desired: &Desired, report: &mut dyn FnMut(&str)) -> Result<()> {
494 let fingerprint = match &desired.image.server {
495 None => Some(
496 local_image(client, &desired.image.alias)?
497 .ok_or_else(|| images::missing_on(client, &desired.image.alias))?,
498 ),
499 Some(_) => None,
500 };
501 let step = format!("create instance {}", desired.name);
502 let mut last_err = None;
503 for attempt in 0..2 {
504 let token = random_token();
505 let mut config = desired.config.clone();
506 config.insert(CREATE_TOKEN_KEY.into(), token.clone());
507 let devices: BTreeMap<&String, &Props> = desired
508 .devices
509 .iter()
510 .filter(|(_, d)| d.search.is_none())
511 .map(|(k, d)| (k, &d.props))
512 .collect();
513 let body = json!({
514 "name": desired.name,
515 "type": desired.instance_type.as_api(),
516 "source": desired.image.to_api(fingerprint.as_deref()),
517 "config": config,
518 "devices": devices,
519 "profiles": desired.profiles,
520 });
521 match client.mutate(
522 "POST",
523 "/1.0/instances",
524 Some(&body),
525 &step,
526 client.timeouts.create,
527 ) {
528 Ok(_) => return Ok(()),
529 Err(e) if e.is_timeout() => {
530 report(&format!("{e}"));
531 if let Error::OperationTimeout { operation, .. } = &e {
532 if client
535 .wait_operation(operation, &step, client.timeouts.settle)
536 .is_ok()
537 {
538 report(&format!("{step}: finished late; keeping it"));
539 return Ok(());
540 }
541 }
542 if let Err(ce) = cleanup_half_created(client, &desired.name, &token, report) {
543 report(&format!("{}: cleanup failed: {ce}", desired.name));
544 if matches!(ce, Error::AlreadyExists(_)) {
545 return Err(e);
546 }
547 }
548 if attempt == 0 {
549 report(&format!("{step}: retrying once"));
550 }
551 last_err = Some(e);
552 }
553 Err(e) => {
554 if !e.is_conflict() {
556 let _ = cleanup_half_created(client, &desired.name, &token, report);
557 }
558 return Err(e);
559 }
560 }
561 }
562 Err(last_err.expect("loop ran"))
563}
564
565pub fn cleanup_half_created(
570 client: &Client,
571 name: &str,
572 token: &str,
573 report: &mut dyn FnMut(&str),
574) -> Result<()> {
575 let Some(inst) = client.get_opt(&inst_path(name))? else {
576 return Ok(());
577 };
578 let a = Actual::from_api(&inst);
579 if a.config.get(CREATE_TOKEN_KEY).map(String::as_str) != Some(token) {
580 report(&format!(
581 "{name}: exists but was not created by this call; leaving it alone"
582 ));
583 return Err(Error::AlreadyExists(name.to_string()));
584 }
585 report(&format!("{name}: removing half-created instance"));
586 force_delete(client, name)
587}
588
589fn start_instance(client: &Client, name: &str) -> Result<()> {
590 let r = client.mutate(
591 "PUT",
592 &format!("{}/state", inst_path(name)),
593 Some(&json!({"action": "start", "timeout": 30})),
594 &format!("start {name}"),
595 client.timeouts.state,
596 );
597 match r {
598 Ok(_) => Ok(()),
599 Err(e) => match get_actual(client, name) {
601 Ok(Some(a)) if a.running() => Ok(()),
602 _ => Err(e),
603 },
604 }
605}
606
607fn stop_instance(client: &Client, name: &str, force: bool, timeout: Duration) -> Result<()> {
608 client
609 .mutate(
610 "PUT",
611 &format!("{}/state", inst_path(name)),
612 Some(&json!({"action": "stop", "force": force, "timeout": timeout.as_secs().max(1)})),
613 &format!("stop {name}"),
614 client.timeouts.state.max(timeout + Duration::from_secs(10)),
615 )
616 .map(|_| ())
617}
618
619fn force_delete(client: &Client, name: &str) -> Result<()> {
620 if let Some(a) = get_actual(client, name)? {
621 if !a.status.eq_ignore_ascii_case("stopped") {
622 let _ = stop_instance(client, name, true, Duration::from_secs(5));
623 }
624 }
625 let started = Instant::now();
628 loop {
629 let r = client.mutate(
630 "DELETE",
631 &inst_path(name),
632 None,
633 &format!("delete {name}"),
634 client.timeouts.other,
635 );
636 match r {
637 Ok(_) => return Ok(()),
638 Err(e) if e.is_not_found() => return Ok(()),
639 Err(e) if e.is_timeout() || started.elapsed() >= Duration::from_secs(60) => {
640 return Err(e);
641 }
642 Err(_) => {
643 std::thread::sleep(Duration::from_secs(2));
644 if let Ok(Some(a)) = get_actual(client, name) {
645 if a.running() {
646 let _ = stop_instance(client, name, true, Duration::from_secs(5));
647 }
648 }
649 }
650 }
651 }
652}
653
654fn add_port_searching(
656 client: &Client,
657 name: &str,
658 device: &str,
659 props: &Props,
660 search: u16,
661) -> Result<String> {
662 let listen = props
663 .get("listen")
664 .cloned()
665 .ok_or_else(|| Error::invalid("port without listen"))?;
666 let Some((proto, host, port)) = split_addr(&listen) else {
667 return Err(Error::invalid(format!("cannot search from {listen}")));
668 };
669 let (proto, host) = (proto.to_string(), host.to_string());
670 let last = port.saturating_add(search);
671 let mut last_err = None;
672 for p in port..=last {
673 if proto == "tcp" {
677 if let Err(e) = std::net::TcpListener::bind(format!("{host}:{p}")) {
678 if e.kind() == std::io::ErrorKind::AddrInUse {
679 continue;
680 }
681 }
682 }
683 let addr = format!("{proto}:{host}:{p}");
684 let mut dev = props.clone();
685 dev.insert("listen".into(), addr.clone());
686 let r = update_instance(
687 client,
688 name,
689 &format!("add port {device}"),
690 &mut |_, devices| {
691 devices.insert(device.to_string(), json!(dev));
692 Ok(())
693 },
694 );
695 match r {
696 Ok(()) => return Ok(addr),
697 Err(e @ Error::Api { .. }) | Err(e @ Error::OperationFailed { .. }) => {
698 last_err = Some(e)
699 }
700 Err(e) => return Err(e),
701 }
702 }
703 Err(Error::invalid(format!(
704 "{name}: no free port for {device} in {proto}:{host}:{port}-{last}{}",
705 last_err
706 .map(|e| format!(" (last error: {e})"))
707 .unwrap_or_default()
708 )))
709}
710
711const OWNER_SCRIPT: &str = r#"set -e
712owner="$1"; path="$2"
713user="${owner%%:*}"
714group=""
715case "$owner" in *:*) group="${owner#*:}" ;; esac
716home=""
717if ent="$(getent passwd "$user")"; then
718 uid="$(printf %s "$ent" | cut -d: -f3)"
719 gid="$(printf %s "$ent" | cut -d: -f4)"
720 home="$(printf %s "$ent" | cut -d: -f6)"
721else
722 case "$user" in ''|*[!0-9]*) echo "isb: no such user: $user" >&2; exit 1 ;; esac
723 uid="$user"; gid="$user"
724fi
725[ -n "$group" ] || group="$gid"
726chown "$uid:$group" "$path"
727# Parents the mount conjured are root-owned; fix those inside the user's home
728# only, and stop at the first one that is not root's.
729[ -n "$home" ] && [ "$home" != / ] || exit 0
730case "$path" in
731 "$home"/*)
732 d="$(dirname "$path")"
733 while [ "$d" != "$home" ] && [ "$d" != "/" ]; do
734 [ "$(stat -c %u "$d")" = 0 ] || break
735 chown "$uid:$group" "$d"
736 d="$(dirname "$d")"
737 done ;;
738esac
739"#;
740
741fn fix_owner(client: &Client, name: &str, path: &str, owner: &str) -> Result<()> {
742 let argv: Vec<String> = ["sh", "-c", OWNER_SCRIPT, "isb-owner", owner, path]
743 .iter()
744 .map(|s| s.to_string())
745 .collect();
746 let out = exec::run_captured(
747 client,
748 name,
749 &argv,
750 &exec::Request::default(),
751 Stdin::Null,
752 Some(Duration::from_secs(60)),
753 )?;
754 if !out.success() {
755 return Err(Error::OperationFailed {
756 step: format!("chown {owner} {path} in {name}"),
757 message: out.stderr_text().trim().to_string(),
758 });
759 }
760 Ok(())
761}
762
763pub fn has_default_route(route_v4: &str, route_v6: &str) -> bool {
765 let v4 = route_v4.lines().skip(1).any(|l| {
766 let f: Vec<&str> = l.split_whitespace().collect();
767 f.len() > 3
768 && f[1] == "00000000"
769 && u32::from_str_radix(f[3], 16)
770 .map(|fl| fl & 1 == 1)
771 .unwrap_or(false)
772 });
773 let v6 = route_v6.lines().any(|l| {
774 let f: Vec<&str> = l.split_whitespace().collect();
775 f.len() >= 10 && f[0].chars().all(|c| c == '0') && f[1] == "00" && f[9] != "lo"
776 });
777 v4 || v6
778}
779
780#[expect(
782 clippy::excessive_nesting,
783 reason = "predates the lint ratchet; split it when next changed"
784)]
785pub fn wait_ready(
786 client: &Client,
787 name: &str,
788 checks: &[ReadyCheck],
789 timeout: Duration,
790 exec_defaults: &ExecDefaults,
791) -> Result<()> {
792 let started = Instant::now();
793 let mut restarted = false;
794 for check in checks {
795 let mut last;
796 let mut stopped_since: Option<Instant> = None;
797 loop {
798 match run_check(client, name, check, exec_defaults) {
799 Ok(true) => break,
800 Ok(false) => last = "not yet".into(),
801 Err(e) => last = e.to_string(),
802 }
803 if let Ok(Some(a)) = get_actual(client, name) {
810 if a.running() || a.status.eq_ignore_ascii_case("starting") {
811 stopped_since = None;
812 } else if stopped_since.get_or_insert_with(Instant::now).elapsed()
813 >= Duration::from_secs(30)
814 {
815 if !restarted {
816 restarted = true;
817 stopped_since = None;
818 if start_instance(client, name).is_ok() {
819 continue;
820 }
821 }
822 return Err(Error::NotReady {
823 sandbox: name.into(),
824 check: check.to_string(),
825 detail: format!(
826 "instance is {} (it stopped while getting ready{}; see `incus info --show-log {name}`)",
827 a.status,
828 if restarted {
829 ", and a restart did not stick"
830 } else {
831 ""
832 }
833 ),
834 waited: started.elapsed(),
835 });
836 }
837 }
838 if started.elapsed() >= timeout {
839 return Err(Error::NotReady {
840 sandbox: name.into(),
841 check: check.to_string(),
842 detail: last,
843 waited: started.elapsed(),
844 });
845 }
846 std::thread::sleep(Duration::from_millis(250));
847 }
848 }
849 Ok(())
850}
851
852fn run_check(
853 client: &Client,
854 name: &str,
855 check: &ReadyCheck,
856 defaults: &ExecDefaults,
857) -> Result<bool> {
858 let cap = |argv: &[&str]| -> Result<ExecOutput> {
859 let argv: Vec<String> = argv.iter().map(|s| s.to_string()).collect();
860 exec::run_captured(
861 client,
862 name,
863 &argv,
864 &exec::Request::default(),
865 Stdin::Null,
866 Some(Duration::from_secs(20)),
867 )
868 };
869 Ok(match check {
870 ReadyCheck::Running => get_actual(client, name)?.is_some_and(|a| a.running()),
871 ReadyCheck::Agent => cap(&["true"])?.success(),
872 ReadyCheck::DefaultRoute => {
873 let v4 = cap(&["cat", "/proc/net/route"])?;
874 let v6 = cap(&["cat", "/proc/net/ipv6_route"]).unwrap_or_default();
875 has_default_route(&v4.stdout_text(), &v6.stdout_text())
876 }
877 ReadyCheck::UserExists(u) => cap(&["getent", "passwd", u])?.success(),
878 ReadyCheck::PathWritable(p) => {
879 let argv = vec!["test".to_string(), "-w".into(), p.clone()];
880 let req = exec::build_request(
881 client,
882 name,
883 &argv,
884 &ExecDefaults {
885 user: defaults.user.clone(),
886 ..Default::default()
887 },
888 &ExecOptions::default(),
889 )?;
890 exec::start_with_timeout(
891 client,
892 name,
893 req,
894 Stdin::Null,
895 Some(Duration::from_secs(20)),
896 )?
897 .collect_output()?
898 .success()
899 }
900 ReadyCheck::Command(argv) => {
901 let argv: Vec<&str> = argv.iter().map(String::as_str).collect();
902 cap(&argv)?.success()
903 }
904 })
905}
906
907pub fn ensure(
909 client: &Client,
910 desired: &Desired,
911 opts: EnsureOptions,
912 report: Reporter<'_>,
913) -> Result<ApplyReport> {
914 let _lock = NameLock::acquire(
915 client.project_name(),
916 &desired.name,
917 opts.lock_wait,
918 &mut |p| {
919 report(&format!(
920 "{}: waiting for another isb holding {}",
921 desired.name,
922 p.display()
923 ))
924 },
925 )?;
926 let plan = plan_desired(client, desired, opts.diff)?;
927 let mut out = apply(client, desired, &plan, report)?;
928 if desired.devices.values().any(|d| d.search.is_some()) {
931 if let Some(inst) = client.get_opt(&inst_path(&desired.name))? {
932 let actual = Actual::from_api(&inst);
933 for (dev, d) in &desired.devices {
934 if d.search.is_none() {
935 continue;
936 }
937 if let Some(listen) = actual.devices.get(dev).and_then(|p| p.get("listen")) {
938 out.ports.insert(dev.clone(), listen.clone());
939 }
940 }
941 }
942 }
943 if !out.restart_needed.is_empty() {
944 report(&format!(
945 "{}: {} changed; takes effect after `isb restart {}`",
946 desired.name,
947 out.restart_needed.join(", "),
948 desired.name
949 ));
950 }
951 if opts.wait_ready {
952 wait_ready(
953 client,
954 &desired.name,
955 &desired.ready,
956 desired.ready_timeout,
957 &desired.exec,
958 )?;
959 if let Some(p) = &desired.egress {
960 crate::egress::after_ready(client, p)?;
961 }
962 }
963 Ok(out)
964}
965
966#[derive(Debug, Clone)]
968pub struct Sandbox {
969 client: Client,
970 name: String,
971 exec_defaults: ExecDefaults,
972 ready: Vec<ReadyCheck>,
973 ready_timeout: Duration,
974}
975
976impl Sandbox {
977 pub(crate) fn from_desired(client: &Client, d: &Desired) -> Sandbox {
978 Sandbox {
979 client: client.clone(),
980 name: d.name.clone(),
981 exec_defaults: d.exec.clone(),
982 ready: d.ready.clone(),
983 ready_timeout: d.ready_timeout,
984 }
985 }
986
987 pub(crate) fn like(client: &Client, name: &str, d: &Desired) -> Sandbox {
990 Sandbox {
991 client: client.clone(),
992 name: name.to_string(),
993 exec_defaults: d.exec.clone(),
994 ready: d.ready.clone(),
995 ready_timeout: d.ready_timeout,
996 }
997 }
998
999 pub fn create(client: &Client, spec: &SandboxSpec) -> Result<Sandbox> {
1002 Self::create_with(
1003 client,
1004 spec,
1005 &VolumeDefs::new(),
1006 EnsureOptions::default(),
1007 &mut |_| {},
1008 )
1009 }
1010
1011 pub fn create_with(
1012 client: &Client,
1013 spec: &SandboxSpec,
1014 defs: &VolumeDefs,
1015 opts: EnsureOptions,
1016 report: Reporter<'_>,
1017 ) -> Result<Sandbox> {
1018 let d = resolve(client, spec, defs, &std::env::current_dir()?)?;
1019 let _lock = NameLock::acquire(
1020 client.project_name(),
1021 &d.name,
1022 Duration::from_secs(900),
1023 &mut |_| {},
1024 )?;
1025 if get_actual(client, &d.name)?.is_some() {
1026 return Err(Error::AlreadyExists(d.name.clone()));
1027 }
1028 let plan = plan_desired(client, &d, opts.diff)?;
1029 apply(client, &d, &plan, report)?;
1030 if opts.wait_ready {
1031 wait_ready(client, &d.name, &d.ready, d.ready_timeout, &d.exec)?;
1032 if let Some(p) = &d.egress {
1033 crate::egress::after_ready(client, p)?;
1034 }
1035 }
1036 Ok(Self::from_desired(client, &d))
1037 }
1038
1039 pub fn connect_or_create(client: &Client, spec: &SandboxSpec) -> Result<Sandbox> {
1042 Self::connect_or_create_with(
1043 client,
1044 spec,
1045 &VolumeDefs::new(),
1046 EnsureOptions::default(),
1047 &mut |_| {},
1048 )
1049 .map(|(s, _)| s)
1050 }
1051
1052 pub fn connect_or_create_with(
1053 client: &Client,
1054 spec: &SandboxSpec,
1055 defs: &VolumeDefs,
1056 opts: EnsureOptions,
1057 report: Reporter<'_>,
1058 ) -> Result<(Sandbox, ApplyReport)> {
1059 let d = resolve(client, spec, defs, &std::env::current_dir()?)?;
1060 let r = ensure(client, &d, opts, report)?;
1061 Ok((Self::from_desired(client, &d), r))
1062 }
1063
1064 pub fn connect_or_create_with_base(
1067 client: &Client,
1068 spec: &SandboxSpec,
1069 defs: &VolumeDefs,
1070 base: &Path,
1071 opts: EnsureOptions,
1072 report: Reporter<'_>,
1073 ) -> Result<(Sandbox, ApplyReport)> {
1074 let d = resolve(client, spec, defs, base)?;
1075 let r = ensure(client, &d, opts, report)?;
1076 Ok((Self::from_desired(client, &d), r))
1077 }
1078
1079 pub fn get(client: &Client, name: &str) -> Result<Sandbox> {
1081 if get_actual(client, name)?.is_none() {
1082 return Err(Error::NotFound(format!("sandbox {name}")));
1083 }
1084 Ok(Sandbox {
1085 client: client.clone(),
1086 name: name.into(),
1087 exec_defaults: ExecDefaults::default(),
1088 ready: vec![ReadyCheck::Running],
1089 ready_timeout: Duration::from_secs(60),
1090 })
1091 }
1092
1093 pub fn list(client: &Client) -> Result<Vec<SandboxInfo>> {
1095 Self::list_with(client, &[])
1096 }
1097
1098 pub fn list_with(client: &Client, labels: &[LabelFilter]) -> Result<Vec<SandboxInfo>> {
1100 let v = client.get("/1.0/instances?recursion=1")?;
1101 let mut out: Vec<SandboxInfo> = v
1102 .as_array()
1103 .map(|a| a.iter().map(SandboxInfo::from_api).collect())
1104 .unwrap_or_default();
1105 out.retain(|i| i.matches(labels));
1106 out.sort_by(|a, b| a.name.cmp(&b.name));
1107 Ok(out)
1108 }
1109
1110 pub fn remove(client: &Client, name: &str, force: bool) -> Result<()> {
1112 let Some(a) = get_actual(client, name)? else {
1113 return Err(Error::NotFound(format!("sandbox {name}")));
1114 };
1115 if a.running() && !force {
1116 return Err(Error::invalid(format!(
1117 "{name} is running; stop it first or force removal"
1118 )));
1119 }
1120 force_delete(client, name)?;
1121 let _ = crate::egress::plumb::teardown(client, name);
1123 Ok(())
1124 }
1125
1126 pub fn name(&self) -> &str {
1127 &self.name
1128 }
1129
1130 pub fn client(&self) -> &Client {
1131 &self.client
1132 }
1133
1134 pub fn with_exec_defaults(mut self, d: ExecDefaults) -> Self {
1136 self.exec_defaults = d;
1137 self
1138 }
1139
1140 pub fn info(&self) -> Result<SandboxInfo> {
1141 let v = self
1142 .client
1143 .get_opt(&inst_path(&self.name))?
1144 .ok_or_else(|| Error::NotFound(format!("sandbox {}", self.name)))?;
1145 Ok(SandboxInfo::from_api(&v))
1146 }
1147
1148 pub fn labels(&self) -> Result<BTreeMap<String, String>> {
1149 Ok(self.info()?.labels)
1150 }
1151
1152 pub fn start(&self) -> Result<()> {
1153 if self.info()?.status.eq_ignore_ascii_case("running") {
1154 return Ok(());
1155 }
1156 let _lock = NameLock::acquire(
1157 self.client.project_name(),
1158 &self.name,
1159 Duration::from_secs(900),
1160 &mut |_| {},
1161 )?;
1162 retry_once(&mut |_| {}, || start_instance(&self.client, &self.name))?;
1163 self.wait_ready()
1164 }
1165
1166 pub fn stop(&self, force: bool, timeout: Duration) -> Result<()> {
1169 if self.info()?.status.eq_ignore_ascii_case("stopped") {
1170 return Ok(());
1171 }
1172 stop_instance(&self.client, &self.name, force, timeout)
1173 }
1174
1175 pub fn restart(&self) -> Result<()> {
1176 self.stop(false, Duration::from_secs(30))?;
1177 self.start()
1178 }
1179
1180 pub fn wait_ready(&self) -> Result<()> {
1182 wait_ready(
1183 &self.client,
1184 &self.name,
1185 &self.ready,
1186 self.ready_timeout,
1187 &self.exec_defaults,
1188 )
1189 }
1190
1191 pub fn exec<I, S>(&self, argv: I) -> Result<ExecOutput>
1194 where
1195 I: IntoIterator<Item = S>,
1196 S: Into<String>,
1197 {
1198 self.exec_with(argv, ExecOptions::default())
1199 }
1200
1201 pub fn exec_with<I, S>(&self, argv: I, opts: ExecOptions) -> Result<ExecOutput>
1202 where
1203 I: IntoIterator<Item = S>,
1204 S: Into<String>,
1205 {
1206 self.exec_stream(argv, opts)?.collect_output()
1207 }
1208
1209 pub fn exec_stream<I, S>(&self, argv: I, opts: ExecOptions) -> Result<ExecStream>
1211 where
1212 I: IntoIterator<Item = S>,
1213 S: Into<String>,
1214 {
1215 let argv: Vec<String> = argv.into_iter().map(Into::into).collect();
1216 let req = exec::build_request(&self.client, &self.name, &argv, &self.exec_defaults, &opts)?;
1217 exec::start_with_timeout(
1218 &self.client,
1219 &self.name,
1220 req,
1221 opts.stdin.clone(),
1222 opts.timeout,
1223 )
1224 }
1225
1226 pub fn attach<I, S>(&self, argv: I, opts: ExecOptions) -> Result<i32>
1229 where
1230 I: IntoIterator<Item = S>,
1231 S: Into<String>,
1232 {
1233 let argv: Vec<String> = argv.into_iter().map(Into::into).collect();
1234 let req = exec::build_request(&self.client, &self.name, &argv, &self.exec_defaults, &opts)?;
1235 exec::attach(
1236 &self.client,
1237 &self.name,
1238 req,
1239 opts.stdin.clone(),
1240 opts.timeout,
1241 )
1242 }
1243
1244 pub fn add_port(&self, port: &PortSpec) -> Result<String> {
1248 let instance_type = if self.info()?.instance_type == "virtual-machine" {
1250 crate::spec::InstanceType::VirtualMachine
1251 } else {
1252 crate::spec::InstanceType::Container
1253 };
1254 let spec = SandboxSpec {
1255 name: Some(self.name.clone()),
1256 image: "unused".into(),
1257 instance_type,
1258 ports: vec![port.clone()],
1259 ..Default::default()
1260 };
1261 let host = HostFacts {
1262 pools: vec!["unused".into()],
1263 ..Default::default()
1264 };
1265 let d = plan::resolve(&spec, &VolumeDefs::new(), &host, Path::new("/"))?;
1266 let (dname, want): (&String, &DesiredDevice) = d
1267 .devices
1268 .iter()
1269 .find(|(k, _)| k.as_str() != "root")
1270 .expect("one port");
1271 let _lock = NameLock::acquire(
1272 self.client.project_name(),
1273 &self.name,
1274 Duration::from_secs(900),
1275 &mut |_| {},
1276 )?;
1277 let info = self.info()?;
1278 if let Some(have) = info.devices.get(dname) {
1279 if plan::device_matches(want, have) {
1280 return Ok(have.get("listen").cloned().unwrap_or_default());
1281 }
1282 self.remove_device_unlocked(dname)?;
1283 }
1284 match want.search {
1285 Some(n) => add_port_searching(&self.client, &self.name, dname, &want.props, n),
1286 None => {
1287 let props = want.props.clone();
1288 update_instance(
1289 &self.client,
1290 &self.name,
1291 &format!("add port {dname}"),
1292 &mut |_, devices| {
1293 devices.insert(dname.clone(), json!(props));
1294 Ok(())
1295 },
1296 )?;
1297 Ok(want.props["listen"].clone())
1298 }
1299 }
1300 }
1301
1302 pub fn remove_device(&self, device: &str) -> Result<bool> {
1304 let _lock = NameLock::acquire(
1305 self.client.project_name(),
1306 &self.name,
1307 Duration::from_secs(900),
1308 &mut |_| {},
1309 )?;
1310 self.remove_device_unlocked(device)
1311 }
1312
1313 fn remove_device_unlocked(&self, device: &str) -> Result<bool> {
1314 if device == "root" {
1315 return Err(Error::invalid("refusing to remove the root disk"));
1316 }
1317 let mut found = false;
1318 update_instance(
1319 &self.client,
1320 &self.name,
1321 &format!("remove device {device}"),
1322 &mut |_, devices| {
1323 found = devices.remove(device).is_some();
1324 Ok(())
1325 },
1326 )?;
1327 Ok(found)
1328 }
1329}
1330
1331#[derive(Debug, Clone, Serialize)]
1333pub struct PruneItem {
1334 pub name: String,
1335 pub path: String,
1336 pub deleted: bool,
1337}
1338
1339pub fn prune_missing_path(
1343 client: &Client,
1344 label: &str,
1345 dry_run: bool,
1346 report: Reporter<'_>,
1347) -> Result<Vec<PruneItem>> {
1348 let mut out = Vec::new();
1349 for i in Sandbox::list_with(client, &[LabelFilter::parse(label)])? {
1350 let Some(path) = i.labels.get(label) else {
1351 continue;
1352 };
1353 if path.is_empty() || !path.starts_with('/') || Path::new(path).exists() {
1354 continue;
1355 }
1356 if dry_run {
1357 report(&format!(
1358 "would delete {} ({label}={path} is gone)",
1359 i.name
1360 ));
1361 } else {
1362 report(&format!("deleting {} ({label}={path} is gone)", i.name));
1363 force_delete(client, &i.name)?;
1364 }
1365 out.push(PruneItem {
1366 name: i.name,
1367 path: path.clone(),
1368 deleted: !dry_run,
1369 });
1370 }
1371 Ok(out)
1372}
1373
1374#[cfg(test)]
1375mod tests;