This is an automated email from the ASF dual-hosted git repository.
spetz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iggy.git
The following commit(s) were added to refs/heads/master by this push:
new e9f3362e3 feat(connectors)!: add shared retry_async helper to the
connector SDK (#4104)
e9f3362e3 is described below
commit e9f3362e3d79f1becc0cc4232de7098c270d167d
Author: Ryan Huang <[email protected]>
AuthorDate: Fri Sep 11 00:36:41 2026 +0800
feat(connectors)!: add shared retry_async helper to the connector SDK
(#4104)
Relates to #3702 #4084
---
.claude/skills/connector-sdk/SKILL.md | 6 +-
.claude/skills/connector-sink/SKILL.md | 3 +-
Cargo.lock | 1 +
core/connectors/sdk/Cargo.toml | 5 +
core/connectors/sdk/README.md | 13 +
core/connectors/sdk/src/retry.rs | 589 ++++++++++++++++++---
core/connectors/sinks/doris_sink/src/lib.rs | 59 +--
core/connectors/sinks/influxdb_sink/README.md | 6 +-
core/connectors/sinks/influxdb_sink/src/lib.rs | 11 +-
core/connectors/sinks/meilisearch_sink/src/lib.rs | 42 +-
core/connectors/sinks/quickwit_sink/src/lib.rs | 10 +-
core/connectors/sinks/rabbitmq_sink/src/lib.rs | 14 +-
core/connectors/sinks/s3_sink/src/sink.rs | 6 +-
core/connectors/sinks/surrealdb_sink/src/lib.rs | 8 +-
core/connectors/sources/influxdb_source/README.md | 4 +-
core/connectors/sources/influxdb_source/src/lib.rs | 11 +-
16 files changed, 604 insertions(+), 184 deletions(-)
diff --git a/.claude/skills/connector-sdk/SKILL.md
b/.claude/skills/connector-sdk/SKILL.md
index d155b4b54..7d00fd484 100644
--- a/.claude/skills/connector-sdk/SKILL.md
+++ b/.claude/skills/connector-sdk/SKILL.md
@@ -47,7 +47,7 @@ sdk/src/
├── api.rs ConnectorStatus, ConnectorStats (feature = "api").
├── convert.rs owned_value_to_serde_json (simd_json ⇄ serde_json
bridge).
├── log.rs CallbackLayer for tracing across FFI.
-├── retry.rs CircuitBreaker, HttpRetryMiddleware,
exponential_backoff, jitter.
+├── retry.rs retry_async + RetryPolicy, CircuitBreaker,
HttpRetryMiddleware.
├── decoders/ One per schema: json, raw, text, proto, flatbuffer,
avro.
├── encoders/ Mirror of decoders.
└── transforms/ add_fields, delete_fields, update_fields,
filter_fields,
@@ -149,8 +149,10 @@ Plugin authors call this on every consumed message. The
implementation in `lib.r
## Retry helpers (`retry.rs`)
- `CircuitBreaker`: threshold + cooldown, `try_lock()` on the success path to
avoid hot-path contention.
+- `retry_async(policy, context, should_retry, op)`: the retry loop for
anything failing as `Err`. Owns attempt counting, backoff and the per-retry
log; returns `RetryFailure { error, attempts, exhausted }` so callers write
their own terminal log.
+- `retry_backoff(base, retry, max)`: backoff only, for a loop that computes
its own delay. `retry` is 1-based. `exponential_backoff` is the 0-based
primitive underneath and applies no jitter, so call it directly only when the
delay must be exact (`source.rs::nack_retry_delay`, whose tests assert exact
values).
- `HttpRetryMiddleware`: integrates with `reqwest-middleware`. Retries 429 +
5xx + network errors. Honors `Retry-After`.
-- `max_retries` = **total attempts** including the first try, not extra
retries. Document if you change this convention.
+- `max_retries` = **total attempts** including the first try, not extra
retries. Document if you change this convention. `meilisearch_sink` is the
standing exception: its `max_retries` / `max_open_retries` count retries
*after* the first, as its README states.
- New helpers must take `Duration` (not `u64 millis`) on the public API.
Internal computation uses `humantime` parsing of `String`.
## `ConnectorState`
diff --git a/.claude/skills/connector-sink/SKILL.md
b/.claude/skills/connector-sink/SKILL.md
index 80533036b..207e55354 100644
--- a/.claude/skills/connector-sink/SKILL.md
+++ b/.claude/skills/connector-sink/SKILL.md
@@ -82,7 +82,8 @@ for getting them to the external system reliably and
efficiently.
- SDK helpers cover the simple case:
`iggy_connector_sdk::retry::check_connectivity_with_retry(...)` for `open()`,
`HttpRetryMiddleware` for default 429/5xx/network policy on reqwest clients.
- Custom strategies: `reqwest-middleware` +
`reqwest_retry::RetryTransientMiddleware::new_with_policy_and_strategy`.
`http_sink` defines its own `HttpSinkRetryStrategy` (honors
`success_status_codes`, per-status decisions).
- Non-HTTP clients: write `is_transient_error(&e)` mapping driver-specific
codes. `postgres_sink::is_transient_error` maps SQLSTATEs `40001`, `40P01`,
`57P01-03`, `08000/03/06`.
-- Backoff: `iggy_connector_sdk::retry::exponential_backoff(base, attempt,
max)` + `jitter()`.
+- Retry loop: new connectors use
`iggy_connector_sdk::retry::retry_async(policy, context, should_retry, op)` for
anything that fails as `Err`. It owns attempt counting, backoff and the
per-retry log, and returns `RetryFailure { error, attempts, exhausted }`; the
caller logs the terminal failure.
+- Backoff only, for a loop that computes its own delay (it retries on an `Ok`
response, carries a deadline, or reconnects between attempts):
`retry_backoff(base, retry, max)`, where `retry` is 1-based.
`exponential_backoff` is the 0-based primitive underneath, and it applies no
jitter; call it directly only when the delay must be exact, as
`sdk/src/source.rs::nack_retry_delay` needs for its tests.
- Cap retries at 3 total attempts.
### Idempotency
diff --git a/Cargo.lock b/Cargo.lock
index 3afc2f782..4f18e5740 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -7472,6 +7472,7 @@ dependencies = [
"tracing",
"tracing-subscriber",
"uuid",
+ "wiremock",
]
[[package]]
diff --git a/core/connectors/sdk/Cargo.toml b/core/connectors/sdk/Cargo.toml
index 5d6b98daa..3d511d3f5 100644
--- a/core/connectors/sdk/Cargo.toml
+++ b/core/connectors/sdk/Cargo.toml
@@ -68,5 +68,10 @@ tracing = { workspace = true }
tracing-subscriber = { workspace = true }
uuid = { workspace = true }
+[dev-dependencies]
+# `test-util` is not part of `full`; the retry tests need its paused time.
+tokio = { workspace = true, features = ["full", "test-util"] }
+wiremock = { workspace = true }
+
[lints]
workspace = true
diff --git a/core/connectors/sdk/README.md b/core/connectors/sdk/README.md
index 1c72e143d..96f8037ad 100644
--- a/core/connectors/sdk/README.md
+++ b/core/connectors/sdk/README.md
@@ -48,6 +48,19 @@ key = "message"
value.static = "hello"
```
+## Retry helpers
+
+`retry_async` runs an operation that fails with `Err` and retries it while
`should_retry` accepts the error. It owns attempt counting, backoff and the
per-retry log, and returns `RetryFailure { error, attempts, exhausted }` so the
caller logs the terminal failure. `retry_backoff` computes a single delay for a
loop that cannot use `retry_async`, such as `HttpRetryMiddleware`, which
retries on an `Ok` response rather than an `Err`. Its `retry` argument is
1-based.
+
+Two items changed in a way that breaks out-of-tree plugins, so those plugins
must be rebuilt against the current source:
+
+| Removed | Replacement |
+| --- | --- |
+| `ConnectivityConfig` | `RetryPolicy`. `max_open_retries` becomes
`max_attempts`, `retry_delay` becomes `base_delay`, and `open_retry_max_delay`
becomes `max_delay`. |
+| `jitter` (was public) | `retry_backoff`, which applies the jitter itself. |
+
+Both types carry `(u32, Duration, Duration)` and the two delay roles cross
over, so a field-by-field rename compiles and swaps the base delay for the cap.
Map the fields by name.
+
## Protocol Buffers Support
The SDK includes support for Protocol Buffers (protobuf) format with both
encoding and decoding capabilities. Protocol Buffers provide efficient
serialization and are particularly useful for high-performance data streaming
scenarios.
diff --git a/core/connectors/sdk/src/retry.rs b/core/connectors/sdk/src/retry.rs
index 06866e8da..b1d63294d 100644
--- a/core/connectors/sdk/src/retry.rs
+++ b/core/connectors/sdk/src/retry.rs
@@ -24,11 +24,13 @@
//! - [`build_retry_client`] — wraps a `reqwest::Client` with the middleware
//! - [`check_connectivity`] — single health-check probe (GET /health)
//! - [`check_connectivity_with_retry`] — startup probe with exponential
backoff
-//! - [`ConnectivityConfig`] — parameters for the startup retry loop
+//! - [`retry_async`] — generic retry loop for fallible async operations
+//! - [`RetryFailure`] — terminal outcome of [`retry_async`], with attempt
count
+//! - [`RetryPolicy`] — attempt budget and backoff bounds for [`retry_async`]
//! - [`is_transient_status`] — transient HTTP status predicate
//! - [`parse_duration`] — humantime duration parsing with fallback
-//! - [`jitter`] — ±20 % random jitter for retry delays
//! - [`exponential_backoff`] — capped exponential backoff
+//! - [`retry_backoff`] — jittered backoff for a 1-based retry number
//! - [`parse_retry_after`] — HTTP `Retry-After` header parsing
use anyhow::anyhow;
@@ -36,10 +38,12 @@ use http::Extensions;
use humantime::Duration as HumanDuration;
use rand::RngExt as _;
use reqwest_middleware::{ClientBuilder, ClientWithMiddleware, Middleware,
Next};
+use std::fmt;
+use std::future::Future;
use std::str::FromStr;
use std::time::Duration;
use tokio::sync::Mutex;
-use tracing::{info, warn};
+use tracing::{error, info, warn};
// ---------------------------------------------------------------------------
// Circuit breaker
@@ -150,7 +154,7 @@ pub fn parse_duration(value: Option<&str>, default_value:
&str) -> Duration {
}
/// Apply ±20 % random jitter to `base` to spread retry storms.
-pub fn jitter(base: Duration) -> Duration {
+pub(crate) fn jitter(base: Duration) -> Duration {
let millis = base.as_millis() as u64;
let jitter_range = millis / 5; // 20% of base
if jitter_range == 0 {
@@ -161,14 +165,17 @@ pub fn jitter(base: Duration) -> Duration {
}
/// True exponential backoff: `base × 2^attempt`, capped at `max_delay`.
+///
+/// `attempt` is 0-based. Retry loops count from 1, so passing their counter
+/// here makes the first retry wait twice `base`; use [`retry_backoff`], which
+/// takes a 1-based retry number and applies jitter and the cap.
pub fn exponential_backoff(base: Duration, attempt: u32, max_delay: Duration)
-> Duration {
let factor = 2u64.saturating_pow(attempt);
let millis = base
.as_millis()
.saturating_mul(factor as u128)
.min(max_delay.as_millis());
- let millis_u64 = u64::try_from(millis).unwrap_or(u64::MAX);
- Duration::from_millis(millis_u64)
+ Duration::from_millis(u64::try_from(millis).unwrap_or(u64::MAX))
}
/// Parse a `Retry-After` header value (integer seconds).
@@ -181,6 +188,160 @@ pub fn parse_retry_after(value: &str) -> Option<Duration>
{
None
}
+// ---------------------------------------------------------------------------
+// Generic retry loop
+// ---------------------------------------------------------------------------
+
+/// Parameters for [`retry_async`].
+///
+/// `max_attempts` is a *total attempt count*, not a count of extra retries:
+/// `1` runs the operation once and never retries, `3` allows two retries. `0`
+/// behaves as `1`, so a misconfigured value degrades to a single attempt
+/// rather than skipping the operation entirely.
+#[derive(Debug, Clone, Copy)]
+pub struct RetryPolicy {
+ pub max_attempts: u32,
+ pub base_delay: Duration,
+ pub max_delay: Duration,
+}
+
+impl RetryPolicy {
+ /// Jittered backoff before retry number `retry` (1-based). See
+ /// [`retry_backoff`].
+ pub fn backoff(&self, retry: u32) -> Duration {
+ retry_backoff(self.base_delay, retry, self.max_delay)
+ }
+}
+
+/// Jittered, capped backoff before retry number `retry` (1-based): the first
+/// retry waits `base_delay`, the second `2 × base_delay`, and so on, which is
+/// the convention the `retry_delay` config fields document.
+///
+/// The cap is re-applied after jittering, because ±20 % jitter on an
+/// already-capped delay can otherwise land above `max_delay`, which the config
+/// fields document as a strict upper bound.
+///
+/// Prefer [`retry_async`], which calls this for you. Reach for it directly
+/// only in a loop that cannot be expressed as a retried `Result`.
+pub fn retry_backoff(base_delay: Duration, retry: u32, max_delay: Duration) ->
Duration {
+ jitter(exponential_backoff(
+ base_delay,
+ retry.saturating_sub(1),
+ max_delay,
+ ))
+ .min(max_delay)
+}
+
+/// Why [`retry_async`] stopped.
+///
+/// Carries the attempt count so a caller can write its own terminal log: the
+/// helper owns the per-retry line, giving up is the caller's to report.
+#[derive(Debug)]
+pub struct RetryFailure<E> {
+ pub error: E,
+ /// Attempts actually made, including the first.
+ pub attempts: u32,
+ /// `true` when the attempt budget ran out, `false` when `should_retry`
+ /// rejected the error and no further attempt was made.
+ pub exhausted: bool,
+}
+
+impl<E> RetryFailure<E> {
+ /// Discard the attempt bookkeeping and keep the underlying error.
+ pub fn into_error(self) -> E {
+ self.error
+ }
+}
+
+// No `source()`: `Display` already prints the inner error, so returning it
+// here repeats the same text in an error chain.
+impl<E> std::error::Error for RetryFailure<E> where E: std::error::Error +
'static {}
+
+impl<E: fmt::Display> fmt::Display for RetryFailure<E> {
+ fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
+ let reason = if self.exhausted {
+ "ran out of attempts"
+ } else {
+ "hit a non-retryable error"
+ };
+ let plural = if self.attempts == 1 {
+ "attempt"
+ } else {
+ "attempts"
+ };
+ write!(
+ f,
+ "{reason} after {} {plural}: {}",
+ self.attempts, self.error
+ )
+ }
+}
+
+/// Run `operation`, retrying while it fails with an error `should_retry`
+/// accepts and the attempt budget in `policy` is not exhausted.
+///
+/// This is the retry skeleton for connectors whose failures surface as `Err`,
+/// including backends that report failure in-band (a 200 response carrying a
+/// failure status in its body) and so cannot use [`HttpRetryMiddleware`],
+/// which classifies on the status code alone.
+///
+/// `context` identifies the connector and operation in every log line this
+/// emits, e.g. `"Doris sink ID 3 Stream Load (label=abc)"`. Callers build it
+/// before the first attempt, so keep it off per-message paths.
+///
+/// The per-retry line is logged here. Giving up is not: the returned
+/// [`RetryFailure`] carries the attempt count and which condition ended the
+/// loop, so a caller logs the terminal failure at the level and wording that
+/// suit it. The error itself is returned unchanged.
+pub async fn retry_async<T, E, S, Op, Fut>(
+ policy: RetryPolicy,
+ context: &str,
+ should_retry: S,
+ mut operation: Op,
+) -> Result<T, RetryFailure<E>>
+where
+ S: Fn(&E) -> bool,
+ Op: FnMut() -> Fut,
+ Fut: Future<Output = Result<T, E>>,
+ E: fmt::Display,
+{
+ let max_attempts = policy.max_attempts.max(1);
+ let mut attempt = 0u32;
+
+ loop {
+ let error = match operation().await {
+ Ok(value) => {
+ if attempt > 0 {
+ let plural = if attempt == 1 { "retry" } else { "retries"
};
+ info!("{context} succeeded after {attempt} {plural}.");
+ }
+ return Ok(value);
+ }
+ Err(error) => error,
+ };
+
+ attempt += 1;
+ let retryable = should_retry(&error);
+ let budget_spent = attempt >= max_attempts;
+ if !retryable || budget_spent {
+ return Err(RetryFailure {
+ error,
+ attempts: attempt,
+ // A non-retryable error is the reason we stopped even when the
+ // budget happened to run out on the same attempt.
+ exhausted: retryable && budget_spent,
+ });
+ }
+
+ let delay = policy.backoff(attempt);
+ warn!(
+ "{context} failed on attempt {attempt}/{max_attempts}: {error}. \
+ Retrying in {delay:?}..."
+ );
+ tokio::time::sleep(delay).await;
+ }
+}
+
// ---------------------------------------------------------------------------
// reqwest-middleware retry implementation
// ---------------------------------------------------------------------------
@@ -201,18 +362,19 @@ pub fn is_transient_status(status: reqwest::StatusCode)
-> bool {
/// (e.g. `"InfluxDB"`, `"Elasticsearch"`), allowing this middleware to be
/// reused across connectors without misleading log output.
///
-/// The `max_retries` parameter is the *total attempt count* (not the number
-/// of extra attempts), consistent with the rest of the connector retry config:
-/// - `max_retries = 1` → one attempt, no retries on failure
-/// - `max_retries = 3` → up to three attempts (two retries after a failure)
+/// `max_retries` is a total attempt count, as described on [`RetryPolicy`].
///
/// Non-transient error responses (4xx except 429) are returned as-is so
/// callers can inspect the status and body to build a meaningful error.
+///
+/// The loop is hand-written rather than delegating to [`retry_async`] because
+/// a retry here is driven by an `Ok(Response)` carrying a transient status,
+/// and the final response is handed back to the caller instead of being turned
+/// into an `Err`. Backoff still comes from [`RetryPolicy`], so the timing
+/// matches every other connector retry: the first retry waits `retry_delay`.
#[derive(Debug, Clone)]
pub struct HttpRetryMiddleware {
- max_retries: u32,
- retry_delay: Duration,
- max_delay: Duration,
+ policy: RetryPolicy,
log_prefix: &'static str,
}
@@ -224,9 +386,11 @@ impl HttpRetryMiddleware {
log_prefix: &'static str,
) -> Self {
Self {
- max_retries,
- retry_delay,
- max_delay,
+ policy: RetryPolicy {
+ max_attempts: max_retries,
+ base_delay: retry_delay,
+ max_delay,
+ },
log_prefix,
}
}
@@ -256,34 +420,28 @@ impl Middleware for HttpRetryMiddleware {
return Ok(response);
}
- // Parse Retry-After header on 429 before falling back to
- // our own calculated backoff.
+ // Parse Retry-After on 429 before falling back to our own
+ // calculated backoff.
let retry_after = if status ==
reqwest::StatusCode::TOO_MANY_REQUESTS {
response
.headers()
.get("Retry-After")
- .and_then(|v| v.to_str().ok())
+ .and_then(|value| value.to_str().ok())
.and_then(parse_retry_after)
} else {
None
};
attempts += 1;
- if is_transient_status(status) && attempts <
self.max_retries {
+ if is_transient_status(status) && attempts <
self.policy.max_attempts {
// Consume the error body for logging, then retry.
let body_text =
response.text().await.unwrap_or_default();
- let delay = retry_after.unwrap_or_else(|| {
- jitter(exponential_backoff(
- self.retry_delay,
- attempts,
- self.max_delay,
- ))
- });
+ let delay = retry_after.unwrap_or_else(||
self.policy.backoff(attempts));
warn!(
"{} transient error {status} \
(attempt {attempts}/{}): {body_text}. \
Retrying in {delay:?}...",
- self.log_prefix, self.max_retries
+ self.log_prefix, self.policy.max_attempts
);
tokio::time::sleep(delay).await;
current_req = match next_req {
@@ -304,16 +462,12 @@ impl Middleware for HttpRetryMiddleware {
}
Err(e) => {
attempts += 1;
- if attempts < self.max_retries {
- let delay = jitter(exponential_backoff(
- self.retry_delay,
- attempts,
- self.max_delay,
- ));
+ if attempts < self.policy.max_attempts {
+ let delay = self.policy.backoff(attempts);
warn!(
"{} network error (attempt {attempts}/{}): {e}. \
Retrying in {delay:?}...",
- self.log_prefix, self.max_retries
+ self.log_prefix, self.policy.max_attempts
);
tokio::time::sleep(delay).await;
current_req = match next_req {
@@ -358,18 +512,6 @@ pub fn build_retry_client(
// Shared connectivity helper
// ---------------------------------------------------------------------------
-/// Configuration for the startup connectivity retry loop.
-///
-/// This is intentionally separate from per-request retry config so that
-/// startup can wait patiently for a service (e.g. 10 retries over 60 s)
-/// without affecting the shorter per-request retry window used during
-/// normal operation.
-pub struct ConnectivityConfig {
- pub max_open_retries: u32,
- pub open_retry_max_delay: Duration,
- pub retry_delay: Duration,
-}
-
/// Probe `url` with a plain GET and return `Ok(())` if the response is 2xx.
///
/// This is a single, non-retried attempt. The caller is responsible for the
@@ -398,49 +540,334 @@ pub async fn check_connectivity(
/// Retry [`check_connectivity`] with exponential backoff + jitter.
///
-/// `connector_label` is used in log messages (e.g. `"InfluxDB sink connector
ID: 1"`).
+/// `connector_label` names the connector in log messages (e.g. `"InfluxDB
sink"`).
/// `connector_id` is included in log messages for multi-instance deployments.
+///
+/// Startup usually wants a more patient policy than the per-request one (ten
+/// attempts over a minute, say), so callers pass their own [`RetryPolicy`]
+/// rather than reusing the one that governs live traffic.
pub async fn check_connectivity_with_retry(
client: &reqwest::Client,
url: reqwest::Url,
connector_label: &str,
connector_id: u32,
- cfg: &ConnectivityConfig,
+ policy: RetryPolicy,
) -> Result<(), crate::Error> {
- let max_open_retries = cfg.max_open_retries.max(1);
- let mut attempt = 0u32;
+ let context =
+ format!("{connector_label} startup connectivity for connector ID:
{connector_id}");
+
+ retry_async(
+ policy,
+ &context,
+ |_| true,
+ || check_connectivity(client, url.clone(), connector_label),
+ )
+ .await
+ .map_err(|failure| {
+ // `open()`'s Err reaches the FFI boundary and is dropped there, so the
+ // runtime logs only "Plugin initialization failed". Some callers log
+ // the error again themselves.
+ error!("{context} {failure}");
+ failure.into_error()
+ })
+}
- loop {
- match check_connectivity(client, url.clone(), connector_label).await {
- Ok(()) => {
- if attempt > 0 {
- tracing::info!(
- "{connector_label} connectivity established after
{attempt} retries \
- for connector ID: {connector_id}"
- );
- }
- return Ok(());
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use std::cell::Cell;
+ use std::time::Instant;
+ use wiremock::matchers::method;
+ use wiremock::{Mock, MockServer, ResponseTemplate};
+
+ const BASE: Duration = Duration::from_millis(100);
+ const MAX: Duration = Duration::from_secs(10);
+
+ // `jitter` is ±20 %, so every timing assertion below is a band around the
+ // nominal delay rather than an equality.
+ const JITTER_LOW: f64 = 0.8;
+ const JITTER_HIGH: f64 = 1.2;
+
+ #[derive(Debug, Clone, Copy, PartialEq, Eq)]
+ enum TestError {
+ Transient,
+ Permanent,
+ }
+
+ impl fmt::Display for TestError {
+ fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
+ match self {
+ Self::Transient => write!(f, "transient"),
+ Self::Permanent => write!(f, "permanent"),
}
- Err(e) => {
- attempt += 1;
- if attempt >= max_open_retries {
- tracing::error!(
- "{connector_label} connectivity check failed after
{attempt} attempts \
- for connector ID: {connector_id}. Giving up: {e}"
- );
- return Err(e);
- }
- let backoff = jitter(exponential_backoff(
- cfg.retry_delay,
- attempt,
- cfg.open_retry_max_delay,
- ));
- tracing::warn!(
- "{connector_label} health check failed \
- (attempt {attempt}/{max_open_retries}) \
- for connector ID: {connector_id}. Retrying in
{backoff:?}: {e}"
- );
- tokio::time::sleep(backoff).await;
+ }
+ }
+
+ fn should_retry(error: &TestError) -> bool {
+ matches!(error, TestError::Transient)
+ }
+
+ fn policy(max_attempts: u32) -> RetryPolicy {
+ RetryPolicy {
+ max_attempts,
+ base_delay: BASE,
+ max_delay: MAX,
+ }
+ }
+
+ /// Operation that fails with `error` for the first `failures` calls and
then
+ /// succeeds, returning the 1-based number of the call that succeeded.
+ fn failing_times(
+ calls: &Cell<u32>,
+ failures: u32,
+ error: TestError,
+ ) -> impl FnMut() -> std::future::Ready<Result<u32, TestError>> + '_ {
+ move || {
+ let call = calls.get() + 1;
+ calls.set(call);
+ std::future::ready(if call <= failures {
+ Err(error)
+ } else {
+ Ok(call)
+ })
+ }
+ }
+
+ #[tokio::test(start_paused = true)]
+ async fn given_a_transient_failure_should_retry_until_it_succeeds() {
+ let calls = Cell::new(0);
+ let result = retry_async(
+ policy(4),
+ "test",
+ should_retry,
+ failing_times(&calls, 2, TestError::Transient),
+ )
+ .await;
+
+ assert_eq!(result.map_err(RetryFailure::into_error), Ok(3));
+ assert_eq!(calls.get(), 3);
+ }
+
+ #[tokio::test(start_paused = true)]
+ async fn given_a_permanent_failure_should_not_retry() {
+ let calls = Cell::new(0);
+ let result = retry_async(
+ policy(4),
+ "test",
+ should_retry,
+ failing_times(&calls, 2, TestError::Permanent),
+ )
+ .await;
+
+ let failure = result.expect_err("permanent error should not retry");
+ assert_eq!(failure.error, TestError::Permanent);
+ assert_eq!(failure.attempts, 1);
+ assert!(
+ !failure.exhausted,
+ "should_retry rejected it, budget untouched"
+ );
+ assert_eq!(calls.get(), 1);
+ }
+
+ #[tokio::test(start_paused = true)]
+ async fn
given_a_permanent_failure_on_the_last_attempt_should_not_report_exhaustion() {
+ // Budget of 1 makes the last attempt also the first, so both stop
+ // conditions fire at once; the non-retryable one is the real reason.
+ let calls = Cell::new(0);
+ let result = retry_async(
+ policy(1),
+ "test",
+ should_retry,
+ failing_times(&calls, u32::MAX, TestError::Permanent),
+ )
+ .await;
+
+ let failure = result.expect_err("permanent error should fail");
+ assert!(
+ !failure.exhausted,
+ "reported exhaustion for an error that was never retryable"
+ );
+ }
+
+ #[tokio::test(start_paused = true)]
+ async fn given_an_exhausted_budget_should_return_the_last_error() {
+ let calls = Cell::new(0);
+ // Distinct error per attempt, so "last" is actually discriminated.
+ let result: Result<u32, RetryFailure<u32>> = retry_async(
+ policy(3),
+ "test",
+ |_| true,
+ || {
+ let call = calls.get() + 1;
+ calls.set(call);
+ std::future::ready(Err(call))
+ },
+ )
+ .await;
+
+ let failure = result.expect_err("budget should be exhausted");
+ assert_eq!(failure.error, 3, "should surface the final attempt's
error");
+ assert_eq!(failure.attempts, 3);
+ assert!(failure.exhausted);
+ assert_eq!(calls.get(), 3, "max_attempts is a total attempt count");
+ }
+
+ #[tokio::test(start_paused = true)]
+ async fn given_a_single_attempt_budget_should_run_the_operation_once() {
+ for max_attempts in [0, 1] {
+ let calls = Cell::new(0);
+ let result = retry_async(
+ policy(max_attempts),
+ "test",
+ should_retry,
+ failing_times(&calls, u32::MAX, TestError::Transient),
+ )
+ .await;
+
+ let failure = result.expect_err("single attempt should fail");
+ assert_eq!(failure.error, TestError::Transient);
+ assert!(failure.exhausted, "max_attempts = {max_attempts}");
+ assert_eq!(calls.get(), 1, "max_attempts = {max_attempts}");
+ }
+ }
+
+ #[tokio::test(start_paused = true)]
+ async fn given_successive_retries_should_double_the_delay() {
+ let calls = Cell::new(0);
+ let started = tokio::time::Instant::now();
+ let result = retry_async(
+ policy(3),
+ "test",
+ should_retry,
+ failing_times(&calls, 2, TestError::Transient),
+ )
+ .await;
+ let elapsed = started.elapsed();
+
+ assert_eq!(result.map_err(RetryFailure::into_error), Ok(3));
+ // base + 2 × base, each independently jittered. Passing the retry
+ // number straight to `exponential_backoff` would give 2 + 4 × base.
+ let nominal = BASE * 3;
+ assert!(
+ elapsed >= nominal.mul_f64(JITTER_LOW) && elapsed <=
nominal.mul_f64(JITTER_HIGH),
+ "two retries waited {elapsed:?}, expected roughly {nominal:?}"
+ );
+ }
+
+ // `HttpRetryMiddleware` is the path every HTTP connector rides, and
nothing
+ // covered it before: these cases pin the first delay and that a server's
+ // `Retry-After` outranks the computed backoff.
+ fn retry_client(max_attempts: u32, base: Duration, max: Duration) ->
ClientWithMiddleware {
+ build_retry_client(reqwest::Client::new(), max_attempts, base, max,
"test")
+ }
+
+ async fn mock_then_ok(status: u16, headers: &[(&str, &str)]) -> MockServer
{
+ let server = MockServer::start().await;
+ let mut first = ResponseTemplate::new(status);
+ for (name, value) in headers {
+ first = first.insert_header(*name, *value);
+ }
+ Mock::given(method("GET"))
+ .respond_with(first)
+ .up_to_n_times(1)
+ .mount(&server)
+ .await;
+ Mock::given(method("GET"))
+ .respond_with(ResponseTemplate::new(200))
+ .mount(&server)
+ .await;
+ server
+ }
+
+ #[tokio::test]
+ async fn
given_no_retry_after_should_wait_the_base_delay_on_the_first_retry() {
+ let server = mock_then_ok(503, &[]).await;
+ // 1s rather than 200ms so the pass band and the bug's band do not
+ // overlap: the bug produces jitter(2s) in [1.6s, 2.4s], a correct run
+ // produces jitter(1s) in [0.8s, 1.2s], and 1.4s separates them with
+ // room for two loopback round trips.
+ let base = Duration::from_secs(1);
+ let client = retry_client(3, base, Duration::from_secs(30));
+
+ let started = Instant::now();
+ let response = client.get(server.uri()).send().await.unwrap();
+ let elapsed = started.elapsed();
+
+ assert_eq!(response.status(), 200);
+ assert_eq!(
+ server.received_requests().await.unwrap().len(),
+ 2,
+ "expected exactly one retry"
+ );
+ assert!(
+ elapsed >= base.mul_f64(JITTER_LOW),
+ "first retry waited {elapsed:?}, expected roughly {base:?}"
+ );
+ // Guards the middleware's own `policy.backoff(attempts)` wiring, which
+ // the pure `retry_backoff` tests do not reach: feeding a 1-based
+ // counter to the 0-based `exponential_backoff` doubles this delay.
+ // A correct run tops out at 1.2 x base and the bug starts at 1.6 x, so
+ // 1.5 x leaves the widest margin for loopback round trips on a loaded
+ // runner while still failing on the bug.
+ assert!(
+ elapsed < base.mul_f64(1.5),
+ "first retry waited {elapsed:?}, past the {base:?} the config asks
for"
+ );
+ }
+
+ #[tokio::test]
+ async fn
given_a_retry_after_should_take_precedence_over_the_computed_backoff() {
+ let server = mock_then_ok(429, &[("Retry-After", "1")]).await;
+ // Backoff bounded far below the header, so honoring it is visible.
+ let client = retry_client(3, Duration::from_millis(10),
Duration::from_millis(50));
+
+ let started = Instant::now();
+ let response = client.get(server.uri()).send().await.unwrap();
+ let elapsed = started.elapsed();
+
+ assert_eq!(response.status(), 200);
+ assert_eq!(
+ server.received_requests().await.unwrap().len(),
+ 2,
+ "expected exactly one retry"
+ );
+ assert!(
+ elapsed >= Duration::from_millis(900),
+ "used the computed backoff instead of Retry-After: waited
{elapsed:?}"
+ );
+ // The header asks for 1s. Anything far past it means the middleware
+ // added its own backoff on top instead of honoring the header.
+ assert!(
+ elapsed < Duration::from_millis(1500),
+ "waited {elapsed:?}, past the 1s the header asked for"
+ );
+ }
+
+ #[test]
+ fn given_a_retry_number_should_back_off_from_the_base_delay() {
+ let policy = policy(8);
+ for (retry, factor) in [(1u32, 1.0), (2, 2.0), (3, 4.0), (4, 8.0)] {
+ let delay = policy.backoff(retry);
+ let nominal = BASE.mul_f64(factor);
+ assert!(
+ delay >= nominal.mul_f64(JITTER_LOW) && delay <=
nominal.mul_f64(JITTER_HIGH),
+ "retry {retry} backed off {delay:?}, expected roughly
{nominal:?}"
+ );
+ }
+ }
+
+ #[test]
+ fn given_jitter_on_a_capped_delay_should_never_exceed_max_delay() {
+ // Base far above the cap, so every draw starts clamped and only jitter
+ // could push it back over.
+ let policy = RetryPolicy {
+ max_attempts: 8,
+ base_delay: Duration::from_secs(30),
+ max_delay: Duration::from_secs(1),
+ };
+ for retry in 1..=8 {
+ for _ in 0..64 {
+ assert!(policy.backoff(retry) <= policy.max_delay);
}
}
}
diff --git a/core/connectors/sinks/doris_sink/src/lib.rs
b/core/connectors/sinks/doris_sink/src/lib.rs
index b58b112a8..78e2d70d3 100644
--- a/core/connectors/sinks/doris_sink/src/lib.rs
+++ b/core/connectors/sinks/doris_sink/src/lib.rs
@@ -19,7 +19,7 @@ use async_trait::async_trait;
use base64::{Engine as _, engine::general_purpose};
use bytes::Bytes;
use humantime::Duration as HumanDuration;
-use iggy_connector_sdk::retry::{exponential_backoff, jitter};
+use iggy_connector_sdk::retry::{RetryPolicy, retry_async};
use iggy_connector_sdk::{
ConsumedMessage, Error, MessagesMetadata, Payload, Sink, TopicMetadata,
sink_connector,
};
@@ -300,7 +300,7 @@ impl DorisSink {
}
let response = request.send().await.map_err(|e| {
- error!("Doris sink ID {} HTTP request failed: {e}", self.id);
+ warn!("Doris sink ID {} HTTP request failed: {e}", self.id);
Error::HttpRequestFailed(e.to_string())
})?;
@@ -385,7 +385,9 @@ impl DorisSink {
"Doris sink ID {} stream load returned HTTP {status}:
{response_for_log}",
self.id
);
- error!("{msg}");
+ // Per-attempt detail only: `retry_async` logs every retried
+ // attempt, and `consume` logs the terminal error carrying
`msg`.
+ warn!("{msg}");
// 408/429 are 4xx but transient, so include them in the
bounded
// in-request retry path.
return Err(match status {
@@ -417,36 +419,31 @@ impl DorisSink {
))
})?;
- let mut attempt = 0u32;
- loop {
- let error = match self
- .send_stream_load(connected, label, body.clone())
- .await
- .and_then(|response| classify_status(self.id,
&response).map(|()| response))
- {
- Ok(response) => return Ok(response),
- Err(error) => error,
- };
-
- attempt += 1;
- if attempt >= connected.max_retries || !is_transient_error(&error)
{
- return Err(error);
+ let policy = RetryPolicy {
+ max_attempts: connected.max_retries,
+ base_delay: connected.retry_delay,
+ max_delay: connected.max_retry_delay,
+ };
+ let context = format!("Doris sink ID {} Stream Load (label={label})",
self.id);
+
+ retry_async(policy, &context, is_transient_error, || {
+ let body = body.clone();
+ async move {
+ let response = self.send_stream_load(connected, label,
body).await?;
+ classify_status(self.id, &response)?;
+ Ok(response)
}
-
- // `attempt` counts completed attempts. Subtract one so the first
- // retry waits exactly the configured base delay (base * 2^0).
- let delay = jitter(exponential_backoff(
- connected.retry_delay,
- attempt - 1,
- connected.max_retry_delay,
- ))
- .min(connected.max_retry_delay);
+ })
+ .await
+ .map_err(|failure| {
+ // The only place the attempt count and the reason survive:
+ // `consume` logs the error itself, which carries neither.
warn!(
- "Doris sink ID {} transient Stream Load failure on attempt
{attempt}/{} (label={label}): {error}; retrying in {delay:?}",
- self.id, connected.max_retries
+ "Doris sink ID {} Stream Load (label={label}) {failure}",
+ self.id
);
- tokio::time::sleep(delay).await;
- }
+ failure.into_error()
+ })
}
}
@@ -1011,7 +1008,7 @@ impl Sink for DorisSink {
self.id
);
}
- // `exponential_backoff` already caps at the max, but a base above the
cap
+ // `retry_backoff` already caps at the max, but a base above the cap
// is a config mistake worth surfacing rather than silently flattening.
let (retry_delay, max_retry_delay) = if retry_delay > max_retry_delay {
warn!(
diff --git a/core/connectors/sinks/influxdb_sink/README.md
b/core/connectors/sinks/influxdb_sink/README.md
index 37a8e0125..64813f8d8 100644
--- a/core/connectors/sinks/influxdb_sink/README.md
+++ b/core/connectors/sinks/influxdb_sink/README.md
@@ -69,10 +69,10 @@ verbose_logging = false
```toml
timeout = "30s" # per-request timeout
-max_retries = 3 # retries per write on transient errors
(429/5xx)
+max_retries = 3 # total write attempts, including the
first (429/5xx)
retry_delay = "1s" # initial backoff between retries
retry_max_delay = "5s" # backoff cap
-max_open_retries = 10 # retries during open() health check
+max_open_retries = 10 # total open() health-check attempts,
including the first
open_retry_max_delay = "60s" # backoff cap for open() retries
circuit_breaker_threshold = 5 # consecutive failures before circuit trips
circuit_breaker_cool_down = "30s" # how long circuit stays open before
half-open probe
@@ -119,6 +119,6 @@ circuit_breaker_cool_down = "15s"
The sink uses a layered design:
- **Batch accumulator**: messages are serialised to line protocol and buffered
until `batch_size` is reached, then flushed in a single HTTP POST.
-- **Retry middleware**: `reqwest-retry` with exponential backoff handles 429
and 5xx responses automatically before the connector-level retry logic runs.
+- **Retry middleware**: `iggy_connector_sdk::retry::HttpRetryMiddleware`
retries 429, 5xx and network errors with exponential backoff and jitter.
- **Circuit breaker**: after `circuit_breaker_threshold` consecutive failures
the connector stops issuing writes and waits for the cool-down window before
probing again.
- **Precision mapping**: V3's `/api/v3/write_lp` endpoint requires full
English words (`nanosecond`, `microsecond`, `millisecond`, `second`); the
connector maps the short forms automatically.
diff --git a/core/connectors/sinks/influxdb_sink/src/lib.rs
b/core/connectors/sinks/influxdb_sink/src/lib.rs
index 51a6828d1..2018bd90f 100644
--- a/core/connectors/sinks/influxdb_sink/src/lib.rs
+++ b/core/connectors/sinks/influxdb_sink/src/lib.rs
@@ -23,8 +23,7 @@ use base64::{Engine as _, engine::general_purpose};
use bytes::Bytes;
use iggy_common::serde_secret::serialize_secret;
use iggy_connector_sdk::retry::{
- CircuitBreaker, ConnectivityConfig, build_retry_client,
check_connectivity_with_retry,
- parse_duration,
+ CircuitBreaker, RetryPolicy, build_retry_client,
check_connectivity_with_retry, parse_duration,
};
use iggy_connector_sdk::{
ConsumedMessage, Error, MessagesMetadata, Sink, TopicMetadata,
sink_connector,
@@ -749,13 +748,13 @@ impl Sink for InfluxDbSink {
self.config.build_health_url()?,
"InfluxDB sink",
self.id,
- &ConnectivityConfig {
- max_open_retries: self.config.max_open_retries(),
- open_retry_max_delay: parse_duration(
+ RetryPolicy {
+ max_attempts: self.config.max_open_retries(),
+ base_delay: self.retry_delay,
+ max_delay: parse_duration(
self.config.open_retry_max_delay(),
DEFAULT_OPEN_RETRY_MAX_DELAY,
),
- retry_delay: self.retry_delay,
},
)
.await?;
diff --git a/core/connectors/sinks/meilisearch_sink/src/lib.rs
b/core/connectors/sinks/meilisearch_sink/src/lib.rs
index 64b42ab8c..27f0a04ae 100644
--- a/core/connectors/sinks/meilisearch_sink/src/lib.rs
+++ b/core/connectors/sinks/meilisearch_sink/src/lib.rs
@@ -20,7 +20,7 @@ use base64::{Engine as _, engine::general_purpose};
use iggy_common::IggyTimestamp;
use iggy_connector_sdk::{
ConsumedMessage, Error, MessagesMetadata, Payload, Sink, TopicMetadata,
- retry::{exponential_backoff, jitter, parse_duration},
+ retry::{parse_duration, retry_backoff},
sink_connector,
};
use meilisearch_sdk::{
@@ -228,11 +228,7 @@ impl MeilisearchSink {
)));
}
retries += 1;
- let delay = jitter(exponential_backoff(
- self.config.retry_delay,
- retries,
- self.config.max_retry_delay,
- ));
+ let delay = self.backoff(retries);
warn!(
"Meilisearch health check returned status '{}' (retry
{}/{}). Retrying in {:?}...",
health.status, retries, self.config.max_open_retries,
delay
@@ -246,11 +242,7 @@ impl MeilisearchSink {
return Err(map_sdk_error(error));
}
retries += 1;
- let delay = jitter(exponential_backoff(
- self.config.retry_delay,
- retries,
- self.config.max_retry_delay,
- ));
+ let delay = self.backoff(retries);
warn!(
"Meilisearch health check failed (retry {}/{}): {}.
Retrying in {:?}...",
retries, self.config.max_open_retries, error, delay
@@ -265,11 +257,7 @@ impl MeilisearchSink {
)));
}
retries += 1;
- let delay = jitter(exponential_backoff(
- self.config.retry_delay,
- retries,
- self.config.max_retry_delay,
- ));
+ let delay = self.backoff(retries);
warn!(
"Meilisearch health check timed out after {:?} (retry
{}/{}). Retrying in {:?}...",
self.config.timeout, retries,
self.config.max_open_retries, delay
@@ -642,6 +630,16 @@ impl MeilisearchSink {
)))
}
+ /// Backoff before retry number `retry` (1-based), from the configured
bounds.
+ ///
+ /// Only the delay is shared with the SDK. This connector's `max_retries` /
+ /// `max_open_retries` count retries *after* the first request, as the
+ /// README documents, unlike `RetryPolicy::max_attempts` which is a total.
+ /// Aligning them would silently cut every deployed budget by one.
+ fn backoff(&self, retry: u32) -> Duration {
+ retry_backoff(self.config.retry_delay, retry,
self.config.max_retry_delay)
+ }
+
async fn retry_sdk_operation<T, Fut, Op>(
&self,
operation: &str,
@@ -699,11 +697,7 @@ impl MeilisearchSink {
return Err(map_sdk_error(error));
}
retries += 1;
- let delay = jitter(exponential_backoff(
- self.config.retry_delay,
- retries,
- self.config.max_retry_delay,
- ));
+ let delay = self.backoff(retries);
warn!(
"Meilisearch {operation} failed (retry
{retries}/{max_retries}): {error}. Retrying in {delay:?}..."
);
@@ -721,11 +715,7 @@ impl MeilisearchSink {
)));
}
retries += 1;
- let delay = jitter(exponential_backoff(
- self.config.retry_delay,
- retries,
- self.config.max_retry_delay,
- ));
+ let delay = self.backoff(retries);
warn!(
"Meilisearch {operation} timed out after {:?} (retry
{retries}/{max_retries}). Retrying in {delay:?}...",
self.config.timeout
diff --git a/core/connectors/sinks/quickwit_sink/src/lib.rs
b/core/connectors/sinks/quickwit_sink/src/lib.rs
index cd5392d58..fadcb8b55 100644
--- a/core/connectors/sinks/quickwit_sink/src/lib.rs
+++ b/core/connectors/sinks/quickwit_sink/src/lib.rs
@@ -20,7 +20,7 @@ use std::time::Duration;
use async_trait::async_trait;
use base64::{Engine as _, engine::general_purpose};
use iggy_connector_sdk::retry::{
- ConnectivityConfig, build_retry_client, check_connectivity_with_retry,
is_transient_status,
+ RetryPolicy, build_retry_client, check_connectivity_with_retry,
is_transient_status,
};
use iggy_connector_sdk::{
ConsumedMessage, Error, MessagesMetadata, Payload, Sink, TopicMetadata,
sink_connector,
@@ -374,14 +374,14 @@ impl Sink for QuickwitSink {
endpoint_url(&base_url, &["health", "readyz"])?,
"Quickwit sink",
self.id,
- &ConnectivityConfig {
- max_open_retries: self
+ RetryPolicy {
+ max_attempts: self
.config
.max_open_retries
.unwrap_or(DEFAULT_MAX_OPEN_RETRIES)
.max(1),
- open_retry_max_delay,
- retry_delay,
+ base_delay: retry_delay,
+ max_delay: open_retry_max_delay,
},
)
.await?;
diff --git a/core/connectors/sinks/rabbitmq_sink/src/lib.rs
b/core/connectors/sinks/rabbitmq_sink/src/lib.rs
index 69504096e..a428d1250 100644
--- a/core/connectors/sinks/rabbitmq_sink/src/lib.rs
+++ b/core/connectors/sinks/rabbitmq_sink/src/lib.rs
@@ -17,7 +17,7 @@
use async_trait::async_trait;
use iggy::prelude::HeaderKind;
-use iggy_connector_sdk::retry::{exponential_backoff, jitter};
+use iggy_connector_sdk::retry::retry_backoff;
use iggy_connector_sdk::{
ConsumedMessage, Error, MessagesMetadata, Sink, TopicMetadata,
sink_connector,
};
@@ -229,11 +229,7 @@ impl RabbitMQSink {
"failed to reconnect: {reconnect_error}"
)));
}
- let delay = jitter(exponential_backoff(
- self.retry_delay,
- attempts.saturating_sub(1),
- self.max_retry_delay,
- ));
+ let delay = retry_backoff(self.retry_delay, attempts,
self.max_retry_delay);
warn!(
"RabbitMQ not connected for connector ID: {} (attempt
{attempts}/{}). Retrying in {:?}.",
self.id, self.max_retries, delay
@@ -372,11 +368,7 @@ impl RabbitMQSink {
}
}
- let delay = jitter(exponential_backoff(
- self.retry_delay,
- attempts.saturating_sub(1),
- self.max_retry_delay,
- ));
+ let delay = retry_backoff(self.retry_delay, attempts,
self.max_retry_delay);
warn!(
"Transient RabbitMQ publish error for connector ID: {}
(attempt {attempts}/{}): {error}. Retrying in {:?}.",
self.id, self.max_retries, delay
diff --git a/core/connectors/sinks/s3_sink/src/sink.rs
b/core/connectors/sinks/s3_sink/src/sink.rs
index 59b1a59f8..7ca3e1296 100644
--- a/core/connectors/sinks/s3_sink/src/sink.rs
+++ b/core/connectors/sinks/s3_sink/src/sink.rs
@@ -20,7 +20,7 @@ use crate::formatter;
use crate::path::{PathContext, render_s3_key};
use crate::{BufferKey, S3Sink};
use async_trait::async_trait;
-use iggy_connector_sdk::retry::{exponential_backoff, jitter};
+use iggy_connector_sdk::retry::retry_backoff;
use iggy_connector_sdk::{ConsumedMessage, Error, MessagesMetadata, Sink,
TopicMetadata};
use std::sync::Arc;
use std::time::Duration;
@@ -392,9 +392,7 @@ impl S3Sink {
);
}
}
- // exponential_backoff expects a 0-based retry index
- let retry_index = attempt - 1;
- let delay = jitter(exponential_backoff(base_delay, retry_index,
MAX_BACKOFF));
+ let delay = retry_backoff(base_delay, attempt, MAX_BACKOFF);
tokio::time::sleep(delay).await;
}
}
diff --git a/core/connectors/sinks/surrealdb_sink/src/lib.rs
b/core/connectors/sinks/surrealdb_sink/src/lib.rs
index 6e354e9ae..0d774083f 100644
--- a/core/connectors/sinks/surrealdb_sink/src/lib.rs
+++ b/core/connectors/sinks/surrealdb_sink/src/lib.rs
@@ -20,7 +20,7 @@ use base64::Engine;
use base64::engine::general_purpose;
use bytes::Bytes;
use iggy_connector_sdk::convert::owned_value_to_serde_json;
-use iggy_connector_sdk::retry::{exponential_backoff, jitter, parse_duration};
+use iggy_connector_sdk::retry::{parse_duration, retry_backoff};
use iggy_connector_sdk::{
ConsumedMessage, Error, MessagesMetadata, Payload, Sink, TopicMetadata,
sink_connector,
};
@@ -773,11 +773,7 @@ impl SurrealDbSink {
}
}
- let delay = jitter(exponential_backoff(
- self.retry_delay,
- attempts.saturating_sub(1),
- self.max_retry_delay,
- ));
+ let delay = retry_backoff(self.retry_delay, attempts,
self.max_retry_delay);
warn!(
"Transient SurrealDB write error for connector ID: {}
(attempt {attempts}/{}): {error}. Retrying in {:?}.",
self.id, self.max_retries, delay
diff --git a/core/connectors/sources/influxdb_source/README.md
b/core/connectors/sources/influxdb_source/README.md
index f7110fc3a..333e4f3d7 100644
--- a/core/connectors/sources/influxdb_source/README.md
+++ b/core/connectors/sources/influxdb_source/README.md
@@ -108,10 +108,10 @@ verbose_logging = false
```toml
timeout = "10s" # per-request timeout
-max_retries = 3 # retries per query on transient errors
(429/5xx)
+max_retries = 3 # total query attempts, including the
first (429/5xx)
retry_delay = "1s" # initial backoff
retry_max_delay = "5s" # backoff cap
-max_open_retries = 10 # retries during open() health check
+max_open_retries = 10 # total open() health-check attempts,
including the first
open_retry_max_delay = "60s" # backoff cap for open() retries
circuit_breaker_threshold = 5 # consecutive failures before circuit trips
circuit_breaker_cool_down = "30s" # cool-down before half-open probe
diff --git a/core/connectors/sources/influxdb_source/src/lib.rs
b/core/connectors/sources/influxdb_source/src/lib.rs
index ac0ed9fda..59a1b5e12 100644
--- a/core/connectors/sources/influxdb_source/src/lib.rs
+++ b/core/connectors/sources/influxdb_source/src/lib.rs
@@ -31,8 +31,7 @@ use common::{
validate_cursor_field,
};
use iggy_connector_sdk::retry::{
- CircuitBreaker, ConnectivityConfig, build_retry_client,
check_connectivity_with_retry,
- parse_duration,
+ CircuitBreaker, RetryPolicy, build_retry_client,
check_connectivity_with_retry, parse_duration,
};
use iggy_connector_sdk::{
ConnectorState, Error, ProducedMessages, Schema, Source, source_connector,
@@ -420,13 +419,13 @@ impl Source for InfluxDbSource {
health_url,
CONNECTOR_NAME,
self.id,
- &ConnectivityConfig {
- max_open_retries: self.config.max_open_retries(),
- open_retry_max_delay: parse_duration(
+ RetryPolicy {
+ max_attempts: self.config.max_open_retries(),
+ base_delay: self.retry_delay,
+ max_delay: parse_duration(
self.config.open_retry_max_delay(),
DEFAULT_OPEN_RETRY_MAX_DELAY,
),
- retry_delay: self.retry_delay,
},
)
.await?;