diff --git a/src/client/client_send.c b/src/client/client_send.c index 7a0046c..f1d4516 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -923,15 +923,27 @@ static bool receive_stats_record(int fd, ReceiverStats* stats, ArrayList* would_ int count = 0; if (!receive_int(fd, &count) || count < 0 || count > MAX_MANIFEST_ENTRIES) return false; + /* Mirror the delete-plan parser: every retained path must be a valid + destination-relative path, and the whole list shares one MAX_MANIFEST_BYTES + budget so a hostile peer cannot make the client retain unbounded memory. */ + size_t bytes = 0; for (int i = 0; i < count; i++) { char* path = receive_wire_str(fd); if (!path) return false; - if (would_delete) { - char* copy = str_dup(path); + if (path[0] == '\0' || path[0] == '/' || has_path_traversal(path)) { free(path); - if (!copy || !array_list_add(would_delete, copy)) { - free(copy); + return false; + } + if (would_delete) { + size_t entry_size = strlen(path) + sizeof(char*) + 16; + if (entry_size > MAX_MANIFEST_BYTES - bytes) { + free(path); + return false; + } + bytes += entry_size; + if (!array_list_add(would_delete, path)) { + free(path); return false; } } else { @@ -1455,6 +1467,19 @@ static bool scan_paths_only(const Config* config, const ScannerOptions* options, } chunk_destroy(chunk); } + if (ok) { + /* Keep every traversed source directory, including empty ones, so a plan + no longer removes the destination directory itself. Their own plans are + emitted after the data stream (no file frame triggers them). */ + if (plans && options->plan_dirs) { + for (int i = 0; i < options->plan_dirs->size; i++) { + if (!delete_plan_sender_add(plans, (const char*)options->plan_dirs->items[i], true)) { + ok = false; + break; + } + } + } + } if (ok && directory_scanner_failed(scanner)) ok = false; if (io_error_out) @@ -1929,7 +1954,12 @@ static int send_dry_run_remote(Config* config) { event.path = path; char* line = change_render_format(config->out_format, config, &event); if (line) { - printf("%s\n", line); + /* Escape the whole rendered line, exactly like change_emit() does + for a real transfer, so a control byte in the peer-supplied path + cannot forge output. */ + char* escaped = output_escape(line, config->eight_bit_output); + printf("%s\n", escaped ? escaped : line); + free(escaped); free(line); } } else { @@ -2474,6 +2504,13 @@ static int send_chunks_multithreaded(void* pipeline_context) { NULL) != 0) goto send_fail; } + /* Emit the plans for source directories the data stream never triggered + (empty directories): their extras are still cleared while the directory + itself is kept. */ + if (!context->scan_stopped_early && context->delete_plans && context->plan_dirs && + delete_plan_send_remaining(client->file_descriptor, context->delete_plans, + context->plan_dirs) != 0) + goto send_fail; /* P7 Wave D: transmit the captured directory times last. The scanner thread (and all parallel workers) has been joined before scanner_done was set, so the list is complete and race-free; on an early stop the list may be @@ -2808,6 +2845,8 @@ int send_files(Config* config) { /* Size-pruned prefixes (always protected) and synchronized directories. */ ArrayList* size_skipped = NULL; ArrayList* synced_dirs = NULL; + /* Traversed source directories for the per-directory delete keep set. */ + ArrayList* plan_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; @@ -2908,13 +2947,16 @@ int send_files(Config* config) { 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) + plan_dirs = array_list_create(free); + if (!plan_sender || !plan_dirs) goto send_fail; + prepared.options.plan_dirs = plan_dirs; bool prescan_ok = scan_paths_only(config, &prepared.options, NULL, plan_sender, &had_scan_io); bool plans_ok = false; if (prescan_ok) { const char* walk_root = delete_plan_walk_root(config, synced_dirs); - const ArrayList* scope = config->files_from_set ? synced_dirs : (walk_root ? synced_dirs : NULL); + const ArrayList* scope = + config->files_from_set ? synced_dirs : (walk_root ? synced_dirs : NULL); delete_plan_sender_finalize(plan_sender, scope, walk_root); delete_plan_sender_set_config(plan_sender, excluded, size_skipped, missing_args); if (had_scan_io && delete_plan_sender_empty(plan_sender)) { @@ -2929,6 +2971,7 @@ int send_files(Config* config) { prepared.options.excluded_paths = NULL; prepared.options.size_skipped_paths = NULL; prepared.options.synced_dirs = NULL; + prepared.options.plan_dirs = NULL; if (!prescan_ok || !plans_ok) goto send_fail; } else if (config->use_delete) { @@ -3085,6 +3128,12 @@ int send_files(Config* config) { } } } + /* Emit the plans for any source directories the data stream never triggered + (an empty directory has no file frame). Sending them now still clears that + directory's destination extras while keeping the directory itself. */ + if (!scan_stopped_early && plan_sender && plan_dirs && + delete_plan_send_remaining(client->file_descriptor, plan_sender, plan_dirs) != 0) + goto send_fail; /* P7 Wave D: every directory has now been traversed (or the scan stopped early), so transmit the captured directory times last. The receiver defers applying them until after its own deletion/publication phase. */ @@ -3127,6 +3176,8 @@ send_fail: array_list_delete(size_skipped); if (synced_dirs) array_list_delete(synced_dirs); + if (plan_dirs) + array_list_delete(plan_dirs); if (missing_args) array_list_delete(missing_args); if (remove_sources) @@ -3268,7 +3319,10 @@ int send_files_multithreaded(Config** config_ptr) { } if (per_dir) { context->delete_plans = delete_plan_sender_create(); - prepared_ok = prepared_ok && context->delete_plans != NULL; + context->plan_dirs = array_list_create(free); + prepared_ok = prepared_ok && context->delete_plans != NULL && context->plan_dirs != NULL; + if (prepared_ok) + prepared.options.plan_dirs = context->plan_dirs; } else { context->manifest = array_list_create(free); prepared_ok = prepared_ok && context->manifest != NULL; @@ -3279,8 +3333,8 @@ int send_files_multithreaded(Config** config_ptr) { prepared_scanner_destroy(&prepared); if (per_dir && prebuilt) { const char* walk_root = delete_plan_walk_root(config, context->synced_dirs); - const ArrayList* scope = - config->files_from_set ? context->synced_dirs : (walk_root ? context->synced_dirs : NULL); + const ArrayList* scope = config->files_from_set ? context->synced_dirs + : (walk_root ? context->synced_dirs : NULL); delete_plan_sender_finalize(context->delete_plans, scope, walk_root); delete_plan_sender_set_config(context->delete_plans, context->excluded_paths, context->size_skipped_paths, context->missing_args); diff --git a/src/client/scanner.c b/src/client/scanner.c index ce27f2c..61e7dd9 100644 --- a/src/client/scanner.c +++ b/src/client/scanner.c @@ -449,7 +449,7 @@ static void scanner_record_size_skipped(DirectoryScanner* scanner, const char* f Returns false on allocation failure. */ static bool scanner_record_synced_dir(const ScannerOptions* options, const char* fs_path, const char* rel, bool relative_mode) { - if (!options->synced_dirs) + if (!options->synced_dirs && !options->plan_dirs) return true; if (!file_list_dir_in_scope(options->file_list, rel)) return true; @@ -469,7 +469,14 @@ static bool scanner_record_synced_dir(const ScannerOptions* options, const char* dest++; if (dest[0] == '\0') dest = "."; - bool ok = excluded_sink_append(options->synced_dirs, options->excluded_mutex, dest); + bool ok = true; + if (options->synced_dirs) + ok = excluded_sink_append(options->synced_dirs, options->excluded_mutex, dest); + /* The delete-plan keep set needs an entry for every traversed source + directory, including empty ones, so its destination mirror is kept rather + than deleted as an extra; the receive root (".") is implicit. */ + if (ok && options->plan_dirs && strcmp(dest, ".") != 0) + ok = excluded_sink_append(options->plan_dirs, options->excluded_mutex, dest); free(prefixed); return ok; } @@ -528,11 +535,11 @@ static int open_directory_filter_context(DirectoryScanner* scanner, const Filter FilterRuleList* own = read_dir_filters(&scanner->options, scanner->current_path, scanner->current_rel ? scanner->current_rel : "", &any_exists, err, sizeof(err)); - if (!own && any_exists) { - scanner->current_node = (FilterNode*)inherited; - return 0; - } if (!own) { + /* read_dir_filters() leaves `err` set on a parse/allocation failure even + when an earlier merge file in the same directory existed (any_exists true); + key off the error text rather than any_exists so an invalid per-directory + filter file can never be silently ignored. */ if (err[0] == '\0') { scanner->current_node = (FilterNode*)inherited; return 0; @@ -1429,7 +1436,12 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) { wire paths are never recorded (see ScannerOptions.excluded_paths). */ bool files_from_prune = scanner->options.file_list && !file_list_affects(scanner->options.file_list, rel); - if (!files_from_prune && !scanner->relative_mode) { + if (protect && scanner->relative_mode) { + /* -R + --files-from: the destination/wire path is the bare relative + name, so the protected mirror prefix must be `rel` (not the source + path) for the delete walker to match it. */ + scanner_record_excluded(scanner, rel); + } else if (!files_from_prune && !scanner->relative_mode) { if (scanner->options.relative_prefix) { char* wrel = scanner_prefix_send_path(scanner->options.relative_prefix, rel); if (!wrel) { @@ -1856,9 +1868,13 @@ static void scan_root_entry(const ScannerOptions* options, const FilterNode* roo exclusions are never recorded (see ScannerOptions.excluded_paths). */ bool files_from_prune = options->file_list && !file_list_affects(options->file_list, rel); if ((!files_from_prune && !use_rel) || protect) { - const char* rel_path = *cur_path == '/' ? cur_path + 1 : cur_path; + const char* rel_path; char* prefixed = NULL; - if (options->relative_prefix) { + if (use_rel) { + /* -R + --files-from: the destination/wire path is the bare relative + name, not the source path. */ + rel_path = rel; + } else if (options->relative_prefix) { prefixed = scanner_prefix_send_path(options->relative_prefix, entry->d_name); if (!prefixed) { free(rel); @@ -1867,6 +1883,8 @@ static void scan_root_entry(const ScannerOptions* options, const FilterNode* roo return; } rel_path = prefixed; + } else { + rel_path = *cur_path == '/' ? cur_path + 1 : cur_path; } if (options->excluded_paths && !excluded_sink_append(options->excluded_paths, options->excluded_mutex, rel_path)) @@ -2148,9 +2166,9 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory bool any_exists = false; FilterRuleList* own = read_dir_filters(options, root_directory, "", &any_exists, err, sizeof(err)); - if (!own && any_exists) { - /* no files exist: leave root_node NULL */ - } else if (!own) { + if (!own) { + /* A parse/allocation failure must fail the scan even when an earlier + merge file in the same directory existed (see the sequential scanner). */ if (err[0] != '\0') { log_message(LOG_LEVEL_ERROR, "invalid per-directory filter in %s: %s", root_directory, err); array_list_delete(root_files); @@ -2158,6 +2176,7 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory parallel_scanner_destroy(ps); return NULL; } + /* no files exist: leave root_node NULL */ } else if (any_exists && (own->count > 0 || own->dir_merge_count > 0)) { root_node = filter_node_alloc(NULL, own); if (!root_node) { diff --git a/src/client/scanner.h b/src/client/scanner.h index 8788a34..3aeab8c 100644 --- a/src/client/scanner.h +++ b/src/client/scanner.h @@ -114,6 +114,13 @@ typedef struct { * directories, exactly like rsync; the receive root is the "." sentinel. * Guarded by `excluded_mutex`. */ ArrayList* synced_dirs; + /* Delete-plan directory sink (optional): when non-NULL the scanner appends + * the destination-relative path of every directory it traverses (except the + * receive root). The per-directory --delete-during/--delete-delay plan + * builder uses this to keep an empty in-scope source directory (rsync keeps + * it) and to emit its plan after the data stream, when no file frame would + * otherwise trigger it. Guarded by `excluded_mutex`. */ + ArrayList* plan_dirs; /* --ignore-errors: an unreadable directory during the scan is recorded as an * I/O error and skipped instead of aborting the scan. Client-only. */ bool ignore_io_errors; diff --git a/src/shared/delete_plan.c b/src/shared/delete_plan.c index 6b711d3..f0e8fb8 100644 --- a/src/shared/delete_plan.c +++ b/src/shared/delete_plan.c @@ -20,8 +20,6 @@ * 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 */ @@ -387,6 +385,17 @@ int delete_plan_send_for_path(int fd, DeletePlanSender* sender, const char* path return rc; } +int delete_plan_send_remaining(int fd, DeletePlanSender* sender, const ArrayList* dirs) { + if (!sender || !dirs) + return 0; + for (int i = 0; i < dirs->size; i++) { + const char* dir = (const char*)dirs->items[i]; + if (delete_plan_send_for_path(fd, sender, dir, true) != 0) + return -1; + } + return 0; +} + /* ------------------------------------------------------------------ */ /* Receiver: delete session */ /* ------------------------------------------------------------------ */ @@ -398,6 +407,7 @@ struct DeletePlanSession { size_t deleted; size_t skipped; bool limit_hit; + bool limit_logged; bool config_seen; bool missing_applied; ArrayList* protected_prefixes; @@ -811,9 +821,11 @@ int delete_plan_session_receive(DeletePlanSession* session, const Config* config send_status(fd, STATUS_ERROR); return -1; } - if (session->limit_hit) + if (session->limit_hit && !session->limit_logged) { + session->limit_logged = true; log_message(LOG_LEVEL_WARNING, "Deletions stopped due to the delete limit (%zu skipped)", session->skipped); + } return 0; } diff --git a/src/shared/delete_plan.h b/src/shared/delete_plan.h index fb2b97b..1f02a68 100644 --- a/src/shared/delete_plan.h +++ b/src/shared/delete_plan.h @@ -54,6 +54,10 @@ 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); +/* Send the plan for every directory in `dirs` that has not been transmitted + * yet. Called after the data stream so an empty source directory's plan still + * clears its destination extras even though no file frame triggered it. */ +int delete_plan_send_remaining(int fd, DeletePlanSender* sender, const ArrayList* dirs); /* ---- Receiver: delete session ---- */ diff --git a/src/shared/filter.c b/src/shared/filter.c index 1033962..85c5d41 100644 --- a/src/shared/filter.c +++ b/src/shared/filter.c @@ -4,10 +4,23 @@ #include #include #include +#include #include #include #include +/* Write a diagnostic message into the caller's optional buffer. A NULL `err` + * (or a zero size) is a no-op, so a caller that only needs the boolean status + * may pass NULL without the snprintf-on-NULL undefined behaviour. */ +static void filter_set_error(char* err, size_t err_size, const char* fmt, ...) { + if (!err || err_size == 0) + return; + va_list ap; + va_start(ap, fmt); + vsnprintf(err, err_size, fmt, ap); + va_end(ap); +} + /* ---- Ordered rule lists ---- */ void filter_rule_free(FilterRule* rule) { @@ -268,7 +281,7 @@ FilterRule* filter_rule_parse(const char* line, const FilterParseOptions* opts, while (*p == ' ' || *p == '\t') p++; if (*p == '\0' || *p == '\n' || *p == '\r') { - snprintf(err, err_size, "empty filter rule"); + filter_set_error(err, err_size, "empty filter rule"); return NULL; } @@ -279,27 +292,27 @@ FilterRule* filter_rule_parse(const char* line, const FilterParseOptions* opts, size_t pat_len; if (!parse_rule_syntax(p, &kind, &sides, &sides_explicit, &negate, &anchored_mod, &perishable, &xattr, &cvs_inject, &pat, &pat_len)) { - snprintf(err, err_size, "unrecognized filter rule syntax"); + filter_set_error(err, err_size, "unrecognized filter rule syntax"); return NULL; } if (cvs_inject) { /* The C modifier expands to the CVS defaults in place; the rule itself carries no pattern and is handled by the caller. */ - snprintf(err, err_size, "the C modifier is handled by the rule-list parser"); + filter_set_error(err, err_size, "the C modifier is handled by the rule-list parser"); return NULL; } if (kind == RULE_KIND_MERGE || kind == RULE_KIND_DIR_MERGE) { - snprintf(err, err_size, "merge/dir-merge rules are handled by the rule-list parser"); + filter_set_error(err, err_size, "merge/dir-merge rules are handled by the rule-list parser"); return NULL; } if (kind == RULE_KIND_CLEAR) { if (pat_len != 0) { - snprintf(err, err_size, "clear takes no pattern"); + filter_set_error(err, err_size, "clear takes no pattern"); return NULL; } FilterRule* rule = calloc(1, sizeof(FilterRule)); if (!rule) { - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return NULL; } rule->action = FILTER_ACTION_NONE; /* clear marker: no pattern */ @@ -335,7 +348,7 @@ FilterRule* filter_rule_parse(const char* line, const FilterParseOptions* opts, sides = FILTER_SIDE_SENDER; if (pat_len == 0) { - snprintf(err, err_size, "filter rule has no pattern"); + filter_set_error(err, err_size, "filter rule has no pattern"); return NULL; } @@ -352,7 +365,7 @@ FilterRule* filter_rule_parse(const char* line, const FilterParseOptions* opts, pat_len--; } if (pat_len == 0) { - snprintf(err, err_size, "filter rule has no pattern after '/' anchor"); + filter_set_error(err, err_size, "filter rule has no pattern after '/' anchor"); return NULL; } bool dir_only = false; @@ -361,19 +374,19 @@ FilterRule* filter_rule_parse(const char* line, const FilterParseOptions* opts, pat_len--; } if (pat_len == 0) { - snprintf(err, err_size, "filter rule has no pattern"); + filter_set_error(err, err_size, "filter rule has no pattern"); return NULL; } FilterRule* rule = calloc(1, sizeof(FilterRule)); if (!rule) { - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return NULL; } rule->pattern = malloc(pat_len + 1); if (!rule->pattern) { free(rule); - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return NULL; } memcpy(rule->pattern, pat_begin, pat_len); @@ -450,18 +463,18 @@ static bool filter_list_merge_file(FilterRuleList* list, const char* name, const FilterParseOptions* opts, const char* base_dir, int depth, char* err, size_t err_size) { if (name[0] == '\0') { - snprintf(err, err_size, "merge requires a filename"); + filter_set_error(err, err_size, "merge requires a filename"); return false; } char* path = (base_dir && base_dir[0] && name[0] != '/') ? path_cat(base_dir, name) : str_dup(name); if (!path) { - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return false; } FILE* fp = fopen(path, "r"); if (!fp) { - snprintf(err, err_size, "could not read merge file '%s': %s", path, strerror(errno)); + filter_set_error(err, err_size, "could not read merge file '%s': %s", path, strerror(errno)); free(path); return false; } @@ -471,7 +484,7 @@ static bool filter_list_merge_file(FilterRuleList* list, const char* name, while (true) { ssize_t n = utils_getdelim_bounded(fp, &line, &cap, '\n', UTILS_MAX_LINE_LEN); if (n < 0) { - snprintf(err, err_size, "error reading merge file '%s'", path); + filter_set_error(err, err_size, "error reading merge file '%s'", path); ok = false; break; } @@ -499,7 +512,7 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin const FilterParseOptions* opts, const char* base_dir, int depth, char* err, size_t err_size) { if (depth > FILTER_MAX_MERGE_DEPTH) { - snprintf(err, err_size, "merge files nested too deeply"); + filter_set_error(err, err_size, "merge files nested too deeply"); return false; } const char* p = line; @@ -515,7 +528,7 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin size_t pat_len; if (!parse_rule_syntax(p, &kind, &sides, &sides_explicit, &negate, &anchored_mod, &perishable, &xattr, &cvs_inject, &pat, &pat_len)) { - snprintf(err, err_size, "unrecognized filter rule syntax: %s", p); + filter_set_error(err, err_size, "unrecognized filter rule syntax: %s", p); return false; } (void)sides_explicit; @@ -530,7 +543,7 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin } if (kind == RULE_KIND_CLEAR) { if (pat_len != 0) { - snprintf(err, err_size, "clear takes no pattern"); + filter_set_error(err, err_size, "clear takes no pattern"); return false; } for (int i = 0; i < list->count; i++) @@ -540,12 +553,12 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin } if (kind == RULE_KIND_MERGE) { if (pat_len == 0) { - snprintf(err, err_size, "merge requires a filename"); + filter_set_error(err, err_size, "merge requires a filename"); return false; } char* name = malloc(pat_len + 1); if (!name) { - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return false; } memcpy(name, pat, pat_len); @@ -556,12 +569,12 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin } if (kind == RULE_KIND_DIR_MERGE) { if (pat_len == 0) { - snprintf(err, err_size, "dir-merge requires a filename"); + filter_set_error(err, err_size, "dir-merge requires a filename"); return false; } char* name = malloc(pat_len + 1); if (!name) { - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return false; } memcpy(name, pat, pat_len); @@ -569,7 +582,7 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin bool ok = filter_rule_list_add_dir_merge(list, name); free(name); if (!ok) { - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return false; } return true; @@ -580,7 +593,7 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin return false; if (!filter_rule_list_add(list, rule)) { filter_rule_free(rule); - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return false; } return true; @@ -602,7 +615,7 @@ FilterRuleList* filter_base_build(const char* const* rule_texts, int rule_count, err[0] = '\0'; FilterRuleList* list = filter_rule_list_create(); if (!list) { - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return NULL; } FilterParseOptions opts = {.delete_excluded = delete_excluded, .cvs_exclude = cvs_exclude}; @@ -616,7 +629,7 @@ FilterRuleList* filter_base_build(const char* const* rule_texts, int rule_count, } if (cvs_exclude && !filter_list_append_cvs(list, FILTER_SIDE_SENDER | FILTER_SIDE_RECEIVER)) { filter_rule_list_free(list); - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return NULL; } return list; @@ -635,7 +648,7 @@ bool filter_file_append(FilterRuleList* list, const char* dir_path, const char* return false; char* filter_path = path_cat(dir_path, name); if (!filter_path) { - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return false; } FILE* fp = fopen(filter_path, "r"); @@ -652,6 +665,7 @@ bool filter_file_append(FilterRuleList* list, const char* dir_path, const char* if (exists) *exists = true; int rules_before = list->count; + int dir_merges_before = list->dir_merge_count; char* line = NULL; size_t line_cap = 0; bool ok = true; @@ -659,9 +673,10 @@ bool filter_file_append(FilterRuleList* list, const char* dir_path, const char* ssize_t n = utils_getdelim_bounded(fp, &line, &line_cap, '\n', UTILS_MAX_LINE_LEN); if (n < 0) { if (errno == EFBIG) { - snprintf(err, err_size, "line in %s exceeds %d bytes", name, (int)UTILS_MAX_LINE_LEN); + filter_set_error(err, err_size, "line in %s exceeds %d bytes", name, + (int)UTILS_MAX_LINE_LEN); } else { - snprintf(err, err_size, "error reading %s: %s", name, strerror(errno)); + filter_set_error(err, err_size, "error reading %s: %s", name, strerror(errno)); } ok = false; break; @@ -683,16 +698,19 @@ bool filter_file_append(FilterRuleList* list, const char* dir_path, const char* free(line); fclose(fp); if (!ok) { - /* Drop only the rules this file appended, leaving the caller's earlier - content untouched. */ + /* Drop only the rules and dir-merge registrations this file appended, + leaving the caller's earlier content untouched. */ for (int i = rules_before; i < list->count; i++) filter_rule_free(list->items[i]); list->count = rules_before; + for (int i = dir_merges_before; i < list->dir_merge_count; i++) + free(list->dir_merge_names[i]); + list->dir_merge_count = dir_merges_before; return false; } for (int i = rules_before; i < list->count; i++) { if (!set_rule_owner(list->items[i], owner_rel)) { - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return false; } } @@ -705,7 +723,7 @@ FilterRuleList* filter_file_read_named(const char* dir_path, const char* name, FilterRuleList* list = filter_rule_list_create(); if (!list) { if (err && err_size > 0) - snprintf(err, err_size, "memory allocation failed"); + filter_set_error(err, err_size, "memory allocation failed"); return NULL; } if (!filter_file_append(list, dir_path, name, owner_rel, opts, exists, err, err_size)) { diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 5e31e64..a4f56bb 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -32,6 +32,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que context->excluded_paths = NULL; context->size_skipped_paths = NULL; context->synced_dirs = NULL; + context->plan_dirs = NULL; context->missing_args = NULL; context->scan_had_io_error = false; context->remove_source_files = NULL; @@ -196,6 +197,8 @@ void pipeline_context_sender_destroy(PipelineContextSender* context) { array_list_delete(context->size_skipped_paths); if (context->synced_dirs) array_list_delete(context->synced_dirs); + if (context->plan_dirs) + array_list_delete(context->plan_dirs); if (context->missing_args) array_list_delete(context->missing_args); if (context->remove_source_files) diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index 39d2c38..d2e5284 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -55,6 +55,11 @@ typedef struct { "delete only in synchronized directories" (notably for --files-from). Populated by the scanner thread or the early pre-scan. */ ArrayList* synced_dirs; + /* Destination-relative paths of every traversed source directory, for the + per-directory delete plan keep set (so an empty source directory survives + --delete rather than being removed as an extra). Prebuilt by the path-only + pre-scan on the calling thread. */ + ArrayList* plan_dirs; /* --delete-missing-args: the destination-relative mirrors of the --files-from entries that are missing under the source. Computed by the preflight on the calling thread before the pipeline starts; the sender thread transmits diff --git a/tests/integration/test_parity_blockers.py b/tests/integration/test_parity_blockers.py index 6b689e6..5be98bc 100644 --- a/tests/integration/test_parity_blockers.py +++ b/tests/integration/test_parity_blockers.py @@ -73,10 +73,10 @@ class TestRelativePerDirDeleteScope: server.start(extra_args=["--allow-delete"]) result, _ = run_client(spec, dest, flags=["-a", "-R", timing], port=server.port) assert result.returncode == 0, (result.stderr or result.stdout)[:300] - # The prefix's parent-directory sibling survives on both sides. +#The prefix's parent-directory sibling survives on both sides. assert os.path.isfile(os.path.join(dest, "unrelated", "keep.txt")) assert os.path.isfile(os.path.join(rdst, "unrelated", "keep.txt")) - # The in-scope extra is removed on both sides. +#The in - scope extra is removed on both sides. assert not os.path.exists(os.path.join(dest, "foo", "extra.txt")) assert not os.path.exists(os.path.join(rdst, "foo", "extra.txt")) assert _tree(dest) == _tree(rdst) @@ -100,7 +100,7 @@ def _seed_delta_pair(tag): clean_dir(rdst) payload = (b"0123456789abcdef" * 16384)[:200000] _write(os.path.join(source, "f.bin"), payload) - # Destination basis: same length, one byte changed, deliberately older. +#Destination basis : same length, one byte changed, deliberately older. basis = bytearray(payload) basis[100000] ^= 0xFF received = get_dest_received_dir(dest, source) @@ -166,3 +166,130 @@ class TestReceiverWireStats: line for line in result.stdout.splitlines() if line.startswith("*deleting") ) assert fast_del and fast_del == rsync_del, f"rsync={rsync_del}\nfastsync={fast_del}" + + +class TestRelativeFilesFromProtect: + """Blocker #9: a -R + --files-from receiver-protect rule must record the bare + relative wire path so the protected destination mirror survives --delete.""" + + @pytest.mark.ci + @pytest.mark.parametrize("mt", [False, True]) + def test_hidden_protected_mirror_survives_delete(self, mt): + source = os.path.join(TEST_DATA_DIR, "rfprot_src") + dest = os.path.join(TEST_DATA_DIR, "rfprot_dst") + clean_dir(source) + clean_dir(dest) +#Root - level entry exercises the parallel root scanner; the nested one +#exercises the sequential worker scanner. + _write(os.path.join(source, "root_secret.tmp"), b"root\n") + _write(os.path.join(source, "sub", "nested_secret.tmp"), b"nested\n") + _write(os.path.join(source, "sub", "keep.txt"), b"keep\n") + listfile = os.path.join(TEST_DATA_DIR, "rfprot.list") + with open(listfile, "w") as fh: + fh.write(".\n") +#H hides from the sender, P protects the receiver mirror from-- delete. + filters = ["--filter=H root_secret.tmp", "--filter=P root_secret.tmp", + "--filter=H sub/nested_secret.tmp", "--filter=P sub/nested_secret.tmp"] + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + seed = ["--files-from", listfile, "-R"] + (["--threads"] if mt else []) + result, _ = run_client(source, dest, flags=seed, port=server.port) + assert result.returncode == 0, (result.stderr or result.stdout)[:300] + assert os.path.isfile(os.path.join(dest, "root_secret.tmp")) + assert os.path.isfile(os.path.join(dest, "sub", "nested_secret.tmp")) + _write(os.path.join(dest, "extra.txt"), b"extra\n") + _write(os.path.join(dest, "sub", "extra.txt"), b"extra\n") + flags = seed + ["--delete"] + filters + result, _ = run_client(source, dest, flags=flags, port=server.port) + assert result.returncode == 0, (result.stderr or result.stdout)[:300] + assert os.path.isfile(os.path.join(dest, "root_secret.tmp")), \ + "root-level protected mirror was deleted" + assert os.path.isfile(os.path.join(dest, "sub", "nested_secret.tmp")), \ + "nested protected mirror was deleted" + assert not os.path.exists(os.path.join(dest, "extra.txt")) + assert not os.path.exists(os.path.join(dest, "sub", "extra.txt")) + + +class TestInvalidPerDirFilter: + """Blocker #8: a per-directory filter file that fails to parse must fail the + scan even when an earlier merge file in the same directory existed.""" + + @pytest.mark.ci + @pytest.mark.parametrize("mt", [False, True]) + def test_invalid_dir_filter_fails_scan(self, mt): + source = os.path.join(TEST_DATA_DIR, "badfilter_src") + dest = os.path.join(TEST_DATA_DIR, "badfilter_dst") + clean_dir(source) + clean_dir(dest) +#A valid.rsync - filter makes any_exists true for the directory; the +#invalid.rules must not then be silently ignored. + _write(os.path.join(source, ".rsync-filter"), b"- *.bak\n") + _write(os.path.join(source, ".rules"), b"protect\n") + _write(os.path.join(source, "a.txt"), b"a\n") + flags = ["-a", "-F", "--filter=: .rules"] + if mt: + flags.append("--threads") + with ServerManager() as server: + result, _ = run_client(source, dest, flags=flags, port=server.port) + assert result.returncode != 0, "invalid per-directory filter was silently ignored" + assert "invalid per-directory filter" in (result.stderr + result.stdout) + + +class TestWouldDeleteEscaping: + """Blocker #5: -n --delete --out-format must escape control bytes in a + peer-supplied would-delete path so it cannot forge output lines.""" + + @pytest.mark.ci + def test_out_format_escapes_control_chars(self): + source = os.path.join(TEST_DATA_DIR, "esc_src") + dest = os.path.join(TEST_DATA_DIR, "esc_dst") + clean_dir(source) + _write(os.path.join(source, "a.txt"), b"a\n") + received = get_dest_received_dir(dest, source) + clean_dir(received) + _write(os.path.join(received, "a.txt"), b"a\n") +#A newline in a destination filename must not split the printed line. + with open(os.path.join(received, "evil\nname.txt"), "wb") as fh: + fh.write(b"x\n") + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, + flags=["-a", "-n", "--delete", "--out-format=%n"], + port=server.port) + assert result.returncode == 0, (result.stderr or result.stdout)[:300] + assert "\\#012" in result.stdout, result.stdout + assert "evil\nname.txt" not in result.stdout, result.stdout + + +class TestEmptySourceDirectoryDelete: + """Blocker #10: an empty in-scope source directory must survive + --delete-during/--delete-delay (rsync keeps it) while its extras are still + removed.""" + + @requires_rsync + @pytest.mark.ci + @pytest.mark.parametrize("timing", ["--delete-during", "--delete-delay"]) + def test_empty_source_dir_survives_matches_rsync(self, timing): + source = os.path.join(TEST_DATA_DIR, "emptydir_src") + dest = os.path.join(TEST_DATA_DIR, "emptydir_dst") + rdst = os.path.join(TEST_DATA_DIR, "emptydir_rdst") + clean_dir(source) + os.makedirs(os.path.join(source, "empty")) + _write(os.path.join(source, "keep.txt"), b"keep\n") + received = get_dest_received_dir(dest, source) + for root in (rdst, received): + clean_dir(root) + _write(os.path.join(root, "keep.txt"), b"keep\n") + _write(os.path.join(root, "empty", "extra.txt"), b"extra\n") + rsync_result = _rsync(["-a", timing, source + "/", rdst + "/"]) + assert rsync_result.returncode == 0, rsync_result.stderr + assert os.path.isdir(os.path.join(rdst, "empty")) + assert not os.path.exists(os.path.join(rdst, "empty", "extra.txt")) + with ServerManager() as server: + server.start(extra_args=["--allow-delete"]) + result, _ = run_client(source, dest, flags=["-a", timing], port=server.port) + assert result.returncode == 0, (result.stderr or result.stdout)[:300] + assert os.path.isdir(os.path.join(received, "empty")), \ + "empty source directory was removed" + assert not os.path.exists(os.path.join(received, "empty", "extra.txt")) + assert _tree(received) == _tree(rdst) diff --git a/tests/test_server.c b/tests/test_server.c index 76ee4c4..501c9a3 100644 --- a/tests/test_server.c +++ b/tests/test_server.c @@ -1153,10 +1153,10 @@ static void test_dry_run_delete_plan_commit_does_not_delete() { io_set_bwlimit(0); EXPECT_TRUE(send_status(p[1], STATUS_DELETE_PLAN)); - EXPECT_TRUE(send_int(p[1], 1)); /* first frame carries the config sections */ - EXPECT_TRUE(send_int(p[1], 0)); /* protected prefixes */ - EXPECT_TRUE(send_int(p[1], 0)); /* size-skipped prefixes */ - EXPECT_TRUE(send_int(p[1], 1)); /* missing-args exact deletions */ + EXPECT_TRUE(send_int(p[1], 1)); /* first frame carries the config sections */ + EXPECT_TRUE(send_int(p[1], 0)); /* protected prefixes */ + EXPECT_TRUE(send_int(p[1], 0)); /* size-skipped prefixes */ + EXPECT_TRUE(send_int(p[1], 1)); /* missing-args exact deletions */ EXPECT_TRUE(send_str(p[1], "victim.txt")); EXPECT_TRUE(send_str(p[1], ".")); /* receive root plan */ EXPECT_TRUE(send_int(p[1], 0)); /* kept child directories */