diff --git a/.claude/skills/connector-runtime/SKILL.md b/.claude/skills/connector-runtime/SKILL.md index 8fb2695a65..b764c31fda 100644 --- a/.claude/skills/connector-runtime/SKILL.md +++ b/.claude/skills/connector-runtime/SKILL.md @@ -110,13 +110,15 @@ Don't mix. 1. `iggy_source_handle(id, send_callback)` - plugin registers itself. 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>` (`pub(crate)`). `SourceSenderEntry` wraps the sender + a pre-extracted owned `Counter` (the `errors` series, `Arc` inside). The FFI callback bumps errors on deserialize or channel-closed failure with one relaxed atomic - no `Family` lookup, no `Arc` handle. -4. `source_forwarding_loop` pulls from the channel, deserializes, applies transforms, encodes via `StreamEncoder`, sends to Iggy producer. -5. On success, save returned `ConnectorState` via `FileStateProvider`. +4. The runtime registers the optional source-stop callback, then `source_forwarding_loop` pulls from the channel, deserializes, applies transforms, encodes via `StreamEncoder`, and sends to the Iggy producer. The SDK warns after 30 seconds without a result but keeps the same batch pending; it does not poll or enqueue another copy. +5. On success, save returned `ConnectorState` via `FileStateProvider`, then return the batch result to the plugin. A result NACK may cause replay with capped backoff. +6. A stop-callback reason ends the forwarding loop and reports an unexpected stop. Closing the stop channel on normal shutdown disables that select arm, allowing already queued batches to receive results before the forwarding channel disconnects. +7. At loop exit, clean up `SOURCE_SENDERS` and update status and metrics. A self-stop reports `Error` and removes the source from the running gauge; normal shutdown reports `Stopped`. **Shutdown ordering (`manager/source.rs::stop_connector`):** 1. Call `iggy_source_close` FIRST. It blocks until the plugin's polling task stops, so no new send callbacks fire after it returns. -2. `cleanup_sender(plugin_id)` NEXT - dropping the channel sender makes the forwarding task's `recv_async()` resolve with `Disconnected` and exit cleanly, instead of blocking until the abort timeout. +2. `cleanup_sender(plugin_id)` NEXT - dropping the channel sender lets the forwarding task finish any queued batches before `recv_async()` resolves with `Disconnected` and exits cleanly, instead of blocking until the abort timeout. 3. Finally await spawned handlers with `tokio::time::timeout`. On timeout, `handle.abort()` + drain - prevents leaked tasks colliding with the next `start_connector` (a late `file.save()` could otherwise race the new instance). The silent-drop branch in `handle_produced_messages` only covers the window between close and cleanup. Gotchas: diff --git a/.claude/skills/connector-source/SKILL.md b/.claude/skills/connector-source/SKILL.md index 23613c5b05..9c9eb1ac64 100644 --- a/.claude/skills/connector-source/SKILL.md +++ b/.claude/skills/connector-source/SKILL.md @@ -61,13 +61,15 @@ let persisted = ConnectorState::serialize(&candidate, CONNECTOR_NAME, self.id) `on_batch_result` (added by #3855) is how a source learns what happened to the batch it just returned. The SDK keeps exactly one batch in flight: it will not call `poll()` again until this -returns, and it stops the source after `MAX_CONSECUTIVE_NACKS` (5) consecutive NACKs, roughly 1.5s -of backoff, without calling `close()`. +returns. By default it stops after `MAX_CONSECUTIVE_NACKS` (5) consecutive NACKs; a source can +override `batch_policy()` to disable the limit when its accepted input cannot be replayed after +restart. A stop is reported to the runtime only when the plugin exports the current SDK's stop +callback, so rebuild older source plugins. - `Ack` means the runtime sent the batch **and** persisted its state. `Nack` means it could not confirm both, which is **not** the same as neither happening: a batch that reached the topic but - whose state save failed is NACKed, and the SDK NACKs on its own result timeout while the send may - still have landed. A source that replays on `Nack` is at-least-once, not exactly-once. + whose state save failed is NACKed. A missing runtime result leaves the batch pending; it does not + produce a synthetic NACK. A source that replays on `Nack` is at-least-once, not exactly-once. - The trait has a **default no-op**, which suits only a source with no staged cursor and no destructive work. If `poll()` advances a cursor, deletes rows, or drains an in-memory buffer, omitting this loses data silently and nothing will tell you. `random_source` and `http_source` @@ -100,13 +102,14 @@ of backoff, without calling `close()`. state needed to resume them safely. - Keep `State` small - rewritten every batch. No unbounded vecs. -The SDK allows one in-flight batch. Five consecutive NACKs stop the source and -require a manual restart. Returning `Err` from `on_batch_result` is fatal, so +The SDK allows one in-flight batch. Five consecutive NACKs stop a source using +the default policy and require a manual restart. Returning `Err` from +`on_batch_result` is fatal regardless of the NACK limit, so retry transient backend failures inside the callback before returning an error. -The runtime must report ACK or NACK within the SDK's 30-second batch-result -window. Once the result is received, the SDK waits for `on_batch_result` to -finish, so the callback must bound its own connection acquisition and retry -budget rather than relying on the SDK deadline. +If the runtime has not reported ACK or NACK after 30 seconds, the SDK warns +and keeps waiting for that result without polling or replaying the batch. +Once the result is received, the SDK waits for `on_batch_result` to finish, +so the callback must bound its own connection acquisition and retry budget. ### Sleep first diff --git a/Cargo.lock b/Cargo.lock index 0b28866872..2138a222e6 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -7635,7 +7635,7 @@ dependencies = [ [[package]] name = "iggy_connector_sdk" -version = "0.5.1-edge.1" +version = "0.5.1-edge.2" dependencies = [ "anyhow", "apache-avro 0.22.0", diff --git a/Cargo.toml b/Cargo.toml index e868a9421c..01c3f0acd7 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -217,7 +217,7 @@ iggy-gateway-kafka = { path = "gateways/kafka" } iggy_binary_protocol = { path = "core/binary_protocol", version = "0.11.1-edge.1" } iggy_common = { path = "core/common", version = "0.11.1-edge.1" } iggy_connector_doris_sink = { path = "core/connectors/sinks/doris_sink" } -iggy_connector_sdk = { path = "core/connectors/sdk", version = "0.5.1-edge.1" } +iggy_connector_sdk = { path = "core/connectors/sdk", version = "0.5.1-edge.2" } indexmap = "2.14.2" integration = { path = "core/integration" } ipnet = "2.12.2" diff --git a/core/connectors/runtime/src/main.rs b/core/connectors/runtime/src/main.rs index 55f1665523..c53e28f209 100644 --- a/core/connectors/runtime/src/main.rs +++ b/core/connectors/runtime/src/main.rs @@ -37,10 +37,11 @@ use iggy_connector_sdk::{ StreamDecoder, StreamEncoder, api::ConnectorStatus, sink::ConsumeCallback, - source::{BatchResultCallback, HandleCallback, SendCallback}, + source::{BatchResultCallback, HandleCallback, SendCallback, SourceStoppedCallback}, transforms::Transform, }; use mimalloc::MiMalloc; +use source::RegisterStopCallback; use state::StateStorage; use std::{ collections::HashMap, @@ -109,6 +110,8 @@ pub(crate) struct SourceApi { log_callback: iggy_connector_sdk::LogCallback, ) -> i32, iggy_source_handle_v2: extern "C" fn(id: u32, callback: SendCallback) -> i32, + iggy_source_register_stop_callback: + Option i32>, iggy_source_batch_result: extern "C" fn(plugin_id: u32, batch_id: u64, result: u8) -> i32, iggy_source_close: extern "C" fn(id: u32) -> i32, iggy_source_version: extern "C" fn() -> *const std::ffi::c_char, @@ -266,12 +269,14 @@ async fn run() -> Result<(), RuntimeError> { for (_path, source) in sources { let container = source.container; let handle_callback = container.iggy_source_handle_v2; + let register_stop_callback = container.iggy_source_register_stop_callback; let batch_result_callback = container.iggy_source_batch_result; for plugin in &source.plugins { source_containers_by_key.insert(plugin.key.clone(), container.clone()); } source_wrappers.push(SourceConnectorWrapper { handle_callback, + register_stop_callback, batch_result_callback, plugins: source.plugins, }); @@ -583,6 +588,7 @@ struct SourceConnectorProducer { struct SourceConnectorWrapper { handle_callback: HandleCallback, + register_stop_callback: Option, batch_result_callback: BatchResultCallback, plugins: Vec, } diff --git a/core/connectors/runtime/src/manager/source.rs b/core/connectors/runtime/src/manager/source.rs index fe962b1ebf..44e0cf667c 100644 --- a/core/connectors/runtime/src/manager/source.rs +++ b/core/connectors/runtime/src/manager/source.rs @@ -35,6 +35,13 @@ use tokio::sync::Mutex; use tokio::task::JoinHandle; use tracing::{info, warn}; +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +pub enum StopReport { + Ignored, + AlreadyError, + NewError, +} + #[derive(Debug)] pub struct SourceManager { sources: DashMap>>, @@ -88,18 +95,32 @@ impl SourceManager { pub async fn set_error(&self, key: &str, error_message: &str, metrics: Option<&Arc>) { if let Some(source) = self.sources.get(key) { - let mut source = source.lock().await; - // Through the shared transition, so leaving `Running` moves the - // gauge. Skipping it left an errored instance counted as running, - // and the loop's later `Stopped` could not correct that either, - // because by then the old status was `Error` and neither branch - // fires. - // - // The message is assigned after the transition, and that ordering is - // what preserves it. `Error` being outside the set that clears - // `last_error` is belt and braces here, not the mechanism. - source.apply_status(ConnectorStatus::Error, metrics); - source.info.last_error = Some(ConnectorError::new(error_message)); + source.lock().await.record_error(error_message, metrics); + } + } + + pub async fn report_unexpected_stop( + &self, + key: &str, + error_message: &str, + metrics: &Arc, + ) -> StopReport { + let Some(source) = self.sources.get(key).map(|entry| entry.value().clone()) else { + return StopReport::Ignored; + }; + let mut source = source.lock().await; + if matches!( + source.info.status, + ConnectorStatus::Stopping | ConnectorStatus::Stopped + ) { + return StopReport::Ignored; + } + let already_error = source.info.status == ConnectorStatus::Error; + source.record_error(error_message, Some(metrics)); + if already_error { + StopReport::AlreadyError + } else { + StopReport::NewError } } @@ -258,6 +279,7 @@ impl SourceManager { }; let handle_callback = container.iggy_source_handle_v2; + let register_stop_callback = container.iggy_source_register_stop_callback; let batch_result_callback = container.iggy_source_batch_result; // The lock is taken before the spawn so nothing can await between @@ -286,6 +308,7 @@ impl SourceManager { transforms, state_storage, handle_callback, + register_stop_callback, batch_result_callback, context.clone(), ); @@ -362,6 +385,13 @@ pub struct SourceDetails { } impl SourceDetails { + fn record_error(&mut self, message: &str, metrics: Option<&Arc>) { + // Update the gauge before storing the error; `apply_status` clears old + // error text only when entering Running or Stopped. + self.apply_status(ConnectorStatus::Error, metrics); + self.info.last_error = Some(ConnectorError::new(message)); + } + /// Applies a status transition and the gauge move that belongs with it. /// /// On `&mut self` rather than behind a key, so it can run inside a lock the @@ -526,6 +556,80 @@ mod tests { assert_eq!(metrics.get_sources_running(), 1); } + #[tokio::test] + async fn given_poll_task_exits_when_running_should_report_error_and_decrement_gauge() { + let metrics = Arc::new(Metrics::init()); + let mut details = create_test_source_details("pg", 1); + details.info.status = ConnectorStatus::Stopped; + let manager = SourceManager::new(vec![details]); + manager + .update_status("pg", ConnectorStatus::Running, Some(&metrics)) + .await; + + assert_eq!( + manager + .report_unexpected_stop("pg", "polling stopped", &metrics) + .await, + StopReport::NewError + ); + let source = manager.get("pg").await.expect("source exists"); + let source = source.lock().await; + assert_eq!(source.info.status, ConnectorStatus::Error); + assert_eq!(metrics.get_sources_running(), 0); + assert!(source.info.last_error.is_some()); + } + + #[tokio::test] + async fn given_batch_error_before_stop_should_not_count_a_second_error() { + let metrics = Arc::new(Metrics::init()); + let mut details = create_test_source_details("pg", 1); + details.info.status = ConnectorStatus::Stopped; + let manager = SourceManager::new(vec![details]); + manager + .update_status("pg", ConnectorStatus::Running, Some(&metrics)) + .await; + manager + .set_error("pg", "batch rejected", Some(&metrics)) + .await; + + assert_eq!( + manager + .report_unexpected_stop("pg", "NACK limit reached", &metrics) + .await, + StopReport::AlreadyError + ); + assert_eq!(metrics.get_sources_running(), 0); + let source = manager.get("pg").await.expect("source exists"); + let source = source.lock().await; + assert_eq!(source.info.status, ConnectorStatus::Error); + assert!(source.info.last_error.is_some()); + } + + #[tokio::test] + async fn given_user_shutdown_when_poll_task_exits_should_not_report_error() { + let metrics = Arc::new(Metrics::init()); + let mut details = create_test_source_details("pg", 1); + details.info.status = ConnectorStatus::Stopped; + let manager = SourceManager::new(vec![details]); + manager + .update_status("pg", ConnectorStatus::Running, Some(&metrics)) + .await; + manager + .update_status("pg", ConnectorStatus::Stopping, Some(&metrics)) + .await; + + assert_eq!( + manager + .report_unexpected_stop("pg", "polling stopped", &metrics) + .await, + StopReport::Ignored + ); + let source = manager.get("pg").await.expect("source exists"); + let source = source.lock().await; + assert_eq!(source.info.status, ConnectorStatus::Stopping); + assert!(source.info.last_error.is_none()); + } + #[tokio::test] async fn should_not_double_count_when_an_error_falls_between_two_running_reports() { // The interleaving that double counts: the forwarding loop reports diff --git a/core/connectors/runtime/src/source.rs b/core/connectors/runtime/src/source.rs index 5722f2faae..36f4eb2396 100644 --- a/core/connectors/runtime/src/source.rs +++ b/core/connectors/runtime/src/source.rs @@ -27,7 +27,10 @@ use iggy_connector_sdk::encoders::avro::{AvroEncoderConfig, AvroStreamEncoder}; use iggy_connector_sdk::{ ConnectorState, DecodedMessage, Error as SdkError, ProducedMessages, Schema, StreamEncoder, TopicMetadata, - source::{BatchResultCallback, HandleCallback, SourceBatchResult}, + source::{ + BatchResultCallback, HandleCallback, SourceBatchResult, SourceStopReason, + SourceStoppedCallback, + }, transforms::Transform, }; use std::{ @@ -43,6 +46,7 @@ use crate::benchmark; use crate::configs::connectors::SourceConfig; use crate::context::RuntimeContext; use crate::log::LOG_CALLBACK; +use crate::manager::source::StopReport; use crate::metrics::ConnectorType; use crate::metrics::SourceLabels; use crate::{ @@ -54,13 +58,17 @@ use crate::{ use iggy_connector_sdk::api::ConnectorStatus; use prometheus_client::metrics::counter::Counter; use tokio::runtime::Handle; +use tokio::sync::mpsc; use tokio::task::JoinHandle; const MAX_FAILED_TAIL_RETRIES: u32 = 3; const SOURCE_TOPIC_MESSAGES_REQUIRED_TO_SAVE: u32 = 1; +pub(crate) type RegisterStopCallback = extern "C" fn(u32, SourceStoppedCallback) -> i32; + pub(crate) struct SourceSenderEntry { pub(crate) sender: Sender, + stop_sender: mpsc::UnboundedSender, // Owned errors counter (Arc inside) so the FFI callback bumps // it with one relaxed atomic - no Family RwLock + HashMap lookup per call. pub(crate) error_counter: Counter, @@ -79,6 +87,12 @@ pub(crate) fn cleanup_sender(plugin_id: u32) { SOURCE_SENDERS.remove(&plugin_id); } +extern "C" fn source_stopped(plugin_id: u32, reason: u8) { + if let Some(entry) = SOURCE_SENDERS.get(&plugin_id) { + let _ = entry.stop_sender.send(SourceStopReason::from(reason)); + } +} + /// Initializes all enabled source connectors. /// /// Per-connector failures (path resolution, dlopen, plugin init, @@ -562,6 +576,7 @@ pub(crate) async fn source_forwarding_loop( transforms: Vec>, state_storage: StateStorage, receiver: Receiver, + mut stop_receiver: mpsc::UnboundedReceiver, batch_result_callback: BatchResultCallback, context: Arc, labels: Arc, @@ -588,7 +603,24 @@ pub(crate) async fn source_forwarding_loop( topic: producer.topic().to_string(), }; - while let Ok(produced_batch) = receiver.recv_async().await { + let mut stop_reason = None; + let mut stop_receiver_open = true; + loop { + let produced_batch = tokio::select! { + biased; + stopped = stop_receiver.recv(), if stop_receiver_open => { + if let Some(reason) = stopped { + stop_reason = Some(reason); + break; + } + stop_receiver_open = false; + continue; + }, + batch = receiver.recv_async() => batch, + }; + let Ok(produced_batch) = produced_batch else { + break; + }; let total_start = Instant::now(); let batch_id = produced_batch.id; let produced_messages = produced_batch.messages; @@ -816,14 +848,34 @@ pub(crate) async fn source_forwarding_loop( } info!("Source connector with ID: {plugin_id} stopped."); - context - .sources - .update_status( - &plugin_key, - ConnectorStatus::Stopped, - Some(&context.metrics), - ) - .await; + // A self-stopped source must release its sender even when no manager stop follows. + 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" + ); + match context + .sources + .report_unexpected_stop(&plugin_key, &error_msg, &context.metrics) + .await + { + StopReport::NewError => { + error!("{error_msg}"); + context.metrics.inc_errors_with_labels(&labels.counter); + } + StopReport::AlreadyError => error!("{error_msg}"), + StopReport::Ignored => {} + } + } else { + context + .sources + .update_status( + &plugin_key, + ConnectorStatus::Stopped, + Some(&context.metrics), + ) + .await; + } } fn should_recover_source(batch_result: SourceBatchResult, sent_count: usize) -> bool { @@ -841,22 +893,25 @@ pub(crate) fn spawn_source_handler( transforms: Vec>, state_storage: StateStorage, handle_callback: HandleCallback, + register_stop_callback: Option, batch_result_callback: BatchResultCallback, context: Arc, ) -> Vec> { let (sender, receiver) = flume::unbounded(); + let (stop_sender, stop_receiver) = mpsc::unbounded_channel(); let plugin_key = plugin_key.to_string(); let labels = Arc::new(SourceLabels::new(&plugin_key)); SOURCE_SENDERS.insert( plugin_id, SourceSenderEntry { sender, + stop_sender, error_counter: context.metrics.error_counter(&labels.counter), }, ); let blocking_handle = tokio::task::spawn_blocking(move || { - handle_callback(plugin_id, handle_produced_messages); + start_source_polling(plugin_id, handle_callback, register_stop_callback) }); let handler_task = tokio::spawn(async move { source_forwarding_loop( @@ -869,6 +924,7 @@ pub(crate) fn spawn_source_handler( transforms, state_storage, receiver, + stop_receiver, batch_result_callback, context, labels, @@ -879,6 +935,32 @@ pub(crate) fn spawn_source_handler( vec![blocking_handle, handler_task] } +fn start_source_polling( + plugin_id: u32, + handle_callback: HandleCallback, + register_stop_callback: Option, +) { + let registration_failed = match register_stop_callback { + Some(register) => register(plugin_id, source_stopped) != 0, + None => { + warn!( + "Source connector with ID: {plugin_id} has no stop callback export; rebuild it with the current connector SDK to report poll-task stops" + ); + false + } + }; + let stop_reason = if registration_failed { + Some(SourceStopReason::RegistrationFailed) + } else if handle_callback(plugin_id, handle_produced_messages) != 0 { + Some(SourceStopReason::HandlerFailed) + } else { + None + }; + if let Some(reason) = stop_reason { + source_stopped(plugin_id, reason as u8); + } +} + pub fn handle( sources: Vec, context: Arc, @@ -912,6 +994,7 @@ pub fn handle( plugin.transforms, plugin.state_storage, source.handle_callback, + source.register_stop_callback, source.batch_result_callback, context.clone(), ); @@ -1126,13 +1209,162 @@ fn build_iggy_message( #[cfg(test)] mod tests { use super::*; + use crate::configs::connectors::create_connectors_config_provider; + use crate::configs::runtime::{ConnectorsConfig, LocalConnectorsConfig}; + use crate::manager::sink::SinkManager; + use crate::manager::source::{SourceDetails, SourceInfo, SourceManager}; + use crate::metrics::Metrics; + use crate::state::FileStateFactory; + use crate::stream::IggyClients; + use iggy_common::IggyTimestamp; + use iggy_connector_sdk::source::SendCallback; + use secrecy::SecretString; use std::collections::VecDeque; use std::future::ready; use std::sync::Mutex; - use std::sync::atomic::{AtomicU32, Ordering}; + use std::sync::atomic::{AtomicU32, AtomicU64, AtomicUsize, Ordering}; use std::time::Duration; static TEST_PLUGIN_ID: AtomicU32 = AtomicU32::new(u32::MAX / 2); + static HANDLE_CALLS: AtomicUsize = AtomicUsize::new(0); + static DRAINED_BATCH_ID: AtomicU64 = AtomicU64::new(0); + static DRAINED_BATCH_RESULT: AtomicU32 = AtomicU32::new(u32::MAX); + + extern "C" fn reject_stop_registration(_: u32, _: SourceStoppedCallback) -> i32 { + -1 + } + + extern "C" fn accept_stop_registration(_: u32, _: SourceStoppedCallback) -> i32 { + 0 + } + + extern "C" fn reject_source_handle(_: u32, _: SendCallback) -> i32 { + HANDLE_CALLS.fetch_add(1, Ordering::SeqCst); + -1 + } + + extern "C" fn reject_source_handle_untracked(_: u32, _: SendCallback) -> i32 { + -1 + } + + extern "C" fn ignore_batch_result(_: u32, _: u64, _: u8) -> i32 { + 0 + } + + extern "C" fn record_drained_batch_result(_: u32, batch_id: u64, result: u8) -> i32 { + DRAINED_BATCH_RESULT.store(u32::from(result), Ordering::SeqCst); + DRAINED_BATCH_ID.store(batch_id, Ordering::SeqCst); + 0 + } + + async fn test_source_runtime( + plugin_id: u32, + plugin_key: &str, + directory: &std::path::Path, + ) -> (Arc, StateStorage, IggyProducer) { + let config_provider = + create_connectors_config_provider(&ConnectorsConfig::Local(LocalConnectorsConfig { + config_dir: directory.display().to_string(), + })) + .await + .expect("local config provider should initialize"); + let state_factory = Arc::new(FileStateFactory::new(directory.display().to_string())); + let state_storage = state_factory + .storage_for(plugin_key) + .expect("file state storage should initialize"); + let producer = IggyClient::default() + .producer("stream", "topic") + .expect("producer builder should initialize") + .build(); + let context = Arc::new(RuntimeContext { + sinks: SinkManager::new(vec![]), + sources: SourceManager::new(vec![SourceDetails { + info: SourceInfo { + id: plugin_id, + key: plugin_key.to_string(), + name: plugin_key.to_string(), + path: "test".to_string(), + version: "test".to_string(), + enabled: true, + status: ConnectorStatus::Stopped, + last_error: None, + plugin_config_format: None, + }, + config: SourceConfig { + key: plugin_key.to_string(), + ..SourceConfig::default() + }, + handler_tasks: vec![], + container: None, + restart_guard: Arc::new(tokio::sync::Mutex::new(())), + }]), + api_key: SecretString::from("test".to_string()), + config_provider: Arc::from(config_provider), + metrics: Arc::new(Metrics::init()), + start_time: IggyTimestamp::now(), + iggy_clients: Arc::new(IggyClients { + producer: IggyClient::default(), + consumer: IggyClient::default(), + }), + state_factory, + }); + (context, state_storage, producer) + } + + async fn assert_start_failure( + register_stop_callback: Option, + expected_reason: SourceStopReason, + ) { + let plugin_id = next_plugin_id(); + let plugin_key = format!("source_{plugin_id}"); + let directory = tempfile::tempdir().expect("test directory should exist"); + let (context, state_storage, producer) = + test_source_runtime(plugin_id, &plugin_key, directory.path()).await; + + let tasks = spawn_source_handler( + plugin_id, + &plugin_key, + false, + false, + producer, + Schema::Raw.encoder(), + vec![], + state_storage, + reject_source_handle_untracked, + register_stop_callback, + ignore_batch_result, + Arc::clone(&context), + ); + for task in tasks { + tokio::time::timeout(Duration::from_secs(2), task) + .await + .expect("source task should stop") + .expect("source task should complete"); + } + let source = context + .sources + .get(&plugin_key) + .await + .expect("source should be registered"); + let source = source.lock().await; + assert_eq!(source.info.status, ConnectorStatus::Error); + let message = &source + .info + .last_error + .as_ref() + .expect("stop should record the reason") + .message; + assert!(message.contains(&expected_reason.to_string())); + assert!(message.contains("restart required")); + assert_eq!(context.metrics.get_sources_running(), 0); + assert_eq!( + context + .metrics + .error_counter(&SourceLabels::new(&plugin_key).counter) + .get(), + 1 + ); + } fn next_plugin_id() -> u32 { TEST_PLUGIN_ID.fetch_add(1, Ordering::Relaxed) @@ -1320,10 +1552,12 @@ mod tests { let plugin_id = next_plugin_id(); let batch_id = 73; let (sender, receiver) = flume::unbounded(); + let (stop_sender, _stop_receiver) = mpsc::unbounded_channel(); SOURCE_SENDERS.insert( plugin_id, SourceSenderEntry { sender, + stop_sender, error_counter: Counter::default(), }, ); @@ -1352,15 +1586,276 @@ mod tests { cleanup_sender(plugin_id); } + #[test] + fn given_source_stop_when_callback_runs_should_signal_forwarder() { + let plugin_id = next_plugin_id(); + let (sender, receiver) = flume::unbounded(); + let (stop_sender, mut stop_receiver) = mpsc::unbounded_channel(); + SOURCE_SENDERS.insert( + plugin_id, + SourceSenderEntry { + sender, + stop_sender, + error_counter: Counter::default(), + }, + ); + + source_stopped(plugin_id, SourceStopReason::NackLimit as u8); + assert_eq!(stop_receiver.try_recv(), Ok(SourceStopReason::NackLimit)); + + cleanup_sender(plugin_id); + assert!(receiver.is_disconnected()); + } + + #[test] + fn given_registration_failure_should_stop_without_starting_handle() { + let plugin_id = next_plugin_id(); + let (sender, _receiver) = flume::unbounded(); + let (stop_sender, mut stop_receiver) = mpsc::unbounded_channel(); + SOURCE_SENDERS.insert( + plugin_id, + SourceSenderEntry { + sender, + stop_sender, + error_counter: Counter::default(), + }, + ); + HANDLE_CALLS.store(0, Ordering::SeqCst); + + start_source_polling( + plugin_id, + reject_source_handle, + Some(reject_stop_registration), + ); + assert_eq!( + stop_receiver.try_recv(), + Ok(SourceStopReason::RegistrationFailed) + ); + assert_eq!(HANDLE_CALLS.load(Ordering::SeqCst), 0); + cleanup_sender(plugin_id); + } + + #[test] + fn given_handle_failure_should_signal_stop_reason() { + let plugin_id = next_plugin_id(); + let (sender, _receiver) = flume::unbounded(); + let (stop_sender, mut stop_receiver) = mpsc::unbounded_channel(); + SOURCE_SENDERS.insert( + plugin_id, + SourceSenderEntry { + sender, + stop_sender, + error_counter: Counter::default(), + }, + ); + + start_source_polling( + plugin_id, + reject_source_handle_untracked, + Some(accept_stop_registration), + ); + assert_eq!( + stop_receiver.try_recv(), + Ok(SourceStopReason::HandlerFailed) + ); + cleanup_sender(plugin_id); + } + + #[test] + fn given_old_plugin_without_stop_export_should_still_start_handle() { + let plugin_id = next_plugin_id(); + let (sender, _receiver) = flume::unbounded(); + let (stop_sender, mut stop_receiver) = mpsc::unbounded_channel(); + SOURCE_SENDERS.insert( + plugin_id, + SourceSenderEntry { + sender, + stop_sender, + error_counter: Counter::default(), + }, + ); + + start_source_polling(plugin_id, reject_source_handle_untracked, None); + assert_eq!( + stop_receiver.try_recv(), + Ok(SourceStopReason::HandlerFailed) + ); + cleanup_sender(plugin_id); + } + + #[tokio::test] + async fn given_stop_registration_failure_should_report_restart_required() { + assert_start_failure( + Some(reject_stop_registration), + SourceStopReason::RegistrationFailed, + ) + .await; + } + + #[tokio::test] + async fn given_source_handle_failure_should_report_restart_required() { + assert_start_failure( + Some(accept_stop_registration), + SourceStopReason::HandlerFailed, + ) + .await; + } + + #[tokio::test] + async fn given_queued_batch_and_existing_error_when_stopped_should_not_process_or_recount() { + let plugin_id = next_plugin_id(); + let plugin_key = format!("source_{plugin_id}"); + let directory = tempfile::tempdir().expect("test directory should exist"); + let (context, state_storage, producer) = + test_source_runtime(plugin_id, &plugin_key, directory.path()).await; + let (sender, receiver) = flume::unbounded(); + let queued = receiver.clone(); + let (stop_sender, stop_receiver) = mpsc::unbounded_channel(); + SOURCE_SENDERS.insert( + plugin_id, + SourceSenderEntry { + sender: sender.clone(), + stop_sender: stop_sender.clone(), + error_counter: Counter::default(), + }, + ); + let labels = Arc::new(SourceLabels::new(&plugin_key)); + let forwarding = tokio::spawn(source_forwarding_loop( + plugin_id, + plugin_key.clone(), + false, + false, + producer, + Schema::Raw.encoder(), + vec![], + state_storage, + receiver, + stop_receiver, + ignore_batch_result, + Arc::clone(&context), + Arc::clone(&labels), + )); + + tokio::time::timeout(Duration::from_secs(2), async { + while context.metrics.get_sources_running() != 1 { + tokio::task::yield_now().await; + } + }) + .await + .expect("forwarding loop should reach Running"); + context + .sources + .set_error(&plugin_key, "prior send failure", Some(&context.metrics)) + .await; + sender + .send(ProducedBatch { + id: 1, + messages: ProducedMessages { + schema: Schema::Raw, + messages: vec![], + state: None, + }, + }) + .expect("batch should enter the queue"); + stop_sender + .send(SourceStopReason::NackLimit) + .expect("stop should reach the forwarding loop"); + + tokio::time::timeout(Duration::from_secs(2), forwarding) + .await + .expect("forwarding loop should stop") + .expect("forwarding task should complete"); + assert_eq!(queued.try_recv().expect("queued batch must remain").id, 1); + let source = context + .sources + .get(&plugin_key) + .await + .expect("source should remain registered"); + let source = source.lock().await; + assert_eq!(source.info.status, ConnectorStatus::Error); + assert!( + source + .info + .last_error + .as_ref() + .expect("stop reason should replace the prior error") + .message + .contains("consecutive NACK limit reached") + ); + assert_eq!(context.metrics.get_sources_running(), 0); + assert_eq!(context.metrics.error_counter(&labels.counter).get(), 0); + assert!(!SOURCE_SENDERS.contains_key(&plugin_id)); + } + + #[tokio::test] + async fn given_queued_batch_when_stop_channel_closes_should_deliver_result_before_stopping() { + DRAINED_BATCH_ID.store(0, Ordering::SeqCst); + DRAINED_BATCH_RESULT.store(u32::MAX, Ordering::SeqCst); + let plugin_id = next_plugin_id(); + let plugin_key = format!("source_{plugin_id}"); + let directory = tempfile::tempdir().expect("test directory should exist"); + let (context, state_storage, producer) = + test_source_runtime(plugin_id, &plugin_key, directory.path()).await; + let (sender, receiver) = flume::unbounded(); + let (stop_sender, stop_receiver) = mpsc::unbounded_channel(); + sender + .send(ProducedBatch { + id: 41, + messages: ProducedMessages { + schema: Schema::Raw, + messages: vec![], + state: None, + }, + }) + .expect("batch should enter the queue"); + drop(sender); + drop(stop_sender); + + tokio::time::timeout( + Duration::from_secs(2), + source_forwarding_loop( + plugin_id, + plugin_key.clone(), + false, + false, + producer, + Schema::Raw.encoder(), + vec![], + state_storage, + receiver, + stop_receiver, + record_drained_batch_result, + Arc::clone(&context), + Arc::new(SourceLabels::new(&plugin_key)), + ), + ) + .await + .expect("forwarding loop should drain and stop"); + + assert_eq!(DRAINED_BATCH_ID.load(Ordering::SeqCst), 41); + assert_eq!( + DRAINED_BATCH_RESULT.load(Ordering::SeqCst), + u32::from(SourceBatchResult::Ack as u8) + ); + let source = context + .sources + .get(&plugin_key) + .await + .expect("source should remain registered"); + assert_eq!(source.lock().await.info.status, ConnectorStatus::Stopped); + } + #[test] fn given_invalid_payload_when_callback_runs_should_reject_batch() { let plugin_id = next_plugin_id(); let (sender, _receiver) = flume::unbounded(); + let (stop_sender, _stop_receiver) = mpsc::unbounded_channel(); let error_counter = Counter::default(); SOURCE_SENDERS.insert( plugin_id, SourceSenderEntry { sender, + stop_sender, error_counter: error_counter.clone(), }, ); diff --git a/core/connectors/sdk/Cargo.toml b/core/connectors/sdk/Cargo.toml index aad4b621ff..0cef19204b 100644 --- a/core/connectors/sdk/Cargo.toml +++ b/core/connectors/sdk/Cargo.toml @@ -17,7 +17,7 @@ [package] name = "iggy_connector_sdk" -version = "0.5.1-edge.1" +version = "0.5.1-edge.2" description = "Iggy is the persistent message streaming platform written in Rust, supporting QUIC, TCP and HTTP transport protocols, capable of processing millions of messages per second." edition = "2024" license = "Apache-2.0" diff --git a/core/connectors/sdk/README.md b/core/connectors/sdk/README.md index f046397dc1..a0e053ec56 100644 --- a/core/connectors/sdk/README.md +++ b/core/connectors/sdk/README.md @@ -29,11 +29,19 @@ 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. -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. +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 SDK warns after 30 seconds without a batch result, but keeps that batch in flight until the runtime replies or the source shuts down. -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. +The source receives only ACK or NACK, not the runtime's failure cause. It cannot classify a NACK as transient. Disable the limit only when repeated replay is safer than stopping; a permanent failure still needs operator attention. If `on_batch_result()` returns an error, polling 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. + +SDK 0.5.1-edge.2 adds `Source::batch_policy()` with a default implementation, 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. + +Rebuild source plugins to export `iggy_source_register_stop_callback`; the runtime warns when an older plugin lacks it. The callback reports why polling stopped so the runtime can mark the source `Error` and update its running gauge. + +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.5.1-edge.2 also adds the optional stop-callback export; older plugins cannot report an unexpected poll-task exit. Moreover, it contains both, the `decoders` and `encoders` modules, implementing either `StreamDecoder` or `StreamEncoder` traits, which are used when consuming or producing data from/to Iggy streams. diff --git a/core/connectors/sdk/src/lib.rs b/core/connectors/sdk/src/lib.rs index 1fe5a67e46..cde3797db1 100644 --- a/core/connectors/sdk/src/lib.rs +++ b/core/connectors/sdk/src/lib.rs @@ -109,6 +109,15 @@ 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 whether repeated NACKs stop this source. + /// + /// 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. + /// The SDK reads this policy once after `open()`; later changes have no effect. + fn batch_policy(&self) -> source::BatchPolicy { + source::BatchPolicy::default() + } + /// Retrieves the next batch for the runtime to process and deliver. async fn poll(&self) -> Result; diff --git a/core/connectors/sdk/src/source.rs b/core/connectors/sdk/src/source.rs index 09d18bf673..7a29412949 100644 --- a/core/connectors/sdk/src/source.rs +++ b/core/connectors/sdk/src/source.rs @@ -17,11 +17,12 @@ use crate::log::{CallbackLayer, LogCallback}; use crate::retry::exponential_backoff; -use crate::{ConnectorState, Source, get_runtime}; +use crate::{ConnectorState, ProducedMessages, Source, get_runtime}; use serde::de::DeserializeOwned; +use std::num::NonZeroU32; use std::sync::{ Arc, Mutex, MutexGuard, PoisonError, - atomic::{AtomicU32, Ordering}, + atomic::{AtomicBool, AtomicU32, Ordering}, }; use std::time::Duration; use tokio::{ @@ -50,12 +51,89 @@ 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, reason: u8); -/// Maximum time the runtime may take to report a source batch result. +/// Reason source polling stopped, reported by the SDK callback or the runtime. +#[derive(Debug, Clone, Copy, PartialEq, Eq)] +#[repr(u8)] +pub enum SourceStopReason { + /// The SDK poll task panicked, or the runtime received an unknown reason byte. + Unexpected = 0, + /// The SDK reached the source's configured consecutive-NACK limit. + NackLimit = 1, + /// The SDK could not apply a batch result to the source. + BatchResultError = 2, + /// The SDK lost the pending batch's result channel. + ResultChannelClosed = 3, + /// The runtime could not register the SDK's stop callback. + RegistrationFailed = 4, + /// The runtime could not start the source handler. + HandlerFailed = 5, + /// Normal SDK shutdown; this reason is not sent to the runtime callback. + Shutdown = 6, +} + +impl std::fmt::Display for SourceStopReason { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter.write_str(match self { + Self::Unexpected => "source poll task exited unexpectedly", + Self::NackLimit => "consecutive NACK limit reached", + Self::BatchResultError => "batch result could not be applied", + Self::ResultChannelClosed => "batch result channel closed", + Self::RegistrationFailed => "source stop callback registration failed", + Self::HandlerFailed => "source handler failed to start", + Self::Shutdown => "source shut down", + }) + } +} + +impl From for SourceStopReason { + fn from(value: u8) -> Self { + match value { + 1 => Self::NackLimit, + 2 => Self::BatchResultError, + 3 => Self::ResultChannelClosed, + 4 => Self::RegistrationFailed, + 5 => Self::HandlerFailed, + 6 => Self::Shutdown, + _ => Self::Unexpected, + } + } +} + +struct StopNotification { + plugin_id: u32, + callback: Option, + closing: Arc, +} + +impl StopNotification { + fn report(mut self, reason: SourceStopReason) { + let callback = self.callback.take(); + if reason != SourceStopReason::Shutdown + && !self.closing.load(Ordering::Acquire) + && let Some(callback) = callback + { + callback(self.plugin_id, reason as u8); + } + } +} + +impl Drop for StopNotification { + fn drop(&mut self) { + if !self.closing.load(Ordering::Acquire) + && let Some(callback) = self.callback + { + callback(self.plugin_id, SourceStopReason::Unexpected as u8); + } + } +} + +/// Time after which a missing source batch result is logged as a warning. pub const BATCH_RESULT_TIMEOUT: Duration = Duration::from_secs(30); const NACK_RETRY_DELAY: Duration = Duration::from_millis(100); const MAX_NACK_RETRY_DELAY: Duration = Duration::from_secs(5); -/// Number of consecutive rejected batches after which a source is stopped. +/// Default number of consecutive rejected batches after which a source is stopped. pub const MAX_CONSECUTIVE_NACKS: u32 = 5; /// Delivery result for the single batch currently in flight from a source plugin. @@ -97,15 +175,28 @@ enum PendingBatchClearResult { #[derive(Debug, Clone, Copy, PartialEq, Eq)] enum BatchCompletion { Applied(SourceBatchResult), - Stop, + Stop(SourceStopReason), } +/// Polling policy for one source instance. The default stops after five NACKs. #[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, +} + +impl BatchPolicy { + pub fn max_consecutive_nacks(&self) -> Option { + self.max_consecutive_nacks + } + + /// Sets the NACK limit. `None` keeps polling with capped backoff until shutdown. + pub fn with_max_consecutive_nacks(mut self, max_consecutive_nacks: Option) -> Self { + self.max_consecutive_nacks = max_consecutive_nacks; + self + } } impl Default for BatchPolicy { @@ -114,7 +205,7 @@ impl Default for BatchPolicy { result_timeout: BATCH_RESULT_TIMEOUT, nack_retry_delay: NACK_RETRY_DELAY, max_nack_retry_delay: MAX_NACK_RETRY_DELAY, - max_consecutive_nacks: MAX_CONSECUTIVE_NACKS, + max_consecutive_nacks: NonZeroU32::new(MAX_CONSECUTIVE_NACKS), } } } @@ -127,6 +218,9 @@ pub struct SourceContainer { task: Option>, pending_batch: Arc>>, consecutive_nacks: Arc, + batch_policy: BatchPolicy, + stop_callback: Option, + closing: Arc, } impl SourceContainer { @@ -138,6 +232,9 @@ impl SourceContainer { task: None, pending_batch: Arc::new(Mutex::new(None)), consecutive_nacks: Arc::new(AtomicU32::new(0)), + batch_policy: BatchPolicy::default(), + stop_callback: None, + closing: Arc::new(AtomicBool::new(false)), } } @@ -185,6 +282,7 @@ impl SourceContainer { let mut source = factory(id, config, state); let runtime = get_runtime(); let result = runtime.block_on(source.open()); + self.batch_policy = source.batch_policy(); self.id = id; self.source = Some(Arc::new(source)); if result.is_ok() { 0 } else { 1 } @@ -203,6 +301,7 @@ impl SourceContainer { }; info!("Closing source connector with ID: {}...", self.id); + self.closing.store(true, Ordering::Release); if let Some(sender) = self.shutdown.take() { let _ = sender.send(()); } @@ -246,8 +345,14 @@ impl SourceContainer { let source = Arc::clone(source); let pending_batch = Arc::clone(&self.pending_batch); let consecutive_nacks = Arc::clone(&self.consecutive_nacks); + let batch_policy = self.batch_policy; + let stop_notification = StopNotification { + plugin_id, + callback: self.stop_callback, + closing: Arc::clone(&self.closing), + }; let handle = runtime.spawn(async move { - handle_messages( + let reason = handle_messages( plugin_id, source, move |plugin_id, batch_id, messages_ptr, messages_len| { @@ -256,9 +361,10 @@ impl SourceContainer { shutdown_rx, pending_batch, consecutive_nacks, - BatchPolicy::default(), + batch_policy, ) .await; + stop_notification.report(reason); }); self.shutdown = Some(shutdown_tx); @@ -266,6 +372,15 @@ impl SourceContainer { 0 } + #[doc(hidden)] + pub fn register_stop_callback(&mut self, callback: SourceStoppedCallback) -> i32 { + if self.source.is_none() || self.task.is_some() { + return -1; + } + self.stop_callback = Some(callback); + 0 + } + #[doc(hidden)] pub fn complete_batch(&self, batch_id: u64, result: u8) -> i32 { let Some(source) = self.source.as_ref() else { @@ -280,12 +395,39 @@ impl SourceContainer { batch_id, result, self.id, - MAX_CONSECUTIVE_NACKS, + self.batch_policy.max_consecutive_nacks, ) } } async fn handle_messages( + plugin_id: u32, + source: Arc, + callback: F, + shutdown: watch::Receiver<()>, + pending_batch: Arc>>, + consecutive_nacks: Arc, + policy: BatchPolicy, +) -> SourceStopReason +where + T: Source, + F: Fn(u32, u64, *const u8, usize) -> i32, +{ + handle_messages_with_serializer( + plugin_id, + source, + callback, + shutdown, + pending_batch, + consecutive_nacks, + policy, + postcard::to_allocvec, + ) + .await +} + +#[allow(clippy::too_many_arguments)] +async fn handle_messages_with_serializer( plugin_id: u32, source: Arc, callback: F, @@ -293,16 +435,19 @@ async fn handle_messages( pending_batch: Arc>>, consecutive_nacks: Arc, policy: BatchPolicy, -) where + serialize: S, +) -> SourceStopReason +where T: Source, F: Fn(u32, u64, *const u8, usize) -> i32, + S: Fn(&ProducedMessages) -> Result, postcard::Error>, { let mut batch_id = 1u64; loop { tokio::select! { _ = shutdown.changed() => { info!("Shutting down source connector with ID: {plugin_id}"); - break; + break SourceStopReason::Shutdown; } messages = source.poll() => { let messages = match messages { @@ -313,24 +458,24 @@ async fn handle_messages( } }; - let messages = match postcard::to_allocvec(&messages) { + let messages = match serialize(&messages) { Ok(messages) => messages, Err(err) => { error!("Failed to serialize messages for source connector with ID: {plugin_id}. {err}"); - if matches!( - apply_batch_result( + let completion = apply_batch_result( &source, &consecutive_nacks, SourceBatchResult::Nack, plugin_id, policy.max_consecutive_nacks, ) - .await, - BatchCompletion::Stop - ) { - break; + .await; + if let BatchCompletion::Stop(reason) = completion { + break reason; + } + if sleep_after_nack(&consecutive_nacks, policy, &mut shutdown).await { + break SourceStopReason::Shutdown; } - sleep_after_nack(&consecutive_nacks, policy).await; continue; } }; @@ -352,7 +497,7 @@ async fn handle_messages( if clear_pending_batch(&pending_batch, batch_id, plugin_id) != PendingBatchClearResult::Cleared { - break; + break SourceStopReason::Unexpected; } let completion = apply_batch_result( &source, @@ -362,81 +507,71 @@ async fn handle_messages( policy.max_consecutive_nacks, ) .await; - if matches!(completion, BatchCompletion::Stop) { - break; + if let BatchCompletion::Stop(reason) = completion { + break reason; + } + if sleep_after_nack(&consecutive_nacks, policy, &mut shutdown).await { + break SourceStopReason::Shutdown; } - sleep_after_nack(&consecutive_nacks, policy).await; batch_id += 1; continue; } let mut result_receiver = Box::pin(result_receiver); - let (completion, shutting_down) = tokio::select! { - biased; - result = &mut result_receiver => { - (result.unwrap_or(BatchCompletion::Stop), false) - }, - _ = shutdown.changed() => { - let completion = match clear_pending_batch( - &pending_batch, - batch_id, - plugin_id, - ) { - PendingBatchClearResult::Cleared => apply_batch_result( - &source, - &consecutive_nacks, - SourceBatchResult::Nack, - plugin_id, - policy.max_consecutive_nacks, - ).await, - PendingBatchClearResult::ResultReceived - | PendingBatchClearResult::Unavailable => result_receiver - .as_mut() - .await - .unwrap_or(BatchCompletion::Stop), - }; - (completion, true) - }, - _ = tokio::time::sleep(policy.result_timeout) => { - warn!( - "Timed out waiting for batch result for source connector with ID: {plugin_id}, batch ID: {batch_id}" - ); - let completion = match clear_pending_batch( - &pending_batch, - batch_id, - plugin_id, - ) { - PendingBatchClearResult::Cleared => apply_batch_result( - &source, - &consecutive_nacks, - SourceBatchResult::Nack, + let mut timeout_warned = false; + let (completion, shutting_down) = loop { + let outcome = tokio::select! { + biased; + result = &mut result_receiver => { + (result.unwrap_or(BatchCompletion::Stop(SourceStopReason::ResultChannelClosed)), false) + }, + _ = shutdown.changed() => { + let completion = match clear_pending_batch( + &pending_batch, + batch_id, plugin_id, - policy.max_consecutive_nacks, - ).await, - PendingBatchClearResult::ResultReceived - | PendingBatchClearResult::Unavailable => result_receiver - .as_mut() - .await - .unwrap_or(BatchCompletion::Stop), - }; - (completion, false) - } + ) { + PendingBatchClearResult::Cleared => apply_batch_result( + &source, + &consecutive_nacks, + SourceBatchResult::Nack, + plugin_id, + policy.max_consecutive_nacks, + ).await, + PendingBatchClearResult::ResultReceived + | PendingBatchClearResult::Unavailable => result_receiver + .as_mut() + .await + .unwrap_or(BatchCompletion::Stop(SourceStopReason::ResultChannelClosed)), + }; + (completion, true) + }, + _ = tokio::time::sleep(policy.result_timeout), if !timeout_warned => { + warn!( + "Still waiting for batch result for source connector with ID: {plugin_id}, batch ID: {batch_id}" + ); + timeout_warned = true; + continue; + } + }; + break outcome; }; - if matches!(completion, BatchCompletion::Stop) { - break; + if let BatchCompletion::Stop(reason) = completion { + break reason; } if shutting_down { info!("Shutting down source connector with ID: {plugin_id}"); - break; + break SourceStopReason::Shutdown; } if matches!( completion, BatchCompletion::Applied(SourceBatchResult::Nack) - ) { - sleep_after_nack(&consecutive_nacks, policy).await; + ) && sleep_after_nack(&consecutive_nacks, policy, &mut shutdown).await + { + break SourceStopReason::Shutdown; } batch_id += 1; } @@ -451,7 +586,7 @@ fn complete_pending_batch( batch_id: u64, result_code: u8, plugin_id: u32, - max_consecutive_nacks: u32, + max_consecutive_nacks: Option, ) -> i32 where T: Source + 'static, @@ -487,7 +622,7 @@ where return -1; } - if invalid_result || matches!(completion, BatchCompletion::Stop) { + if invalid_result || matches!(completion, BatchCompletion::Stop(_)) { -1 } else { 0 @@ -578,11 +713,11 @@ async fn apply_batch_result( consecutive_nacks: &AtomicU32, result: SourceBatchResult, plugin_id: u32, - max_consecutive_nacks: u32, + max_consecutive_nacks: Option, ) -> BatchCompletion { if let Err(err) = source.on_batch_result(result).await { error!("Failed to process {result:?} for source connector with ID: {plugin_id}. {err}"); - return BatchCompletion::Stop; + return BatchCompletion::Stop(SourceStopReason::BatchResultError); } let consecutive_nacks = match result { @@ -590,20 +725,31 @@ async fn apply_batch_result( consecutive_nacks.store(0, Ordering::Relaxed); 0 } - SourceBatchResult::Nack => consecutive_nacks.fetch_add(1, Ordering::Relaxed) + 1, + SourceBatchResult::Nack => consecutive_nacks + .try_update(Ordering::Relaxed, Ordering::Relaxed, |count| { + Some(count.saturating_add(1)) + }) + .map_or(u32::MAX, |previous| previous.saturating_add(1)), }; - if consecutive_nacks >= max_consecutive_nacks { + if max_consecutive_nacks.is_some_and(|limit| consecutive_nacks >= limit.get()) { error!( "Stopping source connector with ID: {plugin_id} after {consecutive_nacks} consecutive NACKs" ); - return BatchCompletion::Stop; + return BatchCompletion::Stop(SourceStopReason::NackLimit); } BatchCompletion::Applied(result) } -async fn sleep_after_nack(consecutive_nacks: &AtomicU32, policy: BatchPolicy) { - tokio::time::sleep(nack_retry_delay(consecutive_nacks, policy)).await; +async fn sleep_after_nack( + consecutive_nacks: &AtomicU32, + policy: BatchPolicy, + shutdown: &mut watch::Receiver<()>, +) -> bool { + tokio::select! { + _ = shutdown.changed() => true, + _ = tokio::time::sleep(nack_retry_delay(consecutive_nacks, policy)) => false, + } } fn nack_retry_delay(consecutive_nacks: &AtomicU32, policy: BatchPolicy) -> Duration { @@ -627,6 +773,7 @@ macro_rules! source_connector { use std::sync::LazyLock; use $crate::LogCallback; use $crate::source::SendCallback; + use $crate::source::SourceStoppedCallback; use $crate::source::SourceContainer; static INSTANCES: LazyLock>> = @@ -685,6 +832,18 @@ macro_rules! source_connector { instance.handle(callback) } + #[cfg(not(test))] + #[unsafe(no_mangle)] + extern "C" fn iggy_source_register_stop_callback( + id: u32, + callback: SourceStoppedCallback, + ) -> i32 { + let Some(mut instance) = INSTANCES.get_mut(&id) else { + return -1; + }; + instance.register_stop_callback(callback) + } + #[cfg(not(test))] #[unsafe(no_mangle)] extern "C" fn iggy_source_batch_result(id: u32, batch_id: u64, result: u8) -> i32 { @@ -727,23 +886,298 @@ mod tests { use crate::{ProducedMessages, Schema}; use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering}; use std::time::Duration; - use tokio::sync::mpsc; + use tokio::sync::{Notify, mpsc}; + + static STOPPED_PLUGIN: AtomicU32 = AtomicU32::new(0); + static STOP_REASON: AtomicU32 = AtomicU32::new(0); + static STOP_COUNT: AtomicUsize = AtomicUsize::new(0); + static CLOSE_STOP_COUNT: AtomicUsize = AtomicUsize::new(0); + static SHUTDOWN_STOP_COUNT: AtomicUsize = AtomicUsize::new(0); + static CLOSING_NACK_STOP_COUNT: AtomicUsize = AtomicUsize::new(0); + static PANIC_STOP_COUNT: AtomicUsize = AtomicUsize::new(0); + static PANIC_STOP_REASON: AtomicU32 = AtomicU32::new(u32::MAX); + + extern "C" fn record_stop(plugin_id: u32, reason: u8) { + STOPPED_PLUGIN.store(plugin_id, Ordering::SeqCst); + STOP_REASON.store(u32::from(reason), Ordering::SeqCst); + STOP_COUNT.fetch_add(1, Ordering::SeqCst); + } + + extern "C" fn reject_batch(_: u32, _: u64, _: *const u8, _: usize) -> i32 { + -1 + } + + extern "C" fn accept_batch(_: u32, _: u64, _: *const u8, _: usize) -> i32 { + 0 + } + + extern "C" fn record_close_stop(_: u32, _: u8) { + CLOSE_STOP_COUNT.fetch_add(1, Ordering::SeqCst); + } + + extern "C" fn record_shutdown_stop(_: u32, _: u8) { + SHUTDOWN_STOP_COUNT.fetch_add(1, Ordering::SeqCst); + } + + extern "C" fn record_closing_nack_stop(_: u32, _: u8) { + CLOSING_NACK_STOP_COUNT.fetch_add(1, Ordering::SeqCst); + } + + extern "C" fn record_panic_stop(_: u32, reason: u8) { + PANIC_STOP_REASON.store(u32::from(reason), Ordering::SeqCst); + PANIC_STOP_COUNT.fetch_add(1, Ordering::SeqCst); + } + + #[test] + fn given_unopened_source_when_registering_stop_callback_should_reject() { + let mut container = SourceContainer::::new(39); + assert_eq!(container.register_stop_callback(record_panic_stop), -1); + } + + #[test] + fn given_nack_limit_when_polling_stops_should_notify_runtime() { + STOPPED_PLUGIN.store(0, Ordering::SeqCst); + STOP_REASON.store(0, Ordering::SeqCst); + STOP_COUNT.store(0, Ordering::SeqCst); + let mut container = SourceContainer::::new(32); + let config = b"{}"; + assert_eq!( + unsafe { + container.open( + 32, + config.as_ptr(), + config.len(), + std::ptr::null(), + 0, + ignore_log, + |_, _: serde_json::Value, _| TestSource { + batch_policy: test_policy().with_max_consecutive_nacks(NonZeroU32::new(2)), + ..TestSource::default() + }, + ) + }, + 0 + ); + assert_eq!(container.register_stop_callback(record_stop), 0); + assert_eq!(unsafe { container.handle(reject_batch) }, 0); + get_runtime().block_on(async { + tokio::time::timeout(Duration::from_secs(2), async { + while STOPPED_PLUGIN.load(Ordering::SeqCst) != 32 { + tokio::time::sleep(Duration::from_millis(1)).await; + } + }) + .await + .expect("source did not report its terminal NACK limit"); + }); + let source = Arc::clone(container.source.as_ref().expect("source should be open")); + assert_eq!(source.polls.load(Ordering::SeqCst), 2); + drop(source); + assert_eq!(unsafe { container.close() }, 0); + assert_eq!(STOPPED_PLUGIN.load(Ordering::SeqCst), 32); + assert_eq!( + STOP_REASON.load(Ordering::SeqCst), + SourceStopReason::NackLimit as u32 + ); + assert_eq!(STOP_COUNT.load(Ordering::SeqCst), 1); + assert_eq!(container.consecutive_nacks.load(Ordering::SeqCst), 2); + + drop(StopNotification { + plugin_id: 33, + callback: Some(record_stop), + closing: Arc::new(AtomicBool::new(true)), + }); + assert_eq!(STOP_COUNT.load(Ordering::SeqCst), 1); + } + + #[test] + fn given_running_source_when_closed_should_not_notify_runtime() { + CLOSE_STOP_COUNT.store(0, Ordering::SeqCst); + let mut container = SourceContainer::::new(35); + let config = b"{}"; + assert_eq!( + unsafe { + container.open( + 35, + config.as_ptr(), + config.len(), + std::ptr::null(), + 0, + ignore_log, + |_, _: serde_json::Value, _| TestSource::default(), + ) + }, + 0 + ); + assert_eq!(container.register_stop_callback(record_close_stop), 0); + assert_eq!(unsafe { container.handle(accept_batch) }, 0); + let source = Arc::clone(container.source.as_ref().expect("source should be open")); + get_runtime().block_on(async { + tokio::time::timeout(Duration::from_secs(1), async { + while source.polls.load(Ordering::SeqCst) == 0 { + tokio::time::sleep(Duration::from_millis(1)).await; + } + }) + .await + .expect("source did not begin polling"); + }); + drop(source); + + assert_eq!(unsafe { container.close() }, 0); + assert_eq!(CLOSE_STOP_COUNT.load(Ordering::SeqCst), 0); + } + + #[test] + fn given_shutdown_reason_without_close_should_not_notify_runtime() { + SHUTDOWN_STOP_COUNT.store(0, Ordering::SeqCst); + StopNotification { + plugin_id: 36, + callback: Some(record_shutdown_stop), + closing: Arc::new(AtomicBool::new(false)), + } + .report(SourceStopReason::Shutdown); + assert_eq!(SHUTDOWN_STOP_COUNT.load(Ordering::SeqCst), 0); + } + + #[test] + fn given_nack_limit_during_close_should_not_notify_runtime() { + CLOSING_NACK_STOP_COUNT.store(0, Ordering::SeqCst); + let nack_entered = Arc::new(Notify::new()); + let finish_nack = Arc::new(Notify::new()); + let mut container = SourceContainer::::new(37); + let config = b"{}"; + assert_eq!( + unsafe { + container.open( + 37, + config.as_ptr(), + config.len(), + std::ptr::null(), + 0, + ignore_log, + |_, _: serde_json::Value, _| TestSource { + batch_policy: test_policy().with_max_consecutive_nacks(NonZeroU32::new(1)), + nack_pause: Some((Arc::clone(&nack_entered), Arc::clone(&finish_nack))), + ..TestSource::default() + }, + ) + }, + 0 + ); + assert_eq!( + container.register_stop_callback(record_closing_nack_stop), + 0 + ); + assert_eq!(unsafe { container.handle(reject_batch) }, 0); + get_runtime().block_on(async { + tokio::time::timeout(Duration::from_secs(2), nack_entered.notified()) + .await + .expect("source did not enter NACK handling"); + }); + + let closing = Arc::clone(&container.closing); + let close_thread = std::thread::spawn(move || unsafe { container.close() }); + get_runtime().block_on(async { + tokio::time::timeout(Duration::from_secs(2), async { + while !closing.load(Ordering::Acquire) { + tokio::task::yield_now().await; + } + }) + .await + .expect("close did not mark the source as closing"); + }); + finish_nack.notify_one(); + assert_eq!( + close_thread.join().expect("source close thread panicked"), + 0 + ); + assert_eq!(CLOSING_NACK_STOP_COUNT.load(Ordering::SeqCst), 0); + } + + #[test] + fn given_poll_panic_should_notify_runtime_once_as_unexpected() { + PANIC_STOP_COUNT.store(0, Ordering::SeqCst); + PANIC_STOP_REASON.store(u32::MAX, Ordering::SeqCst); + let mut container = SourceContainer::::new(38); + let config = b"{}"; + assert_eq!( + unsafe { + container.open( + 38, + config.as_ptr(), + config.len(), + std::ptr::null(), + 0, + ignore_log, + |_, _: serde_json::Value, _| TestSource { + panic_on_poll: AtomicBool::new(true), + ..TestSource::default() + }, + ) + }, + 0 + ); + assert_eq!(container.register_stop_callback(record_panic_stop), 0); + assert_eq!(unsafe { container.handle(accept_batch) }, 0); + assert_eq!(container.register_stop_callback(record_panic_stop), -1); + get_runtime().block_on(async { + tokio::time::timeout(Duration::from_secs(1), async { + while PANIC_STOP_COUNT.load(Ordering::SeqCst) == 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("poll panic should report a stop"); + }); + assert_eq!( + PANIC_STOP_REASON.load(Ordering::SeqCst), + SourceStopReason::Unexpected as u32 + ); + assert_eq!(unsafe { container.close() }, 0); + assert_eq!(PANIC_STOP_COUNT.load(Ordering::SeqCst), 1); + } + + #[test] + fn source_stop_reasons_should_keep_their_wire_bytes() { + for (reason, byte) in [ + (SourceStopReason::Unexpected, 0), + (SourceStopReason::NackLimit, 1), + (SourceStopReason::BatchResultError, 2), + (SourceStopReason::ResultChannelClosed, 3), + (SourceStopReason::RegistrationFailed, 4), + (SourceStopReason::HandlerFailed, 5), + (SourceStopReason::Shutdown, 6), + ] { + assert_eq!(reason as u8, byte); + assert_eq!(SourceStopReason::from(byte), reason); + } + assert_eq!( + SourceStopReason::from(u8::MAX), + SourceStopReason::Unexpected + ); + } fn test_policy() -> BatchPolicy { BatchPolicy { result_timeout: Duration::from_millis(500), nack_retry_delay: Duration::from_millis(1), max_nack_retry_delay: Duration::from_millis(2), - max_consecutive_nacks: MAX_CONSECUTIVE_NACKS, + max_consecutive_nacks: test_nack_limit(), } } + fn test_nack_limit() -> Option { + NonZeroU32::new(MAX_CONSECUTIVE_NACKS) + } + #[derive(Debug, Default)] struct TestSource { polls: AtomicUsize, results: Mutex>, fail_batch_result: AtomicBool, + panic_on_poll: AtomicBool, batch_result_delay: Duration, + batch_policy: BatchPolicy, + nack_pause: Option<(Arc, Arc)>, } #[async_trait::async_trait] @@ -752,8 +1186,13 @@ mod tests { Ok(()) } + fn batch_policy(&self) -> BatchPolicy { + self.batch_policy + } + async fn poll(&self) -> Result { self.polls.fetch_add(1, Ordering::SeqCst); + assert!(!self.panic_on_poll.load(Ordering::SeqCst)); Ok(ProducedMessages { schema: Schema::Raw, messages: Vec::new(), @@ -762,6 +1201,12 @@ mod tests { } async fn on_batch_result(&self, result: SourceBatchResult) -> Result<(), crate::Error> { + if result == SourceBatchResult::Nack + && let Some((entered, finish)) = &self.nack_pause + { + entered.notify_one(); + finish.notified().await; + } if !self.batch_result_delay.is_zero() { tokio::time::sleep(self.batch_result_delay).await; } @@ -782,6 +1227,83 @@ mod tests { } } + extern "C" fn ignore_log( + _level: u8, + _target_ptr: *const u8, + _target_len: usize, + _message_ptr: *const u8, + _message_len: usize, + ) { + } + + #[test] + fn given_custom_source_policy_when_opened_should_capture_nack_limit() { + let mut container = SourceContainer::::new(31); + let config = b"{}"; + let result = unsafe { + container.open( + 31, + config.as_ptr(), + config.len(), + std::ptr::null(), + 0, + ignore_log, + |_, _: serde_json::Value, _| TestSource { + batch_policy: BatchPolicy::default().with_max_consecutive_nacks(None), + ..TestSource::default() + }, + ) + }; + + assert_eq!(result, 0); + assert_eq!(container.batch_policy.max_consecutive_nacks(), None); + assert_eq!(unsafe { container.close() }, 0); + } + + #[test] + fn given_no_nack_limit_when_completing_six_nacks_should_not_stop() { + let mut container = SourceContainer::::new(34); + let config = b"{}"; + assert_eq!( + unsafe { + container.open( + 34, + config.as_ptr(), + config.len(), + std::ptr::null(), + 0, + ignore_log, + |_, _: serde_json::Value, _| TestSource { + batch_policy: test_policy().with_max_consecutive_nacks(None), + ..TestSource::default() + }, + ) + }, + 0 + ); + + for batch_id in 1..=6 { + let (result_sender, result_receiver) = oneshot::channel(); + *lock_pending_batch(&container.pending_batch) = Some(PendingBatch { + id: batch_id, + result_sender, + result_received: false, + }); + assert_eq!( + container.complete_batch(batch_id, SourceBatchResult::Nack as u8), + 0 + ); + assert_eq!( + get_runtime() + .block_on(result_receiver) + .expect("batch completion should be delivered"), + BatchCompletion::Applied(SourceBatchResult::Nack) + ); + } + assert_eq!(container.consecutive_nacks.load(Ordering::Relaxed), 6); + assert_eq!(unsafe { container.close() }, 0); + } + async fn complete_test_batch( pending_batch: Arc>>, source: Arc, @@ -798,7 +1320,7 @@ mod tests { batch_id, result as u8, plugin_id, - MAX_CONSECUTIVE_NACKS, + test_nack_limit(), ) }) .await @@ -962,7 +1484,7 @@ mod tests { 42, SourceBatchResult::Ack as u8, 11, - MAX_CONSECUTIVE_NACKS, + test_nack_limit(), ), -1 ); @@ -974,7 +1496,7 @@ mod tests { 41, SourceBatchResult::Ack as u8, 11, - MAX_CONSECUTIVE_NACKS, + test_nack_limit(), ), 0 ); @@ -1008,7 +1530,7 @@ mod tests { 51, 99, 17, - MAX_CONSECUTIVE_NACKS, + test_nack_limit(), ), -1 ); @@ -1030,7 +1552,7 @@ mod tests { } #[test] - fn given_result_timeout_should_nack_and_poll_next_batch() { + fn given_result_timeout_should_keep_batch_pending_until_result() { let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); runtime.block_on(async { let source = Arc::new(TestSource::default()); @@ -1059,18 +1581,41 @@ mod tests { )); assert_eq!(batch_receiver.recv().await, Some(1)); - let next_batch_id = tokio::time::timeout(Duration::from_secs(1), batch_receiver.recv()) - .await - .expect("source did not poll after timed-out batch") - .expect("batch channel closed"); - assert_eq!(next_batch_id, 2); + assert!( + tokio::time::timeout(Duration::from_millis(50), batch_receiver.recv()) + .await + .is_err(), + "source must not poll again without a runtime result" + ); + assert_eq!(source.polls.load(Ordering::Relaxed), 1); + assert_eq!(consecutive_nacks.load(Ordering::Relaxed), 0); + assert_eq!( + lock_pending_batch(&pending_batch) + .as_ref() + .map(|batch| batch.id), + Some(1) + ); assert_eq!( *source .results .lock() .unwrap_or_else(PoisonError::into_inner), - vec![SourceBatchResult::Nack] + Vec::::new() + ); + + assert_eq!( + complete_test_batch( + Arc::clone(&pending_batch), + Arc::clone(&source), + Arc::clone(&consecutive_nacks), + 1, + SourceBatchResult::Ack, + 19, + ) + .await, + 0 ); + assert_eq!(batch_receiver.recv().await, Some(2)); task.abort(); let _ = task.await; @@ -1213,7 +1758,7 @@ mod tests { &consecutive_nacks, SourceBatchResult::Nack, 23, - MAX_CONSECUTIVE_NACKS, + test_nack_limit(), ) .await, BatchCompletion::Applied(SourceBatchResult::Nack) @@ -1227,10 +1772,10 @@ mod tests { &consecutive_nacks, SourceBatchResult::Nack, 23, - MAX_CONSECUTIVE_NACKS, + test_nack_limit(), ) .await, - BatchCompletion::Stop + BatchCompletion::Stop(SourceStopReason::NackLimit) ); assert_eq!( consecutive_nacks.load(Ordering::Relaxed), @@ -1239,6 +1784,295 @@ mod tests { }); } + #[test] + fn given_no_nack_limit_when_repeated_nacks_should_keep_polling() { + let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); + runtime.block_on(async { + let source = Arc::new(TestSource::default()); + let consecutive_nacks = AtomicU32::new(0); + + for expected_count in 1..=MAX_CONSECUTIVE_NACKS + 1 { + assert_eq!( + apply_batch_result( + &source, + &consecutive_nacks, + SourceBatchResult::Nack, + 23, + None, + ) + .await, + BatchCompletion::Applied(SourceBatchResult::Nack) + ); + assert_eq!(consecutive_nacks.load(Ordering::Relaxed), expected_count); + } + }); + } + + #[test] + fn given_result_hook_failure_should_stop_even_without_nack_limit() { + let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); + runtime.block_on(async { + let source = Arc::new(TestSource { + fail_batch_result: AtomicBool::new(true), + ..TestSource::default() + }); + let consecutive_nacks = AtomicU32::new(0); + + assert_eq!( + apply_batch_result( + &source, + &consecutive_nacks, + SourceBatchResult::Nack, + 23, + None, + ) + .await, + BatchCompletion::Stop(SourceStopReason::BatchResultError) + ); + assert_eq!(consecutive_nacks.load(Ordering::Relaxed), 0); + }); + } + + #[test] + fn given_no_nack_limit_when_callback_rejects_batches_should_poll_past_default_limit() { + let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); + runtime.block_on(async { + let source = Arc::new(TestSource::default()); + let pending_batch = Arc::new(Mutex::new(None)); + let consecutive_nacks = Arc::new(AtomicU32::new(0)); + let (shutdown_sender, shutdown_receiver) = watch::channel(()); + let (batch_sender, mut batch_receiver) = mpsc::unbounded_channel(); + let policy = test_policy().with_max_consecutive_nacks(None); + + let task = tokio::spawn(handle_messages( + 23, + Arc::clone(&source), + move |_, batch_id, _, _| { + batch_sender + .send(batch_id) + .expect("batch receiver should remain open"); + -1 + }, + shutdown_receiver, + pending_batch, + Arc::clone(&consecutive_nacks), + policy, + )); + + for expected_batch_id in 1..=u64::from(MAX_CONSECUTIVE_NACKS + 1) { + let batch_id = tokio::time::timeout(Duration::from_secs(1), batch_receiver.recv()) + .await + .expect("source stopped polling after repeated NACKs"); + assert_eq!(batch_id, Some(expected_batch_id)); + } + + shutdown_sender + .send(()) + .expect("source task should still be running"); + task.await.expect("source task should shut down cleanly"); + assert!( + consecutive_nacks.load(Ordering::Relaxed) >= MAX_CONSECUTIVE_NACKS, + "the source should have handled more NACKs than the default limit" + ); + }); + } + + #[test] + fn given_unlimited_nacks_when_shutdown_during_backoff_should_stop_promptly() { + let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); + runtime.block_on(async { + let source = Arc::new(TestSource::default()); + let pending_batch = Arc::new(Mutex::new(None)); + let consecutive_nacks = Arc::new(AtomicU32::new(0)); + let (shutdown_sender, shutdown_receiver) = watch::channel(()); + let (batch_sender, mut batch_receiver) = mpsc::unbounded_channel(); + let policy = BatchPolicy { + nack_retry_delay: Duration::from_secs(10), + max_nack_retry_delay: Duration::from_secs(10), + ..test_policy().with_max_consecutive_nacks(None) + }; + + let task = tokio::spawn(handle_messages( + 23, + source, + move |_, batch_id, _, _| { + batch_sender + .send(batch_id) + .expect("batch receiver should remain open"); + -1 + }, + shutdown_receiver, + pending_batch, + consecutive_nacks, + policy, + )); + + let batch_id = tokio::time::timeout(Duration::from_secs(1), batch_receiver.recv()) + .await + .expect("source should submit a batch"); + assert_eq!(batch_id, Some(1)); + shutdown_sender + .send(()) + .expect("source task should still be running"); + tokio::time::timeout(Duration::from_millis(500), task) + .await + .expect("shutdown should interrupt NACK backoff") + .expect("source task should shut down cleanly"); + }); + } + + #[test] + fn given_runtime_nack_when_shutdown_during_backoff_should_stop_promptly() { + let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); + runtime.block_on(async { + let source = Arc::new(TestSource::default()); + let pending_batch = Arc::new(Mutex::new(None)); + let consecutive_nacks = Arc::new(AtomicU32::new(0)); + let (shutdown_sender, shutdown_receiver) = watch::channel(()); + let (batch_sender, mut batch_receiver) = mpsc::unbounded_channel(); + let policy = BatchPolicy { + nack_retry_delay: Duration::from_secs(10), + max_nack_retry_delay: Duration::from_secs(10), + ..test_policy().with_max_consecutive_nacks(None) + }; + + let task = tokio::spawn(handle_messages( + 24, + Arc::clone(&source), + move |_, batch_id, _, _| { + batch_sender + .send(batch_id) + .expect("batch receiver should remain open"); + 0 + }, + shutdown_receiver, + Arc::clone(&pending_batch), + Arc::clone(&consecutive_nacks), + policy, + )); + + assert_eq!(batch_receiver.recv().await, Some(1)); + assert_eq!( + complete_test_batch( + Arc::clone(&pending_batch), + Arc::clone(&source), + Arc::clone(&consecutive_nacks), + 1, + SourceBatchResult::Nack, + 24, + ) + .await, + 0 + ); + shutdown_sender + .send(()) + .expect("source task should still be running"); + assert_eq!( + tokio::time::timeout(Duration::from_millis(500), task) + .await + .expect("shutdown should interrupt result-NACK backoff") + .expect("source task should complete"), + SourceStopReason::Shutdown + ); + assert_eq!(source.polls.load(Ordering::SeqCst), 1); + }); + } + + #[test] + fn given_nack_backoff_when_shutdown_should_interrupt_shared_delay() { + let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); + runtime.block_on(async { + let consecutive_nacks = AtomicU32::new(1); + let policy = BatchPolicy { + nack_retry_delay: Duration::from_secs(10), + max_nack_retry_delay: Duration::from_secs(10), + ..test_policy() + }; + let (shutdown_sender, mut shutdown_receiver) = watch::channel(()); + let waiting = tokio::spawn(async move { + sleep_after_nack(&consecutive_nacks, policy, &mut shutdown_receiver).await + }); + shutdown_sender + .send(()) + .expect("delay should still be active"); + assert!( + tokio::time::timeout(Duration::from_millis(500), waiting) + .await + .expect("shutdown should interrupt shared NACK delay") + .expect("delay task should complete") + ); + }); + } + + #[test] + fn given_serialization_failure_when_shutdown_during_backoff_should_stop_promptly() { + let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); + runtime.block_on(async { + let source = Arc::new(TestSource::default()); + let pending_batch = Arc::new(Mutex::new(None)); + let consecutive_nacks = Arc::new(AtomicU32::new(0)); + let (shutdown_sender, shutdown_receiver) = watch::channel(()); + let policy = BatchPolicy { + nack_retry_delay: Duration::from_secs(10), + max_nack_retry_delay: Duration::from_secs(10), + ..test_policy().with_max_consecutive_nacks(None) + }; + + let source_for_task = Arc::clone(&source); + let nacks_for_task = Arc::clone(&consecutive_nacks); + let task = tokio::spawn(handle_messages_with_serializer( + 25, + source_for_task, + |_, _, _, _| panic!("serialization failure must prevent the callback"), + shutdown_receiver, + pending_batch, + nacks_for_task, + policy, + |_| Err(postcard::Error::SerdeSerCustom), + )); + + tokio::time::timeout(Duration::from_secs(1), async { + while consecutive_nacks.load(Ordering::Relaxed) == 0 { + tokio::task::yield_now().await; + } + }) + .await + .expect("serialization failure should NACK the batch"); + shutdown_sender + .send(()) + .expect("source task should still be running"); + assert_eq!( + tokio::time::timeout(Duration::from_millis(500), task) + .await + .expect("shutdown should interrupt serialization-failure backoff") + .expect("source task should complete"), + SourceStopReason::Shutdown + ); + assert_eq!(source.polls.load(Ordering::SeqCst), 1); + }); + } + + #[test] + fn given_custom_nack_limit_when_reached_should_stop() { + let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); + runtime.block_on(async { + let source = Arc::new(TestSource::default()); + let consecutive_nacks = AtomicU32::new(1); + + assert_eq!( + apply_batch_result( + &source, + &consecutive_nacks, + SourceBatchResult::Nack, + 23, + NonZeroU32::new(2), + ) + .await, + BatchCompletion::Stop(SourceStopReason::NackLimit) + ); + }); + } + #[test] fn given_ack_after_nack_should_reset_consecutive_nacks() { let runtime = tokio::runtime::Runtime::new().expect("failed to create test runtime"); @@ -1252,7 +2086,7 @@ mod tests { &consecutive_nacks, SourceBatchResult::Ack, 29, - MAX_CONSECUTIVE_NACKS, + test_nack_limit(), ) .await, BatchCompletion::Applied(SourceBatchResult::Ack) diff --git a/core/connectors/sinks/influxdb_sink/dependencies.md b/core/connectors/sinks/influxdb_sink/dependencies.md index b056d46e4d..c718745985 100644 --- a/core/connectors/sinks/influxdb_sink/dependencies.md +++ b/core/connectors/sinks/influxdb_sink/dependencies.md @@ -41,7 +41,7 @@ to `cargo tree -p iggy_connector_influxdb_sink` for the full graph. | `base64` | `^0.22.1` | MIT / Apache-2.0 | Encodes raw message payloads as base64 when `payload_format = "base64"` is configured. Uses `Engine::encode_string` to write directly into the line-protocol output buffer with no intermediate allocation. | | `bytes` | `^1.11.1` | MIT | Zero-copy `Bytes::from(body.into_bytes())` converts the line-protocol string body into the request body sent to InfluxDB's write endpoint without an extra copy. | | `iggy_common` | `^0.11.0` | Apache-2.0 | Shared Iggy types: `serde_secret` for safe token serialisation in config structs. | -| `iggy_connector_sdk` | `^0.5.0` | Apache-2.0 | Core connector abstractions: `Sink` trait, `ConsumedMessage`, `MessagesMetadata`, `TopicMetadata`, `Error`, retry/circuit-breaker utilities, and the `sink_connector!` registration macro. | +| `iggy_connector_sdk` | `^0.5.1-edge.2` | Apache-2.0 | Core connector abstractions: `Sink` trait, `ConsumedMessage`, `MessagesMetadata`, `TopicMetadata`, `Error`, retry/circuit-breaker utilities, and the `sink_connector!` registration macro. | | `reqwest` | `^0.13.3` | MIT / Apache-2.0 | Async HTTP client used to POST line-protocol batches to `/api/v2/write` (V2) and `/api/v3/write_lp` (V3). | | `reqwest-middleware` | `^0.5.1` | MIT | Middleware wrapper around `reqwest::Client` that attaches the retry and tracing layers built by `iggy_connector_sdk::retry::build_retry_client`. | | `secrecy` | `^0.10` | MIT / Apache-2.0 | `SecretString` / `SecretBox` wrappers that prevent accidental logging of API tokens in config structs and the `auth_header` field. | diff --git a/core/connectors/sources/http_source/README.md b/core/connectors/sources/http_source/README.md index 4ce82ea377..035943764a 100644 --- a/core/connectors/sources/http_source/README.md +++ b/core/connectors/sources/http_source/README.md @@ -13,7 +13,9 @@ Loss windows, explicitly: 1. **Process crash** between HTTP 200 and the producer send. Buffered messages are volatile. 2. **Producer failure.** Narrowed by #3855, which gave sources a delivery result. The runtime NACKs a batch it could not send, or whose state it could not persist, and neither is abandoned here: a message batch is held and replayed on the next `poll()`, and a state-only batch re-arms the flush so the mutation is handed out again. Answering 200 already told the sender this gateway owns the event, and the only honest way to shed load is the 429 handlers return once the bridge fills. 3. **Shutdown.** Narrowed by #3321, which closes the plugin before tearing down the forwarding channel so in-flight batches drain. Messages still in the bridge when the poll task stops are lost; the connector logs the count and increments `http_source_dropped_on_close_total`. -4. **Poll task stopped by the SDK.** A source is stopped after five consecutive NACKs, roughly 1.5s of backoff plus five send rounds, so a broker outage longer than that ends the poll task while the listener keeps accepting. Whatever the bridge holds is then lost on the next restart, and the stop is unobservable: the runtime's forwarding loop stays parked, so the source keeps being counted as running while nothing polls. Tracked in #3941. +4. **Poll task stopped by the SDK.** The HTTP source retries NACKed batches without a fixed stop limit by default, with capped backoff. A configured `max_consecutive_nacks` can still stop it; the runtime reports the stop as `Error` and removes it from the running gauge. Rebuild the plugin with the current SDK to enable stop reporting. + + The listener keeps answering 200 until its bridge fills, even though a stopped source cannot drain it. Restarting closes the old instance and loses its staged batch and buffered messages; `http_source_dropped_on_close_total` counts them. Do not set a NACK limit unless those accepted events can be replayed outside the connector. What mitigates this in practice is the caller: webhook senders such as GitHub, Stripe, and Twilio retry on timeout and 5xx, so the sender side is at-least-once up to the moment this connector returns 200. The connector's job is to make the post-200 window as small and as observable as possible. @@ -21,7 +23,7 @@ What mitigates this in practice is the caller: webhook senders such as GitHub, S 1. A caller times out after the connector enqueued the message but before the 200 reached it, retries, and the payload lands twice. 2. A load balancer or proxy retries a POST after a hiccup downstream of a successful enqueue. -3. A NACK does not mean the batch never landed. The runtime NACKs a batch it sent but whose state it could not persist, and the SDK NACKs on its own result timeout while the send may still have succeeded, so the replay on the next `poll()` re-sends messages that are already on the topic. This is the window the connector itself creates, and it is the price of not dropping post-200 data. +3. A NACK does not mean the batch never landed. The runtime NACKs a batch it sent but whose state it could not persist, so the replay on the next `poll()` may re-send messages already on the topic. This is the window the connector itself creates, and it is the price of not dropping post-200 data. So the guarantee is best-effort in both directions: no silent-loss guarantee and no dedup guarantee. @@ -107,7 +109,13 @@ 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. + +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. + +Set `max_consecutive_nacks` in `[plugin_config]` to a positive integer only if an external replay mechanism makes stopping safe. This does not make a mismatched stream schema valid: `schema` lives under `[[streams]]` and the plugin only receives `[plugin_config]`. + +Rebuild the HTTP source plugin with SDK 0.5.1-edge.2 to get this behavior and the `iggy_source_register_stop_callback` export. An older plugin binary still uses the five-NACK limit and cannot notify the runtime when its poll task stops; the runtime warns when the export is absent. ### Options @@ -122,6 +130,7 @@ The batch is then replayed on every poll and the SDK stops the poll task after f | `max_body_size_bytes` | usize | `1048576` | Request body limit, applied by the handlers rather than an extractor. Routing wins over it: an oversized POST to an unknown or revoked path answers 404 without the body being read. Must match across instances sharing a listener. **Max 64000000**, Iggy's message payload cap, since a body becomes the payload unchanged; a larger value fails `open()`. | | `buffer_capacity` | usize | `10000` | Messages the instance bridge holds. A full bridge answers 429, which since #3855 signals either an arrival burst or a slow Iggy, since the poll loop stalls waiting for the previous batch to be acknowledged. **Max 1000000**; a larger value fails `open()`. | | `max_batch_size` | usize | `500` | Maximum messages a single `poll()` returns. **Max 100000**; a larger value fails `open()`. | +| `max_consecutive_nacks` | positive integer | disabled | Optional SDK breaker limit. Omit to keep retrying NACKed batches and refused state flushes with capped backoff; set only when accepted events can be replayed after a restart. An idle gateway can hit this limit on state-flush NACKs alone. Zero is invalid. | | `include_http_metadata` | bool | `true` | Adds instance, peer address, and receive time as message headers. | | `forward_headers` | array | `[]` | Request headers copied onto the message. Invalid names fail `open()`, as do `Authorization`, `Proxy-Authorization`, and `Cookie`, because forwarding a reusable credential would copy it onto every message and persist it in the log. | | `endpoints` | array | `[]` | Static secret-path endpoints. | @@ -188,7 +197,7 @@ Content-Type: application/json | 413 | Body over `max_body_size_bytes` | `{"error":"payload too large"}` | | 405 | A known path with the wrong method | `{"error":"method not allowed"}` | | 429 | Bridge full | `{"error":"too many requests"}` plus `Retry-After: 1` | -| 503 | `GET /health` when an instance on the listener has stopped polling, having previously polled; a named-path POST whose route changed hands while the body was still arriving; a POST whose instance left the listener mid-request; or a POST whose instance bridge has no receiver | `{"status":"unavailable"}`, `{"error":"route unavailable"}`, `{"error":"instance is closing"}` or `{"error":"service unavailable"}` | +| 503 | `GET /health` when an instance's poll path has stopped or stalled for more than 60 seconds, or its staged batch has kept receiving NACKs for that long; a named-path POST whose route changed hands while the body was still arriving; a POST whose instance left the listener mid-request; or a POST whose instance bridge has no receiver | `{"status":"unavailable"}`, `{"error":"route unavailable"}`, `{"error":"instance is closing"}` or `{"error":"service unavailable"}` | Revoked and expired endpoints both answer 404 rather than 410 or 403 on purpose: a leaked URL must not be usable to confirm that it was once live. The lookup runs before any credential is checked, so anything other than 404 would answer that question for an unauthenticated caller. Error bodies carry no internals; diagnostics live on the admin listener. @@ -198,7 +207,9 @@ A 404 there would tell a conventional client the resource is gone, and it would The named path's 503 covers two conditions and says `route unavailable` for both: the path was withdrawn, or another instance on the same listener took it over. The caller retries either way, and the body deliberately does not say which, since that would describe the listener's topology to a sender that has no use for it. -`GET /health` on the public listener answers 200 only while every instance on it is serving, and 503 otherwise, which is what a load balancer should watch. It is deliberately all rather than any: one address fronts every instance sharing the listener, so a sibling whose poll task has stopped would otherwise keep receiving traffic into a bridge nothing drains. Shedding the healthy siblings costs availability the sender recovers by retrying, where the alternative loses requests already answered 200. +`GET /health` on the public listener answers 200 only while every instance on it is serving and its poll path is healthy, and 503 otherwise, which is what a load balancer should watch. It is deliberately all rather than any: one address fronts every instance sharing the listener, so a sibling whose poll task has stopped or whose staged batch is stuck would otherwise keep receiving traffic into a bridge nothing drains. + +Shedding the healthy siblings costs availability the sender recovers by retrying, where the alternative loses requests already answered 200. An instance that has not polled yet is not counted, which is a different condition from one that has stopped. Between `open()` and its first `poll()` the runtime finishes the sources it has not reached yet and then every sink, in series, and only then starts the poll tasks; counting that window meant restarting one instance took every healthy sibling out of rotation for the length of it, and on boot the length of it depends on connectors that have nothing to do with this one. @@ -294,7 +305,7 @@ The chain below holds end to end. It did not always: the runtime's forwarding ch What closed the gap instead was #3855. The SDK now keeps one batch in flight and will not call `poll()` again until the runtime acknowledges the last one, so a slow Iggy stalls the poll loop directly. The bridge then fills on arrival and the handlers answer 429, which is the coupling that was missing. -That coupling holds only while the runtime answers inside the SDK's batch-result timeout, 30s. Past it the SDK stops waiting, NACKs, and polls again, so the bridge drains and the pressure moves into the runtime's unbounded forwarding channel instead of reaching the sender. Tracked in #3981. +After 30 seconds without a batch result, the SDK warns but continues waiting on the same batch. It does not poll again or enqueue a second copy. The bridge fills and handlers return 429 until the runtime replies. A prolonged stall makes readiness fail after 60 seconds. ```text Iggy slow -> forwarding loop blocks -> batch stays unacknowledged -> poll() stalls @@ -319,7 +330,9 @@ Size all of them together. ## Observability -`GET /admin/health` returns per-instance JSON: queue depth and capacity, serving endpoint counts by origin plus expired and revoked counts, `state_submitted`, and header loss counters. +`GET /admin/health` returns per-instance JSON: queue depth and capacity, serving endpoint counts by origin plus expired and revoked counts, `state_submitted`, and header loss counters. It answers 200 with `"status":"degraded"` when an instance is not ready. + +`poll_is_live` describes the poll path; `staged_batch_is_stuck` separately flags a staged batch that has kept receiving NACKs for more than 60 seconds. In that case polling continues, but public `/health` answers 503. A batch still awaiting its first runtime result instead makes `poll_is_live` false after 60 seconds without marking `staged_batch_is_stuck`. `GET /admin/metrics` returns Prometheus text format. The runtime's own stage histograms begin at `poll()`, so they see nothing a sender experiences; these cover the gateway's own handling. The clock starts after the request body has been read, so a slow or large upload is not counted in `http_source_request_duration_seconds`. diff --git a/core/connectors/sources/http_source/config.toml b/core/connectors/sources/http_source/config.toml index 6ac501d525..a242a27382 100644 --- a/core/connectors/sources/http_source/config.toml +++ b/core/connectors/sources/http_source/config.toml @@ -53,6 +53,9 @@ management_token = "replace_with_admin_token" max_body_size_bytes = 1048576 buffer_capacity = 10000 max_batch_size = 500 +# Omitted: keep retrying NACKed webhook batches. Set a positive value only +# if accepted webhooks can be replayed after a connector restart. +# max_consecutive_nacks = 5 include_http_metadata = true # The provider's delivery id is the dedup key a consumer needs, because # delivery is best-effort in both directions. See the README. diff --git a/core/connectors/sources/http_source/src/lib.rs b/core/connectors/sources/http_source/src/lib.rs index a0961ef63e..75caed60e7 100644 --- a/core/connectors/sources/http_source/src/lib.rs +++ b/core/connectors/sources/http_source/src/lib.rs @@ -36,6 +36,7 @@ use secrecy::{ExposeSecret, SecretString}; use serde::{Deserialize, Serialize}; use std::collections::BTreeSet; use std::net::SocketAddr; +use std::num::NonZeroU32; use std::str::FromStr; use std::sync::Arc; use std::sync::atomic::{AtomicBool, AtomicU32, AtomicU64, AtomicUsize, Ordering}; @@ -78,10 +79,9 @@ pub const MAX_BODY_SIZE_BYTES_LIMIT: usize = MAX_PAYLOAD_SIZE as usize; pub const DEFAULT_HMAC_HEADER: &str = "X-Hub-Signature-256"; pub const DEFAULT_HMAC_PREFIX: &str = "sha256="; -/// How long after a poll returns the source still counts as live. Generous -/// against the SDK's 30s batch-result timeout, which is the longest a healthy -/// source sits between polls. -const POLL_LIVENESS_SECONDS: u64 = 60; +/// How long after a poll returns the source still counts as live. A stalled +/// runtime can keep one batch awaiting its result beyond this window. +pub(crate) const POLL_LIVENESS_SECONDS: u64 = 60; /// HTTP handler side of an instance's bridge. The pair comes from /// `bounded_async` because the `poll()` side needs `recv().await`; the handler @@ -253,6 +253,8 @@ pub struct SharedState { /// When the last `poll()` returned. Covers the gap between polls, where /// the SDK is awaiting a batch result and nothing is in flight. last_poll_at: AtomicU64, + /// First NACK of the currently staged webhook batch, or zero if none. + staged_nack_since: AtomicU64, /// Serializes registry writers, which are all control-plane. The request /// path only ever loads the `ArcSwap` and never touches this. registry_writer: StdMutex<()>, @@ -381,18 +383,25 @@ impl SharedState { } } - /// Whether the poll task still looks alive. + /// Whether the poll path has advanced within the readiness window. /// /// An in-flight poll counts even when it has been blocked on an empty /// bridge for hours, which is the normal state of a quiet gateway. The /// timestamp covers the other case, where the SDK is between polls waiting - /// for a batch result. Neither advances once the poll task has stopped. + /// for a batch result. A stopped task or a result stalled beyond the + /// readiness window both make this false. pub fn poll_is_live(&self, now_seconds: u64) -> bool { self.poll_active.load(Ordering::Acquire) || now_seconds.saturating_sub(self.last_poll_at.load(Ordering::Acquire)) <= POLL_LIVENESS_SECONDS } + /// Whether a retained webhook batch has been NACKed beyond the readiness window. + pub(crate) fn staged_batch_is_stuck(&self, now_seconds: u64) -> bool { + let since = self.staged_nack_since.load(Ordering::Acquire); + since != 0 && now_seconds.saturating_sub(since) > POLL_LIVENESS_SECONDS + } + /// Marks the instance as having left the route table, before anything /// counts what its bridge still holds. pub(crate) fn mark_departed(&self) { @@ -438,10 +447,9 @@ impl SharedState { } // Re-posting immediately turns a latched state store into a tight // loop: `poll()` produces another state-only batch, the runtime - // refuses it for the same reason, and the SDK stops the source after - // five of those without calling `close()`. Backing off keeps the - // mutation deliverable without spending that budget on a store that - // is not going to answer yet. Dropping the permit instead would leave + // refuses it for the same reason. Backing off keeps the mutation + // deliverable without hammering a store that is not going to answer + // yet. Dropping the permit instead would leave // an idle gateway, one with no traffic and no further mutations, with // no wakeup at all. let state_flush = Arc::clone(&self.state_flush); @@ -539,6 +547,9 @@ pub struct HttpSourceConfig { /// Maximum messages returned by a single `poll()`. #[serde(default = "default_max_batch_size")] pub max_batch_size: usize, + /// Omitted to keep retrying NACKed batches; a nonzero value enables the SDK breaker. + #[serde(default)] + pub max_consecutive_nacks: Option, /// Path segment exposed as `POST /topics/{topic_path}`. Unset disables /// the named path, leaving only secret-path endpoints. #[serde(default)] @@ -608,6 +619,11 @@ impl EndpointAuthType { impl HttpSourceConfig { fn validate(&self) -> Result<(), Error> { + if self.max_consecutive_nacks == Some(0) { + return Err(Error::InvalidConfigValue( + "max_consecutive_nacks must be greater than zero".to_string(), + )); + } self.listen_addr.parse::().map_err(|error| { Error::InvalidConfigValue(format!("listen_addr '{}': {error}", self.listen_addr)) })?; @@ -827,6 +843,7 @@ impl HttpSource { has_polled: AtomicBool::new(false), departed: AtomicBool::new(false), last_poll_at: AtomicU64::new(0), + staged_nack_since: AtomicU64::new(0), registry_writer: StdMutex::new(()), state_flush: Arc::new(Notify::new()), state_flush_nacks: AtomicU32::new(0), @@ -906,10 +923,9 @@ impl HttpSource { /// that bounded, visible backpressure for silent loss growing with the /// length of the outage, and no downstream can detect it. /// - /// A batch that can never be delivered would replay forever, but this - /// connector rejects malformed work at the door rather than mid-stream: - /// an empty body gets 400 and an oversized one 413, headers are clamped on - /// accept, and `Schema::Raw` cannot fail to decode. + /// Malformed work is rejected before enqueueing, but downstream failures + /// can still leave an accepted batch replaying indefinitely. Readiness + /// fails once it has been NACKed longer than the liveness window. fn on_nack(&self) -> Result<(), Error> { let mut staged = self.lock_staged(); let Some(batch) = staged.as_mut() else { @@ -941,6 +957,12 @@ impl HttpSource { let nacks = batch.nacks; let count = batch.messages.len(); drop(staged); + let _ = self.shared.staged_nack_since.compare_exchange( + 0, + unix_now_seconds(), + Ordering::AcqRel, + Ordering::Acquire, + ); warn!( "Runtime NACKed {count} messages ({nacks} so far) for {CONNECTOR_NAME} connector ID: {}", self.shared.id @@ -968,6 +990,15 @@ impl Source for HttpSource { opened } + fn batch_policy(&self) -> source::BatchPolicy { + source::BatchPolicy::default().with_max_consecutive_nacks( + self.shared + .config + .max_consecutive_nacks + .and_then(NonZeroU32::new), + ) + } + async fn poll(&self) -> Result { let _polling = self.shared.enter_poll(); // Published here rather than in `join`, because this is the first @@ -1079,15 +1110,15 @@ impl Source for HttpSource { /// the staged copy is free. /// /// A Nack does not mean neither happened. The runtime NACKs a batch it sent - /// successfully but whose state it could not persist, and the SDK NACKs on - /// its own result timeout while that send may still have landed. Holding - /// the copy and replaying it is therefore at-least-once by construction: + /// successfully but whose state it could not persist. Holding the copy + /// and replaying it is therefore at-least-once by construction: /// the alternative, dropping it, would turn every one of those into /// silent loss. async fn on_batch_result(&self, result: source::SourceBatchResult) -> Result<(), Error> { match result { source::SourceBatchResult::Ack => { self.lock_staged().take(); + self.shared.staged_nack_since.store(0, Ordering::Release); // Only the batch that carried the flush clears its backoff, // the same rule `on_nack` already applies. Clearing on every // Ack let traffic zero a counter that traffic knows nothing @@ -1152,16 +1183,11 @@ impl Source for HttpSource { /// Backoff before re-posting a state flush the runtime refused. /// /// The first retry is immediate, so an isolated refusal costs nothing. Past -/// that it doubles to an eight second ceiling, which is inside the SDK's -/// result timeout, so a store that recovers is picked up promptly. +/// that it doubles to an eight second ceiling, so a store that recovers is +/// picked up promptly. /// -/// What a store that never recovers costs depends on whether the gateway is -/// busy, and only the idle case ends. Idle, every batch is the flush, so its -/// refusals are consecutive and the SDK stops the source after five. Under -/// traffic they are not consecutive: each traffic batch Acks and resets the -/// SDK's counter, so the source runs on indefinitely, retrying at the ceiling. -/// #3941 is the subject there, not something this connector can fix from the -/// inside. +/// A store that never recovers is retried at the ceiling. The HTTP source has +/// no default NACK stop limit because its accepted webhooks exist only here. fn state_flush_retry_delay(attempt: u32) -> Duration { match attempt { 0 => Duration::ZERO, @@ -1227,6 +1253,7 @@ pub(crate) mod test_support { max_body_size_bytes: DEFAULT_MAX_BODY_SIZE_BYTES, buffer_capacity: DEFAULT_BUFFER_CAPACITY, max_batch_size: DEFAULT_MAX_BATCH_SIZE, + max_consecutive_nacks: None, topic_path: topic_path.map(str::to_string), instance_name: None, auth_bearer_token: None, @@ -1311,6 +1338,7 @@ mod tests { assert_eq!(config.max_body_size_bytes, DEFAULT_MAX_BODY_SIZE_BYTES); assert_eq!(config.buffer_capacity, DEFAULT_BUFFER_CAPACITY); assert_eq!(config.max_batch_size, DEFAULT_MAX_BATCH_SIZE); + assert_eq!(config.max_consecutive_nacks, None); assert!(config.include_http_metadata); assert!(config.topic_path.is_none()); assert!(config.auth_bearer_token.is_none()); @@ -1324,6 +1352,76 @@ mod tests { assert!(parse(minimal_config_json()).validate().is_ok()); } + #[test] + fn given_default_config_when_policy_read_should_disable_nack_breaker() { + let source = HttpSource::new(1, parse(minimal_config_json()), None); + assert_eq!(source.batch_policy().max_consecutive_nacks(), None); + } + + #[test] + fn given_custom_nack_limit_when_policy_read_should_apply_limit() { + let config = parse(r#"{"listen_addr": "127.0.0.1:9090", "max_consecutive_nacks": 12}"#); + let source = HttpSource::new(1, config, None); + assert_eq!( + source.batch_policy().max_consecutive_nacks(), + NonZeroU32::new(12) + ); + } + + #[test] + fn given_zero_nack_limit_when_validated_should_name_field() { + let config = parse(r#"{"listen_addr": "127.0.0.1:9090", "max_consecutive_nacks": 0}"#); + assert!(matches!( + config.validate(), + Err(Error::InvalidConfigValue(message)) if message.contains("max_consecutive_nacks") + )); + } + + #[tokio::test] + async fn given_repeated_nacks_should_keep_first_stamp_and_poll_liveness() { + let source = HttpSource::new(1, parse(minimal_config_json()), None); + source.shared.poll_active.store(true, Ordering::Release); + source.stage(vec![queued("event")]); + source + .on_nack() + .expect("NACK should retain the staged batch"); + assert_ne!(source.shared.staged_nack_since.load(Ordering::Acquire), 0); + source + .shared + .staged_nack_since + .store(100, Ordering::Release); + + source + .on_nack() + .expect("a repeated NACK should retain the staged batch"); + assert_eq!(source.shared.staged_nack_since.load(Ordering::Acquire), 100); + + assert!(source.shared.poll_is_live(100 + POLL_LIVENESS_SECONDS)); + assert!(source.shared.poll_is_live(101 + POLL_LIVENESS_SECONDS)); + assert!( + !source + .shared + .staged_batch_is_stuck(100 + POLL_LIVENESS_SECONDS) + ); + assert!( + source + .shared + .staged_batch_is_stuck(101 + POLL_LIVENESS_SECONDS) + ); + + source + .on_batch_result(SourceBatchResult::Ack) + .await + .expect("ACK should clear the staged batch"); + assert_eq!(source.shared.staged_nack_since.load(Ordering::Acquire), 0); + assert!(source.shared.poll_is_live(101 + POLL_LIVENESS_SECONDS)); + assert!( + !source + .shared + .staged_batch_is_stuck(101 + POLL_LIVENESS_SECONDS) + ); + } + #[test] fn given_unparsable_listen_addr_when_validated_should_reject() { let config = parse(r#"{"listen_addr": "not-an-address"}"#); @@ -1638,9 +1736,8 @@ mod tests { #[test] fn given_repeated_refusals_when_retried_should_back_off_to_a_ceiling() { - // Growth is what stops a latched store burning the SDK's five-NACK - // budget in a tight loop; the ceiling is what stops a recovered store - // waiting minutes to be noticed. + // Growth avoids a tight retry loop against a latched store; the ceiling + // lets a recovered store be noticed promptly. let delays: Vec = (1..=10).map(state_flush_retry_delay).collect(); for pair in delays.windows(2) { assert!( @@ -1654,7 +1751,7 @@ mod tests { assert_eq!( *delays.last().expect("range is not empty"), Duration::from_secs(8), - "the delay must stop growing well inside the SDK's result timeout" + "the delay must stop growing at the configured ceiling" ); } @@ -2437,7 +2534,7 @@ mod tests { !source .shared .poll_is_live(returned_at + POLL_LIVENESS_SECONDS + 1), - "but a poll task the SDK stopped after five NACKs eventually goes stale, which is what stops readiness reporting ok" + "a stopped poll task eventually goes stale and fails readiness" ); } diff --git a/core/connectors/sources/http_source/src/management.rs b/core/connectors/sources/http_source/src/management.rs index 3387b21ee8..679848b3a4 100644 --- a/core/connectors/sources/http_source/src/management.rs +++ b/core/connectors/sources/http_source/src/management.rs @@ -216,7 +216,7 @@ async fn register_endpoint( ); } } - warn_if_poll_stopped(&instance, "Registered", endpoint_id.as_str()); + warn_if_poll_cannot_forward(&instance, "Registered", endpoint_id.as_str()); if let Some(failure) = republish(&state).await { // Undo the insert. Left in place it would be persisted on the next // flush and come back live after a restart, despite the caller having @@ -337,7 +337,7 @@ async fn rotate_secret( return error_response(StatusCode::SERVICE_UNAVAILABLE, "instance is closing"); } - warn_if_poll_stopped(&instance, "Rotated", &endpoint_id); + warn_if_poll_cannot_forward(&instance, "Rotated", &endpoint_id); info!( "Rotated the secret for endpoint {} on {CONNECTOR_NAME} connector ID: {}", EndpointId::log_prefix_of(&endpoint_id), @@ -386,7 +386,7 @@ async fn revoke_endpoint( return error_response(StatusCode::SERVICE_UNAVAILABLE, "instance is closing"); } - warn_if_poll_stopped(&instance, "Revoked", &endpoint_id); + warn_if_poll_cannot_forward(&instance, "Revoked", &endpoint_id); info!( "Revoked endpoint {} on {CONNECTOR_NAME} connector ID: {}", EndpointId::log_prefix_of(&endpoint_id), @@ -452,18 +452,27 @@ fn owner_of(state: &ServerState, endpoint_id: &str) -> Option> .find(|instance| instance.registry().endpoint(endpoint_id).is_some()) } -/// Warns when a mutation landed on an instance whose poll task looks stopped. +/// Warns when a mutation cannot reach the runtime through this poll path. /// -/// [`still_joined`] proves only that the instance is registered. The SDK stops -/// the poll task after five consecutive NACKs without calling `close()`, so an -/// instance can stay registered and keep answering 201 and 202 for changes -/// nothing will ever carry to the runtime. `/admin/health` reports the same -/// pair, as `poll_is_live` and `has_polled`. See #3941. -fn warn_if_poll_stopped(instance: &Arc, action: &str, endpoint_id: &str) { - if instance.poll_is_live(unix_now_seconds()) { +/// [`still_joined`] proves only that the instance is registered. A configured +/// NACK limit can stop polling without calling +/// `close()`, so an instance can stay registered while changes cannot reach +/// the runtime. A stuck batch also blocks the next state flush while polling +/// continues. `/admin/health` reports both conditions separately. +fn warn_if_poll_cannot_forward(instance: &Arc, action: &str, endpoint_id: &str) { + let now = unix_now_seconds(); + let poll_is_live = instance.poll_is_live(now); + if poll_is_live && !instance.staged_batch_is_stuck(now) { return; } let endpoint = EndpointId::log_prefix_of(endpoint_id); + if poll_is_live { + warn!( + "{action} endpoint {endpoint} on {CONNECTOR_NAME} connector ID: {} while a staged batch is still being NACKed; polling continues, but the change cannot reach the runtime until that batch succeeds", + instance.id + ); + return; + } // Two ways to be not live, and they want different words rather than one // message or none. Saying "stopped" for an instance that has not started // was a false alarm on ordinary startup timing, but staying silent about @@ -474,7 +483,7 @@ fn warn_if_poll_stopped(instance: &Arc, action: &str, endpoint_id: // Without this line nothing at all reports that. if instance.has_polled() { warn!( - "{action} endpoint {endpoint} on {CONNECTOR_NAME} connector ID: {} while its poll task looks stopped; the change is in memory but nothing is carrying it to the runtime", + "{action} endpoint {endpoint} on {CONNECTOR_NAME} connector ID: {} while its poll path has not advanced; the task may be stopped or awaiting a runtime batch result, and the change cannot reach the runtime yet", instance.id ); } else { diff --git a/core/connectors/sources/http_source/src/server.rs b/core/connectors/sources/http_source/src/server.rs index 4f4d7ed6fc..609b37a70f 100644 --- a/core/connectors/sources/http_source/src/server.rs +++ b/core/connectors/sources/http_source/src/server.rs @@ -337,7 +337,10 @@ impl Published { .instances .iter() .filter(|instance| instance.has_polled()) - .all(|instance| instance.poll_is_live(now_seconds)) + .all(|instance| { + instance.poll_is_live(now_seconds) + && !instance.staged_batch_is_stuck(now_seconds) + }) } } @@ -991,12 +994,10 @@ async fn handle_admin_metrics(State(state): State>) -> Response /// Readiness for a load balancer: unavailable until an instance is serving. /// -/// A joined instance is not the same as a polling one. The SDK stops the poll -/// task after five consecutive NACKs without calling `close()`, so the instance -/// stays joined and this kept answering 200 to the load balancer while handlers -/// accepted into a bridge nobody drained, until it filled and 429d forever. -/// Readiness needs all three: an instance, a route it can reach, and a poll -/// task still running behind it. +/// A joined instance is not the same as a polling one. A configured NACK +/// limit can stop polling without calling `close()`; +/// repeated NACKs can also leave a staged batch stuck while polling continues. +/// Readiness needs an instance, a reachable route, and a healthy poll path. async fn handle_health(State(state): State>) -> Response { let now = unix_now_seconds(); // One guard, so readiness cannot be computed from a route table and an @@ -1024,29 +1025,7 @@ async fn handle_admin_health(State(state): State>) -> Response let instances = published .instances .iter() - .map(|instance| { - let registry = instance.registry(); - InstanceHealth { - instance: instance.instance_name.clone(), - topic_path: instance.config.topic_path.clone(), - buffer_used: instance.sender.len(), - buffer_capacity: instance.config.buffer_capacity, - endpoints_static: registry.serving_count_by_origin(EndpointOrigin::Static, now), - endpoints_dynamic: registry.serving_count_by_origin(EndpointOrigin::Dynamic, now), - endpoints_expired: registry.expired_count(now), - endpoints_revoked: registry.revoked_count(), - named_path: instance.config.topic_path.is_some(), - poll_is_live: instance.poll_is_live(now), - has_polled: instance.has_polled(), - // The registry is handed over whole, so an owed flush is the - // whole answer. Still not `persisted`: the flag clears when the - // state leaves the plugin, and the runtime's write landing is - // something no poll return value reports back. - state_submitted: !instance.has_pending_state(), - headers_dropped: state.metrics.headers_dropped(&instance.instance_name), - headers_clamped: state.metrics.headers_clamped(&instance.instance_name), - } - }) + .map(|instance| InstanceHealth::from_instance(instance, &state.metrics, now)) .collect(); // Derived, not a constant: this answered "ok" while `/health` was @@ -1360,18 +1339,21 @@ struct InstanceHealth { /// Whether a poll has run recently enough to believe one still will. /// /// `state_submitted` says a change was handed over; this says whether - /// anything is still there to hand the next one to. The SDK stops the poll - /// task after five consecutive NACKs without calling `close()`, so the - /// instance stays registered and keeps accepting mutations that will never - /// be persisted. See #3941. + /// the poll path has advanced recently. A configured NACK limit can stop + /// polling without calling `close()`, leaving registered routes that + /// cannot drain or persist mutations. A pending runtime result can also + /// stall the path beyond the readiness window. /// /// Read it with `has_polled`, which separates the two ways this can be /// false. `has_polled` false alongside it is an instance still starting /// up, and the listener is not held to that one. `has_polled` true - /// alongside it is a poll task that stopped, which is what takes the whole + /// alongside it is a stalled or stopped poll path, which takes the whole /// listener out of rotation. This one true with `has_polled` false is /// simply a first poll in flight. poll_is_live: bool, + /// A retained batch is still being retried after the readiness window. + /// Polling may continue even when this is true. + staged_batch_is_stuck: bool, has_polled: bool, state_submitted: bool, /// Named to mirror the metric families exactly, so an operator reading @@ -1381,6 +1363,33 @@ struct InstanceHealth { headers_clamped: u64, } +impl InstanceHealth { + fn from_instance(instance: &SharedState, metrics: &Metrics, now: u64) -> Self { + let registry = instance.registry(); + Self { + instance: instance.instance_name.clone(), + topic_path: instance.config.topic_path.clone(), + buffer_used: instance.sender.len(), + buffer_capacity: instance.config.buffer_capacity, + endpoints_static: registry.serving_count_by_origin(EndpointOrigin::Static, now), + endpoints_dynamic: registry.serving_count_by_origin(EndpointOrigin::Dynamic, now), + endpoints_expired: registry.expired_count(now), + endpoints_revoked: registry.revoked_count(), + named_path: instance.config.topic_path.is_some(), + poll_is_live: instance.poll_is_live(now), + staged_batch_is_stuck: instance.staged_batch_is_stuck(now), + has_polled: instance.has_polled(), + // The registry is handed over whole, so an owed flush is the + // whole answer. Still not `persisted`: the flag clears when the + // state leaves the plugin, and the runtime's write landing is + // something no poll return value reports back. + state_submitted: !instance.has_pending_state(), + headers_dropped: metrics.headers_dropped(&instance.instance_name), + headers_clamped: metrics.headers_clamped(&instance.instance_name), + } + } +} + #[cfg(test)] mod tests { use super::*; @@ -1912,6 +1921,7 @@ mod tests { assert_eq!(instance["endpoints_dynamic"], 0); assert_eq!(instance["endpoints_revoked"], 0); assert_eq!(instance["named_path"], true); + assert_eq!(instance["staged_batch_is_stuck"], false); assert_eq!(instance["state_submitted"], true); close(&mut source).await; } @@ -2746,6 +2756,43 @@ mod tests { close(&mut source).await; } + #[tokio::test] + async fn given_stuck_batch_when_polling_continues_should_fail_readiness() { + let mut source = open(1, config(free_port(), free_port(), &[ENDPOINT_ONE])).await; + let shared = Arc::clone(&source.shared); + let response = post_signed(&base_url(&source), ENDPOINT_ONE, "{}").await; + assert_eq!(response.status(), StatusCode::OK); + let batch = source + .poll() + .await + .expect("the accepted webhook should be polled"); + assert_eq!(batch.messages.len(), 1); + source + .on_batch_result(iggy_connector_sdk::source::SourceBatchResult::Nack) + .await + .expect("the NACK should retain the batch"); + + let polling = shared.enter_poll(); + let later = unix_now_seconds() + crate::POLL_LIVENESS_SECONDS + 1; + let state = { + let servers = SERVERS.lock().await; + Arc::clone( + &servers + .get(&source.shared.config.listen_addr) + .expect("source listener should be registered") + .state, + ) + }; + assert!(shared.poll_is_live(later)); + assert!(shared.staged_batch_is_stuck(later)); + assert!(!state.published().is_ready(later)); + let health = InstanceHealth::from_instance(&shared, &state.metrics, later); + assert!(health.poll_is_live); + assert!(health.staged_batch_is_stuck); + drop(polling); + close(&mut source).await; + } + #[tokio::test] async fn given_one_stopped_instance_when_health_checked_should_report_unavailable() { // One address fronts both instances, so a sibling whose poll task has diff --git a/core/connectors/sources/influxdb_source/dependencies.md b/core/connectors/sources/influxdb_source/dependencies.md index 84d0ed0e94..f6058a3b8c 100644 --- a/core/connectors/sources/influxdb_source/dependencies.md +++ b/core/connectors/sources/influxdb_source/dependencies.md @@ -43,7 +43,7 @@ to `cargo tree -p iggy_connector_influxdb_source` for the full graph. | `csv` | `^1.4.0` | MIT / Apache-2.0 | Parses InfluxDB V2 annotated-CSV query responses in `row.rs::parse_csv_rows`. | | `dashmap` | `^6.1.0` | MIT | Concurrent hash map; injected into this crate's namespace by the `source_connector!` macro expansion in the SDK. Not used directly in source files. | | `iggy_common` | `^0.11.0` | Apache-2.0 | Shared Iggy types: `DateTime`, `Utc`, and `serde_secret` used for safe token serialisation in config structs. | -| `iggy_connector_sdk` | `^0.5.0` | Apache-2.0 | Core connector abstractions: `Source` trait, `ProducedMessage`, `ProducedMessages`, `ConnectorState`, `Schema`, `Error`, retry/circuit-breaker utilities, and the `source_connector!` registration macro. | +| `iggy_connector_sdk` | `^0.5.1-edge.2` | Apache-2.0 | Core connector abstractions: `Source` trait, `ProducedMessage`, `ProducedMessages`, `ConnectorState`, `Schema`, `Error`, retry/circuit-breaker utilities, and the `source_connector!` registration macro. | | `regex` | `^1.12.3` | MIT / Apache-2.0 | Compiles the RFC 3339 cursor-validation regex once via `OnceLock` in `common.rs::cursor_re()`. | | `reqwest` | `^0.13.3` | MIT / Apache-2.0 | Async HTTP client used to issue Flux (V2) and SQL (V3) query requests to InfluxDB. | | `reqwest-middleware` | `^0.5.1` | MIT | Middleware wrapper around `reqwest::Client` that attaches the retry and tracing layers built by `iggy_connector_sdk::retry::build_retry_client`. | diff --git a/core/integration/tests/connectors/runtime/http_state.rs b/core/integration/tests/connectors/runtime/http_state.rs index fe461b96bc..af69ac627e 100644 --- a/core/integration/tests/connectors/runtime/http_state.rs +++ b/core/integration/tests/connectors/runtime/http_state.rs @@ -334,6 +334,37 @@ async fn given_conflict_mid_stream_should_nack_and_latch( fixture.store.conflict_mode.store(true, Ordering::SeqCst); wait_for_status(harness, ConnectorStatus::Error).await; + let version_after_conflict = fixture.store.version.load(Ordering::SeqCst); + let puts_after_conflict = fixture.store.put_count.load(Ordering::SeqCst); + sleep(Duration::from_millis(500)).await; + assert_eq!( + fixture.store.put_count.load(Ordering::SeqCst), + puts_after_conflict, + "a latched provider must not send further PUTs" + ); + assert_eq!( + fixture.store.version.load(Ordering::SeqCst), + version_after_conflict, + "the checkpoint must not advance after a 412" + ); + + let deadline = Instant::now() + WAIT_DEADLINE; + loop { + let source = fetch_source(harness).await; + if source + .last_error + .as_ref() + .is_some_and(|error| error.message.contains("consecutive NACK limit reached")) + { + break; + } + assert!( + Instant::now() < deadline, + "source did not report a restart-required NACK-limit stop" + ); + sleep(POLL_INTERVAL).await; + } + let api_url = harness .connectors_runtime() .expect("connector runtime should be available") @@ -353,20 +384,6 @@ async fn given_conflict_mid_stream_should_nack_and_latch( stats.sources_running, 0, "a latched source has Error status and must not count as running" ); - - let version_after_conflict = fixture.store.version.load(Ordering::SeqCst); - let puts_after_conflict = fixture.store.put_count.load(Ordering::SeqCst); - sleep(Duration::from_millis(500)).await; - assert_eq!( - fixture.store.put_count.load(Ordering::SeqCst), - puts_after_conflict, - "a latched provider must not send further PUTs" - ); - assert_eq!( - fixture.store.version.load(Ordering::SeqCst), - version_after_conflict, - "the checkpoint must not advance after a 412" - ); } #[iggy_harness(