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