Skip to content

fix(connectors): make source NACK retries configurable and observable - #4269

Open
rohankumardubey wants to merge 10 commits into
apache:masterfrom
rohankumardubey:fix/connectors-source-nack-policy
Open

rohankumardubey wants to merge 10 commits into
apache:masterfrom
rohankumardubey:fix/connectors-source-nack-policy

Conversation

@rohankumardubey

@rohankumardubey rohankumardubey commented Sep 22, 2026 •

Copy link
Copy Markdown
Contributor

Which issue does this PR address?

Closes #3941

Rationale

Five consecutive NACKs could permanently stop a source after a short broker outage. For sources holding accepted input only in memory, stopping can lead to data loss, and the runtime did not reliably report that the poll task had ended.

What changed?

The HTTP source now disables the consecutive-NACK stop by default, while other sources retain the existing five-NACK limit and can configure it through batch_policy(). When a batch result is missing, the SDK warns after 30 seconds and keeps waiting instead of NACKing and replaying a batch that may still be queued in the runtime.

The runtime records a stable reason when a source stops unexpectedly, updates its status and running gauge, and drains queued batches on normal close. The SDK version is now 0.5.1-edge.2; tests and documentation cover the stop, timeout, and health behavior.

Local Execution

  • Passed: cargo clippy --all-features --all-targets -- -D warnings.
  • Passed: cargo test -p iggy_connector_sdk --all-features, cargo test -p iggy-connectors, and cargo test -p iggy_connector_http_source.
  • Passed: cargo build --bin iggy-server --bin iggy-connectors, cargo build -p iggy_connector_random_source, and cargo test -p integration -- connectors::runtime::http_state::given_conflict_mid_stream_should_nack_and_latch (1 passed).
  • Pre-commit hooks passed via prek run on the staged changes, including Markdown lint, license headers, typos, Taplo, cargo fmt, and cargo sort. The pre-push hook has not run because this branch has not been pushed since these changes.

AI Usage

Codex (ChatGPT) assisted with implementation, tests, and documentation. The diff was reviewed against the maintainer comments; format, workspace Clippy, connector tests, the targeted integration test, and prek run passed locally. The author remains responsible for reviewing and explaining the final diff.

@github-actions

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
  • /pin - exempt the PR from the stale bot, /unpin to undo

See CONTRIBUTING.md for details.

@github-actions github-actions Bot added the S-waiting-on-review PR is waiting on a reviewer label Sep 22, 2026
@rohankumardubey rohankumardubey changed the title Fix/connectors source nack policy fix(connectors): make source NACK retries configurable and observable Sep 22, 2026
@codecov

codecov Bot commented Sep 22, 2026 •

Copy link
Copy Markdown

Codecov Report

❌ Patch coverage is 96.55870% with 34 lines in your changes missing coverage. Please review.
✅ Project coverage is 68.95%. Comparing base (2d7fddb) to head (76a7d54).

Files with missing lines Patch % Lines
core/connectors/sdk/src/source.rs 94.85% 19 Missing and 2 partials ⚠️
core/connectors/runtime/src/source.rs 97.31% 8 Missing and 1 partial ⚠️
...e/connectors/sources/http_source/src/management.rs 75.00% 3 Missing ⚠️
core/connectors/runtime/src/manager/source.rs 98.83% 1 Missing ⚠️
Additional details and impacted files
@@              Coverage Diff              @@
##             master    #4269       +/-   ##
=============================================
- Coverage     87.96%   68.95%   -19.01%     
  Complexity     1579     1579               
=============================================
  Files          1290     1288        -2     
  Lines        230145   187749    -42396     
  Branches     193480   151086    -42394     
=============================================
- Hits         202440   129467    -72973     
- Misses        22980    53436    +30456     
- Partials       4725     4846      +121     
Components Coverage Δ
Rust Core 65.36% <96.55%> (-23.71%) ⬇️
Java SDK 68.74% <ø> (ø)
C# SDK 77.67% <ø> (+0.05%) ⬆️
Python SDK 91.24% <ø> (ø)
PHP SDK 85.67% <ø> (ø)
Node SDK 96.56% <ø> (-0.04%) ⬇️
Go SDK 70.26% <ø> (+0.07%) ⬆️
Files with missing lines Coverage Δ
core/connectors/runtime/src/main.rs 87.24% <100.00%> (-0.59%) ⬇️
core/connectors/sdk/src/lib.rs 85.11% <100.00%> (+0.17%) ⬆️
core/connectors/sources/http_source/src/lib.rs 97.55% <100.00%> (+0.15%) ⬆️
core/connectors/sources/http_source/src/server.rs 95.05% <100.00%> (+0.04%) ⬆️
core/connectors/runtime/src/manager/source.rs 94.07% <98.83%> (+0.92%) ⬆️
...e/connectors/sources/http_source/src/management.rs 94.64% <75.00%> (-0.28%) ⬇️
core/connectors/runtime/src/source.rs 88.51% <97.31%> (+0.81%) ⬆️
core/connectors/sdk/src/source.rs 93.65% <94.85%> (+3.31%) ⬆️

... and 424 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.

@rohankumardubey

Copy link
Copy Markdown
Contributor Author

/ready

@mlevkov mlevkov 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.

Thanks for picking up #3941. I traced every way the poll task can end: each one reports once, a user stop or restart never reports, and the running gauge moves once. The new export also loads both ways, old plugin on the new runtime and new plugin on an old runtime.

The main points, with details inline:

  • NackDisposition cannot do what its docs say, because no NACK reason reaches the source. I would drop it for now.
  • With the breaker off by default, a batch that always fails keeps http_source's /health at 200.
  • Seven changes to the new code keep every unit test in these crates green. I ran them as mutations.

Where #3941 stands with this PR:

  • Item 1, policy per connector: done for sources that override batch_policy(). Only http_source does. postgres_source still stops after five NACKs.
  • Item 2, keep polling across a transient NACK: None and Retry both keep polling. Neither can pick out a transient NACK, because no reason reaches the source.
  • Item 3, the stop is visible: done for plugins built with this SDK.
  • Item 4, clear Error after a good send: already on master from #3957. The PR body lists it as a change here.
  • With max_consecutive_nacks set, steps 6 and 7 still happen. After the stop the instance keeps answering 200 until the bridge fills. Fine as a follow-up.

Small things:

  • sdk/Cargo.toml:20: no release tag has SDK 0.5.0 yet, and the changes here are additive, so 0.6.0 may not be needed. If it goes back to 0.5.0, the "SDK 0.6" text in both READMEs changes too. If it stays, influxdb_sink/dependencies.md:44 and influxdb_source/dependencies.md:46 still say ^0.5.0.
  • http_source/src/lib.rs:542: max_consecutive_nacks = 0 fails in the SDK's generic config parse, so the log never names the field. Option<u32> plus a check in validate() gives a named error, as max_batch_size has.
  • sdk/src/source.rs:131: BatchPolicy is public now but has no doc, and MAX_CONSECUTIVE_NACKS and BATCH_RESULT_TIMEOUT still read as fixed limits.
  • Optional simplifications: a RegisterStopCallback alias at runtime/src/main.rs:513 and runtime/src/source.rs:880 (the SourceApi field must stay inline, because dlopen2's derive rejects an alias inside Option). report_unexpected_stop can share set_error's two lines and their ordering comment. The spawn_blocking closure fits in one if with is_none_or. The SDK README says three times that a hook error stops polling.

Comment thread core/connectors/sdk/src/lib.rs Outdated
Comment thread core/connectors/sources/http_source/src/lib.rs
Comment thread core/connectors/sdk/src/source.rs Outdated
Comment thread core/connectors/sdk/src/source.rs Outdated
Comment thread core/connectors/sdk/README.md Outdated
Comment thread core/connectors/sdk/src/source.rs
Comment thread core/connectors/runtime/src/source.rs Outdated
Comment thread core/connectors/sources/http_source/README.md
@hubcio

hubcio commented Sep 28, 2026

Copy link
Copy Markdown
Contributor

/skill team-review-slim

@github-actions github-actions Bot left a comment

Copy link
Copy Markdown

Choose a reason for hiding this comment

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

Summary: The stop-notification path and the configurable NACK policy hold on every traced path, but the PR leaves HTTP source docs describing the removed five-NACK stop as unobservable, counts a NACK-limit stop twice in the error metric, and leaves two dependency tables pinning the SDK at 0.5.0.

Counts: critical 0, warning 1, nit 3, simplification 1


This review was generated by Claude Code 2.1.284 on deepseek-flash[1m]. Review the output before you act on it.

Comment thread core/connectors/sources/http_source/README.md
Comment thread core/connectors/runtime/src/source.rs Outdated
Comment thread core/connectors/sdk/Cargo.toml
Comment thread core/connectors/sdk/src/source.rs
Comment thread core/connectors/runtime/src/source.rs Outdated
@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 Sep 28, 2026
@rohankumardubey

Copy link
Copy Markdown
Contributor Author

@mlevkov I’ve pushed the changes addressing the code feedback. I removed the configurable result timeout; the remaining late-result/replay risk is for #3981. Could you take another look when you have time?

@mlevkov mlevkov 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.

Thanks, I re-checked at de05b6ace. The fixes hold: NackDisposition and the timeout knob are gone, stops carry a reason, the double count is gone, and old plugins get a warning. Every mutation from my first review that still applies now fails a test. The stop path and the FFI in both directions hold up.

Three things inline: how a stuck batch shows up in health and logs, a gap in the readiness test, and loss window 4.

The PR body still describes the first version. It lists a per-source batch-result timeout and a transient-NACK option, which are gone. It also credits this PR with clearing a stale error after a good send, which is #3957. It does not mention the stop reason, the readiness rule, max_consecutive_nacks or the SDK 0.6.0 bump.

Closing #3941 is fine with me. Proposal 2 is dropped on purpose, since no NACK cause reaches the source. Steps 6 and 7 now happen only with a configured limit or a panic, which is what inline 3 asks the README to say. The timeout half lives in #3981, where I added a note.

Smaller points:

  • Before 0.6.0 ships: SourceStopReason::Shutdown never reaches the runtime on a real stop, because close() sets closing first. Only a dropped watch sender would send it, as "(Shutdown); restart required". BatchPolicy::result_timeout() has no caller now, and BATCH_RESULT_TIMEOUT says "Default" though nothing outside the SDK can change it.
  • Loop tests: no unit test drives the forwarding loop through a stop that finds the status already Error, or with a batch still queued. An extra error increment in the AlreadyError arm, or the select! arms swapped back, passes every test in the SDK, runtime and http_source crates. The None stop-export path and reason bytes 0, 2, 3 and 6 are unpinned too. A table of literal bytes would also catch a renumbered enum.
  • Other tests: in http_state.rs the new NackLimit wait runs before the 500 ms no-PUT check, so that check can no longer fail. server.rs:1479 duplicates the test #4302 added (:1594 in this branch).
  • Docs: lib.rs:1719 still mentions the five-NACK budget. server.rs:995, :1362 and management.rs:458 name a result-hook failure, which this source cannot produce. With a limit set, refused state flushes count too, so an idle gateway can stop (lib.rs:1189, README line 131). README line 306, which I wrote, says the bridge drains after a timeout. Replay comes first, so with no limit the same batch is queued again after every timeout. .claude/skills/connector-runtime/SKILL.md:119 still describes the drain-first exit.

Comment thread core/connectors/sources/http_source/src/lib.rs
Comment thread core/connectors/sources/http_source/src/lib.rs Outdated
Comment thread core/connectors/sources/http_source/README.md Outdated
}

fn batch_policy(&self) -> source::BatchPolicy {
source::BatchPolicy::default().with_max_consecutive_nacks(

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.

with no limit by default, the old 5-copy cap on timeout replays is gone. When the broker accepts the connection but never answers, the SDK timeout arm (sdk/src/source.rs:497) fakes a NACK, the staged batch is replayed, and another full copy lands in the unbounded flume channel (runtime/src/source.rs:894) about every 35s. Each copy can be ~500 MiB at defaults, held in the runtime every connector shares, and on recovery every copy is sent again with status flapping Error/Running.

On timeout, warn and keep waiting on result_receiver instead of faking a NACK. That keeps 429 backpressure, and /health still turns 503 after 60s. bounded(1) is not a fix: the blocking send under the DashMap guard (runtime/src/source.rs:1157) hangs close().

loop {
let produced_batch = tokio::select! {
biased;
stopped = stop_receiver.recv() => {

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 arm breaks on None too. stop_connector runs close and then cleanup_sender (manager/source.rs:196), which drops stop_sender, so the loop exits before flume yields batches already queued. Before this PR, while let Ok(..) = recv_async() drained them to Iggy. For http_source that batch was already NACKed and close() dropped the staged copy, so a webhook answered 200 is lost.

Break only on Some(reason), disable the arm on None and keep draining until Disconnected, plus a test that a queued batch still gets its batch_result on stop.

[package]
name = "iggy_connector_sdk"
version = "0.5.1-edge.1"
version = "0.6.0"

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.

0.5.1-edge.1 -> 0.6.0 plus the new iggy_source_register_stop_callback export (src/source.rs:801) needs explicit maintainer sign-off per CLAUDE.md. The bump also looks unneeded: the API change is additive and batch_policy is a defaulted method (README:40 says no code change). A bare 0.6.0 also breaks the -edge convention from #4315.

Use 0.5.1-edge.2 (or scripts/bump-version.sh rust-connector-sdk --minor --edge), then fix the "SDK 0.6" text and the influxdb ^0.6.0 pins, which won't match a prerelease.

if source
.last_error
.as_ref()
.is_some_and(|error| error.message.contains("NackLimit"))

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 changed, but the PR says the integration harness never reached the connector tests and prek run was not run. Please run both locally and post the result. Also, this pins the {reason:?} Debug text that runtime/src/source.rs:849 puts into REST last_error.

cleanup_sender(plugin_id);
if let Some(reason) = stop_reason {
let error_msg = format!(
"Source polling stopped for connector with ID: {plugin_id} ({reason:?}); restart required"

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.

nit: {reason:?} exposes Debug output through the status API. An impl Display for SourceStopReason gives operators stable wording, and the test above can match it.

}
SourceBatchResult::Nack => consecutive_nacks.fetch_add(1, Ordering::Relaxed) + 1,
SourceBatchResult::Nack => {
let _ = consecutive_nacks.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |count| {

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.

nit: fetch_update and then a separate load is two atomics. fetch_update(..).map_or(u32::MAX, |prev| prev.saturating_add(1)) gives the new count in one.

pub type BatchResultCallback = extern "C" fn(plugin_id: u32, batch_id: u64, result: u8) -> i32;
pub type SourceStoppedCallback = extern "C" fn(plugin_id: u32, reason: u8);

/// Reason a source poll task ended, passed to the runtime stop callback.

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.

nit: not every variant goes through the stop callback. Shutdown is filtered at :92, and RegistrationFailed/HandlerFailed are emitted only by the runtime (runtime/src/source.rs:946). document who emits each one.

/// Invoked when the source is initialized, allowing it to perform any necessary setup.
async fn open(&mut self) -> Result<(), Error>;

/// Controls whether repeated NACKs stop this source.

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.

nit: say that this is read once, right after open() (src/source.rs:264), so returning a different policy later has no effect.

The batch is then replayed on every poll and the SDK stops the poll task after five consecutive NACKs, while the listener keeps answering 200. The connector cannot guard against this itself: `schema` lives under `[[streams]]` and the plugin only ever receives `[plugin_config]`, so nothing in `open()` can see it.
The batch is then replayed on every poll. This source disables the SDK's consecutive-NACK breaker by default because accepted webhooks exist only in its in-memory bridge. Repeated failures back off to a five-second retry delay, and the listener answers 429 once the bridge fills.

Empty bodies are rejected with 400. `max_body_size_bytes` cannot exceed Iggy's 64,000,000-byte payload cap; requests above the configured limit are rejected with 413. A staged batch that keeps failing for more than 60 seconds makes `/health` answer 503 while retries continue. `/admin/health` still answers 200, but reports `"status":"degraded"` and marks that instance's `staged_batch_is_stuck` true while `poll_is_live` remains true.

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.

nit: the stuck-batch and /admin/health behavior is buried in the schema = raw section next to body-size facts repeated elsewhere. move it near the Observability list, which also doesn't mention staged_batch_is_stuck.

2. Plugin polls + invokes `send_callback(plugin_id, ptr, len)`.
3. Callback runs in the SDK macro's spawned async task. Pushes postcard `ProducedMessages` into a `flume` channel keyed by `plugin_id` in `SOURCE_SENDERS: Lazy<DashMap<u32, SourceSenderEntry>>` (`pub(crate)`). `SourceSenderEntry` wraps the sender + a pre-extracted owned `Counter` (the `errors` series, `Arc<AtomicU64>` inside). The FFI callback bumps errors on deserialize or channel-closed failure with one relaxed atomic - no `Family` lookup, no `Arc<Metrics>` handle.
4. `source_forwarding_loop` pulls from the channel, deserializes, applies transforms, encodes via `StreamEncoder`, sends to Iggy producer.
4. `source_forwarding_loop` pulls from the channel, deserializes, applies transforms, encodes via `StreamEncoder`, sends to Iggy producer. If the SDK times out awaiting a result, it NACKs and polls again. A source that retains its NACKed batch, such as `http_source`, replays it before draining fresh input; repeated timeouts can queue copies of that batch in this unbounded channel.

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.

nit: only step 4 changed. the loop steps still leave out register-stop-callback (runtime/src/source.rs:932-956), the stop_receiver exit, report_unexpected_stop and the loop-end cleanup.

@rohankumardubey

Copy link
Copy Markdown
Contributor Author

@numinnex Thanks for the detailed review. I’ll address the queued-batch stop behavior and the test/docs feedback. Before changing the timeout path, should that fix be included in #4269 or remain in #3981? Also, do you prefer 0.5.1-edge.2 for the SDK, and is the additive stop-callback export acceptable?

@hubcio

hubcio commented Oct 5, 2026

Copy link
Copy Markdown
Contributor

you can include this fix. also, you can bump the version to edge.2. yes, its acceptable

This branch has not been deployed

No deployments
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.

Source NACK breaker stops non-replayable sources, and the stop is unobservable

4 participants