Author SHA1 Message Date
TapTap 7b4164d3ba fix: accept numeric chmod special bits
CI / lint (pull_request) Successful in 12s
CI / sanitizers (address) (pull_request) Successful in 38s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 14s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-09-03 22:14:47 +02:00
TapTap e093a3c5b6 feat: add rsync-compatible chmod option
CI / lint (pull_request) Successful in 11s
CI / sanitizers (address) (pull_request) Successful in 36s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-09-03 17:49:52 +02:00
TapTap 190fc5d300 Merge pull request 'fix: harden network security boundaries' (#213) from security-fixes into dev
CI / lint (push) Successful in 11s
CI / sanitizers (undefined) (push) Successful in 37s
CI / sanitizers (address) (push) Successful in 37s
CI / fuzz-build (push) Successful in 14s
CI / coverage (push) Successful in 32s
CI / build-and-test (push) Successful in 1m15s
CI / valgrind (push) Successful in 33s
2026-09-01 20:52:55 +02:00
TapTap 8ae5f82013 style: clang-format file_receive.c
CI / lint (pull_request) Successful in 11s
CI / sanitizers (address) (pull_request) Successful in 35s
CI / sanitizers (undefined) (pull_request) Successful in 36s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 31s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 32s
2026-09-01 20:49:39 +02:00
TapTap 0c4855e344 merge: resolve dev (quality refactor) into security-fixes
- file.c split layout retained; security-hardened secure-fs helpers
  (open_secure_parent/to_disk_secure/rename_secure/stat_secure) now live in
  file.c with file_ prefix and are shared with file_receive.c
- file_receive.c takes the security branch's bounded allocations
  (receive_data_limited, data_decompress_limited, size checks) and
  STATUS_ERROR signaling
- file_send.c gains the data consistency check on file->data
- client_validation.c: stricter --tls requiring --ca, log_message style
- utils.c: hardened openat/mkdirat mkdir_r from security branch
2026-09-01 20:46:07 +02:00
TapTap 265f0c0c09 Merge pull request 'refactor: quality cleanup — table-driven CLI, file.c split, scanner helpers, unified error reporting' (#216) from refactor/quality-cleanup into dev
CI / lint (push) Successful in 12s
CI / sanitizers (address) (push) Successful in 35s
CI / sanitizers (undefined) (push) Successful in 35s
CI / fuzz-build (push) Successful in 14s
CI / coverage (push) Successful in 32s
CI / build-and-test (push) Successful in 1m15s
CI / valgrind (push) Successful in 32s
2026-09-01 20:39:31 +02:00
TapTap 047a9a1906 style: apply clang-format 18
CI / lint (pull_request) Successful in 23s
CI / sanitizers (address) (pull_request) Successful in 54s
CI / sanitizers (undefined) (pull_request) Successful in 54s
CI / fuzz-build (pull_request) Successful in 14s
CI / build-and-test (pull_request) Successful in 1m17s
CI / coverage (pull_request) Successful in 31s
CI / valgrind (pull_request) Successful in 33s
2026-08-30 14:20:33 +02:00
TapTap 6d82967c68 refactor: standardize error reporting on the log module
- Add log_perror() helper (context + strerror(errno)) to the log module
- Replace all bare perror() calls with log_perror() so errors are routed
  through the unified logger (stderr sink + optional --log-file sink)
- Convert fprintf(stderr, "Error:/Warning: ...") in client code to
  log_message(); raw fprintf kept only for progress/stats output
2026-08-30 14:06:48 +02:00
TapTap ae95211a05 refactor: unify send_files cleanup and share progress printing
- Extract print_transfer_progress() shared by single-threaded loop and
  the multithreaded progress thread
- Route all send_files() exits through a single send_fail cleanup path
- Fix pre-existing manifest leak on success without --delete
2026-08-30 14:02:12 +02:00
TapTap d639dfdc07 refactor: split file.c into file_send.c and file_receive.c
- file.c: File/FileMetadata lifecycle and local disk helpers (~110 lines)
- file_send.c: client-side send path (file_send_single_calls, file_send_sendfile)
- file_receive.c: server-side receive/save path (file_receive, receive_incremental_check,
  receive_manifest, file_save_to_disk)
- file_types.h holds shared struct definitions; file.h remains an umbrella header
  so existing includes are unaffected
Completes the transfer/protocol separation started in PR #212
2026-08-30 13:59:43 +02:00
TapTap 3fa3e150ce refactor: split parallel_scanner_create_with_options into focused helpers
- parallel_scanner_init(): result queue + sync primitive setup with unwinding
- batch_files(): root-file chunk batching, reusable by other scan paths
- scan_root_directory()/scan_root_entry(): root-dir scanning
- spawn_parallel_workers(): worker thread creation with per-thread arg setup
Main function reduced from ~230 to ~40 lines
2026-08-30 13:56:27 +02:00
TapTap e3e766ba3d refactor: table-driven CLI option parsing in client_cli
- Add OPTION_TABLE for options that map directly to Config fields
  (flag/string/pos-int/nonneg-int/ull kinds)
- Extract parse_ull_arg() replacing 5 duplicated strtoull blocks
- Extract config_add_pattern() replacing duplicated --exclude/--include
  append logic, also reused by read_patterns_from_file()
- parse_args() reduced from ~275 to ~160 lines
2026-08-30 13:53:59 +02:00
TapTap f2917eb163 refactor: remove dead API, rename to_disk/config_is_remote_dest, fix perror newlines
- Delete unused public array_list_extend (made static)
- Delete legacy 16-parameter parallel_scanner_create wrapper; migrate test to parallel_scanner_create_with_options
- Rename to_disk -> file_write_to_disk and is_remote_dest -> config_is_remote_dest for module_action naming convention
- Remove stray newlines in perror calls (perror already appends one)
2026-08-30 13:50:42 +02:00
TapTap 6c4d0246bc merge: resolve origin dev refactor conflicts
CI / lint (pull_request) Successful in 28s
CI / sanitizers (address) (pull_request) Successful in 36s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 14s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-08-16 09:56:15 +02:00
TapTap 1391564ce7 test: harden allocation checks
CI / lint (pull_request) Successful in 31s
CI / sanitizers (address) (pull_request) Successful in 37s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 14s
CI / coverage (pull_request) Successful in 31s
CI / build-and-test (pull_request) Successful in 1m14s
CI / valgrind (pull_request) Successful in 33s
2026-08-16 09:37:28 +02:00
TapTap d102622e28 fix: close remaining PR review gaps
CI / lint (pull_request) Failing after 31s
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-16 09:35:26 +02:00
TapTap 1abc17142f Merge pull request 'refactor: split transfer and protocol responsibilities' (#212) from refactor/codebase-structure into dev
CI / lint (push) Successful in 30s
CI / sanitizers (undefined) (push) Successful in 38s
CI / sanitizers (address) (push) Successful in 39s
CI / fuzz-build (push) Successful in 14s
CI / coverage (push) Successful in 31s
CI / build-and-test (push) Successful in 1m15s
CI / valgrind (push) Successful in 33s
Reviewed-on: #212
2026-08-16 09:11:58 +02:00
TapTap 059dc2ac75 fix: close remaining PR review gaps
CI / lint (pull_request) Successful in 32s
CI / sanitizers (undefined) (pull_request) Successful in 36s
CI / sanitizers (address) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 13s
CI / coverage (pull_request) Successful in 31s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 20:38:04 +02:00
TapTap 7cffff0b8b fix: fail closed on authorization setup failure
CI / lint (pull_request) Successful in 30s
CI / sanitizers (address) (pull_request) Successful in 36s
CI / sanitizers (undefined) (pull_request) Successful in 35s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 31s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 20:02:09 +02:00
TapTap d3b4c62f51 Merge branch 'dev' into refactor/codebase-structure
CI / lint (pull_request) Successful in 30s
CI / sanitizers (address) (pull_request) Successful in 36s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 31s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 32s
2026-08-15 15:09:00 +02:00
TapTap 2f804ecfdd fix: address remaining PR 212 review issues
CI / lint (pull_request) Successful in 31s
CI / sanitizers (address) (pull_request) Successful in 37s
CI / sanitizers (undefined) (pull_request) Successful in 36s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 13:52:38 +02:00
TapTap 0c24136052 Merge pull request 'docs: clarify rsync compatibility roadmap' (#215) from docs/rsync-compatible-readme-clean into dev
CI / lint (push) Successful in 34s
CI / sanitizers (address) (push) Successful in 35s
CI / sanitizers (undefined) (push) Successful in 35s
CI / fuzz-build (push) Successful in 15s
CI / build-and-test (push) Successful in 1m15s
CI / coverage (push) Successful in 31s
CI / valgrind (push) Successful in 33s
Reviewed-on: #215
2026-08-15 13:48:01 +02:00
TapTap d19fc68803 Merge branch 'dev' into docs/rsync-compatible-readme-clean
CI / lint (pull_request) Successful in 33s
CI / sanitizers (address) (pull_request) Successful in 36s
CI / sanitizers (undefined) (pull_request) Successful in 36s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m16s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 13:47:55 +02:00
TapTap 78876b110e fix: harden remaining review findings
CI / lint (pull_request) Successful in 34s
CI / sanitizers (address) (pull_request) Successful in 37s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 33s
CI / build-and-test (pull_request) Successful in 1m16s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 13:44:31 +02:00
TapTap b3a724afc8 fix: preserve scanner file ownership on failure
CI / lint (pull_request) Successful in 32s
CI / sanitizers (address) (pull_request) Successful in 37s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 14s
CI / coverage (pull_request) Successful in 31s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 34s
2026-08-15 13:40:31 +02:00
TapTap daf662cc2d fix: address remaining PR review findings
CI / lint (pull_request) Successful in 33s
CI / sanitizers (address) (pull_request) Successful in 36s
CI / sanitizers (undefined) (pull_request) Successful in 36s
CI / fuzz-build (pull_request) Successful in 14s
CI / coverage (pull_request) Successful in 31s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 13:24:13 +02:00
TapTap 98f833980d fix: handle pipeline cancellation failures
CI / lint (pull_request) Successful in 30s
CI / sanitizers (address) (pull_request) Successful in 37s
CI / sanitizers (undefined) (pull_request) Successful in 38s
CI / fuzz-build (pull_request) Successful in 16s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m16s
CI / valgrind (pull_request) Successful in 32s
2026-08-15 13:20:50 +02:00
TapTap 298bfe0d0f fix: address delegated review findings
CI / lint (pull_request) Successful in 33s
CI / sanitizers (address) (pull_request) Successful in 35s
CI / sanitizers (undefined) (pull_request) Successful in 36s
CI / fuzz-build (pull_request) Successful in 14s
CI / build-and-test (pull_request) Successful in 1m14s
CI / coverage (pull_request) Successful in 30s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 13:13:57 +02:00
TapTap 604a14f0be fix: close remaining PR review gaps
CI / lint (pull_request) Successful in 30s
CI / sanitizers (address) (pull_request) Successful in 36s
CI / sanitizers (undefined) (pull_request) Successful in 36s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 33s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 13:13:50 +02:00
TapTap 6f8847899c fix: validate remaining wire booleans
CI / lint (pull_request) Successful in 34s
CI / sanitizers (address) (pull_request) Successful in 35s
CI / fuzz-build (pull_request) Successful in 14s
CI / sanitizers (undefined) (pull_request) Successful in 35s
CI / build-and-test (pull_request) Successful in 1m14s
CI / coverage (pull_request) Successful in 31s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 12:49:30 +02:00
TapTap 09454413d4 fix: address PR 212 review issues
CI / lint (pull_request) Successful in 30s
CI / sanitizers (address) (pull_request) Successful in 38s
CI / sanitizers (undefined) (pull_request) Successful in 38s
CI / fuzz-build (pull_request) Successful in 15s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 12:48:31 +02:00
TapTap ff189aeef2 fix: harden network security boundaries
CI / lint (pull_request) Successful in 33s
CI / sanitizers (address) (pull_request) Successful in 37s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 13s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m16s
CI / valgrind (pull_request) Successful in 32s
2026-08-15 12:34:41 +02:00
TapTap 298976aefb fix: link receiver into fuzz targets
CI / lint (pull_request) Successful in 29s
CI / sanitizers (address) (pull_request) Successful in 37s
CI / sanitizers (undefined) (pull_request) Successful in 36s
CI / fuzz-build (pull_request) Successful in 16s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m16s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 12:34:19 +02:00
TapTap 3b7f97853c fix: satisfy cppcheck const analysis
CI / lint (pull_request) Successful in 31s
CI / sanitizers (address) (pull_request) Successful in 38s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Failing after 14s
CI / coverage (pull_request) Successful in 32s
CI / build-and-test (pull_request) Successful in 1m15s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 12:30:52 +02:00
TapTap 8cca891b9a refactor: split transfer and protocol responsibilities
CI / lint (pull_request) Failing after 30s
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:27:18 +02:00
TapTap 41b624db5a Merge pull request 'fix: satisfy clang-format in CLI test' (#211) from fix/dev-ci-clang-format into dev
CI / lint (push) Successful in 34s
CI / sanitizers (address) (push) Successful in 37s
CI / sanitizers (undefined) (push) Successful in 37s
CI / fuzz-build (push) Successful in 14s
CI / coverage (push) Successful in 31s
CI / build-and-test (push) Successful in 1m16s
CI / valgrind (push) Successful in 34s
Reviewed-on: #211
2026-08-15 11:55:41 +02:00
TapTap 786e686563 fix: satisfy clang-format in CLI test
CI / lint (pull_request) Successful in 33s
CI / sanitizers (address) (pull_request) Successful in 37s
CI / sanitizers (undefined) (pull_request) Successful in 37s
CI / fuzz-build (pull_request) Successful in 14s
CI / coverage (pull_request) Successful in 31s
CI / build-and-test (pull_request) Successful in 1m16s
CI / valgrind (pull_request) Successful in 33s
2026-08-15 11:52:05 +02:00
70 changed files with 3882 additions and 2578 deletions
+7 -4
View File
@@ -58,15 +58,18 @@ endif()
find_package(OpenSSL REQUIRED)
file(GLOB SHARED_SRCS "src/shared/*.c")
set(FILE_STORE_SRCS "${CMAKE_CURRENT_SOURCE_DIR}/src/shared/file_store.c")
list(REMOVE_ITEM SHARED_SRCS ${FILE_STORE_SRCS})
file(GLOB SERVER_SRCS "src/server/*.c")
set(SERVER_RECEIVER_SRCS src/server/receiver.c)
file(GLOB CLIENT_SRCS "src/client/*.c")
# --- Main executables ---
add_executable(server ${SERVER_SRCS} ${SHARED_SRCS})
add_executable(server ${SERVER_SRCS} ${SHARED_SRCS} ${FILE_STORE_SRCS})
target_include_directories(server PRIVATE src/shared src/server src/client)
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} ${FILE_STORE_SRCS} ${SERVER_RECEIVER_SRCS})
target_include_directories(client PRIVATE src/shared src/server src/client)
target_link_libraries(client PRIVATE Threads::Threads ${ZSTD_LIBRARY} OpenSSL::SSL OpenSSL::Crypto xxhash)
@@ -79,7 +82,7 @@ set(TEST_INCLUDES tests src/shared src/server src/client)
# Monolithic test binary (backward compatible)
file(GLOB TEST_SRCS "tests/test_*.c" "tests/runner.c")
add_executable(tests ${TEST_SRCS} ${SHARED_SRCS} src/client/scanner.c src/client/client_cli.c)
add_executable(tests ${TEST_SRCS} ${SHARED_SRCS} ${FILE_STORE_SRCS} ${SERVER_RECEIVER_SRCS} src/client/scanner.c src/client/client_cli.c src/client/client_validation.c src/client/usage.c)
target_include_directories(tests PRIVATE ${TEST_INCLUDES})
target_compile_definitions(tests PRIVATE FASTSYNC_TEST_BUILD)
target_link_libraries(tests PRIVATE ${TEST_LIBS})
@@ -94,7 +97,7 @@ if(ENABLE_FUZZ)
file(GLOB FUZZ_SRCS "tests/fuzz/*.c")
foreach(FUZZ_SRC ${FUZZ_SRCS})
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} ${FILE_STORE_SRCS} ${SERVER_RECEIVER_SRCS})
target_include_directories(${FUZZ_NAME} PRIVATE ${TEST_INCLUDES})
target_compile_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined -fno-omit-frame-pointer)
target_link_options(${FUZZ_NAME} PRIVATE -fsanitize=fuzzer,address,undefined)
+188 -25
View File
@@ -1,4 +1,4 @@
# FastSync
#FastSync
FastSync is a high-performance file synchronization tool designed to become a
drop-in replacement for common `rsync` workflows. It keeps the familiar
@@ -79,10 +79,166 @@ partial, alternate, and planned behavior.
### Build
Requirements: C11 compiler, CMake 3.22 or newer, xxHash, zstd, OpenSSL,
pthreads, and an SSH client for SSH transport. The first CMake configure fetches
xxHash from GitHub, so network access is required unless the dependency is
already cached.
### Client
| Argument | Description |
|----------|-------------|
| 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) |
| `--client-cn <name>` | Required TLS client certificate common name |
### 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) |
| `--destination-root <path>` | Authorized destination root (default: `.`) |
| `--allow-delete` | Permit manifest deletion |
| `--allow-unauthenticated` | Permit plaintext TCP clients |
| `-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,
mutual 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
TLS requires `--ca` and performs mutual TLS verification (`SSL_VERIFY_PEER` with depth 4). Connections without certificate verification are rejected.
### 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
cmake -B build -S .
@@ -105,6 +261,7 @@ working directory, so use a destination below that directory unless the
remote server is otherwise configured with a matching authorized root.
```bash
ssh user@host 'mkdir -p destination'
./build/client /path/to/source user@host:destination
```
@@ -124,8 +281,10 @@ Then run the client:
--save-to-disk
```
### TLS transfer
Plain TCP requires the explicit `--allow-unauthenticated` server option. Use TLS for
authenticated network connections.
### TLS transfer
```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 \
@@ -140,32 +299,32 @@ These examples show the intended rsync-style workflow. Options marked as
FastSync-native are optional performance or transport extensions.
```bash
# Basic synchronization
#Basic synchronization
./build/client /source/ /destination/
# Archive-style synchronization (current FastSync archive behavior)
#Archive - style synchronization(current FastSync archive behavior)
./build/client -a /source/ user@host:destination/
# Preview a transfer without changing the destination
#Preview a transfer without changing the destination
./build/client -n /source/ /destination/
# Exclude temporary and object files
#Exclude temporary and object files
./build/client --exclude '*.tmp' --exclude '*.o' \
/source/ user@host:destination/
# Remove destination entries not present in the source
#Remove destination entries not present in the source
./build/client --delete /source/ user@host:destination/
# Skip unchanged files using size and modification time
#Skip unchanged files using size and modification time
./build/client --incremental /source/ user@host:destination/
# Verify content when size and time are not sufficient
#Verify content when size and time are not sufficient
./build/client --incremental --checksum /source/ user@host:destination/
# Preserve supported mode and timestamp metadata
#Preserve supported mode and timestamp metadata
./build/client -M /source/ user@host:destination/
# Keep backups of overwritten destination files
#Keep backups of overwritten destination files
./build/client --backup --backup-dir backups \
/source/ user@host:destination/
```
@@ -215,14 +374,16 @@ before FastSync can claim full rsync CLI compatibility.
| `--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`. |
| `--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
@@ -230,7 +391,8 @@ before FastSync can claim full rsync CLI compatibility.
| Option | Description |
|---|---|
| `-M`, `--preserve` | Preserve supported file metadata, currently mode and modification time. |
| `-l`, `--links` | Request symlink preservation; link-target transfer remains incomplete. |
| `-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. |
@@ -273,7 +435,8 @@ before FastSync can claim full rsync CLI compatibility.
| `--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. |
| `--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. |
+1 -1
View File
@@ -120,7 +120,7 @@ This document maps rsync's full feature set to FastSync's current implementation
| `-g`, `--group` | Preserve group | ✅ Implemented | Part of -M |
| `-t`, `--times` | Preserve modification times | ✅ Implemented | Part of -M |
| `-E`, `--executability` | Preserve executability | ❌ Not Implemented | |
| `--chmod=CHMOD` | Affect file permissions | ❌ Not Implemented | |
| `--chmod=CHMOD` | Affect file permissions | ✅ Implemented | Supports numeric and symbolic `ugo` `rwx` changes; retains receiver safety masking |
| `-A`, `--acls` | Preserve ACLs | ❌ Not Implemented | Removed because it had no effect |
| `-X`, `--xattrs` | Preserve extended attributes | ❌ Not Implemented | Removed because it had no effect |
| `-H`, `--hard-links` | Preserve hard links | ❌ Not Implemented | Removed because it had no effect |
+221 -246
View File
@@ -1,18 +1,21 @@
#include "client_send.h"
#include "client_validation.h"
#include "chmod.h"
#include "config.h"
#include "delta.h"
#include "log.h"
#include "protocol.h"
#include "transport_tcp.h"
#include "transport_tls.h"
#include "usage.h"
#include "utils.h"
#include <errno.h>
#include <limits.h>
#include <stdbool.h>
#include <stddef.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include "usage.c"
#ifndef FASTSYNC_TEST_BUILD
/* Parse environment variables for source/destination directories and save-to-disk flag. */
@@ -58,7 +61,7 @@ static bool parse_positive_int(const char* s, int* out_val) {
static int set_string_option(char** dest, const char* value, const char* option_name) {
char* dup = str_dup(value);
if (!dup) {
fprintf(stderr, "Error: memory allocation failed for %s\n", option_name);
log_message(LOG_LEVEL_ERROR, "memory allocation failed for %s", option_name);
return -1;
}
free(*dest);
@@ -69,7 +72,7 @@ static int set_string_option(char** dest, const char* value, const char* option_
/* Parse a string as a positive integer into *dest. Returns true on success, false on error. */
static int set_positive_int_option(int* dest, const char* value, const char* option_name) {
if (!parse_positive_int(value, dest)) {
fprintf(stderr, "Error: %s must be a positive integer\n", option_name);
log_message(LOG_LEVEL_ERROR, "%s must be a positive integer", option_name);
return -1;
}
return 0;
@@ -78,7 +81,7 @@ static int set_positive_int_option(int* dest, const char* value, const char* opt
/* Parse a string as a non-negative integer into *dest. Returns true on success, false on error. */
static int set_nonneg_int_option(int* dest, const char* value, const char* option_name) {
if (!parse_nonneg_int(value, dest)) {
fprintf(stderr, "Error: %s must be a non-negative integer\n", option_name);
log_message(LOG_LEVEL_ERROR, "%s must be a non-negative integer", option_name);
return -1;
}
return 0;
@@ -86,105 +89,208 @@ static int set_nonneg_int_option(int* dest, const char* value, const char* optio
static int read_patterns_from_file(const char* filepath, char*** patterns, int* count);
/* Parse a string as an unsigned long long. Returns 0 on success, -1 on error. */
static int parse_ull_arg(const char* val, unsigned long long* out, const char* optname) {
char* end;
errno = 0;
unsigned long long v = strtoull(val, &end, 10);
if (errno != 0 || *end != '\0') {
log_message(LOG_LEVEL_ERROR, "%s must be a non-negative integer", optname);
return -1;
}
*out = v;
return 0;
}
/* Append a duplicated pattern to a growable pattern array. Returns 0 on success, -1 on error. */
static int config_add_pattern(char*** patterns, int* count, const char* value,
const char* optname) {
char** tmp = realloc(*patterns, (*count + 1) * sizeof(char*));
if (!tmp) {
log_message(LOG_LEVEL_ERROR, "memory allocation failed for %s", optname);
return -1;
}
*patterns = tmp;
char* dup = str_dup(value);
if (!dup) {
log_message(LOG_LEVEL_ERROR, "memory allocation failed for %s", optname);
return -1;
}
(*patterns)[(*count)++] = dup;
return 0;
}
typedef enum {
OPT_FLAG,
OPT_STRING,
OPT_POS_INT,
OPT_NONNEG_INT,
OPT_ULL,
} OptKind;
typedef struct {
const char* name;
const char* alias;
OptKind kind;
size_t offset; /* offsetof of the target field in Config */
} OptionEntry;
/* Options that map directly onto a Config field with no side effects. */
static const OptionEntry OPTION_TABLE[] = {
{"--dry-run", "-n", OPT_FLAG, offsetof(Config, dry_run)},
{"--delete", NULL, OPT_FLAG, offsetof(Config, use_delete)},
{"--incremental", NULL, OPT_FLAG, offsetof(Config, use_incremental)},
{"--delta", NULL, OPT_FLAG, offsetof(Config, use_delta)},
{"--save-to-disk", NULL, OPT_FLAG, offsetof(Config, save_to_disk)},
{"--progress", NULL, OPT_FLAG, offsetof(Config, show_progress)},
{"--tls", NULL, OPT_FLAG, offsetof(Config, use_tls)},
{"--backup", NULL, OPT_FLAG, offsetof(Config, backup)},
{"--stats", NULL, OPT_FLAG, offsetof(Config, stats)},
{"--partial", NULL, OPT_FLAG, offsetof(Config, partial)},
{"--links", "-l", OPT_FLAG, offsetof(Config, follow_symlinks)},
{"--copy-links", NULL, OPT_FLAG, offsetof(Config, copy_links)},
{"--safe-links", NULL, OPT_FLAG, offsetof(Config, safe_links)},
{"--copy-unsafe-links", NULL, OPT_FLAG, offsetof(Config, copy_unsafe_links)},
{"--sparse", "-S", OPT_FLAG, offsetof(Config, preserve_sparse)},
{"--inplace", NULL, OPT_FLAG, offsetof(Config, inplace)},
{"--checksum", NULL, OPT_FLAG, offsetof(Config, checksum)},
{"--chmod", NULL, OPT_STRING, offsetof(Config, chmod_spec)},
{"--source-dir", NULL, OPT_STRING, offsetof(Config, send_directory)},
{"--dest-dir", NULL, OPT_STRING, offsetof(Config, receive_root_directory)},
{"--server-host", NULL, OPT_STRING, offsetof(Config, server_host)},
{"--cert", NULL, OPT_STRING, offsetof(Config, tls_cert)},
{"--key", NULL, OPT_STRING, offsetof(Config, tls_key)},
{"--ca", NULL, OPT_STRING, offsetof(Config, tls_ca)},
{"--backup-dir", NULL, OPT_STRING, offsetof(Config, backup_dir)},
{"--fastsync-server-path", NULL, OPT_STRING, offsetof(Config, fastsync_server_path)},
{"--partial-dir", NULL, OPT_STRING, offsetof(Config, partial_dir)},
{"--suffix", NULL, OPT_STRING, offsetof(Config, suffix)},
{"--timeout", NULL, OPT_POS_INT, offsetof(Config, timeout)},
{"--contimeout", NULL, OPT_POS_INT, offsetof(Config, contimeout)},
{"--max-depth", NULL, OPT_NONNEG_INT, offsetof(Config, max_depth)},
{"--max-size", NULL, OPT_ULL, offsetof(Config, max_size)},
{"--min-size", NULL, OPT_ULL, offsetof(Config, min_size)},
};
static bool opt_is(const char* arg, const char* name, const char* alias) {
return strcmp(arg, name) == 0 || (alias && strcmp(arg, alias) == 0);
}
static const OptionEntry* find_table_option(const char* arg) {
for (size_t i = 0; i < sizeof(OPTION_TABLE) / sizeof(OPTION_TABLE[0]); i++)
if (opt_is(arg, OPTION_TABLE[i].name, OPTION_TABLE[i].alias))
return &OPTION_TABLE[i];
return NULL;
}
static int apply_table_option(Config* config, const OptionEntry* entry, const char* value) {
void* field = (char*)config + entry->offset;
switch (entry->kind) {
case OPT_FLAG:
*(bool*)field = true;
return 0;
case OPT_STRING:
return set_string_option((char**)field, value, entry->name);
case OPT_POS_INT:
return set_positive_int_option((int*)field, value, entry->name);
case OPT_NONNEG_INT:
return set_nonneg_int_option((int*)field, value, entry->name);
case OPT_ULL: {
unsigned long long v;
if (parse_ull_arg(value, &v, entry->name) != 0)
return -1;
*(unsigned long long*)field = v;
return 0;
}
}
return -1;
}
/* Parse CLI arguments into config. Returns 0 on success, -1 on error, 1 for help/clean-exit. */
int parse_args(Config* config, int argc, char* argv[], int* positional_args,
int* positional_count) {
for (int i = 1; i < argc; i++) {
if (strcmp(argv[i], "--help") == 0) {
const OptionEntry* entry = find_table_option(argv[i]);
if (entry) {
if (entry->kind != OPT_FLAG) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", entry->name);
return -1;
}
if (apply_table_option(config, entry, argv[++i]) != 0)
return -1;
if (entry->offset == offsetof(Config, chmod_spec)) {
mode_t ignored;
if (!chmod_apply(0, config->chmod_spec, &ignored)) {
log_message(LOG_LEVEL_ERROR, "--chmod has invalid permission changes");
return -1;
}
config->use_metadata = true;
}
} else if (apply_table_option(config, entry, NULL) != 0) {
return -1;
}
continue;
}
if (strncmp(argv[i], "--chmod=", 8) == 0) {
if (set_string_option(&config->chmod_spec, argv[i] + 8, "--chmod") != 0)
return -1;
mode_t ignored;
if (!chmod_apply(0, config->chmod_spec, &ignored)) {
log_message(LOG_LEVEL_ERROR, "--chmod has invalid permission changes");
return -1;
}
config->use_metadata = true;
continue;
}
if (opt_is(argv[i], "--help", NULL)) {
print_usage();
return 1;
} else if (strcmp(argv[i], "-V") == 0 || strcmp(argv[i], "--version") == 0) {
} else if (opt_is(argv[i], "-V", "--version")) {
printf("fastsync version %s\n", PROTOCOL_VERSION);
return 1;
} else if (strcmp(argv[i], "-a") == 0 || strcmp(argv[i], "--archive") == 0) {
} else if (opt_is(argv[i], "-a", "--archive")) {
config->use_compression = true;
config->use_multithreading = true;
config->use_metadata = true;
log_message(LOG_LEVEL_INFO, "Enabled archive mode (-c -m -M)");
} else if (strcmp(argv[i], "-n") == 0 || strcmp(argv[i], "--dry-run") == 0) {
config->dry_run = true;
} else if (strcmp(argv[i], "-p") == 0 && i + 1 < argc) {
} else if (opt_is(argv[i], "-p", NULL) && i + 1 < argc) {
if (set_positive_int_option(&config->ssh_port, argv[++i], "-p") != 0)
return -1;
if (config->ssh_port > 65535) {
log_message(LOG_LEVEL_ERROR, "SSH port must be 1-65535\n");
log_message(LOG_LEVEL_ERROR, "SSH port must be 1-65535");
return -1;
}
} else if (strcmp(argv[i], "--delete") == 0) {
config->use_delete = true;
} else if (strcmp(argv[i], "--exclude") == 0 && i + 1 < argc) {
char** tmp = realloc(config->exclude_patterns, (config->exclude_count + 1) * sizeof(char*));
if (!tmp) {
fprintf(stderr, "Error: memory allocation failed for --exclude\n");
} else if (opt_is(argv[i], "--exclude", NULL) && i + 1 < argc) {
if (config_add_pattern(&config->exclude_patterns, &config->exclude_count, argv[++i],
"--exclude") != 0)
return -1;
}
config->exclude_patterns = tmp;
char* dup = str_dup(argv[++i]);
if (!dup) {
fprintf(stderr, "Error: memory allocation failed for --exclude\n");
} else if (opt_is(argv[i], "--include", NULL) && i + 1 < argc) {
if (config_add_pattern(&config->include_patterns, &config->include_count, argv[++i],
"--include") != 0)
return -1;
}
config->exclude_patterns[config->exclude_count++] = dup;
} else if (strcmp(argv[i], "--include") == 0 && i + 1 < argc) {
char** tmp = realloc(config->include_patterns, (config->include_count + 1) * sizeof(char*));
if (!tmp) {
fprintf(stderr, "Error: memory allocation failed for --include\n");
} else if (opt_is(argv[i], "--delta-block", NULL) && i + 1 < argc) {
unsigned long long val;
if (parse_ull_arg(argv[++i], &val, "--delta-block") != 0)
return -1;
}
config->include_patterns = tmp;
char* dup = str_dup(argv[++i]);
if (!dup) {
fprintf(stderr, "Error: memory allocation failed for --include\n");
return -1;
}
config->include_patterns[config->include_count++] = dup;
} else if (strcmp(argv[i], "--max-size") == 0 && i + 1 < argc) {
char* end;
errno = 0;
unsigned long long val = strtoull(argv[++i], &end, 10);
if (errno != 0 || *end != '\0') {
fprintf(stderr, "Error: --max-size must be a non-negative integer\n");
return -1;
}
config->max_size = val;
} else if (strcmp(argv[i], "--min-size") == 0 && i + 1 < argc) {
char* end;
errno = 0;
unsigned long long val = strtoull(argv[++i], &end, 10);
if (errno != 0 || *end != '\0') {
fprintf(stderr, "Error: --min-size must be a non-negative integer\n");
return -1;
}
config->min_size = val;
} else if (strcmp(argv[i], "--incremental") == 0) {
config->use_incremental = true;
} else if (strcmp(argv[i], "--delta") == 0) {
config->use_delta = true;
} else if (strcmp(argv[i], "--delta-block") == 0 && i + 1 < argc) {
char* end;
errno = 0;
unsigned long long val = strtoull(argv[++i], &end, 10);
if (errno != 0 || *end != '\0') {
fprintf(stderr, "Error: --delta-block must be a positive integer\n");
return -1;
}
if (val >= DELTA_BLOCK_SIZE_MIN && val <= DELTA_BLOCK_SIZE_MAX)
config->delta_block_size = (uint32_t)val;
else
fprintf(stderr, "Warning: --delta-block value %llu out of range, using default\n", val);
} else if (strcmp(argv[i], "--delta-max") == 0 && i + 1 < argc) {
char* end;
errno = 0;
unsigned long long val = strtoull(argv[++i], &end, 10);
if (errno != 0 || *end != '\0') {
fprintf(stderr, "Error: --delta-max must be a positive integer\n");
log_message(LOG_LEVEL_WARNING, "--delta-block value %llu out of range, using default", val);
} else if (opt_is(argv[i], "--delta-max", NULL) && i + 1 < argc) {
unsigned long long val;
if (parse_ull_arg(argv[++i], &val, "--delta-max") != 0)
return -1;
}
if (val >= DELTA_MIN_FILE_SIZE)
config->delta_max_file_size = val;
else
fprintf(stderr, "Warning: --delta-max value %llu too small, using default\n", val);
} else if (strcmp(argv[i], "-c") == 0 || strcmp(argv[i], "-z") == 0) {
log_message(LOG_LEVEL_WARNING, "--delta-max value %llu too small, using default", val);
} else if (opt_is(argv[i], "-c", "-z")) {
config->use_compression = true;
log_message(LOG_LEVEL_INFO, "Enabled Compression");
if (i + 1 < argc) {
@@ -192,7 +298,7 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
long level = strtol(argv[i + 1], &end_ptr, 10);
if (*end_ptr == '\0') {
if (level < 1 || level > 22) {
fprintf(stderr, "Error: compression level must be 1-22\n");
log_message(LOG_LEVEL_ERROR, "compression level must be 1-22");
return -1;
}
config->compression_level = (int)level;
@@ -200,91 +306,51 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
i++;
}
}
} else if (strcmp(argv[i], "--source-dir") == 0 && i + 1 < argc) {
if (set_string_option(&config->send_directory, argv[++i], "--source-dir") != 0)
return -1;
} else if (strcmp(argv[i], "--dest-dir") == 0 && i + 1 < argc) {
if (set_string_option(&config->receive_root_directory, argv[++i], "--dest-dir") != 0)
return -1;
} else if (strcmp(argv[i], "--save-to-disk") == 0) {
config->save_to_disk = true;
} else if (strcmp(argv[i], "-M") == 0 || strcmp(argv[i], "--preserve") == 0) {
} else if (opt_is(argv[i], "-M", "--preserve")) {
config->use_metadata = true;
log_message(LOG_LEVEL_INFO, "Enabled metadata preservation");
} else if (strcmp(argv[i], "-f") == 0 || strcmp(argv[i], "--sendfile") == 0) {
} else if (opt_is(argv[i], "-f", "--sendfile")) {
config->use_sendfile = true;
log_message(LOG_LEVEL_INFO, "Enabled sendfile");
} else if (strcmp(argv[i], "-m") == 0) {
} else if (opt_is(argv[i], "-m", NULL)) {
config->use_multithreading = true;
log_message(LOG_LEVEL_INFO, "Enabled Multithreading");
} else if (strcmp(argv[i], "-s") == 0) {
} else if (opt_is(argv[i], "-s", NULL)) {
config->use_chunk_serialization = true;
log_message(LOG_LEVEL_INFO, "Enabled Chunk Serialization");
} else if (strcmp(argv[i], "--server-host") == 0 && i + 1 < argc) {
if (set_string_option(&config->server_host, argv[++i], "--server-host") != 0)
return -1;
} else if (strcmp(argv[i], "--server-port") == 0 && i + 1 < argc) {
} else if (opt_is(argv[i], "--server-port", NULL) && i + 1 < argc) {
if (!parse_positive_int(argv[++i], &config->server_port)) {
fprintf(stderr, "Error: invalid --server-port value: %s\n", argv[i]);
log_message(LOG_LEVEL_ERROR, "invalid --server-port value: %s", argv[i]);
return -1;
}
if (config->server_port > 65535) {
fprintf(stderr, "Error: server port must be 1-65535\n");
log_message(LOG_LEVEL_ERROR, "server port must be 1-65535");
return -1;
}
} else if (strcmp(argv[i], "--bwlimit") == 0 && i + 1 < argc) {
char* end;
errno = 0;
unsigned long long kbps = strtoull(argv[++i], &end, 10);
if (errno != 0 || *end != '\0' || kbps == 0) {
fprintf(stderr, "Error: --bwlimit must be a positive integer\n");
} else if (opt_is(argv[i], "--bwlimit", NULL) && i + 1 < argc) {
unsigned long long kbps;
if (parse_ull_arg(argv[++i], &kbps, "--bwlimit") != 0)
return -1;
if (kbps == 0) {
log_message(LOG_LEVEL_ERROR, "--bwlimit must be a positive integer");
return -1;
}
if (kbps > ULLONG_MAX / 1024) {
fprintf(stderr, "Error: --bwlimit value too large\n");
log_message(LOG_LEVEL_ERROR, "--bwlimit value too large");
return -1;
}
io_set_bwlimit(kbps * 1024);
log_message(LOG_LEVEL_INFO, "Set bandwidth limit to %llu KB/s", kbps);
} else if (strcmp(argv[i], "--progress") == 0) {
config->show_progress = true;
} else if (strcmp(argv[i], "--chunk-size") == 0 && i + 1 < argc) {
char* end;
errno = 0;
unsigned long long val = strtoull(argv[++i], &end, 10);
if (errno != 0 || *end != '\0' || val == 0) {
fprintf(stderr, "Error: --chunk-size must be a positive integer\n");
} else if (opt_is(argv[i], "--chunk-size", NULL) && i + 1 < argc) {
unsigned long long val;
if (parse_ull_arg(argv[++i], &val, "--chunk-size") != 0)
return -1;
if (val == 0) {
log_message(LOG_LEVEL_ERROR, "--chunk-size must be a positive integer");
return -1;
}
config->chunk_size = val;
} else if (strcmp(argv[i], "--tls") == 0) {
config->use_tls = true;
} else if (strcmp(argv[i], "--cert") == 0 && i + 1 < argc) {
if (set_string_option(&config->tls_cert, argv[++i], "--cert") != 0)
return -1;
} else if (strcmp(argv[i], "--key") == 0 && i + 1 < argc) {
if (set_string_option(&config->tls_key, argv[++i], "--key") != 0)
return -1;
} else if (strcmp(argv[i], "--ca") == 0 && i + 1 < argc) {
if (set_string_option(&config->tls_ca, argv[++i], "--ca") != 0)
return -1;
} else if (strcmp(argv[i], "--timeout") == 0 && i + 1 < argc) {
if (set_positive_int_option(&config->timeout, argv[++i], "--timeout") != 0)
return -1;
} else if (strcmp(argv[i], "--contimeout") == 0 && i + 1 < argc) {
if (set_positive_int_option(&config->contimeout, argv[++i], "--contimeout") != 0)
return -1;
} else if (strcmp(argv[i], "--backup") == 0) {
config->backup = true;
} else if (strcmp(argv[i], "--backup-dir") == 0 && i + 1 < argc) {
if (set_string_option(&config->backup_dir, argv[++i], "--backup-dir") != 0)
return -1;
} else if (strcmp(argv[i], "--stats") == 0) {
config->stats = true;
} else if (strcmp(argv[i], "--max-depth") == 0 && i + 1 < argc) {
if (set_nonneg_int_option(&config->max_depth, argv[++i], "--max-depth") != 0)
return -1;
} else if (strcmp(argv[i], "--log-file") == 0 && i + 1 < argc) {
} else if (opt_is(argv[i], "--log-file", NULL) && i + 1 < argc) {
if (config->log_file) {
fclose(config->log_file);
config->log_file = NULL;
@@ -292,55 +358,29 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
}
FILE* lf = fopen(argv[++i], "a");
if (!lf) {
fprintf(stderr, "Error: could not open log file '%s': %s\n", argv[i], strerror(errno));
log_message(LOG_LEVEL_ERROR, "could not open log file '%s': %s", argv[i], strerror(errno));
return -1;
}
config->log_file = lf;
log_set_file(lf);
} else if (strcmp(argv[i], "--exclude-from") == 0 && i + 1 < argc) {
} else if (opt_is(argv[i], "--exclude-from", NULL) && i + 1 < argc) {
if (read_patterns_from_file(argv[++i], &config->exclude_patterns, &config->exclude_count) !=
0)
return -1;
} else if (strcmp(argv[i], "--include-from") == 0 && i + 1 < argc) {
} else if (opt_is(argv[i], "--include-from", NULL) && i + 1 < argc) {
if (read_patterns_from_file(argv[++i], &config->include_patterns, &config->include_count) !=
0)
return -1;
} else if (strcmp(argv[i], "--partial") == 0) {
config->partial = true;
} else if (strcmp(argv[i], "--fastsync-server-path") == 0 && i + 1 < argc) {
if (set_string_option(&config->fastsync_server_path, argv[++i], "--fastsync-server-path") !=
0)
return -1;
} else if (strcmp(argv[i], "-v") == 0 || strcmp(argv[i], "--verbose") == 0) {
} else if (opt_is(argv[i], "-v", "--verbose")) {
set_log_level(LOG_LEVEL_DEBUG);
} else if (strcmp(argv[i], "-l") == 0 || strcmp(argv[i], "--links") == 0) {
config->follow_symlinks = true;
} else if (strcmp(argv[i], "--copy-links") == 0) {
config->copy_links = true;
} else if (strcmp(argv[i], "--safe-links") == 0) {
config->safe_links = true;
} else if (strcmp(argv[i], "--copy-unsafe-links") == 0) {
config->copy_unsafe_links = true;
} else if (strcmp(argv[i], "-S") == 0 || strcmp(argv[i], "--sparse") == 0) {
config->preserve_sparse = true;
} else if (strcmp(argv[i], "--inplace") == 0) {
config->inplace = true;
} else if (strcmp(argv[i], "--partial-dir") == 0 && i + 1 < argc) {
if (set_string_option(&config->partial_dir, argv[++i], "--partial-dir") != 0)
return -1;
} else if (strcmp(argv[i], "--suffix") == 0 && i + 1 < argc) {
if (set_string_option(&config->suffix, argv[++i], "--suffix") != 0)
return -1;
} else if (strcmp(argv[i], "-T") == 0 && i + 1 < argc) {
} else if (opt_is(argv[i], "-T", NULL) && i + 1 < argc) {
if (set_positive_int_option(&config->timeout, argv[++i], "-T") != 0)
return -1;
} else if (strcmp(argv[i], "--checksum") == 0) {
config->checksum = true;
} else if (strcmp(argv[i], "--compress-level") == 0 && i + 1 < argc) {
} else if (opt_is(argv[i], "--compress-level", NULL) && i + 1 < argc) {
if (set_positive_int_option(&config->compression_level, argv[++i], "--compress-level") != 0)
return -1;
if (config->compression_level < 1 || config->compression_level > 22) {
fprintf(stderr, "Error: --compress-level must be between 1 and 22\n");
log_message(LOG_LEVEL_ERROR, "--compress-level must be between 1 and 22");
return -1;
}
} else if (argv[i][0] == '-') {
@@ -360,59 +400,10 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
return 0;
}
#ifndef FASTSYNC_TEST_BUILD
/* Validate config after parsing. Returns true if valid. */
static bool validate_config(const Config* config) {
if (!config->send_directory || !config->receive_root_directory) {
fprintf(stderr, "Error: source and destination directories are required\n");
print_usage();
return false;
}
if (config->use_sendfile && (config->use_chunk_serialization || config->use_compression)) {
fprintf(stderr, "Error: -f/--sendfile cannot be combined with -c (compression) or -s (chunk "
"serialization)\n");
return false;
}
if (config->transport == TRANSPORT_SSH && config->use_sendfile) {
fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n");
return false;
}
if (config->use_incremental && config->use_chunk_serialization) {
fprintf(stderr, "Error: --incremental is not supported with -s (chunk serialization)\n");
return false;
}
if (config->use_delta && !config->use_incremental) {
fprintf(stderr, "Error: --delta requires --incremental\n");
return false;
}
if (config->use_delta && config->use_chunk_serialization) {
fprintf(stderr, "Error: --delta cannot be combined with -s (chunk serialization)\n");
return false;
}
if (config->use_delta && config->use_sendfile) {
fprintf(stderr, "Error: --delta cannot be combined with -f (sendfile)\n");
return false;
}
if (config->append || config->append_verify) {
fprintf(
stderr,
"Error: --append and --append-verify are not supported yet; refusing to ignore option\n");
return false;
}
if (config->use_tls) {
if (!config->tls_cert || !config->tls_key) {
fprintf(stderr, "Error: --tls requires --cert and --key\n");
return false;
}
}
return true;
}
#endif /* FASTSYNC_TEST_BUILD */
static int read_patterns_from_file(const char* filepath, char*** patterns, int* count) {
FILE* fp = fopen(filepath, "r");
if (!fp) {
fprintf(stderr, "Error: could not open pattern file '%s': %s\n", filepath, strerror(errno));
log_message(LOG_LEVEL_ERROR, "could not open pattern file '%s': %s", filepath, strerror(errno));
return -1;
}
char* line = NULL;
@@ -429,22 +420,11 @@ static int read_patterns_from_file(const char* filepath, char*** patterns, int*
p[--len] = '\0';
if (len == 0)
continue;
char** tmp = realloc(*patterns, (*count + 1) * sizeof(char*));
if (!tmp) {
fprintf(stderr, "Error: memory allocation failed for pattern file\n");
if (config_add_pattern(patterns, count, p, "pattern file") != 0) {
free(line);
fclose(fp);
return -1;
}
*patterns = tmp;
char* dup = str_dup(p);
if (!dup) {
fprintf(stderr, "Error: memory allocation failed for pattern file\n");
free(line);
fclose(fp);
return -1;
}
(*patterns)[(*count)++] = dup;
}
free(line);
fclose(fp);
@@ -459,10 +439,9 @@ int main(int argc, char* argv[]) {
parse_environment(&env_source, &env_dest, &save_to_disk);
int exit_code = 0;
bool config_owned_by_pipeline = false;
Config* config = config_create();
if (!config) {
fprintf(stderr, "Error: failed to allocate config\n");
log_message(LOG_LEVEL_ERROR, "failed to allocate config");
return 1;
}
config->save_to_disk = save_to_disk;
@@ -483,20 +462,20 @@ int main(int argc, char* argv[]) {
free(config->receive_root_directory);
config->send_directory = str_dup(argv[positional_args[0]]);
if (!config->send_directory) {
fprintf(stderr, "Error: memory allocation failed\n");
log_message(LOG_LEVEL_ERROR, "memory allocation failed");
exit_code = 1;
goto cleanup;
}
config->receive_root_directory = str_dup(argv[positional_args[1]]);
if (!config->receive_root_directory) {
fprintf(stderr, "Error: memory allocation failed\n");
log_message(LOG_LEVEL_ERROR, "memory allocation failed");
exit_code = 1;
goto cleanup;
}
config->save_to_disk = true;
config_parse_ssh_dest(config);
} else if (positional_count == 1) {
fprintf(stderr, "Error: missing destination argument\n");
log_message(LOG_LEVEL_ERROR, "missing destination argument");
print_usage();
exit_code = 1;
goto cleanup;
@@ -504,7 +483,7 @@ int main(int argc, char* argv[]) {
if (!config->send_directory && env_source) {
config->send_directory = str_dup(env_source);
if (!config->send_directory) {
fprintf(stderr, "Error: memory allocation failed\n");
log_message(LOG_LEVEL_ERROR, "memory allocation failed");
exit_code = 1;
goto cleanup;
}
@@ -512,7 +491,7 @@ int main(int argc, char* argv[]) {
if (!config->receive_root_directory && env_dest) {
config->receive_root_directory = str_dup(env_dest);
if (!config->receive_root_directory) {
fprintf(stderr, "Error: memory allocation failed\n");
log_message(LOG_LEVEL_ERROR, "memory allocation failed");
exit_code = 1;
goto cleanup;
}
@@ -542,18 +521,14 @@ int main(int argc, char* argv[]) {
/* Execute transfer */
if (config->use_multithreading) {
config_owned_by_pipeline = true;
exit_code = send_files_multithreaded(config);
exit_code = send_files_multithreaded(&config);
} else {
exit_code = send_files(config);
}
cleanup:
if (config) {
if (config->log_file)
fclose(config->log_file);
if (!config_owned_by_pipeline)
config_delete(config);
config_delete(config);
}
return exit_code;
}
+221 -216
View File
@@ -28,6 +28,88 @@
/* Forward declaration for progress-reporting thread used in multithreaded send. */
static int progress_thread_fn(void* arg);
static ScannerOptions scanner_options_from_config(const Config* config, int num_threads) {
ScannerOptions options = {
config->use_metadata, config->chunk_size, config->exclude_patterns,
config->exclude_count, config->include_patterns, config->include_count,
config->max_size, config->min_size, config->max_depth,
num_threads, config->follow_symlinks, config->copy_links,
config->safe_links, config->copy_unsafe_links, config->checksum};
return options;
}
/* Select the configured transport for both transfer execution paths. */
static Client* connect_transfer_client(const Config* config) {
if (config->transport == TRANSPORT_SSH) {
if (config->use_sendfile) {
log_message(LOG_LEVEL_ERROR, "-f/--sendfile is not supported with SSH transport");
return NULL;
}
return client_connect_ssh(config->ssh_destination, config->ssh_port,
config->fastsync_server_path);
}
Client* client = client_create();
if (!client)
return NULL;
bool connected;
if (config->use_tls) {
connected = client_connect_tls(client, config->server_host, config->server_port,
config->tls_cert, config->tls_key, config->tls_ca);
} else {
connected = client_connect(client, config->server_host, config->server_port);
}
if (!connected) {
client_disconnect(client);
client_delete(client);
return NULL;
}
return client;
}
static void disconnect_transfer_client(Client* client) {
if (!client)
return;
client_disconnect(client);
client_delete(client);
}
static ArrayList* create_transfer_manifest(const Config* config) {
return config->use_delete ? array_list_create(free) : NULL;
}
static bool add_chunk_to_manifest(ArrayList* manifest, const Chunk* chunk) {
if (!manifest)
return true;
for (int i = 0; i < chunk->element_count; i++) {
const char* path = chunk->items[i]->path;
if (*path == '/')
path++;
char* entry = str_dup(path);
if (!entry) {
log_message(LOG_LEVEL_ERROR, "Failed to allocate manifest entry");
return false;
}
if (!array_list_add(manifest, entry)) {
free(entry);
return false;
}
}
return true;
}
static bool finalize_transfer(Client* client) {
Status status;
return send_status(client->file_descriptor, STATUS_FINISHED) &&
receive_status(client->file_descriptor, &status) && status == STATUS_OK;
}
static void mark_sender_done(PipelineContextSender* context) {
mtx_lock(&context->mutex_progress);
context->sender_done = true;
mtx_unlock(&context->mutex_progress);
}
static void pipeline_cancel(PipelineContextSender* context) {
mtx_lock(&context->mutex_scanner);
mtx_lock(&context->mutex_loader);
@@ -43,12 +125,10 @@ static void pipeline_cancel(PipelineContextSender* context) {
}
/* Print dry-run manifest showing files that would be transferred. Returns 0 on success. */
static int send_dry_run_manifest(Config* config) {
DirectoryScanner* scanner = directory_scanner_create(
config->send_directory, config->use_metadata, config->chunk_size, config->exclude_patterns,
config->exclude_count, config->include_patterns, config->include_count, config->max_size,
config->min_size, config->max_depth, config->follow_symlinks, config->copy_links,
config->safe_links, config->copy_unsafe_links, config->checksum);
static int send_dry_run_manifest(const Config* config) {
ScannerOptions options = scanner_options_from_config(config, 0);
DirectoryScanner* scanner =
directory_scanner_create_with_options(config->send_directory, &options);
if (!scanner)
return -1;
Chunk* chunk;
@@ -70,6 +150,8 @@ static int send_dry_run_manifest(Config* config) {
/* Send the delete manifest (list of files) to the server. Returns 0 on success, -1 on failure. */
static int send_delete_manifest(int fd, ArrayList* manifest) {
if (!manifest)
return -1;
if (!send_status(fd, STATUS_MANIFEST))
return -1;
if (!send_int(fd, manifest->size))
@@ -111,17 +193,22 @@ static int incremental_check(Client* client, File* file, const Config* config,
return 1;
if (s == STATUS_DELTA_SIGNATURE) {
Data* sig_data = receive_data(client->file_descriptor);
if (!sig_data)
if (!sig_data) {
send_status(client->file_descriptor, STATUS_ERROR);
return -1;
}
DeltaSignature* sig = delta_signature_deserialize(sig_data);
data_destroy(sig_data);
if (!sig)
if (!sig) {
send_status(client->file_descriptor, STATUS_ERROR);
return -1;
}
*out_sig = sig;
return 2;
}
if (s != STATUS_NEXT) {
log_message(LOG_LEVEL_ERROR, "Unexpected server status");
send_status(client->file_descriptor, STATUS_ERROR);
return -1;
}
return 0;
@@ -129,8 +216,10 @@ static int incremental_check(Client* client, File* file, const Config* config,
static int send_delta(Client* client, File* file, DeltaSignature* sig, Config* config) {
Delta* delta = delta_compute(file->data->data, file->data->size, sig, config->delta_block_size);
/* The receiver is blocked after sending the signature. Every local
fallback therefore needs the explicit NEXT response before full data. */
if (!delta)
return 1;
return send_status(client->file_descriptor, STATUS_NEXT) ? 1 : -1;
if (!delta_is_worthwhile(delta, file->data->size)) {
delta_destroy(delta);
@@ -142,14 +231,14 @@ static int send_delta(Client* client, File* file, DeltaSignature* sig, Config* c
Data* delta_data = delta_serialize(delta);
delta_destroy(delta);
if (!delta_data)
return -1;
return send_status(client->file_descriptor, STATUS_NEXT) ? 1 : -1;
Data* to_send = delta_data;
if (config->use_compression) {
to_send = data_compress(delta_data, config->compression_level);
data_destroy(delta_data);
if (!to_send)
return -1;
return send_status(client->file_descriptor, STATUS_NEXT) ? 1 : -1;
}
bool ok = send_status(client->file_descriptor, STATUS_DELTA_DATA) &&
@@ -279,7 +368,8 @@ int send_chunk(Client* client, Chunk* chunk, Config* config) {
if (f == NULL)
continue;
bool stream = f->data->data == NULL && f->data->size > 0;
bool use_sendfile = (config->use_sendfile && !config->use_compression) || stream;
bool use_sendfile =
(config->use_sendfile && !config->use_compression) || (stream && !config->use_compression);
int rc = send_single_file(client, f, config, config->use_incremental, use_sendfile);
if (rc == 1)
continue;
@@ -291,49 +381,24 @@ int send_chunk(Client* client, Chunk* chunk, Config* config) {
static int send_chunks_multithreaded(void* pipeline_context) {
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
Client* client;
if (context->config->transport == TRANSPORT_SSH) {
if (context->config->use_sendfile) {
fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n");
mtx_lock(&context->mutex_progress);
context->sender_done = true;
mtx_unlock(&context->mutex_progress);
return 1;
}
client = client_connect_ssh(context->config->ssh_destination, context->config->ssh_port,
context->config->fastsync_server_path);
} else if (context->config->use_tls) {
client = client_create();
if (!client || !client_connect_tls(client, context->config->server_host,
context->config->server_port, context->config->tls_cert,
context->config->tls_key, context->config->tls_ca)) {
if (client)
client_delete(client);
fprintf(stderr, "Error: could not connect to server via TLS\n");
mtx_lock(&context->mutex_progress);
context->sender_done = true;
mtx_unlock(&context->mutex_progress);
return thrd_error;
}
} else {
client = client_create();
if (!client ||
!client_connect(client, context->config->server_host, context->config->server_port)) {
if (client)
client_delete(client);
fprintf(stderr, "Error: could not connect to server\n");
mtx_lock(&context->mutex_progress);
context->sender_done = true;
mtx_unlock(&context->mutex_progress);
return thrd_error;
}
Client* client = connect_transfer_client(context->config);
if (!client) {
if (context->config->transport == TRANSPORT_TCP)
log_message(LOG_LEVEL_ERROR, "could not connect to server%s",
context->config->use_tls ? " via TLS" : "");
pipeline_cancel(context);
mark_sender_done(context);
return thrd_error;
}
ProtocolSession session;
protocol_session_init(&session, client->file_descriptor, client->file_descriptor);
protocol_session_set_ssl(&session, (SSL*)client->ssl);
protocol_session_bind(&session);
if (!config_send(client->file_descriptor, context->config)) {
client_disconnect(client);
client_delete(client);
mtx_lock(&context->mutex_progress);
context->sender_done = true;
mtx_unlock(&context->mutex_progress);
pipeline_cancel(context);
disconnect_transfer_client(client);
mark_sender_done(context);
protocol_session_unbind();
return thrd_error;
}
@@ -342,39 +407,37 @@ static int send_chunks_multithreaded(void* pipeline_context) {
context->queue_loader, &context->mutex_loader, &context->condition_not_empty_loader,
&context->condition_not_full_loader, &context->loader_done);
if (current_chunk == NULL) {
if (atomic_load(&context->cancelled)) {
pipeline_cancel(context);
disconnect_transfer_client(client);
mark_sender_done(context);
protocol_session_unbind();
return thrd_error;
}
if (context->config->use_delete) {
if (send_delete_manifest(client->file_descriptor, context->manifest) != 0)
goto send_fail;
}
if (!send_status(client->file_descriptor, STATUS_FINISHED))
goto send_fail;
Status s;
int ok = receive_status(client->file_descriptor, &s) && s == STATUS_OK;
client_disconnect(client);
client_delete(client);
mtx_lock(&context->mutex_progress);
context->sender_done = true;
mtx_unlock(&context->mutex_progress);
bool ok = finalize_transfer(client);
disconnect_transfer_client(client);
mark_sender_done(context);
protocol_session_unbind();
return ok ? thrd_success : thrd_error;
send_fail:
pipeline_cancel(context);
client_disconnect(client);
client_delete(client);
mtx_lock(&context->mutex_progress);
context->sender_done = true;
mtx_unlock(&context->mutex_progress);
disconnect_transfer_client(client);
mark_sender_done(context);
protocol_session_unbind();
return thrd_error;
}
if (send_chunk(client, current_chunk, context->config) != 0) {
fprintf(stderr, "Error: unexpected error while sending chunk\n");
log_message(LOG_LEVEL_ERROR, "unexpected error while sending chunk");
chunk_destroy(current_chunk);
pipeline_cancel(context);
client_disconnect(client);
client_delete(client);
mtx_lock(&context->mutex_progress);
context->sender_done = true;
mtx_unlock(&context->mutex_progress);
disconnect_transfer_client(client);
mark_sender_done(context);
protocol_session_unbind();
return thrd_error;
}
if (context->config->show_progress) {
@@ -393,13 +456,9 @@ static int send_chunks_multithreaded(void* pipeline_context) {
static int scan_directory_multithreaded(void* pipeline_context) {
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
ParallelScanner* scanner = parallel_scanner_create(
context->config->send_directory, context->config->use_metadata, context->config->chunk_size,
context->config->exclude_patterns, context->config->exclude_count,
context->config->include_patterns, context->config->include_count, context->config->max_size,
context->config->min_size, context->config->max_depth, 4, context->config->follow_symlinks,
context->config->copy_links, context->config->safe_links, context->config->copy_unsafe_links,
context->config->checksum);
ScannerOptions options = scanner_options_from_config(context->config, 4);
ParallelScanner* scanner =
parallel_scanner_create_with_options(context->config->send_directory, &options);
Chunk* current_chunk;
if (scanner == NULL) {
@@ -410,28 +469,14 @@ static int scan_directory_multithreaded(void* pipeline_context) {
while ((current_chunk = parallel_scanner_next(scanner)) != NULL) {
if (context->config->use_delete) {
mtx_lock(&context->mutex_scanner);
for (int i = 0; i < current_chunk->element_count; i++) {
const char* p = current_chunk->items[i]->path;
if (*p == '/')
p++;
char* manifest_entry = str_dup(p);
if (!manifest_entry) {
log_message(LOG_LEVEL_ERROR, "Failed to allocate manifest entry");
mtx_unlock(&context->mutex_scanner);
pipeline_cancel(context);
parallel_scanner_destroy(scanner);
return thrd_error;
}
if (!array_list_add(context->manifest, manifest_entry)) {
free(manifest_entry);
mtx_unlock(&context->mutex_scanner);
pipeline_cancel(context);
chunk_destroy(current_chunk);
parallel_scanner_destroy(scanner);
return thrd_error;
}
}
bool manifest_ok = add_chunk_to_manifest(context->manifest, current_chunk);
mtx_unlock(&context->mutex_scanner);
if (!manifest_ok) {
pipeline_cancel(context);
chunk_destroy(current_chunk);
parallel_scanner_destroy(scanner);
return thrd_error;
}
}
if (!queue_enqueue_multithreaded_cancel(
context->queue_scanner, current_chunk, &context->mutex_scanner,
@@ -478,12 +523,13 @@ static int load_files_multithreaded(void* pipeline_context) {
if (!context->config->use_sendfile) {
for (int i = 0; i < chunk->element_count; i++) {
File* f = chunk->items[i];
if (f->data->size > STREAM_THRESHOLD)
if (f->data->size > STREAM_THRESHOLD && !context->config->use_compression)
continue;
if (!file_load_data(f)) {
log_message(LOG_LEVEL_ERROR, "Failed to load file data, skipping");
file_destroy(f);
chunk->items[i] = NULL;
log_message(LOG_LEVEL_ERROR, "Failed to load file data");
chunk_destroy(chunk);
pipeline_cancel(context);
return thrd_error;
}
}
}
@@ -492,14 +538,23 @@ static int load_files_multithreaded(void* pipeline_context) {
&context->condition_not_full_loader,
&context->cancelled)) {
chunk_destroy(chunk);
atomic_store(&context->cancelled, true);
cnd_broadcast(&context->condition_not_full_loader);
cnd_broadcast(&context->condition_not_empty_loader);
pipeline_cancel(context);
return thrd_error;
}
}
}
/* Print a one-line transfer progress report to stderr. `suffix` ends the
line (e.g. "Done.\n") or is "" for in-place refresh. Shared by the
single-threaded loop and the multithreaded progress thread. */
static void print_transfer_progress(unsigned long long total_bytes, time_t start,
const char* suffix) {
double elapsed = difftime(time(NULL), start);
double rate = elapsed > 0.0 ? total_bytes / (1048576.0 * elapsed) : 0.0;
fprintf(stderr, "\rSent %.1f MB (%.1f MB/s) %s", total_bytes / 1048576.0, rate, suffix);
fflush(stderr);
}
/* Progress-reporting thread for multithreaded send. Runs in parallel with
the scanner/loader/sender threads and prints periodic progress to stderr. */
static int progress_thread_fn(void* arg) {
@@ -514,20 +569,14 @@ static int progress_thread_fn(void* arg) {
mtx_unlock(&context->mutex_progress);
if (done) {
time_t now = time(NULL);
double elapsed = difftime(now, start);
double rate = elapsed > 0.0 ? total / (1048576.0 * elapsed) : 0.0;
fprintf(stderr, "\rSent %.1f MB (%.1f MB/s) Done.\n", total / 1048576.0, rate);
print_transfer_progress(total, start, "Done.\n");
break;
}
time_t now = time(NULL);
if (now - last_progress >= 1) {
last_progress = now;
double elapsed = difftime(now, start);
double rate = elapsed > 0.0 ? total / (1048576.0 * elapsed) : 0.0;
fprintf(stderr, "\rSent %.1f MB (%.1f MB/s) ", total / 1048576.0, rate);
fflush(stderr);
print_transfer_progress(total, start, "");
}
struct timespec ts = {0, 100 * 1000000L}; /* 100 ms */
@@ -540,103 +589,58 @@ int send_files(Config* config) {
if (config->dry_run)
return send_dry_run_manifest(config);
Client* client;
if (config->transport == TRANSPORT_SSH) {
if (config->use_sendfile) {
fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n");
return 1;
}
client =
client_connect_ssh(config->ssh_destination, config->ssh_port, config->fastsync_server_path);
if (!client)
return 1;
} else if (config->use_tls) {
client = client_create();
if (!client || !client_connect_tls(client, config->server_host, config->server_port,
config->tls_cert, config->tls_key, config->tls_ca)) {
if (client) {
client_disconnect(client);
client_delete(client);
}
fprintf(stderr, "Error: could not connect to server via TLS\n");
return 1;
}
} else {
client = client_create();
if (!client || !client_connect(client, config->server_host, config->server_port)) {
if (client) {
client_disconnect(client);
client_delete(client);
}
fprintf(stderr, "Error: could not connect to server\n");
return 1;
}
}
if (!config_send(client->file_descriptor, config)) {
client_disconnect(client);
client_delete(client);
Client* client = connect_transfer_client(config);
if (!client) {
if (config->transport == TRANSPORT_TCP)
log_message(LOG_LEVEL_ERROR, "could not connect to server%s",
config->use_tls ? " via TLS" : "");
return 1;
}
DirectoryScanner* scanner = directory_scanner_create(
config->send_directory, config->use_metadata, config->chunk_size, config->exclude_patterns,
config->exclude_count, config->include_patterns, config->include_count, config->max_size,
config->min_size, config->max_depth, config->follow_symlinks, config->copy_links,
config->safe_links, config->copy_unsafe_links, config->checksum);
ProtocolSession session;
protocol_session_init(&session, client->file_descriptor, client->file_descriptor);
protocol_session_set_ssl(&session, (SSL*)client->ssl);
protocol_session_bind(&session);
int ret = 1;
DirectoryScanner* scanner = NULL;
ArrayList* manifest = NULL;
if (!config_send(client->file_descriptor, config))
goto send_fail;
ScannerOptions scanner_options = scanner_options_from_config(config, 0);
scanner = directory_scanner_create_with_options(config->send_directory, &scanner_options);
manifest = create_transfer_manifest(config);
if (!scanner || (config->use_delete && !manifest))
goto send_fail;
Chunk* current_chunk;
unsigned long long total_bytes = 0;
int total_files = 0;
time_t last_progress = 0;
time_t start = time(NULL);
ArrayList* manifest = config->use_delete ? array_list_create(free) : NULL;
if (!scanner || (config->use_delete && !manifest)) {
if (scanner)
directory_scanner_destroy(scanner);
if (manifest)
array_list_delete(manifest);
client_disconnect(client);
client_delete(client);
return 1;
}
while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
unsigned long long chunk_bytes = 0;
for (int i = 0; i < current_chunk->element_count; i++) {
chunk_bytes += current_chunk->items[i]->data->size;
total_files++;
if (manifest) {
const char* p = current_chunk->items[i]->path;
if (*p == '/')
p++;
char* manifest_entry = str_dup(p);
if (!manifest_entry) {
log_message(LOG_LEVEL_ERROR, "Failed to allocate manifest entry");
chunk_destroy(current_chunk);
array_list_delete(manifest);
directory_scanner_destroy(scanner);
client_disconnect(client);
client_delete(client);
return 1;
}
if (!array_list_add(manifest, manifest_entry)) {
free(manifest_entry);
chunk_destroy(current_chunk);
array_list_delete(manifest);
directory_scanner_destroy(scanner);
client_disconnect(client);
client_delete(client);
return 1;
}
}
}
if (!add_chunk_to_manifest(manifest, current_chunk)) {
chunk_destroy(current_chunk);
goto send_fail;
}
if (!config->use_sendfile) {
bool load_ok = true;
for (int i = 0; i < current_chunk->element_count; i++) {
File* f = current_chunk->items[i];
if (f->data->size > STREAM_THRESHOLD)
if (f->data->size > STREAM_THRESHOLD && !config->use_compression)
continue;
if (!file_load_data(f)) {
log_message(LOG_LEVEL_ERROR, "Failed to load file data");
continue;
load_ok = false;
break;
}
}
if (!load_ok) {
chunk_destroy(current_chunk);
goto send_fail;
}
}
if (send_chunk(client, current_chunk, config) != 0) {
log_message(LOG_LEVEL_ERROR, "Failed to send chunk");
@@ -651,10 +655,7 @@ int send_files(Config* config) {
time_t now = time(NULL);
if (now - last_progress >= 1) {
last_progress = now;
double elapsed = difftime(now, start);
double rate = elapsed > 0 ? total_bytes / (1048576.0 * elapsed) : 0;
fprintf(stderr, "\rSent %.1f MB (%.1f MB/s) ", total_bytes / 1048576.0, rate);
fflush(stderr);
print_transfer_progress(total_bytes, start, "");
}
}
chunk_destroy(current_chunk);
@@ -664,40 +665,39 @@ int send_files(Config* config) {
if (config->use_delete) {
if (send_delete_manifest(client->file_descriptor, manifest) != 0) {
array_list_delete(manifest);
manifest = NULL;
goto send_fail;
}
array_list_delete(manifest);
manifest = NULL;
}
if (!send_status(client->file_descriptor, STATUS_FINISHED))
goto send_fail;
Status s;
int ok = receive_status(client->file_descriptor, &s) && s == STATUS_OK;
double elapsed_total = difftime(time(NULL), start);
if (config->show_progress) {
double rate = elapsed_total > 0 ? total_bytes / (1048576.0 * elapsed_total) : 0;
fprintf(stderr, "\rSent %.1f MB (%.1f MB/s) Done.\n", total_bytes / 1048576.0, rate);
}
bool ok = finalize_transfer(client);
if (config->show_progress)
print_transfer_progress(total_bytes, start, "Done.\n");
if (config->stats) {
double elapsed_total = difftime(time(NULL), start);
double rate = elapsed_total > 0 ? total_bytes / (1048576.0 * elapsed_total) : 0;
fprintf(stderr, "Stats: %d files, %.1f MB, %.1f MB/s\n", total_files, total_bytes / 1048576.0,
rate);
}
directory_scanner_destroy(scanner);
client_disconnect(client);
client_delete(client);
return ok ? 0 : 1;
ret = ok ? 0 : 1;
send_fail:
/* Single cleanup path for all exits. The manifest is intentionally deleted
here even on success without --delete, fixing a pre-existing leak. */
if (manifest)
array_list_delete(manifest);
directory_scanner_destroy(scanner);
client_disconnect(client);
client_delete(client);
return 1;
if (scanner)
directory_scanner_destroy(scanner);
disconnect_transfer_client(client);
protocol_session_unbind();
return ret;
}
int send_files_multithreaded(Config* config) {
int send_files_multithreaded(Config** config_ptr) {
if (!config_ptr || !*config_ptr)
return 1;
Config* config = *config_ptr;
if (config->dry_run)
return send_dry_run_manifest(config);
@@ -728,8 +728,13 @@ int send_files_multithreaded(Config* config) {
queue_destroy(q2);
return 1;
}
*config_ptr = NULL; /* context now owns config through all remaining paths */
if (config->use_delete)
context->manifest = array_list_create(free);
if (config->use_delete && !context->manifest) {
pipeline_context_sender_destroy(context);
return 1;
}
thrd_t scanner, loader, sender;
bool scanner_created = false;
@@ -743,7 +748,7 @@ int send_files_multithreaded(Config* config) {
sender_created = (thrd_create(&sender, send_chunks_multithreaded, context) == thrd_success);
if (!scanner_created || !loader_created || !sender_created) {
perror("Error creating threads.\n");
log_perror("Error creating threads");
pipeline_cancel(context);
mtx_lock(&context->mutex_progress);
context->sender_done = true;
@@ -763,7 +768,7 @@ int send_files_multithreaded(Config* config) {
if (config->show_progress) {
progress_created = (thrd_create(&progress, progress_thread_fn, context) == thrd_success);
if (!progress_created) {
perror("Error creating progress thread.\n");
log_perror("Error creating progress thread");
/* Non-fatal; continue without progress reporting */
}
}
+2 -1
View File
@@ -7,6 +7,7 @@
int send_chunk(Client* client, Chunk* chunk, Config* config);
int send_files(Config* config);
int send_files_multithreaded(Config* config);
/* Takes ownership only when *config is set to NULL on return. */
int send_files_multithreaded(Config** config);
#endif
+51
View File
@@ -0,0 +1,51 @@
#include "client_validation.h"
#include "log.h"
#include "usage.h"
#include <stdio.h>
/* Validate config after parsing. Returns true if valid. */
bool validate_config(const Config* config) {
if (!config->send_directory || !config->receive_root_directory) {
log_message(LOG_LEVEL_ERROR, "source and destination directories are required");
print_usage();
return false;
}
if (config->use_sendfile && (config->use_chunk_serialization || config->use_compression)) {
log_message(LOG_LEVEL_ERROR, "-f/--sendfile cannot be combined with -c (compression) or -s "
"(chunk serialization)");
return false;
}
if (config->transport == TRANSPORT_SSH && config->use_sendfile) {
log_message(LOG_LEVEL_ERROR, "-f/--sendfile is not supported with SSH transport");
return false;
}
if (config->use_incremental && config->use_chunk_serialization) {
log_message(LOG_LEVEL_ERROR, "--incremental is not supported with -s (chunk serialization)");
return false;
}
if (config->use_delta && !config->use_incremental) {
log_message(LOG_LEVEL_ERROR, "--delta requires --incremental");
return false;
}
if (config->use_delta && config->use_chunk_serialization) {
log_message(LOG_LEVEL_ERROR, "--delta cannot be combined with -s (chunk serialization)");
return false;
}
if (config->use_delta && config->use_sendfile) {
log_message(LOG_LEVEL_ERROR, "--delta cannot be combined with -f (sendfile)");
return false;
}
if (config->append || config->append_verify) {
fprintf(
stderr,
"Error: --append and --append-verify are not supported yet; refusing to ignore option\n");
return false;
}
if (config->use_tls) {
if (!config->tls_cert || !config->tls_key || !config->tls_ca) {
log_message(LOG_LEVEL_ERROR, "--tls requires --cert, --key, and --ca");
return false;
}
}
return true;
}
+9
View File
@@ -0,0 +1,9 @@
#ifndef CLIENT_VALIDATION_H
#define CLIENT_VALIDATION_H
#include "config.h"
#include <stdbool.h>
bool validate_config(const Config* config);
#endif
+367 -411
View File
@@ -1,3 +1,4 @@
#include "log.h"
#include "scanner.h"
#include "array_list.h"
#include "chunk.h"
@@ -52,13 +53,84 @@ static bool safe_relative_link(const char* source_root, const char* containing_d
return safe;
}
DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_metadata,
unsigned long long chunk_size, char** exclude_patterns,
int exclude_count, char** include_patterns,
int include_count, unsigned long long max_size,
unsigned long long min_size, int max_depth,
bool follow_symlinks, bool copy_links, bool safe_links,
bool copy_unsafe_links, bool checksum) {
typedef struct {
char* path;
struct stat stats;
bool is_directory;
} ScannerEntry;
/* Inspect symlinks, resolve the entry type, and apply file filters once for both scanners. */
static int scanner_inspect_entry(const ScannerOptions* options, const char* source_root,
const char* containing_dir, const char* name,
ScannerEntry* entry) {
entry->path = path_cat(containing_dir, name);
if (!entry->path)
return -1;
struct stat link_stats;
if (lstat(entry->path, &link_stats) != 0) {
free(entry->path);
return 0;
}
bool is_symlink = S_ISLNK(link_stats.st_mode);
if (is_symlink && !options->follow_symlinks && !options->copy_links && !options->safe_links &&
!options->copy_unsafe_links)
goto skip;
if (is_symlink && options->safe_links) {
char link_target[4096];
ssize_t length = readlink(entry->path, link_target, sizeof(link_target) - 1);
if (length < 0)
goto skip;
link_target[length] = '\0';
if (link_target[0] == '/' || !safe_relative_link(source_root, containing_dir, link_target))
goto skip;
}
if (is_symlink && options->copy_unsafe_links && !options->copy_links) {
char link_target[4096];
ssize_t length = readlink(entry->path, link_target, sizeof(link_target) - 1);
if (length < 0)
goto skip;
link_target[length] = '\0';
if (link_target[0] != '/')
goto skip;
}
if (is_symlink && options->follow_symlinks && !options->copy_links)
entry->stats = link_stats;
else if (stat(entry->path, &entry->stats) != 0)
goto skip;
entry->is_directory = S_ISDIR(entry->stats.st_mode);
if (entry->is_directory)
return 1;
for (int i = 0; i < options->exclude_count; i++)
if (glob_match(options->exclude_patterns[i], name))
goto skip;
if (options->include_count > 0) {
bool included = false;
for (int i = 0; i < options->include_count; i++)
if (glob_match(options->include_patterns[i], name))
included = true;
if (!included)
goto skip;
}
if ((options->max_size > 0 && (unsigned long long)entry->stats.st_size > options->max_size) ||
(options->min_size > 0 && (unsigned long long)entry->stats.st_size < options->min_size))
goto skip;
return 1;
skip:
free(entry->path);
entry->path = NULL;
return 0;
}
DirectoryScanner* directory_scanner_create_with_options(const char* root_directory,
const ScannerOptions* options) {
if (!root_directory || !options)
return NULL;
DirectoryScanner* scanner = calloc(1, sizeof(DirectoryScanner));
if (scanner == NULL)
return NULL;
@@ -69,21 +141,21 @@ DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_
}
scanner->current_dir = NULL;
scanner->current_path = NULL;
scanner->use_metadata = use_metadata;
scanner->chunk_size = chunk_size > 0 ? chunk_size : DESIRED_CHUNK_SIZE;
scanner->exclude_patterns = exclude_patterns;
scanner->exclude_count = exclude_count;
scanner->include_patterns = include_patterns;
scanner->include_count = include_count;
scanner->max_size = max_size;
scanner->min_size = min_size;
scanner->max_depth = max_depth;
scanner->use_metadata = options->use_metadata;
scanner->chunk_size = options->chunk_size > 0 ? options->chunk_size : DESIRED_CHUNK_SIZE;
scanner->exclude_patterns = options->exclude_patterns;
scanner->exclude_count = options->exclude_count;
scanner->include_patterns = options->include_patterns;
scanner->include_count = options->include_count;
scanner->max_size = options->max_size;
scanner->min_size = options->min_size;
scanner->max_depth = options->max_depth;
scanner->current_depth = 0;
scanner->follow_symlinks = follow_symlinks;
scanner->copy_links = copy_links;
scanner->safe_links = safe_links;
scanner->copy_unsafe_links = copy_unsafe_links;
scanner->checksum = checksum;
scanner->follow_symlinks = options->follow_symlinks;
scanner->copy_links = options->copy_links;
scanner->safe_links = options->safe_links;
scanner->copy_unsafe_links = options->copy_unsafe_links;
scanner->checksum = options->checksum;
scanner->failed = false;
DirEntry* root = dir_entry_create(root_directory, 0);
if (!root) {
@@ -100,6 +172,20 @@ DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_
return scanner;
}
DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_metadata,
unsigned long long chunk_size, char** exclude_patterns,
int exclude_count, char** include_patterns,
int include_count, unsigned long long max_size,
unsigned long long min_size, int max_depth,
bool follow_symlinks, bool copy_links, bool safe_links,
bool copy_unsafe_links, bool checksum) {
ScannerOptions options = {
use_metadata, chunk_size, exclude_patterns, exclude_count, include_patterns,
include_count, max_size, min_size, max_depth, 0,
follow_symlinks, copy_links, safe_links, copy_unsafe_links, checksum};
return directory_scanner_create_with_options(root_directory, &options);
}
void directory_scanner_destroy(DirectoryScanner* scanner) {
if (scanner == NULL)
return;
@@ -141,7 +227,7 @@ static int open_next_directory(DirectoryScanner* scanner) {
free(de);
scanner->current_dir = opendir(scanner->current_path);
if (scanner->current_dir == NULL) {
perror("Could not open directory");
log_perror("Could not open directory");
free(scanner->current_path);
scanner->current_path = NULL;
scanner->failed = true;
@@ -167,7 +253,7 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
break;
}
struct dirent* entry = readdir(scanner->current_dir);
const struct dirent* entry = readdir(scanner->current_dir);
if (entry == NULL) {
closedir(scanner->current_dir);
scanner->current_dir = NULL;
@@ -179,67 +265,27 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
char* cur_path = path_cat(scanner->current_path, entry->d_name);
if (!cur_path) {
ScannerOptions options = {scanner->use_metadata, scanner->chunk_size,
scanner->exclude_patterns, scanner->exclude_count,
scanner->include_patterns, scanner->include_count,
scanner->max_size, scanner->min_size,
scanner->max_depth, 0,
scanner->follow_symlinks, scanner->copy_links,
scanner->safe_links, scanner->copy_unsafe_links,
scanner->checksum};
ScannerEntry inspected;
int inspection = scanner_inspect_entry(&options, scanner->current_path, scanner->current_path,
entry->d_name, &inspected);
if (inspection < 0) {
scanner->failed = true;
break;
}
struct stat stats;
struct stat lstats;
bool is_symlink = false;
if (lstat(cur_path, &lstats) != 0) {
free(cur_path);
if (inspection == 0)
continue;
}
is_symlink = S_ISLNK(lstats.st_mode);
char* cur_path = inspected.path;
struct stat stats = inspected.stats;
if (is_symlink && !scanner->follow_symlinks && !scanner->copy_links && !scanner->safe_links &&
!scanner->copy_unsafe_links) {
free(cur_path);
continue;
}
if (is_symlink && scanner->safe_links) {
char link_target[4096];
ssize_t len = readlink(cur_path, link_target, sizeof(link_target) - 1);
if (len < 0) {
free(cur_path);
continue;
}
link_target[len] = '\0';
if (link_target[0] == '/' ||
!safe_relative_link(scanner->current_path, scanner->current_path, link_target)) {
free(cur_path);
continue;
}
}
if (is_symlink && scanner->copy_unsafe_links && !scanner->copy_links) {
char link_target[4096];
ssize_t len = readlink(cur_path, link_target, sizeof(link_target) - 1);
if (len < 0) {
free(cur_path);
continue;
}
link_target[len] = '\0';
bool unsafe = (link_target[0] == '/');
if (!unsafe) {
free(cur_path);
continue;
}
}
bool use_lstat = is_symlink && scanner->follow_symlinks && !scanner->copy_links;
if (use_lstat) {
stats = lstats;
} else {
if (stat(cur_path, &stats) != 0) {
free(cur_path);
continue;
}
}
if (S_ISDIR(stats.st_mode)) {
if (inspected.is_directory) {
int next_depth = scanner->current_depth + 1;
if (scanner->max_depth <= 0 || next_depth < scanner->max_depth) {
DirEntry* de = dir_entry_create(cur_path, next_depth);
@@ -254,38 +300,6 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
free(cur_path);
continue;
}
bool excluded = false;
for (int i = 0; i < scanner->exclude_count; i++) {
if (glob_match(scanner->exclude_patterns[i], entry->d_name)) {
excluded = true;
break;
}
}
if (excluded) {
free(cur_path);
continue;
}
if (scanner->include_count > 0) {
bool included = false;
for (int i = 0; i < scanner->include_count; i++) {
if (glob_match(scanner->include_patterns[i], entry->d_name)) {
included = true;
break;
}
}
if (!included) {
free(cur_path);
continue;
}
}
if ((scanner->max_size > 0 && (unsigned long long)stats.st_size > scanner->max_size) ||
(scanner->min_size > 0 && (unsigned long long)stats.st_size < scanner->min_size)) {
free(cur_path);
continue;
}
File* file = file_create(cur_path);
if (file == NULL) {
free(cur_path);
@@ -336,29 +350,13 @@ typedef struct {
ParallelScanner* ps;
char** dirs;
int dir_count;
bool use_metadata;
unsigned long long chunk_size;
char** exclude_patterns;
int exclude_count;
char** include_patterns;
int include_count;
unsigned long long max_size;
unsigned long long min_size;
int max_depth;
bool follow_symlinks;
bool copy_links;
bool safe_links;
bool copy_unsafe_links;
bool checksum;
ScannerOptions options;
} ParallelWorkerArg;
static int parallel_worker_thread(void* arg) {
ParallelWorkerArg* wa = (ParallelWorkerArg*)arg;
for (int i = 0; i < wa->dir_count; i++) {
DirectoryScanner* ds = directory_scanner_create(
wa->dirs[i], wa->use_metadata, wa->chunk_size, wa->exclude_patterns, wa->exclude_count,
wa->include_patterns, wa->include_count, wa->max_size, wa->min_size, wa->max_depth,
wa->follow_symlinks, wa->copy_links, wa->safe_links, wa->copy_unsafe_links, wa->checksum);
DirectoryScanner* ds = directory_scanner_create_with_options(wa->dirs[i], &wa->options);
if (!ds) {
mtx_lock(&wa->ps->result_mutex);
wa->ps->failed = true;
@@ -366,6 +364,8 @@ static int parallel_worker_thread(void* arg) {
cnd_broadcast(&wa->ps->result_not_empty);
cnd_broadcast(&wa->ps->result_not_full);
mtx_unlock(&wa->ps->result_mutex);
for (int j = i; j < wa->dir_count; j++)
free(wa->dirs[j]);
break;
}
Chunk* chunk;
@@ -401,21 +401,23 @@ static int parallel_worker_thread(void* arg) {
return thrd_success;
}
ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata,
unsigned long long chunk_size, char** exclude_patterns,
int exclude_count, char** include_patterns,
int include_count, unsigned long long max_size,
unsigned long long min_size, int max_depth,
int num_threads, bool follow_symlinks, bool copy_links,
bool safe_links, bool copy_unsafe_links, bool checksum) {
ParallelScanner* ps = calloc(1, sizeof(ParallelScanner));
if (!ps)
return NULL;
static void parallel_scanner_creation_failed(ParallelScanner* ps) {
mtx_lock(&ps->result_mutex);
ps->failed = true;
atomic_store(&ps->cancelled, true);
ps->expected_threads = ps->created_threads;
if (ps->completed >= ps->expected_threads)
ps->done = true;
cnd_broadcast(&ps->result_not_empty);
cnd_broadcast(&ps->result_not_full);
mtx_unlock(&ps->result_mutex);
}
/* Initialize result queue and synchronization primitives. Returns true on success. */
static bool parallel_scanner_init(ParallelScanner* ps) {
ps->result_queue = queue_create(100, chunk_destroy);
if (!ps->result_queue) {
free(ps);
return NULL;
}
if (!ps->result_queue)
return false;
atomic_init(&ps->cancelled, false);
int init = 0;
bool ok = true;
@@ -440,14 +442,220 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
if (init >= 1)
mtx_destroy(&ps->result_mutex);
queue_destroy(ps->result_queue);
free(ps);
ps->result_queue = NULL;
return false;
}
return true;
}
/* Split files into chunks of roughly chunk_size bytes. Returns the first chunk (also stored
* chunks beyond the first are enqueued on `queue`). Nulls out consumed entries in `files`.
* Sets *failed on allocation/enqueue errors. */
static Chunk* batch_files(ArrayList* files, unsigned long long chunk_size, Queue* queue,
bool* failed) {
Chunk* first = NULL;
if (files->size <= 0)
return NULL;
ArrayList* batch = array_list_create(NULL);
if (!batch) {
*failed = true;
return NULL;
}
unsigned long long batch_size = 0;
for (int i = 0; i < files->size; i++) {
File* f = (File*)files->items[i];
if (!array_list_add(batch, f)) {
*failed = true;
break;
}
batch_size += f->data->size;
if (batch_size >= chunk_size || i == files->size - 1) {
void** items = array_list_to_array(batch);
if (!items) {
*failed = true;
array_list_delete(batch);
batch = NULL;
break;
}
Chunk* c = chunk_create((File**)items, batch->size);
free(items);
if (!c) {
*failed = true;
array_list_delete(batch);
batch = NULL;
break;
}
int batch_start = i - batch->size + 1;
for (int j = batch_start; j <= i; j++)
files->items[j] = NULL;
batch->item_destroyer = NULL;
array_list_delete(batch);
batch = NULL;
if (!first) {
first = c;
} else {
if (!queue_enqueue(queue, c)) {
chunk_destroy(c);
*failed = true;
}
}
if (i < files->size - 1) {
batch = array_list_create(NULL);
if (!batch) {
*failed = true;
break;
}
batch_size = 0;
}
}
}
if (batch) {
batch->item_destroyer = NULL;
array_list_delete(batch);
}
return first;
}
/* Scan one root-directory entry into either the subdirs or files list. */
static void scan_root_entry(const ScannerOptions* options, const char* root_directory,
const struct dirent* entry, ArrayList* root_files, ArrayList* subdirs,
ParallelScanner* ps) {
ScannerEntry inspected;
int inspection =
scanner_inspect_entry(options, root_directory, root_directory, entry->d_name, &inspected);
if (inspection < 0) {
ps->failed = true;
return;
}
if (inspection == 0)
return;
char* cur_path = inspected.path;
struct stat st = inspected.stats;
if (inspected.is_directory) {
if (!array_list_add(subdirs, cur_path)) {
free(cur_path);
ps->failed = true;
}
return;
}
File* file = file_create(cur_path);
free(cur_path);
if (!file) {
ps->failed = true;
return;
}
file->data->size = st.st_size;
if (options->use_metadata)
file->metadata = file_metadata_create(&st);
if (options->use_metadata && !file->metadata) {
file_destroy(file);
ps->failed = true;
return;
}
if (!array_list_add(root_files, file)) {
file_destroy(file);
ps->failed = true;
}
}
/* Scan the root directory itself, collecting root files and subdirectories.
* Returns false if the root directory could not be opened. */
static bool scan_root_directory(ParallelScanner* ps, const char* root_directory,
const ScannerOptions* options, ArrayList* root_files,
ArrayList* subdirs) {
DIR* dir = opendir(root_directory);
if (!dir) {
perror("Could not open root directory for parallel scan");
parallel_scanner_destroy(ps);
log_perror("Could not open root directory for parallel scan");
return false;
}
const struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
scan_root_entry(options, root_directory, entry, root_files, subdirs, ps);
}
closedir(dir);
return true;
}
/* Spawn worker threads, one per group of subdirectories. */
static void spawn_parallel_workers(ParallelScanner* ps, ArrayList* subdirs,
const ScannerOptions* options, unsigned long long cs) {
if (subdirs->size <= 0)
return;
int n = options->num_threads > 0 ? options->num_threads : 4;
if (n > subdirs->size)
n = subdirs->size;
ps->num_threads = n;
ps->expected_threads = n;
ps->threads = calloc(n, sizeof(thrd_t));
if (!ps->threads) {
ps->num_threads = 0;
ps->expected_threads = 0;
ps->failed = true;
return;
}
int dirs_per_thread = subdirs->size / n;
int remainder = subdirs->size % n;
int start = 0;
ps->num_threads = 0;
for (int t = 0; t < n; t++) {
int count = dirs_per_thread + (t < remainder ? 1 : 0);
if (count == 0)
break;
ParallelWorkerArg* wa = calloc(1, sizeof(ParallelWorkerArg));
if (!wa) {
parallel_scanner_creation_failed(ps);
break;
}
wa->ps = ps;
wa->dirs = calloc(count, sizeof(char*));
if (!wa->dirs) {
free(wa);
parallel_scanner_creation_failed(ps);
break;
}
bool dup_ok = true;
for (int j = 0; j < count; j++) {
wa->dirs[j] = str_dup((char*)subdirs->items[start + j]);
if (!wa->dirs[j])
dup_ok = false;
}
if (!dup_ok) {
for (int j = 0; j < count; j++)
free(wa->dirs[j]);
free(wa->dirs);
free(wa);
parallel_scanner_creation_failed(ps);
break;
}
wa->dir_count = count;
wa->options = *options;
wa->options.chunk_size = cs;
start += count;
if (thrd_create(&ps->threads[t], parallel_worker_thread, wa) != thrd_success) {
for (int j = 0; j < count; j++)
free(wa->dirs[j]);
free(wa->dirs);
free(wa);
parallel_scanner_creation_failed(ps);
break;
}
ps->num_threads++;
ps->created_threads++;
}
}
ParallelScanner* parallel_scanner_create_with_options(const char* root_directory,
const ScannerOptions* options) {
if (!root_directory || !options)
return NULL;
ParallelScanner* ps = calloc(1, sizeof(ParallelScanner));
if (!ps)
return NULL;
if (!parallel_scanner_init(ps)) {
free(ps);
return NULL;
}
@@ -456,279 +664,22 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
if (!root_files || !subdirs) {
array_list_delete(root_files);
array_list_delete(subdirs);
closedir(dir);
parallel_scanner_destroy(ps);
return NULL;
}
struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
char* cur_path = path_cat(root_directory, entry->d_name);
if (!cur_path)
continue;
struct stat lstats;
if (lstat(cur_path, &lstats) != 0) {
free(cur_path);
continue;
}
bool is_symlink = S_ISLNK(lstats.st_mode);
// Skip symlinks unless the user explicitly enabled following/copying them.
if (is_symlink && !follow_symlinks && !copy_links && !safe_links && !copy_unsafe_links) {
free(cur_path);
continue;
}
// --safe-links: reject symlinks pointing outside the source tree.
if (is_symlink && safe_links) {
char link_target[4096];
ssize_t len = readlink(cur_path, link_target, sizeof(link_target) - 1);
if (len < 0) {
free(cur_path);
continue;
}
link_target[len] = 0;
if (link_target[0] == '/' ||
!safe_relative_link(root_directory, root_directory, link_target)) {
free(cur_path);
continue;
}
}
// --copy-unsafe-links (without --copy-links): only copy absolute symlinks.
if (is_symlink && copy_unsafe_links && !copy_links) {
char link_target[4096];
ssize_t len = readlink(cur_path, link_target, sizeof(link_target) - 1);
if (len < 0) {
free(cur_path);
continue;
}
link_target[len] = 0;
bool unsafe = (link_target[0] == '/');
if (!unsafe) {
free(cur_path);
continue;
}
}
// Determine whether to use lstat or stat results for the entry.
struct stat st;
bool use_lstat_res = is_symlink && follow_symlinks && !copy_links;
if (use_lstat_res) {
st = lstats;
} else {
if (stat(cur_path, &st) != 0) {
free(cur_path);
continue;
}
}
if (S_ISDIR(st.st_mode)) {
if (!array_list_add(subdirs, cur_path)) {
free(cur_path);
ps->failed = true;
}
} else {
bool excluded = false;
for (int i = 0; i < exclude_count; i++) {
if (glob_match(exclude_patterns[i], entry->d_name)) {
excluded = true;
break;
}
}
if (excluded) {
free(cur_path);
continue;
}
if (include_count > 0) {
bool included = false;
for (int i = 0; i < include_count; i++) {
if (glob_match(include_patterns[i], entry->d_name)) {
included = true;
break;
}
}
if (!included) {
free(cur_path);
continue;
}
}
if ((max_size > 0 && (unsigned long long)st.st_size > max_size) ||
(min_size > 0 && (unsigned long long)st.st_size < min_size)) {
free(cur_path);
continue;
}
File* file = file_create(cur_path);
free(cur_path);
if (!file) {
ps->failed = true;
continue;
}
file->data->size = st.st_size;
if (use_metadata)
file->metadata = file_metadata_create(&st);
if (use_metadata && !file->metadata) {
file_destroy(file);
ps->failed = true;
continue;
}
if (!array_list_add(root_files, file)) {
file_destroy(file);
ps->failed = true;
}
}
if (!scan_root_directory(ps, root_directory, options, root_files, subdirs)) {
array_list_delete(root_files);
array_list_delete(subdirs);
parallel_scanner_destroy(ps);
return NULL;
}
closedir(dir);
unsigned long long cs = chunk_size > 0 ? chunk_size : DESIRED_CHUNK_SIZE;
if (root_files->size > 0) {
ArrayList* batch = array_list_create(NULL);
if (!batch) {
ps->failed = true;
array_list_delete(root_files);
array_list_delete(subdirs);
parallel_scanner_destroy(ps);
return NULL;
}
unsigned long long batch_size = 0;
Chunk* first = NULL;
for (int i = 0; i < root_files->size; i++) {
File* f = (File*)root_files->items[i];
if (!array_list_add(batch, f)) {
ps->failed = true;
break;
}
batch_size += f->data->size;
if (batch_size >= cs || i == root_files->size - 1) {
void** items = array_list_to_array(batch);
if (!items) {
ps->failed = true;
batch->item_destroyer = file_destroy;
array_list_delete(batch);
batch = NULL;
break;
}
Chunk* c = chunk_create((File**)items, batch->size);
free(items);
if (!c) {
ps->failed = true;
batch->item_destroyer = file_destroy;
array_list_delete(batch);
batch = NULL;
break;
}
batch->item_destroyer = NULL;
array_list_delete(batch);
batch = NULL;
if (!first) {
first = c;
} else {
if (!queue_enqueue(ps->result_queue, c)) {
chunk_destroy(c);
ps->failed = true;
}
}
if (i < root_files->size - 1) {
batch = array_list_create(NULL);
if (!batch) {
ps->failed = true;
break;
}
batch_size = 0;
}
}
}
if (batch) {
batch->item_destroyer = NULL;
array_list_delete(batch);
}
ps->initial_chunk = first;
root_files->item_destroyer = NULL;
}
unsigned long long cs = options->chunk_size > 0 ? options->chunk_size : DESIRED_CHUNK_SIZE;
ps->initial_chunk = batch_files(root_files, cs, ps->result_queue, &ps->failed);
array_list_delete(root_files);
int n = num_threads > 0 ? num_threads : 4;
if (n > subdirs->size)
n = subdirs->size > 0 ? subdirs->size : 1;
if (subdirs->size > 0) {
ps->num_threads = n;
ps->expected_threads = n;
ps->threads = calloc(n, sizeof(thrd_t));
if (!ps->threads) {
array_list_delete(subdirs);
parallel_scanner_destroy(ps);
return NULL;
}
int dirs_per_thread = subdirs->size / n;
int remainder = subdirs->size % n;
int start = 0;
ps->num_threads = 0;
for (int t = 0; t < n; t++) {
int count = dirs_per_thread + (t < remainder ? 1 : 0);
if (count == 0)
break;
ParallelWorkerArg* wa = calloc(1, sizeof(ParallelWorkerArg));
if (!wa) {
ps->failed = true;
break;
}
wa->ps = ps;
wa->dirs = calloc(count, sizeof(char*));
if (!wa->dirs) {
free(wa);
ps->failed = true;
break;
}
bool dup_ok = true;
for (int j = 0; j < count; j++) {
wa->dirs[j] = str_dup((char*)subdirs->items[start + j]);
if (!wa->dirs[j])
dup_ok = false;
}
if (!dup_ok) {
for (int j = 0; j < count; j++)
free(wa->dirs[j]);
free(wa->dirs);
free(wa);
ps->failed = true;
break;
}
wa->dir_count = count;
wa->use_metadata = use_metadata;
wa->chunk_size = cs;
wa->exclude_patterns = exclude_patterns;
wa->exclude_count = exclude_count;
wa->include_patterns = include_patterns;
wa->include_count = include_count;
wa->max_size = max_size;
wa->min_size = min_size;
wa->max_depth = max_depth;
wa->follow_symlinks = follow_symlinks;
wa->copy_links = copy_links;
wa->safe_links = safe_links;
wa->copy_unsafe_links = copy_unsafe_links;
wa->checksum = checksum;
start += count;
if (thrd_create(&ps->threads[t], parallel_worker_thread, wa) != thrd_success) {
for (int j = 0; j < count; j++)
free(wa->dirs[j]);
free(wa->dirs);
free(wa);
ps->failed = true;
atomic_store(&ps->cancelled, true);
ps->expected_threads = ps->created_threads;
mtx_lock(&ps->result_mutex);
cnd_broadcast(&ps->result_not_empty);
cnd_broadcast(&ps->result_not_full);
mtx_unlock(&ps->result_mutex);
break;
}
ps->num_threads++;
ps->created_threads++;
}
}
spawn_parallel_workers(ps, subdirs, options, cs);
array_list_delete(subdirs);
return ps;
}
@@ -741,6 +692,11 @@ Chunk* parallel_scanner_next(ParallelScanner* ps) {
}
if (ps->num_threads == 0) {
mtx_lock(&ps->result_mutex);
if (!queue_is_empty(ps->result_queue)) {
Chunk* chunk = queue_dequeue(ps->result_queue);
mtx_unlock(&ps->result_mutex);
return chunk;
}
ps->done = true;
mtx_unlock(&ps->result_mutex);
return NULL;
+22 -7
View File
@@ -8,6 +8,24 @@
#include <threads.h>
#include <stdatomic.h>
typedef struct {
bool use_metadata;
unsigned long long chunk_size;
char** exclude_patterns;
int exclude_count;
char** include_patterns;
int include_count;
unsigned long long max_size;
unsigned long long min_size;
int max_depth;
int num_threads;
bool follow_symlinks;
bool copy_links;
bool safe_links;
bool copy_unsafe_links;
bool checksum;
} ScannerOptions;
typedef struct {
Queue* directories;
DIR* current_dir;
@@ -53,17 +71,14 @@ DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_
unsigned long long min_size, int max_depth,
bool follow_symlinks, bool copy_links, bool safe_links,
bool copy_unsafe_links, bool checksum);
DirectoryScanner* directory_scanner_create_with_options(const char* root_directory,
const ScannerOptions* options);
Chunk* directory_scanner_next(DirectoryScanner* scanner);
bool directory_scanner_failed(const DirectoryScanner* scanner);
void directory_scanner_destroy(DirectoryScanner* scanner);
ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata,
unsigned long long chunk_size, char** exclude_patterns,
int exclude_count, char** include_patterns,
int include_count, unsigned long long max_size,
unsigned long long min_size, int max_depth,
int num_threads, bool follow_symlinks, bool copy_links,
bool safe_links, bool copy_unsafe_links, bool checksum);
ParallelScanner* parallel_scanner_create_with_options(const char* root_directory,
const ScannerOptions* options);
Chunk* parallel_scanner_next(ParallelScanner* scanner);
bool parallel_scanner_failed(const ParallelScanner* scanner);
void parallel_scanner_destroy(ParallelScanner* scanner);
+4 -2
View File
@@ -1,8 +1,9 @@
#include "stdio.h"
#include "usage.h"
#include <stdio.h>
#include <delta.h>
#include <chunk.h>
static __attribute__((unused)) void print_usage() {
void print_usage(void) {
printf("Usage:\n");
printf(" fastsync [options] <source> <destination>\n");
printf(" fastsync [options] --source-dir <src> --dest-dir <dst>\n");
@@ -37,6 +38,7 @@ static __attribute__((unused)) void print_usage() {
printf(" -f Enable sendfile (TCP only, not with -c or -s)\n");
printf(" -v, --verbose Enable debug logging\n");
printf(" -M, --preserve Preserve file metadata\n");
printf(" --chmod <changes> Modify transferred permissions (rsync syntax)\n");
printf(" --chunk-size <n> Chunk size in bytes (default: %d)\n", DEFAULT_CHUNK_SIZE);
printf(" --source-dir <path> Source directory\n");
printf(" --dest-dir <path> Destination directory\n");
+6
View File
@@ -0,0 +1,6 @@
#ifndef USAGE_H
#define USAGE_H
void print_usage(void);
#endif
+142
View File
@@ -0,0 +1,142 @@
#include "receiver.h"
#include "chunk.h"
#include "log.h"
#include "protocol.h"
#include "utils.h"
#include <stdlib.h>
#include <sys/stat.h>
static bool receiver_process_chunk(Chunk* chunk, const ReceiverSink* sink) {
if (!chunk || !sink || !sink->store_file)
return false;
for (int i = 0; i < chunk->element_count; i++) {
File* file = chunk->items[i];
if (!file) {
chunk_destroy(chunk);
return false;
}
chunk->items[i] = NULL;
if (!sink->store_file(file, sink->context)) {
chunk_destroy(chunk);
return false;
}
}
chunk_destroy(chunk);
return true;
}
static bool receiver_process_batch(Config* config, int file_descriptor) {
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;
if (!receive_n_data(file_descriptor, &check_size, sizeof(check_size)) ||
!receive_n_data(file_descriptor, &check_mtime, sizeof(check_mtime))) {
free(check_path);
return false;
}
if (!utils_valid_batch_path(check_path)) {
free(check_path);
send_status(file_descriptor, STATUS_ERROR);
return false;
}
if (check_size > MAX_RECEIVE_FILE_SIZE) {
free(check_path);
send_status(file_descriptor, STATUS_ERROR);
return false;
}
char* full_path = path_cat(config->receive_root_directory, check_path);
if (!full_path) {
free(check_path);
send_status(file_descriptor, STATUS_ERROR);
return false;
}
struct stat st;
bool has_old = file_stat_secure(full_path, &st);
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 receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink) {
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;
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(file_descriptor, config, &skipped);
if (!skipped && (!file || !sink->store_file(file, sink->context)))
goto receive_error;
} else if (status == STATUS_CHUNK) {
Chunk* chunk = receive_chunk_data(file_descriptor, config);
if (!chunk || !receiver_process_chunk(chunk, sink))
goto receive_error;
} else if (status == STATUS_CHECK_BATCH) {
if (!receiver_process_batch(config, file_descriptor))
return -1;
goto next;
} else {
File* file = file_receive(config, file_descriptor);
if (!file) {
log_message(LOG_LEVEL_ERROR, "Failed to receive file");
goto receive_error;
}
if (!sink->store_file(file, sink->context))
goto receive_error;
}
next:
if (!receive_status(file_descriptor, &status))
goto receive_error;
}
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");
goto receive_error;
}
if (sink->send_success && !send_status(file_descriptor, STATUS_OK))
return -1;
return 0;
receive_error:
if (sink->send_error)
send_status(file_descriptor, STATUS_ERROR);
return -1;
}
static bool receiver_save_file(File* file, void* context) {
Config* config = context;
bool success =
!config->save_to_disk || file_save_to_disk(config->receive_root_directory, file, config);
file_destroy(file);
return success;
}
int receiver_receive_files(Config* config, int file_descriptor) {
ReceiverSink sink = {receiver_save_file, config, true, true};
return receiver_process(config, file_descriptor, &sink);
}
+19
View File
@@ -0,0 +1,19 @@
#ifndef RECEIVER_H
#define RECEIVER_H
#include "config.h"
#include "file.h"
typedef bool (*ReceiverFileSink)(File* file, void* context);
typedef struct {
ReceiverFileSink store_file;
void* context;
bool send_error;
bool send_success;
} ReceiverSink;
int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink);
int receiver_receive_files(Config* config, int file_descriptor);
#endif
+143 -63
View File
@@ -1,54 +1,102 @@
#include "array_list.h"
#include "chunk.h"
#include "config.h"
#include "data.h"
#include "chunk.h"
#include "file.h"
#include "log.h"
#include "multiprocessing.h"
#include "protocol.h"
#include "queue.h"
#include "receiver.h"
#include "transport_tcp.h"
#include "transport_tls.h"
#include "unistd.h"
#include "utils.h"
#include <fcntl.h>
#include <limits.h>
#include <signal.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <limits.h>
#include <fcntl.h>
#include <sys/stat.h>
#include <unistd.h>
#include <openssl/x509.h>
static char* authorized_root;
static int authorized_root_fd = -1;
static bool allow_delete;
static bool allow_unauthenticated;
static const char* required_client_cn;
static bool tls_client_identity_allowed(SSL* ssl) {
if (!ssl || !required_client_cn)
return false;
X509* certificate = SSL_get1_peer_certificate(ssl);
if (!certificate)
return false;
char common_name[256];
int length = X509_NAME_get_text_by_NID(X509_get_subject_name(certificate), NID_commonName,
common_name, sizeof(common_name));
size_t required_length = strlen(required_client_cn);
bool allowed = length >= 0 && (size_t)length == required_length &&
required_length < sizeof(common_name) &&
memcmp(common_name, required_client_cn, required_length) == 0;
X509_free(certificate);
return allowed;
}
static void release_authorization(void) {
file_set_authorized_root(-1, NULL);
utils_set_authorized_root_fd(-1);
if (authorized_root_fd >= 0)
close(authorized_root_fd);
authorized_root_fd = -1;
free(authorized_root);
authorized_root = NULL;
}
static bool path_is_within(const char* root, const char* path) {
size_t n = strlen(root);
return strncmp(root, path, n) == 0 && (path[n] == '\0' || path[n] == '/');
}
static bool valid_batch_path(const char* path) {
return path && path[0] != '\0' && path[0] != '/' && !has_path_traversal(path) &&
strchr(path, '\0') == path + strlen(path);
}
static bool __attribute__((unused)) configure_authorization(const char* root) {
char resolved[PATH_MAX];
if (!root || !realpath(root, resolved))
if (!root) {
file_set_authorized_root(-1, NULL);
utils_set_authorized_root(-1, NULL);
return false;
}
int root_fd = open(root, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (root_fd < 0) {
file_set_authorized_root(-1, NULL);
utils_set_authorized_root(-1, NULL);
return false;
}
char fd_path[64];
int fd_path_length = snprintf(fd_path, sizeof(fd_path), "/proc/self/fd/%d", root_fd);
if (fd_path_length < 0 || (size_t)fd_path_length >= sizeof(fd_path) ||
!realpath(fd_path, resolved)) {
close(root_fd);
file_set_authorized_root(-1, NULL);
utils_set_authorized_root(-1, NULL);
return false;
}
authorized_root = str_dup(resolved);
if (!authorized_root)
if (!authorized_root) {
close(root_fd);
file_set_authorized_root(-1, NULL);
utils_set_authorized_root(-1, NULL);
return false;
authorized_root_fd = open(resolved, O_RDONLY | O_DIRECTORY | O_CLOEXEC);
if (authorized_root_fd < 0) {
}
authorized_root_fd = root_fd;
if (!file_set_authorized_root(authorized_root_fd, authorized_root) ||
!utils_set_authorized_root(authorized_root_fd, authorized_root)) {
file_set_authorized_root(-1, NULL);
utils_set_authorized_root(-1, NULL);
close(authorized_root_fd);
authorized_root_fd = -1;
free(authorized_root);
authorized_root = NULL;
return false;
}
file_set_authorized_root(authorized_root_fd, authorized_root);
utils_set_authorized_root_fd(authorized_root_fd);
return true;
}
@@ -113,14 +161,19 @@ int receive_files(Config* config, int fd) {
free(check_path);
return -1;
}
if (!valid_batch_path(check_path)) {
if (!utils_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;
char* full_path = path_cat(config->receive_root_directory, check_path);
if (!full_path) {
free(check_path);
send_status(fd, STATUS_ERROR);
return -1;
}
bool has_old = full_path && file_stat_secure(full_path, &st);
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);
@@ -153,8 +206,9 @@ int receive_files(Config* config, int fd) {
}
if (status == STATUS_MANIFEST) {
if (receive_manifest(fd, config, &status) != 0)
if (receive_manifest(fd, config, &status) != 0) {
return -1;
}
}
if (status != STATUS_FINISHED) {
log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status");
@@ -167,39 +221,59 @@ int receive_files(Config* config, int fd) {
void handler(int file_descriptor) {
SSL* ssl = io_get_ssl();
ProtocolSession session;
protocol_session_init(&session, file_descriptor, file_descriptor);
protocol_session_set_ssl(&session, ssl);
protocol_session_bind(&session);
Config* config = config_receive(file_descriptor);
if (config == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to receive config");
close(file_descriptor);
protocol_session_unbind();
return;
}
if (!authorized_root) {
log_message(LOG_LEVEL_ERROR, "No server-side destination root configured");
config_delete(config);
close(file_descriptor);
protocol_session_unbind();
return;
}
char resolved_destination[PATH_MAX];
char* canonical_destination = realpath(config->receive_root_directory, NULL);
const char* destination =
canonical_destination ? canonical_destination : config->receive_root_directory;
if (has_path_traversal(destination) || !path_is_within(authorized_root, destination)) {
log_message(LOG_LEVEL_ERROR, "Rejected destination outside authorized root");
free(canonical_destination);
if (!allow_unauthenticated && ssl == NULL) {
log_message(LOG_LEVEL_ERROR, "Rejected unauthenticated plaintext connection");
config_delete(config);
close(file_descriptor);
protocol_session_unbind();
return;
}
if (ssl && required_client_cn && !tls_client_identity_allowed(ssl)) {
log_message(LOG_LEVEL_ERROR, "Rejected TLS client with unauthorized identity");
config_delete(config);
close(file_descriptor);
return;
}
if (canonical_destination)
snprintf(resolved_destination, sizeof(resolved_destination), "%s", canonical_destination);
else
snprintf(resolved_destination, sizeof(resolved_destination), "%s", destination);
free(canonical_destination);
free(config->receive_root_directory);
config->receive_root_directory = str_dup(resolved_destination);
char* destination = config->receive_root_directory;
char* joined_destination = NULL;
if (destination && destination[0] != '/')
joined_destination = path_cat(authorized_root, destination);
if (joined_destination)
destination = joined_destination;
if (!destination || has_path_traversal(destination) ||
!path_is_within(authorized_root, destination)) {
log_message(LOG_LEVEL_ERROR, "Rejected destination outside authorized root");
free(joined_destination);
config_delete(config);
close(file_descriptor);
return;
}
if (joined_destination) {
free(config->receive_root_directory);
config->receive_root_directory = joined_destination;
}
if (!config->receive_root_directory) {
config_delete(config);
close(file_descriptor);
protocol_session_unbind();
return;
}
config->use_delete = config->use_delete && allow_delete;
@@ -208,6 +282,7 @@ void handler(int file_descriptor) {
if (q == NULL) {
config_delete(config);
close(file_descriptor);
protocol_session_unbind();
return;
}
PipelineContextReceiver* context =
@@ -216,16 +291,17 @@ void handler(int file_descriptor) {
queue_destroy(q);
config_delete(config);
close(file_descriptor);
protocol_session_unbind();
return;
}
context->session.total_allocated_bytes = session.total_allocated_bytes;
thrd_t receiver, writer;
bool receiver_created = false;
bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success;
bool writer_created = false;
receiver_created = (thrd_create(&receiver, receive_thread, context) == thrd_success);
if (receiver_created)
writer_created = (thrd_create(&writer, write_thread, context) == thrd_success);
writer_created = thrd_create(&writer, write_thread, context) == thrd_success;
if (!receiver_created || !writer_created) {
perror("Error creating Threads");
log_perror("Error creating Threads");
if (receiver_created) {
mtx_lock(&context->mutex);
atomic_store(&context->cancelled, true);
@@ -240,22 +316,23 @@ void handler(int file_descriptor) {
if (writer_created)
thrd_join(writer, NULL);
pipeline_context_receiver_destroy(context);
protocol_session_unbind();
return;
}
int receiver_result;
int writer_result;
thrd_join(receiver, &receiver_result);
thrd_join(writer, &writer_result);
if (receiver_result == thrd_success && writer_result == thrd_success)
send_status(file_descriptor, STATUS_OK);
else
send_status(file_descriptor, STATUS_ERROR);
send_status(file_descriptor, receiver_result == thrd_success && writer_result == thrd_success
? STATUS_OK
: STATUS_ERROR);
pipeline_context_receiver_destroy(context);
} else {
if (receive_files(config, file_descriptor) != 0)
if (receiver_receive_files(config, file_descriptor) != 0)
log_message(LOG_LEVEL_ERROR, "Transfer failed");
config_delete(config);
}
protocol_session_unbind();
close(file_descriptor);
}
@@ -264,16 +341,14 @@ static Server* g_server = NULL;
static void cleanup(int sig) {
(void)sig;
if (g_server) {
if (g_server)
server_delete(&g_server);
}
_exit(0);
}
static void print_server_usage(void) {
printf("FastSync Server\n");
printf("Usage: fastsync-server [options]\n");
printf("\n");
printf("Usage: fastsync-server [options]\n\n");
printf("Options:\n");
printf(" --stdio Run in stdio mode (SSH transport)\n");
printf(" -p <port> TCP port (default: 8080, range: 1-65535)\n");
@@ -281,17 +356,17 @@ static void print_server_usage(void) {
printf(" --cert <path> TLS certificate file (PEM)\n");
printf(" --key <path> TLS private key file (PEM)\n");
printf(" --ca <path> TLS CA certificate file (PEM)\n");
printf(" --client-cn <name> Required TLS client certificate CN\n");
printf(" --destination-root <path> Authorized destination root (default: .)\n");
printf(" --allow-delete Permit manifest deletion\n");
printf(" --allow-unauthenticated Allow plaintext/anonymous network clients\n");
printf(" -v, --verbose Enable debug logging\n");
printf(" --help Show this help\n");
}
int main(int argc, char* argv[]) {
bool use_tls = false;
char* tls_cert = NULL;
char* tls_key = NULL;
char* tls_ca = NULL;
char *tls_cert = NULL, *tls_key = NULL, *tls_ca = NULL;
int port = 8080;
const char* destination_root = ".";
bool stdio_mode = false;
@@ -313,10 +388,14 @@ int main(int argc, char* argv[]) {
tls_key = argv[++i];
} else if (strcmp(argv[i], "--ca") == 0 && i + 1 < argc) {
tls_ca = argv[++i];
} else if (strcmp(argv[i], "--client-cn") == 0 && i + 1 < argc) {
required_client_cn = argv[++i];
} else if (strcmp(argv[i], "--destination-root") == 0 && i + 1 < argc) {
destination_root = argv[++i];
} else if (strcmp(argv[i], "--allow-delete") == 0) {
allow_delete = true;
} else if (strcmp(argv[i], "--allow-unauthenticated") == 0) {
allow_unauthenticated = true;
} else if (strcmp(argv[i], "-p") == 0 && i + 1 < argc) {
char* end;
long p = strtol(argv[++i], &end, 10);
@@ -331,11 +410,8 @@ int main(int argc, char* argv[]) {
return 1;
}
}
if (tls_ca && !use_tls) {
if (tls_ca && !use_tls)
log_message(LOG_LEVEL_WARNING, "--ca has no effect without --tls");
}
signal(SIGINT, cleanup);
signal(SIGTERM, cleanup);
if (!configure_authorization(destination_root)) {
@@ -343,35 +419,39 @@ int main(int argc, char* argv[]) {
return 1;
}
if (stdio_mode) {
/* SSH authenticates the stdio transport outside of FastSync. */
allow_unauthenticated = true;
io_set_fds(STDIN_FILENO, STDOUT_FILENO);
handler(STDIN_FILENO);
file_set_authorized_root(-1, NULL);
utils_set_authorized_root_fd(-1);
close(authorized_root_fd);
free(authorized_root);
release_authorization();
return 0;
}
g_server = server_create(port);
if (g_server == NULL) {
if (!g_server) {
log_message(LOG_LEVEL_ERROR, "Failed to create server");
release_authorization();
return 1;
}
if (use_tls) {
if (!tls_cert || !tls_key) {
fprintf(stderr, "Error: --tls requires --cert and --key\n");
if (!tls_cert || !tls_key || !tls_ca || !required_client_cn) {
fprintf(stderr, "Error: --tls requires --cert, --key, --ca, and --client-cn\n");
server_delete(&g_server);
release_authorization();
return 1;
}
tls_global_init();
if (!server_create_tls(g_server, tls_cert, tls_key, tls_ca)) {
log_message(LOG_LEVEL_ERROR, "Failed to set up TLS");
server_delete(&g_server);
release_authorization();
return 1;
}
server_listen_tls(g_server, handler);
} else {
server_listen(g_server, handler);
}
server_delete(&g_server);
release_authorization();
return 0;
}
#endif /* !FASTSYNC_SERVER_AS_LIB */
#endif
+5 -4
View File
@@ -1,3 +1,4 @@
#include "log.h"
#include "array_list.h"
#include <stdio.h>
#include <stdlib.h>
@@ -6,7 +7,7 @@
ArrayList* array_list_create(void (*item_destroyer)(void* item)) {
ArrayList* list = (ArrayList*)malloc(sizeof(ArrayList));
if (list == NULL) {
perror("ERROR: Could not allocate memory for array list struct");
log_perror("ERROR: Could not allocate memory for array list struct");
return NULL;
}
@@ -34,7 +35,7 @@ void array_list_delete(ArrayList* array_list) {
free(array_list);
}
bool array_list_extend(ArrayList* array_list) {
static bool array_list_extend(ArrayList* array_list) {
if (array_list == NULL)
return false;
int new_capacity = array_list->capacity * 2;
@@ -42,7 +43,7 @@ bool array_list_extend(ArrayList* array_list) {
new_capacity = INITIAL_ARRAY_SIZE;
void* new_items = realloc(array_list->items, new_capacity * sizeof(void*));
if (new_items == NULL) {
perror("ERROR: Could not reallocate memory for array list items");
log_perror("ERROR: Could not reallocate memory for array list items");
return false;
}
array_list->items = new_items;
@@ -68,7 +69,7 @@ void** array_list_to_array(const ArrayList* array_list) {
}
void** array = malloc(array_list->size * sizeof(void*));
if (array == NULL) {
perror("Could not malloc space for array from array list!");
log_perror("Could not malloc space for array from array list!");
return NULL;
}
memcpy(array, array_list->items, array_list->size * sizeof(void*));
-1
View File
@@ -14,7 +14,6 @@ typedef struct ArrayList {
ArrayList* array_list_create(void (*item_destroyer)(void* item));
void array_list_delete(ArrayList* array_list);
bool array_list_extend(ArrayList* array_list);
bool array_list_add(ArrayList* array_list, void* item);
void** array_list_to_array(const ArrayList* array_list);
+90
View File
@@ -0,0 +1,90 @@
#include "chmod.h"
#include <stddef.h>
#include <string.h>
static bool parse_clause(mode_t* mode, const char* begin, const char* end) {
const char* p = begin;
unsigned who = 0;
while (p < end && strchr("ugoa", *p)) {
if (*p == 'a')
who = 7;
else
who |= *p == 'u' ? 1U : (*p == 'g' ? 2U : 4U);
p++;
}
if (who == 0)
who = 7;
if (p == end || (*p != '+' && *p != '-' && *p != '='))
return false;
char operation = *p++;
mode_t bits = 0;
while (p < end) {
mode_t bit;
switch (*p++) {
case 'r':
bit = 4;
break;
case 'w':
bit = 2;
break;
case 'x':
bit = 1;
break;
default:
return false;
}
bits |= bit;
}
for (unsigned class_index = 0; class_index < 3; class_index++) {
unsigned class_bit = 1U << class_index;
if (!(who & class_bit))
continue;
mode_t shift = (mode_t)((2U - class_index) * 3U);
mode_t mask = (mode_t)(7U << shift);
mode_t class_bits = (mode_t)(bits << shift);
if (operation == '+')
*mode |= class_bits;
else if (operation == '-')
*mode &= ~class_bits;
else
*mode = (*mode & ~mask) | class_bits;
}
return true;
}
bool chmod_apply(mode_t mode, const char* spec, mode_t* result) {
if (!spec || !*spec || !result)
return false;
bool numeric = true;
size_t length = strlen(spec);
if (length > 4)
numeric = false;
for (size_t i = 0; i < length && numeric; i++)
numeric = spec[i] >= '0' && spec[i] <= '7';
if (numeric) {
if (length == 0 || length > 4)
return false;
mode_t parsed = 0;
for (size_t i = 0; i < length; i++)
parsed = (mode_t)((parsed << 3) | (spec[i] - '0'));
*result = parsed;
return true;
}
mode_t changed = mode;
const char* begin = spec;
while (*begin) {
const char* end = strchr(begin, ',');
if (!end)
end = begin + strlen(begin);
if (!parse_clause(&changed, begin, end))
return false;
if (*end == '\0')
break;
begin = end + 1;
if (!*begin)
return false;
}
*result = changed;
return true;
}
+10
View File
@@ -0,0 +1,10 @@
#ifndef CHMOD_H
#define CHMOD_H
#include <stdbool.h>
#include <sys/stat.h>
/* Apply the supported rsync --chmod syntax to a permission mode. */
bool chmod_apply(mode_t mode, const char* spec, mode_t* result);
#endif
+115 -19
View File
@@ -1,4 +1,6 @@
#include <stddef.h>
#include <stdint.h>
#include <limits.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
@@ -11,21 +13,33 @@
#include "log.h"
#include "metadata.h"
#include "protocol.h"
#include "utils.h"
/* Maximum individual file data size within a chunk (64 MB) */
#define MAX_FILE_DATA_SIZE (64ULL * 1024 * 1024)
#define MAX_FILES_PER_CHUNK 65536U
Chunk* chunk_create(File** items, int element_count) {
if (element_count < 0 || (element_count > 0 && items == NULL))
return NULL;
Chunk* chunk = (Chunk*)malloc(sizeof(Chunk));
if (chunk == NULL) {
perror("ERROR: Could not allocate memory for chunk structure");
log_perror("ERROR: Could not allocate memory for chunk structure");
return NULL;
}
chunk->items = (File**)malloc(element_count * sizeof(File*));
if (chunk->items == NULL) {
free(chunk);
return NULL;
if (element_count == 0) {
chunk->items = NULL;
} else {
if ((size_t)element_count > SIZE_MAX / sizeof(File*)) {
free(chunk);
return NULL;
}
chunk->items = (File**)malloc((size_t)element_count * sizeof(File*));
if (chunk->items == NULL) {
free(chunk);
return NULL;
}
}
for (int i = 0; i < element_count; i++) {
@@ -50,15 +64,37 @@ void chunk_destroy(void* item) {
}
static unsigned long long per_file_serialize_size(File* file, bool use_metadata) {
return sizeof(size_t) + strlen(file->path) +
(use_metadata ? sizeof(int) + (file->metadata ? FILE_METADATA_WIRE_SIZE : 0) : 0) +
sizeof(size_t) + file->data->size;
unsigned long long size = sizeof(size_t);
size_t path_len = strlen(file->path);
unsigned long long metadata_size =
use_metadata ? sizeof(int) + (file->metadata ? FILE_METADATA_WIRE_SIZE : 0) : 0;
if ((unsigned long long)path_len > ULLONG_MAX - size)
return 0;
size += path_len;
if (metadata_size > ULLONG_MAX - size)
return 0;
size += metadata_size;
if (sizeof(size_t) > ULLONG_MAX - size)
return 0;
size += sizeof(size_t);
if ((unsigned long long)file->data->size > ULLONG_MAX - size)
return 0;
return size + file->data->size;
}
Data* chunk_serialize(Chunk* chunk, bool use_metadata) {
if (!chunk || chunk->element_count < 0 || (chunk->element_count > 0 && chunk->items == NULL))
return NULL;
unsigned long long data_size = 0;
for (int i = 0; i < chunk->element_count; i++) {
data_size += per_file_serialize_size(chunk->items[i], use_metadata);
if (!chunk->items[i] || !chunk->items[i]->path || !chunk->items[i]->data ||
(chunk->items[i]->data->size > 0 && !chunk->items[i]->data->data) ||
chunk->items[i]->path[0] == '\0' || has_path_traversal(chunk->items[i]->path))
return NULL;
unsigned long long file_size = per_file_serialize_size(chunk->items[i], use_metadata);
if (file_size == 0 || file_size > ULLONG_MAX - data_size || data_size + file_size > SIZE_MAX)
return NULL;
data_size += file_size;
}
Data* data = data_create_empty(data_size);
if (data == NULL) {
@@ -87,11 +123,20 @@ Data* chunk_serialize(Chunk* chunk, bool use_metadata) {
}
Chunk* chunk_deserialize(Data* data, bool use_metadata) {
if (!data || (!data->data && data->size != 0))
return NULL;
ArrayList* files = array_list_create(file_destroy);
if (files == NULL)
return NULL;
char* data_pointer = data->data;
size_t remaining_size = data->size;
while (remaining_size > 0) {
if ((unsigned int)files->size >= MAX_FILES_PER_CHUNK) {
log_message(LOG_LEVEL_ERROR, "Chunk contains too many files");
array_list_delete(files);
return NULL;
}
if (remaining_size < sizeof(size_t)) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for path length");
array_list_delete(files);
@@ -103,48 +148,77 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
data_pointer += sizeof(size_t);
remaining_size -= sizeof(size_t);
if (remaining_size < path_len) {
if (path_len > SIZE_MAX - 1 || remaining_size < path_len) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for path");
array_list_delete(files);
return NULL;
}
if (path_len == SIZE_MAX) {
array_list_delete(files);
return NULL;
}
char* path = malloc(path_len + 1);
if (path == NULL) {
perror("Could not allocate memory for file path");
log_perror("Could not allocate memory for file path");
array_list_delete(files);
return NULL;
}
memcpy(path, data_pointer, path_len);
path[path_len] = '\0';
if (memchr(path, '\0', path_len) != NULL) {
free(path);
array_list_delete(files);
return NULL;
}
data_pointer += path_len;
remaining_size -= path_len;
if (path_len == 0 || has_path_traversal(path)) {
free(path);
array_list_delete(files);
return NULL;
}
File* file = file_create(path);
free(path);
if (file == NULL) {
array_list_delete(files);
return NULL;
}
if (use_metadata) {
if (remaining_size < sizeof(int)) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for metadata");
file_destroy(file);
array_list_delete(files);
return NULL;
}
// Peek at present flag to determine total size needed before reading
int present_flag;
memcpy(&present_flag, data_pointer, sizeof(int));
if (present_flag && remaining_size < sizeof(int) + FILE_METADATA_WIRE_SIZE) {
if ((present_flag != 0 && present_flag != 1) ||
(present_flag == 1 && remaining_size < sizeof(int) + FILE_METADATA_WIRE_SIZE)) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for metadata body");
file_destroy(file);
array_list_delete(files);
return NULL;
}
file->metadata = metadata_from_buf(&data_pointer);
remaining_size -= sizeof(int);
if (file->metadata)
if (present_flag == 1) {
if (file->metadata == NULL) {
file_destroy(file);
array_list_delete(files);
return NULL;
}
remaining_size -= FILE_METADATA_WIRE_SIZE;
}
}
if (remaining_size < sizeof(size_t)) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for data size");
file_destroy(file);
array_list_delete(files);
return NULL;
}
@@ -156,6 +230,7 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
if (remaining_size < file_data_size) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for file content");
file_destroy(file);
array_list_delete(files);
return NULL;
}
@@ -164,29 +239,50 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
if (file_data_size > MAX_FILE_DATA_SIZE) {
log_message(LOG_LEVEL_ERROR, "File data size %zu exceeds maximum %llu", file_data_size,
(unsigned long long)MAX_FILE_DATA_SIZE);
file_destroy(file);
array_list_delete(files);
return NULL;
}
void* file_data = malloc(file_data_size);
size_t allocation_size = file_data_size > 0 ? file_data_size : 1;
void* file_data = malloc(allocation_size);
if (file_data == NULL) {
perror("Could not allocate memory for file data");
log_perror("Could not allocate memory for file data");
file_destroy(file);
array_list_delete(files);
return NULL;
}
memcpy(file_data, data_pointer, file_data_size);
Data* replacement = data_create(file_data, file_data_size);
if (replacement == NULL) {
file_destroy(file);
array_list_delete(files);
return NULL;
}
data_destroy(file->data);
file->data = data_create(file_data, file_data_size);
file->data = replacement;
data_pointer += file_data_size;
remaining_size -= file_data_size;
array_list_add(files, file);
if (!array_list_add(files, file)) {
file_destroy(file);
array_list_delete(files);
return NULL;
}
}
File** file_array = (File**)array_list_to_array(files);
if (files->size > 0 && file_array == NULL) {
array_list_delete(files);
return NULL;
}
Chunk* chunk = chunk_create(file_array, files->size);
free(file_array);
if (chunk == NULL) {
array_list_delete(files);
return NULL;
}
files->item_destroyer = NULL;
array_list_delete(files);
@@ -207,14 +303,14 @@ Data* chunk_compress(Chunk* chunk, int compression_level, bool use_metadata) {
}
Chunk* receive_chunk_data(int fd, const Config* config) {
Data* chunk_data = receive_data(fd);
Data* chunk_data = receive_data_limited(fd, MAX_CHUNK_SIZE);
if (chunk_data == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to receive chunk data");
return NULL;
}
Data* data_to_process = chunk_data;
if (config->use_compression) {
data_to_process = data_decompress(chunk_data);
data_to_process = data_decompress_limited(chunk_data, MAX_CHUNK_SIZE);
data_destroy(chunk_data);
if (data_to_process == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to decompress chunk");
+21 -9
View File
@@ -2,6 +2,8 @@
#include "data.h"
#include "log.h"
#include <stdlib.h>
#include <limits.h>
#include <stdint.h>
#include <string.h>
#include <strings.h>
#include <zstd.h>
@@ -69,7 +71,10 @@ Data* data_compress(Data* data_to_compress, int compression_level) {
return compressed_data;
}
Data* data_decompress(Data* compressed_data) {
Data* data_decompress_limited(Data* compressed_data, size_t maximum_size) {
if (!compressed_data || (!compressed_data->data && compressed_data->size != 0) ||
maximum_size == 0)
return NULL;
log_message(LOG_LEVEL_DEBUG, "Start to decompress data");
unsigned long long dst_size =
ZSTD_getFrameContentSize(compressed_data->data, compressed_data->size);
@@ -82,15 +87,16 @@ Data* data_decompress(Data* compressed_data) {
// ZSTD_CONTENTSIZE_UNKNOWN (~2^64) can cause massive allocation;
// fall back to a conservative estimate (3x compressed size) when unknown.
if (dst_size == ZSTD_CONTENTSIZE_UNKNOWN) {
if (compressed_data->size > ULLONG_MAX / 3)
return NULL;
dst_size = compressed_data->size * 3;
if (dst_size < INITIAL_DECOMPRESS_BUF_SIZE)
dst_size = INITIAL_DECOMPRESS_BUF_SIZE;
if (dst_size > MAX_DECOMPRESSED_SIZE)
dst_size = MAX_DECOMPRESSED_SIZE;
}
if (dst_size > MAX_DECOMPRESSED_SIZE) {
log_message(LOG_LEVEL_ERROR, "Declared decompressed size exceeds %llu bytes",
(unsigned long long)MAX_DECOMPRESSED_SIZE);
unsigned long long hard_limit =
maximum_size < MAX_DECOMPRESSED_SIZE ? maximum_size : MAX_DECOMPRESSED_SIZE;
if (dst_size > hard_limit) {
log_message(LOG_LEVEL_ERROR, "Declared decompressed size exceeds %llu bytes", hard_limit);
return NULL;
}
@@ -101,6 +107,8 @@ Data* data_decompress(Data* compressed_data) {
}
size_t buf_size = (dst_size > 0) ? (size_t)dst_size : INITIAL_DECOMPRESS_BUF_SIZE;
if (buf_size > maximum_size)
buf_size = maximum_size;
Data* uncompressed_data = data_create_empty(buf_size);
if (!uncompressed_data) {
log_message(LOG_LEVEL_ERROR, "Failed to allocate decompression buffer");
@@ -121,7 +129,7 @@ Data* data_decompress(Data* compressed_data) {
return NULL;
}
if (ret > 0 && output.pos == output.size) {
if (buf_size >= MAX_DECOMPRESSED_SIZE) {
if (buf_size >= hard_limit || buf_size > SIZE_MAX / 2) {
log_message(LOG_LEVEL_ERROR, "Decompressed data exceeds %llu bytes",
(unsigned long long)MAX_DECOMPRESSED_SIZE);
ZSTD_freeDCtx(dctx);
@@ -129,8 +137,8 @@ Data* data_decompress(Data* compressed_data) {
return NULL;
}
buf_size *= 2;
if (buf_size > MAX_DECOMPRESSED_SIZE)
buf_size = MAX_DECOMPRESSED_SIZE;
if (buf_size > hard_limit)
buf_size = (size_t)hard_limit;
void* new_data = realloc(uncompressed_data->data, buf_size);
if (!new_data) {
log_message(LOG_LEVEL_ERROR, "Failed to grow decompression buffer");
@@ -150,3 +158,7 @@ Data* data_decompress(Data* compressed_data) {
log_message(LOG_LEVEL_DEBUG, "Decompressed data successfully");
return uncompressed_data;
}
Data* data_decompress(Data* compressed_data) {
return data_decompress_limited(compressed_data, MAX_DECOMPRESSED_SIZE);
}
+1
View File
@@ -6,6 +6,7 @@
Data* data_compress(Data* data_to_compress, int compression_level);
Data* data_decompress(Data* compressed_data);
Data* data_decompress_limited(Data* compressed_data, size_t maximum_size);
bool compression_should_skip(const char* path);
#endif
+191 -251
View File
@@ -1,4 +1,5 @@
#include "config.h"
#include "chmod.h"
#include "delta.h"
#include "log.h"
#include "protocol.h"
@@ -98,6 +99,46 @@ static void config_set_defaults(Config* config) {
config->server_mode = false;
config->checksum = false;
config->compress_choice = NULL;
config->chmod_spec = NULL;
}
static bool valid_wire_bool(int value) {
return value == 0 || value == 1;
}
static bool receive_wire_bool(int fd, bool* value) {
int wire_value;
if (!receive_int(fd, &wire_value) || !valid_wire_bool(wire_value))
return false;
*value = wire_value != 0;
return true;
}
static bool validate_received_config(const Config* config) {
return valid_wire_bool(config->save_to_disk) && valid_wire_bool(config->use_multithreading) &&
valid_wire_bool(config->use_chunk_serialization) &&
valid_wire_bool(config->use_compression) && valid_wire_bool(config->use_metadata) &&
valid_wire_bool(config->use_sendfile) && valid_wire_bool(config->use_delete) &&
valid_wire_bool(config->use_incremental) && valid_wire_bool(config->use_delta) &&
valid_wire_bool(config->backup) && valid_wire_bool(config->follow_symlinks) &&
valid_wire_bool(config->copy_links) && valid_wire_bool(config->safe_links) &&
valid_wire_bool(config->copy_unsafe_links) &&
valid_wire_bool(config->preserve_hard_links) && valid_wire_bool(config->preserve_acls) &&
valid_wire_bool(config->preserve_xattrs) && valid_wire_bool(config->preserve_devices) &&
valid_wire_bool(config->preserve_sparse) && valid_wire_bool(config->update) &&
valid_wire_bool(config->inplace) && valid_wire_bool(config->append) &&
valid_wire_bool(config->append_verify) && valid_wire_bool(config->delete_excluded) &&
valid_wire_bool(config->delete_after) && valid_wire_bool(config->relative) &&
valid_wire_bool(config->prune_empty_dirs) && valid_wire_bool(config->partial) &&
valid_wire_bool(config->delete_before) && valid_wire_bool(config->checksum) &&
(!config->use_compression ||
(config->compression_level >= 1 && config->compression_level <= 22)) &&
config->chunk_size > 0 && config->chunk_size <= MAX_CHUNK_SIZE &&
config->delta_block_size >= DELTA_BLOCK_SIZE_MIN &&
config->delta_block_size <= DELTA_BLOCK_SIZE_MAX &&
config->delta_max_file_size <= DELTA_MAX_FILE_SIZE && config->max_delete >= 0 &&
(!config->chmod_spec || !*config->chmod_spec ||
chmod_apply(0, config->chmod_spec, &(mode_t){0}));
}
Config* config_create(void) {
@@ -108,7 +149,7 @@ Config* config_create(void) {
return config;
}
bool is_remote_dest(const char* s) {
bool config_is_remote_dest(const char* s) {
if (s == NULL)
return false;
const char* colon = strchr(s, ':');
@@ -124,7 +165,7 @@ bool is_remote_dest(const char* s) {
}
void config_parse_ssh_dest(Config* config) {
if (!is_remote_dest(config->receive_root_directory))
if (!config_is_remote_dest(config->receive_root_directory))
return;
config->transport = TRANSPORT_SSH;
config->ssh_destination = str_dup(config->receive_root_directory);
@@ -137,6 +178,10 @@ void config_parse_ssh_dest(Config* config) {
void config_delete(Config* config) {
if (config == NULL)
return;
if (config->log_file) {
fclose(config->log_file);
config->log_file = NULL;
}
free(config->version);
free(config->send_directory);
free(config->receive_root_directory);
@@ -167,108 +212,141 @@ void config_delete(Config* config) {
free(config->bind_address);
free(config->daemon_config);
free(config->compress_choice);
free(config->chmod_spec);
if (config->filters) {
array_list_delete(config->filters);
}
free(config);
}
/* Wire format order (must match config_receive and be updated when PROTOCOL_VERSION bumps):
* version, send_directory, receive_root_directory, save_to_disk, use_multithreading,
* 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,
* preserve_hard_links, preserve_acls, preserve_xattrs, preserve_devices, preserve_sparse,
* update, inplace, append, append_verify, delete_excluded, delete_after, max_delete, relative,
* prune_empty_dirs, temp_dir, partial, partial_dir, suffix, delete_before, checksum,
* compress_choice, status
*/
/* Each helper is deliberately ordered to match the wire format. Keep the
* helper call order in config_send and config_receive unchanged when adding
* fields. */
static bool send_core_fields(int fd, const Config* c) {
return send_str(fd, c->version) && send_str(fd, c->send_directory) &&
send_str(fd, c->receive_root_directory) && send_int(fd, c->save_to_disk) &&
send_int(fd, c->use_multithreading) && send_int(fd, c->use_chunk_serialization) &&
send_int(fd, c->use_compression) && send_int(fd, c->use_metadata) &&
send_int(fd, c->compression_level) &&
send_n_data(fd, &c->chunk_size, sizeof(c->chunk_size)) && send_int(fd, c->use_sendfile);
}
static bool send_delta_fields(int fd, const Config* c) {
return send_int(fd, c->use_delete) && send_int(fd, c->use_incremental) &&
send_int(fd, c->use_delta) &&
send_n_data(fd, &c->delta_block_size, sizeof(c->delta_block_size)) &&
send_n_data(fd, &c->delta_max_file_size, sizeof(unsigned long long));
}
static bool send_file_options(int fd, const Config* c) {
return send_int(fd, c->backup) && send_str(fd, c->backup_dir ? c->backup_dir : "") &&
send_int(fd, c->follow_symlinks) && send_int(fd, c->copy_links) &&
send_int(fd, c->safe_links) && send_int(fd, c->copy_unsafe_links) &&
send_int(fd, c->preserve_hard_links) && send_int(fd, c->preserve_acls) &&
send_int(fd, c->preserve_xattrs) && send_int(fd, c->preserve_devices) &&
send_int(fd, c->preserve_sparse);
}
static bool send_selection_options(int fd, const Config* c) {
return send_int(fd, c->update) && send_int(fd, c->inplace) && send_int(fd, c->append) &&
send_int(fd, c->append_verify) && send_int(fd, c->delete_excluded) &&
send_int(fd, c->delete_after) && send_n_data(fd, &c->max_delete, sizeof(c->max_delete)) &&
send_int(fd, c->relative) && send_int(fd, c->prune_empty_dirs);
}
static bool send_resume_options(int fd, const Config* c) {
return send_str(fd, c->temp_dir ? c->temp_dir : "") && send_int(fd, c->partial) &&
send_str(fd, c->partial_dir ? c->partial_dir : "") &&
send_str(fd, c->suffix ? c->suffix : "") && send_int(fd, c->delete_before) &&
send_int(fd, c->checksum) && send_str(fd, c->compress_choice ? c->compress_choice : "") &&
send_str(fd, c->chmod_spec ? c->chmod_spec : "");
}
static bool receive_core_fields(int fd, Config* c) {
int value;
c->send_directory = receive_str(fd);
c->receive_root_directory = receive_str(fd);
if (!c->send_directory || !c->receive_root_directory)
return false;
if (!receive_wire_bool(fd, &c->save_to_disk) || !receive_wire_bool(fd, &c->use_multithreading) ||
!receive_wire_bool(fd, &c->use_chunk_serialization) ||
!receive_wire_bool(fd, &c->use_compression) || !receive_wire_bool(fd, &c->use_metadata))
return false;
if (!receive_int(fd, &value))
return false;
c->compression_level = value;
if (!receive_n_data(fd, &c->chunk_size, sizeof(c->chunk_size)))
return false;
if (!receive_wire_bool(fd, &c->use_sendfile))
return false;
return true;
}
static bool receive_delta_fields(int fd, Config* c) {
if (!receive_wire_bool(fd, &c->use_delete))
return false;
if (!receive_wire_bool(fd, &c->use_incremental))
return false;
if (!receive_wire_bool(fd, &c->use_delta))
return false;
return receive_n_data(fd, &c->delta_block_size, sizeof(c->delta_block_size)) &&
receive_n_data(fd, &c->delta_max_file_size, sizeof(unsigned long long));
}
static bool receive_file_options(int fd, Config* c) {
if (!receive_wire_bool(fd, &c->backup))
return false;
c->backup_dir = receive_str(fd);
if (!c->backup_dir)
return false;
bool* flags[] = {&c->follow_symlinks, &c->copy_links, &c->safe_links,
&c->copy_unsafe_links, &c->preserve_hard_links, &c->preserve_acls,
&c->preserve_xattrs, &c->preserve_devices, &c->preserve_sparse};
for (size_t i = 0; i < sizeof(flags) / sizeof(flags[0]); i++) {
if (!receive_wire_bool(fd, flags[i]))
return false;
}
return true;
}
static bool receive_selection_options(int fd, Config* c) {
bool* flags[] = {&c->update, &c->inplace, &c->append,
&c->append_verify, &c->delete_excluded, &c->delete_after};
for (size_t i = 0; i < sizeof(flags) / sizeof(flags[0]); i++) {
if (!receive_wire_bool(fd, flags[i]))
return false;
}
if (!receive_n_data(fd, &c->max_delete, sizeof(c->max_delete)))
return false;
if (!receive_wire_bool(fd, &c->relative))
return false;
if (!receive_wire_bool(fd, &c->prune_empty_dirs))
return false;
return true;
}
static bool receive_resume_options(int fd, Config* c) {
c->temp_dir = receive_str(fd);
if (!c->temp_dir || !receive_wire_bool(fd, &c->partial))
return false;
c->partial_dir = receive_str(fd);
c->suffix = c->partial_dir ? receive_str(fd) : NULL;
if (!c->partial_dir || !c->suffix || !receive_wire_bool(fd, &c->delete_before))
return false;
if (!receive_wire_bool(fd, &c->checksum))
return false;
c->compress_choice = receive_str(fd);
if (!c->compress_choice)
return false;
c->chmod_spec = receive_str(fd);
return c->chmod_spec != NULL;
}
bool config_send(int file_descriptor, const Config* config) {
if (!send_str(file_descriptor, config->version))
return false;
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 : ""))
if (!send_core_fields(file_descriptor, config) || !send_delta_fields(file_descriptor, config) ||
!send_file_options(file_descriptor, config) ||
!send_selection_options(file_descriptor, config) ||
!send_resume_options(file_descriptor, config))
return false;
Status status;
if (!receive_status(file_descriptor, &status))
@@ -280,159 +358,25 @@ bool config_send(int file_descriptor, const Config* config) {
return true;
}
/* Wire format order: see the comment above config_send. */
Config* config_receive(int file_descriptor) {
Config* config = (Config*)malloc(sizeof(Config));
if (config == NULL)
Config* config = config_create();
if (!config)
return NULL;
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 (!config->version)
goto error;
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;
goto error;
}
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;
if (!receive_int(file_descriptor, &tmp))
goto error;
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)
if (!receive_core_fields(file_descriptor, config) ||
!receive_delta_fields(file_descriptor, config) ||
!receive_file_options(file_descriptor, config) ||
!receive_selection_options(file_descriptor, config) ||
!receive_resume_options(file_descriptor, config))
goto error;
if (config->compress_choice[0] != '\0' && strcmp(config->compress_choice, "zstd") != 0 &&
strcmp(config->compress_choice, "none") != 0) {
@@ -440,20 +384,16 @@ Config* config_receive(int file_descriptor) {
send_status(file_descriptor, STATUS_ERROR);
goto error;
}
if (!validate_received_config(config)) {
fprintf(stderr, "Invalid configuration received from client\n");
send_status(file_descriptor, STATUS_ERROR);
goto error;
}
if (!send_status(file_descriptor, STATUS_OK))
goto error;
return config;
error:
free(config->version);
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);
config_delete(config);
return NULL;
}
+3 -2
View File
@@ -126,16 +126,17 @@ typedef struct Config {
// PR #184: Compression algorithm negotiation
char* compress_choice;
char* chmod_spec;
} Config;
#define PROTOCOL_VERSION "2.2.0"
#define PROTOCOL_VERSION "2.3.0"
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
Config* config_create(void);
void config_delete(Config* config);
bool config_send(int file_descriptor, const Config* config);
Config* config_receive(int file_descriptor);
bool is_remote_dest(const char* s);
bool config_is_remote_dest(const char* s);
void config_parse_ssh_dest(Config* config);
#endif
+4
View File
@@ -21,6 +21,7 @@ Data* data_create_reserve(size_t size) {
}
d->data = NULL;
d->size = size;
d->protocol_charge = 0;
return d;
}
@@ -33,12 +34,15 @@ Data* data_create(void* data, size_t data_size) {
}
new_data->data = data;
new_data->size = data_size;
new_data->protocol_charge = 0;
return new_data;
}
void data_destroy(Data* data) {
if (data == NULL)
return;
if (data->protocol_charge != 0)
protocol_release_memory(data->protocol_charge);
free(data->data);
free(data);
}
+3
View File
@@ -6,11 +6,14 @@
typedef struct {
void* data;
size_t size;
/* Non-zero only for a buffer charged to the protocol connection budget. */
size_t protocol_charge;
} Data;
Data* data_create_empty(size_t data_size);
Data* data_create_reserve(size_t size);
Data* data_create(void* data, size_t data_size);
void data_destroy(Data* data);
void protocol_release_memory(size_t charge);
#endif
+85 -56
View File
@@ -1,6 +1,7 @@
#include "delta.h"
#include "log.h"
#include <stdint.h>
#include <limits.h>
#include <stdlib.h>
#include <string.h>
@@ -36,6 +37,10 @@ DeltaSignature* delta_signature_create(const void* old_file_data, uint64_t old_f
if (old_file_data == NULL || old_file_size == 0 || block_size == 0)
return NULL;
if (old_file_size > DELTA_MAX_FILE_SIZE || block_size > DELTA_BLOCK_SIZE_MAX ||
old_file_size > UINT32_MAX * (uint64_t)block_size)
return NULL;
uint32_t block_count = (uint32_t)((old_file_size + block_size - 1) / block_size);
DeltaSignature* sig = malloc(sizeof(DeltaSignature));
@@ -45,7 +50,11 @@ DeltaSignature* delta_signature_create(const void* old_file_data, uint64_t old_f
sig->file_size = old_file_size;
sig->block_size = block_size;
sig->block_count = block_count;
sig->blocks = malloc(block_count * sizeof(DeltaBlockSig));
if (block_count == 0) {
free(sig);
return NULL;
}
sig->blocks = malloc((size_t)block_count * sizeof(DeltaBlockSig));
if (!sig->blocks) {
free(sig);
return NULL;
@@ -67,8 +76,11 @@ Data* delta_signature_serialize(const DeltaSignature* sig) {
if (!sig)
return NULL;
uint64_t total = sizeof(uint64_t) + sizeof(uint32_t) + sizeof(uint32_t) +
(uint64_t)sig->block_count * (sizeof(uint32_t) + sizeof(uint32_t));
uint64_t block_bytes = (uint64_t)sig->block_count * (sizeof(uint32_t) + sizeof(uint32_t));
uint64_t total = sizeof(uint64_t) + sizeof(uint32_t) + sizeof(uint32_t) + block_bytes;
if (block_bytes > UINT64_MAX - (sizeof(uint64_t) + sizeof(uint32_t) + sizeof(uint32_t)) ||
total > SIZE_MAX)
return NULL;
uint8_t* buf = malloc((size_t)total);
if (!buf)
@@ -118,6 +130,13 @@ DeltaSignature* delta_signature_deserialize(const Data* data) {
return NULL;
}
if (sig->block_size == 0 || sig->block_size > DELTA_BLOCK_SIZE_MAX ||
sig->file_size > DELTA_MAX_FILE_SIZE || sig->file_size == 0 ||
(sig->file_size + sig->block_size - 1) / sig->block_size != sig->block_count) {
free(sig);
return NULL;
}
uint64_t expected = sizeof(uint64_t) + sizeof(uint32_t) + sizeof(uint32_t) +
(uint64_t)sig->block_count * (sizeof(uint32_t) + sizeof(uint32_t));
if (data->size < expected) {
@@ -156,8 +175,10 @@ void delta_signature_destroy(DeltaSignature* sig) {
static bool ensure_capacity(DeltaInstruction** instrs, uint32_t* capacity, uint32_t count) {
if (count < *capacity)
return true;
if (*capacity > MAX_DELTA_INSTRUCTIONS / 2)
return false;
uint32_t new_cap = *capacity * 2;
DeltaInstruction* tmp = realloc(*instrs, new_cap * sizeof(DeltaInstruction));
DeltaInstruction* tmp = realloc(*instrs, (size_t)new_cap * sizeof(DeltaInstruction));
if (!tmp)
return false;
*instrs = tmp;
@@ -169,6 +190,8 @@ static bool flush_literal(DeltaInstruction** instrs, uint32_t* capacity, uint32_
const uint8_t* data, uint64_t start, uint64_t end) {
if (start >= end)
return true;
if (end - start > UINT32_MAX || *count >= MAX_DELTA_INSTRUCTIONS)
return false;
uint32_t lit_len = (uint32_t)(end - start);
if (!ensure_capacity(instrs, capacity, *count))
return false;
@@ -183,16 +206,26 @@ static bool flush_literal(DeltaInstruction** instrs, uint32_t* capacity, uint32_
return true;
}
static void free_instructions(DeltaInstruction* instrs, uint32_t count) {
if (!instrs)
return;
for (uint32_t i = 0; i < count; i++)
if (instrs[i].type == DELTA_INSTR_LITERAL)
free(instrs[i].literal.data);
free(instrs);
}
Delta* delta_compute(const void* new_file_data, uint64_t new_file_size, const DeltaSignature* sig,
uint32_t block_size) {
if (!new_file_data || !sig || new_file_size == 0 || block_size == 0)
if (!new_file_data || !sig || !sig->blocks || new_file_size == 0 || block_size == 0 ||
block_size > DELTA_BLOCK_SIZE_MAX || sig->block_size != block_size)
return NULL;
const uint8_t* new_data = (const uint8_t*)new_file_data;
uint32_t capacity = 64;
uint32_t count = 0;
DeltaInstruction* instrs = malloc(capacity * sizeof(DeltaInstruction));
DeltaInstruction* instrs = malloc((size_t)capacity * sizeof(DeltaInstruction));
if (!instrs)
return NULL;
@@ -236,14 +269,14 @@ Delta* delta_compute(const void* new_file_data, uint64_t new_file_size, const De
if (xxh == sig->blocks[j].xxhash) {
if (has_literal) {
if (!flush_literal(&instrs, &capacity, &count, new_data, literal_start, i)) {
free(instrs);
free_instructions(instrs, count);
return NULL;
}
has_literal = false;
}
if (!ensure_capacity(&instrs, &capacity, count)) {
free(instrs);
free_instructions(instrs, count);
return NULL;
}
instrs[count].type = DELTA_INSTR_BLOCK_MATCH;
@@ -271,18 +304,14 @@ Delta* delta_compute(const void* new_file_data, uint64_t new_file_size, const De
if (has_literal) {
if (!flush_literal(&instrs, &capacity, &count, new_data, literal_start, new_file_size)) {
free(instrs);
free_instructions(instrs, count);
return NULL;
}
}
Delta* delta = malloc(sizeof(Delta));
if (!delta) {
for (uint32_t k = 0; k < count; k++) {
if (instrs[k].type == DELTA_INSTR_LITERAL)
free(instrs[k].literal.data);
}
free(instrs);
free_instructions(instrs, count);
return NULL;
}
@@ -292,11 +321,24 @@ Delta* delta_compute(const void* new_file_data, uint64_t new_file_size, const De
delta->delta_size = 0;
for (uint32_t k = 0; k < count; k++) {
if (delta->delta_size == UINT64_MAX) {
delta_destroy(delta);
return NULL;
}
delta->delta_size += 1;
if (instrs[k].type == DELTA_INSTR_BLOCK_MATCH) {
if (delta->delta_size > UINT64_MAX - sizeof(uint32_t) * 3) {
delta_destroy(delta);
return NULL;
}
delta->delta_size += sizeof(uint32_t) * 3;
} else {
delta->delta_size += sizeof(uint32_t) + instrs[k].literal.length;
uint64_t extra = sizeof(uint32_t) + instrs[k].literal.length;
if (delta->delta_size > UINT64_MAX - extra) {
delta_destroy(delta);
return NULL;
}
delta->delta_size += extra;
}
}
@@ -307,7 +349,12 @@ Data* delta_serialize(const Delta* delta) {
if (!delta)
return NULL;
uint64_t total = sizeof(uint64_t) + sizeof(uint32_t) + delta->delta_size;
if (delta->instruction_count > 0 && !delta->instructions)
return NULL;
uint64_t header_size = sizeof(uint64_t) + sizeof(uint32_t);
if (delta->delta_size > UINT64_MAX - header_size || header_size + delta->delta_size > SIZE_MAX)
return NULL;
uint64_t total = header_size + delta->delta_size;
uint8_t* buf = malloc((size_t)total);
if (!buf)
return NULL;
@@ -365,8 +412,10 @@ Delta* delta_deserialize(const Data* data) {
return NULL;
}
delta->instructions = malloc(delta->instruction_count * sizeof(DeltaInstruction));
if (!delta->instructions) {
delta->instructions = delta->instruction_count == 0
? NULL
: malloc((size_t)delta->instruction_count * sizeof(DeltaInstruction));
if (delta->instruction_count > 0 && !delta->instructions) {
free(delta);
return NULL;
}
@@ -375,11 +424,7 @@ Delta* delta_deserialize(const Data* data) {
for (uint32_t i = 0; i < delta->instruction_count; i++) {
if (pos >= data->size) {
for (uint32_t k = 0; k < i; k++) {
if (delta->instructions[k].type == DELTA_INSTR_LITERAL)
free(delta->instructions[k].literal.data);
}
free(delta->instructions);
free_instructions(delta->instructions, i);
free(delta);
return NULL;
}
@@ -391,12 +436,8 @@ Delta* delta_deserialize(const Data* data) {
delta->delta_size += 1;
if (type == DELTA_OP_BLOCK_MATCH) {
if (pos + sizeof(uint32_t) * 3 > data->size) {
for (uint32_t k = 0; k < i; k++) {
if (delta->instructions[k].type == DELTA_INSTR_LITERAL)
free(delta->instructions[k].literal.data);
}
free(delta->instructions);
if (data->size - pos < sizeof(uint32_t) * 3) {
free_instructions(delta->instructions, i);
free(delta);
return NULL;
}
@@ -409,12 +450,8 @@ Delta* delta_deserialize(const Data* data) {
pos += sizeof(uint32_t);
delta->delta_size += sizeof(uint32_t) * 3;
} else if (type == DELTA_OP_LITERAL) {
if (pos + sizeof(uint32_t) > data->size) {
for (uint32_t k = 0; k < i; k++) {
if (delta->instructions[k].type == DELTA_INSTR_LITERAL)
free(delta->instructions[k].literal.data);
}
free(delta->instructions);
if (data->size - pos < sizeof(uint32_t)) {
free_instructions(delta->instructions, i);
free(delta);
return NULL;
}
@@ -423,23 +460,15 @@ Delta* delta_deserialize(const Data* data) {
pos += sizeof(uint32_t);
uint32_t lit_len = delta->instructions[i].literal.length;
if (pos + lit_len > data->size) {
for (uint32_t k = 0; k < i; k++) {
if (delta->instructions[k].type == DELTA_INSTR_LITERAL)
free(delta->instructions[k].literal.data);
}
free(delta->instructions);
if (lit_len > data->size - pos) {
free_instructions(delta->instructions, i);
free(delta);
return NULL;
}
delta->instructions[i].literal.data = malloc(lit_len);
delta->instructions[i].literal.data = malloc(lit_len ? lit_len : 1);
if (!delta->instructions[i].literal.data) {
log_message(LOG_LEVEL_ERROR, "Failed to allocate %u bytes for literal data", lit_len);
for (uint32_t k = 0; k < i; k++) {
if (delta->instructions[k].type == DELTA_INSTR_LITERAL)
free(delta->instructions[k].literal.data);
}
free(delta->instructions);
free_instructions(delta->instructions, i);
free(delta);
return NULL;
}
@@ -447,11 +476,7 @@ Delta* delta_deserialize(const Data* data) {
pos += lit_len;
delta->delta_size += sizeof(uint32_t) + lit_len;
} else {
for (uint32_t k = 0; k < i; k++) {
if (delta->instructions[k].type == DELTA_INSTR_LITERAL)
free(delta->instructions[k].literal.data);
}
free(delta->instructions);
free_instructions(delta->instructions, i);
free(delta);
return NULL;
}
@@ -463,10 +488,11 @@ Delta* delta_deserialize(const Data* data) {
void* delta_apply(const void* old_data, uint64_t old_size, const Delta* delta,
uint32_t block_size) {
if (!old_data || !delta || (delta->new_file_size > 0 && delta->instructions == NULL) ||
(delta->instruction_count > 0 && block_size == 0))
(delta->instruction_count > 0 && block_size == 0) ||
delta->new_file_size > DELTA_MAX_FILE_SIZE || delta->new_file_size > SIZE_MAX)
return NULL;
void* output = malloc((size_t)delta->new_file_size);
void* output = malloc(delta->new_file_size ? (size_t)delta->new_file_size : 1);
if (!output)
return NULL;
@@ -491,7 +517,7 @@ void* delta_apply(const void* old_data, uint64_t old_size, const Delta* delta,
}
memcpy(out + out_pos, old + src_offset, len);
out_pos += len;
} else {
} else if (delta->instructions[i].type == DELTA_INSTR_LITERAL) {
uint32_t len = delta->instructions[i].literal.length;
if (out_pos > delta->new_file_size || (uint64_t)len > delta->new_file_size - out_pos) {
free(output);
@@ -499,6 +525,9 @@ void* delta_apply(const void* old_data, uint64_t old_size, const Delta* delta,
}
memcpy(out + out_pos, delta->instructions[i].literal.data, len);
out_pos += len;
} else {
free(output);
return NULL;
}
}
@@ -534,7 +563,7 @@ bool delta_should_attempt(uint64_t old_size, uint64_t new_size, uint64_t max_fil
}
bool delta_is_worthwhile(const Delta* delta, uint64_t new_file_size) {
if (!delta || delta->instruction_count == 0)
if (!delta || delta->instruction_count == 0 || new_file_size == 0)
return false;
bool has_match = false;
+122 -841
View File
File diff suppressed because it is too large Load Diff
+20 -29
View File
@@ -1,45 +1,36 @@
#ifndef FILE_H
#define FILE_H
#include "config.h"
#include "data.h"
#include "file_send.h"
#include "file_receive.h"
#include "file_types.h"
#include <stdbool.h>
#include <stdint.h>
#include <sys/stat.h>
typedef enum { FILE_TYPE_REGULAR, FILE_TYPE_SYMLINK, FILE_TYPE_DIR } FileType;
typedef struct {
mode_t mode;
uid_t uid;
gid_t gid;
time_t mtime_sec;
long mtime_nsec;
} FileMetadata;
typedef struct {
char* path;
Data* data;
FileMetadata* metadata;
bool skip;
} File;
/* File/FileMetadata lifecycle, local disk helpers, and secure filesystem
primitives shared by the send/receive pipelines. */
File* file_create(const char* path);
void file_destroy(void* item);
bool file_load_data(File* file);
bool file_checksum(File* file, uint64_t* checksum);
File* file_receive(const Config* config, int file_descriptor);
bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata,
int compression_level, bool send_path);
bool file_send_sendfile(File* file, int file_descriptor, bool use_metadata, int compression_level,
bool send_path);
size_t file_content_to_buffer(File* file);
FileMetadata* file_metadata_create(const struct stat* stats);
void file_metadata_destroy(void* metadata);
bool to_disk(const char* path, const void* data, unsigned long long data_size, bool inplace,
bool sparse);
bool file_save_to_disk(const char* root_directory, const File* file, const Config* config);
void file_set_authorized_root(int fd, const char* canonical_path);
File* receive_incremental_check(int fd, const Config* config, bool* skipped);
int receive_manifest(int fd, const Config* config, int* next_status);
bool file_write_to_disk(const char* path, const void* data, unsigned long long data_size,
bool inplace, bool sparse);
/* A configured fd without a canonical identity deliberately rejects paths. */
bool file_set_authorized_root(int fd, const char* canonical_path);
/* Secure path/filesystem primitives (symlink-safe, O_NOFOLLOW, root-confined). */
bool file_path_exists_secure(const char* path);
bool file_stat_secure(const char* path, struct stat* st);
int file_open_secure_parent(const char* path, char** leaf_out, bool create_dirs);
bool file_ensure_directory_secure(const char* path);
bool file_rename_secure(const char* old_path, const char* new_path);
bool file_to_disk_secure(const char* path, const void* data, unsigned long long data_size,
bool inplace, bool sparse, const FileMetadata* metadata);
#endif
+585
View File
@@ -0,0 +1,585 @@
#include <errno.h>
#include <fcntl.h>
#include <libgen.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/stat.h>
#include <unistd.h>
#include "array_list.h"
#include "chmod.h"
#include "compression.h"
#include "config.h"
#include "data.h"
#include "delta.h"
#include "file.h"
#include "log.h"
#include "metadata.h"
#include "protocol.h"
#include "utils.h"
#define MAX_SERVER_DELETE_COUNT 100000U
#define MAX_FILE_DATA_SIZE MAX_RECEIVE_FILE_SIZE
bool file_save_to_disk(const char* root_directory, const File* file, const Config* config) {
bool backup_enabled = config && config->backup;
bool inplace = config && config->inplace;
bool sparse = config && config->preserve_sparse;
const char* backup_suffix = (config && config->suffix) ? config->suffix : "~";
const char* backup_dir = (config && config->backup_dir) ? config->backup_dir : NULL;
const char* partial_dir = (config && config->partial_dir) ? config->partial_dir : NULL;
char *confined_backup = NULL, *confined_partial = NULL, *disk_path = NULL;
char *backup_path = NULL, *parent_copy = NULL;
if (!file || !file->path || !file->data || (file->data->size != 0 && !file->data->data) ||
has_path_traversal(file->path) ||
(backup_enabled &&
(!backup_suffix || backup_suffix[0] == '\0' || strchr(backup_suffix, '/') != NULL ||
strcmp(backup_suffix, ".") == 0 || strcmp(backup_suffix, "..") == 0))) {
log_message(LOG_LEVEL_ERROR, "Invalid file or path received");
return false;
}
/* These options arrive from the client. They are names below the server
root, never independent filesystem roots. */
if ((backup_dir && (backup_dir[0] == '/' || has_path_traversal(backup_dir))) ||
(partial_dir && (partial_dir[0] == '/' || has_path_traversal(partial_dir))))
return false;
if (backup_dir && !(confined_backup = path_cat(root_directory, backup_dir)))
return false;
if (partial_dir && !(confined_partial = path_cat(root_directory, partial_dir))) {
free(confined_backup);
return false;
}
const char* actual_root =
(partial_dir && config && config->partial) ? confined_partial : root_directory;
disk_path = path_cat(actual_root, file->path);
if (disk_path == NULL) {
free(confined_backup);
free(confined_partial);
return false;
}
/* --update is receiver-side policy: never replace a newer destination. */
if (config && config->update) {
struct stat destination_stat;
if (file_stat_secure(disk_path, &destination_stat) && file->metadata &&
destination_stat.st_mtime > file->metadata->mtime_sec) {
free(confined_backup);
free(confined_partial);
free(disk_path);
return true;
}
}
if (backup_enabled) {
struct stat backup_stat;
if (file_stat_secure(disk_path, &backup_stat)) {
if (backup_dir) {
backup_path = path_cat(confined_backup, file->path);
} else {
size_t path_len = strlen(disk_path);
size_t suffix_len = strlen(backup_suffix);
if (path_len > SIZE_MAX - suffix_len - 1)
goto fail;
backup_path = malloc(path_len + suffix_len + 1);
if (backup_path) {
memcpy(backup_path, disk_path, path_len);
memcpy(backup_path + path_len, backup_suffix, suffix_len + 1);
}
}
if (!backup_path)
goto fail;
parent_copy = str_dup(backup_path);
if (!parent_copy || !file_ensure_directory_secure(dirname(parent_copy)))
goto fail;
free(parent_copy);
parent_copy = NULL;
if (!file_rename_secure(disk_path, backup_path))
goto fail;
free(backup_path);
backup_path = NULL;
}
}
FileMetadata adjusted_metadata;
const FileMetadata* metadata = file->metadata;
if (metadata && config && config->chmod_spec && *config->chmod_spec) {
adjusted_metadata = *metadata;
if (!chmod_apply(adjusted_metadata.mode, config->chmod_spec, &adjusted_metadata.mode))
goto fail;
metadata = &adjusted_metadata;
}
bool ok =
file_to_disk_secure(disk_path, file->data->data, file->data->size, inplace, sparse, metadata);
free(parent_copy);
free(backup_path);
free(confined_backup);
free(confined_partial);
free(disk_path);
return ok;
fail:
free(parent_copy);
free(backup_path);
free(confined_backup);
free(confined_partial);
free(disk_path);
return false;
}
static File* receive_delta_file(int fd, const Config* config, const char* check_path,
void* old_data, unsigned long long old_size, bool* failed) {
if (!old_data)
return NULL;
DeltaSignature* sig = delta_signature_create(old_data, old_size, config->delta_block_size);
if (!sig) {
free(old_data);
*failed = true;
return NULL;
}
Data* sig_data = delta_signature_serialize(sig);
if (!sig_data) {
delta_signature_destroy(sig);
free(old_data);
*failed = true;
return NULL;
}
bool sig_sent = send_status(fd, STATUS_DELTA_SIGNATURE) && send_data(fd, sig_data);
data_destroy(sig_data);
if (!sig_sent) {
delta_signature_destroy(sig);
free(old_data);
*failed = true;
return NULL;
}
Status resp;
if (!receive_status(fd, &resp)) {
delta_signature_destroy(sig);
free(old_data);
*failed = true;
return NULL;
}
if (resp == STATUS_DELTA_DATA) {
Data* delta_data = receive_data_limited(fd, MAX_RECEIVE_FILE_SIZE);
if (!delta_data) {
delta_signature_destroy(sig);
free(old_data);
*failed = true;
return NULL;
}
Data* raw_delta = delta_data;
if (config->use_compression) {
raw_delta = data_decompress_limited(delta_data, MAX_RECEIVE_FILE_SIZE);
data_destroy(delta_data);
if (!raw_delta) {
free(old_data);
delta_signature_destroy(sig);
*failed = true;
return NULL;
}
}
Delta* delta = delta_deserialize(raw_delta);
data_destroy(raw_delta);
if (!delta) {
free(old_data);
delta_signature_destroy(sig);
*failed = true;
return NULL;
}
uint64_t new_size = delta->new_file_size;
if (new_size > MAX_RECEIVE_FILE_SIZE || new_size > SIZE_MAX) {
delta_destroy(delta);
free(old_data);
delta_signature_destroy(sig);
send_status(fd, STATUS_ERROR);
return NULL;
}
void* new_data = delta_apply(old_data, old_size, delta, config->delta_block_size);
delta_destroy(delta);
if (!new_data) {
free(old_data);
delta_signature_destroy(sig);
*failed = true;
return NULL;
}
File* file = file_create(check_path);
if (!file) {
free(new_data);
free(old_data);
delta_signature_destroy(sig);
*failed = true;
return NULL;
}
if (config->use_metadata) {
int meta_ok = 1;
file->metadata = metadata_receive(fd, &meta_ok);
if (!meta_ok) {
file_destroy(file);
free(new_data);
free(old_data);
delta_signature_destroy(sig);
*failed = true;
return NULL;
}
}
Data* replacement = data_create(new_data, (size_t)new_size);
if (replacement == NULL) {
file_destroy(file);
free(old_data);
delta_signature_destroy(sig);
send_status(fd, STATUS_ERROR);
return NULL;
}
data_destroy(file->data);
file->data = replacement;
free(old_data);
delta_signature_destroy(sig);
return file;
}
if (resp == STATUS_NEXT) {
delta_signature_destroy(sig);
free(old_data);
File* file = file_create(check_path);
if (!file) {
*failed = true;
return NULL;
}
if (config->use_metadata) {
int meta_ok = 1;
file->metadata = metadata_receive(fd, &meta_ok);
if (!meta_ok) {
file_destroy(file);
*failed = true;
return NULL;
}
}
Data* file_data = receive_data_limited(fd, MAX_RECEIVE_FILE_SIZE);
if (file_data == NULL) {
file_destroy(file);
*failed = true;
return NULL;
}
if (config->use_compression) {
Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_FILE_SIZE);
data_destroy(file_data);
if (uncompressed == NULL) {
file_destroy(file);
*failed = true;
return NULL;
}
if (uncompressed->size > MAX_FILE_DATA_SIZE) {
data_destroy(uncompressed);
file_destroy(file);
send_status(fd, STATUS_ERROR);
return NULL;
}
file_data = uncompressed;
}
data_destroy(file->data);
file->data = file_data;
return file;
}
delta_signature_destroy(sig);
free(old_data);
send_status(fd, STATUS_ERROR);
*failed = true;
return NULL;
}
File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
if (!config || !skipped) {
send_status(fd, STATUS_ERROR);
return NULL;
}
*skipped = false;
char* check_path = receive_str(fd);
if (check_path == NULL) {
return NULL;
}
unsigned long long check_size;
long long check_mtime;
uint64_t check_checksum = 0;
if (!receive_n_data(fd, &check_size, sizeof(check_size)) ||
!receive_n_data(fd, &check_mtime, sizeof(check_mtime))) {
free(check_path);
return NULL;
}
if (config->checksum && !receive_n_data(fd, &check_checksum, sizeof(check_checksum))) {
free(check_path);
return NULL;
}
if (check_size > MAX_RECEIVE_FILE_SIZE) {
free(check_path);
send_status(fd, STATUS_ERROR);
return NULL;
}
if (has_path_traversal(check_path)) {
log_message(LOG_LEVEL_ERROR, "Path traversal detected: %s", check_path);
free(check_path);
return NULL;
}
char* full_path = path_cat(config->receive_root_directory, check_path);
if (!full_path) {
free(check_path);
send_status(fd, STATUS_ERROR);
return NULL;
}
struct stat st;
bool has_old_file = false;
int old_fd = -1;
char* leaf = NULL;
int parent_fd = file_open_secure_parent(full_path, &leaf, false);
if (parent_fd >= 0) {
old_fd = openat(parent_fd, leaf, O_RDONLY | O_CLOEXEC | O_NOFOLLOW);
free(leaf);
close(parent_fd);
has_old_file = old_fd >= 0 && fstat(old_fd, &st) == 0 && S_ISREG(st.st_mode);
}
unsigned long long old_size = has_old_file ? (unsigned long long)st.st_size : 0;
void* old_data = NULL;
if (has_old_file && old_size > 0 && old_size <= MAX_RECEIVE_FILE_SIZE && old_size <= SIZE_MAX) {
old_data = malloc((size_t)old_size);
if (old_data) {
size_t got = 0;
while (got < (size_t)old_size) {
ssize_t n = read(old_fd, (char*)old_data + got, (size_t)old_size - got);
if (n <= 0) {
free(old_data);
old_data = NULL;
break;
}
got += (size_t)n;
}
}
}
if (old_fd >= 0) {
close(old_fd);
}
bool match = has_old_file && (unsigned long long)st.st_size == check_size;
if (match && config->checksum) {
uint64_t old_checksum = old_size == 0 ? delta_xxhash64("", 0) : 0;
if (old_data)
old_checksum = delta_xxhash64(old_data, (size_t)old_size);
match = (old_size == 0 || old_data) && old_checksum == check_checksum;
free(old_data);
old_data = NULL;
} else if (match) {
match = (long long)st.st_mtime == check_mtime;
}
if (match) {
free(old_data);
if (!send_status(fd, STATUS_OK)) {
free(full_path);
free(check_path);
return NULL;
}
free(full_path);
free(check_path);
*skipped = true;
return NULL;
}
bool try_delta = config->use_delta && has_old_file && old_data != NULL &&
delta_should_attempt(old_size, check_size, config->delta_max_file_size);
if (try_delta) {
bool delta_failed = false;
File* delta_file =
receive_delta_file(fd, config, check_path, old_data, old_size, &delta_failed);
old_data = NULL; /* receive_delta_file consumes the snapshot on every path */
if (delta_file) {
free(full_path);
free(check_path);
return delta_file;
}
if (delta_failed) {
free(full_path);
free(check_path);
return NULL;
}
free(old_data);
old_data = NULL;
try_delta = false;
}
if (!try_delta) {
if (!send_status(fd, STATUS_NEXT)) {
free(full_path);
free(check_path);
return NULL;
}
}
File* file = file_create(check_path);
free(check_path);
free(full_path);
if (file == NULL) {
return NULL;
}
if (config->use_metadata) {
int meta_ok = 1;
file->metadata = metadata_receive(fd, &meta_ok);
if (!meta_ok) {
file_destroy(file);
return NULL;
}
}
Data* file_data = receive_data_limited(fd, MAX_RECEIVE_FILE_SIZE);
if (file_data == NULL) {
file_destroy(file);
return NULL;
}
if (config->use_compression) {
Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_FILE_SIZE);
data_destroy(file_data);
if (uncompressed == NULL) {
file_destroy(file);
return NULL;
}
if (uncompressed->size > MAX_FILE_DATA_SIZE) {
data_destroy(uncompressed);
file_destroy(file);
send_status(fd, STATUS_ERROR);
return NULL;
}
file_data = uncompressed;
}
data_destroy(file->data);
file->data = file_data;
return file;
}
File* file_receive(const Config* config, int file_descriptor) {
char* path = receive_str(file_descriptor);
if (path == NULL)
return NULL;
if (path[0] == '\0' || has_path_traversal(path)) {
log_message(LOG_LEVEL_ERROR, "Invalid received file path: %s", path);
free(path);
return NULL;
}
File* file = file_create(path);
free(path);
if (file == NULL)
return NULL;
if (config->use_metadata) {
int meta_ok = 1;
file->metadata = metadata_receive(file_descriptor, &meta_ok);
if (!meta_ok) {
file_destroy(file);
return NULL;
}
}
Data* file_data = receive_data_limited(file_descriptor, MAX_RECEIVE_FILE_SIZE);
if (file_data == NULL) {
file_destroy(file);
return NULL;
}
if (config->use_compression && !compression_should_skip(file->path)) {
Data* file_data_uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_FILE_SIZE);
data_destroy(file_data);
if (file_data_uncompressed == NULL) {
file_destroy(file);
return NULL;
}
if (file_data_uncompressed->size > MAX_FILE_DATA_SIZE) {
data_destroy(file_data_uncompressed);
file_destroy(file);
return NULL;
}
file_data = file_data_uncompressed;
}
data_destroy(file->data);
file->data = file_data;
return file;
}
int receive_manifest(int fd, const Config* config, int* next_status) {
if (!config) {
send_status(fd, STATUS_ERROR);
return -1;
}
int received_status = STATUS_ERROR;
int* status_out = next_status ? next_status : &received_status;
int count;
if (!receive_int(fd, &count)) {
send_status(fd, STATUS_ERROR);
return -1;
}
if (count < 0 || count > MAX_MANIFEST_ENTRIES) {
send_status(fd, STATUS_ERROR);
return -1;
}
ArrayList* manifest = array_list_create(free);
if (!manifest) {
send_status(fd, STATUS_ERROR);
return -1;
}
size_t manifest_bytes = 0;
for (int i = 0; i < count; i++) {
char* s = receive_str(fd);
size_t entry_size = s ? strlen(s) : 0;
if (!s || s[0] == '\0' || s[0] == '/' || has_path_traversal(s) ||
entry_size > MAX_MANIFEST_BYTES - manifest_bytes ||
(manifest_bytes += entry_size) > MAX_MANIFEST_BYTES || !array_list_add(manifest, s)) {
free(s);
array_list_delete(manifest);
send_status(fd, STATUS_ERROR);
return -1;
}
}
if (!receive_status(fd, status_out)) {
array_list_delete(manifest);
send_status(fd, STATUS_ERROR);
return -1;
}
/* Deletion is a commit operation: never perform it until the sender has
completed the manifest frame successfully. */
if (*status_out != STATUS_FINISHED || !config->use_delete) {
array_list_delete(manifest);
if (*status_out != STATUS_FINISHED)
send_status(fd, STATUS_ERROR);
return *status_out == STATUS_FINISHED ? 0 : -1;
}
fprintf(stderr, "Deleting files not in manifest...\n");
bool deletion_ok =
delete_extras_limited(config->receive_root_directory, manifest, MAX_SERVER_DELETE_COUNT);
array_list_delete(manifest);
if (!deletion_ok)
send_status(fd, STATUS_ERROR);
return deletion_ok ? 0 : -1;
}
+15
View File
@@ -0,0 +1,15 @@
#ifndef FILE_RECEIVE_H
#define FILE_RECEIVE_H
#include "config.h"
#include "file_types.h"
#include <stdbool.h>
/* Server-side file receive/save path. */
File* file_receive(const Config* config, int file_descriptor);
File* receive_incremental_check(int fd, const Config* config, bool* skipped);
int receive_manifest(int fd, const Config* config, int* next_status);
bool file_save_to_disk(const char* root_directory, const File* file, const Config* config);
#endif
+136
View File
@@ -0,0 +1,136 @@
#include <errno.h>
#include <fcntl.h>
#include <limits.h>
#include <poll.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/sendfile.h>
#include <sys/stat.h>
#include <time.h>
#include <unistd.h>
#include "compression.h"
#include "data.h"
#include "file.h"
#include "log.h"
#include "metadata.h"
#include "protocol.h"
bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata,
int compression_level, bool send_path) {
if (!file || !file->path || !file->data || (file->data->size != 0 && !file->data->data))
return false;
const Data* data_to_send = file->data;
Data* compressed_data = NULL;
if (compression_level > 0 && !compression_should_skip(file->path)) {
compressed_data = data_compress(file->data, compression_level);
if (compressed_data == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to compress file data");
return false;
}
data_to_send = compressed_data;
}
if (send_path && !send_str(file_descriptor, file->path)) {
data_destroy(compressed_data);
return false;
}
if (use_metadata && !metadata_send(file_descriptor, file->metadata)) {
data_destroy(compressed_data);
return false;
}
if (!send_data(file_descriptor, data_to_send)) {
data_destroy(compressed_data);
return false;
}
data_destroy(compressed_data);
return true;
}
bool file_send_sendfile(File* file, int file_descriptor, bool use_metadata, int compression_level,
bool send_path) {
if (!file || !file->path || !file->data)
return false;
if (compression_level > 0)
return file_send_single_calls(file, file_descriptor, use_metadata, compression_level,
send_path);
if (send_path && !send_str(file_descriptor, file->path))
return false;
if (use_metadata && !metadata_send(file_descriptor, file->metadata))
return false;
int fd = open(file->path, O_RDONLY);
if (fd == -1) {
log_perror("Could not open file for sendfile");
return false;
}
unsigned long long file_size = file->data->size;
struct stat source_stat;
if (fstat(fd, &source_stat) != 0 || !S_ISREG(source_stat.st_mode) ||
(unsigned long long)source_stat.st_size < file_size) {
close(fd);
return false;
}
if (!send_n_data(file_descriptor, &file_size, sizeof(unsigned long long))) {
close(fd);
return false;
}
/* sendfile cannot encrypt TLS records. Keep the framing identical but
route encrypted transfers through the deadline-aware IO layer. */
if (io_get_ssl() != NULL) {
unsigned char buffer[64 * 1024];
unsigned long long remaining = file_size;
bool ok = true;
while (remaining > 0) {
size_t want = remaining > sizeof(buffer) ? sizeof(buffer) : (size_t)remaining;
ssize_t got = read(fd, buffer, want);
if (got <= 0 || !send_n_data(file_descriptor, buffer, (size_t)got)) {
ok = false;
break;
}
remaining -= (unsigned long long)got;
}
close(fd);
return ok;
}
off_t offset = 0;
struct timespec deadline;
clock_gettime(CLOCK_MONOTONIC, &deadline);
deadline.tv_sec += 60;
while ((unsigned long long)offset < file_size) {
struct timespec now;
clock_gettime(CLOCK_MONOTONIC, &now);
long long remaining = (long long)(deadline.tv_sec - now.tv_sec) * 1000LL +
(deadline.tv_nsec - now.tv_nsec) / 1000000LL;
if (remaining <= 0) {
close(fd);
return false;
}
struct pollfd pfd = {.fd = file_descriptor, .events = POLLOUT};
int timeout = remaining > INT_MAX ? INT_MAX : (int)remaining;
int polled = poll(&pfd, 1, timeout);
if (polled <= 0 || (pfd.revents & (POLLERR | POLLHUP | POLLNVAL))) {
close(fd);
return false;
}
ssize_t sent = sendfile(file_descriptor, fd, &offset, file_size - offset);
if (sent == -1) {
if (errno == EAGAIN || errno == EINTR)
continue;
log_perror("sendfile failed");
close(fd);
return false;
}
if (sent == 0) {
close(fd);
return false;
}
}
close(fd);
return true;
}
+14
View File
@@ -0,0 +1,14 @@
#ifndef FILE_SEND_H
#define FILE_SEND_H
#include "file_types.h"
#include <stdbool.h>
/* Client-side file send path. */
bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata,
int compression_level, bool send_path);
bool file_send_sendfile(File* file, int file_descriptor, bool use_metadata, int compression_level,
bool send_path);
#endif
+197
View File
@@ -0,0 +1,197 @@
#include <errno.h>
#include <fcntl.h>
#include <libgen.h>
#include <limits.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include "file_store.h"
#include "metadata.h"
#include "utils.h"
static int authorized_root_fd = -1;
static char* authorized_root_path;
static bool path_is_within_root(const char* root, const char* path) {
size_t root_length = strlen(root);
return strncmp(root, path, root_length) == 0 &&
(path[root_length] == '\0' || path[root_length] == '/');
}
bool file_store_set_authorized_root(int fd, const char* canonical_path) {
char* new_path = canonical_path ? str_dup(canonical_path) : NULL;
if (canonical_path && !new_path) {
authorized_root_fd = -1;
free(authorized_root_path);
authorized_root_path = NULL;
return false;
}
free(authorized_root_path);
authorized_root_path = new_path;
authorized_root_fd = fd;
return true;
}
int file_store_open_secure_parent(const char* path, char** leaf_out) {
char* copy = str_dup(path);
if (!copy)
return -1;
char* parent = dirname(copy);
const char* slash = strrchr(path, '/');
char* leaf = str_dup(slash ? slash + 1 : path);
if (!leaf) {
free(copy);
return -1;
}
int fd;
if (authorized_root_fd >= 0) {
if (!authorized_root_path || path[0] != '/' ||
!path_is_within_root(authorized_root_path, path)) {
free(copy);
free(leaf);
return -1;
}
fd = dup(authorized_root_fd);
if (fd < 0) {
free(copy);
free(leaf);
return -1;
}
size_t root_length = strlen(authorized_root_path);
char* relative = str_dup(path + root_length);
if (!relative) {
free(copy);
free(leaf);
close(fd);
return -1;
}
free(copy);
copy = relative;
parent = dirname(copy);
} else {
fd = (parent[0] == '/') ? open("/", O_RDONLY | O_DIRECTORY | O_CLOEXEC)
: open(".", O_RDONLY | O_DIRECTORY | O_CLOEXEC);
}
if (fd < 0) {
free(copy);
free(leaf);
return -1;
}
char* save = NULL;
char* component = strtok_r(parent, "/", &save);
while (component) {
if (strcmp(component, "..") == 0) {
close(fd);
free(copy);
free(leaf);
return -1;
}
if (strcmp(component, ".") != 0) {
int next = openat(fd, component, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (next < 0 && errno == ENOENT) {
if (mkdirat(fd, component, 0755) == 0 || errno == EEXIST)
next = openat(fd, component, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
}
if (next < 0) {
close(fd);
free(copy);
free(leaf);
return -1;
}
close(fd);
fd = next;
}
component = strtok_r(NULL, "/", &save);
}
free(copy);
*leaf_out = leaf;
return fd;
}
bool file_store_rename_secure(const char* old_path, const char* new_path) {
char *old_leaf = NULL, *new_leaf = NULL;
int old_parent = file_store_open_secure_parent(old_path, &old_leaf);
int new_parent = file_store_open_secure_parent(new_path, &new_leaf);
bool ok = old_parent >= 0 && new_parent >= 0 &&
renameat(old_parent, old_leaf, new_parent, new_leaf) == 0;
if (old_parent >= 0)
close(old_parent);
if (new_parent >= 0)
close(new_parent);
free(old_leaf);
free(new_leaf);
return ok;
}
static bool write_all(int fd, const void* data, unsigned long long size) {
const unsigned char* p = data;
unsigned long long done = 0;
while (done < size) {
ssize_t n = write(fd, p + done, (size_t)(size - done));
if (n < 0 && errno == EINTR)
continue;
if (n <= 0)
return false;
done += (unsigned long long)n;
}
return true;
}
bool file_store_write_secure(const char* path, const void* data, unsigned long long data_size,
bool inplace, bool sparse, const FileMetadata* metadata) {
char* leaf = NULL;
int dirfd = file_store_open_secure_parent(path, &leaf);
if (dirfd < 0)
return false;
int fd = -1;
bool ok = false;
if (inplace) {
fd = openat(dirfd, leaf, O_WRONLY | O_CREAT | O_TRUNC | O_CLOEXEC | O_NOFOLLOW, 0644);
if (fd >= 0) {
if (!sparse || data_size == 0 || ftruncate(fd, (off_t)data_size) == 0)
ok = write_all(fd, data, data_size);
if (ok && metadata)
ok = file_restore_metadata_fd(fd, metadata);
}
} else {
int tmp_size = snprintf(NULL, 0, ".%s.tmp.%ld.%u", leaf, (long)getpid(), 99U);
if (tmp_size < 0) {
close(dirfd);
free(leaf);
return false;
}
char* tmp = malloc((size_t)tmp_size + 1);
if (!tmp) {
close(dirfd);
free(leaf);
return false;
}
for (unsigned int i = 0; i < 100 && !ok; ++i) {
snprintf(tmp, (size_t)tmp_size + 1, ".%s.tmp.%ld.%u", leaf, (long)getpid(), i);
fd = openat(dirfd, tmp, O_WRONLY | O_CREAT | O_EXCL | O_CLOEXEC | O_NOFOLLOW, 0600);
if (fd < 0)
continue;
if (sparse && data_size > 0)
ok = ftruncate(fd, (off_t)data_size) == 0;
if (ok || (!sparse || data_size == 0))
ok = write_all(fd, data, data_size);
if (ok && metadata)
ok = file_restore_metadata_fd(fd, metadata);
if (close(fd) != 0)
ok = false;
fd = -1;
if (ok && renameat(dirfd, tmp, dirfd, leaf) != 0)
ok = false;
if (!ok)
unlinkat(dirfd, tmp, 0);
}
free(tmp);
}
if (fd >= 0)
close(fd);
close(dirfd);
free(leaf);
return ok;
}
+13
View File
@@ -0,0 +1,13 @@
#ifndef FILE_STORE_H
#define FILE_STORE_H
#include "file.h"
#include <stdbool.h>
bool file_store_set_authorized_root(int fd, const char* canonical_path);
int file_store_open_secure_parent(const char* path, char** leaf_out);
bool file_store_rename_secure(const char* old_path, const char* new_path);
bool file_store_write_secure(const char* path, const void* data, unsigned long long data_size,
bool inplace, bool sparse, const FileMetadata* metadata);
#endif
+25
View File
@@ -0,0 +1,25 @@
#ifndef FILE_TYPES_H
#define FILE_TYPES_H
#include "data.h"
#include <stdbool.h>
#include <sys/stat.h>
typedef enum { FILE_TYPE_REGULAR, FILE_TYPE_SYMLINK, FILE_TYPE_DIR } FileType;
typedef struct {
mode_t mode;
uid_t uid;
gid_t gid;
time_t mtime_sec;
long mtime_nsec;
} FileMetadata;
typedef struct {
char* path;
Data* data;
FileMetadata* metadata;
bool skip;
} File;
#endif
+6
View File
@@ -1,6 +1,8 @@
#include "log.h"
#include <errno.h>
#include <stdarg.h>
#include <stdio.h>
#include <string.h>
#include <time.h>
static const char* log_level_strings[] = {"DEBUG", "INFO", "WARN", "ERROR"};
@@ -50,3 +52,7 @@ void log_message(LogLevel log_level, const char* format, ...) {
va_end(args);
}
}
void log_perror(const char* context) {
log_message(LOG_LEVEL_ERROR, "%s: %s", context, strerror(errno));
}
+1
View File
@@ -6,6 +6,7 @@
typedef enum { LOG_LEVEL_DEBUG, LOG_LEVEL_INFO, LOG_LEVEL_WARNING, LOG_LEVEL_ERROR } LogLevel;
void log_message(LogLevel log_level, const char* message, ...);
void log_perror(const char* context);
void set_log_level(LogLevel level);
void log_set_file(FILE* fp);
+12 -3
View File
@@ -51,6 +51,8 @@ FileMetadata* metadata_from_buf(char** buf) {
int32_t present;
memcpy(&present, *buf, sizeof(present));
*buf += sizeof(present);
if (present != 0 && present != 1)
return NULL;
if (!present)
return NULL;
FileMetadata* m = malloc(sizeof(FileMetadata));
@@ -76,10 +78,15 @@ FileMetadata* metadata_from_buf(char** buf) {
memcpy(&mtime_nsec, *buf, sizeof(mtime_nsec));
*buf += sizeof(mtime_nsec);
m->mtime_nsec = (long)mtime_nsec;
if (present != 1 || mtime_nsec < 0 || mtime_nsec >= 1000000000LL || mode < 0 || uid < 0 ||
gid < 0) {
free(m);
return NULL;
}
return m;
}
bool metadata_send(int file_descriptor, FileMetadata* m) {
bool metadata_send(int file_descriptor, const FileMetadata* m) {
if (m == NULL) {
int32_t zero = 0;
return send_n_data(file_descriptor, &zero, sizeof(zero));
@@ -175,7 +182,8 @@ FileMetadata* metadata_receive(int file_descriptor, int* ok) {
void file_restore_metadata(const char* path, const FileMetadata* metadata) {
if (metadata == NULL)
return;
if (chmod(path, metadata->mode & 07777 & ~(S_ISUID | S_ISGID)) != 0)
mode_t safe_mode = metadata->mode & 0777 & ~(S_IWGRP | S_IWOTH);
if (chmod(path, safe_mode) != 0)
log_message(LOG_LEVEL_WARNING, "Failed to chmod %s: %s", path, strerror(errno));
/* Never apply client-supplied ownership. The descriptor API below is the
receiver write path; retain this legacy API only for compatibility. */
@@ -192,7 +200,8 @@ bool file_restore_metadata_fd(int fd, const FileMetadata* metadata) {
if (fd < 0 || metadata == NULL)
return metadata == NULL;
bool ok = true;
if (fchmod(fd, metadata->mode & 07777 & ~(S_ISUID | S_ISGID)) != 0)
mode_t safe_mode = metadata->mode & 0777 & ~(S_IWGRP | S_IWOTH);
if (fchmod(fd, safe_mode) != 0)
ok = false;
/* Client uid/gid values are deliberately not authoritative. */
struct timespec times[2] = {{.tv_sec = 0, .tv_nsec = UTIME_OMIT},
+1 -1
View File
@@ -27,7 +27,7 @@
void metadata_to_buf(char** buf, const FileMetadata* m);
FileMetadata* metadata_from_buf(char** buf);
bool metadata_send(int file_descriptor, FileMetadata* m);
bool metadata_send(int file_descriptor, const FileMetadata* m);
FileMetadata* metadata_receive(int file_descriptor, int* ok);
void file_restore_metadata(const char* path, const FileMetadata* metadata);
bool file_restore_metadata_fd(int fd, const FileMetadata* metadata);
+29 -121
View File
@@ -1,4 +1,5 @@
#include "multiprocessing.h"
#include "receiver.h"
#include "array_list.h"
#include "chunk.h"
@@ -14,10 +15,6 @@
#include <string.h>
#include <threads.h>
static bool valid_batch_path(const char* path) {
return path && path[0] != '\0' && path[0] != '/' && !has_path_traversal(path);
}
PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner,
Queue* queue_loader) {
PipelineContextSender* context = malloc(sizeof(PipelineContextSender));
@@ -58,7 +55,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
return context;
fail:
perror("Error initializing synchronization objects");
log_perror("Error initializing synchronization objects");
if (init >= 6)
cnd_destroy(&context->condition_not_empty_loader);
if (init >= 5)
@@ -101,6 +98,8 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
context->queue = queue;
context->file_descriptor = file_descriptor;
context->ssl = ssl;
protocol_session_init(&context->session, file_descriptor, file_descriptor);
protocol_session_set_ssl(&context->session, ssl);
context->receiver_done = false;
atomic_init(&context->cancelled, false);
int init = 0;
@@ -117,7 +116,7 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
return context;
fail:
perror("Error initializing synchronization objects");
log_perror("Error initializing synchronization objects");
if (init >= 3)
cnd_destroy(&context->condition_not_empty);
if (init >= 2)
@@ -137,24 +136,14 @@ void pipeline_context_receiver_destroy(PipelineContextReceiver* context) {
free(context);
}
static bool receive_chunk_enqueue(int file_descriptor, PipelineContextReceiver* context) {
Chunk* chunk = receive_chunk_data(file_descriptor, context->config);
if (chunk == NULL)
return false;
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 bool receiver_enqueue_file(File* file, void* context_pointer) {
PipelineContextReceiver* context = context_pointer;
if (queue_enqueue_multithreaded_cancel(context->queue, file, &context->mutex,
&context->condition_not_empty,
&context->condition_not_full, &context->cancelled))
return true;
file_destroy(file);
return false;
}
static void receiver_thread_fail(PipelineContextReceiver* context) {
@@ -167,123 +156,42 @@ static void receiver_thread_fail(PipelineContextReceiver* 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;
if (context->ssl)
io_set_ssl(context->ssl);
protocol_session_bind(&context->session);
mtx_lock(&context->mutex);
int file_descriptor = context->file_descriptor;
const Config* config = context->config;
mtx_unlock(&context->mutex);
Status status;
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();
ReceiverSink sink = {receiver_enqueue_file, context, false, false};
if (receiver_process((Config*)config, file_descriptor, &sink) != 0) {
receiver_thread_fail(context);
protocol_session_unbind();
return thrd_error;
}
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);
context->receiver_done = true;
cnd_signal(&context->condition_not_empty);
mtx_unlock(&context->mutex);
#undef RECEIVE_THREAD_FAIL
protocol_session_unbind();
return thrd_success;
}
int write_thread(void* pipeline_context) {
PipelineContextReceiver* context = (PipelineContextReceiver*)pipeline_context;
if (context->ssl)
io_set_ssl(context->ssl);
mtx_lock(&context->mutex);
bool save_to_disk = context->config->save_to_disk;
char* root_directory = str_dup(context->config->receive_root_directory);
mtx_unlock(&context->mutex);
if (save_to_disk && !root_directory) {
mtx_lock(&context->mutex);
atomic_store(&context->cancelled, true);
context->receiver_done = true;
cnd_broadcast(&context->condition_not_full);
cnd_broadcast(&context->condition_not_empty);
mtx_unlock(&context->mutex);
return thrd_error;
}
while (true) {
File* file =
+1
View File
@@ -35,6 +35,7 @@ typedef struct PipelineContextReceiver {
Config* config;
int file_descriptor;
SSL* ssl;
ProtocolSession session;
mtx_t mutex;
cnd_t condition_not_full;
cnd_t condition_not_empty;
+219 -86
View File
@@ -13,75 +13,130 @@
#define RECEIVE_TIMEOUT_SEC 60 /* 60 second per-message timeout */
#define SEND_TIMEOUT_SEC 60
#define MAX_CONNECTION_MEMORY (1024ULL * 1024 * 1024) /* 1 GB total per connection */
#define MAX_CONNECTION_MEMORY (256ULL * 1024 * 1024) /* bounded cumulative receive budget */
static __thread int io_read_fd = -1;
static __thread int io_write_fd = -1;
static __thread SSL* io_ssl;
static __thread ProtocolSession* bound_session;
static __thread ProtocolSession legacy_io_session = {.read_fd = -1, .write_fd = -1};
static unsigned long long io_bwlimit = 0;
static long long bw_tokens = 0;
static struct timespec bw_last_refill = {0, 0};
static mtx_t bw_mutex;
static once_flag bw_mutex_once = ONCE_FLAG_INIT;
static __thread unsigned long long total_allocated_bytes = 0;
static unsigned long long global_bwlimit(void);
void protocol_release_memory(size_t charge) {
ProtocolSession* session = bound_session ? bound_session : &legacy_io_session;
if ((unsigned long long)charge >= session->total_allocated_bytes)
session->total_allocated_bytes = 0;
else
session->total_allocated_bytes -= charge;
}
void io_set_fds(int read_fd, int write_fd) {
bound_session = NULL;
io_read_fd = read_fd;
io_write_fd = write_fd;
/* A descriptor switch starts a new transport; never reuse a TLS object
belonging to a previous connection or test pipe. */
io_ssl = NULL;
total_allocated_bytes = 0;
legacy_io_session.read_fd = read_fd;
legacy_io_session.write_fd = write_fd;
legacy_io_session.ssl = NULL;
legacy_io_session.total_allocated_bytes = 0;
protocol_session_set_bwlimit(&legacy_io_session, global_bwlimit());
}
void protocol_session_init(ProtocolSession* session, int read_fd, int write_fd) {
if (!session)
return;
memset(session, 0, sizeof(*session));
session->read_fd = read_fd;
session->write_fd = write_fd;
protocol_session_set_bwlimit(session, global_bwlimit());
}
void protocol_session_bind(ProtocolSession* session) {
bound_session = session;
}
void protocol_session_unbind(void) {
bound_session = NULL;
}
void protocol_session_set_ssl(ProtocolSession* session, SSL* ssl) {
if (session)
session->ssl = ssl;
}
static void bw_mutex_init(void) {
mtx_init(&bw_mutex, mtx_plain);
}
static unsigned long long global_bwlimit(void) {
unsigned long long limit;
call_once(&bw_mutex_once, bw_mutex_init);
mtx_lock(&bw_mutex);
limit = io_bwlimit;
mtx_unlock(&bw_mutex);
return limit;
}
void io_set_bwlimit(unsigned long long bytes_per_sec) {
call_once(&bw_mutex_once, bw_mutex_init);
mtx_lock(&bw_mutex);
io_bwlimit = bytes_per_sec;
bw_tokens = (long long)io_bwlimit;
clock_gettime(CLOCK_MONOTONIC, &bw_last_refill);
io_bwlimit =
bytes_per_sec > (unsigned long long)LLONG_MAX ? (unsigned long long)LLONG_MAX : bytes_per_sec;
mtx_unlock(&bw_mutex);
}
static void bw_throttle(size_t bytes_written) {
if (io_bwlimit == 0)
void protocol_session_set_bwlimit(ProtocolSession* session, unsigned long long bytes_per_sec) {
if (!session)
return;
session->bwlimit =
bytes_per_sec > (unsigned long long)LLONG_MAX ? (unsigned long long)LLONG_MAX : bytes_per_sec;
session->bw_tokens = (long long)session->bwlimit;
struct timespec now;
clock_gettime(CLOCK_MONOTONIC, &now);
session->bw_last_refill_sec = now.tv_sec;
session->bw_last_refill_nsec = now.tv_nsec;
}
static void bw_throttle_session(ProtocolSession* session, size_t bytes_written) {
if (session->bwlimit == 0)
return;
call_once(&bw_mutex_once, bw_mutex_init);
mtx_lock(&bw_mutex);
struct timespec now;
clock_gettime(CLOCK_MONOTONIC, &now);
long long elapsed_ns =
(now.tv_sec - bw_last_refill.tv_sec) * 1000000000LL + (now.tv_nsec - bw_last_refill.tv_nsec);
bw_last_refill = now;
long long elapsed_ns = (now.tv_sec - session->bw_last_refill_sec) * 1000000000LL +
(now.tv_nsec - session->bw_last_refill_nsec);
session->bw_last_refill_sec = now.tv_sec;
session->bw_last_refill_nsec = now.tv_nsec;
long long tokens_to_add = (long long)((double)io_bwlimit * elapsed_ns / 1000000000.0);
bw_tokens += tokens_to_add;
if (bw_tokens > (long long)io_bwlimit)
bw_tokens = (long long)io_bwlimit;
long long tokens_to_add = (long long)((double)session->bwlimit * elapsed_ns / 1000000000.0);
session->bw_tokens += tokens_to_add;
if (session->bw_tokens > (long long)session->bwlimit)
session->bw_tokens = (long long)session->bwlimit;
bw_tokens -= (long long)bytes_written;
session->bw_tokens -= bytes_written;
if (bw_tokens < 0) {
long long deficit_us = (long long)((double)(-bw_tokens) / io_bwlimit * 1000000.0);
if (session->bw_tokens < 0) {
long long deficit_us =
(long long)((double)(-session->bw_tokens) / session->bwlimit * 1000000.0);
if (deficit_us >= 1000)
poll(NULL, 0, (int)(deficit_us / 1000));
else
usleep((useconds_t)deficit_us);
bw_tokens = 0;
clock_gettime(CLOCK_MONOTONIC, &bw_last_refill);
session->bw_tokens = 0;
session->bw_last_refill_sec = now.tv_sec;
session->bw_last_refill_nsec = now.tv_nsec;
}
mtx_unlock(&bw_mutex);
}
void io_set_ssl(SSL* ssl) {
bound_session = NULL;
io_ssl = ssl;
}
@@ -89,8 +144,30 @@ SSL* io_get_ssl(void) {
return io_ssl;
}
static int io_fd(int dir_fd, int file_descriptor) {
return (dir_fd != -1) ? dir_fd : file_descriptor;
static ProtocolSession* legacy_session(int read_fd, int write_fd) {
if (bound_session)
return bound_session;
int target_read_fd = io_read_fd != -1 ? io_read_fd : read_fd;
int target_write_fd = io_write_fd != -1 ? io_write_fd : write_fd;
if (legacy_io_session.read_fd != target_read_fd ||
legacy_io_session.write_fd != target_write_fd) {
legacy_io_session.read_fd = target_read_fd;
legacy_io_session.write_fd = target_write_fd;
legacy_io_session.total_allocated_bytes = 0;
protocol_session_set_bwlimit(&legacy_io_session, global_bwlimit());
} else if (legacy_io_session.bwlimit != global_bwlimit()) {
protocol_session_set_bwlimit(&legacy_io_session, global_bwlimit());
}
legacy_io_session.ssl = io_ssl;
return &legacy_io_session;
}
bool send_n_data(int file_descriptor, const void* data, size_t data_size) {
return protocol_send_n_data(legacy_session(-1, file_descriptor), data, data_size);
}
bool receive_n_data(int file_descriptor, void* data, size_t data_size) {
return protocol_receive_n_data(legacy_session(file_descriptor, -1), data, data_size);
}
static int deadline_remaining_ms(const struct timespec* deadline) {
@@ -104,9 +181,13 @@ static int deadline_remaining_ms(const struct timespec* deadline) {
return ms > INT_MAX ? INT_MAX : (int)ms;
}
bool send_n_data(int file_descriptor, const void* data, size_t data_size) {
bool protocol_send_n_data(ProtocolSession* session, const void* data, size_t data_size) {
if (!data && data_size != 0)
return false;
log_message(LOG_LEVEL_DEBUG, " Sending n Data: %zu", data_size);
int fd = io_fd(io_write_fd, file_descriptor);
if (!session)
return false;
int fd = session->write_fd;
struct timespec deadline;
clock_gettime(CLOCK_MONOTONIC, &deadline);
deadline.tv_sec += SEND_TIMEOUT_SEC;
@@ -114,7 +195,7 @@ bool send_n_data(int file_descriptor, const void* data, size_t data_size) {
ssize_t total_bytes_send = 0;
while ((size_t)total_bytes_send < data_size) {
size_t chunk = data_size - total_bytes_send;
if (io_bwlimit > 0 && chunk > 65536)
if (session->bwlimit > 0 && chunk > 65536)
chunk = 65536;
struct pollfd pfd = {.fd = fd, .events = wait_events};
int poll_result = poll(&pfd, 1, deadline_remaining_ms(&deadline));
@@ -127,13 +208,13 @@ bool send_n_data(int file_descriptor, const void* data, size_t data_size) {
if (pfd.revents & (POLLERR | POLLNVAL))
return false;
ssize_t bytes_send;
if (io_ssl)
bytes_send = SSL_write(io_ssl, (const char*)data + total_bytes_send, chunk);
if (session->ssl)
bytes_send = SSL_write(session->ssl, (const char*)data + total_bytes_send, chunk);
else
bytes_send = write(fd, (const char*)data + total_bytes_send, chunk);
if (bytes_send <= 0) {
if (io_ssl) {
int ssl_err = SSL_get_error(io_ssl, (int)bytes_send);
if (session->ssl) {
int ssl_err = SSL_get_error(session->ssl, (int)bytes_send);
if (ssl_err == SSL_ERROR_WANT_WRITE || ssl_err == SSL_ERROR_WANT_READ) {
wait_events = ssl_err == SSL_ERROR_WANT_WRITE ? POLLOUT : POLLIN;
continue;
@@ -142,16 +223,20 @@ bool send_n_data(int file_descriptor, const void* data, size_t data_size) {
log_message(LOG_LEVEL_ERROR, "Could not send data");
return false;
}
bw_throttle((size_t)bytes_send);
bw_throttle_session(session, (size_t)bytes_send);
total_bytes_send += bytes_send;
if (session->ssl)
wait_events = POLLOUT;
}
log_message(LOG_LEVEL_DEBUG, " Send n Data: %zu", total_bytes_send);
return true;
}
bool receive_n_data(int file_descriptor, void* data, size_t data_size) {
bool protocol_receive_n_data(ProtocolSession* session, void* data, size_t data_size) {
log_message(LOG_LEVEL_DEBUG, " Receiving n Data: %zu", data_size);
int fd = io_fd(io_read_fd, file_descriptor);
if (!session)
return false;
int fd = session->read_fd;
struct timespec deadline;
clock_gettime(CLOCK_MONOTONIC, &deadline);
@@ -160,31 +245,33 @@ bool receive_n_data(int file_descriptor, void* data, size_t data_size) {
size_t total_bytes_received = 0;
short wait_events = POLLIN;
while (total_bytes_received < data_size) {
struct pollfd pfd = {.fd = fd, .events = wait_events};
int poll_result = poll(&pfd, 1, deadline_remaining_ms(&deadline));
if (poll_result == 0) {
log_message(LOG_LEVEL_ERROR, "Receive timeout after %ds", RECEIVE_TIMEOUT_SEC);
return false;
if (!session->ssl || SSL_pending(session->ssl) == 0) {
struct pollfd pfd = {.fd = fd, .events = wait_events};
int poll_result = poll(&pfd, 1, deadline_remaining_ms(&deadline));
if (poll_result == 0) {
log_message(LOG_LEVEL_ERROR, "Receive timeout after %ds", RECEIVE_TIMEOUT_SEC);
return false;
}
if (poll_result < 0) {
if (errno == EINTR)
continue;
return false;
}
/* POLLHUP may accompany the final readable bytes on pipes/sockets. */
if (pfd.revents & (POLLERR | POLLNVAL))
return false;
}
if (poll_result < 0) {
if (errno == EINTR)
continue;
return false;
}
/* POLLHUP may accompany the final readable bytes on pipes/sockets. */
if (pfd.revents & (POLLERR | POLLNVAL))
return false;
ssize_t bytes_received;
if (io_ssl)
bytes_received =
SSL_read(io_ssl, (char*)data + total_bytes_received, data_size - total_bytes_received);
if (session->ssl)
bytes_received = SSL_read(session->ssl, (char*)data + total_bytes_received,
data_size - total_bytes_received);
else
bytes_received =
read(fd, (char*)data + total_bytes_received, data_size - total_bytes_received);
if (bytes_received <= 0) {
if (io_ssl) {
int ssl_err = SSL_get_error(io_ssl, (int)bytes_received);
if (session->ssl) {
int ssl_err = SSL_get_error(session->ssl, (int)bytes_received);
if (ssl_err == SSL_ERROR_WANT_WRITE || ssl_err == SSL_ERROR_WANT_READ) {
wait_events = ssl_err == SSL_ERROR_WANT_WRITE ? POLLOUT : POLLIN;
continue;
@@ -196,7 +283,9 @@ bool receive_n_data(int file_descriptor, void* data, size_t data_size) {
log_message(LOG_LEVEL_ERROR, "Could not receive bytes");
return false;
}
total_bytes_received += bytes_received;
total_bytes_received += (size_t)bytes_received;
if (session->ssl)
wait_events = POLLIN;
}
log_message(LOG_LEVEL_DEBUG, " Received n Data: %zu", total_bytes_received);
return true;
@@ -231,24 +320,24 @@ static const char* status_to_string(Status status) {
}
}
bool send_str(int file_descriptor, const char* data) {
bool protocol_send_str(ProtocolSession* session, const char* data) {
if (data == NULL)
return false;
size_t size = strlen(data);
if (!send_n_data(file_descriptor, &size, sizeof(size_t)))
if (!protocol_send_n_data(session, &size, sizeof(size_t)))
return false;
if (!send_n_data(file_descriptor, data, size))
if (!protocol_send_n_data(session, data, size))
return false;
log_message(LOG_LEVEL_DEBUG, "Send String: %s", data);
return true;
}
char* receive_str(int file_descriptor) {
char* protocol_receive_str(ProtocolSession* session) {
size_t size;
if (!receive_n_data(file_descriptor, &size, sizeof(size_t)))
if (!protocol_receive_n_data(session, &size, sizeof(size_t)))
return NULL;
if (size > MAX_STRING_SIZE || size > SIZE_MAX - 1 ||
size + 1 > MAX_CONNECTION_MEMORY - total_allocated_bytes) {
size + 1 > MAX_CONNECTION_MEMORY - session->total_allocated_bytes) {
log_message(LOG_LEVEL_ERROR, "String size %zu exceeds maximum %llu", size,
(unsigned long long)MAX_STRING_SIZE);
return NULL;
@@ -256,83 +345,127 @@ char* receive_str(int file_descriptor) {
char* data = (char*)malloc(size + 1);
if (data == NULL)
return NULL;
if (!receive_n_data(file_descriptor, data, size)) {
if (!protocol_receive_n_data(session, data, size)) {
free(data);
return NULL;
}
if (memchr(data, '\0', size) != NULL) {
free(data);
log_message(LOG_LEVEL_ERROR, "Received string contains an embedded NUL");
return NULL;
}
data[size] = '\0';
total_allocated_bytes += size + 1;
session->total_allocated_bytes += size + 1;
log_message(LOG_LEVEL_DEBUG, "Received String: %s", data);
return data;
}
bool send_data(int file_descriptor, const Data* data) {
unsigned long long data_size = data->size;
if (!send_n_data(file_descriptor, &data_size, sizeof(unsigned long long)))
bool protocol_send_data(ProtocolSession* session, const Data* data) {
if (!data || (!data->data && data->size != 0))
return false;
if (!send_n_data(file_descriptor, data->data, data_size))
if (!session)
return false;
unsigned long long data_size = data->size;
if (!protocol_send_n_data(session, &data_size, sizeof(unsigned long long)))
return false;
if (!protocol_send_n_data(session, data->data, data_size))
return false;
log_message(LOG_LEVEL_DEBUG, "Send %lld data", data_size);
return true;
}
Data* receive_data(int file_descriptor) {
unsigned long long size = 0;
if (!receive_n_data(file_descriptor, &size, sizeof(unsigned long long)))
Data* protocol_receive_data_limited(ProtocolSession* session, unsigned long long maximum_size) {
if (!session)
return NULL;
if (size > MAX_DATA_PAYLOAD_SIZE) {
unsigned long long size = 0;
if (!protocol_receive_n_data(session, &size, sizeof(unsigned long long)))
return NULL;
if (size > MAX_DATA_PAYLOAD_SIZE || size > maximum_size) {
log_message(LOG_LEVEL_ERROR, "Data size %llu exceeds maximum %llu", size,
(unsigned long long)MAX_DATA_PAYLOAD_SIZE);
return NULL;
}
size_t allocation_size = size == 0 ? 1 : (size_t)size;
if (allocation_size > MAX_CONNECTION_MEMORY - total_allocated_bytes) {
if (allocation_size > MAX_CONNECTION_MEMORY - session->total_allocated_bytes) {
log_message(LOG_LEVEL_ERROR, "Per-connection memory limit exceeded (%llu + %llu > %llu)",
(unsigned long long)total_allocated_bytes, size,
(unsigned long long)session->total_allocated_bytes, size,
(unsigned long long)MAX_CONNECTION_MEMORY);
return NULL;
}
void* data = malloc(allocation_size);
if (data == NULL)
return NULL;
if (!receive_n_data(file_descriptor, data, (size_t)size)) {
if (!protocol_receive_n_data(session, data, (size_t)size)) {
free(data);
return NULL;
}
total_allocated_bytes += allocation_size;
session->total_allocated_bytes += allocation_size;
log_message(LOG_LEVEL_DEBUG, "Received %lld data", size);
Data* result = data_create(data, (size_t)size);
if (!result) {
free(data);
total_allocated_bytes -= allocation_size;
session->total_allocated_bytes -= allocation_size;
return NULL;
}
result->protocol_charge = allocation_size;
return result;
}
bool send_int(int file_descriptor, int data) {
if (!send_n_data(file_descriptor, &data, sizeof(int)))
Data* protocol_receive_data(ProtocolSession* session) {
return protocol_receive_data_limited(session, MAX_DATA_PAYLOAD_SIZE);
}
bool protocol_send_int(ProtocolSession* session, int data) {
if (!protocol_send_n_data(session, &data, sizeof(int)))
return false;
log_message(LOG_LEVEL_DEBUG, "Send Int: %d", data);
return true;
}
bool receive_int(int file_descriptor, int* data) {
if (!receive_n_data(file_descriptor, data, sizeof(int)))
bool protocol_receive_int(ProtocolSession* session, int* data) {
if (!protocol_receive_n_data(session, data, sizeof(int)))
return false;
log_message(LOG_LEVEL_DEBUG, "Received Int: %d", *data);
return true;
}
bool send_status(int file_descriptor, Status status) {
if (!send_n_data(file_descriptor, &status, sizeof(Status)))
bool protocol_send_status(ProtocolSession* session, Status status) {
if (!protocol_send_n_data(session, &status, sizeof(Status)))
return false;
log_message(LOG_LEVEL_DEBUG, "Send Status: %s", status_to_string(status));
return true;
}
bool receive_status(int file_descriptor, Status* status) {
if (!receive_n_data(file_descriptor, status, sizeof(Status)))
bool protocol_receive_status(ProtocolSession* session, Status* status) {
if (!protocol_receive_n_data(session, status, sizeof(Status)))
return false;
log_message(LOG_LEVEL_DEBUG, "Received Status: %s", status_to_string(*status));
return true;
}
bool send_str(int fd, const char* data) {
return protocol_send_str(legacy_session(-1, fd), data);
}
char* receive_str(int fd) {
return protocol_receive_str(legacy_session(fd, -1));
}
bool send_data(int fd, const Data* data) {
return protocol_send_data(legacy_session(-1, fd), data);
}
Data* receive_data(int fd) {
return protocol_receive_data_limited(legacy_session(fd, -1), MAX_DATA_PAYLOAD_SIZE);
}
Data* receive_data_limited(int fd, unsigned long long maximum_size) {
return protocol_receive_data_limited(legacy_session(fd, -1), maximum_size);
}
bool send_int(int fd, int data) {
return protocol_send_int(legacy_session(-1, fd), data);
}
bool receive_int(int fd, int* data) {
return protocol_receive_int(legacy_session(fd, -1), data);
}
bool send_status(int fd, Status status) {
return protocol_send_status(legacy_session(-1, fd), status);
}
bool receive_status(int fd, Status* status) {
return protocol_receive_status(legacy_session(fd, -1), status);
}
+38
View File
@@ -10,6 +10,8 @@
/* Maximum allowed data payload size for receive_data (100 MB) */
#define MAX_DATA_PAYLOAD_SIZE (100ULL * 1024 * 1024)
/* Maximum uncompressed file payload accepted by the receiver. */
#define MAX_RECEIVE_FILE_SIZE (64ULL * 1024 * 1024)
/* Maximum chunk size (64 MB) — prevents unbounded allocation from the wire */
#define MAX_CHUNK_SIZE (64ULL * 1024 * 1024)
@@ -19,6 +21,23 @@
typedef struct ssl_st SSL;
/*
* Explicit owner of protocol I/O. A session does not own the descriptors or
* SSL object; it only describes the transport used by a transfer. This makes
* it safe to pass the transport to a worker without relying on inherited
* thread-local state.
*/
typedef struct ProtocolSession {
int read_fd;
int write_fd;
SSL* ssl;
unsigned long long bwlimit;
long long bw_tokens;
long long bw_last_refill_sec;
long bw_last_refill_nsec;
unsigned long long total_allocated_bytes;
} ProtocolSession;
typedef int Status;
enum NET_STATUS {
STATUS_OK,
@@ -39,6 +58,24 @@ void io_set_fds(int read_fd, int write_fd);
void io_set_bwlimit(unsigned long long bytes_per_sec);
void io_set_ssl(SSL* ssl);
SSL* io_get_ssl(void);
void protocol_session_init(ProtocolSession* session, int read_fd, int write_fd);
/* Transitional bridge for helpers whose signatures still carry only an fd. */
void protocol_session_bind(ProtocolSession* session);
void protocol_session_unbind(void);
void protocol_session_set_ssl(ProtocolSession* session, SSL* ssl);
void protocol_session_set_bwlimit(ProtocolSession* session, unsigned long long bytes_per_sec);
bool protocol_send_n_data(ProtocolSession* session, const void* data, size_t data_size);
bool protocol_receive_n_data(ProtocolSession* session, void* data, size_t data_size);
bool protocol_send_str(ProtocolSession* session, const char* data);
char* protocol_receive_str(ProtocolSession* session);
bool protocol_send_data(ProtocolSession* session, const Data* data);
Data* protocol_receive_data(ProtocolSession* session);
Data* protocol_receive_data_limited(ProtocolSession* session, unsigned long long maximum_size);
bool protocol_send_int(ProtocolSession* session, int data);
bool protocol_receive_int(ProtocolSession* session, int* data);
bool protocol_send_status(ProtocolSession* session, Status status);
bool protocol_receive_status(ProtocolSession* session, Status* status);
bool send_n_data(int file_descriptor, const void* data, size_t data_size);
bool receive_n_data(int file_descriptor, void* data, size_t data_size);
@@ -46,6 +83,7 @@ bool send_str(int file_descriptor, const char* data);
char* receive_str(int file_descriptor);
bool send_data(int file_descriptor, const Data* data);
Data* receive_data(int file_descriptor);
Data* receive_data_limited(int file_descriptor, unsigned long long maximum_size);
bool send_int(int file_descriptor, int data);
bool receive_int(int file_descriptor, int* data);
bool send_status(int file_descriptor, Status status);
+4 -3
View File
@@ -1,3 +1,4 @@
#include "log.h"
#include <stdbool.h>
#include <limits.h>
#include <stdio.h>
@@ -13,7 +14,7 @@ Queue* queue_create(int capacity, void (*destroyer)(void* item)) {
Queue* queue = (Queue*)malloc(sizeof(Queue));
if (queue == NULL) {
perror("ERROR: Could not allocate memory for queue structure");
log_perror("ERROR: Could not allocate memory for queue structure");
return NULL;
}
@@ -72,7 +73,7 @@ static bool queue_double_capacity(Queue* queue) {
new_capacity = 100;
void** new_items = malloc(new_capacity * sizeof(void*));
if (new_items == NULL) {
perror("ERROR: Could not allocate memory for doubling capacity of queue.");
log_perror("ERROR: Could not allocate memory for doubling capacity of queue.");
return false;
}
for (int i = 0; i < queue->size; i++)
@@ -127,7 +128,7 @@ bool queue_enqueue_multithreaded_cancel(Queue* queue, void* item, mtx_t* mutex,
void* queue_dequeue(Queue* queue) {
if (queue == NULL || queue_is_empty(queue)) {
perror("ERROR: Could not dequeue from null or empty queue.");
log_perror("ERROR: Could not dequeue from null or empty queue.");
return NULL;
}
+5 -4
View File
@@ -1,3 +1,4 @@
#include "log.h"
#include "transport_ssh.h"
#include "utils.h"
#include <fcntl.h>
@@ -76,7 +77,7 @@ Client* client_connect_ssh(const char* destination, int port, const char* server
int sv[2];
if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) {
perror("socketpair failed");
log_perror("socketpair failed");
remote_dest_destroy(&r);
return NULL;
}
@@ -89,7 +90,7 @@ Client* client_connect_ssh(const char* destination, int port, const char* server
int exec_pipe[2];
if (pipe(exec_pipe) < 0) {
perror("pipe failed");
log_perror("pipe failed");
close(sv[0]);
close(sv[1]);
remote_dest_destroy(&r);
@@ -98,7 +99,7 @@ Client* client_connect_ssh(const char* destination, int port, const char* server
pid_t pid = fork();
if (pid < 0) {
perror("fork failed");
log_perror("fork failed");
close(sv[0]);
close(sv[1]);
close(exec_pipe[0]);
@@ -151,7 +152,7 @@ Client* client_connect_ssh(const char* destination, int port, const char* server
ssh_argv[ac++] = "--stdio";
ssh_argv[ac] = NULL;
execvp("ssh", ssh_argv);
perror("exec of ssh failed");
log_perror("exec of ssh failed");
ssize_t wret = write(exec_pipe[1], "x", 1);
(void)wret;
_exit(1);
+18 -8
View File
@@ -30,20 +30,20 @@ static void sigchld_handler(int sig) {
Server* server_create(int port) {
Server* server = (Server*)malloc(sizeof(Server));
if (server == NULL) {
perror("Could not allocate space for Server");
log_perror("Could not allocate space for Server");
return NULL;
}
int file_descriptor = socket(AF_INET, SOCK_STREAM, 0);
if (file_descriptor < 0) {
perror("Could not create Socket!");
log_perror("Could not create Socket!");
free(server);
return NULL;
}
server->file_descriptor = file_descriptor;
int opt = 1;
if (setsockopt(server->file_descriptor, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt))) {
perror("Error setting a socket option!");
log_perror("Error setting a socket option!");
close(server->file_descriptor);
free(server);
return NULL;
@@ -59,7 +59,7 @@ Server* server_create(int port) {
if (bind(server->file_descriptor, (struct sockaddr*)&server->address, server->address_length) <
0) {
perror("Could not bind server");
log_perror("Could not bind server");
close(server->file_descriptor);
free(server);
return NULL;
@@ -83,7 +83,7 @@ void server_delete(Server** server) {
static void accept_loop(Server* server, void (*child_fn)(int, void*), void* child_ctx,
const char* log_fmt) {
if (listen(server->file_descriptor, SOMAXCONN) < 0) {
perror("Could not listen on port!");
log_perror("Could not listen on port!");
return;
}
signal(SIGCHLD, sigchld_handler);
@@ -92,7 +92,7 @@ static void accept_loop(Server* server, void (*child_fn)(int, void*), void* chil
socklen_t client_len = sizeof(client_addr);
int fd = accept(server->file_descriptor, (struct sockaddr*)&client_addr, &client_len);
if (fd < 0) {
perror("Could not accept the connection");
log_perror("Could not accept the connection");
continue;
}
tcp_apply_socket_timeout(fd);
@@ -151,6 +151,10 @@ int tcp_get_contimeout_sec(void) {
return g_contimeout_sec;
}
int tcp_get_timeout_sec(void) {
return g_timeout_sec;
}
static void tcp_apply_socket_timeout(int fd) {
struct timeval tv;
tv.tv_sec = g_timeout_sec;
@@ -173,7 +177,7 @@ Client* client_create() {
return client;
}
bool client_connect(Client* client, char* host, int port) {
bool tcp_connect_socket(Client* client, char* host, int port) {
struct addrinfo hints;
struct addrinfo* result;
memset(&hints, 0, sizeof(hints));
@@ -218,10 +222,16 @@ bool client_connect(Client* client, char* host, int port) {
freeaddrinfo(result);
if (!connected) {
perror("Could not connect to Server!");
log_perror("Could not connect to Server!");
return false;
}
return true;
}
bool client_connect(Client* client, char* host, int port) {
if (!tcp_connect_socket(client, host, port))
return false;
tcp_apply_socket_timeout(client->file_descriptor);
return true;
}
+2
View File
@@ -30,9 +30,11 @@ void server_accept_loop(Server* server, void (*child_fn)(int, void*), void* chil
void server_delete(Server** server);
Client* client_create();
bool client_connect(Client* client, char* host, int port);
bool tcp_connect_socket(Client* client, char* host, int port);
void client_disconnect(Client* client);
void client_delete(Client* client);
void tcp_set_timeouts(int timeout_sec, int contimeout_sec);
int tcp_get_contimeout_sec(void);
int tcp_get_timeout_sec(void);
#endif
+43 -50
View File
@@ -3,7 +3,6 @@
#include "protocol.h"
#include "transport_tcp.h"
#include <arpa/inet.h>
#include <netdb.h>
#include <openssl/err.h>
#include <openssl/ssl.h>
#include <signal.h>
@@ -11,7 +10,9 @@
#include <stdlib.h>
#include <string.h>
#include <sys/socket.h>
#include <sys/stat.h>
#include <sys/wait.h>
#include <time.h>
#include <unistd.h>
bool tls_global_init(void) {
@@ -34,6 +35,10 @@ static void log_ssl_errors(void) {
static SSL_CTX* create_ssl_ctx(bool is_server, const char* cert, const char* key,
const char* ca_path) {
if (!is_server && !ca_path) {
log_message(LOG_LEVEL_ERROR, "TLS clients require a CA certificate path");
return NULL;
}
const SSL_METHOD* method = is_server ? TLS_server_method() : TLS_client_method();
SSL_CTX* ctx = SSL_CTX_new(method);
if (!ctx) {
@@ -42,9 +47,23 @@ static SSL_CTX* create_ssl_ctx(bool is_server, const char* cert, const char* key
return NULL;
}
SSL_CTX_set_min_proto_version(ctx, TLS1_2_VERSION);
if (SSL_CTX_set_min_proto_version(ctx, TLS1_2_VERSION) != 1) {
SSL_CTX_free(ctx);
return NULL;
}
if (SSL_CTX_set_cipher_list(ctx, "HIGH:!aNULL:!eNULL:!MD5:!RC4:!3DES") != 1) {
SSL_CTX_free(ctx);
return NULL;
}
if (cert && key) {
struct stat key_stat;
if (stat(key, &key_stat) != 0 || !S_ISREG(key_stat.st_mode) || key_stat.st_uid != geteuid() ||
(key_stat.st_mode & (S_IRGRP | S_IWGRP | S_IROTH | S_IWOTH))) {
log_message(LOG_LEVEL_ERROR, "TLS private key must be owned by the current user and private");
SSL_CTX_free(ctx);
return NULL;
}
if (SSL_CTX_use_certificate_file(ctx, cert, SSL_FILETYPE_PEM) <= 0) {
log_message(LOG_LEVEL_ERROR, "Failed to load certificate: %s", cert);
log_ssl_errors();
@@ -86,15 +105,22 @@ static SSL* wrap_fd_with_ssl(int fd, SSL_CTX* ctx, bool is_server, const char* h
log_message(LOG_LEVEL_ERROR, "Failed to create SSL object");
return NULL;
}
SSL_set_fd(ssl, fd);
if (SSL_set_fd(ssl, fd) != 1) {
SSL_free(ssl);
return NULL;
}
// Enable hostname verification for client connections when a hostname is provided.
// Must be done before SSL_connect to take effect during the handshake.
if (!is_server && hostname) {
SSL_set1_host(ssl, hostname);
if (SSL_set1_host(ssl, hostname) != 1) {
SSL_free(ssl);
return NULL;
}
}
// Retry SSL_accept/SSL_connect on WANT_READ/WANT_WRITE (non-blocking handshake)
time_t deadline = time(NULL) + (is_server ? tcp_get_timeout_sec() : tcp_get_contimeout_sec());
int ret;
do {
if (is_server)
@@ -104,7 +130,8 @@ static SSL* wrap_fd_with_ssl(int fd, SSL_CTX* ctx, bool is_server, const char* h
if (ret <= 0) {
int ssl_err = SSL_get_error(ssl, ret);
if (ssl_err == SSL_ERROR_WANT_READ || ssl_err == SSL_ERROR_WANT_WRITE)
if ((ssl_err == SSL_ERROR_WANT_READ || ssl_err == SSL_ERROR_WANT_WRITE) &&
time(NULL) < deadline)
continue;
log_message(LOG_LEVEL_ERROR, "SSL %s failed", is_server ? "accept" : "connect");
log_ssl_errors();
@@ -132,8 +159,10 @@ struct tls_child_ctx {
static void tls_child_fn(int fd, void* arg) {
struct tls_child_ctx* ctx = (struct tls_child_ctx*)arg;
SSL* ssl = wrap_fd_with_ssl(fd, ctx->ssl_ctx, true, NULL);
if (!ssl)
if (!ssl) {
io_set_ssl(NULL);
return;
}
io_set_ssl(ssl);
ctx->handler(fd);
SSL_shutdown(ssl);
@@ -149,57 +178,19 @@ 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,
const char* key_path, const char* ca_path) {
struct addrinfo hints;
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;
}
struct addrinfo* rp;
bool connected = false;
for (rp = result; rp != NULL; rp = rp->ai_next) {
if (!tcp_connect_socket(client, host, port)) {
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!");
client->file_descriptor = -1;
return false;
}
SSL_CTX* ctx = create_ssl_ctx(false, cert_path, key_path, ca_path);
if (!ctx)
if (!ctx) {
close(client->file_descriptor);
client->file_descriptor = -1;
return false;
}
client->ssl_ctx = ctx;
// Pass the server hostname for TLS hostname verification (SSL_set1_host
@@ -209,6 +200,8 @@ bool client_connect_tls(Client* client, char* host, int port, const char* cert_p
if (!ssl) {
SSL_CTX_free(ctx);
client->ssl_ctx = NULL;
close(client->file_descriptor);
client->file_descriptor = -1;
return false;
}
+152 -54
View File
@@ -7,62 +7,117 @@
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <stdint.h>
#include <sys/stat.h>
#include <unistd.h>
static int authorized_root_fd = -1;
static char* authorized_root_path;
bool utils_set_authorized_root(int fd, const char* canonical_path) {
char* path_copy = canonical_path ? str_dup(canonical_path) : NULL;
if (canonical_path && !path_copy) {
authorized_root_fd = -1;
free(authorized_root_path);
authorized_root_path = NULL;
return false;
}
authorized_root_fd = fd;
free(authorized_root_path);
authorized_root_path = path_copy;
return true;
}
void utils_set_authorized_root_fd(int fd) {
authorized_root_fd = fd;
(void)utils_set_authorized_root(fd, NULL);
}
static bool path_is_within_root(const char* root, const char* path) {
size_t root_len = strlen(root);
return strncmp(root, path, root_len) == 0 && (path[root_len] == '\0' || path[root_len] == '/');
}
static int open_authorized_destination(const char* dest_root) {
if (authorized_root_fd < 0 || !authorized_root_path || !dest_root ||
!path_is_within_root(authorized_root_path, dest_root))
return -1;
int dirfd = dup(authorized_root_fd);
if (dirfd < 0)
return -1;
const char* relative_path = dest_root + strlen(authorized_root_path);
while (*relative_path == '/')
relative_path++;
char* relative = str_dup(*relative_path ? relative_path : ".");
if (!relative) {
close(dirfd);
return -1;
}
char* saveptr = NULL;
char* component = strtok_r(relative, "/", &saveptr);
while (component) {
if (strcmp(component, "..") == 0) {
free(relative);
close(dirfd);
return -1;
}
if (strcmp(component, ".") == 0) {
component = strtok_r(NULL, "/", &saveptr);
continue;
}
int next = openat(dirfd, component, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
if (next < 0) {
free(relative);
close(dirfd);
return -1;
}
close(dirfd);
dirfd = next;
component = strtok_r(NULL, "/", &saveptr);
}
free(relative);
return dirfd;
}
bool mkdir_r(const char* path) {
size_t path_len = strlen(path);
char* path_duplicate = malloc(path_len + 1);
if (!path_duplicate)
if (!path || *path == '\0')
return false;
memcpy(path_duplicate, path, path_len + 1);
size_t capacity = path_len + 2;
char* path_current = (char*)malloc(capacity * sizeof(char));
if (!path_current) {
free(path_duplicate);
char* duplicate = str_dup(path);
if (!duplicate)
return false;
int dirfd = open(path[0] == '/' ? "/" : ".", O_RDONLY | O_DIRECTORY | O_CLOEXEC | O_NOFOLLOW);
if (dirfd < 0) {
free(duplicate);
return false;
}
char* path_current_position = path_current;
if (path[0] == '/') {
path_current[0] = '/';
path_current[1] = '\0';
path_current_position += 1;
} else {
path_current[0] = '\0';
}
const char* delimiter = "/";
char* saveptr;
const char* part = strtok_r(path_duplicate, delimiter, &saveptr);
bool ok = true;
while (part != NULL) {
size_t part_len = strlen(part);
if ((size_t)(path_current_position - path_current) + part_len + 2 > capacity) {
char* saveptr = NULL;
char* component = strtok_r(duplicate, "/", &saveptr);
while (component) {
if (strcmp(component, "..") == 0) {
ok = false;
break;
}
memcpy(path_current_position, part, part_len);
path_current_position += part_len;
path_current_position[0] = '/';
path_current_position[1] = '\0';
path_current_position++;
struct stat st;
if (stat(path_current, &st) != 0) {
if (mkdir(path_current, 0755) != 0) {
perror("Could not create directory");
if (strcmp(component, ".") != 0) {
int next = openat(dirfd, component, O_RDONLY | O_DIRECTORY | O_CLOEXEC | O_NOFOLLOW);
if (next < 0 && errno == ENOENT) {
if (mkdirat(dirfd, component, 0755) == 0 || errno == EEXIST)
next = openat(dirfd, component, O_RDONLY | O_DIRECTORY | O_CLOEXEC | O_NOFOLLOW);
}
if (next < 0) {
ok = false;
break;
}
close(dirfd);
dirfd = next;
}
part = strtok_r(NULL, delimiter, &saveptr);
component = strtok_r(NULL, "/", &saveptr);
}
free(path_duplicate);
free(path_current);
close(dirfd);
free(duplicate);
return ok;
}
@@ -143,7 +198,8 @@ static bool is_dir_in_manifest(const char* rel_path, ArrayList* manifest) {
return false;
}
static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifest) {
static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifest,
size_t max_delete, size_t* deleted_count) {
int scanfd = dup(dirfd);
if (scanfd < 0)
return false;
@@ -152,15 +208,20 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifes
close(scanfd);
return false;
}
bool all_removed = true;
bool operation_ok = true;
const struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
char* child_rel = path_cat((char*)rel_path, entry->d_name);
if (!child_rel) {
operation_ok = false;
continue;
}
struct stat st;
if (fstatat(dirfd, entry->d_name, &st, AT_SYMLINK_NOFOLLOW) != 0) {
if (errno != ENOENT)
operation_ok = false;
free(child_rel);
continue;
}
@@ -173,14 +234,24 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifes
int childfd = openat(dirfd, entry->d_name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool child_removed = false;
if (childfd >= 0) {
child_removed = delete_extras_fd(childfd, child_rel, manifest);
child_removed = delete_extras_fd(childfd, child_rel, manifest, max_delete, deleted_count);
if (!child_removed)
operation_ok = false;
close(childfd);
}
if (child_removed && !is_dir_in_manifest(child_rel, manifest) &&
unlinkat(dirfd, entry->d_name, AT_REMOVEDIR) != 0 && errno != ENOENT) {
} else if (errno != ENOENT) {
operation_ok = false;
} else if (!child_removed) {
all_removed = false;
}
if (child_removed && !is_dir_in_manifest(child_rel, manifest)) {
if (*deleted_count >= max_delete) {
operation_ok = false;
} else {
if (unlinkat(dirfd, entry->d_name, AT_REMOVEDIR) != 0) {
if (errno != ENOENT)
operation_ok = false;
} else {
(*deleted_count)++;
}
}
}
} else {
// Check if relative path is in manifest
@@ -192,38 +263,59 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, ArrayList* manifes
}
}
if (!found) {
if (unlinkat(dirfd, entry->d_name, 0) != 0 && errno != ENOENT)
if (*deleted_count >= max_delete) {
operation_ok = false;
free(child_rel);
continue;
}
if (unlinkat(dirfd, entry->d_name, 0) != 0) {
if (errno != ENOENT)
operation_ok = false;
} else {
(*deleted_count)++;
}
fprintf(stderr, " Deleted: %s\n", child_rel);
} else {
all_removed = false;
}
}
free(child_rel);
}
closedir(dir);
(void)all_removed;
return operation_ok;
}
bool delete_extras(const char* dest_root, ArrayList* manifest) {
int rootfd = authorized_root_fd >= 0
? dup(authorized_root_fd)
: open(dest_root, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
bool delete_extras_limited(const char* dest_root, ArrayList* manifest, size_t max_delete) {
if (!manifest)
return false;
int rootfd;
if (authorized_root_fd >= 0) {
if (authorized_root_path)
rootfd = open_authorized_destination(dest_root);
else if (dest_root == NULL)
rootfd = dup(authorized_root_fd);
else
rootfd = -1;
} else {
rootfd = open(dest_root, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC);
}
if (rootfd < 0)
return false;
bool ok = delete_extras_fd(rootfd, "", manifest);
size_t deleted_count = 0;
bool ok = delete_extras_fd(rootfd, "", manifest, max_delete, &deleted_count);
if (close(rootfd) != 0)
ok = false;
return ok;
}
bool delete_extras(const char* dest_root, ArrayList* manifest) {
return delete_extras_limited(dest_root, manifest, SIZE_MAX);
}
bool has_path_traversal(const char* path) {
if (!path)
return false;
return true;
char* dup = str_dup(path);
if (!dup)
return false;
return true;
char* saveptr;
const char* part = strtok_r(dup, "/", &saveptr);
while (part) {
@@ -237,6 +329,10 @@ bool has_path_traversal(const char* path) {
return false;
}
bool utils_valid_batch_path(const char* path) {
return path && path[0] != '\0' && path[0] != '/' && !has_path_traversal(path);
}
char* path_cat(const char* path1, const char* path2) {
if (path1 == NULL || *path1 == '\0')
return str_dup(path2);
@@ -251,6 +347,8 @@ char* path_cat(const char* path1, const char* path2) {
offset = 1;
path2_len -= 1;
}
if (path1_len > SIZE_MAX - path2_len - 2)
return NULL;
char* new_path = malloc(path1_len + path2_len + 2);
if (new_path == NULL)
return NULL;
+6
View File
@@ -2,6 +2,7 @@
#define UTILS_H
#include "array_list.h"
#include <stddef.h>
#include <stdbool.h>
bool mkdir_r(const char* path);
@@ -9,7 +10,12 @@ char* str_dup(const char* string);
char* path_cat(const char* path1, const char* path2);
bool glob_match(const char* pattern, const char* str);
bool delete_extras(const char* dest_root, ArrayList* manifest);
bool delete_extras_limited(const char* dest_root, ArrayList* manifest, size_t max_delete);
bool utils_set_authorized_root(int fd, const char* canonical_path);
/* The fd-only compatibility form is fail-closed for path-based operations;
* callers should use utils_set_authorized_root with the canonical identity. */
void utils_set_authorized_root_fd(int fd);
bool has_path_traversal(const char* path);
bool utils_valid_batch_path(const char* path);
#endif
+3 -1
View File
@@ -25,7 +25,9 @@ class ServerManager:
def start(self, extra_args=None):
self.stop()
self._port = _find_free_port()
cmd = SERVER_CMD + ["-p", str(self._port)]
# Plain TCP is intentionally explicit in the server; integration tests
# exercise that opt-in mode rather than relying on the secure default.
cmd = SERVER_CMD + ["-p", str(self._port), "--allow-unauthenticated"]
if extra_args:
cmd += extra_args
self._proc = subprocess.Popen(cmd, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL)
+16
View File
@@ -52,6 +52,22 @@ class TestArchiveMode:
assert not mismatches, f"Mismatch: {mismatches}"
class TestChmod:
def test_chmod_applies_to_transferred_files(self, shared_server):
clean_dir(DEST_DIR)
source_file = os.path.join(SOURCE_DIR, "small.txt")
os.chmod(source_file, 0o777)
result, dur = run_client(
SOURCE_DIR, DEST_DIR,
flags=["--chmod=u=rw,go=r"],
port=shared_server.port,
)
if result.returncode != 0:
pytest.fail(f"Exit {result.returncode}: {(result.stderr or result.stdout)[:200]}")
received = get_dest_received_dir(DEST_DIR, SOURCE_DIR)
assert (os.stat(os.path.join(received, "small.txt")).st_mode & 0o777) == 0o644
class TestExclude:
def test_exclude_single(self, shared_server):
clean_dir(DEST_DIR)
+3
View File
@@ -102,6 +102,7 @@ class TestTLSBasic:
with ServerManager() as server:
server.start(extra_args=[
"--tls", "--cert", certs["server_cert"], "--key", certs["server_key"],
"--ca", certs["ca"], "--client-cn", "fastsync-client",
])
result, dur = run_client(
SOURCE_DIR, DEST_DIR,
@@ -125,6 +126,7 @@ class TestTLSBasic:
with ServerManager() as server:
server.start(extra_args=[
"--tls", "--cert", certs["server_cert"], "--key", certs["server_key"],
"--ca", certs["ca"], "--client-cn", "fastsync-client",
])
result, dur = run_client(
SOURCE_DIR, DEST_DIR,
@@ -149,6 +151,7 @@ class TestTLSBasic:
with ServerManager() as server:
server.start(extra_args=[
"--tls", "--cert", certs["server_cert"], "--key", certs["server_key"],
"--ca", certs["ca"], "--client-cn", "fastsync-client",
])
result, dur = run_client(
SOURCE_DIR, DEST_DIR,
+4
View File
@@ -17,6 +17,8 @@ void test_array_list() {
// Test adding
int* val1 = malloc(sizeof(int));
if (!val1)
return;
*val1 = 42;
array_list_add(list, val1);
EXPECT_EQ_INT(list->size, 1);
@@ -26,6 +28,8 @@ void test_array_list() {
// Initial capacity is 100. Let's add 105 elements.
for (int i = 0; i < 105; i++) {
int* val = malloc(sizeof(int));
if (!val)
return;
*val = i;
array_list_add(list, val);
}
+3 -3
View File
@@ -11,7 +11,7 @@ static void test_file_operations() {
char* test_content = "Hello, Chunk System!";
unsigned long long test_len = strlen(test_content);
to_disk(test_path, test_content, test_len, false, false);
file_write_to_disk(test_path, test_content, test_len, false, false);
File* f = file_create(test_path);
EXPECT_NOT_NULL(f);
@@ -43,8 +43,8 @@ static void test_chunk_operations() {
char* content2 = "chunk item number 2";
unsigned long long len2 = strlen(content2);
to_disk(path1, content1, len1, false, false);
to_disk(path2, content2, len2, false, false);
file_write_to_disk(path1, content1, len1, false, false);
file_write_to_disk(path2, content2, len2, false, false);
struct stat st1, st2;
stat(path1, &st1);
+97 -2
View File
@@ -1,4 +1,6 @@
#include "test_client_cli.h"
#include "client_validation.h"
#include "chmod.h"
#include "config.h"
#include "test_utils.h"
#include "utils.h"
@@ -9,6 +11,58 @@
/* Declaration of parse_args from client_cli.c */
int parse_args(Config* config, int argc, char* argv[], int* positional_args, int* positional_count);
static Config* valid_client_config() {
Config* cfg = config_create();
if (!cfg)
return NULL;
cfg->send_directory = str_dup("/src");
cfg->receive_root_directory = str_dup("/dst");
return cfg;
}
static void test_validate_config_required_paths() {
Config* cfg = config_create();
EXPECT_FALSE(validate_config(cfg));
cfg->send_directory = str_dup("/src");
EXPECT_FALSE(validate_config(cfg));
config_delete(cfg);
}
static void test_validate_config_incompatible_options() {
Config* cfg = valid_client_config();
cfg->use_sendfile = true;
cfg->use_compression = true;
EXPECT_FALSE(validate_config(cfg));
cfg->use_compression = false;
cfg->use_incremental = true;
cfg->use_chunk_serialization = true;
EXPECT_FALSE(validate_config(cfg));
config_delete(cfg);
}
static void test_validate_config_tls_requirements() {
Config* cfg = valid_client_config();
cfg->use_tls = true;
EXPECT_FALSE(validate_config(cfg));
cfg->tls_cert = str_dup("cert.pem");
EXPECT_FALSE(validate_config(cfg));
cfg->tls_key = str_dup("key.pem");
EXPECT_FALSE(validate_config(cfg));
cfg->tls_ca = str_dup("ca.pem");
EXPECT_TRUE(validate_config(cfg));
config_delete(cfg);
}
static void test_validate_config_delta_sendfile_constraints() {
Config* cfg = valid_client_config();
cfg->use_delta = true;
EXPECT_FALSE(validate_config(cfg));
cfg->use_incremental = true;
cfg->use_sendfile = true;
EXPECT_FALSE(validate_config(cfg));
config_delete(cfg);
}
/* Test main() with --help flag (early return path, no server connection needed) */
static void test_cli_help() {
/* We can't easily call main() because it calls send_files which needs a server.
@@ -128,6 +182,40 @@ static void test_parse_args_valid_port() {
config_delete(cfg);
}
static void test_parse_args_chmod() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--chmod=u=rw,go=r", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0);
EXPECT_EQ_STR(cfg->chmod_spec, "u=rw,go=r");
EXPECT_TRUE(cfg->use_metadata);
mode_t result;
EXPECT_TRUE(chmod_apply(0777, cfg->chmod_spec, &result));
EXPECT_EQ_INT(result, 0644);
config_delete(cfg);
}
static void test_parse_args_numeric_chmod() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--chmod", "7777", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0);
EXPECT_EQ_STR(cfg->chmod_spec, "7777");
EXPECT_TRUE(cfg->use_metadata);
config_delete(cfg);
}
static void test_parse_args_rejects_invalid_chmod() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--chmod=a+X", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), -1);
config_delete(cfg);
}
/* Test parse_args rejects port > 65535 */
static void test_parse_args_invalid_port() {
Config* cfg = config_create();
@@ -256,8 +344,8 @@ static void test_parse_args_rejects_unimplemented_options() {
"--ipv4",
"--daemon",
"--config",
"--server",
"--compress-choice"};
"--server",
"--compress-choice"};
for (size_t i = 0; i < sizeof(options) / sizeof(options[0]); i++) {
Config* cfg = config_create();
@@ -287,6 +375,10 @@ static void test_parse_args_archive() {
}
void test_client_cli() {
test_validate_config_required_paths();
test_validate_config_incompatible_options();
test_validate_config_tls_requirements();
test_validate_config_delta_sendfile_constraints();
test_cli_help();
test_cli_archive_flags();
test_cli_dry_run();
@@ -295,6 +387,9 @@ void test_client_cli() {
test_parse_args_help();
test_parse_args_version();
test_parse_args_valid_port();
test_parse_args_chmod();
test_parse_args_numeric_chmod();
test_parse_args_rejects_invalid_chmod();
test_parse_args_invalid_port();
test_parse_args_non_numeric_port();
test_parse_args_invalid_server_port();
+4 -2
View File
@@ -13,6 +13,8 @@ static void test_data_compress_decompress_roundtrip() {
size_t len = strlen(original);
char* buf = malloc(len);
if (!buf)
return;
memcpy(buf, original, len);
Data* original_data = data_create(buf, len);
EXPECT_NOT_NULL(original_data);
@@ -62,8 +64,8 @@ static void test_chunk_compress_decompress_roundtrip() {
char* content2 = "chunk compression test file 2 with more data";
unsigned long long len2 = strlen(content2);
to_disk(path1, content1, len1, false, false);
to_disk(path2, content2, len2, false, false);
file_write_to_disk(path1, content1, len1, false, false);
file_write_to_disk(path2, content2, len2, false, false);
struct stat st1, st2;
EXPECT_EQ_INT(stat(path1, &st1), 0);
+36 -15
View File
@@ -224,26 +224,46 @@ static void test_config_send_receive_version_mismatch() {
}
}
static void test_is_remote_dest() {
static void test_config_receive_truncated() {
int p[2];
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
io_set_fds(p[0], p[0]);
io_set_bwlimit(0);
/* A valid prefix exercises cleanup after allocated wire strings and a
* partially received scalar field. */
EXPECT_TRUE(send_str(p[1], PROTOCOL_VERSION));
EXPECT_TRUE(send_str(p[1], "/src"));
EXPECT_TRUE(send_str(p[1], "/dst"));
EXPECT_TRUE(send_int(p[1], 1));
shutdown(p[1], SHUT_WR);
const Config* cfg = config_receive(p[0]);
EXPECT_NULL(cfg);
close(p[0]);
close(p[1]);
}
static void test_config_is_remote_dest() {
/* Valid SSH-style destinations */
EXPECT_TRUE(is_remote_dest("user@host:/path"));
EXPECT_TRUE(is_remote_dest("host:/path"));
EXPECT_TRUE(is_remote_dest("user@192.168.1.1:/remote/path"));
EXPECT_TRUE(config_is_remote_dest("user@host:/path"));
EXPECT_TRUE(config_is_remote_dest("host:/path"));
EXPECT_TRUE(config_is_remote_dest("user@192.168.1.1:/remote/path"));
/* Invalid destinations */
EXPECT_FALSE(is_remote_dest(NULL));
EXPECT_FALSE(is_remote_dest(""));
EXPECT_FALSE(is_remote_dest(":"));
EXPECT_FALSE(is_remote_dest("/local/path"));
EXPECT_FALSE(is_remote_dest("relative/path"));
EXPECT_FALSE(config_is_remote_dest(NULL));
EXPECT_FALSE(config_is_remote_dest(""));
EXPECT_FALSE(config_is_remote_dest(":"));
EXPECT_FALSE(config_is_remote_dest("/local/path"));
EXPECT_FALSE(config_is_remote_dest("relative/path"));
/* C:/windows/path is treated as remote (colon with no preceding slash) */
EXPECT_TRUE(is_remote_dest("C:/windows/path"));
EXPECT_TRUE(config_is_remote_dest("C:/windows/path"));
/* Edge cases */
EXPECT_FALSE(is_remote_dest("noslash"));
EXPECT_FALSE(is_remote_dest("/"));
EXPECT_TRUE(is_remote_dest("host:"));
EXPECT_TRUE(is_remote_dest("user@host:"));
EXPECT_FALSE(config_is_remote_dest("noslash"));
EXPECT_FALSE(config_is_remote_dest("/"));
EXPECT_TRUE(config_is_remote_dest("host:"));
EXPECT_TRUE(config_is_remote_dest("user@host:"));
}
void test_config() {
@@ -256,6 +276,7 @@ void test_config() {
if (!is_running_under_valgrind()) {
test_config_send_receive();
test_config_send_receive_version_mismatch();
test_config_receive_truncated();
}
test_is_remote_dest();
test_config_is_remote_dest();
}
+23 -20
View File
@@ -34,7 +34,8 @@ static void test_file_destroy_normal() {
static void test_file_load_data() {
const char* content = "Hello Load Test";
EXPECT_TRUE(to_disk("test_file_load_data.txt", content, strlen(content), false, false));
EXPECT_TRUE(
file_write_to_disk("test_file_load_data.txt", content, strlen(content), false, false));
struct stat st;
EXPECT_EQ_INT(stat("test_file_load_data.txt", &st), 0);
@@ -87,15 +88,16 @@ static void test_file_save_to_disk() {
rmdir("test_save_tmp");
}
static void test_to_disk_basic() {
const char* content = "Basic to_disk test";
EXPECT_TRUE(to_disk("test_to_disk_basic.txt", content, strlen(content), false, false));
static void test_file_write_to_disk_basic() {
const char* content = "Basic file_write_to_disk test";
EXPECT_TRUE(file_write_to_disk("test_file_write_to_disk_basic.txt", content, strlen(content),
false, false));
struct stat st;
EXPECT_EQ_INT(stat("test_to_disk_basic.txt", &st), 0);
EXPECT_EQ_INT(stat("test_file_write_to_disk_basic.txt", &st), 0);
EXPECT_EQ_INT((int)st.st_size, (int)strlen(content));
FILE* fp = fopen("test_to_disk_basic.txt", "rb");
FILE* fp = fopen("test_file_write_to_disk_basic.txt", "rb");
EXPECT_NOT_NULL(fp);
char buf[100];
size_t nread = fread(buf, 1, sizeof(buf), fp);
@@ -103,12 +105,13 @@ static void test_to_disk_basic() {
EXPECT_EQ_INT((int)nread, (int)strlen(content));
EXPECT_EQ_INT(memcmp(buf, content, strlen(content)), 0);
unlink("test_to_disk_basic.txt");
unlink("test_file_write_to_disk_basic.txt");
}
static void test_to_disk_creates_dirs() {
static void test_file_write_to_disk_creates_dirs() {
const char* content = "Nested dir test";
EXPECT_TRUE(to_disk("test_nested_tmp/nested/file.txt", content, strlen(content), false, false));
EXPECT_TRUE(file_write_to_disk("test_nested_tmp/nested/file.txt", content, strlen(content), false,
false));
struct stat st;
EXPECT_EQ_INT(stat("test_nested_tmp/nested/file.txt", &st), 0);
@@ -126,15 +129,15 @@ static void test_to_disk_creates_dirs() {
rmdir("test_nested_tmp");
}
static void test_to_disk_does_not_follow_symlink() {
const char* outside = "test_to_disk_outside.txt";
const char* link = "test_to_disk_link.txt";
static void test_file_write_to_disk_does_not_follow_symlink() {
const char* outside = "test_file_write_to_disk_outside.txt";
const char* link = "test_file_write_to_disk_link.txt";
const char* content = "confined";
unlink(outside);
unlink(link);
EXPECT_TRUE(to_disk(outside, "outside", 7, false, false));
EXPECT_TRUE(file_write_to_disk(outside, "outside", 7, false, false));
EXPECT_EQ_INT(symlink(outside, link), 0);
EXPECT_TRUE(to_disk(link, content, strlen(content), false, false));
EXPECT_TRUE(file_write_to_disk(link, content, strlen(content), false, false));
FILE* fp = fopen(outside, "rb");
char buf[16] = {0};
EXPECT_NOT_NULL(fp);
@@ -151,7 +154,7 @@ static void test_to_disk_does_not_follow_symlink() {
static void test_file_content_to_buffer() {
const char* content = "Buffer content test";
EXPECT_TRUE(to_disk("test_buffer_file.txt", content, strlen(content), false, false));
EXPECT_TRUE(file_write_to_disk("test_buffer_file.txt", content, strlen(content), false, false));
File* f = file_create("test_buffer_file.txt");
EXPECT_NOT_NULL(f);
@@ -273,7 +276,7 @@ static void test_file_send_no_path() {
}
static void test_file_metadata_create() {
EXPECT_TRUE(to_disk("test_meta_file.txt", "metadata test", 13, false, false));
EXPECT_TRUE(file_write_to_disk("test_meta_file.txt", "metadata test", 13, false, false));
struct stat st;
EXPECT_EQ_INT(stat("test_meta_file.txt", &st), 0);
@@ -382,7 +385,7 @@ static void test_file_send_single_calls_metadata_and_path() {
/* Create a real file on disk so we can have metadata */
const char* content = "File with metadata";
size_t len = strlen(content);
EXPECT_TRUE(to_disk("test_meta_send.txt", content, len, false, false));
EXPECT_TRUE(file_write_to_disk("test_meta_send.txt", content, len, false, false));
struct stat st;
EXPECT_EQ_INT(stat("test_meta_send.txt", &st), 0);
@@ -455,9 +458,9 @@ void test_file() {
test_file_load_data();
test_file_load_data_missing_file();
test_file_save_to_disk();
test_to_disk_basic();
test_to_disk_creates_dirs();
test_to_disk_does_not_follow_symlink();
test_file_write_to_disk_basic();
test_file_write_to_disk_creates_dirs();
test_file_write_to_disk_does_not_follow_symlink();
test_file_content_to_buffer();
test_file_save_to_disk_path_traversal();
test_file_save_to_disk_deep_traversal();
+4 -4
View File
@@ -15,7 +15,7 @@
static void test_sendfile_basic() {
const char* content = "Hello from sendfile test!";
size_t len = strlen(content);
EXPECT_TRUE(to_disk("test_sendfile_basic.txt", content, len, false, false));
EXPECT_TRUE(file_write_to_disk("test_sendfile_basic.txt", content, len, false, false));
File* file = file_create("test_sendfile_basic.txt");
EXPECT_NOT_NULL(file);
@@ -77,7 +77,7 @@ static void test_sendfile_basic() {
static void test_sendfile_empty_file() {
const char* content = "";
size_t len = 0;
EXPECT_TRUE(to_disk("test_sendfile_empty.txt", content, len, false, false));
EXPECT_TRUE(file_write_to_disk("test_sendfile_empty.txt", content, len, false, false));
File* file = file_create("test_sendfile_empty.txt");
EXPECT_NOT_NULL(file);
@@ -156,7 +156,7 @@ static void test_sendfile_missing_file() {
static void test_sendfile_compression_fallback() {
const char* content = "Compression fallback content";
size_t len = strlen(content);
EXPECT_TRUE(to_disk("test_sendfile_comp.txt", content, len, false, false));
EXPECT_TRUE(file_write_to_disk("test_sendfile_comp.txt", content, len, false, false));
struct stat st;
EXPECT_EQ_INT(stat("test_sendfile_comp.txt", &st), 0);
@@ -222,7 +222,7 @@ static void test_sendfile_compression_fallback() {
static void test_sendfile_no_path() {
const char* content = "No path sendfile test";
size_t len = strlen(content);
EXPECT_TRUE(to_disk("test_sendfile_nopath.txt", content, len, false, false));
EXPECT_TRUE(file_write_to_disk("test_sendfile_nopath.txt", content, len, false, false));
File* file = file_create("test_sendfile_nopath.txt");
EXPECT_NOT_NULL(file);
+1 -1
View File
@@ -101,7 +101,7 @@ static void test_fuzz_delta_deserialize() {
/* Smoke test for metadata_from_buf fuzz target */
static void test_fuzz_metadata_from_buf() {
/* Create a real file to get metadata from */
EXPECT_TRUE(to_disk("fuzz_meta_test.txt", "metadata test", 13, false, false));
EXPECT_TRUE(file_write_to_disk("fuzz_meta_test.txt", "metadata test", 13, false, false));
struct stat st;
EXPECT_EQ_INT(stat("fuzz_meta_test.txt", &st), 0);
+24 -1
View File
@@ -1,4 +1,5 @@
#include "test_metadata.h"
#include "chmod.h"
#include "metadata.h"
#include "protocol.h"
#include "test_utils.h"
@@ -125,7 +126,7 @@ static void test_metadata_rejects_invalid_values() {
static void test_file_restore_metadata() {
const char* path = "temp_meta_restore_test.txt";
const char* content = "test content";
EXPECT_TRUE(to_disk(path, content, strlen(content), false, false));
EXPECT_TRUE(file_write_to_disk(path, content, strlen(content), false, false));
FileMetadata m;
m.mode = 0644;
@@ -144,6 +145,27 @@ static void test_file_restore_metadata() {
unlink(path);
}
static void test_chmod_changes() {
mode_t result;
EXPECT_TRUE(chmod_apply(0777, "u=rw,go=r", &result));
EXPECT_EQ_INT(result, 0644);
EXPECT_TRUE(chmod_apply(0644, "a+x", &result));
EXPECT_EQ_INT(result, 0755);
result = 0777;
EXPECT_TRUE(chmod_apply(0777, "0000", &result));
EXPECT_EQ_INT(result, 0000);
result = 0777;
EXPECT_TRUE(chmod_apply(0777, "7777", &result));
EXPECT_EQ_INT(result, 07777);
result = 0777;
EXPECT_TRUE(chmod_apply(0777, "755", &result));
EXPECT_EQ_INT(result, 0755);
EXPECT_FALSE(chmod_apply(0777, "888", &result));
EXPECT_FALSE(chmod_apply(0777, "10000", &result));
EXPECT_FALSE(chmod_apply(0777, "a+X", &result));
EXPECT_FALSE(chmod_apply(0777, "a+r,", &result));
}
void test_metadata() {
test_metadata_to_from_buf_roundtrip();
test_metadata_to_buf_null();
@@ -152,4 +174,5 @@ void test_metadata() {
test_metadata_send_null();
test_metadata_rejects_invalid_values();
test_file_restore_metadata();
test_chmod_changes();
}
+5 -1
View File
@@ -14,6 +14,8 @@
static Data* random_data(int min_size, int max_size) {
int size = min_size + rand() % (max_size - min_size + 1);
char* buf = malloc(size);
if (!buf)
return NULL;
for (int i = 0; i < size; i++)
buf[i] = (char)(rand() % 256);
return data_create(buf, size);
@@ -81,10 +83,12 @@ static void test_property_chunk_roundtrip() {
int content_len = 1 + rand() % 4096;
char* content = malloc(content_len);
if (!content)
return;
for (int i = 0; i < content_len; i++)
content[i] = (char)(rand() % 256);
to_disk(path, content, content_len, false, false);
file_write_to_disk(path, content, content_len, false, false);
struct stat st;
stat(path, &st);
+18
View File
@@ -38,6 +38,23 @@ static void test_send_receive_n_data_zero() {
close(p[1]);
}
static void test_explicit_session_context() {
int p[2];
EXPECT_EQ_INT(pipe(p), 0);
ProtocolSession session;
protocol_session_init(&session, p[0], p[1]);
protocol_session_set_bwlimit(&session, 0);
const char payload[] = "explicit context";
char received[sizeof(payload)] = {0};
EXPECT_TRUE(protocol_send_n_data(&session, payload, sizeof(payload)));
EXPECT_TRUE(protocol_receive_n_data(&session, received, sizeof(received)));
EXPECT_EQ_INT(memcmp(payload, received, sizeof(payload)), 0);
close(p[0]);
close(p[1]);
}
static void test_send_receive_str() {
int p[2];
EXPECT_EQ_INT(pipe(p), 0);
@@ -173,6 +190,7 @@ static void test_receive_str_truncated() {
void test_protocol() {
test_send_receive_n_data();
test_send_receive_n_data_zero();
test_explicit_session_context();
test_send_receive_str();
test_send_receive_str_normal();
test_send_receive_data();
+4
View File
@@ -122,6 +122,8 @@ static void test_queue_destroyer() {
for (int i = 0; i < 3; i++) {
int* val = malloc(sizeof(int));
if (!val)
break;
*val = i;
queue_enqueue(q, val);
}
@@ -181,6 +183,8 @@ static void test_queue_multithreaded() {
for (int i = 1; i <= 100; i++) {
int* val = malloc(sizeof(int));
if (!val)
break;
*val = i;
queue_enqueue_multithreaded(q, val, &mutex, &cnd_empty, &cnd_full);
}
+1 -1
View File
@@ -13,7 +13,7 @@
static void test_chunk_deserialize_truncated() {
char* path = "test_rob_trunc.txt";
char* content = "hello";
to_disk(path, content, strlen(content), false, false);
file_write_to_disk(path, content, strlen(content), false, false);
struct stat st;
stat(path, &st);
+31 -1
View File
@@ -7,7 +7,7 @@
#include <unistd.h>
static void create_test_file(const char* path, const char* content) {
(void)to_disk(path, content, strlen(content), false, false);
(void)file_write_to_disk(path, content, strlen(content), false, false);
}
static void test_scanner_single_file() {
@@ -385,6 +385,35 @@ static void test_scanner_no_patterns() {
rmdir(dir);
}
static void test_parallel_scanner_root_chunks_without_workers() {
const char* dir = "test_parallel_scan_root";
const char* file1 = "test_parallel_scan_root/a.txt";
const char* file2 = "test_parallel_scan_root/b.txt";
EXPECT_EQ_INT(mkdir(dir, 0755), 0);
create_test_file(file1, "a");
create_test_file(file2, "b");
ScannerOptions options = {false, 1, NULL, 0, NULL, 0, 0, 0,
0, 0, false, false, false, false, false};
ParallelScanner* scanner = parallel_scanner_create_with_options(dir, &options);
EXPECT_NOT_NULL(scanner);
int total_files = 0;
Chunk* chunk;
while ((chunk = parallel_scanner_next(scanner)) != NULL) {
total_files += chunk->element_count;
chunk_destroy(chunk);
}
EXPECT_EQ_INT(total_files, 2);
EXPECT_FALSE(parallel_scanner_failed(scanner));
parallel_scanner_destroy(scanner);
unlink(file1);
unlink(file2);
rmdir(dir);
}
void test_scanner() {
test_scanner_single_file();
test_scanner_multiple_files();
@@ -399,4 +428,5 @@ void test_scanner() {
test_scanner_size_range();
test_scanner_mixed_patterns();
test_scanner_no_patterns();
test_parallel_scanner_root_chunks_without_workers();
}
+4 -8
View File
@@ -11,11 +11,7 @@
#include <sys/wait.h>
#include <unistd.h>
/* Include server.c but rename main to avoid conflict with test runner's main */
#define main server_main_
#define FASTSYNC_SERVER_AS_LIB
#include "server.c"
#undef main
#include "receiver.h"
/* Test receive_files with immediate FINISHED status */
static void test_receive_files_finished() {
@@ -36,7 +32,7 @@ static void test_receive_files_finished() {
/* Child: use p[0] for both read and write */
close(p[1]);
io_set_fds(p[0], p[0]);
int ret = receive_files(cfg, p[0]);
int ret = receiver_receive_files(cfg, p[0]);
close(p[0]);
config_delete(cfg);
_exit(ret == 0 ? 0 : 1);
@@ -88,7 +84,7 @@ static void test_receive_files_single_file() {
/* Child: use p[0] for both read and write */
close(p[1]);
io_set_fds(p[0], p[0]);
int ret = receive_files(cfg, p[0]);
int ret = receiver_receive_files(cfg, p[0]);
close(p[0]);
config_delete(cfg);
_exit(ret == 0 ? 0 : 1);
@@ -149,7 +145,7 @@ static void test_receive_files_abort() {
if (pid == 0) {
close(p[1]);
io_set_fds(p[0], p[0]);
int ret = receive_files(cfg, p[0]);
int ret = receiver_receive_files(cfg, p[0]);
close(p[0]);
config_delete(cfg);
/* Should return -1 on abort */
+6
View File
@@ -32,6 +32,8 @@ static int mpmc_producer_func(void* arg) {
ProducerCtx* ctx = (ProducerCtx*)arg;
for (int i = 1; i <= ITEMS_PER_PRODUCER; i++) {
int* val = malloc(sizeof(int));
if (!val)
return thrd_error;
*val = ctx->producer_id * ITEMS_PER_PRODUCER + i;
queue_enqueue_multithreaded(ctx->q, val, ctx->mutex, ctx->cnd_empty, ctx->cnd_full);
}
@@ -121,6 +123,8 @@ static int bp_producer_func(void* arg) {
BackpressureCtx* ctx = (BackpressureCtx*)arg;
for (int i = 0; i < 5; i++) {
int* val = malloc(sizeof(int));
if (!val)
return thrd_error;
*val = i + 1;
queue_enqueue_multithreaded(ctx->q, val, ctx->mutex, ctx->cnd_empty, ctx->cnd_full);
ctx->items_sent++;
@@ -194,6 +198,8 @@ static void test_queue_rapid_create_destroy() {
for (int j = 0; j < 3; j++) {
int* val = malloc(sizeof(int));
if (!val)
break;
*val = j;
queue_enqueue(q, val);
}