feat(stats): populate receiver wire counters on both receive paths
The single-threaded and -m receivers never populated ReceiverStats.matched_data or .deleted_files, so --stats always printed 0 for both even when rsync reported nonzero. Track the bytes reconstructed from the basis file while applying a delta, and tally the delete-commit counts (manifest and per-directory sessions) into the receiver stats. The -m pipeline now carries its own stats/would-delete fields and emits the STATUS_STATS frame before the terminal success, so --threads finally reports the counters and renders -n --delete lines. Also normalize the -n --delete would-delete enumeration's absolute basis prefixes exactly like the real commit path (fixing an over-report) and fix the basis_delete_relative off-by-one when the receive root is '/'. Unit tests cover the root mapping and the basis protection; integration tests cover matched/deleted stats for both receivers and the --threads dry-run delete lines.
This commit is contained in:
+22
-4
@@ -83,6 +83,13 @@ bool receiver_send_stats_frame(int fd, const Config* config, const ReceiverStats
|
||||
return true;
|
||||
}
|
||||
|
||||
/* Add a delete commit's tally to the sink's end-of-transfer wire counters (when
|
||||
the sink reports them). Runs on the receiving thread, so no locking. */
|
||||
static void receiver_tally_deleted(const ReceiverSink* sink, size_t deleted) {
|
||||
if (sink && sink->stats && deleted > 0)
|
||||
sink->stats->deleted_files += deleted;
|
||||
}
|
||||
|
||||
static bool receiver_process_chunk(Chunk* chunk, const ReceiverSink* sink) {
|
||||
if (!chunk || !sink || !sink->store_file)
|
||||
return false;
|
||||
@@ -390,9 +397,12 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
||||
A later transfer failure does not restore these deletions. A
|
||||
--max-delete-capped commit still succeeds and the transfer proceeds;
|
||||
the terminal success frame reports the cap. */
|
||||
DeleteCommitResult deletion = (config->use_delete || config->delete_missing_args)
|
||||
? manifest_delete_all(config, manifest)
|
||||
: DELETE_COMMIT_OK;
|
||||
size_t deleted = 0;
|
||||
DeleteCommitResult deletion =
|
||||
(config->use_delete || config->delete_missing_args)
|
||||
? manifest_delete_all_counted(config, manifest, &deleted)
|
||||
: DELETE_COMMIT_OK;
|
||||
receiver_tally_deleted(sink, deleted);
|
||||
delete_manifest_free(manifest);
|
||||
if (deletion == DELETE_COMMIT_ERROR) {
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
@@ -469,7 +479,10 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
||||
*pending_manifest = deferred_manifest;
|
||||
deferred_manifest = NULL;
|
||||
} else {
|
||||
DeleteCommitResult deletion = manifest_delete_all(config, deferred_manifest);
|
||||
size_t deleted = 0;
|
||||
DeleteCommitResult deletion =
|
||||
manifest_delete_all_counted(config, deferred_manifest, &deleted);
|
||||
receiver_tally_deleted(sink, deleted);
|
||||
delete_manifest_free(deferred_manifest);
|
||||
deferred_manifest = NULL;
|
||||
if (deletion == DELETE_COMMIT_ERROR) {
|
||||
@@ -496,6 +509,7 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver
|
||||
} else {
|
||||
DeleteCommitResult deletion = delete_plan_session_commit(plan_session, config);
|
||||
bool limit = delete_plan_session_limit_reached(plan_session);
|
||||
receiver_tally_deleted(sink, delete_plan_session_deleted(plan_session));
|
||||
delete_plan_session_destroy(plan_session);
|
||||
plan_session = NULL;
|
||||
if (deletion == DELETE_COMMIT_ERROR) {
|
||||
@@ -572,6 +586,10 @@ static bool receiver_save_file(File* file, void* context_pointer) {
|
||||
} else {
|
||||
result = file_save_to_disk_full(context->config->receive_root_directory, file, context->config);
|
||||
}
|
||||
/* Wire-stats tally: bytes reconstructed from the basis file (delta matches)
|
||||
count as matched data in the end-of-transfer report. */
|
||||
if (result != FILE_SAVE_ERROR && file->matched_bytes > 0)
|
||||
context->stats.matched_data += file->matched_bytes;
|
||||
/* A directory's metadata is deferred, never applied inline: collect it now
|
||||
and apply it at the end. -O/--omit-dir-times and --preserve_perms/-times
|
||||
are honored by dir_metadata_list_apply's caller (see
|
||||
|
||||
@@ -29,6 +29,8 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
|
||||
context->deferred_manifest = NULL;
|
||||
context->deferred_plans = NULL;
|
||||
context->delete_limit_reached = false;
|
||||
memset(&context->stats, 0, sizeof(context->stats));
|
||||
context->would_delete = NULL;
|
||||
atomic_init(&context->cancelled, false);
|
||||
int init = 0;
|
||||
if (mtx_init(&context->mutex, mtx_plain) != thrd_success)
|
||||
@@ -41,6 +43,9 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
|
||||
goto fail;
|
||||
// cppcheck-suppress unreadVariable
|
||||
init++;
|
||||
context->would_delete = array_list_create(free);
|
||||
if (!context->would_delete)
|
||||
goto fail;
|
||||
return context;
|
||||
|
||||
fail:
|
||||
@@ -64,6 +69,8 @@ void pipeline_context_receiver_destroy(PipelineContextReceiver* context) {
|
||||
queue_destroy(context->queue);
|
||||
receiver_outcomes_destroy(&context->outcomes);
|
||||
dir_time_list_free(&context->dir_times);
|
||||
if (context->would_delete)
|
||||
array_list_delete(context->would_delete);
|
||||
mtx_destroy(&context->mutex);
|
||||
cnd_destroy(&context->condition_not_full);
|
||||
cnd_destroy(&context->condition_not_empty);
|
||||
@@ -136,6 +143,11 @@ bool pipeline_context_receiver_enqueue_file(PipelineContextReceiver* context, Fi
|
||||
|
||||
static bool receiver_enqueue_file(File* file, void* context_pointer) {
|
||||
PipelineContextReceiver* context = (PipelineContextReceiver*)context_pointer;
|
||||
if (file && file->matched_bytes > 0) {
|
||||
mtx_lock(&context->mutex);
|
||||
context->stats.matched_data += file->matched_bytes;
|
||||
mtx_unlock(&context->mutex);
|
||||
}
|
||||
return pipeline_context_receiver_enqueue_file(context, file);
|
||||
}
|
||||
|
||||
@@ -171,8 +183,8 @@ int receive_thread(void* pipeline_context) {
|
||||
false,
|
||||
NULL,
|
||||
receiver_pipeline_note_delete_limit,
|
||||
NULL,
|
||||
NULL};
|
||||
&context->stats,
|
||||
context->would_delete};
|
||||
if (receiver_process_pending((Config*)config, file_descriptor, &sink, &context->deferred_manifest,
|
||||
&context->deferred_plans) != 0) {
|
||||
receiver_thread_fail(context);
|
||||
|
||||
@@ -54,6 +54,13 @@ typedef struct PipelineContextReceiver {
|
||||
directory entries. Only write_thread mutates it (before it joins); the
|
||||
caller (server.c) applies it after the delete/delay-updates phase. */
|
||||
DirTimeList dir_times;
|
||||
/* End-of-transfer wire counters (protocol 2.25.0). receive_thread accumulates
|
||||
matched_data under `mutex`; server.c adds the delete-commit tallies after
|
||||
both threads join and emits the STATUS_STATS frame. */
|
||||
ReceiverStats stats;
|
||||
/* -n/--dry-run --delete would-delete path list, collected by receive_thread
|
||||
and reported in the STATUS_STATS frame. */
|
||||
struct ArrayList* would_delete;
|
||||
} PipelineContextReceiver;
|
||||
|
||||
PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue_receiver,
|
||||
|
||||
+10
-2
@@ -955,7 +955,10 @@ void handler(int file_descriptor) {
|
||||
--delay-updates run; the walker skips the staging directory. A
|
||||
server-contacting --dry-run deletes nothing (no manifest is sent). */
|
||||
if (context->deferred_manifest) {
|
||||
DeleteCommitResult deletion = manifest_delete_all(config, context->deferred_manifest);
|
||||
size_t deleted = 0;
|
||||
DeleteCommitResult deletion =
|
||||
manifest_delete_all_counted(config, context->deferred_manifest, &deleted);
|
||||
context->stats.deleted_files += deleted;
|
||||
if (deletion == DELETE_COMMIT_ERROR) {
|
||||
transfer_ok = false;
|
||||
} else if (deletion == DELETE_COMMIT_LIMIT_REACHED) {
|
||||
@@ -975,6 +978,7 @@ void handler(int file_descriptor) {
|
||||
DeleteCommitResult deletion = config->dry_run ? DELETE_COMMIT_OK
|
||||
: delete_plan_session_commit(
|
||||
context->deferred_plans, config);
|
||||
context->stats.deleted_files += delete_plan_session_deleted(context->deferred_plans);
|
||||
if (deletion == DELETE_COMMIT_ERROR) {
|
||||
transfer_ok = false;
|
||||
} else if (deletion == DELETE_COMMIT_LIMIT_REACHED) {
|
||||
@@ -1003,7 +1007,11 @@ void handler(int file_descriptor) {
|
||||
}
|
||||
if (transfer_ok) {
|
||||
Status final_status = context->delete_limit_reached ? STATUS_DELETE_LIMIT : STATUS_OK;
|
||||
if (!receiver_send_final_success(file_descriptor, config, &context->outcomes, final_status))
|
||||
/* Emit the optional wire-stats record first (protocol 2.25.0), then the
|
||||
success/outcome frame, exactly like the single-threaded receiver. */
|
||||
if (!receiver_send_stats_frame(file_descriptor, config, &context->stats,
|
||||
context->would_delete) ||
|
||||
!receiver_send_final_success(file_descriptor, config, &context->outcomes, final_status))
|
||||
transfer_ok = false;
|
||||
} else {
|
||||
send_error_detail(file_descriptor, "transfer failed on receiver");
|
||||
|
||||
Reference in New Issue
Block a user