Release v2.29.0 #312
@@ -1148,6 +1148,43 @@ send_fail:
|
|||||||
static int scan_directory_multithreaded(void* pipeline_context) {
|
static int scan_directory_multithreaded(void* pipeline_context) {
|
||||||
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
|
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
|
||||||
protocol_session_bind(&context->allocation_session);
|
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;
|
PreparedScanner prepared;
|
||||||
/* -j/--threads=N sizes the parallel scanner's worker pool; 0 (bare -j) lets
|
/* -j/--threads=N sizes the parallel scanner's worker pool; 0 (bare -j) lets
|
||||||
* the scanner apply its built-in default. */
|
* the scanner apply its built-in default. */
|
||||||
@@ -2032,13 +2069,27 @@ int send_files_multithreaded(Config* config) {
|
|||||||
if (prepared_ok)
|
if (prepared_ok)
|
||||||
prepared.options.plan_dirs = context->plan_dirs;
|
prepared.options.plan_dirs = context->plan_dirs;
|
||||||
} else {
|
} 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);
|
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 =
|
bool prebuilt =
|
||||||
prepared_ok &&
|
prepared_ok &&
|
||||||
scan_paths_only(config, &prepared.options, context->manifest, context->delete_plans,
|
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);
|
prepared_scanner_destroy(&prepared);
|
||||||
if (per_dir && prebuilt) {
|
if (per_dir && prebuilt) {
|
||||||
const char* walk_root = delete_plan_walk_root(config, context->synced_dirs);
|
const char* walk_root = delete_plan_walk_root(config, context->synced_dirs);
|
||||||
|
|||||||
@@ -37,6 +37,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
|
|||||||
context->scan_had_io_error = false;
|
context->scan_had_io_error = false;
|
||||||
context->remove_source_files = NULL;
|
context->remove_source_files = NULL;
|
||||||
context->early_delete = false;
|
context->early_delete = false;
|
||||||
|
context->prescan_chunks = NULL;
|
||||||
context->delete_plans = NULL;
|
context->delete_plans = NULL;
|
||||||
context->delete_suppressed = false;
|
context->delete_suppressed = false;
|
||||||
context->scan_stopped_early = false;
|
context->scan_stopped_early = false;
|
||||||
@@ -194,6 +195,8 @@ void pipeline_context_sender_destroy(PipelineContextSender* context) {
|
|||||||
if (context->manifest) {
|
if (context->manifest) {
|
||||||
array_list_delete(context->manifest);
|
array_list_delete(context->manifest);
|
||||||
}
|
}
|
||||||
|
if (context->prescan_chunks)
|
||||||
|
array_list_delete(context->prescan_chunks);
|
||||||
if (context->delete_plans)
|
if (context->delete_plans)
|
||||||
delete_plan_sender_destroy(context->delete_plans);
|
delete_plan_sender_destroy(context->delete_plans);
|
||||||
if (context->excluded_paths)
|
if (context->excluded_paths)
|
||||||
|
|||||||
@@ -78,6 +78,13 @@ typedef struct {
|
|||||||
path-only pre-scan on the calling thread and the pipeline scanner must not
|
path-only pre-scan on the calling thread and the pipeline scanner must not
|
||||||
append to it. Set once before the worker threads start. */
|
append to it. Set once before the worker threads start. */
|
||||||
bool early_delete;
|
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
|
/* 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
|
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
|
thread transmits the root plan before any data and the remaining plans
|
||||||
|
|||||||
@@ -559,56 +559,72 @@ class TestDeleteDelayMaxDeleteRefilledDir:
|
|||||||
|
|
||||||
class TestDeleteBeforeLateFileParity:
|
class TestDeleteBeforeLateFileParity:
|
||||||
"""rsync builds its file list once, so a source file created after that scan
|
"""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
|
is NOT transferred and its destination extra is deleted. FastSync used to
|
||||||
single-threaded --delete-before used to re-scan the source in its data pass
|
re-scan the source in its data pass (single-threaded) or pipeline a fresh
|
||||||
and would transfer the late file (a safe superset); it now replays the
|
re-scan against the pre-scan keep-set (``--threads``) and would transfer the
|
||||||
pre-scan file list instead, matching rsync.
|
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
|
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
|
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
|
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
|
@requires_rsync
|
||||||
def test_late_source_file_not_transferred_and_extra_deleted(self):
|
@pytest.mark.parametrize("mt", [False, True])
|
||||||
source = os.path.join(TEST_DATA_DIR, "dblate_src")
|
def test_late_source_file_not_transferred_and_extra_deleted(self, mt):
|
||||||
dest = os.path.join(TEST_DATA_DIR, "dblate_dst")
|
tag = f"dblate_mt{int(mt)}"
|
||||||
rsync_dst = os.path.join(TEST_DATA_DIR, "dblate_rsync_dst")
|
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(source)
|
||||||
clean_dir(dest)
|
clean_dir(dest)
|
||||||
clean_dir(rsync_dst)
|
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.
|
# Both destinations carry the would-be late file as an extra.
|
||||||
for root in (dest, rsync_dst):
|
for root in (dest, rsync_dst):
|
||||||
_write(os.path.join(get_dest_received_dir(root, source), "d", "late.txt"),
|
received = get_dest_received_dir(root, source)
|
||||||
b"stale extra\n")
|
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
|
# rsync reference: the same source with no late file; the extras are
|
||||||
# and nothing is transferred for the (never-scanned) late path.
|
# removed and nothing is transferred for the (never-scanned) late paths.
|
||||||
rsync_result = _rsync(["-a", "--delete-before", source + "/", rsync_dst + "/"])
|
rsync_result = _rsync(["-a", "--delete-before", source + "/", rsync_dst + "/"])
|
||||||
assert rsync_result.returncode == 0, rsync_result.stderr
|
assert rsync_result.returncode == 0, rsync_result.stderr
|
||||||
rsync_tree = _tree(rsync_dst)
|
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)
|
received = get_dest_received_dir(dest, source)
|
||||||
late_source = os.path.join(source, "d", "late.txt")
|
|
||||||
|
|
||||||
def hook():
|
def hook():
|
||||||
# Runs after the pre-scan and before the data pass begins.
|
# Runs after the pre-scan and before the receiver's delete ack.
|
||||||
_write(late_source, b"created after the scan\n")
|
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:
|
with ServerManager() as server:
|
||||||
server.start(extra_args=["--allow-delete"])
|
server.start(extra_args=["--allow-delete"])
|
||||||
proxy = _SlicingProxy(server.port, hook=hook, hook_after_config_ack=True)
|
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()
|
proxy.finish()
|
||||||
assert result.returncode == 0, (result.stderr or result.stdout)[:400]
|
assert result.returncode == 0, (result.stderr or result.stdout)[:400]
|
||||||
assert proxy.hook_called.is_set(), "late-file hook never fired"
|
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, (
|
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"
|
||||||
)
|
)
|
||||||
|
|||||||
Reference in New Issue
Block a user