Release v2.29.0 #312
@@ -139,6 +139,8 @@ set(CLIENT_CORE_SRCS
|
|||||||
src/client/client_send.c
|
src/client/client_send.c
|
||||||
src/client/client_validation.c
|
src/client/client_validation.c
|
||||||
src/client/scanner.c
|
src/client/scanner.c
|
||||||
|
src/client/scanner_filter.c
|
||||||
|
src/client/scanner_parallel.c
|
||||||
src/client/usage.c
|
src/client/usage.c
|
||||||
)
|
)
|
||||||
set(CLIENT_MAIN_SRCS src/client/client_cli.c)
|
set(CLIENT_MAIN_SRCS src/client/client_cli.c)
|
||||||
|
|||||||
+306
-1623
File diff suppressed because it is too large
Load Diff
@@ -0,0 +1,671 @@
|
|||||||
|
#include "log.h"
|
||||||
|
#include "scanner.h"
|
||||||
|
#include "scanner_internal.h"
|
||||||
|
#include "array_list.h"
|
||||||
|
#include "chunk.h"
|
||||||
|
#include "file.h"
|
||||||
|
#include "queue.h"
|
||||||
|
#include "utils.h"
|
||||||
|
#include <dirent.h>
|
||||||
|
#include <stdio.h>
|
||||||
|
#include <stdlib.h>
|
||||||
|
#include <string.h>
|
||||||
|
#include <sys/stat.h>
|
||||||
|
#include <sys/sysmacros.h>
|
||||||
|
#include <threads.h>
|
||||||
|
#include <unistd.h>
|
||||||
|
#include <limits.h>
|
||||||
|
|
||||||
|
#include "xattr.h"
|
||||||
|
|
||||||
|
/* 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;
|
||||||
|
};
|
||||||
|
|
||||||
|
void filter_node_destroy(void* item) {
|
||||||
|
if (item) {
|
||||||
|
FilterNode* node = (FilterNode*)item;
|
||||||
|
if (node->own)
|
||||||
|
filter_rule_list_free(node->own);
|
||||||
|
free(node);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
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 one entry. 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). The
|
||||||
|
* sender-side verdict decides whether the entry is hidden from the transfer;
|
||||||
|
* the receiver-side verdict decides whether its destination mirror is protected
|
||||||
|
* from --delete. Each side takes the FIRST matching rule independently. */
|
||||||
|
typedef struct {
|
||||||
|
bool hide; /* sender-side exclude matched */
|
||||||
|
bool protect; /* receiver-side exclude matched */
|
||||||
|
} FilterOutcome;
|
||||||
|
|
||||||
|
static void chain_rules_outcome(const FilterRuleList* base, const FilterNode* node, const char* rel,
|
||||||
|
const char* leaf, bool is_dir, FilterOutcome* out) {
|
||||||
|
memset(out, 0, sizeof(*out));
|
||||||
|
bool sender_decided = false;
|
||||||
|
bool receiver_decided = false;
|
||||||
|
const FilterNode* n = node;
|
||||||
|
while (!sender_decided || !receiver_decided) {
|
||||||
|
const FilterRuleList* list = n ? n->own : base;
|
||||||
|
if (list) {
|
||||||
|
if (!sender_decided) {
|
||||||
|
FilterAction action = filter_rules_apply_side(list, rel, leaf, is_dir, FILTER_SIDE_SENDER);
|
||||||
|
if (action != FILTER_ACTION_NONE) {
|
||||||
|
out->hide = action == FILTER_ACTION_EXCLUDE;
|
||||||
|
sender_decided = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (!receiver_decided) {
|
||||||
|
FilterAction action =
|
||||||
|
filter_rules_apply_side(list, rel, leaf, is_dir, FILTER_SIDE_RECEIVER);
|
||||||
|
if (action != FILTER_ACTION_NONE) {
|
||||||
|
out->protect = action == FILTER_ACTION_PROTECT;
|
||||||
|
receiver_decided = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (!n)
|
||||||
|
break;
|
||||||
|
n = n->parent;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
static bool entry_allowed(const FilterRuleList* base, const FilterNode* node, const char* rel,
|
||||||
|
const char* leaf, bool is_dir, bool exclude_filter_files,
|
||||||
|
bool* protect_out) {
|
||||||
|
/* -FF: per-directory .rsync-filter files are never transferred (single -F
|
||||||
|
transfers them, matching rsync). */
|
||||||
|
if (exclude_filter_files && !is_dir && strcmp(leaf, ".rsync-filter") == 0) {
|
||||||
|
if (protect_out)
|
||||||
|
*protect_out = false;
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
FilterOutcome outcome;
|
||||||
|
chain_rules_outcome(base, node, rel, leaf, is_dir, &outcome);
|
||||||
|
if (protect_out)
|
||||||
|
*protect_out = outcome.protect;
|
||||||
|
return !outcome.hide;
|
||||||
|
}
|
||||||
|
|
||||||
|
void dir_entry_destroy(void* item) {
|
||||||
|
if (item) {
|
||||||
|
DirEntry* de = (DirEntry*)item;
|
||||||
|
free(de->path);
|
||||||
|
free(de);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
DirEntry* dir_entry_create(const char* path, int depth, FilterNode* context) {
|
||||||
|
DirEntry* de = malloc(sizeof(DirEntry));
|
||||||
|
if (!de)
|
||||||
|
return NULL;
|
||||||
|
de->path = str_dup(path);
|
||||||
|
if (!de->path) {
|
||||||
|
free(de);
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
de->depth = depth;
|
||||||
|
de->context = context;
|
||||||
|
return de;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Apply rsync's symlink-resolution precedence to one S_ISLNK entry:
|
||||||
|
* --copy-links dereferences every symlink;
|
||||||
|
* --copy-unsafe-links dereferences only targets unsafe_symlink() flags;
|
||||||
|
* -k/--copy-dirlinks dereferences only a symlink whose referent is a dir;
|
||||||
|
* --safe-links (receiver-side in rsync; modelled here) ignores an unsafe
|
||||||
|
* target that would otherwise be carried; with --munge-links
|
||||||
|
* every stored target becomes absolute, so --safe-links then
|
||||||
|
* ignores every symlink, exactly as rsync documents;
|
||||||
|
* -l/--links carries the link.
|
||||||
|
* `link_rel` is the symlink's transfer-relative path (incl. name) and is used
|
||||||
|
* only for the lexical unsafe test. `target` receives the raw link value. */
|
||||||
|
LinkAction scanner_link_action(const ScannerOptions* options, const char* path,
|
||||||
|
const char* link_rel, char* target, size_t target_size) {
|
||||||
|
if (!options->follow_symlinks && !options->copy_links && !options->safe_links &&
|
||||||
|
!options->copy_unsafe_links && !options->copy_dirlinks)
|
||||||
|
return LINK_ACTION_SKIP;
|
||||||
|
ssize_t length = readlink(path, target, target_size - 1);
|
||||||
|
if (length < 0)
|
||||||
|
return LINK_ACTION_SKIP;
|
||||||
|
target[length] = '\0';
|
||||||
|
|
||||||
|
bool unsafe = file_symlink_unsafe(target, link_rel);
|
||||||
|
if (options->copy_links || (options->copy_unsafe_links && unsafe))
|
||||||
|
return LINK_ACTION_DEREF;
|
||||||
|
if (options->copy_dirlinks) {
|
||||||
|
struct stat ref;
|
||||||
|
if (stat(path, &ref) == 0 && S_ISDIR(ref.st_mode))
|
||||||
|
return LINK_ACTION_DEREF;
|
||||||
|
}
|
||||||
|
if (options->safe_links && (unsafe || options->munge_links))
|
||||||
|
return LINK_ACTION_SKIP_PROTECTED;
|
||||||
|
if (!options->follow_symlinks || target[0] == '\0')
|
||||||
|
return LINK_ACTION_SKIP;
|
||||||
|
return LINK_ACTION_CARRY;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* --one-file-system (-x) decision. Only directories can carry a different
|
||||||
|
* device than their parent (mount points), so this is checked when a child
|
||||||
|
* directory is about to be descended into. */
|
||||||
|
bool scanner_same_filesystem(int one_file_system, dev_t root_device, dev_t entry_device) {
|
||||||
|
return one_file_system <= 0 || entry_device == root_device;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Build a payload-less directory File carrying the captured metadata (when
|
||||||
|
* requested). Used by -x mount-point emission and --list-only directory
|
||||||
|
* entries. Returns NULL on allocation failure. */
|
||||||
|
File* scanner_build_dir_file(const char* path, const struct stat* stats,
|
||||||
|
const ScannerOptions* options) {
|
||||||
|
File* dir = file_create(path);
|
||||||
|
if (dir == NULL)
|
||||||
|
return NULL;
|
||||||
|
dir->is_dir = true;
|
||||||
|
if (options->use_metadata) {
|
||||||
|
dir->metadata =
|
||||||
|
file_metadata_create(dir->path, stats, options->preserve_atimes, options->preserve_crtimes);
|
||||||
|
if (!dir->metadata) {
|
||||||
|
file_destroy(dir);
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return dir;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* 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);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* -R/--relative destination-relative prefix reconstructed from a source spec:
|
||||||
|
* everything after the first '.' path component (rsync's '/./' cut point),
|
||||||
|
* with leading/trailing slashes removed; or the whole spec (normalized) when
|
||||||
|
* there is no cut. Returns "" for the receive root. Exposed for tests. */
|
||||||
|
char* scanner_relative_prefix(const char* spec) {
|
||||||
|
if (!spec || spec[0] == '\0')
|
||||||
|
return NULL;
|
||||||
|
const char* after = spec;
|
||||||
|
if (spec[0] == '.' && spec[1] == '/') {
|
||||||
|
after = spec + 2;
|
||||||
|
} else {
|
||||||
|
const char* cut = strstr(spec, "/./");
|
||||||
|
if (cut)
|
||||||
|
after = cut + 3;
|
||||||
|
}
|
||||||
|
size_t cap = strlen(spec) + 1;
|
||||||
|
char* out = malloc(cap);
|
||||||
|
if (!out)
|
||||||
|
return NULL;
|
||||||
|
size_t len = 0;
|
||||||
|
for (const char* s = after; *s;) {
|
||||||
|
while (*s == '/')
|
||||||
|
s++;
|
||||||
|
const char* comp = s;
|
||||||
|
while (*s && *s != '/')
|
||||||
|
s++;
|
||||||
|
size_t clen = (size_t)(s - comp);
|
||||||
|
if (clen == 0 || (clen == 1 && comp[0] == '.'))
|
||||||
|
continue;
|
||||||
|
if (len)
|
||||||
|
out[len++] = '/';
|
||||||
|
memcpy(out + len, comp, clen);
|
||||||
|
len += clen;
|
||||||
|
}
|
||||||
|
out[len] = '\0';
|
||||||
|
return out;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Relative path of a child entry below the current directory. */
|
||||||
|
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);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Destination-relative wire path for an entry under an -R prefix. */
|
||||||
|
char* scanner_prefix_send_path(const char* prefix, const char* rel) {
|
||||||
|
if (prefix[0] == '\0')
|
||||||
|
return str_dup(rel);
|
||||||
|
if (rel[0] == '\0')
|
||||||
|
return str_dup(prefix);
|
||||||
|
return path_cat(prefix, rel);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Apply the --files-from allow-set and the filter layer to one entry. On
|
||||||
|
* return `*protect_out` is true when a receiver-side rule protects the entry's
|
||||||
|
* destination mirror from deletion. */
|
||||||
|
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, bool exclude_filter_files, bool* protect_out) {
|
||||||
|
if (protect_out)
|
||||||
|
*protect_out = false;
|
||||||
|
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, exclude_filter_files, protect_out);
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Best-effort capture of the file's whitelisted xattrs (-X/-A). A failure to
|
||||||
|
* read xattrs is non-fatal: the file is transferred without them. */
|
||||||
|
void scanner_capture_xattrs(const DirectoryScanner* scanner, File* file) {
|
||||||
|
if (!scanner || !file || !(scanner->options.preserve_xattrs || scanner->options.preserve_acls))
|
||||||
|
return;
|
||||||
|
file->xattrs = xattr_capture_path(file->path, scanner->options.preserve_acls);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Apply --hard-links (-H) detection to one regular File. On a sibling (a
|
||||||
|
* later member of an already-seen source inode) the File keeps the group id
|
||||||
|
* and the first member's wire path but carries NO data payload (size 0); the
|
||||||
|
* first member is left untouched (data present, link_first). Allocation
|
||||||
|
* failure is fatal: the scanner is marked failed. */
|
||||||
|
void scanner_assign_hardlink(DirectoryScanner* scanner, HardLinkTable* table, File* file,
|
||||||
|
const struct stat* stats) {
|
||||||
|
if (!table || !file || !stats)
|
||||||
|
return;
|
||||||
|
int gid;
|
||||||
|
bool is_first;
|
||||||
|
char* first_path = NULL;
|
||||||
|
if (!hardlink_table_assign(table, file_wire_path(file), stats->st_dev, stats->st_ino, &gid,
|
||||||
|
&is_first, &first_path)) {
|
||||||
|
if (scanner)
|
||||||
|
scanner->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
file->link_group = gid;
|
||||||
|
file->link_first = is_first;
|
||||||
|
if (!is_first) {
|
||||||
|
file->hardlink_target = first_path;
|
||||||
|
file->data->size = 0;
|
||||||
|
} else {
|
||||||
|
free(first_path);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Phase 4 special/devices decision for one non-regular entry, matching rsync:
|
||||||
|
- a char/block device is RECREATED as a node under -D/--devices, unless
|
||||||
|
--copy-devices asks for its content to be copied into a regular file;
|
||||||
|
- a FIFO/socket is RECREATED under --specials;
|
||||||
|
- when the matching flag is absent the entry is SKIPPED ("skipping
|
||||||
|
non-regular file"), exactly like rsync's default, instead of being
|
||||||
|
silently copied as a zero-length regular file;
|
||||||
|
- anything else (regular/directory) is left to the normal data path. */
|
||||||
|
ScannerSpecial scanner_prepare_special(bool preserve_devices, bool preserve_specials,
|
||||||
|
bool copy_devices, File* file, const struct stat* stats) {
|
||||||
|
if (!file || !stats)
|
||||||
|
return SCANNER_SPECIAL_REGULAR;
|
||||||
|
bool is_device = S_ISCHR(stats->st_mode) || S_ISBLK(stats->st_mode);
|
||||||
|
bool is_fifo = S_ISFIFO(stats->st_mode);
|
||||||
|
bool is_socket = S_ISSOCK(stats->st_mode);
|
||||||
|
if (!is_device && !is_fifo && !is_socket)
|
||||||
|
return SCANNER_SPECIAL_REGULAR;
|
||||||
|
if (is_device && copy_devices)
|
||||||
|
return SCANNER_SPECIAL_REGULAR; /* copy device content as a regular file */
|
||||||
|
bool preserve = is_device ? preserve_devices : preserve_specials;
|
||||||
|
if (!preserve)
|
||||||
|
return SCANNER_SPECIAL_SKIP;
|
||||||
|
file->is_special = true;
|
||||||
|
file->data->size = 0;
|
||||||
|
file->data->data = NULL;
|
||||||
|
if (is_device) {
|
||||||
|
file->rdev_major = (int32_t)major(stats->st_rdev);
|
||||||
|
file->rdev_minor = (int32_t)minor(stats->st_rdev);
|
||||||
|
}
|
||||||
|
return SCANNER_SPECIAL_RECREATE;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Append `rel` to the caller's exclusion sink, taking `mtx` when shared across
|
||||||
|
parallel worker threads. Returns false on allocation failure (list left
|
||||||
|
unchanged). */
|
||||||
|
bool excluded_sink_append(ArrayList* list, mtx_t* mtx, const char* rel) {
|
||||||
|
if (!list)
|
||||||
|
return true;
|
||||||
|
char* dup = str_dup(rel);
|
||||||
|
if (!dup)
|
||||||
|
return false;
|
||||||
|
if (mtx)
|
||||||
|
mtx_lock(mtx);
|
||||||
|
bool ok = array_list_add(list, dup);
|
||||||
|
if (mtx)
|
||||||
|
mtx_unlock(mtx);
|
||||||
|
if (!ok)
|
||||||
|
free(dup);
|
||||||
|
return ok;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Record one pruned filesystem path in a delete-protection sink. The stored
|
||||||
|
form is the entry's wire/destination-relative path (a single leading '/'
|
||||||
|
removed, exactly how manifest keep entries are stored), so the receiver's
|
||||||
|
walker prefixes match the destination layout. An allocation failure is a
|
||||||
|
fatal scan error. */
|
||||||
|
static void scanner_record_protected(DirectoryScanner* scanner, const char* fs_path,
|
||||||
|
ArrayList* sink) {
|
||||||
|
if (!sink || !fs_path)
|
||||||
|
return;
|
||||||
|
const char* rel = *fs_path == '/' ? fs_path + 1 : fs_path;
|
||||||
|
if (!excluded_sink_append(sink, scanner->options.excluded_mutex, rel))
|
||||||
|
scanner->failed = true;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* rsync's `--info=nonreg` line for a non-regular entry that is not being
|
||||||
|
* preserved: `skipping non-regular file "NAME"`. The name is the path relative
|
||||||
|
* to the transfer root, so it matches rsync's displayed name. */
|
||||||
|
void scanner_note_nonreg(const ScannerOptions* options, const char* fs_path) {
|
||||||
|
if (!options || !options->note_nonreg || !fs_path)
|
||||||
|
return;
|
||||||
|
const char* rel = utils_strip_transfer_root(fs_path, options->send_directory);
|
||||||
|
char* escaped = output_escape(rel, options->eight_bit_output);
|
||||||
|
printf("skipping non-regular file \"%s\"\n", escaped ? escaped : rel);
|
||||||
|
free(escaped);
|
||||||
|
fflush(stdout);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* rsync 3.4.1's `--info=mount` line, emitted when `-xx` drops a mount-point
|
||||||
|
* directory: `[sender] skipping mount-point dir NAME` (the client is the
|
||||||
|
* sender). Plain `-x` keeps the empty directory and prints nothing, matching
|
||||||
|
* rsync. */
|
||||||
|
void scanner_note_mount(const ScannerOptions* options, const char* fs_path) {
|
||||||
|
if (!options || !options->note_mount || !fs_path)
|
||||||
|
return;
|
||||||
|
const char* rel = utils_strip_transfer_root(fs_path, options->send_directory);
|
||||||
|
char* escaped = output_escape(rel, options->eight_bit_output);
|
||||||
|
printf("[sender] skipping mount-point dir %s\n", escaped ? escaped : rel);
|
||||||
|
free(escaped);
|
||||||
|
fflush(stdout);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* --debug=filter: a selection/filter decision dropped an entry. */
|
||||||
|
void scanner_note_filter(const ScannerOptions* options, const char* name) {
|
||||||
|
if (!options || !log_debug_enabled(LOG_DEBUG_FILTER) || !name)
|
||||||
|
return;
|
||||||
|
log_debug_message(LOG_DEBUG_FILTER, "filter: excluded %s", name);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Account for a directory that will not be represented by an inline directory
|
||||||
|
* entry. Paired with scanner_dir_count_uncount for empty directories that are
|
||||||
|
* emitted inline, so every traversed directory is counted exactly once. */
|
||||||
|
void scanner_dir_count_count(const ScannerOptions* options) {
|
||||||
|
if (options && options->dir_count)
|
||||||
|
atomic_fetch_add(options->dir_count, 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
void scanner_dir_count_uncount(const ScannerOptions* options) {
|
||||||
|
if (options && options->dir_count)
|
||||||
|
atomic_fetch_sub(options->dir_count, 1);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* A user-selection exclusion (--filter/-C/per-dir or --exclude/--include). */
|
||||||
|
void scanner_record_excluded(DirectoryScanner* scanner, const char* fs_path) {
|
||||||
|
scanner_record_protected(scanner, fs_path, scanner->options.excluded_paths);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* A --max-size/--min-size prune (always protected, even under --delete-excluded). */
|
||||||
|
void scanner_record_size_skipped(DirectoryScanner* scanner, const char* fs_path) {
|
||||||
|
scanner_record_protected(scanner, fs_path, scanner->options.size_skipped_paths);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Record a directory the scan synchronized. `fs_path` is its absolute path and
|
||||||
|
`rel` its path relative to the transfer root ("" for the root); the stored
|
||||||
|
form matches the wire layout (the bare relative path in -R+--files-from, else
|
||||||
|
the source path with a leading '/' removed, with "." for the receive root).
|
||||||
|
Returns false on allocation failure. */
|
||||||
|
bool scanner_record_synced_dir(const ScannerOptions* options, const char* fs_path, const char* rel,
|
||||||
|
bool relative_mode) {
|
||||||
|
if (!options->synced_dirs && !options->plan_dirs)
|
||||||
|
return true;
|
||||||
|
if (!file_list_dir_in_scope(options->file_list, rel))
|
||||||
|
return true;
|
||||||
|
char* prefixed = NULL;
|
||||||
|
const char* dest;
|
||||||
|
if (relative_mode) {
|
||||||
|
dest = rel;
|
||||||
|
} else if (options->relative_prefix) {
|
||||||
|
prefixed = scanner_prefix_send_path(options->relative_prefix, rel);
|
||||||
|
if (!prefixed)
|
||||||
|
return false;
|
||||||
|
dest = prefixed;
|
||||||
|
} else {
|
||||||
|
dest = fs_path;
|
||||||
|
}
|
||||||
|
if (dest[0] == '/')
|
||||||
|
dest++;
|
||||||
|
if (dest[0] == '\0')
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Read every per-directory filter file that applies to `dir_path` (its
|
||||||
|
* .rsync-filter when -F is active, plus each registered "dir-merge NAME") into a
|
||||||
|
* fresh list. Returns NULL on allocation/parse failure (message in `err`);
|
||||||
|
* returns an empty list (and *any_exists=false) when no file exists. */
|
||||||
|
FilterRuleList* read_dir_filters(const ScannerOptions* options, const char* dir_path,
|
||||||
|
const char* rel, bool* any_exists, char* err, size_t err_size) {
|
||||||
|
if (err && err_size > 0)
|
||||||
|
err[0] = '\0';
|
||||||
|
const FilterRuleList* base = options->base_filters;
|
||||||
|
bool have_names = options->per_dir_filters || (base && base->dir_merge_count > 0);
|
||||||
|
if (any_exists)
|
||||||
|
*any_exists = false;
|
||||||
|
if (!have_names)
|
||||||
|
return NULL;
|
||||||
|
FilterRuleList* own = filter_rule_list_create();
|
||||||
|
if (!own) {
|
||||||
|
snprintf(err, err_size, "memory allocation failed");
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
FilterParseOptions opts = {.delete_excluded = options->delete_excluded, .cvs_exclude = false};
|
||||||
|
bool exists = false;
|
||||||
|
if (options->per_dir_filters) {
|
||||||
|
if (!filter_file_append(own, dir_path, ".rsync-filter", rel, &opts, &exists, err, err_size))
|
||||||
|
goto fail;
|
||||||
|
if (exists && any_exists)
|
||||||
|
*any_exists = true;
|
||||||
|
}
|
||||||
|
if (base) {
|
||||||
|
for (int i = 0; i < base->dir_merge_count; i++) {
|
||||||
|
if (!filter_file_append(own, dir_path, base->dir_merge_names[i], rel, &opts, &exists, err,
|
||||||
|
err_size))
|
||||||
|
goto fail;
|
||||||
|
if (exists && any_exists)
|
||||||
|
*any_exists = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return own;
|
||||||
|
fail:
|
||||||
|
filter_rule_list_free(own);
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Merge the open directory's own per-directory filter files (the default
|
||||||
|
* .rsync-filter when -F is active, plus every "dir-merge NAME" registered on the
|
||||||
|
* base rule list) 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. */
|
||||||
|
int open_directory_filter_context(DirectoryScanner* scanner, const FilterNode* inherited) {
|
||||||
|
char err[256];
|
||||||
|
bool any_exists = false;
|
||||||
|
FilterRuleList* own = read_dir_filters(&scanner->options, scanner->current_path,
|
||||||
|
scanner->current_rel ? scanner->current_rel : "",
|
||||||
|
&any_exists, err, sizeof(err));
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
char* escaped_path = output_escape(scanner->current_path, log_get_8_bit_output());
|
||||||
|
log_message(LOG_LEVEL_ERROR, "invalid per-directory filter in %s: %s",
|
||||||
|
escaped_path ? escaped_path : "<allocation failed>", err);
|
||||||
|
free(escaped_path);
|
||||||
|
scanner->failed = true;
|
||||||
|
return -1;
|
||||||
|
}
|
||||||
|
if (any_exists && (own->count > 0 || own->dir_merge_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.
|
||||||
|
* `link_rel` is the entry's path relative to the transfer root (including its
|
||||||
|
* name), used for the lexical rsync unsafe-symlink test. */
|
||||||
|
int scanner_inspect_entry(const ScannerOptions* options, const char* containing_dir,
|
||||||
|
const char* link_rel, const char* name, ScannerEntry* entry) {
|
||||||
|
entry->excluded = false;
|
||||||
|
entry->size_excluded = false;
|
||||||
|
entry->referent_error = false;
|
||||||
|
entry->is_symlink = false;
|
||||||
|
entry->link_target = NULL;
|
||||||
|
entry->path = path_cat(containing_dir, name);
|
||||||
|
if (!entry->path)
|
||||||
|
return -1;
|
||||||
|
|
||||||
|
struct stat link_stats;
|
||||||
|
if (lstat(entry->path, &link_stats) != 0) {
|
||||||
|
free(entry->path);
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
if (!S_ISLNK(link_stats.st_mode))
|
||||||
|
goto regular;
|
||||||
|
|
||||||
|
char link_target[4096];
|
||||||
|
switch (scanner_link_action(options, entry->path, link_rel, link_target, sizeof(link_target))) {
|
||||||
|
case LINK_ACTION_SKIP:
|
||||||
|
goto skip;
|
||||||
|
case LINK_ACTION_SKIP_PROTECTED:
|
||||||
|
/* --safe-links ignored the link, but rsync still counts it as present in
|
||||||
|
the transfer, so its destination mirror survives --delete. Record it as
|
||||||
|
an excluded path (the same delete-protection channel as a filter prune). */
|
||||||
|
entry->excluded = true;
|
||||||
|
goto skip;
|
||||||
|
case LINK_ACTION_DEREF:
|
||||||
|
if (stat(entry->path, &entry->stats) != 0) {
|
||||||
|
/* rsync reports "symlink has no referent" and continues with a partial
|
||||||
|
transfer (exit 23); record the error so the run exits 23 too. */
|
||||||
|
char* escaped = output_escape(entry->path, log_get_8_bit_output());
|
||||||
|
log_message(LOG_LEVEL_WARNING, "symlink has no referent: %s",
|
||||||
|
escaped ? escaped : "<allocation failed>");
|
||||||
|
free(escaped);
|
||||||
|
entry->referent_error = true;
|
||||||
|
goto skip;
|
||||||
|
}
|
||||||
|
entry->is_directory = S_ISDIR(entry->stats.st_mode);
|
||||||
|
if (entry->is_directory)
|
||||||
|
return 1;
|
||||||
|
goto apply_filters;
|
||||||
|
case LINK_ACTION_CARRY:
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Carry the link as a symlink. --munge-links is applied by the RECEIVER (it
|
||||||
|
prefixes every stored target with /rsyncd-munged/); when the SOURCE already
|
||||||
|
holds a munged value the sender strips it so the receiver re-munges a clean
|
||||||
|
target, round-tripping a munged tree exactly like rsync. */
|
||||||
|
entry->is_symlink = true;
|
||||||
|
entry->stats = link_stats;
|
||||||
|
entry->is_directory = false;
|
||||||
|
entry->link_target = str_dup(link_target);
|
||||||
|
if (!entry->link_target)
|
||||||
|
goto skip;
|
||||||
|
if (options->munge_links)
|
||||||
|
file_symlink_unmunge(entry->link_target);
|
||||||
|
goto apply_filters;
|
||||||
|
|
||||||
|
regular:
|
||||||
|
/* Not a symlink: the lstat() above already described this entry, and lstat
|
||||||
|
and stat are identical for every non-symlink, so reuse that result instead
|
||||||
|
of issuing a redundant stat() on the scanner hot path. stat() is still
|
||||||
|
used on the dereference paths above/below for actual symlinks (copy-links,
|
||||||
|
safe/copy-unsafe links, and -k symlinks-to-directories). */
|
||||||
|
entry->stats = link_stats;
|
||||||
|
entry->is_directory = S_ISDIR(link_stats.st_mode);
|
||||||
|
if (entry->is_directory)
|
||||||
|
return 1;
|
||||||
|
|
||||||
|
apply_filters:
|
||||||
|
for (int i = 0; i < options->exclude_count; i++)
|
||||||
|
if (glob_match(options->exclude_patterns[i], name)) {
|
||||||
|
entry->excluded = true;
|
||||||
|
goto skip;
|
||||||
|
}
|
||||||
|
if (options->include_count > 0) {
|
||||||
|
bool included = false;
|
||||||
|
for (int i = 0; i < options->include_count; i++)
|
||||||
|
if (glob_match(options->include_patterns[i], name))
|
||||||
|
included = true;
|
||||||
|
if (!included) {
|
||||||
|
entry->excluded = true;
|
||||||
|
goto skip;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if ((options->max_size > 0 && (unsigned long long)entry->stats.st_size > options->max_size) ||
|
||||||
|
(options->min_size > 0 && (unsigned long long)entry->stats.st_size < options->min_size)) {
|
||||||
|
entry->excluded = true;
|
||||||
|
entry->size_excluded = true;
|
||||||
|
goto skip;
|
||||||
|
}
|
||||||
|
return 1;
|
||||||
|
|
||||||
|
skip:
|
||||||
|
free(entry->path);
|
||||||
|
entry->path = NULL;
|
||||||
|
free(entry->link_target);
|
||||||
|
entry->link_target = NULL;
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
@@ -0,0 +1,107 @@
|
|||||||
|
#ifndef SCANNER_INTERNAL_H
|
||||||
|
#define SCANNER_INTERNAL_H
|
||||||
|
|
||||||
|
/* Internal declarations shared between the scanner translation units
|
||||||
|
* (scanner_filter.c, scanner.c, scanner_parallel.c). Nothing here is part of
|
||||||
|
* the public scanner façade (scanner.h); every symbol stays internal to the
|
||||||
|
* client module. */
|
||||||
|
|
||||||
|
#include "array_list.h"
|
||||||
|
#include "file.h"
|
||||||
|
#include "scanner.h"
|
||||||
|
#include <stdbool.h>
|
||||||
|
#include <stddef.h>
|
||||||
|
#include <sys/stat.h>
|
||||||
|
|
||||||
|
typedef struct {
|
||||||
|
char* path;
|
||||||
|
int depth;
|
||||||
|
FilterNode* context; /* inherited per-directory filter context */
|
||||||
|
} DirEntry;
|
||||||
|
|
||||||
|
/* How rsync's readlink_stat()/generator resolves one source symlink. */
|
||||||
|
typedef enum {
|
||||||
|
LINK_ACTION_SKIP, /* not transferred (no link option) */
|
||||||
|
LINK_ACTION_SKIP_PROTECTED, /* ignored as unsafe by --safe-links; rsync keeps
|
||||||
|
it in the transfer, so its destination mirror
|
||||||
|
must be protected from --delete */
|
||||||
|
LINK_ACTION_DEREF, /* follow the referent (--copy-links, an unsafe
|
||||||
|
target under --copy-unsafe-links, or -k dir) */
|
||||||
|
LINK_ACTION_CARRY, /* transmit the link itself (-l) */
|
||||||
|
} LinkAction;
|
||||||
|
|
||||||
|
typedef struct {
|
||||||
|
char* path;
|
||||||
|
struct stat stats;
|
||||||
|
bool is_directory;
|
||||||
|
/* True when the entry should be carried through as a SYMLINK (is_symlink)
|
||||||
|
rather than a dereferenced file/directory. When true, `link_target` holds
|
||||||
|
the owned target string to transmit (sender-munged under --munge-links);
|
||||||
|
ownership transfers to the File built from this entry. */
|
||||||
|
bool is_symlink;
|
||||||
|
char* link_target;
|
||||||
|
/* True when the entry was pruned by a user selection rule (--filter/-C/per-dir
|
||||||
|
rules or the --exclude/--include layer) rather than skipped for another
|
||||||
|
reason (unreadable, symlink policy, not applicable). */
|
||||||
|
bool excluded;
|
||||||
|
/* True when the entry was skipped specifically by --max-size/--min-size.
|
||||||
|
Size pruning protects the destination mirror even under --delete-excluded,
|
||||||
|
so it is recorded into a separate sink from `excluded`. */
|
||||||
|
bool size_excluded;
|
||||||
|
/* True when a symlink selected for dereferencing (-L/--copy-links or an
|
||||||
|
unsafe target under --copy-unsafe-links) had no usable referent (a broken
|
||||||
|
link or a stat() failure). rsync still reports this as a partial transfer
|
||||||
|
(exit 23) even though the entry is skipped, so the scanner records it as a
|
||||||
|
non-fatal I/O error. */
|
||||||
|
bool referent_error;
|
||||||
|
} ScannerEntry;
|
||||||
|
|
||||||
|
typedef enum {
|
||||||
|
SCANNER_SPECIAL_REGULAR, /* ordinary file: transfer content */
|
||||||
|
SCANNER_SPECIAL_RECREATE, /* is_special node to recreate on the receiver */
|
||||||
|
SCANNER_SPECIAL_SKIP, /* non-regular entry not requested: skip */
|
||||||
|
} ScannerSpecial;
|
||||||
|
|
||||||
|
/* scanner_filter.c */
|
||||||
|
void filter_node_destroy(void* item);
|
||||||
|
FilterNode* filter_node_alloc(FilterNode* parent, FilterRuleList* own);
|
||||||
|
void dir_entry_destroy(void* item);
|
||||||
|
DirEntry* dir_entry_create(const char* path, int depth, FilterNode* context);
|
||||||
|
LinkAction scanner_link_action(const ScannerOptions* options, const char* path,
|
||||||
|
const char* link_rel, char* target, size_t target_size);
|
||||||
|
File* scanner_build_dir_file(const char* path, const struct stat* stats,
|
||||||
|
const ScannerOptions* options);
|
||||||
|
char* child_rel_path(const char* parent_rel, const char* name);
|
||||||
|
char* scanner_prefix_send_path(const char* prefix, const char* rel);
|
||||||
|
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, bool exclude_filter_files, bool* protect_out);
|
||||||
|
void scanner_capture_xattrs(const DirectoryScanner* scanner, File* file);
|
||||||
|
void scanner_assign_hardlink(DirectoryScanner* scanner, HardLinkTable* table, File* file,
|
||||||
|
const struct stat* stats);
|
||||||
|
ScannerSpecial scanner_prepare_special(bool preserve_devices, bool preserve_specials,
|
||||||
|
bool copy_devices, File* file, const struct stat* stats);
|
||||||
|
bool excluded_sink_append(ArrayList* list, mtx_t* mtx, const char* rel);
|
||||||
|
void scanner_note_nonreg(const ScannerOptions* options, const char* fs_path);
|
||||||
|
void scanner_note_mount(const ScannerOptions* options, const char* fs_path);
|
||||||
|
void scanner_note_filter(const ScannerOptions* options, const char* name);
|
||||||
|
void scanner_dir_count_count(const ScannerOptions* options);
|
||||||
|
void scanner_dir_count_uncount(const ScannerOptions* options);
|
||||||
|
void scanner_record_excluded(DirectoryScanner* scanner, const char* fs_path);
|
||||||
|
void scanner_record_size_skipped(DirectoryScanner* scanner, const char* fs_path);
|
||||||
|
bool scanner_record_synced_dir(const ScannerOptions* options, const char* fs_path, const char* rel,
|
||||||
|
bool relative_mode);
|
||||||
|
FilterRuleList* read_dir_filters(const ScannerOptions* options, const char* dir_path,
|
||||||
|
const char* rel, bool* any_exists, char* err, size_t err_size);
|
||||||
|
int open_directory_filter_context(DirectoryScanner* scanner, const FilterNode* inherited);
|
||||||
|
int scanner_inspect_entry(const ScannerOptions* options, const char* containing_dir,
|
||||||
|
const char* link_rel, const char* name, ScannerEntry* entry);
|
||||||
|
|
||||||
|
/* scanner.c */
|
||||||
|
bool scanner_capture_dir_time(ArrayList* dir_entries, mtx_t* mutex, const char* root_path,
|
||||||
|
const char* fs_path, bool relative_mode, const char* relative_prefix,
|
||||||
|
bool preserve_atimes, bool preserve_crtimes, bool preserve_xattrs,
|
||||||
|
bool preserve_acls, bool no_implied_dirs,
|
||||||
|
const FileListSet* file_list);
|
||||||
|
|
||||||
|
#endif
|
||||||
@@ -0,0 +1,703 @@
|
|||||||
|
#include "log.h"
|
||||||
|
#include "scanner.h"
|
||||||
|
#include "scanner_internal.h"
|
||||||
|
#include "array_list.h"
|
||||||
|
#include "chunk.h"
|
||||||
|
#include "file.h"
|
||||||
|
#include "queue.h"
|
||||||
|
#include "utils.h"
|
||||||
|
#include <dirent.h>
|
||||||
|
#include <stdio.h>
|
||||||
|
#include <stdlib.h>
|
||||||
|
#include <string.h>
|
||||||
|
#include <sys/stat.h>
|
||||||
|
#include <sys/sysmacros.h>
|
||||||
|
#include <threads.h>
|
||||||
|
#include <unistd.h>
|
||||||
|
#include <limits.h>
|
||||||
|
|
||||||
|
#include "xattr.h"
|
||||||
|
|
||||||
|
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;
|
||||||
|
|
||||||
|
static int parallel_worker_thread(void* arg) {
|
||||||
|
ParallelWorkerArg* wa = (ParallelWorkerArg*)arg;
|
||||||
|
ProtocolSession* allocation_session = wa->allocation_session;
|
||||||
|
if (allocation_session)
|
||||||
|
protocol_session_bind(allocation_session);
|
||||||
|
for (int i = 0; i < wa->dir_count; i++) {
|
||||||
|
DirectoryScanner* ds = directory_scanner_create_with_options(wa->dirs[i], &wa->options);
|
||||||
|
if (!ds) {
|
||||||
|
mtx_lock(&wa->ps->result_mutex);
|
||||||
|
wa->ps->failed = true;
|
||||||
|
atomic_store(&wa->ps->cancelled, true);
|
||||||
|
cnd_broadcast(&wa->ps->result_not_empty);
|
||||||
|
cnd_broadcast(&wa->ps->result_not_full);
|
||||||
|
mtx_unlock(&wa->ps->result_mutex);
|
||||||
|
for (int j = i; j < wa->dir_count; j++)
|
||||||
|
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. Exclusion
|
||||||
|
* recording shares one caller-owned list across the workers. */
|
||||||
|
free(ds->root_path);
|
||||||
|
ds->root_path = str_dup(wa->root_dir);
|
||||||
|
ds->seed_node = wa->ps->root_filter_node;
|
||||||
|
ds->options.excluded_mutex = &wa->ps->result_mutex;
|
||||||
|
Chunk* chunk;
|
||||||
|
while ((chunk = directory_scanner_next(ds)) != NULL) {
|
||||||
|
if (!queue_enqueue_multithreaded_cancel(wa->ps->result_queue, chunk, &wa->ps->result_mutex,
|
||||||
|
&wa->ps->result_not_empty, &wa->ps->result_not_full,
|
||||||
|
&wa->ps->cancelled)) {
|
||||||
|
chunk_destroy(chunk);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (directory_scanner_failed(ds)) {
|
||||||
|
mtx_lock(&wa->ps->result_mutex);
|
||||||
|
wa->ps->failed = true;
|
||||||
|
atomic_store(&wa->ps->cancelled, true);
|
||||||
|
cnd_broadcast(&wa->ps->result_not_empty);
|
||||||
|
cnd_broadcast(&wa->ps->result_not_full);
|
||||||
|
mtx_unlock(&wa->ps->result_mutex);
|
||||||
|
} else if (directory_scanner_had_io_error(ds)) {
|
||||||
|
/* --ignore-errors path: an unreadable directory was skipped, not fatal. */
|
||||||
|
mtx_lock(&wa->ps->result_mutex);
|
||||||
|
wa->ps->io_error = true;
|
||||||
|
mtx_unlock(&wa->ps->result_mutex);
|
||||||
|
}
|
||||||
|
directory_scanner_destroy(ds);
|
||||||
|
free(wa->dirs[i]);
|
||||||
|
}
|
||||||
|
ParallelScanner* ps = wa->ps;
|
||||||
|
free(wa->root_dir);
|
||||||
|
free(wa->dirs);
|
||||||
|
free(wa);
|
||||||
|
mtx_lock(&ps->result_mutex);
|
||||||
|
ps->completed++;
|
||||||
|
if (ps->completed >= ps->expected_threads) {
|
||||||
|
ps->done = true;
|
||||||
|
cnd_signal(&ps->result_not_empty);
|
||||||
|
}
|
||||||
|
mtx_unlock(&ps->result_mutex);
|
||||||
|
if (allocation_session)
|
||||||
|
protocol_session_unbind();
|
||||||
|
return thrd_success;
|
||||||
|
}
|
||||||
|
|
||||||
|
static void parallel_scanner_creation_failed(ParallelScanner* ps) {
|
||||||
|
mtx_lock(&ps->result_mutex);
|
||||||
|
ps->failed = true;
|
||||||
|
atomic_store(&ps->cancelled, true);
|
||||||
|
ps->expected_threads = ps->created_threads;
|
||||||
|
if (ps->completed >= ps->expected_threads)
|
||||||
|
ps->done = true;
|
||||||
|
cnd_broadcast(&ps->result_not_empty);
|
||||||
|
cnd_broadcast(&ps->result_not_full);
|
||||||
|
mtx_unlock(&ps->result_mutex);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Initialize result queue and synchronization primitives. Returns true on success. */
|
||||||
|
static bool parallel_scanner_init(ParallelScanner* ps) {
|
||||||
|
ps->result_queue = queue_create(100, chunk_destroy);
|
||||||
|
if (!ps->result_queue)
|
||||||
|
return false;
|
||||||
|
atomic_init(&ps->cancelled, false);
|
||||||
|
int init = 0;
|
||||||
|
bool ok = true;
|
||||||
|
if (mtx_init(&ps->result_mutex, mtx_plain) != thrd_success)
|
||||||
|
ok = false;
|
||||||
|
if (ok) {
|
||||||
|
init++;
|
||||||
|
if (cnd_init(&ps->result_not_empty) != thrd_success)
|
||||||
|
ok = false;
|
||||||
|
}
|
||||||
|
if (ok) {
|
||||||
|
// cppcheck-suppress unreadVariable
|
||||||
|
init++;
|
||||||
|
if (cnd_init(&ps->result_not_full) != thrd_success)
|
||||||
|
ok = false;
|
||||||
|
}
|
||||||
|
if (!ok) {
|
||||||
|
if (init >= 3)
|
||||||
|
cnd_destroy(&ps->result_not_full);
|
||||||
|
if (init >= 2)
|
||||||
|
cnd_destroy(&ps->result_not_empty);
|
||||||
|
if (init >= 1)
|
||||||
|
mtx_destroy(&ps->result_mutex);
|
||||||
|
queue_destroy(ps->result_queue);
|
||||||
|
ps->result_queue = NULL;
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Split files into chunks of roughly chunk_size bytes. Returns the first chunk (also stored
|
||||||
|
* chunks beyond the first are enqueued on `queue`). Nulls out consumed entries in `files`.
|
||||||
|
* Sets *failed on allocation/enqueue errors. */
|
||||||
|
static Chunk* batch_files(ArrayList* files, unsigned long long chunk_size, Queue* queue,
|
||||||
|
bool* failed) {
|
||||||
|
Chunk* first = NULL;
|
||||||
|
if (files->size <= 0)
|
||||||
|
return NULL;
|
||||||
|
ArrayList* batch = array_list_create(NULL);
|
||||||
|
if (!batch) {
|
||||||
|
*failed = true;
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
unsigned long long batch_size = 0;
|
||||||
|
for (int i = 0; i < files->size; i++) {
|
||||||
|
File* f = (File*)files->items[i];
|
||||||
|
if (!array_list_add(batch, f)) {
|
||||||
|
*failed = true;
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
batch_size += f->data->size;
|
||||||
|
if (batch_size >= chunk_size || i == files->size - 1) {
|
||||||
|
void** items = array_list_to_array(batch);
|
||||||
|
if (!items) {
|
||||||
|
*failed = true;
|
||||||
|
array_list_delete(batch);
|
||||||
|
batch = NULL;
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
Chunk* c = chunk_create((File**)items, batch->size);
|
||||||
|
free(items);
|
||||||
|
if (!c) {
|
||||||
|
*failed = true;
|
||||||
|
array_list_delete(batch);
|
||||||
|
batch = NULL;
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
int batch_start = i - batch->size + 1;
|
||||||
|
for (int j = batch_start; j <= i; j++)
|
||||||
|
files->items[j] = NULL;
|
||||||
|
batch->item_destroyer = NULL;
|
||||||
|
array_list_delete(batch);
|
||||||
|
batch = NULL;
|
||||||
|
if (!first) {
|
||||||
|
first = c;
|
||||||
|
} else {
|
||||||
|
if (!queue_enqueue(queue, c)) {
|
||||||
|
chunk_destroy(c);
|
||||||
|
*failed = true;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (i < files->size - 1) {
|
||||||
|
batch = array_list_create(NULL);
|
||||||
|
if (!batch) {
|
||||||
|
*failed = true;
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
batch_size = 0;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (batch) {
|
||||||
|
batch->item_destroyer = NULL;
|
||||||
|
array_list_delete(batch);
|
||||||
|
}
|
||||||
|
return first;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Scan one root-directory entry into either the subdirs or files list. */
|
||||||
|
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, entry->d_name, entry->d_name, &inspected);
|
||||||
|
if (inspection < 0) {
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (inspection == 0) {
|
||||||
|
if (inspected.referent_error)
|
||||||
|
ps->io_error = true;
|
||||||
|
ArrayList* sink = NULL;
|
||||||
|
if (inspected.excluded)
|
||||||
|
sink = inspected.size_excluded ? options->size_skipped_paths : options->excluded_paths;
|
||||||
|
if (sink) {
|
||||||
|
/* A root-level prune protects the destination mirror of the entry's wire
|
||||||
|
path: under -R + --files-from that is the bare relative name, otherwise
|
||||||
|
it is the full source path with a leading '/' removed (matching the
|
||||||
|
send_path/file_wire_path the scanner hands the sender). */
|
||||||
|
if (options->relative && options->file_list != NULL) {
|
||||||
|
if (!excluded_sink_append(sink, options->excluded_mutex, entry->d_name))
|
||||||
|
ps->failed = true;
|
||||||
|
} else if (options->relative_prefix) {
|
||||||
|
char* wrel = scanner_prefix_send_path(options->relative_prefix, entry->d_name);
|
||||||
|
if (!wrel) {
|
||||||
|
ps->failed = true;
|
||||||
|
} else {
|
||||||
|
if (!excluded_sink_append(sink, options->excluded_mutex, wrel))
|
||||||
|
ps->failed = true;
|
||||||
|
free(wrel);
|
||||||
|
}
|
||||||
|
} else {
|
||||||
|
char* abs_path = path_cat(root_directory, entry->d_name);
|
||||||
|
if (!abs_path) {
|
||||||
|
ps->failed = true;
|
||||||
|
} else {
|
||||||
|
const char* rel = *abs_path == '/' ? abs_path + 1 : abs_path;
|
||||||
|
if (!excluded_sink_append(sink, options->excluded_mutex, rel))
|
||||||
|
ps->failed = true;
|
||||||
|
free(abs_path);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
char* cur_path = inspected.path;
|
||||||
|
struct stat st = inspected.stats;
|
||||||
|
bool is_dir = inspected.is_directory;
|
||||||
|
char* rel = str_dup(entry->d_name);
|
||||||
|
if (!rel) {
|
||||||
|
free(cur_path);
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
bool protect = false;
|
||||||
|
bool passes = entry_passes_selection(options->file_list, options->base_filters, root_node, rel,
|
||||||
|
entry->d_name, is_dir, options->per_dir_filters,
|
||||||
|
options->exclude_per_dir_filter_files, &protect);
|
||||||
|
/* -R + --files-from: root-level files keep their bare relative send path. */
|
||||||
|
bool use_rel = options->relative && options->file_list != NULL;
|
||||||
|
if (!passes || protect) {
|
||||||
|
/* --files-from subset pruning is not a filter exclusion; -R bare-wire-path
|
||||||
|
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;
|
||||||
|
char* prefixed = NULL;
|
||||||
|
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);
|
||||||
|
free(cur_path);
|
||||||
|
ps->failed = true;
|
||||||
|
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))
|
||||||
|
ps->failed = true;
|
||||||
|
free(prefixed);
|
||||||
|
}
|
||||||
|
if (!passes) {
|
||||||
|
scanner_note_filter(options, entry->d_name);
|
||||||
|
free(rel);
|
||||||
|
free(cur_path);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (is_dir) {
|
||||||
|
if (!scanner_same_filesystem(options->one_file_system, root_dev, st.st_dev)) {
|
||||||
|
if (options->one_file_system > 1) {
|
||||||
|
/* -xx: drop the mount-point directory entirely (rsync) and print the
|
||||||
|
--info=mount line when enabled. */
|
||||||
|
scanner_note_mount(options, cur_path);
|
||||||
|
free(rel);
|
||||||
|
free(cur_path);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
/* -x/--one-file-system: emit the mount-point directory entry (empty) but
|
||||||
|
do not descend into it (see the sequential scanner for the same rule). */
|
||||||
|
File* mount = file_create(cur_path);
|
||||||
|
free(cur_path);
|
||||||
|
if (mount == NULL) {
|
||||||
|
free(rel);
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
mount->is_dir = true;
|
||||||
|
if (options->use_metadata) {
|
||||||
|
mount->metadata = file_metadata_create(mount->path, &st, options->preserve_atimes,
|
||||||
|
options->preserve_crtimes);
|
||||||
|
if (!mount->metadata) {
|
||||||
|
free(rel);
|
||||||
|
file_destroy(mount);
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (options->relative_prefix) {
|
||||||
|
mount->send_path = scanner_prefix_send_path(options->relative_prefix, rel);
|
||||||
|
if (!mount->send_path) {
|
||||||
|
free(rel);
|
||||||
|
file_destroy(mount);
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
free(rel);
|
||||||
|
if (!array_list_add(root_files, mount)) {
|
||||||
|
file_destroy(mount);
|
||||||
|
ps->failed = true;
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
free(rel);
|
||||||
|
if (!array_list_add(subdirs, cur_path)) {
|
||||||
|
free(cur_path);
|
||||||
|
ps->failed = true;
|
||||||
|
}
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
File* file = file_create(cur_path);
|
||||||
|
free(cur_path);
|
||||||
|
if (!file) {
|
||||||
|
free(rel);
|
||||||
|
free(inspected.link_target);
|
||||||
|
inspected.link_target = NULL;
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (inspected.is_symlink) {
|
||||||
|
file->is_symlink = true;
|
||||||
|
file->symlink_target = inspected.link_target;
|
||||||
|
inspected.link_target = NULL;
|
||||||
|
} else {
|
||||||
|
file->data->size = st.st_size;
|
||||||
|
}
|
||||||
|
if (use_rel) {
|
||||||
|
file->send_path = rel;
|
||||||
|
rel = NULL;
|
||||||
|
} else if (options->relative_prefix) {
|
||||||
|
file->send_path = scanner_prefix_send_path(options->relative_prefix, rel);
|
||||||
|
free(rel);
|
||||||
|
rel = NULL;
|
||||||
|
if (!file->send_path) {
|
||||||
|
file_destroy(file);
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
ScannerSpecial special = scanner_prepare_special(
|
||||||
|
options->preserve_devices, options->preserve_specials, options->copy_devices, file, &st);
|
||||||
|
if (special == SCANNER_SPECIAL_SKIP) {
|
||||||
|
scanner_note_nonreg(ps->options, file->path);
|
||||||
|
free(rel);
|
||||||
|
file_destroy(file);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (options->hardlinks && S_ISREG(st.st_mode)) {
|
||||||
|
int gid;
|
||||||
|
bool is_first;
|
||||||
|
char* first_path = NULL;
|
||||||
|
if (!hardlink_table_assign((HardLinkTable*)options->hardlinks, file_wire_path(file), st.st_dev,
|
||||||
|
st.st_ino, &gid, &is_first, &first_path)) {
|
||||||
|
ps->failed = true;
|
||||||
|
} else {
|
||||||
|
file->link_group = gid;
|
||||||
|
file->link_first = is_first;
|
||||||
|
if (!is_first) {
|
||||||
|
file->hardlink_target = first_path;
|
||||||
|
file->data->size = 0;
|
||||||
|
} else {
|
||||||
|
free(first_path);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if (options->use_metadata)
|
||||||
|
file->metadata =
|
||||||
|
file_metadata_create(file->path, &st, options->preserve_atimes, options->preserve_crtimes);
|
||||||
|
if (options->use_metadata && !file->metadata) {
|
||||||
|
free(rel);
|
||||||
|
file_destroy(file);
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if ((options->preserve_xattrs || options->preserve_acls) &&
|
||||||
|
!(file->link_group != 0 && !file->link_first))
|
||||||
|
file->xattrs = xattr_capture_path(file->path, options->preserve_acls);
|
||||||
|
if (!array_list_add(root_files, file)) {
|
||||||
|
free(rel);
|
||||||
|
file_destroy(file);
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
free(rel);
|
||||||
|
}
|
||||||
|
|
||||||
|
/* 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, 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");
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
/* The parallel scanner opens the transfer root directly (not through
|
||||||
|
open_next_directory), so record it as synchronized here. */
|
||||||
|
if (!scanner_record_synced_dir(options, root_directory, "",
|
||||||
|
options->relative && options->file_list != NULL)) {
|
||||||
|
closedir(dir);
|
||||||
|
ps->failed = true;
|
||||||
|
return false;
|
||||||
|
}
|
||||||
|
log_debug_message(LOG_DEBUG_FLIST, "flist: scanning %s", root_directory);
|
||||||
|
const struct dirent* entry;
|
||||||
|
while ((entry = readdir(dir)) != NULL) {
|
||||||
|
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
|
||||||
|
continue;
|
||||||
|
scan_root_entry(options, root_node, root_directory, entry, root_files, subdirs, root_dev, ps);
|
||||||
|
}
|
||||||
|
closedir(dir);
|
||||||
|
return true;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Spawn worker threads, one per group of subdirectories. */
|
||||||
|
static void spawn_parallel_workers(ParallelScanner* ps, ArrayList* subdirs,
|
||||||
|
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;
|
||||||
|
if (n > subdirs->size)
|
||||||
|
n = subdirs->size;
|
||||||
|
|
||||||
|
ps->num_threads = n;
|
||||||
|
ps->expected_threads = n;
|
||||||
|
ps->threads = calloc(n, sizeof(thrd_t));
|
||||||
|
if (!ps->threads) {
|
||||||
|
ps->num_threads = 0;
|
||||||
|
ps->expected_threads = 0;
|
||||||
|
ps->failed = true;
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
int dirs_per_thread = subdirs->size / n;
|
||||||
|
int remainder = subdirs->size % n;
|
||||||
|
int start = 0;
|
||||||
|
ps->num_threads = 0;
|
||||||
|
for (int t = 0; t < n; t++) {
|
||||||
|
int count = dirs_per_thread + (t < remainder ? 1 : 0);
|
||||||
|
if (count == 0)
|
||||||
|
break;
|
||||||
|
ParallelWorkerArg* wa = calloc(1, sizeof(ParallelWorkerArg));
|
||||||
|
if (!wa) {
|
||||||
|
parallel_scanner_creation_failed(ps);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
wa->ps = ps;
|
||||||
|
wa->dirs = calloc(count, sizeof(char*));
|
||||||
|
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;
|
||||||
|
}
|
||||||
|
bool dup_ok = true;
|
||||||
|
for (int j = 0; j < count; j++) {
|
||||||
|
wa->dirs[j] = str_dup((char*)subdirs->items[start + j]);
|
||||||
|
if (!wa->dirs[j])
|
||||||
|
dup_ok = false;
|
||||||
|
}
|
||||||
|
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);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
wa->dir_count = count;
|
||||||
|
wa->options = *options;
|
||||||
|
wa->options.chunk_size = cs;
|
||||||
|
wa->allocation_session = ps->allocation_session;
|
||||||
|
start += count;
|
||||||
|
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);
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
ps->num_threads++;
|
||||||
|
ps->created_threads++;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
ParallelScanner* parallel_scanner_create_with_options(const char* root_directory,
|
||||||
|
const ScannerOptions* options,
|
||||||
|
ProtocolSession* allocation_session) {
|
||||||
|
if (!root_directory || !options)
|
||||||
|
return NULL;
|
||||||
|
ParallelScanner* ps = calloc(1, sizeof(ParallelScanner));
|
||||||
|
if (!ps)
|
||||||
|
return NULL;
|
||||||
|
if (!parallel_scanner_init(ps)) {
|
||||||
|
free(ps);
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
ps->allocation_session = allocation_session;
|
||||||
|
ps->options = options;
|
||||||
|
|
||||||
|
ArrayList* root_files = array_list_create(file_destroy);
|
||||||
|
ArrayList* subdirs = array_list_create(free);
|
||||||
|
if (!root_files || !subdirs) {
|
||||||
|
array_list_delete(root_files);
|
||||||
|
array_list_delete(subdirs);
|
||||||
|
parallel_scanner_destroy(ps);
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
|
||||||
|
dev_t root_dev = 0;
|
||||||
|
if (options->one_file_system) {
|
||||||
|
struct stat root_stats;
|
||||||
|
if (stat(root_directory, &root_stats) != 0) {
|
||||||
|
log_perror("Could not stat source directory");
|
||||||
|
array_list_delete(root_files);
|
||||||
|
array_list_delete(subdirs);
|
||||||
|
parallel_scanner_destroy(ps);
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
root_dev = root_stats.st_dev;
|
||||||
|
}
|
||||||
|
|
||||||
|
/* Build the root directory's per-directory filter context once; workers seed
|
||||||
|
* their scanners with it so per-dir rules behave identically to the sequential
|
||||||
|
* scanner. */
|
||||||
|
FilterNode* root_node = NULL;
|
||||||
|
{
|
||||||
|
char err[256];
|
||||||
|
bool any_exists = false;
|
||||||
|
FilterRuleList* own =
|
||||||
|
read_dir_filters(options, root_directory, "", &any_exists, err, sizeof(err));
|
||||||
|
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);
|
||||||
|
array_list_delete(subdirs);
|
||||||
|
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) {
|
||||||
|
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);
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
/* The root itself is a traversed directory (rsync counts it in
|
||||||
|
`Number of files`); the worker DirectoryScanners account for every
|
||||||
|
subdirectory below it. */
|
||||||
|
scanner_dir_count_count(options);
|
||||||
|
/* P7 Wave D: the parallel scanner never runs a DirectoryScanner over the
|
||||||
|
transfer root itself (it hands the root's immediate subdirectories to
|
||||||
|
workers), so capture the root's directory time here. */
|
||||||
|
if (options->capture_dir_times &&
|
||||||
|
!scanner_capture_dir_time(
|
||||||
|
options->dir_entries, options->dir_entries_mutex, root_directory, root_directory,
|
||||||
|
options->relative && options->file_list != NULL, options->relative_prefix,
|
||||||
|
options->preserve_atimes, options->preserve_crtimes, options->preserve_xattrs,
|
||||||
|
options->preserve_acls, options->no_implied_dirs, options->file_list)) {
|
||||||
|
array_list_delete(root_files);
|
||||||
|
array_list_delete(subdirs);
|
||||||
|
parallel_scanner_destroy(ps);
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
|
||||||
|
unsigned long long cs = options->chunk_size > 0 ? options->chunk_size : DESIRED_CHUNK_SIZE;
|
||||||
|
ps->initial_chunk = batch_files(root_files, cs, ps->result_queue, &ps->failed);
|
||||||
|
array_list_delete(root_files);
|
||||||
|
|
||||||
|
spawn_parallel_workers(ps, subdirs, options, root_directory, cs);
|
||||||
|
array_list_delete(subdirs);
|
||||||
|
return ps;
|
||||||
|
}
|
||||||
|
|
||||||
|
Chunk* parallel_scanner_next(ParallelScanner* ps) {
|
||||||
|
if (ps->initial_chunk) {
|
||||||
|
Chunk* c = ps->initial_chunk;
|
||||||
|
ps->initial_chunk = NULL;
|
||||||
|
return c;
|
||||||
|
}
|
||||||
|
if (ps->num_threads == 0) {
|
||||||
|
mtx_lock(&ps->result_mutex);
|
||||||
|
if (!queue_is_empty(ps->result_queue)) {
|
||||||
|
Chunk* chunk = queue_dequeue(ps->result_queue);
|
||||||
|
mtx_unlock(&ps->result_mutex);
|
||||||
|
return chunk;
|
||||||
|
}
|
||||||
|
ps->done = true;
|
||||||
|
mtx_unlock(&ps->result_mutex);
|
||||||
|
return NULL;
|
||||||
|
}
|
||||||
|
Chunk* chunk = queue_dequeue_multithreaded(
|
||||||
|
ps->result_queue, &ps->result_mutex, &ps->result_not_empty, &ps->result_not_full, &ps->done);
|
||||||
|
return chunk;
|
||||||
|
}
|
||||||
|
|
||||||
|
bool parallel_scanner_failed(const ParallelScanner* ps) {
|
||||||
|
return ps == NULL || ps->failed;
|
||||||
|
}
|
||||||
|
|
||||||
|
bool parallel_scanner_had_io_error(const ParallelScanner* ps) {
|
||||||
|
return ps != NULL && ps->io_error;
|
||||||
|
}
|
||||||
|
|
||||||
|
void parallel_scanner_destroy(ParallelScanner* ps) {
|
||||||
|
if (!ps)
|
||||||
|
return;
|
||||||
|
mtx_lock(&ps->result_mutex);
|
||||||
|
ps->done = true;
|
||||||
|
atomic_store(&ps->cancelled, true);
|
||||||
|
cnd_broadcast(&ps->result_not_empty);
|
||||||
|
cnd_broadcast(&ps->result_not_full);
|
||||||
|
mtx_unlock(&ps->result_mutex);
|
||||||
|
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);
|
||||||
|
mtx_destroy(&ps->result_mutex);
|
||||||
|
cnd_destroy(&ps->result_not_empty);
|
||||||
|
cnd_destroy(&ps->result_not_full);
|
||||||
|
free(ps);
|
||||||
|
}
|
||||||
Reference in New Issue
Block a user