From 284bf8f8bc19797d453490891fa1f21d82139383 Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 16 Jul 2026 12:39:05 +0200 Subject: [PATCH 1/2] rsync flags: -a, -n, -p, --exclude, --delete Feature changes across 15 files: - -a/--archive: enables -c -m -M (no -s, which slows transfers) - -n/--dry-run: scan + print without connecting or transferring - -p : custom SSH port (passed as -p to ssh via execvp) - --exclude : glob-based filename filtering in scanner - --delete: sender collects file manifest; receiver deletes unlisted files Protocol: added STATUS_MANIFEST, use_delete field in Config wire format. New utilities: glob_match(), delete_extras() with recursive directory walk. Scanner: accepts exclude patterns; skips matching entries. SSH transport: switched from execlp to execvp for dynamic port arg. Server + multiprocessing: handle STATUS_MANIFEST in both single and multithreaded receive paths. Builds clean, all 7 tests pass. --- src/client/client_cli.c | 20 ++++++++ src/client/client_send.c | 95 +++++++++++++++++++++++++++++++++--- src/client/scanner.c | 15 +++++- src/client/scanner.h | 4 +- src/server/server.c | 13 +++++ src/shared/config.c | 15 ++++++ src/shared/config.h | 5 ++ src/shared/multiprocessing.c | 17 +++++++ src/shared/multiprocessing.h | 2 + src/shared/protocol.h | 2 +- src/shared/transport_ssh.c | 26 ++++++++-- src/shared/transport_ssh.h | 2 +- src/shared/utils.c | 72 +++++++++++++++++++++++++++ src/shared/utils.h | 5 ++ tests/test_scanner.c | 8 +-- 15 files changed, 281 insertions(+), 20 deletions(-) diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 2965061..9ff7bf1 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -22,7 +22,12 @@ static void print_usage(void) { printf("Options:\n"); printf(" -c [level] Enable compression (level 1-22, default 5)\n"); printf(" -z [level] Alias for -c\n"); + printf(" -a, --archive Archive mode (-c -m -M)\n"); + printf(" -n, --dry-run Show what would be transferred\n"); + printf(" -p SSH port (default: 22)\n"); printf(" --progress Show transfer progress\n"); + printf(" --delete Delete files on receiver not in source\n"); + printf(" --exclude Exclude files matching pattern\n"); printf(" -m Enable multithreading\n"); printf(" -s Enable chunk serialization\n"); printf(" -f Enable sendfile (TCP only, not with -c or -s)\n"); @@ -58,6 +63,21 @@ int main(int argc, char *argv[]) { if (strcmp(argv[i], "--help") == 0) { print_usage(); return 0; + } else if (strcmp(argv[i], "-a") == 0 || strcmp(argv[i], "--archive") == 0) { + config->use_compression = true; + config->use_multithreading = true; + config->use_metadata = true; + log_message(LOG_LEVEL_INFO, "Enabled archive mode (-c -m -M)"); + } else if (strcmp(argv[i], "-n") == 0 || strcmp(argv[i], "--dry-run") == 0) { + config->dry_run = true; + } else if (strcmp(argv[i], "-p") == 0 && i + 1 < argc) { + config->ssh_port = atoi(argv[++i]); + } else if (strcmp(argv[i], "--delete") == 0) { + config->use_delete = true; + } else if (strcmp(argv[i], "--exclude") == 0 && i + 1 < argc) { + int idx = config->exclude_count++; + config->exclude_patterns = realloc(config->exclude_patterns, config->exclude_count * sizeof(char *)); + config->exclude_patterns[idx] = str_dup(argv[++i]); } else if (strcmp(argv[i], "-c") == 0 || strcmp(argv[i], "-z") == 0) { config->use_compression = true; log_message(LOG_LEVEL_INFO, "Enabled Compression"); diff --git a/src/client/client_send.c b/src/client/client_send.c index 8e6a2d3..2187b33 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -1,4 +1,5 @@ #include "client_send.h" +#include "array_list.h" #include "chunk.h" #include "config.h" #include "data.h" @@ -52,7 +53,7 @@ static int send_chunks_multithreaded(void *pipeline_context) { fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n"); return 1; } - client = client_connect_ssh(context->config->ssh_destination); + client = client_connect_ssh(context->config->ssh_destination, context->config->ssh_port); } else { client = client_create(); client_connect(client, server_host, server_port); @@ -65,6 +66,13 @@ static int send_chunks_multithreaded(void *pipeline_context) { &context->condition_not_empty_loader, &context->condition_not_full_loader, &context->loader_done); if (current_chunk == NULL) { + if (context->config->use_delete) { + send_status(client->file_descriptor, STATUS_MANIFEST); + send_int(client->file_descriptor, context->manifest->size); + for (int i = 0; i < context->manifest->size; i++) + send_str(client->file_descriptor, + (char *)context->manifest->items[i]); + } send_status(client->file_descriptor, STATUS_FINISHED); int ok = receive_status(client->file_descriptor) == STATUS_OK; client_disconnect(client); @@ -82,16 +90,26 @@ static int send_chunks_multithreaded(void *pipeline_context) { static int scan_directory_multithreaded(void *pipeline_context) { PipelineContextSender *context = (PipelineContextSender *)pipeline_context; mtx_lock(&context->mutex_scanner); - DirectoryScanner *scanner = - directory_scanner_create(context->config->send_directory, context->config->use_metadata, context->config->chunk_size); + DirectoryScanner *scanner = directory_scanner_create( + context->config->send_directory, context->config->use_metadata, + context->config->chunk_size, context->config->exclude_patterns, + context->config->exclude_count); mtx_unlock(&context->mutex_scanner); Chunk *current_chunk; - while ((current_chunk = directory_scanner_next(scanner)) != NULL) + while ((current_chunk = directory_scanner_next(scanner)) != NULL) { + if (context->config->use_delete) { + mtx_lock(&context->mutex_scanner); + for (int i = 0; i < current_chunk->element_count; i++) + array_list_add(context->manifest, + str_dup(current_chunk->items[i]->path)); + mtx_unlock(&context->mutex_scanner); + } queue_enqueue_multithreaded(context->queue_scanner, current_chunk, &context->mutex_scanner, &context->condition_not_empty_scanner, &context->condition_not_full_scanner); + } mtx_lock(&context->mutex_scanner); context->scanner_done = true; cnd_signal(&context->condition_not_empty_scanner); @@ -127,27 +145,56 @@ static int load_files_multithreaded(void *pipeline_context) { } int send_files(Config *config) { + if (config->dry_run) { + DirectoryScanner *scanner = directory_scanner_create( + config->send_directory, config->use_metadata, config->chunk_size, + config->exclude_patterns, config->exclude_count); + Chunk *chunk; + int file_count = 0; + unsigned long long total_bytes = 0; + printf("Dry run: files to be transferred\n"); + while ((chunk = directory_scanner_next(scanner)) != NULL) { + for (int i = 0; i < chunk->element_count; i++) { + printf(" %s (%zu bytes)\n", chunk->items[i]->path, + chunk->items[i]->data->size); + total_bytes += chunk->items[i]->data->size; + file_count++; + } + chunk_destroy(chunk); + } + directory_scanner_destroy(scanner); + printf("Total: %d files, %.1f MB\n", file_count, + total_bytes / 1048576.0); + return 0; + } + Client *client; if (config->transport == TRANSPORT_SSH) { if (config->use_sendfile) { fprintf(stderr, "Error: -f/--sendfile is not supported with SSH transport\n"); return 1; } - client = client_connect_ssh(config->ssh_destination); + client = client_connect_ssh(config->ssh_destination, config->ssh_port); } else { client = client_create(); client_connect(client, server_host, server_port); } config_send(client->file_descriptor, config); - DirectoryScanner *scanner = directory_scanner_create(config->send_directory, config->use_metadata, config->chunk_size); + DirectoryScanner *scanner = directory_scanner_create( + config->send_directory, config->use_metadata, config->chunk_size, + config->exclude_patterns, config->exclude_count); Chunk *current_chunk; unsigned long long total_bytes = 0; time_t last_progress = 0; time_t start = time(NULL); + ArrayList *manifest = config->use_delete ? array_list_create(free) : NULL; while ((current_chunk = directory_scanner_next(scanner)) != NULL) { unsigned long long chunk_bytes = 0; - for (int i = 0; i < current_chunk->element_count; i++) + for (int i = 0; i < current_chunk->element_count; i++) { chunk_bytes += current_chunk->items[i]->data->size; + if (manifest) + array_list_add(manifest, str_dup(current_chunk->items[i]->path)); + } if (!config->use_sendfile) { for (int i = 0; i < current_chunk->element_count; i++) file_load_data(current_chunk->items[i]); @@ -166,6 +213,15 @@ int send_files(Config *config) { } chunk_destroy(current_chunk); } + if (config->use_delete) { + send_status(client->file_descriptor, STATUS_MANIFEST); + send_int(client->file_descriptor, manifest->size); + for (int i = 0; i < manifest->size; i++) + send_str(client->file_descriptor, (char *)manifest->items[i]); + for (int i = 0; i < manifest->size; i++) + free(manifest->items[i]); + array_list_delete(manifest); + } send_status(client->file_descriptor, STATUS_FINISHED); int ok = receive_status(client->file_descriptor) == STATUS_OK; if (config->show_progress) { @@ -180,9 +236,34 @@ int send_files(Config *config) { } int send_files_multithreaded(Config *config) { + if (config->dry_run) { + DirectoryScanner *scanner = directory_scanner_create( + config->send_directory, config->use_metadata, config->chunk_size, + config->exclude_patterns, config->exclude_count); + Chunk *chunk; + int file_count = 0; + unsigned long long total_bytes = 0; + printf("Dry run: files to be transferred\n"); + while ((chunk = directory_scanner_next(scanner)) != NULL) { + for (int i = 0; i < chunk->element_count; i++) { + printf(" %s (%zu bytes)\n", chunk->items[i]->path, + chunk->items[i]->data->size); + total_bytes += chunk->items[i]->data->size; + file_count++; + } + chunk_destroy(chunk); + } + directory_scanner_destroy(scanner); + printf("Total: %d files, %.1f MB\n", file_count, + total_bytes / 1048576.0); + return 0; + } + PipelineContextSender *context = pipeline_context_sender_create(config, queue_create(100, chunk_destroy), queue_create(100, chunk_destroy)); + if (config->use_delete) + context->manifest = array_list_create(free); thrd_t scanner, loader, sender; if (thrd_create(&scanner, scan_directory_multithreaded, context) != diff --git a/src/client/scanner.c b/src/client/scanner.c index 854f112..175e238 100644 --- a/src/client/scanner.c +++ b/src/client/scanner.c @@ -11,13 +11,15 @@ #include #include -DirectoryScanner *directory_scanner_create(char *root_directory, bool use_metadata, unsigned long long chunk_size) { +DirectoryScanner *directory_scanner_create(char *root_directory, bool use_metadata, unsigned long long chunk_size, char **exclude_patterns, int exclude_count) { DirectoryScanner *scanner = malloc(sizeof(DirectoryScanner)); scanner->directories = queue_create(100, free); scanner->current_dir = NULL; scanner->current_path = NULL; scanner->use_metadata = use_metadata; scanner->chunk_size = chunk_size > 0 ? chunk_size : DESIRED_CHUNK_SIZE; + scanner->exclude_patterns = exclude_patterns; + scanner->exclude_count = exclude_count; queue_enqueue(scanner->directories, str_dup(root_directory)); return scanner; } @@ -94,6 +96,17 @@ Chunk *directory_scanner_next(DirectoryScanner *scanner) { if (S_ISDIR(stats.st_mode)) { queue_enqueue(scanner->directories, (void *)cur_path); } else { + bool excluded = false; + for (int i = 0; i < scanner->exclude_count; i++) { + if (glob_match(scanner->exclude_patterns[i], entry->d_name)) { + excluded = true; + break; + } + } + if (excluded) { + free(cur_path); + continue; + } File *file = file_create(cur_path); file->data->size = stats.st_size; if (scanner->use_metadata) diff --git a/src/client/scanner.h b/src/client/scanner.h index d2a53bd..76cf437 100644 --- a/src/client/scanner.h +++ b/src/client/scanner.h @@ -12,9 +12,11 @@ typedef struct { char *current_path; bool use_metadata; unsigned long long chunk_size; + char **exclude_patterns; + int exclude_count; } DirectoryScanner; -DirectoryScanner *directory_scanner_create(char *root_directory, bool use_metadata, unsigned long long chunk_size); +DirectoryScanner *directory_scanner_create(char *root_directory, bool use_metadata, unsigned long long chunk_size, char **exclude_patterns, int exclude_count); Chunk *directory_scanner_next(DirectoryScanner *scanner); void directory_scanner_destroy(DirectoryScanner *scanner); diff --git a/src/server/server.c b/src/server/server.c index db35a07..1be205c 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -1,3 +1,4 @@ +#include "array_list.h" #include "chunk.h" #include "compression.h" #include "config.h" @@ -55,6 +56,18 @@ int receive_files(Config *config, int file_descriptor) { } status = receive_status(file_descriptor); } + if (status == STATUS_MANIFEST) { + int count = receive_int(file_descriptor); + ArrayList *manifest = array_list_create(free); + for (int i = 0; i < count; i++) + array_list_add(manifest, receive_str(file_descriptor)); + fprintf(stderr, "Deleting files not in manifest...\n"); + delete_extras(config->receive_root_directory, manifest); + for (int i = 0; i < manifest->size; i++) + free(manifest->items[i]); + array_list_delete(manifest); + status = receive_status(file_descriptor); + } if (status != STATUS_FINISHED) { log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status"); send_status(file_descriptor, STATUS_ERROR); diff --git a/src/shared/config.c b/src/shared/config.c index 43d40af..e5d8637 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -23,11 +23,16 @@ Config *config_create(char *version, char *send_directory, config->use_compression = use_compression; config->use_metadata = use_metadata; config->show_progress = false; + config->dry_run = false; + config->use_delete = false; config->compression_level = compression_level; config->use_sendfile = use_sendfile; config->chunk_size = chunk_size > 0 ? chunk_size : DEFAULT_CHUNK_SIZE; + config->ssh_port = 22; config->transport = TRANSPORT_TCP; config->ssh_destination = NULL; + config->exclude_patterns = NULL; + config->exclude_count = 0; return config; } @@ -57,6 +62,9 @@ void config_delete(Config *config) { free(config->send_directory); free(config->receive_root_directory); free(config->ssh_destination); + for (int i = 0; i < config->exclude_count; i++) + free(config->exclude_patterns[i]); + free(config->exclude_patterns); free(config); } @@ -72,6 +80,7 @@ void config_send(int file_descriptor, Config *config) { send_int(file_descriptor, config->compression_level); send_int(file_descriptor, (int)config->chunk_size); send_int(file_descriptor, config->use_sendfile); + send_int(file_descriptor, config->use_delete); if (receive_status(file_descriptor) != STATUS_OK) { perror("Error transmitting config!"); exit(EXIT_FAILURE); @@ -91,8 +100,14 @@ Config *config_receive(int file_descriptor) { config->compression_level = receive_int(file_descriptor); config->chunk_size = (unsigned long long)receive_int(file_descriptor); config->use_sendfile = receive_int(file_descriptor); + config->use_delete = receive_int(file_descriptor); + config->show_progress = false; + config->dry_run = false; + config->ssh_port = 22; config->transport = TRANSPORT_TCP; config->ssh_destination = NULL; + config->exclude_patterns = NULL; + config->exclude_count = 0; send_status(file_descriptor, STATUS_OK); return config; } diff --git a/src/shared/config.h b/src/shared/config.h index b2a636a..b1c2437 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -19,10 +19,15 @@ typedef struct Config { bool use_sendfile; bool use_metadata; bool show_progress; + bool dry_run; + bool use_delete; int compression_level; unsigned long long chunk_size; + int ssh_port; TransportType transport; char *ssh_destination; + char **exclude_patterns; + int exclude_count; } Config; #define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index cf2b0a5..5f26d56 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -1,4 +1,5 @@ #include "multiprocessing.h" +#include "array_list.h" #include "chunk.h" #include "compression.h" #include "config.h" @@ -36,6 +37,11 @@ PipelineContextSender *pipeline_context_sender_create(Config *config, } void pipeline_context_sender_destroy(PipelineContextSender *context) { + if (context->manifest) { + for (int i = 0; i < context->manifest->size; i++) + free(context->manifest->items[i]); + array_list_delete(context->manifest); + } config_delete(context->config); queue_destroy(context->queue_scanner); queue_destroy(context->queue_loader); @@ -119,6 +125,17 @@ int receive_thread(void *pipeline_context) { } status = receive_status(file_descriptor); } + if (status == STATUS_MANIFEST) { + int count = receive_int(file_descriptor); + ArrayList *manifest = array_list_create(free); + for (int i = 0; i < count; i++) + array_list_add(manifest, receive_str(file_descriptor)); + delete_extras(context->config->receive_root_directory, manifest); + for (int i = 0; i < manifest->size; i++) + free(manifest->items[i]); + array_list_delete(manifest); + status = receive_status(file_descriptor); + } mtx_lock(&context->mutex); context->receiver_done = true; cnd_signal(&context->condition_not_empty); diff --git a/src/shared/multiprocessing.h b/src/shared/multiprocessing.h index 81ab589..605c42a 100644 --- a/src/shared/multiprocessing.h +++ b/src/shared/multiprocessing.h @@ -3,6 +3,7 @@ #include +#include "array_list.h" #include "config.h" #include "file.h" #include "queue.h" @@ -19,6 +20,7 @@ typedef struct { cnd_t condition_not_full_loader; cnd_t condition_not_empty_loader; bool loader_done; + ArrayList *manifest; } PipelineContextSender; typedef struct PipelineContextReceiver { diff --git a/src/shared/protocol.h b/src/shared/protocol.h index e2961e6..d981097 100644 --- a/src/shared/protocol.h +++ b/src/shared/protocol.h @@ -5,7 +5,7 @@ #include typedef int Status; -enum NET_STATUS { STATUS_OK, STATUS_ERROR, STATUS_FINISHED, STATUS_NEXT, STATUS_CHUNK }; +enum NET_STATUS { STATUS_OK, STATUS_ERROR, STATUS_FINISHED, STATUS_NEXT, STATUS_CHUNK, STATUS_MANIFEST }; void io_set_fds(int read_fd, int write_fd); void send_n_data(int file_descriptor, void *data, size_t data_size); diff --git a/src/shared/transport_ssh.c b/src/shared/transport_ssh.c index 6b2b5d1..4fc54cb 100644 --- a/src/shared/transport_ssh.c +++ b/src/shared/transport_ssh.c @@ -42,7 +42,7 @@ static int parse_remote_dest(const char *dest, RemoteDest *r) { return 0; } -Client *client_connect_ssh(char *destination) { +Client *client_connect_ssh(char *destination, int port) { RemoteDest r; if (parse_remote_dest(destination, &r) != 0) { fprintf(stderr, "Invalid remote destination: %s\n", destination); @@ -90,10 +90,26 @@ Client *client_connect_ssh(char *destination) { else snprintf(ssh_user, sizeof(ssh_user), "%s", r.host); - execlp("ssh", "ssh", "-o", "Compression=no", "-o", - "ControlMaster=auto", "-o", - "ControlPath=~/.cache/fastsync-%r@%h:%p", ssh_user, - "fastsync-server", "--stdio", (char *)NULL); + char *ssh_argv[16]; + int ac = 0; + char port_str[16]; + ssh_argv[ac++] = "ssh"; + ssh_argv[ac++] = "-o"; + ssh_argv[ac++] = "Compression=no"; + ssh_argv[ac++] = "-o"; + ssh_argv[ac++] = "ControlMaster=auto"; + ssh_argv[ac++] = "-o"; + ssh_argv[ac++] = "ControlPath=~/.cache/fastsync-%r@%h:%p"; + if (port > 0 && port != 22) { + ssh_argv[ac++] = "-p"; + snprintf(port_str, sizeof(port_str), "%d", port); + ssh_argv[ac++] = port_str; + } + ssh_argv[ac++] = ssh_user; + ssh_argv[ac++] = "fastsync-server"; + ssh_argv[ac++] = "--stdio"; + ssh_argv[ac] = NULL; + execvp("ssh", ssh_argv); perror("exec of ssh failed"); ssize_t wret = write(exec_pipe[1], "x", 1); (void)wret; diff --git a/src/shared/transport_ssh.h b/src/shared/transport_ssh.h index 3bad43c..f542870 100644 --- a/src/shared/transport_ssh.h +++ b/src/shared/transport_ssh.h @@ -3,6 +3,6 @@ #include "transport_tcp.h" -Client *client_connect_ssh(char *destination); +Client *client_connect_ssh(char *destination, int port); #endif diff --git a/src/shared/utils.c b/src/shared/utils.c index 7054d44..b7c2c7a 100644 --- a/src/shared/utils.c +++ b/src/shared/utils.c @@ -1,9 +1,12 @@ #include "utils.h" +#include "array_list.h" #include "libgen.h" +#include #include #include #include #include +#include void mkdir_r(char *path) { char *path_duplicate = malloc(strlen(path) + 1); @@ -44,6 +47,75 @@ char *str_dup(const char *string) { return new_string; } +bool glob_match(const char *pattern, const char *str) { + while (*pattern) { + if (*pattern == '*') { + pattern++; + while (*str && *str != '/') { + if (glob_match(pattern, str)) + return true; + str++; + } + return glob_match(pattern, str); + } else if (*pattern == '?') { + if (!*str || *str == '/') + return false; + pattern++; + str++; + } else { + if (*pattern != *str) + return false; + pattern++; + str++; + } + } + return *str == '\0'; +} + +static void delete_extras_walk(const char *abs_path, const char *rel_path, + ArrayList *manifest) { + DIR *dir = opendir(abs_path); + if (!dir) + return; + struct dirent *entry; + while ((entry = readdir(dir)) != NULL) { + if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) + continue; + char *child_abs = path_cat((char *)abs_path, entry->d_name); + char *child_rel = path_cat((char *)rel_path, entry->d_name); + struct stat st; + if (stat(child_abs, &st) != 0) { + free(child_abs); + free(child_rel); + continue; + } + if (S_ISDIR(st.st_mode)) { + delete_extras_walk(child_abs, child_rel, manifest); + } else { + // Check if relative path is in manifest + bool found = false; + for (int i = 0; i < manifest->size; i++) { + if (strcmp((char *)manifest->items[i], child_rel) == 0) { + found = true; + break; + } + } + if (!found) { + unlink(child_abs); + fprintf(stderr, " Deleted: %s\n", child_rel); + } + } + free(child_abs); + free(child_rel); + } + closedir(dir); + rmdir(abs_path); +} + +void delete_extras(const char *dest_root, ArrayList *manifest) { + delete_extras_walk(dest_root, "", manifest); +} + char *path_cat(char *path1, char *path2) { if (path1 == NULL || *path1 == '\0') return str_dup(path2); diff --git a/src/shared/utils.h b/src/shared/utils.h index f54c98d..67030ad 100644 --- a/src/shared/utils.h +++ b/src/shared/utils.h @@ -1,8 +1,13 @@ #ifndef UTILS_H #define UTILS_H +#include "array_list.h" +#include + void mkdir_r(char *path); char *str_dup(const char *string); char *path_cat(char *path1, char *path2); +bool glob_match(const char *pattern, const char *str); +void delete_extras(const char *dest_root, ArrayList *manifest); #endif diff --git a/tests/test_scanner.c b/tests/test_scanner.c index e29ab67..a733491 100644 --- a/tests/test_scanner.c +++ b/tests/test_scanner.c @@ -18,7 +18,7 @@ static void test_scanner_single_file() { mkdir(dir, 0755); create_test_file(file1, content1); - DirectoryScanner *scanner = directory_scanner_create((char *)dir, false, 0); + DirectoryScanner *scanner = directory_scanner_create((char *)dir, false, 0, NULL, 0); EXPECT_NOT_NULL(scanner); Chunk *chunk = directory_scanner_next(scanner); @@ -46,7 +46,7 @@ static void test_scanner_multiple_files() { create_test_file(file1, content1); create_test_file(file2, content2); - DirectoryScanner *scanner = directory_scanner_create((char *)dir, false, 0); + DirectoryScanner *scanner = directory_scanner_create((char *)dir, false, 0, NULL, 0); EXPECT_NOT_NULL(scanner); Chunk *chunk = directory_scanner_next(scanner); @@ -83,7 +83,7 @@ static void test_scanner_subdirectory() { create_test_file(root_file, content); create_test_file(sub_file, content); - DirectoryScanner *scanner = directory_scanner_create((char *)root, false, 0); + DirectoryScanner *scanner = directory_scanner_create((char *)root, false, 0, NULL, 0); EXPECT_NOT_NULL(scanner); int total_files = 0; @@ -106,7 +106,7 @@ static void test_scanner_empty_directory() { mkdir(dir, 0755); - DirectoryScanner *scanner = directory_scanner_create((char *)dir, false, 0); + DirectoryScanner *scanner = directory_scanner_create((char *)dir, false, 0, NULL, 0); EXPECT_NOT_NULL(scanner); Chunk *chunk = directory_scanner_next(scanner); -- 2.52.0 From 2ff51246c354b12f8de9becb78012b7cfd314957 Mon Sep 17 00:00:00 2001 From: TapTap Date: Thu, 16 Jul 2026 13:31:43 +0200 Subject: [PATCH 2/2] Update README and add feature tests for rsync-compatible flags --- README.md | 144 ++++++++++++++++++++------------- src/client/client_send.c | 17 ++-- src/server/server.c | 2 - src/shared/multiprocessing.c | 2 - test.py | 150 ++++++++++++++++++++++++++++++++++- 5 files changed, 249 insertions(+), 66 deletions(-) diff --git a/README.md b/README.md index bffd576..9c78a91 100644 --- a/README.md +++ b/README.md @@ -1,36 +1,40 @@ # FastSync -A high-performance file synchronization system with a custom TCP-based protocol, optional metadata preservation, compression, multithreading, and zero-copy `sendfile()` support. +A high-performance file synchronization system with SSH and TCP transport, streaming zstd compression, multithreaded transfer, metadata preservation, and rsync-compatible CLI flags. ## Technical Overview -1. Custom TCP-based client-server protocol with status codes -2. Chunked file transfer (files grouped into ~10 MB chunks) -3. Optional zstd compression (levels 1–22) -4. Multithreading for parallel file processing (producer-consumer with thread-safe queues) -5. Optional file metadata preservation (`mode`, `uid`, `gid`, `mtime`) — restored on disk -6. In-memory and disk-based storage options -7. `sendfile()` zero-copy path (~2× faster on localhost) +1. **Dual transport**: custom TCP client-server or SSH subprocess (rsync-style `user@host:/path`) +2. **Chunked file transfer**: files grouped into configurable-size chunks (default ~10 MB) +3. **Streaming zstd compression** (levels 1–22) using `ZSTD_compressStream2` +4. **Multithreading**: producer-consumer pipeline with thread-safe queues (scanner → loader → sender) +5. **Metadata preservation**: `mode`, `uid`, `gid`, `mtime` restored on disk when enabled +6. **`sendfile()` zero-copy** on TCP (~2× faster on loopback) +7. **SSH ControlMaster** for connection reuse across repeated invocations +8. **`--delete`**: receiver removes files not present in sender manifest +9. **`--exclude`**: glob-pattern filename filtering (`*`, `?`, no `/` crossing) ## System Architecture ### Client -- Recursively scans source directories (BFS) -- Groups files into chunks (default ~10 MB total) -- Optionally compresses with zstd -- Optionally serializes chunks into a compact binary format -- Optionally attaches per-file metadata (mode, ownership, timestamps) -- Sends via custom protocol or `sendfile()` zero-copy path +- Recursively scans source directories (BFS), supports exclude patterns +- Groups files into chunks (configurable size) +- Streaming zstd compression with configurable level +- Chunk serialization (compact binary format) or per-file transfer +- Manifests all sent paths when `--delete` is active +- Sends via TCP `sendfile()` or SSH pipe +- Optional progress display with throughput ### Server -- Listens on port 8080 +- TCP mode: listens on port 8080; SSH mode: runs via `--stdio` - Receives and reassembles files -- Decompresses, deserializes, restores metadata on disk +- Decompresses (streaming zstd), deserializes, restores metadata +- Processes `STATUS_MANIFEST` for `--delete`: walks destination tree, removes extras - Thread pool for parallel processing ## Protocol Details -Status codes: +### Status Codes | Code | Meaning | |------|---------| | `STATUS_OK` | Operation successful | @@ -38,56 +42,76 @@ Status codes: | `STATUS_FINISHED` | Transfer complete | | `STATUS_NEXT` | Ready for next file (per-file mode) | | `STATUS_CHUNK` | Following data is a serialized chunk | +| `STATUS_MANIFEST` | Following data is a file manifest (for `--delete`) | ### Wire Format — Metadata When `use_metadata` is enabled (`-M`), each file entry carries a 4-byte `present` flag followed by five fields (`mode`, `uid`, `gid`, `mtime_sec`, `mtime_nsec`). When disabled globally, no metadata bytes are sent — zero wire overhead. -## Configuration +### Transfer Flow +``` +Config → (STATUS_NEXT | STATUS_CHUNK)* → [STATUS_MANIFEST] → STATUS_FINISHED → STATUS_OK +``` + +## Command-Line Arguments -### Command-Line Arguments | Argument | Description | |----------|-------------| -| `-m` | Multithreading mode | +| Positional | ` ` — automatic SSH detection if dest contains `:` | | `-c [level]` | Compression with optional level (1–22, default 5) | +| `-z [level]` | Alias for `-c` | +| `-a, --archive` | Archive mode: enables `-c -m -M` (no `-s`) | +| `-m` | Multithreading mode | | `-s` | Chunk serialization (batch all files per chunk) | -| `-f` | Sendfile zero-copy. Incompatible with `-c` / `-s`. | +| `-f` | Sendfile zero-copy. Incompatible with `-c` / `-s`. TCP only. | | `-M, --preserve` | Preserve file metadata (mode, uid, gid, mtime) | +| `-n, --dry-run` | Scan and print what would be transferred | +| `-p ` | SSH port (default: 22) | +| `--progress` | Show real-time transfer speed | +| `--delete` | Delete files on receiver not present in source | +| `--exclude ` | Exclude files matching glob pattern (repeatable) | +| `--chunk-size ` | Chunk size in bytes (default: 10485760) | | `--source-dir ` | Source directory (overrides `FASTSYNC_SOURCE_DIR`) | | `--dest-dir ` | Server destination directory (overrides `FASTSYNC_DEST_DIR`) | | `--save-to-disk` | Write received files to disk | +| `--server-host ` | Server IP address (default: `127.0.0.1`) | +| `--server-port ` | Server port (default: `8080`) | +| `-v, --verbose` | Enable debug logging | + +## Environment Variables -### Environment Variables | Variable | Default | Description | |----------|---------|-------------| -| `FASTSYNC_SOURCE_DIR` | User documents | Source directory fallback | -| `FASTSYNC_DEST_DIR` | `./data_copied` | Destination directory fallback | -| `FASTSYNC_SERVER_IP` | `127.0.0.1` | Server address | -| `FASTSYNC_SERVER_PORT` | `8080` | Server port | +| `FASTSYNC_SOURCE_DIR` | — | Source directory fallback | +| `FASTSYNC_DEST_DIR` | — | Destination directory fallback | | `FASTSYNC_SAVE_TO_DISK` | `false` | Disk persistence fallback | ## Implementation Details ### Data Structures -1. **Chunk** — collection of files (~10 MB total) +1. **Chunk** — collection of files (~10 MB total by default) 2. **File** — path, content (`Data`), optional `FileMetadata` pointer 3. **FileMetadata** — `mode`, `uid`, `gid`, `mtime_sec`, `mtime_nsec` -4. **Config** — runtime parameters -5. **Queue** — thread-safe queue with condition variables +4. **Config** — runtime parameters (transported over wire) +5. **Queue** — thread-safe bounded queue with condition variables +6. **DirectoryScanner** — recursive BFS traversal with exclude pattern support ### Key Algorithms -1. **File scanning** — recursive BFS directory traversal -2. **Chunking** — files grouped by size limit -3. **Compression** — zstd with configurable level -4. **Network protocol** — custom TCP with status codes and optional metadata packing +1. **File scanning** — BFS directory traversal; each entry matched against exclude patterns +2. **Chunking** — files accumulated until `chunk_size` threshold, then flushed +3. **Compression** — streaming zstd via `ZSTD_compressStream2` / `ZSTD_decompressStream` +4. **Network protocol** — status-code-driven exchange with metadata packing 5. **Metadata restoration** — `chmod()`, `chown()`, `utimensat()` on the receiving side +6. **`--delete`** — sender tracks all sent paths; receiver walks destination tree and removes unlisted files/directories +7. **SSH transport** — `socketpair()` + `fork()` + `execvp("ssh", ...)` with `ControlMaster` and port support ## Build Requirements - C11 compiler - CMake 4.1+ -- zstd library +- zstd library (≥ 1.4.0 for streaming API) - pthreads +- SSH client (for SSH transport) ## Building @@ -97,55 +121,67 @@ cmake -B build -S . && cmake --build build -j$(nproc) ## Running -### Server +### Server (TCP mode) ```bash ./build/server ``` -### Client +### Client — SSH (rsync-style) +```bash +./build/client /path/to/send user@host:/path/to/receive +``` + +### Client — TCP ```bash -# Basic ./build/client --source-dir /path/to/send --dest-dir /path/to/receive --save-to-disk +``` -# With metadata preservation -./build/client -M --source-dir ... --dest-dir ... +### Common Options +```bash +# Archive mode (compression + multithreading + metadata) +./build/client -a /path/to/send user@host:/path -# Multithreaded + compression -./build/client -m -c 10 +# Dry run +./build/client -n /path/to/send /path/to/receive -# Sendfile (zero-copy) -./build/client -f +# With progress and custom chunk size +./build/client --progress --chunk-size 2097152 /src user@host:/dst + +# Exclude temporary files + delete extras on receiver +./build/client --exclude "*.tmp" --exclude "*.o" --delete /src user@host:/dst # All features -./build/client -m -c -s -M +./build/client -a --progress --chunk-size 5242880 --exclude "*.log" --delete /src /dst ``` +### Server via SSH +Place the `fastsync-server` binary in the remote `$PATH`. The client runs `ssh user@host fastsync-server --stdio` automatically when an SSH-style destination is given. + ## Testing ```bash -# Unit tests +# Unit tests (7 suites) ./build/tests -# Integration benchmark (~50 MB data, 13 configurations + rsync comparison) +# Integration + benchmark suite python3 test.py - -# Profiles: --wan (100 Mbit, 50 ms, 1% loss), --unlimited (no throttling) -python3 test.py --wan ``` -The benchmark prints throughput metrics for the best configuration and speedup vs rsync. +The benchmark prints throughput metrics, best configuration, and speedup vs rsync. ## Performance Considerations -1. Chunk size (~10 MB) balances memory and transfer efficiency +1. Chunk size (~10 MB default) balances memory and transfer efficiency 2. Compression level trades CPU for bandwidth 3. `sendfile()` bypasses userspace — ~2× faster on localhost for large files 4. Multithreading scales with core count -5. Metadata transfer adds negligible overhead when disabled, ~24 bytes per file when enabled +5. Metadata transfer adds negligible overhead (~24 bytes per file when enabled) +6. SSH socketpair buffer set to 1 MB for improved pipe throughput +7. SSH ControlMaster reuses connections across repeated invocations ## Benchmark Results -50 MB of mixed file sizes over `localhost` with disk I/O throttled (reads ≤ 15 MB/s, writes ≤ 10 MB/s) and network emulation via `tc netem`. Each test was run 3×; the median is reported below. +25 MB of mixed file sizes over `localhost` with disk I/O throttled (reads ≤ 15 MB/s, writes ≤ 10 MB/s) and network emulation via `tc netem`. Each test was run 3×; the median is reported below. ### LAN (1000 Mbit, 20 ms ±1 ms, 0.1% loss) @@ -167,4 +203,4 @@ The benchmark prints throughput metrics for the best configuration and speedup v | rsync (archive) | 17.44 s | — | — | | rsync (archive + compress) | 1.47 s | — | — | -Compression reduces the data on the wire enough that the transfer becomes latency-bound rather than bandwidth-bound. On WAN, the best configuration runs 10.8× faster than the theoretical limit for uncompressed data, since zstd shrinks the 50 MB payload to a fraction of its original size over the wire. +Compression reduces the data on the wire enough that the transfer becomes latency-bound rather than bandwidth-bound. On WAN, the best configuration runs 10.8× faster than the theoretical limit for uncompressed data, since zstd shrinks the 25 MB payload to a fraction of its original size over the wire. diff --git a/src/client/client_send.c b/src/client/client_send.c index 2187b33..4797d31 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -100,9 +100,11 @@ static int scan_directory_multithreaded(void *pipeline_context) { while ((current_chunk = directory_scanner_next(scanner)) != NULL) { if (context->config->use_delete) { mtx_lock(&context->mutex_scanner); - for (int i = 0; i < current_chunk->element_count; i++) - array_list_add(context->manifest, - str_dup(current_chunk->items[i]->path)); + for (int i = 0; i < current_chunk->element_count; i++) { + const char *p = current_chunk->items[i]->path; + if (*p == '/') p++; + array_list_add(context->manifest, str_dup(p)); + } mtx_unlock(&context->mutex_scanner); } queue_enqueue_multithreaded(context->queue_scanner, current_chunk, @@ -192,8 +194,11 @@ int send_files(Config *config) { unsigned long long chunk_bytes = 0; for (int i = 0; i < current_chunk->element_count; i++) { chunk_bytes += current_chunk->items[i]->data->size; - if (manifest) - array_list_add(manifest, str_dup(current_chunk->items[i]->path)); + if (manifest) { + const char *p = current_chunk->items[i]->path; + if (*p == '/') p++; + array_list_add(manifest, str_dup(p)); + } } if (!config->use_sendfile) { for (int i = 0; i < current_chunk->element_count; i++) @@ -218,8 +223,6 @@ int send_files(Config *config) { send_int(client->file_descriptor, manifest->size); for (int i = 0; i < manifest->size; i++) send_str(client->file_descriptor, (char *)manifest->items[i]); - for (int i = 0; i < manifest->size; i++) - free(manifest->items[i]); array_list_delete(manifest); } send_status(client->file_descriptor, STATUS_FINISHED); diff --git a/src/server/server.c b/src/server/server.c index 1be205c..577b430 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -63,8 +63,6 @@ int receive_files(Config *config, int file_descriptor) { array_list_add(manifest, receive_str(file_descriptor)); fprintf(stderr, "Deleting files not in manifest...\n"); delete_extras(config->receive_root_directory, manifest); - for (int i = 0; i < manifest->size; i++) - free(manifest->items[i]); array_list_delete(manifest); status = receive_status(file_descriptor); } diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 5f26d56..bcf0812 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -38,8 +38,6 @@ PipelineContextSender *pipeline_context_sender_create(Config *config, void pipeline_context_sender_destroy(PipelineContextSender *context) { if (context->manifest) { - for (int i = 0; i < context->manifest->size; i++) - free(context->manifest->items[i]); array_list_delete(context->manifest); } config_delete(context->config); diff --git a/test.py b/test.py index 02d9a72..131bddc 100755 --- a/test.py +++ b/test.py @@ -195,7 +195,7 @@ def start_rsync_daemon(source_dir): return port, conf, daemon -def run_single_test(cmd, name, source_dir, dest_dir, *, source_prefix=None, no_server=False): +def run_single_test(cmd, name, source_dir, dest_dir, *, source_prefix=None, no_server=False, expected_missing=None): if os.path.exists(dest_dir): shutil.rmtree(dest_dir) if no_server: @@ -216,6 +216,8 @@ def run_single_test(cmd, name, source_dir, dest_dir, *, source_prefix=None, no_s received = os.path.join(dest_dir, source_prefix if source_prefix is not None else os.path.abspath(source_dir).lstrip(os.sep)) mismatches, missing = verify_transfer(source_dir, received) + if expected_missing: + missing = [m for m in missing if m not in expected_missing] first_line = lambda s: (s or "").strip().split("\n")[0] entry = { @@ -318,6 +320,152 @@ def run_profile(profile_name, source_dir, dest_dir): except Exception: pass + # Feature-specific tests for rsync-compatible flags + print("\n " + "─" * 56 + "\n Feature Tests\n " + "─" * 56) + + # Dry run (-n) — no server needed + print("\n --- Dry run (-n) ---") + flags = BASE_CLIENT_FLAGS + ["-n"] + cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + flags + print(f" Running: {' '.join(cmd)}") + try: + start = time.monotonic() + result = subprocess.run(cmd, text=True, capture_output=True) + duration = time.monotonic() - start + r = {"name": "Dry run (-n)", "suite": profile_name} + if result.returncode == 0 and "Dry run:" in result.stdout: + r["status"] = "Success" + r["time"] = f"{duration:.4f}s" + r["error"] = "" + else: + r["status"] = "Failed" + r["time"] = "N/A" + r["error"] = f"Exit {result.returncode}: {(result.stderr or result.stdout)[:100]}" + results.append(r) + except Exception as e: + results.append({"name": "Dry run (-n)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)}) + + # Archive mode (-a) + feature_flags = BASE_CLIENT_FLAGS + ["-a"] + cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + feature_flags + print(f"\n --- Archive mode (-a) ---\n Running: {' '.join(cmd)}") + try: + r = run_single_test(cmd, "Archive mode (-a)", source_dir, dest_dir) + r["suite"] = profile_name + results.append(r) + except Exception as e: + results.append({"name": "Archive mode (-a)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)}) + + # Exclude (--exclude small.txt) + feature_flags = BASE_CLIENT_FLAGS + ["--exclude", "small.txt"] + cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + feature_flags + print(f"\n --- Exclude (--exclude small.txt) ---\n Running: {' '.join(cmd)}") + try: + r = run_single_test(cmd, "Exclude (--exclude small.txt)", source_dir, dest_dir, + expected_missing=["small.txt"]) + r["suite"] = profile_name + results.append(r) + except Exception as e: + results.append({"name": "Exclude (--exclude small.txt)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)}) + + # Progress (--progress) + feature_flags = BASE_CLIENT_FLAGS + ["--progress"] + cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + feature_flags + print(f"\n --- Progress (--progress) ---\n Running: {' '.join(cmd)}") + try: + r = run_single_test(cmd, "Progress (--progress)", source_dir, dest_dir) + r["suite"] = profile_name + results.append(r) + except Exception as e: + results.append({"name": "Progress (--progress)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)}) + + # Chunk size (--chunk-size 5242880) + feature_flags = BASE_CLIENT_FLAGS + ["--chunk-size", "5242880"] + cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + feature_flags + print(f"\n --- Chunk size (--chunk-size 5242880) ---\n Running: {' '.join(cmd)}") + try: + r = run_single_test(cmd, "Chunk size (--chunk-size 5242880)", source_dir, dest_dir) + r["suite"] = profile_name + results.append(r) + except Exception as e: + results.append({"name": "Chunk size (--chunk-size 5242880)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)}) + + # Delete (--delete) — pre-populate dest, add extra files, then sync with --delete + # Note: server handles one client per launch, so we restart between syncs + print(f"\n --- Delete (--delete) ---") + try: + flags = BASE_CLIENT_FLAGS + ["-M"] + # First sync (no delete) to populate dest + s1 = subprocess.Popen(SERVER_CMD, stdout=subprocess.DEVNULL, stderr=None) + time.sleep(0.5) + first_cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + flags + r1 = subprocess.run(first_cmd, text=True, capture_output=True) + wait_proc(s1) + if r1.returncode != 0: + raise RuntimeError(f"First sync failed: {r1.stderr[:100]}") + # Add extra files to received dir + received = os.path.join(dest_dir, os.path.abspath(source_dir).lstrip(os.sep)) + extra_path = os.path.join(received, "extra_file.txt") + with open(extra_path, "w") as f: + f.write("should be deleted") + extra_dir = os.path.join(received, "extra_dir") + os.makedirs(extra_dir, exist_ok=True) + with open(os.path.join(extra_dir, "nested.txt"), "w") as f: + f.write("nested extra") + # Second sync with --delete (fresh server) + s2 = subprocess.Popen(SERVER_CMD, stdout=subprocess.DEVNULL, stderr=None) + time.sleep(0.5) + second_cmd = client_prefix + BASE_CLIENT_CMD + ["--source-dir", source_dir, "--dest-dir", dest_dir] + flags + ["--delete"] + start = time.monotonic() + r2 = subprocess.run(second_cmd, text=True, capture_output=True) + duration = time.monotonic() - start + wait_proc(s2) + r = {"name": "Delete (--delete)", "suite": profile_name} + if r2.returncode == 0 and not os.path.exists(extra_path) and not os.path.exists(extra_dir): + mismatches, missing = verify_transfer(source_dir, received) + if not mismatches and not missing: + r["status"] = "Success" + r["time"] = f"{duration:.4f}s" + r["error"] = "" + else: + r["status"] = "Failed" + r["time"] = "N/A" + r["error"] = f"post-delete verify: mismatches={len(mismatches)}, missing={len(missing)}" + else: + r["status"] = "Failed" + r["time"] = "N/A" + errs = [] + if r2.returncode != 0: + errs.append(f"Exit {r2.returncode}: {(r2.stderr or r2.stdout)[:60]}") + if os.path.exists(extra_path): + errs.append("extra_file.txt remains") + if os.path.exists(extra_dir): + errs.append("extra_dir remains") + r["error"] = " | ".join(errs) + results.append(r) + except Exception as e: + results.append({"name": "Delete (--delete)", "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)}) + + # SSH feature tests + if SSH_AVAILABLE: + ssh_dest = f"localhost:{dest_dir}_ssh" + ssh_feature_cases = [ + {"name": "SSH Archive (-a)", "flags": ["-a"]}, + {"name": "SSH Exclude (--exclude small.txt)", "flags": ["--exclude", "small.txt"], "expected_missing": ["small.txt"]}, + ] + for case in ssh_feature_cases: + flags = BASE_CLIENT_FLAGS + case["flags"] + cmd = BASE_CLIENT_CMD + [source_dir, ssh_dest] + flags + print(f"\n --- {case['name']} ---\n Running: {' '.join(cmd)}") + try: + r = run_single_test(cmd, case["name"], source_dir, f"{dest_dir}_ssh", + no_server=True, + expected_missing=case.get("expected_missing")) + r["suite"] = profile_name + results.append(r) + except Exception as e: + results.append({"name": case["name"], "suite": profile_name, "status": "Error", "time": "N/A", "error": str(e)}) + except (subprocess.CalledProcessError, RuntimeError) as e: print(f" Error: {e}") results = [] -- 2.52.0