Packages
Typed distributed messaging for Gleam on the BEAM.
Retired package: Deprecated - The project needs to be redesigned around a much smaller and clearer core.
Current section
Files
Jump to
Current section
Files
src/distribute@retry.erl
-module(distribute@retry).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/distribute/retry.gleam").
-export([default/0, default_with_jitter/0, no_retry/0, aggressive/0, conservative/0, with_max_attempts/2, with_base_delay_ms/2, with_max_delay_ms/2, with_multiplier/2, with_jitter/2, with_full_jitter/1, should_retry/2, is_final_attempt/2, total_attempts/1, calculate_delay/2, delay_ms/2, jitter_to_string/1, policy_to_string/1]).
-export_type([jitter_strategy/0, retry_policy/0, delay_result/0]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
?MODULEDOC(
" Retry policy with exponential backoff and jitter.\n"
"\n"
" This module provides type-safe retry logic for distributed operations.\n"
" It implements industry best practices for handling transient failures\n"
" in distributed systems.\n"
"\n"
" ## Jitter Strategies\n"
"\n"
" Jitter is crucial to prevent the \"thundering herd\" problem where many\n"
" clients retry simultaneously after a failure. This module supports\n"
" multiple jitter strategies as recommended by AWS and Google Cloud.\n"
"\n"
" - `NoJitter`: Deterministic exponential backoff (not recommended)\n"
" - `FullJitter`: `random(0, delay)` - Best for reducing contention\n"
" - `EqualJitter`: `delay/2 + random(0, delay/2)` - Balanced approach\n"
" - `DecorrelatedJitter`: `random(base, prev_delay * 3)` - Good for APIs\n"
"\n"
" ## References\n"
"\n"
" - AWS: https://aws.amazon.com/blogs/architecture/exponential-backoff-and-jitter/\n"
" - Google Cloud: https://cloud.google.com/storage/docs/retry-strategy\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" import distribute/retry\n"
"\n"
" // Create a policy with full jitter (recommended)\n"
" let policy = retry.default_with_jitter()\n"
"\n"
" // Calculate delay for attempt 3\n"
" let delay = retry.calculate_delay(policy, 3)\n"
" // delay will be random(0, min(5000, 100 * 2^2)) = random(0, 400)\n"
" ```\n"
).
-type jitter_strategy() :: no_jitter |
full_jitter |
equal_jitter |
decorrelated_jitter.
-type retry_policy() :: {retry_policy,
integer(),
integer(),
integer(),
float(),
jitter_strategy()}.
-type delay_result() :: {delay_result,
integer(),
integer(),
integer(),
boolean()}.
-file("src/distribute/retry.gleam", 130).
?DOC(
" Default retry policy without jitter.\n"
"\n"
" Conservative settings suitable for general use:\n"
" - 3 attempts total\n"
" - 100ms base delay\n"
" - 5000ms max delay\n"
" - 2.0 multiplier (exponential)\n"
" - No jitter\n"
"\n"
" For distributed systems, prefer `default_with_jitter()`.\n"
).
-spec default() -> retry_policy().
default() ->
{retry_policy, 3, 100, 5000, 2.0, no_jitter}.
-file("src/distribute/retry.gleam", 150).
?DOC(
" Default retry policy with full jitter (RECOMMENDED).\n"
"\n"
" Same as `default()` but with FullJitter enabled.\n"
" This is the recommended policy for distributed systems to prevent\n"
" thundering herd problems.\n"
"\n"
" Delay progression (example, actual values are randomized):\n"
" - Attempt 1: random(0, 100ms)\n"
" - Attempt 2: random(0, 200ms)\n"
" - Attempt 3: random(0, 400ms)\n"
).
-spec default_with_jitter() -> retry_policy().
default_with_jitter() ->
{retry_policy, 3, 100, 5000, 2.0, full_jitter}.
-file("src/distribute/retry.gleam", 164).
?DOC(
" No retry policy (single attempt only).\n"
"\n"
" Use when retry is handled at a different layer or for operations\n"
" that should not be retried (e.g., non-idempotent operations).\n"
).
-spec no_retry() -> retry_policy().
no_retry() ->
{retry_policy, 1, 0, 0, 1.0, no_jitter}.
-file("src/distribute/retry.gleam", 181).
?DOC(
" Aggressive retry policy for critical operations.\n"
"\n"
" More retries with shorter base delays:\n"
" - 5 attempts total\n"
" - 50ms base delay\n"
" - 3000ms max delay\n"
" - Full jitter enabled\n"
).
-spec aggressive() -> retry_policy().
aggressive() ->
{retry_policy, 5, 50, 3000, 2.0, full_jitter}.
-file("src/distribute/retry.gleam", 198).
?DOC(
" Conservative retry policy for non-critical operations.\n"
"\n"
" Fewer retries with longer delays:\n"
" - 2 attempts total\n"
" - 500ms base delay\n"
" - 10000ms max delay\n"
" - Full jitter enabled\n"
).
-spec conservative() -> retry_policy().
conservative() ->
{retry_policy, 2, 500, 10000, 2.0, full_jitter}.
-file("src/distribute/retry.gleam", 218).
?DOC(
" Set the maximum number of retry attempts.\n"
"\n"
" ```gleam\n"
" retry.default()\n"
" |> retry.with_max_attempts(5)\n"
" ```\n"
).
-spec with_max_attempts(retry_policy(), integer()) -> retry_policy().
with_max_attempts(Policy, Attempts) ->
{retry_policy,
gleam@int:max(1, Attempts),
erlang:element(3, Policy),
erlang:element(4, Policy),
erlang:element(5, Policy),
erlang:element(6, Policy)}.
-file("src/distribute/retry.gleam", 228).
?DOC(
" Set the base delay in milliseconds.\n"
"\n"
" ```gleam\n"
" retry.default()\n"
" |> retry.with_base_delay_ms(200)\n"
" ```\n"
).
-spec with_base_delay_ms(retry_policy(), integer()) -> retry_policy().
with_base_delay_ms(Policy, Delay_ms) ->
{retry_policy,
erlang:element(2, Policy),
gleam@int:max(0, Delay_ms),
erlang:element(4, Policy),
erlang:element(5, Policy),
erlang:element(6, Policy)}.
-file("src/distribute/retry.gleam", 238).
?DOC(
" Set the maximum delay cap in milliseconds.\n"
"\n"
" ```gleam\n"
" retry.default()\n"
" |> retry.with_max_delay_ms(10_000)\n"
" ```\n"
).
-spec with_max_delay_ms(retry_policy(), integer()) -> retry_policy().
with_max_delay_ms(Policy, Delay_ms) ->
{retry_policy,
erlang:element(2, Policy),
erlang:element(3, Policy),
gleam@int:max(0, Delay_ms),
erlang:element(5, Policy),
erlang:element(6, Policy)}.
-file("src/distribute/retry.gleam", 248).
?DOC(
" Set the backoff multiplier.\n"
"\n"
" ```gleam\n"
" retry.default()\n"
" |> retry.with_multiplier(1.5) // Slower growth\n"
" ```\n"
).
-spec with_multiplier(retry_policy(), float()) -> retry_policy().
with_multiplier(Policy, Multiplier) ->
Safe_multiplier = case Multiplier < 1.0 of
true ->
1.0;
false ->
Multiplier
end,
{retry_policy,
erlang:element(2, Policy),
erlang:element(3, Policy),
erlang:element(4, Policy),
Safe_multiplier,
erlang:element(6, Policy)}.
-file("src/distribute/retry.gleam", 262).
?DOC(
" Set the jitter strategy.\n"
"\n"
" ```gleam\n"
" retry.default()\n"
" |> retry.with_jitter(retry.FullJitter)\n"
" ```\n"
).
-spec with_jitter(retry_policy(), jitter_strategy()) -> retry_policy().
with_jitter(Policy, Jitter) ->
{retry_policy,
erlang:element(2, Policy),
erlang:element(3, Policy),
erlang:element(4, Policy),
erlang:element(5, Policy),
Jitter}.
-file("src/distribute/retry.gleam", 272).
?DOC(
" Enable full jitter (convenience method).\n"
"\n"
" ```gleam\n"
" retry.default()\n"
" |> retry.with_full_jitter()\n"
" ```\n"
).
-spec with_full_jitter(retry_policy()) -> retry_policy().
with_full_jitter(Policy) ->
with_jitter(Policy, full_jitter).
-file("src/distribute/retry.gleam", 349).
?DOC(
" Check if we should retry after this attempt.\n"
"\n"
" Returns True if `attempt < max_attempts`.\n"
"\n"
" ```gleam\n"
" case retry.should_retry(policy, attempt) {\n"
" True -> {\n"
" process.sleep(retry.delay_ms(policy, attempt))\n"
" try_operation(attempt + 1)\n"
" }\n"
" False -> Error(MaxRetriesExceeded)\n"
" }\n"
" ```\n"
).
-spec should_retry(retry_policy(), integer()) -> boolean().
should_retry(Policy, Attempt) ->
Attempt < erlang:element(2, Policy).
-file("src/distribute/retry.gleam", 361).
?DOC(
" Check if this is the final attempt.\n"
"\n"
" ```gleam\n"
" case retry.is_final_attempt(policy, attempt) {\n"
" True -> log.error(\"Final attempt, no more retries\")\n"
" False -> log.warn(\"Retrying...\")\n"
" }\n"
" ```\n"
).
-spec is_final_attempt(retry_policy(), integer()) -> boolean().
is_final_attempt(Policy, Attempt) ->
Attempt >= erlang:element(2, Policy).
-file("src/distribute/retry.gleam", 371).
?DOC(
" Get the total number of attempts that will be made.\n"
"\n"
" ```gleam\n"
" let total = retry.total_attempts(policy)\n"
" log.info(\"Will try up to \" <> int.to_string(total) <> \" times\")\n"
" ```\n"
).
-spec total_attempts(retry_policy()) -> integer().
total_attempts(Policy) ->
erlang:element(2, Policy).
-file("src/distribute/retry.gleam", 422).
?DOC(" Generate a random integer in range [min, max] (inclusive).\n").
-spec random_int(integer(), integer()) -> integer().
random_int(Min, Max) ->
case Max =< Min of
true ->
Min;
false ->
(Min + rand:uniform((Max - Min) + 1)) - 1
end.
-file("src/distribute/retry.gleam", 380).
?DOC(" Apply jitter strategy to a delay value.\n").
-spec apply_jitter(jitter_strategy(), integer(), integer()) -> integer().
apply_jitter(Strategy, Delay_ms, Base_delay_ms) ->
case Strategy of
no_jitter ->
Delay_ms;
full_jitter ->
case Delay_ms =< 0 of
true ->
0;
false ->
random_int(0, Delay_ms)
end;
equal_jitter ->
Half = Delay_ms div 2,
case Half =< 0 of
true ->
Delay_ms;
false ->
Half + random_int(0, Half)
end;
decorrelated_jitter ->
Upper = gleam@int:min(Delay_ms * 3, Delay_ms + (Base_delay_ms * 10)),
case Upper =< Base_delay_ms of
true ->
Base_delay_ms;
false ->
random_int(Base_delay_ms, Upper)
end
end.
-file("src/distribute/retry.gleam", 434).
?DOC(" Float power helper for small integer exponents.\n").
-spec float_power(float(), float()) -> float().
float_power(Base, Exponent) ->
Exp_int = erlang:trunc(Exponent),
case Exp_int of
0 ->
1.0;
1 ->
Base;
2 ->
Base * Base;
3 ->
(Base * Base) * Base;
4 ->
((Base * Base) * Base) * Base;
5 ->
(((Base * Base) * Base) * Base) * Base;
6 ->
((((Base * Base) * Base) * Base) * Base) * Base;
7 ->
(((((Base * Base) * Base) * Base) * Base) * Base) * Base;
8 ->
((((((Base * Base) * Base) * Base) * Base) * Base) * Base) * Base;
_ ->
((((((Base * Base) * Base) * Base) * Base) * Base) * Base) * Base
end.
-file("src/distribute/retry.gleam", 299).
?DOC(
" Calculate the delay for a given attempt number.\n"
"\n"
" Returns a `DelayResult` containing the delay in milliseconds and metadata.\n"
" The attempt number should be 1-indexed (first retry is attempt 1).\n"
"\n"
" ## Algorithm\n"
"\n"
" 1. Calculate base exponential delay: `base_delay * (multiplier ^ (attempt - 1))`\n"
" 2. Cap at `max_delay_ms`\n"
" 3. Apply jitter strategy\n"
"\n"
" ## Example\n"
"\n"
" ```gleam\n"
" let policy = retry.default_with_jitter()\n"
" let result = retry.calculate_delay(policy, 1)\n"
" // result.delay_ms will be random(0, 100)\n"
" process.sleep(result.delay_ms)\n"
" ```\n"
).
-spec calculate_delay(retry_policy(), integer()) -> delay_result().
calculate_delay(Policy, Attempt) ->
Safe_attempt = gleam@int:max(1, Attempt),
Base = erlang:float(erlang:element(3, Policy)),
Exponent = erlang:float(Safe_attempt - 1),
Multiplier_pow = float_power(erlang:element(5, Policy), Exponent),
Delay_float = Base * Multiplier_pow,
Exponential_delay = erlang:trunc(Delay_float),
Capped_delay = gleam@int:min(Exponential_delay, erlang:element(4, Policy)),
Final_delay = apply_jitter(
erlang:element(6, Policy),
Capped_delay,
erlang:element(3, Policy)
),
{delay_result,
Final_delay,
Capped_delay,
Safe_attempt,
Safe_attempt >= erlang:element(2, Policy)}.
-file("src/distribute/retry.gleam", 332).
?DOC(
" Get just the delay value in milliseconds.\n"
"\n"
" Convenience function when you don't need the full DelayResult.\n"
"\n"
" ```gleam\n"
" let delay = retry.delay_ms(policy, attempt)\n"
" process.sleep(delay)\n"
" ```\n"
).
-spec delay_ms(retry_policy(), integer()) -> integer().
delay_ms(Policy, Attempt) ->
erlang:element(2, calculate_delay(Policy, Attempt)).
-file("src/distribute/retry.gleam", 460).
?DOC(" Convert jitter strategy to string for logging.\n").
-spec jitter_to_string(jitter_strategy()) -> binary().
jitter_to_string(Jitter) ->
case Jitter of
no_jitter ->
<<"none"/utf8>>;
full_jitter ->
<<"full"/utf8>>;
equal_jitter ->
<<"equal"/utf8>>;
decorrelated_jitter ->
<<"decorrelated"/utf8>>
end.
-file("src/distribute/retry.gleam", 470).
?DOC(" Convert retry policy to a loggable string representation.\n").
-spec policy_to_string(retry_policy()) -> binary().
policy_to_string(Policy) ->
<<<<<<<<<<<<<<<<"RetryPolicy(max_attempts="/utf8,
(erlang:integer_to_binary(
erlang:element(2, Policy)
))/binary>>/binary,
", base_delay_ms="/utf8>>/binary,
(erlang:integer_to_binary(erlang:element(3, Policy)))/binary>>/binary,
", max_delay_ms="/utf8>>/binary,
(erlang:integer_to_binary(erlang:element(4, Policy)))/binary>>/binary,
", jitter="/utf8>>/binary,
(jitter_to_string(erlang:element(6, Policy)))/binary>>/binary,
")"/utf8>>.