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;
15pub mod images;
16use crate::plan::{
17    self, Action, Actual, Desired, DesiredDevice, DiffOptions, HostFacts, Props, SandboxPlan,
18    VolumeDefs, split_addr,
19};
20use crate::spec::{ExecDefaults, PortSpec, ReadyCheck, SandboxSpec};
21
22/// Config key recording which isb call created an instance. A half-created
23/// instance is only ever cleaned up by the call whose token it carries.
24pub const CREATE_TOKEN_KEY: &str = "user.isb.create-token";
25
26/// Progress callback: one human-readable line per step.
27pub type Reporter<'a> = &'a mut dyn FnMut(&str);
28
29/// A summary of an instance, as listed.
30#[derive(Debug, Clone, Serialize)]
31pub struct SandboxInfo {
32    pub name: String,
33    pub status: String,
34    #[serde(rename = "type")]
35    pub instance_type: String,
36    /// `user.*` config keys with the prefix stripped.
37    pub labels: BTreeMap<String, String>,
38    pub config: BTreeMap<String, String>,
39    /// Instance-local devices.
40    pub devices: BTreeMap<String, Props>,
41    pub profiles: Vec<String>,
42    pub created_at: String,
43    pub description: String,
44}
45
46impl SandboxInfo {
47    pub fn from_api(v: &Value) -> SandboxInfo {
48        let a = Actual::from_api(v);
49        let labels = a
50            .config
51            .iter()
52            .filter_map(|(k, v)| k.strip_prefix("user.").map(|k| (k.to_string(), v.clone())))
53            // isb's own bookkeeping (user.isb.create-token) is not a label.
54            .filter(|(k, _): &(String, String)| !k.starts_with("isb."))
55            .collect();
56        SandboxInfo {
57            name: v.get("name").and_then(Value::as_str).unwrap_or("").into(),
58            status: a.status.clone(),
59            instance_type: a.instance_type.clone(),
60            labels,
61            config: a.config,
62            devices: a.devices,
63            profiles: a.profiles,
64            created_at: v
65                .get("created_at")
66                .and_then(Value::as_str)
67                .unwrap_or("")
68                .into(),
69            description: v
70                .get("description")
71                .and_then(Value::as_str)
72                .unwrap_or("")
73                .into(),
74        }
75    }
76
77    /// Whether this instance matches every filter.
78    pub fn matches(&self, filters: &[LabelFilter]) -> bool {
79        filters
80            .iter()
81            .all(|f| match (&f.value, self.labels.get(&f.key)) {
82                (_, None) => false,
83                (None, Some(_)) => true,
84                (Some(want), Some(have)) => want == have,
85            })
86    }
87}
88
89/// `key` (present) or `key=value` (equal).
90#[derive(Debug, Clone, PartialEq, Eq)]
91pub struct LabelFilter {
92    pub key: String,
93    pub value: Option<String>,
94}
95
96impl LabelFilter {
97    pub fn parse(s: &str) -> LabelFilter {
98        match s.split_once('=') {
99            Some((k, v)) => LabelFilter {
100                key: k.into(),
101                value: Some(v.into()),
102            },
103            None => LabelFilter {
104                key: s.into(),
105                value: None,
106            },
107        }
108    }
109}
110
111/// What `apply` did.
112#[derive(Debug, Clone, Default, Serialize)]
113pub struct ApplyReport {
114    pub name: String,
115    pub created: bool,
116    pub applied: Vec<Action>,
117    /// Device name -> listen address chosen for searched ports.
118    pub ports: BTreeMap<String, String>,
119    /// Config keys changed that only take effect after a restart.
120    pub restart_needed: Vec<String>,
121}
122
123/// Options for ensure / up.
124#[derive(Debug, Clone, Copy)]
125pub struct EnsureOptions {
126    pub diff: DiffOptions,
127    /// Run readiness checks after applying.
128    pub wait_ready: bool,
129    /// How long to wait for another isb holding this sandbox's lock.
130    pub lock_wait: Duration,
131}
132
133impl Default for EnsureOptions {
134    fn default() -> Self {
135        EnsureOptions {
136            diff: DiffOptions::default(),
137            wait_ready: true,
138            lock_wait: Duration::from_secs(900),
139        }
140    }
141}
142
143fn inst_path(name: &str) -> String {
144    format!("/1.0/instances/{}", encode_segment(name))
145}
146
147/// Gather host facts (pools, subids, path map).
148pub fn host_facts(client: &Client) -> Result<HostFacts> {
149    let pools = client
150        .get("/1.0/storage-pools")?
151        .as_array()
152        .map(|a| {
153            a.iter()
154                .filter_map(|u| u.as_str())
155                .filter_map(|u| u.rsplit('/').next())
156                .map(|s| s.split('?').next().unwrap_or(s).to_string())
157                .collect()
158        })
159        .unwrap_or_default();
160    let info = client.server_info()?;
161    let ext = |n: &str| {
162        info["api_extensions"]
163            .as_array()
164            .is_some_and(|a| a.iter().any(|e| e == n))
165    };
166    Ok(HostFacts {
167        subids: SubIds::read_host(),
168        pools,
169        path_map: HostFacts::detect_path_map(),
170        initial_copy: ext("disk_initial_copy"),
171        initial_owner: ext("storage_initial_owner"),
172        incus_version: client.server_version()?,
173        invoking_ids: crate::idmap::invoking_ids(),
174        shared_root: shared_root(),
175        org: crate::org::OrgId::from_incus_project(client.project_name()),
176        registry: crate::registry::info(client)?.map(|i| i.addr),
177        project: client.project_name().to_string(),
178    })
179}
180
181/// On macOS incusd runs in the `isb machine` VM, which sees only `$HOME`.
182fn shared_root() -> Option<String> {
183    if !cfg!(target_os = "macos") {
184        return None;
185    }
186    let home = std::env::var_os("HOME").filter(|h| !h.is_empty())?;
187    let p = std::path::PathBuf::from(home);
188    Some(p.canonicalize().unwrap_or(p).to_string_lossy().into_owned())
189}
190
191/// Resolve a spec against this host. Relative bind paths anchor at `base`.
192pub fn resolve(
193    client: &Client,
194    spec: &SandboxSpec,
195    defs: &VolumeDefs,
196    base: &Path,
197) -> Result<Desired> {
198    let host = host_facts(client)?;
199    plan::resolve(spec, defs, &host, base)
200}
201
202fn get_actual(client: &Client, name: &str) -> Result<Option<Actual>> {
203    Ok(client
204        .get_opt(&inst_path(name))?
205        .map(|v| Actual::from_api(&v)))
206}
207
208fn volume_exists(client: &Client, pool: &str, name: &str) -> Result<bool> {
209    Ok(client
210        .get_opt(&format!(
211            "/1.0/storage-pools/{}/volumes/custom/{}",
212            encode_segment(pool),
213            encode_segment(name)
214        ))?
215        .is_some())
216}
217
218/// The fingerprint of a local image, by alias or fingerprint (prefix).
219fn local_image(client: &Client, alias: &str) -> Result<Option<String>> {
220    if let Some(a) = client.get_opt(&format!("/1.0/images/aliases/{}", encode_segment(alias)))? {
221        return Ok(a.get("target").and_then(Value::as_str).map(String::from));
222    }
223    if alias.len() >= 12 && alias.chars().all(|c| c.is_ascii_hexdigit()) {
224        if let Some(i) = client.get_opt(&format!("/1.0/images/{}", encode_segment(alias)))? {
225            return Ok(i
226                .get("fingerprint")
227                .and_then(Value::as_str)
228                .map(String::from));
229        }
230    }
231    Ok(None)
232}
233
234/// Compute the plan for one resolved sandbox.
235pub fn plan_desired(client: &Client, desired: &Desired, opts: DiffOptions) -> Result<SandboxPlan> {
236    let actual = get_actual(client, &desired.name)?;
237    if actual.is_none()
238        && desired.image.server.is_none()
239        && local_image(client, &desired.image.alias)?.is_none()
240    {
241        return Err(images::missing_on(client, &desired.image.alias));
242    }
243    let mut missing = Vec::new();
244    for v in &desired.volumes {
245        if !volume_exists(client, &v.pool, &v.name)? {
246            missing.push((v.pool.clone(), v.name.clone()));
247        }
248    }
249    let mut plan = plan::diff(desired, actual.as_ref(), &missing, opts)?;
250    // The image is fixed at creation; say so when the local image has moved on
251    // (e.g. dev-base was republished), so a recreate is a visible choice.
252    if let (Some(a), None) = (&actual, &desired.image.server) {
253        if let (Some(built), Some(now)) = (
254            a.config.get("volatile.base_image"),
255            local_image(client, &desired.image.alias)?,
256        ) {
257            if *built != now {
258                plan.actions.push(Action::Note {
259                    message: format!(
260                        "image {} is now {} but this instance was built from {}; recreate to pick it up",
261                        desired.image.alias,
262                        &now[..12.min(now.len())],
263                        &built[..12.min(built.len())]
264                    ),
265                });
266            }
267        }
268    }
269    Ok(plan)
270}
271
272fn random_token() -> String {
273    let mut b = [0u8; 12];
274    if let Ok(mut f) = std::fs::File::open("/dev/urandom") {
275        use std::io::Read;
276        let _ = f.read_exact(&mut b);
277    }
278    let t = std::time::SystemTime::now()
279        .duration_since(std::time::UNIX_EPOCH)
280        .map(|d| d.as_nanos())
281        .unwrap_or(0);
282    format!(
283        "{}{:x}",
284        b.iter().map(|x| format!("{x:02x}")).collect::<String>(),
285        t & 0xffff
286    )
287}
288
289/// Retry a step once after a deadline.
290fn retry_once<T>(report: &mut dyn FnMut(&str), mut f: impl FnMut() -> Result<T>) -> Result<T> {
291    match f() {
292        Err(e) if e.is_timeout() => {
293            report(&format!("{e}; retrying once"));
294            f()
295        }
296        r => r,
297    }
298}
299
300/// Apply a plan. Changes only what the plan says; a correct device is never touched.
301pub fn apply(
302    client: &Client,
303    desired: &Desired,
304    plan: &SandboxPlan,
305    report: Reporter<'_>,
306) -> Result<ApplyReport> {
307    let name = &desired.name;
308    let mut out = ApplyReport {
309        name: name.clone(),
310        ..Default::default()
311    };
312    let mut pending: Vec<&Action> = Vec::new();
313    if let Some(p) = &desired.egress {
314        crate::egress::before_apply(client, p, report)?;
315    }
316    for action in &plan.actions {
317        match action {
318            Action::SetConfig { .. }
319            | Action::AddDevice { .. }
320            | Action::ReplaceDevice { .. }
321            | Action::RemoveDevice { .. } => {
322                pending.push(action);
323                continue;
324            }
325            _ => {}
326        }
327        // Batched config/device changes land before whatever comes next (a start,
328        // most importantly, so a stopped instance boots with the right devices).
329        flush_updates(client, desired, &mut pending, &mut out, report)?;
330        match action {
331            Action::Note { message } => report(&format!("{name}: note: {message}")),
332            Action::CreateVolume {
333                pool,
334                volume,
335                config,
336            } => {
337                report(&format!("{name}: creating volume {volume} on {pool}"));
338                retry_once(report, || {
339                    crate::volume::ensure(client, pool, volume, config).map(|_| ())
340                })?;
341            }
342            Action::CreateInstance { .. } => {
343                report(&format!("{name}: creating from {}", desired.image.spec));
344                create_instance(client, desired, report)?;
345                out.created = true;
346            }
347            Action::StartInstance => {
348                if let (true, Some(h)) = (out.created, &desired.before_start) {
349                    (h.0)(client, name)?;
350                }
351                report(&format!("{name}: starting"));
352                retry_once(report, || start_instance(client, name))?;
353            }
354            Action::AddPort {
355                device,
356                props,
357                search,
358            } => {
359                let listen = add_port_searching(client, name, device, props, *search)?;
360                report(&format!("{name}: port {device} listening on {listen}"));
361                out.ports.insert(device.clone(), listen);
362            }
363            Action::FixOwner {
364                path,
365                owner,
366                mode,
367                fresh_only,
368            } => {
369                if desired.instance_type == crate::spec::InstanceType::VirtualMachine {
370                    // In-guest work needs the VM's agent, which starts after boot.
371                    wait_ready(
372                        client,
373                        name,
374                        &[ReadyCheck::Agent],
375                        desired.ready_timeout,
376                        &desired.exec,
377                    )?;
378                }
379                let fix = crate::owner::Fix {
380                    path,
381                    owner: owner.as_deref(),
382                    mode: mode
383                        .as_deref()
384                        .map(crate::owner::parse_mode)
385                        .transpose()
386                        .map_err(Error::invalid)?,
387                    fresh_only: *fresh_only,
388                };
389                if crate::owner::apply(client, name, &fix)? {
390                    report(&format!("{name}: {action}"));
391                } else {
392                    report(&format!(
393                        "{name}: {path} came seeded from the image; owner left as is"
394                    ));
395                }
396            }
397            _ => unreachable!("batched above"),
398        }
399        out.applied.push(action.clone());
400    }
401    flush_updates(client, desired, &mut pending, &mut out, report)?;
402    Ok(out)
403}
404
405fn flush_updates(
406    client: &Client,
407    desired: &Desired,
408    pending: &mut Vec<&Action>,
409    out: &mut ApplyReport,
410    report: &mut dyn FnMut(&str),
411) -> Result<()> {
412    if pending.is_empty() {
413        return Ok(());
414    }
415    let name = desired.name.as_str();
416    for a in pending.iter() {
417        report(&format!("{name}: {a}"));
418    }
419    let actions: Vec<Action> = pending.iter().map(|a| (*a).clone()).collect();
420    retry_once(report, || {
421        update_instance(
422            client,
423            name,
424            &format!("update {name}"),
425            &mut |config, devices| {
426                for a in &actions {
427                    match a {
428                        Action::SetConfig {
429                            key, to, secret, ..
430                        } => {
431                            // A secret's action carries a placeholder.
432                            let v = match (secret, desired.config.get(key)) {
433                                (true, Some(v)) => v,
434                                _ => to,
435                            };
436                            config.insert(key.clone(), json!(v));
437                        }
438                        Action::AddDevice { device, props } => {
439                            devices.insert(device.clone(), json!(props));
440                        }
441                        Action::ReplaceDevice {
442                            device,
443                            replaces,
444                            to,
445                            ..
446                        } => {
447                            devices.remove(replaces);
448                            devices.insert(device.clone(), json!(to));
449                        }
450                        Action::RemoveDevice { device, .. } => {
451                            devices.remove(device);
452                        }
453                        _ => {}
454                    }
455                }
456                Ok(())
457            },
458        )
459    })?;
460    for a in pending.drain(..) {
461        if let Action::SetConfig {
462            key, restart: true, ..
463        } = a
464        {
465            out.restart_needed.push(key.clone());
466        }
467        out.applied.push(a.clone());
468    }
469    Ok(())
470}
471
472type Obj = serde_json::Map<String, Value>;
473
474/// Read-modify-write of an instance under `If-Match`, so a concurrent change by
475/// another tool is never silently overwritten (a 412 re-reads and retries).
476#[doc(hidden)]
477pub fn update_instance(
478    client: &Client,
479    name: &str,
480    step: &str,
481    modify: &mut dyn FnMut(&mut Obj, &mut Obj) -> Result<()>,
482) -> Result<()> {
483    let path = inst_path(name);
484    for attempt in 0..5 {
485        let (inst, etag) = client.get_etag(&path)?;
486        let mut config = inst
487            .get("config")
488            .and_then(Value::as_object)
489            .cloned()
490            .unwrap_or_default();
491        let mut devices = inst
492            .get("devices")
493            .and_then(Value::as_object)
494            .cloned()
495            .unwrap_or_default();
496        modify(&mut config, &mut devices)?;
497        let body = json!({
498            "architecture": inst.get("architecture"),
499            "config": config,
500            "devices": devices,
501            "ephemeral": inst.get("ephemeral"),
502            "profiles": inst.get("profiles"),
503            "stateful": inst.get("stateful"),
504            "description": inst.get("description"),
505        });
506        match client.mutate_if_match(
507            "PUT",
508            &path,
509            &body,
510            etag.as_deref(),
511            step,
512            client.timeouts.other,
513        ) {
514            Err(Error::Api { status: 412, .. }) if attempt < 4 => continue,
515            r => return r.map(|_| ()),
516        }
517    }
518    unreachable!()
519}
520
521fn create_instance(client: &Client, desired: &Desired, report: &mut dyn FnMut(&str)) -> Result<()> {
522    let fingerprint = match &desired.image.server {
523        None => Some(
524            local_image(client, &desired.image.alias)?
525                .ok_or_else(|| images::missing_on(client, &desired.image.alias))?,
526        ),
527        Some(_) => None,
528    };
529    let step = format!("create instance {}", desired.name);
530    let mut last_err = None;
531    for attempt in 0..2 {
532        let token = random_token();
533        let mut config = desired.config.clone();
534        config.insert(CREATE_TOKEN_KEY.into(), token.clone());
535        let devices: BTreeMap<&String, &Props> = desired
536            .devices
537            .iter()
538            .filter(|(_, d)| d.search.is_none())
539            .map(|(k, d)| (k, &d.props))
540            .collect();
541        let body = json!({
542            "name": desired.name,
543            "type": desired.instance_type.as_api(),
544            "source": desired.image.to_api(fingerprint.as_deref()),
545            "config": config,
546            "devices": devices,
547            "profiles": desired.profiles,
548        });
549        let body = crate::org::disk::sized_root(client, body);
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;