diff --git a/src/client/client_manifest.c b/src/client/client_manifest.c index 1d4215f..04fae0b 100644 --- a/src/client/client_manifest.c +++ b/src/client/client_manifest.c @@ -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) diff --git a/src/client/client_report.c b/src/client/client_report.c index 9c1e4c7..837f9b9 100644 --- a/src/client/client_report.c +++ b/src/client/client_report.c @@ -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); - log_set_client_msg_sink(client_msg_enqueue); + /* 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); - client_msg_bytes = 0; + 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; diff --git a/src/client/client_send.c b/src/client/client_send.c index 5337ccb..354b86f 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -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) diff --git a/src/server/receiver.c b/src/server/receiver.c index a815391..25a183b 100644 --- a/src/server/receiver.c +++ b/src/server/receiver.c @@ -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; diff --git a/src/server/receiver.h b/src/server/receiver.h index 2f68527..1c6b341 100644 --- a/src/server/receiver.h +++ b/src/server/receiver.h @@ -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) ---- diff --git a/src/server/receiver_pipeline.c b/src/server/receiver_pipeline.c index 6feab4e..b28ab1b 100644 --- a/src/server/receiver_pipeline.c +++ b/src/server/receiver_pipeline.c @@ -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; diff --git a/src/server/receiver_pipeline.h b/src/server/receiver_pipeline.h index 33e9f73..84d4142 100644 --- a/src/server/receiver_pipeline.h +++ b/src/server/receiver_pipeline.h @@ -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. */ diff --git a/src/server/server.c b/src/server/server.c index 29faaea..f374ecc 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -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); diff --git a/src/shared/log.c b/src/shared/log.c index 1467828..f3db46b 100644 --- a/src/shared/log.c +++ b/src/shared/log.c @@ -1,4 +1,5 @@ #include "log.h" +#include "utils.h" #include #include #include @@ -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, ...) { diff --git a/tests/integration/test_features.py b/tests/integration/test_features.py index cf28bb5..a93cfac 100644 --- a/tests/integration/test_features.py +++ b/tests/integration/test_features.py @@ -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).""" diff --git a/tests/test_log.c b/tests/test_log.c index 75829c6..59e72f0 100644 --- a/tests/test_log.c +++ b/tests/test_log.c @@ -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();