Merge Wave 6: client features (--port, --threads=N, abort, keepalive) and test coverage
CI / lint (push) Successful in 1m30s
CI / sanitizers (undefined) (push) Successful in 1m4s
CI / sanitizers (address) (push) Successful in 1m9s
CI / fuzz-build (push) Successful in 36s
CI / coverage (push) Successful in 51s
CI / valgrind (push) Successful in 3m14s
CI / build-and-test (push) Successful in 5m37s
CI / lint (push) Successful in 1m30s
CI / sanitizers (undefined) (push) Successful in 1m4s
CI / sanitizers (address) (push) Successful in 1m9s
CI / fuzz-build (push) Successful in 36s
CI / coverage (push) Successful in 51s
CI / valgrind (push) Successful in 3m14s
CI / build-and-test (push) Successful in 5m37s
This commit is contained in:
@@ -104,7 +104,7 @@ partial, alternate, and planned behavior.
|
||||
| `-c, --checksum` | Verify content by checksum instead of size+mtime |
|
||||
| `-z, --compress [level]` | Enable streaming zstd compression (level 1–22, default 5) |
|
||||
| `-a, --archive` | rsync archive mode (`-rlptgoD`): links, metadata, devices and specials (not compression/multithreading) |
|
||||
| `-j, --threads` | Multithreading mode |
|
||||
| `-j, --threads[=N]` | Multithreading mode; `N` (1–256) sets the parallel scanner worker count, bare `-j`/`--threads` uses the default |
|
||||
| `-m` | rsync `--prune-empty-dirs` (short form now rsync-parity) |
|
||||
| `--chunk-serialization` | Chunk serialization (batch all files per chunk; long form only) |
|
||||
| `-s` | rsync `--secluded-args` compatibility no-op (remote SSH argv is already injection-safe) |
|
||||
@@ -202,7 +202,7 @@ transfer is never aborted.
|
||||
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`.
|
||||
4. **Config** — runtime parameters (transported over wire, TLS settings excluded). Includes `timeout`, `contimeout`, `quiet`, `backup`, `backup_dir`, `stats`, `max_depth`, `log_file`.
|
||||
5. **Queue** — thread-safe bounded queue with condition variables
|
||||
6. **DirectoryScanner** — recursive BFS traversal with exclude and include pattern support, max-depth enforcement
|
||||
|
||||
@@ -378,7 +378,7 @@ features without changing the meaning of ordinary compatibility options.
|
||||
|
||||
| Option | Purpose |
|
||||
|---|---|
|
||||
| `-j`, `--threads` | Enable the multithreaded scanner/loader/sender pipeline. |
|
||||
| `-j`, `--threads[=N]` | Enable the multithreaded scanner/loader/sender pipeline. `N` (1–256) sets the parallel scanner worker count; bare `-j`/`--threads` uses the default. |
|
||||
| `-z [level]`, `--compress [level]` | Enable streaming zstd compression, levels 1-22. |
|
||||
| `--compress-level <n>` | Set the zstd compression level. |
|
||||
| `--zc <alg>` | Alias for `--compress-choice`. FastSync supports `zstd` and `none`. |
|
||||
@@ -392,7 +392,7 @@ features without changing the meaning of ordinary compatibility options.
|
||||
| `--delta-block <bytes>` | Set the FastSync delta block size (`--block-size` is an alias). |
|
||||
| `--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. |
|
||||
| `--server-port <port>` | Select the TCP server port (`--port <port>` and `--port=<port>` are rsync-friendly aliases). |
|
||||
| `--tls` | Enable TLS for TCP transport. |
|
||||
| `--bwlimit <KB/s>` | Apply token-bucket bandwidth limiting. |
|
||||
| `--progress` | Show transfer progress and throughput. |
|
||||
@@ -478,7 +478,7 @@ link-target transfer remains incomplete. |
|
||||
| `--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. |
|
||||
| `--server-port <port>` | TCP server port. `--port <port>` / `--port=<port>` is an alias. |
|
||||
| `--tls` | Enable TLS. Requires `--cert` and `--key`. |
|
||||
| `--cert <path>` | TLS certificate file. |
|
||||
| `--key <path>` | TLS private key file. |
|
||||
|
||||
+2
-2
@@ -612,7 +612,7 @@ now transmits targets (the prior behavior was broken/partial); its status moved
|
||||
|------|-------------------|-----------------|-------|
|
||||
| `-e`, `--rsh=COMMAND` | Remote shell to use | ✅ Implemented | `-e`/`--rsh` (and `--rsh=COMMAND`) select the remote-shell program used to build the SSH child argv, overriding the default `ssh`. The command is whitespace-split into the leading argv words so rsync's `-e "ssh -p 2222"` works; the standard `-o` family, an optional `-p` port, `user@host` and the quoted remote command (`fastsync-server --stdio`) follow. Stored in the `rsh_command` config field. **Client-only, never crosses the wire** (it is a launch concern, not a handshake property) |
|
||||
| `--rsync-path=PROGRAM` | rsync binary on remote | ✅ Implemented | Alias for `--fastsync-server-path`: both write the `fastsync_server_path` config field used as the remote-side server program (always quoted as one remote-shell word), which CROSSES the wire as before. Kept separate from `--rsh`, which names the local connecting program |
|
||||
| `--port=PORT` | Alternate daemon port | ✅ Implemented | rsync's daemon-port flag maps to the client-side `server_port` config field: a client connects to a TCP/TLS server (incl. `host::module/path` daemon destinations) with `--server-port`, and the `fastsync-server --daemon` listener's port is taken from its config's `port` key (default 873) or overridden by `--dparam port=` / `-p` |
|
||||
| `--port=PORT`, `--port PORT` | Alternate daemon port | ✅ Implemented | rsync's daemon-port flag is an alias for `--server-port`: both spellings (and `--server-port=PORT`) map to the client-side `server_port` config field. The client connects to a TCP/TLS server (incl. `host::module/path` daemon destinations) on that port, and the `fastsync-server --daemon` listener's port is taken from its config's `port` key (default 873) or overridden by `--dparam port=` / `-p` |
|
||||
| `--sockopts=OPTIONS` | Custom TCP options | ✅ Implemented | Comma-separated allowlist of `OPT=VAL` applied via `setsockopt` after `socket()` before `connect()`/`bind()`. Only `TCP_NODELAY`, `SO_KEEPALIVE`, `SO_REUSEADDR` (0/1) and `SO_RCVBUF`/`SO_SNDBUF` (byte count) are accepted; an unknown option name or a bad value is rejected up front, never silently ignored. A value is required for every option (`OPT=VAL`; a bare name is an error). Applied to the outgoing TCP and TLS client socket; absent by default. `SockOptEntry`/`sockopts` config fields. Local socket concern: never crosses the wire |
|
||||
| `--blocking-io` | Use blocking I/O for remote shell | ✅ Implemented | With `--blocking-io` the SSH-transport socketpair socket is left without `SO_RCVTIMEO`/`SO_SNDTIMEO`, so the transfer blocks naturally; by default it gets the same read/write timeout as the TCP transport (see `--timeout`). `blocking_io` config bool. **Client-only, never crosses the wire** |
|
||||
| `--outbuf=N\|L\|B` | Set output buffering | ✅ Implemented | `N` (none/unbuffered) → `_IONBF`, `L` (line) → `_IOLBF`, `B` (block, the default) → `_IOFBF` via `setvbuf` on stdout and stderr. Garbage values are rejected. `outbuf` config field (`OutbufMode`). **Client-only, never crosses the wire** |
|
||||
@@ -893,7 +893,7 @@ Ranked by user demand, implementation complexity, and interoperability impact (_
|
||||
|
||||
| Feature | Description |
|
||||
|---------|-------------|
|
||||
| `-j` / `--threads` | Multithreaded pipeline (scanner/loader/sender) (renamed from `-m` in Phase 7 Wave A; `-m` is now rsync `--prune-empty-dirs`) |
|
||||
| `-j` / `--threads[=N]` | Multithreaded pipeline (scanner/loader/sender); `N` (1–256) sizes the parallel scanner worker pool, bare `-j`/`--threads` uses the built-in default (renamed from `-m` in Phase 7 Wave A; `-m` is now rsync `--prune-empty-dirs`) |
|
||||
| `--chunk-serialization` | Chunk serialization mode (long form only; `-s` is now rsync `--secluded-args`) |
|
||||
| `--sendfile` | Zero-copy sendfile() syscall (TCP only) (long form only; `-f` is now rsync `--filter`) |
|
||||
| `-z [level]` / `--compress` | zstd compression level (1-22) (`-c` is now rsync `--checksum`) |
|
||||
|
||||
+107
-13
@@ -12,6 +12,7 @@
|
||||
#include "identity.h"
|
||||
#include "log.h"
|
||||
#include "protocol.h"
|
||||
#include "scanner.h"
|
||||
#include "stop_condition.h"
|
||||
#include "transport_tcp.h"
|
||||
#include "transport_tls.h"
|
||||
@@ -27,6 +28,41 @@
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
|
||||
/* Async-signal-safe abort flag set by the SIGINT/SIGTERM handler. Exposed via
|
||||
* client_send.h so the send loops can poll it. Defined here (not in
|
||||
* client_send.c) so the unit-test binary, which compiles this file but not
|
||||
* client_send.c, still links the symbol. */
|
||||
volatile sig_atomic_t client_abort_requested = 0;
|
||||
|
||||
/* Only armed while a network transfer is in flight. Outside that window the
|
||||
* handler restores the default disposition and re-raises, so purely local modes
|
||||
* (--list-only/--dry-run/--read-batch/--only-write-batch and the batch-emission
|
||||
* pass) keep terminating on Ctrl-C/SIGTERM instead of silently swallowing it. */
|
||||
volatile sig_atomic_t client_abort_armed = 0;
|
||||
|
||||
void client_set_abort_armed(bool armed) {
|
||||
client_abort_armed = armed ? 1 : 0;
|
||||
}
|
||||
|
||||
bool client_abort_pending(void) {
|
||||
return client_abort_requested != 0;
|
||||
}
|
||||
|
||||
#ifndef FASTSYNC_TEST_BUILD
|
||||
/* Signal handler: perform NO work beyond storing the flag. Logging, protocol
|
||||
* I/O and the STATUS_ABORT frame are all done later on the normal send path,
|
||||
* which is not async-signal-safe. When no transfer is armed, fall back to the
|
||||
* default action so local-only modes remain interruptible. */
|
||||
static void client_signal_handler(int signo) {
|
||||
if (!client_abort_armed) {
|
||||
signal(signo, SIG_DFL);
|
||||
raise(signo);
|
||||
return;
|
||||
}
|
||||
client_abort_requested = 1;
|
||||
}
|
||||
#endif
|
||||
|
||||
#ifndef FASTSYNC_TEST_BUILD
|
||||
/* Parse environment variables for source/destination directories and save-to-disk flag. */
|
||||
static void parse_environment(const char** out_env_source, const char** out_env_dest,
|
||||
@@ -141,6 +177,23 @@ static int set_checksum_seed(Config* config, const char* value) {
|
||||
return 0;
|
||||
}
|
||||
|
||||
/* --threads=N: enable the -m pipeline and size its parallel scanner pool.
|
||||
* Rejects a non-positive, non-numeric or oversized value up front. */
|
||||
static int set_scanner_threads_option(Config* config, const char* value) {
|
||||
int threads;
|
||||
if (!parse_positive_int(value, &threads)) {
|
||||
log_message(LOG_LEVEL_ERROR, "--threads must be a positive integer");
|
||||
return -1;
|
||||
}
|
||||
if (threads > MAX_SCANNER_THREADS) {
|
||||
log_message(LOG_LEVEL_ERROR, "--threads must be between 1 and %d", MAX_SCANNER_THREADS);
|
||||
return -1;
|
||||
}
|
||||
config->use_multithreading = true;
|
||||
config->scanner_threads = threads;
|
||||
return 0;
|
||||
}
|
||||
|
||||
static int set_compression_threads_option(int* dest, const char* value) {
|
||||
if (set_positive_int_option(dest, value, "--compress-threads") != 0)
|
||||
return -1;
|
||||
@@ -1288,10 +1341,21 @@ static bool cli_handle_transfer_flags(CliParseCtx* ctx) {
|
||||
return true;
|
||||
}
|
||||
if (opt_is(arg, "-j", "--threads")) {
|
||||
/* Bare -j/--threads: enable the pipeline with the scanner's built-in
|
||||
worker default (scanner_threads stays 0). */
|
||||
config->use_multithreading = true;
|
||||
log_info_message(LOG_INFO_MISC, "Enabled Multithreading");
|
||||
return true;
|
||||
}
|
||||
if (strncmp(arg, "--threads=", 10) == 0) {
|
||||
if (set_scanner_threads_option(config, arg + 10) != 0) {
|
||||
ctx->exit_code = -1;
|
||||
} else {
|
||||
log_info_message(LOG_INFO_MISC, "Enabled Multithreading with %d scanner threads",
|
||||
config->scanner_threads);
|
||||
}
|
||||
return true;
|
||||
}
|
||||
if (opt_is(arg, "--chunk-serialization", NULL)) {
|
||||
config->use_chunk_serialization = true;
|
||||
log_info_message(LOG_INFO_MISC, "Enabled Chunk Serialization");
|
||||
@@ -1300,29 +1364,48 @@ static bool cli_handle_transfer_flags(CliParseCtx* ctx) {
|
||||
return false;
|
||||
}
|
||||
|
||||
/* Network/IO options: --server-port, --bwlimit, --chunk-size, --log-file and
|
||||
* --stderr. Returns true when the argument was consumed. */
|
||||
/* Parse and validate a TCP server port (--server-port, or its rsync-friendly
|
||||
* alias --port). Returns 0 on success, -1 (with a message) on a malformed or
|
||||
* out-of-range value. */
|
||||
static int set_server_port_option(Config* config, const char* value, const char* option_name) {
|
||||
int port;
|
||||
if (!parse_positive_int(value, &port)) {
|
||||
char* escaped = output_escape(value, false);
|
||||
log_message(LOG_LEVEL_ERROR, "invalid %s value: %s", option_name,
|
||||
escaped ? escaped : "<allocation failed>");
|
||||
free(escaped);
|
||||
return -1;
|
||||
}
|
||||
if (port > 65535) {
|
||||
log_message(LOG_LEVEL_ERROR, "server port must be 1-65535");
|
||||
return -1;
|
||||
}
|
||||
config->server_port = port;
|
||||
return 0;
|
||||
}
|
||||
|
||||
/* Network/IO options: --server-port/--port, --bwlimit, --chunk-size, --log-file
|
||||
* and --stderr. Returns true when the argument was consumed. */
|
||||
static bool cli_handle_io_options(CliParseCtx* ctx) {
|
||||
Config* config = ctx->config;
|
||||
const char* arg = ctx->argv[ctx->i];
|
||||
if (opt_is(arg, "--server-port", NULL)) {
|
||||
if (opt_is(arg, "--server-port", "--port")) {
|
||||
if (ctx->i + 1 >= ctx->argc) {
|
||||
log_message(LOG_LEVEL_ERROR, "missing argument for %s", arg);
|
||||
ctx->exit_code = -1;
|
||||
return true;
|
||||
}
|
||||
if (!parse_positive_int(ctx->argv[++ctx->i], &config->server_port)) {
|
||||
char* escaped = output_escape(ctx->argv[ctx->i], false);
|
||||
log_message(LOG_LEVEL_ERROR, "invalid --server-port value: %s",
|
||||
escaped ? escaped : "<allocation failed>");
|
||||
free(escaped);
|
||||
if (set_server_port_option(config, ctx->argv[++ctx->i], arg) != 0)
|
||||
ctx->exit_code = -1;
|
||||
return true;
|
||||
}
|
||||
if (config->server_port > 65535) {
|
||||
log_message(LOG_LEVEL_ERROR, "server port must be 1-65535");
|
||||
return true;
|
||||
}
|
||||
/* rsync users commonly write --port=NNNN; --server-port=NNNN is accepted too
|
||||
* so both spellings behave identically. */
|
||||
if (strncmp(arg, "--server-port=", 14) == 0 || strncmp(arg, "--port=", 7) == 0) {
|
||||
const char* option_name = strncmp(arg, "--server-port=", 14) == 0 ? "--server-port" : "--port";
|
||||
const char* value = arg + (strncmp(arg, "--server-port=", 14) == 0 ? 14 : 7);
|
||||
if (set_server_port_option(config, value, option_name) != 0)
|
||||
ctx->exit_code = -1;
|
||||
}
|
||||
return true;
|
||||
}
|
||||
if (opt_is(arg, "--bwlimit", NULL)) {
|
||||
@@ -1975,6 +2058,17 @@ int main(int argc, char* argv[]) {
|
||||
oversized delta). Ignore SIGPIPE so that a broken TCP connection
|
||||
surfaces as a clean write error instead of killing the client. */
|
||||
signal(SIGPIPE, SIG_IGN);
|
||||
/* Ctrl-C / SIGTERM: set the abort flag so the send loops can send
|
||||
* STATUS_ABORT and let the receiver clean up, instead of dying abruptly.
|
||||
* No SA_RESTART so an in-flight poll()/read() is interrupted (EINTR), which
|
||||
* lets the keepalive/abort checks observe the flag promptly. */
|
||||
struct sigaction abort_action;
|
||||
memset(&abort_action, 0, sizeof(abort_action));
|
||||
abort_action.sa_handler = client_signal_handler;
|
||||
sigemptyset(&abort_action.sa_mask);
|
||||
abort_action.sa_flags = 0;
|
||||
sigaction(SIGINT, &abort_action, NULL);
|
||||
sigaction(SIGTERM, &abort_action, NULL);
|
||||
const char* env_source = NULL;
|
||||
const char* env_dest = NULL;
|
||||
bool save_to_disk = false;
|
||||
|
||||
@@ -885,6 +885,10 @@ static int send_delete_manifest(int fd, ArrayList* manifest, ArrayList* protecte
|
||||
unlinks) before replying, so the wait uses a generous explicit deadline
|
||||
instead of the default 60 s receive window. */
|
||||
#define DELETE_ACK_TIMEOUT_SEC 3600
|
||||
/* While waiting for the (potentially slow) receiver-side deletion, send a
|
||||
* STATUS_KEEPALIVE at most this often so the connection is demonstrably alive
|
||||
* and neither side's per-message timeout trips. */
|
||||
#define DELETE_ACK_KEEPALIVE_SEC 10
|
||||
|
||||
static bool send_delete_manifest_early(Client* client, ArrayList* manifest,
|
||||
ArrayList* protected_prefixes, ArrayList* missing_args) {
|
||||
@@ -894,8 +898,21 @@ static bool send_delete_manifest_early(Client* client, ArrayList* manifest,
|
||||
0)
|
||||
return false;
|
||||
Status ack;
|
||||
if (!receive_status_timed(client->file_descriptor, &ack, DELETE_ACK_TIMEOUT_SEC))
|
||||
/* The wait is long (up to an hour) and runs inline on this thread: a helper
|
||||
* thread would race the non-thread-safe protocol send path, so keepalives are
|
||||
* emitted from this wait loop itself. A Ctrl-C/SIGTERM abort flag also ends
|
||||
* the wait; the caller then best-effort sends STATUS_ABORT. */
|
||||
if (!receive_status_keepalive(client->file_descriptor, &ack, DELETE_ACK_TIMEOUT_SEC,
|
||||
DELETE_ACK_KEEPALIVE_SEC, client_abort_pending)) {
|
||||
/* A Ctrl-C/SIGTERM abort ends the wait above; tell the receiver before the
|
||||
caller tears the connection down (best-effort). */
|
||||
if (client_abort_pending()) {
|
||||
log_info_message(LOG_INFO_MISC,
|
||||
"Abort requested while awaiting delete ack; sending STATUS_ABORT");
|
||||
send_status(client->file_descriptor, STATUS_ABORT);
|
||||
}
|
||||
return false;
|
||||
}
|
||||
if (ack != STATUS_OK) {
|
||||
log_message(LOG_LEVEL_ERROR, "Server failed to delete files before the transfer");
|
||||
return false;
|
||||
@@ -1490,6 +1507,19 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
||||
}
|
||||
|
||||
while (true) {
|
||||
/* Graceful abort (Ctrl-C/SIGTERM): tell the receiver to clean up instead of
|
||||
dying abruptly. Best-effort: a failed send just means the peer is gone.
|
||||
Only reached while the session is active (config_send already succeeded). */
|
||||
if (client_abort_pending()) {
|
||||
log_info_message(LOG_INFO_MISC,
|
||||
"Abort requested; sending STATUS_ABORT to server and disconnecting");
|
||||
send_status(client->file_descriptor, STATUS_ABORT);
|
||||
pipeline_cancel(context);
|
||||
disconnect_transfer_client(client);
|
||||
mark_sender_done(context);
|
||||
protocol_session_unbind();
|
||||
return thrd_error;
|
||||
}
|
||||
/* Phase 6: stop-elegantly at the next chunk boundary once the --stop-after
|
||||
/ --stop-at deadline has passed. Everything already sent is finalized by
|
||||
the completion tail below; the run still returns success. */
|
||||
@@ -1624,7 +1654,9 @@ static int scan_directory_multithreaded(void* pipeline_context) {
|
||||
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
|
||||
protocol_session_bind(&context->allocation_session);
|
||||
PreparedScanner prepared;
|
||||
if (!prepare_scanner(context->config, 4, &prepared)) {
|
||||
/* -j/--threads=N sizes the parallel scanner's worker pool; 0 (bare -j) lets
|
||||
* the scanner apply its built-in default. */
|
||||
if (!prepare_scanner(context->config, context->config->scanner_threads, &prepared)) {
|
||||
pipeline_cancel(context);
|
||||
protocol_session_unbind();
|
||||
return thrd_error;
|
||||
@@ -1917,6 +1949,9 @@ int send_files(Config* config) {
|
||||
return 1;
|
||||
}
|
||||
|
||||
/* From here on a server session may be live, so Ctrl-C/SIGTERM should set the
|
||||
abort flag (and be forwarded as STATUS_ABORT) instead of terminating. */
|
||||
client_set_abort_armed(true);
|
||||
Client* client = connect_transfer_client(config);
|
||||
if (!client) {
|
||||
if (config->transport == TRANSPORT_TCP)
|
||||
@@ -2036,6 +2071,15 @@ int send_files(Config* config) {
|
||||
only a prefix of the source. */
|
||||
bool scan_stopped_early = false;
|
||||
while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
|
||||
/* Graceful abort (Ctrl-C/SIGTERM): notify the receiver and clean up. The
|
||||
session is active (config_send already succeeded); a send failure here is
|
||||
fine because the client is exiting anyway. */
|
||||
if (client_abort_pending()) {
|
||||
log_info_message(LOG_INFO_MISC, "Abort requested; sending STATUS_ABORT to server");
|
||||
chunk_destroy(current_chunk);
|
||||
send_status(client->file_descriptor, STATUS_ABORT);
|
||||
goto send_fail;
|
||||
}
|
||||
/* Phase 6: stop-elegantly at the next chunk boundary once the deadline has
|
||||
passed. The scanner may also have stopped early itself; either way the
|
||||
completion tail below keeps everything already sent. */
|
||||
@@ -2099,6 +2143,13 @@ int send_files(Config* config) {
|
||||
goto send_fail;
|
||||
if (directory_scanner_had_io_error(scanner))
|
||||
had_scan_io = true;
|
||||
/* An abort that arrived after the last chunk must still stop the completion
|
||||
tail (manifest/finalize) rather than let it run to success. */
|
||||
if (client_abort_pending()) {
|
||||
log_info_message(LOG_INFO_MISC, "Abort requested; sending STATUS_ABORT to server");
|
||||
send_status(client->file_descriptor, STATUS_ABORT);
|
||||
goto send_fail;
|
||||
}
|
||||
/* Phase 6: the scanner may have stopped early (returning NULL without a
|
||||
failure) as soon as the deadline passed, so reflect that here too. A
|
||||
deadline that cut the scan short leaves an incomplete keep-set; transmitting
|
||||
@@ -2180,6 +2231,7 @@ send_fail:
|
||||
prepared_scanner_destroy(&prepared);
|
||||
disconnect_transfer_client(client);
|
||||
protocol_session_unbind();
|
||||
client_set_abort_armed(false);
|
||||
return ret;
|
||||
}
|
||||
|
||||
@@ -2209,6 +2261,9 @@ int send_files_multithreaded(Config** config_ptr) {
|
||||
return 1;
|
||||
}
|
||||
|
||||
/* Armed only once a session may go live (see send_files). */
|
||||
client_set_abort_armed(true);
|
||||
|
||||
long pages = sysconf(_SC_AVPHYS_PAGES);
|
||||
long page_size = sysconf(_SC_PAGE_SIZE);
|
||||
unsigned long long available_memory =
|
||||
@@ -2276,7 +2331,7 @@ int send_files_multithreaded(Config** config_ptr) {
|
||||
fills the protected excluded prefixes. */
|
||||
PreparedScanner prepared;
|
||||
memset(&prepared, 0, sizeof(prepared));
|
||||
bool prepared_ok = prepare_scanner(config, 4, &prepared);
|
||||
bool prepared_ok = prepare_scanner(config, config->scanner_threads, &prepared);
|
||||
if (prepared_ok && context->excluded_paths)
|
||||
prepared.options.excluded_paths = context->excluded_paths;
|
||||
bool prebuilt = prepared_ok && scan_paths_only(config, &prepared.options, context->manifest,
|
||||
@@ -2363,5 +2418,6 @@ int send_files_multithreaded(Config** config_ptr) {
|
||||
/* --ignore-errors: the run completed (and deleted) past an unreadable source
|
||||
directory; report it as errored like rsync does. */
|
||||
pipeline_context_sender_destroy(context);
|
||||
client_set_abort_armed(false);
|
||||
return sender_ok && !scan_io ? 0 : 1;
|
||||
}
|
||||
|
||||
@@ -4,6 +4,19 @@
|
||||
#include "chunk.h"
|
||||
#include "config.h"
|
||||
#include "transport_tcp.h"
|
||||
#include <signal.h>
|
||||
#include <stdbool.h>
|
||||
|
||||
/* Set ONLY by the client's SIGINT/SIGTERM handler (async-signal-safe: the
|
||||
* handler stores 1 and does nothing else). The send loops poll it via
|
||||
* client_abort_pending() and, when set, best-effort send STATUS_ABORT so the
|
||||
* receiver can clean up before the client exits. */
|
||||
extern volatile sig_atomic_t client_abort_requested;
|
||||
bool client_abort_pending(void);
|
||||
/* Arm/disarm abort handling around the network phase. While disarmed, a
|
||||
* SIGINT/SIGTERM takes the default action (immediate termination) so local-only
|
||||
* modes are not left unresponsive. Defined in client_cli.c. */
|
||||
void client_set_abort_armed(bool armed);
|
||||
|
||||
/* Both sender entry points BORROW `config` for the duration of the call; they
|
||||
* never free it, and the caller retains ownership (freeing it with
|
||||
|
||||
@@ -14,6 +14,10 @@
|
||||
#include <sys/types.h>
|
||||
#include <threads.h>
|
||||
|
||||
/* Upper bound on the configurable parallel scanner worker count (--threads=N):
|
||||
* keeps one transfer from spawning an unbounded pool on a very large machine. */
|
||||
#define MAX_SCANNER_THREADS 256
|
||||
|
||||
typedef struct {
|
||||
bool use_metadata;
|
||||
/* Phase 4 metadata capture: -U/--atimes and -N/--crtimes tell the scanner to
|
||||
|
||||
+6
-1
@@ -2,6 +2,7 @@
|
||||
#include <stdio.h>
|
||||
#include <delta.h>
|
||||
#include <chunk.h>
|
||||
#include "scanner.h"
|
||||
|
||||
void print_usage(void) {
|
||||
printf("Usage:\n");
|
||||
@@ -144,7 +145,10 @@ void print_usage(void) {
|
||||
printf(" Delta block size in bytes (default: %d)\n", DELTA_BLOCK_SIZE_DEFAULT);
|
||||
printf(" --delta-max <n> Max file size for delta transfer (default: %llu)\n",
|
||||
DELTA_MAX_FILE_SIZE);
|
||||
printf(" -j, --threads Enable multithreading\n");
|
||||
printf(" -j, --threads[=N] Enable the multithreaded scanner/loader/sender\n");
|
||||
printf(" pipeline; N (1-%d) sets the parallel scanner worker\n",
|
||||
MAX_SCANNER_THREADS);
|
||||
printf(" count (bare -j/--threads uses the default)\n");
|
||||
printf(" --chunk-serialization Enable chunk serialization (long form only)\n");
|
||||
printf(" -s, --secluded-args Protect-args compatibility option (no effect; remote\n");
|
||||
printf(" SSH argv is already built injection-safe)\n");
|
||||
@@ -203,6 +207,7 @@ void print_usage(void) {
|
||||
printf(" --save-to-disk Write received files to disk\n");
|
||||
printf(" --server-host <ip> Server IP address (default: 127.0.0.1)\n");
|
||||
printf(" --server-port <n> Server port (default: 8080)\n");
|
||||
printf(" --port <n> Alias for --server-port\n");
|
||||
printf(" --password-file <f> Authenticate a host::module/path daemon destination.\n");
|
||||
printf(" The file's first user:password line supplies the\n");
|
||||
printf(" username and password (only a SHA-256 digest of the\n");
|
||||
|
||||
+1
-6
@@ -23,6 +23,7 @@ static void config_set_defaults(Config* config) {
|
||||
config->receive_root_directory = NULL;
|
||||
config->save_to_disk = false;
|
||||
config->use_multithreading = false;
|
||||
config->scanner_threads = 0;
|
||||
config->use_chunk_serialization = false;
|
||||
config->use_compression = false;
|
||||
config->use_metadata = false;
|
||||
@@ -79,7 +80,6 @@ static void config_set_defaults(Config* config) {
|
||||
config->stats = false;
|
||||
config->max_depth = 0;
|
||||
config->log_file = NULL;
|
||||
config->queue_size = 100;
|
||||
config->follow_symlinks = false;
|
||||
config->partial = false;
|
||||
config->copy_links = false;
|
||||
@@ -147,14 +147,11 @@ static void config_set_defaults(Config* config) {
|
||||
config->delete_during = false;
|
||||
config->delete_delay = false;
|
||||
config->address = NULL;
|
||||
config->bind_address = NULL;
|
||||
config->ipv6 = false;
|
||||
config->ipv4 = false;
|
||||
config->sockopts = NULL;
|
||||
config->sockopt_count = 0;
|
||||
config->daemon = false;
|
||||
config->daemon_config = NULL;
|
||||
config->server_mode = false;
|
||||
config->no_motd = false;
|
||||
config->checksum = false;
|
||||
config->checksum_algo = CHECKSUM_ALGO_XXH64;
|
||||
@@ -789,9 +786,7 @@ void config_delete(Config* config) {
|
||||
free(config->partial_dir);
|
||||
free(config->suffix);
|
||||
free(config->address);
|
||||
free(config->bind_address);
|
||||
free(config->sockopts);
|
||||
free(config->daemon_config);
|
||||
free(config->compress_choice);
|
||||
free(config->chmod_spec);
|
||||
if (config->skip_compress_suffixes) {
|
||||
|
||||
+6
-6
@@ -81,6 +81,11 @@ typedef struct Config {
|
||||
char* receive_root_directory;
|
||||
bool save_to_disk;
|
||||
bool use_multithreading;
|
||||
/* -j/--threads=N: number of parallel scanner worker threads for the -m
|
||||
* pipeline. 0 (the default, also set by bare -j/--threads) means "use the
|
||||
* scanner's built-in default" (4). CLIENT-ONLY: it is a local scheduling
|
||||
* concern and is NEVER serialized into the wire config frame. */
|
||||
int scanner_threads;
|
||||
bool use_chunk_serialization;
|
||||
bool use_compression;
|
||||
bool use_sendfile;
|
||||
@@ -170,7 +175,6 @@ typedef struct Config {
|
||||
bool stats;
|
||||
int max_depth;
|
||||
FILE* log_file;
|
||||
int queue_size;
|
||||
bool follow_symlinks;
|
||||
bool partial;
|
||||
|
||||
@@ -348,21 +352,17 @@ typedef struct Config {
|
||||
|
||||
// PR #181: IPv6 and bind address
|
||||
char* address;
|
||||
char* bind_address;
|
||||
bool ipv6;
|
||||
bool ipv4;
|
||||
/* --sockopts=OPTIONS (Phase 5, Wave B): strict allowlist of TCP/socket
|
||||
* options applied via setsockopt after socket() and before connect()/bind().
|
||||
* These are LOCAL socket concerns: they never cross the wire config frame.
|
||||
* .address is the outgoing/source bind address (--address); .bind_address is
|
||||
* reserved for daemon-side binding and is not wired yet. */
|
||||
* .address is the outgoing/source bind address (--address). */
|
||||
SockOptEntry* sockopts;
|
||||
int sockopt_count;
|
||||
|
||||
// PR #182: Daemon/server mode
|
||||
bool daemon;
|
||||
char* daemon_config;
|
||||
bool server_mode;
|
||||
/* --no-motd (Wave C): CLIENT-ONLY, never crosses the wire. Suppresses
|
||||
* DISPLAY of the daemon's MOTD; the daemon still sends the MOTD frame, so
|
||||
* the client reads and discards it to keep the stream in sync. rsync's
|
||||
|
||||
@@ -303,6 +303,12 @@ bool protocol_send_n_data(ProtocolSession* session, const void* data, size_t dat
|
||||
wait_events = ssl_err == SSL_ERROR_WANT_WRITE ? POLLOUT : POLLIN;
|
||||
continue;
|
||||
}
|
||||
/* A signal (e.g. Ctrl-C) interrupts the blocking TLS write: retry so
|
||||
the send loop can observe the abort flag at the next checkpoint. */
|
||||
if (ssl_err == SSL_ERROR_SYSCALL && errno == EINTR)
|
||||
continue;
|
||||
} else if (errno == EINTR) {
|
||||
continue;
|
||||
}
|
||||
log_message(LOG_LEVEL_ERROR, "Could not send data");
|
||||
return false;
|
||||
@@ -609,6 +615,151 @@ bool protocol_receive_status_timed(ProtocolSession* session, Status* status, int
|
||||
return true;
|
||||
}
|
||||
|
||||
/* Read exactly one Status frame within `deadline` (CLOCK_MONOTONIC). Unlike
|
||||
* protocol_receive_status_keepalive this never emits a keepalive: it is used
|
||||
* to consume the first byte(s) of an already-signalled frame and to drain the
|
||||
* peer's outstanding keepalive replies, where injecting a write could split a
|
||||
* reply across a frame boundary. Returns false on timeout/EOF/error. */
|
||||
static bool protocol_read_status_until(ProtocolSession* session, Status* status,
|
||||
const struct timespec* deadline) {
|
||||
Status received = STATUS_ERROR;
|
||||
size_t got = 0;
|
||||
short wait_events = POLLIN;
|
||||
while (got < sizeof(Status)) {
|
||||
if (!session->ssl || SSL_pending(session->ssl) == 0) {
|
||||
int remaining_ms = deadline_remaining_ms(deadline);
|
||||
if (remaining_ms <= 0) {
|
||||
log_message(LOG_LEVEL_ERROR, "Receive timeout while reading status");
|
||||
return false;
|
||||
}
|
||||
struct pollfd pfd = {.fd = session->read_fd, .events = wait_events};
|
||||
int poll_result = poll(&pfd, 1, remaining_ms);
|
||||
if (poll_result == 0) {
|
||||
log_message(LOG_LEVEL_ERROR, "Receive timeout while reading status");
|
||||
return false;
|
||||
}
|
||||
if (poll_result < 0) {
|
||||
if (errno == EINTR)
|
||||
continue;
|
||||
return false;
|
||||
}
|
||||
if (pfd.revents & (POLLERR | POLLNVAL))
|
||||
return false;
|
||||
}
|
||||
ssize_t bytes_received;
|
||||
if (session->ssl)
|
||||
bytes_received = SSL_read(session->ssl, (char*)&received + got, sizeof(Status) - got);
|
||||
else
|
||||
bytes_received = read(session->read_fd, (char*)&received + got, sizeof(Status) - got);
|
||||
if (bytes_received <= 0) {
|
||||
if (session->ssl) {
|
||||
int ssl_err = SSL_get_error(session->ssl, (int)bytes_received);
|
||||
if (ssl_err == SSL_ERROR_WANT_READ || ssl_err == SSL_ERROR_WANT_WRITE) {
|
||||
wait_events = ssl_err == SSL_ERROR_WANT_WRITE ? POLLOUT : POLLIN;
|
||||
continue;
|
||||
}
|
||||
}
|
||||
if (bytes_received < 0 && errno == EINTR)
|
||||
continue;
|
||||
log_message(LOG_LEVEL_ERROR, "Connection closed while receiving status");
|
||||
return false;
|
||||
}
|
||||
got += (size_t)bytes_received;
|
||||
}
|
||||
*status = received;
|
||||
return true;
|
||||
}
|
||||
|
||||
bool protocol_receive_status_keepalive(ProtocolSession* session, Status* status, int timeout_sec,
|
||||
int keepalive_interval_sec, ProtocolWaitAbort abort_check) {
|
||||
if (!session || !status)
|
||||
return false;
|
||||
if (timeout_sec <= 0)
|
||||
timeout_sec = RECEIVE_TIMEOUT_SEC;
|
||||
if (keepalive_interval_sec <= 0)
|
||||
keepalive_interval_sec = timeout_sec;
|
||||
|
||||
struct timespec deadline;
|
||||
clock_gettime(CLOCK_MONOTONIC, &deadline);
|
||||
deadline.tv_sec += timeout_sec;
|
||||
|
||||
unsigned long keepalives_sent = 0;
|
||||
unsigned long replies_seen = 0;
|
||||
Status final = STATUS_ERROR;
|
||||
while (true) {
|
||||
if (abort_check && abort_check())
|
||||
return false;
|
||||
if (!session->ssl || SSL_pending(session->ssl) == 0) {
|
||||
int remaining_ms = deadline_remaining_ms(&deadline);
|
||||
if (remaining_ms <= 0) {
|
||||
log_message(LOG_LEVEL_ERROR, "Receive timeout after %ds", timeout_sec);
|
||||
return false;
|
||||
}
|
||||
/* Only interleave a keepalive while waiting for the FIRST byte of a
|
||||
* frame; once part of a frame is buffered a write could race the peer's
|
||||
* reply into the middle of it. */
|
||||
long long interval_ms_ll = (long long)keepalive_interval_sec * 1000LL;
|
||||
int interval_ms = interval_ms_ll > INT_MAX ? INT_MAX : (int)interval_ms_ll;
|
||||
int wait_ms = interval_ms < remaining_ms ? interval_ms : remaining_ms;
|
||||
struct pollfd pfd = {.fd = session->read_fd, .events = POLLIN};
|
||||
int poll_result = poll(&pfd, 1, wait_ms);
|
||||
if (poll_result == 0) {
|
||||
if (abort_check && abort_check())
|
||||
return false;
|
||||
if (!protocol_send_status(session, STATUS_KEEPALIVE))
|
||||
return false;
|
||||
keepalives_sent++;
|
||||
continue;
|
||||
}
|
||||
if (poll_result < 0) {
|
||||
if (errno == EINTR)
|
||||
continue;
|
||||
return false;
|
||||
}
|
||||
if (pfd.revents & (POLLERR | POLLNVAL))
|
||||
return false;
|
||||
}
|
||||
Status received;
|
||||
if (!protocol_read_status_until(session, &received, &deadline))
|
||||
return false;
|
||||
if (received == STATUS_KEEPALIVE) {
|
||||
/* The receiver's answer to one of our keepalives. */
|
||||
replies_seen++;
|
||||
continue;
|
||||
}
|
||||
final = received;
|
||||
break;
|
||||
}
|
||||
/* Drain the replies the receiver still owes for keepalives we sent while it
|
||||
* was busy. It answers them only after the real status, so leaving them
|
||||
* unread would put stale KEEPALIVE frames ahead of the next exchange and
|
||||
* desynchronize the protocol. */
|
||||
if (replies_seen < keepalives_sent) {
|
||||
/* A short separate grace, not the (possibly exhausted) main deadline: the
|
||||
terminal status already arrived, so a peer that never answers its owed
|
||||
keepalives must not turn a successful ack into a reported failure. */
|
||||
struct timespec drain_deadline;
|
||||
clock_gettime(CLOCK_MONOTONIC, &drain_deadline);
|
||||
drain_deadline.tv_sec += 1;
|
||||
while (replies_seen < keepalives_sent) {
|
||||
Status drained;
|
||||
if (!protocol_read_status_until(session, &drained, &drain_deadline)) {
|
||||
log_message(LOG_LEVEL_WARNING, "peer did not answer %lu keepalive(s); continuing",
|
||||
keepalives_sent - replies_seen);
|
||||
break;
|
||||
}
|
||||
if (drained != STATUS_KEEPALIVE) {
|
||||
log_message(LOG_LEVEL_ERROR, "Unexpected status while draining keepalive replies");
|
||||
return false;
|
||||
}
|
||||
replies_seen++;
|
||||
}
|
||||
}
|
||||
*status = final;
|
||||
log_debug_message(LOG_DEBUG_PROTO, "Received Status: %s", status_to_string(*status));
|
||||
return true;
|
||||
}
|
||||
|
||||
bool send_str(int fd, const char* data) {
|
||||
return protocol_send_str(legacy_session(-1, fd), data);
|
||||
}
|
||||
@@ -648,3 +799,8 @@ bool receive_status(int fd, Status* status) {
|
||||
bool receive_status_timed(int fd, Status* status, int timeout_sec) {
|
||||
return protocol_receive_status_timed(legacy_session(fd, -1), status, timeout_sec);
|
||||
}
|
||||
bool receive_status_keepalive(int fd, Status* status, int timeout_sec, int keepalive_interval_sec,
|
||||
ProtocolWaitAbort abort_check) {
|
||||
return protocol_receive_status_keepalive(legacy_session(fd, -1), status, timeout_sec,
|
||||
keepalive_interval_sec, abort_check);
|
||||
}
|
||||
|
||||
@@ -198,4 +198,23 @@ bool receive_status(int file_descriptor, Status* status);
|
||||
this so the sender does not abort after the deletion already committed. */
|
||||
bool receive_status_timed(int file_descriptor, Status* status, int timeout_sec);
|
||||
|
||||
/* Callback polled by protocol_receive_status_keepalive once per keepalive
|
||||
interval. Return true to stop waiting (e.g. a SIGINT/SIGTERM abort flag was
|
||||
set). Kept as a function pointer so the protocol layer does not depend on
|
||||
client signal state. */
|
||||
typedef bool (*ProtocolWaitAbort)(void);
|
||||
|
||||
/* Like receive_status_timed, but while the peer is silent it emits
|
||||
STATUS_KEEPALIVE every keepalive_interval_sec (the receiver answers each with
|
||||
STATUS_KEEPALIVE, which this function consumes and skips) so a long
|
||||
server-side operation does not look like a dead connection. The total wait
|
||||
is still bounded by timeout_sec; abort_check (may be NULL) is polled every
|
||||
interval and, when it returns true, ends the wait immediately with false.
|
||||
Runs entirely on the calling thread: the protocol send path is NOT safe for
|
||||
concurrent writers, so this must not be paired with a helper thread. */
|
||||
bool receive_status_keepalive(int file_descriptor, Status* status, int timeout_sec,
|
||||
int keepalive_interval_sec, ProtocolWaitAbort abort_check);
|
||||
bool protocol_receive_status_keepalive(ProtocolSession* session, Status* status, int timeout_sec,
|
||||
int keepalive_interval_sec, ProtocolWaitAbort abort_check);
|
||||
|
||||
#endif
|
||||
|
||||
@@ -0,0 +1,126 @@
|
||||
/*
|
||||
* Fuzz the delete-manifest parser: receive_manifest_entries(int fd).
|
||||
*
|
||||
* The parser reads three length-delimited sections (keeps, protected prefixes,
|
||||
* missing-args paths) from the connection. Feeding raw bytes alone exercises
|
||||
* the "reject the first malformed count/string" fast paths, but because each
|
||||
* section is self-delimiting a single bad value hides every later section.
|
||||
*
|
||||
* To reach the protected-prefix and missing-args parsers (the paths that drive
|
||||
* actual destination deletion) we build one canonical, fully-valid manifest
|
||||
* with hand-written wire framing and then feed the receiver several shapes:
|
||||
*
|
||||
* 1. raw : the raw fuzz bytes as the whole manifest.
|
||||
* 2. keeps : the valid keep count only + the fuzz bytes, so the fuzzer
|
||||
* drives the keep count and entries directly.
|
||||
* 3. prot : the valid keeps section + the fuzz bytes, so the fuzzer drives
|
||||
* the protected count and prefixes.
|
||||
* 4. missing: the valid keeps+protected sections + the fuzz bytes, so the
|
||||
* fuzzer drives the trailing missing-args section, including the
|
||||
* aggregate MAX_MANIFEST_BYTES budget.
|
||||
*
|
||||
* The wire encoding matches receive_int (native int) and receive_wire_str
|
||||
* (native size_t length prefix + body); no charset conversion is configured in
|
||||
* the fuzz process, so receive_wire_str is receive_str.
|
||||
*/
|
||||
#include "file_receive.h"
|
||||
#include "protocol.h"
|
||||
#include <errno.h>
|
||||
#include <fcntl.h>
|
||||
#include <stdint.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/socket.h>
|
||||
#include <unistd.h>
|
||||
|
||||
static unsigned char g_manifest[512];
|
||||
static size_t g_len_after_count; /* offset of the first keep entry */
|
||||
static size_t g_len_after_keeps; /* offset of the protected count */
|
||||
static size_t g_len_after_protected; /* offset of the missing count */
|
||||
static int g_manifest_ready;
|
||||
|
||||
static void append_int32(unsigned char* buf, size_t* off, int32_t value) {
|
||||
memcpy(buf + *off, &value, sizeof(value));
|
||||
*off += sizeof(value);
|
||||
}
|
||||
|
||||
static void append_wire_str(unsigned char* buf, size_t* off, const char* s) {
|
||||
size_t n = strlen(s);
|
||||
memcpy(buf + *off, &n, sizeof(n));
|
||||
*off += sizeof(n);
|
||||
memcpy(buf + *off, s, n);
|
||||
*off += n;
|
||||
}
|
||||
|
||||
static void build_canonical_manifest(void) {
|
||||
g_manifest_ready = 1;
|
||||
size_t off = 0;
|
||||
append_int32(g_manifest, &off, 2);
|
||||
g_len_after_count = off;
|
||||
append_wire_str(g_manifest, &off, "keep/a");
|
||||
append_wire_str(g_manifest, &off, "keep/b");
|
||||
g_len_after_keeps = off;
|
||||
append_int32(g_manifest, &off, 1);
|
||||
append_wire_str(g_manifest, &off, "excluded/prefix");
|
||||
g_len_after_protected = off;
|
||||
append_int32(g_manifest, &off, 1);
|
||||
append_wire_str(g_manifest, &off, "missing/path");
|
||||
}
|
||||
|
||||
/* Best-effort non-blocking write: an oversized fuzz input is truncated rather
|
||||
* than stalling the harness. */
|
||||
static void write_best_effort(int fd, const void* data, size_t size) {
|
||||
const unsigned char* p = data;
|
||||
size_t off = 0;
|
||||
while (off < size) {
|
||||
ssize_t n = write(fd, p + off, size - off);
|
||||
if (n > 0) {
|
||||
off += (size_t)n;
|
||||
continue;
|
||||
}
|
||||
if (n < 0 && errno == EINTR)
|
||||
continue;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
/* Build prefix ++ data as a stream and drive receive_manifest_entries over it.
|
||||
* The write half is shut down first so the parser always sees EOF instead of
|
||||
* blocking on a missing frame tail. */
|
||||
static void receive_stream(const unsigned char* prefix, size_t prefix_len, const uint8_t* data,
|
||||
size_t size) {
|
||||
int sv[2];
|
||||
if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) != 0)
|
||||
return;
|
||||
|
||||
int flags = fcntl(sv[0], F_GETFL, 0);
|
||||
if (flags != -1)
|
||||
(void)fcntl(sv[0], F_SETFL, flags | O_NONBLOCK);
|
||||
|
||||
if (prefix_len > 0)
|
||||
write_best_effort(sv[0], prefix, prefix_len);
|
||||
if (size > 0)
|
||||
write_best_effort(sv[0], data, size);
|
||||
shutdown(sv[0], SHUT_WR);
|
||||
|
||||
DeleteManifest* manifest = receive_manifest_entries(sv[1]);
|
||||
delete_manifest_free(manifest);
|
||||
|
||||
close(sv[0]);
|
||||
close(sv[1]);
|
||||
}
|
||||
|
||||
int LLVMFuzzerTestOneInput(const uint8_t* data, size_t size) {
|
||||
if (!g_manifest_ready)
|
||||
build_canonical_manifest();
|
||||
|
||||
/* Raw bytes as the whole manifest. */
|
||||
receive_stream(NULL, 0, data, size);
|
||||
|
||||
/* Keep the valid framing so the fuzzer reaches each later section. */
|
||||
receive_stream(g_manifest, g_len_after_protected, data, size);
|
||||
receive_stream(g_manifest, g_len_after_keeps, data, size);
|
||||
receive_stream(g_manifest, g_len_after_count, data, size);
|
||||
|
||||
return 0;
|
||||
}
|
||||
@@ -0,0 +1,130 @@
|
||||
/*
|
||||
* Fuzz the base protocol framing: receive_str / receive_data / receive_status
|
||||
* (plus the redacted string, the size-limited data and the timed-status
|
||||
* variants) fed arbitrary bytes over an in-memory socketpair.
|
||||
*
|
||||
* Every receive primitive reads a fixed-width header (a size_t string length,
|
||||
* an unsigned long long data length, an int status/int value) and then a body.
|
||||
* The fuzzer attacks:
|
||||
* - oversized length headers (the MAX_STRING_SIZE / MAX_DATA_PAYLOAD_SIZE
|
||||
* gates must reject before allocating),
|
||||
* - truncated bodies (a declared body larger than the stream must fail
|
||||
* cleanly at EOF, never read uninitialised memory or leak),
|
||||
* - embedded NUL bytes in strings (must be refused),
|
||||
* - out-of-range status enum values (status_to_string must stay in bounds).
|
||||
*
|
||||
* Each entry point gets its own socketpair because a single receive consumes a
|
||||
* variable number of bytes from the stream; reusing one would make the later
|
||||
* calls meaningless. The write half is shut down first so a truncated frame
|
||||
* always terminates at EOF instead of blocking.
|
||||
*/
|
||||
#include "data.h"
|
||||
#include "protocol.h"
|
||||
#include <errno.h>
|
||||
#include <fcntl.h>
|
||||
#include <stdint.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/socket.h>
|
||||
#include <unistd.h>
|
||||
|
||||
static void write_best_effort(int fd, const void* data, size_t size) {
|
||||
const unsigned char* p = data;
|
||||
size_t off = 0;
|
||||
while (off < size) {
|
||||
ssize_t n = write(fd, p + off, size - off);
|
||||
if (n > 0) {
|
||||
off += (size_t)n;
|
||||
continue;
|
||||
}
|
||||
if (n < 0 && errno == EINTR)
|
||||
continue;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
/* Create a socketpair pre-loaded with `data`, shut down the write half and
|
||||
* return the read end (which the receiver reads from). `*write_end` is also
|
||||
* returned so the caller can close it. */
|
||||
static int make_stream(const uint8_t* data, size_t size, int* write_end) {
|
||||
int sv[2];
|
||||
if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) != 0) {
|
||||
*write_end = -1;
|
||||
return -1;
|
||||
}
|
||||
int flags = fcntl(sv[0], F_GETFL, 0);
|
||||
if (flags != -1)
|
||||
(void)fcntl(sv[0], F_SETFL, flags | O_NONBLOCK);
|
||||
if (size > 0)
|
||||
write_best_effort(sv[0], data, size);
|
||||
shutdown(sv[0], SHUT_WR);
|
||||
*write_end = sv[0];
|
||||
return sv[1];
|
||||
}
|
||||
|
||||
int LLVMFuzzerTestOneInput(const uint8_t* data, size_t size) {
|
||||
int w;
|
||||
|
||||
int rd = make_stream(data, size, &w);
|
||||
if (rd >= 0) {
|
||||
char* s = receive_str(rd);
|
||||
free(s);
|
||||
close(rd);
|
||||
close(w);
|
||||
}
|
||||
|
||||
rd = make_stream(data, size, &w);
|
||||
if (rd >= 0) {
|
||||
char* s = receive_str_redacted(rd);
|
||||
free(s);
|
||||
close(rd);
|
||||
close(w);
|
||||
}
|
||||
|
||||
/* Bounded so a crafted 256 MiB length header cannot make each iteration
|
||||
allocate the full MAX_DATA_PAYLOAD_SIZE under ASan; the framing logic is
|
||||
identical to receive_data(), which delegates to the limited variant. */
|
||||
rd = make_stream(data, size, &w);
|
||||
if (rd >= 0) {
|
||||
Data* d = receive_data_limited(rd, 1u << 20);
|
||||
data_destroy(d);
|
||||
close(rd);
|
||||
close(w);
|
||||
}
|
||||
|
||||
/* The size-limited variant must reject anything beyond its explicit bound
|
||||
* before allocating the body buffer. */
|
||||
rd = make_stream(data, size, &w);
|
||||
if (rd >= 0) {
|
||||
Data* d = receive_data_limited(rd, 256);
|
||||
data_destroy(d);
|
||||
close(rd);
|
||||
close(w);
|
||||
}
|
||||
|
||||
rd = make_stream(data, size, &w);
|
||||
if (rd >= 0) {
|
||||
Status status = STATUS_OK;
|
||||
(void)receive_status(rd, &status);
|
||||
close(rd);
|
||||
close(w);
|
||||
}
|
||||
|
||||
rd = make_stream(data, size, &w);
|
||||
if (rd >= 0) {
|
||||
Status status = STATUS_OK;
|
||||
(void)receive_status_timed(rd, &status, 1);
|
||||
close(rd);
|
||||
close(w);
|
||||
}
|
||||
|
||||
rd = make_stream(data, size, &w);
|
||||
if (rd >= 0) {
|
||||
int value = 0;
|
||||
(void)receive_int(rd, &value);
|
||||
close(rd);
|
||||
close(w);
|
||||
}
|
||||
|
||||
return 0;
|
||||
}
|
||||
@@ -0,0 +1,115 @@
|
||||
/*
|
||||
* Fuzz the xattr wire block parser: xattr_receive(int fd, int* ok).
|
||||
*
|
||||
* The block is a count followed by that many (name_len, name, value_len, value)
|
||||
* records. The receiver must reject an invalid count, an out-of-range or
|
||||
* negative name/value length, an embedded NUL or non-whitelisted namespace in
|
||||
* the name, an oversized value, and an aggregate payload beyond
|
||||
* XATTR_TOTAL_MAX -- all without over-allocating or leaking.
|
||||
*
|
||||
* Raw bytes mostly stop at the first invalid count/length, so we also build a
|
||||
* canonical, fully-valid two-entry block by hand and feed the receiver valid
|
||||
* prefixes of it followed by the fuzz bytes. That drives the deep value-
|
||||
* parsing and per-entry namespace/budget checks with attacker-controlled input.
|
||||
*/
|
||||
#include "protocol.h"
|
||||
#include "xattr.h"
|
||||
#include <errno.h>
|
||||
#include <fcntl.h>
|
||||
#include <stdint.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/socket.h>
|
||||
#include <unistd.h>
|
||||
|
||||
static unsigned char g_block[512];
|
||||
static size_t g_off_after_count; /* start of entry 0 */
|
||||
static size_t g_off_after_entry0; /* start of entry 1 */
|
||||
static size_t g_off_value0; /* start of the first value length */
|
||||
static int g_block_ready;
|
||||
|
||||
static void append_int32(unsigned char* buf, size_t* off, int32_t value) {
|
||||
memcpy(buf + *off, &value, sizeof(value));
|
||||
*off += sizeof(value);
|
||||
}
|
||||
|
||||
static void append_bytes(unsigned char* buf, size_t* off, const void* p, size_t n) {
|
||||
if (n > 0)
|
||||
memcpy(buf + *off, p, n);
|
||||
*off += n;
|
||||
}
|
||||
|
||||
static void build_canonical_block(void) {
|
||||
g_block_ready = 1;
|
||||
size_t off = 0;
|
||||
append_int32(g_block, &off, 2);
|
||||
g_off_after_count = off;
|
||||
|
||||
int32_t name0_len = (int32_t)strlen("user.foo");
|
||||
append_int32(g_block, &off, name0_len);
|
||||
append_bytes(g_block, &off, "user.foo", (size_t)name0_len);
|
||||
g_off_value0 = off;
|
||||
append_int32(g_block, &off, 3);
|
||||
append_bytes(g_block, &off, "bar", 3);
|
||||
|
||||
g_off_after_entry0 = off;
|
||||
int32_t name1_len = (int32_t)strlen("user.empty");
|
||||
append_int32(g_block, &off, name1_len);
|
||||
append_bytes(g_block, &off, "user.empty", (size_t)name1_len);
|
||||
append_int32(g_block, &off, 0);
|
||||
}
|
||||
|
||||
static void write_best_effort(int fd, const void* data, size_t size) {
|
||||
const unsigned char* p = data;
|
||||
size_t off = 0;
|
||||
while (off < size) {
|
||||
ssize_t n = write(fd, p + off, size - off);
|
||||
if (n > 0) {
|
||||
off += (size_t)n;
|
||||
continue;
|
||||
}
|
||||
if (n < 0 && errno == EINTR)
|
||||
continue;
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
static void receive_stream(const unsigned char* prefix, size_t prefix_len, const uint8_t* data,
|
||||
size_t size) {
|
||||
int sv[2];
|
||||
if (socketpair(AF_UNIX, SOCK_STREAM, 0, sv) != 0)
|
||||
return;
|
||||
|
||||
int flags = fcntl(sv[0], F_GETFL, 0);
|
||||
if (flags != -1)
|
||||
(void)fcntl(sv[0], F_SETFL, flags | O_NONBLOCK);
|
||||
|
||||
if (prefix_len > 0)
|
||||
write_best_effort(sv[0], prefix, prefix_len);
|
||||
if (size > 0)
|
||||
write_best_effort(sv[0], data, size);
|
||||
shutdown(sv[0], SHUT_WR);
|
||||
|
||||
int ok = 0;
|
||||
FileXattrList* list = xattr_receive(sv[1], &ok);
|
||||
xattr_list_free(list);
|
||||
|
||||
close(sv[0]);
|
||||
close(sv[1]);
|
||||
}
|
||||
|
||||
int LLVMFuzzerTestOneInput(const uint8_t* data, size_t size) {
|
||||
if (!g_block_ready)
|
||||
build_canonical_block();
|
||||
|
||||
/* Raw bytes as the whole block. */
|
||||
receive_stream(NULL, 0, data, size);
|
||||
|
||||
/* Valid framing so the fuzzer mutates the entry list, the first value and
|
||||
* the second entry respectively instead of stopping at the count. */
|
||||
receive_stream(g_block, g_off_after_entry0, data, size);
|
||||
receive_stream(g_block, g_off_value0, data, size);
|
||||
receive_stream(g_block, g_off_after_count, data, size);
|
||||
|
||||
return 0;
|
||||
}
|
||||
@@ -0,0 +1,367 @@
|
||||
"""Fault injection: the server must survive truncated / corrupted protocol
|
||||
frames and abrupt mid-frame disconnects, and keep serving later connections.
|
||||
|
||||
These tests deliberately speak raw bytes to a real server process:
|
||||
|
||||
* malformed frames before/inside the config handshake (oversized length
|
||||
headers, truncated string bodies, outright garbage),
|
||||
* a captured *valid* config frame replayed so the connection reaches the
|
||||
operation loop, followed by a partial ``STATUS_MANIFEST`` frame that is cut
|
||||
mid-body and dropped, and
|
||||
* a real client run relayed through a proxy that truncates the stream at a
|
||||
range of byte offsets and resets both ends.
|
||||
|
||||
After every fault the server process is asserted alive and a subsequent
|
||||
ordinary transfer must complete and verify, proving the accept loop and
|
||||
per-connection children recovered cleanly. All interactions are bounded by
|
||||
short socket timeouts (no sleeps).
|
||||
"""
|
||||
import os
|
||||
import select
|
||||
import shutil
|
||||
import socket
|
||||
import struct
|
||||
import subprocess
|
||||
import sys
|
||||
import threading
|
||||
|
||||
import pytest
|
||||
|
||||
sys.path.insert(0, os.path.dirname(__file__))
|
||||
from common import ( # noqa: E402
|
||||
ServerManager,
|
||||
TEST_DATA_DIR,
|
||||
get_dest_received_dir,
|
||||
run_client,
|
||||
verify_transfer,
|
||||
)
|
||||
|
||||
PROTOCOL_VERSION = b"2.20.0"
|
||||
STATUS_MANIFEST = 5
|
||||
STATUS_OK = 0
|
||||
|
||||
SOURCE_DIR = os.path.join(TEST_DATA_DIR, "fault_src")
|
||||
DEST_DIR = os.path.join(TEST_DATA_DIR, "fault_dst")
|
||||
|
||||
|
||||
@pytest.fixture(scope="module")
|
||||
def fault_server():
|
||||
"""A dedicated server so the aliveness assertions observe exactly the
|
||||
process these faults were sent to."""
|
||||
server = ServerManager()
|
||||
server.start()
|
||||
yield server
|
||||
server.stop()
|
||||
|
||||
|
||||
@pytest.fixture(scope="module", autouse=True)
|
||||
def _seed_source():
|
||||
if os.path.exists(SOURCE_DIR):
|
||||
shutil.rmtree(SOURCE_DIR)
|
||||
os.makedirs(os.path.join(SOURCE_DIR, "nested"))
|
||||
with open(os.path.join(SOURCE_DIR, "hello.txt"), "wb") as fh:
|
||||
fh.write(b"fault injection payload\n" * 64)
|
||||
with open(os.path.join(SOURCE_DIR, "nested", "deep.bin"), "wb") as fh:
|
||||
fh.write(bytes(range(256)) * 16)
|
||||
yield
|
||||
shutil.rmtree(SOURCE_DIR, ignore_errors=True)
|
||||
shutil.rmtree(DEST_DIR, ignore_errors=True)
|
||||
|
||||
|
||||
def _assert_alive(server):
|
||||
assert server._proc is not None, "server process missing"
|
||||
assert server._proc.poll() is None, (
|
||||
f"server exited with {server._proc.returncode} after fault injection"
|
||||
)
|
||||
|
||||
|
||||
def _recover(server, label):
|
||||
"""Run one ordinary transfer and verify it end-to-end."""
|
||||
shutil.rmtree(DEST_DIR, ignore_errors=True)
|
||||
os.makedirs(DEST_DIR)
|
||||
result, _ = run_client(SOURCE_DIR, DEST_DIR, flags=["--preserve"], port=server.port)
|
||||
assert result.returncode == 0, (
|
||||
f"{label}: recovery transfer failed rc={result.returncode}: "
|
||||
f"{(result.stderr or result.stdout)[:200]}"
|
||||
)
|
||||
received = get_dest_received_dir(DEST_DIR, SOURCE_DIR)
|
||||
mismatches, missing = verify_transfer(SOURCE_DIR, received)
|
||||
assert not missing, f"{label}: recovery missing {missing}"
|
||||
assert not mismatches, f"{label}: recovery mismatch {mismatches}"
|
||||
|
||||
|
||||
def _abrupt_close(sock):
|
||||
"""Force an RST instead of a graceful FIN, the nastier mid-frame drop."""
|
||||
try:
|
||||
sock.setsockopt(socket.SOL_SOCKET, socket.SO_LINGER, struct.pack("ii", 1, 0))
|
||||
except OSError:
|
||||
pass
|
||||
try:
|
||||
sock.close()
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
def _raw_connect(server):
|
||||
sock = socket.create_connection(("127.0.0.1", server.port), timeout=5)
|
||||
sock.settimeout(5)
|
||||
return sock
|
||||
|
||||
|
||||
def _recv_exact(sock, n):
|
||||
buf = b""
|
||||
while len(buf) < n:
|
||||
chunk = sock.recv(n - len(buf))
|
||||
if not chunk:
|
||||
return None
|
||||
buf += chunk
|
||||
return buf
|
||||
|
||||
|
||||
# --- faults before/inside the config handshake -----------------------------
|
||||
|
||||
CONFIG_HANDSHAKE_FAULTS = {
|
||||
"empty": b"",
|
||||
# Length header claims a 1 EiB string body that never arrives.
|
||||
"oversized_length": struct.pack("<Q", 1 << 60),
|
||||
# A truncated 8-byte length header (only 3 bytes of it are sent).
|
||||
"truncated_length_header": b"\x10\x00\x00",
|
||||
# A valid version string followed by a string-length header whose body is
|
||||
# deliberately truncated (mid-config-frame disconnect).
|
||||
"truncated_config_body": struct.pack("<Q", len(PROTOCOL_VERSION)) + PROTOCOL_VERSION
|
||||
+ struct.pack("<Q", 4096)
|
||||
+ b"partial",
|
||||
# Pure garbage that is not a valid frame at any offset.
|
||||
"garbage": b"\xff" * 32,
|
||||
}
|
||||
|
||||
|
||||
class TestConfigHandshakeFaults:
|
||||
def test_truncated_and_corrupt_config_frames(self, fault_server):
|
||||
for name, payload in CONFIG_HANDSHAKE_FAULTS.items():
|
||||
sock = _raw_connect(fault_server)
|
||||
if payload:
|
||||
sock.sendall(payload)
|
||||
_abrupt_close(sock)
|
||||
_assert_alive(fault_server)
|
||||
_recover(fault_server, "config handshake faults")
|
||||
|
||||
|
||||
# --- capture a valid config frame, then truncate a STATUS_MANIFEST ----------
|
||||
|
||||
|
||||
class _CaptureProxy:
|
||||
"""Relay one client<->server connection and record the client's config
|
||||
frame (all client bytes forwarded before the server's first reply)."""
|
||||
|
||||
def __init__(self, target_port):
|
||||
self.target = ("127.0.0.1", target_port)
|
||||
self.listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
self.listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
self.listener.bind(("127.0.0.1", 0))
|
||||
self.listener.listen(1)
|
||||
self.listener.settimeout(20)
|
||||
self.port = self.listener.getsockname()[1]
|
||||
self.config_frame = None
|
||||
|
||||
def run(self, cmd):
|
||||
def serve():
|
||||
try:
|
||||
client, _ = self.listener.accept()
|
||||
except OSError:
|
||||
return
|
||||
try:
|
||||
backend = socket.create_connection(self.target, timeout=10)
|
||||
except OSError:
|
||||
client.close()
|
||||
return
|
||||
client.settimeout(20)
|
||||
backend.settimeout(20)
|
||||
buf_c = bytearray()
|
||||
seen_server = False
|
||||
try:
|
||||
while True:
|
||||
ready, _, _ = select.select([client, backend], [], [], 20)
|
||||
if not ready:
|
||||
break
|
||||
done = False
|
||||
for sock in ready:
|
||||
data = sock.recv(65536)
|
||||
if not data:
|
||||
done = True
|
||||
continue
|
||||
if sock is client:
|
||||
buf_c += data
|
||||
backend.sendall(data)
|
||||
else:
|
||||
if not seen_server:
|
||||
seen_server = True
|
||||
self.config_frame = bytes(buf_c)
|
||||
client.sendall(data)
|
||||
if done:
|
||||
break
|
||||
except OSError:
|
||||
pass
|
||||
finally:
|
||||
client.close()
|
||||
backend.close()
|
||||
|
||||
thread = threading.Thread(target=serve)
|
||||
thread.start()
|
||||
result = subprocess.run(cmd, capture_output=True, text=True, timeout=60)
|
||||
thread.join(20)
|
||||
return result
|
||||
|
||||
def close(self):
|
||||
try:
|
||||
self.listener.close()
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
@pytest.fixture(scope="module")
|
||||
def captured_config(fault_server):
|
||||
"""Capture the config frame of one real client run through a relay."""
|
||||
proxy = _CaptureProxy(fault_server.port)
|
||||
cmd = [
|
||||
os.path.join(os.path.dirname(__file__), "..", "..", "build", "client"),
|
||||
"--source-dir",
|
||||
SOURCE_DIR,
|
||||
"--dest-dir",
|
||||
DEST_DIR,
|
||||
"--save-to-disk",
|
||||
"--server-port",
|
||||
str(proxy.port),
|
||||
]
|
||||
try:
|
||||
result = proxy.run(cmd)
|
||||
assert result.returncode == 0, (
|
||||
f"capture run failed rc={result.returncode}: "
|
||||
f"{(result.stderr or result.stdout)[:200]}"
|
||||
)
|
||||
assert proxy.config_frame, "failed to capture the client config frame"
|
||||
yield proxy.config_frame
|
||||
finally:
|
||||
proxy.close()
|
||||
|
||||
|
||||
class TestTruncatedStatusFrame:
|
||||
def test_partial_manifest_frame_then_drop(self, fault_server, captured_config):
|
||||
sock = _raw_connect(fault_server)
|
||||
sock.sendall(captured_config)
|
||||
ack = _recv_exact(sock, 4)
|
||||
assert ack is not None, "server closed before the config ack"
|
||||
(status,) = struct.unpack("<i", ack)
|
||||
assert status == STATUS_OK, f"expected STATUS_OK, got {status}"
|
||||
|
||||
# STATUS_MANIFEST, then only half of the keep-count int, then an RST.
|
||||
sock.sendall(struct.pack("<i", STATUS_MANIFEST) + b"\x02\x00")
|
||||
_abrupt_close(sock)
|
||||
|
||||
_assert_alive(fault_server)
|
||||
_recover(fault_server, "truncated manifest frame")
|
||||
|
||||
def test_manifest_count_without_sections(self, fault_server, captured_config):
|
||||
"""A syntactically valid STATUS_MANIFEST whose bodies never arrive."""
|
||||
sock = _raw_connect(fault_server)
|
||||
sock.sendall(captured_config)
|
||||
assert _recv_exact(sock, 4) is not None
|
||||
|
||||
sock.sendall(struct.pack("<i", STATUS_MANIFEST) + struct.pack("<i", 3))
|
||||
# Announce three keeps but send none; then drop.
|
||||
_abrupt_close(sock)
|
||||
|
||||
_assert_alive(fault_server)
|
||||
_recover(fault_server, "manifest body truncation")
|
||||
|
||||
|
||||
# --- abrupt truncation of a real transfer ----------------------------------
|
||||
|
||||
|
||||
class _TruncatingProxy:
|
||||
"""Forward at most ``max_client_bytes`` from client to server, then reset
|
||||
both ends mid-stream. Runs one client command (which is expected to fail)."""
|
||||
|
||||
def __init__(self, target_port, max_client_bytes):
|
||||
self.target = ("127.0.0.1", target_port)
|
||||
self.max_client_bytes = max_client_bytes
|
||||
self.listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM)
|
||||
self.listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1)
|
||||
self.listener.bind(("127.0.0.1", 0))
|
||||
self.listener.listen(1)
|
||||
self.listener.settimeout(20)
|
||||
self.port = self.listener.getsockname()[1]
|
||||
|
||||
def run(self, cmd):
|
||||
def serve():
|
||||
try:
|
||||
client, _ = self.listener.accept()
|
||||
except OSError:
|
||||
return
|
||||
try:
|
||||
backend = socket.create_connection(self.target, timeout=10)
|
||||
except OSError:
|
||||
client.close()
|
||||
return
|
||||
# A short receive timeout bounds the case where the client has
|
||||
# nothing left to send and is waiting on the server: the proxy then
|
||||
# cuts the stream anyway instead of stalling the test.
|
||||
client.settimeout(2)
|
||||
backend.settimeout(20)
|
||||
forwarded = 0
|
||||
try:
|
||||
while forwarded < self.max_client_bytes:
|
||||
data = client.recv(65536)
|
||||
if not data:
|
||||
break
|
||||
room = self.max_client_bytes - forwarded
|
||||
take = data[:room]
|
||||
backend.sendall(take)
|
||||
forwarded += len(take)
|
||||
if forwarded >= self.max_client_bytes:
|
||||
break
|
||||
except (OSError, socket.timeout):
|
||||
pass
|
||||
for sock in (client, backend):
|
||||
try:
|
||||
sock.setsockopt(socket.SOL_SOCKET, socket.SO_LINGER, struct.pack("ii", 1, 0))
|
||||
except OSError:
|
||||
pass
|
||||
try:
|
||||
sock.close()
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
thread = threading.Thread(target=serve)
|
||||
thread.start()
|
||||
try:
|
||||
subprocess.run(cmd, capture_output=True, text=True, timeout=30)
|
||||
finally:
|
||||
thread.join(20)
|
||||
self.listener.close()
|
||||
|
||||
|
||||
class TestAbruptMidTransferDisconnect:
|
||||
def test_client_stream_cut_at_offsets(self, fault_server, captured_config):
|
||||
"""Cut the real client stream at offsets anchored to the config frame's
|
||||
actual size: mid-config, right after the config, and into the operation
|
||||
stream -- each followed by an RST of both ends."""
|
||||
config_len = len(captured_config)
|
||||
cuts = sorted({max(1, config_len // 2), max(1, config_len - 1), config_len + 8,
|
||||
config_len + 256})
|
||||
for cut in cuts:
|
||||
proxy = _TruncatingProxy(fault_server.port, cut)
|
||||
cmd = [
|
||||
os.path.join(os.path.dirname(__file__), "..", "..", "build", "client"),
|
||||
"--source-dir",
|
||||
SOURCE_DIR,
|
||||
"--dest-dir",
|
||||
DEST_DIR,
|
||||
"--save-to-disk",
|
||||
"--server-port",
|
||||
str(proxy.port),
|
||||
]
|
||||
# The client is expected to fail; what matters is the server survives.
|
||||
proxy.run(cmd)
|
||||
_assert_alive(fault_server)
|
||||
_recover(fault_server, "abrupt mid-transfer disconnects")
|
||||
@@ -16,6 +16,7 @@
|
||||
#include "test_file_sendfile.h"
|
||||
#include "test_fuzz_smoke.h"
|
||||
#include "test_glob.h"
|
||||
#include "test_hardlink.h"
|
||||
#include "test_iconv.h"
|
||||
#include "test_log.h"
|
||||
#include "test_metadata.h"
|
||||
@@ -88,6 +89,7 @@ int main() {
|
||||
RUN_TEST(test_server_cli);
|
||||
RUN_TEST(test_fuzz_smoke);
|
||||
RUN_TEST(test_xattr);
|
||||
RUN_TEST(test_hardlink);
|
||||
|
||||
printf("\n\033[1;36m=== TEST SUMMARY ===\033[0m\n");
|
||||
printf("Total Tests Run: %d\n", tests_run);
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
#include "test_client_cli.h"
|
||||
#include "checksum.h"
|
||||
#include "client_send.h"
|
||||
#include "client_validation.h"
|
||||
#include "chmod.h"
|
||||
#include "config.h"
|
||||
@@ -579,6 +580,80 @@ static void test_parse_args_invalid_server_port() {
|
||||
config_delete(cfg);
|
||||
}
|
||||
|
||||
/* --port is a documented rsync-style alias for --server-port; both the
|
||||
* two-argument and the inline "=" spellings must work. */
|
||||
static void test_parse_args_port_alias() {
|
||||
Config* cfg = config_create();
|
||||
char* argv_space[] = {"fastsync", "--port", "9000", "/src", "/dst"};
|
||||
int positional_args[2];
|
||||
int positional_count = 0;
|
||||
EXPECT_EQ_INT(parse_args(cfg, 5, argv_space, positional_args, &positional_count), 0);
|
||||
EXPECT_EQ_INT(cfg->server_port, 9000);
|
||||
config_delete(cfg);
|
||||
|
||||
cfg = config_create();
|
||||
char* argv_inline[] = {"fastsync", "--port=9001", "/src", "/dst"};
|
||||
positional_count = 0;
|
||||
EXPECT_EQ_INT(parse_args(cfg, 4, argv_inline, positional_args, &positional_count), 0);
|
||||
EXPECT_EQ_INT(cfg->server_port, 9001);
|
||||
config_delete(cfg);
|
||||
|
||||
cfg = config_create();
|
||||
char* argv_long[] = {"fastsync", "--server-port=9002", "/src", "/dst"};
|
||||
positional_count = 0;
|
||||
EXPECT_EQ_INT(parse_args(cfg, 4, argv_long, positional_args, &positional_count), 0);
|
||||
EXPECT_EQ_INT(cfg->server_port, 9002);
|
||||
config_delete(cfg);
|
||||
}
|
||||
|
||||
/* --threads=N sizes the pipeline scanner; bare -j/--threads keeps the default
|
||||
* (scanner_threads == 0), and invalid values are rejected. */
|
||||
static void test_parse_args_threads() {
|
||||
Config* cfg = config_create();
|
||||
char* argv_eq[] = {"fastsync", "--threads=8", "/src", "/dst"};
|
||||
int positional_args[2];
|
||||
int positional_count = 0;
|
||||
EXPECT_EQ_INT(parse_args(cfg, 4, argv_eq, positional_args, &positional_count), 0);
|
||||
EXPECT_TRUE(cfg->use_multithreading);
|
||||
EXPECT_EQ_INT(cfg->scanner_threads, 8);
|
||||
config_delete(cfg);
|
||||
|
||||
cfg = config_create();
|
||||
char* argv_short[] = {"fastsync", "-j", "/src", "/dst"};
|
||||
positional_count = 0;
|
||||
EXPECT_EQ_INT(parse_args(cfg, 4, argv_short, positional_args, &positional_count), 0);
|
||||
EXPECT_TRUE(cfg->use_multithreading);
|
||||
EXPECT_EQ_INT(cfg->scanner_threads, 0);
|
||||
config_delete(cfg);
|
||||
|
||||
cfg = config_create();
|
||||
char* argv_long[] = {"fastsync", "--threads", "/src", "/dst"};
|
||||
positional_count = 0;
|
||||
EXPECT_EQ_INT(parse_args(cfg, 4, argv_long, positional_args, &positional_count), 0);
|
||||
EXPECT_TRUE(cfg->use_multithreading);
|
||||
EXPECT_EQ_INT(cfg->scanner_threads, 0);
|
||||
config_delete(cfg);
|
||||
|
||||
const char* bad[] = {"--threads=0", "--threads=-3", "--threads=abc", "--threads=257"};
|
||||
for (size_t i = 0; i < sizeof(bad) / sizeof(bad[0]); i++) {
|
||||
cfg = config_create();
|
||||
char* argv_bad[] = {"fastsync", (char*)bad[i], "/src", "/dst"};
|
||||
positional_count = 0;
|
||||
EXPECT_EQ_INT(parse_args(cfg, 4, argv_bad, positional_args, &positional_count), -1);
|
||||
config_delete(cfg);
|
||||
}
|
||||
}
|
||||
|
||||
/* The graceful-abort flag is a plain sig_atomic_t toggled by the handler. */
|
||||
static void test_client_abort_flag() {
|
||||
client_abort_requested = 0;
|
||||
EXPECT_FALSE(client_abort_pending());
|
||||
client_abort_requested = 1;
|
||||
EXPECT_TRUE(client_abort_pending());
|
||||
client_abort_requested = 0;
|
||||
EXPECT_FALSE(client_abort_pending());
|
||||
}
|
||||
|
||||
/* Test parse_args rejects invalid compression level (-z/--compress) */
|
||||
static void test_parse_args_invalid_compression_level() {
|
||||
Config* cfg = config_create();
|
||||
@@ -3224,6 +3299,9 @@ void test_client_cli() {
|
||||
test_parse_args_invalid_port();
|
||||
test_parse_args_non_numeric_port();
|
||||
test_parse_args_invalid_server_port();
|
||||
test_parse_args_port_alias();
|
||||
test_parse_args_threads();
|
||||
test_client_abort_flag();
|
||||
test_parse_args_invalid_compression_level();
|
||||
test_parse_args_valid_compression_level();
|
||||
test_parse_args_debug_flags();
|
||||
|
||||
@@ -0,0 +1,203 @@
|
||||
#include "test_hardlink.h"
|
||||
#include "hardlink.h"
|
||||
#include "test_utils.h"
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/types.h>
|
||||
|
||||
/* A fresh table starts empty and destroy accepts NULL / a fresh table. */
|
||||
static void test_hardlink_create_destroy() {
|
||||
hardlink_table_destroy(NULL);
|
||||
|
||||
HardLinkTable* table = hardlink_table_create();
|
||||
EXPECT_NOT_NULL(table);
|
||||
EXPECT_EQ_INT((int)table->count, 0);
|
||||
EXPECT_EQ_INT((int)table->capacity, 0);
|
||||
EXPECT_EQ_INT(table->next_gid, 1);
|
||||
hardlink_table_destroy(table);
|
||||
}
|
||||
|
||||
/* The first member of an (dev, ino) group is data-carrying and owns the group;
|
||||
* every later member gets the SAME gid, is not first, and points back at the
|
||||
* first member's wire path. */
|
||||
static void test_hardlink_grouping() {
|
||||
HardLinkTable* table = hardlink_table_create();
|
||||
EXPECT_NOT_NULL(table);
|
||||
|
||||
int gid_a = -1, gid_b = -1;
|
||||
bool first_a = false, first_b = false;
|
||||
char* first_path_a = NULL;
|
||||
char* first_path_b = NULL;
|
||||
|
||||
EXPECT_TRUE(
|
||||
hardlink_table_assign(table, "dir/first.txt", 7, 42, &gid_a, &first_a, &first_path_a));
|
||||
EXPECT_TRUE(first_a);
|
||||
EXPECT_EQ_INT(gid_a, 1);
|
||||
EXPECT_NOT_NULL(first_path_a);
|
||||
EXPECT_EQ_STR(first_path_a, "dir/first.txt");
|
||||
|
||||
EXPECT_TRUE(
|
||||
hardlink_table_assign(table, "dir/second.txt", 7, 42, &gid_b, &first_b, &first_path_b));
|
||||
EXPECT_FALSE(first_b);
|
||||
EXPECT_EQ_INT(gid_b, gid_a);
|
||||
EXPECT_NOT_NULL(first_path_b);
|
||||
EXPECT_EQ_STR(first_path_b, "dir/first.txt");
|
||||
|
||||
/* Two members map onto a single stored group. */
|
||||
EXPECT_EQ_INT((int)table->count, 1);
|
||||
EXPECT_EQ_INT(table->next_gid, 2);
|
||||
|
||||
free(first_path_a);
|
||||
free(first_path_b);
|
||||
hardlink_table_destroy(table);
|
||||
}
|
||||
|
||||
/* A different inode on the same device is a distinct group with a fresh gid. */
|
||||
static void test_hardlink_distinct_inode() {
|
||||
HardLinkTable* table = hardlink_table_create();
|
||||
EXPECT_NOT_NULL(table);
|
||||
|
||||
int gid1 = -1, gid2 = -1;
|
||||
bool first1 = false, first2 = false;
|
||||
char* path1 = NULL;
|
||||
char* path2 = NULL;
|
||||
|
||||
EXPECT_TRUE(hardlink_table_assign(table, "a", 7, 100, &gid1, &first1, &path1));
|
||||
EXPECT_TRUE(first1);
|
||||
EXPECT_TRUE(hardlink_table_assign(table, "b", 7, 101, &gid2, &first2, &path2));
|
||||
EXPECT_TRUE(first2);
|
||||
EXPECT_TRUE(gid1 != gid2);
|
||||
EXPECT_EQ_INT(gid1, 1);
|
||||
EXPECT_EQ_INT(gid2, 2);
|
||||
EXPECT_EQ_STR(path1, "a");
|
||||
EXPECT_EQ_STR(path2, "b");
|
||||
|
||||
free(path1);
|
||||
free(path2);
|
||||
hardlink_table_destroy(table);
|
||||
}
|
||||
|
||||
/* Identical (dev, ino) on a DIFFERENT device must never be conflated: inode
|
||||
* numbers are only unique per filesystem, so grouping is scoped by st_dev. */
|
||||
static void test_hardlink_distinct_device() {
|
||||
HardLinkTable* table = hardlink_table_create();
|
||||
EXPECT_NOT_NULL(table);
|
||||
|
||||
int gid1 = -1, gid2 = -1;
|
||||
bool first1 = false, first2 = false;
|
||||
char* path1 = NULL;
|
||||
char* path2 = NULL;
|
||||
|
||||
EXPECT_TRUE(hardlink_table_assign(table, "dev_a/one", 1, 55, &gid1, &first1, &path1));
|
||||
EXPECT_TRUE(hardlink_table_assign(table, "dev_b/one", 2, 55, &gid2, &first2, &path2));
|
||||
EXPECT_TRUE(first1);
|
||||
EXPECT_TRUE(first2);
|
||||
EXPECT_TRUE(gid1 != gid2);
|
||||
EXPECT_EQ_INT((int)table->count, 2);
|
||||
|
||||
free(path1);
|
||||
free(path2);
|
||||
hardlink_table_destroy(table);
|
||||
}
|
||||
|
||||
/* The table owns deep copies of every path: mutating (or freeing) the caller's
|
||||
* buffer after assign must not affect the stored / returned paths. */
|
||||
static void test_hardlink_path_ownership() {
|
||||
HardLinkTable* table = hardlink_table_create();
|
||||
EXPECT_NOT_NULL(table);
|
||||
|
||||
char caller[] = "owned/path";
|
||||
int gid = -1;
|
||||
bool is_first = false;
|
||||
char* out = NULL;
|
||||
|
||||
EXPECT_TRUE(hardlink_table_assign(table, caller, 3, 9, &gid, &is_first, &out));
|
||||
EXPECT_TRUE(is_first);
|
||||
/* The returned pointer is a distinct allocation, not the caller's buffer. */
|
||||
EXPECT_TRUE(out != caller);
|
||||
EXPECT_TRUE(table->items[0].first_path != caller);
|
||||
|
||||
memset(caller, 'X', sizeof(caller) - 1);
|
||||
EXPECT_EQ_STR(out, "owned/path");
|
||||
EXPECT_EQ_STR(table->items[0].first_path, "owned/path");
|
||||
|
||||
/* Later members get their own independent copy of the first path. */
|
||||
char second_caller[] = "owned/second";
|
||||
int gid2 = -1;
|
||||
bool first2 = true;
|
||||
char* out2 = NULL;
|
||||
EXPECT_TRUE(hardlink_table_assign(table, second_caller, 3, 9, &gid2, &first2, &out2));
|
||||
EXPECT_FALSE(first2);
|
||||
EXPECT_EQ_STR(out2, "owned/path");
|
||||
EXPECT_TRUE(out2 != table->items[0].first_path);
|
||||
EXPECT_TRUE(out2 != out);
|
||||
|
||||
free(out);
|
||||
free(out2);
|
||||
hardlink_table_destroy(table);
|
||||
}
|
||||
|
||||
/* Bad arguments must be rejected without touching the table. */
|
||||
static void test_hardlink_reject_bad_args() {
|
||||
HardLinkTable* table = hardlink_table_create();
|
||||
EXPECT_NOT_NULL(table);
|
||||
|
||||
int gid = 0;
|
||||
bool is_first = false;
|
||||
char* out = NULL;
|
||||
|
||||
EXPECT_FALSE(hardlink_table_assign(NULL, "x", 1, 1, &gid, &is_first, &out));
|
||||
EXPECT_FALSE(hardlink_table_assign(table, NULL, 1, 1, &gid, &is_first, &out));
|
||||
EXPECT_FALSE(hardlink_table_assign(table, "x", 1, 1, NULL, &is_first, &out));
|
||||
EXPECT_FALSE(hardlink_table_assign(table, "x", 1, 1, &gid, NULL, &out));
|
||||
EXPECT_FALSE(hardlink_table_assign(table, "x", 1, 1, &gid, &is_first, NULL));
|
||||
EXPECT_EQ_INT((int)table->count, 0);
|
||||
|
||||
hardlink_table_destroy(table);
|
||||
}
|
||||
|
||||
/* Many distinct groups grow the item array through its realloc path and keep
|
||||
* gid assignment stable and monotonic. */
|
||||
static void test_hardlink_many_groups() {
|
||||
HardLinkTable* table = hardlink_table_create();
|
||||
EXPECT_NOT_NULL(table);
|
||||
|
||||
const int n = 200;
|
||||
for (int i = 0; i < n; i++) {
|
||||
int gid = -1;
|
||||
bool is_first = false;
|
||||
char* out = NULL;
|
||||
char path[32];
|
||||
snprintf(path, sizeof(path), "file_%d", i);
|
||||
EXPECT_TRUE(hardlink_table_assign(table, path, 1, (ino_t)(1000 + i), &gid, &is_first, &out));
|
||||
EXPECT_TRUE(is_first);
|
||||
EXPECT_EQ_INT(gid, i + 1);
|
||||
EXPECT_EQ_STR(out, path);
|
||||
free(out);
|
||||
}
|
||||
EXPECT_EQ_INT((int)table->count, n);
|
||||
EXPECT_EQ_INT(table->next_gid, n + 1);
|
||||
|
||||
/* Re-querying an existing inode still reports the original gid. */
|
||||
int gid = -1;
|
||||
bool is_first = true;
|
||||
char* out = NULL;
|
||||
EXPECT_TRUE(hardlink_table_assign(table, "file_7_again", 1, 1007, &gid, &is_first, &out));
|
||||
EXPECT_FALSE(is_first);
|
||||
EXPECT_EQ_INT(gid, 8);
|
||||
EXPECT_EQ_STR(out, "file_7");
|
||||
free(out);
|
||||
|
||||
hardlink_table_destroy(table);
|
||||
}
|
||||
|
||||
void test_hardlink() {
|
||||
test_hardlink_create_destroy();
|
||||
test_hardlink_grouping();
|
||||
test_hardlink_distinct_inode();
|
||||
test_hardlink_distinct_device();
|
||||
test_hardlink_path_ownership();
|
||||
test_hardlink_reject_bad_args();
|
||||
test_hardlink_many_groups();
|
||||
}
|
||||
@@ -0,0 +1,6 @@
|
||||
#ifndef TEST_HARDLINK_H
|
||||
#define TEST_HARDLINK_H
|
||||
|
||||
void test_hardlink(void);
|
||||
|
||||
#endif
|
||||
@@ -460,6 +460,96 @@ static void test_send_receive_status_timed() {
|
||||
close(p[0]);
|
||||
}
|
||||
|
||||
static bool keepalive_always_abort(void) {
|
||||
return true;
|
||||
}
|
||||
|
||||
/* A pre-buffered KEEPALIVE reply from the peer must be consumed transparently,
|
||||
leaving the first real status visible to the caller. */
|
||||
static void test_receive_status_keepalive_skips_reply() {
|
||||
int p[2];
|
||||
EXPECT_EQ_INT(pipe(p), 0);
|
||||
ProtocolSession session;
|
||||
protocol_session_init(&session, p[0], p[1]);
|
||||
|
||||
EXPECT_TRUE(protocol_send_status(&session, STATUS_KEEPALIVE));
|
||||
EXPECT_TRUE(protocol_send_status(&session, STATUS_OK));
|
||||
|
||||
Status received = STATUS_ERROR;
|
||||
EXPECT_TRUE(protocol_receive_status_keepalive(&session, &received, 5, 1, NULL));
|
||||
EXPECT_EQ_INT((int)received, (int)STATUS_OK);
|
||||
|
||||
close(p[0]);
|
||||
close(p[1]);
|
||||
}
|
||||
|
||||
/* The abort callback ends the wait immediately, before any keepalive traffic. */
|
||||
static void test_receive_status_keepalive_aborts() {
|
||||
int p[2];
|
||||
EXPECT_EQ_INT(pipe(p), 0);
|
||||
ProtocolSession session;
|
||||
protocol_session_init(&session, p[0], p[1]);
|
||||
|
||||
Status received = STATUS_ERROR;
|
||||
EXPECT_FALSE(
|
||||
protocol_receive_status_keepalive(&session, &received, 5, 1, keepalive_always_abort));
|
||||
|
||||
close(p[0]);
|
||||
close(p[1]);
|
||||
}
|
||||
|
||||
typedef struct {
|
||||
int peer_read_fd;
|
||||
int peer_write_fd;
|
||||
bool replied;
|
||||
} KeepalivePeerArg;
|
||||
|
||||
static int keepalive_peer(void* arg) {
|
||||
KeepalivePeerArg* peer = arg;
|
||||
ProtocolSession session;
|
||||
protocol_session_init(&session, peer->peer_read_fd, peer->peer_write_fd);
|
||||
Status status = STATUS_ERROR;
|
||||
if (protocol_receive_status(&session, &status) && status == STATUS_KEEPALIVE) {
|
||||
/* Model the busy receiver: it sends the real ack first, then the keepalive
|
||||
reply it owes for the queued keepalive (which the client must drain so it
|
||||
does not desynchronize the stream). */
|
||||
peer->replied = protocol_send_status(&session, STATUS_OK) &&
|
||||
protocol_send_status(&session, STATUS_KEEPALIVE);
|
||||
}
|
||||
return thrd_success;
|
||||
}
|
||||
|
||||
/* While the peer is silent the helper must emit STATUS_KEEPALIVE, then consume
|
||||
the peer's ack and drain the keepalive reply that follows it -- proving the
|
||||
inline keepalive loop works without a second writer racing the send path. */
|
||||
static void test_receive_status_keepalive_emits() {
|
||||
int to_client[2];
|
||||
int to_peer[2];
|
||||
EXPECT_EQ_INT(pipe(to_client), 0);
|
||||
EXPECT_EQ_INT(pipe(to_peer), 0);
|
||||
|
||||
ProtocolSession session;
|
||||
protocol_session_init(&session, to_client[0], to_peer[1]);
|
||||
|
||||
KeepalivePeerArg peer = {.peer_read_fd = to_peer[0], .peer_write_fd = to_client[1]};
|
||||
thrd_t thread;
|
||||
EXPECT_EQ_INT(thrd_create(&thread, keepalive_peer, &peer), thrd_success);
|
||||
|
||||
Status received = STATUS_ERROR;
|
||||
EXPECT_TRUE(protocol_receive_status_keepalive(&session, &received, 10, 1, NULL));
|
||||
EXPECT_EQ_INT((int)received, (int)STATUS_OK);
|
||||
|
||||
int result = 0;
|
||||
EXPECT_EQ_INT(thrd_join(thread, &result), thrd_success);
|
||||
EXPECT_EQ_INT(result, thrd_success);
|
||||
EXPECT_TRUE(peer.replied);
|
||||
|
||||
close(to_client[0]);
|
||||
close(to_client[1]);
|
||||
close(to_peer[0]);
|
||||
close(to_peer[1]);
|
||||
}
|
||||
|
||||
void test_protocol() {
|
||||
test_send_receive_n_data();
|
||||
test_send_receive_n_data_zero();
|
||||
@@ -471,6 +561,9 @@ void test_protocol() {
|
||||
test_send_receive_status();
|
||||
test_protocol_session_io_timeout();
|
||||
test_send_receive_status_timed();
|
||||
test_receive_status_keepalive_skips_reply();
|
||||
test_receive_status_keepalive_aborts();
|
||||
test_receive_status_keepalive_emits();
|
||||
test_receive_n_data_truncated();
|
||||
test_receive_str_truncated();
|
||||
test_max_alloc_rejects_single_buffer();
|
||||
|
||||
Reference in New Issue
Block a user