Merge fix/security-hardening into dev
fixes #254 (receiver queue byte-budget) #258 (inplace setuid/truncate) # Conflicts: # src/shared/multiprocessing.c
This commit is contained in:
@@ -4,6 +4,7 @@
|
||||
#include "log.h"
|
||||
#include "metadata.h"
|
||||
#include "multiprocessing.h"
|
||||
#include "protocol.h"
|
||||
#include "queue.h"
|
||||
#include "receiver.h"
|
||||
#include "transport_tcp.h"
|
||||
@@ -26,6 +27,13 @@ static bool allow_delete;
|
||||
static bool allow_unauthenticated;
|
||||
static const char* required_client_cn;
|
||||
|
||||
/* Aggregate payload bytes the multithreaded receiver may buffer ahead of the
|
||||
slow disk writer. Receiving one more chunk adds up to ~2 * MAX_CHUNK_SIZE
|
||||
of transient wire/decompression buffers on top of the queued payloads, so
|
||||
this ceiling keeps total per-connection receive memory (decompressed and
|
||||
per-file copied chunk buffers included) within MAX_CONNECTION_MEMORY. */
|
||||
#define RECEIVER_QUEUE_MAX_BYTES (MAX_CONNECTION_MEMORY - 2 * MAX_CHUNK_SIZE)
|
||||
|
||||
static bool tls_client_identity_allowed(SSL* ssl) {
|
||||
if (!ssl || !required_client_cn)
|
||||
return false;
|
||||
@@ -314,6 +322,7 @@ void handler(int file_descriptor) {
|
||||
protocol_session_set_max_alloc(&context->session, config->max_alloc);
|
||||
atomic_store(&context->session.total_allocated_bytes,
|
||||
atomic_load(&session.total_allocated_bytes));
|
||||
pipeline_context_receiver_set_queue_byte_limit(context, RECEIVER_QUEUE_MAX_BYTES);
|
||||
thrd_t receiver, writer;
|
||||
bool receiver_created = thrd_create(&receiver, receive_thread, context) == thrd_success;
|
||||
bool writer_created = false;
|
||||
|
||||
+17
-2
@@ -348,10 +348,25 @@ static bool file_to_disk_secure_impl(const char* path, const void* data,
|
||||
if (newer) {
|
||||
ok = true;
|
||||
} else {
|
||||
if (!sparse || data_size == 0 || ftruncate(fd, (off_t)data_size) == 0)
|
||||
/* In-place overwrites: pre-size sparse targets and always trim the
|
||||
file to the new payload length afterwards so shorter payloads can
|
||||
never leave stale trailing bytes from a previous version. */
|
||||
if (sparse && data_size > 0)
|
||||
ok = ftruncate(fd, (off_t)data_size) == 0;
|
||||
if (ok || !sparse || data_size == 0)
|
||||
ok = write_all(fd, data, data_size);
|
||||
if (ok && metadata)
|
||||
if (ok)
|
||||
ok = ftruncate(fd, (off_t)data_size) == 0;
|
||||
/* Normalize the mode: apply the metadata-derived safe mode when the
|
||||
sender supplied metadata (setuid/setgid/sticky are never honored);
|
||||
otherwise fall back to a safe default so dangerous bits on an
|
||||
existing destination cannot survive an overwrite. */
|
||||
if (ok) {
|
||||
if (metadata)
|
||||
ok = file_restore_metadata_fd(fd, metadata, preserve_executability);
|
||||
else if (fchmod(fd, S_IRUSR | S_IWUSR | S_IRGRP | S_IROTH) != 0)
|
||||
ok = false;
|
||||
}
|
||||
if (ok && use_fsync)
|
||||
ok = fsync(fd) == 0;
|
||||
}
|
||||
|
||||
@@ -111,6 +111,8 @@ PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue*
|
||||
protocol_session_init(&context->session, file_descriptor, file_descriptor);
|
||||
protocol_session_set_ssl(&context->session, ssl);
|
||||
context->receiver_done = false;
|
||||
context->queued_bytes = 0;
|
||||
context->max_queue_bytes = 0;
|
||||
atomic_init(&context->cancelled, false);
|
||||
int init = 0;
|
||||
if (mtx_init(&context->mutex, mtx_plain) != thrd_success)
|
||||
@@ -147,14 +149,73 @@ void pipeline_context_receiver_destroy(PipelineContextReceiver* context) {
|
||||
free(context);
|
||||
}
|
||||
|
||||
static bool receiver_enqueue_file(File* file, void* context_pointer) {
|
||||
PipelineContextReceiver* context = context_pointer;
|
||||
if (queue_enqueue_multithreaded_cancel(context->queue, file, &context->mutex,
|
||||
&context->condition_not_empty,
|
||||
&context->condition_not_full, &context->cancelled))
|
||||
return true;
|
||||
void pipeline_context_receiver_set_queue_byte_limit(PipelineContextReceiver* context,
|
||||
size_t max_bytes) {
|
||||
if (context == NULL)
|
||||
return;
|
||||
mtx_lock(&context->mutex);
|
||||
context->max_queue_bytes = max_bytes;
|
||||
context->queued_bytes = 0;
|
||||
cnd_broadcast(&context->condition_not_full);
|
||||
mtx_unlock(&context->mutex);
|
||||
}
|
||||
|
||||
void pipeline_context_receiver_note_bytes_released(PipelineContextReceiver* context,
|
||||
size_t released_bytes) {
|
||||
if (context == NULL || context->max_queue_bytes == 0 || released_bytes == 0)
|
||||
return;
|
||||
mtx_lock(&context->mutex);
|
||||
if (released_bytes >= context->queued_bytes)
|
||||
context->queued_bytes = 0;
|
||||
else
|
||||
context->queued_bytes -= released_bytes;
|
||||
cnd_signal(&context->condition_not_full);
|
||||
mtx_unlock(&context->mutex);
|
||||
}
|
||||
|
||||
bool pipeline_context_receiver_enqueue_file(PipelineContextReceiver* context, File* file) {
|
||||
if (context == NULL || file == NULL)
|
||||
return false;
|
||||
size_t file_bytes = file->data ? file->data->size : 0;
|
||||
mtx_lock(&context->mutex);
|
||||
while (!atomic_load(&context->cancelled)) {
|
||||
bool blocked_by_count = queue_is_full(context->queue);
|
||||
bool blocked_by_budget = false;
|
||||
if (context->max_queue_bytes > 0) {
|
||||
size_t budget = context->max_queue_bytes;
|
||||
size_t used = context->queued_bytes;
|
||||
if (used >= budget) {
|
||||
blocked_by_budget = true;
|
||||
} else if (file_bytes > budget - used) {
|
||||
/* A single payload larger than the whole budget (not possible with
|
||||
the per-file receive cap) is only admitted to an empty pipeline so
|
||||
the wait can never deadlock. */
|
||||
blocked_by_budget = used != 0;
|
||||
}
|
||||
}
|
||||
if (!blocked_by_count && !blocked_by_budget)
|
||||
break;
|
||||
cnd_wait(&context->condition_not_full, &context->mutex);
|
||||
}
|
||||
if (atomic_load(&context->cancelled)) {
|
||||
mtx_unlock(&context->mutex);
|
||||
file_destroy(file);
|
||||
return false;
|
||||
}
|
||||
if (!queue_enqueue(context->queue, file)) {
|
||||
mtx_unlock(&context->mutex);
|
||||
file_destroy(file);
|
||||
return false;
|
||||
}
|
||||
context->queued_bytes += file_bytes;
|
||||
cnd_signal(&context->condition_not_empty);
|
||||
mtx_unlock(&context->mutex);
|
||||
return true;
|
||||
}
|
||||
|
||||
static bool receiver_enqueue_file(File* file, void* context_pointer) {
|
||||
PipelineContextReceiver* context = (PipelineContextReceiver*)context_pointer;
|
||||
return pipeline_context_receiver_enqueue_file(context, file);
|
||||
}
|
||||
|
||||
static void receiver_thread_fail(PipelineContextReceiver* context) {
|
||||
@@ -215,11 +276,13 @@ int write_thread(void* pipeline_context) {
|
||||
protocol_session_unbind();
|
||||
return thrd_success;
|
||||
}
|
||||
size_t file_bytes = file->data ? file->data->size : 0;
|
||||
FileSaveResult result = FILE_SAVE_SKIPPED;
|
||||
if (save_to_disk) {
|
||||
result = file_save_to_disk_full(root_directory, file, context->config);
|
||||
if (result == FILE_SAVE_ERROR) {
|
||||
file_destroy(file);
|
||||
pipeline_context_receiver_note_bytes_released(context, file_bytes);
|
||||
mtx_lock(&context->mutex);
|
||||
atomic_store(&context->cancelled, true);
|
||||
context->receiver_done = true;
|
||||
@@ -236,6 +299,7 @@ int write_thread(void* pipeline_context) {
|
||||
if (context->config->remove_source_files &&
|
||||
!receiver_outcomes_append(&context->outcomes, (unsigned char)result)) {
|
||||
file_destroy(file);
|
||||
pipeline_context_receiver_note_bytes_released(context, file_bytes);
|
||||
mtx_lock(&context->mutex);
|
||||
atomic_store(&context->cancelled, true);
|
||||
context->receiver_done = true;
|
||||
@@ -247,5 +311,6 @@ int write_thread(void* pipeline_context) {
|
||||
return thrd_error;
|
||||
}
|
||||
file_destroy(file);
|
||||
pipeline_context_receiver_note_bytes_released(context, file_bytes);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -47,6 +47,14 @@ typedef struct PipelineContextReceiver {
|
||||
cnd_t condition_not_empty;
|
||||
bool receiver_done;
|
||||
atomic_bool cancelled;
|
||||
/* Aggregate payload bytes that have been received but not yet released by
|
||||
the disk writer (queued or in the writer's hand). Guarded by `mutex`.
|
||||
When `max_queue_bytes` is non-zero the receiver blocks before enqueuing
|
||||
once this total would exceed it, so decompressed/copied file payloads
|
||||
buffered ahead of a slow disk writer respect the per-connection memory
|
||||
budget instead of growing without bound. */
|
||||
size_t queued_bytes;
|
||||
size_t max_queue_bytes;
|
||||
} PipelineContextReceiver;
|
||||
|
||||
PipelineContextSender* pipeline_context_sender_create(Config* config, Queue* queue_scanner,
|
||||
@@ -55,6 +63,18 @@ void pipeline_context_sender_destroy(PipelineContextSender* context);
|
||||
PipelineContextReceiver* pipeline_context_receiver_create(Config* config, Queue* queue_receiver,
|
||||
int file_descriptor, SSL* ssl);
|
||||
void pipeline_context_receiver_destroy(PipelineContextReceiver* context);
|
||||
/* Bound the bytes buffered ahead of the disk writer (see max_queue_bytes). */
|
||||
void pipeline_context_receiver_set_queue_byte_limit(PipelineContextReceiver* context,
|
||||
size_t max_bytes);
|
||||
/* Blocking enqueue used by the receive pipeline sink. Blocks while the queue
|
||||
is full by element count or when adding `file` would push queued_bytes over
|
||||
the configured byte limit; waits until the disk writer releases bytes.
|
||||
Takes ownership of `file` on success and destroys it on failure/cancel. */
|
||||
bool pipeline_context_receiver_enqueue_file(PipelineContextReceiver* context, File* file);
|
||||
/* Account for `released_bytes` of payload memory that has been freed by the
|
||||
disk writer, unblocking a receiver that is waiting on the byte limit. */
|
||||
void pipeline_context_receiver_note_bytes_released(PipelineContextReceiver* context,
|
||||
size_t released_bytes);
|
||||
int receive_thread(void* pipeline_context);
|
||||
int write_thread(void* pipeline_context);
|
||||
#endif
|
||||
|
||||
@@ -14,7 +14,6 @@
|
||||
|
||||
#define RECEIVE_TIMEOUT_SEC 60 /* 60 second per-message timeout */
|
||||
#define SEND_TIMEOUT_SEC 60
|
||||
#define MAX_CONNECTION_MEMORY (256ULL * 1024 * 1024) /* bounded cumulative receive budget */
|
||||
|
||||
static __thread int io_read_fd = -1;
|
||||
static __thread int io_write_fd = -1;
|
||||
|
||||
@@ -28,6 +28,10 @@
|
||||
#define DEFAULT_MAX_ALLOC (1ULL * 1024 * 1024 * 1024)
|
||||
/* Server policy ceiling for a client-provided allocation limit. */
|
||||
#define MAX_SERVER_ALLOC (256ULL * 1024 * 1024)
|
||||
/* Bounded cumulative per-connection receive budget. In-flight wire buffers,
|
||||
decompression buffers and queued (not yet written) file payloads for a
|
||||
connection must stay within this ceiling. */
|
||||
#define MAX_CONNECTION_MEMORY (256ULL * 1024 * 1024)
|
||||
|
||||
typedef struct ssl_st SSL;
|
||||
|
||||
|
||||
@@ -5,6 +5,7 @@
|
||||
#include "utils.h"
|
||||
#include "protocol.h"
|
||||
#include "test_utils.h"
|
||||
#include <fcntl.h>
|
||||
#include <stdlib.h>
|
||||
#include <string.h>
|
||||
#include <sys/stat.h>
|
||||
@@ -730,6 +731,161 @@ static void test_file_send_single_calls_metadata_and_path() {
|
||||
}
|
||||
}
|
||||
|
||||
static void test_inplace_overwrite_clears_special_mode_bits() {
|
||||
const char* root = "test_inplace_tmp";
|
||||
const char* path = "test_inplace_tmp/priv.txt";
|
||||
const char* content = "olddata";
|
||||
unlink(path);
|
||||
rmdir(root);
|
||||
EXPECT_EQ_INT(mkdir(root, 0700), 0);
|
||||
|
||||
/* Create a destination carrying setuid + sticky bits. */
|
||||
int fd = open(path, O_WRONLY | O_CREAT | O_TRUNC | O_CLOEXEC, 0644);
|
||||
EXPECT_TRUE(fd >= 0);
|
||||
// cppcheck-suppress knownConditionTrueFalse
|
||||
if (fd < 0) {
|
||||
rmdir(root);
|
||||
return;
|
||||
}
|
||||
EXPECT_EQ_INT((int)write(fd, content, strlen(content)), (int)strlen(content));
|
||||
EXPECT_EQ_INT(fchmod(fd, S_ISUID | S_ISVTX | 0755), 0);
|
||||
EXPECT_EQ_INT(close(fd), 0);
|
||||
|
||||
/* Overwrite in place without metadata: the mode must be normalized to a
|
||||
safe default (0644) and the setuid/sticky bits must be gone. */
|
||||
File* f = file_create("priv.txt");
|
||||
EXPECT_NOT_NULL(f);
|
||||
const char* new_content = "newdata";
|
||||
f->data->data = malloc(strlen(new_content));
|
||||
EXPECT_NOT_NULL(f->data->data);
|
||||
memcpy(f->data->data, new_content, strlen(new_content));
|
||||
f->data->size = strlen(new_content);
|
||||
|
||||
Config* cfg = config_create();
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
cfg->inplace = true;
|
||||
EXPECT_TRUE(file_save_to_disk(root, f, cfg));
|
||||
file_destroy(f);
|
||||
config_delete(cfg);
|
||||
|
||||
struct stat st;
|
||||
EXPECT_EQ_INT(stat(path, &st), 0);
|
||||
EXPECT_EQ_INT((int)(st.st_mode & (S_ISUID | S_ISGID | S_ISVTX)), 0);
|
||||
EXPECT_EQ_INT((int)(st.st_mode & 0777), 0644);
|
||||
FILE* stream = fopen(path, "rb");
|
||||
char buf[16] = {0};
|
||||
EXPECT_NOT_NULL(stream);
|
||||
// cppcheck-suppress knownConditionTrueFalse
|
||||
if (stream) {
|
||||
size_t nread = fread(buf, 1, sizeof(buf) - 1, stream);
|
||||
fclose(stream);
|
||||
EXPECT_EQ_INT((int)nread, (int)strlen(new_content));
|
||||
}
|
||||
EXPECT_EQ_STR(buf, new_content);
|
||||
|
||||
unlink(path);
|
||||
rmdir(root);
|
||||
}
|
||||
|
||||
static void test_inplace_overwrite_metadata_strips_special_bits() {
|
||||
const char* root = "test_inplace_meta_tmp";
|
||||
const char* path = "test_inplace_meta_tmp/meta.txt";
|
||||
const char* source = "test_inplace_meta_source.txt";
|
||||
unlink(path);
|
||||
unlink(source);
|
||||
rmdir(root);
|
||||
EXPECT_EQ_INT(mkdir(root, 0700), 0);
|
||||
|
||||
/* Existing destination with setuid+sticky set. */
|
||||
int fd = open(path, O_WRONLY | O_CREAT | O_TRUNC | O_CLOEXEC, 0644);
|
||||
EXPECT_TRUE(fd >= 0);
|
||||
// cppcheck-suppress knownConditionTrueFalse
|
||||
if (fd < 0) {
|
||||
rmdir(root);
|
||||
return;
|
||||
}
|
||||
EXPECT_EQ_INT((int)write(fd, "olddata", 7), 7);
|
||||
EXPECT_EQ_INT(fchmod(fd, S_ISUID | S_ISVTX | 0755), 0);
|
||||
EXPECT_EQ_INT(close(fd), 0);
|
||||
|
||||
/* Build source metadata carrying a plain executable mode (no specials). */
|
||||
EXPECT_TRUE(file_write_to_disk(source, "source", 6, false, false));
|
||||
EXPECT_EQ_INT(chmod(source, 0755), 0);
|
||||
struct stat source_st;
|
||||
EXPECT_EQ_INT(stat(source, &source_st), 0);
|
||||
|
||||
File* f = file_create("meta.txt");
|
||||
EXPECT_NOT_NULL(f);
|
||||
const char* new_content = "meta";
|
||||
f->data->data = malloc(strlen(new_content));
|
||||
EXPECT_NOT_NULL(f->data->data);
|
||||
memcpy(f->data->data, new_content, strlen(new_content));
|
||||
f->data->size = strlen(new_content);
|
||||
f->metadata = file_metadata_create(&source_st);
|
||||
EXPECT_NOT_NULL(f->metadata);
|
||||
|
||||
Config* cfg = config_create();
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
cfg->inplace = true;
|
||||
EXPECT_TRUE(file_save_to_disk(root, f, cfg));
|
||||
file_destroy(f);
|
||||
config_delete(cfg);
|
||||
unlink(source);
|
||||
|
||||
struct stat st;
|
||||
EXPECT_EQ_INT(stat(path, &st), 0);
|
||||
/* Metadata-derived mode is applied and never includes setuid/setgid/sticky. */
|
||||
EXPECT_EQ_INT((int)(st.st_mode & (S_ISUID | S_ISGID | S_ISVTX)), 0);
|
||||
EXPECT_EQ_INT((int)(st.st_mode & 0777), 0755);
|
||||
|
||||
unlink(path);
|
||||
rmdir(root);
|
||||
}
|
||||
|
||||
static void test_inplace_overwrite_truncates_shorter_payload() {
|
||||
const char* root = "test_inplace_trunc_tmp";
|
||||
const char* path = "test_inplace_trunc_tmp/big.txt";
|
||||
unlink(path);
|
||||
rmdir(root);
|
||||
EXPECT_EQ_INT(mkdir(root, 0700), 0);
|
||||
|
||||
const char* old_content = "0123456789abcdef"; /* 16 bytes */
|
||||
EXPECT_TRUE(file_write_to_disk(path, old_content, strlen(old_content), false, false));
|
||||
|
||||
File* f = file_create("big.txt");
|
||||
EXPECT_NOT_NULL(f);
|
||||
const char* new_content = "hi";
|
||||
f->data->data = malloc(strlen(new_content));
|
||||
EXPECT_NOT_NULL(f->data->data);
|
||||
memcpy(f->data->data, new_content, strlen(new_content));
|
||||
f->data->size = strlen(new_content);
|
||||
|
||||
Config* cfg = config_create();
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
cfg->inplace = true;
|
||||
EXPECT_TRUE(file_save_to_disk(root, f, cfg));
|
||||
file_destroy(f);
|
||||
config_delete(cfg);
|
||||
|
||||
/* A shorter payload must truncate the file: no stale trailing bytes. */
|
||||
struct stat st;
|
||||
EXPECT_EQ_INT(stat(path, &st), 0);
|
||||
EXPECT_EQ_INT((int)st.st_size, (int)strlen(new_content));
|
||||
FILE* stream = fopen(path, "rb");
|
||||
char buf[32] = {0};
|
||||
EXPECT_NOT_NULL(stream);
|
||||
// cppcheck-suppress knownConditionTrueFalse
|
||||
if (stream) {
|
||||
size_t nread = fread(buf, 1, sizeof(buf) - 1, stream);
|
||||
fclose(stream);
|
||||
EXPECT_EQ_INT((int)nread, (int)strlen(new_content));
|
||||
}
|
||||
EXPECT_EQ_STR(buf, new_content);
|
||||
|
||||
unlink(path);
|
||||
rmdir(root);
|
||||
}
|
||||
|
||||
void test_file() {
|
||||
test_file_create();
|
||||
test_file_destroy_null();
|
||||
@@ -762,4 +918,7 @@ void test_file() {
|
||||
test_file_send_single_calls_metadata_and_path();
|
||||
}
|
||||
test_file_metadata_create();
|
||||
test_inplace_overwrite_clears_special_mode_bits();
|
||||
test_inplace_overwrite_metadata_strips_special_bits();
|
||||
test_inplace_overwrite_truncates_shorter_payload();
|
||||
}
|
||||
|
||||
@@ -250,6 +250,87 @@ static void test_write_thread_done() {
|
||||
config_delete(cfg);
|
||||
}
|
||||
|
||||
typedef struct {
|
||||
PipelineContextReceiver* context;
|
||||
File* file;
|
||||
atomic_bool* done;
|
||||
atomic_bool* result;
|
||||
} ByteBudgetEnqueueArg;
|
||||
|
||||
static int byte_budget_enqueue_worker(void* arg) {
|
||||
ByteBudgetEnqueueArg* worker = arg;
|
||||
bool ok = pipeline_context_receiver_enqueue_file(worker->context, worker->file);
|
||||
atomic_store(worker->result, ok);
|
||||
atomic_store(worker->done, true);
|
||||
return thrd_success;
|
||||
}
|
||||
|
||||
/* A receiver must not buffer more decompressed/copied payload bytes ahead of
|
||||
the (slow) disk writer than the configured byte budget: an enqueue that
|
||||
would exceed the budget blocks until the writer releases bytes. */
|
||||
static void test_receiver_enqueue_byte_budget() {
|
||||
Config* cfg = config_create();
|
||||
EXPECT_NOT_NULL(cfg);
|
||||
free(cfg->version);
|
||||
cfg->version = str_dup(PROTOCOL_VERSION);
|
||||
cfg->send_directory = str_dup("/src");
|
||||
cfg->receive_root_directory = str_dup("/dst");
|
||||
cfg->save_to_disk = false;
|
||||
|
||||
Queue* q = queue_create(16, file_destroy);
|
||||
EXPECT_NOT_NULL(q);
|
||||
PipelineContextReceiver* ctx = pipeline_context_receiver_create(cfg, q, -1, NULL);
|
||||
EXPECT_NOT_NULL(ctx);
|
||||
pipeline_context_receiver_set_queue_byte_limit(ctx, 3000);
|
||||
ctx->receiver_done = false;
|
||||
|
||||
File* first = file_create("budget_file_1");
|
||||
EXPECT_NOT_NULL(first);
|
||||
first->data->size = 2000;
|
||||
EXPECT_TRUE(pipeline_context_receiver_enqueue_file(ctx, first));
|
||||
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000);
|
||||
|
||||
/* Second 2000-byte payload would push the pipeline to 4000 > 3000 budget,
|
||||
so the enqueue must block until the first payload is released. */
|
||||
File* second = file_create("budget_file_2");
|
||||
EXPECT_NOT_NULL(second);
|
||||
second->data->size = 2000;
|
||||
atomic_bool done;
|
||||
atomic_bool result;
|
||||
atomic_init(&done, false);
|
||||
atomic_init(&result, false);
|
||||
ByteBudgetEnqueueArg arg = {ctx, second, &done, &result};
|
||||
thrd_t enqueuer;
|
||||
EXPECT_EQ_INT(thrd_create(&enqueuer, byte_budget_enqueue_worker, &arg), thrd_success);
|
||||
|
||||
/* Give a broken (unbounded) implementation every chance to enqueue. */
|
||||
struct timespec wait = {0, 200 * 1000000L};
|
||||
thrd_sleep(&wait, NULL);
|
||||
EXPECT_FALSE(atomic_load(&done));
|
||||
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* budget still honored */
|
||||
|
||||
/* Simulate the disk writer: dequeue + destroy + release the first file. */
|
||||
File* drained = queue_dequeue_multithreaded(q, &ctx->mutex, &ctx->condition_not_empty,
|
||||
&ctx->condition_not_full, &ctx->receiver_done);
|
||||
EXPECT_NOT_NULL(drained);
|
||||
file_destroy(drained);
|
||||
pipeline_context_receiver_note_bytes_released(ctx, 2000);
|
||||
EXPECT_EQ_INT((int)ctx->queued_bytes, 0);
|
||||
|
||||
EXPECT_EQ_INT(thrd_join(enqueuer, NULL), thrd_success);
|
||||
EXPECT_TRUE(atomic_load(&done));
|
||||
EXPECT_TRUE(atomic_load(&result));
|
||||
EXPECT_EQ_INT((int)ctx->queued_bytes, 2000); /* second payload now in flight */
|
||||
|
||||
/* Tear down: the second file is still queued and is freed by queue_destroy. */
|
||||
mtx_destroy(&ctx->mutex);
|
||||
cnd_destroy(&ctx->condition_not_full);
|
||||
cnd_destroy(&ctx->condition_not_empty);
|
||||
free(ctx);
|
||||
queue_destroy(q);
|
||||
config_delete(cfg);
|
||||
}
|
||||
|
||||
void test_multiprocessing() {
|
||||
test_sender_create_destroy();
|
||||
test_receiver_create_destroy();
|
||||
@@ -261,4 +342,5 @@ void test_multiprocessing() {
|
||||
test_receive_thread_failure_wakes_writer();
|
||||
}
|
||||
test_write_thread_done();
|
||||
test_receiver_enqueue_byte_budget();
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user