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
|
the closest safe approximation their engine mode allows; the divergences are
|
||||||
noted in the rows above.
|
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
|
## 8. Metadata Preservation
|
||||||
|
|
||||||
| Flag | Rsync Description | FastSync Status | Notes |
|
| 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
|
--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
|
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
|
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) {
|
static bool send_delete_manifest_early(Client* client, ArrayList* manifest) {
|
||||||
if (!client || !manifest)
|
if (!client || !manifest)
|
||||||
return false;
|
return false;
|
||||||
if (send_delete_manifest(client->file_descriptor, manifest) != 0)
|
if (send_delete_manifest(client->file_descriptor, manifest) != 0)
|
||||||
return false;
|
return false;
|
||||||
Status ack;
|
Status ack;
|
||||||
if (!receive_status(client->file_descriptor, &ack))
|
if (!receive_status_timed(client->file_descriptor, &ack, DELETE_ACK_TIMEOUT_SEC))
|
||||||
return false;
|
return false;
|
||||||
if (ack != STATUS_OK) {
|
if (ack != STATUS_OK) {
|
||||||
log_message(LOG_LEVEL_ERROR, "Server failed to delete files before the transfer");
|
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(" transfer has succeeded)\n");
|
||||||
printf(" --delete-before Delete extras before the transfer starts\n");
|
printf(" --delete-before Delete extras before the transfer starts\n");
|
||||||
printf(" (implies --delete)\n");
|
printf(" (implies --delete)\n");
|
||||||
printf(" --delete-during Delete extras once the keep-set is known, before\n");
|
printf(" --delete-during Delete extras once the keep-set manifest is known,\n");
|
||||||
printf(" --del data is applied (alias --del; implies --delete)\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(" --delete-delay Delete extras only after a successful transfer\n");
|
||||||
printf(" (implies --delete)\n");
|
printf(" (implies --delete)\n");
|
||||||
printf(" --delete-after Alias of the default --delete timing: delete only\n");
|
printf(" --delete-after Delete only after the whole transfer succeeded\n");
|
||||||
printf(" after the transfer succeeded (implies --delete)\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(" --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(" --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");
|
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))
|
if (!receive_status(file_descriptor, &status))
|
||||||
return -1;
|
return -1;
|
||||||
bool early_delete = config_delete_timing_early(config);
|
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;
|
ArrayList* deferred_manifest = NULL;
|
||||||
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
|
while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK ||
|
||||||
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH ||
|
status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH ||
|
||||||
status == STATUS_MKDIR || status == STATUS_MANIFEST) {
|
status == STATUS_MKDIR || status == STATUS_MANIFEST) {
|
||||||
if (status == STATUS_KEEPALIVE) {
|
if (status == STATUS_KEEPALIVE) {
|
||||||
if (!send_status(file_descriptor, STATUS_KEEPALIVE))
|
if (!send_status(file_descriptor, STATUS_KEEPALIVE))
|
||||||
return -1;
|
goto fail;
|
||||||
goto next_status;
|
goto next_status;
|
||||||
}
|
}
|
||||||
if (status == STATUS_ABORT) {
|
if (status == STATUS_ABORT) {
|
||||||
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
|
log_message(LOG_LEVEL_INFO, "Received abort from client, cleaning up");
|
||||||
return -1;
|
goto fail;
|
||||||
}
|
}
|
||||||
if (status == STATUS_CHECK) {
|
if (status == STATUS_CHECK) {
|
||||||
bool skipped;
|
bool skipped;
|
||||||
@@ -172,7 +175,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
|||||||
goto receive_error;
|
goto receive_error;
|
||||||
} else if (status == STATUS_CHECK_BATCH) {
|
} else if (status == STATUS_CHECK_BATCH) {
|
||||||
if (!receiver_process_batch(config, file_descriptor))
|
if (!receiver_process_batch(config, file_descriptor))
|
||||||
return -1;
|
goto fail;
|
||||||
goto next_status;
|
goto next_status;
|
||||||
} else if (status == STATUS_MKDIR) {
|
} else if (status == STATUS_MKDIR) {
|
||||||
File* dir = file_receive_directory(file_descriptor);
|
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) {
|
} else if (status == STATUS_MANIFEST) {
|
||||||
ArrayList* manifest = receive_manifest_entries(file_descriptor);
|
ArrayList* manifest = receive_manifest_entries(file_descriptor);
|
||||||
if (!manifest)
|
if (!manifest)
|
||||||
return -1; /* receive_manifest_entries already sent STATUS_ERROR */
|
goto fail; /* receive_manifest_entries already sent STATUS_ERROR */
|
||||||
if (early_delete) {
|
if (early_delete) {
|
||||||
/* --delete-before / --delete-during: the manifest is authoritative the
|
/* --delete-before / --delete-during: the manifest is authoritative the
|
||||||
moment it arrives, before any file data. Delete now and acknowledge
|
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);
|
array_list_delete(manifest);
|
||||||
if (!deletion_ok) {
|
if (!deletion_ok) {
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
send_status(file_descriptor, STATUS_ERROR);
|
||||||
return -1;
|
goto fail;
|
||||||
}
|
}
|
||||||
if (!send_status(file_descriptor, STATUS_OK))
|
if (!send_status(file_descriptor, STATUS_OK))
|
||||||
return -1;
|
goto fail;
|
||||||
} else {
|
} else if (config->use_delete) {
|
||||||
/* Plain --delete / --delete-after / --delete-delay: hold the keep-set
|
/* Plain --delete / --delete-after / --delete-delay: hold the keep-set
|
||||||
and commit the deletion only after STATUS_FINISHED. */
|
and commit the deletion only after STATUS_FINISHED. */
|
||||||
if (config->use_delete) {
|
if (deferred_manifest) {
|
||||||
if (deferred_manifest) {
|
log_message(LOG_LEVEL_ERROR, "Received a second delete manifest");
|
||||||
log_message(LOG_LEVEL_ERROR, "Received a second delete manifest");
|
array_list_delete(deferred_manifest);
|
||||||
array_list_delete(manifest);
|
deferred_manifest = NULL;
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
|
||||||
return -1;
|
|
||||||
}
|
|
||||||
deferred_manifest = manifest;
|
|
||||||
} else {
|
|
||||||
array_list_delete(manifest);
|
array_list_delete(manifest);
|
||||||
|
send_status(file_descriptor, STATUS_ERROR);
|
||||||
|
goto fail;
|
||||||
}
|
}
|
||||||
|
deferred_manifest = manifest;
|
||||||
|
} else {
|
||||||
|
array_list_delete(manifest);
|
||||||
}
|
}
|
||||||
goto next_status;
|
goto next_status;
|
||||||
} else {
|
} else {
|
||||||
@@ -241,26 +244,41 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
|||||||
if (deferred_manifest) {
|
if (deferred_manifest) {
|
||||||
if (pending_manifest) {
|
if (pending_manifest) {
|
||||||
*pending_manifest = deferred_manifest;
|
*pending_manifest = deferred_manifest;
|
||||||
|
deferred_manifest = NULL;
|
||||||
} else {
|
} else {
|
||||||
bool deletion_ok = manifest_delete_extras(config, deferred_manifest);
|
bool deletion_ok = manifest_delete_extras(config, deferred_manifest);
|
||||||
array_list_delete(deferred_manifest);
|
array_list_delete(deferred_manifest);
|
||||||
|
deferred_manifest = NULL;
|
||||||
if (!deletion_ok) {
|
if (!deletion_ok) {
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
send_status(file_descriptor, STATUS_ERROR);
|
||||||
return -1;
|
goto fail;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
if (sink->send_success) {
|
if (sink->send_success) {
|
||||||
if (sink->send_success_frame) {
|
if (sink->send_success_frame) {
|
||||||
if (!sink->send_success_frame(file_descriptor, sink->context))
|
if (!sink->send_success_frame(file_descriptor, sink->context))
|
||||||
return -1;
|
goto fail;
|
||||||
} else if (!send_status(file_descriptor, STATUS_OK)) {
|
} else if (!send_status(file_descriptor, STATUS_OK)) {
|
||||||
return -1;
|
goto fail;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
return 0;
|
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:
|
receive_error:
|
||||||
|
if (deferred_manifest) {
|
||||||
|
array_list_delete(deferred_manifest);
|
||||||
|
deferred_manifest = NULL;
|
||||||
|
}
|
||||||
if (sink->send_error)
|
if (sink->send_error)
|
||||||
send_status(file_descriptor, STATUS_ERROR);
|
send_status(file_descriptor, STATUS_ERROR);
|
||||||
return -1;
|
return -1;
|
||||||
|
|||||||
+26
-2
@@ -302,15 +302,25 @@ bool protocol_send_n_data(ProtocolSession* session, const void* data, size_t dat
|
|||||||
return true;
|
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) {
|
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);
|
log_debug_message(LOG_DEBUG_IO, " Receiving n Data: %zu", data_size);
|
||||||
if (!session)
|
if (!session)
|
||||||
return false;
|
return false;
|
||||||
int fd = session->read_fd;
|
int fd = session->read_fd;
|
||||||
|
if (timeout_sec <= 0)
|
||||||
|
timeout_sec = RECEIVE_TIMEOUT_SEC;
|
||||||
|
|
||||||
struct timespec deadline;
|
struct timespec deadline;
|
||||||
clock_gettime(CLOCK_MONOTONIC, &deadline);
|
clock_gettime(CLOCK_MONOTONIC, &deadline);
|
||||||
deadline.tv_sec += RECEIVE_TIMEOUT_SEC;
|
deadline.tv_sec += timeout_sec;
|
||||||
|
|
||||||
size_t total_bytes_received = 0;
|
size_t total_bytes_received = 0;
|
||||||
short wait_events = POLLIN;
|
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};
|
struct pollfd pfd = {.fd = fd, .events = wait_events};
|
||||||
int poll_result = poll(&pfd, 1, deadline_remaining_ms(&deadline));
|
int poll_result = poll(&pfd, 1, deadline_remaining_ms(&deadline));
|
||||||
if (poll_result == 0) {
|
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;
|
return false;
|
||||||
}
|
}
|
||||||
if (poll_result < 0) {
|
if (poll_result < 0) {
|
||||||
@@ -516,6 +526,17 @@ bool protocol_receive_status(ProtocolSession* session, Status* status) {
|
|||||||
return true;
|
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) {
|
bool send_str(int fd, const char* data) {
|
||||||
return protocol_send_str(legacy_session(-1, fd), 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) {
|
bool receive_status(int fd, Status* status) {
|
||||||
return protocol_receive_status(legacy_session(fd, -1), 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 receive_int(int file_descriptor, int* data);
|
||||||
bool send_status(int file_descriptor, Status status);
|
bool send_status(int file_descriptor, Status status);
|
||||||
bool receive_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
|
#endif
|
||||||
@@ -2052,10 +2052,14 @@ class TestDeleteTiming:
|
|||||||
f"{flag}: nested file was not written after the early deletion"
|
f"{flag}: nested file was not written after the early deletion"
|
||||||
|
|
||||||
@pytest.mark.parametrize("flag", ["--delete", "--delete-after", "--delete-delay"])
|
@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
|
"""Plain --delete/--delete-after/--delete-delay defer deletion until the
|
||||||
whole transfer succeeds: a mid-transfer write failure must leave every
|
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")
|
source = self._seed("late")
|
||||||
dest = os.path.join(TEST_DATA_DIR, "deltiming_late_dst")
|
dest = os.path.join(TEST_DATA_DIR, "deltiming_late_dst")
|
||||||
clean_dir(dest)
|
clean_dir(dest)
|
||||||
@@ -2072,13 +2076,14 @@ class TestDeleteTiming:
|
|||||||
with open(blocker, "wb") as fh:
|
with open(blocker, "wb") as fh:
|
||||||
fh.write(b"blocks the nested destination directory")
|
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, \
|
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), \
|
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), \
|
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):
|
def test_early_flag_respected_when_server_refuses_delete(self, shared_server):
|
||||||
"""With an --allow-delete-less server the client's early timing still
|
"""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();
|
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() {
|
void test_protocol() {
|
||||||
test_send_receive_n_data();
|
test_send_receive_n_data();
|
||||||
test_send_receive_n_data_zero();
|
test_send_receive_n_data_zero();
|
||||||
@@ -421,6 +440,7 @@ void test_protocol() {
|
|||||||
test_send_receive_data();
|
test_send_receive_data();
|
||||||
test_send_receive_int();
|
test_send_receive_int();
|
||||||
test_send_receive_status();
|
test_send_receive_status();
|
||||||
|
test_send_receive_status_timed();
|
||||||
test_receive_n_data_truncated();
|
test_receive_n_data_truncated();
|
||||||
test_receive_str_truncated();
|
test_receive_str_truncated();
|
||||||
test_max_alloc_rejects_single_buffer();
|
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() {
|
void test_server() {
|
||||||
if (!is_running_under_valgrind()) {
|
if (!is_running_under_valgrind()) {
|
||||||
test_receive_files_finished();
|
test_receive_files_finished();
|
||||||
@@ -470,5 +559,8 @@ void test_server() {
|
|||||||
test_incremental_check_quick_skip_by_mtime();
|
test_incremental_check_quick_skip_by_mtime();
|
||||||
test_incremental_check_size_mismatch_full_transfer();
|
test_incremental_check_size_mismatch_full_transfer();
|
||||||
test_incremental_check_delta_oversize_reports_failure();
|
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