Author SHA1 Message Date
TapTap 1aebcfd55d refactor: consolidate transfer infrastructure
CI / lint (pull_request) Failing after 3s
CI / build-and-test (pull_request) Skipped
CI / sanitizers (address) (pull_request) Skipped
CI / sanitizers (undefined) (pull_request) Skipped
CI / fuzz-build (pull_request) Skipped
CI / coverage (pull_request) Skipped
CI / valgrind (pull_request) Skipped
2026-09-01 21:21:24 +02:00
TapTap 6b7213db5c docs: align compatibility guidance
CI / lint (pull_request) Failing after 2s
CI / build-and-test (pull_request) Has been skipped
CI / sanitizers (address) (pull_request) Has been skipped
CI / sanitizers (undefined) (pull_request) Has been skipped
CI / fuzz-build (pull_request) Has been skipped
CI / coverage (pull_request) Has been skipped
CI / valgrind (pull_request) Has been skipped
2026-08-15 13:41:23 +02:00
TapTap de537810e5 docs: fix server destination root examples
CI / lint (pull_request) Failing after 2s
CI / build-and-test (pull_request) Has been skipped
CI / sanitizers (address) (pull_request) Has been skipped
CI / sanitizers (undefined) (pull_request) Has been skipped
CI / fuzz-build (pull_request) Has been skipped
CI / coverage (pull_request) Has been skipped
CI / valgrind (pull_request) Has been skipped
2026-08-15 13:14:28 +02:00
TapTap def58554d9 docs: clarify rsync compatibility roadmap
CI / lint (pull_request) Failing after 2s
CI / build-and-test (pull_request) Has been skipped
CI / sanitizers (address) (pull_request) Has been skipped
CI / sanitizers (undefined) (pull_request) Has been skipped
CI / fuzz-build (pull_request) Has been skipped
CI / coverage (pull_request) Has been skipped
CI / valgrind (pull_request) Has been skipped
2026-08-15 12:54:38 +02:00
11 changed files with 666 additions and 886 deletions
+72 -29
View File
@@ -8,10 +8,8 @@ set(CMAKE_C_STANDARD_REQUIRED ON)
add_compile_options(-Wall -g -O3) add_compile_options(-Wall -g -O3)
# --- Sanitizer option ---
set(SANITIZER "none" CACHE STRING "Sanitizer to enable (address, thread, undefined, none)") set(SANITIZER "none" CACHE STRING "Sanitizer to enable (address, thread, undefined, none)")
set_property(CACHE SANITIZER PROPERTY STRINGS address thread undefined none) set_property(CACHE SANITIZER PROPERTY STRINGS address thread undefined none)
if(SANITIZER STREQUAL "address") if(SANITIZER STREQUAL "address")
add_compile_options(-fsanitize=address -fno-omit-frame-pointer -g) add_compile_options(-fsanitize=address -fno-omit-frame-pointer -g)
add_link_options(-fsanitize=address) add_link_options(-fsanitize=address)
@@ -25,13 +23,11 @@ elseif(NOT SANITIZER STREQUAL "none")
message(FATAL_ERROR "Unknown sanitizer: ${SANITIZER}. Supported values: address, thread, undefined, none") message(FATAL_ERROR "Unknown sanitizer: ${SANITIZER}. Supported values: address, thread, undefined, none")
endif() endif()
# --- Strict warnings option ---
option(STRICT_WARNINGS "Enable strict warnings (Wextra, Wpedantic, Werror)" OFF) option(STRICT_WARNINGS "Enable strict warnings (Wextra, Wpedantic, Werror)" OFF)
if(STRICT_WARNINGS) if(STRICT_WARNINGS)
add_compile_options(-Wextra -Wpedantic -Werror) add_compile_options(-Wextra -Wpedantic -Werror)
endif() endif()
# --- Coverage option ---
option(ENABLE_COVERAGE "Enable gcov coverage" OFF) option(ENABLE_COVERAGE "Enable gcov coverage" OFF)
if(ENABLE_COVERAGE) if(ENABLE_COVERAGE)
add_compile_options(--coverage -fprofile-arcs -ftest-coverage -O0 -g) add_compile_options(--coverage -fprofile-arcs -ftest-coverage -O0 -g)
@@ -49,55 +45,102 @@ FetchContent_MakeAvailable(xxhash)
set(THREADS_PREFER_PTHREAD_FLAG ON) set(THREADS_PREFER_PTHREAD_FLAG ON)
find_package(Threads REQUIRED) find_package(Threads REQUIRED)
find_library(ZSTD_LIBRARY zstd) find_library(ZSTD_LIBRARY zstd)
if(NOT ZSTD_LIBRARY) if(NOT ZSTD_LIBRARY)
message(FATAL_ERROR "zstd library not found. Ensure it is in your nix-shell!") message(FATAL_ERROR "zstd library not found. Ensure it is in your nix-shell!")
endif() endif()
find_package(OpenSSL REQUIRED) find_package(OpenSSL REQUIRED)
file(GLOB SHARED_SRCS "src/shared/*.c") set(SHARED_SRCS
file(GLOB SERVER_SRCS "src/server/*.c") src/shared/array_list.c
file(GLOB CLIENT_SRCS "src/client/*.c") src/shared/chunk.c
src/shared/compression.c
src/shared/config.c
src/shared/data.c
src/shared/delta.c
src/shared/file.c
src/shared/log.c
src/shared/metadata.c
src/shared/multiprocessing.c
src/shared/protocol.c
src/shared/queue.c
src/shared/transport_ssh.c
src/shared/transport_tcp.c
src/shared/transport_tls.c
src/shared/utils.c
)
set(SERVER_SRCS src/server/server.c)
set(CLIENT_SRCS
src/client/client_cli.c
src/client/client_send.c
src/client/scanner.c
src/client/usage.c
)
function(configure_fastsync_target target)
target_include_directories(${target} PRIVATE src/shared src/server src/client)
target_link_libraries(${target} PRIVATE Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash)
endfunction()
# --- Main executables ---
add_executable(server ${SERVER_SRCS} ${SHARED_SRCS}) add_executable(server ${SERVER_SRCS} ${SHARED_SRCS})
target_include_directories(server PRIVATE src/shared src/server src/client) configure_fastsync_target(server)
target_link_libraries(server PRIVATE Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash)
add_executable(client ${CLIENT_SRCS} ${SHARED_SRCS}) add_executable(client ${CLIENT_SRCS} ${SHARED_SRCS})
target_include_directories(client PRIVATE src/shared src/server src/client) configure_fastsync_target(client)
target_link_libraries(client PRIVATE Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash)
# --- Testing ---
enable_testing() enable_testing()
set(TEST_SRCS
# Common test libraries tests/runner.c
set(TEST_LIBS Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash) tests/test_array_list.c
set(TEST_INCLUDES tests src/shared src/server src/client) tests/test_chunk.c
tests/test_client_cli.c
# Monolithic test binary (backward compatible) tests/test_compression.c
file(GLOB TEST_SRCS "tests/test_*.c" "tests/runner.c") tests/test_config.c
tests/test_data.c
tests/test_delta.c
tests/test_file.c
tests/test_file_sendfile.c
tests/test_fuzz_smoke.c
tests/test_glob.c
tests/test_log.c
tests/test_metadata.c
tests/test_multiprocessing.c
tests/test_property.c
tests/test_protocol.c
tests/test_queue.c
tests/test_robustness.c
tests/test_scanner.c
tests/test_server.c
tests/test_shared_utils.c
tests/test_stress.c
tests/test_transport_ssh.c
tests/test_transport_tcp.c
tests/test_transport_tls.c
)
add_executable(tests ${TEST_SRCS} ${SHARED_SRCS} src/client/scanner.c src/client/client_cli.c) add_executable(tests ${TEST_SRCS} ${SHARED_SRCS} src/client/scanner.c src/client/client_cli.c)
target_include_directories(tests PRIVATE ${TEST_INCLUDES}) configure_fastsync_target(tests)
target_include_directories(tests PRIVATE tests)
target_compile_definitions(tests PRIVATE FASTSYNC_TEST_BUILD) target_compile_definitions(tests PRIVATE FASTSYNC_TEST_BUILD)
target_link_libraries(tests PRIVATE ${TEST_LIBS})
add_test(NAME unit_all COMMAND tests) add_test(NAME unit_all COMMAND tests)
# --- Fuzz targets (requires clang) ---
option(ENABLE_FUZZ "Build fuzz targets (requires clang)" OFF) option(ENABLE_FUZZ "Build fuzz targets (requires clang)" OFF)
if(ENABLE_FUZZ) if(ENABLE_FUZZ)
if(NOT CMAKE_C_COMPILER_ID MATCHES "Clang") if(NOT CMAKE_C_COMPILER_ID MATCHES "Clang")
message(FATAL_ERROR "ENABLE_FUZZ requires Clang (compiler is ${CMAKE_C_COMPILER_ID})") message(FATAL_ERROR "ENABLE_FUZZ requires Clang (compiler is ${CMAKE_C_COMPILER_ID})")
endif() endif()
file(GLOB FUZZ_SRCS "tests/fuzz/*.c") set(FUZZ_SRCS
tests/fuzz/fuzz_chunk_deserialize.c
tests/fuzz/fuzz_compress_decompress.c
tests/fuzz/fuzz_delta_deserialize.c
tests/fuzz/fuzz_delta_signature_deserialize.c
tests/fuzz/fuzz_glob_match.c
tests/fuzz/fuzz_metadata_from_buf.c
)
foreach(FUZZ_SRC ${FUZZ_SRCS}) foreach(FUZZ_SRC ${FUZZ_SRCS})
get_filename_component(FUZZ_NAME ${FUZZ_SRC} NAME_WE) get_filename_component(FUZZ_NAME ${FUZZ_SRC} NAME_WE)
add_executable(${FUZZ_NAME} ${FUZZ_SRC} ${SHARED_SRCS}) add_executable(${FUZZ_NAME} ${FUZZ_SRC} ${SHARED_SRCS})
target_include_directories(${FUZZ_NAME} PRIVATE ${TEST_INCLUDES}) configure_fastsync_target(${FUZZ_NAME})
target_include_directories(${FUZZ_NAME} PRIVATE tests)
target_compile_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined -fno-omit-frame-pointer) target_compile_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined -fno-omit-frame-pointer)
target_link_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined) target_link_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined)
target_link_libraries(${FUZZ_NAME} PRIVATE ${TEST_LIBS})
endforeach() endforeach()
endif() endif()
+316 -297
View File
@@ -1,363 +1,382 @@
# FastSync # FastSync
A high-performance file synchronization system with SSH and TCP transport, TLS encryption, streaming zstd compression, multithreaded transfer, incremental sync, metadata preservation, and rsync-compatible CLI flags. FastSync is a high-performance file synchronization tool designed to become a
drop-in replacement for common `rsync` workflows. It keeps the familiar
source/destination model and rsync-style options while adding optional
multithreading, streaming zstd compression, chunking, zero-copy TCP transfers,
and native TCP/TLS transports.
## Technical Overview The compatibility target is straightforward:
1. **Dual transport**: custom TCP client-server or SSH subprocess (rsync-style `user@host:/path`) - Existing rsync commands should keep the same meaning.
2. **TLS encryption**: OpenSSL-based TLS 1.2+ for encrypted TCP connections with optional CA verification - FastSync-only performance options should be additive and optional.
3. **Chunked file transfer**: files grouped into configurable-size chunks (default ~10 MB) - A normal compatibility-mode transfer should prioritize rsync filesystem
4. **Streaming zstd compression** (levels 1–22) using `ZSTD_compressStream2` semantics over maximum throughput.
5. **Multithreading**: producer-consumer pipeline with thread-safe queues (scanner → loader → sender)
6. **Incremental sync**: skip files unchanged since last transfer (compares size + mtime)
7. **Batch incremental**: send incremental checks in batched groups for reduced round-trips
8. **Metadata preservation**: file mode and mtime are restored when enabled; ownership and atime are intentionally not restored
9. **`sendfile()` zero-copy** on TCP (~2× faster on loopback)
10. **SSH ControlMaster** for connection reuse across repeated invocations
11. **Bandwidth limiting**: token-bucket throttling (`--bwlimit`)
12. **`--delete`**: receiver removes files not present in sender manifest
13. **`--exclude` / `--include`**: glob-pattern filename filtering
14. **Path traversal protection**: `..` sequences in file paths are rejected automatically
15. **Connection limits**: server enforces maximum concurrent connections (default 100)
16. **Keep-alive**: periodic `STATUS_KEEPALIVE` messages detect stalled connections
17. **Abort handling**: `SIGINT` sends `STATUS_ABORT` for clean server-side teardown
18. **Atomic writes**: received files are written to a temporary name then atomically renamed
19. **Backup mode**: `--backup` preserves overwritten files with optional `--backup-dir`
20. **Log file**: `--log-file` redirects log output to a file instead of stderr
21. **Transfer statistics**: `--stats` prints summary of transferred bytes, files, and timing
## System Architecture FastSync currently speaks its own protocol to `fastsync-server`. SSH mode
starts that server remotely; it does not yet interoperate with an unmodified
rsync client or rsync daemon. See [Compatibility Status](#compatibility-status)
for the current boundary.
### Client ## Why FastSync
- Recursively scans source directories (BFS), supports exclude and include patterns
- Groups files into chunks (configurable size)
- Streaming zstd compression with configurable level
- Chunk serialization (compact binary format) or per-file transfer
- Incremental transfer: sends file metadata to server, skips unchanged files
- Batch incremental: groups incremental checks to minimize round-trips
- Manifests all sent paths when `--delete` is active
- Sends via TCP `sendfile()` or SSH pipe
- Optional progress display with throughput
- Bandwidth limiting via token-bucket algorithm
- Configurable I/O and connection timeouts (`--timeout`, `--contimeout`)
- Backup overwritten files (`--backup`) with optional directory (`--backup-dir`)
- Transfer statistics summary (`--stats`)
- Maximum directory depth control (`--max-depth`)
- Log file output (`--log-file`)
- Exclude patterns from file (`--exclude-from`)
### Server FastSync uses a producer-consumer transfer pipeline and can combine several
- TCP mode: listens on configurable port (default 8080); SSH mode: runs via `--stdio` optimizations for large or high-latency transfers:
- TLS mode: wraps TCP connections with OpenSSL with optional CA verification
- Receives and reassembles files
- Decompresses (streaming zstd), deserializes, restores metadata
- Handles incremental checks: compares size + mtime against destination files
- Handles batch incremental checks for reduced round-trips
- Processes `STATUS_MANIFEST` for `--delete`: walks destination tree, removes extras
- Per-connection concurrency via `fork()` with configurable connection limit (default 100)
- Thread pool for parallel processing
- Atomic writes: files written to `.tmp` path then atomically renamed on success
- Abort handling: cleanly shuts down on `STATUS_ABORT` from client
- Path traversal protection: rejects file paths containing `..`
## Protocol Details - Multithreaded scanning, loading, and sending.
- Streaming zstd compression with levels 1 through 22.
- Configurable file chunking and compact chunk serialization.
- `sendfile()` zero-copy transfers over TCP.
- Batched incremental checks to reduce round trips.
- Optional block-level delta transfer for FastSync peers.
- Bandwidth limiting, progress reporting, statistics, and backups.
- TCP, SSH, and TLS transports.
- Atomic temporary-file writes by default.
### Status Codes These optimizations are disabled or selected independently. Users can start
| Code | Meaning | with rsync-style commands and add FastSync options when they are useful.
|------|---------|
| `STATUS_OK` | Operation successful |
| `STATUS_ERROR` | Error occurred |
| `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`) |
| `STATUS_CHECK` | Incremental check: client sends file path + size + mtime and, when negotiated, checksum; server responds with OK (skip) or NEXT (send) |
| `STATUS_CHECK_BATCH` | Batch incremental check: multiple file checks sent in one message |
| `STATUS_KEEPALIVE` | Keep-alive heartbeat to detect stalled connections |
| `STATUS_ABORT` | Abort signal: client interrupts, server cleans up and exits |
| `STATUS_DELTA_SIGNATURE` | Delta sync: following data is a file signature (rsync-style rolling hash) |
| `STATUS_DELTA_DATA` | Delta sync: following data is a delta patch for a file |
### Wire Format — Metadata ## Compatibility Status
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. FastSync is currently an rsync-compatible CLI in progress, not a complete
replacement for every rsync feature or protocol mode.
### Transfer Flow ### Working today
```
Config → (STATUS_NEXT | STATUS_CHUNK | STATUS_CHECK | STATUS_CHECK_BATCH)* → [STATUS_MANIFEST] → STATUS_FINISHED → STATUS_OK
```
Keep-alive (`STATUS_KEEPALIVE`) may be sent at any point during the transfer. The receiver resets its inactivity timer on receipt. If no data arrives within the receive timeout, the connection is aborted. - Recursive directory scanning.
- Rsync-style source and destination arguments.
- SSH transport using `user@host:destination` paths below the remote authorized root.
- TCP client/server transfers.
- Dry runs, excludes, includes, size filters, backups, statistics, and
bandwidth limiting.
- Incremental size/mtime checks and optional xxHash64 content checks.
- FastSync-native delta transfer for changed files.
- Optional mode and timestamp preservation.
- Delete manifests with server-side delete authorization.
- Temporary-file writes with atomic rename by default.
- Path traversal checks and destination-root confinement.
Abort (`STATUS_ABORT`) may be sent at any point. On receipt the server cleans up temporary files and exits the child process. ### Not yet equivalent to rsync
### Protocol Version - The FastSync wire protocol is not the rsync wire protocol.
- SSH mode requires `fastsync-server` on the remote host.
- Archive mode does not yet provide all of rsync's `-rlptgoD` behavior.
- Symlink transfer is incomplete; link targets are not yet recreated in all
modes.
- Owner/group, ACL, xattr, hard-link, device, and special-file handling is
incomplete or unavailable.
- Sparse-file handling does not yet preserve all holes correctly.
- `--partial`, `--partial-dir`, `--append`, and `--append-verify` are not yet
full rsync-style resumable transfers.
- Several rsync short options currently have FastSync-specific meanings. Do
not assume every short option is interchangeable yet.
`2.2.0` — server and client must match. This version adds a 64-bit XXH64 checksum to checksum-enabled `STATUS_CHECK` messages and validates the negotiated compression choice (`zstd` or `none`). Older clients and servers must not be mixed with this version; mismatch results in `STATUS_ERROR`. The detailed flag matrix is maintained in
[`RSYNC_COMPAT.md`](RSYNC_COMPAT.md). It distinguishes implemented,
partial, alternate, and planned behavior.
Config negotiation is sender-driven: the client serializes transfer options and the server applies them while receiving and writing files. `--checksum` compares size and content checksum instead of timestamps. `--compress-choice zstd` enables zstd; `none` disables it. Unsupported choices are rejected during config exchange. ## Quick Start
## Command-Line Arguments ### Build
### Client Requirements: C11 compiler, CMake 3.22 or newer, xxHash, zstd, OpenSSL,
pthreads, and an SSH client for SSH transport. The first CMake configure fetches
| Argument | Description | xxHash from GitHub, so network access is required unless the dependency is
|----------|-------------| already cached.
| Positional | `<source> <dest>` — 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` | Sendfile zero-copy. Incompatible with `-c` / `-s`. TCP only. |
| `-M, --preserve` | Preserve supported file metadata (mode and mtime; ownership and atime are unsupported) |
| `-n, --dry-run` | Scan and print what would be transferred |
| `-p <port>` | SSH port (default: 22) |
| `-v, --verbose` | Enable debug logging |
| `--progress` | Show real-time transfer speed |
| `--delete` | Delete files on receiver not present in source |
| `--exclude <pattern>` | Exclude files matching glob pattern (repeatable) |
| `--exclude-from <file>` | Read exclude patterns from a file (one per line) |
| `--include <pattern>` | Only transfer files matching glob pattern (repeatable, whitelist) |
| `--max-size <n>` | Skip files larger than n bytes |
| `--min-size <n>` | Skip files smaller than n bytes |
| `--incremental` | Skip files unchanged since last transfer (size + mtime). Auto-enables `--preserve`. Incompatible with `-s`. |
| `--bwlimit <KB/s>` | Bandwidth limit in kilobytes per second |
| `--chunk-size <n>` | Chunk size in bytes (default: 10485760) |
| `--timeout <sec>` | I/O timeout in seconds (default: 30) |
| `--contimeout <sec>` | Connection timeout in seconds (default: 10) |
| `--backup` | Backup existing destination files before overwriting |
| `--backup-dir <dir>` | Target directory for backups (requires `--backup`) |
| `--stats` | Print transfer statistics at end (bytes, files, timing) |
| `--max-depth <n>` | Maximum directory depth to recurse (0 = unlimited, default: 0) |
| `--log-file <path>` | Write log messages to file instead of stderr |
| `--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`) |
| `--tls` | Enable TLS encryption |
| `--cert <path>` | TLS certificate file (PEM) |
| `--key <path>` | TLS private key file (PEM) |
| `--ca <path>` | TLS CA certificate file for verification (PEM) |
### Server
| Argument | Description |
|----------|-------------|
| `--stdio` | Run in stdio mode (for SSH transport; single connection then exits) |
| `-p <port>` | TCP listen port (default: 8080, range: 1–65535) |
| `--tls` | Enable TLS encryption |
| `--cert <path>` | TLS certificate file (PEM) |
| `--key <path>` | TLS private key file (PEM) |
| `--ca <path>` | TLS CA certificate file for verification (PEM) |
| `-v, --verbose` | Enable debug logging |
| `--help` | Show help |
## Environment Variables
| Variable | Default | Description |
|----------|---------|-------------|
| `FASTSYNC_SOURCE_DIR` | — | Source directory fallback |
| `FASTSYNC_DEST_DIR` | — | Destination directory fallback |
| `FASTSYNC_SAVE_TO_DISK` | `false` | Disk persistence fallback |
| `FASTSYNC_SSH_PORT` | `22` | Default SSH port |
| `FASTSYNC_SERVER_HOST` | `127.0.0.1` | Default server host |
| `FASTSYNC_SERVER_PORT` | `8080` | Default server port |
| `FASTSYNC_TLS_CERT` | — | Default TLS certificate path |
| `FASTSYNC_TLS_KEY` | — | Default TLS private key path |
| `FASTSYNC_TLS_CA` | — | Default TLS CA certificate path |
## Implementation Details
### Data Structures
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`; uid/gid are advisory wire fields and are never applied by the receiver; atime is unsupported
4. **Config** — runtime parameters (transported over wire, TLS settings excluded). Includes `timeout`, `contimeout`, `quiet`, `backup`, `backup_dir`, `stats`, `max_depth`, `log_file`, `queue_size`.
5. **Queue** — thread-safe bounded queue with condition variables
6. **DirectoryScanner** — recursive BFS traversal with exclude and include pattern support, max-depth enforcement
### Key Algorithms
1. **File scanning** — BFS directory traversal; entries matched against exclude and include patterns, max-depth enforced
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, keep-alive, and abort support
5. **Incremental check** — client sends `STATUS_CHECK` + path + size + mtime and, with `--checksum`, XXH64 content checksum; server compares against destination. Can be batched via `STATUS_CHECK_BATCH` for reduced round-trips.
6. **Bandwidth limiting** — token-bucket algorithm with `nanosleep` throttling on 64 KB write chunks
7. **Metadata restoration** — `chmod()`, `chown()`, `utimensat()` on the receiving side
8. **`--delete`** — sender tracks all sent paths; receiver walks destination tree and removes unlisted files/directories
9. **SSH transport** — `socketpair()` + `fork()` + `execvp("ssh", ...)` with `ControlMaster` and port support
10. **TLS transport** — OpenSSL `SSL_CTX` with TLS 1.2 minimum, optional CA verification, transparent `SSL_read`/`SSL_write` via `io_set_ssl()`
11. **Path traversal protection** — `has_path_traversal()` rejects any file path containing `..` components, preventing directory escape attacks
12. **Connection limiting** — server tracks active connections and rejects new ones beyond `max_connections` (default 100)
13. **Keep-alive** — idle connections receive periodic `STATUS_KEEPALIVE` to detect half-open TCP connections
14. **Abort handling** — `SIGINT` sets an abort flag; the next protocol operation sends `STATUS_ABORT` for clean server cleanup
15. **Atomic writes** — files are written to a `.tmp` suffix then atomically renamed via `rename()`, preventing partial files
16. **Backup** — before overwriting, existing files are moved to `--backup-dir` (or same directory with `~` suffix) preserving the original
## Security Features
### Path Traversal Protection
All received file paths are validated by `has_path_traversal()` before any disk operation. Any path containing `..` components is rejected with `STATUS_ERROR`, preventing directory escape attacks.
### TLS Certificate Verification
When `--ca` is provided, the server performs mutual TLS verification (`SSL_VERIFY_PEER` with depth 4). Without `--ca`, TLS is still encrypted but peer certificates are not verified.
### Connection Limits
The server enforces a maximum of 100 concurrent connections (configurable via `max_connections` in `Server`). When the limit is reached, new connections are immediately rejected and closed.
### Abort Handling
If the client receives `SIGINT` (Ctrl+C) during a transfer, it sends `STATUS_ABORT` to the server. The server then cleans up temporary files and exits the child process, preventing incomplete files from remaining on disk.
### Atomic Writes
Received files are written to a temporary path (suffixed with `.tmp`) and then atomically renamed to the final filename via `rename()`. This prevents partial or corrupted files from appearing at the destination if the transfer is interrupted.
## Build Requirements
- C11 compiler
- CMake >= 3.22
- zstd library
- OpenSSL (development headers and libraries)
- pthreads
- SSH client (for SSH transport mode only)
### Installing Dependencies
**Ubuntu/Debian:**
```bash
sudo apt install cmake build-essential libzstd-dev libssl-dev openssh-client
```
**Nix:**
```bash
nix-shell # provides zstd, openssl, cmake, gcc
```
## Building
```bash ```bash
cmake -B build -S . && cmake --build build -j$(nproc) cmake -B build -S .
cmake --build build -j$(nproc)
``` ```
## Running With Nix:
### Server (TCP mode)
```bash ```bash
./build/server nix-shell
cmake -B build -S .
cmake --build build -j$(nproc)
``` ```
### Server with TLS ### SSH transfer
The remote host must have `fastsync-server` available in `PATH`, or use
`--fastsync-server-path`. SSH starts `fastsync-server --stdio` in its remote
working directory, so use a destination below that directory unless the
remote server is otherwise configured with a matching authorized root.
```bash ```bash
./build/server --tls --cert server.pem --key server-key.pem ./build/client /path/to/source user@host:destination
``` ```
### Server via SSH ### TCP transfer
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.
Start the FastSync server:
### Client — SSH (rsync-style)
```bash ```bash
./build/client /path/to/send user@host:/path/to/receive ./build/server --destination-root /path/to -p 8080
``` ```
### Client — TCP Then run the client:
```bash ```bash
./build/client --source-dir /path/to/send --dest-dir /path/to/receive --save-to-disk ./build/client --server-host 127.0.0.1 --server-port 8080 \
--source-dir /path/to/source --dest-dir /path/to/destination \
--save-to-disk
``` ```
### Client — TCP with TLS ### TLS transfer
```bash ```bash
./build/server --destination-root /path/to --tls --cert server.pem --key server-key.pem -p 8443
./build/client --tls --cert client.pem --key client-key.pem --ca ca.pem \ ./build/client --tls --cert client.pem --key client-key.pem --ca ca.pem \
--source-dir /path/to/send --dest-dir /path/to/receive --save-to-disk --server-host example.com --server-port 8443 \
--source-dir /path/to/source --dest-dir /path/to/destination \
--save-to-disk
``` ```
### Common Options ## Common Workflows
These examples show the intended rsync-style workflow. Options marked as
FastSync-native are optional performance or transport extensions.
```bash ```bash
# Archive mode (compression + multithreading + metadata) # Basic synchronization
./build/client -a /path/to/send user@host:/path ./build/client /source/ /destination/
# Dry run # Archive-style synchronization (current FastSync archive behavior)
./build/client -n /path/to/send /path/to/receive ./build/client -a /source/ user@host:destination/
# With progress and custom chunk size # Preview a transfer without changing the destination
./build/client --progress --chunk-size 2097152 /src user@host:/dst ./build/client -n /source/ /destination/
# Exclude temporary files + delete extras on receiver # Exclude temporary and object files
./build/client --exclude "*.tmp" --exclude "*.o" --delete /src user@host:/dst ./build/client --exclude '*.tmp' --exclude '*.o' \
/source/ user@host:destination/
# Incremental sync (skip unchanged files) # Remove destination entries not present in the source
./build/client --incremental /src user@host:/dst ./build/client --delete /source/ user@host:destination/
# Bandwidth limit to 1 MB/s # Skip unchanged files using size and modification time
./build/client --bwlimit 1024 /src user@host:/dst ./build/client --incremental /source/ user@host:destination/
# With timeouts and stats # Verify content when size and time are not sufficient
./build/client --timeout 60 --contimeout 15 --stats /src user@host:/dst ./build/client --incremental --checksum /source/ user@host:destination/
# Backup overwritten files to a directory # Preserve supported mode and timestamp metadata
./build/client --backup --backup-dir /backups /src user@host:/dst ./build/client -M /source/ user@host:destination/
# Exclude patterns from file, limit depth # Keep backups of overwritten destination files
./build/client --exclude-from ignore.txt --max-depth 3 /src user@host:/dst ./build/client --backup --backup-dir backups \
/source/ user@host:destination/
# Log to file
./build/client --log-file /tmp/fastsync.log /src user@host:/dst
# All features
./build/client -a --progress --chunk-size 5242880 --exclude "*.log" --delete /src /dst
``` ```
## FastSync Extensions
FastSync-native options are intended to add performance or operational
features without changing the meaning of ordinary compatibility options.
| Option | Purpose |
|---|---|
| `-m` | Enable the multithreaded scanner/loader/sender pipeline. |
| `-c [level]`, `-z [level]` | Enable streaming zstd compression, levels 1-22. |
| `--compress-level <n>` | Set the zstd compression level. |
| `--chunk-size <bytes>` | Set the transfer chunk size. |
| `-s` | Enable FastSync chunk serialization. |
| `-f`, `--sendfile` | Use TCP `sendfile()` zero-copy transfer. Incompatible with compression and chunk serialization. |
| `--delta` | Use FastSync-native block delta transfer. Requires `--incremental`. |
| `--delta-block <bytes>` | Set the FastSync delta block size. |
| `--delta-max <bytes>` | Limit files eligible for FastSync delta transfer. |
| `--server-host <host>` | Select the TCP server host. |
| `--server-port <port>` | Select the TCP server port. |
| `--tls` | Enable TLS for TCP transport. |
| `--bwlimit <KB/s>` | Apply token-bucket bandwidth limiting. |
| `--progress` | Show transfer progress and throughput. |
| `--stats` | Print transfer statistics. |
| `--timeout <seconds>` | Set I/O timeout. |
| `--contimeout <seconds>` | Set connection timeout. |
Current short-option conflicts are tracked as compatibility work. In
particular, FastSync currently uses `-p` for SSH port, `-s` for chunk
serialization, and `-S` for sparse handling. These meanings must be reconciled
before FastSync can claim full rsync CLI compatibility.
## Client Options
### Selection and transfer
| Option | Description |
|---|---|
| `-a`, `--archive` | Enable current archive preset. Full rsync archive semantics are planned. |
| `-n`, `--dry-run` | Scan and report without writing files. |
| `--delete` | Request removal of destination entries absent from the source. The server must allow deletion. |
| `--exclude <pattern>` | Exclude matching paths. Repeatable. |
| `--include <pattern>` | Include matching paths. Repeatable. |
| `--exclude-from <file>` | Read exclude patterns from a file. |
| `--include-from <file>` | Read include patterns from a file. |
| `--max-size <bytes>` | Skip files larger than the limit. |
| `--min-size <bytes>` | Skip files smaller than the limit. |
| `--max-depth <n>` | Limit recursive scanning depth; zero means unlimited. |
| `--incremental` | Skip files matching destination size and mtime. |
| `--checksum` | Include xxHash64 content checks in incremental comparisons. |
| `--backup` | Back up overwritten files. |
| `--backup-dir <dir>` | Store backups under a separate directory. |
| `--suffix <suffix>` | Set the backup filename suffix. |
| `--partial` | Select partial-transfer handling. With `--partial-dir`, completed files are written there; resumable transfers are not implemented. |
| `--partial-dir <dir>` | Set a relative partial-transfer directory below the server destination root; use with `--partial`. |
| `--inplace` | Write directly to the destination instead of using a temporary file. |
### Metadata and links
| Option | Description |
|---|---|
| `-M`, `--preserve` | Preserve supported file metadata, currently mode and modification time. |
| `-l`, `--links` | Request symlink preservation; link-target transfer remains incomplete. |
| `--copy-links` | Copy symlink referents. |
| `--safe-links` | Skip symlinks that point outside the transfer tree. |
| `--copy-unsafe-links` | Copy unsafe symlink referents. |
| `-S`, `--sparse` | Request sparse-file handling; full hole preservation is planned. |
### Output and logging
| Option | Description |
|---|---|
| `-v`, `--verbose` | Enable debug logging. |
| `--progress` | Show live transfer progress. |
| `--stats` | Print transfer statistics. |
| `--log-file <path>` | Write log output to a file. |
| `-V`, `--version` | Print the FastSync protocol version. |
| `--help` | Print command usage. |
### Paths and transport
| Option | Description |
|---|---|
| `-p <port>` | SSH port in the current CLI. This conflicts with rsync's `-p` permissions option and is planned for correction. |
| `--fastsync-server-path <path>` | Remote FastSync server path for SSH mode. |
| `--source-dir <path>` | Set the source directory explicitly. |
| `--dest-dir <path>` | Set the destination directory explicitly. |
| `--save-to-disk` | Enable server-side disk persistence. |
| `--server-host <host>` | TCP server address. |
| `--server-port <port>` | TCP server port. |
| `--tls` | Enable TLS. Requires `--cert` and `--key`. |
| `--cert <path>` | TLS certificate file. |
| `--key <path>` | TLS private key file. |
| `--ca <path>` | CA file for peer verification. |
## Server Options
| Option | Description |
|---|---|
| `--stdio` | Serve one SSH connection over standard input/output. |
| `-p <port>` | TCP listen port. |
| `--tls` | Enable TLS. |
| `--cert <path>` | TLS certificate file. |
| `--key <path>` | TLS private key file. |
| `--ca <path>` | CA file for peer verification. |
| `--destination-root <path>` | Confine received files to this server-side root; defaults to the current directory. |
| `--allow-delete` | Permit client delete manifests. Deletion is refused by default. |
| `-v`, `--verbose` | Enable debug logging. |
| `--help` | Print server usage. |
## Architecture
### Client
- Recursively scans the source tree with include, exclude, size, and depth
filters.
- Sends individual files or serialized chunks.
- Performs incremental checks and optional content checksums.
- Uses a multithreaded producer-consumer pipeline when requested.
- Sends over TCP, TLS-wrapped TCP, or an SSH subprocess.
- Supports progress, statistics, backups, timeouts, and bandwidth limiting.
### Server
- Runs as a TCP listener or one-shot SSH `--stdio` server.
- Receives and reassembles files and decompresses streaming zstd data.
- Applies supported metadata and writes files through a confined destination
root.
- Uses temporary files and atomic rename by default.
- Handles delete manifests only when explicitly authorized.
- Enforces connection, message-size, and path-safety limits.
## Protocol and Security
FastSync protocol version `2.2.0` is shared by the client and server. The
current protocol is sender-driven and includes configuration negotiation,
incremental checks, checksums, manifests, keep-alives, abort handling, and
FastSync-native delta messages. Client and server versions must currently
match exactly.
TLS provides encrypted TCP transport. Supplying `--ca` enables certificate
verification; without it, traffic is encrypted but peer identity is not
verified. Use certificate verification for deployments where authentication
matters. The default TCP transport is not encrypted.
The receiver protects its destination root with path validation, `openat()`
directory traversal, `O_NOFOLLOW`, temporary files, and atomic renames. Delete
operations require the server's explicit `--allow-delete` policy.
## Compatibility Roadmap
The project will reach the drop-in replacement goal in stages:
1. Correct rsync option meanings, including short options, combined options,
and `--option=value` syntax.
2. Add differential tests that compare FastSync and rsync contents, metadata,
links, deletes, filters, dry runs, and exit codes.
3. Make `-a` implement the expected recursive, links, permissions, times,
owner/group, and supported special-file behavior.
4. Complete symlink, sparse-file, metadata, delete-policy, and resumable-write
semantics.
5. Add rsync remote-shell and daemon protocol interoperability.
6. Keep FastSync performance options as negotiated, optional extensions.
The exhaustive implementation matrix and compatibility notes are in
[`RSYNC_COMPAT.md`](RSYNC_COMPAT.md).
## Testing ## Testing
```bash Run the unit test binary:
# Unit tests (18 suites — array_list, chunk, compression, config, data, delta, file, glob,
# metadata, property, protocol, queue, robustness, scanner,
# shared_utils, stress, transport_tcp, transport_ssh, transport_tls)
./build/tests
# Integration + benchmark suite ```bash
python3 test.py ./build/tests
``` ```
The benchmark prints throughput metrics, best configuration, and speedup vs rsync. Run the Python integration suite:
## Performance Considerations ```bash
python3 -m pytest tests/
```
1. Chunk size (~10 MB default) balances memory and transfer efficiency For stricter local validation:
2. Compression level trades CPU for bandwidth
3. `sendfile()` bypasses userspace — ~2× faster on localhost for large files
4. Multithreading scales with core count and uses memory-based pipeline sizing
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
8. Incremental sync eliminates redundant transfers entirely
9. Batch incremental reduces round-trips by grouping multiple checks into one message
10. Bandwidth limiting uses token-bucket with nanosleep for accurate throttling
11. Atomic writes add a single `rename()` per file — negligible overhead
12. Path traversal check is O(n) in path length with negligible cost
## Benchmark Results ```bash
cmake -B build-strict -S . -DSTRICT_WARNINGS=ON
cmake --build build-strict -j$(nproc)
cmake -B build-asan -S . -DSANITIZER=address
cmake --build build-asan -j$(nproc)
```
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. The benchmark tool compares FastSync configurations with rsync under
controlled local and network conditions:
### LAN (1000 Mbit, 20 ms ±1 ms, 0.1% loss) ```bash
python3 benchmark/bench.py --help
```
| Configuration | Time | vs rsync (archive) | vs rsync (compress) | Benchmark results measure transfer performance only. They do not establish
|---|---|---|---| rsync protocol or filesystem-semantic compatibility.
| **Best: `-m -c`** | **0.20 s** | **11.2× faster** | **3.6× faster** |
| Compression (`-c`) | 0.31 s | 7.3× faster | 2.3× faster |
| Standard | 1.27 s | 1.8× faster | — |
| rsync (archive) | 2.27 s | — | — |
| rsync (archive + compress) | 0.72 s | — | — |
### WAN (100 Mbit, 50 ms ±10 ms, 1% loss) ## Performance Guidance
| Configuration | Time | vs rsync (archive) | vs rsync (compress) | - Use `-m` for workloads with many files or enough CPU parallelism.
|---|---|---|---| - Use `-c` or `-z` when network bandwidth is more constrained than CPU.
| **Best: `-m -c`** | **0.39 s** | **44.8× faster** | **3.8× faster** | - Tune `--chunk-size` for file sizes, memory limits, and network latency.
| Compression (`-c`) | 0.64 s | 27.3× faster | 2.3× faster | - Use `-f` for large uncompressed TCP transfers where zero-copy I/O helps.
| Standard | 7.12 s | 2.4× faster | — | - Use `--incremental` to avoid retransmitting unchanged files.
| rsync (archive) | 17.44 s | — | — | - Use `--delta` for changed files when both endpoints are FastSync peers.
| rsync (archive + compress) | 1.47 s | — | — | - Use `--bwlimit` when sharing a link with other traffic.
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. Always validate the compatibility behavior required by a deployment before
replacing an existing rsync job.
+2 -2
View File
@@ -152,14 +152,14 @@ This document maps rsync's full feature set to FastSync's current implementation
| Flag | Rsync Description | FastSync Status | Notes | | Flag | Rsync Description | FastSync Status | Notes |
|------|-------------------|-----------------|-------| |------|-------------------|-----------------|-------|
| `-S`, `--sparse` | Sparse block handling | ✅ Implemented | `preserve_sparse` config field | | `-S`, `--sparse` | Sparse block handling | ⚠️ Partial | Flag is accepted, but full hole preservation is not implemented |
| `--preallocate` | Allocate dest files before writing | ❌ Not Implemented | | | `--preallocate` | Allocate dest files before writing | ❌ Not Implemented | |
## 11. Checksum & Comparison ## 11. Checksum & Comparison
| Flag | Rsync Description | FastSync Status | Notes | | Flag | Rsync Description | FastSync Status | Notes |
|------|-------------------|-----------------|-------| |------|-------------------|-----------------|-------|
| `--checksum` | Skip based on checksum | ❌ Not Implemented | Removed because it had no effect; `-c` means compression | | `--checksum` | Skip based on checksum | ✅ Implemented | With `--incremental`, compares xxHash64 content checksums; `-c` remains compression |
| `--checksum-choice=STR` | Choose checksum algorithm | ❌ Not Implemented | xxHash used internally | | `--checksum-choice=STR` | Choose checksum algorithm | ❌ Not Implemented | xxHash used internally |
| `--compare-dest=DIR` | Compare dest files relative to DIR | ❌ Not Implemented | Removed because it had no effect | | `--compare-dest=DIR` | Compare dest files relative to DIR | ❌ Not Implemented | Removed because it had no effect |
| `--copy-dest=DIR` | Include copies of unchanged files | ❌ Not Implemented | Removed because it had no effect | | `--copy-dest=DIR` | Include copies of unchanged files | ❌ Not Implemented | Removed because it had no effect |
+7 -116
View File
@@ -1,7 +1,4 @@
#include "array_list.h"
#include "chunk.h"
#include "config.h" #include "config.h"
#include "data.h"
#include "file.h" #include "file.h"
#include "log.h" #include "log.h"
#include "multiprocessing.h" #include "multiprocessing.h"
@@ -29,9 +26,12 @@ static bool path_is_within(const char* root, const char* path) {
return strncmp(root, path, n) == 0 && (path[n] == '\0' || path[n] == '/'); return strncmp(root, path, n) == 0 && (path[n] == '\0' || path[n] == '/');
} }
static bool valid_batch_path(const char* path) { static bool save_received_file(File* file, void* context) {
return path && path[0] != '\0' && path[0] != '/' && !has_path_traversal(path) && Config* config = context;
strchr(path, '\0') == path + strlen(path); if (config->save_to_disk && !file_save_to_disk(config->receive_root_directory, file, config))
return false;
file_destroy(file);
return true;
} }
static bool __attribute__((unused)) configure_authorization(const char* root) { static bool __attribute__((unused)) configure_authorization(const char* root) {
@@ -53,116 +53,7 @@ static bool __attribute__((unused)) configure_authorization(const char* root) {
} }
int receive_files(Config* config, int fd) { int receive_files(Config* config, int fd) {
Status status; return receive_files_common(config, fd, save_received_file, config, true);
if (!receive_status(fd, &status))
return -1;
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH) {
if (status == STATUS_KEEPALIVE) {
send_status(fd, STATUS_KEEPALIVE);
goto next;
}
if (status == STATUS_ABORT) {
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
return -1;
}
if (status == STATUS_CHECK) {
bool skipped;
File* file = receive_incremental_check(fd, config, &skipped);
if (skipped)
goto next;
if (file == NULL && !skipped)
return -1;
if (config->save_to_disk &&
!file_save_to_disk(config->receive_root_directory, file, config)) {
file_destroy(file);
send_status(fd, STATUS_ERROR);
return -1;
}
file_destroy(file);
} else if (status == STATUS_CHUNK) {
Chunk* chunk = receive_chunk_data(fd, config);
if (chunk == NULL) {
send_status(fd, STATUS_ERROR);
return -1;
}
for (int i = 0; i < chunk->element_count; i++) {
if (config->save_to_disk &&
!file_save_to_disk(config->receive_root_directory, chunk->items[i], config)) {
chunk_destroy(chunk);
send_status(fd, STATUS_ERROR);
return -1;
}
}
chunk_destroy(chunk);
} else if (status == STATUS_CHECK_BATCH) {
int count;
/* Batch framing has no checksum field yet; never silently downgrade a
checksum-enabled transfer into mtime-only matching. */
if (config->checksum || !receive_int(fd, &count) || count < 0 || count > MAX_MANIFEST_ENTRIES)
return -1;
for (int i = 0; i < count; i++) {
char* check_path = receive_str(fd);
if (!check_path)
return -1;
unsigned long long check_size;
long long check_mtime;
if (!receive_n_data(fd, &check_size, sizeof(check_size)) ||
!receive_n_data(fd, &check_mtime, sizeof(check_mtime))) {
free(check_path);
return -1;
}
if (!valid_batch_path(check_path)) {
free(check_path);
send_status(fd, STATUS_ERROR);
return -1;
}
char* full_path = path_cat(config->receive_root_directory, check_path);
struct stat st;
bool has_old = full_path && lstat(full_path, &st) == 0;
bool match = has_old && (unsigned long long)st.st_size == check_size &&
(long long)st.st_mtime == check_mtime;
bool sent = send_status(fd, match ? STATUS_OK : STATUS_NEXT);
free(full_path);
free(check_path);
if (!sent)
return -1;
}
goto next;
} else {
File* file = file_receive(config, fd);
if (file == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to receive file");
send_status(fd, STATUS_ERROR);
return -1;
}
if (config->save_to_disk &&
!file_save_to_disk(config->receive_root_directory, file, config)) {
file_destroy(file);
send_status(fd, STATUS_ERROR);
return -1;
}
file_destroy(file);
}
next:
if (!receive_status(fd, &status)) {
send_status(fd, STATUS_ERROR);
return -1;
}
}
if (status == STATUS_MANIFEST) {
if (receive_manifest(fd, config, &status) != 0)
return -1;
}
if (status != STATUS_FINISHED) {
log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status");
send_status(fd, STATUS_ERROR);
return -1;
}
send_status(fd, STATUS_OK);
return 0;
} }
void handler(int file_descriptor) { void handler(int file_descriptor) {
+87 -254
View File
@@ -173,103 +173,39 @@ void config_delete(Config* config) {
free(config); free(config);
} }
/* Wire format order (must match config_receive and be updated when PROTOCOL_VERSION bumps): static bool config_send_string(int file_descriptor, const char* value, bool optional) {
* version, send_directory, receive_root_directory, save_to_disk, use_multithreading, return send_str(file_descriptor, value ? value : (optional ? "" : NULL));
* use_chunk_serialization, use_compression, use_metadata, compression_level, chunk_size, }
* use_sendfile, use_delete, use_incremental, use_delta, delta_block_size, delta_max_file_size,
* backup, backup_dir, follow_symlinks, copy_links, safe_links, copy_unsafe_links, #define CONFIG_SEND_VERSION(field) \
* preserve_hard_links, preserve_acls, preserve_xattrs, preserve_devices, preserve_sparse, do { \
* update, inplace, append, append_verify, delete_excluded, delete_after, max_delete, relative, if (!config_send_string(file_descriptor, config->field, false)) \
* prune_empty_dirs, temp_dir, partial, partial_dir, suffix, delete_before, checksum, return false; \
* compress_choice, status } while (0);
*/ #define CONFIG_SEND_STRING(field) \
do { \
if (!config_send_string(file_descriptor, config->field, false)) \
return false; \
} while (0);
#define CONFIG_SEND_OPTIONAL_STRING(field) \
do { \
if (!config_send_string(file_descriptor, config->field, true)) \
return false; \
} while (0);
#define CONFIG_SEND_INTEGER(field) \
do { \
if (!send_int(file_descriptor, config->field)) \
return false; \
} while (0);
#define CONFIG_SEND_DATA(field) \
do { \
if (!send_n_data(file_descriptor, &config->field, sizeof(config->field))) \
return false; \
} while (0);
bool config_send(int file_descriptor, const Config* config) { bool config_send(int file_descriptor, const Config* config) {
if (!send_str(file_descriptor, config->version)) CONFIG_WIRE_FIELDS(CONFIG_SEND_VERSION, CONFIG_SEND_STRING, CONFIG_SEND_OPTIONAL_STRING,
return false; CONFIG_SEND_INTEGER, CONFIG_SEND_DATA)
if (!send_str(file_descriptor, config->send_directory))
return false;
if (!send_str(file_descriptor, config->receive_root_directory))
return false;
if (!send_int(file_descriptor, config->save_to_disk))
return false;
if (!send_int(file_descriptor, config->use_multithreading))
return false;
if (!send_int(file_descriptor, config->use_chunk_serialization))
return false;
if (!send_int(file_descriptor, config->use_compression))
return false;
if (!send_int(file_descriptor, config->use_metadata))
return false;
if (!send_int(file_descriptor, config->compression_level))
return false;
if (!send_n_data(file_descriptor, &config->chunk_size, sizeof(config->chunk_size)))
return false;
if (!send_int(file_descriptor, config->use_sendfile))
return false;
if (!send_int(file_descriptor, config->use_delete))
return false;
if (!send_int(file_descriptor, config->use_incremental))
return false;
if (!send_int(file_descriptor, config->use_delta))
return false;
if (!send_n_data(file_descriptor, &config->delta_block_size, sizeof(config->delta_block_size)))
return false;
if (!send_n_data(file_descriptor, &config->delta_max_file_size, sizeof(unsigned long long)))
return false;
if (!send_int(file_descriptor, config->backup))
return false;
if (!send_str(file_descriptor, config->backup_dir ? config->backup_dir : ""))
return false;
if (!send_int(file_descriptor, config->follow_symlinks))
return false;
if (!send_int(file_descriptor, config->copy_links))
return false;
if (!send_int(file_descriptor, config->safe_links))
return false;
if (!send_int(file_descriptor, config->copy_unsafe_links))
return false;
if (!send_int(file_descriptor, config->preserve_hard_links))
return false;
if (!send_int(file_descriptor, config->preserve_acls))
return false;
if (!send_int(file_descriptor, config->preserve_xattrs))
return false;
if (!send_int(file_descriptor, config->preserve_devices))
return false;
if (!send_int(file_descriptor, config->preserve_sparse))
return false;
if (!send_int(file_descriptor, config->update))
return false;
if (!send_int(file_descriptor, config->inplace))
return false;
if (!send_int(file_descriptor, config->append))
return false;
if (!send_int(file_descriptor, config->append_verify))
return false;
if (!send_int(file_descriptor, config->delete_excluded))
return false;
if (!send_int(file_descriptor, config->delete_after))
return false;
if (!send_n_data(file_descriptor, &config->max_delete, sizeof(config->max_delete)))
return false;
if (!send_int(file_descriptor, config->relative))
return false;
if (!send_int(file_descriptor, config->prune_empty_dirs))
return false;
if (!send_str(file_descriptor, config->temp_dir ? config->temp_dir : ""))
return false;
if (!send_int(file_descriptor, config->partial))
return false;
if (!send_str(file_descriptor, config->partial_dir ? config->partial_dir : ""))
return false;
if (!send_str(file_descriptor, config->suffix ? config->suffix : ""))
return false;
if (!send_int(file_descriptor, config->delete_before))
return false;
if (!send_int(file_descriptor, config->checksum))
return false;
if (!send_str(file_descriptor, config->compress_choice ? config->compress_choice : ""))
return false;
Status status; Status status;
if (!receive_status(file_descriptor, &status)) if (!receive_status(file_descriptor, &status))
return false; return false;
@@ -280,160 +216,60 @@ bool config_send(int file_descriptor, const Config* config) {
return true; return true;
} }
/* Wire format order: see the comment above config_send. */ #undef CONFIG_SEND_VERSION
#undef CONFIG_SEND_STRING
#undef CONFIG_SEND_OPTIONAL_STRING
#undef CONFIG_SEND_INTEGER
#undef CONFIG_SEND_DATA
static bool config_receive_string(int file_descriptor, char** destination) {
char* value = receive_str(file_descriptor);
if (!value)
return false;
*destination = value;
return true;
}
#define CONFIG_RECEIVE_VERSION(field) \
do { \
free(config->field); \
config->field = receive_str(file_descriptor); \
if (!config->field) \
goto error; \
if (strcmp(config->field, PROTOCOL_VERSION) != 0) { \
fprintf(stderr, "Protocol version mismatch: client=%s, server=%s\n", config->field, \
PROTOCOL_VERSION); \
config_delete(config); \
send_status(file_descriptor, STATUS_ERROR); \
return NULL; \
} \
} while (0);
#define CONFIG_RECEIVE_STRING(field) \
do { \
if (!config_receive_string(file_descriptor, &config->field)) \
goto error; \
} while (0);
#define CONFIG_RECEIVE_OPTIONAL_STRING(field) CONFIG_RECEIVE_STRING(field)
#define CONFIG_RECEIVE_INTEGER(field) \
do { \
if (!receive_int(file_descriptor, &tmp)) \
goto error; \
config->field = tmp; \
} while (0);
#define CONFIG_RECEIVE_DATA(field) \
do { \
if (!receive_n_data(file_descriptor, &config->field, sizeof(config->field))) \
goto error; \
} while (0);
Config* config_receive(int file_descriptor) { Config* config_receive(int file_descriptor) {
Config* config = (Config*)malloc(sizeof(Config)); Config* config = (Config*)malloc(sizeof(Config));
if (config == NULL) if (config == NULL)
return NULL; return NULL;
config_set_defaults(config); config_set_defaults(config);
free(config->version);
config->version = receive_str(file_descriptor);
if (!config->version) {
free(config->server_host);
free(config);
return NULL;
}
if (strcmp(config->version, PROTOCOL_VERSION) != 0) {
fprintf(stderr, "Protocol version mismatch: client=%s, server=%s\n", config->version,
PROTOCOL_VERSION);
free(config->version);
free(config->server_host);
free(config);
send_status(file_descriptor, STATUS_ERROR);
return NULL;
}
config->send_directory = receive_str(file_descriptor);
if (!config->send_directory) {
free(config->version);
free(config->server_host);
free(config);
return NULL;
}
config->receive_root_directory = receive_str(file_descriptor);
if (!config->receive_root_directory) {
free(config->version);
free(config->send_directory);
free(config->server_host);
free(config);
return NULL;
}
int tmp; int tmp;
if (!receive_int(file_descriptor, &tmp)) CONFIG_WIRE_FIELDS(CONFIG_RECEIVE_VERSION, CONFIG_RECEIVE_STRING, CONFIG_RECEIVE_OPTIONAL_STRING,
goto error; CONFIG_RECEIVE_INTEGER, CONFIG_RECEIVE_DATA)
config->save_to_disk = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->use_multithreading = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->use_chunk_serialization = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->use_compression = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->use_metadata = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->compression_level = tmp;
if (!receive_n_data(file_descriptor, &config->chunk_size, sizeof(config->chunk_size)))
goto error;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->use_sendfile = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->use_delete = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->use_incremental = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->use_delta = tmp;
if (!receive_n_data(file_descriptor, &config->delta_block_size, sizeof(config->delta_block_size)))
goto error;
if (!receive_n_data(file_descriptor, &config->delta_max_file_size, sizeof(unsigned long long)))
goto error;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->backup = tmp;
config->backup_dir = receive_str(file_descriptor);
if (config->backup_dir == NULL)
goto error;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->follow_symlinks = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->copy_links = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->safe_links = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->copy_unsafe_links = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->preserve_hard_links = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->preserve_acls = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->preserve_xattrs = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->preserve_devices = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->preserve_sparse = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->update = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->inplace = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->append = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->append_verify = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->delete_excluded = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->delete_after = tmp;
if (!receive_n_data(file_descriptor, &config->max_delete, sizeof(config->max_delete)))
goto error;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->relative = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->prune_empty_dirs = tmp;
config->temp_dir = receive_str(file_descriptor);
if (config->temp_dir == NULL)
goto error;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->partial = tmp;
config->partial_dir = receive_str(file_descriptor);
if (config->partial_dir == NULL)
goto error;
config->suffix = receive_str(file_descriptor);
if (config->suffix == NULL)
goto error;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->delete_before = tmp;
if (!receive_int(file_descriptor, &tmp))
goto error;
config->checksum = tmp;
config->compress_choice = receive_str(file_descriptor);
if (config->compress_choice == NULL)
goto error;
if (config->compress_choice[0] != '\0' && strcmp(config->compress_choice, "zstd") != 0 && if (config->compress_choice[0] != '\0' && strcmp(config->compress_choice, "zstd") != 0 &&
strcmp(config->compress_choice, "none") != 0) { strcmp(config->compress_choice, "none") != 0) {
fprintf(stderr, "Unsupported compression choice: %s\n", config->compress_choice); fprintf(stderr, "Unsupported compression choice: %s\n", config->compress_choice);
@@ -445,15 +281,12 @@ Config* config_receive(int file_descriptor) {
return config; return config;
error: error:
free(config->version); config_delete(config);
free(config->send_directory);
free(config->receive_root_directory);
free(config->server_host);
free(config->backup_dir);
free(config->temp_dir);
free(config->partial_dir);
free(config->suffix);
free(config->compress_choice);
free(config);
return NULL; return NULL;
} }
#undef CONFIG_RECEIVE_VERSION
#undef CONFIG_RECEIVE_STRING
#undef CONFIG_RECEIVE_OPTIONAL_STRING
#undef CONFIG_RECEIVE_INTEGER
#undef CONFIG_RECEIVE_DATA
+47
View File
@@ -128,6 +128,53 @@ typedef struct Config {
char* compress_choice; char* compress_choice;
} Config; } Config;
/* Keep the on-wire field order in one place. The first three strings require
* values when sent; optional strings are encoded as empty strings when NULL. */
#define CONFIG_WIRE_FIELDS(VERSION, STRING, OPTIONAL_STRING, INTEGER, DATA) \
VERSION(version) \
STRING(send_directory) \
STRING(receive_root_directory) \
INTEGER(save_to_disk) \
INTEGER(use_multithreading) \
INTEGER(use_chunk_serialization) \
INTEGER(use_compression) \
INTEGER(use_metadata) \
INTEGER(compression_level) \
DATA(chunk_size) \
INTEGER(use_sendfile) \
INTEGER(use_delete) \
INTEGER(use_incremental) \
INTEGER(use_delta) \
DATA(delta_block_size) \
DATA(delta_max_file_size) \
INTEGER(backup) \
OPTIONAL_STRING(backup_dir) \
INTEGER(follow_symlinks) \
INTEGER(copy_links) \
INTEGER(safe_links) \
INTEGER(copy_unsafe_links) \
INTEGER(preserve_hard_links) \
INTEGER(preserve_acls) \
INTEGER(preserve_xattrs) \
INTEGER(preserve_devices) \
INTEGER(preserve_sparse) \
INTEGER(update) \
INTEGER(inplace) \
INTEGER(append) \
INTEGER(append_verify) \
INTEGER(delete_excluded) \
INTEGER(delete_after) \
DATA(max_delete) \
INTEGER(relative) \
INTEGER(prune_empty_dirs) \
OPTIONAL_STRING(temp_dir) \
INTEGER(partial) \
OPTIONAL_STRING(partial_dir) \
OPTIONAL_STRING(suffix) \
INTEGER(delete_before) \
INTEGER(checksum) \
OPTIONAL_STRING(compress_choice)
#define PROTOCOL_VERSION "2.2.0" #define PROTOCOL_VERSION "2.2.0"
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
+106 -122
View File
@@ -13,11 +13,107 @@
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <threads.h> #include <threads.h>
#include <sys/stat.h>
static bool valid_batch_path(const char* path) { static bool valid_batch_path(const char* path) {
return path && path[0] != '\0' && path[0] != '/' && !has_path_traversal(path); return path && path[0] != '\0' && path[0] != '/' && !has_path_traversal(path);
} }
static bool handle_batch_checks(int file_descriptor, const Config* config) {
int count;
if (config->checksum || !receive_int(file_descriptor, &count) || count < 0 ||
count > MAX_MANIFEST_ENTRIES)
return false;
for (int i = 0; i < count; i++) {
char* check_path = receive_str(file_descriptor);
if (!check_path)
return false;
unsigned long long check_size;
long long check_mtime;
bool received = receive_n_data(file_descriptor, &check_size, sizeof(check_size)) &&
receive_n_data(file_descriptor, &check_mtime, sizeof(check_mtime));
if (!received || !valid_batch_path(check_path)) {
free(check_path);
if (received)
send_status(file_descriptor, STATUS_ERROR);
return false;
}
char* full_path = path_cat(config->receive_root_directory, check_path);
struct stat st;
bool has_old = full_path && lstat(full_path, &st) == 0;
bool match = has_old && (unsigned long long)st.st_size == check_size &&
(long long)st.st_mtime == check_mtime;
bool sent = send_status(file_descriptor, match ? STATUS_OK : STATUS_NEXT);
free(full_path);
free(check_path);
if (!sent)
return false;
}
return true;
}
int receive_files_common(const Config* config, int file_descriptor, ReceivedFileHandler handler,
void* context, bool send_completion_status) {
Status status;
if (!receive_status(file_descriptor, &status))
return -1;
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH) {
if (status == STATUS_KEEPALIVE) {
if (!send_status(file_descriptor, STATUS_KEEPALIVE))
return -1;
} else if (status == STATUS_ABORT) {
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
return -1;
} else if (status == STATUS_CHECK_BATCH) {
if (!handle_batch_checks(file_descriptor, config))
return -1;
} else {
if (status == STATUS_CHUNK) {
Chunk* chunk = receive_chunk_data(file_descriptor, config);
if (!chunk)
return -1;
for (int i = 0; i < chunk->element_count; i++) {
File* file = chunk->items[i];
chunk->items[i] = NULL;
if (!handler(file, context)) {
file_destroy(file);
chunk_destroy(chunk);
return -1;
}
}
chunk_destroy(chunk);
} else {
bool skipped = false;
File* file = status == STATUS_CHECK
? receive_incremental_check(file_descriptor, config, &skipped)
: file_receive(config, file_descriptor);
if (!skipped) {
if (!file || !handler(file, context)) {
file_destroy(file);
if (status == STATUS_NEXT)
log_message(LOG_LEVEL_ERROR, "Failed to receive file");
return -1;
}
}
}
}
if (!receive_status(file_descriptor, &status))
return -1;
}
if (status == STATUS_MANIFEST && receive_manifest(file_descriptor, config, &status) != 0)
return -1;
if (status != STATUS_FINISHED) {
log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status");
if (send_completion_status)
send_status(file_descriptor, STATUS_ERROR);
return -1;
}
if (send_completion_status && !send_status(file_descriptor, STATUS_OK))
return -1;
return 0;
}
PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner,
Queue* queue_loader) { Queue* queue_loader) {
PipelineContextSender* context = malloc(sizeof(PipelineContextSender)); PipelineContextSender* context = malloc(sizeof(PipelineContextSender));
@@ -137,41 +233,14 @@ void pipeline_context_receiver_destroy(PipelineContextReceiver* context) {
free(context); free(context);
} }
static bool receive_chunk_enqueue(int file_descriptor, PipelineContextReceiver* context) { static bool enqueue_received_file(File* file, void* context) {
Chunk* chunk = receive_chunk_data(file_descriptor, context->config); PipelineContextReceiver* receiver = context;
if (chunk == NULL) return queue_enqueue_multithreaded_cancel(receiver->queue, file, &receiver->mutex,
return false; &receiver->condition_not_empty,
&receiver->condition_not_full, &receiver->cancelled);
for (int i = 0; i < chunk->element_count; i++) {
File* file = chunk->items[i];
chunk->items[i] = NULL;
if (!queue_enqueue_multithreaded_cancel(context->queue, file, &context->mutex,
&context->condition_not_empty,
&context->condition_not_full, &context->cancelled)) {
file_destroy(file);
chunk_destroy(chunk);
return false;
}
}
chunk_destroy(chunk);
return true;
}
static void receiver_thread_fail(PipelineContextReceiver* context) {
mtx_lock(&context->mutex);
atomic_store(&context->cancelled, true);
context->receiver_done = true;
cnd_broadcast(&context->condition_not_empty);
cnd_broadcast(&context->condition_not_full);
mtx_unlock(&context->mutex);
} }
int receive_thread(void* pipeline_context) { int receive_thread(void* pipeline_context) {
#define RECEIVE_THREAD_FAIL() \
do { \
receiver_thread_fail(context); \
return thrd_error; \
} while (0)
PipelineContextReceiver* context = (PipelineContextReceiver*)pipeline_context; PipelineContextReceiver* context = (PipelineContextReceiver*)pipeline_context;
if (context->ssl) if (context->ssl)
io_set_ssl(context->ssl); io_set_ssl(context->ssl);
@@ -180,100 +249,15 @@ int receive_thread(void* pipeline_context) {
const Config* config = context->config; const Config* config = context->config;
mtx_unlock(&context->mutex); mtx_unlock(&context->mutex);
Status status; int result = receive_files_common(config, file_descriptor, enqueue_received_file, context, false);
if (!receive_status(file_descriptor, &status))
RECEIVE_THREAD_FAIL();
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH) {
if (status == STATUS_KEEPALIVE) {
if (!send_status(file_descriptor, STATUS_KEEPALIVE))
RECEIVE_THREAD_FAIL();
goto next;
}
if (status == STATUS_ABORT) {
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
RECEIVE_THREAD_FAIL();
}
if (status == STATUS_CHECK) {
bool skipped;
File* file = receive_incremental_check(file_descriptor, config, &skipped);
if (!skipped) {
if (file == NULL)
RECEIVE_THREAD_FAIL();
if (!queue_enqueue_multithreaded_cancel(
context->queue, file, &context->mutex, &context->condition_not_empty,
&context->condition_not_full, &context->cancelled)) {
file_destroy(file);
RECEIVE_THREAD_FAIL();
}
}
} else if (status == STATUS_CHUNK) {
if (!receive_chunk_enqueue(file_descriptor, context))
RECEIVE_THREAD_FAIL();
} else if (status == STATUS_CHECK_BATCH) {
int count;
if (config->checksum || !receive_int(file_descriptor, &count) || count < 0 ||
count > MAX_MANIFEST_ENTRIES)
RECEIVE_THREAD_FAIL();
for (int i = 0; i < count; i++) {
char* check_path = receive_str(file_descriptor);
if (!check_path)
RECEIVE_THREAD_FAIL();
unsigned long long check_size;
long long check_mtime;
if (!receive_n_data(file_descriptor, &check_size, sizeof(check_size)) ||
!receive_n_data(file_descriptor, &check_mtime, sizeof(check_mtime))) {
free(check_path);
RECEIVE_THREAD_FAIL();
}
if (!valid_batch_path(check_path)) {
free(check_path);
if (!send_status(file_descriptor, STATUS_ERROR))
RECEIVE_THREAD_FAIL();
RECEIVE_THREAD_FAIL();
}
char* full_path = path_cat(config->receive_root_directory, check_path);
struct stat st;
bool has_old = full_path && lstat(full_path, &st) == 0;
bool match = has_old && (unsigned long long)st.st_size == check_size &&
(long long)st.st_mtime == check_mtime;
if (!send_status(file_descriptor, match ? STATUS_OK : STATUS_NEXT))
RECEIVE_THREAD_FAIL();
free(full_path);
free(check_path);
}
goto next;
} else {
File* file = file_receive(config, file_descriptor);
if (file) {
if (!queue_enqueue_multithreaded_cancel(
context->queue, file, &context->mutex, &context->condition_not_empty,
&context->condition_not_full, &context->cancelled)) {
file_destroy(file);
receiver_thread_fail(context);
return thrd_error;
}
} else {
log_message(LOG_LEVEL_ERROR, "Failed to receive file");
RECEIVE_THREAD_FAIL();
}
}
next:
if (!receive_status(file_descriptor, &status))
RECEIVE_THREAD_FAIL();
}
if (status == STATUS_MANIFEST) {
if (receive_manifest(file_descriptor, config, &status) != 0)
RECEIVE_THREAD_FAIL();
}
if (status != STATUS_FINISHED)
RECEIVE_THREAD_FAIL();
mtx_lock(&context->mutex); mtx_lock(&context->mutex);
if (result != 0)
atomic_store(&context->cancelled, true);
context->receiver_done = true; context->receiver_done = true;
cnd_signal(&context->condition_not_empty); cnd_signal(&context->condition_not_empty);
cnd_broadcast(&context->condition_not_full);
mtx_unlock(&context->mutex); mtx_unlock(&context->mutex);
#undef RECEIVE_THREAD_FAIL return result == 0 ? thrd_success : thrd_error;
return thrd_success;
} }
int write_thread(void* pipeline_context) { int write_thread(void* pipeline_context) {
+5
View File
@@ -42,6 +42,11 @@ typedef struct PipelineContextReceiver {
atomic_bool cancelled; atomic_bool cancelled;
} PipelineContextReceiver; } PipelineContextReceiver;
typedef bool (*ReceivedFileHandler)(File* file, void* context);
int receive_files_common(const Config* config, int file_descriptor, ReceivedFileHandler handler,
void* context, bool send_completion_status);
PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner,
Queue* queue_loader); Queue* queue_loader);
void pipeline_context_sender_destroy(PipelineContextSender* context); void pipeline_context_sender_destroy(PipelineContextSender* context);
+5 -1
View File
@@ -173,7 +173,7 @@ Client* client_create() {
return client; return client;
} }
bool client_connect(Client* client, char* host, int port) { bool tcp_connect_socket(Client* client, const char* host, int port) {
struct addrinfo hints; struct addrinfo hints;
struct addrinfo* result; struct addrinfo* result;
memset(&hints, 0, sizeof(hints)); memset(&hints, 0, sizeof(hints));
@@ -226,6 +226,10 @@ bool client_connect(Client* client, char* host, int port) {
return true; return true;
} }
bool client_connect(Client* client, char* host, int port) {
return tcp_connect_socket(client, host, port);
}
void client_disconnect(Client* client) { void client_disconnect(Client* client) {
if (client->ssl) { if (client->ssl) {
SSL_shutdown(client->ssl); SSL_shutdown(client->ssl);
+1
View File
@@ -29,6 +29,7 @@ void server_accept_loop(Server* server, void (*child_fn)(int, void*), void* chil
const char* log_fmt); const char* log_fmt);
void server_delete(Server** server); void server_delete(Server** server);
Client* client_create(); Client* client_create();
bool tcp_connect_socket(Client* client, const char* host, int port);
bool client_connect(Client* client, char* host, int port); bool client_connect(Client* client, char* host, int port);
void client_disconnect(Client* client); void client_disconnect(Client* client);
void client_delete(Client* client); void client_delete(Client* client);
+1 -48
View File
@@ -2,8 +2,6 @@
#include "log.h" #include "log.h"
#include "protocol.h" #include "protocol.h"
#include "transport_tcp.h" #include "transport_tcp.h"
#include <arpa/inet.h>
#include <netdb.h>
#include <openssl/err.h> #include <openssl/err.h>
#include <openssl/ssl.h> #include <openssl/ssl.h>
#include <signal.h> #include <signal.h>
@@ -149,53 +147,8 @@ bool server_listen_tls(Server* server, void (*handler)(int file_descriptor)) {
bool client_connect_tls(Client* client, char* host, int port, const char* cert_path, bool client_connect_tls(Client* client, char* host, int port, const char* cert_path,
const char* key_path, const char* ca_path) { const char* key_path, const char* ca_path) {
struct addrinfo hints; if (!tcp_connect_socket(client, host, port))
struct addrinfo* result;
memset(&hints, 0, sizeof(hints));
hints.ai_family = AF_UNSPEC;
hints.ai_socktype = SOCK_STREAM;
hints.ai_protocol = IPPROTO_TCP;
char port_str[16];
snprintf(port_str, sizeof(port_str), "%d", port);
int err = getaddrinfo(host, port_str, &hints, &result);
if (err != 0 || result == NULL) {
fprintf(stderr, "Could not resolve host: %s (%s)\n", host, gai_strerror(err));
return false; return false;
}
struct addrinfo* rp;
bool connected = false;
for (rp = result; rp != NULL; rp = rp->ai_next) {
if (client->file_descriptor >= 0)
close(client->file_descriptor);
client->file_descriptor = socket(rp->ai_family, rp->ai_socktype, rp->ai_protocol);
if (client->file_descriptor < 0)
continue;
struct timeval ct;
ct.tv_sec = tcp_get_contimeout_sec();
ct.tv_usec = 0;
setsockopt(client->file_descriptor, SOL_SOCKET, SO_RCVTIMEO, &ct, sizeof(ct));
setsockopt(client->file_descriptor, SOL_SOCKET, SO_SNDTIMEO, &ct, sizeof(ct));
memcpy(&client->address, rp->ai_addr, rp->ai_addrlen);
client->address_length = rp->ai_addrlen;
if (connect(client->file_descriptor, (struct sockaddr*)&client->address,
client->address_length) == 0) {
connected = true;
break;
}
}
freeaddrinfo(result);
if (!connected) {
perror("Could not connect to Server!");
return false;
}
SSL_CTX* ctx = create_ssl_ctx(false, cert_path, key_path, ca_path); SSL_CTX* ctx = create_ssl_ctx(false, cert_path, key_path, ca_path);
if (!ctx) if (!ctx)