# Model Gateway — 智能模型网关 统一的 AI 模型网关服务,负责多模态(文本、图像、音频、向量、视频、全模态)模型调用的路由、执行与管理。基于 Go (GoFrame v2) 构建,使用 PostgreSQL、Redis 和 Consul 服务发现。 --- ## 架构概览 ``` 客户端 ──> HTTP API ──> Controller ──> Service ──> DAO ──> PostgreSQL │ └─> Redis(信号量/队列) │ ├──> AI 模型服务商(OpenAI、阿里云、火山引擎等) ├──> OSS 文件服务 ├──> Admin-go(租户/余额) ├──> Prompts-core(会话回调) └──> Skill 技能服务 ``` 系统分为六个层次: | 层次 | 目录 | 职责 | |---|---|---| | **控制器层** | `controller/` | HTTP 路由注册与请求绑定 | | **业务逻辑层** | `service/` | 模型配置 CRUD、异步任务执行、提示词构建、队列管理、统计 | | **数据访问层** | `dao/` | PostgreSQL 数据访问(GoFrame ORM) | | **模型层** | `model/` | DTO 请求/响应结构与数据库实体定义 | | **通用工具层** | `common/` | 计费、类型转换、文件处理、请求头、映射、提示词、流式解析等工具 | | **常量层** | `consts/` | 公共常量和表名定义 | --- ## 核心功能 ### 1. 模型配置管理 (`model_gateway_models`) 提供 AI 模型服务配置的完整增删改查: | 字段 | 说明 | |---|---| | `model_name` | 模型名称(唯一标识) | | `model_type` | 模型分类(100=推理, 200=图像, 300=音频, 400=向量, 500=全模态, 600=视频) | | `operator_name` | 运营商标识(OpenAI、阿里云、火山引擎等) | | `base_url` | 模型服务地址 | | `http_method` | GET / POST | | `head_msg` | 每次调用注入的请求头 | | `form_json` | 动态表单定义(用于前端按模型渲染参数表单) | | `request_mapping` / `response_mapping` | 标准格式与提供商 API 之间的参数映射 | | `call_mode` | 0=同步, 1=异步, 2=流式 | | `max_concurrency` | 单模型最大并发数(按租户维度限流) | | `timeout_seconds` | 请求超时时间(秒) | | `retry_times` | 失败重试次数 | | `billing_config` | 计费规则(推理阶梯计价 / 视频分辨率计价) | | `stream_config` | SSE 流式解析配置 | | `query_config` | 异步任务轮询/查询配置 | **支持的模型类型:** | 类型码 | 类别 | 子类型 | |---|---|---| | 100 | 推理模型 | 文本生成、对话 | | 200-205 | 图像模型 | 文生图、图生图、图片编辑、图片变体、图文生图 | | 300-303 | 音频模型 | 文生语音、语音转文字、语音转语音 | | 400-402 | 向量模型 | 文本嵌入、重排序 | | 500-502 | 全模态模型 | 文图音频、视觉理解 | | 600-604 | 视频模型 | 文生视频、图生视频、图文生视频、视频生视频 | ### 2. 异步任务执行 (`model_gateway_task`) 任务生命周期如下: ``` CreateTask ──> state=0(排队中) │ Worker 抢占(state=1,执行中) │ ├── 同步模式:直接调用模型 ├── 异步模式:提交任务,通过 QueryConfig 轮询 └── 流式模式:通过 StreamConfig 解析 SSE 事件 │ ├── 成功(state=2)──> 上传结果到 OSS │ └──> 触发回调(如有配置) └── 失败(state=3)──> 重试(最多 retry_times 次) │ 客户端下载(state=4,已下载) ``` **关键特性:** - **任务创建**:立即返回 `taskId`,执行在 goroutine 中异步完成 - **并发控制**:基于 Redis 的分布式信号量,按模型维度限制并发数 - **队列门控**:基于 Redis Lua 脚本的严格队列插槽机制,防止分布式创建下超限 - **自动重试**:临时性错误(超时、内部错误)自动重试;硬错误直接失败 - **OSS 上传**:结果以 JSON 格式上传到 OSS 文件服务,OSS URL 存入任务记录 - **回调通知**:任务完成时可选触发 HTTP 回调(`TriggerCallback`、`TriggerPromptsCallback`、`CallbackBuildResult`) - **阶梯计费**:支持推理 token 阶梯计价和视频分辨率计价两种模型 ### 3. 提示词构建 (`/buildMessages`) 异步构建推理模型的结构化消息: - 从配置合并系统提示词(`modelPrompts.types`) - 拉取技能 Markdown 内容(`SkillMdContent`) - 通过 `GetSessionHistory` 注入会话历史 - 根据附件数量自动拆分多轮(`SplitByAttachment`) - 支持视频模型通过对话模型编排的多轮生成 ### 4. 动态调参 (`/autoTune`) 周期性自动调参(建议每小时触发),基于近期 P90 执行耗时和到达率动态调整 `max_concurrency` 和队列上限: - 读取模型配置(数据库中的上限值 cap) - 根据时间窗口内已完成任务统计 P90 执行耗时 - 使用 Little 定律(到达率 × P90 / 利用率)计算新的并发数 - 单次调整幅度限制在 ±50% - 运行时参数写入 Redis(2 小时 TTL),不修改数据库中的上限值 ### 5. 日统计 (`model_gateway_log_stat`) 按天、租户、创建人、模型维度的请求计数,使用原子 upsert(`ON DUPLICATE KEY UPDATE`)。 按天、租户、创建人、模型维度的请求计数,使用原子 upsert(`ON DUPLICATE KEY UPDATE`)。 ### 6. 核心接口详细说明 #### 6.1 构建消息结构 `/buildMessages` 构建推理模型的提示词消息体,将系统提示词、技能知识、会话历史、用户自定义提示词合并为最终的消息结构。 **请求示例(JSON):** ```json { "modelName": "gpt-4o", "buildType": 1, "skillName": "code-review", "callbackUrl": "http://callback.example.com/result", "nodeId": "node_001", "sessionId": "session_abc", "customPrompt": "请用中文回答", "messages": { "model": "gpt-4o", "max_tokens": 4096, "messages": [ {"role": "system", "content": [{"type": "text", "text": "你是一名助手"}]}, {"role": "user", "content": [{"type": "text", "text": "写一封邮件"}]} ] } } ``` | 请求字段 | 类型 | 必填 | 说明 | |---|---|---|---| | `modelName` | string | 是 | 网关模型名称 | | `buildType` | int | 是 | 构建类型:1=单轮构建, 2=多轮构建 | | `skillName` | string | 否 | 技能名称,用于拉取技能 MD 知识 | | `callbackUrl` | string | 否 | 构建完成后的回调地址 | | `nodeId` | string | 否 | 节点 ID(用于查询会话历史) | | `sessionId` | string | 否 | 会话 ID | | `messages` | object | 是 | 前端构建的消息结构 | | `customPrompt` | string | 否 | 用户自定义提示词 | **响应:** ```json { "taskId": "550e8400-e29b-41d4-a716-446655440000" } ``` **处理流程:** 1. **查模型配置** — 根据 `modelName` 查询 `model_gateway_models` 获取模型定义 2. **创建构建记录** — 写入 `model_gateway_build_record`,状态为处理中 3. **异步执行构建**(goroutine): - **推理模型(type=100)**: - 从配置读取系统提示词(`modelPrompts.types.100`) - 若有 `skillName`,通过 `GetSkillUser` 拉取技能 ZIP,提取 MD 内容 - 合并系统提示词 + 技能知识 + 用户自定义提示词到 `messages` - 通过 `GetSessionHistory` 注入会话历史 - 按附件数量自动拆分多轮(`SplitByAttachment`,按 video/image/audio 的 `maxCount` 约束) - 返回单轮或多轮消息结构 - **视频模型(type=600-699)**: - 查找当前用户的对话模型(`is_chat_model=1`) - 获取该对话模型的协议模板(`prompts_provider_protocol`) - 构建对话模型请求体,调用对话模型生成视频分镜的 JSON rounds - 将 rounds 上传 OSS,存入构建记录 4. **回写结果** — 更新构建记录状态(成功/失败 + 耗时) 5. **触发回调** — 若有 `callbackUrl`,POST 回调通知任务完成 **对应源码:** - Controller:[controller/model_gateway_task_controller.go](`ModelGatewayTask.BuildMessages`) - Service:[service/task/task_service.go](`taskService.BuildMessages` → `executeBuild` → `buildResult`) - Util(提示词合并):[common/util/prompt.go](`MergePrompt`、`InjectHistory`、`SplitByAttachment`) - Util(技能知识):[service/prompt/prompt_files_handle_service.go](`SkillMdContent`) #### 6.2 创建异步任务 `/createTask` 创建异步模型调用任务,立即返回 `taskId`,后端 goroutine 异步执行模型调用、结果上传与计费。 **请求示例(JSON):** ```json { "modelName": "dall-e-3", "bizName": "image-generator", "callbackUrl": "http://callback.example.com/result", "epicycleId": 12345, "buildModelName": "gpt-4o", "requestPayload": { "model": "dall-e-3", "prompt": "一只站在树枝上的猫头鹰,水墨风格", "n": 1, "size": "1024x1024" } } ``` | 请求字段 | 类型 | 必填 | 说明 | |---|---|---|---| | `modelName` | string | 是 | 模型名称 | | `bizName` | string | 否 | 业务名称,用于统计区分 | | `callbackUrl` | string | 否 | 任务完成后的回调地址 | | `requestPayload` | object | 是 | 透传给模型服务的请求参数 | | `epicycleId` | int64 | 否 | 轮次 ID(prompts-core 场景) | | `buildModelName` | string | 否 | 构建阶段使用的模型名(prompts-core 场景) | **响应:** ```json { "taskId": "550e8400-e29b-41d4-a716-446655440000" } ``` **处理流程:** 1. **鉴权与参数校验** — 获取用户信息,校验模型配置是否存在且已启用 2. **队列门控(可选)** — 通过 Redis Lua 脚本做严格的分布式队列插槽校验,避免并发超限 3. **写入任务记录** — 状态为 `1(执行中)`,记录 `requestPayload`、`callbackUrl`、`bizName` 等 4. **操作日志** — 记录创建任务的审计日志到 `model_gateway_logs_op` 5. **计费预处理** — 若模型配置了 `billing_config`,从请求参数中提取计费维度数据 6. **异步执行**(goroutine `AsyncWorker.handleOne`): ``` 检查租户余额 ──> 调用模型服务 ──> 解析响应映射 ──> 计费计算并扣费 │ ┌── 同步模式(call_mode=0):直接调用,等待 HTTP 响应 ├── 异步模式(call_mode=1):提交任务后通过 QueryConfig 轮询结果 └── 流式模式(call_mode=2):通过 StreamConfig 解析 SSE 事件流 │ ┌── prompts-core 场景:对模型输出做 ParseAndValidate + 重试机制 └── 普通场景:直接映射响应 │ 上传 OSS ──> 更新任务为成功 ──> 触发回调 ``` 7. **重试机制**:超时、内部错误等可重试错误会按 `retry_times` 重试;硬错误(参数错误等)直接失败 8. **回调通知**:任务成功或失败后,若有 `callbackUrl` 则 POST 回调;prompts-core 场景额外触发 `TriggerPromptsCallback` **调用模式对比:** | 模式 | 值 | 行为 | |---|---|---| | 同步 | 0 | 直接 HTTP 调用模型,等待完整响应后返回 | | 异步 | 1 | 提交任务后通过 `QueryConfig` 配置的轮询地址定时查询结果 | | 流式 | 2 | 对 SSE 事件流按 `StreamConfig` 规则解析(支持 concat/base64_concat/collect/final) | **对应源码:** - Controller:[controller/model_gateway_task_controller.go](`ModelGatewayTask.CreateTask`) - Service:[service/task/task_service.go](`taskService.Create`) - Worker:[service/task/worker.go](`asyncWorker.handleOne`、`InvokeModel`) - 队列门控:[service/queue/queue_gate.go](`AcquireQueueSlot`、`ReleaseQueueSlot`) - 运行时调参:[service/queue/runtime_tune.go](`GetRuntimeMaxConcurrency`、`GetRuntimeQueueLimit`) - 异步轮询:[common/util/pull_task.go](`PullTaskResult`) - 流式解析:[common/util/streaming.go](`ParseStreamResponse`) - 计费:[common/util/billing.go](`CalculateBilling`、`ExtractRequestBilling`) ### 7. 全量 API 接口一览 | 分类 | 接口路径 | 方法 | 说明 | |---|---|---|---| | **模型配置** | `/createModel` | POST | 创建模型配置 | | | `/updateModel` | PUT | 更新模型配置 | | | `/deleteModel` | DELETE | 删除模型配置 | | | `/getModel` | GET | 获取模型详情 | | | `/listModel` | GET | 模型列表(分页+筛选) | | | `/listType` | GET | 模型类型列表 | | | `/listOperator` | GET | 运营商列表 | | | `/updateChatModel` | POST | 设置当前用户的对话模型 | | | `/getIsChatModel` | GET | 获取当前对话模型 | | | `/autoTune` | POST | 触发动态调参 | | **任务管理** | `/createTask` | POST | 创建异步任务 | | | `/jobTask` | POST | 定时批量任务处理器 | | | `/getTaskResult` | GET | 获取单条任务结果 | | | `/getTaskBatch` | POST | 批量查询任务(成功任务自动标记为已下载) | | | `/listTask` | GET | 任务列表分页查询 | | | `/modelCallback` | POST | 接收异步模型回调通知 | | | `/queryPending` | GET | 轮询进行中的异步任务 | | **提示词** | `/buildMessages` | POST | 构建结构化提示词消息(异步) | | **统计** | `/listModelStat` | GET | 日使用量统计 | -- ## 技术栈 | 组件 | 技术选型 | |---|---| | 语言 | Go 1.26 | | 框架 | GoFrame v2(`github.com/gogf/gf/v2`) | | 数据库 | PostgreSQL(GoFrame pgsql 驱动) | | 缓存/队列 | Redis(分布式信号量、队列门控、运行时调参存储) | | 服务发现 | Consul | | 链路追踪 | Jaeger(OTLP HTTP) | | 模型调用 | 原生 HTTP(可定制请求头、鉴权、超时) | | 文件存储 | OSS 文件上传服务(multipart 表单) | ## 数据库表 | 表名 | 说明 | |---|---| | `model_gateway_models` | 模型服务配置(动态表单、请求/响应映射、计费规则) | | `model_gateway_task` | 异步任务记录(状态机、重试、OSS 结果文件) | | `model_gateway_build_record` | 提示词构建记录 | | `model_gateway_logs_op` | 操作审计日志 | | `model_gateway_logs_stat` | 日维度请求量统计 | | `prompts_provider_protocol` | 视频模型编排的提供商协议模板 | ## 配置说明 详见 `config.yml`。主要配置块: - `database` — PostgreSQL 连接(双数据源:`default` + `model_gateway`) - `redis` — Redis 连接 - `consul` — Consul 地址 - `jaeger` — Jaeger OTLP HTTP 端点 - `queryPending` — 异步任务自动轮询配置(调试开关) - `jobTask` — 批量处理器间隔、批大小、协程池大小 - `modelPrompts.types` — 各模型类型的系统提示词(100/200/300/400/500) - `nodePrompts` — 节点路由提示词模板 ## 快速开始 ### 环境要求 - Go 1.26+ - PostgreSQL - Redis - Consul(可选,用于服务发现) ### 本地开发 1. 克隆仓库 2. 在 PostgreSQL 中执行 `update.sql` 创建所有表 3. 修改 `config.yml` 中的数据库、Redis 和 Consul 配置 4. 启动服务: ```bash go run main.go ``` 5. 服务默认监听 `3004` 端口(可在 `config.yml` 中修改) ### Docker 构建 ```bash docker build -t model-gateway . docker run -p 3004:3004 model-gateway ``` Dockerfile 使用多阶段构建,基于 `golang:alpine` 镜像,配置了国内 Go 代理,输出剥离调试信息的精简二进制。 --- ## 项目目录结构 ``` . ├── main.go # 入口:路由注册、自动执行器启动 ├── config.yml # 应用配置 ├── Dockerfile # 多阶段 Docker 构建 ├── go.mod / go.sum # Go 模块依赖 ├── update.sql # 数据库 DDL │ ├── controller/ # HTTP 接口控制器 │ ├── model_gateway_models_controller.go │ ├── model_gateway_task_controller.go │ └── model_gateway_logs_stat_controller.go │ ├── service/ # 业务逻辑层 │ ├── gateway/ # OSS 上传、回调、余额、技能/会话查询 │ ├── model/ # 模型 CRUD、对话模型管理 │ ├── prompt/ # 文件拉取、技能 Markdown、提示词构建 │ ├── queue/ # 动态调参、信号量、队列门控、运行时调参 │ ├── stat/ # 使用量统计 │ └── task/ # 任务创建、Worker 执行、异步结果处理 │ ├── dao/ # 数据访问层 │ ├── model_gateway_models_dao.go │ ├── model_gateway_task_dao.go │ ├── model_gateway_build_record_dao.go │ ├── model_gateway_logs_stat_dao.go │ ├── model_gateway_logs_op_dao.go │ └── provider_protocol_dao.go │ ├── model/ # 数据模型 │ ├── dto/ # 请求/响应结构体 │ └── entity/ # 数据库实体 │ ├── common/ # 通用工具 │ └── util/ # 计费、类型转换、文件处理、请求头、映射、提示词、流式、异步轮询 │ └── consts/ # 常量定义 └── public/ # 模型类型、任务状态、运营商列表、表名 ``` --- ## 任务状态机 | 状态码 | 含义 | 转换来源 | |---|---:|---| | 0 | 排队中 | 创建任务或重试入队 | | 1 | 执行中 | Worker 抢占成功 | | 2 | 成功(已上传 OSS) | Worker 执行完成 | | 3 | 失败 | Worker 执行出错或超时 | | 4 | 已下载 | 批量查询接口标记 | --- ## 接口文档 地址:https://s.apifox.cn/42764602-cb82-45fa-887f-c751311ba406 ## License 内部项目 — 红未来科技