diff --git a/docs/integrations/weatherapi.md b/docs/integrations/weatherapi.md index 9da4faa..ef12b3c5 100644 --- a/docs/integrations/weatherapi.md +++ b/docs/integrations/weatherapi.md @@ -138,9 +138,10 @@ It does not retry other HTTP statuses, malformed envelopes, missing data, or payload decoding failures. A canceled context also stops an in-progress retry delay. -The adapter reads at most 10 MiB from one response body. A non-2xx response, -request construction failure, read failure, or decode failure includes endpoint -context in its error. +The adapter accepts response bodies up to 10 MiB and rejects larger bodies +before decoding. A non-2xx response reports its relative endpoint and status, +without including upstream response text. Request construction, response-limit, +read, and decode failures include endpoint context in their errors. Retry counts and delays are adapter behavior rather than Weather API request parameters. Do not depend on a particular attempt count when implementing the diff --git a/docs/internal/collect.md b/docs/internal/collect.md index dce2ccc..5fbd7d8 100644 --- a/docs/internal/collect.md +++ b/docs/internal/collect.md @@ -16,6 +16,10 @@ API base URL, as weather-collection setup errors and fetch failures as bundle-collection errors. It does not retry, persist, select reports, derive facts, build modules, invoke Promptkit, or notify Distributor. +Fetch failures retain Weather API endpoint context but do not project upstream +response bodies into application-facing errors. Oversized response bodies fail +collection before source decoding. + ## Application Composition `internal/app` owns the narrow `Collector` interface used by workflow tests; diff --git a/internal/adapters/weatherapi/client.go b/internal/adapters/weatherapi/client.go index 1ad8a69..148ab27 100644 --- a/internal/adapters/weatherapi/client.go +++ b/internal/adapters/weatherapi/client.go @@ -7,6 +7,7 @@ import ( "crypto/sha256" "encoding/hex" "encoding/json" + "errors" "fmt" "io" "net/http" @@ -30,8 +31,11 @@ const ( defaultWarmupDelay = time.Second defaultFetchAttempts = 2 defaultFetchRetryDelay = time.Second + maxResponseBodyBytes = 10 << 20 ) +var errResponseBodyTooLarge = errors.New("response exceeds 10 MiB limit") + type Client struct { baseURL *url.URL httpClient *http.Client @@ -497,16 +501,12 @@ func (c *Client) warmupOnce(ctx context.Context, endpoint string) error { } defer resp.Body.Close() - body, err := io.ReadAll(io.LimitReader(resp.Body, 10<<20)) + _, err = readResponseBody(resp.Body) if err != nil { - err = fmt.Errorf("read %s response: %w", endpoint, err) - if ctx.Err() != nil { - return err - } - return retryableRequestError{err: err} + return responseReadError(ctx, endpoint, err) } if resp.StatusCode < 200 || resp.StatusCode >= 300 { - err := fmt.Errorf("fetch %s: unexpected HTTP status %d: %s", endpoint, resp.StatusCode, strings.TrimSpace(string(body))) + err := fmt.Errorf("fetch %s: unexpected HTTP status %d", endpoint, resp.StatusCode) if isRetryableHTTPStatus(resp.StatusCode) { return retryableRequestError{err: err} } @@ -559,16 +559,12 @@ func (c *Client) fetchHTTPOnce(ctx context.Context, endpoint string, opts queryO } defer resp.Body.Close() - body, err := io.ReadAll(io.LimitReader(resp.Body, 10<<20)) + body, err := readResponseBody(resp.Body) if err != nil { - err = fmt.Errorf("read %s response: %w", endpoint, err) - if ctx.Err() != nil { - return reqURL, nil, err - } - return reqURL, nil, retryableRequestError{err: err} + return reqURL, nil, responseReadError(ctx, endpoint, err) } if resp.StatusCode < 200 || resp.StatusCode >= 300 { - err := fmt.Errorf("fetch %s: unexpected HTTP status %d: %s", endpoint, resp.StatusCode, strings.TrimSpace(string(body))) + err := fmt.Errorf("fetch %s: unexpected HTTP status %d", endpoint, resp.StatusCode) if isRetryableHTTPStatus(resp.StatusCode) { return reqURL, nil, retryableRequestError{err: err} } @@ -577,6 +573,25 @@ func (c *Client) fetchHTTPOnce(ctx context.Context, endpoint string, opts queryO return reqURL, body, nil } +func readResponseBody(body io.Reader) ([]byte, error) { + data, err := io.ReadAll(io.LimitReader(body, maxResponseBodyBytes+1)) + if err != nil { + return nil, err + } + if int64(len(data)) > maxResponseBodyBytes { + return nil, errResponseBodyTooLarge + } + return data, nil +} + +func responseReadError(ctx context.Context, endpoint string, err error) error { + err = fmt.Errorf("read %s response: %w", endpoint, err) + if errors.Is(err, errResponseBodyTooLarge) || ctx.Err() != nil { + return err + } + return retryableRequestError{err: err} +} + type retryableRequestError struct { err error } diff --git a/internal/adapters/weatherapi/client_test.go b/internal/adapters/weatherapi/client_test.go index 9be004e..b66ddea 100644 --- a/internal/adapters/weatherapi/client_test.go +++ b/internal/adapters/weatherapi/client_test.go @@ -201,8 +201,9 @@ func TestFetchBundleRecordsSourceHash(t *testing.T) { } func TestHTTPErrorIsActionable(t *testing.T) { + const marker = "upstream-secret-marker" server := fixtureServer(t, map[string]handlerOverride{ - "/forecast/hourly": {status: http.StatusBadGateway, body: `upstream failed`}, + "/forecast/hourly": {status: http.StatusBadGateway, body: marker + strings.Repeat("x", 4096)}, }, nil) client := newTestClient(t, server.URL+"/", nil) @@ -213,6 +214,9 @@ func TestHTTPErrorIsActionable(t *testing.T) { if !strings.Contains(err.Error(), "/forecast/hourly") || !strings.Contains(err.Error(), "502") { t.Fatalf("error = %q, want endpoint and status", err.Error()) } + if strings.Contains(err.Error(), marker) { + t.Fatalf("error = %q, must not contain upstream response text", err.Error()) + } } func TestWarmupRetriesBeforeFetchBundle(t *testing.T) { @@ -291,6 +295,92 @@ func TestWarmupDoesNotRetryPermanentStatus(t *testing.T) { } } +func TestWarmupErrorDiagnosticsRedactResponseBody(t *testing.T) { + const marker = "upstream-secret-marker" + server := fixtureServer(t, map[string]handlerOverride{ + defaultWarmupEndpoint: {status: http.StatusNotFound, body: marker + strings.Repeat("x", 4096)}, + }, nil) + client := newTestClient(t, server.URL+"/", nil) + + _, err := client.FetchBundle(context.Background()) + if err == nil { + t.Fatal("FetchBundle() error = nil, want warmup error") + } + if !strings.Contains(err.Error(), defaultWarmupEndpoint) || !strings.Contains(err.Error(), "404") { + t.Fatalf("error = %q, want warmup endpoint and status", err.Error()) + } + if strings.Contains(err.Error(), marker) { + t.Fatalf("error = %q, must not contain upstream response text", err.Error()) + } +} + +func TestFetchAcceptsResponseAtBodyLimit(t *testing.T) { + body := paddedJSON(t, `{"data":null}`, int(maxResponseBodyBytes)) + server := fixtureServer(t, map[string]handlerOverride{ + "/forecast/narrative": {handler: func(w http.ResponseWriter, r *http.Request) { + _, _ = w.Write([]byte(body)) + }}, + }, nil) + client := newTestClient(t, server.URL+"/", nil) + + if _, err := client.FetchBundle(context.Background()); err != nil { + t.Fatalf("FetchBundle() error = %v", err) + } +} + +func TestFetchRejectsOversizedResponseWithoutRetry(t *testing.T) { + var requested []string + var narrativeCalls int + oversizedBody := paddedJSON(t, `{"data":null}`, int(maxResponseBodyBytes)) + "x" + server := fixtureServer(t, map[string]handlerOverride{ + "/forecast/narrative": {handler: func(w http.ResponseWriter, r *http.Request) { + narrativeCalls++ + _, _ = w.Write([]byte(oversizedBody)) + }}, + }, &requested) + client := newTestClient(t, server.URL+"/", nil) + + _, err := client.FetchBundle(context.Background()) + if err == nil { + t.Fatal("FetchBundle() error = nil, want oversized response error") + } + if !strings.Contains(err.Error(), "/forecast/narrative") || !strings.Contains(err.Error(), errResponseBodyTooLarge.Error()) { + t.Fatalf("error = %q, want endpoint and response limit", err.Error()) + } + if narrativeCalls != 1 { + t.Fatalf("narrative calls = %d, want no retry", narrativeCalls) + } + if containsPath(requested, "/alerts/active") { + t.Fatalf("requested paths = %v, want source failure before later fetches", requested) + } +} + +func TestWarmupRejectsOversizedResponseWithoutRetry(t *testing.T) { + var requested []string + oversizedBody := paddedJSON(t, `{"data":{}}`, int(maxResponseBodyBytes)) + "x" + server := fixtureServer(t, map[string]handlerOverride{ + defaultWarmupEndpoint: {handler: func(w http.ResponseWriter, r *http.Request) { + _, _ = w.Write([]byte(oversizedBody)) + }}, + }, &requested) + client := newTestClient(t, server.URL+"/", nil) + client.warmupAttempts = 2 + + _, err := client.FetchBundle(context.Background()) + if err == nil { + t.Fatal("FetchBundle() error = nil, want oversized warmup response error") + } + if !strings.Contains(err.Error(), defaultWarmupEndpoint) || !strings.Contains(err.Error(), errResponseBodyTooLarge.Error()) { + t.Fatalf("error = %q, want warmup endpoint and response limit", err.Error()) + } + if got := countPath(requested, defaultWarmupEndpoint); got != 1 { + t.Fatalf("warmup requests = %d, want no retry; all requests = %v", got, requested) + } + if containsPath(requested, "/observations") { + t.Fatalf("requested paths = %v, want warmup failure before source fetches", requested) + } +} + func TestFetchRetriesRetryableStatus(t *testing.T) { var hourlyCalls int server := fixtureServer(t, map[string]handlerOverride{ @@ -766,6 +856,14 @@ func fixedNow() time.Time { return time.Date(2026, 5, 29, 15, 0, 0, 0, time.UTC) } +func paddedJSON(t *testing.T, value string, size int) string { + t.Helper() + if len(value) > size { + t.Fatalf("JSON value length = %d, exceeds requested size %d", len(value), size) + } + return value + strings.Repeat(" ", size-len(value)) +} + func containsPath(requested []string, path string) bool { for _, rawURL := range requested { if strings.HasPrefix(rawURL, path+"?") || rawURL == path {