From 92e328a575aadb494307a2e473acf318ff974d38 Mon Sep 17 00:00:00 2001 From: wangbo Date: Sat, 1 Aug 2026 19:29:08 +0800 Subject: [PATCH] =?UTF-8?q?fix(acceptance):=20=E9=9A=94=E7=A6=BB=E5=AE=B9?= =?UTF-8?q?=E9=87=8F=E5=8E=8B=E6=B5=8B=E5=B9=B6=E7=A8=B3=E5=AE=9A=E6=95=B0?= =?UTF-8?q?=E6=8D=AE=E5=BA=93=E8=BF=9E=E6=8E=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 将协议模拟验收的供应商并发限额与 Worker 容量门禁分离,真实金丝雀继续保留生产限流。readyz 改用关键连接池并预热关键及 River 池,延长生产连接空闲周期,降低跨地域连接抖动。失败 Run 现在可以原子替换且失败报告记录任务数量,避免 validation 之间短暂放开正式流量。\n\n验证:Go 全量测试、go vet、PostgreSQL 集成测试、ShellCheck、迁移安全、发布脚本、pnpm lint/test/build、OpenAPI 无漂移均通过。 --- apps/api/cmd/gateway/main.go | 4 +- .../availability_timeout_integration_test.go | 43 +++++++++++ apps/api/internal/httpapi/handlers.go | 6 +- apps/api/internal/runner/admission.go | 40 +++++++++- apps/api/internal/runner/admission_test.go | 44 ++++++++++- apps/api/internal/runner/service.go | 6 +- apps/api/internal/store/acceptance.go | 37 +++++++++- .../store/acceptance_integration_test.go | 74 +++++++++++++++++++ .../easyai-ai-gateway-cluster-release | 2 +- ...ai-ai-gateway-cluster-release.conf.example | 2 +- scripts/cluster/run-production-acceptance.sh | 61 ++++++++------- 11 files changed, 277 insertions(+), 42 deletions(-) diff --git a/apps/api/cmd/gateway/main.go b/apps/api/cmd/gateway/main.go index 9bac2bd..a17228f 100644 --- a/apps/api/cmd/gateway/main.go +++ b/apps/api/cmd/gateway/main.go @@ -62,7 +62,7 @@ func main() { if cfg.DatabaseCriticalMaxConns > 0 { coordinationDB, err = store.ConnectWithPoolOptions(ctx, cfg.DatabaseURL, store.PostgresPoolOptions{ MaxConns: cfg.DatabaseCriticalMaxConns, - MinIdleConns: min(cfg.DatabaseCriticalMaxConns, 1), + MinIdleConns: min(cfg.DatabaseCriticalMaxConns, max(cfg.DatabaseMinIdleConns, 1)), MaxConnIdleTime: time.Duration(cfg.DatabaseMaxConnIdleSeconds) * time.Second, IdleInTransactionTimeout: time.Duration(cfg.DatabaseIdleInTransactionTimeoutSeconds) * time.Second, LockTimeout: time.Duration(cfg.DatabaseLockTimeoutSeconds) * time.Second, @@ -78,7 +78,7 @@ func main() { if cfg.DatabaseRiverMaxConns > 0 { riverDB, err = store.ConnectWithPoolOptions(ctx, cfg.DatabaseURL, store.PostgresPoolOptions{ MaxConns: cfg.DatabaseRiverMaxConns, - MinIdleConns: min(cfg.DatabaseRiverMaxConns, 1), + MinIdleConns: min(cfg.DatabaseRiverMaxConns, max(cfg.DatabaseMinIdleConns, 1)), MaxConnIdleTime: time.Duration(cfg.DatabaseMaxConnIdleSeconds) * time.Second, IdleInTransactionTimeout: time.Duration(cfg.DatabaseIdleInTransactionTimeoutSeconds) * time.Second, LockTimeout: time.Duration(cfg.DatabaseLockTimeoutSeconds) * time.Second, diff --git a/apps/api/internal/httpapi/availability_timeout_integration_test.go b/apps/api/internal/httpapi/availability_timeout_integration_test.go index 964627c..8fb806d 100644 --- a/apps/api/internal/httpapi/availability_timeout_integration_test.go +++ b/apps/api/internal/httpapi/availability_timeout_integration_test.go @@ -35,6 +35,24 @@ func TestReadyReturnsPostgresUnavailableWithinTwoSeconds(t *testing.T) { } } +func TestReadyUsesReservedCriticalPoolWhenExecutionPoolIsBusy(t *testing.T) { + executionDB := newExhaustedPostgresStore(t) + criticalDB := newAvailablePostgresStore(t) + server := &Server{ + store: executionDB, + coordinationStore: criticalDB, + logger: slog.New(slog.NewJSONHandler(io.Discard, nil)), + } + request := httptest.NewRequest(http.MethodGet, "/readyz", nil) + recorder := httptest.NewRecorder() + + server.ready(recorder, request) + + if recorder.Code != http.StatusOK { + t.Fatalf("readiness status=%d, want 200 while critical pool is available; body=%s", recorder.Code, recorder.Body.String()) + } +} + func TestLoginReturnsAuthStoreUnavailableWithinFiveSeconds(t *testing.T) { db := newExhaustedPostgresStore(t) var logs bytes.Buffer @@ -93,6 +111,31 @@ func newExhaustedPostgresStore(t *testing.T) *store.Store { return db } +func newAvailablePostgresStore(t *testing.T) *store.Store { + t.Helper() + databaseURL := strings.TrimSpace(os.Getenv("AI_GATEWAY_TEST_DATABASE_URL")) + if databaseURL == "" { + t.Skip("set AI_GATEWAY_TEST_DATABASE_URL to run PostgreSQL availability timeout tests") + } + parsed, err := url.Parse(databaseURL) + if err != nil { + t.Fatalf("parse test database URL: %v", err) + } + query := parsed.Query() + query.Set("pool_max_conns", "1") + query.Set("pool_min_conns", "0") + parsed.RawQuery = query.Encode() + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + db, err := store.Connect(ctx, parsed.String()) + if err != nil { + t.Fatalf("connect available test store: %v", err) + } + t.Cleanup(db.Close) + return db +} + func assertUnavailableResponse(t *testing.T, recorder *httptest.ResponseRecorder, expectedCode, expectedMessage string) { t.Helper() if recorder.Code != http.StatusServiceUnavailable { diff --git a/apps/api/internal/httpapi/handlers.go b/apps/api/internal/httpapi/handlers.go index c82a9a7..1f9373d 100644 --- a/apps/api/internal/httpapi/handlers.go +++ b/apps/api/internal/httpapi/handlers.go @@ -53,7 +53,11 @@ func (s *Server) health(w http.ResponseWriter, r *http.Request) { func (s *Server) ready(w http.ResponseWriter, r *http.Request) { ctx, cancel := context.WithTimeout(r.Context(), postgresReadinessTimeout) defer cancel() - if err := s.store.Ping(ctx); err != nil { + readinessStore := s.coordinationStore + if readinessStore == nil { + readinessStore = s.store + } + if err := readinessStore.Ping(ctx); err != nil { s.logPostgresUnavailable("postgres readiness check failed") writeError(w, http.StatusServiceUnavailable, "postgres unavailable", errorCodePostgresDown) return diff --git a/apps/api/internal/runner/admission.go b/apps/api/internal/runner/admission.go index 6954f9d..0c46349 100644 --- a/apps/api/internal/runner/admission.go +++ b/apps/api/internal/runner/admission.go @@ -146,7 +146,11 @@ func (s *Service) buildTaskAdmissionPlanForCurrentBinding( ModelType: modelType, }, nil } - if err := s.store.CheckRateLimits(ctx, s.rateLimitReservations(ctx, user, candidate, body)); err != nil { + reservations := acceptanceInfrastructureReservations( + task, + s.rateLimitReservations(ctx, user, candidate, body), + ) + if err := s.store.CheckRateLimits(ctx, reservations); err != nil { return taskAdmissionPlan{}, err } return taskAdmissionPlan{ @@ -206,8 +210,17 @@ func acceptanceAdmissionScopes(task store.GatewayTask, scopes []store.AdmissionS } out := append([]store.AdmissionScope(nil), scopes...) for index := range out { + // Protocol-emulated acceptance measures Gateway and Worker capacity, so + // the isolated Run is bounded by the worker_capacity scope instead of a + // production supplier quota. The real acceptance_canary path deliberately + // retains the production platform-model concurrency limit. + if task.RunMode == "acceptance" && out[index].ScopeType == "platform_model" { + out[index].ConcurrentLimit = 0 + } if out[index].ConcurrentLimit <= 0 { - continue + if task.RunMode != "acceptance" || out[index].ScopeType != "platform_model" { + continue + } } out[index].QueueLimit = acceptanceQueueLimit out[index].MaxWaitSeconds = acceptanceQueueMaxWait @@ -215,6 +228,23 @@ func acceptanceAdmissionScopes(task store.GatewayTask, scopes []store.AdmissionS return out } +func acceptanceInfrastructureReservations( + task store.GatewayTask, + reservations []store.RateLimitReservation, +) []store.RateLimitReservation { + if task.RunMode != "acceptance" { + return reservations + } + out := make([]store.RateLimitReservation, 0, len(reservations)) + for _, reservation := range reservations { + if reservation.ScopeType == "platform_model" && reservation.Metric == "concurrent" { + continue + } + out = append(out, reservation) + } + return out +} + func (s *Service) loadAsyncTaskAdmission(ctx context.Context, task store.GatewayTask) (*store.TaskAdmission, error) { if !task.AsyncMode { return nil, nil @@ -420,7 +450,11 @@ func (s *Service) ensureCandidateAdmission( } return store.TaskAdmissionResult{}, false, nil } - if err := s.store.CheckRateLimits(ctx, s.rateLimitReservations(ctx, user, candidate, body)); err != nil { + reservations := acceptanceInfrastructureReservations( + task, + s.rateLimitReservations(ctx, user, candidate, body), + ) + if err := s.store.CheckRateLimits(ctx, reservations); err != nil { return store.TaskAdmissionResult{}, true, err } plan := taskAdmissionPlan{ diff --git a/apps/api/internal/runner/admission_test.go b/apps/api/internal/runner/admission_test.go index c1de5d2..2a7b842 100644 --- a/apps/api/internal/runner/admission_test.go +++ b/apps/api/internal/runner/admission_test.go @@ -24,28 +24,66 @@ func TestDistributedAdmissionModelTypeBoundary(t *testing.T) { } } -func TestAcceptanceAdmissionScopesEnableBoundedQueueWithoutChangingConcurrency(t *testing.T) { +func TestAcceptanceAdmissionScopesUseWorkerCapacityInsteadOfSupplierConcurrency(t *testing.T) { input := []store.AdmissionScope{{ ScopeType: "platform_model", ScopeKey: "model-1", ConcurrentLimit: 10, QueueLimit: 0, MaxWaitSeconds: 0, + }, { + ScopeType: "worker_capacity", + ScopeKey: "global", + ConcurrentLimit: 48, }} got := acceptanceAdmissionScopes(store.GatewayTask{RunMode: "acceptance"}, input) - if got[0].ConcurrentLimit != 10 { - t.Fatalf("acceptance must preserve the production concurrency limit, got %+v", got[0]) + if got[0].ConcurrentLimit != 0 { + t.Fatalf("protocol-emulated acceptance must defer to worker capacity, got %+v", got[0]) } if got[0].QueueLimit != acceptanceQueueLimit || got[0].MaxWaitSeconds != acceptanceQueueMaxWait { t.Fatalf("acceptance must enable a bounded queue, got %+v", got[0]) } + if got[1].ConcurrentLimit != 48 { + t.Fatalf("acceptance worker capacity changed, got %+v", got[1]) + } if input[0].QueueLimit != 0 || input[0].MaxWaitSeconds != 0 { t.Fatalf("acceptance queue overlay must not mutate the production scopes, got %+v", input[0]) } } +func TestAcceptanceCanaryPreservesProductionConcurrency(t *testing.T) { + input := []store.AdmissionScope{{ + ScopeType: "platform_model", + ScopeKey: "model-1", + ConcurrentLimit: 10, + }} + + got := acceptanceAdmissionScopes(store.GatewayTask{RunMode: "acceptance_canary"}, input) + + if got[0].ConcurrentLimit != 10 { + t.Fatalf("real acceptance canary must preserve production concurrency, got %+v", got[0]) + } +} + +func TestAcceptanceInfrastructureReservationsOnlyRemoveSupplierConcurrency(t *testing.T) { + input := []store.RateLimitReservation{ + {ScopeType: "platform_model", Metric: "concurrent", Limit: 10}, + {ScopeType: "platform_model", Metric: "rpm", Limit: 600}, + {ScopeType: "user_group", Metric: "concurrent", Limit: 20}, + } + + got := acceptanceInfrastructureReservations(store.GatewayTask{RunMode: "acceptance"}, input) + if len(got) != 2 || got[0].Metric != "rpm" || got[1].ScopeType != "user_group" { + t.Fatalf("unexpected acceptance reservations: %+v", got) + } + canary := acceptanceInfrastructureReservations(store.GatewayTask{RunMode: "acceptance_canary"}, input) + if len(canary) != len(input) { + t.Fatalf("real canary reservations changed: %+v", canary) + } +} + func TestAcceptanceAdmissionScopesLeaveProductionPolicyUnchanged(t *testing.T) { input := []store.AdmissionScope{{ ScopeType: "platform_model", diff --git a/apps/api/internal/runner/service.go b/apps/api/internal/runner/service.go index 9c66131..d7ce8a3 100644 --- a/apps/api/internal/runner/service.go +++ b/apps/api/internal/runner/service.go @@ -601,7 +601,11 @@ func (s *Service) executeWithToken(ctx context.Context, task store.GatewayTask, } } if hasConcurrentLimit { - if err := s.store.CheckRateLimits(ctx, s.rateLimitReservations(ctx, user, candidates[0], body)); err != nil { + reservations := acceptanceInfrastructureReservations( + task, + s.rateLimitReservations(ctx, user, candidates[0], body), + ) + if err := s.store.CheckRateLimits(ctx, reservations); err != nil { if task.AsyncMode && errors.Is(err, store.ErrRateLimited) && store.RateLimitRetryable(err) { queued, delay, queueErr := s.requeueRateLimitedTask(ctx, task, err, candidates[0]) if queueErr != nil { diff --git a/apps/api/internal/store/acceptance.go b/apps/api/internal/store/acceptance.go index 2f15291..c5a50c6 100644 --- a/apps/api/internal/store/acceptance.go +++ b/apps/api/internal/store/acceptance.go @@ -209,9 +209,6 @@ FOR UPDATE`, SystemSettingGatewayTrafficMode).Scan(¤tValue); err != nil { if err := json.Unmarshal(currentValue, ¤t); err != nil { return GatewayTrafficMode{}, err } - if normalizeTrafficMode(current.Mode) != "live" { - return GatewayTrafficMode{}, ErrAcceptanceStateConflict - } run, err := scanAcceptanceRun(tx.QueryRow(ctx, ` SELECT `+acceptanceRunColumns+` FROM gateway_acceptance_runs @@ -223,6 +220,40 @@ FOR UPDATE`, strings.TrimSpace(runID))) if run.Status != "pending" && run.Status != "failed" { return GatewayTrafficMode{}, ErrAcceptanceStateConflict } + switch normalizeTrafficMode(current.Mode) { + case "live": + case "validation": + if current.RunID == "" || current.RunID == run.ID { + return GatewayTrafficMode{}, ErrAcceptanceStateConflict + } + previous, previousErr := scanAcceptanceRun(tx.QueryRow(ctx, ` +SELECT `+acceptanceRunColumns+` +FROM gateway_acceptance_runs +WHERE id = $1::uuid +FOR UPDATE`, current.RunID)) + if previousErr != nil { + return GatewayTrafficMode{}, previousErr + } + if previous.Status != "failed" || + previous.ReleaseSHA != current.ReleaseSHA || + previous.APIImageDigest != current.APIImageDigest || + previous.WorkerImageDigest != current.WorkerImageDigest { + return GatewayTrafficMode{}, ErrAcceptanceStateConflict + } + var activeTasks int + if err := tx.QueryRow(ctx, ` +SELECT count(*) +FROM gateway_tasks +WHERE acceptance_run_id = $1::uuid + AND status IN ('queued', 'running')`, previous.ID).Scan(&activeTasks); err != nil { + return GatewayTrafficMode{}, err + } + if activeTasks != 0 { + return GatewayTrafficMode{}, ErrAcceptanceStateConflict + } + default: + return GatewayTrafficMode{}, ErrAcceptanceStateConflict + } next := GatewayTrafficMode{ Mode: "validation", RunID: run.ID, diff --git a/apps/api/internal/store/acceptance_integration_test.go b/apps/api/internal/store/acceptance_integration_test.go index 01021b7..94eb708 100644 --- a/apps/api/internal/store/acceptance_integration_test.go +++ b/apps/api/internal/store/acceptance_integration_test.go @@ -323,3 +323,77 @@ FROM gateway_tasks`, t.Fatalf("delete acceptance cleanup billing user: %v", err) } } + +func TestActivateAcceptanceRunAtomicallyReplacesDrainedFailedRun(t *testing.T) { + databaseURL := strings.TrimSpace(os.Getenv("AI_GATEWAY_TEST_DATABASE_URL")) + if databaseURL == "" { + t.Skip("set AI_GATEWAY_TEST_DATABASE_URL to run acceptance traffic integration test") + } + ctx, cancel := context.WithTimeout(context.Background(), time.Minute) + defer cancel() + applyOIDCJITTestMigrations(t, ctx, databaseURL) + db, err := Connect(ctx, databaseURL) + if err != nil { + t.Fatalf("connect store: %v", err) + } + defer db.Close() + if _, err := db.pool.Exec(ctx, ` +UPDATE system_settings +SET value = '{"mode":"live","revision":0}'::jsonb, updated_at = now() +WHERE setting_key = $1`, SystemSettingGatewayTrafficMode); err != nil { + t.Fatalf("reset traffic mode: %v", err) + } + t.Cleanup(func() { + _, _ = db.pool.Exec(context.Background(), ` +UPDATE system_settings +SET value = '{"mode":"live","revision":0}'::jsonb, updated_at = now() +WHERE setting_key = $1`, SystemSettingGatewayTrafficMode) + }) + + createRun := func(releaseCharacter, digestCharacter string) AcceptanceRun { + t.Helper() + run, createErr := db.CreateAcceptanceRun(ctx, CreateAcceptanceRunInput{ + ReleaseSHA: strings.Repeat(releaseCharacter, 40), + APIImageDigest: "sha256:" + strings.Repeat(digestCharacter, 64), + WorkerImageDigest: "sha256:" + strings.Repeat(digestCharacter, 64), + APIKeyID: "atomic-replace-api-key", + UserID: "atomic-replace-user", + Token: strings.Repeat("t", 32), + EmulatorBaseURL: "http://acceptance-emulator:8090", + CallbackURL: "http://acceptance-emulator:8090/callbacks", + CapacityProfile: "P24", + }) + if createErr != nil { + t.Fatalf("create acceptance run: %v", createErr) + } + return run + } + + previous := createRun("1", "2") + firstMode, err := db.ActivateAcceptanceRun(ctx, previous.ID) + if err != nil { + t.Fatalf("activate previous run: %v", err) + } + if _, err := db.FinishAcceptanceRun(ctx, FinishAcceptanceRunInput{ + RunID: previous.ID, Passed: false, FailureReason: "capacity gate failed", + }); err != nil { + t.Fatalf("fail previous run: %v", err) + } + nextRun := createRun("3", "4") + nextMode, err := db.ActivateAcceptanceRun(ctx, nextRun.ID) + if err != nil { + t.Fatalf("atomically replace failed run: %v", err) + } + if firstMode.Mode != "validation" || nextMode.Mode != "validation" || + nextMode.RunID != nextRun.ID || nextMode.Revision != firstMode.Revision+1 { + t.Fatalf("unexpected replacement modes: first=%+v next=%+v", firstMode, nextMode) + } + storedPrevious, err := db.GetAcceptanceRun(ctx, previous.ID) + if err != nil || storedPrevious.Status != "failed" { + t.Fatalf("previous run changed during replacement: run=%+v err=%v", storedPrevious, err) + } + storedNext, err := db.GetAcceptanceRun(ctx, nextRun.ID) + if err != nil || storedNext.Status != "running" { + t.Fatalf("next run was not activated: run=%+v err=%v", storedNext, err) + } +} diff --git a/deploy/kubernetes/easyai-ai-gateway-cluster-release b/deploy/kubernetes/easyai-ai-gateway-cluster-release index 847bfd5..6607002 100755 --- a/deploy/kubernetes/easyai-ai-gateway-cluster-release +++ b/deploy/kubernetes/easyai-ai-gateway-cluster-release @@ -28,7 +28,7 @@ source "$config_file" : "${AI_GATEWAY_DATABASE_RIVER_MAX_CONNS:=8}" : "${AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS:=4}" : "${AI_GATEWAY_DATABASE_MIN_IDLE_CONNS:=4}" -: "${AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS:=30}" +: "${AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS:=300}" : "${AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY:=24}" : "${AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY:=24}" : "${AI_GATEWAY_WORKER_AUTOSCALING_ENABLED:=false}" diff --git a/deploy/kubernetes/easyai-ai-gateway-cluster-release.conf.example b/deploy/kubernetes/easyai-ai-gateway-cluster-release.conf.example index 802e082..ac40d5f 100644 --- a/deploy/kubernetes/easyai-ai-gateway-cluster-release.conf.example +++ b/deploy/kubernetes/easyai-ai-gateway-cluster-release.conf.example @@ -19,7 +19,7 @@ AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=4 AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=8 AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS=4 AI_GATEWAY_DATABASE_MIN_IDLE_CONNS=4 -AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS=30 +AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS=300 AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY=24 AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY=24 AI_GATEWAY_WORKER_AUTOSCALING_ENABLED=false diff --git a/scripts/cluster/run-production-acceptance.sh b/scripts/cluster/run-production-acceptance.sh index fa60385..a49a39f 100755 --- a/scripts/cluster/run-production-acceptance.sh +++ b/scripts/cluster/run-production-acceptance.sh @@ -88,6 +88,7 @@ require_commands curl git go jq node openssl sed shasum : "${AI_GATEWAY_ACCEPTANCE_API_DATABASE_MAX_CONNS:=31}" : "${AI_GATEWAY_ACCEPTANCE_WORKER_RIVER_MAX_CONNS:=8}" : "${AI_GATEWAY_ACCEPTANCE_API_RIVER_MAX_CONNS:=4}" +: "${AI_GATEWAY_ACCEPTANCE_DATABASE_MAX_CONN_IDLE_SECONDS:=300}" : "${AI_GATEWAY_ACCEPTANCE_API_MEDIA_REQUEST_CONCURRENCY:=128}" : "${AI_GATEWAY_ACCEPTANCE_IDENTITY_SHARDS:=32}" : "${AI_GATEWAY_ACCEPTANCE_WORKER_MEMORY_REQUEST_MIB:=1536}" @@ -122,6 +123,11 @@ fi echo 'AI_GATEWAY_ACCEPTANCE_API_DATABASE_MAX_CONNS must be between 1 and 256' >&2 exit 1 } +[[ $AI_GATEWAY_ACCEPTANCE_DATABASE_MAX_CONN_IDLE_SECONDS =~ ^[1-9][0-9]*$ && + $AI_GATEWAY_ACCEPTANCE_DATABASE_MAX_CONN_IDLE_SECONDS -le 3600 ]] || { + echo 'AI_GATEWAY_ACCEPTANCE_DATABASE_MAX_CONN_IDLE_SECONDS must be between 1 and 3600' >&2 + exit 1 +} [[ $AI_GATEWAY_ACCEPTANCE_WORKER_RIVER_MAX_CONNS =~ ^[1-9][0-9]*$ && $AI_GATEWAY_ACCEPTANCE_API_RIVER_MAX_CONNS =~ ^[1-9][0-9]*$ && $AI_GATEWAY_ACCEPTANCE_WORKER_RIVER_MAX_CONNS -lt 28 && @@ -1344,7 +1350,6 @@ create_and_activate_run() { local activate_response=$temporary_root/activate-run.json local traffic_response=$temporary_root/pre-activate-traffic.json local previous_run_response=$temporary_root/previous-run.json - local abort_response=$temporary_root/abort-previous-run.json local body body=$(jq -cn \ --arg releaseSha "$release_sha" \ @@ -1421,24 +1426,9 @@ create_and_activate_run() { echo 'refusing to replace validation mode owned by a non-ancestor release' >&2 return 1 fi - body=$(jq -cn \ - --argjson revision "$(jq -r '.revision' "$traffic_response")" \ - --arg releaseSha "$traffic_release" \ - --arg apiDigest "$traffic_api_digest" \ - --arg workerDigest "$traffic_worker_digest" \ - '{ - revision:$revision, - releaseSha:$releaseSha, - apiImageDigest:$apiDigest, - workerImageDigest:$workerDigest - }') - admin_request POST \ - "/api/admin/system/acceptance/runs/$previous_run_id/abort" \ - "$body" "$abort_response" - [[ $(jq -r '.mode // empty' "$abort_response") == live ]] || { - echo 'previous failed acceptance Run did not return traffic to live' >&2 - return 1 - } + # ActivateAcceptanceRun replaces a drained failed Run while holding the + # traffic-mode row lock. Formal traffic therefore never observes a brief + # live interval between two validation Runs. wait_for_terminal_acceptance_tasks_to_drain elif [[ $traffic_mode != live ]]; then echo "unsupported production traffic mode before acceptance: $traffic_mode" >&2 @@ -1455,13 +1445,23 @@ mark_run_failed() { local response=$temporary_root/failed-run.json local body [[ -n $run_id ]] || return 0 - local artifact_names + local artifact_names task_count gemini_count video_count queue_state + task_count=$(database_query "SELECT count(*) FROM gateway_tasks WHERE acceptance_run_id='$run_id'::uuid;" 2>/dev/null || echo 0) + gemini_count=$(database_query "SELECT count(*) FROM gateway_tasks WHERE acceptance_run_id='$run_id'::uuid AND model='$AI_GATEWAY_ACCEPTANCE_GEMINI_MODEL';" 2>/dev/null || echo 0) + video_count=$(database_query "SELECT count(*) FROM gateway_tasks WHERE acceptance_run_id='$run_id'::uuid AND model='$AI_GATEWAY_ACCEPTANCE_VIDEO_MODEL';" 2>/dev/null || echo 0) + queue_state=$(database_query "SELECT count(*) FILTER (WHERE status='queued')||':'||count(*) FILTER (WHERE status='running') FROM gateway_tasks WHERE acceptance_run_id='$run_id'::uuid;" 2>/dev/null || echo '0:0') + [[ $task_count =~ ^[0-9]+$ ]] || task_count=0 + [[ $gemini_count =~ ^[0-9]+$ ]] || gemini_count=0 + [[ $video_count =~ ^[0-9]+$ ]] || video_count=0 + [[ $queue_state =~ ^[0-9]+:[0-9]+$ ]] || queue_state=0:0 artifact_names=$(find "$report_root" -maxdepth 1 -type f -exec basename {} \; 2>/dev/null | LC_ALL=C sort | jq -Rsc 'split("\n") | map(select(length > 0))') body=$(jq -cn \ --arg reason "$reason" \ --arg gateId "${failure_gate_id:-unknown_gate}" \ --arg activeProfile "$active_profile" \ + --arg queueState "$queue_state" \ + --argjson taskCount "$task_count" \ --argjson artifacts "$artifact_names" \ '{ passed:false, @@ -1469,7 +1469,8 @@ mark_run_failed() { report:{ gateId:$gateId, activeCapacityProfile:$activeProfile, - queueFinal:null, + tasks:$taskCount, + queueAtFailure:$queueState, completedArtifacts:$artifacts } }') @@ -1482,6 +1483,10 @@ mark_run_failed() { --arg stableCapacityProfile "$stable_profile" \ --arg gateId "${failure_gate_id:-unknown_gate}" \ --arg reason "$reason" \ + --arg queueState "$queue_state" \ + --argjson taskCount "$task_count" \ + --argjson geminiCount "$gemini_count" \ + --argjson videoCount "$video_count" \ --argjson artifacts "$artifact_names" \ '{ runId:$runId, @@ -1490,9 +1495,10 @@ mark_run_failed() { snapshotConfigHash:$snapshotConfigHash, snapshotSha256:$snapshotSha256, stableCapacityProfile:$stableCapacityProfile, - tasks:null, - geminiTasks:null, - videoTasks:null, + tasks:$taskCount, + geminiTasks:$geminiCount, + videoTasks:$videoCount, + queueAtFailure:$queueState, realCanaryPassed:false, failureReason:$reason, completedArtifacts:$artifacts, @@ -1530,7 +1536,7 @@ apply_capacity_profile() { local slots pool media global baseline_replicas local api_pool=$AI_GATEWAY_ACCEPTANCE_API_DATABASE_MAX_CONNS local min_idle=4 - local max_idle_seconds=30 + local max_idle_seconds=$AI_GATEWAY_ACCEPTANCE_DATABASE_MAX_CONN_IDLE_SECONDS case $profile in P24) slots=24; pool=32; media=24 ;; P28) slots=28; pool=36; media=28 ;; @@ -1715,7 +1721,7 @@ apply_autoscaling_profile() { global=$((slots * (max_replicas_ningbo + max_replicas_hongkong))) local command_text capacity_output printf -v command_text \ - 'AI_GATEWAY_WORKER_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MIN_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MIN_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MAX_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MAX_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_AUTOSCALING_ENABLED=true AI_GATEWAY_WORKER_SCALE_UP_WINDOW_SECONDS=20 AI_GATEWAY_WORKER_SCALE_DOWN_STABILIZATION_SECONDS=600 AI_GATEWAY_WORKER_DRAIN_TIMEOUT_SECONDS=600 AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=%q AI_GATEWAY_DATABASE_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_MAX_CONNS=%q AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=4 AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_DATABASE_MIN_IDLE_CONNS=4 AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS=30 AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY=%q AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY=%q AI_GATEWAY_WORKER_TARGET_OUTSTANDING_PER_REPLICA=%q AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB=%q AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES=%q %q capacity' \ + 'AI_GATEWAY_WORKER_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MIN_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MIN_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_MAX_REPLICAS_NINGBO=%q AI_GATEWAY_WORKER_MAX_REPLICAS_HONGKONG=%q AI_GATEWAY_WORKER_AUTOSCALING_ENABLED=true AI_GATEWAY_WORKER_SCALE_UP_WINDOW_SECONDS=20 AI_GATEWAY_WORKER_SCALE_DOWN_STABILIZATION_SECONDS=600 AI_GATEWAY_WORKER_DRAIN_TIMEOUT_SECONDS=600 AI_GATEWAY_ASYNC_WORKER_INSTANCE_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT=%q AI_GATEWAY_ASYNC_WORKER_GLOBAL_HARD_LIMIT=%q AI_GATEWAY_DATABASE_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_MAX_CONNS=%q AI_GATEWAY_DATABASE_CRITICAL_MAX_CONNS=4 AI_GATEWAY_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_API_DATABASE_RIVER_MAX_CONNS=%q AI_GATEWAY_DATABASE_MIN_IDLE_CONNS=4 AI_GATEWAY_DATABASE_MAX_CONN_IDLE_SECONDS=%q AI_GATEWAY_MEDIA_MATERIALIZATION_CONCURRENCY=%q AI_GATEWAY_MEDIA_REQUEST_CONCURRENCY=%q AI_GATEWAY_WORKER_TARGET_OUTSTANDING_PER_REPLICA=%q AI_GATEWAY_WORKER_MEMORY_REQUEST_MIB=%q AI_GATEWAY_WORKER_CPU_REQUEST_MILLICORES=%q %q capacity' \ "$AI_GATEWAY_ACCEPTANCE_AUTOSCALING_MIN_REPLICAS_NINGBO" \ "$AI_GATEWAY_ACCEPTANCE_AUTOSCALING_MIN_REPLICAS_HONGKONG" \ "$AI_GATEWAY_ACCEPTANCE_AUTOSCALING_MIN_REPLICAS_NINGBO" \ @@ -1726,6 +1732,7 @@ apply_autoscaling_profile() { "$AI_GATEWAY_ACCEPTANCE_API_DATABASE_MAX_CONNS" \ "$AI_GATEWAY_ACCEPTANCE_WORKER_RIVER_MAX_CONNS" \ "$AI_GATEWAY_ACCEPTANCE_API_RIVER_MAX_CONNS" \ + "$AI_GATEWAY_ACCEPTANCE_DATABASE_MAX_CONN_IDLE_SECONDS" \ "$media" "$AI_GATEWAY_ACCEPTANCE_API_MEDIA_REQUEST_CONCURRENCY" "$((slots * 2))" \ "$AI_GATEWAY_ACCEPTANCE_WORKER_MEMORY_REQUEST_MIB" \ "$AI_GATEWAY_ACCEPTANCE_WORKER_CPU_REQUEST_MILLICORES" \ @@ -2560,7 +2567,7 @@ workload_metrics_for_site() { fi ;; esac - cluster_ssh "$host" "curl -fsS --max-time 5 http://$pod_ip:8088/metrics" \ + cluster_ssh "$host" "curl -fsS --max-time 10 http://$pod_ip:8088/metrics" \