1use serde::{Deserialize, Serialize};
12use serde_json::{Value, json};
13
14use super::{VolumeBackups, incus, model, org_pool};
15use crate::backup::{Compression, Destination};
16use crate::client::Client;
17use crate::error::{Error, Result};
18use crate::jobs::{Run, RunLog, RunStatus, RunTrigger};
19use crate::org::OrgId;
20
21#[derive(Debug, Clone, Default, Deserialize)]
23#[serde(deny_unknown_fields)]
24pub struct VolumeRestoreRequest {
25 pub name: String,
27 #[serde(default)]
28 pub snapshot: Option<String>,
29 #[serde(default)]
31 pub backup: Option<String>,
32 #[serde(default)]
34 pub destination: Option<String>,
35 #[serde(default)]
36 pub key: Option<String>,
37 #[serde(default)]
39 pub instance: Option<String>,
40}
41
42#[derive(Debug, Clone, Serialize)]
44pub struct StagedRestore {
45 pub volume: String,
47 pub of: String,
49 pub from: String,
51 pub stamp: String,
52 pub by: String,
53 pub instance: Option<String>,
55 pub path: String,
56 pub attached: bool,
57}
58
59enum Origin {
60 Snapshot(String),
61 Object(Destination, String, Compression),
62}
63
64impl Origin {
65 fn label(&self) -> String {
66 match self {
67 Origin::Snapshot(s) => format!("snapshot:{s}"),
68 Origin::Object(_, k, _) => format!("backup:{k}"),
69 }
70 }
71}
72
73impl VolumeBackups {
74 fn restore_source(
76 &self,
77 org: &OrgId,
78 req: &VolumeRestoreRequest,
79 oc: &Client,
80 pool: &str,
81 ) -> Result<Origin> {
82 let bk = &self.inner.backups;
83 let (dest, prefix) = match (&req.snapshot, &req.backup, &req.destination) {
84 (Some(s), None, None) => {
85 if !incus::snapshots(oc, pool, &req.name)?
86 .iter()
87 .any(|x| &x.name == s)
88 {
89 return Err(Error::NotFound(format!(
90 "snapshot {s} of volume {}",
91 req.name
92 )));
93 }
94 return Ok(Origin::Snapshot(s.clone()));
95 }
96 (None, Some(b), None) => {
97 let b = bk.get(org, b)?;
98 if b.spec.volume.as_deref() != Some(req.name.as_str()) {
99 return Err(Error::invalid(format!(
100 "backup {} does not back up volume {}",
101 b.spec.name, req.name
102 )));
103 }
104 let d = bk.destination_get(org, &b.spec.destination)?;
105 let p = crate::backup::backup_prefix(&d, org, &b.spec.name);
106 (d, p)
107 }
108 (None, None, Some(d)) => {
109 let key = req.key.as_deref().ok_or_else(|| {
110 Error::invalid("with `destination`, name the object with `key`")
111 })?;
112 let p = key
113 .rfind('/')
114 .map(|i| key[..=i].to_string())
115 .unwrap_or_default();
116 (bk.destination_get(org, d)?, p)
117 }
118 _ => {
119 return Err(Error::invalid(
120 "restore from a `snapshot`, a `backup` (and maybe `key`), or a `destination` and `key`: one of them",
121 ));
122 }
123 };
124 let key = match &req.key {
125 Some(k) => k.clone(),
126 None => super::export::files(bk, org, &dest, &prefix)?
127 .into_iter()
128 .next()
129 .map(|f| f.key)
130 .ok_or_else(|| Error::invalid(format!("no volume backups under {prefix}")))?,
131 };
132 let (_, _, c) = model::parse_volume_key(&prefix, &key)
133 .ok_or_else(|| Error::invalid(format!("{key}: not a volume backup isb wrote")))?;
134 Ok(Origin::Object(dest, key, c))
135 }
136
137 fn restore_instance(
140 &self,
141 oc: &Client,
142 pool: &str,
143 req: &VolumeRestoreRequest,
144 ) -> Result<Option<String>> {
145 if let Some(i) = &req.instance {
146 incus::is_running(oc, i)?;
147 return Ok(Some(i.clone()));
148 }
149 let Some(info) = crate::volume::get(oc, pool, &req.name)? else {
150 return Ok(None);
151 };
152 let users = incus::instances_using(&info);
153 let running = users
154 .iter()
155 .find(|i| incus::is_running(oc, i).unwrap_or(false));
156 Ok(running.or(users.first()).cloned())
157 }
158
159 pub fn restore(
162 &self,
163 org: &OrgId,
164 req: VolumeRestoreRequest,
165 by: &str,
166 ) -> Result<(Run, model::Staged)> {
167 model::validate_volume_name(&req.name)?;
168 let (oc, pool) = org_pool(self.client(), org)?;
169 let from = self.restore_source(org, &req, &oc, &pool)?;
170 let instance = self.restore_instance(&oc, &pool, &req)?;
171 let stamp = model::stamp(super::now());
172 let staged = model::staged(&req.name, &stamp);
173 if crate::volume::get(&oc, &pool, &staged.volume)?.is_some() {
174 return Err(Error::AlreadyExists(format!(
175 "volume {} (a restore started this second; try again)",
176 staged.volume
177 )));
178 }
179 let Some(guard) = self
180 .inner
181 .running
182 .enter(org, "restore", &staged.volume, true)
183 else {
184 return Err(Error::invalid(format!(
185 "a restore into {} is running",
186 staged.volume
187 )));
188 };
189 let store = self.inner.backups.restore_runs(org);
190 let (mut r, log) = store.start("restore", RunTrigger::Manual, by, None, 50)?;
191 r.detail = json!({
192 "volume": req.name, "target": staged.volume, "new": true, "staged": true,
193 "from": from.label(), "instance": instance, "path": staged.path, "stamp": stamp,
194 });
195 if let Origin::Object(d, k, _) = &from {
196 r.detail["key"] = json!(k);
197 r.detail["destination"] = json!(d.name);
198 }
199 store.save(&r)?;
200 let me = self.clone();
201 let (org2, run, st, by) = (org.clone(), r.clone(), staged.clone(), by.to_string());
202 std::thread::spawn(move || {
203 let _guard = guard;
204 let job = Job {
205 org: &org2,
206 volume: &req.name,
207 staged: &st,
208 from: &from,
209 instance: instance.as_deref(),
210 by: &by,
211 stamp: &stamp,
212 };
213 me.restore_run(&job, run, log);
214 });
215 Ok((r, staged))
216 }
217
218 fn restore_run(&self, job: &Job, mut r: Run, mut log: RunLog) {
219 let res = self.restore_once(job, &mut log);
220 let (kind, level, msg) = match res {
221 Ok(attached) => {
222 r.detail["attached"] = json!(attached);
223 r.exit_code = Some(0);
224 r.finish(RunStatus::Succeeded);
225 let at = match (attached, job.instance) {
226 (true, Some(i)) => format!("mounted at {} in {i}", job.staged.path),
227 _ => "detached".to_string(),
228 };
229 (
230 "volume.restore.staged",
231 "info",
232 format!(
233 "{} of {} staged as {} ({at})",
234 job.from.label(),
235 job.volume,
236 job.staged.volume
237 ),
238 )
239 }
240 Err(e) => {
241 log.line(&format!("isb: {e}"));
242 r.error = Some(e.to_string());
243 r.finish(RunStatus::Failed);
244 (
245 "volume.restore.failed",
246 "error",
247 format!("restore of {} failed: {e}", job.volume),
248 )
249 }
250 };
251 if let Err(e) = self
252 .inner
253 .backups
254 .restore_runs(job.org)
255 .finish(&mut r, &mut log)
256 {
257 eprintln!("isb serve: restore {}: {e}", r.id);
258 }
259 self.event(job.org, job.volume, kind, level, msg);
260 }
261
262 fn restore_once(&self, job: &Job, log: &mut RunLog) -> Result<bool> {
264 let (oc, pool) = org_pool(self.client(), job.org)?;
265 let mut labels = vec![
266 (model::KEY_RESTORE_OF, job.volume.to_string()),
267 (model::KEY_RESTORE_FROM, job.from.label()),
268 (model::KEY_RESTORE_STAMP, job.stamp.to_string()),
269 (model::KEY_RESTORE_BY, job.by.to_string()),
270 ];
271 let made = self.make_staged(job, &oc, &pool, &labels, log);
272 if let Err(e) = made {
273 if crate::volume::get(&oc, &pool, &job.staged.volume)?.is_some() {
274 log.line(&format!("isb: deleting the partial {}", job.staged.volume));
275 let _ = incus::delete_volume(&oc, &pool, &job.staged.volume);
276 }
277 return Err(e);
278 }
279 let Some(inst) = job.instance else {
280 log.line("isb: no instance uses the volume; the restore is left detached");
281 return Ok(false);
282 };
283 if !incus::is_running(&oc, inst)? {
284 log.line(&format!(
285 "isb: {inst} is stopped; {} is left detached (discard it, or restore again once {inst} runs)",
286 job.staged.volume
287 ));
288 return Ok(false);
289 }
290 incus::attach(
291 &oc,
292 inst,
293 &job.staged.device,
294 (&pool, &job.staged.volume),
295 &job.staged.path,
296 )?;
297 labels.push((model::KEY_RESTORE_INSTANCE, inst.to_string()));
298 incus::set_config(&oc, &pool, &job.staged.volume, &labels)?;
299 log.line(&format!(
300 "isb: mounted at {} in {inst}, read-write; the live volume is untouched",
301 job.staged.path
302 ));
303 Ok(true)
304 }
305
306 fn make_staged(
307 &self,
308 job: &Job,
309 oc: &Client,
310 pool: &str,
311 labels: &[(&str, String)],
312 log: &mut RunLog,
313 ) -> Result<()> {
314 match job.from {
315 Origin::Snapshot(s) => {
316 log.line(&format!(
317 "isb: copying snapshot {}/{s} to {}",
318 job.volume, job.staged.volume
319 ));
320 incus::copy_from_snapshot(oc, pool, (job.volume, s), &job.staged.volume, labels)
321 }
322 Origin::Object(d, key, c) => {
323 let s3 = self.inner.backups.client(job.org, d)?;
324 let (len, body) = s3.get(key)?;
325 log.line(&format!(
326 "isb: importing s3://{}/{key} ({len} bytes) as {}",
327 d.bucket, job.staged.volume
328 ));
329 let mut reader = crate::backup::decompressor(*c, body)?;
330 incus::import(oc, pool, &job.staged.volume, &mut reader)?;
331 incus::set_config(oc, pool, &job.staged.volume, labels)
332 }
333 }
334 }
335
336 pub fn staged_list(&self, org: &OrgId, volume: Option<&str>) -> Result<Vec<StagedRestore>> {
338 let (oc, pool) = org_pool(self.client(), org)?;
339 let mut out: Vec<StagedRestore> = crate::volume::list(&oc, &pool)?
340 .into_iter()
341 .filter_map(|v| {
342 let of = v.config.get(model::KEY_RESTORE_OF)?.clone();
343 if volume.is_some_and(|x| x != of) {
344 return None;
345 }
346 let get = |k: &str| v.config.get(k).cloned().unwrap_or_default();
347 let stamp = get(model::KEY_RESTORE_STAMP);
348 Some(StagedRestore {
349 from: get(model::KEY_RESTORE_FROM),
350 by: get(model::KEY_RESTORE_BY),
351 instance: v.config.get(model::KEY_RESTORE_INSTANCE).cloned(),
352 path: model::staged(&of, &stamp).path,
353 attached: !incus::instances_using(&v).is_empty(),
354 volume: v.name,
355 of,
356 stamp,
357 })
358 })
359 .collect();
360 out.sort_by(|a, b| b.stamp.cmp(&a.stamp));
361 Ok(out)
362 }
363
364 pub fn staged_discard(&self, org: &OrgId, volume: &str, stamp: &str) -> Result<Value> {
366 model::validate_volume_name(volume)?;
367 let staged = model::staged(volume, stamp);
368 model::validate_volume_name(&staged.volume)?;
369 let (oc, pool) = org_pool(self.client(), org)?;
370 let info = incus::volume(&oc, &pool, &staged.volume)?;
371 if info.config.get(model::KEY_RESTORE_OF).map(String::as_str) != Some(volume) {
372 return Err(Error::invalid(format!(
373 "{} is not a staged restore of {volume}; refusing to delete it",
374 staged.volume
375 )));
376 }
377 let users = incus::instances_using(&info);
378 for i in &users {
379 incus::detach(&oc, i, &staged.volume)?;
380 if incus::is_running(&oc, i).unwrap_or(false) {
383 let _ = crate::sandbox::Sandbox::get(&oc, i).and_then(|sb| {
384 sb.exec_with(
385 ["rmdir", staged.path.as_str()],
386 crate::exec::ExecOptions::default()
387 .user("0")
388 .timeout(std::time::Duration::from_secs(10)),
389 )
390 });
391 }
392 }
393 incus::delete_volume(&oc, &pool, &staged.volume)?;
394 self.event(
395 org,
396 volume,
397 "volume.restore.discarded",
398 "info",
399 format!("staged restore {} discarded", staged.volume),
400 );
401 Ok(json!({"ok": true, "volume": staged.volume, "detached_from": users}))
402 }
403}
404
405struct Job<'a> {
407 org: &'a OrgId,
408 volume: &'a str,
409 staged: &'a model::Staged,
410 from: &'a Origin,
411 instance: Option<&'a str>,
412 by: &'a str,
413 stamp: &'a str,
414}