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