From bf6856d6a5beaaaf030791a03b70624b4072bb6c Mon Sep 17 00:00:00 2001 From: ranxi2001 Date: Sat, 3 Oct 2026 04:47:59 +0800 Subject: [PATCH] fix(server): make readiness dependency budget configurable --- backend/cmd/server/wire_gen.go | 2 +- backend/internal/config/config.go | 7 +- backend/internal/config/readiness_test.go | 49 ++++++++ backend/internal/server/lifecycle.go | 45 ++++--- .../server/lifecycle_readiness_test.go | 114 ++++++++++++++++++ backend/internal/server/lifecycle_test.go | 25 ++++ deploy/config.example.yaml | 3 + deploy/kubernetes/README.md | 19 ++- 8 files changed, 246 insertions(+), 18 deletions(-) create mode 100644 backend/internal/config/readiness_test.go create mode 100644 backend/internal/server/lifecycle_readiness_test.go diff --git a/backend/cmd/server/wire_gen.go b/backend/cmd/server/wire_gen.go index db27b019d12..9b727e262d2 100644 --- a/backend/cmd/server/wire_gen.go +++ b/backend/cmd/server/wire_gen.go @@ -387,7 +387,7 @@ func initializeApplication(buildInfo handler.BuildInfo) (*Application, error) { apiKeyAuthMiddleware := middleware.NewAPIKeyAuthMiddleware(apiKeyService, subscriptionService, configConfig) auditLogMiddleware := middleware.NewAuditLogMiddleware(auditLogService) stepUpAuthMiddleware := middleware.NewStepUpAuthMiddleware(totpService, userService, settingService) - lifecycle := server.ProvideLifecycle(db, redisClient) + lifecycle := server.ProvideLifecycle(configConfig, db, redisClient) engine := server.ProvideRouter(configConfig, handlers, jwtAuthMiddleware, optionalJWTAuthMiddleware, adminAuthMiddleware, apiKeyAuthMiddleware, auditLogMiddleware, stepUpAuthMiddleware, apiKeyService, subscriptionService, opsService, settingService, compositeRouteResolver, redisClient, lifecycle) httpServer := server.ProvideHTTPServer(configConfig, engine, lifecycle) opsMetricsCollector := service.ProvideOpsMetricsCollector(opsRepository, settingRepository, accountRepository, concurrencyService, db, redisClient, configConfig) diff --git a/backend/internal/config/config.go b/backend/internal/config/config.go index 8ec6c2ce08e..d19a2a904ea 100644 --- a/backend/internal/config/config.go +++ b/backend/internal/config/config.go @@ -690,6 +690,7 @@ type PricingConfig struct { type ServerConfig struct { GracefulShutdownTimeout int `mapstructure:"graceful_shutdown_timeout"` // seconds; 0 preserves the legacy 5s budget ShutdownDrainDelay int `mapstructure:"shutdown_drain_delay"` // seconds to withdraw from load balancers before closing the listener + ReadinessTimeoutSeconds int `mapstructure:"readiness_timeout_seconds"` // dependency probe budget; 0 preserves the 1s default Host string `mapstructure:"host"` Port int `mapstructure:"port"` Mode string `mapstructure:"mode"` // debug/release @@ -2166,7 +2167,8 @@ func setDefaults() { viper.SetDefault("server.mode", "release") viper.SetDefault("server.enable_server_timing", false) viper.SetDefault("server.frontend_url", "") - viper.SetDefault("server.read_header_timeout", 10) // 10秒读取请求头 + viper.SetDefault("server.read_header_timeout", 10) // 10秒读取请求头 + viper.SetDefault("server.readiness_timeout_seconds", 0) // 依赖探测超时,0 保持 1 秒默认值 viper.SetDefault("server.max_header_bytes", 64*1024) viper.SetDefault("server.idle_timeout", 120) // 120秒空闲超时 viper.SetDefault("server.max_request_body_size", int64(256*1024*1024)) @@ -2862,6 +2864,9 @@ func (c *Config) Validate() error { if c.Server.ReadHeaderTimeout < 1 || c.Server.ReadHeaderTimeout > 60 { return fmt.Errorf("server.read_header_timeout must be between 1 and 60 seconds") } + if c.Server.ReadinessTimeoutSeconds < 0 || c.Server.ReadinessTimeoutSeconds > 60 { + return fmt.Errorf("server.readiness_timeout_seconds must be between 0 and 60 seconds") + } if c.Server.MaxHeaderBytes < 8*1024 || c.Server.MaxHeaderBytes > 1024*1024 { return fmt.Errorf("server.max_header_bytes must be between 8192 and 1048576 bytes") } diff --git a/backend/internal/config/readiness_test.go b/backend/internal/config/readiness_test.go new file mode 100644 index 00000000000..46bd8a9ad27 --- /dev/null +++ b/backend/internal/config/readiness_test.go @@ -0,0 +1,49 @@ +package config + +import ( + "os" + "path/filepath" + "testing" + + "github.com/stretchr/testify/require" +) + +func TestLoadReadinessTimeout(t *testing.T) { + for _, tc := range []struct { + name string + file string + env string + want int + }{ + {name: "omitted keeps default", want: 0}, + {name: "YAML", file: "3", want: 3}, + {name: "environment only", env: "2", want: 2}, + {name: "environment overrides YAML", file: "3", env: "4", want: 4}, + {name: "zero resets YAML to default", file: "3", env: "0", want: 0}, + {name: "upper bound", env: "60", want: 60}, + } { + t.Run(tc.name, func(t *testing.T) { + resetViperWithJWTSecret(t) + t.Setenv("SERVER_READINESS_TIMEOUT_SECONDS", tc.env) + if tc.file != "" { + path := filepath.Join(t.TempDir(), "config.yaml") + require.NoError(t, os.WriteFile(path, []byte("server:\n readiness_timeout_seconds: "+tc.file+"\n"), 0o600)) + t.Setenv("CONFIG_FILE", path) + } + cfg, err := Load() + require.NoError(t, err) + require.Equal(t, tc.want, cfg.Server.ReadinessTimeoutSeconds) + }) + } +} + +func TestLoadReadinessTimeoutRejectsInvalidBudgets(t *testing.T) { + for _, value := range []string{"-1", "61"} { + t.Run(value, func(t *testing.T) { + resetViperWithJWTSecret(t) + t.Setenv("SERVER_READINESS_TIMEOUT_SECONDS", value) + _, err := Load() + require.ErrorContains(t, err, "server.readiness_timeout_seconds") + }) + } +} diff --git a/backend/internal/server/lifecycle.go b/backend/internal/server/lifecycle.go index d77f331d4ac..a9f85e949cb 100644 --- a/backend/internal/server/lifecycle.go +++ b/backend/internal/server/lifecycle.go @@ -9,28 +9,43 @@ import ( "sync" "time" + "github.com/Wei-Shaw/sub2api/internal/config" "github.com/redis/go-redis/v9" ) // Lifecycle counts handlers, including hijacked WebSocket handlers, without // wrapping ResponseWriter and losing its streaming or hijacking interfaces. type Lifecycle struct { - onDrain []func() - mu sync.Mutex - draining bool - active int - idle chan struct{} - checks []func(context.Context) error - probeMu sync.Mutex - probe *readinessProbe + onDrain []func() + mu sync.Mutex + draining bool + active int + idle chan struct{} + checks []func(context.Context) error + probeMu sync.Mutex + probe *readinessProbe + probeTimeout time.Duration } func NewLifecycle(checks ...func(context.Context) error) *Lifecycle { - return &Lifecycle{idle: make(chan struct{}), checks: checks} + return newLifecycleWithTimeout(time.Second, checks...) } -func ProvideLifecycle(db *sql.DB, cache *redis.Client) *Lifecycle { - return NewLifecycle(func(ctx context.Context) error { +// Both HTTP probes and internal readiness callers share this dependency budget. +// It is immutable after construction and does not affect liveness or draining. +func newLifecycleWithTimeout(timeout time.Duration, checks ...func(context.Context) error) *Lifecycle { + if timeout <= 0 { + timeout = time.Second + } + return &Lifecycle{idle: make(chan struct{}), checks: checks, probeTimeout: timeout} +} + +func ProvideLifecycle(cfg *config.Config, db *sql.DB, cache *redis.Client) *Lifecycle { + timeout := time.Second + if cfg != nil && cfg.Server.ReadinessTimeoutSeconds > 0 { + timeout = time.Duration(cfg.Server.ReadinessTimeoutSeconds) * time.Second + } + return newLifecycleWithTimeout(timeout, func(ctx context.Context) error { if db == nil || cache == nil { return errors.New("readiness dependencies unavailable") } @@ -118,9 +133,7 @@ func (l *Lifecycle) serveReadiness(w http.ResponseWriter, r *http.Request) { writeLifecycleStatus(w, http.StatusServiceUnavailable, "draining") return } - ctx, cancel := context.WithTimeout(r.Context(), time.Second) - defer cancel() - if err := l.checkWithinBudget(ctx); err != nil { + if err := l.checkWithinBudget(r.Context()); err != nil { writeLifecycleStatus(w, http.StatusServiceUnavailable, "not_ready") return } @@ -162,13 +175,15 @@ type readinessProbe struct { } func (l *Lifecycle) checkWithinBudget(ctx context.Context) error { + ctx, cancel := context.WithTimeout(ctx, l.probeTimeout) + defer cancel() l.probeMu.Lock() probe := l.probe if probe == nil { probe = &readinessProbe{done: make(chan struct{})} l.probe = probe go func() { - probeCtx, cancel := context.WithTimeout(context.Background(), time.Second) + probeCtx, cancel := context.WithTimeout(context.Background(), l.probeTimeout) defer cancel() err := l.checkDependencies(probeCtx) l.probeMu.Lock() diff --git a/backend/internal/server/lifecycle_readiness_test.go b/backend/internal/server/lifecycle_readiness_test.go new file mode 100644 index 00000000000..2f4e91165e7 --- /dev/null +++ b/backend/internal/server/lifecycle_readiness_test.go @@ -0,0 +1,114 @@ +package server + +import ( + "context" + "errors" + "net/http" + "net/http/httptest" + "sync/atomic" + "testing" + "time" + + "github.com/DATA-DOG/go-sqlmock" + "github.com/alicebob/miniredis/v2" + "github.com/redis/go-redis/v9" + "github.com/stretchr/testify/require" + + "github.com/Wei-Shaw/sub2api/internal/config" +) + +// Exercise the actual provider and HTTP wrapper: changing only the background +// dependency timeout would still let the old HTTP deadline reject the probe. +func TestProvideLifecycleSlowDependency(t *testing.T) { + for _, tc := range []struct { + name string + seconds int + want int + }{ + {name: "default rejects after one second", want: http.StatusServiceUnavailable}, + {name: "configured budget permits slow success", seconds: 2, want: http.StatusOK}, + } { + t.Run(tc.name, func(t *testing.T) { + t.Parallel() + db, mock, err := sqlmock.New(sqlmock.MonitorPingsOption(true)) + require.NoError(t, err) + t.Cleanup(func() { _ = db.Close() }) + mock.ExpectPing().WillDelayFor(1200 * time.Millisecond) + mr := miniredis.RunT(t) + cache := redis.NewClient(&redis.Options{Addr: mr.Addr()}) + t.Cleanup(func() { _ = cache.Close() }) + cfg := &config.Config{Server: config.ServerConfig{ReadinessTimeoutSeconds: tc.seconds}} + lifecycle := ProvideLifecycle(cfg, db, cache) + response := httptest.NewRecorder() + lifecycle.Wrap(http.NotFoundHandler()).ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + require.Equal(t, tc.want, response.Code) + require.NoError(t, mock.ExpectationsWereMet()) + }) + } +} + +func TestProvideLifecycleReadinessDefaults(t *testing.T) { + for _, cfg := range []*config.Config{nil, {}, {Server: config.ServerConfig{ReadinessTimeoutSeconds: 0}}} { + lifecycle := ProvideLifecycle(cfg, nil, nil) + require.Equal(t, time.Second, lifecycle.probeTimeout) + response := httptest.NewRecorder() + lifecycle.Wrap(http.NotFoundHandler()).ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + require.Equal(t, http.StatusServiceUnavailable, response.Code, "an unavailable dependency must still fail closed") + require.JSONEq(t, `{"status":"not_ready"}`, response.Body.String()) + } +} + +func TestProvideLifecycleConfiguredBudgetStillRejectsDependencyFailures(t *testing.T) { + for _, dependency := range []string{"postgres", "redis"} { + t.Run(dependency, func(t *testing.T) { + db, mock, err := sqlmock.New(sqlmock.MonitorPingsOption(true)) + require.NoError(t, err) + t.Cleanup(func() { _ = db.Close() }) + ping := mock.ExpectPing() + mr := miniredis.RunT(t) + const privateError = "synthetic-password-must-stay-private" + if dependency == "postgres" { + ping.WillReturnError(errors.New(privateError)) + } else { + mr.SetError(privateError) + } + cache := redis.NewClient(&redis.Options{Addr: mr.Addr(), MaxRetries: -1}) + t.Cleanup(func() { _ = cache.Close() }) + cfg := &config.Config{Server: config.ServerConfig{ReadinessTimeoutSeconds: 3}} + lifecycle := ProvideLifecycle(cfg, db, cache) + handler := lifecycle.Wrap(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + w.WriteHeader(http.StatusOK) + })) + response := httptest.NewRecorder() + handler.ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + require.Equal(t, http.StatusServiceUnavailable, response.Code) + require.JSONEq(t, `{"status":"not_ready"}`, response.Body.String()) + require.NotContains(t, response.Body.String(), privateError) + response = httptest.NewRecorder() + handler.ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/health", nil)) + require.Equal(t, http.StatusOK, response.Code) + require.NoError(t, mock.ExpectationsWereMet()) + }) + } +} + +func TestLifecycleInternalReadinessWaitIsBounded(t *testing.T) { + release := make(chan struct{}) + t.Cleanup(func() { close(release) }) + var calls atomic.Int32 + lifecycle := newLifecycleWithTimeout(25*time.Millisecond, func(context.Context) error { + calls.Add(1) + <-release // model a driver that ignores cancellation + return nil + }) + done := make(chan error, 1) + go func() { done <- lifecycle.checkWithinBudget(context.Background()) }() + select { + case err := <-done: + require.ErrorIs(t, err, context.DeadlineExceeded) + case <-time.After(time.Second): + t.Fatal("internal readiness caller remained blocked past the configured budget") + } + require.ErrorIs(t, lifecycle.checkWithinBudget(context.Background()), context.DeadlineExceeded) + require.Equal(t, int32(1), calls.Load(), "overlapping callers must share the stalled check") +} diff --git a/backend/internal/server/lifecycle_test.go b/backend/internal/server/lifecycle_test.go index b35a0243cc9..04269288c92 100644 --- a/backend/internal/server/lifecycle_test.go +++ b/backend/internal/server/lifecycle_test.go @@ -49,6 +49,31 @@ func TestLifecycleReadinessAndSetupNeverReturnSPA(t *testing.T) { require.Contains(t, response.Body.String(), "needs_setup") } +func TestLifecycleReadinessUsesConfiguredDependencyTimeout(t *testing.T) { + entered := make(chan struct{}) + release := make(chan struct{}) + lifecycle := newLifecycleWithTimeout(50*time.Millisecond, func(context.Context) error { + close(entered) + <-release + return nil + }) + t.Cleanup(func() { close(release) }) + handler := lifecycle.Wrap(http.NotFoundHandler()) + response := httptest.NewRecorder() + done := make(chan struct{}) + go func() { + handler.ServeHTTP(response, httptest.NewRequest(http.MethodGet, "/readyz", nil)) + close(done) + }() + <-entered + select { + case <-done: + require.Equal(t, http.StatusServiceUnavailable, response.Code) + case <-time.After(time.Second): + t.Fatal("configured readiness probe budget was not enforced") + } +} + func TestLifecycleDrainPreservesStreamAndRejectsNewRequests(t *testing.T) { lifecycle := NewLifecycle() release := make(chan struct{}) diff --git a/deploy/config.example.yaml b/deploy/config.example.yaml index e0242395bee..e6d120680ad 100644 --- a/deploy/config.example.yaml +++ b/deploy/config.example.yaml @@ -17,6 +17,9 @@ runtime: # 服务器配置 # ============================================================================= server: + # Shared PostgreSQL/Redis readiness budget: 0 = 1s (default), allowed 1..60s. + # If changed, Kubernetes readinessProbe.timeoutSeconds must exceed this budget. + readiness_timeout_seconds: 0 # HTTP/WS request drain budget after load-balancer withdrawal (seconds). # Default 5 preserves existing deployments; Kubernetes example uses 240. graceful_shutdown_timeout: 5 diff --git a/deploy/kubernetes/README.md b/deploy/kubernetes/README.md index c5b50c7c355..af90731cff4 100644 --- a/deploy/kubernetes/README.md +++ b/deploy/kubernetes/README.md @@ -42,11 +42,28 @@ server: 对应环境变量为 `RUNTIME_ROLE`、`SERVER_GRACEFUL_SHUTDOWN_TIMEOUT`、`SERVER_SHUTDOWN_DRAIN_DELAY`。timeout 范围 0–3600,drain delay 范围 0–300。需要通过 Compose 使用时显式加入 service.environment,宿主 .env 文件本身不会自动进入容器。 - `/health` 用作 liveness,只反映进程 HTTP 服务存活。 -- `/readyz` 用作 readiness:正常模式检查 PostgreSQL 和 Redis(共享 1 秒预算),依赖不可用时返回 503,不返回内部错误或凭据。 +- `/readyz` 用作 readiness:正常模式检查 PostgreSQL 和 Redis。默认共享 1 秒预算;使用远程共享依赖的 gateway 可将 `server.readiness_timeout_seconds` 设置为有界值(例如 3),依赖不可用时仍返回 503,不返回内部错误或凭据。 - 未完成配置的 setup 模式对 /readyz 返回 503,避免把前端 SPA 的 200 当成准备就绪。 - SIGTERM 到达后立即撤销 readiness,新业务请求返回 503/Retry-After,已有处理器继续运行;先等待 drain delay,再在 graceful timeout 内排空 HTTP 和仍在处理中的 hijacked/WS handler。到期后进程进入关闭流程。 - 示例为 10 秒撤流等待、240 秒请求排空、320 秒 Pod 宽限期。宽限期要覆盖两个阶段和后台清理余量。更长的流需要更大的预算;本功能不保证无限长连接不中断,也不会将已经输出的请求自动重放到另一个 Pod。 +`server.readiness_timeout_seconds` 支持 YAML 或 `SERVER_READINESS_TIMEOUT_SECONDS` 环境变量,重启生效;省略或设为 `0` 保持 1 秒,显式值允许 `1..60` 秒。这是 PostgreSQL 与 Redis 顺序检查的总预算。先测量真实依赖操作的耗时,再决定是否调整;仅连接 SSH 本地转发监听端口的耗时无法反映远端数据库往返。 + +例如依赖检查通常需要 1–2 秒时,可以为相应 gateway 显式配置 3 秒,并给 kubelet 留出余量: + +```yaml +env: + - name: SERVER_READINESS_TIMEOUT_SECONDS + value: "3" +readinessProbe: + httpGet: {path: /readyz, port: http} + timeoutSeconds: 5 + periodSeconds: 5 + failureThreshold: 2 +``` + +`readinessProbe.timeoutSeconds` 必须大于应用探测预算;只修改 kubelet 的超时不会改变应用内的预算。Compose 部署需在 `service.environment` 显式传入环境变量。调大预算会延长依赖故障时的撤流等待,不能用来掩盖持续故障;本配置不改变 scheduler 重建或 outbox 超时。内部调用者(例如 Serverless heartbeat)自身更短的 deadline 仍优先。`/health` 的存活检查、draining 时立即返回 503 和 setup 模式行为保持原样。 + ## 配置和存储 示例使用两个 gateway Pod 的 StatefulSet,提供稳定实例名和各自的 10 GiB RWO PVC。没有把两个 Pod 同时挂到一个本地数据目录。K3s 可使用默认 local-path,其他集群使用已配置的默认 StorageClass;local-path 不能使数据跟随 Pod 跨主机迁移,也不构成存储高可用。拓扑分布是软约束,单节点 k3s 同样可以调度两个副本。