diff --git a/shell.nix b/shell.nix index 362207d..0f994e2 100644 --- a/shell.nix +++ b/shell.nix @@ -1,4 +1,6 @@ -{ pkgs ? import { } }: +{ + pkgs ? import { }, +}: pkgs.mkShell { nativeBuildInputs = with pkgs; [ @@ -13,6 +15,6 @@ pkgs.mkShell { ]; shellHook = '' - echo "fastSync development environment loaded" + ./tmux.sh ''; } diff --git a/src/client/client.c b/src/client/client.c index 5c48c2e..21c54c0 100644 --- a/src/client/client.c +++ b/src/client/client.c @@ -14,16 +14,15 @@ #include "utils.h" #include -int send_chunk(Client *client, Chunk *chunk, bool use_compression) { - if (use_compression) { - Data *data = chunk_compress(chunk); - send_data(client->file_descriptor, data->data, data->data_size); +int send_chunk(Client *client, Chunk *chunk, Config *config) { + if (config->use_compression) { + Data *data = chunk_compress(chunk, config->compression_level); + send_data(client->file_descriptor, data->data, data->size); } else { for (int i = 0; i < chunk->element_count; i++) { send_status(client->file_descriptor, NEXT); File *file = chunk->items[i]; - send_str(client->file_descriptor, file->path); - send_data(client->file_descriptor, file->data, file->stats.st_size); + file_send_single_calls(file, client->file_descriptor); } } return 0; @@ -76,7 +75,6 @@ int load_files_multithreaded(void *pipeline_context) { int send_chunks_multithreaded(void *pipeline_context) { PipelineContextSender *context = (PipelineContextSender *)pipeline_context; - bool use_compression = context->config->use_compression; Client *client = client_create(); client_connect(client, "127.0.0.1", 8080); config_send(client->file_descriptor, context->config); @@ -92,7 +90,7 @@ int send_chunks_multithreaded(void *pipeline_context) { client_delete(client); return thrd_success; } - if (send_chunk(client, current_chunk, use_compression) != 0) { + if (send_chunk(client, current_chunk, context->config) != 0) { perror("Something unexpected happend while sending the chunk"); exit(EXIT_FAILURE); } @@ -109,7 +107,7 @@ int send_files(Config *config) { while ((current_chunk = directory_scanner_next(scanner)) != NULL) { for (int i = 0; i < current_chunk->element_count; i++) file_load_data(current_chunk->items[i]); - send_chunk(client, current_chunk, config->use_compression); + send_chunk(client, current_chunk, config); chunk_destroy(current_chunk); } send_status(client->file_descriptor, FINISHED); @@ -175,7 +173,7 @@ int send_files_multithreaded(Config *config) { int main(int argc, char *argv[]) { Config *config = config_create( str_dup("1.0.0"), str_dup("/home/taptap/Nextcloud/Uni/moodle/MINT-Raum"), - str_dup("./data_copied"), false, false, false, false, 1); + str_dup("./data_copied"), false, false, false, false, 5, 1); for (int i = 1; i < argc; i++) { handle_arg(argv[i], "-m", &config->use_multithreading, "Enabled Multithreading"); diff --git a/src/server/server.c b/src/server/server.c index ba0f760..744c41a 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -11,8 +11,8 @@ FileReceive *receive_file_receive(int file_descriptor) { char *path = (char *)receive_str(file_descriptor); - DataFragment *file_data_fragment = receive_data(file_descriptor); - FileReceive *file = file_receive_create(path, file_data_fragment); + Data *file_data = receive_data(file_descriptor); + FileReceive *file = file_receive_create(path, file_data); return file; } @@ -53,8 +53,8 @@ int write_thread(void *pipeline_context) { return thrd_success; } if (save_to_disk) - to_disk(path_cat(root_directory, file->path), file->data_fragment->data, - file->data_fragment->size); + to_disk(path_cat(root_directory, file->path), file->data->data, + file->data->size); } } @@ -64,7 +64,7 @@ int receive_files(Config *config, int file_descriptor) { FileReceive *file = receive_file_receive(file_descriptor); if (config->save_to_disk) to_disk(path_cat(config->receive_root_directory, file->path), - file->data_fragment->data, file->data_fragment->size); + file->data->data, file->data->size); file_receive_destroy(file); status = receive_status(file_descriptor); } diff --git a/src/shared/chunk.c b/src/shared/chunk.c index e9d882e..9199dbb 100644 --- a/src/shared/chunk.c +++ b/src/shared/chunk.c @@ -1,11 +1,13 @@ #include #include +#include #include #include #include #include #include "chunk.h" +#include "data.h" #include "log.h" #include "socket.h" @@ -59,6 +61,11 @@ void file_print(void *item) { printf("%s\n", ((File *)item)->path); } +void file_send_single_calls(File *file, int file_descriptor) { + send_str(file_descriptor, file->path); + send_data(file_descriptor, file->data, file->stats.st_size); +} + void file_content_to_buffer(File *file, char *buffer) { if (buffer == NULL) { perror("Buffer is to write file content to is NULL!"); @@ -77,10 +84,10 @@ void file_content_to_buffer(File *file, char *buffer) { fclose(file_pointer); } -FileReceive *file_receive_create(char *path, DataFragment *data_fragment) { +FileReceive *file_receive_create(char *path, Data *data) { FileReceive *file = malloc(sizeof(FileReceive)); file->path = path; - file->data_fragment = data_fragment; + file->data_fragment = data; return file; } @@ -88,7 +95,7 @@ void file_receive_destroy(void *file_receive) { if (file_receive == NULL) return; FileReceive *file = (FileReceive *)file_receive; - data_fragment_delete(file->data_fragment); + data_destroy(file->data_fragment); free(file->path); free(file); } @@ -177,7 +184,8 @@ Data *chunk_format(Chunk *chunk) { return chunk_data_create(data, buffer_size); } -Data *chunk_compress(Chunk *chunk) { +Data *chunk_compress(Chunk *chunk, int compression_level) { + log_message(LOG_LEVEL_DEBUG, "Starting to gather data for chunk compression"); unsigned long long data_size = 0; for (int i = 0; i < chunk->element_count; i++) { data_size += sizeof(unsigned long long); @@ -185,8 +193,13 @@ Data *chunk_compress(Chunk *chunk) { data_size += sizeof(unsigned long long); data_size += chunk->items[i]->stats.st_size; } - char *data = malloc(data_size); - char *data_pointer = data; + Data *data = data_create_empty(data_size); + if (data == NULL) { + log_message(LOG_LEVEL_ERROR, + "Could not allocate memory for chunk compression"); + exit(EXIT_FAILURE); + } + char *data_pointer = data->data; for (int i = 0; i < chunk->element_count; i++) { // path length unsigned long long path_len = strlen(chunk->items[i]->path); @@ -197,19 +210,19 @@ Data *chunk_compress(Chunk *chunk) { // file data unsigned long long data_size = chunk->items[i]->stats.st_size; memcpy(data_pointer, &data_size, sizeof(unsigned long long)); - data_pointer += data_size; + data_pointer += sizeof(unsigned long long); memcpy(data_pointer, chunk->items[i]->data, data_size); data_pointer += data_size; } - size_t compressed_data_size = ZSTD_compressBound(data_size); - void *compressed_data = malloc(compressed_data_size); - if (compressed_data == NULL) { - log_message(ERROR, "Could not allocate memory for compressed Chunk"); - exit(EXIT_FAILURE); - } - // TODO - Data *tmp = malloc(sizeof(Data)); - return tmp; + + log_message(LOG_LEVEL_DEBUG, "Chunk succesfully compressed"); + return compress_data(data, compression_level); +} + +Chunk *chunk_decompress(Data *data, int compression_level) { + // ArrayList *files = array_list_create(file_destroy); + // size_t data_size = ZSTD_getFrameContentSize(const void *src, size_t + // srcSize); } Data *chunk_data_create(void *data, unsigned long long data_size) { @@ -219,7 +232,7 @@ Data *chunk_data_create(void *data, unsigned long long data_size) { exit(EXIT_FAILURE); } chunk_formated->data = data; - chunk_formated->data_size = data_size; + chunk_formated->size = data_size; return chunk_formated; } diff --git a/src/shared/chunk.h b/src/shared/chunk.h index 2557356..7c2e7cd 100644 --- a/src/shared/chunk.h +++ b/src/shared/chunk.h @@ -1,7 +1,7 @@ #ifndef CHUNK_H #define CHUNK_H -#include "socket.h" +#include "data.h" #include #define DESIRED_CHUNK_SIZE 10 * 1024 * 1024 @@ -16,7 +16,7 @@ typedef struct { typedef struct { char *path; - DataFragment *data_fragment; + Data *data; } FileReceive; typedef struct { @@ -24,25 +24,22 @@ typedef struct { int element_count; } Chunk; -typedef struct { - void *data; - unsigned long long data_size; -} Data; - File *file_create(const char *path, struct stat *stats); void file_destroy(void *item); void file_load_data(File *file); void file_print(void *item); +void file_send_single_calls(File *file, int file_descriptor); void file_content_to_buffer(File *file, char *buffer); -FileReceive *file_receive_create(char *path, DataFragment *data_fragment); +FileReceive *file_receive_create(char *path, Data *data); void file_receive_destroy(void *file_receive); Chunk *chunk_create(File **items, int element_count); void chunk_destroy(void *chunk); void chunk_print(void *chunk); Data *chunk_format(Chunk *chunk); -Data *chunk_compress(Chunk *chunk); +Data *chunk_compress(Chunk *chunk, int compression_level); +Chunk *chunk_decompress(Data *data, int compression_level); Data *chunk_data_create(void *data, unsigned long long data_size); void chunk_data_delete(void *chunk); diff --git a/src/shared/config.c b/src/shared/config.c index ce56ec6..a00164f 100644 --- a/src/shared/config.c +++ b/src/shared/config.c @@ -7,7 +7,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, int num_connections) { + bool use_compression, int compression_level, + int num_connections) { Config *config = malloc(sizeof(Config)); config->version = version; @@ -17,6 +18,7 @@ Config *config_create(char *version, char *send_directory, config->use_multithreading = use_multithreading; config->use_chunk_serialization = use_chunk_serialization; config->use_compression = use_compression; + config->compression_level = compression_level; config->num_connections = num_connections; return config; } @@ -36,6 +38,7 @@ void config_send(int file_descriptor, Config *config) { send_int(file_descriptor, config->use_multithreading); send_int(file_descriptor, config->use_chunk_serialization); send_int(file_descriptor, config->use_compression); + send_int(file_descriptor, config->use_compression); send_int(file_descriptor, config->num_connections); if (receive_status(file_descriptor) != OK) { perror("Error transmitting config!"); @@ -52,6 +55,7 @@ Config *config_receive(int file_descriptor) { config->use_multithreading = receive_int(file_descriptor); config->use_chunk_serialization = receive_int(file_descriptor); config->use_compression = receive_int(file_descriptor); + config->compression_level = receive_int(file_descriptor); config->num_connections = receive_int(file_descriptor); send_status(file_descriptor, OK); return config; diff --git a/src/shared/config.h b/src/shared/config.h index 8b846c5..cd8a881 100644 --- a/src/shared/config.h +++ b/src/shared/config.h @@ -11,13 +11,16 @@ typedef struct Config { bool use_multithreading; bool use_chunk_serialization; bool use_compression; + bool use_single_send_per_file; + int compression_level; int num_connections; } Config; 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, int num_connections); + bool use_compression, int compression_level, + int num_connections); void config_delete(Config *config); void config_send(int file_descriptor, Config *config); Config *config_receive(int file_descriptor); diff --git a/src/shared/data.c b/src/shared/data.c new file mode 100644 index 0000000..babef80 --- /dev/null +++ b/src/shared/data.c @@ -0,0 +1,70 @@ + +#include "data.h" +#include "log.h" +#include "stdlib.h" +#include "zstd.h" + +Data *data_create_empty(size_t data_size) { + void *data = malloc(data_size); + if (data == NULL) { + log_message(LOG_LEVEL_ERROR, "Could not allocate memory for empty data"); + exit(EXIT_FAILURE); + } + return data_create(data, data_size); +} + +Data *data_create(void *data, size_t data_size) { + Data *new_data = malloc(sizeof(Data)); + if (new_data == NULL) { + log_message(LOG_LEVEL_ERROR, "Could not allocate memory for data"); + exit(EXIT_FAILURE); + } + new_data->data = data; + new_data->size = data_size; + return data; +} + +void data_destroy(Data *data) { + free(data->data); + free(data); +} + +Data *compress_data(Data *data_to_compress, int compression_level) { + log_message(LOG_LEVEL_DEBUG, "Starting to compress data"); + Data *compressed_data = + data_create_empty(ZSTD_compressBound(data_to_compress->size)); + + compressed_data->size = ZSTD_compress( + compressed_data->data, compressed_data->size, data_to_compress->data, + data_to_compress->size, compression_level); + if (ZSTD_isError(compressed_data->size)) { + log_message(LOG_LEVEL_ERROR, "Compression failed: %s", + ZSTD_getErrorName(compressed_data->size)); + exit(EXIT_FAILURE); + } + + log_message(LOG_LEVEL_DEBUG, "Data succesfully compressed"); + return compressed_data; +} + +Data *decompress_data(Data *compressed_data) { + log_message(LOG_LEVEL_DEBUG, "Start to decompress data"); + Data *uncompressed_data = data_create_empty( + ZSTD_getFrameContentSize(compressed_data->data, compressed_data->size)); + if (ZSTD_isError(uncompressed_data->size)) { + log_message(LOG_LEVEL_ERROR, "Decompression failed: %s", + ZSTD_getErrorName(uncompressed_data->size)); + exit(EXIT_FAILURE); + } + + uncompressed_data->size = + ZSTD_decompress(uncompressed_data->data, uncompressed_data->size, + compressed_data->data, compressed_data->size); + if (ZSTD_isError(uncompressed_data->size)) { + log_message(LOG_LEVEL_ERROR, "Decompression failed: %s", + ZSTD_getErrorName(uncompressed_data->size)); + exit(EXIT_FAILURE); + } + log_message(LOG_LEVEL_DEBUG, "Decompressed data successfully"); + return uncompressed_data; +} diff --git a/src/shared/data.h b/src/shared/data.h new file mode 100644 index 0000000..28d1ce1 --- /dev/null +++ b/src/shared/data.h @@ -0,0 +1,17 @@ +#ifndef DATA_H +#define DATA_H + +#include "stdlib.h" + +typedef struct { + void *data; + size_t size; +} Data; + +Data *data_create_empty(size_t data_size); +Data *data_create(void *data, size_t data_size); +void data_destroy(Data *data); +Data *compress_data(Data *data_to_compress, int compression_level); +Data *decompress_data(Data *compressed_data); + +#endif diff --git a/src/shared/socket.c b/src/shared/socket.c index a747cfe..43918a4 100644 --- a/src/shared/socket.c +++ b/src/shared/socket.c @@ -1,6 +1,7 @@ #include "socket.h" #include "log.h" #include +#include #include #include #include @@ -107,24 +108,9 @@ void client_delete(Client *client) { free(client); } -DataFragment *data_fragment_create(void *data, unsigned long long size) { - DataFragment *data_fragment = malloc(sizeof(DataFragment)); - data_fragment->data = data; - data_fragment->size = size; - return data_fragment; -} - -void data_fragment_delete(void *data_fragment) { - if (data_fragment == NULL) - return; - DataFragment *fragment = (DataFragment *)data_fragment; - free(fragment->data); - free(fragment); -} - -void send_n_data(int file_descriptor, void *data, NET_SIZE data_size) { +void send_n_data(int file_descriptor, void *data, size_t data_size) { log_message(LOG_LEVEL_DEBUG, " Sending n Data: %d", data_size); - NET_SIZE total_bytes_send = 0; + size_t total_bytes_send = 0; while (total_bytes_send < data_size) { long long bytes_send = send(file_descriptor, (char *)data + total_bytes_send, @@ -138,9 +124,9 @@ void send_n_data(int file_descriptor, void *data, NET_SIZE data_size) { log_message(LOG_LEVEL_DEBUG, " Send n Data: %d", total_bytes_send); } -void receive_n_data(int file_descriptor, void *data, NET_SIZE data_size) { +void receive_n_data(int file_descriptor, void *data, size_t data_size) { log_message(LOG_LEVEL_DEBUG, " Receiving n Data: %d", data_size); - NET_SIZE total_bytes_received = 0; + size_t total_bytes_received = 0; while (total_bytes_received < data_size) { long long bytes_received = recv(file_descriptor, data + total_bytes_received, @@ -155,15 +141,15 @@ void receive_n_data(int file_descriptor, void *data, NET_SIZE data_size) { } void send_str(int file_descriptor, char *data) { - NET_SIZE size = strlen(data); - send_n_data(file_descriptor, &size, sizeof(NET_SIZE)); + size_t size = strlen(data); + send_n_data(file_descriptor, &size, sizeof(size_t)); send_n_data(file_descriptor, data, size); log_message(LOG_LEVEL_DEBUG, "Send String: %s", data); } char *receive_str(int file_descriptor) { - NET_SIZE size; - receive_n_data(file_descriptor, &size, sizeof(NET_SIZE)); + size_t size; + receive_n_data(file_descriptor, &size, sizeof(size_t)); char *data = (char *)malloc(size + 1); receive_n_data(file_descriptor, data, size); data[size] = '\0'; @@ -177,13 +163,13 @@ void send_data(int file_descriptor, void *data, unsigned long long data_size) { log_message(LOG_LEVEL_DEBUG, "Send %lld data", data_size); } -DataFragment *receive_data(int file_descriptor) { - unsigned long long size = 0; +Data *receive_data(int file_descriptor) { + size_t size = 0; receive_n_data(file_descriptor, &size, sizeof(unsigned long long)); void *data = malloc(size); receive_n_data(file_descriptor, data, size); log_message(LOG_LEVEL_DEBUG, "Received %lld data", size); - return data_fragment_create(data, size); + return data_create(data, size); } void send_int(int file_descriptor, int data) { diff --git a/src/shared/socket.h b/src/shared/socket.h index 2ca98c7..dc199a0 100644 --- a/src/shared/socket.h +++ b/src/shared/socket.h @@ -1,10 +1,9 @@ #ifndef SOCKET_H #define SOCKET_H -#include "array_list.h" +#include "data.h" #include -typedef unsigned long long NET_SIZE; typedef int Status; enum NET_STATUS { OK, ERROR, FINISHED, NEXT }; @@ -24,25 +23,17 @@ typedef struct Client { int file_descriptor; } Client; -typedef struct DataFragment { - void *data; - unsigned long long size; -} DataFragment; - Client *client_create(); void client_disconnect(Client *client); void client_delete(Client *client); void client_connect(Client *client, char *host, int port); -DataFragment *data_fragment_create(void *data, unsigned long long size); -void data_fragment_delete(void *data_fragment); - -void send_n_data(int file_descriptor, void *data, NET_SIZE data_size); -void receive_n_data(int file_descriptor, void *data, NET_SIZE data_size); +void send_n_data(int file_descriptor, void *data, size_t data_size); +void receive_n_data(int file_descriptor, void *data, size_t data_size); void send_str(int file_descriptor, char *data); char *receive_str(int file_descriptor); void send_data(int file_descriptor, void *data, unsigned long long data_size); -DataFragment *receive_data(int file_descriptor); +Data *receive_data(int file_descriptor); void send_int(int file_descriptor, int data); int receive_int(int file_descriptor); void send_status(int file_descriptor, Status status); diff --git a/tests/test_chunk.c b/tests/test_chunk.c index 118a772..bda51de 100644 --- a/tests/test_chunk.c +++ b/tests/test_chunk.c @@ -1,8 +1,7 @@ #include "test_chunk.h" #include "chunk.h" -#include "utils.h" #include "test_utils.h" -#include +#include "utils.h" #include #include #include @@ -37,7 +36,7 @@ static void test_file_receive_operations() { char *data = str_dup("receive data content"); unsigned long long size = strlen(data); - DataFragment *df = data_fragment_create(data, size); + Data *df = data_create(data, size); EXPECT_NOT_NULL(df); EXPECT_EQ_INT((int)df->size, (int)size); EXPECT_EQ_STR(df->data, "receive data content"); @@ -80,10 +79,10 @@ static void test_chunk_operations() { Data *formatted = chunk_format(chunk); EXPECT_NOT_NULL(formatted); - unsigned long long expected_size = - (sizeof(int) + strlen(path1) + sizeof(unsigned long long) + len1) + - (sizeof(int) + strlen(path2) + sizeof(unsigned long long) + len2); - EXPECT_EQ_INT((int)formatted->data_size, (int)expected_size); + unsigned long long expected_size = + (sizeof(int) + strlen(path1) + sizeof(unsigned long long) + len1) + + (sizeof(int) + strlen(path2) + sizeof(unsigned long long) + len2); + EXPECT_EQ_INT((int)formatted->size, (int)expected_size); char *ptr = (char *)formatted->data; diff --git a/tests/test_config.c b/tests/test_config.c index 4adb110..f5aa760 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, 4); + true, true, false, false, 1, 4); EXPECT_NOT_NULL(cfg); EXPECT_EQ_STR(cfg->version, "1.0"); EXPECT_EQ_STR(cfg->send_directory, "/src"); @@ -23,7 +23,7 @@ static void test_config_lifecycle() { static void test_pipeline_sender_lifecycle() { Config *cfg = config_create(str_dup("2.0"), str_dup("/src2"), - str_dup("/dst2"), false, false, true, true, 8); + str_dup("/dst2"), false, false, true, true, 1, 8); Queue *q1 = queue_create(5, NULL); Queue *q2 = queue_create(15, NULL); @@ -40,7 +40,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, 2); + str_dup("/dst3"), true, true, true, true, 1, 2); Queue *q = queue_create(20, NULL); PipelineContextReceiver *pcr = pipeline_context_receiver_create(cfg, q, 42); diff --git a/tmux.sh b/tmux.sh new file mode 100755 index 0000000..4240838 --- /dev/null +++ b/tmux.sh @@ -0,0 +1,15 @@ +SESSION="fastSync" + +tmux has-session -t $SESSION 2>/dev/null + +if [ $? != 0 ]; then + tmux new-session -d -s $SESSION -n "Neovim" + tmux send-keys -t $SESSION:0 'nvim .' C-m + tmux new-window -t $SESSION -n "Console" + tmux send-keys -t $SESSION:1 'cd ./build' C-m + tmux split-window -h -t $SESSION:1 + tmux send-keys -t $SESSION:1.1 'cd ./build' C-m + tmux select-window -t $SESSION:0 +fi + +tmux attach-session -t $SESSION