Release v2.26.0 #284

Merged
TapTap merged 210 commits from dev into main 2026-09-18 19:05:52 +02:00
10 changed files with 314 additions and 65 deletions
Showing only changes of commit 6c6f02e5dd - Show all commits
+63 -9
View File
@@ -923,15 +923,27 @@ static bool receive_stats_record(int fd, ReceiverStats* stats, ArrayList* would_
int count = 0; int count = 0;
if (!receive_int(fd, &count) || count < 0 || count > MAX_MANIFEST_ENTRIES) if (!receive_int(fd, &count) || count < 0 || count > MAX_MANIFEST_ENTRIES)
return false; 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++) { for (int i = 0; i < count; i++) {
char* path = receive_wire_str(fd); char* path = receive_wire_str(fd);
if (!path) if (!path)
return false; return false;
if (path[0] == '\0' || path[0] == '/' || has_path_traversal(path)) {
free(path);
return false;
}
if (would_delete) { if (would_delete) {
char* copy = str_dup(path); 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); free(path);
if (!copy || !array_list_add(would_delete, copy)) {
free(copy);
return false; return false;
} }
} else { } else {
@@ -1455,6 +1467,19 @@ static bool scan_paths_only(const Config* config, const ScannerOptions* options,
} }
chunk_destroy(chunk); 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)) if (ok && directory_scanner_failed(scanner))
ok = false; ok = false;
if (io_error_out) if (io_error_out)
@@ -1929,7 +1954,12 @@ static int send_dry_run_remote(Config* config) {
event.path = path; event.path = path;
char* line = change_render_format(config->out_format, config, &event); char* line = change_render_format(config->out_format, config, &event);
if (line) { 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); free(line);
} }
} else { } else {
@@ -2474,6 +2504,13 @@ static int send_chunks_multithreaded(void* pipeline_context) {
NULL) != 0) NULL) != 0)
goto send_fail; 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 /* 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 (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 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. */ /* Size-pruned prefixes (always protected) and synchronized directories. */
ArrayList* size_skipped = NULL; ArrayList* size_skipped = NULL;
ArrayList* synced_dirs = 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); 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 /* -d/--dirs does not recurse, so a per-directory plan would carry no child
information and could delete the contents of an untraversed directory; 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 receive root's extras are handled exactly like rsync's first generator
directory. The remaining plans are streamed with the data below. */ directory. The remaining plans are streamed with the data below. */
plan_sender = delete_plan_sender_create(); plan_sender = delete_plan_sender_create();
if (!plan_sender) plan_dirs = array_list_create(free);
if (!plan_sender || !plan_dirs)
goto send_fail; 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 prescan_ok = scan_paths_only(config, &prepared.options, NULL, plan_sender, &had_scan_io);
bool plans_ok = false; bool plans_ok = false;
if (prescan_ok) { if (prescan_ok) {
const char* walk_root = delete_plan_walk_root(config, synced_dirs); 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_finalize(plan_sender, scope, walk_root);
delete_plan_sender_set_config(plan_sender, excluded, size_skipped, missing_args); delete_plan_sender_set_config(plan_sender, excluded, size_skipped, missing_args);
if (had_scan_io && delete_plan_sender_empty(plan_sender)) { 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.excluded_paths = NULL;
prepared.options.size_skipped_paths = NULL; prepared.options.size_skipped_paths = NULL;
prepared.options.synced_dirs = NULL; prepared.options.synced_dirs = NULL;
prepared.options.plan_dirs = NULL;
if (!prescan_ok || !plans_ok) if (!prescan_ok || !plans_ok)
goto send_fail; goto send_fail;
} else if (config->use_delete) { } 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 /* P7 Wave D: every directory has now been traversed (or the scan stopped
early), so transmit the captured directory times last. The receiver defers early), so transmit the captured directory times last. The receiver defers
applying them until after its own deletion/publication phase. */ applying them until after its own deletion/publication phase. */
@@ -3127,6 +3176,8 @@ send_fail:
array_list_delete(size_skipped); array_list_delete(size_skipped);
if (synced_dirs) if (synced_dirs)
array_list_delete(synced_dirs); array_list_delete(synced_dirs);
if (plan_dirs)
array_list_delete(plan_dirs);
if (missing_args) if (missing_args)
array_list_delete(missing_args); array_list_delete(missing_args);
if (remove_sources) if (remove_sources)
@@ -3268,7 +3319,10 @@ int send_files_multithreaded(Config** config_ptr) {
} }
if (per_dir) { if (per_dir) {
context->delete_plans = delete_plan_sender_create(); 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 { } else {
context->manifest = array_list_create(free); context->manifest = array_list_create(free);
prepared_ok = prepared_ok && context->manifest != NULL; prepared_ok = prepared_ok && context->manifest != NULL;
@@ -3279,8 +3333,8 @@ int send_files_multithreaded(Config** config_ptr) {
prepared_scanner_destroy(&prepared); prepared_scanner_destroy(&prepared);
if (per_dir && prebuilt) { if (per_dir && prebuilt) {
const char* walk_root = delete_plan_walk_root(config, context->synced_dirs); const char* walk_root = delete_plan_walk_root(config, context->synced_dirs);
const ArrayList* scope = const ArrayList* scope = config->files_from_set ? context->synced_dirs
config->files_from_set ? context->synced_dirs : (walk_root ? context->synced_dirs : NULL); : (walk_root ? context->synced_dirs : NULL);
delete_plan_sender_finalize(context->delete_plans, scope, walk_root); delete_plan_sender_finalize(context->delete_plans, scope, walk_root);
delete_plan_sender_set_config(context->delete_plans, context->excluded_paths, delete_plan_sender_set_config(context->delete_plans, context->excluded_paths,
context->size_skipped_paths, context->missing_args); context->size_skipped_paths, context->missing_args);
+31 -12
View File
@@ -449,7 +449,7 @@ static void scanner_record_size_skipped(DirectoryScanner* scanner, const char* f
Returns false on allocation failure. */ Returns false on allocation failure. */
static bool scanner_record_synced_dir(const ScannerOptions* options, const char* fs_path, static bool scanner_record_synced_dir(const ScannerOptions* options, const char* fs_path,
const char* rel, bool relative_mode) { const char* rel, bool relative_mode) {
if (!options->synced_dirs) if (!options->synced_dirs && !options->plan_dirs)
return true; return true;
if (!file_list_dir_in_scope(options->file_list, rel)) if (!file_list_dir_in_scope(options->file_list, rel))
return true; return true;
@@ -469,7 +469,14 @@ static bool scanner_record_synced_dir(const ScannerOptions* options, const char*
dest++; dest++;
if (dest[0] == '\0') if (dest[0] == '\0')
dest = "."; 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); free(prefixed);
return ok; 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, FilterRuleList* own = read_dir_filters(&scanner->options, scanner->current_path,
scanner->current_rel ? scanner->current_rel : "", scanner->current_rel ? scanner->current_rel : "",
&any_exists, err, sizeof(err)); &any_exists, err, sizeof(err));
if (!own && any_exists) {
scanner->current_node = (FilterNode*)inherited;
return 0;
}
if (!own) { 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') { if (err[0] == '\0') {
scanner->current_node = (FilterNode*)inherited; scanner->current_node = (FilterNode*)inherited;
return 0; return 0;
@@ -1429,7 +1436,12 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
wire paths are never recorded (see ScannerOptions.excluded_paths). */ wire paths are never recorded (see ScannerOptions.excluded_paths). */
bool files_from_prune = bool files_from_prune =
scanner->options.file_list && !file_list_affects(scanner->options.file_list, rel); 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) { if (scanner->options.relative_prefix) {
char* wrel = scanner_prefix_send_path(scanner->options.relative_prefix, rel); char* wrel = scanner_prefix_send_path(scanner->options.relative_prefix, rel);
if (!wrel) { 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). */ exclusions are never recorded (see ScannerOptions.excluded_paths). */
bool files_from_prune = options->file_list && !file_list_affects(options->file_list, rel); bool files_from_prune = options->file_list && !file_list_affects(options->file_list, rel);
if ((!files_from_prune && !use_rel) || protect) { 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; 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); prefixed = scanner_prefix_send_path(options->relative_prefix, entry->d_name);
if (!prefixed) { if (!prefixed) {
free(rel); free(rel);
@@ -1867,6 +1883,8 @@ static void scan_root_entry(const ScannerOptions* options, const FilterNode* roo
return; return;
} }
rel_path = prefixed; rel_path = prefixed;
} else {
rel_path = *cur_path == '/' ? cur_path + 1 : cur_path;
} }
if (options->excluded_paths && if (options->excluded_paths &&
!excluded_sink_append(options->excluded_paths, options->excluded_mutex, rel_path)) !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; bool any_exists = false;
FilterRuleList* own = FilterRuleList* own =
read_dir_filters(options, root_directory, "", &any_exists, err, sizeof(err)); read_dir_filters(options, root_directory, "", &any_exists, err, sizeof(err));
if (!own && any_exists) { if (!own) {
/* no files exist: leave root_node NULL */ /* A parse/allocation failure must fail the scan even when an earlier
} else if (!own) { merge file in the same directory existed (see the sequential scanner). */
if (err[0] != '\0') { if (err[0] != '\0') {
log_message(LOG_LEVEL_ERROR, "invalid per-directory filter in %s: %s", root_directory, err); log_message(LOG_LEVEL_ERROR, "invalid per-directory filter in %s: %s", root_directory, err);
array_list_delete(root_files); array_list_delete(root_files);
@@ -2158,6 +2176,7 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
parallel_scanner_destroy(ps); parallel_scanner_destroy(ps);
return NULL; return NULL;
} }
/* no files exist: leave root_node NULL */
} else if (any_exists && (own->count > 0 || own->dir_merge_count > 0)) { } else if (any_exists && (own->count > 0 || own->dir_merge_count > 0)) {
root_node = filter_node_alloc(NULL, own); root_node = filter_node_alloc(NULL, own);
if (!root_node) { if (!root_node) {
+7
View File
@@ -114,6 +114,13 @@ typedef struct {
* directories, exactly like rsync; the receive root is the "." sentinel. * directories, exactly like rsync; the receive root is the "." sentinel.
* Guarded by `excluded_mutex`. */ * Guarded by `excluded_mutex`. */
ArrayList* synced_dirs; 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 /* --ignore-errors: an unreadable directory during the scan is recorded as an
* I/O error and skipped instead of aborting the scan. Client-only. */ * I/O error and skipped instead of aborting the scan. Client-only. */
bool ignore_io_errors; bool ignore_io_errors;
+15 -3
View File
@@ -20,8 +20,6 @@
* the number of entries one deletion commit may remove. A client * the number of entries one deletion commit may remove. A client
* --max-delete=NUM smaller than this replaces it for the run. */ * --max-delete=NUM smaller than this replaces it for the run. */
#define DELETE_PLAN_SERVER_LIMIT 100000U #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 */ /* Sender: plan builder */
@@ -387,6 +385,17 @@ int delete_plan_send_for_path(int fd, DeletePlanSender* sender, const char* path
return rc; 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 */ /* Receiver: delete session */
/* ------------------------------------------------------------------ */ /* ------------------------------------------------------------------ */
@@ -398,6 +407,7 @@ struct DeletePlanSession {
size_t deleted; size_t deleted;
size_t skipped; size_t skipped;
bool limit_hit; bool limit_hit;
bool limit_logged;
bool config_seen; bool config_seen;
bool missing_applied; bool missing_applied;
ArrayList* protected_prefixes; ArrayList* protected_prefixes;
@@ -811,9 +821,11 @@ int delete_plan_session_receive(DeletePlanSession* session, const Config* config
send_status(fd, STATUS_ERROR); send_status(fd, STATUS_ERROR);
return -1; 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)", log_message(LOG_LEVEL_WARNING, "Deletions stopped due to the delete limit (%zu skipped)",
session->skipped); session->skipped);
}
return 0; return 0;
} }
+4
View File
@@ -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, /* Send the plans for every ancestor of `path` (root-first) and, when is_dir,
* for `path` itself; already-sent plans are skipped. */ * for `path` itself; already-sent plans are skipped. */
int delete_plan_send_for_path(int fd, DeletePlanSender* sender, const char* path, bool is_dir); 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 ---- */ /* ---- Receiver: delete session ---- */
+51 -33
View File
@@ -4,10 +4,23 @@
#include <ctype.h> #include <ctype.h>
#include <errno.h> #include <errno.h>
#include <limits.h> #include <limits.h>
#include <stdarg.h>
#include <stdio.h> #include <stdio.h>
#include <stdlib.h> #include <stdlib.h>
#include <string.h> #include <string.h>
/* 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 ---- */ /* ---- Ordered rule lists ---- */
void filter_rule_free(FilterRule* rule) { void filter_rule_free(FilterRule* rule) {
@@ -268,7 +281,7 @@ FilterRule* filter_rule_parse(const char* line, const FilterParseOptions* opts,
while (*p == ' ' || *p == '\t') while (*p == ' ' || *p == '\t')
p++; p++;
if (*p == '\0' || *p == '\n' || *p == '\r') { 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; return NULL;
} }
@@ -279,27 +292,27 @@ FilterRule* filter_rule_parse(const char* line, const FilterParseOptions* opts,
size_t pat_len; size_t pat_len;
if (!parse_rule_syntax(p, &kind, &sides, &sides_explicit, &negate, &anchored_mod, &perishable, if (!parse_rule_syntax(p, &kind, &sides, &sides_explicit, &negate, &anchored_mod, &perishable,
&xattr, &cvs_inject, &pat, &pat_len)) { &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; return NULL;
} }
if (cvs_inject) { if (cvs_inject) {
/* The C modifier expands to the CVS defaults in place; the rule itself /* The C modifier expands to the CVS defaults in place; the rule itself
carries no pattern and is handled by the caller. */ 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; return NULL;
} }
if (kind == RULE_KIND_MERGE || kind == RULE_KIND_DIR_MERGE) { 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; return NULL;
} }
if (kind == RULE_KIND_CLEAR) { if (kind == RULE_KIND_CLEAR) {
if (pat_len != 0) { if (pat_len != 0) {
snprintf(err, err_size, "clear takes no pattern"); filter_set_error(err, err_size, "clear takes no pattern");
return NULL; return NULL;
} }
FilterRule* rule = calloc(1, sizeof(FilterRule)); FilterRule* rule = calloc(1, sizeof(FilterRule));
if (!rule) { if (!rule) {
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return NULL; return NULL;
} }
rule->action = FILTER_ACTION_NONE; /* clear marker: no pattern */ 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; sides = FILTER_SIDE_SENDER;
if (pat_len == 0) { 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; return NULL;
} }
@@ -352,7 +365,7 @@ FilterRule* filter_rule_parse(const char* line, const FilterParseOptions* opts,
pat_len--; pat_len--;
} }
if (pat_len == 0) { 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; return NULL;
} }
bool dir_only = false; bool dir_only = false;
@@ -361,19 +374,19 @@ FilterRule* filter_rule_parse(const char* line, const FilterParseOptions* opts,
pat_len--; pat_len--;
} }
if (pat_len == 0) { 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; return NULL;
} }
FilterRule* rule = calloc(1, sizeof(FilterRule)); FilterRule* rule = calloc(1, sizeof(FilterRule));
if (!rule) { if (!rule) {
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return NULL; return NULL;
} }
rule->pattern = malloc(pat_len + 1); rule->pattern = malloc(pat_len + 1);
if (!rule->pattern) { if (!rule->pattern) {
free(rule); free(rule);
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return NULL; return NULL;
} }
memcpy(rule->pattern, pat_begin, pat_len); 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, const FilterParseOptions* opts, const char* base_dir, int depth,
char* err, size_t err_size) { char* err, size_t err_size) {
if (name[0] == '\0') { if (name[0] == '\0') {
snprintf(err, err_size, "merge requires a filename"); filter_set_error(err, err_size, "merge requires a filename");
return false; return false;
} }
char* path = char* path =
(base_dir && base_dir[0] && name[0] != '/') ? path_cat(base_dir, name) : str_dup(name); (base_dir && base_dir[0] && name[0] != '/') ? path_cat(base_dir, name) : str_dup(name);
if (!path) { if (!path) {
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return false; return false;
} }
FILE* fp = fopen(path, "r"); FILE* fp = fopen(path, "r");
if (!fp) { 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); free(path);
return false; return false;
} }
@@ -471,7 +484,7 @@ static bool filter_list_merge_file(FilterRuleList* list, const char* name,
while (true) { while (true) {
ssize_t n = utils_getdelim_bounded(fp, &line, &cap, '\n', UTILS_MAX_LINE_LEN); ssize_t n = utils_getdelim_bounded(fp, &line, &cap, '\n', UTILS_MAX_LINE_LEN);
if (n < 0) { 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; ok = false;
break; break;
} }
@@ -499,7 +512,7 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin
const FilterParseOptions* opts, const char* base_dir, const FilterParseOptions* opts, const char* base_dir,
int depth, char* err, size_t err_size) { int depth, char* err, size_t err_size) {
if (depth > FILTER_MAX_MERGE_DEPTH) { 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; return false;
} }
const char* p = line; const char* p = line;
@@ -515,7 +528,7 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin
size_t pat_len; size_t pat_len;
if (!parse_rule_syntax(p, &kind, &sides, &sides_explicit, &negate, &anchored_mod, &perishable, if (!parse_rule_syntax(p, &kind, &sides, &sides_explicit, &negate, &anchored_mod, &perishable,
&xattr, &cvs_inject, &pat, &pat_len)) { &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; return false;
} }
(void)sides_explicit; (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 (kind == RULE_KIND_CLEAR) {
if (pat_len != 0) { if (pat_len != 0) {
snprintf(err, err_size, "clear takes no pattern"); filter_set_error(err, err_size, "clear takes no pattern");
return false; return false;
} }
for (int i = 0; i < list->count; i++) 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 (kind == RULE_KIND_MERGE) {
if (pat_len == 0) { if (pat_len == 0) {
snprintf(err, err_size, "merge requires a filename"); filter_set_error(err, err_size, "merge requires a filename");
return false; return false;
} }
char* name = malloc(pat_len + 1); char* name = malloc(pat_len + 1);
if (!name) { if (!name) {
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return false; return false;
} }
memcpy(name, pat, pat_len); 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 (kind == RULE_KIND_DIR_MERGE) {
if (pat_len == 0) { 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; return false;
} }
char* name = malloc(pat_len + 1); char* name = malloc(pat_len + 1);
if (!name) { if (!name) {
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return false; return false;
} }
memcpy(name, pat, pat_len); 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); bool ok = filter_rule_list_add_dir_merge(list, name);
free(name); free(name);
if (!ok) { if (!ok) {
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return false; return false;
} }
return true; return true;
@@ -580,7 +593,7 @@ static bool filter_list_parse_append_depth(FilterRuleList* list, const char* lin
return false; return false;
if (!filter_rule_list_add(list, rule)) { if (!filter_rule_list_add(list, rule)) {
filter_rule_free(rule); filter_rule_free(rule);
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return false; return false;
} }
return true; return true;
@@ -602,7 +615,7 @@ FilterRuleList* filter_base_build(const char* const* rule_texts, int rule_count,
err[0] = '\0'; err[0] = '\0';
FilterRuleList* list = filter_rule_list_create(); FilterRuleList* list = filter_rule_list_create();
if (!list) { if (!list) {
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return NULL; return NULL;
} }
FilterParseOptions opts = {.delete_excluded = delete_excluded, .cvs_exclude = cvs_exclude}; 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)) { if (cvs_exclude && !filter_list_append_cvs(list, FILTER_SIDE_SENDER | FILTER_SIDE_RECEIVER)) {
filter_rule_list_free(list); filter_rule_list_free(list);
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return NULL; return NULL;
} }
return list; return list;
@@ -635,7 +648,7 @@ bool filter_file_append(FilterRuleList* list, const char* dir_path, const char*
return false; return false;
char* filter_path = path_cat(dir_path, name); char* filter_path = path_cat(dir_path, name);
if (!filter_path) { if (!filter_path) {
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return false; return false;
} }
FILE* fp = fopen(filter_path, "r"); FILE* fp = fopen(filter_path, "r");
@@ -652,6 +665,7 @@ bool filter_file_append(FilterRuleList* list, const char* dir_path, const char*
if (exists) if (exists)
*exists = true; *exists = true;
int rules_before = list->count; int rules_before = list->count;
int dir_merges_before = list->dir_merge_count;
char* line = NULL; char* line = NULL;
size_t line_cap = 0; size_t line_cap = 0;
bool ok = true; 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); ssize_t n = utils_getdelim_bounded(fp, &line, &line_cap, '\n', UTILS_MAX_LINE_LEN);
if (n < 0) { if (n < 0) {
if (errno == EFBIG) { 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 { } 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; ok = false;
break; break;
@@ -683,16 +698,19 @@ bool filter_file_append(FilterRuleList* list, const char* dir_path, const char*
free(line); free(line);
fclose(fp); fclose(fp);
if (!ok) { if (!ok) {
/* Drop only the rules this file appended, leaving the caller's earlier /* Drop only the rules and dir-merge registrations this file appended,
content untouched. */ leaving the caller's earlier content untouched. */
for (int i = rules_before; i < list->count; i++) for (int i = rules_before; i < list->count; i++)
filter_rule_free(list->items[i]); filter_rule_free(list->items[i]);
list->count = rules_before; 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; return false;
} }
for (int i = rules_before; i < list->count; i++) { for (int i = rules_before; i < list->count; i++) {
if (!set_rule_owner(list->items[i], owner_rel)) { 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; return false;
} }
} }
@@ -705,7 +723,7 @@ FilterRuleList* filter_file_read_named(const char* dir_path, const char* name,
FilterRuleList* list = filter_rule_list_create(); FilterRuleList* list = filter_rule_list_create();
if (!list) { if (!list) {
if (err && err_size > 0) if (err && err_size > 0)
snprintf(err, err_size, "memory allocation failed"); filter_set_error(err, err_size, "memory allocation failed");
return NULL; return NULL;
} }
if (!filter_file_append(list, dir_path, name, owner_rel, opts, exists, err, err_size)) { if (!filter_file_append(list, dir_path, name, owner_rel, opts, exists, err, err_size)) {
+3
View File
@@ -32,6 +32,7 @@ PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* que
context->excluded_paths = NULL; context->excluded_paths = NULL;
context->size_skipped_paths = NULL; context->size_skipped_paths = NULL;
context->synced_dirs = NULL; context->synced_dirs = NULL;
context->plan_dirs = NULL;
context->missing_args = NULL; context->missing_args = NULL;
context->scan_had_io_error = false; context->scan_had_io_error = false;
context->remove_source_files = NULL; context->remove_source_files = NULL;
@@ -196,6 +197,8 @@ void pipeline_context_sender_destroy(PipelineContextSender* context) {
array_list_delete(context->size_skipped_paths); array_list_delete(context->size_skipped_paths);
if (context->synced_dirs) if (context->synced_dirs)
array_list_delete(context->synced_dirs); array_list_delete(context->synced_dirs);
if (context->plan_dirs)
array_list_delete(context->plan_dirs);
if (context->missing_args) if (context->missing_args)
array_list_delete(context->missing_args); array_list_delete(context->missing_args);
if (context->remove_source_files) if (context->remove_source_files)
+5
View File
@@ -55,6 +55,11 @@ typedef struct {
"delete only in synchronized directories" (notably for --files-from). "delete only in synchronized directories" (notably for --files-from).
Populated by the scanner thread or the early pre-scan. */ Populated by the scanner thread or the early pre-scan. */
ArrayList* synced_dirs; 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 /* --delete-missing-args: the destination-relative mirrors of the --files-from
entries that are missing under the source. Computed by the preflight on entries that are missing under the source. Computed by the preflight on
the calling thread before the pipeline starts; the sender thread transmits the calling thread before the pipeline starts; the sender thread transmits
+130 -3
View File
@@ -73,10 +73,10 @@ class TestRelativePerDirDeleteScope:
server.start(extra_args=["--allow-delete"]) server.start(extra_args=["--allow-delete"])
result, _ = run_client(spec, dest, flags=["-a", "-R", timing], port=server.port) result, _ = run_client(spec, dest, flags=["-a", "-R", timing], port=server.port)
assert result.returncode == 0, (result.stderr or result.stdout)[:300] 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(dest, "unrelated", "keep.txt"))
assert os.path.isfile(os.path.join(rdst, "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(dest, "foo", "extra.txt"))
assert not os.path.exists(os.path.join(rdst, "foo", "extra.txt")) assert not os.path.exists(os.path.join(rdst, "foo", "extra.txt"))
assert _tree(dest) == _tree(rdst) assert _tree(dest) == _tree(rdst)
@@ -100,7 +100,7 @@ def _seed_delta_pair(tag):
clean_dir(rdst) clean_dir(rdst)
payload = (b"0123456789abcdef" * 16384)[:200000] payload = (b"0123456789abcdef" * 16384)[:200000]
_write(os.path.join(source, "f.bin"), payload) _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 = bytearray(payload)
basis[100000] ^= 0xFF basis[100000] ^= 0xFF
received = get_dest_received_dir(dest, source) received = get_dest_received_dir(dest, source)
@@ -166,3 +166,130 @@ class TestReceiverWireStats:
line for line in result.stdout.splitlines() if line.startswith("*deleting") 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}" 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)