From 14c064a093dbf75c6f8b8ea0535379126f6ef420 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sat, 5 Sep 2026 12:46:44 +0200 Subject: [PATCH 1/2] fix: #251 #252 #253 #255 #256 #257 receiver & remove-source correctness #251 --remove-source-files deletes sources that were skipped receiver-side (--existing/--ignore-existing/--update). The receiver now reports a per-file outcome for every processed data file when the sender requests removal; the client only unlinks sources the receiver actually wrote. add remove_source_files to the wire config and bump the protocol to 2.5.0. #252 --backup/--suffix/--backup-dir broken by NULL-vs-empty wire loss. Receivers canonicalize the empty wire string back to NULL for backup_dir, temp_dir, partial_dir and suffix, and --suffix is received unconditionally. #253 --partial --partial-dir never installed completed files. file_save_to_disk now renames a fully written partial-dir file into the real destination. #255 STATUS_CHECK read the entire old file before the size/mtime quick check. Old contents are only read when a checksum compare or delta needs them. #256 receive_delta_file failure paths did not set *failed, so the caller sent STATUS_NEXT and waited for a body that never came. Every NULL return now marks the transfer failed. #257 files >64 MiB could not transfer. Whole-file receive caps raised to the 256 MiB connection/allocation ceiling (chunk caps stay 64 MiB) and the client ignores SIGPIPE so a server-side close surfaces as a clean error. Unit tests added: config NULL-vs-empty round trip, incremental quick-check skip/NEXT paths, delta oversize failure, partial-dir install, save-result skip reporting. --- README.md | 5 +- src/client/client_cli.c | 5 + src/client/client_send.c | 43 ++++-- src/server/receiver.c | 95 +++++++++++-- src/server/receiver.h | 21 +++ src/server/server.c | 12 +- src/shared/config.c | 63 +++++++-- src/shared/config.h | 2 +- src/shared/file_receive.c | 179 +++++++++++++++---------- src/shared/file_receive.h | 8 ++ src/shared/multiprocessing.c | 27 +++- src/shared/multiprocessing.h | 2 + src/shared/protocol.h | 14 +- tests/test_config.c | 100 ++++++++++++++ tests/test_file.c | 108 +++++++++++++++ tests/test_server.c | 250 +++++++++++++++++++++++++++++++++++ 16 files changed, 823 insertions(+), 111 deletions(-) diff --git a/README.md b/README.md index 9c35016..9a174b8 100644 --- a/README.md +++ b/README.md @@ -484,10 +484,11 @@ defaults to the current directory. | ## Protocol and Security -FastSync protocol version `2.4.0` is shared by the client and server. The +FastSync protocol version `2.5.0` is shared by the client and server. The current protocol is sender-driven and includes configuration negotiation, including the maximum allocation limit, incremental checks, checksums, -manifests, keep-alives, abort handling, and FastSync-native delta messages. +manifests, keep-alives, abort handling, per-file remove-source results, and +FastSync-native delta messages. Client and server versions must currently match exactly. TLS provides encrypted TCP transport. Supplying `--ca` enables certificate diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 4caf6fb..ad5b77c 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -12,6 +12,7 @@ #include "utils.h" #include #include +#include #include #include #include @@ -903,6 +904,10 @@ static int read_patterns_from_file(const char* filepath, char*** patterns, int* #ifndef FASTSYNC_TEST_BUILD int main(int argc, char* argv[]) { + /* The server may close a connection mid-stream (e.g. when it rejects an + oversized delta). Ignore SIGPIPE so that a broken TCP connection + surfaces as a clean write error instead of killing the client. */ + signal(SIGPIPE, SIG_IGN); const char* env_source = NULL; const char* env_dest = NULL; bool save_to_disk = false; diff --git a/src/client/client_send.c b/src/client/client_send.c index 5adfa11..c63e024 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -109,16 +109,13 @@ static bool add_chunk_to_manifest(ArrayList* manifest, const Chunk* chunk) { return true; } -static bool finalize_transfer(Client* client) { - Status status; - return send_status(client->file_descriptor, STATUS_FINISHED) && - receive_status(client->file_descriptor, &status) && status == STATUS_OK; -} +/* (finalize_transfer is defined after the SourceFile helpers below.) */ -typedef struct { +typedef struct SourceFile { char* path; dev_t device; ino_t inode; + bool skipped; /* receiver reported the file was not written */ } SourceFile; static void source_file_destroy(void* item) { @@ -135,6 +132,8 @@ static void remove_transferred_sources(const Config* config, ArrayList* paths) { return; for (int i = 0; i < paths->size; i++) { SourceFile* source = paths->items[i]; + if (source->skipped) + continue; const char* slash = strrchr(source->path, '/'); const char* leaf = slash ? slash + 1 : source->path; char parent[PATH_MAX]; @@ -177,6 +176,7 @@ static SourceFile* source_file_create(const File* file) { source->path = str_dup(file->path); source->device = st.st_dev; source->inode = st.st_ino; + source->skipped = false; if (!source->path) { source_file_destroy(source); return NULL; @@ -203,6 +203,33 @@ static void mark_sender_done(PipelineContextSender* context) { mtx_unlock(&context->mutex_progress); } +/* Send the final STATUS_FINISHED frame and await the receiver's verdict. + When --remove-source-files is active the receiver acknowledges each data + file it processed, in send order: STATUS_NEXT means the file was written, + STATUS_OK means the file was skipped/unchanged. Skipped sources are marked + so the later removal pass keeps them. */ +static bool finalize_transfer(Client* client, const Config* config, ArrayList* remove_sources) { + if (!send_status(client->file_descriptor, STATUS_FINISHED)) + return false; + if (config->remove_source_files && remove_sources) { + for (int i = 0; i < remove_sources->size; i++) { + Status per_file; + if (!receive_status(client->file_descriptor, &per_file)) + return false; + if (per_file == STATUS_ERROR) + return false; + if (per_file == STATUS_OK) { + ((SourceFile*)remove_sources->items[i])->skipped = true; + } else if (per_file != STATUS_NEXT) { + log_message(LOG_LEVEL_ERROR, "Unexpected per-file status from receiver"); + return false; + } + } + } + Status status; + return receive_status(client->file_descriptor, &status) && status == STATUS_OK; +} + static void pipeline_cancel(PipelineContextSender* context) { mtx_lock(&context->mutex_scanner); mtx_lock(&context->mutex_loader); @@ -575,7 +602,7 @@ static int send_chunks_multithreaded(void* pipeline_context) { if (send_delete_manifest(client->file_descriptor, context->manifest) != 0) goto send_fail; } - bool ok = finalize_transfer(client); + bool ok = finalize_transfer(client, context->config, context->remove_source_files); if (ok) remove_transferred_sources(context->config, context->remove_source_files); mtx_lock(&context->mutex_progress); @@ -866,7 +893,7 @@ int send_files(Config* config) { array_list_delete(manifest); manifest = NULL; } - bool ok = finalize_transfer(client); + bool ok = finalize_transfer(client, config, remove_sources); if (ok) remove_transferred_sources(config, remove_sources); if (config->show_progress && !config->quiet) diff --git a/src/server/receiver.c b/src/server/receiver.c index 502eea2..35fadb2 100644 --- a/src/server/receiver.c +++ b/src/server/receiver.c @@ -1,6 +1,7 @@ #include "receiver.h" #include "chunk.h" +#include "file_receive.h" #include "log.h" #include "metadata.h" #include "protocol.h" @@ -8,6 +9,48 @@ #include #include +bool receiver_outcomes_append(ReceiverOutcomes* outcomes, unsigned char code) { + if (!outcomes) + return false; + if (outcomes->count == outcomes->capacity) { + size_t new_capacity = outcomes->capacity == 0 ? 64 : outcomes->capacity * 2; + if (new_capacity < outcomes->capacity) + return false; + unsigned char* grown = realloc(outcomes->entries, new_capacity); + if (!grown) + return false; + outcomes->entries = grown; + outcomes->capacity = new_capacity; + } + outcomes->entries[outcomes->count++] = code; + return true; +} + +void receiver_outcomes_destroy(ReceiverOutcomes* outcomes) { + if (!outcomes) + return; + free(outcomes->entries); + outcomes->entries = NULL; + outcomes->count = 0; + outcomes->capacity = 0; +} + +/* End-of-transfer success frame. When --remove-source-files was negotiated + each processed data file is acknowledged first (STATUS_NEXT = written, + STATUS_OK = skipped) so the sender never removes a source the receiver did + not actually store. The frame always ends with a plain STATUS_OK. */ +bool receiver_send_final_success(int fd, const Config* config, const ReceiverOutcomes* outcomes) { + if (!config->remove_source_files) + return send_status(fd, STATUS_OK); + size_t count = outcomes ? outcomes->count : 0; + for (size_t i = 0; i < count; i++) { + Status per_file = outcomes->entries[i] == FILE_SAVE_WRITTEN ? STATUS_NEXT : STATUS_OK; + if (!send_status(fd, per_file)) + return false; + } + return send_status(fd, STATUS_OK); +} + static bool receiver_process_chunk(Chunk* chunk, const ReceiverSink* sink) { if (!chunk || !sink || !sink->store_file) return false; @@ -52,7 +95,7 @@ static bool receiver_process_batch(Config* config, int file_descriptor) { send_status(file_descriptor, STATUS_ERROR); return false; } - if (check_size > MAX_RECEIVE_FILE_SIZE) { + if (check_size > MAX_RECEIVE_WHOLE_FILE_SIZE) { free(check_path); send_status(file_descriptor, STATUS_ERROR); return false; @@ -130,8 +173,14 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status"); goto receive_error; } - if (sink->send_success && !send_status(file_descriptor, STATUS_OK)) - return -1; + if (sink->send_success) { + if (sink->send_success_frame) { + if (!sink->send_success_frame(file_descriptor, sink->context)) + return -1; + } else if (!send_status(file_descriptor, STATUS_OK)) { + return -1; + } + } return 0; receive_error: @@ -140,15 +189,41 @@ receive_error: return -1; } -static bool receiver_save_file(File* file, void* context) { - Config* config = context; - bool success = - !config->save_to_disk || file_save_to_disk(config->receive_root_directory, file, config); +/* ---- Single-threaded sink (used by receiver_receive_files) ---- */ + +typedef struct { + Config* config; + ReceiverOutcomes outcomes; +} ReceiverSaveContext; + +static bool receiver_save_file(File* file, void* context_pointer) { + ReceiverSaveContext* context = context_pointer; + FileSaveResult result = FILE_SAVE_ERROR; + if (!context->config->save_to_disk) { + /* Nothing is stored; report the file as not-written so a + --remove-source-files sender keeps its source. */ + result = FILE_SAVE_SKIPPED; + } else { + result = file_save_to_disk_full(context->config->receive_root_directory, file, context->config); + } + if (result != FILE_SAVE_ERROR && context->config->remove_source_files && + !receiver_outcomes_append(&context->outcomes, (unsigned char)result)) { + file_destroy(file); + return false; + } file_destroy(file); - return success; + return result != FILE_SAVE_ERROR; +} + +static bool receiver_send_success_frame(int fd, void* context_pointer) { + ReceiverSaveContext* context = context_pointer; + return receiver_send_final_success(fd, context->config, &context->outcomes); } int receiver_receive_files(Config* config, int file_descriptor) { - ReceiverSink sink = {receiver_save_file, config, true, true}; - return receiver_process(config, file_descriptor, &sink); + 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); + receiver_outcomes_destroy(&context.outcomes); + return ret; } diff --git a/src/server/receiver.h b/src/server/receiver.h index d619d2e..4da64b0 100644 --- a/src/server/receiver.h +++ b/src/server/receiver.h @@ -3,16 +3,37 @@ #include "config.h" #include "file.h" +#include "file_receive.h" typedef bool (*ReceiverFileSink)(File* file, void* context); +/* Ordered per-file save outcomes for one connection. One entry is appended + for every data-bearing file the receiver processes (in the order the files + were sent) so the sender of a --remove-source-files transfer can be told + which sources were actually written versus skipped on the receiver. */ +typedef struct { + unsigned char* entries; /* FILE_SAVE_WRITTEN or FILE_SAVE_SKIPPED */ + size_t count; + size_t capacity; +} ReceiverOutcomes; + +typedef bool (*ReceiverSuccessFrame)(int fd, void* context); + typedef struct { ReceiverFileSink store_file; void* context; bool send_error; bool send_success; + /* Emits the end-of-transfer success frame. When the sender requested + --remove-source-files this includes one per-file status per processed + data file followed by the final STATUS_OK; otherwise just STATUS_OK. */ + ReceiverSuccessFrame send_success_frame; } ReceiverSink; +bool receiver_outcomes_append(ReceiverOutcomes* outcomes, unsigned char code); +void receiver_outcomes_destroy(ReceiverOutcomes* outcomes); +bool receiver_send_final_success(int fd, const Config* config, const ReceiverOutcomes* outcomes); + int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink); int receiver_receive_files(Config* config, int file_descriptor); diff --git a/src/server/server.c b/src/server/server.c index c968115..c96612d 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -342,9 +342,15 @@ void handler(int file_descriptor) { int writer_result; thrd_join(receiver, &receiver_result); thrd_join(writer, &writer_result); - send_status(file_descriptor, receiver_result == thrd_success && writer_result == thrd_success - ? STATUS_OK - : STATUS_ERROR); + bool transfer_ok = receiver_result == thrd_success && writer_result == thrd_success; + 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) + log_message(LOG_LEVEL_ERROR, "Transfer failed"); pipeline_context_receiver_destroy(context); } else { if (receiver_receive_files(config, file_descriptor) != 0) diff --git a/src/shared/config.c b/src/shared/config.c index f088c3c..50c6f05 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -139,8 +139,9 @@ static bool validate_received_config(const Config* config) { valid_wire_bool(config->use_delete) && valid_wire_bool(config->use_incremental) && valid_wire_bool(config->size_only) && valid_wire_bool(config->ignore_times) && valid_wire_bool(config->use_delta) && valid_wire_bool(config->backup) && - valid_wire_bool(config->follow_symlinks) && valid_wire_bool(config->copy_links) && - valid_wire_bool(config->safe_links) && valid_wire_bool(config->copy_unsafe_links) && + valid_wire_bool(config->remove_source_files) && valid_wire_bool(config->follow_symlinks) && + valid_wire_bool(config->copy_links) && valid_wire_bool(config->safe_links) && + valid_wire_bool(config->copy_unsafe_links) && valid_wire_bool(config->preserve_hard_links) && valid_wire_bool(config->preserve_acls) && valid_wire_bool(config->preserve_xattrs) && valid_wire_bool(config->preserve_devices) && valid_wire_bool(config->preserve_sparse) && valid_wire_bool(config->ignore_existing) && @@ -274,11 +275,11 @@ static bool send_delta_fields(int fd, const Config* c) { static bool send_file_options(int fd, const Config* c) { return send_int(fd, c->backup) && send_str(fd, c->backup_dir ? c->backup_dir : "") && - send_int(fd, c->follow_symlinks) && send_int(fd, c->copy_links) && - send_int(fd, c->safe_links) && send_int(fd, c->copy_unsafe_links) && - send_int(fd, c->preserve_hard_links) && send_int(fd, c->preserve_acls) && - send_int(fd, c->preserve_xattrs) && send_int(fd, c->preserve_devices) && - send_int(fd, c->preserve_sparse); + send_int(fd, c->remove_source_files) && send_int(fd, c->follow_symlinks) && + send_int(fd, c->copy_links) && send_int(fd, c->safe_links) && + send_int(fd, c->copy_unsafe_links) && send_int(fd, c->preserve_hard_links) && + send_int(fd, c->preserve_acls) && send_int(fd, c->preserve_xattrs) && + send_int(fd, c->preserve_devices) && send_int(fd, c->preserve_sparse); } static bool send_selection_options(int fd, const Config* c) { @@ -355,8 +356,17 @@ static bool receive_delta_fields(int fd, Config* c) { static bool receive_file_options(int fd, Config* c) { if (!receive_wire_bool(fd, &c->backup)) return false; - c->backup_dir = receive_str(fd); - if (!c->backup_dir) + char* backup_dir = receive_str(fd); + if (!backup_dir) + return false; + if (*backup_dir != '\0') { + c->backup_dir = backup_dir; + } else { + /* The sender serializes an unset (NULL) string as "", so canonicalize the + empty wire value back to NULL to preserve NULL-vs-empty semantics. */ + free(backup_dir); + } + if (!receive_wire_bool(fd, &c->remove_source_files)) return false; bool* flags[] = {&c->follow_symlinks, &c->copy_links, &c->safe_links, &c->copy_unsafe_links, &c->preserve_hard_links, &c->preserve_acls, @@ -386,12 +396,37 @@ static bool receive_selection_options(int fd, Config* c) { } static bool receive_resume_options(int fd, Config* c) { - c->temp_dir = receive_str(fd); - if (!c->temp_dir || !receive_wire_bool(fd, &c->partial)) + char* temp_dir = receive_str(fd); + if (!temp_dir) return false; - c->partial_dir = receive_str(fd); - c->suffix = c->partial_dir ? receive_str(fd) : NULL; - if (!c->partial_dir || !c->suffix || !receive_wire_bool(fd, &c->delete_before)) + if (*temp_dir != '\0') { + c->temp_dir = temp_dir; + } else { + free(temp_dir); + } + if (!receive_wire_bool(fd, &c->partial)) + return false; + /* These options have NULL client defaults, so the sender transmits an empty + string for "unset". Canonicalize the empty wire value back to NULL so + receivers observe exactly what the client configured (plain --backup, for + example, must not look like --backup-dir ""). */ + char* partial_dir = receive_str(fd); + if (!partial_dir) + return false; + if (*partial_dir != '\0') { + c->partial_dir = partial_dir; + } else { + free(partial_dir); + } + char* suffix = receive_str(fd); + if (!suffix) + return false; + if (*suffix != '\0') { + c->suffix = suffix; + } else { + free(suffix); + } + if (!receive_wire_bool(fd, &c->delete_before)) return false; if (!receive_wire_bool(fd, &c->checksum)) return false; diff --git a/src/shared/config.h b/src/shared/config.h index 1bd6a46..de7fa4e 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -147,7 +147,7 @@ typedef struct Config { bool skip_compress_set; } Config; -#define PROTOCOL_VERSION "2.4.0" +#define PROTOCOL_VERSION "2.5.0" #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) Config* config_create(void); diff --git a/src/shared/file_receive.c b/src/shared/file_receive.c index a7e5973..30cecb4 100644 --- a/src/shared/file_receive.c +++ b/src/shared/file_receive.c @@ -20,9 +20,14 @@ #include "utils.h" #define MAX_SERVER_DELETE_COUNT 100000U -#define MAX_FILE_DATA_SIZE MAX_RECEIVE_FILE_SIZE +#define MAX_FILE_DATA_SIZE MAX_RECEIVE_WHOLE_FILE_SIZE bool file_save_to_disk(const char* root_directory, const File* file, const Config* config) { + return file_save_to_disk_full(root_directory, file, config) != FILE_SAVE_ERROR; +} + +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 would make a concurrent no-replace commit overwrite its old name. */ bool backup_enabled = config && config->backup && !config->ignore_existing; @@ -32,6 +37,7 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi const char* backup_suffix = (config && config->suffix) ? config->suffix : "~"; const char* backup_dir = (config && config->backup_dir) ? config->backup_dir : NULL; const char* partial_dir = (config && config->partial_dir) ? config->partial_dir : NULL; + bool use_partial_root = partial_dir && config && config->partial; char *confined_backup = NULL, *confined_partial = NULL, *disk_path = NULL; char* destination_path = NULL; char *backup_path = NULL, *parent_copy = NULL; @@ -42,23 +48,22 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi (!backup_suffix || backup_suffix[0] == '\0' || strchr(backup_suffix, '/') != NULL || strcmp(backup_suffix, ".") == 0 || strcmp(backup_suffix, "..") == 0))) { log_message(LOG_LEVEL_ERROR, "Invalid file or path received"); - return false; + return FILE_SAVE_ERROR; } /* These options arrive from the client. They are names below the server root, never independent filesystem roots. */ if ((backup_dir && (backup_dir[0] == '/' || has_path_traversal(backup_dir))) || (partial_dir && (partial_dir[0] == '/' || has_path_traversal(partial_dir)))) - return false; + return FILE_SAVE_ERROR; if (backup_dir && !(confined_backup = path_cat(root_directory, backup_dir))) - return false; + return FILE_SAVE_ERROR; if (partial_dir && !(confined_partial = path_cat(root_directory, partial_dir))) { free(confined_backup); - return false; + return FILE_SAVE_ERROR; } - const char* actual_root = - (partial_dir && config && config->partial) ? confined_partial : root_directory; + const char* actual_root = use_partial_root ? confined_partial : root_directory; destination_path = path_cat(root_directory, file->path); disk_path = path_cat(actual_root, file->path); if (destination_path == NULL || disk_path == NULL) { @@ -66,7 +71,7 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi free(confined_partial); free(destination_path); free(disk_path); - return false; + return FILE_SAVE_ERROR; } /* --existing checks the final destination, not a temporary partial path. */ @@ -75,7 +80,7 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi free(confined_partial); free(destination_path); free(disk_path); - return true; + return FILE_SAVE_SKIPPED; } /* --ignore-existing checks the final destination before partial files or @@ -87,34 +92,40 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi free(confined_partial); free(destination_path); free(disk_path); - return true; + return FILE_SAVE_SKIPPED; } } - free(destination_path); - destination_path = NULL; /* --update is receiver-side policy: never replace a newer destination. - The secure stat does not require read permission on the destination. */ - if (config && config->update && file_destination_is_newer_secure(disk_path, file->metadata)) { + In partial-dir mode the entry that would be replaced is the real + destination, not the temporary partial file. The secure stat does not + require read permission on the destination. */ + const char* update_target = use_partial_root ? destination_path : disk_path; + if (config && config->update && file_destination_is_newer_secure(update_target, file->metadata)) { free(confined_backup); free(confined_partial); + free(destination_path); free(disk_path); - return true; + return FILE_SAVE_SKIPPED; } if (backup_enabled) { + /* Back up the entry that the incoming write will replace. When writing + through a partial dir the pre-existing destination file is the one to + preserve; any stale partial file is overwritten without a backup. */ + const char* replace_target = use_partial_root ? destination_path : disk_path; struct stat backup_stat; - if (file_stat_secure(disk_path, &backup_stat)) { + if (file_stat_secure(replace_target, &backup_stat)) { if (backup_dir) { backup_path = path_cat(confined_backup, file->path); } else { - size_t path_len = strlen(disk_path); + size_t path_len = strlen(replace_target); size_t suffix_len = strlen(backup_suffix); if (path_len > SIZE_MAX - suffix_len - 1) goto fail; backup_path = malloc(path_len + suffix_len + 1); if (backup_path) { - memcpy(backup_path, disk_path, path_len); + memcpy(backup_path, replace_target, path_len); memcpy(backup_path + path_len, backup_suffix, suffix_len + 1); } } @@ -125,7 +136,7 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi goto fail; free(parent_copy); parent_copy = NULL; - if (!file_rename_secure(disk_path, backup_path)) + if (!file_rename_secure(replace_target, backup_path)) goto fail; free(backup_path); backup_path = NULL; @@ -149,13 +160,25 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi : file_to_disk_secure_with_fsync(disk_path, file->data->data, file->data->size, inplace, sparse, metadata, preserve_executability, config && config->use_fsync); + if (!ok) + goto fail; + + /* --partial --partial-dir writes the complete file under the partial dir so + interrupted transfers leave a resumable copy there. Once the file is + fully written it must be atomically installed at the real destination; + otherwise completed transfers would linger under the partial dir. */ + if (use_partial_root) { + if (!file_rename_secure(disk_path, destination_path)) + goto fail; + } + free(parent_copy); free(backup_path); free(confined_backup); free(confined_partial); free(destination_path); free(disk_path); - return ok; + return FILE_SAVE_WRITTEN; fail: free(parent_copy); @@ -164,13 +187,15 @@ fail: free(confined_partial); free(destination_path); free(disk_path); - return false; + return FILE_SAVE_ERROR; } static File* receive_delta_file(int fd, const Config* config, const char* check_path, void* old_data, unsigned long long old_size, bool* failed) { - if (!old_data) + if (!old_data) { + *failed = true; return NULL; + } DeltaSignature* sig = delta_signature_create(old_data, old_size, config->delta_block_size); if (!sig) { @@ -206,7 +231,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_ } if (resp == STATUS_DELTA_DATA) { - Data* delta_data = receive_data_limited(fd, MAX_RECEIVE_FILE_SIZE); + Data* delta_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE); if (!delta_data) { delta_signature_destroy(sig); free(old_data); @@ -219,7 +244,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_ !compression_should_skip_with_suffixes( check_path, config->skip_compress_suffixes, config->skip_compress_set ? config->skip_compress_count : -1)) { - raw_delta = data_decompress_limited(delta_data, MAX_RECEIVE_FILE_SIZE); + raw_delta = data_decompress_limited(delta_data, MAX_RECEIVE_WHOLE_FILE_SIZE); data_destroy(delta_data); if (!raw_delta) { free(old_data); @@ -239,11 +264,12 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_ } uint64_t new_size = delta->new_file_size; - if (new_size > MAX_RECEIVE_FILE_SIZE || new_size > SIZE_MAX) { + if (new_size > MAX_RECEIVE_WHOLE_FILE_SIZE || new_size > SIZE_MAX) { delta_destroy(delta); free(old_data); delta_signature_destroy(sig); send_status(fd, STATUS_ERROR); + *failed = true; return NULL; } void* new_data = delta_apply(old_data, old_size, delta, config->delta_block_size); @@ -284,6 +310,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_ free(old_data); delta_signature_destroy(sig); send_status(fd, STATUS_ERROR); + *failed = true; return NULL; } data_destroy(file->data); @@ -314,7 +341,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_ } } - Data* file_data = receive_data_limited(fd, MAX_RECEIVE_FILE_SIZE); + Data* file_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE); if (file_data == NULL) { file_destroy(file); *failed = true; @@ -325,7 +352,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_ !compression_should_skip_with_suffixes( file->path, config->skip_compress_suffixes, config->skip_compress_set ? config->skip_compress_count : -1)) { - Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_FILE_SIZE); + Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE); data_destroy(file_data); if (uncompressed == NULL) { file_destroy(file); @@ -384,7 +411,7 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { return NULL; } - if (check_size > MAX_RECEIVE_FILE_SIZE) { + if (check_size > MAX_RECEIVE_WHOLE_FILE_SIZE) { free(check_path); send_status(fd, STATUS_ERROR); return NULL; @@ -405,6 +432,9 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { send_status(fd, STATUS_ERROR); return NULL; } + + /* Open the existing destination entry (if any) once and keep the descriptor + until the quick-check below decides whether the old contents are needed. */ struct stat st; bool has_old_file = false; int old_fd = -1; @@ -416,9 +446,34 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { close(parent_fd); has_old_file = old_fd >= 0 && fstat(old_fd, &st) == 0 && S_ISREG(st.st_mode); } + if (!has_old_file && old_fd >= 0) { + close(old_fd); + old_fd = -1; + } unsigned long long old_size = has_old_file ? (unsigned long long)st.st_size : 0; + + /* Decide from metadata alone whether the receiver already holds the file + the sender is offering. The old contents are only read into memory when + a checksum comparison or a delta transfer actually requires them. */ + bool size_equal = has_old_file && old_size == check_size; + bool match_by_metadata = false; + if (size_equal && !config->ignore_times && !config->size_only) { + long long old_mtime_nsec = 0; +#ifdef __linux__ + old_mtime_nsec = st.st_mtim.tv_nsec; +#endif + match_by_metadata = metadata_mtime_matches(st.st_mtime, old_mtime_nsec, (time_t)check_mtime, + (long)check_mtime_nsec, config->modify_window); + } + + bool try_delta = config->use_delta && !config->whole_file && has_old_file && + delta_should_attempt(old_size, check_size, config->delta_max_file_size); + bool checksum_needs_read = size_equal && !config->ignore_times && config->checksum; + bool need_old_data = checksum_needs_read || try_delta; + void* old_data = NULL; - if (has_old_file && old_size > 0 && old_size <= MAX_RECEIVE_FILE_SIZE && old_size <= SIZE_MAX) { + if (need_old_data && has_old_file && old_size > 0 && old_size <= MAX_RECEIVE_WHOLE_FILE_SIZE && + old_size <= SIZE_MAX) { old_data = protocol_alloc((size_t)old_size); if (old_data) { size_t got = 0; @@ -433,73 +488,63 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { } } } - if (old_fd >= 0) { - close(old_fd); - } - bool match = - !config->ignore_times && has_old_file && (unsigned long long)st.st_size == check_size; - if (match && config->checksum) { - uint64_t old_checksum = old_size == 0 ? delta_xxhash64("", 0) : 0; - if (old_data) - old_checksum = delta_xxhash64(old_data, (size_t)old_size); - match = (old_size == 0 || old_data) && old_checksum == check_checksum; - free(old_data); - old_data = NULL; - } else if (match && !config->size_only) { - long long old_mtime_nsec = 0; -#ifdef __linux__ - old_mtime_nsec = st.st_mtim.tv_nsec; -#endif - match = metadata_mtime_matches(st.st_mtime, old_mtime_nsec, (time_t)check_mtime, - (long)check_mtime_nsec, config->modify_window); + /* Quick-skip decision. If no content comparison is required this is final + and the old file was never read; if the read failed the file is not + skipped and the transfer proceeds with the full new contents. */ + bool match = false; + if (checksum_needs_read) { + if (old_size == 0) + match = delta_xxhash64("", 0) == check_checksum; + else + match = old_data != NULL && delta_xxhash64(old_data, (size_t)old_size) == check_checksum; + } else if (size_equal && !config->ignore_times) { + match = config->size_only || match_by_metadata; } if (match) { free(old_data); if (!send_status(fd, STATUS_OK)) { + close(old_fd); free(full_path); free(check_path); return NULL; } + close(old_fd); free(full_path); free(check_path); *skipped = true; return NULL; } - bool try_delta = config->use_delta && !config->whole_file && has_old_file && old_data != NULL && - delta_should_attempt(old_size, check_size, config->delta_max_file_size); - - if (try_delta) { + if (try_delta && old_data != NULL) { bool delta_failed = false; File* delta_file = receive_delta_file(fd, config, check_path, old_data, old_size, &delta_failed); old_data = NULL; /* receive_delta_file consumes the snapshot on every path */ if (delta_file) { + close(old_fd); free(full_path); free(check_path); return delta_file; } if (delta_failed) { + close(old_fd); free(full_path); free(check_path); return NULL; } - free(old_data); - old_data = NULL; - try_delta = false; } + free(old_data); + old_data = NULL; - if (!try_delta) { - free(old_data); - old_data = NULL; - if (!send_status(fd, STATUS_NEXT)) { - free(full_path); - free(check_path); - return NULL; - } + if (!send_status(fd, STATUS_NEXT)) { + close(old_fd); + free(full_path); + free(check_path); + return NULL; } + close(old_fd); File* file = file_create(check_path); free(check_path); @@ -517,7 +562,7 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { } } - Data* file_data = receive_data_limited(fd, MAX_RECEIVE_FILE_SIZE); + Data* file_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE); if (file_data == NULL) { file_destroy(file); return NULL; @@ -527,7 +572,7 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes, config->skip_compress_set ? config->skip_compress_count : -1)) { - Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_FILE_SIZE); + Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE); data_destroy(file_data); if (uncompressed == NULL) { file_destroy(file); @@ -571,7 +616,7 @@ File* file_receive(const Config* config, int file_descriptor) { return NULL; } } - Data* file_data = receive_data_limited(file_descriptor, MAX_RECEIVE_FILE_SIZE); + Data* file_data = receive_data_limited(file_descriptor, MAX_RECEIVE_WHOLE_FILE_SIZE); if (file_data == NULL) { file_destroy(file); return NULL; @@ -580,7 +625,7 @@ File* file_receive(const Config* config, int file_descriptor) { !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes, config->skip_compress_set ? config->skip_compress_count : -1)) { - Data* file_data_uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_FILE_SIZE); + Data* file_data_uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE); data_destroy(file_data); if (file_data_uncompressed == NULL) { file_destroy(file); diff --git a/src/shared/file_receive.h b/src/shared/file_receive.h index 5b2d2a4..179bdfd 100644 --- a/src/shared/file_receive.h +++ b/src/shared/file_receive.h @@ -10,6 +10,14 @@ File* file_receive(const Config* config, int file_descriptor); File* receive_incremental_check(int fd, const Config* config, bool* skipped); int receive_manifest(int fd, const Config* config, int* next_status); + +/* Outcome of a single file_save_to_disk operation. The receiver needs to + distinguish "written" from "skipped" so --remove-source-files can be told + which sources were actually stored. */ +typedef enum { FILE_SAVE_ERROR = 0, FILE_SAVE_WRITTEN = 1, FILE_SAVE_SKIPPED = 2 } FileSaveResult; + +FileSaveResult file_save_to_disk_full(const char* root_directory, const File* file, + const Config* config); bool file_save_to_disk(const char* root_directory, const File* file, const Config* config); #endif diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index cc357e1..ace0c9d 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -105,6 +105,9 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* context->queue = queue; context->file_descriptor = file_descriptor; context->ssl = ssl; + context->outcomes.entries = NULL; + context->outcomes.count = 0; + context->outcomes.capacity = 0; protocol_session_init(&context->session, file_descriptor, file_descriptor); protocol_session_set_ssl(&context->session, ssl); context->receiver_done = false; @@ -137,6 +140,7 @@ fail: void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { config_delete(context->config); queue_destroy(context->queue); + receiver_outcomes_destroy(&context->outcomes); mtx_destroy(&context->mutex); cnd_destroy(&context->condition_not_full); cnd_destroy(&context->condition_not_empty); @@ -170,7 +174,7 @@ int receive_thread(void* pipeline_context) { const Config* config = context->config; mtx_unlock(&context->mutex); - ReceiverSink sink = {receiver_enqueue_file, context, false, false}; + ReceiverSink sink = {receiver_enqueue_file, context, false, false, NULL}; if (receiver_process((Config*)config, file_descriptor, &sink) != 0) { receiver_thread_fail(context); protocol_session_unbind(); @@ -211,7 +215,26 @@ int write_thread(void* pipeline_context) { protocol_session_unbind(); return thrd_success; } - if (save_to_disk && !file_save_to_disk(root_directory, file, context->config)) { + FileSaveResult result = FILE_SAVE_SKIPPED; + if (save_to_disk) { + result = file_save_to_disk_full(root_directory, file, context->config); + if (result == FILE_SAVE_ERROR) { + file_destroy(file); + mtx_lock(&context->mutex); + atomic_store(&context->cancelled, true); + context->receiver_done = true; + cnd_broadcast(&context->condition_not_full); + cnd_broadcast(&context->condition_not_empty); + mtx_unlock(&context->mutex); + free(root_directory); + protocol_session_unbind(); + return thrd_error; + } + } + /* Record the per-file outcome so a --remove-source-files sender learns + which sources were actually written versus skipped on the receiver. */ + if (context->config->remove_source_files && + !receiver_outcomes_append(&context->outcomes, (unsigned char)result)) { file_destroy(file); mtx_lock(&context->mutex); atomic_store(&context->cancelled, true); diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index 5431c1f..d0e01ea 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -9,6 +9,7 @@ #include "file.h" #include "protocol.h" #include "queue.h" +#include "receiver.h" #include typedef struct { @@ -40,6 +41,7 @@ typedef struct PipelineContextReceiver { int file_descriptor; SSL* ssl; ProtocolSession session; + ReceiverOutcomes outcomes; mtx_t mutex; cnd_t condition_not_full; cnd_t condition_not_empty; diff --git a/src/shared/protocol.h b/src/shared/protocol.h index 54b3326..ff79e25 100644 --- a/src/shared/protocol.h +++ b/src/shared/protocol.h @@ -9,10 +9,16 @@ /* Maximum allowed string size for receive_str (64 KB) */ #define MAX_STRING_SIZE (64 * 1024) -/* Maximum allowed data payload size for receive_data (100 MB) */ -#define MAX_DATA_PAYLOAD_SIZE (100ULL * 1024 * 1024) -/* Maximum uncompressed file payload accepted by the receiver. */ -#define MAX_RECEIVE_FILE_SIZE (64ULL * 1024 * 1024) +/* Maximum uncompressed file payload accepted by the receiver's whole-file + * paths. A single whole file is charged against the per-connection memory + * reservation (MAX_CONNECTION_MEMORY) and against the server allocation + * ceiling (MAX_SERVER_ALLOC), so this mirrors those 256 MB bounds rather than + * the older 64 MB chunk-era cap. Chunk-serialized payloads keep their own + * 64 MB cap (MAX_CHUNK_SIZE). */ +#define MAX_RECEIVE_WHOLE_FILE_SIZE (256ULL * 1024 * 1024) + +/* Maximum allowed data payload size for receive_data (whole-file bound) */ +#define MAX_DATA_PAYLOAD_SIZE MAX_RECEIVE_WHOLE_FILE_SIZE /* Maximum chunk size (64 MB) — prevents unbounded allocation from the wire */ #define MAX_CHUNK_SIZE (64ULL * 1024 * 1024) diff --git a/tests/test_config.c b/tests/test_config.c index 6a41eb4..6a06791 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -282,6 +282,105 @@ static void test_config_receive_truncated() { close(p[1]); } +static bool config_string_roundtrip_matches(const Config* send_cfg, Config* recv) { + /* The sender serializes NULL strings as "" on the wire. Receivers must + canonicalize those empty values back to NULL for the options whose client + default is NULL (backup_dir, temp_dir, partial_dir, suffix), while a real + non-empty value round-trips unchanged. */ + const char* fields[4]; + char* const* recv_fields[4]; + fields[0] = send_cfg->backup_dir; + recv_fields[0] = &recv->backup_dir; + fields[1] = send_cfg->temp_dir; + recv_fields[1] = &recv->temp_dir; + fields[2] = send_cfg->partial_dir; + recv_fields[2] = &recv->partial_dir; + fields[3] = send_cfg->suffix; + recv_fields[3] = &recv->suffix; + for (int i = 0; i < 4; i++) { + const char* sent = fields[i]; + const char* got = *recv_fields[i]; + if (sent == NULL || sent[0] == '\0') { + if (got != NULL) + return false; + } else if (got == NULL || strcmp(sent, got) != 0) { + return false; + } + } + return true; +} + +static bool roundtrip_config_ok(const Config* send_cfg) { + int p[2]; + if (socketpair(AF_UNIX, SOCK_STREAM, 0, p) != 0) + return false; + pid_t pid = fork(); + if (pid == 0) { + close(p[1]); + io_set_fds(p[0], p[0]); + Config* recv = config_receive(p[0]); + bool ok = recv != NULL; + if (ok) { + ok = recv->version != NULL && strcmp(recv->version, PROTOCOL_VERSION) == 0; + ok = ok && recv->send_directory && recv->receive_root_directory; + ok = ok && config_string_roundtrip_matches(send_cfg, recv); + } + config_delete(recv); + close(p[0]); + _exit(ok ? 0 : 1); + } else { + close(p[0]); + io_set_fds(p[1], p[1]); + bool sent = config_send(p[1], send_cfg); + int status; + waitpid(pid, &status, 0); + close(p[1]); + return sent && WIFEXITED(status) && WEXITSTATUS(status) == 0; + } +} + +/* Issue #252: NULL-vs-empty must survive the wire for backup_dir, temp_dir, + partial_dir, and suffix. NULL and explicitly-empty client values are both + serialized as "" and must be reconstructed as NULL so plain --backup (with + no --suffix/--backup-dir) works exactly like the client configured it. */ +static void test_config_string_null_vs_empty_roundtrip() { + if (is_running_under_valgrind()) + return; + + /* NULL values on the wire must come back as NULL. */ + Config* a = config_create(); + EXPECT_NOT_NULL(a); + a->send_directory = str_dup("/src"); + a->receive_root_directory = str_dup("/dst"); + EXPECT_TRUE(roundtrip_config_ok(a)); + config_delete(a); + + /* Explicitly empty strings (indistinguishable on the wire from NULL) must + be canonicalized to NULL by the receiver. */ + Config* b = config_create(); + EXPECT_NOT_NULL(b); + b->send_directory = str_dup("/src"); + b->receive_root_directory = str_dup("/dst"); + b->backup_dir = str_dup(""); + b->temp_dir = str_dup(""); + b->partial_dir = str_dup(""); + b->suffix = str_dup(""); + EXPECT_TRUE(roundtrip_config_ok(b)); + config_delete(b); + + /* Non-empty values must round-trip unchanged. */ + Config* c = config_create(); + EXPECT_NOT_NULL(c); + c->send_directory = str_dup("/src"); + c->receive_root_directory = str_dup("/dst"); + c->backup_dir = str_dup("backups"); + c->temp_dir = str_dup("/tmp/fast"); + c->partial_dir = str_dup(".partial"); + c->suffix = str_dup(".bak"); + 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")); @@ -315,6 +414,7 @@ void test_config() { test_config_send_receive(); test_config_send_receive_version_mismatch(); test_config_receive_truncated(); + test_config_string_null_vs_empty_roundtrip(); } test_config_is_remote_dest(); } diff --git a/tests/test_file.c b/tests/test_file.c index 787a60e..0d416be 100644 --- a/tests/test_file.c +++ b/tests/test_file.c @@ -9,6 +9,7 @@ #include #include #include +#include #include static void test_file_create() { @@ -250,6 +251,111 @@ static void test_file_save_to_disk_ignore_existing_entry_types() { rmdir(root); } +/* Issue #253: with --partial --partial-dir a completed write must be installed + at the real destination rather than left under the partial directory. */ +static void test_file_save_to_disk_partial_install() { + const char* root = "test_partial_install_tmp"; + const char* dest_file = "test_partial_install_tmp/file.txt"; + const char* partial_file = "test_partial_install_tmp/.partial/file.txt"; + unlink(dest_file); + unlink(partial_file); + rmdir("test_partial_install_tmp/.partial"); + rmdir(root); + + File* f = file_create("file.txt"); + EXPECT_NOT_NULL(f); + const char* content = "partial-dir content"; + f->data->data = malloc(strlen(content)); + EXPECT_NOT_NULL(f->data->data); + memcpy(f->data->data, content, strlen(content)); + f->data->size = strlen(content); + + Config* config = config_create(); + EXPECT_NOT_NULL(config); + config->partial = true; + config->partial_dir = str_dup(".partial"); + + EXPECT_EQ_INT(file_save_to_disk_full(root, f, config), FILE_SAVE_WRITTEN); + + FILE* fp = fopen(dest_file, "rb"); + EXPECT_NOT_NULL(fp); + // cppcheck-suppress knownConditionTrueFalse + if (fp) { + char buf[64] = {0}; + size_t nread = fread(buf, 1, sizeof(buf) - 1, fp); + fclose(fp); + EXPECT_EQ_INT((int)nread, (int)strlen(content)); + EXPECT_EQ_INT(memcmp(buf, content, strlen(content)), 0); + } + /* A completed transfer must not linger under the partial dir. */ + EXPECT_EQ_INT(access(partial_file, F_OK), -1); + + file_destroy(f); + config_delete(config); + unlink(dest_file); + rmdir(root); +} + +/* Issue #251: file_save_to_disk_full must distinguish receiver-side skips + (--existing/--ignore-existing/--update) from real writes so the sender can + decide whether --remove-source-files may unlink its source. */ +static void test_file_save_to_disk_reports_skips() { + const char* root = "test_save_skip_tmp"; + const char* existing_path = "test_save_skip_tmp/existing.txt"; + unlink(existing_path); + rmdir(root); + EXPECT_TRUE(file_write_to_disk(existing_path, "old", 3, false, false)); + + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + + File* new_file = file_create("missing.txt"); + EXPECT_NOT_NULL(new_file); + new_file->data->data = malloc(7); + EXPECT_NOT_NULL(new_file->data->data); + memcpy(new_file->data->data, "skipped", 7); + new_file->data->size = 7; + + /* --existing: destination is missing -> skipped, not an error. */ + cfg->existing = true; + EXPECT_EQ_INT(file_save_to_disk_full(root, new_file, cfg), FILE_SAVE_SKIPPED); + cfg->existing = false; + + /* --ignore-existing: destination present -> skipped. */ + File* present = file_create("existing.txt"); + EXPECT_NOT_NULL(present); + present->data->data = malloc(3); + EXPECT_NOT_NULL(present->data->data); + memcpy(present->data->data, "new", 3); + present->data->size = 3; + cfg->ignore_existing = true; + EXPECT_EQ_INT(file_save_to_disk_full(root, present, cfg), FILE_SAVE_SKIPPED); + cfg->ignore_existing = false; + + /* A normal overwrite of an existing file is a real write. */ + EXPECT_EQ_INT(file_save_to_disk_full(root, present, cfg), FILE_SAVE_WRITTEN); + + /* --update: a newer destination is skipped. */ + struct stat st; + EXPECT_EQ_INT(stat(existing_path, &st), 0); + time_t now = time(NULL); + FileMetadata metadata = {.mode = st.st_mode, + .uid = st.st_uid, + .gid = st.st_gid, + .mtime_sec = now - 100, + .mtime_nsec = 0}; + present->metadata = &metadata; + cfg->update = true; + EXPECT_EQ_INT(file_save_to_disk_full(root, present, cfg), FILE_SAVE_SKIPPED); + present->metadata = NULL; + + file_destroy(new_file); + file_destroy(present); + config_delete(cfg); + unlink(existing_path); + rmdir(root); +} + static void test_file_write_to_disk_basic() { const char* content = "Basic file_write_to_disk test"; EXPECT_TRUE(file_write_to_disk("test_file_write_to_disk_basic.txt", content, strlen(content), @@ -635,6 +741,8 @@ void test_file() { test_file_save_to_disk_existing(); test_file_save_to_disk_ignore_existing(); test_file_save_to_disk_ignore_existing_entry_types(); + test_file_save_to_disk_partial_install(); + test_file_save_to_disk_reports_skips(); test_file_write_to_disk_basic(); test_file_write_to_disk_with_fsync(); test_file_write_to_disk_creates_dirs(); diff --git a/tests/test_server.c b/tests/test_server.c index d17137c..bbcd76d 100644 --- a/tests/test_server.c +++ b/tests/test_server.c @@ -1,14 +1,18 @@ #include "test_server.h" #include "config.h" +#include "delta.h" #include "file.h" #include "protocol.h" #include "test_utils.h" #include "utils.h" +#include #include #include #include #include +#include #include +#include #include #include "receiver.h" @@ -210,6 +214,249 @@ static void test_receive_incremental_check_rejects_invalid_nanoseconds() { config_delete(cfg); } +static char* make_check_root(const char* tag) { + char tmpl[128]; + snprintf(tmpl, sizeof(tmpl), "/tmp/fastsync_%s_XXXXXX", tag); + char* path = str_dup(tmpl); + if (!path) + return NULL; + if (!mkdtemp(path)) { + free(path); + return NULL; + } + return path; +} + +static void write_check_file(const char* dir, const char* name, const char* content) { + char path[1024]; + snprintf(path, sizeof(path), "%s/%s", dir, name); + int fd = open(path, O_WRONLY | O_CREAT | O_TRUNC, 0644); + if (fd >= 0) { + size_t len = strlen(content); + if (write(fd, content, len) != (ssize_t)len) { + /* intentionally ignored in tests */ + } + close(fd); + } +} + +/* Issue #255: a same-size/mtime match is decided from metadata alone, so the + receiver answers STATUS_OK (skip) and never asks for a data body. */ +static void test_incremental_check_quick_skip_by_mtime() { + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + char* root = make_check_root("qskip"); + EXPECT_NOT_NULL(root); + cfg->receive_root_directory = str_dup(root); + write_check_file(root, "file.txt", "0123456789abcdef"); + + char path[1024]; + snprintf(path, sizeof(path), "%s/file.txt", root); + struct stat st; + EXPECT_EQ_INT(stat(path, &st), 0); + + int p[2]; + EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0); + io_set_fds(p[0], p[1]); + io_set_bwlimit(0); + + pid_t pid = fork(); + if (pid == 0) { + alarm(30); + close(p[1]); + io_set_fds(p[0], p[0]); + bool skipped = false; + File* file = receive_incremental_check(p[0], cfg, &skipped); + bool ok = file == NULL && skipped; + file_destroy(file); + config_delete(cfg); + close(p[0]); + _exit(ok ? 0 : 1); + } else { + close(p[0]); + io_set_fds(p[1], p[1]); + EXPECT_TRUE(send_str(p[1], "file.txt")); + unsigned long long size = (unsigned long long)st.st_size; + long long mtime = (long long)st.st_mtime; + long long mtime_nsec = 0; +#ifdef __linux__ + mtime_nsec = (long long)st.st_mtim.tv_nsec; +#endif + EXPECT_TRUE(send_n_data(p[1], &size, sizeof(size))); + EXPECT_TRUE(send_n_data(p[1], &mtime, sizeof(mtime))); + EXPECT_TRUE(send_n_data(p[1], &mtime_nsec, sizeof(mtime_nsec))); + Status s; + EXPECT_TRUE(receive_status(p[1], &s)); + EXPECT_EQ_INT(s, STATUS_OK); + + int status; + waitpid(pid, &status, 0); + close(p[1]); + config_delete(cfg); + unlink(path); + rmdir(root); + free(root); + EXPECT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0); + } +} + +/* Issue #255: a size mismatch cannot be a skip, so the receiver answers + STATUS_NEXT and consumes the full data body that follows. */ +static void test_incremental_check_size_mismatch_full_transfer() { + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + char* root = make_check_root("qnext"); + EXPECT_NOT_NULL(root); + cfg->receive_root_directory = str_dup(root); + write_check_file(root, "file.txt", "0123456789abcdef"); + + char path[1024]; + snprintf(path, sizeof(path), "%s/file.txt", root); + struct stat st; + EXPECT_EQ_INT(stat(path, &st), 0); + + int p[2]; + EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0); + io_set_fds(p[0], p[1]); + io_set_bwlimit(0); + + pid_t pid = fork(); + if (pid == 0) { + alarm(30); + close(p[1]); + io_set_fds(p[0], p[0]); + bool skipped = false; + File* file = receive_incremental_check(p[0], cfg, &skipped); + bool ok = file != NULL && !skipped && file->path != NULL && strcmp(file->path, "file.txt") == 0; + file_destroy(file); + config_delete(cfg); + close(p[0]); + _exit(ok ? 0 : 1); + } else { + close(p[0]); + io_set_fds(p[1], p[1]); + EXPECT_TRUE(send_str(p[1], "file.txt")); + unsigned long long size = (unsigned long long)st.st_size + 1; + long long mtime = (long long)st.st_mtime; + long long mtime_nsec = 0; +#ifdef __linux__ + mtime_nsec = (long long)st.st_mtim.tv_nsec; +#endif + EXPECT_TRUE(send_n_data(p[1], &size, sizeof(size))); + EXPECT_TRUE(send_n_data(p[1], &mtime, sizeof(mtime))); + EXPECT_TRUE(send_n_data(p[1], &mtime_nsec, sizeof(mtime_nsec))); + Status s; + EXPECT_TRUE(receive_status(p[1], &s)); + EXPECT_EQ_INT(s, STATUS_NEXT); + + Data* body = data_create_reserve(8); + EXPECT_NOT_NULL(body); + body->data = malloc(8); + EXPECT_NOT_NULL(body->data); + memcpy(body->data, "replaced", 8); + body->size = 8; + EXPECT_TRUE(send_data(p[1], body)); + data_destroy(body); + + int status; + waitpid(pid, &status, 0); + close(p[1]); + config_delete(cfg); + unlink(path); + rmdir(root); + free(root); + EXPECT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0); + } +} + +/* Issue #256: when a received delta claims a result above the whole-file cap, + receive_delta_file must mark the operation failed so the caller aborts with + STATUS_ERROR instead of emitting STATUS_NEXT and waiting for a body that + never arrives. */ +static void test_incremental_check_delta_oversize_reports_failure() { + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + char* root = make_check_root("qdelta"); + EXPECT_NOT_NULL(root); + cfg->receive_root_directory = str_dup(root); + cfg->use_delta = true; + + char content[20000]; + memset(content, 'a', sizeof(content)); + content[sizeof(content) - 1] = '\0'; + write_check_file(root, "file.txt", content); + + char path[1024]; + snprintf(path, sizeof(path), "%s/file.txt", root); + struct stat st; + EXPECT_EQ_INT(stat(path, &st), 0); + EXPECT_EQ_INT((int)st.st_size, 19999); + + int p[2]; + EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0); + io_set_fds(p[0], p[1]); + io_set_bwlimit(0); + + pid_t pid = fork(); + if (pid == 0) { + alarm(30); + close(p[1]); + io_set_fds(p[0], p[0]); + bool skipped = false; + File* file = receive_incremental_check(p[0], cfg, &skipped); + bool ok = file == NULL && !skipped; + if (ok) + send_status(p[0], STATUS_ERROR); /* mirror the server error path */ + file_destroy(file); + config_delete(cfg); + close(p[0]); + _exit(ok ? 0 : 1); + } else { + close(p[0]); + io_set_fds(p[1], p[1]); + EXPECT_TRUE(send_str(p[1], "file.txt")); + unsigned long long size = (unsigned long long)st.st_size; + long long mtime = 1; /* different from the file mtime: force a transfer */ + long long mtime_nsec = 0; + EXPECT_TRUE(send_n_data(p[1], &size, sizeof(size))); + EXPECT_TRUE(send_n_data(p[1], &mtime, sizeof(mtime))); + EXPECT_TRUE(send_n_data(p[1], &mtime_nsec, sizeof(mtime_nsec))); + + Status s; + EXPECT_TRUE(receive_status(p[1], &s)); + EXPECT_EQ_INT(s, STATUS_DELTA_SIGNATURE); + Data* sig_data = receive_data(p[1]); + EXPECT_NOT_NULL(sig_data); + DeltaSignature* sig = delta_signature_deserialize(sig_data); + EXPECT_NOT_NULL(sig); + delta_signature_destroy(sig); + data_destroy(sig_data); + + /* Send a delta whose claimed output size exceeds the whole-file cap. */ + Delta delta; + memset(&delta, 0, sizeof(delta)); + delta.new_file_size = MAX_RECEIVE_WHOLE_FILE_SIZE + 1; + Data* bogus = delta_serialize(&delta); + EXPECT_NOT_NULL(bogus); + EXPECT_TRUE(send_status(p[1], STATUS_DELTA_DATA)); + EXPECT_TRUE(bogus != NULL && send_data(p[1], bogus)); + data_destroy(bogus); + + /* The receiver must answer with an error, never with STATUS_NEXT. */ + EXPECT_TRUE(receive_status(p[1], &s)); + EXPECT_EQ_INT(s, STATUS_ERROR); + + int status; + waitpid(pid, &status, 0); + close(p[1]); + config_delete(cfg); + unlink(path); + rmdir(root); + free(root); + EXPECT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0); + } +} + void test_server() { if (!is_running_under_valgrind()) { test_receive_files_finished(); @@ -217,5 +464,8 @@ void test_server() { test_receive_files_abort(); test_receive_manifest_rejects_traversal(); test_receive_incremental_check_rejects_invalid_nanoseconds(); + test_incremental_check_quick_skip_by_mtime(); + test_incremental_check_size_mismatch_full_transfer(); + test_incremental_check_delta_oversize_reports_failure(); } } From 3704a6adaaf92a667c91a352a8da4a98daaddb9a Mon Sep 17 00:00:00 2001 From: TapTap Date: Sat, 5 Sep 2026 12:46:47 +0200 Subject: [PATCH 2/2] test: integration coverage for #251 #252 #253 #257 - remove-source-files keeps sources skipped by --existing/--ignore-existing/ --update (rsync reference behavior) - --backup keeps ~, --suffix .bak, and --backup-dir backups - --partial --partial-dir installs completed files in the destination - a 100 MB file transfers end to end (>64 MiB whole-file cap regression) --- tests/integration/test_features.py | 184 +++++++++++++++++++++++++++++ 1 file changed, 184 insertions(+) diff --git a/tests/integration/test_features.py b/tests/integration/test_features.py index 1d8a5c8..69b8c90 100644 --- a/tests/integration/test_features.py +++ b/tests/integration/test_features.py @@ -1,4 +1,5 @@ """Feature tests: incremental sync, bandwidth limiting, dry run, metadata, filters.""" +import filecmp import os import shutil import sys @@ -813,3 +814,186 @@ class TestBandwidthLimit: mismatches, missing = verify_transfer(SOURCE_DIR, received) assert not missing, f"Missing: {missing}" assert not mismatches, f"Mismatch: {mismatches}" + + +def _read_file(path): + with open(path, "rb") as fh: + return fh.read() + + +class TestRemoveSourceFilesSkips: + """--remove-source-files must not delete sources the receiver skipped + (rsync reference behavior).""" + + def test_existing_first_sync_keeps_new_source(self, shared_server): + source = os.path.join(TEST_DATA_DIR, "remove_rsf_existing_src") + dest = os.path.join(TEST_DATA_DIR, "remove_rsf_existing_dst") + clean_dir(source) + clean_dir(dest) + with open(os.path.join(source, "only.txt"), "wb") as f: + f.write(b"keep me") + + result, _ = run_client(source, dest, flags=["--remove-source-files", "--existing"], + port=shared_server.port) + assert result.returncode == 0, f"Sync failed: {result.stderr[:200]}" + # The file exists only on the source side, so --existing makes the + # receiver skip it; the source must therefore not be removed. + assert os.path.isfile(os.path.join(source, "only.txt")) + received = get_dest_received_dir(dest, source) + assert not os.path.exists(os.path.join(received, "only.txt")) + + def test_ignore_existing_keeps_skipped_source(self, shared_server): + source = os.path.join(TEST_DATA_DIR, "remove_rsf_ignore_src") + dest = os.path.join(TEST_DATA_DIR, "remove_rsf_ignore_dst") + clean_dir(source) + clean_dir(dest) + source_file = os.path.join(source, "file.txt") + with open(source_file, "wb") as f: + f.write(b"payload") + + result, _ = run_client(source, dest, port=shared_server.port) + assert result.returncode == 0 + + result, _ = run_client(source, dest, flags=["--remove-source-files", "--ignore-existing"], + port=shared_server.port) + assert result.returncode == 0, f"Sync failed: {result.stderr[:200]}" + # Destination already has the file, so the second run is a receiver + # skip; the source file must survive. + assert os.path.isfile(source_file) + received = get_dest_received_dir(dest, source) + assert _read_file(os.path.join(received, "file.txt")) == b"payload" + + def test_update_newer_destination_keeps_source(self, shared_server): + source = os.path.join(TEST_DATA_DIR, "remove_rsf_update_src") + dest = os.path.join(TEST_DATA_DIR, "remove_rsf_update_dst") + clean_dir(source) + clean_dir(dest) + source_file = os.path.join(source, "file.txt") + with open(source_file, "wb") as f: + f.write(b"source payload") + + result, _ = run_client(source, dest, port=shared_server.port) + assert result.returncode == 0 + + received = get_dest_received_dir(dest, source) + received_file = os.path.join(received, "file.txt") + with open(received_file, "wb") as f: + f.write(b"newer destination payload") + os.utime(received_file, ns=(time.time_ns() + 10**9, time.time_ns() + 10**9)) + + result, _ = run_client(source, dest, flags=["--remove-source-files", "--update"], + port=shared_server.port) + assert result.returncode == 0, f"Sync failed: {result.stderr[:200]}" + # --update skips a destination that is newer than the source, so the + # source must not be removed. + assert os.path.isfile(source_file) + assert _read_file(received_file) == b"newer destination payload" + + +class TestBackup: + def _sync(self, source, dest, flags, port): + return run_client(source, dest, flags=flags, port=port) + + def test_plain_backup_keeps_previous_version(self, shared_server): + source = os.path.join(TEST_DATA_DIR, "backup_src") + dest = os.path.join(TEST_DATA_DIR, "backup_dst") + clean_dir(source) + clean_dir(dest) + source_file = os.path.join(source, "f.txt") + with open(source_file, "wb") as f: + f.write(b"AAAA") + + result, _ = self._sync(source, dest, ["--backup"], shared_server.port) + assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}" + + with open(source_file, "wb") as f: + f.write(b"BBBB") + result, _ = self._sync(source, dest, ["--backup"], shared_server.port) + assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}" + + received = get_dest_received_dir(dest, source) + assert _read_file(os.path.join(received, "f.txt")) == b"BBBB" + # rsync default suffix "~" keeps the overwritten version. + assert _read_file(os.path.join(received, "f.txt~")) == b"AAAA" + + def test_backup_custom_suffix(self, shared_server): + source = os.path.join(TEST_DATA_DIR, "backup_suffix_src") + dest = os.path.join(TEST_DATA_DIR, "backup_suffix_dst") + clean_dir(source) + clean_dir(dest) + source_file = os.path.join(source, "f.txt") + with open(source_file, "wb") as f: + f.write(b"AAAA") + + flags = ["--backup", "--suffix", ".bak"] + result, _ = self._sync(source, dest, flags, shared_server.port) + assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}" + with open(source_file, "wb") as f: + f.write(b"BBBB") + result, _ = self._sync(source, dest, flags, shared_server.port) + assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}" + + received = get_dest_received_dir(dest, source) + assert _read_file(os.path.join(received, "f.txt")) == b"BBBB" + assert _read_file(os.path.join(received, "f.txt.bak")) == b"AAAA" + + def test_backup_dir_stores_backups_separately(self, shared_server): + source = os.path.join(TEST_DATA_DIR, "backup_dir_src") + dest = os.path.join(TEST_DATA_DIR, "backup_dir_dst") + clean_dir(source) + clean_dir(dest) + source_file = os.path.join(source, "f.txt") + with open(source_file, "wb") as f: + f.write(b"AAAA") + + flags = ["--backup", "--backup-dir", "backups"] + result, _ = self._sync(source, dest, flags, shared_server.port) + assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}" + with open(source_file, "wb") as f: + f.write(b"BBBB") + result, _ = self._sync(source, dest, flags, shared_server.port) + assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}" + + received = get_dest_received_dir(dest, source) + assert _read_file(os.path.join(received, "f.txt")) == b"BBBB" + backup = os.path.join(dest, "backups", os.path.relpath(source_file, os.path.sep)) + assert _read_file(backup) == b"AAAA" + + +class TestPartialDir: + def test_completed_transfer_installed_in_destination(self, shared_server): + source = os.path.join(TEST_DATA_DIR, "partial_src") + dest = os.path.join(TEST_DATA_DIR, "partial_dst") + clean_dir(source) + clean_dir(dest) + source_file = os.path.join(source, "f.txt") + with open(source_file, "wb") as f: + f.write(b"partial payload") + + result, _ = run_client(source, dest, flags=["--partial", "--partial-dir", ".partial"], + port=shared_server.port) + assert result.returncode == 0, f"Partial sync failed: {result.stderr[:200]}" + + received = get_dest_received_dir(dest, source) + assert _read_file(os.path.join(received, "f.txt")) == b"partial payload" + # A completed transfer must not remain under the partial directory. + partial = os.path.join(dest, ".partial", os.path.relpath(source_file, os.path.sep)) + assert not os.path.exists(partial) + + +class TestLargeFile: + def test_transfer_100mb_file(self, shared_server): + source = os.path.join(TEST_DATA_DIR, "large_src") + dest = os.path.join(TEST_DATA_DIR, "large_dst") + clean_dir(source) + clean_dir(dest) + source_file = os.path.join(source, "big.bin") + chunk = os.urandom(1024 * 1024) + with open(source_file, "wb") as f: + for _ in range(100): + f.write(chunk) + + result, _ = run_client(source, dest, port=shared_server.port) + assert result.returncode == 0, f"Large-file sync failed: {result.stderr[:200]}" + received = get_dest_received_dir(dest, source) + assert filecmp.cmp(source_file, os.path.join(received, "big.bin"), shallow=False)