diff --git a/CMakeLists.txt b/CMakeLists.txt index 192afdc..e985d79 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -1,6 +1,6 @@ cmake_minimum_required(VERSION 3.22) -project(FastFileTransfer VERSION 2.26.0) +project(FastFileTransfer VERSION 2.27.0) set(CMAKE_EXPORT_COMPILE_COMMANDS ON) set(CMAKE_C_STANDARD 11) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index a0b93ad..53bbc31 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -2648,9 +2648,16 @@ static int cli_finalize_config(Config* config, bool verbose, bool no_delta, bool } } } + /* --info=del on a real --delete run asks the receiver to report the paths it + actually removed; the report rides the STATUS_STATS path list, so the wire + stats frame must be negotiated too. */ + config->report_deletes = + config->use_delete && !config->dry_run && + ((config->info_level & LOG_INFO_DEL) != 0 || config->itemize_changes || + config->out_format != NULL); config->report_stats = config->stats || config->show_progress || (config->info_level & LOG_INFO_PROGRESS) || format_needs_wire || - (config->dry_run && config->use_delete); + config->report_deletes || (config->dry_run && config->use_delete); return 0; } diff --git a/src/client/client_send.c b/src/client/client_send.c index 253a2c2..22d4aad 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -307,6 +307,45 @@ static bool info_flag_enabled(const Config* config, LogInfoFlag flag) { return config != NULL && (config->info_level & flag) != 0; } +/* Print rsync's deletion lines for a received list of destination-relative + * paths: `*deleting PATH` when itemizing, the --out-format expansion when a + * format is set, else `deleting PATH` for --info=del. Used by both the dry-run + * would-delete report and the real --info=del report. */ +static void print_delete_reports(const Config* config, const ArrayList* paths) { + if (!config || !paths || config->quiet) + return; + if (!(config->itemize_changes || config->out_format != NULL || + info_flag_enabled(config, LOG_INFO_DEL))) + return; + for (int i = 0; i < paths->size; i++) { + const char* raw = (const char*)paths->items[i]; + const char* path = delete_display_path(config, raw); + if (config->out_format != NULL) { + ChangeEvent event; + memset(&event, 0, sizeof(event)); + event.decision = CHANGE_SENT; + event.deleted = true; + event.name = path; + event.path = path; + char* line = change_render_format(config->out_format, config, &event); + if (line) { + char* escaped = output_escape(line, config->eight_bit_output); + printf("%s\n", escaped ? escaped : line); + free(escaped); + free(line); + } + } else { + char* escaped = output_escape(path, config->eight_bit_output); + if (config->itemize_changes) + printf("*deleting %s\n", escaped ? escaped : path); + else + printf("deleting %s\n", escaped ? escaped : path); + free(escaped); + } + } + fflush(stdout); +} + static void client_progress_begin(const Config* config) { g_progress_active = (config->show_progress || info_flag_enabled(config, LOG_INFO_PROGRESS)) && !config->quiet; @@ -1129,8 +1168,19 @@ static bool finalize_transfer(Client* client, const Config* config, ArrayList* r return false; if (status == STATUS_STATS) { ReceiverStats scratch; - if (!receive_stats_record(client->file_descriptor, stats_out ? stats_out : &scratch, NULL)) + /* A real --info=del run carries the actually-removed paths in the stats + frame's path list; collect and print them in rsync's format. */ + ArrayList* deleted = config->report_deletes ? array_list_create(free) : NULL; + if (config->report_deletes && !deleted) return false; + if (!receive_stats_record(client->file_descriptor, stats_out ? stats_out : &scratch, deleted)) { + array_list_delete(deleted); + return false; + } + if (deleted) { + print_delete_reports(config, deleted); + array_list_delete(deleted); + } if (!receive_status(client->file_descriptor, &status)) return false; } @@ -2081,40 +2131,7 @@ static int send_dry_run_remote(Config* config) { array_list_delete(would_delete); goto dry_fail; } - /* rsync prints `*deleting PATH` when itemizing, `deleting PATH` under - --info=del/--info=remove, and the --out-format expansion when set. */ - if (!config->quiet && (config->itemize_changes || config->out_format != NULL || - info_flag_enabled(config, LOG_INFO_DEL))) { - for (int i = 0; i < would_delete->size; i++) { - const char* raw = (const char*)would_delete->items[i]; - const char* path = delete_display_path(config, raw); - if (config->out_format != NULL) { - ChangeEvent event; - memset(&event, 0, sizeof(event)); - event.decision = CHANGE_SENT; - event.deleted = true; - event.name = path; - event.path = path; - char* line = change_render_format(config->out_format, config, &event); - if (line) { - /* Escape the whole rendered line, exactly like change_emit() does - for a real transfer, so a control byte in the peer-supplied path - cannot forge output. */ - char* escaped = output_escape(line, config->eight_bit_output); - printf("%s\n", escaped ? escaped : line); - free(escaped); - free(line); - } - } else { - char* escaped = output_escape(path, config->eight_bit_output); - if (config->itemize_changes) - printf("*deleting %s\n", escaped ? escaped : path); - else - printf("deleting %s\n", escaped ? escaped : path); - free(escaped); - } - } - } + print_delete_reports(config, would_delete); array_list_delete(would_delete); if (!receive_status(client->file_descriptor, &status)) goto dry_fail; diff --git a/src/server/receiver.c b/src/server/receiver.c index d90e934..e74b3d9 100644 --- a/src/server/receiver.c +++ b/src/server/receiver.c @@ -61,13 +61,18 @@ bool receiver_send_final_success(int fd, const Config* config, const ReceiverOut } bool receiver_send_stats_frame(int fd, const Config* config, const ReceiverStats* stats, - const struct ArrayList* would_delete) { + const struct ArrayList* would_delete, + const struct ArrayList* deleted_paths) { if (!config->report_stats) return true; ReceiverStats local; memset(&local, 0, sizeof(local)); const ReceiverStats* out = stats ? stats : &local; - size_t count = would_delete ? (size_t)would_delete->size : 0; + /* The path list carries the dry-run would-delete set for a -n run and the + actually-removed set for a real --info=del run. */ + const struct ArrayList* paths = + config->dry_run ? would_delete : (config->report_deletes ? deleted_paths : NULL); + size_t count = paths ? (size_t)paths->size : 0; if (count > (size_t)MAX_MANIFEST_ENTRIES) count = MAX_MANIFEST_ENTRIES; ReceiverStats record = *out; @@ -76,7 +81,7 @@ bool receiver_send_stats_frame(int fd, const Config* config, const ReceiverStats !send_int(fd, (int)count)) return false; for (size_t i = 0; i < count; i++) { - const char* path = (const char*)would_delete->items[i]; + const char* path = (const char*)paths->items[i]; if (!send_wire_str(fd, path ? path : "")) return false; } @@ -90,6 +95,20 @@ static void receiver_tally_deleted(const ReceiverSink* sink, size_t deleted) { sink->stats->deleted_files += deleted; } +/* Observer for --info=del: record each truly-removed destination-relative path + in the ArrayList passed as the observer context, so the terminal STATUS_STATS + frame can list it. A failed append is best-effort (the deletion already + happened; output is cosmetic). Shared by the single-threaded receiver and + the -m pipeline's deferred commit. */ +void receiver_record_deleted_path(void* context, const char* rel_path) { + ArrayList* paths = context; + if (!paths || !rel_path) + return; + char* copy = str_dup(rel_path); + if (copy && !array_list_add(paths, copy)) + free(copy); +} + static bool receiver_process_chunk(Chunk* chunk, const ReceiverSink* sink) { if (!chunk || !sink || !sink->store_file) return false; @@ -398,9 +417,13 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver --max-delete-capped commit still succeeds and the transfer proceeds; the terminal success frame reports the cap. */ size_t deleted = 0; - DeleteCommitResult deletion = (config->use_delete || config->delete_missing_args) - ? manifest_delete_all_counted(config, manifest, &deleted) - : DELETE_COMMIT_OK; + DeletePathObserver observer = + sink->deleted_paths ? receiver_record_deleted_path : NULL; + DeleteCommitResult deletion = + (config->use_delete || config->delete_missing_args) + ? manifest_delete_all_observed(config, manifest, &deleted, observer, + (void*)sink->deleted_paths) + : DELETE_COMMIT_OK; receiver_tally_deleted(sink, deleted); delete_manifest_free(manifest); if (deletion == DELETE_COMMIT_ERROR) { @@ -435,8 +458,12 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver send_status(file_descriptor, STATUS_ERROR); goto fail; } - if (!plan_session) + if (!plan_session) { plan_session = delete_plan_session_create(config); + if (plan_session && sink->deleted_paths) + delete_plan_session_set_delete_observer(plan_session, receiver_record_deleted_path, + (void*)sink->deleted_paths); + } if (!plan_session || delete_plan_session_receive(plan_session, config, file_descriptor) != 0) goto fail; if (delete_plan_session_limit_reached(plan_session) && !delete_limit_noted && @@ -479,8 +506,10 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver deferred_manifest = NULL; } else { size_t deleted = 0; - DeleteCommitResult deletion = - manifest_delete_all_counted(config, deferred_manifest, &deleted); + DeletePathObserver observer = + sink->deleted_paths ? receiver_record_deleted_path : NULL; + DeleteCommitResult deletion = manifest_delete_all_observed( + config, deferred_manifest, &deleted, observer, (void*)sink->deleted_paths); receiver_tally_deleted(sink, deleted); delete_manifest_free(deferred_manifest); deferred_manifest = NULL; @@ -498,6 +527,9 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver hands the session to its caller instead, which commits after the disk writer drained. */ if (plan_session) { + if (sink->deleted_paths) + delete_plan_session_set_delete_observer(plan_session, receiver_record_deleted_path, + (void*)sink->deleted_paths); if (pending_plans) { *pending_plans = plan_session; plan_session = NULL; @@ -569,6 +601,8 @@ typedef struct { --delete would-delete path list collected while processing the manifest. */ ReceiverStats stats; ArrayList* would_delete; + /* --info=del: actually-removed paths collected during the delete commit. */ + ArrayList* deleted_paths; } ReceiverSaveContext; static bool receiver_save_file(File* file, void* context_pointer) { @@ -621,7 +655,8 @@ static void receiver_note_delete_limit(void* context_pointer) { static bool receiver_send_success_frame(int fd, void* context_pointer) { ReceiverSaveContext* context = context_pointer; Status final_status = context->delete_limit_reached ? STATUS_DELETE_LIMIT : STATUS_OK; - if (!receiver_send_stats_frame(fd, context->config, &context->stats, context->would_delete)) + if (!receiver_send_stats_frame(fd, context->config, &context->stats, context->would_delete, + context->deleted_paths)) return false; /* Server-contacting --dry-run: nothing was staged or written, so there is nothing to publish and no directory times to stamp. */ @@ -651,8 +686,12 @@ int receiver_receive_files(Config* config, int file_descriptor) { ReceiverSaveContext context = {.config = config, .outcomes = {0}}; dir_time_list_init(&context.dir_times); context.would_delete = array_list_create(free); - if (!context.would_delete) + context.deleted_paths = array_list_create(free); + if (!context.would_delete || !context.deleted_paths) { + array_list_delete(context.would_delete); + array_list_delete(context.deleted_paths); return -1; + } ReceiverSink sink = {receiver_save_file, &context, true, @@ -660,12 +699,14 @@ int receiver_receive_files(Config* config, int file_descriptor) { receiver_send_success_frame, receiver_note_delete_limit, &context.stats, - context.would_delete}; + context.would_delete, + context.deleted_paths}; int ret = receiver_process(config, file_descriptor, &sink); if (ret != 0 && config->delay_updates && config->delay_context) delay_updates_cleanup(config->delay_context); receiver_outcomes_destroy(&context.outcomes); dir_time_list_free(&context.dir_times); array_list_delete(context.would_delete); + array_list_delete(context.deleted_paths); return ret; } diff --git a/src/server/receiver.h b/src/server/receiver.h index 0dc2277..3b2dcff 100644 --- a/src/server/receiver.h +++ b/src/server/receiver.h @@ -46,11 +46,20 @@ typedef struct { carries the -n/--dry-run --delete path list. */ ReceiverStats* stats; struct ArrayList* would_delete; + /* When --info=del requested it, receiver-owned strings for every path the + deletion commit ACTUALLY removed, sent in the terminal STATUS_STATS frame's + path list so the sender can print rsync's `deleting PATH` lines. */ + struct ArrayList* deleted_paths; } ReceiverSink; bool receiver_outcomes_append(ReceiverOutcomes* outcomes, unsigned char code); void receiver_outcomes_destroy(ReceiverOutcomes* outcomes); +/* DeletePathObserver implementation for --info=del: `context` is an ArrayList* + that receives owned copies of every truly-removed destination-relative path. + Shared by the single-threaded receiver and the -m pipeline's deferred commit. */ +void receiver_record_deleted_path(void* context, const char* rel_path); + /* Send the terminal success frame. `final_status` is usually STATUS_OK, or STATUS_DELETE_LIMIT when a --max-delete commit was capped. */ bool receiver_send_final_success(int fd, const Config* config, const ReceiverOutcomes* outcomes, @@ -60,7 +69,8 @@ bool receiver_send_final_success(int fd, const Config* config, const ReceiverOut non-NULL, a count and that many wire strings) when the wire config requested report_stats. A no-op otherwise. */ bool receiver_send_stats_frame(int fd, const Config* config, const ReceiverStats* stats, - const struct ArrayList* would_delete); + const struct ArrayList* would_delete, + const struct ArrayList* deleted_paths); int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink); /* receiver_process with an escape hatch for the commit-style (late) deletion: diff --git a/src/server/receiver_pipeline.c b/src/server/receiver_pipeline.c index 40784ed..d8b7d0d 100644 --- a/src/server/receiver_pipeline.c +++ b/src/server/receiver_pipeline.c @@ -31,6 +31,7 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* context->delete_limit_reached = false; memset(&context->stats, 0, sizeof(context->stats)); context->would_delete = NULL; + context->deleted_paths = NULL; atomic_init(&context->cancelled, false); int init = 0; if (mtx_init(&context->mutex, mtx_plain) != thrd_success) @@ -46,6 +47,9 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* context->would_delete = array_list_create(free); if (!context->would_delete) goto fail; + context->deleted_paths = array_list_create(free); + if (!context->deleted_paths) + goto fail; return context; fail: @@ -71,6 +75,8 @@ void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { dir_time_list_free(&context->dir_times); if (context->would_delete) array_list_delete(context->would_delete); + if (context->deleted_paths) + array_list_delete(context->deleted_paths); mtx_destroy(&context->mutex); cnd_destroy(&context->condition_not_full); cnd_destroy(&context->condition_not_empty); @@ -184,7 +190,8 @@ int receive_thread(void* pipeline_context) { NULL, receiver_pipeline_note_delete_limit, &context->stats, - context->would_delete}; + context->would_delete, + context->deleted_paths}; if (receiver_process_pending((Config*)config, file_descriptor, &sink, &context->deferred_manifest, &context->deferred_plans) != 0) { receiver_thread_fail(context); diff --git a/src/server/receiver_pipeline.h b/src/server/receiver_pipeline.h index 9ad6110..099bd9c 100644 --- a/src/server/receiver_pipeline.h +++ b/src/server/receiver_pipeline.h @@ -61,6 +61,9 @@ typedef struct PipelineContextReceiver { /* -n/--dry-run --delete would-delete path list, collected by receive_thread and reported in the STATUS_STATS frame. */ struct ArrayList* would_delete; + /* --info=del actually-removed path list, collected by the deferred delete + commit in server.c and reported in the STATUS_STATS frame. */ + struct ArrayList* deleted_paths; } PipelineContextReceiver; PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue_receiver, diff --git a/src/server/server.c b/src/server/server.c index 6172f5e..47f94b2 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -956,8 +956,11 @@ void handler(int file_descriptor) { server-contacting --dry-run deletes nothing (no manifest is sent). */ if (context->deferred_manifest) { size_t deleted = 0; + DeletePathObserver observer = + config->report_deletes ? receiver_record_deleted_path : NULL; DeleteCommitResult deletion = - manifest_delete_all_counted(config, context->deferred_manifest, &deleted); + manifest_delete_all_observed(config, context->deferred_manifest, &deleted, observer, + (void*)context->deleted_paths); context->stats.deleted_files += deleted; if (deletion == DELETE_COMMIT_ERROR) { transfer_ok = false; @@ -975,6 +978,10 @@ void handler(int file_descriptor) { if (context->deferred_plans) { /* Defence in depth (the enclosing block already excludes dry-run): a -n run never commits a deletion. */ + if (config->report_deletes) + delete_plan_session_set_delete_observer(context->deferred_plans, + receiver_record_deleted_path, + (void*)context->deleted_paths); DeleteCommitResult deletion = config->dry_run ? DELETE_COMMIT_OK : delete_plan_session_commit(context->deferred_plans, config); @@ -1010,7 +1017,7 @@ void handler(int file_descriptor) { /* 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) || + context->would_delete, context->deleted_paths) || !receiver_send_final_success(file_descriptor, config, &context->outcomes, final_status)) transfer_ok = false; } else { diff --git a/src/shared/config.h b/src/shared/config.h index e9fe715..1f2d9a1 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -259,9 +259,17 @@ typedef enum SuperMode { SUPER_MODE_AUTO = 0, SUPER_MODE_ON = 1, SUPER_MODE_OFF * for -n/--dry-run --delete, the destination-relative paths it WOULD have * deleted. It is set by the client only when --stats, --progress/-P, an * --out-format token needs a wire counter (%b/%c), or a dry-run carries + * --delete; the transfer decision itself is unchanged. + * + * --info wave (protocol 2.27.0). report_deletes tells the receiver to include + * the destination-relative paths it ACTUALLY removed in its terminal + * STATUS_STATS record (the same path-list field the dry-run would-delete report + * uses), so the sender can print rsync's `deleting PATH`/`*deleting` lines for a + * real (non-dry-run) deletion. It is set only when --info=del is requested with * --delete; the transfer decision itself is unchanged. */ #define CONFIG_WIRE_OUTPUT_FIELDS(X) \ - X(report_dest_info, bool, false, BOOL) X(report_stats, bool, false, BOOL) + X(report_dest_info, bool, false, BOOL) X(report_stats, bool, false, BOOL) \ + X(report_deletes, bool, false, BOOL) /* Codec-negotiation wave (protocol 2.26.0). compression_algo is the concrete * codec the client selected for this transfer (a CompressionAlgo id) and is the @@ -988,7 +996,14 @@ typedef struct Config { * boundary, and the strict same-version handshake (config_receive rejects a * mismatched version before parsing anything else) keeps mixed deployments from * ever reaching that state. */ -#define PROTOCOL_VERSION "2.26.0" +/* (7) --info=del report (protocol 2.27.0): the config frame gains one trailing + * bool, report_deletes, appended after report_stats. When set, the receiver + * lists the paths it actually removed in the terminal STATUS_STATS path list + * (the same count-delimited list the -n/--dry-run would-delete report uses), so + * the sender can print rsync's `deleting PATH` lines for a real deletion. No + * change to the fixed STATUS_STATS record itself; only a new trailing config + * bool, which still requires the version bump for the strict lockstep. */ +#define PROTOCOL_VERSION "2.27.0" #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) /* Upper bound on total basis-dir entries (rsync caps --link-dest at 20). */ #define MAX_BASIS_DIRS 64 diff --git a/src/shared/delete_plan.c b/src/shared/delete_plan.c index ef579ac..d14d7be 100644 --- a/src/shared/delete_plan.c +++ b/src/shared/delete_plan.c @@ -426,8 +426,16 @@ struct DeletePlanSession { ArrayList* size_skipped; ArrayList* missing; ArrayList* deferred; + DeletePathObserver observer; + void* observer_context; }; +/* Report one path the session truly removed (no-op without an observer). */ +static void notify_deleted(DeletePlanSession* session, const char* rel) { + if (session && session->observer && rel) + session->observer(session->observer_context, rel); +} + DeletePlanSession* delete_plan_session_create(const Config* config) { if (!config) return NULL; @@ -642,6 +650,7 @@ static bool process_extra_dir(int dirfd, const char* name, const char* child_rel session->deleted++; session->planned++; log_deleted(child_rel); + notify_deleted(session, child_rel); *removed = true; return true; } @@ -667,6 +676,7 @@ static bool process_extra_file(int dirfd, const char* name, const char* child_re session->deleted++; session->planned++; log_deleted(child_rel); + notify_deleted(session, child_rel); } else if (errno != ENOENT) { return false; } @@ -781,8 +791,9 @@ static bool apply_missing(DeletePlanSession* session, const Config* config) { size_t deleted = 0; size_t skipped = 0; bool limit = false; - bool ok = manifest_delete_missing_args_limited(config, &manifest, remaining, &deleted, &skipped, - &limit); + bool ok = manifest_delete_missing_args_limited_observed( + config, &manifest, remaining, &deleted, &skipped, &limit, session->observer, + session->observer_context); session->deleted += deleted; session->planned += deleted; session->skipped += skipped; @@ -875,12 +886,21 @@ static bool apply_deferred_path(DeletePlanSession* session, const Config* config if (rc == 0) { session->deleted++; log_deleted(rel); + notify_deleted(session, rel); } close(parent_fd); free(leaf); return ok; } +void delete_plan_session_set_delete_observer(DeletePlanSession* session, DeletePathObserver observer, + void* context) { + if (!session) + return; + session->observer = observer; + session->observer_context = context; +} + DeleteCommitResult delete_plan_session_commit(DeletePlanSession* session, const Config* config) { if (!session || !config) return DELETE_COMMIT_ERROR; diff --git a/src/shared/delete_plan.h b/src/shared/delete_plan.h index 0f4644b..6b99ccd 100644 --- a/src/shared/delete_plan.h +++ b/src/shared/delete_plan.h @@ -5,6 +5,7 @@ #include "config.h" #include "file_receive.h" #include "protocol.h" +#include "utils.h" #include /* Per-directory delete plans (protocol 2.24.0). @@ -79,5 +80,11 @@ bool delete_plan_session_limit_reached(const DeletePlanSession* session); /* Number of destination entries the session's plans removed (or, for --delete-delay, snapshotted for removal), for the end-of-transfer stats. */ size_t delete_plan_session_deleted(const DeletePlanSession* session); +/* Install an observer invoked for every destination-relative path the session + truly removes (including the deferred --delete-delay commit), so the receiver + can report rsync's `deleting PATH` lines through the terminal STATUS_STATS + record. Pass NULL/0 to clear. */ +void delete_plan_session_set_delete_observer(DeletePlanSession* session, DeletePathObserver observer, + void* context); #endif diff --git a/src/shared/file_receive.c b/src/shared/file_receive.c index c9f9c18..1a1adc5 100644 --- a/src/shared/file_receive.c +++ b/src/shared/file_receive.c @@ -3195,8 +3195,9 @@ char* file_receive_basis_delete_relative(const Config* config, const char* path) alternate basis directories are never destination content and are skipped at any depth. Returns true unless a traversal/unlink error aborted the walk; the budget's limit_hit/skipped fields report a cap-stopped run. */ -static bool delete_extras_budgeted(const Config* config, DeleteManifest* manifest, - DeleteBudgetState* budget) { +static bool delete_extras_budgeted_observed(const Config* config, DeleteManifest* manifest, + DeleteBudgetState* budget, DeletePathObserver observer, + void* observer_context) { if (!config || !manifest || !manifest->keeps) return false; fprintf(stderr, "Deleting files not in manifest...\n"); @@ -3261,8 +3262,9 @@ static bool delete_extras_budgeted(const Config* config, DeleteManifest* manifes size_t deleted = 0; size_t skipped = 0; DeleteWalkResult result = - delete_extras_limited(config->receive_root_directory, manifest->keeps, manifest->dirs, - remaining, skips, used, &deleted, &skipped); + delete_extras_limited_observed(config->receive_root_directory, manifest->keeps, manifest->dirs, + remaining, skips, used, &deleted, &skipped, observer, + observer_context); if (owned_prefixes) { for (int i = 0; i < config->basis_count; i++) free(owned_prefixes[i]); @@ -3282,6 +3284,31 @@ static bool delete_extras_budgeted(const Config* config, DeleteManifest* manifes return true; } +static bool delete_extras_budgeted(const Config* config, DeleteManifest* manifest, + DeleteBudgetState* budget) { + return delete_extras_budgeted_observed(config, manifest, budget, NULL, NULL); +} + +/* Prefixes every observed path with a fixed subtree root, so a nested walk + (a recursively removed missing-arg directory) reports receive-root-relative + names like the rest of the delete output. */ +typedef struct { + DeletePathObserver inner; + void* inner_context; + const char* prefix; +} PrefixedDeleteObserver; + +static void prefixed_delete_observer(void* context, const char* rel) { + PrefixedDeleteObserver* prefixed = context; + if (!prefixed->inner || !rel) + return; + char* joined = path_cat((char*)prefixed->prefix, rel); + if (joined) { + prefixed->inner(prefixed->inner_context, joined); + free(joined); + } +} + /* --delete-missing-args exact-path deletions: each destination mirror in manifest->missing is an explicit user request, so it is removed even when the ordinary extras walk (with its protected prefixes) would leave it alone. The @@ -3295,8 +3322,10 @@ static bool delete_extras_budgeted(const Config* config, DeleteManifest* manifes --max-delete budget: once it is exhausted the remaining requests are skipped and counted. Returns false only on a genuine error (a confinement failure on a validated path or an I/O error), which fails the run. */ -static bool delete_missing_args_budgeted(const Config* config, DeleteManifest* manifest, - DeleteBudgetState* budget) { +static bool delete_missing_args_budgeted_observed(const Config* config, DeleteManifest* manifest, + DeleteBudgetState* budget, + DeletePathObserver observer, + void* observer_context) { if (!config || !manifest) return false; if (!manifest->missing || manifest->missing->size == 0) @@ -3413,9 +3442,12 @@ static bool delete_missing_args_budgeted(const Config* config, DeleteManifest* m budget->deleted >= budget->max_delete ? 0 : budget->max_delete - budget->deleted; size_t contents_deleted = 0; size_t contents_skipped = 0; + PrefixedDeleteObserver nested = {observer, observer_context, rel}; DeleteWalkResult walk = - no_keeps ? delete_extras_limited(full, no_keeps, NULL, remaining, NULL, 0, - &contents_deleted, &contents_skipped) + no_keeps ? delete_extras_limited_observed(full, no_keeps, NULL, remaining, NULL, 0, + &contents_deleted, &contents_skipped, + observer ? prefixed_delete_observer : NULL, + observer ? &nested : NULL) : DELETE_WALK_ERROR; if (no_keeps) array_list_delete(no_keeps); @@ -3455,6 +3487,8 @@ static bool delete_missing_args_budgeted(const Config* config, DeleteManifest* m } if (removed) { budget->deleted++; + if (observer) + observer(observer_context, rel); char* escaped = output_escape(rel, log_get_8_bit_output()); fprintf(stderr, " Deleted: %s\n", escaped ? escaped : ""); free(escaped); @@ -3541,15 +3575,25 @@ bool manifest_delete_extras(const Config* config, DeleteManifest* manifest) { bool manifest_delete_missing_args(const Config* config, DeleteManifest* manifest) { DeleteBudgetState budget = { .max_delete = SIZE_MAX, .deleted = 0, .skipped = 0, .limit_hit = false}; - return delete_missing_args_budgeted(config, manifest, &budget); + return delete_missing_args_budgeted_observed(config, manifest, &budget, NULL, NULL); } bool manifest_delete_missing_args_limited(const Config* config, DeleteManifest* manifest, size_t max_delete, size_t* deleted, size_t* skipped, bool* limit_hit) { + return manifest_delete_missing_args_limited_observed(config, manifest, max_delete, deleted, + skipped, limit_hit, NULL, NULL); +} + +bool manifest_delete_missing_args_limited_observed(const Config* config, DeleteManifest* manifest, + size_t max_delete, size_t* deleted, + size_t* skipped, bool* limit_hit, + DeletePathObserver observer, + void* observer_context) { DeleteBudgetState budget = { .max_delete = max_delete, .deleted = 0, .skipped = 0, .limit_hit = false}; - bool ok = delete_missing_args_budgeted(config, manifest, &budget); + bool ok = delete_missing_args_budgeted_observed(config, manifest, &budget, observer, + observer_context); if (deleted) *deleted = budget.deleted; if (skipped) @@ -3572,6 +3616,12 @@ DeleteCommitResult manifest_delete_all(const Config* config, DeleteManifest* man DeleteCommitResult manifest_delete_all_counted(const Config* config, DeleteManifest* manifest, size_t* deleted) { + return manifest_delete_all_observed(config, manifest, deleted, NULL, NULL); +} + +DeleteCommitResult manifest_delete_all_observed(const Config* config, DeleteManifest* manifest, + size_t* deleted, DeletePathObserver observer, + void* observer_context) { if (deleted) *deleted = 0; if (!config || !manifest) @@ -3590,9 +3640,11 @@ DeleteCommitResult manifest_delete_all_counted(const Config* config, DeleteManif .deleted = 0, .skipped = 0, .limit_hit = false}; - if (config->delete_missing_args && !delete_missing_args_budgeted(config, manifest, &budget)) + if (config->delete_missing_args && + !delete_missing_args_budgeted_observed(config, manifest, &budget, observer, observer_context)) return DELETE_COMMIT_ERROR; - if (config->use_delete && !delete_extras_budgeted(config, manifest, &budget)) + if (config->use_delete && + !delete_extras_budgeted_observed(config, manifest, &budget, observer, observer_context)) return DELETE_COMMIT_ERROR; if (deleted) *deleted = budget.deleted; diff --git a/src/shared/file_receive.h b/src/shared/file_receive.h index 2bf0c4e..a22ca5e 100644 --- a/src/shared/file_receive.h +++ b/src/shared/file_receive.h @@ -3,6 +3,7 @@ #include "config.h" #include "file_types.h" +#include "utils.h" #include /* Server-side file receive/save path. */ @@ -126,6 +127,12 @@ bool manifest_delete_missing_args(const Config* config, DeleteManifest* manifest bool manifest_delete_missing_args_limited(const Config* config, DeleteManifest* manifest, size_t max_delete, size_t* deleted, size_t* skipped, bool* limit_hit); +/* Observer-aware form of manifest_delete_missing_args_limited: `observer` (may + be NULL) is invoked for every destination-relative path truly removed. */ +bool manifest_delete_missing_args_limited_observed(const Config* config, DeleteManifest* manifest, + size_t max_delete, size_t* deleted, + size_t* skipped, bool* limit_hit, + DeletePathObserver observer, void* observer_context); /* Outcome of committing a delete manifest. LIMIT_REACHED reports rsync's partial --max-delete result: the budget allowed some deletions and the rest were skipped (the run still stores all file data but the client exits 25). */ @@ -146,6 +153,11 @@ DeleteCommitResult manifest_delete_all(const Config* config, DeleteManifest* man removed (for the end-of-transfer wire stats). `deleted` may be NULL. */ DeleteCommitResult manifest_delete_all_counted(const Config* config, DeleteManifest* manifest, size_t* deleted); +/* Observer-aware form of manifest_delete_all_counted: `observer` (may be NULL) + is invoked for every destination-relative path truly removed. */ +DeleteCommitResult manifest_delete_all_observed(const Config* config, DeleteManifest* manifest, + size_t* deleted, DeletePathObserver observer, + void* observer_context); /* -n/--dry-run --delete would-delete reporting: walk the destination exactly as the delete pass would and append (strdup'd) destination-relative paths that diff --git a/src/shared/protocol.c b/src/shared/protocol.c index b5ceb0b..2a3f844 100644 --- a/src/shared/protocol.c +++ b/src/shared/protocol.c @@ -187,12 +187,24 @@ unsigned long long io_get_bwlimit(void) { return global_bwlimit(); } +/* rsync's throttle (io.c sleep_for_bwlimit) sleeps once its unslept debt + * reaches ~100 ms of bandwidth, so its effective initial burst is about 0.1 s + * worth of bytes, not a full second. FastSync models the same with a token + * bucket whose capacity is bwlimit/10, so a throttled run paces like rsync + * instead of sending a full second's worth up front. */ +static long long bw_burst_capacity(unsigned long long bwlimit) { + if (bwlimit == 0) + return 0; + long long burst = (long long)(bwlimit / 10); + return burst > 0 ? burst : 1; +} + void protocol_session_set_bwlimit(ProtocolSession* session, unsigned long long bytes_per_sec) { if (!session) return; session->bwlimit = bytes_per_sec > (unsigned long long)LLONG_MAX ? (unsigned long long)LLONG_MAX : bytes_per_sec; - session->bw_tokens = (long long)session->bwlimit; + session->bw_tokens = bw_burst_capacity(session->bwlimit); struct timespec now; clock_gettime(CLOCK_MONOTONIC, &now); session->bw_last_refill_sec = now.tv_sec; @@ -226,8 +238,9 @@ static void bw_throttle_session(ProtocolSession* session, size_t bytes_written) long long tokens_to_add = (long long)((double)session->bwlimit * elapsed_ns / 1000000000.0); session->bw_tokens += tokens_to_add; - if (session->bw_tokens > (long long)session->bwlimit) - session->bw_tokens = (long long)session->bwlimit; + long long burst = bw_burst_capacity(session->bwlimit); + if (session->bw_tokens > burst) + session->bw_tokens = burst; session->bw_tokens -= bytes_written; @@ -238,7 +251,11 @@ static void bw_throttle_session(ProtocolSession* session, size_t bytes_written) poll(NULL, 0, (int)(deficit_us / 1000)); else usleep((useconds_t)deficit_us); + /* Reset the bucket AFTER the sleep: crediting the sleep duration as elapsed + refill time would cancel half the throttle (the next call would see the + whole sleep as refill and immediately grant a fresh burst). */ session->bw_tokens = 0; + clock_gettime(CLOCK_MONOTONIC, &now); session->bw_last_refill_sec = now.tv_sec; session->bw_last_refill_nsec = now.tv_nsec; } diff --git a/src/shared/utils.c b/src/shared/utils.c index 14bdd7d..ad5db18 100644 --- a/src/shared/utils.c +++ b/src/shared/utils.c @@ -616,7 +616,8 @@ static bool is_synced_dir(const PathIndex* dirs, const char* rel) { static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* keep, const PathIndex* dirs, DeleteBudget* budget, const DeleteSkipEntry* skips, int skip_count, bool parent_deletable, - bool* all_removed) { + bool* all_removed, DeletePathObserver observer, + void* observer_context) { /* openat(dirfd, ".") opens an independent file description: a dup() would share dirfd's file offset and a prior pass could leave the stream drained. */ int scanfd = openat(dirfd, ".", O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC); @@ -666,7 +667,7 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* k bool child_all_removed = false; if (childfd >= 0) { if (!delete_extras_fd(childfd, child_rel, keep, dirs, budget, skips, skip_count, deletable, - &child_all_removed)) + &child_all_removed, observer, observer_context)) operation_ok = false; close(childfd); } else if (errno != ENOENT) { @@ -693,6 +694,8 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* k local_survives = true; } else { budget->deleted++; + if (observer) + observer(observer_context, child_rel); } } else { local_survives = true; @@ -713,6 +716,8 @@ static bool delete_extras_fd(int dirfd, const char* rel_path, const PathIndex* k local_survives = true; } else { budget->deleted++; + if (observer) + observer(observer_context, child_rel); char* escaped_path = output_escape(child_rel, log_get_8_bit_output()); fprintf(stderr, " Deleted: %s\n", escaped_path ? escaped_path : ""); free(escaped_path); @@ -867,10 +872,11 @@ bool delete_extras_list(const char* dest_root, const ArrayList* manifest, return ok; } -DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* manifest, - const ArrayList* synced_dirs, size_t max_delete, - const DeleteSkipEntry* skips, int skip_count, - size_t* deleted_out, size_t* skipped_out) { +DeleteWalkResult delete_extras_limited_observed(const char* dest_root, const ArrayList* manifest, + const ArrayList* synced_dirs, size_t max_delete, + const DeleteSkipEntry* skips, int skip_count, + size_t* deleted_out, size_t* skipped_out, + DeletePathObserver observer, void* observer_context) { if (deleted_out) *deleted_out = 0; if (skipped_out) @@ -911,7 +917,7 @@ DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* m DeleteBudget budget = {.max_delete = max_delete, .deleted = 0, .skipped = 0, .limit_hit = false}; bool all_removed = false; bool ok = delete_extras_fd(rootfd, "", &keep, have_dirs ? &dirs : NULL, &budget, skips, - skip_count, false, &all_removed); + skip_count, false, &all_removed, observer, observer_context); if (close(rootfd) != 0) ok = false; path_index_free(&keep); @@ -926,6 +932,14 @@ DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* m return budget.limit_hit ? DELETE_WALK_LIMIT_REACHED : DELETE_WALK_OK; } +DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* manifest, + const ArrayList* synced_dirs, size_t max_delete, + const DeleteSkipEntry* skips, int skip_count, + size_t* deleted_out, size_t* skipped_out) { + return delete_extras_limited_observed(dest_root, manifest, synced_dirs, max_delete, skips, + skip_count, deleted_out, skipped_out, NULL, NULL); +} + bool delete_extras(const char* dest_root, const ArrayList* manifest) { return delete_extras_limited(dest_root, manifest, NULL, SIZE_MAX, NULL, 0, NULL, NULL) == DELETE_WALK_OK; diff --git a/src/shared/utils.h b/src/shared/utils.h index 578706f..e56f019 100644 --- a/src/shared/utils.h +++ b/src/shared/utils.h @@ -134,6 +134,19 @@ bool path_under_skip_prefix(const char* child_rel, bool at_root, const DeleteSki cap and returns DELETE_WALK_LIMIT_REACHED when more extras remained. `deleted_out`/`skipped_out` optionally receive the number of entries removed and the number skipped because of the cap. */ +/* Optional per-deletion observer: called for each destination-relative path + actually removed (a file, symlink, or directory), in removal order, so the + receiver can stream rsync's `--info=del`/`--info=remove` lines. */ +typedef void (*DeletePathObserver)(void* context, const char* rel_path); + +/* `delete_extras_limited_observed` is delete_extras_limited with an optional + observer; the observer is invoked only for entries truly removed. */ +DeleteWalkResult delete_extras_limited_observed(const char* dest_root, const ArrayList* manifest, + const ArrayList* synced_dirs, size_t max_delete, + const DeleteSkipEntry* skips, int skip_count, + size_t* deleted_out, size_t* skipped_out, + DeletePathObserver observer, void* observer_context); + DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* manifest, const ArrayList* synced_dirs, size_t max_delete, const DeleteSkipEntry* skips, int skip_count, diff --git a/tests/integration/test_fault_injection.py b/tests/integration/test_fault_injection.py index 80d6a8d..2694b31 100644 --- a/tests/integration/test_fault_injection.py +++ b/tests/integration/test_fault_injection.py @@ -36,7 +36,7 @@ from common import ( # noqa: E402 verify_transfer, ) -PROTOCOL_VERSION = b"2.26.0" +PROTOCOL_VERSION = b"2.27.0" STATUS_MANIFEST = 5 STATUS_OK = 0 diff --git a/tests/integration/test_preflight.py b/tests/integration/test_preflight.py index b78c808..db3c62c 100644 --- a/tests/integration/test_preflight.py +++ b/tests/integration/test_preflight.py @@ -94,14 +94,14 @@ def _seed_protocol_source(source): class TestProtocol: @pytest.mark.ci def test_protocol_current_version_accepted(self, shared_server): - """--protocol=2.26.0 (the current PROTOCOL_VERSION) is accepted and the + """--protocol=2.27.0 (the current PROTOCOL_VERSION) is accepted and the transfer completes normally.""" source = os.path.join(TEST_DATA_DIR, "proto_ok_src") dest = os.path.join(TEST_DATA_DIR, "proto_ok_dst") shutil.rmtree(dest, ignore_errors=True) os.makedirs(dest) _seed_protocol_source(source) - result, _ = run_client(source, dest, flags=["--protocol=2.26.0"], + result, _ = run_client(source, dest, flags=["--protocol=2.27.0"], port=shared_server.port) assert result.returncode == 0, \ f"--protocol current run failed: {(result.stderr or result.stdout)[:400]}" diff --git a/tests/test_client_cli.c b/tests/test_client_cli.c index a9ae488..3b16032 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -318,7 +318,7 @@ static void test_parse_args_protocol_accept_current() { Config* cfg = valid_client_config(); EXPECT_NOT_NULL(cfg); char* argv_equals[] = {"fastsync", "--source-dir", "/src", - "--dest-dir", "/dst", "--protocol=2.26.0"}; + "--dest-dir", "/dst", "--protocol=2.27.0"}; int positional_args[2]; int positional_count = 0; EXPECT_EQ_INT(parse_args(cfg, 6, argv_equals, positional_args, &positional_count), 0); @@ -328,7 +328,7 @@ static void test_parse_args_protocol_accept_current() { cfg = valid_client_config(); EXPECT_NOT_NULL(cfg); char* argv_space[] = {"fastsync", "--source-dir", "/src", "--dest-dir", - "/dst", "--protocol", "2.26.0"}; + "/dst", "--protocol", "2.27.0"}; positional_count = 0; EXPECT_EQ_INT(parse_args(cfg, 7, argv_space, positional_args, &positional_count), 0); EXPECT_EQ_STR(cfg->version, PROTOCOL_VERSION); diff --git a/tests/test_config.c b/tests/test_config.c index b7e6c7b..c25ce36 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -2831,14 +2831,15 @@ static void golden_config_populate(Config* c) { c->copy_as_gid = 222; } -/* The pinned golden frame (protocol 2.26.0). The values below are the only +/* The pinned golden frame (protocol 2.27.0). The values below are the only * thing that ties the generated table to the historical wire format; update * them ONLY with a PROTOCOL_VERSION bump and a documented reason. The 2.24.0 * delete-plan wave changed only the version string; 2.25.0 appended the - * report_stats bool and 2.26.0 appended the compression_algo int. The - * byte-exact values are recomputed for the merged layout. */ -#define GOLDEN_WIRE_LEN 705 -#define GOLDEN_WIRE_HASH 4673424031554175633ULL + * report_stats bool, 2.26.0 appended the compression_algo int, and 2.27.0 + * appended the report_deletes bool. The byte-exact values are recomputed for + * the merged layout. */ +#define GOLDEN_WIRE_LEN 709 +#define GOLDEN_WIRE_HASH 14423869696887880000ULL static unsigned long long fnv1a_64(const unsigned char* buf, size_t len) { unsigned long long h = 1469598103934665603ULL;