From add99e542126442bef3fb0a6a97fd86525be10c1 Mon Sep 17 00:00:00 2001 From: wangbo Date: Fri, 31 Jul 2026 21:25:42 +0800 Subject: [PATCH] =?UTF-8?q?fix(acceptance):=20=E7=89=A9=E5=8C=96=E9=AA=8C?= =?UTF-8?q?=E6=94=B6=E4=BB=BB=E5=8A=A1=E7=9A=84=E6=9C=80=E7=BB=88=E5=AA=92?= =?UTF-8?q?=E4=BD=93?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 仅对已认证 Acceptance Run 强制下载并持久化 URL 型结果,确保集群外压测端能验证最终视频;普通生产任务继续遵循现有文件存储策略。压测器仅对 Gateway 自有地址做双入口重写,避免向第三方媒体域名泄露验收凭据。\n\n同时修复本地集群重复使用同一快照文件时的幂等复制失败。验证:Go 全量测试、go vet、gofmt、bash -n、ShellCheck 和 git diff --check 通过。 --- apps/api/cmd/acceptance-load/main.go | 80 +++++++++++++++++++++-- apps/api/cmd/acceptance-load/main_test.go | 35 +++++++++- apps/api/internal/runner/service.go | 2 +- apps/api/internal/runner/upload.go | 14 +++- apps/api/internal/runner/upload_test.go | 11 ++++ scripts/acceptance/local-cluster.sh | 4 +- 6 files changed, 136 insertions(+), 10 deletions(-) diff --git a/apps/api/cmd/acceptance-load/main.go b/apps/api/cmd/acceptance-load/main.go index 045ee78..cb8f27f 100644 --- a/apps/api/cmd/acceptance-load/main.go +++ b/apps/api/cmd/acceptance-load/main.go @@ -760,7 +760,7 @@ func pollVideoTask(ctx context.Context, client *http.Client, opts options, taskI if mediaURL == "" { return fmt.Errorf("video task %s succeeded without a media URL", taskID) } - return validateVideoAsset(ctx, client, mediaURL) + return validateVideoAsset(ctx, client, opts, mediaURL, index) case "failed", "cancelled", "canceled": return fmt.Errorf("video task %s finished with status %s", taskID, status) } @@ -780,15 +780,34 @@ func pollVideoTask(ctx context.Context, client *http.Client, opts options, taskI func findMediaURL(value any) string { switch typed := value.(type) { case map[string]any: - for key, item := range typed { + preferredKeys := []string{"video_url", "url", "data", "content", "output", "result", "upload"} + visited := make(map[string]struct{}, len(typed)) + for _, key := range preferredKeys { + item, ok := typed[key] + if !ok { + continue + } + visited[key] = struct{}{} normalized := strings.ToLower(strings.TrimSpace(key)) - if (normalized == "url" || normalized == "video_url") && strings.HasPrefix(strings.TrimSpace(fmt.Sprint(item)), "http") { + if (normalized == "url" || normalized == "video_url") && acceptanceMediaURL(strings.TrimSpace(fmt.Sprint(item))) { return strings.TrimSpace(fmt.Sprint(item)) } if mediaURL := findMediaURL(item); mediaURL != "" { return mediaURL } } + keys := make([]string, 0, len(typed)-len(visited)) + for key := range typed { + if _, ok := visited[key]; !ok { + keys = append(keys, key) + } + } + sort.Strings(keys) + for _, key := range keys { + if mediaURL := findMediaURL(typed[key]); mediaURL != "" { + return mediaURL + } + } case []any: for _, item := range typed { if mediaURL := findMediaURL(item); mediaURL != "" { @@ -799,11 +818,22 @@ func findMediaURL(value any) string { return "" } -func validateVideoAsset(ctx context.Context, client *http.Client, mediaURL string) error { - request, err := http.NewRequestWithContext(ctx, http.MethodGet, mediaURL, nil) +func acceptanceMediaURL(value string) bool { + return strings.HasPrefix(value, "/") || strings.HasPrefix(value, "http://") || strings.HasPrefix(value, "https://") +} + +func validateVideoAsset(ctx context.Context, client *http.Client, opts options, mediaURL string, index int) error { + requestURL, gatewayRequest, err := acceptanceMediaRequestURL(opts, mediaURL, index) + if err != nil { + return fmt.Errorf("resolve final video: %w", err) + } + request, err := http.NewRequestWithContext(ctx, http.MethodGet, requestURL, nil) if err != nil { return err } + if gatewayRequest { + opts.setHeaders(request, index, false) + } response, err := client.Do(request) if err != nil { return fmt.Errorf("download final video: %w", err) @@ -824,6 +854,46 @@ func validateVideoAsset(ctx context.Context, client *http.Client, mediaURL strin return nil } +func acceptanceMediaRequestURL(opts options, mediaURL string, index int) (string, bool, error) { + mediaURL = strings.TrimSpace(mediaURL) + if mediaURL == "" { + return "", false, errors.New("empty media URL") + } + if strings.HasPrefix(mediaURL, "/") { + if len(opts.gateways) == 0 { + return "", false, errors.New("no gateway URL configured") + } + base, err := url.Parse(opts.gateways[index%len(opts.gateways)]) + if err != nil { + return "", false, err + } + reference, err := url.Parse(mediaURL) + if err != nil { + return "", false, err + } + return base.ResolveReference(reference).String(), true, nil + } + parsed, err := url.Parse(mediaURL) + if err != nil { + return "", false, err + } + if opts.gatewayTLSName == "" || !strings.EqualFold(parsed.Hostname(), opts.gatewayTLSName) { + return parsed.String(), false, nil + } + if len(opts.gateways) == 0 { + return "", false, errors.New("no gateway URL configured") + } + base, err := url.Parse(opts.gateways[index%len(opts.gateways)]) + if err != nil { + return "", false, err + } + base.Path = parsed.Path + base.RawPath = parsed.RawPath + base.RawQuery = parsed.RawQuery + base.Fragment = parsed.Fragment + return base.String(), true, nil +} + func (o options) setHeaders(request *http.Request, requestIndex int, realUpstream bool) { if o.gatewayTLSName != "" { request.Host = o.gatewayTLSName diff --git a/apps/api/cmd/acceptance-load/main_test.go b/apps/api/cmd/acceptance-load/main_test.go index 46c2a37..bcb4432 100644 --- a/apps/api/cmd/acceptance-load/main_test.go +++ b/apps/api/cmd/acceptance-load/main_test.go @@ -158,12 +158,15 @@ func TestGeminiLoadIsSplitAcrossTwoGatewayAPIs(t *testing.T) { } func TestValidateVideoAssetDownloadsFinalMedia(t *testing.T) { - server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, request *http.Request) { + if request.Header.Get("Authorization") != "" { + t.Fatal("acceptance credentials leaked to external media host") + } w.Header().Set("Content-Type", "video/mp4") _, _ = w.Write([]byte{0, 0, 0, 16, 'f', 't', 'y', 'p', 'i', 's', 'o', 'm'}) })) defer server.Close() - if err := validateVideoAsset(t.Context(), server.Client(), server.URL+"/result.mp4"); err != nil { + if err := validateVideoAsset(t.Context(), server.Client(), options{}, server.URL+"/result.mp4", 0); err != nil { t.Fatalf("validate video asset: %v", err) } if got := findMediaURL(map[string]any{"content": map[string]any{"video_url": server.URL}}); got != server.URL { @@ -171,6 +174,34 @@ func TestValidateVideoAssetDownloadsFinalMedia(t *testing.T) { } } +func TestValidateVideoAssetUsesGatewayForMaterializedPath(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, request *http.Request) { + if request.Host != "gateway.easyai.local" || request.Header.Get("Authorization") != "Bearer key-1" { + t.Fatalf("host=%q authorization=%q", request.Host, request.Header.Get("Authorization")) + } + if request.URL.Path != "/static/generated/result.mp4" { + t.Fatalf("path=%q", request.URL.Path) + } + w.Header().Set("Content-Type", "video/mp4") + _, _ = w.Write([]byte{0, 0, 0, 16, 'f', 't', 'y', 'p', 'i', 's', 'o', 'm'}) + })) + defer server.Close() + opts := options{ + gateways: []string{server.URL}, apiKeys: []string{"key-1"}, + runID: "run-1", runToken: "token-1", gatewayTLSName: "gateway.easyai.local", + } + if err := validateVideoAsset(t.Context(), server.Client(), opts, "/static/generated/result.mp4", 0); err != nil { + t.Fatalf("validate materialized video: %v", err) + } + got := findMediaURL(map[string]any{ + "raw": map[string]any{"video_url": "http://internal.invalid/video.mp4"}, + "data": []any{map[string]any{"video_url": "/static/generated/result.mp4"}}, + }) + if got != "/static/generated/result.mp4" { + t.Fatalf("preferred media URL=%q", got) + } +} + func TestAcceptanceReportErrorRedactsSecretsAndURLs(t *testing.T) { got := redactError( `token-1 failed at https://example.invalid/video.mp4?token=signed`, diff --git a/apps/api/internal/runner/service.go b/apps/api/internal/runner/service.go index 7839bfe..1c59203 100644 --- a/apps/api/internal/runner/service.go +++ b/apps/api/internal/runner/service.go @@ -1605,7 +1605,7 @@ func (s *Service) runCandidate( } defer releaseMediaSlot() } - uploadedResult, err := s.uploadGeneratedAssets(ctx, task.ID, task.Kind, response.Result) + uploadedResult, err := s.uploadGeneratedAssets(ctx, task.ID, task.Kind, task.AcceptanceRunID, response.Result) if err != nil { metrics := mergeMetrics(taskMetrics(task, user, body, candidate, response, simulated), parameterPreprocessingMetrics(preprocessing), map[string]any{ "error": err.Error(), diff --git a/apps/api/internal/runner/upload.go b/apps/api/internal/runner/upload.go index 2ef1b34..8025fac 100644 --- a/apps/api/internal/runner/upload.go +++ b/apps/api/internal/runner/upload.go @@ -112,7 +112,7 @@ func mediaTaskNeedsPreUpstreamMaterializationSlot(modelType string) bool { } } -func (s *Service) uploadGeneratedAssets(ctx context.Context, taskID string, taskKind string, result map[string]any) (map[string]any, error) { +func (s *Service) uploadGeneratedAssets(ctx context.Context, taskID string, taskKind string, acceptanceRunID string, result map[string]any) (map[string]any, error) { data, _ := result["data"].([]any) rawNeedsUpload := generatedRawValueHasInlineMedia(result["raw"], "", nil) hasInlineBinary := TaskResultHasInlineBinary(result) @@ -124,6 +124,11 @@ func (s *Service) uploadGeneratedAssets(ctx context.Context, taskID string, task if err != nil { return nil, &clients.ClientError{Code: "upload_config_failed", Message: err.Error(), Retryable: true} } + // Acceptance runs must prove that a final URL returned by the protocol + // emulator can be persisted and consumed outside the cluster. Keep this + // override scoped to the authenticated run so normal production tasks retain + // the configured file-storage policy. + policy = generatedAssetUploadPolicyForAcceptanceRun(policy, acceptanceRunID) if policy.LocalizeInlineMedia { next, _, err := s.materializeLocalBinaryResult(ctx, taskID, result) if err != nil { @@ -260,6 +265,13 @@ func (s *Service) uploadGeneratedAssets(ctx context.Context, taskID string, task return s.finalizeGeneratedAssets(ctx, taskID, taskKind, next, policy, channels, channelsLoaded, len(nextData)) } +func generatedAssetUploadPolicyForAcceptanceRun(policy generatedAssetUploadPolicy, acceptanceRunID string) generatedAssetUploadPolicy { + if strings.TrimSpace(acceptanceRunID) != "" { + policy.UploadURLMedia = true + } + return policy +} + func (s *Service) finalizeGeneratedAssets( ctx context.Context, taskID string, diff --git a/apps/api/internal/runner/upload_test.go b/apps/api/internal/runner/upload_test.go index f731e47..0bf7778 100644 --- a/apps/api/internal/runner/upload_test.go +++ b/apps/api/internal/runner/upload_test.go @@ -290,6 +290,17 @@ func TestGeneratedAssetUploadPolicyFromName(t *testing.T) { } } +func TestAcceptanceRunForcesURLMaterializationOnlyForAcceptanceTask(t *testing.T) { + configured := generatedAssetUploadPolicy{UploadInlineMedia: true} + if got := generatedAssetUploadPolicyForAcceptanceRun(configured, ""); got.UploadURLMedia { + t.Fatal("ordinary task unexpectedly enabled URL materialization") + } + got := generatedAssetUploadPolicyForAcceptanceRun(configured, "acceptance-run-id") + if !got.UploadInlineMedia || !got.UploadURLMedia { + t.Fatalf("acceptance policy=%+v", got) + } +} + func TestFinalizeGeneratedAssetsUploadsNestedInlineBinaryUnderDefaultPolicy(t *testing.T) { storageDir := t.TempDir() service := &Service{cfg: config.Config{LocalGeneratedStorageDir: storageDir}} diff --git a/scripts/acceptance/local-cluster.sh b/scripts/acceptance/local-cluster.sh index 348a8e6..77e6659 100755 --- a/scripts/acceptance/local-cluster.sh +++ b/scripts/acceptance/local-cluster.sh @@ -542,7 +542,9 @@ up_cluster() { cd "$repository_root/apps/api" go run ./cmd/acceptance-snapshot validate --input "$snapshot" ) - cp "$snapshot" "$state_root/snapshot.json" + if [[ ! -e $state_root/snapshot.json || ! $snapshot -ef $state_root/snapshot.json ]]; then + cp "$snapshot" "$state_root/snapshot.json" + fi chmod 0600 "$state_root/snapshot.json" local k3d config native_arch release_sha api_image web_image netem_image api_digest k3d=$(k3d_binary)