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//! the stack's events while it ran, and the output of each service's last
14//! replica that failed to come up during it.
15
16use 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
30/// Deployment records kept per stack, as for an app.
31pub const KEEP_DEPLOYMENTS: usize = 30;
32
33/// Events kept per record.
34const 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    /// A newer deploy came before this one's rollout settled.
44    Superseded,
45}
46
47impl Status {
48    pub fn finished(self) -> bool {
49        !matches!(self, Status::Queued | Status::Deploying)
50    }
51}
52
53/// One line of a deployment's log: an event of its stack.
54#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
55pub struct LogLine {
56    /// Unix milliseconds.
57    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/// One deploy of a compose stack.
65#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
66pub struct Deployment {
67    pub id: u64,
68    pub stack: String,
69    /// `manual` (the local CLI) or `api` (a tool call over MCP or REST).
70    pub trigger: String,
71    /// What deployed: `deploy` (a compose file), `rollback`, `env` (an
72    /// environment change) or `domains` (a managed domains change).
73    pub action: String,
74    /// The caller.
75    pub actor: String,
76    pub status: Status,
77    /// A rollback to a kept deployment: its id.
78    #[serde(default, skip_serializing_if = "Option::is_none")]
79    pub rollback_of: Option<u64>,
80    /// The services it created, changed, scaled or removed.
81    #[serde(default)]
82    pub services: Vec<String>,
83    /// `file:`/`environment:` secrets it deployed with the value an
84    /// earlier deploy stored, because it was given none.
85    #[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    /// Unix seconds, as are `started_at` and `finished_at`.
90    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    /// The compose source it deployed, as stack_export gives it.
96    #[serde(default)]
97    pub source: String,
98    /// Where the source's relative paths resolved.
99    #[serde(default)]
100    pub base_dir: PathBuf,
101    /// The stack's environment at the time (`.env` text).
102    #[serde(default)]
103    pub env: String,
104    /// The managed domains at the time, per service.
105    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
106    pub domains: BTreeMap<String, Vec<DomainSpec>>,
107    /// The stack's events while it ran.
108    #[serde(default)]
109    pub events: Vec<LogLine>,
110    /// The event feed's sequence number it has read up to.
111    #[serde(default)]
112    pub events_seq: u64,
113    /// Per service, the last replica that failed to come up while it rolled
114    /// out, with its output (each bounded by [`super::failure::OUTPUT_BYTES`]).
115    #[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
116    pub failed_attempts: BTreeMap<String, super::failure::FailedAttempt>,
117}
118
119impl Deployment {
120    /// The record for listings: without the source, environment, domains
121    /// and events.
122    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    /// The events as text, one line each.
141    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/// The per-stack settings and records, for every org.
168#[derive(Clone)]
169pub struct StackMeta {
170    dir: PathBuf,
171    /// Held across read-modify-write of a record.
172    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    /// Under the daemon's state directory.
204    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    /// Whether anything is kept for the stack.
225    pub fn exists(&self, org: &OrgId, name: &str) -> bool {
226        self.stack_dir(org, name).is_dir()
227    }
228
229    /// Forget everything kept for a removed stack.
230    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    // --- environment and domains -------------------------------------------
238
239    /// The stack's environment as stored (`.env` text; empty: none).
240    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    /// The stack's managed domains per service.
251    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    // --- deployments --------------------------------------------------------
273
274    /// The stack's deployments, newest first.
275    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    /// Record a new deployment, `deploying` from now: numbered after the
312    /// newest, which is superseded if it has not settled. Old records past
313    /// [`KEEP_DEPLOYMENTS`] are pruned.
314    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    /// Change a record under the lock; a finished one is left alone.
338    /// Returns it as saved.
339    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    /// Mark a deployment failed (the controller refused it).
359    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    /// Deployments a stopped daemon left unfinished are marked failed.
367    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    /// Follow deployment `id` of `org`/`name` in a thread until its rollout
403    /// settles, collecting the stack's events into it.
404    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            // After `settled`: a replica's attempt is kept before its
445            // service says it failed.
446            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
471/// Whether the stack's current deployment has settled: `None` while it
472/// rolls out, else what failed or paused (empty: all converged). As
473/// `wait_settled` judges it: only a status of the current revision and
474/// replica count counts.
475fn 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        // A finished record is not changed again.
543        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        // The summary leaves the bulk out.
556        let s = all[0].summary();
557        assert!(s.get("source").is_none() && s.get("events").is_none());
558        assert_eq!(s["status"], "deploying");
559        // Reused secrets are named only when there are some, and kept.
560        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        // Settings live beside the records and go with them.
578        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        // A record written before the field reads back without it.
600        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}