spetz commented on code in PR #4104:
URL: https://github.com/apache/iggy/pull/4104#discussion_r3978742187


##########
core/connectors/sdk/src/retry.rs:
##########
@@ -181,6 +190,147 @@ 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
+    }
+}
+
+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"
+        };
+        write!(
+            f,
+            "{reason} after {} attempts: {}",
+            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, Op, Fut>(
+    policy: RetryPolicy,
+    context: &str,
+    should_retry: impl Fn(&E) -> bool,
+    mut operation: Op,
+) -> Result<T, RetryFailure<E>>
+where
+    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 {
+                    info!("{context} succeeded after {attempt} retries.");
+                }
+                return Ok(value);
+            }
+            Err(error) => error,
+        };
+
+        attempt += 1;
+        let retryable = should_retry(&error);
+        let exhausted = attempt >= max_attempts;

Review Comment:
    `exhausted` is computed without `retryable`, so a permanent error on the 
last attempt reports `exhausted: true` and `Display` says "ran out of attempts".



##########
core/connectors/sdk/src/retry.rs:
##########
@@ -161,16 +165,21 @@ 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))
 }
 
+/// Ceiling for a server-supplied `Retry-After`.

Review Comment:
   doc line "Ceiling for a server-supplied `Retry-After`." refers to a constant 
that does not exist in this PR



##########
core/connectors/sinks/doris_sink/src/lib.rs:
##########
@@ -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 each attempt and

Review Comment:
   Comment says `retry_async` logs the terminal outcome. It does not 
(`retry.rs:283-286`), and line 438 in this file says `consume` does.



##########
core/connectors/sdk/src/retry.rs:
##########
@@ -398,49 +529,296 @@ 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| {
+        // Sole record of the cause: `open()`'s Err is dropped at the FFI
+        // boundary and the runtime logs only "Plugin initialization failed".
+        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_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:?}"
+        );
+    }
+
+    /// The regression test for the convention this helper exists to unify: the

Review Comment:
   Comment block on `retry_client` describes tests not present in this module 
(Retry-After bound, statuses the header is read on).



##########
core/connectors/sdk/src/retry.rs:
##########
@@ -181,6 +190,147 @@ 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
+    }
+}
+
+impl<E: fmt::Display> fmt::Display for RetryFailure<E> {

Review Comment:
    `RetryFailure<E>` implements `Display` but not `std::error::Error`, so it 
cannot be propagated with `?` into `anyhow::Error` or `Box<dyn Error>`.



##########
core/connectors/sinks/doris_sink/src/lib.rs:
##########
@@ -417,36 +419,25 @@ 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);
-            warn!(
-                "Doris sink ID {} transient Stream Load failure on attempt 
{attempt}/{} (label={label}): {error}; retrying in {delay:?}",
-                self.id, connected.max_retries
-            );
-            tokio::time::sleep(delay).await;
-        }
+        })
+        .await
+        // `consume` already logs the terminal error for the batch, so the
+        // attempt bookkeeping is dropped here rather than logged twice.
+        .map_err(RetryFailure::into_error)

Review Comment:
    `into_error` drops `attempts` and `exhausted`, terminal log at line 1178 no 
longer says whether retries ran out or the error was permanent.



##########
core/connectors/sdk/src/retry.rs:
##########
@@ -181,6 +190,147 @@ 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
+    }
+}
+
+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"
+        };
+        write!(
+            f,
+            "{reason} after {} attempts: {}",
+            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, Op, Fut>(
+    policy: RetryPolicy,
+    context: &str,
+    should_retry: impl Fn(&E) -> bool,

Review Comment:
   `impl Fn(&E) -> bool` in argument position mixed with explicit generics `<T, 
E, Op, Fut>` makes `retry_async` impossible to turbofish (E0632)



##########
core/connectors/sdk/src/retry.rs:
##########
@@ -398,49 +529,296 @@ 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| {
+        // Sole record of the cause: `open()`'s Err is dropped at the FFI
+        // boundary and the runtime logs only "Plugin initialization failed".
+        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_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:?}"
+        );
+    }
+
+    /// The regression test for the convention this helper exists to unify: the
+    /// first retry waits the configured base delay, not twice it.
+    // `HttpRetryMiddleware` is the path every HTTP connector rides, and the
+    // behaviours it covers are the ones that regressed unnoticed: the bound on
+    // `Retry-After`, the statuses the header is read on, and the first delay.
+    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!(
+            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.
+        assert!(
+            elapsed < base.mul_f64(1.4),

Review Comment:
   Upper bound `1.4 × base` leaves 200ms over max jitter (1.2s) for two 
wiremock round trips; flaky on loaded CI. Bug band starts at 1.6s.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to