Author SHA1 Message Date
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 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 6b7213db5c docs: align compatibility guidance
CI / lint (pull_request) Failing after 2s
CI / build-and-test (pull_request) Has been skipped
CI / sanitizers (address) (pull_request) Has been skipped
CI / sanitizers (undefined) (pull_request) Has been skipped
CI / fuzz-build (pull_request) Has been skipped
CI / coverage (pull_request) Has been skipped
CI / valgrind (pull_request) Has been skipped
2026-08-15 13:41:23 +02:00
TapTap 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 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 de537810e5 docs: fix server destination root examples
CI / lint (pull_request) Failing after 2s
CI / build-and-test (pull_request) Has been skipped
CI / sanitizers (address) (pull_request) Has been skipped
CI / sanitizers (undefined) (pull_request) Has been skipped
CI / fuzz-build (pull_request) Has been skipped
CI / coverage (pull_request) Has been skipped
CI / valgrind (pull_request) Has been skipped
2026-08-15 13:14:28 +02:00
TapTap 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 def58554d9 docs: clarify rsync compatibility roadmap
CI / lint (pull_request) Failing after 2s
CI / build-and-test (pull_request) Has been skipped
CI / sanitizers (address) (pull_request) Has been skipped
CI / sanitizers (undefined) (pull_request) Has been skipped
CI / fuzz-build (pull_request) Has been skipped
CI / coverage (pull_request) Has been skipped
CI / valgrind (pull_request) Has been skipped
2026-08-15 12:54:38 +02:00
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 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
43 changed files with 1811 additions and 1549 deletions

No files matched your search

+1 -1
View File
@@ -58,7 +58,7 @@ endif()
find_package(OpenSSL REQUIRED) find_package(OpenSSL REQUIRED)
file(GLOB SHARED_SRCS "src/shared/*.c") file(GLOB SHARED_SRCS "src/shared/*.c")
set(FILE_STORE_SRCS src/shared/file_store.c) set(FILE_STORE_SRCS "${CMAKE_CURRENT_SOURCE_DIR}/src/shared/file_store.c")
list(REMOVE_ITEM SHARED_SRCS ${FILE_STORE_SRCS}) list(REMOVE_ITEM SHARED_SRCS ${FILE_STORE_SRCS})
file(GLOB SERVER_SRCS "src/server/*.c") file(GLOB SERVER_SRCS "src/server/*.c")
set(SERVER_RECEIVER_SRCS src/server/receiver.c) set(SERVER_RECEIVER_SRCS src/server/receiver.c)
+317 -297
View File
@@ -1,363 +1,383 @@
# FastSync # FastSync
A high-performance file synchronization system with SSH and TCP transport, TLS encryption, streaming zstd compression, multithreaded transfer, incremental sync, metadata preservation, and rsync-compatible CLI flags. FastSync is a high-performance file synchronization tool designed to become a
drop-in replacement for common `rsync` workflows. It keeps the familiar
source/destination model and rsync-style options while adding optional
multithreading, streaming zstd compression, chunking, zero-copy TCP transfers,
and native TCP/TLS transports.
## Technical Overview The compatibility target is straightforward:
1. **Dual transport**: custom TCP client-server or SSH subprocess (rsync-style `user@host:/path`) - Existing rsync commands should keep the same meaning.
2. **TLS encryption**: OpenSSL-based TLS 1.2+ for encrypted TCP connections with optional CA verification - FastSync-only performance options should be additive and optional.
3. **Chunked file transfer**: files grouped into configurable-size chunks (default ~10 MB) - A normal compatibility-mode transfer should prioritize rsync filesystem
4. **Streaming zstd compression** (levels 1–22) using `ZSTD_compressStream2` semantics over maximum throughput.
5. **Multithreading**: producer-consumer pipeline with thread-safe queues (scanner → loader → sender)
6. **Incremental sync**: skip files unchanged since last transfer (compares size + mtime)
7. **Batch incremental**: send incremental checks in batched groups for reduced round-trips
8. **Metadata preservation**: file mode and mtime are restored when enabled; ownership and atime are intentionally not restored
9. **`sendfile()` zero-copy** on TCP (~2× faster on loopback)
10. **SSH ControlMaster** for connection reuse across repeated invocations
11. **Bandwidth limiting**: token-bucket throttling (`--bwlimit`)
12. **`--delete`**: receiver removes files not present in sender manifest
13. **`--exclude` / `--include`**: glob-pattern filename filtering
14. **Path traversal protection**: `..` sequences in file paths are rejected automatically
15. **Connection limits**: server enforces maximum concurrent connections (default 100)
16. **Keep-alive**: periodic `STATUS_KEEPALIVE` messages detect stalled connections
17. **Abort handling**: `SIGINT` sends `STATUS_ABORT` for clean server-side teardown
18. **Atomic writes**: received files are written to a temporary name then atomically renamed
19. **Backup mode**: `--backup` preserves overwritten files with optional `--backup-dir`
20. **Log file**: `--log-file` redirects log output to a file instead of stderr
21. **Transfer statistics**: `--stats` prints summary of transferred bytes, files, and timing
## System Architecture FastSync currently speaks its own protocol to `fastsync-server`. SSH mode
starts that server remotely; it does not yet interoperate with an unmodified
rsync client or rsync daemon. See [Compatibility Status](#compatibility-status)
for the current boundary.
### Client ## Why FastSync
- Recursively scans source directories (BFS), supports exclude and include patterns
- Groups files into chunks (configurable size)
- Streaming zstd compression with configurable level
- Chunk serialization (compact binary format) or per-file transfer
- Incremental transfer: sends file metadata to server, skips unchanged files
- Batch incremental: groups incremental checks to minimize round-trips
- Manifests all sent paths when `--delete` is active
- Sends via TCP `sendfile()` or SSH pipe
- Optional progress display with throughput
- Bandwidth limiting via token-bucket algorithm
- Configurable I/O and connection timeouts (`--timeout`, `--contimeout`)
- Backup overwritten files (`--backup`) with optional directory (`--backup-dir`)
- Transfer statistics summary (`--stats`)
- Maximum directory depth control (`--max-depth`)
- Log file output (`--log-file`)
- Exclude patterns from file (`--exclude-from`)
### Server FastSync uses a producer-consumer transfer pipeline and can combine several
- TCP mode: listens on configurable port (default 8080); SSH mode: runs via `--stdio` optimizations for large or high-latency transfers:
- TLS mode: wraps TCP connections with OpenSSL with optional CA verification
- Receives and reassembles files
- Decompresses (streaming zstd), deserializes, restores metadata
- Handles incremental checks: compares size + mtime against destination files
- Handles batch incremental checks for reduced round-trips
- Processes `STATUS_MANIFEST` for `--delete`: walks destination tree, removes extras
- Per-connection concurrency via `fork()` with configurable connection limit (default 100)
- Thread pool for parallel processing
- Atomic writes: files written to `.tmp` path then atomically renamed on success
- Abort handling: cleanly shuts down on `STATUS_ABORT` from client
- Path traversal protection: rejects file paths containing `..`
## Protocol Details - Multithreaded scanning, loading, and sending.
- Streaming zstd compression with levels 1 through 22.
- Configurable file chunking and compact chunk serialization.
- `sendfile()` zero-copy transfers over TCP.
- Batched incremental checks to reduce round trips.
- Optional block-level delta transfer for FastSync peers.
- Bandwidth limiting, progress reporting, statistics, and backups.
- TCP, SSH, and TLS transports.
- Atomic temporary-file writes by default.
### Status Codes These optimizations are disabled or selected independently. Users can start
| Code | Meaning | with rsync-style commands and add FastSync options when they are useful.
|------|---------|
| `STATUS_OK` | Operation successful |
| `STATUS_ERROR` | Error occurred |
| `STATUS_FINISHED` | Transfer complete |
| `STATUS_NEXT` | Ready for next file (per-file mode) |
| `STATUS_CHUNK` | Following data is a serialized chunk |
| `STATUS_MANIFEST` | Following data is a file manifest (for `--delete`) |
| `STATUS_CHECK` | Incremental check: client sends file path + size + mtime and, when negotiated, checksum; server responds with OK (skip) or NEXT (send) |
| `STATUS_CHECK_BATCH` | Batch incremental check: multiple file checks sent in one message |
| `STATUS_KEEPALIVE` | Keep-alive heartbeat to detect stalled connections |
| `STATUS_ABORT` | Abort signal: client interrupts, server cleans up and exits |
| `STATUS_DELTA_SIGNATURE` | Delta sync: following data is a file signature (rsync-style rolling hash) |
| `STATUS_DELTA_DATA` | Delta sync: following data is a delta patch for a file |
### Wire Format — Metadata ## Compatibility Status
When `use_metadata` is enabled (`-M`), each file entry carries a 4-byte `present` flag followed by five fields (`mode`, `uid`, `gid`, `mtime_sec`, `mtime_nsec`). When disabled globally, no metadata bytes are sent — zero wire overhead. FastSync is currently an rsync-compatible CLI in progress, not a complete
replacement for every rsync feature or protocol mode.
### Transfer Flow ### Working today
```
Config → (STATUS_NEXT | STATUS_CHUNK | STATUS_CHECK | STATUS_CHECK_BATCH)* → [STATUS_MANIFEST] → STATUS_FINISHED → STATUS_OK
```
Keep-alive (`STATUS_KEEPALIVE`) may be sent at any point during the transfer. The receiver resets its inactivity timer on receipt. If no data arrives within the receive timeout, the connection is aborted. - Recursive directory scanning.
- Rsync-style source and destination arguments.
- SSH transport using `user@host:destination` paths below the remote authorized root.
- TCP client/server transfers.
- Dry runs, excludes, includes, size filters, backups, statistics, and
bandwidth limiting.
- Incremental size/mtime checks and optional xxHash64 content checks.
- FastSync-native delta transfer for changed files.
- Optional mode and timestamp preservation.
- Delete manifests with server-side delete authorization.
- Temporary-file writes with atomic rename by default.
- Path traversal checks and destination-root confinement.
Abort (`STATUS_ABORT`) may be sent at any point. On receipt the server cleans up temporary files and exits the child process. ### Not yet equivalent to rsync
### Protocol Version - The FastSync wire protocol is not the rsync wire protocol.
- SSH mode requires `fastsync-server` on the remote host.
- Archive mode does not yet provide all of rsync's `-rlptgoD` behavior.
- Symlink transfer is incomplete; link targets are not yet recreated in all
modes.
- Owner/group, ACL, xattr, hard-link, device, and special-file handling is
incomplete or unavailable.
- Sparse-file handling does not yet preserve all holes correctly.
- `--partial`, `--partial-dir`, `--append`, and `--append-verify` are not yet
full rsync-style resumable transfers.
- Several rsync short options currently have FastSync-specific meanings. Do
not assume every short option is interchangeable yet.
`2.2.0` — server and client must match. This version adds a 64-bit XXH64 checksum to checksum-enabled `STATUS_CHECK` messages and validates the negotiated compression choice (`zstd` or `none`). Older clients and servers must not be mixed with this version; mismatch results in `STATUS_ERROR`. The detailed flag matrix is maintained in
[`RSYNC_COMPAT.md`](RSYNC_COMPAT.md). It distinguishes implemented,
partial, alternate, and planned behavior.
Config negotiation is sender-driven: the client serializes transfer options and the server applies them while receiving and writing files. `--checksum` compares size and content checksum instead of timestamps. `--compress-choice zstd` enables zstd; `none` disables it. Unsupported choices are rejected during config exchange. ## Quick Start
## Command-Line Arguments ### Build
### Client Requirements: C11 compiler, CMake 3.22 or newer, xxHash, zstd, OpenSSL,
pthreads, and an SSH client for SSH transport. The first CMake configure fetches
| Argument | Description | xxHash from GitHub, so network access is required unless the dependency is
|----------|-------------| already cached.
| Positional | `<source> <dest>` — automatic SSH detection if dest contains `:` |
| `-c [level]` | Compression with optional level (1–22, default 5) |
| `-z [level]` | Alias for `-c` |
| `-a, --archive` | Archive mode: enables `-c -m -M` (no `-s`) |
| `-m` | Multithreading mode |
| `-s` | Chunk serialization (batch all files per chunk) |
| `-f, --sendfile` | Sendfile zero-copy. Incompatible with `-c` / `-s`. TCP only. |
| `-M, --preserve` | Preserve supported file metadata (mode and mtime; ownership and atime are unsupported) |
| `-n, --dry-run` | Scan and print what would be transferred |
| `-p <port>` | SSH port (default: 22) |
| `-v, --verbose` | Enable debug logging |
| `--progress` | Show real-time transfer speed |
| `--delete` | Delete files on receiver not present in source |
| `--exclude <pattern>` | Exclude files matching glob pattern (repeatable) |
| `--exclude-from <file>` | Read exclude patterns from a file (one per line) |
| `--include <pattern>` | Only transfer files matching glob pattern (repeatable, whitelist) |
| `--max-size <n>` | Skip files larger than n bytes |
| `--min-size <n>` | Skip files smaller than n bytes |
| `--incremental` | Skip files unchanged since last transfer (size + mtime). Auto-enables `--preserve`. Incompatible with `-s`. |
| `--bwlimit <KB/s>` | Bandwidth limit in kilobytes per second |
| `--chunk-size <n>` | Chunk size in bytes (default: 10485760) |
| `--timeout <sec>` | I/O timeout in seconds (default: 30) |
| `--contimeout <sec>` | Connection timeout in seconds (default: 10) |
| `--backup` | Backup existing destination files before overwriting |
| `--backup-dir <dir>` | Target directory for backups (requires `--backup`) |
| `--stats` | Print transfer statistics at end (bytes, files, timing) |
| `--max-depth <n>` | Maximum directory depth to recurse (0 = unlimited, default: 0) |
| `--log-file <path>` | Write log messages to file instead of stderr |
| `--source-dir <path>` | Source directory (overrides `FASTSYNC_SOURCE_DIR`) |
| `--dest-dir <path>` | Server destination directory (overrides `FASTSYNC_DEST_DIR`) |
| `--save-to-disk` | Write received files to disk |
| `--server-host <ip>` | Server IP address (default: `127.0.0.1`) |
| `--server-port <n>` | Server port (default: `8080`) |
| `--tls` | Enable TLS encryption |
| `--cert <path>` | TLS certificate file (PEM) |
| `--key <path>` | TLS private key file (PEM) |
| `--ca <path>` | TLS CA certificate file for verification (PEM) |
### Server
| Argument | Description |
|----------|-------------|
| `--stdio` | Run in stdio mode (for SSH transport; single connection then exits) |
| `-p <port>` | TCP listen port (default: 8080, range: 1–65535) |
| `--tls` | Enable TLS encryption |
| `--cert <path>` | TLS certificate file (PEM) |
| `--key <path>` | TLS private key file (PEM) |
| `--ca <path>` | TLS CA certificate file for verification (PEM) |
| `-v, --verbose` | Enable debug logging |
| `--help` | Show help |
## Environment Variables
| Variable | Default | Description |
|----------|---------|-------------|
| `FASTSYNC_SOURCE_DIR` | — | Source directory fallback |
| `FASTSYNC_DEST_DIR` | — | Destination directory fallback |
| `FASTSYNC_SAVE_TO_DISK` | `false` | Disk persistence fallback |
| `FASTSYNC_SSH_PORT` | `22` | Default SSH port |
| `FASTSYNC_SERVER_HOST` | `127.0.0.1` | Default server host |
| `FASTSYNC_SERVER_PORT` | `8080` | Default server port |
| `FASTSYNC_TLS_CERT` | — | Default TLS certificate path |
| `FASTSYNC_TLS_KEY` | — | Default TLS private key path |
| `FASTSYNC_TLS_CA` | — | Default TLS CA certificate path |
## Implementation Details
### Data Structures
1. **Chunk** — collection of files (~10 MB total by default)
2. **File** — path, content (`Data`), optional `FileMetadata` pointer
3. **FileMetadata** — `mode`, `uid`, `gid`, `mtime_sec`, `mtime_nsec`; uid/gid are advisory wire fields and are never applied by the receiver; atime is unsupported
4. **Config** — runtime parameters (transported over wire, TLS settings excluded). Includes `timeout`, `contimeout`, `quiet`, `backup`, `backup_dir`, `stats`, `max_depth`, `log_file`, `queue_size`.
5. **Queue** — thread-safe bounded queue with condition variables
6. **DirectoryScanner** — recursive BFS traversal with exclude and include pattern support, max-depth enforcement
### Key Algorithms
1. **File scanning** — BFS directory traversal; entries matched against exclude and include patterns, max-depth enforced
2. **Chunking** — files accumulated until `chunk_size` threshold, then flushed
3. **Compression** — streaming zstd via `ZSTD_compressStream2` / `ZSTD_decompressStream`
4. **Network protocol** — status-code-driven exchange with metadata packing, keep-alive, and abort support
5. **Incremental check** — client sends `STATUS_CHECK` + path + size + mtime and, with `--checksum`, XXH64 content checksum; server compares against destination. Can be batched via `STATUS_CHECK_BATCH` for reduced round-trips.
6. **Bandwidth limiting** — token-bucket algorithm with `nanosleep` throttling on 64 KB write chunks
7. **Metadata restoration** — `chmod()`, `chown()`, `utimensat()` on the receiving side
8. **`--delete`** — sender tracks all sent paths; receiver walks destination tree and removes unlisted files/directories
9. **SSH transport** — `socketpair()` + `fork()` + `execvp("ssh", ...)` with `ControlMaster` and port support
10. **TLS transport** — OpenSSL `SSL_CTX` with TLS 1.2 minimum, optional CA verification, transparent `SSL_read`/`SSL_write` via `io_set_ssl()`
11. **Path traversal protection** — `has_path_traversal()` rejects any file path containing `..` components, preventing directory escape attacks
12. **Connection limiting** — server tracks active connections and rejects new ones beyond `max_connections` (default 100)
13. **Keep-alive** — idle connections receive periodic `STATUS_KEEPALIVE` to detect half-open TCP connections
14. **Abort handling** — `SIGINT` sets an abort flag; the next protocol operation sends `STATUS_ABORT` for clean server cleanup
15. **Atomic writes** — files are written to a `.tmp` suffix then atomically renamed via `rename()`, preventing partial files
16. **Backup** — before overwriting, existing files are moved to `--backup-dir` (or same directory with `~` suffix) preserving the original
## Security Features
### Path Traversal Protection
All received file paths are validated by `has_path_traversal()` before any disk operation. Any path containing `..` components is rejected with `STATUS_ERROR`, preventing directory escape attacks.
### TLS Certificate Verification
When `--ca` is provided, the server performs mutual TLS verification (`SSL_VERIFY_PEER` with depth 4). Without `--ca`, TLS is still encrypted but peer certificates are not verified.
### Connection Limits
The server enforces a maximum of 100 concurrent connections (configurable via `max_connections` in `Server`). When the limit is reached, new connections are immediately rejected and closed.
### Abort Handling
If the client receives `SIGINT` (Ctrl+C) during a transfer, it sends `STATUS_ABORT` to the server. The server then cleans up temporary files and exits the child process, preventing incomplete files from remaining on disk.
### Atomic Writes
Received files are written to a temporary path (suffixed with `.tmp`) and then atomically renamed to the final filename via `rename()`. This prevents partial or corrupted files from appearing at the destination if the transfer is interrupted.
## Build Requirements
- C11 compiler
- CMake >= 3.22
- zstd library
- OpenSSL (development headers and libraries)
- pthreads
- SSH client (for SSH transport mode only)
### Installing Dependencies
**Ubuntu/Debian:**
```bash
sudo apt install cmake build-essential libzstd-dev libssl-dev openssh-client
```
**Nix:**
```bash
nix-shell # provides zstd, openssl, cmake, gcc
```
## Building
```bash ```bash
cmake -B build -S . && cmake --build build -j$(nproc) cmake -B build -S .
cmake --build build -j$(nproc)
``` ```
## Running With Nix:
### Server (TCP mode)
```bash ```bash
./build/server nix-shell
cmake -B build -S .
cmake --build build -j$(nproc)
``` ```
### Server with TLS ### SSH transfer
The remote host must have `fastsync-server` available in `PATH`, or use
`--fastsync-server-path`. SSH starts `fastsync-server --stdio` in its remote
working directory, so use a destination below that directory unless the
remote server is otherwise configured with a matching authorized root.
```bash ```bash
./build/server --tls --cert server.pem --key server-key.pem ssh user@host 'mkdir -p destination'
./build/client /path/to/source user@host:destination
``` ```
### Server via SSH ### TCP transfer
Place the `fastsync-server` binary in the remote `$PATH`. The client runs `ssh user@host fastsync-server --stdio` automatically when an SSH-style destination is given.
Start the FastSync server:
### Client — SSH (rsync-style)
```bash ```bash
./build/client /path/to/send user@host:/path/to/receive ./build/server --destination-root /path/to -p 8080
``` ```
### Client — TCP Then run the client:
```bash ```bash
./build/client --source-dir /path/to/send --dest-dir /path/to/receive --save-to-disk ./build/client --server-host 127.0.0.1 --server-port 8080 \
--source-dir /path/to/source --dest-dir /path/to/destination \
--save-to-disk
``` ```
### Client — TCP with TLS ### TLS transfer
```bash ```bash
./build/server --destination-root /path/to --tls --cert server.pem --key server-key.pem -p 8443
./build/client --tls --cert client.pem --key client-key.pem --ca ca.pem \ ./build/client --tls --cert client.pem --key client-key.pem --ca ca.pem \
--source-dir /path/to/send --dest-dir /path/to/receive --save-to-disk --server-host example.com --server-port 8443 \
--source-dir /path/to/source --dest-dir /path/to/destination \
--save-to-disk
``` ```
### Common Options ## Common Workflows
These examples show the intended rsync-style workflow. Options marked as
FastSync-native are optional performance or transport extensions.
```bash ```bash
# Archive mode (compression + multithreading + metadata) # Basic synchronization
./build/client -a /path/to/send user@host:/path ./build/client /source/ /destination/
# Dry run # Archive-style synchronization (current FastSync archive behavior)
./build/client -n /path/to/send /path/to/receive ./build/client -a /source/ user@host:destination/
# With progress and custom chunk size # Preview a transfer without changing the destination
./build/client --progress --chunk-size 2097152 /src user@host:/dst ./build/client -n /source/ /destination/
# Exclude temporary files + delete extras on receiver # Exclude temporary and object files
./build/client --exclude "*.tmp" --exclude "*.o" --delete /src user@host:/dst ./build/client --exclude '*.tmp' --exclude '*.o' \
/source/ user@host:destination/
# Incremental sync (skip unchanged files) # Remove destination entries not present in the source
./build/client --incremental /src user@host:/dst ./build/client --delete /source/ user@host:destination/
# Bandwidth limit to 1 MB/s # Skip unchanged files using size and modification time
./build/client --bwlimit 1024 /src user@host:/dst ./build/client --incremental /source/ user@host:destination/
# With timeouts and stats # Verify content when size and time are not sufficient
./build/client --timeout 60 --contimeout 15 --stats /src user@host:/dst ./build/client --incremental --checksum /source/ user@host:destination/
# Backup overwritten files to a directory # Preserve supported mode and timestamp metadata
./build/client --backup --backup-dir /backups /src user@host:/dst ./build/client -M /source/ user@host:destination/
# Exclude patterns from file, limit depth # Keep backups of overwritten destination files
./build/client --exclude-from ignore.txt --max-depth 3 /src user@host:/dst ./build/client --backup --backup-dir backups \
/source/ user@host:destination/
# Log to file
./build/client --log-file /tmp/fastsync.log /src user@host:/dst
# All features
./build/client -a --progress --chunk-size 5242880 --exclude "*.log" --delete /src /dst
``` ```
## FastSync Extensions
FastSync-native options are intended to add performance or operational
features without changing the meaning of ordinary compatibility options.
| Option | Purpose |
|---|---|
| `-m` | Enable the multithreaded scanner/loader/sender pipeline. |
| `-c [level]`, `-z [level]` | Enable streaming zstd compression, levels 1-22. |
| `--compress-level <n>` | Set the zstd compression level. |
| `--chunk-size <bytes>` | Set the transfer chunk size. |
| `-s` | Enable FastSync chunk serialization. |
| `-f`, `--sendfile` | Use TCP `sendfile()` zero-copy transfer. Incompatible with compression and chunk serialization. |
| `--delta` | Use FastSync-native block delta transfer. Requires `--incremental`. |
| `--delta-block <bytes>` | Set the FastSync delta block size. |
| `--delta-max <bytes>` | Limit files eligible for FastSync delta transfer. |
| `--server-host <host>` | Select the TCP server host. |
| `--server-port <port>` | Select the TCP server port. |
| `--tls` | Enable TLS for TCP transport. |
| `--bwlimit <KB/s>` | Apply token-bucket bandwidth limiting. |
| `--progress` | Show transfer progress and throughput. |
| `--stats` | Print transfer statistics. |
| `--timeout <seconds>` | Set I/O timeout. |
| `--contimeout <seconds>` | Set connection timeout. |
Current short-option conflicts are tracked as compatibility work. In
particular, FastSync currently uses `-p` for SSH port, `-s` for chunk
serialization, and `-S` for sparse handling. These meanings must be reconciled
before FastSync can claim full rsync CLI compatibility.
## Client Options
### Selection and transfer
| Option | Description |
|---|---|
| `-a`, `--archive` | Enable current archive preset. Full rsync archive semantics are planned. |
| `-n`, `--dry-run` | Scan and report without writing files. |
| `--delete` | Request removal of destination entries absent from the source. The server must allow deletion. |
| `--exclude <pattern>` | Exclude matching paths. Repeatable. |
| `--include <pattern>` | Include matching paths. Repeatable. |
| `--exclude-from <file>` | Read exclude patterns from a file. |
| `--include-from <file>` | Read include patterns from a file. |
| `--max-size <bytes>` | Skip files larger than the limit. |
| `--min-size <bytes>` | Skip files smaller than the limit. |
| `--max-depth <n>` | Limit recursive scanning depth; zero means unlimited. |
| `--incremental` | Skip files matching destination size and mtime. |
| `--checksum` | Include xxHash64 content checks in incremental comparisons. |
| `--backup` | Back up overwritten files. |
| `--backup-dir <dir>` | Store backups under a separate directory. |
| `--suffix <suffix>` | Set the backup filename suffix. |
| `--partial` | Select partial-transfer handling. With `--partial-dir`, completed files are written there; resumable transfers are not implemented. |
| `--partial-dir <dir>` | Set a relative partial-transfer directory below the server destination root; use with `--partial`. |
| `--inplace` | Write directly to the destination instead of using a temporary file. |
### Metadata and links
| Option | Description |
|---|---|
| `-M`, `--preserve` | Preserve supported file metadata, currently mode and modification time. |
| `-l`, `--links` | Request symlink preservation; link-target transfer remains incomplete. |
| `--copy-links` | Copy symlink referents. |
| `--safe-links` | Skip symlinks that point outside the transfer tree. |
| `--copy-unsafe-links` | Copy unsafe symlink referents. |
| `-S`, `--sparse` | Request sparse-file handling; full hole preservation is planned. |
### Output and logging
| Option | Description |
|---|---|
| `-v`, `--verbose` | Enable debug logging. |
| `--progress` | Show live transfer progress. |
| `--stats` | Print transfer statistics. |
| `--log-file <path>` | Write log output to a file. |
| `-V`, `--version` | Print the FastSync protocol version. |
| `--help` | Print command usage. |
### Paths and transport
| Option | Description |
|---|---|
| `-p <port>` | SSH port in the current CLI. This conflicts with rsync's `-p` permissions option and is planned for correction. |
| `--fastsync-server-path <path>` | Remote FastSync server path for SSH mode. |
| `--source-dir <path>` | Set the source directory explicitly. |
| `--dest-dir <path>` | Set the destination directory explicitly. |
| `--save-to-disk` | Enable server-side disk persistence. |
| `--server-host <host>` | TCP server address. |
| `--server-port <port>` | TCP server port. |
| `--tls` | Enable TLS. Requires `--cert` and `--key`. |
| `--cert <path>` | TLS certificate file. |
| `--key <path>` | TLS private key file. |
| `--ca <path>` | CA file for peer verification. |
## Server Options
| Option | Description |
|---|---|
| `--stdio` | Serve one SSH connection over standard input/output. |
| `-p <port>` | TCP listen port. |
| `--tls` | Enable TLS. |
| `--cert <path>` | TLS certificate file. |
| `--key <path>` | TLS private key file. |
| `--ca <path>` | CA file for peer verification. |
| `--destination-root <path>` | Confine received files to this server-side root; defaults to the current directory. |
| `--allow-delete` | Permit client delete manifests. Deletion is refused by default. |
| `-v`, `--verbose` | Enable debug logging. |
| `--help` | Print server usage. |
## Architecture
### Client
- Recursively scans the source tree with include, exclude, size, and depth
filters.
- Sends individual files or serialized chunks.
- Performs incremental checks and optional content checksums.
- Uses a multithreaded producer-consumer pipeline when requested.
- Sends over TCP, TLS-wrapped TCP, or an SSH subprocess.
- Supports progress, statistics, backups, timeouts, and bandwidth limiting.
### Server
- Runs as a TCP listener or one-shot SSH `--stdio` server.
- Receives and reassembles files and decompresses streaming zstd data.
- Applies supported metadata and writes files through a confined destination
root.
- Uses temporary files and atomic rename by default.
- Handles delete manifests only when explicitly authorized.
- Enforces connection, message-size, and path-safety limits.
## Protocol and Security
FastSync protocol version `2.2.0` is shared by the client and server. The
current protocol is sender-driven and includes configuration negotiation,
incremental checks, checksums, manifests, keep-alives, abort handling, and
FastSync-native delta messages. Client and server versions must currently
match exactly.
TLS provides encrypted TCP transport. Supplying `--ca` enables certificate
verification; without it, traffic is encrypted but peer identity is not
verified. Use certificate verification for deployments where authentication
matters. The default TCP transport is not encrypted.
The receiver protects its destination root with path validation, `openat()`
directory traversal, `O_NOFOLLOW`, temporary files, and atomic renames. Delete
operations require the server's explicit `--allow-delete` policy.
## Compatibility Roadmap
The project will reach the drop-in replacement goal in stages:
1. Correct rsync option meanings, including short options, combined options,
and `--option=value` syntax.
2. Add differential tests that compare FastSync and rsync contents, metadata,
links, deletes, filters, dry runs, and exit codes.
3. Make `-a` implement the expected recursive, links, permissions, times,
owner/group, and supported special-file behavior.
4. Complete symlink, sparse-file, metadata, delete-policy, and resumable-write
semantics.
5. Add rsync remote-shell and daemon protocol interoperability.
6. Keep FastSync performance options as negotiated, optional extensions.
The exhaustive implementation matrix and compatibility notes are in
[`RSYNC_COMPAT.md`](RSYNC_COMPAT.md).
## Testing ## Testing
```bash Run the unit test binary:
# Unit tests (18 suites — array_list, chunk, compression, config, data, delta, file, glob,
# metadata, property, protocol, queue, robustness, scanner,
# shared_utils, stress, transport_tcp, transport_ssh, transport_tls)
./build/tests
# Integration + benchmark suite ```bash
python3 test.py ./build/tests
``` ```
The benchmark prints throughput metrics, best configuration, and speedup vs rsync. Run the Python integration suite:
## Performance Considerations ```bash
python3 -m pytest tests/
```
1. Chunk size (~10 MB default) balances memory and transfer efficiency For stricter local validation:
2. Compression level trades CPU for bandwidth
3. `sendfile()` bypasses userspace — ~2× faster on localhost for large files
4. Multithreading scales with core count and uses memory-based pipeline sizing
5. Metadata transfer adds negligible overhead (~24 bytes per file when enabled)
6. SSH socketpair buffer set to 1 MB for improved pipe throughput
7. SSH ControlMaster reuses connections across repeated invocations
8. Incremental sync eliminates redundant transfers entirely
9. Batch incremental reduces round-trips by grouping multiple checks into one message
10. Bandwidth limiting uses token-bucket with nanosleep for accurate throttling
11. Atomic writes add a single `rename()` per file — negligible overhead
12. Path traversal check is O(n) in path length with negligible cost
## Benchmark Results ```bash
cmake -B build-strict -S . -DSTRICT_WARNINGS=ON
cmake --build build-strict -j$(nproc)
cmake -B build-asan -S . -DSANITIZER=address
cmake --build build-asan -j$(nproc)
```
25 MB of mixed file sizes over `localhost` with disk I/O throttled (reads ≤ 15 MB/s, writes ≤ 10 MB/s) and network emulation via `tc netem`. Each test was run 3×; the median is reported below. The benchmark tool compares FastSync configurations with rsync under
controlled local and network conditions:
### LAN (1000 Mbit, 20 ms ±1 ms, 0.1% loss) ```bash
python3 benchmark/bench.py --help
```
| Configuration | Time | vs rsync (archive) | vs rsync (compress) | Benchmark results measure transfer performance only. They do not establish
|---|---|---|---| rsync protocol or filesystem-semantic compatibility.
| **Best: `-m -c`** | **0.20 s** | **11.2× faster** | **3.6× faster** |
| Compression (`-c`) | 0.31 s | 7.3× faster | 2.3× faster |
| Standard | 1.27 s | 1.8× faster | — |
| rsync (archive) | 2.27 s | — | — |
| rsync (archive + compress) | 0.72 s | — | — |
### WAN (100 Mbit, 50 ms ±10 ms, 1% loss) ## Performance Guidance
| Configuration | Time | vs rsync (archive) | vs rsync (compress) | - Use `-m` for workloads with many files or enough CPU parallelism.
|---|---|---|---| - Use `-c` or `-z` when network bandwidth is more constrained than CPU.
| **Best: `-m -c`** | **0.39 s** | **44.8× faster** | **3.8× faster** | - Tune `--chunk-size` for file sizes, memory limits, and network latency.
| Compression (`-c`) | 0.64 s | 27.3× faster | 2.3× faster | - Use `-f` for large uncompressed TCP transfers where zero-copy I/O helps.
| Standard | 7.12 s | 2.4× faster | — | - Use `--incremental` to avoid retransmitting unchanged files.
| rsync (archive) | 17.44 s | — | — | - Use `--delta` for changed files when both endpoints are FastSync peers.
| rsync (archive + compress) | 1.47 s | — | — | - Use `--bwlimit` when sharing a link with other traffic.
Compression reduces the data on the wire enough that the transfer becomes latency-bound rather than bandwidth-bound. On WAN, the best configuration runs 10.8× faster than the theoretical limit for uncompressed data, since zstd shrinks the 25 MB payload to a fraction of its original size over the wire. Always validate the compatibility behavior required by a deployment before
replacing an existing rsync job.
+2 -2
View File
@@ -152,14 +152,14 @@ This document maps rsync's full feature set to FastSync's current implementation
| Flag | Rsync Description | FastSync Status | Notes | | Flag | Rsync Description | FastSync Status | Notes |
|------|-------------------|-----------------|-------| |------|-------------------|-----------------|-------|
| `-S`, `--sparse` | Sparse block handling | ✅ Implemented | `preserve_sparse` config field | | `-S`, `--sparse` | Sparse block handling | ⚠️ Partial | Flag is accepted, but full hole preservation is not implemented |
| `--preallocate` | Allocate dest files before writing | ❌ Not Implemented | | | `--preallocate` | Allocate dest files before writing | ❌ Not Implemented | |
## 11. Checksum & Comparison ## 11. Checksum & Comparison
| Flag | Rsync Description | FastSync Status | Notes | | Flag | Rsync Description | FastSync Status | Notes |
|------|-------------------|-----------------|-------| |------|-------------------|-----------------|-------|
| `--checksum` | Skip based on checksum | ❌ Not Implemented | Removed because it had no effect; `-c` means compression | | `--checksum` | Skip based on checksum | ✅ Implemented | With `--incremental`, compares xxHash64 content checksums; `-c` remains compression |
| `--checksum-choice=STR` | Choose checksum algorithm | ❌ Not Implemented | xxHash used internally | | `--checksum-choice=STR` | Choose checksum algorithm | ❌ Not Implemented | xxHash used internally |
| `--compare-dest=DIR` | Compare dest files relative to DIR | ❌ Not Implemented | Removed because it had no effect | | `--compare-dest=DIR` | Compare dest files relative to DIR | ❌ Not Implemented | Removed because it had no effect |
| `--copy-dest=DIR` | Include copies of unchanged files | ❌ Not Implemented | Removed because it had no effect | | `--copy-dest=DIR` | Include copies of unchanged files | ❌ Not Implemented | Removed because it had no effect |
+195 -189
View File
@@ -11,6 +11,7 @@
#include <errno.h> #include <errno.h>
#include <limits.h> #include <limits.h>
#include <stdbool.h> #include <stdbool.h>
#include <stddef.h>
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
@@ -59,7 +60,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) { static int set_string_option(char** dest, const char* value, const char* option_name) {
char* dup = str_dup(value); char* dup = str_dup(value);
if (!dup) { 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; return -1;
} }
free(*dest); free(*dest);
@@ -70,7 +71,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. */ /* 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) { static int set_positive_int_option(int* dest, const char* value, const char* option_name) {
if (!parse_positive_int(value, dest)) { 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 -1;
} }
return 0; return 0;
@@ -79,7 +80,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. */ /* 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) { static int set_nonneg_int_option(int* dest, const char* value, const char* option_name) {
if (!parse_nonneg_int(value, dest)) { 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 -1;
} }
return 0; return 0;
@@ -87,105 +88,187 @@ 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); 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)},
{"--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. */ /* 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 parse_args(Config* config, int argc, char* argv[], int* positional_args,
int* positional_count) { int* positional_count) {
for (int i = 1; i < argc; i++) { 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;
} else if (apply_table_option(config, entry, NULL) != 0) {
return -1;
}
continue;
}
if (opt_is(argv[i], "--help", NULL)) {
print_usage(); print_usage();
return 1; 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); printf("fastsync version %s\n", PROTOCOL_VERSION);
return 1; 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_compression = true;
config->use_multithreading = true; config->use_multithreading = true;
config->use_metadata = true; config->use_metadata = true;
log_message(LOG_LEVEL_INFO, "Enabled archive mode (-c -m -M)"); log_message(LOG_LEVEL_INFO, "Enabled archive mode (-c -m -M)");
} else if (strcmp(argv[i], "-n") == 0 || strcmp(argv[i], "--dry-run") == 0) { } else if (opt_is(argv[i], "-p", NULL) && i + 1 < argc) {
config->dry_run = true;
} else if (strcmp(argv[i], "-p") == 0 && i + 1 < argc) {
if (set_positive_int_option(&config->ssh_port, argv[++i], "-p") != 0) if (set_positive_int_option(&config->ssh_port, argv[++i], "-p") != 0)
return -1; return -1;
if (config->ssh_port > 65535) { 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; return -1;
} }
} else if (strcmp(argv[i], "--delete") == 0) { } else if (opt_is(argv[i], "--exclude", NULL) && i + 1 < argc) {
config->use_delete = true; if (config_add_pattern(&config->exclude_patterns, &config->exclude_count, argv[++i],
} else if (strcmp(argv[i], "--exclude") == 0 && i + 1 < argc) { "--exclude") != 0)
char** tmp = realloc(config->exclude_patterns, (config->exclude_count + 1) * sizeof(char*));
if (!tmp) {
fprintf(stderr, "Error: memory allocation failed for --exclude\n");
return -1; return -1;
} } else if (opt_is(argv[i], "--include", NULL) && i + 1 < argc) {
config->exclude_patterns = tmp; if (config_add_pattern(&config->include_patterns, &config->include_count, argv[++i],
char* dup = str_dup(argv[++i]); "--include") != 0)
if (!dup) {
fprintf(stderr, "Error: memory allocation failed for --exclude\n");
return -1; return -1;
} } else if (opt_is(argv[i], "--delta-block", NULL) && i + 1 < argc) {
config->exclude_patterns[config->exclude_count++] = dup; unsigned long long val;
} else if (strcmp(argv[i], "--include") == 0 && i + 1 < argc) { if (parse_ull_arg(argv[++i], &val, "--delta-block") != 0)
char** tmp = realloc(config->include_patterns, (config->include_count + 1) * sizeof(char*));
if (!tmp) {
fprintf(stderr, "Error: memory allocation failed for --include\n");
return -1; 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) if (val >= DELTA_BLOCK_SIZE_MIN && val <= DELTA_BLOCK_SIZE_MAX)
config->delta_block_size = (uint32_t)val; config->delta_block_size = (uint32_t)val;
else else
fprintf(stderr, "Warning: --delta-block value %llu out of range, using default\n", val); log_message(LOG_LEVEL_WARNING, "--delta-block value %llu out of range, using default", val);
} else if (strcmp(argv[i], "--delta-max") == 0 && i + 1 < argc) { } else if (opt_is(argv[i], "--delta-max", NULL) && i + 1 < argc) {
char* end; unsigned long long val;
errno = 0; if (parse_ull_arg(argv[++i], &val, "--delta-max") != 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");
return -1; return -1;
}
if (val >= DELTA_MIN_FILE_SIZE) if (val >= DELTA_MIN_FILE_SIZE)
config->delta_max_file_size = val; config->delta_max_file_size = val;
else else
fprintf(stderr, "Warning: --delta-max value %llu too small, using default\n", val); log_message(LOG_LEVEL_WARNING, "--delta-max value %llu too small, using default", val);
} else if (strcmp(argv[i], "-c") == 0 || strcmp(argv[i], "-z") == 0) { } else if (opt_is(argv[i], "-c", "-z")) {
config->use_compression = true; config->use_compression = true;
log_message(LOG_LEVEL_INFO, "Enabled Compression"); log_message(LOG_LEVEL_INFO, "Enabled Compression");
if (i + 1 < argc) { if (i + 1 < argc) {
@@ -193,7 +276,7 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
long level = strtol(argv[i + 1], &end_ptr, 10); long level = strtol(argv[i + 1], &end_ptr, 10);
if (*end_ptr == '\0') { if (*end_ptr == '\0') {
if (level < 1 || level > 22) { 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; return -1;
} }
config->compression_level = (int)level; config->compression_level = (int)level;
@@ -201,91 +284,51 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
i++; i++;
} }
} }
} else if (strcmp(argv[i], "--source-dir") == 0 && i + 1 < argc) { } else if (opt_is(argv[i], "-M", "--preserve")) {
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) {
config->use_metadata = true; config->use_metadata = true;
log_message(LOG_LEVEL_INFO, "Enabled metadata preservation"); log_message(LOG_LEVEL_INFO, "Enabled metadata preservation");
} else if (strcmp(argv[i], "-f") == 0 || strcmp(argv[i], "--sendfile") == 0) { } else if (opt_is(argv[i], "-f", "--sendfile")) {
config->use_sendfile = true; config->use_sendfile = true;
log_message(LOG_LEVEL_INFO, "Enabled sendfile"); 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; config->use_multithreading = true;
log_message(LOG_LEVEL_INFO, "Enabled Multithreading"); 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; config->use_chunk_serialization = true;
log_message(LOG_LEVEL_INFO, "Enabled Chunk Serialization"); log_message(LOG_LEVEL_INFO, "Enabled Chunk Serialization");
} else if (strcmp(argv[i], "--server-host") == 0 && i + 1 < argc) { } else if (opt_is(argv[i], "--server-port", NULL) && 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) {
if (!parse_positive_int(argv[++i], &config->server_port)) { 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; return -1;
} }
if (config->server_port > 65535) { 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; return -1;
} }
} else if (strcmp(argv[i], "--bwlimit") == 0 && i + 1 < argc) { } else if (opt_is(argv[i], "--bwlimit", NULL) && i + 1 < argc) {
char* end; unsigned long long kbps;
errno = 0; if (parse_ull_arg(argv[++i], &kbps, "--bwlimit") != 0)
unsigned long long kbps = strtoull(argv[++i], &end, 10); return -1;
if (errno != 0 || *end != '\0' || kbps == 0) { if (kbps == 0) {
fprintf(stderr, "Error: --bwlimit must be a positive integer\n"); log_message(LOG_LEVEL_ERROR, "--bwlimit must be a positive integer");
return -1; return -1;
} }
if (kbps > ULLONG_MAX / 1024) { if (kbps > ULLONG_MAX / 1024) {
fprintf(stderr, "Error: --bwlimit value too large\n"); log_message(LOG_LEVEL_ERROR, "--bwlimit value too large");
return -1; return -1;
} }
io_set_bwlimit(kbps * 1024); io_set_bwlimit(kbps * 1024);
log_message(LOG_LEVEL_INFO, "Set bandwidth limit to %llu KB/s", kbps); log_message(LOG_LEVEL_INFO, "Set bandwidth limit to %llu KB/s", kbps);
} else if (strcmp(argv[i], "--progress") == 0) { } else if (opt_is(argv[i], "--chunk-size", NULL) && i + 1 < argc) {
config->show_progress = true; unsigned long long val;
} else if (strcmp(argv[i], "--chunk-size") == 0 && i + 1 < argc) { if (parse_ull_arg(argv[++i], &val, "--chunk-size") != 0)
char* end; return -1;
errno = 0; if (val == 0) {
unsigned long long val = strtoull(argv[++i], &end, 10); log_message(LOG_LEVEL_ERROR, "--chunk-size must be a positive integer");
if (errno != 0 || *end != '\0' || val == 0) {
fprintf(stderr, "Error: --chunk-size must be a positive integer\n");
return -1; return -1;
} }
config->chunk_size = val; config->chunk_size = val;
} else if (strcmp(argv[i], "--tls") == 0) { } else if (opt_is(argv[i], "--log-file", NULL) && i + 1 < argc) {
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) {
if (config->log_file) { if (config->log_file) {
fclose(config->log_file); fclose(config->log_file);
config->log_file = NULL; config->log_file = NULL;
@@ -293,55 +336,29 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
} }
FILE* lf = fopen(argv[++i], "a"); FILE* lf = fopen(argv[++i], "a");
if (!lf) { 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; return -1;
} }
config->log_file = lf; config->log_file = lf;
log_set_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) != if (read_patterns_from_file(argv[++i], &config->exclude_patterns, &config->exclude_count) !=
0) 0)
return -1; 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) != if (read_patterns_from_file(argv[++i], &config->include_patterns, &config->include_count) !=
0) 0)
return -1; return -1;
} else if (strcmp(argv[i], "--partial") == 0) { } else if (opt_is(argv[i], "-v", "--verbose")) {
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) {
set_log_level(LOG_LEVEL_DEBUG); set_log_level(LOG_LEVEL_DEBUG);
} else if (strcmp(argv[i], "-l") == 0 || strcmp(argv[i], "--links") == 0) { } else if (opt_is(argv[i], "-T", NULL) && i + 1 < argc) {
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) {
if (set_positive_int_option(&config->timeout, argv[++i], "-T") != 0) if (set_positive_int_option(&config->timeout, argv[++i], "-T") != 0)
return -1; return -1;
} else if (strcmp(argv[i], "--checksum") == 0) { } else if (opt_is(argv[i], "--compress-level", NULL) && i + 1 < argc) {
config->checksum = true;
} else if (strcmp(argv[i], "--compress-level") == 0 && i + 1 < argc) {
if (set_positive_int_option(&config->compression_level, argv[++i], "--compress-level") != 0) if (set_positive_int_option(&config->compression_level, argv[++i], "--compress-level") != 0)
return -1; return -1;
if (config->compression_level < 1 || config->compression_level > 22) { 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; return -1;
} }
} else if (argv[i][0] == '-') { } else if (argv[i][0] == '-') {
@@ -364,7 +381,7 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
static int read_patterns_from_file(const char* filepath, char*** patterns, int* count) { static int read_patterns_from_file(const char* filepath, char*** patterns, int* count) {
FILE* fp = fopen(filepath, "r"); FILE* fp = fopen(filepath, "r");
if (!fp) { 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; return -1;
} }
char* line = NULL; char* line = NULL;
@@ -381,22 +398,11 @@ static int read_patterns_from_file(const char* filepath, char*** patterns, int*
p[--len] = '\0'; p[--len] = '\0';
if (len == 0) if (len == 0)
continue; continue;
char** tmp = realloc(*patterns, (*count + 1) * sizeof(char*)); if (config_add_pattern(patterns, count, p, "pattern file") != 0) {
if (!tmp) {
fprintf(stderr, "Error: memory allocation failed for pattern file\n");
free(line); free(line);
fclose(fp); fclose(fp);
return -1; 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); free(line);
fclose(fp); fclose(fp);
@@ -414,7 +420,7 @@ int main(int argc, char* argv[]) {
bool config_owned_by_pipeline = false; bool config_owned_by_pipeline = false;
Config* config = config_create(); Config* config = config_create();
if (!config) { if (!config) {
fprintf(stderr, "Error: failed to allocate config\n"); log_message(LOG_LEVEL_ERROR, "failed to allocate config");
return 1; return 1;
} }
config->save_to_disk = save_to_disk; config->save_to_disk = save_to_disk;
@@ -435,20 +441,20 @@ int main(int argc, char* argv[]) {
free(config->receive_root_directory); free(config->receive_root_directory);
config->send_directory = str_dup(argv[positional_args[0]]); config->send_directory = str_dup(argv[positional_args[0]]);
if (!config->send_directory) { if (!config->send_directory) {
fprintf(stderr, "Error: memory allocation failed\n"); log_message(LOG_LEVEL_ERROR, "memory allocation failed");
exit_code = 1; exit_code = 1;
goto cleanup; goto cleanup;
} }
config->receive_root_directory = str_dup(argv[positional_args[1]]); config->receive_root_directory = str_dup(argv[positional_args[1]]);
if (!config->receive_root_directory) { if (!config->receive_root_directory) {
fprintf(stderr, "Error: memory allocation failed\n"); log_message(LOG_LEVEL_ERROR, "memory allocation failed");
exit_code = 1; exit_code = 1;
goto cleanup; goto cleanup;
} }
config->save_to_disk = true; config->save_to_disk = true;
config_parse_ssh_dest(config); config_parse_ssh_dest(config);
} else if (positional_count == 1) { } else if (positional_count == 1) {
fprintf(stderr, "Error: missing destination argument\n"); log_message(LOG_LEVEL_ERROR, "missing destination argument");
print_usage(); print_usage();
exit_code = 1; exit_code = 1;
goto cleanup; goto cleanup;
@@ -456,7 +462,7 @@ int main(int argc, char* argv[]) {
if (!config->send_directory && env_source) { if (!config->send_directory && env_source) {
config->send_directory = str_dup(env_source); config->send_directory = str_dup(env_source);
if (!config->send_directory) { if (!config->send_directory) {
fprintf(stderr, "Error: memory allocation failed\n"); log_message(LOG_LEVEL_ERROR, "memory allocation failed");
exit_code = 1; exit_code = 1;
goto cleanup; goto cleanup;
} }
@@ -464,7 +470,7 @@ int main(int argc, char* argv[]) {
if (!config->receive_root_directory && env_dest) { if (!config->receive_root_directory && env_dest) {
config->receive_root_directory = str_dup(env_dest); config->receive_root_directory = str_dup(env_dest);
if (!config->receive_root_directory) { if (!config->receive_root_directory) {
fprintf(stderr, "Error: memory allocation failed\n"); log_message(LOG_LEVEL_ERROR, "memory allocation failed");
exit_code = 1; exit_code = 1;
goto cleanup; goto cleanup;
} }
+69 -53
View File
@@ -42,7 +42,7 @@ static ScannerOptions scanner_options_from_config(const Config* config, int num_
static Client* connect_transfer_client(const Config* config) { static Client* connect_transfer_client(const Config* config) {
if (config->transport == TRANSPORT_SSH) { if (config->transport == TRANSPORT_SSH) {
if (config->use_sendfile) { if (config->use_sendfile) {
fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n"); log_message(LOG_LEVEL_ERROR, "-f/--sendfile is not supported with SSH transport");
return NULL; return NULL;
} }
return client_connect_ssh(config->ssh_destination, config->ssh_port, return client_connect_ssh(config->ssh_destination, config->ssh_port,
@@ -60,6 +60,7 @@ static Client* connect_transfer_client(const Config* config) {
connected = client_connect(client, config->server_host, config->server_port); connected = client_connect(client, config->server_host, config->server_port);
} }
if (!connected) { if (!connected) {
client_disconnect(client);
client_delete(client); client_delete(client);
return NULL; return NULL;
} }
@@ -208,8 +209,11 @@ static int incremental_check(Client* client, File* file, const Config* config,
static int send_delta(Client* client, File* file, DeltaSignature* sig, 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); Delta* delta = delta_compute(file->data->data, file->data->size, sig, config->delta_block_size);
if (!delta) if (!delta) {
if (!send_status(client->file_descriptor, STATUS_NEXT))
return -1;
return 1; return 1;
}
if (!delta_is_worthwhile(delta, file->data->size)) { if (!delta_is_worthwhile(delta, file->data->size)) {
delta_destroy(delta); delta_destroy(delta);
@@ -373,8 +377,9 @@ static int send_chunks_multithreaded(void* pipeline_context) {
Client* client = connect_transfer_client(context->config); Client* client = connect_transfer_client(context->config);
if (!client) { if (!client) {
if (context->config->transport == TRANSPORT_TCP) if (context->config->transport == TRANSPORT_TCP)
fprintf(stderr, "Error: could not connect to server%s\n", log_message(LOG_LEVEL_ERROR, "could not connect to server%s",
context->config->use_tls ? " via TLS" : ""); context->config->use_tls ? " via TLS" : "");
pipeline_cancel(context);
mark_sender_done(context); mark_sender_done(context);
return thrd_error; return thrd_error;
} }
@@ -383,6 +388,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
protocol_session_set_ssl(&session, (SSL*)client->ssl); protocol_session_set_ssl(&session, (SSL*)client->ssl);
protocol_session_bind(&session); protocol_session_bind(&session);
if (!config_send(client->file_descriptor, context->config)) { if (!config_send(client->file_descriptor, context->config)) {
pipeline_cancel(context);
disconnect_transfer_client(client); disconnect_transfer_client(client);
mark_sender_done(context); mark_sender_done(context);
protocol_session_unbind(); protocol_session_unbind();
@@ -394,6 +400,13 @@ static int send_chunks_multithreaded(void* pipeline_context) {
context->queue_loader, &context->mutex_loader, &context->condition_not_empty_loader, context->queue_loader, &context->mutex_loader, &context->condition_not_empty_loader,
&context->condition_not_full_loader, &context->loader_done); &context->condition_not_full_loader, &context->loader_done);
if (current_chunk == NULL) { 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 (context->config->use_delete) {
if (send_delete_manifest(client->file_descriptor, context->manifest) != 0) if (send_delete_manifest(client->file_descriptor, context->manifest) != 0)
goto send_fail; goto send_fail;
@@ -412,7 +425,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
return thrd_error; return thrd_error;
} }
if (send_chunk(client, current_chunk, context->config) != 0) { 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); chunk_destroy(current_chunk);
pipeline_cancel(context); pipeline_cancel(context);
disconnect_transfer_client(client); disconnect_transfer_client(client);
@@ -506,9 +519,10 @@ static int load_files_multithreaded(void* pipeline_context) {
if (f->data->size > STREAM_THRESHOLD) if (f->data->size > STREAM_THRESHOLD)
continue; continue;
if (!file_load_data(f)) { if (!file_load_data(f)) {
log_message(LOG_LEVEL_ERROR, "Failed to load file data, skipping"); log_message(LOG_LEVEL_ERROR, "Failed to load file data");
file_destroy(f); chunk_destroy(chunk);
chunk->items[i] = NULL; pipeline_cancel(context);
return thrd_error;
} }
} }
} }
@@ -517,14 +531,23 @@ static int load_files_multithreaded(void* pipeline_context) {
&context->condition_not_full_loader, &context->condition_not_full_loader,
&context->cancelled)) { &context->cancelled)) {
chunk_destroy(chunk); chunk_destroy(chunk);
atomic_store(&context->cancelled, true); pipeline_cancel(context);
cnd_broadcast(&context->condition_not_full_loader);
cnd_broadcast(&context->condition_not_empty_loader);
return thrd_error; 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 /* Progress-reporting thread for multithreaded send. Runs in parallel with
the scanner/loader/sender threads and prints periodic progress to stderr. */ the scanner/loader/sender threads and prints periodic progress to stderr. */
static int progress_thread_fn(void* arg) { static int progress_thread_fn(void* arg) {
@@ -539,20 +562,14 @@ static int progress_thread_fn(void* arg) {
mtx_unlock(&context->mutex_progress); mtx_unlock(&context->mutex_progress);
if (done) { if (done) {
time_t now = time(NULL); print_transfer_progress(total, start, "Done.\n");
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);
break; break;
} }
time_t now = time(NULL); time_t now = time(NULL);
if (now - last_progress >= 1) { if (now - last_progress >= 1) {
last_progress = now; last_progress = now;
double elapsed = difftime(now, start); print_transfer_progress(total, 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);
} }
struct timespec ts = {0, 100 * 1000000L}; /* 100 ms */ struct timespec ts = {0, 100 * 1000000L}; /* 100 ms */
@@ -568,36 +585,29 @@ int send_files(Config* config) {
Client* client = connect_transfer_client(config); Client* client = connect_transfer_client(config);
if (!client) { if (!client) {
if (config->transport == TRANSPORT_TCP) if (config->transport == TRANSPORT_TCP)
fprintf(stderr, "Error: could not connect to server%s\n", config->use_tls ? " via TLS" : ""); log_message(LOG_LEVEL_ERROR, "could not connect to server%s",
config->use_tls ? " via TLS" : "");
return 1; return 1;
} }
ProtocolSession session; ProtocolSession session;
protocol_session_init(&session, client->file_descriptor, client->file_descriptor); protocol_session_init(&session, client->file_descriptor, client->file_descriptor);
protocol_session_set_ssl(&session, (SSL*)client->ssl); protocol_session_set_ssl(&session, (SSL*)client->ssl);
protocol_session_bind(&session); protocol_session_bind(&session);
if (!config_send(client->file_descriptor, config)) { int ret = 1;
disconnect_transfer_client(client); DirectoryScanner* scanner = NULL;
protocol_session_unbind(); ArrayList* manifest = NULL;
return 1; if (!config_send(client->file_descriptor, config))
} goto send_fail;
ScannerOptions scanner_options = scanner_options_from_config(config, 0); ScannerOptions scanner_options = scanner_options_from_config(config, 0);
DirectoryScanner* scanner = scanner = directory_scanner_create_with_options(config->send_directory, &scanner_options);
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; Chunk* current_chunk;
unsigned long long total_bytes = 0; unsigned long long total_bytes = 0;
int total_files = 0; int total_files = 0;
time_t last_progress = 0; time_t last_progress = 0;
time_t start = time(NULL); time_t start = time(NULL);
ArrayList* manifest = create_transfer_manifest(config);
if (!scanner || (config->use_delete && !manifest)) {
if (scanner)
directory_scanner_destroy(scanner);
if (manifest)
array_list_delete(manifest);
disconnect_transfer_client(client);
protocol_session_unbind();
return 1;
}
while ((current_chunk = directory_scanner_next(scanner)) != NULL) { while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
unsigned long long chunk_bytes = 0; unsigned long long chunk_bytes = 0;
for (int i = 0; i < current_chunk->element_count; i++) { for (int i = 0; i < current_chunk->element_count; i++) {
@@ -609,15 +619,21 @@ int send_files(Config* config) {
goto send_fail; goto send_fail;
} }
if (!config->use_sendfile) { if (!config->use_sendfile) {
bool load_ok = true;
for (int i = 0; i < current_chunk->element_count; i++) { for (int i = 0; i < current_chunk->element_count; i++) {
File* f = current_chunk->items[i]; File* f = current_chunk->items[i];
if (f->data->size > STREAM_THRESHOLD) if (f->data->size > STREAM_THRESHOLD)
continue; continue;
if (!file_load_data(f)) { if (!file_load_data(f)) {
log_message(LOG_LEVEL_ERROR, "Failed to load file data"); 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) { if (send_chunk(client, current_chunk, config) != 0) {
log_message(LOG_LEVEL_ERROR, "Failed to send chunk"); log_message(LOG_LEVEL_ERROR, "Failed to send chunk");
@@ -632,10 +648,7 @@ int send_files(Config* config) {
time_t now = time(NULL); time_t now = time(NULL);
if (now - last_progress >= 1) { if (now - last_progress >= 1) {
last_progress = now; last_progress = now;
double elapsed = difftime(now, start); print_transfer_progress(total_bytes, 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);
} }
} }
chunk_destroy(current_chunk); chunk_destroy(current_chunk);
@@ -645,34 +658,33 @@ int send_files(Config* config) {
if (config->use_delete) { if (config->use_delete) {
if (send_delete_manifest(client->file_descriptor, manifest) != 0) { if (send_delete_manifest(client->file_descriptor, manifest) != 0) {
array_list_delete(manifest); array_list_delete(manifest);
manifest = NULL;
goto send_fail; goto send_fail;
} }
array_list_delete(manifest); array_list_delete(manifest);
manifest = NULL; manifest = NULL;
} }
bool ok = finalize_transfer(client); bool ok = finalize_transfer(client);
double elapsed_total = difftime(time(NULL), start); if (config->show_progress)
if (config->show_progress) { print_transfer_progress(total_bytes, start, "Done.\n");
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);
}
if (config->stats) { if (config->stats) {
double elapsed_total = difftime(time(NULL), start);
double rate = elapsed_total > 0 ? total_bytes / (1048576.0 * elapsed_total) : 0; 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, fprintf(stderr, "Stats: %d files, %.1f MB, %.1f MB/s\n", total_files, total_bytes / 1048576.0,
rate); rate);
} }
directory_scanner_destroy(scanner); ret = ok ? 0 : 1;
disconnect_transfer_client(client);
protocol_session_unbind();
return ok ? 0 : 1;
send_fail: 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) if (manifest)
array_list_delete(manifest); array_list_delete(manifest);
if (scanner)
directory_scanner_destroy(scanner); directory_scanner_destroy(scanner);
disconnect_transfer_client(client); disconnect_transfer_client(client);
protocol_session_unbind(); protocol_session_unbind();
return 1; return ret;
} }
int send_files_multithreaded(Config* config) { int send_files_multithreaded(Config* config) {
@@ -708,6 +720,10 @@ int send_files_multithreaded(Config* config) {
} }
if (config->use_delete) if (config->use_delete)
context->manifest = create_transfer_manifest(config); context->manifest = create_transfer_manifest(config);
if (config->use_delete && !context->manifest) {
pipeline_context_sender_destroy(context);
return 1;
}
thrd_t scanner, loader, sender; thrd_t scanner, loader, sender;
bool scanner_created = false; bool scanner_created = false;
@@ -721,7 +737,7 @@ int send_files_multithreaded(Config* config) {
sender_created = (thrd_create(&sender, send_chunks_multithreaded, context) == thrd_success); sender_created = (thrd_create(&sender, send_chunks_multithreaded, context) == thrd_success);
if (!scanner_created || !loader_created || !sender_created) { if (!scanner_created || !loader_created || !sender_created) {
perror("Error creating threads.\n"); log_perror("Error creating threads");
pipeline_cancel(context); pipeline_cancel(context);
mtx_lock(&context->mutex_progress); mtx_lock(&context->mutex_progress);
context->sender_done = true; context->sender_done = true;
@@ -741,7 +757,7 @@ int send_files_multithreaded(Config* config) {
if (config->show_progress) { if (config->show_progress) {
progress_created = (thrd_create(&progress, progress_thread_fn, context) == thrd_success); progress_created = (thrd_create(&progress, progress_thread_fn, context) == thrd_success);
if (!progress_created) { if (!progress_created) {
perror("Error creating progress thread.\n"); log_perror("Error creating progress thread");
/* Non-fatal; continue without progress reporting */ /* Non-fatal; continue without progress reporting */
} }
} }
+10 -9
View File
@@ -1,37 +1,38 @@
#include "client_validation.h" #include "client_validation.h"
#include "log.h"
#include "usage.h" #include "usage.h"
#include <stdio.h> #include <stdio.h>
/* Validate config after parsing. Returns true if valid. */ /* Validate config after parsing. Returns true if valid. */
bool validate_config(const Config* config) { bool validate_config(const Config* config) {
if (!config->send_directory || !config->receive_root_directory) { if (!config->send_directory || !config->receive_root_directory) {
fprintf(stderr, "Error: source and destination directories are required\n"); log_message(LOG_LEVEL_ERROR, "source and destination directories are required");
print_usage(); print_usage();
return false; return false;
} }
if (config->use_sendfile && (config->use_chunk_serialization || config->use_compression)) { 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 " log_message(LOG_LEVEL_ERROR, "-f/--sendfile cannot be combined with -c (compression) or -s "
"serialization)\n"); "(chunk serialization)");
return false; return false;
} }
if (config->transport == TRANSPORT_SSH && config->use_sendfile) { if (config->transport == TRANSPORT_SSH && config->use_sendfile) {
fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n"); log_message(LOG_LEVEL_ERROR, "-f/--sendfile is not supported with SSH transport");
return false; return false;
} }
if (config->use_incremental && config->use_chunk_serialization) { if (config->use_incremental && config->use_chunk_serialization) {
fprintf(stderr, "Error: --incremental is not supported with -s (chunk serialization)\n"); log_message(LOG_LEVEL_ERROR, "--incremental is not supported with -s (chunk serialization)");
return false; return false;
} }
if (config->use_delta && !config->use_incremental) { if (config->use_delta && !config->use_incremental) {
fprintf(stderr, "Error: --delta requires --incremental\n"); log_message(LOG_LEVEL_ERROR, "--delta requires --incremental");
return false; return false;
} }
if (config->use_delta && config->use_chunk_serialization) { if (config->use_delta && config->use_chunk_serialization) {
fprintf(stderr, "Error: --delta cannot be combined with -s (chunk serialization)\n"); log_message(LOG_LEVEL_ERROR, "--delta cannot be combined with -s (chunk serialization)");
return false; return false;
} }
if (config->use_delta && config->use_sendfile) { if (config->use_delta && config->use_sendfile) {
fprintf(stderr, "Error: --delta cannot be combined with -f (sendfile)\n"); log_message(LOG_LEVEL_ERROR, "--delta cannot be combined with -f (sendfile)");
return false; return false;
} }
if (config->append || config->append_verify) { if (config->append || config->append_verify) {
@@ -42,7 +43,7 @@ bool validate_config(const Config* config) {
} }
if (config->use_tls) { if (config->use_tls) {
if (!config->tls_cert || !config->tls_key) { if (!config->tls_cert || !config->tls_key) {
fprintf(stderr, "Error: --tls requires --cert and --key\n"); log_message(LOG_LEVEL_ERROR, "--tls requires --cert and --key");
return false; return false;
} }
} }
+163 -125
View File
@@ -1,3 +1,4 @@
#include "log.h"
#include "scanner.h" #include "scanner.h"
#include "array_list.h" #include "array_list.h"
#include "chunk.h" #include "chunk.h"
@@ -226,7 +227,7 @@ static int open_next_directory(DirectoryScanner* scanner) {
free(de); free(de);
scanner->current_dir = opendir(scanner->current_path); scanner->current_dir = opendir(scanner->current_path);
if (scanner->current_dir == NULL) { if (scanner->current_dir == NULL) {
perror("Could not open directory"); log_perror("Could not open directory");
free(scanner->current_path); free(scanner->current_path);
scanner->current_path = NULL; scanner->current_path = NULL;
scanner->failed = true; scanner->failed = true;
@@ -363,6 +364,8 @@ static int parallel_worker_thread(void* arg) {
cnd_broadcast(&wa->ps->result_not_empty); cnd_broadcast(&wa->ps->result_not_empty);
cnd_broadcast(&wa->ps->result_not_full); cnd_broadcast(&wa->ps->result_not_full);
mtx_unlock(&wa->ps->result_mutex); mtx_unlock(&wa->ps->result_mutex);
for (int j = i; j < wa->dir_count; j++)
free(wa->dirs[j]);
break; break;
} }
Chunk* chunk; Chunk* chunk;
@@ -398,18 +401,23 @@ static int parallel_worker_thread(void* arg) {
return thrd_success; return thrd_success;
} }
ParallelScanner* parallel_scanner_create_with_options(const char* root_directory, static void parallel_scanner_creation_failed(ParallelScanner* ps) {
const ScannerOptions* options) { mtx_lock(&ps->result_mutex);
if (!root_directory || !options) ps->failed = true;
return NULL; atomic_store(&ps->cancelled, true);
ParallelScanner* ps = calloc(1, sizeof(ParallelScanner)); ps->expected_threads = ps->created_threads;
if (!ps) if (ps->completed >= ps->expected_threads)
return NULL; 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); ps->result_queue = queue_create(100, chunk_destroy);
if (!ps->result_queue) { if (!ps->result_queue)
free(ps); return false;
return NULL;
}
atomic_init(&ps->cancelled, false); atomic_init(&ps->cancelled, false);
int init = 0; int init = 0;
bool ok = true; bool ok = true;
@@ -434,93 +442,37 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
if (init >= 1) if (init >= 1)
mtx_destroy(&ps->result_mutex); mtx_destroy(&ps->result_mutex);
queue_destroy(ps->result_queue); queue_destroy(ps->result_queue);
free(ps); ps->result_queue = NULL;
return NULL; return false;
} }
return true;
}
DIR* dir = opendir(root_directory); /* Split files into chunks of roughly chunk_size bytes. Returns the first chunk (also stored
if (!dir) { * chunks beyond the first are enqueued on `queue`). Nulls out consumed entries in `files`.
perror("Could not open root directory for parallel scan"); * Sets *failed on allocation/enqueue errors. */
parallel_scanner_destroy(ps); static Chunk* batch_files(ArrayList* files, unsigned long long chunk_size, Queue* queue,
bool* failed) {
Chunk* first = NULL;
if (files->size <= 0)
return NULL; return NULL;
}
ArrayList* root_files = array_list_create(file_destroy);
ArrayList* subdirs = array_list_create(free);
if (!root_files || !subdirs) {
array_list_delete(root_files);
array_list_delete(subdirs);
closedir(dir);
parallel_scanner_destroy(ps);
return NULL;
}
const struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
ScannerEntry inspected;
int inspection =
scanner_inspect_entry(options, root_directory, root_directory, entry->d_name, &inspected);
if (inspection < 0) {
ps->failed = true;
continue;
}
if (inspection == 0)
continue;
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;
}
} else {
File* file = file_create(cur_path);
free(cur_path);
if (!file) {
ps->failed = true;
continue;
}
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;
continue;
}
if (!array_list_add(root_files, file)) {
file_destroy(file);
ps->failed = true;
}
}
}
closedir(dir);
unsigned long long cs = options->chunk_size > 0 ? options->chunk_size : DESIRED_CHUNK_SIZE;
if (root_files->size > 0) {
ArrayList* batch = array_list_create(NULL); ArrayList* batch = array_list_create(NULL);
if (!batch) { if (!batch) {
ps->failed = true; *failed = true;
array_list_delete(root_files);
array_list_delete(subdirs);
parallel_scanner_destroy(ps);
return NULL; return NULL;
} }
unsigned long long batch_size = 0; unsigned long long batch_size = 0;
Chunk* first = NULL; for (int i = 0; i < files->size; i++) {
for (int i = 0; i < root_files->size; i++) { File* f = (File*)files->items[i];
File* f = (File*)root_files->items[i];
if (!array_list_add(batch, f)) { if (!array_list_add(batch, f)) {
ps->failed = true; *failed = true;
break; break;
} }
batch_size += f->data->size; batch_size += f->data->size;
if (batch_size >= cs || i == root_files->size - 1) { if (batch_size >= chunk_size || i == files->size - 1) {
void** items = array_list_to_array(batch); void** items = array_list_to_array(batch);
if (!items) { if (!items) {
ps->failed = true; *failed = true;
batch->item_destroyer = file_destroy;
array_list_delete(batch); array_list_delete(batch);
batch = NULL; batch = NULL;
break; break;
@@ -528,27 +480,29 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
Chunk* c = chunk_create((File**)items, batch->size); Chunk* c = chunk_create((File**)items, batch->size);
free(items); free(items);
if (!c) { if (!c) {
ps->failed = true; *failed = true;
batch->item_destroyer = file_destroy;
array_list_delete(batch); array_list_delete(batch);
batch = NULL; batch = NULL;
break; break;
} }
int batch_start = i - batch->size + 1;
for (int j = batch_start; j <= i; j++)
files->items[j] = NULL;
batch->item_destroyer = NULL; batch->item_destroyer = NULL;
array_list_delete(batch); array_list_delete(batch);
batch = NULL; batch = NULL;
if (!first) { if (!first) {
first = c; first = c;
} else { } else {
if (!queue_enqueue(ps->result_queue, c)) { if (!queue_enqueue(queue, c)) {
chunk_destroy(c); chunk_destroy(c);
ps->failed = true; *failed = true;
} }
} }
if (i < root_files->size - 1) { if (i < files->size - 1) {
batch = array_list_create(NULL); batch = array_list_create(NULL);
if (!batch) { if (!batch) {
ps->failed = true; *failed = true;
break; break;
} }
batch_size = 0; batch_size = 0;
@@ -559,23 +513,88 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
batch->item_destroyer = NULL; batch->item_destroyer = NULL;
array_list_delete(batch); array_list_delete(batch);
} }
ps->initial_chunk = first; return first;
root_files->item_destroyer = NULL; }
}
array_list_delete(root_files);
/* 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) {
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; int n = options->num_threads > 0 ? options->num_threads : 4;
if (n > subdirs->size) if (n > subdirs->size)
n = subdirs->size > 0 ? subdirs->size : 1; n = subdirs->size;
if (subdirs->size > 0) {
ps->num_threads = n; ps->num_threads = n;
ps->expected_threads = n; ps->expected_threads = n;
ps->threads = calloc(n, sizeof(thrd_t)); ps->threads = calloc(n, sizeof(thrd_t));
if (!ps->threads) { if (!ps->threads) {
array_list_delete(subdirs); ps->num_threads = 0;
parallel_scanner_destroy(ps); ps->expected_threads = 0;
return NULL; ps->failed = true;
return;
} }
int dirs_per_thread = subdirs->size / n; int dirs_per_thread = subdirs->size / n;
int remainder = subdirs->size % n; int remainder = subdirs->size % n;
@@ -587,14 +606,14 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
break; break;
ParallelWorkerArg* wa = calloc(1, sizeof(ParallelWorkerArg)); ParallelWorkerArg* wa = calloc(1, sizeof(ParallelWorkerArg));
if (!wa) { if (!wa) {
ps->failed = true; parallel_scanner_creation_failed(ps);
break; break;
} }
wa->ps = ps; wa->ps = ps;
wa->dirs = calloc(count, sizeof(char*)); wa->dirs = calloc(count, sizeof(char*));
if (!wa->dirs) { if (!wa->dirs) {
free(wa); free(wa);
ps->failed = true; parallel_scanner_creation_failed(ps);
break; break;
} }
bool dup_ok = true; bool dup_ok = true;
@@ -608,7 +627,7 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
free(wa->dirs[j]); free(wa->dirs[j]);
free(wa->dirs); free(wa->dirs);
free(wa); free(wa);
ps->failed = true; parallel_scanner_creation_failed(ps);
break; break;
} }
wa->dir_count = count; wa->dir_count = count;
@@ -620,35 +639,49 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
free(wa->dirs[j]); free(wa->dirs[j]);
free(wa->dirs); free(wa->dirs);
free(wa); free(wa);
ps->failed = true; parallel_scanner_creation_failed(ps);
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; break;
} }
ps->num_threads++; ps->num_threads++;
ps->created_threads++; ps->created_threads++;
} }
}
array_list_delete(subdirs);
return ps;
} }
ParallelScanner* parallel_scanner_create(const char* root_directory, bool use_metadata, ParallelScanner* parallel_scanner_create_with_options(const char* root_directory,
unsigned long long chunk_size, char** exclude_patterns, const ScannerOptions* options) {
int exclude_count, char** include_patterns, if (!root_directory || !options)
int include_count, unsigned long long max_size, return NULL;
unsigned long long min_size, int max_depth, ParallelScanner* ps = calloc(1, sizeof(ParallelScanner));
int num_threads, bool follow_symlinks, bool copy_links, if (!ps)
bool safe_links, bool copy_unsafe_links, bool checksum) { return NULL;
ScannerOptions options = {use_metadata, chunk_size, exclude_patterns, exclude_count, if (!parallel_scanner_init(ps)) {
include_patterns, include_count, max_size, min_size, free(ps);
max_depth, num_threads, follow_symlinks, copy_links, return NULL;
safe_links, copy_unsafe_links, checksum}; }
return parallel_scanner_create_with_options(root_directory, &options);
ArrayList* root_files = array_list_create(file_destroy);
ArrayList* subdirs = array_list_create(free);
if (!root_files || !subdirs) {
array_list_delete(root_files);
array_list_delete(subdirs);
parallel_scanner_destroy(ps);
return NULL;
}
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;
}
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);
spawn_parallel_workers(ps, subdirs, options, cs);
array_list_delete(subdirs);
return ps;
} }
Chunk* parallel_scanner_next(ParallelScanner* ps) { Chunk* parallel_scanner_next(ParallelScanner* ps) {
@@ -659,6 +692,11 @@ Chunk* parallel_scanner_next(ParallelScanner* ps) {
} }
if (ps->num_threads == 0) { if (ps->num_threads == 0) {
mtx_lock(&ps->result_mutex); 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; ps->done = true;
mtx_unlock(&ps->result_mutex); mtx_unlock(&ps->result_mutex);
return NULL; return NULL;
-7
View File
@@ -77,13 +77,6 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner);
bool directory_scanner_failed(const DirectoryScanner* scanner); bool directory_scanner_failed(const DirectoryScanner* scanner);
void directory_scanner_destroy(DirectoryScanner* scanner); void directory_scanner_destroy(DirectoryScanner* scanner);
ParallelScanner* parallel_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,
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, ParallelScanner* parallel_scanner_create_with_options(const char* root_directory,
const ScannerOptions* options); const ScannerOptions* options);
Chunk* parallel_scanner_next(ParallelScanner* scanner); Chunk* parallel_scanner_next(ParallelScanner* scanner);
+25 -6
View File
@@ -20,6 +20,16 @@ static char* authorized_root;
static int authorized_root_fd = -1; static int authorized_root_fd = -1;
static bool allow_delete; static bool allow_delete;
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) { static bool path_is_within(const char* root, const char* path) {
size_t n = strlen(root); size_t n = strlen(root);
return strncmp(root, path, n) == 0 && (path[n] == '\0' || path[n] == '/'); return strncmp(root, path, n) == 0 && (path[n] == '\0' || path[n] == '/');
@@ -38,7 +48,13 @@ static bool __attribute__((unused)) configure_authorization(const char* root) {
authorized_root = NULL; authorized_root = NULL;
return false; return false;
} }
file_set_authorized_root(authorized_root_fd, authorized_root); if (!file_set_authorized_root(authorized_root_fd, authorized_root)) {
close(authorized_root_fd);
authorized_root_fd = -1;
free(authorized_root);
authorized_root = NULL;
return false;
}
utils_set_authorized_root_fd(authorized_root_fd); utils_set_authorized_root_fd(authorized_root_fd);
return true; return true;
} }
@@ -106,13 +122,14 @@ void handler(int file_descriptor) {
protocol_session_unbind(); protocol_session_unbind();
return; return;
} }
context->session.total_allocated_bytes = session.total_allocated_bytes;
thrd_t receiver, writer; thrd_t receiver, writer;
bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success; bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success;
bool writer_created = false; bool writer_created = false;
if (receiver_created) 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) { if (!receiver_created || !writer_created) {
perror("Error creating Threads"); log_perror("Error creating Threads");
if (receiver_created) { if (receiver_created) {
mtx_lock(&context->mutex); mtx_lock(&context->mutex);
atomic_store(&context->cancelled, true); atomic_store(&context->cancelled, true);
@@ -226,33 +243,35 @@ int main(int argc, char* argv[]) {
if (stdio_mode) { if (stdio_mode) {
io_set_fds(STDIN_FILENO, STDOUT_FILENO); io_set_fds(STDIN_FILENO, STDOUT_FILENO);
handler(STDIN_FILENO); handler(STDIN_FILENO);
file_set_authorized_root(-1, NULL); release_authorization();
utils_set_authorized_root_fd(-1);
close(authorized_root_fd);
free(authorized_root);
return 0; return 0;
} }
g_server = server_create(port); g_server = server_create(port);
if (!g_server) { if (!g_server) {
log_message(LOG_LEVEL_ERROR, "Failed to create server"); log_message(LOG_LEVEL_ERROR, "Failed to create server");
release_authorization();
return 1; return 1;
} }
if (use_tls) { if (use_tls) {
if (!tls_cert || !tls_key) { if (!tls_cert || !tls_key) {
fprintf(stderr, "Error: --tls requires --cert and --key\n"); fprintf(stderr, "Error: --tls requires --cert and --key\n");
server_delete(&g_server); server_delete(&g_server);
release_authorization();
return 1; return 1;
} }
tls_global_init(); tls_global_init();
if (!server_create_tls(g_server, tls_cert, tls_key, tls_ca)) { if (!server_create_tls(g_server, tls_cert, tls_key, tls_ca)) {
log_message(LOG_LEVEL_ERROR, "Failed to set up TLS"); log_message(LOG_LEVEL_ERROR, "Failed to set up TLS");
server_delete(&g_server); server_delete(&g_server);
release_authorization();
return 1; return 1;
} }
server_listen_tls(g_server, handler); server_listen_tls(g_server, handler);
} else { } else {
server_listen(g_server, handler); server_listen(g_server, handler);
} }
server_delete(&g_server);
release_authorization();
return 0; return 0;
} }
#endif #endif
+5 -4
View File
@@ -1,3 +1,4 @@
#include "log.h"
#include "array_list.h" #include "array_list.h"
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
@@ -6,7 +7,7 @@
ArrayList* array_list_create(void (*item_destroyer)(void* item)) { ArrayList* array_list_create(void (*item_destroyer)(void* item)) {
ArrayList* list = (ArrayList*)malloc(sizeof(ArrayList)); ArrayList* list = (ArrayList*)malloc(sizeof(ArrayList));
if (list == NULL) { 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; return NULL;
} }
@@ -34,7 +35,7 @@ void array_list_delete(ArrayList* array_list) {
free(array_list); free(array_list);
} }
bool array_list_extend(ArrayList* array_list) { static bool array_list_extend(ArrayList* array_list) {
if (array_list == NULL) if (array_list == NULL)
return false; return false;
int new_capacity = array_list->capacity * 2; int new_capacity = array_list->capacity * 2;
@@ -42,7 +43,7 @@ bool array_list_extend(ArrayList* array_list) {
new_capacity = INITIAL_ARRAY_SIZE; new_capacity = INITIAL_ARRAY_SIZE;
void* new_items = realloc(array_list->items, new_capacity * sizeof(void*)); void* new_items = realloc(array_list->items, new_capacity * sizeof(void*));
if (new_items == NULL) { 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; return false;
} }
array_list->items = new_items; 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*)); void** array = malloc(array_list->size * sizeof(void*));
if (array == NULL) { 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; return NULL;
} }
memcpy(array, array_list->items, array_list->size * sizeof(void*)); 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)); ArrayList* array_list_create(void (*item_destroyer)(void* item));
void array_list_delete(ArrayList* array_list); void array_list_delete(ArrayList* array_list);
bool array_list_extend(ArrayList* array_list);
bool array_list_add(ArrayList* array_list, void* item); bool array_list_add(ArrayList* array_list, void* item);
void** array_list_to_array(const ArrayList* array_list); void** array_list_to_array(const ArrayList* array_list);
+45 -6
View File
@@ -1,4 +1,5 @@
#include <stddef.h> #include <stddef.h>
#include <stdint.h>
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
@@ -14,11 +15,12 @@
/* Maximum individual file data size within a chunk (64 MB) */ /* Maximum individual file data size within a chunk (64 MB) */
#define MAX_FILE_DATA_SIZE (64ULL * 1024 * 1024) #define MAX_FILE_DATA_SIZE (64ULL * 1024 * 1024)
#define MAX_CHUNK_FILES (1024 * 1024)
Chunk* chunk_create(File** items, int element_count) { Chunk* chunk_create(File** items, int element_count) {
Chunk* chunk = (Chunk*)malloc(sizeof(Chunk)); Chunk* chunk = (Chunk*)malloc(sizeof(Chunk));
if (chunk == NULL) { if (chunk == NULL) {
perror("ERROR: Could not allocate memory for chunk structure"); log_perror("ERROR: Could not allocate memory for chunk structure");
return NULL; return NULL;
} }
@@ -92,6 +94,11 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
size_t remaining_size = data->size; size_t remaining_size = data->size;
while (remaining_size > 0) { while (remaining_size > 0) {
if (files->size >= MAX_CHUNK_FILES) {
log_message(LOG_LEVEL_ERROR, "Chunk contains too many files");
array_list_delete(files);
return NULL;
}
if (remaining_size < sizeof(size_t)) { if (remaining_size < sizeof(size_t)) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for path length"); log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for path length");
array_list_delete(files); array_list_delete(files);
@@ -103,7 +110,7 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
data_pointer += sizeof(size_t); data_pointer += sizeof(size_t);
remaining_size -= 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"); log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for path");
array_list_delete(files); array_list_delete(files);
return NULL; return NULL;
@@ -111,7 +118,7 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
char* path = malloc(path_len + 1); char* path = malloc(path_len + 1);
if (path == NULL) { if (path == NULL) {
perror("Could not allocate memory for file path"); log_perror("Could not allocate memory for file path");
array_list_delete(files); array_list_delete(files);
return NULL; return NULL;
} }
@@ -122,10 +129,15 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
File* file = file_create(path); File* file = file_create(path);
free(path); free(path);
if (file == NULL) {
array_list_delete(files);
return NULL;
}
if (use_metadata) { if (use_metadata) {
if (remaining_size < sizeof(int)) { if (remaining_size < sizeof(int)) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for metadata"); log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for metadata");
file_destroy(file);
array_list_delete(files); array_list_delete(files);
return NULL; return NULL;
} }
@@ -134,6 +146,7 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
memcpy(&present_flag, data_pointer, sizeof(int)); memcpy(&present_flag, data_pointer, sizeof(int));
if (present_flag && remaining_size < sizeof(int) + FILE_METADATA_WIRE_SIZE) { if (present_flag && remaining_size < sizeof(int) + FILE_METADATA_WIRE_SIZE) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for metadata body"); log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for metadata body");
file_destroy(file);
array_list_delete(files); array_list_delete(files);
return NULL; return NULL;
} }
@@ -141,10 +154,16 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
remaining_size -= sizeof(int); remaining_size -= sizeof(int);
if (file->metadata) if (file->metadata)
remaining_size -= FILE_METADATA_WIRE_SIZE; remaining_size -= FILE_METADATA_WIRE_SIZE;
else if (present_flag) {
file_destroy(file);
array_list_delete(files);
return NULL;
}
} }
if (remaining_size < sizeof(size_t)) { if (remaining_size < sizeof(size_t)) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for data size"); log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for data size");
file_destroy(file);
array_list_delete(files); array_list_delete(files);
return NULL; return NULL;
} }
@@ -156,6 +175,7 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
if (remaining_size < file_data_size) { if (remaining_size < file_data_size) {
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for file content"); log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for file content");
file_destroy(file);
array_list_delete(files); array_list_delete(files);
return NULL; return NULL;
} }
@@ -164,29 +184,48 @@ Chunk* chunk_deserialize(Data* data, bool use_metadata) {
if (file_data_size > MAX_FILE_DATA_SIZE) { if (file_data_size > MAX_FILE_DATA_SIZE) {
log_message(LOG_LEVEL_ERROR, "File data size %zu exceeds maximum %llu", 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); (unsigned long long)MAX_FILE_DATA_SIZE);
file_destroy(file);
array_list_delete(files); array_list_delete(files);
return NULL; return NULL;
} }
void* file_data = malloc(file_data_size); void* file_data = malloc(file_data_size > 0 ? file_data_size : 1);
if (file_data == NULL) { 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); array_list_delete(files);
return NULL; return NULL;
} }
memcpy(file_data, data_pointer, file_data_size); memcpy(file_data, data_pointer, file_data_size);
data_destroy(file->data); data_destroy(file->data);
file->data = data_create(file_data, file_data_size); file->data = data_create(file_data, file_data_size);
if (file->data == NULL) {
file_destroy(file);
array_list_delete(files);
return NULL;
}
data_pointer += file_data_size; data_pointer += file_data_size;
remaining_size -= 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); File** file_array = (File**)array_list_to_array(files);
if (file_array == NULL) {
array_list_delete(files);
return NULL;
}
Chunk* chunk = chunk_create(file_array, files->size); Chunk* chunk = chunk_create(file_array, files->size);
free(file_array); free(file_array);
if (chunk == NULL) {
array_list_delete(files);
return NULL;
}
files->item_destroyer = NULL; files->item_destroyer = NULL;
array_list_delete(files); array_list_delete(files);
+2 -2
View File
@@ -108,7 +108,7 @@ Config* config_create(void) {
return config; return config;
} }
bool is_remote_dest(const char* s) { bool config_is_remote_dest(const char* s) {
if (s == NULL) if (s == NULL)
return false; return false;
const char* colon = strchr(s, ':'); const char* colon = strchr(s, ':');
@@ -124,7 +124,7 @@ bool is_remote_dest(const char* s) {
} }
void config_parse_ssh_dest(Config* config) { 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; return;
config->transport = TRANSPORT_SSH; config->transport = TRANSPORT_SSH;
config->ssh_destination = str_dup(config->receive_root_directory); config->ssh_destination = str_dup(config->receive_root_directory);
+1 -1
View File
@@ -135,7 +135,7 @@ Config* config_create(void);
void config_delete(Config* config); void config_delete(Config* config);
bool config_send(int file_descriptor, const Config* config); bool config_send(int file_descriptor, const Config* config);
Config* config_receive(int file_descriptor); 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); void config_parse_ssh_dest(Config* config);
#endif #endif
+11 -706
View File
@@ -1,27 +1,14 @@
#include <dirent.h>
#include <errno.h>
#include <fcntl.h>
#include <libgen.h>
#include <limits.h>
#include <poll.h>
#include <stddef.h>
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <sys/sendfile.h>
#include <sys/stat.h> #include <sys/stat.h>
#include <unistd.h> #include <unistd.h>
#include <time.h>
#include "compression.h"
#include "delta.h"
#include "log.h"
#include "config.h"
#include "data.h" #include "data.h"
#include "delta.h"
#include "file.h" #include "file.h"
#include "file_store.h" #include "file_store.h"
#include "metadata.h" #include "log.h"
#include "protocol.h"
#include "utils.h" #include "utils.h"
bool file_checksum(File* file, uint64_t* checksum) { bool file_checksum(File* file, uint64_t* checksum) {
@@ -42,7 +29,7 @@ File* file_create(const char* path) {
return NULL; return NULL;
File* file = (File*)malloc(sizeof(File)); File* file = (File*)malloc(sizeof(File));
if (file == NULL) { if (file == NULL) {
perror("ERROR: Could not allocate memory for file struct"); log_perror("ERROR: Could not allocate memory for file struct");
return NULL; return NULL;
} }
@@ -82,7 +69,7 @@ void file_destroy(void* item) {
FileMetadata* file_metadata_create(const struct stat* stats) { FileMetadata* file_metadata_create(const struct stat* stats) {
FileMetadata* m = malloc(sizeof(FileMetadata)); FileMetadata* m = malloc(sizeof(FileMetadata));
if (m == NULL) { if (m == NULL) {
perror("ERROR: Could not allocate memory for file metadata"); log_perror("ERROR: Could not allocate memory for file metadata");
return NULL; return NULL;
} }
m->mode = stats->st_mode; m->mode = stats->st_mode;
@@ -109,7 +96,7 @@ bool file_load_data(File* file) {
return true; return true;
file->data->data = malloc(file->data->size); file->data->data = malloc(file->data->size);
if (file->data->data == NULL) { if (file->data->data == NULL) {
perror("Could not allocate memory for file data"); log_perror("Could not allocate memory for file data");
return false; return false;
} }
} }
@@ -124,711 +111,29 @@ bool file_load_data(File* file) {
return true; return true;
} }
bool file_send_single_calls(File* file, int file_descriptor, bool use_metadata, bool file_set_authorized_root(int fd, const char* canonical_path) {
int compression_level, bool send_path) { return file_store_set_authorized_root(fd, canonical_path);
if (!file || !file->path || !file->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;
} }
void file_set_authorized_root(int fd, const char* canonical_path) { bool file_write_to_disk(const char* path, const void* data, unsigned long long data_size,
file_store_set_authorized_root(fd, canonical_path); bool inplace, bool sparse) {
}
static bool path_is_within_root(const char* root, const char* path) {
size_t n = strlen(root);
return strncmp(root, path, n) == 0 && (path[n] == '\0' || path[n] == '/');
}
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;
if (!file || !file->path || !file->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;
}
char* resolved_root = NULL;
const char* actual_root =
(partial_dir && config && config->partial) ? confined_partial : root_directory;
resolved_root = realpath(actual_root, NULL);
if (resolved_root == NULL) {
if (mkdir_r(actual_root)) {
resolved_root = realpath(actual_root, NULL);
}
}
if (resolved_root == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to resolve destination root: %s", actual_root);
free(confined_backup);
free(confined_partial);
return false;
}
char* resolved_base = realpath(root_directory, NULL);
if (resolved_base == NULL || !path_is_within_root(resolved_base, resolved_root)) {
free(resolved_base);
free(confined_backup);
free(confined_partial);
free(resolved_root);
return false;
}
free(resolved_base);
char* disk_path = path_cat(resolved_root, file->path);
if (disk_path == NULL) {
free(confined_backup);
free(confined_partial);
free(resolved_root);
return false;
}
/* --update is receiver-side policy: never replace a newer destination. */
if (config && config->update) {
struct stat destination_stat;
if (stat(disk_path, &destination_stat) == 0 && file->metadata &&
destination_stat.st_mtime > file->metadata->mtime_sec) {
free(resolved_root);
free(confined_backup);
free(confined_partial);
free(disk_path);
return true;
}
}
if (backup_enabled) {
struct stat backup_stat;
if (stat(disk_path, &backup_stat) == 0) {
char* backup_path = NULL;
if (backup_dir) {
char* resolved_backup_dir = realpath(confined_backup, NULL);
if (!resolved_backup_dir) {
mkdir_r(confined_backup);
resolved_backup_dir = realpath(confined_backup, NULL);
}
if (resolved_backup_dir) {
char* backup_base = realpath(root_directory, NULL);
if (backup_base && path_is_within_root(backup_base, resolved_backup_dir))
backup_path = path_cat(resolved_backup_dir, file->path);
free(backup_base);
free(resolved_backup_dir);
}
}
if (!backup_path) {
size_t path_len = strlen(disk_path);
size_t suffix_len = strlen(backup_suffix);
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) {
char* backup_dir_path = str_dup(backup_path);
if (backup_dir_path) {
const char* bdir = dirname(backup_dir_path);
mkdir_r(bdir);
free(backup_dir_path);
}
if (!file_store_rename_secure(disk_path, backup_path)) {
free(backup_path);
free(resolved_root);
free(confined_backup);
free(confined_partial);
free(disk_path);
return false;
}
free(backup_path);
}
}
}
char* dir_dup = str_dup(disk_path);
if (!dir_dup) {
free(confined_backup);
free(confined_partial);
free(resolved_root);
free(disk_path);
return false;
}
char* dir_str = dirname(dir_dup);
if (!mkdir_r(dir_str)) {
free(dir_dup);
free(confined_backup);
free(confined_partial);
free(resolved_root);
free(disk_path);
return false;
}
char* resolved_dir = realpath(dir_str, NULL);
free(dir_dup);
if (resolved_dir == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to resolve directory for: %s", disk_path);
free(confined_backup);
free(confined_partial);
free(resolved_root);
free(disk_path);
return false;
}
size_t root_len = strlen(resolved_root);
if (strncmp(resolved_dir, resolved_root, root_len) != 0 ||
(resolved_dir[root_len] != '\0' && resolved_dir[root_len] != '/')) {
log_message(LOG_LEVEL_ERROR, "Path escape detected: %s is outside %s", disk_path, actual_root);
free(resolved_dir);
free(confined_backup);
free(confined_partial);
free(resolved_root);
free(disk_path);
return false;
}
free(resolved_dir);
free(resolved_root);
bool ok = file_store_write_secure(disk_path, file->data->data, file->data->size, inplace, sparse,
file->metadata);
free(confined_backup);
free(confined_partial);
free(disk_path);
return ok;
}
static File* receive_delta_file(int fd, const Config* config, const char* check_path,
void* old_data, unsigned long long old_size) {
if (!old_data)
return NULL;
DeltaSignature* sig = delta_signature_create(old_data, old_size, config->delta_block_size);
if (!sig) {
free(old_data);
return NULL;
}
Data* sig_data = delta_signature_serialize(sig);
if (!sig_data) {
delta_signature_destroy(sig);
free(old_data);
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);
return NULL;
}
Status resp;
if (!receive_status(fd, &resp)) {
delta_signature_destroy(sig);
free(old_data);
return NULL;
}
if (resp == STATUS_DELTA_DATA) {
Data* delta_data = receive_data(fd);
if (!delta_data) {
delta_signature_destroy(sig);
free(old_data);
send_status(fd, STATUS_ERROR);
return NULL;
}
Data* raw_delta = delta_data;
if (config->use_compression) {
raw_delta = data_decompress(delta_data);
data_destroy(delta_data);
if (!raw_delta) {
free(old_data);
delta_signature_destroy(sig);
send_status(fd, STATUS_ERROR);
return NULL;
}
}
Delta* delta = delta_deserialize(raw_delta);
data_destroy(raw_delta);
if (!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);
uint64_t new_size = delta->new_file_size;
delta_destroy(delta);
if (!new_data) {
free(old_data);
delta_signature_destroy(sig);
send_status(fd, STATUS_ERROR);
return NULL;
}
File* file = file_create(check_path);
if (!file) {
free(new_data);
free(old_data);
delta_signature_destroy(sig);
send_status(fd, STATUS_ERROR);
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);
send_status(fd, STATUS_ERROR);
return NULL;
}
}
data_destroy(file->data);
file->data = data_create(new_data, (size_t)new_size);
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) {
send_status(fd, STATUS_ERROR);
return NULL;
}
if (config->use_metadata) {
int meta_ok = 1;
file->metadata = metadata_receive(fd, &meta_ok);
if (!meta_ok) {
file_destroy(file);
send_status(fd, STATUS_ERROR);
return NULL;
}
}
Data* file_data = receive_data(fd);
if (file_data == NULL) {
file_destroy(file);
send_status(fd, STATUS_ERROR);
return NULL;
}
if (config->use_compression) {
Data* uncompressed = data_decompress(file_data);
data_destroy(file_data);
if (uncompressed == NULL) {
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);
return NULL;
}
File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
*skipped = false;
char* check_path = receive_str(fd);
if (check_path == NULL) {
send_status(fd, STATUS_ERROR);
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);
send_status(fd, STATUS_ERROR);
return NULL;
}
if (config->checksum && !receive_n_data(fd, &check_checksum, sizeof(check_checksum))) {
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);
send_status(fd, STATUS_ERROR);
return NULL;
}
char* full_path = path_cat(config->receive_root_directory, check_path);
struct stat st;
bool has_old_file = false;
int old_fd = -1;
if (full_path) {
char* leaf = NULL;
int parent_fd = file_store_open_secure_parent(full_path, &leaf);
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_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 &&
delta_should_attempt(old_size, check_size, config->delta_max_file_size);
if (try_delta) {
File* delta_file = receive_delta_file(fd, config, check_path, old_data, old_size);
old_data = NULL; /* receive_delta_file consumes the snapshot on every path */
if (delta_file) {
free(full_path);
free(check_path);
return delta_file;
}
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) {
send_status(fd, STATUS_ERROR);
return NULL;
}
if (config->use_metadata) {
int meta_ok = 1;
file->metadata = metadata_receive(fd, &meta_ok);
if (!meta_ok) {
file_destroy(file);
send_status(fd, STATUS_ERROR);
return NULL;
}
}
Data* file_data = receive_data(fd);
if (file_data == NULL) {
file_destroy(file);
send_status(fd, STATUS_ERROR);
return NULL;
}
if (config->use_compression) {
Data* uncompressed = data_decompress(file_data);
data_destroy(file_data);
if (uncompressed == NULL) {
file_destroy(file);
send_status(fd, STATUS_ERROR);
return NULL;
}
file_data = uncompressed;
}
data_destroy(file->data);
file->data = file_data;
return file;
}
bool to_disk(const char* path, const void* data, unsigned long long data_size, bool inplace,
bool sparse) {
if (!path || (!data && data_size != 0) || has_path_traversal(path)) if (!path || (!data && data_size != 0) || has_path_traversal(path))
return false; return false;
return file_store_write_secure(path, data, data_size, inplace, sparse, NULL); return file_store_write_secure(path, data, data_size, inplace, sparse, NULL);
} }
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) {
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;
perror("sendfile failed");
close(fd);
return false;
}
if (sent == 0) {
close(fd);
return false;
}
}
close(fd);
return true;
}
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(file_descriptor);
if (file_data == NULL) {
file_destroy(file);
return NULL;
}
if (config->use_compression) {
Data* file_data_uncompressed = data_decompress(file_data);
data_destroy(file_data);
if (file_data_uncompressed == NULL) {
file_destroy(file);
return NULL;
}
file_data = file_data_uncompressed;
}
data_destroy(file->data);
file->data = file_data;
return file;
}
size_t file_content_to_buffer(File* file) { size_t file_content_to_buffer(File* file) {
FILE* file_pointer = fopen(file->path, "rb"); FILE* file_pointer = fopen(file->path, "rb");
if (file_pointer == NULL) { if (file_pointer == NULL) {
perror("Could not open the file!"); log_perror("Could not open the file!");
return 0; return 0;
} }
size_t bytes_read = fread(file->data->data, 1, file->data->size, file_pointer); size_t bytes_read = fread(file->data->data, 1, file->data->size, file_pointer);
if (bytes_read != (size_t)file->data->size) { if (bytes_read != (size_t)file->data->size) {
fclose(file_pointer); fclose(file_pointer);
perror("Read unexpected number of bytes from File!"); log_perror("Read unexpected number of bytes from File!");
return 0; return 0;
} }
fclose(file_pointer); fclose(file_pointer);
return bytes_read; return bytes_read;
} }
int receive_manifest(int fd, const Config* config, int* next_status) {
int received_status = STATUS_ERROR;
int* status_out = next_status ? next_status : &received_status;
int count;
if (!receive_int(fd, &count))
return -1;
if (count < 0 || count > MAX_MANIFEST_ENTRIES)
return -1;
ArrayList* manifest = array_list_create(free);
if (!manifest)
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);
return -1;
}
}
if (!receive_status(fd, status_out)) {
array_list_delete(manifest);
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);
return *status_out == STATUS_FINISHED ? 0 : -1;
}
fprintf(stderr, "Deleting files not in manifest...\n");
bool deletion_ok = delete_extras(config->receive_root_directory, manifest);
array_list_delete(manifest);
return deletion_ok ? 0 : -1;
}
+8 -30
View File
@@ -1,45 +1,23 @@
#ifndef FILE_H #ifndef FILE_H
#define FILE_H #define FILE_H
#include "config.h" #include "file_send.h"
#include "data.h" #include "file_receive.h"
#include "file_types.h"
#include <stdbool.h> #include <stdbool.h>
#include <sys/stat.h> #include <stdint.h>
typedef enum { FILE_TYPE_REGULAR, FILE_TYPE_SYMLINK, FILE_TYPE_DIR } FileType; /* File/FileMetadata lifecycle and local disk helpers. */
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* file_create(const char* path); File* file_create(const char* path);
void file_destroy(void* item); void file_destroy(void* item);
bool file_load_data(File* file); bool file_load_data(File* file);
bool file_checksum(File* file, uint64_t* checksum); 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); size_t file_content_to_buffer(File* file);
FileMetadata* file_metadata_create(const struct stat* stats); FileMetadata* file_metadata_create(const struct stat* stats);
void file_metadata_destroy(void* metadata); void file_metadata_destroy(void* metadata);
bool to_disk(const char* path, const void* data, unsigned long long data_size, bool inplace, bool file_write_to_disk(const char* path, const void* data, unsigned long long data_size,
bool sparse); bool inplace, bool sparse);
bool file_save_to_disk(const char* root_directory, const File* file, const Config* config); bool file_set_authorized_root(int fd, const char* canonical_path);
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);
#endif #endif
+588
View File
@@ -0,0 +1,588 @@
#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 "compression.h"
#include "config.h"
#include "data.h"
#include "delta.h"
#include "file.h"
#include "file_store.h"
#include "log.h"
#include "metadata.h"
#include "protocol.h"
#include "utils.h"
static bool path_is_within_root(const char* root, const char* path) {
size_t n = strlen(root);
return strncmp(root, path, n) == 0 && (path[n] == '\0' || path[n] == '/');
}
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;
if (!file || !file->path || !file->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;
}
char* resolved_root = NULL;
const char* actual_root =
(partial_dir && config && config->partial) ? confined_partial : root_directory;
resolved_root = realpath(actual_root, NULL);
if (resolved_root == NULL) {
if (mkdir_r(actual_root)) {
resolved_root = realpath(actual_root, NULL);
}
}
if (resolved_root == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to resolve destination root: %s", actual_root);
free(confined_backup);
free(confined_partial);
return false;
}
char* resolved_base = realpath(root_directory, NULL);
if (resolved_base == NULL || !path_is_within_root(resolved_base, resolved_root)) {
free(resolved_base);
free(confined_backup);
free(confined_partial);
free(resolved_root);
return false;
}
free(resolved_base);
char* disk_path = path_cat(resolved_root, file->path);
if (disk_path == NULL) {
free(confined_backup);
free(confined_partial);
free(resolved_root);
return false;
}
/* --update is receiver-side policy: never replace a newer destination. */
if (config && config->update) {
struct stat destination_stat;
if (stat(disk_path, &destination_stat) == 0 && file->metadata &&
destination_stat.st_mtime > file->metadata->mtime_sec) {
free(resolved_root);
free(confined_backup);
free(confined_partial);
free(disk_path);
return true;
}
}
if (backup_enabled) {
struct stat backup_stat;
if (stat(disk_path, &backup_stat) == 0) {
char* backup_path = NULL;
if (backup_dir) {
char* resolved_backup_dir = realpath(confined_backup, NULL);
if (!resolved_backup_dir) {
mkdir_r(confined_backup);
resolved_backup_dir = realpath(confined_backup, NULL);
}
if (resolved_backup_dir) {
char* backup_base = realpath(root_directory, NULL);
if (backup_base && path_is_within_root(backup_base, resolved_backup_dir))
backup_path = path_cat(resolved_backup_dir, file->path);
free(backup_base);
free(resolved_backup_dir);
}
}
if (!backup_path) {
size_t path_len = strlen(disk_path);
size_t suffix_len = strlen(backup_suffix);
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) {
char* backup_dir_path = str_dup(backup_path);
if (backup_dir_path) {
const char* bdir = dirname(backup_dir_path);
mkdir_r(bdir);
free(backup_dir_path);
}
if (!file_store_rename_secure(disk_path, backup_path)) {
free(backup_path);
free(resolved_root);
free(confined_backup);
free(confined_partial);
free(disk_path);
return false;
}
free(backup_path);
}
}
}
char* dir_dup = str_dup(disk_path);
if (!dir_dup) {
free(confined_backup);
free(confined_partial);
free(resolved_root);
free(disk_path);
return false;
}
char* dir_str = dirname(dir_dup);
if (!mkdir_r(dir_str)) {
free(dir_dup);
free(confined_backup);
free(confined_partial);
free(resolved_root);
free(disk_path);
return false;
}
char* resolved_dir = realpath(dir_str, NULL);
free(dir_dup);
if (resolved_dir == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to resolve directory for: %s", disk_path);
free(confined_backup);
free(confined_partial);
free(resolved_root);
free(disk_path);
return false;
}
size_t root_len = strlen(resolved_root);
if (strncmp(resolved_dir, resolved_root, root_len) != 0 ||
(resolved_dir[root_len] != '\0' && resolved_dir[root_len] != '/')) {
log_message(LOG_LEVEL_ERROR, "Path escape detected: %s is outside %s", disk_path, actual_root);
free(resolved_dir);
free(confined_backup);
free(confined_partial);
free(resolved_root);
free(disk_path);
return false;
}
free(resolved_dir);
free(resolved_root);
bool ok = file_store_write_secure(disk_path, file->data->data, file->data->size, inplace, sparse,
file->metadata);
free(confined_backup);
free(confined_partial);
free(disk_path);
return ok;
}
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(fd);
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(delta_data);
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;
}
void* new_data = delta_apply(old_data, old_size, delta, config->delta_block_size);
uint64_t new_size = delta->new_file_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_destroy(file->data);
file->data = data_create(new_data, (size_t)new_size);
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(fd);
if (file_data == NULL) {
file_destroy(file);
*failed = true;
return NULL;
}
if (config->use_compression) {
Data* uncompressed = data_decompress(file_data);
data_destroy(file_data);
if (uncompressed == NULL) {
file_destroy(file);
*failed = true;
return NULL;
}
file_data = uncompressed;
}
data_destroy(file->data);
file->data = file_data;
return file;
}
delta_signature_destroy(sig);
free(old_data);
*failed = true;
return NULL;
}
File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
*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 (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);
struct stat st;
bool has_old_file = false;
int old_fd = -1;
if (full_path) {
char* leaf = NULL;
int parent_fd = file_store_open_secure_parent(full_path, &leaf);
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_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 &&
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(fd);
if (file_data == NULL) {
file_destroy(file);
return NULL;
}
if (config->use_compression) {
Data* uncompressed = data_decompress(file_data);
data_destroy(file_data);
if (uncompressed == NULL) {
file_destroy(file);
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(file_descriptor);
if (file_data == NULL) {
file_destroy(file);
return NULL;
}
if (config->use_compression) {
Data* file_data_uncompressed = data_decompress(file_data);
data_destroy(file_data);
if (file_data_uncompressed == NULL) {
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) {
int received_status = STATUS_ERROR;
int* status_out = next_status ? next_status : &received_status;
int count;
if (!receive_int(fd, &count))
return -1;
if (count < 0 || count > MAX_MANIFEST_ENTRIES)
return -1;
ArrayList* manifest = array_list_create(free);
if (!manifest)
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);
return -1;
}
}
if (!receive_status(fd, status_out)) {
array_list_delete(manifest);
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);
return *status_out == STATUS_FINISHED ? 0 : -1;
}
fprintf(stderr, "Deleting files not in manifest...\n");
bool deletion_ok = delete_extras(config->receive_root_directory, manifest);
array_list_delete(manifest);
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)
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
+21 -5
View File
@@ -20,10 +20,14 @@ static bool path_is_within_root(const char* root, const char* path) {
(path[root_length] == '\0' || path[root_length] == '/'); (path[root_length] == '\0' || path[root_length] == '/');
} }
void file_store_set_authorized_root(int fd, const char* canonical_path) { bool file_store_set_authorized_root(int fd, const char* canonical_path) {
authorized_root_fd = fd; char* new_path = canonical_path ? str_dup(canonical_path) : NULL;
if (canonical_path && !new_path)
return false;
free(authorized_root_path); free(authorized_root_path);
authorized_root_path = canonical_path ? str_dup(canonical_path) : NULL; authorized_root_path = new_path;
authorized_root_fd = fd;
return true;
} }
int file_store_open_secure_parent(const char* path, char** leaf_out) { int file_store_open_secure_parent(const char* path, char** leaf_out) {
@@ -130,9 +134,20 @@ bool file_store_write_secure(const char* path, const void* data, unsigned long l
ok = file_restore_metadata_fd(fd, metadata); ok = file_restore_metadata_fd(fd, metadata);
} }
} else { } else {
char tmp[NAME_MAX]; 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) { for (unsigned int i = 0; i < 100 && !ok; ++i) {
snprintf(tmp, sizeof(tmp), ".%s.tmp.%ld.%u", leaf, (long)getpid(), 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); fd = openat(dirfd, tmp, O_WRONLY | O_CREAT | O_EXCL | O_CLOEXEC | O_NOFOLLOW, 0600);
if (fd < 0) if (fd < 0)
continue; continue;
@@ -150,6 +165,7 @@ bool file_store_write_secure(const char* path, const void* data, unsigned long l
if (!ok) if (!ok)
unlinkat(dirfd, tmp, 0); unlinkat(dirfd, tmp, 0);
} }
free(tmp);
} }
if (fd >= 0) if (fd >= 0)
close(fd); close(fd);
+1 -1
View File
@@ -4,7 +4,7 @@
#include "file.h" #include "file.h"
#include <stdbool.h> #include <stdbool.h>
void file_store_set_authorized_root(int fd, const char* canonical_path); bool file_store_set_authorized_root(int fd, const char* canonical_path);
int file_store_open_secure_parent(const char* path, char** leaf_out); 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_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 file_store_write_secure(const char* path, const void* data, unsigned long long data_size,
+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 "log.h"
#include <errno.h>
#include <stdarg.h> #include <stdarg.h>
#include <stdio.h> #include <stdio.h>
#include <string.h>
#include <time.h> #include <time.h>
static const char* log_level_strings[] = {"DEBUG", "INFO", "WARN", "ERROR"}; 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); 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; 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_message(LogLevel log_level, const char* message, ...);
void log_perror(const char* context);
void set_log_level(LogLevel level); void set_log_level(LogLevel level);
void log_set_file(FILE* fp); void log_set_file(FILE* fp);
+11 -2
View File
@@ -55,7 +55,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
return context; return context;
fail: fail:
perror("Error initializing synchronization objects"); log_perror("Error initializing synchronization objects");
if (init >= 6) if (init >= 6)
cnd_destroy(&context->condition_not_empty_loader); cnd_destroy(&context->condition_not_empty_loader);
if (init >= 5) if (init >= 5)
@@ -116,7 +116,7 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
return context; return context;
fail: fail:
perror("Error initializing synchronization objects"); log_perror("Error initializing synchronization objects");
if (init >= 3) if (init >= 3)
cnd_destroy(&context->condition_not_empty); cnd_destroy(&context->condition_not_empty);
if (init >= 2) if (init >= 2)
@@ -183,6 +183,15 @@ int write_thread(void* pipeline_context) {
bool save_to_disk = context->config->save_to_disk; bool save_to_disk = context->config->save_to_disk;
char* root_directory = str_dup(context->config->receive_root_directory); char* root_directory = str_dup(context->config->receive_root_directory);
mtx_unlock(&context->mutex); 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) { while (true) {
File* file = File* file =
+30 -37
View File
@@ -19,6 +19,7 @@ static __thread int io_read_fd = -1;
static __thread int io_write_fd = -1; static __thread int io_write_fd = -1;
static __thread SSL* io_ssl; static __thread SSL* io_ssl;
static __thread ProtocolSession* bound_session; 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 unsigned long long io_bwlimit = 0;
static long long bw_tokens = 0; static long long bw_tokens = 0;
@@ -26,8 +27,6 @@ static struct timespec bw_last_refill = {0, 0};
static mtx_t bw_mutex; static mtx_t bw_mutex;
static once_flag bw_mutex_once = ONCE_FLAG_INIT; static once_flag bw_mutex_once = ONCE_FLAG_INIT;
static __thread unsigned long long total_allocated_bytes = 0;
void io_set_fds(int read_fd, int write_fd) { void io_set_fds(int read_fd, int write_fd) {
bound_session = NULL; bound_session = NULL;
io_read_fd = read_fd; io_read_fd = read_fd;
@@ -35,7 +34,11 @@ void io_set_fds(int read_fd, int write_fd) {
/* A descriptor switch starts a new transport; never reuse a TLS object /* A descriptor switch starts a new transport; never reuse a TLS object
belonging to a previous connection or test pipe. */ belonging to a previous connection or test pipe. */
io_ssl = NULL; 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, io_bwlimit);
} }
void protocol_session_init(ProtocolSession* session, int read_fd, int write_fd) { void protocol_session_init(ProtocolSession* session, int read_fd, int write_fd) {
@@ -126,32 +129,30 @@ SSL* io_get_ssl(void) {
return io_ssl; return io_ssl;
} }
static ProtocolSession* legacy_session(void) { static ProtocolSession* legacy_session(int read_fd, int write_fd) {
static __thread ProtocolSession session;
if (bound_session) if (bound_session)
return bound_session; return bound_session;
session.read_fd = io_read_fd; int target_read_fd = io_read_fd != -1 ? io_read_fd : read_fd;
session.write_fd = io_write_fd; int target_write_fd = io_write_fd != -1 ? io_write_fd : write_fd;
session.ssl = io_ssl; if (legacy_io_session.read_fd != target_read_fd ||
session.bwlimit = io_bwlimit; legacy_io_session.write_fd != target_write_fd) {
session.bw_tokens = (unsigned long long)(bw_tokens < 0 ? 0 : bw_tokens); legacy_io_session.read_fd = target_read_fd;
session.bw_last_refill_sec = bw_last_refill.tv_sec; legacy_io_session.write_fd = target_write_fd;
session.bw_last_refill_nsec = bw_last_refill.tv_nsec; legacy_io_session.total_allocated_bytes = 0;
return &session; protocol_session_set_bwlimit(&legacy_io_session, io_bwlimit);
} else if (legacy_io_session.bwlimit != io_bwlimit) {
protocol_session_set_bwlimit(&legacy_io_session, io_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) { bool send_n_data(int file_descriptor, const void* data, size_t data_size) {
ProtocolSession* session = legacy_session(); return protocol_send_n_data(legacy_session(-1, file_descriptor), data, data_size);
if (session->write_fd == -1)
session->write_fd = file_descriptor;
return protocol_send_n_data(session, data, data_size);
} }
bool receive_n_data(int file_descriptor, void* data, size_t data_size) { bool receive_n_data(int file_descriptor, void* data, size_t data_size) {
ProtocolSession* session = legacy_session(); return protocol_receive_n_data(legacy_session(file_descriptor, -1), data, data_size);
if (session->read_fd == -1)
session->read_fd = file_descriptor;
return protocol_receive_n_data(session, data, data_size);
} }
static int deadline_remaining_ms(const struct timespec* deadline) { static int deadline_remaining_ms(const struct timespec* deadline) {
@@ -403,34 +404,26 @@ bool protocol_receive_status(ProtocolSession* session, Status* status) {
} }
bool send_str(int fd, const char* data) { bool send_str(int fd, const char* data) {
(void)fd; return protocol_send_str(legacy_session(-1, fd), data);
return protocol_send_str(legacy_session(), data);
} }
char* receive_str(int fd) { char* receive_str(int fd) {
(void)fd; return protocol_receive_str(legacy_session(fd, -1));
return protocol_receive_str(legacy_session());
} }
bool send_data(int fd, const Data* data) { bool send_data(int fd, const Data* data) {
(void)fd; return protocol_send_data(legacy_session(-1, fd), data);
return protocol_send_data(legacy_session(), data);
} }
Data* receive_data(int fd) { Data* receive_data(int fd) {
(void)fd; return protocol_receive_data(legacy_session(fd, -1));
return protocol_receive_data(legacy_session());
} }
bool send_int(int fd, int data) { bool send_int(int fd, int data) {
(void)fd; return protocol_send_int(legacy_session(-1, fd), data);
return protocol_send_int(legacy_session(), data);
} }
bool receive_int(int fd, int* data) { bool receive_int(int fd, int* data) {
(void)fd; return protocol_receive_int(legacy_session(fd, -1), data);
return protocol_receive_int(legacy_session(), data);
} }
bool send_status(int fd, Status status) { bool send_status(int fd, Status status) {
(void)fd; return protocol_send_status(legacy_session(-1, fd), status);
return protocol_send_status(legacy_session(), status);
} }
bool receive_status(int fd, Status* status) { bool receive_status(int fd, Status* status) {
(void)fd; return protocol_receive_status(legacy_session(fd, -1), status);
return protocol_receive_status(legacy_session(), status);
} }
+4 -3
View File
@@ -1,3 +1,4 @@
#include "log.h"
#include <stdbool.h> #include <stdbool.h>
#include <limits.h> #include <limits.h>
#include <stdio.h> #include <stdio.h>
@@ -13,7 +14,7 @@ Queue* queue_create(int capacity, void (*destroyer)(void* item)) {
Queue* queue = (Queue*)malloc(sizeof(Queue)); Queue* queue = (Queue*)malloc(sizeof(Queue));
if (queue == NULL) { if (queue == NULL) {
perror("ERROR: Could not allocate memory for queue structure"); log_perror("ERROR: Could not allocate memory for queue structure");
return NULL; return NULL;
} }
@@ -72,7 +73,7 @@ static bool queue_double_capacity(Queue* queue) {
new_capacity = 100; new_capacity = 100;
void** new_items = malloc(new_capacity * sizeof(void*)); void** new_items = malloc(new_capacity * sizeof(void*));
if (new_items == NULL) { 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; return false;
} }
for (int i = 0; i < queue->size; i++) 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) { void* queue_dequeue(Queue* queue) {
if (queue == NULL || queue_is_empty(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; return NULL;
} }
+5 -4
View File
@@ -1,3 +1,4 @@
#include "log.h"
#include "transport_ssh.h" #include "transport_ssh.h"
#include "utils.h" #include "utils.h"
#include <fcntl.h> #include <fcntl.h>
@@ -76,7 +77,7 @@ Client* client_connect_ssh(const char* destination, int port, const char* server
int sv[2]; int sv[2];
if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) { if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) < 0) {
perror("socketpair failed"); log_perror("socketpair failed");
remote_dest_destroy(&r); remote_dest_destroy(&r);
return NULL; return NULL;
} }
@@ -89,7 +90,7 @@ Client* client_connect_ssh(const char* destination, int port, const char* server
int exec_pipe[2]; int exec_pipe[2];
if (pipe(exec_pipe) < 0) { if (pipe(exec_pipe) < 0) {
perror("pipe failed"); log_perror("pipe failed");
close(sv[0]); close(sv[0]);
close(sv[1]); close(sv[1]);
remote_dest_destroy(&r); remote_dest_destroy(&r);
@@ -98,7 +99,7 @@ Client* client_connect_ssh(const char* destination, int port, const char* server
pid_t pid = fork(); pid_t pid = fork();
if (pid < 0) { if (pid < 0) {
perror("fork failed"); log_perror("fork failed");
close(sv[0]); close(sv[0]);
close(sv[1]); close(sv[1]);
close(exec_pipe[0]); 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++] = "--stdio";
ssh_argv[ac] = NULL; ssh_argv[ac] = NULL;
execvp("ssh", ssh_argv); execvp("ssh", ssh_argv);
perror("exec of ssh failed"); log_perror("exec of ssh failed");
ssize_t wret = write(exec_pipe[1], "x", 1); ssize_t wret = write(exec_pipe[1], "x", 1);
(void)wret; (void)wret;
_exit(1); _exit(1);
+11 -7
View File
@@ -30,20 +30,20 @@ static void sigchld_handler(int sig) {
Server* server_create(int port) { Server* server_create(int port) {
Server* server = (Server*)malloc(sizeof(Server)); Server* server = (Server*)malloc(sizeof(Server));
if (server == NULL) { if (server == NULL) {
perror("Could not allocate space for Server"); log_perror("Could not allocate space for Server");
return NULL; return NULL;
} }
int file_descriptor = socket(AF_INET, SOCK_STREAM, 0); int file_descriptor = socket(AF_INET, SOCK_STREAM, 0);
if (file_descriptor < 0) { if (file_descriptor < 0) {
perror("Could not create Socket!"); log_perror("Could not create Socket!");
free(server); free(server);
return NULL; return NULL;
} }
server->file_descriptor = file_descriptor; server->file_descriptor = file_descriptor;
int opt = 1; int opt = 1;
if (setsockopt(server->file_descriptor, SOL_SOCKET, SO_REUSEADDR, &opt, sizeof(opt))) { 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); close(server->file_descriptor);
free(server); free(server);
return NULL; return NULL;
@@ -59,7 +59,7 @@ Server* server_create(int port) {
if (bind(server->file_descriptor, (struct sockaddr*)&server->address, server->address_length) < if (bind(server->file_descriptor, (struct sockaddr*)&server->address, server->address_length) <
0) { 0) {
perror("Could not bind server"); log_perror("Could not bind server");
close(server->file_descriptor); close(server->file_descriptor);
free(server); free(server);
return NULL; 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, static void accept_loop(Server* server, void (*child_fn)(int, void*), void* child_ctx,
const char* log_fmt) { const char* log_fmt) {
if (listen(server->file_descriptor, SOMAXCONN) < 0) { if (listen(server->file_descriptor, SOMAXCONN) < 0) {
perror("Could not listen on port!"); log_perror("Could not listen on port!");
return; return;
} }
signal(SIGCHLD, sigchld_handler); 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); socklen_t client_len = sizeof(client_addr);
int fd = accept(server->file_descriptor, (struct sockaddr*)&client_addr, &client_len); int fd = accept(server->file_descriptor, (struct sockaddr*)&client_addr, &client_len);
if (fd < 0) { if (fd < 0) {
perror("Could not accept the connection"); log_perror("Could not accept the connection");
continue; continue;
} }
tcp_apply_socket_timeout(fd); tcp_apply_socket_timeout(fd);
@@ -151,6 +151,10 @@ int tcp_get_contimeout_sec(void) {
return g_contimeout_sec; return g_contimeout_sec;
} }
int tcp_get_timeout_sec(void) {
return g_timeout_sec;
}
static void tcp_apply_socket_timeout(int fd) { static void tcp_apply_socket_timeout(int fd) {
struct timeval tv; struct timeval tv;
tv.tv_sec = g_timeout_sec; tv.tv_sec = g_timeout_sec;
@@ -218,7 +222,7 @@ bool tcp_connect_socket(Client* client, char* host, int port) {
freeaddrinfo(result); freeaddrinfo(result);
if (!connected) { if (!connected) {
perror("Could not connect to Server!"); log_perror("Could not connect to Server!");
return false; return false;
} }
+1
View File
@@ -35,5 +35,6 @@ void client_disconnect(Client* client);
void client_delete(Client* client); void client_delete(Client* client);
void tcp_set_timeouts(int timeout_sec, int contimeout_sec); void tcp_set_timeouts(int timeout_sec, int contimeout_sec);
int tcp_get_contimeout_sec(void); int tcp_get_contimeout_sec(void);
int tcp_get_timeout_sec(void);
#endif #endif
+4 -1
View File
@@ -9,6 +9,7 @@
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <time.h>
#include <unistd.h> #include <unistd.h>
bool tls_global_init(void) { bool tls_global_init(void) {
@@ -92,6 +93,7 @@ static SSL* wrap_fd_with_ssl(int fd, SSL_CTX* ctx, bool is_server, const char* h
} }
// Retry SSL_accept/SSL_connect on WANT_READ/WANT_WRITE (non-blocking handshake) // 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; int ret;
do { do {
if (is_server) if (is_server)
@@ -101,7 +103,8 @@ static SSL* wrap_fd_with_ssl(int fd, SSL_CTX* ctx, bool is_server, const char* h
if (ret <= 0) { if (ret <= 0) {
int ssl_err = SSL_get_error(ssl, ret); 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; continue;
log_message(LOG_LEVEL_ERROR, "SSL %s failed", is_server ? "accept" : "connect"); log_message(LOG_LEVEL_ERROR, "SSL %s failed", is_server ? "accept" : "connect");
log_ssl_errors(); log_ssl_errors();
+2 -1
View File
@@ -1,3 +1,4 @@
#include "log.h"
#include "utils.h" #include "utils.h"
#include "array_list.h" #include "array_list.h"
#include "libgen.h" #include "libgen.h"
@@ -54,7 +55,7 @@ bool mkdir_r(const char* path) {
struct stat st; struct stat st;
if (stat(path_current, &st) != 0) { if (stat(path_current, &st) != 0) {
if (mkdir(path_current, 0755) != 0) { if (mkdir(path_current, 0755) != 0) {
perror("Could not create directory"); log_perror("Could not create directory");
ok = false; ok = false;
break; break;
} }
+3 -3
View File
@@ -11,7 +11,7 @@ static void test_file_operations() {
char* test_content = "Hello, Chunk System!"; char* test_content = "Hello, Chunk System!";
unsigned long long test_len = strlen(test_content); 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); File* f = file_create(test_path);
EXPECT_NOT_NULL(f); EXPECT_NOT_NULL(f);
@@ -43,8 +43,8 @@ static void test_chunk_operations() {
char* content2 = "chunk item number 2"; char* content2 = "chunk item number 2";
unsigned long long len2 = strlen(content2); unsigned long long len2 = strlen(content2);
to_disk(path1, content1, len1, false, false); file_write_to_disk(path1, content1, len1, false, false);
to_disk(path2, content2, len2, false, false); file_write_to_disk(path2, content2, len2, false, false);
struct stat st1, st2; struct stat st1, st2;
stat(path1, &st1); stat(path1, &st1);
+2 -2
View File
@@ -62,8 +62,8 @@ static void test_chunk_compress_decompress_roundtrip() {
char* content2 = "chunk compression test file 2 with more data"; char* content2 = "chunk compression test file 2 with more data";
unsigned long long len2 = strlen(content2); unsigned long long len2 = strlen(content2);
to_disk(path1, content1, len1, false, false); file_write_to_disk(path1, content1, len1, false, false);
to_disk(path2, content2, len2, false, false); file_write_to_disk(path2, content2, len2, false, false);
struct stat st1, st2; struct stat st1, st2;
EXPECT_EQ_INT(stat(path1, &st1), 0); EXPECT_EQ_INT(stat(path1, &st1), 0);
+15 -15
View File
@@ -244,26 +244,26 @@ static void test_config_receive_truncated() {
close(p[1]); close(p[1]);
} }
static void test_is_remote_dest() { static void test_config_is_remote_dest() {
/* Valid SSH-style destinations */ /* Valid SSH-style destinations */
EXPECT_TRUE(is_remote_dest("user@host:/path")); EXPECT_TRUE(config_is_remote_dest("user@host:/path"));
EXPECT_TRUE(is_remote_dest("host:/path")); EXPECT_TRUE(config_is_remote_dest("host:/path"));
EXPECT_TRUE(is_remote_dest("user@192.168.1.1:/remote/path")); EXPECT_TRUE(config_is_remote_dest("user@192.168.1.1:/remote/path"));
/* Invalid destinations */ /* Invalid destinations */
EXPECT_FALSE(is_remote_dest(NULL)); EXPECT_FALSE(config_is_remote_dest(NULL));
EXPECT_FALSE(is_remote_dest("")); EXPECT_FALSE(config_is_remote_dest(""));
EXPECT_FALSE(is_remote_dest(":")); EXPECT_FALSE(config_is_remote_dest(":"));
EXPECT_FALSE(is_remote_dest("/local/path")); EXPECT_FALSE(config_is_remote_dest("/local/path"));
EXPECT_FALSE(is_remote_dest("relative/path")); EXPECT_FALSE(config_is_remote_dest("relative/path"));
/* C:/windows/path is treated as remote (colon with no preceding slash) */ /* 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 */ /* Edge cases */
EXPECT_FALSE(is_remote_dest("noslash")); EXPECT_FALSE(config_is_remote_dest("noslash"));
EXPECT_FALSE(is_remote_dest("/")); EXPECT_FALSE(config_is_remote_dest("/"));
EXPECT_TRUE(is_remote_dest("host:")); EXPECT_TRUE(config_is_remote_dest("host:"));
EXPECT_TRUE(is_remote_dest("user@host:")); EXPECT_TRUE(config_is_remote_dest("user@host:"));
} }
void test_config() { void test_config() {
@@ -278,5 +278,5 @@ void test_config() {
test_config_send_receive_version_mismatch(); test_config_send_receive_version_mismatch();
test_config_receive_truncated(); 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() { static void test_file_load_data() {
const char* content = "Hello Load Test"; 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; struct stat st;
EXPECT_EQ_INT(stat("test_file_load_data.txt", &st), 0); 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"); rmdir("test_save_tmp");
} }
static void test_to_disk_basic() { static void test_file_write_to_disk_basic() {
const char* content = "Basic to_disk test"; const char* content = "Basic file_write_to_disk test";
EXPECT_TRUE(to_disk("test_to_disk_basic.txt", content, strlen(content), false, false)); EXPECT_TRUE(file_write_to_disk("test_file_write_to_disk_basic.txt", content, strlen(content),
false, false));
struct stat st; 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)); 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); EXPECT_NOT_NULL(fp);
char buf[100]; char buf[100];
size_t nread = fread(buf, 1, sizeof(buf), fp); 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((int)nread, (int)strlen(content));
EXPECT_EQ_INT(memcmp(buf, content, strlen(content)), 0); 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"; 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; struct stat st;
EXPECT_EQ_INT(stat("test_nested_tmp/nested/file.txt", &st), 0); 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"); rmdir("test_nested_tmp");
} }
static void test_to_disk_does_not_follow_symlink() { static void test_file_write_to_disk_does_not_follow_symlink() {
const char* outside = "test_to_disk_outside.txt"; const char* outside = "test_file_write_to_disk_outside.txt";
const char* link = "test_to_disk_link.txt"; const char* link = "test_file_write_to_disk_link.txt";
const char* content = "confined"; const char* content = "confined";
unlink(outside); unlink(outside);
unlink(link); 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_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"); FILE* fp = fopen(outside, "rb");
char buf[16] = {0}; char buf[16] = {0};
EXPECT_NOT_NULL(fp); EXPECT_NOT_NULL(fp);
@@ -151,7 +154,7 @@ static void test_to_disk_does_not_follow_symlink() {
static void test_file_content_to_buffer() { static void test_file_content_to_buffer() {
const char* content = "Buffer content test"; 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"); File* f = file_create("test_buffer_file.txt");
EXPECT_NOT_NULL(f); EXPECT_NOT_NULL(f);
@@ -273,7 +276,7 @@ static void test_file_send_no_path() {
} }
static void test_file_metadata_create() { 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; struct stat st;
EXPECT_EQ_INT(stat("test_meta_file.txt", &st), 0); 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 */ /* Create a real file on disk so we can have metadata */
const char* content = "File with metadata"; const char* content = "File with metadata";
size_t len = strlen(content); 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; struct stat st;
EXPECT_EQ_INT(stat("test_meta_send.txt", &st), 0); 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();
test_file_load_data_missing_file(); test_file_load_data_missing_file();
test_file_save_to_disk(); test_file_save_to_disk();
test_to_disk_basic(); test_file_write_to_disk_basic();
test_to_disk_creates_dirs(); test_file_write_to_disk_creates_dirs();
test_to_disk_does_not_follow_symlink(); test_file_write_to_disk_does_not_follow_symlink();
test_file_content_to_buffer(); test_file_content_to_buffer();
test_file_save_to_disk_path_traversal(); test_file_save_to_disk_path_traversal();
test_file_save_to_disk_deep_traversal(); test_file_save_to_disk_deep_traversal();
+4 -4
View File
@@ -15,7 +15,7 @@
static void test_sendfile_basic() { static void test_sendfile_basic() {
const char* content = "Hello from sendfile test!"; const char* content = "Hello from sendfile test!";
size_t len = strlen(content); 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"); File* file = file_create("test_sendfile_basic.txt");
EXPECT_NOT_NULL(file); EXPECT_NOT_NULL(file);
@@ -77,7 +77,7 @@ static void test_sendfile_basic() {
static void test_sendfile_empty_file() { static void test_sendfile_empty_file() {
const char* content = ""; const char* content = "";
size_t len = 0; 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"); File* file = file_create("test_sendfile_empty.txt");
EXPECT_NOT_NULL(file); EXPECT_NOT_NULL(file);
@@ -156,7 +156,7 @@ static void test_sendfile_missing_file() {
static void test_sendfile_compression_fallback() { static void test_sendfile_compression_fallback() {
const char* content = "Compression fallback content"; const char* content = "Compression fallback content";
size_t len = strlen(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; struct stat st;
EXPECT_EQ_INT(stat("test_sendfile_comp.txt", &st), 0); 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() { static void test_sendfile_no_path() {
const char* content = "No path sendfile test"; const char* content = "No path sendfile test";
size_t len = strlen(content); 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"); File* file = file_create("test_sendfile_nopath.txt");
EXPECT_NOT_NULL(file); 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 */ /* Smoke test for metadata_from_buf fuzz target */
static void test_fuzz_metadata_from_buf() { static void test_fuzz_metadata_from_buf() {
/* Create a real file to get metadata from */ /* 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; struct stat st;
EXPECT_EQ_INT(stat("fuzz_meta_test.txt", &st), 0); EXPECT_EQ_INT(stat("fuzz_meta_test.txt", &st), 0);
+1 -1
View File
@@ -125,7 +125,7 @@ static void test_metadata_rejects_invalid_values() {
static void test_file_restore_metadata() { static void test_file_restore_metadata() {
const char* path = "temp_meta_restore_test.txt"; const char* path = "temp_meta_restore_test.txt";
const char* content = "test content"; 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; FileMetadata m;
m.mode = 0644; m.mode = 0644;
+1 -1
View File
@@ -84,7 +84,7 @@ static void test_property_chunk_roundtrip() {
for (int i = 0; i < content_len; i++) for (int i = 0; i < content_len; i++)
content[i] = (char)(rand() % 256); 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; struct stat st;
stat(path, &st); stat(path, &st);
+1 -1
View File
@@ -13,7 +13,7 @@
static void test_chunk_deserialize_truncated() { static void test_chunk_deserialize_truncated() {
char* path = "test_rob_trunc.txt"; char* path = "test_rob_trunc.txt";
char* content = "hello"; char* content = "hello";
to_disk(path, content, strlen(content), false, false); file_write_to_disk(path, content, strlen(content), false, false);
struct stat st; struct stat st;
stat(path, &st); stat(path, &st);
+31 -1
View File
@@ -7,7 +7,7 @@
#include <unistd.h> #include <unistd.h>
static void create_test_file(const char* path, const char* content) { 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() { static void test_scanner_single_file() {
@@ -385,6 +385,35 @@ static void test_scanner_no_patterns() {
rmdir(dir); 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() { void test_scanner() {
test_scanner_single_file(); test_scanner_single_file();
test_scanner_multiple_files(); test_scanner_multiple_files();
@@ -399,4 +428,5 @@ void test_scanner() {
test_scanner_size_range(); test_scanner_size_range();
test_scanner_mixed_patterns(); test_scanner_mixed_patterns();
test_scanner_no_patterns(); test_scanner_no_patterns();
test_parallel_scanner_root_chunks_without_workers();
} }