feat: implement --delay-updates receiver staging and publication
CI / lint (pull_request) Failing after 22s
CI / build-and-test (pull_request) Skipped
CI / sanitizers (address) (pull_request) Skipped
CI / sanitizers (undefined) (pull_request) Skipped
CI / fuzz-build (pull_request) Skipped
CI / coverage (pull_request) Skipped
CI / valgrind (pull_request) Skipped
CI / lint (pull_request) Failing after 22s
CI / build-and-test (pull_request) Skipped
CI / sanitizers (address) (pull_request) Skipped
CI / sanitizers (undefined) (pull_request) Skipped
CI / fuzz-build (pull_request) Skipped
CI / coverage (pull_request) Skipped
CI / valgrind (pull_request) Skipped
Stage every successfully written file under a private 0700 .fastsync-stage directory inside the receive root and atomically publish all staged files only after the whole protocol stream (manifest/delete handling included) has completed, immediately before the success/outcome frame. On any abort/error before publication nothing is installed and staging is removed; a publish failure aborts the transfer with best-effort cleanup of the remainder (already-published files are not rolled back). Crash leftovers are wiped when the next delayed transfer starts. Wire: new delay_updates config flag (selection-options block), protocol version bumped to 2.6.0, client/server validation rejects --inplace. CLI/usage/validation updated. Works in single-threaded and -m modes (exactly one write_thread stages files; the staged-file registry is mutex-protected; publication runs once after both threads join). --existing/--ignore-existing/--update decide against the final destination at stage time; --backup is deferred to publication. remove_source_files outcomes are only sent after publication so skipped/unpublished sources are never deleted. Default (no flag) behavior is unchanged. Tests: config wire round-trip, CLI parse, --inplace rejection, new test_delay_updates unit suite (27 suites total), and integration TestDelayUpdates covering single/-m parity, incremental reruns, remove source files, receiver-skip ordering, and a deterministic publish-failure abort path.
This commit is contained in:
@@ -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},
|
||||
|
||||
@@ -60,5 +60,9 @@ 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;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
|
||||
@@ -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");
|
||||
|
||||
@@ -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
@@ -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)
|
||||
|
||||
+18
-7
@@ -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,7 @@ 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) &&
|
||||
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 +252,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 +297,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 +393,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
@@ -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);
|
||||
|
||||
@@ -0,0 +1,279 @@
|
||||
#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/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;
|
||||
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);
|
||||
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);
|
||||
}
|
||||
|
||||
/* 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;
|
||||
}
|
||||
bool ok = delay_wipe_dir_fd(fd);
|
||||
if (close(fd) != 0)
|
||||
ok = false;
|
||||
if (!ok) {
|
||||
log_message(LOG_LEVEL_ERROR, "could not clear stale --delay-updates staging files under '%s'",
|
||||
context->staging_root);
|
||||
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(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;
|
||||
}
|
||||
|
||||
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. */
|
||||
int fd = open(context->staging_root, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
|
||||
if (fd >= 0) {
|
||||
delay_wipe_dir_fd(fd);
|
||||
close(fd);
|
||||
rmdir(context->staging_root);
|
||||
}
|
||||
return ok;
|
||||
}
|
||||
|
||||
void delay_updates_cleanup(DelayUpdatesContext* context) {
|
||||
if (!context)
|
||||
return;
|
||||
int fd = open(context->staging_root, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
|
||||
if (fd < 0)
|
||||
return;
|
||||
delay_wipe_dir_fd(fd);
|
||||
close(fd);
|
||||
rmdir(context->staging_root);
|
||||
}
|
||||
@@ -0,0 +1,62 @@
|
||||
#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 and stale leftovers wiped once */
|
||||
} DelayUpdatesContext;
|
||||
|
||||
/* Name of the private staging subdirectory created under the receive root. */
|
||||
#define DELAY_UPDATES_STAGING_DIR ".fastsync-stage"
|
||||
|
||||
/* 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
@@ -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",
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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,71 @@ 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) {
|
||||
bool sparse = config && config->preserve_sparse;
|
||||
bool preserve_executability = config && config->use_executability;
|
||||
|
||||
if (config && config->existing && !file_path_exists_secure(destination_path))
|
||||
return FILE_SAVE_SKIPPED;
|
||||
if (config && config->ignore_existing && file_path_exists_secure(destination_path))
|
||||
return FILE_SAVE_SKIPPED;
|
||||
if (config && 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 && 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 && 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 +145,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);
|
||||
|
||||
Reference in New Issue
Block a user