Skip to main content

isb_daemon/daemon/
mod.rs

1//! `isb serve`: the stack controller behind an MCP server.
2//!
3//! Two doors, one set of tools:
4//! - the unix socket, for the local CLI (`isb stack ...`), trusted as the
5//!   daemon's own user;
6//! - loopback HTTP at `/mcp`, for remote agents through a cloudflared tunnel
7//!   and Cloudflare Access, held to [`policy::RemotePolicy`].
8//!
9//! The tools manage stacks (long-running, replicated, load-balanced
10//! services), apps over them ([`crate::app`]), plain sandboxes (an isolated
11//! machine for an agent), and each org's secrets ([`crate::secrets`]).
12
13/// Register one tool: `tool!(registry, ctx, name, title, description,
14/// input_schema, annotations, handler)`. The handler is called with its own
15/// clone of `ctx`, the arguments and the caller. Defined before the
16/// submodules so their tool tables use it too.
17macro_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 ssh;
49mod stack_deploy;
50pub mod superadmin;
51pub mod templates;
52mod terminal;
53mod tools;
54pub mod volumes;
55pub mod workspaces;
56
57use std::collections::BTreeMap;
58use std::path::PathBuf;
59use std::sync::Arc;
60use std::time::{Duration, Instant};
61
62use serde::Deserialize;
63use serde::de::DeserializeOwned;
64use serde_json::{Value, json};
65
66use crate::auth::AuthConfig;
67use crate::auth::AuthStore;
68use crate::auth::http::{ApiConfig, AuthApi};
69use crate::client::Client;
70use crate::error::{Error, Result};
71use crate::exec::{ExecOptions, Stdin};
72use crate::sandbox::{EnsureOptions, Sandbox, SandboxInfo};
73use crate::server::{AccessValidator, Caller, Listener, Registry, Tool, ToolPolicy};
74use crate::spec::SandboxSpec;
75use crate::stack::{Controller, Store, now_secs};
76use authorize::{
77    CROSS_ORG_READS, PLATFORM_TOOLS, arg_org, authorize_class, event_visible, read_orgs,
78    tool_listed, visible_orgs,
79};
80use policy::RemotePolicy;
81use stack_deploy::{DeployArgs, How, deploy, stack_deploy};
82#[cfg(test)]
83use stack_deploy::{reused_note, stack_owner};
84use tools::Ann;
85
86/// Marks an instance a remote caller created with `sandbox_create`.
87pub const LABEL_OWNER: &str = "isb.owner";
88
89/// How `isb serve` runs.
90#[derive(Debug, Clone)]
91pub struct ServeConfig {
92    /// `host:port` addresses for remote MCP and the web UI: loopback, or a
93    /// tailnet address with `superadmin_tailnet`. Empty serves the socket
94    /// only.
95    pub listen: Vec<String>,
96    pub socket: PathBuf,
97    /// Cloudflare Access team domain and application audience.
98    pub access: Option<(String, String)>,
99    /// Serve the TCP listener with no Access (local testing only).
100    pub allow_unauthenticated: bool,
101    /// Which tools remote callers see.
102    pub remote_tools: ToolPolicy,
103    pub policy: RemotePolicy,
104    pub state_dir: PathBuf,
105    pub interval: Duration,
106    /// Where the secrets key is looked for (and generated).
107    pub keys: crate::secrets::KeySources,
108    /// `~/.config/isb/secrets.toml`: break-glass recipients.
109    pub secrets_config: PathBuf,
110    /// Session lifetimes and the rest of the identity store's settings.
111    pub auth: AuthConfig,
112    /// Where users reach isb, for invitation and reset links, provider
113    /// callbacks and the passkey relying party.
114    pub public_url: Option<String>,
115    /// External sign-in providers (GitHub, Google, generic OIDC).
116    pub oauth: crate::auth::oauth::OAuthSettings,
117    /// Accounts without an invitation, for verified provider emails.
118    pub open_signup: bool,
119    /// The HTTP(S) edge for stack domains; `None` leaves domains unserved.
120    pub ingress: Option<crate::ingress::IngressConfig>,
121    /// The port each org's workspace reaches the org-bound MCP on, on the
122    /// org bridge's address.
123    pub workspace_mcp_port: u16,
124    /// `--workspace-pool`: the storage pool new workspace homes go in,
125    /// unless the org sets its own; none: the org's default pool.
126    pub workspace_pool: Option<String>,
127    /// `--workspace-home-root`: workspace homes are host folders
128    /// `<root>/<org>/home` instead of managed volumes.
129    pub workspace_home_root: Option<PathBuf>,
130    /// `--preview-domain`: where workspace ports' previews get their origins.
131    pub preview_domain: Option<workspaces::PreviewBase>,
132    /// How long audit rows are kept.
133    pub audit_retention: Duration,
134    /// Record read-only tool calls too (secret reads always are).
135    pub audit_all: bool,
136    /// How long, and how many, history rows are kept.
137    pub history_retention: Duration,
138    pub history_max_rows: i64,
139    /// `--superadmin-tailnet`: tailnet logins and tags with the unix
140    /// socket's reach.
141    pub superadmin_tailnet: Option<crate::server::tailnet::AllowList>,
142    /// `--superadmin-access`: Access emails and service token client ids
143    /// with the unix socket's reach.
144    pub superadmin_access: Option<superadmin::AccessAllowList>,
145    /// `ISB_DEV_SUPERADMIN` ([`crate::auth::dev`]; debug builds only):
146    /// every loopback request with no credential is this superadmin.
147    pub dev_superadmin: Option<String>,
148    /// `--heartbeat-url`: a dead man's switch pinged every interval.
149    pub heartbeat: Option<crate::monitor::heartbeat::Heartbeat>,
150    /// `--egress-pin NAME=IP[:PORT]`: names the egress proxy connects to
151    /// at a fixed address instead of resolving.
152    pub egress_pins: Vec<String>,
153    /// `--egress-ca FILE`: roots the egress proxy trusts besides the system's.
154    pub egress_ca: Vec<PathBuf>,
155}
156
157/// The identity endpoints over `<state>/isb.db`, and the web UI. Provider
158/// client secrets not in the environment are read from the default org's
159/// secrets.
160fn auth_routes(
161    cfg: &ServeConfig,
162    store: Arc<AuthStore>,
163    secrets: &Arc<crate::secrets::Secrets>,
164    log: &Arc<crate::audit::AuditLog>,
165    gate: Arc<superadmin::Gate>,
166    orgs: crate::auth::ops::OrgsFn,
167) -> Result<crate::server::Routes> {
168    use crate::auth::oauth::SecretFn;
169    let default_org = crate::org::OrgId::default_org();
170    let lookup = |name: &str| -> Option<SecretFn> {
171        secrets.inspect(&default_org, name).ok()?;
172        let (s, org, name) = (secrets.clone(), default_org.clone(), name.to_string());
173        Some(Arc::new(move || {
174            let (v, _) = s.get(&org, &name).map_err(|e| e.to_string())?;
175            String::from_utf8(v)
176                .map(|v| v.trim().to_string())
177                .map_err(|_| "the secret is not UTF-8".to_string())
178        }))
179    };
180    let (providers, notes) = cfg.oauth.providers(&lookup);
181    for n in notes {
182        eprintln!("isb serve: {n}");
183    }
184    let path = crate::auth::db_path(&cfg.state_dir);
185    let agent_ways = gate.agent_ways().with_public_url(cfg.public_url.clone());
186    let api = AuthApi::new(
187        store.clone(),
188        ApiConfig {
189            agent: Some(gate.agent_fn()),
190            agent_ways,
191            public_url: cfg.public_url.clone(),
192            notifier: None,
193            setup_token_file: Some(cfg.state_dir.join("setup-token")),
194            edge: Some(gate.edge_fn()),
195            providers,
196            open_signup: cfg.open_signup,
197            audit: Some(log.clone()),
198            superadmin: Some(Arc::new(move |r: &crate::server::http::Request| match gate
199                .resolve(r, None)
200            {
201                superadmin::Resolved::Superadmin(s) => Some(s),
202                _ => None,
203            })),
204            orgs: Some(orgs),
205        },
206    )?;
207    eprintln!("isb serve: identity store {}", path.display());
208    if !crate::web::BUILT {
209        eprintln!(
210            "isb serve: this binary was built without the web UI (a placeholder page is served)"
211        );
212    }
213    // The identity endpoints first, the audit tail, then the web UI, which
214    // answers every other non-API GET.
215    let (auth, web) = (Arc::new(api).router(), crate::web::routes());
216    let tail = audit::stream_route(log.clone(), store);
217    Ok(Arc::new(move |r| {
218        auth(r).or_else(|| tail(r)).or_else(|| web(r))
219    }))
220}
221
222struct Daemon {
223    client: Client,
224    ctl: Controller,
225    policy: RemotePolicy,
226    state_dir: PathBuf,
227    secrets: Arc<crate::secrets::Secrets>,
228    apps: crate::app::Apps,
229    ingress: Option<Arc<crate::ingress::Manager>>,
230    /// The identity store: the org list memberships hang off.
231    users: Arc<AuthStore>,
232    notifier: crate::notify::Notifier,
233    monitors: crate::monitor::Monitors,
234    history: crate::metrics_history::History,
235    /// Databases' backups and scheduled jobs.
236    data: data::Ctx,
237    /// Named volumes' snapshots and staged restores.
238    volumes: crate::volume_backup::VolumeBackups,
239    audit: Arc<crate::audit::AuditLog>,
240    /// Who is a superadmin, and what `host_policy` reports.
241    gate: Arc<superadmin::Gate>,
242    host: Value,
243    /// The template catalogs, shared by the tools and the logo route.
244    catalogs: Arc<crate::template::catalog::Catalogs>,
245    /// Each org's workspace, its token and its bridge listener; the
246    /// sandbox reaper.
247    workspaces: Arc<workspaces::Workspaces>,
248    /// Where users reach isb, for invitation links.
249    public_url: Option<String>,
250    /// One proxy per egress network (docs/guides/egress.md).
251    egress: Arc<isb_egress::Manager>,
252    /// Compose stacks' environments, managed domains and deployments.
253    meta: crate::stack::deployments::StackMeta,
254}
255
256/// Run the daemon until SIGINT/SIGTERM. Apps keep running when it stops.
257#[expect(
258    clippy::too_many_lines,
259    reason = "predates the lint ratchet; split it when next changed"
260)]
261pub fn serve(client: Client, cfg: ServeConfig) -> Result<()> {
262    client
263        .server_info()
264        .map_err(|e| Error::invalid(format!("isb serve needs incusd: {e}")))?;
265    default_org::warn_old_incus(&client);
266    let store = Store::open(&cfg.state_dir)?;
267    dns::open_dns_path(&cfg.state_dir);
268    // The default org is the incus project `isb-default`, made here when
269    // it is missing.
270    default_org::ensure(&client, &store);
271    let secrets_config = crate::secrets::SecretsConfig::load(&cfg.secrets_config)?;
272    let opened = crate::secrets::Secrets::open(&cfg.state_dir, &cfg.keys, &secrets_config)?;
273    for n in &opened.notes {
274        eprintln!("isb serve: {n}");
275    }
276    let secrets = Arc::new(with_external_drivers(opened.secrets, &cfg.state_dir)?);
277    // Definitions from before secrets were references carry values: move
278    // them into the store before anything reads the definitions.
279    for r in crate::stack::migrate::run(&store, &secrets, Some(&client)) {
280        match r {
281            Ok(m) => eprintln!("isb serve: {m}"),
282            Err(e) => eprintln!("isb serve: WARNING: {e}"),
283        }
284    }
285    // Open the identity store before anything starts, so a bad one fails
286    // startup cleanly. Its endpoints ride on the TCP listener.
287    let db = crate::auth::db_path(&cfg.state_dir);
288    let users = Arc::new(
289        AuthStore::open_with(&db, cfg.auth.clone())
290            .map_err(|e| Error::invalid(format!("open {}: {e}", db.display())))?,
291    );
292    let audit_db = crate::audit::db_path(&cfg.state_dir);
293    let audit_log = Arc::new(
294        crate::audit::AuditLog::open(&audit_db, cfg.audit_retention)?
295            .with_history_limits(cfg.history_retention, cfg.history_max_rows),
296    );
297    // The history: markers for the time nobody was watching and for this
298    // start, then incus' lifecycle events from now on.
299    let recorder = crate::history::Recorder::start(audit_log.clone());
300    let stop_history = Arc::new(std::sync::atomic::AtomicBool::new(false));
301    history_start(&audit_log, &recorder, &client, &stop_history);
302    eprintln!(
303        "isb serve: audit log {} (kept {} days{})",
304        audit_db.display(),
305        cfg.audit_retention.as_secs() / 86400,
306        if cfg.audit_all {
307            ", reads included"
308        } else {
309            ""
310        }
311    );
312    let access = match &cfg.access {
313        Some((team, aud)) => Some(Arc::new(AccessValidator::new(team, aud)?)),
314        None => None,
315    };
316    let gate = Arc::new(superadmin::gate(&cfg, users.clone(), access.clone())?);
317    let auth = if cfg.listen.is_empty() {
318        None
319    } else {
320        Some(auth_routes(
321            &cfg,
322            users.clone(),
323            &secrets,
324            &audit_log,
325            gate.clone(),
326            orgs::existing_fn(client.clone()),
327        )?)
328    };
329    // The local registry, when set up: this daemon pushes to it and keeps
330    // its push index under the state directory.
331    match crate::registry::Registry::open(&client, Some(&cfg.state_dir)) {
332        Ok(Some(r)) => {
333            eprintln!("isb serve: local registry {}", r.info().url());
334            crate::registry::install(Arc::new(r));
335        }
336        Ok(None) => {}
337        Err(e) => eprintln!("isb serve: WARNING: local registry: {e}"),
338    }
339    // The ingress follows rotation from the first replica the controller
340    // resumes, so it exists before the controller does.
341    let ingress = match &cfg.ingress {
342        Some(ic) => Some(crate::ingress::Manager::new(
343            ic.clone(),
344            client.clone(),
345            secrets.clone(),
346            &cfg.state_dir,
347        )?),
348        None => None,
349    };
350    let observer = ingress
351        .clone()
352        .map(|m| m as Arc<dyn crate::stack::controller::Observer>);
353    let ctl = Controller::start_with(
354        client.clone(),
355        store,
356        cfg.interval,
357        secrets.clone(),
358        observer,
359    )?;
360    if let Some(m) = &ingress {
361        m.start(ctl.clone())?;
362    }
363    let apps = crate::app::Apps::new(&cfg.state_dir, client.clone(), ctl.clone(), secrets.clone());
364    // Every compose stack belongs to a project environment: stacks from
365    // before that (or whose adoption failed last time) get one now. Only
366    // records change; nothing is redeployed.
367    let stacks: Vec<(crate::org::OrgId, String)> = ctl
368        .definitions()
369        .iter()
370        .map(|d| (d.org.clone(), d.name.clone()))
371        .collect();
372    for r in apps.adopt_compose_stacks(&stacks) {
373        match r {
374            Ok(m) => eprintln!("isb serve: {m}"),
375            Err(e) => eprintln!("isb serve: WARNING: {e}"),
376        }
377    }
378    // Services of those stacks are named in their environment too.
379    ctl.set_dns_scope(apps.scope_fn());
380    // Notifications follow the event feed from the start of this run.
381    let ra = apps.clone();
382    let resolve: crate::notify::Resolve = Arc::new(move |org, stack, service| {
383        let a = ra.get(org, service).ok()?;
384        // The app's own stack, or one of its previews' (`<project>-...-pr-<n>`).
385        let own = a.spec.stack().ok()? == stack;
386        let preview = stack.starts_with(&format!("{}-", a.spec.project))
387            && crate::app::preview::is_pr_suffix(stack);
388        (own || preview).then_some(a.spec.project)
389    });
390    let notifier = crate::notify::Notifier::new(&cfg.state_dir, secrets.clone(), resolve)?;
391    notifier.start(ctl.clone());
392    let monitors = monitors::start(&cfg, &apps, &secrets, &notifier);
393    // Every controller event goes to the history (the ones already emitted at startup first).
394    ctl.set_event_sink(recorder.controller_sink());
395    // Every metrics sample also goes to the history.
396    let history = crate::metrics_history::History::new(&cfg.state_dir);
397    ctl.set_metrics_sink(history.start());
398    // Previews past their TTL, and removals that did not finish.
399    apps.start_preview_upkeep();
400    // Jobs and backups share one scheduler thread.
401    let jobs = crate::jobs::Jobs::new(&cfg.state_dir, apps.clone());
402    let backups = crate::backup::Backups::new(&cfg.state_dir, apps.clone());
403    let volumes =
404        crate::volume_backup::VolumeBackups::new(&cfg.state_dir, apps.clone(), backups.clone());
405    let scheduler = crate::jobs::Scheduler::start(vec![
406        Arc::new(jobs.clone()) as Arc<dyn crate::jobs::Scheduled>,
407        Arc::new(backups.clone()),
408        Arc::new(volumes.clone()),
409    ]);
410    jobs.set_scheduler(scheduler.clone());
411    backups.set_scheduler(scheduler.clone());
412    volumes.set_scheduler(scheduler.clone());
413    if let Some(ic) = &cfg.ingress {
414        if ic.tunnel_port == cfg.workspace_mcp_port {
415            return Err(Error::invalid(format!(
416                "--workspace-mcp-port {} is the ingress's tunnel port; pick another",
417                cfg.workspace_mcp_port
418            )));
419        }
420    }
421    let workspaces = workspaces::Workspaces::new(
422        &cfg.state_dir,
423        client.clone(),
424        secrets.clone(),
425        recorder.clone(),
426        cfg.workspace_mcp_port,
427        cfg.workspace_pool.clone(),
428        cfg.workspace_home_root.clone(),
429    );
430    workspaces.previews.set_base(cfg.preview_domain.clone());
431    let (egress, stop_egress) = egress::start(&cfg, &client, &secrets)?;
432    let meta = crate::stack::deployments::StackMeta::new(&cfg.state_dir);
433    meta.recover();
434    let d = Arc::new(Daemon {
435        client,
436        ctl: ctl.clone(),
437        policy: cfg.policy.clone(),
438        state_dir: cfg.state_dir.clone(),
439        secrets,
440        apps: apps.clone(),
441        ingress: ingress.clone(),
442        users: users.clone(),
443        notifier: notifier.clone(),
444        monitors: monitors.clone(),
445        history,
446        data: data::Ctx {
447            apps: apps.clone(),
448            jobs,
449            backups,
450        },
451        volumes,
452        audit: audit_log.clone(),
453        gate: gate.clone(),
454        host: superadmin::host_summary(&cfg, &gate),
455        catalogs: Arc::new(crate::template::catalog::Catalogs::new(&cfg.state_dir)),
456        workspaces: workspaces.clone(),
457        public_url: cfg.public_url.clone(),
458        egress,
459        meta,
460    });
461    let registry = registry(d.clone())?;
462    let mut hooks = hooks(d.clone(), users.clone(), cfg.allow_unauthenticated);
463    superadmin::announce(&cfg, &gate, &users);
464    hooks.audit = Some(audit::hook(audit_log.clone(), cfg.audit_all));
465    // Webhooks carry their own credential (a signature), and come from
466    // senders that hold no session.
467    let webhooks = audit::audited_webhooks(apps::webhook_routes(apps.clone()), audit_log.clone());
468    // Template logos from isb's own cache, ahead of the web UI.
469    let auth = auth.map(|a| -> crate::server::Routes {
470        let logo = templates::logo::route(
471            d.catalogs.clone(),
472            Arc::new(templates::logo::Logos::new(&cfg.state_dir)),
473            templates::logo::admit(
474                hooks.authn.clone().expect("the daemon authenticates"),
475                access.clone(),
476                cfg.allow_unauthenticated,
477            ),
478        );
479        Arc::new(move |r| logo(r).or_else(|| a(r)))
480    });
481    let mut listeners = vec![Listener::unix(&cfg.socket).hooks(hooks.clone())];
482    for addr in &cfg.listen {
483        let tailnet = superadmin::is_tailnet_listen(addr);
484        let mut l = Listener::tcp(addr.clone())
485            .policy(cfg.remote_tools.clone())
486            .hooks(hooks.clone())
487            .public_routes(webhooks.clone())
488            .preview(workspaces::preview_route(d.clone()))
489            .tailnet(tailnet);
490        if let Some(r) = &auth {
491            l = l.routes(r.clone());
492        }
493        listeners.push(match &access {
494            // Access guards the loopback listeners (the tunnel's end).
495            Some(v) if !tailnet => l.access_shared(v.clone()),
496            // Callers sign in with an API token or a session (or are tailnet
497            // superadmins); the authorizer refuses anonymous ones unless
498            // --allow-unauthenticated.
499            _ => l.allow_unauthenticated(true),
500        });
501    }
502    let hd = d.clone();
503    let healthz: crate::server::Healthz = Arc::new(move || {
504        let stacks: Vec<Value> = hd
505            .ctl
506            .list()
507            .into_iter()
508            .map(|s| json!({"name": s.name, "converged": s.converged}))
509            .collect();
510        (
511            true,
512            json!({"ok": true, "isb": env!("CARGO_PKG_VERSION"), "stacks": stacks}),
513        )
514    });
515    // Each org's workspace reaches the org-bound surface on its bridge:
516    // the hooks and tool policy of the TCP listeners, no Access (only the
517    // org's own subnet, with bearer tokens, gets in).
518    let registry = Arc::new(registry);
519    workspaces.set_serving(
520        Listener::tcp("org-bridge")
521            .policy(cfg.remote_tools.clone())
522            .hooks(hooks.clone())
523            .allow_unauthenticated(true),
524        registry.clone(),
525        healthz.clone(),
526    );
527    workspaces.start(ctl.clone());
528    workspaces::start_ports(d.clone());
529    let r = crate::server::serve_shared(listeners, registry, healthz);
530    workspaces.shutdown();
531    stop_egress.store(true, std::sync::atomic::Ordering::Relaxed);
532    stop_history.store(true, std::sync::atomic::Ordering::Relaxed);
533    recorder.record(crate::history::marker(
534        "serve.stopped",
535        "isb serve stopped: incus events from now on are not observed".into(),
536        json!({"version": env!("CARGO_PKG_VERSION")}),
537    ));
538    recorder.shutdown();
539    notifier.shutdown();
540    monitors.shutdown();
541    scheduler.shutdown();
542    ctl.shutdown();
543    if let Some(m) = &ingress {
544        m.shutdown();
545    }
546    r
547}
548
549/// Record the gap since the history last heard anything and this start,
550/// then follow incus' lifecycle events until `stop`.
551fn history_start(
552    log: &Arc<crate::audit::AuditLog>,
553    rec: &Arc<crate::history::Recorder>,
554    client: &Client,
555    stop: &Arc<std::sync::atomic::AtomicBool>,
556) {
557    use crate::history::{HistoryQuery, marker};
558    let now = crate::audit::now_ms();
559    let last = log
560        .history_list(
561            &HistoryQuery::default(),
562            &crate::audit::Visibility::All,
563            None,
564            None,
565            1,
566        )
567        .ok()
568        .and_then(|v| v.into_iter().next());
569    if let Some(l) = last {
570        let clean = l.kind == "serve.stopped";
571        let reason = if clean {
572            "isb serve was not running"
573        } else {
574            "isb serve was not running (it did not stop cleanly)"
575        };
576        rec.record(marker(
577            "incus.gap",
578            format!(
579                "incus events between {} and {} were not observed: {reason}",
580                crate::history::fmt_ms(l.time),
581                crate::history::fmt_ms(now),
582            ),
583            json!({"from": l.time, "to": now, "reason": reason}),
584        ));
585    }
586    rec.record(marker(
587        "serve.started",
588        format!("isb serve {} started", env!("CARGO_PKG_VERSION")),
589        json!({"version": env!("CARGO_PKG_VERSION"), "pid": std::process::id()}),
590    ));
591    let (c, r, s) = (client.clone(), rec.clone(), stop.clone());
592    let _ = std::thread::Builder::new()
593        .name("isb-incus-events".into())
594        .spawn(move || crate::history::watch_incus(c, r, s));
595}
596
597/// The external secret drivers, each reading its credentials from the org's own `local` secrets.
598fn with_external_drivers(
599    secrets: crate::secrets::Secrets,
600    state_dir: &std::path::Path,
601) -> Result<crate::secrets::Secrets> {
602    use crate::secrets::{Driver, local::LocalDriver, onepassword};
603    let local = Arc::new(LocalDriver::new(state_dir, secrets.keyring().clone()));
604    let token: onepassword::TokenSource =
605        Arc::new(move |org| match local.get(org, onepassword::TOKEN_SECRET) {
606            Ok((v, _)) => Ok(Some(
607                String::from_utf8(v)
608                    .map_err(|_| Error::invalid("the 1Password token is not text"))?
609                    .trim()
610                    .to_string(),
611            )),
612            Err(e) if e.is_not_found() => Ok(None),
613            Err(e) => Err(e),
614        });
615    secrets.with_driver(Arc::new(onepassword::OnePasswordDriver::new(token)))
616}
617
618/// Authentication and authorization for every listener.
619fn hooks(d: Arc<Daemon>, users: Arc<AuthStore>, allow_anonymous: bool) -> crate::server::Hooks {
620    use crate::server::Authenticated;
621    let term = terminal::terminal(d.clone());
622    let ssh = ssh::ssh(d.clone(), users.clone());
623    let u = users.clone();
624    let gate = d.gate.clone();
625    let wsa = d.workspaces.clone();
626    let authn: crate::server::mcp::Authn = Arc::new(move |req, id| {
627        match gate.resolve(req, id) {
628            superadmin::Resolved::Superadmin(s) => return Authenticated::Superadmin(s),
629            superadmin::Resolved::Refused => return Authenticated::Refused,
630            superadmin::Resolved::None => {}
631        }
632        // An org's workspace token: judged by the workspaces, which keep it.
633        if let Some(t) = bearer(req) {
634            if t.starts_with(crate::auth::secret::TokenKind::Workspace.prefix()) {
635                return match wsa.authenticate(t) {
636                    Some(p) => Authenticated::User(Arc::new(p)),
637                    None => Authenticated::Refused,
638                };
639            }
640        }
641        if req.header("authorization").is_some()
642            || req
643                .header("cookie")
644                .is_some_and(|c| c.contains("isb_session="))
645        {
646            return match u.principal_from_request(req) {
647                Some(p) => Authenticated::User(Arc::new(p)),
648                None => Authenticated::Refused,
649            };
650        }
651        // Access vouches for the email; the isb account decides the orgs.
652        if let Some(email) = id.and_then(|i| i.email.as_deref()) {
653            if let Ok(Some(p)) = u.principal_for_email(email) {
654                return Authenticated::User(Arc::new(p));
655            }
656        }
657        // A tailnet or Access caller an org mapped to a role.
658        if let Some(p) = gate.agent(req, id) {
659            return Authenticated::User(Arc::new(p));
660        }
661        Authenticated::None
662    });
663    let authorize: crate::server::mcp::Authorize = Arc::new(move |c, tool, args, scope| {
664        // `app` for `name`, `command` for `argv`, before anything reads them.
665        let args = crate::server::aliases::alias_args(tool, args);
666        authorize_class(
667            c,
668            &tool.name,
669            audit::class_for(tool, &args),
670            args,
671            scope,
672            allow_anonymous,
673        )
674    });
675    let events: crate::server::mcp::Events = Arc::new(move |c, since| {
676        if let (Caller::Unauthenticated { .. }, false) = (c, allow_anonymous) {
677            return Err(Error::Forbidden("sign in to follow events".into()));
678        }
679        if let Caller::Access(id) = c {
680            return Err(Error::Forbidden(format!(
681                "{} has no isb account",
682                id.name()
683            )));
684        }
685        let orgs = visible_orgs(c);
686        let ctl = d.ctl.clone();
687        Ok(Box::new(move |w: &mut dyn std::io::Write| {
688            let mut since = ctl.resume_from(since);
689            loop {
690                let (seq, evs) = ctl.wait_events(since, 200, Duration::from_secs(15));
691                let mut wrote = false;
692                for e in evs {
693                    if !event_visible(&orgs, &e.stack) {
694                        continue;
695                    }
696                    let data = serde_json::to_string(&e).unwrap_or_default();
697                    write!(w, "id: {}\nevent: {}\ndata: {data}\n\n", e.seq, e.level)?;
698                    wrote = true;
699                }
700                if !wrote {
701                    // Keeps proxies from closing an idle stream.
702                    w.write_all(b": keepalive\n\n")?;
703                }
704                w.flush()?;
705                since = seq.max(since);
706            }
707        }))
708    });
709    crate::server::Hooks {
710        authn: Some(authn),
711        authorize: Some(authorize),
712        events: Some(events),
713        terminal: Some(term),
714        ssh: Some(ssh),
715        audit: None,
716        listed: Some(Arc::new(tool_listed)),
717        refuse_anonymous: !allow_anonymous,
718    }
719}
720
721/// The bearer token a request carries, if any.
722fn bearer(req: &crate::server::http::Request) -> Option<&str> {
723    let (scheme, token) = req.header("authorization")?.trim().split_once(' ')?;
724    scheme.eq_ignore_ascii_case("bearer").then(|| token.trim())
725}
726
727fn args<T: DeserializeOwned>(v: Value) -> Result<T> {
728    serde_json::from_value(v).map_err(|e| Error::invalid(format!("bad arguments: {e}")))
729}
730
731fn obj(mut props: Value, required: &[&str]) -> Value {
732    // Every tool works within one org.
733    props["org"] =
734        json!({"type": "string", "description": "The org to act in (default: default)."});
735    json!({"type": "object", "properties": props, "required": required, "additionalProperties": false})
736}
737
738/// A stack's qualified name from a tool's `org` and `name`.
739fn qname(org: &Option<String>, name: &str) -> Result<String> {
740    let org = match org {
741        Some(o) => crate::org::OrgId::new(o.clone())?,
742        None => crate::org::OrgId::default_org(),
743    };
744    Ok(crate::stack::qualified(&org, name))
745}
746
747fn caller_name(c: &Caller) -> String {
748    c.to_string()
749}
750
751/// Build the tool registry.
752fn registry(d: Arc<Daemon>) -> Result<Registry> {
753    let mut r = Registry::new().instructions(guide::INSTRUCTIONS);
754    superadmin::register(&mut r, d.clone())?;
755    host_monitor::register(&mut r, d.clone())?;
756    let ann = Ann {
757        ro: json!({"readOnlyHint": true, "openWorldHint": false}),
758        destructive: json!({"destructiveHint": true, "openWorldHint": false}),
759        write: json!({"destructiveHint": false, "openWorldHint": false}),
760    };
761
762    guide::register(&mut r, &d, &ann)?;
763    tools::stack_deploy_tool(&mut r, &d, &ann)?;
764    tools::overview_tool(&mut r, &d, &ann)?;
765    tools::events_tool(&mut r, &d, &ann)?;
766    tools::ingress_status_tool(&mut r, &d, &ann)?;
767    tools::stack_list_tool(&mut r, &d, &ann)?;
768    tools::stack_status_tool(&mut r, &d, &ann)?;
769    tools::stack_config_tool(&mut r, &d, &ann)?;
770    tools::stack_logs_tool(&mut r, &d, &ann)?;
771    tools::stack_scale_tool(&mut r, &d, &ann)?;
772    tools::stack_edit_tools(&mut r, &d, &ann)?;
773    tools::sandbox_create_tool(&mut r, &d, &ann)?;
774    secret_hooks::register(&mut r, &d)?;
775    builds::register(
776        &mut r,
777        builds::Ctx {
778            client: d.client.clone(),
779            policy: d.policy.clone(),
780            ctl: d.ctl.clone(),
781        },
782    )?;
783    tools::sandbox_tools(&mut r, &d, &ann)?;
784    apps::register(&mut r, d.apps.clone(), d.ingress.is_some())?;
785    apps::delete_tools(&mut r, &d)?;
786    previews::register(&mut r, d.apps.clone())?;
787    let mut t = templates::Templates::new(
788        &d.state_dir,
789        d.apps.clone(),
790        d.secrets.clone(),
791        d.ingress.as_ref().and_then(|m| m.public_ip()),
792    );
793    t.catalogs = d.catalogs.clone();
794    templates::register(&mut r, t)?;
795    data::register(&mut r, d.data.clone())?;
796    volumes::register(&mut r, d.volumes.clone())?;
797    tools::server_status_tool(&mut r, &d, &ann)?;
798    orgs::register(&mut r, d.clone())?;
799    notify::register(
800        &mut r,
801        d.notifier.clone(),
802        d.history.clone(),
803        d.apps.clone(),
804    )?;
805    monitors::register(&mut r, d.monitors.clone())?;
806    audit::register(&mut r, d.audit.clone())?;
807    audit::register_history(&mut r, d.audit.clone())?;
808    ssh::register(&mut r, d.clone())?;
809    workspaces::register(&mut r, d.clone())?;
810    accounts::register(&mut r, d.clone())?;
811    Ok(r)
812}
813
814impl Daemon {
815    /// May this caller touch this instance? Local callers: always. Remote:
816    /// only what isb serve manages, unless the operator allowed any.
817    fn reachable(&self, c: &Caller, i: &SandboxInfo) -> bool {
818        // A signed-in user reaches everything in an org they belong to (the
819        // authorizer already checked the org): the org is the boundary.
820        c.is_trusted()
821            || c.principal().is_some()
822            || self.policy.any_instance
823            || i.config.contains_key("user.isb.stack")
824            || i.config.contains_key(&format!("user.{LABEL_OWNER}"))
825    }
826
827    /// A client on the org a tool call names (default: the default org),
828    /// refused up front when that org does not exist.
829    fn oc(&self, org: &Option<String>) -> Result<Client> {
830        let org = crate::org::OrgId::new(org.as_deref().unwrap_or(crate::org::DEFAULT_ORG))?;
831        crate::org::check_exists(&self.client, &org)?;
832        Ok(crate::org::client(&self.client, &org))
833    }
834
835    fn reach(&self, c: &Caller, oc: &Client, name: &str) -> Result<SandboxInfo> {
836        let info = Sandbox::get(oc, name)?.info()?;
837        if !self.reachable(c, &info) {
838            // Indistinguishable from absent, so a remote caller cannot map
839            // the host's other instances.
840            return Err(Error::NotFound(format!("sandbox {name}")));
841        }
842        Ok(info)
843    }
844
845    /// Where a remote stack's relative paths resolve by default.
846    /// An org's workspace definition, by name.
847    fn workspaces_def(
848        &self,
849        org: &crate::org::OrgId,
850        name: &str,
851    ) -> Result<crate::workspace::Workspace> {
852        crate::workspace::Store::new(&self.state_dir)
853            .get(org, name)?
854            .ok_or_else(|| Error::NotFound(format!("org {org} has no workspace {name}")))
855    }
856
857    fn files_dir(&self, stack: &str) -> Result<PathBuf> {
858        let p = self.state_dir.join("files").join(stack);
859        std::fs::create_dir_all(&p)?;
860        Ok(p)
861    }
862}
863
864/// Poll until every service is converged, or one is paused or failing (its
865/// message says why), or `timeout`.
866pub fn wait_settled(
867    ctl: &Controller,
868    name: &str,
869    timeout: Duration,
870) -> Result<crate::stack::controller::StackStatus> {
871    let started = Instant::now();
872    let def = ctl.definition(name)?;
873    loop {
874        let st = ctl.status(name)?;
875        // A status from before the worker saw this deployment has an older
876        // revision or replica count; only one that matches counts.
877        let settled = st.services.iter().all(|s| {
878            let current = def.revision(&s.service).is_ok_and(|r| r == s.rev)
879                && def
880                    .service(&s.service)
881                    .is_ok_and(|d| d.replicas() == s.replicas);
882            current && matches!(s.state.as_str(), "converged" | "paused" | "failing")
883        });
884        if settled || started.elapsed() >= timeout {
885            return Ok(st);
886        }
887        std::thread::sleep(Duration::from_secs(1));
888    }
889}
890
891#[derive(Deserialize)]
892#[serde(untagged)]
893enum SpecArg {
894    Text(String),
895    Object(Box<SandboxSpec>),
896}
897
898#[expect(
899    clippy::too_many_lines,
900    reason = "predates the lint ratchet; split it when next changed"
901)]
902fn sandbox_create(d: &Daemon, a: Value, c: &Caller) -> Result<Value> {
903    #[derive(Deserialize)]
904    struct A {
905        spec: Value,
906        #[serde(default)]
907        wait_ready: Option<bool>,
908        #[serde(default)]
909        expires: Option<String>,
910        #[serde(default)]
911        idle_timeout: Option<String>,
912        #[serde(default)]
913        org: Option<String>,
914    }
915    let org = arg_org(&a)?;
916    let a: A = args(a)?;
917    let mut spec = match serde_json::from_value::<SpecArg>(a.spec)
918        .map_err(|e| Error::invalid(format!("spec: {e}")))?
919    {
920        SpecArg::Text(t) => serde_yaml_ng::from_str::<SandboxSpec>(&t)
921            .map_err(|e| Error::invalid(format!("spec: {e}")))?,
922        SpecArg::Object(s) => *s,
923    };
924    let name = spec
925        .name
926        .clone()
927        .ok_or_else(|| Error::invalid("spec needs container_name"))?;
928    egress::check_secrets(&d.secrets, &org, &spec)?;
929    let base = if c.is_local() {
930        std::env::current_dir()?
931    } else {
932        d.files_dir("_sandboxes")?
933    };
934    if let Caller::Superadmin(s) = c {
935        // The socket's reach, under the superadmin's own name.
936        spec.labels.insert(LABEL_OWNER.into(), s.label());
937    }
938    if !c.is_trusted() {
939        d.policy.check_spec(&spec, &base)?;
940        if let Ok(sb) = Sandbox::get(&d.oc(&a.org)?, &name) {
941            // Reconciling someone else's instance would be taking it over.
942            if !d.reachable(c, &sb.info()?) {
943                return Err(Error::AlreadyExists(name));
944            }
945        }
946        spec.labels.insert(LABEL_OWNER.into(), owner_label(c));
947    }
948    if let Ok(sb) = Sandbox::get(&d.oc(&a.org)?, &name) {
949        let info = sb.info()?;
950        // A sandbox spec over the workspace would replace the org's machine.
951        if info.config.contains_key(crate::workspace::KEY_WORKSPACE) {
952            return Err(Error::invalid(format!(
953                "{name} is the org's workspace; pick another name"
954            )));
955        }
956    }
957    // incus counts every disk against an org's disk quota, and refuses a
958    // root disk without a size there.
959    if !spec
960        .raw_devices
961        .get("root")
962        .is_some_and(|r| r.contains_key("size"))
963        && workspaces::project_has_disk_limit(&d.client, &org)
964    {
965        spec.raw_devices
966            .entry("root".into())
967            .or_default()
968            .insert("size".into(), workspaces::SANDBOX_ROOT_SIZE.into());
969    }
970    // Short-lived by rule: the org's defaults unless the call says otherwise.
971    let settings = d.workspaces.settings(&org)?;
972    let (expires_at, idle) = crate::workspace::sandbox_deadlines(
973        &settings,
974        a.expires.as_deref(),
975        a.idle_timeout.as_deref(),
976        now_secs(),
977    )?;
978    spec.labels
979        .insert("isb.expires_at".into(), expires_at.to_string());
980    match idle {
981        Some(s) => {
982            spec.labels.insert("isb.idle_timeout".into(), s.to_string());
983        }
984        None => {
985            spec.labels.insert("isb.idle_timeout".into(), "0".into());
986        }
987    }
988    let opts = EnsureOptions {
989        wait_ready: a.wait_ready.unwrap_or(true),
990        ..Default::default()
991    };
992    let mut log: Vec<String> = Vec::new();
993    let (sb, report) = Sandbox::connect_or_create_with_base(
994        &d.oc(&a.org)?,
995        &spec,
996        &Default::default(),
997        &base,
998        opts,
999        &mut |m| log.push(m.to_string()),
1000    )?;
1001    d.workspaces.mark_active(&org.incus_project(), &name);
1002    d.egress.kick();
1003    Ok(json!({
1004        "info": sb.info()?,
1005        "report": report,
1006        "log": log,
1007        "expires_at": expires_at,
1008        "idle_timeout": idle,
1009        "message": format!(
1010            "{name} expires {} from now{}; sandbox_extend pushes it out.",
1011            crate::workspace::human(expires_at.saturating_sub(now_secs())),
1012            match idle {
1013                Some(s) => format!(" and is deleted after {} idle", crate::workspace::human(s)),
1014                None => String::new(),
1015            }
1016        ),
1017    }))
1018}
1019
1020/// What `isb.owner` says about a sandbox this caller creates.
1021fn owner_label(c: &Caller) -> String {
1022    match c {
1023        Caller::Superadmin(s) => s.label(),
1024        Caller::User { principal } if principal.is_workspace() => {
1025            crate::auth::WORKSPACE_ACTOR.to_string()
1026        }
1027        _ => format!("mcp:{}", caller_name(c)),
1028    }
1029}
1030
1031fn sandbox_extend(d: &Daemon, a: Value, c: &Caller) -> Result<Value> {
1032    #[derive(Deserialize)]
1033    #[serde(deny_unknown_fields)]
1034    struct A {
1035        name: String,
1036        #[serde(default)]
1037        by: Option<String>,
1038        #[serde(default)]
1039        idle_timeout: Option<String>,
1040        #[serde(default)]
1041        org: Option<String>,
1042    }
1043    let org = arg_org(&a)?;
1044    let a: A = args(a)?;
1045    let oc = d.oc(&a.org)?;
1046    let info = d.reach(c, &oc, &a.name)?;
1047    let labels: BTreeMap<String, String> = info
1048        .config
1049        .iter()
1050        .filter_map(|(k, v)| k.strip_prefix("user.").map(|k| (k.to_string(), v.clone())))
1051        .collect();
1052    if crate::workspace::kind_of(&labels) != "sandbox" {
1053        return Err(Error::invalid(format!(
1054            "{} is a {}, not a sandbox: it does not expire",
1055            a.name,
1056            crate::workspace::kind_of(&labels)
1057        )));
1058    }
1059    // Its creator, or the org's admins.
1060    let mine = labels
1061        .get("isb.owner")
1062        .is_some_and(|o| *o == owner_label(c));
1063    let admin = match c {
1064        Caller::Local { .. } | Caller::Superadmin(_) => true,
1065        Caller::User { principal } => {
1066            principal.platform_admin
1067                || principal
1068                    .role_in(&org)
1069                    .is_some_and(|r| r >= crate::auth::Role::Admin)
1070        }
1071        _ => false,
1072    };
1073    if !mine && !admin {
1074        return Err(Error::Forbidden(format!(
1075            "{} was created by {}; its creator or the org's admins extend it",
1076            a.name,
1077            labels
1078                .get("isb.owner")
1079                .map(String::as_str)
1080                .unwrap_or("someone else")
1081        )));
1082    }
1083    let now = now_secs();
1084    let mut patch = serde_json::Map::new();
1085    let current = labels
1086        .get("isb.expires_at")
1087        .and_then(|v| v.parse::<u64>().ok());
1088    let by = match &a.by {
1089        Some(b) => crate::flex::parse_duration(b).map_err(Error::invalid)?,
1090        None if a.idle_timeout.is_some() => Duration::ZERO,
1091        None => Duration::from_secs(86400),
1092    };
1093    let mut expires_at = current;
1094    if !by.is_zero() {
1095        let e = crate::workspace::extended(current, by, now)?;
1096        patch.insert(
1097            crate::workspace::KEY_EXPIRES_AT.into(),
1098            json!(e.to_string()),
1099        );
1100        expires_at = Some(e);
1101    }
1102    let mut idle = labels
1103        .get("isb.idle_timeout")
1104        .and_then(|v| v.parse::<u64>().ok())
1105        .filter(|s| *s > 0);
1106    if let Some(t) = &a.idle_timeout {
1107        idle = crate::workspace::idle(t)?.map(|d| d.as_secs());
1108        patch.insert(
1109            crate::workspace::KEY_IDLE_TIMEOUT.into(),
1110            json!(idle.unwrap_or(0).to_string()),
1111        );
1112    }
1113    oc.mutate(
1114        "PATCH",
1115        &format!("/1.0/instances/{}", crate::client::encode_segment(&a.name)),
1116        Some(&json!({"config": patch})),
1117        &format!("extend sandbox {}", a.name),
1118        oc.timeouts.other,
1119    )?;
1120    d.workspaces.mark_active(&org.incus_project(), &a.name);
1121    Ok(json!({
1122        "name": a.name,
1123        "expires_at": expires_at,
1124        "idle_timeout": idle,
1125        "message": format!(
1126            "{} now expires {} from now.",
1127            a.name,
1128            crate::workspace::human(expires_at.unwrap_or(now).saturating_sub(now))
1129        ),
1130    }))
1131}
1132
1133const OUTPUT_CAP: usize = 256 * 1024;
1134
1135fn cap(b: &[u8]) -> (String, bool) {
1136    if b.len() <= OUTPUT_CAP {
1137        return (String::from_utf8_lossy(b).into_owned(), false);
1138    }
1139    (
1140        String::from_utf8_lossy(&b[b.len() - OUTPUT_CAP..]).into_owned(),
1141        true,
1142    )
1143}
1144
1145fn sandbox_exec(d: &Daemon, a: Value, c: &Caller) -> Result<Value> {
1146    #[derive(Deserialize)]
1147    struct A {
1148        name: String,
1149        #[serde(default)]
1150        org: Option<String>,
1151        argv: Vec<String>,
1152        cwd: Option<String>,
1153        user: Option<String>,
1154        #[serde(default)]
1155        env: BTreeMap<String, String>,
1156        stdin: Option<String>,
1157        timeout: Option<String>,
1158    }
1159    let org = arg_org(&a)?;
1160    let a: A = args(a)?;
1161    let oc = d.oc(&a.org)?;
1162    d.reach(c, &oc, &a.name)?;
1163    d.workspaces.mark_active(&org.incus_project(), &a.name);
1164    let timeout = match &a.timeout {
1165        Some(t) => crate::flex::parse_duration(t).map_err(Error::invalid)?,
1166        None => Duration::from_secs(600),
1167    };
1168    let mut opts = ExecOptions::default().timeout(timeout);
1169    opts.cwd = a.cwd;
1170    opts.user = a.user;
1171    opts.env = a.env;
1172    if let Some(s) = a.stdin {
1173        opts.stdin = Stdin::Bytes(s.into_bytes());
1174    }
1175    let sb = Sandbox::get(&oc, &a.name)?;
1176    let out = match sb.exec_with(a.argv, opts) {
1177        Err(Error::ExecTimeout { .. }) => {
1178            return Err(Error::invalid(format!(
1179                "timed out after {timeout:?} and was killed"
1180            )));
1181        }
1182        r => r?,
1183    };
1184    let (stdout, t1) = cap(&out.stdout);
1185    let (stderr, t2) = cap(&out.stderr);
1186    Ok(
1187        json!({"exit_code": out.exit_code, "stdout": stdout, "stderr": stderr, "truncated": t1 || t2}),
1188    )
1189}
1190
1191pub use crate::stack::local_deploy_args;
1192
1193/// Default state directory, exported for the CLI.
1194pub fn default_state_dir() -> PathBuf {
1195    Store::default_dir()
1196}
1197
1198#[cfg(test)]
1199#[path = "agent_tests.rs"]
1200mod agent_tests;
1201
1202#[cfg(test)]
1203mod tests;
1204
1205#[cfg(test)]
1206mod downscope_tests;