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.
This commit is contained in:
2026-09-07 15:43:29 +02:00
parent 6701c103cb
commit 6ad3887aa1
11 changed files with 470 additions and 51 deletions
+13
View File
@@ -435,6 +435,8 @@ static const OptionEntry OPTION_TABLE[] = {
{"--copy-unsafe-links", NULL, OPT_FLAG, offsetof(Config, copy_unsafe_links)}, {"--copy-unsafe-links", NULL, OPT_FLAG, offsetof(Config, copy_unsafe_links)},
{"--sparse", "-S", OPT_FLAG, offsetof(Config, preserve_sparse)}, {"--sparse", "-S", OPT_FLAG, offsetof(Config, preserve_sparse)},
{"--inplace", NULL, OPT_FLAG, offsetof(Config, inplace)}, {"--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)}, {"--fsync", NULL, OPT_FLAG, offsetof(Config, use_fsync)},
{"--checksum", NULL, OPT_FLAG, offsetof(Config, checksum)}, {"--checksum", NULL, OPT_FLAG, offsetof(Config, checksum)},
{"--8-bit-output", "-8", OPT_FLAG, offsetof(Config, eight_bit_output)}, {"--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; 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. */ /* Incremental and delta transfers need metadata unless the user disabled it. */
if ((config->use_incremental || config->use_delta) && !config->use_metadata && if ((config->use_incremental || config->use_delta) && !config->use_metadata &&
!config->metadata_explicitly_disabled) { !config->metadata_explicitly_disabled) {
+124 -3
View File
@@ -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, 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; *out_sig = NULL;
if (resume_offset)
*resume_offset = 0;
if (!send_status(client->file_descriptor, STATUS_CHECK)) if (!send_status(client->file_descriptor, STATUS_CHECK))
return -1; return -1;
if (!send_str(client->file_descriptor, file_wire_path(file))) 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; *out_sig = sig;
return 2; 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) { if (s != STATUS_NEXT) {
log_message(LOG_LEVEL_ERROR, "Unexpected server status"); log_message(LOG_LEVEL_ERROR, "Unexpected server status");
send_status(client->file_descriptor, STATUS_ERROR); 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; 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). // Send a single file directly (non-incremental path).
static bool send_file_direct(File* file, int fd, bool use_metadata, int compression_level, static bool send_file_direct(File* file, int fd, bool use_metadata, int compression_level,
const Config* config) { 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 // Incremental path: use sendfile for the actual data if enabled and no compression
if (use_sendfile) { if (use_sendfile) {
DeltaSignature* sig = NULL; 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) { if (rc == 1) {
log_info_message(LOG_INFO_SKIP, "Skipping unchanged %s", file->path); log_info_message(LOG_INFO_SKIP, "Skipping unchanged %s", file->path);
delta_signature_destroy(sig); 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); delta_signature_destroy(sig);
return -1; 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 == 0: unchanged file, skip
// rc == 2: server sent delta signature but sendfile doesn't support delta // rc == 2: server sent delta signature but sendfile doesn't support delta
delta_signature_destroy(sig); 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) // Incremental path with single_calls (supports compression and delta)
DeltaSignature* sig = NULL; 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) { if (rc < 0) {
delta_signature_destroy(sig); delta_signature_destroy(sig);
return -1; return -1;
@@ -906,6 +1016,17 @@ static int send_single_file(Client* client, File* file, Config* config, bool use
delta_signature_destroy(sig); delta_signature_destroy(sig);
return 1; 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) { if (rc == 2 && config->use_delta && !config->whole_file) {
int drc = send_delta(client, file, sig, config); int drc = send_delta(client, file, sig, config);
delta_signature_destroy(sig); delta_signature_destroy(sig);
+18 -6
View File
@@ -51,14 +51,26 @@ bool validate_config(const Config* config) {
log_message(LOG_LEVEL_ERROR, "--delta cannot be combined with -f (sendfile)"); log_message(LOG_LEVEL_ERROR, "--delta cannot be combined with -f (sendfile)");
return false; return false;
} }
if (config->log_file_format && !config->log_file) { /* --append / --append-verify resume a shorter existing destination by
log_message(LOG_LEVEL_ERROR, "--log-file-format requires --log-file"); 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; return false;
} }
if (config->append || config->append_verify) { if ((config->append || config->append_verify) && config->whole_file) {
fprintf( log_message(LOG_LEVEL_ERROR,
stderr, "--append/--append-verify are incompatible with --whole-file (which forces a "
"Error: --append and --append-verify are not supported yet; refusing to ignore option\n"); "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; return false;
} }
if (config->use_tls) { if (config->use_tls) {
+4
View File
@@ -159,6 +159,10 @@ void print_usage(void) {
printf(" --copy-unsafe-links Only transform unsafe symlinks into referent files\n"); printf(" --copy-unsafe-links Only transform unsafe symlinks into referent files\n");
printf(" -S, --sparse Handle sparse files efficiently\n"); printf(" -S, --sparse Handle sparse files efficiently\n");
printf(" --inplace Update files in-place (no temp+rename)\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(" --fsync Fsync every written file before publication\n");
printf(" --compress-level <n> Compression level (default: 5)\n"); printf(" --compress-level <n> Compression level (default: 5)\n");
printf(" --zl <n> Alias for --compress-level\n"); printf(" --zl <n> Alias for --compress-level\n");
+4
View File
@@ -174,6 +174,10 @@ static bool validate_received_config(const Config* config) {
valid_wire_bool(config->checksum) && valid_wire_bool(config->eight_bit_output) && valid_wire_bool(config->checksum) && valid_wire_bool(config->eight_bit_output) &&
config_has_valid_delete_timing(config) && config_has_valid_delete_timing(config) &&
!(config->skip_compress_set && config->use_chunk_serialization) && !(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->use_compression ||
(config->compression_level >= 1 && config->compression_level <= 22)) && (config->compression_level >= 1 && config->compression_level <= 22)) &&
config->chunk_size > 0 && config->chunk_size <= MAX_CHUNK_SIZE && config->chunk_size > 0 && config->chunk_size <= MAX_CHUNK_SIZE &&
+1 -1
View File
@@ -231,7 +231,7 @@ typedef struct Config {
DelayUpdatesContext* delay_context; DelayUpdatesContext* delay_context;
} Config; } Config;
#define PROTOCOL_VERSION "2.9.0" #define PROTOCOL_VERSION "2.10.0"
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
/* Upper bound on total basis-dir entries (rsync caps --link-dest at 20). */ /* Upper bound on total basis-dir entries (rsync caps --link-dest at 20). */
#define MAX_BASIS_DIRS 64 #define MAX_BASIS_DIRS 64
+263 -40
View File
@@ -983,6 +983,50 @@ static void* fuzzy_basis_find_and_load(const Config* config, const char* check_p
return basis; 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) { File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
if (!config || !skipped) { if (!config || !skipped) {
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
@@ -1193,6 +1237,224 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
basis_match_free(&basis); 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) { if (try_delta && old_data != NULL) {
bool delta_failed = false; bool delta_failed = false;
File* delta_file = File* delta_file =
@@ -1256,48 +1518,9 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
} }
close(old_fd); close(old_fd);
File* file = file_create(check_path); File* file = receive_full_file(fd, config, check_path);
free(check_path); free(check_path);
free(full_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; return file;
} }
+8
View File
@@ -397,6 +397,14 @@ static const char* status_to_string(Status status) {
return "CHECK_BATCH"; return "CHECK_BATCH";
case STATUS_MKDIR: case STATUS_MKDIR:
return "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: default:
return "UNKNOWN"; return "UNKNOWN";
} }
+15 -1
View File
@@ -70,7 +70,21 @@ enum NET_STATUS {
STATUS_CHECK_BATCH, STATUS_CHECK_BATCH,
/* An explicit directory entry (--dirs): the sender transmits only the path; /* An explicit directory entry (--dirs): the sender transmits only the path;
* the receiver creates the directory below the receive root. */ * 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); void io_set_fds(int read_fd, int write_fd);
+12
View File
@@ -531,3 +531,15 @@ char* path_cat(const char* path1, const char* path2) {
new_path[path1_len + path2_len + 1] = '\0'; new_path[path1_len + path2_len + 1] = '\0';
return new_path; 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;
}
+8
View File
@@ -53,5 +53,13 @@ void utils_set_authorized_root_fd(int fd);
bool has_path_traversal(const char* path); bool has_path_traversal(const char* path);
bool utils_valid_batch_path(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); 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 #endif