Skip to main content

isb_apps/app/
deploy.rs

1//! The app service: records on disk, the deploy queue and pipeline, and
2//! webhooks.
3//!
4//! A deploy is a record ([`Deployment`]) that moves `queued` → `building` →
5//! `deploying` → `done` | `failed`. One runs at a time per app; a deploy
6//! asked for while one runs waits behind it, and a newer request
7//! supersedes a waiting one that has not started (`superseded`), so a burst
8//! of pushes builds the latest commit once.
9//!
10//! The pipeline: resolve the image (an image source's digest, or fetch the
11//! git source and build it with [`crate::build::run`]), render the app into
12//! its project environment's stack in place of its old service, hand the
13//! stack to the controller, and wait for that service to converge.
14
15use std::collections::{BTreeMap, VecDeque};
16use std::io::Write;
17use std::path::{Path, PathBuf};
18use std::sync::{Arc, Mutex, OnceLock};
19use std::time::{Duration, Instant};
20
21use serde::{Deserialize, Serialize};
22use serde_json::{Value, json};
23
24use super::git::{self, Credentials, GitAuth};
25use super::{App, AppSpec, Rendered, Source, webhook};
26use crate::build::{BuildRequest, BuiltImage};
27use crate::client::Client;
28use crate::error::{Error, Result};
29use crate::org::OrgId;
30use crate::secrets::Secrets;
31use crate::stack::{Controller, StackDef};
32
33#[path = "db_secrets.rs"]
34mod db_secrets;
35
36/// Turns a checkout into an image: [`crate::build::run`], or a stand-in.
37pub type BuildFn =
38    Arc<dyn Fn(&Client, &BuildRequest, &mut dyn FnMut(&str)) -> Result<BuiltImage> + Send + Sync>;
39
40/// Asks an image's registry about it, waiting at most the given time:
41/// whether it exists, and its manifest digest
42/// ([`crate::image_check::probe`], or a stand-in).
43pub type ImageProbe = Arc<dyn Fn(&str, Duration) -> crate::image_check::Probe + Send + Sync>;
44
45/// Deployment records kept per app (older ones and their logs are pruned).
46const KEEP_DEPLOYMENTS: usize = 30;
47
48/// How long a deploy waits for its service to converge.
49const DEPLOY_TIMEOUT: Duration = Duration::from_secs(900);
50
51/// Who asked for a deployment.
52#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
53#[serde(rename_all = "lowercase")]
54pub enum Trigger {
55    /// The local CLI.
56    Manual,
57    /// A tool call over MCP or REST.
58    Api,
59    Webhook,
60}
61
62#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
63#[serde(rename_all = "lowercase")]
64pub enum Status {
65    Queued,
66    Building,
67    Deploying,
68    Done,
69    Failed,
70    /// A newer deploy replaced it before it started.
71    Superseded,
72    /// Closed before it started, because what queued it stopped.
73    Cancelled,
74}
75
76impl Status {
77    pub fn finished(self) -> bool {
78        !matches!(self, Status::Queued | Status::Building | Status::Deploying)
79    }
80
81    /// The moves the pipeline makes; anything else is a bug.
82    pub fn can_become(self, next: Status) -> bool {
83        use Status::*;
84        matches!(
85            (self, next),
86            (Queued, Building)
87                | (Queued, Superseded)
88                | (Queued, Cancelled)
89                | (Queued, Failed)
90                | (Building, Deploying)
91                | (Building, Failed)
92                | (Deploying, Done)
93                | (Deploying, Failed)
94        )
95    }
96}
97
98#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
99pub struct Commit {
100    pub sha: String,
101    pub message: String,
102}
103
104/// One deploy of an app.
105#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
106pub struct Deployment {
107    pub id: u64,
108    pub app: String,
109    pub trigger: Trigger,
110    /// The caller (a user, an API token, `local`, `webhook:github`).
111    pub by: String,
112    pub status: Status,
113    /// What a webhook said was pushed (`refs/heads/main abc123`).
114    #[serde(default, skip_serializing_if = "Option::is_none")]
115    pub requested: Option<String>,
116    /// A rollback: the deployment whose image and settings it restores.
117    #[serde(default, skip_serializing_if = "Option::is_none")]
118    pub rollback_of: Option<u64>,
119    #[serde(default, skip_serializing_if = "Option::is_none")]
120    pub commit: Option<Commit>,
121    #[serde(default, skip_serializing_if = "Option::is_none")]
122    pub image: Option<String>,
123    #[serde(default, skip_serializing_if = "Option::is_none")]
124    pub digest: Option<String>,
125    #[serde(default, skip_serializing_if = "Option::is_none")]
126    pub error: Option<String>,
127    /// Unix milliseconds.
128    pub created_at: u64,
129    #[serde(default, skip_serializing_if = "Option::is_none")]
130    pub started_at: Option<u64>,
131    #[serde(default, skip_serializing_if = "Option::is_none")]
132    pub finished_at: Option<u64>,
133    /// The service as deployed: what a rollback puts back.
134    #[serde(default, skip_serializing_if = "Option::is_none")]
135    pub rendered: Option<Rendered>,
136}
137
138impl Deployment {
139    /// Move to `next`, refusing a move the state machine does not have.
140    pub fn advance(&mut self, next: Status) -> Result<()> {
141        if !self.status.can_become(next) {
142            return Err(Error::invalid(format!(
143                "deployment {}: cannot go from {:?} to {next:?}",
144                self.id, self.status
145            )));
146        }
147        self.status = next;
148        let now = crate::stack::controller::now_ms();
149        match next {
150            Status::Building => self.started_at = Some(now),
151            s if s.finished() => self.finished_at = Some(now),
152            _ => {}
153        }
154        Ok(())
155    }
156
157    /// The record without the rendered service, for listings.
158    pub fn summary(&self) -> Value {
159        let mut v = serde_json::to_value(self).unwrap_or_default();
160        if let Some(o) = v.as_object_mut() {
161            o.remove("rendered");
162        }
163        v
164    }
165}
166
167#[derive(Default)]
168pub(super) struct Queue {
169    pub(super) running: bool,
170    pub(super) next: Option<u64>,
171}
172
173pub(super) struct Inner {
174    pub(super) state: PathBuf,
175    pub(super) client: Client,
176    pub(super) ctl: Controller,
177    pub(super) secrets: Arc<Secrets>,
178    pub(super) build: BuildFn,
179    pub(super) probe: ImageProbe,
180    pub(super) timeout: Duration,
181    /// Held across read-modify-write of app and project records.
182    pub(super) edit: Mutex<()>,
183    /// Held across read-splice-deploy of a stack, so two apps deploying
184    /// into one environment never drop each other's service.
185    pub(super) stacks: Mutex<()>,
186    /// Keyed by app, or `<app>#pr-<n>` for a preview.
187    pub(super) queues: Mutex<BTreeMap<(OrgId, String), Queue>>,
188    /// Recent webhook delivery ids, to ignore a replayed delivery.
189    deliveries: Mutex<VecDeque<String>>,
190    /// Told when a secret gets a new value outside `secret_set` (a
191    /// database's URL secrets, written at deploy), for what stacks don't
192    /// cover: the workspaces.
193    pub(super) secret_hook: OnceLock<SecretHook>,
194    /// Compose stack owners per org, read from the project records.
195    pub(super) owners: Mutex<BTreeMap<OrgId, Arc<super::compose::Owners>>>,
196}
197
198/// Called with an org and a secret's name after it got a new value.
199pub type SecretHook = Arc<dyn Fn(&OrgId, &str) + Send + Sync>;
200
201/// Projects, environments, apps and their deployments, for every org.
202#[derive(Clone)]
203pub struct Apps {
204    pub(super) inner: Arc<Inner>,
205}
206
207impl Apps {
208    /// Open the app store under `state`. Deployments a stopped daemon left
209    /// unfinished are marked failed.
210    pub fn new(state: &Path, client: Client, ctl: Controller, secrets: Arc<Secrets>) -> Apps {
211        let a = Apps {
212            inner: Arc::new(Inner {
213                state: state.to_path_buf(),
214                client,
215                ctl,
216                secrets,
217                build: Arc::new(|c: &Client, r: &BuildRequest, l: &mut dyn FnMut(&str)| {
218                    crate::build::run(c, r, l)
219                }),
220                probe: Arc::new(crate::image_check::probe),
221                timeout: DEPLOY_TIMEOUT,
222                edit: Mutex::new(()),
223                stacks: Mutex::new(()),
224                queues: Mutex::new(BTreeMap::new()),
225                deliveries: Mutex::new(VecDeque::new()),
226                secret_hook: OnceLock::new(),
227                owners: Mutex::new(BTreeMap::new()),
228            }),
229        };
230        a.recover();
231        a
232    }
233
234    fn with(self, f: impl FnOnce(&mut Inner)) -> Apps {
235        let mut inner = Arc::try_unwrap(self.inner)
236            .unwrap_or_else(|_| panic!("Apps::with_* before the Apps is shared"));
237        f(&mut inner);
238        Apps {
239            inner: Arc::new(inner),
240        }
241    }
242
243    /// Use `f` instead of [`crate::build::run`] (tests).
244    pub fn with_build(self, f: BuildFn) -> Apps {
245        self.with(|i| i.build = f)
246    }
247
248    /// Use `f` to ask registries about images instead of `skopeo inspect`.
249    pub fn with_probe(self, f: ImageProbe) -> Apps {
250        self.with(|i| i.probe = f)
251    }
252
253    pub fn with_timeout(self, t: Duration) -> Apps {
254        self.with(|i| i.timeout = t)
255    }
256
257    // --- paths -----------------------------------------------------------
258
259    pub(super) fn apps_dir(&self, org: &OrgId) -> PathBuf {
260        super::org_root(&self.inner.state, org).join("apps")
261    }
262
263    pub(super) fn app_dir(&self, org: &OrgId, app: &str) -> PathBuf {
264        self.apps_dir(org).join(app)
265    }
266
267    fn deployments_dir(&self, org: &OrgId, app: &str) -> PathBuf {
268        self.app_dir(org, app).join("deployments")
269    }
270
271    /// Where an app's git checkout lives.
272    pub fn source_dir(&self, org: &OrgId, app: &str) -> PathBuf {
273        super::org_root(&self.inner.state, org)
274            .join("sources")
275            .join(app)
276    }
277
278    // --- apps ------------------------------------------------------------
279
280    pub fn get(&self, org: &OrgId, name: &str) -> Result<App> {
281        super::validate_app_name(name)?;
282        match std::fs::read(self.app_dir(org, name).join("app.json")) {
283            Ok(b) => Ok(serde_json::from_slice(&b)?),
284            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
285                Err(Error::NotFound(format!("app {name} in org {org}")))
286            }
287            Err(e) => Err(e.into()),
288        }
289    }
290
291    pub fn list(&self, org: &OrgId) -> Result<Vec<App>> {
292        let mut out = Vec::new();
293        let Ok(rd) = std::fs::read_dir(self.apps_dir(org)) else {
294            return Ok(out);
295        };
296        for e in rd.flatten() {
297            let p = e.path().join("app.json");
298            if !p.is_file() {
299                continue;
300            }
301            match serde_json::from_slice::<App>(&std::fs::read(&p)?) {
302                Ok(a) => out.push(a),
303                Err(e) => eprintln!("isb serve: skipping {}: {e}", p.display()),
304            }
305        }
306        out.sort_by(|a, b| a.spec.name.cmp(&b.spec.name));
307        Ok(out)
308    }
309
310    fn save(&self, org: &OrgId, app: &App) -> Result<()> {
311        super::write_atomic(
312            &self.app_dir(org, &app.spec.name).join("app.json"),
313            &serde_json::to_vec_pretty(app)?,
314        )
315    }
316
317    /// Create an app (not deployed yet) and its webhook secret. Returns the
318    /// app and that secret.
319    pub fn create(&self, org: &OrgId, mut spec: AppSpec) -> Result<(App, String)> {
320        if let Source::Database(db) = &mut spec.source {
321            db.normalize(&spec.name);
322        }
323        spec.validate()?;
324        let _g = self.inner.edit.lock().unwrap();
325        let proj = self.project_get(org, &spec.project)?;
326        if !proj.environments.contains(&spec.environment) {
327            return Err(Error::NotFound(format!(
328                "environment {} in project {} (it has {})",
329                spec.environment,
330                spec.project,
331                proj.environments.join(", ")
332            )));
333        }
334        if self.app_dir(org, &spec.name).join("app.json").exists() {
335            return Err(Error::AlreadyExists(format!("app {}", spec.name)));
336        }
337        self.check_app_name_free(org, &spec.project, &spec.environment, &spec.name)?;
338        self.check_secrets(org, &spec)?;
339        let now = crate::stack::now_secs();
340        let app = App {
341            spec,
342            created_at: now,
343            updated_at: now,
344            next_deployment: 1,
345            current: None,
346        };
347        if let Source::Database(db) = &app.spec.source {
348            self.database_credentials(org, &app.spec, db)?;
349        }
350        let secret = git::random_hex(32);
351        self.inner
352            .secrets
353            .set(org, &app.spec.webhook_secret(), secret.as_bytes())?;
354        self.save(org, &app)?;
355        Ok((app, secret))
356    }
357
358    /// The org's secret store (backups read destination credentials and
359    /// database passwords from it).
360    pub fn secrets(&self) -> &Arc<Secrets> {
361        &self.inner.secrets
362    }
363
364    /// The incus client the apps run on.
365    pub fn client(&self) -> &Client {
366        &self.inner.client
367    }
368
369    /// The stack controller (events, definitions).
370    pub fn controller(&self) -> &Controller {
371        &self.inner.ctl
372    }
373
374    /// Wait (up to the deploy timeout) for an app's service to converge at
375    /// its current revision: `(converged, why not)`.
376    pub fn wait_converged(&self, org: &OrgId, name: &str) -> Result<(bool, String)> {
377        let app = self.get(org, name)?;
378        let q = crate::stack::qualified(org, &app.spec.stack()?);
379        self.wait_service(&q, name)
380    }
381
382    /// Secrets an app names must exist (a typo is better caught now than
383    /// at deploy).
384    fn check_secrets(&self, org: &OrgId, spec: &AppSpec) -> Result<()> {
385        let mut names = spec.env.secret_names();
386        if let Some(p) = &spec.previews {
387            names.extend(p.secret_names());
388        }
389        names.extend(spec.files.iter().map(|f| f.secret.clone()));
390        if let Source::Git(g) = &spec.source {
391            names.extend(g.auth.secret().map(String::from));
392        }
393        for n in names {
394            self.inner.secrets.inspect(org, &n).map_err(|e| match e {
395                Error::NotFound(_) => Error::invalid(format!(
396                    "secret {n} does not exist in org {org}; create it with `isb secret create {n} --org {org}`"
397                )),
398                e => e,
399            })?;
400        }
401        Ok(())
402    }
403
404    /// What `create` and `update` check about the secrets a spec names,
405    /// for callers that check without writing (`app_apply` with
406    /// `dry_run`).
407    pub fn check_spec(&self, org: &OrgId, spec: &AppSpec) -> Result<()> {
408        self.check_secrets(org, spec)
409    }
410
411    /// Change an app's settings with a JSON merge patch (`null` clears a
412    /// setting). Name, project and environment are fixed. Takes effect at
413    /// the next deploy.
414    pub fn update(&self, org: &OrgId, name: &str, patch: &Value) -> Result<App> {
415        let _g = self.inner.edit.lock().unwrap();
416        let mut app = self.get(org, name)?;
417        let mut v = serde_json::to_value(&app.spec)?;
418        super::merge_patch(&mut v, patch);
419        let mut spec: AppSpec =
420            serde_json::from_value(v).map_err(|e| Error::invalid(format!("app {name}: {e}")))?;
421        super::manifest::check_update(&app.spec, &mut spec)?;
422        spec.validate()?;
423        self.check_secrets(org, &spec)?;
424        app.spec = spec;
425        app.updated_at = crate::stack::now_secs();
426        self.save(org, &app)?;
427        Ok(app)
428    }
429
430    pub fn env_get(&self, org: &OrgId, name: &str) -> Result<String> {
431        Ok(self.get(org, name)?.spec.env.render())
432    }
433
434    /// Replace an app's environment with `.env` text.
435    pub fn env_set(&self, org: &OrgId, name: &str, text: &str) -> Result<App> {
436        let env = super::EnvFile::parse(text)?;
437        self.update(org, name, &json!({"env": env.render()}))
438    }
439
440    /// The app's webhook secret; with `rotate`, a new one first.
441    pub fn webhook_secret(&self, org: &OrgId, name: &str, rotate: bool) -> Result<String> {
442        let app = self.get(org, name)?;
443        let key = app.spec.webhook_secret();
444        if !rotate {
445            if let Ok((v, _)) = self.inner.secrets.get(org, &key) {
446                return String::from_utf8(v)
447                    .map_err(|_| Error::invalid("webhook secret is not text"));
448            }
449        }
450        let s = git::random_hex(32);
451        self.inner.secrets.set(org, &key, s.as_bytes())?;
452        Ok(s)
453    }
454
455    /// A new ed25519 deploy key for a git app over SSH: stored as the org
456    /// secret `app.<app>.deploy-key`, made the app's credential, and its
457    /// public half returned to paste into the repository's deploy keys.
458    pub fn deploy_key(&self, org: &OrgId, name: &str) -> Result<String> {
459        let app = self.get(org, name)?;
460        let Source::Git(g) = &app.spec.source else {
461            return Err(Error::invalid(format!("app {name} has no git source")));
462        };
463        if git::transport(&g.url)? != git::Transport::Ssh {
464            return Err(Error::invalid(format!(
465                "a deploy key needs an SSH URL; {} is not one (git@host:owner/repo)",
466                g.url
467            )));
468        }
469        let (private, public) =
470            git::generate_deploy_key(&self.apps_dir(org), &format!("isb-deploy-{org}-{name}"))?;
471        let key = super::deploy_key_secret(name);
472        self.inner.secrets.set(org, &key, &private)?;
473        self.update(
474            org,
475            name,
476            &json!({"source": {"git": {"auth": {"ssh_key_secret": key}}}}),
477        )?;
478        Ok(public)
479    }
480
481    /// Delete an app: its service leaves the stack (the stack goes when it
482    /// was the last), then its records, checkout and the secrets isb made
483    /// for it. Named volumes are kept.
484    pub fn delete(&self, org: &OrgId, name: &str) -> Result<()> {
485        let app = self.get(org, name)?;
486        if self
487            .inner
488            .queues
489            .lock()
490            .unwrap()
491            .contains_key(&(org.clone(), name.to_string()))
492        {
493            return Err(Error::invalid(format!(
494                "app {name} is deploying; delete it once that finishes"
495            )));
496        }
497        // Its previews go first, with their stacks, volumes and images.
498        self.previews_remove_all(org, name)?;
499        let stack = app.spec.stack()?;
500        let q = crate::stack::qualified(org, &stack);
501        {
502            let _s = self.inner.stacks.lock().unwrap();
503            if let Ok(cur) = self.inner.ctl.definition(&q) {
504                if cur.file.services.contains_key(name) {
505                    let file = super::splice(Some(&cur.file), &stack, name, None);
506                    if file.services.is_empty() {
507                        self.inner.ctl.remove(&q, false, Duration::from_secs(300))?;
508                    } else {
509                        let def = self.stack_def(org, &stack, file, "app delete")?;
510                        self.inner.ctl.deploy(def)?;
511                    }
512                }
513            }
514        }
515        let _g = self.inner.edit.lock().unwrap();
516        for s in [super::webhook_secret(name), super::deploy_key_secret(name)] {
517            match self.inner.secrets.delete(org, &s) {
518                Ok(()) => {}
519                Err(e) if e.is_not_found() => {}
520                Err(e) => return Err(e),
521            }
522        }
523        let _ = std::fs::remove_dir_all(self.source_dir(org, name));
524        std::fs::remove_dir_all(self.app_dir(org, name))?;
525        Ok(())
526    }
527
528    // --- deployments -----------------------------------------------------
529
530    fn dep_path(&self, org: &OrgId, app: &str, id: u64) -> PathBuf {
531        self.deployments_dir(org, app).join(format!("{id}.json"))
532    }
533
534    fn log_path(&self, org: &OrgId, app: &str, id: u64) -> PathBuf {
535        self.deployments_dir(org, app).join(format!("{id}.log"))
536    }
537
538    pub(super) fn save_dep(&self, org: &OrgId, d: &Deployment) -> Result<()> {
539        super::write_atomic(
540            &self.dep_path(org, &d.app, d.id),
541            &serde_json::to_vec_pretty(d)?,
542        )
543    }
544
545    pub fn deployment(&self, org: &OrgId, app: &str, id: u64) -> Result<Deployment> {
546        super::validate_app_name(app)?;
547        match std::fs::read(self.dep_path(org, app, id)) {
548            Ok(b) => Ok(serde_json::from_slice(&b)?),
549            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
550                Err(Error::NotFound(format!("deployment {id} of app {app}")))
551            }
552            Err(e) => Err(e.into()),
553        }
554    }
555
556    /// Newest first.
557    pub fn deployments(&self, org: &OrgId, app: &str) -> Result<Vec<Deployment>> {
558        self.get(org, app)?;
559        let mut ids: Vec<u64> = match std::fs::read_dir(self.deployments_dir(org, app)) {
560            Ok(rd) => rd
561                .flatten()
562                .filter_map(|e| {
563                    let n = e.file_name().into_string().ok()?;
564                    n.strip_suffix(".json")?.parse().ok()
565                })
566                .collect(),
567            Err(_) => vec![],
568        };
569        ids.sort_unstable_by(|a, b| b.cmp(a));
570        ids.into_iter()
571            .map(|id| self.deployment(org, app, id))
572            .collect()
573    }
574
575    /// A deployment's log from byte `offset`: `(text, next offset, finished)`.
576    pub fn log(&self, org: &OrgId, app: &str, id: u64, offset: u64) -> Result<(String, u64, bool)> {
577        let (text, next, d) = self.log_and_record(org, app, id, offset)?;
578        Ok((text, next, d.status.finished()))
579    }
580
581    /// [`Apps::log`] with the deployment record read just before the log,
582    /// so a client follows the text and the status in one call: a record
583    /// that says finished means the text holds every line.
584    pub fn log_and_record(
585        &self,
586        org: &OrgId,
587        app: &str,
588        id: u64,
589        offset: u64,
590    ) -> Result<(String, u64, Deployment)> {
591        let d = self.deployment(org, app, id)?;
592        let b = std::fs::read(self.log_path(org, app, id)).unwrap_or_default();
593        let start = (offset as usize).min(b.len());
594        Ok((
595            String::from_utf8_lossy(&b[start..]).into_owned(),
596            b.len() as u64,
597            d,
598        ))
599    }
600
601    /// Wait until a deployment finishes, or `timeout`.
602    pub fn wait(&self, org: &OrgId, app: &str, id: u64, timeout: Duration) -> Result<Deployment> {
603        let started = Instant::now();
604        loop {
605            let d = self.deployment(org, app, id)?;
606            if d.status.finished() || started.elapsed() >= timeout {
607                return Ok(d);
608            }
609            std::thread::sleep(Duration::from_millis(500));
610        }
611    }
612
613    /// Queue a deploy of the app's current settings.
614    pub fn deploy(
615        &self,
616        org: &OrgId,
617        name: &str,
618        trigger: Trigger,
619        by: &str,
620        requested: Option<String>,
621    ) -> Result<Deployment> {
622        self.enqueue(org, name, trigger, by, requested, None)
623    }
624
625    /// Queue a rollback to deployment `to` (default: the one before the
626    /// current): its image and settings, without building.
627    pub fn rollback(
628        &self,
629        org: &OrgId,
630        name: &str,
631        to: Option<u64>,
632        trigger: Trigger,
633        by: &str,
634    ) -> Result<Deployment> {
635        let app = self.get(org, name)?;
636        let target = match to {
637            Some(id) => self.deployment(org, name, id)?,
638            None => self
639                .deployments(org, name)?
640                .into_iter()
641                .find(|d| d.status == Status::Done && Some(d.id) != app.current)
642                .ok_or_else(|| {
643                    Error::invalid(format!(
644                        "app {name} has no earlier successful deployment to roll back to"
645                    ))
646                })?,
647        };
648        if target.status != Status::Done || target.rendered.is_none() {
649            return Err(Error::invalid(format!(
650                "deployment {} did not finish done; only those can be rolled back to",
651                target.id
652            )));
653        }
654        self.enqueue(org, name, trigger, by, None, Some(target.id))
655    }
656
657    fn enqueue(
658        &self,
659        org: &OrgId,
660        name: &str,
661        trigger: Trigger,
662        by: &str,
663        requested: Option<String>,
664        rollback_of: Option<u64>,
665    ) -> Result<Deployment> {
666        let dep = {
667            let _g = self.inner.edit.lock().unwrap();
668            let mut app = self.get(org, name)?;
669            let id = app.next_deployment;
670            app.next_deployment += 1;
671            self.save(org, &app)?;
672            let d = Deployment {
673                id,
674                app: name.into(),
675                trigger,
676                by: by.into(),
677                status: Status::Queued,
678                requested,
679                rollback_of,
680                commit: None,
681                image: None,
682                digest: None,
683                error: None,
684                created_at: crate::stack::controller::now_ms(),
685                started_at: None,
686                finished_at: None,
687                rendered: None,
688            };
689            self.save_dep(org, &d)?;
690            self.prune(org, &app);
691            d
692        };
693        self.event(
694            org,
695            name,
696            "info",
697            format!("deployment {} queued by {by}", dep.id),
698        );
699        let start = {
700            let mut qs = self.inner.queues.lock().unwrap();
701            let q = qs.entry((org.clone(), name.to_string())).or_default();
702            if let Some(old) = q.next.replace(dep.id) {
703                if let Ok(mut d) = self.deployment(org, name, old) {
704                    if d.advance(Status::Superseded).is_ok() {
705                        d.error = Some(format!("superseded by deployment {}", dep.id));
706                        let _ = self.save_dep(org, &d);
707                    }
708                }
709            }
710            !std::mem::replace(&mut q.running, true)
711        };
712        if start {
713            let me = self.clone();
714            let (org, name) = (org.clone(), name.to_string());
715            std::thread::spawn(move || me.drain(&org, &name));
716        }
717        Ok(dep)
718    }
719
720    /// Run queued deployments of one app until none is left.
721    fn drain(&self, org: &OrgId, name: &str) {
722        loop {
723            let next = {
724                let mut qs = self.inner.queues.lock().unwrap();
725                let key = (org.clone(), name.to_string());
726                let q = qs.entry(key.clone()).or_default();
727                match q.next.take() {
728                    Some(id) => id,
729                    None => {
730                        qs.remove(&key);
731                        return;
732                    }
733                }
734            };
735            self.run(org, name, next);
736        }
737    }
738
739    fn run(&self, org: &OrgId, name: &str, id: u64) {
740        let Ok(mut dep) = self.deployment(org, name, id) else {
741            return;
742        };
743        let mut log = match DeployLog::open(self, org, name, id) {
744            Ok(l) => l,
745            Err(e) => {
746                eprintln!("isb serve: app {name}: deployment {id}: no log: {e}");
747                return;
748            }
749        };
750        let r = self.pipeline(org, name, &mut dep, &mut log);
751        let (kind, level, msg) = match &r {
752            Ok(()) => ("deploy.succeeded", "info", format!("deployment {id} done")),
753            Err(e) => (
754                "deploy.failed",
755                "error",
756                format!("deployment {id} failed: {e}"),
757            ),
758        };
759        log.line(&msg);
760        if let Err(e) = r {
761            dep.error = Some(e.to_string());
762            let _ = dep.advance(Status::Failed);
763        }
764        if let Err(e) = self.save_dep(org, &dep) {
765            eprintln!("isb serve: app {name}: deployment {id}: {e}");
766        }
767        self.kind_event(org, name, kind, level, msg);
768    }
769
770    fn set_status(&self, org: &OrgId, dep: &mut Deployment, s: Status) -> Result<()> {
771        dep.advance(s)?;
772        self.save_dep(org, dep)?;
773        self.event(
774            org,
775            &dep.app,
776            "info",
777            format!("deployment {}: {s:?}", dep.id).to_lowercase(),
778        );
779        Ok(())
780    }
781
782    #[expect(
783        clippy::too_many_lines,
784        reason = "predates the lint ratchet; split it when next changed"
785    )]
786    fn pipeline(
787        &self,
788        org: &OrgId,
789        name: &str,
790        dep: &mut Deployment,
791        log: &mut DeployLog,
792    ) -> Result<()> {
793        self.set_status(org, dep, Status::Building)?;
794        let app = self.get(org, name)?;
795        let mut notes = Vec::new();
796        let rendered = match dep.rollback_of {
797            Some(of) => {
798                let t = self.deployment(org, name, of)?;
799                log.line(&format!(
800                    "rolling back to deployment {of} ({})",
801                    t.image.as_deref().unwrap_or("?")
802                ));
803                dep.image = t.image;
804                dep.digest = t.digest;
805                dep.commit = t.commit;
806                t.rendered
807                    .ok_or_else(|| Error::invalid(format!("deployment {of} kept no settings")))?
808            }
809            None => {
810                let (image, digest) = match &app.spec.source {
811                    Source::Image(i) => {
812                        log.line(&format!("image {i}"));
813                        self.resolve(i, log)?
814                    }
815                    Source::Database(db) => {
816                        let i = db.image();
817                        log.line(&format!(
818                            "database {} {}: image {i}",
819                            db.engine,
820                            db.version()
821                        ));
822                        // A `urls` entry added since the last deploy is
823                        // written now, not at the next password change, and
824                        // what uses a value that moved follows it.
825                        for n in self.write_urls(org, &app.spec, db)? {
826                            log.line(&format!("secret {n}: the connection URL"));
827                            self.url_secret_moved(org, &n, log);
828                        }
829                        self.resolve(&i, log)?
830                    }
831                    Source::Git(g) => {
832                        let creds = self.credentials(org, &g.auth)?;
833                        let co = git::fetch(g, &creds, &self.source_dir(org, name), &mut |l| {
834                            log.line(l)
835                        })?;
836                        dep.commit = Some(Commit {
837                            sha: co.sha.clone(),
838                            message: co.message.clone(),
839                        });
840                        self.save_dep(org, dep)?;
841                        let b =
842                            app.spec.build.as_ref().ok_or_else(|| {
843                                Error::invalid("a git source needs build settings")
844                            })?;
845                        let req = BuildRequest {
846                            org: org.clone(),
847                            app: name.into(),
848                            context: co.dir.clone(),
849                            subdir: g.subdir.clone(),
850                            builder: b.builder.clone(),
851                            args: b.args.iter().map(|(k, v)| (k.clone(), v.clone())).collect(),
852                            tag: co.sha.clone(),
853                            untrusted: b.untrusted,
854                            cache: None,
855                        };
856                        log.line(&format!("building {} with {:?}", co.sha, b.builder));
857                        let built =
858                            (self.inner.build)(&self.inner.client, &req, &mut |l| log.line(l))?;
859                        log.line(&format!("built {} ({})", built.image, built.digest));
860                        (built.image, Some(built.digest))
861                    }
862                };
863                dep.image = Some(image.clone());
864                dep.digest = digest;
865                super::render(&app.spec, &image, &mut notes)?
866            }
867        };
868        for n in &notes {
869            log.line(&format!("note: {n}"));
870        }
871        dep.rendered = Some(rendered.clone());
872        self.set_status(org, dep, Status::Deploying)?;
873        let stack = app.spec.stack()?;
874        let q = crate::stack::qualified(org, &stack);
875        let mark = self.event_mark();
876        {
877            let _s = self.inner.stacks.lock().unwrap();
878            let cur = self.inner.ctl.definition(&q).ok();
879            let file = super::splice(cur.as_ref().map(|d| &d.file), &stack, name, Some(&rendered));
880            let def = self.stack_def(org, &stack, file, &dep.by)?;
881            let changes = self.inner.ctl.deploy(def)?;
882            let mine = changes.iter().find(|c| c.service == name);
883            match mine {
884                Some(c) if c.change == "unchanged" => {
885                    // Deploy means fresh instances, as in Dokploy: a moved
886                    // tag or a restart the settings do not show.
887                    log.line("settings unchanged: replacing the instances anyway");
888                    self.inner.ctl.redeploy(&q, name)?;
889                }
890                Some(c) => log.line(&format!(
891                    "stack {stack}: {} {} (rev {}, {} replicas)",
892                    c.service, c.change, c.rev, c.replicas
893                )),
894                None => {}
895            }
896        }
897        let (ok, msg) = self.wait_service_with(&q, name, Some(mark), |m| log.line(m))?;
898        if !ok {
899            return Err(Error::invalid(format!("service {name}: {msg}")));
900        }
901        log.line(&format!("service {name}.{stack} converged"));
902        self.set_status(org, dep, Status::Done)?;
903        let _g = self.inner.edit.lock().unwrap();
904        let mut app = self.get(org, name)?;
905        app.current = Some(dep.id);
906        self.save(org, &app)?;
907        Ok(())
908    }
909
910    /// The stack definition for `file`, with its secrets bound.
911    pub(super) fn stack_def(
912        &self,
913        org: &OrgId,
914        stack: &str,
915        file: crate::spec::ComposeFile,
916        by: &str,
917    ) -> Result<StackDef> {
918        let base = self.apps_dir(org);
919        std::fs::create_dir_all(&base)?;
920        let secrets = crate::stack::secrets::bind(
921            &self.inner.secrets,
922            org,
923            stack,
924            &file,
925            &BTreeMap::new(),
926            false,
927        )?;
928        Ok(StackDef {
929            source: None,
930            domains: Default::default(),
931            name: stack.into(),
932            org: org.clone(),
933            file,
934            base_dir: base,
935            secrets,
936            force: BTreeMap::new(),
937            images: BTreeMap::new(),
938            deployed_at: crate::stack::now_secs(),
939            deployed_by: by.into(),
940            previous: None,
941        })
942    }
943
944    /// Wait for one service of a stack to settle at its current revision:
945    /// `(converged, why not)`.
946    pub(super) fn wait_service(&self, q: &str, svc: &str) -> Result<(bool, String)> {
947        self.wait_service_with(q, svc, None, |_| {})
948    }
949
950    /// The controller's latest event number: rollout events after it
951    /// belong to a deploy that starts now.
952    pub(super) fn event_mark(&self) -> u64 {
953        self.inner.ctl.events(u64::MAX, 0).0
954    }
955
956    /// [`Apps::wait_service`], handing each controller event about the
957    /// service after `since` (the rollout: slots created, probed, serving,
958    /// old ones drained) to `seen` as it happens, so a deployment's log
959    /// tells the whole story. App events (`app NAME: ...`) are skipped:
960    /// they are the log's own lines.
961    pub(super) fn wait_service_with(
962        &self,
963        q: &str,
964        svc: &str,
965        mut since: Option<u64>,
966        mut seen: impl FnMut(&str),
967    ) -> Result<(bool, String)> {
968        let started = Instant::now();
969        let own = format!("app {svc}");
970        let mut relay = |since: &mut Option<u64>| {
971            let Some(after) = *since else { return };
972            let (last, evs) = self.inner.ctl.events(after, 500);
973            for e in evs.iter().filter(|e| e.stack == q && e.service == svc) {
974                if !e.message.starts_with(&own) {
975                    seen(&e.message);
976                }
977            }
978            *since = Some(last.max(after));
979        };
980        loop {
981            relay(&mut since);
982            let def = self.inner.ctl.definition(q)?;
983            let rev = def.revision(svc)?;
984            let replicas = def.service(svc)?.replicas();
985            let st = self.inner.ctl.status(q)?;
986            if let Some(s) = st.services.iter().find(|s| s.service == svc) {
987                if s.rev == rev && s.replicas == replicas {
988                    match s.state.as_str() {
989                        "converged" => {
990                            relay(&mut since);
991                            return Ok((true, String::new()));
992                        }
993                        "paused" | "failing" => {
994                            relay(&mut since);
995                            return Ok((
996                                false,
997                                format!("{}: {}", s.state, s.message.clone().unwrap_or_default()),
998                            ));
999                        }
1000                        _ => {}
1001                    }
1002                }
1003            }
1004            if started.elapsed() >= self.inner.timeout {
1005                return Ok((
1006                    false,
1007                    format!("not converged after {:?}", self.inner.timeout),
1008                ));
1009            }
1010            std::thread::sleep(Duration::from_millis(500));
1011        }
1012    }
1013
1014    pub(super) fn credentials(&self, org: &OrgId, auth: &GitAuth) -> Result<Credentials> {
1015        let read = |n: &str| {
1016            self.inner
1017                .secrets
1018                .get(org, n)
1019                .map(|(v, _)| v)
1020                .map_err(|e| Error::invalid(format!("git credential {n}: {e}")))
1021        };
1022        Ok(match auth {
1023            GitAuth::None => Credentials::None,
1024            GitAuth::Token {
1025                token_secret,
1026                username,
1027            } => Credentials::Token {
1028                username: username.clone().unwrap_or_else(|| "x-access-token".into()),
1029                token: String::from_utf8(read(token_secret)?)
1030                    .map_err(|_| Error::invalid("the git token is not text"))?
1031                    .trim()
1032                    .to_string(),
1033            },
1034            GitAuth::SshKey { ssh_key_secret } => Credentials::SshKey {
1035                private_key: read(ssh_key_secret)?,
1036            },
1037        })
1038    }
1039
1040    /// Drop the oldest records beyond [`KEEP_DEPLOYMENTS`], never the
1041    /// current one.
1042    fn prune(&self, org: &OrgId, app: &App) {
1043        let Ok(rd) = std::fs::read_dir(self.deployments_dir(org, &app.spec.name)) else {
1044            return;
1045        };
1046        let mut ids: Vec<u64> = rd
1047            .flatten()
1048            .filter_map(|e| {
1049                e.file_name()
1050                    .into_string()
1051                    .ok()?
1052                    .strip_suffix(".json")?
1053                    .parse()
1054                    .ok()
1055            })
1056            .collect();
1057        ids.sort_unstable();
1058        let excess = ids.len().saturating_sub(KEEP_DEPLOYMENTS);
1059        for id in ids.into_iter().take(excess) {
1060            if Some(id) == app.current {
1061                continue;
1062            }
1063            let _ = std::fs::remove_file(self.dep_path(org, &app.spec.name, id));
1064            let _ = std::fs::remove_file(self.log_path(org, &app.spec.name, id));
1065        }
1066    }
1067
1068    /// An event about an app on the daemon's feed, under its stack.
1069    pub(super) fn event(&self, org: &OrgId, app: &str, level: &str, message: String) {
1070        let stack = self
1071            .get(org, app)
1072            .ok()
1073            .and_then(|a| a.spec.stack().ok())
1074            .unwrap_or_default();
1075        let q = crate::stack::qualified(org, &stack);
1076        self.inner
1077            .ctl
1078            .note_service(level, &q, app, format!("app {app}: {message}"));
1079    }
1080
1081    /// [`Apps::event`] with a kind ([`crate::stack::controller::Event::kind`]).
1082    fn kind_event(&self, org: &OrgId, app: &str, kind: &str, level: &str, message: String) {
1083        let stack = self
1084            .get(org, app)
1085            .ok()
1086            .and_then(|a| a.spec.stack().ok())
1087            .unwrap_or_default();
1088        let q = crate::stack::qualified(org, &stack);
1089        self.inner
1090            .ctl
1091            .event(kind, level, &q, app, format!("app {app}: {message}"));
1092    }
1093
1094    // --- webhooks --------------------------------------------------------
1095
1096    /// Answer `POST /api/v1/webhooks/<org>/<app>`: `(status, body)`. Nothing
1097    /// happens without a valid signature or token; an unknown org or app
1098    /// answers as a bad signature does, so the endpoint maps nothing.
1099    pub fn webhook(
1100        &self,
1101        org: &str,
1102        app: &str,
1103        header: &dyn Fn(&str) -> Option<String>,
1104        token: Option<&str>,
1105        body: &[u8],
1106    ) -> (u16, Value) {
1107        let refuse = |why: &str| (401, json!({"error": "unauthorized", "message": why}));
1108        let (Ok(org), Ok(())) = (OrgId::new(org), super::validate_app_name(app)) else {
1109            return refuse("invalid signature");
1110        };
1111        let Ok(found) = self.get(&org, app) else {
1112            return refuse("invalid signature");
1113        };
1114        let Ok((secret, _)) = self.inner.secrets.get(&org, &found.spec.webhook_secret()) else {
1115            return refuse("invalid signature");
1116        };
1117        let provider = match webhook::verify(trim_ascii(&secret), header, token, body) {
1118            Ok(p) => p,
1119            Err(webhook::Refusal::Missing) => {
1120                return refuse(
1121                    "sign the request (X-Hub-Signature-256, X-Gitea-Signature, X-Gitlab-Token) or pass ?token=",
1122                );
1123            }
1124            Err(webhook::Refusal::Invalid) => return refuse("invalid signature"),
1125        };
1126        if let Some(id) = webhook::delivery_id(header) {
1127            let key = format!("{org}/{app}/{id}");
1128            let mut seen = self.inner.deliveries.lock().unwrap();
1129            if seen.contains(&key) {
1130                return (200, json!({"ignored": "delivery already received"}));
1131            }
1132            if seen.len() >= 1000 {
1133                seen.pop_front();
1134            }
1135            seen.push_back(key);
1136        }
1137        let by = format!("webhook:{provider:?}").to_lowercase();
1138        let requested = match webhook::event(provider, header, body) {
1139            webhook::Event::Ping => return (200, json!({"ok": true, "ping": true})),
1140            webhook::Event::Other(kind) => {
1141                return (200, json!({"ignored": format!("event {kind}")}));
1142            }
1143            webhook::Event::Trigger => None,
1144            webhook::Event::PullRequest(pr) => {
1145                return self.preview_webhook(&org, &found, provider, &by, &pr);
1146            }
1147            webhook::Event::Push {
1148                reference,
1149                after,
1150                message,
1151                deleted,
1152            } => {
1153                if deleted {
1154                    return (200, json!({"ignored": format!("{reference} was deleted")}));
1155                }
1156                if let Source::Git(g) = &found.spec.source {
1157                    if !g.matches_push(&reference) {
1158                        return (
1159                            200,
1160                            json!({"ignored": format!("{reference} is not {}", g.reference)}),
1161                        );
1162                    }
1163                }
1164                let mut r = reference;
1165                if let Some(a) = after {
1166                    r = format!("{r} {a}");
1167                }
1168                if let Some(m) = message {
1169                    r = format!("{r}: {m}");
1170                }
1171                Some(r)
1172            }
1173        };
1174        match self.deploy(&org, app, Trigger::Webhook, &by, requested) {
1175            Ok(d) => (202, json!({"deployment": d.id, "status": d.status})),
1176            Err(e) => (500, json!({"error": "deploy", "message": e.to_string()})),
1177        }
1178    }
1179}
1180
1181fn trim_ascii(b: &[u8]) -> &[u8] {
1182    let s = b
1183        .iter()
1184        .position(|c| !c.is_ascii_whitespace())
1185        .unwrap_or(b.len());
1186    let e = b
1187        .iter()
1188        .rposition(|c| !c.is_ascii_whitespace())
1189        .map_or(s, |i| i + 1);
1190    &b[s..e]
1191}
1192
1193/// A deployment's log: a file, and each line on the events feed.
1194pub(super) struct DeployLog {
1195    file: std::fs::File,
1196    apps: Apps,
1197    org: OrgId,
1198    app: String,
1199    id: u64,
1200}
1201
1202impl DeployLog {
1203    fn open(apps: &Apps, org: &OrgId, app: &str, id: u64) -> Result<DeployLog> {
1204        use std::os::unix::fs::OpenOptionsExt;
1205        let p = apps.log_path(org, app, id);
1206        if let Some(d) = p.parent() {
1207            std::fs::create_dir_all(d)?;
1208        }
1209        let file = std::fs::OpenOptions::new()
1210            .create(true)
1211            .append(true)
1212            .mode(0o600)
1213            .open(p)?;
1214        Ok(DeployLog {
1215            file,
1216            apps: apps.clone(),
1217            org: org.clone(),
1218            app: app.into(),
1219            id,
1220        })
1221    }
1222
1223    pub(super) fn line(&mut self, l: &str) {
1224        let _ = writeln!(self.file, "{l}");
1225        self.apps
1226            .event(&self.org, &self.app, "log", format!("#{}: {l}", self.id));
1227    }
1228}
1229
1230#[cfg(test)]
1231pub(super) mod tests {
1232    use super::*;
1233    use crate::app::EnvValue;
1234    use crate::image_check::Probe;
1235
1236    #[test]
1237    fn state_machine() {
1238        use Status::*;
1239        let mut d = Deployment {
1240            id: 1,
1241            app: "web".into(),
1242            trigger: Trigger::Api,
1243            by: "t".into(),
1244            status: Queued,
1245            requested: None,
1246            rollback_of: None,
1247            commit: None,
1248            image: None,
1249            digest: None,
1250            error: None,
1251            created_at: 0,
1252            started_at: None,
1253            finished_at: None,
1254            rendered: None,
1255        };
1256        assert!(d.advance(Done).is_err(), "queued cannot jump to done");
1257        d.advance(Building).unwrap();
1258        assert!(d.started_at.is_some());
1259        assert!(
1260            d.advance(Superseded).is_err(),
1261            "only a waiting one is superseded"
1262        );
1263        d.advance(Deploying).unwrap();
1264        d.advance(Done).unwrap();
1265        assert!(d.finished_at.is_some());
1266        for s in [Queued, Building, Deploying, Done, Failed, Superseded] {
1267            assert!(!Done.can_become(s) && !Failed.can_become(s) && !Superseded.can_become(s));
1268        }
1269        assert!(Queued.can_become(Superseded));
1270        assert!(Building.can_become(Failed));
1271        let v = d.summary();
1272        assert!(v.get("rendered").is_none());
1273        assert_eq!(v["status"], "done");
1274        assert_eq!(v["trigger"], "api");
1275    }
1276
1277    /// An `Apps` over a controller with no incusd behind it: everything up
1278    /// to the stack deploy works, which then fails.
1279    pub(in crate::app) fn apps(dir: &Path, gate: Arc<(Mutex<bool>, std::sync::Condvar)>) -> Apps {
1280        let k = crate::secrets::Keyring::new(age::x25519::Identity::generate(), vec![]);
1281        let secrets = Arc::new(Secrets::new(crate::secrets::LocalDriver::new(
1282            dir,
1283            Arc::new(k),
1284        )));
1285        let client = Client::with_socket("/nonexistent/isb-test/incus.sock");
1286        let store = crate::stack::Store::open(dir).unwrap();
1287        let ctl = Controller::start(
1288            client.clone(),
1289            store,
1290            Duration::from_secs(60),
1291            secrets.clone(),
1292        )
1293        .unwrap();
1294        // `docker:slow` holds its deploy in `building` until the gate opens.
1295        // `docker:traefik:whoami` is not on the registry; `docker:traefik/whoami` is.
1296        let probe: ImageProbe = Arc::new(move |image: &str, _| {
1297            if image == "docker:traefik:whoami" {
1298                return Probe::NotFound("manifest unknown".into());
1299            }
1300            if image == "docker:slow" {
1301                let (m, cv) = &*gate;
1302                let mut open = m.lock().unwrap();
1303                while !*open {
1304                    open = cv.wait(open).unwrap();
1305                }
1306            }
1307            Probe::Found(None)
1308        });
1309        Apps::new(dir, client, ctl, secrets).with_probe(probe)
1310    }
1311
1312    fn spec(v: Value) -> AppSpec {
1313        serde_json::from_value(v).unwrap()
1314    }
1315
1316    fn hdrs(h: &[(&str, &str)]) -> Box<webhook::Headers<'static>> {
1317        let h: Vec<(String, String)> = h
1318            .iter()
1319            .map(|(k, v)| (k.to_ascii_lowercase(), v.to_string()))
1320            .collect();
1321        Box::new(move |k: &str| {
1322            h.iter()
1323                .find(|(n, _)| *n == k.to_ascii_lowercase())
1324                .map(|(_, v)| v.clone())
1325        })
1326    }
1327
1328    #[test]
1329    fn rollout_events_reach_the_deployment_log() {
1330        let dir = tempfile::tempdir().unwrap();
1331        let gate = Arc::new((Mutex::new(true), std::sync::Condvar::new()));
1332        let ap = apps(dir.path(), gate);
1333        let ctl = ap.controller();
1334        ctl.note_service(
1335            "info",
1336            "acme/shop-production",
1337            "web",
1338            "before the deploy".into(),
1339        );
1340        let mark = ap.event_mark();
1341        ctl.note_service(
1342            "info",
1343            "acme/shop-production",
1344            "web",
1345            "app web: #1: own line".into(),
1346        );
1347        ctl.note_service(
1348            "info",
1349            "acme/shop-production",
1350            "web",
1351            "rolling out rev 1 to 1 slot(s)".into(),
1352        );
1353        ctl.note_service(
1354            "info",
1355            "acme/shop-production",
1356            "api",
1357            "another service".into(),
1358        );
1359        ctl.note_service("info", "acme/other", "web", "another stack".into());
1360        let mut seen = Vec::new();
1361        // No definition here, so the wait itself fails, after relaying.
1362        let r = ap.wait_service_with("acme/shop-production", "web", Some(mark), |m| {
1363            seen.push(m.to_string())
1364        });
1365        assert!(r.is_err());
1366        assert_eq!(seen, ["rolling out rev 1 to 1 slot(s)"]);
1367        // Without a mark, nothing is relayed.
1368        let mut none = Vec::new();
1369        let _ = ap.wait_service_with("acme/shop-production", "web", None, |m| {
1370            none.push(m.to_string())
1371        });
1372        assert!(none.is_empty());
1373    }
1374
1375    #[test]
1376    #[expect(
1377        clippy::too_many_lines,
1378        clippy::cognitive_complexity,
1379        reason = "predates the lint ratchet; split it when next changed"
1380    )]
1381    fn apps_records_webhooks_and_queue() {
1382        let dir = tempfile::tempdir().unwrap();
1383        let gate = Arc::new((Mutex::new(false), std::sync::Condvar::new()));
1384        let ap = apps(dir.path(), gate.clone());
1385        let org = OrgId::new("acme").unwrap();
1386
1387        // Projects and environments.
1388        let p = ap.project_create(&org, "shop", "", &[]).unwrap();
1389        assert_eq!(p.environments, ["production"]);
1390        assert!(ap.project_create(&org, "shop", "", &[]).is_err());
1391        ap.environment_create(&org, "shop", "staging").unwrap();
1392        assert!(
1393            dir.path()
1394                .join("orgs/acme/apps/projects/shop.json")
1395                .is_file()
1396        );
1397
1398        // An app naming a secret that does not exist is refused.
1399        let web = json!({
1400            "name": "web", "project": "shop",
1401            "source": {"image": "docker:traefik/whoami"},
1402            "env": "# greeting\nA=1\nT=${{secret.tok}}\n",
1403        });
1404        assert!(ap.create(&org, spec(web.clone())).is_err());
1405        ap.inner.secrets.set(&org, "tok", b"v").unwrap();
1406        let (app, secret) = ap.create(&org, spec(web.clone())).unwrap();
1407        assert_eq!(secret.len(), 64);
1408        assert_eq!(app.spec.environment, "production");
1409        assert!(ap.create(&org, spec(web)).is_err(), "names are unique");
1410        let mut bad = spec(
1411            json!({"name": "x", "project": "shop", "environment": "qa", "source": {"image": "x"}}),
1412        );
1413        assert!(ap.create(&org, bad.clone()).is_err(), "no such environment");
1414        bad.project = "nope".into();
1415        assert!(ap.create(&org, bad).is_err(), "no such project");
1416
1417        // The env editor keeps comments and never shows a secret's value.
1418        let text = ap.env_get(&org, "web").unwrap();
1419        assert_eq!(text, "# greeting\nA=1\nT=${{secret.tok}}\n");
1420        ap.env_set(
1421            &org,
1422            "web",
1423            "# greeting\nA=2\nT=${{secret.tok}}\nB=\"two words\"\n",
1424        )
1425        .unwrap();
1426        assert!(
1427            ap.env_get(&org, "web")
1428                .unwrap()
1429                .contains("A=2\nT=${{secret.tok}}\nB=\"two words\"")
1430        );
1431        assert!(ap.env_set(&org, "web", "T=${{secret.missing}}\n").is_err());
1432        let u = ap
1433            .update(&org, "web", &json!({"replicas": 3, "port": 80}))
1434            .unwrap();
1435        assert_eq!(u.spec.replicas, 3);
1436        assert_eq!(u.spec.env.get("A"), Some(&EnvValue::Plain("2".into())));
1437        assert!(
1438            ap.update(&org, "web", &json!({"project": "other"}))
1439                .is_err()
1440        );
1441        assert!(
1442            ap.update(&org, "web", &json!({"port": null}))
1443                .unwrap()
1444                .spec
1445                .port
1446                .is_none()
1447        );
1448        assert!(ap.rollback(&org, "web", None, Trigger::Api, "t").is_err());
1449
1450        // Webhooks: refused without the secret, whatever the app.
1451        let push = br#"{"ref":"refs/heads/main","after":"abc","head_commit":{"message":"m"}}"#;
1452        let none = hdrs(&[]);
1453        assert_eq!(ap.webhook("acme", "web", &none, None, push).0, 401);
1454        assert_eq!(ap.webhook("acme", "web", &none, Some("wrong"), push).0, 401);
1455        assert_eq!(
1456            ap.webhook("acme", "nope", &none, Some(&secret), push).0,
1457            401
1458        );
1459        assert_eq!(
1460            ap.webhook("Bad Org", "web", &none, Some(&secret), push).0,
1461            401
1462        );
1463        let forged = hdrs(&[
1464            ("X-GitHub-Event", "push"),
1465            (
1466                "X-Hub-Signature-256",
1467                &format!("sha256={}", webhook::sign(b"other", push)),
1468            ),
1469        ]);
1470        assert_eq!(ap.webhook("acme", "web", &forged, None, push).0, 401);
1471        assert!(ap.deployments(&org, "web").unwrap().is_empty());
1472        let signed = |event: &str, delivery: &str, body: &[u8]| {
1473            hdrs(&[
1474                ("X-GitHub-Event", event),
1475                ("X-GitHub-Delivery", delivery),
1476                (
1477                    "X-Hub-Signature-256",
1478                    &format!("sha256={}", webhook::sign(secret.as_bytes(), body)),
1479                ),
1480            ])
1481        };
1482        let (st, v) = ap.webhook("acme", "web", &signed("ping", "d0", b"{}"), None, b"{}");
1483        assert_eq!((st, v["ping"].as_bool()), (200, Some(true)));
1484        let (st, v) = ap.webhook("acme", "web", &signed("issues", "d1", push), None, push);
1485        assert_eq!(st, 200);
1486        assert!(v["ignored"].as_str().unwrap().contains("issues"), "{v}");
1487        let (st, v) = ap.webhook("acme", "web", &signed("push", "d2", push), None, push);
1488        assert_eq!(st, 202, "{v}");
1489        assert_eq!(v["deployment"], 1);
1490        // The same delivery again is not a second deploy.
1491        let (st, v) = ap.webhook("acme", "web", &signed("push", "d2", push), None, push);
1492        assert_eq!(st, 200, "{v}");
1493        let d = ap.wait(&org, "web", 1, Duration::from_secs(30)).unwrap();
1494        assert_eq!(d.trigger, Trigger::Webhook);
1495        assert_eq!(d.by, "webhook:github");
1496        assert_eq!(d.requested.as_deref(), Some("refs/heads/main abc: m"));
1497        // No incusd here: the stack deploy fails, and says so.
1498        assert_eq!(d.status, Status::Failed, "{d:?}");
1499        assert!(d.image.is_some());
1500        let (log, _, done) = ap.log(&org, "web", 1, 0).unwrap();
1501        assert!(done && log.contains("image docker:traefik/whoami"), "{log}");
1502        // The record comes with the text, so one call follows both.
1503        let (rest, next, rec) = ap.log_and_record(&org, "web", 1, 4).unwrap();
1504        assert_eq!((rec.id, rec.status), (1, Status::Failed));
1505        assert_eq!(next as usize, log.len());
1506        assert_eq!(rest, log[4..]);
1507        assert!(rec.summary().get("rendered").is_none());
1508        let ds = ap.deployments(&org, "web").unwrap();
1509        assert_eq!(ds.len(), 1);
1510
1511        // A git app deploys only on pushes to its branch.
1512        ap.inner.secrets.set(&org, "gh", b"t").unwrap();
1513        let (_, gsecret) = ap
1514            .create(
1515                &org,
1516                spec(json!({
1517                    "name": "api", "project": "shop", "environment": "staging",
1518                    "source": {"git": {"url": "https://example.invalid/o/r.git", "ref": "main", "auth": {"token_secret": "gh"}}},
1519                    "build": {"builder": {"type": "railpack"}},
1520                })),
1521            )
1522            .unwrap();
1523        let dev = br#"{"ref":"refs/heads/dev","after":"abc"}"#;
1524        let gl = |body: &[u8]| {
1525            let _ = body;
1526            hdrs(&[
1527                ("X-Gitlab-Event", "Push Hook"),
1528                ("X-Gitlab-Token", &gsecret),
1529            ])
1530        };
1531        let (st, v) = ap.webhook("acme", "api", &gl(dev), None, dev);
1532        assert_eq!(st, 200);
1533        assert!(
1534            v["ignored"].as_str().unwrap().contains("refs/heads/dev"),
1535            "{v}"
1536        );
1537        assert!(ap.deployments(&org, "api").unwrap().is_empty());
1538        let (st, _) = ap.webhook("acme", "api", &gl(push), None, push);
1539        assert_eq!(st, 202);
1540        let d = ap.wait(&org, "api", 1, Duration::from_secs(60)).unwrap();
1541        assert_eq!(d.status, Status::Failed);
1542        assert!(d.error.as_deref().unwrap_or("").contains("git"), "{d:?}");
1543        // A rotated secret: the old one stops working.
1544        let s2 = ap.webhook_secret(&org, "api", true).unwrap();
1545        assert_ne!(s2, gsecret);
1546        assert_eq!(ap.webhook("acme", "api", &gl(push), None, push).0, 401);
1547
1548        // The queue: one at a time; a newer request replaces a waiting one.
1549        ap.create(
1550            &org,
1551            spec(json!({"name": "slow", "project": "shop", "source": {"image": "docker:slow"}})),
1552        )
1553        .unwrap();
1554        let d1 = ap.deploy(&org, "slow", Trigger::Api, "t", None).unwrap();
1555        let started = Instant::now();
1556        while ap.deployment(&org, "slow", d1.id).unwrap().status != Status::Building {
1557            assert!(started.elapsed() < Duration::from_secs(10));
1558            std::thread::sleep(Duration::from_millis(20));
1559        }
1560        let d2 = ap.deploy(&org, "slow", Trigger::Api, "t", None).unwrap();
1561        let d3 = ap.deploy(&org, "slow", Trigger::Manual, "t", None).unwrap();
1562        assert_eq!(
1563            ap.deployment(&org, "slow", d2.id).unwrap().status,
1564            Status::Superseded
1565        );
1566        assert_eq!(
1567            ap.deployment(&org, "slow", d3.id).unwrap().status,
1568            Status::Queued
1569        );
1570        assert!(ap.delete(&org, "slow").is_err(), "not while deploying");
1571        {
1572            let (m, cv) = &*gate;
1573            *m.lock().unwrap() = true;
1574            cv.notify_all();
1575        }
1576        let d3 = ap
1577            .wait(&org, "slow", d3.id, Duration::from_secs(30))
1578            .unwrap();
1579        assert!(d3.status.finished());
1580        assert!(
1581            ap.wait(&org, "slow", d1.id, Duration::from_secs(1))
1582                .unwrap()
1583                .status
1584                .finished()
1585        );
1586
1587        // Deleting: projects with apps stay; an app takes its secrets along.
1588        assert!(ap.project_delete(&org, "shop").is_err());
1589        assert!(ap.environment_delete(&org, "shop", "staging").is_err());
1590        let started = Instant::now();
1591        for a in ["web", "api", "slow"] {
1592            while let Err(e) = ap.delete(&org, a) {
1593                assert!(started.elapsed() < Duration::from_secs(10), "{a}: {e}");
1594                std::thread::sleep(Duration::from_millis(50));
1595            }
1596        }
1597        assert!(ap.inner.secrets.inspect(&org, "app.web.webhook").is_err());
1598        assert!(
1599            ap.inner.secrets.inspect(&org, "tok").is_ok(),
1600            "not isb's to delete"
1601        );
1602        ap.environment_delete(&org, "shop", "staging").unwrap();
1603        ap.project_delete(&org, "shop").unwrap();
1604        assert!(ap.project_list(&org).unwrap().is_empty());
1605    }
1606
1607    #[test]
1608    #[expect(
1609        clippy::too_many_lines,
1610        clippy::cognitive_complexity,
1611        reason = "predates the lint ratchet; split it when next changed"
1612    )]
1613    fn previews_from_webhooks() {
1614        let dir = tempfile::tempdir().unwrap();
1615        let gate = Arc::new((Mutex::new(true), std::sync::Condvar::new()));
1616        let ap = apps(dir.path(), gate);
1617        let org = OrgId::new("acme").unwrap();
1618        ap.project_create(&org, "shop", "", &[]).unwrap();
1619        assert!(
1620            ap.environment_create(&org, "shop", "prod-pr-1").is_err(),
1621            "preview stack names are kept"
1622        );
1623        let (_, secret) = ap
1624            .create(
1625                &org,
1626                spec(json!({
1627                    "name": "web", "project": "shop",
1628                    "source": {"git": {"url": "https://example.invalid/acme/web.git", "ref": "main"}},
1629                    "build": {"builder": {"type": "dockerfile"}},
1630                    "port": 8080,
1631                })),
1632            )
1633            .unwrap();
1634        let pr = |action: &str, n: u64, head_repo: &str, base: &str| {
1635            format!(
1636                r#"{{"action":"{action}","number":{n},"pull_request":{{"title":"t{n}","merged":false,
1637                "head":{{"ref":"f{n}","sha":"{}","repo":{{"full_name":"{head_repo}"}}}},
1638                "base":{{"ref":"{base}","repo":{{"full_name":"acme/web"}}}}}}}}"#,
1639                "a".repeat(40)
1640            )
1641            .into_bytes()
1642        };
1643        let mut delivery = 0;
1644        let mut send = |body: &[u8]| {
1645            delivery += 1;
1646            let h = hdrs(&[
1647                ("X-Gitea-Event", "pull_request"),
1648                ("X-Gitea-Delivery", &format!("p{delivery}")),
1649                ("X-Gitea-Signature", &webhook::sign(secret.as_bytes(), body)),
1650            ]);
1651            ap.webhook("acme", "web", &h, None, body)
1652        };
1653        // Previews are opt-in.
1654        let (st, v) = send(&pr("opened", 1, "acme/web", "main"));
1655        assert_eq!(st, 200);
1656        assert!(v["ignored"].as_str().unwrap().contains("off"), "{v}");
1657        ap.update(
1658            &org,
1659            "web",
1660            &json!({"previews": {"enabled": true, "max": 2, "env": "MODE=preview\n"}}),
1661        )
1662        .unwrap();
1663        // Another base branch, a fork: ignored.
1664        let (st, v) = send(&pr("opened", 1, "acme/web", "dev"));
1665        assert_eq!(st, 200);
1666        assert!(
1667            v["ignored"].as_str().unwrap().contains("targets dev"),
1668            "{v}"
1669        );
1670        let (st, v) = send(&pr("opened", 1, "mallory/web", "main"));
1671        assert_eq!(st, 200);
1672        assert!(v["ignored"].as_str().unwrap().contains("fork"), "{v}");
1673        assert!(ap.preview_list(&org, "web").unwrap().is_empty());
1674        // Opened: a preview and a deploy (which fails here: no git host).
1675        let (st, v) = send(&pr("opened", 1, "acme/web", "main"));
1676        assert_eq!(st, 202, "{v}");
1677        assert_eq!(
1678            (v["preview"].as_u64(), v["deployment"].as_u64()),
1679            (Some(1), Some(1))
1680        );
1681        let p = ap.preview_get(&org, "web", 1).unwrap();
1682        assert_eq!(p.stack, "shop-production-pr-1");
1683        assert_eq!((p.provider.as_str(), p.fork), ("gitea", false));
1684        let d = ap
1685            .preview_wait(&org, "web", 1, 1, Duration::from_secs(60))
1686            .unwrap();
1687        assert_eq!(d.status, Status::Failed, "{d:?}");
1688        let (log, _, done) = ap.preview_log(&org, "web", 1, 1, 0).unwrap();
1689        assert!(done && log.contains("refs/pull/1/head"), "{log}");
1690        // Synchronize to the commit it has (Gitea sends one after opening):
1691        // nothing; to a new one: the same preview, a second deployment.
1692        let (st, v) = send(&pr("synchronize", 1, "acme/web", "main"));
1693        assert_eq!(st, 200, "{v}");
1694        assert!(v["ignored"].as_str().unwrap().contains("already"), "{v}");
1695        let moved = String::from_utf8(pr("synchronize", 1, "acme/web", "main"))
1696            .unwrap()
1697            .replace(&"a".repeat(40), &"b".repeat(40));
1698        let (st, v) = send(moved.as_bytes());
1699        assert_eq!((st, v["deployment"].as_u64()), (202, Some(2)), "{v}");
1700        ap.preview_wait(&org, "web", 1, 2, Duration::from_secs(60))
1701            .unwrap();
1702        // The limit: two at once.
1703        assert_eq!(send(&pr("opened", 2, "acme/web", "main")).0, 202);
1704        let (st, v) = send(&pr("opened", 3, "acme/web", "main"));
1705        assert_eq!(st, 200);
1706        assert!(v["ignored"].as_str().unwrap().contains("2 previews"), "{v}");
1707        assert!(ap.preview_get(&org, "web", 3).is_err());
1708        ap.preview_wait(&org, "web", 2, 1, Duration::from_secs(60))
1709            .unwrap();
1710        // Forks when allowed: marked, so they build in a VM without the
1711        // app's secrets.
1712        ap.update(&org, "web", &json!({"previews": {"forks": true, "max": 5}}))
1713            .unwrap();
1714        assert_eq!(send(&pr("opened", 4, "mallory/web", "main")).0, 202);
1715        assert!(ap.preview_get(&org, "web", 4).unwrap().fork);
1716        ap.preview_wait(&org, "web", 4, 1, Duration::from_secs(60))
1717            .unwrap();
1718        // Closed: removed with its records.
1719        let (st, v) = send(&pr("closed", 1, "acme/web", "main"));
1720        assert_eq!(st, 202, "{v}");
1721        let started = Instant::now();
1722        while ap.preview_get(&org, "web", 1).is_ok() {
1723            assert!(started.elapsed() < Duration::from_secs(30));
1724            std::thread::sleep(Duration::from_millis(50));
1725        }
1726        assert!(!dir.path().join("orgs/acme/apps/web/previews/1").exists());
1727        let (st, v) = send(&pr("closed", 1, "acme/web", "main"));
1728        assert_eq!(st, 200, "{v}");
1729        // The tools' path: redeploy, then delete.
1730        let d = ap
1731            .preview_redeploy(&org, "web", 2, Trigger::Api, "t")
1732            .unwrap();
1733        assert_eq!(d.id, 2);
1734        ap.preview_wait(&org, "web", 2, 2, Duration::from_secs(60))
1735            .unwrap();
1736        ap.preview_remove(&org, "web", 2, "test", true).unwrap();
1737        assert!(ap.preview_get(&org, "web", 2).is_err());
1738        // Deleting the app takes its previews along first; here the fork's
1739        // build cache cannot be checked without incusd, so it stops there.
1740        assert!(ap.delete(&org, "web").is_err());
1741        assert!(ap.get(&org, "web").is_ok());
1742        assert!(ap.preview_get(&org, "web", 4).unwrap().removing);
1743    }
1744}