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