jiengup commented on code in PR #4104:
URL: https://github.com/apache/iggy/pull/4104#discussion_r3971448886
##########
core/connectors/sdk/src/retry.rs:
##########
@@ -256,34 +376,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));
Review Comment:
This logic seems to lack test cases. Let's add `Retry-After` takes
precedence over local backoff tests.
##########
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:
As a common util, I think this interface should be designed to be simpler
and easier to use. Consider renaming the `is_transient` parameter and removing
the `context` parameter, just letting users to print logs in `Op` themselves.
--
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]