From 6ad3887aa10240de4f82393a25c7f9acb2a3d215 Mon Sep 17 00:00:00 2001 From: TapTap Date: Mon, 7 Sep 2026 15:43:29 +0200 Subject: [PATCH 1/3] feat: --append / --append-verify tail-only resume Implement rsync's append modes: when an existing destination file is SHORTER than the source, the receiver negotiates a resume offset and only the tail is transferred; the full file (retained prefix + tail) is rebuilt and installed through the normal atomic store path, so the result is byte-identical to the source whenever the prefix matches. - --append: sends the tail without content-verifying the retained prefix (rsync parity; the documented prefix-trust risk). - --append-verify: verifies the retained prefix against the source's prefix xxHash64 before appending and, on a mismatch, falls back to a clean full transfer (never a corrupt prefix+tail blend). New wire frames STATUS_APPEND / STATUS_APPEND_SIG / STATUS_APPEND_OK / STATUS_APPEND_DATA; PROTOCOL_VERSION bumped 2.9.0 -> 2.10.0 (peers must match). Both flags imply --incremental and are incompatible with -s (chunk serialization) and --whole-file (rejected up front). Respects --inplace, --partial/--partial-dir and --delay-updates via the shared store engine. --- src/client/client_cli.c | 13 ++ src/client/client_send.c | 127 +++++++++++++- src/client/client_validation.c | 24 ++- src/client/usage.c | 4 + src/shared/config.c | 4 + src/shared/config.h | 2 +- src/shared/file_receive.c | 303 ++++++++++++++++++++++++++++----- src/shared/protocol.c | 8 + src/shared/protocol.h | 16 +- src/shared/utils.c | 12 ++ src/shared/utils.h | 8 + 11 files changed, 470 insertions(+), 51 deletions(-) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 995a5cc..4a8a91e 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -435,6 +435,8 @@ static const OptionEntry OPTION_TABLE[] = { {"--copy-unsafe-links", NULL, OPT_FLAG, offsetof(Config, copy_unsafe_links)}, {"--sparse", "-S", OPT_FLAG, offsetof(Config, preserve_sparse)}, {"--inplace", NULL, OPT_FLAG, offsetof(Config, inplace)}, + {"--append", NULL, OPT_FLAG, offsetof(Config, append)}, + {"--append-verify", NULL, OPT_FLAG, offsetof(Config, append_verify)}, {"--fsync", NULL, OPT_FLAG, offsetof(Config, use_fsync)}, {"--checksum", NULL, OPT_FLAG, offsetof(Config, checksum)}, {"--8-bit-output", "-8", OPT_FLAG, offsetof(Config, eight_bit_output)}, @@ -1102,6 +1104,17 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args, config->use_delta = true; } + /* --append / --append-verify resume a shorter existing destination file. + * The receiver must run the per-file STATUS_CHECK handshake to learn the + * destination length and reply STATUS_APPEND, so an append mode forces + * --incremental on (exactly like the basis-dir options: the handshake is + * required, not optional). The resume itself is a dedicated tail-only + * exchange, not the block delta, so no delta implication is made. When both + * spelling are given the safer --append-verify semantics win. */ + if (config->append || config->append_verify) { + config->use_incremental = true; + } + /* Incremental and delta transfers need metadata unless the user disabled it. */ if ((config->use_incremental || config->use_delta) && !config->use_metadata && !config->metadata_explicitly_disabled) { diff --git a/src/client/client_send.c b/src/client/client_send.c index 03cb828..3fb59aa 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -717,8 +717,10 @@ static bool scan_paths_only(const Config* config, const ScannerOptions* options, } static int incremental_check(Client* client, File* file, const Config* config, - DeltaSignature** out_sig) { + DeltaSignature** out_sig, unsigned long long* resume_offset) { *out_sig = NULL; + if (resume_offset) + *resume_offset = 0; if (!send_status(client->file_descriptor, STATUS_CHECK)) return -1; if (!send_str(client->file_descriptor, file_wire_path(file))) @@ -765,6 +767,19 @@ static int incremental_check(Client* client, File* file, const Config* config, *out_sig = sig; return 2; } + if (s == STATUS_APPEND) { + /* --append / --append-verify tail resume: the receiver found an existing + destination SHORTER than the source and wants only the tail from this + offset (the bytes it already holds). */ + unsigned long long offset; + if (!receive_n_data(client->file_descriptor, &offset, sizeof(offset))) { + send_status(client->file_descriptor, STATUS_ERROR); + return -1; + } + if (resume_offset) + *resume_offset = offset; + return 3; + } if (s != STATUS_NEXT) { log_message(LOG_LEVEL_ERROR, "Unexpected server status"); send_status(client->file_descriptor, STATUS_ERROR); @@ -813,6 +828,89 @@ static int send_delta(Client* client, File* file, DeltaSignature* sig, Config* c return ok ? 0 : -1; } +/* --append / --append-verify tail resume. The receiver learned the existing + * destination is SHORTER than the source and replied STATUS_APPEND with the + * resume offset (prefix bytes it already holds). For plain --append we send + * the tail immediately (the prefix is not content-verified, matching rsync). + * For --append-verify we first send the source prefix xxHash64; the receiver + * compares it to the retained prefix and replies STATUS_APPEND_OK (send the + * tail) or STATUS_NEXT (prefix mismatch -> full transfer, never corrupt). + * Returns 0 on success, 1 when a full transfer was done instead, -1 on error. */ +static int send_append(const Client* client, File* file, Config* config, + unsigned long long offset) { + int fd = client->file_descriptor; + const unsigned long long fsize = file->data->size; + if (offset >= fsize) { + send_status(fd, STATUS_ERROR); + return -1; + } + size_t off = (size_t)offset; + size_t tail_len = (size_t)(fsize - off); + int compression_level = config->use_compression ? config->compression_level : 0; + int skip_count = config->skip_compress_set ? config->skip_compress_count : -1; + bool compress = compression_level > 0 && + !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes, + skip_count); + + /* --append-verify: exchange the source prefix checksum and await the verdict. */ + if (config->append_verify) { + uint64_t prefix_hash = delta_xxhash64(file->data->data, off); + if (!send_status(fd, STATUS_APPEND_SIG) || !send_n_data(fd, &prefix_hash, sizeof(prefix_hash))) + return -1; + Status resp; + if (!receive_status(fd, &resp)) + return -1; + if (resp == STATUS_NEXT) { + /* Retained prefix does not match the source: fall back to the atomic full + transfer (byte-identical, never a corrupt prefix+tail blend). */ + int rc = file_send_single_calls_with_skip(file, fd, config->use_metadata, compression_level, + false, config->skip_compress_suffixes, skip_count, + config->compression_threads) + ? 1 + : -1; + return rc; + } + if (resp != STATUS_APPEND_OK) { + send_status(fd, STATUS_ERROR); + return -1; + } + } + + if (!send_status(fd, STATUS_APPEND_DATA)) { + return -1; + } + if (config->use_metadata && !metadata_send(fd, file->metadata)) { + return -1; + } + bool ok; + if (compress) { + /* Compression needs an owned copy of the tail to compress. */ + Data* tail = data_create_empty(tail_len); + if (!tail) { + send_status(fd, STATUS_ERROR); + return -1; + } + memcpy(tail->data, (const char*)file->data->data + off, tail_len); + Data* comp = data_compress_with_threads(tail, compression_level, config->compression_threads); + data_destroy(tail); + if (!comp) { + send_status(fd, STATUS_ERROR); + return -1; + } + ok = send_data(fd, comp); + data_destroy(comp); + } else { + /* Uncompressed: send directly from the source buffer (no per-file copy; + send_data is synchronous, so the view outlives the call). */ + Data tail_view; + tail_view.data = (char*)file->data->data + off; + tail_view.size = tail_len; + tail_view.protocol_charge = 0; + ok = send_data(fd, &tail_view); + } + return ok ? 0 : -1; +} + // Send a single file directly (non-incremental path). static bool send_file_direct(File* file, int fd, bool use_metadata, int compression_level, const Config* config) { @@ -867,7 +965,8 @@ static int send_single_file(Client* client, File* file, Config* config, bool use // Incremental path: use sendfile for the actual data if enabled and no compression if (use_sendfile) { DeltaSignature* sig = NULL; - int rc = incremental_check(client, file, config, &sig); + unsigned long long resume_offset = 0; + int rc = incremental_check(client, file, config, &sig, &resume_offset); if (rc == 1) { log_info_message(LOG_INFO_SKIP, "Skipping unchanged %s", file->path); delta_signature_destroy(sig); @@ -877,6 +976,16 @@ static int send_single_file(Client* client, File* file, Config* config, bool use delta_signature_destroy(sig); return -1; } + // rc == 3: append resume (tail-only) -- send_append uses the data path. + if (rc == 3) { + delta_signature_destroy(sig); + int arc = send_append(client, file, config, resume_offset); + if (arc == 1) { + log_info_message(LOG_INFO_COPY, "Append prefix mismatch; full transfer of %s", file->path); + return 0; + } + return arc == 0 ? 0 : -1; + } // rc == 0: unchanged file, skip // rc == 2: server sent delta signature but sendfile doesn't support delta delta_signature_destroy(sig); @@ -896,7 +1005,8 @@ static int send_single_file(Client* client, File* file, Config* config, bool use // Incremental path with single_calls (supports compression and delta) DeltaSignature* sig = NULL; - int rc = incremental_check(client, file, config, &sig); + unsigned long long resume_offset = 0; + int rc = incremental_check(client, file, config, &sig, &resume_offset); if (rc < 0) { delta_signature_destroy(sig); return -1; @@ -906,6 +1016,17 @@ static int send_single_file(Client* client, File* file, Config* config, bool use delta_signature_destroy(sig); return 1; } + if (rc == 3) { + /* --append / --append-verify tail resume. send_append reports 1 when the + verified prefix mismatched and a full transfer was sent instead. */ + delta_signature_destroy(sig); + int arc = send_append(client, file, config, resume_offset); + if (arc == 1) { + log_info_message(LOG_INFO_COPY, "Append prefix mismatch; full transfer of %s", file->path); + return 0; + } + return arc == 0 ? 0 : -1; + } if (rc == 2 && config->use_delta && !config->whole_file) { int drc = send_delta(client, file, sig, config); delta_signature_destroy(sig); diff --git a/src/client/client_validation.c b/src/client/client_validation.c index f068699..c307b02 100644 --- a/src/client/client_validation.c +++ b/src/client/client_validation.c @@ -51,14 +51,26 @@ bool validate_config(const Config* config) { log_message(LOG_LEVEL_ERROR, "--delta cannot be combined with -f (sendfile)"); return false; } - if (config->log_file_format && !config->log_file) { - log_message(LOG_LEVEL_ERROR, "--log-file-format requires --log-file"); + /* --append / --append-verify resume a shorter existing destination by + transmitting only the tail. The resume needs the per-file STATUS_CHECK + handshake (so the dest length is learned), which chunk serialization -s + disables; and whole-file is the opposite intent (send everything), so the + two would silently make the resume pointless. Both are rejected up front + rather than silently degrading to a full transfer. */ + if ((config->append || config->append_verify) && config->use_chunk_serialization) { + log_message(LOG_LEVEL_ERROR, + "--append/--append-verify require the per-file incremental check and cannot be " + "combined with -s (chunk serialization)"); return false; } - if (config->append || config->append_verify) { - fprintf( - stderr, - "Error: --append and --append-verify are not supported yet; refusing to ignore option\n"); + if ((config->append || config->append_verify) && config->whole_file) { + log_message(LOG_LEVEL_ERROR, + "--append/--append-verify are incompatible with --whole-file (which forces a " + "full transfer)"); + return false; + } + if (config->log_file_format && !config->log_file) { + log_message(LOG_LEVEL_ERROR, "--log-file-format requires --log-file"); return false; } if (config->use_tls) { diff --git a/src/client/usage.c b/src/client/usage.c index f4ee960..1628c83 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -159,6 +159,10 @@ void print_usage(void) { printf(" --copy-unsafe-links Only transform unsafe symlinks into referent files\n"); printf(" -S, --sparse Handle sparse files efficiently\n"); printf(" --inplace Update files in-place (no temp+rename)\n"); + printf(" --append Resume a shorter destination by appending only its tail\n"); + printf(" (prefix is not verified; requires --incremental)\n"); + printf(" --append-verify Like --append, but verifies the retained prefix checksum\n"); + printf(" before appending (falls back to a full transfer on mismatch)\n"); printf(" --fsync Fsync every written file before publication\n"); printf(" --compress-level Compression level (default: 5)\n"); printf(" --zl Alias for --compress-level\n"); diff --git a/src/shared/config.c b/src/shared/config.c index 49c25f0..dfd86a8 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -174,6 +174,10 @@ static bool validate_received_config(const Config* config) { valid_wire_bool(config->checksum) && valid_wire_bool(config->eight_bit_output) && config_has_valid_delete_timing(config) && !(config->skip_compress_set && config->use_chunk_serialization) && + /* --append / --append-verify tail resume needs the per-file check, + which chunk serialization -s disables: reject on the receiver too + so a -s sender cannot negotiate an inert append mode. */ + !((config->append || config->append_verify) && config->use_chunk_serialization) && (!config->use_compression || (config->compression_level >= 1 && config->compression_level <= 22)) && config->chunk_size > 0 && config->chunk_size <= MAX_CHUNK_SIZE && diff --git a/src/shared/config.h b/src/shared/config.h index 34b1552..b3b2d6c 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -231,7 +231,7 @@ typedef struct Config { DelayUpdatesContext* delay_context; } Config; -#define PROTOCOL_VERSION "2.9.0" +#define PROTOCOL_VERSION "2.10.0" #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) /* Upper bound on total basis-dir entries (rsync caps --link-dest at 20). */ #define MAX_BASIS_DIRS 64 diff --git a/src/shared/file_receive.c b/src/shared/file_receive.c index a937787..0c182a1 100644 --- a/src/shared/file_receive.c +++ b/src/shared/file_receive.c @@ -983,6 +983,50 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p return basis; } +/* Read the remainder of a full-file transfer after the receiver has already + * sent STATUS_NEXT: receive the metadata frame (when enabled) followed by the + * data frame, and return an owned File. Shared by the plain full-transfer path + * and the --append-verify prefix-mismatch fallback (a clean full transfer + * instead of a corrupt prefix+tail blend). */ +static File* receive_full_file(int fd, const Config* config, const char* path) { + File* file = file_create(path); + if (!file) + return NULL; + if (config->use_metadata) { + int meta_ok = 1; + file->metadata = metadata_receive(fd, &meta_ok); + if (!meta_ok) { + file_destroy(file); + return NULL; + } + } + Data* file_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE); + if (file_data == NULL) { + file_destroy(file); + return NULL; + } + if (config->use_compression && + !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes, + config->skip_compress_set ? config->skip_compress_count + : -1)) { + Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE); + data_destroy(file_data); + if (uncompressed == NULL) { + file_destroy(file); + return NULL; + } + if (uncompressed->size > MAX_FILE_DATA_SIZE) { + data_destroy(uncompressed); + file_destroy(file); + return NULL; + } + file_data = uncompressed; + } + data_destroy(file->data); + file->data = file_data; + return file; +} + File* receive_incremental_check(int fd, const Config* config, bool* skipped) { if (!config || !skipped) { send_status(fd, STATUS_ERROR); @@ -1193,6 +1237,224 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { basis_match_free(&basis); } + /* ---- --append / --append-verify tail resume ---- + * When the existing destination file is SHORTER than the source, an append + * mode resumes it by negotiating a resume offset (the prefix length already + * present) from the receiver and transferring ONLY the tail. The receiver + * then reconstructs the full file (prefix + tail) and installs it through the + * normal atomic store path, so the result is byte-identical to the source. + * This takes precedence over block delta (a growing file is cheapest as a + * pure tail), and falls through to delta/full only when no shorter old file + * makes a resume possible. */ + bool append_resume = (config->append || config->append_verify) && has_old_file && + append_resume_eligible(old_size, check_size); + if (append_resume) { + /* Ensure the retained prefix (== the whole, shorter destination file) is + in memory; it is needed both to rebuild the full file and, for + --append-verify, to checksum it. A load failure is not fatal: the + resume is simply not possible and we fall through to the other paths. */ + if (old_data == NULL && old_size > 0 && old_size <= MAX_RECEIVE_WHOLE_FILE_SIZE && + old_size <= SIZE_MAX) { + old_data = protocol_alloc((size_t)old_size); + if (old_data) { + size_t got = 0; + while (got < (size_t)old_size) { + ssize_t n = read(old_fd, (char*)old_data + got, (size_t)old_size - got); + if (n <= 0) { + free(old_data); + old_data = NULL; + break; + } + got += (size_t)n; + } + } + } + if (old_data != NULL || old_size == 0) { + if (!send_status(fd, STATUS_APPEND) || !send_n_data(fd, &old_size, sizeof(old_size))) { + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + bool verify = config->append_verify; + bool full_fallback = false; + if (verify) { + Status sig_status; + if (!receive_status(fd, &sig_status)) { + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + if (sig_status != STATUS_APPEND_SIG) { + send_status(fd, STATUS_ERROR); + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + uint64_t src_prefix_hash; + if (!receive_n_data(fd, &src_prefix_hash, sizeof(src_prefix_hash))) { + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + /* Compare the retained prefix against the source prefix. A mismatch + must never be silently appended to: fall back to a full transfer so + the result is a byte-identical source copy. */ + uint64_t dst_prefix_hash = + old_size == 0 ? delta_xxhash64("", 0) : delta_xxhash64(old_data, (size_t)old_size); + if (dst_prefix_hash == src_prefix_hash) { + if (!send_status(fd, STATUS_APPEND_OK)) { + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + } else { + if (!send_status(fd, STATUS_NEXT)) { + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + full_fallback = true; + } + } + + if (full_fallback) { + /* Retained prefix differed: receive the sender's full transfer. */ + free(old_data); + old_data = NULL; + close(old_fd); + File* file = receive_full_file(fd, config, check_path); + free(check_path); + free(full_path); + return file; + } + + /* Receive the tail (STATUS_APPEND_DATA + metadata + tail bytes). */ + Status tail_status; + if (!receive_status(fd, &tail_status)) { + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + if (tail_status != STATUS_APPEND_DATA) { + send_status(fd, STATUS_ERROR); + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + FileMetadata* meta = NULL; + if (config->use_metadata) { + int meta_ok = 1; + meta = metadata_receive(fd, &meta_ok); + if (!meta_ok) { + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + } + Data* tail = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE); + if (tail == NULL) { + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + if (config->use_compression && + !compression_should_skip_with_suffixes( + check_path, config->skip_compress_suffixes, + config->skip_compress_set ? config->skip_compress_count : -1)) { + Data* uncompressed = data_decompress_limited(tail, MAX_RECEIVE_WHOLE_FILE_SIZE); + data_destroy(tail); + if (uncompressed == NULL) { + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + if (uncompressed->size > MAX_FILE_DATA_SIZE) { + data_destroy(uncompressed); + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + tail = uncompressed; + } + /* The tail must complete the file exactly; anything else is a protocol + violation (never a truncated or overrun file). */ + unsigned long long expected_tail; + if (!append_tail_length(old_size, check_size, &expected_tail) || + tail->size != (size_t)expected_tail) { + send_status(fd, STATUS_ERROR); + data_destroy(tail); + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + size_t full_size = (size_t)check_size; + void* full = protocol_alloc(full_size ? full_size : 1); + if (!full) { + data_destroy(tail); + close(old_fd); + free(full_path); + free(check_path); + free(old_data); + return NULL; + } + if (old_size > 0 && old_data) + memcpy(full, old_data, (size_t)old_size); + if (tail->size > 0) + memcpy((char*)full + old_size, tail->data, tail->size); + data_destroy(tail); + free(old_data); + old_data = NULL; + + File* file = file_create(check_path); + if (!file) { + free(full); + close(old_fd); + free(full_path); + free(check_path); + return NULL; + } + file->metadata = meta; + file->data = data_create(full, full_size); + if (!file->data) { /* data_create already freed full on failure */ + file_destroy(file); + close(old_fd); + free(full_path); + free(check_path); + return NULL; + } + close(old_fd); + free(full_path); + free(check_path); + return file; + } + } + if (try_delta && old_data != NULL) { bool delta_failed = false; File* delta_file = @@ -1256,48 +1518,9 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) { } close(old_fd); - File* file = file_create(check_path); + File* file = receive_full_file(fd, config, check_path); free(check_path); free(full_path); - if (file == NULL) { - return NULL; - } - - if (config->use_metadata) { - int meta_ok = 1; - file->metadata = metadata_receive(fd, &meta_ok); - if (!meta_ok) { - file_destroy(file); - return NULL; - } - } - - Data* file_data = receive_data_limited(fd, MAX_RECEIVE_WHOLE_FILE_SIZE); - if (file_data == NULL) { - file_destroy(file); - return NULL; - } - - if (config->use_compression && - !compression_should_skip_with_suffixes(file->path, config->skip_compress_suffixes, - config->skip_compress_set ? config->skip_compress_count - : -1)) { - Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE); - data_destroy(file_data); - if (uncompressed == NULL) { - file_destroy(file); - return NULL; - } - if (uncompressed->size > MAX_FILE_DATA_SIZE) { - data_destroy(uncompressed); - file_destroy(file); - return NULL; - } - file_data = uncompressed; - } - - data_destroy(file->data); - file->data = file_data; return file; } diff --git a/src/shared/protocol.c b/src/shared/protocol.c index 1d4a421..9729e9b 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -397,6 +397,14 @@ static const char* status_to_string(Status status) { return "CHECK_BATCH"; case STATUS_MKDIR: return "MKDIR"; + case STATUS_APPEND: + return "APPEND"; + case STATUS_APPEND_SIG: + return "APPEND_SIG"; + case STATUS_APPEND_OK: + return "APPEND_OK"; + case STATUS_APPEND_DATA: + return "APPEND_DATA"; default: return "UNKNOWN"; } diff --git a/src/shared/protocol.h b/src/shared/protocol.h index 67a929f..28c6f9f 100644 --- a/src/shared/protocol.h +++ b/src/shared/protocol.h @@ -70,7 +70,21 @@ enum NET_STATUS { STATUS_CHECK_BATCH, /* An explicit directory entry (--dirs): the sender transmits only the path; * the receiver creates the directory below the receive root. */ - STATUS_MKDIR + STATUS_MKDIR, + /* --append / --append-verify tail resume. STATUS_APPEND is sent by the + * receiver after a per-file STATUS_CHECK when the existing destination file + * is SHORTER than the source and an append mode is negotiated: its payload is + * the resume offset (the number of prefix bytes already present), after which + * the sender answers either directly with STATUS_APPEND_DATA (plain --append, + * prefix not verified) or, for --append-verify, first with STATUS_APPEND_SIG + * carrying the xxHash64 of the source prefix; the receiver then replies + * STATUS_APPEND_OK (prefix matched -> sender transmits the tail) or + * STATUS_NEXT (prefix mismatch -> sender falls back to a full transfer). + * STATUS_APPEND_DATA carries the tail bytes (compressed data frame). */ + STATUS_APPEND, + STATUS_APPEND_SIG, + STATUS_APPEND_OK, + STATUS_APPEND_DATA }; void io_set_fds(int read_fd, int write_fd); diff --git a/src/shared/utils.c b/src/shared/utils.c index d2dbf89..83f69f9 100644 --- a/src/shared/utils.c +++ b/src/shared/utils.c @@ -531,3 +531,15 @@ char* path_cat(const char* path1, const char* path2) { new_path[path1_len + path2_len + 1] = '\0'; return new_path; } + +bool append_resume_eligible(unsigned long long old_size, unsigned long long check_size) { + return old_size < check_size; +} + +bool append_tail_length(unsigned long long old_size, unsigned long long check_size, + unsigned long long* tail_out) { + if (!tail_out || !append_resume_eligible(old_size, check_size)) + return false; + *tail_out = check_size - old_size; + return true; +} diff --git a/src/shared/utils.h b/src/shared/utils.h index 4206095..9adfd6d 100644 --- a/src/shared/utils.h +++ b/src/shared/utils.h @@ -53,5 +53,13 @@ void utils_set_authorized_root_fd(int fd); bool has_path_traversal(const char* path); bool utils_valid_batch_path(const char* path); bool format_human_bytes(unsigned long long bytes, char* buffer, size_t buffer_size); +/* --append / --append-verify tail-resume math (pure). A resume is eligible only + when an existing destination file is SHORTER than the source; the tail length + is then the difference. append_resume_eligible answers whether the shorter + file makes a resume possible; append_tail_length additionally returns that + tail length, refusing (false) the degenerate old_size >= check_size case. */ +bool append_resume_eligible(unsigned long long old_size, unsigned long long check_size); +bool append_tail_length(unsigned long long old_size, unsigned long long check_size, + unsigned long long* tail_out); #endif From ffdbb6568ff4e44745e34893b7b889defafa645d Mon Sep 17 00:00:00 2001 From: TapTap Date: Mon, 7 Sep 2026 15:43:33 +0200 Subject: [PATCH 2/3] test: append/append-verify resume coverage - CLI: --append/--append-verify acceptance (parse + imply --incremental, validate) and incompatibility rejection with -s and --whole-file; both removed from the unimplemented reject list. - Config: on-the-wire append/append_verify round-trip. - Unit: append_resume_eligible / append_tail_length pure resume math. - Integration (test_append.py): matching-prefix resume is byte-identical and tail-only (wire bytes << source size); --append with a wrong prefix keeps prefix+tail (rsync parity) while --append-verify detects the mismatch and falls back to a byte-exact full transfer; --append with --inplace and -m. --- tests/integration/test_append.py | 141 +++++++++++++++++++++++++++++++ tests/test_client_cli.c | 98 ++++++++++++++++++++- tests/test_config.c | 44 ++++++++++ tests/test_shared_utils.c | 16 ++++ 4 files changed, 297 insertions(+), 2 deletions(-) create mode 100644 tests/integration/test_append.py diff --git a/tests/integration/test_append.py b/tests/integration/test_append.py new file mode 100644 index 0000000..be84a5a --- /dev/null +++ b/tests/integration/test_append.py @@ -0,0 +1,141 @@ +"""--append / --append-verify tail-resume integration tests. + +A shorter existing destination file is resumed by transferring only the tail: +--append sends it without verifying the retained prefix (rsync parity: a wrong +prefix is kept, so the result can differ from the source), while --append-verify +checksums the retained prefix against the source and, on a mismatch, falls back +to a clean full transfer so the result is always a byte-identical source copy. +""" +import os +import random +import shutil + +from common import ( + TEST_DATA_DIR, + run_client, CountingProxy, clean_dir, + get_dest_received_dir, CLIENT_CMD, +) + +REL = "sub/grow.dat" + + +def _grow_payload(prefix_size, added_size, seed=99): + r = random.Random(seed) + return bytes(r.randbytes(prefix_size)), bytes(r.randbytes(added_size)) + + +class TestAppend: + def _make(self, tag): + source = os.path.join(TEST_DATA_DIR, f"append_{tag}_src") + dest = os.path.join(TEST_DATA_DIR, f"append_{tag}_dst") + clean_dir(source) + shutil.rmtree(dest, ignore_errors=True) + return source, dest + + def _place(self, root, rel, data): + p = os.path.join(root, rel) + os.makedirs(os.path.dirname(p), exist_ok=True) + with open(p, "wb") as fh: + fh.write(data) + return p + + def _read(self, root, rel): + with open(os.path.join(root, rel), "rb") as fh: + return fh.read() + + def _dest_file(self, source, dest, rel): + return os.path.join(get_dest_received_dir(dest, source), rel) + + def test_append_resumes_short_dest_atomically(self, shared_server): + """A shorter dest with a MATCHING prefix is resumed; the reconstructed + file is byte-identical to the source.""" + source, dest = self._make("atomic") + prefix, added = _grow_payload(1 * 1024 * 1024, 64 * 1024) + self._place(source, REL, prefix + added) + self._place(self._dest_file(source, dest, ""), REL, prefix) + + result, _ = run_client(source, dest, flags=["--append"], port=shared_server.port) + assert result.returncode == 0, \ + f"--append failed: {(result.stderr or result.stdout)[:400]}" + assert self._read(self._dest_file(source, dest, ""), REL) == prefix + added + + def test_append_verify_matching_prefix_succeeds(self, shared_server): + source, dest = self._make("verify_ok") + prefix, added = _grow_payload(512 * 1024, 32 * 1024) + self._place(source, REL, prefix + added) + self._place(self._dest_file(source, dest, ""), REL, prefix) + + result, _ = run_client(source, dest, flags=["--append-verify"], port=shared_server.port) + assert result.returncode == 0, \ + f"--append-verify failed: {(result.stderr or result.stdout)[:400]}" + assert self._read(self._dest_file(source, dest, ""), REL) == prefix + added + + def test_append_sends_only_tail(self, shared_server): + """Sorted transfer moves only the tail: wire bytes stay well below the + full source size (incompressible payload, no -c).""" + source, dest = self._make("tail") + prefix, added = _grow_payload(4 * 1024 * 1024, 8 * 1024, seed=7) + full = prefix + added + self._place(source, REL, full) + self._place(self._dest_file(source, dest, ""), REL, prefix) + + proxy = CountingProxy(shared_server.port) + cmd = (CLIENT_CMD + ["--source-dir", source, "--dest-dir", dest, + "--save-to-disk", "--server-port", str(proxy.port), "--append"]) + result = proxy.run(cmd) + assert result.returncode == 0, \ + f"--append failed: {(result.stderr or result.stdout)[:400]}" + assert self._read(self._dest_file(source, dest, ""), REL) == full + assert proxy.client_to_server < full.__len__() // 2, \ + f"expected a tail-only transfer, sent {proxy.client_to_server}B for {full.__len__()}B" + + def test_plain_append_wrong_prefix_is_rsync_parity(self, shared_server): + """--append does NOT verify the retained prefix: a wrong prefix is kept, + so the result is prefix+tail (differs from the source). This is the + documented rsync-parity risk of plain --append.""" + source, dest = self._make("plain_wrong") + correct_prefix, added = _grow_payload(256 * 1024, 32 * 1024, seed=1) + wrong_prefix = bytes(b ^ 0xFF for b in correct_prefix) + self._place(source, REL, correct_prefix + added) + self._place(self._dest_file(source, dest, ""), REL, wrong_prefix) + + result, _ = run_client(source, dest, flags=["--append"], port=shared_server.port) + assert result.returncode == 0 + assert self._read(self._dest_file(source, dest, ""), REL) == wrong_prefix + added + + def test_append_verify_wrong_prefix_never_corrupts(self, shared_server): + """--append-verify detects the retained prefix mismatch and falls back to + a full transfer, so the result is a byte-identical source copy.""" + source, dest = self._make("verify_wrong") + correct_prefix, added = _grow_payload(256 * 1024, 32 * 1024, seed=2) + wrong_prefix = bytes(b ^ 0xFF for b in correct_prefix) + self._place(source, REL, correct_prefix + added) + self._place(self._dest_file(source, dest, ""), REL, wrong_prefix) + + result, _ = run_client(source, dest, flags=["--append-verify"], port=shared_server.port) + assert result.returncode == 0, \ + f"--append-verify mismatch fallback failed: {(result.stderr or result.stdout)[:400]}" + assert self._read(self._dest_file(source, dest, ""), REL) == correct_prefix + added + + def test_append_with_inplace(self, shared_server): + source, dest = self._make("inplace") + prefix, added = _grow_payload(128 * 1024, 16 * 1024, seed=3) + self._place(source, REL, prefix + added) + self._place(self._dest_file(source, dest, ""), REL, prefix) + + result, _ = run_client(source, dest, flags=["--append", "--inplace"], + port=shared_server.port) + assert result.returncode == 0, \ + f"--append --inplace failed: {(result.stderr or result.stdout)[:400]}" + assert self._read(self._dest_file(source, dest, ""), REL) == prefix + added + + def test_append_multithreaded(self, shared_server): + source, dest = self._make("mthread") + prefix, added = _grow_payload(512 * 1024, 32 * 1024, seed=4) + self._place(source, REL, prefix + added) + self._place(self._dest_file(source, dest, ""), REL, prefix) + + result, _ = run_client(source, dest, flags=["--append", "-m"], port=shared_server.port) + assert result.returncode == 0, \ + f"--append -m failed: {(result.stderr or result.stdout)[:400]}" + assert self._read(self._dest_file(source, dest, ""), REL) == prefix + added \ No newline at end of file diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index 5882bae..cc38b5a 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -773,8 +773,6 @@ static void test_parse_args_rejects_unimplemented_options() { "--xattrs", "-D", "--devices", - "--append", - "--append-verify", "--delete-excluded", "--max-delete", "--prune-empty-dirs", @@ -1855,8 +1853,104 @@ static void test_parse_args_max_delete_inert_without_delete() { config_delete(cfg); } +/* --append is accepted and implies the per-file incremental check a tail resume + * needs; it validates cleanly on its own. */ +static void test_parse_args_append() { + Config* cfg = config_create(); + cfg->send_directory = str_dup("/src"); + cfg->receive_root_directory = str_dup("/dst"); + char* argv[] = {"fastsync", "--append", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0); + EXPECT_TRUE(cfg->append); + EXPECT_FALSE(cfg->append_verify); + EXPECT_TRUE(cfg->use_incremental); + EXPECT_TRUE(validate_config(cfg)); + config_delete(cfg); +} + +static void test_parse_args_append_verify() { + Config* cfg = config_create(); + cfg->send_directory = str_dup("/src"); + cfg->receive_root_directory = str_dup("/dst"); + char* argv[] = {"fastsync", "--append-verify", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 4, argv, positional_args, &positional_count), 0); + EXPECT_TRUE(cfg->append_verify); + EXPECT_FALSE(cfg->append); + EXPECT_TRUE(cfg->use_incremental); + EXPECT_TRUE(validate_config(cfg)); + config_delete(cfg); +} + +/* Both spellings are accepted; the safer --append-verify semantics win on the + * wire (the sender checks append_verify first), so neither flag is silently + * dropped but the run is still valid. */ +static void test_parse_args_append_both() { + Config* cfg = config_create(); + cfg->send_directory = str_dup("/src"); + cfg->receive_root_directory = str_dup("/dst"); + char* argv[] = {"fastsync", "--append", "--append-verify", "/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->append); + EXPECT_TRUE(cfg->append_verify); + EXPECT_TRUE(validate_config(cfg)); + config_delete(cfg); +} + +static void test_validate_config_append_rejects_chunk_serialization() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--append", "-s", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0); + EXPECT_FALSE(validate_config(cfg)); + config_delete(cfg); +} + +static void test_validate_config_append_verify_rejects_chunk_serialization() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--append-verify", "-s", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0); + EXPECT_FALSE(validate_config(cfg)); + config_delete(cfg); +} + +static void test_validate_config_append_rejects_whole_file() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--append", "-W", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0); + EXPECT_FALSE(validate_config(cfg)); + config_delete(cfg); +} + +static void test_validate_config_append_verify_rejects_whole_file() { + Config* cfg = config_create(); + char* argv[] = {"fastsync", "--append-verify", "-W", "/src", "/dst"}; + int positional_args[2]; + int positional_count = 0; + EXPECT_EQ_INT(parse_args(cfg, 5, argv, positional_args, &positional_count), 0); + EXPECT_FALSE(validate_config(cfg)); + config_delete(cfg); +} + void test_client_cli() { test_validate_config_required_paths(); + test_parse_args_append(); + test_parse_args_append_verify(); + test_parse_args_append_both(); + test_validate_config_append_rejects_chunk_serialization(); + test_validate_config_append_verify_rejects_chunk_serialization(); + test_validate_config_append_rejects_whole_file(); + test_validate_config_append_verify_rejects_whole_file(); test_validate_config_incompatible_options(); test_validate_config_tls_requirements(); test_validate_config_delta_sendfile_constraints(); diff --git a/tests/test_config.c b/tests/test_config.c index ced2318..b6f4b17 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -754,6 +754,49 @@ static void test_config_is_remote_dest() { EXPECT_TRUE(config_is_remote_dest("user@host:")); } +/* The append-mode fields cross the wire unchanged: --append and --append-verify + are negotiated to the receiver so it knows to reply STATUS_APPEND on a + shorter destination. */ +static void test_config_append_wire_roundtrip() { + if (is_running_under_valgrind()) + return; + struct { + bool append, append_verify; + } cases[] = {{true, false}, {false, true}, {true, true}, {false, false}}; + for (size_t i = 0; i < sizeof(cases) / sizeof(cases[0]); i++) { + int p[2]; + EXPECT_EQ_INT(socketpair(AF_UNIX, SOCK_STREAM, 0, p), 0); + 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->append == cases[i].append && recv->append_verify == cases[i].append_verify; + config_delete(recv); + close(p[0]); + _exit(ok ? 0 : 1); + } else { + close(p[0]); + io_set_fds(p[1], p[1]); + Config* send_cfg = config_create(); + EXPECT_NOT_NULL(send_cfg); + send_cfg->send_directory = str_dup("/src"); + send_cfg->receive_root_directory = str_dup("/dst"); + send_cfg->append = cases[i].append; + send_cfg->append_verify = cases[i].append_verify; + bool sent = config_send(p[1], send_cfg); + int status; + waitpid(pid, &status, 0); + close(p[1]); + config_delete(send_cfg); + EXPECT_TRUE(sent); + EXPECT_TRUE(WIFEXITED(status) && WEXITSTATUS(status) == 0); + } + } +} + void test_config() { test_config_lifecycle(); test_config_ssh_dest(); @@ -771,6 +814,7 @@ void test_config() { test_config_delete_timing_wire_roundtrip(); test_config_delete_timing_conflict_rejected(); test_config_delete_policy_wire_roundtrip(); + test_config_append_wire_roundtrip(); test_config_basis_roundtrip(); test_config_basis_wire_rejects_escaping(); test_config_basis_normalization(); diff --git a/tests/test_shared_utils.c b/tests/test_shared_utils.c index 9828620..a2cbc1c 100644 --- a/tests/test_shared_utils.c +++ b/tests/test_shared_utils.c @@ -276,6 +276,22 @@ void test_shared_utils() { test_walker_unlimited_deletes_all(); test_walker_hard_bound_all_or_nothing(); + /* --append / --append-verify tail-resume math: a resume is eligible only for + a shorter existing destination, and the tail length is then the difference. */ + EXPECT_TRUE(append_resume_eligible(0, 10)); + EXPECT_TRUE(append_resume_eligible(7, 10)); + EXPECT_FALSE(append_resume_eligible(10, 10)); + EXPECT_FALSE(append_resume_eligible(11, 10)); + + unsigned long long tail; + EXPECT_TRUE(append_tail_length(0, 10, &tail)); + EXPECT_EQ_INT((int)tail, 10); + EXPECT_TRUE(append_tail_length(7, 10, &tail)); + EXPECT_EQ_INT((int)tail, 3); + EXPECT_FALSE(append_tail_length(10, 10, &tail)); + EXPECT_FALSE(append_tail_length(11, 10, &tail)); + EXPECT_FALSE(append_tail_length(7, 10, NULL)); + char formatted[32]; EXPECT_TRUE(format_human_bytes(0, formatted, sizeof(formatted))); EXPECT_EQ_STR(formatted, "0 B"); From bfb3a531e9754c8f6e2eb85b80861038d4344101 Mon Sep 17 00:00:00 2001 From: TapTap Date: Mon, 7 Sep 2026 15:43:52 +0200 Subject: [PATCH 3/3] docs: mark --append / --append-verify implemented in RSYNC_COMPAT Flip both rows to Implemented with precise Notes covering the tail-resume model, the prefix-verify semantics (and the plain-append rsync-parity safety statement), the new wire frames, and the PROTOCOL_VERSION 2.9.0 -> 2.10.0 bump. Add a Phase-3 append-wave implementation-notes block. The Summary count line is deliberately left untouched. --- RSYNC_COMPAT.md | 30 ++++++++++++++++++++++++++++-- 1 file changed, 28 insertions(+), 2 deletions(-) diff --git a/RSYNC_COMPAT.md b/RSYNC_COMPAT.md index 8baf6b6..59514d5 100644 --- a/RSYNC_COMPAT.md +++ b/RSYNC_COMPAT.md @@ -83,8 +83,8 @@ This document maps rsync's full feature set to FastSync's current implementation |------|-------------------|-----------------|-------| | `-u`, `--update` | Skip files newer on receiver | ❌ Not Implemented | Removed because it had no effect | | `--inplace` | Update files in-place | ✅ Implemented | Direct write mode | -| `--append` | Append data to shorter files | ❌ Not Implemented | Removed because it had no effect | -| `--append-verify` | Append with old-data checksum | ❌ Not Implemented | Removed because it had no effect | +| `--append` | Append data to shorter files | ✅ Implemented | Tail-only resume. When an existing destination file is SHORTER than the source, the receiver negotiates a resume offset with the sender and only the tail is transferred; the receiver rebuilds the full file (retained prefix + tail) and installs it through the normal atomic store path, so the result is byte-identical to the source whenever the retained prefix matches. Plain `--append` does NOT content-verify that prefix (rsync parity): a destination whose prefix differs from the source is resumed anyway, so the result (wrong prefix + correct tail) is NOT byte-identical and the file is effectively left corrupt — the documented rsync-parity risk (use `--append-verify` when the prefix cannot be trusted). Non-content attributes (permissions/ownership/mtime, via `-M`) are still applied. Requires the per-file `STATUS_CHECK` handshake, so it implies `--incremental`; it takes precedence over block delta for a growing file and falls back to delta/full when the destination is not shorter. Incompatible with `-s` (chunk serialization) and `--whole-file` (both rejected up front so the mode never silently degrades to a full transfer). Combines with `--inplace`, `--partial`/`--partial-dir`, and `--delay-updates` (the reconstructed full file flows through those paths unchanged). Divergence: rsync appends in place; FastSync reconstructs and atomically installs, so an interrupted or failed resume never leaves a half-written file at the destination (no corruption window), and `--append` is thus safe to use with the normal atomic path — not only with in-place writes | +| `--append-verify` | Append with old-data checksum | ✅ Implemented | Like `--append`, but the retained prefix IS verified before resuming: the sender transmits the source prefix checksum and the receiver compares it to the xxHash64 of the retained destination prefix; on a match only the tail is transferred, on a MISMATCH the run falls back to a clean full transfer so the result is always a byte-identical source copy (never a corrupt prefix+tail blend). Wire/protocol: the append handshake adds `STATUS_APPEND` / `STATUS_APPEND_SIG` / `STATUS_APPEND_OK` / `STATUS_APPEND_DATA` frames and `PROTOCOL_VERSION` was bumped **2.9.0 → 2.10.0** (peers must match, and both must be 2.10.0 or the run fails the version check). Same implications/incompatibilities as `--append`; when both spellings are given `--append-verify` wins (the safer semantics). See the Phase-3 append notes below | | `-W`, `--whole-file` | Copy whole file (no delta) | ❌ Not Implemented | | | `--block-size=SIZE` | Force checksum block-size | ⚠️ Partial | Parsed as `--delta-block`; controls delta transfer block size | @@ -189,6 +189,32 @@ order-independent because it runs over the fully parsed config. The deletion POLICY flags (`--delete-excluded`, `--max-delete`, `--ignore-errors`, `--force`) do NOT imply `--delete`; without `--delete` they are inert (matching rsync). +**Append-resume notes (Phase 3, append wave):** `--append` and `--append-verify` +are real. Both are negotiated when an existing destination file is found to be +**shorter** than the source during the per-file `STATUS_CHECK`; the receiver +replies with a new `STATUS_APPEND` frame carrying the resume offset (the prefix +length it already holds) instead of `STATUS_NEXT`/`STATUS_DELTA_SIGNATURE`. +The sender transmits ONLY the tail. For `--append-verify` it first sends the +source's prefix xxHash64 in a `STATUS_APPEND_SIG` frame; the receiver compares +it to the retained prefix and answers `STATUS_APPEND_OK` (transfer the tail) or +`STATUS_NEXT` (prefix mismatch → the sender falls back to a byte-exact full +transfer). The tail arrives in a `STATUS_APPEND_DATA` frame (compression and +metadata still apply). The receiver then rebuilds the full file in memory +(prefix + tail) and routes it through the existing atomic store engine, so all +of `--inplace`, `--partial`/`--partial-dir`, `--delay-updates`, `--backup`, +`--existing`/`--ignore-existing`/`--update` and delete-manifest behaviour is +unchanged and the result is a byte-identical source copy (given a matching +prefix). These new frames changed the wire, so `PROTOCOL_VERSION` was bumped +**2.9.0 → 2.10.0** (peers must match; the pre-existing `append`/`append_verify` +config booleans already crossed the wire). CLI: both flags imply `--incremental` +(the handshake needs it); they are incompatible with `-s` (chunk serialization) +and `--whole-file` (both rejected up front, never a silent full transfer); when +both spellings are given `--append-verify` wins. The FastSync divergence from +rsync is intentional and safer: rsync appends in place, whereas FastSync +reconstructs the whole file and atomically installs it, so an interrupted or +failed resume never leaves a partial/corrupt file at the destination — this is +why plain `--append` works on the normal atomic path, not only with `--inplace`. + ## 8. Metadata Preservation | Flag | Rsync Description | FastSync Status | Notes |