From a5d45ef266fca8f66b5cc6907f2e2d71ced16fa3 Mon Sep 17 00:00:00 2001 From: TapTap Date: Fri, 18 Sep 2026 21:56:44 +0200 Subject: [PATCH] fix(parity): init delete_suppressed; gate deleted-path retention - Initialize PipelineContextSender.delete_suppressed=false: an uninitialized true silently skipped the --delete keep-set manifest under -m/--threads, so destination extras were never removed. - Allocate/install the receiver deleted-path observer only when report_deletes is set (--info=del / -i / --out-format under --delete), cap the retained list at MAX_MANIFEST_ENTRIES, and free already-created lists on the receiver-pipeline create failure path. - Validate report_deletes/report_stats/report_dest_info on receive. - Correct stale comments (config.h report_deletes, delete_plan.h deleted count, multiprocessing.h stats locking, utils.h observer placement). - Tests: sender delete_suppressed init, report_deletes gating (unit), and -m/--delete default delete-after keep-set removal (integration). --- src/server/receiver.c | 22 +++++-- src/server/receiver_pipeline.c | 18 +++++- src/shared/config.c | 3 +- src/shared/config.h | 5 +- src/shared/delete_plan.h | 6 +- src/shared/multiprocessing.c | 1 + src/shared/multiprocessing.h | 6 +- src/shared/utils.h | 10 ++-- .../integration/test_delete_timing_parity.py | 21 +++++++ tests/test_multiprocessing.c | 58 +++++++++++++++++++ 10 files changed, 130 insertions(+), 20 deletions(-) diff --git a/src/server/receiver.c b/src/server/receiver.c index 7fd6f46..a6cd5a1 100644 --- a/src/server/receiver.c +++ b/src/server/receiver.c @@ -104,6 +104,11 @@ void receiver_record_deleted_path(void* context, const char* rel_path) { ArrayList* paths = context; if (!paths || !rel_path) return; + /* Bound the retained list like the keep-set manifest: only MAX_MANIFEST_ENTRIES + paths are ever transmitted in the terminal STATUS_STATS frame, so recording + more only grows memory. A hostile/huge deletion set is therefore capped. */ + if ((size_t)paths->size >= (size_t)MAX_MANIFEST_ENTRIES) + return; char* copy = str_dup(rel_path); if (copy && !array_list_add(paths, copy)) free(copy); @@ -417,7 +422,8 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver --max-delete-capped commit still succeeds and the transfer proceeds; the terminal success frame reports the cap. */ size_t deleted = 0; - DeletePathObserver observer = sink->deleted_paths ? receiver_record_deleted_path : NULL; + DeletePathObserver observer = + (config->report_deletes && sink->deleted_paths) ? receiver_record_deleted_path : NULL; DeleteCommitResult deletion = (config->use_delete || config->delete_missing_args) ? manifest_delete_all_observed(config, manifest, &deleted, observer, @@ -459,7 +465,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver } if (!plan_session) { plan_session = delete_plan_session_create(config); - if (plan_session && sink->deleted_paths) + if (plan_session && config->report_deletes && sink->deleted_paths) delete_plan_session_set_delete_observer(plan_session, receiver_record_deleted_path, (void*)sink->deleted_paths); } @@ -505,7 +511,8 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver deferred_manifest = NULL; } else { size_t deleted = 0; - DeletePathObserver observer = sink->deleted_paths ? receiver_record_deleted_path : NULL; + DeletePathObserver observer = + (config->report_deletes && sink->deleted_paths) ? receiver_record_deleted_path : NULL; DeleteCommitResult deletion = manifest_delete_all_observed( config, deferred_manifest, &deleted, observer, (void*)sink->deleted_paths); receiver_tally_deleted(sink, deleted); @@ -525,7 +532,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver hands the session to its caller instead, which commits after the disk writer drained. */ if (plan_session) { - if (sink->deleted_paths) + if (config->report_deletes && sink->deleted_paths) delete_plan_session_set_delete_observer(plan_session, receiver_record_deleted_path, (void*)sink->deleted_paths); if (pending_plans) { @@ -684,8 +691,11 @@ int receiver_receive_files(Config* config, int file_descriptor) { ReceiverSaveContext context = {.config = config, .outcomes = {0}}; dir_time_list_init(&context.dir_times); context.would_delete = array_list_create(free); - context.deleted_paths = array_list_create(free); - if (!context.would_delete || !context.deleted_paths) { + /* report_deletes (--info=del / -i / --out-format under --delete) is the only + reason to retain the actually-removed paths; a plain --delete must not + str_dup every removal. NULL is handled by every consumer. */ + context.deleted_paths = config->report_deletes ? array_list_create(free) : NULL; + if (!context.would_delete || (config->report_deletes && !context.deleted_paths)) { array_list_delete(context.would_delete); array_list_delete(context.deleted_paths); return -1; diff --git a/src/server/receiver_pipeline.c b/src/server/receiver_pipeline.c index d8b7d0d..5fb2c68 100644 --- a/src/server/receiver_pipeline.c +++ b/src/server/receiver_pipeline.c @@ -47,9 +47,15 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* context->would_delete = array_list_create(free); if (!context->would_delete) goto fail; - context->deleted_paths = array_list_create(free); - if (!context->deleted_paths) - goto fail; + /* The actually-removed path list is only needed to render rsync's + `deleting PATH` lines, which the client requests via report_deletes + (--info=del / -i / --out-format under --delete). A plain --delete run must + not allocate it or observe every removal. */ + if (config->report_deletes) { + context->deleted_paths = array_list_create(free); + if (!context->deleted_paths) + goto fail; + } return context; fail: @@ -60,6 +66,12 @@ fail: cnd_destroy(&context->condition_not_full); if (init >= 1) mtx_destroy(&context->mutex); + /* Free every list that was already created before the failing allocation: + `context` itself is freed below, so they would otherwise leak. */ + if (context->would_delete) + array_list_delete(context->would_delete); + if (context->deleted_paths) + array_list_delete(context->deleted_paths); free(context); return NULL; } diff --git a/src/shared/config.c b/src/shared/config.c index 03dcebd..6b065c1 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -206,7 +206,8 @@ static bool validate_received_config(const Config* config) { valid_wire_bool(config->preserve_perms) && valid_wire_bool(config->preserve_times) && valid_wire_bool(config->preserve_owner) && valid_wire_bool(config->preserve_group) && valid_wire_bool(config->munge_links) && valid_wire_bool(config->keep_dirlinks) && - valid_wire_bool(config->fake_super) && + valid_wire_bool(config->fake_super) && valid_wire_bool(config->report_dest_info) && + valid_wire_bool(config->report_stats) && valid_wire_bool(config->report_deletes) && (!config->copy_as_set || (config->copy_as_uid >= 0 && config->copy_as_gid >= 0)) && (!config->use_compression || (config->compression_level >= 1 && config->compression_level <= 22)) && diff --git a/src/shared/config.h b/src/shared/config.h index 9cad6aa..7d67662 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -265,8 +265,9 @@ typedef enum SuperMode { SUPER_MODE_AUTO = 0, SUPER_MODE_ON = 1, SUPER_MODE_OFF * the destination-relative paths it ACTUALLY removed in its terminal * STATUS_STATS record (the same path-list field the dry-run would-delete report * uses), so the sender can print rsync's `deleting PATH`/`*deleting` lines for a - * real (non-dry-run) deletion. It is set only when --info=del is requested with - * --delete; the transfer decision itself is unchanged. */ + * real (non-dry-run) deletion. It is set when --delete is active and any of + * --info=del, -i/--itemize-changes or --out-format requests per-file change + * output; the transfer decision itself is unchanged. */ #define CONFIG_WIRE_OUTPUT_FIELDS(X) \ X(report_dest_info, bool, false, BOOL) \ X(report_stats, bool, false, BOOL) X(report_deletes, bool, false, BOOL) diff --git a/src/shared/delete_plan.h b/src/shared/delete_plan.h index eda33d3..b984a53 100644 --- a/src/shared/delete_plan.h +++ b/src/shared/delete_plan.h @@ -77,8 +77,10 @@ int delete_plan_session_receive(DeletePlanSession* session, const Config* config DeleteCommitResult delete_plan_session_commit(DeletePlanSession* session, const Config* config); /* True once the shared --max-delete budget stopped part of a deletion. */ bool delete_plan_session_limit_reached(const DeletePlanSession* session); -/* Number of destination entries the session's plans removed (or, for - --delete-delay, snapshotted for removal), for the end-of-transfer stats. */ +/* Number of destination entries the session's plans ACTUALLY removed, for the + end-of-transfer stats. A --delete-delay entry only snapshotted for removal + that survives the commit (e.g. a directory refilled before it is removed, so + unlink/rmdir fails with ENOTEMPTY) is not counted. */ size_t delete_plan_session_deleted(const DeletePlanSession* session); /* Install an observer invoked for every destination-relative path the session truly removes (including the deferred --delete-delay commit), so the receiver diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 72faa2f..c45109a 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -38,6 +38,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que context->remove_source_files = NULL; context->early_delete = false; context->delete_plans = NULL; + context->delete_suppressed = false; context->scan_stopped_early = false; context->total_files = 0; context->progress_bytes = 0; diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index 3886d84..a6fda8c 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -93,7 +93,11 @@ typedef struct { unsigned long long progress_bytes; unsigned long long total_bytes; /* Per-type flist / transferred accounting for the rsync --stats breakdown and - the progress `to-chk` denominator. Guarded by mutex_progress. */ + the progress `to-chk` denominator. Owned by the sender thread: it is the + only writer (the entry/transfer notes in send_chunks_multithreaded) and it + reads the totals in its completion tail, so no lock is needed. This is NOT + guarded by mutex_progress (which covers total_files/progress_bytes/ + total_bytes/sender_done). */ TransferStats stats; bool sender_done; atomic_bool cancelled; diff --git a/src/shared/utils.h b/src/shared/utils.h index 642cda3..7484b96 100644 --- a/src/shared/utils.h +++ b/src/shared/utils.h @@ -134,6 +134,11 @@ bool path_under_skip_prefix(const char* child_rel, bool at_root, const DeleteSki cap and returns DELETE_WALK_LIMIT_REACHED when more extras remained. `deleted_out`/`skipped_out` optionally receive the number of entries removed and the number skipped because of the cap. */ +DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* manifest, + const ArrayList* synced_dirs, size_t max_delete, + const DeleteSkipEntry* skips, int skip_count, + size_t* deleted_out, size_t* skipped_out); + /* Optional per-deletion observer: called for each destination-relative path actually removed (a file, symlink, or directory), in removal order, so the receiver can stream rsync's `--info=del`/`--info=remove` lines. */ @@ -147,11 +152,6 @@ DeleteWalkResult delete_extras_limited_observed(const char* dest_root, const Arr size_t* deleted_out, size_t* skipped_out, DeletePathObserver observer, void* observer_context); - -DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* manifest, - const ArrayList* synced_dirs, size_t max_delete, - const DeleteSkipEntry* skips, int skip_count, - size_t* deleted_out, size_t* skipped_out); /* Read-only companion to delete_extras_limited: walk the destination exactly as the delete pass would and APPEND (strdup'd) destination-relative paths that WOULD be removed, without touching disk. Used for -n/--dry-run --delete diff --git a/tests/integration/test_delete_timing_parity.py b/tests/integration/test_delete_timing_parity.py index de831e4..a60d260 100644 --- a/tests/integration/test_delete_timing_parity.py +++ b/tests/integration/test_delete_timing_parity.py @@ -441,3 +441,24 @@ class TestDeleteDelayVsAfterSnapshot: f"{timing} (mt={mt}): new_extra present=" f"{os.path.exists(new_extra)}, expected survives={new_survives}" ) + + +class TestDeleteAfterThreadsKeepSet: + """Regression: -m/--threads with the default delete-after timing (plain + --delete) must still transmit the keep-set manifest and remove destination + extras. PipelineContextSender.delete_suppressed was left uninitialized, so a + garbage true silently skipped the manifest under --threads.""" + + @pytest.mark.parametrize("delete_flag", ["--delete", "--delete-after"]) + def test_threads_delete_after_sends_keep_set(self, delete_flag): + source, dest, received = _seed_pair("mtkeep") + extra = os.path.join(received, "d", "old_extra") + assert os.path.exists(extra) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, flags=["--threads", delete_flag], + port=server.port) + assert result.returncode == 0, (result.stderr or result.stdout)[:300] + assert not os.path.exists(extra), ( + f"{delete_flag} --threads did not remove an extra: keep-set manifest was suppressed" + ) diff --git a/tests/test_multiprocessing.c b/tests/test_multiprocessing.c index e35ab7b..5aa3be3 100644 --- a/tests/test_multiprocessing.c +++ b/tests/test_multiprocessing.c @@ -103,6 +103,62 @@ static void test_sender_zero_capacity() { config_delete(cfg); } +/* Regression (blocker): PipelineContextSender.delete_suppressed must be + initialized false. A garbage true silently suppresses the --delete keep-set + manifest under -m/--threads, so destination extras would never be removed. */ +static void test_sender_delete_suppressed_initialized() { + Config* cfg = config_create(); + EXPECT_NOT_NULL(cfg); + free(cfg->version); + cfg->version = str_dup(PROTOCOL_VERSION); + cfg->send_directory = str_dup("/src"); + cfg->receive_root_directory = str_dup("/dst"); + + Queue* q1 = queue_create(5, NULL); + Queue* q2 = queue_create(5, NULL); + EXPECT_NOT_NULL(q1); + EXPECT_NOT_NULL(q2); + PipelineContextSender* ctx = pipeline_context_sender_create(cfg, q1, q2); + EXPECT_NOT_NULL(ctx); + EXPECT_FALSE(ctx->delete_suppressed); + pipeline_context_sender_destroy(ctx); + config_delete(cfg); +} + +/* --info=del (report_deletes) is the only reason the receiver retains the + actually-removed paths: a plain --delete receiver must not allocate the list, + and with report_deletes armed the observer records into it. */ +static void test_receiver_deleted_paths_gated_by_report_deletes() { + Config* plain = config_create(); + EXPECT_NOT_NULL(plain); + free(plain->version); + plain->version = str_dup(PROTOCOL_VERSION); + plain->receive_root_directory = str_dup("/dst"); + plain->report_deletes = false; + Queue* q_plain = queue_create(5, file_destroy); + EXPECT_NOT_NULL(q_plain); + PipelineContextReceiver* ctx_plain = pipeline_context_receiver_create(plain, q_plain, -1, NULL); + EXPECT_NOT_NULL(ctx_plain); + EXPECT_NULL(ctx_plain->deleted_paths); + pipeline_context_receiver_destroy(ctx_plain); + + Config* info = config_create(); + EXPECT_NOT_NULL(info); + free(info->version); + info->version = str_dup(PROTOCOL_VERSION); + info->receive_root_directory = str_dup("/dst"); + info->report_deletes = true; + Queue* q_info = queue_create(5, file_destroy); + EXPECT_NOT_NULL(q_info); + PipelineContextReceiver* ctx_info = pipeline_context_receiver_create(info, q_info, -1, NULL); + EXPECT_NOT_NULL(ctx_info); + EXPECT_NOT_NULL(ctx_info->deleted_paths); + receiver_record_deleted_path(ctx_info->deleted_paths, "d/old_extra"); + EXPECT_EQ_INT(ctx_info->deleted_paths->size, 1); + EXPECT_EQ_STR((const char*)ctx_info->deleted_paths->items[0], "d/old_extra"); + pipeline_context_receiver_destroy(ctx_info); +} + /* Test receiver with zero file_descriptor */ static void test_receiver_fd_zero() { Config* cfg = config_create(); @@ -462,6 +518,8 @@ void test_multiprocessing() { test_receiver_create_destroy(); test_sender_queue_capacities(); test_sender_zero_capacity(); + test_sender_delete_suppressed_initialized(); + test_receiver_deleted_paths_gated_by_report_deletes(); test_receiver_fd_zero(); if (!is_running_under_valgrind()) { test_receive_thread_finished();