1use std::collections::BTreeMap;
16use std::path::{Path, PathBuf};
17use std::sync::{Arc, Mutex};
18use std::time::{Duration, Instant};
19
20use serde::{Deserialize, Serialize};
21use serde_json::Value;
22
23use super::Controller;
24use super::now_secs;
25use crate::error::{Error, Result};
26use crate::org::OrgId;
27use crate::spec::DomainSpec;
28
29pub const KEEP_DEPLOYMENTS: usize = 30;
31
32const EVENTS_KEPT: usize = 1000;
34
35#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
36#[serde(rename_all = "lowercase")]
37pub enum Status {
38 Queued,
39 Deploying,
40 Done,
41 Failed,
42 Superseded,
44}
45
46impl Status {
47 pub fn finished(self) -> bool {
48 !matches!(self, Status::Queued | Status::Deploying)
49 }
50}
51
52#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
54pub struct LogLine {
55 pub at: u64,
57 pub level: String,
58 #[serde(default, skip_serializing_if = "String::is_empty")]
59 pub service: String,
60 pub message: String,
61}
62
63#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
65pub struct Deployment {
66 pub id: u64,
67 pub stack: String,
68 pub trigger: String,
70 pub action: String,
73 pub actor: String,
75 pub status: Status,
76 #[serde(default, skip_serializing_if = "Option::is_none")]
78 pub rollback_of: Option<u64>,
79 #[serde(default)]
81 pub services: Vec<String>,
82 #[serde(default, skip_serializing_if = "Vec::is_empty")]
85 pub reused_secrets: Vec<String>,
86 #[serde(default, skip_serializing_if = "Option::is_none")]
87 pub error: Option<String>,
88 pub created_at: u64,
90 #[serde(default, skip_serializing_if = "Option::is_none")]
91 pub started_at: Option<u64>,
92 #[serde(default, skip_serializing_if = "Option::is_none")]
93 pub finished_at: Option<u64>,
94 #[serde(default)]
96 pub source: String,
97 #[serde(default)]
99 pub base_dir: PathBuf,
100 #[serde(default)]
102 pub env: String,
103 #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
105 pub domains: BTreeMap<String, Vec<DomainSpec>>,
106 #[serde(default)]
108 pub events: Vec<LogLine>,
109 #[serde(default)]
111 pub events_seq: u64,
112}
113
114impl Deployment {
115 pub fn summary(&self) -> Value {
118 let mut v = serde_json::to_value(self).unwrap_or_default();
119 if let Some(o) = v.as_object_mut() {
120 for k in [
121 "source",
122 "base_dir",
123 "env",
124 "domains",
125 "events",
126 "events_seq",
127 ] {
128 o.remove(k);
129 }
130 }
131 v
132 }
133
134 pub fn log(&self) -> String {
136 let mut out = String::new();
137 for l in &self.events {
138 let secs = (l.at / 1000) % 86_400;
139 out.push_str(&format!(
140 "{:02}:{:02}:{:02} {}{}{}\n",
141 secs / 3600,
142 (secs / 60) % 60,
143 secs % 60,
144 if l.level == "info" || l.level == "log" {
145 String::new()
146 } else {
147 format!("[{}] ", l.level)
148 },
149 if l.service.is_empty() {
150 String::new()
151 } else {
152 format!("{}: ", l.service)
153 },
154 l.message
155 ));
156 }
157 out
158 }
159}
160
161#[derive(Clone)]
163pub struct StackMeta {
164 dir: PathBuf,
165 edit: Arc<Mutex<()>>,
167}
168
169fn write_private(path: &Path, bytes: &[u8]) -> Result<()> {
170 use std::io::Write;
171 use std::os::unix::fs::OpenOptionsExt;
172 if let Some(p) = path.parent() {
173 std::fs::create_dir_all(p)?;
174 }
175 let tmp = path.with_extension("tmp");
176 let mut f = std::fs::OpenOptions::new()
177 .write(true)
178 .create(true)
179 .truncate(true)
180 .mode(0o600)
181 .open(&tmp)?;
182 f.write_all(bytes)?;
183 f.sync_all()?;
184 std::fs::rename(&tmp, path)?;
185 Ok(())
186}
187
188fn read_opt(path: &Path) -> Result<Option<Vec<u8>>> {
189 match std::fs::read(path) {
190 Ok(b) => Ok(Some(b)),
191 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
192 Err(e) => Err(e.into()),
193 }
194}
195
196impl StackMeta {
197 pub fn new(state: &Path) -> StackMeta {
199 StackMeta {
200 dir: state.to_path_buf(),
201 edit: Arc::new(Mutex::new(())),
202 }
203 }
204
205 fn stack_dir(&self, org: &OrgId, name: &str) -> PathBuf {
206 let base = if org.is_default() {
207 self.dir.clone()
208 } else {
209 org.dir(&self.dir)
210 };
211 base.join("stack-meta").join(name)
212 }
213
214 fn deployments_dir(&self, org: &OrgId, name: &str) -> PathBuf {
215 self.stack_dir(org, name).join("deployments")
216 }
217
218 pub fn exists(&self, org: &OrgId, name: &str) -> bool {
220 self.stack_dir(org, name).is_dir()
221 }
222
223 pub fn remove(&self, org: &OrgId, name: &str) -> Result<()> {
225 match std::fs::remove_dir_all(self.stack_dir(org, name)) {
226 Err(e) if e.kind() != std::io::ErrorKind::NotFound => Err(e.into()),
227 _ => Ok(()),
228 }
229 }
230
231 pub fn env(&self, org: &OrgId, name: &str) -> Result<String> {
235 let b = read_opt(&self.stack_dir(org, name).join("env"))?;
236 Ok(b.map(|b| String::from_utf8_lossy(&b).into_owned())
237 .unwrap_or_default())
238 }
239
240 pub fn set_env(&self, org: &OrgId, name: &str, text: &str) -> Result<()> {
241 write_private(&self.stack_dir(org, name).join("env"), text.as_bytes())
242 }
243
244 pub fn domains(&self, org: &OrgId, name: &str) -> Result<BTreeMap<String, Vec<DomainSpec>>> {
246 match read_opt(&self.stack_dir(org, name).join("domains.json"))? {
247 Some(b) => Ok(serde_json::from_slice(&b)?),
248 None => Ok(BTreeMap::new()),
249 }
250 }
251
252 pub fn set_domains(
253 &self,
254 org: &OrgId,
255 name: &str,
256 domains: &BTreeMap<String, Vec<DomainSpec>>,
257 ) -> Result<()> {
258 let kept: BTreeMap<&String, &Vec<DomainSpec>> =
259 domains.iter().filter(|(_, d)| !d.is_empty()).collect();
260 write_private(
261 &self.stack_dir(org, name).join("domains.json"),
262 serde_json::to_string_pretty(&kept)?.as_bytes(),
263 )
264 }
265
266 pub fn deployments(&self, org: &OrgId, name: &str) -> Result<Vec<Deployment>> {
270 let mut out = Vec::new();
271 let Ok(rd) = std::fs::read_dir(self.deployments_dir(org, name)) else {
272 return Ok(out);
273 };
274 for e in rd.flatten() {
275 let p = e.path();
276 if p.extension().is_some_and(|x| x == "json") {
277 if let Ok(d) = serde_json::from_slice::<Deployment>(&std::fs::read(&p)?) {
278 out.push(d);
279 }
280 }
281 }
282 out.sort_by_key(|a| std::cmp::Reverse(a.id));
283 Ok(out)
284 }
285
286 pub fn deployment(&self, org: &OrgId, name: &str, id: u64) -> Result<Deployment> {
287 let p = self.deployments_dir(org, name).join(format!("{id}.json"));
288 match read_opt(&p)? {
289 Some(b) => Ok(serde_json::from_slice(&b)?),
290 None => Err(Error::NotFound(format!(
291 "deployment {id} of stack {name} (the last {KEEP_DEPLOYMENTS} are kept)"
292 ))),
293 }
294 }
295
296 fn save(&self, org: &OrgId, d: &Deployment) -> Result<()> {
297 write_private(
298 &self
299 .deployments_dir(org, &d.stack)
300 .join(format!("{}.json", d.id)),
301 serde_json::to_string_pretty(d)?.as_bytes(),
302 )
303 }
304
305 pub fn start(&self, org: &OrgId, mut d: Deployment) -> Result<Deployment> {
309 let _g = self.edit.lock().unwrap();
310 let all = self.deployments(org, &d.stack)?;
311 d.id = all.first().map_or(1, |x| x.id + 1);
312 d.status = Status::Deploying;
313 d.created_at = now_secs();
314 d.started_at = Some(d.created_at);
315 for mut old in all.iter().filter(|x| !x.status.finished()).cloned() {
316 old.status = Status::Superseded;
317 old.finished_at = Some(d.created_at);
318 old.error = Some(format!("superseded by deployment {}", d.id));
319 self.save(org, &old)?;
320 }
321 self.save(org, &d)?;
322 for old in all.iter().skip(KEEP_DEPLOYMENTS.saturating_sub(1)) {
323 let _ = std::fs::remove_file(
324 self.deployments_dir(org, &d.stack)
325 .join(format!("{}.json", old.id)),
326 );
327 }
328 Ok(d)
329 }
330
331 fn update(
334 &self,
335 org: &OrgId,
336 name: &str,
337 id: u64,
338 f: impl FnOnce(&mut Deployment),
339 ) -> Result<Deployment> {
340 let _g = self.edit.lock().unwrap();
341 let mut d = self.deployment(org, name, id)?;
342 if !d.status.finished() {
343 f(&mut d);
344 if d.status.finished() && d.finished_at.is_none() {
345 d.finished_at = Some(now_secs());
346 }
347 self.save(org, &d)?;
348 }
349 Ok(d)
350 }
351
352 pub fn fail(&self, org: &OrgId, name: &str, id: u64, error: &str) -> Result<Deployment> {
354 self.update(org, name, id, |d| {
355 d.status = Status::Failed;
356 d.error = Some(error.to_string());
357 })
358 }
359
360 pub fn recover(&self) {
362 let mut roots = vec![self.dir.join("stack-meta")];
363 if let Ok(rd) = std::fs::read_dir(self.dir.join("orgs")) {
364 roots.extend(rd.flatten().map(|e| e.path().join("stack-meta")));
365 }
366 for root in roots {
367 let Ok(stacks) = std::fs::read_dir(&root) else {
368 continue;
369 };
370 for s in stacks.flatten() {
371 let Ok(rd) = std::fs::read_dir(s.path().join("deployments")) else {
372 continue;
373 };
374 for e in rd.flatten() {
375 let p = e.path();
376 let Ok(mut d) = std::fs::read(&p)
377 .map_err(Error::from)
378 .and_then(|b| Ok(serde_json::from_slice::<Deployment>(&b)?))
379 else {
380 continue;
381 };
382 if d.status.finished() {
383 continue;
384 }
385 d.status = Status::Failed;
386 d.finished_at = Some(now_secs());
387 d.error = Some("isb serve stopped before the rollout settled".into());
388 if let Ok(t) = serde_json::to_string_pretty(&d) {
389 let _ = write_private(&p, t.as_bytes());
390 }
391 }
392 }
393 }
394 }
395
396 pub fn watch(&self, ctl: Controller, org: OrgId, name: String, id: u64, timeout: Duration) {
399 let meta = self.clone();
400 std::thread::spawn(move || meta.follow(&ctl, &org, &name, id, timeout));
401 }
402
403 fn follow(&self, ctl: &Controller, org: &OrgId, name: &str, id: u64, timeout: Duration) {
404 let q = super::qualified(org, name);
405 let started = Instant::now();
406 loop {
407 let Ok(d) = self.deployment(org, name, id) else {
408 return;
409 };
410 if d.status.finished() {
411 return;
412 }
413 let (_, events) = ctl.events(d.events_seq, EVENTS_KEPT);
414 let seq = events.last().map_or(d.events_seq, |e| e.seq);
415 let lines: Vec<LogLine> = events
416 .into_iter()
417 .filter(|e| e.stack == q)
418 .map(|e| LogLine {
419 at: e.at,
420 level: e.level,
421 service: e.service,
422 message: e.message,
423 })
424 .collect();
425 let outcome = match settled(ctl, &q) {
426 Err(e) => Some((Status::Failed, Some(e.to_string()))),
427 Ok(Some(problems)) if problems.is_empty() => Some((Status::Done, None)),
428 Ok(Some(problems)) => Some((Status::Failed, Some(problems.join("; ")))),
429 Ok(None) if started.elapsed() >= timeout => Some((
430 Status::Failed,
431 Some(format!(
432 "the rollout did not settle in {}s",
433 timeout.as_secs()
434 )),
435 )),
436 Ok(None) => None,
437 };
438 let r = self.update(org, name, id, |d| {
439 d.events.extend(lines);
440 let over = d.events.len().saturating_sub(EVENTS_KEPT);
441 d.events.drain(..over);
442 d.events_seq = seq;
443 if let Some((s, e)) = outcome {
444 d.status = s;
445 d.error = e;
446 }
447 });
448 match r {
449 Ok(d) if !d.status.finished() => {}
450 _ => return,
451 }
452 std::thread::sleep(Duration::from_secs(1));
453 }
454 }
455}
456
457fn settled(ctl: &Controller, q: &str) -> Result<Option<Vec<String>>> {
462 let def = ctl.definition(q)?;
463 let st = ctl.status(q)?;
464 let mut problems = Vec::new();
465 for s in &st.services {
466 let current = def.revision(&s.service).is_ok_and(|r| r == s.rev)
467 && def
468 .service(&s.service)
469 .is_ok_and(|d| d.replicas() == s.replicas);
470 match s.state.as_str() {
471 "converged" if current => {}
472 "paused" | "failing" if current => problems.push(format!(
473 "{}: {}{}",
474 s.service,
475 s.state,
476 s.message
477 .as_deref()
478 .map(|m| format!(" ({m})"))
479 .unwrap_or_default()
480 )),
481 _ => return Ok(None),
482 }
483 }
484 Ok(Some(problems))
485}
486
487#[cfg(test)]
488mod tests {
489 use super::*;
490
491 fn rec(stack: &str) -> Deployment {
492 Deployment {
493 id: 0,
494 stack: stack.into(),
495 trigger: "api".into(),
496 action: "deploy".into(),
497 actor: "u".into(),
498 status: Status::Queued,
499 rollback_of: None,
500 services: vec!["web".into()],
501 reused_secrets: vec![],
502 error: None,
503 created_at: 0,
504 started_at: None,
505 finished_at: None,
506 source: "services: {}\n".into(),
507 base_dir: "/srv".into(),
508 env: "A=1\n".into(),
509 domains: BTreeMap::new(),
510 events: vec![],
511 events_seq: 0,
512 }
513 }
514
515 #[test]
516 fn records_number_supersede_prune_and_recover() {
517 let dir = tempfile::tempdir().unwrap();
518 let m = StackMeta::new(dir.path());
519 let org = OrgId::new("acme").unwrap();
520 let a = m.start(&org, rec("shop")).unwrap();
521 assert_eq!((a.id, a.status), (1, Status::Deploying));
522 let b = m.start(&org, rec("shop")).unwrap();
523 assert_eq!(b.id, 2);
524 let a = m.deployment(&org, "shop", 1).unwrap();
525 assert_eq!(a.status, Status::Superseded);
526 assert_eq!(a.error.as_deref(), Some("superseded by deployment 2"));
527 m.fail(&org, "shop", 1, "late").unwrap();
529 assert_eq!(
530 m.deployment(&org, "shop", 1).unwrap().status,
531 Status::Superseded
532 );
533 for _ in 0..KEEP_DEPLOYMENTS + 3 {
534 m.start(&org, rec("shop")).unwrap();
535 }
536 let all = m.deployments(&org, "shop").unwrap();
537 assert_eq!(all.len(), KEEP_DEPLOYMENTS);
538 assert_eq!(all[0].id, KEEP_DEPLOYMENTS as u64 + 5);
539 assert!(m.deployment(&org, "shop", 1).is_err());
540 let s = all[0].summary();
542 assert!(s.get("source").is_none() && s.get("events").is_none());
543 assert_eq!(s["status"], "deploying");
544 assert!(s.get("reused_secrets").is_none());
546 let mut r = rec("shop");
547 r.reused_secrets = vec!["db_password".into()];
548 let r = m.start(&org, r).unwrap();
549 assert_eq!(
550 r.summary()["reused_secrets"],
551 serde_json::json!(["db_password"])
552 );
553 assert_eq!(
554 m.deployment(&org, "shop", r.id).unwrap().reused_secrets,
555 ["db_password"]
556 );
557 let all = m.deployments(&org, "shop").unwrap();
558 m.recover();
559 let top = m.deployment(&org, "shop", all[0].id).unwrap();
560 assert_eq!(top.status, Status::Failed);
561 assert!(top.finished_at.is_some());
562 m.set_env(&org, "shop", "A=1\n").unwrap();
564 assert_eq!(m.env(&org, "shop").unwrap(), "A=1\n");
565 assert_eq!(m.env(&org, "other").unwrap(), "");
566 let doms = BTreeMap::from([(
567 "web".to_string(),
568 vec![DomainSpec {
569 host: "a.io".into(),
570 ..Default::default()
571 }],
572 )]);
573 m.set_domains(&org, "shop", &doms).unwrap();
574 assert_eq!(m.domains(&org, "shop").unwrap(), doms);
575 assert!(dir.path().join("orgs/acme/stack-meta/shop/env").is_file());
576 m.remove(&org, "shop").unwrap();
577 assert!(!m.exists(&org, "shop"));
578 assert!(m.deployments(&org, "shop").unwrap().is_empty());
579 }
580
581 #[test]
582 fn the_log_reads_as_lines() {
583 let mut d = rec("shop");
584 d.events = vec![
585 LogLine {
586 at: 3_600_000 + 61_000,
587 level: "info".into(),
588 service: String::new(),
589 message: "deployed by u".into(),
590 },
591 LogLine {
592 at: 3_600_000 + 62_000,
593 level: "warn".into(),
594 service: "web".into(),
595 message: "unhealthy".into(),
596 },
597 ];
598 assert_eq!(
599 d.log(),
600 "01:01:01 deployed by u\n01:01:02 [warn] web: unhealthy\n"
601 );
602 }
603}