From c2df0347ef6423acef7a911e337aef3c52d36f0e Mon Sep 17 00:00:00 2001 From: TapTap Date: Sun, 13 Sep 2026 06:59:02 +0200 Subject: [PATCH] fix(client,protocol): EINTR-safe sends, armed abort, keepalive drain grace, TLS WANT_WRITE --- src/client/client_cli.c | 31 +++++++++++++++------ src/client/client_send.c | 8 ++++++ src/client/client_send.h | 4 +++ src/shared/protocol.c | 43 ++++++++++++++++++++++-------- tests/fuzz/fuzz_protocol_framing.c | 5 +++- 5 files changed, 71 insertions(+), 20 deletions(-) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 73812ef..dd644c4 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -34,20 +34,35 @@ * client_send.c, still links the symbol. */ volatile sig_atomic_t 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. Only the production client installs it. */ -static void client_signal_handler(int signo) { - (void)signo; - client_abort_requested = 1; +/* Only armed while a network transfer is in flight. Outside that window the + * handler restores the default disposition and re-raises, so purely local modes + * (--list-only/--dry-run/--read-batch/--only-write-batch and the batch-emission + * pass) keep terminating on Ctrl-C/SIGTERM instead of silently swallowing it. */ +volatile sig_atomic_t client_abort_armed = 0; + +void client_set_abort_armed(bool armed) { + client_abort_armed = armed ? 1 : 0; } -#endif bool client_abort_pending(void) { 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 /* 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, diff --git a/src/client/client_send.c b/src/client/client_send.c index 1febec8..db7a560 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -1949,6 +1949,9 @@ int send_files(Config* config) { 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); if (!client) { if (config->transport == TRANSPORT_TCP) @@ -2228,6 +2231,7 @@ send_fail: prepared_scanner_destroy(&prepared); disconnect_transfer_client(client); protocol_session_unbind(); + client_set_abort_armed(false); return ret; } @@ -2257,6 +2261,9 @@ int send_files_multithreaded(Config** config_ptr) { 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 page_size = sysconf(_SC_PAGE_SIZE); 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 directory; report it as errored like rsync does. */ pipeline_context_sender_destroy(context); + client_set_abort_armed(false); return sender_ok && !scan_io ? 0 : 1; } diff --git a/src/client/client_send.h b/src/client/client_send.h index 06d6c6d..2605899 100644 --- a/src/client/client_send.h +++ b/src/client/client_send.h @@ -13,6 +13,10 @@ * receiver can clean up before the client exits. */ extern volatile sig_atomic_t client_abort_requested; 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 * never free it, and the caller retains ownership (freeing it with diff --git a/src/shared/protocol.c b/src/shared/protocol.c index 0a8dac9..57fd2e2 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -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; 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"); return false; @@ -618,6 +624,7 @@ static bool protocol_read_status_until(ProtocolSession* session, Status* status, const struct timespec* deadline) { Status received = STATUS_ERROR; size_t got = 0; + short wait_events = POLLIN; while (got < sizeof(Status)) { if (!session->ssl || SSL_pending(session->ssl) == 0) { 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"); 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); if (poll_result == 0) { log_message(LOG_LEVEL_ERROR, "Receive timeout while reading status"); @@ -647,8 +654,10 @@ static bool protocol_read_status_until(ProtocolSession* session, Status* status, if (bytes_received <= 0) { if (session->ssl) { 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; + } } if (bytes_received < 0 && errno == EINTR) continue; @@ -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 * frame; once part of a frame is buffered a write could race the peer's * 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; struct pollfd pfd = {.fd = session->read_fd, .events = POLLIN}; int poll_result = poll(&pfd, 1, wait_ms); @@ -724,15 +734,26 @@ bool protocol_receive_status_keepalive(ProtocolSession* session, Status* status, * 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 * desynchronize the protocol. */ - while (replies_seen < keepalives_sent) { - Status drained; - if (!protocol_read_status_until(session, &drained, &deadline)) - return false; - if (drained != STATUS_KEEPALIVE) { - log_message(LOG_LEVEL_ERROR, "Unexpected status while draining keepalive replies"); - return false; + 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) { + Status drained; + if (!protocol_read_status_until(session, &drained, &drain_deadline)) { + log_message(LOG_LEVEL_WARNING, "peer did not answer %lu keepalive(s); continuing", + keepalives_sent - replies_seen); + break; + } + if (drained != STATUS_KEEPALIVE) { + log_message(LOG_LEVEL_ERROR, "Unexpected status while draining keepalive replies"); + return false; + } + replies_seen++; } - replies_seen++; } *status = final; log_debug_message(LOG_DEBUG_PROTO, "Received Status: %s", status_to_string(*status)); diff --git a/tests/fuzz/fuzz_protocol_framing.c b/tests/fuzz/fuzz_protocol_framing.c index 6bb0359..a80350a 100644 --- a/tests/fuzz/fuzz_protocol_framing.c +++ b/tests/fuzz/fuzz_protocol_framing.c @@ -81,9 +81,12 @@ int LLVMFuzzerTestOneInput(const uint8_t* data, size_t size) { 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); if (rd >= 0) { - Data* d = receive_data(rd); + Data* d = receive_data_limited(rd, 1u << 20); data_destroy(d); close(rd); close(w);