rsync 3.4.1 drop-in parity (#285-#297) + parity completion (protocol 2.26.0) #298
@@ -1,15 +0,0 @@
|
||||
B=/workspace/build-ci
|
||||
D=/workspace/.probe
|
||||
rm -rf $D && mkdir -p $D/src $D/dst
|
||||
head -c 200000 /dev/urandom > $D/src/f.bin
|
||||
$B/server -p 45995 --allow-unauthenticated >$D/srv.log 2>&1 &
|
||||
SRV=$!; sleep 0.7
|
||||
for flags in "" "--incremental" "--incremental --delta"; do
|
||||
rm -rf $D/dst; mkdir -p $D/dst
|
||||
echo "=== fastsync flags='$flags' fresh: %b %c %l %n ==="
|
||||
$B/client --source-dir $D/src --dest-dir $D/dst --server-port 45995 --save-to-disk -a $flags --out-format="%b %c %l %n" 2>&1 | grep -v ERROR
|
||||
done
|
||||
echo "=== rsync whole-file: ==="
|
||||
rm -rf $D/rsrc $D/rdst; mkdir -p $D/rsrc $D/rdst; head -c 200000 /dev/urandom > $D/rsrc/f.bin
|
||||
rsync -a --out-format="%b %c %l %n" $D/rsrc/ $D/rdst/
|
||||
kill $SRV 2>/dev/null || true
|
||||
@@ -101,15 +101,24 @@ class _SlicingProxy:
|
||||
"""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
|
||||
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,
|
||||
throttle=0.0):
|
||||
throttle=0.0, wait_for_reply=False):
|
||||
self.target = ("127.0.0.1", target_port)
|
||||
self.forward_limit = forward_limit
|
||||
self.hook = hook
|
||||
self.hook_after = hook_after
|
||||
self.throttle = throttle
|
||||
self.wait_for_reply = wait_for_reply
|
||||
self.server_replied = threading.Event()
|
||||
self.hook_called = threading.Event()
|
||||
self.listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
self.listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
@@ -158,15 +167,7 @@ class _SlicingProxy:
|
||||
data = data[:room]
|
||||
backend.sendall(data)
|
||||
forwarded += len(data)
|
||||
if (self.hook is not None and not self.hook_called.is_set()
|
||||
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()
|
||||
self._maybe_hook(forwarded)
|
||||
if self.forward_limit is not None and forwarded >= self.forward_limit:
|
||||
socks = []
|
||||
break
|
||||
@@ -174,6 +175,10 @@ class _SlicingProxy:
|
||||
time.sleep(self.throttle)
|
||||
else:
|
||||
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:
|
||||
pass
|
||||
for sock in (client, backend):
|
||||
@@ -190,6 +195,19 @@ class _SlicingProxy:
|
||||
except OSError:
|
||||
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):
|
||||
self._thread.join(30)
|
||||
try:
|
||||
@@ -317,8 +335,16 @@ class TestDeleteDelayVsAfterSnapshot:
|
||||
# after the directory's plan (delay) has been processed.
|
||||
_write(new_extra, b"created mid-transfer\n")
|
||||
|
||||
proxy = _SlicingProxy(server.port, hook=hook, hook_after=MID_TRANSFER_BYTES, throttle=PROXY_THROTTLE)
|
||||
flags = [timing] + (["--threads"] if mt else [])
|
||||
# --incremental gives the receiver a mid-transfer handshake
|
||||
# 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)
|
||||
proxy.finish()
|
||||
assert result.returncode == 0, (
|
||||
|
||||
@@ -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.
|
||||
"""
|
||||
import os
|
||||
import re
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
@@ -383,19 +384,22 @@ class TestWireStatsParity:
|
||||
|
||||
@requires_rsync
|
||||
@pytest.mark.ci
|
||||
def test_progress_first_frame_matches_rsync(self, shared_server):
|
||||
"""For a sub-32 KiB file the first --progress frame is deterministic
|
||||
(0.00 kB/s, 0:00:00) and must be byte-identical to rsync's."""
|
||||
@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")
|
||||
dest = os.path.join(TEST_DATA_DIR, "wire_pg_dst")
|
||||
rdst = os.path.join(TEST_DATA_DIR, "wire_pg_rdst")
|
||||
_make_one_file(source, "f.bin", 100)
|
||||
clean_dir(dest)
|
||||
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
|
||||
result, _ = run_client(source, dest, flags=["-a", "--progress"],
|
||||
port=shared_server.port)
|
||||
flags = ["-a", progress_flag] + (["--threads"] if mt else [])
|
||||
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
|
||||
assert result.returncode == 0, result.stderr[:300]
|
||||
|
||||
def frames(text):
|
||||
@@ -410,8 +414,10 @@ class TestWireStatsParity:
|
||||
|
||||
@requires_rsync
|
||||
@pytest.mark.ci
|
||||
def test_stats_selected_lines_match_rsync(self, shared_server):
|
||||
"""The protocol-independent --stats lines must match rsync exactly."""
|
||||
@pytest.mark.parametrize("mt", [False, True])
|
||||
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")
|
||||
dest = os.path.join(TEST_DATA_DIR, "wire_st_dst")
|
||||
rdst = os.path.join(TEST_DATA_DIR, "wire_st_rdst")
|
||||
@@ -420,8 +426,8 @@ class TestWireStatsParity:
|
||||
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)
|
||||
flags = ["-a", "--stats"] + (["--threads"] if mt else [])
|
||||
result, _ = run_client(source, dest, flags=flags, port=shared_server.port)
|
||||
assert result.returncode == 0, result.stderr[:300]
|
||||
keys = (
|
||||
"Number of regular files transferred",
|
||||
@@ -430,6 +436,7 @@ class TestWireStatsParity:
|
||||
"Literal data",
|
||||
"Matched data",
|
||||
"Number of deleted files",
|
||||
"File list size",
|
||||
)
|
||||
|
||||
def pick(text):
|
||||
@@ -446,8 +453,68 @@ class TestWireStatsParity:
|
||||
|
||||
@requires_rsync
|
||||
@pytest.mark.ci
|
||||
def test_dry_run_delete_lines_match_rsync(self):
|
||||
"""-n --delete emits transfer-relative `*deleting` lines like rsync."""
|
||||
def test_stats_file_count_breakdown_residual(self, shared_server):
|
||||
"""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")
|
||||
dest = os.path.join(TEST_DATA_DIR, "wire_del_dst")
|
||||
rdst = os.path.join(TEST_DATA_DIR, "wire_del_rdst")
|
||||
@@ -483,10 +550,10 @@ class TestWireStatsParity:
|
||||
line for line in rsync_result.stdout.splitlines() if line.startswith("*deleting")
|
||||
)
|
||||
# The shared session server refuses deletion; start one that allows it.
|
||||
flags = ["-a", "-n", "--delete", "-i"] + (["--threads"] if mt else [])
|
||||
with ServerManager() as server:
|
||||
server.start(extra_args=["--allow-delete"])
|
||||
result, _ = run_client(source, dest, flags=["-a", "-n", "--delete", "-i"],
|
||||
port=server.port)
|
||||
result, _ = run_client(source, dest, flags=flags, port=server.port)
|
||||
assert result.returncode == 0, result.stderr[:300]
|
||||
fast_del = sorted(
|
||||
line for line in result.stdout.splitlines() if line.startswith("*deleting")
|
||||
|
||||
Reference in New Issue
Block a user