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