Skip to main content

isb_core/stack/
deployments.rs

1//! What the daemon keeps beside a compose stack's definition: its
2//! environment (`.env` text that `${VAR}` resolves against), its managed
3//! domains (domain records per service, merged in at deploy), and a record
4//! of each deploy ([`Deployment`]), kept per stack under
5//! `stack-meta/<stack>/` in the stack's org directory.
6//!
7//! A record is written when a deploy is handed to the controller
8//! (`deploying`) and finished by a watcher thread that follows the rollout:
9//! `done` when every service converged, `failed` when one fails or pauses
10//! or the rollout outlasts its timeout, `superseded` when a newer deploy
11//! came first. It keeps the compose source it deployed, the environment
12//! and managed domains of the moment (secret references, never values),
13//! and the stack's events while it ran.
14
15use 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
29/// Deployment records kept per stack, as for an app.
30pub const KEEP_DEPLOYMENTS: usize = 30;
31
32/// Events kept per record.
33const 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    /// A newer deploy came before this one's rollout settled.
43    Superseded,
44}
45
46impl Status {
47    pub fn finished(self) -> bool {
48        !matches!(self, Status::Queued | Status::Deploying)
49    }
50}
51
52/// One line of a deployment's log: an event of its stack.
53#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
54pub struct LogLine {
55    /// Unix milliseconds.
56    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/// One deploy of a compose stack.
64#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
65pub struct Deployment {
66    pub id: u64,
67    pub stack: String,
68    /// `manual` (the local CLI) or `api` (a tool call over MCP or REST).
69    pub trigger: String,
70    /// What deployed: `deploy` (a compose file), `rollback`, `env` (an
71    /// environment change) or `domains` (a managed domains change).
72    pub action: String,
73    /// The caller.
74    pub actor: String,
75    pub status: Status,
76    /// A rollback to a kept deployment: its id.
77    #[serde(default, skip_serializing_if = "Option::is_none")]
78    pub rollback_of: Option<u64>,
79    /// The services it created, changed, scaled or removed.
80    #[serde(default)]
81    pub services: Vec<String>,
82    /// `file:`/`environment:` secrets it deployed with the value an
83    /// earlier deploy stored, because it was given none.
84    #[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    /// Unix seconds, as are `started_at` and `finished_at`.
89    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    /// The compose source it deployed, as stack_export gives it.
95    #[serde(default)]
96    pub source: String,
97    /// Where the source's relative paths resolved.
98    #[serde(default)]
99    pub base_dir: PathBuf,
100    /// The stack's environment at the time (`.env` text).
101    #[serde(default)]
102    pub env: String,
103    /// The managed domains at the time, per service.
104    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
105    pub domains: BTreeMap<String, Vec<DomainSpec>>,
106    /// The stack's events while it ran.
107    #[serde(default)]
108    pub events: Vec<LogLine>,
109    /// The event feed's sequence number it has read up to.
110    #[serde(default)]
111    pub events_seq: u64,
112}
113
114impl Deployment {
115    /// The record for listings: without the source, environment, domains
116    /// and events.
117    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    /// The events as text, one line each.
135    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/// The per-stack settings and records, for every org.
162#[derive(Clone)]
163pub struct StackMeta {
164    dir: PathBuf,
165    /// Held across read-modify-write of a record.
166    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    /// Under the daemon's state directory.
198    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    /// Whether anything is kept for the stack.
219    pub fn exists(&self, org: &OrgId, name: &str) -> bool {
220        self.stack_dir(org, name).is_dir()
221    }
222
223    /// Forget everything kept for a removed stack.
224    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    // --- environment and domains -------------------------------------------
232
233    /// The stack's environment as stored (`.env` text; empty: none).
234    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    /// The stack's managed domains per service.
245    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    // --- deployments --------------------------------------------------------
267
268    /// The stack's deployments, newest first.
269    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    /// Record a new deployment, `deploying` from now: numbered after the
306    /// newest, which is superseded if it has not settled. Old records past
307    /// [`KEEP_DEPLOYMENTS`] are pruned.
308    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    /// Change a record under the lock; a finished one is left alone.
332    /// Returns it as saved.
333    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    /// Mark a deployment failed (the controller refused it).
353    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    /// Deployments a stopped daemon left unfinished are marked failed.
361    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    /// Follow deployment `id` of `org`/`name` in a thread until its rollout
397    /// settles, collecting the stack's events into it.
398    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
457/// Whether the stack's current deployment has settled: `None` while it
458/// rolls out, else what failed or paused (empty: all converged). As
459/// `wait_settled` judges it: only a status of the current revision and
460/// replica count counts.
461fn 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        // A finished record is not changed again.
528        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        // The summary leaves the bulk out.
541        let s = all[0].summary();
542        assert!(s.get("source").is_none() && s.get("events").is_none());
543        assert_eq!(s["status"], "deploying");
544        // Reused secrets are named only when there are some, and kept.
545        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        // Settings live beside the records and go with them.
563        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}