Merge pull request 'feat: add -f/--sendfile for zero-copy file transfer' (#8) from sendfile-support into main
Reviewed-on: #8
This commit is contained in:
@@ -12,6 +12,7 @@ FastFileTransfer is a C implementation of a file synchronization system that:
|
|||||||
4. Utilizes multithreading for parallel file processing
|
4. Utilizes multithreading for parallel file processing
|
||||||
5. Implements producer-consumer patterns with thread-safe queues
|
5. Implements producer-consumer patterns with thread-safe queues
|
||||||
6. Provides both in-memory and disk-based storage options
|
6. Provides both in-memory and disk-based storage options
|
||||||
|
7. Supports `sendfile()` for zero-copy file transfer
|
||||||
|
|
||||||
## System Architecture
|
## System Architecture
|
||||||
|
|
||||||
@@ -23,6 +24,7 @@ The system consists of two main components:
|
|||||||
- Compresses data using zstd algorithm
|
- Compresses data using zstd algorithm
|
||||||
- Serializes chunks into a compact binary format for batch transfer
|
- Serializes chunks into a compact binary format for batch transfer
|
||||||
- Sends files to server using custom protocol
|
- Sends files to server using custom protocol
|
||||||
|
- Supports sendfile for zero-copy file transfer (`-f`)
|
||||||
- Supports both single-threaded and multi-threaded operation
|
- Supports both single-threaded and multi-threaded operation
|
||||||
|
|
||||||
### Server
|
### Server
|
||||||
@@ -49,6 +51,7 @@ The client-server communication uses the following status codes:
|
|||||||
| `-m` | Enable multithreading mode |
|
| `-m` | Enable multithreading mode |
|
||||||
| `-c [level]` | Enable compression with optional level (1-22, default: 5) |
|
| `-c [level]` | Enable compression with optional level (1-22, default: 5) |
|
||||||
| `-s` | Enable chunk serialization (batch-transfer all files per chunk) |
|
| `-s` | Enable chunk serialization (batch-transfer all files per chunk) |
|
||||||
|
| `-f` | Enable sendfile (zero-copy file transfer, bypasses userspace memory). Can be combined with `-m`. Incompatible with `-c` and `-s`. |
|
||||||
| `--source-dir <path>` | Source directory to sync (overrides `FASTSYNC_SOURCE_DIR`) |
|
| `--source-dir <path>` | Source directory to sync (overrides `FASTSYNC_SOURCE_DIR`) |
|
||||||
| `--dest-dir <path>` | Server-side destination directory (overrides `FASTSYNC_DEST_DIR`) |
|
| `--dest-dir <path>` | Server-side destination directory (overrides `FASTSYNC_DEST_DIR`) |
|
||||||
| `--save-to-disk` | Persist received files to disk |
|
| `--save-to-disk` | Persist received files to disk |
|
||||||
@@ -118,6 +121,12 @@ make
|
|||||||
|
|
||||||
# Multithreaded with compressed chunk serialization
|
# Multithreaded with compressed chunk serialization
|
||||||
./build/client -m -s -c 3
|
./build/client -m -s -c 3
|
||||||
|
|
||||||
|
# Sendfile (zero-copy, bypasses userspace for large files)
|
||||||
|
./build/client -f
|
||||||
|
|
||||||
|
# Sendfile with multithreading
|
||||||
|
./build/client -f -m
|
||||||
```
|
```
|
||||||
|
|
||||||
## Testing
|
## Testing
|
||||||
@@ -150,8 +159,9 @@ tests/ # Unit tests
|
|||||||
|
|
||||||
1. Chunk size (10MB default) affects memory usage and transfer efficiency
|
1. Chunk size (10MB default) affects memory usage and transfer efficiency
|
||||||
2. Compression level (1-22) trades CPU usage for space savings
|
2. Compression level (1-22) trades CPU usage for space savings
|
||||||
3. Multithreading improves performance on multi-core systems
|
3. `sendfile()` (`-f`) bypasses userspace memory, ~2x faster on localhost for large files
|
||||||
4. Thread-safe queues minimize contention between producer/consumer threads
|
4. Multithreading improves performance on multi-core systems
|
||||||
|
5. Thread-safe queues minimize contention between producer/consumer threads
|
||||||
|
|
||||||
## Extensibility
|
## Extensibility
|
||||||
|
|
||||||
|
|||||||
+22
-5
@@ -29,6 +29,11 @@ int send_chunk(Client *client, Chunk *chunk, Config *config) {
|
|||||||
}
|
}
|
||||||
send_data(client->file_descriptor, data->data, data->size);
|
send_data(client->file_descriptor, data->data, data->size);
|
||||||
data_destroy(data);
|
data_destroy(data);
|
||||||
|
} else if (config->use_sendfile && !config->use_compression) {
|
||||||
|
for (int i = 0; i < chunk->element_count; i++) {
|
||||||
|
send_status(client->file_descriptor, STATUS_NEXT);
|
||||||
|
file_send_sendfile(chunk->items[i], client->file_descriptor);
|
||||||
|
}
|
||||||
} else {
|
} else {
|
||||||
for (int i = 0; i < chunk->element_count; i++) {
|
for (int i = 0; i < chunk->element_count; i++) {
|
||||||
send_status(client->file_descriptor, STATUS_NEXT);
|
send_status(client->file_descriptor, STATUS_NEXT);
|
||||||
@@ -81,8 +86,10 @@ int load_files_multithreaded(void *pipeline_context) {
|
|||||||
mtx_unlock(&context->mutex_loader);
|
mtx_unlock(&context->mutex_loader);
|
||||||
return thrd_success;
|
return thrd_success;
|
||||||
}
|
}
|
||||||
for (int i = 0; i < chunk->element_count; i++)
|
if (!context->config->use_sendfile) {
|
||||||
file_load_data(chunk->items[i]);
|
for (int i = 0; i < chunk->element_count; i++)
|
||||||
|
file_load_data(chunk->items[i]);
|
||||||
|
}
|
||||||
queue_enqueue_multithreaded(context->queue_loader, chunk,
|
queue_enqueue_multithreaded(context->queue_loader, chunk,
|
||||||
&context->mutex_loader,
|
&context->mutex_loader,
|
||||||
&context->condition_not_empty_loader,
|
&context->condition_not_empty_loader,
|
||||||
@@ -130,8 +137,10 @@ int send_files(Config *config) {
|
|||||||
DirectoryScanner *scanner = directory_scanner_create(config->send_directory);
|
DirectoryScanner *scanner = directory_scanner_create(config->send_directory);
|
||||||
Chunk *current_chunk;
|
Chunk *current_chunk;
|
||||||
while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
|
while ((current_chunk = directory_scanner_next(scanner)) != NULL) {
|
||||||
for (int i = 0; i < current_chunk->element_count; i++)
|
if (!config->use_sendfile) {
|
||||||
file_load_data(current_chunk->items[i]);
|
for (int i = 0; i < current_chunk->element_count; i++)
|
||||||
|
file_load_data(current_chunk->items[i]);
|
||||||
|
}
|
||||||
send_chunk(client, current_chunk, config);
|
send_chunk(client, current_chunk, config);
|
||||||
chunk_destroy(current_chunk);
|
chunk_destroy(current_chunk);
|
||||||
}
|
}
|
||||||
@@ -192,7 +201,7 @@ int main(int argc, char *argv[]) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
Config *config = config_create(str_dup("1.0.0"), source_dir, dest_dir,
|
Config *config = config_create(str_dup("1.0.0"), source_dir, dest_dir,
|
||||||
save_to_disk, false, false, false, 5, 20);
|
save_to_disk, false, false, false, 5, 20, false);
|
||||||
for (int i = 1; i < argc; i++) {
|
for (int i = 1; i < argc; i++) {
|
||||||
if (strcmp(argv[i], "-c") == 0) {
|
if (strcmp(argv[i], "-c") == 0) {
|
||||||
config->use_compression = true;
|
config->use_compression = true;
|
||||||
@@ -215,6 +224,9 @@ int main(int argc, char *argv[]) {
|
|||||||
config->receive_root_directory = str_dup(argv[++i]);
|
config->receive_root_directory = str_dup(argv[++i]);
|
||||||
} else if (strcmp(argv[i], "--save-to-disk") == 0) {
|
} else if (strcmp(argv[i], "--save-to-disk") == 0) {
|
||||||
config->save_to_disk = true;
|
config->save_to_disk = true;
|
||||||
|
} else if (strcmp(argv[i], "-f") == 0 || strcmp(argv[i], "--sendfile") == 0) {
|
||||||
|
config->use_sendfile = true;
|
||||||
|
log_message(LOG_LEVEL_INFO, "Enabled sendfile");
|
||||||
} else {
|
} else {
|
||||||
handle_arg(argv[i], "-m", &config->use_multithreading,
|
handle_arg(argv[i], "-m", &config->use_multithreading,
|
||||||
"Enabled Multithreading");
|
"Enabled Multithreading");
|
||||||
@@ -223,6 +235,11 @@ int main(int argc, char *argv[]) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
if (config->use_sendfile && (config->use_chunk_serialization || config->use_compression)) {
|
||||||
|
fprintf(stderr, "Error: -f/--sendfile cannot be combined with -c (compression) or -s (chunk serialization)\n");
|
||||||
|
return 1;
|
||||||
|
}
|
||||||
|
|
||||||
if (config->use_multithreading)
|
if (config->use_multithreading)
|
||||||
return send_files_multithreaded(config);
|
return send_files_multithreaded(config);
|
||||||
return send_files(config);
|
return send_files(config);
|
||||||
|
|||||||
+2
-1
@@ -8,7 +8,7 @@ Config *config_create(char *version, char *send_directory,
|
|||||||
char *receive_directory, bool save_to_disk,
|
char *receive_directory, bool save_to_disk,
|
||||||
bool use_multithreading, bool use_chunk_serialization,
|
bool use_multithreading, bool use_chunk_serialization,
|
||||||
bool use_compression, int compression_level,
|
bool use_compression, int compression_level,
|
||||||
int num_connections) {
|
int num_connections, bool use_sendfile) {
|
||||||
|
|
||||||
Config *config = malloc(sizeof(Config));
|
Config *config = malloc(sizeof(Config));
|
||||||
config->version = version;
|
config->version = version;
|
||||||
@@ -20,6 +20,7 @@ Config *config_create(char *version, char *send_directory,
|
|||||||
config->use_compression = use_compression;
|
config->use_compression = use_compression;
|
||||||
config->compression_level = compression_level;
|
config->compression_level = compression_level;
|
||||||
config->num_connections = num_connections;
|
config->num_connections = num_connections;
|
||||||
|
config->use_sendfile = use_sendfile;
|
||||||
return config;
|
return config;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
+2
-1
@@ -11,6 +11,7 @@ typedef struct Config {
|
|||||||
bool use_multithreading;
|
bool use_multithreading;
|
||||||
bool use_chunk_serialization;
|
bool use_chunk_serialization;
|
||||||
bool use_compression;
|
bool use_compression;
|
||||||
|
bool use_sendfile;
|
||||||
bool use_single_send_per_file;
|
bool use_single_send_per_file;
|
||||||
int compression_level;
|
int compression_level;
|
||||||
int num_connections;
|
int num_connections;
|
||||||
@@ -20,7 +21,7 @@ Config *config_create(char *version, char *send_directory,
|
|||||||
char *receive_directory, bool save_to_disk,
|
char *receive_directory, bool save_to_disk,
|
||||||
bool use_multithreading, bool use_chunk_serialization,
|
bool use_multithreading, bool use_chunk_serialization,
|
||||||
bool use_compression, int compression_level,
|
bool use_compression, int compression_level,
|
||||||
int num_connections);
|
int num_connections, bool use_sendfile);
|
||||||
void config_delete(Config *config);
|
void config_delete(Config *config);
|
||||||
void config_send(int file_descriptor, Config *config);
|
void config_send(int file_descriptor, Config *config);
|
||||||
Config *config_receive(int file_descriptor);
|
Config *config_receive(int file_descriptor);
|
||||||
|
|||||||
@@ -1,9 +1,12 @@
|
|||||||
#include <dirent.h>
|
#include <dirent.h>
|
||||||
|
#include <fcntl.h>
|
||||||
#include <libgen.h>
|
#include <libgen.h>
|
||||||
#include <stddef.h>
|
#include <stddef.h>
|
||||||
#include <stdio.h>
|
#include <stdio.h>
|
||||||
#include <stdlib.h>
|
#include <stdlib.h>
|
||||||
#include <string.h>
|
#include <string.h>
|
||||||
|
#include <sys/sendfile.h>
|
||||||
|
#include <unistd.h>
|
||||||
#include <zstd.h>
|
#include <zstd.h>
|
||||||
|
|
||||||
#include "data.h"
|
#include "data.h"
|
||||||
@@ -68,6 +71,31 @@ void file_send_single_calls(File *file, int file_descriptor) {
|
|||||||
send_data(file_descriptor, file->data->data, file->data->size);
|
send_data(file_descriptor, file->data->data, file->data->size);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
void file_send_sendfile(File *file, int file_descriptor) {
|
||||||
|
send_str(file_descriptor, file->path);
|
||||||
|
|
||||||
|
int fd = open(file->path, O_RDONLY);
|
||||||
|
if (fd == -1) {
|
||||||
|
perror("Could not open file for sendfile");
|
||||||
|
exit(EXIT_FAILURE);
|
||||||
|
}
|
||||||
|
|
||||||
|
unsigned long long file_size = file->stats.st_size;
|
||||||
|
send_n_data(file_descriptor, &file_size, sizeof(unsigned long long));
|
||||||
|
|
||||||
|
off_t offset = 0;
|
||||||
|
while (offset < file_size) {
|
||||||
|
ssize_t sent = sendfile(file_descriptor, fd, &offset, file_size - offset);
|
||||||
|
if (sent == -1) {
|
||||||
|
perror("sendfile failed");
|
||||||
|
close(fd);
|
||||||
|
exit(EXIT_FAILURE);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
close(fd);
|
||||||
|
}
|
||||||
|
|
||||||
size_t file_content_to_buffer(File *file) {
|
size_t file_content_to_buffer(File *file) {
|
||||||
FILE *file_pointer = fopen(file->path, "rb");
|
FILE *file_pointer = fopen(file->path, "rb");
|
||||||
if (file_pointer == NULL) {
|
if (file_pointer == NULL) {
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ void file_destroy(void *item);
|
|||||||
void file_load_data(File *file);
|
void file_load_data(File *file);
|
||||||
void file_print(void *item);
|
void file_print(void *item);
|
||||||
void file_send_single_calls(File *file, int file_descriptor);
|
void file_send_single_calls(File *file, int file_descriptor);
|
||||||
|
void file_send_sendfile(File *file, int file_descriptor);
|
||||||
size_t file_content_to_buffer(File *file);
|
size_t file_content_to_buffer(File *file);
|
||||||
|
|
||||||
FileReceive *file_receive_create(char *path, Data *data);
|
FileReceive *file_receive_create(char *path, Data *data);
|
||||||
|
|||||||
@@ -48,6 +48,8 @@ TEST_CASES = [
|
|||||||
"name": "Multithreading + Compression + Chunk Serialization (-m -c -s)",
|
"name": "Multithreading + Compression + Chunk Serialization (-m -c -s)",
|
||||||
"flags": ["-m", "-c", "-s"],
|
"flags": ["-m", "-c", "-s"],
|
||||||
},
|
},
|
||||||
|
{"name": "Sendfile (-f)", "flags": ["-f"]},
|
||||||
|
{"name": "Sendfile + Multithreading (-f -m)", "flags": ["-f", "-m"]},
|
||||||
]
|
]
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
+3
-3
@@ -8,7 +8,7 @@
|
|||||||
|
|
||||||
static void test_config_lifecycle() {
|
static void test_config_lifecycle() {
|
||||||
Config *cfg = config_create(str_dup("1.0"), str_dup("/src"), str_dup("/dst"),
|
Config *cfg = config_create(str_dup("1.0"), str_dup("/src"), str_dup("/dst"),
|
||||||
true, true, false, false, 1, 4);
|
true, true, false, false, 1, 4, false);
|
||||||
EXPECT_NOT_NULL(cfg);
|
EXPECT_NOT_NULL(cfg);
|
||||||
EXPECT_EQ_STR(cfg->version, "1.0");
|
EXPECT_EQ_STR(cfg->version, "1.0");
|
||||||
EXPECT_EQ_STR(cfg->send_directory, "/src");
|
EXPECT_EQ_STR(cfg->send_directory, "/src");
|
||||||
@@ -23,7 +23,7 @@ static void test_config_lifecycle() {
|
|||||||
|
|
||||||
static void test_pipeline_sender_lifecycle() {
|
static void test_pipeline_sender_lifecycle() {
|
||||||
Config *cfg = config_create(str_dup("2.0"), str_dup("/src2"),
|
Config *cfg = config_create(str_dup("2.0"), str_dup("/src2"),
|
||||||
str_dup("/dst2"), false, false, true, true, 1, 8);
|
str_dup("/dst2"), false, false, true, true, 1, 8, false);
|
||||||
Queue *q1 = queue_create(5, NULL);
|
Queue *q1 = queue_create(5, NULL);
|
||||||
Queue *q2 = queue_create(15, NULL);
|
Queue *q2 = queue_create(15, NULL);
|
||||||
|
|
||||||
@@ -40,7 +40,7 @@ static void test_pipeline_sender_lifecycle() {
|
|||||||
|
|
||||||
static void test_pipeline_receiver_lifecycle() {
|
static void test_pipeline_receiver_lifecycle() {
|
||||||
Config *cfg = config_create(str_dup("3.0"), str_dup("/src3"),
|
Config *cfg = config_create(str_dup("3.0"), str_dup("/src3"),
|
||||||
str_dup("/dst3"), true, true, true, true, 1, 2);
|
str_dup("/dst3"), true, true, true, true, 1, 2, false);
|
||||||
Queue *q = queue_create(20, NULL);
|
Queue *q = queue_create(20, NULL);
|
||||||
|
|
||||||
PipelineContextReceiver *pcr = pipeline_context_receiver_create(cfg, q, 42);
|
PipelineContextReceiver *pcr = pipeline_context_receiver_create(cfg, q, 42);
|
||||||
|
|||||||
Reference in New Issue
Block a user