Compare commits
5 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 9414ba3ade | |||
| ded5527543 | |||
| 78d2037314 | |||
| 37ae474b31 | |||
| e55b1f921a |
@@ -1,2 +1,4 @@
|
||||
build
|
||||
data_copied
|
||||
test_data/
|
||||
__pycache__/
|
||||
|
||||
+19
-2
@@ -17,9 +17,18 @@
|
||||
#include <dirent.h>
|
||||
|
||||
int send_chunk(Client *client, Chunk *chunk, Config *config) {
|
||||
if (config->use_compression && config->use_chunk_serialization) {
|
||||
Data *data = chunk_compress(chunk, config->compression_level);
|
||||
if (config->use_chunk_serialization) {
|
||||
send_status(client->file_descriptor, STATUS_CHUNK);
|
||||
Data *data;
|
||||
if (config->use_compression) {
|
||||
data = chunk_compress(chunk, config->compression_level);
|
||||
} else {
|
||||
for (int i = 0; i < chunk->element_count; i++)
|
||||
file_load_data(chunk->items[i]);
|
||||
data = chunk_serialize(chunk);
|
||||
}
|
||||
send_data(client->file_descriptor, data->data, data->size);
|
||||
data_destroy(data);
|
||||
} else {
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
send_status(client->file_descriptor, STATUS_NEXT);
|
||||
@@ -198,6 +207,14 @@ int main(int argc, char *argv[]) {
|
||||
config->compression_level);
|
||||
}
|
||||
}
|
||||
} else if (strcmp(argv[i], "--source-dir") == 0 && i + 1 < argc) {
|
||||
free(config->send_directory);
|
||||
config->send_directory = str_dup(argv[++i]);
|
||||
} else if (strcmp(argv[i], "--dest-dir") == 0 && i + 1 < argc) {
|
||||
free(config->receive_root_directory);
|
||||
config->receive_root_directory = str_dup(argv[++i]);
|
||||
} else if (strcmp(argv[i], "--save-to-disk") == 0) {
|
||||
config->save_to_disk = true;
|
||||
} else {
|
||||
handle_arg(argv[i], "-m", &config->use_multithreading,
|
||||
"Enabled Multithreading");
|
||||
|
||||
+53
-18
@@ -7,10 +7,14 @@
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/stat.h>
|
||||
#include <unistd.h>
|
||||
|
||||
DirectoryScanner *directory_scanner_create(char *root_directory) {
|
||||
DirectoryScanner *scanner = malloc(sizeof(DirectoryScanner));
|
||||
scanner->directories = queue_create(100, free);
|
||||
scanner->current_dir = NULL;
|
||||
scanner->current_path = NULL;
|
||||
queue_enqueue(scanner->directories, str_dup(root_directory));
|
||||
return scanner;
|
||||
}
|
||||
@@ -18,11 +22,16 @@ DirectoryScanner *directory_scanner_create(char *root_directory) {
|
||||
void directory_scanner_destroy(DirectoryScanner *scanner) {
|
||||
if (scanner == NULL)
|
||||
return;
|
||||
if (scanner->current_dir) {
|
||||
closedir(scanner->current_dir);
|
||||
scanner->current_dir = NULL;
|
||||
}
|
||||
free(scanner->current_path);
|
||||
queue_destroy(scanner->directories);
|
||||
free(scanner);
|
||||
}
|
||||
|
||||
Chunk *chunk_data_to_chunk(ArrayList *chunk_data) {
|
||||
static Chunk *chunk_data_to_chunk(ArrayList *chunk_data) {
|
||||
void **chunk_items = array_list_to_array(chunk_data);
|
||||
Chunk *chunk = chunk_create((File **)chunk_items, chunk_data->size);
|
||||
free(chunk_items);
|
||||
@@ -31,29 +40,57 @@ Chunk *chunk_data_to_chunk(ArrayList *chunk_data) {
|
||||
return chunk;
|
||||
}
|
||||
|
||||
static int open_next_directory(DirectoryScanner *scanner) {
|
||||
if (scanner->current_dir) {
|
||||
closedir(scanner->current_dir);
|
||||
scanner->current_dir = NULL;
|
||||
}
|
||||
free(scanner->current_path);
|
||||
|
||||
if (queue_is_empty(scanner->directories))
|
||||
return 0;
|
||||
|
||||
scanner->current_path = (char *)queue_dequeue(scanner->directories);
|
||||
scanner->current_dir = opendir(scanner->current_path);
|
||||
if (scanner->current_dir == NULL) {
|
||||
perror("Could not open directory!");
|
||||
exit(EXIT_FAILURE);
|
||||
}
|
||||
return 1;
|
||||
}
|
||||
|
||||
Chunk *directory_scanner_next(DirectoryScanner *scanner) {
|
||||
ArrayList *chunk_data = array_list_create(file_destroy);
|
||||
unsigned long long chunk_data_size = 0;
|
||||
|
||||
while (!queue_is_empty(scanner->directories)) {
|
||||
char *path = (char *)queue_dequeue(scanner->directories);
|
||||
DIR *dir;
|
||||
struct dirent *entry;
|
||||
dir = opendir(path);
|
||||
if (dir == NULL) {
|
||||
perror("Could not open directory!");
|
||||
exit(EXIT_FAILURE);
|
||||
while (1) {
|
||||
if (scanner->current_dir == NULL) {
|
||||
if (!open_next_directory(scanner))
|
||||
break;
|
||||
}
|
||||
while ((entry = readdir(dir)) != NULL) {
|
||||
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0) {
|
||||
|
||||
struct dirent *entry = readdir(scanner->current_dir);
|
||||
if (entry == NULL) {
|
||||
closedir(scanner->current_dir);
|
||||
scanner->current_dir = NULL;
|
||||
free(scanner->current_path);
|
||||
scanner->current_path = NULL;
|
||||
continue;
|
||||
}
|
||||
char *cur_path = path_cat(path, entry->d_name);
|
||||
|
||||
if (strcmp(entry->d_name, ".") == 0 || strcmp(entry->d_name, "..") == 0)
|
||||
continue;
|
||||
|
||||
char *cur_path = path_cat(scanner->current_path, entry->d_name);
|
||||
struct stat stats;
|
||||
stat(cur_path, &stats);
|
||||
if (!S_ISREG(stats.st_mode))
|
||||
if (stat(cur_path, &stats) != 0) {
|
||||
free(cur_path);
|
||||
continue;
|
||||
}
|
||||
|
||||
if (!S_ISREG(stats.st_mode)) {
|
||||
queue_enqueue(scanner->directories, (void *)cur_path);
|
||||
else {
|
||||
} else {
|
||||
File *file = file_create(cur_path, &stats);
|
||||
array_list_add(chunk_data, file);
|
||||
chunk_data_size += file->stats.st_size;
|
||||
@@ -62,9 +99,7 @@ Chunk *directory_scanner_next(DirectoryScanner *scanner) {
|
||||
free(cur_path);
|
||||
}
|
||||
}
|
||||
closedir(dir);
|
||||
free(path);
|
||||
}
|
||||
|
||||
if (chunk_data->size > 0)
|
||||
return chunk_data_to_chunk(chunk_data);
|
||||
return NULL;
|
||||
|
||||
@@ -3,8 +3,12 @@
|
||||
|
||||
#include "chunk.h"
|
||||
#include "queue.h"
|
||||
#include <dirent.h>
|
||||
|
||||
typedef struct {
|
||||
Queue *directories;
|
||||
DIR *current_dir;
|
||||
char *current_path;
|
||||
} DirectoryScanner;
|
||||
|
||||
DirectoryScanner *directory_scanner_create(char *root_directory);
|
||||
|
||||
+52
-5
@@ -1,3 +1,4 @@
|
||||
#include "chunk.h"
|
||||
#include "config.h"
|
||||
#include "data.h"
|
||||
#include "file.h"
|
||||
@@ -16,13 +17,36 @@ FileReceive *receive_file_receive(Config *config, int file_descriptor) {
|
||||
Data *file_data = receive_data(file_descriptor);
|
||||
if (config->use_compression) {
|
||||
Data *file_data_uncompressed = data_decompress(file_data);
|
||||
free(file_data);
|
||||
data_destroy(file_data);
|
||||
file_data = file_data_uncompressed;
|
||||
}
|
||||
FileReceive *file = file_receive_create(path, file_data);
|
||||
return file;
|
||||
}
|
||||
|
||||
static void receive_chunk_enqueue(int file_descriptor, Config *config,
|
||||
PipelineContextReceiver *context) {
|
||||
Data *chunk_data = receive_data(file_descriptor);
|
||||
Data *data_to_process = chunk_data;
|
||||
if (config->use_compression) {
|
||||
data_to_process = data_decompress(chunk_data);
|
||||
data_destroy(chunk_data);
|
||||
}
|
||||
Chunk *chunk = chunk_deserialize(data_to_process);
|
||||
data_destroy(data_to_process);
|
||||
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
FileReceive *file = file_receive_create(chunk->items[i]->path,
|
||||
chunk->items[i]->data);
|
||||
chunk->items[i]->path = NULL;
|
||||
chunk->items[i]->data = NULL;
|
||||
queue_enqueue_multithreaded(context->queue, file, &context->mutex,
|
||||
&context->condition_not_empty,
|
||||
&context->condition_not_full);
|
||||
}
|
||||
chunk_destroy(chunk);
|
||||
}
|
||||
|
||||
int receive_thread(void *pipeline_context) {
|
||||
PipelineContextReceiver *context =
|
||||
(PipelineContextReceiver *)pipeline_context;
|
||||
@@ -31,12 +55,18 @@ int receive_thread(void *pipeline_context) {
|
||||
Config *config = context->config;
|
||||
mtx_unlock(&context->mutex);
|
||||
|
||||
while (receive_status(file_descriptor) == STATUS_NEXT) {
|
||||
Status status = receive_status(file_descriptor);
|
||||
while (status == STATUS_NEXT || status == STATUS_CHUNK) {
|
||||
if (status == STATUS_CHUNK) {
|
||||
receive_chunk_enqueue(file_descriptor, config, context);
|
||||
} else {
|
||||
FileReceive *file = receive_file_receive(config, file_descriptor);
|
||||
queue_enqueue_multithreaded(context->queue, file, &context->mutex,
|
||||
&context->condition_not_empty,
|
||||
&context->condition_not_full);
|
||||
}
|
||||
status = receive_status(file_descriptor);
|
||||
}
|
||||
mtx_lock(&context->mutex);
|
||||
context->receiver_done = true;
|
||||
cnd_signal(&context->condition_not_empty);
|
||||
@@ -68,17 +98,34 @@ int write_thread(void *pipeline_context) {
|
||||
|
||||
int receive_files(Config *config, int file_descriptor) {
|
||||
Status status = receive_status(file_descriptor);
|
||||
while (status == STATUS_NEXT) {
|
||||
while (status == STATUS_NEXT || status == STATUS_CHUNK) {
|
||||
if (status == STATUS_CHUNK) {
|
||||
Data *chunk_data = receive_data(file_descriptor);
|
||||
Data *data_to_process = chunk_data;
|
||||
if (config->use_compression) {
|
||||
data_to_process = data_decompress(chunk_data);
|
||||
data_destroy(chunk_data);
|
||||
}
|
||||
Chunk *chunk = chunk_deserialize(data_to_process);
|
||||
data_destroy(data_to_process);
|
||||
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
if (config->save_to_disk)
|
||||
to_disk(path_cat(config->receive_root_directory, chunk->items[i]->path),
|
||||
chunk->items[i]->data->data, chunk->items[i]->data->size);
|
||||
}
|
||||
chunk_destroy(chunk);
|
||||
} else {
|
||||
FileReceive *file = receive_file_receive(config, file_descriptor);
|
||||
if (config->save_to_disk)
|
||||
to_disk(path_cat(config->receive_root_directory, file->path),
|
||||
file->data->data, file->data->size);
|
||||
file_receive_destroy(file);
|
||||
// send_status(file_descriptor, STATUS_OK);
|
||||
}
|
||||
status = receive_status(file_descriptor);
|
||||
}
|
||||
if (status != STATUS_FINISHED) {
|
||||
log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED or NEXT Status");
|
||||
log_message(LOG_LEVEL_ERROR, "Did not receive FINISHED Status");
|
||||
send_status(file_descriptor, STATUS_ERROR);
|
||||
return -1;
|
||||
}
|
||||
|
||||
+50
-46
@@ -4,7 +4,6 @@
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <zstd.h>
|
||||
|
||||
#include "chunk.h"
|
||||
#include "array_list.h"
|
||||
@@ -99,60 +98,46 @@ Data *chunk_format(Chunk *chunk) {
|
||||
return chunk_data_create(data, buffer_size);
|
||||
}
|
||||
|
||||
Data *chunk_compress(Chunk *chunk, int compression_level) {
|
||||
log_message(LOG_LEVEL_DEBUG, "Starting to gather data for chunk compression");
|
||||
Data *chunk_serialize(Chunk *chunk) {
|
||||
unsigned long long data_size = 0;
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
data_size += sizeof(unsigned long long);
|
||||
data_size += sizeof(size_t);
|
||||
data_size += strlen(chunk->items[i]->path);
|
||||
data_size += sizeof(unsigned long long);
|
||||
data_size += sizeof(size_t);
|
||||
data_size += chunk->items[i]->stats.st_size;
|
||||
}
|
||||
Data *data = data_create_empty(data_size);
|
||||
if (data == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR,
|
||||
"Could not allocate memory for chunk compression");
|
||||
"Could not allocate memory for chunk serialization");
|
||||
exit(EXIT_FAILURE);
|
||||
}
|
||||
char *data_pointer = data->data;
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
// path length
|
||||
size_t path_len = strlen(chunk->items[i]->path);
|
||||
memcpy(data_pointer, &path_len, sizeof(size_t));
|
||||
data_pointer += sizeof(size_t);
|
||||
memcpy(data_pointer, chunk->items[i]->path, path_len);
|
||||
data_pointer += path_len;
|
||||
// file data
|
||||
unsigned long long data_size = chunk->items[i]->stats.st_size;
|
||||
memcpy(data_pointer, &data_size, sizeof(size_t));
|
||||
|
||||
size_t file_data_size = chunk->items[i]->stats.st_size;
|
||||
memcpy(data_pointer, &file_data_size, sizeof(size_t));
|
||||
data_pointer += sizeof(size_t);
|
||||
memcpy(data_pointer, chunk->items[i]->data->data, data_size);
|
||||
data_pointer += data_size;
|
||||
}
|
||||
|
||||
log_message(LOG_LEVEL_DEBUG, "Chunk succesfully compressed");
|
||||
Data *compressed = data_compress(data, compression_level);
|
||||
data_destroy(data);
|
||||
return compressed;
|
||||
}
|
||||
|
||||
Chunk *chunk_decompress(Data *compressed_data) {
|
||||
log_message(LOG_LEVEL_DEBUG, "Starting to decompress chunk");
|
||||
Data *uncompressed_data = data_decompress(compressed_data);
|
||||
if (uncompressed_data == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to decompress chunk data");
|
||||
return NULL;
|
||||
memcpy(data_pointer, chunk->items[i]->data->data, file_data_size);
|
||||
data_pointer += file_data_size;
|
||||
}
|
||||
return data;
|
||||
}
|
||||
|
||||
Chunk *chunk_deserialize(Data *data) {
|
||||
ArrayList *files = array_list_create(file_destroy);
|
||||
char *data_pointer = uncompressed_data->data;
|
||||
size_t remaining_size = uncompressed_data->size;
|
||||
char *data_pointer = data->data;
|
||||
size_t remaining_size = data->size;
|
||||
|
||||
while (remaining_size > 0) {
|
||||
if (remaining_size < sizeof(size_t)) {
|
||||
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for path length");
|
||||
array_list_delete(files);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
@@ -163,7 +148,6 @@ Chunk *chunk_decompress(Data *compressed_data) {
|
||||
if (remaining_size < path_len) {
|
||||
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for path");
|
||||
array_list_delete(files);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
@@ -171,7 +155,6 @@ Chunk *chunk_decompress(Data *compressed_data) {
|
||||
if (path == NULL) {
|
||||
perror("Could not allocate memory for file path");
|
||||
array_list_delete(files);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
}
|
||||
memcpy(path, data_pointer, path_len);
|
||||
@@ -183,59 +166,80 @@ Chunk *chunk_decompress(Data *compressed_data) {
|
||||
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for data size");
|
||||
free(path);
|
||||
array_list_delete(files);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
size_t data_size = *(size_t *)data_pointer;
|
||||
size_t file_data_size = *(size_t *)data_pointer;
|
||||
data_pointer += sizeof(size_t);
|
||||
remaining_size -= sizeof(size_t);
|
||||
|
||||
if (remaining_size < data_size) {
|
||||
if (remaining_size < file_data_size) {
|
||||
log_message(LOG_LEVEL_ERROR, "Invalid chunk format: not enough data for file content");
|
||||
free(path);
|
||||
array_list_delete(files);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
struct stat st = {0};
|
||||
st.st_size = data_size;
|
||||
st.st_size = file_data_size;
|
||||
File *file = file_create(path, &st);
|
||||
if (file == NULL) {
|
||||
free(path);
|
||||
array_list_delete(files);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
void *file_data = malloc(data_size);
|
||||
void *file_data = malloc(file_data_size);
|
||||
if (file_data == NULL) {
|
||||
perror("Could not allocate memory for file data");
|
||||
free(path);
|
||||
array_list_delete(files);
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
}
|
||||
memcpy(file_data, data_pointer, data_size);
|
||||
file->data = data_create(file_data, data_size);
|
||||
data_pointer += data_size;
|
||||
remaining_size -= data_size;
|
||||
memcpy(file_data, data_pointer, file_data_size);
|
||||
file->data = data_create(file_data, file_data_size);
|
||||
data_pointer += file_data_size;
|
||||
remaining_size -= file_data_size;
|
||||
|
||||
array_list_add(files, file);
|
||||
free(path);
|
||||
}
|
||||
|
||||
// Create the chunk from the files
|
||||
File **file_array = (File **)array_list_to_array(files);
|
||||
Chunk *chunk = chunk_create(file_array, files->size);
|
||||
|
||||
// Clean up - files are now owned by the chunk
|
||||
free(file_array);
|
||||
files->item_destroyer = NULL;
|
||||
array_list_delete(files);
|
||||
data_destroy(uncompressed_data);
|
||||
|
||||
return chunk;
|
||||
}
|
||||
|
||||
Data *chunk_compress(Chunk *chunk, int compression_level) {
|
||||
log_message(LOG_LEVEL_DEBUG, "Starting to compress chunk");
|
||||
Data *serialized = chunk_serialize(chunk);
|
||||
Data *compressed = data_compress(serialized, compression_level);
|
||||
data_destroy(serialized);
|
||||
log_message(LOG_LEVEL_DEBUG, "Chunk successfully compressed");
|
||||
return compressed;
|
||||
}
|
||||
|
||||
Chunk *chunk_decompress(Data *compressed_data) {
|
||||
log_message(LOG_LEVEL_DEBUG, "Starting to decompress chunk");
|
||||
Data *uncompressed_data = data_decompress(compressed_data);
|
||||
if (uncompressed_data == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to decompress chunk data");
|
||||
return NULL;
|
||||
}
|
||||
|
||||
Chunk *chunk = chunk_deserialize(uncompressed_data);
|
||||
if (chunk == NULL) {
|
||||
log_message(LOG_LEVEL_ERROR, "Failed to deserialize chunk data");
|
||||
data_destroy(uncompressed_data);
|
||||
return NULL;
|
||||
}
|
||||
|
||||
data_destroy(uncompressed_data);
|
||||
log_message(LOG_LEVEL_DEBUG, "Chunk successfully decompressed");
|
||||
return chunk;
|
||||
}
|
||||
|
||||
+2
-3
@@ -6,8 +6,6 @@
|
||||
#include <sys/stat.h>
|
||||
|
||||
#define DESIRED_CHUNK_SIZE 10 * 1024 * 1024
|
||||
#define FILE_PATH_SEPERATOR "#&&SEPP&&#"
|
||||
#define FILE_PATH_DATA_SEPERATOR "#&&SEPD&&#"
|
||||
|
||||
typedef struct {
|
||||
File **items;
|
||||
@@ -18,10 +16,11 @@ 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_serialize(Chunk *chunk);
|
||||
Chunk *chunk_deserialize(Data *data);
|
||||
Data *chunk_compress(Chunk *chunk, int compression_level);
|
||||
Chunk *chunk_decompress(Data *compressed_data);
|
||||
|
||||
Data *chunk_data_create(void *data, unsigned long long data_size);
|
||||
void chunk_data_delete(void *chunk);
|
||||
void chunk_data_to_disk(Data *chunk, char *root_directory);
|
||||
#endif
|
||||
|
||||
+1
-1
@@ -100,4 +100,4 @@ void file_receive_destroy(void *file_receive) {
|
||||
free(file);
|
||||
}
|
||||
|
||||
FileReceive *file_receive_from_buffer(void *buffer) {}
|
||||
|
||||
|
||||
@@ -21,11 +21,8 @@ void file_load_data(File *file);
|
||||
void file_print(void *item);
|
||||
void file_send_single_calls(File *file, int file_descriptor);
|
||||
size_t file_content_to_buffer(File *file);
|
||||
Data *file_compress(File *file);
|
||||
|
||||
FileReceive *file_receive_create(char *path, Data *data);
|
||||
void file_receive_destroy(void *file_receive);
|
||||
FileReceive *file_receive_from_buffer(void *buffer);
|
||||
FileReceive *file_receive_decompress(void *FileReceive);
|
||||
|
||||
#endif
|
||||
|
||||
@@ -195,6 +195,8 @@ const char *status_to_string(Status status) {
|
||||
return "FINISHED";
|
||||
case STATUS_NEXT:
|
||||
return "NEXT";
|
||||
case STATUS_CHUNK:
|
||||
return "CHUNK";
|
||||
default:
|
||||
return "UNKNOWN";
|
||||
}
|
||||
|
||||
+1
-1
@@ -5,7 +5,7 @@
|
||||
#include <netinet/in.h>
|
||||
|
||||
typedef int Status;
|
||||
enum NET_STATUS { STATUS_OK, STATUS_ERROR, STATUS_FINISHED, STATUS_NEXT };
|
||||
enum NET_STATUS { STATUS_OK, STATUS_ERROR, STATUS_FINISHED, STATUS_NEXT, STATUS_CHUNK };
|
||||
|
||||
typedef struct Server {
|
||||
struct sockaddr_in address;
|
||||
|
||||
@@ -1,27 +1,31 @@
|
||||
import argparse
|
||||
import filecmp
|
||||
import os
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
import tempfile
|
||||
import time
|
||||
|
||||
# --- Configuration ---
|
||||
TEST_DIR = os.path.join(os.path.dirname(os.path.abspath(__file__)), "test_data")
|
||||
DEFAULT_SOURCE_DIR = os.path.join(TEST_DIR, "source")
|
||||
DEFAULT_DEST_DIR = os.path.join(TEST_DIR, "dest")
|
||||
|
||||
SERVER_CMD = ["./build/server"]
|
||||
base_client_cmd = ["./build/client"]
|
||||
|
||||
# --- Resource Limit Configuration ---
|
||||
# 💾 Disk throttling settings
|
||||
DISK_DEVICE = (
|
||||
"/dev/nvme0n1p5" # IMPORTANT: Change this to your disk (e.g., /dev/nvme0n1)
|
||||
)
|
||||
READ_BPS_MAX = "15M" # Max read speed (M for megabytes)
|
||||
WRITE_BPS_MAX = "10M" # Max write speed
|
||||
DISK_DEVICE = "/dev/nvme0n1p5"
|
||||
READ_BPS_MAX = "15M"
|
||||
WRITE_BPS_MAX = "10M"
|
||||
|
||||
# 🐢 Network throttling settings (Linux tc)
|
||||
NET_LIMIT = "100mbit"
|
||||
NET_DELAY = "100ms"
|
||||
NETWORK_INTERFACE = "lo"
|
||||
NET_LIMIT_CMD = f"sudo tc qdisc add dev {NETWORK_INTERFACE} root netem rate {NET_LIMIT} delay {NET_DELAY}".split()
|
||||
NET_LIMIT_CMD = (
|
||||
f"sudo tc qdisc add dev {NETWORK_INTERFACE} root netem rate {NET_LIMIT} delay {NET_DELAY}".split()
|
||||
)
|
||||
NET_RESET_CMD = f"sudo tc qdisc del dev {NETWORK_INTERFACE} root".split()
|
||||
|
||||
# --- Build the client command prefix with throttling ---
|
||||
CLIENT_CMD_PREFIX = [
|
||||
"sudo",
|
||||
"systemd-run",
|
||||
@@ -32,89 +36,161 @@ CLIENT_CMD_PREFIX = [
|
||||
f"IOWriteBandwidthMax={DISK_DEVICE} {WRITE_BPS_MAX}",
|
||||
]
|
||||
|
||||
# --- Test Cases ---
|
||||
TEST_CASES = [
|
||||
{"name": "Standard (Single-threaded)", "flags": []},
|
||||
{"name": "Multithreading (-m)", "flags": ["-m"]},
|
||||
# {"name": "Compression (-c -5)", "flags": ["-c -5"]},
|
||||
{"name": "Compression (-c 0)", "flags": ["-c 0"]},
|
||||
# {"name": "Compression (-c 10)", "flags": ["-c 10"]},
|
||||
# {"name": "Compression (-c 20)", "flags": ["-c 20"]},
|
||||
# {"name": "Chunk Serialization (-s)", "flags": ["-s"]},
|
||||
{"name": "Compression (-c 0)", "flags": ["-c", "0"]},
|
||||
{"name": "Chunk Serialization (-s)", "flags": ["-s"]},
|
||||
{"name": "Compression + Chunk Serialization (-c -s)", "flags": ["-c", "-s"]},
|
||||
{"name": "Multithreading + Compression (-m -c)", "flags": ["-m", "-c"]},
|
||||
# {"name": "Multithreading + Chunk Serialization (-m -s)", "flags": ["-m", "-s"]},
|
||||
# {"name": "Compression + Chunk Serialization (-c -s)", "flags": ["-c", "-s"]},
|
||||
# {
|
||||
# "name": "Multithreading + Compression + Chunk Serialization (-m -c -s)",
|
||||
# "flags": ["-m", "-c", "-s"],
|
||||
# },
|
||||
{"name": "Multithreading + Chunk Serialization (-m -s)", "flags": ["-m", "-s"]},
|
||||
{
|
||||
"name": "Multithreading + Compression + Chunk Serialization (-m -c -s)",
|
||||
"flags": ["-m", "-c", "-s"],
|
||||
},
|
||||
]
|
||||
|
||||
|
||||
def run_suite(env_name, apply_limits):
|
||||
def generate_test_files(source_dir):
|
||||
if os.path.exists(source_dir):
|
||||
shutil.rmtree(source_dir)
|
||||
os.makedirs(source_dir)
|
||||
|
||||
target_total = 50 * 1024 * 1024
|
||||
written = 0
|
||||
|
||||
files = {
|
||||
"small.txt": b"hello world\n",
|
||||
"medium.txt": b"the quick brown fox jumps over the lazy dog\n" * 5000,
|
||||
"binary.bin": bytes(range(256)) * 1000,
|
||||
"nested/subdir/deep.txt": b"deeply nested file\n",
|
||||
"nested/another.txt": b"another nested file\n" * 50,
|
||||
}
|
||||
for rel_path, content in files.items():
|
||||
full_path = os.path.join(source_dir, rel_path)
|
||||
os.makedirs(os.path.dirname(full_path), exist_ok=True)
|
||||
with open(full_path, "wb") as f:
|
||||
f.write(content)
|
||||
written += len(content)
|
||||
|
||||
i = 0
|
||||
while written < target_total:
|
||||
chunk_size = min(5 * 1024 * 1024, target_total - written)
|
||||
rel_path = f"bulk/file_{i}.dat"
|
||||
full_path = os.path.join(source_dir, rel_path)
|
||||
os.makedirs(os.path.dirname(full_path), exist_ok=True)
|
||||
with open(full_path, "wb") as f:
|
||||
f.write(b"0" * chunk_size)
|
||||
written += chunk_size
|
||||
i += 1
|
||||
|
||||
total_mb = written / (1024 * 1024)
|
||||
print(f" Generated {total_mb:.1f} MB of test data in {source_dir}")
|
||||
|
||||
|
||||
def verify_transfer(source_dir, dest_dir):
|
||||
source_dir = os.path.abspath(source_dir)
|
||||
dest_dir = os.path.abspath(dest_dir)
|
||||
|
||||
received_prefix = os.path.join(dest_dir, source_dir.lstrip(os.sep))
|
||||
if not os.path.exists(received_prefix):
|
||||
return [], ["no received files found"]
|
||||
|
||||
mismatches = []
|
||||
missing = []
|
||||
|
||||
for root, dirs, files in os.walk(source_dir):
|
||||
for f in files:
|
||||
src_path = os.path.join(root, f)
|
||||
rel = os.path.relpath(src_path, source_dir)
|
||||
dst_path = os.path.join(received_prefix, rel)
|
||||
|
||||
if not os.path.exists(dst_path):
|
||||
missing.append(rel)
|
||||
elif not filecmp.cmp(src_path, dst_path, shallow=False):
|
||||
mismatches.append(rel)
|
||||
|
||||
return mismatches, missing
|
||||
|
||||
|
||||
def run_suite(env_name, apply_limits, source_dir, dest_dir):
|
||||
results = []
|
||||
print(f"\n{'=' * 60}")
|
||||
print(f"🚀 Starting Suite: {env_name}")
|
||||
print(f"Suite: {env_name}")
|
||||
print(f"{'=' * 60}")
|
||||
|
||||
if apply_limits:
|
||||
print(
|
||||
f"Applying Disk I/O Limits: Reads <= {READ_BPS_MAX}, Writes <= {WRITE_BPS_MAX}"
|
||||
)
|
||||
print(f"Applying Network Limits: {NET_LIMIT}, {NET_DELAY} delay")
|
||||
print(f" Disk I/O: Reads <= {READ_BPS_MAX}, Writes <= {WRITE_BPS_MAX}")
|
||||
print(f" Network: {NET_LIMIT}, {NET_DELAY} delay")
|
||||
client_prefix = CLIENT_CMD_PREFIX
|
||||
else:
|
||||
print("Running Baseline (No limits applied)")
|
||||
client_prefix = [] # Run normally without systemd-run/limits
|
||||
print(" Baseline (no limits)")
|
||||
client_prefix = []
|
||||
|
||||
try:
|
||||
# SETUP: Apply or ensure clean network limits
|
||||
if apply_limits:
|
||||
subprocess.run(NET_LIMIT_CMD, check=True)
|
||||
else:
|
||||
# Silently attempt to clear any leftover rules just to ensure a clean baseline
|
||||
subprocess.run(NET_RESET_CMD, capture_output=True)
|
||||
|
||||
for case in TEST_CASES:
|
||||
name = case["name"]
|
||||
flags = case["flags"]
|
||||
print(f"\n --- {name} ---")
|
||||
|
||||
print(f"\n--- Running: {name} ---")
|
||||
if os.path.exists(dest_dir):
|
||||
shutil.rmtree(dest_dir)
|
||||
|
||||
server_process = None
|
||||
try:
|
||||
# 1. Start the server
|
||||
print(" Starting server...")
|
||||
server_process = subprocess.Popen(
|
||||
SERVER_CMD, stdout=subprocess.DEVNULL, stderr=None
|
||||
)
|
||||
time.sleep(0.5) # Allow server to bind to port
|
||||
time.sleep(0.5)
|
||||
|
||||
# 2. Build and run the client
|
||||
client_cmd = client_prefix + base_client_cmd + flags
|
||||
print(f" Running client: {' '.join(client_cmd)}")
|
||||
env = os.environ.copy()
|
||||
|
||||
client_cmd = (
|
||||
client_prefix
|
||||
+ base_client_cmd
|
||||
+ ["--source-dir", source_dir, "--dest-dir", dest_dir, "--save-to-disk"]
|
||||
+ flags
|
||||
)
|
||||
print(f" Running: {' '.join(client_cmd)}")
|
||||
|
||||
start_time = time.monotonic()
|
||||
client_result = subprocess.run(
|
||||
client_cmd, text=True, capture_output=True
|
||||
client_cmd, env=env, text=True, capture_output=True
|
||||
)
|
||||
end_time = time.monotonic()
|
||||
|
||||
duration = end_time - start_time
|
||||
|
||||
if server_process:
|
||||
try:
|
||||
server_process.wait(timeout=5)
|
||||
except subprocess.TimeoutExpired:
|
||||
server_process.kill()
|
||||
server_process.wait()
|
||||
server_process = None
|
||||
|
||||
mismatches, missing = [], []
|
||||
if client_result.returncode == 0:
|
||||
results.append(
|
||||
{
|
||||
"environment": env_name,
|
||||
mismatches, missing = verify_transfer(source_dir, dest_dir)
|
||||
|
||||
entry = {
|
||||
"name": name,
|
||||
"status": "Success",
|
||||
"time": f"{duration:.4f}s",
|
||||
"error": "",
|
||||
"suite": env_name,
|
||||
"time": f"{duration:.4f}s" if client_result.returncode == 0 else "N/A",
|
||||
}
|
||||
)
|
||||
|
||||
if client_result.returncode == 0 and not mismatches and not missing:
|
||||
entry["status"] = "Success"
|
||||
entry["error"] = ""
|
||||
else:
|
||||
print(f" ⚠️ Failed (code: {client_result.returncode})")
|
||||
err_msg = (
|
||||
entry["status"] = "Failed"
|
||||
errors = []
|
||||
if client_result.returncode != 0:
|
||||
err = (
|
||||
client_result.stderr.strip().split("\n")[0]
|
||||
if client_result.stderr
|
||||
else (
|
||||
@@ -123,107 +199,101 @@ def run_suite(env_name, apply_limits):
|
||||
else "No output"
|
||||
)
|
||||
)
|
||||
results.append(
|
||||
{
|
||||
"environment": env_name,
|
||||
"name": name,
|
||||
"status": "Failed",
|
||||
"time": "N/A",
|
||||
"error": f"Exit code {client_result.returncode}: {err_msg[:40]}",
|
||||
}
|
||||
)
|
||||
errors.append(f"Exit code {client_result.returncode}: {err[:80]}")
|
||||
if missing:
|
||||
errors.append(f"Missing ({len(missing)}): {', '.join(missing[:5])}")
|
||||
if mismatches:
|
||||
errors.append(f"Mismatch ({len(mismatches)}): {', '.join(mismatches[:3])}")
|
||||
entry["error"] = " | ".join(errors)
|
||||
|
||||
results.append(entry)
|
||||
|
||||
except subprocess.TimeoutExpired:
|
||||
print(" ⚠️ Timeout (exceeded 15s)")
|
||||
results.append(
|
||||
{
|
||||
"environment": env_name,
|
||||
"name": name,
|
||||
"status": "Timeout",
|
||||
"time": "N/A",
|
||||
"error": "Exceeded 15 seconds",
|
||||
}
|
||||
{"name": name, "suite": env_name, "status": "Timeout", "time": "N/A", "error": "Exceeded 15s"}
|
||||
)
|
||||
except Exception as e:
|
||||
print(f" ❌ Error: {e}")
|
||||
results.append(
|
||||
{
|
||||
"environment": env_name,
|
||||
"name": name,
|
||||
"status": "Error",
|
||||
"time": "N/A",
|
||||
"error": str(e),
|
||||
}
|
||||
{"name": name, "suite": env_name, "status": "Error", "time": "N/A", "error": str(e)}
|
||||
)
|
||||
finally:
|
||||
# Clean up the server for this test case
|
||||
if server_process:
|
||||
print(" Stopping server...")
|
||||
try:
|
||||
server_process.terminate()
|
||||
server_process.wait(timeout=5)
|
||||
except subprocess.TimeoutExpired:
|
||||
server_process.kill()
|
||||
server_process.wait()
|
||||
|
||||
except subprocess.CalledProcessError as e:
|
||||
print(f"❌ Error running system limit command: {' '.join(e.cmd)}")
|
||||
print("Are you running this script with 'sudo' privileges?")
|
||||
|
||||
print(f" Error running limit command: {' '.join(e.cmd)}")
|
||||
finally:
|
||||
# TEARDOWN: Remove network limits if they were applied
|
||||
if apply_limits:
|
||||
print("\nCleaning up limits for this suite...")
|
||||
try:
|
||||
subprocess.run(NET_RESET_CMD, check=True, capture_output=True)
|
||||
print("Network limits removed.")
|
||||
except Exception as e:
|
||||
print(f"⚠️ Could not reset network settings: {e}")
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
return results
|
||||
|
||||
|
||||
os.system("cmake -B build -S .")
|
||||
os.system("cd build && make")
|
||||
def main():
|
||||
parser = argparse.ArgumentParser(description="FastSync integration test / benchmark")
|
||||
parser.add_argument("--source-dir", default=DEFAULT_SOURCE_DIR,
|
||||
help="Source directory for test files (default: %(default)s)")
|
||||
parser.add_argument("--dest-dir", default=DEFAULT_DEST_DIR,
|
||||
help="Destination directory for received files (default: %(default)s)")
|
||||
parser.add_argument("--keep-data", action="store_true",
|
||||
help="Keep test_data directory after run")
|
||||
parser.add_argument("--no-throttled", action="store_true",
|
||||
help="Skip throttled suite (requires sudo)")
|
||||
parser.add_argument("--no-unlimited", action="store_true",
|
||||
help="Skip unlimited suite")
|
||||
args = parser.parse_args()
|
||||
|
||||
# --- Main Execution ---
|
||||
os.system("cmake -B build -S . > /dev/null 2>&1")
|
||||
ret = os.system("cd build && make -j$(nproc) 2>&1 | tail -3")
|
||||
if ret != 0:
|
||||
print("Build failed")
|
||||
sys.exit(1)
|
||||
|
||||
generate_test_files(args.source_dir)
|
||||
os.makedirs(args.dest_dir, exist_ok=True)
|
||||
|
||||
try:
|
||||
all_results = []
|
||||
|
||||
# 1. Run Baseline (No Limits)
|
||||
all_results.extend(run_suite("Unlimited", apply_limits=False))
|
||||
|
||||
# # 2. Run Throttled (With Limits)
|
||||
all_results.extend(run_suite("Throttled", apply_limits=True))
|
||||
|
||||
# --- Print Comparison Table ---
|
||||
print("\n" + "=" * 105)
|
||||
print(f"{'fastSync BENCHMARK RESULTS (COMPARISON)':^105}")
|
||||
print("=" * 105)
|
||||
print(
|
||||
f"{'Configuration':<45} | {'Environment':<12} | {'Status':<10} | {'Time':<10} | {'Details/Error':<20}"
|
||||
if not args.no_unlimited:
|
||||
all_results.extend(
|
||||
run_suite("Unlimited", False, args.source_dir, args.dest_dir)
|
||||
)
|
||||
print("-" * 105)
|
||||
|
||||
# Sort results by test case name first, then environment to easily compare
|
||||
# This groups the baseline and throttled results for the same test next to each other
|
||||
# sorted_results = sorted(
|
||||
# all_results,
|
||||
# key=lambda x: (
|
||||
# TEST_CASES.index(
|
||||
# next(item for item in TEST_CASES if item["name"] == x["name"])
|
||||
# ),
|
||||
# x["environment"],
|
||||
# ),
|
||||
# )
|
||||
if not args.no_throttled:
|
||||
all_results.extend(
|
||||
run_suite("Throttled", True, args.source_dir, args.dest_dir)
|
||||
)
|
||||
|
||||
print("\n" + "=" * 110)
|
||||
print(f"{'RESULTS':^110}")
|
||||
print("=" * 110)
|
||||
print(f"{'Configuration':<45} | {'Suite':<12} | {'Status':<8} | {'Time':<10} | {'Details'}")
|
||||
print("-" * 110)
|
||||
|
||||
for res in all_results:
|
||||
status_symbol = (
|
||||
"✅"
|
||||
if res["status"] == "Success"
|
||||
else ("⏳" if res["status"] == "Timeout" else "❌")
|
||||
)
|
||||
status_str = f"{status_symbol} {res['status']}"
|
||||
print(
|
||||
f"{res['name']:<45} | {res['environment']:<12} | {status_str:<10} | {res['time']:<10} | {res['error']:<20}"
|
||||
f"{res['name']:<45} | {res['suite']:<12} | {res['status']:<8} | {res['time']:<10} | {res['error']}"
|
||||
)
|
||||
print("=" * 105)
|
||||
|
||||
failed = [r for r in all_results if r["status"] != "Success"]
|
||||
if failed:
|
||||
print(f"\n {len(failed)} test(s) FAILED")
|
||||
sys.exit(1)
|
||||
else:
|
||||
print(f"\n ALL {len(all_results)} TESTS PASSED")
|
||||
|
||||
finally:
|
||||
if not args.keep_data:
|
||||
shutil.rmtree(TEST_DIR, ignore_errors=True)
|
||||
|
||||
|
||||
if __name__ == "__main__":
|
||||
main()
|
||||
|
||||
Reference in New Issue
Block a user