mirror of
https://github.com/immich-app/yucca.git
synced 2026-09-30 13:33:00 +08:00
fix(michael): retry loop (#367)
This commit is contained in:
@@ -152,7 +152,10 @@ func (s *Server) getBlob(w http.ResponseWriter, r *http.Request) {
|
||||
rangeHeader := r.Header.Get("Range")
|
||||
obj, err := s.Storage.GetObject(r.Context(), a.Repository, key, rangeHeader)
|
||||
if err != nil {
|
||||
hlog.FromRequest(r).Error().Err(err).Msg("get blob failed")
|
||||
hlog.FromRequest(r).Error().Err(err).Str("repository", a.Repository).Str("key", key).Msg("get blob failed: backend storage error")
|
||||
if s.Metrics != nil {
|
||||
s.Metrics.StorageErrors.Add(r.Context(), 1, metrics.StorageErrorOption("get", metrics.BlobType(r)))
|
||||
}
|
||||
writeError(w, r,http.StatusInternalServerError, "An error occurred with the storage server")
|
||||
return
|
||||
}
|
||||
@@ -176,7 +179,10 @@ func (s *Server) saveBlob(w http.ResponseWriter, r *http.Request) {
|
||||
writeError(w, r,http.StatusBadRequest, "Content hash does not match blob name")
|
||||
return
|
||||
}
|
||||
hlog.FromRequest(r).Error().Err(err).Msg("save blob failed")
|
||||
hlog.FromRequest(r).Error().Err(err).Str("repository", a.Repository).Str("key", key).Msg("save blob failed: backend storage error")
|
||||
if s.Metrics != nil {
|
||||
s.Metrics.StorageErrors.Add(r.Context(), 1, metrics.StorageErrorOption("put", metrics.BlobType(r)))
|
||||
}
|
||||
writeError(w, r,http.StatusInternalServerError, "An error occurred with the storage server")
|
||||
return
|
||||
}
|
||||
|
||||
@@ -48,6 +48,7 @@ type Metrics struct {
|
||||
RequestTTFB otelmetric.Float64Histogram
|
||||
RequestCount otelmetric.Int64Counter
|
||||
RequestErrors otelmetric.Int64Counter
|
||||
StorageErrors otelmetric.Int64Counter
|
||||
AuthCacheHits otelmetric.Int64Counter
|
||||
AuthCacheMisses otelmetric.Int64Counter
|
||||
}
|
||||
@@ -115,6 +116,16 @@ func NewMetrics(meter otelmetric.Meter) (*Metrics, error) {
|
||||
return nil, fmt.Errorf("creating request_errors counter: %w", err)
|
||||
}
|
||||
|
||||
// Backend (RGW) storage-operation failures, split out from the generic HTTP
|
||||
// error counter so a gateway write-error spike is directly visible and
|
||||
// alertable — the signal that a client-facing 500 wave is the storage
|
||||
// backend, not michael or auth.
|
||||
storageErrors, err := meter.Int64Counter("storage.backend.errors",
|
||||
otelmetric.WithDescription("Backend S3 storage operation failures, by operation and blob type"))
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("creating storage_errors counter: %w", err)
|
||||
}
|
||||
|
||||
authCacheHits, err := meter.Int64Counter("auth.cache.hits",
|
||||
otelmetric.WithDescription("JWT verifications served from the token cache"))
|
||||
if err != nil {
|
||||
@@ -136,11 +147,33 @@ func NewMetrics(meter otelmetric.Meter) (*Metrics, error) {
|
||||
RequestTTFB: requestTTFB,
|
||||
RequestCount: requestCount,
|
||||
RequestErrors: requestErrors,
|
||||
StorageErrors: storageErrors,
|
||||
AuthCacheHits: authCacheHits,
|
||||
AuthCacheMisses: authCacheMisses,
|
||||
}, nil
|
||||
}
|
||||
|
||||
type storageErrAttrKey struct{ operation, blobType string }
|
||||
|
||||
var storageErrAttrCache sync.Map
|
||||
|
||||
// StorageErrorOption labels a backend storage failure by operation ("put",
|
||||
// "get", …) and blob type. Deliberately low-cardinality (no per-user labels):
|
||||
// this is a fleet-health/alerting signal, and the per-request identity is
|
||||
// already on the logged error line.
|
||||
func StorageErrorOption(operation, blobType string) otelmetric.MeasurementOption {
|
||||
key := storageErrAttrKey{operation, blobType}
|
||||
if v, ok := storageErrAttrCache.Load(key); ok {
|
||||
return v.(otelmetric.MeasurementOption)
|
||||
}
|
||||
opt := otelmetric.WithAttributeSet(attribute.NewSet(
|
||||
attribute.String("operation", operation),
|
||||
attribute.String("type", blobType),
|
||||
))
|
||||
storageErrAttrCache.Store(key, opt)
|
||||
return opt
|
||||
}
|
||||
|
||||
func SetupMeterProvider(cfg config.Config) (*sdkmetric.MeterProvider, error) {
|
||||
ctx := context.Background()
|
||||
|
||||
|
||||
@@ -262,9 +262,18 @@ func (s *S3Storage) PutObject(ctx context.Context, bucket, key string, body io.R
|
||||
}
|
||||
|
||||
// Use unsigned payload so the SDK doesn't need to seek the body for signing.
|
||||
// RetryMaxAttempts=1 (no SDK retry): the body is restic's non-seekable
|
||||
// proxied stream, so any retry would try to rewind it and fail with "failed
|
||||
// to rewind transport stream for retry, request stream is not seekable" —
|
||||
// masking the real backend error and, under load, stalling clients on a
|
||||
// large fraction of PUTs. With retries off, a transient gateway error
|
||||
// surfaces cleanly and restic retries the pack itself (its body IS
|
||||
// seekable). See TestPutObject_NonSeekableBodyOn503_NoRewindRetry.
|
||||
_, err := s.client.PutObject(ctx, input, s3.WithAPIOptions(
|
||||
v4.SwapComputePayloadSHA256ForUnsignedPayloadMiddleware,
|
||||
))
|
||||
), func(o *s3.Options) {
|
||||
o.RetryMaxAttempts = 1
|
||||
})
|
||||
if err != nil {
|
||||
if isPreconditionFailed(err) {
|
||||
return ErrPreconditionFailed
|
||||
|
||||
@@ -3,8 +3,11 @@ package storage
|
||||
import (
|
||||
"context"
|
||||
"errors"
|
||||
"io"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
"testing"
|
||||
|
||||
"michael/internal/config"
|
||||
@@ -12,6 +15,14 @@ import (
|
||||
"github.com/aws/aws-sdk-go-v2/service/s3/types"
|
||||
)
|
||||
|
||||
// readOnly hides any io.Seeker the underlying reader implements, so the value
|
||||
// handed to the S3 SDK is a plain, non-seekable stream — exactly what michael
|
||||
// proxies in production (restic's r.Body, further wrapped in an io.TeeReader
|
||||
// for hashing). The SDK cannot rewind it to retry.
|
||||
type readOnly struct{ r io.Reader }
|
||||
|
||||
func (ro readOnly) Read(p []byte) (int, error) { return ro.r.Read(p) }
|
||||
|
||||
func TestIsPreconditionFailed_NilError(t *testing.T) {
|
||||
if isPreconditionFailed(nil) {
|
||||
t.Error("expected false for nil error")
|
||||
@@ -166,6 +177,49 @@ func TestListObjects_EmptyListing(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestPutObject_NonSeekableBodyOn503_NoRewindRetry reproduces the production
|
||||
// throughput collapse: michael streams restic's non-seekable pack body straight
|
||||
// into S3 PutObject, and when the gateway returns a retryable 5xx the SDK's
|
||||
// default retryer tries to rewind the body to resend it — which fails with
|
||||
// "failed to rewind transport stream for retry, request stream is not seekable"
|
||||
// and surfaces as an opaque 500 to restic. Under load this hit ~28% of requests
|
||||
// and stalled the whole fleet.
|
||||
//
|
||||
// The desired behaviour (asserted here) is: no rewind is ever attempted, exactly
|
||||
// one upload is made, and the caller gets the clean underlying backend error —
|
||||
// which restic retries at the pack level (its own body IS seekable). RED before
|
||||
// the s3.go RetryMaxAttempts=1 fix, GREEN after.
|
||||
func TestPutObject_NonSeekableBodyOn503_NoRewindRetry(t *testing.T) {
|
||||
var puts int32
|
||||
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
if r.Method == http.MethodPut {
|
||||
atomic.AddInt32(&puts, 1)
|
||||
_, _ = io.Copy(io.Discard, r.Body) // let the client finish sending
|
||||
w.WriteHeader(http.StatusServiceUnavailable)
|
||||
return
|
||||
}
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}))
|
||||
defer srv.Close()
|
||||
|
||||
s := NewS3StorageForEndpoint(probeConfig(), srv.URL)
|
||||
body := readOnly{strings.NewReader("restic-pack-bytes")}
|
||||
err := s.PutObject(context.Background(), "bucket", "data/deadbeef", body, int64(len("restic-pack-bytes")), true, "")
|
||||
|
||||
if err == nil {
|
||||
t.Fatal("expected an error from the 503 gateway, got nil")
|
||||
}
|
||||
// The bug: a retry on the non-seekable body fails to rewind and masks the
|
||||
// real backend error.
|
||||
if msg := err.Error(); strings.Contains(msg, "rewind") || strings.Contains(msg, "not seekable") {
|
||||
t.Fatalf("PutObject attempted a rewind-for-retry on a non-seekable body: %v", err)
|
||||
}
|
||||
// And with retries off, the gateway must see exactly one upload attempt.
|
||||
if n := atomic.LoadInt32(&puts); n != 1 {
|
||||
t.Fatalf("expected exactly 1 upload attempt (no retry on a non-seekable body), got %d", n)
|
||||
}
|
||||
}
|
||||
|
||||
func probeConfig() config.Config {
|
||||
return config.Config{
|
||||
S3AccessKeyID: "test",
|
||||
|
||||
@@ -215,24 +215,32 @@ func fromGodo(d godo.Droplet) Droplet {
|
||||
return out
|
||||
}
|
||||
|
||||
// CreateDroplets creates one batch of droplets in a single region (DO's
|
||||
// multi-create is per-region) with the fleet tag and ssh key.
|
||||
// maxMultiCreate is DO's hard cap on names per multi-create request; larger
|
||||
// batches 422 ("cannot create more than 10 droplets at a time"), so we chunk.
|
||||
const maxMultiCreate = 10
|
||||
|
||||
// CreateDroplets creates droplets in a single region (DO's multi-create is
|
||||
// per-region) with the fleet tag and ssh key, chunked to DO's per-request cap.
|
||||
func (c *Client) CreateDroplets(ctx context.Context, names []string, region, size, image, tag string, keyID int) ([]Droplet, error) {
|
||||
req := &godo.DropletMultiCreateRequest{
|
||||
Names: names,
|
||||
Region: region,
|
||||
Size: size,
|
||||
Image: godo.DropletCreateImage{Slug: image},
|
||||
SSHKeys: []godo.DropletCreateSSHKey{{ID: keyID}},
|
||||
Tags: []string{tag},
|
||||
}
|
||||
droplets, _, err := c.do.Droplets.CreateMultiple(ctx, req)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("create %d droplets in %s: %w", len(names), region, err)
|
||||
}
|
||||
out := make([]Droplet, len(droplets))
|
||||
for i, d := range droplets {
|
||||
out[i] = fromGodo(d)
|
||||
var out []Droplet
|
||||
for start := 0; start < len(names); start += maxMultiCreate {
|
||||
end := min(start+maxMultiCreate, len(names))
|
||||
batch := names[start:end]
|
||||
req := &godo.DropletMultiCreateRequest{
|
||||
Names: batch,
|
||||
Region: region,
|
||||
Size: size,
|
||||
Image: godo.DropletCreateImage{Slug: image},
|
||||
SSHKeys: []godo.DropletCreateSSHKey{{ID: keyID}},
|
||||
Tags: []string{tag},
|
||||
}
|
||||
droplets, _, err := c.do.Droplets.CreateMultiple(ctx, req)
|
||||
if err != nil {
|
||||
return out, fmt.Errorf("create %d droplets in %s: %w", len(batch), region, err)
|
||||
}
|
||||
for _, d := range droplets {
|
||||
out = append(out, fromGodo(d))
|
||||
}
|
||||
}
|
||||
return out, nil
|
||||
}
|
||||
|
||||
@@ -1,32 +1,30 @@
|
||||
# Cilium for the prod cluster. Talos sets cni:none + proxy:disabled, so nodes
|
||||
# stay NotReady until this lands the datapath.
|
||||
#
|
||||
# TUNNEL (geneve) routing at jumbo MTU. The cluster spans two L2 domains — CPs
|
||||
# on the kube-cp VLAN (11), workers on the kube VLAN (10) — routed by the spine
|
||||
# IRBs. A geneve overlay carries pod↔pod (and host↔pod, e.g. apiserver→pod
|
||||
# webhooks) over the routed node-to-node path with the pod identity in-band —
|
||||
# which autoDirectNodeRoutes (same-L2 only) can't span. The jumbo MTU (below),
|
||||
# the spine LB ECMP (core-fabric), and DSR (loadBalancer, below) are retained
|
||||
# and independent of the routing mode.
|
||||
# TUNNEL (geneve) routing. The cluster spans two L2 domains — CPs on the kube-cp
|
||||
# VLAN (11), workers on the kube VLAN (10) — routed by the spine IRBs. A geneve
|
||||
# overlay carries pod↔pod (and host↔pod, e.g. apiserver→pod webhooks) over the
|
||||
# routed node-to-node path with the pod identity in-band — which
|
||||
# autoDirectNodeRoutes (same-L2 only) can't span. The spine LB ECMP (core-fabric)
|
||||
# and DSR (loadBalancer, below) are independent of the routing mode and retained.
|
||||
#
|
||||
# MTU: auto-detected (NOT set). Cilium takes the lowest MTU across its devices,
|
||||
# which is wt0 (NetBird's WireGuard, 1280) → ~1230 pod MTU. This is DELIBERATE:
|
||||
# a single cluster-wide MTU governs the whole pod datapath, and pods egress the
|
||||
# 1280 mesh tunnel to reach o11y (vmauth on 10.69.0.10 via wt0). Forcing a jumbo
|
||||
# MTU (tried 2026-07-28) black-holes every pod's large frames over WireGuard
|
||||
# (PMTUD doesn't survive WG+NAT) and takes o11y remote-write down. The fabric's
|
||||
# 1500-bound flows (client↔michael, michael↔RGW) don't benefit from jumbo pods
|
||||
# anyway, so there is no reason to raise it past the mesh ceiling.
|
||||
#
|
||||
# securityContext + cgroup blocks are MANDATORY on Talos. Ref: Talos "Deploying Cilium".
|
||||
|
||||
ipam:
|
||||
mode: kubernetes
|
||||
|
||||
# Tunnel (geneve) routing — see the header. East-west (incl. host↔pod across
|
||||
# the two VLANs) rides the fabric encapsulated; geneve overhead is ~50B, so at
|
||||
# the jumbo MTU below pods still get ~8950 (vs the 1230 the old auto-detected
|
||||
# MTU clamped us to over wt0).
|
||||
routingMode: tunnel
|
||||
tunnelProtocol: geneve
|
||||
|
||||
# Explicit datapath MTU (the fabric bond/VLAN MTU; switch L2 is 9216, IRBs
|
||||
# 9202). NEVER auto-detect: Cilium's detection takes the lowest MTU across its
|
||||
# devices, and wt0 (NetBird's WireGuard, 1280) once clamped every pod route to
|
||||
# 1230. In tunnel mode pods get this minus the ~50B geneve header.
|
||||
MTU: ${fabric_mtu}
|
||||
|
||||
# Masquerade pod traffic to non-pod destinations (internet, node IPs); pod↔pod is
|
||||
# tunnelled so it isn't masqueraded. ip-masq-agent excludes only the pod CIDR, so
|
||||
# node-IP traffic egressing wt0 (NetBird, the operator plane) is still SNAT'd to
|
||||
|
||||
@@ -1,8 +1,8 @@
|
||||
# Cilium CNI — installed post-bootstrap in the same apply (helm provider bound to
|
||||
# the bootstrap CP, providers.tf). Talos set cni:none + proxy:disabled, so nodes
|
||||
# go Ready only once this lands the datapath. Native routing at the fabric MTU
|
||||
# spans the two routed VLANs (kube/kube-cp) via the spine IRBs + Cilium iBGP
|
||||
# PodCIDR advertisements (cilium-values.yaml.tftpl).
|
||||
# go Ready only once this lands the datapath. Tunnel (geneve) routing spans the
|
||||
# two routed VLANs (kube/kube-cp) via the spine IRBs; MTU is auto-detected
|
||||
# (wt0-limited) — see cilium-values.yaml.tftpl.
|
||||
resource "helm_release" "cilium" {
|
||||
name = "cilium"
|
||||
namespace = "kube-system"
|
||||
@@ -12,7 +12,6 @@ resource "helm_release" "cilium" {
|
||||
|
||||
values = [templatefile("${path.module}/cilium-values.yaml.tftpl", {
|
||||
pod_cidr = local.pod_cidr
|
||||
fabric_mtu = local.fabric_mtu
|
||||
kube_proxy_replacement = true
|
||||
hubble = var.cluster.hubble
|
||||
})]
|
||||
@@ -21,14 +20,7 @@ resource "helm_release" "cilium" {
|
||||
timeout = 600
|
||||
cleanup_on_fail = true
|
||||
|
||||
# The machineconfig applies must land first: Cilium's explicit MTU assumes
|
||||
# the bonds are already at fabric_mtu (a 9000 datapath over still-1500 bonds
|
||||
# blackholes east-west until the node config catches up).
|
||||
depends_on = [
|
||||
talos_cluster_kubeconfig.this,
|
||||
talos_machine_configuration_apply.cp,
|
||||
talos_machine_configuration_apply.worker,
|
||||
]
|
||||
depends_on = [talos_cluster_kubeconfig.this]
|
||||
}
|
||||
|
||||
# Full health gate AFTER the CNI is in — now node-Ready is achievable.
|
||||
|
||||
Reference in New Issue
Block a user