Compare commits
4
Commits
36c2c04910
...
02679fe335
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
02679fe335 | ||
|
|
c0c315cf48 | ||
|
|
ebfaced5c2 | ||
|
|
bbf982dc0b |
No files matched your search
@@ -129,6 +129,30 @@ manifest ack. `--delete-delay` and `--delete-during` are each implemented as
|
||||
the closest safe approximation their engine mode allows; the divergences are
|
||||
noted in the rows above.
|
||||
|
||||
Manifest size: the sender's keep-set collection (streaming or early pre-scan)
|
||||
is unbounded, but the receiver rejects any manifest beyond `MAX_MANIFEST_ENTRIES`
|
||||
(1 048 576 entries) / `MAX_MANIFEST_BYTES` (16 MB of paths) as a hard protocol
|
||||
error. In the commit modes this only means the deletion is refused after the
|
||||
data already arrived; in the NEW early modes (`--delete-before`/`--delete-during`)
|
||||
the manifest is the first frame, so an oversized keep-set now aborts the whole
|
||||
transfer BEFORE any data is sent (previously all data transferred and only the
|
||||
deletion step failed). Keep the source tree small enough for the receiver's
|
||||
manifest caps when using the early timing.
|
||||
|
||||
Early-delete ACK wait: after committing a large deletion (up to
|
||||
`MAX_SERVER_DELETE_COUNT` unlinks) the receiver's `STATUS_OK`/`STATUS_ERROR`
|
||||
reply can legitimately take much longer than a normal round trip, so the sender
|
||||
waits for that single ACK with an extended explicit deadline (1 hour) instead
|
||||
of the default 60 s per-message receive window. A receiver that is genuinely
|
||||
gone still aborts the wait via connection close/error; the extended bound only
|
||||
protects against aborting after the deletion already committed on the receiver.
|
||||
|
||||
Flag-conflict policy: unlike rsync's last-one-wins behaviour, every deletion
|
||||
timing flag implies `--delete`, and combining a timing flag with `--no-delete`
|
||||
(in either argument order) — or more than one timing flag — is rejected as a
|
||||
configuration error rather than silently resolved. Note the check is
|
||||
order-independent because it runs over the fully parsed config.
|
||||
|
||||
## 8. Metadata Preservation
|
||||
|
||||
| Flag | Rsync Description | FastSync Status | Notes |
|
||||
|
||||
@@ -597,14 +597,20 @@ static int send_delete_manifest(int fd, ArrayList* manifest) {
|
||||
--delete-before/--delete-during, where the extras are removed on the receiver
|
||||
BEFORE the first byte of file data is sent: the receiver acknowledges with
|
||||
STATUS_OK once the bounded delete committed, or STATUS_ERROR if it could not
|
||||
(in which case the sender aborts without streaming any data). */
|
||||
(in which case the sender aborts without streaming any data). The ACK may
|
||||
take much longer than an ordinary per-message round trip because the receiver
|
||||
performs the whole bounded deletion walk (up to MAX_SERVER_DELETE_COUNT
|
||||
unlinks) before replying, so the wait uses a generous explicit deadline
|
||||
instead of the default 60 s receive window. */
|
||||
#define DELETE_ACK_TIMEOUT_SEC 3600
|
||||
|
||||
static bool send_delete_manifest_early(Client* client, ArrayList* manifest) {
|
||||
if (!client || !manifest)
|
||||
return false;
|
||||
if (send_delete_manifest(client->file_descriptor, manifest) != 0)
|
||||
return false;
|
||||
Status ack;
|
||||
if (!receive_status(client->file_descriptor, &ack))
|
||||
if (!receive_status_timed(client->file_descriptor, &ack, DELETE_ACK_TIMEOUT_SEC))
|
||||
return false;
|
||||
if (ack != STATUS_OK) {
|
||||
log_message(LOG_LEVEL_ERROR, "Server failed to delete files before the transfer");
|
||||
|
||||
+7
-4
@@ -28,12 +28,15 @@ void print_usage(void) {
|
||||
printf(" transfer has succeeded)\n");
|
||||
printf(" --delete-before Delete extras before the transfer starts\n");
|
||||
printf(" (implies --delete)\n");
|
||||
printf(" --delete-during Delete extras once the keep-set is known, before\n");
|
||||
printf(" --del data is applied (alias --del; implies --delete)\n");
|
||||
printf(" --delete-during Delete extras once the keep-set manifest is known,\n");
|
||||
printf(" before the data is applied (implies --delete)\n");
|
||||
printf(" --del Alias for --delete-during\n");
|
||||
printf(" --delete-delay Delete extras only after a successful transfer\n");
|
||||
printf(" (implies --delete)\n");
|
||||
printf(" --delete-after Alias of the default --delete timing: delete only\n");
|
||||
printf(" after the transfer succeeded (implies --delete)\n");
|
||||
printf(" --delete-after Delete only after the whole transfer succeeded\n");
|
||||
printf(" (the default --delete timing; implies --delete)\n");
|
||||
printf(" Note: each timing flag implies --delete. Combining a timing flag with\n");
|
||||
printf(" --no-delete (in either order) is rejected as a config error.\n");
|
||||
printf(" --ignore-existing Skip files that already exist on receiver\n");
|
||||
printf(" --delay-updates Put updated files into place only at the end of transfer\n");
|
||||
printf(" --dirs, -d, --old-dirs, --old-d Transfer the named directory entries without\n");
|
||||
|
||||
+37
-19
@@ -148,18 +148,21 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
||||
if (!receive_status(file_descriptor, &status))
|
||||
return -1;
|
||||
bool early_delete = config_delete_timing_early(config);
|
||||
/* Parked keep-set for the late/commit timing. Every exit path below frees it
|
||||
exactly once; the only exception is the successful FINISHED handoff, which
|
||||
transfers ownership to *pending_manifest (used by the -m receiver). */
|
||||
ArrayList* deferred_manifest = NULL;
|
||||
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
|
||||
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH ||
|
||||
status == STATUS_MKDIR || status == STATUS_MANIFEST) {
|
||||
if (status == STATUS_KEEPALIVE) {
|
||||
if (!send_status(file_descriptor, STATUS_KEEPALIVE))
|
||||
return -1;
|
||||
goto fail;
|
||||
goto next_status;
|
||||
}
|
||||
if (status == STATUS_ABORT) {
|
||||
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
|
||||
return -1;
|
||||
goto fail;
|
||||
}
|
||||
if (status == STATUS_CHECK) {
|
||||
bool skipped;
|
||||
@@ -172,7 +175,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
||||
goto receive_error;
|
||||
} else if (status == STATUS_CHECK_BATCH) {
|
||||
if (!receiver_process_batch(config, file_descriptor))
|
||||
return -1;
|
||||
goto fail;
|
||||
goto next_status;
|
||||
} else if (status == STATUS_MKDIR) {
|
||||
File* dir = file_receive_directory(file_descriptor);
|
||||
@@ -181,7 +184,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
||||
} else if (status == STATUS_MANIFEST) {
|
||||
ArrayList* manifest = receive_manifest_entries(file_descriptor);
|
||||
if (!manifest)
|
||||
return -1; /* receive_manifest_entries already sent STATUS_ERROR */
|
||||
goto fail; /* receive_manifest_entries already sent STATUS_ERROR */
|
||||
if (early_delete) {
|
||||
/* --delete-before / --delete-during: the manifest is authoritative the
|
||||
moment it arrives, before any file data. Delete now and acknowledge
|
||||
@@ -192,24 +195,24 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
||||
array_list_delete(manifest);
|
||||
if (!deletion_ok) {
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
return -1;
|
||||
goto fail;
|
||||
}
|
||||
if (!send_status(file_descriptor, STATUS_OK))
|
||||
return -1;
|
||||
} else {
|
||||
goto fail;
|
||||
} else if (config->use_delete) {
|
||||
/* Plain --delete / --delete-after / --delete-delay: hold the keep-set
|
||||
and commit the deletion only after STATUS_FINISHED. */
|
||||
if (config->use_delete) {
|
||||
if (deferred_manifest) {
|
||||
log_message(LOG_LEVEL_ERROR, "Received a second delete manifest");
|
||||
array_list_delete(manifest);
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
return -1;
|
||||
}
|
||||
deferred_manifest = manifest;
|
||||
} else {
|
||||
if (deferred_manifest) {
|
||||
log_message(LOG_LEVEL_ERROR, "Received a second delete manifest");
|
||||
array_list_delete(deferred_manifest);
|
||||
deferred_manifest = NULL;
|
||||
array_list_delete(manifest);
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
goto fail;
|
||||
}
|
||||
deferred_manifest = manifest;
|
||||
} else {
|
||||
array_list_delete(manifest);
|
||||
}
|
||||
goto next_status;
|
||||
} else {
|
||||
@@ -241,26 +244,41 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
||||
if (deferred_manifest) {
|
||||
if (pending_manifest) {
|
||||
*pending_manifest = deferred_manifest;
|
||||
deferred_manifest = NULL;
|
||||
} else {
|
||||
bool deletion_ok = manifest_delete_extras(config, deferred_manifest);
|
||||
array_list_delete(deferred_manifest);
|
||||
deferred_manifest = NULL;
|
||||
if (!deletion_ok) {
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
return -1;
|
||||
goto fail;
|
||||
}
|
||||
}
|
||||
}
|
||||
if (sink->send_success) {
|
||||
if (sink->send_success_frame) {
|
||||
if (!sink->send_success_frame(file_descriptor, sink->context))
|
||||
return -1;
|
||||
goto fail;
|
||||
} else if (!send_status(file_descriptor, STATUS_OK)) {
|
||||
return -1;
|
||||
goto fail;
|
||||
}
|
||||
}
|
||||
return 0;
|
||||
|
||||
fail:
|
||||
/* Failure exits that must not (or already did) report a STATUS_ERROR. The
|
||||
parked keep-set is dropped: never commit a deletion for a failed stream. */
|
||||
if (deferred_manifest) {
|
||||
array_list_delete(deferred_manifest);
|
||||
deferred_manifest = NULL;
|
||||
}
|
||||
return -1;
|
||||
|
||||
receive_error:
|
||||
if (deferred_manifest) {
|
||||
array_list_delete(deferred_manifest);
|
||||
deferred_manifest = NULL;
|
||||
}
|
||||
if (sink->send_error)
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
return -1;
|
||||
|
||||
+26
-2
@@ -302,15 +302,25 @@ bool protocol_send_n_data(ProtocolSession* session, const void* data, size_t dat
|
||||
return true;
|
||||
}
|
||||
|
||||
bool protocol_receive_n_data_timed(ProtocolSession* session, void* data, size_t data_size,
|
||||
int timeout_sec);
|
||||
|
||||
bool protocol_receive_n_data(ProtocolSession* session, void* data, size_t data_size) {
|
||||
return protocol_receive_n_data_timed(session, data, data_size, RECEIVE_TIMEOUT_SEC);
|
||||
}
|
||||
|
||||
bool protocol_receive_n_data_timed(ProtocolSession* session, void* data, size_t data_size,
|
||||
int timeout_sec) {
|
||||
log_debug_message(LOG_DEBUG_IO, " Receiving n Data: %zu", data_size);
|
||||
if (!session)
|
||||
return false;
|
||||
int fd = session->read_fd;
|
||||
if (timeout_sec <= 0)
|
||||
timeout_sec = RECEIVE_TIMEOUT_SEC;
|
||||
|
||||
struct timespec deadline;
|
||||
clock_gettime(CLOCK_MONOTONIC, &deadline);
|
||||
deadline.tv_sec += RECEIVE_TIMEOUT_SEC;
|
||||
deadline.tv_sec += timeout_sec;
|
||||
|
||||
size_t total_bytes_received = 0;
|
||||
short wait_events = POLLIN;
|
||||
@@ -319,7 +329,7 @@ bool protocol_receive_n_data(ProtocolSession* session, void* data, size_t data_s
|
||||
struct pollfd pfd = {.fd = fd, .events = wait_events};
|
||||
int poll_result = poll(&pfd, 1, deadline_remaining_ms(&deadline));
|
||||
if (poll_result == 0) {
|
||||
log_message(LOG_LEVEL_ERROR, "Receive timeout after %ds", RECEIVE_TIMEOUT_SEC);
|
||||
log_message(LOG_LEVEL_ERROR, "Receive timeout after %ds", timeout_sec);
|
||||
return false;
|
||||
}
|
||||
if (poll_result < 0) {
|
||||
@@ -516,6 +526,17 @@ bool protocol_receive_status(ProtocolSession* session, Status* status) {
|
||||
return true;
|
||||
}
|
||||
|
||||
/* protocol_receive_status with an explicit per-message deadline (seconds).
|
||||
Used where a single reply may legitimately take far longer than the default
|
||||
60 s receive window - e.g. the sender waiting for the early-delete ACK after
|
||||
the receiver committed a large (up to MAX_SERVER_DELETE_COUNT) deletion. */
|
||||
bool protocol_receive_status_timed(ProtocolSession* session, Status* status, int timeout_sec) {
|
||||
if (!protocol_receive_n_data_timed(session, status, sizeof(Status), timeout_sec))
|
||||
return false;
|
||||
log_debug_message(LOG_DEBUG_PROTO, "Received Status: %s", status_to_string(*status));
|
||||
return true;
|
||||
}
|
||||
|
||||
bool send_str(int fd, const char* data) {
|
||||
return protocol_send_str(legacy_session(-1, fd), data);
|
||||
}
|
||||
@@ -543,3 +564,6 @@ bool send_status(int fd, Status status) {
|
||||
bool receive_status(int fd, Status* status) {
|
||||
return protocol_receive_status(legacy_session(fd, -1), status);
|
||||
}
|
||||
bool receive_status_timed(int fd, Status* status, int timeout_sec) {
|
||||
return protocol_receive_status_timed(legacy_session(fd, -1), status, timeout_sec);
|
||||
}
|
||||
@@ -112,5 +112,10 @@ bool send_int(int file_descriptor, int data);
|
||||
bool receive_int(int file_descriptor, int* data);
|
||||
bool send_status(int file_descriptor, Status status);
|
||||
bool receive_status(int file_descriptor, Status* status);
|
||||
/* receive_status with an explicit per-message deadline in seconds, instead of
|
||||
the default RECEIVE_TIMEOUT_SEC. A reply that may legitimately take longer
|
||||
(e.g. the early-delete ACK after a large receiver-side deletion) must use
|
||||
this so the sender does not abort after the deletion already committed. */
|
||||
bool receive_status_timed(int file_descriptor, Status* status, int timeout_sec);
|
||||
|
||||
#endif
|
||||
@@ -2052,10 +2052,14 @@ class TestDeleteTiming:
|
||||
f"{flag}: nested file was not written after the early deletion"
|
||||
|
||||
@pytest.mark.parametrize("flag", ["--delete", "--delete-after", "--delete-delay"])
|
||||
def test_late_flags_commit_only_after_success(self, flag):
|
||||
@pytest.mark.parametrize("mt", [False, True])
|
||||
def test_late_flags_commit_only_after_success(self, flag, mt):
|
||||
"""Plain --delete/--delete-after/--delete-delay defer deletion until the
|
||||
whole transfer succeeds: a mid-transfer write failure must leave every
|
||||
extra in place (commit-style safety)."""
|
||||
extra in place (commit-style safety). The -m receiver must also keep
|
||||
the extras: the deferred keep-set is committed by the server only after
|
||||
the disk-writer thread has finished, and a failing writer means the
|
||||
manifest is freed, never applied."""
|
||||
source = self._seed("late")
|
||||
dest = os.path.join(TEST_DATA_DIR, "deltiming_late_dst")
|
||||
clean_dir(dest)
|
||||
@@ -2072,13 +2076,14 @@ class TestDeleteTiming:
|
||||
with open(blocker, "wb") as fh:
|
||||
fh.write(b"blocks the nested destination directory")
|
||||
|
||||
result, _ = run_client(source, dest, flags=[flag], port=server.port)
|
||||
flags = [flag] + (["-m"] if mt else [])
|
||||
result, _ = run_client(source, dest, flags=flags, port=server.port)
|
||||
assert result.returncode != 0, \
|
||||
f"{flag} unexpectedly succeeded (deletion must be deferred)"
|
||||
f"{flag} (mt={mt}) unexpectedly succeeded (deletion must be deferred)"
|
||||
assert os.path.exists(extra), \
|
||||
f"{flag} removed an extra although the transfer failed"
|
||||
f"{flag} (mt={mt}) removed an extra although the transfer failed"
|
||||
assert os.path.isfile(blocker), \
|
||||
f"{flag} deleted the blocker although the transfer failed"
|
||||
f"{flag} (mt={mt}) deleted the blocker although the transfer failed"
|
||||
|
||||
def test_early_flag_respected_when_server_refuses_delete(self, shared_server):
|
||||
"""With an --allow-delete-less server the client's early timing still
|
||||
|
||||
@@ -412,6 +412,25 @@ static void test_protocol_accounting_release_does_not_underflow() {
|
||||
protocol_session_unbind();
|
||||
}
|
||||
|
||||
static void test_send_receive_status_timed() {
|
||||
int p[2];
|
||||
EXPECT_EQ_INT(pipe(p), 0);
|
||||
io_set_fds(p[0], p[1]);
|
||||
io_set_bwlimit(0);
|
||||
|
||||
/* The extended-deadline variant must read an ordinary status just like the
|
||||
default window, and must fail cleanly on EOF rather than block. */
|
||||
EXPECT_TRUE(send_status(0, STATUS_OK));
|
||||
Status received = -1;
|
||||
EXPECT_TRUE(receive_status_timed(0, &received, 5));
|
||||
EXPECT_EQ_INT((int)received, (int)STATUS_OK);
|
||||
|
||||
close(p[1]);
|
||||
EXPECT_FALSE(receive_status_timed(0, &received, 5));
|
||||
|
||||
close(p[0]);
|
||||
}
|
||||
|
||||
void test_protocol() {
|
||||
test_send_receive_n_data();
|
||||
test_send_receive_n_data_zero();
|
||||
@@ -421,6 +440,7 @@ void test_protocol() {
|
||||
test_send_receive_data();
|
||||
test_send_receive_int();
|
||||
test_send_receive_status();
|
||||
test_send_receive_status_timed();
|
||||
test_receive_n_data_truncated();
|
||||
test_receive_str_truncated();
|
||||
test_max_alloc_rejects_single_buffer();
|
||||
|
||||
@@ -460,6 +460,95 @@ static void test_incremental_check_delta_oversize_reports_failure() {
|
||||
}
|
||||
}
|
||||
|
||||
/* Late-timing keep-set leak guard: a manifest parked by the commit path must
|
||||
be freed on every error exit, never leaked. These tests drive
|
||||
receiver_process_pending() through an error AFTER the manifest was parked and
|
||||
are exercised under ASan/valgrind to prove the list is released. */
|
||||
|
||||
static Config* make_late_delete_config(const char* root) {
|
||||
Config* cfg = config_create();
|
||||
if (!cfg)
|
||||
return NULL;
|
||||
cfg->send_directory = str_dup("/src");
|
||||
cfg->receive_root_directory = str_dup(root);
|
||||
cfg->use_delete = true;
|
||||
cfg->delete_after = true;
|
||||
return cfg;
|
||||
}
|
||||
|
||||
static int run_pending_receiver(Config* cfg, int fd, ArrayList** pending) {
|
||||
ReceiverSink sink = {0};
|
||||
return receiver_process_pending(cfg, fd, &sink, pending);
|
||||
}
|
||||
|
||||
static void test_late_manifest_abort_frees_keepset() {
|
||||
Config* cfg = make_late_delete_config("/tmp/fastsync_late_abort");
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
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);
|
||||
|
||||
EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST));
|
||||
EXPECT_TRUE(send_int(p[1], 1));
|
||||
EXPECT_TRUE(send_str(p[1], "keep.txt"));
|
||||
EXPECT_TRUE(send_status(p[1], STATUS_ABORT));
|
||||
|
||||
ArrayList* pending = NULL;
|
||||
EXPECT_EQ_INT(run_pending_receiver(cfg, p[0], &pending), -1);
|
||||
EXPECT_NULL(pending);
|
||||
|
||||
close(p[0]);
|
||||
close(p[1]);
|
||||
config_delete(cfg);
|
||||
}
|
||||
|
||||
static void test_late_manifest_eof_frees_keepset() {
|
||||
Config* cfg = make_late_delete_config("/tmp/fastsync_late_eof");
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
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);
|
||||
|
||||
EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST));
|
||||
EXPECT_TRUE(send_int(p[1], 1));
|
||||
EXPECT_TRUE(send_str(p[1], "keep.txt"));
|
||||
shutdown(p[1], SHUT_WR);
|
||||
|
||||
ArrayList* pending = NULL;
|
||||
EXPECT_EQ_INT(run_pending_receiver(cfg, p[0], &pending), -1);
|
||||
EXPECT_NULL(pending);
|
||||
|
||||
close(p[0]);
|
||||
close(p[1]);
|
||||
config_delete(cfg);
|
||||
}
|
||||
|
||||
static void test_late_second_manifest_frees_both() {
|
||||
Config* cfg = make_late_delete_config("/tmp/fastsync_late_second");
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
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);
|
||||
|
||||
EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST));
|
||||
EXPECT_TRUE(send_int(p[1], 1));
|
||||
EXPECT_TRUE(send_str(p[1], "first.txt"));
|
||||
EXPECT_TRUE(send_status(p[1], STATUS_MANIFEST));
|
||||
EXPECT_TRUE(send_int(p[1], 1));
|
||||
EXPECT_TRUE(send_str(p[1], "second.txt"));
|
||||
|
||||
ArrayList* pending = NULL;
|
||||
EXPECT_EQ_INT(run_pending_receiver(cfg, p[0], &pending), -1);
|
||||
EXPECT_NULL(pending);
|
||||
|
||||
close(p[0]);
|
||||
close(p[1]);
|
||||
config_delete(cfg);
|
||||
}
|
||||
|
||||
void test_server() {
|
||||
if (!is_running_under_valgrind()) {
|
||||
test_receive_files_finished();
|
||||
@@ -470,5 +559,8 @@ void test_server() {
|
||||
test_incremental_check_quick_skip_by_mtime();
|
||||
test_incremental_check_size_mismatch_full_transfer();
|
||||
test_incremental_check_delta_oversize_reports_failure();
|
||||
test_late_manifest_abort_frees_keepset();
|
||||
test_late_manifest_eof_frees_keepset();
|
||||
test_late_second_manifest_frees_both();
|
||||
}
|
||||
}
|
||||
Reference in new issue
Block a user