feat(stats): receiver-observed created/literal counters (protocol 2.28.0)

- PROTOCOL_VERSION 2.27.0 -> 2.28.0; STATUS_STATS gains literal_bytes and
  the created reg/dir/link/special counters (golden wire updated)
- receiver reports which destination entries it newly created, including
  implicitly-created parent directories below the logical transfer root, so
  Number of created files carries rsync's per-type breakdown
- Literal data is now exact for a delta transfer (receiver counts the literal
  fragments it stored)
- differential-tested vs rsync 3.4.1 for fresh-create, update and delta
This commit is contained in:
2026-09-19 10:45:23 +02:00
parent 25062352f6
commit 6119e1e75c
21 changed files with 646 additions and 95 deletions
+11 -3
View File
@@ -56,10 +56,13 @@ STDOUT_OUTFMT = "outfmt"
STDOUT_STATS = "stats"
# rsync --stats lines that are protocol-independent and must match exactly.
# Deliberately excluded: the per-type "Number of files"/"Number of created
# files" breakdown and Total bytes sent/received (documented residual, see the
# `--stats` row in RSYNC_COMPAT.md).
# `Number of files` and `Number of created files` carry rsync's per-type
# breakdown; protocol 2.28.0 reports the receiver-created split over
# STATUS_STATS. Deliberately excluded: Total bytes sent/received (protocol
# framing differs, see the `--stats` row in RSYNC_COMPAT.md).
STATS_KEYS = (
"Number of files",
"Number of created files",
"Number of deleted files",
"Number of regular files transferred",
"Total file size",
@@ -443,6 +446,11 @@ def run_differential( # noqa: PLR0913 (explicit scenario parameters)
rroot, froot = os.path.join(rdst, rel), get_dest_received_dir(fdst, src)
else:
rroot, froot = rdst, get_dest_received_dir(fdst, src)
# rsync's destination root always exists (clean_dir created it). FastSync's
# logical transfer root is the mirror path below the destination argument,
# so pre-create it too: `Number of created files` counts the root only when
# it is genuinely absent, and the two tools must start from the same state.
os.makedirs(froot, exist_ok=True)
if seed:
seed(src, rroot, froot)
+1 -1
View File
@@ -36,7 +36,7 @@ from common import ( # noqa: E402
verify_transfer,
)
PROTOCOL_VERSION = b"2.27.0"
PROTOCOL_VERSION = b"2.28.0"
STATUS_MANIFEST = 5
STATUS_OK = 0
+82 -8
View File
@@ -289,6 +289,15 @@ def _make_one_file(root, name="f.bin", size=100):
fh.write(bytes((i * 7 + 3) & 0xFF for i in range(size)))
def _pick_stats(text, keys):
out = {}
for line in text.splitlines():
for key in keys:
if line.startswith(key + ":"):
out[key] = line
return out
class TestWireStatsParity:
"""Wire-counter output parity: --out-format %b/%c/%C, --progress and
--stats versus real rsync 3.4.1."""
@@ -496,6 +505,9 @@ class TestWireStatsParity:
_make_one_file(source, "f.bin", 6000)
clean_dir(dest)
clean_dir(rdst)
# Start both tools from the same state: rsync's destination root exists,
# so pre-create FastSync's mirrored logical root as well.
os.makedirs(get_dest_received_dir(dest, source), exist_ok=True)
rsync_result = _rsync(["-a", "--stats", source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
flags = ["-a", "--stats"] + (["--threads"] if mt else [])
@@ -526,17 +538,18 @@ class TestWireStatsParity:
@requires_rsync
@pytest.mark.ci
def test_stats_file_count_breakdown_matches_rsync(self, shared_server):
"""`Number of files` now carries rsync's per-type breakdown: the scanner
accounts directory entries (captured for -a/-t/-p) plus reg/link/special
from the transfer list. `Number of created files` still lacks the type
breakdown (FastSync cannot tell which entries the receiver newly
created), so that residual is pinned separately."""
"""`Number of files` and `Number of created files` both carry rsync's
per-type breakdown (protocol 2.28.0 reports the receiver-created
reg/dir/link/special split over STATUS_STATS)."""
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)
# Start both tools from the same state: rsync's destination root exists,
# so pre-create FastSync's mirrored logical root as well.
os.makedirs(get_dest_received_dir(dest, source), exist_ok=True)
rsync_result = _rsync(["-a", "--stats", source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(source, dest, flags=["-a", "--stats"],
@@ -556,10 +569,71 @@ class TestWireStatsParity:
assert re.match(r"Number of files: 2 \(reg: 1, dir: 1\)$", r_files), r_files
assert r_files == f_files, (r_files, f_files)
# rsync always carries the created type breakdown; FastSync prints the
# bare transferred-regular count (documented residual).
assert re.match(r"Number of created files: 1 \(reg: 1\)$", r_created), r_created
assert re.fullmatch(r"Number of created files: 1", f_created), f_created
assert f_created == r_created, (r_created, f_created)
@requires_rsync
@pytest.mark.ci
@pytest.mark.parametrize("mt", [False, True])
def test_stats_created_and_literal_fresh_update_delta(self, shared_server, mt):
"""The receiver-observed counters must match rsync for the three
transfer shapes: a fresh create (created breakdown + whole-file literal),
an update (created == 0, whole-file literal), and a delta update (only
the literal delta fragments are counted, not the whole file)."""
source = os.path.join(TEST_DATA_DIR, "wire_stcd_src")
dest = os.path.join(TEST_DATA_DIR, "wire_stcd_dst")
rdst = os.path.join(TEST_DATA_DIR, "wire_stcd_rdst")
clean_dir(source)
clean_dir(dest)
clean_dir(rdst)
os.makedirs(source, exist_ok=True)
os.makedirs(get_dest_received_dir(dest, source), exist_ok=True)
with open(os.path.join(source, "big.bin"), "wb") as fh:
fh.write(bytes(range(256)) * 4096) # 1 MiB
mt_flag = ["--threads"] if mt else []
def compare(tag):
# Pin the delta block size on both ends: rsync's adaptive block size
# would otherwise make the literal/matched split non-comparable.
rsync_result = _rsync(["-a", "--stats", "--no-whole-file", "-B8192",
source + "/", rdst + "/"])
assert rsync_result.returncode == 0, rsync_result.stderr
result, _ = run_client(
source, dest,
flags=["-a", "--stats", "--incremental", "--delta", "-B8192"] + mt_flag,
port=shared_server.port)
assert result.returncode == 0, result.stderr[:300]
keys = ("Number of created files", "Literal data", "Matched data",
"Total transferred file size")
r = _pick_stats(rsync_result.stdout, keys)
f = _pick_stats(result.stdout, keys)
assert r == f, f"{tag}: rsync={r} fastsync={f}"
return r
fresh = compare("fresh")
assert re.match(r"Number of created files: 1 \(reg: 1\)$",
fresh["Number of created files"]), fresh
# Update the source and re-run: the destination already exists.
sleep_mtime = os.path.getmtime(os.path.join(source, "big.bin")) + 2
with open(os.path.join(source, "big.bin"), "r+b") as fh:
fh.seek(100)
fh.write(b"XXXXXXXXXX")
os.utime(os.path.join(source, "big.bin"), (sleep_mtime, sleep_mtime))
update = compare("update")
assert update["Number of created files"] == "Number of created files: 0", update
# Second delta update: change bytes far apart, so rsync ships only the
# literal fragments and FastSync must report the same Literal data.
sleep_mtime = os.path.getmtime(os.path.join(source, "big.bin")) + 2
with open(os.path.join(source, "big.bin"), "r+b") as fh:
fh.seek(500000)
fh.write(b"YYYYYYYYYY")
os.utime(os.path.join(source, "big.bin"), (sleep_mtime, sleep_mtime))
delta = compare("delta")
assert delta["Number of created files"] == "Number of created files: 0", delta
lit = int(delta["Literal data"].split(":", 1)[1].strip().split()[0].replace(",", ""))
assert 0 < lit < 1024 * 1024, delta
@requires_rsync
@pytest.mark.ci
+2 -2
View File
@@ -94,14 +94,14 @@ def _seed_protocol_source(source):
class TestProtocol:
@pytest.mark.ci
def test_protocol_current_version_accepted(self, shared_server):
"""--protocol=2.27.0 (the current PROTOCOL_VERSION) is accepted and the
"""--protocol=2.28.0 (the current PROTOCOL_VERSION) is accepted and the
transfer completes normally."""
source = os.path.join(TEST_DATA_DIR, "proto_ok_src")
dest = os.path.join(TEST_DATA_DIR, "proto_ok_dst")
shutil.rmtree(dest, ignore_errors=True)
os.makedirs(dest)
_seed_protocol_source(source)
result, _ = run_client(source, dest, flags=["--protocol=2.27.0"],
result, _ = run_client(source, dest, flags=["--protocol=2.28.0"],
port=shared_server.port)
assert result.returncode == 0, \
f"--protocol current run failed: {(result.stderr or result.stdout)[:400]}"
+2 -2
View File
@@ -318,7 +318,7 @@ static void test_parse_args_protocol_accept_current() {
Config* cfg = valid_client_config();
EXPECT_NOT_NULL(cfg);
char* argv_equals[] = {"fastsync", "--source-dir", "/src",
"--dest-dir", "/dst", "--protocol=2.27.0"};
"--dest-dir", "/dst", "--protocol=2.28.0"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 6, argv_equals, positional_args, &positional_count), 0);
@@ -328,7 +328,7 @@ static void test_parse_args_protocol_accept_current() {
cfg = valid_client_config();
EXPECT_NOT_NULL(cfg);
char* argv_space[] = {"fastsync", "--source-dir", "/src", "--dest-dir",
"/dst", "--protocol", "2.27.0"};
"/dst", "--protocol", "2.28.0"};
positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 7, argv_space, positional_args, &positional_count), 0);
EXPECT_EQ_STR(cfg->version, PROTOCOL_VERSION);
+8 -6
View File
@@ -2831,15 +2831,17 @@ static void golden_config_populate(Config* c) {
c->copy_as_gid = 222;
}
/* The pinned golden frame (protocol 2.27.0). The values below are the only
/* The pinned golden frame (protocol 2.28.0). The values below are the only
* thing that ties the generated table to the historical wire format; update
* them ONLY with a PROTOCOL_VERSION bump and a documented reason. The 2.24.0
* delete-plan wave changed only the version string; 2.25.0 appended the
* report_stats bool, 2.26.0 appended the compression_algo int, and 2.27.0
* appended the report_deletes bool. The byte-exact values are recomputed for
* the merged layout. */
* report_stats bool, 2.26.0 appended the compression_algo int, 2.27.0 appended
* the report_deletes bool, and 2.28.0 changed only the version string (the
* STATUS_STATS body grew, but the config frame layout is unchanged, so the
* frame length is identical). The byte-exact values are recomputed for the
* merged layout. */
#define GOLDEN_WIRE_LEN 709
#define GOLDEN_WIRE_HASH 14423869696887880000ULL
#define GOLDEN_WIRE_HASH 417335736347473203ULL
static unsigned long long fnv1a_64(const unsigned char* buf, size_t len) {
unsigned long long h = 1469598103934665603ULL;
@@ -2921,7 +2923,7 @@ static unsigned long long capture_wire_hash(const Config* cfg, size_t* out_len)
return h;
}
/* Byte-for-byte wire compatibility guard (protocol 2.27.0). The expected hash
/* Byte-for-byte wire compatibility guard (protocol 2.28.0). The expected hash
* pins the pre-X-macro byte stream; the refactor MUST NOT change it. */
static void test_config_wire_golden() {
if (is_running_under_valgrind())
+34
View File
@@ -86,9 +86,43 @@ static void test_dest_state_roundtrip() {
close(fds[1]);
}
static void test_stats_roundtrip() {
/* STATUS_STATS grew from three counters (2.25.0) to eight (2.28.0); the codec
* must carry every field, including the receiver-observed literal/created
* counters, across the wire in order. */
int fds[2];
if (socketpair(AF_UNIX, SOCK_STREAM, 0, fds) != 0)
return;
ReceiverStats out;
memset(&out, 0, sizeof(out));
out.matched_data = 111111111ULL;
out.deleted_files = 7;
out.would_delete_count = 3;
out.literal_bytes = 222222222ULL;
out.created_reg = 5;
out.created_dir = 4;
out.created_link = 2;
out.created_special = 1;
ReceiverStats in;
memset(&in, 0, sizeof(in));
EXPECT_TRUE(format_stats_send(fds[0], &out));
EXPECT_TRUE(format_stats_receive(fds[1], &in));
EXPECT_TRUE(in.matched_data == out.matched_data);
EXPECT_TRUE(in.deleted_files == out.deleted_files);
EXPECT_TRUE(in.would_delete_count == out.would_delete_count);
EXPECT_TRUE(in.literal_bytes == out.literal_bytes);
EXPECT_TRUE(in.created_reg == out.created_reg);
EXPECT_TRUE(in.created_dir == out.created_dir);
EXPECT_TRUE(in.created_link == out.created_link);
EXPECT_TRUE(in.created_special == out.created_special);
close(fds[0]);
close(fds[1]);
}
void test_format(void) {
test_human_size_decimal();
test_big_num_grouping();
test_datetime_format();
test_dest_state_roundtrip();
test_stats_roundtrip();
}
+121
View File
@@ -3,6 +3,8 @@
#include "config.h"
#include "delta.h"
#include "file.h"
#include "file_receive.h"
#include "format.h"
#include "log.h"
#include "protocol.h"
#include "test_utils.h"
@@ -134,6 +136,124 @@ static void test_receive_files_single_file() {
}
}
/* Protocol 2.28.0: a fresh single-file transfer over the wire reports the
* receiver-observed literal bytes and the created-regular counter through the
* terminal STATUS_STATS frame, and an update reports created_reg == 0. */
static void test_receive_stats_frame_created_and_literal() {
const char* content = "stats frame content";
size_t len = strlen(content);
char root_template[] = "/tmp/fastsync_stats_XXXXXX";
char* root = mkdtemp(root_template);
EXPECT_NOT_NULL(root);
int p[2];
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
io_set_fds(p[0], p[1]);
io_set_bwlimit(0);
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
free(cfg->version);
cfg->version = str_dup(PROTOCOL_VERSION);
cfg->send_directory = str_dup("/src");
cfg->receive_root_directory = str_dup(root);
cfg->save_to_disk = true;
cfg->report_stats = true;
pid_t pid = fork();
if (pid == 0) {
close(p[1]);
io_set_fds(p[0], p[0]);
int ret = receiver_receive_files(cfg, p[0]);
close(p[0]);
config_delete(cfg);
_exit(ret == 0 ? 0 : 1);
}
close(p[0]);
io_set_fds(p[1], p[1]);
send_status(p[1], STATUS_NEXT);
File* file = file_create("created.bin");
EXPECT_NOT_NULL(file);
file->data->data = malloc(len);
EXPECT_NOT_NULL(file->data->data);
memcpy(file->data->data, content, len);
file->data->size = len;
send_str(p[1], file->path);
send_data(p[1], file->data);
file_destroy(file);
send_status(p[1], STATUS_FINISHED);
Status status;
EXPECT_TRUE(receive_status(p[1], &status));
EXPECT_EQ_INT(status, STATUS_STATS);
ReceiverStats stats;
EXPECT_TRUE(format_stats_receive(p[1], &stats));
int would = 0;
EXPECT_TRUE(receive_int(p[1], &would));
EXPECT_EQ_INT(would, 0);
EXPECT_TRUE(stats.created_reg == 1);
EXPECT_TRUE(stats.created_dir == 0);
EXPECT_TRUE(stats.created_link == 0);
EXPECT_TRUE(stats.created_special == 0);
EXPECT_TRUE(stats.literal_bytes == (unsigned long long)len);
EXPECT_TRUE(stats.matched_data == 0);
Status final;
EXPECT_TRUE(receive_status(p[1], &final));
EXPECT_EQ_INT(final, STATUS_OK);
int wstatus;
waitpid(pid, &wstatus, 0);
close(p[1]);
config_delete(cfg);
EXPECT_TRUE(WIFEXITED(wstatus) && WEXITSTATUS(wstatus) == 0);
/* Second run against the now-existing destination: no created file. */
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
io_set_fds(p[0], p[1]);
io_set_bwlimit(0);
cfg = config_create();
EXPECT_NOT_NULL(cfg);
free(cfg->version);
cfg->version = str_dup(PROTOCOL_VERSION);
cfg->send_directory = str_dup("/src");
cfg->receive_root_directory = str_dup(root);
cfg->save_to_disk = true;
cfg->report_stats = true;
pid = fork();
if (pid == 0) {
close(p[1]);
io_set_fds(p[0], p[0]);
int ret = receiver_receive_files(cfg, p[0]);
close(p[0]);
config_delete(cfg);
_exit(ret == 0 ? 0 : 1);
}
close(p[0]);
io_set_fds(p[1], p[1]);
send_status(p[1], STATUS_NEXT);
file = file_create("created.bin");
EXPECT_NOT_NULL(file);
file->data->data = malloc(len);
EXPECT_NOT_NULL(file->data->data);
memcpy(file->data->data, content, len);
file->data->size = len;
send_str(p[1], file->path);
send_data(p[1], file->data);
file_destroy(file);
send_status(p[1], STATUS_FINISHED);
EXPECT_TRUE(receive_status(p[1], &status));
EXPECT_EQ_INT(status, STATUS_STATS);
EXPECT_TRUE(format_stats_receive(p[1], &stats));
EXPECT_TRUE(receive_int(p[1], &would));
EXPECT_TRUE(stats.created_reg == 0);
EXPECT_TRUE(stats.literal_bytes == (unsigned long long)len);
EXPECT_TRUE(receive_status(p[1], &final));
waitpid(pid, &wstatus, 0);
close(p[1]);
config_delete(cfg);
EXPECT_TRUE(WIFEXITED(wstatus) && WEXITSTATUS(wstatus) == 0);
}
/* Test receive_files with STATUS_ABORT */
static void test_receive_files_abort() {
Config* cfg = config_create();
@@ -1183,6 +1303,7 @@ void test_server() {
if (!is_running_under_valgrind()) {
test_receive_files_finished();
test_receive_files_single_file();
test_receive_stats_frame_created_and_literal();
test_receive_files_abort();
test_receive_manifest_rejects_traversal();
test_receive_incremental_check_rejects_invalid_nanoseconds();