From 2ff51246c354b12f8de9becb78012b7cfd314957 Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 16 Jul 2026 13:31:43 +0200 Subject: [PATCH] Update README and add feature tests for rsync-compatible flags --- README.md | 144 ++++++++++++++++++++------------- src/client/client_send.c | 17 ++-- src/server/server.c | 2 - src/shared/multiprocessing.c | 2 - test.py | 150 ++++++++++++++++++++++++++++++++++- 5 files changed, 249 insertions(+), 66 deletions(-) diff --git a/README.md b/README.md index bffd576..9c78a91 100644 --- a/README.md +++ b/README.md @@ -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 1–22) -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 1–22) 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 | ` ` — automatic SSH detection if dest contains `:` | | `-c [level]` | Compression with optional level (1–22, 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 ` | SSH port (default: 22) | +| `--progress` | Show real-time transfer speed | +| `--delete` | Delete files on receiver not present in source | +| `--exclude ` | Exclude files matching glob pattern (repeatable) | +| `--chunk-size ` | Chunk size in bytes (default: 10485760) | | `--source-dir ` | Source directory (overrides `FASTSYNC_SOURCE_DIR`) | | `--dest-dir ` | Server destination directory (overrides `FASTSYNC_DEST_DIR`) | | `--save-to-disk` | Write received files to disk | +| `--server-host ` | Server IP address (default: `127.0.0.1`) | +| `--server-port ` | 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. diff --git a/src/client/client_send.c b/src/client/client_send.c index 2187b33..4797d31 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -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); diff --git a/src/server/server.c b/src/server/server.c index 1be205c..577b430 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -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); } diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 5f26d56..bcf0812 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -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); diff --git a/test.py b/test.py index 02d9a72..131bddc 100755 --- a/test.py +++ b/test.py @@ -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 = [] -- 2.52.0