From 78767219062e0ee86b2d0574938755368b61d91c Mon Sep 17 00:00:00 2001 From: TapTap Date: Sat, 18 Jul 2026 17:00:11 +0200 Subject: [PATCH] refactor: extract receive_chunk_data and receive_manifest shared helpers - receive_chunk_data(fd, config) -> Chunk*: shared receive/decompress/deserialize (eliminates ~25 lines of duplication between server.c and multiprocessing.c) - receive_manifest(fd, config, next_status): moved from server.c to file.c, fixes missing 'Deleting files not in manifest...' log in multiprocessing.c - Removes redundant manual free loop in multiprocessing.c manifest handling (array_list_create(free) destructor already handles this) - server.c and multiprocessing.c: remove local chunk/manifest helpers, remove compression.h include (no longer needed) --- src/server/server.c | 56 +++++------------------------------- src/shared/chunk.c | 23 +++++++++++++++ src/shared/chunk.h | 2 ++ src/shared/file.c | 20 +++++++++++-- src/shared/file.h | 1 + src/shared/multiprocessing.c | 40 ++------------------------ 6 files changed, 54 insertions(+), 88 deletions(-) diff --git a/src/server/server.c b/src/server/server.c index 1cfbedc..a0cf8d2 100644 --- a/src/server/server.c +++ b/src/server/server.c @@ -1,6 +1,5 @@ #include "array_list.h" #include "chunk.h" -#include "compression.h" #include "config.h" #include "data.h" #include "file.h" @@ -16,53 +15,6 @@ #include #include -static int receive_chunk(int fd, Config *config) { - Data *chunk_data = receive_data(fd); - if (chunk_data == NULL) { - log_message(LOG_LEVEL_ERROR, "Failed to receive chunk data"); - return -1; - } - Data *data_to_process = chunk_data; - if (config->use_compression) { - data_to_process = data_decompress(chunk_data); - data_destroy(chunk_data); - if (data_to_process == NULL) { - log_message(LOG_LEVEL_ERROR, "Failed to decompress chunk"); - return -1; - } - } - Chunk *chunk = chunk_deserialize(data_to_process, config->use_metadata); - data_destroy(data_to_process); - if (chunk == NULL) { - log_message(LOG_LEVEL_ERROR, "Failed to deserialize chunk, skipping"); - return -1; - } - - for (int i = 0; i < chunk->element_count; i++) { - if (config->save_to_disk) - file_save_to_disk(config->receive_root_directory, chunk->items[i]); - } - chunk_destroy(chunk); - return 0; -} - -static int receive_manifest(int fd, Config *config, Status *next_status) { - int count; - if (!receive_int(fd, &count)) return -1; - ArrayList *manifest = array_list_create(free); - if (manifest) { - for (int i = 0; i < count; i++) { - char *s = receive_str(fd); - if (s) array_list_add(manifest, s); - } - fprintf(stderr, "Deleting files not in manifest...\n"); - delete_extras(config->receive_root_directory, manifest); - array_list_delete(manifest); - } - if (!receive_status(fd, next_status)) return -1; - return 0; -} - int receive_files(Config *config, int fd) { Status status; if (!receive_status(fd, &status)) return -1; @@ -77,10 +29,16 @@ int receive_files(Config *config, int fd) { file_save_to_disk(config->receive_root_directory, file); file_destroy(file); } else if (status == STATUS_CHUNK) { - if (receive_chunk(fd, config) != 0) { + Chunk *chunk = receive_chunk_data(fd, config); + if (chunk == NULL) { send_status(fd, STATUS_ERROR); return -1; } + for (int i = 0; i < chunk->element_count; i++) { + if (config->save_to_disk) + file_save_to_disk(config->receive_root_directory, chunk->items[i]); + } + chunk_destroy(chunk); } else { File *file = file_receive(config, fd); if (file == NULL) { diff --git a/src/shared/chunk.c b/src/shared/chunk.c index ebc6434..18b7edc 100644 --- a/src/shared/chunk.c +++ b/src/shared/chunk.c @@ -10,6 +10,7 @@ #include "file.h" #include "log.h" #include "metadata.h" +#include "protocol.h" Chunk *chunk_create(File **items, int element_count) { Chunk *chunk = (Chunk *)malloc(sizeof(Chunk)); @@ -178,5 +179,27 @@ Data *chunk_compress(Chunk *chunk, int compression_level, bool use_metadata) { return compressed; } +Chunk *receive_chunk_data(int fd, Config *config) { + Data *chunk_data = receive_data(fd); + if (chunk_data == NULL) { + log_message(LOG_LEVEL_ERROR, "Failed to receive chunk data"); + return NULL; + } + Data *data_to_process = chunk_data; + if (config->use_compression) { + data_to_process = data_decompress(chunk_data); + data_destroy(chunk_data); + if (data_to_process == NULL) { + log_message(LOG_LEVEL_ERROR, "Failed to decompress chunk"); + return NULL; + } + } + Chunk *chunk = chunk_deserialize(data_to_process, config->use_metadata); + data_destroy(data_to_process); + if (chunk == NULL) + log_message(LOG_LEVEL_ERROR, "Failed to deserialize chunk, skipping"); + return chunk; +} + diff --git a/src/shared/chunk.h b/src/shared/chunk.h index 7837936..fc07792 100644 --- a/src/shared/chunk.h +++ b/src/shared/chunk.h @@ -1,6 +1,7 @@ #ifndef CHUNK_H #define CHUNK_H +#include "config.h" #include "data.h" #include "file.h" #include @@ -18,5 +19,6 @@ void chunk_destroy(void *chunk); Data *chunk_serialize(Chunk *chunk, bool use_metadata); Chunk *chunk_deserialize(Data *data, bool use_metadata); Data *chunk_compress(Chunk *chunk, int compression_level, bool use_metadata); +Chunk *receive_chunk_data(int fd, Config *config); #endif diff --git a/src/shared/file.c b/src/shared/file.c index d757431..1cfb569 100644 --- a/src/shared/file.c +++ b/src/shared/file.c @@ -184,8 +184,7 @@ bool file_save_to_disk(const char *root_directory, File *file) { return ok; } -File *receive_incremental_check(int fd, Config *config, bool *skipped) { - *skipped = false; +File *receive_incremental_check(int fd, Config *config, bool *skipped) { *skipped = false; char *check_path = receive_str(fd); if (check_path == NULL) { send_status(fd, STATUS_ERROR); return NULL; } @@ -348,4 +347,21 @@ size_t file_content_to_buffer(File *file) { return bytes_read; } +int receive_manifest(int fd, Config *config, int *next_status) { + int count; + if (!receive_int(fd, &count)) return -1; + ArrayList *manifest = array_list_create(free); + if (manifest) { + for (int i = 0; i < count; i++) { + char *s = receive_str(fd); + if (s) array_list_add(manifest, s); + } + fprintf(stderr, "Deleting files not in manifest...\n"); + delete_extras(config->receive_root_directory, manifest); + array_list_delete(manifest); + } + if (!receive_status(fd, next_status)) return -1; + return 0; +} + diff --git a/src/shared/file.h b/src/shared/file.h index d7bf752..7eb4d42 100644 --- a/src/shared/file.h +++ b/src/shared/file.h @@ -34,5 +34,6 @@ void file_metadata_destroy(void *metadata); bool to_disk(const char *path, const void *data, unsigned long long data_size); bool file_save_to_disk(const char *root_directory, File *file); File *receive_incremental_check(int fd, Config *config, bool *skipped); +int receive_manifest(int fd, Config *config, int *next_status); #endif diff --git a/src/shared/multiprocessing.c b/src/shared/multiprocessing.c index 8e11890..ffb6148 100644 --- a/src/shared/multiprocessing.c +++ b/src/shared/multiprocessing.c @@ -1,7 +1,6 @@ #include "multiprocessing.h" #include "array_list.h" #include "chunk.h" -#include "compression.h" #include "config.h" #include "data.h" #include "file.h" @@ -84,26 +83,8 @@ void pipeline_context_receiver_destroy(PipelineContextReceiver *context) { static void receive_chunk_enqueue(int file_descriptor, PipelineContextReceiver *context) { - Data *chunk_data = receive_data(file_descriptor); - if (chunk_data == NULL) { - log_message(LOG_LEVEL_ERROR, "Failed to receive chunk data"); - return; - } - Data *data_to_process = chunk_data; - if (context->config->use_compression) { - data_to_process = data_decompress(chunk_data); - data_destroy(chunk_data); - if (data_to_process == NULL) { - log_message(LOG_LEVEL_ERROR, "Failed to decompress chunk"); - return; - } - } - Chunk *chunk = chunk_deserialize(data_to_process, context->config->use_metadata); - data_destroy(data_to_process); - if (chunk == NULL) { - log_message(LOG_LEVEL_ERROR, "Failed to deserialize chunk, skipping"); - return; - } + Chunk *chunk = receive_chunk_data(file_descriptor, context->config); + if (chunk == NULL) return; for (int i = 0; i < chunk->element_count; i++) { File *file = chunk->items[i]; @@ -150,22 +131,7 @@ int receive_thread(void *pipeline_context) { if (!receive_status(file_descriptor, &status)) return thrd_error; } if (status == STATUS_MANIFEST) { - int count; - if (!receive_int(file_descriptor, &count)) return thrd_error; - ArrayList *manifest = array_list_create(free); - if (manifest) { - for (int i = 0; i < count; i++) { - char *s = receive_str(file_descriptor); - if (s) { - array_list_add(manifest, s); - } - } - delete_extras(context->config->receive_root_directory, manifest); - for (int i = 0; i < manifest->size; i++) - free(manifest->items[i]); - array_list_delete(manifest); - } - if (!receive_status(file_descriptor, &status)) return thrd_error; + if (receive_manifest(file_descriptor, config, &status) != 0) return thrd_error; } mtx_lock(&context->mutex); context->receiver_done = true;