feat: wire backlog — client-msg, partial exit 23, dir itemize, stats, per-dir merge (protocol 2.30.0) #327

Merged
TapTap merged 7 commits from wire/backlog-230 into dev 2026-09-24 00:37:48 +02:00
11 changed files with 201 additions and 44 deletions
Showing only changes of commit 6d9cdb83ba - Show all commits
+9 -2
View File
@@ -495,6 +495,7 @@ int send_dry_run_remote(Config* config) {
protocol_session_bind(&session); protocol_session_bind(&session);
int ret = 1; int ret = 1;
bool partial = false;
time_t dry_start = time(NULL); time_t dry_start = time(NULL);
ReceiverStats dry_stats; ReceiverStats dry_stats;
memset(&dry_stats, 0, sizeof(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)) if (!receive_status(client->file_descriptor, &status))
goto dry_fail; 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; goto dry_fail;
}
if (!config->quiet) { if (!config->quiet) {
if (config->human_readable) if (config->human_readable)
printf("Total: %d files, %s\n", file_count, 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; dry_transfer.literal_data = total_bytes;
report_transfer_stats(config, &dry_transfer, dry_start, &dry_stats); report_transfer_stats(config, &dry_transfer, dry_start, &dry_stats);
} }
ret = io_error ? 1 : 0; ret = io_error ? 1 : (partial ? 23 : 0);
dry_fail: dry_fail:
if (dry_manifest) if (dry_manifest)
+29 -9
View File
@@ -1213,8 +1213,11 @@ static ArrayList* client_msg_queue = NULL; /* owns char* */
static size_t client_msg_bytes = 0; static size_t client_msg_bytes = 0;
/* True only while a live transfer session exists: before the connection is up /* 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 (or after it drops) the sink declines so log_message falls back to local
output, matching rsync's documented fallback. */ output, matching rsync's documented fallback. Written by the sender thread
static bool client_msg_active = false; (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) { static void client_msg_mutex_init(void) {
mtx_init(&client_msg_mutex, mtx_plain); mtx_init(&client_msg_mutex, mtx_plain);
@@ -1229,20 +1232,26 @@ void client_messages_install(void) {
mtx_lock(&client_msg_mutex); mtx_lock(&client_msg_mutex);
if (!client_msg_queue) if (!client_msg_queue)
client_msg_queue = array_list_create(free); client_msg_queue = array_list_create(free);
bool ready = client_msg_queue != NULL;
mtx_unlock(&client_msg_mutex); 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) { 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 /* log_message sink: takes ownership (queues) the message when a session is
* live; returns false otherwise so the caller writes it locally. */ * live; returns false otherwise so the caller writes it locally. */
static bool client_msg_enqueue(const char* message) { static bool client_msg_enqueue(const char* message) {
bool active = atomic_load(&client_msg_active);
if (!message || message[0] == '\0') if (!message || message[0] == '\0')
return client_msg_active; return active;
if (!client_msg_active) if (!active)
return false; return false;
size_t len = strlen(message); size_t len = strlen(message);
call_once(&client_msg_mutex_once, client_msg_mutex_init); 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); call_once(&client_msg_mutex_once, client_msg_mutex_init);
mtx_lock(&client_msg_mutex); mtx_lock(&client_msg_mutex);
ArrayList* pending = client_msg_queue; ArrayList* pending = client_msg_queue;
client_msg_queue = array_list_create(free); if (pending) {
client_msg_bytes = 0; 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); mtx_unlock(&client_msg_mutex);
if (!pending) if (!pending)
return; return;
@@ -1289,7 +1309,7 @@ void client_flush_client_messages(int fd) {
/* Tear down the sink after a transfer and free anything still queued. */ /* Tear down the sink after a transfer and free anything still queued. */
void client_messages_end(void) { void client_messages_end(void) {
log_set_client_msg_sink(NULL); 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); call_once(&client_msg_mutex_once, client_msg_mutex_init);
mtx_lock(&client_msg_mutex); mtx_lock(&client_msg_mutex);
ArrayList* pending = client_msg_queue; 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; *delete_limit_out = false;
if (partial_out) if (partial_out)
*partial_out = false; *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)) if (!send_status(client->file_descriptor, STATUS_FINISHED))
return false; 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, /* 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. */ then any per-file --remove-source-files acks, then the terminal status. */
Status status; Status status;
@@ -899,6 +909,13 @@ static int send_chunk_with_removal(Client* client, Chunk* chunk, Config* config,
if (chunk->items[i] == NULL) if (chunk->items[i] == NULL)
continue; continue;
transfer_stats_note_entry(stats, chunk->items[i]); 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 /* The chunk-serialization path emits no --progress name lines, so only
feed -i/--out-format its ancestor directory lines here. */ feed -i/--out-format its ancestor directory lines here. */
if (config->itemize_changes || config->out_format != NULL) 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) { 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 /* 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. */ *pending_plans so the -m caller commits after its disk writer drained. */
DeletePlanSession* plan_session; DeletePlanSession* plan_session;
/* Observer context for the per-directory delete session, which outlives the /* Observer context for the per-directory delete session, which outlives the
frame handler; must stay alive until the session commits. */ frame handler; must stay alive until the session commits. For a session
ReceiverDeleteContext delete_ctx; 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 early_delete;
bool per_dir_delete; bool per_dir_delete;
bool delete_limit_noted; 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 /* 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 * bounded string and log it through the normal destination/level gate. The
* destination). The body is peer-controlled text, so it is logged verbatim * body is peer-controlled text: log_client_message() escapes every
* (log_client_message adds the standard prefix); trailing newlines are stripped * non-printable byte (newlines, CR, ANSI ESC, ...) before writing, so a hostile
* so a message cannot inject a blank line. A malformed string (over-long or * client cannot forge log lines or inject terminal control sequences. A
* embedded NUL) is a framing error and tears the connection down. */ * malformed string (over-long or embedded NUL) is a framing error and tears the
* connection down. */
static ReceiverStep receiver_handle_client_msg(ReceiverPendingState* state) { static ReceiverStep receiver_handle_client_msg(ReceiverPendingState* state) {
char* message = receive_str(state->fd); char* message = receive_str(state->fd);
if (!message) if (!message)
return RECEIVER_STEP_FAIL; 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') if (message[0] != '\0')
log_client_message(message); log_client_message(message);
free(message); free(message);
@@ -652,7 +658,7 @@ static ReceiverStep receiver_handle_delete_plan(ReceiverPendingState* state) {
state->plan_session = delete_plan_session_create(config); state->plan_session = delete_plan_session_create(config);
if (state->plan_session && (sink->stats || sink->deleted_paths)) if (state->plan_session && (sink->stats || sink->deleted_paths))
delete_plan_session_set_delete_observer(state->plan_session, receiver_record_deleted_path, 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) if (!state->plan_session || delete_plan_session_receive(state->plan_session, config, fd) != 0)
return RECEIVER_STEP_FAIL; 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 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 above for how the -m receiver defers that commit until its disk writer has
drained. */ drained. */
int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink, int receiver_process_pending_ctx(Config* config, int file_descriptor, const ReceiverSink* sink,
DeleteManifest** pending_manifest, DeletePlanSession** pending_plans) { DeleteManifest** pending_manifest,
DeletePlanSession** pending_plans,
ReceiverDeleteContext* observer_ctx) {
Status status; Status status;
if (!receive_status(file_descriptor, &status)) if (!receive_status(file_descriptor, &status))
return -1; return -1;
@@ -749,6 +757,13 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
last_progress = session_start; last_progress = session_start;
if (!receiver_note_status(&session_start, &last_progress, status, file_descriptor, sink)) if (!receiver_note_status(&session_start, &last_progress, status, file_descriptor, sink))
return -1; 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 = { ReceiverPendingState state = {
.config = config, .config = config,
.fd = file_descriptor, .fd = file_descriptor,
@@ -757,7 +772,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
.pending_plans = pending_plans, .pending_plans = pending_plans,
.deferred_manifest = NULL, .deferred_manifest = NULL,
.plan_session = 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), .early_delete = config_delete_timing_early(config),
.per_dir_delete = config_delete_timing_per_dir(config), .per_dir_delete = config_delete_timing_per_dir(config),
.delete_limit_noted = false, .delete_limit_noted = false,
@@ -797,9 +812,9 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
state.deferred_manifest = NULL; state.deferred_manifest = NULL;
} else { } else {
size_t deleted = 0; 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( 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); receiver_tally_deleted(sink, deleted);
delete_manifest_free(state.deferred_manifest); delete_manifest_free(state.deferred_manifest);
state.deferred_manifest = NULL; state.deferred_manifest = NULL;
@@ -819,7 +834,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
if (state.plan_session) { if (state.plan_session) {
if (sink->stats || sink->deleted_paths) if (sink->stats || sink->deleted_paths)
delete_plan_session_set_delete_observer(state.plan_session, receiver_record_deleted_path, delete_plan_session_set_delete_observer(state.plan_session, receiver_record_deleted_path,
&state.delete_ctx); state.delete_ctx);
if (state.pending_plans) { if (state.pending_plans) {
*state.pending_plans = state.plan_session; *state.pending_plans = state.plan_session;
state.plan_session = NULL; 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). */ for either to keep the default behaviour (delete before the success frame). */
int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink, int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink,
DeleteManifest** pending_manifest, DeletePlanSession** pending_plans); 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); int receiver_receive_files(Config* config, int file_descriptor);
/* ---- Connection time bounds (anti-slowloris) ---- /* ---- 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->max_queue_bytes = 0;
context->deferred_manifest = NULL; context->deferred_manifest = NULL;
context->deferred_plans = NULL; context->deferred_plans = NULL;
context->delete_ctx.stats = NULL;
context->delete_ctx.deleted_paths = NULL;
context->delete_limit_reached = false; context->delete_limit_reached = false;
context->failed_entries = 0; context->failed_entries = 0;
memset(&context->stats, 0, sizeof(context->stats)); memset(&context->stats, 0, sizeof(context->stats));
@@ -205,8 +207,9 @@ int receive_thread(void* pipeline_context) {
&context->stats, &context->stats,
context->would_delete, context->would_delete,
context->deleted_paths}; context->deleted_paths};
if (receiver_process_pending((Config*)config, file_descriptor, &sink, &context->deferred_manifest, if (receiver_process_pending_ctx((Config*)config, file_descriptor, &sink,
&context->deferred_plans) != 0) { &context->deferred_manifest, &context->deferred_plans,
&context->delete_ctx) != 0) {
receiver_thread_fail(context); receiver_thread_fail(context);
protocol_session_unbind(); protocol_session_unbind();
return thrd_error; 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 committing while the disk writer may still be draining; server.c commits it
after both threads joined. NULL for every other timing. */ after both threads joined. NULL for every other timing. */
DeletePlanSession* deferred_plans; 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 /* Set by server.c when the deferred delete commit hit the --max-delete
budget; the terminal success frame then carries STATUS_DELETE_LIMIT budget; the terminal success frame then carries STATUS_DELETE_LIMIT
(rsync exit 25) while the transfer itself still succeeds. */ (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. */ --delete-during already applied its plans on the receive thread. */
if (context->deferred_plans) { if (context->deferred_plans) {
/* Defence in depth (the enclosing block already excludes dry-run): a /* Defence in depth (the enclosing block already excludes dry-run): a
-n run never commits a deletion. */ -n run never commits a deletion. The session's observer context was
ReceiverDeleteContext delctx = {&context->stats, context->deleted_paths}; installed by receive_thread from context->delete_ctx, which outlives
if (delctx.stats || delctx.deleted_paths) both threads, so no stack context is needed here. */
delete_plan_session_set_delete_observer(context->deferred_plans,
receiver_record_deleted_path, &delctx);
DeleteCommitResult deletion = DeleteCommitResult deletion =
config->dry_run ? DELETE_COMMIT_OK config->dry_run ? DELETE_COMMIT_OK
: delete_plan_session_commit(context->deferred_plans, config); : delete_plan_session_commit(context->deferred_plans, config);
+15 -8
View File
@@ -1,4 +1,5 @@
#include "log.h" #include "log.h"
#include "utils.h"
#include <errno.h> #include <errno.h>
#include <stdbool.h> #include <stdbool.h>
#include <stdarg.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) { void log_client_message(const char* message) {
if (!message) if (!message)
return; return;
time_t now = time(NULL); /* The body is peer-controlled: escape every non-printable byte (newlines,
struct tm t; CR, ANSI ESC, ...) so a hostile client cannot forge log lines or inject
if (!localtime_r(&now, &t)) 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; return;
char* line = format_log_line_from_body(LOG_LEVEL_INFO, &t, message); /* Route through the ordinary log level / destination gate (log_message):
if (!line) this respects --log-file, the configured stderr mode and the level
return; threshold instead of always writing to stderr. The wire body carries no
emit_log_line(stderr, line); severity, so the forwarded diagnostic is emitted as a warning -- the
free(line); 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, ...) { 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"]) result, _ = run_client(SOURCE_DIR, DEST_DIR, flags=["-i", "--dry-run"])
assert result.returncode == 0, f"dry-run -i failed: {result.stderr[:200]}" 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): 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 """A changed file itemizes exactly once on an incremental rerun while
unchanged files print nothing (no double emission).""" 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); 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() { void test_log() {
test_log_message_debug(); test_log_message_debug();
test_log_message_info(); test_log_message_info();
@@ -320,6 +362,7 @@ void test_log() {
test_log_filtering(); test_log_filtering();
test_log_stderr_mode_all(); test_log_stderr_mode_all();
test_log_stderr_mode_client(); test_log_stderr_mode_client();
test_log_client_message_sanitized();
test_log_message_formats(); test_log_message_formats();
test_log_debug_enabled_matches_gate(); test_log_debug_enabled_matches_gate();
test_log_concurrent_no_torn_lines(); test_log_concurrent_no_torn_lines();