From ee265d78ba71fa9e6f7358204fde27c73ca5dccf Mon Sep 17 00:00:00 2001 From: TapTap Date: Tue, 22 Sep 2026 23:19:55 +0200 Subject: [PATCH] fix(delete-before): replay pre-scan list in --threads path --- src/client/client_send.c | 55 +++++++++++++++- src/shared/multiprocessing.c | 3 + src/shared/multiprocessing.h | 7 ++ .../integration/test_delete_timing_parity.py | 64 ++++++++++++------- 4 files changed, 103 insertions(+), 26 deletions(-) diff --git a/src/client/client_send.c b/src/client/client_send.c index 8c6212e..c8dd8bd 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -1148,6 +1148,43 @@ send_fail: static int scan_directory_multithreaded(void* pipeline_context) { PipelineContextSender* context = (PipelineContextSender*)pipeline_context; protocol_session_bind(&context->allocation_session); + if (context->prescan_chunks != NULL) { + /* --delete-before replays the pre-scan that built the early keep-set as the + data pass (rsync builds one file list). Feed the retained chunks straight + into the pipeline instead of re-reading the source, so a file created + after the pre-scan is neither transferred nor kept. The chunk also + carries the directory times captured by that scan (there is no later + scan), so no scanner is created here. */ + bool failed = false; + for (int i = 0; i < context->prescan_chunks->size; i++) { + Chunk* chunk = (Chunk*)context->prescan_chunks->items[i]; + /* Move ownership out of the retained list so a cleanup here never + double-frees a chunk the queue now owns. */ + context->prescan_chunks->items[i] = NULL; + if (chunk == NULL) + continue; + if (!queue_enqueue_multithreaded_cancel( + context->queue_scanner, chunk, &context->mutex_scanner, + &context->condition_not_empty_scanner, &context->condition_not_full_scanner, + &context->cancelled)) { + chunk_destroy(chunk); + failed = true; + break; + } + } + mtx_lock(&context->mutex_scanner); + context->scanner_done = true; + cnd_broadcast(&context->condition_not_empty_scanner); + cnd_broadcast(&context->condition_not_full_scanner); + mtx_unlock(&context->mutex_scanner); + if (failed) { + pipeline_cancel(context); + protocol_session_unbind(); + return thrd_error; + } + protocol_session_unbind(); + return thrd_success; + } PreparedScanner prepared; /* -j/--threads=N sizes the parallel scanner's worker pool; 0 (bare -j) lets * the scanner apply its built-in default. */ @@ -2032,13 +2069,27 @@ int send_files_multithreaded(Config* config) { if (prepared_ok) prepared.options.plan_dirs = context->plan_dirs; } else { + /* --delete-before: retain the pre-scan chunks as the pipeline's data + pass (rsync's single file list) so a source file created after the + scan is not transferred. No later scan runs, so this pass must also + capture the deferred directory times and the --stats directory + count. */ context->manifest = array_list_create(free); - prepared_ok = prepared_ok && context->manifest != NULL; + context->prescan_chunks = array_list_create(chunk_destroy); + prepared_ok = prepared_ok && context->manifest != NULL && context->prescan_chunks != NULL; + if (prepared_ok) { + prepared.options.dir_entries = context->dir_entries; + prepared.options.dir_entries_mutex = &context->dir_entries_mutex; + prepared.options.dir_count = config->stats ? &context->dir_count : NULL; + if (!append_implied_dir_times(config, context->dir_entries)) + prepared_ok = false; + } } bool prebuilt = prepared_ok && scan_paths_only(config, &prepared.options, context->manifest, context->delete_plans, - &context->scan_had_io_error, &pre_scan_non_dir, NULL, false); + &context->scan_had_io_error, &pre_scan_non_dir, context->prescan_chunks, + context->prescan_chunks != NULL); prepared_scanner_destroy(&prepared); if (per_dir && prebuilt) { const char* walk_root = delete_plan_walk_root(config, context->synced_dirs); diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 60f87ee..3e54e87 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -37,6 +37,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que context->scan_had_io_error = false; context->remove_source_files = NULL; context->early_delete = false; + context->prescan_chunks = NULL; context->delete_plans = NULL; context->delete_suppressed = false; context->scan_stopped_early = false; @@ -194,6 +195,8 @@ void pipeline_context_sender_destroy(PipelineContextSender* context) { if (context->manifest) { array_list_delete(context->manifest); } + if (context->prescan_chunks) + array_list_delete(context->prescan_chunks); if (context->delete_plans) delete_plan_sender_destroy(context->delete_plans); if (context->excluded_paths) diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index e53c542..92ae4cc 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -78,6 +78,13 @@ typedef struct { 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; + /* --delete-before: the path-only pre-scan that built the early keep-set, + retained as the pipeline's file list (owning Chunk*; consumed and NULLed by + the scanner thread) so the data pass replays rsync's single file list + instead of re-reading the source. NULL in every other mode, where the + scanner thread scans normally. Set once before the worker threads start + and freed with the context. */ + ArrayList* prescan_chunks; /* Non-NULL for --delete-during/--delete-delay: the per-directory plan set prebuilt by the path-only pre-scan on the calling thread. The sender thread transmits the root plan before any data and the remaining plans diff --git a/tests/integration/test_delete_timing_parity.py b/tests/integration/test_delete_timing_parity.py index 6757f76..a638c03 100644 --- a/tests/integration/test_delete_timing_parity.py +++ b/tests/integration/test_delete_timing_parity.py @@ -559,56 +559,72 @@ class TestDeleteDelayMaxDeleteRefilledDir: class TestDeleteBeforeLateFileParity: """rsync builds its file list once, so a source file created after that scan - is NOT transferred and its destination extra is deleted. FastSync's - single-threaded --delete-before used to re-scan the source in its data pass - and would transfer the late file (a safe superset); it now replays the - pre-scan file list instead, matching rsync. + is NOT transferred and its destination extra is deleted. FastSync used to + re-scan the source in its data pass (single-threaded) or pipeline a fresh + re-scan against the pre-scan keep-set (``--threads``) and would transfer the + late file (a safe superset); both paths now replay the pre-scan file list + instead, matching rsync. The late file is injected through the config-ack barrier: the first client bytes after the config ack are the pre-scan keep-set manifest, so the hook runs causally after the source scan and before the receiver's delete ack - releases the client into its data pass -- deterministic, no timing guess. + releases the client into its data pass. + + For ``--threads`` the pipeline scanner runs concurrently with the sender, so + the injection must land while that re-scan is still in flight to be observed + by it. The source is therefore a tree of ``_N_DIRS`` directories: the + injection writes the late file into EVERY directory, so it is enough that + any one directory is still unscanned when the hook fires. The tree is sized + so the hook (a localhost round trip) lands long before a full scan finishes; + a re-scanning pipeline then transfers the late files for the directories it + has not yet reached, which the tree comparison catches. """ + _N_DIRS = 2000 + @requires_rsync - def test_late_source_file_not_transferred_and_extra_deleted(self): - source = os.path.join(TEST_DATA_DIR, "dblate_src") - dest = os.path.join(TEST_DATA_DIR, "dblate_dst") - rsync_dst = os.path.join(TEST_DATA_DIR, "dblate_rsync_dst") + @pytest.mark.parametrize("mt", [False, True]) + def test_late_source_file_not_transferred_and_extra_deleted(self, mt): + tag = f"dblate_mt{int(mt)}" + source = os.path.join(TEST_DATA_DIR, f"{tag}_src") + dest = os.path.join(TEST_DATA_DIR, f"{tag}_dst") + rsync_dst = os.path.join(TEST_DATA_DIR, f"{tag}_rsync_dst") clean_dir(source) clean_dir(dest) clean_dir(rsync_dst) - _write(os.path.join(source, "d", "keep.txt"), b"kept payload\n") + for i in range(self._N_DIRS): + _write(os.path.join(source, f"dir{i:05d}", "keep.txt"), b"kept payload\n") # Both destinations carry the would-be late file as an extra. for root in (dest, rsync_dst): - _write(os.path.join(get_dest_received_dir(root, source), "d", "late.txt"), - b"stale extra\n") + received = get_dest_received_dir(root, source) + for i in range(self._N_DIRS): + _write(os.path.join(received, f"dir{i:05d}", "late.txt"), b"stale extra\n") - # rsync reference: the same source with no late file; the extra is removed - # and nothing is transferred for the (never-scanned) late path. + # rsync reference: the same source with no late file; the extras are + # removed and nothing is transferred for the (never-scanned) late paths. rsync_result = _rsync(["-a", "--delete-before", source + "/", rsync_dst + "/"]) assert rsync_result.returncode == 0, rsync_result.stderr rsync_tree = _tree(rsync_dst) - assert "d/late.txt" not in rsync_tree + assert "dir00000/late.txt" not in rsync_tree received = get_dest_received_dir(dest, source) - late_source = os.path.join(source, "d", "late.txt") def hook(): - # Runs after the pre-scan and before the data pass begins. - _write(late_source, b"created after the scan\n") + # Runs after the pre-scan and before the receiver's delete ack. + for i in range(self._N_DIRS): + _write(os.path.join(source, f"dir{i:05d}", "late.txt"), + b"created after the scan\n") with ServerManager() as server: server.start(extra_args=["--allow-delete"]) proxy = _SlicingProxy(server.port, hook=hook, hook_after_config_ack=True) - result, _ = run_client(source, dest, flags=["--delete-before"], port=proxy.port) + flags = ["--delete-before"] + (["--threads=4"] if mt else []) + result, _ = run_client(source, dest, flags=flags, port=proxy.port) proxy.finish() assert result.returncode == 0, (result.stderr or result.stdout)[:400] assert proxy.hook_called.is_set(), "late-file hook never fired" - assert os.path.exists(late_source), "the source late file unexpectedly vanished" - assert not os.path.exists(os.path.join(received, "d", "late.txt")), ( - "late source file was transferred: the single-threaded data pass re-scanned" - ) assert _tree(received) == rsync_tree, ( - f"fastsync tree {_tree(received)} != rsync tree {rsync_tree}" + f"late source {'multithreaded' if mt else 'single-threaded'} data pass re-scanned: " + f"{sum(1 for p in _tree(received) if p.endswith('late.txt'))} late files were " + "transferred" )