test: --fuzzy coverage (CLI, config wire, integration)
Unit: parse -y/--fuzzy and --no-fuzzy negation (order-independent), -W and --no-delta leave fuzzy inert, --fuzzy implies --incremental/--delta, and validate_config rejects fuzzy with -s (chunk serialization) and -f (sendfile). Config: fuzzy survives the socketpair wire round-trip. Integration (TestFuzzy, byte counts via a new CountingProxy helper): the rename case transfers a 2 MiB file in a few percent of its size with byte- exact output under --fuzzy, while the no-fuzzy, no-candidate, dissimilar- sibling and --whole-file runs send the whole file; -y alias; -m parity; dest-holds-unsuitable-file falls back to the sibling; --delay-updates and --remove-source-files keep their semantics on fuzzy transfers.
This commit is contained in:
@@ -6,6 +6,7 @@ import socket
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import threading
|
||||
import time
|
||||
|
||||
PROJECT_ROOT = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", ".."))
|
||||
@@ -53,6 +54,76 @@ class ServerManager:
|
||||
self.stop()
|
||||
|
||||
|
||||
class CountingProxy:
|
||||
"""One-shot TCP forwarder that counts the bytes flowing in each direction
|
||||
between one client and the real server.
|
||||
|
||||
Client output and --stats report SOURCE lengths, so a delta/fuzzy transfer
|
||||
that moves only a few percent of the file is invisible in normal output.
|
||||
Routing the client through this proxy makes the actual wire usage
|
||||
observable: client_to_server counts every byte the client sent (config,
|
||||
paths, and file/delta payloads), server_to_client counts the reply bytes
|
||||
(including the receiver's delta signatures).
|
||||
"""
|
||||
|
||||
def __init__(self, target_port):
|
||||
self.target_port = target_port
|
||||
self._listener = socket.socket()
|
||||
self._listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
self._listener.bind(("127.0.0.1", 0))
|
||||
self._listener.listen(1)
|
||||
self._listener.settimeout(30)
|
||||
self.port = self._listener.getsockname()[1]
|
||||
self.client_to_server = 0
|
||||
self.server_to_client = 0
|
||||
|
||||
@staticmethod
|
||||
def _pump(src, dst, counter):
|
||||
while True:
|
||||
try:
|
||||
data = src.recv(65536)
|
||||
except OSError:
|
||||
return
|
||||
if not data:
|
||||
try:
|
||||
dst.shutdown(socket.SHUT_WR)
|
||||
except OSError:
|
||||
pass
|
||||
return
|
||||
try:
|
||||
dst.sendall(data)
|
||||
except OSError:
|
||||
return
|
||||
counter[0] += len(data)
|
||||
|
||||
def run(self, cmd):
|
||||
"""Forward one client run (the full command list) to the real server and
|
||||
return the CompletedProcess after the counts have settled."""
|
||||
|
||||
def serve():
|
||||
try:
|
||||
client_sock, _ = self._listener.accept()
|
||||
server_sock = socket.create_connection(("127.0.0.1", self.target_port),
|
||||
timeout=10)
|
||||
except OSError:
|
||||
return
|
||||
c2s, s2c = [0], [0]
|
||||
a = threading.Thread(target=self._pump, args=(client_sock, server_sock, c2s))
|
||||
b = threading.Thread(target=self._pump, args=(server_sock, client_sock, s2c))
|
||||
a.start()
|
||||
b.start()
|
||||
a.join()
|
||||
b.join()
|
||||
self.client_to_server = c2s[0]
|
||||
self.server_to_client = s2c[0]
|
||||
|
||||
thread = threading.Thread(target=serve)
|
||||
thread.start()
|
||||
result = subprocess.run(cmd, capture_output=True, text=True, timeout=180)
|
||||
thread.join(20)
|
||||
return result
|
||||
|
||||
|
||||
def run_client(source_dir, dest_dir, flags=None, port=None, extra_args=None):
|
||||
"""Run the client and return (result, duration)."""
|
||||
cmd = CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir, "--save-to-disk"]
|
||||
|
||||
Reference in New Issue
Block a user