diff --git a/packages/michael/internal/handlers/blob.go b/packages/michael/internal/handlers/blob.go index 652839a4..8c5ac408 100644 --- a/packages/michael/internal/handlers/blob.go +++ b/packages/michael/internal/handlers/blob.go @@ -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 } diff --git a/packages/michael/internal/metrics/metrics.go b/packages/michael/internal/metrics/metrics.go index 0c4edfca..62277328 100644 --- a/packages/michael/internal/metrics/metrics.go +++ b/packages/michael/internal/metrics/metrics.go @@ -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() diff --git a/packages/michael/internal/storage/s3.go b/packages/michael/internal/storage/s3.go index ce66e4a1..2b6451b8 100644 --- a/packages/michael/internal/storage/s3.go +++ b/packages/michael/internal/storage/s3.go @@ -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 diff --git a/packages/michael/internal/storage/s3_test.go b/packages/michael/internal/storage/s3_test.go index f7cbf58a..b7ec2b9f 100644 --- a/packages/michael/internal/storage/s3_test.go +++ b/packages/michael/internal/storage/s3_test.go @@ -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", diff --git a/packages/yuctl/internal/do/do.go b/packages/yuctl/internal/do/do.go index e9aed3f9..60124beb 100644 --- a/packages/yuctl/internal/do/do.go +++ b/packages/yuctl/internal/do/do.go @@ -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 } diff --git a/tf/deployment/prod/htz-fsn1/talos/cilium-values.yaml.tftpl b/tf/deployment/prod/htz-fsn1/talos/cilium-values.yaml.tftpl index e3fdb81c..5833165a 100644 --- a/tf/deployment/prod/htz-fsn1/talos/cilium-values.yaml.tftpl +++ b/tf/deployment/prod/htz-fsn1/talos/cilium-values.yaml.tftpl @@ -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 diff --git a/tf/deployment/prod/htz-fsn1/talos/cilium.tf b/tf/deployment/prod/htz-fsn1/talos/cilium.tf index e265b025..eb021065 100644 --- a/tf/deployment/prod/htz-fsn1/talos/cilium.tf +++ b/tf/deployment/prod/htz-fsn1/talos/cilium.tf @@ -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.