Skip to main content

isb_core/
sandbox.rs

1//! The sandbox API and the apply engine.
2
3use 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
23/// Config key recording which isb call created an instance. A half-created
24/// instance is only ever cleaned up by the call whose token it carries.
25pub const CREATE_TOKEN_KEY: &str = "user.isb.create-token";
26
27/// Progress callback: one human-readable line per step.
28pub type Reporter<'a> = &'a mut dyn FnMut(&str);
29
30/// A summary of an instance, as listed.
31#[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    /// `user.*` config keys with the prefix stripped.
38    pub labels: BTreeMap<String, String>,
39    pub config: BTreeMap<String, String>,
40    /// Instance-local devices.
41    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            // isb's own bookkeeping (user.isb.create-token) is not a label.
55            .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    /// Whether this instance matches every filter.
79    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/// `key` (present) or `key=value` (equal).
91#[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/// What `apply` did.
113#[derive(Debug, Clone, Default, Serialize)]
114pub struct ApplyReport {
115    pub name: String,
116    pub created: bool,
117    pub applied: Vec<Action>,
118    /// Device name -> listen address chosen for searched ports.
119    pub ports: BTreeMap<String, String>,
120    /// Config keys changed that only take effect after a restart.
121    pub restart_needed: Vec<String>,
122}
123
124/// Options for ensure / up.
125#[derive(Debug, Clone, Copy)]
126pub struct EnsureOptions {
127    pub diff: DiffOptions,
128    /// Run readiness checks after applying.
129    pub wait_ready: bool,
130    /// How long to wait for another isb holding this sandbox's lock.
131    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
148/// Gather host facts (pools, subids, path map).
149pub 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 info = client.server_info()?;
162    let ext = |n: &str| {
163        info["api_extensions"]
164            .as_array()
165            .is_some_and(|a| a.iter().any(|e| e == n))
166    };
167    Ok(HostFacts {
168        subids: SubIds::read_host(),
169        pools,
170        path_map: HostFacts::detect_path_map(),
171        initial_copy: ext("disk_initial_copy"),
172        initial_owner: ext("storage_initial_owner"),
173        incus_version: client.server_version()?,
174        invoking_ids: crate::idmap::invoking_ids(),
175        shared_root: shared_root(),
176        org: crate::org::OrgId::from_incus_project(client.project_name()),
177        registry: crate::registry::info(client)?.map(|i| i.addr),
178        project: client.project_name().to_string(),
179    })
180}
181
182/// On macOS incusd runs in the `isb machine` VM, which sees only `$HOME`.
183fn shared_root() -> Option<String> {
184    if !cfg!(target_os = "macos") {
185        return None;
186    }
187    let home = std::env::var_os("HOME").filter(|h| !h.is_empty())?;
188    let p = std::path::PathBuf::from(home);
189    Some(p.canonicalize().unwrap_or(p).to_string_lossy().into_owned())
190}
191
192/// Resolve a spec against this host. Relative bind paths anchor at `base`.
193pub fn resolve(
194    client: &Client,
195    spec: &SandboxSpec,
196    defs: &VolumeDefs,
197    base: &Path,
198) -> Result<Desired> {
199    let host = host_facts(client)?;
200    plan::resolve(spec, defs, &host, base)
201}
202
203fn get_actual(client: &Client, name: &str) -> Result<Option<Actual>> {
204    Ok(client
205        .get_opt(&inst_path(name))?
206        .map(|v| Actual::from_api(&v)))
207}
208
209fn volume_exists(client: &Client, pool: &str, name: &str) -> Result<bool> {
210    Ok(client
211        .get_opt(&format!(
212            "/1.0/storage-pools/{}/volumes/custom/{}",
213            encode_segment(pool),
214            encode_segment(name)
215        ))?
216        .is_some())
217}
218
219/// The fingerprint of a local image, by alias or fingerprint (prefix).
220fn local_image(client: &Client, alias: &str) -> Result<Option<String>> {
221    if let Some(a) = client.get_opt(&format!("/1.0/images/aliases/{}", encode_segment(alias)))? {
222        return Ok(a.get("target").and_then(Value::as_str).map(String::from));
223    }
224    if alias.len() >= 12 && alias.chars().all(|c| c.is_ascii_hexdigit()) {
225        if let Some(i) = client.get_opt(&format!("/1.0/images/{}", encode_segment(alias)))? {
226            return Ok(i
227                .get("fingerprint")
228                .and_then(Value::as_str)
229                .map(String::from));
230        }
231    }
232    Ok(None)
233}
234
235/// Compute the plan for one resolved sandbox.
236pub fn plan_desired(client: &Client, desired: &Desired, opts: DiffOptions) -> Result<SandboxPlan> {
237    let actual = get_actual(client, &desired.name)?;
238    if actual.is_none()
239        && desired.image.server.is_none()
240        && local_image(client, &desired.image.alias)?.is_none()
241    {
242        return Err(images::missing_on(client, &desired.image.alias));
243    }
244    let mut missing = Vec::new();
245    for v in &desired.volumes {
246        if !volume_exists(client, &v.pool, &v.name)? {
247            missing.push((v.pool.clone(), v.name.clone()));
248        }
249    }
250    let mut plan = plan::diff(desired, actual.as_ref(), &missing, opts)?;
251    // The image is fixed at creation; say so when the local image has moved on
252    // (e.g. dev-base was republished), so a recreate is a visible choice.
253    if let (Some(a), None) = (&actual, &desired.image.server) {
254        if let (Some(built), Some(now)) = (
255            a.config.get("volatile.base_image"),
256            local_image(client, &desired.image.alias)?,
257        ) {
258            if *built != now {
259                plan.actions.push(Action::Note {
260                    message: format!(
261                        "image {} is now {} but this instance was built from {}; recreate to pick it up",
262                        desired.image.alias,
263                        &now[..12.min(now.len())],
264                        &built[..12.min(built.len())]
265                    ),
266                });
267            }
268        }
269    }
270    Ok(plan)
271}
272
273fn random_token() -> String {
274    let mut b = [0u8; 12];
275    if let Ok(mut f) = std::fs::File::open("/dev/urandom") {
276        use std::io::Read;
277        let _ = f.read_exact(&mut b);
278    }
279    let t = std::time::SystemTime::now()
280        .duration_since(std::time::UNIX_EPOCH)
281        .map(|d| d.as_nanos())
282        .unwrap_or(0);
283    format!(
284        "{}{:x}",
285        b.iter().map(|x| format!("{x:02x}")).collect::<String>(),
286        t & 0xffff
287    )
288}
289
290/// Retry a step once after a deadline.
291fn retry_once<T>(report: &mut dyn FnMut(&str), mut f: impl FnMut() -> Result<T>) -> Result<T> {
292    match f() {
293        Err(e) if e.is_timeout() => {
294            report(&format!("{e}; retrying once"));
295            f()
296        }
297        r => r,
298    }
299}
300
301/// Apply a plan. Changes only what the plan says; a correct device is never touched.
302pub fn apply(
303    client: &Client,
304    desired: &Desired,
305    plan: &SandboxPlan,
306    report: Reporter<'_>,
307) -> Result<ApplyReport> {
308    let name = &desired.name;
309    let mut out = ApplyReport {
310        name: name.clone(),
311        ..Default::default()
312    };
313    let mut pending: Vec<&Action> = Vec::new();
314    if let Some(p) = &desired.egress {
315        crate::egress::before_apply(client, p, report)?;
316    }
317    for action in &plan.actions {
318        match action {
319            Action::SetConfig { .. }
320            | Action::AddDevice { .. }
321            | Action::ReplaceDevice { .. }
322            | Action::RemoveDevice { .. } => {
323                pending.push(action);
324                continue;
325            }
326            _ => {}
327        }
328        // Batched config/device changes land before whatever comes next (a start,
329        // most importantly, so a stopped instance boots with the right devices).
330        flush_updates(client, desired, &mut pending, &mut out, report)?;
331        match action {
332            Action::Note { message } => report(&format!("{name}: note: {message}")),
333            Action::CreateVolume {
334                pool,
335                volume,
336                config,
337            } => {
338                report(&format!("{name}: creating volume {volume} on {pool}"));
339                retry_once(report, || {
340                    crate::volume::ensure(client, pool, volume, config).map(|_| ())
341                })?;
342            }
343            Action::CreateInstance { .. } => {
344                report(&format!("{name}: creating from {}", desired.image.spec));
345                create_instance(client, desired, report)?;
346                out.created = true;
347            }
348            Action::StartInstance => {
349                if let (true, Some(h)) = (out.created, &desired.before_start) {
350                    (h.0)(client, name)?;
351                }
352                report(&format!("{name}: starting"));
353                retry_once(report, || start_instance(client, name))?;
354            }
355            Action::AddPort {
356                device,
357                props,
358                search,
359            } => {
360                let listen = add_port_searching(client, name, device, props, *search)?;
361                report(&format!("{name}: port {device} listening on {listen}"));
362                out.ports.insert(device.clone(), listen);
363            }
364            Action::FixOwner {
365                path,
366                owner,
367                mode,
368                fresh_only,
369            } => {
370                if desired.instance_type == crate::spec::InstanceType::VirtualMachine {
371                    // In-guest work needs the VM's agent, which starts after boot.
372                    wait_ready(
373                        client,
374                        name,
375                        &[ReadyCheck::Agent],
376                        desired.ready_timeout,
377                        &desired.exec,
378                    )?;
379                }
380                let fix = crate::owner::Fix {
381                    path,
382                    owner: owner.as_deref(),
383                    mode: mode
384                        .as_deref()
385                        .map(crate::owner::parse_mode)
386                        .transpose()
387                        .map_err(Error::invalid)?,
388                    fresh_only: *fresh_only,
389                };
390                if crate::owner::apply(client, name, &fix)? {
391                    report(&format!("{name}: {action}"));
392                } else {
393                    report(&format!(
394                        "{name}: {path} came seeded from the image; owner left as is"
395                    ));
396                }
397            }
398            _ => unreachable!("batched above"),
399        }
400        out.applied.push(action.clone());
401    }
402    flush_updates(client, desired, &mut pending, &mut out, report)?;
403    Ok(out)
404}
405
406fn flush_updates(
407    client: &Client,
408    desired: &Desired,
409    pending: &mut Vec<&Action>,
410    out: &mut ApplyReport,
411    report: &mut dyn FnMut(&str),
412) -> Result<()> {
413    if pending.is_empty() {
414        return Ok(());
415    }
416    let name = desired.name.as_str();
417    for a in pending.iter() {
418        report(&format!("{name}: {a}"));
419    }
420    let actions: Vec<Action> = pending.iter().map(|a| (*a).clone()).collect();
421    retry_once(report, || {
422        update_instance(
423            client,
424            name,
425            &format!("update {name}"),
426            &mut |config, devices| {
427                for a in &actions {
428                    match a {
429                        Action::SetConfig {
430                            key, to, secret, ..
431                        } => {
432                            // A secret's action carries a placeholder.
433                            let v = match (secret, desired.config.get(key)) {
434                                (true, Some(v)) => v,
435                                _ => to,
436                            };
437                            config.insert(key.clone(), json!(v));
438                        }
439                        Action::AddDevice { device, props } => {
440                            devices.insert(device.clone(), json!(props));
441                        }
442                        Action::ReplaceDevice {
443                            device,
444                            replaces,
445                            to,
446                            ..
447                        } => {
448                            devices.remove(replaces);
449                            devices.insert(device.clone(), json!(to));
450                        }
451                        Action::RemoveDevice { device, .. } => {
452                            devices.remove(device);
453                        }
454                        _ => {}
455                    }
456                }
457                Ok(())
458            },
459        )
460    })?;
461    for a in pending.drain(..) {
462        if let Action::SetConfig {
463            key, restart: true, ..
464        } = a
465        {
466            out.restart_needed.push(key.clone());
467        }
468        out.applied.push(a.clone());
469    }
470    Ok(())
471}
472
473type Obj = serde_json::Map<String, Value>;
474
475/// Read-modify-write of an instance under `If-Match`, so a concurrent change by
476/// another tool is never silently overwritten (a 412 re-reads and retries).
477#[doc(hidden)]
478pub fn update_instance(
479    client: &Client,
480    name: &str,
481    step: &str,
482    modify: &mut dyn FnMut(&mut Obj, &mut Obj) -> Result<()>,
483) -> Result<()> {
484    let path = inst_path(name);
485    for attempt in 0..5 {
486        let (inst, etag) = client.get_etag(&path)?;
487        let mut config = inst
488            .get("config")
489            .and_then(Value::as_object)
490            .cloned()
491            .unwrap_or_default();
492        let mut devices = inst
493            .get("devices")
494            .and_then(Value::as_object)
495            .cloned()
496            .unwrap_or_default();
497        modify(&mut config, &mut devices)?;
498        let body = json!({
499            "architecture": inst.get("architecture"),
500            "config": config,
501            "devices": devices,
502            "ephemeral": inst.get("ephemeral"),
503            "profiles": inst.get("profiles"),
504            "stateful": inst.get("stateful"),
505            "description": inst.get("description"),
506        });
507        match client.mutate_if_match(
508            "PUT",
509            &path,
510            &body,
511            etag.as_deref(),
512            step,
513            client.timeouts.other,
514        ) {
515            Err(Error::Api { status: 412, .. }) if attempt < 4 => continue,
516            r => return r.map(|_| ()),
517        }
518    }
519    unreachable!()
520}
521
522fn create_instance(client: &Client, desired: &Desired, report: &mut dyn FnMut(&str)) -> Result<()> {
523    let fingerprint = match &desired.image.server {
524        None => Some(
525            local_image(client, &desired.image.alias)?
526                .ok_or_else(|| images::missing_on(client, &desired.image.alias))?,
527        ),
528        Some(_) => None,
529    };
530    let step = format!("create instance {}", desired.name);
531    let mut last_err = None;
532    for attempt in 0..2 {
533        let token = random_token();
534        let mut config = desired.config.clone();
535        config.insert(CREATE_TOKEN_KEY.into(), token.clone());
536        let devices: BTreeMap<&String, &Props> = desired
537            .devices
538            .iter()
539            .filter(|(_, d)| d.search.is_none())
540            .map(|(k, d)| (k, &d.props))
541            .collect();
542        let body = json!({
543            "name": desired.name,
544            "type": desired.instance_type.as_api(),
545            "source": desired.image.to_api(fingerprint.as_deref()),
546            "config": config,
547            "devices": devices,
548            "profiles": desired.profiles,
549        });
550        match client.mutate(
551            "POST",
552            "/1.0/instances",
553            Some(&body),
554            &step,
555            client.timeouts.create,
556        ) {
557            Ok(_) => return Ok(()),
558            Err(e) if e.is_timeout() => {
559                report(&format!("{e}"));
560                if let Error::OperationTimeout { operation, .. } = &e {
561                    // Most create operations cannot be cancelled. Let it settle so
562                    // the cleanup sees what it actually did.
563                    if client
564                        .wait_operation(operation, &step, client.timeouts.settle)
565                        .is_ok()
566                    {
567                        report(&format!("{step}: finished late; keeping it"));
568                        return Ok(());
569                    }
570                }
571                if let Err(ce) = cleanup_half_created(client, &desired.name, &token, report) {
572                    report(&format!("{}: cleanup failed: {ce}", desired.name));
573                    if matches!(ce, Error::AlreadyExists(_)) {
574                        return Err(e);
575                    }
576                }
577                if attempt == 0 {
578                    report(&format!("{step}: retrying once"));
579                }
580                last_err = Some(e);
581            }
582            Err(e) => {
583                // A failed create may still have left a stub behind.
584                if !e.is_conflict() {
585                    let _ = cleanup_half_created(client, &desired.name, &token, report);
586                }
587                return Err(e);
588            }
589        }
590    }
591    Err(last_err.expect("loop ran"))
592}
593
594/// Delete a half-created instance, but only if it carries `token` in
595/// `user.isb.create-token`, i.e. only if the call holding that token created it.
596/// Anything else (another owner's instance, one created by hand) is refused with
597/// [`Error::AlreadyExists`] and left alone. Missing is fine.
598pub fn cleanup_half_created(
599    client: &Client,
600    name: &str,
601    token: &str,
602    report: &mut dyn FnMut(&str),
603) -> Result<()> {
604    let Some(inst) = client.get_opt(&inst_path(name))? else {
605        return Ok(());
606    };
607    let a = Actual::from_api(&inst);
608    if a.config.get(CREATE_TOKEN_KEY).map(String::as_str) != Some(token) {
609        report(&format!(
610            "{name}: exists but was not created by this call; leaving it alone"
611        ));
612        return Err(Error::AlreadyExists(name.to_string()));
613    }
614    report(&format!("{name}: removing half-created instance"));
615    force_delete(client, name)
616}
617
618fn start_instance(client: &Client, name: &str) -> Result<()> {
619    let r = client.mutate(
620        "PUT",
621        &format!("{}/state", inst_path(name)),
622        Some(&json!({"action": "start", "timeout": 30})),
623        &format!("start {name}"),
624        client.timeouts.state,
625    );
626    match r {
627        Ok(_) => Ok(()),
628        // Someone else started it in between: that is the state we wanted.
629        Err(e) => match get_actual(client, name) {
630            Ok(Some(a)) if a.running() => Ok(()),
631            _ => Err(e),
632        },
633    }
634}
635
636fn stop_instance(client: &Client, name: &str, force: bool, timeout: Duration) -> Result<()> {
637    client
638        .mutate(
639            "PUT",
640            &format!("{}/state", inst_path(name)),
641            Some(&json!({"action": "stop", "force": force, "timeout": timeout.as_secs().max(1)})),
642            &format!("stop {name}"),
643            client.timeouts.state.max(timeout + Duration::from_secs(10)),
644        )
645        .map(|_| ())
646}
647
648fn force_delete(client: &Client, name: &str) -> Result<()> {
649    if let Some(a) = get_actual(client, name)? {
650        if !a.status.eq_ignore_ascii_case("stopped") {
651            let _ = stop_instance(client, name, true, Duration::from_secs(5));
652        }
653    }
654    // An instance mid-transition (rebooting, stopping) refuses deletion for a
655    // moment; retry for a bounded while before giving up.
656    let started = Instant::now();
657    loop {
658        let r = client.mutate(
659            "DELETE",
660            &inst_path(name),
661            None,
662            &format!("delete {name}"),
663            client.timeouts.other,
664        );
665        match r {
666            Ok(_) => return Ok(()),
667            Err(e) if e.is_not_found() => return Ok(()),
668            Err(e) if e.is_timeout() || started.elapsed() >= Duration::from_secs(60) => {
669                return Err(e);
670            }
671            Err(_) => {
672                std::thread::sleep(Duration::from_secs(2));
673                if let Ok(Some(a)) = get_actual(client, name) {
674                    if a.running() {
675                        let _ = stop_instance(client, name, true, Duration::from_secs(5));
676                    }
677                }
678            }
679        }
680    }
681}
682
683/// Add a host-bound proxy, stepping past taken ports. Returns the listen address.
684fn add_port_searching(
685    client: &Client,
686    name: &str,
687    device: &str,
688    props: &Props,
689    search: u16,
690) -> Result<String> {
691    let listen = props
692        .get("listen")
693        .cloned()
694        .ok_or_else(|| Error::invalid("port without listen"))?;
695    let Some((proto, host, port)) = split_addr(&listen) else {
696        return Err(Error::invalid(format!("cannot search from {listen}")));
697    };
698    let (proto, host) = (proto.to_string(), host.to_string());
699    let last = port.saturating_add(search);
700    let mut last_err = None;
701    for p in port..=last {
702        // Probe locally first: cheaper than a failed device add, and it steps
703        // around listeners that are not incus devices. An address this process
704        // cannot bind at all (not local here) is left to incus to judge.
705        if proto == "tcp" {
706            if let Err(e) = std::net::TcpListener::bind(format!("{host}:{p}")) {
707                if e.kind() == std::io::ErrorKind::AddrInUse {
708                    continue;
709                }
710            }
711        }
712        let addr = format!("{proto}:{host}:{p}");
713        let mut dev = props.clone();
714        dev.insert("listen".into(), addr.clone());
715        let r = update_instance(
716            client,
717            name,
718            &format!("add port {device}"),
719            &mut |_, devices| {
720                devices.insert(device.to_string(), json!(dev));
721                Ok(())
722            },
723        );
724        match r {
725            Ok(()) => return Ok(addr),
726            Err(e @ Error::Api { .. }) | Err(e @ Error::OperationFailed { .. }) => {
727                last_err = Some(e)
728            }
729            Err(e) => return Err(e),
730        }
731    }
732    Err(Error::invalid(format!(
733        "{name}: no free port for {device} in {proto}:{host}:{port}-{last}{}",
734        last_err
735            .map(|e| format!(" (last error: {e})"))
736            .unwrap_or_default()
737    )))
738}
739
740/// Parse `/proc/net/route` and `/proc/net/ipv6_route` for a default route.
741pub fn has_default_route(route_v4: &str, route_v6: &str) -> bool {
742    let v4 = route_v4.lines().skip(1).any(|l| {
743        let f: Vec<&str> = l.split_whitespace().collect();
744        f.len() > 3
745            && f[1] == "00000000"
746            && u32::from_str_radix(f[3], 16)
747                .map(|fl| fl & 1 == 1)
748                .unwrap_or(false)
749    });
750    let v6 = route_v6.lines().any(|l| {
751        let f: Vec<&str> = l.split_whitespace().collect();
752        f.len() >= 10 && f[0].chars().all(|c| c == '0') && f[1] == "00" && f[9] != "lo"
753    });
754    v4 || v6
755}
756
757/// Run readiness checks in order, each polled until the shared deadline.
758#[expect(
759    clippy::excessive_nesting,
760    reason = "predates the lint ratchet; split it when next changed"
761)]
762pub fn wait_ready(
763    client: &Client,
764    name: &str,
765    checks: &[ReadyCheck],
766    timeout: Duration,
767    exec_defaults: &ExecDefaults,
768) -> Result<()> {
769    let started = Instant::now();
770    let mut restarted = false;
771    for check in checks {
772        let mut last;
773        let mut stopped_since: Option<Instant> = None;
774        loop {
775            match run_check(client, name, check, exec_defaults) {
776                Ok(true) => break,
777                Ok(false) => last = "not yet".into(),
778                Err(e) => last = e.to_string(),
779            }
780            // An instance that stopped (crashed, powered off) will not get ready
781            // by waiting. A reboot (common on a VM's first boot) passes through
782            // Stopped briefly, so only a sustained stop counts. incus 7.0.1 sometimes
783            // fails to complete a guest-initiated reboot (two concurrent onStop
784            // hooks; fixed upstream in lxc/incus#3997), leaving it Stopped, so
785            // start it once more before failing fast instead of at the deadline.
786            if let Ok(Some(a)) = get_actual(client, name) {
787                if a.running() || a.status.eq_ignore_ascii_case("starting") {
788                    stopped_since = None;
789                } else if stopped_since.get_or_insert_with(Instant::now).elapsed()
790                    >= Duration::from_secs(30)
791                {
792                    if !restarted {
793                        restarted = true;
794                        stopped_since = None;
795                        if start_instance(client, name).is_ok() {
796                            continue;
797                        }
798                    }
799                    return Err(Error::NotReady {
800                        sandbox: name.into(),
801                        check: check.to_string(),
802                        detail: format!(
803                            "instance is {} (it stopped while getting ready{}; see `incus info --show-log {name}`)",
804                            a.status,
805                            if restarted {
806                                ", and a restart did not stick"
807                            } else {
808                                ""
809                            }
810                        ),
811                        waited: started.elapsed(),
812                    });
813                }
814            }
815            if started.elapsed() >= timeout {
816                return Err(Error::NotReady {
817                    sandbox: name.into(),
818                    check: check.to_string(),
819                    detail: last,
820                    waited: started.elapsed(),
821                });
822            }
823            std::thread::sleep(Duration::from_millis(250));
824        }
825    }
826    Ok(())
827}
828
829fn run_check(
830    client: &Client,
831    name: &str,
832    check: &ReadyCheck,
833    defaults: &ExecDefaults,
834) -> Result<bool> {
835    let cap = |argv: &[&str]| -> Result<ExecOutput> {
836        let argv: Vec<String> = argv.iter().map(|s| s.to_string()).collect();
837        exec::run_captured(
838            client,
839            name,
840            &argv,
841            &exec::Request::default(),
842            Stdin::Null,
843            Some(Duration::from_secs(20)),
844        )
845    };
846    Ok(match check {
847        ReadyCheck::Running => get_actual(client, name)?.is_some_and(|a| a.running()),
848        ReadyCheck::Agent => cap(&["true"])?.success(),
849        ReadyCheck::DefaultRoute => {
850            let v4 = cap(&["cat", "/proc/net/route"])?;
851            let v6 = cap(&["cat", "/proc/net/ipv6_route"]).unwrap_or_default();
852            has_default_route(&v4.stdout_text(), &v6.stdout_text())
853        }
854        ReadyCheck::UserExists(u) => cap(&["getent", "passwd", u])?.success(),
855        ReadyCheck::PathWritable(p) => {
856            let argv = vec!["test".to_string(), "-w".into(), p.clone()];
857            let req = exec::build_request(
858                client,
859                name,
860                &argv,
861                &ExecDefaults {
862                    user: defaults.user.clone(),
863                    ..Default::default()
864                },
865                &ExecOptions::default(),
866            )?;
867            exec::start_with_timeout(
868                client,
869                name,
870                req,
871                Stdin::Null,
872                Some(Duration::from_secs(20)),
873            )?
874            .collect_output()?
875            .success()
876        }
877        ReadyCheck::Command(argv) => {
878            let argv: Vec<&str> = argv.iter().map(String::as_str).collect();
879            cap(&argv)?.success()
880        }
881    })
882}
883
884/// Lock, plan, apply, wait for readiness.
885pub fn ensure(
886    client: &Client,
887    desired: &Desired,
888    opts: EnsureOptions,
889    report: Reporter<'_>,
890) -> Result<ApplyReport> {
891    let _lock = NameLock::acquire(
892        client.project_name(),
893        &desired.name,
894        opts.lock_wait,
895        &mut |p| {
896            report(&format!(
897                "{}: waiting for another isb holding {}",
898                desired.name,
899                p.display()
900            ))
901        },
902    )?;
903    let plan = plan_desired(client, desired, opts.diff)?;
904    let mut out = apply(client, desired, &plan, report)?;
905    // Report the settled listen address of every searched port, whether this
906    // call added it or found it already correct somewhere in its range.
907    if desired.devices.values().any(|d| d.search.is_some()) {
908        if let Some(inst) = client.get_opt(&inst_path(&desired.name))? {
909            let actual = Actual::from_api(&inst);
910            for (dev, d) in &desired.devices {
911                if d.search.is_none() {
912                    continue;
913                }
914                if let Some(listen) = actual.devices.get(dev).and_then(|p| p.get("listen")) {
915                    out.ports.insert(dev.clone(), listen.clone());
916                }
917            }
918        }
919    }
920    if !out.restart_needed.is_empty() {
921        report(&format!(
922            "{}: {} changed; takes effect after `isb restart {}`",
923            desired.name,
924            out.restart_needed.join(", "),
925            desired.name
926        ));
927    }
928    if opts.wait_ready {
929        wait_ready(
930            client,
931            &desired.name,
932            &desired.ready,
933            desired.ready_timeout,
934            &desired.exec,
935        )?;
936        if let Some(p) = &desired.egress {
937            crate::egress::after_ready(client, p)?;
938        }
939    }
940    Ok(out)
941}
942
943/// A handle on one sandbox.
944#[derive(Debug, Clone)]
945pub struct Sandbox {
946    client: Client,
947    name: String,
948    exec_defaults: ExecDefaults,
949    ready: Vec<ReadyCheck>,
950    ready_timeout: Duration,
951}
952
953impl Sandbox {
954    pub(crate) fn from_desired(client: &Client, d: &Desired) -> Sandbox {
955        Sandbox {
956            client: client.clone(),
957            name: d.name.clone(),
958            exec_defaults: d.exec.clone(),
959            ready: d.ready.clone(),
960            ready_timeout: d.ready_timeout,
961        }
962    }
963
964    /// A handle on `name` with the exec defaults and readiness of `d` (a
965    /// sibling built from the same spec, such as another replica).
966    pub(crate) fn like(client: &Client, name: &str, d: &Desired) -> Sandbox {
967        Sandbox {
968            client: client.clone(),
969            name: name.to_string(),
970            exec_defaults: d.exec.clone(),
971            ready: d.ready.clone(),
972            ready_timeout: d.ready_timeout,
973        }
974    }
975
976    /// Create and start a new sandbox; fails if it already exists. Relative bind
977    /// paths resolve against the current directory.
978    pub fn create(client: &Client, spec: &SandboxSpec) -> Result<Sandbox> {
979        Self::create_with(
980            client,
981            spec,
982            &VolumeDefs::new(),
983            EnsureOptions::default(),
984            &mut |_| {},
985        )
986    }
987
988    pub fn create_with(
989        client: &Client,
990        spec: &SandboxSpec,
991        defs: &VolumeDefs,
992        opts: EnsureOptions,
993        report: Reporter<'_>,
994    ) -> Result<Sandbox> {
995        let d = resolve(client, spec, defs, &std::env::current_dir()?)?;
996        let _lock = NameLock::acquire(
997            client.project_name(),
998            &d.name,
999            Duration::from_secs(900),
1000            &mut |_| {},
1001        )?;
1002        if get_actual(client, &d.name)?.is_some() {
1003            return Err(Error::AlreadyExists(d.name.clone()));
1004        }
1005        let plan = plan_desired(client, &d, opts.diff)?;
1006        apply(client, &d, &plan, report)?;
1007        if opts.wait_ready {
1008            wait_ready(client, &d.name, &d.ready, d.ready_timeout, &d.exec)?;
1009            if let Some(p) = &d.egress {
1010                crate::egress::after_ready(client, p)?;
1011            }
1012        }
1013        Ok(Self::from_desired(client, &d))
1014    }
1015
1016    /// Reconcile a sandbox to `spec`, creating it if needed: only what differs is
1017    /// changed, and a correct device is never touched.
1018    pub fn connect_or_create(client: &Client, spec: &SandboxSpec) -> Result<Sandbox> {
1019        Self::connect_or_create_with(
1020            client,
1021            spec,
1022            &VolumeDefs::new(),
1023            EnsureOptions::default(),
1024            &mut |_| {},
1025        )
1026        .map(|(s, _)| s)
1027    }
1028
1029    pub fn connect_or_create_with(
1030        client: &Client,
1031        spec: &SandboxSpec,
1032        defs: &VolumeDefs,
1033        opts: EnsureOptions,
1034        report: Reporter<'_>,
1035    ) -> Result<(Sandbox, ApplyReport)> {
1036        let d = resolve(client, spec, defs, &std::env::current_dir()?)?;
1037        let r = ensure(client, &d, opts, report)?;
1038        Ok((Self::from_desired(client, &d), r))
1039    }
1040
1041    /// [`Sandbox::connect_or_create_with`], with relative bind paths
1042    /// resolving against `base` instead of the current directory.
1043    pub fn connect_or_create_with_base(
1044        client: &Client,
1045        spec: &SandboxSpec,
1046        defs: &VolumeDefs,
1047        base: &Path,
1048        opts: EnsureOptions,
1049        report: Reporter<'_>,
1050    ) -> Result<(Sandbox, ApplyReport)> {
1051        let d = resolve(client, spec, defs, base)?;
1052        let r = ensure(client, &d, opts, report)?;
1053        Ok((Self::from_desired(client, &d), r))
1054    }
1055
1056    /// Handle on an existing sandbox.
1057    pub fn get(client: &Client, name: &str) -> Result<Sandbox> {
1058        if get_actual(client, name)?.is_none() {
1059            return Err(Error::NotFound(format!("sandbox {name}")));
1060        }
1061        Ok(Sandbox {
1062            client: client.clone(),
1063            name: name.into(),
1064            exec_defaults: ExecDefaults::default(),
1065            ready: vec![ReadyCheck::Running],
1066            ready_timeout: Duration::from_secs(60),
1067        })
1068    }
1069
1070    /// All instances in the client's project.
1071    pub fn list(client: &Client) -> Result<Vec<SandboxInfo>> {
1072        Self::list_with(client, &[])
1073    }
1074
1075    /// Instances whose labels match every filter.
1076    pub fn list_with(client: &Client, labels: &[LabelFilter]) -> Result<Vec<SandboxInfo>> {
1077        let v = client.get("/1.0/instances?recursion=1")?;
1078        let mut out: Vec<SandboxInfo> = v
1079            .as_array()
1080            .map(|a| a.iter().map(SandboxInfo::from_api).collect())
1081            .unwrap_or_default();
1082        out.retain(|i| i.matches(labels));
1083        out.sort_by(|a, b| a.name.cmp(&b.name));
1084        Ok(out)
1085    }
1086
1087    /// Delete a sandbox. A running one needs `force` (it is stopped first).
1088    pub fn remove(client: &Client, name: &str, force: bool) -> Result<()> {
1089        let Some(a) = get_actual(client, name)? else {
1090            return Err(Error::NotFound(format!("sandbox {name}")));
1091        };
1092        if a.running() && !force {
1093            return Err(Error::invalid(format!(
1094                "{name} is running; stop it first or force removal"
1095            )));
1096        }
1097        force_delete(client, name)?;
1098        // The egress bridge and ACL go with the sandbox; a daemon sweeps any left behind.
1099        let _ = crate::egress::plumb::teardown(client, name);
1100        Ok(())
1101    }
1102
1103    pub fn name(&self) -> &str {
1104        &self.name
1105    }
1106
1107    pub fn client(&self) -> &Client {
1108        &self.client
1109    }
1110
1111    /// Use these exec defaults for this handle.
1112    pub fn with_exec_defaults(mut self, d: ExecDefaults) -> Self {
1113        self.exec_defaults = d;
1114        self
1115    }
1116
1117    pub fn info(&self) -> Result<SandboxInfo> {
1118        let v = self
1119            .client
1120            .get_opt(&inst_path(&self.name))?
1121            .ok_or_else(|| Error::NotFound(format!("sandbox {}", self.name)))?;
1122        Ok(SandboxInfo::from_api(&v))
1123    }
1124
1125    pub fn labels(&self) -> Result<BTreeMap<String, String>> {
1126        Ok(self.info()?.labels)
1127    }
1128
1129    pub fn start(&self) -> Result<()> {
1130        if self.info()?.status.eq_ignore_ascii_case("running") {
1131            return Ok(());
1132        }
1133        let _lock = NameLock::acquire(
1134            self.client.project_name(),
1135            &self.name,
1136            Duration::from_secs(900),
1137            &mut |_| {},
1138        )?;
1139        retry_once(&mut |_| {}, || start_instance(&self.client, &self.name))?;
1140        self.wait_ready()
1141    }
1142
1143    /// Stop. `force` kills instead of a clean shutdown; `timeout` bounds the
1144    /// clean shutdown.
1145    pub fn stop(&self, force: bool, timeout: Duration) -> Result<()> {
1146        if self.info()?.status.eq_ignore_ascii_case("stopped") {
1147            return Ok(());
1148        }
1149        stop_instance(&self.client, &self.name, force, timeout)
1150    }
1151
1152    pub fn restart(&self) -> Result<()> {
1153        self.stop(false, Duration::from_secs(30))?;
1154        self.start()
1155    }
1156
1157    /// Run this handle's readiness checks.
1158    pub fn wait_ready(&self) -> Result<()> {
1159        wait_ready(
1160            &self.client,
1161            &self.name,
1162            &self.ready,
1163            self.ready_timeout,
1164            &self.exec_defaults,
1165        )
1166    }
1167
1168    /// Run a command and capture its output. `argv[0]` is the program; nothing is
1169    /// joined into a shell string.
1170    pub fn exec<I, S>(&self, argv: I) -> Result<ExecOutput>
1171    where
1172        I: IntoIterator<Item = S>,
1173        S: Into<String>,
1174    {
1175        self.exec_with(argv, ExecOptions::default())
1176    }
1177
1178    pub fn exec_with<I, S>(&self, argv: I, opts: ExecOptions) -> Result<ExecOutput>
1179    where
1180        I: IntoIterator<Item = S>,
1181        S: Into<String>,
1182    {
1183        self.exec_stream(argv, opts)?.collect_output()
1184    }
1185
1186    /// Start a command and stream its output as it is produced.
1187    pub fn exec_stream<I, S>(&self, argv: I, opts: ExecOptions) -> Result<ExecStream>
1188    where
1189        I: IntoIterator<Item = S>,
1190        S: Into<String>,
1191    {
1192        let argv: Vec<String> = argv.into_iter().map(Into::into).collect();
1193        let req = exec::build_request(&self.client, &self.name, &argv, &self.exec_defaults, &opts)?;
1194        exec::start_with_timeout(
1195            &self.client,
1196            &self.name,
1197            req,
1198            opts.stdin.clone(),
1199            opts.timeout,
1200        )
1201    }
1202
1203    /// Run attached to this process's stdio (and terminal, with `opts.tty`),
1204    /// forwarding signals. Returns the exit code.
1205    pub fn attach<I, S>(&self, argv: I, opts: ExecOptions) -> Result<i32>
1206    where
1207        I: IntoIterator<Item = S>,
1208        S: Into<String>,
1209    {
1210        let argv: Vec<String> = argv.into_iter().map(Into::into).collect();
1211        let req = exec::build_request(&self.client, &self.name, &argv, &self.exec_defaults, &opts)?;
1212        exec::attach(
1213            &self.client,
1214            &self.name,
1215            req,
1216            opts.stdin.clone(),
1217            opts.timeout,
1218        )
1219    }
1220
1221    /// Add (or correct) a proxy device. Host-bound ports with `search` step past
1222    /// taken ports. Returns the listen address in use. A correct device is left
1223    /// untouched.
1224    pub fn add_port(&self, port: &PortSpec) -> Result<String> {
1225        // Same rules as the spec path: a VM gets NAT mode and no bind=guest.
1226        let instance_type = if self.info()?.instance_type == "virtual-machine" {
1227            crate::spec::InstanceType::VirtualMachine
1228        } else {
1229            crate::spec::InstanceType::Container
1230        };
1231        let spec = SandboxSpec {
1232            name: Some(self.name.clone()),
1233            image: "unused".into(),
1234            instance_type,
1235            ports: vec![port.clone()],
1236            ..Default::default()
1237        };
1238        let host = HostFacts {
1239            pools: vec!["unused".into()],
1240            ..Default::default()
1241        };
1242        let d = plan::resolve(&spec, &VolumeDefs::new(), &host, Path::new("/"))?;
1243        let (dname, want): (&String, &DesiredDevice) = d
1244            .devices
1245            .iter()
1246            .find(|(k, _)| k.as_str() != "root")
1247            .expect("one port");
1248        let _lock = NameLock::acquire(
1249            self.client.project_name(),
1250            &self.name,
1251            Duration::from_secs(900),
1252            &mut |_| {},
1253        )?;
1254        let info = self.info()?;
1255        if let Some(have) = info.devices.get(dname) {
1256            if plan::device_matches(want, have) {
1257                return Ok(have.get("listen").cloned().unwrap_or_default());
1258            }
1259            self.remove_device_unlocked(dname)?;
1260        }
1261        match want.search {
1262            Some(n) => add_port_searching(&self.client, &self.name, dname, &want.props, n),
1263            None => {
1264                let props = want.props.clone();
1265                update_instance(
1266                    &self.client,
1267                    &self.name,
1268                    &format!("add port {dname}"),
1269                    &mut |_, devices| {
1270                        devices.insert(dname.clone(), json!(props));
1271                        Ok(())
1272                    },
1273                )?;
1274                Ok(want.props["listen"].clone())
1275            }
1276        }
1277    }
1278
1279    /// Remove an instance-local device. `Ok(false)` if it was not there.
1280    pub fn remove_device(&self, device: &str) -> Result<bool> {
1281        let _lock = NameLock::acquire(
1282            self.client.project_name(),
1283            &self.name,
1284            Duration::from_secs(900),
1285            &mut |_| {},
1286        )?;
1287        self.remove_device_unlocked(device)
1288    }
1289
1290    fn remove_device_unlocked(&self, device: &str) -> Result<bool> {
1291        if device == "root" {
1292            return Err(Error::invalid("refusing to remove the root disk"));
1293        }
1294        let mut found = false;
1295        update_instance(
1296            &self.client,
1297            &self.name,
1298            &format!("remove device {device}"),
1299            &mut |_, devices| {
1300                found = devices.remove(device).is_some();
1301                Ok(())
1302            },
1303        )?;
1304        Ok(found)
1305    }
1306}
1307
1308/// One instance `prune` considered.
1309#[derive(Debug, Clone, Serialize)]
1310pub struct PruneItem {
1311    pub name: String,
1312    pub path: String,
1313    pub deleted: bool,
1314}
1315
1316/// Delete instances whose `label` value is a host path that no longer exists.
1317/// Instances without the label, or whose path exists, are never touched.
1318/// With `dry_run`, nothing is deleted.
1319pub fn prune_missing_path(
1320    client: &Client,
1321    label: &str,
1322    dry_run: bool,
1323    report: Reporter<'_>,
1324) -> Result<Vec<PruneItem>> {
1325    let mut out = Vec::new();
1326    for i in Sandbox::list_with(client, &[LabelFilter::parse(label)])? {
1327        let Some(path) = i.labels.get(label) else {
1328            continue;
1329        };
1330        if path.is_empty() || !path.starts_with('/') || Path::new(path).exists() {
1331            continue;
1332        }
1333        if dry_run {
1334            report(&format!(
1335                "would delete {}  ({label}={path} is gone)",
1336                i.name
1337            ));
1338        } else {
1339            report(&format!("deleting {}  ({label}={path} is gone)", i.name));
1340            force_delete(client, &i.name)?;
1341        }
1342        out.push(PruneItem {
1343            name: i.name,
1344            path: path.clone(),
1345            deleted: !dry_run,
1346        });
1347    }
1348    Ok(out)
1349}
1350
1351#[cfg(test)]
1352mod tests;