From 76a7702925c0bcd92d62db345f35e38151b9b875 Mon Sep 17 00:00:00 2001 From: wangbo Date: Thu, 23 Jul 2026 22:35:41 +0800 Subject: [PATCH] =?UTF-8?q?fix(media):=20=E5=AF=B9=E9=BD=90=20EasyAI=20?= =?UTF-8?q?=E5=AA=92=E4=BD=93=E5=93=8D=E5=BA=94=E4=B8=8E=E8=BD=AC=E5=AD=98?= =?UTF-8?q?=E7=AD=96=E7=95=A5?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 统一媒体异步提交和轮询响应,扩展 /ai/result 到当前用户的 Gateway 媒体任务,并保持既有 Gateway 字段兼容。\n\n启用转存但无可用渠道时回退到 24 小时本地静态资源;关闭转存时保留图片、音频等上游 Base64 字段。同步更新文件上传兼容结构、OpenAPI、管理端说明和回归测试。\n\n验证:cd apps/api && env -u AI_GATEWAY_TEST_DATABASE_URL go test ./... -count=1;gofmt -l 无输出;pnpm nx run web:typecheck。 --- apps/api/docs/swagger.json | 298 +++++++++-- apps/api/docs/swagger.yaml | 237 +++++++-- .../httpapi/advanced_task_handlers.go | 4 +- .../httpapi/core_flow_integration_test.go | 97 +++- apps/api/internal/httpapi/easyai_compat.go | 489 ++++++++++++++++++ .../internal/httpapi/easyai_compat_test.go | 311 +++++++++++ .../internal/httpapi/file_upload_handlers.go | 2 +- apps/api/internal/httpapi/handlers.go | 55 +- apps/api/internal/httpapi/openapi_models.go | 70 ++- apps/api/internal/httpapi/server.go | 2 +- .../internal/httpapi/task_idempotency_test.go | 30 ++ .../httpapi/volces_compat_handlers.go | 49 +- apps/api/internal/runner/upload.go | 34 +- apps/api/internal/runner/upload_test.go | 43 +- .../src/pages/admin/SystemSettingsPanel.tsx | 4 +- 15 files changed, 1562 insertions(+), 163 deletions(-) create mode 100644 apps/api/internal/httpapi/easyai_compat.go create mode 100644 apps/api/internal/httpapi/easyai_compat_test.go diff --git a/apps/api/docs/swagger.json b/apps/api/docs/swagger.json index a337bba..d7c7808 100644 --- a/apps/api/docs/swagger.json +++ b/apps/api/docs/swagger.json @@ -4709,13 +4709,14 @@ "BearerAuth": [] } ], + "description": "查询当前用户的图片、视频、音频、视频高清和矢量化任务,并返回 EasyAI GeneratedResponse 兼容结构。", "produces": [ "application/json" ], "tags": [ - "volces-compatible" + "tasks" ], - "summary": "查询 server-main 兼容视频结果", + "summary": "查询 server-main EasyAIClient 兼容任务结果", "parameters": [ { "type": "string", @@ -4729,8 +4730,13 @@ "200": { "description": "OK", "schema": { - "type": "object", - "additionalProperties": true + "$ref": "#/definitions/httpapi.EasyAIGeneratedResponse" + } + }, + "404": { + "description": "Not Found", + "schema": { + "$ref": "#/definitions/httpapi.EasyAIGeneratedResponse" } } } @@ -5976,7 +5982,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/httpapi.CompatibleResponse" + "$ref": "#/definitions/httpapi.OpenAIImagesResponse" } }, "202": { @@ -6051,7 +6057,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/httpapi.CompatibleResponse" + "$ref": "#/definitions/httpapi.OpenAIImagesResponse" } }, "202": { @@ -6126,7 +6132,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/httpapi.CompatibleResponse" + "$ref": "#/definitions/httpapi.EasyAIGeneratedResponse" } }, "202": { @@ -6710,7 +6716,7 @@ "BearerAuth": [] } ], - "description": "统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。", + "description": "默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway 与 server-main EasyAIClient 的异步提交结构。", "consumes": [ "application/json" ], @@ -6718,9 +6724,9 @@ "application/json" ], "tags": [ - "tasks" + "media" ], - "summary": "创建或执行 AI 任务", + "summary": "创建 EasyAI 媒体任务", "parameters": [ { "type": "boolean", @@ -6735,7 +6741,7 @@ "in": "header" }, { - "description": "AI 任务请求,字段随任务类型变化", + "description": "媒体任务请求,字段随任务类型变化", "name": "input", "in": "body", "required": true, @@ -6748,7 +6754,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/httpapi.CompatibleResponse" + "$ref": "#/definitions/httpapi.EasyAIGeneratedResponse" } }, "202": { @@ -7548,7 +7554,7 @@ "BearerAuth": [] } ], - "description": "统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。", + "description": "默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway 与 server-main EasyAIClient 的异步提交结构。", "consumes": [ "application/json" ], @@ -7556,9 +7562,9 @@ "application/json" ], "tags": [ - "tasks" + "media" ], - "summary": "创建或执行 AI 任务", + "summary": "创建 EasyAI 媒体任务", "parameters": [ { "type": "boolean", @@ -7573,7 +7579,7 @@ "in": "header" }, { - "description": "AI 任务请求,字段随任务类型变化", + "description": "媒体任务请求,字段随任务类型变化", "name": "input", "in": "body", "required": true, @@ -7586,7 +7592,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/httpapi.CompatibleResponse" + "$ref": "#/definitions/httpapi.EasyAIGeneratedResponse" } }, "202": { @@ -7647,7 +7653,7 @@ "BearerAuth": [] } ], - "description": "统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。", + "description": "默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway 与 server-main EasyAIClient 的异步提交结构。", "consumes": [ "application/json" ], @@ -7655,9 +7661,9 @@ "application/json" ], "tags": [ - "tasks" + "media" ], - "summary": "创建或执行 AI 任务", + "summary": "创建 EasyAI 媒体任务", "parameters": [ { "type": "boolean", @@ -7672,7 +7678,7 @@ "in": "header" }, { - "description": "AI 任务请求,字段随任务类型变化", + "description": "媒体任务请求,字段随任务类型变化", "name": "input", "in": "body", "required": true, @@ -7685,7 +7691,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/httpapi.CompatibleResponse" + "$ref": "#/definitions/httpapi.EasyAIGeneratedResponse" } }, "202": { @@ -8063,15 +8069,14 @@ "application/json" ], "tags": [ - "volces-compatible" + "media" ], "summary": "创建 server-main 兼容视频任务", "responses": { "200": { "description": "OK", "schema": { - "type": "object", - "additionalProperties": true + "$ref": "#/definitions/httpapi.TaskAcceptedResponse" } } } @@ -8084,7 +8089,7 @@ "BearerAuth": [] } ], - "description": "统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。", + "description": "默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway 与 server-main EasyAIClient 的异步提交结构。", "consumes": [ "application/json" ], @@ -8092,9 +8097,9 @@ "application/json" ], "tags": [ - "tasks" + "media" ], - "summary": "创建或执行 AI 任务", + "summary": "创建 EasyAI 媒体任务", "parameters": [ { "type": "boolean", @@ -8109,7 +8114,7 @@ "in": "header" }, { - "description": "AI 任务请求,字段随任务类型变化", + "description": "媒体任务请求,字段随任务类型变化", "name": "input", "in": "body", "required": true, @@ -8122,7 +8127,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/httpapi.CompatibleResponse" + "$ref": "#/definitions/httpapi.EasyAIGeneratedResponse" } }, "202": { @@ -8348,7 +8353,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/httpapi.CompatibleResponse" + "$ref": "#/definitions/httpapi.EasyAIGeneratedResponse" } }, "202": { @@ -8397,7 +8402,7 @@ "BearerAuth": [] } ], - "description": "统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。", + "description": "默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway 与 server-main EasyAIClient 的异步提交结构。", "consumes": [ "application/json" ], @@ -8405,9 +8410,9 @@ "application/json" ], "tags": [ - "tasks" + "media" ], - "summary": "创建或执行 AI 任务", + "summary": "创建 EasyAI 媒体任务", "parameters": [ { "type": "boolean", @@ -8422,7 +8427,7 @@ "in": "header" }, { - "description": "AI 任务请求,字段随任务类型变化", + "description": "媒体任务请求,字段随任务类型变化", "name": "input", "in": "body", "required": true, @@ -8435,7 +8440,7 @@ "200": { "description": "OK", "schema": { - "$ref": "#/definitions/httpapi.CompatibleResponse" + "$ref": "#/definitions/httpapi.EasyAIGeneratedResponse" } }, "202": { @@ -9813,6 +9818,143 @@ } } }, + "httpapi.EasyAIGeneratedResponse": { + "type": "object", + "properties": { + "cancellable": { + "type": "boolean" + }, + "code": { + "type": "string" + }, + "created": { + "type": "integer", + "example": 1784772000000 + }, + "data": { + "type": "array", + "items": { + "$ref": "#/definitions/httpapi.EasyAIMediaOutput" + } + }, + "message": { + "type": "string" + }, + "output": { + "type": "array", + "items": { + "type": "string" + } + }, + "output_content": { + "type": "array", + "items": { + "$ref": "#/definitions/httpapi.EasyAIMediaOutput" + } + }, + "query_url": { + "type": "string", + "example": "/api/v1/ai/result/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25" + }, + "result": { + "type": "object", + "additionalProperties": true + }, + "status": { + "type": "string", + "enum": [ + "submitted", + "process", + "success", + "failed" + ], + "example": "success" + }, + "submitted": { + "type": "boolean" + }, + "taskId": { + "type": "string", + "example": "9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25" + }, + "task_id": { + "type": "string", + "example": "9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25" + }, + "upstream_task_id": { + "type": "string", + "example": "provider-task-123" + }, + "usage": { + "type": "object", + "additionalProperties": true + } + } + }, + "httpapi.EasyAIMediaOutput": { + "type": "object", + "properties": { + "audio_url": { + "type": "string" + }, + "b64_json": { + "description": "B64JSON is retained when result media transfer is disabled.", + "type": "string" + }, + "content": { + "type": "string" + }, + "duration": { + "type": "number" + }, + "format": { + "type": "string" + }, + "height": { + "type": "integer" + }, + "image_url": { + "type": "string" + }, + "metadata": { + "type": "object", + "additionalProperties": true + }, + "mime_type": { + "type": "string" + }, + "source_resolution": { + "type": "string" + }, + "target_resolution": { + "type": "string" + }, + "text": { + "type": "string" + }, + "type": { + "type": "string", + "enum": [ + "image", + "video", + "audio", + "file", + "text" + ], + "example": "video" + }, + "url": { + "type": "string", + "example": "https://cdn.example.com/output.mp4" + }, + "video_url": { + "type": "string" + }, + "width": { + "type": "integer" + } + } + }, "httpapi.ErrorEnvelope": { "type": "object", "properties": { @@ -9854,7 +9996,7 @@ } } }, - "httpapi.FileUploadResponse": { + "httpapi.FileUploadData": { "type": "object", "properties": { "assetStorage": { @@ -9883,6 +10025,42 @@ } } }, + "httpapi.FileUploadResponse": { + "type": "object", + "properties": { + "assetStorage": { + "type": "object", + "additionalProperties": true + }, + "contentType": { + "type": "string", + "example": "image/png" + }, + "data": { + "$ref": "#/definitions/httpapi.FileUploadData" + }, + "filename": { + "type": "string", + "example": "image.png" + }, + "id": { + "type": "string", + "example": "file_abc123" + }, + "size": { + "type": "integer", + "example": 1024 + }, + "status": { + "type": "string", + "example": "success" + }, + "url": { + "type": "string", + "example": "/static/uploaded/upload-abc123.png" + } + } + }, "httpapi.GeminiErrorEnvelope": { "type": "object", "properties": { @@ -10514,6 +10692,36 @@ } } }, + "httpapi.OpenAIImageData": { + "type": "object", + "properties": { + "b64_json": { + "type": "string" + }, + "revised_prompt": { + "type": "string" + }, + "url": { + "type": "string", + "example": "https://cdn.example.com/output.png" + } + } + }, + "httpapi.OpenAIImagesResponse": { + "type": "object", + "properties": { + "created": { + "type": "integer", + "example": 1710000000 + }, + "data": { + "type": "array", + "items": { + "$ref": "#/definitions/httpapi.OpenAIImageData" + } + } + } + }, "httpapi.PlatformListResponse": { "type": "object", "properties": { @@ -10894,12 +11102,24 @@ "next": { "$ref": "#/definitions/httpapi.TaskNextLinks" }, + "query_url": { + "type": "string", + "example": "/api/v1/ai/result/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25" + }, + "status": { + "type": "string", + "example": "submitted" + }, "task": { "$ref": "#/definitions/store.GatewayTask" }, "taskId": { "type": "string", "example": "9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25" + }, + "task_id": { + "type": "string", + "example": "9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25" } } }, @@ -10961,6 +11181,10 @@ "events": { "type": "string", "example": "/api/v1/tasks/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25/events" + }, + "result": { + "type": "string", + "example": "/api/v1/ai/result/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25" } } }, diff --git a/apps/api/docs/swagger.yaml b/apps/api/docs/swagger.yaml index 3ff0350..37a9eb0 100644 --- a/apps/api/docs/swagger.yaml +++ b/apps/api/docs/swagger.yaml @@ -346,6 +346,103 @@ definitions: $ref: '#/definitions/store.DailyTokenUsage' type: array type: object + httpapi.EasyAIGeneratedResponse: + properties: + cancellable: + type: boolean + code: + type: string + created: + example: 1784772000000 + type: integer + data: + items: + $ref: '#/definitions/httpapi.EasyAIMediaOutput' + type: array + message: + type: string + output: + items: + type: string + type: array + output_content: + items: + $ref: '#/definitions/httpapi.EasyAIMediaOutput' + type: array + query_url: + example: /api/v1/ai/result/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25 + type: string + result: + additionalProperties: true + type: object + status: + enum: + - submitted + - process + - success + - failed + example: success + type: string + submitted: + type: boolean + task_id: + example: 9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25 + type: string + taskId: + example: 9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25 + type: string + upstream_task_id: + example: provider-task-123 + type: string + usage: + additionalProperties: true + type: object + type: object + httpapi.EasyAIMediaOutput: + properties: + audio_url: + type: string + b64_json: + description: B64JSON is retained when result media transfer is disabled. + type: string + content: + type: string + duration: + type: number + format: + type: string + height: + type: integer + image_url: + type: string + metadata: + additionalProperties: true + type: object + mime_type: + type: string + source_resolution: + type: string + target_resolution: + type: string + text: + type: string + type: + enum: + - image + - video + - audio + - file + - text + example: video + type: string + url: + example: https://cdn.example.com/output.mp4 + type: string + video_url: + type: string + width: + type: integer + type: object httpapi.ErrorEnvelope: properties: error: @@ -374,7 +471,7 @@ definitions: $ref: '#/definitions/store.FileStorageChannel' type: array type: object - httpapi.FileUploadResponse: + httpapi.FileUploadData: properties: assetStorage: additionalProperties: true @@ -395,6 +492,32 @@ definitions: example: /static/uploaded/upload-abc123.png type: string type: object + httpapi.FileUploadResponse: + properties: + assetStorage: + additionalProperties: true + type: object + contentType: + example: image/png + type: string + data: + $ref: '#/definitions/httpapi.FileUploadData' + filename: + example: image.png + type: string + id: + example: file_abc123 + type: string + size: + example: 1024 + type: integer + status: + example: success + type: string + url: + example: /static/uploaded/upload-abc123.png + type: string + type: object httpapi.GeminiErrorEnvelope: properties: error: @@ -826,6 +949,26 @@ definitions: example: invalid_request_error type: string type: object + httpapi.OpenAIImageData: + properties: + b64_json: + type: string + revised_prompt: + type: string + url: + example: https://cdn.example.com/output.png + type: string + type: object + httpapi.OpenAIImagesResponse: + properties: + created: + example: 1710000000 + type: integer + data: + items: + $ref: '#/definitions/httpapi.OpenAIImageData' + type: array + type: object httpapi.PlatformListResponse: properties: items: @@ -1092,8 +1235,17 @@ definitions: properties: next: $ref: '#/definitions/httpapi.TaskNextLinks' + query_url: + example: /api/v1/ai/result/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25 + type: string + status: + example: submitted + type: string task: $ref: '#/definitions/store.GatewayTask' + task_id: + example: 9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25 + type: string taskId: example: 9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25 type: string @@ -1140,6 +1292,9 @@ definitions: events: example: /api/v1/tasks/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25/events type: string + result: + example: /api/v1/ai/result/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25 + type: string type: object httpapi.TaskParamPreprocessingLogListResponse: properties: @@ -6638,6 +6793,7 @@ paths: - playground /api/v1/ai/result/{taskID}: get: + description: 查询当前用户的图片、视频、音频、视频高清和矢量化任务,并返回 EasyAI GeneratedResponse 兼容结构。 parameters: - description: 任务 ID in: path @@ -6650,13 +6806,16 @@ paths: "200": description: OK schema: - additionalProperties: true - type: object + $ref: '#/definitions/httpapi.EasyAIGeneratedResponse' + "404": + description: Not Found + schema: + $ref: '#/definitions/httpapi.EasyAIGeneratedResponse' security: - BearerAuth: [] - summary: 查询 server-main 兼容视频结果 + summary: 查询 server-main EasyAIClient 兼容任务结果 tags: - - volces-compatible + - tasks /api/v1/api-keys: get: description: 返回当前用户创建的 API Key 元数据,secret 只在创建时返回。 @@ -7451,7 +7610,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/httpapi.CompatibleResponse' + $ref: '#/definitions/httpapi.OpenAIImagesResponse' "202": description: Accepted schema: @@ -7499,7 +7658,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/httpapi.CompatibleResponse' + $ref: '#/definitions/httpapi.OpenAIImagesResponse' "202": description: Accepted schema: @@ -7548,7 +7707,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/httpapi.CompatibleResponse' + $ref: '#/definitions/httpapi.EasyAIGeneratedResponse' "202": description: Accepted schema: @@ -7926,7 +8085,8 @@ paths: post: consumes: - application/json - description: 统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。 + description: 默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway + 与 server-main EasyAIClient 的异步提交结构。 parameters: - description: true 时异步创建任务并返回 202 in: header @@ -7936,7 +8096,7 @@ paths: in: header name: Idempotency-Key type: string - - description: AI 任务请求,字段随任务类型变化 + - description: 媒体任务请求,字段随任务类型变化 in: body name: input required: true @@ -7948,7 +8108,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/httpapi.CompatibleResponse' + $ref: '#/definitions/httpapi.EasyAIGeneratedResponse' "202": description: Accepted schema: @@ -7983,9 +8143,9 @@ paths: $ref: '#/definitions/httpapi.ErrorEnvelope' security: - BearerAuth: [] - summary: 创建或执行 AI 任务 + summary: 创建 EasyAI 媒体任务 tags: - - tasks + - media /api/v1/openapi.json: get: description: 返回当前构建内嵌的完整机器可读 Swagger JSON,供 Agent 在 SKILL references 未覆盖接口时查询。 @@ -8470,7 +8630,8 @@ paths: post: consumes: - application/json - description: 统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。 + description: 默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway + 与 server-main EasyAIClient 的异步提交结构。 parameters: - description: true 时异步创建任务并返回 202 in: header @@ -8480,7 +8641,7 @@ paths: in: header name: Idempotency-Key type: string - - description: AI 任务请求,字段随任务类型变化 + - description: 媒体任务请求,字段随任务类型变化 in: body name: input required: true @@ -8492,7 +8653,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/httpapi.CompatibleResponse' + $ref: '#/definitions/httpapi.EasyAIGeneratedResponse' "202": description: Accepted schema: @@ -8527,14 +8688,15 @@ paths: $ref: '#/definitions/httpapi.ErrorEnvelope' security: - BearerAuth: [] - summary: 创建或执行 AI 任务 + summary: 创建 EasyAI 媒体任务 tags: - - tasks + - media /api/v1/speech/generations: post: consumes: - application/json - description: 统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。 + description: 默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway + 与 server-main EasyAIClient 的异步提交结构。 parameters: - description: true 时异步创建任务并返回 202 in: header @@ -8544,7 +8706,7 @@ paths: in: header name: Idempotency-Key type: string - - description: AI 任务请求,字段随任务类型变化 + - description: 媒体任务请求,字段随任务类型变化 in: body name: input required: true @@ -8556,7 +8718,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/httpapi.CompatibleResponse' + $ref: '#/definitions/httpapi.EasyAIGeneratedResponse' "202": description: Accepted schema: @@ -8591,9 +8753,9 @@ paths: $ref: '#/definitions/httpapi.ErrorEnvelope' security: - BearerAuth: [] - summary: 创建或执行 AI 任务 + summary: 创建 EasyAI 媒体任务 tags: - - tasks + - media /api/v1/tasks: get: description: 按当前用户列出任务,支持关键字、模型类型、时间范围和分页过滤。 @@ -8803,18 +8965,18 @@ paths: "200": description: OK schema: - additionalProperties: true - type: object + $ref: '#/definitions/httpapi.TaskAcceptedResponse' security: - BearerAuth: [] summary: 创建 server-main 兼容视频任务 tags: - - volces-compatible + - media /api/v1/videos/generations: post: consumes: - application/json - description: 统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。 + description: 默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway + 与 server-main EasyAIClient 的异步提交结构。 parameters: - description: true 时异步创建任务并返回 202 in: header @@ -8824,7 +8986,7 @@ paths: in: header name: Idempotency-Key type: string - - description: AI 任务请求,字段随任务类型变化 + - description: 媒体任务请求,字段随任务类型变化 in: body name: input required: true @@ -8836,7 +8998,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/httpapi.CompatibleResponse' + $ref: '#/definitions/httpapi.EasyAIGeneratedResponse' "202": description: Accepted schema: @@ -8871,9 +9033,9 @@ paths: $ref: '#/definitions/httpapi.ErrorEnvelope' security: - BearerAuth: [] - summary: 创建或执行 AI 任务 + summary: 创建 EasyAI 媒体任务 tags: - - tasks + - media /api/v1/videos/omni-video: post: consumes: @@ -8983,7 +9145,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/httpapi.CompatibleResponse' + $ref: '#/definitions/httpapi.EasyAIGeneratedResponse' "202": description: Accepted schema: @@ -9017,7 +9179,8 @@ paths: post: consumes: - application/json - description: 统一公开入口按 model 选择平台模型并默认同步返回兼容响应;设置 X-Async=true 时异步创建任务并返回 202。 + description: 默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway + 与 server-main EasyAIClient 的异步提交结构。 parameters: - description: true 时异步创建任务并返回 202 in: header @@ -9027,7 +9190,7 @@ paths: in: header name: Idempotency-Key type: string - - description: AI 任务请求,字段随任务类型变化 + - description: 媒体任务请求,字段随任务类型变化 in: body name: input required: true @@ -9039,7 +9202,7 @@ paths: "200": description: OK schema: - $ref: '#/definitions/httpapi.CompatibleResponse' + $ref: '#/definitions/httpapi.EasyAIGeneratedResponse' "202": description: Accepted schema: @@ -9074,9 +9237,9 @@ paths: $ref: '#/definitions/httpapi.ErrorEnvelope' security: - BearerAuth: [] - summary: 创建或执行 AI 任务 + summary: 创建 EasyAI 媒体任务 tags: - - tasks + - media /api/v1/voice_clone/voices: get: description: 返回当前用户在网关中维护的克隆音色,以及克隆时绑定的平台与平台模型。 diff --git a/apps/api/internal/httpapi/advanced_task_handlers.go b/apps/api/internal/httpapi/advanced_task_handlers.go index b4dc5db..d3e7c92 100644 --- a/apps/api/internal/httpapi/advanced_task_handlers.go +++ b/apps/api/internal/httpapi/advanced_task_handlers.go @@ -11,7 +11,7 @@ import "net/http" // @Security BearerAuth // @Param X-Async header bool false "true 时异步创建任务并返回 202" // @Param input body ImageVectorizeRequest true "图片矢量化请求" -// @Success 200 {object} CompatibleResponse +// @Success 200 {object} EasyAIGeneratedResponse // @Success 202 {object} TaskAcceptedResponse // @Failure 400 {object} ErrorEnvelope // @Failure 401 {object} ErrorEnvelope @@ -32,7 +32,7 @@ func (s *Server) createImageVectorizeTask() http.Handler { // @Security BearerAuth // @Param X-Async header bool false "true 时异步创建任务并返回 202" // @Param input body VideoUpscaleRequest true "视频超分请求" -// @Success 200 {object} CompatibleResponse +// @Success 200 {object} EasyAIGeneratedResponse // @Success 202 {object} TaskAcceptedResponse // @Failure 400 {object} ErrorEnvelope // @Failure 401 {object} ErrorEnvelope diff --git a/apps/api/internal/httpapi/core_flow_integration_test.go b/apps/api/internal/httpapi/core_flow_integration_test.go index 6ccc52f..707aab5 100644 --- a/apps/api/internal/httpapi/core_flow_integration_test.go +++ b/apps/api/internal/httpapi/core_flow_integration_test.go @@ -487,6 +487,7 @@ VALUES ($1, 5, '{"purpose":"core-flow"}'::jsonb)`, inviteCode); err != nil { cancelReq.Header.Set("Content-Type", "application/json") cancelErrCh := make(chan error, 1) cancelTaskIDCh := make(chan string, 1) + cancelRequestStartedAt := time.Now().UTC().Add(-time.Second) go func() { resp, err := http.DefaultClient.Do(cancelReq) if resp != nil { @@ -505,7 +506,7 @@ VALUES ($1, 5, '{"purpose":"core-flow"}'::jsonb)`, inviteCode); err != nil { t.Fatal("cancelled stream did not return response headers") } if cancelTaskID == "" { - t.Fatal("cancelled stream response did not expose X-Gateway-Task-Id") + cancelTaskID = waitForLatestTaskIDSince(t, ctx, testPool, "chat.completions", cancelRequestStartedAt, 2*time.Second) } cancelRequest() select { @@ -519,11 +520,20 @@ VALUES ($1, 5, '{"purpose":"core-flow"}'::jsonb)`, inviteCode); err != nil { } var imageResponse struct { - Task struct { + Status string `json:"status"` + LegacyTaskID string `json:"task_id"` + TaskID string `json:"taskId"` + QueryURL string `json:"query_url"` + Task struct { ID string `json:"id"` Status string `json:"status"` Result map[string]any `json:"result"` } `json:"task"` + Next struct { + Detail string `json:"detail"` + Events string `json:"events"` + Result string `json:"result"` + } `json:"next"` } doJSONWithHeaders(t, server.URL, http.MethodPost, "/api/v1/images/generations", apiKeyResponse.Secret, map[string]any{ "model": defaultImageModel, @@ -534,14 +544,59 @@ VALUES ($1, 5, '{"purpose":"core-flow"}'::jsonb)`, inviteCode); err != nil { "simulation": true, "simulationDurationMs": 5, }, map[string]string{"X-Async": "true"}, http.StatusAccepted, &imageResponse) + if imageResponse.Status != "submitted" || + imageResponse.TaskID == "" || + imageResponse.LegacyTaskID != imageResponse.TaskID || + imageResponse.Task.ID != imageResponse.TaskID || + imageResponse.QueryURL != "/api/v1/ai/result/"+imageResponse.TaskID || + imageResponse.Next.Result != imageResponse.QueryURL { + t.Fatalf("unexpected EasyAI-compatible image submission: %+v", imageResponse) + } waitForTaskStatus(t, server.URL, apiKeyResponse.Secret, imageResponse.Task.ID, []string{"succeeded"}, 10*time.Second) doJSON(t, server.URL, http.MethodGet, "/api/v1/tasks/"+imageResponse.Task.ID, apiKeyResponse.Secret, nil, http.StatusOK, &imageResponse.Task) if imageResponse.Task.Status != "succeeded" || imageResponse.Task.Result["id"] == "" { t.Fatalf("unexpected image generation task: %+v", imageResponse.Task) } + var imageEasyAIResult struct { + Status string `json:"status"` + TaskID string `json:"task_id"` + Data []map[string]any `json:"data"` + Output []string `json:"output"` + OutputContent []map[string]any `json:"output_content"` + } + doJSON(t, server.URL, http.MethodGet, imageResponse.QueryURL, apiKeyResponse.Secret, nil, http.StatusOK, &imageEasyAIResult) + if imageEasyAIResult.Status != "success" || + imageEasyAIResult.TaskID != imageResponse.TaskID || + len(imageEasyAIResult.Data) == 0 || + len(imageEasyAIResult.Output) == 0 || + len(imageEasyAIResult.OutputContent) != len(imageEasyAIResult.Data) || + imageEasyAIResult.Data[0]["type"] != "image" { + t.Fatalf("unexpected EasyAI-compatible image result: %+v", imageEasyAIResult) + } + var ordinaryLoginResponse struct { + AccessToken string `json:"accessToken"` + } + doJSON(t, server.URL, http.MethodPost, "/api/v1/auth/login", "", map[string]any{ + "account": ordinaryUsername, + "password": password, + }, http.StatusOK, &ordinaryLoginResponse) + var inaccessibleImageResult struct { + Status string `json:"status"` + Code string `json:"code"` + Data []any `json:"data"` + } + doJSON(t, server.URL, http.MethodGet, imageResponse.QueryURL, ordinaryLoginResponse.AccessToken, nil, http.StatusNotFound, &inaccessibleImageResult) + if inaccessibleImageResult.Status != "failed" || + inaccessibleImageResult.Code != "not_found" || + len(inaccessibleImageResult.Data) != 0 { + t.Fatalf("EasyAI task result must stay owner-isolated: %+v", inaccessibleImageResult) + } var imageEditResponse struct { - Task struct { + Status string `json:"status"` + LegacyTaskID string `json:"task_id"` + TaskID string `json:"taskId"` + Task struct { ID string `json:"id"` Status string `json:"status"` Result map[string]any `json:"result"` @@ -556,6 +611,12 @@ VALUES ($1, 5, '{"purpose":"core-flow"}'::jsonb)`, inviteCode); err != nil { "simulation": true, "simulationDurationMs": 5, }, map[string]string{"X-Async": "true"}, http.StatusAccepted, &imageEditResponse) + if imageEditResponse.Status != "submitted" || + imageEditResponse.TaskID == "" || + imageEditResponse.LegacyTaskID != imageEditResponse.TaskID || + imageEditResponse.Task.ID != imageEditResponse.TaskID { + t.Fatalf("unexpected EasyAI-compatible image edit submission: %+v", imageEditResponse) + } waitForTaskStatus(t, server.URL, apiKeyResponse.Secret, imageEditResponse.Task.ID, []string{"succeeded"}, 10*time.Second) doJSON(t, server.URL, http.MethodGet, "/api/v1/tasks/"+imageEditResponse.Task.ID, apiKeyResponse.Secret, nil, http.StatusOK, &imageEditResponse.Task) if imageEditResponse.Task.Status != "succeeded" || imageEditResponse.Task.Result["id"] == "" { @@ -2069,16 +2130,17 @@ func doJSONWithHeaders(t *testing.T, baseURL string, method string, path string, func doAPIV1ChatCompletionAndLoadTask(t *testing.T, ctx context.Context, pool *pgxpool.Pool, baseURL string, token string, payload map[string]any, marker string, expectedStatus int, responseOut any, taskDetailOut any) string { t.Helper() - _ = ctx - _ = pool - _ = marker if responseOut == nil { responseOut = &map[string]any{} } + requestStartedAt := time.Now().UTC().Add(-time.Second) responseHeaders := doJSONWithHeaders(t, baseURL, http.MethodPost, "/api/v1/chat/completions", token, payload, nil, expectedStatus, responseOut) taskID := strings.TrimSpace(responseHeaders.Get("X-Gateway-Task-Id")) if taskID == "" { - t.Fatalf("chat completion response did not expose X-Gateway-Task-Id: headers=%v body=%+v", responseHeaders, responseOut) + taskID = waitForLatestTaskIDSince(t, ctx, pool, "chat.completions", requestStartedAt, 2*time.Second) + } + if taskID == "" { + t.Fatalf("load chat completion task %s failed: headers=%v body=%+v", marker, responseHeaders, responseOut) } if taskDetailOut != nil { doJSON(t, baseURL, http.MethodGet, "/api/v1/tasks/"+taskID, token, nil, http.StatusOK, taskDetailOut) @@ -2086,6 +2148,27 @@ func doAPIV1ChatCompletionAndLoadTask(t *testing.T, ctx context.Context, pool *p return taskID } +func waitForLatestTaskIDSince(t *testing.T, ctx context.Context, pool *pgxpool.Pool, kind string, since time.Time, timeout time.Duration) string { + t.Helper() + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + var taskID string + err := pool.QueryRow(ctx, ` +SELECT id::text +FROM gateway_tasks +WHERE kind = $1 + AND created_at >= $2 +ORDER BY created_at DESC +LIMIT 1`, kind, since).Scan(&taskID) + if err == nil && taskID != "" { + return taskID + } + time.Sleep(10 * time.Millisecond) + } + t.Fatalf("latest %s task was not persisted within %s", kind, timeout) + return "" +} + type taskWaitDetail struct { ID string `json:"id"` Status string `json:"status"` diff --git a/apps/api/internal/httpapi/easyai_compat.go b/apps/api/internal/httpapi/easyai_compat.go new file mode 100644 index 0000000..2a25a0e --- /dev/null +++ b/apps/api/internal/httpapi/easyai_compat.go @@ -0,0 +1,489 @@ +package httpapi + +import ( + "encoding/json" + "net/http" + "strings" + + "github.com/easyai/easyai-ai-gateway/apps/api/internal/runner" + "github.com/easyai/easyai-ai-gateway/apps/api/internal/store" +) + +func easyAIAsyncMediaRequest(kind string, r *http.Request) bool { + if r == nil || !asyncRequest(r) { + return false + } + switch kind { + case "images.generations", + "images.edits", + "images.vectorize", + "videos.generations", + "videos.upscales", + "song.generations", + "music.generations", + "speech.generations", + "voice.clone": + return true + default: + return false + } +} + +func easyAITaskAcceptedResponse(task store.GatewayTask) map[string]any { + resultPath := "/api/v1/ai/result/" + task.ID + return map[string]any{ + "status": "submitted", + "task_id": task.ID, + "taskId": task.ID, + "query_url": resultPath, + "task": task, + "next": map[string]string{ + "events": "/api/v1/tasks/" + task.ID + "/events", + "detail": "/api/v1/tasks/" + task.ID, + "result": resultPath, + }, + } +} + +func easyAITaskAcceptedResponseWithExtensions(task store.GatewayTask, extensions map[string]any) map[string]any { + response := cloneEasyAIMap(extensions) + for key, value := range easyAITaskAcceptedResponse(task) { + response[key] = value + } + return response +} + +func writeEasyAIAsyncError(w http.ResponseWriter, status int, message string, details map[string]any, codes ...string) { + code := "" + if len(codes) > 0 { + code = strings.TrimSpace(codes[0]) + } + errorPayload := map[string]any{ + "message": message, + "status": status, + } + if code != "" { + errorPayload["code"] = code + } + for key, value := range details { + errorPayload[key] = value + } + response := map[string]any{ + "status": "failed", + "message": message, + "data": []any{}, + "error": errorPayload, + } + if code != "" { + response["code"] = code + } + writeJSON(w, status, response) +} + +func easyAISynchronousTaskResponse(task store.GatewayTask, output map[string]any) map[string]any { + if task.Kind != "images.vectorize" { + return output + } + response := cloneEasyAIMap(output) + response["taskId"] = task.ID + response["task_id"] = task.ID + response["query_url"] = "/api/v1/ai/result/" + task.ID + return response +} + +func easyAIFileUploadResponse(upload map[string]any) map[string]any { + response := cloneEasyAIMap(upload) + response["status"] = "success" + response["data"] = cloneEasyAIMap(upload) + return response +} + +func easyAITaskResultResponse(task store.GatewayTask) map[string]any { + sourceResult := cloneEasyAIMap(task.Result) + cleanResult := cloneEasyAIMap(sourceResult) + delete(cleanResult, "raw") + delete(cleanResult, "raw_data") + normalizeEasyAIInlineMediaFields(cleanResult) + + data := easyAITaskResultData(task, sourceResult) + status := easyAITaskResultStatus(task.Status) + if status == "failed" { + data = []any{} + } + cleanResult["data"] = data + cleanResult["output_content"] = data + + response := cloneEasyAIMap(cleanResult) + response["status"] = status + response["task_id"] = task.ID + response["taskId"] = task.ID + response["query_url"] = "/api/v1/ai/result/" + task.ID + response["created"] = task.CreatedAt.UnixMilli() + response["data"] = data + response["output_content"] = data + response["output"] = easyAIOutputURLs(data) + response["result"] = cleanResult + + upstreamTaskID := firstNonEmpty( + easyAIString(cleanResult["upstream_task_id"]), + task.RemoteTaskID, + ) + if upstreamTaskID != "" { + response["upstream_task_id"] = upstreamTaskID + } + if len(task.Usage) > 0 && response["usage"] == nil { + response["usage"] = task.Usage + } + + cancelState := runner.DescribeTaskCancellation(task) + response["cancellable"] = cancelState.Cancellable + response["submitted"] = cancelState.Submitted + + message := firstNonEmpty( + easyAIString(cleanResult["message"]), + task.ErrorMessage, + task.Error, + task.Message, + ) + if message == "" { + switch status { + case "submitted": + message = "任务已提交,请稍后查看结果" + case "process": + message = "任务处理中" + case "failed": + message = "任务执行失败" + } + } + if message != "" { + response["message"] = message + } + if code := firstNonEmpty(task.ErrorCode, easyAIString(cleanResult["code"])); code != "" { + response["code"] = code + } + return response +} + +func easyAITaskResultStatus(status string) string { + switch strings.ToLower(strings.TrimSpace(status)) { + case "succeeded", "success", "completed": + return "success" + case "failed", "cancelled", "canceled", "rejected": + return "failed" + case "running", "processing", "in_progress": + return "process" + default: + return "submitted" + } +} + +func easyAITaskResultData(task store.GatewayTask, result map[string]any) []any { + for _, source := range easyAIResultSources(result) { + if items := easyAIOutputArray(source); len(items) > 0 { + return normalizeEasyAIOutputItems(task, items) + } + } + if item := easyAIOutputMap(result); len(item) > 0 { + normalizeEasyAIInlineMediaFields(item) + if easyAIOutputURL(item) != "" || + easyAIOutputHasInlineMedia(item) || + (strings.EqualFold(easyAIString(item["type"]), "text") && (item["text"] != nil || item["content"] != nil)) { + return normalizeEasyAIOutputItems(task, []any{item}) + } + } + return []any{} +} + +func easyAIResultSources(result map[string]any) []any { + sources := []any{ + result["data"], + result["output_content"], + easyAIStructuredOutput(result["content"]), + result["images"], + result["videos"], + result["audios"], + result["files"], + result["output"], + } + if raw, ok := result["raw"].(map[string]any); ok { + sources = append(sources, + raw["data"], + raw["output_content"], + easyAIStructuredOutput(raw["content"]), + raw["images"], + raw["videos"], + raw["audios"], + raw["files"], + raw["output"], + ) + } + return sources +} + +func easyAIStructuredOutput(value any) any { + switch value.(type) { + case []any, []string, map[string]any: + return value + default: + return nil + } +} + +func easyAIOutputArray(value any) []any { + switch typed := value.(type) { + case []any: + return typed + case []string: + items := make([]any, 0, len(typed)) + for _, item := range typed { + items = append(items, item) + } + return items + case map[string]any, string: + return []any{typed} + default: + return nil + } +} + +func normalizeEasyAIOutputItems(task store.GatewayTask, items []any) []any { + normalized := make([]any, 0, len(items)) + for _, item := range items { + output := easyAIOutputMap(item) + if len(output) == 0 { + continue + } + delete(output, "raw_data") + normalizeEasyAIInlineMediaFields(output) + + mediaURL := easyAIOutputURL(output) + if mediaURL != "" { + output["url"] = mediaURL + } + if strings.TrimSpace(easyAIString(output["type"])) == "" { + if outputType := easyAIOutputType(task, output, mediaURL); outputType != "" { + output["type"] = outputType + } + } + normalized = append(normalized, output) + } + return normalized +} + +func easyAIOutputMap(value any) map[string]any { + switch typed := value.(type) { + case map[string]any: + return cloneEasyAIMap(typed) + case string: + trimmed := strings.TrimSpace(typed) + if trimmed == "" { + return nil + } + if strings.HasPrefix(strings.ToLower(trimmed), "data:") { + return map[string]any{"url": trimmed} + } + if strings.HasPrefix(trimmed, "/") || + strings.HasPrefix(strings.ToLower(trimmed), "http://") || + strings.HasPrefix(strings.ToLower(trimmed), "https://") { + return map[string]any{"url": trimmed} + } + return map[string]any{"type": "text", "content": trimmed} + default: + return nil + } +} + +func easyAIOutputURL(item map[string]any) string { + for _, key := range []string{ + "url", + "generated_url", + "generatedUrl", + "output_url", + "outputUrl", + "download_url", + "downloadUrl", + "image_url", + "video_url", + "audio_url", + } { + if value := easyAINestedURL(item[key]); value != "" { + return value + } + } + return "" +} + +func easyAINestedURL(value any) string { + var candidate string + switch typed := value.(type) { + case string: + candidate = strings.TrimSpace(typed) + case map[string]any: + candidate = strings.TrimSpace(easyAIString(typed["url"])) + default: + return "" + } + if strings.HasPrefix(strings.ToLower(candidate), "data:") { + return "" + } + return candidate +} + +func normalizeEasyAIInlineMediaFields(item map[string]any) { + for _, key := range []string{ + "url", + "generated_url", + "generatedUrl", + "output_url", + "outputUrl", + "download_url", + "downloadUrl", + "image_url", + "video_url", + "audio_url", + } { + value := easyAINestedURLAllowData(item[key]) + contentType, payload, ok := easyAIBase64DataURL(value) + if !ok { + continue + } + if strings.HasPrefix(contentType, "image/") { + if easyAIString(item["b64_json"]) == "" { + item["b64_json"] = payload + } + } else if easyAIString(item["content"]) == "" { + item["content"] = value + } + if easyAIString(item["mime_type"]) == "" { + item["mime_type"] = contentType + } + delete(item, key) + } +} + +func easyAIBase64DataURL(value string) (string, string, bool) { + value = strings.TrimSpace(value) + if !strings.HasPrefix(strings.ToLower(value), "data:") { + return "", "", false + } + meta, payload, ok := strings.Cut(value, ",") + if !ok || strings.TrimSpace(payload) == "" { + return "", "", false + } + parts := strings.Split(strings.TrimSpace(meta[5:]), ";") + if len(parts) < 2 || !strings.EqualFold(strings.TrimSpace(parts[len(parts)-1]), "base64") { + return "", "", false + } + contentType := strings.ToLower(strings.TrimSpace(parts[0])) + if !strings.HasPrefix(contentType, "image/") && + !strings.HasPrefix(contentType, "audio/") && + !strings.HasPrefix(contentType, "video/") && + !strings.HasPrefix(contentType, "application/") { + return "", "", false + } + return contentType, strings.TrimSpace(payload), true +} + +func easyAIOutputHasInlineMedia(item map[string]any) bool { + if easyAIString(item["b64_json"]) != "" { + return true + } + content := strings.ToLower(easyAIString(item["content"])) + return strings.HasPrefix(content, "data:") && strings.Contains(content, ";base64,") +} + +func easyAINestedURLAllowData(value any) string { + switch typed := value.(type) { + case string: + return strings.TrimSpace(typed) + case map[string]any: + return strings.TrimSpace(easyAIString(typed["url"])) + default: + return "" + } +} + +func easyAIOutputType(task store.GatewayTask, item map[string]any, mediaURL string) string { + kind := strings.ToLower(strings.TrimSpace(task.Kind)) + switch { + case kind == "images.vectorize": + format := strings.ToLower(firstNonEmpty( + easyAIString(item["format"]), + easyAIString(task.Request["format"]), + )) + if format == "svg" || format == "png" { + return "image" + } + if format != "" { + return "file" + } + case strings.HasPrefix(kind, "images."): + return "image" + case strings.HasPrefix(kind, "videos.") || kind == "video.upscale": + return "video" + case kind == "song.generations" || kind == "music.generations" || kind == "speech.generations": + return "audio" + } + + lowerURL := strings.ToLower(mediaURL) + if index := strings.IndexAny(lowerURL, "?#"); index >= 0 { + lowerURL = lowerURL[:index] + } + switch { + case strings.HasSuffix(lowerURL, ".svg"), + strings.HasSuffix(lowerURL, ".png"), + strings.HasSuffix(lowerURL, ".jpg"), + strings.HasSuffix(lowerURL, ".jpeg"), + strings.HasSuffix(lowerURL, ".webp"), + strings.HasSuffix(lowerURL, ".gif"): + return "image" + case strings.HasSuffix(lowerURL, ".mp4"), + strings.HasSuffix(lowerURL, ".mov"), + strings.HasSuffix(lowerURL, ".webm"): + return "video" + case strings.HasSuffix(lowerURL, ".mp3"), + strings.HasSuffix(lowerURL, ".wav"), + strings.HasSuffix(lowerURL, ".m4a"), + strings.HasSuffix(lowerURL, ".aac"), + strings.HasSuffix(lowerURL, ".flac"): + return "audio" + case mediaURL != "": + return "file" + case item["text"] != nil || item["content"] != nil: + return "text" + default: + return "" + } +} + +func easyAIOutputURLs(data []any) []string { + output := make([]string, 0, len(data)) + for _, item := range data { + if typed, ok := item.(map[string]any); ok { + if mediaURL := easyAIOutputURL(typed); mediaURL != "" { + output = append(output, mediaURL) + } + } + } + return output +} + +func easyAIString(value any) string { + typed, _ := value.(string) + return strings.TrimSpace(typed) +} + +func cloneEasyAIMap(source map[string]any) map[string]any { + if len(source) == 0 { + return map[string]any{} + } + raw, err := json.Marshal(source) + if err != nil { + return map[string]any{} + } + var result map[string]any + if err := json.Unmarshal(raw, &result); err != nil { + return map[string]any{} + } + return result +} diff --git a/apps/api/internal/httpapi/easyai_compat_test.go b/apps/api/internal/httpapi/easyai_compat_test.go new file mode 100644 index 0000000..5a64eee --- /dev/null +++ b/apps/api/internal/httpapi/easyai_compat_test.go @@ -0,0 +1,311 @@ +package httpapi + +import ( + "encoding/json" + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/easyai/easyai-ai-gateway/apps/api/internal/store" +) + +func TestEasyAITaskAcceptedResponseKeepsGatewayFieldsAndAddsLegacyFields(t *testing.T) { + task := store.GatewayTask{ + ID: "task-accepted-1", + Kind: "speech.generations", + Status: "queued", + AsyncMode: true, + } + + got := easyAITaskAcceptedResponse(task) + if got["status"] != "submitted" || got["task_id"] != task.ID || got["taskId"] != task.ID { + t.Fatalf("unexpected EasyAI submission identity: %+v", got) + } + if _, ok := got["task"].(store.GatewayTask); !ok { + t.Fatalf("Gateway task envelope was removed: %+v", got) + } + next, _ := got["next"].(map[string]string) + if next["detail"] != "/api/v1/tasks/"+task.ID || + next["events"] != "/api/v1/tasks/"+task.ID+"/events" || + next["result"] != "/api/v1/ai/result/"+task.ID { + t.Fatalf("unexpected next links: %+v", next) + } +} + +func TestEasyAITaskResultStatus(t *testing.T) { + for _, item := range []struct { + input string + want string + }{ + {input: "queued", want: "submitted"}, + {input: "pending", want: "submitted"}, + {input: "running", want: "process"}, + {input: "processing", want: "process"}, + {input: "succeeded", want: "success"}, + {input: "completed", want: "success"}, + {input: "failed", want: "failed"}, + {input: "cancelled", want: "failed"}, + } { + t.Run(item.input, func(t *testing.T) { + if got := easyAITaskResultStatus(item.input); got != item.want { + t.Fatalf("easyAITaskResultStatus(%q)=%q, want %q", item.input, got, item.want) + } + }) + } +} + +func TestEasyAITaskResultResponseNormalizesMediaOutputs(t *testing.T) { + createdAt := time.Date(2026, 7, 23, 10, 0, 0, 0, time.UTC) + for _, item := range []struct { + name string + task store.GatewayTask + wantType string + wantURL string + wantCount int + }{ + { + name: "image alias", + task: store.GatewayTask{ + ID: "image-1", Kind: "images.generations", Status: "succeeded", CreatedAt: createdAt, + Result: map[string]any{"data": []any{map[string]any{ + "image_url": "https://example.com/output.png", + }}}, + }, + wantType: "image", wantURL: "https://example.com/output.png", wantCount: 1, + }, + { + name: "video raw content", + task: store.GatewayTask{ + ID: "video-1", Kind: "videos.generations", Status: "succeeded", CreatedAt: createdAt, + Result: map[string]any{"raw": map[string]any{"content": map[string]any{ + "video_url": "https://example.com/output.mp4", + "duration": 5, + }}}, + }, + wantType: "video", wantURL: "https://example.com/output.mp4", wantCount: 1, + }, + { + name: "audio alias", + task: store.GatewayTask{ + ID: "audio-1", Kind: "speech.generations", Status: "succeeded", CreatedAt: createdAt, + Result: map[string]any{"data": []any{map[string]any{ + "audio_url": "https://example.com/output.wav", + }}}, + }, + wantType: "audio", wantURL: "https://example.com/output.wav", wantCount: 1, + }, + { + name: "voice clone without media", + task: store.GatewayTask{ + ID: "voice-1", Kind: "voice.clone", Status: "succeeded", CreatedAt: createdAt, + Result: map[string]any{"voice_id": "voice-test-1", "cloned_voice": map[string]any{"voiceId": "voice-test-1"}}, + }, + wantCount: 0, + }, + } { + t.Run(item.name, func(t *testing.T) { + got := easyAITaskResultResponse(item.task) + if got["status"] != "success" || got["task_id"] != item.task.ID || got["created"] != createdAt.UnixMilli() { + t.Fatalf("unexpected task result identity: %+v", got) + } + data, _ := got["data"].([]any) + if len(data) != item.wantCount { + t.Fatalf("unexpected output count: got=%d response=%+v", len(data), got) + } + outputContent, _ := got["output_content"].([]any) + if len(outputContent) != item.wantCount { + t.Fatalf("output_content is not synchronized: %+v", got) + } + if item.wantCount == 0 { + if got["voice_id"] != "voice-test-1" { + t.Fatalf("voice clone fields were lost: %+v", got) + } + return + } + output, _ := data[0].(map[string]any) + if output["type"] != item.wantType || output["url"] != item.wantURL { + t.Fatalf("unexpected normalized output: %+v", output) + } + if _, exposed := output["b64_json"]; exposed { + t.Fatalf("base64 output was exposed: %+v", output) + } + urls, _ := got["output"].([]string) + if len(urls) != 1 || urls[0] != item.wantURL { + t.Fatalf("unexpected output URL list: %+v", got["output"]) + } + }) + } +} + +func TestEasyAITaskResultResponsePreservesBase64ImageOutput(t *testing.T) { + base64Payload := "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAusB9Y9ZlB8AAAAASUVORK5CYII=" + task := store.GatewayTask{ + ID: "image-base64-1", Kind: "images.generations", Status: "succeeded", + Result: map[string]any{"data": []any{map[string]any{ + "b64_json": base64Payload, + "mime_type": "image/png", + }}}, + } + + got := easyAITaskResultResponse(task) + data, _ := got["data"].([]any) + if len(data) != 1 { + t.Fatalf("unexpected base64 output count: %+v", got) + } + output, _ := data[0].(map[string]any) + if output["type"] != "image" || output["b64_json"] != base64Payload || output["mime_type"] != "image/png" { + t.Fatalf("base64 image output was not preserved: %+v", output) + } + if _, exposedAsURL := output["url"]; exposedAsURL { + t.Fatalf("base64 image must not be exposed as a URL: %+v", output) + } + outputContent, _ := got["output_content"].([]any) + content, _ := outputContent[0].(map[string]any) + if content["b64_json"] != base64Payload { + t.Fatalf("output_content lost base64 image: %+v", outputContent) + } + urls, _ := got["output"].([]string) + if len(urls) != 0 { + t.Fatalf("base64-only output must not create URL aliases: %+v", urls) + } +} + +func TestEasyAITaskResultResponseMovesImageDataURLToB64JSON(t *testing.T) { + base64Payload := "iVBORw0KGgoAAAANSUhEUgAAAAEAAAABCAQAAAC1HAwCAAAAC0lEQVR42mP8/x8AAusB9Y9ZlB8AAAAASUVORK5CYII=" + task := store.GatewayTask{ + ID: "image-data-url-1", Kind: "images.generations", Status: "succeeded", + Result: map[string]any{"data": []any{map[string]any{ + "url": "data:image/png;base64," + base64Payload, + }}}, + } + + got := easyAITaskResultResponse(task) + data, _ := got["data"].([]any) + output, _ := data[0].(map[string]any) + if output["b64_json"] != base64Payload || output["mime_type"] != "image/png" || output["type"] != "image" { + t.Fatalf("image data URL was not normalized: %+v", output) + } + if _, exists := output["url"]; exists { + t.Fatalf("image data URL should move out of the URL field: %+v", output) + } +} + +func TestEasyAITaskResultResponsePreservesAudioDataURLAsContent(t *testing.T) { + base64Payload := "SUQzBAAAAAAAI1RTU0UAAAAPAAADTGF2ZjYwLjMuMTAwAAAAAAAAAAAAAAD/" + task := store.GatewayTask{ + ID: "audio-data-url-1", Kind: "speech.generations", Status: "succeeded", + Result: map[string]any{"data": []any{map[string]any{ + "audio_url": "data:audio/mpeg;base64," + base64Payload, + }}}, + } + + got := easyAITaskResultResponse(task) + data, _ := got["data"].([]any) + output, _ := data[0].(map[string]any) + if output["content"] != "data:audio/mpeg;base64,"+base64Payload || + output["mime_type"] != "audio/mpeg" || + output["type"] != "audio" { + t.Fatalf("audio data URL was not normalized to inline content: %+v", output) + } + if _, exists := output["audio_url"]; exists { + t.Fatalf("audio data URL should move out of the URL field: %+v", output) + } + urls, _ := got["output"].([]string) + if len(urls) != 0 { + t.Fatalf("base64-only audio must not create URL aliases: %+v", urls) + } +} + +func TestEasyAITaskResultResponseMapsVectorFormats(t *testing.T) { + for _, item := range []struct { + format string + wantType string + }{ + {format: "svg", wantType: "image"}, + {format: "png", wantType: "image"}, + {format: "pdf", wantType: "file"}, + {format: "eps", wantType: "file"}, + } { + t.Run(item.format, func(t *testing.T) { + task := store.GatewayTask{ + ID: "vector-" + item.format, Kind: "images.vectorize", Status: "succeeded", + Request: map[string]any{"format": item.format}, + Result: map[string]any{"data": []any{map[string]any{"url": "https://example.com/output." + item.format}}}, + } + got := easyAITaskResultResponse(task) + data, _ := got["data"].([]any) + output, _ := data[0].(map[string]any) + if output["type"] != item.wantType { + t.Fatalf("format %s mapped to %v, want %s", item.format, output["type"], item.wantType) + } + }) + } +} + +func TestEasyAITaskResultResponseClearsFailedOutputs(t *testing.T) { + task := store.GatewayTask{ + ID: "failed-1", Kind: "videos.upscales", Status: "cancelled", + ErrorCode: "cancelled", ErrorMessage: "任务已取消", + Result: map[string]any{"data": []any{map[string]any{ + "url": "https://example.com/partial.mp4", + }}}, + } + got := easyAITaskResultResponse(task) + data, _ := got["data"].([]any) + outputContent, _ := got["output_content"].([]any) + if got["status"] != "failed" || got["code"] != "cancelled" || got["message"] != "任务已取消" || + len(data) != 0 || len(outputContent) != 0 { + t.Fatalf("unexpected failed response: %+v", got) + } +} + +func TestEasyAISynchronousVectorAndFileUploadResponsesAreAdditive(t *testing.T) { + vectorTask := store.GatewayTask{ID: "vector-sync-1", Kind: "images.vectorize"} + vector := easyAISynchronousTaskResponse(vectorTask, map[string]any{ + "status": "success", + "data": []any{map[string]any{"url": "https://example.com/output.svg"}}, + }) + if vector["status"] != "success" || vector["taskId"] != vectorTask.ID || vector["task_id"] != vectorTask.ID { + t.Fatalf("unexpected vector response: %+v", vector) + } + + upload := easyAIFileUploadResponse(map[string]any{ + "id": "file-1", "url": "https://example.com/upload.png", "filename": "upload.png", + }) + data, _ := upload["data"].(map[string]any) + if upload["status"] != "success" || upload["url"] != data["url"] || upload["id"] != data["id"] { + t.Fatalf("unexpected upload response: %+v", upload) + } +} + +func TestEasyAIAsyncMediaRequestRequiresExplicitHeader(t *testing.T) { + syncRequest := httptest.NewRequest(http.MethodPost, "/api/v1/images/generations", nil) + if easyAIAsyncMediaRequest("images.generations", syncRequest) { + t.Fatal("image request without X-Async must stay synchronous") + } + asyncRequest := httptest.NewRequest(http.MethodPost, "/api/v1/images/generations", nil) + asyncRequest.Header.Set("X-Async", "true") + if !easyAIAsyncMediaRequest("images.generations", asyncRequest) { + t.Fatal("image request with X-Async must enable EasyAI async compatibility") + } + if easyAIAsyncMediaRequest("chat.completions", asyncRequest) { + t.Fatal("chat completions must not enter EasyAI media async compatibility") + } +} + +func TestWriteEasyAIAsyncErrorKeepsNestedGatewayError(t *testing.T) { + recorder := httptest.NewRecorder() + writeEasyAIAsyncError(recorder, http.StatusBadRequest, "invalid request", map[string]any{"param": "model"}, "invalid_parameter") + + var got map[string]any + if err := json.NewDecoder(recorder.Body).Decode(&got); err != nil { + t.Fatal(err) + } + errorPayload, _ := got["error"].(map[string]any) + data, _ := got["data"].([]any) + if got["status"] != "failed" || got["code"] != "invalid_parameter" || + errorPayload["message"] != "invalid request" || len(data) != 0 { + t.Fatalf("unexpected async error response: %+v", got) + } +} diff --git a/apps/api/internal/httpapi/file_upload_handlers.go b/apps/api/internal/httpapi/file_upload_handlers.go index 07f49ef..0a02b14 100644 --- a/apps/api/internal/httpapi/file_upload_handlers.go +++ b/apps/api/internal/httpapi/file_upload_handlers.go @@ -62,7 +62,7 @@ func (s *Server) uploadFile(w http.ResponseWriter, r *http.Request) { writeError(w, status, err.Error()) return } - writeJSON(w, http.StatusOK, upload) + writeJSON(w, http.StatusOK, easyAIFileUploadResponse(upload)) } func firstNonEmptyFormValue(r *http.Request, key string, fallback string) string { diff --git a/apps/api/internal/httpapi/handlers.go b/apps/api/internal/httpapi/handlers.go index 2a69032..a6879b8 100644 --- a/apps/api/internal/httpapi/handlers.go +++ b/apps/api/internal/httpapi/handlers.go @@ -1039,15 +1039,14 @@ func (s *Server) listModelRateLimitStatuses(w http.ResponseWriter, r *http.Reque // @Failure 429 {object} ErrorEnvelope // @Failure 502 {object} ErrorEnvelope // @Router /api/v1/reranks [post] -// @Router /api/v1/videos/generations [post] -// @Router /api/v1/song/generations [post] -// @Router /api/v1/music/generations [post] -// @Router /api/v1/speech/generations [post] -// @Router /api/v1/voice_clone [post] func (s *Server) createTask(kind string, compatible bool) http.Handler { return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { targetProtocol := targetProtocolForTaskRequest(kind, r) writeTaskError := func(status int, message string, details map[string]any, codes ...string) { + if easyAIAsyncMediaRequest(kind, r) { + writeEasyAIAsyncError(w, status, message, details, codes...) + return + } if targetProtocol == "" { writeErrorWithDetails(w, status, message, details, codes...) return @@ -1093,9 +1092,10 @@ func (s *Server) createTask(kind string, compatible bool) http.Handler { return } responsePlan := planTaskResponse(kind, compatible, body, r) - if targetProtocol != "" { + if targetProtocol != "" && !easyAIAsyncMediaRequest(kind, r) { // OpenAI compatibility endpoints are synchronous APIs. Gateway async - // task envelopes belong only to native Gateway routes. + // task envelopes belong only to native Gateway routes, except for the + // explicit EasyAI media X-Async extension. responsePlan.asyncMode = false } idempotencyKey, hasIdempotencyKey, err := optionalTaskIdempotencyKey(r) @@ -1150,7 +1150,7 @@ func (s *Server) createTask(kind string, compatible bool) http.Handler { writeTaskError(http.StatusConflict, "streaming idempotent replay is not supported", nil, "idempotency_stream_replay_unsupported") return } - if targetProtocol != "" { + if targetProtocol != "" && !responsePlan.asyncMode { if task.Status == "succeeded" { writeJSON(w, http.StatusOK, task.Result) return @@ -1274,7 +1274,7 @@ func openAIEmbeddingsDoc() {} // @Security BearerAuth // @Param input body TaskRequest true "OpenAI Images 请求" // @Param X-Async header bool false "true 时异步创建任务并返回 202" -// @Success 200 {object} CompatibleResponse +// @Success 200 {object} OpenAIImagesResponse // @Success 202 {object} TaskAcceptedResponse // @Failure 400 {object} OpenAIErrorEnvelope // @Failure 401 {object} OpenAIErrorEnvelope @@ -1284,6 +1284,32 @@ func openAIEmbeddingsDoc() {} // @Router /api/v1/images/edits [post] func openAIImagesDoc() {} +// easyAIMediaTasksDoc godoc +// @Summary 创建 EasyAI 媒体任务 +// @Description 默认同步返回 EasyAI GeneratedResponse;设置 X-Async=true 时返回同时兼容 Gateway 与 server-main EasyAIClient 的异步提交结构。 +// @Tags media +// @Accept json +// @Produce json +// @Security BearerAuth +// @Param X-Async header bool false "true 时异步创建任务并返回 202" +// @Param Idempotency-Key header string false "可选请求幂等键;同一用户范围内唯一" +// @Param input body TaskRequest true "媒体任务请求,字段随任务类型变化" +// @Success 200 {object} EasyAIGeneratedResponse +// @Success 202 {object} TaskAcceptedResponse +// @Failure 400 {object} ErrorEnvelope +// @Failure 401 {object} ErrorEnvelope +// @Failure 402 {object} ErrorEnvelope +// @Failure 403 {object} ErrorEnvelope +// @Failure 404 {object} ErrorEnvelope +// @Failure 429 {object} ErrorEnvelope +// @Failure 502 {object} ErrorEnvelope +// @Router /api/v1/videos/generations [post] +// @Router /api/v1/song/generations [post] +// @Router /api/v1/music/generations [post] +// @Router /api/v1/speech/generations [post] +// @Router /api/v1/voice_clone [post] +func easyAIMediaTasksDoc() {} + func (s *Server) requestExecutionContext(r *http.Request) (context.Context, context.CancelFunc) { base := context.WithoutCancel(r.Context()) if s.ctx == nil { @@ -1418,7 +1444,7 @@ func writeProtocolCompatibleTaskResponse(runCtx context.Context, w http.Response writeWireResponse(w, result.Wire) return } - writeJSON(w, http.StatusOK, result.Output) + writeJSON(w, http.StatusOK, easyAISynchronousTaskResponse(result.Task, result.Output)) } func streamIncludeUsage(body map[string]any) bool { @@ -1461,14 +1487,7 @@ func streamCompatibleKind(kind string) bool { } func writeTaskAccepted(w http.ResponseWriter, task store.GatewayTask) { - writeJSON(w, http.StatusAccepted, map[string]any{ - "taskId": task.ID, - "task": task, - "next": map[string]string{ - "events": fmt.Sprintf("/api/v1/tasks/%s/events", task.ID), - "detail": fmt.Sprintf("/api/v1/tasks/%s", task.ID), - }, - }) + writeJSON(w, http.StatusAccepted, easyAITaskAcceptedResponse(task)) } func apiKeyScopeAllowed(user *auth.User, kind string) bool { diff --git a/apps/api/internal/httpapi/openapi_models.go b/apps/api/internal/httpapi/openapi_models.go index 24701de..9141dd4 100644 --- a/apps/api/internal/httpapi/openapi_models.go +++ b/apps/api/internal/httpapi/openapi_models.go @@ -208,6 +208,17 @@ type FileUploadResponse struct { ContentType string `json:"contentType,omitempty" example:"image/png"` Size int `json:"size,omitempty" example:"1024"` AssetStorage map[string]interface{} `json:"assetStorage,omitempty"` + Status string `json:"status" example:"success"` + Data FileUploadData `json:"data"` +} + +type FileUploadData struct { + ID string `json:"id,omitempty" example:"file_abc123"` + URL string `json:"url,omitempty" example:"/static/uploaded/upload-abc123.png"` + Filename string `json:"filename,omitempty" example:"image.png"` + ContentType string `json:"contentType,omitempty" example:"image/png"` + Size int `json:"size,omitempty" example:"1024"` + AssetStorage map[string]interface{} `json:"assetStorage,omitempty"` } type ReplacePlatformModelsRequest struct { @@ -215,9 +226,12 @@ type ReplacePlatformModelsRequest struct { } type TaskAcceptedResponse struct { - TaskID string `json:"taskId" example:"9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25"` - Task store.GatewayTask `json:"task"` - Next TaskNextLinks `json:"next"` + Status string `json:"status" example:"submitted"` + LegacyTaskID string `json:"task_id" example:"9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25"` + TaskID string `json:"taskId" example:"9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25"` + QueryURL string `json:"query_url" example:"/api/v1/ai/result/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25"` + Task store.GatewayTask `json:"task"` + Next TaskNextLinks `json:"next"` } type TaskCancelResponse struct { @@ -231,6 +245,56 @@ type TaskCancelResponse struct { type TaskNextLinks struct { Events string `json:"events" example:"/api/v1/tasks/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25/events"` Detail string `json:"detail" example:"/api/v1/tasks/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25"` + Result string `json:"result" example:"/api/v1/ai/result/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25"` +} + +type EasyAIGeneratedResponse struct { + Status string `json:"status" example:"success" enums:"submitted,process,success,failed"` + TaskID string `json:"task_id" example:"9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25"` + CamelTaskID string `json:"taskId,omitempty" example:"9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25"` + UpstreamTaskID string `json:"upstream_task_id,omitempty" example:"provider-task-123"` + QueryURL string `json:"query_url,omitempty" example:"/api/v1/ai/result/9f4d8f3d-5f5f-4bb7-a4be-344a9f930e25"` + Created int64 `json:"created" example:"1784772000000"` + Message string `json:"message,omitempty"` + Code string `json:"code,omitempty"` + Data []EasyAIMediaOutput `json:"data"` + Output []string `json:"output"` + OutputContent []EasyAIMediaOutput `json:"output_content"` + Cancellable bool `json:"cancellable"` + Submitted bool `json:"submitted"` + Result map[string]interface{} `json:"result,omitempty"` + Usage map[string]interface{} `json:"usage,omitempty"` +} + +type EasyAIMediaOutput struct { + Type string `json:"type,omitempty" example:"video" enums:"image,video,audio,file,text"` + URL string `json:"url,omitempty" example:"https://cdn.example.com/output.mp4"` + // B64JSON is retained when result media transfer is disabled. + B64JSON string `json:"b64_json,omitempty"` + ImageURL string `json:"image_url,omitempty"` + VideoURL string `json:"video_url,omitempty"` + AudioURL string `json:"audio_url,omitempty"` + Content string `json:"content,omitempty"` + Text string `json:"text,omitempty"` + MimeType string `json:"mime_type,omitempty"` + Format string `json:"format,omitempty"` + Duration float64 `json:"duration,omitempty"` + Width int `json:"width,omitempty"` + Height int `json:"height,omitempty"` + SourceResolution string `json:"source_resolution,omitempty"` + TargetResolution string `json:"target_resolution,omitempty"` + Metadata map[string]interface{} `json:"metadata,omitempty"` +} + +type OpenAIImagesResponse struct { + Created int64 `json:"created" example:"1710000000"` + Data []OpenAIImageData `json:"data"` +} + +type OpenAIImageData struct { + URL string `json:"url,omitempty" example:"https://cdn.example.com/output.png"` + B64JSON string `json:"b64_json,omitempty"` + RevisedPrompt string `json:"revised_prompt,omitempty"` } type TaskAcceptedEvent struct { diff --git a/apps/api/internal/httpapi/server.go b/apps/api/internal/httpapi/server.go index 13a0734..818d69c 100644 --- a/apps/api/internal/httpapi/server.go +++ b/apps/api/internal/httpapi/server.go @@ -268,7 +268,7 @@ func NewServerWithContext(ctx context.Context, cfg config.Config, db *store.Stor mux.Handle("POST /api/v1/videos/upscales", server.requireUser(auth.PermissionBasic, server.createVideoUpscaleTask())) mux.Handle("POST /api/v1/video/upscale", server.requireUser(auth.PermissionBasic, server.createVideoUpscaleTask())) mux.Handle("POST /api/v1/video/generations", server.requireUser(auth.PermissionBasic, http.HandlerFunc(server.createLegacyVolcesVideoGeneration))) - mux.Handle("GET /api/v1/ai/result/{taskID}", server.requireUser(auth.PermissionBasic, http.HandlerFunc(server.getLegacyVolcesVideoResult))) + mux.Handle("GET /api/v1/ai/result/{taskID}", server.requireUser(auth.PermissionBasic, http.HandlerFunc(server.getEasyAITaskResult))) mux.Handle("POST /api/v1/song/generations", server.requireUser(auth.PermissionBasic, server.createTask("song.generations", true))) mux.Handle("POST /api/v1/music/generations", server.requireUser(auth.PermissionBasic, server.createTask("music.generations", true))) mux.Handle("POST /api/v1/speech/generations", server.requireUser(auth.PermissionBasic, server.createTask("speech.generations", true))) diff --git a/apps/api/internal/httpapi/task_idempotency_test.go b/apps/api/internal/httpapi/task_idempotency_test.go index 4d49609..ddf8bcb 100644 --- a/apps/api/internal/httpapi/task_idempotency_test.go +++ b/apps/api/internal/httpapi/task_idempotency_test.go @@ -1,8 +1,12 @@ package httpapi import ( + "encoding/json" + "net/http" "net/http/httptest" "testing" + + "github.com/easyai/easyai-ai-gateway/apps/api/internal/store" ) func TestOptionalTaskIdempotencyKeyRejectsMultipleValues(t *testing.T) { @@ -29,3 +33,29 @@ func TestTaskIdempotencyRequestHashIsCanonical(t *testing.T) { t.Fatal("async response semantics must be part of the request hash") } } + +func TestWriteIdempotentAsyncTaskReplayKeepsEasyAICompatibilityFields(t *testing.T) { + t.Parallel() + task := store.GatewayTask{ + ID: "async-replay-1", + Kind: "images.generations", + Status: "queued", + AsyncMode: true, + } + recorder := httptest.NewRecorder() + + writeIdempotentTaskReplay(recorder, task, true) + + if recorder.Code != http.StatusAccepted || + recorder.Header().Get("Idempotent-Replayed") != "true" || + recorder.Header().Get("X-Gateway-Task-Id") != task.ID { + t.Fatalf("unexpected replay status/headers: code=%d headers=%v", recorder.Code, recorder.Header()) + } + var got map[string]any + if err := json.NewDecoder(recorder.Body).Decode(&got); err != nil { + t.Fatal(err) + } + if got["status"] != "submitted" || got["task_id"] != task.ID || got["taskId"] != task.ID { + t.Fatalf("unexpected replay body: %+v", got) + } +} diff --git a/apps/api/internal/httpapi/volces_compat_handlers.go b/apps/api/internal/httpapi/volces_compat_handlers.go index 96be443..7a8e9a9 100644 --- a/apps/api/internal/httpapi/volces_compat_handlers.go +++ b/apps/api/internal/httpapi/volces_compat_handlers.go @@ -143,11 +143,11 @@ func (s *Server) deleteVolcesContentsGenerationTask(w http.ResponseWriter, r *ht // createLegacyVolcesVideoGeneration godoc // @Summary 创建 server-main 兼容视频任务 // @Description 兼容 server-main 的 /api/v1/video/generations,返回 submitted 和 task_id;额外保留火山任务字段。 -// @Tags volces-compatible +// @Tags media // @Accept json // @Produce json // @Security BearerAuth -// @Success 200 {object} map[string]any +// @Success 200 {object} TaskAcceptedResponse // @Router /api/v1/video/generations [post] func (s *Server) createLegacyVolcesVideoGeneration(w http.ResponseWriter, r *http.Request) { user, ok := auth.UserFromContext(r.Context()) @@ -166,35 +166,40 @@ func (s *Server) createLegacyVolcesVideoGeneration(w http.ResponseWriter, r *htt return } response := volcesCompatibleTask(task) - response["status"] = "submitted" - response["task_id"] = task.ID - writeJSON(w, http.StatusOK, response) + writeJSON(w, http.StatusOK, easyAITaskAcceptedResponseWithExtensions(task, response)) } -// getLegacyVolcesVideoResult godoc -// @Summary 查询 server-main 兼容视频结果 -// @Tags volces-compatible +// getEasyAITaskResult godoc +// @Summary 查询 server-main EasyAIClient 兼容任务结果 +// @Description 查询当前用户的图片、视频、音频、视频高清和矢量化任务,并返回 EasyAI GeneratedResponse 兼容结构。 +// @Tags tasks // @Produce json // @Security BearerAuth // @Param taskID path string true "任务 ID" -// @Success 200 {object} map[string]any +// @Success 200 {object} EasyAIGeneratedResponse +// @Failure 404 {object} EasyAIGeneratedResponse // @Router /api/v1/ai/result/{taskID} [get] -func (s *Server) getLegacyVolcesVideoResult(w http.ResponseWriter, r *http.Request) { - task, ok := s.volcesCompatibleTaskForUser(w, r) - if !ok { +func (s *Server) getEasyAITaskResult(w http.ResponseWriter, r *http.Request) { + user, ok := auth.UserFromContext(r.Context()) + if !ok || user == nil { + writeEasyAIAsyncError(w, http.StatusUnauthorized, "unauthorized", nil, "unauthorized") return } - compat := volcesCompatibleTask(task) - legacyStatus := "process" - switch compat["status"] { - case "succeeded": - legacyStatus = "success" - case "failed", "cancelled": - legacyStatus = "failed" + task, err := s.store.GetTask(r.Context(), strings.TrimSpace(r.PathValue("taskID"))) + if err != nil { + if store.IsNotFound(err) { + writeEasyAIAsyncError(w, http.StatusNotFound, "task not found", nil, "not_found") + return + } + s.logger.Error("get EasyAI-compatible task failed", "error", err) + writeError(w, http.StatusInternalServerError, "get task failed", "internal_error") + return } - writeJSON(w, http.StatusOK, map[string]any{ - "status": legacyStatus, "task_id": task.ID, "data": compat["content"], "result": compat, - }) + if !runner.TaskAccessibleToUser(task, user) { + writeEasyAIAsyncError(w, http.StatusNotFound, "task not found", nil, "not_found") + return + } + writeJSON(w, http.StatusOK, easyAITaskResultResponse(task)) } func (s *Server) createVolcesCompatibleTask(r *http.Request, user *auth.User, body map[string]any) (store.GatewayTask, error) { diff --git a/apps/api/internal/runner/upload.go b/apps/api/internal/runner/upload.go index a3fc312..801bf7e 100644 --- a/apps/api/internal/runner/upload.go +++ b/apps/api/internal/runner/upload.go @@ -44,9 +44,9 @@ type FileUploadPayload struct { } type generatedAssetUploadPolicy struct { - UploadInlineMedia bool - UploadURLMedia bool - StoreInlineMediaLocally bool + UploadInlineMedia bool + UploadURLMedia bool + PreserveInlineMedia bool } type generatedAssetDecision struct { @@ -118,7 +118,7 @@ func (s *Service) uploadGeneratedAssets(ctx context.Context, taskID string, task return result, nil } var channels []store.FileStorageChannel - if (needsUpload && generatedAssetNeedsChannelLookup(policy, decisions)) || (rawNeedsUpload && !policy.StoreInlineMediaLocally) { + if needsUpload || rawNeedsUpload { channels, err = s.activeFileStorageChannels(ctx, store.FileStorageSceneImageResult) if err != nil { return nil, &clients.ClientError{Code: "upload_config_failed", Message: err.Error(), Retryable: true} @@ -151,7 +151,7 @@ func (s *Service) uploadGeneratedAssets(ctx context.Context, taskID string, task var contentType string var err error if decision.Inline != nil { - upload, contentType, kind, strategy, err = s.uploadGeneratedAsset(ctx, taskID, decision.Inline, index, channels, policy.StoreInlineMediaLocally) + upload, contentType, kind, strategy, err = s.uploadGeneratedAsset(ctx, taskID, decision.Inline, index, channels) sourceKey = decision.Inline.SourceKey } else { upload, contentType, kind, strategy, err = s.uploadGeneratedURLAsset(ctx, taskID, decision.URL, index, channels) @@ -269,7 +269,7 @@ func (s *Service) uploadGeneratedRawMediaValue(ctx context.Context, taskID strin if !ok { return value, false, nil } - upload, contentType, kind, strategy, err := s.uploadGeneratedAsset(ctx, taskID, asset, *index, channels, policy.StoreInlineMediaLocally) + upload, contentType, kind, strategy, err := s.uploadGeneratedAsset(ctx, taskID, asset, *index, channels) if err != nil { return nil, false, err } @@ -542,25 +542,13 @@ func generatedAssetUploadPolicyFromName(policyName string) generatedAssetUploadP case store.FileStorageResultUploadPolicyUploadAll: return generatedAssetUploadPolicy{UploadInlineMedia: true, UploadURLMedia: true} case store.FileStorageResultUploadPolicyUploadNone: - return generatedAssetUploadPolicy{UploadInlineMedia: true, UploadURLMedia: false, StoreInlineMediaLocally: true} + return generatedAssetUploadPolicy{UploadInlineMedia: false, UploadURLMedia: false, PreserveInlineMedia: true} default: return defaultGeneratedAssetUploadPolicy() } } -func generatedAssetNeedsChannelLookup(policy generatedAssetUploadPolicy, decisions []generatedAssetDecision) bool { - for _, decision := range decisions { - if decision.URL != nil { - return true - } - if decision.Inline != nil && !policy.StoreInlineMediaLocally { - return true - } - } - return false -} - -func (s *Service) uploadGeneratedAsset(ctx context.Context, taskID string, asset *generatedInlineAsset, index int, channels []store.FileStorageChannel, forceLocal bool) (map[string]any, string, string, string, error) { +func (s *Service) uploadGeneratedAsset(ctx context.Context, taskID string, asset *generatedInlineAsset, index int, channels []store.FileStorageChannel) (map[string]any, string, string, string, error) { contentType := resolvedGeneratedAssetContentType(asset.ContentType, asset.Kind, asset.Bytes) kind := generatedAssetKindFromContentType(asset.Kind, contentType) payload := FileUploadPayload{ @@ -570,7 +558,7 @@ func (s *Service) uploadGeneratedAsset(ctx context.Context, taskID string, asset Scene: store.FileStorageSceneImageResult, Source: "ai-gateway", } - if forceLocal || len(channels) == 0 { + if len(channels) == 0 { upload, err := s.storeFileLocally(payload, s.cfg.LocalGeneratedStorageDir, config.DefaultLocalGeneratedStorageDir, localStaticGeneratedPathPrefix) return upload, contentType, kind, "local_static_inline_media", err } @@ -1005,7 +993,9 @@ func generatedAssetDecisionForItem(taskKind string, item map[string]any, policy urlKey, mediaURL := mediaURLSourceFromItem(item) if mediaURL != "" { if !policy.UploadURLMedia { - decision.StripKeys = inlineMediaKeys(item) + if !policy.PreserveInlineMedia { + decision.StripKeys = inlineMediaKeys(item) + } return decision, nil } contentType := mediaContentTypeFromItem(item) diff --git a/apps/api/internal/runner/upload_test.go b/apps/api/internal/runner/upload_test.go index 97ee6fc..7e470bf 100644 --- a/apps/api/internal/runner/upload_test.go +++ b/apps/api/internal/runner/upload_test.go @@ -8,6 +8,7 @@ import ( "path/filepath" "strings" "testing" + "time" "github.com/easyai/easyai-ai-gateway/apps/api/internal/config" "github.com/easyai/easyai-ai-gateway/apps/api/internal/store" @@ -170,7 +171,7 @@ func TestGeneratedAssetDecisionUploadsURLWhenPolicyUploadAll(t *testing.T) { } } -func TestGeneratedAssetDecisionStoresInlineLocallyWhenPolicyUploadNone(t *testing.T) { +func TestGeneratedAssetDecisionPreservesInlineWhenPolicyUploadNone(t *testing.T) { item := map[string]any{ "b64_json": base64.StdEncoding.EncodeToString([]byte("inline image")), } @@ -179,11 +180,26 @@ func TestGeneratedAssetDecisionStoresInlineLocallyWhenPolicyUploadNone(t *testin if err != nil { t.Fatalf("unexpected error: %v", err) } - if decision.Inline == nil || decision.URL != nil { - t.Fatalf("upload_none should still turn inline payloads into static URLs: %+v", decision) + if decision.Inline != nil || decision.URL != nil { + t.Fatalf("upload_none should not transfer inline payloads: %+v", decision) } - if !containsString(decision.StripKeys, "b64_json") { - t.Fatalf("inline payload should be stripped before persistence: %+v", decision.StripKeys) + if len(decision.StripKeys) != 0 { + t.Fatalf("upload_none should preserve b64_json for the caller: %+v", decision.StripKeys) + } +} + +func TestGeneratedAssetDecisionPreservesInlineAlongsideURLWhenPolicyUploadNone(t *testing.T) { + item := map[string]any{ + "url": "https://cdn.example.com/generated.png", + "b64_json": base64.StdEncoding.EncodeToString([]byte("inline image")), + } + + decision, err := generatedAssetDecisionForItem("images.generations", item, generatedAssetUploadPolicyFromName(store.FileStorageResultUploadPolicyUploadNone)) + if err != nil { + t.Fatalf("unexpected error: %v", err) + } + if decision.Inline != nil || decision.URL != nil || len(decision.StripKeys) != 0 { + t.Fatalf("upload_none should preserve the upstream URL and base64 fields unchanged: %+v", decision) } } @@ -196,17 +212,17 @@ func TestGeneratedAssetUploadPolicyFromName(t *testing.T) { { name: "default", policyName: store.FileStorageResultUploadPolicyDefault, - want: generatedAssetUploadPolicy{UploadInlineMedia: true, UploadURLMedia: false, StoreInlineMediaLocally: false}, + want: generatedAssetUploadPolicy{UploadInlineMedia: true, UploadURLMedia: false, PreserveInlineMedia: false}, }, { name: "upload all", policyName: store.FileStorageResultUploadPolicyUploadAll, - want: generatedAssetUploadPolicy{UploadInlineMedia: true, UploadURLMedia: true, StoreInlineMediaLocally: false}, + want: generatedAssetUploadPolicy{UploadInlineMedia: true, UploadURLMedia: true, PreserveInlineMedia: false}, }, { name: "upload none", policyName: store.FileStorageResultUploadPolicyUploadNone, - want: generatedAssetUploadPolicy{UploadInlineMedia: true, UploadURLMedia: false, StoreInlineMediaLocally: true}, + want: generatedAssetUploadPolicy{UploadInlineMedia: false, UploadURLMedia: false, PreserveInlineMedia: true}, }, } @@ -261,7 +277,7 @@ func TestUploadGeneratedAssetStoresLocalWhenNoChannels(t *testing.T) { SourceKey: "b64_json", } - upload, contentType, kind, strategy, err := service.uploadGeneratedAsset(context.Background(), "task-123", asset, 0, nil, false) + upload, contentType, kind, strategy, err := service.uploadGeneratedAsset(context.Background(), "task-123", asset, 0, nil) if err != nil { t.Fatalf("unexpected error: %v", err) } @@ -272,9 +288,14 @@ func TestUploadGeneratedAssetStoresLocalWhenNoChannels(t *testing.T) { if !strings.HasPrefix(urlValue, "/static/generated/gateway-result-task-123-01-") || !strings.HasSuffix(urlValue, ".png") { t.Fatalf("unexpected local static URL: %s", urlValue) } - if stringFromAny(upload["expiresAt"]) == "" { + expiresAt, err := time.Parse(time.RFC3339, stringFromAny(upload["expiresAt"])) + if err != nil { t.Fatalf("local static upload should expose expiresAt: %+v", upload) } + remaining := time.Until(expiresAt) + if remaining < 23*time.Hour+59*time.Minute || remaining > 24*time.Hour+time.Minute { + t.Fatalf("local static upload should expire after one day, remaining=%s", remaining) + } entries, err := os.ReadDir(storageDir) if err != nil { t.Fatalf("failed to read local static dir: %v", err) @@ -301,7 +322,7 @@ func TestUploadGeneratedAssetStoresAudioLocalWhenNoChannels(t *testing.T) { SourceKey: "content", } - upload, contentType, kind, strategy, err := service.uploadGeneratedAsset(context.Background(), "task-tts", asset, 0, nil, false) + upload, contentType, kind, strategy, err := service.uploadGeneratedAsset(context.Background(), "task-tts", asset, 0, nil) if err != nil { t.Fatalf("unexpected error: %v", err) } diff --git a/apps/web/src/pages/admin/SystemSettingsPanel.tsx b/apps/web/src/pages/admin/SystemSettingsPanel.tsx index 83822a2..7a5663a 100644 --- a/apps/web/src/pages/admin/SystemSettingsPanel.tsx +++ b/apps/web/src/pages/admin/SystemSettingsPanel.tsx @@ -51,7 +51,7 @@ const sceneOptions = [ const resultUploadPolicyOptions = [ { value: 'default', label: '默认:仅非链接资源转存', description: 'URL 结果直接保存;base64 / buffer 等结果转存后保存 URL' }, { value: 'upload_all', label: '全部转存', description: 'URL、base64、buffer 等生成媒体结果都会转存到当前文件渠道' }, - { value: 'upload_none', label: '全部不转存', description: '链接结果直接保存;base64 / buffer 结果写入 24 小时本地静态托管后保存 URL' }, + { value: 'upload_none', label: '全部不转存', description: '链接结果直接保存;base64 / buffer 结果保留在对应响应字段中' }, ]; export function SystemSettingsPanel(props: { @@ -176,7 +176,7 @@ export function SystemSettingsPanel(props: {
全局生成媒体转存策略 - 对图片、音频和视频生成结果统一生效;网关本地静态资源仅保留 24 小时。 + 对图片、音频和视频生成结果统一生效;启用转存但没有可用渠道时,回退到保留 24 小时的网关本地静态资源。