//! Remote servers in the daemon (docs/guides/servers.md).
//!
//! On a **control plane**: the `server_*` tools, and the route that sends
//! every call for an org placed on a server to that server's agent (after
//! this daemon's authentication, authorization and audit), merges
//! cross-org reads across servers, and forwards app webhooks.
//!
//! On an **agent** (`isb serve --agent`): the mTLS listener the control
//! plane calls, its assertion of the caller, the agent's own check that a
//! call is for an org placed on it, and the internal routes (heartbeat,
//! placement, certificate rotation).
use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use serde::Deserialize;
use serde_json::{Value, json};
use super::accounts::TOOLS as ACCOUNTS;
use super::{Daemon, arg_org, args, read_orgs, visible_orgs};
use crate::error::{Error, Result};
use crate::org::OrgId;
use crate::server::http::{Request, Response};
use crate::server::{Caller, Registry, Routes, Tool};
#[cfg(test)]
use crate::server::{Hooks, Listener};
use crate::servers::provision::{Kind, Provision};
use crate::servers::wire::Assertion;
use crate::servers::{Servers, bootstrap, health, merge, vm};
mod agent;
mod upgrade;
pub(super) use agent::agent_listener;
#[cfg(test)]
use agent::{AgentState, agent_hooks, org_bound, rotate};
/// How long one server may take to answer a merged read.
const FAN_OUT_TIMEOUT: Duration = Duration::from_secs(15);
/// Tools a control plane always runs itself, whatever org they name (and ACCOUNTS).
const LOCAL_ONLY: &[&str] = &[
"events",
"audit_list",
"audit_verify",
"registry_gc",
"notification_settings",
"template_catalog_add",
"template_catalog_remove",
"template_catalog_list",
"template_list",
"template_get",
"server_status",
];
impl Daemon {
/// The server `org` lives on, when this daemon is a control plane and
/// the org is not local.
pub(super) fn remote(&self, org: &OrgId) -> Option<(&Arc<Servers>, String)> {
let s = self.servers.as_ref()?;
s.placement(org).map(|p| (s, p))
}
}
fn platform_only(c: &Caller) -> Result<()> {
match c {
Caller::Local { .. } | Caller::Superadmin(_) => Ok(()),
Caller::User { principal } if principal.platform_admin => Ok(()),
_ => Err(Error::Forbidden("servers are for platform admins".into())),
}
}
fn servers(d: &Daemon) -> Result<&Arc<Servers>> {
d.servers.as_ref().ok_or_else(|| {
Error::invalid("this daemon is an agent: servers are managed on its control plane")
})
}
#[expect(
clippy::too_many_lines,
reason = "predates the lint ratchet; split it when next changed"
)]
pub(super) fn register(r: &mut Registry, d: Arc<Daemon>) -> Result<()> {
let ro = json!({"readOnlyHint": true, "openWorldHint": false});
let destructive = json!({"destructiveHint": true, "openWorldHint": false});
let write = json!({"destructiveHint": false, "openWorldHint": true});
let name_only = json!({"type": "object", "properties": {"name": {"type": "string"}}, "required": ["name"], "additionalProperties": false});
macro_rules! tool {
($name:expr, $title:expr, $desc:expr, $schema:expr, $ann:expr, $f:expr) => {{
let d = d.clone();
let f = $f;
r.register(
Tool::new($name, $desc, $schema, move |a, c| f(&d, a, c))
.title($title)
.annotations($ann.clone()),
)?;
}};
}
tool!(
"server_add",
"Add a server",
"Platform admins: make a Linux box (Ubuntu/Debian, x86_64 or aarch64) a server orgs can be placed on. Over SSH (root or passwordless sudo) it installs incus from Zabbly's stable channel and the isb binary (checksum checked), runs `isb host setup`, issues the agent a certificate from this control plane's CA and starts `isb serve --agent` as a systemd unit; with allow_from it closes the box's firewall to SSH and the agent port from those addresses. The SSH key is used for this only. Takes minutes on a fresh box.",
json!({"type": "object", "properties": {
"name": {"type": "string", "description": "[a-z0-9-], a letter first."},
"ssh": {"type": "string", "description": "user@host"},
"ssh_port": {"type": "integer", "minimum": 1, "maximum": 65535},
"key": {"type": "string", "description": "Private key file on this host (local CLI only)."},
"ssh_key": {"type": "string", "description": "The private key itself (kept only for the bootstrap)."},
"address": {"type": "string", "description": "What this control plane dials (default: the SSH host)."},
"agent_port": {"type": "integer", "minimum": 1, "maximum": 65535, "description": "The agent's mTLS port (default 7443)."},
"allow_from": {"type": "array", "items": {"type": "string"}, "description": "Addresses or CIDRs that may reach the agent port (this control plane's egress address); the box's firewall then allows only SSH and these."},
"isb_binary": {"type": "string", "description": "A Linux isb binary on this host to install (local CLI only); default the release of this version."},
"version": {"type": "string", "description": "The isb release to install (default this daemon's)."},
"self_binary": {"type": "boolean", "description": "Install this control plane's own isb executable instead of a release (same version and build; the box must have the same architecture)."},
"public_ingress": {"type": "boolean", "description": "Serve the server's orgs' domains on its own ports 80 and 443 (opened in its firewall)."},
"wait": {"type": "boolean", "description": "Wait for the bootstrap to finish (default true). false answers at once with `provision`; follow it with server_provision_get."}
}, "required": ["name", "ssh"], "additionalProperties": false}),
write,
|d: &Daemon, a: Value, c: &Caller| -> Result<Value> {
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct A {
name: String,
ssh: String,
ssh_port: Option<u16>,
key: Option<PathBuf>,
ssh_key: Option<String>,
address: Option<String>,
agent_port: Option<u16>,
#[serde(default)]
allow_from: Vec<String>,
isb_binary: Option<PathBuf>,
version: Option<String>,
#[serde(default)]
self_binary: bool,
#[serde(default)]
public_ingress: bool,
wait: Option<bool>,
}
platform_only(c)?;
let s = servers(d)?;
let a: A = args(a)?;
if !c.is_trusted() && (a.key.is_some() || a.isb_binary.is_some()) {
return Err(Error::Forbidden(
"key and isb_binary are paths on the control plane's host: local CLI only (send ssh_key)".into(),
));
}
if a.self_binary && (a.isb_binary.is_some() || a.version.is_some()) {
return Err(Error::invalid(
"self_binary installs this daemon's own build: drop isb_binary and version",
));
}
bootstrap::validate_name(&a.name)?;
let request = json!({
"name": a.name, "ssh": a.ssh, "ssh_port": a.ssh_port.unwrap_or(22),
"address": a.address, "agent_port": a.agent_port.unwrap_or(bootstrap::DEFAULT_AGENT_PORT),
"allow_from": a.allow_from, "version": a.version, "self_binary": a.self_binary,
"public_ingress": a.public_ingress,
});
let mut o = bootstrap::AddOptions {
name: a.name.clone(),
ssh: a.ssh,
ssh_port: a.ssh_port.unwrap_or(22),
key: PathBuf::new(),
address: a.address,
agent_port: a.agent_port.unwrap_or(bootstrap::DEFAULT_AGENT_PORT),
allow_from: a.allow_from,
isb_binary: a.isb_binary,
version: a.version,
self_binary: a.self_binary,
public_ingress: a.public_ingress,
};
s.check_add(&o)?;
let p = s
.runs
.begin(Provision::new(&o.name, Kind::Ssh, None, request))?;
let mut temp_key = None;
o.key = match (a.key, a.ssh_key) {
(Some(k), None) => k,
(None, Some(text)) => {
let path = d
.state_dir
.join("servers")
.join(format!("bootstrap-key-{}", a.name));
let r = std::fs::create_dir_all(d.state_dir.join("servers"))
.map_err(Error::from)
.and_then(|()| {
let mut text = text;
if !text.ends_with('\n') {
text.push('\n');
}
crate::servers::pki::write_private(&path, &text)
});
if let Err(e) = r {
p.finish(Err(e.to_string()));
return Err(e);
}
temp_key = Some(path.clone());
path
}
_ => {
let e = Error::invalid("pass exactly one of key or ssh_key");
p.finish(Err(e.to_string()));
return Err(e);
}
};
let run = {
let (s, ctl, p) = (s.clone(), d.ctl.clone(), p.clone());
move || -> Result<Value> {
let r = s.add(&o, Some(&ctl), &p);
if let Some(k) = &temp_key {
let _ = std::fs::remove_file(k);
}
let rec = r?;
if o.allow_from.is_empty() {
p.log(&format!(
"the agent port {} is open to any address (mTLS still required); pass allow_from to firewall it",
o.agent_port
));
}
Ok(s.view(&rec))
}
};
if a.wait == Some(false) {
spawn_run(p.clone(), run);
return Ok(json!({"provision": p.view()}));
}
let r = run();
p.finish(r.as_ref().map(Value::clone).map_err(|e| e.to_string()));
let mut v = r?;
v["log"] = json!(p.view().log);
Ok(v)
}
);
tool!(
"server_list",
"List servers",
"Platform admins: the servers orgs can be placed on, with their health (up, unreachable, unknown), last heartbeat (versions, CPU, memory, disk) and the orgs on each; `provisions`, servers being added (and recent failures); `dedicated_vm`, whether this host can run dedicated VMs; `suggested_allow_from`, addresses this control plane's traffic leaves from.",
json!({"type": "object", "properties": {}, "additionalProperties": false}),
ro,
|d: &Daemon, _a: Value, c: &Caller| -> Result<Value> {
platform_only(c)?;
let s = servers(d)?;
let v: Vec<Value> = s.records().iter().map(|r| s.view(r)).collect();
let vm = match vm::support(&d.client) {
Ok(()) => json!({"supported": true}),
Err(reason) => json!({"supported": false, "reason": reason}),
};
Ok(json!({
"servers": v,
"provisions": s.runs.list(),
"dedicated_vm": vm,
"suggested_allow_from": suggested_allow_from(),
}))
}
);
tool!(
"server_provision_get",
"Follow a server being added",
"Platform admins: how far adding a server (server_add with wait=false) or making an org's dedicated VM (org_create with placement vm) got: its steps, log, state (running, done, failed) and error. Kept for an hour after a failure.",
name_only.clone(),
ro,
|d: &Daemon, a: Value, c: &Caller| -> Result<Value> {
#[derive(Deserialize)]
struct A {
name: String,
}
platform_only(c)?;
let a: A = args(a)?;
let s = servers(d)?;
s.runs
.get(&a.name)
.map(|v| serde_json::to_value(v).unwrap_or_default())
.ok_or_else(|| Error::NotFound(format!("no server {} being added", a.name)))
}
);
tool!(
"server_show",
"Show a server",
"Platform admins: one server: address, port, how it was bootstrapped, its certificate's fingerprint and expiry, health with the last heartbeat, and its orgs.",
name_only.clone(),
ro,
|d: &Daemon, a: Value, c: &Caller| -> Result<Value> {
#[derive(Deserialize)]
struct A {
name: String,
}
platform_only(c)?;
let s = servers(d)?;
let a: A = args(a)?;
Ok(s.view(&s.record(&a.name)?))
}
);
tool!(
"server_remove",
"Remove a server",
"Platform admins: forget a server. Refused while orgs are placed on it (delete them first). A box added over SSH keeps running its agent until it is stopped there (systemctl disable --now isb-agent); a dedicated VM this control plane made is deleted with it.",
name_only.clone(),
destructive,
|d: &Daemon, a: Value, c: &Caller| -> Result<Value> {
#[derive(Deserialize)]
struct A {
name: String,
}
platform_only(c)?;
let a: A = args(a)?;
remove_server(d, servers(d)?, &a.name)
}
);
tool!(
"server_rotate_cert",
"Rotate a server's certificate",
"Platform admins: issue the server's agent a new certificate (and key) over the current mTLS connection; the agent switches to it for new connections, and the control plane checks it does.",
name_only,
write,
|d: &Daemon, a: Value, c: &Caller| -> Result<Value> {
#[derive(Deserialize)]
struct A {
name: String,
}
platform_only(c)?;
let a: A = args(a)?;
let s = servers(d)?;
Ok(s.view(&s.rotate_cert(&a.name)?))
}
);
upgrade::register(r, d)
}
/// Run a provisioning in the background, recording how it ends.
fn spawn_run(p: Provision, run: impl FnOnce() -> Result<Value> + Send + 'static) {
let name = p.view().name;
let p2 = p.clone();
let spawned = std::thread::Builder::new()
.name(format!("isb-provision-{name}"))
.spawn(move || {
let r = run();
p2.finish(r.map_err(|e| e.to_string()));
});
if let Err(e) = spawned {
p.finish(Err(format!("start the provisioning thread: {e}")));
}
}
/// Forget a server; a dedicated VM is deleted first (refused while orgs
/// are placed on it, like any server).
fn remove_server(d: &Daemon, s: &Arc<Servers>, name: &str) -> Result<Value> {
let rec = s.record(name)?;
let orgs = s.orgs_on(name);
if !orgs.is_empty() {
// The same refusal remove() gives, before a VM is touched.
return s.remove(name).map(|_| Value::Null);
}
if let Some(v) = &rec.vm {
vm::delete(&d.client, &v.project, &v.instance)?;
s.remove(name)?;
return Ok(json!({"ok": true, "removed": rec.name, "note": format!(
"deleted VM {} in project {}", v.instance, v.project
)}));
}
let r = s.remove(name)?;
Ok(json!({"ok": true, "removed": r.name, "note": format!(
"the agent still runs on {}: stop it there with `systemctl disable --now {}`", r.ssh, bootstrap::AGENT_UNIT
)}))
}
/// Addresses this control plane's traffic leaves from, for a new server's
/// `allow_from`: toward the internet, and toward the tailnet when there is
/// one. Behind NAT a box sees the public address instead.
fn suggested_allow_from() -> Vec<Value> {
let from = |dest: &str| -> Option<std::net::IpAddr> {
let s = std::net::UdpSocket::bind("0.0.0.0:0").ok()?;
s.connect(dest).ok()?;
let ip = s.local_addr().ok()?.ip();
(!ip.is_unspecified() && !ip.is_loopback()).then_some(ip)
};
let mut out: Vec<Value> = Vec::new();
if let Some(ip) = from("1.1.1.1:53") {
let private = matches!(ip, std::net::IpAddr::V4(v) if v.is_private());
out.push(json!({
"address": ip.to_string(),
"via": if private { "default route (private: behind NAT, a box on the internet sees your public address instead)" } else { "default route" },
}));
}
if let Some(ip @ std::net::IpAddr::V4(v)) = from("100.100.100.100:53") {
if v.octets()[0] == 100
&& (64..128).contains(&v.octets()[1])
&& !out.iter().any(|o| o["address"] == ip.to_string())
{
out.push(json!({"address": ip.to_string(), "via": "tailnet"}));
}
}
out
}
/// A call as forwarded: a compose file the local CLI resolved goes as YAML
/// text, since the agent holds the control plane's callers to the remote
/// rules; the local CLI's host paths mean nothing there.
pub(super) fn forwarded_args(tool: &str, mut a: Value, c: &Caller) -> Result<Value> {
if let Some(o) = a.as_object_mut() {
if tool == "stack_deploy" {
if let Some(f) = o.remove("file") {
let yaml = serde_yaml_ng::to_string(&f)
.map_err(|e| Error::invalid(format!("compose file: {e}")))?;
// Already interpolated: keep `$` literal through the
// agent's interpolation.
o.insert("compose".into(), json!(yaml.replace('$', "$$")));
}
}
if c.is_trusted() {
o.remove("base_dir");
}
}
Ok(a)
}
/// Where a control plane runs a call.
#[derive(Debug, PartialEq, Eq)]
pub(super) enum Way {
Here,
/// To this server, in this org's endpoint.
Forward(String, OrgId),
/// Here and on the servers, merged.
FanOut,
/// Creating an org on this server.
OrgCreate(String),
/// Creating an org in a dedicated VM this control plane makes.
OrgCreateVm(vm::VmSize),
/// `org_get`, `org_update`, `org_delete`: forwarded when the org is on
/// a server (and a move is refused).
OrgOther,
}
pub(super) fn decide(name: &str, a: &Value, placement: &dyn Fn(&OrgId) -> Option<String>) -> Way {
let all = name == "secret_reencrypt" && a.get("all").and_then(Value::as_bool) == Some(true);
if name.starts_with("server_") || LOCAL_ONLY.contains(&name) || all || ACCOUNTS.contains(&name)
{
return Way::Here;
}
match name {
"org_create" => {
// A bad placement is the local tool's to refuse.
return match vm::placement(a) {
Ok(vm::Placement::Server(s)) => Way::OrgCreate(s),
Ok(vm::Placement::Vm(size)) => Way::OrgCreateVm(size),
Ok(vm::Placement::Local) | Err(_) => Way::Here,
};
}
"org_update" | "org_delete" | "org_get" => return Way::OrgOther,
_ => {}
}
if merge::FAN_OUT.contains(&name) {
return Way::FanOut;
}
// A bad org name fails in the tool, here.
match arg_org(a) {
Ok(org) => match placement(&org) {
Some(s) => Way::Forward(s, org),
None => Way::Here,
},
Err(_) => Way::Here,
}
}
/// The control plane's router: see the module docs.
pub(super) fn route(d: Arc<Daemon>) -> crate::server::mcp::Route {
Arc::new(move |tool, a, caller, origin| {
let s = d.servers.as_ref()?;
let name = tool.name.as_str();
let rid = origin.request_id.as_deref();
match decide(name, a, &|o| s.placement(o)) {
Way::Here => None,
Way::OrgCreate(server) => Some(org_create(&d, s, &server, a, caller, rid)),
Way::OrgCreateVm(size) => Some(org_create_vm(&d, s, size, a, caller, rid)),
Way::OrgOther => org_other(&d, s, name, a, caller, rid),
Way::FanOut => Some(fan_out(&d, s, tool, a, caller, rid)),
Way::Forward(server, org) => Some(
forwarded_args(name, a.clone(), caller)
.and_then(|a| s.call(&server, name, &a, caller, Some(&org), rid)),
),
}
})
}
/// The arguments an agent takes: where the org runs is the control
/// plane's business.
fn without_server(a: &Value) -> Value {
let mut a = a.clone();
if let Some(o) = a.as_object_mut() {
for k in ["server", "placement", "wait", "delete_vm"] {
o.remove(k);
}
}
a
}
/// Where an org runs and how it is kept apart, for org_get and org_list.
pub(super) fn placement_view(s: Option<&Arc<Servers>>, org: &OrgId) -> Value {
let Some(name) = s.and_then(|s| s.placement(org)) else {
return json!({"kind": "local", "server": "local", "isolation": "shared-kernel"});
};
match s.and_then(|s| s.record(&name).ok()).and_then(|r| r.vm) {
Some(v) => json!({
"kind": "vm", "server": name, "isolation": "own-kernel",
"vm": {"cpus": v.cpus, "memory": v.memory, "disk": v.disk, "project": v.project, "instance": v.instance},
}),
None => json!({"kind": "server", "server": name, "isolation": "own-host"}),
}
}
fn org_create(
d: &Daemon,
s: &Arc<Servers>,
server: &str,
a: &Value,
c: &Caller,
rid: Option<&str>,
) -> Result<Value> {
let org = new_org(d, s, a)?;
s.record(server)?;
create_on(d, s, server, &org, a, c, rid)
}
/// The org a create is for, if it is free: not the default, not placed,
/// not on this host.
fn new_org(d: &Daemon, s: &Servers, a: &Value) -> Result<OrgId> {
let org = arg_org(a)?;
if org.is_default() {
return Err(Error::invalid("the default org stays on the control plane"));
}
if s.placement(&org).is_some() {
return Err(Error::AlreadyExists(format!("org {org}")));
}
match crate::org::get(&d.client, &org) {
Ok(_) => Err(Error::AlreadyExists(format!("org {org} (on this host)"))),
Err(e) if e.is_not_found() => Ok(org),
Err(e) => Err(e),
}
}
/// Place `org` on `server` and create it there.
fn create_on(
d: &Daemon,
s: &Arc<Servers>,
server: &str,
org: &OrgId,
a: &Value,
c: &Caller,
rid: Option<&str>,
) -> Result<Value> {
s.place(org, Some(server))?;
let r = s.call(server, "org_create", &without_server(a), c, Some(org), rid);
let mut v = match r {
Ok(v) => v,
Err(e) => {
let _ = s.place(org, None);
return Err(e);
}
};
d.users
.ensure_org(org)
.map_err(|e| Error::invalid(e.to_string()))?;
v["server"] = json!(server);
v["placement"] = placement_view(Some(s), org);
Ok(v)
}
/// `org_create` with `placement: {vm: ...}`: make the org a VM of its own
/// on this host (server `vm-<org>`), then create it there. With
/// `wait: false` it answers at once with the provisioning to follow.
fn org_create_vm(
d: &Arc<Daemon>,
s: &Arc<Servers>,
size: vm::VmSize,
a: &Value,
c: &Caller,
rid: Option<&str>,
) -> Result<Value> {
let org = new_org(d, s, a)?;
// The org's own settings are checked now, not after the VM is made.
let st: super::orgs::Settings = args(without_server(a))?;
st.options(None)?;
vm::support(&d.client)
.map_err(|e| Error::invalid(format!("a dedicated VM cannot be made here: {e}")))?;
let name = vm::server_name(&org);
if let Ok(r) = s.record(&name) {
if r.vm.as_ref().is_none_or(|v| v.org != org) {
return Err(Error::AlreadyExists(format!(
"server {name} exists and is not org {org}'s dedicated VM"
)));
}
}
let mut request = without_server(a);
request["placement"] = json!({"vm": size});
let p = s
.runs
.begin(Provision::new(&name, Kind::Vm, Some(org.as_str()), request))?;
let run = {
let (d, s, p, c) = (d.clone(), s.clone(), p.clone(), c.clone());
let (a, rid) = (a.clone(), rid.map(str::to_string));
move || -> Result<Value> {
let rec = s.add_vm(&d.client, &org, &size, Some(&d.ctl), &p)?;
p.step("org");
p.log(&format!("creating org {org} on server {}", rec.name));
let v = create_on(&d, &s, &rec.name, &org, &a, &c, rid.as_deref())?;
p.log(&format!("org {org} runs in its own VM"));
Ok(v)
}
};
if a.get("wait").and_then(Value::as_bool) == Some(false) {
spawn_run(p.clone(), run);
return Ok(json!({"provision": p.view()}));
}
let r = run();
p.finish(r.as_ref().map(Value::clone).map_err(|e| e.to_string()));
r
}
fn org_other(
d: &Daemon,
s: &Arc<Servers>,
tool: &str,
a: &Value,
c: &Caller,
rid: Option<&str>,
) -> Option<Result<Value>> {
let org = arg_org(a).ok()?;
let here = s.placement(&org);
if tool == "org_update" && (a.get("server").is_some() || a.get("placement").is_some()) {
let now = here.as_deref().unwrap_or("local");
let want = match vm::placement(a) {
Ok(vm::Placement::Local) => "local".to_string(),
Ok(vm::Placement::Server(s)) => s,
Ok(vm::Placement::Vm(_)) => vm::server_name(&org),
Err(e) => return Some(Err(e)),
};
if want != now {
return Some(Err(Error::invalid(format!(
"org {org} runs on {now}; moving an org is not supported \
(docs/concepts/placement.md#moving-an-org: remove its workloads, delete it, create it \
again where it should run, restore its data)"
))));
}
}
let server = here?;
Some((|| {
let mut v = s.call(&server, tool, &without_server(a), c, Some(&org), rid)?;
match tool {
"org_delete" => {
d.users
.delete_org(&org)
.map_err(|e| Error::invalid(e.to_string()))?;
s.place(&org, None)?;
let dedicated = s
.record(&server)
.ok()
.and_then(|r| r.vm)
.filter(|v| v.org == org);
if let Some(vm) = dedicated {
if a.get("delete_vm").and_then(Value::as_bool) == Some(true) {
if s.orgs_on(&server).is_empty() {
remove_server(d, s, &server)?;
v["deleted_vm"] = json!(vm.instance);
}
} else {
v["notes"] = json!([format!(
"its VM {} keeps running as server {server}: `isb server rm {server}` deletes it",
vm.instance
)]);
}
}
}
_ => {
v["server"] = json!(server);
v["placement"] = placement_view(Some(s), &org);
v["members"] = json!(d.users.list_members(&org).map(|m| m.len()).unwrap_or(0));
}
}
Ok(v)
})())
}
/// Run a cross-org read here and on every server holding an org the
/// caller sees, and merge the answers.
#[expect(
clippy::excessive_nesting,
reason = "predates the lint ratchet; split it when next changed"
)]
fn fan_out(
d: &Daemon,
s: &Arc<Servers>,
tool: &Tool,
a: &Value,
c: &Caller,
rid: Option<&str>,
) -> Result<Value> {
let local = (tool.handler)(a.clone(), c)?;
// A read that names its org (all but org_list) asks only the servers
// holding that org, and passes the org on for them to narrow by too.
let named = if tool.name == "org_list" {
None
} else {
a.get("org")
.and_then(Value::as_str)
.map(OrgId::new)
.transpose()?
};
let visible = if named.is_some() {
read_orgs(c, a)?
} else {
visible_orgs(c)
};
let targets: Vec<String> = s
.records()
.into_iter()
.map(|r| r.name)
.filter(|n| match &visible {
None => true,
Some(v) => s.orgs_on(n).iter().any(|o| v.contains(o)),
})
.collect();
if targets.is_empty() {
return Ok(local);
}
let who = Assertion::for_caller(c)
.ok_or_else(|| Error::Forbidden(format!("{c} cannot read servers")))?;
let mut unscoped = a.clone();
if named.is_none() {
if let Some(o) = unscoped.as_object_mut() {
o.remove("org");
}
}
let results: Vec<(String, std::result::Result<Value, String>)> = std::thread::scope(|sc| {
let hs: Vec<_> = targets
.iter()
.map(|n| {
let (who, args) = (&who, &unscoped);
sc.spawn(move || {
if s.health(n).state == health::State::Unreachable {
return (n.clone(), Err("unreachable".to_string()));
}
let r = s
.client(n)
.and_then(|cl| cl.call(&tool.name, args, who, None, rid, FAN_OUT_TIMEOUT));
(n.clone(), r.map_err(|e| e.to_string()))
})
})
.collect();
hs.into_iter()
.map(|h| {
h.join()
.unwrap_or_else(|_| (String::new(), Err("panicked".into())))
})
.collect()
});
let mut results = results;
if tool.name == "org_list" {
// Only the orgs this control plane placed there, with its members.
for (n, r) in results.iter_mut() {
if let Ok(v) = r {
let placed = s.orgs_on(n);
if let Some(arr) = v.get_mut("orgs").and_then(Value::as_array_mut) {
arr.retain(|o| {
o["name"]
.as_str()
.is_some_and(|x| placed.iter().any(|p| p.as_str() == x))
});
for o in arr.iter_mut() {
if let Some(Ok(id)) = o["name"].as_str().map(OrgId::new) {
o["members"] =
json!(d.users.list_members(&id).map(|m| m.len()).unwrap_or(0));
o["placement"] = placement_view(Some(s), &id);
}
}
}
}
}
}
Ok(merge::merge(&tool.name, local, results))
}
/// App webhooks for an org on a server go to that server's agent as they
/// came (the agent checks the signature: it holds the app's secret).
pub(super) fn forward_webhooks(inner: Routes, s: Arc<Servers>) -> Routes {
Arc::new(move |req: &Request| {
let rest = req.path.strip_prefix("/api/v1/webhooks/")?;
let org = rest.split_once('/').map(|(o, _)| o)?;
let Some(server) = OrgId::new(org).ok().and_then(|o| s.placement(&o)) else {
return inner(req);
};
let keep = |k: &str| {
let k = k.to_ascii_lowercase();
!matches!(
k.as_str(),
"host" | "content-length" | "connection" | "transfer-encoding" | "authorization"
)
};
let headers: Vec<(String, String)> = req
.headers
.iter()
.filter(|(k, _)| keep(k))
.cloned()
.collect();
let path = match &req.query {
Some(q) => format!("{}?{q}", req.path),
None => req.path.clone(),
};
Some(
match s.client(&server).and_then(|c| {
c.request(
&req.method,
&path,
&headers,
&req.body,
Duration::from_secs(60),
)
}) {
Ok(a) => {
let mut r = Response::new(a.status).body(a.body);
if let Some((_, ct)) = a
.headers
.iter()
.find(|(k, _)| k.eq_ignore_ascii_case("content-type"))
{
r = r.header("Content-Type", ct.clone());
}
r.header("Cache-Control", "no-store")
}
Err(e) => {
eprintln!("isb serve: webhook for {org} on server {server}: {e}");
Response::json(
502,
&json!({"error": "bad_gateway", "message": format!("server {server} did not take the webhook")}),
)
}
},
)
})
}
#[cfg(test)]
mod tests;