Merge branch 'main' into fix/improved-error-handling-import

This commit is contained in:
izzy
2026-09-14 14:21:17 +01:00
43 changed files with 610 additions and 51 deletions
+4
View File
@@ -3,4 +3,8 @@
set -e
source "$(dirname "$0")/../restic-api/env"
# No load balancer in front of a local run, so there is nothing to drain away
# from — skip the fleet's 15s SIGTERM drain and let Ctrl-C exit immediately.
export DRAIN_DELAY_MS=${DRAIN_DELAY_MS:-0}
cd packages/michael && go run .
+1 -1
View File
@@ -1,3 +1,3 @@
{
".": "0.40.2"
".": "0.41.0"
}
+18
View File
@@ -1,5 +1,23 @@
# Changelog
## [0.41.0](https://github.com/immich-app/yucca/compare/v0.40.2...v0.41.0) (2026-09-14)
### Features
* first batch of closed beta improvements ([#657](https://github.com/immich-app/yucca/issues/657)) ([42a9324](https://github.com/immich-app/yucca/commit/42a932406c4a65a7783c69444513922879936514))
* **monk:** measure scrub lateness and mirror the mgr not-scrubbed check ([#644](https://github.com/immich-app/yucca/issues/644)) ([5cafca4](https://github.com/immich-app/yucca/commit/5cafca4853b33413fb9c702af00b2888f6918287))
### Bug Fixes
* **michael:** drain in-flight restic requests before closing the listener on rollout ([#660](https://github.com/immich-app/yucca/issues/660)) ([0b15733](https://github.com/immich-app/yucca/commit/0b1573300509c4022755d880e2a3f6950a3c63f1))
* **monk:** shut down cleanly on SIGTERM and drop the per-PG map allocation ([5626011](https://github.com/immich-app/yucca/commit/5626011081b0680a1762b2aa713a55b16a75d98f))
* **monk:** stop cleanly on SIGTERM and drop a per-PG map allocation ([#646](https://github.com/immich-app/yucca/issues/646)) ([5626011](https://github.com/immich-app/yucca/commit/5626011081b0680a1762b2aa713a55b16a75d98f))
* **o11y:** keep sflow at its native 5s resolution ([#659](https://github.com/immich-app/yucca/issues/659)) ([33d0bb9](https://github.com/immich-app/yucca/commit/33d0bb948c8280e5636a968f845124d045f8fd70))
* **o11y:** key the deep scrub verdict on lateness, not interval age ([#645](https://github.com/immich-app/yucca/issues/645)) ([85273ac](https://github.com/immich-app/yucca/commit/85273ac719adc7ac47b250fa0e81cc3fd7c9afb4))
* **yucca sdk:** report interrupted backup runs to the backend on restart ([#651](https://github.com/immich-app/yucca/issues/651)) ([a744273](https://github.com/immich-app/yucca/commit/a744273730e3cbf492c9a2d17dc7213299129250))
## [0.40.2](https://github.com/immich-app/yucca/compare/v0.40.1...v0.40.2) (2026-09-11)
+1 -1
View File
@@ -3,7 +3,7 @@ name: columbo
description: Ticket investigation agent (Go)
type: application
version: 0.1.0
appVersion: "0.40.2" # x-release-please-version
appVersion: "0.41.0" # x-release-please-version
dependencies:
- name: yucca-common
version: 0.2.0
+1 -1
View File
@@ -3,7 +3,7 @@ name: futo-backups-bot
description: FUTO Backups Discord support bot (NestJS)
type: application
version: 0.1.0
appVersion: "0.40.2" # x-release-please-version
appVersion: "0.41.0" # x-release-please-version
dependencies:
- name: yucca-common
version: 0.2.0
+1 -1
View File
@@ -3,7 +3,7 @@ name: michael
description: Yucca backup/restore service (Go)
type: application
version: 0.1.0
appVersion: "0.40.2" # x-release-please-version
appVersion: "0.41.0" # x-release-please-version
dependencies:
- name: yucca-common
version: 0.2.0
+32 -2
View File
@@ -3,6 +3,27 @@
# Stateless S3 proxy, safe to run N-way.
replicas: 2
# Never dip below the declared replica count mid-rollout. The Deployment default
# (maxUnavailable 25%) would take 3 of father's 12 gateways down at once, on top
# of the 3 it surges in — a quarter of the restic data plane gone while the
# replacements are still probing their RGW backends.
strategy:
type: RollingUpdate
rollingUpdate:
maxSurge: 25%
maxUnavailable: 0
# Half the fleet, not the lib chart's minAvailable: 1 — a node drain may
# otherwise evict 11 of 12 gateways simultaneously.
pdb:
minAvailable: 50%
# Must exceed DRAIN_DELAY_MS + SHUTDOWN_TIMEOUT_MS (15s + 120s in the base
# HelmRelease), or the kubelet SIGKILLs the pod while a restic blob is still on
# the wire — precisely the failure the drain exists to avoid. It is a ceiling,
# not a wait: an idle pod exits as soon as the drain delay elapses.
terminationGracePeriodSeconds: 150
# Generous requests, deliberately huge memory limits (never CPU-limit).
resources:
requests: { cpu: 250m, memory: 256Mi }
@@ -52,6 +73,10 @@ env:
value: "3010"
- name: LOG_LEVEL
value: debug
# Dev only: the fleet default is 15s, which would make every `tilt down` and
# pod delete sit for a quarter minute with nothing to drain.
- name: DRAIN_DELAY_MS
value: "1000"
# S3 object store = Rook-Ceph RGW. michael creates one bucket per restic
# repository, so it needs a full RGW user (CephObjectStoreUser via
# charts/ceph-objectuser), NOT a bucket-scoped ObjectBucketClaim. Rook writes
@@ -115,6 +140,11 @@ startupProbe:
tcpSocket: { port: http }
periodSeconds: 2
failureThreshold: 90
# Readiness is an HTTP endpoint, not a TCP connect, because it has to report
# one thing a bound socket cannot: that the process has begun draining. Probed
# every 2s so a SIGTERMed pod leaves the Service's endpoints well inside its
# drain delay, taking new restic requests with it.
readinessProbe:
tcpSocket: { port: http }
periodSeconds: 10
httpGet: { path: /readyz, port: http }
periodSeconds: 2
failureThreshold: 2
+1 -1
View File
@@ -3,7 +3,7 @@ name: web
description: Yucca web UI (SvelteKit SSR)
type: application
version: 0.1.0
appVersion: "0.40.2" # x-release-please-version
appVersion: "0.41.0" # x-release-please-version
dependencies:
- name: yucca-common
version: 0.2.0
+1 -1
View File
@@ -3,7 +3,7 @@ name: yucca-admin-api
description: Yucca admin API (NestJS)
type: application
version: 0.1.0
appVersion: "0.40.2" # x-release-please-version
appVersion: "0.41.0" # x-release-please-version
dependencies:
- name: yucca-common
version: 0.2.0
+1 -1
View File
@@ -3,7 +3,7 @@ name: yucca-api
description: Yucca REST API (NestJS)
type: application
version: 0.1.0
appVersion: "0.40.2" # x-release-please-version
appVersion: "0.41.0" # x-release-please-version
dependencies:
- name: yucca-common
version: 0.2.0
+1 -1
View File
@@ -3,7 +3,7 @@ name: yucca-metrics-worker
description: Yucca metrics worker (NestJS)
type: application
version: 0.1.0
appVersion: "0.40.2" # x-release-please-version
appVersion: "0.41.0" # x-release-please-version
dependencies:
- name: yucca-common
version: 0.2.0
@@ -7,6 +7,10 @@ metadata:
{{- include "yucca-common.labels" . | nindent 4 }}
spec:
replicas: {{ .Values.replicas | default 1 }}
{{- with .Values.strategy }}
strategy:
{{- toYaml . | nindent 4 }}
{{- end }}
selector:
matchLabels:
{{- include "yucca-common.selectorLabels" . | nindent 6 }}
@@ -19,6 +23,9 @@ spec:
{{- toYaml . | nindent 8 }}
{{- end }}
spec:
{{- with .Values.terminationGracePeriodSeconds }}
terminationGracePeriodSeconds: {{ . }}
{{- end }}
{{- with .Values.serviceAccountName }}
serviceAccountName: {{ . }}
{{- end }}
@@ -6,6 +6,10 @@ metadata:
name: michael
spec:
interval: 1h
# Well above helm's 5m default: with maxUnavailable 0 a terminating pod holds
# its slot until its in-flight restic requests drain, so father's 12 gateways
# roll at the pace of DRAIN_DELAY_MS + SHUTDOWN_TIMEOUT_MS per wave.
timeout: 15m
chart:
spec:
chart: charts/apps/michael
@@ -42,6 +46,17 @@ spec:
value: "3010"
- name: LOG_LEVEL
value: info
# Rollout drain. DRAIN_DELAY_MS is the window between /readyz going 503
# and the listener closing — it has to outlast readiness detection (2s
# probe x 2 failures) plus EndpointSlice and Envoy EDS propagation, or the
# gateway is still steering fresh restic requests at a closing socket.
# SHUTDOWN_TIMEOUT_MS then caps how long an in-flight blob transfer gets
# to finish. The chart's terminationGracePeriodSeconds must stay above
# their sum.
- name: DRAIN_DELAY_MS
value: "15000"
- name: SHUTDOWN_TIMEOUT_MS
value: "120000"
# Must match the cluster's STORAGE_CLUSTER_CODE (topology ConfigMap):
# yucca-api mints restic tokens with a storageCluster claim, and michael
# fails closed on cluster codes it does not front.
@@ -429,3 +429,8 @@ spec:
bearerTokenSecret:
name: vmagent-remote-write
key: token
# o11y's vmstorage runs no dedup (only vmselect, at 20s on read), so
# a cumulative series re-exported unchanged is paid for in mesh
# bandwidth and disk before anything drops it. Collapse it here.
streamAggrConfig:
dedupInterval: 20s
@@ -19,7 +19,7 @@ data:
url: http://victoria-metrics.netops.svc:8428
isDefault: true
jsonData:
timeInterval: 10s
timeInterval: 5s
---
apiVersion: v1
kind: ConfigMap
@@ -1,5 +1,5 @@
# VictoriaMetrics single — the SHORT-TERM, HIGH-GRANULARITY tier of the hybrid
# design: 30d retention at 10s scrape resolution, in-cluster. vmagent remote-writes
# design: 30d retention, 5s dedup floor, in-cluster. vmagent remote-writes
# here today and will FAN OUT to the o11y cluster later (second -remoteWrite.url in
# vmagent.yaml) — this instance stays the fast local buffer, o11y the long-term one.
#
@@ -38,7 +38,10 @@ spec:
args:
- -storageDataPath=/storage
- -retentionPeriod=30d
- -dedup.minScrapeInterval=10s
# Tracks the FINEST scrape in the tier: sflow at 5s, itself pinned to
# the fabric's sFlow polling-interval (tf core-fabric). At 10s this
# dropped every other counter sample from that seconds-granularity feed.
- -dedup.minScrapeInterval=5s
- -httpListenAddr=:8428
ports:
- containerPort: 8428
@@ -19,5 +19,5 @@ metadata:
spec:
interval: 1m
ref:
tag: v0.40.2 # x-release-please-version
tag: v0.41.0 # x-release-please-version
url: https://github.com/immich-app/yucca.git
@@ -13,4 +13,4 @@ metadata:
name: image-versions
namespace: flux-system
data:
YUCCA_IMAGE_TAG: v0.40.2 # x-release-please-version
YUCCA_IMAGE_TAG: v0.41.0 # x-release-please-version
+2 -1
View File
@@ -17,7 +17,8 @@ spec:
namespace: yucca
interval: 1h
retryInterval: 2m
timeout: 10m
# Must outlast the HelmRelease's own 15m rollout budget (draining gateways).
timeout: 20m
path: ./kubernetes/apps/base/michael
prune: true
wait: true
+2 -2
View File
@@ -1,12 +1,12 @@
{
"name": "yucca-monorepo",
"version": "0.40.2",
"version": "0.41.0",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"name": "yucca-monorepo",
"version": "0.40.2",
"version": "0.41.0",
"devDependencies": {
"@eslint/js": "^9.39.2",
"eslint": "^9.39.2",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "yucca-monorepo",
"version": "0.40.2",
"version": "0.41.0",
"description": "Monorepo for yucca",
"private": true,
"devDependencies": {
+1 -1
View File
@@ -2,4 +2,4 @@
// by release-please (extra-files in release-please-config.json).
package version
const Version = "0.40.2" // x-release-please-version
const Version = "0.41.0" // x-release-please-version
+1 -1
View File
@@ -1,7 +1,7 @@
{
"private": true,
"name": "@common/server",
"version": "0.40.2",
"version": "0.41.0",
"description": "Common code for the NestJS applications",
"keywords": [],
"author": "",
+3 -3
View File
@@ -4,9 +4,9 @@ export const otelEnv = z
.object({
NODE_ENV: z.enum(['development', 'production', 'test', 'provision']).default('development'),
OTEL_DEBUG: z.coerce.boolean().default(false),
OTEL_SAMPLE_RATE: z.number().min(0).max(1).default(1),
OTEL_METRICS_EXPORT_INTERVAL: z.number().default(10_000),
OTEL_DEBUG: z.stringbool().default(false),
OTEL_SAMPLE_RATE: z.coerce.number().min(0).max(1).default(1),
OTEL_METRICS_EXPORT_INTERVAL: z.coerce.number().default(10_000),
OTEL_METRICS: z.string().default('http://localhost:8428/opentelemetry/v1/metrics'),
OTEL_TRACING: z.string().default('http://localhost:10428/insert/opentelemetry/v1/traces'),
OTEL_LOGGING: z.string().default('http://localhost:9428/insert/opentelemetry/v1/logs'),
+1 -1
View File
@@ -26,7 +26,7 @@ const otelSDK = new NodeSDK({
url: otelEnv.OTEL_METRICS,
temporalityPreference: metrics.AggregationTemporality.CUMULATIVE,
}),
exportIntervalMillis: 1000,
exportIntervalMillis: otelEnv.OTEL_METRICS_EXPORT_INTERVAL,
}),
// tracing
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "e2e",
"version": "0.40.2",
"version": "0.41.0",
"description": "",
"type": "module",
"scripts": {
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "futo-backups-bot",
"version": "0.40.2",
"version": "0.41.0",
"description": "",
"author": "",
"private": true,
+30 -1
View File
@@ -64,6 +64,18 @@ type Config struct {
// address into. Only its LAST entry is trusted — see geoip.ClientAddr.
ClientIPHeader string
// DrainDelay is how long michael keeps serving normally after SIGTERM while
// /readyz already answers 503. It must outlast readiness detection plus
// EndpointSlice and Envoy EDS propagation: cut short, the gateway is still
// routing fresh restic requests at a socket that has begun closing.
DrainDelay time.Duration
// ShutdownTimeout caps how long in-flight requests get to finish once the
// drain delay has elapsed. One restic blob can be tens of megabytes over a
// slow client uplink, so this is minutes — and the pod's
// terminationGracePeriodSeconds must exceed DrainDelay + ShutdownTimeout, or
// the kubelet SIGKILLs mid-upload anyway.
ShutdownTimeout time.Duration
OTLPMetricsEndpoint string
OTLPMetricsURLPath string
OTLPMetricsInterval time.Duration
@@ -239,9 +251,12 @@ func LoadConfig() Config {
asnDatabasePath := envOr("ASN_DB_PATH", "/etc/michael/asn.mmdb")
clientIPHeader := envOr("CLIENT_IP_HEADER", "X-Forwarded-For")
drainDelay := envDurationMS("DRAIN_DELAY_MS", 15*time.Second)
shutdownTimeout := envDurationMS("SHUTDOWN_TIMEOUT_MS", 2*time.Minute)
otlpEndpoint := os.Getenv("OTLP_METRICS_ENDPOINT")
otlpURLPath := os.Getenv("OTLP_METRICS_URL_PATH")
otlpInterval := 1000 * time.Millisecond
otlpInterval := 10 * time.Second
if v := os.Getenv("OTLP_METRICS_INTERVAL_MS"); v != "" {
ms, err := strconv.Atoi(v)
if err != nil {
@@ -296,6 +311,8 @@ func LoadConfig() Config {
S3DefaultCluster: defaultCluster,
ASNDatabasePath: asnDatabasePath,
ClientIPHeader: clientIPHeader,
DrainDelay: drainDelay,
ShutdownTimeout: shutdownTimeout,
OTLPMetricsEndpoint: otlpEndpoint,
OTLPMetricsURLPath: otlpURLPath,
OTLPMetricsInterval: otlpInterval,
@@ -610,6 +627,18 @@ func boolOr(v *bool, fallback bool) bool {
return fallback
}
func envDurationMS(key string, fallback time.Duration) time.Duration {
v := os.Getenv(key)
if v == "" {
return fallback
}
ms, err := strconv.Atoi(v)
if err != nil || ms < 0 {
log.Fatal().Msgf("%s must be a number >= 0", key)
}
return time.Duration(ms) * time.Millisecond
}
func envOr(key, fallback string) string {
if v := os.Getenv(key); v != "" {
return v
+36 -1
View File
@@ -5,6 +5,7 @@ import (
"crypto/ecdsa"
"fmt"
"net/http"
"sync/atomic"
"time"
"michael/internal/auth"
@@ -33,8 +34,21 @@ type Server struct {
// metrics and the access log. Optional: nil leaves both unattributed, which
// is what the tests and any deployment without an ASN database get.
ResolveClient func(*http.Request) geoip.Client
draining atomic.Bool
}
// readyPath is the kubelet's readiness probe. A repository path is always a
// UUID (auth rejects anything else), so no repository can collide with it.
const readyPath = "/readyz"
// BeginDrain fails readiness from here on. Nothing already running is touched
// and requests that still arrive are still served in full — the point is only
// to get this replica out of the gateway's endpoint set before the listener
// closes, so a rollout moves restic to another replica instead of resetting it
// mid-backup.
func (s *Server) BeginDrain() { s.draining.Store(true) }
// NewServer builds a single-cluster server: everything is served from s under
// the default cluster code.
func NewServer(s storage.Storage, jwtPublicKey *ecdsa.PublicKey, m *metrics.Metrics) *Server {
@@ -170,7 +184,28 @@ func (s *Server) Handler() http.Handler {
})
})
return r
return s.withReadiness(r)
}
// withReadiness answers the readiness probe ahead of the router, keeping a
// probe every couple of seconds out of the access log and the request metrics.
//
// Backend health is deliberately not part of the answer: an RGW outage is
// shared fate across every replica, so failing readiness on it would empty the
// Service's endpoints and turn a degraded data plane into an unreachable one.
// The pool sheds with 503 + Retry-After for that case instead.
func (s *Server) withReadiness(next http.Handler) http.Handler {
return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
if r.URL.Path != readyPath {
next.ServeHTTP(w, r)
return
}
if s.draining.Load() {
http.Error(w, "draining", http.StatusServiceUnavailable)
return
}
w.WriteHeader(http.StatusOK)
})
}
// op stamps the operation name and the resolved route pattern onto the request
@@ -418,3 +418,46 @@ func TestClientNetworkLogOutput(t *testing.T) {
}
}
}
func TestReadyzReportsDrainState(t *testing.T) {
srv := newTestServer(&mockStorage{})
handler := srv.Handler()
rec := httptest.NewRecorder()
handler.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, readyPath, nil))
if rec.Code != http.StatusOK {
t.Fatalf("expected 200 before draining, got %d", rec.Code)
}
srv.BeginDrain()
rec = httptest.NewRecorder()
handler.ServeHTTP(rec, httptest.NewRequest(http.MethodGet, readyPath, nil))
if rec.Code != http.StatusServiceUnavailable {
t.Fatalf("expected 503 while draining, got %d", rec.Code)
}
}
// A draining instance must keep serving: readiness only steers the gateway
// away, it never rejects a request that still arrives.
func TestDrainingStillServesRequests(t *testing.T) {
store := &mockStorage{
checkBucketFn: func(context.Context, string) (bool, error) { return false, nil },
createBucketFn: func(context.Context, string) error { return nil },
}
srv := newTestServer(store)
srv.BeginDrain()
req := httptest.NewRequest(http.MethodPost, "/"+testRepository+"/?create=true", nil)
req.Header.Set("Authorization", makeBasicAuth(makeJWT(t, jwt.MapClaims{
"user": testUser,
"repository": testRepository,
"writeOnce": false,
})))
rec := httptest.NewRecorder()
srv.Handler().ServeHTTP(rec, req)
if rec.Code != http.StatusOK {
t.Fatalf("expected the repository create to succeed while draining, got %d: %s", rec.Code, rec.Body)
}
}
+1 -1
View File
@@ -2,4 +2,4 @@
// by release-please (extra-files in release-please-config.json).
package version
const Version = "0.40.2" // x-release-please-version
const Version = "0.41.0" // x-release-please-version
+37 -11
View File
@@ -119,13 +119,15 @@ func main() {
Handler: srv.Handler(),
}
// Graceful shutdown
ctx, stop := signal.NotifyContext(context.Background(), syscall.SIGINT, syscall.SIGTERM)
defer stop()
// Drive each backend pool's reconcile/probe loop until shutdown.
// The pools outlive the signal context on purpose: requests still draining
// after SIGTERM need a backend set someone is still resolving and probing.
poolCtx, stopPools := context.WithCancel(context.Background())
defer stopPools()
for _, pool := range pools {
go pool.Run(ctx)
go pool.Run(poolCtx)
}
go func() {
@@ -136,27 +138,33 @@ func main() {
}()
<-ctx.Done()
log.Info().Msg("shutting down")
// Restore default signal handling: a second SIGTERM/SIGINT from an impatient
// operator now kills the process instead of being swallowed by the drain.
stop()
shutdownCtx, cancel := context.WithTimeout(context.Background(), 10*time.Second)
defer cancel()
if err := httpSrv.Shutdown(shutdownCtx); err != nil {
log.Fatal().Err(err).Msg("shutdown error")
if err := drain(srv, httpSrv, cfg.DrainDelay, cfg.ShutdownTimeout); err != nil {
log.Error().Err(err).Msg("shutdown deadline hit with requests still in flight")
}
stopPools()
if err := asnDB.Close(); err != nil {
log.Error().Err(err).Msg("ASN database close error")
}
// Telemetry gets its own budget: a drain that used up shutdownCtx would
// otherwise hand the exporters an already-expired deadline and lose the
// final flush — exactly the window worth having logs for.
telemetryCtx, cancelTelemetry := context.WithTimeout(context.Background(), 10*time.Second)
defer cancelTelemetry()
if meterProvider != nil {
if err := meterProvider.Shutdown(shutdownCtx); err != nil {
if err := meterProvider.Shutdown(telemetryCtx); err != nil {
log.Error().Err(err).Msg("meter provider shutdown error")
}
}
if otelLogWriter != nil {
if err := otelLogWriter.Shutdown(shutdownCtx); err != nil {
if err := otelLogWriter.Shutdown(telemetryCtx); err != nil {
log.Error().Err(err).Msg("OTLP log provider shutdown error")
}
}
@@ -164,6 +172,24 @@ func main() {
log.Info().Msg("shutdown complete")
}
// drain retires this instance in two phases. Phase one fails /readyz while
// serving exactly as before, giving the readiness probe, the EndpointSlice and
// Envoy's EDS push time to move new restic requests onto another replica. Only
// then does phase two close the listener and wait out whatever is still in
// flight — closing it up front would reset live blob transfers and fail the
// backup that was running.
func drain(srv *handlers.Server, httpSrv *http.Server, delay, timeout time.Duration) error {
srv.BeginDrain()
log.Info().Dur("drain_delay", delay).Msg("draining: readiness failing, still serving")
time.Sleep(delay)
ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()
log.Info().Dur("timeout", timeout).Msg("closing listener, waiting for in-flight requests")
return httpSrv.Shutdown(ctx)
}
// buildClusters builds one Storage per storage cluster michael fronts, keyed by
// cluster code, plus the subset of those that are load-balancing pools (returned
// concretely so the caller can register their metrics and run their reconcile
+343
View File
@@ -0,0 +1,343 @@
package main
import (
"bufio"
"context"
"crypto/ecdsa"
"crypto/elliptic"
"crypto/rand"
"encoding/base64"
"errors"
"io"
"net"
"net/http"
"strconv"
"strings"
"sync"
"testing"
"time"
"michael/internal/handlers"
"michael/internal/storage"
"github.com/golang-jwt/jwt/v5"
)
var drainKey, _ = ecdsa.GenerateKey(elliptic.P256(), rand.Reader)
const (
drainUser = "00000000-0000-0000-0000-0000000000d1"
drainRepo = "00000000-0000-0000-0000-0000000000d2"
drainBlob = "1111111111111111111111111111111111111111111111111111111111111111"
)
// drainStorage stands in for the S3 backend. PutObject drains the request body
// and then blocks until release is closed, so a test can hold one blob upload
// open across an entire drain the way a slow restic client would.
type drainStorage struct {
started chan string
release chan struct{}
releaseOnce sync.Once
mu sync.Mutex
stored map[string]int
}
func (d *drainStorage) releaseAll() { d.releaseOnce.Do(func() { close(d.release) }) }
func newDrainStorage() *drainStorage {
return &drainStorage{
started: make(chan string, 8),
release: make(chan struct{}),
stored: map[string]int{},
}
}
func (d *drainStorage) PutObject(ctx context.Context, bucket, key string, body io.Reader, _ int64, _ bool, _ string) error {
n, err := io.Copy(io.Discard, body)
if err != nil {
return err
}
d.started <- key
// "locks" blobs return immediately: a test needs some request that completes
// on its own while another is still pinned open.
if !strings.HasPrefix(key, "locks/") {
<-d.release
}
d.mu.Lock()
d.stored[key] = int(n)
d.mu.Unlock()
return nil
}
func (d *drainStorage) storedBytes(key string) (int, bool) {
d.mu.Lock()
defer d.mu.Unlock()
n, ok := d.stored[key]
return n, ok
}
func (d *drainStorage) CheckBucket(context.Context, string) (bool, error) { return true, nil }
func (d *drainStorage) CreateBucket(context.Context, string) error { return nil }
func (d *drainStorage) HeadObject(context.Context, string, string) (int64, error) {
return 0, nil
}
func (d *drainStorage) GetObject(context.Context, string, string, string) (*storage.S3Object, error) {
return nil, errors.New("not used")
}
func (d *drainStorage) ListObjects(context.Context, string, string, func(storage.BlobInfo) error) error {
return nil
}
func (d *drainStorage) DeleteObject(context.Context, string, string) error { return nil }
func drainAuthHeader(t *testing.T) string {
t.Helper()
token := jwt.NewWithClaims(jwt.SigningMethodES256, jwt.MapClaims{
"user": drainUser,
"repository": drainRepo,
"writeOnce": false,
})
signed, err := token.SignedString(drainKey)
if err != nil {
t.Fatalf("sign token: %v", err)
}
return "Basic " + base64.StdEncoding.EncodeToString([]byte("restic:"+signed))
}
// drainHarness is michael's real handler stack on a real TCP listener, reached
// over a real HTTP client — the drain is about socket lifecycle, so httptest's
// in-memory shortcuts would test the wrong thing.
type drainHarness struct {
store *drainStorage
srv *handlers.Server
httpSrv *http.Server
addr string
client *http.Client
}
func newDrainHarness(t *testing.T) *drainHarness {
t.Helper()
store := newDrainStorage()
srv := handlers.NewServer(store, &drainKey.PublicKey, nil)
ln, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
t.Fatalf("listen: %v", err)
}
httpSrv := &http.Server{Handler: srv.Handler()}
go func() { _ = httpSrv.Serve(ln) }()
h := &drainHarness{
store: store,
srv: srv,
httpSrv: httpSrv,
addr: ln.Addr().String(),
client: &http.Client{Transport: &http.Transport{}},
}
t.Cleanup(func() {
store.releaseAll()
_ = httpSrv.Close()
})
return h
}
func (h *drainHarness) url(path string) string { return "http://" + h.addr + path }
func (h *drainHarness) putBlob(t *testing.T, blobType, name, body string) (*http.Response, error) {
t.Helper()
req, err := http.NewRequest(http.MethodPost, h.url("/"+drainRepo+"/"+blobType+"/"+name), strings.NewReader(body))
if err != nil {
t.Fatalf("build request: %v", err)
}
req.Header.Set("Authorization", drainAuthHeader(t))
return h.client.Do(req)
}
func (h *drainHarness) readyStatus(t *testing.T) int {
t.Helper()
resp, err := h.client.Get(h.url("/readyz"))
if err != nil {
t.Fatalf("GET /readyz: %v", err)
}
defer resp.Body.Close()
_, _ = io.Copy(io.Discard, resp.Body)
return resp.StatusCode
}
// The whole point of the two phases: while the drain delay runs the instance is
// unready but fully functional, and only once the listener closes does it stop
// taking connections — with the upload that was already running still finishing.
func TestDrainServesThroughTheDelayThenFinishesInFlight(t *testing.T) {
h := newDrainHarness(t)
payload := strings.Repeat("x", 4096)
if got := h.readyStatus(t); got != http.StatusOK {
t.Fatalf("ready before drain: got %d, want 200", got)
}
uploadDone := make(chan *http.Response, 1)
uploadErr := make(chan error, 1)
go func() {
resp, err := h.putBlob(t, "data", drainBlob, payload)
if err != nil {
uploadErr <- err
return
}
uploadDone <- resp
}()
select {
case key := <-h.store.started:
if key != "data/"+drainBlob {
t.Fatalf("unexpected in-flight key %q", key)
}
case err := <-uploadErr:
t.Fatalf("upload failed before drain started: %v", err)
case <-time.After(5 * time.Second):
t.Fatal("upload never reached the storage layer")
}
drainErr := make(chan error, 1)
drainStarted := time.Now()
go func() { drainErr <- drain(h.srv, h.httpSrv, 750*time.Millisecond, 10*time.Second) }()
// Phase one: unready, but new work is still accepted and served in full.
waitFor(t, 2*time.Second, "readyz to report draining", func() bool {
return h.readyStatus(t) == http.StatusServiceUnavailable
})
if elapsed := time.Since(drainStarted); elapsed > 700*time.Millisecond {
t.Fatalf("readiness only flipped after %s — the delay had already elapsed", elapsed)
}
lockName := strings.Repeat("a", 64)
resp, err := h.putBlob(t, "locks", lockName, "lock")
if err != nil {
t.Fatalf("request during drain delay failed: %v", err)
}
resp.Body.Close()
if resp.StatusCode != http.StatusOK {
t.Fatalf("request during drain delay: got %d, want 200", resp.StatusCode)
}
<-h.store.started
// Phase two: the listener closes, so nothing new can connect.
waitFor(t, 5*time.Second, "listener to stop accepting", func() bool {
c, err := net.DialTimeout("tcp", h.addr, 500*time.Millisecond)
if err != nil {
return true
}
_ = c.Close()
return false
})
select {
case <-uploadDone:
t.Fatal("the in-flight upload completed before it was released")
default:
}
h.store.releaseAll()
select {
case resp := <-uploadDone:
defer resp.Body.Close()
if resp.StatusCode != http.StatusOK {
t.Fatalf("in-flight upload finished with %d, want 200", resp.StatusCode)
}
case err := <-uploadErr:
t.Fatalf("in-flight upload was cut off by the shutdown: %v", err)
case <-time.After(5 * time.Second):
t.Fatal("in-flight upload never completed")
}
if n, ok := h.store.storedBytes("data/" + drainBlob); !ok || n != len(payload) {
t.Fatalf("blob stored as %d bytes (present=%v), want %d", n, ok, len(payload))
}
select {
case err := <-drainErr:
if err != nil {
t.Fatalf("drain returned %v, want nil once the upload finished", err)
}
case <-time.After(5 * time.Second):
t.Fatal("drain never returned after the upload finished")
}
}
// A pooled-but-idle keep-alive is what restic leaves behind between blobs. It
// must not hold the pod open, and the server has to close it rather than leave
// restic holding a socket it will never answer again. Driven over a raw
// connection so the assertion is about the server's FIN and not about whatever
// http.Transport decides to do with its pool.
func TestDrainClosesIdleKeepAlives(t *testing.T) {
h := newDrainHarness(t)
conn, err := net.Dial("tcp", h.addr)
if err != nil {
t.Fatalf("dial: %v", err)
}
defer func() { _ = conn.Close() }()
body := "lock"
req := "POST /" + drainRepo + "/locks/" + strings.Repeat("b", 64) + " HTTP/1.1\r\n" +
"Host: " + h.addr + "\r\n" +
"Authorization: " + drainAuthHeader(t) + "\r\n" +
"Content-Length: " + strconv.Itoa(len(body)) + "\r\n\r\n" + body
if _, err := conn.Write([]byte(req)); err != nil {
t.Fatalf("write request: %v", err)
}
br := bufio.NewReader(conn)
resp, err := http.ReadResponse(br, nil)
if err != nil {
t.Fatalf("read response: %v", err)
}
_, _ = io.Copy(io.Discard, resp.Body)
resp.Body.Close()
if resp.StatusCode != http.StatusOK {
t.Fatalf("warm-up request: got %d, want 200", resp.StatusCode)
}
if resp.Close {
t.Fatal("server closed the connection before the drain; the test proves nothing")
}
<-h.store.started
started := time.Now()
if err := drain(h.srv, h.httpSrv, 0, 10*time.Second); err != nil {
t.Fatalf("drain returned %v, want nil with only idle connections open", err)
}
if elapsed := time.Since(started); elapsed > 3*time.Second {
t.Fatalf("drain waited %s on an idle keep-alive", elapsed)
}
_ = conn.SetReadDeadline(time.Now().Add(2 * time.Second))
if _, err := br.Read(make([]byte, 1)); !errors.Is(err, io.EOF) {
t.Fatalf("idle keep-alive read returned %v, want EOF (server should have closed it)", err)
}
}
// The shutdown timeout is a ceiling, not a promise: a transfer that outlives it
// is reported rather than waited on forever, so the kubelet's grace period is
// never what ends the process.
func TestDrainReportsRequestsPastTheDeadline(t *testing.T) {
h := newDrainHarness(t)
go func() { _, _ = h.putBlob(t, "data", drainBlob, "payload") }()
<-h.store.started
err := drain(h.srv, h.httpSrv, 0, 250*time.Millisecond)
if !errors.Is(err, context.DeadlineExceeded) {
t.Fatalf("drain returned %v, want context.DeadlineExceeded", err)
}
}
func waitFor(t *testing.T, limit time.Duration, what string, cond func() bool) {
t.Helper()
deadline := time.Now().Add(limit)
for time.Now().Before(deadline) {
if cond() {
return
}
time.Sleep(10 * time.Millisecond)
}
t.Fatalf("timed out after %s waiting for %s", limit, what)
}
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "mock-oidc-provider",
"version": "0.40.2",
"version": "0.41.0",
"description": "In-repo OIDC mock provider for dev and integration tests",
"private": true,
"license": "UNLICENSED",
+1 -1
View File
@@ -2,4 +2,4 @@
// by release-please (extra-files in release-please-config.json).
package version
const Version = "0.40.2" // x-release-please-version
const Version = "0.41.0" // x-release-please-version
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "restic-api",
"version": "0.40.2",
"version": "0.41.0",
"description": "",
"author": "",
"private": true,
+1 -1
View File
@@ -1,7 +1,7 @@
{
"name": "web",
"private": true,
"version": "0.40.2",
"version": "0.41.0",
"type": "module",
"files": [
"build"
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "yucca-admin-api",
"version": "0.40.2",
"version": "0.41.0",
"description": "",
"author": "",
"private": true,
+1 -1
View File
@@ -1,7 +1,7 @@
{
"name": "@futo-org/backups-api-client",
"repository": "https://github.com/immich-app/yucca",
"version": "0.40.2",
"version": "0.41.0",
"description": "Auto-generated SDK for backups API",
"type": "module",
"main": "dist/index.js",
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "yucca-api",
"version": "0.40.2",
"version": "0.41.0",
"description": "",
"author": "",
"private": true,
+1 -1
View File
@@ -1,6 +1,6 @@
{
"name": "yucca-metrics-worker",
"version": "0.40.2",
"version": "0.41.0",
"description": "",
"author": "",
"private": true,
@@ -1,7 +1,7 @@
{
"name": "@futo-org/backups-orchestrator-api",
"repository": "https://github.com/immich-app/yucca",
"version": "0.40.2",
"version": "0.41.0",
"description": "Backups orchestrator (API)",
"main": "dist/index.js",
"types": "dist/index.d.ts",
@@ -2,7 +2,7 @@
"name": "@futo-org/backups-orchestrator-ui",
"repository": "https://github.com/immich-app/yucca",
"description": "Backups orchestrator (UI library)",
"version": "0.40.2",
"version": "0.41.0",
"license": "Source First License 1.1",
"scripts": {
"dev": "vite dev --port 5174",