Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 5 additions & 3 deletions .claude/skills/connector-runtime/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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<DashMap<u32, SourceSenderEntry>>` (`pub(crate)`). `SourceSenderEntry` wraps the sender + a pre-extracted owned `Counter` (the `errors` series, `Arc<AtomicU64>` inside). The FFI callback bumps errors on deserialize or channel-closed failure with one relaxed atomic - no `Family` lookup, no `Arc<Metrics>` handle.
4. `source_forwarding_loop` pulls from the channel, deserializes, applies transforms, encodes via `StreamEncoder`, sends to Iggy producer.
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:
Expand Down
23 changes: 13 additions & 10 deletions .claude/skills/connector-source/SKILL.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`
Expand Down Expand Up @@ -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

Expand Down
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -208,7 +208,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"
Expand Down
8 changes: 7 additions & 1 deletion core/connectors/runtime/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -32,10 +32,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,
Expand Down Expand Up @@ -84,6 +85,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<extern "C" fn(id: u32, callback: SourceStoppedCallback) -> 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,
Expand Down Expand Up @@ -189,12 +192,14 @@ async fn main() -> 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,
});
Expand Down Expand Up @@ -506,6 +511,7 @@ struct SourceConnectorProducer {

struct SourceConnectorWrapper {
handle_callback: HandleCallback,
register_stop_callback: Option<RegisterStopCallback>,
batch_result_callback: BatchResultCallback,
plugins: Vec<SourceConnectorPlugin>,
}
Expand Down
128 changes: 116 additions & 12 deletions core/connectors/runtime/src/manager/source.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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<String, Arc<Mutex<SourceDetails>>>,
Expand Down Expand Up @@ -88,18 +95,32 @@ impl SourceManager {

pub async fn set_error(&self, key: &str, error_message: &str, metrics: Option<&Arc<Metrics>>) {
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<Metrics>,
) -> 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
}
}

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -286,6 +308,7 @@ impl SourceManager {
transforms,
state_storage,
handle_callback,
register_stop_callback,
batch_result_callback,
context.clone(),
);
Expand Down Expand Up @@ -362,6 +385,13 @@ pub struct SourceDetails {
}

impl SourceDetails {
fn record_error(&mut self, message: &str, metrics: Option<&Arc<Metrics>>) {
// 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
Expand Down Expand Up @@ -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
Expand Down
Loading
Loading