Merge feat/p2-delay-updates: implement --delay-updates

Receiver stages writes and publishes only after the full transfer succeeds
(single and -m, before the outcomes frame). delete walker skips staging;
backup-dir name reserved; flock serializes concurrent delayed sessions;
protocol bump to 2.6.0. c-review REQUEST CHANGES -> blockers fixed; PR #264.
This commit is contained in:
2026-09-06 14:19:38 +02:00
21 changed files with 1221 additions and 33 deletions
+3 -3
View File
@@ -6,11 +6,11 @@ This document maps rsync's full feature set to FastSync's current implementation
| Status | Count | Description |
|--------|-------|-------------|
| ✅ Implemented | 57 | Feature works end-to-end |
| ✅ Implemented | 58 | Feature works end-to-end |
| 🔀 Alt Arg | 3 | Functionality exists but under different flag/semantics |
| ⚠️ Partial | 5 | Flag parsed/stored but behavior incomplete |
| 🔄 Compatibility No-op | 1 | Flag is accepted for CLI compatibility but has no effect |
| ❌ Not Implemented | 81 | Flag not recognized or no behavior |
| ❌ Not Implemented | 80 | Flag not recognized or no behavior |
| **Total** | **147** | |
---
@@ -96,7 +96,7 @@ This document maps rsync's full feature set to FastSync's current implementation
| `-b`, `--backup` | Make backups of overwritten files | ✅ Implemented | Backup before overwrite |
| `--backup-dir=DIR` | Backup directory hierarchy | ✅ Implemented | `backup_dir` config field |
| `--suffix=SUFFIX` | Backup suffix (default ~) | ✅ Implemented | `suffix` config field |
| `--delay-updates` | Put updated files in place at end | ❌ Not Implemented | |
| `--delay-updates` | Put updated files in place at end | ✅ Implemented | Successfully received files are staged under a private 0700 `.fastsync-stage` dir inside the receive root and atomically renamed into their final destinations only after the whole transfer (manifest/delete handling included) succeeds, just before the success/outcome frame is sent. The delete walker deliberately skips the staging dir at the receive root, so `--delete` removes genuine extras but never the staged files (deletion runs before publication; rsync's delete-after ordering is not implemented). `--existing`/`--ignore-existing`/`--update` decide against the final destination path at stage time; `--backup` moves the old file aside at publication. Incompatible with `--inplace` and with `--backup-dir=.fastsync-stage` (the internal staging name is reserved; both are rejected). The staging dir name is fixed, so two simultaneous delayed transfers to the same destination root are serialized with an exclusive advisory lock held for the whole transfer: the second session fails cleanly instead of corrupting the first. Aborting or failing before publication installs nothing and removes the staging tree; a crash between stage and publish leaves staged leftovers that the next delayed run wipes at start (process death releases the lock). A stage→publish failure aborts the transfer (best-effort cleanup of the not-yet-published staged files; already-published files are not rolled back). Works in single-threaded and `-m` modes |
| `-T`, `--temp-dir=DIR` | Create temporary files in DIR | ✅ Implemented | `--temp-dir` only; `-T` stays FastSync's `--timeout` alias. Scratch dir is resolved under the receive root; temp copies use a unique name there and are atomically renamed into place. If the scratch dir and destination are on different filesystems the atomic rename fails with EXDEV and the file save fails, which aborts the whole transfer (FastSync has no per-file skip/resume on a save error; rsync's non-atomic copy fallback is deliberately not used). `--inplace` and `--partial-dir` writes bypass the scratch dir |
## 7. Deletion
+1
View File
@@ -396,6 +396,7 @@ static const OptionEntry OPTION_TABLE[] = {
{"--log-file-format", NULL, OPT_STRING, offsetof(Config, log_file_format)},
{"--existing", NULL, OPT_FLAG, offsetof(Config, existing)},
{"--ignore-existing", NULL, OPT_FLAG, offsetof(Config, ignore_existing)},
{"--delay-updates", NULL, OPT_FLAG, offsetof(Config, delay_updates)},
{"--chmod", NULL, OPT_STRING, offsetof(Config, chmod_spec)},
{"--dirs", "--old-dirs", OPT_UNSUPPORTED, 0},
{"--old-d", NULL, OPT_UNSUPPORTED, 0},
+11
View File
@@ -1,4 +1,5 @@
#include "client_validation.h"
#include "delay_updates.h"
#include "log.h"
#include "usage.h"
#include <stdio.h>
@@ -60,5 +61,15 @@ bool validate_config(const Config* config) {
return false;
}
}
if (config->delay_updates && config->inplace) {
log_message(LOG_LEVEL_ERROR, "--delay-updates does not work with --inplace");
return false;
}
if (config->delay_updates && delay_updates_staging_name_conflict(config->backup_dir)) {
log_message(LOG_LEVEL_ERROR,
"--backup-dir is reserved when --delay-updates is active (used for the internal "
"staging directory)");
return false;
}
return true;
}
+1
View File
@@ -25,6 +25,7 @@ void print_usage(void) {
printf(" -8, --8-bit-output Leave high-bit characters unescaped in output\n");
printf(" --delete Delete files on receiver not in source\n");
printf(" --ignore-existing Skip files that already exist on receiver\n");
printf(" --delay-updates Put updated files into place only at the end of transfer\n");
printf(
" --dirs, --old-dirs, --old-d Transfer directories without recursing (not implemented)\n");
printf(" --del Alias for --delete-during (not implemented)\n");
+15
View File
@@ -1,6 +1,8 @@
#include "receiver.h"
#include "chunk.h"
#include "config.h"
#include "delay_updates.h"
#include "file_receive.h"
#include "log.h"
#include "metadata.h"
@@ -217,6 +219,17 @@ static bool receiver_save_file(File* file, void* context_pointer) {
static bool receiver_send_success_frame(int fd, void* context_pointer) {
ReceiverSaveContext* context = context_pointer;
/* --delay-updates: the whole protocol stream (including manifest/delete
handling, which ran inside receiver_process) has succeeded and every
staged file was fully written. Publish them atomically now, before the
success/outcome frame tells a --remove-source-files sender it may delete
its sources. */
if (context->config->delay_updates && context->config->delay_context) {
if (!delay_updates_publish(context->config->delay_context, context->config)) {
send_status(fd, STATUS_ERROR);
return false;
}
}
return receiver_send_final_success(fd, context->config, &context->outcomes);
}
@@ -224,6 +237,8 @@ int receiver_receive_files(Config* config, int file_descriptor) {
ReceiverSaveContext context = {.config = config, .outcomes = {0}};
ReceiverSink sink = {receiver_save_file, &context, true, true, receiver_send_success_frame};
int ret = receiver_process(config, file_descriptor, &sink);
if (ret != 0 && config->delay_updates && config->delay_context)
delay_updates_cleanup(config->delay_context);
receiver_outcomes_destroy(&context.outcomes);
return ret;
}
+30 -1
View File
@@ -1,4 +1,5 @@
#include "config.h"
#include "delay_updates.h"
#include "file.h"
#include "log.h"
#include "multiprocessing.h"
@@ -164,6 +165,20 @@ void handler(int file_descriptor) {
return;
}
config->use_delete = config->use_delete && allow_delete;
/* A --delay-updates transfer stages under a private 0700 directory inside
the receive root. Create it up front (wiping leftovers of any previously
interrupted delayed transfer) so a fully-skipped run also starts clean. */
if (config->delay_updates) {
config->delay_context = delay_updates_context_create(config->receive_root_directory);
if (!config->delay_context || !delay_updates_prepare(config->delay_context)) {
log_message(LOG_LEVEL_ERROR, "Failed to initialize --delay-updates staging area");
delay_updates_cleanup(config->delay_context);
config_delete(config);
close(file_descriptor);
protocol_session_unbind();
return;
}
}
if (config->use_multithreading) {
Queue* q = queue_create(100, file_destroy);
if (q == NULL) {
@@ -214,14 +229,28 @@ void handler(int file_descriptor) {
thrd_join(receiver, &receiver_result);
thrd_join(writer, &writer_result);
bool transfer_ok = receiver_result == thrd_success && writer_result == thrd_success;
if (transfer_ok) {
/* --delay-updates: receive_thread has finished the whole protocol stream
(including manifest/delete handling) and write_thread has drained its
queue, so every staged file is complete. Publish atomically before the
success/outcome frame so a --remove-source-files sender only learns of
files that were actually installed. */
if (config->delay_updates && config->delay_context &&
!delay_updates_publish(config->delay_context, config)) {
transfer_ok = false;
}
}
if (transfer_ok) {
if (!receiver_send_final_success(file_descriptor, config, &context->outcomes))
transfer_ok = false;
} else {
send_status(file_descriptor, STATUS_ERROR);
}
if (!transfer_ok)
if (!transfer_ok) {
log_message(LOG_LEVEL_ERROR, "Transfer failed");
if (config->delay_updates && config->delay_context)
delay_updates_cleanup(config->delay_context);
}
pipeline_context_receiver_destroy(context);
} else {
if (receiver_receive_files(config, file_descriptor) != 0)
+19 -7
View File
@@ -1,5 +1,6 @@
#include "config.h"
#include "chmod.h"
#include "delay_updates.h"
#include "delta.h"
#include "log.h"
#include "protocol.h"
@@ -84,6 +85,7 @@ static void config_set_defaults(Config* config) {
config->ignore_existing = false;
config->update = false;
config->inplace = false;
config->delay_updates = false;
config->use_fsync = false;
config->append = false;
config->append_verify = false;
@@ -119,6 +121,7 @@ static void config_set_defaults(Config* config) {
config->skip_compress_suffixes = NULL;
config->skip_compress_count = 0;
config->skip_compress_set = false;
config->delay_context = NULL;
}
static bool valid_wire_bool(int value) {
@@ -152,6 +155,8 @@ static bool validate_received_config(const Config* config) {
valid_wire_bool(config->use_fsync) && valid_wire_bool(config->append_verify) &&
valid_wire_bool(config->delete_excluded) && valid_wire_bool(config->delete_after) &&
valid_wire_bool(config->relative) && valid_wire_bool(config->prune_empty_dirs) &&
valid_wire_bool(config->delay_updates) && !(config->delay_updates && config->inplace) &&
!(config->delay_updates && delay_updates_staging_name_conflict(config->backup_dir)) &&
valid_wire_bool(config->partial) && valid_wire_bool(config->delete_before) &&
valid_wire_bool(config->checksum) && valid_wire_bool(config->eight_bit_output) &&
!(config->skip_compress_set && config->use_chunk_serialization) &&
@@ -248,6 +253,12 @@ void config_delete(Config* config) {
if (config->filters) {
array_list_delete(config->filters);
}
/* A --delay-updates staging tree is transient receiver state: remove any
leftovers on every exit path (success already emptied it). */
if (config->delay_context)
delay_updates_cleanup(config->delay_context);
delay_updates_context_destroy(config->delay_context);
config->delay_context = NULL;
free(config);
}
@@ -287,10 +298,11 @@ static bool send_file_options(int fd, const Config* c) {
static bool send_selection_options(int fd, const Config* c) {
return send_int(fd, c->ignore_existing) && send_int(fd, c->existing) && send_int(fd, c->update) &&
send_int(fd, c->inplace) && send_int(fd, c->append) && send_int(fd, c->use_fsync) &&
send_int(fd, c->append_verify) && send_int(fd, c->delete_excluded) &&
send_int(fd, c->delete_after) && send_n_data(fd, &c->max_delete, sizeof(c->max_delete)) &&
send_int(fd, c->relative) && send_int(fd, c->prune_empty_dirs);
send_int(fd, c->inplace) && send_int(fd, c->delay_updates) && send_int(fd, c->append) &&
send_int(fd, c->use_fsync) && send_int(fd, c->append_verify) &&
send_int(fd, c->delete_excluded) && send_int(fd, c->delete_after) &&
send_n_data(fd, &c->max_delete, sizeof(c->max_delete)) && send_int(fd, c->relative) &&
send_int(fd, c->prune_empty_dirs);
}
static bool send_skip_compress_options(int fd, const Config* c) {
@@ -382,9 +394,9 @@ static bool receive_file_options(int fd, Config* c) {
}
static bool receive_selection_options(int fd, Config* c) {
bool* flags[] = {&c->ignore_existing, &c->existing, &c->update,
&c->inplace, &c->append, &c->use_fsync,
&c->append_verify, &c->delete_excluded, &c->delete_after};
bool* flags[] = {&c->ignore_existing, &c->existing, &c->update, &c->inplace,
&c->delay_updates, &c->append, &c->use_fsync, &c->append_verify,
&c->delete_excluded, &c->delete_after};
for (size_t i = 0; i < sizeof(flags) / sizeof(flags[0]); i++) {
if (!receive_wire_bool(fd, flags[i]))
return false;
+10 -1
View File
@@ -8,6 +8,10 @@
typedef enum { TRANSPORT_TCP, TRANSPORT_SSH } TransportType;
/* Receiver-side staging state for --delay-updates. Forward-declared here so
Config can carry it; the concrete type lives in delay_updates.h. */
typedef struct DelayUpdatesContext DelayUpdatesContext;
typedef struct Config {
char* version;
char* send_directory;
@@ -91,6 +95,7 @@ typedef struct Config {
bool ignore_existing;
bool update;
bool inplace;
bool delay_updates;
bool use_fsync;
bool append;
bool append_verify;
@@ -147,9 +152,13 @@ typedef struct Config {
char** skip_compress_suffixes;
int skip_compress_count;
bool skip_compress_set;
// Receiver-side runtime staging registry for --delay-updates. Never sent
// over the wire and never set on the sender side.
DelayUpdatesContext* delay_context;
} Config;
#define PROTOCOL_VERSION "2.5.0"
#define PROTOCOL_VERSION "2.6.0"
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
Config* config_create(void);
+338
View File
@@ -0,0 +1,338 @@
#include "delay_updates.h"
#include "config.h"
#include "file.h"
#include "log.h"
#include "utils.h"
#include <dirent.h>
#include <errno.h>
#include <fcntl.h>
#include <libgen.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/file.h>
#include <sys/stat.h>
#include <unistd.h>
DelayUpdatesContext* delay_updates_context_create(const char* root_directory) {
if (!root_directory)
return NULL;
DelayUpdatesContext* context = calloc(1, sizeof(DelayUpdatesContext));
if (!context)
return NULL;
context->root_directory = str_dup(root_directory);
if (!context->root_directory) {
free(context);
return NULL;
}
context->staging_root = path_cat(root_directory, DELAY_UPDATES_STAGING_DIR);
if (!context->staging_root) {
free(context->root_directory);
free(context);
return NULL;
}
context->entries = NULL;
context->count = 0;
context->capacity = 0;
context->prepared = false;
context->lock_fd = -1;
if (mtx_init(&context->mutex, mtx_plain) != thrd_success) {
free(context->staging_root);
free(context->root_directory);
free(context);
return NULL;
}
return context;
}
void delay_updates_context_destroy(DelayUpdatesContext* context) {
if (!context)
return;
mtx_destroy(&context->mutex);
if (context->lock_fd >= 0)
close(context->lock_fd);
context->lock_fd = -1;
free(context->staging_root);
free(context->root_directory);
for (size_t i = 0; i < context->count; i++) {
free(context->entries[i].staged_path);
free(context->entries[i].final_path);
free(context->entries[i].file_path);
}
free(context->entries);
free(context);
}
bool delay_updates_staging_name_conflict(const char* dir) {
if (!dir || !*dir)
return false;
size_t length = strlen(dir);
while (length > 0 && dir[length - 1] == '/')
length--;
size_t reserved_length = strlen(DELAY_UPDATES_STAGING_DIR);
if (length != reserved_length)
return false;
return strncmp(dir, DELAY_UPDATES_STAGING_DIR, length) == 0;
}
/* Recursively delete every entry inside an open directory (never following
symlinks). The directory itself is left in place. Mirrors the fd-relative
walk used by the delete code so a symlink planted inside the staging tree
can never redirect removal outside of it. */
static bool delay_wipe_dir_fd(int dirfd) {
int scanfd = dup(dirfd);
if (scanfd < 0)
return false;
DIR* dir = fdopendir(scanfd);
if (!dir) {
close(scanfd);
return false;
}
bool operation_ok = true;
const struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
struct stat st;
if (fstatat(dirfd, entry->d_name, &st, AT_SYMLINK_NOFOLLOW) != 0) {
if (errno != ENOENT)
operation_ok = false;
continue;
}
if (S_ISDIR(st.st_mode)) {
int childfd = openat(dirfd, entry->d_name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool child_removed = false;
if (childfd >= 0) {
child_removed = delay_wipe_dir_fd(childfd);
close(childfd);
} else if (errno != ENOENT) {
operation_ok = false;
}
if (child_removed && unlinkat(dirfd, entry->d_name, AT_REMOVEDIR) != 0 && errno != ENOENT)
operation_ok = false;
} else {
if (unlinkat(dirfd, entry->d_name, 0) != 0 && errno != ENOENT)
operation_ok = false;
}
}
closedir(dir);
return operation_ok;
}
bool delay_updates_prepare(DelayUpdatesContext* context) {
if (!context)
return false;
if (context->prepared)
return true;
int fd = file_open_private_dir(context->staging_root);
if (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);
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. */
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));
}
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;
}
bool delay_updates_record(DelayUpdatesContext* context, const char* staged_path,
const char* final_path, const char* file_path) {
if (!context || !staged_path || !final_path || !file_path)
return false;
char* staged_copy = str_dup(staged_path);
char* final_copy = str_dup(final_path);
char* file_copy = str_dup(file_path);
if (!staged_copy || !final_copy || !file_copy) {
free(staged_copy);
free(final_copy);
free(file_copy);
return false;
}
mtx_lock(&context->mutex);
bool ok = true;
if (context->count == context->capacity) {
size_t new_capacity = context->capacity == 0 ? 64 : context->capacity * 2;
if (new_capacity < context->capacity) {
ok = false;
} else {
StagedFileEntry* grown = realloc(context->entries, new_capacity * sizeof(StagedFileEntry));
if (!grown) {
ok = false;
} else {
context->entries = grown;
context->capacity = new_capacity;
}
}
}
if (ok) {
context->entries[context->count].staged_path = staged_copy;
context->entries[context->count].final_path = final_copy;
context->entries[context->count].file_path = file_copy;
context->count++;
}
mtx_unlock(&context->mutex);
if (!ok) {
free(staged_copy);
free(final_copy);
free(file_copy);
}
return ok;
}
/* Move an existing final destination file aside before the staged replacement
is installed. Deferred from stage time so the final destination is not
modified until publication. Mirrors the immediate-mode backup logic. */
static bool delay_publish_backup(const DelayUpdatesContext* context, const Config* config,
const StagedFileEntry* entry) {
bool backup_enabled = config && config->backup && !config->ignore_existing;
if (!backup_enabled)
return true;
const char* backup_suffix = (config && config->suffix) ? config->suffix : "~";
struct stat backup_stat;
if (!file_stat_secure(entry->final_path, &backup_stat))
return true; /* nothing to back up */
char* backup_path = NULL;
if (config->backup_dir) {
char* confined_backup = path_cat(context->root_directory, config->backup_dir);
if (!confined_backup)
return false;
backup_path = path_cat(confined_backup, entry->file_path);
free(confined_backup);
} else {
size_t path_len = strlen(entry->final_path);
size_t suffix_len = strlen(backup_suffix);
if (path_len > SIZE_MAX - suffix_len - 1)
return false;
backup_path = malloc(path_len + suffix_len + 1);
if (backup_path) {
memcpy(backup_path, entry->final_path, path_len);
memcpy(backup_path + path_len, backup_suffix, suffix_len + 1);
}
}
if (!backup_path)
return false;
char* parent_copy = str_dup(backup_path);
if (!parent_copy || !file_ensure_directory_secure(dirname(parent_copy))) {
free(parent_copy);
free(backup_path);
return false;
}
free(parent_copy);
bool ok = file_rename_secure(entry->final_path, backup_path);
free(backup_path);
return ok;
}
static bool delay_publish_entry(DelayUpdatesContext* context, const Config* config,
const StagedFileEntry* entry) {
if (!delay_publish_backup(context, config, entry))
return false;
if (!file_rename_secure(entry->staged_path, entry->final_path)) {
if (errno == EXDEV) {
char* escaped = output_escape(entry->final_path, false);
log_message(LOG_LEVEL_ERROR,
"staging directory is on a different filesystem than the destination; cannot "
"atomically install file (EXDEV): %s",
escaped ? escaped : "<allocation failed>");
free(escaped);
} else {
char* escaped = output_escape(entry->final_path, false);
log_message(LOG_LEVEL_ERROR, "could not install staged file '%s': %s",
escaped ? escaped : "<allocation failed>", strerror(errno));
free(escaped);
}
return false;
}
return true;
}
/* Remove the staging tree (contents plus the directory itself). Returns true
when nothing is left behind (including the case where it never existed). */
static bool delay_updates_remove_staging_tree(DelayUpdatesContext* context) {
int fd = open(context->staging_root, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (fd < 0)
return errno == ENOENT;
bool ok = delay_wipe_dir_fd(fd);
if (close(fd) != 0)
ok = false;
if (ok && rmdir(context->staging_root) != 0 && errno != ENOENT)
ok = false;
return ok;
}
bool delay_updates_publish(DelayUpdatesContext* context, const Config* config) {
if (!context)
return false;
mtx_lock(&context->mutex);
bool ok = true;
for (size_t i = 0; i < context->count; i++) {
if (!delay_publish_entry(context, config, &context->entries[i])) {
ok = false;
break;
}
}
mtx_unlock(&context->mutex);
/* Renaming files out of the staging tree leaves the mirrored directories
behind, and a mid-publish failure leaves the remaining staged files.
Remove whatever is left so a later run starts from a clean staging area
and no staged content can linger after a failed publish. If that cleanup
fails, tell the operator: a stale staging directory would otherwise
silently accumulate and make the next transfer's prepare-wipe fail. */
if (!delay_updates_remove_staging_tree(context)) {
log_message(LOG_LEVEL_WARNING,
"could not fully remove --delay-updates staging directory '%s' after publish; a "
"later --delay-updates transfer to this destination will try to clear it",
context->staging_root);
}
return ok;
}
void delay_updates_cleanup(DelayUpdatesContext* context) {
if (!context)
return;
/* Only a context that gained exclusive ownership may touch the shared
staging directory. If prepare never succeeded (e.g. lock contention with
another live session) the directory belongs to that other session and must
be left alone. */
if (!context->prepared)
return;
delay_updates_remove_staging_tree(context);
}
+68
View File
@@ -0,0 +1,68 @@
#ifndef DELAY_UPDATES_H
#define DELAY_UPDATES_H
#include <stdbool.h>
#include <stddef.h>
#include <threads.h>
/* Forward-declared in config.h; full type needed by file_save_to_disk. */
typedef struct Config Config;
/* One staged file awaiting publication. */
typedef struct {
char* staged_path; /* full path inside the staging tree */
char* final_path; /* full final destination path */
char* file_path; /* the file path as received on the wire */
} StagedFileEntry;
/* Receiver-side --delay-updates staging registry. All successfully written
files land under a private staging directory inside the receive root and are
atomically renamed into their final destination only at the very end of the
transfer. A single PipelineContextReceiver has exactly one writer thread,
but the registry is still mutex-protected so the same object can be safely
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_root; /* root_directory/<staging dir name> */
mtx_t mutex;
StagedFileEntry* entries;
size_t count;
size_t capacity;
bool prepared; /* staging dir created, wiped, and exclusively locked */
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. */
#define DELAY_UPDATES_STAGING_DIR ".fastsync-stage"
/* True when `dir` (ignoring a trailing "/") is the reserved staging directory
name. Used to reject a --backup-dir that would collide with the internal
staging area. */
bool delay_updates_staging_name_conflict(const char* dir);
/* Create an empty staging context rooted below root_directory. Does not touch
the filesystem yet. */
DelayUpdatesContext* delay_updates_context_create(const char* root_directory);
void delay_updates_context_destroy(DelayUpdatesContext* context);
/* Create the private 0700 staging directory (on first call) and wipe any
leftovers from a previously interrupted delayed transfer. Idempotent. */
bool delay_updates_prepare(DelayUpdatesContext* context);
/* Record a fully-written staged file for later publication. Copies all three
paths. Returns false on allocation failure. */
bool delay_updates_record(DelayUpdatesContext* context, const char* staged_path,
const char* final_path, const char* file_path);
/* Atomically rename every staged file into its final destination. Deferred
--backup handling runs immediately before each rename. On any failure the
remaining staged files are removed (best effort); already-published files
are not rolled back. Afterwards the staging tree is removed so a successful
or failed publish leaves no staging leftovers. */
bool delay_updates_publish(DelayUpdatesContext* context, const Config* config);
/* Best-effort removal of every staged file and the staging directory itself.
Safe to call when nothing was staged or after a successful publish. */
void delay_updates_cleanup(DelayUpdatesContext* context);
#endif
+11 -12
View File
@@ -338,23 +338,22 @@ bool file_rename_secure(const char* old_path, const char* new_path) {
return ok;
}
/* Open the configured --temp-dir scratch directory, creating it (and any
missing path components) on demand. scratch_path is expected to already be
confined below the authorized root by the caller; file_open_secure_parent
re-checks that confinement and rejects `..` components, so a scratch
directory can never be created or opened outside the destination root.
Returns an O_DIRECTORY|O_NOFOLLOW fd, or -1 on error. */
static int file_open_scratch_dir(const char* scratch_path) {
if (!scratch_path)
/* Open a private staging/scratch directory, creating it (and any missing path
components) on demand. dir_path is expected to already be confined below
the authorized root by the caller; file_open_secure_parent re-checks that
confinement and rejects `..` components, so a scratch directory can never be
created or opened outside the destination root. The directory itself is
created 0700 so other users cannot race on names inside it. Returns an
O_DIRECTORY|O_NOFOLLOW fd, or -1 on error. */
int file_open_private_dir(const char* dir_path) {
if (!dir_path)
return -1;
char* leaf = NULL;
int parent_fd = file_open_secure_parent(scratch_path, &leaf, true);
int parent_fd = file_open_secure_parent(dir_path, &leaf, true);
if (parent_fd < 0)
return -1;
int fd = openat(parent_fd, leaf, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (fd < 0 && errno == ENOENT) {
/* A scratch directory holds transient working copies only; keep it
private (0700) so other users cannot race on temp names inside it. */
if (mkdirat(parent_fd, leaf, 0700) == 0 || errno == EEXIST)
fd = openat(parent_fd, leaf, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
}
@@ -429,7 +428,7 @@ static bool file_to_disk_secure_impl(const char* path, const void* data,
file is created in the destination directory, exactly as historically. */
int scratch_dirfd = -1;
if (temp_dir) {
scratch_dirfd = file_open_scratch_dir(temp_dir);
scratch_dirfd = file_open_private_dir(temp_dir);
if (scratch_dirfd < 0) {
int saved_errno = errno;
log_message(LOG_LEVEL_ERROR, "could not open --temp-dir scratch directory '%s': %s",
+4
View File
@@ -31,6 +31,10 @@ bool file_destination_is_newer_secure(const char* path, const FileMetadata* meta
int file_open_secure_parent(const char* path, char** leaf_out, bool create_dirs);
bool file_ensure_directory_secure(const char* path);
bool file_rename_secure(const char* old_path, const char* new_path);
/* Open a private 0700 directory (creating it on demand) that must live below
the authorized root. Used for the --temp-dir scratch directory and the
--delay-updates staging directory. */
int file_open_private_dir(const char* dir_path);
/* The file_to_disk_secure* variants write a temporary copy in the destination
directory and atomically rename it over `path`. temp_dir is an absolute,
+86 -2
View File
@@ -12,6 +12,7 @@
#include "compression.h"
#include "config.h"
#include "data.h"
#include "delay_updates.h"
#include "delta.h"
#include "file.h"
#include "log.h"
@@ -26,6 +27,72 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi
return file_save_to_disk_full(root_directory, file, config) != FILE_SAVE_ERROR;
}
/* --delay-updates receiver path: write the file into a private staging tree
below the receive root instead of its final destination, and remember it so
it can be atomically renamed into place only once the whole transfer has
succeeded. Existence/update policies (--existing/--ignore-existing/--update)
are decided against the FINAL destination path at stage time so the run
decides exactly what an immediate (non-delayed) run would decide; the staged
file is then never re-checked at publication. Backups are deferred to
publication so the final destination is untouched until the transfer ends. */
static FileSaveResult file_stage_delayed_update(const char* root_directory,
const char* destination_path, const File* file,
Config* config) {
if (!config)
return FILE_SAVE_ERROR;
bool sparse = config->preserve_sparse;
bool preserve_executability = config->use_executability;
if (config->existing && !file_path_exists_secure(destination_path))
return FILE_SAVE_SKIPPED;
if (config->ignore_existing && file_path_exists_secure(destination_path))
return FILE_SAVE_SKIPPED;
if (config->update && file_destination_is_newer_secure(destination_path, file->metadata))
return FILE_SAVE_SKIPPED;
FileMetadata adjusted_metadata;
const FileMetadata* metadata = file->metadata;
if (metadata && config->chmod_spec && *config->chmod_spec) {
adjusted_metadata = *metadata;
if (!chmod_apply(adjusted_metadata.mode, config->chmod_spec, &adjusted_metadata.mode))
return FILE_SAVE_ERROR;
metadata = &adjusted_metadata;
}
if (!config->delay_context) {
config->delay_context = delay_updates_context_create(root_directory);
if (!config->delay_context)
return FILE_SAVE_ERROR;
}
DelayUpdatesContext* context = config->delay_context;
if (!delay_updates_prepare(context))
return FILE_SAVE_ERROR;
char* staged_path = path_cat(context->staging_root, file->path);
if (!staged_path)
return FILE_SAVE_ERROR;
/* The staged location is brand new (stale leftovers from a prior crash were
wiped by prepare), so the plain atomic temp+rename engine installs the
complete file there. --temp-dir scratch is deliberately not layered on
top of the delay-updates staging tree. */
bool ok =
file_to_disk_secure_with_fsync(staged_path, file->data->data, file->data->size, false, sparse,
metadata, preserve_executability, config->use_fsync, NULL);
if (!ok) {
free(staged_path);
return FILE_SAVE_ERROR;
}
if (!delay_updates_record(context, staged_path, destination_path, file->path)) {
unlink(staged_path);
free(staged_path);
return FILE_SAVE_ERROR;
}
free(staged_path);
return FILE_SAVE_WRITTEN;
}
FileSaveResult file_save_to_disk_full(const char* root_directory, const File* file,
const Config* config) {
/* Backups are incompatible with ignore-existing: moving the entry first
@@ -79,6 +146,18 @@ FileSaveResult file_save_to_disk_full(const char* root_directory, const File* fi
return FILE_SAVE_ERROR;
}
/* --delay-updates diverts the whole write into the staging tree; the rest of
this function is the immediate-install path. */
if (config && config->delay_updates) {
FileSaveResult result =
file_stage_delayed_update(root_directory, destination_path, file, (Config*)config);
free(confined_backup);
free(confined_partial);
free(destination_path);
free(disk_path);
return result;
}
/* --existing checks the final destination, not a temporary partial path. */
if (config && config->existing && !file_path_exists_secure(destination_path)) {
free(confined_backup);
@@ -718,8 +797,13 @@ int receive_manifest(int fd, const Config* config, int* next_status) {
return *status_out == STATUS_FINISHED ? 0 : -1;
}
fprintf(stderr, "Deleting files not in manifest...\n");
bool deletion_ok =
delete_extras_limited(config->receive_root_directory, manifest, MAX_SERVER_DELETE_COUNT);
/* With --delay-updates the staged (not yet published) files live directly
under the receive root in the staging directory; the delete walker must
not treat them as extras or it would remove every staged file before it
can be published. */
const char* skip_staging = config->delay_updates ? DELAY_UPDATES_STAGING_DIR : NULL;
bool deletion_ok = delete_extras_limited(config->receive_root_directory, manifest,
MAX_SERVER_DELETE_COUNT, skip_staging);
array_list_delete(manifest);
if (!deletion_ok)
send_status(fd, STATUS_ERROR);
+15 -5
View File
@@ -205,7 +205,8 @@ static bool is_dir_in_manifest(const char* rel_path, ArrayList* manifest) {
}
static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifest,
size_t max_delete, size_t* deleted_count) {
size_t max_delete, size_t* deleted_count,
const char* skip_root_child) {
int scanfd = dup(dirfd);
if (scanfd < 0)
return false;
@@ -219,6 +220,13 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifes
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
/* A --delay-updates run keeps its staging directory as a direct child of
the receive root. Its contents are not manifest entries yet (they are
published after deletion), so descending into it would delete every
staged file as an "extra". Skip only the top-level staging name; nested
directories with the same name are ordinary destination content. */
if (rel_path[0] == '\0' && skip_root_child && strcmp(entry->d_name, skip_root_child) == 0)
continue;
char* child_rel = path_cat((char*)rel_path, entry->d_name);
if (!child_rel) {
operation_ok = false;
@@ -240,7 +248,8 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifes
int childfd = openat(dirfd, entry->d_name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool child_removed = false;
if (childfd >= 0) {
child_removed = delete_extras_fd(childfd, child_rel, manifest, max_delete, deleted_count);
child_removed = delete_extras_fd(childfd, child_rel, manifest, max_delete, deleted_count,
skip_root_child);
if (!child_removed)
operation_ok = false;
close(childfd);
@@ -291,7 +300,8 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifes
return operation_ok;
}
bool delete_extras_limited(const char* dest_root, ArrayList* manifest, size_t max_delete) {
bool delete_extras_limited(const char* dest_root, ArrayList* manifest, size_t max_delete,
const char* skip_root_child) {
if (!manifest)
return false;
int rootfd;
@@ -308,14 +318,14 @@ bool delete_extras_limited(const char* dest_root, ArrayList* manifest, size_t ma
if (rootfd < 0)
return false;
size_t deleted_count = 0;
bool ok = delete_extras_fd(rootfd, "", manifest, max_delete, &deleted_count);
bool ok = delete_extras_fd(rootfd, "", manifest, max_delete, &deleted_count, skip_root_child);
if (close(rootfd) != 0)
ok = false;
return ok;
}
bool delete_extras(const char* dest_root, ArrayList* manifest) {
return delete_extras_limited(dest_root, manifest, SIZE_MAX);
return delete_extras_limited(dest_root, manifest, SIZE_MAX, NULL);
}
bool has_path_traversal(const char* path) {
+6 -1
View File
@@ -10,7 +10,12 @@ char* output_escape(const char* string, bool eight_bit_output);
char* path_cat(const char* path1, const char* path2);
bool glob_match(const char* pattern, const char* str);
bool delete_extras(const char* dest_root, ArrayList* manifest);
bool delete_extras_limited(const char* dest_root, ArrayList* manifest, size_t max_delete);
/* Remove files/dirs under dest_root that are not listed in manifest. When
skip_root_child is non-NULL, a direct child of dest_root with that exact
name is left untouched (used to protect the --delay-updates staging
directory, which holds files that are still to be published). */
bool delete_extras_limited(const char* dest_root, ArrayList* manifest, size_t max_delete,
const char* skip_root_child);
bool utils_set_authorized_root(int fd, const char* canonical_path);
/* The fd-only compatibility form is fail-closed for path-based operations;
* callers should use utils_set_authorized_root with the canonical identity. */
+225 -1
View File
@@ -12,7 +12,7 @@ from common import (
PROJECT_ROOT, BUILD_DIR, TEST_DATA_DIR,
run_client,
generate_test_files, verify_transfer, clean_dir, make_result,
get_dest_received_dir, CLIENT_CMD,
get_dest_received_dir, CLIENT_CMD, ServerManager,
)
SOURCE_DIR = os.path.join(TEST_DATA_DIR, "feature_source")
@@ -1405,3 +1405,227 @@ class TestLogFileFormat:
expected = {f"{os.path.join(source, rel)} {len(data)}" for rel, data in files.items()}
for line in expected:
assert line in content, f"log file (-m) missing {line!r}"
class TestDelayUpdates:
"""--delay-updates stages every updated file under a private 0700 staging
directory inside the receive root and atomically publishes all of them only
after the whole transfer succeeds."""
STAGING = ".fastsync-stage"
def _make_source(self, name):
source = os.path.join(TEST_DATA_DIR, name)
clean_dir(source)
entries = {
"top.txt": b"top level\n",
"sub/deep.txt": b"deeply nested file\n",
"sub/another.txt": b"another nested file\n" * 20,
"binary.bin": bytes(range(256)) * 4,
}
for rel, content in entries.items():
full = os.path.join(source, rel)
os.makedirs(os.path.dirname(full), exist_ok=True)
with open(full, "wb") as fh:
fh.write(content)
return source
@pytest.mark.parametrize("mt", [False, True])
def test_delay_updates_matches_plain_transfer(self, shared_server, mt):
source = self._make_source("delay_match_src")
plain_dest = os.path.join(TEST_DATA_DIR, "delay_match_plain_dst")
delay_dest = os.path.join(TEST_DATA_DIR, "delay_match_delay_dst")
clean_dir(plain_dest)
clean_dir(delay_dest)
result, _ = run_client(source, plain_dest, port=shared_server.port)
assert result.returncode == 0, f"plain sync failed: {result.stderr[:200]}"
flags = ["--delay-updates"] + (["-m"] if mt else [])
result, _ = run_client(source, delay_dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, f"delay-updates sync failed: {result.stderr[:200]}"
plain_received = get_dest_received_dir(plain_dest, source)
delay_received = get_dest_received_dir(delay_dest, source)
mismatches, missing = verify_transfer(source, delay_received)
assert not missing, f"Missing: {missing}"
assert not mismatches, f"Mismatch: {mismatches}"
for root, _dirs, files in os.walk(delay_received):
for name in files:
rel = os.path.relpath(os.path.join(root, name), delay_received)
assert filecmp.cmp(os.path.join(plain_received, rel),
os.path.join(delay_received, rel), shallow=False), rel
assert not os.path.isdir(os.path.join(delay_dest, self.STAGING)), \
"staging directory left behind after a successful delayed transfer"
@pytest.mark.parametrize("mt", [False, True])
def test_delay_updates_incremental_rerun_no_leftovers(self, shared_server, mt):
source = self._make_source("delay_rerun_src")
dest = os.path.join(TEST_DATA_DIR, "delay_rerun_dst")
clean_dir(dest)
flags = ["--delay-updates", "-M", "--incremental"] + (["-m"] if mt else [])
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, f"first delayed sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
mismatches, missing = verify_transfer(source, received)
assert not missing and not mismatches
assert not os.path.isdir(os.path.join(dest, self.STAGING))
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, f"second delayed sync failed: {result.stderr[:200]}"
assert not os.path.isdir(os.path.join(dest, self.STAGING)), \
"fully-skipped delayed run left a staging directory"
@pytest.mark.parametrize("mt", [False, True])
def test_remove_source_files_with_delay_updates(self, shared_server, mt):
source = self._make_source("delay_rsf_src")
dest = os.path.join(TEST_DATA_DIR, "delay_rsf_dst")
clean_dir(dest)
flags = ["--remove-source-files", "--delay-updates"] + (["-m"] if mt else [])
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, f"delayed remove-source sync failed: {result.stderr[:200]}"
# Sources are removed only after the receiver published every file.
for root, _dirs, files in os.walk(source):
assert files == [], f"source files survived delayed remove-source-files: {files}"
received = get_dest_received_dir(dest, source)
assert os.path.isfile(os.path.join(received, "top.txt"))
assert os.path.isfile(os.path.join(received, "sub", "deep.txt"))
assert not os.path.isdir(os.path.join(dest, self.STAGING))
@pytest.mark.parametrize("mt", [False, True])
def test_delete_with_delay_updates(self, mt):
"""--delete runs before publication, so the delete walker must not treat
the staging directory as a set of extras: a changed file must still be
published after genuine extras are removed. Uses its own server started
with --allow-delete (the shared session server refuses deletion)."""
source = os.path.join(TEST_DATA_DIR, "delay_delete_src")
dest = os.path.join(TEST_DATA_DIR, "delay_delete_dst")
clean_dir(source)
clean_dir(dest)
with open(os.path.join(source, "f.txt"), "wb") as fh:
fh.write(b"AAAA")
with open(os.path.join(source, "extra.txt"), "wb") as fh:
fh.write(b"seed extra")
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
result, _ = run_client(source, dest, port=server.port)
assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
assert _read_file(os.path.join(received, "extra.txt")) == b"seed extra"
# Second source state: f.txt changed, extra.txt removed from source.
with open(os.path.join(source, "f.txt"), "wb") as fh:
fh.write(b"BBBB")
os.remove(os.path.join(source, "extra.txt"))
flags = ["--delete", "--delay-updates"] + (["-m"] if mt else [])
result, _ = run_client(source, dest, flags=flags, port=server.port)
assert result.returncode == 0, \
f"delete+delay-updates sync failed: {result.stderr[:200]}"
assert _read_file(os.path.join(received, "f.txt")) == b"BBBB", \
"changed file was not published after deletion"
assert not os.path.exists(os.path.join(received, "extra.txt")), \
"genuine extra file was not deleted"
assert not os.path.isdir(os.path.join(dest, self.STAGING))
def test_delay_updates_rejects_reserved_backup_dir(self):
"""--backup-dir equal to the internal staging name must be rejected so
an old backup can never be silently installed as the "new" file."""
source = self._make_source("delay_reserved_bak_src")
for variant, suffix in (("bare", ""), ("slash", "/")):
dest = os.path.join(TEST_DATA_DIR, f"delay_reserved_bak_{variant}_dst")
clean_dir(dest)
flags = ["--delay-updates", "--backup", "--backup-dir",
".fastsync-stage" + suffix]
result, _ = run_client(source, dest, flags=flags, port=None)
assert result.returncode != 0, \
f"reserved --backup-dir '{suffix}' was accepted"
assert not os.path.isdir(os.path.join(dest, self.STAGING)), \
"staging directory created by a rejected run"
@pytest.mark.parametrize("remove_source_files", [False, True])
@pytest.mark.parametrize("mt", [False, True])
def test_mid_publish_failure_keeps_published_no_rollback(self, shared_server, mt,
remove_source_files):
"""A stage->publish rename failing part way through publication must
fail the whole transfer, keep the already-published top-level file (no
rollback), leave the not-yet-published nested file absent, and clean up
the staging area. A regular file is planted where the final "sub"
directory must be created, so the nested rename fails (mkdir over a
file is impossible even for root) while the top-level file, which is
always staged first, publishes. With --remove-source-files the sender
must keep every source because no success/outcome frame is ever sent."""
source = os.path.join(TEST_DATA_DIR, "delay_mid_src")
dest = os.path.join(TEST_DATA_DIR, "delay_mid_dst")
clean_dir(source)
clean_dir(dest)
top_path = os.path.join(source, "top.txt")
deep_path = os.path.join(source, "sub", "deep.txt")
with open(top_path, "wb") as fh:
fh.write(b"top payload\n")
os.makedirs(os.path.dirname(deep_path))
with open(deep_path, "wb") as fh:
fh.write(b"deep payload\n")
received = get_dest_received_dir(dest, source)
os.makedirs(received)
with open(os.path.join(received, "sub"), "wb") as fh:
fh.write(b"blocks the nested destination directory")
flags = ["--delay-updates"] + (["-m"] if mt else [])
if remove_source_files:
flags += ["--remove-source-files"]
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode != 0, "blocked nested publish did not fail"
# The top-level file was published before the nested rename failed and
# is intentionally NOT rolled back.
assert _read_file(os.path.join(received, "top.txt")) == b"top payload\n"
# The nested file was never published.
assert not os.path.lexists(os.path.join(received, "sub", "deep.txt")), \
"nested file appeared despite a failed publish"
assert not os.path.isdir(os.path.join(dest, self.STAGING)), \
"staging leftovers after a failed mid-publish"
# Sources survive: no success frame was sent, so a remove-source-files
# sender must not delete anything.
assert os.path.isfile(top_path)
assert os.path.isfile(deep_path)
@pytest.mark.parametrize("mt", [False, True])
def test_remove_source_files_keeps_receiver_skipped_source(self, shared_server, mt):
"""With --delay-updates + --ignore-existing a receiver-skipped source
must survive (its outcome is sent only after publication) while a
freshly delivered file is published and its source removed."""
source = os.path.join(TEST_DATA_DIR, "delay_rsf_skip_src")
dest = os.path.join(TEST_DATA_DIR, "delay_rsf_skip_dst")
clean_dir(source)
clean_dir(dest)
with open(os.path.join(source, "keep.txt"), "wb") as fh:
fh.write(b"existing on dest")
result, _ = run_client(source, dest, port=shared_server.port)
assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}"
with open(os.path.join(source, "keep.txt"), "wb") as fh:
fh.write(b"changed on source")
with open(os.path.join(source, "deliver.txt"), "wb") as fh:
fh.write(b"new file")
flags = ["--remove-source-files", "--ignore-existing", "--delay-updates"] + (["-m"] if mt else [])
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, f"delayed skip sync failed: {result.stderr[:200]}"
# keep.txt already existed at the destination: receiver skip -> source stays.
assert os.path.isfile(os.path.join(source, "keep.txt")), \
"receiver-skipped source was removed despite --ignore-existing"
# deliver.txt was new: staged, published, and its source removed.
assert not os.path.isfile(os.path.join(source, "deliver.txt")), \
"published source was not removed"
received = get_dest_received_dir(dest, source)
assert not os.path.isdir(os.path.join(dest, self.STAGING))
def test_delay_updates_rejects_inplace(self):
source = self._make_source("delay_inplace_src")
dest = os.path.join(TEST_DATA_DIR, "delay_inplace_dst")
clean_dir(dest)
result, _ = run_client(source, dest, flags=["--delay-updates", "--inplace"])
assert result.returncode != 0, "--inplace with --delay-updates was accepted"
assert not os.path.isdir(os.path.join(dest, self.STAGING))
+2
View File
@@ -5,6 +5,7 @@
#include "test_compression.h"
#include "test_config.h"
#include "test_data.h"
#include "test_delay_updates.h"
#include "test_delta.h"
#include "test_file.h"
#include "test_file_sendfile.h"
@@ -51,6 +52,7 @@ int main() {
RUN_TEST(test_metadata);
RUN_TEST(test_glob);
RUN_TEST(test_file);
RUN_TEST(test_delay_updates);
RUN_TEST(test_file_sendfile);
RUN_TEST(test_multiprocessing);
RUN_TEST(test_log);
+46
View File
@@ -1265,6 +1265,49 @@ static void test_parse_args_log_file_format() {
config_delete(cfg);
}
/* --delay-updates is a plain boolean receiver option. */
static void test_parse_args_delay_updates() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--delay-updates", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0);
EXPECT_TRUE(cfg->delay_updates);
EXPECT_EQ_INT(positional_count, 2);
config_delete(cfg);
}
/* rsync rejects --delay-updates with --inplace; FastSync must too. */
static void test_validate_config_delay_updates_rejects_inplace() {
Config* cfg = valid_client_config();
cfg->delay_updates = true;
cfg->inplace = true;
EXPECT_FALSE(validate_config(cfg));
config_delete(cfg);
}
/* --backup-dir may not collide with the internal --delay-updates staging
directory (with or without a trailing slash), or old backups would silently
be installed as the "new" file. */
static void test_validate_config_delay_updates_rejects_reserved_backup_dir() {
static const char* const reserved[] = {".fastsync-stage", ".fastsync-stage/"};
for (size_t i = 0; i < sizeof(reserved) / sizeof(reserved[0]); i++) {
Config* cfg = valid_client_config();
cfg->delay_updates = true;
cfg->backup_dir = str_dup(reserved[i]);
EXPECT_FALSE(validate_config(cfg));
config_delete(cfg);
}
/* A non-colliding backup dir is fine alongside --delay-updates. */
Config* ok = valid_client_config();
ok->delay_updates = true;
ok->backup_dir = str_dup("backups");
EXPECT_TRUE(validate_config(ok));
config_delete(ok);
}
void test_client_cli() {
test_validate_config_required_paths();
test_validate_config_incompatible_options();
@@ -1344,4 +1387,7 @@ void test_client_cli() {
test_parse_args_checksum_choice_aliases();
test_parse_args_checksum_choice_requires_value();
test_parse_args_temp_dir();
test_parse_args_delay_updates();
test_validate_config_delay_updates_rejects_inplace();
test_validate_config_delay_updates_rejects_reserved_backup_dir();
}
+28
View File
@@ -136,6 +136,7 @@ static void test_config_send_receive() {
send_cfg->modify_window = 4;
send_cfg->existing = true;
send_cfg->ignore_existing = true;
send_cfg->delay_updates = true;
send_cfg->skip_compress_set = true;
send_cfg->skip_compress_count = 1;
send_cfg->skip_compress_suffixes = calloc(1, sizeof(char*));
@@ -191,6 +192,8 @@ static void test_config_send_receive() {
ok = false;
if (!recv_cfg->ignore_existing)
ok = false;
if (!recv_cfg->delay_updates)
ok = false;
if (!recv_cfg->skip_compress_set || recv_cfg->skip_compress_count != 1 ||
strcmp(recv_cfg->skip_compress_suffixes[0], ".zip") != 0)
ok = false;
@@ -405,6 +408,30 @@ static void test_config_temp_dir_roundtrip() {
config_delete(c);
}
static void test_config_delay_updates_reserved_backup_rejected() {
if (is_running_under_valgrind())
return;
Config* c = config_create();
EXPECT_NOT_NULL(c);
c->send_directory = str_dup("/src");
c->receive_root_directory = str_dup("/dst");
c->delay_updates = true;
c->backup_dir = str_dup(".fastsync-stage");
/* The receiver-side wire validation must reject a --backup-dir that collides
with the internal delay-updates staging directory. */
EXPECT_FALSE(roundtrip_config_ok(c));
config_delete(c);
c = config_create();
EXPECT_NOT_NULL(c);
c->send_directory = str_dup("/src");
c->receive_root_directory = str_dup("/dst");
c->delay_updates = true;
c->backup_dir = str_dup("backups");
EXPECT_TRUE(roundtrip_config_ok(c));
config_delete(c);
}
static void test_config_is_remote_dest() {
/* Valid SSH-style destinations */
EXPECT_TRUE(config_is_remote_dest("user@host:/path"));
@@ -440,6 +467,7 @@ void test_config() {
test_config_receive_truncated();
test_config_string_null_vs_empty_roundtrip();
test_config_temp_dir_roundtrip();
test_config_delay_updates_reserved_backup_rejected();
}
test_config_is_remote_dest();
}
+296
View File
@@ -0,0 +1,296 @@
#include "test_delay_updates.h"
#include "config.h"
#include "delay_updates.h"
#include "file.h"
#include "file_receive.h"
#include "test_utils.h"
#include "utils.h"
#include <dirent.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/stat.h>
#include <unistd.h>
/* Recursively remove a test tree (never follows symlinks). */
static void remove_tree(const char* path) {
struct stat st;
if (lstat(path, &st) != 0)
return;
if (S_ISDIR(st.st_mode)) {
DIR* dir = opendir(path);
if (!dir)
return;
const struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
char* child = path_cat(path, entry->d_name);
if (child) {
remove_tree(child);
free(child);
}
}
closedir(dir);
rmdir(path);
} else {
unlink(path);
}
}
/* Build a File that carries `content`. */
static File* make_file(const char* path, const char* content) {
File* f = file_create(path);
if (!f)
return NULL;
f->data->data = malloc(strlen(content));
if (!f->data->data) {
file_destroy(f);
return NULL;
}
memcpy(f->data->data, content, strlen(content));
f->data->size = strlen(content);
return f;
}
static char* read_all(const char* path) {
FILE* fp = fopen(path, "rb");
if (!fp)
return NULL;
char buf[256] = {0};
size_t n = fread(buf, 1, sizeof(buf) - 1, fp);
fclose(fp);
char* out = malloc(n + 1);
if (!out)
return NULL;
memcpy(out, buf, n);
out[n] = '\0';
return out;
}
static void test_delay_updates_no_final_before_publish() {
const char* root = "test_delay_tmp";
remove_tree(root);
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
cfg->delay_updates = true;
File* f = make_file("sub/file.txt", "staged payload");
EXPECT_NOT_NULL(f);
// cppcheck-suppress knownConditionTrueFalse
if (!cfg || !f)
goto out;
EXPECT_EQ_INT(file_save_to_disk_full(root, f, cfg), FILE_SAVE_WRITTEN);
EXPECT_NOT_NULL(cfg->delay_context);
const char* final_path = "test_delay_tmp/sub/file.txt";
/* Before publication the final destination must not contain the file. */
EXPECT_FALSE(file_path_exists_secure(final_path));
/* The complete staged copy must live inside the staging tree. */
char* staged = path_cat("test_delay_tmp/.fastsync-stage", "/sub/file.txt");
EXPECT_NOT_NULL(staged);
// cppcheck-suppress knownConditionTrueFalse
if (staged) {
char* content = read_all(staged);
EXPECT_NOT_NULL(content);
// cppcheck-suppress knownConditionTrueFalse
if (content) {
EXPECT_EQ_STR(content, "staged payload");
free(content);
}
free(staged);
}
out:
file_destroy(f);
config_delete(cfg);
remove_tree(root);
}
static void test_delay_updates_publish_installs_files() {
const char* root = "test_delay_pub_tmp";
remove_tree(root);
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
cfg->delay_updates = true;
File* f = make_file("sub/file.txt", "published payload");
EXPECT_NOT_NULL(f);
// cppcheck-suppress knownConditionTrueFalse
if (!cfg || !f)
goto out;
EXPECT_EQ_INT(file_save_to_disk_full(root, f, cfg), FILE_SAVE_WRITTEN);
const char* final_path = "test_delay_pub_tmp/sub/file.txt";
EXPECT_FALSE(file_path_exists_secure(final_path));
EXPECT_TRUE(delay_updates_publish(cfg->delay_context, cfg));
/* After a successful publish the file is installed and staging is gone. */
char* content = read_all(final_path);
EXPECT_NOT_NULL(content);
// cppcheck-suppress knownConditionTrueFalse
if (content) {
EXPECT_EQ_STR(content, "published payload");
free(content);
}
EXPECT_FALSE(file_path_exists_secure("test_delay_pub_tmp/.fastsync-stage"));
out:
file_destroy(f);
config_delete(cfg);
remove_tree(root);
}
/* The staged tree is cleaned on the error/abort path and final files that were
never published do not appear at the destination. */
static void test_delay_updates_cleanup_removes_staged() {
const char* root = "test_delay_clean_tmp";
remove_tree(root);
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
cfg->delay_updates = true;
File* f = make_file("sub/file.txt", "never installed");
EXPECT_NOT_NULL(f);
// cppcheck-suppress knownConditionTrueFalse
if (!cfg || !f)
goto out;
EXPECT_EQ_INT(file_save_to_disk_full(root, f, cfg), FILE_SAVE_WRITTEN);
EXPECT_TRUE(file_path_exists_secure("test_delay_clean_tmp/.fastsync-stage/sub/file.txt"));
delay_updates_cleanup(cfg->delay_context);
EXPECT_FALSE(file_path_exists_secure("test_delay_clean_tmp/.fastsync-stage"));
EXPECT_FALSE(file_path_exists_secure("test_delay_clean_tmp/sub/file.txt"));
out:
file_destroy(f);
config_delete(cfg);
remove_tree(root);
}
/* With --backup the previous version is only moved aside at publication. */
static void test_delay_updates_backup_deferred_to_publish() {
const char* root = "test_delay_bak_tmp";
remove_tree(root);
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
cfg->delay_updates = true;
cfg->backup = true;
EXPECT_TRUE(file_write_to_disk("test_delay_bak_tmp/file.txt", "AAAA", 4, false, false));
File* f = make_file("file.txt", "BBBB");
EXPECT_NOT_NULL(f);
// cppcheck-suppress knownConditionTrueFalse
if (!cfg || !f)
goto out;
EXPECT_EQ_INT(file_save_to_disk_full(root, f, cfg), FILE_SAVE_WRITTEN);
/* Stage time must not touch the final file or create the backup yet. */
char* before = read_all("test_delay_bak_tmp/file.txt");
EXPECT_NOT_NULL(before);
// cppcheck-suppress knownConditionTrueFalse
if (before) {
EXPECT_EQ_STR(before, "AAAA");
free(before);
}
EXPECT_FALSE(file_path_exists_secure("test_delay_bak_tmp/file.txt~"));
EXPECT_TRUE(delay_updates_publish(cfg->delay_context, cfg));
char* after = read_all("test_delay_bak_tmp/file.txt");
char* backup = read_all("test_delay_bak_tmp/file.txt~");
EXPECT_NOT_NULL(after);
EXPECT_NOT_NULL(backup);
// cppcheck-suppress knownConditionTrueFalse
if (after) {
EXPECT_EQ_STR(after, "BBBB");
free(after);
}
// cppcheck-suppress knownConditionTrueFalse
if (backup) {
EXPECT_EQ_STR(backup, "AAAA");
free(backup);
}
out:
file_destroy(f);
config_delete(cfg);
remove_tree(root);
}
/* Skip/update policy checks run against the final path at stage time, matching
what an immediate run would decide. */
static void test_delay_updates_skip_semantics() {
const char* root = "test_delay_skip_tmp";
remove_tree(root);
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
cfg->delay_updates = true;
/* --existing: final destination missing -> skipped, nothing staged. */
File* missing = make_file("missing.txt", "new");
EXPECT_NOT_NULL(missing);
// cppcheck-suppress knownConditionTrueFalse
if (!cfg || !missing)
goto out;
cfg->existing = true;
EXPECT_EQ_INT(file_save_to_disk_full(root, missing, cfg), FILE_SAVE_SKIPPED);
cfg->existing = false;
/* --ignore-existing: final destination present -> skipped. */
EXPECT_TRUE(file_write_to_disk("test_delay_skip_tmp/existing.txt", "old", 3, false, false));
File* present = make_file("existing.txt", "new");
EXPECT_NOT_NULL(present);
// cppcheck-suppress knownConditionTrueFalse
if (!present)
goto out;
cfg->ignore_existing = true;
EXPECT_EQ_INT(file_save_to_disk_full(root, present, cfg), FILE_SAVE_SKIPPED);
cfg->ignore_existing = false;
/* Without a skip flag the file is staged and later published. */
File* fresh = make_file("fresh.txt", "content");
EXPECT_NOT_NULL(fresh);
// cppcheck-suppress knownConditionTrueFalse
if (!fresh)
goto out;
EXPECT_EQ_INT(file_save_to_disk_full(root, fresh, cfg), FILE_SAVE_WRITTEN);
EXPECT_TRUE(delay_updates_publish(cfg->delay_context, cfg));
char* content = read_all("test_delay_skip_tmp/fresh.txt");
EXPECT_NOT_NULL(content);
// cppcheck-suppress knownConditionTrueFalse
if (content) {
EXPECT_EQ_STR(content, "content");
free(content);
}
out:
file_destroy(missing);
file_destroy(present);
file_destroy(fresh);
config_delete(cfg);
remove_tree(root);
}
/* The reserved staging name must be recognizable for validation, including
with a trailing slash. */
static void test_delay_updates_reserved_name_helper() {
EXPECT_TRUE(delay_updates_staging_name_conflict(".fastsync-stage"));
EXPECT_TRUE(delay_updates_staging_name_conflict(".fastsync-stage/"));
EXPECT_TRUE(delay_updates_staging_name_conflict(".fastsync-stage///"));
EXPECT_FALSE(delay_updates_staging_name_conflict(NULL));
EXPECT_FALSE(delay_updates_staging_name_conflict(""));
EXPECT_FALSE(delay_updates_staging_name_conflict("backups"));
EXPECT_FALSE(delay_updates_staging_name_conflict(".fastsync-stage.bak"));
}
void test_delay_updates() {
test_delay_updates_reserved_name_helper();
test_delay_updates_no_final_before_publish();
test_delay_updates_publish_installs_files();
test_delay_updates_cleanup_removes_staged();
test_delay_updates_backup_deferred_to_publish();
test_delay_updates_skip_semantics();
}
+6
View File
@@ -0,0 +1,6 @@
#ifndef TEST_DELAY_UPDATES_H
#define TEST_DELAY_UPDATES_H
void test_delay_updates(void);
#endif