Merge pull request 'fixup: address PR review comments' (#3) from fixup-pr-comments into cleanup
Reviewed-on: #3
This commit was merged in pull request #3.
This commit is contained in:
@@ -1,6 +1,5 @@
|
||||
#include "client_send.h"
|
||||
#include "chunk.h"
|
||||
#include "compression.h"
|
||||
#include "config.h"
|
||||
#include "data.h"
|
||||
#include "file.h"
|
||||
@@ -24,11 +23,9 @@ int send_chunk(Client *client, Chunk *chunk, Config *config) {
|
||||
if (config->use_compression) {
|
||||
data = chunk_compress(chunk, config->compression_level, config->use_metadata);
|
||||
} else {
|
||||
for (int i = 0; i < chunk->element_count; i++)
|
||||
file_load_data(chunk->items[i]);
|
||||
data = chunk_serialize(chunk, config->use_metadata);
|
||||
}
|
||||
send_data(client->file_descriptor, data->data, data->size);
|
||||
send_data(client->file_descriptor, data);
|
||||
data_destroy(data);
|
||||
} else if (config->use_sendfile && !config->use_compression) {
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
@@ -38,14 +35,9 @@ int send_chunk(Client *client, Chunk *chunk, Config *config) {
|
||||
} else {
|
||||
for (int i = 0; i < chunk->element_count; i++) {
|
||||
send_status(client->file_descriptor, STATUS_NEXT);
|
||||
File *file = chunk->items[i];
|
||||
if (config->use_compression) {
|
||||
Data *compressed_data =
|
||||
data_compress(file->data, config->compression_level);
|
||||
data_destroy(file->data);
|
||||
file->data = compressed_data;
|
||||
}
|
||||
file_send_single_calls(file, client->file_descriptor, config->use_metadata);
|
||||
file_send_single_calls(chunk->items[i], client->file_descriptor,
|
||||
config->use_metadata,
|
||||
config->use_compression ? config->compression_level : 0);
|
||||
}
|
||||
}
|
||||
return 0;
|
||||
|
||||
@@ -3,7 +3,6 @@
|
||||
#include "config.h"
|
||||
#include "data.h"
|
||||
#include "file.h"
|
||||
#include "io.h"
|
||||
#include "log.h"
|
||||
#include "metadata.h"
|
||||
#include "multiprocessing.h"
|
||||
|
||||
+26
-3
@@ -9,9 +9,10 @@
|
||||
#include <sys/stat.h>
|
||||
#include <unistd.h>
|
||||
|
||||
#include "compression.h"
|
||||
#include "config.h"
|
||||
#include "data.h"
|
||||
#include "file.h"
|
||||
#include "io.h"
|
||||
#include "log.h"
|
||||
#include "metadata.h"
|
||||
#include "protocol.h"
|
||||
@@ -90,11 +91,16 @@ void file_load_data(File *file) {
|
||||
}
|
||||
}
|
||||
|
||||
void file_send_single_calls(File *file, int file_descriptor, bool use_metadata) {
|
||||
void file_send_single_calls(File *file, int file_descriptor, bool use_metadata, int compression_level) {
|
||||
if (compression_level > 0) {
|
||||
Data *compressed_data = data_compress(file->data, compression_level);
|
||||
data_destroy(file->data);
|
||||
file->data = compressed_data;
|
||||
}
|
||||
send_str(file_descriptor, file->path);
|
||||
if (use_metadata)
|
||||
metadata_send(file_descriptor, file->metadata);
|
||||
send_data(file_descriptor, file->data->data, file->data->size);
|
||||
send_data(file_descriptor, file->data);
|
||||
}
|
||||
|
||||
void to_disk(const char *path, const void *data, unsigned long long data_size) {
|
||||
@@ -139,6 +145,23 @@ void file_send_sendfile(File *file, int file_descriptor, bool use_metadata) {
|
||||
close(fd);
|
||||
}
|
||||
|
||||
File *file_receive(Config *config, int file_descriptor) {
|
||||
char *path = (char *)receive_str(file_descriptor);
|
||||
File *file = file_create(path);
|
||||
free(path);
|
||||
if (config->use_metadata)
|
||||
file->metadata = metadata_receive(file_descriptor);
|
||||
Data *file_data = receive_data(file_descriptor);
|
||||
if (config->use_compression) {
|
||||
Data *file_data_uncompressed = data_decompress(file_data);
|
||||
data_destroy(file_data);
|
||||
file_data = file_data_uncompressed;
|
||||
}
|
||||
data_destroy(file->data);
|
||||
file->data = file_data;
|
||||
return file;
|
||||
}
|
||||
|
||||
size_t file_content_to_buffer(File *file) {
|
||||
FILE *file_pointer = fopen(file->path, "rb");
|
||||
if (file_pointer == NULL) {
|
||||
|
||||
+3
-1
@@ -1,6 +1,7 @@
|
||||
#ifndef FILE_H
|
||||
#define FILE_H
|
||||
|
||||
#include "config.h"
|
||||
#include "data.h"
|
||||
#include <stdbool.h>
|
||||
#include <sys/stat.h>
|
||||
@@ -22,7 +23,8 @@ typedef struct {
|
||||
File *file_create(const char *path);
|
||||
void file_destroy(void *item);
|
||||
void file_load_data(File *file);
|
||||
void file_send_single_calls(File *file, int file_descriptor, bool use_metadata);
|
||||
File *file_receive(Config *config, int file_descriptor);
|
||||
void file_send_single_calls(File *file, int file_descriptor, bool use_metadata, int compression_level);
|
||||
void file_send_sendfile(File *file, int file_descriptor, bool use_metadata);
|
||||
size_t file_content_to_buffer(File *file);
|
||||
FileMetadata *file_metadata_create(struct stat *stats);
|
||||
|
||||
@@ -1,49 +0,0 @@
|
||||
#include "io.h"
|
||||
#include "log.h"
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <unistd.h>
|
||||
|
||||
static __thread int io_read_fd = -1;
|
||||
static __thread int io_write_fd = -1;
|
||||
|
||||
void io_set_fds(int read_fd, int write_fd) {
|
||||
io_read_fd = read_fd;
|
||||
io_write_fd = write_fd;
|
||||
}
|
||||
|
||||
static int io_fd(int dir_fd, int file_descriptor) {
|
||||
return (dir_fd != -1) ? dir_fd : file_descriptor;
|
||||
}
|
||||
|
||||
void send_n_data(int file_descriptor, void *data, size_t data_size) {
|
||||
log_message(LOG_LEVEL_DEBUG, " Sending n Data: %zu", data_size);
|
||||
int fd = io_fd(io_write_fd, file_descriptor);
|
||||
ssize_t total_bytes_send = 0;
|
||||
while (total_bytes_send < data_size) {
|
||||
ssize_t bytes_send =
|
||||
write(fd, (char *)data + total_bytes_send, data_size - total_bytes_send);
|
||||
if (bytes_send <= 0) {
|
||||
perror("Could not send data!");
|
||||
exit(EXIT_FAILURE);
|
||||
}
|
||||
total_bytes_send += bytes_send;
|
||||
}
|
||||
log_message(LOG_LEVEL_DEBUG, " Send n Data: %zu", total_bytes_send);
|
||||
}
|
||||
|
||||
void receive_n_data(int file_descriptor, void *data, size_t data_size) {
|
||||
log_message(LOG_LEVEL_DEBUG, " Receiving n Data: %zu", data_size);
|
||||
int fd = io_fd(io_read_fd, file_descriptor);
|
||||
size_t total_bytes_received = 0;
|
||||
while (total_bytes_received < data_size) {
|
||||
ssize_t bytes_received =
|
||||
read(fd, (char *)data + total_bytes_received, data_size - total_bytes_received);
|
||||
if (bytes_received == -1 || bytes_received == 0) {
|
||||
perror("Could not receive bytes!");
|
||||
exit(EXIT_FAILURE);
|
||||
}
|
||||
total_bytes_received += bytes_received;
|
||||
}
|
||||
log_message(LOG_LEVEL_DEBUG, " Received n Data: %zu", total_bytes_received);
|
||||
}
|
||||
@@ -1,10 +0,0 @@
|
||||
#ifndef IO_H
|
||||
#define IO_H
|
||||
|
||||
#include <stddef.h>
|
||||
|
||||
void io_set_fds(int read_fd, int write_fd);
|
||||
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);
|
||||
|
||||
#endif
|
||||
@@ -1,6 +1,6 @@
|
||||
#include "metadata.h"
|
||||
#include "file.h"
|
||||
#include "io.h"
|
||||
#include "protocol.h"
|
||||
#include <fcntl.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
|
||||
@@ -74,23 +74,6 @@ void pipeline_context_receiver_destroy(PipelineContextReceiver *context) {
|
||||
free(context);
|
||||
}
|
||||
|
||||
File *file_receive(Config *config, int file_descriptor) {
|
||||
char *path = (char *)receive_str(file_descriptor);
|
||||
File *file = file_create(path);
|
||||
free(path);
|
||||
if (config->use_metadata)
|
||||
file->metadata = metadata_receive(file_descriptor);
|
||||
Data *file_data = receive_data(file_descriptor);
|
||||
if (config->use_compression) {
|
||||
Data *file_data_uncompressed = data_decompress(file_data);
|
||||
data_destroy(file_data);
|
||||
file_data = file_data_uncompressed;
|
||||
}
|
||||
data_destroy(file->data);
|
||||
file->data = file_data;
|
||||
return file;
|
||||
}
|
||||
|
||||
static void receive_chunk_enqueue(int file_descriptor,
|
||||
PipelineContextReceiver *context) {
|
||||
Data *chunk_data = receive_data(file_descriptor);
|
||||
|
||||
@@ -39,7 +39,6 @@ PipelineContextReceiver *pipeline_context_receiver_create(Config *config,
|
||||
Queue *queue_receiver,
|
||||
int file_descriptor);
|
||||
void pipeline_context_receiver_destroy(PipelineContextReceiver *context);
|
||||
File *file_receive(Config *config, int file_descriptor);
|
||||
int receive_thread(void *pipeline_context);
|
||||
int write_thread(void *pipeline_context);
|
||||
#endif
|
||||
|
||||
+48
-3
@@ -1,9 +1,53 @@
|
||||
#include "protocol.h"
|
||||
#include "io.h"
|
||||
#include "log.h"
|
||||
#include <stdio.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <unistd.h>
|
||||
|
||||
static __thread int io_read_fd = -1;
|
||||
static __thread int io_write_fd = -1;
|
||||
|
||||
void io_set_fds(int read_fd, int write_fd) {
|
||||
io_read_fd = read_fd;
|
||||
io_write_fd = write_fd;
|
||||
}
|
||||
|
||||
static int io_fd(int dir_fd, int file_descriptor) {
|
||||
return (dir_fd != -1) ? dir_fd : file_descriptor;
|
||||
}
|
||||
|
||||
void send_n_data(int file_descriptor, void *data, size_t data_size) {
|
||||
log_message(LOG_LEVEL_DEBUG, " Sending n Data: %zu", data_size);
|
||||
int fd = io_fd(io_write_fd, file_descriptor);
|
||||
ssize_t total_bytes_send = 0;
|
||||
while (total_bytes_send < data_size) {
|
||||
ssize_t bytes_send =
|
||||
write(fd, (char *)data + total_bytes_send, data_size - total_bytes_send);
|
||||
if (bytes_send <= 0) {
|
||||
perror("Could not send data!");
|
||||
exit(EXIT_FAILURE);
|
||||
}
|
||||
total_bytes_send += bytes_send;
|
||||
}
|
||||
log_message(LOG_LEVEL_DEBUG, " Send n Data: %zu", total_bytes_send);
|
||||
}
|
||||
|
||||
void receive_n_data(int file_descriptor, void *data, size_t data_size) {
|
||||
log_message(LOG_LEVEL_DEBUG, " Receiving n Data: %zu", data_size);
|
||||
int fd = io_fd(io_read_fd, file_descriptor);
|
||||
size_t total_bytes_received = 0;
|
||||
while (total_bytes_received < data_size) {
|
||||
ssize_t bytes_received =
|
||||
read(fd, (char *)data + total_bytes_received, data_size - total_bytes_received);
|
||||
if (bytes_received == -1 || bytes_received == 0) {
|
||||
perror("Could not receive bytes!");
|
||||
exit(EXIT_FAILURE);
|
||||
}
|
||||
total_bytes_received += bytes_received;
|
||||
}
|
||||
log_message(LOG_LEVEL_DEBUG, " Received n Data: %zu", total_bytes_received);
|
||||
}
|
||||
|
||||
static const char *status_to_string(Status status) {
|
||||
switch (status) {
|
||||
@@ -39,9 +83,10 @@ char *receive_str(int file_descriptor) {
|
||||
return data;
|
||||
}
|
||||
|
||||
void send_data(int file_descriptor, void *data, unsigned long long data_size) {
|
||||
void send_data(int file_descriptor, Data *data) {
|
||||
unsigned long long data_size = data->size;
|
||||
send_n_data(file_descriptor, &data_size, sizeof(unsigned long long));
|
||||
send_n_data(file_descriptor, data, data_size);
|
||||
send_n_data(file_descriptor, data->data, data_size);
|
||||
log_message(LOG_LEVEL_DEBUG, "Send %lld data", data_size);
|
||||
}
|
||||
|
||||
|
||||
@@ -2,15 +2,18 @@
|
||||
#define PROTOCOL_H
|
||||
|
||||
#include "data.h"
|
||||
#include "io.h"
|
||||
#include <stddef.h>
|
||||
|
||||
typedef int Status;
|
||||
enum NET_STATUS { STATUS_OK, STATUS_ERROR, STATUS_FINISHED, STATUS_NEXT, STATUS_CHUNK };
|
||||
|
||||
void io_set_fds(int read_fd, int write_fd);
|
||||
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);
|
||||
void send_data(int file_descriptor, Data *data);
|
||||
Data *receive_data(int file_descriptor);
|
||||
void send_int(int file_descriptor, int data);
|
||||
int receive_int(int file_descriptor);
|
||||
|
||||
Reference in New Issue
Block a user