1pub mod export;
24pub mod hook;
25pub mod incus;
26pub mod model;
27pub mod restore;
28mod schedule;
29
30use std::path::{Path, PathBuf};
31use std::sync::{Arc, Mutex};
32
33use serde_json::{Value, json};
34
35use crate::app::Apps;
36use crate::backup::Backups;
37use crate::client::Client;
38use crate::error::{Error, Result};
39use crate::jobs::scheduler::Scheduler;
40use crate::jobs::{Run, RunLog, RunStatus, RunStore, RunTrigger, Running};
41use crate::org::OrgId;
42pub use model::{VolumeRecord, VolumeSettings};
43
44struct Inner {
45 state: PathBuf,
46 apps: Apps,
47 backups: Backups,
48 running: Arc<Running>,
49 edit: Mutex<()>,
50 scheduler: Mutex<Scheduler>,
51}
52
53#[derive(Clone)]
55pub struct VolumeBackups {
56 inner: Arc<Inner>,
57}
58
59#[derive(Debug, Clone)]
61pub enum SnapshotName {
62 Auto,
64 Manual(Option<String>),
66}
67
68pub fn settings_path(state: &Path, org: &OrgId, volume: &str) -> PathBuf {
70 crate::app::org_root(state, org)
71 .join("volumes")
72 .join(volume)
73 .join("volume.json")
74}
75
76pub fn load_settings(state: &Path, org: &OrgId, volume: &str) -> VolumeSettings {
78 std::fs::read(settings_path(state, org, volume))
79 .ok()
80 .and_then(|b| serde_json::from_slice::<VolumeRecord>(&b).ok())
81 .map(|r| r.settings)
82 .unwrap_or_default()
83}
84
85pub fn org_pool(client: &Client, org: &OrgId) -> Result<(Client, String)> {
87 let oc = crate::org::client(client, org);
88 let pool = crate::sandbox::host_facts(&oc)?.pick_pool(None)?;
89 Ok((oc, pool))
90}
91
92const PROJECT_ALLOWS: [&str; 2] = ["restricted.snapshots", "restricted.backups"];
95
96pub fn allow_snapshots(client: &Client, org: &OrgId) -> Result<()> {
101 let project = org.incus_project();
102 let path = format!("/1.0/projects/{}", crate::client::encode_segment(&project));
103 let (p, etag) = match client.get_etag(&path) {
104 Ok(x) => x,
105 Err(e) if e.is_not_found() => return Ok(()),
106 Err(e) => return Err(e),
107 };
108 let mut cfg = p["config"].clone();
109 if cfg["restricted"].as_str() != Some("true")
110 || PROJECT_ALLOWS
111 .iter()
112 .all(|k| cfg[*k].as_str() == Some("allow"))
113 {
114 return Ok(());
115 }
116 for k in PROJECT_ALLOWS {
118 cfg[k] = json!("allow");
119 }
120 client
121 .mutate_if_match(
122 "PUT",
123 &path,
124 &json!({"description": p["description"], "config": cfg}),
125 etag.as_deref(),
126 &format!("allow snapshots in {project}"),
127 client.get_timeouts().other,
128 )
129 .map(|_| ())
130}
131
132fn now() -> i64 {
133 crate::stack::now_secs() as i64
134}
135
136impl VolumeBackups {
137 pub fn new(state: &Path, apps: Apps, backups: Backups) -> VolumeBackups {
138 let v = VolumeBackups {
139 inner: Arc::new(Inner {
140 state: state.to_path_buf(),
141 apps,
142 backups,
143 running: Arc::default(),
144 edit: Mutex::new(()),
145 scheduler: Mutex::new(Scheduler::idle()),
146 }),
147 };
148 for org in crate::jobs::orgs(state) {
149 for name in v.configured(&org) {
150 v.runs(&org, &name).recover();
151 }
152 }
153 v
154 }
155
156 pub fn set_scheduler(&self, s: Scheduler) {
157 *self.inner.scheduler.lock().unwrap() = s;
158 }
159
160 pub fn backups(&self) -> &Backups {
161 &self.inner.backups
162 }
163
164 fn client(&self) -> &Client {
165 self.inner.apps.client()
166 }
167
168 fn dir(&self, org: &OrgId, volume: &str) -> PathBuf {
169 crate::app::org_root(&self.inner.state, org)
170 .join("volumes")
171 .join(volume)
172 }
173
174 pub fn runs(&self, org: &OrgId, volume: &str) -> RunStore {
176 RunStore::new(self.dir(org, volume).join("runs"))
177 }
178
179 pub fn configured(&self, org: &OrgId) -> Vec<String> {
181 let root = crate::app::org_root(&self.inner.state, org).join("volumes");
182 let mut out: Vec<String> = std::fs::read_dir(root)
183 .into_iter()
184 .flatten()
185 .flatten()
186 .filter(|e| e.path().join("volume.json").exists())
187 .filter_map(|e| e.file_name().to_str().map(String::from))
188 .collect();
189 out.sort();
190 out
191 }
192
193 pub fn record(&self, org: &OrgId, volume: &str) -> Result<Option<VolumeRecord>> {
194 model::validate_volume_name(volume)?;
195 match std::fs::read(settings_path(&self.inner.state, org, volume)) {
196 Ok(b) => Ok(Some(serde_json::from_slice(&b)?)),
197 Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
198 Err(e) => Err(e.into()),
199 }
200 }
201
202 fn save(&self, org: &OrgId, volume: &str, r: &VolumeRecord) -> Result<()> {
203 crate::app::write_atomic(
204 &settings_path(&self.inner.state, org, volume),
205 &serde_json::to_vec_pretty(r)?,
206 )
207 }
208
209 pub fn info(&self, org: &OrgId, volume: &str) -> Result<crate::volume::VolumeInfo> {
211 model::validate_volume_name(volume)?;
212 let (oc, pool) = org_pool(self.client(), org)?;
213 incus::volume(&oc, &pool, volume)
214 }
215
216 pub fn update(&self, org: &OrgId, volume: &str, patch: &Value) -> Result<VolumeRecord> {
219 self.info(org, volume)?;
220 let _g = self.inner.edit.lock().unwrap();
221 let now = now();
222 let mut r = self.record(org, volume)?.unwrap_or(VolumeRecord {
223 settings: VolumeSettings::default(),
224 anchor: now,
225 created_at: now as u64,
226 updated_at: now as u64,
227 });
228 let mut v = serde_json::to_value(&r.settings)?;
229 crate::app::merge_patch(&mut v, patch);
230 let mut s: VolumeSettings = serde_json::from_value(v)
231 .map_err(|e| Error::invalid(format!("volume {volume}: {e}")))?;
232 if s.schedule.as_deref().is_some_and(|x| x.trim().is_empty()) {
233 s.schedule = None;
234 }
235 s.validate()?;
236 if s.schedule != r.settings.schedule
237 || s.timezone != r.settings.timezone
238 || (s.enabled && !r.settings.enabled)
239 {
240 r.anchor = now;
241 }
242 r.settings = s;
243 r.updated_at = now as u64;
244 self.save(org, volume, &r)?;
245 self.inner.scheduler.lock().unwrap().wake();
246 Ok(r)
247 }
248
249 pub fn next_run(&self, r: &VolumeRecord) -> Option<i64> {
250 if !r.settings.enabled {
251 return None;
252 }
253 r.settings.schedule().ok()??.next_after(r.anchor)
254 }
255
256 pub fn snapshots(&self, org: &OrgId, volume: &str) -> Result<Vec<incus::Snapshot>> {
258 model::validate_volume_name(volume)?;
259 let (oc, pool) = org_pool(self.client(), org)?;
260 let mut s = incus::snapshots(&oc, &pool, volume)?;
261 s.reverse();
262 Ok(s)
263 }
264
265 pub fn snapshot_delete(&self, org: &OrgId, volume: &str, snapshot: &str) -> Result<()> {
266 model::validate_volume_name(volume)?;
267 let (oc, pool) = org_pool(self.client(), org)?;
268 if !incus::snapshots(&oc, &pool, volume)?
269 .iter()
270 .any(|s| s.name == snapshot)
271 {
272 return Err(Error::NotFound(format!(
273 "snapshot {snapshot} of volume {volume}"
274 )));
275 }
276 incus::snapshot_delete(&oc, &pool, volume, snapshot)?;
277 self.event(
278 org,
279 volume,
280 "volume.snapshot.deleted",
281 "info",
282 format!("snapshot {volume}/{snapshot} deleted"),
283 );
284 Ok(())
285 }
286
287 pub fn create(&self, org: &OrgId, volume: &str, size: Option<&str>) -> Result<bool> {
291 model::validate_volume_name(volume)?;
292 let (oc, pool) = org_pool(self.client(), org)?;
293 let config = size
294 .map(|s| [("size".to_string(), s.to_string())].into())
295 .unwrap_or_default();
296 let created = crate::volume::ensure(&oc, &pool, volume, &config)?;
297 if created {
298 self.event(
299 org,
300 volume,
301 "volume.created",
302 "info",
303 format!("volume {volume} created"),
304 );
305 }
306 Ok(created)
307 }
308
309 pub fn delete(&self, org: &OrgId, volume: &str) -> Result<()> {
314 let info = self.info(org, volume)?;
315 if info.config.contains_key(model::KEY_RESTORE_OF) {
316 return Err(Error::invalid(format!(
317 "{volume} is a staged restore: volume_restore_discard removes it"
318 )));
319 }
320 if info.config.contains_key(model::KEY_TEMPORARY) {
321 return Err(Error::invalid(format!(
322 "{volume} is isb's own, for a backup or restore in progress"
323 )));
324 }
325 if let Some(b) = self
327 .backups()
328 .list(org)?
329 .into_iter()
330 .find(|b| b.spec.volume.as_deref() == Some(volume))
331 {
332 return Err(Error::invalid(format!(
333 "backup {} backs {volume} up; delete it first (backup_delete)",
334 b.spec.name
335 )));
336 }
337 let Some(_guard) = self.inner.running.enter(org, "volume", volume, true) else {
338 return Err(Error::invalid(format!(
339 "a snapshot of {volume} is being taken; try again when it is done"
340 )));
341 };
342 let (oc, pool) = org_pool(self.client(), org)?;
343 crate::volume::remove(&oc, &pool, volume)?;
344 let _g = self.inner.edit.lock().unwrap();
345 match std::fs::remove_dir_all(self.dir(org, volume)) {
346 Err(e) if e.kind() != std::io::ErrorKind::NotFound => {
347 eprintln!("isb serve: volume {volume}: settings not removed: {e}")
348 }
349 _ => {}
350 }
351 self.event(
352 org,
353 volume,
354 "volume.deleted",
355 "info",
356 format!("volume {volume} deleted"),
357 );
358 Ok(())
359 }
360
361 pub fn snapshot(
364 &self,
365 org: &OrgId,
366 volume: &str,
367 name: SnapshotName,
368 (trigger, by, slot): (RunTrigger, &str, Option<i64>),
369 ) -> Result<Option<Run>> {
370 let info = self.info(org, volume)?;
371 if let SnapshotName::Manual(Some(n)) = &name {
372 model::validate_snapshot_name(n)?;
373 }
374 let store = self.runs(org, volume);
375 let Some(guard) = self.inner.running.enter(org, "volume", volume, true) else {
376 if trigger != RunTrigger::Manual {
377 let (mut r, mut log) = store.start("snapshot", trigger, by, slot, 50)?;
378 r.error = Some("the previous snapshot of this volume was still running".into());
379 r.finish(RunStatus::Skipped);
380 store.finish(&mut r, &mut log)?;
381 }
382 return Ok(None);
383 };
384 let (r, log) = store.start("snapshot", trigger, by, slot, 50)?;
385 let me = self.clone();
386 let (org2, vol, run) = (org.clone(), volume.to_string(), r.clone());
387 std::thread::spawn(move || {
388 let _guard = guard;
389 me.snapshot_run(&org2, &vol, &info, name, run, log);
390 });
391 Ok(Some(r))
392 }
393
394 fn snapshot_run(
395 &self,
396 org: &OrgId,
397 volume: &str,
398 info: &crate::volume::VolumeInfo,
399 name: SnapshotName,
400 mut r: Run,
401 mut log: RunLog,
402 ) {
403 let res = self.snapshot_once(org, volume, info, &name, &mut log);
404 let (kind, level, msg) = match res {
405 Ok(detail) => {
406 let msg = format!(
407 "snapshot {volume}/{} taken",
408 detail["snapshot"].as_str().unwrap_or("")
409 );
410 r.detail = detail;
411 r.exit_code = Some(0);
412 r.finish(RunStatus::Succeeded);
413 ("volume.snapshot.created", "info", msg)
414 }
415 Err(e) => {
416 log.line(&format!("isb: {e}"));
417 r.error = Some(e.to_string());
418 r.finish(RunStatus::Failed);
419 (
420 "volume.snapshot.failed",
421 "error",
422 format!("snapshot of {volume} failed: {e}"),
423 )
424 }
425 };
426 if let Err(e) = self.runs(org, volume).finish(&mut r, &mut log) {
427 eprintln!("isb serve: volume {volume}: run {}: {e}", r.id);
428 }
429 self.event(org, volume, kind, level, msg);
430 }
431
432 fn snapshot_once(
433 &self,
434 org: &OrgId,
435 volume: &str,
436 info: &crate::volume::VolumeInfo,
437 name: &SnapshotName,
438 log: &mut RunLog,
439 ) -> Result<Value> {
440 let settings = load_settings(&self.inner.state, org, volume);
441 let (oc, pool) = org_pool(self.client(), org)?;
442 let st = model::stamp(now());
443 let snap = match name {
444 SnapshotName::Auto => format!("{}{st}", model::AUTO),
445 SnapshotName::Manual(Some(n)) => n.clone(),
446 SnapshotName::Manual(None) => format!("{}{st}", model::MANUAL),
447 };
448 allow_snapshots(self.client(), org)?;
449 let hooks = pre_snapshot(&oc, info, &settings, ("snapshot", &snap), log)?;
450 log.line(&format!("isb: snapshotting {volume} as {snap}"));
451 incus::snapshot_create(&oc, &pool, volume, &snap, "isb")?;
452 let mut pruned = Vec::new();
453 if matches!(name, SnapshotName::Auto) {
454 let names: Vec<String> = incus::snapshots(&oc, &pool, volume)?
455 .into_iter()
456 .map(|s| s.name)
457 .collect();
458 for old in model::select_snapshot_prune(&names, settings.keep as usize) {
459 match incus::snapshot_delete(&oc, &pool, volume, &old) {
460 Ok(()) => {
461 log.line(&format!("isb: pruned {old}"));
462 pruned.push(old);
463 }
464 Err(e) => log.line(&format!("isb: prune {old}: {e}")),
465 }
466 }
467 }
468 Ok(json!({"volume": volume, "snapshot": snap, "hooks": hooks, "pruned": pruned}))
469 }
470
471 pub(crate) fn event(&self, org: &OrgId, volume: &str, kind: &str, level: &str, msg: String) {
473 let q = crate::stack::qualified(org, volume);
474 self.inner
475 .apps
476 .controller()
477 .event(kind, level, &q, volume, msg);
478 }
479}
480
481pub fn pre_snapshot(
484 oc: &Client,
485 info: &crate::volume::VolumeInfo,
486 settings: &VolumeSettings,
487 what: (&str, &str),
488 log: &mut RunLog,
489) -> Result<Vec<hook::HookResult>> {
490 let mut running = Vec::new();
491 for i in incus::instances_using(info) {
492 match incus::is_running(oc, &i) {
493 Ok(true) => running.push(i),
494 Ok(false) => log.line(&format!("isb: {i} is stopped; no hook needed")),
495 Err(e) => log.line(&format!("isb: {i}: {e}")),
496 }
497 }
498 let env = [
499 ("ISB_VOLUME", info.name.as_str()),
500 ("ISB_REASON", what.0),
501 ("ISB_SNAPSHOT", what.1),
502 ];
503 hook::run_hooks(
504 &hook::IncusHooks(oc),
505 &running,
506 &env,
507 (settings.hook_timeout()?, settings.hook_required),
508 log,
509 )
510}