fix(delay-updates): unique staging dir and implied --delete-after ordering (#317)

This commit is contained in:
2026-09-24 01:18:13 +02:00
parent f7d5dda93b
commit 96e02f52c0
14 changed files with 429 additions and 104 deletions
+118 -36
View File
@@ -8,13 +8,49 @@
#include <errno.h>
#include <fcntl.h>
#include <libgen.h>
#include <stdatomic.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/file.h>
#include <sys/stat.h>
#include <time.h>
#include <unistd.h>
/* Process-wide counter so two staging contexts created in the same process (or
within the same clock tick) can never pick the same name. */
static unsigned long long delay_updates_next_sequence(void) {
static atomic_ullong sequence;
return atomic_fetch_add_explicit(&sequence, 1, memory_order_relaxed);
}
/* Build the per-run staging directory basename: the reserved prefix plus the
pid and an entropy token. A fixed name could collide with a genuine
destination entry; the token makes such a collision vanishingly unlikely and,
if it ever happens, prepare() refuses to touch the existing directory. */
static char* delay_updates_make_staging_name(void) {
unsigned long long entropy = 0;
int fd = open("/dev/urandom", O_RDONLY | O_CLOEXEC);
if (fd >= 0) {
ssize_t got = read(fd, &entropy, sizeof(entropy));
close(fd);
if (got != (ssize_t)sizeof(entropy))
entropy = 0;
}
if (entropy == 0)
entropy = ((unsigned long long)time(NULL) << 20) ^ ((unsigned long long)getpid() << 8) ^
delay_updates_next_sequence();
int length = snprintf(NULL, 0, DELAY_UPDATES_STAGING_DIR ".%ld.%llx", (long)getpid(), entropy);
if (length < 0)
return NULL;
char* name = malloc((size_t)length + 1);
if (!name)
return NULL;
snprintf(name, (size_t)length + 1, DELAY_UPDATES_STAGING_DIR ".%ld.%llx", (long)getpid(),
entropy);
return name;
}
DelayUpdatesContext* delay_updates_context_create(const char* root_directory) {
if (!root_directory)
return NULL;
@@ -26,8 +62,15 @@ DelayUpdatesContext* delay_updates_context_create(const char* root_directory) {
free(context);
return NULL;
}
context->staging_root = path_cat(root_directory, DELAY_UPDATES_STAGING_DIR);
context->staging_name = delay_updates_make_staging_name();
if (!context->staging_name) {
free(context->root_directory);
free(context);
return NULL;
}
context->staging_root = path_cat(root_directory, context->staging_name);
if (!context->staging_root) {
free(context->staging_name);
free(context->root_directory);
free(context);
return NULL;
@@ -39,6 +82,7 @@ DelayUpdatesContext* delay_updates_context_create(const char* root_directory) {
context->lock_fd = -1;
if (mtx_init(&context->mutex, mtx_plain) != thrd_success) {
free(context->staging_root);
free(context->staging_name);
free(context->root_directory);
free(context);
return NULL;
@@ -54,6 +98,7 @@ void delay_updates_context_destroy(DelayUpdatesContext* context) {
close(context->lock_fd);
context->lock_fd = -1;
free(context->staging_root);
free(context->staging_name);
free(context->root_directory);
for (size_t i = 0; i < context->count; i++) {
free(context->entries[i].staged_path);
@@ -125,48 +170,81 @@ bool delay_updates_prepare(DelayUpdatesContext* context) {
return false;
if (context->prepared)
return true;
int fd = file_open_private_dir(context->staging_root);
if (fd < 0) {
/* Create the per-run staging directory with O_EXCL semantics. The name is
unique to this transfer, so if the path already exists it is NOT ours:
either a genuine destination entry that happens to share the name or a
leftover from another session. Refuse rather than wipe it -- the old
fixed-name design could destroy a real destination entry. A crash
leftover is never reused (the next run picks a fresh name). */
char* leaf = NULL;
int parent_fd = file_open_secure_parent(context->staging_root, &leaf, true);
if (parent_fd < 0) {
int saved_errno = errno;
char* escaped = output_escape(context->staging_root, false);
log_message(LOG_LEVEL_ERROR, "could not create --delay-updates staging directory '%s': %s",
escaped ? escaped : "<allocation failed>", strerror(saved_errno));
free(escaped);
free(leaf);
return false;
}
/* Hold an exclusive advisory lock on the staging directory for the whole
transfer. The staging directory name is fixed, so two simultaneous
delayed transfers to the same destination root would otherwise share it
and destroy each other's staged files. The lock makes the second session
fail cleanly instead of corrupting the first. The lock is released when
the context (and its file descriptor) is destroyed. */
int fd = openat(parent_fd, leaf, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (fd >= 0) {
close(fd);
close(parent_fd);
char* escaped = output_escape(context->staging_root, false);
log_message(LOG_LEVEL_ERROR,
"--delay-updates staging directory '%s' already exists and is not owned by this "
"transfer; refusing to overwrite it",
escaped ? escaped : "<allocation failed>");
free(escaped);
free(leaf);
return false;
}
if (errno != ENOENT) {
int saved_errno = errno;
close(parent_fd);
char* escaped = output_escape(context->staging_root, false);
log_message(LOG_LEVEL_ERROR, "could not open --delay-updates staging directory '%s': %s",
escaped ? escaped : "<allocation failed>", strerror(saved_errno));
free(escaped);
free(leaf);
return false;
}
if (mkdirat(parent_fd, leaf, 0700) != 0) {
int saved_errno = errno;
close(parent_fd);
char* escaped = output_escape(context->staging_root, false);
log_message(LOG_LEVEL_ERROR, "could not create --delay-updates staging directory '%s': %s",
escaped ? escaped : "<allocation failed>", strerror(saved_errno));
free(escaped);
free(leaf);
return false;
}
fd = openat(parent_fd, leaf, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
close(parent_fd);
free(leaf);
if (fd < 0) {
int saved_errno = errno;
char* escaped = output_escape(context->staging_root, false);
log_message(LOG_LEVEL_ERROR, "could not open --delay-updates staging directory '%s': %s",
escaped ? escaped : "<allocation failed>", strerror(saved_errno));
free(escaped);
return false;
}
/* Keep the exclusive advisory lock as defense in depth: the unique name
already prevents two sessions from sharing a staging directory, but the
lock also catches an improbable same-name collision that raced between the
existence check above and the open. */
if (flock(fd, LOCK_EX | LOCK_NB) != 0) {
int saved_errno = errno;
close(fd);
if (saved_errno == EWOULDBLOCK || saved_errno == EAGAIN) {
char* escaped = output_escape(context->staging_root, false);
log_message(LOG_LEVEL_ERROR,
"another --delay-updates transfer to '%s' is already in progress; refusing to "
"share the staging directory",
escaped ? escaped : "<allocation failed>");
free(escaped);
} else {
log_message(LOG_LEVEL_ERROR, "could not lock --delay-updates staging directory '%s': %s",
context->staging_root, strerror(saved_errno));
}
char* escaped = output_escape(context->staging_root, false);
log_message(LOG_LEVEL_ERROR, "could not lock --delay-updates staging directory '%s': %s",
escaped ? escaped : "<allocation failed>", strerror(saved_errno));
free(escaped);
return false;
}
context->lock_fd = fd;
/* Only now, with exclusive ownership, wipe leftovers from an interrupted
earlier transfer; this can never race with a live session. */
bool ok = delay_wipe_dir_fd(fd);
if (!ok) {
log_message(LOG_LEVEL_ERROR, "could not clear stale --delay-updates staging files under '%s'",
context->staging_root);
close(context->lock_fd);
context->lock_fd = -1;
return false;
}
context->prepared = true;
return true;
}
@@ -264,12 +342,16 @@ static bool delay_publish_entry(DelayUpdatesContext* context, const Config* conf
const StagedFileEntry* entry) {
if (!delay_publish_backup(context, config, entry))
return false;
/* --force: an incoming regular file/symlink may replace a destination
DIRECTORY (possibly non-empty). The immediate-install path handles this in
file_receive; a --delay-updates run stages elsewhere and only discovers the
blocking directory here, so clear it before the rename (rsync's
"could not make way for new regular file" without --force). */
if (config && config->force_delete && file_directory_exists_secure(entry->final_path)) {
/* An incoming regular file/symlink may replace a destination DIRECTORY that
blocks it. rsync removes the blocker recursively when --delete or --force
is active (its generator's "make way" deletion), and a --delay-updates run
stages elsewhere so it only discovers the blocker here. FastSync's
immediate-install path clears it too; without --delete/--force a non-empty
blocker fails the run (rsync's "could not make way for new regular file").
use_delete is gated by the server --allow-delete policy, so a client can
never use this to bypass deletion authorization. */
if (config && (config->force_delete || config->use_delete) &&
file_directory_exists_secure(entry->final_path)) {
if (!file_remove_tree_secure(entry->final_path)) {
char* escaped = output_escape(entry->final_path, false);
log_message(LOG_LEVEL_ERROR, "could not remove destination directory blocking '%s': %s",
+6 -1
View File
@@ -24,6 +24,7 @@ typedef struct {
shared with the publish/cleanup phase that runs after the threads join. */
typedef struct DelayUpdatesContext {
char* root_directory; /* receive root the staging dir lives under */
char* staging_name; /* per-run unique staging dir basename */
char* staging_root; /* root_directory/<staging dir name> */
mtx_t mutex;
StagedFileEntry* entries;
@@ -33,7 +34,11 @@ typedef struct DelayUpdatesContext {
int lock_fd; /* advisory exclusive flock held on the staging dir, or -1 */
} DelayUpdatesContext;
/* Name of the private staging subdirectory created under the receive root. */
/* Reserved prefix for the private staging subdirectory created under the
receive root. The actual directory name is per-run unique (the prefix plus a
pid/entropy token) so it can never clobber a genuine destination entry that
happens to share the name; the bare prefix is still what a --backup-dir must
not collide with. */
#define DELAY_UPDATES_STAGING_DIR ".fastsync-stage"
/* True when `dir` (ignoring a trailing "/") is the reserved staging directory
+26 -1
View File
@@ -39,12 +39,31 @@ FilterAction delete_protect_verdict(const DeleteProtectRules* protect, const cha
const char* leaf, bool is_dir) {
if (!protect)
return FILTER_ACTION_NONE;
/* rsync protects its own --backup files from the delete pass: a name ending
in the backup suffix is never an extra. Checked before the filter rules so
an explicit exclude cannot be bypassed (the suffix is always a shield). */
if (protect->backup_suffix && protect->backup_suffix[0] != '\0') {
size_t name_len = strlen(leaf);
size_t suffix_len = strlen(protect->backup_suffix);
if (name_len > suffix_len &&
strcmp(leaf + (name_len - suffix_len), protect->backup_suffix) == 0)
return FILTER_ACTION_PROTECT;
}
FilterAction action = filter_dir_rules_apply_side(protect->dir_rules, rel_path, leaf, is_dir);
if (action != FILTER_ACTION_NONE)
return action;
return filter_rules_apply_side(protect->base_rules, rel_path, leaf, is_dir, FILTER_SIDE_RECEIVER);
}
const char* delete_backup_suffix(const Config* config) {
if (!config || !config->backup || config->ignore_existing)
return NULL;
const char* suffix = config->suffix ? config->suffix : "~";
if (!suffix[0] || strchr(suffix, '/'))
return NULL;
return suffix;
}
/* Classify a removed entry from its st_mode for the per-type delete counters. */
DeleteEntryType delete_entry_type_of_mode(mode_t mode) {
if (S_ISDIR(mode))
@@ -631,7 +650,13 @@ bool delete_skips_build(const Config* config, const ArrayList* protected_paths,
}
int idx = 0;
if (config->delay_updates) {
out->entries[idx].prefix = DELAY_UPDATES_STAGING_DIR;
/* Protect this transfer's actual (per-run unique) staging directory. The
runtime name is only known to the receiver-side context; fall back to the
reserved prefix for a context that was never created (e.g. a dry run). */
const char* staging_name = (config->delay_context && config->delay_context->staging_name)
? config->delay_context->staging_name
: DELAY_UPDATES_STAGING_DIR;
out->entries[idx].prefix = staging_name;
out->entries[idx].top_level_only = true;
idx++;
}
+10
View File
@@ -37,6 +37,11 @@ typedef enum {
typedef struct {
const FilterRuleList* base_rules;
const FilterRuleList* dir_rules;
/* When non-NULL and non-empty, a destination entry whose name ends with this
suffix is protected from deletion. rsync never treats a --backup file as
an extra, so a backup created at --delay-updates publication (or a
pre-existing one) survives the delete-after pass. */
const char* backup_suffix;
} DeleteProtectRules;
/* rsync's first-match-wins receiver verdict for one candidate extra: the
@@ -48,6 +53,11 @@ typedef struct {
FilterAction delete_protect_verdict(const DeleteProtectRules* protect, const char* rel_path,
const char* leaf, bool is_dir);
/* The backup suffix the delete walker must shield from deletion, or NULL when
--backup is inactive or the configured suffix is unusable (empty, or holding
a path separator). Matches the suffix file_save uses for backups. */
const char* delete_backup_suffix(const Config* config);
/* One protected entry for the delete walker. When top_level_only is true the
prefix is skipped only as a DIRECT child of dest_root (the --delay-updates
staging directory, which must not hide genuine extras inside a nested
+4 -2
View File
@@ -170,7 +170,8 @@ static bool delete_extras_budgeted_observed(const Config* config, const DeleteMa
size_t deleted = 0;
size_t skipped = 0;
DeleteProtectRules protect = {.base_rules = config->protect_rules,
.dir_rules = manifest->per_dir_rules};
.dir_rules = manifest->per_dir_rules,
.backup_suffix = delete_backup_suffix(config)};
DeleteWalkResult result = delete_extras_limited_observed(
config->receive_root_directory, manifest->keeps, manifest->dirs, remaining, skips.entries,
skips.count, &protect, &deleted, &skipped, observer, observer_context);
@@ -398,7 +399,8 @@ bool manifest_would_delete_list(const Config* config, const DeleteManifest* mani
if (!delete_skips_build(config, manifest->protected, NULL, true, &skips))
return false;
DeleteProtectRules protect = {.base_rules = config->protect_rules,
.dir_rules = manifest->per_dir_rules};
.dir_rules = manifest->per_dir_rules,
.backup_suffix = delete_backup_suffix(config)};
bool ok = delete_extras_list(config->receive_root_directory, manifest->keeps, manifest->dirs,
skips.entries, skips.count, &protect, out, count_out);
delete_skips_free(&skips);
+1
View File
@@ -801,6 +801,7 @@ static bool build_plan_skips(const Config* config, const DeletePlanSession* sess
PlanSkips* out) {
out->protect.base_rules = config->protect_rules;
out->protect.dir_rules = session->per_dir_rules;
out->protect.backup_suffix = delete_backup_suffix(config);
/* The per-directory plan walk keeps each basis path verbatim (it does not
convert an absolute under-root path to its root-relative form, unlike the
whole-tree commit walk). */