Repository navigation
fix(connectors): make source NACK retries configurable and observable - #4269
rohankumardubey wants to merge 12 commits into
Conversation
|
Thanks for the PR. It is labeled Slash commands (own line, regular comment) move it around the queue:
See CONTRIBUTING.md for details. |
Codecov Report✅ All modified and coverable lines are covered by tests. Additional details and impacted files@@ Coverage Diff @@
## master #4269 +/- ##
============================================
- Coverage 87.82% 87.67% -0.16%
- Complexity 1607 1608 +1
============================================
Files 1302 1300 -2
Lines 251971 240386 -11585
Branches 212264 200681 -11583
============================================
- Hits 221299 210751 -10548
+ Misses 25386 24641 -745
+ Partials 5286 4994 -292
🚀 New features to boost your workflow:
|
|
/ready |
mlevkov
left a comment
There was a problem hiding this comment.
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:
NackDispositioncannot 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
/healthat 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_sourcestill stops after five NACKs. - Item 2, keep polling across a transient NACK:
NoneandRetryboth 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
Errorafter a good send: already on master from #3957. The PR body lists it as a change here. - With
max_consecutive_nacksset, 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:44andinfluxdb_source/dependencies.md:46still say^0.5.0.http_source/src/lib.rs:542:max_consecutive_nacks = 0fails in the SDK's generic config parse, so the log never names the field.Option<u32>plus a check invalidate()gives a named error, asmax_batch_sizehas.sdk/src/source.rs:131:BatchPolicyis public now but has no doc, andMAX_CONSECUTIVE_NACKSandBATCH_RESULT_TIMEOUTstill read as fixed limits.- Optional simplifications: a
RegisterStopCallbackalias atruntime/src/main.rs:513andruntime/src/source.rs:880(theSourceApifield must stay inline, because dlopen2's derive rejects an alias insideOption).report_unexpected_stopcan shareset_error's two lines and their ordering comment. Thespawn_blockingclosure fits in oneifwithis_none_or. The SDK README says three times that a hook error stops polling.
|
/skill team-review-slim |
There was a problem hiding this comment.
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.
mlevkov
left a comment
There was a problem hiding this comment.
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::Shutdownnever reaches the runtime on a real stop, becauseclose()setsclosingfirst. Only a dropped watch sender would send it, as "(Shutdown); restart required".BatchPolicy::result_timeout()has no caller now, andBATCH_RESULT_TIMEOUTsays "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 theAlreadyErrorarm, or theselect!arms swapped back, passes every test in the SDK, runtime and http_source crates. TheNonestop-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.rsthe newNackLimitwait runs before the 500 ms no-PUT check, so that check can no longer fail.server.rs:1479duplicates the test #4302 added (:1594in this branch). - Docs:
lib.rs:1719still mentions the five-NACK budget.server.rs:995,:1362andmanagement.rs:458name 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:119still describes the drain-first exit.
|
@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? |
|
you can include this fix. also, you can bump the version to |
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
cargo clippy --all-features --all-targets -- -D warnings.cargo test -p iggy_connector_sdk --all-features,cargo test -p iggy-connectors, andcargo test -p iggy_connector_http_source.cargo build --bin iggy-server --bin iggy-connectors,cargo build -p iggy_connector_random_source, andcargo test -p integration -- connectors::runtime::http_state::given_conflict_mid_stream_should_nack_and_latch(1 passed).prek runon 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 runpassed locally. The author remains responsible for reviewing and explaining the final diff.