From 5c7d6ac9aab194208036ac5592bc1bb3506f3d9e Mon Sep 17 00:00:00 2001 From: wangbo Date: Tue, 4 Aug 2026 14:33:04 +0800 Subject: [PATCH] =?UTF-8?q?fix(runner):=20=E4=BF=AE=E5=A4=8D=E8=BD=AC?= =?UTF-8?q?=E5=AD=98=E5=90=8E=E6=AE=8B=E7=95=99=E4=BA=8C=E8=BF=9B=E5=88=B6?= =?UTF-8?q?=E8=AF=AF=E5=88=A4?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 统一结果二进制探测与上传规则,避免将 thinking_bytes 等上游元数据误判为待转存媒体,同时保留显式媒体字段和签名识别。\n\n补充不含原始内容的安全诊断,并兼容对象存储前缀旧字段;修正真实 OSS 验收脚本使用的正式字段。\n\n验证:Go 全量测试、go vet、迁移安全检查、真实 Gemini 上游响应及阿里云 OSS 转存均通过。 --- apps/api/internal/runner/binary_results.go | 54 ++++++++++++++++++- .../internal/runner/binary_results_test.go | 11 +++- apps/api/internal/runner/object_storage.go | 5 +- .../internal/runner/object_storage_test.go | 17 ++++++ apps/api/internal/runner/upload.go | 37 +++++++------ apps/api/internal/runner/upload_test.go | 38 +++++++++++++ .../acceptance/run-live-object-storage.mjs | 2 +- 7 files changed, 142 insertions(+), 22 deletions(-) diff --git a/apps/api/internal/runner/binary_results.go b/apps/api/internal/runner/binary_results.go index 420abd0..a2b2748 100644 --- a/apps/api/internal/runner/binary_results.go +++ b/apps/api/internal/runner/binary_results.go @@ -10,6 +10,7 @@ import ( "io" "os" "path/filepath" + "sort" "strconv" "strings" "time" @@ -106,6 +107,52 @@ func TaskResultHasInlineBinary(result map[string]any) bool { return localBinaryValueHasPayload(result, "", nil, 0) } +func taskResultInlineBinaryDiagnostics(result map[string]any) []string { + diagnostics := make([]string, 0, 4) + appendInlineBinaryDiagnostics(result, "", nil, "$", 0, &diagnostics) + return diagnostics +} + +func appendInlineBinaryDiagnostics(value any, key string, siblings map[string]any, path string, depth int, diagnostics *[]string) { + if depth >= localBinaryMaxDepth || len(*diagnostics) >= 8 { + return + } + switch typed := value.(type) { + case map[string]any: + if payload, contentType, ok := localBufferObjectBytes(typed); ok { + *diagnostics = append(*diagnostics, fmt.Sprintf("%s buffer bytes=%d contentType=%s", path, len(payload), contentType)) + return + } + keys := make([]string, 0, len(typed)) + for childKey := range typed { + keys = append(keys, childKey) + } + sort.Strings(keys) + for _, childKey := range keys { + appendInlineBinaryDiagnostics(typed[childKey], childKey, typed, path+"."+childKey, depth+1, diagnostics) + } + case []any: + if localBinaryKey(key) { + if payload, ok := bytesFromNumberArray(typed); ok { + *diagnostics = append(*diagnostics, fmt.Sprintf("%s number-array bytes=%d", path, len(payload))) + return + } + } + for index, child := range typed { + appendInlineBinaryDiagnostics(child, key, siblings, fmt.Sprintf("%s[%d]", path, index), depth+1, diagnostics) + } + case []byte: + if len(typed) > 0 { + *diagnostics = append(*diagnostics, fmt.Sprintf("%s bytes=%d", path, len(typed))) + } + case string: + payload, contentType, encoding, ok := localBinaryStringBytes(key, typed, siblings) + if ok { + *diagnostics = append(*diagnostics, fmt.Sprintf("%s string bytes=%d contentType=%s encoding=%s", path, len(payload), contentType, encoding)) + } + } +} + func localBinaryValueHasPayload(value any, key string, siblings map[string]any, depth int) bool { if depth >= localBinaryMaxDepth { return false @@ -518,7 +565,12 @@ func localBinaryStringBytes(key string, value string, siblings map[string]any) ( return nil, "", "", false } contentType := firstNonEmptyString(mediaContentTypeFromItem(siblings), defaultContentTypeForRawMediaKey(key)) - if !strict && contentType == "" { + // Keys such as thinking_bytes and signature_buffer can carry opaque provider + // metadata. Treat them as generated media only when a sibling content type or + // the payload signature proves that they are media/document bytes. Explicit + // Base64 media keys (b64_json, image_data, binary_data_base64, ...) retain the + // strict behavior expected by compatible image protocols. + if contentType == "" && (!strict || !generatedRawDataMediaPayloadKey(key)) { contentType = detectGeneratedAssetContentType(payload) if !generatedContentTypeIsMedia(contentType) && !generatedContentTypeIsDocument(contentType) { return nil, "", "", false diff --git a/apps/api/internal/runner/binary_results_test.go b/apps/api/internal/runner/binary_results_test.go index 18d7a0f..0bc9de1 100644 --- a/apps/api/internal/runner/binary_results_test.go +++ b/apps/api/internal/runner/binary_results_test.go @@ -314,10 +314,13 @@ func TestHistoricalLocalBinaryStorageUnavailableDoesNotRetryProvider(t *testing. func TestTaskResultIgnoresLongOpaqueBase64Metadata(t *testing.T) { opaque := base64.StdEncoding.EncodeToString([]byte(strings.Repeat("opaque metadata ", 400))) - result := map[string]any{"thought_signature": opaque} + result := map[string]any{ + "thought_signature": opaque, + "thinking_bytes": opaque, + } if TaskResultHasInlineBinary(result) { - t.Fatal("opaque Base64 metadata must not be mistaken for generated media") + t.Fatal("opaque Base64 metadata, including bytes/buffer-style keys, must not be mistaken for generated media") } } @@ -328,6 +331,10 @@ func TestTaskResultDetectsLongBase64MediaWithoutSemanticKey(t *testing.T) { if !TaskResultHasInlineBinary(result) { t.Fatal("signature-detected generated media must still be materialized") } + diagnostics := taskResultInlineBinaryDiagnostics(result) + if len(diagnostics) != 1 || !strings.Contains(diagnostics[0], "$.provider_payload") || !strings.Contains(diagnostics[0], "contentType=image/png") { + t.Fatalf("unexpected safe inline binary diagnostics: %+v", diagnostics) + } } func bytesToAny(payload []byte) []any { diff --git a/apps/api/internal/runner/object_storage.go b/apps/api/internal/runner/object_storage.go index 0120bf5..82ec7c1 100644 --- a/apps/api/internal/runner/object_storage.go +++ b/apps/api/internal/runner/object_storage.go @@ -291,7 +291,10 @@ func (a *objectStorageAdapter) objectKey(payload FileUploadPayload) string { now := a.now().UTC() digest := sha256.Sum256(payload.Bytes) extension := uploadFileExtension(payload.ContentType, path.Ext(payload.FileName)) - prefix := strings.Trim(objectStorageConfigString(a.channel.Config, "objectPrefix"), "/") + prefix := strings.Trim(firstNonEmptyString( + objectStorageConfigString(a.channel.Config, "objectPrefix"), + objectStorageConfigString(a.channel.Config, "objectKeyPrefix"), + ), "/") parts := make([]string, 0, 7) if prefix != "" { parts = append(parts, prefix) diff --git a/apps/api/internal/runner/object_storage_test.go b/apps/api/internal/runner/object_storage_test.go index 9bbf315..03bc446 100644 --- a/apps/api/internal/runner/object_storage_test.go +++ b/apps/api/internal/runner/object_storage_test.go @@ -74,6 +74,23 @@ func TestS3ObjectStorageUsesSigV4AndDeterministicObjectKey(t *testing.T) { } } +func TestObjectStorageObjectKeySupportsLegacyObjectKeyPrefix(t *testing.T) { + channel := testObjectStorageChannel("s3", "https://s3.example.com", "s3-legacy-prefix") + delete(channel.Config, "objectPrefix") + channel.Config["objectKeyPrefix"] = "/legacy/gateway/" + adapter, err := newObjectStorageAdapter(channel) + if err != nil { + t.Fatal(err) + } + adapter.now = func() time.Time { return time.Date(2026, time.August, 4, 1, 2, 3, 0, time.UTC) } + objectKey := adapter.objectKey(FileUploadPayload{ + Bytes: []byte("legacy-prefix"), ContentType: "image/png", Scene: store.FileStorageSceneImageResult, + }) + if !strings.HasPrefix(objectKey, "legacy/gateway/image_result/2026/08/04/") { + t.Fatalf("legacy objectKeyPrefix was not honored: %s", objectKey) + } +} + func TestRequestAssetUsesPrivateSignedURLEvenWhenPublicBaseURLExists(t *testing.T) { server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusOK) diff --git a/apps/api/internal/runner/upload.go b/apps/api/internal/runner/upload.go index c20cb9e..6eb7d3a 100644 --- a/apps/api/internal/runner/upload.go +++ b/apps/api/internal/runner/upload.go @@ -306,9 +306,17 @@ func (s *Service) finalizeGeneratedAssets( if !TaskResultHasInlineBinary(next) { return next, nil } + diagnostics := taskResultInlineBinaryDiagnostics(next) + if s.logger != nil { + s.logger.Error("generated result still contains inline binary after object storage materialization", + "taskID", taskID, + "remainingInlineBinary", diagnostics, + ) + } return nil, &clients.ClientError{ Code: "storage_write_failed", Message: "generated binary result could not be written to object storage", + Details: map[string]any{"scope": "gateway_result_storage", "remainingInlineBinary": diagnostics}, StatusCode: http.StatusServiceUnavailable, Retryable: true, } @@ -417,24 +425,19 @@ func (s *Service) uploadGeneratedBinaryValue(ctx context.Context, taskID string, } func generatedRawInlineMediaAsset(key string, value string, siblings map[string]any, taskKind string) (*generatedInlineAsset, bool) { - raw := strings.TrimSpace(value) - if raw == "" { - return nil, false - } - keyLooksLikeMediaPayload := generatedRawDataMediaPayloadKey(key) - contentType := firstNonEmptyString(mediaContentTypeFromItem(siblings), defaultContentTypeForRawMediaKey(key)) - if !keyLooksLikeMediaPayload && !generatedContentTypeIsMedia(contentType) && !generatedContentTypeIsDocument(contentType) && - !strings.HasPrefix(strings.ToLower(raw), "data:") && len(raw) < localBinaryGenericBase64MinLength { - return nil, false - } - payload, payloadContentType, ok, err := inlineMediaPayload(raw, keyLooksLikeMediaPayload) - if err != nil || !ok || len(payload) == 0 { - return nil, false - } - contentType = firstNonEmptyString(payloadContentType, contentType, detectGeneratedAssetContentType(payload)) - if !generatedContentTypeIsMedia(contentType) && !generatedContentTypeIsDocument(contentType) { - return nil, false + payload, contentType, _, ok := localBinaryStringBytes(key, value, siblings) + if !ok { + contentType = firstNonEmptyString(mediaContentTypeFromItem(siblings), defaultContentTypeForRawMediaKey(key)) + if !generatedContentTypeIsMedia(contentType) && !generatedContentTypeIsDocument(contentType) { + return nil, false + } + var err error + payload, _, ok, err = inlineMediaPayload(value, false) + if err != nil || !ok || len(payload) == 0 { + return nil, false + } } + contentType = firstNonEmptyString(contentType, detectGeneratedAssetContentType(payload)) kind := mediaKindForAsset(taskKind, siblings, key, contentType) return &generatedInlineAsset{ Bytes: payload, diff --git a/apps/api/internal/runner/upload_test.go b/apps/api/internal/runner/upload_test.go index 2a92f16..46be563 100644 --- a/apps/api/internal/runner/upload_test.go +++ b/apps/api/internal/runner/upload_test.go @@ -348,6 +348,7 @@ func TestFinalizeGeneratedAssetsMaterializesDetectedMediaAndKeepsOpaqueMetadata( result := map[string]any{ "provider_payload": base64.StdEncoding.EncodeToString(payload), "thought_signature": opaque, + "thinking_bytes": opaque, } finalized, err := service.finalizeGeneratedAssets( @@ -373,6 +374,43 @@ func TestFinalizeGeneratedAssetsMaterializesDetectedMediaAndKeepsOpaqueMetadata( if finalized["thought_signature"] != opaque { t.Fatal("opaque provider metadata must remain unchanged") } + if finalized["thinking_bytes"] != opaque { + t.Fatal("opaque provider bytes metadata must remain unchanged") + } +} + +func TestFinalizeGeneratedAssetsMaterializesGeminiImageData(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, _ *http.Request) { w.WriteHeader(http.StatusOK) })) + defer server.Close() + service := &Service{} + channels := []store.FileStorageChannel{testObjectStorageChannel("s3", server.URL, "s3-result")} + payload := append([]byte{0x89, 'P', 'N', 'G', 0x0d, 0x0a, 0x1a, 0x0a}, bytes.Repeat([]byte{0}, 256)...) + result := map[string]any{ + "id": "gemini-image", + "created": float64(1), + "model": "gemini-3-pro-image", + "data": []any{map[string]any{ + "b64_json": base64.StdEncoding.EncodeToString(payload), + "mime_type": "image/png", + }}, + } + + finalized, err := service.finalizeGeneratedAssets( + t.Context(), "task-gemini-image", "images.generations", result, + defaultGeneratedAssetUploadPolicy(), channels, true, 0, + ) + if err != nil { + t.Fatal(err) + } + if TaskResultHasInlineBinary(finalized) { + t.Fatalf("Gemini image result still contains inline binary: %+v", taskResultInlineBinaryDiagnostics(finalized)) + } + data := finalized["data"].([]any) + item := data[0].(map[string]any) + reference, ok := item["b64_json"].(map[string]any) + if !ok || reference["assetRef"] == nil || reference["upload"] == nil { + t.Fatalf("Gemini b64_json was not replaced by an object reference: %#v", item) + } } func TestResolvedGeneratedAssetContentTypePrefersDetectedMedia(t *testing.T) { diff --git a/scripts/acceptance/run-live-object-storage.mjs b/scripts/acceptance/run-live-object-storage.mjs index 633ed54..90703ff 100755 --- a/scripts/acceptance/run-live-object-storage.mjs +++ b/scripts/acceptance/run-live-object-storage.mjs @@ -149,7 +149,7 @@ function channelInput() { endpoint: endpoint.toString().replace(/\/$/, ''), region, bucket, - objectKeyPrefix: 'easyai-ai-gateway/live-acceptance', + objectPrefix: 'easyai-ai-gateway/live-acceptance', accessScope: 'private', temporaryFileExpirePolicy: expirationPolicy, signedUrlExpiresSeconds: signedURLExpiresSeconds,