Skip to content

in_emitter: bound pending batches and drain them on shutdown - #12521

Draft
terrorobe wants to merge 2 commits into
fluent:masterfrom
terrorobe:terrorobe/emitter-buffer-rollover-master-pr
Draft

terrorobe wants to merge 2 commits into
fluent:masterfrom
terrorobe:terrorobe/emitter-buffer-rollover-master-pr

Conversation

@terrorobe

Copy link
Copy Markdown
Contributor

This draft PR targets master and depends on the group-framing change below. The branch contains only emitter changes and must not merge before the prerequisite.

Currently, the emitter appends each tag's entire pending buffer to the engine in one piece, so a rewrite_tag burst can produce chunks far larger than the 2 MB chunk target.

Batching changes:

  • Start a new pending buffer for a tag once the next record would push the current one past 256 KiB.
  • Leave completed buffers queued for the collector, and append new records only to the newest buffer for their tag. This preserves per-tag order without re-entering the filter pipeline.
  • Keep each record intact, even when it exceeds 256 KiB. Complete group envelopes come from rewrite_tag: mitigate dropping metadata on rewrite tag #12500; the emitter does not decode, copy, or reconstruct group markers.

With smaller batches, a backpressure pause can leave accepted records queued in the emitter, and records still queued in a rewrite_tag emitter at shutdown are currently discarded. Shutdown changes:

  • Drain accepted buffers during graceful shutdown, even if backpressure has already paused the emitter.
  • Detach each buffer before submitting it and prevent nested drains, so a backpressure-triggered pause cannot resubmit, modify, or free it.
  • Keep internal emitter collectors running during the grace period so chained emitters can flush late records in either creation order.
  • New records still respect the emitter memory limit during shutdown, even if a source delivers them after being paused. The existing allowance for callback-owning filters, such as buffered multiline, still applies.

Fixes #12519.

Testing

Tested on Linux/amd64 with master a415d1b and the prerequisite at 8795652. This PR branch intentionally excludes the prerequisite; its grouped-record tests are expected to fail against master until #12500 merges. Once it lands, we will update this branch and rerun CI.

The new runtime test, tests/runtime/filter_rewrite_tag_burst.c, sends bursts through the real engine and verifies every delivered record: identity, order, payload, timestamp, and group context. It covers:

  • chunk size bounds, including oversized records and partially filled core chunks
  • keep=true duplication, chained rewrites, and native output processors
  • memory and filesystem backpressure, including shutdown while paused
  • shutdown with records pending in chained emitters, in either creation order
  • late records that hit memory or filesystem limits during shutdown
  • memory admission for new records during shutdown

Without the emitter changes, the size checks fail. The new test and the existing multiline tests (the other emitter user) pass under ASan/UBSan and Valgrind; the rewrite_tag and multiline integration tests pass under strict Valgrind.

  • Example configuration: see the manual reproduction below
  • Debug log output from testing the change (below)
  • Valgrind output showing no leaks or memory errors (below)
  • [N/A] Local packaging tests: packaging is unchanged
  • [N/A] ok-package-test label: packaging is unchanged
Regression output and verification commands

The automated burst regression uses an in-process collector and library output callbacks. To reproduce the burst manually with a local build, run this dummy → rewrite_tag → null pipeline from the repository root:

record=$(python3 -c 'import json; print(json.dumps({"log": "x" * 4096}))')
./build/bin/fluent-bit \
  -i dummy -t source -p "dummy=$record" -p copies=3000 -p samples=1 \
  -F rewrite_tag -m source -p 'rule=$log ^x routed false' -p emitter_mem_buf_limit=64M \
  -o null -m routed -p alias=routed -p log_level=debug -p workers=0

Wait for the output to drain, then press Ctrl-C. With only the prerequisite applied, the whole burst arrives as one chunk:

[2026/10/05 21:36:28.715] [debug] [output:null:routed] discarding 12351000 bytes

With this change:

[2026/10/06 00:31:46.515] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.515] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.515] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.516] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.516] [debug] [output:null:routed] discarding 2074968 bytes
[2026/10/06 00:31:46.517] [debug] [output:null:routed] discarding 1976160 bytes
  • Total after rollover: 5 × 2,074,968 + 1,976,160 = 12,351,000 bytes, matching the baseline.
  • Bound asserted for small framed records without filter expansion: one core chunk target plus one emitter batch, 2,048,000 + 262,144 = 2,310,144 bytes.

Equivalent build and CTest commands, run from the repository root:

cmake -S . -B build -DFLB_TESTS_RUNTIME=On -DFLB_TESTS_INTERNAL=On
cmake --build build --target fluent-bit-bin flb-rt-filter_rewrite_tag \
  flb-rt-filter_rewrite_tag_burst flb-rt-filter_multiline \
  flb-it-log_event_decoder flb-it-log_event_encoder flb-it-output -j8
ctest --test-dir build \
  -R '^flb-(rt-filter_(rewrite_tag(_burst)?|multiline)|it-(log_event_(decoder|encoder)|output))$' \
  --output-on-failure --timeout 180

To run the existing group-preservation and multiline integration scenarios from the repository root:

./tests/integration/setup-venv.sh
FLUENT_BIT_BINARY="$(pwd)/build/bin/fluent-bit"
cd tests/integration
FLUENT_BIT_BINARY="$FLUENT_BIT_BINARY" .venv/bin/python -m pytest \
  scenarios/filter_rewrite_tag scenarios/filter_multiline -q
VALGRIND=1 VALGRIND_STRICT=1 FLUENT_BIT_BINARY="$FLUENT_BIT_BINARY" \
  .venv/bin/python -m pytest scenarios/filter_rewrite_tag scenarios/filter_multiline -q

Valgrind command, run from the repository root, and output excerpt:

valgrind --leak-check=full --show-leak-kinds=all \
  --errors-for-leak-kinds=definite,indirect --error-exitcode=99 \
  ./build/bin/flb-rt-filter_rewrite_tag_burst --exec=never
==7949==     in use at exit: 0 bytes in 0 blocks
==7949== All heap blocks were freed -- no leaks are possible
==7949== ERROR SUMMARY: 0 errors from 0 contexts (suppressed: 0 from 0)

Documentation

  • [N/A] No configuration or public API changes

Backporting

  • Backport to the latest stable release

Fluent Bit is licensed under Apache 2.0, by submitting this pull request I understand that this code will be released under the terms of that license.

Roll deferred per-tag buffers at a 256 KiB target without splitting
records. Keep completed buffers ordered and hand them to core only from
the collector.

Drain accepted buffers during graceful shutdown, including after a pause.
Detach each in-flight buffer and prevent nested drains when backpressure
re-enters the pause callback. Keep ordinary new records subject to memory
admission while preserving the callback-owner shutdown allowance.

Fixes fluent#12519.

Signed-off-by: Michael Renner <terrorobe@github.com>
Exercise chunk bounds, full payload and group context, per-tag order,
retained originals, and chained rewrite_tag paths through the real engine.

Cover oversized records, partially filled core chunks, backpressure,
late shutdown deliveries and new-record admission. The grouped cases
require the prerequisite group-framing fix in fluent#12500.

Signed-off-by: Michael Renner <terrorobe@github.com>
@coderabbitai

coderabbitai Bot commented Oct 6, 2026

Copy link
Copy Markdown

Important

Draft PR not reviewed

Draft PRs are not automatically reviewed by default.

  • Trigger a manual review

To automatically review draft PRs, update your CodeRabbit configuration:

reviews:
  auto_review:
    drafts: true
  • Autopilot · Keep fixing CodeRabbit findings and required CI, and resolving merge conflicts

Thanks for using CodeRabbit! It's free for OSS, and your support helps us grow. If you like it, consider giving us a shout-out.

❤️ Share

Comment @coderabbitai help to get the list of available commands.

This branch had an error being deployed

1 failed deployment
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

in_emitter: rewrite_tag chunks can reach emitter_mem_buf_limit instead of the 2 MB target

1 participant