1use std::collections::BTreeMap;
10use std::sync::{Arc, Mutex};
11
12use serde::Serialize;
13use serde_json::Value;
14
15use crate::stack::now_secs;
16
17const MAX_LOG: usize = 400;
19const 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#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
52#[serde(rename_all = "snake_case")]
53pub enum Kind {
54 Ssh,
56 Vm,
58}
59
60pub 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#[derive(Debug, Clone, Serialize)]
84pub struct View {
85 pub name: String,
87 pub kind: Kind,
88 #[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 pub log_start: usize,
98 pub error: Option<String>,
99 pub request: Value,
101 #[serde(skip_serializing_if = "Option::is_none")]
103 pub result: Option<Value>,
104}
105
106#[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 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 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 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#[derive(Debug, Default)]
206pub struct Runs(Mutex<BTreeMap<String, Provision>>);
207
208impl Runs {
209 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 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}