Skip to main content

isb_server/servers/
provision.rs

1//! Progress of a server being added: over SSH (`server_add`) or as a
2//! dedicated VM this control plane makes itself (`org_create` with
3//! `placement: {vm: ...}`). Each run is a fixed list of steps, the lines it
4//! logged and how it ended, kept in memory so the web UI and the CLI can
5//! follow one that runs in the background (`server_provision_get`, and
6//! `provisions` in `server_list`). Nothing secret goes in: the SSH key never
7//! reaches a log line.
8
9use std::collections::BTreeMap;
10use std::sync::{Arc, Mutex};
11
12use serde::Serialize;
13use serde_json::Value;
14
15use crate::stack::now_secs;
16
17/// Lines of log kept per run.
18const MAX_LOG: usize = 400;
19/// How long a finished run stays listed: failed ones longer, so they can be
20/// read and retried.
21const KEEP_DONE: u64 = 10 * 60;
22const KEEP_FAILED: u64 = 60 * 60;
23
24#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
25#[serde(rename_all = "snake_case")]
26pub enum StepState {
27    Pending,
28    Running,
29    Done,
30    Failed,
31}
32
33#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
34#[serde(rename_all = "snake_case")]
35pub enum RunState {
36    Running,
37    Done,
38    Failed,
39}
40
41#[derive(Debug, Clone, Serialize)]
42pub struct Step {
43    pub id: &'static str,
44    pub title: &'static str,
45    pub state: StepState,
46    pub started_at: Option<u64>,
47    pub finished_at: Option<u64>,
48}
49
50/// How a server is being made.
51#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
52#[serde(rename_all = "snake_case")]
53pub enum Kind {
54    /// `server_add`: a box reached over SSH.
55    Ssh,
56    /// A dedicated VM on this control plane's host, for one org.
57    Vm,
58}
59
60/// The steps of each kind, in order.
61pub fn steps(kind: Kind) -> &'static [(&'static str, &'static str)] {
62    match kind {
63        Kind::Ssh => &[
64            ("check", "Check the box over SSH"),
65            ("binary", "Get the isb binary"),
66            ("upload", "Upload isb"),
67            ("install", "Install incus, the agent and its unit"),
68            ("agent", "Wait for the agent's heartbeat"),
69        ],
70        Kind::Vm => &[
71            ("support", "Check this host can run VMs"),
72            ("vm", "Create and start the VM"),
73            ("boot", "Wait for the VM to boot"),
74            ("upload", "Copy isb into the VM"),
75            ("install", "Install incus, the agent and its firewall"),
76            ("agent", "Wait for the agent's heartbeat"),
77            ("org", "Create the org on it"),
78        ],
79    }
80}
81
82/// One run, as the tools answer it.
83#[derive(Debug, Clone, Serialize)]
84pub struct View {
85    /// The server's name (`vm-<org>` for a dedicated VM).
86    pub name: String,
87    pub kind: Kind,
88    /// The org a dedicated VM is for.
89    #[serde(skip_serializing_if = "Option::is_none")]
90    pub org: Option<String>,
91    pub state: RunState,
92    pub started_at: u64,
93    pub finished_at: Option<u64>,
94    pub steps: Vec<Step>,
95    pub log: Vec<String>,
96    /// Lines dropped off the top of `log` (the first kept line's number).
97    pub log_start: usize,
98    pub error: Option<String>,
99    /// What was asked, without secrets: enough for "retry" to ask again.
100    pub request: Value,
101    /// What the run produced (the server, or the org it created).
102    #[serde(skip_serializing_if = "Option::is_none")]
103    pub result: Option<Value>,
104}
105
106/// A handle on one run: cheap to clone, shared with the thread doing it.
107#[derive(Debug, Clone)]
108pub struct Provision(Arc<Mutex<View>>);
109
110impl Provision {
111    pub fn new(name: &str, kind: Kind, org: Option<&str>, request: Value) -> Provision {
112        Provision(Arc::new(Mutex::new(View {
113            name: name.to_string(),
114            kind,
115            org: org.map(str::to_string),
116            state: RunState::Running,
117            started_at: now_secs(),
118            finished_at: None,
119            steps: steps(kind)
120                .iter()
121                .map(|(id, title)| Step {
122                    id,
123                    title,
124                    state: StepState::Pending,
125                    started_at: None,
126                    finished_at: None,
127                })
128                .collect(),
129            log: Vec::new(),
130            log_start: 0,
131            error: None,
132            request,
133            result: None,
134        })))
135    }
136
137    pub fn view(&self) -> View {
138        self.0.lock().unwrap().clone()
139    }
140
141    /// Start step `id`: the one running before it is done.
142    pub fn step(&self, id: &str) {
143        let now = now_secs();
144        let mut v = self.0.lock().unwrap();
145        for s in v.steps.iter_mut() {
146            if s.state == StepState::Running {
147                s.state = StepState::Done;
148                s.finished_at = Some(now);
149            }
150        }
151        if let Some(s) = v.steps.iter_mut().find(|s| s.id == id) {
152            s.state = StepState::Running;
153            s.started_at = Some(now);
154        }
155    }
156
157    /// A line of progress (also the daemon's log).
158    pub fn log(&self, line: &str) {
159        let mut v = self.0.lock().unwrap();
160        eprintln!("isb serve: server {}: {line}", v.name);
161        v.log.push(line.to_string());
162        let over = v.log.len().saturating_sub(MAX_LOG);
163        if over > 0 {
164            v.log.drain(..over);
165            v.log_start += over;
166        }
167    }
168
169    /// The run ended: every step done, or the running one failed.
170    pub fn finish(&self, r: std::result::Result<Value, String>) {
171        let now = now_secs();
172        let mut v = self.0.lock().unwrap();
173        v.finished_at = Some(now);
174        match r {
175            Ok(result) => {
176                for s in v.steps.iter_mut() {
177                    if s.state != StepState::Done {
178                        s.state = StepState::Done;
179                        s.started_at.get_or_insert(now);
180                        s.finished_at = Some(now);
181                    }
182                }
183                v.state = RunState::Done;
184                v.result = Some(result);
185            }
186            Err(e) => {
187                for s in v.steps.iter_mut() {
188                    if s.state == StepState::Running {
189                        s.state = StepState::Failed;
190                        s.finished_at = Some(now);
191                    }
192                }
193                v.state = RunState::Failed;
194                v.error = Some(e);
195            }
196        }
197    }
198
199    pub fn running(&self) -> bool {
200        self.0.lock().unwrap().state == RunState::Running
201    }
202}
203
204/// Every run by server name: one at a time per name.
205#[derive(Debug, Default)]
206pub struct Runs(Mutex<BTreeMap<String, Provision>>);
207
208impl Runs {
209    /// Start a run for `name`, refused while one is running.
210    pub fn begin(&self, p: Provision) -> crate::error::Result<Provision> {
211        let mut m = self.0.lock().unwrap();
212        let name = p.view().name;
213        if m.get(&name).is_some_and(Provision::running) {
214            return Err(crate::error::Error::AlreadyExists(format!(
215                "server {name} is being added already (server_provision_get shows how far it got)"
216            )));
217        }
218        m.insert(name, p.clone());
219        Ok(p)
220    }
221
222    pub fn get(&self, name: &str) -> Option<View> {
223        self.0.lock().unwrap().get(name).map(Provision::view)
224    }
225
226    /// Runs going on, and recent finished ones (failed ones for an hour).
227    pub fn list(&self) -> Vec<View> {
228        let now = now_secs();
229        let mut m = self.0.lock().unwrap();
230        m.retain(|_, p| {
231            let v = p.view();
232            match (v.state, v.finished_at) {
233                (RunState::Running, _) | (_, None) => true,
234                (RunState::Done, Some(t)) => now.saturating_sub(t) < KEEP_DONE,
235                (RunState::Failed, Some(t)) => now.saturating_sub(t) < KEEP_FAILED,
236            }
237        });
238        m.values().map(Provision::view).collect()
239    }
240}
241
242#[cfg(test)]
243mod tests {
244    use super::*;
245    use serde_json::json;
246
247    #[test]
248    fn steps_advance_and_a_failure_marks_the_running_one() {
249        let p = Provision::new("vm-acme", Kind::Vm, Some("acme"), json!({"org": "acme"}));
250        p.step("support");
251        p.step("vm");
252        let v = p.view();
253        assert_eq!(v.steps[0].state, StepState::Done);
254        assert_eq!(v.steps[1].state, StepState::Running);
255        assert_eq!(v.steps[2].state, StepState::Pending);
256        p.finish(Err("no kvm".into()));
257        let v = p.view();
258        assert_eq!(v.state, RunState::Failed);
259        assert_eq!(v.steps[1].state, StepState::Failed);
260        assert_eq!(v.steps[2].state, StepState::Pending);
261        assert_eq!(v.error.as_deref(), Some("no kvm"));
262
263        let ok = Provision::new("box", Kind::Ssh, None, json!({}));
264        ok.step("check");
265        ok.finish(Ok(json!({"name": "box"})));
266        assert!(ok.view().steps.iter().all(|s| s.state == StepState::Done));
267    }
268
269    #[test]
270    fn one_run_per_name_and_the_log_is_bounded() {
271        let runs = Runs::default();
272        let a = runs
273            .begin(Provision::new("box", Kind::Ssh, None, json!({})))
274            .unwrap();
275        assert!(
276            runs.begin(Provision::new("box", Kind::Ssh, None, json!({})))
277                .is_err()
278        );
279        a.finish(Err("x".into()));
280        runs.begin(Provision::new("box", Kind::Ssh, None, json!({})))
281            .unwrap();
282        for i in 0..(MAX_LOG + 10) {
283            a.log(&format!("line {i}"));
284        }
285        assert_eq!(a.view().log.len(), MAX_LOG);
286        assert_eq!(a.view().log_start, 10);
287        assert_eq!(a.view().log[0], "line 10");
288        assert_eq!(runs.list().len(), 1);
289    }
290}