Release v2.26.0 #284
@@ -104,7 +104,7 @@ partial, alternate, and planned behavior.
|
|||||||
| `-c, --checksum` | Verify content by checksum instead of size+mtime |
|
| `-c, --checksum` | Verify content by checksum instead of size+mtime |
|
||||||
| `-z, --compress [level]` | Enable streaming zstd compression (level 1–22, default 5) |
|
| `-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) |
|
| `-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) |
|
| `-m` | rsync `--prune-empty-dirs` (short form now rsync-parity) |
|
||||||
| `--chunk-serialization` | Chunk serialization (batch all files per chunk; long form only) |
|
| `--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) |
|
| `-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`;
|
3. **FileMetadata** — `mode`, `uid`, `gid`, `mtime_sec`, `mtime_nsec`;
|
||||||
uid / gid are advisory wire fields and are never applied by the receiver;
|
uid / gid are advisory wire fields and are never applied by the receiver;
|
||||||
atime is unsupported
|
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
|
5. **Queue** — thread-safe bounded queue with condition variables
|
||||||
6. **DirectoryScanner** — recursive BFS traversal with exclude and include pattern support, max-depth enforcement
|
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 |
|
| 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. |
|
| `-z [level]`, `--compress [level]` | Enable streaming zstd compression, levels 1-22. |
|
||||||
| `--compress-level <n>` | Set the zstd compression level. |
|
| `--compress-level <n>` | Set the zstd compression level. |
|
||||||
| `--zc <alg>` | Alias for `--compress-choice`. FastSync supports `zstd` and `none`. |
|
| `--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-block <bytes>` | Set the FastSync delta block size (`--block-size` is an alias). |
|
||||||
| `--delta-max <bytes>` | Limit files eligible for FastSync delta transfer. |
|
| `--delta-max <bytes>` | Limit files eligible for FastSync delta transfer. |
|
||||||
| `--server-host <host>` | Select the TCP server host. |
|
| `--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. |
|
| `--tls` | Enable TLS for TCP transport. |
|
||||||
| `--bwlimit <KB/s>` | Apply token-bucket bandwidth limiting. |
|
| `--bwlimit <KB/s>` | Apply token-bucket bandwidth limiting. |
|
||||||
| `--progress` | Show transfer progress and throughput. |
|
| `--progress` | Show transfer progress and throughput. |
|
||||||
@@ -478,7 +478,7 @@ link-target transfer remains incomplete. |
|
|||||||
| `--dest-dir <path>` | Set the destination directory explicitly. |
|
| `--dest-dir <path>` | Set the destination directory explicitly. |
|
||||||
| `--save-to-disk` | Enable server-side disk persistence. |
|
| `--save-to-disk` | Enable server-side disk persistence. |
|
||||||
| `--server-host <host>` | TCP server address. |
|
| `--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`. |
|
| `--tls` | Enable TLS. Requires `--cert` and `--key`. |
|
||||||
| `--cert <path>` | TLS certificate file. |
|
| `--cert <path>` | TLS certificate file. |
|
||||||
| `--key <path>` | TLS private key 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) |
|
| `-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 |
|
| `--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 |
|
| `--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** |
|
| `--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** |
|
| `--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 |
|
| 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`) |
|
| `--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`) |
|
| `--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`) |
|
| `-z [level]` / `--compress` | zstd compression level (1-22) (`-c` is now rsync `--checksum`) |
|
||||||
|
|||||||
+92
-13
@@ -12,6 +12,7 @@
|
|||||||
#include "identity.h"
|
#include "identity.h"
|
||||||
#include "log.h"
|
#include "log.h"
|
||||||
#include "protocol.h"
|
#include "protocol.h"
|
||||||
|
#include "scanner.h"
|
||||||
#include "stop_condition.h"
|
#include "stop_condition.h"
|
||||||
#include "transport_tcp.h"
|
#include "transport_tcp.h"
|
||||||
#include "transport_tls.h"
|
#include "transport_tls.h"
|
||||||
@@ -27,6 +28,26 @@
|
|||||||
#include <stdlib.h>
|
#include <stdlib.h>
|
||||||
#include <string.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;
|
||||||
|
|
||||||
|
#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
|
#ifndef FASTSYNC_TEST_BUILD
|
||||||
/* Parse environment variables for source/destination directories and save-to-disk flag. */
|
/* 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,
|
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;
|
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) {
|
static int set_compression_threads_option(int* dest, const char* value) {
|
||||||
if (set_positive_int_option(dest, value, "--compress-threads") != 0)
|
if (set_positive_int_option(dest, value, "--compress-threads") != 0)
|
||||||
return -1;
|
return -1;
|
||||||
@@ -1288,10 +1326,21 @@ static bool cli_handle_transfer_flags(CliParseCtx* ctx) {
|
|||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
if (opt_is(arg, "-j", "--threads")) {
|
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;
|
config->use_multithreading = true;
|
||||||
log_info_message(LOG_INFO_MISC, "Enabled Multithreading");
|
log_info_message(LOG_INFO_MISC, "Enabled Multithreading");
|
||||||
return true;
|
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)) {
|
if (opt_is(arg, "--chunk-serialization", NULL)) {
|
||||||
config->use_chunk_serialization = true;
|
config->use_chunk_serialization = true;
|
||||||
log_info_message(LOG_INFO_MISC, "Enabled Chunk Serialization");
|
log_info_message(LOG_INFO_MISC, "Enabled Chunk Serialization");
|
||||||
@@ -1300,29 +1349,48 @@ static bool cli_handle_transfer_flags(CliParseCtx* ctx) {
|
|||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
|
|
||||||
/* Network/IO options: --server-port, --bwlimit, --chunk-size, --log-file and
|
/* Parse and validate a TCP server port (--server-port, or its rsync-friendly
|
||||||
* --stderr. Returns true when the argument was consumed. */
|
* 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) {
|
static bool cli_handle_io_options(CliParseCtx* ctx) {
|
||||||
Config* config = ctx->config;
|
Config* config = ctx->config;
|
||||||
const char* arg = ctx->argv[ctx->i];
|
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) {
|
if (ctx->i + 1 >= ctx->argc) {
|
||||||
log_message(LOG_LEVEL_ERROR, "missing argument for %s", arg);
|
log_message(LOG_LEVEL_ERROR, "missing argument for %s", arg);
|
||||||
ctx->exit_code = -1;
|
ctx->exit_code = -1;
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
if (!parse_positive_int(ctx->argv[++ctx->i], &config->server_port)) {
|
if (set_server_port_option(config, ctx->argv[++ctx->i], arg) != 0)
|
||||||
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);
|
|
||||||
ctx->exit_code = -1;
|
ctx->exit_code = -1;
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
if (config->server_port > 65535) {
|
/* rsync users commonly write --port=NNNN; --server-port=NNNN is accepted too
|
||||||
log_message(LOG_LEVEL_ERROR, "server port must be 1-65535");
|
* 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;
|
ctx->exit_code = -1;
|
||||||
}
|
|
||||||
return true;
|
return true;
|
||||||
}
|
}
|
||||||
if (opt_is(arg, "--bwlimit", NULL)) {
|
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
|
oversized delta). Ignore SIGPIPE so that a broken TCP connection
|
||||||
surfaces as a clean write error instead of killing the client. */
|
surfaces as a clean write error instead of killing the client. */
|
||||||
signal(SIGPIPE, SIG_IGN);
|
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_source = NULL;
|
||||||
const char* env_dest = NULL;
|
const char* env_dest = NULL;
|
||||||
bool save_to_disk = false;
|
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
|
unlinks) before replying, so the wait uses a generous explicit deadline
|
||||||
instead of the default 60 s receive window. */
|
instead of the default 60 s receive window. */
|
||||||
#define DELETE_ACK_TIMEOUT_SEC 3600
|
#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,
|
static bool send_delete_manifest_early(Client* client, ArrayList* manifest,
|
||||||
ArrayList* protected_prefixes, ArrayList* missing_args) {
|
ArrayList* protected_prefixes, ArrayList* missing_args) {
|
||||||
@@ -894,8 +898,21 @@ static bool send_delete_manifest_early(Client* client, ArrayList* manifest,
|
|||||||
0)
|
0)
|
||||||
return false;
|
return false;
|
||||||
Status ack;
|
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;
|
return false;
|
||||||
|
}
|
||||||
if (ack != STATUS_OK) {
|
if (ack != STATUS_OK) {
|
||||||
log_message(LOG_LEVEL_ERROR, "Server failed to delete files before the transfer");
|
log_message(LOG_LEVEL_ERROR, "Server failed to delete files before the transfer");
|
||||||
return false;
|
return false;
|
||||||
@@ -1490,6 +1507,19 @@ static int send_chunks_multithreaded(void* pipeline_context) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
while (true) {
|
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
|
/* Phase 6: stop-elegantly at the next chunk boundary once the --stop-after
|
||||||
/ --stop-at deadline has passed. Everything already sent is finalized by
|
/ --stop-at deadline has passed. Everything already sent is finalized by
|
||||||
the completion tail below; the run still returns success. */
|
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;
|
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
|
||||||
protocol_session_bind(&context->allocation_session);
|
protocol_session_bind(&context->allocation_session);
|
||||||
PreparedScanner prepared;
|
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);
|
pipeline_cancel(context);
|
||||||
protocol_session_unbind();
|
protocol_session_unbind();
|
||||||
return thrd_error;
|
return thrd_error;
|
||||||
@@ -2036,6 +2068,15 @@ int send_files(Config* config) {
|
|||||||
only a prefix of the source. */
|
only a prefix of the source. */
|
||||||
bool scan_stopped_early = false;
|
bool scan_stopped_early = false;
|
||||||
while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
|
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
|
/* 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
|
passed. The scanner may also have stopped early itself; either way the
|
||||||
completion tail below keeps everything already sent. */
|
completion tail below keeps everything already sent. */
|
||||||
@@ -2099,6 +2140,13 @@ int send_files(Config* config) {
|
|||||||
goto send_fail;
|
goto send_fail;
|
||||||
if (directory_scanner_had_io_error(scanner))
|
if (directory_scanner_had_io_error(scanner))
|
||||||
had_scan_io = true;
|
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
|
/* 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
|
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
|
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. */
|
fills the protected excluded prefixes. */
|
||||||
PreparedScanner prepared;
|
PreparedScanner prepared;
|
||||||
memset(&prepared, 0, sizeof(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)
|
if (prepared_ok && context->excluded_paths)
|
||||||
prepared.options.excluded_paths = context->excluded_paths;
|
prepared.options.excluded_paths = context->excluded_paths;
|
||||||
bool prebuilt = prepared_ok && scan_paths_only(config, &prepared.options, context->manifest,
|
bool prebuilt = prepared_ok && scan_paths_only(config, &prepared.options, context->manifest,
|
||||||
|
|||||||
@@ -4,6 +4,15 @@
|
|||||||
#include "chunk.h"
|
#include "chunk.h"
|
||||||
#include "config.h"
|
#include "config.h"
|
||||||
#include "transport_tcp.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);
|
||||||
|
|
||||||
/* Both sender entry points BORROW `config` for the duration of the call; they
|
/* Both sender entry points BORROW `config` for the duration of the call; they
|
||||||
* never free it, and the caller retains ownership (freeing it with
|
* never free it, and the caller retains ownership (freeing it with
|
||||||
|
|||||||
@@ -14,6 +14,10 @@
|
|||||||
#include <sys/types.h>
|
#include <sys/types.h>
|
||||||
#include <threads.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 {
|
typedef struct {
|
||||||
bool use_metadata;
|
bool use_metadata;
|
||||||
/* Phase 4 metadata capture: -U/--atimes and -N/--crtimes tell the scanner to
|
/* Phase 4 metadata capture: -U/--atimes and -N/--crtimes tell the scanner to
|
||||||
|
|||||||
+6
-1
@@ -2,6 +2,7 @@
|
|||||||
#include <stdio.h>
|
#include <stdio.h>
|
||||||
#include <delta.h>
|
#include <delta.h>
|
||||||
#include <chunk.h>
|
#include <chunk.h>
|
||||||
|
#include "scanner.h"
|
||||||
|
|
||||||
void print_usage(void) {
|
void print_usage(void) {
|
||||||
printf("Usage:\n");
|
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 block size in bytes (default: %d)\n", DELTA_BLOCK_SIZE_DEFAULT);
|
||||||
printf(" --delta-max <n> Max file size for delta transfer (default: %llu)\n",
|
printf(" --delta-max <n> Max file size for delta transfer (default: %llu)\n",
|
||||||
DELTA_MAX_FILE_SIZE);
|
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(" --chunk-serialization Enable chunk serialization (long form only)\n");
|
||||||
printf(" -s, --secluded-args Protect-args compatibility option (no effect; remote\n");
|
printf(" -s, --secluded-args Protect-args compatibility option (no effect; remote\n");
|
||||||
printf(" SSH argv is already built injection-safe)\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(" --save-to-disk Write received files to disk\n");
|
||||||
printf(" --server-host <ip> Server IP address (default: 127.0.0.1)\n");
|
printf(" --server-host <ip> Server IP address (default: 127.0.0.1)\n");
|
||||||
printf(" --server-port <n> Server port (default: 8080)\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(" --password-file <f> Authenticate a host::module/path daemon destination.\n");
|
||||||
printf(" The file's first user:password line supplies the\n");
|
printf(" The file's first user:password line supplies the\n");
|
||||||
printf(" username and password (only a SHA-256 digest of 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->receive_root_directory = NULL;
|
||||||
config->save_to_disk = false;
|
config->save_to_disk = false;
|
||||||
config->use_multithreading = false;
|
config->use_multithreading = false;
|
||||||
|
config->scanner_threads = 0;
|
||||||
config->use_chunk_serialization = false;
|
config->use_chunk_serialization = false;
|
||||||
config->use_compression = false;
|
config->use_compression = false;
|
||||||
config->use_metadata = false;
|
config->use_metadata = false;
|
||||||
@@ -79,7 +80,6 @@ static void config_set_defaults(Config* config) {
|
|||||||
config->stats = false;
|
config->stats = false;
|
||||||
config->max_depth = 0;
|
config->max_depth = 0;
|
||||||
config->log_file = NULL;
|
config->log_file = NULL;
|
||||||
config->queue_size = 100;
|
|
||||||
config->follow_symlinks = false;
|
config->follow_symlinks = false;
|
||||||
config->partial = false;
|
config->partial = false;
|
||||||
config->copy_links = false;
|
config->copy_links = false;
|
||||||
@@ -147,14 +147,11 @@ static void config_set_defaults(Config* config) {
|
|||||||
config->delete_during = false;
|
config->delete_during = false;
|
||||||
config->delete_delay = false;
|
config->delete_delay = false;
|
||||||
config->address = NULL;
|
config->address = NULL;
|
||||||
config->bind_address = NULL;
|
|
||||||
config->ipv6 = false;
|
config->ipv6 = false;
|
||||||
config->ipv4 = false;
|
config->ipv4 = false;
|
||||||
config->sockopts = NULL;
|
config->sockopts = NULL;
|
||||||
config->sockopt_count = 0;
|
config->sockopt_count = 0;
|
||||||
config->daemon = false;
|
config->daemon = false;
|
||||||
config->daemon_config = NULL;
|
|
||||||
config->server_mode = false;
|
|
||||||
config->no_motd = false;
|
config->no_motd = false;
|
||||||
config->checksum = false;
|
config->checksum = false;
|
||||||
config->checksum_algo = CHECKSUM_ALGO_XXH64;
|
config->checksum_algo = CHECKSUM_ALGO_XXH64;
|
||||||
@@ -789,9 +786,7 @@ void config_delete(Config* config) {
|
|||||||
free(config->partial_dir);
|
free(config->partial_dir);
|
||||||
free(config->suffix);
|
free(config->suffix);
|
||||||
free(config->address);
|
free(config->address);
|
||||||
free(config->bind_address);
|
|
||||||
free(config->sockopts);
|
free(config->sockopts);
|
||||||
free(config->daemon_config);
|
|
||||||
free(config->compress_choice);
|
free(config->compress_choice);
|
||||||
free(config->chmod_spec);
|
free(config->chmod_spec);
|
||||||
if (config->skip_compress_suffixes) {
|
if (config->skip_compress_suffixes) {
|
||||||
|
|||||||
+6
-6
@@ -81,6 +81,11 @@ typedef struct Config {
|
|||||||
char* receive_root_directory;
|
char* receive_root_directory;
|
||||||
bool save_to_disk;
|
bool save_to_disk;
|
||||||
bool use_multithreading;
|
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_chunk_serialization;
|
||||||
bool use_compression;
|
bool use_compression;
|
||||||
bool use_sendfile;
|
bool use_sendfile;
|
||||||
@@ -170,7 +175,6 @@ typedef struct Config {
|
|||||||
bool stats;
|
bool stats;
|
||||||
int max_depth;
|
int max_depth;
|
||||||
FILE* log_file;
|
FILE* log_file;
|
||||||
int queue_size;
|
|
||||||
bool follow_symlinks;
|
bool follow_symlinks;
|
||||||
bool partial;
|
bool partial;
|
||||||
|
|
||||||
@@ -348,21 +352,17 @@ typedef struct Config {
|
|||||||
|
|
||||||
// PR #181: IPv6 and bind address
|
// PR #181: IPv6 and bind address
|
||||||
char* address;
|
char* address;
|
||||||
char* bind_address;
|
|
||||||
bool ipv6;
|
bool ipv6;
|
||||||
bool ipv4;
|
bool ipv4;
|
||||||
/* --sockopts=OPTIONS (Phase 5, Wave B): strict allowlist of TCP/socket
|
/* --sockopts=OPTIONS (Phase 5, Wave B): strict allowlist of TCP/socket
|
||||||
* options applied via setsockopt after socket() and before connect()/bind().
|
* options applied via setsockopt after socket() and before connect()/bind().
|
||||||
* These are LOCAL socket concerns: they never cross the wire config frame.
|
* These are LOCAL socket concerns: they never cross the wire config frame.
|
||||||
* .address is the outgoing/source bind address (--address); .bind_address is
|
* .address is the outgoing/source bind address (--address). */
|
||||||
* reserved for daemon-side binding and is not wired yet. */
|
|
||||||
SockOptEntry* sockopts;
|
SockOptEntry* sockopts;
|
||||||
int sockopt_count;
|
int sockopt_count;
|
||||||
|
|
||||||
// PR #182: Daemon/server mode
|
// PR #182: Daemon/server mode
|
||||||
bool daemon;
|
bool daemon;
|
||||||
char* daemon_config;
|
|
||||||
bool server_mode;
|
|
||||||
/* --no-motd (Wave C): CLIENT-ONLY, never crosses the wire. Suppresses
|
/* --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
|
* 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
|
* the client reads and discards it to keep the stream in sync. rsync's
|
||||||
|
|||||||
@@ -609,6 +609,136 @@ bool protocol_receive_status_timed(ProtocolSession* session, Status* status, int
|
|||||||
return true;
|
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) {
|
bool send_str(int fd, const char* data) {
|
||||||
return protocol_send_str(legacy_session(-1, fd), 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) {
|
bool receive_status_timed(int fd, Status* status, int timeout_sec) {
|
||||||
return protocol_receive_status_timed(legacy_session(fd, -1), status, 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. */
|
this so the sender does not abort after the deletion already committed. */
|
||||||
bool receive_status_timed(int file_descriptor, Status* status, int timeout_sec);
|
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
|
#endif
|
||||||
|
|||||||
@@ -1,5 +1,6 @@
|
|||||||
#include "test_client_cli.h"
|
#include "test_client_cli.h"
|
||||||
#include "checksum.h"
|
#include "checksum.h"
|
||||||
|
#include "client_send.h"
|
||||||
#include "client_validation.h"
|
#include "client_validation.h"
|
||||||
#include "chmod.h"
|
#include "chmod.h"
|
||||||
#include "config.h"
|
#include "config.h"
|
||||||
@@ -579,6 +580,80 @@ static void test_parse_args_invalid_server_port() {
|
|||||||
config_delete(cfg);
|
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) */
|
/* Test parse_args rejects invalid compression level (-z/--compress) */
|
||||||
static void test_parse_args_invalid_compression_level() {
|
static void test_parse_args_invalid_compression_level() {
|
||||||
Config* cfg = config_create();
|
Config* cfg = config_create();
|
||||||
@@ -3224,6 +3299,9 @@ void test_client_cli() {
|
|||||||
test_parse_args_invalid_port();
|
test_parse_args_invalid_port();
|
||||||
test_parse_args_non_numeric_port();
|
test_parse_args_non_numeric_port();
|
||||||
test_parse_args_invalid_server_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_invalid_compression_level();
|
||||||
test_parse_args_valid_compression_level();
|
test_parse_args_valid_compression_level();
|
||||||
test_parse_args_debug_flags();
|
test_parse_args_debug_flags();
|
||||||
|
|||||||
@@ -460,6 +460,96 @@ static void test_send_receive_status_timed() {
|
|||||||
close(p[0]);
|
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() {
|
void test_protocol() {
|
||||||
test_send_receive_n_data();
|
test_send_receive_n_data();
|
||||||
test_send_receive_n_data_zero();
|
test_send_receive_n_data_zero();
|
||||||
@@ -471,6 +561,9 @@ void test_protocol() {
|
|||||||
test_send_receive_status();
|
test_send_receive_status();
|
||||||
test_protocol_session_io_timeout();
|
test_protocol_session_io_timeout();
|
||||||
test_send_receive_status_timed();
|
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_n_data_truncated();
|
||||||
test_receive_str_truncated();
|
test_receive_str_truncated();
|
||||||
test_max_alloc_rejects_single_buffer();
|
test_max_alloc_rejects_single_buffer();
|
||||||
|
|||||||
Reference in New Issue
Block a user