Bound Distributor response diagnostics

This commit is contained in:
2026-08-13 02:55:25 +00:00
parent 0b57d99a97
commit 4b748c2e53
10 changed files with 497 additions and 46 deletions

View File

@@ -2,10 +2,12 @@
package distributor
import (
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net/http"
"os"
"strings"
@@ -107,6 +109,10 @@ type runStatus struct {
const statusPollInterval = 250 * time.Millisecond
const maxDistributorResponseBytes int64 = 1 << 20
var errDistributorResponseTooLarge = fmt.Errorf("distributor response exceeds the %d-byte limit", maxDistributorResponseBytes)
func New(cfg config.DistributorNotifyConfig) *Client {
return newClient(cfg, newDistributorUploadClient)
}
@@ -202,6 +208,7 @@ func (c *Client) Upload(ctx context.Context, req UploadRequest) (UploadResult, e
UploadStatus: result.Status,
}
status, statusErr := waitForRunStatus(runCtx, uploadClient, result.RunID, c.Timeout > 0)
status = sanitizeRunStatus(status)
if status.RunID != "" || status.Status != "" {
uploadResult.RunStatus = &RunStatus{
RunID: status.RunID,
@@ -218,7 +225,7 @@ func (c *Client) Upload(ctx context.Context, req UploadRequest) (UploadResult, e
}
}
if statusErr != nil {
uploadResult.StatusError = redactTokenString(statusErr.Error(), token)
uploadResult.StatusError = safeDistributorDiagnostic(statusErr, token).Error()
return uploadResult, nil
}
if status.Status == "failed" {
@@ -261,10 +268,52 @@ type distributorUploadClient struct {
client *distributorupload.Client
}
type boundedResponseTransport struct {
base http.RoundTripper
limit int64
}
func (t boundedResponseTransport) RoundTrip(req *http.Request) (*http.Response, error) {
base := t.base
if base == nil {
base = http.DefaultTransport
}
response, err := base.RoundTrip(req)
if err != nil {
return nil, err
}
defer response.Body.Close()
data, err := io.ReadAll(io.LimitReader(response.Body, t.limit+1))
if err != nil {
return nil, err
}
if int64(len(data)) > t.limit {
return nil, errDistributorResponseTooLarge
}
response.Body = io.NopCloser(bytes.NewReader(data))
response.ContentLength = int64(len(data))
return response, nil
}
type RemoteResponseError struct {
StatusCode int
Retryable bool
}
func (e *RemoteResponseError) Error() string {
if e == nil || e.StatusCode == 0 {
return "distributor request failed"
}
return fmt.Sprintf("distributor request failed with HTTP status %d", e.StatusCode)
}
func newDistributorUploadClient(endpoint, token string, timeout time.Duration) (uploadClient, error) {
httpClient := (*http.Client)(nil)
httpClient := &http.Client{
Transport: boundedResponseTransport{base: http.DefaultTransport, limit: maxDistributorResponseBytes},
}
if timeout > 0 {
httpClient = &http.Client{Timeout: timeout}
httpClient.Timeout = timeout
}
client, err := distributorupload.NewClient(distributorupload.ClientOptions{
Endpoint: endpoint,
@@ -306,7 +355,7 @@ func (c distributorUploadClient) Status(ctx context.Context, runID string) (runS
if err != nil {
return runStatus{}, err
}
return runStatus{
return sanitizeRunStatus(runStatus{
RunID: status.RunID,
PipelineID: status.PipelineID,
Status: status.Status,
@@ -315,7 +364,7 @@ func (c distributorUploadClient) Status(ctx context.Context, runID string) (runS
FinishedAt: status.FinishedAt,
Report: append(json.RawMessage(nil), status.Report...),
Error: status.Error,
}, nil
}), nil
}
type uploadErrorContext struct {
@@ -331,7 +380,7 @@ type uploadErrorContext struct {
func wrapUploadError(err error, ctx uploadErrorContext) error {
var conflict *distributorupload.IdempotencyConflictError
isConflict := errors.As(err, &conflict)
err = redactToken(err, ctx.Token)
err = safeDistributorDiagnostic(err, ctx.Token)
if isConflict {
return &IdempotencyConflictError{
Err: fmt.Errorf("upload distributor bundle %q to pipeline %q at endpoint %q with idempotency key %q from sources %q as bundle paths %q: idempotency conflict: %w", ctx.BundleID, ctx.PipelineID, ctx.Endpoint, ctx.IdempotencyKey, ctx.SourcePaths, ctx.BundlePaths, err),
@@ -340,6 +389,31 @@ func wrapUploadError(err error, ctx uploadErrorContext) error {
return fmt.Errorf("upload distributor bundle %q to pipeline %q at endpoint %q with idempotency key %q from sources %q as bundle paths %q: %w", ctx.BundleID, ctx.PipelineID, ctx.Endpoint, ctx.IdempotencyKey, ctx.SourcePaths, ctx.BundlePaths, err)
}
func safeDistributorDiagnostic(err error, token string) error {
if err == nil {
return nil
}
if errors.Is(err, errDistributorResponseTooLarge) {
return errDistributorResponseTooLarge
}
if errors.Is(err, context.Canceled) || errors.Is(err, context.DeadlineExceeded) {
return redactToken(err, token)
}
var httpErr *distributorupload.HTTPError
if errors.As(err, &httpErr) {
return &RemoteResponseError{StatusCode: httpErr.StatusCode, Retryable: httpErr.Retryable}
}
return errors.New("distributor request failed")
}
func sanitizeRunStatus(status runStatus) runStatus {
status.Report = nil
if status.Error != "" {
status.Error = "distributor reported a failed run"
}
return status
}
func uploadSourcePaths(files []UploadFile) []string {
paths := make([]string, 0, len(files))
for _, file := range files {