Skip to main content

isb_apps/
backup.rs

1//! Database backups to S3-compatible storage, and restores from them.
2//!
3//! A **destination** is a bucket (endpoint, region, bucket, prefix,
4//! path-style addressing) whose access key pair is kept as org secrets. A
5//! **backup** is a schedule for one database: cron, destination, how many
6//! to keep, compression. Each run:
7//!
8//! 1. runs the engine's native dump inside the database's own instance
9//!    (`pg_dump`, `mysqldump`, `mariadb-dump`, `mongodump`, a Redis RDB),
10//! 2. streams its output through gzip or zstd and up to the bucket (one
11//!    part buffered at a time, multipart beyond [`crate::s3::PART_SIZE`];
12//!    nothing touches the host's disk, and the database's org needs no
13//!    network path to the bucket: the daemon carries the bytes),
14//! 3. checks the object with a `HEAD`, then deletes the oldest beyond
15//!    `keep`, and emits `backup.succeeded` or `backup.failed`.
16//!
17//! Objects: `<prefix>/<org>/<backup>/<database>-<YYYYMMDDTHHMMSSZ>.<engine>[.gz|.zst]`.
18//!
19//! A **restore** streams an object back the other way, into the same
20//! database or a new one created for it.
21
22use std::io::{Read, Write};
23use std::path::{Path, PathBuf};
24use std::sync::mpsc::{Receiver, SyncSender, sync_channel};
25use std::sync::{Arc, Mutex};
26use std::time::Duration;
27
28use serde::{Deserialize, Serialize};
29use serde_json::{Value, json};
30
31use crate::app::{Apps, DatabaseSource, Engine, Source};
32use crate::cron::Schedule;
33use crate::error::{Error, Result};
34use crate::exec::{ExecEvent, ExecOptions, Stdin};
35use crate::jobs::scheduler::{Entry, Scheduled, Scheduler};
36use crate::jobs::{Run, RunLog, RunStatus, RunStore, RunTrigger, Running};
37use crate::org::OrgId;
38use crate::sandbox::Sandbox;
39
40/// How long a dump or a restore may take.
41const DUMP_TIMEOUT: Duration = Duration::from_secs(12 * 3600);
42
43// --- compression -----------------------------------------------------------
44
45#[derive(Debug, Clone, Copy, Default, PartialEq, Eq, Serialize, Deserialize)]
46#[serde(rename_all = "lowercase")]
47pub enum Compression {
48    #[default]
49    Gzip,
50    Zstd,
51    None,
52}
53
54impl Compression {
55    pub fn extension(self) -> &'static str {
56        match self {
57            Compression::Gzip => ".gz",
58            Compression::Zstd => ".zst",
59            Compression::None => "",
60        }
61    }
62
63    pub fn content_type(self) -> &'static str {
64        match self {
65            Compression::Gzip => "application/gzip",
66            Compression::Zstd => "application/zstd",
67            Compression::None => "application/octet-stream",
68        }
69    }
70
71    fn from_suffix(name: &str) -> Compression {
72        if name.ends_with(".gz") {
73            Compression::Gzip
74        } else if name.ends_with(".zst") {
75            Compression::Zstd
76        } else {
77            Compression::None
78        }
79    }
80}
81
82/// A writer that compresses into `W`, finished explicitly.
83pub trait Compressor<W>: Write + Send {
84    fn finish(self: Box<Self>) -> std::io::Result<W>;
85}
86
87struct Plain<W>(W);
88
89impl<W: Write + Send> Write for Plain<W> {
90    fn write(&mut self, b: &[u8]) -> std::io::Result<usize> {
91        self.0.write(b)
92    }
93    fn flush(&mut self) -> std::io::Result<()> {
94        self.0.flush()
95    }
96}
97
98impl<W: Write + Send> Compressor<W> for Plain<W> {
99    fn finish(self: Box<Self>) -> std::io::Result<W> {
100        Ok(self.0)
101    }
102}
103
104impl<W: Write + Send> Compressor<W> for flate2::write::GzEncoder<W> {
105    fn finish(self: Box<Self>) -> std::io::Result<W> {
106        flate2::write::GzEncoder::finish(*self)
107    }
108}
109
110/// ruzstd's encoder pulls from a reader, so it runs on its own thread fed
111/// through a bounded channel. Its sink never fails it (the first error is
112/// kept and reported at finish), so it never panics on a failed upload.
113struct Zstd<W> {
114    tx: Option<SyncSender<Vec<u8>>>,
115    failed: Arc<Mutex<Option<String>>>,
116    thread: Option<std::thread::JoinHandle<W>>,
117}
118
119struct ChanReader {
120    rx: Receiver<Vec<u8>>,
121    buf: Vec<u8>,
122    pos: usize,
123}
124
125impl Read for ChanReader {
126    fn read(&mut self, out: &mut [u8]) -> std::io::Result<usize> {
127        while self.pos == self.buf.len() {
128            match self.rx.recv() {
129                Ok(b) => {
130                    self.buf = b;
131                    self.pos = 0;
132                }
133                Err(_) => return Ok(0),
134            }
135        }
136        let n = out.len().min(self.buf.len() - self.pos);
137        out[..n].copy_from_slice(&self.buf[self.pos..self.pos + n]);
138        self.pos += n;
139        Ok(n)
140    }
141}
142
143struct Sink<W> {
144    w: W,
145    failed: Arc<Mutex<Option<String>>>,
146}
147
148impl<W: Write> Write for Sink<W> {
149    fn write(&mut self, b: &[u8]) -> std::io::Result<usize> {
150        let mut f = self.failed.lock().unwrap();
151        if f.is_none() {
152            if let Err(e) = self.w.write_all(b) {
153                *f = Some(e.to_string());
154            }
155        }
156        Ok(b.len())
157    }
158    fn flush(&mut self) -> std::io::Result<()> {
159        Ok(())
160    }
161}
162
163impl<W: Write + Send + 'static> Zstd<W> {
164    fn new(w: W) -> Zstd<W> {
165        let (tx, rx) = sync_channel::<Vec<u8>>(8);
166        let failed = Arc::new(Mutex::new(None));
167        let f2 = failed.clone();
168        let thread = std::thread::spawn(move || {
169            let mut sink = Sink { w, failed: f2 };
170            ruzstd::encoding::compress(
171                ChanReader {
172                    rx,
173                    buf: Vec::new(),
174                    pos: 0,
175                },
176                &mut sink,
177                ruzstd::encoding::CompressionLevel::Fastest,
178            );
179            sink.w
180        });
181        Zstd {
182            tx: Some(tx),
183            failed,
184            thread: Some(thread),
185        }
186    }
187}
188
189impl<W: Write + Send + 'static> Write for Zstd<W> {
190    fn write(&mut self, b: &[u8]) -> std::io::Result<usize> {
191        if let Some(e) = self.failed.lock().unwrap().clone() {
192            return Err(std::io::Error::other(e));
193        }
194        let tx = self
195            .tx
196            .as_ref()
197            .ok_or_else(|| std::io::Error::other("finished"))?;
198        tx.send(b.to_vec())
199            .map_err(|_| std::io::Error::other("the zstd encoder stopped"))?;
200        Ok(b.len())
201    }
202    fn flush(&mut self) -> std::io::Result<()> {
203        Ok(())
204    }
205}
206
207impl<W: Write + Send + 'static> Compressor<W> for Zstd<W> {
208    fn finish(mut self: Box<Self>) -> std::io::Result<W> {
209        drop(self.tx.take());
210        let w = self
211            .thread
212            .take()
213            .ok_or_else(|| std::io::Error::other("finished twice"))?
214            .join()
215            .map_err(|_| std::io::Error::other("the zstd encoder panicked"))?;
216        if let Some(e) = self.failed.lock().unwrap().take() {
217            return Err(std::io::Error::other(e));
218        }
219        Ok(w)
220    }
221}
222
223/// Compress into `w`.
224pub fn compressor<W: Write + Send + 'static>(c: Compression, w: W) -> Box<dyn Compressor<W>> {
225    match c {
226        Compression::Gzip => Box::new(flate2::write::GzEncoder::new(
227            w,
228            flate2::Compression::default(),
229        )),
230        Compression::Zstd => Box::new(Zstd::new(w)),
231        Compression::None => Box::new(Plain(w)),
232    }
233}
234
235/// Decompress what `r` reads.
236pub fn decompressor<R: Read + Send + 'static>(
237    c: Compression,
238    r: R,
239) -> Result<Box<dyn Read + Send>> {
240    Ok(match c {
241        Compression::Gzip => Box::new(flate2::read::MultiGzDecoder::new(r)),
242        Compression::Zstd => Box::new(
243            ruzstd::decoding::StreamingDecoder::new(r)
244                .map_err(|e| Error::invalid(format!("zstd: {e}")))?,
245        ),
246        Compression::None => Box::new(r),
247    })
248}
249
250// --- records ---------------------------------------------------------------
251
252fn default_region() -> String {
253    "us-east-1".into()
254}
255
256/// An S3-compatible bucket backups go to.
257#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
258#[serde(deny_unknown_fields)]
259pub struct Destination {
260    pub name: String,
261    /// `https://s3.eu-central-1.amazonaws.com`, `https://<account>.r2.cloudflarestorage.com`,
262    /// `http://minio.internal:9000`.
263    pub endpoint: String,
264    #[serde(default = "default_region")]
265    pub region: String,
266    pub bucket: String,
267    /// Keys start with it (`isb/backups`).
268    #[serde(default, skip_serializing_if = "String::is_empty")]
269    pub prefix: String,
270    /// `endpoint/bucket/key` (MinIO, most self-hosted stores) rather than
271    /// `bucket.endpoint/key`.
272    #[serde(default)]
273    pub path_style: bool,
274    /// Org secrets holding the key pair.
275    pub access_key_secret: String,
276    pub secret_key_secret: String,
277    /// Loopback endpoints are refused unless the destination was made by
278    /// the local CLI or a platform admin.
279    #[serde(default, skip_serializing_if = "std::ops::Not::not")]
280    pub allow_local: bool,
281    #[serde(default)]
282    pub created_at: u64,
283}
284
285fn default_keep() -> u32 {
286    7
287}
288
289fn yes() -> bool {
290    true
291}
292
293/// What a user sets on a backup.
294#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
295#[serde(deny_unknown_fields)]
296pub struct BackupSpec {
297    pub name: String,
298    /// The database app, or empty for a volume backup.
299    #[serde(default, skip_serializing_if = "String::is_empty")]
300    pub database: String,
301    /// Or a named volume in the org ([`crate::volume_backup`]).
302    #[serde(default, skip_serializing_if = "Option::is_none")]
303    pub volume: Option<String>,
304    pub destination: String,
305    /// Cron (five fields or an alias).
306    pub schedule: String,
307    #[serde(default, skip_serializing_if = "Option::is_none")]
308    pub timezone: Option<String>,
309    /// Backups kept in the bucket (older ones are deleted).
310    #[serde(default = "default_keep")]
311    pub keep: u32,
312    #[serde(default)]
313    pub compression: Compression,
314    #[serde(default = "yes")]
315    pub enabled: bool,
316    #[serde(default, skip_serializing_if = "Option::is_none")]
317    pub missed_grace: Option<String>,
318}
319
320#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
321pub struct Backup {
322    pub spec: BackupSpec,
323    pub created_at: u64,
324    pub updated_at: u64,
325    pub anchor: i64,
326}
327
328impl BackupSpec {
329    pub fn schedule(&self) -> Result<Schedule> {
330        let off = crate::cron::parse_offset(self.timezone.as_deref().unwrap_or("UTC"))?;
331        Schedule::parse_in(&self.schedule, off)
332    }
333
334    pub fn validate(&self) -> Result<()> {
335        crate::jobs::validate_name("backup", &self.name)?;
336        match &self.volume {
337            None => crate::app::validate_app_name(&self.database)?,
338            Some(v) if self.database.is_empty() => {
339                crate::volume_backup::model::validate_volume_name(v)?
340            }
341            Some(_) => {
342                return Err(Error::invalid(
343                    "back up a `database` or a `volume`, not both",
344                ));
345            }
346        }
347        crate::jobs::validate_name("destination", &self.destination)?;
348        self.schedule()?;
349        crate::jobs::parse_grace(&self.missed_grace)?;
350        if self.keep == 0 || self.keep > 1000 {
351            return Err(Error::invalid("keep: 1 to 1000 backups"));
352        }
353        Ok(())
354    }
355}
356
357/// The key prefix a backup's objects live under (ends with `/`).
358pub fn backup_prefix(dest: &Destination, org: &OrgId, backup: &str) -> String {
359    let p = dest.prefix.trim_matches('/');
360    if p.is_empty() {
361        format!("{org}/{backup}/")
362    } else {
363        format!("{p}/{org}/{backup}/")
364    }
365}
366
367/// An object key for a backup taken at `t`.
368pub fn object_key(prefix: &str, database: &str, t: i64, engine: Engine, c: Compression) -> String {
369    format!(
370        "{prefix}{database}-{}.{}{}",
371        crate::cron::compact_utc(t),
372        engine.name(),
373        c.extension()
374    )
375}
376
377/// A key this module wrote: its time, engine and compression.
378pub fn parse_key(prefix: &str, key: &str) -> Option<(String, i64, Engine, Compression)> {
379    let file = key.strip_prefix(prefix)?;
380    if file.contains('/') {
381        return None;
382    }
383    // <database>-<16-char timestamp>.<engine>[.gz|.zst]
384    let dot = file.find('.')?;
385    let (stem, rest) = file.split_at(dot);
386    if stem.len() < 18 {
387        return None;
388    }
389    let (db, ts) = stem.split_at(stem.len() - 16);
390    let db = db.strip_suffix('-')?;
391    let t = crate::cron::parse_compact_utc(ts)?;
392    let rest = &rest[1..];
393    let engine = rest.split('.').next()?;
394    let engine = Engine::parse(engine).ok()?;
395    let c = Compression::from_suffix(rest);
396    if rest != format!("{}{}", engine.name(), c.extension()) {
397        return None;
398    }
399    Some((db.to_string(), t, engine, c))
400}
401
402/// Which objects retention deletes: everything this module wrote under
403/// `prefix` beyond the newest `keep`. Keys it did not write are never
404/// touched.
405pub fn select_prune(prefix: &str, keys: &[crate::s3::Object], keep: usize) -> Vec<String> {
406    let mut ours: Vec<(i64, &str)> = keys
407        .iter()
408        .filter_map(|o| parse_key(prefix, &o.key).map(|(_, t, _, _)| (t, o.key.as_str())))
409        .collect();
410    ours.sort_by(|a, b| b.cmp(a));
411    ours.into_iter()
412        .skip(keep)
413        .map(|(_, k)| k.to_string())
414        .collect()
415}
416
417/// A backup object as listed.
418#[derive(Debug, Clone, Serialize)]
419pub struct BackupFile {
420    pub key: String,
421    pub size: u64,
422    pub taken_at: String,
423    #[serde(skip_serializing_if = "Option::is_none")]
424    pub engine: Option<Engine>,
425    /// A volume backup's volume.
426    #[serde(skip_serializing_if = "Option::is_none")]
427    pub volume: Option<String>,
428    pub compression: Compression,
429}
430
431/// What a restore takes from.
432#[derive(Debug, Clone, Default, Deserialize)]
433#[serde(deny_unknown_fields)]
434pub struct RestoreRequest {
435    /// The backup whose objects to restore from (its destination and
436    /// prefix), or `destination` and `key`.
437    #[serde(default)]
438    pub backup: Option<String>,
439    #[serde(default)]
440    pub destination: Option<String>,
441    /// The object (default: the backup's newest).
442    #[serde(default)]
443    pub key: Option<String>,
444    /// An existing database to restore into (its data is replaced).
445    #[serde(default)]
446    pub target: Option<String>,
447    /// Or a new database to create for it.
448    #[serde(default)]
449    pub new: Option<NewDatabase>,
450    /// Required to restore over an existing database.
451    #[serde(default)]
452    pub confirm: bool,
453}
454
455#[derive(Debug, Clone, Default, Deserialize)]
456#[serde(deny_unknown_fields)]
457pub struct NewDatabase {
458    pub name: String,
459    /// Default: the backed-up database's project and environment.
460    #[serde(default)]
461    pub project: Option<String>,
462    #[serde(default)]
463    pub environment: Option<String>,
464    /// Default: the backed-up database's version.
465    #[serde(default)]
466    pub version: Option<String>,
467}
468
469// --- the service -----------------------------------------------------------
470
471struct Inner {
472    state: PathBuf,
473    apps: Apps,
474    running: Arc<Running>,
475    edit: Mutex<()>,
476    scheduler: Mutex<Scheduler>,
477}
478
479/// Every org's destinations and backups.
480#[derive(Clone)]
481pub struct Backups {
482    inner: Arc<Inner>,
483}
484
485/// Whether `host` (a URL's host, maybe with a port) resolves to loopback,
486/// unspecified or link-local addresses: places on the daemon's own host.
487fn is_local_endpoint(endpoint: &str) -> bool {
488    use std::net::ToSocketAddrs;
489    let rest = endpoint
490        .split_once("://")
491        .map(|(_, r)| r)
492        .unwrap_or(endpoint);
493    let hostport = rest.trim_end_matches('/');
494    let with_port = if hostport
495        .rsplit_once(':')
496        .is_some_and(|(_, p)| p.parse::<u16>().is_ok())
497        && !hostport.ends_with(']')
498    {
499        hostport.to_string()
500    } else {
501        format!("{hostport}:443")
502    };
503    let Ok(addrs) = with_port.to_socket_addrs() else {
504        return false;
505    };
506    addrs.into_iter().any(|a| {
507        let ip = a.ip();
508        ip.is_loopback()
509            || ip.is_unspecified()
510            || match ip {
511                std::net::IpAddr::V4(v) => v.is_link_local(),
512                std::net::IpAddr::V6(v) => (v.segments()[0] & 0xffc0) == 0xfe80,
513            }
514    })
515}
516
517impl Backups {
518    pub fn new(state: &Path, apps: Apps) -> Backups {
519        let b = Backups {
520            inner: Arc::new(Inner {
521                state: state.to_path_buf(),
522                apps,
523                running: Arc::default(),
524                edit: Mutex::new(()),
525                scheduler: Mutex::new(Scheduler::idle()),
526            }),
527        };
528        for org in crate::jobs::orgs(state) {
529            for bk in b.list(&org).unwrap_or_default() {
530                b.runs(&org, &bk.spec.name).recover();
531            }
532            b.restore_runs(&org).recover();
533        }
534        b
535    }
536
537    pub fn set_scheduler(&self, s: Scheduler) {
538        *self.inner.scheduler.lock().unwrap() = s;
539    }
540
541    pub fn apps(&self) -> &Apps {
542        &self.inner.apps
543    }
544
545    pub fn state_dir(&self) -> &Path {
546        &self.inner.state
547    }
548
549    fn root(&self, org: &OrgId) -> PathBuf {
550        crate::app::org_root(&self.inner.state, org).join("backups")
551    }
552
553    fn dest_path(&self, org: &OrgId, name: &str) -> PathBuf {
554        self.root(org)
555            .join("destinations")
556            .join(format!("{name}.json"))
557    }
558
559    fn backup_path(&self, org: &OrgId, name: &str) -> PathBuf {
560        self.root(org)
561            .join("schedules")
562            .join(name)
563            .join("backup.json")
564    }
565
566    pub fn runs(&self, org: &OrgId, name: &str) -> RunStore {
567        RunStore::new(self.root(org).join("schedules").join(name).join("runs"))
568    }
569
570    pub fn restore_runs(&self, org: &OrgId) -> RunStore {
571        RunStore::new(self.root(org).join("restores").join("runs"))
572    }
573
574    // --- destinations ------------------------------------------------------
575
576    /// Create a destination. `access_key`/`secret_key` values are stored
577    /// as the org secrets `backup.<name>.access-key`/`.secret-key`;
578    /// otherwise `dest` names existing secrets. `trusted`: the caller may
579    /// point it at this host (loopback).
580    pub fn destination_create(
581        &self,
582        org: &OrgId,
583        mut dest: Destination,
584        access_key: Option<String>,
585        secret_key: Option<String>,
586        trusted: bool,
587    ) -> Result<Destination> {
588        crate::jobs::validate_name("destination", &dest.name)?;
589        if is_local_endpoint(&dest.endpoint) {
590            if !trusted {
591                return Err(Error::Forbidden(format!(
592                    "endpoint {}: loopback and link-local endpoints are for the local CLI and platform admins",
593                    dest.endpoint
594                )));
595            }
596            dest.allow_local = true;
597        } else {
598            dest.allow_local = false;
599        }
600        let _g = self.inner.edit.lock().unwrap();
601        let p = self.dest_path(org, &dest.name);
602        if p.exists() {
603            return Err(Error::AlreadyExists(format!("destination {}", dest.name)));
604        }
605        let secrets = self.inner.apps.secrets();
606        if let Some(a) = access_key {
607            dest.access_key_secret = format!("backup.{}.access-key", dest.name);
608            secrets.set(org, &dest.access_key_secret, a.trim().as_bytes())?;
609        }
610        if let Some(s) = secret_key {
611            dest.secret_key_secret = format!("backup.{}.secret-key", dest.name);
612            secrets.set(org, &dest.secret_key_secret, s.trim().as_bytes())?;
613        }
614        for n in [&dest.access_key_secret, &dest.secret_key_secret] {
615            if n.is_empty() {
616                return Err(Error::invalid(
617                    "give the access key pair (access_key and secret_key), or the secrets holding it",
618                ));
619            }
620            secrets.inspect(org, n).map_err(|e| match e {
621                Error::NotFound(_) => {
622                    Error::invalid(format!("secret {n} does not exist in org {org}"))
623                }
624                e => e,
625            })?;
626        }
627        // Checks the endpoint and bucket without a request.
628        self.client(org, &dest)?;
629        dest.created_at = crate::stack::now_secs();
630        crate::app::write_atomic(&p, &serde_json::to_vec_pretty(&dest)?)?;
631        Ok(dest)
632    }
633
634    pub fn destination_get(&self, org: &OrgId, name: &str) -> Result<Destination> {
635        crate::jobs::validate_name("destination", name)?;
636        match std::fs::read(self.dest_path(org, name)) {
637            Ok(b) => Ok(serde_json::from_slice(&b)?),
638            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
639                Err(Error::NotFound(format!("destination {name} in org {org}")))
640            }
641            Err(e) => Err(e.into()),
642        }
643    }
644
645    pub fn destination_list(&self, org: &OrgId) -> Result<Vec<Destination>> {
646        let mut out = Vec::new();
647        let Ok(rd) = std::fs::read_dir(self.root(org).join("destinations")) else {
648            return Ok(out);
649        };
650        for e in rd.flatten() {
651            let p = e.path();
652            if p.extension().is_some_and(|x| x == "json") {
653                if let Ok(d) = serde_json::from_slice::<Destination>(&std::fs::read(&p)?) {
654                    out.push(d);
655                }
656            }
657        }
658        out.sort_by(|a, b| a.name.cmp(&b.name));
659        Ok(out)
660    }
661
662    /// Delete a destination no backup uses, and the key secrets isb made.
663    /// Objects in the bucket stay.
664    pub fn destination_delete(&self, org: &OrgId, name: &str) -> Result<()> {
665        let _g = self.inner.edit.lock().unwrap();
666        let d = self.destination_get(org, name)?;
667        let users: Vec<String> = self
668            .list(org)?
669            .into_iter()
670            .filter(|b| b.spec.destination == name)
671            .map(|b| b.spec.name)
672            .collect();
673        if !users.is_empty() {
674            return Err(Error::invalid(format!(
675                "destination {name} is used by backups {}; delete them first",
676                users.join(", ")
677            )));
678        }
679        for s in [&d.access_key_secret, &d.secret_key_secret] {
680            if s.starts_with(&format!("backup.{name}.")) {
681                match self.inner.apps.secrets().delete(org, s) {
682                    Ok(()) => {}
683                    Err(e) if e.is_not_found() => {}
684                    Err(e) => return Err(e),
685                }
686            }
687        }
688        std::fs::remove_file(self.dest_path(org, name))?;
689        Ok(())
690    }
691
692    /// An S3 client for a destination, with its key pair read now.
693    pub fn client(&self, org: &OrgId, d: &Destination) -> Result<crate::s3::Client> {
694        if !d.allow_local && is_local_endpoint(&d.endpoint) {
695            return Err(Error::Forbidden(format!(
696                "destination {}: endpoint {} is on this host; only the local CLI or a platform admin may create such a destination",
697                d.name, d.endpoint
698            )));
699        }
700        let read = |n: &str| -> Result<String> {
701            let (v, _) =
702                self.inner.apps.secrets().get(org, n).map_err(|e| {
703                    Error::invalid(format!("destination {}: secret {n}: {e}", d.name))
704                })?;
705            String::from_utf8(v)
706                .map(|s| s.trim().to_string())
707                .map_err(|_| Error::invalid(format!("secret {n} is not text")))
708        };
709        crate::s3::Client::new(crate::s3::Bucket {
710            endpoint: d.endpoint.clone(),
711            region: d.region.clone(),
712            bucket: d.bucket.clone(),
713            path_style: d.path_style,
714            creds: crate::s3::Credentials {
715                access_key: read(&d.access_key_secret)?,
716                secret_key: read(&d.secret_key_secret)?,
717            },
718        })
719    }
720
721    /// Create the destination's bucket (self-hosted stores; S3 itself
722    /// usually wants buckets made in its console).
723    pub fn destination_create_bucket(&self, org: &OrgId, name: &str) -> Result<()> {
724        let d = self.destination_get(org, name)?;
725        self.client(org, &d)?.create_bucket()
726    }
727
728    /// Write, read back and delete a small object under the destination's
729    /// prefix: the key pair, the bucket and the network all work.
730    pub fn destination_test(&self, org: &OrgId, name: &str) -> Result<Value> {
731        let d = self.destination_get(org, name)?;
732        let c = self.client(org, &d)?;
733        let started = std::time::Instant::now();
734        let p = d.prefix.trim_matches('/');
735        let test = format!(".isb-test-{}", crate::app::git::random_hex(6));
736        let key = if p.is_empty() {
737            format!("{org}/{test}")
738        } else {
739            format!("{p}/{org}/{test}")
740        };
741        let body = b"isb destination test\n";
742        c.put(&key, body)?;
743        let size = c.head(&key)?;
744        c.delete(&key)?;
745        if size != Some(body.len() as u64) {
746            return Err(Error::invalid(format!(
747                "destination {name}: wrote {} bytes, HEAD reported {size:?}",
748                body.len()
749            )));
750        }
751        Ok(json!({"ok": true, "key": key, "ms": started.elapsed().as_millis() as u64}))
752    }
753
754    // --- backups -----------------------------------------------------------
755
756    pub fn get(&self, org: &OrgId, name: &str) -> Result<Backup> {
757        crate::jobs::validate_name("backup", name)?;
758        match std::fs::read(self.backup_path(org, name)) {
759            Ok(b) => Ok(serde_json::from_slice(&b)?),
760            Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
761                Err(Error::NotFound(format!("backup {name} in org {org}")))
762            }
763            Err(e) => Err(e.into()),
764        }
765    }
766
767    pub fn list(&self, org: &OrgId) -> Result<Vec<Backup>> {
768        let mut out = Vec::new();
769        let Ok(rd) = std::fs::read_dir(self.root(org).join("schedules")) else {
770            return Ok(out);
771        };
772        for e in rd.flatten() {
773            let p = e.path().join("backup.json");
774            if let Ok(b) = std::fs::read(&p) {
775                match serde_json::from_slice::<Backup>(&b) {
776                    Ok(b) => out.push(b),
777                    Err(e) => eprintln!("isb serve: skipping {}: {e}", p.display()),
778                }
779            }
780        }
781        out.sort_by(|a, b| a.spec.name.cmp(&b.spec.name));
782        Ok(out)
783    }
784
785    fn save(&self, org: &OrgId, b: &Backup) -> Result<()> {
786        crate::app::write_atomic(
787            &self.backup_path(org, &b.spec.name),
788            &serde_json::to_vec_pretty(b)?,
789        )
790    }
791
792    /// The database app and its source.
793    fn database(&self, org: &OrgId, name: &str) -> Result<(crate::app::App, DatabaseSource)> {
794        let app = self.inner.apps.get(org, name)?;
795        match &app.spec.source {
796            Source::Database(db) => {
797                let db = db.clone();
798                Ok((app, db))
799            }
800            _ => Err(Error::invalid(format!("app {name} is not a database"))),
801        }
802    }
803
804    pub fn create(&self, org: &OrgId, spec: BackupSpec) -> Result<Backup> {
805        spec.validate()?;
806        match &spec.volume {
807            Some(v) => crate::volume_backup::export::check_source(self, org, v).map(|_| ())?,
808            None => self.database(org, &spec.database).map(|_| ())?,
809        }
810        self.destination_get(org, &spec.destination)?;
811        let _g = self.inner.edit.lock().unwrap();
812        if self.backup_path(org, &spec.name).exists() {
813            return Err(Error::AlreadyExists(format!("backup {}", spec.name)));
814        }
815        let now = crate::stack::now_secs();
816        let b = Backup {
817            spec,
818            created_at: now,
819            updated_at: now,
820            anchor: now as i64,
821        };
822        self.save(org, &b)?;
823        self.inner.scheduler.lock().unwrap().wake();
824        Ok(b)
825    }
826
827    /// A merge patch of a backup's settings (name and database fixed).
828    pub fn update(&self, org: &OrgId, name: &str, patch: &Value) -> Result<Backup> {
829        let _g = self.inner.edit.lock().unwrap();
830        let mut b = self.get(org, name)?;
831        let mut v = serde_json::to_value(&b.spec)?;
832        crate::app::merge_patch(&mut v, patch);
833        let spec: BackupSpec =
834            serde_json::from_value(v).map_err(|e| Error::invalid(format!("backup {name}: {e}")))?;
835        if spec.name != b.spec.name
836            || spec.database != b.spec.database
837            || spec.volume != b.spec.volume
838        {
839            return Err(Error::invalid(
840                "a backup's name and what it backs up are fixed",
841            ));
842        }
843        spec.validate()?;
844        self.destination_get(org, &spec.destination)?;
845        let now = crate::stack::now_secs();
846        if spec.schedule != b.spec.schedule
847            || spec.timezone != b.spec.timezone
848            || (spec.enabled && !b.spec.enabled)
849        {
850            b.anchor = now as i64;
851        }
852        b.spec = spec;
853        b.updated_at = now;
854        self.save(org, &b)?;
855        self.inner.scheduler.lock().unwrap().wake();
856        Ok(b)
857    }
858
859    /// Delete a backup schedule and its run records; its objects stay in
860    /// the bucket (restore them with `destination` and `key`).
861    pub fn delete(&self, org: &OrgId, name: &str) -> Result<()> {
862        let _g = self.inner.edit.lock().unwrap();
863        self.get(org, name)?;
864        if self.inner.running.is_running(org, "backup", name) {
865            return Err(Error::invalid(format!("backup {name} is running")));
866        }
867        std::fs::remove_dir_all(self.root(org).join("schedules").join(name))?;
868        Ok(())
869    }
870
871    pub fn next_run(&self, b: &Backup) -> Option<i64> {
872        if !b.spec.enabled {
873            return None;
874        }
875        b.spec.schedule().ok()?.next_after(b.anchor)
876    }
877
878    /// The backup's objects in its bucket, newest first.
879    pub fn files(&self, org: &OrgId, name: &str) -> Result<Vec<BackupFile>> {
880        let b = self.get(org, name)?;
881        let d = self.destination_get(org, &b.spec.destination)?;
882        let prefix = backup_prefix(&d, org, name);
883        if b.spec.volume.is_some() {
884            return crate::volume_backup::export::files(self, org, &d, &prefix);
885        }
886        self.files_at(org, &d, &prefix)
887    }
888
889    fn files_at(&self, org: &OrgId, d: &Destination, prefix: &str) -> Result<Vec<BackupFile>> {
890        let c = self.client(org, d)?;
891        let mut out: Vec<(i64, BackupFile)> = c
892            .list(prefix)?
893            .into_iter()
894            .filter_map(|o| {
895                let (_, t, engine, compression) = parse_key(prefix, &o.key)?;
896                Some((
897                    t,
898                    BackupFile {
899                        key: o.key,
900                        size: o.size,
901                        taken_at: crate::cron::rfc3339(t),
902                        engine: Some(engine),
903                        volume: None,
904                        compression,
905                    },
906                ))
907            })
908            .collect();
909        out.sort_by_key(|a| std::cmp::Reverse(a.0));
910        Ok(out.into_iter().map(|(_, f)| f).collect())
911    }
912
913    /// Back up now, in the background.
914    pub fn run_now(&self, org: &OrgId, name: &str, by: &str) -> Result<Run> {
915        let b = self.get(org, name)?;
916        self.start(org, b, RunTrigger::Manual, by, None)?
917            .ok_or_else(|| Error::invalid(format!("backup {name} is already running")))
918    }
919
920    fn start(
921        &self,
922        org: &OrgId,
923        b: Backup,
924        trigger: RunTrigger,
925        by: &str,
926        slot: Option<i64>,
927    ) -> Result<Option<Run>> {
928        let name = b.spec.name.clone();
929        let store = self.runs(org, &name);
930        let Some(guard) = self.inner.running.enter(org, "backup", &name, true) else {
931            if trigger != RunTrigger::Manual {
932                let (mut r, mut log) =
933                    store.start("backup", trigger, by, slot, b.spec.keep.max(10) as usize)?;
934                r.error = Some("the previous backup was still running".into());
935                r.finish(RunStatus::Skipped);
936                store.finish(&mut r, &mut log)?;
937            }
938            return Ok(None);
939        };
940        // Records: at least as many as objects kept, at least ten.
941        let (r, log) = store.start("backup", trigger, by, slot, b.spec.keep.max(10) as usize)?;
942        let me = self.clone();
943        let (org2, run) = (org.clone(), r.clone());
944        std::thread::spawn(move || {
945            let _guard = guard;
946            me.execute(&org2, &b, run, log);
947        });
948        Ok(Some(r))
949    }
950
951    fn execute(&self, org: &OrgId, b: &Backup, mut r: Run, mut log: RunLog) {
952        let store = self.runs(org, &b.spec.name);
953        let res = match &b.spec.volume {
954            Some(v) => crate::volume_backup::export::backup_once(self, org, &b.spec, v, &mut log),
955            None => self.backup_once(org, b, &mut log),
956        };
957        let stack = match &b.spec.volume {
958            Some(v) => v.clone(),
959            None => self
960                .inner
961                .apps
962                .get(org, &b.spec.database)
963                .ok()
964                .and_then(|a| a.spec.stack().ok())
965                .unwrap_or_default(),
966        };
967        let what = b.spec.volume.as_deref().unwrap_or(&b.spec.database);
968        let q = crate::stack::qualified(org, &stack);
969        let (kind, level, msg) = match res {
970            Ok(detail) => {
971                let msg = format!(
972                    "backup {} of {}: {} ({} bytes)",
973                    b.spec.name,
974                    what,
975                    detail["key"].as_str().unwrap_or(""),
976                    detail["size"]
977                );
978                r.detail = detail;
979                r.exit_code = Some(0);
980                r.finish(RunStatus::Succeeded);
981                ("backup.succeeded", "info", msg)
982            }
983            Err(e) => {
984                log.line(&format!("isb: {e}"));
985                r.error = Some(e.to_string());
986                r.finish(RunStatus::Failed);
987                (
988                    "backup.failed",
989                    "error",
990                    format!("backup {} of {} failed: {e}", b.spec.name, what),
991                )
992            }
993        };
994        if let Err(e) = store.finish(&mut r, &mut log) {
995            eprintln!("isb serve: backup {}: run {}: {e}", b.spec.name, r.id);
996        }
997        self.inner
998            .apps
999            .controller()
1000            .event(kind, level, &q, what, msg);
1001    }
1002
1003    /// One backup: dump, compress, upload, verify, prune.
1004    fn backup_once(&self, org: &OrgId, b: &Backup, log: &mut RunLog) -> Result<Value> {
1005        let (app, db) = self.database(org, &b.spec.database)?;
1006        let d = self.destination_get(org, &b.spec.destination)?;
1007        let c = self.client(org, &d)?;
1008        let stack = crate::stack::qualified(org, &app.spec.stack()?);
1009        let client = self.inner.apps.client().clone();
1010        let inst = crate::jobs::running_instance(&client, org, &stack, &app.spec.name)?;
1011        let prefix = backup_prefix(&d, org, &b.spec.name);
1012        let now = std::time::SystemTime::now()
1013            .duration_since(std::time::UNIX_EPOCH)
1014            .map(|x| x.as_secs() as i64)
1015            .unwrap_or(0);
1016        let key = object_key(&prefix, &app.spec.name, now, db.engine, b.spec.compression);
1017        log.line(&format!(
1018            "isb: dumping {} {} in {inst} to s3://{}/{key}",
1019            db.engine,
1020            db.database_name(&app.spec.name),
1021            d.bucket
1022        ));
1023        let started = std::time::Instant::now();
1024        let (raw, size) = dump_to(
1025            &client,
1026            org,
1027            &inst,
1028            db.engine,
1029            &c,
1030            &key,
1031            b.spec.compression,
1032            log,
1033        )?;
1034        log.line(&format!(
1035            "isb: uploaded {size} bytes ({raw} before compression) in {:.1}s",
1036            started.elapsed().as_secs_f64()
1037        ));
1038        match c.head(&key)? {
1039            Some(n) if n == size => log.line(&format!("isb: verified: {key} is {n} bytes")),
1040            other => {
1041                return Err(Error::invalid(format!(
1042                    "verify {key}: uploaded {size} bytes, HEAD says {other:?}"
1043                )));
1044            }
1045        }
1046        let objects = c.list(&prefix)?;
1047        let mut pruned = Vec::new();
1048        for k in select_prune(&prefix, &objects, b.spec.keep as usize) {
1049            match c.delete(&k) {
1050                Ok(()) => {
1051                    log.line(&format!("isb: pruned {k}"));
1052                    pruned.push(k);
1053                }
1054                Err(e) => log.line(&format!("isb: prune {k}: {e}")),
1055            }
1056        }
1057        Ok(json!({
1058            "key": key,
1059            "size": size,
1060            "dump_bytes": raw,
1061            "engine": db.engine,
1062            "database": db.database_name(&app.spec.name),
1063            "destination": d.name,
1064            "pruned": pruned,
1065        }))
1066    }
1067
1068    // --- restores ----------------------------------------------------------
1069
1070    /// Restore a backup object into an existing database (its data is
1071    /// replaced: `confirm` is required) or into a new database created for
1072    /// it. Runs in the background; returns the run record.
1073    pub fn restore(&self, org: &OrgId, req: RestoreRequest, by: &str) -> Result<Run> {
1074        if req.target.is_some() == req.new.is_some() {
1075            return Err(Error::invalid(
1076                "restore into `target` (an existing database) or `new` (a database to create), one of them",
1077            ));
1078        }
1079        let (dest, prefix, source_db) = match (&req.backup, &req.destination) {
1080            (Some(bn), _) => {
1081                let b = self.get(org, bn)?;
1082                if b.spec.volume.is_some() {
1083                    return Err(Error::invalid(format!(
1084                        "{bn} backs up a volume: restore it with volume_restore (isb volume restore), into a new volume"
1085                    )));
1086                }
1087                let d = self.destination_get(org, &b.spec.destination)?;
1088                let p = backup_prefix(&d, org, bn);
1089                (d, p, Some(b.spec.database.clone()))
1090            }
1091            (None, Some(dn)) => {
1092                let d = self.destination_get(org, dn)?;
1093                let key = req.key.as_deref().ok_or_else(|| {
1094                    Error::invalid("with `destination`, name the object with `key`")
1095                })?;
1096                let p = match key.rfind('/') {
1097                    Some(i) => key[..=i].to_string(),
1098                    None => String::new(),
1099                };
1100                (d, p, None)
1101            }
1102            (None, None) => {
1103                return Err(Error::invalid(
1104                    "name a `backup`, or a `destination` and `key`",
1105                ));
1106            }
1107        };
1108        let key = match &req.key {
1109            Some(k) => k.clone(),
1110            None => self
1111                .files_at(org, &dest, &prefix)?
1112                .into_iter()
1113                .next()
1114                .map(|f| f.key)
1115                .ok_or_else(|| Error::invalid(format!("no backups under {prefix}")))?,
1116        };
1117        let (dumped_app, _, engine, compression) = parse_key(&prefix, &key)
1118            .ok_or_else(|| Error::invalid(format!("{key}: not a backup isb wrote")))?;
1119        let source_db = source_db.unwrap_or(dumped_app);
1120        let target = match (&req.target, &req.new) {
1121            (Some(t), _) => {
1122                let (_, tdb) = self.database(org, t)?;
1123                if !tdb.engine.restores_from(engine) {
1124                    return Err(Error::invalid(format!(
1125                        "{key} is a {engine} dump; {t} is {}",
1126                        tdb.engine
1127                    )));
1128                }
1129                if !req.confirm {
1130                    return Err(Error::invalid(format!(
1131                        "restoring replaces the data in database {t}; pass confirm: true"
1132                    )));
1133                }
1134                t.clone()
1135            }
1136            (None, Some(n)) => n.name.clone(),
1137            _ => unreachable!(),
1138        };
1139        let store = self.restore_runs(org);
1140        let Some(guard) = self.inner.running.enter(org, "restore", &target, true) else {
1141            return Err(Error::invalid(format!(
1142                "a restore into {target} is already running"
1143            )));
1144        };
1145        let (mut r, log) = store.start("restore", RunTrigger::Manual, by, None, 50)?;
1146        r.detail = json!({
1147            "key": key, "destination": dest.name, "target": target,
1148            "new": req.new.is_some(), "engine": engine,
1149        });
1150        store.save(&r)?;
1151        let me = self.clone();
1152        let (org2, run) = (org.clone(), r.clone());
1153        let source = self.database(org, &source_db).ok();
1154        std::thread::spawn(move || {
1155            let _guard = guard;
1156            me.restore_run(&org2, run, log, dest, key, engine, compression, req, source);
1157        });
1158        Ok(r)
1159    }
1160
1161    #[expect(clippy::too_many_arguments)]
1162    fn restore_run(
1163        &self,
1164        org: &OrgId,
1165        mut r: Run,
1166        mut log: RunLog,
1167        dest: Destination,
1168        key: String,
1169        engine: Engine,
1170        compression: Compression,
1171        req: RestoreRequest,
1172        source: Option<(crate::app::App, DatabaseSource)>,
1173    ) {
1174        let store = self.restore_runs(org);
1175        let res = self.restore_once(
1176            org,
1177            &mut log,
1178            &dest,
1179            &key,
1180            engine,
1181            compression,
1182            &req,
1183            source.as_ref(),
1184        );
1185        let target = r.detail["target"].as_str().unwrap_or_default().to_string();
1186        let stack = self
1187            .inner
1188            .apps
1189            .get(org, &target)
1190            .ok()
1191            .and_then(|a| a.spec.stack().ok())
1192            .unwrap_or_default();
1193        let q = crate::stack::qualified(org, &stack);
1194        let (kind, level, msg) = match res {
1195            Ok(bytes) => {
1196                r.detail["bytes"] = json!(bytes);
1197                r.exit_code = Some(0);
1198                r.finish(RunStatus::Succeeded);
1199                (
1200                    "restore.succeeded",
1201                    "info",
1202                    format!("restore of {key} into {target} done ({bytes} bytes)"),
1203                )
1204            }
1205            Err(e) => {
1206                log.line(&format!("isb: {e}"));
1207                r.error = Some(e.to_string());
1208                r.finish(RunStatus::Failed);
1209                (
1210                    "restore.failed",
1211                    "error",
1212                    format!("restore of {key} into {target} failed: {e}"),
1213                )
1214            }
1215        };
1216        if let Err(e) = store.finish(&mut r, &mut log) {
1217            eprintln!("isb serve: restore {}: {e}", r.id);
1218        }
1219        self.inner
1220            .apps
1221            .controller()
1222            .event(kind, level, &q, &target, msg);
1223    }
1224
1225    #[expect(clippy::too_many_arguments)]
1226    fn restore_once(
1227        &self,
1228        org: &OrgId,
1229        log: &mut RunLog,
1230        dest: &Destination,
1231        key: &str,
1232        engine: Engine,
1233        compression: Compression,
1234        req: &RestoreRequest,
1235        source: Option<&(crate::app::App, DatabaseSource)>,
1236    ) -> Result<u64> {
1237        let apps = &self.inner.apps;
1238        let target = match (&req.target, &req.new) {
1239            (Some(t), _) => t.clone(),
1240            (None, Some(n)) => {
1241                let (project, environment, version) = match source {
1242                    Some((a, db)) => (
1243                        n.project.clone().unwrap_or_else(|| a.spec.project.clone()),
1244                        n.environment
1245                            .clone()
1246                            .unwrap_or_else(|| a.spec.environment.clone()),
1247                        n.version
1248                            .clone()
1249                            .unwrap_or_else(|| db.version().to_string()),
1250                    ),
1251                    None => (
1252                        n.project
1253                            .clone()
1254                            .ok_or_else(|| Error::invalid("new: give the project"))?,
1255                        n.environment
1256                            .clone()
1257                            .unwrap_or_else(|| crate::app::DEFAULT_ENVIRONMENT.into()),
1258                        n.version
1259                            .clone()
1260                            .unwrap_or_else(|| engine.default_version().into()),
1261                    ),
1262                };
1263                let spec: crate::app::AppSpec = serde_json::from_value(json!({
1264                    "name": n.name, "project": project, "environment": environment,
1265                    "source": {"database": {"engine": engine, "version": version}},
1266                }))?;
1267                log.line(&format!(
1268                    "isb: creating database {} ({engine} {version}) in {project}/{environment}",
1269                    n.name
1270                ));
1271                apps.create(org, spec)?;
1272                let d = apps.deploy(
1273                    org,
1274                    &n.name,
1275                    crate::app::deploy::Trigger::Api,
1276                    "restore",
1277                    None,
1278                )?;
1279                let d = apps.wait(org, &n.name, d.id, Duration::from_secs(900))?;
1280                if d.status != crate::app::deploy::Status::Done {
1281                    return Err(Error::invalid(format!(
1282                        "deploying {}: {}",
1283                        n.name,
1284                        d.error.unwrap_or_else(|| format!("{:?}", d.status))
1285                    )));
1286                }
1287                log.line(&format!("isb: database {} is up", n.name));
1288                n.name.clone()
1289            }
1290            _ => return Err(Error::invalid("no restore target")),
1291        };
1292        let (tapp, tdb) = self.database(org, &target)?;
1293        let stack = crate::stack::qualified(org, &tapp.spec.stack()?);
1294        let client = apps.client().clone();
1295        let inst = crate::jobs::running_instance(&client, org, &stack, &target)?;
1296        let c = self.client(org, dest)?;
1297        let (len, body) = c.get(key)?;
1298        log.line(&format!(
1299            "isb: restoring s3://{}/{key} ({len} bytes) into {target} ({inst})",
1300            dest.bucket
1301        ));
1302        let source_name = source
1303            .map(|(a, db)| db.database_name(&a.spec.name))
1304            .unwrap_or_else(|| tdb.database_name(&target));
1305        let reader = decompressor(compression, body)?;
1306        let bytes = restore_from(&client, org, &inst, tdb.engine, &source_name, reader, log)?;
1307        log.line(&format!("isb: restored {bytes} bytes of dump"));
1308        if tdb.engine == Engine::Redis {
1309            // The server stopped to load the snapshot; wait for it back.
1310            std::thread::sleep(Duration::from_secs(3));
1311            let (ok, why) = apps.wait_converged(org, &target)?;
1312            if !ok {
1313                return Err(Error::invalid(format!("{target} did not come back: {why}")));
1314            }
1315        }
1316        Ok(bytes)
1317    }
1318}
1319
1320/// Stream a dump out of `instance`, compressed, into `key`. Returns the
1321/// dump's size and the object's.
1322#[expect(clippy::too_many_arguments)]
1323pub fn dump_to(
1324    client: &crate::client::Client,
1325    org: &OrgId,
1326    instance: &str,
1327    engine: Engine,
1328    c: &crate::s3::Client,
1329    key: &str,
1330    compression: Compression,
1331    log: &mut RunLog,
1332) -> Result<(u64, u64)> {
1333    let oc = crate::org::client(client, org);
1334    let sb = Sandbox::get(&oc, instance)?;
1335    let mut s = sb.exec_stream(
1336        engine.dump_command(),
1337        ExecOptions::default().timeout(DUMP_TIMEOUT),
1338    )?;
1339    let mut w = compressor(compression, c.upload(key, compression.content_type()));
1340    let mut raw = 0u64;
1341    let mut failed: Option<String> = None;
1342    while let Some(ev) = s.next_event() {
1343        match ev {
1344            ExecEvent::Stdout(b) => {
1345                raw += b.len() as u64;
1346                if failed.is_none() {
1347                    if let Err(e) = w.write_all(&b) {
1348                        // Stop the dump; its upload is aborted below.
1349                        failed = Some(e.to_string());
1350                        let _ = s.signal(15);
1351                    }
1352                }
1353            }
1354            ExecEvent::Stderr(b) => log.write(&b),
1355        }
1356    }
1357    let code = match s.wait() {
1358        Err(Error::ExecTimeout { .. }) => {
1359            return Err(Error::invalid(format!(
1360                "the dump took longer than {DUMP_TIMEOUT:?}"
1361            )));
1362        }
1363        r => r?,
1364    };
1365    if let Some(e) = failed {
1366        return Err(Error::invalid(format!("upload: {e}")));
1367    }
1368    if code != 0 {
1369        return Err(Error::invalid(format!("the dump exited {code}")));
1370    }
1371    if raw == 0 {
1372        return Err(Error::invalid("the dump was empty"));
1373    }
1374    let upload = w
1375        .finish()
1376        .map_err(|e| Error::invalid(format!("compress: {e}")))?;
1377    let size = upload.finish()?;
1378    Ok((raw, size))
1379}
1380
1381/// Feed `reader` to the engine's restore command in `instance`. Returns
1382/// the bytes fed.
1383pub fn restore_from(
1384    client: &crate::client::Client,
1385    org: &OrgId,
1386    instance: &str,
1387    engine: Engine,
1388    source_db: &str,
1389    mut reader: Box<dyn Read + Send>,
1390    log: &mut RunLog,
1391) -> Result<u64> {
1392    let oc = crate::org::client(client, org);
1393    let sb = Sandbox::get(&oc, instance)?;
1394    let mut s = sb.exec_stream(
1395        engine.restore_command(),
1396        ExecOptions::default()
1397            .timeout(DUMP_TIMEOUT)
1398            .env("ISB_SOURCE_DB", source_db)
1399            .stdin(Stdin::Piped),
1400    )?;
1401    let ctl = s.controller();
1402    // The writer feeds stdin while this thread reads the output, so
1403    // neither side's bounded queue can stall the other.
1404    let writer = std::thread::spawn(move || -> Result<u64> {
1405        let mut buf = vec![0u8; 256 << 10];
1406        let mut n = 0u64;
1407        loop {
1408            let k = match reader.read(&mut buf) {
1409                Ok(0) => break,
1410                Ok(k) => k,
1411                Err(e) => {
1412                    let _ = ctl.close_stdin();
1413                    return Err(Error::invalid(format!("reading the backup: {e}")));
1414                }
1415            };
1416            if ctl.write_stdin(&buf[..k]).is_err() {
1417                // The command stopped reading; its exit code says why.
1418                break;
1419            }
1420            n += k as u64;
1421        }
1422        let _ = ctl.close_stdin();
1423        Ok(n)
1424    });
1425    while let Some(ev) = s.next_event() {
1426        match ev {
1427            ExecEvent::Stdout(b) | ExecEvent::Stderr(b) => log.write(&b),
1428        }
1429    }
1430    let code = s.wait();
1431    let fed = writer
1432        .join()
1433        .map_err(|_| Error::invalid("the restore writer panicked"))??;
1434    let code = code?;
1435    if code != 0 {
1436        return Err(Error::invalid(format!("the restore exited {code}")));
1437    }
1438    Ok(fed)
1439}
1440
1441impl Scheduled for Backups {
1442    fn entries(&self) -> Vec<Entry> {
1443        let mut out = Vec::new();
1444        for org in crate::jobs::orgs(&self.inner.state) {
1445            for b in self.list(&org).unwrap_or_default() {
1446                if !b.spec.enabled {
1447                    continue;
1448                }
1449                let (Ok(schedule), Ok(grace)) = (
1450                    b.spec.schedule(),
1451                    crate::jobs::parse_grace(&b.spec.missed_grace),
1452                ) else {
1453                    continue;
1454                };
1455                out.push(Entry {
1456                    org: org.clone(),
1457                    name: b.spec.name.clone(),
1458                    schedule,
1459                    anchor: b.anchor,
1460                    grace,
1461                });
1462            }
1463        }
1464        out
1465    }
1466
1467    fn fire(&self, e: &Entry, slot: i64, late: bool) {
1468        let b = {
1469            let _g = self.inner.edit.lock().unwrap();
1470            let Ok(mut b) = self.get(&e.org, &e.name) else {
1471                return;
1472            };
1473            if b.anchor >= slot {
1474                return;
1475            }
1476            b.anchor = slot;
1477            if let Err(err) = self.save(&e.org, &b) {
1478                eprintln!("isb serve: backup {}: {err}", e.name);
1479                return;
1480            }
1481            b
1482        };
1483        let trigger = if late {
1484            RunTrigger::Missed
1485        } else {
1486            RunTrigger::Schedule
1487        };
1488        if let Err(err) = self.start(&e.org, b, trigger, "schedule", Some(slot)) {
1489            eprintln!("isb serve: backup {}: {err}", e.name);
1490        }
1491    }
1492
1493    fn advance(&self, e: &Entry, to: i64) {
1494        let _g = self.inner.edit.lock().unwrap();
1495        if let Ok(mut b) = self.get(&e.org, &e.name) {
1496            b.anchor = to;
1497            let _ = self.save(&e.org, &b);
1498        }
1499    }
1500}
1501
1502#[cfg(test)]
1503mod tests {
1504    use super::*;
1505
1506    fn obj(key: &str) -> crate::s3::Object {
1507        crate::s3::Object {
1508            key: key.into(),
1509            size: 1,
1510            last_modified: String::new(),
1511        }
1512    }
1513
1514    #[test]
1515    fn keys_round_trip() {
1516        let d = Destination {
1517            name: "s3".into(),
1518            endpoint: "https://s3.example.com".into(),
1519            region: "us-east-1".into(),
1520            bucket: "b".into(),
1521            prefix: "/isb/".into(),
1522            path_style: true,
1523            access_key_secret: "a".into(),
1524            secret_key_secret: "s".into(),
1525            allow_local: false,
1526            created_at: 0,
1527        };
1528        let org = OrgId::new("acme").unwrap();
1529        let p = backup_prefix(&d, &org, "nightly");
1530        assert_eq!(p, "isb/acme/nightly/");
1531        let t = crate::cron::parse_compact_utc("20261003T040506Z").unwrap();
1532        let k = object_key(&p, "main-db", t, Engine::Postgres, Compression::Gzip);
1533        assert_eq!(k, "isb/acme/nightly/main-db-20261003T040506Z.postgres.gz");
1534        assert_eq!(
1535            parse_key(&p, &k),
1536            Some(("main-db".into(), t, Engine::Postgres, Compression::Gzip))
1537        );
1538        let z = object_key(&p, "c", t, Engine::Redis, Compression::Zstd);
1539        assert_eq!(parse_key(&p, &z).unwrap().3, Compression::Zstd);
1540        let n = object_key(&p, "c", t, Engine::Mysql, Compression::None);
1541        assert_eq!(parse_key(&p, &n).unwrap().2, Engine::Mysql);
1542        for other in [
1543            "isb/acme/nightly/notes.txt",
1544            "isb/acme/nightly/main-db-2026.postgres.gz",
1545            "isb/acme/nightly/main-db-20261003T040506Z.oracle.gz",
1546            "isb/acme/nightly/main-db-20261003T040506Z.postgres.gz.bak",
1547            "isb/acme/nightly/sub/main-db-20261003T040506Z.postgres.gz",
1548            "isb/acme/other/main-db-20261003T040506Z.postgres.gz",
1549        ] {
1550            assert_eq!(parse_key(&p, other), None, "{other}");
1551        }
1552        let empty = Destination {
1553            prefix: String::new(),
1554            ..d
1555        };
1556        assert_eq!(backup_prefix(&empty, &org, "x"), "acme/x/");
1557    }
1558
1559    #[test]
1560    fn retention_keeps_the_newest() {
1561        let p = "acme/nightly/";
1562        let keys: Vec<_> = [
1563            "acme/nightly/db-20261001T000000Z.postgres.gz",
1564            "acme/nightly/db-20261003T000000Z.postgres.gz",
1565            "acme/nightly/db-20261002T000000Z.postgres.gz",
1566            "acme/nightly/db-20260930T000000Z.postgres.zst",
1567            "acme/nightly/README",
1568            "acme/nightly/db-20200101T000000Z.postgres.gz.keep",
1569        ]
1570        .iter()
1571        .map(|k| obj(k))
1572        .collect();
1573        let mut gone = select_prune(p, &keys, 2);
1574        gone.sort();
1575        assert_eq!(
1576            gone,
1577            [
1578                "acme/nightly/db-20260930T000000Z.postgres.zst",
1579                "acme/nightly/db-20261001T000000Z.postgres.gz",
1580            ]
1581        );
1582        assert!(select_prune(p, &keys, 10).is_empty());
1583        assert_eq!(select_prune(p, &keys, 1).len(), 3);
1584    }
1585
1586    #[test]
1587    fn compression_round_trips() {
1588        let data: Vec<u8> = (0..300_000u32)
1589            .flat_map(|i| (i % 251).to_le_bytes())
1590            .collect();
1591        for c in [Compression::Gzip, Compression::Zstd, Compression::None] {
1592            let mut w = compressor(c, Vec::new());
1593            for chunk in data.chunks(10_000) {
1594                w.write_all(chunk).unwrap();
1595            }
1596            let out = w.finish().unwrap();
1597            if c != Compression::None {
1598                assert!(out.len() < data.len() / 2, "{c:?}: {}", out.len());
1599            }
1600            let mut back = Vec::new();
1601            decompressor(c, std::io::Cursor::new(out))
1602                .unwrap()
1603                .read_to_end(&mut back)
1604                .unwrap();
1605            assert_eq!(back, data, "{c:?}");
1606        }
1607    }
1608
1609    /// A sink that fails after some bytes, like a broken upload.
1610    struct Failing(usize);
1611    impl Write for Failing {
1612        fn write(&mut self, b: &[u8]) -> std::io::Result<usize> {
1613            if self.0 < b.len() {
1614                return Err(std::io::Error::other("upload broke"));
1615            }
1616            self.0 -= b.len();
1617            Ok(b.len())
1618        }
1619        fn flush(&mut self) -> std::io::Result<()> {
1620            Ok(())
1621        }
1622    }
1623
1624    #[test]
1625    fn a_failing_sink_is_an_error_not_a_panic() {
1626        // Incompressible, so the sink sees more than it takes.
1627        let mut x = 0x2545_f491_4f6c_dd1du64;
1628        let data: Vec<u8> = (0..4 << 20)
1629            .map(|_| {
1630                x ^= x << 13;
1631                x ^= x >> 7;
1632                x ^= x << 17;
1633                x as u8
1634            })
1635            .collect();
1636        for c in [Compression::Gzip, Compression::Zstd] {
1637            let mut w = compressor(c, Failing(1000));
1638            let wrote = data.chunks(64 << 10).try_for_each(|ch| w.write_all(ch));
1639            let finished = w.finish();
1640            assert!(wrote.is_err() || finished.is_err(), "{c:?}");
1641        }
1642    }
1643
1644    #[test]
1645    fn local_endpoints() {
1646        assert!(is_local_endpoint("http://127.0.0.1:9000"));
1647        assert!(is_local_endpoint("http://localhost:9000"));
1648        assert!(is_local_endpoint("http://[::1]:9000"));
1649        assert!(is_local_endpoint("http://169.254.169.254"));
1650        assert!(!is_local_endpoint("http://10.1.2.3:9000"));
1651        assert!(!is_local_endpoint("https://192.0.2.10"));
1652    }
1653
1654    #[test]
1655    fn spec_validation() {
1656        let b: BackupSpec = serde_json::from_value(json!({
1657            "name": "nightly", "database": "main-db", "destination": "s3", "schedule": "@daily",
1658        }))
1659        .unwrap();
1660        b.validate().unwrap();
1661        assert_eq!(
1662            (b.keep, b.compression, b.enabled),
1663            (7, Compression::Gzip, true)
1664        );
1665        let mut bad = b.clone();
1666        bad.schedule = "* * *".into();
1667        assert!(bad.validate().is_err());
1668        let mut bad = b.clone();
1669        bad.keep = 0;
1670        assert!(bad.validate().is_err());
1671        assert!(
1672            serde_json::from_value::<BackupSpec>(json!({"name": "x", "database": "d", "destination": "s", "schedule": "@daily", "compression": "lz4"})).is_err()
1673        );
1674    }
1675}