1macro_rules! tool {
18 ($r:expr, $ctx:expr, $name:expr, $title:expr, $desc:expr, $schema:expr, $ann:expr, $f:expr) => {{
19 let ctx = $ctx.clone();
20 let f = $f;
21 $r.register(
22 Tool::new($name, $desc, $schema, move |a, c| f(&ctx, a, c))
23 .title($title)
24 .annotations($ann.clone()),
25 )?;
26 }};
27}
28
29mod accounts;
30pub mod apps;
31pub mod audit;
32mod authorize;
33pub mod builds;
34pub mod data;
35mod default_org;
36mod dns;
37mod egress;
38mod guide;
39mod host_monitor;
40mod kube;
41mod monitors;
42mod notify;
43mod orgs;
44pub mod policy;
45pub mod previews;
46mod registry;
47mod secret_hooks;
48pub mod secrets;
49mod ssh;
50mod stack_deploy;
51pub mod superadmin;
52pub mod templates;
53mod terminal;
54mod tools;
55pub mod volumes;
56pub mod workspaces;
57
58use std::collections::BTreeMap;
59use std::path::PathBuf;
60use std::sync::Arc;
61use std::time::{Duration, Instant};
62
63use serde::Deserialize;
64use serde::de::DeserializeOwned;
65use serde_json::{Value, json};
66
67use crate::auth::AuthConfig;
68use crate::auth::AuthStore;
69use crate::auth::http::{ApiConfig, AuthApi};
70use crate::client::Client;
71use crate::error::{Error, Result};
72use crate::exec::{ExecOptions, Stdin};
73use crate::sandbox::{EnsureOptions, Sandbox, SandboxInfo};
74use crate::server::{AccessValidator, Caller, Listener, Registry, Tool, ToolPolicy};
75use crate::spec::SandboxSpec;
76use crate::stack::{Controller, Store, now_secs};
77use authorize::{
78 CROSS_ORG_READS, PLATFORM_TOOLS, arg_org, authorize_class, event_visible, read_orgs,
79 tool_listed, visible_orgs,
80};
81use policy::RemotePolicy;
82use stack_deploy::{DeployArgs, How, deploy, stack_deploy};
83#[cfg(test)]
84use stack_deploy::{reused_note, stack_owner};
85use tools::Ann;
86
87pub const LABEL_OWNER: &str = "isb.owner";
89
90#[derive(Debug, Clone)]
92pub struct ServeConfig {
93 pub listen: Vec<String>,
97 pub socket: PathBuf,
98 pub access: Option<(String, String)>,
100 pub allow_unauthenticated: bool,
102 pub remote_tools: ToolPolicy,
104 pub policy: RemotePolicy,
105 pub state_dir: PathBuf,
106 pub interval: Duration,
107 pub keys: crate::secrets::KeySources,
109 pub secrets_config: PathBuf,
111 pub auth: AuthConfig,
113 pub public_url: Option<String>,
116 pub oauth: crate::auth::oauth::OAuthSettings,
118 pub open_signup: bool,
120 pub ingress: Option<crate::ingress::IngressConfig>,
122 pub workspace_mcp_port: u16,
125 pub workspace_pool: Option<String>,
128 pub workspace_home_root: Option<PathBuf>,
131 pub preview_domain: Option<workspaces::PreviewBase>,
133 pub skills: Vec<crate::server::Skill>,
135 pub audit_retention: Duration,
137 pub audit_all: bool,
139 pub history_retention: Duration,
141 pub history_max_rows: i64,
142 pub superadmin_tailnet: Option<crate::server::tailnet::AllowList>,
145 pub superadmin_access: Option<superadmin::AccessAllowList>,
148 pub dev_superadmin: Option<String>,
151 pub heartbeat: Option<crate::monitor::heartbeat::Heartbeat>,
153 pub egress_pins: Vec<String>,
156 pub egress_ca: Vec<PathBuf>,
158}
159
160fn auth_routes(
164 cfg: &ServeConfig,
165 store: Arc<AuthStore>,
166 secrets: &Arc<crate::secrets::Secrets>,
167 log: &Arc<crate::audit::AuditLog>,
168 gate: Arc<superadmin::Gate>,
169 orgs: crate::auth::ops::OrgsFn,
170) -> Result<crate::server::Routes> {
171 use crate::auth::oauth::SecretFn;
172 let default_org = crate::org::OrgId::default_org();
173 let lookup = |name: &str| -> Option<SecretFn> {
174 secrets.inspect(&default_org, name).ok()?;
175 let (s, org, name) = (secrets.clone(), default_org.clone(), name.to_string());
176 Some(Arc::new(move || {
177 let (v, _) = s.get(&org, &name).map_err(|e| e.to_string())?;
178 String::from_utf8(v)
179 .map(|v| v.trim().to_string())
180 .map_err(|_| "the secret is not UTF-8".to_string())
181 }))
182 };
183 let (providers, notes) = cfg.oauth.providers(&lookup);
184 for n in notes {
185 eprintln!("isb serve: {n}");
186 }
187 let path = crate::auth::db_path(&cfg.state_dir);
188 let agent_ways = gate.agent_ways().with_public_url(cfg.public_url.clone());
189 let api = AuthApi::new(
190 store.clone(),
191 ApiConfig {
192 agent: Some(gate.agent_fn()),
193 agent_ways,
194 public_url: cfg.public_url.clone(),
195 notifier: None,
196 setup_token_file: Some(cfg.state_dir.join("setup-token")),
197 edge: Some(gate.edge_fn()),
198 providers,
199 open_signup: cfg.open_signup,
200 audit: Some(log.clone()),
201 superadmin: Some(Arc::new(move |r: &crate::server::http::Request| match gate
202 .resolve(r, None)
203 {
204 superadmin::Resolved::Superadmin(s) => Some(s),
205 _ => None,
206 })),
207 orgs: Some(orgs),
208 },
209 )?;
210 eprintln!("isb serve: identity store {}", path.display());
211 if !crate::web::BUILT {
212 eprintln!(
213 "isb serve: this binary was built without the web UI (a placeholder page is served)"
214 );
215 }
216 let (auth, web) = (Arc::new(api).router(), crate::web::routes());
219 let tail = audit::stream_route(log.clone(), store);
220 Ok(Arc::new(move |r| {
221 auth(r).or_else(|| tail(r)).or_else(|| web(r))
222 }))
223}
224
225struct Daemon {
226 client: Client,
227 ctl: Controller,
228 policy: RemotePolicy,
229 state_dir: PathBuf,
230 secrets: Arc<crate::secrets::Secrets>,
231 apps: crate::app::Apps,
232 ingress: Option<Arc<crate::ingress::Manager>>,
233 users: Arc<AuthStore>,
235 notifier: crate::notify::Notifier,
236 monitors: crate::monitor::Monitors,
237 history: crate::metrics_history::History,
238 data: data::Ctx,
240 volumes: crate::volume_backup::VolumeBackups,
242 audit: Arc<crate::audit::AuditLog>,
243 gate: Arc<superadmin::Gate>,
245 host: Value,
246 catalogs: Arc<crate::template::catalog::Catalogs>,
248 workspaces: Arc<workspaces::Workspaces>,
251 public_url: Option<String>,
253 egress: Arc<isb_egress::Manager>,
255 meta: crate::stack::deployments::StackMeta,
257}
258
259#[expect(
261 clippy::too_many_lines,
262 reason = "predates the lint ratchet; split it when next changed"
263)]
264pub fn serve(client: Client, cfg: ServeConfig) -> Result<()> {
265 client
266 .server_info()
267 .map_err(|e| Error::invalid(format!("isb serve needs incusd: {e}")))?;
268 default_org::warn_old_incus(&client);
269 let store = Store::open(&cfg.state_dir)?;
270 dns::open_dns_path(&cfg.state_dir);
271 default_org::ensure(&client, &store);
274 let secrets_config = crate::secrets::SecretsConfig::load(&cfg.secrets_config)?;
275 let opened = crate::secrets::Secrets::open(&cfg.state_dir, &cfg.keys, &secrets_config)?;
276 for n in &opened.notes {
277 eprintln!("isb serve: {n}");
278 }
279 let secrets = Arc::new(with_external_drivers(opened.secrets, &cfg.state_dir)?);
280 for r in crate::stack::migrate::run(&store, &secrets, Some(&client)) {
283 match r {
284 Ok(m) => eprintln!("isb serve: {m}"),
285 Err(e) => eprintln!("isb serve: WARNING: {e}"),
286 }
287 }
288 let db = crate::auth::db_path(&cfg.state_dir);
291 let users = Arc::new(
292 AuthStore::open_with(&db, cfg.auth.clone())
293 .map_err(|e| Error::invalid(format!("open {}: {e}", db.display())))?,
294 );
295 let audit_db = crate::audit::db_path(&cfg.state_dir);
296 let audit_log = Arc::new(
297 crate::audit::AuditLog::open(&audit_db, cfg.audit_retention)?
298 .with_history_limits(cfg.history_retention, cfg.history_max_rows),
299 );
300 let recorder = crate::history::Recorder::start(audit_log.clone());
303 let stop_history = Arc::new(std::sync::atomic::AtomicBool::new(false));
304 history_start(&audit_log, &recorder, &client, &stop_history);
305 eprintln!(
306 "isb serve: audit log {} (kept {} days{})",
307 audit_db.display(),
308 cfg.audit_retention.as_secs() / 86400,
309 if cfg.audit_all {
310 ", reads included"
311 } else {
312 ""
313 }
314 );
315 let access = match &cfg.access {
316 Some((team, aud)) => Some(Arc::new(AccessValidator::new(team, aud)?)),
317 None => None,
318 };
319 let gate = Arc::new(superadmin::gate(&cfg, users.clone(), access.clone())?);
320 let auth = if cfg.listen.is_empty() {
321 None
322 } else {
323 Some(auth_routes(
324 &cfg,
325 users.clone(),
326 &secrets,
327 &audit_log,
328 gate.clone(),
329 orgs::existing_fn(client.clone()),
330 )?)
331 };
332 match crate::registry::Registry::open(&client, Some(&cfg.state_dir)) {
335 Ok(Some(r)) => {
336 eprintln!("isb serve: local registry {}", r.info().url());
337 crate::registry::install(Arc::new(r));
338 }
339 Ok(None) => {}
340 Err(e) => eprintln!("isb serve: WARNING: local registry: {e}"),
341 }
342 let ingress = match &cfg.ingress {
345 Some(ic) => Some(crate::ingress::Manager::new(
346 ic.clone(),
347 client.clone(),
348 secrets.clone(),
349 &cfg.state_dir,
350 )?),
351 None => None,
352 };
353 let observer = ingress
354 .clone()
355 .map(|m| m as Arc<dyn crate::stack::controller::Observer>);
356 let ctl = Controller::start_with(
357 client.clone(),
358 store,
359 cfg.interval,
360 secrets.clone(),
361 observer,
362 )?;
363 if let Some(m) = &ingress {
364 m.start(ctl.clone())?;
365 }
366 let apps = crate::app::Apps::new(&cfg.state_dir, client.clone(), ctl.clone(), secrets.clone());
367 let stacks: Vec<(crate::org::OrgId, String)> = ctl
371 .definitions()
372 .iter()
373 .map(|d| (d.org.clone(), d.name.clone()))
374 .collect();
375 for r in apps.adopt_compose_stacks(&stacks) {
376 match r {
377 Ok(m) => eprintln!("isb serve: {m}"),
378 Err(e) => eprintln!("isb serve: WARNING: {e}"),
379 }
380 }
381 ctl.set_dns_scope(apps.scope_fn());
383 let ra = apps.clone();
385 let resolve: crate::notify::Resolve = Arc::new(move |org, stack, service| {
386 let a = ra.get(org, service).ok()?;
387 let own = a.spec.stack().ok()? == stack;
389 let preview = stack.starts_with(&format!("{}-", a.spec.project))
390 && crate::app::preview::is_pr_suffix(stack);
391 (own || preview).then_some(a.spec.project)
392 });
393 let notifier = crate::notify::Notifier::new(&cfg.state_dir, secrets.clone(), resolve)?;
394 notifier.start(ctl.clone());
395 let monitors = monitors::start(&cfg, &apps, &secrets, ¬ifier);
396 ctl.set_event_sink(recorder.controller_sink());
398 let history = crate::metrics_history::History::new(&cfg.state_dir);
400 ctl.set_metrics_sink(history.start());
401 apps.start_preview_upkeep();
403 let jobs = crate::jobs::Jobs::new(&cfg.state_dir, apps.clone());
405 let backups = crate::backup::Backups::new(&cfg.state_dir, apps.clone());
406 let volumes =
407 crate::volume_backup::VolumeBackups::new(&cfg.state_dir, apps.clone(), backups.clone());
408 let scheduler = crate::jobs::Scheduler::start(vec![
409 Arc::new(jobs.clone()) as Arc<dyn crate::jobs::Scheduled>,
410 Arc::new(backups.clone()),
411 Arc::new(volumes.clone()),
412 ]);
413 jobs.set_scheduler(scheduler.clone());
414 backups.set_scheduler(scheduler.clone());
415 volumes.set_scheduler(scheduler.clone());
416 if let Some(ic) = &cfg.ingress {
417 if ic.tunnel_port == cfg.workspace_mcp_port {
418 return Err(Error::invalid(format!(
419 "--workspace-mcp-port {} is the ingress's tunnel port; pick another",
420 cfg.workspace_mcp_port
421 )));
422 }
423 }
424 let workspaces = workspaces::Workspaces::new(
425 &cfg.state_dir,
426 client.clone(),
427 secrets.clone(),
428 recorder.clone(),
429 cfg.workspace_mcp_port,
430 cfg.workspace_pool.clone(),
431 cfg.workspace_home_root.clone(),
432 );
433 workspaces.previews.set_base(cfg.preview_domain.clone());
434 let (egress, stop_egress) = egress::start(&cfg, &client, &secrets)?;
435 let meta = crate::stack::deployments::StackMeta::new(&cfg.state_dir);
436 meta.recover();
437 let d = Arc::new(Daemon {
438 client,
439 ctl: ctl.clone(),
440 policy: cfg.policy.clone(),
441 state_dir: cfg.state_dir.clone(),
442 secrets,
443 apps: apps.clone(),
444 ingress: ingress.clone(),
445 users: users.clone(),
446 notifier: notifier.clone(),
447 monitors: monitors.clone(),
448 history,
449 data: data::Ctx {
450 apps: apps.clone(),
451 jobs,
452 backups,
453 },
454 volumes,
455 audit: audit_log.clone(),
456 gate: gate.clone(),
457 host: superadmin::host_summary(&cfg, &gate),
458 catalogs: Arc::new(crate::template::catalog::Catalogs::new(&cfg.state_dir)),
459 workspaces: workspaces.clone(),
460 public_url: cfg.public_url.clone(),
461 egress,
462 meta,
463 });
464 let registry = registry::registry(d.clone(), &cfg.skills)?;
465 let mut hooks = hooks(d.clone(), users.clone(), cfg.allow_unauthenticated);
466 superadmin::announce(&cfg, &gate, &users);
467 hooks.audit = Some(audit::hook(audit_log.clone(), cfg.audit_all));
468 let webhooks = audit::audited_webhooks(apps::webhook_routes(apps.clone()), audit_log.clone());
471 let auth = auth.map(|a| -> crate::server::Routes {
473 let logo = templates::logo::route(
474 d.catalogs.clone(),
475 Arc::new(templates::logo::Logos::new(&cfg.state_dir)),
476 templates::logo::admit(
477 hooks.authn.clone().expect("the daemon authenticates"),
478 access.clone(),
479 cfg.allow_unauthenticated,
480 ),
481 );
482 Arc::new(move |r| logo(r).or_else(|| a(r)))
483 });
484 let mut listeners = vec![Listener::unix(&cfg.socket).hooks(hooks.clone())];
485 for addr in &cfg.listen {
486 let tailnet = superadmin::is_tailnet_listen(addr);
487 let mut l = Listener::tcp(addr.clone())
488 .policy(cfg.remote_tools.clone())
489 .hooks(hooks.clone())
490 .public_routes(webhooks.clone())
491 .preview(workspaces::preview_route(d.clone()))
492 .tailnet(tailnet);
493 if let Some(r) = &auth {
494 l = l.routes(r.clone());
495 }
496 listeners.push(match &access {
497 Some(v) if !tailnet => l.access_shared(v.clone()),
499 _ => l.allow_unauthenticated(true),
503 });
504 }
505 let hd = d.clone();
506 let healthz: crate::server::Healthz = Arc::new(move || {
507 let stacks: Vec<Value> = hd
508 .ctl
509 .list()
510 .into_iter()
511 .map(|s| json!({"name": s.name, "converged": s.converged}))
512 .collect();
513 (
514 true,
515 json!({"ok": true, "isb": env!("CARGO_PKG_VERSION"), "stacks": stacks}),
516 )
517 });
518 let registry = Arc::new(registry);
522 workspaces.set_serving(
523 Listener::tcp("org-bridge")
524 .policy(cfg.remote_tools.clone())
525 .hooks(hooks.clone())
526 .allow_unauthenticated(true),
527 registry.clone(),
528 healthz.clone(),
529 );
530 workspaces.start(ctl.clone());
531 workspaces::start_ports(d.clone());
532 let r = crate::server::serve_shared(listeners, registry, healthz);
533 workspaces.shutdown();
534 stop_egress.store(true, std::sync::atomic::Ordering::Relaxed);
535 stop_history.store(true, std::sync::atomic::Ordering::Relaxed);
536 recorder.record(crate::history::marker(
537 "serve.stopped",
538 "isb serve stopped: incus events from now on are not observed".into(),
539 json!({"version": env!("CARGO_PKG_VERSION")}),
540 ));
541 recorder.shutdown();
542 notifier.shutdown();
543 monitors.shutdown();
544 scheduler.shutdown();
545 ctl.shutdown();
546 if let Some(m) = &ingress {
547 m.shutdown();
548 }
549 r
550}
551
552fn history_start(
555 log: &Arc<crate::audit::AuditLog>,
556 rec: &Arc<crate::history::Recorder>,
557 client: &Client,
558 stop: &Arc<std::sync::atomic::AtomicBool>,
559) {
560 use crate::history::{HistoryQuery, marker};
561 let now = crate::audit::now_ms();
562 let last = log
563 .history_list(
564 &HistoryQuery::default(),
565 &crate::audit::Visibility::All,
566 None,
567 None,
568 1,
569 )
570 .ok()
571 .and_then(|v| v.into_iter().next());
572 if let Some(l) = last {
573 let clean = l.kind == "serve.stopped";
574 let reason = if clean {
575 "isb serve was not running"
576 } else {
577 "isb serve was not running (it did not stop cleanly)"
578 };
579 rec.record(marker(
580 "incus.gap",
581 format!(
582 "incus events between {} and {} were not observed: {reason}",
583 crate::history::fmt_ms(l.time),
584 crate::history::fmt_ms(now),
585 ),
586 json!({"from": l.time, "to": now, "reason": reason}),
587 ));
588 }
589 rec.record(marker(
590 "serve.started",
591 format!("isb serve {} started", env!("CARGO_PKG_VERSION")),
592 json!({"version": env!("CARGO_PKG_VERSION"), "pid": std::process::id()}),
593 ));
594 let (c, r, s) = (client.clone(), rec.clone(), stop.clone());
595 let _ = std::thread::Builder::new()
596 .name("isb-incus-events".into())
597 .spawn(move || crate::history::watch_incus(c, r, s));
598}
599
600fn with_external_drivers(
602 secrets: crate::secrets::Secrets,
603 state_dir: &std::path::Path,
604) -> Result<crate::secrets::Secrets> {
605 use crate::secrets::{Driver, local::LocalDriver, onepassword};
606 let local = Arc::new(LocalDriver::new(state_dir, secrets.keyring().clone()));
607 let token: onepassword::TokenSource =
608 Arc::new(move |org| match local.get(org, onepassword::TOKEN_SECRET) {
609 Ok((v, _)) => Ok(Some(
610 String::from_utf8(v)
611 .map_err(|_| Error::invalid("the 1Password token is not text"))?
612 .trim()
613 .to_string(),
614 )),
615 Err(e) if e.is_not_found() => Ok(None),
616 Err(e) => Err(e),
617 });
618 secrets.with_driver(Arc::new(onepassword::OnePasswordDriver::new(token)))
619}
620
621fn hooks(d: Arc<Daemon>, users: Arc<AuthStore>, allow_anonymous: bool) -> crate::server::Hooks {
623 use crate::server::Authenticated;
624 let term = terminal::terminal(d.clone());
625 let ssh = ssh::ssh(d.clone(), users.clone());
626 let u = users.clone();
627 let gate = d.gate.clone();
628 let wsa = d.workspaces.clone();
629 let authn: crate::server::mcp::Authn = Arc::new(move |req, id| {
630 match gate.resolve(req, id) {
631 superadmin::Resolved::Superadmin(s) => return Authenticated::Superadmin(s),
632 superadmin::Resolved::Refused => return Authenticated::Refused,
633 superadmin::Resolved::None => {}
634 }
635 if let Some(t) = bearer(req) {
637 if t.starts_with(crate::auth::secret::TokenKind::Workspace.prefix()) {
638 return match wsa.authenticate(t) {
639 Some(p) => Authenticated::User(Arc::new(p)),
640 None => Authenticated::Refused,
641 };
642 }
643 }
644 if req.header("authorization").is_some()
645 || req
646 .header("cookie")
647 .is_some_and(|c| c.contains("isb_session="))
648 {
649 return match u.principal_from_request(req) {
650 Some(p) => Authenticated::User(Arc::new(p)),
651 None => Authenticated::Refused,
652 };
653 }
654 if let Some(email) = id.and_then(|i| i.email.as_deref()) {
656 if let Ok(Some(p)) = u.principal_for_email(email) {
657 return Authenticated::User(Arc::new(p));
658 }
659 }
660 if let Some(p) = gate.agent(req, id) {
662 return Authenticated::User(Arc::new(p));
663 }
664 Authenticated::None
665 });
666 let authorize: crate::server::mcp::Authorize = Arc::new(move |c, tool, args, scope| {
667 let args = crate::server::aliases::alias_args(tool, args);
669 authorize_class(
670 c,
671 &tool.name,
672 audit::class_for(tool, &args),
673 args,
674 scope,
675 allow_anonymous,
676 )
677 });
678 let events: crate::server::mcp::Events = Arc::new(move |c, since| {
679 if let (Caller::Unauthenticated { .. }, false) = (c, allow_anonymous) {
680 return Err(Error::Forbidden("sign in to follow events".into()));
681 }
682 if let Caller::Access(id) = c {
683 return Err(Error::Forbidden(format!(
684 "{} has no isb account",
685 id.name()
686 )));
687 }
688 let orgs = visible_orgs(c);
689 let ctl = d.ctl.clone();
690 Ok(Box::new(move |w: &mut dyn std::io::Write| {
691 let mut since = ctl.resume_from(since);
692 loop {
693 let (seq, evs) = ctl.wait_events(since, 200, Duration::from_secs(15));
694 let mut wrote = false;
695 for e in evs {
696 if !event_visible(&orgs, &e.stack) {
697 continue;
698 }
699 let data = serde_json::to_string(&e).unwrap_or_default();
700 write!(w, "id: {}\nevent: {}\ndata: {data}\n\n", e.seq, e.level)?;
701 wrote = true;
702 }
703 if !wrote {
704 w.write_all(b": keepalive\n\n")?;
706 }
707 w.flush()?;
708 since = seq.max(since);
709 }
710 }))
711 });
712 crate::server::Hooks {
713 authn: Some(authn),
714 authorize: Some(authorize),
715 events: Some(events),
716 terminal: Some(term),
717 ssh: Some(ssh),
718 audit: None,
719 listed: Some(Arc::new(tool_listed)),
720 refuse_anonymous: !allow_anonymous,
721 }
722}
723
724fn bearer(req: &crate::server::http::Request) -> Option<&str> {
726 let (scheme, token) = req.header("authorization")?.trim().split_once(' ')?;
727 scheme.eq_ignore_ascii_case("bearer").then(|| token.trim())
728}
729
730fn args<T: DeserializeOwned>(v: Value) -> Result<T> {
731 serde_json::from_value(v).map_err(|e| Error::invalid(format!("bad arguments: {e}")))
732}
733
734fn obj(mut props: Value, required: &[&str]) -> Value {
735 props["org"] =
737 json!({"type": "string", "description": "The org to act in (default: default)."});
738 json!({"type": "object", "properties": props, "required": required, "additionalProperties": false})
739}
740
741fn qname(org: &Option<String>, name: &str) -> Result<String> {
743 let org = match org {
744 Some(o) => crate::org::OrgId::new(o.clone())?,
745 None => crate::org::OrgId::default_org(),
746 };
747 Ok(crate::stack::qualified(&org, name))
748}
749
750fn caller_name(c: &Caller) -> String {
751 c.to_string()
752}
753
754impl Daemon {
755 fn reachable(&self, c: &Caller, i: &SandboxInfo) -> bool {
758 c.is_trusted()
761 || c.principal().is_some()
762 || self.policy.any_instance
763 || i.config.contains_key("user.isb.stack")
764 || i.config.contains_key(&format!("user.{LABEL_OWNER}"))
765 }
766
767 fn oc(&self, org: &Option<String>) -> Result<Client> {
770 let org = crate::org::OrgId::new(org.as_deref().unwrap_or(crate::org::DEFAULT_ORG))?;
771 crate::org::check_exists(&self.client, &org)?;
772 Ok(crate::org::client(&self.client, &org))
773 }
774
775 fn reach(&self, c: &Caller, oc: &Client, name: &str) -> Result<SandboxInfo> {
776 let info = Sandbox::get(oc, name)?.info()?;
777 if !self.reachable(c, &info) {
778 return Err(Error::NotFound(format!("sandbox {name}")));
781 }
782 Ok(info)
783 }
784
785 fn workspaces_def(
788 &self,
789 org: &crate::org::OrgId,
790 name: &str,
791 ) -> Result<crate::workspace::Workspace> {
792 crate::workspace::Store::new(&self.state_dir)
793 .get(org, name)?
794 .ok_or_else(|| Error::NotFound(format!("org {org} has no workspace {name}")))
795 }
796
797 fn files_dir(&self, stack: &str) -> Result<PathBuf> {
798 let p = self.state_dir.join("files").join(stack);
799 std::fs::create_dir_all(&p)?;
800 Ok(p)
801 }
802}
803
804pub fn wait_settled(
807 ctl: &Controller,
808 name: &str,
809 timeout: Duration,
810) -> Result<crate::stack::controller::StackStatus> {
811 let started = Instant::now();
812 let def = ctl.definition(name)?;
813 loop {
814 let st = ctl.status(name)?;
815 let settled = st.services.iter().all(|s| {
818 let current = def.revision(&s.service).is_ok_and(|r| r == s.rev)
819 && def
820 .service(&s.service)
821 .is_ok_and(|d| d.replicas() == s.replicas);
822 current && matches!(s.state.as_str(), "converged" | "paused" | "failing")
823 });
824 if settled || started.elapsed() >= timeout {
825 return Ok(st);
826 }
827 std::thread::sleep(Duration::from_secs(1));
828 }
829}
830
831#[derive(Deserialize)]
832#[serde(untagged)]
833enum SpecArg {
834 Text(String),
835 Object(Box<SandboxSpec>),
836}
837
838#[expect(
839 clippy::too_many_lines,
840 reason = "predates the lint ratchet; split it when next changed"
841)]
842fn sandbox_create(d: &Daemon, a: Value, c: &Caller) -> Result<Value> {
843 #[derive(Deserialize)]
844 struct A {
845 spec: Value,
846 #[serde(default)]
847 wait_ready: Option<bool>,
848 #[serde(default)]
849 expires: Option<String>,
850 #[serde(default)]
851 idle_timeout: Option<String>,
852 #[serde(default)]
853 org: Option<String>,
854 }
855 let org = arg_org(&a)?;
856 let a: A = args(a)?;
857 let mut spec = match serde_json::from_value::<SpecArg>(a.spec)
858 .map_err(|e| Error::invalid(format!("spec: {e}")))?
859 {
860 SpecArg::Text(t) => serde_yaml_ng::from_str::<SandboxSpec>(&t)
861 .map_err(|e| Error::invalid(format!("spec: {e}")))?,
862 SpecArg::Object(s) => *s,
863 };
864 let name = spec
865 .name
866 .clone()
867 .ok_or_else(|| Error::invalid("spec needs container_name"))?;
868 egress::check_secrets(&d.secrets, &org, &spec)?;
869 let base = if c.is_local() {
870 std::env::current_dir()?
871 } else {
872 d.files_dir("_sandboxes")?
873 };
874 if let Caller::Superadmin(s) = c {
875 spec.labels.insert(LABEL_OWNER.into(), s.label());
877 }
878 if !c.is_trusted() {
879 d.policy.check_spec(&spec, &base)?;
880 if let Ok(sb) = Sandbox::get(&d.oc(&a.org)?, &name) {
881 if !d.reachable(c, &sb.info()?) {
883 return Err(Error::AlreadyExists(name));
884 }
885 }
886 spec.labels.insert(LABEL_OWNER.into(), owner_label(c));
887 }
888 if let Ok(sb) = Sandbox::get(&d.oc(&a.org)?, &name) {
889 let info = sb.info()?;
890 if info.config.contains_key(crate::workspace::KEY_WORKSPACE) {
892 return Err(Error::invalid(format!(
893 "{name} is the org's workspace; pick another name"
894 )));
895 }
896 }
897 if !spec
900 .raw_devices
901 .get("root")
902 .is_some_and(|r| r.contains_key("size"))
903 && workspaces::project_has_disk_limit(&d.client, &org)
904 {
905 spec.raw_devices
906 .entry("root".into())
907 .or_default()
908 .insert("size".into(), workspaces::SANDBOX_ROOT_SIZE.into());
909 }
910 let settings = d.workspaces.settings(&org)?;
912 let (expires_at, idle) = crate::workspace::sandbox_deadlines(
913 &settings,
914 a.expires.as_deref(),
915 a.idle_timeout.as_deref(),
916 now_secs(),
917 )?;
918 spec.labels
919 .insert("isb.expires_at".into(), expires_at.to_string());
920 match idle {
921 Some(s) => {
922 spec.labels.insert("isb.idle_timeout".into(), s.to_string());
923 }
924 None => {
925 spec.labels.insert("isb.idle_timeout".into(), "0".into());
926 }
927 }
928 let opts = EnsureOptions {
929 wait_ready: a.wait_ready.unwrap_or(true),
930 ..Default::default()
931 };
932 let mut log: Vec<String> = Vec::new();
933 let (sb, report) = Sandbox::connect_or_create_with_base(
934 &d.oc(&a.org)?,
935 &spec,
936 &Default::default(),
937 &base,
938 opts,
939 &mut |m| log.push(m.to_string()),
940 )?;
941 d.workspaces.mark_active(&org.incus_project(), &name);
942 d.egress.kick();
943 Ok(json!({
944 "info": sb.info()?,
945 "report": report,
946 "log": log,
947 "expires_at": expires_at,
948 "idle_timeout": idle,
949 "message": format!(
950 "{name} expires {} from now{}; sandbox_extend pushes it out.",
951 crate::workspace::human(expires_at.saturating_sub(now_secs())),
952 match idle {
953 Some(s) => format!(" and is deleted after {} idle", crate::workspace::human(s)),
954 None => String::new(),
955 }
956 ),
957 }))
958}
959
960fn owner_label(c: &Caller) -> String {
962 match c {
963 Caller::Superadmin(s) => s.label(),
964 Caller::User { principal } if principal.is_workspace() => {
965 crate::auth::WORKSPACE_ACTOR.to_string()
966 }
967 _ => format!("mcp:{}", caller_name(c)),
968 }
969}
970
971fn sandbox_extend(d: &Daemon, a: Value, c: &Caller) -> Result<Value> {
972 #[derive(Deserialize)]
973 #[serde(deny_unknown_fields)]
974 struct A {
975 name: String,
976 #[serde(default)]
977 by: Option<String>,
978 #[serde(default)]
979 idle_timeout: Option<String>,
980 #[serde(default)]
981 org: Option<String>,
982 }
983 let org = arg_org(&a)?;
984 let a: A = args(a)?;
985 let oc = d.oc(&a.org)?;
986 let info = d.reach(c, &oc, &a.name)?;
987 let labels: BTreeMap<String, String> = info
988 .config
989 .iter()
990 .filter_map(|(k, v)| k.strip_prefix("user.").map(|k| (k.to_string(), v.clone())))
991 .collect();
992 if crate::workspace::kind_of(&labels) != "sandbox" {
993 return Err(Error::invalid(format!(
994 "{} is a {}, not a sandbox: it does not expire",
995 a.name,
996 crate::workspace::kind_of(&labels)
997 )));
998 }
999 let mine = labels
1001 .get("isb.owner")
1002 .is_some_and(|o| *o == owner_label(c));
1003 let admin = match c {
1004 Caller::Local { .. } | Caller::Superadmin(_) => true,
1005 Caller::User { principal } => {
1006 principal.platform_admin
1007 || principal
1008 .role_in(&org)
1009 .is_some_and(|r| r >= crate::auth::Role::Admin)
1010 }
1011 _ => false,
1012 };
1013 if !mine && !admin {
1014 return Err(Error::Forbidden(format!(
1015 "{} was created by {}; its creator or the org's admins extend it",
1016 a.name,
1017 labels
1018 .get("isb.owner")
1019 .map(String::as_str)
1020 .unwrap_or("someone else")
1021 )));
1022 }
1023 let now = now_secs();
1024 let mut patch = serde_json::Map::new();
1025 let current = labels
1026 .get("isb.expires_at")
1027 .and_then(|v| v.parse::<u64>().ok());
1028 let by = match &a.by {
1029 Some(b) => crate::flex::parse_duration(b).map_err(Error::invalid)?,
1030 None if a.idle_timeout.is_some() => Duration::ZERO,
1031 None => Duration::from_secs(86400),
1032 };
1033 let mut expires_at = current;
1034 if !by.is_zero() {
1035 let e = crate::workspace::extended(current, by, now)?;
1036 patch.insert(
1037 crate::workspace::KEY_EXPIRES_AT.into(),
1038 json!(e.to_string()),
1039 );
1040 expires_at = Some(e);
1041 }
1042 let mut idle = labels
1043 .get("isb.idle_timeout")
1044 .and_then(|v| v.parse::<u64>().ok())
1045 .filter(|s| *s > 0);
1046 if let Some(t) = &a.idle_timeout {
1047 idle = crate::workspace::idle(t)?.map(|d| d.as_secs());
1048 patch.insert(
1049 crate::workspace::KEY_IDLE_TIMEOUT.into(),
1050 json!(idle.unwrap_or(0).to_string()),
1051 );
1052 }
1053 oc.mutate(
1054 "PATCH",
1055 &format!("/1.0/instances/{}", crate::client::encode_segment(&a.name)),
1056 Some(&json!({"config": patch})),
1057 &format!("extend sandbox {}", a.name),
1058 oc.timeouts.other,
1059 )?;
1060 d.workspaces.mark_active(&org.incus_project(), &a.name);
1061 Ok(json!({
1062 "name": a.name,
1063 "expires_at": expires_at,
1064 "idle_timeout": idle,
1065 "message": format!(
1066 "{} now expires {} from now.",
1067 a.name,
1068 crate::workspace::human(expires_at.unwrap_or(now).saturating_sub(now))
1069 ),
1070 }))
1071}
1072
1073const OUTPUT_CAP: usize = 256 * 1024;
1074
1075fn cap(b: &[u8]) -> (String, bool) {
1076 if b.len() <= OUTPUT_CAP {
1077 return (String::from_utf8_lossy(b).into_owned(), false);
1078 }
1079 (
1080 String::from_utf8_lossy(&b[b.len() - OUTPUT_CAP..]).into_owned(),
1081 true,
1082 )
1083}
1084
1085fn sandbox_exec(d: &Daemon, a: Value, c: &Caller) -> Result<Value> {
1086 #[derive(Deserialize)]
1087 struct A {
1088 name: String,
1089 #[serde(default)]
1090 org: Option<String>,
1091 argv: Vec<String>,
1092 cwd: Option<String>,
1093 user: Option<String>,
1094 #[serde(default)]
1095 env: BTreeMap<String, String>,
1096 stdin: Option<String>,
1097 timeout: Option<String>,
1098 }
1099 let org = arg_org(&a)?;
1100 let a: A = args(a)?;
1101 let oc = d.oc(&a.org)?;
1102 d.reach(c, &oc, &a.name)?;
1103 d.workspaces.mark_active(&org.incus_project(), &a.name);
1104 let timeout = match &a.timeout {
1105 Some(t) => crate::flex::parse_duration(t).map_err(Error::invalid)?,
1106 None => Duration::from_secs(600),
1107 };
1108 let mut opts = ExecOptions::default().timeout(timeout);
1109 opts.cwd = a.cwd;
1110 opts.user = a.user;
1111 opts.env = a.env;
1112 if let Some(s) = a.stdin {
1113 opts.stdin = Stdin::Bytes(s.into_bytes());
1114 }
1115 let sb = Sandbox::get(&oc, &a.name)?;
1116 let out = match sb.exec_with(a.argv, opts) {
1117 Err(Error::ExecTimeout { .. }) => {
1118 return Err(Error::invalid(format!(
1119 "timed out after {timeout:?} and was killed"
1120 )));
1121 }
1122 r => r?,
1123 };
1124 let (stdout, t1) = cap(&out.stdout);
1125 let (stderr, t2) = cap(&out.stderr);
1126 Ok(
1127 json!({"exit_code": out.exit_code, "stdout": stdout, "stderr": stderr, "truncated": t1 || t2}),
1128 )
1129}
1130
1131pub use crate::stack::local_deploy_args;
1132
1133pub fn default_state_dir() -> PathBuf {
1135 Store::default_dir()
1136}
1137
1138#[cfg(test)]
1139#[path = "agent_tests.rs"]
1140mod agent_tests;
1141
1142#[cfg(test)]
1143mod tests;
1144
1145#[cfg(test)]
1146mod downscope_tests;