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