12 Commits
Author SHA1 Message Date
TapTap f7c6c91e13 fix: set *failed on delta oversize branch; add -m RSF skip test
CI / lint (push) Successful in 20s
CI / sanitizers (address) (push) Successful in 41s
CI / sanitizers (undefined) (push) Successful in 39s
CI / fuzz-build (push) Successful in 16s
CI / coverage (push) Successful in 34s
CI / build-and-test (push) Successful in 1m23s
CI / valgrind (push) Successful in 35s
Review nits from independent review of the four fix branches:
- receive_delta_file STATUS_NEXT oversize branch now sets *failed=true
- receive_incremental_check oversize branch returns NULL (receiver sends the
  single STATUS_ERROR) instead of double-sending
- add multithreaded -m --ignore-existing --remove-source-files integration
  coverage so the writer-thread outcome path is exercised
2026-09-05 13:16:42 +02:00
TapTap 39805e4eee Merge fix/quality-cleanup into dev
fixes #260 (dead code, CLI =form/missing-arg diagnostics)
2026-09-05 13:12:33 +02:00
TapTap 9ef78ef48b Merge fix/perf-delta into dev
fixes #259 (delta_compute hash index)
2026-09-05 13:12:33 +02:00
TapTap 68e7ed57d0 Merge fix/security-hardening into dev
fixes #254 (receiver queue byte-budget) #258 (inplace setuid/truncate)

# Conflicts:
#	src/shared/multiprocessing.c
2026-09-05 13:12:31 +02:00
TapTap 11c591c5aa Merge fix/receiver-correctness into dev
fixes #251 #252 #253 #255 #256 #257 (remove-source-files skips, backup NULL-empty,
partial-dir install, lazy STATUS_CHECK, delta *failed, >64MiB files)
2026-09-05 13:11:41 +02:00
TapTap 1d61a1426e fix: #260 accept --opt=value uniformly and report missing arguments
- Correct misleading doc comments on set_string_option /
  set_positive_int_option / set_nonneg_int_option (they return 0/-1,
  not true/false).
- find_table_option_with_equals() now matches every OPTION_TABLE value
  option (OPT_STRING/OPT_POS_INT/OPT_NONNEG_INT/OPT_ULL), so forms such
  as --max-size=2G, --min-size=1K, --suffix=.bak, --timeout=30,
  --max-depth=5, --backup-dir=X parse instead of dying as 'Unknown option'.
  --max-size/--min-size now accept rsync-style binary suffixes (0 remains a
  valid 'no limit' byte count). Existing special handling for
  --compress-choice, --compress-level, --modify-window=, --chmod=,
  --skip-compress=, --compress-threads= is preserved.
- Options that require a separate value (-p, --exclude, --include,
  --delta-block, --delta-max, --server-port, --bwlimit, --chunk-size,
  --log-file, --exclude-from, --include-from, -T, --skip-compress,
  --compress-threads) now emit an explicit 'missing argument' diagnostic
  instead of falling through to the generic 'Unknown option' branch when
  given as the final argv entry.
- Add unit tests covering the = forms (--max-size=2G, --min-size=1K,
  --suffix=.bak, --timeout=30, --max-depth=5, --backup-dir=X) and a clean
  'missing argument' (not 'Unknown option') diagnostic for trailing
  --exclude/--server-port/--skip-compress/-T.
2026-09-05 12:58:30 +02:00
TapTap beb681c2dc fix: #260 remove dead receive_files() and mkdir_r()
receive_files() in src/server/server.c was a non-static, unprototyped,
zero-caller duplicate of receiver_process()/receiver_receive_files() in
src/server/receiver.c. Delete it along with the includes it uniquely pulled
(chunk.h, metadata.h, sys/stat.h, duplicate quoted unistd.h); remaining code
still uses file.h/config.h/protocol.h (via multiprocessing.h) directly.

mkdir_r() in src/shared/utils.c had zero callers in src/ and tests/; remove
the function and its declaration in utils.h, plus the unused libgen.h include.
2026-09-05 12:58:25 +02:00
TapTap 3704a6adaa test: integration coverage for #251 #252 #253 #257
- remove-source-files keeps sources skipped by --existing/--ignore-existing/
  --update (rsync reference behavior)
- --backup keeps <file>~, --suffix .bak, and --backup-dir backups
- --partial --partial-dir installs completed files in the destination
- a 100 MB file transfers end to end (>64 MiB whole-file cap regression)
2026-09-05 12:46:47 +02:00
TapTap 14c064a093 fix: #251 #252 #253 #255 #256 #257 receiver & remove-source correctness
#251 --remove-source-files deletes sources that were skipped receiver-side
     (--existing/--ignore-existing/--update). The receiver now reports a
     per-file outcome for every processed data file when the sender requests
     removal; the client only unlinks sources the receiver actually wrote.
     add remove_source_files to the wire config and bump the protocol to 2.5.0.
#252 --backup/--suffix/--backup-dir broken by NULL-vs-empty wire loss. Receivers
     canonicalize the empty wire string back to NULL for backup_dir, temp_dir,
     partial_dir and suffix, and --suffix is received unconditionally.
#253 --partial --partial-dir never installed completed files. file_save_to_disk
     now renames a fully written partial-dir file into the real destination.
#255 STATUS_CHECK read the entire old file before the size/mtime quick check.
     Old contents are only read when a checksum compare or delta needs them.
#256 receive_delta_file failure paths did not set *failed, so the caller sent
     STATUS_NEXT and waited for a body that never came. Every NULL return now
     marks the transfer failed.
#257 files >64 MiB could not transfer. Whole-file receive caps raised to the
     256 MiB connection/allocation ceiling (chunk caps stay 64 MiB) and the
     client ignores SIGPIPE so a server-side close surfaces as a clean error.

Unit tests added: config NULL-vs-empty round trip, incremental quick-check
skip/NEXT paths, delta oversize failure, partial-dir install, save-result
skip reporting.
2026-09-05 12:46:44 +02:00
TapTap 2234879b2f fix: #258 normalize mode and truncate on --inplace overwrite
The --inplace branch opened the destination with O_WRONLY|O_CREAT (no
O_TRUNC) and only restored metadata when the sender supplied it.  Two
flaws resulted:
1. An existing destination file kept its original mode when no metadata
   was sent, so setuid/setgid/sticky bits survived an overwrite (a root
   sync could leave a root-owned setuid binary controlled by a client).
2. A shorter payload left stale trailing bytes from the previous version
   because the file was never truncated to the new length.

In the inplace branch of file_to_disk_secure_impl:
- Always trim the file to the new payload length (ftruncate after the
  write) so stale trailing bytes can never survive; sparse targets keep
  their pre-size ftruncate.
- Always normalize the mode after a successful overwrite: apply the
  metadata-derived safe mode when metadata is present (as before), else
  fchmod to a safe default 0644, so setuid/setgid/sticky are cleared in
  both cases.
- The --update newer-destination check still runs before any truncation
  or chmod, preserving the skip semantics.

Adds unit tests in test_file.c: (a) setuid/sticky bits on an existing
destination are cleared after an inplace write with and without metadata,
(b) a shorter inplace payload leaves no trailing stale bytes.
2026-09-05 12:29:46 +02:00
TapTap 98120fc722 fix: #254 bound receiver queue by aggregate payload bytes
The per-connection memory budget (MAX_CONNECTION_MEMORY, 256 MiB) only
charged wire buffers via receive_data_limited.  Decompression buffers and
per-file chunk copies were not accounted for, and the multithreaded
receiver could enqueue up to 100 files (each up to 64 MiB uncompressed)
ahead of a slow disk writer, retaining ~6.4 GiB per connection.  A client
sending highly compressible chunks with little bandwidth could OOM the
host while the reserve never tripped.

Bound the receive pipeline by aggregate payload bytes instead of item
count alone:
- Export MAX_CONNECTION_MEMORY from protocol.h.
- PipelineContextReceiver tracks queued_bytes (payload bytes received but
  not yet released by the disk writer, i.e. queued or in the writer's
  hand) under the existing mutex.
- receiver enqueue now blocks while the queue is full by count OR when
  adding the file would push queued_bytes over the configured byte limit,
  applying backpressure to the sender instead of failing the transfer.
- The disk writer releases the byte budget after each file is freed and
  signals the not-full condition.
- The server sets the byte ceiling to
  MAX_CONNECTION_MEMORY - 2*MAX_CHUNK_SIZE so that the queued payloads
  plus the transient wire/decompression buffers of the one in-flight
  chunk stay within the per-connection budget.

The single-threaded receive path is already bounded: it writes files to
disk before reading the next chunk, so its transient is at most one
chunk's wire + decompressed + copied payload (~3 * MAX_CHUNK_SIZE, below
the budget).  Wire buffers remain charged exactly once by
receive_data_limited; this change does not double charge them.

Adds a deterministic unit test in test_multiprocessing.c proving that an
enqueue which would exceed the byte budget blocks until the writer
releases bytes.
2026-09-05 12:29:41 +02:00
TapTap 5fed0888aa perf: #259 delta_compute hash index over signature blocks 2026-09-05 12:21:26 +02:00
25 changed files with 2006 additions and 350 deletions

No files matched your search

+3 -2
View File
@@ -484,10 +484,11 @@ defaults to the current directory. |
## Protocol and Security ## Protocol and Security
FastSync protocol version `2.4.0` is shared by the client and server. The FastSync protocol version `2.5.0` is shared by the client and server. The
current protocol is sender-driven and includes configuration negotiation, current protocol is sender-driven and includes configuration negotiation,
including the maximum allocation limit, incremental checks, checksums, including the maximum allocation limit, incremental checks, checksums,
manifests, keep-alives, abort handling, and FastSync-native delta messages. manifests, keep-alives, abort handling, per-file remove-source results, and
FastSync-native delta messages.
Client and server versions must currently match exactly. Client and server versions must currently match exactly.
TLS provides encrypted TCP transport. Supplying `--ca` enables certificate TLS provides encrypted TCP transport. Supplying `--ca` enables certificate
+97 -22
View File
@@ -12,6 +12,7 @@
#include "utils.h" #include "utils.h"
#include <errno.h> #include <errno.h>
#include <limits.h> #include <limits.h>
#include <signal.h>
#include <stdbool.h> #include <stdbool.h>
#include <stddef.h> #include <stddef.h>
#include <stdio.h> #include <stdio.h>
@@ -57,7 +58,7 @@ static bool parse_positive_int(const char* s, int* out_val) {
return true; return true;
} }
/* Duplicate a string argument into *dest, freeing the old value. Returns true on success, false on /* Duplicate a string argument into *dest, freeing the old value. Returns 0 on success, -1 on
* failure. */ * failure. */
static int set_string_option(char** dest, const char* value, const char* option_name) { static int set_string_option(char** dest, const char* value, const char* option_name) {
char* dup = str_dup(value); char* dup = str_dup(value);
@@ -70,7 +71,7 @@ static int set_string_option(char** dest, const char* value, const char* option_
return 0; return 0;
} }
/* Parse a string as a positive integer into *dest. Returns true on success, false on error. */ /* Parse a string as a positive integer into *dest. Returns 0 on success, -1 on error. */
static int set_positive_int_option(int* dest, const char* value, const char* option_name) { static int set_positive_int_option(int* dest, const char* value, const char* option_name) {
if (!parse_positive_int(value, dest)) { if (!parse_positive_int(value, dest)) {
log_message(LOG_LEVEL_ERROR, "%s must be a positive integer", option_name); log_message(LOG_LEVEL_ERROR, "%s must be a positive integer", option_name);
@@ -102,7 +103,7 @@ static int set_compression_threads_option(int* dest, const char* value) {
return 0; return 0;
} }
/* Parse a string as a non-negative integer into *dest. Returns true on success, false on error. */ /* Parse a string as a non-negative integer into *dest. Returns 0 on success, -1 on error. */
static int set_nonneg_int_option(int* dest, const char* value, const char* option_name) { static int set_nonneg_int_option(int* dest, const char* value, const char* option_name) {
if (!parse_nonneg_int(value, dest)) { if (!parse_nonneg_int(value, dest)) {
log_message(LOG_LEVEL_ERROR, "%s must be a non-negative integer", option_name); log_message(LOG_LEVEL_ERROR, "%s must be a non-negative integer", option_name);
@@ -237,7 +238,10 @@ static int parse_ull_arg(const char* val, unsigned long long* out, const char* o
return 0; return 0;
} }
static int parse_size_arg(const char* value, unsigned long long* out) { /* Parse a byte count with an optional single-letter binary suffix (K/M/G/T/P/E).
* When allow_zero is false, a bare 0 is rejected (size limits use true, since 0
* means "no limit"). Returns 0 on success, -1 on error. */
static int parse_size_arg_allow_zero(const char* value, unsigned long long* out, bool allow_zero) {
if (!value || *value < '0' || *value > '9') if (!value || *value < '0' || *value > '9')
return -1; return -1;
char* end; char* end;
@@ -281,12 +285,16 @@ static int parse_size_arg(const char* value, unsigned long long* out) {
return -1; return -1;
} }
} }
if (number == 0 || number > ULLONG_MAX / multiplier) if ((!allow_zero && number == 0) || number > ULLONG_MAX / multiplier)
return -1; return -1;
*out = number * multiplier; *out = number * multiplier;
return 0; return 0;
} }
static int parse_size_arg(const char* value, unsigned long long* out) {
return parse_size_arg_allow_zero(value, out, false);
}
/* Append a duplicated pattern to a growable pattern array. Returns 0 on success, -1 on error. */ /* Append a duplicated pattern to a growable pattern array. Returns 0 on success, -1 on error. */
static int config_add_pattern(char*** patterns, int* count, const char* value, static int config_add_pattern(char*** patterns, int* count, const char* value,
const char* optname) { const char* optname) {
@@ -450,6 +458,8 @@ static const OptionEntry* find_table_option(const char* arg) {
return NULL; return NULL;
} }
/* Match a "--opt=value" argument against table options that take a value. Flags,
* no-ops, and unsupported options do not accept an inline "=" value. */
static const OptionEntry* find_table_option_with_equals(const char* arg, const char** value) { static const OptionEntry* find_table_option_with_equals(const char* arg, const char** value) {
const char* equals = strchr(arg, '='); const char* equals = strchr(arg, '=');
if (!equals || equals == arg) if (!equals || equals == arg)
@@ -460,8 +470,8 @@ static const OptionEntry* find_table_option_with_equals(const char* arg, const c
if ((strlen(entry->name) == name_len && strncmp(arg, entry->name, name_len) == 0) || if ((strlen(entry->name) == name_len && strncmp(arg, entry->name, name_len) == 0) ||
(entry->alias && strlen(entry->alias) == name_len && (entry->alias && strlen(entry->alias) == name_len &&
strncmp(arg, entry->alias, name_len) == 0)) { strncmp(arg, entry->alias, name_len) == 0)) {
if (strcmp(entry->name, "--compress-choice") == 0 || if (entry->kind == OPT_STRING || entry->kind == OPT_POS_INT ||
strcmp(entry->name, "--compress-level") == 0) { entry->kind == OPT_NONNEG_INT || entry->kind == OPT_ULL) {
*value = equals + 1; *value = equals + 1;
return entry; return entry;
} }
@@ -516,8 +526,13 @@ static int apply_table_option(Config* config, const OptionEntry* entry, const ch
return set_nonneg_int_option((int*)field, value, entry->name); return set_nonneg_int_option((int*)field, value, entry->name);
case OPT_ULL: { case OPT_ULL: {
unsigned long long v; unsigned long long v;
if (parse_ull_arg(value, &v, entry->name) != 0) /* Size-limit options accept rsync-style suffixes (e.g. --max-size=2G); a
* plain byte count, including 0 ("no limit"), stays valid. */
if (parse_size_arg_allow_zero(value, &v, true) != 0) {
log_message(LOG_LEVEL_ERROR, "%s must be a non-negative size (B, K, M, G, T, P, or E)",
entry->name);
return -1; return -1;
}
*(unsigned long long*)field = v; *(unsigned long long*)field = v;
return 0; return 0;
} }
@@ -667,22 +682,38 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
config->use_multithreading = true; config->use_multithreading = true;
config->use_metadata = true; config->use_metadata = true;
log_info_message(LOG_INFO_MISC, "Enabled archive mode (-c -m -M)"); log_info_message(LOG_INFO_MISC, "Enabled archive mode (-c -m -M)");
} else if (opt_is(argv[i], "-p", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "-p", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (set_positive_int_option(&config->ssh_port, argv[++i], "-p") != 0) if (set_positive_int_option(&config->ssh_port, argv[++i], "-p") != 0)
return -1; return -1;
if (config->ssh_port > 65535) { if (config->ssh_port > 65535) {
log_message(LOG_LEVEL_ERROR, "SSH port must be 1-65535"); log_message(LOG_LEVEL_ERROR, "SSH port must be 1-65535");
return -1; return -1;
} }
} else if (opt_is(argv[i], "--exclude", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--exclude", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (config_add_pattern(&config->exclude_patterns, &config->exclude_count, argv[++i], if (config_add_pattern(&config->exclude_patterns, &config->exclude_count, argv[++i],
"--exclude") != 0) "--exclude") != 0)
return -1; return -1;
} else if (opt_is(argv[i], "--include", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--include", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (config_add_pattern(&config->include_patterns, &config->include_count, argv[++i], if (config_add_pattern(&config->include_patterns, &config->include_count, argv[++i],
"--include") != 0) "--include") != 0)
return -1; return -1;
} else if (opt_is(argv[i], "--delta-block", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--delta-block", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
unsigned long long val; unsigned long long val;
if (parse_ull_arg(argv[++i], &val, "--delta-block") != 0) if (parse_ull_arg(argv[++i], &val, "--delta-block") != 0)
return -1; return -1;
@@ -690,7 +721,11 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
config->delta_block_size = (uint32_t)val; config->delta_block_size = (uint32_t)val;
else else
log_message(LOG_LEVEL_WARNING, "--delta-block value %llu out of range, using default", val); log_message(LOG_LEVEL_WARNING, "--delta-block value %llu out of range, using default", val);
} else if (opt_is(argv[i], "--delta-max", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--delta-max", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
unsigned long long val; unsigned long long val;
if (parse_ull_arg(argv[++i], &val, "--delta-max") != 0) if (parse_ull_arg(argv[++i], &val, "--delta-max") != 0)
return -1; return -1;
@@ -731,7 +766,11 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
} else if (opt_is(argv[i], "-s", NULL)) { } else if (opt_is(argv[i], "-s", 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");
} else if (opt_is(argv[i], "--server-port", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--server-port", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (!parse_positive_int(argv[++i], &config->server_port)) { if (!parse_positive_int(argv[++i], &config->server_port)) {
char* escaped = output_escape(argv[i], false); char* escaped = output_escape(argv[i], false);
log_message(LOG_LEVEL_ERROR, "invalid --server-port value: %s", log_message(LOG_LEVEL_ERROR, "invalid --server-port value: %s",
@@ -743,7 +782,11 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
log_message(LOG_LEVEL_ERROR, "server port must be 1-65535"); log_message(LOG_LEVEL_ERROR, "server port must be 1-65535");
return -1; return -1;
} }
} else if (opt_is(argv[i], "--bwlimit", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--bwlimit", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
unsigned long long kbps; unsigned long long kbps;
if (parse_ull_arg(argv[++i], &kbps, "--bwlimit") != 0) if (parse_ull_arg(argv[++i], &kbps, "--bwlimit") != 0)
return -1; return -1;
@@ -757,7 +800,11 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
} }
io_set_bwlimit(kbps * 1024); io_set_bwlimit(kbps * 1024);
log_info_message(LOG_INFO_MISC, "Set bandwidth limit to %llu KB/s", kbps); log_info_message(LOG_INFO_MISC, "Set bandwidth limit to %llu KB/s", kbps);
} else if (opt_is(argv[i], "--chunk-size", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--chunk-size", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
unsigned long long val; unsigned long long val;
if (parse_ull_arg(argv[++i], &val, "--chunk-size") != 0) if (parse_ull_arg(argv[++i], &val, "--chunk-size") != 0)
return -1; return -1;
@@ -766,7 +813,11 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
return -1; return -1;
} }
config->chunk_size = val; config->chunk_size = val;
} else if (opt_is(argv[i], "--log-file", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--log-file", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (config->log_file) { if (config->log_file) {
fclose(config->log_file); fclose(config->log_file);
config->log_file = NULL; config->log_file = NULL;
@@ -788,11 +839,19 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
} else if (opt_is(argv[i], "--stderr", NULL)) { } else if (opt_is(argv[i], "--stderr", NULL)) {
if (i + 1 >= argc || set_stderr_mode(argv[++i]) != 0) if (i + 1 >= argc || set_stderr_mode(argv[++i]) != 0)
return -1; return -1;
} else if (opt_is(argv[i], "--exclude-from", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--exclude-from", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (read_patterns_from_file(argv[++i], &config->exclude_patterns, &config->exclude_count) != if (read_patterns_from_file(argv[++i], &config->exclude_patterns, &config->exclude_count) !=
0) 0)
return -1; return -1;
} else if (opt_is(argv[i], "--include-from", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--include-from", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (read_patterns_from_file(argv[++i], &config->include_patterns, &config->include_count) != if (read_patterns_from_file(argv[++i], &config->include_patterns, &config->include_count) !=
0) 0)
return -1; return -1;
@@ -817,16 +876,28 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
} else if (opt_is(argv[i], "--info", NULL)) { } else if (opt_is(argv[i], "--info", NULL)) {
if (i + 1 >= argc || parse_info_flags(argv[++i], config) != 0) if (i + 1 >= argc || parse_info_flags(argv[++i], config) != 0)
return -1; return -1;
} else if (opt_is(argv[i], "-T", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "-T", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (set_positive_int_option(&config->timeout, argv[++i], "-T") != 0) if (set_positive_int_option(&config->timeout, argv[++i], "-T") != 0)
return -1; return -1;
} else if (strncmp(argv[i], "--skip-compress=", 16) == 0) { } else if (strncmp(argv[i], "--skip-compress=", 16) == 0) {
if (parse_skip_compress(config, argv[i] + 16) != 0) if (parse_skip_compress(config, argv[i] + 16) != 0)
return -1; return -1;
} else if (opt_is(argv[i], "--skip-compress", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--skip-compress", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (parse_skip_compress(config, argv[++i]) != 0) if (parse_skip_compress(config, argv[++i]) != 0)
return -1; return -1;
} else if (opt_is(argv[i], "--compress-threads", NULL) && i + 1 < argc) { } else if (opt_is(argv[i], "--compress-threads", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (set_compression_threads_option(&config->compression_threads, argv[++i]) != 0) if (set_compression_threads_option(&config->compression_threads, argv[++i]) != 0)
return -1; return -1;
} else if (opt_is(argv[i], "--checksum-choice", "--cc")) { } else if (opt_is(argv[i], "--checksum-choice", "--cc")) {
@@ -903,6 +974,10 @@ static int read_patterns_from_file(const char* filepath, char*** patterns, int*
#ifndef FASTSYNC_TEST_BUILD #ifndef FASTSYNC_TEST_BUILD
int main(int argc, char* argv[]) { int main(int argc, char* argv[]) {
/* The server may close a connection mid-stream (e.g. when it rejects an
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);
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;
+35 -8
View File
@@ -109,16 +109,13 @@ static bool add_chunk_to_manifest(ArrayList* manifest, const Chunk* chunk) {
return true; return true;
} }
static bool finalize_transfer(Client* client) { /* (finalize_transfer is defined after the SourceFile helpers below.) */
Status status;
return send_status(client->file_descriptor, STATUS_FINISHED) &&
receive_status(client->file_descriptor, &status) && status == STATUS_OK;
}
typedef struct { typedef struct SourceFile {
char* path; char* path;
dev_t device; dev_t device;
ino_t inode; ino_t inode;
bool skipped; /* receiver reported the file was not written */
} SourceFile; } SourceFile;
static void source_file_destroy(void* item) { static void source_file_destroy(void* item) {
@@ -135,6 +132,8 @@ static void remove_transferred_sources(const Config* config, ArrayList* paths) {
return; return;
for (int i = 0; i < paths->size; i++) { for (int i = 0; i < paths->size; i++) {
SourceFile* source = paths->items[i]; SourceFile* source = paths->items[i];
if (source->skipped)
continue;
const char* slash = strrchr(source->path, '/'); const char* slash = strrchr(source->path, '/');
const char* leaf = slash ? slash + 1 : source->path; const char* leaf = slash ? slash + 1 : source->path;
char parent[PATH_MAX]; char parent[PATH_MAX];
@@ -177,6 +176,7 @@ static SourceFile* source_file_create(const File* file) {
source->path = str_dup(file->path); source->path = str_dup(file->path);
source->device = st.st_dev; source->device = st.st_dev;
source->inode = st.st_ino; source->inode = st.st_ino;
source->skipped = false;
if (!source->path) { if (!source->path) {
source_file_destroy(source); source_file_destroy(source);
return NULL; return NULL;
@@ -203,6 +203,33 @@ static void mark_sender_done(PipelineContextSender* context) {
mtx_unlock(&context->mutex_progress); mtx_unlock(&context->mutex_progress);
} }
/* Send the final STATUS_FINISHED frame and await the receiver's verdict.
When --remove-source-files is active the receiver acknowledges each data
file it processed, in send order: STATUS_NEXT means the file was written,
STATUS_OK means the file was skipped/unchanged. Skipped sources are marked
so the later removal pass keeps them. */
static bool finalize_transfer(Client* client, const Config* config, ArrayList* remove_sources) {
if (!send_status(client->file_descriptor, STATUS_FINISHED))
return false;
if (config->remove_source_files && remove_sources) {
for (int i = 0; i < remove_sources->size; i++) {
Status per_file;
if (!receive_status(client->file_descriptor, &per_file))
return false;
if (per_file == STATUS_ERROR)
return false;
if (per_file == STATUS_OK) {
((SourceFile*)remove_sources->items[i])->skipped = true;
} else if (per_file != STATUS_NEXT) {
log_message(LOG_LEVEL_ERROR, "Unexpected per-file status from receiver");
return false;
}
}
}
Status status;
return receive_status(client->file_descriptor, &status) && status == STATUS_OK;
}
static void pipeline_cancel(PipelineContextSender* context) { static void pipeline_cancel(PipelineContextSender* context) {
mtx_lock(&context->mutex_scanner); mtx_lock(&context->mutex_scanner);
mtx_lock(&context->mutex_loader); mtx_lock(&context->mutex_loader);
@@ -575,7 +602,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
if (send_delete_manifest(client->file_descriptor, context->manifest) != 0) if (send_delete_manifest(client->file_descriptor, context->manifest) != 0)
goto send_fail; goto send_fail;
} }
bool ok = finalize_transfer(client); bool ok = finalize_transfer(client, context->config, context->remove_source_files);
if (ok) if (ok)
remove_transferred_sources(context->config, context->remove_source_files); remove_transferred_sources(context->config, context->remove_source_files);
mtx_lock(&context->mutex_progress); mtx_lock(&context->mutex_progress);
@@ -866,7 +893,7 @@ int send_files(Config* config) {
array_list_delete(manifest); array_list_delete(manifest);
manifest = NULL; manifest = NULL;
} }
bool ok = finalize_transfer(client); bool ok = finalize_transfer(client, config, remove_sources);
if (ok) if (ok)
remove_transferred_sources(config, remove_sources); remove_transferred_sources(config, remove_sources);
if (config->show_progress && !config->quiet) if (config->show_progress && !config->quiet)
+85 -10
View File
@@ -1,6 +1,7 @@
#include "receiver.h" #include "receiver.h"
#include "chunk.h" #include "chunk.h"
#include "file_receive.h"
#include "log.h" #include "log.h"
#include "metadata.h" #include "metadata.h"
#include "protocol.h" #include "protocol.h"
@@ -8,6 +9,48 @@
#include <stdlib.h> #include <stdlib.h>
#include <sys/stat.h> #include <sys/stat.h>
bool receiver_outcomes_append(ReceiverOutcomes* outcomes, unsigned char code) {
if (!outcomes)
return false;
if (outcomes->count == outcomes->capacity) {
size_t new_capacity = outcomes->capacity == 0 ? 64 : outcomes->capacity * 2;
if (new_capacity < outcomes->capacity)
return false;
unsigned char* grown = realloc(outcomes->entries, new_capacity);
if (!grown)
return false;
outcomes->entries = grown;
outcomes->capacity = new_capacity;
}
outcomes->entries[outcomes->count++] = code;
return true;
}
void receiver_outcomes_destroy(ReceiverOutcomes* outcomes) {
if (!outcomes)
return;
free(outcomes->entries);
outcomes->entries = NULL;
outcomes->count = 0;
outcomes->capacity = 0;
}
/* End-of-transfer success frame. When --remove-source-files was negotiated
each processed data file is acknowledged first (STATUS_NEXT = written,
STATUS_OK = skipped) so the sender never removes a source the receiver did
not actually store. The frame always ends with a plain STATUS_OK. */
bool receiver_send_final_success(int fd, const Config* config, const ReceiverOutcomes* outcomes) {
if (!config->remove_source_files)
return send_status(fd, STATUS_OK);
size_t count = outcomes ? outcomes->count : 0;
for (size_t i = 0; i < count; i++) {
Status per_file = outcomes->entries[i] == FILE_SAVE_WRITTEN ? STATUS_NEXT : STATUS_OK;
if (!send_status(fd, per_file))
return false;
}
return send_status(fd, STATUS_OK);
}
static bool receiver_process_chunk(Chunk* chunk, const ReceiverSink* sink) { static bool receiver_process_chunk(Chunk* chunk, const ReceiverSink* sink) {
if (!chunk || !sink || !sink->store_file) if (!chunk || !sink || !sink->store_file)
return false; return false;
@@ -52,7 +95,7 @@ static bool receiver_process_batch(Config* config, int file_descriptor) {
send_status(file_descriptor, STATUS_ERROR); send_status(file_descriptor, STATUS_ERROR);
return false; return false;
} }
if (check_size > MAX_RECEIVE_FILE_SIZE) { if (check_size > MAX_RECEIVE_WHOLE_FILE_SIZE) {
free(check_path); free(check_path);
send_status(file_descriptor, STATUS_ERROR); send_status(file_descriptor, STATUS_ERROR);
return false; return false;
@@ -130,8 +173,14 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si
log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status"); log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status");
goto receive_error; goto receive_error;
} }
if (sink->send_success && !send_status(file_descriptor, STATUS_OK)) if (sink->send_success) {
return -1; if (sink->send_success_frame) {
if (!sink->send_success_frame(file_descriptor, sink->context))
return -1;
} else if (!send_status(file_descriptor, STATUS_OK)) {
return -1;
}
}
return 0; return 0;
receive_error: receive_error:
@@ -140,15 +189,41 @@ receive_error:
return -1; return -1;
} }
static bool receiver_save_file(File* file, void* context) { /* ---- Single-threaded sink (used by receiver_receive_files) ---- */
Config* config = context;
bool success = typedef struct {
!config->save_to_disk || file_save_to_disk(config->receive_root_directory, file, config); Config* config;
ReceiverOutcomes outcomes;
} ReceiverSaveContext;
static bool receiver_save_file(File* file, void* context_pointer) {
ReceiverSaveContext* context = context_pointer;
FileSaveResult result = FILE_SAVE_ERROR;
if (!context->config->save_to_disk) {
/* Nothing is stored; report the file as not-written so a
--remove-source-files sender keeps its source. */
result = FILE_SAVE_SKIPPED;
} else {
result = file_save_to_disk_full(context->config->receive_root_directory, file, context->config);
}
if (result != FILE_SAVE_ERROR && context->config->remove_source_files &&
!receiver_outcomes_append(&context->outcomes, (unsigned char)result)) {
file_destroy(file);
return false;
}
file_destroy(file); file_destroy(file);
return success; return result != FILE_SAVE_ERROR;
}
static bool receiver_send_success_frame(int fd, void* context_pointer) {
ReceiverSaveContext* context = context_pointer;
return receiver_send_final_success(fd, context->config, &context->outcomes);
} }
int receiver_receive_files(Config* config, int file_descriptor) { int receiver_receive_files(Config* config, int file_descriptor) {
ReceiverSink sink = {receiver_save_file, config, true, true}; ReceiverSaveContext context = {.config = config, .outcomes = {0}};
return receiver_process(config, file_descriptor, &sink); ReceiverSink sink = {receiver_save_file, &context, true, true, receiver_send_success_frame};
int ret = receiver_process(config, file_descriptor, &sink);
receiver_outcomes_destroy(&context.outcomes);
return ret;
} }
+21
View File
@@ -3,16 +3,37 @@
#include "config.h" #include "config.h"
#include "file.h" #include "file.h"
#include "file_receive.h"
typedef bool (*ReceiverFileSink)(File* file, void* context); typedef bool (*ReceiverFileSink)(File* file, void* context);
/* Ordered per-file save outcomes for one connection. One entry is appended
for every data-bearing file the receiver processes (in the order the files
were sent) so the sender of a --remove-source-files transfer can be told
which sources were actually written versus skipped on the receiver. */
typedef struct {
unsigned char* entries; /* FILE_SAVE_WRITTEN or FILE_SAVE_SKIPPED */
size_t count;
size_t capacity;
} ReceiverOutcomes;
typedef bool (*ReceiverSuccessFrame)(int fd, void* context);
typedef struct { typedef struct {
ReceiverFileSink store_file; ReceiverFileSink store_file;
void* context; void* context;
bool send_error; bool send_error;
bool send_success; bool send_success;
/* Emits the end-of-transfer success frame. When the sender requested
--remove-source-files this includes one per-file status per processed
data file followed by the final STATUS_OK; otherwise just STATUS_OK. */
ReceiverSuccessFrame send_success_frame;
} ReceiverSink; } ReceiverSink;
bool receiver_outcomes_append(ReceiverOutcomes* outcomes, unsigned char code);
void receiver_outcomes_destroy(ReceiverOutcomes* outcomes);
bool receiver_send_final_success(int fd, const Config* config, const ReceiverOutcomes* outcomes);
int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink); int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink);
int receiver_receive_files(Config* config, int file_descriptor); int receiver_receive_files(Config* config, int file_descriptor);
+18 -141
View File
@@ -1,14 +1,12 @@
#include "config.h" #include "config.h"
#include "chunk.h"
#include "file.h" #include "file.h"
#include "log.h" #include "log.h"
#include "metadata.h"
#include "multiprocessing.h" #include "multiprocessing.h"
#include "protocol.h"
#include "queue.h" #include "queue.h"
#include "receiver.h" #include "receiver.h"
#include "transport_tcp.h" #include "transport_tcp.h"
#include "transport_tls.h" #include "transport_tls.h"
#include "unistd.h"
#include "utils.h" #include "utils.h"
#include <fcntl.h> #include <fcntl.h>
#include <limits.h> #include <limits.h>
@@ -16,7 +14,6 @@
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <sys/stat.h>
#include <unistd.h> #include <unistd.h>
#include <openssl/x509.h> #include <openssl/x509.h>
@@ -26,6 +23,13 @@ static bool allow_delete;
static bool allow_unauthenticated; static bool allow_unauthenticated;
static const char* required_client_cn; static const char* required_client_cn;
/* Aggregate payload bytes the multithreaded receiver may buffer ahead of the
slow disk writer. Receiving one more chunk adds up to ~2 * MAX_CHUNK_SIZE
of transient wire/decompression buffers on top of the queued payloads, so
this ceiling keeps total per-connection receive memory (decompressed and
per-file copied chunk buffers included) within MAX_CONNECTION_MEMORY. */
#define RECEIVER_QUEUE_MAX_BYTES (MAX_CONNECTION_MEMORY - 2 * MAX_CHUNK_SIZE)
static bool tls_client_identity_allowed(SSL* ssl) { static bool tls_client_identity_allowed(SSL* ssl) {
if (!ssl || !required_client_cn) if (!ssl || !required_client_cn)
return false; return false;
@@ -101,140 +105,6 @@ static bool __attribute__((unused)) configure_authorization(const char* root) {
return true; return true;
} }
int receive_files(Config* config, int fd) {
Status status;
if (!receive_status(fd, &status))
return -1;
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH) {
if (status == STATUS_KEEPALIVE) {
send_status(fd, STATUS_KEEPALIVE);
goto next;
}
if (status == STATUS_ABORT) {
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
return -1;
}
if (status == STATUS_CHECK) {
bool skipped;
File* file = receive_incremental_check(fd, config, &skipped);
if (skipped)
goto next;
if (file == NULL && !skipped)
return -1;
if (config->save_to_disk &&
!file_save_to_disk(config->receive_root_directory, file, config)) {
file_destroy(file);
send_status(fd, STATUS_ERROR);
return -1;
}
file_destroy(file);
} else if (status == STATUS_CHUNK) {
Chunk* chunk = receive_chunk_data(fd, config);
if (chunk == NULL) {
send_status(fd, STATUS_ERROR);
return -1;
}
for (int i = 0; i < chunk->element_count; i++) {
if (config->save_to_disk &&
!file_save_to_disk(config->receive_root_directory, chunk->items[i], config)) {
chunk_destroy(chunk);
send_status(fd, STATUS_ERROR);
return -1;
}
}
chunk_destroy(chunk);
} else if (status == STATUS_CHECK_BATCH) {
int count;
/* Batch framing has no checksum field yet; never silently downgrade a
checksum-enabled transfer into mtime-only matching. */
if (config->checksum || !receive_int(fd, &count) || count < 0 || count > MAX_MANIFEST_ENTRIES)
return -1;
for (int i = 0; i < count; i++) {
char* check_path = receive_str(fd);
if (!check_path)
return -1;
unsigned long long check_size;
long long check_mtime;
if (!receive_n_data(fd, &check_size, sizeof(check_size)) ||
!receive_n_data(fd, &check_mtime, sizeof(check_mtime))) {
free(check_path);
return -1;
}
long long check_mtime_nsec;
if (!receive_n_data(fd, &check_mtime_nsec, sizeof(check_mtime_nsec)) ||
check_mtime_nsec < 0 || check_mtime_nsec >= 1000000000LL) {
free(check_path);
send_status(fd, STATUS_ERROR);
return -1;
}
if (!utils_valid_batch_path(check_path)) {
free(check_path);
send_status(fd, STATUS_ERROR);
return -1;
}
struct stat st;
char* full_path = path_cat(config->receive_root_directory, check_path);
if (!full_path) {
free(check_path);
send_status(fd, STATUS_ERROR);
return -1;
}
bool has_old = full_path && file_stat_secure(full_path, &st);
long long old_mtime_nsec = 0;
if (has_old) {
#ifdef __linux__
old_mtime_nsec = st.st_mtim.tv_nsec;
#endif
}
bool match = !config->ignore_times && has_old &&
(unsigned long long)st.st_size == check_size &&
metadata_mtime_matches(st.st_mtime, old_mtime_nsec, (time_t)check_mtime,
(long)check_mtime_nsec, config->modify_window);
bool sent = send_status(fd, match ? STATUS_OK : STATUS_NEXT);
free(full_path);
free(check_path);
if (!sent)
return -1;
}
goto next;
} else {
File* file = file_receive(config, fd);
if (file == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to receive file");
send_status(fd, STATUS_ERROR);
return -1;
}
if (config->save_to_disk &&
!file_save_to_disk(config->receive_root_directory, file, config)) {
file_destroy(file);
send_status(fd, STATUS_ERROR);
return -1;
}
file_destroy(file);
}
next:
if (!receive_status(fd, &status)) {
send_status(fd, STATUS_ERROR);
return -1;
}
}
if (status == STATUS_MANIFEST) {
if (receive_manifest(fd, config, &status) != 0) {
return -1;
}
}
if (status != STATUS_FINISHED) {
log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status");
send_status(fd, STATUS_ERROR);
return -1;
}
send_status(fd, STATUS_OK);
return 0;
}
void handler(int file_descriptor) { void handler(int file_descriptor) {
SSL* ssl = io_get_ssl(); SSL* ssl = io_get_ssl();
ProtocolSession session; ProtocolSession session;
@@ -314,6 +184,7 @@ void handler(int file_descriptor) {
protocol_session_set_max_alloc(&context->session, config->max_alloc); protocol_session_set_max_alloc(&context->session, config->max_alloc);
atomic_store(&context->session.total_allocated_bytes, atomic_store(&context->session.total_allocated_bytes,
atomic_load(&session.total_allocated_bytes)); atomic_load(&session.total_allocated_bytes));
pipeline_context_receiver_set_queue_byte_limit(context, RECEIVER_QUEUE_MAX_BYTES);
thrd_t receiver, writer; thrd_t receiver, writer;
bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success; bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success;
bool writer_created = false; bool writer_created = false;
@@ -342,9 +213,15 @@ void handler(int file_descriptor) {
int writer_result; int writer_result;
thrd_join(receiver, &receiver_result); thrd_join(receiver, &receiver_result);
thrd_join(writer, &writer_result); thrd_join(writer, &writer_result);
send_status(file_descriptor, receiver_result == thrd_success && writer_result == thrd_success bool transfer_ok = receiver_result == thrd_success && writer_result == thrd_success;
? STATUS_OK if (transfer_ok) {
: STATUS_ERROR); if (!receiver_send_final_success(file_descriptor, config, &context->outcomes))
transfer_ok = false;
} else {
send_status(file_descriptor, STATUS_ERROR);
}
if (!transfer_ok)
log_message(LOG_LEVEL_ERROR, "Transfer failed");
pipeline_context_receiver_destroy(context); pipeline_context_receiver_destroy(context);
} else { } else {
if (receiver_receive_files(config, file_descriptor) != 0) if (receiver_receive_files(config, file_descriptor) != 0)
+49 -14
View File
@@ -139,8 +139,9 @@ static bool validate_received_config(const Config* config) {
valid_wire_bool(config->use_delete) && valid_wire_bool(config->use_incremental) && valid_wire_bool(config->use_delete) && valid_wire_bool(config->use_incremental) &&
valid_wire_bool(config->size_only) && valid_wire_bool(config->ignore_times) && valid_wire_bool(config->size_only) && valid_wire_bool(config->ignore_times) &&
valid_wire_bool(config->use_delta) && valid_wire_bool(config->backup) && valid_wire_bool(config->use_delta) && valid_wire_bool(config->backup) &&
valid_wire_bool(config->follow_symlinks) && valid_wire_bool(config->copy_links) && valid_wire_bool(config->remove_source_files) && valid_wire_bool(config->follow_symlinks) &&
valid_wire_bool(config->safe_links) && valid_wire_bool(config->copy_unsafe_links) && valid_wire_bool(config->copy_links) && valid_wire_bool(config->safe_links) &&
valid_wire_bool(config->copy_unsafe_links) &&
valid_wire_bool(config->preserve_hard_links) && valid_wire_bool(config->preserve_acls) && valid_wire_bool(config->preserve_hard_links) && valid_wire_bool(config->preserve_acls) &&
valid_wire_bool(config->preserve_xattrs) && valid_wire_bool(config->preserve_devices) && valid_wire_bool(config->preserve_xattrs) && valid_wire_bool(config->preserve_devices) &&
valid_wire_bool(config->preserve_sparse) && valid_wire_bool(config->ignore_existing) && valid_wire_bool(config->preserve_sparse) && valid_wire_bool(config->ignore_existing) &&
@@ -274,11 +275,11 @@ static bool send_delta_fields(int fd, const Config* c) {
static bool send_file_options(int fd, const Config* c) { static bool send_file_options(int fd, const Config* c) {
return send_int(fd, c->backup) && send_str(fd, c->backup_dir ? c->backup_dir : "") && return send_int(fd, c->backup) && send_str(fd, c->backup_dir ? c->backup_dir : "") &&
send_int(fd, c->follow_symlinks) && send_int(fd, c->copy_links) && send_int(fd, c->remove_source_files) && send_int(fd, c->follow_symlinks) &&
send_int(fd, c->safe_links) && send_int(fd, c->copy_unsafe_links) && send_int(fd, c->copy_links) && send_int(fd, c->safe_links) &&
send_int(fd, c->preserve_hard_links) && send_int(fd, c->preserve_acls) && send_int(fd, c->copy_unsafe_links) && send_int(fd, c->preserve_hard_links) &&
send_int(fd, c->preserve_xattrs) && send_int(fd, c->preserve_devices) && send_int(fd, c->preserve_acls) && send_int(fd, c->preserve_xattrs) &&
send_int(fd, c->preserve_sparse); send_int(fd, c->preserve_devices) && send_int(fd, c->preserve_sparse);
} }
static bool send_selection_options(int fd, const Config* c) { static bool send_selection_options(int fd, const Config* c) {
@@ -355,8 +356,17 @@ static bool receive_delta_fields(int fd, Config* c) {
static bool receive_file_options(int fd, Config* c) { static bool receive_file_options(int fd, Config* c) {
if (!receive_wire_bool(fd, &c->backup)) if (!receive_wire_bool(fd, &c->backup))
return false; return false;
c->backup_dir = receive_str(fd); char* backup_dir = receive_str(fd);
if (!c->backup_dir) if (!backup_dir)
return false;
if (*backup_dir != '\0') {
c->backup_dir = backup_dir;
} else {
/* The sender serializes an unset (NULL) string as "", so canonicalize the
empty wire value back to NULL to preserve NULL-vs-empty semantics. */
free(backup_dir);
}
if (!receive_wire_bool(fd, &c->remove_source_files))
return false; return false;
bool* flags[] = {&c->follow_symlinks, &c->copy_links, &c->safe_links, bool* flags[] = {&c->follow_symlinks, &c->copy_links, &c->safe_links,
&c->copy_unsafe_links, &c->preserve_hard_links, &c->preserve_acls, &c->copy_unsafe_links, &c->preserve_hard_links, &c->preserve_acls,
@@ -386,12 +396,37 @@ static bool receive_selection_options(int fd, Config* c) {
} }
static bool receive_resume_options(int fd, Config* c) { static bool receive_resume_options(int fd, Config* c) {
c->temp_dir = receive_str(fd); char* temp_dir = receive_str(fd);
if (!c->temp_dir || !receive_wire_bool(fd, &c->partial)) if (!temp_dir)
return false; return false;
c->partial_dir = receive_str(fd); if (*temp_dir != '\0') {
c->suffix = c->partial_dir ? receive_str(fd) : NULL; c->temp_dir = temp_dir;
if (!c->partial_dir || !c->suffix || !receive_wire_bool(fd, &c->delete_before)) } else {
free(temp_dir);
}
if (!receive_wire_bool(fd, &c->partial))
return false;
/* These options have NULL client defaults, so the sender transmits an empty
string for "unset". Canonicalize the empty wire value back to NULL so
receivers observe exactly what the client configured (plain --backup, for
example, must not look like --backup-dir ""). */
char* partial_dir = receive_str(fd);
if (!partial_dir)
return false;
if (*partial_dir != '\0') {
c->partial_dir = partial_dir;
} else {
free(partial_dir);
}
char* suffix = receive_str(fd);
if (!suffix)
return false;
if (*suffix != '\0') {
c->suffix = suffix;
} else {
free(suffix);
}
if (!receive_wire_bool(fd, &c->delete_before))
return false; return false;
if (!receive_wire_bool(fd, &c->checksum)) if (!receive_wire_bool(fd, &c->checksum))
return false; return false;
+1 -1
View File
@@ -147,7 +147,7 @@ typedef struct Config {
bool skip_compress_set; bool skip_compress_set;
} Config; } Config;
#define PROTOCOL_VERSION "2.4.0" #define PROTOCOL_VERSION "2.5.0"
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
Config* config_create(void); Config* config_create(void);
+158 -26
View File
@@ -216,6 +216,119 @@ static void free_instructions(DeltaInstruction* instrs, uint32_t count) {
free(instrs); free(instrs);
} }
/* Sentinel meaning "no signature block" in the lookup index chains. Block
* counts are bounded well below UINT32_MAX, so it doubles as a null link. */
#define DELTA_NO_BLOCK UINT32_MAX
/* Avalanche mix for the rolling checksum so blocks do not cluster in the
* bucket table when the weak checksum has little entropy (e.g. all-zero or
* patterned files). */
static uint32_t delta_adler_mix(uint32_t h) {
h ^= h >> 16;
h *= 0x7feb352dU;
h ^= h >> 15;
h *= 0x846ca68bU;
h ^= h >> 16;
return h;
}
/* Smallest power of two >= v. v must be non-zero. */
static uint32_t delta_next_pow2(uint32_t v) {
v--;
v |= v >> 1;
v |= v >> 2;
v |= v >> 4;
v |= v >> 8;
v |= v >> 16;
return v + 1;
}
/* Build a hash index over sig->blocks keyed by the (mixed) rolling checksum.
* All blocks sharing an Adler-32 value land in the same bucket; collisions
* are chained through a single contiguous allocation:
*
* [0, bucket_count) heads (first block per bucket)
* [bucket_count, 2*bucket_count) tails (last block per bucket)
* [2*bucket_count, ...) per-block chain links
*
* Blocks are inserted in ascending index order so every bucket chain is
* ordered exactly like the historical linear scan. Returns the base pointer
* (also the heads array) or NULL when no index could be allocated; callers
* then fall back to the linear scan. */
static uint32_t* delta_build_index(const DeltaSignature* sig, uint32_t bucket_count) {
if (sig->block_count == 0 || bucket_count == 0)
return NULL;
size_t entries = (size_t)2 * bucket_count + sig->block_count;
if (entries > SIZE_MAX / sizeof(uint32_t))
return NULL;
uint32_t* index = protocol_alloc(entries * sizeof(uint32_t));
if (!index)
return NULL;
uint32_t* heads = index;
uint32_t* tails = index + bucket_count;
uint32_t* next = index + 2 * bucket_count;
uint32_t mask = bucket_count - 1;
memset(heads, 0xFF, (size_t)bucket_count * sizeof(uint32_t));
memset(tails, 0xFF, (size_t)bucket_count * sizeof(uint32_t));
for (uint32_t j = 0; j < sig->block_count; j++) {
uint32_t b = delta_adler_mix(sig->blocks[j].adler32) & mask;
if (heads[b] == DELTA_NO_BLOCK)
heads[b] = j;
else
next[tails[b]] = j;
tails[b] = j;
next[j] = DELTA_NO_BLOCK;
}
return index;
}
/* Locate the signature block matching the byte window at new_data[i].
*
* Mirrors the original per-window behaviour exactly: only a full block_size
* window can match, candidates are accepted only when the weak (Adler-32) and
* strong (xxHash32) checksums both agree, and the lowest block index wins so
* the emitted op stream is byte-identical to the linear scan. When heads is
* non-NULL the candidate set is reached through the bucket index (expected
* O(1) per window); otherwise an exact linear scan is used. */
static uint32_t delta_find_match(const uint8_t* window, uint32_t window_len, uint32_t adler,
bool full_window, const DeltaSignature* sig, const uint32_t* heads,
const uint32_t* next, uint32_t mask) {
if (!full_window || sig->block_count == 0)
return DELTA_NO_BLOCK;
if (heads) {
uint32_t b = delta_adler_mix(adler) & mask;
uint32_t window_xxh = 0;
bool have_xxh = false;
for (uint32_t j = heads[b]; j != DELTA_NO_BLOCK; j = next[j]) {
if (sig->blocks[j].adler32 != adler)
continue;
if (!have_xxh) {
window_xxh = delta_xxhash32(window, window_len);
have_xxh = true;
}
if (window_xxh == sig->blocks[j].xxhash)
return j;
}
return DELTA_NO_BLOCK;
}
/* Fallback used when the index could not be allocated. */
for (uint32_t j = 0; j < sig->block_count; j++) {
if (sig->blocks[j].adler32 == adler) {
uint32_t window_xxh = delta_xxhash32(window, window_len);
if (window_xxh == sig->blocks[j].xxhash)
return j;
}
}
return DELTA_NO_BLOCK;
}
Delta* delta_compute(const void* new_file_data, uint64_t new_file_size, const DeltaSignature* sig, Delta* delta_compute(const void* new_file_data, uint64_t new_file_size, const DeltaSignature* sig,
uint32_t block_size) { uint32_t block_size) {
if (!new_file_data || !sig || !sig->blocks || new_file_size == 0 || block_size == 0 || if (!new_file_data || !sig || !sig->blocks || new_file_size == 0 || block_size == 0 ||
@@ -230,6 +343,25 @@ Delta* delta_compute(const void* new_file_data, uint64_t new_file_size, const De
if (!instrs) if (!instrs)
return NULL; return NULL;
/* Build a one-time bucket index over the signature blocks keyed by the weak
* checksum. This turns the per-byte-window candidate lookup from an
* O(block_count) linear scan into an expected O(1) probe, which dominates
* the cost for large mostly-matching files (the diff steps one byte at a
* time through changed regions). On allocation failure the probe falls back
* to the original linear scan, so behaviour is unchanged under memory
* pressure. */
uint32_t* index = NULL;
const uint32_t* chain_next = NULL;
uint32_t mask = 0;
if (sig->block_count > 0) {
uint32_t bucket_count = delta_next_pow2(sig->block_count);
index = delta_build_index(sig, bucket_count);
if (index) {
chain_next = index + 2 * bucket_count;
mask = bucket_count - 1;
}
}
uint64_t literal_start = 0; uint64_t literal_start = 0;
bool has_literal = false; bool has_literal = false;
@@ -264,34 +396,32 @@ Delta* delta_compute(const void* new_file_data, uint64_t new_file_size, const De
} }
bool matched = false; bool matched = false;
for (uint32_t j = 0; j < sig->block_count; j++) { uint32_t match_block = delta_find_match(new_data + i, window_len, adler, full_window, sig,
if (adler == sig->blocks[j].adler32 && full_window) { index, chain_next, mask);
uint32_t xxh = delta_xxhash32(new_data + i, window_len); if (match_block != DELTA_NO_BLOCK) {
if (xxh == sig->blocks[j].xxhash) { if (has_literal) {
if (has_literal) { if (!flush_literal(&instrs, &capacity, &count, new_data, literal_start, i)) {
if (!flush_literal(&instrs, &capacity, &count, new_data, literal_start, i)) { free_instructions(instrs, count);
free_instructions(instrs, count); free(index);
return NULL; return NULL;
}
has_literal = false;
}
if (!ensure_capacity(&instrs, &capacity, count)) {
free_instructions(instrs, count);
return NULL;
}
instrs[count].type = DELTA_INSTR_BLOCK_MATCH;
instrs[count].match.block_index = j;
instrs[count].match.block_offset = 0;
instrs[count].match.length = window_len;
count++;
i += window_len;
rolling_valid = false;
matched = true;
break;
} }
has_literal = false;
} }
if (!ensure_capacity(&instrs, &capacity, count)) {
free_instructions(instrs, count);
free(index);
return NULL;
}
instrs[count].type = DELTA_INSTR_BLOCK_MATCH;
instrs[count].match.block_index = match_block;
instrs[count].match.block_offset = 0;
instrs[count].match.length = window_len;
count++;
i += window_len;
rolling_valid = false;
matched = true;
} }
if (!matched) { if (!matched) {
@@ -303,6 +433,8 @@ Delta* delta_compute(const void* new_file_data, uint64_t new_file_size, const De
} }
} }
free(index);
if (has_literal) { if (has_literal) {
if (!flush_literal(&instrs, &capacity, &count, new_data, literal_start, new_file_size)) { if (!flush_literal(&instrs, &capacity, &count, new_data, literal_start, new_file_size)) {
free_instructions(instrs, count); free_instructions(instrs, count);
+18 -3
View File
@@ -348,10 +348,25 @@ static bool file_to_disk_secure_impl(const char* path, const void* data,
if (newer) { if (newer) {
ok = true; ok = true;
} else { } else {
if (!sparse || data_size == 0 || ftruncate(fd, (off_t)data_size) == 0) /* In-place overwrites: pre-size sparse targets and always trim the
file to the new payload length afterwards so shorter payloads can
never leave stale trailing bytes from a previous version. */
if (sparse && data_size > 0)
ok = ftruncate(fd, (off_t)data_size) == 0;
if (ok || !sparse || data_size == 0)
ok = write_all(fd, data, data_size); ok = write_all(fd, data, data_size);
if (ok && metadata) if (ok)
ok = file_restore_metadata_fd(fd, metadata, preserve_executability); ok = ftruncate(fd, (off_t)data_size) == 0;
/* Normalize the mode: apply the metadata-derived safe mode when the
sender supplied metadata (setuid/setgid/sticky are never honored);
otherwise fall back to a safe default so dangerous bits on an
existing destination cannot survive an overwrite. */
if (ok) {
if (metadata)
ok = file_restore_metadata_fd(fd, metadata, preserve_executability);
else if (fchmod(fd, S_IRUSR | S_IWUSR | S_IRGRP | S_IROTH) != 0)
ok = false;
}
if (ok && use_fsync) if (ok && use_fsync)
ok = fsync(fd) == 0; ok = fsync(fd) == 0;
} }
+113 -68
View File
@@ -20,9 +20,14 @@
#include "utils.h" #include "utils.h"
#define MAX_SERVER_DELETE_COUNT 100000U #define MAX_SERVER_DELETE_COUNT 100000U
#define MAX_FILE_DATA_SIZE MAX_RECEIVE_FILE_SIZE #define MAX_FILE_DATA_SIZE MAX_RECEIVE_WHOLE_FILE_SIZE
bool file_save_to_disk(const char* root_directory, const File* file, const Config* config) { bool file_save_to_disk(const char* root_directory, const File* file, const Config* config) {
return file_save_to_disk_full(root_directory, file, config) != FILE_SAVE_ERROR;
}
FileSaveResult file_save_to_disk_full(const char* root_directory, const File* file,
const Config* config) {
/* Backups are incompatible with ignore-existing: moving the entry first /* Backups are incompatible with ignore-existing: moving the entry first
would make a concurrent no-replace commit overwrite its old name. */ would make a concurrent no-replace commit overwrite its old name. */
bool backup_enabled = config && config->backup && !config->ignore_existing; bool backup_enabled = config && config->backup && !config->ignore_existing;
@@ -32,6 +37,7 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi
const char* backup_suffix = (config && config->suffix) ? config->suffix : "~"; const char* backup_suffix = (config && config->suffix) ? config->suffix : "~";
const char* backup_dir = (config && config->backup_dir) ? config->backup_dir : NULL; const char* backup_dir = (config && config->backup_dir) ? config->backup_dir : NULL;
const char* partial_dir = (config && config->partial_dir) ? config->partial_dir : NULL; const char* partial_dir = (config && config->partial_dir) ? config->partial_dir : NULL;
bool use_partial_root = partial_dir && config && config->partial;
char *confined_backup = NULL, *confined_partial = NULL, *disk_path = NULL; char *confined_backup = NULL, *confined_partial = NULL, *disk_path = NULL;
char* destination_path = NULL; char* destination_path = NULL;
char *backup_path = NULL, *parent_copy = NULL; char *backup_path = NULL, *parent_copy = NULL;
@@ -42,23 +48,22 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi
(!backup_suffix || backup_suffix[0] == '\0' || strchr(backup_suffix, '/') != NULL || (!backup_suffix || backup_suffix[0] == '\0' || strchr(backup_suffix, '/') != NULL ||
strcmp(backup_suffix, ".") == 0 || strcmp(backup_suffix, "..") == 0))) { strcmp(backup_suffix, ".") == 0 || strcmp(backup_suffix, "..") == 0))) {
log_message(LOG_LEVEL_ERROR, "Invalid file or path received"); log_message(LOG_LEVEL_ERROR, "Invalid file or path received");
return false; return FILE_SAVE_ERROR;
} }
/* These options arrive from the client. They are names below the server /* These options arrive from the client. They are names below the server
root, never independent filesystem roots. */ root, never independent filesystem roots. */
if ((backup_dir && (backup_dir[0] == '/' || has_path_traversal(backup_dir))) || if ((backup_dir && (backup_dir[0] == '/' || has_path_traversal(backup_dir))) ||
(partial_dir && (partial_dir[0] == '/' || has_path_traversal(partial_dir)))) (partial_dir && (partial_dir[0] == '/' || has_path_traversal(partial_dir))))
return false; return FILE_SAVE_ERROR;
if (backup_dir && !(confined_backup = path_cat(root_directory, backup_dir))) if (backup_dir && !(confined_backup = path_cat(root_directory, backup_dir)))
return false; return FILE_SAVE_ERROR;
if (partial_dir && !(confined_partial = path_cat(root_directory, partial_dir))) { if (partial_dir && !(confined_partial = path_cat(root_directory, partial_dir))) {
free(confined_backup); free(confined_backup);
return false; return FILE_SAVE_ERROR;
} }
const char* actual_root = const char* actual_root = use_partial_root ? confined_partial : root_directory;
(partial_dir && config && config->partial) ? confined_partial : root_directory;
destination_path = path_cat(root_directory, file->path); destination_path = path_cat(root_directory, file->path);
disk_path = path_cat(actual_root, file->path); disk_path = path_cat(actual_root, file->path);
if (destination_path == NULL || disk_path == NULL) { if (destination_path == NULL || disk_path == NULL) {
@@ -66,7 +71,7 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi
free(confined_partial); free(confined_partial);
free(destination_path); free(destination_path);
free(disk_path); free(disk_path);
return false; return FILE_SAVE_ERROR;
} }
/* --existing checks the final destination, not a temporary partial path. */ /* --existing checks the final destination, not a temporary partial path. */
@@ -75,7 +80,7 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi
free(confined_partial); free(confined_partial);
free(destination_path); free(destination_path);
free(disk_path); free(disk_path);
return true; return FILE_SAVE_SKIPPED;
} }
/* --ignore-existing checks the final destination before partial files or /* --ignore-existing checks the final destination before partial files or
@@ -87,34 +92,40 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi
free(confined_partial); free(confined_partial);
free(destination_path); free(destination_path);
free(disk_path); free(disk_path);
return true; return FILE_SAVE_SKIPPED;
} }
} }
free(destination_path);
destination_path = NULL;
/* --update is receiver-side policy: never replace a newer destination. /* --update is receiver-side policy: never replace a newer destination.
The secure stat does not require read permission on the destination. */ In partial-dir mode the entry that would be replaced is the real
if (config && config->update && file_destination_is_newer_secure(disk_path, file->metadata)) { destination, not the temporary partial file. The secure stat does not
require read permission on the destination. */
const char* update_target = use_partial_root ? destination_path : disk_path;
if (config && config->update && file_destination_is_newer_secure(update_target, file->metadata)) {
free(confined_backup); free(confined_backup);
free(confined_partial); free(confined_partial);
free(destination_path);
free(disk_path); free(disk_path);
return true; return FILE_SAVE_SKIPPED;
} }
if (backup_enabled) { if (backup_enabled) {
/* Back up the entry that the incoming write will replace. When writing
through a partial dir the pre-existing destination file is the one to
preserve; any stale partial file is overwritten without a backup. */
const char* replace_target = use_partial_root ? destination_path : disk_path;
struct stat backup_stat; struct stat backup_stat;
if (file_stat_secure(disk_path, &backup_stat)) { if (file_stat_secure(replace_target, &backup_stat)) {
if (backup_dir) { if (backup_dir) {
backup_path = path_cat(confined_backup, file->path); backup_path = path_cat(confined_backup, file->path);
} else { } else {
size_t path_len = strlen(disk_path); size_t path_len = strlen(replace_target);
size_t suffix_len = strlen(backup_suffix); size_t suffix_len = strlen(backup_suffix);
if (path_len > SIZE_MAX - suffix_len - 1) if (path_len > SIZE_MAX - suffix_len - 1)
goto fail; goto fail;
backup_path = malloc(path_len + suffix_len + 1); backup_path = malloc(path_len + suffix_len + 1);
if (backup_path) { if (backup_path) {
memcpy(backup_path, disk_path, path_len); memcpy(backup_path, replace_target, path_len);
memcpy(backup_path + path_len, backup_suffix, suffix_len + 1); memcpy(backup_path + path_len, backup_suffix, suffix_len + 1);
} }
} }
@@ -125,7 +136,7 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi
goto fail; goto fail;
free(parent_copy); free(parent_copy);
parent_copy = NULL; parent_copy = NULL;
if (!file_rename_secure(disk_path, backup_path)) if (!file_rename_secure(replace_target, backup_path))
goto fail; goto fail;
free(backup_path); free(backup_path);
backup_path = NULL; backup_path = NULL;
@@ -149,13 +160,25 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi
: file_to_disk_secure_with_fsync(disk_path, file->data->data, file->data->size, : file_to_disk_secure_with_fsync(disk_path, file->data->data, file->data->size,
inplace, sparse, metadata, preserve_executability, inplace, sparse, metadata, preserve_executability,
config && config->use_fsync); config && config->use_fsync);
if (!ok)
goto fail;
/* --partial --partial-dir writes the complete file under the partial dir so
interrupted transfers leave a resumable copy there. Once the file is
fully written it must be atomically installed at the real destination;
otherwise completed transfers would linger under the partial dir. */
if (use_partial_root) {
if (!file_rename_secure(disk_path, destination_path))
goto fail;
}
free(parent_copy); free(parent_copy);
free(backup_path); free(backup_path);
free(confined_backup); free(confined_backup);
free(confined_partial); free(confined_partial);
free(destination_path); free(destination_path);
free(disk_path); free(disk_path);
return ok; return FILE_SAVE_WRITTEN;
fail: fail:
free(parent_copy); free(parent_copy);
@@ -164,13 +187,15 @@ fail:
free(confined_partial); free(confined_partial);
free(destination_path); free(destination_path);
free(disk_path); free(disk_path);
return false; return FILE_SAVE_ERROR;
} }
static File* receive_delta_file(int fd, const Config* config, const char* check_path, static File* receive_delta_file(int fd, const Config* config, const char* check_path,
void* old_data, unsigned long long old_size, bool* failed) { void* old_data, unsigned long long old_size, bool* failed) {
if (!old_data) if (!old_data) {
*failed = true;
return NULL; return NULL;
}
DeltaSignature* sig = delta_signature_create(old_data, old_size, config->delta_block_size); DeltaSignature* sig = delta_signature_create(old_data, old_size, config->delta_block_size);
if (!sig) { if (!sig) {
@@ -206,7 +231,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
} }
if (resp == STATUS_DELTA_DATA) { if (resp == STATUS_DELTA_DATA) {
Data* delta_data = receive_data_limited(fd, MAX_RECEIVE_FILE_SIZE); Data* delta_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE);
if (!delta_data) { if (!delta_data) {
delta_signature_destroy(sig); delta_signature_destroy(sig);
free(old_data); free(old_data);
@@ -219,7 +244,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
!compression_should_skip_with_suffixes( !compression_should_skip_with_suffixes(
check_path, config->skip_compress_suffixes, check_path, config->skip_compress_suffixes,
config->skip_compress_set ? config->skip_compress_count : -1)) { config->skip_compress_set ? config->skip_compress_count : -1)) {
raw_delta = data_decompress_limited(delta_data, MAX_RECEIVE_FILE_SIZE); raw_delta = data_decompress_limited(delta_data, MAX_RECEIVE_WHOLE_FILE_SIZE);
data_destroy(delta_data); data_destroy(delta_data);
if (!raw_delta) { if (!raw_delta) {
free(old_data); free(old_data);
@@ -239,11 +264,12 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
} }
uint64_t new_size = delta->new_file_size; uint64_t new_size = delta->new_file_size;
if (new_size > MAX_RECEIVE_FILE_SIZE || new_size > SIZE_MAX) { if (new_size > MAX_RECEIVE_WHOLE_FILE_SIZE || new_size > SIZE_MAX) {
delta_destroy(delta); delta_destroy(delta);
free(old_data); free(old_data);
delta_signature_destroy(sig); delta_signature_destroy(sig);
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
*failed = true;
return NULL; return NULL;
} }
void* new_data = delta_apply(old_data, old_size, delta, config->delta_block_size); void* new_data = delta_apply(old_data, old_size, delta, config->delta_block_size);
@@ -284,6 +310,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
free(old_data); free(old_data);
delta_signature_destroy(sig); delta_signature_destroy(sig);
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
*failed = true;
return NULL; return NULL;
} }
data_destroy(file->data); data_destroy(file->data);
@@ -314,7 +341,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
} }
} }
Data* file_data = receive_data_limited(fd, MAX_RECEIVE_FILE_SIZE); Data* file_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE);
if (file_data == NULL) { if (file_data == NULL) {
file_destroy(file); file_destroy(file);
*failed = true; *failed = true;
@@ -325,7 +352,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
!compression_should_skip_with_suffixes( !compression_should_skip_with_suffixes(
file->path, config->skip_compress_suffixes, file->path, config->skip_compress_suffixes,
config->skip_compress_set ? config->skip_compress_count : -1)) { config->skip_compress_set ? config->skip_compress_count : -1)) {
Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_FILE_SIZE); Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE);
data_destroy(file_data); data_destroy(file_data);
if (uncompressed == NULL) { if (uncompressed == NULL) {
file_destroy(file); file_destroy(file);
@@ -336,6 +363,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
data_destroy(uncompressed); data_destroy(uncompressed);
file_destroy(file); file_destroy(file);
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
*failed = true;
return NULL; return NULL;
} }
file_data = uncompressed; file_data = uncompressed;
@@ -384,7 +412,7 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
return NULL; return NULL;
} }
if (check_size > MAX_RECEIVE_FILE_SIZE) { if (check_size > MAX_RECEIVE_WHOLE_FILE_SIZE) {
free(check_path); free(check_path);
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
return NULL; return NULL;
@@ -405,6 +433,9 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
return NULL; return NULL;
} }
/* Open the existing destination entry (if any) once and keep the descriptor
until the quick-check below decides whether the old contents are needed. */
struct stat st; struct stat st;
bool has_old_file = false; bool has_old_file = false;
int old_fd = -1; int old_fd = -1;
@@ -416,9 +447,34 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
close(parent_fd); close(parent_fd);
has_old_file = old_fd >= 0 && fstat(old_fd, &st) == 0 && S_ISREG(st.st_mode); has_old_file = old_fd >= 0 && fstat(old_fd, &st) == 0 && S_ISREG(st.st_mode);
} }
if (!has_old_file && old_fd >= 0) {
close(old_fd);
old_fd = -1;
}
unsigned long long old_size = has_old_file ? (unsigned long long)st.st_size : 0; unsigned long long old_size = has_old_file ? (unsigned long long)st.st_size : 0;
/* Decide from metadata alone whether the receiver already holds the file
the sender is offering. The old contents are only read into memory when
a checksum comparison or a delta transfer actually requires them. */
bool size_equal = has_old_file && old_size == check_size;
bool match_by_metadata = false;
if (size_equal && !config->ignore_times && !config->size_only) {
long long old_mtime_nsec = 0;
#ifdef __linux__
old_mtime_nsec = st.st_mtim.tv_nsec;
#endif
match_by_metadata = metadata_mtime_matches(st.st_mtime, old_mtime_nsec, (time_t)check_mtime,
(long)check_mtime_nsec, config->modify_window);
}
bool try_delta = config->use_delta && !config->whole_file && has_old_file &&
delta_should_attempt(old_size, check_size, config->delta_max_file_size);
bool checksum_needs_read = size_equal && !config->ignore_times && config->checksum;
bool need_old_data = checksum_needs_read || try_delta;
void* old_data = NULL; void* old_data = NULL;
if (has_old_file && old_size > 0 && old_size <= MAX_RECEIVE_FILE_SIZE && old_size <= SIZE_MAX) { if (need_old_data && has_old_file && old_size > 0 && old_size <= MAX_RECEIVE_WHOLE_FILE_SIZE &&
old_size <= SIZE_MAX) {
old_data = protocol_alloc((size_t)old_size); old_data = protocol_alloc((size_t)old_size);
if (old_data) { if (old_data) {
size_t got = 0; size_t got = 0;
@@ -433,73 +489,63 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
} }
} }
} }
if (old_fd >= 0) {
close(old_fd);
}
bool match = /* Quick-skip decision. If no content comparison is required this is final
!config->ignore_times && has_old_file && (unsigned long long)st.st_size == check_size; and the old file was never read; if the read failed the file is not
if (match && config->checksum) { skipped and the transfer proceeds with the full new contents. */
uint64_t old_checksum = old_size == 0 ? delta_xxhash64("", 0) : 0; bool match = false;
if (old_data) if (checksum_needs_read) {
old_checksum = delta_xxhash64(old_data, (size_t)old_size); if (old_size == 0)
match = (old_size == 0 || old_data) && old_checksum == check_checksum; match = delta_xxhash64("", 0) == check_checksum;
free(old_data); else
old_data = NULL; match = old_data != NULL && delta_xxhash64(old_data, (size_t)old_size) == check_checksum;
} else if (match && !config->size_only) { } else if (size_equal && !config->ignore_times) {
long long old_mtime_nsec = 0; match = config->size_only || match_by_metadata;
#ifdef __linux__
old_mtime_nsec = st.st_mtim.tv_nsec;
#endif
match = metadata_mtime_matches(st.st_mtime, old_mtime_nsec, (time_t)check_mtime,
(long)check_mtime_nsec, config->modify_window);
} }
if (match) { if (match) {
free(old_data); free(old_data);
if (!send_status(fd, STATUS_OK)) { if (!send_status(fd, STATUS_OK)) {
close(old_fd);
free(full_path); free(full_path);
free(check_path); free(check_path);
return NULL; return NULL;
} }
close(old_fd);
free(full_path); free(full_path);
free(check_path); free(check_path);
*skipped = true; *skipped = true;
return NULL; return NULL;
} }
bool try_delta = config->use_delta && !config->whole_file && has_old_file && old_data != NULL && if (try_delta && old_data != NULL) {
delta_should_attempt(old_size, check_size, config->delta_max_file_size);
if (try_delta) {
bool delta_failed = false; bool delta_failed = false;
File* delta_file = File* delta_file =
receive_delta_file(fd, config, check_path, old_data, old_size, &delta_failed); receive_delta_file(fd, config, check_path, old_data, old_size, &delta_failed);
old_data = NULL; /* receive_delta_file consumes the snapshot on every path */ old_data = NULL; /* receive_delta_file consumes the snapshot on every path */
if (delta_file) { if (delta_file) {
close(old_fd);
free(full_path); free(full_path);
free(check_path); free(check_path);
return delta_file; return delta_file;
} }
if (delta_failed) { if (delta_failed) {
close(old_fd);
free(full_path); free(full_path);
free(check_path); free(check_path);
return NULL; return NULL;
} }
free(old_data);
old_data = NULL;
try_delta = false;
} }
free(old_data);
old_data = NULL;
if (!try_delta) { if (!send_status(fd, STATUS_NEXT)) {
free(old_data); close(old_fd);
old_data = NULL; free(full_path);
if (!send_status(fd, STATUS_NEXT)) { free(check_path);
free(full_path); return NULL;
free(check_path);
return NULL;
}
} }
close(old_fd);
File* file = file_create(check_path); File* file = file_create(check_path);
free(check_path); free(check_path);
@@ -517,7 +563,7 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
} }
} }
Data* file_data = receive_data_limited(fd, MAX_RECEIVE_FILE_SIZE); Data* file_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE);
if (file_data == NULL) { if (file_data == NULL) {
file_destroy(file); file_destroy(file);
return NULL; return NULL;
@@ -527,7 +573,7 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
!compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes, !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes,
config->skip_compress_set ? config->skip_compress_count config->skip_compress_set ? config->skip_compress_count
: -1)) { : -1)) {
Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_FILE_SIZE); Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE);
data_destroy(file_data); data_destroy(file_data);
if (uncompressed == NULL) { if (uncompressed == NULL) {
file_destroy(file); file_destroy(file);
@@ -536,7 +582,6 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
if (uncompressed->size > MAX_FILE_DATA_SIZE) { if (uncompressed->size > MAX_FILE_DATA_SIZE) {
data_destroy(uncompressed); data_destroy(uncompressed);
file_destroy(file); file_destroy(file);
send_status(fd, STATUS_ERROR);
return NULL; return NULL;
} }
file_data = uncompressed; file_data = uncompressed;
@@ -571,7 +616,7 @@ File* file_receive(const Config* config, int file_descriptor) {
return NULL; return NULL;
} }
} }
Data* file_data = receive_data_limited(file_descriptor, MAX_RECEIVE_FILE_SIZE); Data* file_data = receive_data_limited(file_descriptor, MAX_RECEIVE_WHOLE_FILE_SIZE);
if (file_data == NULL) { if (file_data == NULL) {
file_destroy(file); file_destroy(file);
return NULL; return NULL;
@@ -580,7 +625,7 @@ File* file_receive(const Config* config, int file_descriptor) {
!compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes, !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes,
config->skip_compress_set ? config->skip_compress_count config->skip_compress_set ? config->skip_compress_count
: -1)) { : -1)) {
Data* file_data_uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_FILE_SIZE); Data* file_data_uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE);
data_destroy(file_data); data_destroy(file_data);
if (file_data_uncompressed == NULL) { if (file_data_uncompressed == NULL) {
file_destroy(file); file_destroy(file);
+8
View File
@@ -10,6 +10,14 @@
File* file_receive(const Config* config, int file_descriptor); File* file_receive(const Config* config, int file_descriptor);
File* receive_incremental_check(int fd, const Config* config, bool* skipped); File* receive_incremental_check(int fd, const Config* config, bool* skipped);
int receive_manifest(int fd, const Config* config, int* next_status); int receive_manifest(int fd, const Config* config, int* next_status);
/* Outcome of a single file_save_to_disk operation. The receiver needs to
distinguish "written" from "skipped" so --remove-source-files can be told
which sources were actually stored. */
typedef enum { FILE_SAVE_ERROR = 0, FILE_SAVE_WRITTEN = 1, FILE_SAVE_SKIPPED = 2 } FileSaveResult;
FileSaveResult file_save_to_disk_full(const char* root_directory, const File* file,
const Config* config);
bool file_save_to_disk(const char* root_directory, const File* file, const Config* config); bool file_save_to_disk(const char* root_directory, const File* file, const Config* config);
#endif #endif
+97 -9
View File
@@ -105,9 +105,14 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
context->queue = queue; context->queue = queue;
context->file_descriptor = file_descriptor; context->file_descriptor = file_descriptor;
context->ssl = ssl; context->ssl = ssl;
context->outcomes.entries = NULL;
context->outcomes.count = 0;
context->outcomes.capacity = 0;
protocol_session_init(&context->session, file_descriptor, file_descriptor); protocol_session_init(&context->session, file_descriptor, file_descriptor);
protocol_session_set_ssl(&context->session, ssl); protocol_session_set_ssl(&context->session, ssl);
context->receiver_done = false; context->receiver_done = false;
context->queued_bytes = 0;
context->max_queue_bytes = 0;
atomic_init(&context->cancelled, false); atomic_init(&context->cancelled, false);
int init = 0; int init = 0;
if (mtx_init(&context->mutex, mtx_plain) != thrd_success) if (mtx_init(&context->mutex, mtx_plain) != thrd_success)
@@ -137,20 +142,80 @@ fail:
void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { void pipeline_context_receiver_destroy(PipelineContextReceiver* context) {
config_delete(context->config); config_delete(context->config);
queue_destroy(context->queue); queue_destroy(context->queue);
receiver_outcomes_destroy(&context->outcomes);
mtx_destroy(&context->mutex); mtx_destroy(&context->mutex);
cnd_destroy(&context->condition_not_full); cnd_destroy(&context->condition_not_full);
cnd_destroy(&context->condition_not_empty); cnd_destroy(&context->condition_not_empty);
free(context); free(context);
} }
void pipeline_context_receiver_set_queue_byte_limit(PipelineContextReceiver* context,
size_t max_bytes) {
if (context == NULL)
return;
mtx_lock(&context->mutex);
context->max_queue_bytes = max_bytes;
context->queued_bytes = 0;
cnd_broadcast(&context->condition_not_full);
mtx_unlock(&context->mutex);
}
void pipeline_context_receiver_note_bytes_released(PipelineContextReceiver* context,
size_t released_bytes) {
if (context == NULL || context->max_queue_bytes == 0 || released_bytes == 0)
return;
mtx_lock(&context->mutex);
if (released_bytes >= context->queued_bytes)
context->queued_bytes = 0;
else
context->queued_bytes -= released_bytes;
cnd_signal(&context->condition_not_full);
mtx_unlock(&context->mutex);
}
bool pipeline_context_receiver_enqueue_file(PipelineContextReceiver* context, File* file) {
if (context == NULL || file == NULL)
return false;
size_t file_bytes = file->data ? file->data->size : 0;
mtx_lock(&context->mutex);
while (!atomic_load(&context->cancelled)) {
bool blocked_by_count = queue_is_full(context->queue);
bool blocked_by_budget = false;
if (context->max_queue_bytes > 0) {
size_t budget = context->max_queue_bytes;
size_t used = context->queued_bytes;
if (used >= budget) {
blocked_by_budget = true;
} else if (file_bytes > budget - used) {
/* A single payload larger than the whole budget (not possible with
the per-file receive cap) is only admitted to an empty pipeline so
the wait can never deadlock. */
blocked_by_budget = used != 0;
}
}
if (!blocked_by_count && !blocked_by_budget)
break;
cnd_wait(&context->condition_not_full, &context->mutex);
}
if (atomic_load(&context->cancelled)) {
mtx_unlock(&context->mutex);
file_destroy(file);
return false;
}
if (!queue_enqueue(context->queue, file)) {
mtx_unlock(&context->mutex);
file_destroy(file);
return false;
}
context->queued_bytes += file_bytes;
cnd_signal(&context->condition_not_empty);
mtx_unlock(&context->mutex);
return true;
}
static bool receiver_enqueue_file(File* file, void* context_pointer) { static bool receiver_enqueue_file(File* file, void* context_pointer) {
PipelineContextReceiver* context = context_pointer; PipelineContextReceiver* context = (PipelineContextReceiver*)context_pointer;
if (queue_enqueue_multithreaded_cancel(context->queue, file, &context->mutex, return pipeline_context_receiver_enqueue_file(context, file);
&context->condition_not_empty,
&context->condition_not_full, &context->cancelled))
return true;
file_destroy(file);
return false;
} }
static void receiver_thread_fail(PipelineContextReceiver* context) { static void receiver_thread_fail(PipelineContextReceiver* context) {
@@ -170,7 +235,7 @@ int receive_thread(void* pipeline_context) {
const Config* config = context->config; const Config* config = context->config;
mtx_unlock(&context->mutex); mtx_unlock(&context->mutex);
ReceiverSink sink = {receiver_enqueue_file, context, false, false}; ReceiverSink sink = {receiver_enqueue_file, context, false, false, NULL};
if (receiver_process((Config*)config, file_descriptor, &sink) != 0) { if (receiver_process((Config*)config, file_descriptor, &sink) != 0) {
receiver_thread_fail(context); receiver_thread_fail(context);
protocol_session_unbind(); protocol_session_unbind();
@@ -211,8 +276,30 @@ int write_thread(void* pipeline_context) {
protocol_session_unbind(); protocol_session_unbind();
return thrd_success; return thrd_success;
} }
if (save_to_disk && !file_save_to_disk(root_directory, file, context->config)) { size_t file_bytes = file->data ? file->data->size : 0;
FileSaveResult result = FILE_SAVE_SKIPPED;
if (save_to_disk) {
result = file_save_to_disk_full(root_directory, file, context->config);
if (result == FILE_SAVE_ERROR) {
file_destroy(file);
pipeline_context_receiver_note_bytes_released(context, file_bytes);
mtx_lock(&context->mutex);
atomic_store(&context->cancelled, true);
context->receiver_done = true;
cnd_broadcast(&context->condition_not_full);
cnd_broadcast(&context->condition_not_empty);
mtx_unlock(&context->mutex);
free(root_directory);
protocol_session_unbind();
return thrd_error;
}
}
/* Record the per-file outcome so a --remove-source-files sender learns
which sources were actually written versus skipped on the receiver. */
if (context->config->remove_source_files &&
!receiver_outcomes_append(&context->outcomes, (unsigned char)result)) {
file_destroy(file); file_destroy(file);
pipeline_context_receiver_note_bytes_released(context, file_bytes);
mtx_lock(&context->mutex); mtx_lock(&context->mutex);
atomic_store(&context->cancelled, true); atomic_store(&context->cancelled, true);
context->receiver_done = true; context->receiver_done = true;
@@ -224,5 +311,6 @@ int write_thread(void* pipeline_context) {
return thrd_error; return thrd_error;
} }
file_destroy(file); file_destroy(file);
pipeline_context_receiver_note_bytes_released(context, file_bytes);
} }
} }
+22
View File
@@ -9,6 +9,7 @@
#include "file.h" #include "file.h"
#include "protocol.h" #include "protocol.h"
#include "queue.h" #include "queue.h"
#include "receiver.h"
#include <openssl/ssl.h> #include <openssl/ssl.h>
typedef struct { typedef struct {
@@ -40,11 +41,20 @@ typedef struct PipelineContextReceiver {
int file_descriptor; int file_descriptor;
SSL* ssl; SSL* ssl;
ProtocolSession session; ProtocolSession session;
ReceiverOutcomes outcomes;
mtx_t mutex; mtx_t mutex;
cnd_t condition_not_full; cnd_t condition_not_full;
cnd_t condition_not_empty; cnd_t condition_not_empty;
bool receiver_done; bool receiver_done;
atomic_bool cancelled; atomic_bool cancelled;
/* Aggregate payload bytes that have been received but not yet released by
the disk writer (queued or in the writer's hand). Guarded by `mutex`.
When `max_queue_bytes` is non-zero the receiver blocks before enqueuing
once this total would exceed it, so decompressed/copied file payloads
buffered ahead of a slow disk writer respect the per-connection memory
budget instead of growing without bound. */
size_t queued_bytes;
size_t max_queue_bytes;
} PipelineContextReceiver; } PipelineContextReceiver;
PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner,
@@ -53,6 +63,18 @@ void pipeline_context_sender_destroy(PipelineContextSender* context);
PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue_receiver, PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue_receiver,
int file_descriptor, SSL* ssl); int file_descriptor, SSL* ssl);
void pipeline_context_receiver_destroy(PipelineContextReceiver* context); void pipeline_context_receiver_destroy(PipelineContextReceiver* context);
/* Bound the bytes buffered ahead of the disk writer (see max_queue_bytes). */
void pipeline_context_receiver_set_queue_byte_limit(PipelineContextReceiver* context,
size_t max_bytes);
/* Blocking enqueue used by the receive pipeline sink. Blocks while the queue
is full by element count or when adding `file` would push queued_bytes over
the configured byte limit; waits until the disk writer releases bytes.
Takes ownership of `file` on success and destroys it on failure/cancel. */
bool pipeline_context_receiver_enqueue_file(PipelineContextReceiver* context, File* file);
/* Account for `released_bytes` of payload memory that has been freed by the
disk writer, unblocking a receiver that is waiting on the byte limit. */
void pipeline_context_receiver_note_bytes_released(PipelineContextReceiver* context,
size_t released_bytes);
int receive_thread(void* pipeline_context); int receive_thread(void* pipeline_context);
int write_thread(void* pipeline_context); int write_thread(void* pipeline_context);
#endif #endif
-1
View File
@@ -14,7 +14,6 @@
#define RECEIVE_TIMEOUT_SEC 60 /* 60 second per-message timeout */ #define RECEIVE_TIMEOUT_SEC 60 /* 60 second per-message timeout */
#define SEND_TIMEOUT_SEC 60 #define SEND_TIMEOUT_SEC 60
#define MAX_CONNECTION_MEMORY (256ULL * 1024 * 1024) /* bounded cumulative receive budget */
static __thread int io_read_fd = -1; static __thread int io_read_fd = -1;
static __thread int io_write_fd = -1; static __thread int io_write_fd = -1;
+14 -4
View File
@@ -9,10 +9,16 @@
/* Maximum allowed string size for receive_str (64 KB) */ /* Maximum allowed string size for receive_str (64 KB) */
#define MAX_STRING_SIZE (64 * 1024) #define MAX_STRING_SIZE (64 * 1024)
/* Maximum allowed data payload size for receive_data (100 MB) */ /* Maximum uncompressed file payload accepted by the receiver's whole-file
#define MAX_DATA_PAYLOAD_SIZE (100ULL * 1024 * 1024) * paths. A single whole file is charged against the per-connection memory
/* Maximum uncompressed file payload accepted by the receiver. */ * reservation (MAX_CONNECTION_MEMORY) and against the server allocation
#define MAX_RECEIVE_FILE_SIZE (64ULL * 1024 * 1024) * ceiling (MAX_SERVER_ALLOC), so this mirrors those 256 MB bounds rather than
* the older 64 MB chunk-era cap. Chunk-serialized payloads keep their own
* 64 MB cap (MAX_CHUNK_SIZE). */
#define MAX_RECEIVE_WHOLE_FILE_SIZE (256ULL * 1024 * 1024)
/* Maximum allowed data payload size for receive_data (whole-file bound) */
#define MAX_DATA_PAYLOAD_SIZE MAX_RECEIVE_WHOLE_FILE_SIZE
/* Maximum chunk size (64 MB) — prevents unbounded allocation from the wire */ /* Maximum chunk size (64 MB) — prevents unbounded allocation from the wire */
#define MAX_CHUNK_SIZE (64ULL * 1024 * 1024) #define MAX_CHUNK_SIZE (64ULL * 1024 * 1024)
@@ -22,6 +28,10 @@
#define DEFAULT_MAX_ALLOC (1ULL * 1024 * 1024 * 1024) #define DEFAULT_MAX_ALLOC (1ULL * 1024 * 1024 * 1024)
/* Server policy ceiling for a client-provided allocation limit. */ /* Server policy ceiling for a client-provided allocation limit. */
#define MAX_SERVER_ALLOC (256ULL * 1024 * 1024) #define MAX_SERVER_ALLOC (256ULL * 1024 * 1024)
/* Bounded cumulative per-connection receive budget. In-flight wire buffers,
decompression buffers and queued (not yet written) file payloads for a
connection must stay within this ceiling. */
#define MAX_CONNECTION_MEMORY (256ULL * 1024 * 1024)
typedef struct ssl_st SSL; typedef struct ssl_st SSL;
-40
View File
@@ -1,6 +1,5 @@
#include "utils.h" #include "utils.h"
#include "array_list.h" #include "array_list.h"
#include "libgen.h"
#include "log.h" #include "log.h"
#include <dirent.h> #include <dirent.h>
#include <errno.h> #include <errno.h>
@@ -83,45 +82,6 @@ static int open_authorized_destination(const char* dest_root) {
return dirfd; return dirfd;
} }
bool mkdir_r(const char* path) {
if (!path || *path == '\0')
return false;
char* duplicate = str_dup(path);
if (!duplicate)
return false;
int dirfd = open(path[0] == '/' ? "/" : ".", O_RDONLY | O_DIRECTORY | O_CLOEXEC | O_NOFOLLOW);
if (dirfd < 0) {
free(duplicate);
return false;
}
bool ok = true;
char* saveptr = NULL;
char* component = strtok_r(duplicate, "/", &saveptr);
while (component) {
if (strcmp(component, "..") == 0) {
ok = false;
break;
}
if (strcmp(component, ".") != 0) {
int next = openat(dirfd, component, O_RDONLY | O_DIRECTORY | O_CLOEXEC | O_NOFOLLOW);
if (next < 0 && errno == ENOENT) {
if (mkdirat(dirfd, component, 0755) == 0 || errno == EEXIST)
next = openat(dirfd, component, O_RDONLY | O_DIRECTORY | O_CLOEXEC | O_NOFOLLOW);
}
if (next < 0) {
ok = false;
break;
}
close(dirfd);
dirfd = next;
}
component = strtok_r(NULL, "/", &saveptr);
}
close(dirfd);
free(duplicate);
return ok;
}
char* str_dup(const char* string) { char* str_dup(const char* string) {
if (string == NULL) if (string == NULL)
return NULL; return NULL;
-1
View File
@@ -5,7 +5,6 @@
#include <stddef.h> #include <stddef.h>
#include <stdbool.h> #include <stdbool.h>
bool mkdir_r(const char* path);
char* str_dup(const char* string); char* str_dup(const char* string);
char* output_escape(const char* string, bool eight_bit_output); char* output_escape(const char* string, bool eight_bit_output);
char* path_cat(const char* path1, const char* path2); char* path_cat(const char* path1, const char* path2);
+208
View File
@@ -1,4 +1,5 @@
"""Feature tests: incremental sync, bandwidth limiting, dry run, metadata, filters.""" """Feature tests: incremental sync, bandwidth limiting, dry run, metadata, filters."""
import filecmp
import os import os
import shutil import shutil
import sys import sys
@@ -813,3 +814,210 @@ class TestBandwidthLimit:
mismatches, missing = verify_transfer(SOURCE_DIR, received) mismatches, missing = verify_transfer(SOURCE_DIR, received)
assert not missing, f"Missing: {missing}" assert not missing, f"Missing: {missing}"
assert not mismatches, f"Mismatch: {mismatches}" assert not mismatches, f"Mismatch: {mismatches}"
def _read_file(path):
with open(path, "rb") as fh:
return fh.read()
class TestRemoveSourceFilesSkips:
"""--remove-source-files must not delete sources the receiver skipped
(rsync reference behavior)."""
def test_existing_first_sync_keeps_new_source(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "remove_rsf_existing_src")
dest = os.path.join(TEST_DATA_DIR, "remove_rsf_existing_dst")
clean_dir(source)
clean_dir(dest)
with open(os.path.join(source, "only.txt"), "wb") as f:
f.write(b"keep me")
result, _ = run_client(source, dest, flags=["--remove-source-files", "--existing"],
port=shared_server.port)
assert result.returncode == 0, f"Sync failed: {result.stderr[:200]}"
# The file exists only on the source side, so --existing makes the
# receiver skip it; the source must therefore not be removed.
assert os.path.isfile(os.path.join(source, "only.txt"))
received = get_dest_received_dir(dest, source)
assert not os.path.exists(os.path.join(received, "only.txt"))
def test_ignore_existing_keeps_skipped_source(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "remove_rsf_ignore_src")
dest = os.path.join(TEST_DATA_DIR, "remove_rsf_ignore_dst")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "file.txt")
with open(source_file, "wb") as f:
f.write(b"payload")
result, _ = run_client(source, dest, port=shared_server.port)
assert result.returncode == 0
result, _ = run_client(source, dest, flags=["--remove-source-files", "--ignore-existing"],
port=shared_server.port)
assert result.returncode == 0, f"Sync failed: {result.stderr[:200]}"
# Destination already has the file, so the second run is a receiver
# skip; the source file must survive.
assert os.path.isfile(source_file)
received = get_dest_received_dir(dest, source)
assert _read_file(os.path.join(received, "file.txt")) == b"payload"
def test_update_newer_destination_keeps_source(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "remove_rsf_update_src")
dest = os.path.join(TEST_DATA_DIR, "remove_rsf_update_dst")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "file.txt")
with open(source_file, "wb") as f:
f.write(b"source payload")
result, _ = run_client(source, dest, port=shared_server.port)
assert result.returncode == 0
received = get_dest_received_dir(dest, source)
received_file = os.path.join(received, "file.txt")
with open(received_file, "wb") as f:
f.write(b"newer destination payload")
os.utime(received_file, ns=(time.time_ns() + 10**9, time.time_ns() + 10**9))
result, _ = run_client(source, dest, flags=["--remove-source-files", "--update"],
port=shared_server.port)
assert result.returncode == 0, f"Sync failed: {result.stderr[:200]}"
# --update skips a destination that is newer than the source, so the
# source must not be removed.
assert os.path.isfile(source_file)
assert _read_file(received_file) == b"newer destination payload"
def test_multithreaded_ignore_existing_keeps_skipped_source(self, shared_server):
"""The multithreaded writer path must also report per-file outcomes so a
--remove-source-files sender does not delete skipped sources."""
source = os.path.join(TEST_DATA_DIR, "remove_rsf_mt_ignore_src")
dest = os.path.join(TEST_DATA_DIR, "remove_rsf_mt_ignore_dst")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "file.txt")
with open(source_file, "wb") as f:
f.write(b"payload")
result, _ = run_client(source, dest, port=shared_server.port)
assert result.returncode == 0
result, _ = run_client(source, dest,
flags=["--remove-source-files", "--ignore-existing", "-m"],
port=shared_server.port)
assert result.returncode == 0, f"Sync failed: {result.stderr[:200]}"
# Destination already has the file, so the receiver (writer thread)
# skips it; the source must survive.
assert os.path.isfile(source_file)
received = get_dest_received_dir(dest, source)
assert _read_file(os.path.join(received, "file.txt")) == b"payload"
class TestBackup:
def _sync(self, source, dest, flags, port):
return run_client(source, dest, flags=flags, port=port)
def test_plain_backup_keeps_previous_version(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "backup_src")
dest = os.path.join(TEST_DATA_DIR, "backup_dst")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "f.txt")
with open(source_file, "wb") as f:
f.write(b"AAAA")
result, _ = self._sync(source, dest, ["--backup"], shared_server.port)
assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}"
with open(source_file, "wb") as f:
f.write(b"BBBB")
result, _ = self._sync(source, dest, ["--backup"], shared_server.port)
assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
assert _read_file(os.path.join(received, "f.txt")) == b"BBBB"
# rsync default suffix "~" keeps the overwritten version.
assert _read_file(os.path.join(received, "f.txt~")) == b"AAAA"
def test_backup_custom_suffix(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "backup_suffix_src")
dest = os.path.join(TEST_DATA_DIR, "backup_suffix_dst")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "f.txt")
with open(source_file, "wb") as f:
f.write(b"AAAA")
flags = ["--backup", "--suffix", ".bak"]
result, _ = self._sync(source, dest, flags, shared_server.port)
assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}"
with open(source_file, "wb") as f:
f.write(b"BBBB")
result, _ = self._sync(source, dest, flags, shared_server.port)
assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
assert _read_file(os.path.join(received, "f.txt")) == b"BBBB"
assert _read_file(os.path.join(received, "f.txt.bak")) == b"AAAA"
def test_backup_dir_stores_backups_separately(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "backup_dir_src")
dest = os.path.join(TEST_DATA_DIR, "backup_dir_dst")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "f.txt")
with open(source_file, "wb") as f:
f.write(b"AAAA")
flags = ["--backup", "--backup-dir", "backups"]
result, _ = self._sync(source, dest, flags, shared_server.port)
assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}"
with open(source_file, "wb") as f:
f.write(b"BBBB")
result, _ = self._sync(source, dest, flags, shared_server.port)
assert result.returncode == 0, f"Backup sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
assert _read_file(os.path.join(received, "f.txt")) == b"BBBB"
backup = os.path.join(dest, "backups", os.path.relpath(source_file, os.path.sep))
assert _read_file(backup) == b"AAAA"
class TestPartialDir:
def test_completed_transfer_installed_in_destination(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "partial_src")
dest = os.path.join(TEST_DATA_DIR, "partial_dst")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "f.txt")
with open(source_file, "wb") as f:
f.write(b"partial payload")
result, _ = run_client(source, dest, flags=["--partial", "--partial-dir", ".partial"],
port=shared_server.port)
assert result.returncode == 0, f"Partial sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
assert _read_file(os.path.join(received, "f.txt")) == b"partial payload"
# A completed transfer must not remain under the partial directory.
partial = os.path.join(dest, ".partial", os.path.relpath(source_file, os.path.sep))
assert not os.path.exists(partial)
class TestLargeFile:
def test_transfer_100mb_file(self, shared_server):
source = os.path.join(TEST_DATA_DIR, "large_src")
dest = os.path.join(TEST_DATA_DIR, "large_dst")
clean_dir(source)
clean_dir(dest)
source_file = os.path.join(source, "big.bin")
chunk = os.urandom(1024 * 1024)
with open(source_file, "wb") as f:
for _ in range(100):
f.write(chunk)
result, _ = run_client(source, dest, port=shared_server.port)
assert result.returncode == 0, f"Large-file sync failed: {result.stderr[:200]}"
received = get_dest_received_dir(dest, source)
assert filecmp.cmp(source_file, os.path.join(received, "big.bin"), shallow=False)
+67
View File
@@ -1088,6 +1088,70 @@ static void test_parse_args_rejects_invalid_compression_choice() {
config_delete(cfg); config_delete(cfg);
} }
/* Every value-taking table option accepts an inline "--opt=value" form. */
static void test_parse_args_table_equals_size_options() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--max-size=2G", "--min-size=1K", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0);
EXPECT_TRUE(cfg->max_size == 2ULL * 1024 * 1024 * 1024);
EXPECT_TRUE(cfg->min_size == 1024ULL);
EXPECT_EQ_INT(positional_count, 2);
config_delete(cfg);
}
static void test_parse_args_table_equals_string_and_int_options() {
Config* cfg = config_create();
char* argv[] = {"fastsync", "--suffix=.bak", "--timeout=30", "--max-depth=5", "/src", "/dst"};
int positional_args[2];
int positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 6, argv, positional_args, &positional_count), 0);
EXPECT_EQ_STR(cfg->suffix, ".bak");
EXPECT_EQ_INT(cfg->timeout, 30);
EXPECT_EQ_INT(cfg->max_depth, 5);
EXPECT_EQ_INT(positional_count, 2);
config_delete(cfg);
cfg = config_create();
char* backup_argv[] = {"fastsync", "--backup-dir=/tmp/bak", "/src", "/dst"};
positional_count = 0;
EXPECT_EQ_INT(parse_args(cfg, 4, backup_argv, positional_args, &positional_count), 0);
EXPECT_EQ_STR(cfg->backup_dir, "/tmp/bak");
config_delete(cfg);
}
/* Options that take a separate value must report "missing argument", not the
* generic "Unknown option", when they are the final argv entry. */
static void test_parse_args_missing_argument_diagnostic() {
static const char* const options[] = {"--exclude", "--server-port", "--skip-compress", "-T"};
for (size_t i = 0; i < sizeof(options) / sizeof(options[0]); i++) {
Config* cfg = config_create();
char* argv[] = {"fastsync", (char*)options[i]};
int positional_args[2];
int positional_count = 0;
FILE* log_file = tmpfile();
char log_buffer[512] = {0};
EXPECT_NOT_NULL(log_file);
log_set_file(log_file);
EXPECT_EQ_INT(parse_args(cfg, 2, argv, positional_args, &positional_count), -1);
fflush(log_file);
rewind(log_file);
EXPECT_TRUE(fread(log_buffer, 1, sizeof(log_buffer) - 1, log_file) > 0);
EXPECT_TRUE(strstr(log_buffer, "missing argument") != NULL);
EXPECT_TRUE(strstr(log_buffer, "Unknown option") == NULL);
log_set_file(NULL);
fclose(log_file);
config_delete(cfg);
}
}
void test_client_cli() { void test_client_cli() {
test_validate_config_required_paths(); test_validate_config_required_paths();
test_validate_config_incompatible_options(); test_validate_config_incompatible_options();
@@ -1155,6 +1219,9 @@ void test_client_cli() {
test_parse_args_compression_alias_equals(); test_parse_args_compression_alias_equals();
test_parse_args_rejects_invalid_compression_level_equals(); test_parse_args_rejects_invalid_compression_level_equals();
test_parse_args_rejects_invalid_compression_choice(); test_parse_args_rejects_invalid_compression_choice();
test_parse_args_table_equals_size_options();
test_parse_args_table_equals_string_and_int_options();
test_parse_args_missing_argument_diagnostic();
test_parse_args_partial_progress(); test_parse_args_partial_progress();
test_parse_args_checksum_choice_aliases(); test_parse_args_checksum_choice_aliases();
test_parse_args_checksum_choice_requires_value(); test_parse_args_checksum_choice_requires_value();
+100
View File
@@ -282,6 +282,105 @@ static void test_config_receive_truncated() {
close(p[1]); close(p[1]);
} }
static bool config_string_roundtrip_matches(const Config* send_cfg, Config* recv) {
/* The sender serializes NULL strings as "" on the wire. Receivers must
canonicalize those empty values back to NULL for the options whose client
default is NULL (backup_dir, temp_dir, partial_dir, suffix), while a real
non-empty value round-trips unchanged. */
const char* fields[4];
char* const* recv_fields[4];
fields[0] = send_cfg->backup_dir;
recv_fields[0] = &recv->backup_dir;
fields[1] = send_cfg->temp_dir;
recv_fields[1] = &recv->temp_dir;
fields[2] = send_cfg->partial_dir;
recv_fields[2] = &recv->partial_dir;
fields[3] = send_cfg->suffix;
recv_fields[3] = &recv->suffix;
for (int i = 0; i < 4; i++) {
const char* sent = fields[i];
const char* got = *recv_fields[i];
if (sent == NULL || sent[0] == '\0') {
if (got != NULL)
return false;
} else if (got == NULL || strcmp(sent, got) != 0) {
return false;
}
}
return true;
}
static bool roundtrip_config_ok(const Config* send_cfg) {
int p[2];
if (socketpair(AF_UNIX, SOCK_STREAM, 0, p) != 0)
return false;
pid_t pid = fork();
if (pid == 0) {
close(p[1]);
io_set_fds(p[0], p[0]);
Config* recv = config_receive(p[0]);
bool ok = recv != NULL;
if (ok) {
ok = recv->version != NULL && strcmp(recv->version, PROTOCOL_VERSION) == 0;
ok = ok && recv->send_directory && recv->receive_root_directory;
ok = ok && config_string_roundtrip_matches(send_cfg, recv);
}
config_delete(recv);
close(p[0]);
_exit(ok ? 0 : 1);
} else {
close(p[0]);
io_set_fds(p[1], p[1]);
bool sent = config_send(p[1], send_cfg);
int status;
waitpid(pid, &status, 0);
close(p[1]);
return sent && WIFEXITED(status) && WEXITSTATUS(status) == 0;
}
}
/* Issue #252: NULL-vs-empty must survive the wire for backup_dir, temp_dir,
partial_dir, and suffix. NULL and explicitly-empty client values are both
serialized as "" and must be reconstructed as NULL so plain --backup (with
no --suffix/--backup-dir) works exactly like the client configured it. */
static void test_config_string_null_vs_empty_roundtrip() {
if (is_running_under_valgrind())
return;
/* NULL values on the wire must come back as NULL. */
Config* a = config_create();
EXPECT_NOT_NULL(a);
a->send_directory = str_dup("/src");
a->receive_root_directory = str_dup("/dst");
EXPECT_TRUE(roundtrip_config_ok(a));
config_delete(a);
/* Explicitly empty strings (indistinguishable on the wire from NULL) must
be canonicalized to NULL by the receiver. */
Config* b = config_create();
EXPECT_NOT_NULL(b);
b->send_directory = str_dup("/src");
b->receive_root_directory = str_dup("/dst");
b->backup_dir = str_dup("");
b->temp_dir = str_dup("");
b->partial_dir = str_dup("");
b->suffix = str_dup("");
EXPECT_TRUE(roundtrip_config_ok(b));
config_delete(b);
/* Non-empty values must round-trip unchanged. */
Config* c = config_create();
EXPECT_NOT_NULL(c);
c->send_directory = str_dup("/src");
c->receive_root_directory = str_dup("/dst");
c->backup_dir = str_dup("backups");
c->temp_dir = str_dup("/tmp/fast");
c->partial_dir = str_dup(".partial");
c->suffix = str_dup(".bak");
EXPECT_TRUE(roundtrip_config_ok(c));
config_delete(c);
}
static void test_config_is_remote_dest() { static void test_config_is_remote_dest() {
/* Valid SSH-style destinations */ /* Valid SSH-style destinations */
EXPECT_TRUE(config_is_remote_dest("user@host:/path")); EXPECT_TRUE(config_is_remote_dest("user@host:/path"));
@@ -315,6 +414,7 @@ void test_config() {
test_config_send_receive(); test_config_send_receive();
test_config_send_receive_version_mismatch(); test_config_send_receive_version_mismatch();
test_config_receive_truncated(); test_config_receive_truncated();
test_config_string_null_vs_empty_roundtrip();
} }
test_config_is_remote_dest(); test_config_is_remote_dest();
} }
+293
View File
@@ -343,6 +343,297 @@ static void test_delta_apply_rejects_output_overflow() {
EXPECT_TRUE(delta_apply(old_data, sizeof(old_data), &delta, 1) == NULL); EXPECT_TRUE(delta_apply(old_data, sizeof(old_data), &delta, 1) == NULL);
} }
/* ---------------------------------------------------------------------------
* Hash-index lookup differential tests.
*
* delta_compute buckets signature blocks by their weak checksum. These tests
* prove the bucket-indexed candidate lookup is behaviour-identical to the
* original per-window linear scan: the emitted instruction stream (types,
* lengths, literal bytes and chosen block indices) must match a naive linear
* reference exactly, and the delta must reconstruct the new buffer.
* ------------------------------------------------------------------------- */
#define REF_NO_MATCH UINT32_MAX
typedef struct {
DeltaInstruction* items;
uint32_t count;
uint32_t cap;
} RefDelta;
static void ref_delta_free(RefDelta* ref) {
if (!ref->items)
return;
for (uint32_t i = 0; i < ref->count; i++)
if (ref->items[i].type == DELTA_INSTR_LITERAL)
free(ref->items[i].literal.data);
free(ref->items);
ref->items = NULL;
ref->count = 0;
ref->cap = 0;
}
static bool ref_delta_push(RefDelta* ref, DeltaInstruction instr) {
if (ref->count == ref->cap) {
uint32_t new_cap = ref->cap ? ref->cap * 2 : 16;
DeltaInstruction* tmp = realloc(ref->items, (size_t)new_cap * sizeof(DeltaInstruction));
if (!tmp)
return false;
ref->items = tmp;
ref->cap = new_cap;
}
ref->items[ref->count++] = instr;
return true;
}
static bool ref_delta_flush_literal(RefDelta* ref, const uint8_t* data, uint64_t start,
uint64_t end) {
if (start >= end)
return true;
uint8_t* lit = malloc((size_t)(end - start));
if (!lit)
return false;
memcpy(lit, data + start, (size_t)(end - start));
DeltaInstruction instr = {
.type = DELTA_INSTR_LITERAL,
.literal = {.data = lit, .length = (uint32_t)(end - start)},
};
return ref_delta_push(ref, instr);
}
/* Naive O(windows x blocks) re-implementation of the historical delta_compute
* candidate scan: only full windows may match, a candidate needs both the weak
* (Adler-32) and strong (xxHash32) checksums to agree, and the lowest matching
* block index is selected. */
static bool ref_delta_build(RefDelta* ref, const uint8_t* new_data, uint64_t new_size,
const DeltaSignature* sig) {
uint32_t block_size = sig->block_size;
uint64_t i = 0;
uint64_t literal_start = 0;
bool has_literal = false;
while (i < new_size) {
uint32_t window_len = (uint32_t)((new_size - i < block_size) ? (new_size - i) : block_size);
bool full_window = (window_len == block_size);
uint32_t matched = REF_NO_MATCH;
if (full_window) {
uint32_t adler = delta_adler32(new_data + i, window_len);
for (uint32_t j = 0; j < sig->block_count; j++) {
if (sig->blocks[j].adler32 == adler &&
delta_xxhash32(new_data + i, window_len) == sig->blocks[j].xxhash) {
matched = j;
break;
}
}
}
if (matched != REF_NO_MATCH) {
if (has_literal) {
if (!ref_delta_flush_literal(ref, new_data, literal_start, i))
return false;
has_literal = false;
}
DeltaInstruction instr = {
.type = DELTA_INSTR_BLOCK_MATCH,
.match = {.block_index = matched, .block_offset = 0, .length = window_len},
};
if (!ref_delta_push(ref, instr))
return false;
i += window_len;
} else {
if (!has_literal) {
literal_start = i;
has_literal = true;
}
i++;
}
}
if (has_literal && !ref_delta_flush_literal(ref, new_data, literal_start, new_size))
return false;
return true;
}
static bool ref_delta_matches(const RefDelta* ref, const Delta* delta) {
if (ref->count != delta->instruction_count)
return false;
for (uint32_t i = 0; i < ref->count; i++) {
const DeltaInstruction* a = &ref->items[i];
const DeltaInstruction* b = &delta->instructions[i];
if (a->type != b->type)
return false;
if (a->type == DELTA_INSTR_BLOCK_MATCH) {
if (a->match.block_index != b->match.block_index ||
a->match.block_offset != b->match.block_offset || a->match.length != b->match.length)
return false;
} else {
if (a->literal.length != b->literal.length ||
memcmp(a->literal.data, b->literal.data, a->literal.length) != 0)
return false;
}
}
return true;
}
static void expect_linear_reference_match(const uint8_t* old_data, uint64_t old_size,
const uint8_t* new_data, uint64_t new_size,
uint32_t block_size, const char* label) {
DeltaSignature* sig = delta_signature_create(old_data, old_size, block_size);
if (!sig) {
printf(" [FAIL] %s: signature creation failed\n", label);
EXPECT_NOT_NULL(sig);
return;
}
Delta* delta = delta_compute(new_data, new_size, sig, block_size);
if (!delta) {
printf(" [FAIL] %s: delta_compute returned NULL\n", label);
delta_signature_destroy(sig);
EXPECT_NOT_NULL(delta);
return;
}
RefDelta ref = {0};
bool ok = ref_delta_build(&ref, new_data, new_size, sig);
if (ok)
ok = ref_delta_matches(&ref, delta);
if (!ok) {
printf(" [FAIL] %s: instruction stream differs from linear reference "
"(linear=%u indexed=%u)\n",
label, ref.count, delta->instruction_count);
}
ref_delta_free(&ref);
delta_destroy(delta);
delta_signature_destroy(sig);
EXPECT_TRUE(ok);
}
static void fill_delta_pattern(uint8_t* buf, uint64_t size, uint32_t seed) {
uint32_t x = seed ? seed : 1;
for (uint64_t i = 0; i < size; i++) {
x ^= x << 13;
x ^= x >> 17;
x ^= x << 5;
buf[i] = (uint8_t)(x >> 24);
}
}
static void test_delta_hash_index_matches_linear_reference() {
/* Identical file (full block alignment). */
uint8_t old_a[32768];
uint8_t new_a[32768];
fill_delta_pattern(old_a, sizeof(old_a), 42);
memcpy(new_a, old_a, sizeof(old_a));
expect_linear_reference_match(old_a, sizeof(old_a), new_a, sizeof(new_a), 2048,
"identical 32KiB @ 2KiB");
/* Scattered single-byte edits in the middle of each block. */
uint8_t new_b[32768];
memcpy(new_b, old_a, sizeof(old_a));
for (size_t p = 100; p < sizeof(new_b); p += 4096)
new_b[p] ^= 0x5A;
expect_linear_reference_match(old_a, sizeof(old_a), new_b, sizeof(new_b), 2048,
"32KiB scattered single-byte edits @ 2KiB");
/* Non-aligned old file (partial final block) with a single edit. */
uint8_t old_c[30000];
uint8_t new_c[30000];
fill_delta_pattern(old_c, sizeof(old_c), 7);
memcpy(new_c, old_c, sizeof(old_c));
new_c[15000] ^= 0x3C;
expect_linear_reference_match(old_c, sizeof(old_c), new_c, sizeof(new_c), 2048,
"30KiB partial-tail single edit @ 2KiB");
/* Growth: appended data after an identical prefix. */
uint8_t old_d[24576];
uint8_t new_d[34576];
fill_delta_pattern(old_d, sizeof(old_d), 11);
memcpy(new_d, old_d, sizeof(old_d));
fill_delta_pattern(new_d + sizeof(old_d), sizeof(new_d) - sizeof(old_d), 23);
expect_linear_reference_match(old_d, sizeof(old_d), new_d, sizeof(new_d), 2048,
"24KiB -> 34KiB appended @ 2KiB");
/* Insertion shifting everything after the edit point (rsync re-sync). */
uint8_t old_e[65536];
uint8_t new_e[65536 + 3000];
fill_delta_pattern(old_e, sizeof(old_e), 99);
memcpy(new_e, old_e, 20000);
fill_delta_pattern(new_e + 20000, 3000, 101);
memcpy(new_e + 23000, old_e + 20000, sizeof(old_e) - 20000);
expect_linear_reference_match(old_e, sizeof(old_e), new_e, sizeof(new_e), 2048,
"64KiB + 3KiB insertion @ 2KiB");
/* Deletion shrinking the file. */
uint8_t new_f[sizeof(old_e) - 5000];
memcpy(new_f, old_e, 30000);
memcpy(new_f + 30000, old_e + 35000, sizeof(old_e) - 35000);
expect_linear_reference_match(old_e, sizeof(old_e), new_f, sizeof(new_f), 2048,
"64KiB - 5KiB deletion @ 2KiB");
/* Repeated identical blocks must resolve to the lowest block index. */
uint8_t old_g[4 * 4096];
uint8_t new_g[4 * 4096];
for (uint32_t b = 0; b < 4; b++)
fill_delta_pattern(old_g + b * 4096, 4096, b % 2 == 0 ? 500 : 501); /* block0==block2 */
memcpy(new_g, old_g, sizeof(old_g));
new_g[4096 + 5] ^= 0x11; /* edit inside the second (duplicated) chunk */
expect_linear_reference_match(old_g, sizeof(old_g), new_g, sizeof(new_g), 4096,
"duplicated chunks @ 4KiB");
/* Block larger than the file: nothing can match, all literal. */
uint8_t old_h[1000];
uint8_t new_h[1000];
fill_delta_pattern(old_h, sizeof(old_h), 3);
memcpy(new_h, old_h, sizeof(old_h));
expect_linear_reference_match(old_h, sizeof(old_h), new_h, sizeof(new_h), 4096,
"1KiB file @ 4KiB block");
}
static void test_delta_hash_index_large_mostly_matching() {
const uint64_t size = 4ULL * 1024 * 1024;
const uint32_t block_size = 8192;
uint8_t* old_data = malloc((size_t)size);
uint8_t* new_data = malloc((size_t)size);
EXPECT_TRUE(old_data != NULL && new_data != NULL);
fill_delta_pattern(old_data, size, 1234);
memcpy(new_data, old_data, (size_t)size);
/* Scattered single-byte changes across the whole buffer. Each change forces
* the diff to re-synchronise by walking one byte at a time through the
* affected block, which is exactly the case that used to cost O(bytes x
* blocks) with the linear scan. */
const uint64_t nchanges = 64;
for (uint64_t c = 0; c < nchanges; c++) {
uint64_t pos = (c * (size / nchanges)) + (c % 17);
new_data[pos] ^= (uint8_t)(0xA0 + (c % 16));
}
DeltaSignature* sig = delta_signature_create(old_data, size, block_size);
EXPECT_NOT_NULL(sig);
EXPECT_EQ_INT((int)sig->block_count, (int)(size / block_size));
Delta* delta = delta_compute(new_data, size, sig, block_size);
EXPECT_NOT_NULL(delta);
EXPECT_EQ_INT((int)delta->new_file_size, (int)size);
uint32_t match_count = 0;
for (uint32_t i = 0; i < delta->instruction_count; i++)
if (delta->instructions[i].type == DELTA_INSTR_BLOCK_MATCH)
match_count++;
EXPECT_TRUE(match_count > 0);
void* result = delta_apply(old_data, size, delta, block_size);
EXPECT_NOT_NULL(result);
EXPECT_EQ_INT(memcmp(result, new_data, (size_t)size), 0);
free(result);
delta_destroy(delta);
delta_signature_destroy(sig);
free(old_data);
free(new_data);
}
void test_delta() { void test_delta() {
test_adler32_basic(); test_adler32_basic();
test_adler32_different_data(); test_adler32_different_data();
@@ -360,4 +651,6 @@ void test_delta() {
test_is_worthwhile(); test_is_worthwhile();
test_large_file_delta(); test_large_file_delta();
test_delta_apply_rejects_output_overflow(); test_delta_apply_rejects_output_overflow();
test_delta_hash_index_matches_linear_reference();
test_delta_hash_index_large_mostly_matching();
} }
+267
View File
@@ -5,10 +5,12 @@
#include "utils.h" #include "utils.h"
#include "protocol.h" #include "protocol.h"
#include "test_utils.h" #include "test_utils.h"
#include <fcntl.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <sys/stat.h> #include <sys/stat.h>
#include <sys/wait.h> #include <sys/wait.h>
#include <time.h>
#include <unistd.h> #include <unistd.h>
static void test_file_create() { static void test_file_create() {
@@ -250,6 +252,111 @@ static void test_file_save_to_disk_ignore_existing_entry_types() {
rmdir(root); rmdir(root);
} }
/* Issue #253: with --partial --partial-dir a completed write must be installed
at the real destination rather than left under the partial directory. */
static void test_file_save_to_disk_partial_install() {
const char* root = "test_partial_install_tmp";
const char* dest_file = "test_partial_install_tmp/file.txt";
const char* partial_file = "test_partial_install_tmp/.partial/file.txt";
unlink(dest_file);
unlink(partial_file);
rmdir("test_partial_install_tmp/.partial");
rmdir(root);
File* f = file_create("file.txt");
EXPECT_NOT_NULL(f);
const char* content = "partial-dir content";
f->data->data = malloc(strlen(content));
EXPECT_NOT_NULL(f->data->data);
memcpy(f->data->data, content, strlen(content));
f->data->size = strlen(content);
Config* config = config_create();
EXPECT_NOT_NULL(config);
config->partial = true;
config->partial_dir = str_dup(".partial");
EXPECT_EQ_INT(file_save_to_disk_full(root, f, config), FILE_SAVE_WRITTEN);
FILE* fp = fopen(dest_file, "rb");
EXPECT_NOT_NULL(fp);
// cppcheck-suppress knownConditionTrueFalse
if (fp) {
char buf[64] = {0};
size_t nread = fread(buf, 1, sizeof(buf) - 1, fp);
fclose(fp);
EXPECT_EQ_INT((int)nread, (int)strlen(content));
EXPECT_EQ_INT(memcmp(buf, content, strlen(content)), 0);
}
/* A completed transfer must not linger under the partial dir. */
EXPECT_EQ_INT(access(partial_file, F_OK), -1);
file_destroy(f);
config_delete(config);
unlink(dest_file);
rmdir(root);
}
/* Issue #251: file_save_to_disk_full must distinguish receiver-side skips
(--existing/--ignore-existing/--update) from real writes so the sender can
decide whether --remove-source-files may unlink its source. */
static void test_file_save_to_disk_reports_skips() {
const char* root = "test_save_skip_tmp";
const char* existing_path = "test_save_skip_tmp/existing.txt";
unlink(existing_path);
rmdir(root);
EXPECT_TRUE(file_write_to_disk(existing_path, "old", 3, false, false));
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
File* new_file = file_create("missing.txt");
EXPECT_NOT_NULL(new_file);
new_file->data->data = malloc(7);
EXPECT_NOT_NULL(new_file->data->data);
memcpy(new_file->data->data, "skipped", 7);
new_file->data->size = 7;
/* --existing: destination is missing -> skipped, not an error. */
cfg->existing = true;
EXPECT_EQ_INT(file_save_to_disk_full(root, new_file, cfg), FILE_SAVE_SKIPPED);
cfg->existing = false;
/* --ignore-existing: destination present -> skipped. */
File* present = file_create("existing.txt");
EXPECT_NOT_NULL(present);
present->data->data = malloc(3);
EXPECT_NOT_NULL(present->data->data);
memcpy(present->data->data, "new", 3);
present->data->size = 3;
cfg->ignore_existing = true;
EXPECT_EQ_INT(file_save_to_disk_full(root, present, cfg), FILE_SAVE_SKIPPED);
cfg->ignore_existing = false;
/* A normal overwrite of an existing file is a real write. */
EXPECT_EQ_INT(file_save_to_disk_full(root, present, cfg), FILE_SAVE_WRITTEN);
/* --update: a newer destination is skipped. */
struct stat st;
EXPECT_EQ_INT(stat(existing_path, &st), 0);
time_t now = time(NULL);
FileMetadata metadata = {.mode = st.st_mode,
.uid = st.st_uid,
.gid = st.st_gid,
.mtime_sec = now - 100,
.mtime_nsec = 0};
present->metadata = &metadata;
cfg->update = true;
EXPECT_EQ_INT(file_save_to_disk_full(root, present, cfg), FILE_SAVE_SKIPPED);
present->metadata = NULL;
file_destroy(new_file);
file_destroy(present);
config_delete(cfg);
unlink(existing_path);
rmdir(root);
}
static void test_file_write_to_disk_basic() { static void test_file_write_to_disk_basic() {
const char* content = "Basic file_write_to_disk test"; const char* content = "Basic file_write_to_disk test";
EXPECT_TRUE(file_write_to_disk("test_file_write_to_disk_basic.txt", content, strlen(content), EXPECT_TRUE(file_write_to_disk("test_file_write_to_disk_basic.txt", content, strlen(content),
@@ -624,6 +731,161 @@ static void test_file_send_single_calls_metadata_and_path() {
} }
} }
static void test_inplace_overwrite_clears_special_mode_bits() {
const char* root = "test_inplace_tmp";
const char* path = "test_inplace_tmp/priv.txt";
const char* content = "olddata";
unlink(path);
rmdir(root);
EXPECT_EQ_INT(mkdir(root, 0700), 0);
/* Create a destination carrying setuid + sticky bits. */
int fd = open(path, O_WRONLY | O_CREAT | O_TRUNC | O_CLOEXEC, 0644);
EXPECT_TRUE(fd >= 0);
// cppcheck-suppress knownConditionTrueFalse
if (fd < 0) {
rmdir(root);
return;
}
EXPECT_EQ_INT((int)write(fd, content, strlen(content)), (int)strlen(content));
EXPECT_EQ_INT(fchmod(fd, S_ISUID | S_ISVTX | 0755), 0);
EXPECT_EQ_INT(close(fd), 0);
/* Overwrite in place without metadata: the mode must be normalized to a
safe default (0644) and the setuid/sticky bits must be gone. */
File* f = file_create("priv.txt");
EXPECT_NOT_NULL(f);
const char* new_content = "newdata";
f->data->data = malloc(strlen(new_content));
EXPECT_NOT_NULL(f->data->data);
memcpy(f->data->data, new_content, strlen(new_content));
f->data->size = strlen(new_content);
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
cfg->inplace = true;
EXPECT_TRUE(file_save_to_disk(root, f, cfg));
file_destroy(f);
config_delete(cfg);
struct stat st;
EXPECT_EQ_INT(stat(path, &st), 0);
EXPECT_EQ_INT((int)(st.st_mode & (S_ISUID | S_ISGID | S_ISVTX)), 0);
EXPECT_EQ_INT((int)(st.st_mode & 0777), 0644);
FILE* stream = fopen(path, "rb");
char buf[16] = {0};
EXPECT_NOT_NULL(stream);
// cppcheck-suppress knownConditionTrueFalse
if (stream) {
size_t nread = fread(buf, 1, sizeof(buf) - 1, stream);
fclose(stream);
EXPECT_EQ_INT((int)nread, (int)strlen(new_content));
}
EXPECT_EQ_STR(buf, new_content);
unlink(path);
rmdir(root);
}
static void test_inplace_overwrite_metadata_strips_special_bits() {
const char* root = "test_inplace_meta_tmp";
const char* path = "test_inplace_meta_tmp/meta.txt";
const char* source = "test_inplace_meta_source.txt";
unlink(path);
unlink(source);
rmdir(root);
EXPECT_EQ_INT(mkdir(root, 0700), 0);
/* Existing destination with setuid+sticky set. */
int fd = open(path, O_WRONLY | O_CREAT | O_TRUNC | O_CLOEXEC, 0644);
EXPECT_TRUE(fd >= 0);
// cppcheck-suppress knownConditionTrueFalse
if (fd < 0) {
rmdir(root);
return;
}
EXPECT_EQ_INT((int)write(fd, "olddata", 7), 7);
EXPECT_EQ_INT(fchmod(fd, S_ISUID | S_ISVTX | 0755), 0);
EXPECT_EQ_INT(close(fd), 0);
/* Build source metadata carrying a plain executable mode (no specials). */
EXPECT_TRUE(file_write_to_disk(source, "source", 6, false, false));
EXPECT_EQ_INT(chmod(source, 0755), 0);
struct stat source_st;
EXPECT_EQ_INT(stat(source, &source_st), 0);
File* f = file_create("meta.txt");
EXPECT_NOT_NULL(f);
const char* new_content = "meta";
f->data->data = malloc(strlen(new_content));
EXPECT_NOT_NULL(f->data->data);
memcpy(f->data->data, new_content, strlen(new_content));
f->data->size = strlen(new_content);
f->metadata = file_metadata_create(&source_st);
EXPECT_NOT_NULL(f->metadata);
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
cfg->inplace = true;
EXPECT_TRUE(file_save_to_disk(root, f, cfg));
file_destroy(f);
config_delete(cfg);
unlink(source);
struct stat st;
EXPECT_EQ_INT(stat(path, &st), 0);
/* Metadata-derived mode is applied and never includes setuid/setgid/sticky. */
EXPECT_EQ_INT((int)(st.st_mode & (S_ISUID | S_ISGID | S_ISVTX)), 0);
EXPECT_EQ_INT((int)(st.st_mode & 0777), 0755);
unlink(path);
rmdir(root);
}
static void test_inplace_overwrite_truncates_shorter_payload() {
const char* root = "test_inplace_trunc_tmp";
const char* path = "test_inplace_trunc_tmp/big.txt";
unlink(path);
rmdir(root);
EXPECT_EQ_INT(mkdir(root, 0700), 0);
const char* old_content = "0123456789abcdef"; /* 16 bytes */
EXPECT_TRUE(file_write_to_disk(path, old_content, strlen(old_content), false, false));
File* f = file_create("big.txt");
EXPECT_NOT_NULL(f);
const char* new_content = "hi";
f->data->data = malloc(strlen(new_content));
EXPECT_NOT_NULL(f->data->data);
memcpy(f->data->data, new_content, strlen(new_content));
f->data->size = strlen(new_content);
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
cfg->inplace = true;
EXPECT_TRUE(file_save_to_disk(root, f, cfg));
file_destroy(f);
config_delete(cfg);
/* A shorter payload must truncate the file: no stale trailing bytes. */
struct stat st;
EXPECT_EQ_INT(stat(path, &st), 0);
EXPECT_EQ_INT((int)st.st_size, (int)strlen(new_content));
FILE* stream = fopen(path, "rb");
char buf[32] = {0};
EXPECT_NOT_NULL(stream);
// cppcheck-suppress knownConditionTrueFalse
if (stream) {
size_t nread = fread(buf, 1, sizeof(buf) - 1, stream);
fclose(stream);
EXPECT_EQ_INT((int)nread, (int)strlen(new_content));
}
EXPECT_EQ_STR(buf, new_content);
unlink(path);
rmdir(root);
}
void test_file() { void test_file() {
test_file_create(); test_file_create();
test_file_destroy_null(); test_file_destroy_null();
@@ -635,6 +897,8 @@ void test_file() {
test_file_save_to_disk_existing(); test_file_save_to_disk_existing();
test_file_save_to_disk_ignore_existing(); test_file_save_to_disk_ignore_existing();
test_file_save_to_disk_ignore_existing_entry_types(); test_file_save_to_disk_ignore_existing_entry_types();
test_file_save_to_disk_partial_install();
test_file_save_to_disk_reports_skips();
test_file_write_to_disk_basic(); test_file_write_to_disk_basic();
test_file_write_to_disk_with_fsync(); test_file_write_to_disk_with_fsync();
test_file_write_to_disk_creates_dirs(); test_file_write_to_disk_creates_dirs();
@@ -654,4 +918,7 @@ void test_file() {
test_file_send_single_calls_metadata_and_path(); test_file_send_single_calls_metadata_and_path();
} }
test_file_metadata_create(); test_file_metadata_create();
test_inplace_overwrite_clears_special_mode_bits();
test_inplace_overwrite_metadata_strips_special_bits();
test_inplace_overwrite_truncates_shorter_payload();
} }
+82
View File
@@ -250,6 +250,87 @@ static void test_write_thread_done() {
config_delete(cfg); config_delete(cfg);
} }
typedef struct {
PipelineContextReceiver* context;
File* file;
atomic_bool* done;
atomic_bool* result;
} ByteBudgetEnqueueArg;
static int byte_budget_enqueue_worker(void* arg) {
ByteBudgetEnqueueArg* worker = arg;
bool ok = pipeline_context_receiver_enqueue_file(worker->context, worker->file);
atomic_store(worker->result, ok);
atomic_store(worker->done, true);
return thrd_success;
}
/* A receiver must not buffer more decompressed/copied payload bytes ahead of
the (slow) disk writer than the configured byte budget: an enqueue that
would exceed the budget blocks until the writer releases bytes. */
static void test_receiver_enqueue_byte_budget() {
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
free(cfg->version);
cfg->version = str_dup(PROTOCOL_VERSION);
cfg->send_directory = str_dup("/src");
cfg->receive_root_directory = str_dup("/dst");
cfg->save_to_disk = false;
Queue* q = queue_create(16, file_destroy);
EXPECT_NOT_NULL(q);
PipelineContextReceiver* ctx = pipeline_context_receiver_create(cfg, q, -1, NULL);
EXPECT_NOT_NULL(ctx);
pipeline_context_receiver_set_queue_byte_limit(ctx, 3000);
ctx->receiver_done = false;
File* first = file_create("budget_file_1");
EXPECT_NOT_NULL(first);
first->data->size = 2000;
EXPECT_TRUE(pipeline_context_receiver_enqueue_file(ctx, first));
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000);
/* Second 2000-byte payload would push the pipeline to 4000 > 3000 budget,
so the enqueue must block until the first payload is released. */
File* second = file_create("budget_file_2");
EXPECT_NOT_NULL(second);
second->data->size = 2000;
atomic_bool done;
atomic_bool result;
atomic_init(&done, false);
atomic_init(&result, false);
ByteBudgetEnqueueArg arg = {ctx, second, &done, &result};
thrd_t enqueuer;
EXPECT_EQ_INT(thrd_create(&enqueuer, byte_budget_enqueue_worker, &arg), thrd_success);
/* Give a broken (unbounded) implementation every chance to enqueue. */
struct timespec wait = {0, 200 * 1000000L};
thrd_sleep(&wait, NULL);
EXPECT_FALSE(atomic_load(&done));
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* budget still honored */
/* Simulate the disk writer: dequeue + destroy + release the first file. */
File* drained = queue_dequeue_multithreaded(q, &ctx->mutex, &ctx->condition_not_empty,
&ctx->condition_not_full, &ctx->receiver_done);
EXPECT_NOT_NULL(drained);
file_destroy(drained);
pipeline_context_receiver_note_bytes_released(ctx, 2000);
EXPECT_EQ_INT((int)ctx->queued_bytes, 0);
EXPECT_EQ_INT(thrd_join(enqueuer, NULL), thrd_success);
EXPECT_TRUE(atomic_load(&done));
EXPECT_TRUE(atomic_load(&result));
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* second payload now in flight */
/* Tear down: the second file is still queued and is freed by queue_destroy. */
mtx_destroy(&ctx->mutex);
cnd_destroy(&ctx->condition_not_full);
cnd_destroy(&ctx->condition_not_empty);
free(ctx);
queue_destroy(q);
config_delete(cfg);
}
void test_multiprocessing() { void test_multiprocessing() {
test_sender_create_destroy(); test_sender_create_destroy();
test_receiver_create_destroy(); test_receiver_create_destroy();
@@ -261,4 +342,5 @@ void test_multiprocessing() {
test_receive_thread_failure_wakes_writer(); test_receive_thread_failure_wakes_writer();
} }
test_write_thread_done(); test_write_thread_done();
test_receiver_enqueue_byte_budget();
} }
+250
View File
@@ -1,14 +1,18 @@
#include "test_server.h" #include "test_server.h"
#include "config.h" #include "config.h"
#include "delta.h"
#include "file.h" #include "file.h"
#include "protocol.h" #include "protocol.h"
#include "test_utils.h" #include "test_utils.h"
#include "utils.h" #include "utils.h"
#include <fcntl.h>
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
#include <sys/socket.h> #include <sys/socket.h>
#include <sys/stat.h>
#include <sys/wait.h> #include <sys/wait.h>
#include <time.h>
#include <unistd.h> #include <unistd.h>
#include "receiver.h" #include "receiver.h"
@@ -210,6 +214,249 @@ static void test_receive_incremental_check_rejects_invalid_nanoseconds() {
config_delete(cfg); config_delete(cfg);
} }
static char* make_check_root(const char* tag) {
char tmpl[128];
snprintf(tmpl, sizeof(tmpl), "/tmp/fastsync_%s_XXXXXX", tag);
char* path = str_dup(tmpl);
if (!path)
return NULL;
if (!mkdtemp(path)) {
free(path);
return NULL;
}
return path;
}
static void write_check_file(const char* dir, const char* name, const char* content) {
char path[1024];
snprintf(path, sizeof(path), "%s/%s", dir, name);
int fd = open(path, O_WRONLY | O_CREAT | O_TRUNC, 0644);
if (fd >= 0) {
size_t len = strlen(content);
if (write(fd, content, len) != (ssize_t)len) {
/* intentionally ignored in tests */
}
close(fd);
}
}
/* Issue #255: a same-size/mtime match is decided from metadata alone, so the
receiver answers STATUS_OK (skip) and never asks for a data body. */
static void test_incremental_check_quick_skip_by_mtime() {
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
char* root = make_check_root("qskip");
EXPECT_NOT_NULL(root);
cfg->receive_root_directory = str_dup(root);
write_check_file(root, "file.txt", "0123456789abcdef");
char path[1024];
snprintf(path, sizeof(path), "%s/file.txt", root);
struct stat st;
EXPECT_EQ_INT(stat(path, &st), 0);
int p[2];
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
io_set_fds(p[0], p[1]);
io_set_bwlimit(0);
pid_t pid = fork();
if (pid == 0) {
alarm(30);
close(p[1]);
io_set_fds(p[0], p[0]);
bool skipped = false;
File* file = receive_incremental_check(p[0], cfg, &skipped);
bool ok = file == NULL && skipped;
file_destroy(file);
config_delete(cfg);
close(p[0]);
_exit(ok ? 0 : 1);
} else {
close(p[0]);
io_set_fds(p[1], p[1]);
EXPECT_TRUE(send_str(p[1], "file.txt"));
unsigned long long size = (unsigned long long)st.st_size;
long long mtime = (long long)st.st_mtime;
long long mtime_nsec = 0;
#ifdef __linux__
mtime_nsec = (long long)st.st_mtim.tv_nsec;
#endif
EXPECT_TRUE(send_n_data(p[1], &size, sizeof(size)));
EXPECT_TRUE(send_n_data(p[1], &mtime, sizeof(mtime)));
EXPECT_TRUE(send_n_data(p[1], &mtime_nsec, sizeof(mtime_nsec)));
Status s;
EXPECT_TRUE(receive_status(p[1], &s));
EXPECT_EQ_INT(s, STATUS_OK);
int status;
waitpid(pid, &status, 0);
close(p[1]);
config_delete(cfg);
unlink(path);
rmdir(root);
free(root);
EXPECT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0);
}
}
/* Issue #255: a size mismatch cannot be a skip, so the receiver answers
STATUS_NEXT and consumes the full data body that follows. */
static void test_incremental_check_size_mismatch_full_transfer() {
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
char* root = make_check_root("qnext");
EXPECT_NOT_NULL(root);
cfg->receive_root_directory = str_dup(root);
write_check_file(root, "file.txt", "0123456789abcdef");
char path[1024];
snprintf(path, sizeof(path), "%s/file.txt", root);
struct stat st;
EXPECT_EQ_INT(stat(path, &st), 0);
int p[2];
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
io_set_fds(p[0], p[1]);
io_set_bwlimit(0);
pid_t pid = fork();
if (pid == 0) {
alarm(30);
close(p[1]);
io_set_fds(p[0], p[0]);
bool skipped = false;
File* file = receive_incremental_check(p[0], cfg, &skipped);
bool ok = file != NULL && !skipped && file->path != NULL && strcmp(file->path, "file.txt") == 0;
file_destroy(file);
config_delete(cfg);
close(p[0]);
_exit(ok ? 0 : 1);
} else {
close(p[0]);
io_set_fds(p[1], p[1]);
EXPECT_TRUE(send_str(p[1], "file.txt"));
unsigned long long size = (unsigned long long)st.st_size + 1;
long long mtime = (long long)st.st_mtime;
long long mtime_nsec = 0;
#ifdef __linux__
mtime_nsec = (long long)st.st_mtim.tv_nsec;
#endif
EXPECT_TRUE(send_n_data(p[1], &size, sizeof(size)));
EXPECT_TRUE(send_n_data(p[1], &mtime, sizeof(mtime)));
EXPECT_TRUE(send_n_data(p[1], &mtime_nsec, sizeof(mtime_nsec)));
Status s;
EXPECT_TRUE(receive_status(p[1], &s));
EXPECT_EQ_INT(s, STATUS_NEXT);
Data* body = data_create_reserve(8);
EXPECT_NOT_NULL(body);
body->data = malloc(8);
EXPECT_NOT_NULL(body->data);
memcpy(body->data, "replaced", 8);
body->size = 8;
EXPECT_TRUE(send_data(p[1], body));
data_destroy(body);
int status;
waitpid(pid, &status, 0);
close(p[1]);
config_delete(cfg);
unlink(path);
rmdir(root);
free(root);
EXPECT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0);
}
}
/* Issue #256: when a received delta claims a result above the whole-file cap,
receive_delta_file must mark the operation failed so the caller aborts with
STATUS_ERROR instead of emitting STATUS_NEXT and waiting for a body that
never arrives. */
static void test_incremental_check_delta_oversize_reports_failure() {
Config* cfg = config_create();
EXPECT_NOT_NULL(cfg);
char* root = make_check_root("qdelta");
EXPECT_NOT_NULL(root);
cfg->receive_root_directory = str_dup(root);
cfg->use_delta = true;
char content[20000];
memset(content, 'a', sizeof(content));
content[sizeof(content) - 1] = '\0';
write_check_file(root, "file.txt", content);
char path[1024];
snprintf(path, sizeof(path), "%s/file.txt", root);
struct stat st;
EXPECT_EQ_INT(stat(path, &st), 0);
EXPECT_EQ_INT((int)st.st_size, 19999);
int p[2];
EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0);
io_set_fds(p[0], p[1]);
io_set_bwlimit(0);
pid_t pid = fork();
if (pid == 0) {
alarm(30);
close(p[1]);
io_set_fds(p[0], p[0]);
bool skipped = false;
File* file = receive_incremental_check(p[0], cfg, &skipped);
bool ok = file == NULL && !skipped;
if (ok)
send_status(p[0], STATUS_ERROR); /* mirror the server error path */
file_destroy(file);
config_delete(cfg);
close(p[0]);
_exit(ok ? 0 : 1);
} else {
close(p[0]);
io_set_fds(p[1], p[1]);
EXPECT_TRUE(send_str(p[1], "file.txt"));
unsigned long long size = (unsigned long long)st.st_size;
long long mtime = 1; /* different from the file mtime: force a transfer */
long long mtime_nsec = 0;
EXPECT_TRUE(send_n_data(p[1], &size, sizeof(size)));
EXPECT_TRUE(send_n_data(p[1], &mtime, sizeof(mtime)));
EXPECT_TRUE(send_n_data(p[1], &mtime_nsec, sizeof(mtime_nsec)));
Status s;
EXPECT_TRUE(receive_status(p[1], &s));
EXPECT_EQ_INT(s, STATUS_DELTA_SIGNATURE);
Data* sig_data = receive_data(p[1]);
EXPECT_NOT_NULL(sig_data);
DeltaSignature* sig = delta_signature_deserialize(sig_data);
EXPECT_NOT_NULL(sig);
delta_signature_destroy(sig);
data_destroy(sig_data);
/* Send a delta whose claimed output size exceeds the whole-file cap. */
Delta delta;
memset(&delta, 0, sizeof(delta));
delta.new_file_size = MAX_RECEIVE_WHOLE_FILE_SIZE + 1;
Data* bogus = delta_serialize(&delta);
EXPECT_NOT_NULL(bogus);
EXPECT_TRUE(send_status(p[1], STATUS_DELTA_DATA));
EXPECT_TRUE(bogus != NULL && send_data(p[1], bogus));
data_destroy(bogus);
/* The receiver must answer with an error, never with STATUS_NEXT. */
EXPECT_TRUE(receive_status(p[1], &s));
EXPECT_EQ_INT(s, STATUS_ERROR);
int status;
waitpid(pid, &status, 0);
close(p[1]);
config_delete(cfg);
unlink(path);
rmdir(root);
free(root);
EXPECT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0);
}
}
void test_server() { void test_server() {
if (!is_running_under_valgrind()) { if (!is_running_under_valgrind()) {
test_receive_files_finished(); test_receive_files_finished();
@@ -217,5 +464,8 @@ void test_server() {
test_receive_files_abort(); test_receive_files_abort();
test_receive_manifest_rejects_traversal(); test_receive_manifest_rejects_traversal();
test_receive_incremental_check_rejects_invalid_nanoseconds(); test_receive_incremental_check_rejects_invalid_nanoseconds();
test_incremental_check_quick_skip_by_mtime();
test_incremental_check_size_mismatch_full_transfer();
test_incremental_check_delta_oversize_reports_failure();
} }
} }