diff --git a/CMakeLists.txt b/CMakeLists.txt index 727150b..0b7e10c 100644 --- a/CMakeLists.txt +++ b/CMakeLists.txt @@ -89,6 +89,7 @@ set(SHARED_SRCS src/shared/daemon_limits.c src/shared/data.c src/shared/delay_updates.c + src/shared/delete_plan.c src/shared/delta.c src/shared/file.c src/shared/file_list.c diff --git a/src/client/client_send.c b/src/client/client_send.c index cdc1cb3..4d900c7 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -7,6 +7,7 @@ #include "compression.h" #include "config.h" #include "data.h" +#include "delete_plan.h" #include "delta.h" #include "file.h" #include "file_list.h" @@ -1250,7 +1251,7 @@ static bool send_delete_manifest_early(Client* client, ArrayList* manifest, directory and *io_error_out reports it (the caller still performs the deletion but reports the run as errored). */ static bool scan_paths_only(const Config* config, const ScannerOptions* options, - ArrayList* manifest, bool* io_error_out) { + ArrayList* manifest, DeletePlanSender* plans, bool* io_error_out) { if (io_error_out) *io_error_out = false; DirectoryScanner* scanner = @@ -1260,11 +1261,27 @@ static bool scan_paths_only(const Config* config, const ScannerOptions* options, bool ok = true; Chunk* chunk; while ((chunk = directory_scanner_next(scanner)) != NULL) { - if (!add_chunk_to_manifest(manifest, chunk)) { + if (manifest && !add_chunk_to_manifest(manifest, chunk)) { ok = false; chunk_destroy(chunk); break; } + if (plans) { + for (int i = 0; i < chunk->element_count; i++) { + File* f = chunk->items[i]; + if (!f) + continue; + const char* path = file_wire_path(f); + if (!delete_plan_sender_add(plans, path, f->is_dir)) { + ok = false; + break; + } + } + if (!ok) { + chunk_destroy(chunk); + break; + } + } chunk_destroy(chunk); } if (ok && directory_scanner_failed(scanner)) @@ -1275,6 +1292,25 @@ static bool scan_paths_only(const Config* config, const ScannerOptions* options, return ok; } +/* Transmit any not-yet-sent per-directory delete plan needed by the entries in + * `chunk` (ancestors root-first, then the entry's own directory for --dirs + * entries) before its data frames go out, so --delete-during/--delete-delay + * clear a directory's extras (and any type conflict) before the directory's + * first write. */ +static int send_chunk_delete_plans(Client* client, DeletePlanSender* plans, const Chunk* chunk) { + if (!plans) + return 0; + for (int i = 0; i < chunk->element_count; i++) { + File* f = chunk->items[i]; + if (!f) + continue; + if (delete_plan_send_for_path(client->file_descriptor, plans, file_wire_path(f), f->is_dir) != + 0) + return -1; + } + return 0; +} + static int incremental_check(Client* client, File* file, const Config* config, DeltaSignature** out_sig, unsigned long long* resume_offset) { *out_sig = NULL; @@ -2048,6 +2084,16 @@ static int send_chunks_multithreaded(void* pipeline_context) { protocol_session_unbind(); return thrd_error; } + } else if (context->delete_plans) { + /* --delete-during/--delete-delay: transmit the receive root's plan before + any data, exactly like rsync's first generator directory. */ + if (delete_plan_send_root(client->file_descriptor, context->delete_plans) != 0) { + pipeline_cancel(context); + disconnect_transfer_client(client); + mark_sender_done(context); + protocol_session_unbind(); + return thrd_error; + } } while (true) { @@ -2087,6 +2133,15 @@ static int send_chunks_multithreaded(void* pipeline_context) { } break; } + if (send_chunk_delete_plans(client, context->delete_plans, current_chunk) != 0) { + log_message(LOG_LEVEL_ERROR, "unexpected error while sending delete plan"); + chunk_destroy(current_chunk); + pipeline_cancel(context); + disconnect_transfer_client(client); + mark_sender_done(context); + protocol_session_unbind(); + return thrd_error; + } if (send_chunk_with_removal(client, current_chunk, context->config, context->remove_source_files) != 0) { log_message(LOG_LEVEL_ERROR, "unexpected error while sending chunk"); @@ -2134,7 +2189,7 @@ static int send_chunks_multithreaded(void* pipeline_context) { "unscanned source mirrors are not deleted"); else log_message(LOG_LEVEL_WARNING, "transfer stopped early (stop deadline)"); - } else if (context->config->use_delete && !context->early_delete) { + } else if (context->config->use_delete && !context->early_delete && !context->delete_plans) { /* Empty keep-set + scan I/O error must not delete the whole destination (the source may not be genuinely empty -- see send_files). */ bool empty_io; @@ -2151,7 +2206,8 @@ static int send_chunks_multithreaded(void* pipeline_context) { context->size_skipped_paths, context->missing_args, context->synced_dirs) != 0) goto send_fail; - } else if (context->config->delete_missing_args && !context->early_delete) { + } else if (context->config->delete_missing_args && !context->early_delete && + !context->delete_plans) { /* --delete-missing-args without --delete: no keep-set is built, but the exact-delete paths still ride the same manifest frame (commit once the transfer succeeded). */ @@ -2223,14 +2279,15 @@ static int scan_directory_multithreaded(void* pipeline_context) { synchronized directories here (the size-prune protection is collected in every mode). The early modes already transmitted the pre-scan keep-set and its protected lists, so the data pass must not append to them again. */ - if (!context->early_delete) { + if (!context->early_delete && !context->delete_plans) { prepared.options.excluded_paths = context->excluded_paths; /* The root marker for a full recursive transfer is already in the list; do not let the scanner append every directory to it. */ if (context->config->files_from_set != NULL) prepared.options.synced_dirs = context->synced_dirs; } - prepared.options.size_skipped_paths = context->size_skipped_paths; + if (!context->delete_plans) + prepared.options.size_skipped_paths = context->size_skipped_paths; bool dirs_mode = prepared.options.dirs; /* -H also selects the sequential scanner (see the comment at the branch), * so the loop below must choose the scanner by which object exists, not by @@ -2268,7 +2325,7 @@ static int scan_directory_multithreaded(void* pipeline_context) { failed = use_dscanner ? directory_scanner_failed(dscanner) : parallel_scanner_failed(scanner); break; } - if (context->config->use_delete && !context->early_delete) { + if (context->config->use_delete && !context->early_delete && !context->delete_plans) { mtx_lock(&context->mutex_scanner); bool manifest_ok = add_chunk_to_manifest(context->manifest, current_chunk); mtx_unlock(&context->mutex_scanner); @@ -2531,6 +2588,7 @@ int send_files(Config* config) { int ret = 1; DirectoryScanner* scanner = NULL; ArrayList* manifest = NULL; + DeletePlanSender* plan_sender = NULL; ArrayList* remove_sources = NULL; /* P7 Wave D: captured source directory times, transmitted in trailing STATUS_DIR_TIMES frame(s) (only when metadata rides the wire). */ @@ -2541,6 +2599,10 @@ int send_files(Config* config) { ArrayList* size_skipped = NULL; ArrayList* synced_dirs = NULL; bool delete_early = config->use_delete && config_delete_timing_early(config); + /* -d/--dirs does not recurse, so a per-directory plan would carry no child + information and could delete the contents of an untraversed directory; + fall back to the whole-tree end-of-transfer commit for that mode. */ + bool delete_per_dir = config->use_delete && config_delete_timing_per_dir(config) && !config->dirs; bool send_failed = false; bool had_scan_io = false; PreparedScanner prepared; @@ -2591,10 +2653,11 @@ int send_files(Config* config) { prepared.options.synced_dirs = synced_dirs; } } - /* The late-timing modes (plain --delete / --delete-after / --delete-delay) - build the manifest while streaming and send it after the last data frame. - The early modes (--delete-before/--delete-during) send it up front from a - dedicated path-only pre-scan, so no manifest is kept during the data pass. */ + /* The late-timing modes (plain --delete / --delete-after) build the manifest + while streaming and send it after the last data frame. --delete-before + sends a whole-tree keep-set up front; --delete-during/--delete-delay build a + per-directory plan set up front (paths only) and stream the plans alongside + the data, so no manifest is kept during the data pass. */ if (delete_early) { /* Pass 1: collect the complete keep-set (paths only, no data loaded) and transmit it now, before any file data. The receiver removes extras and @@ -2602,7 +2665,8 @@ int send_files(Config* config) { ArrayList* early_manifest = array_list_create(free); if (!early_manifest) goto send_fail; - bool prescan_ok = scan_paths_only(config, &prepared.options, early_manifest, &had_scan_io); + bool prescan_ok = + scan_paths_only(config, &prepared.options, early_manifest, NULL, &had_scan_io); bool early_ok = false; if (prescan_ok) { /* A scan that hit an I/O error and produced NO keep entries is ambiguous @@ -2628,6 +2692,33 @@ int send_files(Config* config) { prepared.options.synced_dirs = NULL; if (!prescan_ok || !early_ok) goto send_fail; + } else if (delete_per_dir) { + /* --delete-during/--delete-delay: build one plan per source directory from a + path-only pre-scan and transmit the root plan now, before any data, so the + receive root's extras are handled exactly like rsync's first generator + directory. The remaining plans are streamed with the data below. */ + plan_sender = delete_plan_sender_create(); + if (!plan_sender) + goto send_fail; + bool prescan_ok = scan_paths_only(config, &prepared.options, NULL, plan_sender, &had_scan_io); + bool plans_ok = false; + if (prescan_ok) { + delete_plan_sender_finalize(plan_sender, config->files_from_set ? synced_dirs : NULL); + delete_plan_sender_set_config(plan_sender, excluded, size_skipped, missing_args); + if (had_scan_io && delete_plan_sender_empty(plan_sender)) { + log_message(LOG_LEVEL_ERROR, + "source scan hit an I/O error before finding any file; refusing to delete " + "with an empty keep-set (--delete)"); + prescan_ok = false; + } else { + plans_ok = delete_plan_send_root(client->file_descriptor, plan_sender) == 0; + } + } + prepared.options.excluded_paths = NULL; + prepared.options.size_skipped_paths = NULL; + prepared.options.synced_dirs = NULL; + if (!prescan_ok || !plans_ok) + goto send_fail; } else if (config->use_delete) { manifest = array_list_create(free); if (!manifest) @@ -2706,6 +2797,11 @@ int send_files(Config* config) { goto send_fail; } } + if (send_chunk_delete_plans(client, plan_sender, current_chunk) != 0) { + chunk_destroy(current_chunk); + send_failed = true; + break; + } if (send_chunk_with_removal(client, current_chunk, config, remove_sources) != 0) { log_message(LOG_LEVEL_ERROR, "Failed to send chunk"); chunk_destroy(current_chunk); @@ -2763,12 +2859,13 @@ int send_files(Config* config) { "an empty keep-set (--delete)"); goto send_fail; } - if ((manifest || config->delete_missing_args) && !delete_early) { + if ((manifest || config->delete_missing_args) && !delete_early && !delete_per_dir) { /* Late (commit) ordering: all file data is out; transmit the manifest so the receiver commits the extras walk (--delete) and/or the --delete-missing-args exact-path deletions only after the transfer - succeeds. In the early modes (--delete-before/--delete-during) the - manifest already went out up front, so nothing is re-sent here. */ + succeeds. In the early modes (--delete-before) and the per-directory + modes the deletion already went out with the data, so nothing is + re-sent here. */ if (send_delete_manifest(client->file_descriptor, manifest, excluded, size_skipped, missing_args, synced_dirs) != 0) { if (manifest) { @@ -2817,6 +2914,8 @@ send_fail: here even on success without --delete, fixing a pre-existing leak. */ if (manifest) array_list_delete(manifest); + if (plan_sender) + delete_plan_sender_destroy(plan_sender); if (excluded) array_list_delete(excluded); if (size_skipped) @@ -2916,11 +3015,6 @@ int send_files_multithreaded(Config** config_ptr) { config->stop_at, now_mono); bool collect_excluded = config->use_delete && !config->delete_excluded; if (config->use_delete) { - context->manifest = array_list_create(free); - if (!context->manifest) { - pipeline_context_sender_destroy(context); - return 1; - } if (collect_excluded) { context->excluded_paths = array_list_create(free); if (!context->excluded_paths) { @@ -2946,11 +3040,15 @@ int send_files_multithreaded(Config** config_ptr) { return 1; } } - if (config_delete_timing_early(config)) { - /* --delete-before/--delete-during: build the complete keep-set manifest + /* -d/--dirs does not recurse, so a per-directory plan would carry no child + information and could delete the contents of an untraversed directory; + fall back to the whole-tree end-of-transfer commit for that mode. */ + bool per_dir = config_delete_timing_per_dir(config) && !config->dirs; + if (config_delete_timing_early(config) || per_dir) { + /* --delete-before / --delete-during / --delete-delay: build the keep-set (paths only, nothing loaded or sent) up front so the sender thread can - transmit it before the first data byte. The path-only pre-scan also - fills the protected excluded prefixes and synchronized directories. */ + transmit it before/with the data. The path-only pre-scan also fills the + protected excluded prefixes and synchronized directories. */ PreparedScanner prepared; memset(&prepared, 0, sizeof(prepared)); bool prepared_ok = prepare_scanner(config, config->scanner_threads, &prepared); @@ -2963,12 +3061,29 @@ int send_files_multithreaded(Config** config_ptr) { if (config->files_from_set != NULL) prepared.options.synced_dirs = context->synced_dirs; } - bool prebuilt = prepared_ok && scan_paths_only(config, &prepared.options, context->manifest, - &context->scan_had_io_error); + if (per_dir) { + context->delete_plans = delete_plan_sender_create(); + prepared_ok = prepared_ok && context->delete_plans != NULL; + } else { + context->manifest = array_list_create(free); + prepared_ok = prepared_ok && context->manifest != NULL; + } + bool prebuilt = + prepared_ok && scan_paths_only(config, &prepared.options, context->manifest, + context->delete_plans, &context->scan_had_io_error); prepared_scanner_destroy(&prepared); - if (prebuilt && context->scan_had_io_error && context->manifest->size == 0) { - /* Empty keep-set + scan I/O error: refusing an empty keep-set manifest - would have deleted the whole destination (see send_files). */ + if (per_dir && prebuilt) { + delete_plan_sender_finalize(context->delete_plans, + config->files_from_set ? context->synced_dirs : NULL); + delete_plan_sender_set_config(context->delete_plans, context->excluded_paths, + context->size_skipped_paths, context->missing_args); + } + bool empty = per_dir + ? (context->delete_plans && delete_plan_sender_empty(context->delete_plans)) + : (context->manifest && context->manifest->size == 0); + if (prebuilt && context->scan_had_io_error && empty) { + /* Empty keep-set + scan I/O error: refusing an empty keep-set would + have deleted the whole destination (see send_files). */ log_message(LOG_LEVEL_ERROR, "source scan hit an I/O error before finding any file; refusing to delete " "with an empty keep-set (--delete)"); @@ -2978,12 +3093,19 @@ int send_files_multithreaded(Config** config_ptr) { pipeline_context_sender_destroy(context); return 1; } - context->early_delete = true; + if (!per_dir) + context->early_delete = true; + } else { + context->manifest = array_list_create(free); + if (!context->manifest) { + pipeline_context_sender_destroy(context); + return 1; + } } } if (config->remove_source_files) context->remove_source_files = array_list_create(source_file_destroy); - if ((config->use_delete && !context->manifest) || + if ((config->use_delete && !context->manifest && !context->delete_plans) || (config->remove_source_files && !context->remove_source_files)) { pipeline_context_sender_destroy(context); return 1; diff --git a/src/client/usage.c b/src/client/usage.c index eee2ec4..d0d8a75 100644 --- a/src/client/usage.c +++ b/src/client/usage.c @@ -66,11 +66,11 @@ void print_usage(void) { printf(" transfer has succeeded)\n"); printf(" --delete-before Delete extras before the transfer starts\n"); printf(" (implies --delete)\n"); - printf(" --delete-during Delete extras once the keep-set manifest is known,\n"); - printf(" before the data is applied (implies --delete)\n"); + printf(" --delete-during Delete a directory's extras as that directory is\n"); + printf(" processed (implies --delete)\n"); printf(" --del Alias for --delete-during\n"); - printf(" --delete-delay Delete extras only after a successful transfer\n"); - printf(" (implies --delete)\n"); + printf(" --delete-delay Record the extras during the scan but remove them\n"); + printf(" only after a successful transfer (implies --delete)\n"); printf(" --delete-after Delete only after the whole transfer succeeded\n"); printf(" (the default --delete timing; implies --delete)\n"); printf(" --delete-excluded Also delete destination files that were excluded on\n"); diff --git a/src/server/receiver.c b/src/server/receiver.c index 3d43890..fc2346e 100644 --- a/src/server/receiver.c +++ b/src/server/receiver.c @@ -3,6 +3,7 @@ #include "charset.h" #include "chunk.h" #include "config.h" +#include "delete_plan.h" #include "delay_updates.h" #include "file.h" #include "file_receive.h" @@ -244,7 +245,7 @@ static bool receiver_note_status(const struct timespec* session_start, } int receiver_process(Config* config, int file_descriptor, const ReceiverSink* sink) { - return receiver_process_pending(config, file_descriptor, sink, NULL); + return receiver_process_pending(config, file_descriptor, sink, NULL, NULL); } /* Runs the whole receive loop. The delete manifest may legitimately arrive @@ -258,7 +259,7 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si the whole transfer succeeded. See receiver_process_pending() for how the -m receiver defers that commit until its disk writer has drained. */ int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink, - DeleteManifest** pending_manifest) { + DeleteManifest** pending_manifest, DeletePlanSession** pending_plans) { Status status; if (!receive_status(file_descriptor, &status)) return -1; @@ -272,14 +273,22 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver if (!receiver_note_status(&session_start, &last_progress, status, file_descriptor, sink)) return -1; bool early_delete = config_delete_timing_early(config); + bool per_dir_delete = config_delete_timing_per_dir(config); /* Parked keep-set for the late/commit timing. Every exit path below frees it exactly once; the only exception is the successful FINISHED handoff, which transfers ownership to *pending_manifest (used by the -m receiver). */ DeleteManifest* deferred_manifest = NULL; + /* Per-directory delete session for --delete-during/--delete-delay. During the + loop it applies plans inline (during) or snapshots their extras (delay); on + a successful FINISHED it is either committed here or handed to + *pending_plans so the -m caller commits after its disk writer drained. */ + DeletePlanSession* plan_session = NULL; + bool delete_limit_noted = false; while (status == STATUS_NEXT || status == STATUS_CHUNK || status == STATUS_CHECK || status == STATUS_KEEPALIVE || status == STATUS_ABORT || status == STATUS_CHECK_BATCH || status == STATUS_MKDIR || status == STATUS_MANIFEST || status == STATUS_HARDLINK || - status == STATUS_SYMLINK || status == STATUS_SPECIAL || status == STATUS_DIR_TIMES) { + status == STATUS_SYMLINK || status == STATUS_SPECIAL || status == STATUS_DIR_TIMES || + status == STATUS_DELETE_PLAN) { if (status == STATUS_KEEPALIVE) { if (!send_status(file_descriptor, STATUS_KEEPALIVE)) goto fail; @@ -344,11 +353,10 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver goto next_status; } if (early_delete) { - /* --delete-before / --delete-during: the manifest is authoritative the - moment it arrives, before any file data. Delete now and acknowledge - so the sender only starts streaming once the deletion committed (or - failed). This is the rsync delete-before/delete-during window: a - later transfer failure does not restore these deletions. A + /* --delete-before: the whole-tree manifest is authoritative the moment + it arrives, before any file data. Delete now and acknowledge so the + sender only starts streaming once the deletion committed (or failed). + 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) @@ -364,9 +372,9 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver if (!send_status(file_descriptor, STATUS_OK)) goto fail; } else if (config->use_delete || config->delete_missing_args) { - /* Plain --delete / --delete-after / --delete-delay and the - --delete-missing-args exact-path deletions: hold the manifest and - commit it only after STATUS_FINISHED. */ + /* Plain --delete / --delete-after and the --delete-missing-args + exact-path deletions: hold the manifest and commit it only after + STATUS_FINISHED. The per-directory modes never send this frame. */ if (deferred_manifest) { log_message(LOG_LEVEL_ERROR, "Received a second delete manifest"); delete_manifest_free(deferred_manifest); @@ -380,6 +388,23 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver delete_manifest_free(manifest); } goto next_status; + } else if (status == STATUS_DELETE_PLAN) { + if (!per_dir_delete) { + log_message(LOG_LEVEL_ERROR, "Received a per-directory delete plan without a per-dir " + "delete timing"); + send_status(file_descriptor, STATUS_ERROR); + goto fail; + } + if (!plan_session) + plan_session = delete_plan_session_create(config); + 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 && + sink->note_delete_limit) { + sink->note_delete_limit(sink->context); + delete_limit_noted = true; + } + goto next_status; } else { File* file = file_receive(config, file_descriptor); if (!file) { @@ -424,6 +449,28 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver sink->note_delete_limit(sink->context); } } + /* Per-directory deletion: --delete-during already applied each plan inline, so + this only finishes the missing-args deletions; --delete-delay committed + nothing yet and applies its decompressed snapshot here. The -m receiver + hands the session to its caller instead, which commits after the disk + writer drained. */ + if (plan_session) { + if (pending_plans) { + *pending_plans = plan_session; + plan_session = NULL; + } else { + DeleteCommitResult deletion = delete_plan_session_commit(plan_session, config); + bool limit = delete_plan_session_limit_reached(plan_session); + delete_plan_session_destroy(plan_session); + plan_session = NULL; + if (deletion == DELETE_COMMIT_ERROR) { + send_status(file_descriptor, STATUS_ERROR); + goto fail; + } + if (limit && !delete_limit_noted && sink->note_delete_limit) + sink->note_delete_limit(sink->context); + } + } if (sink->send_success) { if (sink->send_success_frame) { if (!sink->send_success_frame(file_descriptor, sink->context)) @@ -436,11 +483,14 @@ int receiver_process_pending(Config* config, int file_descriptor, const Receiver fail: /* Failure exits that must not (or already did) report a STATUS_ERROR. The - parked keep-set is dropped: never commit a deletion for a failed stream. */ + parked keep-set/session is dropped: never commit a deletion for a failed + stream. */ if (deferred_manifest) { delete_manifest_free(deferred_manifest); deferred_manifest = NULL; } + if (plan_session) + delete_plan_session_destroy(plan_session); return -1; receive_error: @@ -448,6 +498,8 @@ receive_error: delete_manifest_free(deferred_manifest); deferred_manifest = NULL; } + if (plan_session) + delete_plan_session_destroy(plan_session); if (sink->send_error) send_status(file_descriptor, STATUS_ERROR); return -1; diff --git a/src/server/receiver.h b/src/server/receiver.h index e169c41..e2248d5 100644 --- a/src/server/receiver.h +++ b/src/server/receiver.h @@ -2,6 +2,7 @@ #define RECEIVER_H #include "config.h" +#include "delete_plan.h" #include "file.h" #include "file_receive.h" #include "protocol.h" @@ -53,10 +54,12 @@ int receiver_process(Config* config, int file_descriptor, const ReceiverSink* si when `pending_manifest` is non-NULL the receiver does NOT delete at STATUS_FINISHED itself; instead it stores the owned keep-set manifest there (leaving *pending_manifest untouched on early modes/errors) so the caller can - commit the deletion only after its disk writer has fully drained. Pass NULL - to keep the default behaviour (delete before the success frame). */ + commit the deletion only after its disk writer has fully drained. Likewise, + when `pending_plans` is non-NULL the --delete-delay per-directory session is + handed to the caller instead of being committed at STATUS_FINISHED. Pass NULL + for either to keep the default behaviour (delete before the success frame). */ int receiver_process_pending(Config* config, int file_descriptor, const ReceiverSink* sink, - DeleteManifest** pending_manifest); + DeleteManifest** pending_manifest, DeletePlanSession** pending_plans); int receiver_receive_files(Config* config, int file_descriptor); /* ---- Connection time bounds (anti-slowloris) ---- diff --git a/src/server/receiver_pipeline.c b/src/server/receiver_pipeline.c index 0828db8..0f0bc3a 100644 --- a/src/server/receiver_pipeline.c +++ b/src/server/receiver_pipeline.c @@ -27,6 +27,7 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* context->queued_bytes = 0; context->max_queue_bytes = 0; context->deferred_manifest = NULL; + context->deferred_plans = NULL; context->delete_limit_reached = false; atomic_init(&context->cancelled, false); int init = 0; @@ -58,6 +59,8 @@ void pipeline_context_receiver_destroy(PipelineContextReceiver* context) { config_delete(context->config); if (context->deferred_manifest) delete_manifest_free(context->deferred_manifest); + if (context->deferred_plans) + delete_plan_session_destroy(context->deferred_plans); queue_destroy(context->queue); receiver_outcomes_destroy(&context->outcomes); dir_time_list_free(&context->dir_times); @@ -164,8 +167,8 @@ int receive_thread(void* pipeline_context) { ReceiverSink sink = { receiver_enqueue_file, context, false, false, NULL, receiver_pipeline_note_delete_limit}; - if (receiver_process_pending((Config*)config, file_descriptor, &sink, - &context->deferred_manifest) != 0) { + if (receiver_process_pending((Config*)config, file_descriptor, &sink, &context->deferred_manifest, + &context->deferred_plans) != 0) { receiver_thread_fail(context); protocol_session_unbind(); return thrd_error; diff --git a/src/server/receiver_pipeline.h b/src/server/receiver_pipeline.h index 2f9d604..247aa42 100644 --- a/src/server/receiver_pipeline.h +++ b/src/server/receiver_pipeline.h @@ -41,6 +41,11 @@ typedef struct PipelineContextReceiver { transfer truly succeeded. NULL in the early delete modes (which delete at the manifest). */ DeleteManifest* deferred_manifest; + /* Per-directory delete session for --delete-delay: receive_thread snapshots + each plan's extras as it arrives and hands the session here instead of + committing while the disk writer may still be draining; server.c commits it + after both threads joined. NULL for every other timing. */ + DeletePlanSession* deferred_plans; /* Set by server.c when the deferred delete commit hit the --max-delete budget; the terminal success frame then carries STATUS_DELETE_LIMIT (rsync exit 25) while the transfer itself still succeeds. */ diff --git a/src/server/server.c b/src/server/server.c index 6084945..58fa871 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -962,6 +962,19 @@ void handler(int file_descriptor) { delete_manifest_free(context->deferred_manifest); context->deferred_manifest = NULL; } + /* --delete-delay: receive_thread snapshotted each plan's extras as it + arrived; with the disk writer drained, commit the deferred removals. + --delete-during already applied its plans on the receive thread. */ + if (context->deferred_plans) { + DeleteCommitResult deletion = delete_plan_session_commit(context->deferred_plans, config); + if (deletion == DELETE_COMMIT_ERROR) { + transfer_ok = false; + } else if (deletion == DELETE_COMMIT_LIMIT_REACHED) { + context->delete_limit_reached = true; + } + delete_plan_session_destroy(context->deferred_plans); + context->deferred_plans = NULL; + } } if (transfer_ok && !config->dry_run) { /* --delay-updates: receive_thread has finished the whole protocol stream diff --git a/src/shared/config.c b/src/shared/config.c index fc3ba3f..0a0b09d 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -228,7 +228,13 @@ Config* config_create(void) { bool config_delete_timing_early(const Config* config) { if (!config) return false; - return config->delete_before || config->delete_during; + return config->delete_before; +} + +bool config_delete_timing_per_dir(const Config* config) { + if (!config) + return false; + return config->delete_during || config->delete_delay; } /* A delete-timing flag is only meaningful together with --delete. At most one diff --git a/src/shared/config.h b/src/shared/config.h index 1f3f126..c8654e3 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -551,13 +551,17 @@ typedef struct Config { /* rsync deletion-timing family (real from Phase 3). At most one of delete_before / delete_during / delete_delay / delete_after may be set, and only together with use_delete (the CLI implies --delete for each of them). - delete_before and delete_during select the EARLY engine mode: the keep-set + delete_before selects the EARLY engine mode: the whole-tree keep-set manifest is transmitted before any file data and extras are removed then, - acknowledged, before the first data byte. delete_delay and delete_after - select the LATE commit mode: extras are removed only after the whole - transfer has succeeded (plain --delete keeps this mode). The exact - semantics and the divergences from rsync are documented in RSYNC_COMPAT.md - and in config_delete_timing_early() below. */ + acknowledged, before the first data byte. delete_during and delete_delay + select the per-directory delete-plan mode (protocol 2.24.0): one plan per + source directory is streamed in directory order, and the receiver removes + each directory's extras when its plan arrives (during) or snapshots them + and removes them only after a successful transfer (delay). delete_after + (and plain --delete) keep the whole-tree commit mode: extras are removed + from a fresh end-of-transfer destination scan only after the whole transfer + succeeded. See config_delete_timing_early()/config_delete_timing_per_dir() + below. */ /* partial_dir */ // PR #174: Partial transfer resumption /* suffix */ @@ -916,7 +920,7 @@ typedef struct Config { * carries the new report_dest_info bool appended after the --copy-as block. * This is both a config-frame layout change (one trailing bool) and a frame * sequence change (the new status). */ -#define PROTOCOL_VERSION "2.23.0" +#define PROTOCOL_VERSION "2.24.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 @@ -1010,13 +1014,17 @@ int config_parse_daemon_dest(Config* config); * 0. */ int config_parse_transport_dest(Config* config); -/* True when the negotiated delete timing performs the extra-file deletion - * BEFORE the transfer data (--delete-before / --delete-during). The flag is - * a pure function of the config and is used identically on the sender (to pick +/* True for the whole-tree delete-before timing: a complete keep-set manifest is + * transmitted before any data and committed (with an ack) before the first data + * byte. Pure function of the config, used identically on the sender (to pick * the manifest-first frame order) and the receiver (to delete when the early - * manifest arrives). When false the deletion is committed only after the whole - * transfer succeeded (--delete / --delete-after / --delete-delay). */ + * manifest arrives). */ bool config_delete_timing_early(const Config* config); +/* True for the per-directory timings (--delete-during / --delete-delay). The + * sender streams a delete plan per source directory in directory order; the + * receiver applies each plan on arrival (during) or snapshots its extras and + * commits them only after a fully-successful transfer (delay). */ +bool config_delete_timing_per_dir(const Config* config); /* Delete-timing sanity: with deletion enabled at most one timing flag may be * set (none = the default delete-after commit timing); without deletion no * timing flag may be set (each timing flag implies --delete). */ diff --git a/src/shared/delete_plan.c b/src/shared/delete_plan.c new file mode 100644 index 0000000..23e7a2b --- /dev/null +++ b/src/shared/delete_plan.c @@ -0,0 +1,844 @@ +#include "delete_plan.h" + +#include "charset.h" +#include "delay_updates.h" +#include "file.h" +#include "log.h" +#include "utils.h" +#include +#include +#include +#include +#include +#include +#include +#include +#include +#include + +/* Mirrors MAX_SERVER_DELETE_COUNT in file_receive.c: the server's hard bound on + * the number of entries one deletion commit may remove. A client + * --max-delete=NUM smaller than this replaces it for the run. */ +#define DELETE_PLAN_SERVER_LIMIT 100000U +/* Per-frame entry cap for the name sections (the dir/file child lists). */ +#define DELETE_PLAN_MAX_NAMES MAX_MANIFEST_ENTRIES + +/* ------------------------------------------------------------------ */ +/* Sender: plan builder */ +/* ------------------------------------------------------------------ */ + +typedef struct PlanNode { + char* dir; + ArrayList* files; /* basenames kept directly in dir */ + ArrayList* dirs; /* basenames of kept child directories */ + bool sent; + struct PlanNode* hash_next; +} PlanNode; + +struct DeletePlanSender { + PlanNode** buckets; + size_t capacity; + size_t count; + bool config_sent; + bool all_synced; + const ArrayList* synced_dirs; + const ArrayList* protected_prefixes; + const ArrayList* size_skipped; + const ArrayList* missing_args; + size_t entries; +}; + +static size_t plan_hash(const char* key) { + size_t h = 5381; + for (const unsigned char* p = (const unsigned char*)key; *p; p++) + h = ((h << 5) + h) + *p; + return h; +} + +static bool list_contains_str(const ArrayList* list, const char* value) { + if (!list) + return false; + for (int i = 0; i < list->size; i++) { + if (strcmp((const char*)list->items[i], value) == 0) + return true; + } + return false; +} + +static bool list_add_str_unique(ArrayList* list, const char* value) { + if (!list || !value) + return false; + if (list_contains_str(list, value)) + return true; + char* copy = str_dup(value); + if (!copy) + return false; + if (!array_list_add(list, copy)) { + free(copy); + return false; + } + return true; +} + +DeletePlanSender* delete_plan_sender_create(void) { + DeletePlanSender* sender = calloc(1, sizeof(DeletePlanSender)); + if (!sender) + return NULL; + sender->capacity = 64; + sender->buckets = calloc(sender->capacity, sizeof(PlanNode*)); + if (!sender->buckets) { + free(sender); + return NULL; + } + sender->all_synced = true; + return sender; +} + +static void plan_node_destroy(PlanNode* node) { + if (!node) + return; + free(node->dir); + array_list_delete(node->files); + array_list_delete(node->dirs); + free(node); +} + +void delete_plan_sender_destroy(DeletePlanSender* sender) { + if (!sender) + return; + for (size_t i = 0; i < sender->capacity; i++) { + PlanNode* node = sender->buckets[i]; + while (node) { + PlanNode* next = node->hash_next; + plan_node_destroy(node); + node = next; + } + } + free(sender->buckets); + free(sender); +} + +static PlanNode* plan_find(const DeletePlanSender* sender, const char* dir) { + size_t index = plan_hash(dir) & (sender->capacity - 1); + for (PlanNode* node = sender->buckets[index]; node; node = node->hash_next) { + if (strcmp(node->dir, dir) == 0) + return node; + } + return NULL; +} + +static bool plan_grow(DeletePlanSender* sender) { + size_t new_capacity = sender->capacity * 2; + PlanNode** buckets = calloc(new_capacity, sizeof(PlanNode*)); + if (!buckets) + return false; + for (size_t i = 0; i < sender->capacity; i++) { + PlanNode* node = sender->buckets[i]; + while (node) { + PlanNode* next = node->hash_next; + size_t index = plan_hash(node->dir) & (new_capacity - 1); + node->hash_next = buckets[index]; + buckets[index] = node; + node = next; + } + } + free(sender->buckets); + sender->buckets = buckets; + sender->capacity = new_capacity; + return true; +} + +static PlanNode* plan_ensure(DeletePlanSender* sender, const char* dir) { + PlanNode* node = plan_find(sender, dir); + if (node) + return node; + if (sender->count + 1 > sender->capacity * 3 / 4 && !plan_grow(sender)) + return NULL; + node = calloc(1, sizeof(PlanNode)); + if (!node) + return NULL; + node->dir = str_dup(dir); + node->files = array_list_create(free); + node->dirs = array_list_create(free); + if (!node->dir || !node->files || !node->dirs) { + plan_node_destroy(node); + return NULL; + } + size_t index = plan_hash(dir) & (sender->capacity - 1); + node->hash_next = sender->buckets[index]; + sender->buckets[index] = node; + sender->count++; + return node; +} + +static char* path_parent_dir(const char* path) { + const char* slash = strrchr(path, '/'); + if (!slash) + return str_dup("."); + if (slash == path) + return str_dup("."); + size_t len = (size_t)(slash - path); + char* parent = malloc(len + 1); + if (!parent) + return NULL; + memcpy(parent, path, len); + parent[len] = '\0'; + return parent; +} + +static char* path_base_name(const char* path) { + const char* slash = strrchr(path, '/'); + return str_dup(slash ? slash + 1 : path); +} + +/* Copy `path`, stripping a leading '/' and any trailing '/'. */ +static char* plan_clean_path(const char* path) { + while (*path == '/') + path++; + size_t len = strlen(path); + while (len > 0 && path[len - 1] == '/') + len--; + char* clean = malloc(len + 1); + if (!clean) + return NULL; + memcpy(clean, path, len); + clean[len] = '\0'; + return clean; +} + +static bool plan_ensure_ancestors(DeletePlanSender* sender, const char* dir) { + char* current = str_dup(dir); + if (!current) + return false; + bool ok = true; + while (strcmp(current, ".") != 0) { + char* parent = path_parent_dir(current); + char* base = path_base_name(current); + PlanNode* parent_node = parent ? plan_ensure(sender, parent) : NULL; + if (!parent || !base || !parent_node || !list_add_str_unique(parent_node->dirs, base)) { + ok = false; + free(parent); + free(base); + break; + } + free(base); + free(current); + current = parent; + } + free(current); + return ok; +} + +bool delete_plan_sender_add(DeletePlanSender* sender, const char* path, bool is_dir) { + if (!sender || !path) + return false; + char* clean = plan_clean_path(path); + if (!clean) + return false; + if (*clean == '\0') { + free(clean); + return true; + } + char* parent = path_parent_dir(clean); + char* base = path_base_name(clean); + PlanNode* parent_node = parent ? plan_ensure(sender, parent) : NULL; + bool ok = parent && base && parent_node; + if (ok) { + if (is_dir) { + ok = list_add_str_unique(parent_node->dirs, base) && plan_ensure(sender, clean) != NULL; + } else { + ok = list_add_str_unique(parent_node->files, base); + } + } + if (ok) + ok = plan_ensure_ancestors(sender, parent); + if (ok) + sender->entries++; + free(clean); + free(parent); + free(base); + return ok; +} + +void delete_plan_sender_finalize(DeletePlanSender* sender, const ArrayList* synced_dirs) { + if (!sender) + return; + sender->synced_dirs = synced_dirs; + sender->all_synced = synced_dirs == NULL; +} + +bool delete_plan_sender_empty(const DeletePlanSender* sender) { + return !sender || sender->entries == 0; +} + +void delete_plan_sender_set_config(DeletePlanSender* sender, const ArrayList* protected_prefixes, + const ArrayList* size_skipped, const ArrayList* missing_args) { + if (!sender) + return; + sender->protected_prefixes = protected_prefixes; + sender->size_skipped = size_skipped; + sender->missing_args = missing_args; +} + +static bool plan_is_allowed(const DeletePlanSender* sender, const char* dir) { + if (sender->all_synced) + return true; + return list_contains_str(sender->synced_dirs, dir); +} + +static int send_str_section(int fd, const ArrayList* list) { + int count = list ? list->size : 0; + if (!send_int(fd, count)) + return -1; + for (int i = 0; i < count; i++) { + if (!send_wire_str(fd, (const char*)list->items[i])) + return -1; + } + return 0; +} + +static int send_plan_node(int fd, DeletePlanSender* sender, PlanNode* node) { + if (!send_status(fd, STATUS_DELETE_PLAN)) + return -1; + if (!send_int(fd, sender->config_sent ? 0 : 1)) + return -1; + if (!sender->config_sent) { + if (send_str_section(fd, sender->protected_prefixes) != 0 || + send_str_section(fd, sender->size_skipped) != 0 || + send_str_section(fd, sender->missing_args) != 0) + return -1; + sender->config_sent = true; + } + if (!send_wire_str(fd, node->dir)) + return -1; + if (send_str_section(fd, node->dirs) != 0 || send_str_section(fd, node->files) != 0) + return -1; + node->sent = true; + return 0; +} + +static int send_prefix_plan(int fd, DeletePlanSender* sender, const char* dir) { + PlanNode* node = plan_find(sender, dir); + if (!node || node->sent) + return 0; + if (!plan_is_allowed(sender, dir)) + return 0; + return send_plan_node(fd, sender, node); +} + +int delete_plan_send_root(int fd, DeletePlanSender* sender) { + if (!sender) + return -1; + if (!plan_ensure(sender, ".")) + return -1; + return send_prefix_plan(fd, sender, "."); +} + +int delete_plan_send_for_path(int fd, DeletePlanSender* sender, const char* path, bool is_dir) { + if (!sender || !path) + return -1; + char* clean = plan_clean_path(path); + if (!clean) + return -1; + int rc = send_prefix_plan(fd, sender, "."); + if (rc == 0 && *clean != '\0') { + size_t len = strlen(clean); + size_t end = len; + if (!is_dir) { + const char* slash = strrchr(clean, '/'); + end = slash ? (size_t)(slash - clean) : 0; + } + for (size_t i = 1; i <= end && rc == 0; i++) { + if (i == end || clean[i] == '/') { + char* prefix = malloc(i + 1); + if (!prefix) { + rc = -1; + break; + } + memcpy(prefix, clean, i); + prefix[i] = '\0'; + rc = send_prefix_plan(fd, sender, prefix); + free(prefix); + } + } + } + free(clean); + return rc; +} + +/* ------------------------------------------------------------------ */ +/* Receiver: delete session */ +/* ------------------------------------------------------------------ */ + +struct DeletePlanSession { + bool defer; + bool dry_run; + size_t max_delete; + size_t deleted; + size_t skipped; + bool limit_hit; + bool config_seen; + bool missing_applied; + ArrayList* protected_prefixes; + ArrayList* size_skipped; + ArrayList* missing; + ArrayList* deferred; +}; + +DeletePlanSession* delete_plan_session_create(const Config* config) { + if (!config) + return NULL; + DeletePlanSession* session = calloc(1, sizeof(DeletePlanSession)); + if (!session) + return NULL; + session->defer = config->delete_delay; + session->dry_run = config->dry_run; + bool user_limited = + config->max_delete >= 0 && (size_t)config->max_delete < DELETE_PLAN_SERVER_LIMIT; + session->max_delete = + user_limited ? (size_t)config->max_delete : (size_t)DELETE_PLAN_SERVER_LIMIT; + session->protected_prefixes = array_list_create(free); + session->size_skipped = array_list_create(free); + session->missing = array_list_create(free); + session->deferred = array_list_create(free); + if (!session->protected_prefixes || !session->size_skipped || !session->missing || + !session->deferred) { + delete_plan_session_destroy(session); + return NULL; + } + return session; +} + +void delete_plan_session_destroy(DeletePlanSession* session) { + if (!session) + return; + array_list_delete(session->protected_prefixes); + array_list_delete(session->size_skipped); + array_list_delete(session->missing); + array_list_delete(session->deferred); + free(session); +} + +bool delete_plan_session_limit_reached(const DeletePlanSession* session) { + return session && session->limit_hit; +} + +/* True for a destination-relative path section entry (non-empty, relative, + * traversal-free). */ +static bool valid_rel_path(const char* value) { + return value && value[0] != '\0' && value[0] != '/' && !has_path_traversal(value); +} + +/* True for a single child name (non-empty, no slash, not "."/".."). */ +static bool valid_name(const char* value) { + return value && value[0] != '\0' && strcmp(value, ".") != 0 && strcmp(value, "..") != 0 && + strchr(value, '/') == NULL; +} + +/* Read one count-prefixed section. `bytes` is the running per-frame budget, + * shared across every section of the frame so a hostile peer cannot retain more + * than MAX_MANIFEST_BYTES from one STATUS_DELETE_PLAN frame. */ +static bool read_section(int fd, ArrayList* list, bool rel_path, size_t* bytes) { + int count; + if (!receive_int(fd, &count) || count < 0 || count > MAX_MANIFEST_ENTRIES) + return false; + for (int i = 0; i < count; i++) { + char* value = receive_wire_str(fd); + bool ok = value && (rel_path ? valid_rel_path(value) : valid_name(value)); + if (ok) { + size_t entry_size = strlen(value) + sizeof(char*) + 16; + if (entry_size > MAX_MANIFEST_BYTES - *bytes) { + ok = false; + } else { + *bytes += entry_size; + ok = array_list_add(list, value); + } + } + if (!ok) { + free(value); + return false; + } + } + return true; +} + +static int open_plan_dir(const Config* config, const char* dir) { + char* full = (strcmp(dir, ".") == 0) ? str_dup(config->receive_root_directory) + : path_cat(config->receive_root_directory, dir); + if (!full) + return -1; + int root_fd = utils_get_authorized_root_fd(); + int fd = -1; + if (root_fd >= 0) { + if (utils_get_authorized_root_path()) + fd = utils_open_authorized_destination(full); + else if (strcmp(dir, ".") == 0) + fd = dup(root_fd); + } else { + fd = open(full, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC); + } + free(full); + return fd; +} + +typedef struct PlanSkips { + DeleteSkipEntry* entries; + int count; +} PlanSkips; + +static bool build_plan_skips(const Config* config, const DeletePlanSession* session, + PlanSkips* out) { + out->entries = NULL; + out->count = 0; + int count = (config->delay_updates ? 1 : 0) + config->basis_count + + session->protected_prefixes->size + session->size_skipped->size; + if (count == 0) + return true; + out->entries = calloc((size_t)count, sizeof(DeleteSkipEntry)); + if (!out->entries) + return false; + int idx = 0; + if (config->delay_updates) { + out->entries[idx].prefix = DELAY_UPDATES_STAGING_DIR; + out->entries[idx].top_level_only = true; + idx++; + } + for (int i = 0; i < config->basis_count; i++) { + out->entries[idx].prefix = config->basis_dirs[i].path; + out->entries[idx].top_level_only = false; + idx++; + } + for (int i = 0; i < session->protected_prefixes->size; i++) { + out->entries[idx].prefix = (const char*)session->protected_prefixes->items[i]; + out->entries[idx].top_level_only = false; + idx++; + } + for (int i = 0; i < session->size_skipped->size; i++) { + out->entries[idx].prefix = (const char*)session->size_skipped->items[i]; + out->entries[idx].top_level_only = false; + idx++; + } + out->count = idx; + return true; +} + +static bool budget_available(const DeletePlanSession* session) { + return session->deleted < session->max_delete; +} + +static void note_skipped(DeletePlanSession* session) { + session->limit_hit = true; + session->skipped++; +} + +static void log_deleted(const char* rel) { + char* escaped = output_escape(rel, log_get_8_bit_output()); + fprintf(stderr, " Deleted: %s\n", escaped ? escaped : ""); + free(escaped); +} + +/* Append a snapshot path for --delete-delay. */ +static bool defer_add(DeletePlanSession* session, const char* rel) { + char* copy = str_dup(rel); + if (!copy) + return false; + if (!array_list_add(session->deferred, copy)) { + free(copy); + return false; + } + session->deleted++; + return true; +} + +/* Process the direct children of one directory. `keep_dirs`/`keep_files` + * (basenames) are the source entries that must be kept; NULL means every child + * is an extra (the forced path used inside a removed extra directory tree). + * `survives` reports that at least one child remains (kept, protected, or + * skipped by the budget). `force_now` removes even in --delete-delay mode + * (type conflicts must clear before the incoming data). */ +static bool process_children(int dirfd, const char* dir_rel, const ArrayList* keep_dirs, + const ArrayList* keep_files, bool at_root, bool force_now, + const PlanSkips* skips, DeletePlanSession* session, bool* survives); + +static bool process_extra_dir(int dirfd, const char* name, const char* child_rel, bool force_now, + const PlanSkips* skips, DeletePlanSession* session, bool* removed) { + *removed = false; + int childfd = openat(dirfd, name, O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC); + if (childfd < 0) { + if (errno == ENOENT) { + *removed = true; + return true; + } + return false; + } + bool survives = false; + bool ok = + process_children(childfd, child_rel, NULL, NULL, false, force_now, skips, session, &survives); + close(childfd); + if (!ok) + return false; + if (survives) + return true; + if (!budget_available(session)) { + note_skipped(session); + return true; + } + if (session->defer && !force_now) { + if (!defer_add(session, child_rel)) + return false; + *removed = true; + return true; + } + if (unlinkat(dirfd, name, AT_REMOVEDIR) == 0) { + session->deleted++; + log_deleted(child_rel); + *removed = true; + return true; + } + if (errno == ENOENT) { + *removed = true; + return true; + } + /* ENOTEMPTY/EEXIST: a protected entry the walker leaves behind survived, so + the directory stays; any other errno is a genuine failure. */ + return errno == ENOTEMPTY || errno == EEXIST; +} + +static bool process_extra_file(int dirfd, const char* name, const char* child_rel, bool force_now, + DeletePlanSession* session) { + if (!budget_available(session)) { + note_skipped(session); + return true; + } + if (session->defer && !force_now) { + return defer_add(session, child_rel); + } + if (unlinkat(dirfd, name, 0) == 0) { + session->deleted++; + log_deleted(child_rel); + } else if (errno != ENOENT) { + return false; + } + return true; +} + +static bool process_children(int dirfd, const char* dir_rel, const ArrayList* keep_dirs, + const ArrayList* keep_files, bool at_root, bool force_now, + const PlanSkips* skips, DeletePlanSession* session, bool* survives) { + *survives = false; + int scanfd = openat(dirfd, ".", O_RDONLY | O_DIRECTORY | O_NOFOLLOW | O_CLOEXEC); + if (scanfd < 0) + return false; + DIR* dir = fdopendir(scanfd); + if (!dir) { + close(scanfd); + return false; + } + bool operation_ok = true; + bool local_survives = false; + const struct dirent* entry; + while ((entry = readdir(dir)) != NULL) { + if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) + continue; + char* child_rel = + (strcmp(dir_rel, ".") == 0) ? str_dup(entry->d_name) : path_cat(dir_rel, entry->d_name); + if (!child_rel) { + operation_ok = false; + continue; + } + if (path_under_skip_prefix(child_rel, at_root, skips->entries, skips->count)) { + local_survives = true; + free(child_rel); + continue; + } + struct stat st; + if (fstatat(dirfd, entry->d_name, &st, AT_SYMLINK_NOFOLLOW) != 0) { + if (errno != ENOENT) + operation_ok = false; + free(child_rel); + continue; + } + bool is_dir = S_ISDIR(st.st_mode); + bool in_keep_dirs = is_dir && list_contains_str(keep_dirs, entry->d_name); + bool in_keep_files = !is_dir && list_contains_str(keep_files, entry->d_name); + if (in_keep_dirs) { + local_survives = true; + } else if (keep_dirs && !is_dir && list_contains_str(keep_dirs, entry->d_name)) { + /* Destination file blocks a source directory: clear it now, whatever the + delete timing, so the directory can be created. */ + if (!process_extra_file(dirfd, entry->d_name, child_rel, true, session)) + operation_ok = false; + } else if (in_keep_files) { + local_survives = true; + } else if (keep_files && is_dir && list_contains_str(keep_files, entry->d_name)) { + /* Destination directory blocks a source file: remove it now. */ + bool removed = false; + if (!process_extra_dir(dirfd, entry->d_name, child_rel, true, skips, session, &removed)) + operation_ok = false; + else if (!removed) + local_survives = true; + } else if (is_dir) { + bool removed = false; + if (!process_extra_dir(dirfd, entry->d_name, child_rel, force_now, skips, session, &removed)) + operation_ok = false; + else if (!removed) + local_survives = true; + } else { + if (!process_extra_file(dirfd, entry->d_name, child_rel, force_now, session)) + operation_ok = false; + } + free(child_rel); + } + closedir(dir); + *survives = local_survives; + return operation_ok; +} + +static bool apply_plan_dir(DeletePlanSession* session, const Config* config, const char* dir, + const ArrayList* dirs, const ArrayList* files) { + int dirfd = open_plan_dir(config, dir); + if (dirfd < 0) { + /* An absent destination directory has nothing to delete. */ + return errno == ENOENT || errno == ENOTDIR; + } + PlanSkips skips; + if (!build_plan_skips(config, session, &skips)) { + close(dirfd); + return false; + } + bool survives = false; + bool ok = process_children(dirfd, dir, dirs, files, strcmp(dir, ".") == 0, false, &skips, session, + &survives); + free(skips.entries); + close(dirfd); + if (!ok) + log_message(LOG_LEVEL_ERROR, "deletion failed while removing extraneous files"); + return ok; +} + +static bool apply_missing(DeletePlanSession* session, const Config* config) { + if (session->missing_applied) + return true; + session->missing_applied = true; + /* The server clears delete_missing_args when its --allow-delete policy is + off; never honor the client's exact-path requests then. */ + if (!config->delete_missing_args || session->missing->size == 0) + return true; + DeleteManifest manifest = { + .keeps = NULL, .protected = NULL, .missing = session->missing, .dirs = NULL}; + size_t remaining = budget_available(session) ? session->max_delete - session->deleted : 0; + size_t deleted = 0; + size_t skipped = 0; + bool limit = false; + bool ok = manifest_delete_missing_args_limited(config, &manifest, remaining, &deleted, &skipped, + &limit); + session->deleted += deleted; + session->skipped += skipped; + if (limit) + session->limit_hit = true; + return ok; +} + +int delete_plan_session_receive(DeletePlanSession* session, const Config* config, int fd) { + if (!session || !config) { + send_status(fd, STATUS_ERROR); + return -1; + } + int has_config; + if (!receive_int(fd, &has_config) || (has_config != 0 && has_config != 1)) { + send_status(fd, STATUS_ERROR); + return -1; + } + size_t bytes = 0; + if (has_config) { + if (session->config_seen || !read_section(fd, session->protected_prefixes, true, &bytes) || + !read_section(fd, session->size_skipped, true, &bytes) || + !read_section(fd, session->missing, true, &bytes)) { + send_status(fd, STATUS_ERROR); + return -1; + } + session->config_seen = true; + } + char* dir = receive_wire_str(fd); + ArrayList* dirs = array_list_create(free); + ArrayList* files = array_list_create(free); + bool parsed = dir && (strcmp(dir, ".") == 0 || valid_rel_path(dir)) && dirs && files && + read_section(fd, dirs, false, &bytes) && read_section(fd, files, false, &bytes); + if (!parsed) { + free(dir); + array_list_delete(dirs); + array_list_delete(files); + send_status(fd, STATUS_ERROR); + return -1; + } + bool enabled = config->use_delete || config->delete_missing_args; + bool ok = true; + if (!session->dry_run && enabled) { + if (!session->defer && !apply_missing(session, config)) + ok = false; + if (ok && !apply_plan_dir(session, config, dir, dirs, files)) + ok = false; + } + free(dir); + array_list_delete(dirs); + array_list_delete(files); + if (!ok) { + send_status(fd, STATUS_ERROR); + return -1; + } + if (session->limit_hit) + log_message(LOG_LEVEL_WARNING, "Deletions stopped due to the delete limit (%zu skipped)", + session->skipped); + return 0; +} + +/* Apply one snapshotted --delete-delay path (post-order: children precede their + * parent directory). */ +static bool apply_deferred_path(DeletePlanSession* session, const Config* config, const char* rel) { + (void)session; + char* full = path_cat(config->receive_root_directory, rel); + if (!full) + return false; + char* leaf = NULL; + int parent_fd = file_open_secure_parent(full, &leaf, false); + free(full); + if (parent_fd < 0) { + free(leaf); + return errno == ENOENT || errno == ENOTDIR; + } + struct stat st; + if (fstatat(parent_fd, leaf, &st, AT_SYMLINK_NOFOLLOW) != 0) { + bool absent = errno == ENOENT; + close(parent_fd); + free(leaf); + return absent; + } + int rc; + if (S_ISDIR(st.st_mode)) + rc = unlinkat(parent_fd, leaf, AT_REMOVEDIR); + else + rc = unlinkat(parent_fd, leaf, 0); + bool ok = rc == 0 || errno == ENOENT || errno == ENOTEMPTY || errno == EEXIST; + if (rc == 0) + log_deleted(rel); + close(parent_fd); + free(leaf); + return ok; +} + +DeleteCommitResult delete_plan_session_commit(DeletePlanSession* session, const Config* config) { + if (!session || !config) + return DELETE_COMMIT_ERROR; + bool ok = true; + if (session->defer) { + for (int i = 0; i < session->deferred->size && ok; i++) + ok = apply_deferred_path(session, config, (const char*)session->deferred->items[i]); + } + if (ok) + ok = apply_missing(session, config); + if (!ok) + return DELETE_COMMIT_ERROR; + if (session->limit_hit) + return DELETE_COMMIT_LIMIT_REACHED; + return DELETE_COMMIT_OK; +} diff --git a/src/shared/delete_plan.h b/src/shared/delete_plan.h new file mode 100644 index 0000000..8abede3 --- /dev/null +++ b/src/shared/delete_plan.h @@ -0,0 +1,67 @@ +#ifndef DELETE_PLAN_H +#define DELETE_PLAN_H + +#include "array_list.h" +#include "config.h" +#include "file_receive.h" +#include "protocol.h" +#include + +/* Per-directory delete plans (protocol 2.24.0). + * + * rsync's --delete-during removes a directory's extras while the generator + * processes that directory, and --delete-delay records the deletion list during + * the scan but applies it only after a fully-successful transfer. FastSync has + * no per-directory generator pass; instead the sender streams one plan per + * source directory, in directory order, and the receiver applies it when it + * arrives (during) or snapshots its extras and commits them at the end (delay). + * + * The sender side builds a plan set from the path-only pre-scan (it needs every + * directory's complete direct-child list before the first data byte of that + * directory). The receiver side is a session that carries the global protected + * prefixes (filter-excluded and size-skipped source mirrors), the + * --delete-missing-args exact deletions, the shared --max-delete budget and, + * for --delete-delay, the snapshotted extras. */ + +/* ---- Sender: plan builder ---- */ + +typedef struct DeletePlanSender DeletePlanSender; + +DeletePlanSender* delete_plan_sender_create(void); +void delete_plan_sender_destroy(DeletePlanSender* sender); +/* Record one transmitted entry. `path` is the destination-relative wire path; + * is_dir marks an explicit directory entry (--dirs, a -x mount point). */ +bool delete_plan_sender_add(DeletePlanSender* sender, const char* path, bool is_dir); +/* Drop plans for directories outside `synced_dirs` (the --files-from + * synchronization scope; pass NULL when a full recursive transfer synchronized + * every directory). The receive root is the "." sentinel. */ +void delete_plan_sender_finalize(DeletePlanSender* sender, const ArrayList* synced_dirs); +/* True when no transmitted entry was recorded (an ambiguous empty scan). */ +bool delete_plan_sender_empty(const DeletePlanSender* sender); +/* Attach the global config sections advertised on the first plan frame. */ +void delete_plan_sender_set_config(DeletePlanSender* sender, const ArrayList* protected_prefixes, + const ArrayList* size_skipped, const ArrayList* missing_args); +/* Send the root plan (even before any data, so root extras are handled like + * rsync's first generator directory). Returns -1 on I/O error. */ +int delete_plan_send_root(int fd, DeletePlanSender* sender); +/* Send the plans for every ancestor of `path` (root-first) and, when is_dir, + * for `path` itself; already-sent plans are skipped. */ +int delete_plan_send_for_path(int fd, DeletePlanSender* sender, const char* path, bool is_dir); + +/* ---- Receiver: delete session ---- */ + +typedef struct DeletePlanSession DeletePlanSession; + +DeletePlanSession* delete_plan_session_create(const Config* config); +void delete_plan_session_destroy(DeletePlanSession* session); +/* Read one STATUS_DELETE_PLAN frame (the leading status already consumed) and + * act on it. Returns 0 on success (including a dry-run/disabled no-op) and -1 + * after signalling STATUS_ERROR on a malformed frame or a deletion failure. */ +int delete_plan_session_receive(DeletePlanSession* session, const Config* config, int fd); +/* Apply the deferred snapshot (--delete-delay) and the missing-args deletions. + * Safe to call once; returns the commit outcome. */ +DeleteCommitResult delete_plan_session_commit(DeletePlanSession* session, const Config* config); +/* True once the shared --max-delete budget stopped part of a deletion. */ +bool delete_plan_session_limit_reached(const DeletePlanSession* session); + +#endif diff --git a/src/shared/file_receive.c b/src/shared/file_receive.c index dfd554a..07aef4f 100644 --- a/src/shared/file_receive.c +++ b/src/shared/file_receive.c @@ -3472,6 +3472,21 @@ bool manifest_delete_missing_args(const Config* config, DeleteManifest* manifest return delete_missing_args_budgeted(config, manifest, &budget); } +bool manifest_delete_missing_args_limited(const Config* config, DeleteManifest* manifest, + size_t max_delete, size_t* deleted, size_t* skipped, + bool* limit_hit) { + DeleteBudgetState budget = { + .max_delete = max_delete, .deleted = 0, .skipped = 0, .limit_hit = false}; + bool ok = delete_missing_args_budgeted(config, manifest, &budget); + if (deleted) + *deleted = budget.deleted; + if (skipped) + *skipped = budget.skipped; + if (limit_hit) + *limit_hit = budget.limit_hit; + return ok; +} + /* Commit every deletion family the manifest carries. The --delete-missing-args exact-path deletions run FIRST: they are explicit user requests and must not be blocked by the extras walker's filter-exclusion protection (a protected diff --git a/src/shared/file_receive.h b/src/shared/file_receive.h index 4f2b281..5cddd99 100644 --- a/src/shared/file_receive.h +++ b/src/shared/file_receive.h @@ -118,6 +118,14 @@ bool manifest_delete_extras(const Config* config, DeleteManifest* manifest); confinement or I/O error (the run then fails); tolerated per-path cases are reported and skipped. */ bool manifest_delete_missing_args(const Config* config, DeleteManifest* manifest); +/* Budgeted form of manifest_delete_missing_args for the per-directory delete + session: each removed mirror draws from `max_delete` (SIZE_MAX = unlimited) + and the tallies are accumulated into `*deleted`/`*skipped`. `*limit_hit` is set + when the budget stopped the pass with entries left over. Returns false only + on a genuine deletion error. */ +bool manifest_delete_missing_args_limited(const Config* config, DeleteManifest* manifest, + size_t max_delete, size_t* deleted, size_t* skipped, + bool* limit_hit); /* 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). */ diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 8c752aa..5e31e64 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -36,6 +36,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que context->scan_had_io_error = false; context->remove_source_files = NULL; context->early_delete = false; + context->delete_plans = NULL; context->scan_stopped_early = false; context->total_files = 0; context->progress_bytes = 0; @@ -187,6 +188,8 @@ void pipeline_context_sender_destroy(PipelineContextSender* context) { if (context->manifest) { array_list_delete(context->manifest); } + if (context->delete_plans) + delete_plan_sender_destroy(context->delete_plans); if (context->excluded_paths) array_list_delete(context->excluded_paths); if (context->size_skipped_paths) diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index f8b475f..39d2c38 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -7,6 +7,7 @@ #include "array_list.h" #include "chunk.h" #include "config.h" +#include "delete_plan.h" #include "file.h" #include "protocol.h" #include "queue.h" @@ -66,11 +67,16 @@ typedef struct { --ignore-errors kept the run going. */ bool scan_had_io_error; ArrayList* remove_source_files; - /* True when --delete-before/--delete-during require the keep-set manifest to - be transmitted before any file data: context->manifest is then prebuilt by - a path-only pre-scan on the calling thread and the pipeline scanner must - not append to it. Set once before the worker threads start. */ + /* True when --delete-before requires the whole-tree keep-set manifest to be + transmitted before any file data: context->manifest is then prebuilt by a + path-only pre-scan on the calling thread and the pipeline scanner must not + append to it. Set once before the worker threads start. */ bool early_delete; + /* Non-NULL for --delete-during/--delete-delay: the per-directory plan set + prebuilt by the path-only pre-scan on the calling thread. The sender + thread transmits the root plan before any data and the remaining plans + alongside the chunks. Set once before the worker threads start. */ + DeletePlanSender* delete_plans; mtx_t mutex_progress; int total_files; unsigned long long progress_bytes; diff --git a/src/shared/protocol.h b/src/shared/protocol.h index 95d95d4..4dc91a0 100644 --- a/src/shared/protocol.h +++ b/src/shared/protocol.h @@ -181,7 +181,21 @@ enum NET_STATUS { * (new vs modified, and which of size/time/perms/owner/group differ) without * changing the transfer decision itself. Appended after * STATUS_DELETE_LIMIT so no existing status is renumbered. */ - STATUS_DEST_INFO + STATUS_DEST_INFO, + /* Per-directory delete plan (protocol 2.24.0). The sender of a + * --delete-during/--delete-delay transfer streams one frame per source + * directory in directory order instead of a single whole-tree keep-set + * manifest. The receiver applies the plan when it arrives + * (--delete-during removes that directory's extras immediately) or records + * the extras and applies them only after the whole transfer succeeded + * (--delete-delay). Payload: an int32 has_config flag (1 on the first plan + * of the run, 0 afterwards); when set, the three global config sections + * (protected-prefix count+paths, size-skipped count+paths, missing-args + * count+paths); then the destination-relative directory path wire string + * ("." for the receive root); then the child-directory count + names and the + * child-file count + names that must be kept. Appended after + * STATUS_DEST_INFO so no existing status is renumbered. */ + STATUS_DELETE_PLAN }; void io_set_fds(int read_fd, int write_fd); diff --git a/src/shared/utils.c b/src/shared/utils.c index 00704d0..8c1a20d 100644 --- a/src/shared/utils.c +++ b/src/shared/utils.c @@ -60,7 +60,7 @@ bool path_is_within_root(const char* root, const char* path) { * two differ in create-vs-no-create, in what path component they stop at, and * in the extra receiver policies they apply, so they are intentionally kept * separate. Both rely on the shared lexical path_is_within_root check. */ -static int open_authorized_destination(const char* dest_root) { +int utils_open_authorized_destination(const char* dest_root) { int root_fd = utils_get_authorized_root_fd(); const char* root_path = utils_get_authorized_root_path(); if (root_fd < 0 || !root_path || !dest_root || !path_is_within_root(root_path, dest_root)) @@ -752,7 +752,7 @@ DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* m int root_fd = utils_get_authorized_root_fd(); if (root_fd >= 0) { if (utils_get_authorized_root_path()) - rootfd = open_authorized_destination(dest_root); + rootfd = utils_open_authorized_destination(dest_root); else if (dest_root == NULL) rootfd = dup(root_fd); else diff --git a/src/shared/utils.h b/src/shared/utils.h index 0cca144..8321fa3 100644 --- a/src/shared/utils.h +++ b/src/shared/utils.h @@ -139,6 +139,11 @@ DeleteWalkResult delete_extras_limited(const char* dest_root, const ArrayList* m const DeleteSkipEntry* skips, int skip_count, size_t* deleted_out, size_t* skipped_out); bool delete_extras(const char* dest_root, const ArrayList* manifest); +/* Open the existing destination directory at `dest_root`, confined to the + authorized root with an O_NOFOLLOW component walk (the same confinement the + deletion walker uses for its root). Returns a new fd the caller owns, or -1 + on error (including a destination that does not exist). */ +int utils_open_authorized_destination(const char* dest_root); bool utils_set_authorized_root(int fd, const char* canonical_path); /* The fd-only compatibility form is fail-closed for path-based operations; * callers should use utils_set_authorized_root with the canonical identity. */ diff --git a/tests/integration/test_delete_timing_parity.py b/tests/integration/test_delete_timing_parity.py new file mode 100644 index 0000000..1991f87 --- /dev/null +++ b/tests/integration/test_delete_timing_parity.py @@ -0,0 +1,332 @@ +"""Differential + regression coverage for rsync's delete timing. + +``--delete-during``/``--delete-delay`` stream a per-directory delete plan instead +of one whole-tree manifest, so the timing is observable: + + * ``--delete-during`` removes a directory's extras as it processes that + directory (so an interrupted transfer has already removed the extras of the + directories it reached); + * ``--delete-delay`` snapshots those extras while scanning and commits the + removals only after a fully-successful transfer (so an extra created in the + destination after its directory's plan survives, and a failed transfer + removes nothing); + * ``--delete-after`` re-scans the destination at the end (so that same + late-created extra is removed). + +The final-state tests compare against real ``rsync 3.4.1`` where a deterministic +comparison exists; the timing tests use a byte-slicing proxy to force a +mid-transfer failure or to create a destination entry while the transfer is in +flight. +""" +import os +import select +import shutil +import socket +import struct +import subprocess +import sys +import threading +import time + +import pytest + +sys.path.insert(0, os.path.dirname(__file__)) +from common import ( # noqa: E402 + BUILD_DIR, + TEST_DATA_DIR, + ServerManager, + clean_dir, + get_dest_received_dir, + run_client, +) + +# Every test here is deterministic (the proxy throttles until the delete-plan +# frames are processed), so the PR gate runs the whole module. +pytestmark = pytest.mark.ci + +RSYNC = shutil.which("rsync") +requires_rsync = pytest.mark.skipif(RSYNC is None, reason="rsync 3.4.1 not installed") + +BIG_BYTES = 8 * 1024 * 1024 +# Forward/cut this far into the stream: past the (small) delete-plan frames and +# well into the big payload, so the receiver has already processed the plan. +MID_TRANSFER_BYTES = 256 * 1024 +# Throttle the proxy so the receiver keeps up with the (fast) client and the +# plan frames are provably processed before the hook/cut offset is reached. +PROXY_THROTTLE = 0.001 + + +def _write(path, content): + os.makedirs(os.path.dirname(path), exist_ok=True) + with open(path, "wb") as fh: + fh.write(content) + + +def _seed_pair(tag, big=False): + """Create a source tree and a destination mirror seeded with extras. + + The tree is a single directory ``d`` containing the transferred files plus, + in the destination, an extra ``d/old_extra``. + """ + source = os.path.join(TEST_DATA_DIR, f"dtp_{tag}_src") + dest = os.path.join(TEST_DATA_DIR, f"dtp_{tag}_dst") + clean_dir(source) + clean_dir(dest) + _write(os.path.join(source, "d", "keep.txt"), b"kept payload\n") + if big: + _write(os.path.join(source, "d", "big.bin"), b"B" * BIG_BYTES) + received = get_dest_received_dir(dest, source) + os.makedirs(os.path.join(received, "d"), exist_ok=True) + _write(os.path.join(received, "d", "old_extra"), b"stale extra\n") + return source, dest, received + + +def _tree(root): + """Sorted relative paths of every entry below root (files and dirs).""" + out = [] + for dirpath, dirs, files in os.walk(root): + for name in dirs: + out.append(os.path.relpath(os.path.join(dirpath, name), root)) + for name in files: + out.append(os.path.relpath(os.path.join(dirpath, name), root)) + return sorted(out) + + +def _rsync(args): + env = dict(os.environ, LC_ALL="C") + return subprocess.run([RSYNC] + args, capture_output=True, text=True, env=env, timeout=120) + + +class _SlicingProxy: + """Forward the client stream to a server, optionally cutting it or invoking a + hook after a byte threshold. ``forward_limit`` mode resets both ends after + that many client bytes (a mid-transfer failure). ``hook`` mode calls the + hook once and keeps forwarding to completion.""" + + def __init__(self, target_port, forward_limit=None, hook=None, hook_after=0, + throttle=0.0): + self.target = ("127.0.0.1", target_port) + self.forward_limit = forward_limit + self.hook = hook + self.hook_after = hook_after + self.throttle = throttle + self.hook_called = threading.Event() + self.listener = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + self.listener.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + self.listener.bind(("127.0.0.1", 0)) + self.listener.listen(1) + self.listener.settimeout(20) + self.port = self.listener.getsockname()[1] + self._thread = threading.Thread(target=self._serve, daemon=True) + self._thread.start() + + def _serve(self): + try: + client, _ = self.listener.accept() + except OSError: + return + try: + backend = socket.create_connection(self.target, timeout=10) + except OSError: + client.close() + return + client.settimeout(20) + backend.settimeout(20) + forwarded = 0 + socks = [client, backend] + try: + while socks: + ready, _, _ = select.select(socks, [], [], 20) + if not ready: + break + for sock in ready: + data = sock.recv(65536) + if not data: + socks.remove(sock) + peer = backend if sock is client else client + try: + peer.shutdown(socket.SHUT_WR) + except OSError: + pass + continue + if sock is client: + if self.forward_limit is not None: + room = self.forward_limit - forwarded + if room <= 0: + socks = [] + break + data = data[:room] + backend.sendall(data) + forwarded += len(data) + if (self.hook is not None and not self.hook_called.is_set() + and forwarded >= self.hook_after): + # Give the receiver time to process the (tiny) plan + # frames that precede this offset before the hook + # mutates the destination. + if self.throttle > 0: + time.sleep(0.2) + self.hook() + self.hook_called.set() + if self.forward_limit is not None and forwarded >= self.forward_limit: + socks = [] + break + if self.throttle > 0: + time.sleep(self.throttle) + else: + client.sendall(data) + except OSError: + pass + for sock in (client, backend): + try: + sock.setsockopt(socket.SOL_SOCKET, socket.SO_LINGER, struct.pack("ii", 1, 0)) + except OSError: + pass + try: + sock.close() + except OSError: + pass + try: + self.listener.close() + except OSError: + pass + + def finish(self): + self._thread.join(30) + try: + self.listener.close() + except OSError: + pass + + +class TestDeleteTimingFinalStateParity: + """On a successful transfer the per-directory timings match rsync's result.""" + + def _run_fastsync(self, tag, timing): + source, dest, received = _seed_pair(tag) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, flags=[timing], port=server.port) + return result, received + + @pytest.mark.parametrize("timing", ["--delete-during", "--delete-delay"]) + @requires_rsync + def test_success_final_state_matches_rsync(self, timing): + # Build the rsync fixture from the same seed so both sides start equal. + source, dest, received = _seed_pair("parity_rsync") + source2 = source + rsync_dst = os.path.join(TEST_DATA_DIR, "dtp_parity_rsync_dst") + clean_dir(rsync_dst) + # rsync mirrors src/ into dst/; seed the same extra. + _write(os.path.join(rsync_dst, "d", "old_extra"), b"stale extra\n") + + rsync_result = _rsync(["-a", timing, source2 + "/", rsync_dst + "/"]) + assert rsync_result.returncode == 0, rsync_result.stderr + rsync_tree = _tree(rsync_dst) + + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, flags=[timing], port=server.port) + assert result.returncode == 0, (result.stderr or result.stdout)[:300] + fastsync_tree = _tree(received) + assert fastsync_tree == rsync_tree, ( + f"{timing}: fastsync tree {fastsync_tree} != rsync tree {rsync_tree}" + ) + + +class TestDeleteTimingTypeConflictParity: + """A destination entry whose type differs from the source is replaced, in + both per-directory timings and in both directions, exactly like rsync.""" + + @pytest.mark.parametrize("timing", ["--delete-during", "--delete-delay"]) + @requires_rsync + def test_type_conflicts_match_rsync(self, timing): + source = os.path.join(TEST_DATA_DIR, "dtc_src") + clean_dir(source) + _write(os.path.join(source, "foo"), b"now a file\n") + _write(os.path.join(source, "bar", "inner.txt"), b"now a dir\n") + + def seed_dest(root): + clean_dir(root) + _write(os.path.join(root, "foo", "inner.txt"), b"was a dir\n") + _write(os.path.join(root, "bar"), b"was a file\n") + + rsync_dst = os.path.join(TEST_DATA_DIR, "dtc_rsync_dst") + seed_dest(rsync_dst) + rsync_result = _rsync(["-a", timing, source + "/", rsync_dst + "/"]) + assert rsync_result.returncode == 0, rsync_result.stderr + rsync_tree = _tree(rsync_dst) + + dest = os.path.join(TEST_DATA_DIR, "dtc_dst") + clean_dir(dest) + received = get_dest_received_dir(dest, source) + seed_dest(received) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, flags=[timing], port=server.port) + assert result.returncode == 0, (result.stderr or result.stdout)[:300] + assert _tree(received) == rsync_tree, ( + f"{timing}: fastsync tree {_tree(received)} != rsync tree {rsync_tree}" + ) + + +class TestDeleteTimingFailure: + """A mid-transfer failure distinguishes during from delay.""" + + @pytest.mark.parametrize("mt", [False, True]) + def test_during_removes_delay_preserves_on_failure(self, mt): + source, dest, received = _seed_pair("failure", big=True) + extra = os.path.join(received, "d", "old_extra") + assert os.path.exists(extra) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + for timing, expect_removed in (("--delete-during", True), + ("--delete-delay", False)): + # Re-seed the extra before each run. + _write(extra, b"stale extra\n") + proxy = _SlicingProxy(server.port, forward_limit=MID_TRANSFER_BYTES, throttle=PROXY_THROTTLE) + flags = [timing] + (["--threads"] if mt else []) + result, _ = run_client(source, dest, flags=flags, port=proxy.port) + proxy.finish() + assert result.returncode != 0, f"{timing}: truncated transfer succeeded" + present = os.path.exists(extra) + assert present != expect_removed, ( + f"{timing} (mt={mt}): extra present={present}, expected " + f"removed={expect_removed}" + ) + + +class TestDeleteDelayVsAfterSnapshot: + """A destination entry created after its directory's scan survives under + --delete-delay but is removed by --delete-after's fresh end scan.""" + + @pytest.mark.parametrize("mt", [False, True]) + def test_late_created_extra_survives_delay_not_after(self, mt): + source, dest, received = _seed_pair("latecreate", big=True) + old_extra = os.path.join(received, "d", "old_extra") + new_extra = os.path.join(received, "d", "new_extra") + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + for timing, new_survives in (("--delete-delay", True), + ("--delete-after", False)): + _write(old_extra, b"stale extra\n") + if os.path.exists(new_extra): + os.unlink(new_extra) + + def hook(): + # Runs on the proxy thread while the big file is in flight, + # after the directory's plan (delay) has been processed. + _write(new_extra, b"created mid-transfer\n") + + proxy = _SlicingProxy(server.port, hook=hook, hook_after=MID_TRANSFER_BYTES, throttle=PROXY_THROTTLE) + flags = [timing] + (["--threads"] if mt else []) + result, _ = run_client(source, dest, flags=flags, port=proxy.port) + proxy.finish() + assert result.returncode == 0, ( + f"{timing}: {(result.stderr or result.stdout)[:300]}" + ) + assert proxy.hook_called.is_set(), f"{timing}: hook never fired" + assert not os.path.exists(old_extra), f"{timing}: old extra survived" + assert os.path.exists(new_extra) == new_survives, ( + f"{timing} (mt={mt}): new_extra present=" + f"{os.path.exists(new_extra)}, expected survives={new_survives}" + ) diff --git a/tests/integration/test_fault_injection.py b/tests/integration/test_fault_injection.py index 7c7127f..fd93651 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.23.0" +PROTOCOL_VERSION = b"2.24.0" STATUS_MANIFEST = 5 STATUS_OK = 0 diff --git a/tests/integration/test_features.py b/tests/integration/test_features.py index 3784601..de0c4f0 100644 --- a/tests/integration/test_features.py +++ b/tests/integration/test_features.py @@ -3755,15 +3755,15 @@ class TestDeleteTiming: assert _read_file(os.path.join(received, "sub", "deep.txt")) == b"deeply nested file\n", \ f"{flag}: nested file was not written after the early deletion" - @pytest.mark.parametrize("flag", ["--delete", "--delete-after", "--delete-delay"]) + @pytest.mark.parametrize("flag", ["--delete", "--delete-after"]) @pytest.mark.parametrize("mt", [False, True]) def test_late_flags_commit_only_after_success(self, flag, mt): - """Plain --delete/--delete-after/--delete-delay defer deletion until the - whole transfer succeeds: a mid-transfer write failure must leave every - extra in place (commit-style safety). The -m receiver must also keep - the extras: the deferred keep-set is committed by the server only after - the disk-writer thread has finished, and a failing writer means the - manifest is freed, never applied.""" + """Plain --delete/--delete-after defer deletion until the whole transfer + succeeds: a mid-transfer write failure must leave every extra in place + (commit-style safety). The -m receiver must also keep the extras: the + deferred keep-set is committed by the server only after the disk-writer + thread has finished, and a failing writer means the manifest is freed, + never applied.""" source = self._seed("late") dest = os.path.join(TEST_DATA_DIR, "deltiming_late_dst") clean_dir(dest) @@ -3789,10 +3789,42 @@ class TestDeleteTiming: assert os.path.isfile(blocker), \ f"{flag} (mt={mt}) deleted the blocker although the transfer failed" - def test_early_flag_respected_when_server_refuses_delete(self, shared_server): - """With an --allow-delete-less server the client's early timing still - completes (no deadlock on the pre-delete ack) and simply never deletes, - exactly like the plain server policy.""" + @pytest.mark.parametrize("mt", [False, True]) + def test_delete_delay_clears_type_conflict_like_rsync(self, mt): + """rsync clears a destination file that blocks a source directory even + when the deletion itself is deferred (--delete-delay); the type conflict + is resolved immediately so the nested write succeeds. The transfer must + therefore succeed and the unrelated extra must still be removed.""" + source = self._seed("delayconflict") + dest = os.path.join(TEST_DATA_DIR, "deltiming_delayconflict_dst") + clean_dir(dest) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, port=server.port) + assert result.returncode == 0, f"seed sync failed: {result.stderr[:200]}" + received = get_dest_received_dir(dest, source) + extra = os.path.join(received, "extra.txt") + with open(extra, "wb") as fh: + fh.write(b"extra file") + blocker = os.path.join(received, "sub") + shutil.rmtree(blocker) + with open(blocker, "wb") as fh: + fh.write(b"blocks the nested destination directory") + + flags = ["--delete-delay"] + (["--threads"] if mt else []) + result, _ = run_client(source, dest, flags=flags, port=server.port) + assert result.returncode == 0, \ + f"--delete-delay (mt={mt}) did not clear the type conflict: " \ + f"{(result.stderr or result.stdout)[:300]}" + assert os.path.isdir(blocker), "blocker file was not replaced by the source directory" + assert _read_file(os.path.join(received, "sub", "deep.txt")) == b"deeply nested file\n" + assert not os.path.exists(extra), "--delete-delay did not remove the extra" + + @pytest.mark.parametrize("flag", ["--delete-before", "--delete-during", "--delete-delay"]) + def test_early_flag_respected_when_server_refuses_delete(self, flag, shared_server): + """With an --allow-delete-less server the client's timing still completes + (no deadlock on the pre-delete ack, no per-directory deletion) and simply + never deletes, exactly like the plain server policy.""" source = self._seed("refused") dest = os.path.join(TEST_DATA_DIR, "deltiming_refused_dst") clean_dir(dest) @@ -3802,9 +3834,9 @@ class TestDeleteTiming: extra = os.path.join(received, "extra.txt") with open(extra, "wb") as fh: fh.write(b"extra file") - result, _ = run_client(source, dest, flags=["--delete-before"], port=shared_server.port) + result, _ = run_client(source, dest, flags=[flag], port=shared_server.port) assert result.returncode == 0, \ - f"--delete-before against a refuse-delete server failed: {result.stderr[:300]}" + f"{flag} against a refuse-delete server failed: {result.stderr[:300]}" assert os.path.exists(extra), "unauthorized delete removed an extra file" @@ -3899,6 +3931,31 @@ class TestDeleteScope: finally: server.stop() + @pytest.mark.parametrize("mt", [False, True]) + @pytest.mark.parametrize("timing", ["--delete-during", "--delete-delay"]) + @pytest.mark.ci + def test_files_from_per_dir_timing_confined_to_listed_dirs(self, mt, timing): + """The per-directory timings honor the same --files-from scope: an extra + inside a listed directory is removed, while unlisted siblings and the + receive-root extra survive.""" + source, dest, received, server = self._seed(f"pd_{timing.strip('-')}_{mt}") + try: + listed = _write_rel_list(b"listed.txt\nsub/\n") + flags = ["--files-from", listed, timing] + (["--threads"] if mt else []) + result, _ = run_client(source, dest, flags=flags, port=server.port) + assert result.returncode == 0, f"{timing} delete failed: {result.stderr[:300]}" + assert not os.path.exists(os.path.join(received, "sub", "extra.txt")), \ + f"{timing} did not delete the in-scope extra" + assert os.path.isfile(os.path.join(received, "sub", "x.txt")) + assert os.path.exists(os.path.join(received, "unlisted.txt")), \ + f"{timing} deleted an unlisted path (data loss)" + assert os.path.exists(os.path.join(received, "other", "c.txt")), \ + f"{timing} deleted an unlisted sibling directory (data loss)" + assert os.path.exists(os.path.join(received, "rootextra.txt")), \ + f"{timing} deleted the receive-root extra (data loss)" + finally: + server.stop() + class TestDeleteExtraneousSymlinks: """#290 (3): --delete unlinks extraneous destination symlinks (never follows @@ -3986,7 +4043,8 @@ class TestDeletePolicy: @pytest.mark.parametrize("mt", [False, True]) @pytest.mark.parametrize("timing", - ["--delete", "--delete-before", "--delete-after", "--delete-delay"]) + ["--delete", "--delete-before", "--delete-after", "--delete-delay", + "--delete-during"]) def test_delete_protects_excluded_by_default_and_delete_excluded_removes(self, mt, timing): """rsync parity: with a --delete timing the destination mirror path whose source was excluded survives (protected by default); --delete-excluded @@ -4070,7 +4128,7 @@ class TestDeletePolicy: "--delete-excluded did not remove the excluded dir subtree" @pytest.mark.parametrize("mt", [False, True]) - @pytest.mark.parametrize("timing", ["--delete", "--delete-before"]) + @pytest.mark.parametrize("timing", ["--delete", "--delete-before", "--delete-during"]) def test_max_delete_exceeded_deletes_up_to_cap_and_exits_25(self, mt, timing): """rsync parity: --max-delete=N deletes up to N extras, skips the rest and still succeeds as a transfer, exiting 25 with a diagnostic.""" diff --git a/tests/integration/test_preflight.py b/tests/integration/test_preflight.py index d6b601a..547d8f0 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.23.0 (the current PROTOCOL_VERSION) is accepted and the + """--protocol=2.24.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.23.0"], + result, _ = run_client(source, dest, flags=["--protocol=2.24.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 5503814..01bc272 100644 --- a/tests/test_client_cli.c +++ b/tests/test_client_cli.c @@ -317,7 +317,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.23.0"}; + "--dest-dir", "/dst", "--protocol=2.24.0"}; int positional_args[2]; int positional_count = 0; EXPECT_EQ_INT(parse_args(cfg, 6, argv_equals, positional_args, &positional_count), 0); @@ -327,7 +327,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.23.0"}; + "/dst", "--protocol", "2.24.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 a5fe1c7..b506bad 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -852,13 +852,15 @@ static void test_config_delete_timing_early_helper() { cfg->use_delete = true; cfg->delete_before = true; EXPECT_TRUE(config_delete_timing_early(cfg)); + EXPECT_FALSE(config_delete_timing_per_dir(cfg)); EXPECT_TRUE(config_has_valid_delete_timing(cfg)); config_delete(cfg); cfg = config_create(); cfg->use_delete = true; cfg->delete_during = true; - EXPECT_TRUE(config_delete_timing_early(cfg)); + EXPECT_FALSE(config_delete_timing_early(cfg)); + EXPECT_TRUE(config_delete_timing_per_dir(cfg)); EXPECT_TRUE(config_has_valid_delete_timing(cfg)); config_delete(cfg); @@ -866,6 +868,7 @@ static void test_config_delete_timing_early_helper() { cfg->use_delete = true; cfg->delete_delay = true; EXPECT_FALSE(config_delete_timing_early(cfg)); + EXPECT_TRUE(config_delete_timing_per_dir(cfg)); EXPECT_TRUE(config_has_valid_delete_timing(cfg)); config_delete(cfg); @@ -873,6 +876,7 @@ static void test_config_delete_timing_early_helper() { cfg->use_delete = true; cfg->delete_after = true; EXPECT_FALSE(config_delete_timing_early(cfg)); + EXPECT_FALSE(config_delete_timing_per_dir(cfg)); EXPECT_TRUE(config_has_valid_delete_timing(cfg)); config_delete(cfg); @@ -2771,14 +2775,14 @@ static void golden_config_populate(Config* c) { c->copy_as_gid = 222; } -/* The pinned golden frame (protocol 2.23.0). The values below are the only +/* The pinned golden frame (protocol 2.24.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.23.0 - * rsync-parity wave changes the config-frame layout (map-entry range + TO name, - * one report_dest_info bool, and other wire changes landing in this version); - * the byte-exact values are recomputed for the merged layout. */ + * them ONLY with a PROTOCOL_VERSION bump and a documented reason. The 2.24.0 + * per-directory delete-plan wave changes only the version string in the config + * frame (the frame layout itself is unchanged from 2.23.0); the byte-exact hash + * is recomputed for the new version bytes. */ #define GOLDEN_WIRE_LEN 697 -#define GOLDEN_WIRE_HASH 7835017034643051109ULL +#define GOLDEN_WIRE_HASH 13736055061412501670ULL static unsigned long long fnv1a_64(const unsigned char* buf, size_t len) { unsigned long long h = 1469598103934665603ULL; @@ -2860,7 +2864,7 @@ static unsigned long long capture_wire_hash(const Config* cfg, size_t* out_len) return h; } -/* Byte-for-byte wire compatibility guard (protocol 2.23.0). The expected hash +/* Byte-for-byte wire compatibility guard (protocol 2.24.0). The expected hash * pins the pre-X-macro byte stream; the refactor MUST NOT change it. */ static void test_config_wire_golden() { if (is_running_under_valgrind()) diff --git a/tests/test_receiver_timeout.c b/tests/test_receiver_timeout.c index 8066b5c..a291be0 100644 --- a/tests/test_receiver_timeout.c +++ b/tests/test_receiver_timeout.c @@ -65,7 +65,7 @@ static void test_receiver_aborts_idle_keepalive() { ssize_t wrote = write(sv[0], &keepalive, sizeof(keepalive)); int result = -2; if (wrote == (ssize_t)sizeof(keepalive)) - result = receiver_process_pending(config, sv[1], &sink, NULL); + result = receiver_process_pending(config, sv[1], &sink, NULL, NULL); Status reply = STATUS_OK; ssize_t got = -1; if (result == -1) diff --git a/tests/test_server.c b/tests/test_server.c index aa680a4..f4f7fc8 100644 --- a/tests/test_server.c +++ b/tests/test_server.c @@ -566,7 +566,7 @@ static Config* make_late_delete_config(const char* root) { static int run_pending_receiver(Config* cfg, int fd, DeleteManifest** pending) { ReceiverSink sink = {0}; - return receiver_process_pending(cfg, fd, &sink, pending); + return receiver_process_pending(cfg, fd, &sink, pending, NULL); } static void test_late_manifest_abort_frees_keepset() { @@ -797,7 +797,7 @@ static void test_receiver_pending_commits_missing_args() { /* NULL pending: the single-threaded commit path deletes at FINISHED. The sink sends the terminal STATUS_OK success frame. */ ReceiverSink sink = {.send_success = true}; - EXPECT_EQ_INT(receiver_process_pending(cfg, p[0], &sink, NULL), 0); + EXPECT_EQ_INT(receiver_process_pending(cfg, p[0], &sink, NULL, NULL), 0); Status ack; EXPECT_TRUE(receive_status(p[1], &ack)); EXPECT_EQ_INT(ack, STATUS_OK);