Skip to main content

isb_server/servers/
mod.rs

1//! Remote servers (P5.1, P5.2): a control plane places orgs on other hosts,
2//! each running incus and `isb serve --agent`, and forwards their calls over
3//! mutual TLS. Federation, not incus clustering: every server is a whole
4//! isb (controller, ingress, registry, builds, secrets) for the orgs on it.
5//! See docs/guides/servers.md.
6
7pub mod bootstrap;
8pub mod client;
9pub mod health;
10pub mod merge;
11pub mod pki;
12pub mod provision;
13pub mod store;
14pub mod upgrade;
15pub mod vm;
16pub mod wire;
17
18use std::collections::{BTreeMap, BTreeSet};
19use std::path::{Path, PathBuf};
20use std::sync::atomic::{AtomicBool, Ordering};
21use std::sync::{Arc, Mutex};
22use std::time::Duration;
23
24use serde_json::{Value, json};
25
26use crate::error::{Error, Result};
27use crate::org::OrgId;
28use crate::stack::{Controller, now_secs};
29pub use client::AgentClient;
30pub use store::ServerRecord;
31use wire::Assertion;
32
33/// How long a forwarded tool call may take (deploys with `wait`, exec).
34pub const CALL_TIMEOUT: Duration = Duration::from_secs(20 * 60);
35
36/// The control plane's servers: their records, placement, CA, health, and
37/// the threads that watch them.
38pub struct Servers {
39    dir: PathBuf,
40    ca: pki::Ca,
41    tls: Arc<rustls::ClientConfig>,
42    store: Mutex<store::Store>,
43    health: Mutex<BTreeMap<String, health::Health>>,
44    mirrored: Mutex<BTreeSet<String>>,
45    /// Servers being added, and recently added or failed.
46    pub runs: provision::Runs,
47    /// Serializes record and placement changes (held briefly; a bootstrap
48    /// runs outside it, one per name through `runs`).
49    admin: Mutex<()>,
50    /// Servers being upgraded.
51    upgrading: Mutex<BTreeSet<String>>,
52    stop: Arc<AtomicBool>,
53}
54
55impl std::fmt::Debug for Servers {
56    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
57        f.debug_struct("Servers").field("dir", &self.dir).finish()
58    }
59}
60
61impl Servers {
62    /// `<state>/servers/`: the CA under `pki/`, the records, the placement.
63    pub fn open(state_dir: &Path) -> Result<Arc<Servers>> {
64        let dir = state_dir.join("servers");
65        std::fs::create_dir_all(&dir)?;
66        let ca = pki::Ca::open(&dir.join("pki"))?;
67        let tls = ca.client_config()?;
68        let store = store::Store::open(&dir)?;
69        Ok(Arc::new(Servers {
70            dir,
71            ca,
72            tls,
73            store: Mutex::new(store),
74            health: Mutex::new(BTreeMap::new()),
75            mirrored: Mutex::new(BTreeSet::new()),
76            runs: provision::Runs::default(),
77            admin: Mutex::new(()),
78            upgrading: Mutex::new(BTreeSet::new()),
79            stop: Arc::new(AtomicBool::new(false)),
80        }))
81    }
82
83    pub fn record(&self, name: &str) -> Result<ServerRecord> {
84        self.store
85            .lock()
86            .unwrap()
87            .servers
88            .get(name)
89            .cloned()
90            .ok_or_else(|| Error::NotFound(format!("server {name}")))
91    }
92
93    pub fn records(&self) -> Vec<ServerRecord> {
94        self.store
95            .lock()
96            .unwrap()
97            .servers
98            .values()
99            .cloned()
100            .collect()
101    }
102
103    pub fn client(&self, name: &str) -> Result<AgentClient> {
104        let r = self.record(name)?;
105        Ok(AgentClient::new(
106            &r.name,
107            &r.address,
108            r.port,
109            self.tls.clone(),
110        ))
111    }
112
113    /// The server `org` is placed on; `None` is this daemon.
114    pub fn placement(&self, org: &OrgId) -> Option<String> {
115        self.store.lock().unwrap().placement.get(org).cloned()
116    }
117
118    pub fn placements(&self) -> BTreeMap<OrgId, String> {
119        self.store.lock().unwrap().placement.clone()
120    }
121
122    pub fn orgs_on(&self, server: &str) -> Vec<OrgId> {
123        self.store.lock().unwrap().orgs_on(server)
124    }
125
126    pub fn health(&self, name: &str) -> health::Health {
127        self.health
128            .lock()
129            .unwrap()
130            .get(name)
131            .cloned()
132            .unwrap_or_default()
133    }
134
135    /// A server as `server_list` and `server_show` answer it.
136    pub fn view(&self, r: &ServerRecord) -> Value {
137        let mut v = serde_json::to_value(r).unwrap_or_default();
138        let h = self.health(&r.name);
139        v["kind"] = json!(if r.vm.is_some() { "vm" } else { "ssh" });
140        v["orgs"] = json!(self.orgs_on(&r.name));
141        v["version"] = version_view(&r.name, &h.heartbeat, r.vm.is_some());
142        v["health"] = serde_json::to_value(h).unwrap_or_default();
143        v
144    }
145
146    /// Refuse to forward to `server` when its agent speaks a protocol this
147    /// control plane cannot (`need`: the oldest that will do).
148    pub fn check_protocol(&self, server: &str, need: u64) -> Result<()> {
149        upgrade::compatible(server, &self.health(server).heartbeat, need)
150    }
151
152    /// Place `org` on `server` (the agent is told first, so it accepts
153    /// calls for it), or take it off (`None`).
154    pub fn place(&self, org: &OrgId, server: Option<&str>) -> Result<()> {
155        let _g = self.admin.lock().unwrap();
156        match server {
157            Some(s) => {
158                self.client(s)?.internal(
159                    "POST",
160                    "/internal/v1/orgs",
161                    Some(&json!({"org": org, "placed": true})),
162                    health::TIMEOUT,
163                )?;
164                let mut st = self.store.lock().unwrap();
165                st.placement.insert(org.clone(), s.to_string());
166                st.save()
167            }
168            None => {
169                let prev = self.store.lock().unwrap().placement.get(org).cloned();
170                if let Some(s) = prev {
171                    // Best effort: the agent may be gone for good.
172                    if let Ok(c) = self.client(&s) {
173                        if let Err(e) = c.internal(
174                            "POST",
175                            "/internal/v1/orgs",
176                            Some(&json!({"org": org, "placed": false})),
177                            health::TIMEOUT,
178                        ) {
179                            eprintln!("isb serve: server {s}: unplacing {org}: {e}");
180                        }
181                    }
182                }
183                let mut st = self.store.lock().unwrap();
184                st.placement.remove(org);
185                st.save()
186            }
187        }
188    }
189
190    /// Forward a tool call for `org` (or an unscoped cross-org read) to
191    /// `server` as `caller`.
192    pub fn call(
193        &self,
194        server: &str,
195        tool: &str,
196        args: &Value,
197        caller: &crate::server::Caller,
198        org: Option<&OrgId>,
199        request_id: Option<&str>,
200    ) -> Result<Value> {
201        let who = Assertion::for_caller(caller)
202            .ok_or_else(|| Error::Forbidden(format!("{caller} cannot act on server {server}")))?;
203        self.check_protocol(server, upgrade::MIN_PROTOCOL)?;
204        self.client(server)?
205            .call(tool, args, &who, org, request_id, CALL_TIMEOUT)
206    }
207
208    /// What `add` would refuse before it starts: a bad name, target or
209    /// address, or a server of that name already.
210    pub fn check_add(&self, o: &bootstrap::AddOptions) -> Result<String> {
211        bootstrap::validate_name(&o.name)?;
212        let (_, host) = bootstrap::ssh_host(&o.ssh)?;
213        for c in &o.allow_from {
214            bootstrap::check_cidr(c)?;
215        }
216        let address = o.address.clone().unwrap_or_else(|| host.to_string());
217        bootstrap::check_address(&address)?;
218        if self.store.lock().unwrap().servers.contains_key(&o.name) {
219            return Err(Error::AlreadyExists(format!("server {}", o.name)));
220        }
221        Ok(address)
222    }
223
224    /// Bootstrap a new server over SSH and record it, reporting to `p`.
225    pub fn add(
226        self: &Arc<Self>,
227        o: &bootstrap::AddOptions,
228        ctl: Option<&Controller>,
229        p: &provision::Provision,
230    ) -> Result<ServerRecord> {
231        let address = self.check_add(o)?;
232        let leaf = self.ca.issue_server(&o.name, &address)?;
233        let known = self.dir.join("known_hosts");
234        let ssh = bootstrap::Ssh::new(&o.ssh, o.ssh_port, &o.key, &known);
235        let scratch = self.dir.join(format!("tmp-{}", o.name));
236        std::fs::create_dir_all(&scratch)?;
237        let r = bootstrap::run(o, &ssh, &self.ca.cert_pem, &leaf, &scratch, p);
238        let _ = std::fs::remove_dir_all(&scratch);
239        let incus = r?;
240        self.register(
241            ServerRecord {
242                name: o.name.clone(),
243                address,
244                port: o.agent_port,
245                ssh: o.ssh.clone(),
246                ssh_port: o.ssh_port,
247                added_at: now_secs(),
248                fingerprint: leaf.fingerprint()?,
249                cert_not_after: Some(now_secs() + pki::LEAF_DAYS as u64 * 86400),
250                isb_version: String::new(),
251                allow_from: o.allow_from.clone(),
252                vm: None,
253            },
254            &incus,
255            ctl,
256            p,
257        )
258    }
259
260    /// Make a dedicated VM for `org` on this host and record it as server
261    /// `vm-<org>` (docs/guides/servers.md#dedicated-vms). Idempotent: a VM or a
262    /// record left by an earlier attempt is reused.
263    pub fn add_vm(
264        self: &Arc<Self>,
265        client: &crate::Client,
266        org: &OrgId,
267        size: &vm::VmSize,
268        ctl: Option<&Controller>,
269        p: &provision::Provision,
270    ) -> Result<ServerRecord> {
271        let name = vm::server_name(org);
272        if let Ok(r) = self.record(&name) {
273            match &r.vm {
274                Some(v) if &v.org == org => {
275                    let c = self.client(&name)?;
276                    if c.internal("GET", "/internal/v1/heartbeat", None, health::TIMEOUT)
277                        .is_ok()
278                    {
279                        p.log(&format!("server {name} is up already"));
280                        return Ok(r);
281                    }
282                    p.log(&format!(
283                        "server {name} is recorded but does not answer: bootstrapping it again"
284                    ));
285                }
286                _ => {
287                    return Err(Error::AlreadyExists(format!(
288                        "server {name} (not org {org}'s dedicated VM)"
289                    )));
290                }
291            }
292        }
293        let booted = vm::boot(client, org, size, p)?;
294        let leaf = self.ca.issue_server(&name, &booted.address)?;
295        let binary = bootstrap::own_binary(std::env::consts::ARCH)?;
296        let sha = bootstrap::sha256_hex(&binary);
297        let upload = "/root/isb-agent.upload";
298        let allow = vec![booted.host_address.clone()];
299        let script = bootstrap::render_script(
300            upload,
301            &sha,
302            &self.ca.cert_pem,
303            &leaf,
304            bootstrap::DEFAULT_AGENT_PORT,
305            &allow,
306            None,
307            false,
308        );
309        let incus = vm::install(client, org, &binary, &script, upload, p)?;
310        self.register(
311            ServerRecord {
312                name: name.clone(),
313                address: booted.address,
314                port: bootstrap::DEFAULT_AGENT_PORT,
315                ssh: String::new(),
316                ssh_port: 0,
317                added_at: now_secs(),
318                fingerprint: leaf.fingerprint()?,
319                cert_not_after: Some(now_secs() + pki::LEAF_DAYS as u64 * 86400),
320                isb_version: String::new(),
321                allow_from: allow,
322                vm: Some(store::VmRecord {
323                    org: org.clone(),
324                    project: vm::PROJECT.to_string(),
325                    instance: name,
326                    cpus: size.cpus,
327                    memory: size.memory.clone(),
328                    disk: size.disk.clone(),
329                }),
330            },
331            &incus,
332            ctl,
333            p,
334        )
335    }
336
337    /// Wait for a freshly bootstrapped agent, check it presents the
338    /// certificate just issued, and record it.
339    fn register(
340        self: &Arc<Self>,
341        mut rec: ServerRecord,
342        incus: &str,
343        ctl: Option<&Controller>,
344        p: &provision::Provision,
345    ) -> Result<ServerRecord> {
346        p.step("agent");
347        p.log(&format!(
348            "waiting for the agent on {}:{} ({incus})",
349            rec.address, rec.port
350        ));
351        let c = AgentClient::new(&rec.name, &rec.address, rec.port, self.tls.clone());
352        let hb = wait_heartbeat(&c, Duration::from_secs(120))?;
353        if c.peer_fingerprint()? != rec.fingerprint {
354            return Err(Error::invalid(format!(
355                "server {}: the agent answered with a certificate other than the one just issued",
356                rec.name
357            )));
358        }
359        rec.isb_version = hb["isb"].as_str().unwrap_or("").to_string();
360        {
361            let _g = self.admin.lock().unwrap();
362            let mut st = self.store.lock().unwrap();
363            if let Some(old) = st.servers.get(&rec.name) {
364                // Bootstrapped again (a dedicated VM that stopped answering).
365                rec.added_at = old.added_at;
366            }
367            st.servers.insert(rec.name.clone(), rec.clone());
368            st.save()?;
369        }
370        self.health
371            .lock()
372            .unwrap()
373            .entry(rec.name.clone())
374            .or_default()
375            .observe(Ok(hb), now_secs());
376        if let Some(ctl) = ctl {
377            self.mirror(&rec.name, ctl.clone());
378        }
379        p.log(&format!("server {} is up", rec.name));
380        Ok(rec)
381    }
382
383    /// Forget a server. Refused while orgs are placed on it; the agent
384    /// itself keeps running until it is stopped on the box.
385    pub fn remove(&self, name: &str) -> Result<ServerRecord> {
386        let _g = self.admin.lock().unwrap();
387        let mut st = self.store.lock().unwrap();
388        let orgs = st.orgs_on(name);
389        if !orgs.is_empty() {
390            return Err(Error::invalid(format!(
391                "server {name} holds orgs ({}); delete them first",
392                orgs.iter()
393                    .map(|o| o.as_str())
394                    .collect::<Vec<_>>()
395                    .join(", ")
396            )));
397        }
398        let rec = st
399            .servers
400            .remove(name)
401            .ok_or_else(|| Error::NotFound(format!("server {name}")))?;
402        st.save()?;
403        self.health.lock().unwrap().remove(name);
404        Ok(rec)
405    }
406
407    /// Issue the agent a new certificate over the current mTLS connection
408    /// and check it answers with it.
409    pub fn rotate_cert(&self, name: &str) -> Result<ServerRecord> {
410        let _g = self.admin.lock().unwrap();
411        let rec = self.record(name)?;
412        let leaf = self.ca.issue_server(name, &rec.address)?;
413        let c = self.client(name)?;
414        c.internal(
415            "POST",
416            "/internal/v1/cert",
417            Some(&json!({"cert": leaf.cert, "key": leaf.key})),
418            health::TIMEOUT,
419        )?;
420        let fp = leaf.fingerprint()?;
421        let got = c.peer_fingerprint()?;
422        if got != fp {
423            return Err(Error::invalid(format!(
424                "server {name}: still presents {got} after rotation"
425            )));
426        }
427        let mut st = self.store.lock().unwrap();
428        let r = st
429            .servers
430            .get_mut(name)
431            .ok_or_else(|| Error::NotFound(format!("server {name}")))?;
432        r.fingerprint = fp;
433        r.cert_not_after = Some(now_secs() + pki::LEAF_DAYS as u64 * 86400);
434        let r = r.clone();
435        st.save()?;
436        Ok(r)
437    }
438
439    /// Heartbeats for every server, and a mirror of each one's events into
440    /// `ctl`'s feed.
441    pub fn start(self: &Arc<Self>, ctl: Controller) {
442        for r in self.records() {
443            self.mirror(&r.name, ctl.clone());
444        }
445        let me = self.clone();
446        let _ = std::thread::Builder::new()
447            .name("isb-servers".into())
448            .spawn(move || {
449                while !me.stop.load(Ordering::SeqCst) {
450                    for r in me.records() {
451                        me.beat(&r.name, &ctl);
452                    }
453                    let mut slept = Duration::ZERO;
454                    while slept < health::INTERVAL && !me.stop.load(Ordering::SeqCst) {
455                        std::thread::sleep(Duration::from_millis(250));
456                        slept += Duration::from_millis(250);
457                    }
458                }
459            });
460    }
461
462    pub fn shutdown(&self) {
463        self.stop.store(true, Ordering::SeqCst);
464    }
465
466    fn beat(&self, name: &str, ctl: &Controller) {
467        let r = self
468            .client(name)
469            .and_then(|c| c.internal("GET", "/internal/v1/heartbeat", None, health::TIMEOUT))
470            .map_err(|e| e.to_string());
471        let t = self
472            .health
473            .lock()
474            .unwrap()
475            .entry(name.to_string())
476            .or_default()
477            .observe(r, now_secs());
478        let Some(t) = t else { return };
479        let h = self.health(name);
480        let (kind, level, msg) = match t {
481            health::Transition::Unreachable => (
482                "server.unreachable",
483                "error",
484                format!(
485                    "server {name} is unreachable: {}",
486                    h.last_error.as_deref().unwrap_or("no answer")
487                ),
488            ),
489            health::Transition::Recovered => (
490                "server.recovered",
491                "info",
492                format!("server {name} answers again"),
493            ),
494        };
495        let mut stacks: Vec<String> = self
496            .orgs_on(name)
497            .iter()
498            .map(|o| format!("{o}/@servers"))
499            .collect();
500        stacks.push("system/@servers".into());
501        for s in stacks {
502            ctl.relay(Some(kind), level, &s, name, None, msg.clone());
503        }
504        eprintln!("isb serve: {msg}");
505    }
506
507    /// Follow `name`'s event feed into `ctl`, from where it is now. Only
508    /// events of orgs placed on it are taken, so a server cannot speak for
509    /// another's orgs. Ends when the server is removed.
510    fn mirror(self: &Arc<Self>, name: &str, ctl: Controller) {
511        if !self.mirrored.lock().unwrap().insert(name.to_string()) {
512            return;
513        }
514        let (me, name) = (self.clone(), name.to_string());
515        let _ = std::thread::Builder::new()
516            .name(format!("isb-mirror-{name}"))
517            .spawn(move || {
518                let who = Assertion::control_plane();
519                let mut cursor: Option<u64> = None;
520                while !me.stop.load(Ordering::SeqCst) {
521                    let Ok(c) = me.client(&name) else { break };
522                    let since = cursor.unwrap_or(u64::MAX);
523                    let args = match cursor {
524                        Some(s) => json!({"since": s, "limit": 500, "wait": 25}),
525                        // First contact: just learn where the feed is.
526                        None => json!({"since": since, "limit": 1, "wait": 0}),
527                    };
528                    match c.call("events", &args, &who, None, None, Duration::from_secs(40)) {
529                        Ok(v) => {
530                            let seq = v["seq"].as_u64().unwrap_or(0);
531                            let Some(from) = cursor else {
532                                cursor = Some(seq);
533                                continue;
534                            };
535                            if seq < from {
536                                // The agent restarted: its feed starts over.
537                                cursor = Some(0);
538                                continue;
539                            }
540                            let placed = me.orgs_on(&name);
541                            for e in v["events"].as_array().into_iter().flatten() {
542                                relay(&ctl, e, &placed);
543                            }
544                            cursor = Some(seq);
545                        }
546                        Err(_) => std::thread::sleep(Duration::from_secs(5)),
547                    }
548                }
549                me.mirrored.lock().unwrap().remove(&name);
550            });
551    }
552}
553
554/// What a server runs next to what this control plane runs: the version,
555/// the build (a hash of the binary, so two builds of one version differ),
556/// and whether they speak the same protocol.
557fn version_view(name: &str, hb: &Value, vm: bool) -> Value {
558    let known = !hb.is_null();
559    let build = hb["build"].as_str().unwrap_or("");
560    json!({
561        "isb": hb["isb"],
562        "build": hb["build"],
563        "protocol": known.then(|| upgrade::protocol_of(hb)),
564        "control_plane": {
565            "isb": env!("CARGO_PKG_VERSION"),
566            "build": upgrade::build_id(),
567            "protocol": upgrade::PROTOCOL,
568        },
569        // Unknown until it answers; a different build of the same version
570        // is skew too.
571        "skew": known && (hb["isb"].as_str() != Some(env!("CARGO_PKG_VERSION")) || build != upgrade::build_id()),
572        "compatible": upgrade::compatible(name, hb, upgrade::MIN_PROTOCOL).is_ok(),
573        "ssh": upgrade::compatible(name, hb, upgrade::SSH_PROTOCOL).is_ok(),
574        // A dedicated VM is upgraded through incus, helper or not.
575        "upgradable": vm || hb["upgrade"]["helper"] == true,
576        "last_upgrade": hb["upgrade"]["last"],
577    })
578}
579
580/// Re-emit one of a server's events if it belongs to an org placed there.
581fn relay(ctl: &Controller, e: &Value, placed: &[OrgId]) {
582    let stack = e["stack"].as_str().unwrap_or("");
583    let org = stack.split_once('/').map(|(o, _)| o).unwrap_or("");
584    if !placed.iter().any(|o| o.as_str() == org) {
585        return;
586    }
587    ctl.relay(
588        e["kind"].as_str(),
589        e["level"].as_str().unwrap_or("info"),
590        stack,
591        e["service"].as_str().unwrap_or(""),
592        e["instance"].as_str(),
593        e["message"].as_str().unwrap_or("").to_string(),
594    );
595}
596
597fn wait_heartbeat(c: &AgentClient, timeout: Duration) -> Result<Value> {
598    let started = std::time::Instant::now();
599    loop {
600        match c.internal("GET", "/internal/v1/heartbeat", None, health::TIMEOUT) {
601            Ok(v) => return Ok(v),
602            Err(e) if started.elapsed() >= timeout => {
603                return Err(Error::OperationFailed {
604                    step: format!("wait for the agent on server {}", c.name),
605                    message: e.to_string(),
606                });
607            }
608            Err(_) => std::thread::sleep(Duration::from_secs(2)),
609        }
610    }
611}