From 1fa2fbd2669cae18e2087c95073f41b508d66d00 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 13 Sep 2026 06:22:28 +0200 Subject: [PATCH 1/3] feat(client): --port alias, --threads=N, graceful abort, keepalive --- README.md | 10 +-- RSYNC_COMPAT.md | 4 +- src/client/client_cli.c | 105 ++++++++++++++++++++++++++---- src/client/client_send.c | 54 +++++++++++++++- src/client/client_send.h | 9 +++ src/client/scanner.h | 4 ++ src/client/usage.c | 7 +- src/shared/config.c | 7 +- src/shared/config.h | 12 ++-- src/shared/protocol.c | 135 +++++++++++++++++++++++++++++++++++++++ src/shared/protocol.h | 19 ++++++ tests/test_client_cli.c | 78 ++++++++++++++++++++++ tests/test_protocol.c | 93 +++++++++++++++++++++++++++ 13 files changed, 501 insertions(+), 36 deletions(-) diff --git a/README.md b/README.md index 6cd07d0..9820885 100644 --- a/README.md +++ b/README.md @@ -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 ` | Set the zstd compression level. | | `--zc ` | Alias for `--compress-choice`. FastSync supports `zstd` and `none`. | @@ -392,7 +392,7 @@ features without changing the meaning of ordinary compatibility options. | `--delta-block ` | Set the FastSync delta block size (`--block-size` is an alias). | | `--delta-max ` | Limit files eligible for FastSync delta transfer. | | `--server-host ` | Select the TCP server host. | -| `--server-port ` | Select the TCP server port. | +| `--server-port ` | Select the TCP server port (`--port ` and `--port=` are rsync-friendly aliases). | | `--tls` | Enable TLS for TCP transport. | | `--bwlimit ` | Apply token-bucket bandwidth limiting. | | `--progress` | Show transfer progress and throughput. | @@ -478,7 +478,7 @@ link-target transfer remains incomplete. | | `--dest-dir ` | Set the destination directory explicitly. | | `--save-to-disk` | Enable server-side disk persistence. | | `--server-host ` | TCP server address. | -| `--server-port ` | TCP server port. | +| `--server-port ` | TCP server port. `--port ` / `--port=` is an alias. | | `--tls` | Enable TLS. Requires `--cert` and `--key`. | | `--cert ` | TLS certificate file. | | `--key ` | TLS private key file. | diff --git a/RSYNC_COMPAT.md b/RSYNC_COMPAT.md index 6fb91c8..ea171bc 100644 --- a/RSYNC_COMPAT.md +++ b/RSYNC_COMPAT.md @@ -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`) | diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 43136af..73812ef 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -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,26 @@ #include #include +/* 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; + +#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. Only the production client installs it. */ +static void client_signal_handler(int signo) { + (void)signo; + client_abort_requested = 1; +} +#endif + +bool client_abort_pending(void) { + return client_abort_requested != 0; +} + #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 +162,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 +1326,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 +1349,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 : ""); + 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 : ""); - 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 +2043,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; diff --git a/src/client/client_send.c b/src/client/client_send.c index e026e53..1febec8 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -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; @@ -2036,6 +2068,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 +2140,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 @@ -2276,7 +2324,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, diff --git a/src/client/client_send.h b/src/client/client_send.h index 8c85106..06d6c6d 100644 --- a/src/client/client_send.h +++ b/src/client/client_send.h @@ -4,6 +4,15 @@ #include "chunk.h" #include "config.h" #include "transport_tcp.h" +#include +#include + +/* 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); /* Both sender entry points BORROW `config` for the duration of the call; they * never free it, and the caller retains ownership (freeing it with diff --git a/src/client/scanner.h b/src/client/scanner.h index 3db1ce9..51200e2 100644 --- a/src/client/scanner.h +++ b/src/client/scanner.h @@ -14,6 +14,10 @@ #include #include +/* 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 diff --git a/src/client/usage.c b/src/client/usage.c index f63a5c3..0045d26 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -2,6 +2,7 @@ #include #include #include +#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 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 Server IP address (default: 127.0.0.1)\n"); printf(" --server-port Server port (default: 8080)\n"); + printf(" --port Alias for --server-port\n"); printf(" --password-file 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"); diff --git a/src/shared/config.c b/src/shared/config.c index c0f1412..64adf10 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -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) { diff --git a/src/shared/config.h b/src/shared/config.h index b0ae890..0aba9d6 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -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 diff --git a/src/shared/protocol.c b/src/shared/protocol.c index f3dd81b..0a8dac9 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -609,6 +609,136 @@ 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; + 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 = POLLIN}; + 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) + 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. */ + int interval_ms = keepalive_interval_sec * 1000; + 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. */ + while (replies_seen < keepalives_sent) { + Status drained; + if (!protocol_read_status_until(session, &drained, &deadline)) + return false; + 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 +778,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); +} diff --git a/src/shared/protocol.h b/src/shared/protocol.h index b27c2d1..e55a742 100644 --- a/src/shared/protocol.h +++ b/src/shared/protocol.h @@ -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 diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index e5375b9..b2a7021 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -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(); diff --git a/tests/test_protocol.c b/tests/test_protocol.c index 273861e..3a63125 100644 --- a/tests/test_protocol.c +++ b/tests/test_protocol.c @@ -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(); From 2854a9d149be8269e1f76dfb25f7a47c20387a04 Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 13 Sep 2026 06:24:46 +0200 Subject: [PATCH 2/3] test: fuzz manifest/protocol/xattr, hardlink unit, fault injection --- tests/fuzz/fuzz_manifest.c | 126 ++++++++ tests/fuzz/fuzz_protocol_framing.c | 127 ++++++++ tests/fuzz/fuzz_xattr_block.c | 115 +++++++ tests/integration/test_fault_injection.py | 367 ++++++++++++++++++++++ tests/runner.c | 2 + tests/test_hardlink.c | 203 ++++++++++++ tests/test_hardlink.h | 6 + 7 files changed, 946 insertions(+) create mode 100644 tests/fuzz/fuzz_manifest.c create mode 100644 tests/fuzz/fuzz_protocol_framing.c create mode 100644 tests/fuzz/fuzz_xattr_block.c create mode 100644 tests/integration/test_fault_injection.py create mode 100644 tests/test_hardlink.c create mode 100644 tests/test_hardlink.h diff --git a/tests/fuzz/fuzz_manifest.c b/tests/fuzz/fuzz_manifest.c new file mode 100644 index 0000000..57a9f3e --- /dev/null +++ b/tests/fuzz/fuzz_manifest.c @@ -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 +#include +#include +#include +#include +#include +#include + +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; +} diff --git a/tests/fuzz/fuzz_protocol_framing.c b/tests/fuzz/fuzz_protocol_framing.c new file mode 100644 index 0000000..6bb0359 --- /dev/null +++ b/tests/fuzz/fuzz_protocol_framing.c @@ -0,0 +1,127 @@ +/* + * 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 +#include +#include +#include +#include +#include +#include + +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); + } + + rd = make_stream(data, size, &w); + if (rd >= 0) { + Data* d = receive_data(rd); + 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; +} diff --git a/tests/fuzz/fuzz_xattr_block.c b/tests/fuzz/fuzz_xattr_block.c new file mode 100644 index 0000000..26e3775 --- /dev/null +++ b/tests/fuzz/fuzz_xattr_block.c @@ -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 +#include +#include +#include +#include +#include +#include + +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; +} diff --git a/tests/integration/test_fault_injection.py b/tests/integration/test_fault_injection.py new file mode 100644 index 0000000..8978daf --- /dev/null +++ b/tests/integration/test_fault_injection.py @@ -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("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("= 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") diff --git a/tests/runner.c b/tests/runner.c index 51938e1..9e12323 100644 --- a/tests/runner.c +++ b/tests/runner.c @@ -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); diff --git a/tests/test_hardlink.c b/tests/test_hardlink.c new file mode 100644 index 0000000..054e285 --- /dev/null +++ b/tests/test_hardlink.c @@ -0,0 +1,203 @@ +#include "test_hardlink.h" +#include "hardlink.h" +#include "test_utils.h" +#include +#include +#include +#include + +/* 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(); +} diff --git a/tests/test_hardlink.h b/tests/test_hardlink.h new file mode 100644 index 0000000..c271154 --- /dev/null +++ b/tests/test_hardlink.h @@ -0,0 +1,6 @@ +#ifndef TEST_HARDLINK_H +#define TEST_HARDLINK_H + +void test_hardlink(void); + +#endif From c2df0347ef6423acef7a911e337aef3c52d36f0e Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 13 Sep 2026 06:59:02 +0200 Subject: [PATCH 3/3] fix(client,protocol): EINTR-safe sends, armed abort, keepalive drain grace, TLS WANT_WRITE --- src/client/client_cli.c | 31 +++++++++++++++------ src/client/client_send.c | 8 ++++++ src/client/client_send.h | 4 +++ src/shared/protocol.c | 43 ++++++++++++++++++++++-------- tests/fuzz/fuzz_protocol_framing.c | 5 +++- 5 files changed, 71 insertions(+), 20 deletions(-) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 73812ef..dd644c4 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -34,20 +34,35 @@ * client_send.c, still links the symbol. */ volatile sig_atomic_t 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. Only the production client installs it. */ -static void client_signal_handler(int signo) { - (void)signo; - client_abort_requested = 1; +/* 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; } -#endif 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, diff --git a/src/client/client_send.c b/src/client/client_send.c index 1febec8..db7a560 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -1949,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) @@ -2228,6 +2231,7 @@ send_fail: prepared_scanner_destroy(&prepared); disconnect_transfer_client(client); protocol_session_unbind(); + client_set_abort_armed(false); return ret; } @@ -2257,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 = @@ -2411,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; } diff --git a/src/client/client_send.h b/src/client/client_send.h index 06d6c6d..2605899 100644 --- a/src/client/client_send.h +++ b/src/client/client_send.h @@ -13,6 +13,10 @@ * 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 diff --git a/src/shared/protocol.c b/src/shared/protocol.c index 0a8dac9..57fd2e2 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -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; @@ -618,6 +624,7 @@ 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); @@ -625,7 +632,7 @@ static bool protocol_read_status_until(ProtocolSession* session, Status* status, log_message(LOG_LEVEL_ERROR, "Receive timeout while reading status"); return false; } - struct pollfd pfd = {.fd = session->read_fd, .events = POLLIN}; + 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"); @@ -647,8 +654,10 @@ static bool protocol_read_status_until(ProtocolSession* session, Status* status, 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) + 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; @@ -689,7 +698,8 @@ bool protocol_receive_status_keepalive(ProtocolSession* session, Status* status, /* 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. */ - int interval_ms = keepalive_interval_sec * 1000; + 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); @@ -724,15 +734,26 @@ bool protocol_receive_status_keepalive(ProtocolSession* session, Status* status, * 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. */ - while (replies_seen < keepalives_sent) { - Status drained; - if (!protocol_read_status_until(session, &drained, &deadline)) - return false; - if (drained != STATUS_KEEPALIVE) { - log_message(LOG_LEVEL_ERROR, "Unexpected status while draining keepalive replies"); - return false; + 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++; } - replies_seen++; } *status = final; log_debug_message(LOG_DEBUG_PROTO, "Received Status: %s", status_to_string(*status)); diff --git a/tests/fuzz/fuzz_protocol_framing.c b/tests/fuzz/fuzz_protocol_framing.c index 6bb0359..a80350a 100644 --- a/tests/fuzz/fuzz_protocol_framing.c +++ b/tests/fuzz/fuzz_protocol_framing.c @@ -81,9 +81,12 @@ int LLVMFuzzerTestOneInput(const uint8_t* data, size_t size) { 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(rd); + Data* d = receive_data_limited(rd, 1u << 20); data_destroy(d); close(rd); close(w);