diff --git a/apps/api/internal/runner/admission.go b/apps/api/internal/runner/admission.go index 5b63cd0..5299b49 100644 --- a/apps/api/internal/runner/admission.go +++ b/apps/api/internal/runner/admission.go @@ -338,6 +338,16 @@ func (s *Service) activeAsyncTaskAdmission( if len(leases) == 0 { return store.TaskAdmissionResult{}, false, nil } + leaseStore := s.coordinationStore + if leaseStore == nil { + leaseStore = s.store + } + if err := leaseStore.RenewConcurrencyLeases(ctx, leases); err != nil { + if errors.Is(err, store.ErrConcurrencyLeaseLost) { + return store.TaskAdmissionResult{}, false, nil + } + return store.TaskAdmissionResult{}, false, err + } return store.TaskAdmissionResult{ Admission: *admission, Admitted: true, diff --git a/apps/api/internal/runner/binary_results.go b/apps/api/internal/runner/binary_results.go index cd219e1..420abd0 100644 --- a/apps/api/internal/runner/binary_results.go +++ b/apps/api/internal/runner/binary_results.go @@ -517,7 +517,14 @@ func localBinaryStringBytes(key string, value string, siblings map[string]any) ( if err != nil || len(payload) == 0 { return nil, "", "", false } - return payload, firstNonEmptyString(mediaContentTypeFromItem(siblings), defaultContentTypeForRawMediaKey(key)), "raw", true + contentType := firstNonEmptyString(mediaContentTypeFromItem(siblings), defaultContentTypeForRawMediaKey(key)) + if !strict && contentType == "" { + contentType = detectGeneratedAssetContentType(payload) + if !generatedContentTypeIsMedia(contentType) && !generatedContentTypeIsDocument(contentType) { + return nil, "", "", false + } + } + return payload, contentType, "raw", true } func localBufferObjectBytes(value map[string]any) ([]byte, string, bool) { diff --git a/apps/api/internal/runner/binary_results_test.go b/apps/api/internal/runner/binary_results_test.go index 0eea6e7..18d7a0f 100644 --- a/apps/api/internal/runner/binary_results_test.go +++ b/apps/api/internal/runner/binary_results_test.go @@ -1,6 +1,7 @@ package runner import ( + "bytes" "context" "crypto/sha256" "encoding/base64" @@ -311,6 +312,24 @@ 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} + + if TaskResultHasInlineBinary(result) { + t.Fatal("opaque Base64 metadata must not be mistaken for generated media") + } +} + +func TestTaskResultDetectsLongBase64MediaWithoutSemanticKey(t *testing.T) { + payload := append([]byte{0x89, 'P', 'N', 'G', 0x0d, 0x0a, 0x1a, 0x0a}, bytes.Repeat([]byte{0}, 4096)...) + result := map[string]any{"provider_payload": base64.StdEncoding.EncodeToString(payload)} + + if !TaskResultHasInlineBinary(result) { + t.Fatal("signature-detected generated media must still be materialized") + } +} + func bytesToAny(payload []byte) []any { result := make([]any, len(payload)) for index, value := range payload { diff --git a/apps/api/internal/runner/retry_decision.go b/apps/api/internal/runner/retry_decision.go index 9968d06..8b55cb2 100644 --- a/apps/api/internal/runner/retry_decision.go +++ b/apps/api/internal/runner/retry_decision.go @@ -303,16 +303,16 @@ func priorityDemoteDecisionForCandidate(runnerPolicy store.RunnerPolicy, err err } func isResultPersistenceFailure(err error) bool { - switch strings.ToLower(strings.TrimSpace(clients.ErrorCode(err))) { - case "local_result_storage_unavailable", - "binary_result_too_large", - "binary_result_corrupted", - "binary_result_expired", - "result_binary_not_materialized": - return true - default: - return false - } + return isGatewayResultPersistenceCode(clients.ErrorCode(err)) +} + +func isGatewayResultPersistenceCode(code string) bool { + code = strings.ToLower(strings.TrimSpace(code)) + return strings.HasPrefix(code, "storage_") || + strings.HasPrefix(code, "upload_") || + strings.HasPrefix(code, "binary_result_") || + strings.HasPrefix(code, "local_result_") || + strings.HasPrefix(code, "result_binary_") } func effectiveFailoverPolicy(base map[string]any, override map[string]any) map[string]any { @@ -353,6 +353,8 @@ func failureCategory(code string, status int, message string) string { switch { case code == "insufficient_balance": return "insufficient_balance" + case isGatewayResultPersistenceCode(code): + return "gateway_storage" case code == "rate_limit" || status == 429: return "rate_limit" case code == "network": diff --git a/apps/api/internal/runner/retry_decision_test.go b/apps/api/internal/runner/retry_decision_test.go index ac59568..25113cb 100644 --- a/apps/api/internal/runner/retry_decision_test.go +++ b/apps/api/internal/runner/retry_decision_test.go @@ -495,3 +495,37 @@ func TestResolveCandidateFailureUnmatchedRetryableUsesSameThenNext(t *testing.T) t.Fatalf("unmatched exhausted error should rotate without side effects, got %+v", decision) } } + +func TestGatewayStorageFailureNeverRetriesProvider(t *testing.T) { + err := &clients.ClientError{ + Code: "storage_write_failed", + Message: "generated binary result could not be written to object storage", + StatusCode: 503, + Retryable: true, + } + candidate := store.RuntimeModelCandidate{ModelRetryPolicy: map[string]any{ + "enabled": true, + "allowKeywords": []any{"5xx"}, + }} + + retry := retryDecisionForCandidate(candidate, err) + if retry.Retry || retry.Reason != "result_persistence_failed" { + t.Fatalf("Gateway storage failure must not repeat an upstream request: %+v", retry) + } + failover := failoverDecisionForCandidate(store.RunnerPolicy{ + Status: "active", + FailoverPolicy: map[string]any{"enabled": true, "allowCategories": []any{"provider_5xx"}}, + }, candidate, err) + if failover.Retry || failover.Reason != "result_persistence_failed" { + t.Fatalf("Gateway storage failure must not fail over to another provider: %+v", failover) + } + decision := resolveCandidateFailure(resolveCandidateFailureInput{ + RunnerPolicy: store.RunnerPolicy{Status: "active"}, + Err: err, + HasNextCandidate: true, + Async: true, + }) + if decision.Route != "stop" || decision.Reason != "result_persistence_failed" || decision.Info.Category != "gateway_storage" { + t.Fatalf("Gateway storage failure classification is incorrect: %+v", decision) + } +} diff --git a/apps/api/internal/runner/upload.go b/apps/api/internal/runner/upload.go index fdfed64..c20cb9e 100644 --- a/apps/api/internal/runner/upload.go +++ b/apps/api/internal/runner/upload.go @@ -423,18 +423,16 @@ func generatedRawInlineMediaAsset(key string, value string, siblings map[string] } keyLooksLikeMediaPayload := generatedRawDataMediaPayloadKey(key) contentType := firstNonEmptyString(mediaContentTypeFromItem(siblings), defaultContentTypeForRawMediaKey(key)) - if !keyLooksLikeMediaPayload && !generatedContentTypeIsMedia(contentType) { - return nil, false - } - if !keyLooksLikeMediaPayload && !strings.HasPrefix(strings.ToLower(raw), "data:") && len(raw) < 128 { + 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) - if !generatedContentTypeIsMedia(contentType) { + contentType = firstNonEmptyString(payloadContentType, contentType, detectGeneratedAssetContentType(payload)) + if !generatedContentTypeIsMedia(contentType) && !generatedContentTypeIsDocument(contentType) { return nil, false } kind := mediaKindForAsset(taskKind, siblings, key, contentType) diff --git a/apps/api/internal/runner/upload_test.go b/apps/api/internal/runner/upload_test.go index 342b65e..2a92f16 100644 --- a/apps/api/internal/runner/upload_test.go +++ b/apps/api/internal/runner/upload_test.go @@ -338,6 +338,43 @@ func TestFinalizeGeneratedAssetsUploadsNestedInlineBinaryUnderDefaultPolicy(t *t } } +func TestFinalizeGeneratedAssetsMaterializesDetectedMediaAndKeepsOpaqueMetadata(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}, 4096)...) + opaque := base64.StdEncoding.EncodeToString([]byte(strings.Repeat("opaque metadata ", 400))) + result := map[string]any{ + "provider_payload": base64.StdEncoding.EncodeToString(payload), + "thought_signature": opaque, + } + + finalized, err := service.finalizeGeneratedAssets( + t.Context(), + "task-provider-payload", + "images.generations", + result, + defaultGeneratedAssetUploadPolicy(), + channels, + true, + 0, + ) + if err != nil { + t.Fatal(err) + } + if TaskResultHasInlineBinary(finalized) { + t.Fatalf("finalized result still contains generated media: %#v", finalized) + } + reference, ok := finalized["provider_payload"].(map[string]any) + if !ok || reference["assetRef"] == nil || reference["upload"] == nil { + t.Fatalf("detected media was not objectified: %#v", finalized["provider_payload"]) + } + if finalized["thought_signature"] != opaque { + t.Fatal("opaque provider metadata must remain unchanged") + } +} + func TestResolvedGeneratedAssetContentTypePrefersDetectedMedia(t *testing.T) { pngPayload := []byte{0x89, 'P', 'N', 'G', 0x0d, 0x0a, 0x1a, 0x0a, 0, 0, 0, 0} diff --git a/deploy/kubernetes/production/application-config.yaml b/deploy/kubernetes/production/application-config.yaml index a52bdf6..7134bea 100644 --- a/deploy/kubernetes/production/application-config.yaml +++ b/deploy/kubernetes/production/application-config.yaml @@ -15,7 +15,9 @@ data: IDENTITY_SECURITY_EVENTS_STALE_AFTER_SECONDS: "180" IDENTITY_SECURITY_EVENTS_CLOCK_SKEW_SECONDS: "60" SERVER_MAIN_BASE_URL: http://10.77.0.1:3001 - TASK_PROGRESS_CALLBACK_ENABLED: "true" + # server-main 当前没有实现该回调路由。关闭无效投递,任务结果继续通过 + # Gateway task detail/events 接口读取;实现并验收接收端后再显式开启。 + TASK_PROGRESS_CALLBACK_ENABLED: "false" TASK_PROGRESS_CALLBACK_URL: http://10.77.0.1:3001/internal/platform/task-progress-callbacks TASK_PROGRESS_CALLBACK_TIMEOUT_MS: "5000" TASK_PROGRESS_CALLBACK_MAX_ATTEMPTS: "10"