Skip to content

fix(connectors): bound source forwarding channel with backpressure - #3795

Open
mlevkov wants to merge 2 commits into
apache:masterfrom
mlevkov:bounded-source-channel
Open

fix(connectors): bound source forwarding channel with backpressure#3795
mlevkov wants to merge 2 commits into
apache:masterfrom
mlevkov:bounded-source-channel

Conversation

@mlevkov

@mlevkov mlevkov commented Aug 1, 2026

Copy link
Copy Markdown
Contributor

Summary

The channel between a source plugin's send callback and the runtime's forwarding loop was flume::unbounded(), so a slow or hung Iggy meant batches accumulated in memory without bound instead of propagating backpressure into the plugin's polling loop. This is the prerequisite runtime fix requested in the HTTP source discussion (#3039), and it applies to every source connector, including the source PRs currently in flight.

What changed

  • The forwarding channel is now a bounded crossfire channel (crossfire::mpsc::bounded_blocking_async), the same shape shard and server-ng use. flume is no longer a runtime dependency.
  • Capacity comes from a new optional SourceConfig field, channel_capacity, counted in batches (a single batch can be megabytes), defaulting to 1024 and clamped to [1, 65536] since crossfire eagerly allocates the ring and asserts capacity < 2^31. The existing ConfigEnv derive provides IGGY_CONNECTORS_SOURCE_<KEY>_CHANNEL_CAPACITY; configs without the field behave as before apart from the bound.
  • The FFI send callback applies backpressure with a try_send fast path and a send_timeout(10ms) retry loop that re-reads a per-instance shutdown flag between waits. The manager sets that flag before iggy_source_close so a hung Iggy cannot deadlock the close. Process shutdown sets every instance's flag (signal_shutdown_all) before the sequential stops, because instances loaded from one plugin library share a single tokio runtime and a wedged sibling would otherwise hold a worker an earlier close needs.
  • A full channel logs one warn! per backpressure episode (latched, cleared on genuine recovery). A batch that still cannot be enqueued after the stop signal is dropped and counted in iggy_connector_errors_total.
  • Unit tests pin the behavior shutdown relies on: buffered batches drain after the senders drop (crossfire's docs do not promise this), the retry loop unblocks when the flag flips mid-backoff, and shutdown drops are counted without being enqueued.
  • Docs updated across the runtime README, sources README, and the connector skills, which still described flume and the pre-feat(connectors): add agent docs, per-batch observability, atomic state #3321 shutdown order.

One correction to the discussion notes

@hubcio the spec assumed crossfire's blocking sender has no send_timeout. It does: blocking_tx.rs:288 on Tx, reachable from MTx via Deref. The loop is built on it instead of try_send plus sleep, so the sender wakes as soon as capacity frees while shutdown latency stays bounded by the retry interval.

Known limitation

Stopping a single connector via the runtime API while enough same-library sibling instances are saturated can delay that close until the siblings drain, because the callback parks a worker of the shared plugin runtime. The code comment and the connector skill document this. The complete fix is an SDK-side worker handoff (tokio::task::block_in_place around the callback invocation); happy to file it as a follow-up issue.

Test plan

  • cargo clippy -p iggy-connectors --all-targets -- -D warnings clean
  • cargo test -p iggy-connectors: 128 passed, including the new channel and shutdown tests
  • cargo build -p iggy_connector_stdout_sink -p iggy_connector_random_source
  • cargo test -p integration -- connectors::runtime:: could not run on this machine (hwlocality-sys needs pkg-config); relying on CI for the integration suite

@github-actions

github-actions Bot commented Aug 1, 2026

Copy link
Copy Markdown

Thanks for the PR. It is labeled S-waiting-on-review and queued for review.

Slash commands (own line, regular comment) move it around the queue:

  • /ready - back to S-waiting-on-review after addressing feedback
  • /author - flip to S-waiting-on-author while you finish changes
  • /request-review @user-or-team - request a reviewer

See CONTRIBUTING.md for details.

@github-actions github-actions Bot added the S-waiting-on-review PR is waiting on a reviewer label Aug 1, 2026
@codecov

codecov Bot commented Aug 1, 2026

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 93.23077% with 22 lines in your changes missing coverage. Please review.
✅ Project coverage is 20.17%. Comparing base (3e7bf7a) to head (679f6e4).
⚠️ Report is 1 commits behind head on master.

Files with missing lines Patch % Lines
core/connectors/runtime/src/source.rs 94.08% 18 Missing and 1 partial ⚠️
core/connectors/runtime/src/configs/connectors.rs 0.00% 1 Missing ⚠️
core/connectors/runtime/src/main.rs 0.00% 1 Missing ⚠️
core/connectors/runtime/src/manager/source.rs 50.00% 1 Missing ⚠️
Additional details and impacted files
@@              Coverage Diff              @@
##             master    #3795       +/-   ##
=============================================
- Coverage     82.84%   20.17%   -62.68%     
+ Complexity     1299     1291        -8     
=============================================
  Files          1199     1198        -1     
  Lines        161885   135243    -26642     
  Branches     131360   104844    -26516     
=============================================
- Hits         134120    27289   -106831     
- Misses        24225   107329    +83104     
+ Partials       3540      625     -2915     
Components Coverage Δ
Rust Core 2.39% <93.23%> (-80.97%) ⬇️
Java SDK 66.12% <ø> (+0.05%) ⬆️
C# SDK 62.85% <ø> (-13.11%) ⬇️
Python SDK 89.98% <ø> (ø)
PHP SDK 84.26% <ø> (ø)
Node SDK 94.85% <ø> (-1.37%) ⬇️
Go SDK 68.53% <ø> (ø)
Files with missing lines Coverage Δ
core/connectors/runtime/src/configs/connectors.rs 0.00% <0.00%> (-41.08%) ⬇️
core/connectors/runtime/src/main.rs 31.75% <0.00%> (-53.97%) ⬇️
core/connectors/runtime/src/manager/source.rs 70.00% <50.00%> (-23.47%) ⬇️
core/connectors/runtime/src/source.rs 37.48% <94.08%> (-40.21%) ⬇️

... and 680 files with indirect coverage changes

🚀 New features to boost your workflow:
  • ❄️ Test Analytics: Detect flaky tests, report on failures, and find test suite problems.
  • 📦 JS Bundle Analysis: Save yourself from yourself by tracking and limiting bundle sizes in JS merges.

@mlevkov

mlevkov commented Aug 2, 2026

Copy link
Copy Markdown
Contributor Author

/request-review @hubcio

@github-actions
github-actions Bot requested a review from hubcio August 2, 2026 02:21
mlevkov added a commit to mlevkov/iggy that referenced this pull request Aug 2, 2026
Iggy has no way to receive a webhook. Every provider that pushes events
over HTTP needs something in front of it, and today that means running a
separate service whose only job is to accept a POST and republish it.
This connector removes that hop: it runs an embedded HTTP server, accepts
authenticated POST bodies, and produces them to the instance's stream and
topic as raw bytes.

One plugin .so is loaded once no matter how many source entries reference
it, so the listener cannot live on any single instance. It lives in a
process-global registry keyed by listen address: the first open binds the
public and admin ports, later opens validate their body limit, admin
address, management token and instance name against the running listener
before joining, and the last close releases both ports. Mismatches fail
that instance's open rather than silently handing it a listener its
configuration does not describe. A single port can therefore serve many
providers, each routed to its own topic.

Requests resolve against an ArcSwap route table that is rebuilt whole on
every control-plane change, so one atomic load yields both the endpoint's
auth rules and the destination bridge. Secret paths carry 128 bits in the
URL itself, on the model of a Slack webhook, with optional bearer or HMAC
on top; HMAC is verified over the raw body in constant time. Revoked
endpoints answer 404 alongside paths that never existed, so a leaked URL
cannot be used to confirm it was once live.

Endpoints can be registered, re-keyed and revoked at runtime through a
token-guarded API on the admin listener, because revoking a compromised
endpoint is time-critical and provisioning one per tenant is inherently
programmatic. Those endpoints ride the SDK's ConnectorState, and state is
attached only to an empty batch: the runtime saves state solely on the
success branch of the Iggy send, and an empty send always succeeds, so a
mutation cannot be lost to an unrelated send failure. Revocation writes a
tombstone that outranks TOML on restore, so a stale config file cannot
resurrect an endpoint an operator revoked.

Delivery is best-effort in both directions and the README says so first,
before anything else: HTTP 200 means accepted into an in-memory buffer,
and both the loss and duplicate windows are enumerated with what mitigates
each. A full bridge answers 429 with Retry-After rather than blocking,
since holding the connection open would turn a slow Iggy into a retry
storm. Gateway metrics on the admin listener cover accept-to-200 latency,
which the runtime's own stage histograms begin too late to see.

Part of the webhook gateway design accepted in apache#3039. The backpressure
chain is only complete once the bounded runtime forwarding channel from
apache#3795 lands; until then a full bridge signals an arrival burst rather
than a slow Iggy, which the README documents.

Co-authored-by: Claude <noreply@anthropic.com>
@mlevkov
mlevkov force-pushed the bounded-source-channel branch from d7580d5 to bd19da2 Compare August 2, 2026 21:39
The channel between a source plugin's send callback and the runtime's
forwarding loop was flume::unbounded(), so a slow or hung Iggy meant
batches accumulated without bound instead of propagating backpressure
into the plugin's polling loop.

Swap it for a bounded crossfire channel (the shard and server-ng
standard), sized by an optional SourceConfig channel_capacity counted
in batches, defaulting to 1024. The FFI callback retries with
send_timeout while re-reading a shutdown flag, set by the manager
before iggy_source_close and for every instance ahead of the
sequential process-shutdown stops, since same-library instances share
one plugin runtime and a wedged sibling would otherwise hold a worker
an earlier close needs. A unit test pins that buffered batches drain
after the senders drop, which shutdown relies on and crossfire's docs
do not promise. This drops flume from the runtime.

Requested in the HTTP source discussion (apache#3039).

Co-authored-by: Claude <noreply@anthropic.com>
@mlevkov
mlevkov force-pushed the bounded-source-channel branch from bd19da2 to 399e46c Compare August 9, 2026 01:28
mlevkov added a commit to mlevkov/iggy that referenced this pull request Aug 9, 2026
Iggy has no way to receive a webhook. Every provider that pushes events
over HTTP needs something in front of it, and today that means running a
separate service whose only job is to accept a POST and republish it.
This connector removes that hop: it runs an embedded HTTP server, accepts
authenticated POST bodies, and produces them to the instance's stream and
topic as raw bytes.

One plugin .so is loaded once no matter how many source entries reference
it, so the listener cannot live on any single instance. It lives in a
process-global registry keyed by listen address: the first open binds the
public and admin ports, later opens validate their body limit, admin
address, management token and instance name against the running listener
before joining, and the last close releases both ports. Mismatches fail
that instance's open rather than silently handing it a listener its
configuration does not describe. A single port can therefore serve many
providers, each routed to its own topic.

Requests resolve against an ArcSwap route table that is rebuilt whole on
every control-plane change, so one atomic load yields both the endpoint's
auth rules and the destination bridge. Secret paths carry 128 bits in the
URL itself, on the model of a Slack webhook, with optional bearer or HMAC
on top; HMAC is verified over the raw body in constant time. Revoked
endpoints answer 404 alongside paths that never existed, so a leaked URL
cannot be used to confirm it was once live.

Endpoints can be registered, re-keyed and revoked at runtime through a
token-guarded API on the admin listener, because revoking a compromised
endpoint is time-critical and provisioning one per tenant is inherently
programmatic. Those endpoints ride the SDK's ConnectorState, and state is
attached only to an empty batch: the runtime saves state solely on the
success branch of the Iggy send, and an empty send always succeeds, so a
mutation cannot be lost to an unrelated send failure. Revocation writes a
tombstone that outranks TOML on restore, so a stale config file cannot
resurrect an endpoint an operator revoked.

Delivery is best-effort in both directions and the README says so first,
before anything else: HTTP 200 means accepted into an in-memory buffer,
and both the loss and duplicate windows are enumerated with what mitigates
each. A full bridge answers 429 with Retry-After rather than blocking,
since holding the connection open would turn a slow Iggy into a retry
storm. Gateway metrics on the admin listener cover accept-to-200 latency,
which the runtime's own stage histograms begin too late to see.

Part of the webhook gateway design accepted in apache#3039. The backpressure
chain is only complete once the bounded runtime forwarding channel from
apache#3795 lands; until then a full bridge signals an arrival burst rather
than a slow Iggy, which the README documents.

Co-authored-by: Claude <noreply@anthropic.com>

@hubcio hubcio left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

a few things that don't fit a diff line:

  • manager/source.rs:164 - the 5s task-await is per handle (there are two), fixed, and channel_capacity now scales the post-cleanup_sender drain against it (buffered batches still drain after the sender drops, up to capacity send+save iterations). abort mid-drain is a clean tail truncation since the state save is atomic, so it costs lost position, not corruption - but the interaction deserves a line in the readme: raising capacity lengthens shutdown.
  • source.rs:627 - iggy_source_handle's i32 is discarded too, same family as the close-code comment; only reachable in a start-then-stop race today, so hygiene.
  • elasticsearch_source with [state] enabled = true persists its own cursor at close and overrides the runtime state at open, so a runtime-side latch can't cover it. opt-in and off by default; follow-up issue.
  • the pre-existing producer.send() err path has the same cursor-supersession shape (skip save, continue, next batch persists) - the latch here doesn't close that one; separate follow-up.
  • worth a regression test once the latch lands: saturate the channel, stop the connector, restart, assert no row gap - state_persists_across_connector_restart in the postgres integration suite is a natural template.
  • spawn_source_handler is at 11 params; passing &SourceConfig instead needs the resolved-path/version split on SourceConnectorPlugin first, so follow-up sized.

messages: ProducedMessages,
) {
let mut messages = match sender.try_send(messages) {
Ok(()) => {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

the Ok(()) arm never reads shutdown, which makes the shutdown drop below non-terminal: drop batch N, the forwarding loop frees a slot, batch N+1 enqueues, ships, and persists its state - the saved cursor now covers N. sources advance the cursor at poll time and snapshot it into every batch (postgres tracking_offsets, same shape in the other sources), so a restart resumes past the hole. silent mid-stream loss, not tail truncation.

the same supersession already exists on master via the producer.send() err path (skips the save, continues), but that one flips the connector to Error status; this one only logs + counts, and fires on routine paths - SIGTERM (signal_shutdown_all arms every source for the whole sequential-stop window) and connector restart via the api.

fix is a per-instance dropped latch: after the first drop, drop every later batch too, so the persisted cursor can never pass the gap (ring is fifo, buffered older batches still flush). use the same flag to latch the error! in drop_during_shutdown, otherwise after latching it fires once per poll; keep the counter bumping per batch.

return;
}
};
if shutdown.load(Ordering::Acquire) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

zero grace: the first Full after the flag drops immediately, while the forwarding loop keeps draining until cleanup_sender. one bounded send_timeout round before the first drop (skipped once latched) lets most in-flight batches land - the parked sender is woken on drain, so the wait is the drain period, not the full timeout. anything larger only pays off after iggy_source_close moves to spawn_blocking, since it currently blocks a host worker for the duration.

}
}

fn drop_during_shutdown(plugin_id: u32, message_count: usize, error_counter: &Counter) {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

the log reports a message count but the counter moves by 1 per batch, into the same iggy_connector_errors_total series as decode/send/save failures. this is the only signal for a permanently lost batch - a dedicated counter (with inc_by(message_count) where the count exists) keeps loss distinguishable from ordinary errors.

}
}

// Parks a worker thread of the plugin library's shared tokio runtime - a

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this covers the close-delay case but not steady state: the park removes a worker from the plugin library's shared runtime (one runtime per .so, workers = available_parallelism), so saturated instances degrade every sibling from the same .so, and a 1-2 vcpu container can wedge on one instance. worth a sentence here and in the skill doc. also, a runtime-side tokio::task::block_in_place would be a silent no-op on these threads (they belong to the plugin's tokio, not ours), so the handoff fix really does have to be sdk-side.


#[test]
fn given_backoff_in_progress_when_shutdown_signaled_should_unblock_and_drop() {
let (sender, _receiver) = bounded_channel(1);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this test and given_full_channel_when_receiver_frees_capacity_should_deliver_batch are the only ones that enter the backoff loop, and both run capacity 1, which crossfire routes to a different queue impl (OneMpsc) and a different backoff regime (large = capacity >= 10 gates yield-vs-spin). prod default 1024 is ArrayMpsc with large = true, so the park/wake path being pinned isn't the shipped one. capacity >= 10 here fixes both. the capacity-4 drain test is fine - try_send/recv never enter the backoff.

}

#[test]
fn given_registered_entry_when_signal_shutdown_called_should_set_flag() {

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this can pass even if signal_shutdown were a no-op: the sibling test's signal_shutdown_all() sets every entry's flag, including this one, and id partitioning doesn't help against a whole-map op. one merged test with two entries - target set, sibling not set, then signal_shutdown_all sets both - is hermetic. only plain cargo test is affected (nextest is process-per-test), but that's the documented local flow.

source::signal_shutdown(plugin_id);
if let Some(container) = &container {
info!("Closing source connector with ID: {plugin_id} for plugin: {key}");
(container.iggy_source_close)(plugin_id);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

return code dropped, and the new ordering's safety argument leans on close actually stopping callbacks - the sdk returns -1 when it can't join the polling task, and then callbacks can outlive the close and land on the removed-entry branch that drops with no metric. init checks the same call (if close_result != 0 { warn! }) - mirror it here.


Each source configuration accepts an optional `channel_capacity` setting that bounds the channel between the plugin's send callback and the runtime's forwarding loop. Capacity is counted in batches (one `poll()` result each, potentially megabytes), not messages or bytes. The default is 1024 batches.

When the channel is full (Iggy accepts messages more slowly than the plugin produces them), the send callback backs off and retries instead of buffering without bound, so backpressure propagates into the plugin's polling loop. During shutdown, a batch that still cannot be enqueued after the stop signal is dropped and counted in `iggy_connector_errors_total`, so a saturated source may report errors at SIGTERM. Values outside `[1, 65536]` are clamped with a warning.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

this undersells the failure: the drop isn't just an error count - a later batch can still enqueue and persist its state, moving the saved position past the dropped batch, so the data is silently gone (see the send-callback comment). it also happens on connector restart via the api, not only at SIGTERM. once the fix lands this should state the guarantee: replay from the last delivered batch's state - duplicates for offset-mode sources, permanent loss for delete_after_read / processed_column ones. same wording lives in the connector-runtime skill doc.

channel_capacity = 1024
```

Environment override: `IGGY_CONNECTORS_SOURCE_<KEY>_CHANNEL_CAPACITY`.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

only the local config provider wires env overrides - with config_type = "http" this var is silently ignored, and the unknown-var warning is suppressed for the IGGY_CONNECTORS_SOURCE_ prefix, so nothing surfaces it. qualify with "local config provider only".

@@ -72,6 +73,7 @@ path = "libiggy_connector_random_source" # Path to the source connector
config_format = "toml"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

pre-existing, two lines from your change: the example key is config_format but the real field is plugin_config_format - no serde alias and no deny_unknown_fields, so the documented key is silently ignored.

@github-actions github-actions Bot added S-waiting-on-author PR is waiting on author response and removed S-waiting-on-review PR is waiting on a reviewer labels Aug 13, 2026
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

S-waiting-on-author PR is waiting on author response

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants