1use super::deploy::{Apps, Status};
5use crate::error::Result;
6use crate::org::OrgId;
7
8impl Apps {
9 pub fn cancel_queued(&self, org: &OrgId, app: &str, id: u64, reason: &str) -> Result<bool> {
13 if let Some(q) = self
14 .inner
15 .queues
16 .lock()
17 .unwrap()
18 .get_mut(&(org.clone(), app.to_string()))
19 {
20 if q.next == Some(id) {
21 q.next = None;
22 }
23 }
24 let g = self.inner.edit.lock().unwrap();
25 let mut d = self.deployment(org, app, id)?;
26 if d.status != Status::Queued {
27 return Ok(false);
28 }
29 d.advance(Status::Cancelled)?;
30 d.error = Some(reason.to_string());
31 self.save_dep(org, &d)?;
32 drop(g);
33 self.event(
34 org,
35 app,
36 "info",
37 format!("deployment {id} cancelled: {reason}"),
38 );
39 Ok(true)
40 }
41
42 pub(super) fn recover(&self) {
44 let mut orgs = vec![OrgId::default_org()];
45 if let Ok(rd) = std::fs::read_dir(self.inner.state.join("orgs")) {
46 for e in rd.flatten() {
47 if let Some(o) = e.file_name().to_str().and_then(|s| OrgId::new(s).ok()) {
48 orgs.push(o);
49 }
50 }
51 }
52 for org in orgs {
53 for app in self.list(&org).unwrap_or_default() {
54 for mut d in self.deployments(&org, &app.spec.name).unwrap_or_default() {
55 if d.status.finished() {
56 continue;
57 }
58 d.error = Some("interrupted: the daemon stopped during it".into());
59 d.status = Status::Failed;
60 d.finished_at = Some(crate::stack::controller::now_ms());
61 let _ = self.save_dep(&org, &d);
62 }
63 self.recover_previews(&org, &app.spec.name);
64 }
65 }
66 }
67}
68
69#[cfg(test)]
70mod tests {
71 use std::sync::{Arc, Condvar, Mutex};
72 use std::time::{Duration, Instant};
73
74 use serde_json::json;
75
76 use super::*;
77 use crate::app::deploy::Trigger;
78 use crate::client::Client;
79 use crate::stack::Controller;
80
81 fn apps(dir: &std::path::Path, gate: Arc<(Mutex<bool>, Condvar)>) -> Apps {
84 let k = crate::secrets::Keyring::new(age::x25519::Identity::generate(), vec![]);
85 let secrets = Arc::new(crate::secrets::Secrets::new(
86 crate::secrets::LocalDriver::new(dir, Arc::new(k)),
87 ));
88 let client = Client::with_socket("/nonexistent/isb-test/incus.sock");
89 let store = crate::stack::Store::open(dir).unwrap();
90 let ctl = Controller::start(
91 client.clone(),
92 store,
93 Duration::from_secs(60),
94 secrets.clone(),
95 )
96 .unwrap();
97 Apps::new(dir, client, ctl, secrets).with_probe(Arc::new(move |_: &str, _| {
98 let (m, cv) = &*gate;
99 let mut open = m.lock().unwrap();
100 while !*open {
101 open = cv.wait(open).unwrap();
102 }
103 crate::image_check::Probe::Found(None)
104 }))
105 }
106
107 #[test]
108 fn a_queued_deployment_is_cancelled_and_a_restart_closes_unfinished_ones() {
109 let dir = tempfile::tempdir().unwrap();
110 let gate = Arc::new((Mutex::new(false), Condvar::new()));
111 let ap = apps(dir.path(), gate.clone());
112 let org = OrgId::new("acme").unwrap();
113 ap.project_create(&org, "shop", "", &[]).unwrap();
114 let spec = json!({"name": "slow", "project": "shop", "source": {"image": "docker:slow"}});
115 ap.create(&org, serde_json::from_value(spec).unwrap())
116 .unwrap();
117 let d1 = ap.deploy(&org, "slow", Trigger::Api, "t", None).unwrap();
118 let started = Instant::now();
119 while ap.deployment(&org, "slow", d1.id).unwrap().status != Status::Building {
120 assert!(started.elapsed() < Duration::from_secs(10));
121 std::thread::sleep(Duration::from_millis(20));
122 }
123 let d2 = ap
124 .deploy(&org, "slow", Trigger::Api, "t", Some("template".into()))
125 .unwrap();
126 let why = "template deploy stopped: db failed";
128 assert!(ap.cancel_queued(&org, "slow", d2.id, why).unwrap());
129 assert!(!ap.cancel_queued(&org, "slow", d1.id, why).unwrap());
130 let d = ap.deployment(&org, "slow", d2.id).unwrap();
131 assert_eq!(d.status, Status::Cancelled);
132 assert_eq!(d.error.as_deref(), Some(why));
133 assert!(d.status.finished() && !d.status.can_become(Status::Queued));
134 let d3 = ap.deploy(&org, "slow", Trigger::Api, "t", None).unwrap();
136 let open = Arc::new((Mutex::new(true), Condvar::new()));
137 let ap2 = apps(dir.path(), open);
138 for id in [d1.id, d3.id] {
139 let d = ap2.deployment(&org, "slow", id).unwrap();
140 assert_eq!(d.status, Status::Failed, "{d:?}");
141 assert!(d.error.unwrap().contains("interrupted"));
142 }
143 let d = ap2.deployment(&org, "slow", d2.id).unwrap();
144 assert_eq!(d.status, Status::Cancelled);
145 let (m, cv) = &*gate;
146 *m.lock().unwrap() = true;
147 cv.notify_all();
148 let _ = ap.wait(&org, "slow", d3.id, Duration::from_secs(10));
149 }
150}