Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion backend/cmd/server/wire_gen.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

7 changes: 6 additions & 1 deletion backend/internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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")
}
Expand Down
49 changes: 49 additions & 0 deletions backend/internal/config/readiness_test.go
Original file line number Diff line number Diff line change
@@ -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")
})
}
}
45 changes: 30 additions & 15 deletions backend/internal/server/lifecycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
}
Expand Down Expand Up @@ -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
}
Expand Down Expand Up @@ -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()
Expand Down
114 changes: 114 additions & 0 deletions backend/internal/server/lifecycle_readiness_test.go
Original file line number Diff line number Diff line change
@@ -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")
}
25 changes: 25 additions & 0 deletions backend/internal/server/lifecycle_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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{})
Expand Down
3 changes: 3 additions & 0 deletions deploy/config.example.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
19 changes: 18 additions & 1 deletion deploy/kubernetes/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 同样可以调度两个副本。
Expand Down
Loading