From 4e725517f069850d92d8cb477fbde940a09d5d6a Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 6 Sep 2026 18:53:20 +0200 Subject: [PATCH 1/7] feat: implement rsync delete timing (--delete-before/--delete-during/--delete-delay/--delete-after) Deletion timing is now real and selected by the four rsync flags plus the plain --delete default. Wire protocol bumps to 2.8.0: two new config booleans (delete_during, delete_delay) are serialized and validated, joining the existing delete_before/delete_after. - Early modes (--delete-before, --delete-during/--del): the sender pre-scans the whole tree (paths only), transmits the keep-set manifest BEFORE any file data, and the receiver removes extras and acks STATUS_OK; the sender only streams data after the deletion committed. Deletion is thus performed even if a later transfer phase fails (rsync delete-before/during are destructive by definition). FastSync streams in a single scan so it cannot interleave per-directory like rsync delete-during; --delete-during selects the same engine mode as --delete-before (documented divergence). - Late/commit modes (plain --delete, --delete-after, --delete-delay): the manifest closes the data stream and deletion is committed only after STATUS_FINISHED proves the whole transfer succeeded, preserving FastSync's commit-style safety. --delete-delay converges with --delete-after because FastSync never snapshots the destination during data flow (documented). - The STATUS_MANIFEST frame is now self-delimiting and position-independent. Single-threaded receivers delete before the success frame; the -m receiver hands the keep-set to server.c, which commits the deletion only after the disk writer thread has drained (fixes a delete-vs-in-flight-temp race). - Every timing flag implies --delete; at most one timing flag is allowed. - Each timing flag implies --delete, matching rsync; conflicts are rejected. --- src/client/client_cli.c | 40 +++++----- src/client/client_send.c | 142 +++++++++++++++++++++++++++++---- src/client/client_validation.c | 6 ++ src/client/usage.c | 11 ++- src/server/receiver.c | 82 +++++++++++++++++-- src/server/receiver.h | 8 ++ src/server/server.c | 14 ++++ src/shared/config.c | 31 ++++++- src/shared/config.h | 27 ++++++- src/shared/file_receive.c | 59 ++++++-------- src/shared/file_receive.h | 9 ++- src/shared/multiprocessing.c | 7 +- src/shared/multiprocessing.h | 13 +++ tests/test_server.c | 5 +- 14 files changed, 370 insertions(+), 84 deletions(-) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index e537e5e..86c7995 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -370,7 +370,6 @@ typedef enum { OPT_POS_INT, OPT_NONNEG_INT, OPT_ULL, - OPT_UNSUPPORTED, } OptKind; typedef struct { @@ -430,7 +429,10 @@ static const OptionEntry OPTION_TABLE[] = { {"--old-d", NULL, OPT_FLAG, offsetof(Config, dirs)}, {"--relative", "-R", OPT_FLAG, offsetof(Config, relative)}, {"--mkpath", NULL, OPT_FLAG, offsetof(Config, mkpath)}, - {"--delete-during", "--del", OPT_UNSUPPORTED, 0}, + {"--delete-before", NULL, OPT_FLAG, offsetof(Config, delete_before)}, + {"--delete-during", "--del", OPT_FLAG, offsetof(Config, delete_during)}, + {"--delete-delay", NULL, OPT_FLAG, offsetof(Config, delete_delay)}, + {"--delete-after", NULL, OPT_FLAG, offsetof(Config, delete_after)}, {"--source-dir", NULL, OPT_STRING, offsetof(Config, send_directory)}, {"--dest-dir", NULL, OPT_STRING, offsetof(Config, receive_root_directory)}, @@ -547,8 +549,7 @@ static int apply_negation(Config* config, const char* arg) { return 0; } -static int apply_table_option(Config* config, const OptionEntry* entry, const char* option_name, - const char* value) { +static int apply_table_option(Config* config, const OptionEntry* entry, const char* value) { if (entry->kind == OPT_NOOP) return 0; void* field = (char*)config + entry->offset; @@ -578,11 +579,6 @@ static int apply_table_option(Config* config, const OptionEntry* entry, const ch *(unsigned long long*)field = v; return 0; } - case OPT_UNSUPPORTED: { - const char* reason = "delete-during is not implemented"; - log_message(LOG_LEVEL_ERROR, "%s: %s; refusing to ignore option", option_name, reason); - return -1; - } } return -1; } @@ -665,23 +661,20 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args, if (!entry) entry = find_table_option_with_equals(argv[i], &inline_value); if (entry) { - const char* option_name = argv[i]; const char* value = NULL; if (entry->kind != OPT_FLAG) { - if (entry->kind != OPT_UNSUPPORTED) { - value = inline_value; - if (!value && i + 1 < argc) - value = argv[++i]; - if (!value) { - log_message(LOG_LEVEL_ERROR, "missing argument for %s", entry->name); - return -1; - } + value = inline_value; + if (!value && i + 1 < argc) + value = argv[++i]; + if (!value) { + log_message(LOG_LEVEL_ERROR, "missing argument for %s", entry->name); + return -1; } if (strcmp(entry->name, "--compress-choice") == 0) { if (set_compression_choice(config, value) != 0) return -1; } else { - if (apply_table_option(config, entry, option_name, value) != 0) + if (apply_table_option(config, entry, value) != 0) return -1; if (strcmp(entry->name, "--compress-level") == 0 && (config->compression_level < 1 || config->compression_level > 22)) { @@ -697,11 +690,18 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args, config->use_metadata = true; } } - } else if (apply_table_option(config, entry, option_name, NULL) != 0) { + } else if (apply_table_option(config, entry, NULL) != 0) { return -1; } if (entry->offset == offsetof(Config, eight_bit_output)) protocol_set_8_bit_output(true); + /* A delete-timing flag selects when --delete removes extras, so it + implies --delete exactly like the rsync options do. */ + if (entry->offset == offsetof(Config, delete_before) || + entry->offset == offsetof(Config, delete_during) || + entry->offset == offsetof(Config, delete_delay) || + entry->offset == offsetof(Config, delete_after)) + config->use_delete = true; continue; } diff --git a/src/client/client_send.c b/src/client/client_send.c index ee1b04c..df2c22f 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -254,10 +254,6 @@ static void disconnect_transfer_client(Client* client) { client_delete(client); } -static ArrayList* create_transfer_manifest(const Config* config) { - return config->use_delete ? array_list_create(free) : NULL; -} - static bool add_chunk_to_manifest(ArrayList* manifest, const Chunk* chunk) { if (!manifest) return true; @@ -597,6 +593,53 @@ static int send_delete_manifest(int fd, ArrayList* manifest) { return 0; } +/* Transmit the keep-set manifest and wait for the receiver's verdict. Used by + --delete-before/--delete-during, where the extras are removed on the receiver + BEFORE the first byte of file data is sent: the receiver acknowledges with + STATUS_OK once the bounded delete committed, or STATUS_ERROR if it could not + (in which case the sender aborts without streaming any data). */ +static bool send_delete_manifest_early(Client* client, ArrayList* manifest) { + if (!client || !manifest) + return false; + if (send_delete_manifest(client->file_descriptor, manifest) != 0) + return false; + Status ack; + if (!receive_status(client->file_descriptor, &ack)) + return false; + if (ack != STATUS_OK) { + log_message(LOG_LEVEL_ERROR, "Server failed to delete files before the transfer"); + return false; + } + return true; +} + +/* Walk the whole source tree once collecting only destination-relative wire + paths, loading and sending nothing. --delete-before/--delete-during need the + complete keep-set manifest before the first data byte, so it is built by a + dedicated pre-scan pass and transmitted early; the data pass then re-scans + with a fresh scanner. */ +static bool scan_paths_only(const Config* config, const ScannerOptions* options, + ArrayList* manifest) { + DirectoryScanner* scanner = + directory_scanner_create_with_options(config->send_directory, options); + if (!scanner) + return false; + bool ok = true; + Chunk* chunk; + while ((chunk = directory_scanner_next(scanner)) != NULL) { + if (!add_chunk_to_manifest(manifest, chunk)) { + ok = false; + chunk_destroy(chunk); + break; + } + chunk_destroy(chunk); + } + if (ok && directory_scanner_failed(scanner)) + ok = false; + directory_scanner_destroy(scanner); + return ok; +} + static int incremental_check(Client* client, File* file, const Config* config, DeltaSignature** out_sig) { *out_sig = NULL; @@ -906,6 +949,17 @@ static int send_chunks_multithreaded(void* pipeline_context) { protocol_session_unbind(); return thrd_error; } + if (context->early_delete) { + /* The keep-set manifest was prebuilt by a path-only pre-scan. Transmit it + and wait for the receiver to delete extras before streaming any data. */ + if (!send_delete_manifest_early(client, context->manifest)) { + pipeline_cancel(context); + disconnect_transfer_client(client); + mark_sender_done(context); + protocol_session_unbind(); + return thrd_error; + } + } while (true) { Chunk* current_chunk = queue_dequeue_multithreaded( @@ -919,7 +973,7 @@ static int send_chunks_multithreaded(void* pipeline_context) { protocol_session_unbind(); return thrd_error; } - if (context->config->use_delete) { + if (context->config->use_delete && !context->early_delete) { if (send_delete_manifest(client->file_descriptor, context->manifest) != 0) goto send_fail; } @@ -1013,7 +1067,7 @@ static int scan_directory_multithreaded(void* pipeline_context) { failed = dirs_mode ? directory_scanner_failed(dscanner) : parallel_scanner_failed(scanner); break; } - if (context->config->use_delete) { + if (context->config->use_delete && !context->early_delete) { mtx_lock(&context->mutex_scanner); bool manifest_ok = add_chunk_to_manifest(context->manifest, current_chunk); mtx_unlock(&context->mutex_scanner); @@ -1172,19 +1226,46 @@ int send_files(Config* config) { DirectoryScanner* scanner = NULL; ArrayList* manifest = NULL; ArrayList* remove_sources = NULL; + bool delete_early = config->use_delete && config_delete_timing_early(config); + bool send_failed = false; PreparedScanner prepared; memset(&prepared, 0, sizeof(prepared)); if (!config_send(client->file_descriptor, config)) goto send_fail; if (!prepare_scanner(config, 0, &prepared)) goto send_fail; - scanner = directory_scanner_create_with_options(config->send_directory, &prepared.options); - manifest = create_transfer_manifest(config); if (config->remove_source_files) remove_sources = array_list_create(source_file_destroy); - if (!scanner || (config->use_delete && !manifest) || - (config->remove_source_files && !remove_sources)) + if (config->remove_source_files && !remove_sources) goto send_fail; + /* The late-timing modes (plain --delete / --delete-after / --delete-delay) + build the manifest while streaming and send it after the last data frame. + The early modes (--delete-before/--delete-during) send it up front from a + dedicated path-only pre-scan, so no manifest is kept during the data pass. */ + if (delete_early) { + /* Pass 1: collect the complete keep-set (paths only, no data loaded) and + transmit it now, before any file data. The receiver removes extras and + acks; the transfer aborts here if the deletion could not commit. */ + ArrayList* early_manifest = array_list_create(free); + if (!early_manifest) + goto send_fail; + if (!scan_paths_only(config, &prepared.options, early_manifest)) { + array_list_delete(early_manifest); + goto send_fail; + } + bool early_ok = send_delete_manifest_early(client, early_manifest); + array_list_delete(early_manifest); + if (!early_ok) + goto send_fail; + } else if (config->use_delete) { + manifest = array_list_create(free); + if (!manifest) + goto send_fail; + } + scanner = directory_scanner_create_with_options(config->send_directory, &prepared.options); + if (!scanner) + goto send_fail; + Chunk* current_chunk; unsigned long long total_bytes = 0; int total_files = 0; @@ -1196,7 +1277,7 @@ int send_files(Config* config) { chunk_bytes += current_chunk->items[i]->data->size; total_files++; } - if (!add_chunk_to_manifest(manifest, current_chunk)) { + if (manifest && !add_chunk_to_manifest(manifest, current_chunk)) { chunk_destroy(current_chunk); goto send_fail; } @@ -1220,9 +1301,7 @@ int send_files(Config* config) { if (send_chunk_with_removal(client, current_chunk, config, remove_sources) != 0) { log_message(LOG_LEVEL_ERROR, "Failed to send chunk"); chunk_destroy(current_chunk); - if (manifest) - array_list_delete(manifest); - manifest = NULL; + send_failed = true; break; } total_bytes += chunk_bytes; @@ -1235,9 +1314,18 @@ int send_files(Config* config) { } chunk_destroy(current_chunk); } - if (directory_scanner_failed(scanner) || (config->use_delete && manifest == NULL)) + if (send_failed) { + if (manifest) { + array_list_delete(manifest); + manifest = NULL; + } goto send_fail; - if (config->use_delete) { + } + if (directory_scanner_failed(scanner)) + goto send_fail; + if (manifest) { + /* Late (commit) ordering: all file data is out; transmit the keep-set + manifest so the receiver deletes only after the transfer succeeds. */ if (send_delete_manifest(client->file_descriptor, manifest) != 0) { array_list_delete(manifest); manifest = NULL; @@ -1324,8 +1412,28 @@ int send_files_multithreaded(Config** config_ptr) { return 1; } *config_ptr = NULL; /* context now owns config through all remaining paths */ - if (config->use_delete) + if (config->use_delete) { context->manifest = array_list_create(free); + if (!context->manifest) { + pipeline_context_sender_destroy(context); + return 1; + } + if (config_delete_timing_early(config)) { + /* --delete-before/--delete-during: build the complete keep-set manifest + (paths only, nothing loaded or sent) up front so the sender thread can + transmit it before the first data byte. */ + PreparedScanner prepared; + memset(&prepared, 0, sizeof(prepared)); + bool prebuilt = prepare_scanner(config, 4, &prepared) && + scan_paths_only(config, &prepared.options, context->manifest); + prepared_scanner_destroy(&prepared); + if (!prebuilt) { + pipeline_context_sender_destroy(context); + return 1; + } + context->early_delete = true; + } + } if (config->remove_source_files) context->remove_source_files = array_list_create(source_file_destroy); if ((config->use_delete && !context->manifest) || diff --git a/src/client/client_validation.c b/src/client/client_validation.c index 22ac977..d27fded 100644 --- a/src/client/client_validation.c +++ b/src/client/client_validation.c @@ -71,5 +71,11 @@ bool validate_config(const Config* config) { "staging directory)"); return false; } + if (!config_has_valid_delete_timing(config)) { + log_message(LOG_LEVEL_ERROR, + "--delete-before/--delete-during/--delete-delay/--delete-after select the delete " + "timing; at most one may be given and each implies --delete"); + return false; + } return true; } diff --git a/src/client/usage.c b/src/client/usage.c index e9579ba..5634094 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -24,6 +24,16 @@ void print_usage(void) { printf(" -P Partial mode with progress (retention incomplete)\n"); printf(" -8, --8-bit-output Leave high-bit characters unescaped in output\n"); printf(" --delete Delete files on receiver not in source\n"); + printf(" (default timing: delete only after the whole\n"); + printf(" transfer has succeeded)\n"); + printf(" --delete-before Delete extras before the transfer starts\n"); + printf(" (implies --delete)\n"); + printf(" --delete-during Delete extras once the keep-set is known, before\n"); + printf(" --del data is applied (alias --del; implies --delete)\n"); + printf(" --delete-delay Delete extras only after a successful transfer\n"); + printf(" (implies --delete)\n"); + printf(" --delete-after Alias of the default --delete timing: delete only\n"); + printf(" after the transfer succeeded (implies --delete)\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, -d, --old-dirs, --old-d Transfer the named directory entries without\n"); @@ -37,7 +47,6 @@ void print_usage(void) { printf(" parent directory is not itself listed\n"); printf(" --mkpath Create the destination root directory on the server when it\n"); printf(" does not exist yet\n"); - printf(" --del Alias for --delete-during (not implemented)\n"); printf(" --exclude Exclude files matching pattern\n"); printf(" --include Only include files matching pattern\n"); printf(" --exclude-from Read exclude patterns from file\n"); diff --git a/src/server/receiver.c b/src/server/receiver.c index b203be6..be7a04e 100644 --- a/src/server/receiver.c +++ b/src/server/receiver.c @@ -129,16 +129,33 @@ static bool receiver_process_batch(Config* config, int file_descriptor) { } int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink) { + return receiver_process_pending(config, file_descriptor, sink, NULL); +} + +/* Runs the whole receive loop. The delete manifest may legitimately arrive + either FIRST (--delete-before / --delete-during: the sender transmits the + validated keep-set before any file data) or LAST (plain --delete / + --delete-after / --delete-delay: the manifest closes the data stream). In + the early modes the receiver deletes as soon as the manifest has been read + and acknowledges with STATUS_OK so the sender only starts streaming once the + deletion has committed (or failed); in the late modes the manifest is held + and the deletion is committed only after the terminal STATUS_FINISHED proves + the whole transfer succeeded. See receiver_process_pending() for how the -m + receiver defers that commit until its disk writer has drained. */ +int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink, + ArrayList** pending_manifest) { Status status; if (!receive_status(file_descriptor, &status)) return -1; + bool early_delete = config_delete_timing_early(config); + ArrayList* deferred_manifest = NULL; while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK || status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH || - status == STATUS_MKDIR) { + status == STATUS_MKDIR || status == STATUS_MANIFEST) { if (status == STATUS_KEEPALIVE) { if (!send_status(file_descriptor, STATUS_KEEPALIVE)) return -1; - goto next; + goto next_status; } if (status == STATUS_ABORT) { log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up"); @@ -156,11 +173,45 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si } else if (status == STATUS_CHECK_BATCH) { if (!receiver_process_batch(config, file_descriptor)) return -1; - goto next; + goto next_status; } else if (status == STATUS_MKDIR) { File* dir = file_receive_directory(file_descriptor); if (!dir || !sink->store_file(dir, sink->context)) goto receive_error; + } else if (status == STATUS_MANIFEST) { + ArrayList* manifest = receive_manifest_entries(file_descriptor); + if (!manifest) + return -1; /* receive_manifest_entries already sent STATUS_ERROR */ + if (early_delete) { + /* --delete-before / --delete-during: the manifest is authoritative the + moment it arrives, before any file data. Delete now and acknowledge + so the sender only starts streaming once the deletion committed (or + failed). This is the rsync delete-before/delete-during window: a + later transfer failure does not restore these deletions. */ + bool deletion_ok = config->use_delete ? manifest_delete_extras(config, manifest) : true; + array_list_delete(manifest); + if (!deletion_ok) { + send_status(file_descriptor, STATUS_ERROR); + return -1; + } + if (!send_status(file_descriptor, STATUS_OK)) + return -1; + } else { + /* Plain --delete / --delete-after / --delete-delay: hold the keep-set + and commit the deletion only after STATUS_FINISHED. */ + if (config->use_delete) { + if (deferred_manifest) { + log_message(LOG_LEVEL_ERROR, "Received a second delete manifest"); + array_list_delete(manifest); + send_status(file_descriptor, STATUS_ERROR); + return -1; + } + deferred_manifest = manifest; + } else { + array_list_delete(manifest); + } + } + goto next_status; } else { File* file = file_receive(config, file_descriptor); if (!file) { @@ -170,16 +221,35 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si if (!sink->store_file(file, sink->context)) goto receive_error; } - next: + next_status: if (!receive_status(file_descriptor, &status)) goto receive_error; } - if (status == STATUS_MANIFEST && receive_manifest(file_descriptor, config, &status) != 0) - return -1; if (status != STATUS_FINISHED) { log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status"); goto receive_error; } + /* Commit-style (late) deletion: every data frame has been received and the + sender proved the whole tree with STATUS_FINISHED. The single-threaded + receiver stores files synchronously, so everything is on disk here and the + deletion can be committed before the --delay-updates publication in + send_success (the walker skips the staging dir, so staged files are never + treated as extras). The -m receiver passes `pending_manifest` because its + disk writer may still be draining; the caller commits after the writer has + joined so no extra file is removed unless the transfer is known to have + succeeded. */ + if (deferred_manifest) { + if (pending_manifest) { + *pending_manifest = deferred_manifest; + } else { + bool deletion_ok = manifest_delete_extras(config, deferred_manifest); + array_list_delete(deferred_manifest); + if (!deletion_ok) { + send_status(file_descriptor, STATUS_ERROR); + return -1; + } + } + } if (sink->send_success) { if (sink->send_success_frame) { if (!sink->send_success_frame(file_descriptor, sink->context)) diff --git a/src/server/receiver.h b/src/server/receiver.h index 4da64b0..7d41426 100644 --- a/src/server/receiver.h +++ b/src/server/receiver.h @@ -35,6 +35,14 @@ 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); +/* receiver_process with an escape hatch for the commit-style (late) deletion: + when `pending_manifest` is non-NULL the receiver does NOT delete at + STATUS_FINISHED itself; instead it stores the owned keep-set manifest there + (leaving *pending_manifest untouched on early modes/errors) so the caller can + commit the deletion only after its disk writer has fully drained. Pass NULL + to keep the default behaviour (delete before the success frame). */ +int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink, + ArrayList** pending_manifest); int receiver_receive_files(Config* config, int file_descriptor); #endif diff --git a/src/server/server.c b/src/server/server.c index bdfe5ef..34eecd7 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -255,6 +255,20 @@ 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) { + /* Commit-style (late) deletion: receive_thread handed the keep-set + manifest here instead of deleting while write_thread might still be + draining, so by now every file is on disk and the whole transfer is + known to have succeeded. Remove the extras before publishing a + --delay-updates run; the walker skips the staging directory. */ + if (context->deferred_manifest) { + if (!manifest_delete_extras(config, context->deferred_manifest)) { + transfer_ok = false; + } + array_list_delete(context->deferred_manifest); + context->deferred_manifest = NULL; + } + } if (transfer_ok) { /* --delay-updates: receive_thread has finished the whole protocol stream (including manifest/delete handling) and write_thread has drained its diff --git a/src/shared/config.c b/src/shared/config.c index ed5fdfe..6a2aea0 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -115,6 +115,8 @@ static void config_set_defaults(Config* config) { config->partial_dir = NULL; config->suffix = NULL; config->delete_before = false; + config->delete_during = false; + config->delete_delay = false; config->address = NULL; config->bind_address = NULL; config->ipv6 = false; @@ -161,12 +163,14 @@ static bool validate_received_config(const Config* config) { valid_wire_bool(config->inplace) && valid_wire_bool(config->append) && 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->delete_delay) && valid_wire_bool(config->delete_during) && valid_wire_bool(config->relative) && valid_wire_bool(config->prune_empty_dirs) && valid_wire_bool(config->delay_updates) && valid_wire_bool(config->mkpath) && !(config->delay_updates && config->inplace) && !(config->delay_updates && delay_updates_staging_name_conflict(config->backup_dir)) && valid_wire_bool(config->partial) && valid_wire_bool(config->delete_before) && valid_wire_bool(config->checksum) && valid_wire_bool(config->eight_bit_output) && + config_has_valid_delete_timing(config) && !(config->skip_compress_set && config->use_chunk_serialization) && (!config->use_compression || (config->compression_level >= 1 && config->compression_level <= 22)) && @@ -188,6 +192,26 @@ Config* config_create(void) { return config; } +bool config_delete_timing_early(const Config* config) { + if (!config) + return false; + return config->delete_before || config->delete_during; +} + +/* A delete-timing flag is only meaningful together with --delete. At most one + of the four flags may be set; several simultaneous timings are a client bug + and are rejected on both ends. */ +bool config_has_valid_delete_timing(const Config* config) { + if (!config) + return false; + if (!config->use_delete) + return !config->delete_before && !config->delete_during && !config->delete_delay && + !config->delete_after; + int timing_count = (config->delete_before ? 1 : 0) + (config->delete_during ? 1 : 0) + + (config->delete_delay ? 1 : 0) + (config->delete_after ? 1 : 0); + return timing_count <= 1; +} + bool config_is_remote_dest(const char* s) { if (s == NULL) return false; @@ -311,7 +335,8 @@ static bool send_selection_options(int fd, const Config* c) { 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->mkpath); + send_int(fd, c->prune_empty_dirs) && send_int(fd, c->mkpath) && + send_int(fd, c->delete_during) && send_int(fd, c->delete_delay); } static bool send_skip_compress_options(int fd, const Config* c) { @@ -418,7 +443,9 @@ static bool receive_selection_options(int fd, Config* c) { return false; if (!receive_wire_bool(fd, &c->mkpath)) return false; - return true; + if (!receive_wire_bool(fd, &c->delete_during)) + return false; + return receive_wire_bool(fd, &c->delete_delay); } static bool receive_resume_options(int fd, Config* c) { diff --git a/src/shared/config.h b/src/shared/config.h index 1b86c9c..406060b 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -147,6 +147,19 @@ typedef struct Config { // PR #179: Delete policies bool delete_before; + /* rsync deletion-timing family (real from Phase 3). At most one of + delete_before / delete_during / delete_delay / delete_after may be set, and + only together with use_delete (the CLI implies --delete for each of them). + delete_before and delete_during select the EARLY engine mode: the keep-set + manifest is transmitted before any file data and extras are removed then, + acknowledged, before the first data byte. delete_delay and delete_after + select the LATE commit mode: extras are removed only after the whole + transfer has succeeded (plain --delete keeps this mode). The exact + semantics and the divergences from rsync are documented in RSYNC_COMPAT.md + and in config_delete_timing_early() below. */ + bool delete_during; + bool delete_delay; + // PR #181: IPv6 and bind address char* address; char* bind_address; @@ -174,7 +187,7 @@ typedef struct Config { DelayUpdatesContext* delay_context; } Config; -#define PROTOCOL_VERSION "2.7.0" +#define PROTOCOL_VERSION "2.8.0" #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) Config* config_create(void); @@ -184,4 +197,16 @@ Config* config_receive(int file_descriptor); bool config_is_remote_dest(const char* s); void config_parse_ssh_dest(Config* config); +/* True when the negotiated delete timing performs the extra-file deletion + * BEFORE the transfer data (--delete-before / --delete-during). The flag is + * a pure function of the config and is used identically on the sender (to pick + * the manifest-first frame order) and the receiver (to delete when the early + * manifest arrives). When false the deletion is committed only after the whole + * transfer succeeded (--delete / --delete-after / --delete-delay). */ +bool config_delete_timing_early(const Config* config); +/* Delete-timing sanity: with deletion enabled at most one timing flag may be + * set (none = the default delete-after commit timing); without deletion no + * timing flag may be set (each timing flag implies --delete). */ +bool config_has_valid_delete_timing(const Config* config); + #endif diff --git a/src/shared/file_receive.c b/src/shared/file_receive.c index 2c4a9c0..5ccea6b 100644 --- a/src/shared/file_receive.c +++ b/src/shared/file_receive.c @@ -792,26 +792,27 @@ File* file_receive_directory(int file_descriptor) { return file; } -int receive_manifest(int fd, const Config* config, int* next_status) { - if (!config) { - send_status(fd, STATUS_ERROR); - return -1; - } - int received_status = STATUS_ERROR; - int* status_out = next_status ? next_status : &received_status; +/* Read a delete-manifest frame (the STATUS_MANIFEST leading code has already + been consumed): an entry count followed by that many destination-relative + paths. The frame is self-delimiting (the count is authoritative), so the + caller decides what to do next and continues reading the following STATUS_* + frame. Returns an owned ArrayList of validated path strings, or NULL after + sending STATUS_ERROR when the frame is malformed (bad count, empty/absolute + path, path traversal, or an aggregate size beyond MAX_MANIFEST_BYTES). */ +ArrayList* receive_manifest_entries(int fd) { int count; if (!receive_int(fd, &count)) { send_status(fd, STATUS_ERROR); - return -1; + return NULL; } if (count < 0 || count > MAX_MANIFEST_ENTRIES) { send_status(fd, STATUS_ERROR); - return -1; + return NULL; } ArrayList* manifest = array_list_create(free); if (!manifest) { send_status(fd, STATUS_ERROR); - return -1; + return NULL; } size_t manifest_bytes = 0; for (int i = 0; i < count; i++) { @@ -823,32 +824,22 @@ int receive_manifest(int fd, const Config* config, int* next_status) { free(s); array_list_delete(manifest); send_status(fd, STATUS_ERROR); - return -1; + return NULL; } } - if (!receive_status(fd, status_out)) { - array_list_delete(manifest); - send_status(fd, STATUS_ERROR); - return -1; - } - /* Deletion is a commit operation: never perform it until the sender has - completed the manifest frame successfully. */ - if (*status_out != STATUS_FINISHED || !config->use_delete) { - array_list_delete(manifest); - if (*status_out != STATUS_FINISHED) - send_status(fd, STATUS_ERROR); - return *status_out == STATUS_FINISHED ? 0 : -1; - } + return manifest; +} + +/* Remove every destination entry under the receive root that is not listed in + `manifest`, bounded by MAX_SERVER_DELETE_COUNT, using the symlink-safe + delete walker. With --delay-updates the not-yet-published staging directory + is a direct child of the receive root and must not be treated as a set of + extras. Prints a notice and returns true on success. */ +bool manifest_delete_extras(const Config* config, ArrayList* manifest) { + if (!config || !manifest) + return false; fprintf(stderr, "Deleting files not in manifest...\n"); - /* With --delay-updates the staged (not yet published) files live directly - under the receive root in the staging directory; the delete walker must - not treat them as extras or it would remove every staged file before it - can be published. */ const char* skip_staging = config->delay_updates ? DELAY_UPDATES_STAGING_DIR : NULL; - bool deletion_ok = delete_extras_limited(config->receive_root_directory, manifest, - MAX_SERVER_DELETE_COUNT, skip_staging); - array_list_delete(manifest); - if (!deletion_ok) - send_status(fd, STATUS_ERROR); - return deletion_ok ? 0 : -1; + return delete_extras_limited(config->receive_root_directory, manifest, MAX_SERVER_DELETE_COUNT, + skip_staging); } diff --git a/src/shared/file_receive.h b/src/shared/file_receive.h index 2b7e689..f7e9f7f 100644 --- a/src/shared/file_receive.h +++ b/src/shared/file_receive.h @@ -10,7 +10,14 @@ File* file_receive(const Config* config, int file_descriptor); File* file_receive_directory(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); +/* Read a delete-manifest frame: entry count then paths (self-delimiting; the + leading STATUS_MANIFEST code has been consumed). Returns an owned path + ArrayList, or NULL after signalling STATUS_ERROR on a malformed frame. */ +ArrayList* receive_manifest_entries(int fd); +/* Remove destination entries under config->receive_root_directory that are not + in `manifest` (bounded walk, staging-dir skip). The caller decides WHEN to + run it based on the negotiated delete timing. */ +bool manifest_delete_extras(const Config* config, ArrayList* manifest); /* Outcome of a single file_save_to_disk operation. The receiver needs to distinguish "written" from "skipped" so --remove-source-files can be told diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index b573feb..0ec5cd5 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -27,6 +27,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que context->loader_done = false; context->manifest = NULL; context->remove_source_files = NULL; + context->early_delete = false; context->total_files = 0; context->progress_bytes = 0; context->total_bytes = 0; @@ -113,6 +114,7 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* context->receiver_done = false; context->queued_bytes = 0; context->max_queue_bytes = 0; + context->deferred_manifest = NULL; atomic_init(&context->cancelled, false); int init = 0; if (mtx_init(&context->mutex, mtx_plain) != thrd_success) @@ -141,6 +143,8 @@ fail: void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { config_delete(context->config); + if (context->deferred_manifest) + array_list_delete(context->deferred_manifest); queue_destroy(context->queue); receiver_outcomes_destroy(&context->outcomes); mtx_destroy(&context->mutex); @@ -236,7 +240,8 @@ int receive_thread(void* pipeline_context) { mtx_unlock(&context->mutex); ReceiverSink sink = {receiver_enqueue_file, context, false, false, NULL}; - if (receiver_process((Config*)config, file_descriptor, &sink) != 0) { + if (receiver_process_pending((Config*)config, file_descriptor, &sink, + &context->deferred_manifest) != 0) { receiver_thread_fail(context); protocol_session_unbind(); return thrd_error; diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index 28f2fe4..d41cd5a 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -26,6 +26,11 @@ typedef struct { bool loader_done; ArrayList* manifest; ArrayList* remove_source_files; + /* True when --delete-before/--delete-during require the keep-set manifest to + be transmitted before any file data: context->manifest is then prebuilt by + a path-only pre-scan on the calling thread and the pipeline scanner must + not append to it. Set once before the worker threads start. */ + bool early_delete; mtx_t mutex_progress; int total_files; unsigned long long progress_bytes; @@ -55,6 +60,14 @@ typedef struct PipelineContextReceiver { budget instead of growing without bound. */ size_t queued_bytes; size_t max_queue_bytes; + /* Keep-set manifest for the commit-style (late) deletion + (--delete/--delete-after/--delete-delay). receive_thread parses the whole + protocol stream but hands the manifest here instead of deleting while the + disk writer may still be draining; the caller (server.c) commits the + deletion after both threads have joined, so no extra is removed unless the + transfer truly succeeded. NULL in the early delete modes (which delete at + the manifest). */ + ArrayList* deferred_manifest; } PipelineContextReceiver; PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, diff --git a/tests/test_server.c b/tests/test_server.c index bbcd76d..337a7f9 100644 --- a/tests/test_server.c +++ b/tests/test_server.c @@ -180,7 +180,10 @@ static void test_receive_manifest_rejects_traversal() { io_set_fds(p[0], p[1]); EXPECT_TRUE(send_int(p[1], 1)); EXPECT_TRUE(send_str(p[1], "../outside")); - EXPECT_EQ_INT(receive_manifest(p[0], cfg, NULL), -1); + EXPECT_NULL(receive_manifest_entries(p[0])); + Status status; + EXPECT_TRUE(receive_status(p[1], &status)); + EXPECT_EQ_INT(status, STATUS_ERROR); close(p[0]); close(p[1]); config_delete(cfg); From 7eacf7c180e96c8459707da0a3e2b6454bfc61cc Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 6 Sep 2026 18:53:24 +0200 Subject: [PATCH 2/7] test: cover rsync delete-timing flags, config wire, and semantics - CLI: each timing flag (+ --del alias) accepted and implies --delete; conflicting timings and a timing with --no-delete are rejected. - Config: delete_during/delete_delay survive config_send/config_receive; two simultaneous timings are rejected by the receiver-side validation; config_delete_timing_early() mapping is unit-tested. - Integration (TCP, single- and multithreaded): every flag removes extras on a successful transfer; early modes (--delete-before/--delete-during/--del) delete before data is applied so a destination file blocking a nested write is removed and the transfer succeeds, while plain --delete / --delete-after / --delete-delay keep it and fail with every extra intact (commit-style). Early timing also completes (without deleting) when the server refuses deletion. --- tests/integration/test_features.py | 126 +++++++++++++++++++++++++++++ tests/test_client_cli.c | 85 +++++++++++++++++-- tests/test_config.c | 126 +++++++++++++++++++++++++++++ 3 files changed, 330 insertions(+), 7 deletions(-) diff --git a/tests/integration/test_features.py b/tests/integration/test_features.py index f8dd698..802ff5d 100644 --- a/tests/integration/test_features.py +++ b/tests/integration/test_features.py @@ -1971,3 +1971,129 @@ class TestFilters: """--filter/-C/-F rule layer: excludes prune, ordering is first-match-wins, the default with no matching rule is include, and legacy --exclude remains an independent layer.""" + + +class TestDeleteTiming: + """rsync deletion-timing family. --delete-before/--delete-during transmit + the keep-set manifest BEFORE any file data (the receiver deletes extras and + acks first); --delete/--delete-after/--delete-delay commit deletions only + after the whole transfer succeeded. Every timing flag implies --delete.""" + + def _seed(self, tag): + source = os.path.join(TEST_DATA_DIR, f"deltiming_{tag}_src") + clean_dir(source) + entries = { + "top.txt": b"top level\n", + "sub/deep.txt": b"deeply nested file\n", + } + for rel, content in entries.items(): + full = os.path.join(source, rel) + os.makedirs(os.path.dirname(full), exist_ok=True) + with open(full, "wb") as fh: + fh.write(content) + return source + + @pytest.mark.parametrize("flag", ["--delete-before", "--delete-during", "--del", + "--delete-after", "--delete-delay"]) + @pytest.mark.parametrize("mt", [False, True]) + def test_flag_removes_extras_on_success(self, flag, mt): + """Every timing flag is accepted, implies --delete, and on a successful + transfer removes the destination extras exactly like plain --delete.""" + source = self._seed("ok") + dest = os.path.join(TEST_DATA_DIR, "deltiming_ok_dst") + clean_dir(dest) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, port=server.port) + assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}" + received = get_dest_received_dir(dest, source) + extra = os.path.join(received, "extra.txt") + with open(extra, "wb") as fh: + fh.write(b"should be deleted") + + flags = [flag] + (["-m"] if mt else []) + result, _ = run_client(source, dest, flags=flags, port=server.port) + assert result.returncode == 0, \ + f"{flag} sync failed: {(result.stderr or result.stdout)[:300]}" + assert not os.path.exists(extra), f"{flag} did not remove the extra file" + mismatches, missing = verify_transfer(source, received) + assert not missing, f"{flag} missing files: {missing}" + assert not mismatches, f"{flag} mismatched files: {mismatches}" + + @pytest.mark.parametrize("flag", ["--delete-before", "--delete-during", "--del"]) + @pytest.mark.parametrize("mt", [False, True]) + def test_early_flags_delete_before_data(self, flag, mt): + """--delete-before/--delete-during remove extras (and a file blocking a + destination directory) BEFORE data is applied, so a nested write that + would fail while the blocker still exists succeeds.""" + source = self._seed("early") + dest = os.path.join(TEST_DATA_DIR, "deltiming_early_dst") + clean_dir(dest) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, port=server.port) + assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}" + received = get_dest_received_dir(dest, source) + extra = os.path.join(received, "extra.txt") + with open(extra, "wb") as fh: + fh.write(b"extra file") + blocker = os.path.join(received, "sub") + shutil.rmtree(blocker) + with open(blocker, "wb") as fh: + fh.write(b"blocks the nested destination directory") + + flags = [flag] + (["-m"] if mt else []) + result, _ = run_client(source, dest, flags=flags, port=server.port) + assert result.returncode == 0, \ + f"{flag} (early delete) did not remove the blocker in time: " \ + f"{(result.stderr or result.stdout)[:300]}" + assert not os.path.exists(extra), f"{flag} did not delete the extra before data" + assert _read_file(os.path.join(received, "sub", "deep.txt")) == b"deeply nested file\n", \ + f"{flag}: nested file was not written after the early deletion" + + @pytest.mark.parametrize("flag", ["--delete", "--delete-after", "--delete-delay"]) + def test_late_flags_commit_only_after_success(self, flag): + """Plain --delete/--delete-after/--delete-delay defer deletion until the + whole transfer succeeds: a mid-transfer write failure must leave every + extra in place (commit-style safety).""" + source = self._seed("late") + dest = os.path.join(TEST_DATA_DIR, "deltiming_late_dst") + clean_dir(dest) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, port=server.port) + assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}" + received = get_dest_received_dir(dest, source) + extra = os.path.join(received, "extra.txt") + with open(extra, "wb") as fh: + fh.write(b"extra file") + blocker = os.path.join(received, "sub") + shutil.rmtree(blocker) + with open(blocker, "wb") as fh: + fh.write(b"blocks the nested destination directory") + + result, _ = run_client(source, dest, flags=[flag], port=server.port) + assert result.returncode != 0, \ + f"{flag} unexpectedly succeeded (deletion must be deferred)" + assert os.path.exists(extra), \ + f"{flag} removed an extra although the transfer failed" + assert os.path.isfile(blocker), \ + f"{flag} deleted the blocker although the transfer failed" + + def test_early_flag_respected_when_server_refuses_delete(self, shared_server): + """With an --allow-delete-less server the client's early timing still + completes (no deadlock on the pre-delete ack) and simply never deletes, + exactly like the plain server policy.""" + source = self._seed("refused") + dest = os.path.join(TEST_DATA_DIR, "deltiming_refused_dst") + clean_dir(dest) + result, _ = run_client(source, dest, port=shared_server.port) + assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}" + received = get_dest_received_dir(dest, source) + extra = os.path.join(received, "extra.txt") + with open(extra, "wb") as fh: + fh.write(b"extra file") + result, _ = run_client(source, dest, flags=["--delete-before"], port=shared_server.port) + assert result.returncode == 0, \ + f"--delete-before against a refuse-delete server failed: {result.stderr[:300]}" + assert os.path.exists(extra), "unauthorized delete removed an extra file" diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index 425d8aa..250e687 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -595,8 +595,9 @@ static void test_parse_args_relative_no_implied_mkpath() { config_delete(cfg); } -/* --del is recognized as the rsync alias, but its timing mode is not implemented. */ -static void test_parse_args_delete_during_alias_unimplemented() { +/* --del is accepted as the rsync alias for --delete-during: it enables + * deletion with the during (early) timing. */ +static void test_parse_args_delete_during_alias() { static const char* const options[] = {"--del", "--delete-during"}; for (size_t i = 0; i < sizeof(options) / sizeof(options[0]); i++) { @@ -605,12 +606,81 @@ static void test_parse_args_delete_during_alias_unimplemented() { int positional_args[2]; int positional_count = 0; - EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), -1); - EXPECT_FALSE(cfg->use_delete); + EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0); + EXPECT_TRUE(cfg->use_delete); + EXPECT_TRUE(cfg->delete_during); + EXPECT_FALSE(cfg->delete_before); + EXPECT_FALSE(cfg->delete_delay); + EXPECT_FALSE(cfg->delete_after); config_delete(cfg); } } +/* Each rsync deletion-timing flag is accepted and implies --delete. */ +static void test_parse_args_delete_timing_flags() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--delete-before", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0); + EXPECT_TRUE(cfg->use_delete); + EXPECT_TRUE(cfg->delete_before); + config_delete(cfg); + + cfg = config_create(); + char* argv_after[] = {"fastsync", "--delete-after", "/src", "/dst"}; + positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 4, argv_after, positional_args, &positional_count), 0); + EXPECT_TRUE(cfg->use_delete); + EXPECT_TRUE(cfg->delete_after); + EXPECT_FALSE(cfg->delete_before); + config_delete(cfg); + + cfg = config_create(); + char* argv_delay[] = {"fastsync", "--delete-delay", "/src", "/dst"}; + positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 4, argv_delay, positional_args, &positional_count), 0); + EXPECT_TRUE(cfg->use_delete); + EXPECT_TRUE(cfg->delete_delay); + EXPECT_FALSE(cfg->delete_before); + EXPECT_FALSE(cfg->delete_after); + config_delete(cfg); +} + +/* Two different delete-timing flags on one command line are a conflict, not a + * silent last-one-wins choice. */ +static void test_parse_args_delete_timing_conflict_rejected() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--delete-before", "--delete-after", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0); + EXPECT_TRUE(cfg->use_delete); + EXPECT_FALSE(validate_config(cfg)); + config_delete(cfg); + + cfg = config_create(); + char* argv2[] = {"fastsync", "--delete-during", "--delete-delay", "/src", "/dst"}; + positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 5, argv2, positional_args, &positional_count), 0); + EXPECT_FALSE(validate_config(cfg)); + config_delete(cfg); +} + +/* A timing flag whose --delete was then negated away must be rejected: timing + * without deletion is meaningless. */ +static void test_parse_args_delete_timing_without_delete_rejected() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--delete-before", "--no-delete", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0); + EXPECT_FALSE(cfg->use_delete); + EXPECT_TRUE(cfg->delete_before); + EXPECT_FALSE(validate_config(cfg)); + config_delete(cfg); +} + /* Parsed-but-unimplemented options must fail instead of being silently accepted. */ static void test_parse_args_rejects_unimplemented_options() { static const char* const options[] = {"--silent", @@ -626,7 +696,6 @@ static void test_parse_args_rejects_unimplemented_options() { "--append", "--append-verify", "--delete-excluded", - "--delete-after", "--max-delete", "--prune-empty-dirs", "-e", @@ -635,7 +704,6 @@ static void test_parse_args_rejects_unimplemented_options() { "--compare-dest", "--copy-dest", "--link-dest", - "--delete-before", "--address", "--bind-address", "--ipv6", @@ -1542,7 +1610,10 @@ void test_client_cli() { test_parse_args_unknown_option(); test_parse_args_dirs_aliases(); test_parse_args_relative_no_implied_mkpath(); - test_parse_args_delete_during_alias_unimplemented(); + test_parse_args_delete_during_alias(); + test_parse_args_delete_timing_flags(); + test_parse_args_delete_timing_conflict_rejected(); + test_parse_args_delete_timing_without_delete_rejected(); test_parse_args_rejects_unimplemented_options(); test_parse_args_quiet(); test_parse_args_human_readable(); diff --git a/tests/test_config.c b/tests/test_config.c index a89b498..557879b 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -438,6 +438,129 @@ static void test_config_delay_updates_reserved_backup_rejected() { config_delete(c); } +static void test_config_delete_timing_early_helper() { + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + EXPECT_FALSE(config_delete_timing_early(cfg)); + EXPECT_TRUE(config_has_valid_delete_timing(cfg)); + cfg->use_delete = true; + EXPECT_TRUE(config_has_valid_delete_timing(cfg)); + EXPECT_FALSE(config_delete_timing_early(cfg)); + config_delete(cfg); + + cfg = config_create(); + cfg->use_delete = true; + cfg->delete_before = true; + EXPECT_TRUE(config_delete_timing_early(cfg)); + EXPECT_TRUE(config_has_valid_delete_timing(cfg)); + config_delete(cfg); + + cfg = config_create(); + cfg->use_delete = true; + cfg->delete_during = true; + EXPECT_TRUE(config_delete_timing_early(cfg)); + EXPECT_TRUE(config_has_valid_delete_timing(cfg)); + config_delete(cfg); + + cfg = config_create(); + cfg->use_delete = true; + cfg->delete_delay = true; + EXPECT_FALSE(config_delete_timing_early(cfg)); + EXPECT_TRUE(config_has_valid_delete_timing(cfg)); + config_delete(cfg); + + cfg = config_create(); + cfg->use_delete = true; + cfg->delete_after = true; + EXPECT_FALSE(config_delete_timing_early(cfg)); + EXPECT_TRUE(config_has_valid_delete_timing(cfg)); + config_delete(cfg); + + /* Two simultaneous timings are invalid. */ + cfg = config_create(); + cfg->use_delete = true; + cfg->delete_before = true; + cfg->delete_after = true; + EXPECT_TRUE(config_delete_timing_early(cfg)); + EXPECT_FALSE(config_has_valid_delete_timing(cfg)); + config_delete(cfg); + + /* A timing flag without deletion is invalid. */ + cfg = config_create(); + cfg->delete_delay = true; + EXPECT_FALSE(config_has_valid_delete_timing(cfg)); + EXPECT_FALSE(config_delete_timing_early(cfg)); + config_delete(cfg); +} + +/* New delete-timing fields must survive config_send/config_receive unchanged, + and a config carrying two conflicting timings must be rejected. */ +static void test_config_delete_timing_wire_roundtrip() { + if (is_running_under_valgrind()) + return; + + struct { + bool before, during, delay, after; + } cases[] = { + {false, false, false, false}, {true, false, false, false}, {false, true, false, false}, + {false, false, true, false}, {false, false, false, true}, + }; + for (size_t i = 0; i < sizeof(cases) / sizeof(cases[0]); i++) { + int p[2]; + EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0); + 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->use_delete && recv->delete_before == cases[i].before && + recv->delete_during == cases[i].during && recv->delete_delay == cases[i].delay && + recv->delete_after == cases[i].after; + } + config_delete(recv); + close(p[0]); + _exit(ok ? 0 : 1); + } else { + close(p[0]); + io_set_fds(p[1], p[1]); + Config* send_cfg = config_create(); + EXPECT_NOT_NULL(send_cfg); + send_cfg->send_directory = str_dup("/src"); + send_cfg->receive_root_directory = str_dup("/dst"); + send_cfg->use_delete = true; + send_cfg->delete_before = cases[i].before; + send_cfg->delete_during = cases[i].during; + send_cfg->delete_delay = cases[i].delay; + send_cfg->delete_after = cases[i].after; + bool sent = config_send(p[1], send_cfg); + int status; + waitpid(pid, &status, 0); + close(p[1]); + config_delete(send_cfg); + EXPECT_TRUE(sent); + EXPECT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0); + } + } +} + +/* The receiver-side wire validation rejects a keep-set config with two + conflicting delete-timing flags. */ +static void test_config_delete_timing_conflict_rejected() { + if (is_running_under_valgrind()) + return; + Config* c = config_create(); + EXPECT_NOT_NULL(c); + c->send_directory = str_dup("/src"); + c->receive_root_directory = str_dup("/dst"); + c->use_delete = true; + c->delete_before = true; + c->delete_delay = true; + EXPECT_FALSE(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")); @@ -474,6 +597,9 @@ void test_config() { test_config_string_null_vs_empty_roundtrip(); test_config_temp_dir_roundtrip(); test_config_delay_updates_reserved_backup_rejected(); + test_config_delete_timing_wire_roundtrip(); + test_config_delete_timing_conflict_rejected(); } + test_config_delete_timing_early_helper(); test_config_is_remote_dest(); } From 36c2c04910d06fa4872c333332998dce9eb8dfdf Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 6 Sep 2026 18:53:28 +0200 Subject: [PATCH 3/7] docs: mark rsync delete-timing family implemented RSYNC_COMPAT.md: flip --delete-before, --del/--delete-during, --delete-delay and --delete-after to Implemented with precise notes (default-under--delete, safety model, 2.7.0 -> 2.8.0 protocol bump, exact divergences from rsync). README option tables list the new flags and the delete-after default. --- README.md | 12 ++++++++++-- RSYNC_COMPAT.md | 25 ++++++++++++++++++++----- 2 files changed, 30 insertions(+), 7 deletions(-) diff --git a/README.md b/README.md index 9a174b8..11b84ff 100644 --- a/README.md +++ b/README.md @@ -102,7 +102,11 @@ partial, alternate, and planned behavior. | `-q, --quiet` | Suppress non-error output | | `--progress` | Show real-time transfer speed | | `-P` | Enables partial-transfer mode and progress output (partial retention is incomplete) | -| `--delete` | Delete files on receiver not present in source | +| `--delete` | Delete files on receiver not present in source (default timing: delete-after, i.e. only after the whole transfer succeeded) | +| `--delete-before` | Delete extras before the transfer starts (implies `--delete`) | +| `--delete-during`, `--del` | Delete extras once the keep-set is known, before data is applied (implies `--delete`) | +| `--delete-delay` | Delete extras only after a successful transfer (implies `--delete`) | +| `--delete-after` | Explicit delete-after timing (implies `--delete`) | | `--exclude ` | Exclude files matching glob pattern (repeatable) | | `--exclude-from ` | Read exclude patterns from a file (one per line) | | `--include ` | Only transfer files matching glob pattern (repeatable, whitelist) | @@ -386,7 +390,11 @@ option. |---|---| | `-a`, `--archive` | Enable current archive preset. Full rsync archive semantics are planned. | | `-n`, `--dry-run` | Scan and report without writing files. | -| `--delete` | Request removal of destination entries absent from the source. The server must allow deletion. | +| `--delete` | Request removal of destination entries absent from the source. The server must allow deletion. Default timing is delete-after: extras are removed only after the whole transfer succeeded. | +| `--delete-before` | Delete extras before the transfer starts (implies `--delete`). | +| `--delete-during`, `--del` | Delete extras once the keep-set manifest is known, before data is applied (implies `--delete`; early mode, same engine behaviour as `--delete-before`). | +| `--delete-delay` | Delete extras only after a successful transfer (implies `--delete`; commit mode, same behaviour as `--delete-after`). | +| `--delete-after` | Explicit delete-after timing: delete only after the transfer succeeded (implies `--delete`). | | `--exclude ` | Exclude matching paths. Repeatable. | | `--include ` | Include matching paths. Repeatable. | | `--exclude-from ` | Read exclude patterns from a file. | diff --git a/RSYNC_COMPAT.md b/RSYNC_COMPAT.md index 6c301a7..0d2544e 100644 --- a/RSYNC_COMPAT.md +++ b/RSYNC_COMPAT.md @@ -103,17 +103,32 @@ This document maps rsync's full feature set to FastSync's current implementation | Flag | Rsync Description | FastSync Status | Notes | |------|-------------------|-----------------|-------| -| `--delete` | Delete extraneous files from dest | ✅ Implemented | `use_delete` config field | -| `--delete-before` | Delete before transfer | ❌ Not Implemented | Removed because it had no effect | -| `--del`, `--delete-during` | Delete during transfer | ❌ Not Implemented | Both flags are recognized but rejected; delete timing is not implemented | -| `--delete-delay` | Find deletions during, delete after | ❌ Not Implemented | | -| `--delete-after` | Delete after transfer | ❌ Not Implemented | Removed because it had no effect | +| `--delete` | Delete extraneous files from dest | ✅ Implemented | `use_delete` config field. Deletion is always derived from the transmitted keep-set manifest of the paths the sender sent/keeps (never from unchecked input), runs through the symlink-safe walker bounded by `MAX_SERVER_DELETE_COUNT`, and skips the `.fastsync-stage` staging dir under `--delay-updates`. FastSync's default timing when no timing flag is given is **delete-after** (extras are removed only once the whole transfer succeeded) — intentionally NOT rsync's `--del`/delete-during default, to preserve FastSync's commit-style safety | +| `--delete-before` | Delete before transfer | ✅ Implemented | Implies `--delete`. The sender runs a full source pre-scan (paths only) and transmits the keep-set manifest BEFORE any file data; the receiver validates it, removes every destination entry not listed (bounded walk, staging-dir skip), then acks `STATUS_OK`. The sender only starts streaming after the deletion committed, or aborts if the receiver reported a deletion error. By definition the deletions already happened when a later transfer phase fails — rsync's delete-before is destructive the same way; a subsequent failure does not restore the removed files. Divergence: the keep-set is the pre-scan snapshot, so a file that appears on the source between the pre-scan and the data pass is still transferred but was not protected from deletion | +| `--del`, `--delete-during` | Delete during transfer | ✅ Implemented | Both spellings accepted; imply `--delete`. FastSync streams the source in a single directory scan and has no per-directory generator pass, so deletions cannot be interleaved per-directory the way rsync's delete-during does. `--delete-during` therefore selects the same early engine mode as `--delete-before` (manifest transmitted before any data, extras removed and acknowledged before data is applied); observable success/failure behaviour equals `--delete-before`. That is the documented divergence from rsync, where `--del` is the default meaning of `--delete` | +| `--delete-delay` | Find deletions during, delete after | ✅ Implemented | Implies `--delete`. Commit-mode timing: extras are removed only after the whole transfer succeeded. rsync's delete-delay records the deletion list during its scan and applies it at the end; FastSync never snapshots the destination while data flows (the keep-set is the transmitted manifest and the destination is listed only at deletion time), so `--delete-delay` is implemented as the same end-of-transfer commit as `--delete-after` with identical safety. That is the documented divergence | +| `--delete-after` | Delete after transfer | ✅ Implemented | Implies `--delete`. The delete-after timing is also what plain `--delete` does: the keep-set manifest closes the data stream and the receiver commits the bounded deletion only after the terminal `STATUS_FINISHED` proves the whole transfer (every data frame received and stored) succeeded. A failed or aborted transfer removes nothing | | `--delete-excluded` | Also delete excluded files | ❌ Not Implemented | Removed because it had no effect | | `--max-delete=NUM` | Max files to delete | ❌ Not Implemented | Removed because it had no effect | | `--ignore-errors` | Delete even with I/O errors | ❌ Not Implemented | | | `--force` | Force deletion of non-empty dirs | ❌ Not Implemented | | | `--prune-empty-dirs` | Prune empty dir chains | ❌ Not Implemented | Removed because it had no effect | +**Deletion-timing implementation notes (Phase 3):** the delete flags above are +real. Two new config booleans (`delete_during`, `delete_delay`) join the already +serialized `delete_before`/`delete_after`, so the on-the-wire config layout +changed and `PROTOCOL_VERSION` was bumped **2.7.0 → 2.8.0** (peers must match). +The `STATUS_MANIFEST` frame is count-delimited and position-independent: the +receiver commits the deletion either when the manifest arrives (early modes: +`--delete-before`/`--delete-during`, which additionally acknowledge with +`STATUS_OK` before data flows) or after the terminal `STATUS_FINISHED` proves +the whole transfer succeeded (commit modes: plain `--delete`/`--delete-after`/ +`--delete-delay`). Timing is chosen purely from the config, so server policy +(`--allow-delete` off) still disables deletion without deadlocking the early +manifest ack. `--delete-delay` and `--delete-during` are each implemented as +the closest safe approximation their engine mode allows; the divergences are +noted in the rows above. + ## 8. Metadata Preservation | Flag | Rsync Description | FastSync Status | Notes | From bbf982dc0b892cb4554b96249291f20140135b9b Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 6 Sep 2026 19:35:17 +0200 Subject: [PATCH 4/7] fix: free the parked delete manifest on every receiver error exit The late/commit path keeps the received keep-set in a local list until STATUS_FINISHED. Error exits after it was parked (STATUS_ABORT, a failing receive_status / non-FINISHED status, a second manifest frame, or a later file/chunk/store failure) previously dropped the only reference and leaked up to ~16 MB of path strings + pointer array per connection. Both failure labels now discard the parked list exactly once; the successful FINISHED path still hands ownership to *pending_manifest (the -m caller) without freeing it. --- src/server/receiver.c | 56 ++++++++++++++++++++++++++++--------------- 1 file changed, 37 insertions(+), 19 deletions(-) diff --git a/src/server/receiver.c b/src/server/receiver.c index be7a04e..e83c50f 100644 --- a/src/server/receiver.c +++ b/src/server/receiver.c @@ -148,18 +148,21 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver if (!receive_status(file_descriptor, &status)) return -1; bool early_delete = config_delete_timing_early(config); + /* Parked keep-set for the late/commit timing. Every exit path below frees it + exactly once; the only exception is the successful FINISHED handoff, which + transfers ownership to *pending_manifest (used by the -m receiver). */ ArrayList* deferred_manifest = NULL; while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK || status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH || status == STATUS_MKDIR || status == STATUS_MANIFEST) { if (status == STATUS_KEEPALIVE) { if (!send_status(file_descriptor, STATUS_KEEPALIVE)) - return -1; + goto fail; goto next_status; } if (status == STATUS_ABORT) { log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up"); - return -1; + goto fail; } if (status == STATUS_CHECK) { bool skipped; @@ -172,7 +175,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver goto receive_error; } else if (status == STATUS_CHECK_BATCH) { if (!receiver_process_batch(config, file_descriptor)) - return -1; + goto fail; goto next_status; } else if (status == STATUS_MKDIR) { File* dir = file_receive_directory(file_descriptor); @@ -181,7 +184,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver } else if (status == STATUS_MANIFEST) { ArrayList* manifest = receive_manifest_entries(file_descriptor); if (!manifest) - return -1; /* receive_manifest_entries already sent STATUS_ERROR */ + goto fail; /* receive_manifest_entries already sent STATUS_ERROR */ if (early_delete) { /* --delete-before / --delete-during: the manifest is authoritative the moment it arrives, before any file data. Delete now and acknowledge @@ -192,24 +195,24 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver array_list_delete(manifest); if (!deletion_ok) { send_status(file_descriptor, STATUS_ERROR); - return -1; + goto fail; } if (!send_status(file_descriptor, STATUS_OK)) - return -1; - } else { + goto fail; + } else if (config->use_delete) { /* Plain --delete / --delete-after / --delete-delay: hold the keep-set and commit the deletion only after STATUS_FINISHED. */ - if (config->use_delete) { - if (deferred_manifest) { - log_message(LOG_LEVEL_ERROR, "Received a second delete manifest"); - array_list_delete(manifest); - send_status(file_descriptor, STATUS_ERROR); - return -1; - } - deferred_manifest = manifest; - } else { + if (deferred_manifest) { + log_message(LOG_LEVEL_ERROR, "Received a second delete manifest"); + array_list_delete(deferred_manifest); + deferred_manifest = NULL; array_list_delete(manifest); + send_status(file_descriptor, STATUS_ERROR); + goto fail; } + deferred_manifest = manifest; + } else { + array_list_delete(manifest); } goto next_status; } else { @@ -241,26 +244,41 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver if (deferred_manifest) { if (pending_manifest) { *pending_manifest = deferred_manifest; + deferred_manifest = NULL; } else { bool deletion_ok = manifest_delete_extras(config, deferred_manifest); array_list_delete(deferred_manifest); + deferred_manifest = NULL; if (!deletion_ok) { send_status(file_descriptor, STATUS_ERROR); - return -1; + goto fail; } } } if (sink->send_success) { if (sink->send_success_frame) { if (!sink->send_success_frame(file_descriptor, sink->context)) - return -1; + goto fail; } else if (!send_status(file_descriptor, STATUS_OK)) { - return -1; + goto fail; } } return 0; +fail: + /* Failure exits that must not (or already did) report a STATUS_ERROR. The + parked keep-set is dropped: never commit a deletion for a failed stream. */ + if (deferred_manifest) { + array_list_delete(deferred_manifest); + deferred_manifest = NULL; + } + return -1; + receive_error: + if (deferred_manifest) { + array_list_delete(deferred_manifest); + deferred_manifest = NULL; + } if (sink->send_error) send_status(file_descriptor, STATUS_ERROR); return -1; From ebfaced5c2a3c3a7a5628b19b2ef34afca2aa729 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 6 Sep 2026 19:35:21 +0200 Subject: [PATCH 5/7] fix: wait for the early-delete ACK with an extended deadline The receiver performs the whole bounded deletion walk (up to MAX_SERVER_DELETE_COUNT unlinks) before answering the delete-before/during manifest, so its STATUS_OK reply can take far longer than the default 60 s per-message receive window. Waiting with the default would make the sender abort AFTER the deletion had already committed on the receiver. Add a timed receive variant (receive_status_timed / protocol_receive_n_data_timed) and use it for the early-manifest ACK with a 1 h explicit deadline; connection errors and EOF still abort immediately. --- src/client/client_send.c | 10 ++++++++-- src/shared/protocol.c | 28 ++++++++++++++++++++++++++-- src/shared/protocol.h | 5 +++++ 3 files changed, 39 insertions(+), 4 deletions(-) diff --git a/src/client/client_send.c b/src/client/client_send.c index df2c22f..0c33c45 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -597,14 +597,20 @@ static int send_delete_manifest(int fd, ArrayList* manifest) { --delete-before/--delete-during, where the extras are removed on the receiver BEFORE the first byte of file data is sent: the receiver acknowledges with STATUS_OK once the bounded delete committed, or STATUS_ERROR if it could not - (in which case the sender aborts without streaming any data). */ + (in which case the sender aborts without streaming any data). The ACK may + take much longer than an ordinary per-message round trip because the receiver + performs the whole bounded deletion walk (up to MAX_SERVER_DELETE_COUNT + unlinks) before replying, so the wait uses a generous explicit deadline + instead of the default 60 s receive window. */ +#define DELETE_ACK_TIMEOUT_SEC 3600 + static bool send_delete_manifest_early(Client* client, ArrayList* manifest) { if (!client || !manifest) return false; if (send_delete_manifest(client->file_descriptor, manifest) != 0) return false; Status ack; - if (!receive_status(client->file_descriptor, &ack)) + if (!receive_status_timed(client->file_descriptor, &ack, DELETE_ACK_TIMEOUT_SEC)) return false; if (ack != STATUS_OK) { log_message(LOG_LEVEL_ERROR, "Server failed to delete files before the transfer"); diff --git a/src/shared/protocol.c b/src/shared/protocol.c index 17caeb3..1d4a421 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -302,15 +302,25 @@ bool protocol_send_n_data(ProtocolSession* session, const void* data, size_t dat return true; } +bool protocol_receive_n_data_timed(ProtocolSession* session, void* data, size_t data_size, + int timeout_sec); + bool protocol_receive_n_data(ProtocolSession* session, void* data, size_t data_size) { + return protocol_receive_n_data_timed(session, data, data_size, RECEIVE_TIMEOUT_SEC); +} + +bool protocol_receive_n_data_timed(ProtocolSession* session, void* data, size_t data_size, + int timeout_sec) { log_debug_message(LOG_DEBUG_IO, " Receiving n Data: %zu", data_size); if (!session) return false; int fd = session->read_fd; + if (timeout_sec <= 0) + timeout_sec = RECEIVE_TIMEOUT_SEC; struct timespec deadline; clock_gettime(CLOCK_MONOTONIC, &deadline); - deadline.tv_sec += RECEIVE_TIMEOUT_SEC; + deadline.tv_sec += timeout_sec; size_t total_bytes_received = 0; short wait_events = POLLIN; @@ -319,7 +329,7 @@ bool protocol_receive_n_data(ProtocolSession* session, void* data, size_t data_s struct pollfd pfd = {.fd = fd, .events = wait_events}; int poll_result = poll(&pfd, 1, deadline_remaining_ms(&deadline)); if (poll_result == 0) { - log_message(LOG_LEVEL_ERROR, "Receive timeout after %ds", RECEIVE_TIMEOUT_SEC); + log_message(LOG_LEVEL_ERROR, "Receive timeout after %ds", timeout_sec); return false; } if (poll_result < 0) { @@ -516,6 +526,17 @@ bool protocol_receive_status(ProtocolSession* session, Status* status) { return true; } +/* protocol_receive_status with an explicit per-message deadline (seconds). + Used where a single reply may legitimately take far longer than the default + 60 s receive window - e.g. the sender waiting for the early-delete ACK after + the receiver committed a large (up to MAX_SERVER_DELETE_COUNT) deletion. */ +bool protocol_receive_status_timed(ProtocolSession* session, Status* status, int timeout_sec) { + if (!protocol_receive_n_data_timed(session, status, sizeof(Status), timeout_sec)) + return false; + log_debug_message(LOG_DEBUG_PROTO, "Received Status: %s", status_to_string(*status)); + return true; +} + bool send_str(int fd, const char* data) { return protocol_send_str(legacy_session(-1, fd), data); } @@ -543,3 +564,6 @@ bool send_status(int fd, Status status) { bool receive_status(int fd, Status* status) { return protocol_receive_status(legacy_session(fd, -1), status); } +bool receive_status_timed(int fd, Status* status, int timeout_sec) { + return protocol_receive_status_timed(legacy_session(fd, -1), status, timeout_sec); +} diff --git a/src/shared/protocol.h b/src/shared/protocol.h index d29d2fa..67a929f 100644 --- a/src/shared/protocol.h +++ b/src/shared/protocol.h @@ -112,5 +112,10 @@ bool send_int(int file_descriptor, int data); bool receive_int(int file_descriptor, int* data); bool send_status(int file_descriptor, Status status); bool receive_status(int file_descriptor, Status* status); +/* receive_status with an explicit per-message deadline in seconds, instead of + the default RECEIVE_TIMEOUT_SEC. A reply that may legitimately take longer + (e.g. the early-delete ACK after a large receiver-side deletion) must use + this so the sender does not abort after the deletion already committed. */ +bool receive_status_timed(int file_descriptor, Status* status, int timeout_sec); #endif From c0c315cf4813a69fcfb6d174877996c3042a2b6b Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 6 Sep 2026 19:35:25 +0200 Subject: [PATCH 6/7] test: cover -m late-deletion failure, receiver leak exits, timed ACK read - Unit (leak guards): drive receiver_process_pending() past a parked keep-set into STATUS_ABORT, EOF, and a second manifest frame; each must return -1 with no manifest handed out. Verified leak-free under ASan. - Unit: receive_status_timed reads a status and fails cleanly on EOF. - Integration: test_late_flags_commit_only_after_success now parametrizes the -m path, proving a failed -m late-timing run preserves every extra and that the deferred manifest is dropped (never applied) when the writer fails. --- tests/integration/test_features.py | 17 ++++-- tests/test_protocol.c | 20 +++++++ tests/test_server.c | 92 ++++++++++++++++++++++++++++++ 3 files changed, 123 insertions(+), 6 deletions(-) diff --git a/tests/integration/test_features.py b/tests/integration/test_features.py index 802ff5d..191d92d 100644 --- a/tests/integration/test_features.py +++ b/tests/integration/test_features.py @@ -2052,10 +2052,14 @@ class TestDeleteTiming: f"{flag}: nested file was not written after the early deletion" @pytest.mark.parametrize("flag", ["--delete", "--delete-after", "--delete-delay"]) - def test_late_flags_commit_only_after_success(self, flag): + @pytest.mark.parametrize("mt", [False, True]) + def test_late_flags_commit_only_after_success(self, flag, mt): """Plain --delete/--delete-after/--delete-delay defer deletion until the whole transfer succeeds: a mid-transfer write failure must leave every - extra in place (commit-style safety).""" + extra in place (commit-style safety). The -m receiver must also keep + the extras: the deferred keep-set is committed by the server only after + the disk-writer thread has finished, and a failing writer means the + manifest is freed, never applied.""" source = self._seed("late") dest = os.path.join(TEST_DATA_DIR, "deltiming_late_dst") clean_dir(dest) @@ -2072,13 +2076,14 @@ class TestDeleteTiming: with open(blocker, "wb") as fh: fh.write(b"blocks the nested destination directory") - result, _ = run_client(source, dest, flags=[flag], port=server.port) + flags = [flag] + (["-m"] if mt else []) + result, _ = run_client(source, dest, flags=flags, port=server.port) assert result.returncode != 0, \ - f"{flag} unexpectedly succeeded (deletion must be deferred)" + f"{flag} (mt={mt}) unexpectedly succeeded (deletion must be deferred)" assert os.path.exists(extra), \ - f"{flag} removed an extra although the transfer failed" + f"{flag} (mt={mt}) removed an extra although the transfer failed" assert os.path.isfile(blocker), \ - f"{flag} deleted the blocker although the transfer failed" + f"{flag} (mt={mt}) deleted the blocker although the transfer failed" def test_early_flag_respected_when_server_refuses_delete(self, shared_server): """With an --allow-delete-less server the client's early timing still diff --git a/tests/test_protocol.c b/tests/test_protocol.c index c6181f5..79fb8a0 100644 --- a/tests/test_protocol.c +++ b/tests/test_protocol.c @@ -412,6 +412,25 @@ static void test_protocol_accounting_release_does_not_underflow() { protocol_session_unbind(); } +static void test_send_receive_status_timed() { + int p[2]; + EXPECT_EQ_INT(pipe(p), 0); + io_set_fds(p[0], p[1]); + io_set_bwlimit(0); + + /* The extended-deadline variant must read an ordinary status just like the + default window, and must fail cleanly on EOF rather than block. */ + EXPECT_TRUE(send_status(0, STATUS_OK)); + Status received = -1; + EXPECT_TRUE(receive_status_timed(0, &received, 5)); + EXPECT_EQ_INT((int)received, (int)STATUS_OK); + + close(p[1]); + EXPECT_FALSE(receive_status_timed(0, &received, 5)); + + close(p[0]); +} + void test_protocol() { test_send_receive_n_data(); test_send_receive_n_data_zero(); @@ -421,6 +440,7 @@ void test_protocol() { test_send_receive_data(); test_send_receive_int(); test_send_receive_status(); + test_send_receive_status_timed(); test_receive_n_data_truncated(); test_receive_str_truncated(); test_max_alloc_rejects_single_buffer(); diff --git a/tests/test_server.c b/tests/test_server.c index 337a7f9..8fca4e5 100644 --- a/tests/test_server.c +++ b/tests/test_server.c @@ -460,6 +460,95 @@ static void test_incremental_check_delta_oversize_reports_failure() { } } +/* Late-timing keep-set leak guard: a manifest parked by the commit path must + be freed on every error exit, never leaked. These tests drive + receiver_process_pending() through an error AFTER the manifest was parked and + are exercised under ASan/valgrind to prove the list is released. */ + +static Config* make_late_delete_config(const char* root) { + Config* cfg = config_create(); + if (!cfg) + return NULL; + cfg->send_directory = str_dup("/src"); + cfg->receive_root_directory = str_dup(root); + cfg->use_delete = true; + cfg->delete_after = true; + return cfg; +} + +static int run_pending_receiver(Config* cfg, int fd, ArrayList** pending) { + ReceiverSink sink = {0}; + return receiver_process_pending(cfg, fd, &sink, pending); +} + +static void test_late_manifest_abort_frees_keepset() { + Config* cfg = make_late_delete_config("/tmp/fastsync_late_abort"); + EXPECT_NOT_NULL(cfg); + 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); + + EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST)); + EXPECT_TRUE(send_int(p[1], 1)); + EXPECT_TRUE(send_str(p[1], "keep.txt")); + EXPECT_TRUE(send_status(p[1], STATUS_ABORT)); + + ArrayList* pending = NULL; + EXPECT_EQ_INT(run_pending_receiver(cfg, p[0], &pending), -1); + EXPECT_NULL(pending); + + close(p[0]); + close(p[1]); + config_delete(cfg); +} + +static void test_late_manifest_eof_frees_keepset() { + Config* cfg = make_late_delete_config("/tmp/fastsync_late_eof"); + EXPECT_NOT_NULL(cfg); + 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); + + EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST)); + EXPECT_TRUE(send_int(p[1], 1)); + EXPECT_TRUE(send_str(p[1], "keep.txt")); + shutdown(p[1], SHUT_WR); + + ArrayList* pending = NULL; + EXPECT_EQ_INT(run_pending_receiver(cfg, p[0], &pending), -1); + EXPECT_NULL(pending); + + close(p[0]); + close(p[1]); + config_delete(cfg); +} + +static void test_late_second_manifest_frees_both() { + Config* cfg = make_late_delete_config("/tmp/fastsync_late_second"); + EXPECT_NOT_NULL(cfg); + 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); + + EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST)); + EXPECT_TRUE(send_int(p[1], 1)); + EXPECT_TRUE(send_str(p[1], "first.txt")); + EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST)); + EXPECT_TRUE(send_int(p[1], 1)); + EXPECT_TRUE(send_str(p[1], "second.txt")); + + ArrayList* pending = NULL; + EXPECT_EQ_INT(run_pending_receiver(cfg, p[0], &pending), -1); + EXPECT_NULL(pending); + + close(p[0]); + close(p[1]); + config_delete(cfg); +} + void test_server() { if (!is_running_under_valgrind()) { test_receive_files_finished(); @@ -470,5 +559,8 @@ void test_server() { test_incremental_check_quick_skip_by_mtime(); test_incremental_check_size_mismatch_full_transfer(); test_incremental_check_delta_oversize_reports_failure(); + test_late_manifest_abort_frees_keepset(); + test_late_manifest_eof_frees_keepset(); + test_late_second_manifest_frees_both(); } } From 02679fe33545ad9a6e09f05c76b429f54c152bd9 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 6 Sep 2026 19:35:28 +0200 Subject: [PATCH 7/7] docs: clarify --del alias, early keep-set caps, ACK wait, --no-delete conflict Usage help now gives --delete-during a complete description with --del on its own line, and notes that timing flags imply --delete while timing+--no-delete is rejected regardless of argument order. RSYNC_COMPAT.md documents: the receiver's MAX_MANIFEST_ENTRIES/MAX_MANIFEST_BYTES caps now abort an early-mode run before any data (previously only the deletion step failed), the extended early-delete ACK deadline, and the order-independent flag-conflict policy. --- RSYNC_COMPAT.md | 24 ++++++++++++++++++++++++ src/client/usage.c | 11 +++++++---- 2 files changed, 31 insertions(+), 4 deletions(-) diff --git a/RSYNC_COMPAT.md b/RSYNC_COMPAT.md index 0d2544e..5526c39 100644 --- a/RSYNC_COMPAT.md +++ b/RSYNC_COMPAT.md @@ -129,6 +129,30 @@ manifest ack. `--delete-delay` and `--delete-during` are each implemented as the closest safe approximation their engine mode allows; the divergences are noted in the rows above. +Manifest size: the sender's keep-set collection (streaming or early pre-scan) +is unbounded, but the receiver rejects any manifest beyond `MAX_MANIFEST_ENTRIES` +(1 048 576 entries) / `MAX_MANIFEST_BYTES` (16 MB of paths) as a hard protocol +error. In the commit modes this only means the deletion is refused after the +data already arrived; in the NEW early modes (`--delete-before`/`--delete-during`) +the manifest is the first frame, so an oversized keep-set now aborts the whole +transfer BEFORE any data is sent (previously all data transferred and only the +deletion step failed). Keep the source tree small enough for the receiver's +manifest caps when using the early timing. + +Early-delete ACK wait: after committing a large deletion (up to +`MAX_SERVER_DELETE_COUNT` unlinks) the receiver's `STATUS_OK`/`STATUS_ERROR` +reply can legitimately take much longer than a normal round trip, so the sender +waits for that single ACK with an extended explicit deadline (1 hour) instead +of the default 60 s per-message receive window. A receiver that is genuinely +gone still aborts the wait via connection close/error; the extended bound only +protects against aborting after the deletion already committed on the receiver. + +Flag-conflict policy: unlike rsync's last-one-wins behaviour, every deletion +timing flag implies `--delete`, and combining a timing flag with `--no-delete` +(in either argument order) — or more than one timing flag — is rejected as a +configuration error rather than silently resolved. Note the check is +order-independent because it runs over the fully parsed config. + ## 8. Metadata Preservation | Flag | Rsync Description | FastSync Status | Notes | diff --git a/src/client/usage.c b/src/client/usage.c index 5634094..6d18531 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -28,12 +28,15 @@ void print_usage(void) { printf(" transfer has succeeded)\n"); printf(" --delete-before Delete extras before the transfer starts\n"); printf(" (implies --delete)\n"); - printf(" --delete-during Delete extras once the keep-set is known, before\n"); - printf(" --del data is applied (alias --del; implies --delete)\n"); + printf(" --delete-during Delete extras once the keep-set manifest is known,\n"); + printf(" before the data is applied (implies --delete)\n"); + printf(" --del Alias for --delete-during\n"); printf(" --delete-delay Delete extras only after a successful transfer\n"); printf(" (implies --delete)\n"); - printf(" --delete-after Alias of the default --delete timing: delete only\n"); - printf(" after the transfer succeeded (implies --delete)\n"); + printf(" --delete-after Delete only after the whole transfer succeeded\n"); + printf(" (the default --delete timing; implies --delete)\n"); + printf(" Note: each timing flag implies --delete. Combining a timing flag with\n"); + printf(" --no-delete (in either order) is rejected as a config error.\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, -d, --old-dirs, --old-d Transfer the named directory entries without\n");