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