refactor(scanner,send): embed scanner options; unify stats and config ownership

This commit is contained in:
2026-09-13 05:28:32 +02:00
parent 57ce6d04f0
commit 99f8045105
9 changed files with 127 additions and 190 deletions
+38 -29
View File
@@ -44,6 +44,10 @@
receiver's RECEIVER_QUEUE_MAX_BYTES). */ receiver's RECEIVER_QUEUE_MAX_BYTES). */
#define SENDER_QUEUE_MAX_BYTES (MAX_CONNECTION_MEMORY - 2 * MAX_CHUNK_SIZE) #define SENDER_QUEUE_MAX_BYTES (MAX_CONNECTION_MEMORY - 2 * MAX_CHUNK_SIZE)
/* One mebibyte in bytes; the unit used by the --stats/--progress lines.
Always cast to double when dividing so the output stays fractional. */
#define BYTES_PER_MIB (1024ULL * 1024ULL)
/* Forward declaration for progress-reporting thread used in multithreaded send. */ /* Forward declaration for progress-reporting thread used in multithreaded send. */
static int progress_thread_fn(void* arg); static int progress_thread_fn(void* arg);
@@ -51,10 +55,32 @@ static const char* display_bytes(unsigned long long bytes, bool human_readable,
size_t buffer_size) { size_t buffer_size) {
if (human_readable && format_human_bytes(bytes, buffer, buffer_size)) if (human_readable && format_human_bytes(bytes, buffer, buffer_size))
return buffer; return buffer;
snprintf(buffer, buffer_size, "%.1f MB", bytes / 1048576.0); snprintf(buffer, buffer_size, "%.1f MB", (double)bytes / (double)BYTES_PER_MIB);
return buffer; return buffer;
} }
/* Print the canonical `--stats` line. Shared by the single-threaded and
multithreaded send paths so both honor --stats, --human-readable and --quiet
identically; `start` marks the beginning of the transfer for the rate. */
static void report_transfer_stats(const Config* config, int total_files,
unsigned long long total_bytes, time_t start) {
if (!config->stats || config->quiet)
return;
double elapsed = difftime(time(NULL), start);
double rate = elapsed > 0.0 ? (double)total_bytes / ((double)BYTES_PER_MIB * elapsed) : 0.0;
if (config->human_readable) {
char total_buffer[32];
char rate_buffer[32];
fprintf(stderr, "Stats: %d files, %s, %s/s\n", total_files,
display_bytes(total_bytes, true, total_buffer, sizeof(total_buffer)),
display_bytes((unsigned long long)(rate * (double)BYTES_PER_MIB), true, rate_buffer,
sizeof(rate_buffer)));
} else {
fprintf(stderr, "Stats: %d files, %.1f MB, %.1f MB/s\n", total_files,
(double)total_bytes / (double)BYTES_PER_MIB, rate);
}
}
/* Compiled scanner inputs that are shared read-only across scanner instances /* Compiled scanner inputs that are shared read-only across scanner instances
* and, in -m mode, across worker threads. `base_filters` owns the compiled * and, in -m mode, across worker threads. `base_filters` owns the compiled
* command-line + -C rules; the FileListSet allow-set lives in the Config. * command-line + -C rules; the FileListSet allow-set lives in the Config.
@@ -688,7 +714,7 @@ static int send_dry_run_manifest(const Config* config) {
printf("Total: %d files, %s\n", file_count, printf("Total: %d files, %s\n", file_count,
display_bytes(total_bytes, true, size_buffer, sizeof(size_buffer))); display_bytes(total_bytes, true, size_buffer, sizeof(size_buffer)));
else else
printf("Total: %d files, %.1f MB\n", file_count, total_bytes / 1048576.0); printf("Total: %d files, %.1f MB\n", file_count, (double)total_bytes / (double)BYTES_PER_MIB);
} }
return 0; return 0;
} }
@@ -1425,12 +1451,9 @@ static int send_chunk_with_removal(Client* client, Chunk* chunk, Config* config,
return 0; return 0;
} }
int send_chunk(Client* client, Chunk* chunk, Config* config) {
return send_chunk_with_removal(client, chunk, config, NULL);
}
static int send_chunks_multithreaded(void* pipeline_context) { static int send_chunks_multithreaded(void* pipeline_context) {
PipelineContextSender* context = (PipelineContextSender*)pipeline_context; PipelineContextSender* context = (PipelineContextSender*)pipeline_context;
time_t start = time(NULL);
Client* client = connect_transfer_client(context->config); Client* client = connect_transfer_client(context->config);
if (!client) { if (!client) {
if (context->config->transport == TRANSPORT_TCP) if (context->config->transport == TRANSPORT_TCP)
@@ -1578,10 +1601,9 @@ static int send_chunks_multithreaded(void* pipeline_context) {
int total_files = context->total_files; int total_files = context->total_files;
unsigned long long total_bytes = context->total_bytes; unsigned long long total_bytes = context->total_bytes;
mtx_unlock(&context->mutex_progress); mtx_unlock(&context->mutex_progress);
if (context->config->stats) report_transfer_stats(context->config, total_files, total_bytes, start);
fprintf(stderr, "Stats: %d files, %.1f MB\n", total_files, total_bytes / 1048576.0);
log_info_message(LOG_INFO_STATS, "Transfer summary: %d files, %.1f MB", total_files, log_info_message(LOG_INFO_STATS, "Transfer summary: %d files, %.1f MB", total_files,
total_bytes / 1048576.0); (double)total_bytes / (double)BYTES_PER_MIB);
disconnect_transfer_client(client); disconnect_transfer_client(client);
mark_sender_done(context); mark_sender_done(context);
protocol_session_unbind(); protocol_session_unbind();
@@ -1754,17 +1776,18 @@ static int load_files_multithreaded(void* pipeline_context) {
static void print_transfer_progress(unsigned long long total_bytes, time_t start, static void print_transfer_progress(unsigned long long total_bytes, time_t start,
const char* suffix, bool human_readable) { const char* suffix, bool human_readable) {
double elapsed = difftime(time(NULL), start); double elapsed = difftime(time(NULL), start);
double rate = elapsed > 0.0 ? total_bytes / (1048576.0 * elapsed) : 0.0; double rate = elapsed > 0.0 ? (double)total_bytes / ((double)BYTES_PER_MIB * elapsed) : 0.0;
if (human_readable) { if (human_readable) {
char total_buffer[32]; char total_buffer[32];
char rate_buffer[32]; char rate_buffer[32];
fprintf(stderr, "\rSent %s (%s/s) %s", fprintf(stderr, "\rSent %s (%s/s) %s",
display_bytes(total_bytes, true, total_buffer, sizeof(total_buffer)), display_bytes(total_bytes, true, total_buffer, sizeof(total_buffer)),
display_bytes((unsigned long long)(rate * 1048576.0), true, rate_buffer, display_bytes((unsigned long long)(rate * (double)BYTES_PER_MIB), true, rate_buffer,
sizeof(rate_buffer)), sizeof(rate_buffer)),
suffix); suffix);
} else { } else {
fprintf(stderr, "\rSent %.1f MB (%.1f MB/s) %s", total_bytes / 1048576.0, rate, suffix); fprintf(stderr, "\rSent %.1f MB (%.1f MB/s) %s", (double)total_bytes / (double)BYTES_PER_MIB,
rate, suffix);
} }
fflush(stderr); fflush(stderr);
} }
@@ -2132,23 +2155,9 @@ int send_files(Config* config) {
remove_transferred_sources(config, remove_sources); remove_transferred_sources(config, remove_sources);
if (config->show_progress && !config->quiet) if (config->show_progress && !config->quiet)
print_transfer_progress(total_bytes, start, "Done.\n", config->human_readable); print_transfer_progress(total_bytes, start, "Done.\n", config->human_readable);
if (config->stats && !config->quiet) { report_transfer_stats(config, total_files, total_bytes, start);
double elapsed_total = difftime(time(NULL), start);
double rate = elapsed_total > 0 ? total_bytes / (1048576.0 * elapsed_total) : 0;
if (config->human_readable) {
char total_buffer[32];
char rate_buffer[32];
fprintf(stderr, "Stats: %d files, %s, %s/s\n", total_files,
display_bytes(total_bytes, true, total_buffer, sizeof(total_buffer)),
display_bytes((unsigned long long)(rate * 1048576.0), true, rate_buffer,
sizeof(rate_buffer)));
} else {
fprintf(stderr, "Stats: %d files, %.1f MB, %.1f MB/s\n", total_files, total_bytes / 1048576.0,
rate);
}
}
log_info_message(LOG_INFO_STATS, "Transfer summary: %d files, %.1f MB", total_files, log_info_message(LOG_INFO_STATS, "Transfer summary: %d files, %.1f MB", total_files,
total_bytes / 1048576.0); (double)total_bytes / (double)BYTES_PER_MIB);
/* --ignore-errors: an unreadable source directory was skipped but the run /* --ignore-errors: an unreadable source directory was skipped but the run
still completed (and deleted); report the run as errored like rsync does. */ still completed (and deleted); report the run as errored like rsync does. */
ret = (ok && !had_scan_io) ? 0 : 1; ret = (ok && !had_scan_io) ? 0 : 1;
@@ -2237,7 +2246,7 @@ int send_files_multithreaded(Config** config_ptr) {
context->missing_args = missing_args; context->missing_args = missing_args;
missing_args = NULL; /* owned by the context from here on */ missing_args = NULL; /* owned by the context from here on */
pipeline_context_sender_set_queue_byte_limit(context, SENDER_QUEUE_MAX_BYTES); pipeline_context_sender_set_queue_byte_limit(context, SENDER_QUEUE_MAX_BYTES);
*config_ptr = NULL; /* context now owns config through all remaining paths */ /* The context borrows `config`; the caller (main) still owns and frees it. */
struct timespec now_mono; struct timespec now_mono;
if (clock_gettime(CLOCK_MONOTONIC, &now_mono) != 0) { if (clock_gettime(CLOCK_MONOTONIC, &now_mono) != 0) {
now_mono.tv_sec = 0; now_mono.tv_sec = 0;
+3 -2
View File
@@ -5,9 +5,10 @@
#include "config.h" #include "config.h"
#include "transport_tcp.h" #include "transport_tcp.h"
int send_chunk(Client* client, Chunk* chunk, Config* config); /* Both sender entry points BORROW `config` for the duration of the call; they
* never free it, and the caller retains ownership (freeing it with
* config_delete() once the call returns). */
int send_files(Config* config); int send_files(Config* config);
/* Takes ownership only when *config is set to NULL on return. */
int send_files_multithreaded(Config** config); int send_files_multithreaded(Config** config);
/* Phase 6 residual-batch (client-only). See client_send.c. */ /* Phase 6 residual-batch (client-only). See client_send.c. */
int write_batch_from_source(const Config* config, const char* batch_path); int write_batch_from_source(const Config* config, const char* batch_path);
+57 -108
View File
@@ -177,7 +177,7 @@ static bool entry_passes_selection(const FileListSet* file_list, const FilterRul
/* Best-effort capture of the file's whitelisted xattrs (-X/-A). A failure to /* 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. */ * read xattrs is non-fatal: the file is transferred without them. */
static void scanner_capture_xattrs(const DirectoryScanner* scanner, File* file) { static void scanner_capture_xattrs(const DirectoryScanner* scanner, File* file) {
if (!scanner || !file || !(scanner->preserve_xattrs || scanner->preserve_acls)) if (!scanner || !file || !(scanner->options.preserve_xattrs || scanner->options.preserve_acls))
return; return;
file->xattrs = xattr_capture_path(file->path); file->xattrs = xattr_capture_path(file->path);
} }
@@ -263,10 +263,10 @@ static bool excluded_sink_append(ArrayList* list, mtx_t* mtx, const char* rel) {
how manifest keep entries are stored), so the receiver's walker prefixes how manifest keep entries are stored), so the receiver's walker prefixes
match the destination layout. An allocation failure is a fatal scan error. */ match the destination layout. An allocation failure is a fatal scan error. */
static void scanner_record_excluded(DirectoryScanner* scanner, const char* fs_path) { static void scanner_record_excluded(DirectoryScanner* scanner, const char* fs_path) {
if (!scanner->excluded_paths || !fs_path) if (!scanner->options.excluded_paths || !fs_path)
return; return;
const char* rel = *fs_path == '/' ? fs_path + 1 : fs_path; const char* rel = *fs_path == '/' ? fs_path + 1 : fs_path;
if (!excluded_sink_append(scanner->excluded_paths, scanner->excluded_mutex, rel)) if (!excluded_sink_append(scanner->options.excluded_paths, scanner->options.excluded_mutex, rel))
scanner->failed = true; scanner->failed = true;
} }
@@ -274,7 +274,7 @@ static void scanner_record_excluded(DirectoryScanner* scanner, const char* fs_pa
* context, returning the context used for this directory's entries. On a parse * 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. */ * error the scanner is marked failed. Returns 0 on success, -1 on failure. */
static int open_directory_filter_context(DirectoryScanner* scanner, const FilterNode* inherited) { static int open_directory_filter_context(DirectoryScanner* scanner, const FilterNode* inherited) {
if (!scanner->per_dir_filters) { if (!scanner->options.per_dir_filters) {
scanner->current_node = (FilterNode*)inherited; scanner->current_node = (FilterNode*)inherited;
return 0; return 0;
} }
@@ -436,6 +436,11 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo
DirectoryScanner* scanner = calloc(1, sizeof(DirectoryScanner)); DirectoryScanner* scanner = calloc(1, sizeof(DirectoryScanner));
if (scanner == NULL) if (scanner == NULL)
return NULL; return NULL;
/* One copy of the scan inputs; normalize chunk_size as the old field-by-field
copy did. */
scanner->options = *options;
if (scanner->options.chunk_size == 0)
scanner->options.chunk_size = DESIRED_CHUNK_SIZE;
scanner->directories = queue_create(100, dir_entry_destroy); scanner->directories = queue_create(100, dir_entry_destroy);
if (!scanner->directories) { if (!scanner->directories) {
free(scanner); free(scanner);
@@ -443,31 +448,7 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo
} }
scanner->current_dir = NULL; scanner->current_dir = NULL;
scanner->current_path = NULL; scanner->current_path = NULL;
scanner->use_metadata = options->use_metadata;
scanner->preserve_atimes = options->preserve_atimes;
scanner->preserve_crtimes = options->preserve_crtimes;
scanner->preserve_xattrs = options->preserve_xattrs;
scanner->preserve_acls = options->preserve_acls;
scanner->chunk_size = options->chunk_size > 0 ? options->chunk_size : DESIRED_CHUNK_SIZE;
scanner->exclude_patterns = options->exclude_patterns;
scanner->exclude_count = options->exclude_count;
scanner->include_patterns = options->include_patterns;
scanner->include_count = options->include_count;
scanner->max_size = options->max_size;
scanner->min_size = options->min_size;
scanner->max_depth = options->max_depth;
scanner->current_depth = 0; scanner->current_depth = 0;
scanner->follow_symlinks = options->follow_symlinks;
scanner->copy_links = options->copy_links;
scanner->safe_links = options->safe_links;
scanner->copy_unsafe_links = options->copy_unsafe_links;
scanner->copy_dirlinks = options->copy_dirlinks;
scanner->munge_links = options->munge_links;
scanner->checksum = options->checksum;
scanner->one_file_system = options->one_file_system;
scanner->preserve_devices = options->preserve_devices;
scanner->preserve_specials = options->preserve_specials;
scanner->copy_devices = options->copy_devices;
scanner->failed = false; scanner->failed = false;
scanner->root_path = str_dup(root_directory); scanner->root_path = str_dup(root_directory);
if (!scanner->root_path) { if (!scanner->root_path) {
@@ -479,28 +460,14 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo
scanner->at_seed_dir = true; scanner->at_seed_dir = true;
scanner->seed_node = NULL; scanner->seed_node = NULL;
scanner->current_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->excluded_paths = options->excluded_paths;
scanner->excluded_mutex = options->excluded_mutex;
scanner->ignore_io_errors = options->ignore_io_errors;
scanner->ignore_missing_args = options->ignore_missing_args;
scanner->io_error = false; scanner->io_error = false;
scanner->dirs_mode = options->dirs;
scanner->relative_mode = options->relative && options->file_list != NULL; scanner->relative_mode = options->relative && options->file_list != NULL;
scanner->hardlinks = options->hardlinks;
scanner->prune_empty_dirs = options->prune_empty_dirs;
scanner->stop_condition = options->stop_condition;
scanner->capture_dir_times = options->capture_dir_times;
scanner->dir_entries = options->dir_entries;
scanner->dir_entries_mutex = options->dir_entries_mutex;
scanner->dirs_root_emitted = false; scanner->dirs_root_emitted = false;
scanner->list_index = 0; scanner->list_index = 0;
scanner->dirs_batch = NULL; scanner->dirs_batch = NULL;
scanner->dirs_batch_size = 0; scanner->dirs_batch_size = 0;
scanner->filter_nodes = NULL; scanner->filter_nodes = NULL;
if (scanner->base_filters || scanner->per_dir_filters) { if (scanner->options.base_filters || scanner->options.per_dir_filters) {
scanner->filter_nodes = array_list_create(filter_node_destroy); scanner->filter_nodes = array_list_create(filter_node_destroy);
if (!scanner->filter_nodes) { if (!scanner->filter_nodes) {
free(scanner->root_path); free(scanner->root_path);
@@ -509,7 +476,7 @@ DirectoryScanner* directory_scanner_create_with_options(const char* root_directo
return NULL; return NULL;
} }
} }
if (scanner->one_file_system) { if (scanner->options.one_file_system) {
struct stat root_stats; struct stat root_stats;
if (stat(root_directory, &root_stats) != 0) { if (stat(root_directory, &root_stats) != 0) {
log_perror("Could not stat source directory"); log_perror("Could not stat source directory");
@@ -715,7 +682,7 @@ static int open_next_directory(DirectoryScanner* scanner) {
scanner->current_rel = NULL; scanner->current_rel = NULL;
free(scanner->current_path); free(scanner->current_path);
scanner->current_path = NULL; scanner->current_path = NULL;
if (!scanner->ignore_io_errors || is_root_seed) { if (!scanner->options.ignore_io_errors || is_root_seed) {
scanner->failed = true; scanner->failed = true;
return -1; return -1;
} }
@@ -729,10 +696,11 @@ static int open_next_directory(DirectoryScanner* scanner) {
scanner->current_path = NULL; scanner->current_path = NULL;
return -1; return -1;
} }
if (scanner->capture_dir_times && if (scanner->options.capture_dir_times &&
!scanner_capture_dir_time(scanner->dir_entries, scanner->dir_entries_mutex, !scanner_capture_dir_time(scanner->options.dir_entries, scanner->options.dir_entries_mutex,
scanner->root_path, scanner->current_path, scanner->relative_mode, scanner->root_path, scanner->current_path, scanner->relative_mode,
scanner->preserve_atimes, scanner->preserve_crtimes)) { scanner->options.preserve_atimes,
scanner->options.preserve_crtimes)) {
closedir(scanner->current_dir); closedir(scanner->current_dir);
scanner->current_dir = NULL; scanner->current_dir = NULL;
free(scanner->current_path); free(scanner->current_path);
@@ -773,9 +741,9 @@ static File* dirs_root_dir_file(DirectoryScanner* scanner) {
return NULL; return NULL;
} }
file->is_dir = true; file->is_dir = true;
if (scanner->use_metadata) { if (scanner->options.use_metadata) {
file->metadata = file_metadata_create(scanner->root_path, &st, scanner->preserve_atimes, file->metadata = file_metadata_create(scanner->root_path, &st, scanner->options.preserve_atimes,
scanner->preserve_crtimes); scanner->options.preserve_crtimes);
if (!file->metadata) { if (!file->metadata) {
file_destroy(file); file_destroy(file);
scanner->failed = true; scanner->failed = true;
@@ -808,7 +776,7 @@ static File* dirs_file_for_entry(DirectoryScanner* scanner, const char* entry) {
missing argument and is skipped here, exactly as the recursive scan skips missing argument and is skipped here, exactly as the recursive scan skips
nothing (missing entries never appear there). Without the flags it stays nothing (missing entries never appear there). Without the flags it stays
a hard pre-transfer error. */ a hard pre-transfer error. */
if (scanner->ignore_missing_args) { if (scanner->options.ignore_missing_args) {
log_info_message(LOG_INFO_MISC, "skipping missing --files-from entry '%s'", entry); log_info_message(LOG_INFO_MISC, "skipping missing --files-from entry '%s'", entry);
free(abs_path); free(abs_path);
return NULL; return NULL;
@@ -822,8 +790,8 @@ static File* dirs_file_for_entry(DirectoryScanner* scanner, const char* entry) {
if (S_ISLNK(link_stats.st_mode)) { if (S_ISLNK(link_stats.st_mode)) {
/* A symlink is transferred (following its referent) only when a link /* A symlink is transferred (following its referent) only when a link
resolution option is active, mirroring the regular scanner. */ resolution option is active, mirroring the regular scanner. */
bool resolve = scanner->follow_symlinks || scanner->copy_links || scanner->safe_links || bool resolve = scanner->options.follow_symlinks || scanner->options.copy_links ||
scanner->copy_unsafe_links; scanner->options.safe_links || scanner->options.copy_unsafe_links;
if (!resolve || stat(abs_path, &effective) != 0) { if (!resolve || stat(abs_path, &effective) != 0) {
free(abs_path); free(abs_path);
return NULL; return NULL;
@@ -851,9 +819,9 @@ static File* dirs_file_for_entry(DirectoryScanner* scanner, const char* entry) {
return NULL; return NULL;
} }
} }
if (scanner->use_metadata) { if (scanner->options.use_metadata) {
file->metadata = file_metadata_create(file->path, &effective, scanner->preserve_atimes, file->metadata = file_metadata_create(file->path, &effective, scanner->options.preserve_atimes,
scanner->preserve_crtimes); scanner->options.preserve_crtimes);
if (!file->metadata) { if (!file->metadata) {
file_destroy(file); file_destroy(file);
scanner->failed = true; scanner->failed = true;
@@ -884,18 +852,18 @@ static bool dirs_source_dir_is_empty(const char* path) {
/* The next File from the --dirs generator, or NULL when exhausted. */ /* The next File from the --dirs generator, or NULL when exhausted. */
static File* dirs_next_file(DirectoryScanner* scanner) { static File* dirs_next_file(DirectoryScanner* scanner) {
if (!scanner->file_list) { if (!scanner->options.file_list) {
if (scanner->dirs_root_emitted) if (scanner->dirs_root_emitted)
return NULL; return NULL;
scanner->dirs_root_emitted = true; scanner->dirs_root_emitted = true;
/* --prune-empty-dirs: a physically empty source directory's explicit entry /* --prune-empty-dirs: a physically empty source directory's explicit entry
would only create an empty destination directory, so it is omitted. */ would only create an empty destination directory, so it is omitted. */
if (scanner->prune_empty_dirs && dirs_source_dir_is_empty(scanner->root_path)) if (scanner->options.prune_empty_dirs && dirs_source_dir_is_empty(scanner->root_path))
return NULL; return NULL;
return dirs_root_dir_file(scanner); return dirs_root_dir_file(scanner);
} }
while (scanner->list_index < scanner->file_list->count) { while (scanner->list_index < scanner->options.file_list->count) {
const char* entry = scanner->file_list->entries[scanner->list_index++]; const char* entry = scanner->options.file_list->entries[scanner->list_index++];
File* file = dirs_file_for_entry(scanner, entry); File* file = dirs_file_for_entry(scanner, entry);
if (scanner->failed) if (scanner->failed)
return NULL; return NULL;
@@ -922,8 +890,9 @@ static Chunk* dirs_flush_batch(DirectoryScanner* scanner) {
} }
static Chunk* directory_scanner_next_dirs(DirectoryScanner* scanner) { static Chunk* directory_scanner_next_dirs(DirectoryScanner* scanner) {
while (scanner->dirs_batch == NULL || scanner->dirs_batch_size <= scanner->chunk_size) { while (scanner->dirs_batch == NULL || scanner->dirs_batch_size <= scanner->options.chunk_size) {
if (scanner->stop_condition && stop_condition_reached(scanner->stop_condition)) { if (scanner->options.stop_condition &&
stop_condition_reached(scanner->options.stop_condition)) {
Chunk* leftover = dirs_flush_batch(scanner); Chunk* leftover = dirs_flush_batch(scanner);
if (leftover) if (leftover)
chunk_destroy(leftover); chunk_destroy(leftover);
@@ -962,7 +931,7 @@ static Chunk* directory_scanner_next_dirs(DirectoryScanner* scanner) {
} }
Chunk* directory_scanner_next(DirectoryScanner* scanner) { Chunk* directory_scanner_next(DirectoryScanner* scanner) {
if (scanner && scanner->dirs_mode) if (scanner && scanner->options.dirs)
return directory_scanner_next_dirs(scanner); return directory_scanner_next_dirs(scanner);
ArrayList* chunk_data = array_list_create(file_destroy); ArrayList* chunk_data = array_list_create(file_destroy);
if (!chunk_data) { if (!chunk_data) {
@@ -972,7 +941,8 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
unsigned long long chunk_data_size = 0; unsigned long long chunk_data_size = 0;
while (1) { while (1) {
if (scanner->stop_condition && stop_condition_reached(scanner->stop_condition)) { if (scanner->options.stop_condition &&
stop_condition_reached(scanner->options.stop_condition)) {
array_list_delete(chunk_data); array_list_delete(chunk_data);
return NULL; return NULL;
} }
@@ -996,34 +966,9 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
continue; continue;
ScannerOptions options = {
.use_metadata = scanner->use_metadata,
.chunk_size = scanner->chunk_size,
.exclude_patterns = scanner->exclude_patterns,
.exclude_count = scanner->exclude_count,
.include_patterns = scanner->include_patterns,
.include_count = scanner->include_count,
.max_size = scanner->max_size,
.min_size = scanner->min_size,
.max_depth = scanner->max_depth,
.num_threads = 0,
.follow_symlinks = scanner->follow_symlinks,
.copy_links = scanner->copy_links,
.safe_links = scanner->safe_links,
.copy_unsafe_links = scanner->copy_unsafe_links,
.copy_dirlinks = scanner->copy_dirlinks,
.munge_links = scanner->munge_links,
.checksum = scanner->checksum,
.one_file_system = scanner->one_file_system,
.file_list = scanner->file_list,
.base_filters = scanner->base_filters,
.per_dir_filters = scanner->per_dir_filters,
.dirs = false,
.relative = false,
};
ScannerEntry inspected; ScannerEntry inspected;
int inspection = scanner_inspect_entry(&options, scanner->current_path, scanner->current_path, int inspection = scanner_inspect_entry(&scanner->options, scanner->current_path,
entry->d_name, &inspected); scanner->current_path, entry->d_name, &inspected);
if (inspection < 0) { if (inspection < 0) {
scanner->failed = true; scanner->failed = true;
break; break;
@@ -1055,16 +1000,17 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
scanner->failed = true; scanner->failed = true;
break; break;
} }
bool passes_selection = bool passes_selection = entry_passes_selection(
entry_passes_selection(scanner->file_list, scanner->base_filters, scanner->current_node, scanner->options.file_list, scanner->options.base_filters, scanner->current_node, rel,
rel, entry->d_name, is_dir, scanner->per_dir_filters); entry->d_name, is_dir, scanner->options.per_dir_filters);
if (!passes_selection) { if (!passes_selection) {
/* --files-from subset pruning is not a filter exclusion: its delete /* --files-from subset pruning is not a filter exclusion: its delete
semantics stay keep-set-only (an unlisted source path is treated as semantics stay keep-set-only (an unlisted source path is treated as
absent, so its destination mirror is a deletable extra). A rule-based absent, so its destination mirror is a deletable extra). A rule-based
exclusion is recorded as a protected prefix. -R + --files-from bare exclusion is recorded as a protected prefix. -R + --files-from bare
wire paths are never recorded (see ScannerOptions.excluded_paths). */ wire paths are never recorded (see ScannerOptions.excluded_paths). */
bool files_from_prune = scanner->file_list && !file_list_affects(scanner->file_list, rel); bool files_from_prune =
scanner->options.file_list && !file_list_affects(scanner->options.file_list, rel);
if (!files_from_prune && !scanner->relative_mode) if (!files_from_prune && !scanner->relative_mode)
scanner_record_excluded(scanner, cur_path); scanner_record_excluded(scanner, cur_path);
} }
@@ -1085,12 +1031,13 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
if (is_dir) { if (is_dir) {
free(rel_copy); free(rel_copy);
if (!scanner_same_filesystem(scanner->one_file_system, scanner->root_dev, stats.st_dev)) { if (!scanner_same_filesystem(scanner->options.one_file_system, scanner->root_dev,
stats.st_dev)) {
free(cur_path); free(cur_path);
continue; continue;
} }
int next_depth = scanner->current_depth + 1; int next_depth = scanner->current_depth + 1;
if (scanner->max_depth <= 0 || next_depth < scanner->max_depth) { if (scanner->options.max_depth <= 0 || next_depth < scanner->options.max_depth) {
DirEntry* de = dir_entry_create(cur_path, next_depth, scanner->current_node); DirEntry* de = dir_entry_create(cur_path, next_depth, scanner->current_node);
if (!de || !queue_enqueue(scanner->directories, de)) { if (!de || !queue_enqueue(scanner->directories, de)) {
dir_entry_destroy(de); dir_entry_destroy(de);
@@ -1099,7 +1046,8 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
} }
free(cur_path); free(cur_path);
} else { } else {
if (scanner->max_depth > 0 && scanner->current_depth + 1 > scanner->max_depth) { if (scanner->options.max_depth > 0 &&
scanner->current_depth + 1 > scanner->options.max_depth) {
free(rel_copy); free(rel_copy);
free(cur_path); free(cur_path);
continue; continue;
@@ -1126,13 +1074,14 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
} }
/* --devices/--specials: a device/FIFO/socket entry marked for preservation /* --devices/--specials: a device/FIFO/socket entry marked for preservation
becomes a node to recreate (is_special, no data, rdev captured). */ becomes a node to recreate (is_special, no data, rdev captured). */
scanner_prepare_special(scanner->preserve_devices, scanner->preserve_specials, file, &stats); scanner_prepare_special(scanner->options.preserve_devices, scanner->options.preserve_specials,
if (scanner->hardlinks && S_ISREG(stats.st_mode)) file, &stats);
scanner_assign_hardlink(scanner, scanner->hardlinks, file, &stats); if (scanner->options.hardlinks && S_ISREG(stats.st_mode))
if (scanner->use_metadata) scanner_assign_hardlink(scanner, scanner->options.hardlinks, file, &stats);
file->metadata = file_metadata_create(file->path, &stats, scanner->preserve_atimes, if (scanner->options.use_metadata)
scanner->preserve_crtimes); file->metadata = file_metadata_create(file->path, &stats, scanner->options.preserve_atimes,
if (scanner->use_metadata && !file->metadata) { scanner->options.preserve_crtimes);
if (scanner->options.use_metadata && !file->metadata) {
free(rel_copy); free(rel_copy);
file_destroy(file); file_destroy(file);
scanner->failed = true; scanner->failed = true;
@@ -1147,7 +1096,7 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
break; break;
} }
chunk_data_size += file->data->size; chunk_data_size += file->data->size;
if (chunk_data_size > scanner->chunk_size) { if (chunk_data_size > scanner->options.chunk_size) {
free(rel_copy); free(rel_copy);
Chunk* result = chunk_data_to_chunk(chunk_data); Chunk* result = chunk_data_to_chunk(chunk_data);
if (!result) if (!result)
@@ -1211,7 +1160,7 @@ static int parallel_worker_thread(void* arg) {
free(ds->root_path); free(ds->root_path);
ds->root_path = str_dup(wa->root_dir); ds->root_path = str_dup(wa->root_dir);
ds->seed_node = wa->ps->root_filter_node; ds->seed_node = wa->ps->root_filter_node;
ds->excluded_mutex = &wa->ps->result_mutex; ds->options.excluded_mutex = &wa->ps->result_mutex;
Chunk* chunk; Chunk* chunk;
while ((chunk = directory_scanner_next(ds)) != NULL) { while ((chunk = directory_scanner_next(ds)) != NULL) {
if (!queue_enqueue_multithreaded_cancel(wa->ps->result_queue, chunk, &wa->ps->result_mutex, if (!queue_enqueue_multithreaded_cancel(wa->ps->result_queue, chunk, &wa->ps->result_mutex,
+5 -49
View File
@@ -117,37 +117,16 @@ typedef struct {
typedef struct FilterNode FilterNode; typedef struct FilterNode FilterNode;
typedef struct { typedef struct {
/* Scan inputs, copied once at create time. Everything that is also a
ScannerOptions field lives here (with the normalized chunk_size); only
scanner-owned bookkeeping stays as direct members below. */
ScannerOptions options;
Queue* directories; Queue* directories;
DIR* current_dir; DIR* current_dir;
char* current_path; char* current_path;
bool use_metadata;
bool preserve_atimes;
bool preserve_crtimes;
bool preserve_xattrs;
bool preserve_acls;
unsigned long long chunk_size;
char** exclude_patterns;
int exclude_count;
char** include_patterns;
int include_count;
unsigned long long max_size;
unsigned long long min_size;
int max_depth;
int current_depth; int current_depth;
bool follow_symlinks;
bool copy_links;
bool safe_links;
bool copy_unsafe_links;
bool copy_dirlinks;
bool munge_links;
bool checksum;
bool one_file_system;
dev_t root_dev; dev_t root_dev;
bool failed; bool failed;
/* Phase 4 special/devices (see ScannerOptions). */
bool preserve_devices;
bool preserve_specials;
bool copy_devices;
/* Phase 2 (files-from / filter layer). */ /* Phase 2 (files-from / filter layer). */
char* root_path; /* transfer root (fs path) for rel computation */ char* root_path; /* transfer root (fs path) for rel computation */
char* current_rel; /* rel path of the open directory ("" == root) */ char* current_rel; /* rel path of the open directory ("" == root) */
@@ -155,40 +134,17 @@ typedef struct {
FilterNode* seed_node; /* inherited context of the seed dir, or NULL */ FilterNode* seed_node; /* inherited context of the seed dir, or NULL */
FilterNode* current_node; /* filter context of the open directory */ FilterNode* current_node; /* filter context of the open directory */
ArrayList* filter_nodes; /* owned FilterNode arena (may be NULL) */ ArrayList* filter_nodes; /* owned FilterNode arena (may be NULL) */
const FileListSet* file_list; /* --dirs / -R state for the directory-entry generator (options.dirs replaces
const FilterRuleList* base_filters;
bool per_dir_filters;
/* --dirs / -R state for the directory-entry generator (dirs_mode replaces
the recursive scan). */ the recursive scan). */
bool dirs_mode;
bool relative_mode; /* file_list && relative: send bare relative wire paths */ bool relative_mode; /* file_list && relative: send bare relative wire paths */
bool prune_empty_dirs;
bool dirs_root_emitted; bool dirs_root_emitted;
int list_index; int list_index;
ArrayList* dirs_batch; /* owned when non-NULL */ ArrayList* dirs_batch; /* owned when non-NULL */
unsigned long long dirs_batch_size; unsigned long long dirs_batch_size;
/* Excluded-path sink (see ScannerOptions). `excluded_mutex` is shared across
parallel worker threads. */
ArrayList* excluded_paths;
mtx_t* excluded_mutex;
/* --ignore-errors: continue past unreadable directories (records io_error). */
bool ignore_io_errors;
/* --ignore-missing-args: --dirs listed-but-missing entries are skipped, not
fatal (see ScannerOptions.ignore_missing_args). */
bool ignore_missing_args;
/* A directory could not be opened (I/O error, e.g. EACCES). With /* A directory could not be opened (I/O error, e.g. EACCES). With
--ignore-errors the scan continues past it and the caller decides what to --ignore-errors the scan continues past it and the caller decides what to
do; `failed` is reserved for fatal errors that always abort the scan. */ do; `failed` is reserved for fatal errors that always abort the scan. */
bool io_error; bool io_error;
/* --hard-links (-H): shared link-group detection table (see ScannerOptions).
NULL when -H is off. */
HardLinkTable* hardlinks;
/* Phase 6: sender stop deadline (from ScannerOptions). */
const StopCondition* stop_condition;
/* P7 Wave D directory-time capture (see ScannerOptions). */
bool capture_dir_times;
ArrayList* dir_entries;
mtx_t* dir_entries_mutex;
} DirectoryScanner; } DirectoryScanner;
typedef struct { typedef struct {
+3 -1
View File
@@ -179,6 +179,9 @@ bool pipeline_context_sender_enqueue_chunk(PipelineContextSender* context, Chunk
} }
void pipeline_context_sender_destroy(PipelineContextSender* context) { void pipeline_context_sender_destroy(PipelineContextSender* context) {
/* `config` is borrowed: the caller retains ownership and frees it after the
pipeline has been destroyed (the worker threads are already joined, so no
config access can outlive this call). */
if (context->manifest) { if (context->manifest) {
array_list_delete(context->manifest); array_list_delete(context->manifest);
} }
@@ -192,7 +195,6 @@ void pipeline_context_sender_destroy(PipelineContextSender* context) {
array_list_delete(context->dir_entries); array_list_delete(context->dir_entries);
if (context->dir_entries_mutex_init) if (context->dir_entries_mutex_init)
mtx_destroy(&context->dir_entries_mutex); mtx_destroy(&context->dir_entries_mutex);
config_delete(context->config);
queue_destroy(context->queue_scanner); queue_destroy(context->queue_scanner);
queue_destroy(context->queue_loader); queue_destroy(context->queue_loader);
mtx_destroy(&context->mutex_scanner); mtx_destroy(&context->mutex_scanner);
+2
View File
@@ -120,6 +120,8 @@ typedef struct PipelineContextReceiver {
DirTimeList dir_times; DirTimeList dir_times;
} PipelineContextReceiver; } PipelineContextReceiver;
/* `config` is borrowed and must outlive the context: destroy does NOT free it,
so the caller owns it and frees it with config_delete() afterwards. */
PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner, PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner,
Queue* queue_loader); Queue* queue_loader);
void pipeline_context_sender_destroy(PipelineContextSender* context); void pipeline_context_sender_destroy(PipelineContextSender* context);
+14
View File
@@ -1239,6 +1239,20 @@ class TestProgress:
assert "Stats:" in result.stderr assert "Stats:" in result.stderr
assert "KB" in result.stderr assert "KB" in result.stderr
def test_human_readable_stats_multithreaded(self, shared_server):
# The multithreaded sender shares the single-threaded --stats format,
# including --human-readable and the rate suffix.
clean_dir(DEST_DIR)
result, dur = run_client(
SOURCE_DIR, DEST_DIR,
flags=["--threads", "-h", "--stats"],
port=shared_server.port,
)
assert result.returncode == 0, f"Exit {result.returncode}: {result.stderr[:100]}"
assert "Stats:" in result.stderr
assert "KB" in result.stderr
assert "/s" in result.stderr
def test_human_readable_progress_multithreaded(self, shared_server): def test_human_readable_progress_multithreaded(self, shared_server):
clean_dir(DEST_DIR) clean_dir(DEST_DIR)
result, dur = run_client( result, dur = run_client(
+1
View File
@@ -383,6 +383,7 @@ static void test_pipeline_sender_lifecycle() {
EXPECT_EQ_INT((int)pcs->allocation_session.max_alloc, (int)cfg->max_alloc); EXPECT_EQ_INT((int)pcs->allocation_session.max_alloc, (int)cfg->max_alloc);
pipeline_context_sender_destroy(pcs); pipeline_context_sender_destroy(pcs);
config_delete(cfg); /* the context borrows cfg; the caller owns it */
} }
static void test_pipeline_receiver_lifecycle() { static void test_pipeline_receiver_lifecycle() {
+4 -1
View File
@@ -37,6 +37,7 @@ static void test_sender_create_destroy() {
EXPECT_NULL(ctx->manifest); EXPECT_NULL(ctx->manifest);
pipeline_context_sender_destroy(ctx); pipeline_context_sender_destroy(ctx);
config_delete(cfg); /* the context borrows cfg; the caller owns it */
} }
/* Test pipeline_context_receiver_create/destroy with valid arguments */ /* Test pipeline_context_receiver_create/destroy with valid arguments */
@@ -80,6 +81,7 @@ static void test_sender_queue_capacities() {
EXPECT_EQ_INT(ctx->queue_scanner->capacity, 1); EXPECT_EQ_INT(ctx->queue_scanner->capacity, 1);
EXPECT_EQ_INT(ctx->queue_loader->capacity, 1); EXPECT_EQ_INT(ctx->queue_loader->capacity, 1);
pipeline_context_sender_destroy(ctx); pipeline_context_sender_destroy(ctx);
config_delete(cfg); /* the context borrows cfg; the caller owns it */
} }
/* Invalid queue capacities must not create unusable pipeline queues. */ /* Invalid queue capacities must not create unusable pipeline queues. */
@@ -445,8 +447,9 @@ static void test_sender_enqueue_byte_budget() {
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* second payload now in flight */ EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* second payload now in flight */
/* pipeline_context_sender_destroy frees the still-queued second chunk and /* pipeline_context_sender_destroy frees the still-queued second chunk and
owns cfg/q_scanner/q_loader from here on. */ owns q_scanner/q_loader; cfg stays borrowed and is freed by the caller. */
pipeline_context_sender_destroy(ctx); pipeline_context_sender_destroy(ctx);
config_delete(cfg);
} }
void test_multiprocessing() { void test_multiprocessing() {