ryankert01 commented on code in PR #4104:
URL: https://github.com/apache/iggy/pull/4104#discussion_r3975905884
##########
core/connectors/sdk/src/retry.rs:
##########
@@ -181,6 +190,114 @@ 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 [`jitter`], 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)
+}
+/// Run `operation`, retrying while it fails with an error `is_transient`
+/// 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.
+///
+/// Retries, recovery and giving up are all logged here, with the attempt
+/// counts callers would otherwise each have to track. The last error is
+/// returned unchanged so callers keep their own error classification.
+pub async fn retry_async<T, E, Op, Fut>(
Review Comment:
Good idea on the rename — `should_retry` is clearer, will do.
I'd keep `context` though. Plugin logs cross FFI via `CallbackLayer`, which
drops structured fields (checked: `warn!(attempt = 3, "retrying")` arrives as
just `"retrying"`) and sets target to the module path — so retries land as
`iggy_connector_sdk::retry` with no instance identity. Doris puts `label=`
there; two InfluxDB sinks in one runtime would otherwise be indistinguishable.
Agree the helper shouldn't own the terminal log though — I'll return
`RetryFailure { error, attempts, exhausted }` so callers can write their own.
--
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]