From 2c06d6609efc36e19125b39ae7b8a30c9b23ae9e Mon Sep 17 00:00:00 2001 From: Brent Date: Sun, 27 Sep 2026 01:04:45 -0400 Subject: [PATCH] SMOODEV-3375: Make retries method-aware and make timeouts cancel the attempt Retries were method-blind in every port: a POST that timed out or got a 429/5xx was re-sent up to twice more. In TypeScript the per-attempt timeout (mollitia's Timeout module) only raced the request, so the losing attempt kept running server-side while the retry sent it again. Image generation (20-50s) against the 10s default billed three images and returned nothing, and any non-idempotent POST/PATCH could duplicate its side effect. Now, in all five ports, only idempotent methods (RFC 9110 9.2.2: GET, HEAD, OPTIONS, TRACE, PUT, DELETE) retry. Anything else makes exactly one attempt and surfaces its own error, unless the caller opts in with retry.allowNonIdempotent (allow_non_idempotent / AllowNonIdempotent) or the request carries a non-empty Idempotency-Key header. Eligibility is decided after pre-request hooks and the auth provider, and onRejection cannot override it. The client-side rate limiter's retry loop is unaffected: it rejects before anything is sent. TypeScript now aborts each attempt through an AbortController combined with the caller's signal, and still races the attempt so a fetch that ignores the signal times out on schedule. A caller-aborted request is no longer retried. Go waits for the cancelled attempt to unwind before retrying; Python adds a hard per-attempt deadline; Rust and .NET already cancelled and now prove it. spec/retry-idempotency-corpus.json pins the rule; every port runs it against a real local server and counts what the server received, plus a test that the timed-out attempt's connection is actually closed. Major bump: callers relying on POST retries lose them, and the new Rust RetryOptions field breaks struct literals. Co-Authored-By: Claude Opus 5.5 (1M context) Claude-Session: https://claude.ai/code/session_01CfZaWemmpghhtofBauYdti --- .changeset/method-aware-retries.md | 55 +++ README.md | 59 ++- .../RetryIdempotencyTests.cs | 432 ++++++++++++++++++ .../SmooAI.Fetch.Tests.csproj | 1 + dotnet/SmooAI.Fetch/README.md | 36 +- dotnet/SmooAI.Fetch/RetryPolicy.cs | 61 ++- dotnet/SmooAI.Fetch/SmooAI.Fetch.csproj | 2 +- dotnet/SmooAI.Fetch/SmooFetch.cs | 50 +- go/fetch/README.md | 54 ++- go/fetch/client.go | 37 +- go/fetch/options.go | 10 + go/fetch/retry.go | 74 +++ go/fetch/retry_idempotency_test.go | 351 ++++++++++++++ go/fetch/timeout.go | 18 + python/README.md | 43 +- python/src/smooai_fetch/__init__.py | 4 +- python/src/smooai_fetch/_client.py | 55 ++- python/src/smooai_fetch/_retry.py | 40 +- python/src/smooai_fetch/_types.py | 13 + python/tests/test_retry_idempotency.py | 232 ++++++++++ rust/fetch/README.md | 42 +- rust/fetch/src/client.rs | 19 +- rust/fetch/src/defaults.rs | 1 + rust/fetch/src/lib.rs | 1 + rust/fetch/src/retry.rs | 78 +++- rust/fetch/src/types.rs | 25 + rust/fetch/tests/integration_tests.rs | 3 + rust/fetch/tests/rate_limit_retry_tests.rs | 2 + rust/fetch/tests/retry_idempotency_tests.rs | 391 ++++++++++++++++ rust/fetch/tests/retry_tests.rs | 12 + spec/retry-idempotency-corpus.json | 115 +++++ src/fetch.retry-safety.spec.ts | 207 +++++++++ src/fetch.spec.ts | 13 +- src/fetch.ts | 174 +++++-- 34 files changed, 2609 insertions(+), 101 deletions(-) create mode 100644 .changeset/method-aware-retries.md create mode 100644 dotnet/SmooAI.Fetch.Tests/RetryIdempotencyTests.cs create mode 100644 go/fetch/retry_idempotency_test.go create mode 100644 python/tests/test_retry_idempotency.py create mode 100644 rust/fetch/tests/retry_idempotency_tests.rs create mode 100644 spec/retry-idempotency-corpus.json create mode 100644 src/fetch.retry-safety.spec.ts diff --git a/.changeset/method-aware-retries.md b/.changeset/method-aware-retries.md new file mode 100644 index 0000000..6a1b20e --- /dev/null +++ b/.changeset/method-aware-retries.md @@ -0,0 +1,55 @@ +--- +'@smooai/fetch': major +--- + +SMOODEV-3375: Retries never duplicate a side effect: they now check the HTTP method, and timeouts cancel the attempt. This changes a default, so it is a major release. + +**What was wrong.** Retries ignored the HTTP method. A `POST` that timed out or got a 429/5xx was sent again, up to twice more. In TypeScript the per-attempt timeout (mollitia's `Timeout` module) did not stop the losing request, so the first attempt kept running on the server while the retry sent it again. Image generation takes 20-50s against the 10s default timeout, so it billed three images and returned nothing. Any non-idempotent call (CRM writes, message sends, payments) could run its side effect more than once. + +**New default, in all five languages:** + +- **Only idempotent methods are retried:** `GET`, `HEAD`, `OPTIONS`, `TRACE`, `PUT` and `DELETE` (RFC 9110 ยง9.2.2). +- **`POST`, `PATCH` and any other method make exactly one attempt.** You get that attempt's own error (`HTTPResponseError` / `TimeoutError` and the equivalents in other languages), not the retries-exhausted wrapper. +- **A 429 with `Retry-After` on a `POST` is not retried either.** +- **`onRejection` is never called for these requests**, so it cannot turn their retries back on. + +There are two ways to opt back in: + +- Send a non-empty `Idempotency-Key` header, in any casing. The server then deduplicates. +- Set `retry: { allowNonIdempotent: true }`. The name in each language: + - TypeScript: `allowNonIdempotent` + - Python and Rust: `allow_non_idempotent` + - Go and .NET: `AllowNonIdempotent` + +Eligibility is checked after the pre-request hooks and the auth provider run, so a hook can add the key. The client-side rate limiter's retry loop is unchanged, because it rejects requests before anything is sent. + +**Timeouts cancel the attempt.** TypeScript now aborts each attempt with an `AbortController`, combined with the caller's own `signal`, so the connection is closed before any retry. A `fetch` that ignores `signal` still times out on schedule. Rust, Go, Python and .NET already cancelled the attempt. A new test in every language proves the first attempt's connection closes. Go now also waits for the cancelled attempt to finish before it retries. Python now puts a hard deadline on each attempt with `asyncio.timeout`. Before, httpx only had per-phase timeouts, so a server that kept sending bytes slowly never timed out. + +**Also fixed:** in TypeScript, a request the caller aborted is no longer retried. + +**Other API changes:** + +- **TypeScript:** + - New exports: `isIdempotentMethod`, `isRetryEligible`, `IDEMPOTENCY_KEY_HEADER` and the `RetryOptions` type. + - `options.retry` and `FetchBuilder.withRetry` now take a partial. It is merged over the defaults, so `retry: { allowNonIdempotent: true }` keeps every other default. + - The `signal` that `fetch` receives now combines the caller's signal with the timeout's. It is no longer the caller's own object. +- **Rust:** + - `RetryOptions` gains `allow_non_idempotent`. This breaks existing `RetryOptions { .. }` struct literals. + - `RetryOptions` now implements `Default`, so literals can end with `..Default::default()` from now on. + - New: `is_idempotent_method`, `is_retry_eligible` and `IDEMPOTENCY_KEY_HEADER`. +- **Go:** + - New: `RetryOptions.AllowNonIdempotent`, `IsIdempotentMethod` and `IdempotencyKeyHeader`. + - The module path moves to `/v4`. +- **Python:** + - New: `RetryOptions.allow_non_idempotent`. + - `is_idempotent_method` and `IDEMPOTENCY_KEY_HEADER` are exported. +- **.NET:** + - New on `RetryPolicy`: `AllowNonIdempotent`, `IsIdempotentMethod`, `IsRetryEligible` and `IdempotencyKeyHeader`. + - `Microsoft.SourceLink.GitHub` goes to 10.0.303. 8.0.0 pulls in `Microsoft.Build.Tasks.Git` 8.0.0, which has an advisory against it (GHSA-23fw-v26w-5fgq), and a clean restore now fails with NU1902. + +**Why a major version:** + +- A caller that relied on `POST` retries quietly loses them. +- In Rust, adding a field to a struct that callers construct breaks their code. The 3.7.1 entry below explains why that must not ship as a minor. + +The shared `spec/retry-idempotency-corpus.json` pins the rule. Each language's tests run it against a real local server and count the requests the server received. diff --git a/README.md b/README.md index 6e6399f..97b3591 100644 --- a/README.md +++ b/README.md @@ -50,8 +50,8 @@ Traditional `fetch` gives you the request, but leaves you to handle the reality One resilient HTTP client, ported natively to five languages. Every port carries the same core behaviors โ€” verified against the source of each port, not aspirational: -- ๐Ÿ”„ **Smart retries** โ€” exponential backoff with jitter to prevent thundering herds; retries only on network errors and retryable statuses -- โฑ๏ธ **Automatic timeouts** โ€” never hang indefinitely on slow endpoints (10s default, configurable per request) +- ๐Ÿ”„ **Smart retries** โ€” exponential backoff with jitter to prevent thundering herds; retries only on network errors and retryable statuses, and **only for idempotent methods** unless you opt in โ€” a retry never duplicates a POST +- โฑ๏ธ **Automatic timeouts** โ€” never hang indefinitely on slow endpoints (10s default, configurable per request); a timed-out attempt is cancelled, not abandoned - ๐Ÿšฆ **Rate-limit respect** โ€” reads `Retry-After` headers and waits exactly what the server asked, plus a client-side sliding-window rate limiter - ๐Ÿ”Œ **Circuit breaking** โ€” stop hammering services that are clearly down - ๐Ÿ”— **Lifecycle hooks** โ€” pre-request / post-response hooks for auth, logging, and metrics @@ -63,14 +63,15 @@ One resilient HTTP client, ported natively to five languages. Every port carries Each capability in a few lines of real, current API โ€” snippets are verified against [`src/`](./src/) and the language ports, not pseudocode. -| | Capability | What you get | -| --- | ---------------------------------------------------- | ------------------------------------------------------ | -| ๐Ÿ”„ | [**Smart retries**](#-smart-retries) | Backoff + jitter, only on errors worth retrying | -| ๐Ÿšฆ | [**Rate-limit respect**](#-rate-limit-respect) | `Retry-After` honored to the second, in all five ports | -| ๐Ÿ”Œ | [**Circuit breaking**](#-circuit-breaking) | Fail fast when a dependency is down | -| ๐ŸŽฏ | [**Typed responses**](#-typed-responses--validation) | Schema-validated data, typed end to end | -| ๐Ÿ”— | [**Hooks + auth**](#-lifecycle-hooks--auth) | One place for tokens, logging, and response policy | -| ๐Ÿ“ก | [**Trace propagation**](#-trace-context-propagation) | `traceparent` on every request, optional OpenTelemetry | +| | Capability | What you get | +| --- | ----------------------------------------------------------- | ------------------------------------------------------ | +| ๐Ÿ”„ | [**Smart retries**](#-smart-retries) | Backoff + jitter, only on errors worth retrying | +| ๐Ÿงฏ | [**Safe retries**](#-retries-never-duplicate-a-side-effect) | Never re-sends a POST unless you opt in | +| ๐Ÿšฆ | [**Rate-limit respect**](#-rate-limit-respect) | `Retry-After` honored to the second, in all five ports | +| ๐Ÿ”Œ | [**Circuit breaking**](#-circuit-breaking) | Fail fast when a dependency is down | +| ๐ŸŽฏ | [**Typed responses**](#-typed-responses--validation) | Schema-validated data, typed end to end | +| ๐Ÿ”— | [**Hooks + auth**](#-lifecycle-hooks--auth) | One place for tokens, logging, and response policy | +| ๐Ÿ“ก | [**Trace propagation**](#-trace-context-propagation) | `traceparent` on every request, optional OpenTelemetry | ### ๐Ÿ”„ Smart retries @@ -88,6 +89,36 @@ const response = await fetch('https://flaky-api.com/data'); Defaults (TypeScript): 2 automatic retries, exponential backoff starting at 500ms with factor 2, jitter to prevent thundering herds, and retries only on network errors or retryable HTTP statuses. +### ๐Ÿงฏ Retries never duplicate a side effect + +**Changed in 4.0.0.** Only idempotent methods are retried: `GET`, `HEAD`, `OPTIONS`, `TRACE`, `PUT` and `DELETE` ([RFC 9110 ยง9.2.2](https://www.rfc-editor.org/rfc/rfc9110#section-9.2.2)). A `POST` or `PATCH` that timed out or got a 429/5xx may already have done its work on the server โ€” re-sending it bills a second image, sends a second message, charges a card twice. So by default it makes **exactly one attempt** and you get its own error (`HTTPResponseError` / `TimeoutError`), not a `RetryError`. + +Opt a request back in when a duplicate is harmless, in one of two ways: + +```typescript +// 1. Best: send an Idempotency-Key the server deduplicates on. Any request carrying a +// non-empty one is retry-eligible, whatever its method. +await fetch('https://api.example.com/charges', { + method: 'POST', + headers: { 'Idempotency-Key': crypto.randomUUID() }, + body: JSON.stringify(charge), +}); + +// 2. Explicitly: the endpoint tolerates duplicates. +await fetch('https://api.example.com/search', { + method: 'POST', + body: JSON.stringify(query), + options: { retry: { allowNonIdempotent: true } }, +}); +``` + +The same rule holds in every port (`allow_non_idempotent` in Python and Rust, `AllowNonIdempotent` in Go and .NET) and is pinned by the shared [`spec/retry-idempotency-corpus.json`](spec/retry-idempotency-corpus.json), which each suite runs against a real local server and counts what the **server** received. + +- **429 + `Retry-After` on a POST is not retried either** unless you opt in. `Retry-After` says when the server will take a request again, not that the first one did nothing. +- **`onRejection` cannot override this.** An ineligible request never consults it; opting in is the only switch. +- **The client-side rate limiter is unaffected**: it rejects before anything is sent, so its own retry loop still applies to every method. +- **Timeouts cancel the attempt.** When the per-attempt timeout fires, the request is aborted and its connection closed โ€” before any retry โ€” so the server can see the cancel. (Before 4.0.0 the TypeScript timeout only raced the request, leaving it running server-side while the retry sent it again.) + ### ๐Ÿšฆ Rate-limit respect ```typescript @@ -203,10 +234,10 @@ flowchart LR | TypeScript | [`@smooai/fetch`](https://www.npmjs.com/package/@smooai/fetch) | `pnpm add @smooai/fetch` | | Python | [`smooai-fetch`](https://pypi.org/project/smooai-fetch/) | `pip install smooai-fetch` | | Rust | [`smooai-fetch`](https://crates.io/crates/smooai-fetch) | `cargo add smooai-fetch` | -| Go | `github.com/SmooAI/fetch/go/fetch/v3` | `go get github.com/SmooAI/fetch/go/fetch/v3` | +| Go | `github.com/SmooAI/fetch/go/fetch/v4` | `go get github.com/SmooAI/fetch/go/fetch/v4` | | .NET | [`SmooAI.Fetch`](https://www.nuget.org/packages/SmooAI.Fetch) | `dotnet add package SmooAI.Fetch` | -> **Go note:** the module path carries the `/v3` major suffix Go requires above v1, so the `go/fetch/v3.x` tags resolve. The import path is `github.com/SmooAI/fetch/go/fetch/v3`; the package identifier is still `fetch`. Tags minted before this change (through `go/fetch/v3.4.0`) point at commits whose `go.mod` lacked the suffix and will not resolve โ€” use `v3.4.1` or later. +> **Go note:** the module path carries the `/v4` major suffix Go requires above v1 (it was `/v3` before 4.0.0), so the `go/fetch/v4.x` tags resolve. The import path is `github.com/SmooAI/fetch/go/fetch/v4`; the package identifier is still `fetch`. Tags minted before this change (through `go/fetch/v3.4.0`) point at commits whose `go.mod` lacked the suffix and will not resolve โ€” use `v3.4.1` or later. Language-specific source lives in [`src/`](./src/) (TypeScript), [`python/`](./python/), [`rust/`](./rust/), [`go/`](./go/), and [`dotnet/`](./dotnet/). @@ -312,9 +343,9 @@ Where a port leans on a battle-tested ecosystem library (mollitia, Polly), it sa Out of the box, `@smooai/fetch` is configured for the real world: -**Retry strategy** โ€” 2 automatic retries, exponential backoff (500ms โ†’ 1s โ†’ 2s), jitter to prevent thundering herds, and retries only on network errors or retryable responses. +**Retry strategy** โ€” 2 automatic retries, exponential backoff (500ms โ†’ 1s โ†’ 2s), jitter to prevent thundering herds, and retries only on network errors or retryable responses โ€” for idempotent methods only, unless a request opts in ([details](#-retries-never-duplicate-a-side-effect)). -**Timeout protection** โ€” 10-second default timeout, configurable per request, so requests never hang indefinitely. +**Timeout protection** โ€” 10-second default per-attempt timeout, configurable per request, so requests never hang indefinitely. When it fires, the attempt is aborted and its connection closed. **Connect timeout (opt-in)** โ€” `connectTimeoutMs` / `withConnectTimeout` bounds only the connection-establishment phase, in all five ports. A black-holed connect then fails in ~that window and retry lands on a live endpoint, instead of burning the whole-request timeout on a dead one; slow-but-alive handlers are unaffected. Off by default. In TypeScript it needs the optional peer dependency `undici` and applies to Node only. diff --git a/dotnet/SmooAI.Fetch.Tests/RetryIdempotencyTests.cs b/dotnet/SmooAI.Fetch.Tests/RetryIdempotencyTests.cs new file mode 100644 index 0000000..fff5038 --- /dev/null +++ b/dotnet/SmooAI.Fetch.Tests/RetryIdempotencyTests.cs @@ -0,0 +1,432 @@ +using System.Collections.Concurrent; +using System.Diagnostics; +using System.Net; +using System.Net.Sockets; +using System.Text; +using System.Text.Json; +using SmooAI.Fetch; + +namespace SmooAI.Fetch.Tests; + +/// +/// SMOODEV-3375 โ€” retries were method-blind, so a POST that timed out or got a 429/5xx was +/// re-sent up to twice more and its side effect (a charge, a message, a generated image) +/// could execute three times server-side. +/// +/// Every case comes from spec/retry-idempotency-corpus.json, shared with the other four +/// ports (copied next to the test assembly by the csproj). Each one runs against a REAL +/// local server and asserts how many requests the SERVER received: the bug is a side +/// effect repeating server-side, and only the server can count that. +/// +public class RetryIdempotencyTests +{ + private static readonly Lazy CorpusRoot = new(() => + { + using var doc = JsonDocument.Parse(File.ReadAllText("retry-idempotency-corpus.json")); + return doc.RootElement.Clone(); + }); + + private static JsonElement Corpus => CorpusRoot.Value; + + public static IEnumerable Cases() => + Corpus.GetProperty("cases").EnumerateArray().Select(c => new object[] { c.GetProperty("name").GetString()! }); + + // Positive control: a corpus that failed to load, or loaded empty, would make every + // Theory below vacuous โ€” xunit reports zero cases as a pass. + [Fact] + public void Loaded_the_shared_corpus() + { + Assert.True(Corpus.GetProperty("cases").GetArrayLength() > 0, "retry corpus has no cases"); + Assert.True(Corpus.GetProperty("methods").GetProperty("idempotent").GetArrayLength() > 0); + Assert.True(Corpus.GetProperty("methods").GetProperty("nonIdempotent").GetArrayLength() > 0); + Assert.Equal(RetryPolicy.IdempotencyKeyHeader, Corpus.GetProperty("idempotencyKeyHeader").GetString()); + } + + [Fact] + public void Classifies_every_corpus_method() + { + foreach (var m in Corpus.GetProperty("methods").GetProperty("idempotent").EnumerateArray()) + { + Assert.True(RetryPolicy.IsIdempotentMethod(m.GetString()!), $"{m} should be idempotent"); + Assert.True(RetryPolicy.IsIdempotentMethod(new HttpMethod(m.GetString()!)), $"{m} should be idempotent"); + } + + foreach (var m in Corpus.GetProperty("methods").GetProperty("nonIdempotent").EnumerateArray()) + { + Assert.False(RetryPolicy.IsIdempotentMethod(m.GetString()!), $"{m} should NOT be idempotent"); + Assert.False(RetryPolicy.IsIdempotentMethod(new HttpMethod(m.GetString()!)), $"{m} should NOT be idempotent"); + } + } + + [Theory] + [MemberData(nameof(Cases))] + public async Task Server_sees_the_expected_number_of_attempts(string name) + { + var c = Corpus.GetProperty("cases").EnumerateArray().Single(x => x.GetProperty("name").GetString() == name); + var respond = c.GetProperty("respond"); + var hang = respond.TryGetProperty("hang", out var h) && h.GetBoolean(); + var status = hang ? 0 : respond.GetProperty("status").GetInt32(); + var responseHeaders = new Dictionary(); + if (!hang && respond.TryGetProperty("headers", out var rh)) + { + foreach (var p in rh.EnumerateObject()) + { + responseHeaders[p.Name] = p.Value.GetString()!; + } + } + + await using var server = new CountingServer(hang, status, responseHeaders); + + var retry = Corpus.GetProperty("retry"); + var fetch = SmooFetchBuilder.Create() + .WithRetry(new RetryPolicy + { + MaxRetries = retry.GetProperty("attempts").GetInt32(), + BaseDelay = TimeSpan.FromMilliseconds(retry.GetProperty("initialIntervalMs").GetInt32()), + BackoffFactor = retry.GetProperty("factor").GetDouble(), + UseJitter = retry.GetProperty("jitterAdjustment").GetDouble() > 0, + JitterFraction = retry.GetProperty("jitterAdjustment").GetDouble(), + AllowNonIdempotent = c.TryGetProperty("allowNonIdempotent", out var a) && a.GetBoolean(), + }) + .WithTimeout(TimeSpan.FromMilliseconds(Corpus.GetProperty("timeoutMs").GetInt32())) + .Build(); + + using var request = new HttpRequestMessage(new HttpMethod(c.GetProperty("method").GetString()!), server.Url); + if (c.TryGetProperty("body", out var body)) + { + request.Content = new StringContent(body.GetString()!, Encoding.UTF8, "application/json"); + } + + if (c.TryGetProperty("headers", out var headers)) + { + foreach (var p in headers.EnumerateObject()) + { + request.Headers.TryAddWithoutValidation(p.Name, p.Value.GetString()); + } + } + + try + { + using var response = await fetch.SendAsync(request); + Assert.Equal(status, (int)response.StatusCode); + } + catch (Exception ex) when (hang && ex is OperationCanceledException) + { + // Expected: every attempt timed out. + } + + Assert.Equal(c.GetProperty("expectedAttempts").GetInt32(), server.Requests); + } + + [Fact] + public async Task Timed_out_retried_GET_closes_first_connection_before_the_retry_arrives() + { + var (timeoutMs, slackMs) = TimeoutAbortKnobs(); + await using var server = new HangingSocketServer(); + + var fetch = SmooFetchBuilder.Create() + .WithRetry(new RetryPolicy { MaxRetries = 1, BaseDelay = TimeSpan.FromMilliseconds(1), UseJitter = false }) + .WithTimeout(TimeSpan.FromMilliseconds(timeoutMs)) + .Build(); + + await Assert.ThrowsAnyAsync(() => fetch.GetAsync(server.Url)); + + var first = await server.WaitForConnection(0, TimeSpan.FromSeconds(10)); + var second = await server.WaitForConnection(1, TimeSpan.FromSeconds(10)); + var firstClosed = await first.WaitClosed(TimeSpan.FromMilliseconds(timeoutMs + slackMs)); + Assert.NotNull(firstClosed); + Assert.True( + firstClosed!.Value - first.ArrivedAt <= TimeSpan.FromMilliseconds(timeoutMs + slackMs), + $"first attempt's connection closed {(firstClosed.Value - first.ArrivedAt).TotalMilliseconds}ms after it arrived; expected <= {timeoutMs + slackMs}ms"); + + // The FIN goes out before the retry's SYN; the two server-side observations race + // on the thread pool, so allow a small scheduling tolerance. An attempt that was + // abandoned rather than cancelled stays open for the whole test instead. + Assert.True( + firstClosed.Value <= second.ArrivedAt + TimeSpan.FromMilliseconds(100), + $"first connection closed at {firstClosed.Value.TotalMilliseconds}ms, AFTER the retry arrived at {second.ArrivedAt.TotalMilliseconds}ms โ€” the timed-out attempt was not cancelled"); + } + + [Fact] + public async Task Timed_out_POST_closes_its_connection_and_is_not_retried() + { + var (timeoutMs, slackMs) = TimeoutAbortKnobs(); + await using var server = new HangingSocketServer(); + + var fetch = SmooFetchBuilder.Create() + .WithRetry(new RetryPolicy { MaxRetries = 2, BaseDelay = TimeSpan.FromMilliseconds(1), UseJitter = false }) + .WithTimeout(TimeSpan.FromMilliseconds(timeoutMs)) + .Build(); + + await Assert.ThrowsAnyAsync(() => fetch.PostAsync(server.Url, new { a = 1 })); + + var first = await server.WaitForConnection(0, TimeSpan.FromSeconds(10)); + var closed = await first.WaitClosed(TimeSpan.FromMilliseconds(timeoutMs + slackMs)); + Assert.NotNull(closed); + Assert.True(closed!.Value - first.ArrivedAt <= TimeSpan.FromMilliseconds(timeoutMs + slackMs)); + + // Give a (wrongly) retried attempt time to show up before asserting it didn't. + await Task.Delay(200); + Assert.Equal(1, server.ConnectionCount); + } + + private static (int timeoutMs, int slackMs) TimeoutAbortKnobs() + { + var t = Corpus.GetProperty("timeoutAbort"); + return (t.GetProperty("timeoutMs").GetInt32(), t.GetProperty("closeSlackMs").GetInt32()); + } + + private static async Task ReadRequestHeadAsync(NetworkStream stream, CancellationToken ct) + { + // Read the head byte-by-byte up to CRLFCRLF, then drain Content-Length bytes. + var head = new StringBuilder(); + var buf = new byte[1]; + while (!head.ToString().EndsWith("\r\n\r\n", StringComparison.Ordinal)) + { + if (await stream.ReadAsync(buf, ct) == 0) + { + return false; + } + + head.Append((char)buf[0]); + } + + var length = 0; + foreach (var line in head.ToString().Split("\r\n")) + { + if (line.StartsWith("Content-Length:", StringComparison.OrdinalIgnoreCase)) + { + length = int.Parse(line["Content-Length:".Length..].Trim()); + } + } + + var body = new byte[length]; + var read = 0; + while (read < length) + { + var n = await stream.ReadAsync(body.AsMemory(read), ct); + if (n == 0) + { + return false; + } + + read += n; + } + + return true; + } + + /// + /// Minimal HTTP/1.1 responder on 127.0.0.1 that counts every request it reads. Replies + /// with a fixed status and Connection: close, or never replies when hanging. + /// + private sealed class CountingServer : IAsyncDisposable + { + private readonly TcpListener _listener = new(IPAddress.Loopback, 0); + private readonly CancellationTokenSource _cts = new(); + private readonly ConcurrentBag _clients = new(); + private int _requests; + + public CountingServer(bool hang, int status, IReadOnlyDictionary headers) + { + _listener.Start(); + Url = $"http://127.0.0.1:{((IPEndPoint)_listener.LocalEndpoint).Port}/probe"; + _ = Task.Run(() => AcceptLoop(hang, status, headers)); + } + + public string Url { get; } + + public int Requests => Volatile.Read(ref _requests); + + private async Task AcceptLoop(bool hang, int status, IReadOnlyDictionary headers) + { + while (!_cts.IsCancellationRequested) + { + TcpClient client; + try + { + client = await _listener.AcceptTcpClientAsync(_cts.Token); + } + catch + { + return; + } + + _clients.Add(client); + _ = Task.Run(async () => + { + try + { + var stream = client.GetStream(); + if (!await ReadRequestHeadAsync(stream, _cts.Token)) + { + return; + } + + Interlocked.Increment(ref _requests); + if (hang) + { + return; // Never answer; the socket is closed at teardown. + } + + var sb = new StringBuilder($"HTTP/1.1 {status} X\r\nContent-Length: 0\r\nConnection: close\r\n"); + foreach (var kvp in headers) + { + sb.Append($"{kvp.Key}: {kvp.Value}\r\n"); + } + + sb.Append("\r\n"); + await stream.WriteAsync(Encoding.ASCII.GetBytes(sb.ToString()), _cts.Token); + client.Close(); + } + catch + { + // Client went away; nothing to count. + } + }); + } + } + + public ValueTask DisposeAsync() + { + _cts.Cancel(); + _listener.Stop(); + foreach (var c in _clients) + { + c.Dispose(); + } + + _cts.Dispose(); + return ValueTask.CompletedTask; + } + } + + /// + /// Raw socket server that reads each request and never answers, recording when each + /// connection arrived and when the CLIENT closed it (Read returns 0 or the socket resets). + /// + private sealed class HangingSocketServer : IAsyncDisposable + { + private readonly TcpListener _listener = new(IPAddress.Loopback, 0); + private readonly CancellationTokenSource _cts = new(); + private readonly Stopwatch _clock = Stopwatch.StartNew(); + private readonly List _connections = new(); + private readonly SemaphoreSlim _arrived = new(0); + + public HangingSocketServer() + { + _listener.Start(); + Url = $"http://127.0.0.1:{((IPEndPoint)_listener.LocalEndpoint).Port}/hang"; + _ = Task.Run(AcceptLoop); + } + + public string Url { get; } + + public int ConnectionCount + { + get + { + lock (_connections) + { + return _connections.Count; + } + } + } + + public async Task WaitForConnection(int index, TimeSpan timeout) + { + var deadline = _clock.Elapsed + timeout; + while (_clock.Elapsed < deadline) + { + lock (_connections) + { + if (_connections.Count > index) + { + return _connections[index]; + } + } + + await _arrived.WaitAsync(TimeSpan.FromMilliseconds(50)); + } + + throw new TimeoutException($"connection #{index + 1} never arrived"); + } + + private async Task AcceptLoop() + { + while (!_cts.IsCancellationRequested) + { + TcpClient client; + try + { + client = await _listener.AcceptTcpClientAsync(_cts.Token); + } + catch + { + return; + } + + var conn = new Conn(client, _clock.Elapsed); + lock (_connections) + { + _connections.Add(conn); + } + + _arrived.Release(); + _ = Task.Run(async () => + { + var buf = new byte[4096]; + try + { + var stream = client.GetStream(); + while (await stream.ReadAsync(buf, _cts.Token) > 0) + { + // Swallow the request; never answer. + } + } + catch (OperationCanceledException) when (_cts.IsCancellationRequested) + { + return; // Teardown, not a client close. + } + catch + { + // A reset counts as closed too. + } + + conn.Closed.TrySetResult(_clock.Elapsed); + }); + } + } + + public ValueTask DisposeAsync() + { + _cts.Cancel(); + _listener.Stop(); + lock (_connections) + { + foreach (var c in _connections) + { + c.Client.Dispose(); + } + } + + _cts.Dispose(); + return ValueTask.CompletedTask; + } + + public sealed class Conn(TcpClient client, TimeSpan arrivedAt) + { + public TcpClient Client { get; } = client; + + public TimeSpan ArrivedAt { get; } = arrivedAt; + + public TaskCompletionSource Closed { get; } = new(TaskCreationOptions.RunContinuationsAsynchronously); + + public async Task WaitClosed(TimeSpan timeout) + { + var done = await Task.WhenAny(Closed.Task, Task.Delay(timeout)); + return done == Closed.Task ? Closed.Task.Result : null; + } + } + } +} diff --git a/dotnet/SmooAI.Fetch.Tests/SmooAI.Fetch.Tests.csproj b/dotnet/SmooAI.Fetch.Tests/SmooAI.Fetch.Tests.csproj index ba83b35..b4e561c 100644 --- a/dotnet/SmooAI.Fetch.Tests/SmooAI.Fetch.Tests.csproj +++ b/dotnet/SmooAI.Fetch.Tests/SmooAI.Fetch.Tests.csproj @@ -27,6 +27,7 @@ port's suite. Copy rather than hand-copying the values into the test. --> + diff --git a/dotnet/SmooAI.Fetch/README.md b/dotnet/SmooAI.Fetch/README.md index d9f60d4..3b83bc2 100644 --- a/dotnet/SmooAI.Fetch/README.md +++ b/dotnet/SmooAI.Fetch/README.md @@ -16,7 +16,7 @@ dotnet add package SmooAI.Fetch ## What you get - **Typed JSON** โ€” `GetAsync` and `PostAsync`. Your request and response shapes, strongly typed. No `JsonSerializer.Deserialize` boilerplate on every call site. -- **Automatic retries on transient failures** โ€” network blips, timeouts, and `408` / `425` / `429` / `5xx` responses are retried with exponential backoff + jitter. +- **Automatic retries on transient failures, for requests that are safe to replay** โ€” network blips, timeouts, and `408` / `429` / `5xx` responses are retried with exponential backoff + jitter for idempotent methods (`GET`, `HEAD`, `OPTIONS`, `TRACE`, `PUT`, `DELETE`). `POST` and `PATCH` are never retried unless you opt in โ€” see [Retries are method-aware](#retries-are-method-aware). - **`Retry-After` is honored** โ€” when a server tells you "wait 5s", the client waits 5s instead of your default backoff. Never eat a 429 again. - **Async auth tokens** โ€” register an `AuthTokenProvider` once; every request picks up a fresh bearer token without restarting `HttpClient`. - **One typed error per non-2xx** โ€” catch `HttpResponseError` and you've got status, body, headers, URI, and method on one exception. @@ -62,17 +62,37 @@ public class BillingService(SmooFetch fetch) ```csharp options.RetryPolicy = RetryPolicy.ExponentialBackoff( - maxRetries: 3, - initialDelay: TimeSpan.FromMilliseconds(250), - maxDelay: TimeSpan.FromSeconds(10), - backoffFactor: 2.0, - jitter: true); + maxRetries: 3, + baseDelay: TimeSpan.FromMilliseconds(250), + maxDelay: TimeSpan.FromSeconds(10)); ``` -- Retries on transient exceptions (timeouts, socket errors) and on `408` / `425` / `429` / `500` / `502` / `503` / `504`. +- Retries on transient exceptions (timeouts, socket errors) and on `408` / `429` / `500` / `502` / `503` / `504` โ€” for retry-eligible requests only (below). - Honors the `Retry-After` header on `429` / `503` โ€” if the server says "wait 5s", the client waits 5s instead of your backoff. - Exponential backoff with jitter; bounded by `maxDelay`. +## Retries are method-aware + +A `POST` that timed out, or came back `429` / `5xx`, may already have run on the server. Sending it again can charge a card twice, send a message twice, or bill a second generated image. So a failed attempt is retried only when the request is **retry-eligible**: + +- its method is idempotent per RFC 9110 ยง9.2.2 โ€” `GET`, `HEAD`, `OPTIONS`, `TRACE`, `PUT`, `DELETE`; **or** +- the policy sets `AllowNonIdempotent = true`; **or** +- the request carries a non-empty `Idempotency-Key` header (`RetryPolicy.IdempotencyKeyHeader`), i.e. the server has promised to de-duplicate replays. + +An ineligible request makes exactly **one** attempt. `OnRejection` is not consulted, and the outcome surfaces as-is: the non-2xx response or the original exception. That includes a `429` with `Retry-After` on a `POST`: it is still not replayed unless you opt in. Eligibility is decided on the final request, after the `AuthTokenProvider` and `PreRequest` hook run, so a hook can add the `Idempotency-Key`. The in-process rate limiter is unaffected, because it waits before anything is sent. + +```csharp +// This endpoint is safe to replay: opt the whole client in. +options.RetryPolicy = RetryPolicy.Default with { AllowNonIdempotent = true }; + +// Or opt in one request, with a key the server de-duplicates on. +using var request = new HttpRequestMessage(HttpMethod.Post, "/payments") { Content = body }; +request.Headers.Add(RetryPolicy.IdempotencyKeyHeader, paymentId); +using var response = await fetch.SendAsync(request); +``` + +> **Changed in 4.0.0.** Through 3.x every method was retried, `POST` included. + ## Typed errors โ€” one `catch` per layer ```csharp @@ -106,7 +126,7 @@ options.AuthTokenProvider = async ct => ## Cancellation + per-request timeout -Every method accepts a `CancellationToken`. The configured `Timeout` creates a linked `CancellationTokenSource` under the hood: +Every method accepts a `CancellationToken`. The configured `Timeout` applies **per attempt**: each attempt runs under its own `CancellationTokenSource`, linked to yours. When the timeout fires, `HttpClient` aborts the in-flight request and closes its connection before any retry starts, so the server sees the cancel. A timed-out attempt is never left running alongside its retry. ```csharp using var cts = new CancellationTokenSource(TimeSpan.FromSeconds(5)); diff --git a/dotnet/SmooAI.Fetch/RetryPolicy.cs b/dotnet/SmooAI.Fetch/RetryPolicy.cs index 4f8cf33..dd1fcb5 100644 --- a/dotnet/SmooAI.Fetch/RetryPolicy.cs +++ b/dotnet/SmooAI.Fetch/RetryPolicy.cs @@ -136,6 +136,62 @@ public sealed record RetryPolicy /// public OnRejectionCallback? OnRejection { get; init; } + /// + /// Header that, when present with a non-empty value, makes a non-idempotent request + /// (POST, PATCH) retry-eligible: the caller is telling the server to de-duplicate + /// replays, so re-sending it cannot double the side effect. + /// + public const string IdempotencyKeyHeader = "Idempotency-Key"; + + /// + /// Allow retrying non-idempotent requests (POST, PATCH, anything not idempotent per + /// RFC 9110 ยง9.2.2). Defaults to false: a POST that times out or gets a 429/5xx + /// may already have executed server-side, so re-sending it can duplicate the side + /// effect (a second charge, a second message, a second generated image). Set this only + /// when the endpoint is safe to replay; alternatively send an + /// header, which opts that one request in. + /// A 429 with Retry-After on a POST is not retried unless opted in either. + /// + public bool AllowNonIdempotent { get; init; } + + private static readonly HashSet IdempotentMethods = new(StringComparer.OrdinalIgnoreCase) + { + "GET", "HEAD", "OPTIONS", "TRACE", "PUT", "DELETE", + }; + + /// + /// Whether is idempotent per RFC 9110 ยง9.2.2 (GET, HEAD, + /// OPTIONS, TRACE, PUT, DELETE), compared case-insensitively. Anything else โ€” POST, + /// PATCH, CONNECT, or an unrecognised method โ€” is treated as non-idempotent. + /// + public static bool IsIdempotentMethod(string method) => + !string.IsNullOrEmpty(method) && IdempotentMethods.Contains(method); + + /// + public static bool IsIdempotentMethod(HttpMethod method) + { + ArgumentNullException.ThrowIfNull(method); + return IsIdempotentMethod(method.Method); + } + + /// + /// Whether a failed attempt of may be retried: its method is + /// idempotent, OR is set, OR it carries a non-empty + /// . Evaluated on the final request, after the auth + /// provider and PreRequest hook have run (either can add the header). + /// + public bool IsRetryEligible(HttpRequestMessage request) + { + ArgumentNullException.ThrowIfNull(request); + if (IsIdempotentMethod(request.Method) || AllowNonIdempotent) + { + return true; + } + + return request.Headers.TryGetValues(IdempotencyKeyHeader, out var values) + && values.Any(v => !string.IsNullOrWhiteSpace(v)); + } + private static readonly IReadOnlyCollection DefaultRetryStatusCodes = new[] { HttpStatusCode.RequestTimeout, @@ -149,7 +205,10 @@ public sealed record RetryPolicy /// A retry policy that never retries. public static RetryPolicy None { get; } = new() { MaxRetries = 0 }; - /// Default policy: 2 retries, 500 ms base, exponential factor 2, jitter ยฑ50%, honors Retry-After. + /// + /// Default policy: 2 retries, 500 ms base, exponential factor 2, jitter ยฑ50%, honors + /// Retry-After โ€” for idempotent requests only (see ). + /// public static RetryPolicy Default { get; } = new() { MaxRetries = 2, diff --git a/dotnet/SmooAI.Fetch/SmooAI.Fetch.csproj b/dotnet/SmooAI.Fetch/SmooAI.Fetch.csproj index 6b7674e..b8c8034 100644 --- a/dotnet/SmooAI.Fetch/SmooAI.Fetch.csproj +++ b/dotnet/SmooAI.Fetch/SmooAI.Fetch.csproj @@ -35,7 +35,7 @@ - + diff --git a/dotnet/SmooAI.Fetch/SmooFetch.cs b/dotnet/SmooAI.Fetch/SmooFetch.cs index 41639c1..70ba967 100644 --- a/dotnet/SmooAI.Fetch/SmooFetch.cs +++ b/dotnet/SmooAI.Fetch/SmooFetch.cs @@ -24,6 +24,7 @@ public sealed class SmooFetch private readonly SmooFetchOptions _options; private readonly ILogger _logger; private readonly ResiliencePipeline _pipeline; + private readonly ResiliencePipeline _noRetryPipeline; private readonly SlidingWindowLimiterAdapter? _rateLimiter; /// Constant used when registering the typed client with . @@ -35,16 +36,37 @@ internal SmooFetch(HttpClient httpClient, bool ownsHttpClient, SmooFetchOptions _ownsHttpClient = ownsHttpClient; _options = options ?? throw new ArgumentNullException(nameof(options)); _logger = logger ?? NullLogger.Instance; - _pipeline = BuildPipeline(options); + (_pipeline, _noRetryPipeline) = BuildPipelines(options); _rateLimiter = options.RateLimiter is { } rl ? new SlidingWindowLimiterAdapter(rl) : null; } - private static ResiliencePipeline BuildPipeline(SmooFetchOptions options) + /// + /// Builds the pipeline for retry-eligible requests and the one for everything else. + /// Both share ONE circuit breaker instance, so a non-idempotent request still counts + /// towards (and is refused by) the same breaker โ€” it just never gets a second attempt. + /// + private static (ResiliencePipeline Retrying, ResiliencePipeline NoRetry) BuildPipelines(SmooFetchOptions options) { var retryPipeline = options.RetryPolicy.BuildPipeline(); + var breaker = BuildBreaker(options); + if (breaker is null) + { + return (retryPipeline, ResiliencePipeline.Empty); + } + + // Wrap the retry pipeline inside the breaker pipeline. + var retrying = new ResiliencePipelineBuilder() + .AddPipeline(breaker) + .AddPipeline(retryPipeline) + .Build(); + return (retrying, breaker); + } + + private static ResiliencePipeline? BuildBreaker(SmooFetchOptions options) + { if (options.CircuitBreaker is not { } cb) { - return retryPipeline; + return null; } // Compose: caller -> circuit breaker -> retry -> http. @@ -53,7 +75,7 @@ private static ResiliencePipeline BuildPipeline(SmooFetchOp // mirror the retry policy's failure predicate so the breaker trips on the // same conditions (retryable statuses + transient exceptions). var policy = options.RetryPolicy; - var breaker = new ResiliencePipelineBuilder() + return new ResiliencePipelineBuilder() .AddCircuitBreaker(new CircuitBreakerStrategyOptions { FailureRatio = 1.0, @@ -72,12 +94,6 @@ private static ResiliencePipeline BuildPipeline(SmooFetchOp }, }) .Build(); - - // Wrap the retry pipeline inside the breaker pipeline. - return new ResiliencePipelineBuilder() - .AddPipeline(breaker) - .AddPipeline(retryPipeline) - .Build(); } private static bool IsTransient(Exception ex) => @@ -178,10 +194,18 @@ public async Task SendAsync(HttpRequestMessage request, Can await pre(request, cancellationToken).ConfigureAwait(false); } + // Retries are method-aware (RFC 9110 ยง9.2.2). A POST/PATCH that timed out or got a + // 429/5xx may already have executed server-side, so replaying it can duplicate its + // side effect. Decided on the FINAL request โ€” after the auth provider and PreRequest + // hook, either of which can add an Idempotency-Key. An ineligible request makes + // exactly one attempt: OnRejection is not consulted and the underlying outcome + // surfaces as-is. + var pipeline = _options.RetryPolicy.IsRetryEligible(request) ? _pipeline : _noRetryPipeline; + HttpResponseMessage response; try { - response = await _pipeline.ExecuteAsync(async ct => + response = await pipeline.ExecuteAsync(async ct => { // Rate-limit gate: acquire a permit before dispatch. The limiter // waits / queues until a permit is available, so retries are not @@ -192,6 +216,10 @@ public async Task SendAsync(HttpRequestMessage request, Can await _rateLimiter.AcquireAsync(request, ct).ConfigureAwait(false); } + // Per-ATTEMPT timeout, linked to the caller's token. When it fires, HttpClient + // aborts the in-flight request and closes its connection, so the server sees + // the cancel before any retry starts โ€” a timed-out attempt is cancelled, not + // abandoned while a second copy runs alongside it. using var linkedCts = CancellationTokenSource.CreateLinkedTokenSource(ct); linkedCts.CancelAfter(_options.Timeout); var attemptRequest = await CloneRequestAsync(request).ConfigureAwait(false); diff --git a/go/fetch/README.md b/go/fetch/README.md index a33d1ad..80d9d0b 100644 --- a/go/fetch/README.md +++ b/go/fetch/README.md @@ -58,7 +58,7 @@ Ever had a Go microservice pile up goroutines because a downstream API was down? ### Install ```bash -go get github.com/SmooAI/fetch/go/fetch/v3 +go get github.com/SmooAI/fetch/go/fetch/v4 ``` | Language | Package | Install | @@ -66,7 +66,7 @@ go get github.com/SmooAI/fetch/go/fetch/v3 | TypeScript | [`@smooai/fetch`](https://www.npmjs.com/package/@smooai/fetch) | `pnpm add @smooai/fetch` | | Python | [`smooai-fetch`](https://pypi.org/project/smooai-fetch/) | `pip install smooai-fetch` | | Rust | [`smooai-fetch`](https://crates.io/crates/smooai-fetch) | `cargo add smooai-fetch` | -| Go | `github.com/SmooAI/fetch/go/fetch/v3` | `go get github.com/SmooAI/fetch/go/fetch/v3` | +| Go | `github.com/SmooAI/fetch/go/fetch/v4` | `go get github.com/SmooAI/fetch/go/fetch/v4` | ## The Power of Resilient Fetching @@ -75,7 +75,7 @@ go get github.com/SmooAI/fetch/go/fetch/v3 Watch how smooai-fetch handles common failure scenarios: ```go -import "github.com/SmooAI/fetch/go/fetch/v3" +import "github.com/SmooAI/fetch/go/fetch/v4" type ApiData struct { ID string `json:"id"` @@ -121,7 +121,7 @@ resp, err := fetch.Get[Repos](ctx, nil, "https://api.github.com/user/repos", nil ```go import ( - "github.com/SmooAI/fetch/go/fetch/v3" + "github.com/SmooAI/fetch/go/fetch/v4" "time" ) @@ -174,7 +174,7 @@ fmt.Println("Created user:", resp.Data.ID) ```go import ( "errors" - "github.com/SmooAI/fetch/go/fetch/v3" + "github.com/SmooAI/fetch/go/fetch/v4" "time" ) @@ -314,14 +314,52 @@ Out of the box, smooai-fetch is configured for the real world: - 2 automatic retries on failure - Exponential backoff: 500ms -> 1s -> 2s - Jitter to prevent thundering herds -- Only retries on network errors or 5xx responses +- Only retries on network errors, timeouts, 429 or 5xx responses +- **Only idempotent methods are retried** (see below) **Timeout Protection:** -- 10-second default timeout +- 10-second default timeout, per attempt - Prevents indefinite hangs on slow endpoints +- A timed-out attempt is **cancelled**, not abandoned: its context is cancelled, net/http closes the connection so the server sees the request go away, and the next attempt never overlaps it - Configurable per client or per request +### Retries are method-aware (breaking in v4) + +Re-sending a request that timed out or failed with a 429/5xx is only safe when +the request is idempotent โ€” the server may already have acted on the first +attempt. Before v4 a slow `POST` could execute up to three times server-side +(three image generations billed, a message sent twice, a payment captured twice). + +A failed attempt is now retried only when the request is **retry-eligible**: + +- its method is idempotent per RFC 9110 ยง9.2.2 โ€” `GET`, `HEAD`, `OPTIONS`, `TRACE`, `PUT`, `DELETE` (see `fetch.IsIdempotentMethod`), **or** +- the retry options set `AllowNonIdempotent: true`, **or** +- the request carries a non-empty `Idempotency-Key` header (`fetch.IdempotencyKeyHeader`), the server-side contract that makes a replay safe. + +`POST` and `PATCH` get exactly one attempt otherwise, and the error is the +underlying `*HTTPResponseError` / `*TimeoutError` โ€” not a `*RetryError`. That +includes a `429` with `Retry-After`: the server asked the client to come back +later, but without an opt-in the client cannot know the first attempt had no +effect. A custom `OnRejection` callback is not an opt-in; it is not consulted +for an ineligible request. Eligibility is decided on the final request, after +`PreRequest` hooks and the auth-token provider run, so a hook can add the key. + +Rejections raised before anything is sent โ€” the in-process rate limiter and an +open circuit breaker โ€” are still retried for every method. + +```go +// Opt a POST into retries because the endpoint deduplicates on its own. +retry := fetch.DefaultRetryOptions +retry.AllowNonIdempotent = true +client := fetch.NewClientBuilder().WithRetry(&retry).Build() + +// Or send an Idempotency-Key and keep the default options. +resp, err := fetch.Post[Result](ctx, nil, url, payload, &fetch.RequestOptions{ + Headers: http.Header{fetch.IdempotencyKeyHeader: {requestID}}, +}) +``` + **Rate Limit Handling:** - Respects Retry-After headers from 429 responses @@ -364,7 +402,7 @@ Out of the box, smooai-fetch is configured for the real world: ```go import ( "errors" - "github.com/SmooAI/fetch/go/fetch/v3" + "github.com/SmooAI/fetch/go/fetch/v4" ) resp, err := fetch.Get[Data](ctx, client, "https://api.example.com/data", nil) diff --git a/go/fetch/client.go b/go/fetch/client.go index f399bdf..06fc36c 100644 --- a/go/fetch/client.go +++ b/go/fetch/client.go @@ -96,9 +96,28 @@ func Fetch[T any](ctx context.Context, client *Client, method, url string, body hooks = client.hooks } + // Retry-eligibility, seeded from what is known before any hook runs and + // refined by every attempt once hooks and the auth provider have shaped the + // final request (see retryGate). + gate := &retryGate{opts: retryOpts} + { + preHook := make(http.Header) + for key, vals := range client.baseHeaders { + for _, v := range vals { + preHook.Set(key, v) + } + } + for key, vals := range extraHeaders { + for _, v := range vals { + preHook.Set(key, v) + } + } + gate.eligible.Store(isRetryEligible(method, preHook, retryOpts)) + } + // Build the core request-execute function doRequest := func(ctx context.Context) (*FetchResponse[T], error) { - return executeHTTPRequest[T](ctx, client, method, url, body, extraHeaders, hooks, validate) + return executeHTTPRequest[T](ctx, client, method, url, body, extraHeaders, hooks, validate, gate) } // Wrap with timeout if configured @@ -146,10 +165,18 @@ func Fetch[T any](ctx context.Context, client *Client, method, url string, body } } - // Wrap with retry if configured + // Wrap with retry if configured. A request that is not retry-eligible + // (non-idempotent, no opt-in, no Idempotency-Key) makes exactly one attempt + // and surfaces the underlying error: re-sending a POST that timed out or got + // a 429/5xx can execute its side effect twice. Rejections raised before + // anything was sent (rate limiter, open breaker) stay retryable. if retryOpts != nil && retryOpts.Attempts > 0 { return ExecuteWithRetry(ctx, *retryOpts, func(ctx context.Context) (*FetchResponse[T], error) { - return doRequest(ctx) + result, err := doRequest(ctx) + if err != nil && !gate.eligible.Load() && !isPreSendRejection(err) { + return nil, ¬RetryEligibleError{err: err} + } + return result, err }) } @@ -209,6 +236,7 @@ func executeHTTPRequest[T any]( extraHeaders http.Header, hooks *LifecycleHooks, validate func(data any) []string, + gate *retryGate, ) (*FetchResponse[T], error) { // Prepare body var bodyReader io.Reader @@ -284,6 +312,9 @@ func executeHTTPRequest[T any]( req.Header.Set("Authorization", fmt.Sprintf("%s %s", scheme, token)) } + // The request is final now: hooks and the auth provider have run. + gate.record(req) + // Execute the HTTP request httpClient := client.httpClient if httpClient == nil { diff --git a/go/fetch/options.go b/go/fetch/options.go index bd25786..cca662c 100644 --- a/go/fetch/options.go +++ b/go/fetch/options.go @@ -51,7 +51,17 @@ type RetryOptions struct { // FastFirst, when true, fires the first retry with zero delay. FastFirst bool // OnRejection is consulted before each retry. If nil, all attempts use RetryDefault. + // It is never consulted for a request that is not retry-eligible (see + // AllowNonIdempotent): such a request makes exactly one attempt. OnRejection OnRejectionFunc + // AllowNonIdempotent opts a non-idempotent request (POST, PATCH, โ€ฆ) into + // retries. By default only idempotent methods per RFC 9110 ยง9.2.2 (GET, HEAD, + // OPTIONS, TRACE, PUT, DELETE) are retried, because re-sending a POST that + // timed out or failed with a 429/5xx can execute its side effect twice โ€” the + // server may have acted on the first attempt. A request carrying a non-empty + // Idempotency-Key header is retried without this flag, since the key is the + // server-side contract that makes a replay safe. + AllowNonIdempotent bool } // RateLimitRetryOptions aliases RetryOptions for rate-limit-specific retry configuration, diff --git a/go/fetch/retry.go b/go/fetch/retry.go index bfaeabb..99b2ae0 100644 --- a/go/fetch/retry.go +++ b/go/fetch/retry.go @@ -5,9 +5,80 @@ import ( "errors" "math" "math/rand" + "net/http" + "strings" + "sync/atomic" "time" ) +// IdempotencyKeyHeader is the request header that makes a non-idempotent +// request safe to retry: a server that honours it deduplicates replays that +// carry the same key. A request with a non-empty value for it is retried even +// when RetryOptions.AllowNonIdempotent is false. +const IdempotencyKeyHeader = "Idempotency-Key" + +// IsIdempotentMethod reports whether method is idempotent per RFC 9110 ยง9.2.2 +// (GET, HEAD, OPTIONS, TRACE, PUT, DELETE), compared case-insensitively. POST, +// PATCH, CONNECT and anything unrecognised are not. +func IsIdempotentMethod(method string) bool { + switch strings.ToUpper(method) { + case http.MethodGet, http.MethodHead, http.MethodOptions, http.MethodTrace, http.MethodPut, http.MethodDelete: + return true + default: + return false + } +} + +// isRetryEligible reports whether a failed attempt of this request may be +// re-sent: its method is idempotent, the caller opted in via +// AllowNonIdempotent, or it carries a non-empty Idempotency-Key header. +func isRetryEligible(method string, header http.Header, opts *RetryOptions) bool { + if IsIdempotentMethod(method) { + return true + } + if opts != nil && opts.AllowNonIdempotent { + return true + } + // http.Header.Get canonicalises, so any casing the caller used matches. + return strings.TrimSpace(header.Get(IdempotencyKeyHeader)) != "" +} + +// retryGate carries the retry-eligibility of one Fetch call out of the attempt +// that computed it. It is seeded from the pre-hook view of the request and +// overwritten by each attempt once pre-request hooks and the auth provider +// have run, since either can change the method or add an Idempotency-Key. +// Atomic because a timed-out attempt's goroutine may still be writing it. +type retryGate struct { + eligible atomic.Bool + opts *RetryOptions +} + +func (g *retryGate) record(req *http.Request) { + if g == nil { + return + } + g.eligible.Store(isRetryEligible(req.Method, req.Header, g.opts)) +} + +// notRetryEligibleError marks a failure of a request that must not be re-sent. +// ExecuteWithRetry unwraps it and returns the underlying error immediately, +// without consulting OnRejection and without the RetryError wrapper. +type notRetryEligibleError struct { + err error +} + +func (e *notRetryEligibleError) Error() string { return e.err.Error() } +func (e *notRetryEligibleError) Unwrap() error { return e.err } + +// isPreSendRejection reports whether err was raised before anything reached +// the network (the in-process rate limiter or an open circuit breaker). Such +// a request produced no side effect, so retrying it is safe for any method. +func isPreSendRejection(err error) bool { + var rl *RateLimitError + var cb *CircuitBreakerError + return errors.As(err, &rl) || errors.As(err, &cb) +} + // CalculateBackoff computes the backoff duration for a given attempt using exponential backoff with jitter. // // The formula is: interval = initialInterval * (factor ^ attempt) +/- jitter @@ -74,6 +145,9 @@ func ExecuteWithRetry[T any](ctx context.Context, opts RetryOptions, fn func(ctx if err == nil { return result, nil } + if ne, ok := err.(*notRetryEligibleError); ok { + return zero, ne.err + } lastErr = err // If this is the last attempt, don't bother with retry bookkeeping. diff --git a/go/fetch/retry_idempotency_test.go b/go/fetch/retry_idempotency_test.go new file mode 100644 index 0000000..a1c2440 --- /dev/null +++ b/go/fetch/retry_idempotency_test.go @@ -0,0 +1,351 @@ +package fetch + +import ( + "bufio" + "context" + "encoding/json" + "errors" + "net" + "net/http" + "net/http/httptest" + "os" + "strings" + "sync" + "sync/atomic" + "testing" + "time" +) + +// retryIdempotencyCorpus is spec/retry-idempotency-corpus.json, shared with the +// other four ports โ€” see the corpus for why the cases are not inlined here. +type retryIdempotencyCorpus struct { + IdempotencyKeyHeader string `json:"idempotencyKeyHeader"` + Methods struct { + Idempotent []string `json:"idempotent"` + NonIdempotent []string `json:"nonIdempotent"` + } `json:"methods"` + Retry struct { + Attempts int `json:"attempts"` + InitialIntervalMs int `json:"initialIntervalMs"` + Factor float64 `json:"factor"` + JitterAdjustment float64 `json:"jitterAdjustment"` + } `json:"retry"` + TimeoutMs int `json:"timeoutMs"` + Cases []retryIdempotencyCase `json:"cases"` + Abort struct { + TimeoutMs int `json:"timeoutMs"` + CloseSlackMs int `json:"closeSlackMs"` + } `json:"timeoutAbort"` +} + +type retryIdempotencyCase struct { + Name string `json:"name"` + Method string `json:"method"` + Body string `json:"body"` + Headers map[string]string `json:"headers"` + AllowNonIdempotent bool `json:"allowNonIdempotent"` + Respond struct { + Status int `json:"status"` + Headers map[string]string `json:"headers"` + Hang bool `json:"hang"` + } `json:"respond"` + ExpectedAttempts int `json:"expectedAttempts"` +} + +func loadRetryIdempotencyCorpus(t *testing.T) retryIdempotencyCorpus { + t.Helper() + raw, err := os.ReadFile("../../spec/retry-idempotency-corpus.json") + if err != nil { + t.Fatalf("retry-idempotency corpus must be readable: %v", err) + } + var c retryIdempotencyCorpus + if err := json.Unmarshal(raw, &c); err != nil { + t.Fatalf("retry-idempotency corpus must parse: %v", err) + } + return c +} + +// Positive control: a corpus that failed to load would leave every table empty +// and every loop below trivially green. +func TestRetryIdempotencyCorpusLoaded(t *testing.T) { + c := loadRetryIdempotencyCorpus(t) + if len(c.Cases) == 0 || len(c.Methods.Idempotent) == 0 || len(c.Methods.NonIdempotent) == 0 { + t.Fatalf("corpus is missing cases: %+v", c) + } + if c.Retry.Attempts == 0 || c.TimeoutMs == 0 || c.Abort.TimeoutMs == 0 || c.Abort.CloseSlackMs == 0 { + t.Fatalf("corpus is missing knobs: %+v", c) + } + if c.IdempotencyKeyHeader != IdempotencyKeyHeader { + t.Fatalf("corpus idempotency header %q != IdempotencyKeyHeader %q", c.IdempotencyKeyHeader, IdempotencyKeyHeader) + } +} + +func TestIsIdempotentMethodCorpus(t *testing.T) { + c := loadRetryIdempotencyCorpus(t) + for _, m := range c.Methods.Idempotent { + if !IsIdempotentMethod(m) { + t.Errorf("IsIdempotentMethod(%q) = false, want true", m) + } + } + for _, m := range c.Methods.NonIdempotent { + if IsIdempotentMethod(m) { + t.Errorf("IsIdempotentMethod(%q) = true, want false", m) + } + } +} + +func corpusClient(c retryIdempotencyCorpus, allowNonIdempotent bool) *Client { + return NewClientBuilder(). + WithRetry(&RetryOptions{ + Attempts: c.Retry.Attempts, + InitialInterval: time.Duration(c.Retry.InitialIntervalMs) * time.Millisecond, + Factor: c.Retry.Factor, + JitterFraction: c.Retry.JitterAdjustment, + OnRejection: DefaultRetryOptions.OnRejection, + AllowNonIdempotent: allowNonIdempotent, + }). + WithTimeout(time.Duration(c.TimeoutMs) * time.Millisecond). + Build() +} + +// TestRetryIdempotencyCorpus counts what the SERVER received, because the bug +// is a side effect executing more than once server-side. +func TestRetryIdempotencyCorpus(t *testing.T) { + c := loadRetryIdempotencyCorpus(t) + for _, tc := range c.Cases { + t.Run(tc.Name, func(t *testing.T) { + var hits atomic.Int32 + release := make(chan struct{}) + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + hits.Add(1) + if tc.Respond.Hang { + select { + case <-r.Context().Done(): + case <-release: + } + return + } + for k, v := range tc.Respond.Headers { + w.Header().Set(k, v) + } + w.WriteHeader(tc.Respond.Status) + })) + // Release hung handlers before Close, which waits for them. + defer srv.Close() + defer close(release) + + headers := http.Header{} + for k, v := range tc.Headers { + // Raw map assignment keeps the corpus's casing, so the + // case-insensitivity case really sends a lowercase key through. + headers[k] = []string{v} + } + var body any + if tc.Body != "" { + body = tc.Body + } + + _, err := Fetch[any](context.Background(), corpusClient(c, tc.AllowNonIdempotent), tc.Method, srv.URL, body, &RequestOptions{Headers: headers}) + if err == nil { + t.Fatal("expected the request to fail") + } + if got := int(hits.Load()); got != tc.ExpectedAttempts { + t.Fatalf("server received %d requests, want %d (err: %v)", got, tc.ExpectedAttempts, err) + } + + var retryErr *RetryError + if tc.ExpectedAttempts == 1 && errors.As(err, &retryErr) { + t.Fatalf("a request that was not retried must surface the underlying error, got %T: %v", err, err) + } + if tc.ExpectedAttempts > 1 && !errors.As(err, &retryErr) { + t.Fatalf("an exhausted retry must surface *RetryError, got %T: %v", err, err) + } + }) + } +} + +// TestRetryIneligibleSkipsOnRejection pins that the callback is not consulted +// for a request that is not retry-eligible: a custom callback is not an opt-in. +func TestRetryIneligibleSkipsOnRejection(t *testing.T) { + var hits, consulted atomic.Int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + hits.Add(1) + w.WriteHeader(http.StatusServiceUnavailable) + })) + defer srv.Close() + + client := NewClientBuilder().WithRetry(&RetryOptions{ + Attempts: 2, + InitialInterval: time.Millisecond, + OnRejection: func(RetryContext) (RetryDecision, time.Duration) { + consulted.Add(1) + return RetryDefault, 0 + }, + }).Build() + + _, err := Post[any](context.Background(), client, srv.URL, "{}", nil) + var httpErr *HTTPResponseError + if !errors.As(err, &httpErr) || httpErr.StatusCode != http.StatusServiceUnavailable { + t.Fatalf("expected the 503 unwrapped, got %T: %v", err, err) + } + if hits.Load() != 1 || consulted.Load() != 0 { + t.Fatalf("hits=%d consulted=%d, want 1 and 0", hits.Load(), consulted.Load()) + } +} + +// TestRetryEligibilitySeesPreRequestHook: the gate is evaluated on the FINAL +// request, so a hook that adds an Idempotency-Key opts the POST in. +func TestRetryEligibilitySeesPreRequestHook(t *testing.T) { + var hits atomic.Int32 + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + hits.Add(1) + w.WriteHeader(http.StatusServiceUnavailable) + })) + defer srv.Close() + + retry := DefaultRetryOptions + retry.InitialInterval = time.Millisecond + retry.JitterFraction = 0 + client := NewClientBuilder().WithRetry(&retry).WithHooks(&LifecycleHooks{ + PreRequest: func(url string, req *http.Request) (string, *http.Request) { + req.Header.Set(IdempotencyKeyHeader, "from-hook") + return url, req + }, + }).Build() + + _, _ = Post[any](context.Background(), client, srv.URL, "{}", nil) + if got := hits.Load(); got != int32(1+retry.Attempts) { + t.Fatalf("server received %d requests, want %d", got, 1+retry.Attempts) + } +} + +// abortServer is a raw TCP server that reads each request's head and never +// answers, recording when each request arrived and when its connection was +// closed by the client (read returns EOF). Raw sockets, not net/http, so +// "closed" means the peer really closed the TCP connection. +type abortServer struct { + ln net.Listener + mu sync.Mutex + arrivals []time.Time + closes []time.Time // indexed like arrivals; zero until closed + closed chan int +} + +func newAbortServer(t *testing.T) *abortServer { + t.Helper() + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("listen: %v", err) + } + s := &abortServer{ln: ln, closed: make(chan int, 16)} + go func() { + for { + conn, err := ln.Accept() + if err != nil { + return + } + go s.serve(conn) + } + }() + t.Cleanup(func() { _ = ln.Close() }) + return s +} + +func (s *abortServer) serve(conn net.Conn) { + defer conn.Close() + r := bufio.NewReader(conn) + // Read the request head; the body (if any) is drained by the EOF loop. + for { + line, err := r.ReadString('\n') + if err != nil { + return + } + if strings.TrimRight(line, "\r\n") == "" { + break + } + } + s.mu.Lock() + idx := len(s.arrivals) + s.arrivals = append(s.arrivals, time.Now()) + s.closes = append(s.closes, time.Time{}) + s.mu.Unlock() + + buf := make([]byte, 512) + for { + if _, err := r.Read(buf); err != nil { + break + } + } + s.mu.Lock() + s.closes[idx] = time.Now() + s.mu.Unlock() + s.closed <- idx +} + +func (s *abortServer) snapshot() ([]time.Time, []time.Time) { + s.mu.Lock() + defer s.mu.Unlock() + return append([]time.Time(nil), s.arrivals...), append([]time.Time(nil), s.closes...) +} + +// TestTimeoutClosesTheAttemptConnection: a timed-out attempt must be cancelled, +// not abandoned โ€” its connection closes, so the server can stop work โ€” and a +// retry must never overlap the attempt it replaces. +func TestTimeoutClosesTheAttemptConnection(t *testing.T) { + c := loadRetryIdempotencyCorpus(t) + timeout := time.Duration(c.Abort.TimeoutMs) * time.Millisecond + bound := timeout + time.Duration(c.Abort.CloseSlackMs)*time.Millisecond + + for _, method := range []string{http.MethodGet, http.MethodPost} { + t.Run(method, func(t *testing.T) { + srv := newAbortServer(t) + client := NewClientBuilder(). + WithRetry(&RetryOptions{Attempts: 1, InitialInterval: time.Millisecond, OnRejection: DefaultRetryOptions.OnRejection}). + WithTimeout(timeout). + Build() + + _, err := Fetch[any](context.Background(), client, method, "http://"+srv.ln.Addr().String()+"/hang", nil, nil) + var timeoutErr *TimeoutError + if !errors.As(err, &timeoutErr) { + t.Fatalf("expected a *TimeoutError, got %T: %v", err, err) + } + + deadline := time.After(bound) + wait: + for { + select { + case idx := <-srv.closed: + if idx == 0 { + break wait + } + case <-deadline: + break wait + } + } + arrivals, closes := srv.snapshot() + if len(arrivals) == 0 { + t.Fatal("the server never saw the first attempt") + } + if closes[0].IsZero() { + t.Fatalf("the first attempt's connection was still open %v after it arrived", bound) + } + if lag := closes[0].Sub(arrivals[0]); lag > bound { + t.Fatalf("the first attempt's connection closed %v after it arrived, want <= %v", lag, bound) + } + + switch method { + case http.MethodGet: + if len(arrivals) != 2 { + t.Fatalf("GET: server saw %d attempts, want 2", len(arrivals)) + } + if !closes[0].Before(arrivals[1]) { + t.Fatalf("GET: the first connection closed at %v, AFTER the retry arrived at %v", closes[0], arrivals[1]) + } + case http.MethodPost: + if len(arrivals) != 1 { + t.Fatalf("POST: server saw %d attempts, want 1", len(arrivals)) + } + } + }) + } +} diff --git a/go/fetch/timeout.go b/go/fetch/timeout.go index 95b8350..3896eac 100644 --- a/go/fetch/timeout.go +++ b/go/fetch/timeout.go @@ -5,8 +5,21 @@ import ( "time" ) +// attemptUnwindGrace bounds how long a timed-out attempt is given to unwind +// after its context is cancelled. net/http closes the connection synchronously +// on cancellation (microseconds), so this is only ever reached by caller code +// that ignores its context โ€” e.g. a hook or auth provider that blocks โ€” and it +// keeps such code from stretching the timeout indefinitely. +const attemptUnwindGrace = time.Second + // ExecuteWithTimeout runs fn with a timeout derived from the given context. // If the function does not complete within the timeout, a *TimeoutError is returned. +// +// The timeout CANCELS the attempt rather than abandoning it: fn's context is +// cancelled, which makes net/http close the connection so the server can see +// the request go away, and ExecuteWithTimeout waits (up to attemptUnwindGrace) +// for fn to return before reporting the timeout. A retry that follows therefore +// never overlaps the attempt it replaces. func ExecuteWithTimeout[T any](ctx context.Context, timeout time.Duration, fn func(ctx context.Context) (T, error)) (T, error) { ctx, cancel := context.WithTimeout(ctx, timeout) defer cancel() @@ -25,6 +38,11 @@ func ExecuteWithTimeout[T any](ctx context.Context, timeout time.Duration, fn fu select { case <-ctx.Done(): + cancel() + select { + case <-ch: + case <-time.After(attemptUnwindGrace): + } var zero T if ctx.Err() == context.DeadlineExceeded { return zero, &TimeoutError{Timeout: timeout} diff --git a/python/README.md b/python/README.md index e5d91ac..962414e 100644 --- a/python/README.md +++ b/python/README.md @@ -277,14 +277,51 @@ Out of the box, smooai-fetch is configured for the real world: - 2 automatic retries on failure - Exponential backoff: 500ms -> 1s -> 2s - Jitter to prevent thundering herds -- Only retries on network errors or 5xx responses +- Only retries on network errors, timeouts, 429 or 5xx responses +- **Only idempotent requests are retried** (see below) **Timeout Protection:** -- 10-second default timeout -- Prevents indefinite hangs +- 30-second default timeout, applied to each attempt as a whole +- A timed-out attempt is cancelled and its connection closed before any retry, so the server sees it abandoned - Configurable per request +### Retries are method-aware (4.0 behavior change) + +Retrying re-sends the request. For a POST that timed out or got a 429/5xx, that can +run its side effect twice โ€” a second charge, a duplicate message, a second billed +image. So by default only methods that are idempotent per RFC 9110 are retried: +`GET`, `HEAD`, `OPTIONS`, `TRACE`, `PUT`, `DELETE`. A `POST` or `PATCH` makes exactly +one attempt and raises its own error (`HTTPResponseError`, `TimeoutError`, ...), not +`RetryError`. This includes a **429 with `Retry-After`**: it is not retried for a POST +unless you opt in. + +Opt in when the endpoint is safe to replay, either per client/request: + +```python +from smooai_fetch import FetchOptions, RetryOptions, fetch + +await fetch( + "https://api.example.com/jobs", + FetchOptions(method="POST", body=job, retry=RetryOptions(allow_non_idempotent=True)), +) +``` + +or by sending an `Idempotency-Key` header (any casing, non-empty), which lets a +server that honours it deduplicate the replay: + +```python +await fetch( + "https://api.example.com/charges", + FetchOptions(method="POST", body=charge, headers={"Idempotency-Key": charge_id}), +) +``` + +Eligibility is decided after pre-request hooks and the auth provider run, so a hook +can add the key. `is_idempotent_method()` and `IDEMPOTENCY_KEY_HEADER` are exported. +The in-process rate limiter's own retry loop is unaffected: it rejects before anything +is sent. + **Rate Limit Handling:** - Respects Retry-After headers diff --git a/python/src/smooai_fetch/__init__.py b/python/src/smooai_fetch/__init__.py index 307aa02..f6ae9e2 100644 --- a/python/src/smooai_fetch/__init__.py +++ b/python/src/smooai_fetch/__init__.py @@ -39,7 +39,7 @@ from smooai_fetch._response import FetchResponse # Retry utilities -from smooai_fetch._retry import calculate_backoff, is_retryable +from smooai_fetch._retry import IDEMPOTENCY_KEY_HEADER, calculate_backoff, is_idempotent_method, is_retryable # Types from smooai_fetch._types import ( @@ -99,7 +99,9 @@ "TimeoutError", # Utilities "calculate_backoff", + "is_idempotent_method", "is_retryable", + "IDEMPOTENCY_KEY_HEADER", "SlidingWindowRateLimiter", "CircuitBreaker", ] diff --git a/python/src/smooai_fetch/_client.py b/python/src/smooai_fetch/_client.py index 5a75a77..6bc4c9d 100644 --- a/python/src/smooai_fetch/_client.py +++ b/python/src/smooai_fetch/_client.py @@ -2,6 +2,7 @@ from __future__ import annotations +import asyncio import inspect import json from typing import Any, TypeVar @@ -21,7 +22,7 @@ from smooai_fetch._errors import TimeoutError as FetchTimeoutError from smooai_fetch._rate_limit import SlidingWindowRateLimiter from smooai_fetch._response import FetchResponse -from smooai_fetch._retry import execute_with_retry, is_retryable +from smooai_fetch._retry import execute_with_retry, is_retry_eligible, is_retryable from smooai_fetch._timeout import create_timeout from smooai_fetch._types import ( FetchOptions, @@ -265,6 +266,10 @@ async def fetch( Returns: A FetchResponse containing the parsed response data. + Only idempotent requests are retried: a POST/PATCH makes one attempt unless + ``RetryOptions.allow_non_idempotent`` is set or it carries a non-empty + ``Idempotency-Key`` header, and then raises its own error unwrapped. + Raises: HTTPResponseError: For non-2xx responses (after retries exhausted). RetryError: When all retry attempts are exhausted. @@ -335,15 +340,22 @@ async def _do_request() -> FetchResponse[Any]: **request_kwargs, "headers": _inject_trace_context(request_kwargs.get("headers") or {}), } + # httpx's timeout is per PHASE (connect/read/write/pool), so a server + # trickling bytes never trips it. `asyncio.timeout` bounds the whole + # attempt; on expiry it CANCELS the request coroutine, and leaving the + # per-attempt `AsyncClient` block closes the connection โ€” so the server + # sees the attempt abandoned before any retry is sent, rather than the + # timed-out request running on beside it. try: - async with httpx.AsyncClient() as client: - response = await client.request(**attempt_kwargs) - return _parse_response(response, schema) - except httpx.TimeoutException as e: + async with asyncio.timeout(timeout_options.timeout_ms / 1000.0): + async with httpx.AsyncClient() as client: + response = await client.request(**attempt_kwargs) + except (httpx.TimeoutException, TimeoutError) as e: raise FetchTimeoutError( timeout_ms=timeout_options.timeout_ms, url=current_url, ) from e + return _parse_response(response, schema) async def _gated() -> FetchResponse[Any]: # Rate limit check @@ -391,18 +403,31 @@ async def _do_with_rl_retry() -> FetchResponse[Any]: else: return await request_fn() + # Decided AFTER the pre-request hook and auth provider, which can change the + # method or add an Idempotency-Key header. An ineligible request (POST/PATCH + # without an opt-in) gets exactly one attempt and its own error, unwrapped: + # re-sending it could repeat its side effect (SMOODEV-3375). + retry_eligible = is_retry_eligible( + str(request_kwargs.get("method", "GET")), + request_kwargs.get("headers"), + retry_options, + ) + # Execute with error hook wrapping try: - # Wrap with retry logic - def should_retry(error: Exception, attempt: int) -> bool | float: - return _should_retry_default(error, attempt, retry_options) - - result = await execute_with_retry( - func=_execute, - options=retry_options, - should_retry=should_retry, - get_retry_after=_get_retry_after, - ) + if retry_eligible: + # Wrap with retry logic + def should_retry(error: Exception, attempt: int) -> bool | float: + return _should_retry_default(error, attempt, retry_options) + + result = await execute_with_retry( + func=_execute, + options=retry_options, + should_retry=should_retry, + get_retry_after=_get_retry_after, + ) + else: + result = await _execute() except Exception as error: # Apply post-response error hook if hooks and hooks.post_response_error: diff --git a/python/src/smooai_fetch/_retry.py b/python/src/smooai_fetch/_retry.py index 3ae3a9c..0304b92 100644 --- a/python/src/smooai_fetch/_retry.py +++ b/python/src/smooai_fetch/_retry.py @@ -5,8 +5,8 @@ import asyncio import random import time -from collections.abc import Awaitable, Callable -from typing import TypeVar +from collections.abc import Awaitable, Callable, Mapping +from typing import Any, TypeVar from smooai_fetch._errors import HTTPResponseError, RetryError from smooai_fetch._types import ( @@ -17,6 +17,42 @@ T = TypeVar("T") +IDEMPOTENCY_KEY_HEADER = "Idempotency-Key" +"""Header whose non-empty presence makes a non-idempotent request retry-eligible. + +A server honouring it (draft-ietf-httpapi-idempotency-key-header) deduplicates +replays, so re-sending the request cannot repeat its side effect. +""" + +_IDEMPOTENT_METHODS = frozenset({"GET", "HEAD", "OPTIONS", "TRACE", "PUT", "DELETE"}) + + +def is_idempotent_method(method: str) -> bool: + """Whether ``method`` is idempotent per RFC 9110 section 9.2.2 (case-insensitive). + + GET, HEAD, OPTIONS, TRACE, PUT and DELETE are. POST, PATCH, CONNECT and any + unrecognised method are not: re-sending them can repeat a side effect. + """ + return method.upper() in _IDEMPOTENT_METHODS + + +def _has_idempotency_key(headers: Mapping[str, Any] | None) -> bool: + if not headers: + return False + target = IDEMPOTENCY_KEY_HEADER.lower() + return any(str(key).lower() == target and str(value).strip() != "" for key, value in headers.items()) + + +def is_retry_eligible(method: str, headers: Mapping[str, Any] | None, options: RetryOptions) -> bool: + """Whether a failed attempt of this request may be retried at all (SMOODEV-3375). + + Retrying re-sends the request, so it is only safe when doing so cannot + duplicate a side effect: the method is idempotent, the request carries a + non-empty ``Idempotency-Key`` the server can deduplicate on, or the caller + explicitly opted in with ``RetryOptions.allow_non_idempotent``. + """ + return is_idempotent_method(method) or options.allow_non_idempotent or _has_idempotency_key(headers) + def calculate_backoff(attempt: int, options: RetryOptions) -> float: """Calculate the backoff delay in seconds for a given retry attempt. diff --git a/python/src/smooai_fetch/_types.py b/python/src/smooai_fetch/_types.py index 6694951..9dc0495 100644 --- a/python/src/smooai_fetch/_types.py +++ b/python/src/smooai_fetch/_types.py @@ -131,6 +131,19 @@ class RetryOptions: Receives a `RetryContext` and returns an `OnRejectionDecision` that can override the default delay, skip the attempt, or abort retrying entirely. + Never consulted for a request that is not retry-eligible (see + `allow_non_idempotent`). + """ + + allow_non_idempotent: bool = False + """Allow retrying non-idempotent methods (POST, PATCH, ...). Default False. + + By default only idempotent methods (GET, HEAD, OPTIONS, TRACE, PUT, DELETE) + are retried, because re-sending a POST that timed out or got a 429/5xx can + execute its side effect twice (a second charge, a duplicate message, a + second billed image). A request carrying a non-empty ``Idempotency-Key`` + header is retried without this flag. Set it only when the endpoint is safe + to replay. """ diff --git a/python/tests/test_retry_idempotency.py b/python/tests/test_retry_idempotency.py new file mode 100644 index 0000000..573b6eb --- /dev/null +++ b/python/tests/test_retry_idempotency.py @@ -0,0 +1,232 @@ +"""Retry safety: only idempotent (or explicitly opted-in) requests are retried, +and a timed-out attempt is cancelled rather than abandoned (SMOODEV-3375). + +Every case comes from spec/retry-idempotency-corpus.json, shared with the other +four ports -- do not inline cases here. They run against a REAL local server that +counts the requests it receives, because the bug is a side effect executing more +than once server-side, and only the server can count that. +""" + +from __future__ import annotations + +import asyncio +import json +import time +from dataclasses import dataclass, field +from pathlib import Path +from typing import Any + +import pytest + +from smooai_fetch import ( + IDEMPOTENCY_KEY_HEADER, + FetchOptions, + RetryError, + RetryOptions, + TimeoutOptions, + fetch, + is_idempotent_method, +) +from smooai_fetch._retry import is_retry_eligible + +CORPUS: dict[str, Any] = json.loads((Path(__file__).parents[2] / "spec" / "retry-idempotency-corpus.json").read_text()) + + +@dataclass +class _Connection: + arrived_at: float + closed_at: float | None = None + + +@dataclass +class _Server: + """Minimal HTTP/1.1 server: answers each request per `respond`, or never (hang).""" + + respond: dict[str, Any] + requests: int = 0 + connections: list[_Connection] = field(default_factory=list) + _server: asyncio.Server | None = None + _writers: list[asyncio.StreamWriter] = field(default_factory=list) + + async def _handle(self, reader: asyncio.StreamReader, writer: asyncio.StreamWriter) -> None: + conn = _Connection(arrived_at=time.monotonic()) + self.connections.append(conn) + self._writers.append(writer) + try: + head = await reader.readuntil(b"\r\n\r\n") + length = 0 + for line in head.decode("latin-1").split("\r\n")[1:]: + name, _, value = line.partition(":") + if name.strip().lower() == "content-length": + length = int(value.strip()) + if length: + await reader.readexactly(length) + self.requests += 1 + + if not self.respond.get("hang"): + extra = "".join(f"{k}: {v}\r\n" for k, v in self.respond.get("headers", {}).items()) + status_line = f"HTTP/1.1 {self.respond['status']} X\r\n" + writer.write(f"{status_line}Content-Length: 0\r\nConnection: close\r\n{extra}\r\n".encode()) + await writer.drain() + return + + # Hang: never answer. Block until the CLIENT closes the connection -- + # read() returns b"" at EOF -- and record when that happened. + while await reader.read(1024): + pass + conn.closed_at = time.monotonic() + except (asyncio.IncompleteReadError, ConnectionError): + conn.closed_at = time.monotonic() + finally: + writer.close() + + async def __aenter__(self) -> _Server: + self._server = await asyncio.start_server(self._handle, "127.0.0.1", 0) + return self + + async def __aexit__(self, *_: object) -> None: + assert self._server is not None + self._server.close() + for writer in self._writers: + writer.close() + + @property + def url(self) -> str: + assert self._server is not None + port = self._server.sockets[0].getsockname()[1] + return f"http://127.0.0.1:{port}/anything" + + +def _retry_options(allow_non_idempotent: bool = False) -> RetryOptions: + knobs = CORPUS["retry"] + return RetryOptions( + attempts=knobs["attempts"], + initial_interval_ms=knobs["initialIntervalMs"], + factor=knobs["factor"], + jitter=knobs["jitterAdjustment"], + allow_non_idempotent=allow_non_idempotent, + ) + + +def test_corpus_loaded() -> None: + # Positive control: a corpus that failed to load would parametrize nothing + # and every assertion below would pass vacuously. + assert len(CORPUS["cases"]) >= 10 + assert CORPUS["methods"]["idempotent"] and CORPUS["methods"]["nonIdempotent"] + assert CORPUS["idempotencyKeyHeader"] == IDEMPOTENCY_KEY_HEADER + + +@pytest.mark.parametrize("method", CORPUS["methods"]["idempotent"]) +def test_idempotent_methods(method: str) -> None: + assert is_idempotent_method(method) + + +@pytest.mark.parametrize("method", CORPUS["methods"]["nonIdempotent"]) +def test_non_idempotent_methods(method: str) -> None: + assert not is_idempotent_method(method) + + +@pytest.mark.parametrize("case", CORPUS["cases"], ids=[c["name"] for c in CORPUS["cases"]]) +async def test_server_side_attempts(case: dict[str, Any]) -> None: + headers = case.get("headers", {}) + if any(isinstance(v, str) and v and not v.strip() for v in headers.values()): + # h11 refuses to put a whitespace-only header value on the wire (it + # raises LocalProtocolError before connecting), so this case can never + # reach a server from Python. Assert the rule it pins down directly. + assert not is_retry_eligible(case["method"], headers, _retry_options(case.get("allowNonIdempotent", False))) + return + async with _Server(respond=case["respond"]) as server: + options = FetchOptions( + method=case["method"], + headers=dict(case.get("headers", {})), + body=case.get("body"), + retry=_retry_options(case.get("allowNonIdempotent", False)), + timeout=TimeoutOptions(timeout_ms=CORPUS["timeoutMs"]), + ) + with pytest.raises(Exception) as exc_info: + await fetch(server.url, options) + + assert server.requests == case["expectedAttempts"], ( + f"{case['name']}: server received {server.requests} request(s), expected {case['expectedAttempts']}" + ) + if case["expectedAttempts"] == 1: + # A request that was not retried surfaces its own error, not the + # retries-exhausted wrapper. + assert not isinstance(exc_info.value, RetryError) + + +async def _wait_until_closed(conn: _Connection, deadline_s: float) -> None: + end = time.monotonic() + deadline_s + while conn.closed_at is None and time.monotonic() < end: + await asyncio.sleep(0.01) + + +async def test_timed_out_get_closes_first_connection_before_retry() -> None: + knobs = CORPUS["timeoutAbort"] + ceiling_s = (knobs["timeoutMs"] + knobs["closeSlackMs"]) / 1000.0 + async with _Server(respond={"hang": True}) as server: + options = FetchOptions( + method="GET", + retry=_retry_options(), + timeout=TimeoutOptions(timeout_ms=knobs["timeoutMs"]), + ) + with pytest.raises(Exception): + await fetch(server.url, options) + + assert len(server.connections) == CORPUS["retry"]["attempts"] + 1 + first, second = server.connections[0], server.connections[1] + await _wait_until_closed(first, ceiling_s) + + assert first.closed_at is not None, "first attempt's connection was never closed -- the timeout abandoned it" + assert first.closed_at - first.arrived_at <= ceiling_s + assert first.closed_at <= second.arrived_at, "retry was sent while the timed-out attempt was still open" + + +async def test_timed_out_post_closes_its_connection() -> None: + knobs = CORPUS["timeoutAbort"] + ceiling_s = (knobs["timeoutMs"] + knobs["closeSlackMs"]) / 1000.0 + async with _Server(respond={"hang": True}) as server: + options = FetchOptions( + method="POST", + body='{"a":1}', + retry=_retry_options(), + timeout=TimeoutOptions(timeout_ms=knobs["timeoutMs"]), + ) + with pytest.raises(Exception): + await fetch(server.url, options) + + assert len(server.connections) == 1 + conn = server.connections[0] + await _wait_until_closed(conn, ceiling_s) + + assert conn.closed_at is not None, "timed-out POST's connection was never closed" + assert conn.closed_at - conn.arrived_at <= ceiling_s + + +async def test_builder_path_is_gated_too() -> None: + from smooai_fetch import FetchBuilder + + async with _Server(respond={"status": 503}) as server: + builder = FetchBuilder().with_retry(_retry_options()) + with pytest.raises(Exception): + await builder.fetch(server.url, method="POST", body='{"a":1}') + assert server.requests == 1 + + +async def test_pre_request_hook_can_add_idempotency_key() -> None: + from smooai_fetch import LifecycleHooks + + def add_key(url: str, kwargs: dict[str, Any]) -> tuple[str, dict[str, Any]]: + kwargs["headers"] = {**kwargs.get("headers", {}), IDEMPOTENCY_KEY_HEADER: "hook-key"} + return url, kwargs + + async with _Server(respond={"status": 503}) as server: + options = FetchOptions( + method="POST", + body='{"a":1}', + retry=_retry_options(), + hooks=LifecycleHooks(pre_request=add_key), + ) + with pytest.raises(RetryError): + await fetch(server.url, options) + assert server.requests == CORPUS["retry"]["attempts"] + 1 diff --git a/rust/fetch/README.md b/rust/fetch/README.md index 4a767cc..fa2d0ed 100644 --- a/rust/fetch/README.md +++ b/rust/fetch/README.md @@ -120,6 +120,43 @@ let response = fetch::("https://api.github.com/user/repos", i // - Your code continues normally ``` +### Retries Never Duplicate a Side Effect + +Since 4.0.0, only requests that are safe to send twice are retried. A retry +re-executes the request on the server, so a POST that timed out after the server +started work โ€” or came back 5xx/429 โ€” may already have created the record, sent +the message or charged the card. + +- **Retried by default:** `GET`, `HEAD`, `OPTIONS`, `PUT`, `DELETE` (idempotent + per RFC 9110 ยง9.2.2). +- **Never retried by default:** `POST`, `PATCH` โ€” they make exactly one attempt + and return the underlying error (`FetchError::HttpResponse`, + `FetchError::Timeout`, โ€ฆ), not `FetchError::Retry`. That includes a **429 with + `Retry-After`**: the rejected POST is not re-sent unless you opt in. +- **Opting in**, when the endpoint tolerates duplicates or deduplicates them: + +```rust +use smooai_fetch::types::RetryOptions; + +// Every request through this client, including POST/PATCH: +let client = FetchBuilder::::new() + .with_retry(RetryOptions { allow_non_idempotent: true, ..Default::default() }) + .build(); + +// Or just this request: an `Idempotency-Key` header (any casing, non-empty) +// lets a server that honours it deduplicate the replay. +let mut init = RequestInit { method: Method::POST, body: Some(body), ..Default::default() }; +init.headers.insert(smooai_fetch::IDEMPOTENCY_KEY_HEADER.to_string(), key); +``` + +The check runs after the pre-request hook and auth provider, so a hook can add +the header. `is_idempotent_method` and `is_retry_eligible` are exported if you +need the same decision elsewhere. + +**Timeouts cancel the attempt.** When an attempt times out, its request future is +dropped, which closes the connection โ€” the server sees the cancel before any +retry is sent, instead of the abandoned request running on beside it. + ### Production-Ready Examples #### FetchBuilder Pattern @@ -283,12 +320,15 @@ Out of the box, smooai-fetch is configured for the real world: - 2 automatic retries on failure - Exponential backoff: 500ms -> 1s -> 2s - Jitter to prevent thundering herds -- Only retries on network errors, timeouts, or 5xx responses +- Only retries on network errors, timeouts, 429, or 5xx responses +- Only retries idempotent methods (GET, HEAD, OPTIONS, PUT, DELETE) โ€” POST and + PATCH need `allow_non_idempotent` or an `Idempotency-Key` header **Timeout Protection:** - 10-second default timeout - Prevents indefinite hangs on slow endpoints +- A timed-out attempt is cancelled (its connection closed), not abandoned - Configurable per request or per client **Rate Limit Handling:** diff --git a/rust/fetch/src/client.rs b/rust/fetch/src/client.rs index d69f2ea..8b2193f 100644 --- a/rust/fetch/src/client.rs +++ b/rust/fetch/src/client.rs @@ -388,8 +388,23 @@ pub async fn fetch_with_redirect_policy RetryOptions { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, } } diff --git a/rust/fetch/src/lib.rs b/rust/fetch/src/lib.rs index b47423f..fcc1c3b 100644 --- a/rust/fetch/src/lib.rs +++ b/rust/fetch/src/lib.rs @@ -63,6 +63,7 @@ pub use circuit_breaker::{CircuitBreaker, CircuitState, CircuitStateChangeCallba pub use error::FetchError; pub use rate_limit::SlidingWindowRateLimiter; pub use response::FetchResponse; +pub use retry::{is_idempotent_method, is_retry_eligible, IDEMPOTENCY_KEY_HEADER}; pub use types::{ AuthTokenFuture, AuthTokenProvider, FetchContainerOptions, FetchOptions, Method, RateLimitRetryOptions, RequestInit, RetryCallback, RetryContext, RetryDecision, RetryOptions, diff --git a/rust/fetch/src/retry.rs b/rust/fetch/src/retry.rs index 79210f9..4400d56 100644 --- a/rust/fetch/src/retry.rs +++ b/rust/fetch/src/retry.rs @@ -7,7 +7,7 @@ use rand::Rng; use tracing; use crate::error::FetchError; -use crate::types::{RetryContext, RetryDecision, RetryOptions}; +use crate::types::{Method, RequestInit, RetryContext, RetryDecision, RetryOptions}; /// Determine if a given status code is retryable. /// 429 (Too Many Requests) and 5xx are retryable. @@ -15,6 +15,42 @@ pub fn is_retryable(status: u16) -> bool { status == 429 || status >= 500 } +/// The header that lets a single non-idempotent request opt in to retries. +/// +/// A server that honours it (draft-ietf-httpapi-idempotency-key-header โ€” Stripe, +/// Adyen, most payment APIs) deduplicates replays carrying the same key, which +/// is exactly what makes a retried POST safe. +pub const IDEMPOTENCY_KEY_HEADER: &str = "Idempotency-Key"; + +/// Whether `method` is idempotent per RFC 9110 ยง9.2.2 โ€” i.e. sending it twice +/// has the same effect on the server as sending it once. +/// +/// GET, HEAD, OPTIONS, PUT and DELETE are; POST and PATCH are not. (TRACE is +/// idempotent too, and CONNECT is not, but [`Method`] cannot express either.) +pub fn is_idempotent_method(method: &Method) -> bool { + match method { + Method::GET | Method::HEAD | Method::OPTIONS | Method::PUT | Method::DELETE => true, + Method::POST | Method::PATCH => false, + } +} + +/// Whether a failed attempt of this request may be retried at all. +/// +/// True when the method is idempotent, when the caller opted in with +/// [`RetryOptions::allow_non_idempotent`], or when the request carries a +/// non-empty `Idempotency-Key` header (matched case-insensitively). Otherwise +/// the request makes exactly one attempt, whatever the error โ€” a timeout, a +/// 5xx or a 429 with `Retry-After` โ€” because the first attempt may already have +/// had its side effect (SMOODEV-3375: a slow image-generation POST billed three +/// images and returned none). +pub fn is_retry_eligible(init: &RequestInit, options: &RetryOptions) -> bool { + is_idempotent_method(&init.method) + || options.allow_non_idempotent + || init.headers.iter().any(|(name, value)| { + name.eq_ignore_ascii_case(IDEMPOTENCY_KEY_HEADER) && !value.trim().is_empty() + }) +} + /// Calculate the backoff delay for a given attempt using exponential backoff with jitter. /// /// - `attempt`: zero-based attempt index (0 = first retry) @@ -200,6 +236,40 @@ mod tests { assert!(!is_retryable(404)); } + fn init(method: Method, headers: &[(&str, &str)]) -> RequestInit { + RequestInit { + method, + headers: headers + .iter() + .map(|(k, v)| (k.to_string(), v.to_string())) + .collect(), + body: None, + } + } + + #[test] + fn test_retry_eligibility() { + let opts = RetryOptions::default(); + assert!(!opts.allow_non_idempotent, "default must not retry POST"); + assert!(is_retry_eligible(&init(Method::GET, &[]), &opts)); + assert!(is_retry_eligible(&init(Method::DELETE, &[]), &opts)); + assert!(!is_retry_eligible(&init(Method::POST, &[]), &opts)); + assert!(!is_retry_eligible(&init(Method::PATCH, &[]), &opts)); + assert!(is_retry_eligible( + &init(Method::POST, &[("idempotency-key", "k")]), + &opts + )); + assert!(!is_retry_eligible( + &init(Method::POST, &[("Idempotency-Key", " ")]), + &opts + )); + let opted_in = RetryOptions { + allow_non_idempotent: true, + ..Default::default() + }; + assert!(is_retry_eligible(&init(Method::PATCH, &[]), &opted_in)); + } + #[test] fn test_calculate_backoff_attempt_0() { let options = RetryOptions { @@ -210,6 +280,7 @@ mod tests { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }; let delay = calculate_backoff(0, &options); // base * factor^0 = 500 * 1 = 500 @@ -226,6 +297,7 @@ mod tests { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }; let delay = calculate_backoff(1, &options); // base * factor^1 = 500 * 2 = 1000 @@ -242,6 +314,7 @@ mod tests { max_interval_ms: Some(800), fast_first: false, on_rejection: None, + allow_non_idempotent: false, }; let delay = calculate_backoff(2, &options); // base * factor^2 = 500 * 4 = 2000, capped at 800 @@ -258,6 +331,7 @@ mod tests { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }; // With jitter=0.5, delay should be between 500 and 1500 for _ in 0..100 { @@ -277,6 +351,7 @@ mod tests { max_interval_ms: None, fast_first: true, on_rejection: None, + allow_non_idempotent: false, }; let err = FetchError::Timeout { timeout_ms: 1000 }; // attempt 0 (first retry) with fast_first=true โ†’ zero delay @@ -295,6 +370,7 @@ mod tests { max_interval_ms: None, fast_first: true, on_rejection: None, + allow_non_idempotent: false, }; let mut headers = std::collections::HashMap::new(); headers.insert("retry-after".to_string(), "3".to_string()); diff --git a/rust/fetch/src/types.rs b/rust/fetch/src/types.rs index bf52f13..eabe554 100644 --- a/rust/fetch/src/types.rs +++ b/rust/fetch/src/types.rs @@ -87,6 +87,30 @@ pub struct RetryOptions { /// callback can override the default delay, skip the attempt, or abort /// retrying entirely. pub on_rejection: Option, + /// Allow retrying non-idempotent requests (`POST`, `PATCH`). Defaults to + /// `false`. + /// + /// Retrying a request whose side effect already happened executes it + /// again: a POST that times out after the server started work, or gets a + /// 5xx/429 back, may already have created the record, sent the message or + /// charged the card. So by default only methods RFC 9110 ยง9.2.2 calls + /// idempotent (GET, HEAD, OPTIONS, PUT, DELETE) are retried, and a + /// POST/PATCH makes exactly one attempt โ€” even on a 429 with + /// `Retry-After`. + /// + /// Set this to `true` only when the endpoint tolerates duplicates, or send + /// an `Idempotency-Key` header instead, which opts that one request in and + /// lets the server deduplicate. See [`crate::retry::is_retry_eligible`]. + pub allow_non_idempotent: bool, +} + +impl Default for RetryOptions { + /// Same as [`crate::defaults::default_retry_options`]. Construct with + /// `RetryOptions { attempts: 5, ..Default::default() }` so a future field + /// does not break your code. + fn default() -> Self { + crate::defaults::default_retry_options() + } } impl std::fmt::Debug for RetryOptions { @@ -99,6 +123,7 @@ impl std::fmt::Debug for RetryOptions { .field("max_interval_ms", &self.max_interval_ms) .field("fast_first", &self.fast_first) .field("on_rejection", &self.on_rejection.is_some()) + .field("allow_non_idempotent", &self.allow_non_idempotent) .finish() } } diff --git a/rust/fetch/tests/integration_tests.rs b/rust/fetch/tests/integration_tests.rs index c575a9a..3e1c7d1 100644 --- a/rust/fetch/tests/integration_tests.rs +++ b/rust/fetch/tests/integration_tests.rs @@ -73,6 +73,7 @@ async fn test_retry_with_timeout() { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }) .build(); @@ -112,6 +113,7 @@ async fn test_circuit_breaker_with_retry() { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }) .with_circuit_breaker(2, 1, 5000) .build(); @@ -281,6 +283,7 @@ async fn test_full_pipeline_success() { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }) .with_rate_limit(10, 60_000) .with_circuit_breaker(5, 2, 30_000) diff --git a/rust/fetch/tests/rate_limit_retry_tests.rs b/rust/fetch/tests/rate_limit_retry_tests.rs index a7c6791..546ec7d 100644 --- a/rust/fetch/tests/rate_limit_retry_tests.rs +++ b/rust/fetch/tests/rate_limit_retry_tests.rs @@ -43,6 +43,7 @@ async fn test_rate_limit_retry_recovers() { max_interval_ms: Some(40), fast_first: true, on_rejection: None, + allow_non_idempotent: false, }; // limit=1 over the window. The first request burns the slot. @@ -100,6 +101,7 @@ async fn test_rate_limit_retry_exhausts() { max_interval_ms: Some(10), fast_first: false, on_rejection: None, + allow_non_idempotent: false, }; let client = FetchBuilder::::new() diff --git a/rust/fetch/tests/retry_idempotency_tests.rs b/rust/fetch/tests/retry_idempotency_tests.rs new file mode 100644 index 0000000..4a1b117 --- /dev/null +++ b/rust/fetch/tests/retry_idempotency_tests.rs @@ -0,0 +1,391 @@ +//! Retries must be method-aware, and a timed-out attempt must be cancelled. +//! +//! Regression coverage for SMOODEV-3375: retries were method-blind, so a POST +//! that timed out or got a 429/5xx was re-sent up to twice more โ€” image +//! generation billed three images and returned none, and any non-idempotent +//! POST/PATCH could duplicate its side effect. +//! +//! The cases come from spec/retry-idempotency-corpus.json, shared with the +//! TypeScript, Python, Go and .NET suites. Every case runs against a REAL local +//! server and asserts how many requests the SERVER received: the bug is a side +//! effect executing more than once server-side, and only the server can count +//! that. + +use std::collections::HashMap; +use std::sync::atomic::{AtomicUsize, Ordering}; +use std::sync::{Arc, Mutex}; +use std::time::{Duration, Instant}; + +use serde_json::Value; +use smooai_fetch::client; +use smooai_fetch::error::FetchError; +use smooai_fetch::types::{FetchOptions, Method, RequestInit, RetryOptions, TimeoutOptions}; +use smooai_fetch::{is_idempotent_method, is_retry_eligible}; +use tokio::io::{AsyncReadExt, AsyncWriteExt}; +use tokio::net::{TcpListener, TcpStream}; + +/// `include_str!` binds the corpus at compile time, so a corpus change cannot +/// silently miss this suite. +const CORPUS: &str = include_str!("../../../spec/retry-idempotency-corpus.json"); + +fn corpus() -> Value { + serde_json::from_str(CORPUS).expect("retry-idempotency corpus must parse") +} + +/// Parse a corpus method name case-insensitively. `None` for methods the +/// [`Method`] enum cannot express (TRACE, CONNECT, PURGE): those are skipped +/// here, and the other ports classify them. +fn parse_method(name: &str) -> Option { + match name.to_ascii_uppercase().as_str() { + "GET" => Some(Method::GET), + "HEAD" => Some(Method::HEAD), + "OPTIONS" => Some(Method::OPTIONS), + "PUT" => Some(Method::PUT), + "DELETE" => Some(Method::DELETE), + "POST" => Some(Method::POST), + "PATCH" => Some(Method::PATCH), + _ => None, + } +} + +// --- test server ------------------------------------------------------------- + +#[derive(Clone)] +enum Respond { + Status { + status: u16, + headers: Vec<(String, String)>, + }, + /// Read the request, never answer, and record when the client closes. + Hang, +} + +#[derive(Debug, Clone, Copy)] +struct Connection { + arrived: Instant, + closed: Option, +} + +#[derive(Default)] +struct ServerState { + requests: AtomicUsize, + connections: Mutex>, +} + +/// A minimal HTTP/1.1 server on a raw TcpListener โ€” raw so a hung request can +/// observe the client closing its socket (read returns 0), which no mock +/// server exposes. +async fn spawn_server(respond: Respond) -> (String, Arc) { + let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); + let addr = listener.local_addr().unwrap(); + let state = Arc::new(ServerState::default()); + let accept_state = state.clone(); + tokio::spawn(async move { + loop { + let Ok((stream, _)) = listener.accept().await else { + return; + }; + let index = { + let mut conns = accept_state.connections.lock().unwrap(); + conns.push(Connection { + arrived: Instant::now(), + closed: None, + }); + conns.len() - 1 + }; + tokio::spawn(handle(stream, index, respond.clone(), accept_state.clone())); + } + }); + (format!("http://{addr}/op"), state) +} + +async fn handle(mut stream: TcpStream, index: usize, respond: Respond, state: Arc) { + // Read the head, then the body by Content-Length. + let mut buf = Vec::new(); + let mut chunk = [0u8; 4096]; + let head_end = loop { + if let Some(pos) = buf.windows(4).position(|w| w == b"\r\n\r\n") { + break pos + 4; + } + match stream.read(&mut chunk).await { + Ok(0) | Err(_) => return, + Ok(n) => buf.extend_from_slice(&chunk[..n]), + } + }; + let head = String::from_utf8_lossy(&buf[..head_end]).to_ascii_lowercase(); + let content_length = head + .lines() + .find_map(|l| l.strip_prefix("content-length:")) + .and_then(|v| v.trim().parse::().ok()) + .unwrap_or(0); + while buf.len() < head_end + content_length { + match stream.read(&mut chunk).await { + Ok(0) | Err(_) => return, + Ok(n) => buf.extend_from_slice(&chunk[..n]), + } + } + state.requests.fetch_add(1, Ordering::SeqCst); + + match respond { + Respond::Status { status, headers } => { + let mut response = + format!("HTTP/1.1 {status} Test\r\ncontent-length: 0\r\nconnection: close\r\n"); + for (k, v) in headers { + response.push_str(&format!("{k}: {v}\r\n")); + } + response.push_str("\r\n"); + let _ = stream.write_all(response.as_bytes()).await; + let _ = stream.shutdown().await; + } + Respond::Hang => { + // Never answer. A client that cancels closes the socket (EOF or + // reset); one that abandons the attempt leaves it open. + loop { + match stream.read(&mut chunk).await { + Ok(0) | Err(_) => break, + Ok(_) => continue, + } + } + state.connections.lock().unwrap()[index].closed = Some(Instant::now()); + } + } +} + +// --- tests ------------------------------------------------------------------- + +/// Positive control: a corpus that failed to load (or lost its cases) would +/// make every loop below vacuously pass. +#[test] +fn corpus_loaded() { + let c = corpus(); + assert!(c["cases"].as_array().is_some_and(|a| a.len() >= 10)); + assert!(c["methods"]["idempotent"] + .as_array() + .is_some_and(|a| !a.is_empty())); + assert!(c["methods"]["nonIdempotent"] + .as_array() + .is_some_and(|a| !a.is_empty())); + assert!(c["timeoutAbort"]["timeoutMs"].as_u64().is_some()); +} + +#[test] +fn classifies_corpus_methods() { + let c = corpus(); + let mut checked = 0; + for (group, expected) in [("idempotent", true), ("nonIdempotent", false)] { + for name in c["methods"][group].as_array().unwrap() { + let name = name.as_str().unwrap(); + // TRACE, CONNECT and PURGE are not representable by `Method`; the + // other ports classify them. + let Some(method) = parse_method(name) else { + continue; + }; + assert_eq!( + is_idempotent_method(&method), + expected, + "{name} should be idempotent={expected}" + ); + // A bare request of that method is retry-eligible iff idempotent. + let init = RequestInit { + method, + ..Default::default() + }; + assert_eq!(is_retry_eligible(&init, &RetryOptions::default()), expected); + checked += 1; + } + } + assert!(checked >= 10, "only {checked} methods were classified"); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn corpus_cases_hit_the_server_the_expected_number_of_times() { + let c = corpus(); + let retry = &c["retry"]; + let timeout_ms = c["timeoutMs"].as_u64().unwrap(); + let cases = c["cases"].as_array().unwrap(); + let mut failures = Vec::new(); + + for case in cases { + let name = case["name"].as_str().unwrap(); + let respond = if case["respond"]["hang"].as_bool() == Some(true) { + Respond::Hang + } else { + let headers = case["respond"]["headers"] + .as_object() + .map(|h| { + h.iter() + .map(|(k, v)| (k.clone(), v.as_str().unwrap().to_string())) + .collect() + }) + .unwrap_or_default(); + Respond::Status { + status: case["respond"]["status"].as_u64().unwrap() as u16, + headers, + } + }; + let (url, state) = spawn_server(respond).await; + + let headers: HashMap = case["headers"] + .as_object() + .map(|h| { + h.iter() + .map(|(k, v)| (k.clone(), v.as_str().unwrap().to_string())) + .collect() + }) + .unwrap_or_default(); + let init = RequestInit { + method: parse_method(case["method"].as_str().unwrap()).expect("server method"), + headers, + body: case["body"].as_str().map(str::to_string), + }; + let options = FetchOptions { + connect_timeout_ms: None, + timeout: Some(TimeoutOptions { timeout_ms }), + retry: Some(RetryOptions { + attempts: retry["attempts"].as_u64().unwrap() as u32, + initial_interval_ms: retry["initialIntervalMs"].as_u64().unwrap(), + factor: retry["factor"].as_f64().unwrap(), + jitter_adjustment: retry["jitterAdjustment"].as_f64().unwrap(), + allow_non_idempotent: case["allowNonIdempotent"].as_bool().unwrap_or(false), + ..Default::default() + }), + }; + + let result = + client::fetch::(&url, init, Some(options), None, None, None, None).await; + let expected = case["expectedAttempts"].as_u64().unwrap() as usize; + let got = state.requests.load(Ordering::SeqCst); + + let mut problems = Vec::new(); + if got != expected { + problems.push(format!( + "server received {got} request(s), expected {expected}" + )); + } + match &result { + Err(FetchError::Retry { .. }) if expected == 1 => problems.push(format!( + "a request that was not retried must surface the underlying error, got {result:?}" + )), + Err(_) => {} + Ok(r) => problems.push(format!("expected an error, got status {}", r.status)), + } + if !problems.is_empty() { + failures.push(format!("{name}: {}", problems.join("; "))); + } + } + + assert!( + failures.is_empty(), + "{} of {} corpus cases failed:\n {}", + failures.len(), + cases.len(), + failures.join("\n ") + ); +} + +/// Poll until connection `index` has closed, or `deadline` passes. +async fn wait_for_close( + state: &ServerState, + index: usize, + deadline: Instant, +) -> Option { + loop { + let conn = state.connections.lock().unwrap().get(index).copied(); + if let Some(conn) = conn { + if conn.closed.is_some() { + return Some(conn); + } + } + if Instant::now() >= deadline { + return conn; + } + tokio::time::sleep(Duration::from_millis(10)).await; + } +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn timed_out_attempt_closes_its_connection_before_the_retry() { + let c = corpus(); + let timeout_ms = c["timeoutAbort"]["timeoutMs"].as_u64().unwrap(); + let slack_ms = c["timeoutAbort"]["closeSlackMs"].as_u64().unwrap(); + let (url, state) = spawn_server(Respond::Hang).await; + + let options = FetchOptions { + connect_timeout_ms: None, + timeout: Some(TimeoutOptions { timeout_ms }), + retry: Some(RetryOptions { + attempts: 1, + // Wider than the corpus cases' 1ms so "closed before the retry + // arrived" is not decided by which server task the scheduler wakes + // first. An abandoned attempt stays open for the whole test, so + // this cannot mask a missing cancel. + initial_interval_ms: 100, + factor: 1.0, + jitter_adjustment: 0.0, + ..Default::default() + }), + }; + let init = RequestInit { + method: Method::GET, + ..Default::default() + }; + let result = client::fetch::(&url, init, Some(options), None, None, None, None).await; + assert!( + matches!(result, Err(FetchError::Retry { .. })), + "expected the GET to time out twice, got {result:?}" + ); + + let first_arrived = state.connections.lock().unwrap()[0].arrived; + let deadline = first_arrived + Duration::from_millis(timeout_ms + slack_ms); + let first = wait_for_close(&state, 0, deadline).await.unwrap(); + let closed = first + .closed + .expect("the timed-out attempt's connection was never closed โ€” the request was abandoned, not cancelled"); + assert!( + closed - first.arrived <= Duration::from_millis(timeout_ms + slack_ms), + "first connection closed {:?} after arriving", + closed - first.arrived + ); + + let conns = state.connections.lock().unwrap().clone(); + assert_eq!(conns.len(), 2, "a retried GET makes two attempts"); + assert!( + closed <= conns[1].arrived, + "first attempt must be cancelled before the retry is sent" + ); +} + +#[tokio::test(flavor = "multi_thread", worker_threads = 2)] +async fn timed_out_post_closes_its_connection_and_is_not_retried() { + let c = corpus(); + let timeout_ms = c["timeoutAbort"]["timeoutMs"].as_u64().unwrap(); + let slack_ms = c["timeoutAbort"]["closeSlackMs"].as_u64().unwrap(); + let (url, state) = spawn_server(Respond::Hang).await; + + let options = FetchOptions { + connect_timeout_ms: None, + timeout: Some(TimeoutOptions { timeout_ms }), + retry: Some(RetryOptions::default()), + }; + let init = RequestInit { + method: Method::POST, + body: Some("{}".to_string()), + ..Default::default() + }; + let result = client::fetch::(&url, init, Some(options), None, None, None, None).await; + assert!( + matches!(result, Err(FetchError::Timeout { .. })), + "a POST that times out surfaces the timeout unwrapped, got {result:?}" + ); + + let first_arrived = state.connections.lock().unwrap()[0].arrived; + let deadline = first_arrived + Duration::from_millis(timeout_ms + slack_ms); + let first = wait_for_close(&state, 0, deadline).await.unwrap(); + assert!( + first.closed.is_some(), + "the timed-out POST's connection was never closed" + ); + // Give a (wrong) retry the chance to arrive before counting. + tokio::time::sleep(Duration::from_millis(1_500)).await; + assert_eq!(state.requests.load(Ordering::SeqCst), 1); +} diff --git a/rust/fetch/tests/retry_tests.rs b/rust/fetch/tests/retry_tests.rs index 7a0f04c..eed5379 100644 --- a/rust/fetch/tests/retry_tests.rs +++ b/rust/fetch/tests/retry_tests.rs @@ -39,6 +39,7 @@ fn test_calculate_backoff_no_jitter() { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }; assert_eq!(calculate_backoff(0, &options), 100); // 100 * 2^0 @@ -57,6 +58,7 @@ fn test_calculate_backoff_with_max_interval() { max_interval_ms: Some(300), fast_first: false, on_rejection: None, + allow_non_idempotent: false, }; assert_eq!(calculate_backoff(0, &options), 100); @@ -75,6 +77,7 @@ fn test_calculate_backoff_with_jitter_is_bounded() { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }; for _ in 0..100 { @@ -141,6 +144,7 @@ async fn test_retry_succeeds_after_failures() { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }), }; @@ -196,6 +200,7 @@ async fn test_retry_exhausted() { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }), }; @@ -253,6 +258,7 @@ async fn test_non_retryable_error_not_retried() { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }), }; @@ -304,6 +310,7 @@ async fn test_retry_with_retry_after_header() { max_interval_ms: None, fast_first: false, on_rejection: None, + allow_non_idempotent: false, }), }; @@ -367,6 +374,7 @@ async fn test_fast_first_skips_initial_delay() { max_interval_ms: None, fast_first: true, on_rejection: None, + allow_non_idempotent: false, }), }; @@ -443,6 +451,7 @@ async fn test_on_rejection_retry_decision_overrides_delay() { max_interval_ms: None, fast_first: false, on_rejection: Some(callback), + allow_non_idempotent: false, }), }; @@ -502,6 +511,7 @@ async fn test_on_rejection_abort_stops_retry_loop() { max_interval_ms: None, fast_first: false, on_rejection: Some(callback), + allow_non_idempotent: false, }), }; @@ -552,6 +562,7 @@ async fn test_on_rejection_default_falls_through_to_exponential() { max_interval_ms: None, fast_first: false, on_rejection: Some(callback), + allow_non_idempotent: false, }), }; @@ -603,6 +614,7 @@ async fn test_on_rejection_skip_consumes_attempt_without_sleep() { max_interval_ms: None, fast_first: false, on_rejection: Some(callback), + allow_non_idempotent: false, }), }; diff --git a/spec/retry-idempotency-corpus.json b/spec/retry-idempotency-corpus.json new file mode 100644 index 0000000..7ebeb4d --- /dev/null +++ b/spec/retry-idempotency-corpus.json @@ -0,0 +1,115 @@ +{ + "$comment": "Shared retry-safety corpus, loaded by all five ports โ€” do not hand-copy these cases into a language's test file. Each port runs every `cases` entry against a REAL local HTTP server and asserts the number of requests the SERVER received, because the bug this guards against is a side effect executing more than once server-side, and only the server can count that.", + "$regression": "SMOODEV-3375. Retries were method-blind: a POST that timed out or got a 429/5xx was re-sent up to twice more. Image generation (20-50s) against a 10s per-attempt timeout billed three images and returned nothing, and any non-idempotent POST/PATCH (CRM writes, message sends, payments) could duplicate its side effect. On top of that the TypeScript timeout did not abort the in-flight request, so the server never saw the first attempt cancelled.", + "rule": "A failed attempt is retried only when the request is retry-eligible: its method is idempotent per RFC 9110 section 9.2.2, OR the caller set allowNonIdempotent on the retry options, OR the request carries an Idempotency-Key header (any casing) whose value is not empty or whitespace. An ineligible request makes exactly one attempt and surfaces the underlying error, not the retries-exhausted wrapper. The in-process rate limiter's own retry loop is unaffected: it rejects before anything is sent.", + "idempotencyKeyHeader": "Idempotency-Key", + "methods": { + "$comment": "Classification only. Compared case-insensitively. TRACE is idempotent but is not exercised against a server (fetch implementations forbid it); CONNECT and anything unrecognised are treated as non-idempotent. Ports whose method type cannot express a method here skip it.", + "idempotent": ["GET", "HEAD", "OPTIONS", "TRACE", "PUT", "DELETE", "get", "Put"], + "nonIdempotent": ["POST", "PATCH", "CONNECT", "post", "Patch", "PURGE"] + }, + "retry": { + "$comment": "Retry knobs every case runs with. Tiny intervals and no jitter keep the suite fast and deterministic; `attempts` is the number of RETRIES, so an eligible request makes attempts + 1 = 3 calls.", + "attempts": 2, + "initialIntervalMs": 1, + "factor": 1, + "jitterAdjustment": 0 + }, + "timeoutMs": 300, + "cases": [ + { "name": "GET 503 retries", "method": "GET", "respond": { "status": 503 }, "expectedAttempts": 3 }, + { "name": "HEAD 503 retries", "method": "HEAD", "respond": { "status": 503 }, "expectedAttempts": 3 }, + { "name": "OPTIONS 503 retries", "method": "OPTIONS", "respond": { "status": 503 }, "expectedAttempts": 3 }, + { "name": "PUT 503 retries", "method": "PUT", "body": "{\"a\":1}", "respond": { "status": 503 }, "expectedAttempts": 3 }, + { "name": "DELETE 503 retries", "method": "DELETE", "respond": { "status": 503 }, "expectedAttempts": 3 }, + { "name": "POST 503 is not retried by default", "method": "POST", "body": "{\"a\":1}", "respond": { "status": 503 }, "expectedAttempts": 1 }, + { "name": "PATCH 503 is not retried by default", "method": "PATCH", "body": "{\"a\":1}", "respond": { "status": 503 }, "expectedAttempts": 1 }, + { + "name": "POST 503 retries when allowNonIdempotent is set", + "method": "POST", + "body": "{\"a\":1}", + "allowNonIdempotent": true, + "respond": { "status": 503 }, + "expectedAttempts": 3 + }, + { + "name": "PATCH 503 retries when allowNonIdempotent is set", + "method": "PATCH", + "body": "{\"a\":1}", + "allowNonIdempotent": true, + "respond": { "status": 503 }, + "expectedAttempts": 3 + }, + { + "name": "POST 503 retries when it carries an Idempotency-Key", + "method": "POST", + "body": "{\"a\":1}", + "headers": { "Idempotency-Key": "k-1" }, + "respond": { "status": 503 }, + "expectedAttempts": 3 + }, + { + "name": "Idempotency-Key is matched case-insensitively", + "method": "POST", + "body": "{\"a\":1}", + "headers": { "idempotency-key": "k-2" }, + "respond": { "status": 503 }, + "expectedAttempts": 3 + }, + { + "name": "an empty Idempotency-Key does not opt in", + "method": "POST", + "body": "{\"a\":1}", + "headers": { "Idempotency-Key": "" }, + "respond": { "status": 503 }, + "expectedAttempts": 1 + }, + { + "name": "a whitespace-only Idempotency-Key does not opt in", + "method": "POST", + "body": "{\"a\":1}", + "headers": { "Idempotency-Key": " " }, + "respond": { "status": 503 }, + "expectedAttempts": 1 + }, + { + "name": "POST 429 with Retry-After is still not retried by default", + "method": "POST", + "body": "{\"a\":1}", + "respond": { "status": 429, "headers": { "Retry-After": "0" } }, + "expectedAttempts": 1 + }, + { + "name": "POST 429 with Retry-After retries when opted in", + "method": "POST", + "body": "{\"a\":1}", + "allowNonIdempotent": true, + "respond": { "status": 429, "headers": { "Retry-After": "0" } }, + "expectedAttempts": 3 + }, + { "name": "GET 429 with Retry-After retries", "method": "GET", "respond": { "status": 429, "headers": { "Retry-After": "0" } }, "expectedAttempts": 3 }, + { "name": "GET that times out is retried", "method": "GET", "respond": { "hang": true }, "expectedAttempts": 3 }, + { "name": "POST that times out is not retried", "method": "POST", "body": "{\"a\":1}", "respond": { "hang": true }, "expectedAttempts": 1 }, + { + "name": "POST that times out retries when opted in", + "method": "POST", + "body": "{\"a\":1}", + "allowNonIdempotent": true, + "respond": { "hang": true }, + "expectedAttempts": 3 + }, + { + "name": "POST 400 is never retried", + "method": "POST", + "body": "{\"a\":1}", + "allowNonIdempotent": true, + "respond": { "status": 400 }, + "expectedAttempts": 1 + } + ], + "timeoutAbort": { + "$comment": "A timed-out attempt must be CANCELLED, not abandoned: the connection is closed, so the server can see it and stop work. The test server accepts the request and never answers; it asserts the first attempt's connection closes within timeoutMs + closeSlackMs of the request arriving, and โ€” for a retried GET โ€” that it closed before the second attempt arrived. Without cancellation the socket stays open until the process exits.", + "timeoutMs": 300, + "closeSlackMs": 1500 + } +} diff --git a/src/fetch.retry-safety.spec.ts b/src/fetch.retry-safety.spec.ts new file mode 100644 index 0000000..591fe5f --- /dev/null +++ b/src/fetch.retry-safety.spec.ts @@ -0,0 +1,207 @@ +import { readFileSync } from 'node:fs'; +import { createServer, IncomingMessage, Server, ServerResponse } from 'node:http'; +import { AddressInfo } from 'node:net'; +import { dirname, join } from 'node:path'; +import { fileURLToPath } from 'node:url'; +import { TimeoutError } from 'mollitia'; +import { afterEach, describe, expect, test } from 'vitest'; +import fetch, { FetchBuilder, HTTPResponseError, isIdempotentMethod, RetryError } from './fetch'; + +/** + * Retries must never duplicate a side effect (SMOODEV-3375). + * + * Every case comes from spec/retry-idempotency-corpus.json, shared with the + * Python, Rust, Go and .NET suites, and runs against a REAL local server that + * counts what it received โ€” a mocked fetch can only count what the client + * attempted, and the bug is what the server executed. + */ +type CorpusCase = { + name: string; + method: string; + body?: string; + headers?: Record; + allowNonIdempotent?: boolean; + respond: { status?: number; headers?: Record; hang?: boolean }; + expectedAttempts: number; +}; + +const corpus = JSON.parse(readFileSync(join(dirname(fileURLToPath(import.meta.url)), '..', 'spec', 'retry-idempotency-corpus.json'), 'utf8')) as { + methods: { idempotent: string[]; nonIdempotent: string[] }; + retry: { attempts: number; initialIntervalMs: number; factor: number; jitterAdjustment: number }; + timeoutMs: number; + cases: CorpusCase[]; + timeoutAbort: { timeoutMs: number; closeSlackMs: number }; +}; + +type Arrival = { method: string; arrivedAt: number; closedAt?: number }; + +const servers: Server[] = []; +const hanging: ServerResponse[] = []; + +afterEach(async () => { + for (const res of hanging.splice(0)) res.destroy(); + await Promise.all(servers.splice(0).map((server) => new Promise((resolve) => server.close(resolve)))); +}); + +/** A server that answers every request with `respond`, recording each arrival and when its connection closed. */ +async function startServer(respond: CorpusCase['respond']): Promise<{ url: string; arrivals: Arrival[] }> { + const arrivals: Arrival[] = []; + const server = createServer((req: IncomingMessage, res: ServerResponse) => { + const arrival: Arrival = { method: req.method ?? '', arrivedAt: Date.now() }; + arrivals.push(arrival); + req.socket.once('close', () => { + arrival.closedAt = Date.now(); + }); + req.resume(); + if (respond.hang) { + hanging.push(res); + return; + } + res.writeHead(respond.status ?? 200, respond.headers ?? {}); + res.end(); + }); + servers.push(server); + await new Promise((resolve) => server.listen(0, '127.0.0.1', resolve)); + return { url: `http://127.0.0.1:${(server.address() as AddressInfo).port}/`, arrivals }; +} + +const retryOptions = (allowNonIdempotent?: boolean) => ({ + attempts: corpus.retry.attempts, + initialIntervalMs: corpus.retry.initialIntervalMs, + factor: corpus.retry.factor, + jitterAdjustment: corpus.retry.jitterAdjustment, + ...(allowNonIdempotent ? { allowNonIdempotent } : {}), +}); + +describe('retry idempotency corpus', () => { + // Positive control: a corpus that failed to load would leave the case table + // empty and every per-case assertion vacuously satisfied. + test('loaded the shared corpus', () => { + expect(corpus.cases.length).toBeGreaterThan(10); + expect(corpus.methods.idempotent).toContain('GET'); + expect(corpus.methods.nonIdempotent).toContain('POST'); + }); + + test.each(corpus.methods.idempotent)('%s is idempotent', (method) => { + expect(isIdempotentMethod(method)).toBe(true); + }); + + test.each(corpus.methods.nonIdempotent)('%s is not idempotent', (method) => { + expect(isIdempotentMethod(method)).toBe(false); + }); + + test.each(corpus.cases.map((c) => [c.name, c] as const))('%s', async (_name, c) => { + const { url, arrivals } = await startServer(c.respond); + + const promise = fetch(url, { + method: c.method, + headers: c.headers, + body: c.body, + options: { retry: retryOptions(c.allowNonIdempotent), timeout: { timeoutMs: corpus.timeoutMs } }, + }); + + const error = await promise.then( + () => undefined, + (e: unknown) => e, + ); + expect(error).toBeDefined(); + expect(arrivals.length).toBe(c.expectedAttempts); + // An unretried request surfaces its own error, never "ran out of retries". + if (c.expectedAttempts === 1) { + expect(error).not.toBeInstanceOf(RetryError); + expect(error instanceof HTTPResponseError || error instanceof TimeoutError).toBe(true); + } + }); + + test('a pre-request hook that adds an Idempotency-Key makes a POST retry-eligible', async () => { + const { url, arrivals } = await startServer({ status: 503 }); + const client = new FetchBuilder() + .withRetry(retryOptions()) + .withHooks({ + preRequest: (u, init) => [u, { ...init, headers: { ...(init.headers as Record), 'Idempotency-Key': 'from-hook' } }], + }) + .build(); + + await expect(client(url, { method: 'POST', body: '{}' })).rejects.toBeInstanceOf(RetryError); + expect(arrivals.length).toBe(corpus.retry.attempts + 1); + }); + + test('FetchBuilder.withRetry({ allowNonIdempotent }) keeps the other defaults', async () => { + const { url, arrivals } = await startServer({ status: 503 }); + const client = new FetchBuilder().withRetry({ allowNonIdempotent: true, initialIntervalMs: 1 }).build(); + + await expect(client(url, { method: 'POST', body: '{}' })).rejects.toBeInstanceOf(RetryError); + // DEFAULT_RETRY_OPTIONS.attempts (2) survived the partial merge. + expect(arrivals.length).toBe(3); + }); + + test('an ineligible POST never consults onRejection', async () => { + const { url, arrivals } = await startServer({ status: 503 }); + let consulted = 0; + await expect( + fetch(url, { + method: 'POST', + options: { + retry: { + ...retryOptions(), + onRejection: () => { + consulted++; + return true; + }, + }, + }, + }), + ).rejects.toBeInstanceOf(HTTPResponseError); + expect(arrivals.length).toBe(1); + expect(consulted).toBe(0); + }); + + test('a request the caller aborted is not retried', async () => { + const { url, arrivals } = await startServer({ hang: true }); + const controller = new AbortController(); + const promise = fetch(url, { method: 'GET', signal: controller.signal, options: { retry: retryOptions(), timeout: { timeoutMs: 5_000 } } }); + await expect.poll(() => arrivals.length).toBe(1); + controller.abort(); + + await expect(promise).rejects.toThrow(); + await new Promise((resolve) => setTimeout(resolve, 100)); + expect(arrivals.length).toBe(1); + }); +}); + +describe('timeouts abort the attempt', () => { + const { timeoutMs, closeSlackMs } = corpus.timeoutAbort; + + test('a timed-out GET closes its connection before the retry is sent', async () => { + const { url, arrivals } = await startServer({ hang: true }); + + // A 100ms backoff instead of the corpus's 1ms: with 1ms, whether the + // server observes "first socket closed" or "second request arrived" first + // is down to event-loop scheduling. An attempt that is abandoned rather + // than cancelled stays open for the whole test, so the wider gap cannot + // hide a missing abort. + await expect( + fetch(url, { + method: 'GET', + options: { retry: { attempts: 1, initialIntervalMs: 100, jitterAdjustment: 0 }, timeout: { timeoutMs } }, + }), + ).rejects.toBeInstanceOf(TimeoutError); + + expect(arrivals.length).toBe(2); + const [first, second] = arrivals; + expect(first.closedAt, "the timed-out attempt's connection was never closed โ€” it was abandoned, not cancelled").toBeDefined(); + expect(first.closedAt! - first.arrivedAt).toBeLessThanOrEqual(timeoutMs + closeSlackMs); + expect(first.closedAt!).toBeLessThanOrEqual(second.arrivedAt); + }); + + test('a timed-out POST closes its connection and is sent exactly once', async () => { + const { url, arrivals } = await startServer({ hang: true }); + + await expect(fetch(url, { method: 'POST', body: '{}', options: { timeout: { timeoutMs } } })).rejects.toBeInstanceOf(TimeoutError); + + await expect.poll(() => arrivals[0]?.closedAt, { timeout: closeSlackMs }).toBeDefined(); + expect(arrivals[0].closedAt! - arrivals[0].arrivedAt).toBeLessThanOrEqual(timeoutMs + closeSlackMs); + await new Promise((resolve) => setTimeout(resolve, 200)); + expect(arrivals.length).toBe(1); + }); +}); diff --git a/src/fetch.spec.ts b/src/fetch.spec.ts index 6700f82..2b1f691 100644 --- a/src/fetch.spec.ts +++ b/src/fetch.spec.ts @@ -869,7 +869,14 @@ describe('Test fetch', () => { expect(response.ok).toBeTruthy(); expect(response.status).toBe(200); - expect(mockFetch.mock.calls[0][1]?.signal).toBe(controller.signal); + // The caller's signal is combined with the per-attempt timeout signal, + // so it is not passed through by identity โ€” but aborting it must + // still abort what fetch received. + const passedSignal = mockFetch.mock.calls[0][1]?.signal; + expect(passedSignal).toBeInstanceOf(AbortSignal); + expect(passedSignal!.aborted).toBe(false); + controller.abort(); + expect(passedSignal!.aborted).toBe(true); }); test('Test fetch with multiple init options combined', async () => { @@ -905,7 +912,9 @@ describe('Test fetch', () => { expect(init?.mode).toBe('cors'); expect(init?.redirect).toBe('follow'); expect(init?.referrer).toBe('https://example.com'); - expect(init?.signal).toBe(controller.signal); + expect(init?.signal).toBeInstanceOf(AbortSignal); + controller.abort(); + expect(init?.signal?.aborted).toBe(true); }); }); diff --git a/src/fetch.ts b/src/fetch.ts index 0db3005..edc9903 100644 --- a/src/fetch.ts +++ b/src/fetch.ts @@ -1,6 +1,6 @@ import type { StandardSchemaV1 } from '@standard-schema/spec'; import merge from 'lodash.merge'; -import { BreakerError, BreakerState, Circuit, Module, Ratelimit, RatelimitError, Retry, RetryMode, SlidingCountBreaker, Timeout, TimeoutError } from 'mollitia'; +import { BreakerError, BreakerState, Circuit, Module, Ratelimit, RatelimitError, Retry, RetryMode, SlidingCountBreaker, TimeoutError } from 'mollitia'; import { CONTEXT, ContextHeader, ContextKey, ContextKeyHttp, ContextKeyHttpRequest, ContextKeyHttpResponse } from '@smooai/logger/Logger'; import { handleSchemaValidation, HumanReadableSchemaError } from '@smooai/utils/validation/standardSchema'; import { contextLogger } from './logger'; @@ -229,7 +229,7 @@ export type RetryCallback = (err: any, attempt: number) => boolean | number; /** * Configuration options for retry behavior. */ -interface RetryOptions { +export interface RetryOptions { /** Number of retry attempts */ attempts: number; /** Initial delay between retries in milliseconds */ @@ -246,6 +246,56 @@ interface RetryOptions { jitterAdjustment?: number; /** Callback to determine if and when to retry */ onRejection?: RetryCallback; + /** + * Retry non-idempotent methods (POST, PATCH, ...) too. Off by default. + * + * Only idempotent methods (RFC 9110 ยง9.2.2: GET, HEAD, OPTIONS, TRACE, PUT, + * DELETE) are retried unless this is set or the request carries an + * `Idempotency-Key` header. A POST that timed out or got a 429/5xx may + * already have done its work server-side โ€” re-sending it bills a second + * image, sends a second message, charges a card twice. Set this only when + * the endpoint tolerates a duplicate, and prefer an `Idempotency-Key`. + * + * This overrides `onRejection`: an ineligible request makes exactly one + * attempt and `onRejection` is never consulted. + */ + allowNonIdempotent?: boolean; +} + +/** Header that, when set to a non-empty value, makes a non-idempotent request retry-eligible. */ +export const IDEMPOTENCY_KEY_HEADER = 'Idempotency-Key'; + +const IDEMPOTENT_METHODS = new Set(['GET', 'HEAD', 'OPTIONS', 'TRACE', 'PUT', 'DELETE']); + +/** + * Whether an HTTP method is idempotent per RFC 9110 ยง9.2.2 โ€” i.e. sending it + * twice has the same effect on the server as sending it once. Case-insensitive; + * anything unrecognised is treated as non-idempotent. + */ +export function isIdempotentMethod(method: string | undefined): boolean { + return IDEMPOTENT_METHODS.has((method ?? 'GET').toUpperCase()); +} + +function hasIdempotencyKey(headers: RequestInit['headers']): boolean { + if (!headers) return false; + let value: string | null = null; + if (headers instanceof Headers) { + value = headers.get(IDEMPOTENCY_KEY_HEADER); + } else { + const entries = Array.isArray(headers) ? headers : Object.entries(headers); + const found = entries.find(([key]) => key.toLowerCase() === IDEMPOTENCY_KEY_HEADER.toLowerCase()); + value = found ? String(found[1]) : null; + } + return !!value && value.trim().length > 0; +} + +/** + * Whether a failed attempt of this request may be re-sent: its method is + * idempotent, the caller set `retry.allowNonIdempotent`, or it carries a + * non-empty `Idempotency-Key` header (the server then dedupes for us). + */ +export function isRetryEligible(init: Pick, retry?: Partial): boolean { + return isIdempotentMethod(init.method) || !!retry?.allowNonIdempotent || hasIdempotencyKey(init.headers); } export const DEFAULT_RETRY_OPTIONS: RetryOptions = { @@ -327,7 +377,10 @@ export interface LifecycleHooks { export interface RequestOptions { /** Custom logger for request logging */ logger?: LoggerInterface; - /** Timeout configuration */ + /** + * Per-attempt timeout. When it fires the attempt is ABORTED (its + * connection is closed, so the server sees the cancel) before any retry. + */ timeout?: { /** Timeout duration in milliseconds */ timeoutMs: number; @@ -343,8 +396,13 @@ export interface RequestOptions { * Ignored in browser/worker environments (no connect-timeout knob there). */ connectTimeoutMs?: number; - /** Retry configuration */ - retry?: RetryOptions; + /** + * Retry configuration, deep-merged over {@link DEFAULT_RETRY_OPTIONS}, so + * `retry: { allowNonIdempotent: true }` keeps every other default. + * POST/PATCH are not retried unless `allowNonIdempotent` is set or the + * request carries an `Idempotency-Key` โ€” see {@link RetryOptions.allowNonIdempotent}. + */ + retry?: Partial; /** Schema for response validation. Must be a StandardSchemaV1 compatible schema (e.g., Zod schema) */ schema?: Schema; /** Lifecycle hooks for request/response handling */ @@ -466,11 +524,20 @@ function generateRandomName(prefix: string): string { return `${prefix}-${++moduleSequence}`; } -function prepareCircuitModules(options: RequestOptions): Module[] { +function prepareCircuitModules(options: RequestOptions, init: RequestInit): Module[] { const modules: Module[] = []; const logger = options.logger || contextLoggerToUse; + const retryEligible = isRetryEligible(init, options.retry); + const onRejection = options.retry?.onRejection; - if (options.retry) { + // No Retry module at all for an ineligible request: one attempt, and the + // underlying error surfaces as-is rather than as a RetryError. + // + // There is deliberately no mollitia `Timeout` module here either. It only + // RACED the attempt โ€” the losing fetch kept running, so a timed-out POST + // stayed in flight server-side while the retry sent it again. The timeout + // now lives in `doGlobalFetch`, where it can abort the request itself. + if (options.retry && retryEligible) { modules.push( new Retry({ name: generateRandomName('smooai-fetch-retry'), @@ -480,17 +547,8 @@ function prepareCircuitModules(options: mode: options.retry.mode, factor: options.retry.factor, jitterAdjustment: options.retry.jitterAdjustment, - onRejection: options.retry.onRejection, - }), - ); - } - - if (options.timeout) { - modules.push( - new Timeout({ - name: generateRandomName('smooai-fetch-timeout'), - logger: logger, - delay: options.timeout.timeoutMs, + // A request the CALLER aborted is finished, not failed โ€” never re-send it. + onRejection: (error, attempt) => (init.signal?.aborted ? false : onRejection ? onRejection(error, attempt) : true), }), ); } @@ -645,6 +703,60 @@ async function doGlobalFetch( } } + // Per-attempt timeout that ABORTS the request. The caller's own signal still + // works: either one aborting cancels the fetch. + const timeoutMs = options?.timeout?.timeoutMs; + const timeoutController = new AbortController(); + let timer: ReturnType | undefined; + if (timeoutMs !== undefined) { + timer = setTimeout(() => timeoutController.abort(new TimeoutError()), timeoutMs); + useInit.signal = init?.signal ? anySignal([init.signal, timeoutController.signal]) : timeoutController.signal; + } + + try { + const attempt = sendAndRead(url, useInit, options); + if (timeoutMs === undefined) return await attempt; + // Race as well as abort: a fetch implementation that ignores `signal` (a + // polyfill, a test double) must still time out on schedule. With a + // compliant fetch the abort has already torn the request down. + return await Promise.race([ + attempt, + new Promise((_, reject) => { + timeoutController.signal.addEventListener('abort', () => reject(timeoutController.signal.reason), { once: true }); + }), + ]); + } catch (error) { + // undici rejects with the abort reason, a browser with an AbortError โ€” + // normalise both to the TimeoutError the retry policy and callers expect. + if (timeoutController.signal.aborted && !init?.signal?.aborted) { + throw timeoutController.signal.reason instanceof TimeoutError ? timeoutController.signal.reason : new TimeoutError(); + } + throw error; + } finally { + clearTimeout(timer); + } +} + +/** `AbortSignal.any` where available (Node 20+, modern browsers), else a manual fan-in. */ +function anySignal(signals: AbortSignal[]): AbortSignal { + const any = (AbortSignal as unknown as { any?: (signals: AbortSignal[]) => AbortSignal }).any; + if (any) return any(signals); + const controller = new AbortController(); + for (const signal of signals) { + if (signal.aborted) { + controller.abort(signal.reason); + break; + } + signal.addEventListener('abort', () => controller.abort(signal.reason), { once: true }); + } + return controller.signal; +} + +async function sendAndRead( + url: RequestInfo, + useInit: RequestInit, + options?: RequestOptions, +): Promise>> { const response = await globalFetch()(url, useInit); let isJson = false; let data: ResponseType | undefined; @@ -702,13 +814,6 @@ async function doFetch( init: RequestInit, options: RequestOptions, ): Promise>> { - const circuit = new Circuit({ - name: 'node-fetch-circuit', - func: doGlobalFetch, - options: { - modules: prepareCircuitModules(options), - }, - }); const logger = options.logger || contextLoggerToUse; // Apply pre-request hook if present (supports async hooks) @@ -721,6 +826,17 @@ async function doFetch( } } + // Built AFTER the pre-request hook: retry eligibility depends on the final + // method and headers, and a hook may be what adds the Idempotency-Key. + const retryEligible = isRetryEligible(modifiedInit, options.retry); + const circuit = new Circuit({ + name: 'node-fetch-circuit', + func: doGlobalFetch, + options: { + modules: prepareCircuitModules(options, modifiedInit), + }, + }); + const urlObj = new URL(url.toString()); const safeUrl = redactUrl(url.toString()); @@ -777,7 +893,7 @@ async function doFetch( }, }, }); - } else if (options.retry && error instanceof HTTPResponseError) { + } else if (options.retry && retryEligible && error instanceof HTTPResponseError) { if (options.retry.onRejection && options.retry.onRejection(error, 1)) { logger.error(error, `HTTP request "${modifiedInit.method} ${safeUrl}" retries failed after ${options.retry.attempts} retries`, { [ContextKey.Http]: { @@ -972,11 +1088,13 @@ export class FetchBuilder { /** * Configures retry behavior for failed requests. - * If not specified, uses DEFAULT_RETRY_OPTIONS. + * If not specified, uses DEFAULT_RETRY_OPTIONS; a partial is merged over it. + * POST/PATCH are only retried with `allowNonIdempotent: true` or an + * `Idempotency-Key` header. * @param options - Retry configuration options * @returns The builder instance for method chaining */ - withRetry(options: RetryOptions = DEFAULT_RETRY_OPTIONS): FetchBuilder { + withRetry(options: Partial = DEFAULT_RETRY_OPTIONS): FetchBuilder { this._requestOptions = { ...this._requestOptions, retry: options,