Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
27 changes: 15 additions & 12 deletions src/flb_input_chunk.c
Original file line number Diff line number Diff line change
Expand Up @@ -1102,7 +1102,7 @@ int flb_input_chunk_release_space_compound(
* will drop the the oldest chunks when the limitation on local disk is reached.
*/
int flb_input_chunk_find_space_new_data(struct flb_input_chunk *ic,
size_t chunk_size, int overlimit)
size_t chunk_size)
{
int count;
int result;
Expand All @@ -1111,20 +1111,23 @@ int flb_input_chunk_find_space_new_data(struct flb_input_chunk *ic,
size_t local_release_requirement;

/*
* For each output instances that will be over the limit after adding the new chunk,
* we have to determine how many chunks needs to be removed. We will adjust the
* routes_mask to only route to the output plugin that have enough space after
* deleting some chunks fome the queue.
* For each output instance that will be over the limit after adding the new chunk,
* we have to determine how many chunks need to be removed. We will adjust the
* routes_mask to only route to the output plugin that has enough space after
* deleting some chunks from the queue.
*/
count = 0;

mk_list_foreach(head, &ic->in->config->outputs) {
o_ins = mk_list_entry(head, struct flb_output_instance, _head);

if ((o_ins->total_limit_size == -1) || ((1 << o_ins->id) & overlimit) == 0 ||
(flb_routes_mask_get_bit(ic->routes_mask,
o_ins->id,
o_ins->config->router) == 0)) {
if ((o_ins->total_limit_size == -1) ||
(flb_routes_mask_get_bit(ic->routes_mask,
o_ins->id,
o_ins->config->router) == 0) ||
(o_ins->fs_chunks_size +
o_ins->fs_backlog_chunks_size +
chunk_size) <= o_ins->total_limit_size) {
continue;
}

Expand Down Expand Up @@ -1182,7 +1185,7 @@ int flb_input_chunk_has_overlimit_routes(struct flb_input_chunk *ic,
if ((o_ins->fs_chunks_size +
o_ins->fs_backlog_chunks_size +
chunk_size) > o_ins->total_limit_size) {
overlimit |= (1 << o_ins->id);
overlimit = FLB_TRUE;
}
}

Expand All @@ -1195,13 +1198,13 @@ int flb_input_chunk_has_overlimit_routes(struct flb_input_chunk *ic,
int flb_input_chunk_place_new_chunk(struct flb_input_chunk *ic, size_t chunk_size)
{
int result;
int overlimit;
int overlimit;
struct flb_input_instance *i_ins = ic->in;

if (i_ins->storage_type == CIO_STORE_FS) {
overlimit = flb_input_chunk_has_overlimit_routes(ic, chunk_size);
if (overlimit != 0) {
result = flb_input_chunk_find_space_new_data(ic, chunk_size, overlimit);
result = flb_input_chunk_find_space_new_data(ic, chunk_size);

if (result != 0) {
return 0;
Expand Down
169 changes: 169 additions & 0 deletions tests/internal/input_chunk.c
Original file line number Diff line number Diff line change
Expand Up @@ -1123,6 +1123,173 @@ void flb_test_input_chunk_prefers_deletable_files_on_limit(void)
flb_free(root_path);
}

void flb_test_input_chunk_limit_with_many_outputs(void)
{
int i;
int records;
struct flb_input_instance *i_ins;
struct flb_output_instance *o_low;
struct flb_output_instance *o_filler;
struct flb_output_instance *o_high;
struct mk_list *head;
struct mk_list *tmp;
struct flb_input_chunk *ic;
struct flb_task *task;
struct flb_config *cfg;
struct cio_ctx *cio;
struct mk_event_loop *evl;
struct cio_options opts = {0};
char *root_path;
char temp_path[128];
char buf[2048];
ssize_t chunk_real_size;
size_t first_chunk_size;

snprintf(temp_path, sizeof(temp_path) - 1,
"/input-chunk-many-outputs-%i/", getpid());
temp_path[sizeof(temp_path) - 1] = '\0';

root_path = flb_test_tmpdir_cat(temp_path);
TEST_CHECK(root_path != NULL);
if (!root_path) {
return;
}

memset(buf, 0x5A, sizeof(buf));

flb_init_env();
cfg = flb_config_init();
evl = mk_event_loop_create(256);

TEST_CHECK(evl != NULL);
if (!evl) {
flb_config_exit(cfg);
flb_free(root_path);
return;
}

cfg->evl = evl;
flb_log_create(cfg, FLB_LOG_STDERR, FLB_LOG_DEBUG, NULL);

i_ins = flb_input_new(cfg, "dummy", NULL, FLB_TRUE);
TEST_CHECK(i_ins != NULL);
if (!i_ins) {
flb_config_exit(cfg);
flb_free(root_path);
return;
}
i_ins->storage_type = CIO_STORE_FS;

cio_options_init(&opts);
opts.root_path = root_path;
opts.log_cb = log_cb;
opts.log_level = CIO_LOG_DEBUG;
opts.flags = CIO_OPEN;

cio = cio_create(&opts);
TEST_CHECK(cio != NULL);
if (!cio) {
flb_input_exit_all(cfg);
flb_output_exit(cfg);
flb_config_exit(cfg);
flb_free(root_path);
return;
}

flb_storage_input_create(cio, i_ins);
flb_input_init_all(cfg);

o_low = flb_output_new(cfg, "http", NULL, FLB_TRUE);
TEST_CHECK(o_low != NULL);
if (!o_low) {
cio_destroy(cio);
flb_input_exit_all(cfg);
flb_output_exit(cfg);
flb_config_exit(cfg);
flb_free(root_path);
return;
}
flb_output_set_property(o_low, "match", "many.*");
flb_output_set_property(o_low, "storage.total_limit_size", "10M");

for (i = 0; i < 31; i++) {
o_filler = flb_output_new(cfg, "http", NULL, FLB_TRUE);
TEST_CHECK(o_filler != NULL);
if (!o_filler) {
cio_destroy(cio);
flb_input_exit_all(cfg);
flb_output_exit(cfg);
flb_config_exit(cfg);
flb_free(root_path);
return;
}
flb_output_set_property(o_filler, "match", "filler.*");
}

o_high = flb_output_new(cfg, "http", NULL, FLB_TRUE);
TEST_CHECK(o_high != NULL);
if (!o_high) {
cio_destroy(cio);
flb_input_exit_all(cfg);
flb_output_exit(cfg);
flb_config_exit(cfg);
flb_free(root_path);
return;
}
TEST_CHECK(o_low->id == 0);
TEST_CHECK(o_high->id == 32);
flb_output_set_property(o_high, "match", "many.*");
flb_output_set_property(o_high, "storage.total_limit_size", "10M");

TEST_CHECK(flb_routes_mask_set_size(mk_list_size(&cfg->outputs),
cfg->router) == 0);
TEST_CHECK_(flb_router_io_set(cfg) != -1, "unable to router");

records = flb_mp_count(buf, sizeof(buf));

TEST_CHECK(flb_input_chunk_append_raw(i_ins, FLB_INPUT_LOGS,
records, "many.one", 8,
buf, sizeof(buf)) == 0);
ic = mk_list_entry_last(&i_ins->chunks, struct flb_input_chunk, _head);
chunk_real_size = flb_input_chunk_get_real_size(ic);
TEST_CHECK(chunk_real_size > 0);
if (chunk_real_size <= 0) {
goto cleanup;
}
first_chunk_size = (size_t) chunk_real_size;
o_high->total_limit_size = first_chunk_size + (first_chunk_size / 2);
Comment thread
coderabbitai[bot] marked this conversation as resolved.

TEST_CHECK(flb_input_chunk_append_raw(i_ins, FLB_INPUT_LOGS,
records, "many.two", 8,
buf, sizeof(buf)) == 0);

TEST_CHECK(mk_list_size(&i_ins->chunks) == 2);
ic = mk_list_entry_first(&i_ins->chunks, struct flb_input_chunk, _head);
TEST_CHECK(flb_routes_mask_get_bit(ic->routes_mask,
o_low->id, cfg->router) == 1);
TEST_CHECK(flb_routes_mask_get_bit(ic->routes_mask,
o_high->id, cfg->router) == 0);
TEST_CHECK(o_high->fs_chunks_size <= o_high->total_limit_size);

cleanup:
mk_list_foreach_safe(head, tmp, &i_ins->tasks) {
task = mk_list_entry(head, struct flb_task, _head);
flb_task_destroy(task, FLB_TRUE);
}

mk_list_foreach_safe(head, tmp, &i_ins->chunks) {
ic = mk_list_entry(head, struct flb_input_chunk, _head);
flb_input_chunk_destroy(ic, FLB_TRUE);
}

cio_destroy(cio);
flb_router_exit(cfg);
flb_input_exit_all(cfg);
flb_output_exit(cfg);
flb_config_exit(cfg);
flb_free(root_path);
}


/* Test list */
TEST_LIST = {
Expand All @@ -1136,5 +1303,7 @@ TEST_LIST = {
flb_test_input_chunk_grouped_release_space_drop_counters},
{"input_chunk_prefers_deletable_files_on_limit",
flb_test_input_chunk_prefers_deletable_files_on_limit},
{"input_chunk_limit_with_many_outputs",
flb_test_input_chunk_limit_with_many_outputs},
{NULL, NULL}
};