use std::fs;
use std::path::PathBuf;
use yo_common::{Code, Error, Result};
use yo_kv::Snapshot;
use super::args::{self, Args};
use super::server::{REPORTED_VERSION, help};
use super::{DATABASES, Server};
use crate::reply::Out;
pub(super) const DIR_NAME: &str = "backupdir";
const MANIFEST: &str = "appendonly.aof.manifest";
const IN_PROGRESS: &str = "A backup is already in progress, ABORT it first";
const SEALED_EXISTS: &str = "A sealed backup exists, CLEANUP it first";
const NOT_READY: &str = "No backup ready to seal (must be in the incrementing state)";
const NO_BACKUP: &str = "No backup in progress";
const RUNNING: &str = "Backup is in progress";
const ABORTED: &str = "aborted by user";
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
enum Phase {
Idle,
Incrementing,
Sealed,
Failed,
}
impl Phase {
const fn name(self) -> &'static str {
match self {
Phase::Idle => "idle",
Phase::Incrementing => "incrementing",
Phase::Sealed => "sealed",
Phase::Failed => "failed",
}
}
}
pub(super) struct State {
phase: Phase,
start_time: i64,
end_time: i64,
error: String,
seq: u64,
next: u64,
ttl: u64,
}
impl Default for State {
fn default() -> State {
State {
phase: Phase::Idle,
start_time: 0,
end_time: 0,
error: String::new(),
seq: 0,
next: 1,
ttl: 0,
}
}
}
impl State {
pub(super) const fn ttl(&self) -> u64 {
self.ttl
}
pub(super) const fn set_ttl(&mut self, seconds: u64) {
self.ttl = seconds;
}
}
pub(super) fn execute(server: &mut Server, args: Args<'_>, out: &mut Out) -> Result<()> {
let sub = args.get(1);
if args::is(sub, b"start") {
start(server)?;
out.ok();
} else if args::is(sub, b"seal") {
seal(server)?;
out.ok();
} else if args::is(sub, b"abort") {
abort(server)?;
out.ok();
} else if args::is(sub, b"cleanup") {
cleanup(server)?;
out.ok();
} else if args::is(sub, b"status") {
status(server, out);
} else if args::is(sub, b"list") {
list(server, out);
} else if args::is(sub, b"help") {
help(out, BACKUP_HELP);
} else {
return Err(args::unknown_subcommand(sub, "BACKUP"));
}
Ok(())
}
fn start(server: &mut Server) -> Result<()> {
match server.backup.phase {
Phase::Incrementing => return Err(Error::new(Code::Invalid, IN_PROGRESS)),
Phase::Sealed => return Err(Error::new(Code::Invalid, SEALED_EXISTS)),
Phase::Idle | Phase::Failed => {}
}
let seq = server.backup.next;
let now = seconds(server);
yo_alloc::allow(|| {
let dir = dir(server);
let (file, skipped) = image(server);
fs::create_dir_all(&dir)
.and_then(|()| fs::write(dir.join(base_name(seq)), &file))
.map_err(failed)?;
let b = &mut server.backup;
b.phase = Phase::Incrementing;
b.seq = seq;
b.next = seq + 1;
b.start_time = now;
b.end_time = 0;
b.error = if skipped == 0 {
String::new()
} else {
format!("{skipped} key(s) have no RDB shape and are not in this backup")
};
Ok(())
})
}
fn seal(server: &mut Server) -> Result<()> {
if server.backup.phase != Phase::Incrementing {
return Err(Error::new(Code::Invalid, NOT_READY));
}
let seq = server.backup.seq;
let now = seconds(server);
yo_alloc::allow(|| {
let dir = dir(server);
fs::write(dir.join(incr_name(seq)), [])
.and_then(|()| fs::write(dir.join(MANIFEST), manifest(seq)))
.map_err(failed)?;
server.backup.phase = Phase::Sealed;
server.backup.end_time = now;
Ok(())
})
}
fn abort(server: &mut Server) -> Result<()> {
if server.backup.phase != Phase::Incrementing {
return Err(Error::new(Code::Invalid, NO_BACKUP));
}
let seq = server.backup.seq;
yo_alloc::allow(|| {
let _ = fs::remove_file(dir(server).join(base_name(seq)));
server.backup.phase = Phase::Failed;
server.backup.error = ABORTED.to_string();
});
Ok(())
}
fn cleanup(server: &mut Server) -> Result<()> {
if server.backup.phase == Phase::Incrementing {
return Err(Error::new(Code::Invalid, RUNNING));
}
yo_alloc::allow(|| discard(server));
Ok(())
}
fn discard(server: &mut Server) {
let seq = server.backup.seq;
let dir = dir(server);
for name in [base_name(seq), incr_name(seq), MANIFEST.to_string()] {
let _ = fs::remove_file(dir.join(name));
}
let b = &mut server.backup;
b.phase = Phase::Idle;
b.start_time = 0;
b.end_time = 0;
b.error = String::new();
}
pub(super) fn expire(server: &mut Server) {
let b = &server.backup;
if b.phase != Phase::Sealed || b.ttl == 0 {
return;
}
let deadline = b.end_time.saturating_add_unsigned(b.ttl);
if seconds(server) < deadline {
return;
}
yo_alloc::allow(|| discard(server));
}
fn status(server: &Server, out: &mut Out) {
let b = &server.backup;
out.map(4);
out.bulk(b"state");
out.bulk(b.phase.name().as_bytes());
out.bulk(b"error");
out.bulk(b.error.as_bytes());
out.bulk(b"start_time");
out.int(b.start_time);
out.bulk(b"end_time");
out.int(b.end_time);
}
fn list(server: &Server, out: &mut Out) {
let seq = server.backup.seq;
let count = match server.backup.phase {
Phase::Incrementing => 1,
Phase::Sealed => 3,
Phase::Idle | Phase::Failed => 0,
};
out.array(count);
if count == 0 {
return;
}
yo_alloc::allow(|| {
let dir = dir(server);
out.bulk(dir.join(base_name(seq)).to_string_lossy().as_bytes());
if count == 3 {
out.bulk(dir.join(incr_name(seq)).to_string_lossy().as_bytes());
out.bulk(dir.join(MANIFEST).to_string_lossy().as_bytes());
}
});
}
fn image(server: &mut Server) -> (Vec<u8>, usize) {
let bits: &[u8] = if usize::BITS == 64 { b"64" } else { b"32" };
let mut snap = Snapshot::new();
snap.aux(b"redis-ver", REPORTED_VERSION.as_bytes());
snap.aux(b"redis-bits", bits);
snap.aux(b"aof-base", b"1");
for i in 0..DATABASES {
snap.database(i, server.striped(i));
}
let skipped = snap.skipped();
(snap.finish(), skipped)
}
fn manifest(seq: u64) -> String {
format!(
"file {} seq {seq} type b\nfile {} seq {seq} type i startoffset 0 endoffset 0\n",
base_name(seq),
incr_name(seq),
)
}
fn dir(server: &Server) -> PathBuf {
server.dir().join(DIR_NAME)
}
fn base_name(seq: u64) -> String {
format!("appendonly.aof.{seq}.base.rdb")
}
fn incr_name(seq: u64) -> String {
format!("appendonly.aof.{seq}.incr.aof")
}
fn seconds(server: &Server) -> i64 {
(server.clock.now_ms() / 1_000) as i64
}
fn failed(e: std::io::Error) -> Error {
Error::fmt(Code::Io, format_args!("Backup failed: {e}"))
}
const BACKUP_HELP: &[&str] = &[
"BACKUP <subcommand> [<arg> [value] [opt] ...]. Subcommands are:",
"START",
" Start a new backup into the configured 'backupdirname'.",
"SEAL",
" Freeze the current backup (BASE + INCR + manifest).",
"ABORT",
" Cancel a backup that has not been sealed yet.",
"CLEANUP",
" Remove a sealed backup's files and return to idle.",
"STATUS",
" Report the current backup state.",
"LIST",
" List the immutable files pinned so far.",
"HELP",
" Return this help.",
"HELP",
" Print this help.",
];