Merge branch 'main' into metadata-transfer, resolve conflicts

Conflicts resolved:
- config_create signature: added both use_metadata and use_sendfile params
- config_send/config_receive: wire protocol includes both use_metadata and use_sendfile
- client.c: config_create call updated, arg parsing includes both -M/--preserve and -f/--sendfile
- test_config.c: all three config_create calls updated
- file_send_sendfile: use file->data->size instead of removed struct stat
This commit is contained in:
2026-07-05 20:29:04 +02:00
9 changed files with 383 additions and 65 deletions
+12 -2
View File
@@ -12,6 +12,7 @@ FastFileTransfer is a C implementation of a file synchronization system that:
4. Utilizes multithreading for parallel file processing 4. Utilizes multithreading for parallel file processing
5. Implements producer-consumer patterns with thread-safe queues 5. Implements producer-consumer patterns with thread-safe queues
6. Provides both in-memory and disk-based storage options 6. Provides both in-memory and disk-based storage options
7. Supports `sendfile()` for zero-copy file transfer
## System Architecture ## System Architecture
@@ -23,6 +24,7 @@ The system consists of two main components:
- Compresses data using zstd algorithm - Compresses data using zstd algorithm
- Serializes chunks into a compact binary format for batch transfer - Serializes chunks into a compact binary format for batch transfer
- Sends files to server using custom protocol - Sends files to server using custom protocol
- Supports sendfile for zero-copy file transfer (`-f`)
- Supports both single-threaded and multi-threaded operation - Supports both single-threaded and multi-threaded operation
### Server ### Server
@@ -49,6 +51,7 @@ The client-server communication uses the following status codes:
| `-m` | Enable multithreading mode | | `-m` | Enable multithreading mode |
| `-c [level]` | Enable compression with optional level (1-22, default: 5) | | `-c [level]` | Enable compression with optional level (1-22, default: 5) |
| `-s` | Enable chunk serialization (batch-transfer all files per chunk) | | `-s` | Enable chunk serialization (batch-transfer all files per chunk) |
| `-f` | Enable sendfile (zero-copy file transfer, bypasses userspace memory). Can be combined with `-m`. Incompatible with `-c` and `-s`. |
| `--source-dir <path>` | Source directory to sync (overrides `FASTSYNC_SOURCE_DIR`) | | `--source-dir <path>` | Source directory to sync (overrides `FASTSYNC_SOURCE_DIR`) |
| `--dest-dir <path>` | Server-side destination directory (overrides `FASTSYNC_DEST_DIR`) | | `--dest-dir <path>` | Server-side destination directory (overrides `FASTSYNC_DEST_DIR`) |
| `--save-to-disk` | Persist received files to disk | | `--save-to-disk` | Persist received files to disk |
@@ -118,6 +121,12 @@ make
# Multithreaded with compressed chunk serialization # Multithreaded with compressed chunk serialization
./build/client -m -s -c 3 ./build/client -m -s -c 3
# Sendfile (zero-copy, bypasses userspace for large files)
./build/client -f
# Sendfile with multithreading
./build/client -f -m
``` ```
## Testing ## Testing
@@ -150,8 +159,9 @@ tests/ # Unit tests
1. Chunk size (10MB default) affects memory usage and transfer efficiency 1. Chunk size (10MB default) affects memory usage and transfer efficiency
2. Compression level (1-22) trades CPU usage for space savings 2. Compression level (1-22) trades CPU usage for space savings
3. Multithreading improves performance on multi-core systems 3. `sendfile()` (`-f`) bypasses userspace memory, ~2x faster on localhost for large files
4. Thread-safe queues minimize contention between producer/consumer threads 4. Multithreading improves performance on multi-core systems
5. Thread-safe queues minimize contention between producer/consumer threads
## Extensibility ## Extensibility
+1 -3
View File
@@ -4,9 +4,7 @@
"bash": { "bash": {
"*": "allow", "*": "allow",
"git push origin main": "deny", "git push origin main": "deny",
"git push origin master": "deny", "git push main": "ask"
"git push *": "ask",
"git commit *": "ask"
} }
} }
} }
+22 -5
View File
@@ -29,6 +29,11 @@ int send_chunk(Client *client, Chunk *chunk, Config *config) {
} }
send_data(client->file_descriptor, data->data, data->size); send_data(client->file_descriptor, data->data, data->size);
data_destroy(data); data_destroy(data);
} else if (config->use_sendfile && !config->use_compression) {
for (int i = 0; i < chunk->element_count; i++) {
send_status(client->file_descriptor, STATUS_NEXT);
file_send_sendfile(chunk->items[i], client->file_descriptor);
}
} else { } else {
for (int i = 0; i < chunk->element_count; i++) { for (int i = 0; i < chunk->element_count; i++) {
send_status(client->file_descriptor, STATUS_NEXT); send_status(client->file_descriptor, STATUS_NEXT);
@@ -81,8 +86,10 @@ int load_files_multithreaded(void *pipeline_context) {
mtx_unlock(&context->mutex_loader); mtx_unlock(&context->mutex_loader);
return thrd_success; return thrd_success;
} }
for (int i = 0; i < chunk->element_count; i++) if (!context->config->use_sendfile) {
file_load_data(chunk->items[i]); for (int i = 0; i < chunk->element_count; i++)
file_load_data(chunk->items[i]);
}
queue_enqueue_multithreaded(context->queue_loader, chunk, queue_enqueue_multithreaded(context->queue_loader, chunk,
&context->mutex_loader, &context->mutex_loader,
&context->condition_not_empty_loader, &context->condition_not_empty_loader,
@@ -130,8 +137,10 @@ int send_files(Config *config) {
DirectoryScanner *scanner = directory_scanner_create(config->send_directory, config->use_metadata); DirectoryScanner *scanner = directory_scanner_create(config->send_directory, config->use_metadata);
Chunk *current_chunk; Chunk *current_chunk;
while ((current_chunk = directory_scanner_next(scanner)) != NULL) { while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
for (int i = 0; i < current_chunk->element_count; i++) if (!config->use_sendfile) {
file_load_data(current_chunk->items[i]); for (int i = 0; i < current_chunk->element_count; i++)
file_load_data(current_chunk->items[i]);
}
send_chunk(client, current_chunk, config); send_chunk(client, current_chunk, config);
chunk_destroy(current_chunk); chunk_destroy(current_chunk);
} }
@@ -192,7 +201,7 @@ int main(int argc, char *argv[]) {
} }
Config *config = config_create(str_dup("1.0.0"), source_dir, dest_dir, Config *config = config_create(str_dup("1.0.0"), source_dir, dest_dir,
save_to_disk, false, false, false, false, 5, 20); save_to_disk, false, false, false, false, 5, 20, false);
for (int i = 1; i < argc; i++) { for (int i = 1; i < argc; i++) {
if (strcmp(argv[i], "-c") == 0) { if (strcmp(argv[i], "-c") == 0) {
config->use_compression = true; config->use_compression = true;
@@ -218,6 +227,9 @@ int main(int argc, char *argv[]) {
} else if (strcmp(argv[i], "-M") == 0 || strcmp(argv[i], "--preserve") == 0) { } else if (strcmp(argv[i], "-M") == 0 || strcmp(argv[i], "--preserve") == 0) {
config->use_metadata = true; config->use_metadata = true;
log_message(LOG_LEVEL_INFO, "Enabled metadata preservation"); log_message(LOG_LEVEL_INFO, "Enabled metadata preservation");
} else if (strcmp(argv[i], "-f") == 0 || strcmp(argv[i], "--sendfile") == 0) {
config->use_sendfile = true;
log_message(LOG_LEVEL_INFO, "Enabled sendfile");
} else { } else {
handle_arg(argv[i], "-m", &config->use_multithreading, handle_arg(argv[i], "-m", &config->use_multithreading,
"Enabled Multithreading"); "Enabled Multithreading");
@@ -226,6 +238,11 @@ int main(int argc, char *argv[]) {
} }
} }
if (config->use_sendfile && (config->use_chunk_serialization || config->use_compression)) {
fprintf(stderr, "Error: -f/--sendfile cannot be combined with -c (compression) or -s (chunk serialization)\n");
return 1;
}
if (config->use_multithreading) if (config->use_multithreading)
return send_files_multithreaded(config); return send_files_multithreaded(config);
return send_files(config); return send_files(config);
+5 -2
View File
@@ -7,8 +7,8 @@
Config *config_create(char *version, char *send_directory, Config *config_create(char *version, char *send_directory,
char *receive_directory, bool save_to_disk, char *receive_directory, bool save_to_disk,
bool use_multithreading, bool use_chunk_serialization, bool use_multithreading, bool use_chunk_serialization,
bool use_compression, bool use_metadata, bool use_compression, bool use_metadata,
int compression_level, int num_connections) { int compression_level, int num_connections, bool use_sendfile) {
Config *config = malloc(sizeof(Config)); Config *config = malloc(sizeof(Config));
config->version = version; config->version = version;
@@ -21,6 +21,7 @@ Config *config_create(char *version, char *send_directory,
config->use_metadata = use_metadata; config->use_metadata = use_metadata;
config->compression_level = compression_level; config->compression_level = compression_level;
config->num_connections = num_connections; config->num_connections = num_connections;
config->use_sendfile = use_sendfile;
return config; return config;
} }
@@ -42,6 +43,7 @@ void config_send(int file_descriptor, Config *config) {
send_int(file_descriptor, config->use_metadata); send_int(file_descriptor, config->use_metadata);
send_int(file_descriptor, config->compression_level); send_int(file_descriptor, config->compression_level);
send_int(file_descriptor, config->num_connections); send_int(file_descriptor, config->num_connections);
send_int(file_descriptor, config->use_sendfile);
if (receive_status(file_descriptor) != STATUS_OK) { if (receive_status(file_descriptor) != STATUS_OK) {
perror("Error transmitting config!"); perror("Error transmitting config!");
exit(EXIT_FAILURE); exit(EXIT_FAILURE);
@@ -60,6 +62,7 @@ Config *config_receive(int file_descriptor) {
config->use_metadata = receive_int(file_descriptor); config->use_metadata = receive_int(file_descriptor);
config->compression_level = receive_int(file_descriptor); config->compression_level = receive_int(file_descriptor);
config->num_connections = receive_int(file_descriptor); config->num_connections = receive_int(file_descriptor);
config->use_sendfile = receive_int(file_descriptor);
send_status(file_descriptor, STATUS_OK); send_status(file_descriptor, STATUS_OK);
return config; return config;
} }
+3 -2
View File
@@ -11,6 +11,7 @@ typedef struct Config {
bool use_multithreading; bool use_multithreading;
bool use_chunk_serialization; bool use_chunk_serialization;
bool use_compression; bool use_compression;
bool use_sendfile;
bool use_single_send_per_file; bool use_single_send_per_file;
bool use_metadata; bool use_metadata;
int compression_level; int compression_level;
@@ -20,8 +21,8 @@ typedef struct Config {
Config *config_create(char *version, char *send_directory, Config *config_create(char *version, char *send_directory,
char *receive_directory, bool save_to_disk, char *receive_directory, bool save_to_disk,
bool use_multithreading, bool use_chunk_serialization, bool use_multithreading, bool use_chunk_serialization,
bool use_compression, bool use_metadata, bool use_compression, bool use_metadata,
int compression_level, int num_connections); int compression_level, int num_connections, bool use_sendfile);
void config_delete(Config *config); void config_delete(Config *config);
void config_send(int file_descriptor, Config *config); void config_send(int file_descriptor, Config *config);
Config *config_receive(int file_descriptor); Config *config_receive(int file_descriptor);
+28
View File
@@ -1,9 +1,12 @@
#include <dirent.h> #include <dirent.h>
#include <fcntl.h>
#include <libgen.h> #include <libgen.h>
#include <stddef.h> #include <stddef.h>
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <sys/sendfile.h>
#include <unistd.h>
#include <zstd.h> #include <zstd.h>
#include "data.h" #include "data.h"
@@ -122,6 +125,31 @@ void file_send_single_calls(File *file, int file_descriptor) {
send_data(file_descriptor, file->data->data, file->data->size); send_data(file_descriptor, file->data->data, file->data->size);
} }
void file_send_sendfile(File *file, int file_descriptor) {
send_str(file_descriptor, file->path);
int fd = open(file->path, O_RDONLY);
if (fd == -1) {
perror("Could not open file for sendfile");
exit(EXIT_FAILURE);
}
unsigned long long file_size = file->data->size;
send_n_data(file_descriptor, &file_size, sizeof(unsigned long long));
off_t offset = 0;
while (offset < file_size) {
ssize_t sent = sendfile(file_descriptor, fd, &offset, file_size - offset);
if (sent == -1) {
perror("sendfile failed");
close(fd);
exit(EXIT_FAILURE);
}
}
close(fd);
}
size_t file_content_to_buffer(File *file) { size_t file_content_to_buffer(File *file) {
FILE *file_pointer = fopen(file->path, "rb"); FILE *file_pointer = fopen(file->path, "rb");
if (file_pointer == NULL) { if (file_pointer == NULL) {
+1
View File
@@ -23,6 +23,7 @@ void file_destroy(void *item);
void file_load_data(File *file); void file_load_data(File *file);
void file_print(void *item); void file_print(void *item);
void file_send_single_calls(File *file, int file_descriptor); void file_send_single_calls(File *file, int file_descriptor);
void file_send_sendfile(File *file, int file_descriptor);
size_t file_content_to_buffer(File *file); size_t file_content_to_buffer(File *file);
FileMetadata *file_metadata_create(struct stat *stats); FileMetadata *file_metadata_create(struct stat *stats);
void file_metadata_destroy(void *metadata); void file_metadata_destroy(void *metadata);
+308 -48
View File
@@ -1,11 +1,13 @@
import argparse import argparse
import filecmp import filecmp
import os import os
import re
import shutil import shutil
import subprocess import subprocess
import sys import sys
import tempfile import tempfile
import time import time
import socket
TEST_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "test_data") TEST_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "test_data")
DEFAULT_SOURCE_DIR = os.path.join(TEST_DIR, "source") DEFAULT_SOURCE_DIR = os.path.join(TEST_DIR, "source")
@@ -17,14 +19,23 @@ base_client_cmd = ["./build/client"]
DISK_DEVICE = "/dev/nvme0n1p5" DISK_DEVICE = "/dev/nvme0n1p5"
READ_BPS_MAX = "15M" READ_BPS_MAX = "15M"
WRITE_BPS_MAX = "10M" WRITE_BPS_MAX = "10M"
NET_LIMIT = "100mbit"
NET_DELAY = "100ms"
NETWORK_INTERFACE = "lo" NETWORK_INTERFACE = "lo"
NET_LIMIT_CMD = (
f"sudo tc qdisc add dev {NETWORK_INTERFACE} root netem rate {NET_LIMIT} delay {NET_DELAY}".split() NETWORK_PROFILES = {
) "Unlimited": {},
NET_RESET_CMD = f"sudo tc qdisc del dev {NETWORK_INTERFACE} root".split() "LAN": {
"rate": "1000mbit",
"delay": "20ms",
"jitter": "1ms",
"loss": "0.1%",
},
"WAN": {
"rate": "100mbit",
"delay": "50ms",
"jitter": "10ms",
"loss": "1%",
},
}
CLIENT_CMD_PREFIX = [ CLIENT_CMD_PREFIX = [
"sudo", "sudo",
@@ -37,7 +48,7 @@ CLIENT_CMD_PREFIX = [
] ]
TEST_CASES = [ TEST_CASES = [
{"name": "Standard (Single-threaded)", "flags": []}, {"name": "Standard", "flags": []},
{"name": "Multithreading (-m)", "flags": ["-m"]}, {"name": "Multithreading (-m)", "flags": ["-m"]},
{"name": "Compression (-c)", "flags": ["-c"]}, {"name": "Compression (-c)", "flags": ["-c"]},
{"name": "Chunk Serialization (-s)", "flags": ["-s"]}, {"name": "Chunk Serialization (-s)", "flags": ["-s"]},
@@ -48,8 +59,41 @@ TEST_CASES = [
"name": "Multithreading + Compression + Chunk Serialization (-m -c -s)", "name": "Multithreading + Compression + Chunk Serialization (-m -c -s)",
"flags": ["-m", "-c", "-s"], "flags": ["-m", "-c", "-s"],
}, },
{"name": "Sendfile (-f)", "flags": ["-f"]},
{"name": "Sendfile + Multithreading (-f -m)", "flags": ["-f", "-m"]},
] ]
RSYNC_CASES = [
{"name": "rsync (archive)", "args": ["-aH"]},
{"name": "rsync (archive + compress)", "args": ["-aHz"]},
]
def netem_apply(profile):
params = NETWORK_PROFILES[profile]
if not params:
netem_reset()
return
netem_reset()
cmd = ["sudo", "tc", "qdisc", "add", "dev", NETWORK_INTERFACE, "root", "netem"]
cmd += ["rate", params["rate"]]
cmd += ["delay", params["delay"], params["jitter"]]
cmd += ["loss", params["loss"]]
subprocess.run(cmd, check=True, capture_output=True)
def netem_reset():
subprocess.run(
f"sudo tc qdisc del dev {NETWORK_INTERFACE} root".split(),
capture_output=True,
)
def find_free_port():
with socket.socket(socket.AF_INET, socket.SOCK_STREAM) as s:
s.bind(('', 0))
return s.getsockname()[1]
def generate_test_files(source_dir): def generate_test_files(source_dir):
if os.path.exists(source_dir): if os.path.exists(source_dir):
@@ -86,14 +130,14 @@ def generate_test_files(source_dir):
total_mb = written / (1024 * 1024) total_mb = written / (1024 * 1024)
print(f" Generated {total_mb:.1f} MB of test data in {source_dir}") print(f" Generated {total_mb:.1f} MB of test data in {source_dir}")
return written
def verify_transfer(source_dir, dest_dir): def verify_transfer(source_dir, received_dir):
source_dir = os.path.abspath(source_dir) source_dir = os.path.abspath(source_dir)
dest_dir = os.path.abspath(dest_dir) received_dir = os.path.abspath(received_dir)
received_prefix = os.path.join(dest_dir, source_dir.lstrip(os.sep)) if not os.path.exists(received_dir):
if not os.path.exists(received_prefix):
return [], ["no received files found"] return [], ["no received files found"]
mismatches = [] mismatches = []
@@ -103,7 +147,7 @@ def verify_transfer(source_dir, dest_dir):
for f in files: for f in files:
src_path = os.path.join(root, f) src_path = os.path.join(root, f)
rel = os.path.relpath(src_path, source_dir) rel = os.path.relpath(src_path, source_dir)
dst_path = os.path.join(received_prefix, rel) dst_path = os.path.join(received_dir, rel)
if not os.path.exists(dst_path): if not os.path.exists(dst_path):
missing.append(rel) missing.append(rel)
@@ -113,25 +157,28 @@ def verify_transfer(source_dir, dest_dir):
return mismatches, missing return mismatches, missing
def run_suite(env_name, apply_limits, source_dir, dest_dir): def run_profile(profile_name, source_dir, dest_dir):
results = [] results = []
is_limited = profile_name != "Unlimited"
params = NETWORK_PROFILES[profile_name]
print(f"\n{'=' * 60}") print(f"\n{'=' * 60}")
print(f"Suite: {env_name}") print(f"Profile: {profile_name}")
print(f"{'=' * 60}") print(f"{'=' * 60}")
if apply_limits: if is_limited:
p = params
print(f" Network: rate={p['rate']}, delay={p['delay']} ±{p['jitter']}, loss={p['loss']}")
print(f" Disk I/O: Reads <= {READ_BPS_MAX}, Writes <= {WRITE_BPS_MAX}") print(f" Disk I/O: Reads <= {READ_BPS_MAX}, Writes <= {WRITE_BPS_MAX}")
print(f" Network: {NET_LIMIT}, {NET_DELAY} delay")
client_prefix = CLIENT_CMD_PREFIX client_prefix = CLIENT_CMD_PREFIX
else: else:
print(" Baseline (no limits)") print(" No limits applied")
client_prefix = [] client_prefix = []
try: try:
if apply_limits: if is_limited:
subprocess.run(NET_LIMIT_CMD, check=True) netem_apply(profile_name)
else: else:
subprocess.run(NET_RESET_CMD, capture_output=True) netem_reset()
for case in TEST_CASES: for case in TEST_CASES:
name = case["name"] name = case["name"]
@@ -175,11 +222,12 @@ def run_suite(env_name, apply_limits, source_dir, dest_dir):
mismatches, missing = [], [] mismatches, missing = [], []
if client_result.returncode == 0: if client_result.returncode == 0:
mismatches, missing = verify_transfer(source_dir, dest_dir) received = os.path.join(dest_dir, os.path.abspath(source_dir).lstrip(os.sep))
mismatches, missing = verify_transfer(source_dir, received)
entry = { entry = {
"name": name, "name": name,
"suite": env_name, "suite": profile_name,
"time": f"{duration:.4f}s" if client_result.returncode == 0 else "N/A", "time": f"{duration:.4f}s" if client_result.returncode == 0 else "N/A",
} }
@@ -210,11 +258,11 @@ def run_suite(env_name, apply_limits, source_dir, dest_dir):
except subprocess.TimeoutExpired: except subprocess.TimeoutExpired:
results.append( results.append(
{"name": name, "suite": env_name, "status": "Timeout", "time": "N/A", "error": "Exceeded 15s"} {"name": name, "suite": profile_name, "status": "Timeout", "time": "N/A", "error": "Exceeded 15s"}
) )
except Exception as e: except Exception as e:
results.append( results.append(
{"name": name, "suite": env_name, "status": "Error", "time": "N/A", "error": str(e)} {"name": name, "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)}
) )
finally: finally:
if server_process: if server_process:
@@ -224,18 +272,175 @@ def run_suite(env_name, apply_limits, source_dir, dest_dir):
server_process.kill() server_process.kill()
server_process.wait() server_process.wait()
except subprocess.CalledProcessError as e: # Rsync tests (over network via daemon, so tc netem applies)
print(f" Error running limit command: {' '.join(e.cmd)}") rsync_port = find_free_port()
finally: rsyncd_conf = os.path.join(tempfile.gettempdir(), f"rsyncd-{rsync_port}.conf")
if apply_limits: with open(rsyncd_conf, "w") as f:
f.write(f"""port = {rsync_port}
read only = yes
[source]
path = {source_dir}
""")
rsync_daemon = None
try:
rsync_daemon = subprocess.Popen(
["rsync", "--daemon", "--no-detach", f"--config={rsyncd_conf}"],
stdout=subprocess.DEVNULL, stderr=None
)
for _ in range(50):
time.sleep(0.1)
if rsync_daemon.poll() is not None:
print(f" rsync daemon exited (rc={rsync_daemon.returncode})")
raise RuntimeError("rsync daemon failed to start")
try:
with socket.create_connection(("127.0.0.1", rsync_port), timeout=0.3):
break
except (ConnectionRefusedError, OSError):
continue
else:
print(f" timed out waiting for rsync daemon on port {rsync_port}")
raise RuntimeError("rsync daemon did not start")
for case in RSYNC_CASES:
name = case["name"]
rsync_args = case["args"]
print(f"\n --- {name} ---")
if os.path.exists(dest_dir):
shutil.rmtree(dest_dir)
try:
rsync_cmd = (
client_prefix
+ ["rsync"]
+ rsync_args
+ [f"rsync://localhost:{rsync_port}/source/", f"{dest_dir}/"]
)
print(f" Running: {' '.join(rsync_cmd)}")
start_time = time.monotonic()
rsync_result = subprocess.run(
rsync_cmd, capture_output=True, text=True, timeout=120
)
end_time = time.monotonic()
duration = end_time - start_time
mismatches, missing = [], []
if rsync_result.returncode == 0:
mismatches, missing = verify_transfer(source_dir, dest_dir)
entry = {
"name": name,
"suite": profile_name,
"time": f"{duration:.4f}s" if rsync_result.returncode == 0 else "N/A",
}
if rsync_result.returncode == 0 and not mismatches and not missing:
entry["status"] = "Success"
entry["error"] = ""
else:
entry["status"] = "Failed"
errors = []
if rsync_result.returncode != 0:
errs = []
for line in (rsync_result.stderr or "").split("\n"):
line = line.strip()
if not line:
continue
if line.startswith("Running as unit:"):
continue
errs.append(line)
for line in (rsync_result.stdout or "").split("\n"):
line = line.strip()
if not line:
continue
errs.append(line)
if not errs:
errs.append("No output")
err = " | ".join(errs[-3:])
errors.append(f"Exit code {rsync_result.returncode}: {err[:200]}")
# Retry without systemd-run to reveal the actual error
if client_prefix:
tmp_dest = tempfile.mkdtemp()
try:
plain = subprocess.run(
["rsync", "-aH", f"rsync://localhost:{rsync_port}/source/", f"{tmp_dest}/"],
capture_output=True, text=True, timeout=30
)
if plain.returncode != 0:
plain_errs = [l for l in (plain.stderr or "").split("\n") if l.strip()]
if plain_errs:
errors.append(f"raw: {plain_errs[-1][:150]}")
finally:
shutil.rmtree(tmp_dest, ignore_errors=True)
if missing:
errors.append(f"Missing ({len(missing)}): {', '.join(missing[:5])}")
if mismatches:
errors.append(f"Mismatch ({len(mismatches)}): {', '.join(mismatches[:3])}")
entry["error"] = " | ".join(errors)
results.append(entry)
except subprocess.TimeoutExpired:
results.append(
{"name": name, "suite": profile_name, "status": "Timeout", "time": "N/A", "error": "Exceeded 120s"}
)
except Exception as e:
results.append(
{"name": name, "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)}
)
finally:
if rsync_daemon:
try:
rsync_daemon.wait(timeout=5)
except subprocess.TimeoutExpired:
rsync_daemon.kill()
rsync_daemon.wait()
try: try:
subprocess.run(NET_RESET_CMD, check=True, capture_output=True) os.unlink(rsyncd_conf)
except Exception:
pass
except subprocess.CalledProcessError as e:
print(f" Error running netem command: {' '.join(e.cmd)}")
finally:
if is_limited:
try:
netem_reset()
except Exception: except Exception:
pass pass
return results return results
def parse_rate_to_bytes_per_sec(rate_str):
m = re.match(r'(\d+)\s*(mbit|gbit|kbit|bit)', rate_str)
if not m:
return None
val = int(m.group(1))
unit = m.group(2)
bits_per_sec = {
'bit': val,
'kbit': val * 1000,
'mbit': val * 1_000_000,
'gbit': val * 1_000_000_000,
}.get(unit)
return bits_per_sec / 8 if bits_per_sec is not None else None
def format_throughput(bps):
if bps >= 1_000_000_000:
return f"{bps/1_000_000_000:.1f} GB/s"
if bps >= 1_000_000:
return f"{bps/1_000_000:.1f} MB/s"
if bps >= 1_000:
return f"{bps/1_000:.1f} KB/s"
return f"{bps:.0f} B/s"
def main(): def main():
parser = argparse.ArgumentParser(description="FastSync integration test / benchmark") parser = argparse.ArgumentParser(description="FastSync integration test / benchmark")
parser.add_argument("--source-dir", default=DEFAULT_SOURCE_DIR, parser.add_argument("--source-dir", default=DEFAULT_SOURCE_DIR,
@@ -244,10 +449,10 @@ def main():
help="Destination directory for received files (default: %(default)s)") help="Destination directory for received files (default: %(default)s)")
parser.add_argument("--keep-data", action="store_true", parser.add_argument("--keep-data", action="store_true",
help="Keep test_data directory after run") help="Keep test_data directory after run")
parser.add_argument("--no-throttled", action="store_true", parser.add_argument("--unlimited", action="store_true",
help="Skip throttled suite (requires sudo)") help="Run Unlimited profile instead of LAN (no network limits)")
parser.add_argument("--no-unlimited", action="store_true", parser.add_argument("--wan", action="store_true",
help="Skip unlimited suite") help="Run WAN profile instead of LAN (100mbit, 50ms, 1% loss)")
args = parser.parse_args() args = parser.parse_args()
os.system("cmake -B build -S . > /dev/null 2>&1") os.system("cmake -B build -S . > /dev/null 2>&1")
@@ -256,33 +461,88 @@ def main():
print("Build failed") print("Build failed")
sys.exit(1) sys.exit(1)
generate_test_files(args.source_dir) total_bytes = generate_test_files(args.source_dir)
if os.path.exists(args.dest_dir):
shutil.rmtree(args.dest_dir)
os.makedirs(args.dest_dir, exist_ok=True) os.makedirs(args.dest_dir, exist_ok=True)
profiles_to_run = []
if args.unlimited:
profiles_to_run.append("Unlimited")
elif args.wan:
profiles_to_run.append("WAN")
else:
profiles_to_run.append("LAN")
try: try:
all_results = [] all_results = []
for profile in profiles_to_run:
if not args.no_unlimited:
all_results.extend( all_results.extend(
run_suite("Unlimited", False, args.source_dir, args.dest_dir) run_profile(profile, args.source_dir, args.dest_dir)
) )
if not args.no_throttled: print("\n" + "=" * 130)
all_results.extend( print(f"{'RESULTS':^130}")
run_suite("Throttled", True, args.source_dir, args.dest_dir) print("=" * 130)
) print(f"{'Configuration':<45} | {'Profile':<12} | {'Status':<8} | {'Time':<10} | {'Details'}")
print("-" * 130)
print("\n" + "=" * 110)
print(f"{'RESULTS':^110}")
print("=" * 110)
print(f"{'Configuration':<45} | {'Suite':<12} | {'Status':<8} | {'Time':<10} | {'Details'}")
print("-" * 110)
for res in all_results: for res in all_results:
print( print(
f"{res['name']:<45} | {res['suite']:<12} | {res['status']:<8} | {res['time']:<10} | {res['error']}" f"{res['name']:<45} | {res['suite']:<12} | {res['status']:<8} | {res['time']:<10} | {res['error']}"
) )
# --- Additional metrics per profile ---
profiles_results = {}
for res in all_results:
profiles_results.setdefault(res["suite"], []).append(res)
for profile_name, results in profiles_results.items():
params = NETWORK_PROFILES.get(profile_name)
if not params or "rate" not in params:
continue
client_times = []
rsync_times = {}
for r in results:
if r["status"] != "Success" or r["time"] == "N/A":
continue
t = float(r["time"].rstrip("s"))
if r["name"].startswith("rsync"):
rsync_times[r["name"]] = t
else:
client_times.append((t, r["name"]))
if not client_times or len(rsync_times) < 2:
continue
best_time, best_name = min(client_times, key=lambda x: x[0])
rate_Bps = parse_rate_to_bytes_per_sec(params["rate"])
theoretical_max_time = None
speedup_vs_theoretical = None
if rate_Bps is not None:
theoretical_max_time = total_bytes / rate_Bps
speedup_vs_theoretical = theoretical_max_time / best_time
print(f"\n {'' * 90}")
print(f" Profile: {profile_name}")
print(f" {'' * 90}")
print(f" Total data size: {total_bytes / (1024*1024):.1f} MB")
print(f" Network rate: {params['rate']} ({format_throughput(rate_Bps)})" if rate_Bps else "")
print(f" Best client configuration: {best_name}")
print(f" Best client time: {best_time:.4f}s")
if theoretical_max_time is not None:
print(f" Theoretical max (uncompressed): {theoretical_max_time:.4f}s")
print(f" Speedup vs theoretical max: {speedup_vs_theoretical:.2f}x")
rsync_archive = rsync_times.get("rsync (archive)")
rsync_compress = rsync_times.get("rsync (archive + compress)")
if rsync_archive:
print(f" Speedup vs rsync (archive): {rsync_archive / best_time:.2f}x")
if rsync_compress:
print(f" Speedup vs rsync (compress): {rsync_compress / best_time:.2f}x")
failed = [r for r in all_results if r["status"] != "Success"] failed = [r for r in all_results if r["status"] != "Success"]
if failed: if failed:
print(f"\n {len(failed)} test(s) FAILED") print(f"\n {len(failed)} test(s) FAILED")
+3 -3
View File
@@ -8,7 +8,7 @@
static void test_config_lifecycle() { static void test_config_lifecycle() {
Config *cfg = config_create(str_dup("1.0"), str_dup("/src"), str_dup("/dst"), Config *cfg = config_create(str_dup("1.0"), str_dup("/src"), str_dup("/dst"),
true, true, false, false, false, 1, 4); true, true, false, false, false, 1, 4, false);
EXPECT_NOT_NULL(cfg); EXPECT_NOT_NULL(cfg);
EXPECT_EQ_STR(cfg->version, "1.0"); EXPECT_EQ_STR(cfg->version, "1.0");
EXPECT_EQ_STR(cfg->send_directory, "/src"); EXPECT_EQ_STR(cfg->send_directory, "/src");
@@ -23,7 +23,7 @@ static void test_config_lifecycle() {
static void test_pipeline_sender_lifecycle() { static void test_pipeline_sender_lifecycle() {
Config *cfg = config_create(str_dup("2.0"), str_dup("/src2"), Config *cfg = config_create(str_dup("2.0"), str_dup("/src2"),
str_dup("/dst2"), false, false, true, true, false, 1, 8); str_dup("/dst2"), false, false, true, true, false, 1, 8, false);
Queue *q1 = queue_create(5, NULL); Queue *q1 = queue_create(5, NULL);
Queue *q2 = queue_create(15, NULL); Queue *q2 = queue_create(15, NULL);
@@ -40,7 +40,7 @@ static void test_pipeline_sender_lifecycle() {
static void test_pipeline_receiver_lifecycle() { static void test_pipeline_receiver_lifecycle() {
Config *cfg = config_create(str_dup("3.0"), str_dup("/src3"), Config *cfg = config_create(str_dup("3.0"), str_dup("/src3"),
str_dup("/dst3"), true, true, true, true, false, 1, 2); str_dup("/dst3"), true, true, true, true, false, 1, 2, false);
Queue *q = queue_create(20, NULL); Queue *q = queue_create(20, NULL);
PipelineContextReceiver *pcr = pipeline_context_receiver_create(cfg, q, 42); PipelineContextReceiver *pcr = pipeline_context_receiver_create(cfg, q, 42);