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