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 snapshot(
290 &self,
291 org: &OrgId,
292 volume: &str,
293 name: SnapshotName,
294 (trigger, by, slot): (RunTrigger, &str, Option<i64>),
295 ) -> Result<Option<Run>> {
296 let info = self.info(org, volume)?;
297 if let SnapshotName::Manual(Some(n)) = &name {
298 model::validate_snapshot_name(n)?;
299 }
300 let store = self.runs(org, volume);
301 let Some(guard) = self.inner.running.enter(org, "volume", volume, true) else {
302 if trigger != RunTrigger::Manual {
303 let (mut r, mut log) = store.start("snapshot", trigger, by, slot, 50)?;
304 r.error = Some("the previous snapshot of this volume was still running".into());
305 r.finish(RunStatus::Skipped);
306 store.finish(&mut r, &mut log)?;
307 }
308 return Ok(None);
309 };
310 let (r, log) = store.start("snapshot", trigger, by, slot, 50)?;
311 let me = self.clone();
312 let (org2, vol, run) = (org.clone(), volume.to_string(), r.clone());
313 std::thread::spawn(move || {
314 let _guard = guard;
315 me.snapshot_run(&org2, &vol, &info, name, run, log);
316 });
317 Ok(Some(r))
318 }
319
320 fn snapshot_run(
321 &self,
322 org: &OrgId,
323 volume: &str,
324 info: &crate::volume::VolumeInfo,
325 name: SnapshotName,
326 mut r: Run,
327 mut log: RunLog,
328 ) {
329 let res = self.snapshot_once(org, volume, info, &name, &mut log);
330 let (kind, level, msg) = match res {
331 Ok(detail) => {
332 let msg = format!(
333 "snapshot {volume}/{} taken",
334 detail["snapshot"].as_str().unwrap_or("")
335 );
336 r.detail = detail;
337 r.exit_code = Some(0);
338 r.finish(RunStatus::Succeeded);
339 ("volume.snapshot.created", "info", msg)
340 }
341 Err(e) => {
342 log.line(&format!("isb: {e}"));
343 r.error = Some(e.to_string());
344 r.finish(RunStatus::Failed);
345 (
346 "volume.snapshot.failed",
347 "error",
348 format!("snapshot of {volume} failed: {e}"),
349 )
350 }
351 };
352 if let Err(e) = self.runs(org, volume).finish(&mut r, &mut log) {
353 eprintln!("isb serve: volume {volume}: run {}: {e}", r.id);
354 }
355 self.event(org, volume, kind, level, msg);
356 }
357
358 fn snapshot_once(
359 &self,
360 org: &OrgId,
361 volume: &str,
362 info: &crate::volume::VolumeInfo,
363 name: &SnapshotName,
364 log: &mut RunLog,
365 ) -> Result<Value> {
366 let settings = load_settings(&self.inner.state, org, volume);
367 let (oc, pool) = org_pool(self.client(), org)?;
368 let st = model::stamp(now());
369 let snap = match name {
370 SnapshotName::Auto => format!("{}{st}", model::AUTO),
371 SnapshotName::Manual(Some(n)) => n.clone(),
372 SnapshotName::Manual(None) => format!("{}{st}", model::MANUAL),
373 };
374 allow_snapshots(self.client(), org)?;
375 let hooks = pre_snapshot(&oc, info, &settings, ("snapshot", &snap), log)?;
376 log.line(&format!("isb: snapshotting {volume} as {snap}"));
377 incus::snapshot_create(&oc, &pool, volume, &snap, "isb")?;
378 let mut pruned = Vec::new();
379 if matches!(name, SnapshotName::Auto) {
380 let names: Vec<String> = incus::snapshots(&oc, &pool, volume)?
381 .into_iter()
382 .map(|s| s.name)
383 .collect();
384 for old in model::select_snapshot_prune(&names, settings.keep as usize) {
385 match incus::snapshot_delete(&oc, &pool, volume, &old) {
386 Ok(()) => {
387 log.line(&format!("isb: pruned {old}"));
388 pruned.push(old);
389 }
390 Err(e) => log.line(&format!("isb: prune {old}: {e}")),
391 }
392 }
393 }
394 Ok(json!({"volume": volume, "snapshot": snap, "hooks": hooks, "pruned": pruned}))
395 }
396
397 pub(crate) fn event(&self, org: &OrgId, volume: &str, kind: &str, level: &str, msg: String) {
399 let q = crate::stack::qualified(org, volume);
400 self.inner
401 .apps
402 .controller()
403 .event(kind, level, &q, volume, msg);
404 }
405}
406
407pub fn pre_snapshot(
410 oc: &Client,
411 info: &crate::volume::VolumeInfo,
412 settings: &VolumeSettings,
413 what: (&str, &str),
414 log: &mut RunLog,
415) -> Result<Vec<hook::HookResult>> {
416 let mut running = Vec::new();
417 for i in incus::instances_using(info) {
418 match incus::is_running(oc, &i) {
419 Ok(true) => running.push(i),
420 Ok(false) => log.line(&format!("isb: {i} is stopped; no hook needed")),
421 Err(e) => log.line(&format!("isb: {i}: {e}")),
422 }
423 }
424 let env = [
425 ("ISB_VOLUME", info.name.as_str()),
426 ("ISB_REASON", what.0),
427 ("ISB_SNAPSHOT", what.1),
428 ];
429 hook::run_hooks(
430 &hook::IncusHooks(oc),
431 &running,
432 &env,
433 (settings.hook_timeout()?, settings.hook_required),
434 log,
435 )
436}