提交 17e4823d authored 作者: 王鹏飞's avatar 王鹏飞

feat: 新增AI网关

上级 6a6e6197
...@@ -32,3 +32,10 @@ PERMISSION_APP_ID= ...@@ -32,3 +32,10 @@ PERMISSION_APP_ID=
PERMISSION_APP_SECRET= PERMISSION_APP_SECRET=
SSO_USER_CACHE_TTL_SECONDS=180 SSO_USER_CACHE_TTL_SECONDS=180
PERMISSION_CACHE_TTL_SECONDS=180 PERMISSION_CACHE_TTL_SECONDS=180
# AI 网关上游(密钥只放环境变量,不进数据库)
VOLCANO_BASE_URL=https://ark.cn-beijing.volces.com/api/v3
VOLCANO_API_KEY=
DEEPSEEK_BASE_URL=https://api.deepseek.com/v1
DEEPSEEK_API_KEY=
AI_RATE_LIMIT_PER_MINUTE=120
...@@ -25,3 +25,12 @@ PERMISSION_APP_ID=ezijing_82b452761ada1b79fd12b9fcbd73629e ...@@ -25,3 +25,12 @@ PERMISSION_APP_ID=ezijing_82b452761ada1b79fd12b9fcbd73629e
PERMISSION_APP_SECRET=f19d1e788ac9ea0e3c72ef2e3caee8f0 PERMISSION_APP_SECRET=f19d1e788ac9ea0e3c72ef2e3caee8f0
SSO_USER_CACHE_TTL_SECONDS=180 SSO_USER_CACHE_TTL_SECONDS=180
PERMISSION_CACHE_TTL_SECONDS=180 PERMISSION_CACHE_TTL_SECONDS=180
# AI 网关上游(密钥只放环境变量,不进数据库)
VOLCANO_BASE_URL=https://ark.cn-beijing.volces.com/api/v3
VOLCANO_API_KEY=10c95d49-0368-4fb1-b87c-7668eb0ce67d
DEEPSEEK_BASE_URL=https://api.deepseek.com/v1
DEEPSEEK_API_KEY=sk-f1a6f0a7013241de8393cb2cb108e777
AI_RATE_LIMIT_PER_MINUTE=120
...@@ -27,3 +27,6 @@ Thumbs.db ...@@ -27,3 +27,6 @@ Thumbs.db
# Build # Build
build/ build/
dist/ dist/
# 本地数据备份(清理/迁移前的快照),不入库
backup/
差异被折叠。
差异被折叠。
差异被折叠。
-- DMS schema baseline (MySQL 8+) -- DMS schema baseline (MySQL 8+)
-- Execute this SQL manually after reviewing the target database. -- Execute this SQL manually after reviewing the target database.
-- 已有旧版 DMS 表使用 fix_dms_schema.sql;新建库无需再执行 fix。
CREATE TABLE IF NOT EXISTS `product_list` ( CREATE TABLE IF NOT EXISTS `product_list` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT '主键', `id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT '主键',
...@@ -13,6 +14,7 @@ CREATE TABLE IF NOT EXISTS `product_list` ( ...@@ -13,6 +14,7 @@ CREATE TABLE IF NOT EXISTS `product_list` (
PRIMARY KEY (`id`), PRIMARY KEY (`id`),
UNIQUE KEY `uk_products_name` (`name`), UNIQUE KEY `uk_products_name` (`name`),
KEY `idx_products_status` (`status`), KEY `idx_products_status` (`status`),
KEY `idx_products_updated_at` (`updated_at`),
KEY `idx_products_operator_user_id` (`operator_user_id`) KEY `idx_products_operator_user_id` (`operator_user_id`)
) DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='产品信息'; ) DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='产品信息';
...@@ -25,11 +27,11 @@ CREATE TABLE IF NOT EXISTS `project_list` ( ...@@ -25,11 +27,11 @@ CREATE TABLE IF NOT EXISTS `project_list` (
`school_name` varchar(255), `school_name` varchar(255),
`department_name` varchar(255), `department_name` varchar(255),
`product_id` bigint unsigned COMMENT '关联产品ID', `product_id` bigint unsigned COMMENT '关联产品ID',
`product_name` varchar(255),
`contact_name` varchar(120), `contact_name` varchar(120),
`contact_title` varchar(120), `contact_title` varchar(120),
`contact_phone` varchar(64), `contact_phone` varchar(64),
`solution` text, `solution` text,
`attachment_file_url` text COMMENT '阶段附件JSON数组',
`stage` int NOT NULL DEFAULT 10, `stage` int NOT NULL DEFAULT 10,
`status` int NOT NULL DEFAULT 0 COMMENT '项目状态: 0进行中 20已归档', `status` int NOT NULL DEFAULT 0 COMMENT '项目状态: 0进行中 20已归档',
`description` text, `description` text,
...@@ -40,6 +42,7 @@ CREATE TABLE IF NOT EXISTS `project_list` ( ...@@ -40,6 +42,7 @@ CREATE TABLE IF NOT EXISTS `project_list` (
PRIMARY KEY (`id`), PRIMARY KEY (`id`),
UNIQUE KEY `uk_projects_project_code` (`project_code`), UNIQUE KEY `uk_projects_project_code` (`project_code`),
KEY `idx_projects_stage_status` (`stage`, `status`), KEY `idx_projects_stage_status` (`stage`, `status`),
KEY `idx_projects_created_at` (`created_at`),
KEY `idx_projects_product_id` (`product_id`), KEY `idx_projects_product_id` (`product_id`),
KEY `idx_projects_operator_user_id` (`operator_user_id`) KEY `idx_projects_operator_user_id` (`operator_user_id`)
) DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='项目主表'; ) DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='项目主表';
...@@ -49,7 +52,6 @@ CREATE TABLE IF NOT EXISTS `case_list` ( ...@@ -49,7 +52,6 @@ CREATE TABLE IF NOT EXISTS `case_list` (
`name` varchar(255) NOT NULL COMMENT '案例名称', `name` varchar(255) NOT NULL COMMENT '案例名称',
`description` text COMMENT '案例简介', `description` text COMMENT '案例简介',
`product_id` bigint unsigned COMMENT '关联产品ID', `product_id` bigint unsigned COMMENT '关联产品ID',
`product_name` varchar(255) COMMENT '产品名称',
`files` longtext COMMENT '案例附件JSON', `files` longtext COMMENT '案例附件JSON',
`operator_user_id` varchar(64) COMMENT '最后操作人系统用户ID', `operator_user_id` varchar(64) COMMENT '最后操作人系统用户ID',
`operator_name` varchar(120) COMMENT '最后操作人姓名', `operator_name` varchar(120) COMMENT '最后操作人姓名',
...@@ -57,14 +59,14 @@ CREATE TABLE IF NOT EXISTS `case_list` ( ...@@ -57,14 +59,14 @@ CREATE TABLE IF NOT EXISTS `case_list` (
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间', `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`), PRIMARY KEY (`id`),
KEY `idx_cases_product_id` (`product_id`), KEY `idx_cases_product_id` (`product_id`),
KEY `idx_cases_updated_at` (`updated_at`),
KEY `idx_cases_operator_user_id` (`operator_user_id`) KEY `idx_cases_operator_user_id` (`operator_user_id`)
) DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='案例信息'; ) DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='案例信息';
CREATE TABLE IF NOT EXISTS `project_initiations` ( CREATE TABLE IF NOT EXISTS `project_initiations` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT, `id` bigint unsigned NOT NULL AUTO_INCREMENT,
`project_id` bigint unsigned NOT NULL, `project_id` bigint unsigned NOT NULL,
`application_file_url` text COMMENT '附件JSON数组', `attachment_file_url` text COMMENT '阶段附件JSON数组',
`argument_file_url` text COMMENT '附件JSON数组',
`project_amount` decimal(14,2), `project_amount` decimal(14,2),
`fund_source` varchar(255), `fund_source` varchar(255),
`execution_plan` text, `execution_plan` text,
...@@ -86,8 +88,7 @@ CREATE TABLE IF NOT EXISTS `project_procurements` ( ...@@ -86,8 +88,7 @@ CREATE TABLE IF NOT EXISTS `project_procurements` (
`main_bid_owner` varchar(120), `main_bid_owner` varchar(120),
`companion_bidders` text, `companion_bidders` text,
`formal_bid_status` varchar(120), `formal_bid_status` varchar(120),
`winning_notice_file_url` text COMMENT '附件JSON数组', `attachment_file_url` text COMMENT '阶段附件JSON数组',
`bid_archive_file_url` text COMMENT '附件JSON数组',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (`id`), PRIMARY KEY (`id`),
...@@ -102,7 +103,7 @@ CREATE TABLE IF NOT EXISTS `project_contracts` ( ...@@ -102,7 +103,7 @@ CREATE TABLE IF NOT EXISTS `project_contracts` (
`contract_name` varchar(255), `contract_name` varchar(255),
`amount` decimal(14,2), `amount` decimal(14,2),
`drafter` varchar(120), `drafter` varchar(120),
`archive_file_url` text COMMENT '附件JSON数组', `attachment_file_url` text COMMENT '阶段附件JSON数组',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (`id`), PRIMARY KEY (`id`),
...@@ -117,6 +118,7 @@ CREATE TABLE IF NOT EXISTS `project_deliveries` ( ...@@ -117,6 +118,7 @@ CREATE TABLE IF NOT EXISTS `project_deliveries` (
`delivery_contact` varchar(120), `delivery_contact` varchar(120),
`delivery_contact_phone` varchar(64), `delivery_contact_phone` varchar(64),
`delivery_note` text, `delivery_note` text,
`attachment_file_url` text COMMENT '阶段附件JSON数组',
`completed_at` datetime, `completed_at` datetime,
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
...@@ -128,8 +130,8 @@ CREATE TABLE IF NOT EXISTS `project_deliveries` ( ...@@ -128,8 +130,8 @@ CREATE TABLE IF NOT EXISTS `project_deliveries` (
CREATE TABLE IF NOT EXISTS `project_acceptances` ( CREATE TABLE IF NOT EXISTS `project_acceptances` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT, `id` bigint unsigned NOT NULL AUTO_INCREMENT,
`project_id` bigint unsigned NOT NULL, `project_id` bigint unsigned NOT NULL,
`acceptance_report_url` text COMMENT '附件JSON数组',
`acceptance_note` text, `acceptance_note` text,
`attachment_file_url` text COMMENT '阶段附件JSON数组',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
PRIMARY KEY (`id`), PRIMARY KEY (`id`),
......
-- AI 网关建表(MySQL 8+)
-- 六张表:应用与密钥、模型映射与价格、应用额度、调用明细
-- 均为新表,无存量数据迁移。审阅后执行。
-- 接入应用
CREATE TABLE IF NOT EXISTS `ai_apps` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT '主键',
`code` varchar(64) NOT NULL COMMENT '应用标识,如 center-dms',
`name` varchar(120) NOT NULL COMMENT '应用名称',
`billing_mode` varchar(16) NOT NULL DEFAULT 'internal' COMMENT '结算模式: internal内部使用/quota额度控制',
`status` int NOT NULL DEFAULT 1 COMMENT '状态: 0停用 1启用',
`operator_user_id` varchar(64) DEFAULT NULL COMMENT '最后操作人用户ID',
`operator_name` varchar(120) DEFAULT NULL COMMENT '最后操作人姓名',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_ai_apps_code` (`code`),
KEY `idx_ai_apps_status` (`status`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='AI 网关-接入应用';
-- 应用持有的密钥
CREATE TABLE IF NOT EXISTS `ai_api_keys` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT '主键',
`app_id` bigint unsigned NOT NULL COMMENT '归属应用',
`name` varchar(120) NOT NULL COMMENT '备注名',
`key_hash` char(64) NOT NULL COMMENT 'sha256(明文),明文只在创建时返回一次',
`key_prefix` varchar(12) NOT NULL COMMENT '前 8 位,列表展示与排查用',
`status` int NOT NULL DEFAULT 1 COMMENT '状态: 0停用 1启用',
`last_used_at` datetime DEFAULT NULL COMMENT '最近一次调用时间',
`operator_user_id` varchar(64) DEFAULT NULL COMMENT '最后操作人用户ID',
`operator_name` varchar(120) DEFAULT NULL COMMENT '最后操作人姓名',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_ai_api_keys_hash` (`key_hash`),
KEY `idx_ai_api_keys_app` (`app_id`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='AI 网关-API 密钥';
-- 平台模型别名及上游映射
CREATE TABLE IF NOT EXISTS `ai_model_mappings` (
`alias` varchar(80) NOT NULL COMMENT '平台对外暴露的模型别名',
`name` varchar(120) NOT NULL COMMENT '模型展示名称',
`type` varchar(16) NOT NULL COMMENT '模型类型: text/image/video',
`provider` varchar(32) NOT NULL COMMENT '上游: volcano/deepseek',
`upstream_model` varchar(120) NOT NULL COMMENT '上游真实模型,火山为 endpoint id',
`enabled` int NOT NULL DEFAULT 1 COMMENT '是否启用',
`operator_user_id` varchar(64) DEFAULT NULL COMMENT '最后操作人用户ID',
`operator_name` varchar(120) DEFAULT NULL COMMENT '最后操作人姓名',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`alias`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='AI 网关-模型映射';
-- 模型单价(按别名版本化:改价 = 关闭旧版本 + 开新版本)
CREATE TABLE IF NOT EXISTS `ai_model_prices` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT '主键',
`alias` varchar(80) NOT NULL COMMENT '模型别名',
`version` int NOT NULL DEFAULT 1 COMMENT '版本号,改价=新版本',
`pricing_unit` varchar(24) NOT NULL COMMENT '计价单位: per_1m_tokens/per_image/per_second',
`input_price` decimal(12,6) NOT NULL DEFAULT 0 COMMENT '输入 token 单价',
`cached_input_price` decimal(12,6) NOT NULL DEFAULT 0 COMMENT '缓存命中输入单价(0=按 input_price 计)',
`output_price` decimal(12,6) NOT NULL DEFAULT 0 COMMENT '输出 token 单价',
`unit_price` decimal(12,6) NOT NULL DEFAULT 0 COMMENT '每张/每秒单价',
`currency` varchar(3) NOT NULL DEFAULT 'CNY' COMMENT '币种',
`effective_from` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '生效时间',
`effective_to` datetime DEFAULT NULL COMMENT 'NULL=生效中',
`status` int NOT NULL DEFAULT 1 COMMENT '1生效 0停用',
`operator_user_id` varchar(64) DEFAULT NULL COMMENT '最后操作人用户ID',
`operator_name` varchar(120) DEFAULT NULL COMMENT '最后操作人姓名',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
UNIQUE KEY `uk_ai_model_prices_alias_version` (`alias`, `version`),
KEY `idx_ai_model_prices_alias` (`alias`, `status`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='AI 网关-模型价格(按别名版本化)';
-- 应用额度包。合同在商务系统管理,网关只记录可调用额度。
CREATE TABLE IF NOT EXISTS `ai_quota_packages` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT '主键',
`app_id` bigint unsigned NOT NULL COMMENT '所属应用',
`name` varchar(120) NOT NULL COMMENT '额度名称',
`amount_micro` bigint unsigned NOT NULL COMMENT '总额度(微元)',
`used_micro` bigint unsigned NOT NULL DEFAULT 0 COMMENT '已使用额度(微元)',
`start_date` date NOT NULL COMMENT '生效日期',
`end_date` date NOT NULL COMMENT '失效日期',
`status` int NOT NULL DEFAULT 1 COMMENT '状态: 0停用 1启用',
`remark` text DEFAULT NULL COMMENT '备注',
`operator_user_id` varchar(64) DEFAULT NULL COMMENT '最后操作人用户ID',
`operator_name` varchar(120) DEFAULT NULL COMMENT '最后操作人姓名',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
KEY `idx_ai_quota_packages_app` (`app_id`, `status`, `start_date`, `end_date`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='AI 网关-应用额度包';
-- 调用明细
CREATE TABLE IF NOT EXISTS `ai_usage` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT '主键',
`request_id` varchar(64) DEFAULT NULL COMMENT '关联 x-request-id',
`app_id` bigint unsigned DEFAULT NULL COMMENT '应用',
`api_key_id` bigint unsigned DEFAULT NULL COMMENT '密钥',
`quota_package_id` bigint unsigned DEFAULT NULL COMMENT '本次结算使用的额度包',
`user_id` varchar(64) DEFAULT NULL COMMENT '调用用户(来自请求头 X-User-Id)',
`user_name` varchar(120) DEFAULT NULL COMMENT '调用用户姓名(来自 X-User-Name)',
`type` varchar(16) NOT NULL COMMENT '模型类型: text/image/video',
`provider` varchar(32) NOT NULL COMMENT '上游',
`model` varchar(80) NOT NULL COMMENT '客户端请求的别名',
`upstream_model` varchar(120) NOT NULL COMMENT '上游真实模型',
`price_version` int DEFAULT NULL COMMENT '本次计费使用的价格版本',
`unit` varchar(16) NOT NULL COMMENT '计费单位: token/image/second',
`quantity` bigint unsigned NOT NULL DEFAULT 0 COMMENT '计费数量',
`prompt_tokens` bigint unsigned NOT NULL DEFAULT 0 COMMENT '输入 token(仅文本)',
`cached_tokens` bigint unsigned NOT NULL DEFAULT 0 COMMENT '缓存命中的输入 token',
`completion_tokens` bigint unsigned NOT NULL DEFAULT 0 COMMENT '输出 token(仅文本)',
`reasoning_tokens` bigint unsigned NOT NULL DEFAULT 0 COMMENT '思考 token(含在 completion 内)',
`usage_source` varchar(16) NOT NULL DEFAULT 'reported' COMMENT 'reported/missing,missing 表示成本待核对',
`cost` decimal(12,6) NOT NULL DEFAULT 0 COMMENT '成本',
`status` varchar(16) NOT NULL DEFAULT 'pending' COMMENT 'pending/success/error/aborted',
`upstream_task_id` varchar(128) DEFAULT NULL COMMENT '上游异步任务 ID(视频)',
`http_status` int DEFAULT NULL COMMENT '上游 HTTP 状态',
`error_code` varchar(64) DEFAULT NULL COMMENT '错误码',
`latency_ms` int DEFAULT NULL COMMENT '总耗时',
`first_token_ms` int DEFAULT NULL COMMENT '首 token 延迟(流式)',
`stream` int NOT NULL DEFAULT 0 COMMENT '是否流式',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
KEY `idx_ai_usage_app` (`app_id`, `created_at`),
KEY `idx_ai_usage_key` (`api_key_id`, `created_at`),
KEY `idx_ai_usage_quota_package` (`quota_package_id`, `created_at`),
KEY `idx_ai_usage_model` (`provider`, `model`, `created_at`),
KEY `idx_ai_usage_type` (`type`, `created_at`),
KEY `idx_ai_usage_status` (`status`),
KEY `idx_ai_usage_task` (`app_id`, `api_key_id`, `upstream_task_id`),
KEY `idx_ai_usage_pending` (`api_key_id`, `status`, `unit`, `created_at`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='AI 网关-调用明细';
...@@ -20,12 +20,12 @@ CREATE TABLE `project_list` ( ...@@ -20,12 +20,12 @@ CREATE TABLE `project_list` (
`school_name` varchar(255), `school_name` varchar(255),
`department_name` varchar(255), `department_name` varchar(255),
`product_id` bigint unsigned COMMENT '关联产品ID', `product_id` bigint unsigned COMMENT '关联产品ID',
`product_name` varchar(255),
`contact_name` varchar(120), `contact_name` varchar(120),
`contact_title` varchar(120), `contact_title` varchar(120),
`contact_phone` varchar(64), `contact_phone` varchar(64),
`solution` text, `solution` text,
`attachment_file_url` text,
`stage` int NOT NULL DEFAULT 10, `stage` int NOT NULL DEFAULT 10,
`status` int NOT NULL DEFAULT 0 COMMENT '项目状态: 0进行中 20已归档', `status` int NOT NULL DEFAULT 0 COMMENT '项目状态: 0进行中 20已归档',
`description` text, `description` text,
...@@ -42,7 +42,6 @@ CREATE TABLE `case_list` ( ...@@ -42,7 +42,6 @@ CREATE TABLE `case_list` (
`name` varchar(255) NOT NULL COMMENT '案例名称', `name` varchar(255) NOT NULL COMMENT '案例名称',
`description` text COMMENT '案例简介', `description` text COMMENT '案例简介',
`product_id` bigint unsigned COMMENT '关联产品ID', `product_id` bigint unsigned COMMENT '关联产品ID',
`product_name` varchar(255) COMMENT '产品名称',
`files` longtext COMMENT '案例附件JSON', `files` longtext COMMENT '案例附件JSON',
`operator_user_id` varchar(64) COMMENT '最后操作人系统用户ID', `operator_user_id` varchar(64) COMMENT '最后操作人系统用户ID',
`operator_name` varchar(120) COMMENT '最后操作人姓名', `operator_name` varchar(120) COMMENT '最后操作人姓名',
...@@ -54,8 +53,7 @@ CREATE TABLE `case_list` ( ...@@ -54,8 +53,7 @@ CREATE TABLE `case_list` (
CREATE TABLE `project_initiations` ( CREATE TABLE `project_initiations` (
`id` bigint unsigned AUTO_INCREMENT NOT NULL, `id` bigint unsigned AUTO_INCREMENT NOT NULL,
`project_id` bigint unsigned NOT NULL, `project_id` bigint unsigned NOT NULL,
`application_file_url` text, `attachment_file_url` text,
`argument_file_url` text,
`project_amount` decimal(14,2), `project_amount` decimal(14,2),
`fund_source` varchar(255), `fund_source` varchar(255),
`execution_plan` text, `execution_plan` text,
...@@ -76,8 +74,7 @@ CREATE TABLE `project_procurements` ( ...@@ -76,8 +74,7 @@ CREATE TABLE `project_procurements` (
`main_bid_owner` varchar(120), `main_bid_owner` varchar(120),
`companion_bidders` text, `companion_bidders` text,
`formal_bid_status` varchar(120), `formal_bid_status` varchar(120),
`winning_notice_file_url` text, `attachment_file_url` text,
`bid_archive_file_url` text,
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
CONSTRAINT `project_procurements_id` PRIMARY KEY(`id`), CONSTRAINT `project_procurements_id` PRIMARY KEY(`id`),
...@@ -91,7 +88,7 @@ CREATE TABLE `project_contracts` ( ...@@ -91,7 +88,7 @@ CREATE TABLE `project_contracts` (
`contract_name` varchar(255), `contract_name` varchar(255),
`amount` decimal(14,2), `amount` decimal(14,2),
`drafter` varchar(120), `drafter` varchar(120),
`archive_file_url` text, `attachment_file_url` text,
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
CONSTRAINT `project_contracts_id` PRIMARY KEY(`id`), CONSTRAINT `project_contracts_id` PRIMARY KEY(`id`),
...@@ -105,6 +102,7 @@ CREATE TABLE `project_deliveries` ( ...@@ -105,6 +102,7 @@ CREATE TABLE `project_deliveries` (
`delivery_contact` varchar(120), `delivery_contact` varchar(120),
`delivery_contact_phone` varchar(64), `delivery_contact_phone` varchar(64),
`delivery_note` text, `delivery_note` text,
`attachment_file_url` text,
`completed_at` datetime, `completed_at` datetime,
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
...@@ -115,8 +113,8 @@ CREATE TABLE `project_deliveries` ( ...@@ -115,8 +113,8 @@ CREATE TABLE `project_deliveries` (
CREATE TABLE `project_acceptances` ( CREATE TABLE `project_acceptances` (
`id` bigint unsigned AUTO_INCREMENT NOT NULL, `id` bigint unsigned AUTO_INCREMENT NOT NULL,
`project_id` bigint unsigned NOT NULL, `project_id` bigint unsigned NOT NULL,
`acceptance_report_url` text,
`acceptance_note` text, `acceptance_note` text,
`attachment_file_url` text,
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP, `created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP,
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, `updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP,
CONSTRAINT `project_acceptances_id` PRIMARY KEY(`id`), CONSTRAINT `project_acceptances_id` PRIMARY KEY(`id`),
...@@ -175,16 +173,22 @@ ALTER TABLE `project_role_assignments` ADD CONSTRAINT `project_role_assignments_ ...@@ -175,16 +173,22 @@ ALTER TABLE `project_role_assignments` ADD CONSTRAINT `project_role_assignments_
--> statement-breakpoint --> statement-breakpoint
CREATE INDEX `idx_products_status` ON `product_list` (`status`); CREATE INDEX `idx_products_status` ON `product_list` (`status`);
--> statement-breakpoint --> statement-breakpoint
CREATE INDEX `idx_products_updated_at` ON `product_list` (`updated_at`);
--> statement-breakpoint
CREATE INDEX `idx_products_operator_user_id` ON `product_list` (`operator_user_id`); CREATE INDEX `idx_products_operator_user_id` ON `product_list` (`operator_user_id`);
--> statement-breakpoint --> statement-breakpoint
CREATE INDEX `idx_projects_stage_status` ON `project_list` (`stage`,`status`); CREATE INDEX `idx_projects_stage_status` ON `project_list` (`stage`,`status`);
--> statement-breakpoint --> statement-breakpoint
CREATE INDEX `idx_projects_created_at` ON `project_list` (`created_at`);
--> statement-breakpoint
CREATE INDEX `idx_projects_product_id` ON `project_list` (`product_id`); CREATE INDEX `idx_projects_product_id` ON `project_list` (`product_id`);
--> statement-breakpoint --> statement-breakpoint
CREATE INDEX `idx_projects_operator_user_id` ON `project_list` (`operator_user_id`); CREATE INDEX `idx_projects_operator_user_id` ON `project_list` (`operator_user_id`);
--> statement-breakpoint --> statement-breakpoint
CREATE INDEX `idx_cases_product_id` ON `case_list` (`product_id`); CREATE INDEX `idx_cases_product_id` ON `case_list` (`product_id`);
--> statement-breakpoint --> statement-breakpoint
CREATE INDEX `idx_cases_updated_at` ON `case_list` (`updated_at`);
--> statement-breakpoint
CREATE INDEX `idx_cases_operator_user_id` ON `case_list` (`operator_user_id`); CREATE INDEX `idx_cases_operator_user_id` ON `case_list` (`operator_user_id`);
--> statement-breakpoint --> statement-breakpoint
CREATE INDEX `idx_project_timelines_project_id` ON `project_timelines` (`project_id`); CREATE INDEX `idx_project_timelines_project_id` ON `project_timelines` (`project_id`);
......
-- 仅用于已经执行过旧版 002_init_ai_gateway_schema.sql 的开发库。
-- 新库直接执行最新版 002,不需要执行本文件。
ALTER TABLE `ai_model_mappings`
CHANGE COLUMN `capability` `type` varchar(16) NOT NULL COMMENT '模型类型: text/image/video',
ADD COLUMN `name` varchar(120) NULL COMMENT '模型展示名称' AFTER `alias`;
UPDATE `ai_model_mappings`
SET `name` = `alias`
WHERE `name` IS NULL OR `name` = '';
ALTER TABLE `ai_model_mappings`
MODIFY COLUMN `name` varchar(120) NOT NULL COMMENT '模型展示名称';
ALTER TABLE `ai_usage`
DROP INDEX `idx_ai_usage_capability`,
CHANGE COLUMN `capability` `type` varchar(16) NOT NULL COMMENT '模型类型: text/image/video',
ADD INDEX `idx_ai_usage_type` (`type`, `created_at`);
ALTER TABLE `ai_apps`
ADD COLUMN `billing_mode` varchar(16) NOT NULL DEFAULT 'internal' COMMENT '结算模式: internal内部使用/quota额度控制' AFTER `name`;
ALTER TABLE `ai_api_keys`
DROP COLUMN `quota_day`,
DROP COLUMN `quota_month`;
CREATE TABLE `ai_quota_packages` (
`id` bigint unsigned NOT NULL AUTO_INCREMENT COMMENT '主键',
`app_id` bigint unsigned NOT NULL COMMENT '所属应用',
`name` varchar(120) NOT NULL COMMENT '额度名称',
`amount_micro` bigint unsigned NOT NULL COMMENT '总额度(微元)',
`used_micro` bigint unsigned NOT NULL DEFAULT 0 COMMENT '已使用额度(微元)',
`start_date` date NOT NULL COMMENT '生效日期',
`end_date` date NOT NULL COMMENT '失效日期',
`status` int NOT NULL DEFAULT 1 COMMENT '状态: 0停用 1启用',
`remark` text DEFAULT NULL COMMENT '备注',
`operator_user_id` varchar(64) DEFAULT NULL COMMENT '最后操作人用户ID',
`operator_name` varchar(120) DEFAULT NULL COMMENT '最后操作人姓名',
`created_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP COMMENT '创建时间',
`updated_at` datetime NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP COMMENT '更新时间',
PRIMARY KEY (`id`),
KEY `idx_ai_quota_packages_app` (`app_id`, `status`, `start_date`, `end_date`)
) ENGINE=InnoDB DEFAULT CHARSET=utf8mb4 COLLATE=utf8mb4_unicode_ci COMMENT='AI 网关-应用额度包';
ALTER TABLE `ai_usage`
ADD COLUMN `quota_package_id` bigint unsigned DEFAULT NULL COMMENT '本次结算使用的额度包' AFTER `api_key_id`,
ADD INDEX `idx_ai_usage_quota_package` (`quota_package_id`, `created_at`);
DROP TABLE `ai_quota_usage`;
-- 已有旧版 DMS 表的增量修复(MySQL 8+),仅执行一次。
-- 新建库执行最新版 001_init_dms_schema.sql 即可,无需执行本文件。
-- 替代原 004 / 005:若已执行过它们,不要再次执行本文件。
-- 删除 product_name 列及其数据;产品名称统一关联 product_list 获取。
ALTER TABLE `project_list`
DROP COLUMN `product_name`,
ADD INDEX `idx_projects_created_at` (`created_at`);
ALTER TABLE `case_list`
DROP COLUMN `product_name`,
ADD INDEX `idx_cases_updated_at` (`updated_at`);
ALTER TABLE `product_list`
ADD INDEX `idx_products_updated_at` (`updated_at`);
-- 已有 DMS 数据库增加各阶段附件字段(MySQL 8+),仅执行一次。
-- 新建数据库执行最新版 001_init_dms_schema.sql 即可,无需执行本文件。
ALTER TABLE `project_list`
ADD COLUMN `attachment_file_url` text COMMENT '阶段附件JSON数组' AFTER `solution`;
ALTER TABLE `project_initiations`
ADD COLUMN `attachment_file_url` text COMMENT '阶段附件JSON数组' AFTER `project_id`;
ALTER TABLE `project_procurements`
ADD COLUMN `attachment_file_url` text COMMENT '阶段附件JSON数组' AFTER `formal_bid_status`;
ALTER TABLE `project_contracts`
ADD COLUMN `attachment_file_url` text COMMENT '阶段附件JSON数组' AFTER `drafter`;
ALTER TABLE `project_deliveries`
ADD COLUMN `attachment_file_url` text COMMENT '阶段附件JSON数组' AFTER `delivery_note`;
ALTER TABLE `project_acceptances`
ADD COLUMN `attachment_file_url` text COMMENT '阶段附件JSON数组' AFTER `acceptance_note`;
...@@ -14,6 +14,7 @@ ...@@ -14,6 +14,7 @@
"@fastify/cors": "^11.3.0", "@fastify/cors": "^11.3.0",
"@fastify/formbody": "^9.0.0", "@fastify/formbody": "^9.0.0",
"@fastify/http-proxy": "^11.6.2", "@fastify/http-proxy": "^11.6.2",
"@fastify/rate-limit": "^11.2.0",
"@fastify/swagger": "^9.8.1", "@fastify/swagger": "^9.8.1",
"@fastify/swagger-ui": "^6.1.1", "@fastify/swagger-ui": "^6.1.1",
"ali-oss": "^6.23.0", "ali-oss": "^6.23.0",
...@@ -1330,6 +1331,28 @@ ...@@ -1330,6 +1331,28 @@
"ipaddr.js": "^2.1.0" "ipaddr.js": "^2.1.0"
} }
}, },
"node_modules/@fastify/rate-limit": {
"version": "11.2.0",
"resolved": "https://registry.npmjs.org/@fastify/rate-limit/-/rate-limit-11.2.0.tgz",
"integrity": "sha512-X7osJd4XSvMoejYrnJkSZYYjY1eNYoBqhjlzf1RakC2204qExFqZFTKj5+T7VuzA/iUI9Z3UoSqQRkB2HpG0oQ==",
"funding": [
{
"type": "github",
"url": "https://github.com/sponsors/fastify"
},
{
"type": "opencollective",
"url": "https://opencollective.com/fastify"
}
],
"license": "MIT",
"dependencies": {
"@lukeed/ms": "^2.0.2",
"fastify-plugin": "^6.0.0",
"ip-address": "^10.2.0",
"toad-cache": "^3.7.0"
}
},
"node_modules/@fastify/reply-from": { "node_modules/@fastify/reply-from": {
"version": "12.6.5", "version": "12.6.5",
"resolved": "https://registry.npmjs.org/@fastify/reply-from/-/reply-from-12.6.5.tgz", "resolved": "https://registry.npmjs.org/@fastify/reply-from/-/reply-from-12.6.5.tgz",
...@@ -3327,6 +3350,15 @@ ...@@ -3327,6 +3350,15 @@
"integrity": "sha512-k/vGaX4/Yla3WzyMCvTQOXYeIHvqOKtnqBduzTHpzpQZzAskKMhZ2K+EnBiSM9zGSoIFeMpXKxa4dYeZIQqewQ==", "integrity": "sha512-k/vGaX4/Yla3WzyMCvTQOXYeIHvqOKtnqBduzTHpzpQZzAskKMhZ2K+EnBiSM9zGSoIFeMpXKxa4dYeZIQqewQ==",
"license": "ISC" "license": "ISC"
}, },
"node_modules/ip-address": {
"version": "10.7.0",
"resolved": "https://registry.npmjs.org/ip-address/-/ip-address-10.7.0.tgz",
"integrity": "sha512-BGFsyJd5mpXp3rK6jIdADLNgpJUK1jnjzvYF8lK+VyDab9JAmqN0YOKDdP17HlgKb2+ehPgDc8EtnRLbGCAMhA==",
"license": "MIT",
"engines": {
"node": ">= 12"
}
},
"node_modules/ipaddr.js": { "node_modules/ipaddr.js": {
"version": "2.5.0", "version": "2.5.0",
"resolved": "https://registry.npmjs.org/ipaddr.js/-/ipaddr.js-2.5.0.tgz", "resolved": "https://registry.npmjs.org/ipaddr.js/-/ipaddr.js-2.5.0.tgz",
......
...@@ -42,6 +42,7 @@ ...@@ -42,6 +42,7 @@
"@fastify/cors": "^11.3.0", "@fastify/cors": "^11.3.0",
"@fastify/formbody": "^9.0.0", "@fastify/formbody": "^9.0.0",
"@fastify/http-proxy": "^11.6.2", "@fastify/http-proxy": "^11.6.2",
"@fastify/rate-limit": "^11.2.0",
"@fastify/swagger": "^9.8.1", "@fastify/swagger": "^9.8.1",
"@fastify/swagger-ui": "^6.1.1", "@fastify/swagger-ui": "^6.1.1",
"ali-oss": "^6.23.0", "ali-oss": "^6.23.0",
......
...@@ -139,13 +139,12 @@ async function main() { ...@@ -139,13 +139,12 @@ async function main() {
for (const [index, item] of cases.entries()) { for (const [index, item] of cases.entries()) {
const [result] = await connection.execute( const [result] = await connection.execute(
`INSERT INTO case_list `INSERT INTO case_list
(name, description, product_id, product_name, files, operator_user_id, operator_name, created_at, updated_at) (name, description, product_id, files, operator_user_id, operator_name, created_at, updated_at)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)`, VALUES (?, ?, ?, ?, ?, ?, ?, ?)`,
[ [
item.name, item.name,
item.description, item.description,
item.productId, item.productId,
item.productName,
item.files, item.files,
item.operatorUserId, item.operatorUserId,
item.operatorName, item.operatorName,
......
...@@ -54,7 +54,9 @@ export default async function app(fastify, opts) { ...@@ -54,7 +54,9 @@ export default async function app(fastify, opts) {
await fastify.register(autoload, { await fastify.register(autoload, {
dir: path.join(__dirname, 'routes'), dir: path.join(__dirname, 'routes'),
autoHooks: true, autoHooks: true,
cascadeHooks: true, // 不向下级联:routes/ai/autohooks.js 是数据面(API Key 鉴权),
// 级联会把它套到 routes/ai/admin/ 的管理面(TGC 鉴权)上。
cascadeHooks: false,
options: { ...opts }, options: { ...opts },
}) })
} }
import crypto from 'node:crypto' import crypto from 'node:crypto'
import axios from 'axios' import axios from 'axios'
import config from '#src/config.js' import config from '#src/config.js'
import logger from '#src/lib/logger.js'
const MAX_BATCH_SIZE = 500 const MAX_BATCH_SIZE = 500
......
...@@ -38,6 +38,13 @@ const envSchema = z.object({ ...@@ -38,6 +38,13 @@ const envSchema = z.object({
PERMISSION_APP_SECRET: z.string().default(''), PERMISSION_APP_SECRET: z.string().default(''),
SSO_USER_CACHE_TTL_SECONDS: int(180), SSO_USER_CACHE_TTL_SECONDS: int(180),
PERMISSION_CACHE_TTL_SECONDS: int(180), PERMISSION_CACHE_TTL_SECONDS: int(180),
// AI 网关上游
VOLCANO_BASE_URL: z.string().default('https://ark.cn-beijing.volces.com/api/v3'),
VOLCANO_API_KEY: z.string().default(''),
DEEPSEEK_BASE_URL: z.string().default('https://api.deepseek.com/v1'),
DEEPSEEK_API_KEY: z.string().default(''),
AI_RATE_LIMIT_PER_MINUTE: int(120),
}) })
// 启动即校验环境变量,配置错误立刻失败,而不是运行到某个请求时才炸 // 启动即校验环境变量,配置错误立刻失败,而不是运行到某个请求时才炸
...@@ -91,6 +98,14 @@ const config = { ...@@ -91,6 +98,14 @@ const config = {
userInfoCacheTtlSeconds: env.SSO_USER_CACHE_TTL_SECONDS, userInfoCacheTtlSeconds: env.SSO_USER_CACHE_TTL_SECONDS,
permissionCacheTtlSeconds: env.PERMISSION_CACHE_TTL_SECONDS, permissionCacheTtlSeconds: env.PERMISSION_CACHE_TTL_SECONDS,
}, },
ai: {
rateLimitPerMinute: env.AI_RATE_LIMIT_PER_MINUTE,
// 上游密钥只走环境变量,不进数据库
providers: {
volcano: { baseUrl: env.VOLCANO_BASE_URL, apiKey: env.VOLCANO_API_KEY },
deepseek: { baseUrl: env.DEEPSEEK_BASE_URL, apiKey: env.DEEPSEEK_API_KEY },
},
},
wechat: { wechat: {
apps: { apps: {
wxd6109d07f6396e5c: 'd80a330735fc82f3fd6aba425481e8fd', wxd6109d07f6396e5c: 'd80a330735fc82f3fd6aba425481e8fd',
......
import { bigint, decimal, index, int, mysqlTable, varchar } from 'drizzle-orm/mysql-core'
import { timestamps } from '../columns.js'
export const USAGE_STATUSES = {
PENDING: 'pending',
SUCCESS: 'success',
ERROR: 'error',
ABORTED: 'aborted',
}
export const USAGE_SOURCES = {
REPORTED: 'reported',
MISSING: 'missing',
}
/** 「上游成功但没返回用量」——此时用量来源为 missing,成本尚未确认,用它标记出来便于对账 */
export const USAGE_ERROR_CODES = {
USAGE_MISSING: 'usage_missing',
}
// 计费单位:与 ai_usage.unit 一一对应
export const USAGE_UNITS = {
TOKEN: 'token',
IMAGE: 'image',
SECOND: 'second',
}
// 调用明细:用量统计的核心
export const aiUsage = mysqlTable(
'ai_usage',
{
id: bigint('id', { mode: 'number', unsigned: true }).autoincrement().primaryKey(),
request_id: varchar('request_id', { length: 64 }),
app_id: bigint('app_id', { mode: 'number', unsigned: true }),
api_key_id: bigint('api_key_id', { mode: 'number', unsigned: true }),
quota_package_id: bigint('quota_package_id', { mode: 'number', unsigned: true }),
user_id: varchar('user_id', { length: 64 }),
user_name: varchar('user_name', { length: 120 }),
type: varchar('type', { length: 16 }).notNull(),
provider: varchar('provider', { length: 32 }).notNull(),
model: varchar('model', { length: 80 }).notNull(),
upstream_model: varchar('upstream_model', { length: 120 }).notNull(),
// 计价用的价格版本(对账/重算用)
price_version: int('price_version'),
unit: varchar('unit', { length: 16 }).notNull(),
quantity: bigint('quantity', { mode: 'number', unsigned: true }).notNull().default(0),
prompt_tokens: bigint('prompt_tokens', { mode: 'number', unsigned: true }).notNull().default(0),
// 缓存命中的输入 token(单价更低,影响成本计算)
cached_tokens: bigint('cached_tokens', { mode: 'number', unsigned: true }).notNull().default(0),
completion_tokens: bigint('completion_tokens', { mode: 'number', unsigned: true }).notNull().default(0),
// 思考 token(已包含在 completion_tokens 内)
reasoning_tokens: bigint('reasoning_tokens', { mode: 'number', unsigned: true }).notNull().default(0),
usage_source: varchar('usage_source', { length: 16 }).notNull().default(USAGE_SOURCES.REPORTED),
cost: decimal('cost', { precision: 12, scale: 6 }).notNull().default('0'),
status: varchar('status', { length: 16 }).notNull().default(USAGE_STATUSES.PENDING),
upstream_task_id: varchar('upstream_task_id', { length: 128 }),
http_status: int('http_status'),
error_code: varchar('error_code', { length: 64 }),
latency_ms: int('latency_ms'),
first_token_ms: int('first_token_ms'),
stream: int('stream').notNull().default(0),
...timestamps,
},
(table) => ({
appIdx: index('idx_ai_usage_app').on(table.app_id, table.created_at),
keyIdx: index('idx_ai_usage_key').on(table.api_key_id, table.created_at),
quotaPackageIdx: index('idx_ai_usage_quota_package').on(table.quota_package_id, table.created_at),
modelIdx: index('idx_ai_usage_model').on(table.provider, table.model, table.created_at),
typeIdx: index('idx_ai_usage_type').on(table.type, table.created_at),
statusIdx: index('idx_ai_usage_status').on(table.status),
taskIdx: index('idx_ai_usage_task').on(table.app_id, table.api_key_id, table.upstream_task_id),
// 惰性兜底:扫描某 key 过期的 pending 视频任务
pendingIdx: index('idx_ai_usage_pending').on(table.api_key_id, table.status, table.unit, table.created_at),
})
)
import { bigint, char, datetime, index, int, mysqlTable, uniqueIndex, varchar } from 'drizzle-orm/mysql-core'
import { timestamps } from '../columns.js'
export const KEY_STATUSES = {
DISABLED: 0,
ENABLED: 1,
}
// 应用持有的密钥:只负责鉴权,用量控制属于应用额度
export const aiApiKeys = mysqlTable(
'ai_api_keys',
{
id: bigint('id', { mode: 'number', unsigned: true }).autoincrement().primaryKey(),
app_id: bigint('app_id', { mode: 'number', unsigned: true }).notNull(),
name: varchar('name', { length: 120 }).notNull(),
key_hash: char('key_hash', { length: 64 }).notNull(),
key_prefix: varchar('key_prefix', { length: 12 }).notNull(),
status: int('status').notNull().default(KEY_STATUSES.ENABLED),
last_used_at: datetime('last_used_at', { mode: 'string' }),
operator_user_id: varchar('operator_user_id', { length: 64 }),
operator_name: varchar('operator_name', { length: 120 }),
...timestamps,
},
(table) => ({
hashIdx: uniqueIndex('uk_ai_api_keys_hash').on(table.key_hash),
appIdx: index('idx_ai_api_keys_app').on(table.app_id),
})
)
import { bigint, index, int, mysqlTable, uniqueIndex, varchar } from 'drizzle-orm/mysql-core'
import { timestamps } from '../columns.js'
export const APP_STATUSES = {
DISABLED: 0,
ENABLED: 1,
}
export const BILLING_MODES = {
INTERNAL: 'internal',
QUOTA: 'quota',
}
// 接入应用:密钥和用量的业务归属
export const aiApps = mysqlTable(
'ai_apps',
{
id: bigint('id', { mode: 'number', unsigned: true }).autoincrement().primaryKey(),
code: varchar('code', { length: 64 }).notNull(),
name: varchar('name', { length: 120 }).notNull(),
billing_mode: varchar('billing_mode', { length: 16 }).notNull().default(BILLING_MODES.INTERNAL),
status: int('status').notNull().default(APP_STATUSES.ENABLED),
operator_user_id: varchar('operator_user_id', { length: 64 }),
operator_name: varchar('operator_name', { length: 120 }),
...timestamps,
},
(table) => ({
codeIdx: uniqueIndex('uk_ai_apps_code').on(table.code),
statusIdx: index('idx_ai_apps_status').on(table.status),
})
)
import { index, int, mysqlTable, varchar } from 'drizzle-orm/mysql-core'
import { timestamps } from '../columns.js'
export const AI_MODEL_TYPES = {
TEXT: 'text',
IMAGE: 'image',
VIDEO: 'video',
}
// 平台对外暴露的模型别名(上游路由见 ai_model_routes)
export const aiModelMappings = mysqlTable(
'ai_model_mappings',
{
alias: varchar('alias', { length: 80 }).primaryKey(),
name: varchar('name', { length: 120 }).notNull(),
type: varchar('type', { length: 16 }).notNull(),
provider: varchar('provider', { length: 32 }).notNull(),
upstream_model: varchar('upstream_model', { length: 120 }).notNull(),
enabled: int('enabled').notNull().default(1),
operator_user_id: varchar('operator_user_id', { length: 64 }),
operator_name: varchar('operator_name', { length: 120 }),
...timestamps,
},
(table) => ({
providerIdx: index('idx_ai_model_mappings_provider').on(table.provider),
})
)
import { sql } from 'drizzle-orm'
import { datetime, decimal, index, int, mysqlTable, uniqueIndex, varchar } from 'drizzle-orm/mysql-core'
import { timestamps } from '../columns.js'
export const PRICING_UNITS = {
PER_1M_TOKENS: 'per_1m_tokens',
PER_IMAGE: 'per_image',
PER_SECOND: 'per_second',
}
export const PRICE_STATUSES = {
DISABLED: 0,
ACTIVE: 1,
}
// 单价表:每个别名一个价(一个版本一行);改价 = 关闭旧版本 + 开新版本
// 暂不做分辨率/清晰度档位(2026-09-10 决定:先固定一个价格)
export const aiModelPrices = mysqlTable(
'ai_model_prices',
{
// 定价挂在别名上(对客卖的是别名;上游成本在路由层)
alias: varchar('alias', { length: 80 }).notNull(),
// 版本号从 1 递增;改价 = 新增版本 + 关闭旧版本
version: int('version').notNull().default(1),
pricing_unit: varchar('pricing_unit', { length: 24 }).notNull(),
input_price: decimal('input_price', { precision: 12, scale: 6 }).notNull().default('0'),
// 缓存命中单价;为 0 时回退到 input_price
cached_input_price: decimal('cached_input_price', { precision: 12, scale: 6 }).notNull().default('0'),
output_price: decimal('output_price', { precision: 12, scale: 6 }).notNull().default('0'),
unit_price: decimal('unit_price', { precision: 12, scale: 6 }).notNull().default('0'),
currency: varchar('currency', { length: 3 }).notNull().default('CNY'),
effective_from: datetime('effective_from', { mode: 'string' })
.notNull()
.default(sql`CURRENT_TIMESTAMP`),
// NULL 表示仍然生效
effective_to: datetime('effective_to', { mode: 'string' }),
status: int('status').notNull().default(1),
operator_user_id: varchar('operator_user_id', { length: 64 }),
operator_name: varchar('operator_name', { length: 120 }),
...timestamps,
},
(table) => ({
versionIdx: uniqueIndex('uk_ai_model_prices_alias_version').on(table.alias, table.version),
activeIdx: index('idx_ai_model_prices_alias').on(table.alias, table.status),
})
)
import { bigint, date, index, int, mysqlTable, text, varchar } from 'drizzle-orm/mysql-core'
import { timestamps } from '../columns.js'
export const QUOTA_PACKAGE_STATUSES = {
DISABLED: 0,
ENABLED: 1,
}
export const aiQuotaPackages = mysqlTable(
'ai_quota_packages',
{
id: bigint('id', { mode: 'number', unsigned: true }).autoincrement().primaryKey(),
app_id: bigint('app_id', { mode: 'number', unsigned: true }).notNull(),
name: varchar('name', { length: 120 }).notNull(),
amount_micro: bigint('amount_micro', { mode: 'number', unsigned: true }).notNull(),
used_micro: bigint('used_micro', { mode: 'number', unsigned: true }).notNull().default(0),
start_date: date('start_date', { mode: 'string' }).notNull(),
end_date: date('end_date', { mode: 'string' }).notNull(),
status: int('status').notNull().default(QUOTA_PACKAGE_STATUSES.ENABLED),
remark: text('remark'),
operator_user_id: varchar('operator_user_id', { length: 64 }),
operator_name: varchar('operator_name', { length: 120 }),
...timestamps,
},
(table) => ({
appIdx: index('idx_ai_quota_packages_app').on(table.app_id, table.status, table.start_date, table.end_date),
}),
)
...@@ -8,7 +8,6 @@ export const cases = mysqlTable( ...@@ -8,7 +8,6 @@ export const cases = mysqlTable(
name: varchar('name', { length: 255 }).notNull(), name: varchar('name', { length: 255 }).notNull(),
description: text('description'), description: text('description'),
product_id: bigint('product_id', { mode: 'number', unsigned: true }), product_id: bigint('product_id', { mode: 'number', unsigned: true }),
product_name: varchar('product_name', { length: 255 }),
files: longtext('files'), files: longtext('files'),
operator_user_id: varchar('operator_user_id', { length: 64 }), operator_user_id: varchar('operator_user_id', { length: 64 }),
operator_name: varchar('operator_name', { length: 120 }), operator_name: varchar('operator_name', { length: 120 }),
...@@ -16,6 +15,8 @@ export const cases = mysqlTable( ...@@ -16,6 +15,8 @@ export const cases = mysqlTable(
}, },
(table) => ({ (table) => ({
productIdIdx: index('idx_cases_product_id').on(table.product_id), productIdIdx: index('idx_cases_product_id').on(table.product_id),
// 列表按 updated_at 倒序分页
updatedAtIdx: index('idx_cases_updated_at').on(table.updated_at),
operatorUserIdIdx: index('idx_cases_operator_user_id').on(table.operator_user_id), operatorUserIdIdx: index('idx_cases_operator_user_id').on(table.operator_user_id),
}) })
) )
...@@ -20,6 +20,8 @@ export const products = mysqlTable( ...@@ -20,6 +20,8 @@ export const products = mysqlTable(
(table) => ({ (table) => ({
nameIdx: uniqueIndex('uk_products_name').on(table.name), nameIdx: uniqueIndex('uk_products_name').on(table.name),
statusIdx: index('idx_products_status').on(table.status), statusIdx: index('idx_products_status').on(table.status),
// 列表按 updated_at 倒序分页
updatedAtIdx: index('idx_products_updated_at').on(table.updated_at),
operatorUserIdIdx: index('idx_products_operator_user_id').on(table.operator_user_id), operatorUserIdIdx: index('idx_products_operator_user_id').on(table.operator_user_id),
}) })
) )
...@@ -37,11 +37,11 @@ export const projects = mysqlTable( ...@@ -37,11 +37,11 @@ export const projects = mysqlTable(
school_name: varchar('school_name', { length: 255 }), school_name: varchar('school_name', { length: 255 }),
department_name: varchar('department_name', { length: 255 }), department_name: varchar('department_name', { length: 255 }),
product_id: bigint('product_id', { mode: 'number', unsigned: true }), product_id: bigint('product_id', { mode: 'number', unsigned: true }),
product_name: varchar('product_name', { length: 255 }),
contact_name: varchar('contact_name', { length: 120 }), contact_name: varchar('contact_name', { length: 120 }),
contact_title: varchar('contact_title', { length: 120 }), contact_title: varchar('contact_title', { length: 120 }),
contact_phone: varchar('contact_phone', { length: 64 }), contact_phone: varchar('contact_phone', { length: 64 }),
solution: text('solution'), solution: text('solution'),
attachment_file_url: text('attachment_file_url'),
stage: int('stage').notNull().default(PROJECT_STAGES.SOLUTION), stage: int('stage').notNull().default(PROJECT_STAGES.SOLUTION),
status: int('status').notNull().default(PROJECT_STATUSES.ACTIVE), status: int('status').notNull().default(PROJECT_STATUSES.ACTIVE),
description: text('description'), description: text('description'),
...@@ -52,6 +52,8 @@ export const projects = mysqlTable( ...@@ -52,6 +52,8 @@ export const projects = mysqlTable(
(table) => ({ (table) => ({
projectCodeIdx: uniqueIndex('uk_projects_project_code').on(table.project_code), projectCodeIdx: uniqueIndex('uk_projects_project_code').on(table.project_code),
stageStatusIdx: index('idx_projects_stage_status').on(table.stage, table.status), stageStatusIdx: index('idx_projects_stage_status').on(table.stage, table.status),
// 列表按 created_at 倒序分页
createdAtIdx: index('idx_projects_created_at').on(table.created_at),
productIdIdx: index('idx_projects_product_id').on(table.product_id), productIdIdx: index('idx_projects_product_id').on(table.product_id),
operatorUserIdIdx: index('idx_projects_operator_user_id').on(table.operator_user_id), operatorUserIdIdx: index('idx_projects_operator_user_id').on(table.operator_user_id),
}) })
...@@ -64,8 +66,7 @@ export const initiations = mysqlTable( ...@@ -64,8 +66,7 @@ export const initiations = mysqlTable(
project_id: bigint('project_id', { mode: 'number', unsigned: true }) project_id: bigint('project_id', { mode: 'number', unsigned: true })
.notNull() .notNull()
.references(() => projects.id, { onDelete: 'cascade' }), .references(() => projects.id, { onDelete: 'cascade' }),
application_file_url: text('application_file_url'), attachment_file_url: text('attachment_file_url'),
argument_file_url: text('argument_file_url'),
project_amount: decimal('project_amount', { precision: 14, scale: 2 }), project_amount: decimal('project_amount', { precision: 14, scale: 2 }),
fund_source: varchar('fund_source', { length: 255 }), fund_source: varchar('fund_source', { length: 255 }),
execution_plan: text('execution_plan'), execution_plan: text('execution_plan'),
...@@ -91,8 +92,7 @@ export const procurements = mysqlTable( ...@@ -91,8 +92,7 @@ export const procurements = mysqlTable(
main_bid_owner: varchar('main_bid_owner', { length: 120 }), main_bid_owner: varchar('main_bid_owner', { length: 120 }),
companion_bidders: text('companion_bidders'), companion_bidders: text('companion_bidders'),
formal_bid_status: varchar('formal_bid_status', { length: 120 }), formal_bid_status: varchar('formal_bid_status', { length: 120 }),
winning_notice_file_url: text('winning_notice_file_url'), attachment_file_url: text('attachment_file_url'),
bid_archive_file_url: text('bid_archive_file_url'),
...timestamps, ...timestamps,
}, },
(table) => ({ (table) => ({
...@@ -111,7 +111,7 @@ export const contracts = mysqlTable( ...@@ -111,7 +111,7 @@ export const contracts = mysqlTable(
contract_name: varchar('contract_name', { length: 255 }), contract_name: varchar('contract_name', { length: 255 }),
amount: decimal('amount', { precision: 14, scale: 2 }), amount: decimal('amount', { precision: 14, scale: 2 }),
drafter: varchar('drafter', { length: 120 }), drafter: varchar('drafter', { length: 120 }),
archive_file_url: text('archive_file_url'), attachment_file_url: text('attachment_file_url'),
...timestamps, ...timestamps,
}, },
(table) => ({ (table) => ({
...@@ -128,6 +128,7 @@ export const deliveries = mysqlTable('project_deliveries', { ...@@ -128,6 +128,7 @@ export const deliveries = mysqlTable('project_deliveries', {
delivery_contact: varchar('delivery_contact', { length: 120 }), delivery_contact: varchar('delivery_contact', { length: 120 }),
delivery_contact_phone: varchar('delivery_contact_phone', { length: 64 }), delivery_contact_phone: varchar('delivery_contact_phone', { length: 64 }),
delivery_note: text('delivery_note'), delivery_note: text('delivery_note'),
attachment_file_url: text('attachment_file_url'),
completed_at: datetime('completed_at', { mode: 'string' }), completed_at: datetime('completed_at', { mode: 'string' }),
...timestamps, ...timestamps,
}, (table) => ({ }, (table) => ({
...@@ -139,8 +140,8 @@ export const acceptances = mysqlTable('project_acceptances', { ...@@ -139,8 +140,8 @@ export const acceptances = mysqlTable('project_acceptances', {
project_id: bigint('project_id', { mode: 'number', unsigned: true }) project_id: bigint('project_id', { mode: 'number', unsigned: true })
.notNull() .notNull()
.references(() => projects.id, { onDelete: 'cascade' }), .references(() => projects.id, { onDelete: 'cascade' }),
acceptance_report_url: text('acceptance_report_url'),
acceptance_note: text('acceptance_note'), acceptance_note: text('acceptance_note'),
attachment_file_url: text('attachment_file_url'),
...timestamps, ...timestamps,
}, (table) => ({ }, (table) => ({
projectIdIdx: uniqueIndex('uk_acceptances_project_id').on(table.project_id), projectIdIdx: uniqueIndex('uk_acceptances_project_id').on(table.project_id),
......
...@@ -25,11 +25,20 @@ const start = async () => { ...@@ -25,11 +25,20 @@ const start = async () => {
process.exit(1) process.exit(1)
} }
// Graceful shutdown // Graceful shutdown:给 close 设上限
// SSE 流式请求最长可挂 120s,没有上限的话进程会一直不退出(pm2 reload 时请求灰掉)
const SHUTDOWN_TIMEOUT_MS = 10_000
const shutdown = async (signal) => { const shutdown = async (signal) => {
logger.info(`${signal} received, shutting down...`) logger.info(`${signal} received, shutting down...`)
await fastify.close() try {
await Promise.race([
fastify.close(),
new Promise((resolve) => setTimeout(resolve, SHUTDOWN_TIMEOUT_MS).unref()),
])
logger.info('Server closed') logger.info('Server closed')
} catch (error) {
logger.error({ err: error }, '关闭时出错,强制退出')
}
process.exit(0) process.exit(0)
} }
process.on('SIGTERM', () => shutdown('SIGTERM')) process.on('SIGTERM', () => shutdown('SIGTERM'))
......
/**
* 带错误码的业务错误。
*
* 与 lib/http-error.js 的区别:httpError 用于「网关自己拒绝请求」(401/429/503 等),
* 这里的 code 用于标注「上游/流式链路的技术性失败原因」,会被写进 ai_usage.error_code,
* 便于报表按错误类型聚合。两者都会透传 statusCode 给全局 error handler。
*/
export const codedError = (statusCode, message, code) => {
const error = new Error(message)
error.statusCode = statusCode
error.code = code
return error
}
/** 请求已被中止(客户端断开、超时、上游不可达时 fetch 会抛 AbortError/TimeoutError) */
export const isAborted = (error) =>
error?.name === 'AbortError' || error?.name === 'TimeoutError'
...@@ -93,8 +93,9 @@ const logSchema = new mongoose.Schema( ...@@ -93,8 +93,9 @@ const logSchema = new mongoose.Schema(
userAgent: String, userAgent: String,
ip: String, ip: String,
url: String, url: String,
createdAt: String, // immutable:客户端 body 里塞 createdAt 也改不动(写入时由 timestamps 填充)
updatedAt: String, createdAt: { type: String, immutable: true },
updatedAt: { type: String, immutable: true },
}, },
{ {
timestamps: { timestamps: {
......
import fp from 'fastify-plugin'
import rateLimit from '@fastify/rate-limit'
/**
* 限流插件:默认不全局生效(global: false),由各路由按需声明。
* AI 数据面按 API key 限流——鉴权钩子在 onRequest,限流放在 preHandler,
* 所以 keyGenerator 能拿到 request.ai。
*/
export default fp(async (fastify) => {
await fastify.register(rateLimit, {
global: false,
max: 300,
timeWindow: '1 minute',
})
}, { name: 'rate-limit' })
import fp from 'fastify-plugin' import fp from 'fastify-plugin'
import httpProxy from '@fastify/http-proxy' import httpProxy from '@fastify/http-proxy'
// /api/usercenter/* 透传到用户中心 // /proxy/usercenter/* 透传到用户中心(/proxy 前缀表示这类接口是转发的上游,不是我们的业务资源)
export default fp(async (fastify) => { export default fp(async (fastify) => {
await fastify.register(httpProxy, { await fastify.register(httpProxy, {
upstream: 'https://api-usercenter.ezijing.com', upstream: 'https://api-usercenter.ezijing.com',
prefix: '/api/usercenter', prefix: '/proxy/usercenter',
}) })
}, { name: 'usercenter-proxy' }) }, { name: 'usercenter-proxy' })
import { success } from '#src/lib/response.js'
import { getCurrentUser } from '#src/services/dms/auth.service.js'
import * as adminService from '#src/services/ai/admin.service.js'
import { idParam } from '#src/schemas/dms/common.js'
import {
appCreateBody,
appRowSchema,
appUpdateBody,
appsListQuery,
itemResponseSchema,
listResponseSchema,
} from '#src/schemas/ai/admin.js'
const list = async (request, reply) => {
const { keyword, status, page, limit } = request.query
const result = await adminService.listApps({ keyword, status }, { page, limit })
return success(reply, result)
}
const create = async (request, reply) => {
const app = await adminService.createApp(request.body, getCurrentUser(request))
return success(reply, app, 201)
}
const update = async (request, reply) => {
const app = await adminService.updateApp(request.params.id, request.body, getCurrentUser(request))
return success(reply, app)
}
const remove = async (request, reply) => {
const app = await adminService.deleteApp(request.params.id)
return success(reply, app)
}
export default async function aiAdminAppsRoutes(fastify) {
fastify.get(
'/apps',
{ schema: { querystring: appsListQuery, response: { 200: listResponseSchema(appRowSchema) } } },
list,
)
fastify.post(
'/apps',
{ schema: { body: appCreateBody, response: { 201: itemResponseSchema(appRowSchema) } } },
create,
)
fastify.put(
'/apps/:id',
{ schema: { params: idParam, body: appUpdateBody, response: { 200: itemResponseSchema(appRowSchema) } } },
update,
)
fastify.delete(
'/apps/:id',
{ schema: { params: idParam, response: { 200: itemResponseSchema(appRowSchema) } } },
remove,
)
}
import { authenticate, requireRouteAccess } from '#src/services/dms/hooks.js'
// 本目录所有路由统一鉴权:先验登录(onRequest,先于 schema 校验),再验路由权限
export default async function aiAdminAutohooks(fastify) {
fastify.addHook('onRequest', authenticate)
fastify.addHook('preHandler', requireRouteAccess('/dms/ai'))
}
import { success } from '#src/lib/response.js'
import { getCurrentUser } from '#src/services/dms/auth.service.js'
import * as adminService from '#src/services/ai/admin.service.js'
import { idParam } from '#src/schemas/dms/common.js'
import {
itemResponseSchema,
keyCreateBody,
keyCreatedSchema,
keyRowSchema,
keysListQuery,
keyUpdateBody,
listResponseSchema,
} from '#src/schemas/ai/admin.js'
const list = async (request, reply) => {
const { app_id, status, page, limit } = request.query
const result = await adminService.listApiKeys({ app_id, status }, { page, limit })
return success(reply, result)
}
const create = async (request, reply) => {
const key = await adminService.createApiKey(request.body, getCurrentUser(request))
return success(reply, key, 201)
}
const update = async (request, reply) => {
const key = await adminService.updateApiKey(request.params.id, request.body, getCurrentUser(request))
return success(reply, key)
}
const remove = async (request, reply) => {
const key = await adminService.deleteApiKey(request.params.id)
return success(reply, key)
}
export default async function aiAdminKeysRoutes(fastify) {
fastify.get(
'/keys',
{ schema: { querystring: keysListQuery, response: { 200: listResponseSchema(keyRowSchema) } } },
list,
)
fastify.post(
'/keys',
{ schema: { body: keyCreateBody, response: { 201: itemResponseSchema(keyCreatedSchema) } } },
create,
)
fastify.put(
'/keys/:id',
{ schema: { params: idParam, body: keyUpdateBody, response: { 200: itemResponseSchema(keyRowSchema) } } },
update,
)
fastify.delete(
'/keys/:id',
{ schema: { params: idParam, response: { 200: itemResponseSchema(keyRowSchema) } } },
remove,
)
}
import { success } from '#src/lib/response.js'
import { getCurrentUser } from '#src/services/dms/auth.service.js'
import * as adminService from '#src/services/ai/admin.service.js'
import {
itemResponseSchema,
listResponseSchema,
modelAliasParam,
modelCreateBody,
modelRowSchema,
modelsListQuery,
modelUpdateBody,
} from '#src/schemas/ai/admin.js'
const list = async (request, reply) => {
const { keyword, type, provider, page, limit } = request.query
const result = await adminService.listModels({ keyword, type, provider }, { page, limit })
return success(reply, result)
}
const create = async (request, reply) => {
const model = await adminService.createModel(request.body, getCurrentUser(request))
return success(reply, model, 201)
}
const update = async (request, reply) => {
const model = await adminService.updateModel(request.params.alias, request.body, getCurrentUser(request))
return success(reply, model)
}
const remove = async (request, reply) => {
const model = await adminService.deleteModel(request.params.alias)
return success(reply, model)
}
export default async function aiAdminModelsRoutes(fastify) {
fastify.get(
'/models',
{ schema: { querystring: modelsListQuery, response: { 200: listResponseSchema(modelRowSchema) } } },
list,
)
fastify.post(
'/models',
{ schema: { body: modelCreateBody, response: { 201: itemResponseSchema(modelRowSchema) } } },
create,
)
fastify.put(
'/models/:alias',
{ schema: { params: modelAliasParam, body: modelUpdateBody, response: { 200: itemResponseSchema(modelRowSchema) } } },
update,
)
fastify.delete(
'/models/:alias',
{ schema: { params: modelAliasParam, response: { 200: itemResponseSchema(modelRowSchema) } } },
remove,
)
}
import { success } from '#src/lib/response.js'
import { getCurrentUser } from '#src/services/dms/auth.service.js'
import * as adminService from '#src/services/ai/admin.service.js'
import { idParam } from '#src/schemas/dms/common.js'
import {
itemResponseSchema,
listResponseSchema,
quotaPackageCreateBody,
quotaPackageRowSchema,
quotaPackagesListQuery,
quotaPackageUpdateBody,
} from '#src/schemas/ai/admin.js'
export default async function aiAdminQuotaPackageRoutes(fastify) {
fastify.get('/quota-packages', {
schema: { querystring: quotaPackagesListQuery, response: { 200: listResponseSchema(quotaPackageRowSchema) } },
}, async (request, reply) => {
const { app_id, status, page, limit } = request.query
return success(reply, await adminService.listQuotaPackages({ app_id, status }, { page, limit }))
})
fastify.post('/quota-packages', {
schema: { body: quotaPackageCreateBody, response: { 201: itemResponseSchema(quotaPackageRowSchema) } },
}, async (request, reply) => {
return success(reply, await adminService.createQuotaPackage(request.body, getCurrentUser(request)), 201)
})
fastify.put('/quota-packages/:id', {
schema: { params: idParam, body: quotaPackageUpdateBody, response: { 200: itemResponseSchema(quotaPackageRowSchema) } },
}, async (request, reply) => {
return success(reply, await adminService.updateQuotaPackage(request.params.id, request.body, getCurrentUser(request)))
})
fastify.delete('/quota-packages/:id', {
schema: { params: idParam, response: { 200: itemResponseSchema(quotaPackageRowSchema) } },
}, async (request, reply) => {
return success(reply, await adminService.deleteQuotaPackage(request.params.id))
})
}
import { success } from '#src/lib/response.js'
import { usageRecentQuery, usageSummaryQuery } from '#src/schemas/ai/usage.js'
import { summarize, listRecent } from '#src/services/ai/report.service.js'
export default async function aiAdminUsageRoutes(fastify) {
fastify.get('/usage/summary', { schema: { querystring: usageSummaryQuery } }, async (request, reply) => {
const { group_by: groupBy, from, to, app_id: appId, api_key_id: apiKeyId, type, usage_source: usageSource, status, model } = request.query
const result = await summarize({ groupBy, from, to, appId, apiKeyId, type, usageSource, status, model })
return success(reply, result)
})
fastify.get('/usage/recent', { schema: { querystring: usageRecentQuery } }, async (request, reply) => {
const { page, limit, app_id: appId, api_key_id: apiKeyId, type, usage_source: usageSource, billing_status: billingStatus, status, model, request_id: requestId, from, to } = request.query
const result = await listRecent({ appId, apiKeyId, type, usageSource, billingStatus, status, model, requestId, from, to }, { page, limit })
return success(reply, result)
})
}
import { createHash } from 'node:crypto'
import { hasZodFastifySchemaValidationErrors } from 'fastify-type-provider-zod'
import { eq, sql } from 'drizzle-orm'
import { LRUCache } from 'lru-cache'
import { db } from '#src/db/client.js'
import { aiApiKeys, KEY_STATUSES } from '#src/db/schema/ai/api-keys.js'
import { aiApps, APP_STATUSES } from '#src/db/schema/ai/apps.js'
import { httpError } from '#src/lib/http-error.js'
import logger from '#src/lib/logger.js'
import { getPriceByVersion } from '#src/services/ai/pricing.service.js'
import { settleOverduePending } from '#src/services/ai/video.service.js'
import { normalizeGatewayException } from '#src/services/ai/error.service.js'
const LAST_USED_THROTTLE_MS = 60_000
const LAZY_SETTLE_THROTTLE_MS = 60_000
// 每个 key 最多每分钟扫一次过期任务;LRU 防止长期运行下无限增长
const lastSweepAt = new LRUCache({ max: 5_000, ttl: LAZY_SETTLE_THROTTLE_MS })
const hashKey = (plaintext) => createHash('sha256').update(plaintext).digest('hex')
const extractToken = (request) => {
const header = request.headers.authorization || ''
if (!header.toLowerCase().startsWith('bearer ')) return ''
return header.slice(7).trim()
}
/**
* 数据面鉴权:Bearer <api key> -> 应用 + 密钥上下文
* 每个请求查一次库(key_hash 有唯一索引),保证停用/删除立即生效。
*/
export const authenticateApiKey = async (request) => {
const token = extractToken(request)
if (!token) throw httpError(401, '缺少 API key')
const rows = await db
.select({ key: aiApiKeys, app: aiApps })
.from(aiApiKeys)
.innerJoin(aiApps, eq(aiApiKeys.app_id, aiApps.id))
.where(eq(aiApiKeys.key_hash, hashKey(token)))
.limit(1)
const row = rows[0]
if (!row) throw httpError(401, 'API key 无效')
if (row.key.status !== KEY_STATUSES.ENABLED) throw httpError(403, 'API key 已停用')
if (row.app.status !== APP_STATUSES.ENABLED) throw httpError(403, '应用已停用')
// 调用用户来自请求头:一个 key 供多人使用,按 (key, user) 归因
const rawUserId = request.headers['x-user-id']
const rawUserName = request.headers['x-user-name']
const userId = rawUserId ? String(rawUserId).trim().slice(0, 64) : null
// HTTP 头不能直接放中文,调用方通常做 percent-encode,这里解回来
let headerName = rawUserName ? String(rawUserName).trim().slice(0, 360) : null
if (headerName) {
try {
headerName = decodeURIComponent(headerName).slice(0, 120)
} catch {
headerName = headerName.slice(0, 120)
}
}
request.ai = {
appId: row.app.id,
appCode: row.app.code,
keyId: row.key.id,
userId,
userName: headerName,
billingMode: row.app.billing_mode,
}
// 惰性兜底:顺带结算该 key 过期的 pending 视频任务(不阻塞本次请求)
if (!lastSweepAt.has(row.key.id)) {
lastSweepAt.set(row.key.id, true)
settleOverduePending({
apiKeyId: row.key.id,
priceLookup: getPriceByVersion,
}).catch((error) => logger.warn({ err: error }, '惰性结算视频任务失败'))
}
// last_used_at 节流更新,避免每个请求都写库
const lastUsedAt = row.key.last_used_at ? new Date(row.key.last_used_at).getTime() : 0
if (Date.now() - lastUsedAt > LAST_USED_THROTTLE_MS) {
db.update(aiApiKeys)
.set({ last_used_at: sql`CURRENT_TIMESTAMP` })
.where(eq(aiApiKeys.id, row.key.id))
.catch((error) => logger.warn({ err: error }, '更新 last_used_at 失败'))
}
}
export default async function aiV1Autohooks(fastify) {
fastify.setErrorHandler((err, request, reply) => {
const validation = hasZodFastifySchemaValidationErrors(err)
if (validation) {
err.message = err.validation.map((item) => {
const path = item.instancePath.replace(/^\//, '').replaceAll('/', '.')
return `${path ? `${path}: ` : ''}${item.message}`
}).join('; ')
}
const candidate = Number(err.statusCode ?? err.status)
const status = Number.isInteger(candidate) && candidate >= 400 && candidate < 600 ? candidate : 500
const payload = normalizeGatewayException({ error: err, requestId: request.id, validation })
if (status >= 500) request.log.error({ err }, 'AI gateway request failed')
return reply.code(status).send(payload)
})
fastify.addHook('onRequest', authenticateApiKey)
/**
* 每个请求挂一个 AbortSignal,客户端跑掉时中止上游请求。
* 否则图片(300s)/ 非流式文本(120s)会白白把上游跑完并计费。
*
* 监听 reply.raw 的 close(响应结束或客户端提前断开都会触发)。
* 不能监听 request.raw 的 close——请求体读完就会触发,会把正常请求误判成断开。
*/
fastify.addHook('onRequest', async (request, reply) => {
const controller = new AbortController()
request.abortSignal = controller.signal
if (request.raw.aborted || request.raw.destroyed) {
controller.abort()
return
}
reply.raw.once('close', () => controller.abort())
})
}
import { chatCompletionBody } from '#src/schemas/ai/chat.js'
import { resolveBillableModel } from '#src/services/ai/model.service.js'
import { textCostMicro } from '#src/services/ai/pricing.service.js'
import { runCall, upstreamErrorCode } from '#src/services/ai/call-flow.js'
import { providerFor } from '#src/services/ai/providers/index.js'
import { aiRateLimitOptions } from '#src/services/ai/rate-limit.js'
import { AI_MODEL_TYPES } from '#src/db/schema/ai/model-mappings.js'
import { USAGE_UNITS } from '#src/db/schema/ai/ai-usage.js'
import { PRICING_UNITS } from '#src/db/schema/ai/model-prices.js'
const UPSTREAM_TIMEOUT_MS = 120_000
/** 流式特有的空闲超时:两条 chunk 之间的最大间隔 */
const STREAM_IDLE_TIMEOUT_MS = 120_000
/** 取上游上报的 usage;gateway 已统一放在 result.usage(流式与非流式同形) */
const reportedUsage = (result) => result.usage ?? null
/** 客户端断开信号 + 上游总超时;AbortSignal.any 里的 timeout 只兜「上游不响应」(流式另见 idle 超时) */
const upstreamSignal = (request, timeoutMs) => {
const timeout = AbortSignal.timeout(timeoutMs)
return request.abortSignal ? AbortSignal.any([request.abortSignal, timeout]) : timeout
}
export default async function aiChatRoutes(fastify) {
fastify.post(
'/chat/completions',
{ schema: { body: chatCompletionBody }, preHandler: fastify.rateLimit(aiRateLimitOptions) },
async (request, reply) => {
const { model, stream } = request.body
const { mapping, price } = await resolveBillableModel(model, {
type: AI_MODEL_TYPES.TEXT,
pricingUnit: PRICING_UNITS.PER_1M_TOKENS,
})
const provider = providerFor(mapping.provider, stream ? 'streamChat' : 'chat')
const common = {
request,
mapping,
model,
price,
type: AI_MODEL_TYPES.TEXT,
unit: USAGE_UNITS.TOKEN,
// 只根据上游真实 usage 计费。
settleCostMicro: (_result, usage) => textCostMicro(
price,
usage.prompt_tokens,
usage.completion_tokens,
usage.prompt_tokens_details?.cached_tokens ?? 0,
),
quantityOf: (_result, usage) => usage.prompt_tokens + usage.completion_tokens,
onError: (result) => ({ code: upstreamErrorCode(result) }),
}
// ── 流式 ──
if (stream) {
const result = await runCall({
...common,
stream: true,
usageOf: reportedUsage,
execute: () => provider.streamChat({
mapping,
body: request.body,
reply,
signal: request.abortSignal,
requestId: request.id,
idleTimeoutMs: STREAM_IDLE_TIMEOUT_MS,
}),
})
// 已 hijack,不能再 send
if (!result.ok) return reply.code(result.status).send(result.data)
return reply
}
// ── 非流式 ──
const result = await runCall({
...common,
stream: false,
usageOf: reportedUsage,
execute: () => provider.chat({
mapping,
body: request.body,
signal: upstreamSignal(request, UPSTREAM_TIMEOUT_MS),
requestId: request.id,
}),
})
if (!result.ok) return reply.code(result.status).send(result.data)
return reply.send({ ...result.data, model })
},
)
}
import { imageGenerationBody } from '#src/schemas/ai/image.js'
import { resolveBillableModel } from '#src/services/ai/model.service.js'
import { unitCostMicro } from '#src/services/ai/pricing.service.js'
import { runCall, upstreamErrorCode } from '#src/services/ai/call-flow.js'
import { providerFor } from '#src/services/ai/providers/index.js'
import { aiRateLimitOptions } from '#src/services/ai/rate-limit.js'
import { AI_MODEL_TYPES } from '#src/db/schema/ai/model-mappings.js'
import { USAGE_UNITS } from '#src/db/schema/ai/ai-usage.js'
import { PRICING_UNITS } from '#src/db/schema/ai/model-prices.js'
const UPSTREAM_TIMEOUT_MS = 300_000
/**
* 从火山方舟的图片响应里取用量(纯函数,单独可测)。
* 真实响应形如:
* { model, created, data: [{ url, size: "2048x2048" }],
* usage: { generated_images: 1, output_tokens: 16384, total_tokens: 16384 } }
*
* 计费张数优先用 usage.generated_images;只有老版本 / 异常响应没有该字段时才退化为
* 数 data 数组长度。两者都没有就返回 null,由 runCall 标 usage_missing 待核对。
* output_tokens 是我们自己按张计价之外的上游口径,记下来便于与火山账单核对。
*/
export const parseImageUsage = (data) => {
const count = data?.usage?.generated_images
?? (Array.isArray(data?.data) ? data.data.length : null)
if (count === null) return null
return {
generated_images: count,
completion_tokens: data?.usage?.output_tokens ?? 0,
}
}
export default async function aiImageRoutes(fastify) {
fastify.post(
'/images/generations',
{ schema: { body: imageGenerationBody }, preHandler: fastify.rateLimit(aiRateLimitOptions) },
async (request, reply) => {
const { model } = request.body
const { mapping, price } = await resolveBillableModel(model, {
type: AI_MODEL_TYPES.IMAGE,
pricingUnit: PRICING_UNITS.PER_IMAGE,
})
const provider = providerFor(mapping.provider, 'image')
// 真实张数以响应为准。
const result = await runCall({
request,
mapping,
model,
price,
type: AI_MODEL_TYPES.IMAGE,
unit: USAGE_UNITS.IMAGE,
execute: () => provider.image({
mapping,
body: request.body,
requestId: request.id,
signal: request.abortSignal
? AbortSignal.any([request.abortSignal, AbortSignal.timeout(UPSTREAM_TIMEOUT_MS)])
: AbortSignal.timeout(UPSTREAM_TIMEOUT_MS),
}),
// 计费数量 = 上游实际出图张数;连 data 数组都没有才算「没上报用量」,
// 交由 runCall 标记 missing + usage_missing,不估算。
usageOf: (upstream) => parseImageUsage(upstream.data),
quantityOf: (_upstream, usage) => usage.generated_images,
settleCostMicro: (_upstream, usage) => unitCostMicro(price, usage.generated_images),
onError: (upstream) => ({ code: upstreamErrorCode(upstream) }),
})
if (!result.ok) return reply.code(result.status).send(result.data)
return reply.send({ ...result.data, model })
},
)
}
import { listEnabledModels } from '#src/services/ai/model.service.js'
import { aiRateLimitOptions } from '#src/services/ai/rate-limit.js'
import { z } from 'zod'
const modelsQuery = z.object({
type: z.enum(['text', 'image', 'video']).optional(),
}).strict()
export default async function aiModelsRoutes(fastify) {
fastify.get('/models', {
schema: { querystring: modelsQuery },
preHandler: fastify.rateLimit(aiRateLimitOptions),
}, async (request, reply) => {
const rows = await listEnabledModels(request.query.type)
return reply.send({
data: rows.map((row) => ({
id: row.alias,
name: row.name,
type: row.type,
})),
})
})
}
import { appUsageRecentQuery, appUsageSummaryQuery } from '#src/schemas/ai/usage.js'
import { summarize, listRecent } from '#src/services/ai/report.service.js'
import { aiRateLimitOptions } from '#src/services/ai/rate-limit.js'
/**
* 业务系统的用量查询入口。
* 鉴权用的是应用自己的 key,所有查询自动限定在本应用内,查不到别的应用。
*/
export default async function aiAppUsageRoutes(fastify) {
fastify.get(
'/usage/summary',
{ schema: { querystring: appUsageSummaryQuery }, preHandler: fastify.rateLimit(aiRateLimitOptions) },
async (request, reply) => {
const { group_by: groupBy, from, to, user_id: userId, model, type, status } = request.query
const result = await summarize({ groupBy, from, to, userId, model, type, status, appId: request.ai.appId })
return reply.send(result)
},
)
fastify.get(
'/usage',
{ schema: { querystring: appUsageRecentQuery }, preHandler: fastify.rateLimit(aiRateLimitOptions) },
async (request, reply) => {
const { page, limit, user_id: userId, model, type, status, from, to } = request.query
const result = await listRecent(
{ appId: request.ai.appId, userId, model, type, status, from, to },
{ page, limit },
)
return reply.send(result)
},
)
}
import { and, eq } from 'drizzle-orm'
import dayjs from 'dayjs'
import { db } from '#src/db/client.js'
import { httpError } from '#src/lib/http-error.js'
import { aiUsage, USAGE_SOURCES, USAGE_STATUSES, USAGE_UNITS } from '#src/db/schema/ai/ai-usage.js'
import { PRICING_UNITS } from '#src/db/schema/ai/model-prices.js'
import { AI_MODEL_TYPES } from '#src/db/schema/ai/model-mappings.js'
import { videoSubmitBody, videoTaskParam } from '#src/schemas/ai/video.js'
import { resolveBillableModel } from '#src/services/ai/model.service.js'
import { getPriceByVersion } from '#src/services/ai/pricing.service.js'
import { resolveQuotaPackage } from '#src/services/ai/quota.service.js'
import { providerFor } from '#src/services/ai/providers/index.js'
import { buildUsageRecord, recordUsage } from '#src/services/ai/usage.service.js'
import { aiRateLimitOptions } from '#src/services/ai/rate-limit.js'
import { refreshVideoTask } from '#src/services/ai/video.service.js'
const SUBMIT_TIMEOUT_MS = 60_000
const QUERY_TIMEOUT_MS = 30_000
const findTask = async (appId, apiKeyId, taskId) => {
const rows = await db
.select()
.from(aiUsage)
.where(and(eq(aiUsage.app_id, appId), eq(aiUsage.api_key_id, apiKeyId), eq(aiUsage.upstream_task_id, taskId)))
.limit(1)
return rows[0] ?? null
}
const withTimeout = (request, timeoutMs) => {
const timeout = AbortSignal.timeout(timeoutMs)
return request.abortSignal ? AbortSignal.any([request.abortSignal, timeout]) : timeout
}
export default async function aiVideoRoutes(fastify) {
// 提交任务不占额度,完成后按真实时长计费。
fastify.post(
'/videos',
{ schema: { body: videoSubmitBody }, preHandler: fastify.rateLimit(aiRateLimitOptions) },
async (request, reply) => {
const { model, prompt, image, duration, aspect_ratio, resolution, fps, watermark } = request.body
const { mapping, price } = await resolveBillableModel(model, {
type: AI_MODEL_TYPES.VIDEO,
pricingUnit: PRICING_UNITS.PER_SECOND,
})
const provider = providerFor(mapping.provider, 'createVideo')
const submittedAt = new Date()
const quotaPackage = await resolveQuotaPackage(request.ai)
const upstream = await provider.createVideo({
mapping,
body: { model, prompt, image, duration, aspect_ratio, resolution, fps, watermark },
signal: withTimeout(request, SUBMIT_TIMEOUT_MS),
requestId: request.id,
})
const { status, data, latencyMs } = upstream
const upstreamTaskId = data?.id ?? data?.task_id ?? null
const record = (recordStatus, extra = {}) => recordUsage(buildUsageRecord({
request, mapping, model, price,
quota_package_id: quotaPackage?.id ?? null,
type: AI_MODEL_TYPES.VIDEO,
unit: USAGE_UNITS.SECOND,
quantity: 0,
usage_source: USAGE_SOURCES.MISSING,
cost: '0.000000',
status: recordStatus,
http_status: status,
latency_ms: latencyMs,
stream: 0,
created_at: dayjs(submittedAt).format('YYYY-MM-DD HH:mm:ss'),
...extra,
}))
if (!upstream.ok || !upstreamTaskId) {
await record(USAGE_STATUSES.ERROR, { error_code: data?.error?.code ?? null })
return reply.code(status).send(data)
}
await record(USAGE_STATUSES.PENDING, { upstream_task_id: upstreamTaskId })
return reply.code(202).send({
id: upstreamTaskId,
model,
status: 'queued',
video: null,
usage: null,
error: null,
})
},
)
// 查询任务:顺带结算(主路径)
fastify.get(
'/videos/:task_id',
{ schema: { params: videoTaskParam }, preHandler: fastify.rateLimit(aiRateLimitOptions) },
async (request, reply) => {
const row = await findTask(request.ai.appId, request.ai.keyId, request.params.task_id)
if (!row) throw httpError(404, '任务不存在')
// 按提交时锁定的版本取价,避免任务期间改价影响已提交的任务
const price = await getPriceByVersion(row.model, row.price_version)
if (!price) throw httpError(503, `模型未配置价格:${row.model}`)
const result = await refreshVideoTask({
row,
price,
signal: withTimeout(request, QUERY_TIMEOUT_MS),
requestId: request.id,
})
if (result.ok === false) return reply.code(result.status).send(result.data)
return reply.send(result.data ?? {
id: request.params.task_id,
model: row.model,
status: 'processing',
video: null,
usage: null,
error: null,
})
},
)
}
...@@ -31,11 +31,9 @@ const create = async (request, reply) => { ...@@ -31,11 +31,9 @@ const create = async (request, reply) => {
const list = async (request, reply) => { const list = async (request, reply) => {
const { keyword, stage, status, page, limit } = request.query const { keyword, stage, status, page, limit } = request.query
const result = await projectsService.listProjects( // 与其他路径统一传原始 request.user(ssoId + roles),
{ keyword, stage, status }, // 不再在这里另造一个 { userId, roles } 形状——同一概念两种形状正是漂移的来源
{ page, limit }, const result = await projectsService.listProjects({ keyword, stage, status }, { page, limit }, request.user || {})
request.user ? { userId: request.user.ssoId, roles: request.user.roles || [] } : {},
)
return success(reply, result) return success(reply, result)
} }
......
import { z } from 'zod'
import { pagination } from '#src/schemas/dms/common.js'
// ---------- 公共子结构 ----------
// 状态:0 停用 / 1 启用(用 int 而非 boolean,与仓储惯例一致)
const statusField = z.coerce
.number()
.int()
.refine((v) => v === 0 || v === 1, '状态仅支持 0/1')
export const billingModeEnum = z.enum(['internal', 'quota'])
// 网关只管理三类生成模型。
export const modelTypeEnum = z.enum(['text', 'image', 'video'])
// 上游供应商
export const providerEnum = z.enum(['volcano', 'deepseek'])
// 计价单位
export const pricingUnitEnum = z.enum(['per_1m_tokens', 'per_image', 'per_second'])
// 单价:元(小数)
const priceAmount = z.coerce.number().finite().nonnegative()
// 最后操作人(与业务表 operator_* 一致),请求里由登录态注入,不允许入参篡改
export const auditRowFields = {
operator_user_id: z.string().nullable(),
operator_name: z.string().nullable(),
}
// ---------- 应用 ----------
export const appCreateBody = z.object({
code: z.string().trim().min(1).max(64),
name: z.string().trim().min(1).max(120),
billing_mode: billingModeEnum.optional(),
status: statusField.optional(),
})
export const appUpdateBody = appCreateBody.partial()
export const appsListQuery = z.object({
keyword: z.string().trim().optional(),
status: statusField.optional(),
...pagination,
})
export const appRowSchema = z.object({
id: z.number().int(),
code: z.string(),
name: z.string(),
billing_mode: billingModeEnum,
status: z.number().int(),
...auditRowFields,
created_at: z.string(),
updated_at: z.string(),
enabled_key_count: z.number().int().optional(),
month_calls: z.number().int().optional(),
month_cost: z.string().optional(),
last_used_at: z.string().nullable().optional(),
})
// ---------- 密钥 ----------
export const keyCreateBody = z.object({
app_id: z.coerce.number().int().positive(),
name: z.string().trim().min(1).max(120),
status: statusField.optional(),
})
// app_id 不允许改(密钥换应用就等于重签)
export const keyUpdateBody = keyCreateBody.omit({ app_id: true }).partial()
export const keysListQuery = z.object({
app_id: z.coerce.number().int().positive().optional(),
status: statusField.optional(),
...pagination,
})
export const keyRowSchema = z.object({
id: z.number().int(),
app_id: z.number().int(),
name: z.string(),
key_prefix: z.string(),
key_masked: z.string(),
status: z.number().int(),
last_used_at: z.string().nullable(),
...auditRowFields,
created_at: z.string(),
updated_at: z.string(),
})
// 创建密钥响应额外返回一次明文 key(此后不可再取)
export const keyCreatedSchema = keyRowSchema.extend({ key: z.string() })
// ---------- 应用额度 ----------
const quotaPackageFields = {
app_id: z.coerce.number().int().positive(),
name: z.string().trim().min(1).max(120),
amount_micro: z.coerce.number().int().positive(),
start_date: z.string().regex(/^\d{4}-\d{2}-\d{2}$/),
end_date: z.string().regex(/^\d{4}-\d{2}-\d{2}$/),
status: statusField.optional(),
remark: z.string().trim().max(1000).nullish(),
}
export const quotaPackageCreateBody = z.object(quotaPackageFields)
export const quotaPackageUpdateBody = quotaPackageCreateBody.omit({ app_id: true }).partial()
export const quotaPackagesListQuery = z.object({
app_id: z.coerce.number().int().positive().optional(),
status: statusField.optional(),
...pagination,
})
export const quotaPackageRowSchema = z.object({
id: z.number().int(),
app_id: z.number().int(),
name: z.string(),
amount_micro: z.number().int(),
used_micro: z.number().int(),
start_date: z.string(),
end_date: z.string(),
status: z.number().int(),
remark: z.string().nullable(),
...auditRowFields,
created_at: z.string(),
updated_at: z.string(),
})
// ---------- 模型单价 ----------
// ---------- 响应信封(与 src/lib/response.js 的 success() 结构一致)----------
export const itemResponseSchema = (dataSchema) => z.object({ success: z.literal(true), data: dataSchema })
export const listResponseSchema = (rowSchema) =>
z.object({
success: z.literal(true),
data: z.object({
list: z.array(rowSchema),
pagination: z.object({
page: z.number().int(),
limit: z.number().int(),
total: z.number().int(),
pages: z.number().int(),
}),
}),
})
// ---------- 模型(映射 + 价格合并)----------
export const modelCreateBody = z.object({
alias: z.string().trim().min(1).max(80),
name: z.string().trim().min(1).max(120),
type: modelTypeEnum,
provider: providerEnum,
upstream_model: z.string().trim().min(1).max(120),
enabled: statusField.optional(),
pricing_unit: pricingUnitEnum,
input_price: priceAmount.optional(),
cached_input_price: priceAmount.optional(),
output_price: priceAmount.optional(),
unit_price: priceAmount.optional(),
currency: z.string().trim().max(3).optional(),
})
// alias 是主键,不允许改
export const modelUpdateBody = modelCreateBody.omit({ alias: true }).partial()
export const modelsListQuery = z.object({
keyword: z.string().trim().optional(),
type: modelTypeEnum.optional(),
provider: providerEnum.optional(),
...pagination,
})
export const modelAliasParam = z.object({
alias: z.string().trim().min(1).max(80),
})
export const modelRowSchema = z.object({
alias: z.string(),
name: z.string(),
type: modelTypeEnum,
provider: providerEnum,
upstream_model: z.string(),
enabled: z.number().int(),
version: z.number().int().nullable(),
pricing_unit: pricingUnitEnum.nullable(),
input_price: z.string().nullable(),
cached_input_price: z.string().nullable(),
output_price: z.string().nullable(),
unit_price: z.string().nullable(),
currency: z.string().nullable(),
effective_from: z.string().nullable(),
effective_to: z.string().nullable(),
has_price: z.boolean(),
...auditRowFields,
created_at: z.string(),
updated_at: z.string(),
})
import { z } from 'zod'
const message = z
.object({
role: z.string().min(1),
content: z.unknown(),
})
.passthrough()
// 文本接口遵循 OpenAI Chat Completions,只校验网关路由需要的核心字段。
export const chatCompletionBody = z
.object({
model: z.string().min(1).max(80),
messages: z.array(message).min(1),
stream: z.boolean().optional(),
})
.passthrough()
.superRefine((body, ctx) => {
for (const field of ['search', 'web_search']) {
if (field in body) ctx.addIssue({ code: 'custom', path: [field], message: '搜索不属于 AI 网关' })
}
})
import { z } from 'zod'
export const imageGenerationBody = z
.object({
model: z.string().min(1).max(80),
prompt: z.string().min(1),
})
.passthrough()
import { z } from 'zod'
export const USAGE_GROUP_BYS = ['app', 'api_key', 'user', 'model', 'provider', 'type', 'usage_source', 'status', 'day']
export const USAGE_STATUSES = ['pending', 'success', 'error', 'aborted']
export const USAGE_SOURCES = ['reported', 'missing']
export const BILLING_STATUSES = ['pending', 'charged', 'not_charged', 'reconcile']
const dateField = z.string().regex(/^\d{4}-\d{2}-\d{2}$/, '日期格式应为 YYYY-MM-DD')
const statusField = z.enum(USAGE_STATUSES)
const usageSourceField = z.enum(USAGE_SOURCES)
const billingStatusField = z.enum(BILLING_STATUSES)
const modelField = z.string().max(80)
export const usageSummaryQuery = z.object({
group_by: z.enum(USAGE_GROUP_BYS).default('model'),
from: dateField.optional(),
to: dateField.optional(),
app_id: z.coerce.number().int().positive().optional(),
api_key_id: z.coerce.number().int().positive().optional(),
type: z.enum(['text', 'image', 'video']).optional(),
usage_source: usageSourceField.optional(),
status: statusField.optional(),
model: modelField.optional(),
})
export const usageRecentQuery = z.object({
request_id: z.string().trim().max(64).optional(),
app_id: z.coerce.number().int().positive().optional(),
api_key_id: z.coerce.number().int().positive().optional(),
type: z.enum(['text', 'image', 'video']).optional(),
usage_source: usageSourceField.optional(),
billing_status: billingStatusField.optional(),
status: statusField.optional(),
model: modelField.optional(),
// 日期区间:不传时 listRecent 会补默认窗口(见 report.service.withDefaultRange),
// 所以调用方应显式带上,避免「查不到旧日志却不知道原因」
from: dateField.optional(),
to: dateField.optional(),
page: z.coerce.number().int().min(1).optional(),
limit: z.coerce.number().int().min(1).max(200).optional(),
})
// ---------- 数据面(业务系统调用,自动限定在自身应用内)----------
export const appUsageSummaryQuery = z.object({
group_by: z.enum(['user', 'model', 'type', 'day']).default('user'),
from: dateField.optional(),
to: dateField.optional(),
user_id: z.string().max(64).optional(),
type: z.enum(['text', 'image', 'video']).optional(),
status: statusField.optional(),
model: modelField.optional(),
})
export const appUsageRecentQuery = z.object({
user_id: z.string().max(64).optional(),
model: modelField.optional(),
type: z.enum(['text', 'image', 'video']).optional(),
status: statusField.optional(),
from: dateField.optional(),
to: dateField.optional(),
page: z.coerce.number().int().min(1).optional(),
limit: z.coerce.number().int().min(1).max(200).optional(),
})
import { z } from 'zod'
export const videoSubmitBody = z
.object({
model: z.string().min(1).max(80),
prompt: z.string().min(1),
image: z.string().min(1).optional(),
duration: z.coerce.number().int().positive().max(300).optional(),
aspect_ratio: z.string().regex(/^\d{1,2}:\d{1,2}$/).optional(),
resolution: z.string().regex(/^\d{3,4}p$/).optional(),
fps: z.coerce.number().int().positive().max(120).optional(),
watermark: z.boolean().optional(),
})
.strict()
export const videoTaskParam = z.object({
task_id: z.string().min(1).max(128),
})
...@@ -11,7 +11,6 @@ const caseBody = z.object({ ...@@ -11,7 +11,6 @@ const caseBody = z.object({
name: z.string().trim().min(1).max(255), name: z.string().trim().min(1).max(255),
description: z.string().trim().nullish(), description: z.string().trim().nullish(),
product_id: optionalId, product_id: optionalId,
product_name: z.string().max(255).nullish(),
files: z.array(z.unknown()).nullish(), files: z.array(z.unknown()).nullish(),
}) })
......
...@@ -5,9 +5,12 @@ export const pagination = { ...@@ -5,9 +5,12 @@ export const pagination = {
limit: z.coerce.number().int().min(1).max(500).optional(), limit: z.coerce.number().int().min(1).max(500).optional(),
} }
// null / '' 是「清空该字段」,必须解析成 null 保留下来:drizzle 的 update set 会
// 丢弃值为 undefined 的列(drizzle-orm/utils.js mapUpdateSet),解析成 undefined
// 会让清空静默失效——接口 200 但旧值还在。
export const optionalNumber = z.preprocess( export const optionalNumber = z.preprocess(
(value) => value === null || value === '' ? undefined : value, (value) => value === '' ? null : value,
z.coerce.number().optional(), z.coerce.number().nullish(),
) )
export const optionalId = z.preprocess( export const optionalId = z.preprocess(
......
...@@ -19,7 +19,6 @@ export const projectCreateBody = z.object({ ...@@ -19,7 +19,6 @@ export const projectCreateBody = z.object({
school_name: z.string().trim().min(1).max(255), school_name: z.string().trim().min(1).max(255),
department_name: z.string().max(255).nullish(), department_name: z.string().max(255).nullish(),
product_id: z.coerce.number().int().positive(), product_id: z.coerce.number().int().positive(),
product_name: z.string().max(255).nullish(),
contact_name: z.string().max(120).nullish(), contact_name: z.string().max(120).nullish(),
contact_title: z.string().max(120).nullish(), contact_title: z.string().max(120).nullish(),
contact_phone: z.string().max(64).nullish(), contact_phone: z.string().max(64).nullish(),
...@@ -30,11 +29,7 @@ export const projectUpdateBody = projectCreateBody.partial() ...@@ -30,11 +29,7 @@ export const projectUpdateBody = projectCreateBody.partial()
export const solutionUpsertBody = z.object({ export const solutionUpsertBody = z.object({
solution: z.string().nullish(), solution: z.string().nullish(),
}) attachment_file_url: fileListValue,
export const moveStageBody = z.object({
to_stage: z.coerce.number().int(),
description: z.string().nullish(),
}) })
export const rollbackStageBody = z.object({ export const rollbackStageBody = z.object({
...@@ -56,8 +51,7 @@ export const projectTeamBody = z.object({ ...@@ -56,8 +51,7 @@ export const projectTeamBody = z.object({
// ── phase upserts ──────────────────────────────────────────────────── // ── phase upserts ────────────────────────────────────────────────────
export const initiationUpsertBody = z.object({ export const initiationUpsertBody = z.object({
application_file_url: fileListValue, attachment_file_url: fileListValue,
argument_file_url: fileListValue,
project_amount: optionalNumber, project_amount: optionalNumber,
fund_source: z.string().max(255).nullish(), fund_source: z.string().max(255).nullish(),
execution_plan: z.string().nullish(), execution_plan: z.string().nullish(),
...@@ -72,8 +66,7 @@ export const procurementUpsertBody = z.object({ ...@@ -72,8 +66,7 @@ export const procurementUpsertBody = z.object({
main_bid_owner: z.string().max(120).nullish(), main_bid_owner: z.string().max(120).nullish(),
companion_bidders: z.string().nullish(), companion_bidders: z.string().nullish(),
formal_bid_status: z.string().max(120).nullish(), formal_bid_status: z.string().max(120).nullish(),
winning_notice_file_url: fileListValue, attachment_file_url: fileListValue,
bid_archive_file_url: fileListValue,
}) })
export const contractUpsertBody = z.object({ export const contractUpsertBody = z.object({
...@@ -81,7 +74,7 @@ export const contractUpsertBody = z.object({ ...@@ -81,7 +74,7 @@ export const contractUpsertBody = z.object({
contract_name: z.string().max(255).nullish(), contract_name: z.string().max(255).nullish(),
amount: optionalNumber, amount: optionalNumber,
drafter: z.string().max(120).nullish(), drafter: z.string().max(120).nullish(),
archive_file_url: fileListValue, attachment_file_url: fileListValue,
}) })
export const deliveryUpsertBody = z.object({ export const deliveryUpsertBody = z.object({
...@@ -89,9 +82,25 @@ export const deliveryUpsertBody = z.object({ ...@@ -89,9 +82,25 @@ export const deliveryUpsertBody = z.object({
delivery_contact: z.string().max(120).nullish(), delivery_contact: z.string().max(120).nullish(),
delivery_contact_phone: z.string().max(64).nullish(), delivery_contact_phone: z.string().max(64).nullish(),
delivery_note: z.string().nullish(), delivery_note: z.string().nullish(),
attachment_file_url: fileListValue,
}) })
export const acceptanceUpsertBody = z.object({ export const acceptanceUpsertBody = z.object({
acceptance_report_url: fileListValue,
acceptance_note: z.string().nullish(), acceptance_note: z.string().nullish(),
attachment_file_url: fileListValue,
})
const stageCompletionData = z.object({
...solutionUpsertBody.shape,
...initiationUpsertBody.shape,
...procurementUpsertBody.shape,
...contractUpsertBody.shape,
...deliveryUpsertBody.shape,
...acceptanceUpsertBody.shape,
}).partial()
export const moveStageBody = z.object({
to_stage: z.coerce.number().int(),
description: z.string().nullish(),
phase_data: stageCompletionData.optional(),
}) })
差异被折叠。
import { isAborted } from '#src/lib/errors.js'
import { USAGE_ERROR_CODES, USAGE_SOURCES, USAGE_STATUSES } from '#src/db/schema/ai/ai-usage.js'
import { microToCost } from './pricing.service.js'
import { chargeQuotaPackage, resolveQuotaPackage } from './quota.service.js'
import { buildUsageRecord, recordUsage } from './usage.service.js'
import { db } from '#src/db/client.js'
/** 检查软限额 → 转发 → 按真实用量计费。缺失用量只标记待核对,不估算成本。 */
export const runCall = async ({
request, mapping, model, price, type, unit,
execute, usageOf, settleCostMicro, quantityOf, onError, stream = false,
}) => {
const quotaPackage = await resolveQuotaPackage(request.ai)
const record = (values) => recordUsage(buildUsageRecord({
request, mapping, model, price, type, unit, quota_package_id: quotaPackage?.id ?? null,
quantity: 0, prompt_tokens: 0, completion_tokens: 0, cached_tokens: 0, reasoning_tokens: 0,
cost: '0.000000', stream: stream ? 1 : 0,
...values,
}))
let result
try {
result = await execute()
} catch (error) {
await record({
status: isAborted(error) ? USAGE_STATUSES.ABORTED : USAGE_STATUSES.ERROR,
usage_source: USAGE_SOURCES.MISSING,
error_code: error?.code ?? USAGE_ERROR_CODES.USAGE_MISSING,
})
throw error
}
const metrics = { http_status: result.status, latency_ms: result.latencyMs, first_token_ms: result.firstTokenMs }
if (result.ok === false || result.status < 200 || result.status >= 300) {
await record({ ...metrics, status: USAGE_STATUSES.ERROR, error_code: onError?.(result)?.code ?? null })
return result
}
const usage = usageOf?.(result)
if (!usage) {
await record({
...metrics, status: USAGE_STATUSES.SUCCESS,
usage_source: USAGE_SOURCES.MISSING, error_code: USAGE_ERROR_CODES.USAGE_MISSING,
})
return result
}
const costMicro = settleCostMicro(result, usage)
await db.transaction(async (tx) => {
await chargeQuotaPackage(quotaPackage?.id, costMicro, { tx })
await recordUsage(buildUsageRecord({
request, mapping, model, price, type, unit, quota_package_id: quotaPackage?.id ?? null,
...metrics, status: USAGE_STATUSES.SUCCESS, usage_source: USAGE_SOURCES.REPORTED,
quantity: quantityOf?.(result, usage) ?? 0,
prompt_tokens: usage.prompt_tokens ?? 0,
completion_tokens: usage.completion_tokens ?? 0,
cached_tokens: usage.prompt_tokens_details?.cached_tokens ?? 0,
reasoning_tokens: usage.completion_tokens_details?.reasoning_tokens ?? 0,
cost: microToCost(costMicro), stream: stream ? 1 : 0,
}), { tx, throwOnError: true })
})
return result
}
export const upstreamErrorCode = (result) => result?.data?.error?.code ?? null
const TYPES = {
400: 'invalid_request_error',
401: 'authentication_error',
403: 'permission_error',
404: 'not_found_error',
429: 'rate_limit_error',
}
const gatewayError = ({ code, type, message, requestId }) => ({
error: {
source: 'gateway',
code,
type,
message,
request_id: requestId ?? null,
},
})
export const normalizeGatewayException = ({ error, requestId, validation = false }) => {
const status = Number(error.statusCode ?? error.status) || 500
return gatewayError({
code: error.code ?? (validation ? 'invalid_request' : status === 401 ? 'invalid_api_key' : status === 403 ? 'permission_denied' : status === 404 ? 'not_found' : status === 429 ? 'rate_limit_exceeded' : 'gateway_error'),
type: TYPES[status] ?? 'api_error',
message: error.message,
requestId,
})
}
import { httpError } from '#src/lib/http-error.js'
import { codedError } from '#src/lib/errors.js'
import { StringDecoder } from 'node:string_decoder'
import { createUsageParser } from './usage-parser.js'
const STREAM_IDLE_TIMEOUT_CODE = 'stream_idle_timeout'
/**
* 流式空闲超时:上游多久没吐数据就判死。
* 注意不能用 AbortSignal.timeout 做流式总时限——它是从创建那一刻起的绝对超时,
* 不会因为一直在正常收数据而重置,会把超过时限的健康长回答整条掐断。
*/
const createIdleAbort = (timeoutMs) => {
const controller = new AbortController()
let timer = null
const reset = () => {
clearTimeout(timer)
timer = setTimeout(() => {
controller.abort(codedError(504, '上游流式响应超时', STREAM_IDLE_TIMEOUT_CODE))
}, timeoutMs)
timer.unref?.()
}
return {
signal: controller.signal,
reset,
clear: () => clearTimeout(timer),
}
}
const writeChunk = async (raw, chunk, signal) => {
if (raw.write(chunk)) return
await new Promise((resolve, reject) => {
const cleanup = () => {
raw.off('drain', onDrain)
raw.off('close', onClose)
signal?.removeEventListener('abort', onAbort)
}
const onDrain = () => {
cleanup()
resolve()
}
const onClose = () => {
cleanup()
reject(signal?.reason ?? new Error('客户端连接已关闭'))
}
const onAbort = () => {
cleanup()
reject(signal.reason)
}
raw.once('drain', onDrain)
raw.once('close', onClose)
signal?.addEventListener('abort', onAbort, { once: true })
})
}
/** SSE 结构保持不变,只把每个 JSON 事件里的模型改回公共模型 ID。 */
const createModelRewriter = (model) => {
const decoder = new StringDecoder('utf8')
let buffer = ''
const rewriteLine = (line) => {
if (!line.startsWith('data:')) return line
const payload = line.slice(5).trim()
if (!payload || payload === '[DONE]') return line
try {
return `data: ${JSON.stringify({ ...JSON.parse(payload), model })}`
} catch {
return line
}
}
return {
write(chunk) {
buffer += decoder.write(chunk)
const lines = buffer.split('\n')
buffer = lines.pop() ?? ''
return lines.map(rewriteLine).join('\n') + (lines.length ? '\n' : '')
},
end() {
buffer += decoder.end()
const output = buffer ? rewriteLine(buffer) : ''
buffer = ''
return output
},
}
}
/**
* 转发请求到上游。
* POST 时把 model 替换为上游真实模型,其余原样透传;GET 不带 body。
*/
const providerError = ({ mapping, data, status, requestId, providerRequestId }) => {
if (data?.error?.source === 'provider') return data
const original = data?.error ?? data
return {
error: {
source: 'provider',
provider: mapping.provider,
code: original?.code ?? 'provider_error',
type: original?.type ?? 'provider_error',
message: original?.message ?? `供应商请求失败(HTTP ${status})`,
request_id: requestId ?? null,
provider_request_id: providerRequestId ?? null,
details: data,
},
}
}
export const normalizeProviderTaskError = ({ mapping, error, requestId }) => providerError({
mapping,
data: { error },
status: 200,
requestId,
}).error
export const forwardJson = async ({ method = 'POST', path, mapping, body, signal, requestId }) => {
const provider = mapping.account
const apiKey = provider.getApiKey()
const startedAt = Date.now()
if (!apiKey) {
return {
ok: false,
status: 503,
data: { error: { source: 'gateway', code: 'upstream_not_configured', type: 'configuration_error', message: `上游账号未配置密钥:${mapping.provider}`, request_id: requestId ?? null } },
latencyMs: 0,
}
}
let response, text
try {
response = await fetch(`${provider.getBaseUrl()}${path}`, {
method,
headers: {
'content-type': 'application/json',
authorization: `Bearer ${apiKey}`,
},
body: method === 'GET' ? undefined : JSON.stringify({ ...body, model: mapping.upstreamModel }),
signal,
})
// 收到响应头不代表请求完成,读响应体也可能断连或超时。
text = await response.text()
} catch (error) {
const timedOut = error?.name === 'TimeoutError'
return {
ok: false,
status: timedOut ? 504 : 502,
data: {
error: {
code: timedOut ? 'upstream_timeout' : 'upstream_unavailable',
source: 'provider',
provider: mapping.provider,
message: timedOut ? '上游服务响应超时' : '上游服务暂时不可用',
request_id: requestId ?? null,
provider_request_id: null,
details: null,
},
},
latencyMs: Date.now() - startedAt,
}
}
const latencyMs = Date.now() - startedAt
let data
try {
data = JSON.parse(text)
} catch {
data = { error: { message: text || '上游返回非 JSON' } }
}
if (!response.ok) {
data = providerError({
mapping,
data,
status: response.status,
requestId,
providerRequestId: response.headers.get('x-request-id') ?? response.headers.get('x-tt-logid'),
})
}
return { ok: response.ok, status: response.status, data, latencyMs, usage: data?.usage ?? null }
}
/**
* 转发文本补全请求(流式)。
* 透传 SSE 分片的同时旁路解析 usage,并记录首 token 延迟。
* 上游报错时不接管响应,交回调用方按普通 JSON 处理。
*
* 超时语义:
* - idleTimeoutMs 是「两条 chunk 之间的最大间隔」,每收到一片就重置,
* 所以长回答不会被总时长误杀;
* - 连接本身沿用 AbortSignal.timeout 兜住「上游一直不回响应头」的情况。
*/
export const streamChatCompletion = async ({ mapping, path, body, reply, signal, requestId, idleTimeoutMs = 120_000 }) => {
const provider = mapping.account
const apiKey = provider.getApiKey()
if (!apiKey) throw httpError(503, `上游账号未配置密钥:${mapping.provider}`)
const startedAt = Date.now()
const idle = createIdleAbort(idleTimeoutMs)
const upstreamSignal = signal ? AbortSignal.any([signal, idle.signal]) : idle.signal
const upstreamBody = {
...body,
model: mapping.upstreamModel,
// 让上游在流末尾带上 usage;这是计费依据,必须由网关强制打开,
// 不能被调用方的 stream_options 覆盖(否则不返回 usage = 不计费)
stream_options: { ...(body.stream_options ?? {}), include_usage: true },
}
let response
try {
idle.reset()
response = await fetch(`${provider.getBaseUrl()}${path}`, {
method: 'POST',
headers: {
'content-type': 'application/json',
authorization: `Bearer ${apiKey}`,
},
body: JSON.stringify(upstreamBody),
signal: upstreamSignal,
})
} catch (error) {
idle.clear()
// 空闲超时要把原因透出去(落库 error_code),其余交给调用方按中断处理
throw idle.signal.aborted ? idle.signal.reason : error
}
if (!response.ok || !response.body) {
idle.clear()
const text = await response.text()
let data
try {
data = JSON.parse(text)
} catch {
data = { error: { message: text || '上游返回非 JSON' } }
}
return {
ok: false,
status: response.status,
data: providerError({
mapping,
data,
status: response.status,
requestId,
providerRequestId: response.headers.get('x-request-id') ?? response.headers.get('x-tt-logid'),
}),
latencyMs: Date.now() - startedAt,
}
}
// 接管响应,Fastify 不再插手
reply.hijack()
reply.raw.writeHead(response.status, {
'content-type': response.headers.get('content-type') || 'text/event-stream; charset=utf-8',
'cache-control': 'no-cache',
connection: 'keep-alive',
'x-request-id': reply.request.id,
})
const parser = createUsageParser()
const modelRewriter = createModelRewriter(mapping.alias)
let firstTokenMs = null
const iterator = response.body[Symbol.asyncIterator]()
try {
while (true) {
idle.reset()
let onAbort
const next = await Promise.race([
iterator.next(),
new Promise((_, reject) => {
onAbort = () => reject(upstreamSignal.reason)
upstreamSignal.addEventListener('abort', onAbort, { once: true })
if (upstreamSignal.aborted) onAbort()
}),
]).finally(() => upstreamSignal.removeEventListener('abort', onAbort))
if (next.done) break
const chunk = next.value
idle.reset()
const output = modelRewriter.write(chunk)
if (!output) continue
parser.feed(output)
if (firstTokenMs === null && parser.hasContent()) firstTokenMs = Date.now() - startedAt
await writeChunk(reply.raw, output, upstreamSignal)
}
} catch (error) {
// 已接管响应,异常时也必须关闭客户端连接,不能交给 Fastify send。
reply.raw.destroy()
void iterator.return?.().catch(() => {})
throw idle.signal.aborted ? idle.signal.reason : error
} finally {
idle.clear()
}
const tail = modelRewriter.end()
if (tail) {
parser.feed(tail)
await writeChunk(reply.raw, tail, upstreamSignal)
}
reply.raw.end()
return {
ok: true,
status: response.status,
latencyMs: Date.now() - startedAt,
firstTokenMs,
// 统一把上游上报的 usage 放在顶层,调用方不必区分流式 / 非流式
usage: parser.result(),
}
}
import { and, eq } from 'drizzle-orm'
import { LRUCache } from 'lru-cache'
import { db } from '#src/db/client.js'
import { aiModelMappings } from '#src/db/schema/ai/model-mappings.js'
import { httpError } from '#src/lib/http-error.js'
import { accountFor } from './upstream.js'
import { getPrice, invalidatePriceCache } from './pricing.service.js'
const cache = new LRUCache({ max: 500, ttl: 30_000 })
const typeLabels = { text: '文本', image: '图片', video: '视频' }
const pricingUnitLabels = {
per_1m_tokens: '按百万 Token 计价',
per_image: '按张计价',
per_second: '按秒计价',
}
export const resolveModel = async (alias) => {
let mapping = cache.get(alias)
if (mapping === undefined) {
const rows = await db.select().from(aiModelMappings).where(eq(aiModelMappings.alias, alias)).limit(1)
mapping = rows[0] ?? null
cache.set(alias, mapping)
}
if (!mapping) throw httpError(404, `模型不存在:${alias}`)
if (!mapping.enabled) throw httpError(400, `模型已下线:${alias}`)
return {
alias,
type: mapping.type,
provider: mapping.provider,
upstreamModel: mapping.upstream_model,
account: accountFor(mapping.provider),
}
}
export const resolveBillableModel = async (alias, { type, pricingUnit }) => {
const mapping = await resolveModel(alias)
if (mapping.type !== type) {
throw httpError(400, `模型 ${alias} 不是${typeLabels[type]}模型(${mapping.type})`)
}
const price = await getPrice(alias)
if (!price) throw httpError(503, `模型未配置价格:${alias}`)
if (price.pricing_unit !== pricingUnit) {
throw httpError(500, `模型价格配置错误:${alias} 应为${pricingUnitLabels[pricingUnit]}`)
}
return { mapping, price }
}
export const listEnabledModels = async (type) => {
const filters = [eq(aiModelMappings.enabled, 1)]
if (type) filters.push(eq(aiModelMappings.type, type))
return db
.select()
.from(aiModelMappings)
.where(and(...filters))
}
export const invalidateModelCache = (alias) => {
cache.delete(alias)
invalidatePriceCache(alias)
}
import { and, desc, eq, sql } from 'drizzle-orm'
import { LRUCache } from 'lru-cache'
import { db } from '#src/db/client.js'
import { aiModelPrices } from '#src/db/schema/ai/model-prices.js'
const cache = new LRUCache({ max: 500, ttl: 60_000 })
/**
* 「当前生效的价格」判定条件,管理面与数据面必须用同一份。
* 生效 = status=1 且(未设失效时间或失效时间还没到)。
*/
export const activePriceCondition = (alias) => and(
eq(aiModelPrices.alias, alias),
eq(aiModelPrices.status, 1),
sql`(${aiModelPrices.effective_to} IS NULL OR ${aiModelPrices.effective_to} > NOW())`,
)
/**
* 取某别名「当前生效」的价格版本:
* status=1 且 effective_to IS NULL(或未过期)中 version 最大的
*
* 暂不做分辨率档位:每个别名一个价(2026-09-10 决定)。
*/
export const getPrice = async (alias) => {
let price = cache.get(alias)
if (price === undefined) {
const rows = await db
.select()
.from(aiModelPrices)
.where(activePriceCondition(alias))
.orderBy(desc(aiModelPrices.version))
.limit(1)
price = rows[0] ?? null
cache.set(alias, price)
}
return price
}
/** 异步任务结算时按提交时记录的版本取价,避免任务期间改价影响旧任务。 */
export const getPriceByVersion = async (alias, version) => {
const key = `${alias}:${version}`
let price = cache.get(key)
if (price === undefined) {
const rows = await db
.select()
.from(aiModelPrices)
.where(and(eq(aiModelPrices.alias, alias), eq(aiModelPrices.version, version)))
.limit(1)
price = rows[0] ?? null
cache.set(key, price)
}
return price
}
/** 改价 / 停用后清掉缓存,让下一次取价立刻生效 */
export const invalidatePriceCache = (alias) => cache.delete(alias)
/**
* 文本成本(微元,1 元 = 1e6 微元)。
* 单价是「元 / 1M tokens」,所以每个 token 的微元数正好等于单价本身:
* cost(元) = tokens / 1e6 * price → cost(微元) = tokens * price
*
* 缓存命中的输入 token 单价更低;未配置 cached_input_price(为 0)时按 input_price 计。
*/
export const textCostMicro = (price, promptTokens = 0, completionTokens = 0, cachedTokens = 0) => {
if (!price) return 0
const cached = Math.min(Math.max(cachedTokens, 0), promptTokens)
const uncached = promptTokens - cached
const cachedPrice = Number(price.cached_input_price) || Number(price.input_price)
return Math.round(
uncached * Number(price.input_price) + cached * cachedPrice + completionTokens * Number(price.output_price),
)
}
/** 按张 / 按秒的成本(微元) */
export const unitCostMicro = (price, quantity = 0) => {
if (!price) return 0
return Math.round(quantity * Number(price.unit_price) * 1_000_000)
}
/** 微元 -> 入库用的 decimal 字符串 */
export const microToCost = (micro) => (micro / 1_000_000).toFixed(6)
import { forwardJson, streamChatCompletion } from '../gateway.service.js'
export const deepseekProvider = {
chat: (options) => forwardJson({ ...options, path: '/chat/completions' }),
streamChat: (options) => streamChatCompletion({ ...options, path: '/chat/completions' }),
}
import { httpError } from '#src/lib/http-error.js'
import { deepseekProvider } from './deepseek.js'
import { volcanoProvider } from './volcano.js'
const providers = {
volcano: volcanoProvider,
deepseek: deepseekProvider,
}
export const providerFor = (name, operation) => {
const provider = providers[name]
if (!provider) throw httpError(500, `未实现供应商适配器:${name}`)
if (operation && !provider[operation]) throw httpError(500, `供应商 ${name} 不支持 ${operation}`)
return provider
}
export const providerNames = Object.keys(providers)
import { forwardJson, streamChatCompletion } from '../gateway.service.js'
const videoBody = ({ model, prompt, image, duration, aspect_ratio, resolution, fps, watermark }) => {
const content = [{ type: 'text', text: prompt }]
if (image) content.push({ type: 'image_url', image_url: { url: image }, role: 'first_frame' })
return {
model,
content,
...(duration === undefined ? {} : { duration }),
...(aspect_ratio === undefined ? {} : { ratio: aspect_ratio }),
...(resolution === undefined ? {} : { resolution }),
...(fps === undefined ? {} : { framespersecond: fps }),
...(watermark === undefined ? {} : { watermark }),
}
}
export const volcanoProvider = {
chat: (options) => forwardJson({ ...options, path: '/chat/completions' }),
streamChat: (options) => streamChatCompletion({ ...options, path: '/chat/completions' }),
image: (options) => forwardJson({ ...options, path: '/images/generations' }),
createVideo: (options) => forwardJson({
...options,
path: '/contents/generations/tasks',
body: videoBody(options.body),
}),
getVideo: ({ taskId, ...options }) => forwardJson({
...options,
method: 'GET',
path: `/contents/generations/tasks/${encodeURIComponent(taskId)}`,
}),
normalizeVideo(data) {
return {
status: data?.status,
duration: data?.usage?.duration ?? data?.duration ?? null,
url: data?.content?.video_url ?? data?.video?.url ?? data?.video_url ?? null,
expiresAt: data?.expires_at ?? null,
error: data?.error ?? null,
completionTokens: data?.usage?.completion_tokens ?? 0,
}
},
}
import dayjs from 'dayjs'
import { and, eq, gte, lte, sql } from 'drizzle-orm'
import { db } from '#src/db/client.js'
import { BILLING_MODES } from '#src/db/schema/ai/apps.js'
import { aiQuotaPackages, QUOTA_PACKAGE_STATUSES } from '#src/db/schema/ai/quota-packages.js'
import { httpError } from '#src/lib/http-error.js'
/** 内部应用不限额;额度应用必须有一份当前有效的额度包。 */
export const resolveQuotaPackage = async ({ appId, billingMode }, tx = db) => {
if (billingMode === BILLING_MODES.INTERNAL) return null
const today = dayjs().format('YYYY-MM-DD')
const [quotaPackage] = await tx.select().from(aiQuotaPackages).where(and(
eq(aiQuotaPackages.app_id, appId),
eq(aiQuotaPackages.status, QUOTA_PACKAGE_STATUSES.ENABLED),
lte(aiQuotaPackages.start_date, today),
gte(aiQuotaPackages.end_date, today),
)).limit(1)
if (!quotaPackage) throw httpError(429, '当前应用没有生效中的额度')
if (Number(quotaPackage.used_micro) >= Number(quotaPackage.amount_micro)) {
throw httpError(429, '当前应用额度已用完')
}
return quotaPackage
}
/** 按供应商真实用量结算;内部应用或缺失真实用量时不扣额度。 */
export const chargeQuotaPackage = async (quotaPackageId, costMicro, { tx = db } = {}) => {
if (!quotaPackageId || costMicro <= 0) return
await tx.update(aiQuotaPackages)
.set({ used_micro: sql`${aiQuotaPackages.used_micro} + ${costMicro}` })
.where(eq(aiQuotaPackages.id, quotaPackageId))
}
import config from '#src/config.js'
/**
* AI 数据面限流:按 API key 计数(一个 key 共享一份配额)。
* 鉴权钩子在 onRequest,限流挂在 preHandler,所以这里能拿到 request.ai。
* 写接口(chat / images / videos)和读接口(models / usage / quota)共用同一份额度。
*/
export const aiRateLimitOptions = {
max: config.ai.rateLimitPerMinute,
timeWindow: '1 minute',
keyGenerator: (request) => `ai:${request.ai?.keyId ?? request.ip}`,
}
import dayjs from 'dayjs'
import { and, count, desc, eq, getTableColumns, gte, inArray, lt, sql } from 'drizzle-orm'
import { db } from '#src/db/client.js'
import { aiUsage } from '#src/db/schema/ai/ai-usage.js'
import { aiApps } from '#src/db/schema/ai/apps.js'
import { aiApiKeys } from '#src/db/schema/ai/api-keys.js'
// 标签列必须用聚合函数包一层:MySQL 的 ONLY_FULL_GROUP_BY 不允许 select 未参与分组的列
const GROUP_SPECS = {
app: { expr: aiUsage.app_id, label: sql`MAX(${aiApps.code})`, join: 'app' },
api_key: { expr: aiUsage.api_key_id, label: sql`MAX(${aiApiKeys.name})`, join: 'key' },
user: { expr: aiUsage.user_id, label: sql`MAX(${aiUsage.user_name})` },
model: { expr: aiUsage.model },
provider: { expr: aiUsage.provider },
type: { expr: aiUsage.type },
usage_source: { expr: aiUsage.usage_source },
status: { expr: aiUsage.status },
day: { expr: sql`DATE(${aiUsage.created_at})` },
}
const buildWhere = ({ from, to, appId, apiKeyId, userId, type, usageSource, billingStatus, status, model, requestId }) => {
const conditions = []
if (from) conditions.push(gte(aiUsage.created_at, from))
if (to) conditions.push(lt(aiUsage.created_at, dayjs(to).add(1, 'day').format('YYYY-MM-DD')))
if (appId) conditions.push(eq(aiUsage.app_id, appId))
if (apiKeyId) conditions.push(eq(aiUsage.api_key_id, apiKeyId))
if (userId) conditions.push(eq(aiUsage.user_id, userId))
if (type) conditions.push(eq(aiUsage.type, type))
if (usageSource) conditions.push(eq(aiUsage.usage_source, usageSource))
if (billingStatus === 'pending') conditions.push(eq(aiUsage.status, 'pending'))
if (billingStatus === 'charged') conditions.push(and(eq(aiUsage.status, 'success'), eq(aiUsage.usage_source, 'reported')))
if (billingStatus === 'not_charged') conditions.push(eq(aiUsage.status, 'error'))
if (billingStatus === 'reconcile') conditions.push(and(inArray(aiUsage.status, ['success', 'aborted']), eq(aiUsage.usage_source, 'missing')))
if (status) conditions.push(eq(aiUsage.status, status))
if (model) conditions.push(eq(aiUsage.model, model))
if (requestId) conditions.push(eq(aiUsage.request_id, requestId))
return conditions.length ? and(...conditions) : undefined
}
const AGGREGATES = {
calls: count(),
quantity: sql`COALESCE(SUM(${aiUsage.quantity}), 0)`,
token_quantity: sql`COALESCE(SUM(CASE WHEN ${aiUsage.unit} = 'token' THEN ${aiUsage.quantity} ELSE 0 END), 0)`,
image_quantity: sql`COALESCE(SUM(CASE WHEN ${aiUsage.unit} = 'image' THEN ${aiUsage.quantity} ELSE 0 END), 0)`,
second_quantity: sql`COALESCE(SUM(CASE WHEN ${aiUsage.unit} = 'second' THEN ${aiUsage.quantity} ELSE 0 END), 0)`,
prompt_tokens: sql`COALESCE(SUM(${aiUsage.prompt_tokens}), 0)`,
completion_tokens: sql`COALESCE(SUM(${aiUsage.completion_tokens}), 0)`,
cost: sql`COALESCE(SUM(${aiUsage.cost}), 0)`,
}
/** 不传时间范围时默认看最近这么多天——防止一次请求对 ai_usage 全表做聚合 */
export const DEFAULT_RANGE_DAYS = 30
/** 补默认时间窗:任一端缺失都补上,避免退化成全表扫描 */
export const withDefaultRange = (filters, now = new Date()) => {
const { from, to } = filters
if (from && to) return filters
const toDate = to ?? dayjs(now).format('YYYY-MM-DD')
const fromDate = from ?? dayjs(toDate).subtract(DEFAULT_RANGE_DAYS - 1, 'day').format('YYYY-MM-DD')
return { ...filters, from: fromDate, to: toDate }
}
/** 按维度聚合用量 */
export const summarize = async (filters) => {
const spec = GROUP_SPECS[filters.groupBy]
const where = buildWhere(withDefaultRange(filters))
let query = db
.select({ key: spec.expr, label: spec.label ?? spec.expr, ...AGGREGATES })
.from(aiUsage)
if (spec.join === 'app') query = query.leftJoin(aiApps, eq(aiUsage.app_id, aiApps.id))
if (spec.join === 'key') query = query.leftJoin(aiApiKeys, eq(aiUsage.api_key_id, aiApiKeys.id))
const rows = await query
.where(where)
.groupBy(spec.expr)
.orderBy(desc(sql`SUM(${aiUsage.cost})`))
const [total] = await db.select(AGGREGATES).from(aiUsage).where(where)
// 按计费单位拆分总量:token / image / second 无法直接相加,前端要分开展示
const units = await db
.select({ unit: aiUsage.unit, quantity: sql`COALESCE(SUM(${aiUsage.quantity}), 0)` })
.from(aiUsage)
.where(where)
.groupBy(aiUsage.unit)
return { rows, total, units }
}
/** 最近调用明细 */
export const listRecent = async (filters, { page = 1, limit = 50 } = {}) => {
const where = buildWhere(withDefaultRange(filters))
const [list, [totalRow]] = await Promise.all([
db
.select({
...getTableColumns(aiUsage),
app_name: aiApps.name,
app_code: aiApps.code,
})
.from(aiUsage)
.leftJoin(aiApps, eq(aiUsage.app_id, aiApps.id))
.where(where)
.orderBy(desc(aiUsage.id))
.limit(limit)
.offset((page - 1) * limit),
db.select({ value: count() }).from(aiUsage).where(where),
])
return {
list,
pagination: { page, limit, total: Number(totalRow.value), pages: Math.ceil(Number(totalRow.value) / limit) },
}
}
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
差异被折叠。
Markdown 格式
0% 或
您添加了 0 人 到此讨论。请谨慎行事。
请先完成此评论的编辑!
请 注册 或者 后发表评论