From f190af00be5847def7de591c7be12a0b281604db Mon Sep 17 00:00:00 2001 From: wangbo Date: Fri, 31 Jul 2026 20:39:42 +0800 Subject: [PATCH] =?UTF-8?q?fix(acceptance):=20=E4=BF=AE=E5=A4=8D=E5=BF=AB?= =?UTF-8?q?=E7=85=A7=E8=BF=81=E7=A7=BB=E9=87=8D=E6=94=BE=E4=B8=8E=E6=B4=BE?= =?UTF-8?q?=E5=8F=91=E6=AD=BB=E9=94=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 本地同构验收在全量迁移后导入生产快照,导致本次发布的数据迁移被旧能力覆盖;增加仅允许严格本地集群标记启用的导入后迁移重放,并新增幂等 Seedance 约束校准。\n\n批量异步派发改为在事务开始按全局顺序预锁全部任务与容量作用域,同时对 PostgreSQL 死锁和序列化失败做退避重试,避免双 API 重叠批次形成环形等待。\n\n验证:Go 全量测试、gofmt、bash -n、ShellCheck、迁移安全检查。 --- apps/api/cmd/migrate/main.go | 74 ++++++++++---- apps/api/cmd/migrate/main_test.go | 36 +++++++ apps/api/internal/store/admission_lock.go | 14 ++- .../api/internal/store/admission_lock_test.go | 15 +++ apps/api/internal/store/admission_queue.go | 10 ++ ...lces_seedance20_acceptance_constraints.sql | 96 +++++++++++++++++++ scripts/acceptance/local-cluster.sh | 31 ++++++ 7 files changed, 259 insertions(+), 17 deletions(-) create mode 100644 apps/api/migrations/0099_reconcile_volces_seedance20_acceptance_constraints.sql diff --git a/apps/api/cmd/migrate/main.go b/apps/api/cmd/migrate/main.go index 77a7455..0978110 100644 --- a/apps/api/cmd/migrate/main.go +++ b/apps/api/cmd/migrate/main.go @@ -6,6 +6,7 @@ import ( "log/slog" "os" "path/filepath" + "regexp" "sort" "strings" @@ -14,10 +15,15 @@ import ( ) const ( - noTransactionMigrationMarker = "-- easyai:migration:no-transaction" - migrationStatementSeparator = "-- easyai:migration:statement" + noTransactionMigrationMarker = "-- easyai:migration:no-transaction" + migrationStatementSeparator = "-- easyai:migration:statement" + acceptanceImportReplayMarker = "-- easyai:migration:reapply-after-acceptance-import" + acceptanceImportReplayEnvironment = "AI_GATEWAY_MIGRATION_REAPPLY_ACCEPTANCE_IMPORT" + acceptanceLocalClusterSettingKey = "acceptance_local_cluster_id" ) +var acceptanceLocalClusterMarkerPattern = regexp.MustCompile(`^easyai-local-[0-9a-f]{24}$`) + func main() { cfg := config.Load() logger := slog.New(slog.NewTextHandler(os.Stdout, nil)) @@ -42,6 +48,21 @@ CREATE TABLE IF NOT EXISTS schema_migrations ( logger.Error("ensure schema_migrations failed", "error", err) os.Exit(1) } + reapplyAcceptanceImport := strings.EqualFold(strings.TrimSpace(os.Getenv(acceptanceImportReplayEnvironment)), "true") + if reapplyAcceptanceImport { + var localClusterID string + if err := conn.QueryRow(ctx, ` +SELECT COALESCE(value->>'clusterId', '') +FROM system_settings +WHERE setting_key = $1`, acceptanceLocalClusterSettingKey).Scan(&localClusterID); err != nil { + logger.Error("verify local acceptance database failed", "error", err) + os.Exit(1) + } + if !acceptanceLocalClusterMarkerPattern.MatchString(localClusterID) { + logger.Error("refusing acceptance migration replay outside a marked local cluster") + os.Exit(1) + } + } files, err := filepath.Glob("migrations/*.sql") if err != nil { @@ -57,16 +78,16 @@ CREATE TABLE IF NOT EXISTS schema_migrations ( logger.Error("check migration failed", "version", version, "error", err) os.Exit(1) } - if exists { - logger.Info("migration skipped", "version", version) - continue - } - sqlBytes, err := os.ReadFile(file) if err != nil { logger.Error("read migration file failed", "file", file, "error", err) os.Exit(1) } + replaying := exists && reapplyAcceptanceImport && hasMigrationMarker(string(sqlBytes), acceptanceImportReplayMarker) + if exists && !replaying { + logger.Info("migration skipped", "version", version) + continue + } noTransaction, statements := migrationStatements(string(sqlBytes)) if noTransaction { @@ -76,11 +97,17 @@ CREATE TABLE IF NOT EXISTS schema_migrations ( os.Exit(1) } } - if _, err := conn.Exec(ctx, "INSERT INTO schema_migrations(version) VALUES($1)", version); err != nil { - logger.Error("record non-transaction migration failed", "version", version, "error", err) - os.Exit(1) + if !exists { + if _, err := conn.Exec(ctx, "INSERT INTO schema_migrations(version) VALUES($1)", version); err != nil { + logger.Error("record non-transaction migration failed", "version", version, "error", err) + os.Exit(1) + } + } + if replaying { + logger.Info("acceptance migration replayed", "version", version, "transactional", false) + } else { + logger.Info("migration applied", "version", version, "transactional", false) } - logger.Info("migration applied", "version", version, "transactional", false) continue } @@ -102,21 +129,36 @@ CREATE TABLE IF NOT EXISTS schema_migrations ( logger.Error("execute migration failed", "version", version, "error", err) os.Exit(1) } - if _, err := tx.Exec(ctx, "INSERT INTO schema_migrations(version) VALUES($1)", version); err != nil { - _ = tx.Rollback(ctx) - logger.Error("record migration failed", "version", version, "error", err) - os.Exit(1) + if !exists { + if _, err := tx.Exec(ctx, "INSERT INTO schema_migrations(version) VALUES($1)", version); err != nil { + _ = tx.Rollback(ctx) + logger.Error("record migration failed", "version", version, "error", err) + os.Exit(1) + } } if err := tx.Commit(ctx); err != nil { logger.Error("commit migration failed", "version", version, "error", err) os.Exit(1) } - logger.Info("migration applied", "version", version) + if replaying { + logger.Info("acceptance migration replayed", "version", version) + } else { + logger.Info("migration applied", "version", version) + } } fmt.Println("migrations complete") } +func hasMigrationMarker(sql string, marker string) bool { + for _, line := range strings.Split(sql, "\n") { + if strings.TrimSpace(line) == marker { + return true + } + } + return false +} + func migrationStatements(sql string) (bool, []string) { trimmed := strings.TrimSpace(sql) if !strings.HasPrefix(trimmed, noTransactionMigrationMarker) { diff --git a/apps/api/cmd/migrate/main_test.go b/apps/api/cmd/migrate/main_test.go index 7e66538..4f8d7c2 100644 --- a/apps/api/cmd/migrate/main_test.go +++ b/apps/api/cmd/migrate/main_test.go @@ -32,6 +32,21 @@ CREATE INDEX CONCURRENTLY IF NOT EXISTS second_index ON second_table(id); } } +func TestAcceptanceImportReplayMarkerRequiresAnExactCommentLine(t *testing.T) { + if !hasMigrationMarker( + "-- preface\n"+acceptanceImportReplayMarker+"\nSELECT 1;\n", + acceptanceImportReplayMarker, + ) { + t.Fatal("expected exact acceptance import replay marker") + } + if hasMigrationMarker( + "-- mentions "+acceptanceImportReplayMarker+" in prose\nSELECT 1;\n", + acceptanceImportReplayMarker, + ) { + t.Fatal("prose mention must not enable acceptance import replay") + } +} + func TestSeedanceInputImageConstraintMigrationKeepsCatalogSnapshotsInSync(t *testing.T) { payload, err := os.ReadFile("../../migrations/0078_seedance_input_image_constraints.sql") if err != nil { @@ -74,6 +89,27 @@ func TestVolcesSeedanceInputImageConstraintMigrationUsesDocumentedBounds(t *test } } +func TestVolcesSeedanceAcceptanceReconciliationIsReplayable(t *testing.T) { + payload, err := os.ReadFile("../../migrations/0099_reconcile_volces_seedance20_acceptance_constraints.sql") + if err != nil { + t.Fatal(err) + } + content := string(payload) + for _, required := range []string{ + acceptanceImportReplayMarker, + "volces:doubao-seedance-2-0-260128", + "volces:doubao-seedance-2-0-fast-260128", + "volces:doubao-seedance-2-0-mini-260615", + `"long_edge":300`, + `"long_edge":6000`, + `'[0.4,2.5]'::jsonb`, + } { + if !strings.Contains(content, required) { + t.Fatalf("Volces acceptance reconciliation is missing %q", required) + } + } +} + func TestSecurityEventSchemaMigrationsDefineCurrentLifecycle(t *testing.T) { streamPayload, err := os.ReadFile("../../migrations/0063_oidc_security_events.sql") if err != nil { diff --git a/apps/api/internal/store/admission_lock.go b/apps/api/internal/store/admission_lock.go index 02d4bcb..cf433ec 100644 --- a/apps/api/internal/store/admission_lock.go +++ b/apps/api/internal/store/admission_lock.go @@ -8,6 +8,7 @@ import ( "time" "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" ) const ( @@ -153,7 +154,7 @@ func retryAdmissionOperation[T any]( } result, operationErr := operation() release() - if !errors.Is(operationErr, errAdmissionLockBusy) { + if !isRetryableAdmissionTransactionError(operationErr) { return result, operationErr } if err := waitAdmissionLockRetry(ctx, attempt); err != nil { @@ -162,6 +163,17 @@ func retryAdmissionOperation[T any]( } } +func isRetryableAdmissionTransactionError(err error) bool { + if errors.Is(err, errAdmissionLockBusy) { + return true + } + var postgresError *pgconn.PgError + if !errors.As(err, &postgresError) { + return false + } + return postgresError.Code == "40P01" || postgresError.Code == "40001" +} + func waitAdmissionLockRetry(ctx context.Context, attempt int) error { delay := admissionLockRetryMin for index := 0; index < attempt && delay < admissionLockRetryMax; index++ { diff --git a/apps/api/internal/store/admission_lock_test.go b/apps/api/internal/store/admission_lock_test.go index 99e3b83..e297e67 100644 --- a/apps/api/internal/store/admission_lock_test.go +++ b/apps/api/internal/store/admission_lock_test.go @@ -11,8 +11,23 @@ import ( "github.com/google/uuid" "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" ) +func TestAdmissionTransactionRetryClassification(t *testing.T) { + for _, code := range []string{"40P01", "40001"} { + if !isRetryableAdmissionTransactionError(&pgconn.PgError{Code: code}) { + t.Fatalf("PostgreSQL error %s should be retried", code) + } + } + if !isRetryableAdmissionTransactionError(errAdmissionLockBusy) { + t.Fatal("admission lock timeout should be retried") + } + if isRetryableAdmissionTransactionError(&pgconn.PgError{Code: "23505"}) { + t.Fatal("unique violations must not be retried") + } +} + func TestAdmissionLocalLockSetSerializesSharedKeys(t *testing.T) { lockSet := newAdmissionLocalLockSet() ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) diff --git a/apps/api/internal/store/admission_queue.go b/apps/api/internal/store/admission_queue.go index ad34767..b586bda 100644 --- a/apps/api/internal/store/admission_queue.go +++ b/apps/api/internal/store/admission_queue.go @@ -474,12 +474,22 @@ func (s *Store) TryTaskAdmissionBatchWithAdmittedHook( } lockKeys = append(lockKeys, admissionOperationLockKeys(input)...) } + lockKeys = normalizedAdmissionLockKeys(lockKeys) return retryAdmissionOperation(ctx, lockKeys, func() ([]TaskAdmissionBatchOutcome, error) { tx, err := s.pool.Begin(ctx) if err != nil { return nil, err } defer rollbackTransaction(tx) + // A batch touches one shared capacity scope and multiple task keys. Lock + // the complete union in one global order before processing any task. If + // two API processes dispatch overlapping FIFO windows, neither can hold a + // scope while waiting on a task key already owned by the other batch. + for _, lockKey := range lockKeys { + if err := tryAdmissionTransactionLock(ctx, tx, lockKey); err != nil { + return nil, err + } + } outcomes := make([]TaskAdmissionBatchOutcome, 0, len(inputs)) notify := false diff --git a/apps/api/migrations/0099_reconcile_volces_seedance20_acceptance_constraints.sql b/apps/api/migrations/0099_reconcile_volces_seedance20_acceptance_constraints.sql new file mode 100644 index 0000000..b44991a --- /dev/null +++ b/apps/api/migrations/0099_reconcile_volces_seedance20_acceptance_constraints.sql @@ -0,0 +1,96 @@ +-- easyai:migration:reapply-after-acceptance-import +-- A production snapshot is imported after the local schema is migrated. Replay +-- this idempotent catalog reconciliation so the local target state includes the +-- data changes delivered by the release under test. +WITH input_constraints AS ( + SELECT jsonb_build_object( + 'input_image_resolution_range', + '{"min":{"long_edge":300,"short_edge":300},"max":{"long_edge":6000,"short_edge":6000}}'::jsonb, + 'input_image_aspect_ratio_range', + '[0.4,2.5]'::jsonb + ) AS value +), +target_base_models AS ( + SELECT + base_model.id, + COALESCE(base_model.capabilities, '{}'::jsonb) || jsonb_build_object( + 'image_to_video', + COALESCE(base_model.capabilities->'image_to_video', '{}'::jsonb) || input_constraints.value, + 'omni_video', + COALESCE(base_model.capabilities->'omni_video', '{}'::jsonb) || input_constraints.value + ) AS image_capabilities + FROM base_model_catalog base_model + CROSS JOIN input_constraints + WHERE base_model.provider_key = 'volces' + AND ( + base_model.canonical_model_key IN ( + 'volces:doubao-seedance-2-0-260128', + 'volces:doubao-seedance-2-0-fast-260128', + 'volces:doubao-seedance-2-0-mini-260615' + ) + OR base_model.provider_model_name IN ( + 'doubao-seedance-2-0-260128', + 'doubao-seedance-2-0-fast-260128', + 'doubao-seedance-2-0-mini-260615' + ) + ) +), +updated_base_models AS ( + UPDATE base_model_catalog base_model + SET capabilities = target.image_capabilities, + metadata = jsonb_set( + COALESCE(base_model.metadata, '{}'::jsonb) || jsonb_build_object( + 'rawModel', + COALESCE(base_model.metadata->'rawModel', '{}'::jsonb) + ), + '{rawModel,capabilities}', + target.image_capabilities - 'originalTypes', + true + ), + default_snapshot = CASE + WHEN COALESCE(base_model.default_snapshot, '{}'::jsonb) = '{}'::jsonb + THEN base_model.default_snapshot + ELSE jsonb_set( + jsonb_set( + base_model.default_snapshot || jsonb_build_object( + 'metadata', + COALESCE(base_model.default_snapshot->'metadata', '{}'::jsonb) || jsonb_build_object( + 'rawModel', + COALESCE(base_model.default_snapshot#>'{metadata,rawModel}', '{}'::jsonb) + ) + ), + '{capabilities}', + target.image_capabilities, + true + ), + '{metadata,rawModel,capabilities}', + target.image_capabilities - 'originalTypes', + true + ) + END, + updated_at = now() + FROM target_base_models target + WHERE base_model.id = target.id + RETURNING base_model.id +) +UPDATE platform_models platform_model +SET capabilities = COALESCE(platform_model.capabilities, '{}'::jsonb) || jsonb_build_object( + 'image_to_video', + COALESCE(platform_model.capabilities->'image_to_video', '{}'::jsonb) || input_constraints.value, + 'omni_video', + COALESCE(platform_model.capabilities->'omni_video', '{}'::jsonb) || input_constraints.value + ), + updated_at = now() +FROM integration_platforms platform +CROSS JOIN input_constraints +WHERE platform_model.platform_id = platform.id + AND platform.deleted_at IS NULL + AND platform.provider = 'volces' + AND ( + platform_model.base_model_id IN (SELECT id FROM updated_base_models) + OR COALESCE(NULLIF(platform_model.provider_model_name, ''), platform_model.model_name) IN ( + 'doubao-seedance-2-0-260128', + 'doubao-seedance-2-0-fast-260128', + 'doubao-seedance-2-0-mini-260615' + ) + ); diff --git a/scripts/acceptance/local-cluster.sh b/scripts/acceptance/local-cluster.sh index f806cf2..348a8e6 100755 --- a/scripts/acceptance/local-cluster.sh +++ b/scripts/acceptance/local-cluster.sh @@ -357,6 +357,36 @@ EOF --for=condition=complete job/easyai-local-snapshot-import --timeout=5m >/dev/null } +replay_acceptance_import_migrations() { + local api_image=$1 + kubectl --context "$context" -n "$namespace" delete job easyai-local-migration-replay \ + --ignore-not-found --wait=true >/dev/null + cat </dev/null +apiVersion: batch/v1 +kind: Job +metadata: + name: easyai-local-migration-replay + namespace: easyai +spec: + backoffLimit: 0 + template: + spec: + restartPolicy: Never + containers: + - name: migrate + image: $api_image + command: ["/bin/sh", "-ec", "cd /app && exec /app/easyai-ai-gateway-migrate"] + env: + - name: AI_GATEWAY_MIGRATION_REAPPLY_ACCEPTANCE_IMPORT + value: "true" + envFrom: + - secretRef: + name: easyai-ai-gateway-runtime +EOF + kubectl --context "$context" -n "$namespace" wait \ + --for=condition=complete job/easyai-local-migration-replay --timeout=5m >/dev/null +} + render_and_apply_application() { local api_image=$1 web_image=$2 rendered=$state_root/application.rendered.yaml sed \ @@ -564,6 +594,7 @@ up_cluster() { run_migrations "$api_image" mark_local_database import_snapshot "$api_image" "$snapshot" + replay_acceptance_import_migrations "$api_image" sed "s|image: easyai-api|image: $api_image|g; s|image: easyai-acceptance-netem|image: $netem_image|g" \ "$manifest_root/support-services.yaml" |