feat(parity): real --info=del deletion lines + rsync throttle pacing
- Wire: config frame gains report_deletes (protocol 2.26.0 -> 2.27.0); the receiver lists actually-removed paths in the STATUS_STATS path list, so the sender prints rsync's `deleting PATH` / `*deleting PATH` lines for a real --delete run (and -i/out-format). Observers threaded through the manifest, missing-args and per-directory delete engines; golden wire len/hash updated. - bwlimit: throttle now paces like rsync 3.4.1 -- ~100ms burst capacity and the sleep is no longer credited as refill, so 4 MiB at 1024/2048 KiB/s matches rsync within ~4% (was ~2x too fast).
This commit is contained in:
+53
-12
@@ -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;
|
||||
}
|
||||
|
||||
+11
-1
@@ -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:
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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,
|
||||
|
||||
+9
-2
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user