Skip to content
Draft
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
96 changes: 65 additions & 31 deletions plugins/in_emitter/emitter.c
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,10 @@

#define DEFAULT_EMITTER_RING_BUFFER_FLUSH_FREQUENCY 2000

/* Core checks its chunk target after appending. Keep each deferred submission
* small enough that combining it with an existing chunk adds modest overshoot. */
#define EMITTER_BUFFER_TARGET (256 * 1024)

/* return values */
#define FLB_EMITTER_BUSY -2

Expand All @@ -59,11 +63,13 @@ struct flb_emitter {
struct mk_list i_ins_list; /* instance list of linked/sending inputs */
int cycle_reported; /* self-emission cycle already logged ? */
int propagating; /* pause/resume propagation in progress */
int draining; /* queued buffers being submitted */
int owner_flush_enabled; /* owner supports lossless shutdown flush */
int shutdown_flush_done; /* owner shutdown hook already invoked */
int shutdown_flush_active; /* owner shutdown hook is emitting */
};

/* Only callback owners may bypass memory admission for their final flush. */
static inline int shutdown_flush_enabled(struct flb_emitter *ctx)
{
if (ctx->ins->config->is_shutting_down == FLB_TRUE &&
Expand Down Expand Up @@ -109,7 +115,9 @@ static void em_chunk_destroy(struct flb_emitter *ctx, struct em_chunk *ec)
ctx->pending_bytes = 0;
}

mk_list_del(&ec->_head);
if (!mk_list_entry_is_orphan(&ec->_head)) {
mk_list_del(&ec->_head);
}
flb_sds_destroy(ec->tag);
msgpack_sbuffer_destroy(&ec->mp_sbuf);
flb_free(ec);
Expand Down Expand Up @@ -155,13 +163,18 @@ int static do_in_emitter_add_record(struct em_chunk *ec,
struct flb_emitter *ctx = (struct flb_emitter *) in->context;
int ret;

/* Accepted buffers may drain despite a shutdown pause; new records still
* pass through the ordinary admission check in in_emitter_add_record(). */
if (flb_input_buf_paused(ctx->ins) == FLB_TRUE &&
shutdown_flush_enabled(ctx) == FLB_FALSE) {
in->config->is_shutting_down == FLB_FALSE) {
flb_plg_debug(ctx->ins, "_emitter %s paused. Not processing records.",
ctx->ins->name);
return FLB_EMITTER_BUSY;
}

/* Make the in-flight buffer unavailable to rewrites triggered by append. */
mk_list_del(&ec->_head);

/* Associate this backlog chunk to this instance into the engine */
ret = flb_input_log_append(in,
ec->tag, flb_sds_len(ec->tag),
Expand All @@ -177,10 +190,7 @@ int static do_in_emitter_add_record(struct em_chunk *ec,
return 0;
}

/*
* Function used by filters to ingest custom records with custom tags, at the
* moment it's only used by rewrite_tag filter.
*/
/* Queue records emitted by rewrite_tag and buffered multiline filters. */
int in_emitter_add_record(const char *tag, int tag_len,
const char *buf_data, size_t buf_size,
struct flb_input_instance *in,
Expand Down Expand Up @@ -274,11 +284,8 @@ int in_emitter_add_record(const char *tag, int tag_len,
sizeof(struct em_chunk));
}

/*
* A callback-owning filter cannot retry its in-flight batch. Allow one
* input-chunk-sized pending batch while normal chunk accounting pauses
* the emitter, but retain a hard bound under sustained backpressure.
*/
/* A callback-owning filter cannot retry its in-flight batch. Retain its
* pending allowance under pressure, and allow shutdown flushes to finish. */
if (shutdown_flush_enabled(ctx) == FLB_FALSE &&
is_queue_overlimit(ctx, buf_size) == FLB_TRUE &&
(ctx->owner_flush_enabled == FLB_FALSE ||
Expand All @@ -288,8 +295,9 @@ int in_emitter_add_record(const char *tag, int tag_len,
return FLB_EMITTER_BUSY;
}

/* Check if any target chunk already exists */
mk_list_foreach(head, &ctx->chunks) {
/* Only the newest buffer for a tag may accept records. Older buffers were
* rolled over and must not be backfilled with later, smaller records. */
mk_list_foreach_r(head, &ctx->chunks) {
ec = mk_list_entry(head, struct em_chunk, _head);
if (flb_sds_cmp(ec->tag, tag, tag_len) != 0) {
ec = NULL;
Expand All @@ -298,6 +306,15 @@ int in_emitter_add_record(const char *tag, int tag_len,
break;
}

/* Callers supply complete records, including any group framing. Keep each
* submission intact, even when it exceeds the target. Leave older buffers
* queued: ingesting them here would re-enter the filter pipeline. */
if (ec && ec->mp_sbuf.size > 0 &&
(buf_size > EMITTER_BUFFER_TARGET ||
ec->mp_sbuf.size > EMITTER_BUFFER_TARGET - buf_size)) {
ec = NULL;
}

/* No candidate chunk found, so create a new one */
if (!ec) {
ec = em_chunk_create(tag, tag_len, ctx);
Expand All @@ -321,6 +338,10 @@ int in_emitter_add_record(const char *tag, int tag_len,
/* Append raw msgpack data */
ret = msgpack_sbuffer_write(&ec->mp_sbuf, buf_data, buf_size);
if (ret != 0) {
/* Do not leave an empty successor queued after allocation failure. */
if (ec->mp_sbuf.size == 0) {
em_chunk_destroy(ctx, ec);
}
return -1;
}

Expand Down Expand Up @@ -365,6 +386,12 @@ static int cb_queue_chunks(struct flb_input_instance *in,

/* Get context */
ctx = (struct flb_emitter *) data;
/* Core ingestion can pause this emitter synchronously and call back into
* cb_emitter_pause(). Leave the outer drain in sole charge of the queue. */
if (ctx->draining == FLB_TRUE) {
return 0;
}
ctx->draining = FLB_TRUE;

/* Try to enqueue chunks under our limits */
mk_list_foreach_safe(head, tmp, &ctx->chunks) {
Expand All @@ -377,6 +404,7 @@ static int cb_queue_chunks(struct flb_input_instance *in,
}
}

ctx->draining = FLB_FALSE;
return 0;
}

Expand Down Expand Up @@ -430,9 +458,9 @@ static int cb_emitter_init(struct flb_input_instance *in,
callbacks = (struct flb_emitter_callbacks *) in->data;
if (callbacks != NULL && callbacks->pause != NULL) {
ctx->owner_flush_enabled = FLB_TRUE;
in->flags |= FLB_INPUT_SHUTDOWN_FLUSH;
}

/* Accepted buffers must drain even if this emitter was already paused. */
in->flags |= FLB_INPUT_SHUTDOWN_FLUSH;

ret = flb_input_config_map_set(in, (void *) ctx);
if (ret == -1) {
Expand Down Expand Up @@ -506,29 +534,35 @@ static void cb_emitter_pause(void *data, struct flb_config *config)
}
ctx->propagating = FLB_TRUE;

if (shutdown_flush_enabled(ctx) == FLB_TRUE &&
ctx->shutdown_flush_done == FLB_FALSE &&
ctx->ins->data != NULL) {
callbacks = (struct flb_emitter_callbacks *) ctx->ins->data;
ctx->shutdown_flush_done = FLB_TRUE;
if (config->is_shutting_down == FLB_TRUE && config->is_running == FLB_TRUE) {
if (ctx->msgs != NULL) {
in_emitter_ingest_ring_buffer(ctx->ins, config, ctx);
}

if (callbacks->pause != NULL) {
/* Publish any records already waiting in a threaded emitter. */
if (ctx->msgs != NULL) {
in_emitter_ingest_ring_buffer(ctx->ins, config, ctx);
callbacks = (struct flb_emitter_callbacks *) ctx->ins->data;
if (ctx->shutdown_flush_done == FLB_FALSE && callbacks != NULL) {
ctx->shutdown_flush_done = FLB_TRUE;
if (callbacks->pause != NULL) {
ctx->shutdown_flush_active = FLB_TRUE;
callbacks->pause(callbacks->data, config);
ctx->shutdown_flush_active = FLB_FALSE;
}
}

ctx->shutdown_flush_active = FLB_TRUE;
callbacks->pause(callbacks->data, config);
ctx->shutdown_flush_active = FLB_FALSE;

/* Convert owner-emitted records into input chunks before pausing. */
cb_queue_chunks(ctx->ins, config, ctx);
cb_queue_chunks(ctx->ins, config, ctx);
/* Chained emitters can receive more records from another input's
* shutdown hook after their own pause callback has returned. Keep the
* internal collector active during grace, including after backpressure.
* Normal collector resume is disabled once external ingestion stops. */
if (flb_input_collector_start(ctx->coll_fd, ctx->ins) != 0) {
flb_plg_error(ctx->ins, "could not start collector for shutdown flush");
}
}
else {
flb_input_collector_pause(ctx->coll_fd, ctx->ins);
}

/* Pause all known senders */
flb_input_collector_pause(ctx->coll_fd, ctx->ins);
mk_list_foreach_safe(head, tmp, &ctx->i_ins_list) {
i_ref = mk_list_entry(head, struct input_ref, _head);
flb_input_pause(i_ref->i_ins);
Expand Down
1 change: 1 addition & 0 deletions tests/runtime/CMakeLists.txt
Original file line number Diff line number Diff line change
Expand Up @@ -230,6 +230,7 @@ if(FLB_IN_LIB AND FLB_OUT_LIB)
FLB_RT_TEST(FLB_FILTER_THROTTLE_SIZE "filter_throttle_size.c")
FLB_RT_TEST(FLB_FILTER_NEST "filter_nest.c")
FLB_RT_TEST(FLB_FILTER_REWRITE_TAG "filter_rewrite_tag.c")
FLB_RT_TEST(FLB_FILTER_REWRITE_TAG "filter_rewrite_tag_burst.c")
FLB_RT_TEST(FLB_FILTER_KUBERNETES "filter_kubernetes.c")
FLB_RT_TEST(FLB_FILTER_PARSER "filter_parser.c")
FLB_RT_TEST(FLB_FILTER_MODIFY "filter_modify.c")
Expand Down
Loading
Loading