Bound Weather API response diagnostics
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
Reference in New Issue
Block a user