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 secret_hooks;
47pub mod secrets;
48mod servers;
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 audit_retention: Duration,
135 pub audit_all: bool,
137 pub history_retention: Duration,
139 pub history_max_rows: i64,
140 pub agent: Option<AgentConfig>,
143 pub superadmin_tailnet: Option<crate::server::tailnet::AllowList>,
146 pub superadmin_access: Option<superadmin::AccessAllowList>,
149 pub dev_superadmin: Option<String>,
152 pub heartbeat: Option<crate::monitor::heartbeat::Heartbeat>,
154 pub egress_pins: Vec<String>,
157 pub egress_ca: Vec<PathBuf>,
159}
160
161#[derive(Debug, Clone)]
163pub struct AgentConfig {
164 pub listen: String,
166 pub tls_dir: PathBuf,
168}
169
170fn auth_routes(
174 cfg: &ServeConfig,
175 store: Arc<AuthStore>,
176 secrets: &Arc<crate::secrets::Secrets>,
177 log: &Arc<crate::audit::AuditLog>,
178 gate: Arc<superadmin::Gate>,
179 orgs: crate::auth::ops::OrgsFn,
180) -> Result<crate::server::Routes> {
181 use crate::auth::oauth::SecretFn;
182 let default_org = crate::org::OrgId::default_org();
183 let lookup = |name: &str| -> Option<SecretFn> {
184 secrets.inspect(&default_org, name).ok()?;
185 let (s, org, name) = (secrets.clone(), default_org.clone(), name.to_string());
186 Some(Arc::new(move || {
187 let (v, _) = s.get(&org, &name).map_err(|e| e.to_string())?;
188 String::from_utf8(v)
189 .map(|v| v.trim().to_string())
190 .map_err(|_| "the secret is not UTF-8".to_string())
191 }))
192 };
193 let (providers, notes) = cfg.oauth.providers(&lookup);
194 for n in notes {
195 eprintln!("isb serve: {n}");
196 }
197 let path = crate::auth::db_path(&cfg.state_dir);
198 let agent_ways = gate.agent_ways().with_public_url(cfg.public_url.clone());
199 let api = AuthApi::new(
200 store.clone(),
201 ApiConfig {
202 agent: Some(gate.agent_fn()),
203 agent_ways,
204 public_url: cfg.public_url.clone(),
205 notifier: None,
206 setup_token_file: Some(cfg.state_dir.join("setup-token")),
207 edge: Some(gate.edge_fn()),
208 providers,
209 open_signup: cfg.open_signup,
210 audit: Some(log.clone()),
211 superadmin: Some(Arc::new(move |r: &crate::server::http::Request| match gate
212 .resolve(r, None)
213 {
214 superadmin::Resolved::Superadmin(s) => Some(s),
215 _ => None,
216 })),
217 orgs: Some(orgs),
218 },
219 )?;
220 eprintln!("isb serve: identity store {}", path.display());
221 if !crate::web::BUILT {
222 eprintln!(
223 "isb serve: this binary was built without the web UI (a placeholder page is served)"
224 );
225 }
226 let (auth, web) = (Arc::new(api).router(), crate::web::routes());
229 let tail = audit::stream_route(log.clone(), store);
230 Ok(Arc::new(move |r| {
231 auth(r).or_else(|| tail(r)).or_else(|| web(r))
232 }))
233}
234
235struct Daemon {
236 client: Client,
237 ctl: Controller,
238 policy: RemotePolicy,
239 state_dir: PathBuf,
240 secrets: Arc<crate::secrets::Secrets>,
241 apps: crate::app::Apps,
242 ingress: Option<Arc<crate::ingress::Manager>>,
243 users: Arc<AuthStore>,
245 notifier: crate::notify::Notifier,
246 monitors: crate::monitor::Monitors,
247 history: crate::metrics_history::History,
248 data: data::Ctx,
250 volumes: crate::volume_backup::VolumeBackups,
252 audit: Arc<crate::audit::AuditLog>,
253 servers: Option<Arc<crate::servers::Servers>>,
256 gate: Arc<superadmin::Gate>,
258 host: Value,
259 catalogs: Arc<crate::template::catalog::Catalogs>,
261 workspaces: Arc<workspaces::Workspaces>,
264 public_url: Option<String>,
266 egress: Arc<isb_egress::Manager>,
268 meta: crate::stack::deployments::StackMeta,
270}
271
272#[expect(
274 clippy::too_many_lines,
275 reason = "predates the lint ratchet; split it when next changed"
276)]
277pub fn serve(client: Client, cfg: ServeConfig) -> Result<()> {
278 client
279 .server_info()
280 .map_err(|e| Error::invalid(format!("isb serve needs incusd: {e}")))?;
281 default_org::warn_old_incus(&client);
282 let store = Store::open(&cfg.state_dir)?;
283 dns::open_dns_path(&cfg.state_dir);
284 if cfg.agent.is_none() {
287 default_org::ensure(&client, &store);
288 }
289 let secrets_config = crate::secrets::SecretsConfig::load(&cfg.secrets_config)?;
290 let opened = crate::secrets::Secrets::open(&cfg.state_dir, &cfg.keys, &secrets_config)?;
291 for n in &opened.notes {
292 eprintln!("isb serve: {n}");
293 }
294 let secrets = Arc::new(with_external_drivers(opened.secrets, &cfg.state_dir)?);
295 for r in crate::stack::migrate::run(&store, &secrets, Some(&client)) {
298 match r {
299 Ok(m) => eprintln!("isb serve: {m}"),
300 Err(e) => eprintln!("isb serve: WARNING: {e}"),
301 }
302 }
303 let db = crate::auth::db_path(&cfg.state_dir);
306 let users = Arc::new(
307 AuthStore::open_with(&db, cfg.auth.clone())
308 .map_err(|e| Error::invalid(format!("open {}: {e}", db.display())))?,
309 );
310 let audit_db = crate::audit::db_path(&cfg.state_dir);
311 let audit_log = Arc::new(
312 crate::audit::AuditLog::open(&audit_db, cfg.audit_retention)?
313 .with_history_limits(cfg.history_retention, cfg.history_max_rows),
314 );
315 let recorder = crate::history::Recorder::start(audit_log.clone());
318 let stop_history = Arc::new(std::sync::atomic::AtomicBool::new(false));
319 history_start(&audit_log, &recorder, &client, &stop_history);
320 eprintln!(
321 "isb serve: audit log {} (kept {} days{})",
322 audit_db.display(),
323 cfg.audit_retention.as_secs() / 86400,
324 if cfg.audit_all {
325 ", reads included"
326 } else {
327 ""
328 }
329 );
330 if cfg.agent.is_some() && !cfg.listen.is_empty() {
331 return Err(Error::invalid(
332 "--agent serves its control plane only: drop --listen (users reach the control plane)",
333 ));
334 }
335 let access = match &cfg.access {
336 Some((team, aud)) => Some(Arc::new(AccessValidator::new(team, aud)?)),
337 None => None,
338 };
339 let gate = Arc::new(superadmin::gate(&cfg, users.clone(), access.clone())?);
340 let servers = match &cfg.agent {
341 None => Some(crate::servers::Servers::open(&cfg.state_dir)?),
342 Some(_) => None,
343 };
344 let auth = if cfg.listen.is_empty() {
345 None
346 } else {
347 Some(auth_routes(
348 &cfg,
349 users.clone(),
350 &secrets,
351 &audit_log,
352 gate.clone(),
353 orgs::existing_fn(client.clone(), servers.clone()),
354 )?)
355 };
356 match crate::registry::Registry::open(&client, Some(&cfg.state_dir)) {
359 Ok(Some(r)) => {
360 eprintln!("isb serve: local registry {}", r.info().url());
361 crate::registry::install(Arc::new(r));
362 }
363 Ok(None) => {}
364 Err(e) => eprintln!("isb serve: WARNING: local registry: {e}"),
365 }
366 let ingress = match &cfg.ingress {
369 Some(ic) => Some(crate::ingress::Manager::new(
370 ic.clone(),
371 client.clone(),
372 secrets.clone(),
373 &cfg.state_dir,
374 )?),
375 None => None,
376 };
377 let observer = ingress
378 .clone()
379 .map(|m| m as Arc<dyn crate::stack::controller::Observer>);
380 let ctl = Controller::start_with(
381 client.clone(),
382 store,
383 cfg.interval,
384 secrets.clone(),
385 observer,
386 )?;
387 if let Some(m) = &ingress {
388 m.start(ctl.clone())?;
389 }
390 let apps = crate::app::Apps::new(&cfg.state_dir, client.clone(), ctl.clone(), secrets.clone());
391 let stacks: Vec<(crate::org::OrgId, String)> = ctl
395 .definitions()
396 .iter()
397 .map(|d| (d.org.clone(), d.name.clone()))
398 .collect();
399 for r in apps.adopt_compose_stacks(&stacks) {
400 match r {
401 Ok(m) => eprintln!("isb serve: {m}"),
402 Err(e) => eprintln!("isb serve: WARNING: {e}"),
403 }
404 }
405 ctl.set_dns_scope(apps.scope_fn());
407 let ra = apps.clone();
409 let resolve: crate::notify::Resolve = Arc::new(move |org, stack, service| {
410 let a = ra.get(org, service).ok()?;
411 let own = a.spec.stack().ok()? == stack;
413 let preview = stack.starts_with(&format!("{}-", a.spec.project))
414 && crate::app::preview::is_pr_suffix(stack);
415 (own || preview).then_some(a.spec.project)
416 });
417 let notifier = crate::notify::Notifier::new(&cfg.state_dir, secrets.clone(), resolve)?;
418 notifier.start(ctl.clone());
419 let monitors = monitors::start(&cfg, &apps, &secrets, ¬ifier);
420 ctl.set_event_sink(recorder.controller_sink());
422 let history = crate::metrics_history::History::new(&cfg.state_dir);
424 ctl.set_metrics_sink(history.start());
425 apps.start_preview_upkeep();
427 let jobs = crate::jobs::Jobs::new(&cfg.state_dir, apps.clone());
429 let backups = crate::backup::Backups::new(&cfg.state_dir, apps.clone());
430 let volumes =
431 crate::volume_backup::VolumeBackups::new(&cfg.state_dir, apps.clone(), backups.clone());
432 let scheduler = crate::jobs::Scheduler::start(vec![
433 Arc::new(jobs.clone()) as Arc<dyn crate::jobs::Scheduled>,
434 Arc::new(backups.clone()),
435 Arc::new(volumes.clone()),
436 ]);
437 jobs.set_scheduler(scheduler.clone());
438 backups.set_scheduler(scheduler.clone());
439 volumes.set_scheduler(scheduler.clone());
440 if let Some(ic) = &cfg.ingress {
441 if ic.tunnel_port == cfg.workspace_mcp_port {
442 return Err(Error::invalid(format!(
443 "--workspace-mcp-port {} is the ingress's tunnel port; pick another",
444 cfg.workspace_mcp_port
445 )));
446 }
447 }
448 let workspaces = workspaces::Workspaces::new(
449 &cfg.state_dir,
450 client.clone(),
451 secrets.clone(),
452 recorder.clone(),
453 cfg.workspace_mcp_port,
454 cfg.workspace_pool.clone(),
455 cfg.workspace_home_root.clone(),
456 );
457 workspaces.previews.set_base(cfg.preview_domain.clone());
458 let (egress, stop_egress) = egress::start(&cfg, &client, &secrets)?;
459 let meta = crate::stack::deployments::StackMeta::new(&cfg.state_dir);
460 meta.recover();
461 let d = Arc::new(Daemon {
462 client,
463 ctl: ctl.clone(),
464 policy: cfg.policy.clone(),
465 state_dir: cfg.state_dir.clone(),
466 secrets,
467 apps: apps.clone(),
468 ingress: ingress.clone(),
469 users: users.clone(),
470 notifier: notifier.clone(),
471 monitors: monitors.clone(),
472 history,
473 data: data::Ctx {
474 apps: apps.clone(),
475 jobs,
476 backups,
477 },
478 volumes,
479 audit: audit_log.clone(),
480 servers: servers.clone(),
481 gate: gate.clone(),
482 host: superadmin::host_summary(&cfg, &gate),
483 catalogs: Arc::new(crate::template::catalog::Catalogs::new(&cfg.state_dir)),
484 workspaces: workspaces.clone(),
485 public_url: cfg.public_url.clone(),
486 egress,
487 meta,
488 });
489 if let Some(s) = &servers {
490 s.start(ctl.clone());
491 }
492 let registry = registry(d.clone())?;
493 let mut hooks = hooks(d.clone(), users.clone(), cfg.allow_unauthenticated);
494 superadmin::announce(&cfg, &gate, &users);
495 hooks.audit = Some(audit::hook(audit_log.clone(), cfg.audit_all));
496 if servers.is_some() {
497 hooks.route = Some(servers::route(d.clone()));
498 }
499 let webhooks = {
502 let w = apps::webhook_routes(apps.clone());
503 let w = match &servers {
504 Some(s) => servers::forward_webhooks(w, s.clone()),
505 None => w,
506 };
507 audit::audited_webhooks(w, audit_log.clone())
508 };
509 let auth = auth.map(|a| -> crate::server::Routes {
511 let logo = templates::logo::route(
512 d.catalogs.clone(),
513 Arc::new(templates::logo::Logos::new(&cfg.state_dir)),
514 templates::logo::admit(
515 hooks.authn.clone().expect("the daemon authenticates"),
516 access.clone(),
517 cfg.allow_unauthenticated,
518 ),
519 );
520 Arc::new(move |r| logo(r).or_else(|| a(r)))
521 });
522 let mut listeners = vec![Listener::unix(&cfg.socket).hooks(hooks.clone())];
523 if let Some(ac) = &cfg.agent {
524 listeners.push(servers::agent_listener(
525 d.clone(),
526 &hooks,
527 ac,
528 webhooks.clone(),
529 )?);
530 }
531 for addr in &cfg.listen {
532 let tailnet = superadmin::is_tailnet_listen(addr);
533 let mut l = Listener::tcp(addr.clone())
534 .policy(cfg.remote_tools.clone())
535 .hooks(hooks.clone())
536 .public_routes(webhooks.clone())
537 .preview(workspaces::preview_route(d.clone()))
538 .tailnet(tailnet);
539 if let Some(r) = &auth {
540 l = l.routes(r.clone());
541 }
542 listeners.push(match &access {
543 Some(v) if !tailnet => l.access_shared(v.clone()),
545 _ => l.allow_unauthenticated(true),
549 });
550 }
551 let hd = d.clone();
552 let healthz: crate::server::Healthz = Arc::new(move || {
553 let stacks: Vec<Value> = hd
554 .ctl
555 .list()
556 .into_iter()
557 .map(|s| json!({"name": s.name, "converged": s.converged}))
558 .collect();
559 (
560 true,
561 json!({"ok": true, "isb": env!("CARGO_PKG_VERSION"), "stacks": stacks}),
562 )
563 });
564 let registry = Arc::new(registry);
568 workspaces.set_serving(
569 Listener::tcp("org-bridge")
570 .policy(cfg.remote_tools.clone())
571 .hooks(hooks.clone())
572 .allow_unauthenticated(true),
573 registry.clone(),
574 healthz.clone(),
575 );
576 let local: workspaces::LocalOrg = {
577 let d = d.clone();
578 Arc::new(move |o: &crate::org::OrgId| d.remote(o).is_none())
579 };
580 workspaces.start(ctl.clone(), local);
581 workspaces::start_ports(d.clone());
582 let r = crate::server::serve_shared(listeners, registry, healthz);
583 workspaces.shutdown();
584 stop_egress.store(true, std::sync::atomic::Ordering::Relaxed);
585 stop_history.store(true, std::sync::atomic::Ordering::Relaxed);
586 recorder.record(crate::history::marker(
587 "serve.stopped",
588 "isb serve stopped: incus events from now on are not observed".into(),
589 json!({"version": env!("CARGO_PKG_VERSION")}),
590 ));
591 recorder.shutdown();
592 if let Some(s) = &servers {
593 s.shutdown();
594 }
595 notifier.shutdown();
596 monitors.shutdown();
597 scheduler.shutdown();
598 ctl.shutdown();
599 if let Some(m) = &ingress {
600 m.shutdown();
601 }
602 r
603}
604
605fn history_start(
608 log: &Arc<crate::audit::AuditLog>,
609 rec: &Arc<crate::history::Recorder>,
610 client: &Client,
611 stop: &Arc<std::sync::atomic::AtomicBool>,
612) {
613 use crate::history::{HistoryQuery, marker};
614 let now = crate::audit::now_ms();
615 let last = log
616 .history_list(
617 &HistoryQuery::default(),
618 &crate::audit::Visibility::All,
619 None,
620 None,
621 1,
622 )
623 .ok()
624 .and_then(|v| v.into_iter().next());
625 if let Some(l) = last {
626 let clean = l.kind == "serve.stopped";
627 let reason = if clean {
628 "isb serve was not running"
629 } else {
630 "isb serve was not running (it did not stop cleanly)"
631 };
632 rec.record(marker(
633 "incus.gap",
634 format!(
635 "incus events between {} and {} were not observed: {reason}",
636 crate::history::fmt_ms(l.time),
637 crate::history::fmt_ms(now),
638 ),
639 json!({"from": l.time, "to": now, "reason": reason}),
640 ));
641 }
642 rec.record(marker(
643 "serve.started",
644 format!("isb serve {} started", env!("CARGO_PKG_VERSION")),
645 json!({"version": env!("CARGO_PKG_VERSION"), "pid": std::process::id()}),
646 ));
647 let (c, r, s) = (client.clone(), rec.clone(), stop.clone());
648 let _ = std::thread::Builder::new()
649 .name("isb-incus-events".into())
650 .spawn(move || crate::history::watch_incus(c, r, s));
651}
652
653fn with_external_drivers(
655 secrets: crate::secrets::Secrets,
656 state_dir: &std::path::Path,
657) -> Result<crate::secrets::Secrets> {
658 use crate::secrets::{Driver, local::LocalDriver, onepassword};
659 let local = Arc::new(LocalDriver::new(state_dir, secrets.keyring().clone()));
660 let token: onepassword::TokenSource =
661 Arc::new(move |org| match local.get(org, onepassword::TOKEN_SECRET) {
662 Ok((v, _)) => Ok(Some(
663 String::from_utf8(v)
664 .map_err(|_| Error::invalid("the 1Password token is not text"))?
665 .trim()
666 .to_string(),
667 )),
668 Err(e) if e.is_not_found() => Ok(None),
669 Err(e) => Err(e),
670 });
671 secrets.with_driver(Arc::new(onepassword::OnePasswordDriver::new(token)))
672}
673
674fn hooks(d: Arc<Daemon>, users: Arc<AuthStore>, allow_anonymous: bool) -> crate::server::Hooks {
676 use crate::server::Authenticated;
677 let term = terminal::terminal(d.clone());
678 let ssh = ssh::ssh(d.clone(), users.clone());
679 let u = users.clone();
680 let gate = d.gate.clone();
681 let wsa = d.workspaces.clone();
682 let authn: crate::server::mcp::Authn = Arc::new(move |req, id| {
683 match gate.resolve(req, id) {
684 superadmin::Resolved::Superadmin(s) => return Authenticated::Superadmin(s),
685 superadmin::Resolved::Refused => return Authenticated::Refused,
686 superadmin::Resolved::None => {}
687 }
688 if let Some(t) = bearer(req) {
690 if t.starts_with(crate::auth::secret::TokenKind::Workspace.prefix()) {
691 return match wsa.authenticate(t) {
692 Some(p) => Authenticated::User(Arc::new(p)),
693 None => Authenticated::Refused,
694 };
695 }
696 }
697 if req.header("authorization").is_some()
698 || req
699 .header("cookie")
700 .is_some_and(|c| c.contains("isb_session="))
701 {
702 return match u.principal_from_request(req) {
703 Some(p) => Authenticated::User(Arc::new(p)),
704 None => Authenticated::Refused,
705 };
706 }
707 if let Some(email) = id.and_then(|i| i.email.as_deref()) {
709 if let Ok(Some(p)) = u.principal_for_email(email) {
710 return Authenticated::User(Arc::new(p));
711 }
712 }
713 if let Some(p) = gate.agent(req, id) {
715 return Authenticated::User(Arc::new(p));
716 }
717 Authenticated::None
718 });
719 let authorize: crate::server::mcp::Authorize = Arc::new(move |c, tool, args, scope| {
720 let args = crate::server::aliases::alias_args(tool, args);
722 authorize_class(
723 c,
724 &tool.name,
725 audit::class_for(tool, &args),
726 args,
727 scope,
728 allow_anonymous,
729 )
730 });
731 let events: crate::server::mcp::Events = Arc::new(move |c, since| {
732 if let (Caller::Unauthenticated { .. }, false) = (c, allow_anonymous) {
733 return Err(Error::Forbidden("sign in to follow events".into()));
734 }
735 if let Caller::Access(id) = c {
736 return Err(Error::Forbidden(format!(
737 "{} has no isb account",
738 id.name()
739 )));
740 }
741 let orgs = visible_orgs(c);
742 let ctl = d.ctl.clone();
743 Ok(Box::new(move |w: &mut dyn std::io::Write| {
744 let mut since = ctl.resume_from(since);
745 loop {
746 let (seq, evs) = ctl.wait_events(since, 200, Duration::from_secs(15));
747 let mut wrote = false;
748 for e in evs {
749 if !event_visible(&orgs, &e.stack) {
750 continue;
751 }
752 let data = serde_json::to_string(&e).unwrap_or_default();
753 write!(w, "id: {}\nevent: {}\ndata: {data}\n\n", e.seq, e.level)?;
754 wrote = true;
755 }
756 if !wrote {
757 w.write_all(b": keepalive\n\n")?;
759 }
760 w.flush()?;
761 since = seq.max(since);
762 }
763 }))
764 });
765 crate::server::Hooks {
766 authn: Some(authn),
767 authorize: Some(authorize),
768 events: Some(events),
769 terminal: Some(term),
770 ssh: Some(ssh),
771 audit: None,
772 route: None,
773 listed: Some(Arc::new(tool_listed)),
774 refuse_anonymous: !allow_anonymous,
775 }
776}
777
778fn bearer(req: &crate::server::http::Request) -> Option<&str> {
780 let (scheme, token) = req.header("authorization")?.trim().split_once(' ')?;
781 scheme.eq_ignore_ascii_case("bearer").then(|| token.trim())
782}
783
784fn args<T: DeserializeOwned>(v: Value) -> Result<T> {
785 serde_json::from_value(v).map_err(|e| Error::invalid(format!("bad arguments: {e}")))
786}
787
788fn obj(mut props: Value, required: &[&str]) -> Value {
789 props["org"] =
791 json!({"type": "string", "description": "The org to act in (default: default)."});
792 json!({"type": "object", "properties": props, "required": required, "additionalProperties": false})
793}
794
795fn qname(org: &Option<String>, name: &str) -> Result<String> {
797 let org = match org {
798 Some(o) => crate::org::OrgId::new(o.clone())?,
799 None => crate::org::OrgId::default_org(),
800 };
801 Ok(crate::stack::qualified(&org, name))
802}
803
804fn caller_name(c: &Caller) -> String {
805 c.to_string()
806}
807
808fn registry(d: Arc<Daemon>) -> Result<Registry> {
810 let mut r = Registry::new().instructions(guide::INSTRUCTIONS);
811 superadmin::register(&mut r, d.clone())?;
812 host_monitor::register(&mut r, d.clone())?;
813 let ann = Ann {
814 ro: json!({"readOnlyHint": true, "openWorldHint": false}),
815 destructive: json!({"destructiveHint": true, "openWorldHint": false}),
816 write: json!({"destructiveHint": false, "openWorldHint": false}),
817 };
818
819 guide::register(&mut r, &d, &ann)?;
820 tools::stack_deploy_tool(&mut r, &d, &ann)?;
821 tools::overview_tool(&mut r, &d, &ann)?;
822 tools::events_tool(&mut r, &d, &ann)?;
823 tools::ingress_status_tool(&mut r, &d, &ann)?;
824 tools::stack_list_tool(&mut r, &d, &ann)?;
825 tools::stack_status_tool(&mut r, &d, &ann)?;
826 tools::stack_config_tool(&mut r, &d, &ann)?;
827 tools::stack_logs_tool(&mut r, &d, &ann)?;
828 tools::stack_scale_tool(&mut r, &d, &ann)?;
829 tools::stack_edit_tools(&mut r, &d, &ann)?;
830 tools::sandbox_create_tool(&mut r, &d, &ann)?;
831 secret_hooks::register(&mut r, &d)?;
832 builds::register(
833 &mut r,
834 builds::Ctx {
835 client: d.client.clone(),
836 policy: d.policy.clone(),
837 ctl: d.ctl.clone(),
838 },
839 )?;
840 tools::sandbox_tools(&mut r, &d, &ann)?;
841 apps::register(&mut r, d.apps.clone(), d.ingress.is_some())?;
842 apps::delete_tools(&mut r, &d)?;
843 previews::register(&mut r, d.apps.clone())?;
844 let mut t = templates::Templates::new(
845 &d.state_dir,
846 d.apps.clone(),
847 d.secrets.clone(),
848 d.ingress.as_ref().and_then(|m| m.public_ip()),
849 );
850 t.catalogs = d.catalogs.clone();
851 templates::register(&mut r, t)?;
852 data::register(&mut r, d.data.clone())?;
853 volumes::register(&mut r, d.volumes.clone())?;
854 tools::server_status_tool(&mut r, &d, &ann)?;
855 orgs::register(&mut r, d.clone())?;
856 notify::register(
857 &mut r,
858 d.notifier.clone(),
859 d.history.clone(),
860 d.apps.clone(),
861 )?;
862 monitors::register(&mut r, d.monitors.clone())?;
863 audit::register(&mut r, d.audit.clone())?;
864 audit::register_history(&mut r, d.audit.clone())?;
865 servers::register(&mut r, d.clone())?;
866 ssh::register(&mut r, d.clone())?;
867 workspaces::register(&mut r, d.clone())?;
868 accounts::register(&mut r, d.clone())?;
869 Ok(r)
870}
871
872impl Daemon {
873 fn reachable(&self, c: &Caller, i: &SandboxInfo) -> bool {
876 c.is_trusted()
879 || c.principal().is_some()
880 || self.policy.any_instance
881 || i.config.contains_key("user.isb.stack")
882 || i.config.contains_key(&format!("user.{LABEL_OWNER}"))
883 }
884
885 fn oc(&self, org: &Option<String>) -> Result<Client> {
888 let org = crate::org::OrgId::new(org.as_deref().unwrap_or(crate::org::DEFAULT_ORG))?;
889 crate::org::check_exists(&self.client, &org)?;
890 Ok(crate::org::client(&self.client, &org))
891 }
892
893 fn reach(&self, c: &Caller, oc: &Client, name: &str) -> Result<SandboxInfo> {
894 let info = Sandbox::get(oc, name)?.info()?;
895 if !self.reachable(c, &info) {
896 return Err(Error::NotFound(format!("sandbox {name}")));
899 }
900 Ok(info)
901 }
902
903 fn workspaces_def(
906 &self,
907 org: &crate::org::OrgId,
908 name: &str,
909 ) -> Result<crate::workspace::Workspace> {
910 crate::workspace::Store::new(&self.state_dir)
911 .get(org, name)?
912 .ok_or_else(|| Error::NotFound(format!("org {org} has no workspace {name}")))
913 }
914
915 fn files_dir(&self, stack: &str) -> Result<PathBuf> {
916 let p = self.state_dir.join("files").join(stack);
917 std::fs::create_dir_all(&p)?;
918 Ok(p)
919 }
920}
921
922pub fn wait_settled(
925 ctl: &Controller,
926 name: &str,
927 timeout: Duration,
928) -> Result<crate::stack::controller::StackStatus> {
929 let started = Instant::now();
930 let def = ctl.definition(name)?;
931 loop {
932 let st = ctl.status(name)?;
933 let settled = st.services.iter().all(|s| {
936 let current = def.revision(&s.service).is_ok_and(|r| r == s.rev)
937 && def
938 .service(&s.service)
939 .is_ok_and(|d| d.replicas() == s.replicas);
940 current && matches!(s.state.as_str(), "converged" | "paused" | "failing")
941 });
942 if settled || started.elapsed() >= timeout {
943 return Ok(st);
944 }
945 std::thread::sleep(Duration::from_secs(1));
946 }
947}
948
949#[derive(Deserialize)]
950#[serde(untagged)]
951enum SpecArg {
952 Text(String),
953 Object(Box<SandboxSpec>),
954}
955
956#[expect(
957 clippy::too_many_lines,
958 reason = "predates the lint ratchet; split it when next changed"
959)]
960fn sandbox_create(d: &Daemon, a: Value, c: &Caller) -> Result<Value> {
961 #[derive(Deserialize)]
962 struct A {
963 spec: Value,
964 #[serde(default)]
965 wait_ready: Option<bool>,
966 #[serde(default)]
967 expires: Option<String>,
968 #[serde(default)]
969 idle_timeout: Option<String>,
970 #[serde(default)]
971 org: Option<String>,
972 }
973 let org = arg_org(&a)?;
974 let a: A = args(a)?;
975 let mut spec = match serde_json::from_value::<SpecArg>(a.spec)
976 .map_err(|e| Error::invalid(format!("spec: {e}")))?
977 {
978 SpecArg::Text(t) => serde_yaml_ng::from_str::<SandboxSpec>(&t)
979 .map_err(|e| Error::invalid(format!("spec: {e}")))?,
980 SpecArg::Object(s) => *s,
981 };
982 let name = spec
983 .name
984 .clone()
985 .ok_or_else(|| Error::invalid("spec needs container_name"))?;
986 egress::check_secrets(&d.secrets, &org, &spec)?;
987 let base = if c.is_local() {
988 std::env::current_dir()?
989 } else {
990 d.files_dir("_sandboxes")?
991 };
992 if let Caller::Superadmin(s) = c {
993 spec.labels.insert(LABEL_OWNER.into(), s.label());
995 }
996 if !c.is_trusted() {
997 d.policy.check_spec(&spec, &base)?;
998 if let Ok(sb) = Sandbox::get(&d.oc(&a.org)?, &name) {
999 if !d.reachable(c, &sb.info()?) {
1001 return Err(Error::AlreadyExists(name));
1002 }
1003 }
1004 spec.labels.insert(LABEL_OWNER.into(), owner_label(c));
1005 }
1006 if let Ok(sb) = Sandbox::get(&d.oc(&a.org)?, &name) {
1007 let info = sb.info()?;
1008 if info.config.contains_key(crate::workspace::KEY_WORKSPACE) {
1010 return Err(Error::invalid(format!(
1011 "{name} is the org's workspace; pick another name"
1012 )));
1013 }
1014 }
1015 if !spec
1018 .raw_devices
1019 .get("root")
1020 .is_some_and(|r| r.contains_key("size"))
1021 && workspaces::project_has_disk_limit(&d.client, &org)
1022 {
1023 spec.raw_devices
1024 .entry("root".into())
1025 .or_default()
1026 .insert("size".into(), workspaces::SANDBOX_ROOT_SIZE.into());
1027 }
1028 let settings = d.workspaces.settings(&org)?;
1030 let (expires_at, idle) = crate::workspace::sandbox_deadlines(
1031 &settings,
1032 a.expires.as_deref(),
1033 a.idle_timeout.as_deref(),
1034 now_secs(),
1035 )?;
1036 spec.labels
1037 .insert("isb.expires_at".into(), expires_at.to_string());
1038 match idle {
1039 Some(s) => {
1040 spec.labels.insert("isb.idle_timeout".into(), s.to_string());
1041 }
1042 None => {
1043 spec.labels.insert("isb.idle_timeout".into(), "0".into());
1044 }
1045 }
1046 let opts = EnsureOptions {
1047 wait_ready: a.wait_ready.unwrap_or(true),
1048 ..Default::default()
1049 };
1050 let mut log: Vec<String> = Vec::new();
1051 let (sb, report) = Sandbox::connect_or_create_with_base(
1052 &d.oc(&a.org)?,
1053 &spec,
1054 &Default::default(),
1055 &base,
1056 opts,
1057 &mut |m| log.push(m.to_string()),
1058 )?;
1059 d.workspaces.mark_active(&org.incus_project(), &name);
1060 d.egress.kick();
1061 Ok(json!({
1062 "info": sb.info()?,
1063 "report": report,
1064 "log": log,
1065 "expires_at": expires_at,
1066 "idle_timeout": idle,
1067 "message": format!(
1068 "{name} expires {} from now{}; sandbox_extend pushes it out.",
1069 crate::workspace::human(expires_at.saturating_sub(now_secs())),
1070 match idle {
1071 Some(s) => format!(" and is deleted after {} idle", crate::workspace::human(s)),
1072 None => String::new(),
1073 }
1074 ),
1075 }))
1076}
1077
1078fn owner_label(c: &Caller) -> String {
1080 match c {
1081 Caller::Superadmin(s) => s.label(),
1082 Caller::User { principal } if principal.is_workspace() => {
1083 crate::auth::WORKSPACE_ACTOR.to_string()
1084 }
1085 _ => format!("mcp:{}", caller_name(c)),
1086 }
1087}
1088
1089fn sandbox_extend(d: &Daemon, a: Value, c: &Caller) -> Result<Value> {
1090 #[derive(Deserialize)]
1091 #[serde(deny_unknown_fields)]
1092 struct A {
1093 name: String,
1094 #[serde(default)]
1095 by: Option<String>,
1096 #[serde(default)]
1097 idle_timeout: Option<String>,
1098 #[serde(default)]
1099 org: Option<String>,
1100 }
1101 let org = arg_org(&a)?;
1102 let a: A = args(a)?;
1103 let oc = d.oc(&a.org)?;
1104 let info = d.reach(c, &oc, &a.name)?;
1105 let labels: BTreeMap<String, String> = info
1106 .config
1107 .iter()
1108 .filter_map(|(k, v)| k.strip_prefix("user.").map(|k| (k.to_string(), v.clone())))
1109 .collect();
1110 if crate::workspace::kind_of(&labels) != "sandbox" {
1111 return Err(Error::invalid(format!(
1112 "{} is a {}, not a sandbox: it does not expire",
1113 a.name,
1114 crate::workspace::kind_of(&labels)
1115 )));
1116 }
1117 let mine = labels
1119 .get("isb.owner")
1120 .is_some_and(|o| *o == owner_label(c));
1121 let admin = match c {
1122 Caller::Local { .. } | Caller::Superadmin(_) => true,
1123 Caller::User { principal } => {
1124 principal.platform_admin
1125 || principal
1126 .role_in(&org)
1127 .is_some_and(|r| r >= crate::auth::Role::Admin)
1128 }
1129 _ => false,
1130 };
1131 if !mine && !admin {
1132 return Err(Error::Forbidden(format!(
1133 "{} was created by {}; its creator or the org's admins extend it",
1134 a.name,
1135 labels
1136 .get("isb.owner")
1137 .map(String::as_str)
1138 .unwrap_or("someone else")
1139 )));
1140 }
1141 let now = now_secs();
1142 let mut patch = serde_json::Map::new();
1143 let current = labels
1144 .get("isb.expires_at")
1145 .and_then(|v| v.parse::<u64>().ok());
1146 let by = match &a.by {
1147 Some(b) => crate::flex::parse_duration(b).map_err(Error::invalid)?,
1148 None if a.idle_timeout.is_some() => Duration::ZERO,
1149 None => Duration::from_secs(86400),
1150 };
1151 let mut expires_at = current;
1152 if !by.is_zero() {
1153 let e = crate::workspace::extended(current, by, now)?;
1154 patch.insert(
1155 crate::workspace::KEY_EXPIRES_AT.into(),
1156 json!(e.to_string()),
1157 );
1158 expires_at = Some(e);
1159 }
1160 let mut idle = labels
1161 .get("isb.idle_timeout")
1162 .and_then(|v| v.parse::<u64>().ok())
1163 .filter(|s| *s > 0);
1164 if let Some(t) = &a.idle_timeout {
1165 idle = crate::workspace::idle(t)?.map(|d| d.as_secs());
1166 patch.insert(
1167 crate::workspace::KEY_IDLE_TIMEOUT.into(),
1168 json!(idle.unwrap_or(0).to_string()),
1169 );
1170 }
1171 oc.mutate(
1172 "PATCH",
1173 &format!("/1.0/instances/{}", crate::client::encode_segment(&a.name)),
1174 Some(&json!({"config": patch})),
1175 &format!("extend sandbox {}", a.name),
1176 oc.timeouts.other,
1177 )?;
1178 d.workspaces.mark_active(&org.incus_project(), &a.name);
1179 Ok(json!({
1180 "name": a.name,
1181 "expires_at": expires_at,
1182 "idle_timeout": idle,
1183 "message": format!(
1184 "{} now expires {} from now.",
1185 a.name,
1186 crate::workspace::human(expires_at.unwrap_or(now).saturating_sub(now))
1187 ),
1188 }))
1189}
1190
1191const OUTPUT_CAP: usize = 256 * 1024;
1192
1193fn cap(b: &[u8]) -> (String, bool) {
1194 if b.len() <= OUTPUT_CAP {
1195 return (String::from_utf8_lossy(b).into_owned(), false);
1196 }
1197 (
1198 String::from_utf8_lossy(&b[b.len() - OUTPUT_CAP..]).into_owned(),
1199 true,
1200 )
1201}
1202
1203fn sandbox_exec(d: &Daemon, a: Value, c: &Caller) -> Result<Value> {
1204 #[derive(Deserialize)]
1205 struct A {
1206 name: String,
1207 #[serde(default)]
1208 org: Option<String>,
1209 argv: Vec<String>,
1210 cwd: Option<String>,
1211 user: Option<String>,
1212 #[serde(default)]
1213 env: BTreeMap<String, String>,
1214 stdin: Option<String>,
1215 timeout: Option<String>,
1216 }
1217 let org = arg_org(&a)?;
1218 let a: A = args(a)?;
1219 let oc = d.oc(&a.org)?;
1220 d.reach(c, &oc, &a.name)?;
1221 d.workspaces.mark_active(&org.incus_project(), &a.name);
1222 let timeout = match &a.timeout {
1223 Some(t) => crate::flex::parse_duration(t).map_err(Error::invalid)?,
1224 None => Duration::from_secs(600),
1225 };
1226 let mut opts = ExecOptions::default().timeout(timeout);
1227 opts.cwd = a.cwd;
1228 opts.user = a.user;
1229 opts.env = a.env;
1230 if let Some(s) = a.stdin {
1231 opts.stdin = Stdin::Bytes(s.into_bytes());
1232 }
1233 let sb = Sandbox::get(&oc, &a.name)?;
1234 let out = match sb.exec_with(a.argv, opts) {
1235 Err(Error::ExecTimeout { .. }) => {
1236 return Err(Error::invalid(format!(
1237 "timed out after {timeout:?} and was killed"
1238 )));
1239 }
1240 r => r?,
1241 };
1242 let (stdout, t1) = cap(&out.stdout);
1243 let (stderr, t2) = cap(&out.stderr);
1244 Ok(
1245 json!({"exit_code": out.exit_code, "stdout": stdout, "stderr": stderr, "truncated": t1 || t2}),
1246 )
1247}
1248
1249pub use crate::stack::local_deploy_args;
1250
1251pub fn default_state_dir() -> PathBuf {
1253 Store::default_dir()
1254}
1255
1256#[cfg(test)]
1257#[path = "agent_tests.rs"]
1258mod agent_tests;
1259
1260#[cfg(test)]
1261mod tests;
1262
1263#[cfg(test)]
1264mod downscope_tests;