backup.rs
raw
//! The backup pipeline: scanner and chunker thread pools connected by
//! persistent queues, feeding a durable holding queue of manifest
//! entries that finalization turns into a snapshot.
//!
//! ```text
//! scan queue ──> scanner ──┬──> chunk queue ──> chunker ──────┐
//! ^ ^_______________/ ^ | | |
//! |___(subdirectories) |____________/ v |
//! (path became bounded channel|
//! a dir) | |
//! uploader |
//! (chunks to the |
//! backend) v
//! entries queue <───────(directories, symlinks,
//! | fully-uploaded files)
//! v
//! finalization: dedupe -> manifest -> snapshot
//! ```
//!
//! The chunker→uploader handoff is a bounded *in-memory* channel, not a
//! persistent queue: in-flight chunk data is cheaply re-derivable, so a
//! chunker simply waits to acknowledge each file until all of its
//! chunks are confirmed stored. A crash redelivers the file, and dedup
//! re-uploads only what is missing. A full channel blocks the chunkers,
//! which backs the scanner off via its own pacing gate — the pipeline
//! throttles end-to-end to the backend's write speed.
//!
//! Every stage re-validates paths when it dequeues them, so filesystem
//! changes between suspension and resumption are handled by design:
//! vanished paths are skipped, and paths whose type changed are
//! re-dispatched to the right stage. The snapshot's meaning is "each
//! object as it existed when it was read".
use std::{
collections::HashSet,
fs::Metadata,
path::{Path, PathBuf},
sync::{
Arc,
atomic::{AtomicBool, AtomicUsize, Ordering},
},
time::{Duration, Instant},
};
use crossbeam_channel::Sender;
use globset::GlobSet;
use persistent_queue::{Lease, PersistentQueue};
use repository::{ChunkId, Coverage, Entry, EntryKind, LockId, Repository, Snapshot, SnapshotId};
use sha2::{Digest as _, Sha256};
use super::{
Control, EarlyFinish, Error, Event, FINALIZE_REPORT_EVERY, Marker, Outcome, PROGRESS_FILE,
ProgressCounters, QueuedPath, Reporter, SkipReason, Stage, StageBacklog, StoreLimit,
WorkTracker, control_loop, load_progress, peek_paths, read_state_file, report_staffing,
send_event, store_state_file, worker_loop,
};
/// Bounds on how much a run may write between looks at the store's
/// remaining room. Within them it looks at half the headroom the last
/// answer showed, so the looks close in as the floor approaches: rarely
/// while there is room to spare, every few megabytes once there is not.
const FREE_SPACE_CHECK_MIN: u64 = 4 << 20;
/// See [`FREE_SPACE_CHECK_MIN`].
const FREE_SPACE_CHECK_MAX: u64 = 1 << 30;
/// How long a run may go without looking, however little it is writing.
/// A store is rarely a backup's alone, and the floor is about the room
/// left on it, not about what this run put there.
const FREE_SPACE_INTERVAL: Duration = Duration::from_secs(5);
const SCAN_QUEUE: &str = "scan.queue";
const CHUNK_QUEUE: &str = "chunk.queue";
const ENTRIES_QUEUE: &str = "entries.queue";
const MANIFEST_FILE: &str = "manifest";
const CLAIM_FILE: &str = "claim";
/// This run's claim on the repository: the lock it holds, and what the
/// repository had pruned when it last looked.
///
/// The lock says "there are chunks here that no snapshot names yet", and
/// is given back the moment that stops being true — when the backup
/// finishes, or when it suspends and publishes everything it has. A run
/// that was killed leaves its lock behind, which is exactly what should
/// stop a prune.
///
/// The epoch is the other half of the same problem. If someone breaks
/// this run's lock and prunes anyway, the entries it had recorded may
/// name chunks that are now gone; a changed epoch on resume is how the
/// run finds out.
#[derive(Debug, serde::Serialize, serde::Deserialize)]
struct Claim {
/// The lock this run holds, if it currently owes the repository
/// anything. Given back when it publishes what it has recorded, and
/// taken again if it resumes.
lock: Option<LockId>,
/// What the repository had pruned when this run last looked. Outlives
/// the lock, because the question it answers — "has anything been
/// swept since I recorded my work?" — is asked precisely when the run
/// is *not* holding one.
epoch: [u8; 16],
}
/// The manifest this run has stored in the repository so far, and the
/// partial snapshots published from it.
///
/// A backup's chunks are durable the moment they are uploaded, but until
/// a manifest names them they are unreachable. A checkpoint drains the
/// entries recorded since the last one into manifest chunks and
/// publishes a snapshot covering every segment so far; finalization
/// appends the last segment to the same list. Checkpoints happen as the
/// run goes and again as it stands down, so a backup is worth something
/// from early on rather than only once it finishes, and no segment is
/// ever stored twice.
#[derive(Debug, Default, serde::Serialize, serde::Deserialize)]
struct StoredManifest {
/// Manifest chunks, in order. Entries are read back as one
/// continuous CBOR stream, and each segment holds a whole number of
/// them, so appending a segment's chunks appends its entries.
chunks: Vec<ChunkId>,
entries: u64,
content_bytes: u64,
/// Partial snapshots published from these segments. Recorded before
/// the snapshot object is written and cleared only after its
/// replacement is durable, so a crash at any point leaves either a
/// restorable snapshot or a record of one to clean up — never an
/// orphan, and never a moment with nothing to restore from.
published: Vec<SnapshotId>,
}
pub struct BackupOptions {
/// The directory to back up.
pub root: PathBuf,
/// Directory for the pipeline's durable working state. Exclusively
/// owned by the pipeline: it is created, rewritten, and eventually
/// deleted as runs start and finish.
pub state_dir: PathBuf,
/// Paths matching any of these globs (relative to `root`) are not
/// backed up.
pub exclude: GlobSet,
/// Whether the scanner may descend into directories on a different
/// filesystem than `root`. Off by default: mount points are recorded
/// as empty directories but never entered without explicit
/// permission.
pub cross_mount_points: bool,
pub scanner_threads: usize,
pub chunker_threads: usize,
/// Threads writing sealed chunks to the backend. Sized for the
/// backend: for a local disk a few suffice; for a remote this is
/// effectively "how many concurrent transfers".
pub uploader_threads: usize,
/// Capacity, in chunks, of the in-memory chunker→uploader channel.
/// When it is full the chunkers block, which in turn backs the
/// scanner off via `chunk_backlog_limit`: the whole pipeline
/// throttles to the backend's write speed with bounded memory
/// (roughly this many chunks' worth) and bounded state-directory
/// disk.
pub upload_backlog_limit: usize,
/// Additional paths never backed up (compared canonically), on top
/// of the two the pipeline always protects: its own state directory
/// and the repository itself when it lives on the local filesystem.
/// The application registers the beeping state base here so other
/// suspended runs' state is not captured either.
pub protect: Vec<PathBuf>,
/// A file this size or larger is read under admission control, so
/// that a slow backend cannot end up with every chunker parked
/// part-way through an enormous one.
pub large_file_threshold: u64,
/// How many entries may be recorded but unpublished before the run
/// stops to checkpoint, and how long they may sit there whatever
/// their number.
///
/// Between checkpoints, files whose content is in the repository are
/// named by no snapshot: they are safe from a prune, because the run
/// holds the lock, and they come back on a resume, because the
/// entries queue is durable — but they cannot be restored, and they
/// are lost outright if the state directory is what the crash took.
/// These two bound that exposure from both sides: the count for a
/// tree of small files, where thousands go by in a minute, and the
/// interval for one of large files, where a handful can take hours.
pub checkpoint_entries: usize,
pub checkpoint_interval: Duration,
/// Room to leave on the store, in bytes. Once it has less than this
/// left, the backup finishes early: it publishes what it has stored
/// as a partial snapshot and ends, rather than filling the store up.
///
/// A floor, not a fence: see [`SpaceWatch`] for how closely it is
/// held to. `None` — the default — lets a backup use whatever room
/// there is, and a store that cannot say how much room it has left
/// (an object store, or FTP) is never held to a floor either.
pub minimum_free_space: Option<u64>,
/// A size for the repository to stay under, in bytes. Once it has
/// grown to this, the backup finishes early, exactly as it does for
/// [`BackupOptions::minimum_free_space`].
///
/// This is the limit for a store whose room is not the run's to
/// measure — a quota on someone else's server, most of all an FTP
/// account, where "how much is left" has no answer but "how much is
/// there" does. The store is measured once as the run starts, which
/// costs a full listing, and tracked from there by what the run
/// writes; a store that another writer is also filling will drift
/// out from under that.
///
/// Reached, not enforced to the byte: the chunk each uploader is
/// mid-write when the decision lands still lands.
///
/// `None` — the default — sets no size on the repository, and a
/// store that cannot say what it holds is never held to one either.
pub maximum_store_size: Option<u64>,
/// End the run as soon as it starts, as if the files it has left to
/// scan and chunk did not exist: what earlier sessions recorded is
/// published as the run's finished snapshot and the run is over.
///
/// For finishing with a store that cannot hold the whole tree —
/// take what fits, publish it, and narrow the next backup's excludes
/// to match.
pub early_finish: bool,
/// How many large files may be read at once. `None` derives a limit
/// from the chunker pool.
///
/// Concurrency buys nothing once the backend is the bottleneck —
/// the files simply share a fixed budget and all finish late — while
/// every part-read file is work a suspension throws away. Holding a
/// few open at a time keeps completions coming steadily, which is
/// what a backup meant to be interrupted needs.
pub concurrent_large_files: Option<usize>,
/// Scanner pacing: once this many files are waiting to be chunked,
/// the scanner pauses instead of discovering more. This bounds the
/// state directory's intermediate growth (the queue file stays
/// around the compaction threshold instead of growing with the
/// tree), while leaving the chunkers more backlog than they can
/// ever drain between pacing checks. The bound can be overshot by
/// at most one file per scanner thread.
pub chunk_backlog_limit: usize,
}
impl BackupOptions {
pub fn new(root: impl Into<PathBuf>, state_dir: impl Into<PathBuf>) -> BackupOptions {
BackupOptions {
root: root.into(),
state_dir: state_dir.into(),
exclude: GlobSet::empty(),
cross_mount_points: false,
protect: Vec::new(),
scanner_threads: 2,
chunker_threads: std::thread::available_parallelism().map_or(4, usize::from),
uploader_threads: 4,
upload_backlog_limit: 16,
large_file_threshold: 64 << 20,
checkpoint_entries: 10_000,
checkpoint_interval: Duration::from_secs(300),
minimum_free_space: None,
maximum_store_size: None,
early_finish: false,
concurrent_large_files: None,
chunk_backlog_limit: 4096,
}
}
}
/// Tracks how many of one file's chunks are still on their way through
/// the uploaders. The chunker that owns the file waits for zero before
/// queueing the file's manifest entry.
struct UploadTicket {
/// The file these chunks came from, so an uploader can say what it
/// is working on. Held once per file rather than per chunk.
path: PathBuf,
outstanding: AtomicUsize,
}
impl UploadTicket {
fn new(path: PathBuf) -> UploadTicket {
UploadTicket {
path,
outstanding: AtomicUsize::new(0),
}
}
fn expect_one(&self) {
self.outstanding.fetch_add(1, Ordering::AcqRel);
}
fn settle_one(&self) {
self.outstanding.fetch_sub(1, Ordering::AcqRel);
}
/// Waits until every expected upload has settled. Returns `false`
/// if the pipeline began stopping (or failed) first — uploader
/// failures surface through the control, so no separate error path
/// is needed here.
fn await_settled(&self, control: &Control) -> bool {
while self.outstanding.load(Ordering::Acquire) > 0 {
if control.stopping() {
return false;
}
std::thread::sleep(super::IDLE_WAIT);
}
true
}
}
/// A turn to read a large file, given back when dropped. Holds nothing
/// at all for a file below the threshold, which never waits.
struct LargeFile<'a>(Option<&'a AtomicUsize>);
impl Drop for LargeFile<'_> {
fn drop(&mut self) {
if let Some(count) = self.0 {
count.fetch_sub(1, Ordering::AcqRel);
}
}
}
/// One sealed chunk on its way to the backend, tied to the file whose
/// entry is waiting on it.
struct UploadJob {
chunk: repository::PreparedChunk,
ticket: Arc<UploadTicket>,
}
/// Everything the worker threads share.
struct BackupContext {
repository: Arc<Repository>,
root: PathBuf,
root_device: Option<u64>,
exclude: GlobSet,
cross_mount_points: bool,
scan: PersistentQueue<QueuedPath>,
chunk: PersistentQueue<QueuedPath>,
entries: PersistentQueue<Entry>,
tracker: WorkTracker,
control: Control,
events: Sender<Event>,
chunk_backlog_limit: usize,
large_file_threshold: u64,
concurrent_large_files: usize,
/// How many large files are being read right now, which
/// [`BackupContext::admit_large_file`] holds down to the limit.
large_files: AtomicUsize,
uploads: crossbeam_channel::Sender<UploadJob>,
/// Canonical paths the scanner never enters or captures: this run's
/// state directory, caller-registered paths, and a local
/// repository's directory.
protected: Vec<PathBuf>,
/// Running totals, persisted at every sync so a resumed run's
/// progress reporting continues from previous sessions.
progress: ProgressCounters,
}
impl BackupContext {
/// Makes completed work durable, downstream queue first, so that a
/// crash can only ever lose acknowledgements whose products already
/// persist — never the other way around. (The repository itself is
/// durable at every put, making it the most-downstream link.)
fn sync_queues(&self) -> Result<(), Error> {
self.entries.sync()?;
self.chunk.sync()?;
self.scan.sync()?;
Ok(())
}
/// The current exclude rules, applied self-or-ancestor so that an
/// excluded directory covers everything beneath it. Checked when
/// items are enqueued *and* again when they are dequeued: excludes
/// can change between a suspension and its resume, and the rules in
/// force at processing time are the ones that count. (Prospective
/// only — entries already recorded stay in the snapshot.)
fn is_excluded(&self, relative: &Path) -> bool {
crate::patterns::matches_with_ancestors(&self.exclude, relative)
}
/// True for paths holding beeping's own working data, which the
/// scanner must not capture — most importantly a local repository,
/// which would otherwise be backed up into itself.
fn is_protected(&self, absolute: &Path) -> bool {
self.protected
.iter()
.any(|protected| absolute.starts_with(protected))
}
/// Reports an event, folding it into the persistent progress totals.
fn emit(&self, event: Event) {
self.progress.record(&event);
send_event(&self.events, event);
}
/// Reports the run's totals and what is queued at each stage — the
/// state a display needs that no single work item announces.
fn report_status(&self) {
send_event(&self.events, Event::Totals(self.progress.snapshot()));
send_event(
&self.events,
Event::Backlog(vec![
StageBacklog {
stage: Stage::Scan,
waiting: self.scan.pending() as u64,
upcoming: Some(peek_paths(&self.scan, |queued| queued.0)),
},
// Files the scan has already found are work in hand,
// not work still to be looked at, and there are far too
// many of them to be worth naming
StageBacklog {
stage: Stage::Chunk,
waiting: self.chunk.pending() as u64,
upcoming: None,
},
// The uploaders' backlog is sealed chunks in flight
// between the chunkers and the backend, which belong to
// files the chunk stage is already naming
StageBacklog {
stage: Stage::Upload,
waiting: self.uploads.len() as u64,
upcoming: None,
},
// The entries queue is not work waiting for a worker: it
// is the record of everything stored, which finalization
// drains into the manifest. Naming its contents would
// name what the run has already done.
StageBacklog {
stage: Stage::Record,
waiting: self.entries.pending() as u64,
upcoming: None,
},
]),
);
}
fn skip(&self, path: &Path, reason: SkipReason) {
self.emit(Event::Skipped {
path: path.to_owned(),
reason,
});
}
/// Pacing: waits until the chunkers have room for more discovered
/// files, so the scanner never piles up unbounded backlog in the
/// state directory. Returns `false` if the run began stopping while
/// waiting — the caller must then bail out *without* acknowledging
/// its work item, leaving it to be redone if the run resumes.
///
/// This must never gate the scan queue itself (the scanner is its
/// own consumer there — waiting on it could deadlock), and the
/// entries queue must stay ungated because nothing drains it until
/// finalization. The chunk queue is the one safe, useful spot.
fn await_chunk_capacity(&self) -> bool {
while self.chunk.pending() >= self.chunk_backlog_limit {
if self.control.stopping() {
return false;
}
std::thread::sleep(super::IDLE_WAIT);
}
true
}
/// Waits for a turn to read a large file, returning a permit that
/// gives the turn back when the file is done with.
///
/// Small files never wait. Large ones wait *before being opened*,
/// which is the whole point: a worker waiting here holds a lease and
/// nothing else, so a suspension costs nothing, where a worker
/// blocked part-way through a twenty-gigabyte file has that much
/// reading and hashing to throw away. Returns `None` if the run
/// began stopping while waiting — the caller must then bail out
/// without acknowledging its work item.
fn admit_large_file(&self, size: u64) -> Option<LargeFile<'_>> {
if size < self.large_file_threshold {
return Some(LargeFile(None));
}
loop {
let in_flight = self.large_files.load(Ordering::Acquire);
if in_flight < self.concurrent_large_files
&& self
.large_files
.compare_exchange(
in_flight,
in_flight + 1,
Ordering::AcqRel,
Ordering::Acquire,
)
.is_ok()
{
return Some(LargeFile(Some(&self.large_files)));
}
if self.control.stopping() {
return None;
}
std::thread::sleep(super::IDLE_WAIT);
}
}
/// Hands a sealed chunk to the uploaders, blocking while the bounded
/// channel is full — this is the chunker→uploader backpressure.
/// Returns `false` if the run began stopping while waiting; the
/// caller must then abandon its file without acknowledging it.
fn send_upload(&self, chunk: repository::PreparedChunk, ticket: &Arc<UploadTicket>) -> bool {
ticket.expect_one();
let mut job = UploadJob {
chunk,
ticket: ticket.clone(),
};
loop {
match self.uploads.send_timeout(job, super::IDLE_WAIT) {
Ok(()) => return true,
Err(crossbeam_channel::SendTimeoutError::Timeout(returned)) => {
if self.control.stopping() {
return false;
}
job = returned;
}
Err(crossbeam_channel::SendTimeoutError::Disconnected(_)) => {
// Cannot happen while the context is alive; treat it
// as suspension rather than panic in a worker
return false;
}
}
}
}
/// Queues a directory for scanning if permitted, having already
/// recorded its entry. Descending stops at filesystem boundaries
/// unless crossing them was explicitly enabled.
fn queue_directory(&self, relative: &Path, metadata: &Metadata) -> Result<(), Error> {
if self.cross_mount_points || same_device(self.root_device, device_of(metadata)) {
self.tracker.add(1);
self.scan.put(&QueuedPath(relative.to_owned()))?;
} else {
self.skip(relative, SkipReason::MountPoint);
}
Ok(())
}
}
/// The working state a backup left in its state directory, as an
/// operation that wants to edit it finds it.
///
/// Editing one is only safe when the run that owns it stood down
/// cleanly, so what this reports first is whether it did.
pub enum WorkingState {
/// No backup has state here.
Missing,
/// A backup that suspended gracefully: it holds no lock, and
/// everything it recorded is named by a snapshot it published.
Suspended(SuspendedRun),
/// State that must not be edited, and why.
Unsettled(Unsettled),
}
/// Why a working state is not one to edit.
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Unsettled {
/// The run still holds its claim on the repository: it was killed
/// rather than suspended, and what it recorded is not all published.
Killed,
/// The run finished but died before clearing its state away.
Finished,
/// Work recorded but not yet published — a run that stopped somewhere
/// other than a clean suspension.
Unpublished,
/// The state belongs to some other operation entirely.
Foreign,
}
/// A suspended backup's state, held open so that nothing else can take
/// it while it is being edited.
///
/// The queues and the record of what the run has published are one
/// thing: a path taken out of the snapshots must go from the queues too,
/// or the resumed run would put it straight back.
pub struct SuspendedRun {
state_dir: PathBuf,
scan: PersistentQueue<QueuedPath>,
chunk: PersistentQueue<QueuedPath>,
/// Empty, and held open to keep it that way. Its emptiness is the
/// whole of what makes the rest of this safe: everything the run
/// recorded is in a published snapshot, so nothing it owns is
/// invisible to a sweep.
_entries: PersistentQueue<Entry>,
stored: StoredManifest,
claim: Claim,
}
impl SuspendedRun {
/// Opens the working state of a suspended backup of `root`, saying
/// what it found.
///
/// Fails outright if another process holds the state directory: the
/// queues carry a lock of their own, and a run that is alive keeps
/// it. That is what stops this from reaching a backup in progress.
pub fn open(state_dir: &Path, root: &Path) -> Result<WorkingState, Error> {
if !state_dir.exists() {
return Ok(WorkingState::Missing);
}
match Marker::load(state_dir)? {
Some(Marker::BackupRunning { root: existing }) if existing == root => {}
Some(Marker::Completed { .. }) => {
return Ok(WorkingState::Unsettled(Unsettled::Finished));
}
None => return Ok(WorkingState::Missing),
_ => return Ok(WorkingState::Unsettled(Unsettled::Foreign)),
}
let claim: Claim = match read_state_file(state_dir, CLAIM_FILE)? {
Some(claim) => claim,
// Written before the first chunk is stored; without it the
// run never got as far as owing the repository anything
None => return Ok(WorkingState::Missing),
};
if claim.lock.is_some() {
return Ok(WorkingState::Unsettled(Unsettled::Killed));
}
let scan = PersistentQueue::open(state_dir.join(SCAN_QUEUE))?;
let chunk = PersistentQueue::open(state_dir.join(CHUNK_QUEUE))?;
let entries: PersistentQueue<Entry> = PersistentQueue::open(state_dir.join(ENTRIES_QUEUE))?;
if !entries.is_empty() {
return Ok(WorkingState::Unsettled(Unsettled::Unpublished));
}
Ok(WorkingState::Suspended(SuspendedRun {
state_dir: state_dir.to_owned(),
scan,
chunk,
_entries: entries,
stored: read_state_file(state_dir, MANIFEST_FILE)?.unwrap_or_default(),
claim,
}))
}
/// The partial snapshots this run has published, which are the only
/// snapshots of its tree that belong to it.
pub fn published(&self) -> &[SnapshotId] {
&self.stored.published
}
/// How many queued paths the patterns select — what
/// [`SuspendedRun::drop_paths`] would drop.
pub fn count_paths(&self, select: &GlobSet) -> Result<u64, Error> {
let mut count = 0;
for queue in [&self.scan, &self.chunk] {
queue.inspect(|queued| {
if crate::patterns::matches_with_ancestors(select, &queued.0) {
count += 1;
}
})?;
}
Ok(count)
}
/// Drops from the queues every path the patterns select, so that
/// resuming does not put back what has just been taken out of the
/// snapshots.
///
/// Only what is queued: a directory still waiting to be scanned is
/// dropped, but one already scanned may have left children queued
/// under it, and those are dropped by the same rules. What this
/// cannot do is stop the run rediscovering them on disk, which is
/// what an exclude is for.
pub fn drop_paths(&self, select: &GlobSet) -> Result<u64, Error> {
let mut dropped = 0;
for queue in [&self.scan, &self.chunk] {
dropped += queue
.retain(|queued| !crate::patterns::matches_with_ancestors(select, &queued.0))?
as u64;
// Nothing else holds these, so one sync makes both the
// survivors and the drops durable together
queue.sync()?;
}
Ok(dropped)
}
/// Points the run at the snapshot that has replaced the one it
/// published, and at the manifest behind it.
///
/// The run appends to that manifest when it resumes, so this is what
/// keeps an excision out of the snapshot it eventually finishes with.
pub fn republished(
&mut self,
snapshot: SnapshotId,
manifest: Vec<ChunkId>,
entries: u64,
content_bytes: u64,
) -> Result<(), Error> {
self.stored = StoredManifest {
chunks: manifest,
entries,
content_bytes,
published: vec![snapshot],
};
store_state_file(&self.state_dir, MANIFEST_FILE, &self.stored)
}
/// Records the repository's epoch as this run's own, saying that
/// what has been swept since it last looked was nothing of its.
///
/// Only honest from an operation that spared everything the run had
/// recorded — which, its entries queue being empty, is everything
/// its published snapshot names.
pub fn observed(&mut self, epoch: [u8; 16]) -> Result<(), Error> {
self.claim = Claim { lock: None, epoch };
store_state_file(&self.state_dir, CLAIM_FILE, &self.claim)
}
}
/// Runs a backup to completion or suspension.
///
/// If `options.state_dir` holds the state of an interrupted run of the
/// same backup, that run is resumed; otherwise a fresh run starts. Store
/// `true` in `suspend` at any time to request suspension — or don't
/// bother, and kill the process: the next run recovers either way, this
/// path is merely tidier.
pub fn run(
repository: Arc<Repository>,
options: BackupOptions,
events: Sender<Event>,
suspend: Arc<AtomicBool>,
) -> Result<Outcome, Error> {
let root = options.root.canonicalize()?;
let state_dir = options.state_dir.clone();
std::fs::create_dir_all(&state_dir)?;
let fresh = match Marker::load(&state_dir)? {
Some(Marker::Completed {
snapshot,
finished_early,
}) => {
// The previous run finished entirely but died before
// removing its state. Its partial snapshots outlived the one
// they were standing in for, so retire them before the
// record of them goes away with the state.
//
// Never the snapshot the marker names, though. Coverage
// being part of a snapshot's identity is what keeps a
// finished one out of that list, and this is not the place
// to bet a repository's only snapshot on that holding.
if let Some(stored) = read_state_file::<StoredManifest>(&state_dir, MANIFEST_FILE)? {
for id in stored.published {
if id != snapshot {
repository.delete_snapshot(id)?;
}
}
}
// It owes the repository nothing: the snapshot is stored
if let Some(claim) = read_state_file::<Claim>(&state_dir, CLAIM_FILE)?
&& let Some(lock) = claim.lock
{
repository.release_lock(lock)?;
}
std::fs::remove_dir_all(&state_dir)?;
return Ok(match finished_early {
Some(reason) => Outcome::FinishedEarly {
snapshot: Some(snapshot),
reason,
},
None => Outcome::Completed(snapshot),
});
}
Some(Marker::BackupRunning { root: existing }) => {
if existing != root {
return Err(Error::StateMismatch {
existing,
requested: root,
});
}
false
}
Some(Marker::RestoreRunning { target, .. }) => {
return Err(Error::StateMismatch {
existing: target,
requested: root,
});
}
None => {
// No marker means initialization never committed; any state
// files present are debris from a run that died starting up.
std::fs::remove_dir_all(&state_dir)?;
std::fs::create_dir_all(&state_dir)?;
true
}
};
let scan: PersistentQueue<QueuedPath> = PersistentQueue::open(state_dir.join(SCAN_QUEUE))?;
let chunk: PersistentQueue<QueuedPath> = PersistentQueue::open(state_dir.join(CHUNK_QUEUE))?;
let entries: PersistentQueue<Entry> = PersistentQueue::open(state_dir.join(ENTRIES_QUEUE))?;
if fresh {
// Seed, make it durable, and only then commit the initialization
// by writing the marker.
scan.put(&QueuedPath(PathBuf::new()))?;
scan.sync()?;
Marker::BackupRunning { root: root.clone() }.store(&state_dir)?;
}
// What earlier sessions of this run already stored. Load-bearing,
// unlike the progress totals: reading it as absent would finish the
// backup with a snapshot omitting everything they did.
let mut stored: StoredManifest =
read_state_file(&state_dir, MANIFEST_FILE)?.unwrap_or_default();
let claim = claim_repository(&repository, &state_dir, &root, &stored, &entries, &chunk)?;
let tracker = WorkTracker::new((scan.pending() + chunk.pending()) as u64);
// Progress totals accumulated by previous sessions of this run;
// informational, so an unreadable file just means starting from zero
let prior = load_progress(&state_dir);
if !fresh {
send_event(&events, Event::Resumed(prior));
}
// Paths that must never be captured: this run's state, any paths
// the caller registered (canonicalized where they exist), and the
// repository itself when it is a local directory
let mut protected = vec![state_dir.canonicalize()?];
protected.extend(
options
.protect
.iter()
.filter_map(|path| path.canonicalize().ok()),
);
protected.extend(
repository
.backend()
.local_root()
.map(|root| root.to_owned()),
);
let scanner_threads = options.scanner_threads.max(1);
let chunker_threads = options.chunker_threads.max(1);
let uploader_threads = options.uploader_threads.max(1);
let (uploads, upload_jobs) =
crossbeam_channel::bounded::<UploadJob>(options.upload_backlog_limit.max(1));
let context = Arc::new(BackupContext {
root_device: device_of(&std::fs::metadata(&root)?),
repository,
root,
exclude: options.exclude,
cross_mount_points: options.cross_mount_points,
scan,
chunk,
entries,
tracker,
control: Control::new(),
events,
chunk_backlog_limit: options.chunk_backlog_limit.max(1),
large_file_threshold: options.large_file_threshold,
// One large file per eight chunkers, so a modest pool reads one
// at a time and a large one still keeps a couple going
concurrent_large_files: options
.concurrent_large_files
.unwrap_or_else(|| chunker_threads.div_ceil(8))
.max(1),
large_files: AtomicUsize::new(0),
uploads,
protected,
progress: ProgressCounters::new(prior),
});
// Asked for on the command line: the run is over before it takes
// anything up, and publishes what earlier sessions recorded.
if options.early_finish {
finish_early(&context, EarlyFinish::Requested);
}
// Both of these are settled before a worker exists, so a store that
// is already past its limits stops the run without adding to it
let mut store = StoreWatch::start(
&context,
options.minimum_free_space,
options.maximum_store_size,
);
store.look(&context);
let mut workers = Vec::new();
report_staffing(&context.events, Stage::Scan, scanner_threads);
report_staffing(&context.events, Stage::Chunk, chunker_threads);
report_staffing(&context.events, Stage::Upload, uploader_threads);
for index in 0..scanner_threads {
let context = context.clone();
workers.push(
std::thread::Builder::new()
.name(format!("scanner-{index}"))
.spawn(move || {
worker_loop(
&context.scan,
&context.control,
Reporter::new(&context.events, Stage::Scan, index),
|lease, reporter| scan_one(&context, lease, reporter),
)
})?,
);
}
for index in 0..chunker_threads {
let context = context.clone();
workers.push(
std::thread::Builder::new()
.name(format!("chunker-{index}"))
.spawn(move || {
worker_loop(
&context.chunk,
&context.control,
Reporter::new(&context.events, Stage::Chunk, index),
|lease, reporter| chunk_one(&context, lease, reporter),
)
})?,
);
}
for index in 0..uploader_threads {
let context = context.clone();
let upload_jobs = upload_jobs.clone();
workers.push(
std::thread::Builder::new()
.name(format!("uploader-{index}"))
.spawn(move || {
upload_loop(
&context,
&upload_jobs,
Reporter::new(&context.events, Stage::Upload, index),
)
})?,
);
}
let sync_all = || {
// Totals first: they are informational, and this order at worst
// overcounts redone work after a crash, never loses completed
// counts behind durable acknowledgements
store_state_file(&state_dir, PROGRESS_FILE, &context.progress.snapshot())?;
// Advisory, so it rides along rather than being ordered against
// anything: losing recent additions costs round trips, not work
if let Some(cache) = context.repository.chunk_cache() {
cache.flush();
}
context.sync_queues()
};
// Publishing as the run goes is what keeps the backup restorable
// while it is still going: between checkpoints, everything recorded
// is content in the repository that no snapshot names.
let mut last_checkpoint = Instant::now();
let tick = || {
store.look(&context);
let recorded = context.entries.pending();
if recorded == 0
|| (recorded < options.checkpoint_entries
&& last_checkpoint.elapsed() < options.checkpoint_interval)
{
return Ok(());
}
checkpoint(&context, &state_dir, &mut stored)?;
last_checkpoint = Instant::now();
context.emit(Event::Working);
Ok(())
};
control_loop(
&context.control,
&suspend,
|| context.tracker.is_idle(),
|| context.report_status(),
sync_all,
tick,
);
for worker in workers {
if worker.join().is_err() {
context.control.fail(Error::WorkerPanicked);
}
}
sync_all()?;
// The workers have stopped, so this sample is the exact final word
// on what the run accomplished
context.report_status();
if let Some(error) = context.control.take_failure() {
return Err(error);
}
// Whether this run is ending short of the tree on purpose. A person
// asking to suspend outranks that: suspending keeps the work that
// is left, and a store still short of room trips the floor again on
// the next resume.
let finished_early = if context.control.finished() || context.control.suspend_requested() {
None
} else {
context.control.early_finish()
};
if !context.control.finished() && finished_early.is_none() {
// Publish what has been backed up so far before standing down,
// so the interrupted backup is restorable from the repository
// alone rather than only resumable from this machine
let restorable = checkpoint(&context, &state_dir, &mut stored)?;
// Everything this run recorded is now named by a snapshot, so it
// no longer needs to hold off a prune
release_claim(&context.repository, &state_dir, &claim)?;
return Ok(Outcome::Suspended { restorable });
}
let ending = match finished_early {
Some(_) => Ending::Short,
None => Ending::Whole,
};
let published = finalize(&context, &stored, ending)?;
if let Some(snapshot) = published {
Marker::Completed {
snapshot,
finished_early,
}
.store(&state_dir)?;
context.entries.sync()?;
// The finished snapshot stands in for every partial one
// published along the way — including, for a run that stopped
// short just after a checkpoint, one covering this very
// manifest. That the two are distinct objects at all is down to
// coverage being part of a snapshot's identity.
retire_partials(&context, &state_dir, &mut stored, Some(snapshot))?;
}
// Everything recorded is published (or there was nothing to
// publish), so the run owes the repository nothing and its working
// state has no more to say. Dropping the context releases the
// queues' locks before the directory holding them goes.
release_claim(&context.repository, &state_dir, &claim)?;
drop(context);
std::fs::remove_dir_all(&state_dir)?;
Ok(match finished_early {
Some(reason) => Outcome::FinishedEarly {
snapshot: published,
reason,
},
// Only an early finish can have nothing to publish: a run that
// walked the whole tree records a snapshot even of an empty
// one, because "this tree holds nothing" is worth restoring.
None => Outcome::Completed(
published.expect("a run that covered the whole tree publishes a snapshot"),
),
})
}
/// The limits a run is being held to on its store, each watching in the
/// way its question can be answered.
struct StoreWatch {
space: Option<SpaceWatch>,
size: Option<SizeWatch>,
}
impl StoreWatch {
/// Sets up whichever limits the store can answer for, saying so
/// where one cannot be applied. Measuring can be slow — a size
/// limit costs a full listing — so this happens before any worker
/// has taken anything up, and not at all for a run that is already
/// on its way out.
fn start(context: &BackupContext, floor: Option<u64>, limit: Option<u64>) -> StoreWatch {
if context.control.stopping() {
return StoreWatch {
space: None,
size: None,
};
}
StoreWatch {
space: floor.and_then(|floor| SpaceWatch::start(context, floor)),
size: limit.and_then(|limit| SizeWatch::start(context, limit)),
}
}
fn look(&mut self, context: &BackupContext) {
if let Some(space) = &mut self.space {
space.look(context);
}
if let Some(size) = &self.size {
size.look(context);
}
}
}
/// Watches the room left on the store, and ends the run early once
/// there is less of it than the backup was told to leave.
///
/// Looking costs a round trip, so it is paced from both sides: by what
/// the run has written since the last look, which is what can have
/// consumed the room, and by plain elapsed time, since whatever else
/// shares the store is filling it too.
struct SpaceWatch {
floor: u64,
/// When the store was last asked, and what the run had written by
/// then. `None` until the first look, which is due immediately.
last: Option<(Instant, u64)>,
/// How much more may be written before the next look, taken from
/// the headroom the last answer showed.
allowance: u64,
}
impl SpaceWatch {
fn start(context: &BackupContext, floor: u64) -> Option<SpaceWatch> {
if !matches!(context.repository.backend().free_space(), Ok(Some(_))) {
context.emit(Event::LimitNotApplied(StoreLimit::MinimumFreeSpace));
return None;
}
Some(SpaceWatch {
floor,
last: None,
allowance: FREE_SPACE_CHECK_MIN,
})
}
/// Looks if a look is due, and ends the run if the store is down to
/// its floor.
///
/// A store that stops answering is left alone rather than failing
/// the backup: the floor is a courtesy to whatever else shares the
/// disk, and a store that is genuinely out of room fails the next
/// write loudly enough on its own.
fn look(&mut self, context: &BackupContext) {
let written = context.repository.bytes_written();
if let Some((asked, mark)) = self.last
&& written.saturating_sub(mark) < self.allowance
&& asked.elapsed() < FREE_SPACE_INTERVAL
{
return;
}
self.last = Some((Instant::now(), written));
let Ok(Some(free)) = context.repository.backend().free_space() else {
return;
};
// Half of what is left above the floor, so however fast the run
// is filling the store, the next look comes with room to spare
self.allowance =
(free.saturating_sub(self.floor) / 2).clamp(FREE_SPACE_CHECK_MIN, FREE_SPACE_CHECK_MAX);
if free >= self.floor {
return;
}
finish_early(
context,
EarlyFinish::StoreNearlyFull {
free,
floor: self.floor,
},
);
}
}
/// Watches the repository's own size against the limit it was given.
///
/// Measuring a store costs a listing, so it is measured once and
/// followed from there by what this run writes — which is the whole of
/// what grows it, short of another writer working on the same
/// repository at the same time.
struct SizeWatch {
limit: u64,
/// What the store held when it was measured, and what this run had
/// written by then.
measured: u64,
written_at_measure: u64,
}
impl SizeWatch {
fn start(context: &BackupContext, limit: u64) -> Option<SizeWatch> {
let Ok(Some(measured)) = context.repository.backend().used_space() else {
context.emit(Event::LimitNotApplied(StoreLimit::MaximumStoreSize));
return None;
};
Some(SizeWatch {
limit,
measured,
written_at_measure: context.repository.bytes_written(),
})
}
/// Costs nothing to ask, so it is asked every time round.
fn look(&self, context: &BackupContext) {
let used = self.measured.saturating_add(
context
.repository
.bytes_written()
.saturating_sub(self.written_at_measure),
);
if used < self.limit {
return;
}
finish_early(
context,
EarlyFinish::StoreAtMaximum {
used,
limit: self.limit,
},
);
}
}
/// Ends the run, announcing why — once, whatever else reaches the same
/// conclusion afterwards.
fn finish_early(context: &BackupContext, reason: EarlyFinish) {
if context.control.request_early_finish(reason) {
context.emit(Event::FinishingEarly(reason));
}
}
/// Scanner stage: catalogue one directory's immediate children.
fn scan_one(
context: &BackupContext,
lease: Lease<'_, QueuedPath>,
reporter: &Reporter<'_>,
) -> Result<(), Error> {
let relative: &PathBuf = &lease.0;
let absolute = context.root.join(relative);
reporter.began(relative);
// Re-validate against the *current* excludes: this directory may
// have been queued before a suspension, under older rules
if context.is_excluded(relative) {
context.skip(relative, SkipReason::Excluded);
lease.ack();
context.tracker.finish_one();
return Ok(());
}
let directory = match std::fs::read_dir(&absolute) {
Ok(directory) => directory,
Err(err) => {
match err.kind() {
std::io::ErrorKind::NotFound => context.skip(relative, SkipReason::Vanished),
std::io::ErrorKind::NotADirectory => {
// It was a directory when discovered but is a file
// now; hand it to the chunker stage.
if !context.await_chunk_capacity() {
// Stopping: drop the lease unacknowledged so
// this item is redone if the run resumes
return Ok(());
}
context.tracker.add(1);
context.chunk.put(&QueuedPath(relative.clone()))?;
}
_ => context.skip(relative, SkipReason::Unreadable(err.to_string())),
}
lease.ack();
context.tracker.finish_one();
return Ok(());
}
};
for child in directory {
let child = match child {
Ok(child) => child,
Err(err) => {
context.skip(relative, SkipReason::Unreadable(err.to_string()));
continue;
}
};
let child_relative = relative.join(child.file_name());
if context.is_excluded(&child_relative) {
context.skip(&child_relative, SkipReason::Excluded);
continue;
}
if context.is_protected(&child.path()) {
context.skip(&child_relative, SkipReason::Internal);
continue;
}
// DirEntry::metadata does not traverse symlinks, which is what
// we want: symlinks are recorded, never followed.
let metadata = match child.metadata() {
Ok(metadata) => metadata,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
context.skip(&child_relative, SkipReason::Vanished);
continue;
}
Err(err) => {
context.skip(&child_relative, SkipReason::Unreadable(err.to_string()));
continue;
}
};
let file_type = metadata.file_type();
if file_type.is_symlink() {
match std::fs::read_link(child.path()) {
Ok(target) => context.entries.put(&Entry {
path: child_relative,
kind: EntryKind::Symlink { target },
mode: mode_of(&metadata),
mtime: mtime_of(&metadata),
})?,
Err(err) => context.skip(&child_relative, SkipReason::Unreadable(err.to_string())),
}
} else if file_type.is_file() {
// Pacing point: with the chunkers saturated, pause discovery
// mid-directory rather than queue further ahead. Bailing out
// when the run stops leaves this directory unacknowledged; a
// resumed run scans it again, and finalization deduplicates
// the children queued twice.
if !context.await_chunk_capacity() {
return Ok(());
}
context.tracker.add(1);
context.chunk.put(&QueuedPath(child_relative))?;
} else if file_type.is_dir() {
context.entries.put(&Entry {
path: child_relative.clone(),
kind: EntryKind::Directory,
mode: mode_of(&metadata),
mtime: mtime_of(&metadata),
})?;
context.queue_directory(&child_relative, &metadata)?;
} else {
context.skip(&child_relative, SkipReason::SpecialFile);
}
}
context.emit(Event::Scanned {
path: relative.clone(),
});
lease.ack();
context.tracker.finish_one();
Ok(())
}
/// Chunker stage: store one file's content in the repository and emit
/// its manifest entry.
fn chunk_one(
context: &BackupContext,
lease: Lease<'_, QueuedPath>,
reporter: &Reporter<'_>,
) -> Result<(), Error> {
let relative: &PathBuf = &lease.0;
let absolute = context.root.join(relative);
// Re-validate against the *current* excludes (they may have changed
// across a suspension), then against the filesystem: the path was
// an includable file when the scanner saw it, but neither fact is
// guaranteed to still hold.
if context.is_excluded(relative) {
context.skip(relative, SkipReason::Excluded);
lease.ack();
context.tracker.finish_one();
return Ok(());
}
match std::fs::symlink_metadata(&absolute) {
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
context.skip(relative, SkipReason::Vanished);
}
Err(err) => context.skip(relative, SkipReason::Unreadable(err.to_string())),
Ok(metadata) if metadata.file_type().is_symlink() => match std::fs::read_link(&absolute) {
Ok(target) => context.entries.put(&Entry {
path: relative.clone(),
kind: EntryKind::Symlink { target },
mode: mode_of(&metadata),
mtime: mtime_of(&metadata),
})?,
Err(err) => context.skip(relative, SkipReason::Unreadable(err.to_string())),
},
Ok(metadata) if metadata.file_type().is_dir() => {
context.entries.put(&Entry {
path: relative.clone(),
kind: EntryKind::Directory,
mode: mode_of(&metadata),
mtime: mtime_of(&metadata),
})?;
context.queue_directory(relative, &metadata)?;
}
Ok(metadata) if metadata.file_type().is_file() => {
if store_file(context, relative, &absolute, &metadata, reporter)?
== FileOutcome::Interrupted
{
// Stopped mid-file: drop the lease unacknowledged so
// the whole file is redone if the run resumes, when
// dedup will skip whatever chunks already reached the
// backend
return Ok(());
}
}
Ok(_) => context.skip(relative, SkipReason::SpecialFile),
}
lease.ack();
context.tracker.finish_one();
Ok(())
}
/// How [`store_file`] ended: either the file's entry is queued and its
/// chunks are confirmed in the backend, or the run began stopping and
/// the caller must leave its work item unacknowledged.
#[derive(PartialEq, Eq)]
enum FileOutcome {
Recorded,
Interrupted,
}
/// Reads a file, seals its chunks, hands the new ones to the uploaders,
/// and — once every one of them is confirmed stored — emits the file's
/// manifest entry. Errors reading the file skip it (with an event);
/// errors writing to the repository are fatal.
///
/// Waiting for the uploads before returning is what makes the in-memory
/// upload channel safe without any durable staging: the file's queue
/// item is only acknowledged (by the caller) after all of its chunks
/// exist in the backend, so a crash at any point redelivers the file
/// and re-uploads only what is missing.
fn store_file(
context: &BackupContext,
relative: &Path,
absolute: &Path,
metadata: &Metadata,
reporter: &Reporter<'_>,
) -> Result<FileOutcome, Error> {
// Wait for a turn before opening it: until this returns, this
// worker has read nothing and a suspension costs nothing
let Some(_turn) = context.admit_large_file(metadata.len()) else {
return Ok(FileOutcome::Interrupted);
};
reporter.began(relative);
let file = match std::fs::File::open(absolute) {
Ok(file) => file,
Err(err) if err.kind() == std::io::ErrorKind::NotFound => {
context.skip(relative, SkipReason::Vanished);
return Ok(FileOutcome::Recorded);
}
Err(err) => {
context.skip(relative, SkipReason::Unreadable(err.to_string()));
return Ok(FileOutcome::Recorded);
}
};
let ticket = Arc::new(UploadTicket::new(relative.to_owned()));
let mut chunks = Vec::new();
let mut bytes = 0u64;
for chunk in context.repository.chunk_stream(file) {
let chunk = match chunk {
Ok(chunk) => chunk,
Err(repository::Error::IO(err)) => {
// The source file failed mid-read; give up on this file
// but not on the backup. Any chunks already stored just
// become unreferenced.
context.skip(relative, SkipReason::Unreadable(err.to_string()));
return Ok(FileOutcome::Recorded);
}
Err(err) => return Err(err.into()),
};
bytes += chunk.len() as u64;
// A file large enough to hold this stage up is worth showing
// progress through, so a stalled backup can be told from one
// grinding through something enormous
reporter.progressed(bytes, metadata.len());
let (id, prepared) = context.repository.prepare_chunk(&chunk)?;
chunks.push(id);
if let Some(prepared) = prepared
&& !context.send_upload(prepared, &ticket)
{
return Ok(FileOutcome::Interrupted);
}
}
// The entry must not be queued before its chunks are all stored
if !ticket.await_settled(&context.control) {
return Ok(FileOutcome::Interrupted);
}
context.entries.put(&Entry {
path: relative.to_owned(),
kind: EntryKind::File { chunks, len: bytes },
mode: mode_of(metadata),
mtime: mtime_of(metadata),
})?;
context.emit(Event::Stored {
path: relative.to_owned(),
bytes,
});
Ok(FileOutcome::Recorded)
}
/// Uploader stage: drain prepared chunks from the in-memory channel
/// into the backend. Kept separate from the chunkers so that slow
/// (remote) writes and CPU-bound sealing can be scaled independently;
/// the bounded channel between them is the backpressure that stops the
/// chunkers from sealing further ahead than the uploads can absorb.
fn upload_loop(
context: &BackupContext,
uploads: &crossbeam_channel::Receiver<UploadJob>,
reporter: Reporter<'_>,
) {
loop {
if context.control.stopping() {
reporter.idle();
return;
}
match uploads.recv_timeout(super::IDLE_WAIT) {
Ok(job) => {
reporter.began(&job.ticket.path);
if let Err(err) = context.repository.store_prepared(&job.chunk) {
context.control.fail(err.into());
reporter.idle();
return;
}
send_event(
&context.events,
Event::Transferred {
bytes: job.chunk.sealed_len() as u64,
},
);
job.ticket.settle_one();
}
Err(crossbeam_channel::RecvTimeoutError::Timeout) => {
reporter.idle();
if context.control.finished() {
return;
}
}
Err(crossbeam_channel::RecvTimeoutError::Disconnected) => {
reporter.idle();
return;
}
}
}
}
/// Drains everything recorded so far into manifest chunks, appends them
/// to the run's stored manifest, and publishes a partial snapshot
/// covering the lot.
///
/// This is what makes an unfinished backup restorable on its own,
/// whether it is standing down or merely partway through. The order is
/// the pipeline's usual downstream-first rule applied to a new
/// pair: the manifest chunks reach the repository, then the state file
/// naming them is made durable, and only then may the entry
/// acknowledgements be synced. A crash before that last step redelivers
/// the entries and records them a second time, which costs a little
/// manifest and no uploads at all — the file chunks they name are
/// already stored, so they deduplicate.
fn checkpoint(
context: &BackupContext,
state_dir: &Path,
stored: &mut StoredManifest,
) -> Result<Option<SnapshotId>, Error> {
context.emit(Event::Checkpointing);
let segment = store_segment(context)?;
if segment.entries == 0 {
// Nothing new to cover; whatever was published last still stands
return Ok(stored.published.last().copied());
}
stored.chunks.extend(segment.chunks);
stored.entries += segment.entries;
stored.content_bytes += segment.content_bytes;
let id = context
.repository
.snapshot_id(&context.root, &stored.chunks, Coverage::Partial);
// Written before the snapshot object exists, so the cleanup below
// knows about it even if this run dies here
stored.published.push(id);
store_state_file(state_dir, MANIFEST_FILE, stored)?;
// Only now may the drained acknowledgements become durable
context.entries.sync()?;
context.repository.store_snapshot(&Snapshot::partial(
id,
context.root.clone(),
stored.chunks.clone(),
stored.entries,
stored.content_bytes,
))?;
retire_partials(context, state_dir, stored, Some(id))?;
Ok(Some(id))
}
/// Takes this run's claim on the repository, and — if the run is being
/// resumed and the repository has been pruned since it last looked —
/// puts back whatever it had recorded that may no longer be there.
///
/// A prune can only have happened while this run held a lock if someone
/// broke it, which is a deliberate act with a documented consequence:
/// this. The entries queue may name chunks that have been swept, so
/// every file entry in it goes back to the chunker, which will re-read
/// the file and store whatever is missing. Directories and symlinks name
/// no chunks and stay as they are.
fn claim_repository(
repository: &Repository,
state_dir: &Path,
root: &Path,
stored: &StoredManifest,
entries: &PersistentQueue<Entry>,
chunk: &PersistentQueue<QueuedPath>,
) -> Result<Claim, Error> {
let epoch = repository.epoch()?;
let previous: Option<Claim> = read_state_file(state_dir, CLAIM_FILE)?;
if let Some(previous) = &previous
&& previous.epoch != epoch
{
// The manifest segments this run has already stored are kept
// alive by the partial snapshot naming them. If that snapshot
// has been pruned, they may be gone, and nothing here knows what
// was in them — a suspended run's partial snapshot is not
// something to delete lightly.
//
// Only worth checking when the epoch moved: a published id with
// no object behind it is otherwise just the window between
// writing down the intent and creating the snapshot.
for snapshot in &stored.published {
if matches!(
repository.load_snapshot(*snapshot),
Err(repository::Error::MissingSnapshot(_))
) {
return Err(Error::RecordedWorkPruned {
snapshot: *snapshot,
state_dir: state_dir.to_owned(),
});
}
}
requeue_recorded_files(entries, chunk)?;
}
// A run that still holds a lock from an earlier session keeps it,
// rather than leaving it behind and taking another
let lock = match previous.and_then(|previous| previous.lock) {
Some(lock) => lock,
None => {
let lock = repository::Lock::new(root.to_owned(), hostname());
repository.take_lock(&lock)?;
lock.id
}
};
let claim = Claim {
lock: Some(lock),
epoch,
};
store_state_file(state_dir, CLAIM_FILE, &claim)?;
Ok(claim)
}
/// Gives the lock back, once everything this run has recorded is named
/// by a snapshot.
///
/// The epoch stays behind: a suspended run holds no lock, and is exactly
/// the run that needs to notice a prune when it wakes up. Writing that
/// down before releasing the lock means a crash in between leaves the
/// lock object behind rather than forgetting it — a prune held off by a
/// stale lock is a nuisance, one that sweeps live chunks is not.
fn release_claim(repository: &Repository, state_dir: &Path, claim: &Claim) -> Result<(), Error> {
store_state_file(
state_dir,
CLAIM_FILE,
&Claim {
lock: None,
epoch: claim.epoch,
},
)?;
if let Some(lock) = claim.lock {
repository.release_lock(lock)?;
}
Ok(())
}
/// Puts every file entry in the holding queue back on the chunk queue,
/// so its content is stored again rather than assumed to still be there.
///
/// Exactly as many items are taken as the queue held when this started,
/// so the entries put back at its tail are not read again and the pass
/// terminates. Interrupted, it redelivers whatever it had drained and
/// starts over, at worst recording a path twice — which costs a queue
/// item and no uploads.
fn requeue_recorded_files(
entries: &PersistentQueue<Entry>,
chunk: &PersistentQueue<QueuedPath>,
) -> Result<(), Error> {
for _ in 0..entries.pending() {
let Some(lease) = entries.get()? else {
break;
};
let entry = lease.ack();
match entry.kind {
EntryKind::File { .. } => chunk.put(&QueuedPath(entry.path))?,
EntryKind::Directory | EntryKind::Symlink { .. } => entries.put(&entry)?,
}
}
// The puts must outlive the acknowledgements that made room for them
chunk.sync()?;
entries.sync()?;
Ok(())
}
fn hostname() -> String {
std::fs::read_to_string("/etc/hostname")
.map(|name| name.trim().to_string())
.ok()
.filter(|name| !name.is_empty())
.unwrap_or_else(|| "unknown".to_string())
}
/// Removes the partial snapshots that a newer one — or the finished
/// backup — has replaced, keeping `keep` if it is still wanted.
///
/// Only ever called once the replacement is durable, so there is no
/// moment at which the repository holds nothing to restore this backup
/// from. Deleting a record that was never written is not an error.
fn retire_partials(
context: &BackupContext,
state_dir: &Path,
stored: &mut StoredManifest,
keep: Option<SnapshotId>,
) -> Result<(), Error> {
let superseded: Vec<SnapshotId> = stored
.published
.iter()
.copied()
.filter(|id| Some(*id) != keep)
.collect();
if superseded.is_empty() {
return Ok(());
}
for id in superseded {
context.repository.delete_snapshot(id)?;
}
stored.published.retain(|id| Some(*id) == keep);
store_state_file(state_dir, MANIFEST_FILE, stored)?;
Ok(())
}
/// One run of the entries queue into manifest chunks: what a checkpoint
/// publishes, and what finalization appends.
struct Segment {
chunks: Vec<ChunkId>,
entries: u64,
content_bytes: u64,
}
/// Drains the entries queue into manifest chunks.
///
/// Entries are deduplicated by path within the segment (keeping the
/// first occurrence): a resumed run may have scanned a directory twice,
/// and either version of a twice-recorded object is a valid read-time
/// observation. Duplicates *across* segments are not caught — the
/// segments are already stored — but they cost only their own bytes in
/// the manifest, since the file chunks they name deduplicate, and
/// restoring the same entry twice simply writes the same file twice.
///
/// The acknowledgements this drains must not be synced until the chunks
/// it returns are accounted for durably; both callers see to that.
fn store_segment(context: &BackupContext) -> Result<Segment, Error> {
let mut seen: HashSet<[u8; 16]> = HashSet::new();
let mut entry_count = 0u64;
let mut content_bytes = 0u64;
let mut queue_failure: Option<Error> = None;
// Draining is the one worker this stage ever has, so it reports like
// any other: the entry it has in hand names the record row that
// stood empty while the queue was only filling up.
let reporter = Reporter::new(&context.events, Stage::Record, 0);
let mut folded = 0u64;
let entries = std::iter::from_fn(|| {
loop {
let lease = match context.entries.get() {
Ok(Some(lease)) => lease,
Ok(None) => return None,
Err(err) => {
queue_failure = Some(err.into());
return None;
}
};
let entry = lease.ack();
// Folding a large manifest takes a while; name the entry in
// hand and say how much of the holding queue is left,
// starting with the very first one so the row is never
// silent while this runs
if folded.is_multiple_of(FINALIZE_REPORT_EVERY) {
reporter.began(&entry.path);
context.report_status();
}
folded += 1;
if !seen.insert(path_digest(&entry.path)) {
continue;
}
entry_count += 1;
if let EntryKind::File { len, .. } = &entry.kind {
content_bytes += len;
}
return Some(entry);
}
});
let chunks = context.repository.store_manifest(entries)?;
reporter.idle();
if let Some(error) = queue_failure {
return Err(error);
}
Ok(Segment {
chunks,
entries: entry_count,
content_bytes,
})
}
/// Appends the last segment and stores the run's finished snapshot.
///
/// The snapshot is complete however the run ended. Coverage is about
/// the run, not the tree: a backup that stopped short on purpose is as
/// finished as one that walked to the end, and what it stored is its
/// final word rather than a stopgap waiting to be superseded. (An
/// interrupted run's checkpoints are the partial ones.)
///
/// Returns `None` when there is nothing to publish, which only an early
/// finish can arrive at — a run that had recorded nothing before it was
/// told to stop, whose snapshot would assert that an unvisited tree is
/// empty. A run that really did walk the whole tree and found nothing
/// says so, and that is worth storing.
///
/// The tail is deliberately not written to the state file: the entry
/// acknowledgements it drains are only synced by the caller, after the
/// snapshot record is stored, so an interruption anywhere in
/// finalization redoes it in full from state that never mentioned the
/// tail.
fn finalize(
context: &BackupContext,
stored: &StoredManifest,
ending: Ending,
) -> Result<Option<SnapshotId>, Error> {
context.emit(Event::Finalizing);
let segment = store_segment(context)?;
let mut manifest = stored.chunks.clone();
manifest.extend(segment.chunks);
if manifest.is_empty() && ending == Ending::Short {
return Ok(None);
}
let id = context
.repository
.snapshot_id(&context.root, &manifest, Coverage::Complete);
context.repository.store_snapshot(&Snapshot::new(
id,
context.root.clone(),
manifest,
stored.entries + segment.entries,
stored.content_bytes + segment.content_bytes,
))?;
Ok(Some(id))
}
/// Where a run's work stopped, which is all finalization needs to know
/// about how it ended.
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Ending {
/// The whole tree was walked.
Whole,
/// The run stopped with work still in its queues.
Short,
}
/// A fixed-size digest of a path for the deduplication set: 128 bits of
/// SHA-256, small enough to hold millions in memory and wide enough that
/// a collision (which would silently drop an entry) will not happen.
fn path_digest(path: &Path) -> [u8; 16] {
let digest = Sha256::digest(path.as_os_str().as_encoded_bytes());
digest[..16].try_into().expect("SHA-256 yields 32 bytes")
}
fn same_device(root: Option<u64>, child: Option<u64>) -> bool {
match (root, child) {
(Some(root), Some(child)) => root == child,
// Without device information there is no boundary to detect
_ => true,
}
}
#[cfg(unix)]
fn device_of(metadata: &Metadata) -> Option<u64> {
use std::os::unix::fs::MetadataExt as _;
Some(metadata.dev())
}
#[cfg(not(unix))]
fn device_of(_metadata: &Metadata) -> Option<u64> {
None
}
#[cfg(unix)]
fn mode_of(metadata: &Metadata) -> Option<u32> {
use std::os::unix::fs::MetadataExt as _;
Some(metadata.mode() & 0o7777)
}
#[cfg(not(unix))]
fn mode_of(_metadata: &Metadata) -> Option<u32> {
None
}
fn mtime_of(metadata: &Metadata) -> Option<i64> {
let modified = metadata.modified().ok()?;
let since_epoch = modified.duration_since(std::time::UNIX_EPOCH).ok()?;
i64::try_from(since_epoch.as_nanos()).ok()
}
#[cfg(test)]
mod tests {
use super::*;
/// A run that died between publishing its snapshot and retiring the
/// partials it replaced must not, when it is run again to clean up,
/// delete the very snapshot it is about to report — however the
/// list of superseded partials came to name it. Nothing should put
/// it there, since a finished snapshot and a partial one covering
/// the same manifest have different identities; deleting a
/// repository's only snapshot is not a mistake to leave resting on
/// that alone.
#[test]
fn tidying_up_after_a_finished_run_spares_the_snapshot_it_published() {
let dir = tempfile::tempdir().unwrap();
let root = dir.path().join("tree");
let state_dir = dir.path().join("state");
std::fs::create_dir_all(&root).unwrap();
std::fs::create_dir_all(&state_dir).unwrap();
let repository = Arc::new(
Repository::create(
Box::new(repository::LocalBackend::new(dir.path().join("repo")).unwrap()),
"test-password",
)
.unwrap(),
);
// A published snapshot, and a state directory whose list of
// partials to retire names it.
let entry = Entry {
path: PathBuf::from("thing"),
kind: EntryKind::Directory,
mode: None,
mtime: None,
};
let chunks = repository.store_manifest(std::iter::once(entry)).unwrap();
let id = repository.snapshot_id(&root, &chunks, Coverage::Complete);
repository
.store_snapshot(&Snapshot::new(id, root.clone(), chunks.clone(), 1, 0))
.unwrap();
store_state_file(
&state_dir,
MANIFEST_FILE,
&StoredManifest {
chunks,
entries: 1,
content_bytes: 0,
published: vec![id],
},
)
.unwrap();
store_state_file(
&state_dir,
CLAIM_FILE,
&Claim {
lock: None,
epoch: repository.epoch().unwrap(),
},
)
.unwrap();
Marker::Completed {
snapshot: id,
finished_early: Some(EarlyFinish::Requested),
}
.store(&state_dir)
.unwrap();
let outcome = run(
repository.clone(),
BackupOptions::new(&root, &state_dir),
crossbeam_channel::unbounded().0,
Arc::new(AtomicBool::new(false)),
)
.unwrap();
assert!(
matches!(
outcome,
Outcome::FinishedEarly {
snapshot: Some(published),
reason: EarlyFinish::Requested,
} if published == id
),
"{outcome:?}"
);
assert_eq!(
repository.snapshots().unwrap(),
vec![id],
"the snapshot it reported is gone from the repository"
);
assert!(!state_dir.exists());
}
}