Release v2.29.0 #312

Merged
TapTap merged 123 commits from dev into main 2026-09-23 02:05:14 +02:00
27 changed files with 2130 additions and 440 deletions
Showing only changes of commit 0fbb9de915 - Show all commits
+43
View File
@@ -4,6 +4,49 @@ All notable changes to FastSync are documented here. Versions match
`PROTOCOL_VERSION` (printed by `fastsync --version`); the client and server must
run the same version because the handshake is strict.
## [Unreleased]
The rsync-parity cycle 2.29 (no wire change; `PROTOCOL_VERSION` stays 2.28.0).
`RSYNC_COMPAT.md` moves from **116 ✅ / 14 ⚠️ / 27 ❌** to
**120 ✅ / 10 ⚠️ / 27 ❌** of 157 rows.
### Changed
- **rsync-exact traversal order.** The sequential scanner now walks each
directory's entries in rsync 3.4.1's flist order (non-directories ascending,
then directories ascending, depth-first), so `--info=name`, the
`--delete-during`/`--delete-delay`/`-n` would-delete order and the partial
`--max-delete` survivor set match rsync byte-for-byte. `--threads` has no
rsync analogue and stays unordered.
- **Delete timing.** The complete `--delete-during`/`--delete-delay`
per-directory plan set is transmitted before the first data frame, so a
mid-transfer abort has already removed every planned extra like rsync's
generator; `-d/--dirs` uses per-directory plans (shielded untraversed
subdirectories) instead of the end-of-transfer commit. `-n`, `--delete`,
`--del`/`--delete-during` and `--delete-delay` are now ✅ Parity.
- **Basis directories.** A relative `--compare-dest`/`--copy-dest`/`--link-dest`
DIR resolves against the destination directory with the transfer-relative
name appended, exactly like rsync 3.4.1.
- **`-y`/`--fuzzy`.** The candidate search no longer inherits the ordinary delta
engine's 16 KiB minimum or 10× size-ratio bound, so an oversized or
sub-16-KiB sibling is reused exactly as rsync reuses it.
- `--info=mount` prints rsync's mount-point skip line (repeated `-xx` drops the
mount-point directory); `--info=stats` enables the `--stats` block; `-x` is
repeatable. `--stats` counts traversed directories for the `Number of files`
breakdown under a plain `-r` scan. `--debug` emits real output for
`flist`/`del`/`hash`/`deltasum`/`recv`/`filter`/`send`.
### Known residuals
- `--progress` and `--info` still need a receiver→sender event channel for the
root `./` line, ancestor-directory suppression, receiver-side `skip`/`backup`
wording, and symlink/empty-directory quick-checks.
- `--delete-before`'s phase-0 late-file divergence remains (rsync's pre-scan
fixes the file list before the data pass).
- A single file larger than 256 MiB cannot be streamed in the default path
(a general whole-file limit, not basis-specific).
- `--stats` byte totals and `--msgs2stderr` stay documented divergences.
## [2.28.0] - 2026-09-20
The rsync-parity cycle. `PROTOCOL_VERSION` moves `2.26.0 → 2.27.0 → 2.28.0`;
+16 -15
View File
@@ -1,21 +1,22 @@
# FastSync — Session Handoff (2026-09-19)
# FastSync — Session Handoff (2026-09-20)
## Current status
- **rsync-parity tracks 1-6 landed on `dev`** via **PR #303** (`10159dc`,
"feat(parity): rsync parity tracks 1-6 (protocol 2.28.0)"). Dev push CI run
**581** fully green: lint, build-and-test, parity-full, ASan, UBSan,
fuzz-build, coverage, valgrind.
- **Release `v2.28.0`** is tagged and merged to `main` (PR #304, `b4d54504`).
`dev` is at `558782d` (the incremental-check flake fix).
- **`PROTOCOL_VERSION` = `"2.28.0"`** (`src/shared/config.h`); CMake
`project(FastFileTransfer VERSION 2.28.0)`. The cycle batched all wire
changes (stats counters, filter-rule block, `--verify-basis`) under the one
bump.
- **Release `v2.28.0` tagged and merged to `main`** via PR #304
(`b4d54504`); tag `v2.28.0`. Main push CI run **585** fully green (lint,
build-and-test, parity-full, ASan, UBSan, fuzz-build, coverage, valgrind).
Gitea release `v2.28.0` published. `dev` and `main` are at the release
content.
- Parity matrix: **116 ✅ / 14 ⚠️ / 27 ❌ = 157** (was 111/13/33 at cycle start).
- Working tree clean; feature branch deleted; no scratch trees or worktrees.
`project(FastFileTransfer VERSION 2.28.0)`.
- **Parity cycle 2.29 on branch `feat/parity-2.29`** (from `dev` @ `558782d`),
no wire change. It closes the scanner-order, delete-timing, relative-basis and
fuzzy-eligibility residuals and improves the `--info`/`--stats`/`--debug`
partials. Parity matrix: **120 ✅ / 10 ⚠️ / 27 ❌ = 157** (was 116/14/27).
Remaining ⚠️ rows: `--info`, `--debug`, `--msgs2stderr`, `--stats`,
`--progress`, `--delete-before`, `--compare-dest`/`--copy-dest`/`--link-dest`
(over-256-MiB basis MISS), `-y`/`--fuzzy` (256 MiB buffer cap).
- **Deferred (needs a wire bump):** the `--progress`/`--info` receiver→sender
event channel (root `./` line, ancestor suppression, `skip`/`backup` echo,
symlink/empty-dir quick-check); `--delete-before` phase-0 keep-set; and the
general >256 MiB single-file streaming limit (B4).
- Feature branch `feat/parity-2.29`; integration PR to `dev` pending.
## What landed this session
+19 -12
View File
File diff suppressed because one or more lines are too long
+39 -12
View File
@@ -504,9 +504,8 @@ static bool split_flag_level(const char* token, char* name, size_t name_size, in
* of rsync's `symsafe`, `hlink`, and `own`. */
static bool is_accepted_debug_category(const char* name) {
static const char* const categories[] = {
"acl", "backup", "bind", "chdir", "cmd", "connect", "del", "deltasum",
"dup", "exit", "filter", "flist", "fuzzy", "genr", "hash", "hl",
"hlink", "iconv", "nstr", "own", "owner", "recv", "send", "time",
"acl", "backup", "bind", "chdir", "cmd", "connect", "dup", "exit", "fuzzy",
"genr", "hl", "hlink", "iconv", "nstr", "own", "owner", "time",
};
for (size_t i = 0; i < sizeof(categories) / sizeof(categories[0]); i++) {
if (strcmp(name, categories[i]) == 0)
@@ -518,7 +517,6 @@ static bool is_accepted_debug_category(const char* name) {
static bool is_accepted_info_category(const char* name) {
static const char* const categories[] = {
"backup",
"mount",
"syms",
"symsafe",
};
@@ -571,6 +569,18 @@ static int parse_debug_flags(const char* value, Config* config) {
flag = LOG_DEBUG_PACK;
} else if (strcmp(name, "util") == 0) {
flag = LOG_DEBUG_UTIL;
} else if (strcmp(name, "flist") == 0) {
flag = LOG_DEBUG_FLIST;
} else if (strcmp(name, "del") == 0) {
flag = LOG_DEBUG_DEL;
} else if (strcmp(name, "hash") == 0 || strcmp(name, "deltasum") == 0) {
flag = LOG_DEBUG_HASH;
} else if (strcmp(name, "recv") == 0) {
flag = LOG_DEBUG_RECV;
} else if (strcmp(name, "filter") == 0) {
flag = LOG_DEBUG_FILTER;
} else if (strcmp(name, "send") == 0) {
flag = LOG_DEBUG_SEND;
} else if (is_accepted_debug_category(name)) {
continue;
} else {
@@ -645,9 +655,12 @@ static int parse_info_flags(const char* value, Config* config) {
flag = LOG_INFO_MISC;
else if (strcmp(name, "skip") == 0)
flag = LOG_INFO_SKIP;
else if (strcmp(name, "stats") == 0)
else if (strcmp(name, "stats") == 0) {
flag = LOG_INFO_STATS;
else if (strcmp(name, "del") == 0)
/* `--info=stats` requests the same transfer-statistics block as
`--stats`; `--info=stats0` turns it back off. */
config->stats = level > 0;
} else if (strcmp(name, "del") == 0)
flag = LOG_INFO_DEL;
else if (strcmp(name, "remove") == 0)
flag = LOG_INFO_REMOVE;
@@ -655,6 +668,8 @@ static int parse_info_flags(const char* value, Config* config) {
flag = LOG_INFO_FLIST;
else if (strcmp(name, "nonreg") == 0)
flag = LOG_INFO_NONREG;
else if (strcmp(name, "mount") == 0)
flag = LOG_INFO_MOUNT;
else if (strcmp(name, "progress") == 0)
flag = LOG_INFO_PROGRESS;
else if (is_accepted_info_category(name))
@@ -1213,7 +1228,13 @@ static int apply_table_option(Config* config, const OptionEntry* entry, const ch
void* field = (char*)config + entry->offset;
switch (entry->kind) {
case OPT_FLAG:
*(bool*)field = true;
/* -x/--one-file-system is repeatable in rsync: `-xx` increments the level so
the scanner drops mount-point directories instead of recreating them
empty. Everything else is a plain boolean. */
if (entry->offset == offsetof(Config, one_file_system))
(*(int*)field)++;
else
*(bool*)field = true;
return 0;
case OPT_NOOP:
return 0;
@@ -2574,7 +2595,11 @@ static bool cli_handle_outbuf_option(CliParseCtx* ctx) {
* load --files-from once every argument has been seen. Returns 0 on success,
* -1 on error. */
static int cli_finalize_config(Config* config, bool verbose, bool no_delta, bool no_incremental) {
set_log_level(config->quiet ? LOG_LEVEL_ERROR : (verbose ? LOG_LEVEL_DEBUG : LOG_LEVEL_WARNING));
/* An explicit --debug=FLAGS enables the debug log level by itself (rsync
behaviour); -v enables every other INFO-level message. */
bool debug_enabled = verbose || config->debug_level != 0;
set_log_level(config->quiet ? LOG_LEVEL_ERROR
: (debug_enabled ? LOG_LEVEL_DEBUG : LOG_LEVEL_WARNING));
/* rsync's plain --delete defaults to delete-during (--del): each directory's
extras are removed as that directory is processed, so space is freed
progressively and a tight destination never has to hold the whole old+new
@@ -2764,10 +2789,12 @@ static int cli_finalize_config(Config* config, bool verbose, bool no_delta, bool
}
/* --info=del on a real --delete run asks the receiver to report the paths it
actually removed; the report rides the STATUS_STATS path list, so the wire
stats frame must be negotiated too. */
config->report_deletes = config->use_delete && !config->dry_run &&
((config->info_level & LOG_INFO_DEL) != 0 || config->itemize_changes ||
config->out_format != NULL);
stats frame must be negotiated too. --debug=del needs the same paths, so
it opts into the existing report (no new wire field). */
config->report_deletes =
config->use_delete && !config->dry_run &&
((config->info_level & LOG_INFO_DEL) != 0 || config->itemize_changes ||
config->out_format != NULL || (config->debug_level & LOG_DEBUG_DEL) != 0);
config->report_stats = config->stats || config->show_progress ||
(config->info_level & LOG_INFO_PROGRESS) || format_needs_wire ||
config->report_deletes || (config->dry_run && config->use_delete);
+73 -70
View File
@@ -129,6 +129,21 @@ static void stats_type_breakdown(const TransferStats* stats, char* out, size_t o
out_size);
}
/* rsync's `Number of files` counts every directory. A recursive scan that
preserves a directory attribute captures them in `dir_entries`; a `-r` scan
(no -t/-p) captures nothing, so fall back to the scanner's shared counter of
traversed directories that are not already represented by an inline
directory entry. The -d generator counts its explicit directory entries
inline and does not traverse, so it is excluded here. */
static unsigned long long dir_count_for_stats(const Config* config, const ArrayList* dir_entries,
atomic_ullong* counter) {
if (config == NULL || config->dirs || config->list_only)
return 0;
if (dir_metadata_should_capture(config))
return dir_entries != NULL ? (unsigned long long)dir_entries->size : 0;
return counter != NULL ? (unsigned long long)atomic_load(counter) : 0;
}
/* Print the rsync `--stats` block on stdout. The source-side flist and
transferred counters come from `stats` (filled while scanning/sending), the
receiver-only counters from the STATUS_STATS frame, and the wire byte totals
@@ -387,6 +402,15 @@ static bool info_flag_enabled(const Config* config, LogInfoFlag flag) {
static void print_delete_reports(const Config* config, const ArrayList* paths) {
if (!config || !paths || config->quiet)
return;
/* --debug=del is independent of the --info=del/itemize/out-format display:
emit the debug trace even when no deletion line would be printed. */
if (log_debug_enabled(LOG_DEBUG_DEL)) {
for (int i = 0; i < paths->size; i++) {
const char* raw = (const char*)paths->items[i];
const char* path = delete_display_path(config, raw);
log_debug_message(LOG_DEBUG_DEL, "del: %s", path ? path : raw);
}
}
if (!(config->itemize_changes || config->out_format != NULL ||
info_flag_enabled(config, LOG_INFO_DEL)))
return;
@@ -665,6 +689,7 @@ static bool prepare_scanner(const Config* config, int num_threads, PreparedScann
options->ignore_io_errors = config->ignore_errors;
options->ignore_missing_args = config->ignore_missing_args || config->delete_missing_args;
options->note_nonreg = (config->info_level & LOG_INFO_NONREG) != 0 && !config->quiet;
options->note_mount = (config->info_level & LOG_INFO_MOUNT) != 0 && !config->quiet;
options->send_directory = config->send_directory;
options->eight_bit_output = config->eight_bit_output;
options->excluded_paths = NULL;
@@ -672,6 +697,9 @@ static bool prepare_scanner(const Config* config, int num_threads, PreparedScann
options->size_skipped_paths = NULL;
options->synced_dirs = NULL;
options->hardlinks = NULL;
/* Set by the real send paths; NULL for the metadata-only scans (progress
pre-count, batch) that must not perturb the sender's --stats counter. */
options->dir_count = NULL;
/* P7 Wave D: capture source directory metadata when a directory attribute is
requested (-p for modes, -t for times unless -O omits them). Whether they
are APPLIED is decided receiver-side. */
@@ -730,6 +758,8 @@ static bool progress_precount_scan(const Config* config, ProgressPrecount* out)
ScannerOptions local = prepared.options;
local.list_dirs = true;
local.note_nonreg = false;
local.note_mount = false;
local.dir_count = NULL;
local.use_metadata = false;
local.preserve_xattrs = false;
local.preserve_acls = false;
@@ -1336,6 +1366,7 @@ static bool finalize_transfer(Client* client, const Config* config, ArrayList* r
return false;
if (status == STATUS_STATS) {
ReceiverStats scratch;
log_debug_message(LOG_DEBUG_RECV, "recv: receiver stats");
/* A real --info=del run carries the actually-removed paths in the stats
frame's path list; collect and print them in rsync's format. */
ArrayList* deleted = config->report_deletes ? array_list_create(free) : NULL;
@@ -1855,25 +1886,6 @@ static bool scan_paths_only(const Config* config, const ScannerOptions* options,
return ok;
}
/* Transmit any not-yet-sent per-directory delete plan needed by the entries in
* `chunk` (ancestors root-first, then the entry's own directory for --dirs
* entries) before its data frames go out, so --delete-during/--delete-delay
* clear a directory's extras (and any type conflict) before the directory's
* first write. */
static int send_chunk_delete_plans(Client* client, DeletePlanSender* plans, const Chunk* chunk) {
if (!plans)
return 0;
for (int i = 0; i < chunk->element_count; i++) {
File* f = chunk->items[i];
if (!f)
continue;
if (delete_plan_send_for_path(client->file_descriptor, plans, file_wire_path(f), f->is_dir) !=
0)
return -1;
}
return 0;
}
static int incremental_check(Client* client, File* file, const Config* config,
DeltaSignature** out_sig, unsigned long long* resume_offset) {
*out_sig = NULL;
@@ -1904,6 +1916,8 @@ static int incremental_check(Client* client, File* file, const Config* config,
if (!file_checksum(file, (ChecksumAlgo)config->checksum_algo, config->checksum_seed, digest,
sizeof(digest), &digest_len))
return -1;
log_debug_message(LOG_DEBUG_HASH, "hash: %s (algo %d)", file_wire_path(file),
config->checksum_algo);
uint8_t wire_len = (uint8_t)digest_len;
if (!send_n_data(client->file_descriptor, &wire_len, sizeof(wire_len)) ||
!send_n_data(client->file_descriptor, digest, wire_len))
@@ -1938,6 +1952,7 @@ static int incremental_check(Client* client, File* file, const Config* config,
log_server_rejection("Server reported error for file");
return -1;
}
log_debug_message(LOG_DEBUG_RECV, "recv: check reply for %s", file_wire_path(file));
if (s == STATUS_OK)
return 1;
if (s == STATUS_DELTA_SIGNATURE) {
@@ -1952,6 +1967,8 @@ static int incremental_check(Client* client, File* file, const Config* config,
send_status(client->file_descriptor, STATUS_ERROR);
return -1;
}
log_debug_message(LOG_DEBUG_RECV, "recv: delta signature for %s (%d blocks)",
file_wire_path(file), sig->block_count);
*out_sig = sig;
return 2;
}
@@ -1992,6 +2009,7 @@ static int incremental_check(Client* client, File* file, const Config* config,
static int send_delta(Client* client, File* file, DeltaSignature* sig, Config* config) {
Delta* delta = delta_compute_seeded(file->data->data, file->data->size, sig,
config->delta_block_size, (uint32_t)config->checksum_seed);
log_debug_message(LOG_DEBUG_HASH, "deltasum: %s", file_wire_path(file));
/* The receiver is blocked after sending the signature. Every local
fallback therefore needs the explicit NEXT response before full data. */
if (!delta)
@@ -2465,6 +2483,7 @@ static int send_single_file(Client* client, File* file, Config* config, bool use
bool use_sendfile) {
int compression_level = config->use_compression ? config->compression_level : 0;
log_info_message(LOG_INFO_COPY, "Transferring %s", file->path);
log_debug_message(LOG_DEBUG_SEND, "send: %s", file_wire_path(file));
if (!use_incremental) {
if (use_sendfile) {
@@ -2753,9 +2772,12 @@ static int send_chunks_multithreaded(void* pipeline_context) {
return thrd_error;
}
} else if (context->delete_plans) {
/* --delete-during/--delete-delay: transmit the receive root's plan before
any data, exactly like rsync's first generator directory. */
if (delete_plan_send_root(client->file_descriptor, context->delete_plans) != 0) {
/* --delete-during/--delete-delay: transmit the COMPLETE per-directory plan
set before any data, so a mid-transfer abort has already applied every
planned removal exactly like rsync's generator (which runs ahead of its
throttled sender). A completed run is unaffected. */
if (delete_plan_send_all(client->file_descriptor, context->delete_plans, context->plan_dirs) !=
0) {
pipeline_cancel(context);
disconnect_transfer_client(client);
mark_sender_done(context);
@@ -2802,15 +2824,6 @@ static int send_chunks_multithreaded(void* pipeline_context) {
}
break;
}
if (send_chunk_delete_plans(client, context->delete_plans, current_chunk) != 0) {
log_message(LOG_LEVEL_ERROR, "unexpected error while sending delete plan");
chunk_destroy(current_chunk);
pipeline_cancel(context);
disconnect_transfer_client(client);
mark_sender_done(context);
protocol_session_unbind();
return thrd_error;
}
if (send_chunk_with_removal(client, current_chunk, context->config,
context->remove_source_files, &context->stats) != 0) {
log_message(LOG_LEVEL_ERROR, "unexpected error while sending chunk");
@@ -2893,13 +2906,6 @@ static int send_chunks_multithreaded(void* pipeline_context) {
NULL) != 0)
goto send_fail;
}
/* Emit the plans for source directories the data stream never triggered
(empty directories): their extras are still cleared while the directory
itself is kept. */
if (!context->scan_stopped_early && context->delete_plans && context->plan_dirs &&
delete_plan_send_remaining(client->file_descriptor, context->delete_plans,
context->plan_dirs) != 0)
goto send_fail;
/* P7 Wave D: transmit the captured directory times last. The scanner thread
(and all parallel workers) has been joined before scanner_done was set, so
the list is complete and race-free; on an early stop the list may be
@@ -2918,8 +2924,8 @@ static int send_chunks_multithreaded(void* pipeline_context) {
"server reported a deletion failure (--delete); see the server log for the reason");
if (ok)
remove_transferred_sources(context->config, context->remove_source_files);
if (context->dir_entries)
context->stats.flist_dir += (unsigned long long)context->dir_entries->size;
context->stats.flist_dir +=
dir_count_for_stats(context->config, context->dir_entries, &context->dir_count);
report_transfer_stats(context->config, &context->stats, start, &recv_stats);
log_info_message(LOG_INFO_STATS, "Transfer summary: %llu files, %.1f MB",
context->stats.transferred_regular,
@@ -2956,6 +2962,7 @@ static int scan_directory_multithreaded(void* pipeline_context) {
parallel workers append under the context's dedicated mutex. */
prepared.options.dir_entries = context->dir_entries;
prepared.options.dir_entries_mutex = &context->dir_entries_mutex;
prepared.options.dir_count = context->config->stats ? &context->dir_count : NULL;
if (!append_implied_dir_times(context->config, context->dir_entries)) {
pipeline_cancel(context);
protocol_session_unbind();
@@ -3224,6 +3231,9 @@ int send_files(Config* config) {
/* P7 Wave D: captured source directory times, transmitted in trailing
STATUS_DIR_TIMES frame(s) (only when metadata rides the wire). */
ArrayList* dir_entries = NULL;
/* --stats directory accounting for the no-metadata (-r) case. */
atomic_ullong dir_count;
atomic_init(&dir_count, 0);
/* Protected excluded prefixes (delete-excluded default protection). */
ArrayList* excluded = NULL;
/* Size-pruned prefixes (always protected) and synchronized directories. */
@@ -3232,10 +3242,12 @@ int send_files(Config* config) {
/* Traversed source directories for the per-directory delete keep set. */
ArrayList* plan_dirs = NULL;
bool delete_early = config->use_delete && config_delete_timing_early(config);
/* -d/--dirs does not recurse, so a per-directory plan would carry no child
information and could delete the contents of an untraversed directory;
fall back to the whole-tree end-of-transfer commit for that mode. */
bool delete_per_dir = config->use_delete && config_delete_timing_per_dir(config) && !config->dirs;
/* --delete-during/--delete-delay use per-directory plans for every transfer
shape. For -d/--dirs the generator records only the directories whose
direct children it actually enumerated, so the plan removes extras directly
inside a listed directory while an untraversed (kept) subdirectory is
shielded -- rsync's `-d DIR/ --delete`. */
bool delete_per_dir = config->use_delete && config_delete_timing_per_dir(config);
bool send_failed = false;
bool had_scan_io = false;
unsigned long long per_dir_non_dir_count = 0;
@@ -3287,12 +3299,12 @@ int send_files(Config* config) {
prepared.options.synced_dirs = synced_dirs;
}
}
/* The late-timing modes (--delete-after/--delete-commit and a plain --delete
that fell back from per-dir mode because of -d/--dirs) build the manifest
/* The late-timing modes (--delete-after/--delete-commit) build the manifest
while streaming and send it after the last data frame. --delete-before
sends a whole-tree keep-set up front; --delete-during/--delete-delay build a
per-directory plan set up front (paths only) and stream the plans alongside
the data, so no manifest is kept during the data pass. */
sends a whole-tree keep-set up front; --delete-during/--delete-delay build
the complete per-directory plan set up front (paths only) and transmit it
all before the first data frame, so a mid-transfer abort has already
applied every planned removal. */
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
@@ -3335,9 +3347,10 @@ int send_files(Config* config) {
goto send_fail;
} else if (delete_per_dir) {
/* --delete-during/--delete-delay: build one plan per source directory from a
path-only pre-scan and transmit the root plan now, before any data, so the
receive root's extras are handled exactly like rsync's first generator
directory. The remaining plans are streamed with the data below. */
path-only pre-scan and transmit the COMPLETE plan set now, before any data,
so every planned removal has already been applied when a later transfer
phase fails -- exactly like rsync's generator, whose deletion list runs
ahead of its throttled sender. A completed run is unaffected. */
plan_sender = delete_plan_sender_create();
plan_dirs = array_list_create(free);
if (!plan_sender || !plan_dirs)
@@ -3368,7 +3381,7 @@ int send_files(Config* config) {
plan_dirs = NULL;
skip_delete = true;
} else {
plans_ok = delete_plan_send_root(client->file_descriptor, plan_sender) == 0;
plans_ok = delete_plan_send_all(client->file_descriptor, plan_sender, plan_dirs) == 0;
}
}
prepared.options.excluded_paths = NULL;
@@ -3402,6 +3415,7 @@ int send_files(Config* config) {
the directory-time list (otherwise every directory would be captured
twice). */
prepared.options.dir_entries = dir_entries;
prepared.options.dir_count = config->stats ? &dir_count : NULL;
scanner = directory_scanner_create_with_options(config->send_directory, &prepared.options);
if (!scanner)
goto send_fail;
@@ -3455,11 +3469,6 @@ int send_files(Config* config) {
goto send_fail;
}
}
if (send_chunk_delete_plans(client, plan_sender, current_chunk) != 0) {
chunk_destroy(current_chunk);
send_failed = true;
break;
}
if (send_chunk_with_removal(client, current_chunk, config, remove_sources, &transfer_stats) !=
0) {
log_message(LOG_LEVEL_ERROR, "Failed to send chunk");
@@ -3542,12 +3551,6 @@ int send_files(Config* config) {
}
}
}
/* Emit the plans for any source directories the data stream never triggered
(an empty directory has no file frame). Sending them now still clears that
directory's destination extras while keeping the directory itself. */
if (!scan_stopped_early && plan_sender && plan_dirs &&
delete_plan_send_remaining(client->file_descriptor, plan_sender, plan_dirs) != 0)
goto send_fail;
/* P7 Wave D: every directory has now been traversed (or the scan stopped
early), so transmit the captured directory times last. The receiver defers
applying them until after its own deletion/publication phase. */
@@ -3566,8 +3569,7 @@ int send_files(Config* config) {
from the scanner's captured directory list (present whenever a directory
attribute is preserved, e.g. -a/-t/-p). The -d generator counts its
explicit directory entries inline instead. */
if (dir_entries)
transfer_stats.flist_dir += (unsigned long long)dir_entries->size;
transfer_stats.flist_dir += dir_count_for_stats(config, dir_entries, &dir_count);
report_transfer_stats(config, &transfer_stats, start, &recv_stats);
log_info_message(LOG_INFO_STATS, "Transfer summary: %llu files, %.1f MB",
transfer_stats.transferred_regular,
@@ -3714,10 +3716,11 @@ int send_files_multithreaded(Config** config_ptr) {
return 1;
}
}
/* -d/--dirs does not recurse, so a per-directory plan would carry no child
information and could delete the contents of an untraversed directory;
fall back to the whole-tree end-of-transfer commit for that mode. */
bool per_dir = config_delete_timing_per_dir(config) && !config->dirs;
/* --delete-during/--delete-delay use per-directory plans for every transfer
shape. The -d/--dirs generator records only the directories whose direct
children it enumerated, so extras directly inside a listed directory are
removed while an untraversed (kept) subdirectory is shielded. */
bool per_dir = config_delete_timing_per_dir(config);
if (config_delete_timing_early(config) || per_dir) {
/* --delete-before / --delete-during / --delete-delay: build the keep-set
(paths only, nothing loaded or sent) up front so the sender thread can
+248 -49
View File
@@ -205,11 +205,22 @@ typedef struct {
bool referent_error;
} ScannerEntry;
/* One inspected directory entry buffered so the sequential scanner can emit the
stream in rsync's flist order. `name` is the raw dirent name (owned here);
`entry` is the scanner_inspect_entry() result whose path/link_target are owned
when `inspection == 1`; `inspection` is that call's return code (1 keep,
0 skip, <0 fatal). */
typedef struct {
char* name;
ScannerEntry entry;
int inspection;
} SortedEntry;
/* --one-file-system (-x) decision. Only directories can carry a different
* device than their parent (mount points), so this is checked when a child
* directory is about to be descended into. */
bool scanner_same_filesystem(bool one_file_system, dev_t root_device, dev_t entry_device) {
return !one_file_system || entry_device == root_device;
bool scanner_same_filesystem(int one_file_system, dev_t root_device, dev_t entry_device) {
return one_file_system <= 0 || entry_device == root_device;
}
/* Build a payload-less directory File carrying the captured metadata (when
@@ -445,6 +456,40 @@ static void scanner_note_nonreg(const ScannerOptions* options, const char* fs_pa
fflush(stdout);
}
/* rsync 3.4.1's `--info=mount` line, emitted when `-xx` drops a mount-point
* directory: `[sender] skipping mount-point dir NAME` (the client is the
* sender). Plain `-x` keeps the empty directory and prints nothing, matching
* rsync. */
static void scanner_note_mount(const ScannerOptions* options, const char* fs_path) {
if (!options || !options->note_mount || !fs_path)
return;
const char* rel = utils_strip_transfer_root(fs_path, options->send_directory);
char* escaped = output_escape(rel, options->eight_bit_output);
printf("[sender] skipping mount-point dir %s\n", escaped ? escaped : rel);
free(escaped);
fflush(stdout);
}
/* --debug=filter: a selection/filter decision dropped an entry. */
static void scanner_note_filter(const ScannerOptions* options, const char* name) {
if (!options || !log_debug_enabled(LOG_DEBUG_FILTER) || !name)
return;
log_debug_message(LOG_DEBUG_FILTER, "filter: excluded %s", name);
}
/* Account for a directory that will not be represented by an inline directory
* entry. Paired with scanner_dir_count_uncount for empty directories that are
* emitted inline, so every traversed directory is counted exactly once. */
static void scanner_dir_count_count(const ScannerOptions* options) {
if (options && options->dir_count)
atomic_fetch_add(options->dir_count, 1);
}
static void scanner_dir_count_uncount(const ScannerOptions* options) {
if (options && options->dir_count)
atomic_fetch_sub(options->dir_count, 1);
}
/* A user-selection exclusion (--filter/-C/per-dir or --exclude/--include). */
static void scanner_record_excluded(DirectoryScanner* scanner, const char* fs_path) {
scanner_record_protected(scanner, fs_path, scanner->options.excluded_paths);
@@ -687,6 +732,114 @@ skip:
return 0;
}
static void sorted_entry_destroy(void* item) {
SortedEntry* se = (SortedEntry*)item;
if (!se)
return;
free(se->name);
free(se->entry.path);
free(se->entry.link_target);
}
/* rsync flist order within one directory: non-directories first, then
directories, each group by ascending name. strcmp() compares as unsigned
char, matching rsync's f_name_cmp(). */
static int sorted_entry_cmp(const void* a, const void* b) {
const SortedEntry* x = (const SortedEntry*)a;
const SortedEntry* y = (const SortedEntry*)b;
bool x_dir = x->inspection > 0 && x->entry.is_directory;
bool y_dir = y->inspection > 0 && y->entry.is_directory;
if (x_dir != y_dir)
return x_dir ? 1 : -1;
return strcmp(x->name, y->name);
}
static void scanner_free_sorted(DirectoryScanner* scanner) {
SortedEntry* entries = (SortedEntry*)scanner->sorted_entries;
for (size_t i = 0; i < scanner->sorted_count; i++)
sorted_entry_destroy(&entries[i]);
free(entries);
scanner->sorted_entries = NULL;
scanner->sorted_count = 0;
scanner->sorted_index = 0;
}
/* Read every entry of the open directory, inspect it once and store it sorted in
rsync's flist order. Returns 0 on success, -1 on a fatal error (the caller
aborts the scan). */
static int scanner_buffer_current_directory(DirectoryScanner* scanner) {
size_t capacity = 64;
size_t count = 0;
SortedEntry* entries = malloc(capacity * sizeof(*entries));
if (!entries) {
scanner->failed = true;
return -1;
}
const struct dirent* dirent;
while ((dirent = readdir(scanner->current_dir)) != NULL) {
if (strcmp(dirent->d_name, ".") == 0 || strcmp(dirent->d_name, "..") == 0)
continue;
if (count == capacity) {
size_t next = capacity * 2;
SortedEntry* grown = realloc(entries, next * sizeof(*entries));
if (!grown) {
scanner->failed = true;
break;
}
entries = grown;
capacity = next;
}
char* name = str_dup(dirent->d_name);
if (!name) {
scanner->failed = true;
break;
}
char* link_rel = child_rel_path(scanner->current_rel, dirent->d_name);
if (!link_rel) {
free(name);
scanner->failed = true;
break;
}
int inspection = scanner_inspect_entry(&scanner->options, scanner->current_path, link_rel,
dirent->d_name, &entries[count].entry);
free(link_rel);
if (inspection < 0) {
free(name);
scanner->failed = true;
break;
}
entries[count].name = name;
entries[count].inspection = inspection;
count++;
}
if (scanner->failed) {
for (size_t i = 0; i < count; i++)
sorted_entry_destroy(&entries[i]);
free(entries);
return -1;
}
qsort(entries, count, sizeof(*entries), sorted_entry_cmp);
scanner->sorted_entries = entries;
scanner->sorted_count = count;
scanner->sorted_index = 0;
return 0;
}
/* Push this directory's collected child directories onto the LIFO stack in
reverse so the first (ascending) child is popped first (depth-first). */
static void scanner_push_pending_dirs(DirectoryScanner* scanner) {
ArrayList* pending = (ArrayList*)scanner->pending_dirs;
if (!pending)
return;
for (int i = pending->size - 1; i >= 0; i--) {
if (!queue_push(scanner->directories, pending->items[i])) {
dir_entry_destroy(pending->items[i]);
scanner->failed = true;
}
}
pending->size = 0;
}
DirectoryScanner* directory_scanner_create_with_options(const char* root_directory,
const ScannerOptions* options) {
if (!root_directory || !options)
@@ -704,12 +857,22 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo
free(scanner);
return NULL;
}
scanner->pending_dirs = array_list_create(NULL);
if (!scanner->pending_dirs) {
queue_destroy(scanner->directories);
free(scanner);
return NULL;
}
scanner->current_dir = NULL;
scanner->current_path = NULL;
scanner->current_depth = 0;
scanner->failed = false;
scanner->sorted_entries = NULL;
scanner->sorted_count = 0;
scanner->sorted_index = 0;
scanner->root_path = str_dup(root_directory);
if (!scanner->root_path) {
array_list_delete(scanner->pending_dirs);
queue_destroy(scanner->directories);
free(scanner);
return NULL;
@@ -729,6 +892,7 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo
scanner->filter_nodes = array_list_create(filter_node_destroy);
if (!scanner->filter_nodes) {
free(scanner->root_path);
array_list_delete(scanner->pending_dirs);
queue_destroy(scanner->directories);
free(scanner);
return NULL;
@@ -739,6 +903,7 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo
if (stat(root_directory, &root_stats) != 0) {
log_perror("Could not stat source directory");
free(scanner->root_path);
array_list_delete(scanner->pending_dirs);
queue_destroy(scanner->directories);
array_list_delete(scanner->filter_nodes);
free(scanner);
@@ -757,6 +922,7 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo
if (!queue_enqueue(scanner->directories, root)) {
dir_entry_destroy(root);
free(scanner->root_path);
array_list_delete(scanner->pending_dirs);
queue_destroy(scanner->directories);
array_list_delete(scanner->filter_nodes);
free(scanner);
@@ -808,6 +974,13 @@ void directory_scanner_destroy(DirectoryScanner* scanner) {
free(scanner->current_path);
free(scanner->current_rel);
free(scanner->root_path);
scanner_free_sorted(scanner);
ArrayList* pending = (ArrayList*)scanner->pending_dirs;
if (pending) {
for (int i = 0; i < pending->size; i++)
dir_entry_destroy(pending->items[i]);
array_list_delete(pending);
}
array_list_delete(scanner->filter_nodes);
array_list_delete(scanner->dirs_batch);
queue_destroy(scanner->directories);
@@ -957,6 +1130,9 @@ static bool scanner_emit_empty_dir(DirectoryScanner* scanner, ArrayList* chunk_d
file_destroy(dir);
return false;
}
/* The directory was counted when it was opened; this inline entry represents
it, so drop the counter to avoid counting it twice in --stats. */
scanner_dir_count_uncount(&scanner->options);
return true;
}
@@ -975,7 +1151,7 @@ static int open_next_directory(DirectoryScanner* scanner) {
scanner->current_path = NULL;
while (!queue_is_empty(scanner->directories)) {
DirEntry* de = (DirEntry*)queue_dequeue(scanner->directories);
DirEntry* de = (DirEntry*)queue_pop(scanner->directories);
scanner->current_path = de->path;
scanner->current_depth = de->depth;
/* The seed directory inherits the scanner's configured context (the root
@@ -1046,6 +1222,8 @@ static int open_next_directory(DirectoryScanner* scanner) {
scanner->failed = true;
return -1;
}
scanner_dir_count_count(&scanner->options);
log_debug_message(LOG_DEBUG_FLIST, "flist: scanning %s", scanner->current_path);
if (scanner->options.capture_dir_times &&
!scanner_capture_dir_time(
scanner->options.dir_entries, scanner->options.dir_entries_mutex, scanner->root_path,
@@ -1060,6 +1238,14 @@ static int open_next_directory(DirectoryScanner* scanner) {
scanner->failed = true;
return -1;
}
/* Buffer and sort this directory's entries in rsync's flist order. */
if (scanner_buffer_current_directory(scanner) != 0) {
closedir(scanner->current_dir);
scanner->current_dir = NULL;
free(scanner->current_path);
scanner->current_path = NULL;
return -1;
}
return 1;
}
return 0;
@@ -1304,6 +1490,17 @@ static File* dirs_next_file(DirectoryScanner* scanner) {
scanner->dirs_root_emitted = true;
if (scanner->options.prune_empty_dirs && dirs_source_dir_is_empty(scanner->root_path))
return NULL;
/* The listed directory's direct children are about to be enumerated, so
its destination mirror is a synchronized directory: record it for the
per-directory delete plan. The plan keeps the enumerated children and
shields untraversed subdirectories, so --delete-during removes extras
directly inside the listed directory without descending into a kept
(but untraversed) child -- exactly rsync's `-d DIR/ --delete`. */
if (!scanner_record_synced_dir(&scanner->options, scanner->root_path, "",
scanner->relative_mode)) {
scanner->failed = true;
return NULL;
}
scanner->current_dir = opendir(scanner->root_path);
if (!scanner->current_dir) {
scanner->io_error = true;
@@ -1415,8 +1612,7 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
break;
}
const struct dirent* entry = readdir(scanner->current_dir);
if (entry == NULL) {
if (scanner->sorted_index >= scanner->sorted_count) {
/* The directory is exhausted: if nothing was transferred or descended
from it, recreate it at the destination as an explicit entry. */
if (scanner->options.emit_empty_dirs && !scanner->current_dir_produced &&
@@ -1425,10 +1621,12 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
if (!scanner_emit_empty_dir(scanner, chunk_data))
scanner->failed = true;
}
scanner_push_pending_dirs(scanner);
closedir(scanner->current_dir);
scanner->current_dir = NULL;
free(scanner->current_path);
scanner->current_path = NULL;
scanner_free_sorted(scanner);
if (scanner->failed) {
array_list_delete(chunk_data);
return NULL;
@@ -1436,26 +1634,15 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
continue;
}
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
SortedEntry* sorted = &((SortedEntry*)scanner->sorted_entries)[scanner->sorted_index++];
const char* name = sorted->name;
ScannerEntry* inspected = &sorted->entry;
int inspection = sorted->inspection;
ScannerEntry inspected;
char* link_rel = child_rel_path(scanner->current_rel, entry->d_name);
if (!link_rel) {
scanner->failed = true;
break;
}
int inspection = scanner_inspect_entry(&scanner->options, scanner->current_path, link_rel,
entry->d_name, &inspected);
free(link_rel);
if (inspection < 0) {
scanner->failed = true;
break;
}
if (inspection == 0) {
/* A dereferenced symlink with no referent is a partial-transfer error
(rsync exit 23): record it as a non-fatal scan I/O error. */
if (inspected.referent_error)
if (inspected->referent_error)
scanner->io_error = true;
/* A user-selection exclude protects its destination mirror from --delete
unless --delete-excluded; a size prune is always protected. Other
@@ -1463,23 +1650,23 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
--files-from the protected prefix must be the entry's bare relative
wire path, not its source path (which would not match the destination
layout and would leave the mirror deletable). */
if (inspected.excluded) {
if (inspected->excluded) {
char* protected_path;
if (scanner->relative_mode) {
protected_path = child_rel_path(scanner->current_rel, entry->d_name);
protected_path = child_rel_path(scanner->current_rel, name);
} else if (scanner->options.relative_prefix) {
char* relc = child_rel_path(scanner->current_rel, entry->d_name);
char* relc = child_rel_path(scanner->current_rel, name);
protected_path =
relc ? scanner_prefix_send_path(scanner->options.relative_prefix, relc) : NULL;
free(relc);
} else {
protected_path = path_cat(scanner->current_path, entry->d_name);
protected_path = path_cat(scanner->current_path, name);
}
if (!protected_path) {
scanner->failed = true;
break;
}
if (inspected.size_excluded)
if (inspected->size_excluded)
scanner_record_size_skipped(scanner, protected_path);
else
scanner_record_excluded(scanner, protected_path);
@@ -1487,23 +1674,22 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
}
continue;
}
char* cur_path = inspected.path;
struct stat stats = inspected.stats;
char* cur_path = inspected->path;
struct stat stats = inspected->stats;
/* --files-from allow-set and the filter layer apply to files and to
* directories (an excluded directory is not descended into). */
bool is_dir = inspected.is_directory;
char* rel = child_rel_path(scanner->current_rel, entry->d_name);
bool is_dir = inspected->is_directory;
char* rel = child_rel_path(scanner->current_rel, name);
if (!rel) {
free(cur_path);
scanner->failed = true;
break;
}
bool protect = false;
bool passes_selection = entry_passes_selection(
scanner->options.file_list, scanner->options.base_filters, scanner->current_node, rel,
entry->d_name, is_dir, scanner->options.per_dir_filters,
scanner->options.exclude_per_dir_filter_files, &protect);
scanner->options.file_list, scanner->options.base_filters, scanner->current_node, rel, name,
is_dir, scanner->options.per_dir_filters, scanner->options.exclude_per_dir_filter_files,
&protect);
/* A sender-side hide leaves the entry out of the transfer; an independent
receiver-side protect rule keeps a transferred entry's destination mirror
from being deleted. Both are recorded in the same protection set. */
@@ -1525,7 +1711,6 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
char* wrel = scanner_prefix_send_path(scanner->options.relative_prefix, rel);
if (!wrel) {
free(rel);
free(cur_path);
scanner->failed = true;
break;
}
@@ -1542,13 +1727,12 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
char* rel_copy = needs_rel ? str_dup(rel) : NULL;
free(rel);
if (rel_copy == NULL && needs_rel) {
free(cur_path);
scanner->failed = true;
break;
}
if (!passes_selection) {
scanner_note_filter(&scanner->options, name);
free(rel_copy);
free(cur_path);
continue;
}
@@ -1556,6 +1740,13 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
free(rel_copy);
if (!scanner_same_filesystem(scanner->options.one_file_system, scanner->root_dev,
stats.st_dev)) {
if (scanner->options.one_file_system > 1) {
/* rsync's -xx drops the mount-point directory entirely (the plain -x
path below keeps it as an empty directory) and prints the
--info=mount line when that category is enabled. */
scanner_note_mount(&scanner->options, cur_path);
continue;
}
/* rsync's -x/--one-file-system emits the mount-point directory entry
itself (so the destination gets an empty directory) but does NOT
descend into it. Build a payload-less directory File and hand it to
@@ -1563,12 +1754,10 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
File* mount = scanner_build_dir_file(cur_path, &stats, &scanner->options);
if (mount == NULL || !array_list_add(chunk_data, mount)) {
file_destroy(mount);
free(cur_path);
scanner->failed = true;
break;
}
scanner->current_dir_produced = true;
free(cur_path);
continue;
}
/* --list-only: list directory entries too (rsync prints them), even
@@ -1577,7 +1766,6 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
File* dir = scanner_build_dir_file(cur_path, &stats, &scanner->options);
if (dir == NULL || !array_list_add(chunk_data, dir)) {
file_destroy(dir);
free(cur_path);
scanner->failed = true;
break;
}
@@ -1586,32 +1774,29 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
int next_depth = scanner->current_depth + 1;
if (scanner->options.max_depth <= 0 || next_depth < scanner->options.max_depth) {
DirEntry* de = dir_entry_create(cur_path, next_depth, scanner->current_node);
if (!de || !queue_enqueue(scanner->directories, de)) {
if (!de || !array_list_add((ArrayList*)scanner->pending_dirs, de)) {
dir_entry_destroy(de);
scanner->failed = true;
}
}
free(cur_path);
} else {
if (scanner->options.max_depth > 0 &&
scanner->current_depth + 1 > scanner->options.max_depth) {
free(rel_copy);
free(cur_path);
continue;
}
File* file = file_create(cur_path);
free(cur_path);
if (file == NULL) {
free(rel_copy);
free(inspected.link_target);
inspected.link_target = NULL;
free(inspected->link_target);
inspected->link_target = NULL;
scanner->failed = true;
continue;
}
if (inspected.is_symlink) {
if (inspected->is_symlink) {
file->is_symlink = true;
file->symlink_target = inspected.link_target;
inspected.link_target = NULL;
file->symlink_target = inspected->link_target;
inspected->link_target = NULL;
} else {
file->data->size = stats.st_size;
}
@@ -1975,6 +2160,7 @@ static void scan_root_entry(const ScannerOptions* options, const FilterNode* roo
free(prefixed);
}
if (!passes) {
scanner_note_filter(options, entry->d_name);
free(rel);
free(cur_path);
return;
@@ -1982,6 +2168,14 @@ static void scan_root_entry(const ScannerOptions* options, const FilterNode* roo
}
if (is_dir) {
if (!scanner_same_filesystem(options->one_file_system, root_dev, st.st_dev)) {
if (options->one_file_system > 1) {
/* -xx: drop the mount-point directory entirely (rsync) and print the
--info=mount line when enabled. */
scanner_note_mount(options, cur_path);
free(rel);
free(cur_path);
return;
}
/* -x/--one-file-system: emit the mount-point directory entry (empty) but
do not descend into it (see the sequential scanner for the same rule). */
File* mount = file_create(cur_path);
@@ -2119,6 +2313,7 @@ static bool scan_root_directory(ParallelScanner* ps, const char* root_directory,
ps->failed = true;
return false;
}
log_debug_message(LOG_DEBUG_FLIST, "flist: scanning %s", root_directory);
const struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
@@ -2283,6 +2478,10 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
parallel_scanner_destroy(ps);
return NULL;
}
/* The root itself is a traversed directory (rsync counts it in
`Number of files`); the worker DirectoryScanners account for every
subdirectory below it. */
scanner_dir_count_count(options);
/* P7 Wave D: the parallel scanner never runs a DirectoryScanner over the
transfer root itself (it hands the root's immediate subdirectories to
workers), so capture the root's directory time here. */
+25 -2
View File
@@ -10,6 +10,7 @@
#include "stop_condition.h"
#include <dirent.h>
#include <stdbool.h>
#include <stddef.h>
#include <stdatomic.h>
#include <sys/types.h>
#include <threads.h>
@@ -50,7 +51,7 @@ typedef struct {
bool copy_dirlinks;
bool munge_links;
bool checksum;
bool one_file_system;
int one_file_system;
/* Phase 4 special/devices: whether device nodes (--devices) and special files
* (--specials) are preserved via recreation, and whether --copy-devices
* copies a device's content as an ordinary regular file. */
@@ -129,6 +130,17 @@ typedef struct {
/* --info=nonreg: print rsync's `skipping non-regular file "NAME"` line for a
* non-regular entry that is not being preserved. Client-only. */
bool note_nonreg;
/* --info=mount: print rsync's `[sender] skipping mount-point dir NAME` when
* -xx drops a mount-point directory. Client-only. */
bool note_mount;
/* --stats directory accounting for a `-r` run (no -t/-p): a shared counter of
* traversed directories that are NOT otherwise represented by an inline
* directory entry (rsync still counts every directory in `Number of files`).
* Incremented when a directory is opened and decremented when an empty
* directory is emitted inline (so it is counted exactly once). Atomic
* because the parallel scanner's workers share it; NULL disables the
* accounting. Client-only. */
atomic_ullong* dir_count;
/* Source root and 8-bit-output policy used to render a `--info=nonreg` name
* relative to the transfer root. Borrowed read-only. */
const char* send_directory;
@@ -188,6 +200,17 @@ typedef struct {
int current_depth;
dev_t root_dev;
bool failed;
/* rsync-order traversal: each opened directory's entries are inspected once
and buffered (an internal SortedEntry[] owned here) sorted as rsync's flist
orders them -- non-directories ascending, then directories ascending. The
entries are walked in order and child directories are collected in
`pending_dirs` (an ArrayList of DirEntry*, owned here) and pushed onto the
LIFO `directories` stack in reverse at directory exhaustion, so the emitted
stream is depth-first like rsync. `sorted_*` are reset per directory. */
void* sorted_entries;
size_t sorted_count;
size_t sorted_index;
void* pending_dirs;
/* Recursive scan: whether the open directory yielded any transferred or
descended entry. When it did not, closing it emits a directory entry so
the empty source directory is recreated at the destination (rsync
@@ -254,7 +277,7 @@ void directory_scanner_destroy(DirectoryScanner* scanner);
/* --one-file-system (-x) decision: a directory entry may be descended into
* only when the option is disabled or the entry lives on the same device as
* the transfer root. Exposed so tests can exercise the rule directly. */
bool scanner_same_filesystem(bool one_file_system, dev_t root_device, dev_t entry_device);
bool scanner_same_filesystem(int one_file_system, dev_t root_device, dev_t entry_device);
/* Relative path of an on-disk path below `root` ("" == the root itself, NULL
* when `fs_path` is not under `root`). Handles trailing slashes and a root of
+1 -1
View File
@@ -79,7 +79,7 @@ static void config_set_defaults(Config* config) {
config->cvs_exclude = false;
config->per_dir_filter = false;
config->per_dir_filter_count = 0;
config->one_file_system = false;
config->one_file_system = 0;
config->no_implied_dirs = false;
config->dirs = false;
config->rsh_command = NULL;
+3 -1
View File
@@ -463,7 +463,9 @@ typedef struct Config {
* /.rsync-filter' (the .rsync-filter files themselves are transferred); a
* repeated -F adds --filter='- .rsync-filter' so they are excluded too. */
int per_dir_filter_count;
bool one_file_system; /* -x/--one-file-system: do not cross filesystem boundaries */
int one_file_system; /* -x/--one-file-system: do not cross filesystem boundaries.
Repeated -x (rsync's -xx) drops the mount-point
directory entirely instead of recreating it empty. */
/* --no-implied-dirs: client-only. With -R, do not transfer the source
* metadata of the parent directories implied by a listed path; an unlisted
* implied parent is still created (with default attributes) so the listed
+111 -55
View File
@@ -432,6 +432,17 @@ int delete_plan_send_remaining(int fd, DeletePlanSender* sender, const ArrayList
return 0;
}
int delete_plan_send_all(int fd, DeletePlanSender* sender, const ArrayList* dirs) {
if (!sender)
return -1;
/* Root first: this also transmits the one-shot per-run config block on its
own carrier frame (see send_config_only), so it reaches the receiver even
when the scope permits no directory plan at all. */
if (delete_plan_send_root(fd, sender) != 0)
return -1;
return delete_plan_send_remaining(fd, sender, dirs);
}
/* ------------------------------------------------------------------ */
/* Receiver: delete session */
/* ------------------------------------------------------------------ */
@@ -471,6 +482,24 @@ static void notify_deleted(DeletePlanSession* session, const char* rel) {
session->observer(session->observer_context, rel);
}
/* A removed directory is reported with rsync's trailing slash (`deleting dir/`)
while files keep their bare path. */
static void notify_deleted_dir(DeletePlanSession* session, const char* rel) {
if (!session || !session->observer || !rel)
return;
size_t len = strlen(rel);
char* with_slash = malloc(len + 2);
if (!with_slash) {
session->observer(session->observer_context, rel);
return;
}
memcpy(with_slash, rel, len);
with_slash[len] = '/';
with_slash[len + 1] = '\0';
session->observer(session->observer_context, with_slash);
free(with_slash);
}
DeletePlanSession* delete_plan_session_create(const Config* config) {
if (!config)
return NULL;
@@ -696,7 +725,7 @@ static bool process_extra_dir(int dirfd, const char* name, const char* child_rel
session->deleted++;
session->planned++;
log_deleted(child_rel);
notify_deleted(session, child_rel);
notify_deleted_dir(session, child_rel);
*removed = true;
return true;
}
@@ -733,81 +762,108 @@ static bool process_children(int dirfd, const char* dir_rel, const ArrayList* ke
const ArrayList* keep_files, bool at_root, bool force_now,
const PlanSkips* skips, DeletePlanSession* session, bool* survives) {
*survives = false;
int scanfd = openat(dirfd, ".", O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (scanfd < 0)
DeleteDirEntry* entries = NULL;
size_t count = 0;
bool collect_ok = true;
if (!delete_dir_entries_collect(dirfd, &entries, &count, &collect_ok))
return false;
DIR* dir = fdopendir(scanfd);
if (!dir) {
close(scanfd);
bool operation_ok = collect_ok;
bool local_survives = false;
bool* shielded = calloc(count ? count : 1, sizeof(bool));
bool* is_extra = calloc(count ? count : 1, sizeof(bool));
bool* force = calloc(count ? count : 1, sizeof(bool));
if (!shielded || !is_extra || !force) {
free(shielded);
free(is_extra);
free(force);
delete_dir_entries_free(entries, count);
return false;
}
bool operation_ok = true;
bool local_survives = false;
const struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
/* rsync's order: extraneous subdirectories in descending name order, then
extraneous files in descending name order (kept entries survive and are not
touched here — a kept subdirectory gets its own per-directory plan). */
if (count > 1)
qsort(entries, count, sizeof(*entries), delete_dir_entry_cmp_desc);
size_t dir_count = 0;
while (dir_count < count && entries[dir_count].is_dir)
dir_count++;
for (size_t i = 0; i < count; i++) {
char* child_rel =
(strcmp(dir_rel, ".") == 0) ? str_dup(entry->d_name) : path_cat(dir_rel, entry->d_name);
(strcmp(dir_rel, ".") == 0) ? str_dup(entries[i].name) : path_cat(dir_rel, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
}
if (path_under_skip_prefix(child_rel, at_root, skips->entries, skips->count)) {
shielded[i] = true;
local_survives = true;
free(child_rel);
continue;
}
struct stat st;
if (fstatat(dirfd, entry->d_name, &st, AT_SYMLINK_NOFOLLOW) != 0) {
if (errno != ENOENT)
operation_ok = false;
free(child_rel);
continue;
}
bool is_dir = S_ISDIR(st.st_mode);
bool in_keep_dirs = is_dir && list_contains_str(keep_dirs, entry->d_name);
bool in_keep_files = !is_dir && list_contains_str(keep_files, entry->d_name);
bool is_dir = entries[i].is_dir;
bool in_keep_dirs = is_dir && list_contains_str(keep_dirs, entries[i].name);
bool in_keep_files = !is_dir && list_contains_str(keep_files, entries[i].name);
bool rule_protected =
skips->protect_rules &&
filter_rules_apply_side(skips->protect_rules, child_rel, entry->d_name, is_dir,
filter_rules_apply_side(skips->protect_rules, child_rel, entries[i].name, is_dir,
FILTER_SIDE_RECEIVER) == FILTER_ACTION_PROTECT;
if (in_keep_dirs) {
if (in_keep_dirs || in_keep_files || rule_protected) {
shielded[i] = true;
local_survives = true;
} else if (keep_dirs && !is_dir && list_contains_str(keep_dirs, entry->d_name)) {
/* Destination file blocks a source directory: clear it now, whatever the
delete timing, so the directory can be created. */
if (!process_extra_file(dirfd, entry->d_name, child_rel, true, session))
operation_ok = false;
} else if (in_keep_files) {
local_survives = true;
} else if (keep_files && is_dir && list_contains_str(keep_files, entry->d_name)) {
/* Destination directory blocks a source file: remove it now. */
bool removed = false;
if (!process_extra_dir(dirfd, entry->d_name, child_rel, true, skips, session, &removed))
operation_ok = false;
else if (!removed)
local_survives = true;
} else if (is_dir) {
if (rule_protected) {
local_survives = true;
} else {
bool removed = false;
if (!process_extra_dir(dirfd, entry->d_name, child_rel, force_now, skips, session,
&removed))
operation_ok = false;
else if (!removed)
local_survives = true;
}
} else if (rule_protected) {
local_survives = true;
/* A destination directory blocks a source file of the same name: remove
it now, whatever the delete timing, so the file can be created. */
is_extra[i] = true;
force[i] = keep_files && list_contains_str(keep_files, entries[i].name);
} else {
if (!process_extra_file(dirfd, entry->d_name, child_rel, force_now, session))
operation_ok = false;
/* A destination file blocks a source directory of the same name: clear it
now so the directory can be created. */
is_extra[i] = true;
force[i] = keep_dirs && list_contains_str(keep_dirs, entries[i].name);
}
free(child_rel);
}
closedir(dir);
/* Pass 1: extraneous subdirectories, descending. */
for (size_t i = 0; i < dir_count; i++) {
if (!is_extra[i])
continue;
char* child_rel =
(strcmp(dir_rel, ".") == 0) ? str_dup(entries[i].name) : path_cat(dir_rel, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
}
bool removed = false;
if (!process_extra_dir(dirfd, entries[i].name, child_rel, force[i] || force_now, skips, session,
&removed))
operation_ok = false;
else if (!removed)
local_survives = true;
free(child_rel);
}
/* Pass 2: extraneous files, descending. */
for (size_t i = dir_count; i < count; i++) {
if (!is_extra[i])
continue;
char* child_rel =
(strcmp(dir_rel, ".") == 0) ? str_dup(entries[i].name) : path_cat(dir_rel, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
}
if (!process_extra_file(dirfd, entries[i].name, child_rel, force[i] || force_now, session))
operation_ok = false;
free(child_rel);
}
free(shielded);
free(is_extra);
free(force);
delete_dir_entries_free(entries, count);
*survives = local_survives;
return operation_ok;
}
@@ -981,7 +1037,7 @@ static bool apply_deferred_path(DeletePlanSession* session, const Config* config
session->deleted++;
session->planned++;
log_deleted(rel);
notify_deleted(session, rel);
notify_deleted_dir(session, rel);
} else if (errno != ENOENT && errno != ENOTEMPTY && errno != EEXIST) {
close(parent_fd);
free(leaf);
+12 -2
View File
@@ -61,9 +61,19 @@ int delete_plan_send_root(int fd, DeletePlanSender* sender);
* for `path` itself; already-sent plans are skipped. */
int delete_plan_send_for_path(int fd, DeletePlanSender* sender, const char* path, bool is_dir);
/* Send the plan for every directory in `dirs` that has not been transmitted
* yet. Called after the data stream so an empty source directory's plan still
* clears its destination extras even though no file frame triggered it. */
* yet. */
int delete_plan_send_remaining(int fd, DeletePlanSender* sender, const ArrayList* dirs);
/* Transmit the COMPLETE per-directory plan set in one pass, before any data
* frame: the root plan (with the one-shot per-run config block on its carrier
* frame) followed by every directory in `dirs`. Because the whole plan set is
* known from the path-only pre-scan, sending it all up front means a
* mid-transfer abort has already applied every planned removal, matching
* rsync's generator (which runs ahead of its throttled sender). A completed
* run is unaffected. `dirs` is the set of directories whose direct children
* were enumerated (the scanner's plan_dirs sink), so a merely listed but
* untraversed directory never gets a plan and its mirror is left intact.
* Returns -1 on I/O error. */
int delete_plan_send_all(int fd, DeletePlanSender* sender, const ArrayList* dirs);
/* ---- Receiver: delete session ---- */
+86 -41
View File
@@ -1389,6 +1389,43 @@ bool file_basis_content_required(const Config* config) {
return config != NULL && config->verify_basis;
}
/* Probe one candidate basis file: open it (confined, O_NOFOLLOW) and apply
rsync's metadata quick-check; under --verify-basis also hash its bytes and
require the sender's digest. On a hit record `candidate` in `out` and return
true. The caller retains ownership of `candidate`. */
static bool basis_match_probe(const Config* config, const char* candidate,
unsigned long long check_size, time_t check_mtime,
long check_mtime_nsec, const uint8_t* check_digest,
size_t check_digest_len, BasisDestType type, BasisMatch* out) {
int fd;
struct stat st;
if (!basis_open_regular(candidate, check_size, &fd, &st))
return false;
bool hit = false;
if (file_basis_quick_match(config, &st, check_mtime, check_mtime_nsec)) {
hit = true;
if (file_basis_content_required(config)) {
uint8_t basis_digest[CHECKSUM_MAX_DIGEST_LEN];
size_t basis_len = 0;
bool hashed = checksum_digest_fd((ChecksumAlgo)config->checksum_algo, config->checksum_seed,
fd, basis_digest, sizeof(basis_digest), &basis_len);
hit = hashed && basis_len == check_digest_len && check_digest_len > 0 &&
memcmp(basis_digest, check_digest, check_digest_len) == 0;
}
}
close(fd);
if (!hit)
return false;
char* owned = str_dup(candidate);
if (!owned)
return false;
out->hit = true;
out->type = type;
out->basis_path = owned;
out->st = st;
return true;
}
/* Search the basis-dir list in command-line order and return the first match.
By default (no --verify-basis) rsync's metadata quick-check is sufficient:
basis_open_regular has already required an equal size, and
@@ -1403,7 +1440,22 @@ bool file_basis_content_required(const Config* config) {
--dry-run passes false because hashing a basis against a client-supplied
digest would be a 1-bit content oracle. Without --verify-basis a dry-run can
still confirm the metadata-only hit without reading any basis bytes, matching
rsync's read-only quick-check. */
rsync's read-only quick-check.
Path resolution (rsync 3.4.1 parity): rsync resolves a relative
--compare-dest/--copy-dest/--link-dest DIR against the destination directory
(the receiver's cwd) and appends the file's TRANSFER-RELATIVE name, e.g.
`--compare-dest=basis` with `rsync src/ dst/` probes `dst/basis/<name-inside-src>`.
FastSync's receive root IS the destination directory, but its default transfer
mirrors the absolute source path below that root, so check_path carries the
source-root scaffolding rsync would not append. Recover rsync's spelling with
utils_strip_transfer_root for a relative DIR; under -R/--files-from the wire
path is already transfer-relative, so it is used as-is. A relative DIR also
probes the historical mirror-appended spelling as a fallback, so existing
FastSync-laid-out snapshot trees keep resolving. An absolute DIR is used
verbatim and keeps appending the destination-relative check_path (FastSync's
mirrored layout). Every candidate stays confined to the authorized root by
file_open_secure_parent. */
static bool basis_match_find(const Config* config, const char* check_path,
unsigned long long check_size, time_t check_mtime,
long check_mtime_nsec, const uint8_t* check_digest,
@@ -1416,47 +1468,39 @@ static bool basis_match_find(const Config* config, const char* check_path,
the basis bytes. */
if (file_basis_content_required(config) && !hash_content)
return false;
const char* transfer_rel = check_path;
if (!config->relative && config->files_from_set == NULL)
transfer_rel = utils_strip_transfer_root(check_path, config->send_directory);
for (int i = 0; i < config->basis_count; i++) {
const BasisDest* entry = &config->basis_dirs[i];
/* An absolute basis path is used verbatim (rsync semantics); a relative one
is resolved below the receive root. Both remain subject to the receiver's
authorized-root confinement inside file_open_secure_parent. */
char* basis_dir = entry->path[0] == '/' ? str_dup(entry->path)
: path_cat(config->receive_root_directory, entry->path);
bool absolute = entry->path[0] == '/';
char* basis_dir =
absolute ? str_dup(entry->path) : path_cat(config->receive_root_directory, entry->path);
if (!basis_dir)
continue;
char* candidate = path_cat(basis_dir, check_path);
free(basis_dir);
if (!candidate)
continue;
int fd;
struct stat st;
if (basis_open_regular(candidate, check_size, &fd, &st)) {
if (file_basis_quick_match(config, &st, check_mtime, check_mtime_nsec)) {
bool hit = true;
if (file_basis_content_required(config)) {
uint8_t basis_digest[CHECKSUM_MAX_DIGEST_LEN];
size_t basis_len = 0;
bool hashed =
checksum_digest_fd((ChecksumAlgo)config->checksum_algo, config->checksum_seed, fd,
basis_digest, sizeof(basis_digest), &basis_len);
hit = hashed && basis_len == check_digest_len && check_digest_len > 0 &&
memcmp(basis_digest, check_digest, check_digest_len) == 0;
}
if (hit) {
out->hit = true;
out->type = entry->type;
out->basis_path = candidate;
candidate = NULL; /* ownership transferred to out */
out->st = st;
close(fd);
return true;
}
}
close(fd);
const char* names[2];
int name_count = 0;
if (absolute)
names[name_count++] = check_path;
else
names[name_count++] = transfer_rel;
if (!absolute && strcmp(transfer_rel, check_path) != 0)
names[name_count++] = check_path; /* historical mirror-appended spelling */
bool found = false;
for (int n = 0; n < name_count && !found; n++) {
char* candidate = path_cat(basis_dir, names[n]);
if (!candidate)
continue;
found = basis_match_probe(config, candidate, check_size, check_mtime, check_mtime_nsec,
check_digest, check_digest_len, entry->type, out);
free(candidate);
}
free(candidate);
free(basis_dir);
if (found)
return true;
}
return false;
}
@@ -1489,9 +1533,12 @@ static bool basis_match_find(const Config* config, const char* check_path,
* followed and nothing outside the destination root is ever read;
* * dotfiles, directories, the target's own name, and the .fastsync-stage /
* temp scratch names are never candidates;
* * size gate = the delta engine's own bounds (delta_should_attempt: both
* files >= DELTA_MIN_FILE_SIZE, <= delta_max_file_size, ratio <= 10x),
* because FastSync's delta engine cannot use a basis outside them;
* * size gate = rsync's, NOT the ordinary delta engine's bounds: any
* non-empty regular sibling up to the receiver's whole-file buffer cap is
* eligible, regardless of the 16 KiB delta minimum or the 10x delta size
* ratio (rsync's find_fuzzy has no delta-size gate at all). The delta
* engine consumes the fuzzy basis through the same signature handshake
* whether or not it is inside delta_should_attempt's window;
* * first pass = an exact size+mtime match wins regardless of name (rsync's
* "fuzzy size/modtime match");
* * otherwise the winner minimizes rsync's weighted Levenshtein distance
@@ -1639,8 +1686,7 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
long check_mtime_nsec, unsigned long long* out_size) {
*out_size = 0;
if (!config || !config->receive_root_directory || !config->fuzzy || !config->use_delta ||
!check_path || check_size < DELTA_MIN_FILE_SIZE || check_size > config->delta_max_file_size ||
check_size > MAX_RECEIVE_WHOLE_FILE_SIZE)
!check_path || check_size > MAX_RECEIVE_WHOLE_FILE_SIZE)
return NULL;
char* full_path = path_cat(config->receive_root_directory, check_path);
@@ -1717,8 +1763,7 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
if (fstatat(dir_fd, name, &st, AT_SYMLINK_NOFOLLOW) != 0 || !S_ISREG(st.st_mode))
continue;
unsigned long long cand_size = (unsigned long long)st.st_size;
if (cand_size == 0 || cand_size > MAX_RECEIVE_WHOLE_FILE_SIZE ||
!delta_should_attempt(cand_size, check_size, config->delta_max_file_size))
if (cand_size == 0 || cand_size > MAX_RECEIVE_WHOLE_FILE_SIZE)
continue;
long cand_nsec = 0;
#ifdef __linux__
+18 -2
View File
@@ -13,7 +13,19 @@ typedef enum {
LOG_DEBUG_PROTO = 1u << 1,
LOG_DEBUG_PACK = 1u << 2,
LOG_DEBUG_UTIL = 1u << 3,
LOG_DEBUG_ALL = (1u << 4) - 1,
/* rsync --debug categories that now map to a natural FastSync event:
* flist (file-list scan progress), del (deletions), hash/deltasum
* (whole-file hashing and delta-sum generation), recv (receiver
* responses/signatures), filter (selection/exclusion decisions) and send
* (files handed to the sender). Only emitted when the category is
* explicitly enabled; a normal run stays silent. */
LOG_DEBUG_FLIST = 1u << 4,
LOG_DEBUG_DEL = 1u << 5,
LOG_DEBUG_HASH = 1u << 6,
LOG_DEBUG_RECV = 1u << 7,
LOG_DEBUG_FILTER = 1u << 8,
LOG_DEBUG_SEND = 1u << 9,
LOG_DEBUG_ALL = (1u << 10) - 1,
} LogDebugFlag;
typedef enum {
@@ -38,9 +50,13 @@ typedef enum {
the info_level bitset (there is no separate Config field) and is never set
by --info=all (which selects level 1). */
LOG_INFO_NAME_UPTODATE = 1u << 10,
/* --info=mount: print rsync's `[sender] skipping mount-point dir NAME` when
* -xx/--one-file-system drops a mount-point directory (FastSync's client is
* the sender). */
LOG_INFO_MOUNT = 1u << 11,
LOG_INFO_ALL = LOG_INFO_COPY | LOG_INFO_MISC | LOG_INFO_SKIP | LOG_INFO_STATS | LOG_INFO_DEL |
LOG_INFO_REMOVE | LOG_INFO_NAME | LOG_INFO_FLIST | LOG_INFO_NONREG |
LOG_INFO_PROGRESS,
LOG_INFO_PROGRESS | LOG_INFO_MOUNT,
} LogInfoFlag;
void log_message(LogLevel log_level, const char* message, ...);
+1
View File
@@ -50,6 +50,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
protocol_session_set_max_alloc(&context->allocation_session, config->max_alloc);
context->dir_entries = NULL;
context->dir_entries_mutex_init = false;
atomic_init(&context->dir_count, 0);
context->delete_limit = false;
int init = 0;
if (config->use_metadata) {
+4
View File
@@ -119,6 +119,10 @@ typedef struct {
ArrayList* dir_entries;
mtx_t dir_entries_mutex;
bool dir_entries_mutex_init;
/* --stats directory accounting for a `-r` scan (no directory metadata):
shared by the parallel scanner workers, read by the sender thread once the
scanner is done. See ScannerOptions.dir_count. */
atomic_ullong dir_count;
/* Set by the sender thread when the receiver reported a --max-delete-capped
deletion (STATUS_DELETE_LIMIT): the transfer succeeded and the process must
exit 25 like rsync. Read by the caller after the sender thread is joined. */
+17
View File
@@ -139,6 +139,23 @@ void* queue_dequeue(Queue* queue) {
return item;
}
bool queue_push(Queue* queue, void* item) {
return queue_enqueue(queue, item);
}
void* queue_pop(Queue* queue) {
if (queue == NULL || queue_is_empty(queue)) {
log_perror("ERROR: Could not pop from null or empty queue.");
return NULL;
}
queue->rear = (queue->rear - 1 + queue->capacity) % queue->capacity;
void* item = queue->items[queue->rear];
queue->items[queue->rear] = NULL;
queue->size--;
return item;
}
void* queue_dequeue_multithreaded(Queue* queue, mtx_t* mutex, cnd_t* condition_not_empty,
cnd_t* condition_not_full, const bool* other_thread_done) {
mtx_lock(mutex);
+7
View File
@@ -28,4 +28,11 @@ void* queue_dequeue(Queue* queue);
void* queue_dequeue_multithreaded(Queue* queue, mtx_t* mutex, cnd_t* condition_not_empty,
cnd_t* condition_not_full, const bool* other_thread_done);
/* LIFO stack operations over the same ring buffer. queue_push() is the enqueue
primitive; queue_pop() removes from the rear, so a sequence of pushes is
returned in reverse order. Used by the sequential scanner's depth-first
traversal. */
bool queue_push(Queue* queue, void* item);
void* queue_pop(Queue* queue);
#endif
+363 -152
View File
@@ -672,22 +672,24 @@ static bool is_synced_dir(const PathIndex* dirs, const char* rel) {
return path_index_contains(dirs, rel[0] == '\0' ? "." : rel);
}
/* Remove the extras directly inside the directory open on `dirfd`, recursing
into every child directory so kept content below a synchronized prefix is
reached. `all_removed` reports whether every child entry was removed (so the
caller may rmdir this directory). A child directory is never removed when it
is itself a synchronized directory or holds kept content; with a dirs index
supplied, direct children of a non-synchronized directory are never extras at
all (they are left in place but still descended into). Symlinks are unlinked
like any other non-directory extra (never followed). */
static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* keep,
const PathIndex* dirs, DeleteBudget* budget,
const DeleteSkipEntry* skips, int skip_count,
const FilterRuleList* protect_rules, bool parent_deletable,
bool* all_removed, DeletePathObserver observer,
void* observer_context) {
/* openat(dirfd, ".") opens an independent file description: a dup() would
share dirfd's file offset and a prior pass could leave the stream drained. */
/* Unsigned byte-wise string compare, matching rsync's u_strcmp (a signed
strcmp would order bytes >= 0x80 differently). */
static int delete_name_cmp(const char* a, const char* b) {
const unsigned char* pa = (const unsigned char*)a;
const unsigned char* pb = (const unsigned char*)b;
while (*pa != '\0' && *pa == *pb) {
pa++;
pb++;
}
return (int)*pa - (int)*pb;
}
bool delete_dir_entries_collect(int dirfd, DeleteDirEntry** out, size_t* count,
bool* operation_ok) {
*out = NULL;
*count = 0;
if (operation_ok)
*operation_ok = true;
int scanfd = openat(dirfd, ".", O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (scanfd < 0)
return false;
@@ -696,17 +698,127 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* k
close(scanfd);
return false;
}
bool operation_ok = true;
bool local_survives = false;
/* A directory is deletable when it or ANY ancestor is synchronized; the
`parent_deletable` flag carries that down the recursion so dest-only
directories below a synchronized root are removed wholesale. */
bool deletable = parent_deletable || is_synced_dir(dirs, rel_path);
DeleteDirEntry* entries = NULL;
size_t used = 0;
size_t capacity = 0;
bool ok = true;
const struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
char* child_rel = path_cat((char*)rel_path, entry->d_name);
struct stat st;
if (fstatat(dirfd, entry->d_name, &st, AT_SYMLINK_NOFOLLOW) != 0) {
if (errno != ENOENT && operation_ok)
*operation_ok = false;
continue;
}
if (used == capacity) {
size_t next = capacity == 0 ? 16 : capacity * 2;
DeleteDirEntry* grown = realloc(entries, next * sizeof(*grown));
if (!grown) {
ok = false;
break;
}
entries = grown;
capacity = next;
}
entries[used].name = str_dup(entry->d_name);
if (!entries[used].name) {
ok = false;
break;
}
entries[used].is_dir = S_ISDIR(st.st_mode);
used++;
}
closedir(dir);
if (!ok) {
delete_dir_entries_free(entries, used);
return false;
}
*out = entries;
*count = used;
return true;
}
void delete_dir_entries_free(DeleteDirEntry* entries, size_t count) {
if (!entries)
return;
for (size_t i = 0; i < count; i++)
free(entries[i].name);
free(entries);
}
/* rsync's extraneous-entry order: subdirectories before files, each group in
descending name order. */
int delete_dir_entry_cmp_desc(const void* a, const void* b) {
const DeleteDirEntry* ea = a;
const DeleteDirEntry* eb = b;
if (ea->is_dir != eb->is_dir)
return ea->is_dir ? -1 : 1;
return -delete_name_cmp(ea->name, eb->name);
}
/* rsync's kept-subdirectory order: plain ascending name. */
int delete_dir_entry_cmp_asc(const void* a, const void* b) {
const DeleteDirEntry* ea = a;
const DeleteDirEntry* eb = b;
return delete_name_cmp(ea->name, eb->name);
}
/* Remove the extras directly inside the directory open on `dirfd`, recursing
into every child directory so kept content below a synchronized prefix is
reached. `all_removed` reports whether every child entry was removed (so the
caller may rmdir this directory). A child directory is never removed when it
is itself a synchronized directory or holds kept content; with a dirs index
supplied, direct children of a non-synchronized directory are never extras at
all (they are left in place but still descended into). Symlinks are unlinked
like any other non-directory extra (never followed).
Entries are processed in rsync's order (extraneous subdirectories in
descending name order, then extraneous files, then kept subdirectories in
ascending order) rather than readdir() order, so `--max-delete` leaves the
same survivors and the `--info=del`/dry-run line order matches rsync. */
static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* keep,
const PathIndex* dirs, DeleteBudget* budget,
const DeleteSkipEntry* skips, int skip_count,
const FilterRuleList* protect_rules, bool parent_deletable,
bool* all_removed, DeletePathObserver observer,
void* observer_context) {
DeleteDirEntry* entries = NULL;
size_t count = 0;
bool collect_ok = true;
if (!delete_dir_entries_collect(dirfd, &entries, &count, &collect_ok))
return false;
bool operation_ok = collect_ok;
bool local_survives = false;
bool* shielded = calloc(count ? count : 1, sizeof(bool));
bool* is_extra = calloc(count ? count : 1, sizeof(bool));
if (!shielded || !is_extra) {
free(shielded);
free(is_extra);
delete_dir_entries_free(entries, count);
return false;
}
/* A directory is deletable when it or ANY ancestor is synchronized; the
`parent_deletable` flag carries that down the recursion so dest-only
directories below a synchronized root are removed wholesale. */
bool deletable = parent_deletable || is_synced_dir(dirs, rel_path);
bool at_root = rel_path[0] == '\0';
/* Reproduce rsync's traversal order: extraneous subdirectories in descending
name order, then extraneous files in descending name order, and kept
subdirectories only afterwards (ascending). Sorting up front also fixes the
identity of the survivors under a partial --max-delete. */
if (count > 1)
qsort(entries, count, sizeof(*entries), delete_dir_entry_cmp_desc);
size_t dir_count = 0;
while (dir_count < count && entries[dir_count].is_dir)
dir_count++;
/* Classify every entry up front (the verdict does not depend on processing
order) so the ordered passes below can act on it. */
for (size_t i = 0; i < count; i++) {
char* child_rel = path_cat((char*)rel_path, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
@@ -718,93 +830,142 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* k
top-level-only prefix) and the basis prefixes are protected: a nested
destination directory that happens to be called .fastsync-stage is
ordinary content. */
if (path_under_skip_prefix(child_rel, rel_path[0] == '\0', skips, skip_count)) {
if (path_under_skip_prefix(child_rel, at_root, skips, skip_count)) {
shielded[i] = true;
local_survives = true;
free(child_rel);
continue;
}
struct stat st;
if (fstatat(dirfd, entry->d_name, &st, AT_SYMLINK_NOFOLLOW) != 0) {
if (errno != ENOENT)
operation_ok = false;
free(child_rel);
continue;
}
bool is_dir = S_ISDIR(st.st_mode);
if (protect_rules && filter_rules_apply_side(protect_rules, child_rel, entry->d_name, is_dir,
FILTER_SIDE_RECEIVER) == FILTER_ACTION_PROTECT) {
} else if (protect_rules &&
filter_rules_apply_side(protect_rules, child_rel, entries[i].name, entries[i].is_dir,
FILTER_SIDE_RECEIVER) == FILTER_ACTION_PROTECT) {
/* A first-match protect rule shields the extra; for a directory the whole
subtree is shielded (rsync prunes an excluded directory), so do not
descend. */
shielded[i] = true;
local_survives = true;
free(child_rel);
} else if (entries[i].is_dir) {
bool child_synced = dirs && path_index_contains(dirs, child_rel);
is_extra[i] = deletable && !child_synced && !keep_is_dir(keep, child_rel);
if (!is_extra[i])
local_survives = true;
} else {
is_extra[i] = deletable && !keep_is_file(keep, child_rel);
if (!is_extra[i])
local_survives = true;
}
free(child_rel);
}
/* Pass 1: extraneous subdirectories, descending. */
for (size_t i = 0; i < dir_count; i++) {
if (!is_extra[i])
continue;
char* child_rel = path_cat((char*)rel_path, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
}
if (is_dir) {
int childfd = openat(dirfd, entry->d_name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool child_all_removed = false;
if (childfd >= 0) {
if (!delete_extras_fd(childfd, child_rel, keep, dirs, budget, skips, skip_count,
protect_rules, deletable, &child_all_removed, observer,
observer_context))
operation_ok = false;
close(childfd);
} else if (errno != ENOENT) {
int childfd = openat(dirfd, entries[i].name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool child_all_removed = false;
if (childfd >= 0) {
if (!delete_extras_fd(childfd, child_rel, keep, dirs, budget, skips, skip_count,
protect_rules, deletable, &child_all_removed, observer,
observer_context))
operation_ok = false;
}
bool child_synced = dirs && path_index_contains(dirs, child_rel);
if (child_synced || keep_is_dir(keep, child_rel)) {
/* A synchronized directory and a directory holding kept content are
never removed. */
local_survives = true;
} else if (child_all_removed && deletable) {
if (budget->deleted >= budget->max_delete) {
budget->limit_hit = true;
budget->skipped++;
local_survives = true;
} else if (unlinkat(dirfd, entry->d_name, AT_REMOVEDIR) != 0) {
/* ENOENT: already gone (fine). ENOTEMPTY/EEXIST: the directory
still holds entries the walker leaves in place (a protected
excluded prefix, a kept file the manifest protects, a symlink);
rsync leaves such a directory behind, so this is not an error.
Only genuine I/O failures abort the deletion. */
if (errno != ENOENT && errno != ENOTEMPTY && errno != EEXIST)
operation_ok = false;
local_survives = true;
} else {
budget->deleted++;
if (observer)
observer(observer_context, child_rel);
}
} else {
local_survives = true;
}
} else {
bool found = keep_is_file(keep, child_rel);
if (found || !deletable) {
/* Kept file, or a child of a directory that is not synchronized: never
an extra for this run. */
local_survives = true;
} else if (budget->deleted >= budget->max_delete) {
close(childfd);
} else if (errno != ENOENT) {
operation_ok = false;
}
if (child_all_removed && deletable) {
if (budget->deleted >= budget->max_delete) {
budget->limit_hit = true;
budget->skipped++;
local_survives = true;
} else if (unlinkat(dirfd, entry->d_name, 0) != 0) {
if (errno != ENOENT)
} else if (unlinkat(dirfd, entries[i].name, AT_REMOVEDIR) != 0) {
/* ENOENT: already gone (fine). ENOTEMPTY/EEXIST: the directory still
holds entries the walker leaves in place (a protected excluded
prefix, a kept file the manifest protects, a symlink); rsync leaves
such a directory behind, so this is not an error. Only genuine I/O
failures abort the deletion. */
if (errno != ENOENT && errno != ENOTEMPTY && errno != EEXIST)
operation_ok = false;
local_survives = true;
} else {
budget->deleted++;
/* rsync reports a removed directory with a trailing slash. */
if (observer) {
size_t len = strlen(child_rel);
char* with_slash = malloc(len + 2);
if (with_slash) {
memcpy(with_slash, child_rel, len);
with_slash[len] = '/';
with_slash[len + 1] = '\0';
observer(observer_context, with_slash);
free(with_slash);
} else {
observer(observer_context, child_rel);
}
}
}
} else {
local_survives = true;
}
free(child_rel);
}
/* Pass 2: extraneous files, descending. */
for (size_t i = dir_count; i < count; i++) {
if (!is_extra[i])
continue;
if (budget->deleted >= budget->max_delete) {
budget->limit_hit = true;
budget->skipped++;
local_survives = true;
} else if (unlinkat(dirfd, entries[i].name, 0) != 0) {
if (errno != ENOENT)
operation_ok = false;
local_survives = true;
} else {
budget->deleted++;
char* child_rel = path_cat((char*)rel_path, entries[i].name);
if (child_rel) {
if (observer)
observer(observer_context, child_rel);
char* escaped_path = output_escape(child_rel, log_get_8_bit_output());
fprintf(stderr, " Deleted: %s\n", escaped_path ? escaped_path : "<allocation failed>");
free(escaped_path);
}
free(child_rel);
}
}
/* Pass 3: kept subdirectories, ascending (rsync descends into these only
after the parent's own extras have been handled). */
for (size_t i = dir_count; i-- > 0;) {
if (is_extra[i] || shielded[i])
continue;
char* child_rel = path_cat((char*)rel_path, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
}
int childfd = openat(dirfd, entries[i].name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool child_all_removed = false;
if (childfd >= 0) {
if (!delete_extras_fd(childfd, child_rel, keep, dirs, budget, skips, skip_count,
protect_rules, deletable, &child_all_removed, observer,
observer_context))
operation_ok = false;
close(childfd);
} else if (errno != ENOENT) {
operation_ok = false;
}
/* A kept/synchronized directory is never removed. */
local_survives = true;
free(child_rel);
}
closedir(dir);
free(shielded);
free(is_extra);
delete_dir_entries_free(entries, count);
*all_removed = !local_survives;
return operation_ok;
}
@@ -817,97 +978,147 @@ static bool list_extras_fd(int dirfd, const char* rel_path, const PathIndex* kee
const DeleteSkipEntry* skips, int skip_count,
const FilterRuleList* protect_rules, bool parent_deletable,
bool* all_removed) {
int scanfd = openat(dirfd, ".", O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (scanfd < 0)
DeleteDirEntry* entries = NULL;
size_t count = 0;
bool collect_ok = true;
if (!delete_dir_entries_collect(dirfd, &entries, &count, &collect_ok))
return false;
DIR* dir = fdopendir(scanfd);
if (!dir) {
close(scanfd);
bool operation_ok = collect_ok;
bool local_survives = false;
bool* shielded = calloc(count ? count : 1, sizeof(bool));
bool* is_extra = calloc(count ? count : 1, sizeof(bool));
if (!shielded || !is_extra) {
free(shielded);
free(is_extra);
delete_dir_entries_free(entries, count);
return false;
}
bool operation_ok = true;
bool local_survives = false;
bool deletable = parent_deletable || is_synced_dir(dirs, rel_path);
const struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
char* child_rel = path_cat((char*)rel_path, entry->d_name);
bool at_root = rel_path[0] == '\0';
/* Mirror the delete walk's rsync order (extraneous subdirectories descending,
then extraneous files descending, then kept subdirectories ascending). */
if (count > 1)
qsort(entries, count, sizeof(*entries), delete_dir_entry_cmp_desc);
size_t dir_count = 0;
while (dir_count < count && entries[dir_count].is_dir)
dir_count++;
for (size_t i = 0; i < count; i++) {
char* child_rel = path_cat((char*)rel_path, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
}
if (path_under_skip_prefix(child_rel, rel_path[0] == '\0', skips, skip_count)) {
if (path_under_skip_prefix(child_rel, at_root, skips, skip_count)) {
shielded[i] = true;
local_survives = true;
free(child_rel);
continue;
}
struct stat st;
if (fstatat(dirfd, entry->d_name, &st, AT_SYMLINK_NOFOLLOW) != 0) {
if (errno != ENOENT)
operation_ok = false;
free(child_rel);
continue;
}
bool is_dir = S_ISDIR(st.st_mode);
if (protect_rules && filter_rules_apply_side(protect_rules, child_rel, entry->d_name, is_dir,
FILTER_SIDE_RECEIVER) == FILTER_ACTION_PROTECT) {
} else if (protect_rules &&
filter_rules_apply_side(protect_rules, child_rel, entries[i].name, entries[i].is_dir,
FILTER_SIDE_RECEIVER) == FILTER_ACTION_PROTECT) {
/* Mirror the delete walk: a protected entry is never reported as a
would-delete and a protected directory's subtree is not enumerated. */
shielded[i] = true;
local_survives = true;
free(child_rel);
} else if (entries[i].is_dir) {
bool child_synced = dirs && path_index_contains(dirs, child_rel);
is_extra[i] = deletable && !child_synced && !keep_is_dir(keep, child_rel);
if (!is_extra[i])
local_survives = true;
} else {
is_extra[i] = deletable && !keep_is_file(keep, child_rel);
if (!is_extra[i])
local_survives = true;
}
free(child_rel);
}
/* Pass 1: extraneous subdirectories, descending (recorded after contents). */
for (size_t i = 0; i < dir_count; i++) {
if (!is_extra[i])
continue;
char* child_rel = path_cat((char*)rel_path, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
}
if (is_dir) {
int childfd = openat(dirfd, entry->d_name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool child_all_removed = false;
if (childfd >= 0) {
if (!list_extras_fd(childfd, child_rel, keep, dirs, out, recorded, skips, skip_count,
protect_rules, deletable, &child_all_removed))
operation_ok = false;
close(childfd);
} else if (errno != ENOENT) {
int childfd = openat(dirfd, entries[i].name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool child_all_removed = false;
if (childfd >= 0) {
if (!list_extras_fd(childfd, child_rel, keep, dirs, out, recorded, skips, skip_count,
protect_rules, deletable, &child_all_removed))
operation_ok = false;
close(childfd);
} else if (errno != ENOENT) {
operation_ok = false;
}
if (child_all_removed && deletable) {
size_t len = strlen(child_rel);
char* copy = malloc(len + 2);
if (!copy) {
operation_ok = false;
}
bool child_synced = dirs && path_index_contains(dirs, child_rel);
if (child_synced || keep_is_dir(keep, child_rel)) {
local_survives = true;
} else if (child_all_removed && deletable) {
size_t len = strlen(child_rel);
char* copy = malloc(len + 2);
if (!copy) {
operation_ok = false;
} else {
memcpy(copy, child_rel, len);
copy[len] = '/';
copy[len + 1] = '\0';
if (!array_list_add(out, copy)) {
free(copy);
operation_ok = false;
} else {
(*recorded)++;
}
}
} else {
local_survives = true;
}
} else {
bool found = keep_is_file(keep, child_rel);
if (found || !deletable) {
local_survives = true;
} else {
char* copy = str_dup(child_rel);
if (!copy || !array_list_add(out, copy)) {
memcpy(copy, child_rel, len);
copy[len] = '/';
copy[len + 1] = '\0';
if (!array_list_add(out, copy)) {
free(copy);
operation_ok = false;
} else {
(*recorded)++;
}
}
} else {
local_survives = true;
}
free(child_rel);
}
closedir(dir);
/* Pass 2: extraneous files, descending. */
for (size_t i = dir_count; i < count; i++) {
if (!is_extra[i])
continue;
char* child_rel = path_cat((char*)rel_path, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
}
char* copy = str_dup(child_rel);
if (!copy || !array_list_add(out, copy)) {
free(copy);
operation_ok = false;
} else {
(*recorded)++;
}
free(child_rel);
}
/* Pass 3: kept subdirectories, ascending. */
for (size_t i = dir_count; i-- > 0;) {
if (is_extra[i] || shielded[i])
continue;
char* child_rel = path_cat((char*)rel_path, entries[i].name);
if (!child_rel) {
operation_ok = false;
continue;
}
int childfd = openat(dirfd, entries[i].name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool child_all_removed = false;
if (childfd >= 0) {
if (!list_extras_fd(childfd, child_rel, keep, dirs, out, recorded, skips, skip_count,
protect_rules, deletable, &child_all_removed))
operation_ok = false;
close(childfd);
} else if (errno != ENOENT) {
operation_ok = false;
}
local_survives = true;
free(child_rel);
}
free(shielded);
free(is_extra);
delete_dir_entries_free(entries, count);
*all_removed = !local_survives;
return operation_ok;
}
+22
View File
@@ -134,6 +134,28 @@ typedef struct {
only DIRECT children of the destination root, i.e. child_rel has no '/'). */
bool path_under_skip_prefix(const char* child_rel, bool at_root, const DeleteSkipEntry* skips,
int skip_count);
/* One destination-directory entry collected up front so the delete walkers can
reproduce rsync's traversal order instead of readdir() order. rsync processes
a directory's extraneous subdirectories first (descending name, depth-first),
then its extraneous files (descending name), and only afterwards descends into
its kept subdirectories (ascending name). */
typedef struct {
char* name;
bool is_dir;
} DeleteDirEntry;
/* Collect the entries of the directory open on `dirfd` (excluding "." and ".."),
stat'ing each with AT_SYMLINK_NOFOLLOW. On success *out is a malloc'd array of
*count entries whose names the caller frees with delete_dir_entries_free().
Returns false on an allocation/readdir failure; a vanished entry (ENOENT) is
skipped, any other stat failure is reported through *operation_ok while the
walk continues. */
bool delete_dir_entries_collect(int dirfd, DeleteDirEntry** out, size_t* count, bool* operation_ok);
void delete_dir_entries_free(DeleteDirEntry* entries, size_t count);
/* Sort comparators: `_desc` orders subdirectories before files and each group by
descending name (rsync's extraneous-entry order); `_asc` orders plain ascending
name (rsync's kept-subdirectory order). */
int delete_dir_entry_cmp_desc(const void* a, const void* b);
int delete_dir_entry_cmp_asc(const void* a, const void* b);
/* Remove files/dirs/symlinks under dest_root that are not listed in manifest
without ever descending into a protected prefix (see DeleteSkipEntry). When
`synced_dirs` is non-NULL, extras are only removed directly inside a directory
@@ -0,0 +1,291 @@
"""Differential coverage for the delete-timing ABORT BOUNDARY (A9/A10).
rsync's generator runs ahead of its throttled sender, so on a mid-transfer abort
it has already removed every extra it planned. FastSync now transmits the
COMPLETE per-directory plan set before the first data frame, so an abort has the
same effect. Before that change FastSync only removed the extras of the
directories its (slower) data stream had reached, and ``-d/--dirs`` used an
end-of-transfer commit that removed nothing on abort.
These tests abort both tools mid-transfer and assert the destination extras
removed match real ``rsync 3.4.1``. The rsync side is driven locally with
``--bwlimit`` and a small timing window (its generator's delete list is computed
long before the throttled payload finishes); the FastSync side uses the
byte-deterministic slicing proxy from ``test_delete_timing_parity``.
"""
import os
import shutil
import subprocess
import sys
import time
import pytest
sys.path.insert(0, os.path.dirname(__file__))
from common import ( # noqa: E402
TEST_DATA_DIR,
ServerManager,
clean_dir,
get_dest_received_dir,
run_client,
)
from test_delete_timing_parity import _SlicingProxy # noqa: E402
RSYNC = shutil.which("rsync")
requires_rsync = pytest.mark.skipif(RSYNC is None, reason="rsync 3.4.1 not installed")
# Exceeds the 10 MiB scanner chunk, so the next directory lands in a later chunk
# (still unreached when the proxy cuts the stream).
BIG_BYTES = 16 * 1024 * 1024
# Cut well past the (small) config + delete-plan frames and into the big payload,
# so the receiver has provably processed every plan before the abort.
MID_TRANSFER_BYTES = 256 * 1024
PROXY_THROTTLE = 0.001
# Throttle rsync's sender so the generator has deleted long before the payload
# finishes, then interrupt it mid-transfer.
RSYNC_BWLIMIT = 512 # KiB/s -> ~32 s for 16 MiB
RSYNC_ABORT_DELAY = 1.5
def _write(path, content):
os.makedirs(os.path.dirname(path), exist_ok=True)
with open(path, "wb") as fh:
fh.write(content)
def _rsync_aborted(args, delay=RSYNC_ABORT_DELAY):
"""Start rsync, let its generator run, then interrupt it mid-transfer."""
env = dict(os.environ, LC_ALL="C")
proc = subprocess.Popen([RSYNC] + args, stdout=subprocess.PIPE, stderr=subprocess.PIPE,
text=True, env=env)
time.sleep(delay)
proc.terminate()
try:
proc.wait(timeout=10)
except subprocess.TimeoutExpired:
proc.kill()
proc.wait(timeout=5)
return proc
class TestDeleteDuringAbortBoundary:
"""A9: on an abort, every planned removal has already been applied."""
def _seed_recursive(self, tag):
source = os.path.join(TEST_DATA_DIR, f"dab_{tag}_src")
clean_dir(source)
# ``a/keep.bin`` sorts first, so the client streams it (and the proxy
# cuts) before the data pass ever reaches ``z/deep``.
_write(os.path.join(source, "a", "keep.bin"), b"B" * BIG_BYTES)
_write(os.path.join(source, "z", "deep", "keep.txt"), b"keep\n")
return source
@requires_rsync
def test_recursive_abort_removes_all_planned_extras(self):
# ---- FastSync: abort mid ``a/keep.bin``; ``z/deep`` is never reached.
source = self._seed_recursive("rec_fs")
dest = os.path.join(TEST_DATA_DIR, "dab_rec_fs_dst")
clean_dir(dest)
received = get_dest_received_dir(dest, source)
os.makedirs(os.path.join(received, "a"), exist_ok=True)
_write(os.path.join(received, "a", "a_extra"), b"stale\n")
os.makedirs(os.path.join(received, "z", "deep"), exist_ok=True)
_write(os.path.join(received, "z", "deep", "old_extra"), b"stale\n")
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
proxy = _SlicingProxy(server.port, forward_limit=MID_TRANSFER_BYTES,
throttle=PROXY_THROTTLE)
result, _ = run_client(source, dest, flags=["--delete-during"], port=proxy.port)
proxy.finish()
assert result.returncode != 0, "truncated transfer reported success"
assert not os.path.exists(os.path.join(received, "a", "a_extra"))
assert not os.path.exists(os.path.join(received, "z", "deep", "old_extra")), (
"FastSync left an extra in a directory it never reached before the abort"
)
# ---- rsync 3.4.1: same tree, same abort, same delete outcome.
source = self._seed_recursive("rec_rs")
rsync_dst = os.path.join(TEST_DATA_DIR, "dab_rec_rs_dst")
clean_dir(rsync_dst)
os.makedirs(os.path.join(rsync_dst, "a"), exist_ok=True)
_write(os.path.join(rsync_dst, "a", "a_extra"), b"stale\n")
os.makedirs(os.path.join(rsync_dst, "z", "deep"), exist_ok=True)
_write(os.path.join(rsync_dst, "z", "deep", "old_extra"), b"stale\n")
proc = _rsync_aborted(["-a", "--delete-during", f"--bwlimit={RSYNC_BWLIMIT}",
source + "/", rsync_dst + "/"])
assert proc.returncode != 0, "rsync was not actually interrupted"
assert not os.path.exists(os.path.join(rsync_dst, "a", "a_extra"))
assert not os.path.exists(os.path.join(rsync_dst, "z", "deep", "old_extra")), (
"rsync's generator did not delete ahead of its sender"
)
class TestDirsDeleteAbortBoundary:
"""A10: ``-d/--dirs`` uses per-directory plans like rsync.
The listed directory's direct extras are removed by the up-front plan while
a kept but untraversed subdirectory (and its destination content) is
shielded.
"""
def _seed_dirs(self, tag):
source = os.path.join(TEST_DATA_DIR, f"ddb_{tag}_src")
clean_dir(source)
_write(os.path.join(source, "big.bin"), b"B" * BIG_BYTES)
_write(os.path.join(source, "subdir", "keep.txt"), b"inner\n")
return source
@pytest.mark.parametrize("fs_flag,rs_flag", [("--delete-during", "--delete-during"),
("--delete", "--delete")])
@requires_rsync
def test_dirs_abort_removes_direct_extras_only(self, fs_flag, rs_flag):
label = f"{fs_flag.lstrip('-')}_{rs_flag.lstrip('-')}"
# ---- FastSync: ``-d`` lists the immediate children; big.bin streams and
# the abort lands mid-payload.
source = self._seed_dirs(f"dirs_{label}_fs")
dest = os.path.join(TEST_DATA_DIR, f"ddb_{label}_fs_dst")
clean_dir(dest)
received = get_dest_received_dir(dest, source)
_write(os.path.join(received, "old_extra"), b"stale\n")
os.makedirs(os.path.join(received, "subdir"), exist_ok=True)
_write(os.path.join(received, "subdir", "stale.txt"), b"stale inner\n")
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
proxy = _SlicingProxy(server.port, forward_limit=MID_TRANSFER_BYTES,
throttle=PROXY_THROTTLE)
result, _ = run_client(source + "/", dest, flags=["-d", fs_flag], port=proxy.port)
proxy.finish()
assert result.returncode != 0, f"{fs_flag}: truncated transfer reported success"
assert not os.path.exists(os.path.join(received, "old_extra")), (
f"{fs_flag}: the listed directory's direct extra survived the abort"
)
assert os.path.exists(os.path.join(received, "subdir", "stale.txt")), (
f"{fs_flag}: descended into a kept, untraversed subdirectory"
)
# ---- rsync 3.4.1: same shape and same abort.
source = self._seed_dirs(f"dirs_{label}_rs")
rsync_dst = os.path.join(TEST_DATA_DIR, f"ddb_{label}_rs_dst")
clean_dir(rsync_dst)
_write(os.path.join(rsync_dst, "old_extra"), b"stale\n")
os.makedirs(os.path.join(rsync_dst, "subdir"), exist_ok=True)
_write(os.path.join(rsync_dst, "subdir", "stale.txt"), b"stale inner\n")
proc = _rsync_aborted(["-d", rs_flag, f"--bwlimit={RSYNC_BWLIMIT}",
source + "/", rsync_dst + "/"])
assert proc.returncode != 0, "rsync was not actually interrupted"
assert not os.path.exists(os.path.join(rsync_dst, "old_extra")), (
f"rsync {rs_flag}: the listed directory's direct extra survived the abort"
)
assert os.path.exists(os.path.join(rsync_dst, "subdir", "stale.txt")), (
f"rsync {rs_flag}: descended into a kept, untraversed subdirectory"
)
def _tree(root):
out = []
for dirpath, dirs, files in os.walk(root):
for name in dirs:
out.append(os.path.relpath(os.path.join(dirpath, name), root))
for name in files:
out.append(os.path.relpath(os.path.join(dirpath, name), root))
return sorted(out)
class TestDirsDeleteFinalStateParity:
"""A10 completed run: ``-d DIR/ --delete`` (during default) and
``--delete-during`` match rsync's final tree, including a kept but
untraversed subdirectory whose destination content survives."""
@pytest.mark.parametrize("flag", ["--delete", "--delete-during"])
@requires_rsync
def test_dirs_final_state_matches_rsync(self, flag):
source = os.path.join(TEST_DATA_DIR, f"ddf_{flag.lstrip('-')}_src")
clean_dir(source)
_write(os.path.join(source, "keep.txt"), b"new keep\n")
_write(os.path.join(source, "subdir", "inner.txt"), b"inner\n")
def seed_dest(root):
clean_dir(root)
_write(os.path.join(root, "keep.txt"), b"old keep\n")
_write(os.path.join(root, "extra.txt"), b"extra\n")
_write(os.path.join(root, "extrasub", "ex.txt"), b"extra sub\n")
_write(os.path.join(root, "subdir", "stale.txt"), b"stale inner\n")
rsync_dst = os.path.join(TEST_DATA_DIR, f"ddf_{flag.lstrip('-')}_rs_dst")
seed_dest(rsync_dst)
env = dict(os.environ, LC_ALL="C")
rsync_result = subprocess.run(
[RSYNC, "-d", flag, source + "/", rsync_dst + "/"],
capture_output=True, text=True, env=env, timeout=120)
assert rsync_result.returncode == 0, rsync_result.stderr
rsync_tree = _tree(rsync_dst)
dest = os.path.join(TEST_DATA_DIR, f"ddf_{flag.lstrip('-')}_fs_dst")
clean_dir(dest)
received = get_dest_received_dir(dest, source)
seed_dest(received)
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
result, _ = run_client(source + "/", dest, flags=["-d", flag], port=server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
fastsync_tree = _tree(received)
assert fastsync_tree == rsync_tree, (
f"-d {flag}: fastsync tree {fastsync_tree} != rsync tree {rsync_tree}")
class TestOneFileSystemDeleteParity:
"""A9 side effect: the per-directory plan is now emitted only for directories
whose children were enumerated, so a ``-x`` mount-point directory that is
emitted but never traversed is shielded -- its destination content survives,
exactly as rsync keeps a non-descended mount point under ``--delete``."""
@requires_rsync
def test_mountpoint_content_survives_delete(self):
local = os.stat(".")
shm = "/dev/shm"
if not os.path.isdir(shm) or os.stat(shm).st_dev == local.st_dev:
pytest.skip("no cross-device filesystem available")
probe = os.path.join(shm, f"fastsync_dofs_{os.getpid()}")
clean_dir(probe)
_write(os.path.join(probe, "inside.txt"), b"cross\n")
try:
source = os.path.join(TEST_DATA_DIR, "dofs_src")
clean_dir(source)
_write(os.path.join(source, "keep.txt"), b"keep\n")
os.symlink(probe, os.path.join(source, "nested_link"))
def seed_dest(root):
clean_dir(root)
_write(os.path.join(root, "keep.txt"), b"old\n")
_write(os.path.join(root, "nested_link", "stale.txt"), b"stale\n")
rsync_dst = os.path.join(TEST_DATA_DIR, "dofs_rs_dst")
seed_dest(rsync_dst)
env = dict(os.environ, LC_ALL="C")
rsync_result = subprocess.run(
[RSYNC, "-a", "--copy-links", "-x", "--delete-during",
source + "/", rsync_dst + "/"],
capture_output=True, text=True, env=env, timeout=120)
assert rsync_result.returncode == 0, rsync_result.stderr
assert os.path.exists(os.path.join(rsync_dst, "nested_link", "stale.txt")), (
"rsync unexpectedly descended into the mount point")
dest = os.path.join(TEST_DATA_DIR, "dofs_fs_dst")
clean_dir(dest)
received = get_dest_received_dir(dest, source)
seed_dest(received)
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
result, _ = run_client(source, dest,
flags=["-a", "--copy-links", "-x", "--delete-during"],
port=server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
assert os.path.exists(os.path.join(received, "nested_link", "stale.txt")), (
"FastSync descended into a non-traversed mount point under --delete")
finally:
clean_dir(probe)
+40
View File
@@ -662,6 +662,46 @@ class TestWireStatsParity:
assert re.match(r"Number of created files: 1 \(reg: 1\)$", r_created), r_created
assert f_created == r_created, (r_created, f_created)
@requires_rsync
@pytest.mark.ci
@pytest.mark.parametrize("mt", [False, True])
def test_stats_r_directory_breakdown_matches_rsync(self, shared_server, mt):
"""A recursive `-r` scan (no -t/-p) exposes no directory metadata, but
rsync still counts every directory in `Number of files`; the sender's
lightweight directory counter must reproduce the `dir: N` category."""
source = os.path.join(TEST_DATA_DIR, "wire_stdir_src")
dest = os.path.join(TEST_DATA_DIR, "wire_stdir_dst")
rdst = os.path.join(TEST_DATA_DIR, "wire_stdir_rdst")
clean_dir(source)
clean_dir(dest)
clean_dir(rdst)
os.makedirs(os.path.join(source, "sub", "deep"))
os.makedirs(os.path.join(source, "empty"))
for rel in ("a.txt", os.path.join("sub", "b.txt"), os.path.join("sub", "deep", "c.txt")):
with open(os.path.join(source, rel), "wb") as fh:
fh.write(b"x\n")
os.makedirs(get_dest_received_dir(dest, source), exist_ok=True)
rsync_result = _rsync(["-r", "--stats", source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
flags = ["-r", "--stats"] + (["--threads"] if mt else [])
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, result.stderr[:300]
def stats_line(text, key):
for line in text.splitlines():
if line.startswith(key + ":"):
return line
return None
r_files = stats_line(rsync_result.stdout, "Number of files")
f_files = stats_line(result.stdout, "Number of files")
# 3 regular files, 4 directories (root, sub, sub/deep, empty).
assert re.match(r"Number of files: 7 \(reg: 3, dir: 4\)$", r_files), r_files
assert f_files == r_files, (r_files, f_files)
assert (stats_line(result.stdout, "Number of regular files transferred") ==
stats_line(rsync_result.stdout, "Number of regular files transferred"))
@requires_rsync
@pytest.mark.ci
@pytest.mark.parametrize("mt", [False, True])
@@ -0,0 +1,198 @@
"""Differential rsync-parity coverage for two residuals closed on this branch.
* A4 -- ``--compare-dest``/``--copy-dest``/``--link-dest`` relative-DIR
resolution: rsync resolves a relative DIR against the destination directory
and appends the file's TRANSFER-RELATIVE name. FastSync's default transfer
mirrors the absolute source path below its receive root, so a naive relative
DIR used to probe a different tree. These tests seed the basis at rsync's
spelling and assert FastSync finds it (byte-exact / hard-linked / sparse),
matching real rsync 3.4.1.
* A5 -- ``-y``/``--fuzzy`` candidate eligibility: rsync's ``find_fuzzy`` has no
delta-size gate, so it reuses an oversized (>10x) or sub-16-KiB sibling;
FastSync used to decline both. These tests assert FastSync now uses the same
sibling as rsync (observable as ``Matched data``) with a byte-exact result.
Every test skips cleanly when rsync is absent.
"""
import os
import shutil
import subprocess
import sys
import pytest
sys.path.insert(0, os.path.dirname(__file__))
from common import ( # noqa: E402
TEST_DATA_DIR,
clean_dir,
get_dest_received_dir,
run_client,
)
RSYNC = shutil.which("rsync")
requires_rsync = pytest.mark.skipif(RSYNC is None, reason="rsync 3.4.1 not installed")
OLD_MTIME = 1_500_000_000
def _write(path, content, mtime=None):
os.makedirs(os.path.dirname(path), exist_ok=True)
with open(path, "wb") as fh:
fh.write(content)
if mtime is not None:
os.utime(path, (mtime, mtime))
def _read(path):
with open(path, "rb") as fh:
return fh.read()
def _rsync(args):
env = dict(os.environ, LC_ALL="C")
return subprocess.run([RSYNC] + args, capture_output=True, text=True, env=env, timeout=120)
def _stat_bytes(text, label):
for line in text.splitlines():
if line.startswith(label + ":"):
return int(line.split(":", 1)[1].strip().split()[0].replace(",", ""))
return None
class TestRelativeBasisDirResolution:
"""A4: a relative basis DIR must resolve to the same tree as rsync's."""
_FILES = {
"root.txt": b"root-basis-content\n",
"sub/nested.txt": b"nested-basis-content\n",
}
def _seed_source(self, source):
clean_dir(source)
for rel, data in self._FILES.items():
_write(os.path.join(source, rel), data, OLD_MTIME)
return self._FILES
@requires_rsync
@pytest.mark.parametrize("flag", ["--compare-dest", "--link-dest"])
def test_relative_dir_resolves_like_rsync(self, shared_server, flag):
tag = flag.lstrip("-")
source = os.path.join(TEST_DATA_DIR, f"relbasis_{tag}_src")
rdst = os.path.join(TEST_DATA_DIR, f"relbasis_{tag}_rdst")
fdst = os.path.join(TEST_DATA_DIR, f"relbasis_{tag}_fdst")
self._seed_source(source)
# rsync: relative DIR -> dest/basis/<transfer-relative name>.
clean_dir(rdst)
for rel, data in self._FILES.items():
_write(os.path.join(rdst, "basis", rel), data, OLD_MTIME)
rs = _rsync(["-a", f"{flag}=basis", source + "/", rdst + "/"])
assert rs.returncode == 0, rs.stderr
# FastSync: the SAME relative spelling seeded at the SAME
# transfer-relative location under its destination root.
clean_dir(fdst)
for rel, data in self._FILES.items():
_write(os.path.join(fdst, "basis", rel), data, OLD_MTIME)
result, _ = run_client(source, fdst,
flags=["-a", f"{flag}=basis", "--incremental"],
port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
received = get_dest_received_dir(fdst, source)
for rel, data in self._FILES.items():
rfile = os.path.join(rdst, rel)
ffile = os.path.join(received, rel)
basis = os.path.join(fdst, "basis", rel)
if flag == "--compare-dest":
# compare-dest never copies: both destinations stay sparse.
assert not os.path.exists(rfile), f"rsync copied {rel}"
assert not os.path.exists(ffile), (
f"FastSync did not resolve the relative basis DIR at {basis!r} "
f"(expected {rel!r} to stay sparse like rsync)")
else:
# link-dest hard-links; a basis miss would transfer a new file.
assert os.path.exists(ffile), f"FastSync lost {rel}"
assert _read(ffile) == data
assert os.stat(ffile).st_ino == os.stat(basis).st_ino, (
f"FastSync did not hard-link {rel!r} to the relative basis at "
f"{basis!r} (basis not resolved like rsync)")
@requires_rsync
def test_relative_dir_copy_dest_content(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "relbasis_copy_src")
fdst = os.path.join(TEST_DATA_DIR, "relbasis_copy_fdst")
self._seed_source(source)
clean_dir(fdst)
for rel, data in self._FILES.items():
_write(os.path.join(fdst, "basis", rel), data, OLD_MTIME)
result, _ = run_client(source, fdst,
flags=["-a", "--copy-dest=basis", "--incremental"],
port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
received = get_dest_received_dir(fdst, source)
for rel, data in self._FILES.items():
ffile = os.path.join(received, rel)
assert os.path.exists(ffile), f"copy-dest did not materialize {rel}"
assert _read(ffile) == data
assert os.stat(ffile).st_ino != os.stat(os.path.join(fdst, "basis", rel)).st_ino
class TestFuzzyEligibilityWindow:
"""A5: --fuzzy candidate eligibility must match rsync's uncapped window."""
BASE = b"the quick brown fox jumps over the lazy dog\n" * 4000
def _run_pair(self, shared_server, tag, payload, sibling):
source = os.path.join(TEST_DATA_DIR, f"fzw_{tag}_src")
dest = os.path.join(TEST_DATA_DIR, f"fzw_{tag}_dst")
rdst = os.path.join(TEST_DATA_DIR, f"fzw_{tag}_rdst")
clean_dir(source)
clean_dir(dest)
clean_dir(rdst)
_write(os.path.join(source, "report_v2.txt"), payload)
for root in (rdst, get_dest_received_dir(dest, source)):
_write(os.path.join(root, "report_v1.txt"), sibling)
rs = _rsync(["-a", "--no-whole-file", "--fuzzy", "--stats",
source + "/", rdst + "/"])
assert rs.returncode == 0, rs.stderr
result, _ = run_client(
source, dest,
flags=["-a", "--incremental", "--delta", "--fuzzy", "--stats"],
port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
# The reconstructed file is byte-exact in every case.
assert _read(os.path.join(get_dest_received_dir(dest, source),
"report_v2.txt")) == payload
return rs, result
@requires_rsync
def test_oversized_sibling_eligible_like_rsync(self, shared_server):
"""A sibling 20x the source is used by rsync; FastSync must too (its old
10x delta-size gate declined it)."""
n = 65536
payload = (self.BASE * ((n // len(self.BASE)) + 1))[:n]
sibling = (self.BASE * 200)[: n * 20]
rs, result = self._run_pair(shared_server, "big", payload, sibling)
assert _stat_bytes(rs.stdout, "Matched data") > 0, \
"rsync should use a >10x fuzzy basis"
assert _stat_bytes(result.stdout, "Matched data") > 0, (
"FastSync's fuzzy eligibility must accept a >10x sibling like rsync "
f"(Matched data={_stat_bytes(result.stdout, 'Matched data')})")
@requires_rsync
def test_small_source_sibling_eligible_like_rsync(self, shared_server):
"""A sub-16-KiB source with an identical sibling is used by rsync;
FastSync's old 16 KiB delta minimum declined it."""
n = 8192
payload = (self.BASE * ((n // len(self.BASE)) + 1))[:n]
rs, result = self._run_pair(shared_server, "small", payload, payload)
assert _stat_bytes(rs.stdout, "Matched data") > 0, \
"rsync applies --fuzzy below 16 KiB"
assert _stat_bytes(result.stdout, "Matched data") > 0, (
"FastSync's fuzzy eligibility must accept a sub-16-KiB source like "
f"rsync (Matched data={_stat_bytes(result.stdout, 'Matched data')})")
+105
View File
@@ -0,0 +1,105 @@
"""`--debug=FLAGS` natural-event categories (no-wire).
FastSync maps the rsync `--debug` categories that correspond to a real event it
already performs (``flist``, ``del``, ``hash``/``deltasum``, ``recv``,
``filter`` and ``send``) onto debug output. A normal run prints none of it.
"""
import os
import sys
import pytest
sys.path.insert(0, os.path.dirname(__file__))
from common import TEST_DATA_DIR, run_client, clean_dir, get_dest_received_dir, ServerManager
def _make_tree(root):
clean_dir(root)
os.makedirs(os.path.join(root, "sub"))
with open(os.path.join(root, "a.txt"), "wb") as fh:
fh.write(b"alpha\n")
with open(os.path.join(root, "keep.log"), "wb") as fh:
fh.write(b"log\n")
with open(os.path.join(root, "sub", "b.txt"), "wb") as fh:
fh.write(b"beta\n")
@pytest.mark.ci
def test_debug_flist_and_send_emit_output(shared_server):
"""`--debug=flist,send` produces category-tagged debug output."""
source = os.path.join(TEST_DATA_DIR, "dbg_src")
dest = os.path.join(TEST_DATA_DIR, "dbg_dst")
_make_tree(source)
clean_dir(dest)
result, _ = run_client(source, dest, flags=["-a", "--debug=flist,send"],
port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
assert "flist: scanning" in result.stdout, result.stdout
assert "send: " in result.stdout, result.stdout
@pytest.mark.ci
def test_debug_filter_emits_excluded_entry(shared_server):
source = os.path.join(TEST_DATA_DIR, "dbg_filter_src")
dest = os.path.join(TEST_DATA_DIR, "dbg_filter_dst")
_make_tree(source)
clean_dir(dest)
result, _ = run_client(source, dest,
flags=["-a", "--debug=filter", "--exclude=*.log"],
port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
assert "filter: excluded keep.log" in result.stdout, result.stdout
@pytest.mark.ci
def test_debug_hash_and_recv_emit_on_incremental(shared_server):
source = os.path.join(TEST_DATA_DIR, "dbg_hash_src")
dest = os.path.join(TEST_DATA_DIR, "dbg_hash_dst")
_make_tree(source)
clean_dir(dest)
result, _ = run_client(source, dest,
flags=["-a", "--incremental", "--checksum",
"--debug=hash,recv"],
port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
assert "hash: " in result.stdout, result.stdout
assert "recv: " in result.stdout, result.stdout
@pytest.mark.ci
def test_debug_del_emits_deleted_path():
"""`--debug=del` reports the paths the receiver actually removed.
A deletion-capable server is required (the shared fixture refuses
client-requested deletion)."""
source = os.path.join(TEST_DATA_DIR, "dbg_del_src")
dest = os.path.join(TEST_DATA_DIR, "dbg_del_dst")
_make_tree(source)
clean_dir(dest)
seeded = get_dest_received_dir(dest, source)
os.makedirs(seeded)
with open(os.path.join(seeded, "extra.tmp"), "wb") as fh:
fh.write(b"stale\n")
server = ServerManager()
server.start(extra_args=["--allow-super", "--allow-delete"])
try:
result, _ = run_client(source, dest, flags=["-a", "--delete", "--debug=del"],
port=server.port)
finally:
server.stop()
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
assert "del: " in result.stdout and "extra.tmp" in result.stdout, result.stdout
assert not os.path.exists(os.path.join(seeded, "extra.tmp"))
@pytest.mark.ci
def test_normal_run_has_no_debug_output(shared_server):
source = os.path.join(TEST_DATA_DIR, "dbg_quiet_src")
dest = os.path.join(TEST_DATA_DIR, "dbg_quiet_dst")
_make_tree(source)
clean_dir(dest)
result, _ = run_client(source, dest, flags=["-a"], port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
assert "[DEBUG]" not in result.stdout
assert "flist: scanning" not in result.stdout
assert "send: " not in result.stdout
@@ -0,0 +1,174 @@
"""Differential parity for `--info=mount` and `--info=stats` (no-wire).
Both behaviours are compared against real rsync 3.4.1:
* `--info=mount` prints rsync's ``[sender] skipping mount-point dir NAME`` line
when ``-xx`` drops a mount-point directory. Plain ``-x`` keeps the empty
directory and stays silent, exactly like rsync.
* `--info=stats` requests the same transfer-statistics block as `--stats`
(rsync spells the full block ``--info=stats2``/``--stats``).
The tests are skipped when rsync is unavailable.
"""
import os
import re
import shutil
import subprocess
import sys
import pytest
sys.path.insert(0, os.path.dirname(__file__))
from common import TEST_DATA_DIR, run_client, clean_dir, get_dest_received_dir
RSYNC = shutil.which("rsync")
requires_rsync = pytest.mark.skipif(RSYNC is None, reason="rsync 3.4.1 not installed")
def _rsync(args):
env = dict(os.environ, LC_ALL="C")
return subprocess.run([RSYNC] + args, capture_output=True, text=True, env=env, timeout=120)
def _cross_device_mount_tree(source):
"""Build a source whose ``nested_link`` is a symlink onto a tmpfs directory.
``--copy-links`` dereferences it so ``-x`` sees a mount-point directory on a
different device. Returns the probe path to remove, or skips the test when
no cross-device filesystem is available.
"""
local = os.stat(".")
shm = "/dev/shm"
try:
shm_stat = os.stat(shm)
except OSError:
pytest.skip("/dev/shm not available")
if shm_stat.st_dev == local.st_dev:
pytest.skip("no cross-device filesystem available")
clean_dir(source)
with open(os.path.join(source, "keep.txt"), "wb") as fh:
fh.write(b"keep\n")
probe = os.path.join(shm, f"fastsync_info_mount_{os.getpid()}")
shutil.rmtree(probe, ignore_errors=True)
os.makedirs(probe)
with open(os.path.join(probe, "inside.txt"), "wb") as fh:
fh.write(b"cross\n")
try:
os.symlink(probe, os.path.join(source, "nested_link"))
except OSError:
shutil.rmtree(probe, ignore_errors=True)
pytest.skip("cannot create symlink")
return probe
@requires_rsync
@pytest.mark.ci
def test_info_mount_xx_matches_rsync(shared_server):
"""`-xx --info=mount` drops the mount-point dir and prints rsync's line."""
source = os.path.join(TEST_DATA_DIR, "info_mount_src")
dest = os.path.join(TEST_DATA_DIR, "info_mount_dst")
rdst = os.path.join(TEST_DATA_DIR, "info_mount_rdst")
probe = _cross_device_mount_tree(source)
clean_dir(dest)
clean_dir(rdst)
flags = ["-a", "--copy-links", "-xx", "--info=mount"]
try:
rsync_result = _rsync(flags + [source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
expected = "[sender] skipping mount-point dir nested_link"
assert expected in rsync_result.stdout, rsync_result.stdout
assert expected in result.stdout, (result.stdout, result.stderr)
received = get_dest_received_dir(dest, source)
assert os.path.exists(os.path.join(received, "keep.txt"))
# -xx omits the mount-point directory entirely.
assert not os.path.exists(os.path.join(received, "nested_link"))
assert not os.path.exists(os.path.join(rdst, "nested_link"))
finally:
shutil.rmtree(probe, ignore_errors=True)
@requires_rsync
@pytest.mark.ci
def test_info_mount_single_x_is_silent(shared_server):
"""Plain `-x` keeps the empty mount-point directory and prints no line."""
source = os.path.join(TEST_DATA_DIR, "info_mount1_src")
dest = os.path.join(TEST_DATA_DIR, "info_mount1_dst")
rdst = os.path.join(TEST_DATA_DIR, "info_mount1_rdst")
probe = _cross_device_mount_tree(source)
clean_dir(dest)
clean_dir(rdst)
flags = ["-a", "--copy-links", "-x", "--info=mount"]
try:
rsync_result = _rsync(flags + [source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
assert "skipping mount-point dir" not in rsync_result.stdout
assert "skipping mount-point dir" not in result.stdout
received = get_dest_received_dir(dest, source)
assert os.path.isdir(os.path.join(received, "nested_link"))
assert not os.path.exists(os.path.join(received, "nested_link", "inside.txt"))
assert os.path.isdir(os.path.join(rdst, "nested_link"))
assert not os.path.exists(os.path.join(rdst, "nested_link", "inside.txt"))
finally:
shutil.rmtree(probe, ignore_errors=True)
def _make_stats_tree(root):
clean_dir(root)
os.makedirs(os.path.join(root, "sub"))
with open(os.path.join(root, "a.txt"), "wb") as fh:
fh.write(b"alpha\n")
with open(os.path.join(root, "sub", "b.txt"), "wb") as fh:
fh.write(b"beta\n")
def _pick_stats(text):
keys = ("Number of files", "Number of regular files transferred", "Total file size",
"Total transferred file size", "Literal data", "Matched data")
out = {}
for line in text.splitlines():
for key in keys:
if line.startswith(key + ":"):
out[key] = line
return out
@requires_rsync
@pytest.mark.ci
def test_info_stats_emits_full_stats_block(shared_server):
"""`--info=stats` is the same full block as `--stats` and matches rsync."""
source = os.path.join(TEST_DATA_DIR, "info_stats_src")
dest = os.path.join(TEST_DATA_DIR, "info_stats_dst")
rdst = os.path.join(TEST_DATA_DIR, "info_stats_rdst")
dest2 = os.path.join(TEST_DATA_DIR, "info_stats_dst2")
_make_stats_tree(source)
for path in (dest, rdst, dest2):
clean_dir(path)
os.makedirs(get_dest_received_dir(path, source), exist_ok=True)
rsync_result = _rsync(["-a", "--stats", source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
info_result, _ = run_client(source, dest, flags=["-a", "--info=stats"],
port=shared_server.port)
assert info_result.returncode == 0, (info_result.stderr or info_result.stdout)[:300]
stats_result, _ = run_client(source, dest2, flags=["-a", "--stats"],
port=shared_server.port)
assert stats_result.returncode == 0, (stats_result.stderr or stats_result.stdout)[:300]
# --info=stats must print the same block as --stats...
assert _pick_stats(info_result.stdout) == _pick_stats(stats_result.stdout), (
f"info={info_result.stdout} stats={stats_result.stdout}")
# ...and the protocol-independent counters must match real rsync.
assert _pick_stats(info_result.stdout) == _pick_stats(rsync_result.stdout), (
f"rsync={_pick_stats(rsync_result.stdout)} fastsync={_pick_stats(info_result.stdout)}")
assert re.search(r"^Number of files: \d+ \(reg: 2, dir: 2\)$", info_result.stdout,
re.MULTILINE), info_result.stdout
+185
View File
@@ -0,0 +1,185 @@
"""Differential rsync-parity coverage for FastSync's transfer/delete ORDER.
rsync walks a source tree in its sorted flist order: within each directory the
non-directories come first (ascending name), then the subdirectories (ascending
name), each subdirectory immediately followed by its own subtree (depth-first).
The sequential scanner now reproduces that order, which makes both the
``--info=name`` stream and the ``--delete-during`` deletion sequence match real
``rsync 3.4.1`` exactly. ``--threads`` has no rsync analogue and is unordered.
Every test skips cleanly when rsync is absent.
"""
import os
import shutil
import subprocess
import sys
import pytest
sys.path.insert(0, os.path.dirname(__file__))
from common import ( # noqa: E402
TEST_DATA_DIR,
ServerManager,
clean_dir,
get_dest_received_dir,
run_client,
)
RSYNC = shutil.which("rsync")
requires_rsync = pytest.mark.skipif(RSYNC is None, reason="rsync 3.4.1 not installed")
MTIME = 1_500_000_000
_TREE = {
"a.txt": b"a\n",
"b.txt": b"b\n",
"z.txt": b"z\n",
"a_dir/f.txt": b"f\n",
"a_dir/deep/g.txt": b"g\n",
"m_dir/h.txt": b"h\n",
"Z_dir/i.txt": b"i\n",
}
def _write(path, data):
os.makedirs(os.path.dirname(path), exist_ok=True)
with open(path, "wb") as fh:
fh.write(data)
os.utime(path, (MTIME, MTIME))
def _rsync(args):
env = dict(os.environ, LC_ALL="C")
return subprocess.run([RSYNC] + args, capture_output=True, text=True, env=env, timeout=120)
def _deleting(text):
out = []
for line in text.splitlines():
stripped = line.strip()
if stripped.startswith("*deleting") or stripped.startswith("deleting"):
out.append(stripped.split()[-1])
return out
class TestTransferOrderParity:
@requires_rsync
def test_info_name_file_order_matches_rsync(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "order_name_src")
clean_dir(source)
for rel, data in _TREE.items():
_write(os.path.join(source, rel), data)
rdst = os.path.join(TEST_DATA_DIR, "order_name_rdst")
clean_dir(rdst)
rs = _rsync(["-a", "--info=name", source + "/", rdst + "/"])
assert rs.returncode == 0, rs.stderr
# rsync also names the directories (trailing '/'); FastSync names the
# transferred entries. Compare the file/symlink sequence, which is what
# the traversal order determines.
rsync_files = [l for l in rs.stdout.splitlines() if l.strip() and not l.endswith("/")]
fdst = os.path.join(TEST_DATA_DIR, "order_name_fdst")
clean_dir(fdst)
result, _ = run_client(source, fdst, flags=["-a", "--info=name"],
port=shared_server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
fsync_files = [
l for l in result.stdout.splitlines()
if l.strip() and l.strip() != "./" and not l.startswith("sending")
]
assert fsync_files == rsync_files, (
f"transfer order differs\nrsync={rsync_files}\nfastsync={fsync_files}")
class TestDeleteOrderParity:
_EXTRA = {
"a_extra.txt": b"a\n",
"z_extra.txt": b"z\n",
"a_extra_dir/f": b"f\n",
"z_extra_dir/f": b"f\n",
"a_extra_dir/sub/g": b"g\n",
}
def _seed_source(self):
source = os.path.join(TEST_DATA_DIR, "order_del_src")
clean_dir(source)
_write(os.path.join(source, "keep.txt"), b"k\n")
_write(os.path.join(source, "keepdir", "x.txt"), b"x\n")
_write(os.path.join(source, "keep2", "y.txt"), b"y\n")
return source
def _assert_order(self, timing, dry_run=False):
source = self._seed_source()
rdst = os.path.join(TEST_DATA_DIR, f"order_{timing}_rdst")
clean_dir(rdst)
for rel, data in self._EXTRA.items():
_write(os.path.join(rdst, rel), data)
rs_flags = ["-a", "-n"] if dry_run else ["-a"]
rs = _rsync(rs_flags + [timing, "--info=del", source + "/", rdst + "/"])
assert rs.returncode == 0, rs.stderr
fdst = os.path.join(TEST_DATA_DIR, f"order_{timing}_fdst")
clean_dir(fdst)
received = get_dest_received_dir(fdst, source)
for rel, data in self._EXTRA.items():
_write(os.path.join(received, rel), data)
fs_flags = ["-a", "-n"] if dry_run else ["-a"]
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
result, _ = run_client(source, fdst, flags=fs_flags + [timing, "--info=del"],
port=server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
rsync_order = _deleting(rs.stdout)
fsync_order = _deleting(result.stdout)
assert sorted(fsync_order) == sorted(rsync_order), (
f"{timing} deleted set differs\nrsync={rsync_order}\nfastsync={fsync_order}")
assert fsync_order == rsync_order, (
f"{timing} deletion order differs\nrsync={rsync_order}\nfastsync={fsync_order}")
@requires_rsync
def test_delete_during_deletion_order_matches_rsync(self):
self._assert_order("--delete-during")
@requires_rsync
def test_delete_delay_deletion_order_matches_rsync(self):
self._assert_order("--delete-delay")
@requires_rsync
def test_dry_run_delete_order_matches_rsync(self):
self._assert_order("--delete", dry_run=True)
@requires_rsync
def test_partial_max_delete_survivor_order_matches_rsync(self):
"""With the exact removal order matching rsync, a --max-delete cap stops
after the same entries, so the survivor set is identical too."""
source = os.path.join(TEST_DATA_DIR, "order_maxdel_src")
clean_dir(source)
_write(os.path.join(source, "keep.txt"), b"k\n")
extra = {f"e{i}.txt": b"x\n" for i in range(6)}
extra["ed/f"] = b"f\n"
extra["ed/g"] = b"g\n"
rdst = os.path.join(TEST_DATA_DIR, "order_maxdel_rdst")
clean_dir(rdst)
for rel, data in extra.items():
_write(os.path.join(rdst, rel), data)
rs = _rsync(["-a", "--delete-during", "--max-delete=3", "--info=del",
source + "/", rdst + "/"])
assert rs.returncode in (0, 25), (rs.returncode, rs.stderr)
fdst = os.path.join(TEST_DATA_DIR, "order_maxdel_fdst")
clean_dir(fdst)
received = get_dest_received_dir(fdst, source)
for rel, data in extra.items():
_write(os.path.join(received, rel), data)
with ServerManager() as server:
server.start(extra_args=["--allow-delete"])
result, _ = run_client(source, fdst,
flags=["-a", "--delete-during", "--max-delete=3", "--info=del"],
port=server.port)
assert result.returncode in (0, 25), (result.returncode, result.stderr[:300])
assert _deleting(result.stdout) == _deleting(rs.stdout), (
f"partial --max-delete survivor order differs\n"
f"rsync={_deleting(rs.stdout)}\nfastsync={_deleting(result.stdout)}")
+20 -18
View File
@@ -861,11 +861,11 @@ class TestFuzzy:
byte-exact result. FastSync ports rsync 3.4.1's weighted-Levenshtein name
heuristic, so where both delta engines admit the candidate the tools pick
the same basis (the ``fuzzy_basis`` differential asserts the tree and the
Matched/Literal counters match with the block size pinned). The residual is
candidate ELIGIBILITY: FastSync's delta size gate (both files >= 16 KiB and
a <= 10x size ratio) is narrower than rsync's, which empirically uses a
fuzzy basis well beyond 10x and below 16 KiB. These tests pin the window
boundary and prove the byte-exact fallback on both sides of it."""
Matched/Literal counters match with the block size pinned). Candidate
ELIGIBILITY is now rsync's too: the fuzzy search no longer inherits the
ordinary delta engine's 16 KiB minimum or 10x size-ratio bound, so an
oversized or sub-16-KiB sibling is reused exactly as rsync reuses it.
These tests pin that window on both sides."""
_BASE = b"the quick brown fox jumps over the lazy dog\n" * 4000
@@ -900,9 +900,10 @@ class TestFuzzy:
return rs, result
@requires_rsync
def test_fuzzy_above_size_window_declines_but_tree_exact(self, shared_server):
"""A sibling >10x the source is used by rsync but declined by FastSync's
delta size-ratio gate; both destinations stay byte-identical."""
def test_fuzzy_above_size_window_matches_rsync(self, shared_server):
"""A sibling >10x the source is used by rsync and by FastSync: fuzzy
eligibility is rsync's, not the ordinary delta size-ratio gate; both
destinations stay byte-identical and both reuse the basis."""
n = 65536
payload = (self._BASE * ((n // len(self._BASE)) + 1))[:n]
sibling = (self._BASE * 200)[: n * 20]
@@ -911,15 +912,16 @@ class TestFuzzy:
rs, result = self._run_both(shared_server, source, dest, rdst,
payload, sibling)
assert _stat_bytes(rs.stdout, "Matched data") > 0, \
"rsync should still use a >10x fuzzy basis"
assert _stat_bytes(result.stdout, "Matched data") == 0, \
"FastSync's 10x delta size-ratio gate must decline the oversized basis"
assert _stat_bytes(result.stdout, "Literal data") == n
"rsync should use a >10x fuzzy basis"
assert _stat_bytes(result.stdout, "Matched data") > 0, \
"FastSync must accept a >10x fuzzy basis like rsync"
assert _stat_bytes(result.stdout, "Literal data") < n
@requires_rsync
def test_fuzzy_below_delta_minimum_declines_but_tree_exact(self, shared_server):
"""A sibling below the 16 KiB delta minimum is used by rsync but never
enters FastSync's delta/fuzzy path; both trees stay byte-identical."""
def test_fuzzy_below_delta_minimum_matches_rsync(self, shared_server):
"""A sibling below the 16 KiB delta minimum is used by rsync and by
FastSync: fuzzy eligibility no longer inherits the delta engine's
minimum; both trees stay byte-identical and both reuse the basis."""
n = 8192
payload = (self._BASE * ((n // len(self._BASE)) + 1))[:n]
source, dest, rdst = (self._src("small"), self._dst("small"),
@@ -928,9 +930,9 @@ class TestFuzzy:
payload, payload)
assert _stat_bytes(rs.stdout, "Matched data") > 0, \
"rsync applies --fuzzy below 16 KiB"
assert _stat_bytes(result.stdout, "Matched data") == 0, \
"FastSync's 16 KiB delta minimum must bypass the fuzzy basis"
assert _stat_bytes(result.stdout, "Literal data") == n
assert _stat_bytes(result.stdout, "Matched data") > 0, \
"FastSync must apply --fuzzy below 16 KiB like rsync"
assert _stat_bytes(result.stdout, "Literal data") < n
class TestIgnoreExistingShortCircuit:
+9 -8
View File
@@ -762,8 +762,9 @@ static void test_parse_args_debug_flags() {
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0);
EXPECT_EQ_INT(cfg->debug_level, LOG_DEBUG_ALL);
EXPECT_EQ_INT(get_log_debug_flags(), LOG_DEBUG_ALL);
EXPECT_EQ_INT(cfg->debug_level, LOG_DEBUG_IO | LOG_DEBUG_PROTO | LOG_DEBUG_PACK | LOG_DEBUG_UTIL);
EXPECT_EQ_INT(get_log_debug_flags(),
LOG_DEBUG_IO | LOG_DEBUG_PROTO | LOG_DEBUG_PACK | LOG_DEBUG_UTIL);
config_delete(cfg);
}
@@ -1421,10 +1422,9 @@ static void test_parse_args_info_name_and_help() {
config_delete(cfg);
}
/* rsync 3.4.1's full --info/--debug vocabulary parses. The info categories
* with a FastSync event set their flag; the remaining rsync-only categories
* (backup/mount/symsafe/syms) parse but stay silent. Every --debug category
* listed here is FastSync-silent, so debug_level stays 0. */
/* rsync 3.4.1's full --info/--debug vocabulary parses. The categories with a
* FastSync event set their flag; the remaining rsync-only categories
* (backup/symsafe/syms, acl/bind/chdir/...) parse but stay silent. */
static void test_parse_args_rsync_flag_vocabulary_accepted() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--info=backup,del,flist,mount,nonreg,progress,remove,symsafe,syms",
@@ -1436,9 +1436,10 @@ static void test_parse_args_rsync_flag_vocabulary_accepted() {
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0);
EXPECT_EQ_INT(cfg->info_level, LOG_INFO_DEL | LOG_INFO_FLIST | LOG_INFO_NONREG |
EXPECT_EQ_INT(cfg->info_level, LOG_INFO_DEL | LOG_INFO_FLIST | LOG_INFO_MOUNT | LOG_INFO_NONREG |
LOG_INFO_PROGRESS | LOG_INFO_REMOVE);
EXPECT_EQ_INT(cfg->debug_level, 0);
EXPECT_EQ_INT(cfg->debug_level, LOG_DEBUG_DEL | LOG_DEBUG_FLIST | LOG_DEBUG_HASH |
LOG_DEBUG_RECV | LOG_DEBUG_FILTER | LOG_DEBUG_SEND);
config_delete(cfg);
}