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)
This commit is contained in:
2026-09-05 13:11:41 +02:00
17 changed files with 1007 additions and 111 deletions
+3 -2
View File
@@ -484,10 +484,11 @@ defaults to the current directory. |
## 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,
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.
TLS provides encrypted TCP transport. Supplying `--ca` enables certificate
+5
View File
@@ -12,6 +12,7 @@
#include "utils.h"
#include <errno.h>
#include <limits.h>
#include <signal.h>
#include <stdbool.h>
#include <stddef.h>
#include <stdio.h>
@@ -903,6 +904,10 @@ static int read_patterns_from_file(const char* filepath, char*** patterns, int*
#ifndef FASTSYNC_TEST_BUILD
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_dest = NULL;
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;
}
static bool finalize_transfer(Client* client) {
Status status;
return send_status(client->file_descriptor, STATUS_FINISHED) &&
receive_status(client->file_descriptor, &status) && status == STATUS_OK;
}
/* (finalize_transfer is defined after the SourceFile helpers below.) */
typedef struct {
typedef struct SourceFile {
char* path;
dev_t device;
ino_t inode;
bool skipped; /* receiver reported the file was not written */
} SourceFile;
static void source_file_destroy(void* item) {
@@ -135,6 +132,8 @@ static void remove_transferred_sources(const Config* config, ArrayList* paths) {
return;
for (int i = 0; i < paths->size; i++) {
SourceFile* source = paths->items[i];
if (source->skipped)
continue;
const char* slash = strrchr(source->path, '/');
const char* leaf = slash ? slash + 1 : source->path;
char parent[PATH_MAX];
@@ -177,6 +176,7 @@ static SourceFile* source_file_create(const File* file) {
source->path = str_dup(file->path);
source->device = st.st_dev;
source->inode = st.st_ino;
source->skipped = false;
if (!source->path) {
source_file_destroy(source);
return NULL;
@@ -203,6 +203,33 @@ static void mark_sender_done(PipelineContextSender* context) {
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) {
mtx_lock(&context->mutex_scanner);
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)
goto send_fail;
}
bool ok = finalize_transfer(client);
bool ok = finalize_transfer(client, context->config, context->remove_source_files);
if (ok)
remove_transferred_sources(context->config, context->remove_source_files);
mtx_lock(&context->mutex_progress);
@@ -866,7 +893,7 @@ int send_files(Config* config) {
array_list_delete(manifest);
manifest = NULL;
}
bool ok = finalize_transfer(client);
bool ok = finalize_transfer(client, config, remove_sources);
if (ok)
remove_transferred_sources(config, remove_sources);
if (config->show_progress && !config->quiet)
+85 -10
View File
@@ -1,6 +1,7 @@
#include "receiver.h"
#include "chunk.h"
#include "file_receive.h"
#include "log.h"
#include "metadata.h"
#include "protocol.h"
@@ -8,6 +9,48 @@
#include <stdlib.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) {
if (!chunk || !sink || !sink->store_file)
return false;
@@ -52,7 +95,7 @@ static bool receiver_process_batch(Config* config, int file_descriptor) {
send_status(file_descriptor, STATUS_ERROR);
return false;
}
if (check_size > MAX_RECEIVE_FILE_SIZE) {
if (check_size > MAX_RECEIVE_WHOLE_FILE_SIZE) {
free(check_path);
send_status(file_descriptor, STATUS_ERROR);
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");
goto receive_error;
}
if (sink->send_success && !send_status(file_descriptor, STATUS_OK))
return -1;
if (sink->send_success) {
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;
receive_error:
@@ -140,15 +189,41 @@ receive_error:
return -1;
}
static bool receiver_save_file(File* file, void* context) {
Config* config = context;
bool success =
!config->save_to_disk || file_save_to_disk(config->receive_root_directory, file, config);
/* ---- Single-threaded sink (used by receiver_receive_files) ---- */
typedef struct {
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);
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) {
ReceiverSink sink = {receiver_save_file, config, true, true};
return receiver_process(config, file_descriptor, &sink);
ReceiverSaveContext context = {.config = config, .outcomes = {0}};
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 "file.h"
#include "file_receive.h"
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 {
ReceiverFileSink store_file;
void* context;
bool send_error;
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;
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_receive_files(Config* config, int file_descriptor);
+9 -3
View File
@@ -342,9 +342,15 @@ void handler(int file_descriptor) {
int writer_result;
thrd_join(receiver, &receiver_result);
thrd_join(writer, &writer_result);
send_status(file_descriptor, receiver_result == thrd_success && writer_result == thrd_success
? STATUS_OK
: STATUS_ERROR);
bool transfer_ok = receiver_result == thrd_success && writer_result == thrd_success;
if (transfer_ok) {
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);
} else {
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->size_only) && valid_wire_bool(config->ignore_times) &&
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->safe_links) && valid_wire_bool(config->copy_unsafe_links) &&
valid_wire_bool(config->remove_source_files) && valid_wire_bool(config->follow_symlinks) &&
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_xattrs) && valid_wire_bool(config->preserve_devices) &&
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) {
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->safe_links) && send_int(fd, c->copy_unsafe_links) &&
send_int(fd, c->preserve_hard_links) && send_int(fd, c->preserve_acls) &&
send_int(fd, c->preserve_xattrs) && send_int(fd, c->preserve_devices) &&
send_int(fd, c->preserve_sparse);
send_int(fd, c->remove_source_files) && send_int(fd, c->follow_symlinks) &&
send_int(fd, c->copy_links) && send_int(fd, c->safe_links) &&
send_int(fd, c->copy_unsafe_links) && send_int(fd, c->preserve_hard_links) &&
send_int(fd, c->preserve_acls) && send_int(fd, c->preserve_xattrs) &&
send_int(fd, c->preserve_devices) && send_int(fd, c->preserve_sparse);
}
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) {
if (!receive_wire_bool(fd, &c->backup))
return false;
c->backup_dir = receive_str(fd);
if (!c->backup_dir)
char* backup_dir = receive_str(fd);
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;
bool* flags[] = {&c->follow_symlinks, &c->copy_links, &c->safe_links,
&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) {
c->temp_dir = receive_str(fd);
if (!c->temp_dir || !receive_wire_bool(fd, &c->partial))
char* temp_dir = receive_str(fd);
if (!temp_dir)
return false;
c->partial_dir = receive_str(fd);
c->suffix = c->partial_dir ? receive_str(fd) : NULL;
if (!c->partial_dir || !c->suffix || !receive_wire_bool(fd, &c->delete_before))
if (*temp_dir != '\0') {
c->temp_dir = temp_dir;
} 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;
if (!receive_wire_bool(fd, &c->checksum))
return false;
+1 -1
View File
@@ -147,7 +147,7 @@ typedef struct Config {
bool skip_compress_set;
} Config;
#define PROTOCOL_VERSION "2.4.0"
#define PROTOCOL_VERSION "2.5.0"
#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024)
Config* config_create(void);
+112 -67
View File
@@ -20,9 +20,14 @@
#include "utils.h"
#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) {
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
would make a concurrent no-replace commit overwrite its old name. */
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_dir = (config && config->backup_dir) ? config->backup_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* destination_path = 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 ||
strcmp(backup_suffix, ".") == 0 || strcmp(backup_suffix, "..") == 0))) {
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
root, never independent filesystem roots. */
if ((backup_dir && (backup_dir[0] == '/' || has_path_traversal(backup_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)))
return false;
return FILE_SAVE_ERROR;
if (partial_dir && !(confined_partial = path_cat(root_directory, partial_dir))) {
free(confined_backup);
return false;
return FILE_SAVE_ERROR;
}
const char* actual_root =
(partial_dir && config && config->partial) ? confined_partial : root_directory;
const char* actual_root = use_partial_root ? confined_partial : root_directory;
destination_path = path_cat(root_directory, file->path);
disk_path = path_cat(actual_root, file->path);
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(destination_path);
free(disk_path);
return false;
return FILE_SAVE_ERROR;
}
/* --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(destination_path);
free(disk_path);
return true;
return FILE_SAVE_SKIPPED;
}
/* --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(destination_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.
The secure stat does not require read permission on the destination. */
if (config && config->update && file_destination_is_newer_secure(disk_path, file->metadata)) {
In partial-dir mode the entry that would be replaced is the real
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_partial);
free(destination_path);
free(disk_path);
return true;
return FILE_SAVE_SKIPPED;
}
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;
if (file_stat_secure(disk_path, &backup_stat)) {
if (file_stat_secure(replace_target, &backup_stat)) {
if (backup_dir) {
backup_path = path_cat(confined_backup, file->path);
} else {
size_t path_len = strlen(disk_path);
size_t path_len = strlen(replace_target);
size_t suffix_len = strlen(backup_suffix);
if (path_len > SIZE_MAX - suffix_len - 1)
goto fail;
backup_path = malloc(path_len + suffix_len + 1);
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);
}
}
@@ -125,7 +136,7 @@ bool file_save_to_disk(const char* root_directory, const File* file, const Confi
goto fail;
free(parent_copy);
parent_copy = NULL;
if (!file_rename_secure(disk_path, backup_path))
if (!file_rename_secure(replace_target, backup_path))
goto fail;
free(backup_path);
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,
inplace, sparse, metadata, preserve_executability,
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(backup_path);
free(confined_backup);
free(confined_partial);
free(destination_path);
free(disk_path);
return ok;
return FILE_SAVE_WRITTEN;
fail:
free(parent_copy);
@@ -164,13 +187,15 @@ fail:
free(confined_partial);
free(destination_path);
free(disk_path);
return false;
return FILE_SAVE_ERROR;
}
static File* receive_delta_file(int fd, const Config* config, const char* check_path,
void* old_data, unsigned long long old_size, bool* failed) {
if (!old_data)
if (!old_data) {
*failed = true;
return NULL;
}
DeltaSignature* sig = delta_signature_create(old_data, old_size, config->delta_block_size);
if (!sig) {
@@ -206,7 +231,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
}
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) {
delta_signature_destroy(sig);
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(
check_path, config->skip_compress_suffixes,
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);
if (!raw_delta) {
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;
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);
free(old_data);
delta_signature_destroy(sig);
send_status(fd, STATUS_ERROR);
*failed = true;
return NULL;
}
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);
delta_signature_destroy(sig);
send_status(fd, STATUS_ERROR);
*failed = true;
return NULL;
}
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) {
file_destroy(file);
*failed = true;
@@ -325,7 +352,7 @@ static File* receive_delta_file(int fd, const Config* config, const char* check_
!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_FILE_SIZE);
Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE);
data_destroy(file_data);
if (uncompressed == NULL) {
file_destroy(file);
@@ -384,7 +411,7 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
return NULL;
}
if (check_size > MAX_RECEIVE_FILE_SIZE) {
if (check_size > MAX_RECEIVE_WHOLE_FILE_SIZE) {
free(check_path);
send_status(fd, STATUS_ERROR);
return NULL;
@@ -405,6 +432,9 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
send_status(fd, STATUS_ERROR);
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;
bool has_old_file = false;
int old_fd = -1;
@@ -416,9 +446,34 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
close(parent_fd);
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;
/* 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;
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);
if (old_data) {
size_t got = 0;
@@ -433,73 +488,63 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
}
}
}
if (old_fd >= 0) {
close(old_fd);
}
bool match =
!config->ignore_times && has_old_file && (unsigned long long)st.st_size == check_size;
if (match && config->checksum) {
uint64_t old_checksum = old_size == 0 ? delta_xxhash64("", 0) : 0;
if (old_data)
old_checksum = delta_xxhash64(old_data, (size_t)old_size);
match = (old_size == 0 || old_data) && old_checksum == check_checksum;
free(old_data);
old_data = NULL;
} else if (match && !config->size_only) {
long long old_mtime_nsec = 0;
#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);
/* Quick-skip decision. If no content comparison is required this is final
and the old file was never read; if the read failed the file is not
skipped and the transfer proceeds with the full new contents. */
bool match = false;
if (checksum_needs_read) {
if (old_size == 0)
match = delta_xxhash64("", 0) == check_checksum;
else
match = old_data != NULL && delta_xxhash64(old_data, (size_t)old_size) == check_checksum;
} else if (size_equal && !config->ignore_times) {
match = config->size_only || match_by_metadata;
}
if (match) {
free(old_data);
if (!send_status(fd, STATUS_OK)) {
close(old_fd);
free(full_path);
free(check_path);
return NULL;
}
close(old_fd);
free(full_path);
free(check_path);
*skipped = true;
return NULL;
}
bool try_delta = config->use_delta && !config->whole_file && has_old_file && old_data != NULL &&
delta_should_attempt(old_size, check_size, config->delta_max_file_size);
if (try_delta) {
if (try_delta && old_data != NULL) {
bool delta_failed = false;
File* delta_file =
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 */
if (delta_file) {
close(old_fd);
free(full_path);
free(check_path);
return delta_file;
}
if (delta_failed) {
close(old_fd);
free(full_path);
free(check_path);
return NULL;
}
free(old_data);
old_data = NULL;
try_delta = false;
}
free(old_data);
old_data = NULL;
if (!try_delta) {
free(old_data);
old_data = NULL;
if (!send_status(fd, STATUS_NEXT)) {
free(full_path);
free(check_path);
return NULL;
}
if (!send_status(fd, STATUS_NEXT)) {
close(old_fd);
free(full_path);
free(check_path);
return NULL;
}
close(old_fd);
File* file = file_create(check_path);
free(check_path);
@@ -517,7 +562,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) {
file_destroy(file);
return NULL;
@@ -527,7 +572,7 @@ File* receive_incremental_check(int fd, const Config* config, bool* skipped) {
!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_FILE_SIZE);
Data* uncompressed = data_decompress_limited(file_data, MAX_RECEIVE_WHOLE_FILE_SIZE);
data_destroy(file_data);
if (uncompressed == NULL) {
file_destroy(file);
@@ -571,7 +616,7 @@ File* file_receive(const Config* config, int file_descriptor) {
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) {
file_destroy(file);
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,
config->skip_compress_set ? config->skip_compress_count
: -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);
if (file_data_uncompressed == NULL) {
file_destroy(file);
+8
View File
@@ -10,6 +10,14 @@
File* file_receive(const Config* config, int file_descriptor);
File* receive_incremental_check(int fd, const Config* config, bool* skipped);
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);
#endif
+25 -2
View File
@@ -105,6 +105,9 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
context->queue = queue;
context->file_descriptor = file_descriptor;
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_set_ssl(&context->session, ssl);
context->receiver_done = false;
@@ -137,6 +140,7 @@ fail:
void pipeline_context_receiver_destroy(PipelineContextReceiver* context) {
config_delete(context->config);
queue_destroy(context->queue);
receiver_outcomes_destroy(&context->outcomes);
mtx_destroy(&context->mutex);
cnd_destroy(&context->condition_not_full);
cnd_destroy(&context->condition_not_empty);
@@ -170,7 +174,7 @@ int receive_thread(void* pipeline_context) {
const Config* config = context->config;
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) {
receiver_thread_fail(context);
protocol_session_unbind();
@@ -211,7 +215,26 @@ int write_thread(void* pipeline_context) {
protocol_session_unbind();
return thrd_success;
}
if (save_to_disk && !file_save_to_disk(root_directory, file, context->config)) {
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);
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);
mtx_lock(&context->mutex);
atomic_store(&context->cancelled, true);
+2
View File
@@ -9,6 +9,7 @@
#include "file.h"
#include "protocol.h"
#include "queue.h"
#include "receiver.h"
#include <openssl/ssl.h>
typedef struct {
@@ -40,6 +41,7 @@ typedef struct PipelineContextReceiver {
int file_descriptor;
SSL* ssl;
ProtocolSession session;
ReceiverOutcomes outcomes;
mtx_t mutex;
cnd_t condition_not_full;
cnd_t condition_not_empty;
+10 -4
View File
@@ -9,10 +9,16 @@
/* Maximum allowed string size for receive_str (64 KB) */
#define MAX_STRING_SIZE (64 * 1024)
/* Maximum allowed data payload size for receive_data (100 MB) */
#define MAX_DATA_PAYLOAD_SIZE (100ULL * 1024 * 1024)
/* Maximum uncompressed file payload accepted by the receiver. */
#define MAX_RECEIVE_FILE_SIZE (64ULL * 1024 * 1024)
/* Maximum uncompressed file payload accepted by the receiver's whole-file
* paths. A single whole file is charged against the per-connection memory
* reservation (MAX_CONNECTION_MEMORY) and against the server allocation
* 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 */
#define MAX_CHUNK_SIZE (64ULL * 1024 * 1024)
+184
View File
@@ -1,4 +1,5 @@
"""Feature tests: incremental sync, bandwidth limiting, dry run, metadata, filters."""
import filecmp
import os
import shutil
import sys
@@ -813,3 +814,186 @@ class TestBandwidthLimit:
mismatches, missing = verify_transfer(SOURCE_DIR, received)
assert not missing, f"Missing: {missing}"
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"
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)
+100
View File
@@ -282,6 +282,105 @@ static void test_config_receive_truncated() {
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() {
/* Valid SSH-style destinations */
EXPECT_TRUE(config_is_remote_dest("user@host:/path"));
@@ -315,6 +414,7 @@ void test_config() {
test_config_send_receive();
test_config_send_receive_version_mismatch();
test_config_receive_truncated();
test_config_string_null_vs_empty_roundtrip();
}
test_config_is_remote_dest();
}
+108
View File
@@ -9,6 +9,7 @@
#include <string.h>
#include <sys/stat.h>
#include <sys/wait.h>
#include <time.h>
#include <unistd.h>
static void test_file_create() {
@@ -250,6 +251,111 @@ static void test_file_save_to_disk_ignore_existing_entry_types() {
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() {
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),
@@ -635,6 +741,8 @@ void test_file() {
test_file_save_to_disk_existing();
test_file_save_to_disk_ignore_existing();
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_with_fsync();
test_file_write_to_disk_creates_dirs();
+250
View File
@@ -1,14 +1,18 @@
#include "test_server.h"
#include "config.h"
#include "delta.h"
#include "file.h"
#include "protocol.h"
#include "test_utils.h"
#include "utils.h"
#include <fcntl.h>
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <sys/socket.h>
#include <sys/stat.h>
#include <sys/wait.h>
#include <time.h>
#include <unistd.h>
#include "receiver.h"
@@ -210,6 +214,249 @@ static void test_receive_incremental_check_rejects_invalid_nanoseconds() {
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() {
if (!is_running_under_valgrind()) {
test_receive_files_finished();
@@ -217,5 +464,8 @@ void test_server() {
test_receive_files_abort();
test_receive_manifest_rejects_traversal();
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();
}
}