Release v2.26.0 #284
@@ -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
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
@@ -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."""
|
||||||
|
|||||||
+1
-1
@@ -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())
|
||||||
|
|||||||
Reference in New Issue
Block a user