Merge dev into human codebase evaluation
CI / lint (pull_request) Failing after 2s
CI / build-and-test (pull_request) Has been skipped
CI / sanitizers (address) (pull_request) Has been skipped
CI / sanitizers (undefined) (pull_request) Has been skipped
CI / fuzz-build (pull_request) Has been skipped
CI / coverage (pull_request) Has been skipped
CI / valgrind (pull_request) Has been skipped

This commit is contained in:
2026-08-11 21:24:17 +02:00
36 changed files with 1584 additions and 471 deletions
+8
View File
@@ -334,6 +334,8 @@ int parse_args(Config* config, int argc, char* argv[], int* positional_args,
} else if (strcmp(argv[i], "-T") == 0 && i + 1 < argc) {
if (set_positive_int_option(&config->timeout, argv[++i], "-T") != 0)
return -1;
} else if (strcmp(argv[i], "--checksum") == 0) {
config->checksum = true;
} else if (strcmp(argv[i], "--compress-level") == 0 && i + 1 < argc) {
if (set_positive_int_option(&config->compression_level, argv[++i], "--compress-level") != 0)
return -1;
@@ -391,6 +393,12 @@ static bool validate_config(const Config* config) {
fprintf(stderr, "Error: --delta cannot be combined with -f (sendfile)\n");
return false;
}
if (config->append || config->append_verify) {
fprintf(
stderr,
"Error: --append and --append-verify are not supported yet; refusing to ignore option\n");
return false;
}
if (config->use_tls) {
if (!config->tls_cert || !config->tls_key) {
fprintf(stderr, "Error: --tls requires --cert and --key\n");
+85 -22
View File
@@ -28,13 +28,27 @@
/* Forward declaration for progress-reporting thread used in multithreaded send. */
static int progress_thread_fn(void* arg);
static void pipeline_cancel(PipelineContextSender* context) {
mtx_lock(&context->mutex_scanner);
mtx_lock(&context->mutex_loader);
atomic_store(&context->cancelled, true);
context->scanner_done = true;
context->loader_done = true;
cnd_broadcast(&context->condition_not_full_scanner);
cnd_broadcast(&context->condition_not_empty_scanner);
cnd_broadcast(&context->condition_not_full_loader);
cnd_broadcast(&context->condition_not_empty_loader);
mtx_unlock(&context->mutex_loader);
mtx_unlock(&context->mutex_scanner);
}
/* Print dry-run manifest showing files that would be transferred. Returns 0 on success. */
static int send_dry_run_manifest(Config* config) {
DirectoryScanner* scanner = directory_scanner_create(
config->send_directory, config->use_metadata, config->chunk_size, config->exclude_patterns,
config->exclude_count, config->include_patterns, config->include_count, config->max_size,
config->min_size, config->max_depth, config->follow_symlinks, config->copy_links,
config->safe_links, config->copy_unsafe_links);
config->safe_links, config->copy_unsafe_links, config->checksum);
if (!scanner)
return -1;
Chunk* chunk;
@@ -67,7 +81,8 @@ static int send_delete_manifest(int fd, ArrayList* manifest) {
return 0;
}
static int incremental_check(Client* client, File* file, DeltaSignature** out_sig) {
static int incremental_check(Client* client, File* file, const Config* config,
DeltaSignature** out_sig) {
*out_sig = NULL;
if (!send_status(client->file_descriptor, STATUS_CHECK))
return -1;
@@ -79,6 +94,12 @@ static int incremental_check(Client* client, File* file, DeltaSignature** out_si
return -1;
if (!send_n_data(client->file_descriptor, &mtime, sizeof(mtime)))
return -1;
if (config->checksum) {
uint64_t checksum;
if (!file_checksum(file, &checksum) ||
!send_n_data(client->file_descriptor, &checksum, sizeof(checksum)))
return -1;
}
Status s;
if (!receive_status(client->file_descriptor, &s))
return -1;
@@ -176,7 +197,7 @@ static int send_single_file(Client* client, File* file, Config* config, bool use
// Incremental path: use sendfile for the actual data if enabled and no compression
if (use_sendfile) {
DeltaSignature* sig = NULL;
int rc = incremental_check(client, file, &sig);
int rc = incremental_check(client, file, config, &sig);
if (rc == 1) {
delta_signature_destroy(sig);
return 1;
@@ -202,7 +223,7 @@ static int send_single_file(Client* client, File* file, Config* config, bool use
// Incremental path with single_calls (supports compression and delta)
file_send_fn send_fn = (file_send_fn)file_send_single_calls;
DeltaSignature* sig = NULL;
int rc = incremental_check(client, file, &sig);
int rc = incremental_check(client, file, config, &sig);
if (rc < 0) {
delta_signature_destroy(sig);
return -1;
@@ -337,6 +358,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
return ok ? thrd_success : thrd_error;
send_fail:
pipeline_cancel(context);
client_disconnect(client);
client_delete(client);
mtx_lock(&context->mutex_progress);
@@ -347,6 +369,7 @@ static int send_chunks_multithreaded(void* pipeline_context) {
if (send_chunk(client, current_chunk, context->config) != 0) {
fprintf(stderr, "Error: unexpected error while sending chunk\n");
chunk_destroy(current_chunk);
pipeline_cancel(context);
client_disconnect(client);
client_delete(client);
mtx_lock(&context->mutex_progress);
@@ -375,9 +398,15 @@ static int scan_directory_multithreaded(void* pipeline_context) {
context->config->exclude_patterns, context->config->exclude_count,
context->config->include_patterns, context->config->include_count, context->config->max_size,
context->config->min_size, context->config->max_depth, 4, context->config->follow_symlinks,
context->config->copy_links, context->config->safe_links, context->config->copy_unsafe_links);
context->config->copy_links, context->config->safe_links, context->config->copy_unsafe_links,
context->config->checksum);
Chunk* current_chunk;
if (scanner == NULL) {
log_message(LOG_LEVEL_ERROR, "Failed to create parallel scanner");
pipeline_cancel(context);
return thrd_error;
}
while ((current_chunk = parallel_scanner_next(scanner)) != NULL) {
if (context->config->use_delete) {
mtx_lock(&context->mutex_scanner);
@@ -389,13 +418,18 @@ static int scan_directory_multithreaded(void* pipeline_context) {
if (!manifest_entry) {
log_message(LOG_LEVEL_ERROR, "Failed to allocate manifest entry");
mtx_unlock(&context->mutex_scanner);
context->cancelled = true;
cnd_broadcast(&context->condition_not_full_scanner);
cnd_broadcast(&context->condition_not_empty_scanner);
pipeline_cancel(context);
parallel_scanner_destroy(scanner);
return thrd_error;
}
if (!array_list_add(context->manifest, manifest_entry)) {
free(manifest_entry);
mtx_unlock(&context->mutex_scanner);
pipeline_cancel(context);
chunk_destroy(current_chunk);
parallel_scanner_destroy(scanner);
return thrd_error;
}
array_list_add(context->manifest, manifest_entry);
}
mtx_unlock(&context->mutex_scanner);
}
@@ -404,13 +438,21 @@ static int scan_directory_multithreaded(void* pipeline_context) {
&context->condition_not_empty_scanner, &context->condition_not_full_scanner,
&context->cancelled)) {
chunk_destroy(current_chunk);
context->cancelled = true;
cnd_broadcast(&context->condition_not_full_scanner);
cnd_broadcast(&context->condition_not_empty_scanner);
pipeline_cancel(context);
parallel_scanner_destroy(scanner);
return thrd_error;
}
}
if (parallel_scanner_failed(scanner)) {
parallel_scanner_destroy(scanner);
mtx_lock(&context->mutex_scanner);
context->scanner_done = true;
cnd_broadcast(&context->condition_not_empty_scanner);
cnd_broadcast(&context->condition_not_full_scanner);
mtx_unlock(&context->mutex_scanner);
pipeline_cancel(context);
return thrd_error;
}
mtx_lock(&context->mutex_scanner);
context->scanner_done = true;
cnd_signal(&context->condition_not_empty_scanner);
@@ -450,7 +492,7 @@ static int load_files_multithreaded(void* pipeline_context) {
&context->condition_not_full_loader,
&context->cancelled)) {
chunk_destroy(chunk);
context->cancelled = true;
atomic_store(&context->cancelled, true);
cnd_broadcast(&context->condition_not_full_loader);
cnd_broadcast(&context->condition_not_empty_loader);
return thrd_error;
@@ -539,13 +581,22 @@ int send_files(Config* config) {
config->send_directory, config->use_metadata, config->chunk_size, config->exclude_patterns,
config->exclude_count, config->include_patterns, config->include_count, config->max_size,
config->min_size, config->max_depth, config->follow_symlinks, config->copy_links,
config->safe_links, config->copy_unsafe_links);
config->safe_links, config->copy_unsafe_links, config->checksum);
Chunk* current_chunk;
unsigned long long total_bytes = 0;
int total_files = 0;
time_t last_progress = 0;
time_t start = time(NULL);
ArrayList* manifest = config->use_delete ? array_list_create(free) : NULL;
if (!scanner || (config->use_delete && !manifest)) {
if (scanner)
directory_scanner_destroy(scanner);
if (manifest)
array_list_delete(manifest);
client_disconnect(client);
client_delete(client);
return 1;
}
while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
unsigned long long chunk_bytes = 0;
for (int i = 0; i < current_chunk->element_count; i++) {
@@ -565,7 +616,15 @@ int send_files(Config* config) {
client_delete(client);
return 1;
}
array_list_add(manifest, manifest_entry);
if (!array_list_add(manifest, manifest_entry)) {
free(manifest_entry);
chunk_destroy(current_chunk);
array_list_delete(manifest);
directory_scanner_destroy(scanner);
client_disconnect(client);
client_delete(client);
return 1;
}
}
}
if (!config->use_sendfile) {
@@ -582,6 +641,9 @@ int send_files(Config* config) {
if (send_chunk(client, current_chunk, config) != 0) {
log_message(LOG_LEVEL_ERROR, "Failed to send chunk");
chunk_destroy(current_chunk);
if (manifest)
array_list_delete(manifest);
manifest = NULL;
break;
}
if (config->show_progress) {
@@ -597,12 +659,15 @@ int send_files(Config* config) {
}
chunk_destroy(current_chunk);
}
if (directory_scanner_failed(scanner) || (config->use_delete && manifest == NULL))
goto send_fail;
if (config->use_delete) {
if (send_delete_manifest(client->file_descriptor, manifest) != 0) {
array_list_delete(manifest);
goto send_fail;
}
array_list_delete(manifest);
manifest = NULL;
}
if (!send_status(client->file_descriptor, STATUS_FINISHED))
goto send_fail;
@@ -624,6 +689,8 @@ int send_files(Config* config) {
return ok ? 0 : 1;
send_fail:
if (manifest)
array_list_delete(manifest);
directory_scanner_destroy(scanner);
client_disconnect(client);
client_delete(client);
@@ -677,14 +744,10 @@ int send_files_multithreaded(Config* config) {
if (!scanner_created || !loader_created || !sender_created) {
perror("Error creating threads.\n");
context->cancelled = true;
context->scanner_done = true;
context->loader_done = true;
pipeline_cancel(context);
mtx_lock(&context->mutex_progress);
context->sender_done = true;
cnd_broadcast(&context->condition_not_full_scanner);
cnd_broadcast(&context->condition_not_empty_scanner);
cnd_broadcast(&context->condition_not_full_loader);
cnd_broadcast(&context->condition_not_empty_loader);
mtx_unlock(&context->mutex_progress);
if (sender_created)
thrd_join(sender, NULL);
if (loader_created)
+192 -25
View File
@@ -11,6 +11,7 @@
#include <sys/stat.h>
#include <threads.h>
#include <unistd.h>
#include <limits.h>
typedef struct {
char* path;
@@ -38,17 +39,34 @@ static DirEntry* dir_entry_create(const char* path, int depth) {
return de;
}
static bool safe_relative_link(const char* source_root, const char* containing_dir,
const char* link_target) {
char root[PATH_MAX];
if (!realpath(source_root, root))
return false;
char* joined = path_cat(containing_dir, link_target);
char resolved[PATH_MAX];
bool safe = joined && realpath(joined, resolved) && strncmp(root, resolved, strlen(root)) == 0 &&
(resolved[strlen(root)] == '\0' || resolved[strlen(root)] == '/');
free(joined);
return safe;
}
DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_metadata,
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,
bool follow_symlinks, bool copy_links, bool safe_links,
bool copy_unsafe_links) {
DirectoryScanner* scanner = malloc(sizeof(DirectoryScanner));
bool copy_unsafe_links, bool checksum) {
DirectoryScanner* scanner = calloc(1, sizeof(DirectoryScanner));
if (scanner == NULL)
return NULL;
scanner->directories = queue_create(100, dir_entry_destroy);
if (!scanner->directories) {
free(scanner);
return NULL;
}
scanner->current_dir = NULL;
scanner->current_path = NULL;
scanner->use_metadata = use_metadata;
@@ -65,6 +83,8 @@ DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_
scanner->copy_links = copy_links;
scanner->safe_links = safe_links;
scanner->copy_unsafe_links = copy_unsafe_links;
scanner->checksum = checksum;
scanner->failed = false;
DirEntry* root = dir_entry_create(root_directory, 0);
if (!root) {
queue_destroy(scanner->directories);
@@ -94,8 +114,12 @@ void directory_scanner_destroy(DirectoryScanner* scanner) {
static Chunk* chunk_data_to_chunk(ArrayList* chunk_data) {
void** chunk_items = array_list_to_array(chunk_data);
if (!chunk_items)
return NULL;
Chunk* chunk = chunk_create((File**)chunk_items, chunk_data->size);
free(chunk_items);
if (!chunk)
return NULL;
chunk_data->item_destroyer = NULL;
array_list_delete(chunk_data);
return chunk;
@@ -120,6 +144,7 @@ static int open_next_directory(DirectoryScanner* scanner) {
perror("Could not open directory");
free(scanner->current_path);
scanner->current_path = NULL;
scanner->failed = true;
return -1;
}
return 1;
@@ -127,6 +152,10 @@ static int open_next_directory(DirectoryScanner* scanner) {
Chunk* directory_scanner_next(DirectoryScanner* scanner) {
ArrayList* chunk_data = array_list_create(file_destroy);
if (!chunk_data) {
scanner->failed = true;
return NULL;
}
unsigned long long chunk_data_size = 0;
while (1) {
@@ -135,7 +164,7 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
if (ret == 0)
break;
if (ret < 0)
continue;
break;
}
struct dirent* entry = readdir(scanner->current_dir);
@@ -151,6 +180,10 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
continue;
char* cur_path = path_cat(scanner->current_path, entry->d_name);
if (!cur_path) {
scanner->failed = true;
break;
}
struct stat stats;
struct stat lstats;
bool is_symlink = false;
@@ -174,7 +207,8 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
continue;
}
link_target[len] = '\0';
if (link_target[0] == '/') {
if (link_target[0] == '/' ||
!safe_relative_link(scanner->current_path, scanner->current_path, link_target)) {
free(cur_path);
continue;
}
@@ -209,8 +243,10 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
int next_depth = scanner->current_depth + 1;
if (scanner->max_depth <= 0 || next_depth < scanner->max_depth) {
DirEntry* de = dir_entry_create(cur_path, next_depth);
if (!queue_enqueue(scanner->directories, de))
if (!de || !queue_enqueue(scanner->directories, de)) {
dir_entry_destroy(de);
scanner->failed = true;
}
}
free(cur_path);
} else {
@@ -253,27 +289,49 @@ Chunk* directory_scanner_next(DirectoryScanner* scanner) {
File* file = file_create(cur_path);
if (file == NULL) {
free(cur_path);
scanner->failed = true;
continue;
}
file->data->size = stats.st_size;
if (scanner->use_metadata)
file->metadata = file_metadata_create(&stats);
array_list_add(chunk_data, file);
if (scanner->use_metadata && !file->metadata) {
file_destroy(file);
free(cur_path);
scanner->failed = true;
break;
}
if (!array_list_add(chunk_data, file)) {
file_destroy(file);
scanner->failed = true;
break;
}
chunk_data_size += file->data->size;
if (chunk_data_size > scanner->chunk_size) {
free(cur_path);
return chunk_data_to_chunk(chunk_data);
Chunk* result = chunk_data_to_chunk(chunk_data);
if (!result)
scanner->failed = true;
return result;
}
free(cur_path);
}
}
if (chunk_data->size > 0)
return chunk_data_to_chunk(chunk_data);
if (chunk_data->size > 0) {
Chunk* result = chunk_data_to_chunk(chunk_data);
if (!result)
scanner->failed = true;
return result;
}
array_list_delete(chunk_data);
return NULL;
}
bool directory_scanner_failed(const DirectoryScanner* scanner) {
return scanner == NULL || scanner->failed;
}
typedef struct {
ParallelScanner* ps;
char** dirs;
@@ -291,6 +349,7 @@ typedef struct {
bool copy_links;
bool safe_links;
bool copy_unsafe_links;
bool checksum;
} ParallelWorkerArg;
static int parallel_worker_thread(void* arg) {
@@ -299,11 +358,32 @@ static int parallel_worker_thread(void* arg) {
DirectoryScanner* ds = directory_scanner_create(
wa->dirs[i], wa->use_metadata, wa->chunk_size, wa->exclude_patterns, wa->exclude_count,
wa->include_patterns, wa->include_count, wa->max_size, wa->min_size, wa->max_depth,
wa->follow_symlinks, wa->copy_links, wa->safe_links, wa->copy_unsafe_links);
wa->follow_symlinks, wa->copy_links, wa->safe_links, wa->copy_unsafe_links, wa->checksum);
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);
break;
}
Chunk* chunk;
while ((chunk = directory_scanner_next(ds)) != NULL) {
queue_enqueue_multithreaded(wa->ps->result_queue, chunk, &wa->ps->result_mutex,
&wa->ps->result_not_empty, &wa->ps->result_not_full);
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);
}
directory_scanner_destroy(ds);
free(wa->dirs[i]);
@@ -313,7 +393,7 @@ static int parallel_worker_thread(void* arg) {
free(wa);
mtx_lock(&ps->result_mutex);
ps->completed++;
if (ps->completed >= ps->num_threads) {
if (ps->completed >= ps->expected_threads) {
ps->done = true;
cnd_signal(&ps->result_not_empty);
}
@@ -327,7 +407,7 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
int include_count, unsigned long long max_size,
unsigned long long min_size, int max_depth,
int num_threads, bool follow_symlinks, bool copy_links,
bool safe_links, bool copy_unsafe_links) {
bool safe_links, bool copy_unsafe_links, bool checksum) {
ParallelScanner* ps = calloc(1, sizeof(ParallelScanner));
if (!ps)
return NULL;
@@ -336,6 +416,7 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
free(ps);
return NULL;
}
atomic_init(&ps->cancelled, false);
int init = 0;
bool ok = true;
if (mtx_init(&ps->result_mutex, mtx_plain) != thrd_success)
@@ -372,6 +453,13 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
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);
closedir(dir);
parallel_scanner_destroy(ps);
return NULL;
}
struct dirent* entry;
while ((entry = readdir(dir)) != NULL) {
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
@@ -401,7 +489,8 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
continue;
}
link_target[len] = 0;
if (link_target[0] == '/') {
if (link_target[0] == '/' ||
!safe_relative_link(root_directory, root_directory, link_target)) {
free(cur_path);
continue;
}
@@ -436,7 +525,10 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
}
if (S_ISDIR(st.st_mode)) {
array_list_add(subdirs, cur_path);
if (!array_list_add(subdirs, cur_path)) {
free(cur_path);
ps->failed = true;
}
} else {
bool excluded = false;
for (int i = 0; i < exclude_count; i++) {
@@ -469,12 +561,22 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
}
File* file = file_create(cur_path);
free(cur_path);
if (!file)
if (!file) {
ps->failed = true;
continue;
}
file->data->size = st.st_size;
if (use_metadata)
file->metadata = file_metadata_create(&st);
array_list_add(root_files, file);
if (use_metadata && !file->metadata) {
file_destroy(file);
ps->failed = true;
continue;
}
if (!array_list_add(root_files, file)) {
file_destroy(file);
ps->failed = true;
}
}
}
closedir(dir);
@@ -482,27 +584,57 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
unsigned long long cs = chunk_size > 0 ? chunk_size : DESIRED_CHUNK_SIZE;
if (root_files->size > 0) {
ArrayList* batch = array_list_create(NULL);
if (!batch) {
ps->failed = true;
array_list_delete(root_files);
array_list_delete(subdirs);
parallel_scanner_destroy(ps);
return NULL;
}
unsigned long long batch_size = 0;
Chunk* first = NULL;
for (int i = 0; i < root_files->size; i++) {
File* f = (File*)root_files->items[i];
array_list_add(batch, f);
if (!array_list_add(batch, f)) {
ps->failed = true;
break;
}
batch_size += f->data->size;
if (batch_size >= cs || i == root_files->size - 1) {
void** items = array_list_to_array(batch);
if (!items) {
ps->failed = true;
batch->item_destroyer = file_destroy;
array_list_delete(batch);
batch = NULL;
break;
}
Chunk* c = chunk_create((File**)items, batch->size);
free(items);
if (!c) {
ps->failed = true;
batch->item_destroyer = file_destroy;
array_list_delete(batch);
batch = NULL;
break;
}
batch->item_destroyer = NULL;
array_list_delete(batch);
batch = NULL;
if (!first) {
first = c;
} else {
queue_enqueue_multithreaded(ps->result_queue, c, &ps->result_mutex, &ps->result_not_empty,
&ps->result_not_full);
if (!queue_enqueue(ps->result_queue, c)) {
chunk_destroy(c);
ps->failed = true;
}
}
if (i < root_files->size - 1) {
batch = array_list_create(NULL);
if (!batch) {
ps->failed = true;
break;
}
batch_size = 0;
}
}
@@ -522,6 +654,7 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
if (subdirs->size > 0) {
ps->num_threads = n;
ps->expected_threads = n;
ps->threads = calloc(n, sizeof(thrd_t));
if (!ps->threads) {
array_list_delete(subdirs);
@@ -531,21 +664,37 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
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)
if (!wa) {
ps->failed = true;
break;
}
wa->ps = ps;
wa->dirs = calloc(count, sizeof(char*));
if (!wa->dirs) {
free(wa);
ps->failed = true;
break;
}
for (int j = 0; j < count; j++)
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->dirs);
free(wa);
ps->failed = true;
break;
}
wa->dir_count = count;
wa->use_metadata = use_metadata;
wa->chunk_size = cs;
@@ -560,15 +709,24 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
wa->copy_links = copy_links;
wa->safe_links = safe_links;
wa->copy_unsafe_links = copy_unsafe_links;
wa->checksum = checksum;
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->dirs);
free(wa);
ps->num_threads = t;
ps->failed = true;
atomic_store(&ps->cancelled, true);
ps->expected_threads = ps->created_threads;
mtx_lock(&ps->result_mutex);
cnd_broadcast(&ps->result_not_empty);
cnd_broadcast(&ps->result_not_full);
mtx_unlock(&ps->result_mutex);
break;
}
ps->num_threads++;
ps->created_threads++;
}
}
array_list_delete(subdirs);
@@ -582,7 +740,9 @@ Chunk* parallel_scanner_next(ParallelScanner* ps) {
return c;
}
if (ps->num_threads == 0) {
mtx_lock(&ps->result_mutex);
ps->done = true;
mtx_unlock(&ps->result_mutex);
return NULL;
}
Chunk* chunk = queue_dequeue_multithreaded(
@@ -590,12 +750,19 @@ Chunk* parallel_scanner_next(ParallelScanner* ps) {
return chunk;
}
bool parallel_scanner_failed(const ParallelScanner* ps) {
return ps == NULL || ps->failed;
}
void parallel_scanner_destroy(ParallelScanner* ps) {
if (!ps)
return;
mtx_lock(&ps->result_mutex);
ps->done = true;
cnd_signal(&ps->result_not_empty);
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);
+11 -2
View File
@@ -6,6 +6,7 @@
#include <dirent.h>
#include <stdbool.h>
#include <threads.h>
#include <stdatomic.h>
typedef struct {
Queue* directories;
@@ -25,6 +26,8 @@ typedef struct {
bool copy_links;
bool safe_links;
bool copy_unsafe_links;
bool checksum;
bool failed;
} DirectoryScanner;
typedef struct {
@@ -33,8 +36,12 @@ typedef struct {
cnd_t result_not_empty;
cnd_t result_not_full;
int num_threads;
int expected_threads;
int created_threads;
thrd_t* threads;
bool done;
bool failed;
atomic_bool cancelled;
int completed;
Chunk* initial_chunk;
} ParallelScanner;
@@ -45,8 +52,9 @@ DirectoryScanner* directory_scanner_create(const char* root_directory, bool use_
int include_count, unsigned long long max_size,
unsigned long long min_size, int max_depth,
bool follow_symlinks, bool copy_links, bool safe_links,
bool copy_unsafe_links);
bool copy_unsafe_links, bool checksum);
Chunk* directory_scanner_next(DirectoryScanner* scanner);
bool directory_scanner_failed(const DirectoryScanner* scanner);
void directory_scanner_destroy(DirectoryScanner* scanner);
ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata,
@@ -55,8 +63,9 @@ ParallelScanner* parallel_scanner_create(char* root_directory, bool use_metadata
int include_count, unsigned long long max_size,
unsigned long long min_size, int max_depth,
int num_threads, bool follow_symlinks, bool copy_links,
bool safe_links, bool copy_unsafe_links);
bool safe_links, bool copy_unsafe_links, bool checksum);
Chunk* parallel_scanner_next(ParallelScanner* scanner);
bool parallel_scanner_failed(const ParallelScanner* scanner);
void parallel_scanner_destroy(ParallelScanner* scanner);
#endif