diff --git a/src/client/client_cli.c b/src/client/client_cli.c index 8da09eb..d92e452 100644 --- a/src/client/client_cli.c +++ b/src/client/client_cli.c @@ -26,6 +26,7 @@ static void print_usage(void) { printf(" -f Enable sendfile (TCP only, not with -c or -s)\n"); printf(" -v, --verbose Enable debug logging\n"); printf(" -M, --preserve Preserve file metadata\n"); + printf(" --chunk-size Chunk size in bytes (default: %d)\n", DEFAULT_CHUNK_SIZE); printf(" --source-dir Source directory\n"); printf(" --dest-dir Destination directory\n"); printf(" --save-to-disk Write received files to disk\n"); @@ -46,7 +47,7 @@ int main(int argc, char *argv[]) { } Config *config = config_create(str_dup("1.0.0"), NULL, NULL, - save_to_disk, false, false, false, false, 5, false); + save_to_disk, false, false, false, false, 5, false, 0); int positional_args[2]; int positional_count = 0; @@ -93,6 +94,10 @@ int main(int argc, char *argv[]) { server_host = str_dup(argv[++i]); } else if (strcmp(argv[i], "--server-port") == 0 && i + 1 < argc) { server_port = atoi(argv[++i]); + } else if (strcmp(argv[i], "--chunk-size") == 0 && i + 1 < argc) { + unsigned long long val = strtoull(argv[++i], NULL, 10); + if (val > 0) + config->chunk_size = val; } else if (strcmp(argv[i], "-v") == 0 || strcmp(argv[i], "--verbose") == 0) { set_log_level(LOG_LEVEL_DEBUG); } else if (argv[i][0] == '-') { diff --git a/src/client/client_send.c b/src/client/client_send.c index 9556929..81be93a 100644 --- a/src/client/client_send.c +++ b/src/client/client_send.c @@ -82,7 +82,7 @@ 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); + directory_scanner_create(context->config->send_directory, context->config->use_metadata, context->config->chunk_size); mtx_unlock(&context->mutex_scanner); Chunk *current_chunk; @@ -138,7 +138,7 @@ int send_files(Config *config) { client_connect(client, server_host, server_port); } config_send(client->file_descriptor, config); - DirectoryScanner *scanner = directory_scanner_create(config->send_directory, config->use_metadata); + DirectoryScanner *scanner = directory_scanner_create(config->send_directory, config->use_metadata, config->chunk_size); Chunk *current_chunk; while ((current_chunk = directory_scanner_next(scanner)) != NULL) { if (!config->use_sendfile) { diff --git a/src/client/scanner.c b/src/client/scanner.c index d4a331a..854f112 100644 --- a/src/client/scanner.c +++ b/src/client/scanner.c @@ -11,12 +11,13 @@ #include #include -DirectoryScanner *directory_scanner_create(char *root_directory, bool use_metadata) { +DirectoryScanner *directory_scanner_create(char *root_directory, bool use_metadata, unsigned long long chunk_size) { 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; queue_enqueue(scanner->directories, str_dup(root_directory)); return scanner; } @@ -99,7 +100,7 @@ Chunk *directory_scanner_next(DirectoryScanner *scanner) { file->metadata = file_metadata_create(&stats); array_list_add(chunk_data, file); chunk_data_size += file->data->size; - if (chunk_data_size > DESIRED_CHUNK_SIZE) { + if (chunk_data_size > scanner->chunk_size) { free(cur_path); return chunk_data_to_chunk(chunk_data); } diff --git a/src/client/scanner.h b/src/client/scanner.h index 08d1004..d2a53bd 100644 --- a/src/client/scanner.h +++ b/src/client/scanner.h @@ -11,9 +11,10 @@ typedef struct { DIR *current_dir; char *current_path; bool use_metadata; + unsigned long long chunk_size; } DirectoryScanner; -DirectoryScanner *directory_scanner_create(char *root_directory, bool use_metadata); +DirectoryScanner *directory_scanner_create(char *root_directory, bool use_metadata, unsigned long long chunk_size); Chunk *directory_scanner_next(DirectoryScanner *scanner); void directory_scanner_destroy(DirectoryScanner *scanner); diff --git a/src/shared/chunk.h b/src/shared/chunk.h index a411e71..7837936 100644 --- a/src/shared/chunk.h +++ b/src/shared/chunk.h @@ -6,7 +6,7 @@ #include #include -#define DESIRED_CHUNK_SIZE 10 * 1024 * 1024 +#define DESIRED_CHUNK_SIZE (10 * 1024 * 1024) typedef struct { File **items; diff --git a/src/shared/config.c b/src/shared/config.c index a51bdc8..6880be9 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -10,7 +10,8 @@ Config *config_create(char *version, char *send_directory, char *receive_directory, bool save_to_disk, bool use_multithreading, bool use_chunk_serialization, bool use_compression, bool use_metadata, - int compression_level, bool use_sendfile) { + int compression_level, bool use_sendfile, + unsigned long long chunk_size) { Config *config = malloc(sizeof(Config)); config->version = version; @@ -23,6 +24,7 @@ Config *config_create(char *version, char *send_directory, config->use_metadata = use_metadata; config->compression_level = compression_level; config->use_sendfile = use_sendfile; + config->chunk_size = chunk_size > 0 ? chunk_size : DEFAULT_CHUNK_SIZE; config->transport = TRANSPORT_TCP; config->ssh_destination = NULL; return config; @@ -67,6 +69,7 @@ void config_send(int file_descriptor, Config *config) { send_int(file_descriptor, config->use_compression); send_int(file_descriptor, config->use_metadata); send_int(file_descriptor, config->compression_level); + send_int(file_descriptor, (int)config->chunk_size); send_int(file_descriptor, config->use_sendfile); if (receive_status(file_descriptor) != STATUS_OK) { perror("Error transmitting config!"); @@ -85,6 +88,7 @@ Config *config_receive(int file_descriptor) { config->use_compression = receive_int(file_descriptor); config->use_metadata = 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->transport = TRANSPORT_TCP; config->ssh_destination = NULL; diff --git a/src/shared/config.h b/src/shared/config.h index 93358cb..eb30002 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -19,15 +19,19 @@ typedef struct Config { bool use_sendfile; bool use_metadata; int compression_level; + unsigned long long chunk_size; TransportType transport; char *ssh_destination; } Config; +#define DEFAULT_CHUNK_SIZE (10 * 1024 * 1024) + Config *config_create(char *version, char *send_directory, char *receive_directory, bool save_to_disk, bool use_multithreading, bool use_chunk_serialization, bool use_compression, bool use_metadata, - int compression_level, bool use_sendfile); + int compression_level, bool use_sendfile, + unsigned long long chunk_size); void config_delete(Config *config); void config_send(int file_descriptor, Config *config); Config *config_receive(int file_descriptor); diff --git a/src/shared/transport_ssh.c b/src/shared/transport_ssh.c index dd46a00..7d228b9 100644 --- a/src/shared/transport_ssh.c +++ b/src/shared/transport_ssh.c @@ -55,6 +55,12 @@ Client *client_connect_ssh(char *destination) { exit(EXIT_FAILURE); } + int buf_size = 1024 * 1024; + setsockopt(sv[0], SOL_SOCKET, SO_SNDBUF, &buf_size, sizeof(buf_size)); + setsockopt(sv[0], SOL_SOCKET, SO_RCVBUF, &buf_size, sizeof(buf_size)); + setsockopt(sv[1], SOL_SOCKET, SO_SNDBUF, &buf_size, sizeof(buf_size)); + setsockopt(sv[1], SOL_SOCKET, SO_RCVBUF, &buf_size, sizeof(buf_size)); + int exec_pipe[2]; if (pipe(exec_pipe) < 0) { perror("pipe failed"); diff --git a/test.py b/test.py index 26a55b4..02d9a72 100755 --- a/test.py +++ b/test.py @@ -123,7 +123,7 @@ def generate_test_files(source_dir): shutil.rmtree(source_dir) os.makedirs(source_dir) - target_total = 50 * 1024 * 1024 + target_total = 25 * 1024 * 1024 written = 0 files = { diff --git a/tests/test_config.c b/tests/test_config.c index 266fd38..df71e3a 100644 --- a/tests/test_config.c +++ b/tests/test_config.c @@ -8,7 +8,7 @@ static void test_config_lifecycle() { Config *cfg = config_create(str_dup("1.0"), str_dup("/src"), str_dup("/dst"), - true, true, false, false, false, 1, false); + true, true, false, false, false, 1, false, 0); EXPECT_NOT_NULL(cfg); EXPECT_EQ_STR(cfg->version, "1.0"); EXPECT_EQ_STR(cfg->send_directory, "/src"); @@ -24,7 +24,7 @@ static void test_config_lifecycle() { static void test_config_ssh_dest() { Config *cfg = config_create(str_dup("1.0"), str_dup("/src"), str_dup("user@host:/dst"), - true, false, false, false, false, 1, false); + true, false, false, false, false, 1, false, 0); EXPECT_NOT_NULL(cfg); EXPECT_EQ_INT(cfg->transport, TRANSPORT_TCP); EXPECT_NULL(cfg->ssh_destination); @@ -39,7 +39,7 @@ static void test_config_ssh_dest() { static void test_config_ssh_dest_local_path() { Config *cfg = config_create(str_dup("1.0"), str_dup("/src"), str_dup("/local/path"), - true, false, false, false, false, 1, false); + true, false, false, false, false, 1, false, 0); config_parse_ssh_dest(cfg); EXPECT_EQ_INT(cfg->transport, TRANSPORT_TCP); EXPECT_NULL(cfg->ssh_destination); @@ -49,7 +49,7 @@ static void test_config_ssh_dest_local_path() { static void test_config_ssh_dest_no_user() { Config *cfg = config_create(str_dup("1.0"), str_dup("/src"), str_dup("host:/remote"), - true, false, false, false, false, 1, false); + true, false, false, false, false, 1, false, 0); config_parse_ssh_dest(cfg); EXPECT_EQ_INT(cfg->transport, TRANSPORT_SSH); EXPECT_EQ_STR(cfg->ssh_destination, "host:/remote"); @@ -59,7 +59,7 @@ static void test_config_ssh_dest_no_user() { static void test_pipeline_sender_lifecycle() { Config *cfg = config_create(str_dup("2.0"), str_dup("/src2"), - str_dup("/dst2"), false, false, true, true, false, 1, false); + str_dup("/dst2"), false, false, true, true, false, 1, false, 0); Queue *q1 = queue_create(5, NULL); Queue *q2 = queue_create(15, NULL); @@ -76,7 +76,7 @@ static void test_pipeline_sender_lifecycle() { static void test_pipeline_receiver_lifecycle() { Config *cfg = config_create(str_dup("3.0"), str_dup("/src3"), - str_dup("/dst3"), true, true, true, true, false, 1, false); + str_dup("/dst3"), true, true, true, true, false, 1, false, 0); Queue *q = queue_create(20, NULL); PipelineContextReceiver *pcr = pipeline_context_receiver_create(cfg, q, 42); diff --git a/tests/test_scanner.c b/tests/test_scanner.c index 6187cbe..e29ab67 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); + DirectoryScanner *scanner = directory_scanner_create((char *)dir, false, 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); + DirectoryScanner *scanner = directory_scanner_create((char *)dir, false, 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); + DirectoryScanner *scanner = directory_scanner_create((char *)root, false, 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); + DirectoryScanner *scanner = directory_scanner_create((char *)dir, false, 0); EXPECT_NOT_NULL(scanner); Chunk *chunk = directory_scanner_next(scanner);