Release v2.26.0 #284
+23
-8
@@ -34,20 +34,35 @@
|
|||||||
* client_send.c, still links the symbol. */
|
* client_send.c, still links the symbol. */
|
||||||
volatile sig_atomic_t client_abort_requested = 0;
|
volatile sig_atomic_t client_abort_requested = 0;
|
||||||
|
|
||||||
#ifndef FASTSYNC_TEST_BUILD
|
/* Only armed while a network transfer is in flight. Outside that window the
|
||||||
/* Signal handler: perform NO work beyond storing the flag. Logging, protocol
|
* handler restores the default disposition and re-raises, so purely local modes
|
||||||
* I/O and the STATUS_ABORT frame are all done later on the normal send path,
|
* (--list-only/--dry-run/--read-batch/--only-write-batch and the batch-emission
|
||||||
* which is not async-signal-safe. Only the production client installs it. */
|
* pass) keep terminating on Ctrl-C/SIGTERM instead of silently swallowing it. */
|
||||||
static void client_signal_handler(int signo) {
|
volatile sig_atomic_t client_abort_armed = 0;
|
||||||
(void)signo;
|
|
||||||
client_abort_requested = 1;
|
void client_set_abort_armed(bool armed) {
|
||||||
|
client_abort_armed = armed ? 1 : 0;
|
||||||
}
|
}
|
||||||
#endif
|
|
||||||
|
|
||||||
bool client_abort_pending(void) {
|
bool client_abort_pending(void) {
|
||||||
return client_abort_requested != 0;
|
return client_abort_requested != 0;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
#ifndef FASTSYNC_TEST_BUILD
|
||||||
|
/* Signal handler: perform NO work beyond storing the flag. Logging, protocol
|
||||||
|
* I/O and the STATUS_ABORT frame are all done later on the normal send path,
|
||||||
|
* which is not async-signal-safe. When no transfer is armed, fall back to the
|
||||||
|
* default action so local-only modes remain interruptible. */
|
||||||
|
static void client_signal_handler(int signo) {
|
||||||
|
if (!client_abort_armed) {
|
||||||
|
signal(signo, SIG_DFL);
|
||||||
|
raise(signo);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
client_abort_requested = 1;
|
||||||
|
}
|
||||||
|
#endif
|
||||||
|
|
||||||
#ifndef FASTSYNC_TEST_BUILD
|
#ifndef FASTSYNC_TEST_BUILD
|
||||||
/* Parse environment variables for source/destination directories and save-to-disk flag. */
|
/* Parse environment variables for source/destination directories and save-to-disk flag. */
|
||||||
static void parse_environment(const char** out_env_source, const char** out_env_dest,
|
static void parse_environment(const char** out_env_source, const char** out_env_dest,
|
||||||
|
|||||||
@@ -1949,6 +1949,9 @@ int send_files(Config* config) {
|
|||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/* From here on a server session may be live, so Ctrl-C/SIGTERM should set the
|
||||||
|
abort flag (and be forwarded as STATUS_ABORT) instead of terminating. */
|
||||||
|
client_set_abort_armed(true);
|
||||||
Client* client = connect_transfer_client(config);
|
Client* client = connect_transfer_client(config);
|
||||||
if (!client) {
|
if (!client) {
|
||||||
if (config->transport == TRANSPORT_TCP)
|
if (config->transport == TRANSPORT_TCP)
|
||||||
@@ -2228,6 +2231,7 @@ send_fail:
|
|||||||
prepared_scanner_destroy(&prepared);
|
prepared_scanner_destroy(&prepared);
|
||||||
disconnect_transfer_client(client);
|
disconnect_transfer_client(client);
|
||||||
protocol_session_unbind();
|
protocol_session_unbind();
|
||||||
|
client_set_abort_armed(false);
|
||||||
return ret;
|
return ret;
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -2257,6 +2261,9 @@ int send_files_multithreaded(Config** config_ptr) {
|
|||||||
return 1;
|
return 1;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/* Armed only once a session may go live (see send_files). */
|
||||||
|
client_set_abort_armed(true);
|
||||||
|
|
||||||
long pages = sysconf(_SC_AVPHYS_PAGES);
|
long pages = sysconf(_SC_AVPHYS_PAGES);
|
||||||
long page_size = sysconf(_SC_PAGE_SIZE);
|
long page_size = sysconf(_SC_PAGE_SIZE);
|
||||||
unsigned long long available_memory =
|
unsigned long long available_memory =
|
||||||
@@ -2411,5 +2418,6 @@ int send_files_multithreaded(Config** config_ptr) {
|
|||||||
/* --ignore-errors: the run completed (and deleted) past an unreadable source
|
/* --ignore-errors: the run completed (and deleted) past an unreadable source
|
||||||
directory; report it as errored like rsync does. */
|
directory; report it as errored like rsync does. */
|
||||||
pipeline_context_sender_destroy(context);
|
pipeline_context_sender_destroy(context);
|
||||||
|
client_set_abort_armed(false);
|
||||||
return sender_ok && !scan_io ? 0 : 1;
|
return sender_ok && !scan_io ? 0 : 1;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -13,6 +13,10 @@
|
|||||||
* receiver can clean up before the client exits. */
|
* receiver can clean up before the client exits. */
|
||||||
extern volatile sig_atomic_t client_abort_requested;
|
extern volatile sig_atomic_t client_abort_requested;
|
||||||
bool client_abort_pending(void);
|
bool client_abort_pending(void);
|
||||||
|
/* Arm/disarm abort handling around the network phase. While disarmed, a
|
||||||
|
* SIGINT/SIGTERM takes the default action (immediate termination) so local-only
|
||||||
|
* modes are not left unresponsive. Defined in client_cli.c. */
|
||||||
|
void client_set_abort_armed(bool armed);
|
||||||
|
|
||||||
/* Both sender entry points BORROW `config` for the duration of the call; they
|
/* Both sender entry points BORROW `config` for the duration of the call; they
|
||||||
* never free it, and the caller retains ownership (freeing it with
|
* never free it, and the caller retains ownership (freeing it with
|
||||||
|
|||||||
+26
-5
@@ -303,6 +303,12 @@ bool protocol_send_n_data(ProtocolSession* session, const void* data, size_t dat
|
|||||||
wait_events = ssl_err == SSL_ERROR_WANT_WRITE ? POLLOUT : POLLIN;
|
wait_events = ssl_err == SSL_ERROR_WANT_WRITE ? POLLOUT : POLLIN;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
/* A signal (e.g. Ctrl-C) interrupts the blocking TLS write: retry so
|
||||||
|
the send loop can observe the abort flag at the next checkpoint. */
|
||||||
|
if (ssl_err == SSL_ERROR_SYSCALL && errno == EINTR)
|
||||||
|
continue;
|
||||||
|
} else if (errno == EINTR) {
|
||||||
|
continue;
|
||||||
}
|
}
|
||||||
log_message(LOG_LEVEL_ERROR, "Could not send data");
|
log_message(LOG_LEVEL_ERROR, "Could not send data");
|
||||||
return false;
|
return false;
|
||||||
@@ -618,6 +624,7 @@ static bool protocol_read_status_until(ProtocolSession* session, Status* status,
|
|||||||
const struct timespec* deadline) {
|
const struct timespec* deadline) {
|
||||||
Status received = STATUS_ERROR;
|
Status received = STATUS_ERROR;
|
||||||
size_t got = 0;
|
size_t got = 0;
|
||||||
|
short wait_events = POLLIN;
|
||||||
while (got < sizeof(Status)) {
|
while (got < sizeof(Status)) {
|
||||||
if (!session->ssl || SSL_pending(session->ssl) == 0) {
|
if (!session->ssl || SSL_pending(session->ssl) == 0) {
|
||||||
int remaining_ms = deadline_remaining_ms(deadline);
|
int remaining_ms = deadline_remaining_ms(deadline);
|
||||||
@@ -625,7 +632,7 @@ static bool protocol_read_status_until(ProtocolSession* session, Status* status,
|
|||||||
log_message(LOG_LEVEL_ERROR, "Receive timeout while reading status");
|
log_message(LOG_LEVEL_ERROR, "Receive timeout while reading status");
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
struct pollfd pfd = {.fd = session->read_fd, .events = POLLIN};
|
struct pollfd pfd = {.fd = session->read_fd, .events = wait_events};
|
||||||
int poll_result = poll(&pfd, 1, remaining_ms);
|
int poll_result = poll(&pfd, 1, remaining_ms);
|
||||||
if (poll_result == 0) {
|
if (poll_result == 0) {
|
||||||
log_message(LOG_LEVEL_ERROR, "Receive timeout while reading status");
|
log_message(LOG_LEVEL_ERROR, "Receive timeout while reading status");
|
||||||
@@ -647,9 +654,11 @@ static bool protocol_read_status_until(ProtocolSession* session, Status* status,
|
|||||||
if (bytes_received <= 0) {
|
if (bytes_received <= 0) {
|
||||||
if (session->ssl) {
|
if (session->ssl) {
|
||||||
int ssl_err = SSL_get_error(session->ssl, (int)bytes_received);
|
int ssl_err = SSL_get_error(session->ssl, (int)bytes_received);
|
||||||
if (ssl_err == SSL_ERROR_WANT_READ || ssl_err == SSL_ERROR_WANT_WRITE)
|
if (ssl_err == SSL_ERROR_WANT_READ || ssl_err == SSL_ERROR_WANT_WRITE) {
|
||||||
|
wait_events = ssl_err == SSL_ERROR_WANT_WRITE ? POLLOUT : POLLIN;
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
|
}
|
||||||
if (bytes_received < 0 && errno == EINTR)
|
if (bytes_received < 0 && errno == EINTR)
|
||||||
continue;
|
continue;
|
||||||
log_message(LOG_LEVEL_ERROR, "Connection closed while receiving status");
|
log_message(LOG_LEVEL_ERROR, "Connection closed while receiving status");
|
||||||
@@ -689,7 +698,8 @@ bool protocol_receive_status_keepalive(ProtocolSession* session, Status* status,
|
|||||||
/* Only interleave a keepalive while waiting for the FIRST byte of a
|
/* Only interleave a keepalive while waiting for the FIRST byte of a
|
||||||
* frame; once part of a frame is buffered a write could race the peer's
|
* frame; once part of a frame is buffered a write could race the peer's
|
||||||
* reply into the middle of it. */
|
* reply into the middle of it. */
|
||||||
int interval_ms = keepalive_interval_sec * 1000;
|
long long interval_ms_ll = (long long)keepalive_interval_sec * 1000LL;
|
||||||
|
int interval_ms = interval_ms_ll > INT_MAX ? INT_MAX : (int)interval_ms_ll;
|
||||||
int wait_ms = interval_ms < remaining_ms ? interval_ms : remaining_ms;
|
int wait_ms = interval_ms < remaining_ms ? interval_ms : remaining_ms;
|
||||||
struct pollfd pfd = {.fd = session->read_fd, .events = POLLIN};
|
struct pollfd pfd = {.fd = session->read_fd, .events = POLLIN};
|
||||||
int poll_result = poll(&pfd, 1, wait_ms);
|
int poll_result = poll(&pfd, 1, wait_ms);
|
||||||
@@ -724,16 +734,27 @@ bool protocol_receive_status_keepalive(ProtocolSession* session, Status* status,
|
|||||||
* was busy. It answers them only after the real status, so leaving them
|
* was busy. It answers them only after the real status, so leaving them
|
||||||
* unread would put stale KEEPALIVE frames ahead of the next exchange and
|
* unread would put stale KEEPALIVE frames ahead of the next exchange and
|
||||||
* desynchronize the protocol. */
|
* desynchronize the protocol. */
|
||||||
|
if (replies_seen < keepalives_sent) {
|
||||||
|
/* A short separate grace, not the (possibly exhausted) main deadline: the
|
||||||
|
terminal status already arrived, so a peer that never answers its owed
|
||||||
|
keepalives must not turn a successful ack into a reported failure. */
|
||||||
|
struct timespec drain_deadline;
|
||||||
|
clock_gettime(CLOCK_MONOTONIC, &drain_deadline);
|
||||||
|
drain_deadline.tv_sec += 1;
|
||||||
while (replies_seen < keepalives_sent) {
|
while (replies_seen < keepalives_sent) {
|
||||||
Status drained;
|
Status drained;
|
||||||
if (!protocol_read_status_until(session, &drained, &deadline))
|
if (!protocol_read_status_until(session, &drained, &drain_deadline)) {
|
||||||
return false;
|
log_message(LOG_LEVEL_WARNING, "peer did not answer %lu keepalive(s); continuing",
|
||||||
|
keepalives_sent - replies_seen);
|
||||||
|
break;
|
||||||
|
}
|
||||||
if (drained != STATUS_KEEPALIVE) {
|
if (drained != STATUS_KEEPALIVE) {
|
||||||
log_message(LOG_LEVEL_ERROR, "Unexpected status while draining keepalive replies");
|
log_message(LOG_LEVEL_ERROR, "Unexpected status while draining keepalive replies");
|
||||||
return false;
|
return false;
|
||||||
}
|
}
|
||||||
replies_seen++;
|
replies_seen++;
|
||||||
}
|
}
|
||||||
|
}
|
||||||
*status = final;
|
*status = final;
|
||||||
log_debug_message(LOG_DEBUG_PROTO, "Received Status: %s", status_to_string(*status));
|
log_debug_message(LOG_DEBUG_PROTO, "Received Status: %s", status_to_string(*status));
|
||||||
return true;
|
return true;
|
||||||
|
|||||||
@@ -81,9 +81,12 @@ int LLVMFuzzerTestOneInput(const uint8_t* data, size_t size) {
|
|||||||
close(w);
|
close(w);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/* Bounded so a crafted 256 MiB length header cannot make each iteration
|
||||||
|
allocate the full MAX_DATA_PAYLOAD_SIZE under ASan; the framing logic is
|
||||||
|
identical to receive_data(), which delegates to the limited variant. */
|
||||||
rd = make_stream(data, size, &w);
|
rd = make_stream(data, size, &w);
|
||||||
if (rd >= 0) {
|
if (rd >= 0) {
|
||||||
Data* d = receive_data(rd);
|
Data* d = receive_data_limited(rd, 1u << 20);
|
||||||
data_destroy(d);
|
data_destroy(d);
|
||||||
close(rd);
|
close(rd);
|
||||||
close(w);
|
close(w);
|
||||||
|
|||||||
Reference in New Issue
Block a user