diff --git a/.env.example b/.env.example index 3af3e47..8c010a6 100644 --- a/.env.example +++ b/.env.example @@ -84,6 +84,7 @@ AI_GATEWAY_LOCAL_RESULT_TTL_HOURS=24 AI_GATEWAY_LOCAL_RESULT_MIN_FREE_BYTES=10737418240 AI_GATEWAY_LOCAL_RESULT_MAX_BYTES=268435456 AI_GATEWAY_LOCAL_RESULT_MAX_TASK_BYTES=536870912 +AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED=true CORS_ALLOWED_ORIGIN=http://localhost:5178,http://127.0.0.1:5178 VITE_GATEWAY_API_BASE_URL=http://localhost:8088 diff --git a/apps/api/internal/config/config.go b/apps/api/internal/config/config.go index 12a0849..89a7e39 100644 --- a/apps/api/internal/config/config.go +++ b/apps/api/internal/config/config.go @@ -57,6 +57,7 @@ type Config struct { GlobalHTTPProxySource string LogLevel slog.Level BillingEngineMode string + AsyncQueueWorkerEnabled bool AsyncWorkerHardLimit int AsyncWorkerRefreshIntervalSeconds int } @@ -112,6 +113,7 @@ func Load() Config { GlobalHTTPProxySource: globalProxy.Source, LogLevel: logLevel(env("LOG_LEVEL", "info")), BillingEngineMode: strings.ToLower(env("BILLING_ENGINE_MODE", "observe")), + AsyncQueueWorkerEnabled: env("AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED", "true") == "true", AsyncWorkerHardLimit: envIntValidated("AI_GATEWAY_ASYNC_WORKER_HARD_LIMIT", 2048), AsyncWorkerRefreshIntervalSeconds: envIntValidated("AI_GATEWAY_ASYNC_WORKER_REFRESH_INTERVAL_SECONDS", 5), } diff --git a/apps/api/internal/config/config_test.go b/apps/api/internal/config/config_test.go index 481e982..16b6920 100644 --- a/apps/api/internal/config/config_test.go +++ b/apps/api/internal/config/config_test.go @@ -82,6 +82,18 @@ func TestValidateAsyncWorkerSettings(t *testing.T) { } } +func TestLoadAsyncQueueWorkerEnabled(t *testing.T) { + cfg := Load() + if !cfg.AsyncQueueWorkerEnabled { + t.Fatal("async queue worker must be enabled by default") + } + t.Setenv("AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED", "false") + cfg = Load() + if cfg.AsyncQueueWorkerEnabled { + t.Fatal("async queue worker remained enabled after explicit disable") + } +} + func TestValidateTaskHistorySettings(t *testing.T) { cfg := Load() cfg.TaskRetentionDays = 30 diff --git a/apps/api/internal/httpapi/server.go b/apps/api/internal/httpapi/server.go index 00efea0..e0682ee 100644 --- a/apps/api/internal/httpapi/server.go +++ b/apps/api/internal/httpapi/server.go @@ -120,7 +120,11 @@ func NewServerWithContext(ctx context.Context, cfg config.Config, db *store.Stor return server.identityRuntime.LegacyJWTEnabled() } server.auth.LocalAPIKeyVerifier = db.VerifyLocalAPIKey - server.runner.StartAsyncQueueWorker(ctx) + if cfg.AsyncQueueWorkerEnabled { + server.runner.StartAsyncQueueWorker(ctx) + } else if logger != nil { + logger.Info("asynchronous queue worker disabled for this process") + } server.runner.StartBillingSettlementWorker(ctx) server.runner.StartTaskHistoryWorkers(ctx) server.startLocalTempAssetCleanup(ctx) diff --git a/apps/api/internal/store/file_storage_channels.go b/apps/api/internal/store/file_storage_channels.go index 479a4f5..e8c3987 100644 --- a/apps/api/internal/store/file_storage_channels.go +++ b/apps/api/internal/store/file_storage_channels.go @@ -578,7 +578,7 @@ func normalizeFileStorageScene(scene string) string { } func defaultFileStorageScenes() []string { - return []string{FileStorageSceneUpload, FileStorageSceneImageResult} + return []string{FileStorageSceneUpload, FileStorageSceneImageResult, FileStorageSceneRequestAsset} } func defaultFileStorageRetryPolicyIfEmpty(policy map[string]any) map[string]any { diff --git a/apps/api/internal/store/file_storage_channels_test.go b/apps/api/internal/store/file_storage_channels_test.go new file mode 100644 index 0000000..61dbafb --- /dev/null +++ b/apps/api/internal/store/file_storage_channels_test.go @@ -0,0 +1,17 @@ +package store + +import ( + "reflect" + "testing" +) + +func TestDefaultFileStorageScenesIncludeCrossNodeRequestAssets(t *testing.T) { + want := []string{ + FileStorageSceneUpload, + FileStorageSceneImageResult, + FileStorageSceneRequestAsset, + } + if got := defaultFileStorageScenes(); !reflect.DeepEqual(got, want) { + t.Fatalf("default file storage scenes = %#v, want %#v", got, want) + } +} diff --git a/apps/api/migrations/0089_file_storage_request_asset_scene.sql b/apps/api/migrations/0089_file_storage_request_asset_scene.sql new file mode 100644 index 0000000..6b1ad51 --- /dev/null +++ b/apps/api/migrations/0089_file_storage_request_asset_scene.sql @@ -0,0 +1,19 @@ +UPDATE file_storage_channels +SET + config = jsonb_set( + config, + '{scenes}', + CASE + WHEN jsonb_typeof(config->'scenes') = 'array' + THEN (config->'scenes') || '["request_asset"]'::jsonb + ELSE '["upload", "image_result", "request_asset"]'::jsonb + END, + true + ), + updated_at = now() +WHERE deleted_at IS NULL + AND provider = 'server_main_openapi' + AND ( + jsonb_typeof(config->'scenes') IS DISTINCT FROM 'array' + OR NOT (COALESCE(config->'scenes', '[]'::jsonb) ? 'request_asset') + ); diff --git a/apps/web/src/pages/admin/SystemSettingsPanel.tsx b/apps/web/src/pages/admin/SystemSettingsPanel.tsx index c4ac35b..f0c3c11 100644 --- a/apps/web/src/pages/admin/SystemSettingsPanel.tsx +++ b/apps/web/src/pages/admin/SystemSettingsPanel.tsx @@ -42,10 +42,11 @@ const providerOptions = [ { value: 'tencent_cos', label: '腾讯云 COS' }, ]; -const defaultScenes = ['upload', 'image_result']; +const defaultScenes = ['upload', 'image_result', 'request_asset']; const sceneOptions = [ { value: 'upload', label: '上传', description: 'OpenAPI / 管理端主动上传文件' }, { value: 'image_result', label: '生成媒体结果', description: '模型返回 base64 / buffer 图片、音频或视频后的转存' }, + { value: 'request_asset', label: '请求素材', description: '将 multipart / base64 输入转换为跨节点可读取的公共 URL' }, ]; const resultUploadPolicyOptions = [ diff --git a/docker-compose.yml b/docker-compose.yml index ca8df3c..4db52df 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -33,6 +33,7 @@ x-api-environment: &api-environment AI_GATEWAY_LOCAL_RESULT_MIN_FREE_BYTES: ${AI_GATEWAY_LOCAL_RESULT_MIN_FREE_BYTES:-10737418240} AI_GATEWAY_LOCAL_RESULT_MAX_BYTES: ${AI_GATEWAY_LOCAL_RESULT_MAX_BYTES:-268435456} AI_GATEWAY_LOCAL_RESULT_MAX_TASK_BYTES: ${AI_GATEWAY_LOCAL_RESULT_MAX_TASK_BYTES:-536870912} + AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED: ${AI_GATEWAY_ASYNC_QUEUE_WORKER_ENABLED:-true} services: postgres: