Merge branch 'fix/parity-tests' into feat/parity-completion

This commit is contained in:
2026-09-17 01:38:19 +02:00
10 changed files with 485 additions and 92 deletions
+9 -1
View File
@@ -590,7 +590,15 @@ void change_emit_file_sent_bytes(const Config* config, const File* file,
event.bytes_sent = 0; event.bytes_sent = 0;
} else { } else {
event.bytes_sent = bytes_sent; event.bytes_sent = bytes_sent;
event.bytes_read = bytes_read; /* rsync's %c is the block-checksum bytes received for the file. Even a
* whole-file transfer (no basis; --append/--inplace included) receives
* rsync's 16-byte sum header, so rsync reports 16; a dry run transfers
* nothing and reports 0. FastSync's whole-file path has no sum header, so
* report rsync's value for parity. With delta enabled the real received
* bytes are kept, but FastSync's signature framing differs from rsync's so
* those stay numerically divergent. */
bool delta_active = config->use_delta && !config->whole_file;
event.bytes_read = (!config->dry_run && !delta_active) ? 16 : bytes_read;
} }
char* name = NULL; char* name = NULL;
char* path = NULL; char* path = NULL;
+5 -3
View File
@@ -69,7 +69,8 @@ char* change_render_itemize_code(const Config* config, const ChangeEvent* event)
/* Expand an --out-format/--log-file-format template. Supported tokens: /* Expand an --out-format/--log-file-format template. Supported tokens:
* %i itemize code %n transfer-relative name (dir: trailing /) * %i itemize code %n transfer-relative name (dir: trailing /)
* %f long display path %l file length in bytes * %f long display path %l file length in bytes
* %b wire bytes transferred %c wire bytes read back for the file * %b wire bytes transferred %c block-checksum bytes received (rsync: 16
* for a whole-file transfer, 0 for a dry run)
* %C whole-file checksum hex (xxh128 by default; spaces for non-regular) * %C whole-file checksum hex (xxh128 by default; spaces for non-regular)
* %M mtime (YYYY/MM/DD-HH:MM:SS) * %M mtime (YYYY/MM/DD-HH:MM:SS)
* %t current time %o operation ("send"/"del.") * %t current time %o operation ("send"/"del.")
@@ -91,8 +92,9 @@ char* change_render_list_line(const Config* config, const ChangeEvent* event);
void change_emit(const Config* config, const ChangeEvent* event); void change_emit(const Config* config, const ChangeEvent* event);
/* Build and emit a CHANGE_SENT event for a file the client just sent. `bytes_sent` /* Build and emit a CHANGE_SENT event for a file the client just sent. `bytes_sent`
* / `bytes_read` are the process-wide wire-byte deltas for this file (rsync's * is the process-wide wire-byte delta for this file (rsync's %b) and `bytes_read`
* %b / %c); pass 0 when unknown. */ * the received bytes used for the delta handshake; pass 0 when unknown. For a
* whole-file transfer %c is pinned to rsync's 16-byte sum header regardless. */
void change_emit_file_sent_bytes(const Config* config, const File* file, void change_emit_file_sent_bytes(const Config* config, const File* file,
unsigned long long bytes_sent, unsigned long long bytes_read); unsigned long long bytes_sent, unsigned long long bytes_read);
+34
View File
@@ -480,6 +480,36 @@ static bool split_flag_level(const char* token, char* name, size_t name_size, in
return true; return true;
} }
/* rsync --debug/--info categories that FastSync accepts for CLI parity but has
* no output wired to (yet). They must parse successfully so a valid rsync
* invocation is not rejected up front; only categories with a FastSync
* counterpart set a log flag. `pack`/`util` are FastSync-specific (packed
* metadata / general utility logging). `syms`, `hl`, and `owner` are aliases
* 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",
};
for (size_t i = 0; i < sizeof(categories) / sizeof(categories[0]); i++) {
if (strcmp(name, categories[i]) == 0)
return true;
}
return false;
}
static bool is_accepted_info_category(const char* name) {
static const char* const categories[] = {
"backup", "del", "flist", "mount", "nonreg", "progress", "remove", "syms", "symsafe",
};
for (size_t i = 0; i < sizeof(categories) / sizeof(categories[0]); i++) {
if (strcmp(name, categories[i]) == 0)
return true;
}
return false;
}
static int parse_debug_flags(const char* value, Config* config) { static int parse_debug_flags(const char* value, Config* config) {
if (!value || value[0] == '\0' || value[0] == ',' || value[strlen(value) - 1] == ',' || if (!value || value[0] == '\0' || value[0] == ',' || value[strlen(value) - 1] == ',' ||
strstr(value, ",,")) { strstr(value, ",,")) {
@@ -522,6 +552,8 @@ static int parse_debug_flags(const char* value, Config* config) {
flag = LOG_DEBUG_PACK; flag = LOG_DEBUG_PACK;
} else if (strcmp(name, "util") == 0) { } else if (strcmp(name, "util") == 0) {
flag = LOG_DEBUG_UTIL; flag = LOG_DEBUG_UTIL;
} else if (is_accepted_debug_category(name)) {
continue;
} else { } else {
log_message(LOG_LEVEL_ERROR, "unsupported --debug flag: %s", token); log_message(LOG_LEVEL_ERROR, "unsupported --debug flag: %s", token);
free(flags); free(flags);
@@ -584,6 +616,8 @@ static int parse_info_flags(const char* value, Config* config) {
flag = LOG_INFO_SKIP; flag = LOG_INFO_SKIP;
else if (strcmp(name, "stats") == 0) else if (strcmp(name, "stats") == 0)
flag = LOG_INFO_STATS; flag = LOG_INFO_STATS;
else if (is_accepted_info_category(name))
continue;
else { else {
log_message(LOG_LEVEL_ERROR, "unsupported --info flag: %s", token); log_message(LOG_LEVEL_ERROR, "unsupported --info flag: %s", token);
free(flags); free(flags);
+9 -4
View File
@@ -341,15 +341,20 @@ void print_usage(void) {
} }
void print_debug_usage(void) { void print_debug_usage(void) {
printf("Supported debug flags: IO,PROTO,PACK,UTIL,ALL,NONE\n"); printf("Emitting debug flags: IO,PROTO,PACK,UTIL,ALL,NONE\n");
printf("Also accepted for rsync CLI parity (silent): ACL,BACKUP,BIND,CHDIR,\n");
printf("CONNECT,CMD,DEL,DELTASUM,DUP,EXIT,FILTER,FLIST,FUZZY,GENR,HASH,HLINK,\n");
printf("ICONV,NSTR,OWN,RECV,SEND,TIME.\n");
printf("Flags may be comma-separated, for example: --debug=io,proto\n"); printf("Flags may be comma-separated, for example: --debug=io,proto\n");
printf("An optional level suffix is accepted (e.g. --debug=io2); level 0\n"); printf("An optional level suffix is accepted (e.g. --debug=io2); level 0\n");
printf("silences that item. Other rsync debug flags are unsupported and rejected.\n"); printf("silences that item. Unknown names are rejected.\n");
} }
void print_info_usage(void) { void print_info_usage(void) {
printf("Supported info flags: COPY,NAME,MISC,SKIP,STATS,ALL,NONE\n"); printf("Emitting info flags: COPY,NAME,MISC,SKIP,STATS,ALL,NONE\n");
printf("Also accepted for rsync CLI parity (silent): BACKUP,DEL,FLIST,MOUNT,\n");
printf("NONREG,PROGRESS,REMOVE,SYMSAFE.\n");
printf("Flags may be comma-separated, for example: --info=name,stats\n"); printf("Flags may be comma-separated, for example: --info=name,stats\n");
printf("An optional level suffix is accepted (e.g. --info=stats2); level 0\n"); printf("An optional level suffix is accepted (e.g. --info=stats2); level 0\n");
printf("silences that item. Other rsync info flags are unsupported and rejected.\n"); printf("silences that item. Unknown names are rejected.\n");
} }
+11 -5
View File
@@ -101,9 +101,15 @@ class CountingProxy:
return return
counter[0] += len(data) counter[0] += len(data)
def run(self, cmd): def run(self, cmd, join_timeout=20):
"""Forward one client run (the full command list) to the real server and """Forward one client run (the full command list) to the real server and
return the CompletedProcess after the counts have settled.""" return the CompletedProcess after the counts have settled.
``join_timeout`` bounds how long to wait for the forwarding threads. The
client->server count is published as soon as the client side reaches EOF
(i.e. once the client process has exited), so callers that only need that
count can pass a small value instead of waiting for the server to close
its idle socket."""
def serve(): def serve():
try: try:
@@ -119,15 +125,15 @@ class CountingProxy:
a.start() a.start()
b.start() b.start()
a.join() a.join()
b.join()
self.client_to_server = c2s[0] self.client_to_server = c2s[0]
b.join()
self.server_to_client = s2c[0] self.server_to_client = s2c[0]
self._listener.close() self._listener.close()
thread = threading.Thread(target=serve) thread = threading.Thread(target=serve, daemon=True)
thread.start() thread.start()
result = subprocess.run(cmd, capture_output=True, text=True, timeout=180) result = subprocess.run(cmd, capture_output=True, text=True, timeout=180)
thread.join(20) thread.join(join_timeout)
return result return result
+39 -13
View File
@@ -101,15 +101,24 @@ class _SlicingProxy:
"""Forward the client stream to a server, optionally cutting it or invoking a """Forward the client stream to a server, optionally cutting it or invoking a
hook after a byte threshold. ``forward_limit`` mode resets both ends after hook after a byte threshold. ``forward_limit`` mode resets both ends after
that many client bytes (a mid-transfer failure). ``hook`` mode calls the that many client bytes (a mid-transfer failure). ``hook`` mode calls the
hook once and keeps forwarding to completion.""" hook once and keeps forwarding to completion.
With ``wait_for_reply`` the hook is a real barrier, not a timing guess: it
fires only after the server has sent *any* reply, which the receiver does
only after it has consumed the frames that precede the payload (the
per-directory delete plan for ``--delete-delay``). The caller pairs it with
``--incremental`` so a per-file handshake reply is guaranteed mid-transfer.
"""
def __init__(self, target_port, forward_limit=None, hook=None, hook_after=0, def __init__(self, target_port, forward_limit=None, hook=None, hook_after=0,
throttle=0.0): throttle=0.0, wait_for_reply=False):
self.target = ("127.0.0.1", target_port) self.target = ("127.0.0.1", target_port)
self.forward_limit = forward_limit self.forward_limit = forward_limit
self.hook = hook self.hook = hook
self.hook_after = hook_after self.hook_after = hook_after
self.throttle = throttle self.throttle = throttle
self.wait_for_reply = wait_for_reply
self.server_replied = threading.Event()
self.hook_called = threading.Event() self.hook_called = threading.Event()
self.listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM) self.listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
self.listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) self.listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
@@ -158,15 +167,7 @@ class _SlicingProxy:
data = data[:room] data = data[:room]
backend.sendall(data) backend.sendall(data)
forwarded += len(data) forwarded += len(data)
if (self.hook is not None and not self.hook_called.is_set() self._maybe_hook(forwarded)
and forwarded >= self.hook_after):
# Give the receiver time to process the (tiny) plan
# frames that precede this offset before the hook
# mutates the destination.
if self.throttle > 0:
time.sleep(0.2)
self.hook()
self.hook_called.set()
if self.forward_limit is not None and forwarded >= self.forward_limit: if self.forward_limit is not None and forwarded >= self.forward_limit:
socks = [] socks = []
break break
@@ -174,6 +175,10 @@ class _SlicingProxy:
time.sleep(self.throttle) time.sleep(self.throttle)
else: else:
client.sendall(data) client.sendall(data)
# Any server reply proves the receiver consumed the
# frames that precede it, so the hook barrier is met.
self.server_replied.set()
self._maybe_hook(forwarded)
except OSError: except OSError:
pass pass
for sock in (client, backend): for sock in (client, backend):
@@ -190,6 +195,19 @@ class _SlicingProxy:
except OSError: except OSError:
pass pass
def _maybe_hook(self, forwarded):
"""Fire the one-shot hook once its barrier is satisfied: enough client
bytes have been forwarded and, when ``wait_for_reply`` is set, the
server has sent a reply proving it processed the preceding frames."""
if self.hook is None or self.hook_called.is_set():
return
if forwarded < self.hook_after:
return
if self.wait_for_reply and not self.server_replied.is_set():
return
self.hook()
self.hook_called.set()
def finish(self): def finish(self):
self._thread.join(30) self._thread.join(30)
try: try:
@@ -317,8 +335,16 @@ class TestDeleteDelayVsAfterSnapshot:
# after the directory's plan (delay) has been processed. # after the directory's plan (delay) has been processed.
_write(new_extra, b"created mid-transfer\n") _write(new_extra, b"created mid-transfer\n")
proxy = _SlicingProxy(server.port, hook=hook, hook_after=MID_TRANSFER_BYTES, throttle=PROXY_THROTTLE) # --incremental gives the receiver a mid-transfer handshake
flags = [timing] + (["--threads"] if mt else []) # reply; the proxy waits for it (wait_for_reply) so the hook is
# causally after the plan frame, never a timing guess.
# --ignore-times forces the big file to transfer on the second
# timing too (the first run already installed it), keeping the
# mid-transfer reply present in both iterations.
proxy = _SlicingProxy(server.port, hook=hook,
hook_after=MID_TRANSFER_BYTES,
throttle=PROXY_THROTTLE, wait_for_reply=True)
flags = [timing, "--incremental", "--ignore-times"] + (["--threads"] if mt else [])
result, _ = run_client(source, dest, flags=flags, port=proxy.port) result, _ = run_client(source, dest, flags=flags, port=proxy.port)
proxy.finish() proxy.finish()
assert result.returncode == 0, ( assert result.returncode == 0, (
+166 -23
View File
@@ -5,6 +5,7 @@ differential tests run the SAME transfer with real ``rsync 3.4.1`` and with
fastsync and compare stdout, so they are skipped when rsync is unavailable. fastsync and compare stdout, so they are skipped when rsync is unavailable.
""" """
import os import os
import re
import shutil import shutil
import subprocess import subprocess
import sys import sys
@@ -324,37 +325,116 @@ class TestWireStatsParity:
@requires_rsync @requires_rsync
@pytest.mark.ci @pytest.mark.ci
def test_out_format_b_is_wire_bytes(self, shared_server): def test_out_format_b_is_wire_bytes(self, shared_server):
"""%b is true transferred (wire) bytes, not the source length: it must """%b is the bytes actually transferred (wire), not the source length.
differ from %l (the source length) and exceed it for a framed transfer."""
A differential run against rsync confirms both implementations report a
framed value greater than %l. The exact numbers are not compared: each
counts its own protocol framing and checksum trailer, so the two are
protocol-specific and cannot be numerically equal (documented
divergence)."""
source = os.path.join(TEST_DATA_DIR, "wire_b_src") source = os.path.join(TEST_DATA_DIR, "wire_b_src")
dest = os.path.join(TEST_DATA_DIR, "wire_b_dst") dest = os.path.join(TEST_DATA_DIR, "wire_b_dst")
rdst = os.path.join(TEST_DATA_DIR, "wire_b_rdst")
_make_one_file(source, "f.bin", 5000) _make_one_file(source, "f.bin", 5000)
clean_dir(dest) clean_dir(dest)
result, _ = run_client(source, dest, flags=["-a", "--out-format=%b %l %c"], clean_dir(rdst)
fmt = "%b %l"
rsync_result = _rsync(["-a", "--out-format=" + fmt, source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(source, dest, flags=["-a", "--out-format=" + fmt],
port=shared_server.port) port=shared_server.port)
assert result.returncode == 0, result.stderr[:300] assert result.returncode == 0, result.stderr[:300]
line = result.stdout.strip() rb, rl = (int(x) for x in rsync_result.stdout.split()[:2])
parts = line.split() fb, fl = (int(x) for x in result.stdout.split()[:2])
assert len(parts) == 3 and all(p.isdigit() for p in parts), line assert rl == fl == 5000, (rsync_result.stdout, result.stdout)
wire_b, src_l, wire_c = (int(p) for p in parts) assert rb > rl, f"rsync %b must include framing: {rsync_result.stdout!r}"
assert src_l == 5000, line assert fb > fl, f"fastsync %b must include framing: {result.stdout!r}"
assert wire_b > src_l, f"%b must include wire framing: {line}"
@requires_rsync @requires_rsync
@pytest.mark.ci @pytest.mark.ci
def test_progress_first_frame_matches_rsync(self, shared_server): def test_out_format_c_whole_file_matches_rsync(self, shared_server):
"""For a sub-32 KiB file the first --progress frame is deterministic """%c is the block-checksum bytes received. rsync reports its 16-byte
(0.00 kB/s, 0:00:00) and must be byte-identical to rsync's.""" sum header even for a whole-file transfer (no basis), so `%c` must match
rsync exactly for the whole-file case."""
source = os.path.join(TEST_DATA_DIR, "wire_c_src")
dest = os.path.join(TEST_DATA_DIR, "wire_c_dst")
rdst = os.path.join(TEST_DATA_DIR, "wire_c_rdst")
_make_one_file(source, "f.bin", 5000)
clean_dir(dest)
clean_dir(rdst)
fmt = "%c %l %n"
rsync_result = _rsync(["-a", "--out-format=" + fmt, source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(source, dest, flags=["-a", "--out-format=" + fmt],
port=shared_server.port)
assert result.returncode == 0, result.stderr[:300]
def file_lines(text):
return [
line for line in text.splitlines()
if line and not line.rsplit(" ", 1)[-1].endswith("/")
]
assert file_lines(result.stdout) == file_lines(rsync_result.stdout), (
f"rsync={rsync_result.stdout!r} fastsync={result.stdout!r}"
)
assert result.stdout.split()[0] == rsync_result.stdout.split()[0] == "16", (
f"%c must be rsync's 16-byte sum header: {result.stdout!r}"
)
@requires_rsync
@pytest.mark.ci
def test_out_format_c_delta_mode_divergence(self, shared_server):
"""Documented residual: with delta enabled, rsync's %c is its 16-byte sum
header plus one checksum entry per block (protocol-specific, so it grows
with the basis size), while FastSync's %c is the bytes of its own delta
handshake. FastSync's delta %c therefore cannot match rsync numerically;
only the whole-file case is aligned. Pinned here so a future change is
noticed."""
source = os.path.join(TEST_DATA_DIR, "wire_cd_src")
dest = os.path.join(TEST_DATA_DIR, "wire_cd_dst")
rdst = os.path.join(TEST_DATA_DIR, "wire_cd_rdst")
_make_one_file(source, "f.bin", 5000)
clean_dir(dest)
clean_dir(rdst)
fmt = "%c %l"
# rsync local default is whole-file; force the block-delta path.
rsync_result = _rsync(["-a", "--no-whole-file", "--out-format=" + fmt,
source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(source, dest,
flags=["-a", "--incremental", "--delta",
"--out-format=" + fmt],
port=shared_server.port)
assert result.returncode == 0, result.stderr[:300]
rs_c = int(rsync_result.stdout.split()[0])
fs_c = int(result.stdout.split()[0])
# No basis exists, so rsync still reports only its sum header.
assert rs_c == 16, rsync_result.stdout
# FastSync reports its own handshake bytes and is not aligned.
assert fs_c > 16, (
f"FastSync delta %c changed to {fs_c}; the documented divergence "
"may be closable now"
)
@requires_rsync
@pytest.mark.ci
@pytest.mark.parametrize("mt", [False, True])
@pytest.mark.parametrize("progress_flag", ["--progress", "-P"])
def test_progress_first_frame_matches_rsync(self, shared_server, progress_flag, mt):
"""For a sub-32 KiB file the first --progress/-P frame is deterministic
(0.00 kB/s, 0:00:00) and must be byte-identical to rsync's, in both the
single-threaded and --threads send paths."""
source = os.path.join(TEST_DATA_DIR, "wire_pg_src") source = os.path.join(TEST_DATA_DIR, "wire_pg_src")
dest = os.path.join(TEST_DATA_DIR, "wire_pg_dst") dest = os.path.join(TEST_DATA_DIR, "wire_pg_dst")
rdst = os.path.join(TEST_DATA_DIR, "wire_pg_rdst") rdst = os.path.join(TEST_DATA_DIR, "wire_pg_rdst")
_make_one_file(source, "f.bin", 100) _make_one_file(source, "f.bin", 100)
clean_dir(dest) clean_dir(dest)
clean_dir(rdst) clean_dir(rdst)
rsync_result = _rsync(["-a", "--progress", source + "/", rdst + "/"]) rsync_result = _rsync(["-a", progress_flag, source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(source, dest, flags=["-a", "--progress"], flags = ["-a", progress_flag] + (["--threads"] if mt else [])
port=shared_server.port) result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, result.stderr[:300] assert result.returncode == 0, result.stderr[:300]
def frames(text): def frames(text):
@@ -369,8 +449,10 @@ class TestWireStatsParity:
@requires_rsync @requires_rsync
@pytest.mark.ci @pytest.mark.ci
def test_stats_selected_lines_match_rsync(self, shared_server): @pytest.mark.parametrize("mt", [False, True])
"""The protocol-independent --stats lines must match rsync exactly.""" def test_stats_selected_lines_match_rsync(self, shared_server, mt):
"""The protocol-independent --stats lines must match rsync exactly, in
both the single-threaded and --threads (multithreaded) send paths."""
source = os.path.join(TEST_DATA_DIR, "wire_st_src") source = os.path.join(TEST_DATA_DIR, "wire_st_src")
dest = os.path.join(TEST_DATA_DIR, "wire_st_dst") dest = os.path.join(TEST_DATA_DIR, "wire_st_dst")
rdst = os.path.join(TEST_DATA_DIR, "wire_st_rdst") rdst = os.path.join(TEST_DATA_DIR, "wire_st_rdst")
@@ -379,8 +461,8 @@ class TestWireStatsParity:
clean_dir(rdst) clean_dir(rdst)
rsync_result = _rsync(["-a", "--stats", source + "/", rdst + "/"]) rsync_result = _rsync(["-a", "--stats", source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(source, dest, flags=["-a", "--stats"], flags = ["-a", "--stats"] + (["--threads"] if mt else [])
port=shared_server.port) result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
assert result.returncode == 0, result.stderr[:300] assert result.returncode == 0, result.stderr[:300]
keys = ( keys = (
"Number of regular files transferred", "Number of regular files transferred",
@@ -389,6 +471,7 @@ class TestWireStatsParity:
"Literal data", "Literal data",
"Matched data", "Matched data",
"Number of deleted files", "Number of deleted files",
"File list size",
) )
def pick(text): def pick(text):
@@ -405,8 +488,68 @@ class TestWireStatsParity:
@requires_rsync @requires_rsync
@pytest.mark.ci @pytest.mark.ci
def test_dry_run_delete_lines_match_rsync(self): def test_stats_file_count_breakdown_residual(self, shared_server):
"""-n --delete emits transfer-relative `*deleting` lines like rsync.""" """Residual (row #3): rsync prints the `Number of files` and
`Number of created files` lines with a per-type breakdown
(`(reg: X, dir: Y, link: Z)`).
FastSync cannot reproduce it from what the sender currently knows: the
scanner does not put directory entries in the transfer list (directories
are created implicitly), and without a per-entry destination-probe the
sender cannot tell which entries the receiver newly created. So FastSync
prints the bare transferred-entry count. This test pins the divergence
explicitly -- the row must not be marked ✅.
"""
source = os.path.join(TEST_DATA_DIR, "wire_stc_src")
dest = os.path.join(TEST_DATA_DIR, "wire_stc_dst")
rdst = os.path.join(TEST_DATA_DIR, "wire_stc_rdst")
_make_one_file(source, "f.bin", 6000)
clean_dir(dest)
clean_dir(rdst)
rsync_result = _rsync(["-a", "--stats", source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(source, dest, flags=["-a", "--stats"],
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")
r_created = stats_line(rsync_result.stdout, "Number of created files")
f_files = stats_line(result.stdout, "Number of files")
f_created = stats_line(result.stdout, "Number of created files")
# rsync always carries the type breakdown (the source root counts as a
# directory; the single regular file as reg).
assert re.match(r"Number of files: 2 \(reg: 1, dir: 1\)$", r_files), r_files
assert re.match(r"Number of created files: 1 \(reg: 1\)$", r_created), r_created
# FastSync prints only the bare count: no directory accounting and no
# per-entry "created" knowledge.
assert re.fullmatch(r"Number of files: 1", f_files), f_files
assert re.fullmatch(r"Number of created files: 1", f_created), f_created
@requires_rsync
@pytest.mark.ci
@pytest.mark.parametrize("mt", [
False,
pytest.param(
True,
marks=pytest.mark.xfail(
reason="known gap: the --threads dry-run delete path does not "
"consume the receiver's STATUS_STATS delete list yet, so "
"`-n --delete` emits no *deleting lines (tracked by the "
"parity-blockers STATUS_STATS fix)",
strict=False,
),
),
])
def test_dry_run_delete_lines_match_rsync(self, mt):
"""-n --delete emits transfer-relative `*deleting` lines like rsync
(single-threaded; the --threads variant is a documented xfail)."""
source = os.path.join(TEST_DATA_DIR, "wire_del_src") source = os.path.join(TEST_DATA_DIR, "wire_del_src")
dest = os.path.join(TEST_DATA_DIR, "wire_del_dst") dest = os.path.join(TEST_DATA_DIR, "wire_del_dst")
rdst = os.path.join(TEST_DATA_DIR, "wire_del_rdst") rdst = os.path.join(TEST_DATA_DIR, "wire_del_rdst")
@@ -442,10 +585,10 @@ class TestWireStatsParity:
line for line in rsync_result.stdout.splitlines() if line.startswith("*deleting") line for line in rsync_result.stdout.splitlines() if line.startswith("*deleting")
) )
# The shared session server refuses deletion; start one that allows it. # The shared session server refuses deletion; start one that allows it.
flags = ["-a", "-n", "--delete", "-i"] + (["--threads"] if mt else [])
with ServerManager() as server: with ServerManager() as server:
server.start(extra_args=["--allow-delete"]) server.start(extra_args=["--allow-delete"])
result, _ = run_client(source, dest, flags=["-a", "-n", "--delete", "-i"], result, _ = run_client(source, dest, flags=flags, port=server.port)
port=server.port)
assert result.returncode == 0, result.stderr[:300] assert result.returncode == 0, result.stderr[:300]
fast_del = sorted( fast_del = sorted(
line for line in result.stdout.splitlines() if line.startswith("*deleting") line for line in result.stdout.splitlines() if line.startswith("*deleting")
+191 -41
View File
@@ -12,6 +12,8 @@ import pytest
sys.path.insert(0, os.path.dirname(__file__)) sys.path.insert(0, os.path.dirname(__file__))
from common import ( from common import (
CLIENT_CMD,
CountingProxy,
TEST_DATA_DIR, TEST_DATA_DIR,
ServerManager, ServerManager,
run_client, run_client,
@@ -593,6 +595,57 @@ class TestVerifyAndFlip:
assert result.returncode == 0, result.stderr[:300] assert result.returncode == 0, result.stderr[:300]
_assert_same_tree(rdst, get_dest_received_dir(dest, source), "(--preallocate)") _assert_same_tree(rdst, get_dest_received_dir(dest, source), "(--preallocate)")
@requires_rsync
@pytest.mark.ci
def test_sparse_preallocate_blocks_match_rsync(self, shared_server):
"""rsync lets --preallocate win over --sparse: the hole file gets its
full space reserved (st_blocks ~ size/512) even though the sparse writer
seeks over the zero run. FastSync must agree both on the block count and
with rsync's exact value for each flag combination."""
source = self._src("sparsepre")
dest = self._dst("sparsepre")
rdst = self._dst("sparsepre_r")
total = 1024 * 1024
# A zero run >= the 4 KiB sparse threshold in the middle of the image.
with open(os.path.join(source, "hole.bin"), "wb") as fh:
fh.write(os.urandom(64 * 1024))
fh.write(b"\x00" * (768 * 1024))
fh.write(os.urandom(total - 64 * 1024 - 768 * 1024))
def run_both(flags, tag):
rsync_dst = self._dst(f"sparsepre_{tag}_r")
fs_dst = self._dst(f"sparsepre_{tag}_f")
assert _rsync(["-a"] + flags + [source + "/", rsync_dst + "/"]).returncode == 0
result, _ = run_client(source, fs_dst, flags=["-a"] + flags,
port=shared_server.port)
assert result.returncode == 0, result.stderr[:300]
rfile = os.path.join(rsync_dst, "hole.bin")
ffile = os.path.join(get_dest_received_dir(fs_dst, source), "hole.bin")
with open(rfile, "rb") as rfh, open(ffile, "rb") as ffh:
assert rfh.read() == ffh.read(), f"{tag}: content diverged"
return os.stat(rfile).st_blocks, os.stat(ffile).st_blocks
r_sparse, f_sparse = run_both(["--sparse"], "sparse")
r_both, f_both = run_both(["--sparse", "--preallocate"], "both")
assert f_sparse == r_sparse, (
f"--sparse st_blocks: fastsync={f_sparse} rsync={r_sparse}"
)
assert f_both == r_both, (
f"--sparse --preallocate st_blocks: fastsync={f_both} rsync={r_both}"
)
# Only meaningful where the filesystem actually reports holes; otherwise
# both sides simply allocate the full size and the equality above holds.
has_holes = r_sparse * 512 < total
if has_holes:
assert f_both > f_sparse, (
"preallocate must win over sparse: "
f"sparse={f_sparse} both={f_both} blocks"
)
assert r_both > r_sparse, (
f"rsync preallocate must win: sparse={r_sparse} both={r_both}"
)
@requires_rsync @requires_rsync
def test_fuzzy_content_matches_rsync(self, shared_server): def test_fuzzy_content_matches_rsync(self, shared_server):
source = self._src("fuzzy") source = self._src("fuzzy")
@@ -717,6 +770,58 @@ class TestVerifyAndFlip:
"--link-dest must hard-link to the basis file" "--link-dest must hard-link to the basis file"
class TestIgnoreExistingShortCircuit:
"""#9: --ignore-existing is decided by the receiver during the per-file
check, before the sender streams any payload. A large destination file that
is already present must therefore cost almost no wire bytes, not the full
file; rsync short-circuits the same way."""
@requires_rsync
@pytest.mark.ci
def test_skip_answered_before_payload(self):
source = os.path.join(TEST_DATA_DIR, "qw_ie_wire_src")
dest = os.path.join(TEST_DATA_DIR, "qw_ie_wire_dst")
rdst = os.path.join(TEST_DATA_DIR, "qw_ie_wire_rdst")
clean_dir(source)
clean_dir(dest)
clean_dir(rdst)
big = b"S" * (4 * 1024 * 1024)
with open(os.path.join(source, "big.bin"), "wb") as fh:
fh.write(big)
for root in (get_dest_received_dir(dest, source), rdst):
os.makedirs(root, exist_ok=True)
with open(os.path.join(root, "big.bin"), "wb") as fh:
fh.write(b"D" * len(big))
assert _rsync(["-a", "--ignore-existing", source + "/",
rdst + "/"]).returncode == 0
# A dedicated one-shot server keeps the proxy's connection teardown
# deterministic (the session-wide server can linger on an idle socket).
with ServerManager() as server:
cmd = CLIENT_CMD + [
"--source-dir", source, "--dest-dir", dest, "--save-to-disk",
"--server-port", str(server.port), "-a", "--ignore-existing",
]
proxy = CountingProxy(server.port)
# Only the client->server count matters here; it is final once the
# client exits, so do not wait for the server to close its idle
# socket (which can take the full default join timeout).
result = proxy.run(cmd, join_timeout=1.0)
assert result.returncode == 0, (result.stderr or result.stdout)[:300]
# The receiver answered the skip before the sender transmitted the file:
# only config/path/check frames crossed, not the 4 MiB payload.
assert proxy.client_to_server < len(big) // 10, (
f"--ignore-existing transmitted {proxy.client_to_server} bytes for a "
f"skipped {len(big)}-byte file"
)
received = get_dest_received_dir(dest, source)
with open(os.path.join(received, "big.bin"), "rb") as fh:
assert fh.read() == b"D" * len(big), \
"--ignore-existing must preserve the destination content"
_assert_same_tree(rdst, received, "(--ignore-existing wire short-circuit)")
@pytest.mark.skipif(os.geteuid() != 0, reason="ownership mapping requires root") @pytest.mark.skipif(os.geteuid() != 0, reason="ownership mapping requires root")
class TestOwnershipMapping: class TestOwnershipMapping:
"""#33/#34/#35: --usermap/--groupmap/--chown match rsync's numeric result.""" """#33/#34/#35: --usermap/--groupmap/--chown match rsync's numeric result."""
@@ -782,12 +887,33 @@ class TestFakeSuper:
assert fh.read() == b"fake-super-data\n" assert fh.read() == b"fake-super-data\n"
class TestInfoDebugFlagParity: # rsync 3.4.1's full --info/--debug vocabularies (from `rsync --info=help` /
"""#01/#02: rsync's info/debug spellings are either mapped to real output or # `--debug=help`). FastSync must accept every one of them; only the categories
rejected by name (never silently ignored).""" # with an existing FastSync counterpart emit output, the rest are accepted but
# currently silent.
RSYNC_INFO_CATEGORIES = (
"backup", "copy", "del", "flist", "misc", "mount", "name", "nonreg",
"progress", "remove", "skip", "stats", "symsafe", "all", "none",
)
RSYNC_DEBUG_CATEGORIES = (
"acl", "backup", "bind", "chdir", "connect", "cmd", "del", "deltasum",
"dup", "exit", "filter", "flist", "fuzzy", "genr", "hash", "hlink",
"iconv", "io", "nstr", "own", "proto", "recv", "send", "time", "all",
"none",
)
# Extra categories FastSync also accepts: rsync's own help spells these
# `symsafe`/`hlink`/`own`, but the historical aliases are kept working, and
# `pack`/`util` are FastSync-specific debug channels.
FASTSYNC_INFO_ALIASES = ("syms",)
FASTSYNC_DEBUG_ALIASES = ("hl", "owner", "pack", "util")
@requires_rsync
def test_mapped_info_categories_accepted_like_rsync(self, shared_server): class TestInfoDebugFlagParity:
"""#01/#02: FastSync accepts rsync 3.4.1's full --info/--debug vocabulary
(with level suffixes) so a valid rsync invocation is never rejected up
front. Unknown names are still refused by name."""
def _tree(self):
source = os.path.join(TEST_DATA_DIR, "qw_flags_src") source = os.path.join(TEST_DATA_DIR, "qw_flags_src")
dest = os.path.join(TEST_DATA_DIR, "qw_flags_dst") dest = os.path.join(TEST_DATA_DIR, "qw_flags_dst")
rdst = os.path.join(TEST_DATA_DIR, "qw_flags_rdst") rdst = os.path.join(TEST_DATA_DIR, "qw_flags_rdst")
@@ -796,54 +922,78 @@ class TestInfoDebugFlagParity:
clean_dir(rdst) clean_dir(rdst)
with open(os.path.join(source, "a.txt"), "wb") as fh: with open(os.path.join(source, "a.txt"), "wb") as fh:
fh.write(b"a\n") fh.write(b"a\n")
for cat in ("stats2", "name", "copy", "misc", "skip", "STATS2"): return source, dest, rdst
assert _rsync(["-a", "--info=" + cat, source + "/", rdst + "/"]).returncode == 0
@requires_rsync
@pytest.mark.ci
def test_info_vocabulary_accepted_like_rsync(self, shared_server):
source, dest, rdst = self._tree()
for cat in RSYNC_INFO_CATEGORIES:
rs = _rsync(["-a", "--info=" + cat, source + "/", rdst + "/"])
assert rs.returncode == 0, f"rsync rejected --info={cat}: {rs.stderr}"
clean_dir(rdst) clean_dir(rdst)
result, _ = run_client(source, dest, flags=["--info=" + cat], result, _ = run_client(source, dest, flags=["--info=" + cat],
port=shared_server.port) port=shared_server.port)
assert result.returncode == 0, ( assert result.returncode == 0, (
f"--info={cat} must be accepted: {result.stderr[:200]}" f"--info={cat} must be accepted like rsync: {result.stderr[:200]}"
)
@requires_rsync
def test_mapped_debug_categories_accepted_like_rsync(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "qw_dflags_src")
dest = os.path.join(TEST_DATA_DIR, "qw_dflags_dst")
rdst = os.path.join(TEST_DATA_DIR, "qw_dflags_rdst")
clean_dir(source)
clean_dir(dest)
clean_dir(rdst)
with open(os.path.join(source, "a.txt"), "wb") as fh:
fh.write(b"a\n")
for cat in ("io2", "proto0", "all"):
assert _rsync(["-a", "--debug=" + cat, source + "/", rdst + "/"]).returncode == 0
clean_dir(rdst)
result, _ = run_client(source, dest, flags=["--debug=" + cat],
port=shared_server.port)
assert result.returncode == 0, (
f"--debug={cat} must be accepted: {result.stderr[:200]}"
) )
@requires_rsync @requires_rsync
@pytest.mark.ci @pytest.mark.ci
def test_unmapped_categories_rejected_by_name(self, shared_server): def test_debug_vocabulary_accepted_like_rsync(self, shared_server):
"""rsync accepts del/filter; fastsync has no mapping so it must refuse source, dest, rdst = self._tree()
loudly, naming the category, rather than silently ignoring it.""" for cat in RSYNC_DEBUG_CATEGORIES:
source = os.path.join(TEST_DATA_DIR, "qw_umap_src") rs = _rsync(["-a", "--debug=" + cat, source + "/", rdst + "/"])
dest = os.path.join(TEST_DATA_DIR, "qw_umap_dst") assert rs.returncode == 0, f"rsync rejected --debug={cat}: {rs.stderr}"
clean_dir(source) clean_dir(rdst)
clean_dir(dest) result, _ = run_client(source, dest, flags=["--debug=" + cat],
with open(os.path.join(source, "a.txt"), "wb") as fh: port=shared_server.port)
fh.write(b"a\n") assert result.returncode == 0, (
# rsync accepts these (so they are valid rsync invocations). f"--debug={cat} must be accepted like rsync: {result.stderr[:200]}"
assert _rsync(["-a", "--info=del", source + "/", dest + "/"]).returncode == 0 )
assert _rsync(["-a", "--debug=filter", source + "/", dest + "/"]).returncode == 0
for flag, name in (("--info=del", "del"), ("--debug=filter", "filter")): @requires_rsync
@pytest.mark.ci
def test_level_suffixes_accepted_like_rsync(self, shared_server):
source, dest, rdst = self._tree()
for flag in ("--info=stats2", "--info=copy0", "--info=all0",
"--debug=io2", "--debug=proto0", "--debug=all4"):
rs = _rsync(["-a", flag, source + "/", rdst + "/"])
assert rs.returncode == 0, f"rsync rejected {flag}: {rs.stderr}"
clean_dir(rdst)
result, _ = run_client(source, dest, flags=[flag],
port=shared_server.port)
assert result.returncode == 0, (
f"{flag} must be accepted like rsync: {result.stderr[:200]}"
)
def test_fastsync_alias_categories_accepted(self, shared_server):
"""FastSync-specific/alias spellings: accepted (rsync spells them
symsafe/hlink/own) but not emitted."""
source, dest, _ = self._tree()
for flag in (["--info=" + c for c in FASTSYNC_INFO_ALIASES] +
["--debug=" + c for c in FASTSYNC_DEBUG_ALIASES]):
result, _ = run_client(source, dest, flags=[flag],
port=shared_server.port)
assert result.returncode == 0, (
f"{flag} must be accepted: {result.stderr[:200]}"
)
@requires_rsync
@pytest.mark.ci
def test_unknown_categories_rejected_by_name(self, shared_server):
"""Truly unknown names are refused by name, exactly like rsync."""
source, dest, rdst = self._tree()
for flag, name in (("--info=bogus", "bogus"),
("--debug=bogus", "bogus")):
rs = _rsync(["-a", flag, source + "/", rdst + "/"])
assert rs.returncode != 0, f"rsync unexpectedly accepted {flag}"
result, _ = run_client(source, dest, flags=[flag], result, _ = run_client(source, dest, flags=[flag],
port=shared_server.port) port=shared_server.port)
assert result.returncode != 0, f"{flag} must be rejected" assert result.returncode != 0, f"{flag} must be rejected"
assert name in (result.stderr or ""), \ assert name in (result.stderr or ""), (
f"{flag} must be rejected by name, got: {result.stderr[:200]}" f"{flag} must be rejected by name, got: {result.stderr[:200]}"
)
class TestStopAtParity: class TestStopAtParity:
+20 -1
View File
@@ -776,7 +776,7 @@ static void test_parse_args_debug_help() {
} }
static void test_parse_args_debug_flags_validation() { static void test_parse_args_debug_flags_validation() {
static const char* const values[] = {"", "io,", ",io", "io,,proto", "acl", "tls", "unknown"}; static const char* const values[] = {"", "io,", ",io", "io,,proto", "tls", "unknown"};
for (size_t i = 0; i < sizeof(values) / sizeof(values[0]); i++) { for (size_t i = 0; i < sizeof(values) / sizeof(values[0]); i++) {
Config* cfg = config_create(); Config* cfg = config_create();
char option[64]; char option[64];
@@ -1345,6 +1345,24 @@ static void test_parse_args_info_name_and_help() {
config_delete(cfg); config_delete(cfg);
} }
/* rsync 3.4.1's remaining --info/--debug categories parse successfully but
* have no FastSync output wired to them, so they must not set any log flag. */
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",
"--debug=acl,backup,bind,chdir,cmd,connect,del,deltasum,dup,exit,"
"filter,flist,fuzzy,genr,hash,hlink,iconv,nstr,own,recv,send,time,"
"hl,owner",
"/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0);
EXPECT_EQ_INT(cfg->info_level, 0);
EXPECT_EQ_INT(cfg->debug_level, 0);
config_delete(cfg);
}
/* Test parse_args with --archive flag */ /* Test parse_args with --archive flag */
/* rsync accepts a trailing level digit on --debug/--info items (e.g. io2, /* rsync accepts a trailing level digit on --debug/--info items (e.g. io2,
* all4); level 0 silences the item. */ * all4); level 0 silences the item. */
@@ -4441,6 +4459,7 @@ void test_client_cli() {
test_parse_args_debug_flags(); test_parse_args_debug_flags();
test_parse_args_debug_help(); test_parse_args_debug_help();
test_parse_args_debug_flags_validation(); test_parse_args_debug_flags_validation();
test_parse_args_rsync_flag_vocabulary_accepted();
test_parse_args_debug_info_levels(); test_parse_args_debug_info_levels();
test_parse_args_modify_window(); test_parse_args_modify_window();
test_parse_args_rejects_invalid_modify_window(); test_parse_args_rejects_invalid_modify_window();
+1 -1
View File
@@ -2920,7 +2920,7 @@ static unsigned long long capture_wire_hash(const Config* cfg, size_t* out_len)
return h; return h;
} }
/* Byte-for-byte wire compatibility guard (protocol 2.24.0). The expected hash /* Byte-for-byte wire compatibility guard (protocol 2.26.0). The expected hash
* pins the pre-X-macro byte stream; the refactor MUST NOT change it. */ * pins the pre-X-macro byte stream; the refactor MUST NOT change it. */
static void test_config_wire_golden() { static void test_config_wire_golden() {
if (is_running_under_valgrind()) if (is_running_under_valgrind())