1use 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, Project, 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
36pub type BuildFn =
38 Arc<dyn Fn(&Client, &BuildRequest, &mut dyn FnMut(&str)) -> Result<BuiltImage> + Send + Sync>;
39
40pub type DigestFn = Arc<dyn Fn(&str) -> Option<String> + Send + Sync>;
42
43const KEEP_DEPLOYMENTS: usize = 30;
45
46const DEPLOY_TIMEOUT: Duration = Duration::from_secs(900);
48
49#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
51#[serde(rename_all = "lowercase")]
52pub enum Trigger {
53 Manual,
55 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 Superseded,
70 Cancelled,
72}
73
74impl Status {
75 pub fn finished(self) -> bool {
76 !matches!(self, Status::Queued | Status::Building | Status::Deploying)
77 }
78
79 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#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
104pub struct Deployment {
105 pub id: u64,
106 pub app: String,
107 pub trigger: Trigger,
108 pub by: String,
110 pub status: Status,
111 #[serde(default, skip_serializing_if = "Option::is_none")]
113 pub requested: Option<String>,
114 #[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 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 #[serde(default, skip_serializing_if = "Option::is_none")]
133 pub rendered: Option<Rendered>,
134}
135
136impl Deployment {
137 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 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 pub(super) edit: Mutex<()>,
181 pub(super) stacks: Mutex<()>,
184 pub(super) queues: Mutex<BTreeMap<(OrgId, String), Queue>>,
186 deliveries: Mutex<VecDeque<String>>,
188 pub(super) secret_hook: OnceLock<SecretHook>,
192}
193
194pub type SecretHook = Arc<dyn Fn(&OrgId, &str) + Send + Sync>;
196
197#[derive(Clone)]
199pub struct Apps {
200 pub(super) inner: Arc<Inner>,
201}
202
203impl Apps {
204 pub fn new(state: &Path, client: Client, ctl: Controller, secrets: Arc<Secrets>) -> Apps {
207 let a = Apps {
208 inner: Arc::new(Inner {
209 state: state.to_path_buf(),
210 client,
211 ctl,
212 secrets,
213 build: Arc::new(|c: &Client, r: &BuildRequest, l: &mut dyn FnMut(&str)| {
214 crate::build::run(c, r, l)
215 }),
216 digest: Arc::new(skopeo_digest),
217 timeout: DEPLOY_TIMEOUT,
218 edit: Mutex::new(()),
219 stacks: Mutex::new(()),
220 queues: Mutex::new(BTreeMap::new()),
221 deliveries: Mutex::new(VecDeque::new()),
222 secret_hook: OnceLock::new(),
223 }),
224 };
225 a.recover();
226 a
227 }
228
229 fn with(self, f: impl FnOnce(&mut Inner)) -> Apps {
230 let mut inner = Arc::try_unwrap(self.inner)
231 .unwrap_or_else(|_| panic!("Apps::with_* before the Apps is shared"));
232 f(&mut inner);
233 Apps {
234 inner: Arc::new(inner),
235 }
236 }
237
238 pub fn with_build(self, f: BuildFn) -> Apps {
240 self.with(|i| i.build = f)
241 }
242
243 pub fn with_digest(self, f: DigestFn) -> Apps {
245 self.with(|i| i.digest = f)
246 }
247
248 pub fn with_timeout(self, t: Duration) -> Apps {
249 self.with(|i| i.timeout = t)
250 }
251
252 pub(super) fn apps_dir(&self, org: &OrgId) -> PathBuf {
255 super::org_root(&self.inner.state, org).join("apps")
256 }
257
258 fn projects_dir(&self, org: &OrgId) -> PathBuf {
259 self.apps_dir(org).join("projects")
260 }
261
262 pub(super) fn app_dir(&self, org: &OrgId, app: &str) -> PathBuf {
263 self.apps_dir(org).join(app)
264 }
265
266 fn deployments_dir(&self, org: &OrgId, app: &str) -> PathBuf {
267 self.app_dir(org, app).join("deployments")
268 }
269
270 pub fn source_dir(&self, org: &OrgId, app: &str) -> PathBuf {
272 super::org_root(&self.inner.state, org)
273 .join("sources")
274 .join(app)
275 }
276
277 pub fn project_create(
280 &self,
281 org: &OrgId,
282 name: &str,
283 description: &str,
284 environments: &[String],
285 ) -> Result<Project> {
286 super::validate_part("project", name)?;
287 let envs: Vec<String> = if environments.is_empty() {
288 vec![super::DEFAULT_ENVIRONMENT.into()]
289 } else {
290 environments.to_vec()
291 };
292 for e in &envs {
293 super::validate_part("environment", e)?;
294 super::stack_name(name, e)?;
295 }
296 let _g = self.inner.edit.lock().unwrap();
297 let p = self.projects_dir(org).join(format!("{name}.json"));
298 if p.exists() {
299 return Err(Error::AlreadyExists(format!("project {name}")));
300 }
301 let mut envs = envs;
302 envs.dedup();
303 let proj = Project {
304 name: name.into(),
305 description: description.into(),
306 environments: envs,
307 created_at: crate::stack::now_secs(),
308 };
309 super::write_atomic(&p, &serde_json::to_vec_pretty(&proj)?)?;
310 Ok(proj)
311 }
312
313 pub fn project_get(&self, org: &OrgId, name: &str) -> Result<Project> {
314 super::validate_part("project", name)?;
315 let p = self.projects_dir(org).join(format!("{name}.json"));
316 match std::fs::read(&p) {
317 Ok(b) => Ok(serde_json::from_slice(&b)?),
318 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
319 Err(Error::NotFound(format!("project {name} in org {org}")))
320 }
321 Err(e) => Err(e.into()),
322 }
323 }
324
325 pub fn project_list(&self, org: &OrgId) -> Result<Vec<Project>> {
326 let mut out = Vec::new();
327 let Ok(rd) = std::fs::read_dir(self.projects_dir(org)) else {
328 return Ok(out);
329 };
330 for e in rd.flatten() {
331 let p = e.path();
332 if p.extension().is_some_and(|x| x == "json") {
333 match serde_json::from_slice::<Project>(&std::fs::read(&p)?) {
334 Ok(pr) => out.push(pr),
335 Err(e) => eprintln!("isb serve: skipping {}: {e}", p.display()),
336 }
337 }
338 }
339 out.sort_by(|a, b| a.name.cmp(&b.name));
340 Ok(out)
341 }
342
343 pub fn project_delete(&self, org: &OrgId, name: &str) -> Result<()> {
345 let _g = self.inner.edit.lock().unwrap();
346 self.project_get(org, name)?;
347 let apps: Vec<String> = self
348 .list(org)?
349 .into_iter()
350 .filter(|a| a.spec.project == name)
351 .map(|a| a.spec.name)
352 .collect();
353 if !apps.is_empty() {
354 return Err(Error::invalid(format!(
355 "project {name} still has apps: {}; delete them first",
356 apps.join(", ")
357 )));
358 }
359 std::fs::remove_file(self.projects_dir(org).join(format!("{name}.json")))?;
360 Ok(())
361 }
362
363 pub fn environment_create(&self, org: &OrgId, project: &str, env: &str) -> Result<Project> {
364 super::validate_part("environment", env)?;
365 super::stack_name(project, env)?;
366 let _g = self.inner.edit.lock().unwrap();
367 let mut p = self.project_get(org, project)?;
368 if p.environments.iter().any(|e| e == env) {
369 return Err(Error::AlreadyExists(format!(
370 "environment {env} in project {project}"
371 )));
372 }
373 p.environments.push(env.into());
374 self.save_project(org, &p)?;
375 Ok(p)
376 }
377
378 pub fn environment_delete(&self, org: &OrgId, project: &str, env: &str) -> Result<Project> {
379 let _g = self.inner.edit.lock().unwrap();
380 let mut p = self.project_get(org, project)?;
381 if !p.environments.iter().any(|e| e == env) {
382 return Err(Error::NotFound(format!(
383 "environment {env} in project {project}"
384 )));
385 }
386 let apps: Vec<String> = self
387 .list(org)?
388 .into_iter()
389 .filter(|a| a.spec.project == project && a.spec.environment == env)
390 .map(|a| a.spec.name)
391 .collect();
392 if !apps.is_empty() {
393 return Err(Error::invalid(format!(
394 "environment {env} still has apps: {}; delete them first",
395 apps.join(", ")
396 )));
397 }
398 p.environments.retain(|e| e != env);
399 self.save_project(org, &p)?;
400 Ok(p)
401 }
402
403 fn save_project(&self, org: &OrgId, p: &Project) -> Result<()> {
404 super::write_atomic(
405 &self.projects_dir(org).join(format!("{}.json", p.name)),
406 &serde_json::to_vec_pretty(p)?,
407 )
408 }
409
410 pub fn get(&self, org: &OrgId, name: &str) -> Result<App> {
413 super::validate_app_name(name)?;
414 match std::fs::read(self.app_dir(org, name).join("app.json")) {
415 Ok(b) => Ok(serde_json::from_slice(&b)?),
416 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
417 Err(Error::NotFound(format!("app {name} in org {org}")))
418 }
419 Err(e) => Err(e.into()),
420 }
421 }
422
423 pub fn list(&self, org: &OrgId) -> Result<Vec<App>> {
424 let mut out = Vec::new();
425 let Ok(rd) = std::fs::read_dir(self.apps_dir(org)) else {
426 return Ok(out);
427 };
428 for e in rd.flatten() {
429 let p = e.path().join("app.json");
430 if !p.is_file() {
431 continue;
432 }
433 match serde_json::from_slice::<App>(&std::fs::read(&p)?) {
434 Ok(a) => out.push(a),
435 Err(e) => eprintln!("isb serve: skipping {}: {e}", p.display()),
436 }
437 }
438 out.sort_by(|a, b| a.spec.name.cmp(&b.spec.name));
439 Ok(out)
440 }
441
442 fn save(&self, org: &OrgId, app: &App) -> Result<()> {
443 super::write_atomic(
444 &self.app_dir(org, &app.spec.name).join("app.json"),
445 &serde_json::to_vec_pretty(app)?,
446 )
447 }
448
449 pub fn create(&self, org: &OrgId, mut spec: AppSpec) -> Result<(App, String)> {
452 if let Source::Database(db) = &mut spec.source {
453 db.normalize(&spec.name);
454 }
455 spec.validate()?;
456 let _g = self.inner.edit.lock().unwrap();
457 let proj = self.project_get(org, &spec.project)?;
458 if !proj.environments.contains(&spec.environment) {
459 return Err(Error::NotFound(format!(
460 "environment {} in project {} (it has {})",
461 spec.environment,
462 spec.project,
463 proj.environments.join(", ")
464 )));
465 }
466 if self.app_dir(org, &spec.name).join("app.json").exists() {
467 return Err(Error::AlreadyExists(format!("app {}", spec.name)));
468 }
469 self.check_secrets(org, &spec)?;
470 let now = crate::stack::now_secs();
471 let app = App {
472 spec,
473 created_at: now,
474 updated_at: now,
475 next_deployment: 1,
476 current: None,
477 };
478 if let Source::Database(db) = &app.spec.source {
479 self.database_credentials(org, &app.spec, db)?;
480 }
481 let secret = git::random_hex(32);
482 self.inner
483 .secrets
484 .set(org, &app.spec.webhook_secret(), secret.as_bytes())?;
485 self.save(org, &app)?;
486 Ok((app, secret))
487 }
488
489 pub fn secrets(&self) -> &Arc<Secrets> {
492 &self.inner.secrets
493 }
494
495 pub fn client(&self) -> &Client {
497 &self.inner.client
498 }
499
500 pub fn controller(&self) -> &Controller {
502 &self.inner.ctl
503 }
504
505 pub fn wait_converged(&self, org: &OrgId, name: &str) -> Result<(bool, String)> {
508 let app = self.get(org, name)?;
509 let q = crate::stack::qualified(org, &app.spec.stack()?);
510 self.wait_service(&q, name)
511 }
512
513 fn check_secrets(&self, org: &OrgId, spec: &AppSpec) -> Result<()> {
516 let mut names = spec.env.secret_names();
517 if let Some(p) = &spec.previews {
518 names.extend(p.secret_names());
519 }
520 names.extend(spec.files.iter().map(|f| f.secret.clone()));
521 if let Source::Git(g) = &spec.source {
522 names.extend(g.auth.secret().map(String::from));
523 }
524 for n in names {
525 self.inner.secrets.inspect(org, &n).map_err(|e| match e {
526 Error::NotFound(_) => Error::invalid(format!(
527 "secret {n} does not exist in org {org}; create it with `isb secret create {n} --org {org}`"
528 )),
529 e => e,
530 })?;
531 }
532 Ok(())
533 }
534
535 pub fn check_spec(&self, org: &OrgId, spec: &AppSpec) -> Result<()> {
539 self.check_secrets(org, spec)
540 }
541
542 pub fn update(&self, org: &OrgId, name: &str, patch: &Value) -> Result<App> {
546 let _g = self.inner.edit.lock().unwrap();
547 let mut app = self.get(org, name)?;
548 let mut v = serde_json::to_value(&app.spec)?;
549 super::merge_patch(&mut v, patch);
550 let mut spec: AppSpec =
551 serde_json::from_value(v).map_err(|e| Error::invalid(format!("app {name}: {e}")))?;
552 super::manifest::check_update(&app.spec, &mut spec)?;
553 spec.validate()?;
554 self.check_secrets(org, &spec)?;
555 app.spec = spec;
556 app.updated_at = crate::stack::now_secs();
557 self.save(org, &app)?;
558 Ok(app)
559 }
560
561 pub fn env_get(&self, org: &OrgId, name: &str) -> Result<String> {
562 Ok(self.get(org, name)?.spec.env.render())
563 }
564
565 pub fn env_set(&self, org: &OrgId, name: &str, text: &str) -> Result<App> {
567 let env = super::EnvFile::parse(text)?;
568 self.update(org, name, &json!({"env": env.render()}))
569 }
570
571 pub fn webhook_secret(&self, org: &OrgId, name: &str, rotate: bool) -> Result<String> {
573 let app = self.get(org, name)?;
574 let key = app.spec.webhook_secret();
575 if !rotate {
576 if let Ok((v, _)) = self.inner.secrets.get(org, &key) {
577 return String::from_utf8(v)
578 .map_err(|_| Error::invalid("webhook secret is not text"));
579 }
580 }
581 let s = git::random_hex(32);
582 self.inner.secrets.set(org, &key, s.as_bytes())?;
583 Ok(s)
584 }
585
586 pub fn deploy_key(&self, org: &OrgId, name: &str) -> Result<String> {
590 let app = self.get(org, name)?;
591 let Source::Git(g) = &app.spec.source else {
592 return Err(Error::invalid(format!("app {name} has no git source")));
593 };
594 if git::transport(&g.url)? != git::Transport::Ssh {
595 return Err(Error::invalid(format!(
596 "a deploy key needs an SSH URL; {} is not one (git@host:owner/repo)",
597 g.url
598 )));
599 }
600 let (private, public) =
601 git::generate_deploy_key(&self.apps_dir(org), &format!("isb-deploy-{org}-{name}"))?;
602 let key = super::deploy_key_secret(name);
603 self.inner.secrets.set(org, &key, &private)?;
604 self.update(
605 org,
606 name,
607 &json!({"source": {"git": {"auth": {"ssh_key_secret": key}}}}),
608 )?;
609 Ok(public)
610 }
611
612 pub fn delete(&self, org: &OrgId, name: &str) -> Result<()> {
616 let app = self.get(org, name)?;
617 if self
618 .inner
619 .queues
620 .lock()
621 .unwrap()
622 .contains_key(&(org.clone(), name.to_string()))
623 {
624 return Err(Error::invalid(format!(
625 "app {name} is deploying; delete it once that finishes"
626 )));
627 }
628 self.previews_remove_all(org, name)?;
630 let stack = app.spec.stack()?;
631 let q = crate::stack::qualified(org, &stack);
632 {
633 let _s = self.inner.stacks.lock().unwrap();
634 if let Ok(cur) = self.inner.ctl.definition(&q) {
635 if cur.file.services.contains_key(name) {
636 let file = super::splice(Some(&cur.file), &stack, name, None);
637 if file.services.is_empty() {
638 self.inner.ctl.remove(&q, false, Duration::from_secs(300))?;
639 } else {
640 let def = self.stack_def(org, &stack, file, "app delete")?;
641 self.inner.ctl.deploy(def)?;
642 }
643 }
644 }
645 }
646 let _g = self.inner.edit.lock().unwrap();
647 for s in [super::webhook_secret(name), super::deploy_key_secret(name)] {
648 match self.inner.secrets.delete(org, &s) {
649 Ok(()) => {}
650 Err(e) if e.is_not_found() => {}
651 Err(e) => return Err(e),
652 }
653 }
654 let _ = std::fs::remove_dir_all(self.source_dir(org, name));
655 std::fs::remove_dir_all(self.app_dir(org, name))?;
656 Ok(())
657 }
658
659 fn dep_path(&self, org: &OrgId, app: &str, id: u64) -> PathBuf {
662 self.deployments_dir(org, app).join(format!("{id}.json"))
663 }
664
665 fn log_path(&self, org: &OrgId, app: &str, id: u64) -> PathBuf {
666 self.deployments_dir(org, app).join(format!("{id}.log"))
667 }
668
669 pub(super) fn save_dep(&self, org: &OrgId, d: &Deployment) -> Result<()> {
670 super::write_atomic(
671 &self.dep_path(org, &d.app, d.id),
672 &serde_json::to_vec_pretty(d)?,
673 )
674 }
675
676 pub fn deployment(&self, org: &OrgId, app: &str, id: u64) -> Result<Deployment> {
677 super::validate_app_name(app)?;
678 match std::fs::read(self.dep_path(org, app, id)) {
679 Ok(b) => Ok(serde_json::from_slice(&b)?),
680 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
681 Err(Error::NotFound(format!("deployment {id} of app {app}")))
682 }
683 Err(e) => Err(e.into()),
684 }
685 }
686
687 pub fn deployments(&self, org: &OrgId, app: &str) -> Result<Vec<Deployment>> {
689 self.get(org, app)?;
690 let mut ids: Vec<u64> = match std::fs::read_dir(self.deployments_dir(org, app)) {
691 Ok(rd) => rd
692 .flatten()
693 .filter_map(|e| {
694 let n = e.file_name().into_string().ok()?;
695 n.strip_suffix(".json")?.parse().ok()
696 })
697 .collect(),
698 Err(_) => vec![],
699 };
700 ids.sort_unstable_by(|a, b| b.cmp(a));
701 ids.into_iter()
702 .map(|id| self.deployment(org, app, id))
703 .collect()
704 }
705
706 pub fn log(&self, org: &OrgId, app: &str, id: u64, offset: u64) -> Result<(String, u64, bool)> {
708 let (text, next, d) = self.log_and_record(org, app, id, offset)?;
709 Ok((text, next, d.status.finished()))
710 }
711
712 pub fn log_and_record(
716 &self,
717 org: &OrgId,
718 app: &str,
719 id: u64,
720 offset: u64,
721 ) -> Result<(String, u64, Deployment)> {
722 let d = self.deployment(org, app, id)?;
723 let b = std::fs::read(self.log_path(org, app, id)).unwrap_or_default();
724 let start = (offset as usize).min(b.len());
725 Ok((
726 String::from_utf8_lossy(&b[start..]).into_owned(),
727 b.len() as u64,
728 d,
729 ))
730 }
731
732 pub fn wait(&self, org: &OrgId, app: &str, id: u64, timeout: Duration) -> Result<Deployment> {
734 let started = Instant::now();
735 loop {
736 let d = self.deployment(org, app, id)?;
737 if d.status.finished() || started.elapsed() >= timeout {
738 return Ok(d);
739 }
740 std::thread::sleep(Duration::from_millis(500));
741 }
742 }
743
744 pub fn deploy(
746 &self,
747 org: &OrgId,
748 name: &str,
749 trigger: Trigger,
750 by: &str,
751 requested: Option<String>,
752 ) -> Result<Deployment> {
753 self.enqueue(org, name, trigger, by, requested, None)
754 }
755
756 pub fn rollback(
759 &self,
760 org: &OrgId,
761 name: &str,
762 to: Option<u64>,
763 trigger: Trigger,
764 by: &str,
765 ) -> Result<Deployment> {
766 let app = self.get(org, name)?;
767 let target = match to {
768 Some(id) => self.deployment(org, name, id)?,
769 None => self
770 .deployments(org, name)?
771 .into_iter()
772 .find(|d| d.status == Status::Done && Some(d.id) != app.current)
773 .ok_or_else(|| {
774 Error::invalid(format!(
775 "app {name} has no earlier successful deployment to roll back to"
776 ))
777 })?,
778 };
779 if target.status != Status::Done || target.rendered.is_none() {
780 return Err(Error::invalid(format!(
781 "deployment {} did not finish done; only those can be rolled back to",
782 target.id
783 )));
784 }
785 self.enqueue(org, name, trigger, by, None, Some(target.id))
786 }
787
788 fn enqueue(
789 &self,
790 org: &OrgId,
791 name: &str,
792 trigger: Trigger,
793 by: &str,
794 requested: Option<String>,
795 rollback_of: Option<u64>,
796 ) -> Result<Deployment> {
797 let dep = {
798 let _g = self.inner.edit.lock().unwrap();
799 let mut app = self.get(org, name)?;
800 let id = app.next_deployment;
801 app.next_deployment += 1;
802 self.save(org, &app)?;
803 let d = Deployment {
804 id,
805 app: name.into(),
806 trigger,
807 by: by.into(),
808 status: Status::Queued,
809 requested,
810 rollback_of,
811 commit: None,
812 image: None,
813 digest: None,
814 error: None,
815 created_at: crate::stack::controller::now_ms(),
816 started_at: None,
817 finished_at: None,
818 rendered: None,
819 };
820 self.save_dep(org, &d)?;
821 self.prune(org, &app);
822 d
823 };
824 self.event(
825 org,
826 name,
827 "info",
828 format!("deployment {} queued by {by}", dep.id),
829 );
830 let start = {
831 let mut qs = self.inner.queues.lock().unwrap();
832 let q = qs.entry((org.clone(), name.to_string())).or_default();
833 if let Some(old) = q.next.replace(dep.id) {
834 if let Ok(mut d) = self.deployment(org, name, old) {
835 if d.advance(Status::Superseded).is_ok() {
836 d.error = Some(format!("superseded by deployment {}", dep.id));
837 let _ = self.save_dep(org, &d);
838 }
839 }
840 }
841 !std::mem::replace(&mut q.running, true)
842 };
843 if start {
844 let me = self.clone();
845 let (org, name) = (org.clone(), name.to_string());
846 std::thread::spawn(move || me.drain(&org, &name));
847 }
848 Ok(dep)
849 }
850
851 fn drain(&self, org: &OrgId, name: &str) {
853 loop {
854 let next = {
855 let mut qs = self.inner.queues.lock().unwrap();
856 let key = (org.clone(), name.to_string());
857 let q = qs.entry(key.clone()).or_default();
858 match q.next.take() {
859 Some(id) => id,
860 None => {
861 qs.remove(&key);
862 return;
863 }
864 }
865 };
866 self.run(org, name, next);
867 }
868 }
869
870 fn run(&self, org: &OrgId, name: &str, id: u64) {
871 let Ok(mut dep) = self.deployment(org, name, id) else {
872 return;
873 };
874 let mut log = match DeployLog::open(self, org, name, id) {
875 Ok(l) => l,
876 Err(e) => {
877 eprintln!("isb serve: app {name}: deployment {id}: no log: {e}");
878 return;
879 }
880 };
881 let r = self.pipeline(org, name, &mut dep, &mut log);
882 let (kind, level, msg) = match &r {
883 Ok(()) => ("deploy.succeeded", "info", format!("deployment {id} done")),
884 Err(e) => (
885 "deploy.failed",
886 "error",
887 format!("deployment {id} failed: {e}"),
888 ),
889 };
890 log.line(&msg);
891 if let Err(e) = r {
892 dep.error = Some(e.to_string());
893 let _ = dep.advance(Status::Failed);
894 }
895 if let Err(e) = self.save_dep(org, &dep) {
896 eprintln!("isb serve: app {name}: deployment {id}: {e}");
897 }
898 self.kind_event(org, name, kind, level, msg);
899 }
900
901 fn set_status(&self, org: &OrgId, dep: &mut Deployment, s: Status) -> Result<()> {
902 dep.advance(s)?;
903 self.save_dep(org, dep)?;
904 self.event(
905 org,
906 &dep.app,
907 "info",
908 format!("deployment {}: {s:?}", dep.id).to_lowercase(),
909 );
910 Ok(())
911 }
912
913 #[expect(
914 clippy::too_many_lines,
915 reason = "predates the lint ratchet; split it when next changed"
916 )]
917 fn pipeline(
918 &self,
919 org: &OrgId,
920 name: &str,
921 dep: &mut Deployment,
922 log: &mut DeployLog,
923 ) -> Result<()> {
924 self.set_status(org, dep, Status::Building)?;
925 let app = self.get(org, name)?;
926 let mut notes = Vec::new();
927 let rendered = match dep.rollback_of {
928 Some(of) => {
929 let t = self.deployment(org, name, of)?;
930 log.line(&format!(
931 "rolling back to deployment {of} ({})",
932 t.image.as_deref().unwrap_or("?")
933 ));
934 dep.image = t.image;
935 dep.digest = t.digest;
936 dep.commit = t.commit;
937 t.rendered
938 .ok_or_else(|| Error::invalid(format!("deployment {of} kept no settings")))?
939 }
940 None => {
941 let (image, digest) = match &app.spec.source {
942 Source::Image(i) => {
943 log.line(&format!("image {i}"));
944 self.resolve(i, log)
945 }
946 Source::Database(db) => {
947 let i = db.image();
948 log.line(&format!(
949 "database {} {}: image {i}",
950 db.engine,
951 db.version()
952 ));
953 for n in self.write_urls(org, &app.spec, db)? {
957 log.line(&format!("secret {n}: the connection URL"));
958 self.url_secret_moved(org, &n, log);
959 }
960 self.resolve(&i, log)
961 }
962 Source::Git(g) => {
963 let creds = self.credentials(org, &g.auth)?;
964 let co = git::fetch(g, &creds, &self.source_dir(org, name), &mut |l| {
965 log.line(l)
966 })?;
967 dep.commit = Some(Commit {
968 sha: co.sha.clone(),
969 message: co.message.clone(),
970 });
971 self.save_dep(org, dep)?;
972 let b =
973 app.spec.build.as_ref().ok_or_else(|| {
974 Error::invalid("a git source needs build settings")
975 })?;
976 let req = BuildRequest {
977 org: org.clone(),
978 app: name.into(),
979 context: co.dir.clone(),
980 subdir: g.subdir.clone(),
981 builder: b.builder.clone(),
982 args: b.args.iter().map(|(k, v)| (k.clone(), v.clone())).collect(),
983 tag: co.sha.clone(),
984 untrusted: b.untrusted,
985 cache: None,
986 };
987 log.line(&format!("building {} with {:?}", co.sha, b.builder));
988 let built =
989 (self.inner.build)(&self.inner.client, &req, &mut |l| log.line(l))?;
990 log.line(&format!("built {} ({})", built.image, built.digest));
991 (built.image, Some(built.digest))
992 }
993 };
994 dep.image = Some(image.clone());
995 dep.digest = digest;
996 super::render(&app.spec, &image, &mut notes)?
997 }
998 };
999 for n in ¬es {
1000 log.line(&format!("note: {n}"));
1001 }
1002 dep.rendered = Some(rendered.clone());
1003 self.set_status(org, dep, Status::Deploying)?;
1004 let stack = app.spec.stack()?;
1005 let q = crate::stack::qualified(org, &stack);
1006 let mark = self.event_mark();
1007 {
1008 let _s = self.inner.stacks.lock().unwrap();
1009 let cur = self.inner.ctl.definition(&q).ok();
1010 let file = super::splice(cur.as_ref().map(|d| &d.file), &stack, name, Some(&rendered));
1011 let def = self.stack_def(org, &stack, file, &dep.by)?;
1012 let changes = self.inner.ctl.deploy(def)?;
1013 let mine = changes.iter().find(|c| c.service == name);
1014 match mine {
1015 Some(c) if c.change == "unchanged" => {
1016 log.line("settings unchanged: replacing the instances anyway");
1019 self.inner.ctl.redeploy(&q, name)?;
1020 }
1021 Some(c) => log.line(&format!(
1022 "stack {stack}: {} {} (rev {}, {} replicas)",
1023 c.service, c.change, c.rev, c.replicas
1024 )),
1025 None => {}
1026 }
1027 }
1028 let (ok, msg) = self.wait_service_with(&q, name, Some(mark), |m| log.line(m))?;
1029 if !ok {
1030 return Err(Error::invalid(format!("service {name}: {msg}")));
1031 }
1032 log.line(&format!("service {name}.{stack} converged"));
1033 self.set_status(org, dep, Status::Done)?;
1034 let _g = self.inner.edit.lock().unwrap();
1035 let mut app = self.get(org, name)?;
1036 app.current = Some(dep.id);
1037 self.save(org, &app)?;
1038 Ok(())
1039 }
1040
1041 fn resolve(&self, i: &str, log: &mut DeployLog) -> (String, Option<String>) {
1043 match (self.inner.digest)(i) {
1044 Some(d) => {
1045 let pinned = super::pin(i, &d).unwrap_or_else(|| i.to_string());
1046 log.line(&format!("resolved to {pinned}"));
1047 (pinned, Some(d))
1048 }
1049 None => {
1050 log.line("no digest found; deploying by tag");
1051 (i.to_string(), None)
1052 }
1053 }
1054 }
1055
1056 pub(super) fn stack_def(
1058 &self,
1059 org: &OrgId,
1060 stack: &str,
1061 file: crate::spec::ComposeFile,
1062 by: &str,
1063 ) -> Result<StackDef> {
1064 let base = self.apps_dir(org);
1065 std::fs::create_dir_all(&base)?;
1066 let secrets = crate::stack::secrets::bind(
1067 &self.inner.secrets,
1068 org,
1069 stack,
1070 &file,
1071 &BTreeMap::new(),
1072 false,
1073 )?;
1074 Ok(StackDef {
1075 name: stack.into(),
1076 org: org.clone(),
1077 file,
1078 base_dir: base,
1079 secrets,
1080 force: BTreeMap::new(),
1081 images: BTreeMap::new(),
1082 deployed_at: crate::stack::now_secs(),
1083 deployed_by: by.into(),
1084 previous: None,
1085 })
1086 }
1087
1088 pub(super) fn wait_service(&self, q: &str, svc: &str) -> Result<(bool, String)> {
1091 self.wait_service_with(q, svc, None, |_| {})
1092 }
1093
1094 pub(super) fn event_mark(&self) -> u64 {
1097 self.inner.ctl.events(u64::MAX, 0).0
1098 }
1099
1100 pub(super) fn wait_service_with(
1106 &self,
1107 q: &str,
1108 svc: &str,
1109 mut since: Option<u64>,
1110 mut seen: impl FnMut(&str),
1111 ) -> Result<(bool, String)> {
1112 let started = Instant::now();
1113 let own = format!("app {svc}");
1114 let mut relay = |since: &mut Option<u64>| {
1115 let Some(after) = *since else { return };
1116 let (last, evs) = self.inner.ctl.events(after, 500);
1117 for e in evs.iter().filter(|e| e.stack == q && e.service == svc) {
1118 if !e.message.starts_with(&own) {
1119 seen(&e.message);
1120 }
1121 }
1122 *since = Some(last.max(after));
1123 };
1124 loop {
1125 relay(&mut since);
1126 let def = self.inner.ctl.definition(q)?;
1127 let rev = def.revision(svc)?;
1128 let replicas = def.service(svc)?.replicas();
1129 let st = self.inner.ctl.status(q)?;
1130 if let Some(s) = st.services.iter().find(|s| s.service == svc) {
1131 if s.rev == rev && s.replicas == replicas {
1132 match s.state.as_str() {
1133 "converged" => {
1134 relay(&mut since);
1135 return Ok((true, String::new()));
1136 }
1137 "paused" | "failing" => {
1138 relay(&mut since);
1139 return Ok((
1140 false,
1141 format!("{}: {}", s.state, s.message.clone().unwrap_or_default()),
1142 ));
1143 }
1144 _ => {}
1145 }
1146 }
1147 }
1148 if started.elapsed() >= self.inner.timeout {
1149 return Ok((
1150 false,
1151 format!("not converged after {:?}", self.inner.timeout),
1152 ));
1153 }
1154 std::thread::sleep(Duration::from_millis(500));
1155 }
1156 }
1157
1158 pub(super) fn credentials(&self, org: &OrgId, auth: &GitAuth) -> Result<Credentials> {
1159 let read = |n: &str| {
1160 self.inner
1161 .secrets
1162 .get(org, n)
1163 .map(|(v, _)| v)
1164 .map_err(|e| Error::invalid(format!("git credential {n}: {e}")))
1165 };
1166 Ok(match auth {
1167 GitAuth::None => Credentials::None,
1168 GitAuth::Token {
1169 token_secret,
1170 username,
1171 } => Credentials::Token {
1172 username: username.clone().unwrap_or_else(|| "x-access-token".into()),
1173 token: String::from_utf8(read(token_secret)?)
1174 .map_err(|_| Error::invalid("the git token is not text"))?
1175 .trim()
1176 .to_string(),
1177 },
1178 GitAuth::SshKey { ssh_key_secret } => Credentials::SshKey {
1179 private_key: read(ssh_key_secret)?,
1180 },
1181 })
1182 }
1183
1184 fn prune(&self, org: &OrgId, app: &App) {
1187 let Ok(rd) = std::fs::read_dir(self.deployments_dir(org, &app.spec.name)) else {
1188 return;
1189 };
1190 let mut ids: Vec<u64> = rd
1191 .flatten()
1192 .filter_map(|e| {
1193 e.file_name()
1194 .into_string()
1195 .ok()?
1196 .strip_suffix(".json")?
1197 .parse()
1198 .ok()
1199 })
1200 .collect();
1201 ids.sort_unstable();
1202 let excess = ids.len().saturating_sub(KEEP_DEPLOYMENTS);
1203 for id in ids.into_iter().take(excess) {
1204 if Some(id) == app.current {
1205 continue;
1206 }
1207 let _ = std::fs::remove_file(self.dep_path(org, &app.spec.name, id));
1208 let _ = std::fs::remove_file(self.log_path(org, &app.spec.name, id));
1209 }
1210 }
1211
1212 pub(super) fn event(&self, org: &OrgId, app: &str, level: &str, message: String) {
1214 let stack = self
1215 .get(org, app)
1216 .ok()
1217 .and_then(|a| a.spec.stack().ok())
1218 .unwrap_or_default();
1219 let q = crate::stack::qualified(org, &stack);
1220 self.inner
1221 .ctl
1222 .note_service(level, &q, app, format!("app {app}: {message}"));
1223 }
1224
1225 fn kind_event(&self, org: &OrgId, app: &str, kind: &str, level: &str, message: String) {
1227 let stack = self
1228 .get(org, app)
1229 .ok()
1230 .and_then(|a| a.spec.stack().ok())
1231 .unwrap_or_default();
1232 let q = crate::stack::qualified(org, &stack);
1233 self.inner
1234 .ctl
1235 .event(kind, level, &q, app, format!("app {app}: {message}"));
1236 }
1237
1238 pub fn webhook(
1244 &self,
1245 org: &str,
1246 app: &str,
1247 header: &dyn Fn(&str) -> Option<String>,
1248 token: Option<&str>,
1249 body: &[u8],
1250 ) -> (u16, Value) {
1251 let refuse = |why: &str| (401, json!({"error": "unauthorized", "message": why}));
1252 let (Ok(org), Ok(())) = (OrgId::new(org), super::validate_app_name(app)) else {
1253 return refuse("invalid signature");
1254 };
1255 let Ok(found) = self.get(&org, app) else {
1256 return refuse("invalid signature");
1257 };
1258 let Ok((secret, _)) = self.inner.secrets.get(&org, &found.spec.webhook_secret()) else {
1259 return refuse("invalid signature");
1260 };
1261 let provider = match webhook::verify(trim_ascii(&secret), header, token, body) {
1262 Ok(p) => p,
1263 Err(webhook::Refusal::Missing) => {
1264 return refuse(
1265 "sign the request (X-Hub-Signature-256, X-Gitea-Signature, X-Gitlab-Token) or pass ?token=",
1266 );
1267 }
1268 Err(webhook::Refusal::Invalid) => return refuse("invalid signature"),
1269 };
1270 if let Some(id) = webhook::delivery_id(header) {
1271 let key = format!("{org}/{app}/{id}");
1272 let mut seen = self.inner.deliveries.lock().unwrap();
1273 if seen.contains(&key) {
1274 return (200, json!({"ignored": "delivery already received"}));
1275 }
1276 if seen.len() >= 1000 {
1277 seen.pop_front();
1278 }
1279 seen.push_back(key);
1280 }
1281 let by = format!("webhook:{provider:?}").to_lowercase();
1282 let requested = match webhook::event(provider, header, body) {
1283 webhook::Event::Ping => return (200, json!({"ok": true, "ping": true})),
1284 webhook::Event::Other(kind) => {
1285 return (200, json!({"ignored": format!("event {kind}")}));
1286 }
1287 webhook::Event::Trigger => None,
1288 webhook::Event::PullRequest(pr) => {
1289 return self.preview_webhook(&org, &found, provider, &by, &pr);
1290 }
1291 webhook::Event::Push {
1292 reference,
1293 after,
1294 message,
1295 deleted,
1296 } => {
1297 if deleted {
1298 return (200, json!({"ignored": format!("{reference} was deleted")}));
1299 }
1300 if let Source::Git(g) = &found.spec.source {
1301 if !g.matches_push(&reference) {
1302 return (
1303 200,
1304 json!({"ignored": format!("{reference} is not {}", g.reference)}),
1305 );
1306 }
1307 }
1308 let mut r = reference;
1309 if let Some(a) = after {
1310 r = format!("{r} {a}");
1311 }
1312 if let Some(m) = message {
1313 r = format!("{r}: {m}");
1314 }
1315 Some(r)
1316 }
1317 };
1318 match self.deploy(&org, app, Trigger::Webhook, &by, requested) {
1319 Ok(d) => (202, json!({"deployment": d.id, "status": d.status})),
1320 Err(e) => (500, json!({"error": "deploy", "message": e.to_string()})),
1321 }
1322 }
1323}
1324
1325fn trim_ascii(b: &[u8]) -> &[u8] {
1326 let s = b
1327 .iter()
1328 .position(|c| !c.is_ascii_whitespace())
1329 .unwrap_or(b.len());
1330 let e = b
1331 .iter()
1332 .rposition(|c| !c.is_ascii_whitespace())
1333 .map_or(s, |i| i + 1);
1334 &b[s..e]
1335}
1336
1337struct DeployLog {
1339 file: std::fs::File,
1340 apps: Apps,
1341 org: OrgId,
1342 app: String,
1343 id: u64,
1344}
1345
1346impl DeployLog {
1347 fn open(apps: &Apps, org: &OrgId, app: &str, id: u64) -> Result<DeployLog> {
1348 use std::os::unix::fs::OpenOptionsExt;
1349 let p = apps.log_path(org, app, id);
1350 if let Some(d) = p.parent() {
1351 std::fs::create_dir_all(d)?;
1352 }
1353 let file = std::fs::OpenOptions::new()
1354 .create(true)
1355 .append(true)
1356 .mode(0o600)
1357 .open(p)?;
1358 Ok(DeployLog {
1359 file,
1360 apps: apps.clone(),
1361 org: org.clone(),
1362 app: app.into(),
1363 id,
1364 })
1365 }
1366
1367 fn line(&mut self, l: &str) {
1368 let _ = writeln!(self.file, "{l}");
1369 self.apps
1370 .event(&self.org, &self.app, "log", format!("#{}: {l}", self.id));
1371 }
1372}
1373
1374pub fn skopeo_digest(image: &str) -> Option<String> {
1377 let src = crate::plan::ImageSource::parse(image).ok()?;
1378 if !src.is_oci() {
1379 return None;
1380 }
1381 let host = src.server.as_deref()?.strip_prefix("https://")?;
1382 let r = format!("docker://{host}/{}", src.alias);
1383 let mut child = std::process::Command::new("skopeo")
1384 .args(["inspect", "--no-tags", "--format", "{{.Digest}}", &r])
1385 .stdin(std::process::Stdio::null())
1386 .stdout(std::process::Stdio::piped())
1387 .stderr(std::process::Stdio::null())
1388 .spawn()
1389 .ok()?;
1390 let started = Instant::now();
1391 loop {
1392 match child.try_wait() {
1393 Ok(Some(s)) if s.success() => break,
1394 Ok(Some(_)) | Err(_) => return None,
1395 Ok(None) if started.elapsed() > Duration::from_secs(60) => {
1396 let _ = child.kill();
1397 let _ = child.wait();
1398 return None;
1399 }
1400 Ok(None) => std::thread::sleep(Duration::from_millis(100)),
1401 }
1402 }
1403 let mut out = String::new();
1404 use std::io::Read;
1405 child.stdout.take()?.read_to_string(&mut out).ok()?;
1406 let d = out.trim();
1407 (d.starts_with("sha256:") && d.len() == 71).then(|| d.to_string())
1408}
1409
1410#[cfg(test)]
1411mod tests {
1412 use super::*;
1413 use crate::app::EnvValue;
1414
1415 #[test]
1416 fn state_machine() {
1417 use Status::*;
1418 let mut d = Deployment {
1419 id: 1,
1420 app: "web".into(),
1421 trigger: Trigger::Api,
1422 by: "t".into(),
1423 status: Queued,
1424 requested: None,
1425 rollback_of: None,
1426 commit: None,
1427 image: None,
1428 digest: None,
1429 error: None,
1430 created_at: 0,
1431 started_at: None,
1432 finished_at: None,
1433 rendered: None,
1434 };
1435 assert!(d.advance(Done).is_err(), "queued cannot jump to done");
1436 d.advance(Building).unwrap();
1437 assert!(d.started_at.is_some());
1438 assert!(
1439 d.advance(Superseded).is_err(),
1440 "only a waiting one is superseded"
1441 );
1442 d.advance(Deploying).unwrap();
1443 d.advance(Done).unwrap();
1444 assert!(d.finished_at.is_some());
1445 for s in [Queued, Building, Deploying, Done, Failed, Superseded] {
1446 assert!(!Done.can_become(s) && !Failed.can_become(s) && !Superseded.can_become(s));
1447 }
1448 assert!(Queued.can_become(Superseded));
1449 assert!(Building.can_become(Failed));
1450 let v = d.summary();
1451 assert!(v.get("rendered").is_none());
1452 assert_eq!(v["status"], "done");
1453 assert_eq!(v["trigger"], "api");
1454 }
1455
1456 fn apps(dir: &Path, gate: Arc<(Mutex<bool>, std::sync::Condvar)>) -> Apps {
1459 let k = crate::secrets::Keyring::new(age::x25519::Identity::generate(), vec![]);
1460 let secrets = Arc::new(Secrets::new(crate::secrets::LocalDriver::new(
1461 dir,
1462 Arc::new(k),
1463 )));
1464 let client = Client::with_socket("/nonexistent/isb-test/incus.sock");
1465 let store = crate::stack::Store::open(dir).unwrap();
1466 let ctl = Controller::start(
1467 client.clone(),
1468 store,
1469 Duration::from_secs(60),
1470 secrets.clone(),
1471 )
1472 .unwrap();
1473 let digest: DigestFn = Arc::new(move |image: &str| {
1475 if image == "docker:slow" {
1476 let (m, cv) = &*gate;
1477 let mut open = m.lock().unwrap();
1478 while !*open {
1479 open = cv.wait(open).unwrap();
1480 }
1481 }
1482 None
1483 });
1484 Apps::new(dir, client, ctl, secrets).with_digest(digest)
1485 }
1486
1487 fn spec(v: Value) -> AppSpec {
1488 serde_json::from_value(v).unwrap()
1489 }
1490
1491 fn hdrs(h: &[(&str, &str)]) -> Box<webhook::Headers<'static>> {
1492 let h: Vec<(String, String)> = h
1493 .iter()
1494 .map(|(k, v)| (k.to_ascii_lowercase(), v.to_string()))
1495 .collect();
1496 Box::new(move |k: &str| {
1497 h.iter()
1498 .find(|(n, _)| *n == k.to_ascii_lowercase())
1499 .map(|(_, v)| v.clone())
1500 })
1501 }
1502
1503 #[test]
1504 fn rollout_events_reach_the_deployment_log() {
1505 let dir = tempfile::tempdir().unwrap();
1506 let gate = Arc::new((Mutex::new(true), std::sync::Condvar::new()));
1507 let ap = apps(dir.path(), gate);
1508 let ctl = ap.controller();
1509 ctl.note_service(
1510 "info",
1511 "acme/shop-production",
1512 "web",
1513 "before the deploy".into(),
1514 );
1515 let mark = ap.event_mark();
1516 ctl.note_service(
1517 "info",
1518 "acme/shop-production",
1519 "web",
1520 "app web: #1: own line".into(),
1521 );
1522 ctl.note_service(
1523 "info",
1524 "acme/shop-production",
1525 "web",
1526 "rolling out rev 1 to 1 slot(s)".into(),
1527 );
1528 ctl.note_service(
1529 "info",
1530 "acme/shop-production",
1531 "api",
1532 "another service".into(),
1533 );
1534 ctl.note_service("info", "acme/other", "web", "another stack".into());
1535 let mut seen = Vec::new();
1536 let r = ap.wait_service_with("acme/shop-production", "web", Some(mark), |m| {
1538 seen.push(m.to_string())
1539 });
1540 assert!(r.is_err());
1541 assert_eq!(seen, ["rolling out rev 1 to 1 slot(s)"]);
1542 let mut none = Vec::new();
1544 let _ = ap.wait_service_with("acme/shop-production", "web", None, |m| {
1545 none.push(m.to_string())
1546 });
1547 assert!(none.is_empty());
1548 }
1549
1550 #[test]
1551 #[expect(
1552 clippy::too_many_lines,
1553 clippy::cognitive_complexity,
1554 reason = "predates the lint ratchet; split it when next changed"
1555 )]
1556 fn apps_records_webhooks_and_queue() {
1557 let dir = tempfile::tempdir().unwrap();
1558 let gate = Arc::new((Mutex::new(false), std::sync::Condvar::new()));
1559 let ap = apps(dir.path(), gate.clone());
1560 let org = OrgId::new("acme").unwrap();
1561
1562 let p = ap.project_create(&org, "shop", "", &[]).unwrap();
1564 assert_eq!(p.environments, ["production"]);
1565 assert!(ap.project_create(&org, "shop", "", &[]).is_err());
1566 ap.environment_create(&org, "shop", "staging").unwrap();
1567 assert!(
1568 dir.path()
1569 .join("orgs/acme/apps/projects/shop.json")
1570 .is_file()
1571 );
1572
1573 let web = json!({
1575 "name": "web", "project": "shop",
1576 "source": {"image": "docker:traefik/whoami"},
1577 "env": "# greeting\nA=1\nT=${{secret.tok}}\n",
1578 });
1579 assert!(ap.create(&org, spec(web.clone())).is_err());
1580 ap.inner.secrets.set(&org, "tok", b"v").unwrap();
1581 let (app, secret) = ap.create(&org, spec(web.clone())).unwrap();
1582 assert_eq!(secret.len(), 64);
1583 assert_eq!(app.spec.environment, "production");
1584 assert!(ap.create(&org, spec(web)).is_err(), "names are unique");
1585 let mut bad = spec(
1586 json!({"name": "x", "project": "shop", "environment": "qa", "source": {"image": "x"}}),
1587 );
1588 assert!(ap.create(&org, bad.clone()).is_err(), "no such environment");
1589 bad.project = "nope".into();
1590 assert!(ap.create(&org, bad).is_err(), "no such project");
1591
1592 let text = ap.env_get(&org, "web").unwrap();
1594 assert_eq!(text, "# greeting\nA=1\nT=${{secret.tok}}\n");
1595 ap.env_set(
1596 &org,
1597 "web",
1598 "# greeting\nA=2\nT=${{secret.tok}}\nB=\"two words\"\n",
1599 )
1600 .unwrap();
1601 assert!(
1602 ap.env_get(&org, "web")
1603 .unwrap()
1604 .contains("A=2\nT=${{secret.tok}}\nB=\"two words\"")
1605 );
1606 assert!(ap.env_set(&org, "web", "T=${{secret.missing}}\n").is_err());
1607 let u = ap
1608 .update(&org, "web", &json!({"replicas": 3, "port": 80}))
1609 .unwrap();
1610 assert_eq!(u.spec.replicas, 3);
1611 assert_eq!(u.spec.env.get("A"), Some(&EnvValue::Plain("2".into())));
1612 assert!(
1613 ap.update(&org, "web", &json!({"project": "other"}))
1614 .is_err()
1615 );
1616 assert!(
1617 ap.update(&org, "web", &json!({"port": null}))
1618 .unwrap()
1619 .spec
1620 .port
1621 .is_none()
1622 );
1623 assert!(ap.rollback(&org, "web", None, Trigger::Api, "t").is_err());
1624
1625 let push = br#"{"ref":"refs/heads/main","after":"abc","head_commit":{"message":"m"}}"#;
1627 let none = hdrs(&[]);
1628 assert_eq!(ap.webhook("acme", "web", &none, None, push).0, 401);
1629 assert_eq!(ap.webhook("acme", "web", &none, Some("wrong"), push).0, 401);
1630 assert_eq!(
1631 ap.webhook("acme", "nope", &none, Some(&secret), push).0,
1632 401
1633 );
1634 assert_eq!(
1635 ap.webhook("Bad Org", "web", &none, Some(&secret), push).0,
1636 401
1637 );
1638 let forged = hdrs(&[
1639 ("X-GitHub-Event", "push"),
1640 (
1641 "X-Hub-Signature-256",
1642 &format!("sha256={}", webhook::sign(b"other", push)),
1643 ),
1644 ]);
1645 assert_eq!(ap.webhook("acme", "web", &forged, None, push).0, 401);
1646 assert!(ap.deployments(&org, "web").unwrap().is_empty());
1647 let signed = |event: &str, delivery: &str, body: &[u8]| {
1648 hdrs(&[
1649 ("X-GitHub-Event", event),
1650 ("X-GitHub-Delivery", delivery),
1651 (
1652 "X-Hub-Signature-256",
1653 &format!("sha256={}", webhook::sign(secret.as_bytes(), body)),
1654 ),
1655 ])
1656 };
1657 let (st, v) = ap.webhook("acme", "web", &signed("ping", "d0", b"{}"), None, b"{}");
1658 assert_eq!((st, v["ping"].as_bool()), (200, Some(true)));
1659 let (st, v) = ap.webhook("acme", "web", &signed("issues", "d1", push), None, push);
1660 assert_eq!(st, 200);
1661 assert!(v["ignored"].as_str().unwrap().contains("issues"), "{v}");
1662 let (st, v) = ap.webhook("acme", "web", &signed("push", "d2", push), None, push);
1663 assert_eq!(st, 202, "{v}");
1664 assert_eq!(v["deployment"], 1);
1665 let (st, v) = ap.webhook("acme", "web", &signed("push", "d2", push), None, push);
1667 assert_eq!(st, 200, "{v}");
1668 let d = ap.wait(&org, "web", 1, Duration::from_secs(30)).unwrap();
1669 assert_eq!(d.trigger, Trigger::Webhook);
1670 assert_eq!(d.by, "webhook:github");
1671 assert_eq!(d.requested.as_deref(), Some("refs/heads/main abc: m"));
1672 assert_eq!(d.status, Status::Failed, "{d:?}");
1674 assert!(d.image.is_some());
1675 let (log, _, done) = ap.log(&org, "web", 1, 0).unwrap();
1676 assert!(done && log.contains("image docker:traefik/whoami"), "{log}");
1677 let (rest, next, rec) = ap.log_and_record(&org, "web", 1, 4).unwrap();
1679 assert_eq!((rec.id, rec.status), (1, Status::Failed));
1680 assert_eq!(next as usize, log.len());
1681 assert_eq!(rest, log[4..]);
1682 assert!(rec.summary().get("rendered").is_none());
1683 let ds = ap.deployments(&org, "web").unwrap();
1684 assert_eq!(ds.len(), 1);
1685
1686 ap.inner.secrets.set(&org, "gh", b"t").unwrap();
1688 let (_, gsecret) = ap
1689 .create(
1690 &org,
1691 spec(json!({
1692 "name": "api", "project": "shop", "environment": "staging",
1693 "source": {"git": {"url": "https://example.invalid/o/r.git", "ref": "main", "auth": {"token_secret": "gh"}}},
1694 "build": {"builder": {"type": "railpack"}},
1695 })),
1696 )
1697 .unwrap();
1698 let dev = br#"{"ref":"refs/heads/dev","after":"abc"}"#;
1699 let gl = |body: &[u8]| {
1700 let _ = body;
1701 hdrs(&[
1702 ("X-Gitlab-Event", "Push Hook"),
1703 ("X-Gitlab-Token", &gsecret),
1704 ])
1705 };
1706 let (st, v) = ap.webhook("acme", "api", &gl(dev), None, dev);
1707 assert_eq!(st, 200);
1708 assert!(
1709 v["ignored"].as_str().unwrap().contains("refs/heads/dev"),
1710 "{v}"
1711 );
1712 assert!(ap.deployments(&org, "api").unwrap().is_empty());
1713 let (st, _) = ap.webhook("acme", "api", &gl(push), None, push);
1714 assert_eq!(st, 202);
1715 let d = ap.wait(&org, "api", 1, Duration::from_secs(60)).unwrap();
1716 assert_eq!(d.status, Status::Failed);
1717 assert!(d.error.as_deref().unwrap_or("").contains("git"), "{d:?}");
1718 let s2 = ap.webhook_secret(&org, "api", true).unwrap();
1720 assert_ne!(s2, gsecret);
1721 assert_eq!(ap.webhook("acme", "api", &gl(push), None, push).0, 401);
1722
1723 ap.create(
1725 &org,
1726 spec(json!({"name": "slow", "project": "shop", "source": {"image": "docker:slow"}})),
1727 )
1728 .unwrap();
1729 let d1 = ap.deploy(&org, "slow", Trigger::Api, "t", None).unwrap();
1730 let started = Instant::now();
1731 while ap.deployment(&org, "slow", d1.id).unwrap().status != Status::Building {
1732 assert!(started.elapsed() < Duration::from_secs(10));
1733 std::thread::sleep(Duration::from_millis(20));
1734 }
1735 let d2 = ap.deploy(&org, "slow", Trigger::Api, "t", None).unwrap();
1736 let d3 = ap.deploy(&org, "slow", Trigger::Manual, "t", None).unwrap();
1737 assert_eq!(
1738 ap.deployment(&org, "slow", d2.id).unwrap().status,
1739 Status::Superseded
1740 );
1741 assert_eq!(
1742 ap.deployment(&org, "slow", d3.id).unwrap().status,
1743 Status::Queued
1744 );
1745 assert!(ap.delete(&org, "slow").is_err(), "not while deploying");
1746 {
1747 let (m, cv) = &*gate;
1748 *m.lock().unwrap() = true;
1749 cv.notify_all();
1750 }
1751 let d3 = ap
1752 .wait(&org, "slow", d3.id, Duration::from_secs(30))
1753 .unwrap();
1754 assert!(d3.status.finished());
1755 assert!(
1756 ap.wait(&org, "slow", d1.id, Duration::from_secs(1))
1757 .unwrap()
1758 .status
1759 .finished()
1760 );
1761
1762 assert!(ap.project_delete(&org, "shop").is_err());
1764 assert!(ap.environment_delete(&org, "shop", "staging").is_err());
1765 let started = Instant::now();
1766 for a in ["web", "api", "slow"] {
1767 while let Err(e) = ap.delete(&org, a) {
1768 assert!(started.elapsed() < Duration::from_secs(10), "{a}: {e}");
1769 std::thread::sleep(Duration::from_millis(50));
1770 }
1771 }
1772 assert!(ap.inner.secrets.inspect(&org, "app.web.webhook").is_err());
1773 assert!(
1774 ap.inner.secrets.inspect(&org, "tok").is_ok(),
1775 "not isb's to delete"
1776 );
1777 ap.environment_delete(&org, "shop", "staging").unwrap();
1778 ap.project_delete(&org, "shop").unwrap();
1779 assert!(ap.project_list(&org).unwrap().is_empty());
1780 }
1781
1782 #[test]
1783 #[expect(
1784 clippy::too_many_lines,
1785 clippy::cognitive_complexity,
1786 reason = "predates the lint ratchet; split it when next changed"
1787 )]
1788 fn previews_from_webhooks() {
1789 let dir = tempfile::tempdir().unwrap();
1790 let gate = Arc::new((Mutex::new(true), std::sync::Condvar::new()));
1791 let ap = apps(dir.path(), gate);
1792 let org = OrgId::new("acme").unwrap();
1793 ap.project_create(&org, "shop", "", &[]).unwrap();
1794 assert!(
1795 ap.environment_create(&org, "shop", "prod-pr-1").is_err(),
1796 "preview stack names are kept"
1797 );
1798 let (_, secret) = ap
1799 .create(
1800 &org,
1801 spec(json!({
1802 "name": "web", "project": "shop",
1803 "source": {"git": {"url": "https://example.invalid/acme/web.git", "ref": "main"}},
1804 "build": {"builder": {"type": "dockerfile"}},
1805 "port": 8080,
1806 })),
1807 )
1808 .unwrap();
1809 let pr = |action: &str, n: u64, head_repo: &str, base: &str| {
1810 format!(
1811 r#"{{"action":"{action}","number":{n},"pull_request":{{"title":"t{n}","merged":false,
1812 "head":{{"ref":"f{n}","sha":"{}","repo":{{"full_name":"{head_repo}"}}}},
1813 "base":{{"ref":"{base}","repo":{{"full_name":"acme/web"}}}}}}}}"#,
1814 "a".repeat(40)
1815 )
1816 .into_bytes()
1817 };
1818 let mut delivery = 0;
1819 let mut send = |body: &[u8]| {
1820 delivery += 1;
1821 let h = hdrs(&[
1822 ("X-Gitea-Event", "pull_request"),
1823 ("X-Gitea-Delivery", &format!("p{delivery}")),
1824 ("X-Gitea-Signature", &webhook::sign(secret.as_bytes(), body)),
1825 ]);
1826 ap.webhook("acme", "web", &h, None, body)
1827 };
1828 let (st, v) = send(&pr("opened", 1, "acme/web", "main"));
1830 assert_eq!(st, 200);
1831 assert!(v["ignored"].as_str().unwrap().contains("off"), "{v}");
1832 ap.update(
1833 &org,
1834 "web",
1835 &json!({"previews": {"enabled": true, "max": 2, "env": "MODE=preview\n"}}),
1836 )
1837 .unwrap();
1838 let (st, v) = send(&pr("opened", 1, "acme/web", "dev"));
1840 assert_eq!(st, 200);
1841 assert!(
1842 v["ignored"].as_str().unwrap().contains("targets dev"),
1843 "{v}"
1844 );
1845 let (st, v) = send(&pr("opened", 1, "mallory/web", "main"));
1846 assert_eq!(st, 200);
1847 assert!(v["ignored"].as_str().unwrap().contains("fork"), "{v}");
1848 assert!(ap.preview_list(&org, "web").unwrap().is_empty());
1849 let (st, v) = send(&pr("opened", 1, "acme/web", "main"));
1851 assert_eq!(st, 202, "{v}");
1852 assert_eq!(
1853 (v["preview"].as_u64(), v["deployment"].as_u64()),
1854 (Some(1), Some(1))
1855 );
1856 let p = ap.preview_get(&org, "web", 1).unwrap();
1857 assert_eq!(p.stack, "shop-production-pr-1");
1858 assert_eq!((p.provider.as_str(), p.fork), ("gitea", false));
1859 let d = ap
1860 .preview_wait(&org, "web", 1, 1, Duration::from_secs(60))
1861 .unwrap();
1862 assert_eq!(d.status, Status::Failed, "{d:?}");
1863 let (log, _, done) = ap.preview_log(&org, "web", 1, 1, 0).unwrap();
1864 assert!(done && log.contains("refs/pull/1/head"), "{log}");
1865 let (st, v) = send(&pr("synchronize", 1, "acme/web", "main"));
1868 assert_eq!(st, 200, "{v}");
1869 assert!(v["ignored"].as_str().unwrap().contains("already"), "{v}");
1870 let moved = String::from_utf8(pr("synchronize", 1, "acme/web", "main"))
1871 .unwrap()
1872 .replace(&"a".repeat(40), &"b".repeat(40));
1873 let (st, v) = send(moved.as_bytes());
1874 assert_eq!((st, v["deployment"].as_u64()), (202, Some(2)), "{v}");
1875 ap.preview_wait(&org, "web", 1, 2, Duration::from_secs(60))
1876 .unwrap();
1877 assert_eq!(send(&pr("opened", 2, "acme/web", "main")).0, 202);
1879 let (st, v) = send(&pr("opened", 3, "acme/web", "main"));
1880 assert_eq!(st, 200);
1881 assert!(v["ignored"].as_str().unwrap().contains("2 previews"), "{v}");
1882 assert!(ap.preview_get(&org, "web", 3).is_err());
1883 ap.preview_wait(&org, "web", 2, 1, Duration::from_secs(60))
1884 .unwrap();
1885 ap.update(&org, "web", &json!({"previews": {"forks": true, "max": 5}}))
1888 .unwrap();
1889 assert_eq!(send(&pr("opened", 4, "mallory/web", "main")).0, 202);
1890 assert!(ap.preview_get(&org, "web", 4).unwrap().fork);
1891 ap.preview_wait(&org, "web", 4, 1, Duration::from_secs(60))
1892 .unwrap();
1893 let (st, v) = send(&pr("closed", 1, "acme/web", "main"));
1895 assert_eq!(st, 202, "{v}");
1896 let started = Instant::now();
1897 while ap.preview_get(&org, "web", 1).is_ok() {
1898 assert!(started.elapsed() < Duration::from_secs(30));
1899 std::thread::sleep(Duration::from_millis(50));
1900 }
1901 assert!(!dir.path().join("orgs/acme/apps/web/previews/1").exists());
1902 let (st, v) = send(&pr("closed", 1, "acme/web", "main"));
1903 assert_eq!(st, 200, "{v}");
1904 let d = ap
1906 .preview_redeploy(&org, "web", 2, Trigger::Api, "t")
1907 .unwrap();
1908 assert_eq!(d.id, 2);
1909 ap.preview_wait(&org, "web", 2, 2, Duration::from_secs(60))
1910 .unwrap();
1911 ap.preview_remove(&org, "web", 2, "test", true).unwrap();
1912 assert!(ap.preview_get(&org, "web", 2).is_err());
1913 assert!(ap.delete(&org, "web").is_err());
1916 assert!(ap.get(&org, "web").is_ok());
1917 assert!(ap.preview_get(&org, "web", 4).unwrap().removing);
1918 }
1919}