fix(protocol): atomic client-msg flag, sanitize peer messages, flush before terminal, serialization probe

This commit is contained in:
2026-09-23 23:56:31 +02:00
parent b414f197af
commit 729f3ef8e1
11 changed files with 201 additions and 44 deletions
+9 -2
View File
@@ -495,6 +495,7 @@ int send_dry_run_remote(Config* config) {
protocol_session_bind(&session);
int ret = 1;
bool partial = false;
time_t dry_start = time(NULL);
ReceiverStats dry_stats;
memset(&dry_stats, 0, sizeof(dry_stats));
@@ -653,8 +654,14 @@ int send_dry_run_remote(Config* config) {
if (!receive_status(client->file_descriptor, &status))
goto dry_fail;
}
if (status != STATUS_OK)
/* A per-entry receiver failure is rsync's PARTIAL transfer (exit 23), not a
hard failure: a dry run transfers nothing, but keep the verdict consistent
with the normal path instead of treating it as a protocol error. */
if (status == STATUS_PARTIAL) {
partial = true;
} else if (status != STATUS_OK) {
goto dry_fail;
}
if (!config->quiet) {
if (config->human_readable)
printf("Total: %d files, %s\n", file_count,
@@ -672,7 +679,7 @@ int send_dry_run_remote(Config* config) {
dry_transfer.literal_data = total_bytes;
report_transfer_stats(config, &dry_transfer, dry_start, &dry_stats);
}
ret = io_error ? 1 : 0;
ret = io_error ? 1 : (partial ? 23 : 0);
dry_fail:
if (dry_manifest)
+27 -7
View File
@@ -1213,8 +1213,11 @@ static ArrayList* client_msg_queue = NULL; /* owns char* */
static size_t client_msg_bytes = 0;
/* True only while a live transfer session exists: before the connection is up
(or after it drops) the sink declines so log_message falls back to local
output, matching rsync's documented fallback. */
static bool client_msg_active = false;
output, matching rsync's documented fallback. Written by the sender thread
(client_messages_activate) and read by scanner worker threads in
client_msg_enqueue, so it must be atomic: the queue itself stays guarded by
client_msg_mutex, but the flag is polled before taking that lock. */
static _Atomic bool client_msg_active = false;
static void client_msg_mutex_init(void) {
mtx_init(&client_msg_mutex, mtx_plain);
@@ -1229,20 +1232,26 @@ void client_messages_install(void) {
mtx_lock(&client_msg_mutex);
if (!client_msg_queue)
client_msg_queue = array_list_create(free);
bool ready = client_msg_queue != NULL;
mtx_unlock(&client_msg_mutex);
/* Only arm the sink once the queue exists; on allocation failure leave the
sink uninstalled so log_message keeps writing locally instead of handing
messages to a sink that would silently drop them. */
if (ready)
log_set_client_msg_sink(client_msg_enqueue);
}
void client_messages_activate(bool active) {
client_msg_active = active;
atomic_store(&client_msg_active, active);
}
/* log_message sink: takes ownership (queues) the message when a session is
* live; returns false otherwise so the caller writes it locally. */
static bool client_msg_enqueue(const char* message) {
bool active = atomic_load(&client_msg_active);
if (!message || message[0] == '\0')
return client_msg_active;
if (!client_msg_active)
return active;
if (!active)
return false;
size_t len = strlen(message);
call_once(&client_msg_mutex_once, client_msg_mutex_init);
@@ -1273,8 +1282,19 @@ void client_flush_client_messages(int fd) {
call_once(&client_msg_mutex_once, client_msg_mutex_init);
mtx_lock(&client_msg_mutex);
ArrayList* pending = client_msg_queue;
client_msg_queue = array_list_create(free);
if (pending) {
ArrayList* fresh = array_list_create(free);
if (fresh) {
client_msg_queue = fresh;
} else {
/* No memory for a replacement queue: stop queuing new diagnostics (they
fall back to local output) and drain this batch below so nothing is
silently dropped. */
client_msg_queue = NULL;
log_set_client_msg_sink(NULL);
}
client_msg_bytes = 0;
}
mtx_unlock(&client_msg_mutex);
if (!pending)
return;
@@ -1289,7 +1309,7 @@ void client_flush_client_messages(int fd) {
/* Tear down the sink after a transfer and free anything still queued. */
void client_messages_end(void) {
log_set_client_msg_sink(NULL);
client_msg_active = false;
atomic_store(&client_msg_active, false);
call_once(&client_msg_mutex_once, client_msg_mutex_init);
mtx_lock(&client_msg_mutex);
ArrayList* pending = client_msg_queue;
+17
View File
@@ -255,8 +255,18 @@ static bool finalize_transfer(Client* client, const Config* config, ArrayList* r
*delete_limit_out = false;
if (partial_out)
*partial_out = false;
/* --stderr=client: the receiver consumes frames until it reads
STATUS_FINISHED, after which it no longer reads. Flush every diagnostic
queued during the transfer here -- the last frame boundary at which the
peer is still reading -- so nothing is stranded in the queue. */
client_flush_client_messages(client->file_descriptor);
if (!send_status(client->file_descriptor, STATUS_FINISHED))
return false;
/* Past STATUS_FINISHED the receiver has stopped reading, so any diagnostic
logged from here on (notably the STATUS_PARTIAL warning below) can no
longer be forwarded. Deactivate the channel so those messages fall back
to local output instead of being queued for a closed peer and lost. */
client_messages_activate(false);
/* The receiver emits its optional wire-stats frame (protocol 2.25.0) FIRST,
then any per-file --remove-source-files acks, then the terminal status. */
Status status;
@@ -899,6 +909,13 @@ static int send_chunk_with_removal(Client* client, Chunk* chunk, Config* config,
if (chunk->items[i] == NULL)
continue;
transfer_stats_note_entry(stats, chunk->items[i]);
/* Output parity: probe each entry's ancestor directories' destination
state before emitting its itemize line, exactly as the non-serialized
loop does. Without this, dest_state.known stays false and -i/-P
renders an existing dir/symlink as created instead of `.d..t...` (or
suppressing it). */
if (!client_change_probe_ancestors(config, chunk->items[i], client->file_descriptor))
return -1;
/* The chunk-serialization path emits no --progress name lines, so only
feed -i/--out-format its ancestor directory lines here. */
if (config->itemize_changes || config->out_format != NULL)
+33 -18
View File
@@ -341,7 +341,13 @@ static bool receiver_note_status(const struct timespec* session_start,
}
int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink) {
return receiver_process_pending(config, file_descriptor, sink, NULL, NULL);
return receiver_process_pending_ctx(config, file_descriptor, sink, NULL, NULL, NULL);
}
int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink,
DeleteManifest** pending_manifest, DeletePlanSession** pending_plans) {
return receiver_process_pending_ctx(config, file_descriptor, sink, pending_manifest,
pending_plans, NULL);
}
/* Per-connection state threaded through the status handlers below. The parked
@@ -363,8 +369,10 @@ typedef struct {
*pending_plans so the -m caller commits after its disk writer drained. */
DeletePlanSession* plan_session;
/* Observer context for the per-directory delete session, which outlives the
frame handler; must stay alive until the session commits. */
ReceiverDeleteContext delete_ctx;
frame handler; must stay alive until the session commits. For a session
handed to the caller (pending_plans) this points at a caller-owned
long-lived context; otherwise it points at an internal stack context. */
ReceiverDeleteContext* delete_ctx;
bool early_delete;
bool per_dir_delete;
bool delete_limit_noted;
@@ -386,18 +394,16 @@ static ReceiverStep receiver_handle_keepalive(ReceiverPendingState* state) {
}
/* rsync --stderr=client: a client diagnostic forwarded over the wire. Read the
* bounded string and write it to the server's stderr (respecting the server log
* destination). The body is peer-controlled text, so it is logged verbatim
* (log_client_message adds the standard prefix); trailing newlines are stripped
* so a message cannot inject a blank line. A malformed string (over-long or
* embedded NUL) is a framing error and tears the connection down. */
* bounded string and log it through the normal destination/level gate. The
* body is peer-controlled text: log_client_message() escapes every
* non-printable byte (newlines, CR, ANSI ESC, ...) before writing, so a hostile
* client cannot forge log lines or inject terminal control sequences. A
* malformed string (over-long or embedded NUL) is a framing error and tears the
* connection down. */
static ReceiverStep receiver_handle_client_msg(ReceiverPendingState* state) {
char* message = receive_str(state->fd);
if (!message)
return RECEIVER_STEP_FAIL;
size_t len = strlen(message);
while (len > 0 && (message[len - 1] == '\n' || message[len - 1] == '\r'))
message[--len] = '\0';
if (message[0] != '\0')
log_client_message(message);
free(message);
@@ -652,7 +658,7 @@ static ReceiverStep receiver_handle_delete_plan(ReceiverPendingState* state) {
state->plan_session = delete_plan_session_create(config);
if (state->plan_session && (sink->stats || sink->deleted_paths))
delete_plan_session_set_delete_observer(state->plan_session, receiver_record_deleted_path,
&state->delete_ctx);
state->delete_ctx);
}
if (!state->plan_session || delete_plan_session_receive(state->plan_session, config, fd) != 0)
return RECEIVER_STEP_FAIL;
@@ -735,8 +741,10 @@ static void receiver_drop_pending(ReceiverPendingState* state) {
delete-during plan mode (no manifest at all). See the per-frame handlers
above for how the -m receiver defers that commit until its disk writer has
drained. */
int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink,
DeleteManifest** pending_manifest, DeletePlanSession** pending_plans) {
int receiver_process_pending_ctx(Config* config, int file_descriptor, const ReceiverSink* sink,
DeleteManifest** pending_manifest,
DeletePlanSession** pending_plans,
ReceiverDeleteContext* observer_ctx) {
Status status;
if (!receive_status(file_descriptor, &status))
return -1;
@@ -749,6 +757,13 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
last_progress = session_start;
if (!receiver_note_status(&session_start, &last_progress, status, file_descriptor, sink))
return -1;
/* A per-directory delete session handed to the caller outlives this stack
frame, so its observer context must be caller-owned (observer_ctx); only
the default inline-commit case may use the stack context. */
ReceiverDeleteContext local_ctx;
ReceiverDeleteContext* delete_ctx = observer_ctx ? observer_ctx : &local_ctx;
delete_ctx->stats = sink ? sink->stats : NULL;
delete_ctx->deleted_paths = sink ? sink->deleted_paths : NULL;
ReceiverPendingState state = {
.config = config,
.fd = file_descriptor,
@@ -757,7 +772,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
.pending_plans = pending_plans,
.deferred_manifest = NULL,
.plan_session = NULL,
.delete_ctx = {sink ? sink->stats : NULL, sink ? sink->deleted_paths : NULL},
.delete_ctx = delete_ctx,
.early_delete = config_delete_timing_early(config),
.per_dir_delete = config_delete_timing_per_dir(config),
.delete_limit_noted = false,
@@ -797,9 +812,9 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
state.deferred_manifest = NULL;
} else {
size_t deleted = 0;
DeletePathObserver observer = receiver_delete_observer(sink, &state.delete_ctx);
DeletePathObserver observer = receiver_delete_observer(sink, state.delete_ctx);
DeleteCommitResult deletion = manifest_delete_all_observed(
config, state.deferred_manifest, &deleted, observer, &state.delete_ctx);
config, state.deferred_manifest, &deleted, observer, state.delete_ctx);
receiver_tally_deleted(sink, deleted);
delete_manifest_free(state.deferred_manifest);
state.deferred_manifest = NULL;
@@ -819,7 +834,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
if (state.plan_session) {
if (sink->stats || sink->deleted_paths)
delete_plan_session_set_delete_observer(state.plan_session, receiver_record_deleted_path,
&state.delete_ctx);
state.delete_ctx);
if (state.pending_plans) {
*state.pending_plans = state.plan_session;
state.plan_session = NULL;
+11
View File
@@ -93,6 +93,17 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si
for either to keep the default behaviour (delete before the success frame). */
int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink,
DeleteManifest** pending_manifest, DeletePlanSession** pending_plans);
/* receiver_process_pending() with an explicit observer context for a
per-directory delete session that is handed to the caller via
`pending_plans`. The session outlives this call (the -m pipeline commits it
after joining its disk writer), so its observer context must too: pass a
long-lived object such as PipelineContextReceiver.delete_ctx. When
`delete_ctx` is NULL an internal stack context is used, which is only safe
when the session is committed before returning (the default behaviour). */
int receiver_process_pending_ctx(Config* config, int file_descriptor, const ReceiverSink* sink,
DeleteManifest** pending_manifest,
DeletePlanSession** pending_plans,
ReceiverDeleteContext* delete_ctx);
int receiver_receive_files(Config* config, int file_descriptor);
/* ---- Connection time bounds (anti-slowloris) ----
+5 -2
View File
@@ -28,6 +28,8 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
context->max_queue_bytes = 0;
context->deferred_manifest = NULL;
context->deferred_plans = NULL;
context->delete_ctx.stats = NULL;
context->delete_ctx.deleted_paths = NULL;
context->delete_limit_reached = false;
context->failed_entries = 0;
memset(&context->stats, 0, sizeof(context->stats));
@@ -205,8 +207,9 @@ int receive_thread(void* pipeline_context) {
&context->stats,
context->would_delete,
context->deleted_paths};
if (receiver_process_pending((Config*)config, file_descriptor, &sink, &context->deferred_manifest,
&context->deferred_plans) != 0) {
if (receiver_process_pending_ctx((Config*)config, file_descriptor, &sink,
&context->deferred_manifest, &context->deferred_plans,
&context->delete_ctx) != 0) {
receiver_thread_fail(context);
protocol_session_unbind();
return thrd_error;
+5
View File
@@ -46,6 +46,11 @@ typedef struct PipelineContextReceiver {
committing while the disk writer may still be draining; server.c commits it
after both threads joined. NULL for every other timing. */
DeletePlanSession* deferred_plans;
/* Observer context for `deferred_plans`. It must outlive the receive thread
(the session is committed by server.c after both threads join), so it lives
here rather than on receiver_process_pending()'s stack; receive_thread
installs it on the session. */
ReceiverDeleteContext delete_ctx;
/* Set by server.c when the deferred delete commit hit the --max-delete
budget; the terminal success frame then carries STATUS_DELETE_LIMIT
(rsync exit 25) while the transfer itself still succeeds. */
+3 -5
View File
@@ -1011,11 +1011,9 @@ static void server_run_mt_receiver(ServerSession* state) {
--delete-during already applied its plans on the receive thread. */
if (context->deferred_plans) {
/* Defence in depth (the enclosing block already excludes dry-run): a
-n run never commits a deletion. */
ReceiverDeleteContext delctx = {&context->stats, context->deleted_paths};
if (delctx.stats || delctx.deleted_paths)
delete_plan_session_set_delete_observer(context->deferred_plans,
receiver_record_deleted_path, &delctx);
-n run never commits a deletion. The session's observer context was
installed by receive_thread from context->delete_ctx, which outlives
both threads, so no stack context is needed here. */
DeleteCommitResult deletion =
config->dry_run ? DELETE_COMMIT_OK
: delete_plan_session_commit(context->deferred_plans, config);
+15 -8
View File
@@ -1,4 +1,5 @@
#include "log.h"
#include "utils.h"
#include <errno.h>
#include <stdbool.h>
#include <stdarg.h>
@@ -155,15 +156,21 @@ static void emit_log_line(FILE* console, const char* line) {
void log_client_message(const char* message) {
if (!message)
return;
time_t now = time(NULL);
struct tm t;
if (!localtime_r(&now, &t))
/* The body is peer-controlled: escape every non-printable byte (newlines,
CR, ANSI ESC, ...) so a hostile client cannot forge log lines or inject
terminal control sequences. output_escape() is the codebase's canonical
escaper and leaves printable text untouched. */
char* escaped = output_escape(message, log_get_8_bit_output());
if (!escaped)
return;
char* line = format_log_line_from_body(LOG_LEVEL_INFO, &t, message);
if (!line)
return;
emit_log_line(stderr, line);
free(line);
/* Route through the ordinary log level / destination gate (log_message):
this respects --log-file, the configured stderr mode and the level
threshold instead of always writing to stderr. The wire body carries no
severity, so the forwarded diagnostic is emitted as a warning -- the
lowest level the default gate admits, which keeps the peer's messages
visible without bypassing --quiet. */
log_message(LOG_LEVEL_WARNING, "%s", escaped);
free(escaped);
}
void log_message(LogLevel log_level, const char* format, ...) {
+31
View File
@@ -2893,6 +2893,37 @@ class TestItemizeChanges:
result, _ = run_client(SOURCE_DIR, DEST_DIR, flags=["-i", "--dry-run"])
assert result.returncode == 0, f"dry-run -i failed: {result.stderr[:200]}"
def test_chunk_serialization_probes_ancestor_dir_state(self, shared_server):
"""#314: --chunk-serialization + -i must probe ancestor directory state
so a pre-existing directory with a changed mtime itemizes as an
attribute change (`.d..t......`) instead of being rendered as created
(`cd+++++++++`)."""
source = os.path.join(TEST_DATA_DIR, "itemize_chunk_serial_src")
dest = os.path.join(TEST_DATA_DIR, "itemize_chunk_serial_dst")
clean_dir(source)
clean_dir(dest)
subdir = os.path.join(source, "sub")
os.makedirs(subdir)
with open(os.path.join(subdir, "file.txt"), "wb") as fh:
fh.write(b"payload\n")
result, _ = run_client(source, dest, flags=["--preserve"], port=shared_server.port)
assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}"
# Change only the source directory's mtime; its contents stay identical
# so only the directory's time attribute differs on the rerun.
os.utime(subdir, (1_000_000_000, 1_000_000_000))
result, _ = run_client(source, dest,
flags=["--preserve", "-i", "--chunk-serialization"],
port=shared_server.port)
assert result.returncode == 0, f"chunk-serialization -i failed: {result.stderr[:200]}"
dir_lines = [line for line in result.stdout.splitlines() if line.endswith(" sub/")]
assert dir_lines == [".d..t...... sub/"], (
f"expected an attribute-change dir line, got {dir_lines!r}; "
f"full stdout={result.stdout!r}"
)
def test_changed_file_on_second_incremental_run_prints_exactly_one_line(self, shared_server):
"""A changed file itemizes exactly once on an incremental rerun while
unchanged files print nothing (no double emission)."""
+43
View File
@@ -309,6 +309,48 @@ static void test_log_set_file_null_before_fclose(void) {
EXPECT_TRUE(true);
}
/* The server's client-message channel logs a peer-controlled body: it must be
* escaped so an interior newline/CR/ANSI escape cannot forge a log line or
* move the terminal cursor. With the default stderr mode a forwarded message
* is emitted as a warning on stdout. */
static void test_log_client_message_sanitized(void) {
int pipe_fds[2];
EXPECT_EQ_INT(pipe(pipe_fds), 0);
/* Drain anything earlier tests buffered on stdout before redirecting, so the
capture holds only this test's single log line. */
fflush(stdout);
int saved_stdout = dup(STDOUT_FILENO);
EXPECT_TRUE(saved_stdout >= 0);
EXPECT_TRUE(dup2(pipe_fds[1], STDOUT_FILENO) >= 0);
close(pipe_fds[1]);
set_log_level(LOG_LEVEL_WARNING);
log_set_stderr_mode(LOG_STDERR_ERRORS);
log_client_message("forged\n2026-01-01 [ERROR]: fake\x1b[31mred");
fflush(stdout);
EXPECT_TRUE(dup2(saved_stdout, STDOUT_FILENO) >= 0);
close(saved_stdout);
char output[512] = {0};
ssize_t length = read(pipe_fds[0], output, sizeof(output) - 1);
close(pipe_fds[0]);
EXPECT_TRUE(length > 0);
/* The interior newline became an escaped octal, so the body stays on one
physical line: exactly one '\n' (the log terminator) is present. */
int newlines = 0;
for (ssize_t i = 0; i < length; i++) {
if (output[i] == '\n')
newlines++;
}
EXPECT_EQ_INT(newlines, 1);
/* The ESC introducer is escaped too, so no raw ANSI sequence reaches the
terminal. */
EXPECT_TRUE(strchr(output, '\x1b') == NULL);
EXPECT_TRUE(strstr(output, "forged") != NULL);
EXPECT_TRUE(strstr(output, "fake") != NULL);
log_set_stderr_mode(LOG_STDERR_ERRORS);
}
void test_log() {
test_log_message_debug();
test_log_message_info();
@@ -320,6 +362,7 @@ void test_log() {
test_log_filtering();
test_log_stderr_mode_all();
test_log_stderr_mode_client();
test_log_client_message_sanitized();
test_log_message_formats();
test_log_debug_enabled_matches_gate();
test_log_concurrent_no_torn_lines();