1use std::collections::BTreeMap;
4use std::path::Path;
5use std::time::{Duration, Instant};
6
7use serde::Serialize;
8use serde_json::{Value, json};
9
10use crate::client::{Client, encode_segment};
11use crate::error::{Error, Result};
12use crate::exec::{self, ExecOptions, ExecOutput, ExecStream, Stdin};
13use crate::idmap::SubIds;
14use crate::lock::NameLock;
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
22pub const CREATE_TOKEN_KEY: &str = "user.isb.create-token";
25
26pub type Reporter<'a> = &'a mut dyn FnMut(&str);
28
29#[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 pub labels: BTreeMap<String, String>,
38 pub config: BTreeMap<String, String>,
39 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 .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 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#[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#[derive(Debug, Clone, Default, Serialize)]
113pub struct ApplyReport {
114 pub name: String,
115 pub created: bool,
116 pub applied: Vec<Action>,
117 pub ports: BTreeMap<String, String>,
119 pub restart_needed: Vec<String>,
121}
122
123#[derive(Debug, Clone, Copy)]
125pub struct EnsureOptions {
126 pub diff: DiffOptions,
127 pub wait_ready: bool,
129 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
147pub 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
181fn 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
191pub 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
218fn 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
234pub 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 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
289fn 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
300pub 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 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 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 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#[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 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 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
594pub 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 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 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
683fn 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 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
740pub 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#[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 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
884pub 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 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#[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 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 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 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 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 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 pub fn list(client: &Client) -> Result<Vec<SandboxInfo>> {
1072 Self::list_with(client, &[])
1073 }
1074
1075 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 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 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 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 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 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 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 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 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 pub fn add_port(&self, port: &PortSpec) -> Result<String> {
1225 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 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#[derive(Debug, Clone, Serialize)]
1310pub struct PruneItem {
1311 pub name: String,
1312 pub path: String,
1313 pub deleted: bool,
1314}
1315
1316pub 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;