mlevkov commented on code in PR #4269:
URL: https://github.com/apache/iggy/pull/4269#discussion_r4109574420
##########
core/connectors/sdk/src/lib.rs:
##########
@@ -109,6 +109,24 @@ pub trait Source: Send + Sync {
/// Invoked when the source is initialized, allowing it to perform any
necessary setup.
async fn open(&mut self) -> Result<(), Error>;
+ /// Controls how long this source waits for batch results and whether
repeated NACKs stop it.
+ ///
+ /// The default stops after five consecutive NACKs. Sources that hold
accepted input only in
+ /// memory can disable that limit so a transient broker outage does not
strand their input.
+ fn batch_policy(&self) -> source::BatchPolicy {
+ source::BatchPolicy::default()
+ }
+
+ /// Decides whether a successfully handled NACK should trip the configured
stop limit.
+ ///
+ /// `Retry` keeps polling with capped backoff. This is useful when the
source can safely
+ /// replay a rejected batch and knows the failure is transient. The
default applies the
+ /// policy's limit. An error from `on_batch_result` is always terminal:
the source may not
+ /// have rolled back its staged work safely.
+ fn nack_disposition(&self) -> source::NackDisposition {
Review Comment:
`nack_disposition()` gets no NACK reason. `on_batch_result` only sees `Ack`
or `Nack`, and the runtime sends the result as one byte
(`runtime/src/source.rs:791`). So a source cannot tell a broker restart from
`schema = "json"`. Returning `Retry` every time does the same as
`with_max_consecutive_nacks(None)`.
`Retry` NACKs also still count toward the limit (`source.rs:674`). After
four or more `Retry` NACKs in a row, the first `ApplyPolicy` NACK stops the
source at once. The enum doc at `source.rs:89` and README line 36 read as if
they do not count.
I would drop `NackDisposition` and this method from the PR. No plugin
overrides it and it has not shipped, so removing it now breaks no plugin.
##########
core/connectors/sdk/src/source.rs:
##########
@@ -101,11 +128,33 @@ enum BatchCompletion {
}
#[derive(Debug, Clone, Copy)]
-struct BatchPolicy {
+pub struct BatchPolicy {
result_timeout: Duration,
nack_retry_delay: Duration,
max_nack_retry_delay: Duration,
- max_consecutive_nacks: u32,
+ max_consecutive_nacks: Option<NonZeroU32>,
+}
+
+impl BatchPolicy {
+ pub fn result_timeout(&self) -> Duration {
+ self.result_timeout
+ }
+
+ pub fn max_consecutive_nacks(&self) -> Option<NonZeroU32> {
+ self.max_consecutive_nacks
+ }
+
+ /// Changes the result timeout for this source instance.
+ pub fn with_result_timeout(mut self, result_timeout: Duration) -> Self {
Review Comment:
A result timeout NACKs a batch that the runtime may still be sending. The
first send can still land, so the replay is a duplicate. The late result for
the first batch then gets -1, and the runtime sets `Error`. That is #3981. A
shorter timeout makes it happen more often, and nothing calls this yet. Suggest
leaving it out until #3981 is fixed, or accepting only values at or above the
30 s default, so it cannot make this worse.
Related, now that http_source has no limit: during a long Iggy stall, each
30 s timeout adds one more copy of the batch to the runtime's channel. That
channel is unbounded (`runtime/src/source.rs:884`). Before this PR the breaker
stopped that at five copies.
##########
core/connectors/sdk/src/source.rs:
##########
@@ -50,6 +51,23 @@ pub type SendCallback = extern "C" fn(
) -> i32;
pub type BatchResultCallback = extern "C" fn(plugin_id: u32, batch_id: u64,
result: u8) -> i32;
+pub type SourceStoppedCallback = extern "C" fn(plugin_id: u32);
Review Comment:
The callback carries only the plugin id, so `last_error` always ends up as
"Source polling stopped unexpectedly ...; restart required", whatever the
cause. The status API cannot tell a NACK limit from an `on_batch_result` error,
a panic or a refused start. The export is new, so a `u8` reason, like the
result in `BatchResultCallback`, costs nothing now. Later it needs a second
symbol.
##########
core/connectors/runtime/src/source.rs:
##########
@@ -588,7 +597,19 @@ pub(crate) async fn source_forwarding_loop(
topic: producer.topic().to_string(),
};
- while let Ok(produced_batch) = receiver.recv_async().await {
+ let mut stopped_unexpectedly = false;
+ loop {
+ let produced_batch = tokio::select! {
+ biased;
+ batch = receiver.recv_async() => batch,
+ stopped = stop_receiver.recv() => {
+ stopped_unexpectedly = stopped.is_some();
Review Comment:
No unit test runs `source_forwarding_loop` or `spawn_source_handler`.
Setting this to `false` keeps all unit tests green, and so does removing either
`source_stopped(plugin_id)` call at lines 901 and 905. I ran all three as
mutations, and no integration test checks the outcome either. A
`#[tokio::test]` can call `spawn_source_handler` with stub `extern "C"`
functions that return -1, then check for `Error` and "restart required". The
api tests already build a `RuntimeContext` without a server (`api/mod.rs:369`).
For the stop arm end to end, `given_conflict_mid_stream_should_nack_and_latch`
(`core/integration/tests/connectors/runtime/http_state.rs:325`) latches the
state store, so every batch NACKs, but it returns before the fifth NACK.
Waiting there for "restart required" in `last_error` would cover it.
##########
core/connectors/sdk/src/source.rs:
##########
@@ -280,7 +353,7 @@ impl<T: Source + std::fmt::Debug + 'static>
SourceContainer<T> {
batch_id,
result,
self.id,
- MAX_CONSECUTIVE_NACKS,
+ self.batch_policy.max_consecutive_nacks,
Review Comment:
No unit test calls `complete_batch`, and no integration test sends it enough
NACKs for the policy to matter. I reverted this line to the default limit, and
all 684 tests in the SDK, runtime and http_source crates still passed. The same
happens when `handle()` goes back to `BatchPolicy::default()` (line 306).
`given_nack_limit_when_polling_stops_should_notify_runtime` uses the default
limit, and 1.5 s of backoff fits in its 2 s wait. Deleting the `closing` store
in `close()` (line 262) also passes, and so does dropping the stop guard before
the task spawns, which fires it before the first poll.
Suggestions:
- A plain `#[test]` that opens a container with
`with_max_consecutive_nacks(None)`, seeds a pending batch before each call, and
checks that six NACKs through `complete_batch` each return 0.
- A limit of 2 and a poll count in the notify test.
- A test that closes a running container with a callback registered and
checks that the callback never fires.
##########
core/connectors/sources/http_source/src/lib.rs:
##########
@@ -965,6 +969,11 @@ impl Source for HttpSource {
opened
}
+ fn batch_policy(&self) -> source::BatchPolicy {
Review Comment:
With no limit, a batch that the runtime rejects every time replays forever,
and `/health` stays 200. Each replay refreshes `last_poll_at`, so `is_ready`
(`server.rs:334`) stays true. Before this PR the stop took `/health` to 503
after 60 s, and README line 205 says that is what a load balancer should watch.
`schema = "json"` gets there, and so does a single empty POST body with
default limits. A body over Iggy's 64,000,000-byte payload cap does the same
when `max_body_size_bytes` is set above it, and that option allows up to
67,108,864. For both bodies the runtime cannot build the message, so it NACKs
the whole batch.
Suggest failing readiness once the staged batch has kept failing for longer
than the liveness window. `enqueue` also needs to refuse empty and oversize
bodies. That part is my code, so I can send it as a separate PR.
##########
core/connectors/sources/http_source/README.md:
##########
@@ -107,7 +107,11 @@ auth_secret = "partner-token"
The stream's `schema` **must be `raw`**. This connector always produces raw
bodies; with `schema = "json"` the runtime's encoder rejects every message,
which counts as a processing error and becomes a NACK rather than a drop.
-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.
Review Comment:
Loss window 4 at line 16 still says the SDK stops this source after five
NACKs. It also says the stop cannot be seen and that #3941 tracks it. This PR
makes the first two false. The same five-NACK premise is in `server.rs:995`,
`server.rs:1359`, `management.rs:458`, `lib.rs:440`, `lib.rs:1166`,
`lib.rs:1675` and the assert message at `lib.rs:2474`, and in
`.claude/skills/connector-source/SKILL.md` at lines 62-65 and 103-109. The
skills do not mention `batch_policy()` or the stop callback yet.
##########
core/connectors/sdk/README.md:
##########
@@ -29,11 +29,17 @@ The crash behavior is intentionally at-least-once:
An ACK follows Iggy's quorum confirmation. The topic's `durability` policy
decides whether that confirmation also waits for stable storage on the quorum.
Both policies normally write messages to disk.
-Source-side ACK work should be idempotent because process termination can
interrupt it. NACK handling must discard staged cursor changes and staged
delete or mark operations so polling can redeliver the batch. The SDK retries
NACKed batches with capped exponential backoff and stops after repeated NACKs.
+Source-side ACK work should be idempotent because process termination can
interrupt it. NACK handling must discard staged cursor changes and staged
delete or mark operations so polling can redeliver the batch. The SDK retries
NACKed batches with capped exponential backoff and stops after five consecutive
NACKs by default.
+
+A source whose accepted input cannot be re-read after a restart can override
`Source::batch_policy()` and call
`BatchPolicy::with_max_consecutive_nacks(None)` to keep retrying until
shutdown. The policy also offers `with_result_timeout` for sources that need a
different batch-result deadline. `on_batch_result()` errors still stop polling
regardless of the NACK limit.
+
+Sources that can distinguish retryable NACKs from ones that should count
toward the breaker can override `Source::nack_disposition()`. Returning
`NackDisposition::Retry` after a successful `on_batch_result(Nack)` keeps
polling even when the configured limit has been reached, while retaining capped
backoff. The default `ApplyPolicy` preserves the existing breaker. If
`on_batch_result()` returns an error, polling still stops because its
staged-work outcome is unknown.
The default `Source::on_batch_result()` implementation is a no-op for sources
without staged work. Sources that advance cursors, delete rows, or mark rows
must override it. The SDK stops polling if the handler returns an error,
preventing a failed rollback from advancing to another batch.
-This contract is a breaking FFI change. Source plugins must be rebuilt with
the matching SDK. The runtime loads `iggy_source_handle_v2`, which supplies a
batch ID to the runtime callback, and source plugins export
`iggy_source_batch_result` for the corresponding ACK or NACK.
+SDK 0.6 adds `Source::batch_policy()` and `Source::nack_disposition()` with
default implementations, so existing source implementations need no code change
when rebuilt. Sources that opt out of the NACK breaker must retain and replay
their rejected batch; otherwise disabling the stop only turns a visible failure
into a silent drop.
+
+The original batch-acknowledgment contract introduced a breaking FFI change.
Source plugins built before it must be rebuilt with the matching SDK. The
runtime loads `iggy_source_handle_v2`, which supplies a batch ID to the runtime
callback, and source plugins export `iggy_source_batch_result` for the
corresponding ACK or NACK. SDK 0.6 does not change these FFI signatures.
Review Comment:
This is true for the existing exports, but it does not mention the new one.
A plugin built before this SDK has no `iggy_source_register_stop_callback`. The
runtime loads it as `None`, logs nothing, and that plugin's poll-task stops
stay silent as before. Every source plugin from the current release is in that
group. Suggest naming the export and the rebuild here, and one `warn!` at load
when the export is missing.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]