Merge feat/p2-files-from-filter: files-from/from0/filter/-F/-C

Client-side file-list selection and filter-rule layer (client-only; no wire
change). rsync-consistent inner-first per-dir filter precedence; listed-but-
missing files-from entries hard-error; root '/' scan regression fixed;
Summary recounted (63 implemented / 75 not).
c-review REQUEST CHANGES -> blockers fixed; PR #265.
This commit is contained in:
2026-09-06 14:25:51 +02:00
15 changed files with 1944 additions and 63 deletions
+66
View File
@@ -4,6 +4,8 @@
#include "compression.h"
#include "config.h"
#include "delta.h"
#include "file_list.h"
#include "filter.h"
#include "log.h"
#include "protocol.h"
#include "transport_tcp.h"
@@ -313,6 +315,31 @@ static int config_add_pattern(char*** patterns, int* count, const char* value,
return 0;
}
/* Validate and append one --filter=RULE string. Returns 0 on success, -1 on error. */
static int config_add_filter(Config* config, const char* rule) {
char err[160];
FilterRule* parsed = filter_rule_parse(rule, err, sizeof(err));
if (!parsed) {
log_message(LOG_LEVEL_ERROR, "invalid --filter rule '%s': %s", rule, err);
return -1;
}
filter_rule_free(parsed);
if (!config->filters) {
config->filters = array_list_create(free);
if (!config->filters) {
log_message(LOG_LEVEL_ERROR, "memory allocation failed for --filter");
return -1;
}
}
char* dup = str_dup(rule);
if (!dup || !array_list_add(config->filters, dup)) {
free(dup);
log_message(LOG_LEVEL_ERROR, "memory allocation failed for --filter");
return -1;
}
return 0;
}
static int parse_skip_compress(Config* config, const char* value) {
char* list = str_dup(value);
if (!list)
@@ -423,6 +450,9 @@ static const OptionEntry OPTION_TABLE[] = {
{"--max-size", NULL, OPT_ULL, offsetof(Config, max_size)},
{"--min-size", NULL, OPT_ULL, offsetof(Config, min_size)},
{"--one-file-system", "-x", OPT_FLAG, offsetof(Config, one_file_system)},
{"--from0", "-0", OPT_FLAG, offsetof(Config, from0)},
{"--cvs-exclude", "-C", OPT_FLAG, offsetof(Config, cvs_exclude)},
{"-F", NULL, OPT_FLAG, offsetof(Config, per_dir_filter)},
};
/* Only boolean options with no required argument are safe to negate. */
@@ -444,6 +474,8 @@ static const NegatableOption NEGATABLE_OPTIONS[] = {
{"sparse", "S", offsetof(Config, preserve_sparse)},
{"inplace", NULL, offsetof(Config, inplace)},
{"checksum", NULL, offsetof(Config, checksum)},
{"from0", NULL, offsetof(Config, from0)},
{"cvs-exclude", NULL, offsetof(Config, cvs_exclude)},
/* These options are also implied by --archive or handled outside the table. */
{"compress", "c", offsetof(Config, use_compression)},
@@ -862,6 +894,26 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
if (read_patterns_from_file(argv[++i], &config->include_patterns, &config->include_count) !=
0)
return -1;
} else if (strncmp(argv[i], "--filter=", 9) == 0) {
if (config_add_filter(config, argv[i] + 9) != 0)
return -1;
} else if (opt_is(argv[i], "--filter", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (config_add_filter(config, argv[++i]) != 0)
return -1;
} else if (strncmp(argv[i], "--files-from=", 13) == 0) {
if (set_string_option(&config->files_from, argv[i] + 13, "--files-from") != 0)
return -1;
} else if (opt_is(argv[i], "--files-from", NULL)) {
if (i + 1 >= argc) {
log_message(LOG_LEVEL_ERROR, "missing argument for %s", argv[i]);
return -1;
}
if (set_string_option(&config->files_from, argv[++i], "--files-from") != 0)
return -1;
} else if (opt_is(argv[i], "-v", "--verbose")) {
verbose = true;
set_log_level(LOG_LEVEL_DEBUG);
@@ -936,6 +988,20 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
if (config->compress_choice)
config->use_compression = strcmp(config->compress_choice, "zstd") == 0;
/* --files-from is loaded after every argument is seen so that -0/--from0 may
* appear anywhere on the command line. A missing or unreadable file, and
* invalid (absolute / traversal) entries, are hard CLI errors. */
if (config->files_from) {
char err[256];
FileListSet* set = file_list_load(config->files_from, config->from0, err, sizeof(err));
if (!set) {
log_message(LOG_LEVEL_ERROR, "--files-from: %s", err);
return -1;
}
file_list_destroy((FileListSet*)config->files_from_set);
config->files_from_set = set;
}
/* Incremental and delta transfers need metadata unless the user disabled it. */
if ((config->use_incremental || config->use_delta) && !config->use_metadata &&
!config->metadata_explicitly_disabled) {
+154 -23
View File
@@ -7,6 +7,8 @@
#include "data.h"
#include "delta.h"
#include "file.h"
#include "file_list.h"
#include "filter.h"
#include "metadata.h"
#include "log.h"
#include "multiprocessing.h"
@@ -40,16 +42,112 @@ static const char* display_bytes(unsigned long long bytes, bool human_readable,
return buffer;
}
static ScannerOptions scanner_options_from_config(const Config* config, int num_threads) {
ScannerOptions options = {config->use_metadata, config->chunk_size,
config->exclude_patterns, config->exclude_count,
config->include_patterns, config->include_count,
config->max_size, config->min_size,
config->max_depth, num_threads,
config->follow_symlinks, config->copy_links,
config->safe_links, config->copy_unsafe_links,
config->checksum, config->one_file_system};
return options;
/* Compiled scanner inputs that are shared read-only across scanner instances
* and, in -m mode, across worker threads. `base_filters` owns the compiled
* command-line + -C rules; the FileListSet allow-set lives in the Config. */
typedef struct {
ScannerOptions options;
FilterRuleList* base_filters; /* owned; may be NULL */
} PreparedScanner;
/* Build the scanner options for one scan. Returns false and logs on failure. */
static bool prepare_scanner(const Config* config, int num_threads, PreparedScanner* out) {
if (!out)
return false;
out->base_filters = NULL;
memset(&out->options, 0, sizeof(out->options));
int rule_count = config->filters ? config->filters->size : 0;
const char** texts = NULL;
if (rule_count > 0) {
texts = malloc((size_t)rule_count * sizeof(char*));
if (!texts) {
log_message(LOG_LEVEL_ERROR, "memory allocation failed for filter rules");
return false;
}
for (int i = 0; i < rule_count; i++)
texts[i] = (const char*)config->filters->items[i];
}
if (rule_count > 0 || config->cvs_exclude) {
char err[160];
out->base_filters = filter_base_build(texts, rule_count, config->cvs_exclude, err, sizeof(err));
free(texts);
if (!out->base_filters) {
log_message(LOG_LEVEL_ERROR, "invalid filter rule: %s", err);
return false;
}
} else {
free(texts);
}
ScannerOptions* options = &out->options;
options->use_metadata = config->use_metadata;
options->chunk_size = config->chunk_size;
options->exclude_patterns = config->exclude_patterns;
options->exclude_count = config->exclude_count;
options->include_patterns = config->include_patterns;
options->include_count = config->include_count;
options->max_size = config->max_size;
options->min_size = config->min_size;
options->max_depth = config->max_depth;
options->num_threads = num_threads;
options->follow_symlinks = config->follow_symlinks;
options->copy_links = config->copy_links;
options->safe_links = config->safe_links;
options->copy_unsafe_links = config->copy_unsafe_links;
options->checksum = config->checksum;
options->one_file_system = config->one_file_system;
options->file_list = (const FileListSet*)config->files_from_set;
options->base_filters = out->base_filters;
options->per_dir_filters = config->per_dir_filter;
return true;
}
static void prepared_scanner_destroy(PreparedScanner* prepared) {
if (!prepared)
return;
filter_rule_list_free(prepared->base_filters);
prepared->base_filters = NULL;
}
/* --files-from semantics: every listed entry must resolve under the source
* root, otherwise rsync reports a hard error instead of silently transferring
* nothing. An empty list is also an error. An entry of "." (the whole tree)
* and listed-but-empty directories are valid. Runs before any transfer so the
* failure is surfaced uniformly in the single-threaded, -m, dry-run and
* --list-only paths. */
static bool files_from_list_valid(const Config* config) {
const FileListSet* set = (const FileListSet*)config->files_from_set;
if (!set)
return true;
if (!config->send_directory) {
log_message(LOG_LEVEL_ERROR, "--files-from requires a source directory");
return false;
}
if (set->count == 0) {
log_message(LOG_LEVEL_ERROR, "--files-from file '%s' contains no entries; nothing to transfer",
config->files_from ? config->files_from : "");
return false;
}
for (int i = 0; i < set->count; i++) {
const char* entry = set->entries[i];
if (entry[0] == '\0')
continue; /* "." == list the whole tree */
char* full = path_cat(config->send_directory, entry);
if (!full) {
log_message(LOG_LEVEL_ERROR, "memory allocation failed while validating --files-from");
return false;
}
struct stat st;
if (lstat(full, &st) != 0) {
log_message(LOG_LEVEL_ERROR, "--files-from entry '%s' not found in source '%s'", entry,
config->send_directory);
free(full);
return false;
}
free(full);
}
return true;
}
/* Select the configured transport for both transfer execution paths. */
@@ -249,11 +347,17 @@ static void pipeline_cancel(PipelineContextSender* context) {
/* Print dry-run manifest showing files that would be transferred. Returns 0 on success. */
static int send_dry_run_manifest(const Config* config) {
ScannerOptions options = scanner_options_from_config(config, 0);
DirectoryScanner* scanner =
directory_scanner_create_with_options(config->send_directory, &options);
if (!scanner)
if (!files_from_list_valid(config))
return -1;
PreparedScanner prepared;
if (!prepare_scanner(config, 0, &prepared))
return -1;
DirectoryScanner* scanner =
directory_scanner_create_with_options(config->send_directory, &prepared.options);
if (!scanner) {
prepared_scanner_destroy(&prepared);
return -1;
}
Chunk* chunk;
int file_count = 0;
unsigned long long total_bytes = 0;
@@ -267,6 +371,7 @@ static int send_dry_run_manifest(const Config* config) {
if (!escaped_path) {
chunk_destroy(chunk);
directory_scanner_destroy(scanner);
prepared_scanner_destroy(&prepared);
return -1;
}
if (config->human_readable)
@@ -283,6 +388,7 @@ static int send_dry_run_manifest(const Config* config) {
chunk_destroy(chunk);
}
directory_scanner_destroy(scanner);
prepared_scanner_destroy(&prepared);
if (!config->quiet) {
if (config->human_readable)
printf("Total: %d files, %s\n", file_count,
@@ -319,12 +425,18 @@ static int compare_list_entries(const void* left, const void* right) {
* Directory lines are not printed because the scanner only yields regular
* transfer candidates. Returns 0 on success, 1 on error. */
static int send_list_only(const Config* config) {
ScannerOptions options = scanner_options_from_config(config, 0);
options.use_metadata = true; /* capture mode + mtime for the listing */
DirectoryScanner* scanner =
directory_scanner_create_with_options(config->send_directory, &options);
if (!scanner)
if (!files_from_list_valid(config))
return 1;
PreparedScanner prepared;
if (!prepare_scanner(config, 0, &prepared))
return 1;
prepared.options.use_metadata = true; /* capture mode + mtime for the listing */
DirectoryScanner* scanner =
directory_scanner_create_with_options(config->send_directory, &prepared.options);
if (!scanner) {
prepared_scanner_destroy(&prepared);
return 1;
}
ListEntry* entries = NULL;
size_t count = 0;
size_t capacity = 0;
@@ -378,6 +490,7 @@ static int send_list_only(const Config* config) {
}
bool failed = oom || directory_scanner_failed(scanner);
directory_scanner_destroy(scanner);
prepared_scanner_destroy(&prepared);
if (failed) {
list_entries_destroy(entries, count);
if (oom)
@@ -770,14 +883,20 @@ static int send_chunks_multithreaded(void* pipeline_context) {
static int scan_directory_multithreaded(void* pipeline_context) {
PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
protocol_session_bind(&context->allocation_session);
ScannerOptions options = scanner_options_from_config(context->config, 4);
PreparedScanner prepared;
if (!prepare_scanner(context->config, 4, &prepared)) {
pipeline_cancel(context);
protocol_session_unbind();
return thrd_error;
}
ParallelScanner* scanner = parallel_scanner_create_with_options(
context->config->send_directory, &options, &context->allocation_session);
context->config->send_directory, &prepared.options, &context->allocation_session);
Chunk* current_chunk;
if (scanner == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to create parallel scanner");
pipeline_cancel(context);
prepared_scanner_destroy(&prepared);
protocol_session_unbind();
return thrd_error;
}
@@ -790,6 +909,7 @@ static int scan_directory_multithreaded(void* pipeline_context) {
pipeline_cancel(context);
chunk_destroy(current_chunk);
parallel_scanner_destroy(scanner);
prepared_scanner_destroy(&prepared);
protocol_session_unbind();
return thrd_error;
}
@@ -801,12 +921,14 @@ static int scan_directory_multithreaded(void* pipeline_context) {
chunk_destroy(current_chunk);
pipeline_cancel(context);
parallel_scanner_destroy(scanner);
prepared_scanner_destroy(&prepared);
protocol_session_unbind();
return thrd_error;
}
}
if (parallel_scanner_failed(scanner)) {
parallel_scanner_destroy(scanner);
prepared_scanner_destroy(&prepared);
mtx_lock(&context->mutex_scanner);
context->scanner_done = true;
cnd_broadcast(&context->condition_not_empty_scanner);
@@ -822,6 +944,7 @@ static int scan_directory_multithreaded(void* pipeline_context) {
mtx_unlock(&context->mutex_scanner);
parallel_scanner_destroy(scanner);
prepared_scanner_destroy(&prepared);
protocol_session_unbind();
return thrd_success;
}
@@ -923,6 +1046,8 @@ int send_files(Config* config) {
return send_list_only(config);
if (config->dry_run)
return send_dry_run_manifest(config);
if (!files_from_list_valid(config))
return 1;
Client* client = connect_transfer_client(config);
if (!client) {
@@ -939,10 +1064,13 @@ int send_files(Config* config) {
DirectoryScanner* scanner = NULL;
ArrayList* manifest = NULL;
ArrayList* remove_sources = NULL;
PreparedScanner prepared;
memset(&prepared, 0, sizeof(prepared));
if (!config_send(client->file_descriptor, config))
goto send_fail;
ScannerOptions scanner_options = scanner_options_from_config(config, 0);
scanner = directory_scanner_create_with_options(config->send_directory, &scanner_options);
if (!prepare_scanner(config, 0, &prepared))
goto send_fail;
scanner = directory_scanner_create_with_options(config->send_directory, &prepared.options);
manifest = create_transfer_manifest(config);
if (config->remove_source_files)
remove_sources = array_list_create(source_file_destroy);
@@ -1043,6 +1171,7 @@ send_fail:
array_list_delete(remove_sources);
if (scanner)
directory_scanner_destroy(scanner);
prepared_scanner_destroy(&prepared);
disconnect_transfer_client(client);
protocol_session_unbind();
return ret;
@@ -1056,6 +1185,8 @@ int send_files_multithreaded(Config** config_ptr) {
return send_list_only(config);
if (config->dry_run)
return send_dry_run_manifest(config);
if (!files_from_list_valid(config))
return 1;
long pages = sysconf(_SC_AVPHYS_PAGES);
long page_size = sysconf(_SC_PAGE_SIZE);
+296 -24
View File
@@ -17,8 +17,62 @@
typedef struct {
char* path;
int depth;
FilterNode* context; /* inherited per-directory filter context */
} DirEntry;
/* A chain node: `own` holds the .rsync-filter rules of one directory, `parent`
* the context that directory inherited (nearest ancestor with a filter file).
* The chain for a directory's contents runs from that directory's own node up
* to the root; the command-line base rules are evaluated after the whole
* chain. */
struct FilterNode {
FilterNode* parent;
FilterRuleList* own;
};
static void filter_node_destroy(void* item) {
if (item) {
FilterNode* node = (FilterNode*)item;
if (node->own)
filter_rule_list_free(node->own);
free(node);
}
}
static FilterNode* filter_node_alloc(FilterNode* parent, FilterRuleList* own) {
FilterNode* node = malloc(sizeof(FilterNode));
if (!node)
return NULL;
node->parent = parent;
node->own = own;
return node;
}
/* Evaluate a rule chain for an entry inside the directory whose content
* context is `node`. rsync precedence, highest first: the innermost (current)
* directory's .rsync-filter rules, then each ancestor's, then the root's, and
* finally the command-line base rules (--filter/-C). A deeper per-directory
* file therefore overrides a shallower one, and per-directory files override
* the base rules by default. Returns FILTER_ACTION_NONE when nothing matched. */
static FilterAction chain_rules_apply(const FilterRuleList* base, const FilterNode* node,
const char* rel, const char* leaf, bool is_dir) {
if (node) {
FilterAction own_action = filter_rules_apply(node->own, rel, leaf, is_dir);
if (own_action != FILTER_ACTION_NONE)
return own_action;
return chain_rules_apply(base, node->parent, rel, leaf, is_dir);
}
return base ? filter_rules_apply(base, rel, leaf, is_dir) : FILTER_ACTION_NONE;
}
static bool entry_allowed(const FilterRuleList* base, const FilterNode* node, const char* rel,
const char* leaf, bool is_dir, bool per_dir_filters) {
/* -F: per-directory .rsync-filter files are never transferred. */
if (per_dir_filters && !is_dir && strcmp(leaf, ".rsync-filter") == 0)
return false;
return chain_rules_apply(base, node, rel, leaf, is_dir) != FILTER_ACTION_EXCLUDE;
}
static void dir_entry_destroy(void* item) {
if (item) {
DirEntry* de = (DirEntry*)item;
@@ -27,7 +81,7 @@ static void dir_entry_destroy(void* item) {
}
}
static DirEntry* dir_entry_create(const char* path, int depth) {
static DirEntry* dir_entry_create(const char* path, int depth, FilterNode* context) {
DirEntry* de = malloc(sizeof(DirEntry));
if (!de)
return NULL;
@@ -37,6 +91,7 @@ static DirEntry* dir_entry_create(const char* path, int depth) {
return NULL;
}
de->depth = depth;
de->context = context;
return de;
}
@@ -66,6 +121,79 @@ bool scanner_same_filesystem(bool one_file_system, dev_t root_device, dev_t entr
return !one_file_system || entry_device == root_device;
}
/* Relative path of an on-disk path below `root`. The transfer root may be
* given with a trailing slash; the returned rel path never has one and is ""
* for the root itself. A root of "/" is handled (its children start at "/").
* Exposed so tests can exercise the mapping directly. */
char* scanner_path_relative(const char* root, const char* fs_path) {
size_t root_len = strlen(root);
while (root_len > 1 && root[root_len - 1] == '/')
root_len--;
if (strncmp(root, fs_path, root_len) != 0)
return NULL;
if (root_len == 1 && root[0] == '/') {
if (fs_path[1] == '\0')
return str_dup("");
return str_dup(fs_path + 1);
}
if (fs_path[root_len] == '\0')
return str_dup("");
if (fs_path[root_len] != '/')
return NULL;
return str_dup(fs_path + root_len + 1);
}
/* Relative path of a child entry below the current directory. */
static char* child_rel_path(const char* parent_rel, const char* name) {
if (!parent_rel || parent_rel[0] == '\0')
return str_dup(name);
return path_cat(parent_rel, name);
}
/* Apply the --files-from allow-set and the filter layer to one entry. */
static bool entry_passes_selection(const FileListSet* file_list, const FilterRuleList* base,
const FilterNode* node, const char* rel, const char* leaf,
bool is_dir, bool per_dir_filters) {
if (file_list && !file_list_affects(file_list, rel))
return false;
if (base || per_dir_filters)
return entry_allowed(base, node, rel, leaf, is_dir, per_dir_filters);
return true;
}
/* Merge the open directory's own .rsync-filter rules into the inherited
* context, returning the context used for this directory's entries. On a parse
* error the scanner is marked failed. Returns 0 on success, -1 on failure. */
static int open_directory_filter_context(DirectoryScanner* scanner, const FilterNode* inherited) {
if (!scanner->per_dir_filters) {
scanner->current_node = (FilterNode*)inherited;
return 0;
}
char err[256];
bool exists = false;
FilterRuleList* own =
filter_file_read(scanner->current_path, scanner->current_rel ? scanner->current_rel : "",
&exists, err, sizeof(err));
if (!own) {
log_message(LOG_LEVEL_ERROR, "invalid .rsync-filter in %s: %s", scanner->current_path, err);
scanner->failed = true;
return -1;
}
if (exists && own->count > 0) {
FilterNode* node = filter_node_alloc((FilterNode*)inherited, own);
if (!node || !array_list_add(scanner->filter_nodes, node)) {
filter_node_destroy(node);
scanner->failed = true;
return -1;
}
scanner->current_node = node;
} else {
filter_rule_list_free(own);
scanner->current_node = (FilterNode*)inherited;
}
return 0;
}
/* Inspect symlinks, resolve the entry type, and apply file filters once for both scanners. */
static int scanner_inspect_entry(const ScannerOptions* options, const char* source_root,
const char* containing_dir, const char* name,
@@ -165,25 +293,54 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo
scanner->checksum = options->checksum;
scanner->one_file_system = options->one_file_system;
scanner->failed = false;
scanner->root_path = str_dup(root_directory);
if (!scanner->root_path) {
queue_destroy(scanner->directories);
free(scanner);
return NULL;
}
scanner->current_rel = NULL;
scanner->at_seed_dir = true;
scanner->seed_node = NULL;
scanner->current_node = NULL;
scanner->file_list = options->file_list;
scanner->base_filters = options->base_filters;
scanner->per_dir_filters = options->per_dir_filters;
scanner->filter_nodes = NULL;
if (scanner->base_filters || scanner->per_dir_filters) {
scanner->filter_nodes = array_list_create(filter_node_destroy);
if (!scanner->filter_nodes) {
free(scanner->root_path);
queue_destroy(scanner->directories);
free(scanner);
return NULL;
}
}
if (scanner->one_file_system) {
struct stat root_stats;
if (stat(root_directory, &root_stats) != 0) {
log_perror("Could not stat source directory");
free(scanner->root_path);
queue_destroy(scanner->directories);
array_list_delete(scanner->filter_nodes);
free(scanner);
return NULL;
}
scanner->root_dev = root_stats.st_dev;
}
DirEntry* root = dir_entry_create(root_directory, 0);
DirEntry* root = dir_entry_create(root_directory, 0, NULL);
if (!root) {
free(scanner->root_path);
queue_destroy(scanner->directories);
array_list_delete(scanner->filter_nodes);
free(scanner);
return NULL;
}
if (!queue_enqueue(scanner->directories, root)) {
dir_entry_destroy(root);
free(scanner->root_path);
queue_destroy(scanner->directories);
array_list_delete(scanner->filter_nodes);
free(scanner);
return NULL;
}
@@ -197,14 +354,25 @@ DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_
unsigned long long min_size, int max_depth,
bool follow_symlinks, bool copy_links, bool safe_links,
bool copy_unsafe_links, bool checksum) {
ScannerOptions options = {use_metadata, chunk_size,
exclude_patterns, exclude_count,
include_patterns, include_count,
max_size, min_size,
max_depth, 0,
follow_symlinks, copy_links,
safe_links, copy_unsafe_links,
checksum, false};
ScannerOptions options = {use_metadata,
chunk_size,
exclude_patterns,
exclude_count,
include_patterns,
include_count,
max_size,
min_size,
max_depth,
0,
follow_symlinks,
copy_links,
safe_links,
copy_unsafe_links,
checksum,
false,
NULL,
NULL,
false};
return directory_scanner_create_with_options(root_directory, &options);
}
@@ -216,6 +384,9 @@ void directory_scanner_destroy(DirectoryScanner* scanner) {
scanner->current_dir = NULL;
}
free(scanner->current_path);
free(scanner->current_rel);
free(scanner->root_path);
array_list_delete(scanner->filter_nodes);
queue_destroy(scanner->directories);
free(scanner);
}
@@ -246,7 +417,21 @@ static int open_next_directory(DirectoryScanner* scanner) {
DirEntry* de = (DirEntry*)queue_dequeue(scanner->directories);
scanner->current_path = de->path;
scanner->current_depth = de->depth;
/* The seed directory inherits the scanner's configured context (the root
* .rsync-filter context in parallel mode); other dirs inherit the context of
* the directory that enqueued them. */
const FilterNode* inherited = scanner->at_seed_dir ? scanner->seed_node : de->context;
scanner->at_seed_dir = false;
free(de);
free(scanner->current_rel);
scanner->current_rel = scanner_path_relative(scanner->root_path, scanner->current_path);
if (!scanner->current_rel) {
log_message(LOG_LEVEL_ERROR, "Could not compute relative path under %s", scanner->root_path);
scanner->failed = true;
return -1;
}
scanner->current_dir = opendir(scanner->current_path);
if (scanner->current_dir == NULL) {
log_perror("Could not open directory");
@@ -255,6 +440,11 @@ static int open_next_directory(DirectoryScanner* scanner) {
scanner->failed = true;
return -1;
}
if (open_directory_filter_context(scanner, inherited) != 0) {
closedir(scanner->current_dir);
scanner->current_dir = NULL;
return -1;
}
return 1;
}
@@ -294,7 +484,9 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
scanner->max_depth, 0,
scanner->follow_symlinks, scanner->copy_links,
scanner->safe_links, scanner->copy_unsafe_links,
scanner->checksum, scanner->one_file_system};
scanner->checksum, scanner->one_file_system,
scanner->file_list, scanner->base_filters,
scanner->per_dir_filters};
ScannerEntry inspected;
int inspection = scanner_inspect_entry(&options, scanner->current_path, scanner->current_path,
entry->d_name, &inspected);
@@ -307,14 +499,32 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
char* cur_path = inspected.path;
struct stat stats = inspected.stats;
if (inspected.is_directory) {
/* --files-from allow-set and the filter layer apply to files and to
* directories (an excluded directory is not descended into). */
bool is_dir = inspected.is_directory;
char* rel = child_rel_path(scanner->current_rel, entry->d_name);
if (!rel) {
free(cur_path);
scanner->failed = true;
break;
}
bool passes_selection =
entry_passes_selection(scanner->file_list, scanner->base_filters, scanner->current_node,
rel, entry->d_name, is_dir, scanner->per_dir_filters);
free(rel);
if (!passes_selection) {
free(cur_path);
continue;
}
if (is_dir) {
if (!scanner_same_filesystem(scanner->one_file_system, scanner->root_dev, stats.st_dev)) {
free(cur_path);
continue;
}
int next_depth = scanner->current_depth + 1;
if (scanner->max_depth <= 0 || next_depth < scanner->max_depth) {
DirEntry* de = dir_entry_create(cur_path, next_depth);
DirEntry* de = dir_entry_create(cur_path, next_depth, scanner->current_node);
if (!de || !queue_enqueue(scanner->directories, de)) {
dir_entry_destroy(de);
scanner->failed = true;
@@ -376,6 +586,7 @@ typedef struct {
ParallelScanner* ps;
char** dirs;
int dir_count;
char* root_dir; /* the transfer root, for relative-path computation */
ScannerOptions options;
ProtocolSession* allocation_session;
} ParallelWorkerArg;
@@ -398,6 +609,13 @@ static int parallel_worker_thread(void* arg) {
free(wa->dirs[j]);
break;
}
/* Root .rsync-filter rules (parsed by the parallel scanner) apply to the
* contents of every assigned subdirectory. Relative paths (used by the
* allow-set and per-directory rules) are computed against the transfer
* root, not the subdirectory the worker is seeded with. */
free(ds->root_path);
ds->root_path = str_dup(wa->root_dir);
ds->seed_node = wa->ps->root_filter_node;
Chunk* chunk;
while ((chunk = directory_scanner_next(ds)) != NULL) {
if (!queue_enqueue_multithreaded_cancel(wa->ps->result_queue, chunk, &wa->ps->result_mutex,
@@ -419,6 +637,7 @@ static int parallel_worker_thread(void* arg) {
free(wa->dirs[i]);
}
ParallelScanner* ps = wa->ps;
free(wa->root_dir);
free(wa->dirs);
free(wa);
mtx_lock(&ps->result_mutex);
@@ -549,9 +768,10 @@ static Chunk* batch_files(ArrayList* files, unsigned long long chunk_size, Queue
}
/* Scan one root-directory entry into either the subdirs or files list. */
static void scan_root_entry(const ScannerOptions* options, const char* root_directory,
const struct dirent* entry, ArrayList* root_files, ArrayList* subdirs,
dev_t root_dev, ParallelScanner* ps) {
static void scan_root_entry(const ScannerOptions* options, const FilterNode* root_node,
const char* root_directory, const struct dirent* entry,
ArrayList* root_files, ArrayList* subdirs, dev_t root_dev,
ParallelScanner* ps) {
ScannerEntry inspected;
int inspection =
scanner_inspect_entry(options, root_directory, root_directory, entry->d_name, &inspected);
@@ -563,7 +783,21 @@ static void scan_root_entry(const ScannerOptions* options, const char* root_dire
return;
char* cur_path = inspected.path;
struct stat st = inspected.stats;
if (inspected.is_directory) {
bool is_dir = inspected.is_directory;
char* rel = str_dup(entry->d_name);
if (!rel) {
free(cur_path);
ps->failed = true;
return;
}
bool passes = entry_passes_selection(options->file_list, options->base_filters, root_node, rel,
entry->d_name, is_dir, options->per_dir_filters);
free(rel);
if (!passes) {
free(cur_path);
return;
}
if (is_dir) {
if (!scanner_same_filesystem(options->one_file_system, root_dev, st.st_dev)) {
free(cur_path);
return;
@@ -597,8 +831,8 @@ static void scan_root_entry(const ScannerOptions* options, const char* root_dire
/* Scan the root directory itself, collecting root files and subdirectories.
* Returns false if the root directory could not be opened. */
static bool scan_root_directory(ParallelScanner* ps, const char* root_directory,
const ScannerOptions* options, dev_t root_dev,
ArrayList* root_files, ArrayList* subdirs) {
const ScannerOptions* options, const FilterNode* root_node,
dev_t root_dev, ArrayList* root_files, ArrayList* subdirs) {
DIR* dir = opendir(root_directory);
if (!dir) {
log_perror("Could not open root directory for parallel scan");
@@ -608,7 +842,7 @@ static bool scan_root_directory(ParallelScanner* ps, const char* root_directory,
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue;
scan_root_entry(options, root_directory, entry, root_files, subdirs, root_dev, ps);
scan_root_entry(options, root_node, root_directory, entry, root_files, subdirs, root_dev, ps);
}
closedir(dir);
return true;
@@ -616,7 +850,8 @@ static bool scan_root_directory(ParallelScanner* ps, const char* root_directory,
/* Spawn worker threads, one per group of subdirectories. */
static void spawn_parallel_workers(ParallelScanner* ps, ArrayList* subdirs,
const ScannerOptions* options, unsigned long long cs) {
const ScannerOptions* options, const char* root_directory,
unsigned long long cs) {
if (subdirs->size <= 0)
return;
int n = options->num_threads > 0 ? options->num_threads : 4;
@@ -647,7 +882,10 @@ static void spawn_parallel_workers(ParallelScanner* ps, ArrayList* subdirs,
}
wa->ps = ps;
wa->dirs = calloc(count, sizeof(char*));
if (!wa->dirs) {
wa->root_dir = str_dup(root_directory);
if (!wa->dirs || !wa->root_dir) {
free(wa->root_dir);
free(wa->dirs);
free(wa);
parallel_scanner_creation_failed(ps);
break;
@@ -661,6 +899,7 @@ static void spawn_parallel_workers(ParallelScanner* ps, ArrayList* subdirs,
if (!dup_ok) {
for (int j = 0; j < count; j++)
free(wa->dirs[j]);
free(wa->root_dir);
free(wa->dirs);
free(wa);
parallel_scanner_creation_failed(ps);
@@ -674,6 +913,7 @@ static void spawn_parallel_workers(ParallelScanner* ps, ArrayList* subdirs,
if (thrd_create(&ps->threads[t], parallel_worker_thread, wa) != thrd_success) {
for (int j = 0; j < count; j++)
free(wa->dirs[j]);
free(wa->root_dir);
free(wa->dirs);
free(wa);
parallel_scanner_creation_failed(ps);
@@ -720,7 +960,37 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
root_dev = root_stats.st_dev;
}
if (!scan_root_directory(ps, root_directory, options, root_dev, root_files, subdirs)) {
/* Build the root directory's .rsync-filter context once; workers seed their
* scanners with it so per-dir rules behave identically to the sequential
* scanner. */
FilterNode* root_node = NULL;
if (options->per_dir_filters) {
char err[256];
bool exists = false;
FilterRuleList* own = filter_file_read(root_directory, "", &exists, err, sizeof(err));
if (!own) {
log_message(LOG_LEVEL_ERROR, "invalid .rsync-filter in %s: %s", root_directory, err);
array_list_delete(root_files);
array_list_delete(subdirs);
parallel_scanner_destroy(ps);
return NULL;
}
if (exists && own->count > 0) {
root_node = filter_node_alloc(NULL, own);
if (!root_node) {
filter_rule_list_free(own);
array_list_delete(root_files);
array_list_delete(subdirs);
parallel_scanner_destroy(ps);
return NULL;
}
} else {
filter_rule_list_free(own);
}
}
ps->root_filter_node = root_node;
if (!scan_root_directory(ps, root_directory, options, root_node, root_dev, root_files, subdirs)) {
array_list_delete(root_files);
array_list_delete(subdirs);
parallel_scanner_destroy(ps);
@@ -731,7 +1001,7 @@ ParallelScanner* parallel_scanner_create_with_options(const char* root_directory
ps->initial_chunk = batch_files(root_files, cs, ps->result_queue, &ps->failed);
array_list_delete(root_files);
spawn_parallel_workers(ps, subdirs, options, cs);
spawn_parallel_workers(ps, subdirs, options, root_directory, cs);
array_list_delete(subdirs);
return ps;
}
@@ -774,6 +1044,8 @@ void parallel_scanner_destroy(ParallelScanner* ps) {
for (int i = 0; i < ps->num_threads; i++)
thrd_join(ps->threads[i], NULL);
free(ps->threads);
if (ps->root_filter_node)
filter_node_destroy(ps->root_filter_node);
if (ps->initial_chunk)
chunk_destroy(ps->initial_chunk);
queue_destroy(ps->result_queue);
+28
View File
@@ -2,6 +2,8 @@
#define SCANNER_H
#include "chunk.h"
#include "file_list.h"
#include "filter.h"
#include "protocol.h"
#include "queue.h"
#include <dirent.h>
@@ -27,8 +29,18 @@ typedef struct {
bool copy_unsafe_links;
bool checksum;
bool one_file_system;
/* Phase 2 (files-from / filter layer). All pointers are shared read-only
* across scanner instances and worker threads; ownership stays with the
* caller (client_send). */
const FileListSet* file_list; /* --files-from allow-set, or NULL */
const FilterRuleList* base_filters; /* command-line + -C rules, or NULL */
bool per_dir_filters; /* -F: read .rsync-filter per directory */
} ScannerOptions;
/* Internal per-scanner filter state. FilterNode chains represent the ordered
* per-directory .rsync-filter rules that apply below a directory. */
typedef struct FilterNode FilterNode;
typedef struct {
Queue* directories;
DIR* current_dir;
@@ -51,6 +63,16 @@ typedef struct {
bool one_file_system;
dev_t root_dev;
bool failed;
/* Phase 2 (files-from / filter layer). */
char* root_path; /* transfer root (fs path) for rel computation */
char* current_rel; /* rel path of the open directory ("" == root) */
bool at_seed_dir; /* next open is the seed directory */
FilterNode* seed_node; /* inherited context of the seed dir, or NULL */
FilterNode* current_node; /* filter context of the open directory */
ArrayList* filter_nodes; /* owned FilterNode arena (may be NULL) */
const FileListSet* file_list;
const FilterRuleList* base_filters;
bool per_dir_filters;
} DirectoryScanner;
typedef struct {
@@ -68,6 +90,7 @@ typedef struct {
int completed;
Chunk* initial_chunk;
ProtocolSession* allocation_session;
FilterNode* root_filter_node; /* root .rsync-filter context (owned by ps) */
} ParallelScanner;
DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_metadata,
@@ -88,6 +111,11 @@ void directory_scanner_destroy(DirectoryScanner* scanner);
* the transfer root. Exposed so tests can exercise the rule directly. */
bool scanner_same_filesystem(bool one_file_system, dev_t root_device, dev_t entry_device);
/* Relative path of an on-disk path below `root` ("" == the root itself, NULL
* when `fs_path` is not under `root`). Handles trailing slashes and a root of
* "/". Exposed so tests can exercise the mapping directly. */
char* scanner_path_relative(const char* root, const char* fs_path);
ParallelScanner* parallel_scanner_create_with_options(const char* root_directory,
const ScannerOptions* options,
ProtocolSession* allocation_session);
+7
View File
@@ -33,6 +33,13 @@ void print_usage(void) {
printf(" --include <pattern> Only include files matching pattern\n");
printf(" --exclude-from <file> Read exclude patterns from file\n");
printf(" --include-from <file> Read include patterns from file\n");
printf(" --files-from <file> Read the source file list from FILE (paths relative to the "
"source root)\n");
printf(" -0, --from0 Entries in --files-from are NUL-delimited\n");
printf(" --filter=RULE rsync-style filter rule (+/- include/exclude; repeatable; the\n");
printf(" rsync -f short form conflicts with FastSync sendfile -f)\n");
printf(" -C, --cvs-exclude Auto-ignore common CVS/SCM files (.git/, .svn/, *.o, *~, ...)\n");
printf(" -F Apply per-directory .rsync-filter files during the scan\n");
printf(" --max-size <n> Skip files larger than n bytes\n");
printf(" --min-size <n> Skip files smaller than n bytes\n");
printf(" --max-alloc <SIZE> Maximum single allocation (default: 1G)\n");