Update README and add feature tests for rsync-compatible flags #9

Merged
TapTap merged 1 commits from update-readme-and-tests into rsync-flags 2026-07-16 13:38:14 +02:00
5 changed files with 249 additions and 66 deletions
+90 -54
View File
@@ -1,36 +1,40 @@
# FastSync
A high-performance file synchronization system with a custom TCP-based protocol, optional metadata preservation, compression, multithreading, and zero-copy `sendfile()` support.
A high-performance file synchronization system with SSH and TCP transport, streaming zstd compression, multithreaded transfer, metadata preservation, and rsync-compatible CLI flags.
## Technical Overview
1. Custom TCP-based client-server protocol with status codes
2. Chunked file transfer (files grouped into ~10 MB chunks)
3. Optional zstd compression (levels 122)
4. Multithreading for parallel file processing (producer-consumer with thread-safe queues)
5. Optional file metadata preservation (`mode`, `uid`, `gid`, `mtime`) — restored on disk
6. In-memory and disk-based storage options
7. `sendfile()` zero-copy path (~2× faster on localhost)
1. **Dual transport**: custom TCP client-server or SSH subprocess (rsync-style `user@host:/path`)
2. **Chunked file transfer**: files grouped into configurable-size chunks (default ~10 MB)
3. **Streaming zstd compression** (levels 122) using `ZSTD_compressStream2`
4. **Multithreading**: producer-consumer pipeline with thread-safe queues (scanner → loader → sender)
5. **Metadata preservation**: `mode`, `uid`, `gid`, `mtime` restored on disk when enabled
6. **`sendfile()` zero-copy** on TCP (~2× faster on loopback)
7. **SSH ControlMaster** for connection reuse across repeated invocations
8. **`--delete`**: receiver removes files not present in sender manifest
9. **`--exclude`**: glob-pattern filename filtering (`*`, `?`, no `/` crossing)
## System Architecture
### Client
- Recursively scans source directories (BFS)
- Groups files into chunks (default ~10 MB total)
- Optionally compresses with zstd
- Optionally serializes chunks into a compact binary format
- Optionally attaches per-file metadata (mode, ownership, timestamps)
- Sends via custom protocol or `sendfile()` zero-copy path
- Recursively scans source directories (BFS), supports exclude patterns
- Groups files into chunks (configurable size)
- Streaming zstd compression with configurable level
- Chunk serialization (compact binary format) or per-file transfer
- Manifests all sent paths when `--delete` is active
- Sends via TCP `sendfile()` or SSH pipe
- Optional progress display with throughput
### Server
- Listens on port 8080
- TCP mode: listens on port 8080; SSH mode: runs via `--stdio`
- Receives and reassembles files
- Decompresses, deserializes, restores metadata on disk
- Decompresses (streaming zstd), deserializes, restores metadata
- Processes `STATUS_MANIFEST` for `--delete`: walks destination tree, removes extras
- Thread pool for parallel processing
## Protocol Details
Status codes:
### Status Codes
| Code | Meaning |
|------|---------|
| `STATUS_OK` | Operation successful |
@@ -38,56 +42,76 @@ Status codes:
| `STATUS_FINISHED` | Transfer complete |
| `STATUS_NEXT` | Ready for next file (per-file mode) |
| `STATUS_CHUNK` | Following data is a serialized chunk |
| `STATUS_MANIFEST` | Following data is a file manifest (for `--delete`) |
### Wire Format — Metadata
When `use_metadata` is enabled (`-M`), each file entry carries a 4-byte `present` flag followed by five fields (`mode`, `uid`, `gid`, `mtime_sec`, `mtime_nsec`). When disabled globally, no metadata bytes are sent — zero wire overhead.
## Configuration
### Transfer Flow
```
Config → (STATUS_NEXT | STATUS_CHUNK)* → [STATUS_MANIFEST] → STATUS_FINISHED → STATUS_OK
```
## Command-Line Arguments
### Command-Line Arguments
| Argument | Description |
|----------|-------------|
| `-m` | Multithreading mode |
| Positional | `<source> <dest>` — automatic SSH detection if dest contains `:` |
| `-c [level]` | Compression with optional level (122, default 5) |
| `-z [level]` | Alias for `-c` |
| `-a, --archive` | Archive mode: enables `-c -m -M` (no `-s`) |
| `-m` | Multithreading mode |
| `-s` | Chunk serialization (batch all files per chunk) |
| `-f` | Sendfile zero-copy. Incompatible with `-c` / `-s`. |
| `-f` | Sendfile zero-copy. Incompatible with `-c` / `-s`. TCP only. |
| `-M, --preserve` | Preserve file metadata (mode, uid, gid, mtime) |
| `-n, --dry-run` | Scan and print what would be transferred |
| `-p <port>` | SSH port (default: 22) |
| `--progress` | Show real-time transfer speed |
| `--delete` | Delete files on receiver not present in source |
| `--exclude <pattern>` | Exclude files matching glob pattern (repeatable) |
| `--chunk-size <n>` | Chunk size in bytes (default: 10485760) |
| `--source-dir <path>` | Source directory (overrides `FASTSYNC_SOURCE_DIR`) |
| `--dest-dir <path>` | Server destination directory (overrides `FASTSYNC_DEST_DIR`) |
| `--save-to-disk` | Write received files to disk |
| `--server-host <ip>` | Server IP address (default: `127.0.0.1`) |
| `--server-port <n>` | Server port (default: `8080`) |
| `-v, --verbose` | Enable debug logging |
## Environment Variables
### Environment Variables
| Variable | Default | Description |
|----------|---------|-------------|
| `FASTSYNC_SOURCE_DIR` | User documents | Source directory fallback |
| `FASTSYNC_DEST_DIR` | `./data_copied` | Destination directory fallback |
| `FASTSYNC_SERVER_IP` | `127.0.0.1` | Server address |
| `FASTSYNC_SERVER_PORT` | `8080` | Server port |
| `FASTSYNC_SOURCE_DIR` | | Source directory fallback |
| `FASTSYNC_DEST_DIR` | | Destination directory fallback |
| `FASTSYNC_SAVE_TO_DISK` | `false` | Disk persistence fallback |
## Implementation Details
### Data Structures
1. **Chunk** — collection of files (~10 MB total)
1. **Chunk** — collection of files (~10 MB total by default)
2. **File** — path, content (`Data`), optional `FileMetadata` pointer
3. **FileMetadata**`mode`, `uid`, `gid`, `mtime_sec`, `mtime_nsec`
4. **Config** — runtime parameters
5. **Queue** — thread-safe queue with condition variables
4. **Config** — runtime parameters (transported over wire)
5. **Queue** — thread-safe bounded queue with condition variables
6. **DirectoryScanner** — recursive BFS traversal with exclude pattern support
### Key Algorithms
1. **File scanning** recursive BFS directory traversal
2. **Chunking** — files grouped by size limit
3. **Compression**zstd with configurable level
4. **Network protocol** custom TCP with status codes and optional metadata packing
1. **File scanning** — BFS directory traversal; each entry matched against exclude patterns
2. **Chunking** — files accumulated until `chunk_size` threshold, then flushed
3. **Compression**streaming zstd via `ZSTD_compressStream2` / `ZSTD_decompressStream`
4. **Network protocol** — status-code-driven exchange with metadata packing
5. **Metadata restoration**`chmod()`, `chown()`, `utimensat()` on the receiving side
6. **`--delete`** — sender tracks all sent paths; receiver walks destination tree and removes unlisted files/directories
7. **SSH transport**`socketpair()` + `fork()` + `execvp("ssh", ...)` with `ControlMaster` and port support
## Build Requirements
- C11 compiler
- CMake 4.1+
- zstd library
- zstd library (≥ 1.4.0 for streaming API)
- pthreads
- SSH client (for SSH transport)
## Building
@@ -97,55 +121,67 @@ cmake -B build -S . && cmake --build build -j$(nproc)
## Running
### Server
### Server (TCP mode)
```bash
./build/server
```
### Client
### Client — SSH (rsync-style)
```bash
./build/client /path/to/send user@host:/path/to/receive
```
### Client — TCP
```bash
# Basic
./build/client --source-dir /path/to/send --dest-dir /path/to/receive --save-to-disk
```
# With metadata preservation
./build/client -M --source-dir ... --dest-dir ...
### Common Options
```bash
# Archive mode (compression + multithreading + metadata)
./build/client -a /path/to/send user@host:/path
# Multithreaded + compression
./build/client -m -c 10
# Dry run
./build/client -n /path/to/send /path/to/receive
# Sendfile (zero-copy)
./build/client -f
# With progress and custom chunk size
./build/client --progress --chunk-size 2097152 /src user@host:/dst
# Exclude temporary files + delete extras on receiver
./build/client --exclude "*.tmp" --exclude "*.o" --delete /src user@host:/dst
# All features
./build/client -m -c -s -M
./build/client -a --progress --chunk-size 5242880 --exclude "*.log" --delete /src /dst
```
### Server via SSH
Place the `fastsync-server` binary in the remote `$PATH`. The client runs `ssh user@host fastsync-server --stdio` automatically when an SSH-style destination is given.
## Testing
```bash
# Unit tests
# Unit tests (7 suites)
./build/tests
# Integration benchmark (~50 MB data, 13 configurations + rsync comparison)
# Integration + benchmark suite
python3 test.py
# Profiles: --wan (100 Mbit, 50 ms, 1% loss), --unlimited (no throttling)
python3 test.py --wan
```
The benchmark prints throughput metrics for the best configuration and speedup vs rsync.
The benchmark prints throughput metrics, best configuration, and speedup vs rsync.
## Performance Considerations
1. Chunk size (~10 MB) balances memory and transfer efficiency
1. Chunk size (~10 MB default) balances memory and transfer efficiency
2. Compression level trades CPU for bandwidth
3. `sendfile()` bypasses userspace — ~2× faster on localhost for large files
4. Multithreading scales with core count
5. Metadata transfer adds negligible overhead when disabled, ~24 bytes per file when enabled
5. Metadata transfer adds negligible overhead (~24 bytes per file when enabled)
6. SSH socketpair buffer set to 1 MB for improved pipe throughput
7. SSH ControlMaster reuses connections across repeated invocations
## Benchmark Results
50 MB of mixed file sizes over `localhost` with disk I/O throttled (reads ≤ 15 MB/s, writes ≤ 10 MB/s) and network emulation via `tc netem`. Each test was run 3×; the median is reported below.
25 MB of mixed file sizes over `localhost` with disk I/O throttled (reads ≤ 15 MB/s, writes ≤ 10 MB/s) and network emulation via `tc netem`. Each test was run 3×; the median is reported below.
### LAN (1000 Mbit, 20 ms ±1 ms, 0.1% loss)
@@ -167,4 +203,4 @@ The benchmark prints throughput metrics for the best configuration and speedup v
| rsync (archive) | 17.44 s | — | — |
| rsync (archive + compress) | 1.47 s | — | — |
Compression reduces the data on the wire enough that the transfer becomes latency-bound rather than bandwidth-bound. On WAN, the best configuration runs 10.8× faster than the theoretical limit for uncompressed data, since zstd shrinks the 50 MB payload to a fraction of its original size over the wire.
Compression reduces the data on the wire enough that the transfer becomes latency-bound rather than bandwidth-bound. On WAN, the best configuration runs 10.8× faster than the theoretical limit for uncompressed data, since zstd shrinks the 25 MB payload to a fraction of its original size over the wire.
+10 -7
View File
@@ -100,9 +100,11 @@ static int scan_directory_multithreaded(void *pipeline_context) {
while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
if (context->config->use_delete) {
mtx_lock(&context->mutex_scanner);
for (int i = 0; i < current_chunk->element_count; i++)
array_list_add(context->manifest,
str_dup(current_chunk->items[i]->path));
for (int i = 0; i < current_chunk->element_count; i++) {
const char *p = current_chunk->items[i]->path;
if (*p == '/') p++;
array_list_add(context->manifest, str_dup(p));
}
mtx_unlock(&context->mutex_scanner);
}
queue_enqueue_multithreaded(context->queue_scanner, current_chunk,
@@ -192,8 +194,11 @@ int send_files(Config *config) {
unsigned long long chunk_bytes = 0;
for (int i = 0; i < current_chunk->element_count; i++) {
chunk_bytes += current_chunk->items[i]->data->size;
if (manifest)
array_list_add(manifest, str_dup(current_chunk->items[i]->path));
if (manifest) {
const char *p = current_chunk->items[i]->path;
if (*p == '/') p++;
array_list_add(manifest, str_dup(p));
}
}
if (!config->use_sendfile) {
for (int i = 0; i < current_chunk->element_count; i++)
@@ -218,8 +223,6 @@ int send_files(Config *config) {
send_int(client->file_descriptor, manifest->size);
for (int i = 0; i < manifest->size; i++)
send_str(client->file_descriptor, (char *)manifest->items[i]);
for (int i = 0; i < manifest->size; i++)
free(manifest->items[i]);
array_list_delete(manifest);
}
send_status(client->file_descriptor, STATUS_FINISHED);
-2
View File
@@ -63,8 +63,6 @@ int receive_files(Config *config, int file_descriptor) {
array_list_add(manifest, receive_str(file_descriptor));
fprintf(stderr, "Deleting files not in manifest...\n");
delete_extras(config->receive_root_directory, manifest);
for (int i = 0; i < manifest->size; i++)
free(manifest->items[i]);
array_list_delete(manifest);
status = receive_status(file_descriptor);
}
-2
View File
@@ -38,8 +38,6 @@ PipelineContextSender *pipeline_context_sender_create(Config *config,
void pipeline_context_sender_destroy(PipelineContextSender *context) {
if (context->manifest) {
for (int i = 0; i < context->manifest->size; i++)
free(context->manifest->items[i]);
array_list_delete(context->manifest);
}
config_delete(context->config);
+149 -1
View File
@@ -195,7 +195,7 @@ def start_rsync_daemon(source_dir):
return port, conf, daemon
def run_single_test(cmd, name, source_dir, dest_dir, *, source_prefix=None, no_server=False):
def run_single_test(cmd, name, source_dir, dest_dir, *, source_prefix=None, no_server=False, expected_missing=None):
if os.path.exists(dest_dir):
shutil.rmtree(dest_dir)
if no_server:
@@ -216,6 +216,8 @@ def run_single_test(cmd, name, source_dir, dest_dir, *, source_prefix=None, no_s
received = os.path.join(dest_dir, source_prefix if source_prefix is not None
else os.path.abspath(source_dir).lstrip(os.sep))
mismatches, missing = verify_transfer(source_dir, received)
if expected_missing:
missing = [m for m in missing if m not in expected_missing]
first_line = lambda s: (s or "").strip().split("\n")[0]
entry = {
@@ -318,6 +320,152 @@ def run_profile(profile_name, source_dir, dest_dir):
except Exception:
pass
# Feature-specific tests for rsync-compatible flags
print("\n " + "" * 56 + "\n Feature Tests\n " + "" * 56)
# Dry run (-n) — no server needed
print("\n --- Dry run (-n) ---")
flags = BASE_CLIENT_FLAGS + ["-n"]
cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + flags
print(f" Running: {' '.join(cmd)}")
try:
start = time.monotonic()
result = subprocess.run(cmd, text=True, capture_output=True)
duration = time.monotonic() - start
r = {"name": "Dry run (-n)", "suite": profile_name}
if result.returncode == 0 and "Dry run:" in result.stdout:
r["status"] = "Success"
r["time"] = f"{duration:.4f}s"
r["error"] = ""
else:
r["status"] = "Failed"
r["time"] = "N/A"
r["error"] = f"Exit {result.returncode}: {(result.stderr or result.stdout)[:100]}"
results.append(r)
except Exception as e:
results.append({"name": "Dry run (-n)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)})
# Archive mode (-a)
feature_flags = BASE_CLIENT_FLAGS + ["-a"]
cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + feature_flags
print(f"\n --- Archive mode (-a) ---\n Running: {' '.join(cmd)}")
try:
r = run_single_test(cmd, "Archive mode (-a)", source_dir, dest_dir)
r["suite"] = profile_name
results.append(r)
except Exception as e:
results.append({"name": "Archive mode (-a)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)})
# Exclude (--exclude small.txt)
feature_flags = BASE_CLIENT_FLAGS + ["--exclude", "small.txt"]
cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + feature_flags
print(f"\n --- Exclude (--exclude small.txt) ---\n Running: {' '.join(cmd)}")
try:
r = run_single_test(cmd, "Exclude (--exclude small.txt)", source_dir, dest_dir,
expected_missing=["small.txt"])
r["suite"] = profile_name
results.append(r)
except Exception as e:
results.append({"name": "Exclude (--exclude small.txt)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)})
# Progress (--progress)
feature_flags = BASE_CLIENT_FLAGS + ["--progress"]
cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + feature_flags
print(f"\n --- Progress (--progress) ---\n Running: {' '.join(cmd)}")
try:
r = run_single_test(cmd, "Progress (--progress)", source_dir, dest_dir)
r["suite"] = profile_name
results.append(r)
except Exception as e:
results.append({"name": "Progress (--progress)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)})
# Chunk size (--chunk-size 5242880)
feature_flags = BASE_CLIENT_FLAGS + ["--chunk-size", "5242880"]
cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + feature_flags
print(f"\n --- Chunk size (--chunk-size 5242880) ---\n Running: {' '.join(cmd)}")
try:
r = run_single_test(cmd, "Chunk size (--chunk-size 5242880)", source_dir, dest_dir)
r["suite"] = profile_name
results.append(r)
except Exception as e:
results.append({"name": "Chunk size (--chunk-size 5242880)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)})
# Delete (--delete) — pre-populate dest, add extra files, then sync with --delete
# Note: server handles one client per launch, so we restart between syncs
print(f"\n --- Delete (--delete) ---")
try:
flags = BASE_CLIENT_FLAGS + ["-M"]
# First sync (no delete) to populate dest
s1 = subprocess.Popen(SERVER_CMD, stdout=subprocess.DEVNULL, stderr=None)
time.sleep(0.5)
first_cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + flags
r1 = subprocess.run(first_cmd, text=True, capture_output=True)
wait_proc(s1)
if r1.returncode != 0:
raise RuntimeError(f"First sync failed: {r1.stderr[:100]}")
# Add extra files to received dir
received = os.path.join(dest_dir, os.path.abspath(source_dir).lstrip(os.sep))
extra_path = os.path.join(received, "extra_file.txt")
with open(extra_path, "w") as f:
f.write("should be deleted")
extra_dir = os.path.join(received, "extra_dir")
os.makedirs(extra_dir, exist_ok=True)
with open(os.path.join(extra_dir, "nested.txt"), "w") as f:
f.write("nested extra")
# Second sync with --delete (fresh server)
s2 = subprocess.Popen(SERVER_CMD, stdout=subprocess.DEVNULL, stderr=None)
time.sleep(0.5)
second_cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + flags + ["--delete"]
start = time.monotonic()
r2 = subprocess.run(second_cmd, text=True, capture_output=True)
duration = time.monotonic() - start
wait_proc(s2)
r = {"name": "Delete (--delete)", "suite": profile_name}
if r2.returncode == 0 and not os.path.exists(extra_path) and not os.path.exists(extra_dir):
mismatches, missing = verify_transfer(source_dir, received)
if not mismatches and not missing:
r["status"] = "Success"
r["time"] = f"{duration:.4f}s"
r["error"] = ""
else:
r["status"] = "Failed"
r["time"] = "N/A"
r["error"] = f"post-delete verify: mismatches={len(mismatches)}, missing={len(missing)}"
else:
r["status"] = "Failed"
r["time"] = "N/A"
errs = []
if r2.returncode != 0:
errs.append(f"Exit {r2.returncode}: {(r2.stderr or r2.stdout)[:60]}")
if os.path.exists(extra_path):
errs.append("extra_file.txt remains")
if os.path.exists(extra_dir):
errs.append("extra_dir remains")
r["error"] = " | ".join(errs)
results.append(r)
except Exception as e:
results.append({"name": "Delete (--delete)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)})
# SSH feature tests
if SSH_AVAILABLE:
ssh_dest = f"localhost:{dest_dir}_ssh"
ssh_feature_cases = [
{"name": "SSH Archive (-a)", "flags": ["-a"]},
{"name": "SSH Exclude (--exclude small.txt)", "flags": ["--exclude", "small.txt"], "expected_missing": ["small.txt"]},
]
for case in ssh_feature_cases:
flags = BASE_CLIENT_FLAGS + case["flags"]
cmd = BASE_CLIENT_CMD + [source_dir, ssh_dest] + flags
print(f"\n --- {case['name']} ---\n Running: {' '.join(cmd)}")
try:
r = run_single_test(cmd, case["name"], source_dir, f"{dest_dir}_ssh",
no_server=True,
expected_missing=case.get("expected_missing"))
r["suite"] = profile_name
results.append(r)
except Exception as e:
results.append({"name": case["name"], "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)})
except (subprocess.CalledProcessError, RuntimeError) as e:
print(f" Error: {e}")
results = []