1pub 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
33pub const CALL_TIMEOUT: Duration = Duration::from_secs(20 * 60);
35
36pub 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 pub runs: provision::Runs,
47 admin: Mutex<()>,
50 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 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 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 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 pub fn check_protocol(&self, server: &str, need: u64) -> Result<()> {
149 upgrade::compatible(server, &self.health(server).heartbeat, need)
150 }
151
152 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 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 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 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 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 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 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 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 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 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 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 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 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 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
554fn 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 "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 "upgradable": vm || hb["upgrade"]["helper"] == true,
576 "last_upgrade": hb["upgrade"]["last"],
577 })
578}
579
580fn 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}