chore: 更新 .gitignore 并删除旧计划文档

This commit is contained in:
2026-09-03 13:20:59 +08:00
parent 0dde11dbdf
commit 8598b63aa8
29 changed files with 116 additions and 1669 deletions
+3 -1
View File
@@ -1 +1,3 @@
/.idea/*
/.idea/*
/.superpowers/
/docs/superpowers/
+8 -8
View File
@@ -10,7 +10,7 @@ database:
host: "192.168.0.83"
port: "15432"
user: "postgres"
pass: "Bjang09@686^*^"
pass: "Q!P@z#M$1@686^*^.."
name: "model-gateway"
prefix: "" # (可选)表名前缀
role: "master" # (可选)数据库主从角色(master/slave),默认为master。如果不使用应用主从机制请不配置或留空即可。
@@ -28,10 +28,10 @@ database:
timeMaintainDisabled: false # (可选)是否完全关闭时间更新特性,为true时CreatedAt/UpdatedAt/DeletedAt都将失效
model_gateway:
- type: "pgsql"
host: "localhost"
port: "5432"
host: "192.168.0.83"
port: "15432"
user: "postgres"
pass: "123456"
pass: "Q!P@z#M$1@686^*^.."
name: "model-gateway"
prefix: "model_gateway_"
role: "master"
@@ -50,17 +50,17 @@ database:
redis:
default:
address: localhost:6379
address: 192.168.0.83:6379
db: 0
consul:
address: localhost:8500
address: 192.168.0.83:8500
jaeger:
addr: localhost:4318
addr: 192.168.0.83:4318
nats:
addr: localhost
addr: 192.168.0.83
port: 4222
# schema_mapping 自动构建专用 LLM(OpenAI 兼容;密钥不要写死进代码)
+5 -5
View File
@@ -5,9 +5,9 @@ const (
)
const (
TableNameModelManage = "model_manage"
TableNameModelSession = "model_session"
TableNameModelTaskStart = "model_task_start"
TableNameModelTaskEnd = "model_task_end"
TableNameErrorMemory = "error_memory"
TableNameModelManage = "model_manage"
TableNameModelSession = "model_session"
TableNameModelTaskStart = "model_task_start"
TableNameModelTaskEnd = "model_task_end"
TableNameModelErrorMemory = "model_error_memory"
)
-25
View File
@@ -1,25 +0,0 @@
package controller
import (
"context"
"model-gateway/model/dto"
"model-gateway/service"
"gitea.redpowerfuture.com/red-future/common/beans"
)
// ErrorMemory 错误重试记忆控制器
var ErrorMemory = new(errorMemory)
type errorMemory struct{}
// List 错误重试记忆列表
func (c *errorMemory) List(ctx context.Context, req *dto.GetErrorMemoryListReq) (res *dto.GetErrorMemoryListRes, err error) {
return service.ErrorMemory.List(ctx, req)
}
// Delete 删除错误重试记忆
func (c *errorMemory) Delete(ctx context.Context, req *dto.DeleteErrorMemoryReq) (res *beans.ResponseEmpty, err error) {
err = service.ErrorMemory.Delete(ctx, req)
return
}
@@ -0,0 +1,25 @@
package controller
import (
"context"
"model-gateway/model/dto"
"model-gateway/service"
"gitea.redpowerfuture.com/red-future/common/beans"
)
// ModelErrorMemory 错误重试记忆控制器
var ModelErrorMemory = new(modelErrorMemory)
type modelErrorMemory struct{}
// List 错误重试记忆列表
func (c *modelErrorMemory) List(ctx context.Context, req *dto.GetErrorMemoryListReq) (res *dto.GetErrorMemoryListRes, err error) {
return service.ModelErrorMemory.List(ctx, req)
}
// Delete 删除错误重试记忆
func (c *modelErrorMemory) Delete(ctx context.Context, req *dto.DeleteErrorMemoryReq) (res *beans.ResponseEmpty, err error) {
err = service.ModelErrorMemory.Delete(ctx, req)
return
}
-41
View File
@@ -1,41 +0,0 @@
package dao
import (
"context"
"fmt"
"testing"
"time"
"model-gateway/consts/public"
_ "gitea.redpowerfuture.com/red-future/common/consul"
"gitea.redpowerfuture.com/red-future/common/db/gfdb"
_ "github.com/gogf/gf/contrib/drivers/pgsql/v2"
"github.com/gogf/gf/v2/net/gtrace"
)
// TestErrorMemoryDaoGetByKeyMiss 回归:错误记忆未命中必须返回 (nil, nil),
// 而非把空记录丢给 r.Struct(&res) 上浮 sql.ErrNoRows —— 否则 service.shouldRetryWithMemory
// 会把一切查询错误当 fail-closed,导致 分析→落库→重试 整条链路变死代码。
// 依赖真实 PG;不可达时跳过(不在无 PG 环境硬失败)。
func TestErrorMemoryDaoGetByKeyMiss(t *testing.T) {
ctx, span := gtrace.NewSpan(context.Background(), "TestErrorMemoryDaoGetByKeyMiss")
defer span.End()
// 探测连接:不可用则跳过
if _, err := gfdb.DB(ctx, public.DbNameModelGateway).
Model(ctx, public.TableNameErrorMemory).
NoTenantId(ctx).
Count(); err != nil {
t.Skipf("DB不可用,跳过: %v", err)
}
key := fmt.Sprintf("test-miss-%d", time.Now().UnixNano())
row, err := ErrorMemory.GetByKey(ctx, key)
if err != nil {
t.Fatalf("GetByKey 未命中不应报错: %v", err)
}
if row != nil {
t.Fatalf("GetByKey 未命中应返回 nil, got %+v", row)
}
}
@@ -8,17 +8,17 @@ import (
"gitea.redpowerfuture.com/red-future/common/db/gfdb"
)
var ErrorMemory = &errorMemoryDao{}
var ModelErrorMemory = &modelErrorMemoryDao{}
type errorMemoryDao struct{}
type modelErrorMemoryDao struct{}
// GetByKey 按记忆键查询(未命中返回 (nil, nil))
// 错误记忆为全局表:NoTenantId 绕过租户过滤,跨租户共享;r.IsEmpty() 兜底 miss 契约,
// 避免对空记录 r.Struct(&res) 上浮 sql.ErrNoRows 导致调用方 fail-closed。
func (d *errorMemoryDao) GetByKey(ctx context.Context, key string) (res *entity.ErrorMemory, err error) {
r, err := gfdb.DB(ctx, public.DbNameModelGateway).Model(ctx, public.TableNameErrorMemory).
func (d *modelErrorMemoryDao) GetByKey(ctx context.Context, key string) (res *entity.ModelErrorMemory, err error) {
r, err := gfdb.DB(ctx, public.DbNameModelGateway).Model(ctx, public.TableNameModelErrorMemory).
NoTenantId(ctx).
Where(entity.ErrorMemoryCol.MemoryKey, key).
Where(entity.ModelErrorMemoryCol.MemoryKey, key).
One()
if err != nil {
return
@@ -31,18 +31,18 @@ func (d *errorMemoryDao) GetByKey(ctx context.Context, key string) (res *entity.
}
// Upsert 存在则更新 retryable/reason/analyzed_by,不存在则插入
func (d *errorMemoryDao) Upsert(ctx context.Context, m *entity.ErrorMemory) (err error) {
model := gfdb.DB(ctx, public.DbNameModelGateway).Model(ctx, public.TableNameErrorMemory)
func (d *modelErrorMemoryDao) Upsert(ctx context.Context, m *entity.ModelErrorMemory) (err error) {
model := gfdb.DB(ctx, public.DbNameModelGateway).Model(ctx, public.TableNameModelErrorMemory)
// Count 同样全局化:跨租户已存在的记忆键需命中更新分支,而非重复插入
n, err := model.NoTenantId(ctx).Where(entity.ErrorMemoryCol.MemoryKey, m.MemoryKey).Count()
n, err := model.NoTenantId(ctx).Where(entity.ModelErrorMemoryCol.MemoryKey, m.MemoryKey).Count()
if err != nil {
return
}
if n > 0 {
_, err = model.Where(entity.ErrorMemoryCol.MemoryKey, m.MemoryKey).Data(map[string]any{
entity.ErrorMemoryCol.Retryable: m.Retryable,
entity.ErrorMemoryCol.Reason: m.Reason,
entity.ErrorMemoryCol.AnalyzedBy: m.AnalyzedBy,
_, err = model.Where(entity.ModelErrorMemoryCol.MemoryKey, m.MemoryKey).Data(map[string]any{
entity.ModelErrorMemoryCol.Retryable: m.Retryable,
entity.ModelErrorMemoryCol.Reason: m.Reason,
entity.ModelErrorMemoryCol.AnalyzedBy: m.AnalyzedBy,
}).Update()
return
}
@@ -51,20 +51,20 @@ func (d *errorMemoryDao) Upsert(ctx context.Context, m *entity.ErrorMemory) (err
}
// List 分页查询(按 id 倒序);全局表,管理端列表展示所有租户记忆
func (d *errorMemoryDao) List(ctx context.Context, page, pageSize int) (list []entity.ErrorMemory, total int64, err error) {
model := gfdb.DB(ctx, public.DbNameModelGateway).Model(ctx, public.TableNameErrorMemory).NoTenantId(ctx)
func (d *modelErrorMemoryDao) List(ctx context.Context, page, pageSize int) (list []entity.ModelErrorMemory, total int64, err error) {
model := gfdb.DB(ctx, public.DbNameModelGateway).Model(ctx, public.TableNameModelErrorMemory).NoTenantId(ctx)
n, err := model.Count()
if err != nil {
return
}
total = int64(n)
err = model.Page(page, pageSize).OrderDesc(entity.ErrorMemoryCol.Id).Scan(&list)
err = model.Page(page, pageSize).OrderDesc(entity.ModelErrorMemoryCol.Id).Scan(&list)
return
}
// Delete 按 id 删除(软删除)
func (d *errorMemoryDao) Delete(ctx context.Context, id int64) (err error) {
_, err = gfdb.DB(ctx, public.DbNameModelGateway).Model(ctx, public.TableNameErrorMemory).
Where(entity.ErrorMemoryCol.Id, id).Delete()
func (d *modelErrorMemoryDao) Delete(ctx context.Context, id int64) (err error) {
_, err = gfdb.DB(ctx, public.DbNameModelGateway).Model(ctx, public.TableNameModelErrorMemory).
Where(entity.ModelErrorMemoryCol.Id, id).Delete()
return
}
File diff suppressed because it is too large Load Diff
@@ -1,173 +0,0 @@
# 模型网关 LLM 驱动重试决策 + 持久错误记忆 设计
> **日期**:2026-09-01
> **范围**:model-gateway(仅本服务)
> **状态**:已确认
## Goal
请求上游模型失败时,不再用固定错误码清单决定是否重试,改为**调用对话模型分析错误类型**决定是否重试,并持久化"错误 → 结论"的记忆(知识库),命中记忆时不再调用分析模型。
## 背景与现状
当前重试判定(`service/retry.go:isRetryableErrorCode`)基于固定错误码集合 `429/500/501/502/503/InvalidParameter/limit_requests/limit_tokens/rate_limit_exceeded`。判定点两处:
- 同步:`service/session_sync.go:60`(`CreateSession`)
- 流式:`service/session_stream.go:88`,及 `streamRetryCodeOfError` 中对 5xx 的硬编码分支
异步任务启动 `service/model_task_start_service.go`(`CreateTask`)目前**没有重试**。
错误解析已配置驱动:`parseModelError`(按模型 `ErrorMessageMapping` 提取 code/message)。
## 设计决策(已确认)
| # | 决策 |
|---|---|
| 1 | **完全交给 LLM 判断**:删除固定错误码清单,所有错误(记忆未命中时)都调分析模型决定是否重试 |
| 2 | **持久知识库**:结论存 PostgreSQL 新表,跨实例共享、重启不丢 |
| 3 | **记忆键** = 失败上游 BaseURL + 错误码 + 归一化消息指纹(同因同键、跨实例命中) |
| 4 | **二元输出**:`{retryable: bool, reason: str}` |
| 5 | **永久有效**:条目无 TTL,直到同键重新分析覆盖或经管理端点手动删除 |
| 6 | **分析模型 = 对话模型**(不新增配置段):失败模型自身 `chat_model=true` 优先复用;否则取当前用户对话模型(`GetChatModel`);都没有 → fail-closed 不重试 |
| 7 | **三路接入**:同步 / 流式 / 异步任务启动 |
| 8 | **内联分析 + singleflight 去重**:未命中时同步调分析模型;同键并发 miss 合并为一次分析 |
## 架构与数据流
```
请求失败(parseModelError 得 code+msg)
└─ shouldRetryWithMemory(ctx, modelInfo, code, msg) → bool:
1) 归一化消息 → 记忆键 = SHA-256(upstream|code|fingerprint)
2) 查 model_gateway_error_memory
├─ 命中 → 返回存储的 retryable
└─ 未命中 → singleflight 合并后调分析模型
→ 解析 {retryable, reason} → UPSERT 落库 → 返回
3) retryable=true → 现有指数退避重试(modelCallMaxRetries=10 预算)
retryable=false → 记 ErrorMsg 走现有终态逻辑(不重试)
```
## 记忆知识库(新表 `model_gateway_error_memory`)
表前缀 `model_gateway_`,DDL 追加到 `update.sql`,DAO/entity/consts 按现有惯例新建。
| 列 | 类型 | 说明 |
|---|---|---|
| `id` | BIGSERIAL PK | |
| `memory_key` | CHAR(64) | SHA-256(upstream\|code\|fingerprint),**唯一索引** |
| `upstream` | VARCHAR | 失败上游 BaseURL |
| `error_code` | VARCHAR | 解析出的错误码 |
| `msg_fingerprint` | CHAR(32) | 归一化消息 md5 |
| `retryable` | BOOLEAN | LLM 结论 |
| `reason` | VARCHAR | LLM 给出的简短原因(可观测) |
| `analyzed_by` | VARCHAR | 分析所用模型名 |
| `created_at` / `updated_at` | 自动 | 永久有效,无 TTL |
### 归一化 `normalizeErrorMsg`
正则剔除易变片段,使同因不同实例命中同一键:
- UUID(8-4-4-4-12 hex)
- ISO 8601 时间戳 / unix 秒与毫秒数字
- `req-xxx` / `request-xxx` 类请求 ID
- 连续 ≥4 位数字(去掉具体数值,保留位置标记)
### DAO / entity / consts
`dao/model_task_start_dao.go` 模式:包名 `dao`,变量 `var ErrorMemory = &errorMemoryDao{}`;表名/库名常量进 `consts/public`;entity 带列常量(`entity.ErrorMemoryCol.*`)。
## 分析模型与调用(不新增 config)
### 分析模型选择
1. 失败请求的 `modelInfo` 本身 `chat_model=true` → 直接用它(复用其 baseURL/apiKey/HttpMethod)
2. 否则 → `GetChatModel`(当前用户对话模型,现有语义 Creator + chat_model=true)
3. 都没有 → **fail-closed 不重试** + 日志
依赖:重试循环内 `ctx` 需携带用户信息(与现有 X-User-Info 约定一致)。
### 调用方式
**用独立短超时 `http.Client{Timeout: 15s}`** 对分析模型发 OpenAI messages 格式(技术修正:`httpclient.ModelHttpNormalRequest` 底层 `modelDoRaw` 响应头超时 30min 且不区分 HTTP 状态码,用作重试判定会拖死调用,故不复用它):
```json
{
"model": "<分析模型 modelName>",
"messages": [
{"role": "system", "content": "<判定错误是否可重试的系统提示词>"},
{"role": "user", "content": "{error_code, error_message, 截断错误响应体}"}
],
"max_tokens": 256,
"temperature": 0
}
```
### Prompt(要点)
- system:定义任务——判断上游模型错误是否值得指数退避重试;考虑限流/瞬时故障/配置/参数/鉴权等类型;输出严格 JSON。
- user:携带 `error_code``error_message`、截断的错误响应体(≤2000 字符)。
- 输出解析容错:剥 ```json 代码块 / 前后空白 / 首个 `{...}` 提取。
### 失败处理
分析调用超时(15s)/失败/解析失败 → **fail-closed 不重试** + 日志。分析失败不影响记忆表。
## 三路接入
| 路径 | 文件 | 改动 |
|---|---|---|
| 同步 | `service/session_sync.go:60` | `isRetryableErrorCode(errCode)``shouldRetryWithMemory(...)` |
| 流式 | `service/session_stream.go:88` | 同上替换;`streamRetryCodeOfError` 中硬编码 5xx 分支收敛(移除) |
| 异步启动 | `service/model_task_start_service.go` | 新增重试循环(现无重试),错误时走同一判定 |
统一入口落在 `service/retry.go`:
- 删除 `isRetryableErrorCode` 固定清单
- 新增 `shouldRetryWithMemory(ctx, modelInfo, code, msg) (retry bool)`(内部:查记忆 → miss 则分析 → 落库)
- 保留 `modelCallMaxRetries=10``retryWait` 指数退避
- 新增单飞:`golang.org/x/sync/singleflight.Group` 按记忆键合并并发分析(go.sum 已有传递版本,加为直接依赖即可;避免自实现 keyed mutex)
### 管理端点(手动清理记忆)
- `GET /errorMemory/list` — 分页查看记忆条目
- `POST /errorMemory/delete` — 按 id 删除条目(永久记忆的手动纠错途径)
## 边界与降级
| 场景 | 行为 |
|---|---|
| 记忆查询失败(DB 异常) | 不重试 + 日志,不阻塞主流程 |
| 分析模型调用失败/超时/解析失败 | 不重试(fail-closed)+ 日志 |
| 同键并发 miss | singleflight 合并,只调一次分析模型 |
| 消息过长 | 截断到 2000 字符再送分析 |
| 分析模型取不到 | 不重试 + 日志 |
| 重试仍失败(同键) | 烧满 `modelCallMaxRetries` 预算;永久记忆条目保留,靠管理端点手动清理 |
## 已知取舍
- **永久记忆**:上游行为变化后旧结论可能过时。接受,靠管理端点手动纠错;自动"连续失败 N 次降级"明确**不做**(YAGNI,留作未来)。
- **异步 task_start 重试**:重试=重新调用创建任务;首次调用已报错大概率未建任务,接受"响应丢失但任务已建 → 重复建任务"的既有语义(与同步/流式行为一致)。
## 测试策略
**单元测试**
- `normalizeErrorMsg`:UUID / 时间戳 / 请求 ID / 连续数字被剔除;稳定文本不变
- 记忆键构造:同 code+同归一化消息同键;不同上游不同键
- 分析响应解析:纯 JSON / ```json 代码块 / 前后缀 / 非法响应(返回失败)
**集成测试**
- 记忆命中:直接复用存储结论,不调分析模型
- 未命中:调分析模型 → 落库 → 按结论重试/不重试
- 分析失败:不重试,不落库
- 并发同键:多个 goroutine 只触发一次分析调用
- 管理端点:list / delete
**三路径**
- 同步 / 流式:重试预算与退避行为保持现状,仅判定来源替换
- 异步 task_start:新增重试循环行为
## 非目标
- 不做固定错误码快路径(决策 1 已排除)
- 不做记忆自动过期 / 命中计数降级(决策 5 + YAGNI)
- 不新增 config 段(决策 6)
- 不接 ai-agent / prompts-core 做分析(分析在 model-gateway 内完成)
+1 -3
View File
@@ -3,7 +3,7 @@ module model-gateway
go 1.26.1
require (
gitea.redpowerfuture.com/red-future/common v0.0.32
gitea.redpowerfuture.com/red-future/common v0.0.33
github.com/bjang03/gmq v0.0.3
github.com/gogf/gf/contrib/drivers/pgsql/v2 v2.10.2
github.com/gogf/gf/contrib/nosql/redis/v2 v2.10.2
@@ -11,8 +11,6 @@ require (
golang.org/x/sync v0.19.0
)
replace gitea.redpowerfuture.com/red-future/common v0.0.32 => ../common
require (
github.com/BurntSushi/toml v1.5.0 // indirect
github.com/armon/go-metrics v0.4.1 // indirect
+1 -1
View File
@@ -36,7 +36,7 @@ func main() {
http.RouteRegister([]interface{}{
controller.ModelCall,
controller.ModelManage,
controller.ErrorMemory,
controller.ModelErrorMemory,
})
gmq.GmqRegister(public.GmqMsgPluginsName, &mq.NatsConn{
+4
View File
@@ -17,6 +17,8 @@ type ModelCallReq struct {
type ModelCallRes struct {
TaskId int64 `json:"id" dc:"任务ID"`
ModelId int64 `json:"modelId" dc:"生效模型ID(引用行=解析后的系统模型ID,计价按此)"`
MediaType string `json:"mediaType" dc:"输入媒体类型(shop词汇: text/audio/video"`
TotalTokens int64 `json:"totalTokens" dc:"总token"`
PromptTokens int64 `json:"promptTokens" dc:"输入token"`
CompletionTokens int64 `json:"completionTokens" dc:"输出token"`
@@ -61,6 +63,8 @@ type ModelCallStreamReq struct {
type ModelMsg struct {
TaskID int64 `json:"id" dc:"任务ID"`
ModelId int64 `json:"modelId" dc:"生效模型ID(引用行=解析后的系统模型ID,计价按此)"`
MediaType string `json:"mediaType" dc:"输入媒体类型(shop词汇: text/audio/video"`
TotalTokens int64 `json:"totalTokens" dc:"总token"`
PromptTokens int64 `json:"promptTokens" dc:"输入token"`
CompletionTokens int64 `json:"completionTokens" dc:"输出token"`
-24
View File
@@ -1,24 +0,0 @@
package dto
import (
"testing"
"model-gateway/model/entity"
"github.com/gogf/gf/v2/util/gconv"
)
// 引用行字段落地验证:DTO → entity 映射(CreateModelManageReq.RefSystemModelId → entity.RefSystemModelId
func TestRefSystemModelIdRoundtrip(t *testing.T) {
req := CreateModelManageReq{ModelName: "gpt-4o", RefSystemModelId: 100}
var e entity.ModelManage
if err := gconv.Struct(&req, &e); err != nil {
t.Fatalf("gconv dto->entity: %v", err)
}
if e.RefSystemModelId != 100 {
t.Fatalf("refSystemModelId mismatch: got %d want 100", e.RefSystemModelId)
}
if e.ModelName != "gpt-4o" {
t.Fatalf("modelName mismatch: got %s", e.ModelName)
}
}
+1 -1
View File
@@ -19,7 +19,7 @@ type CreateModelTaskStartReq struct {
MsgTopic string `json:"msgTopic" dc:"消息主题(可选,用于后续业务通知)"`
RequestPath string `json:"requestPath" dc:"请求参数保存路径"`
OriginalRequestPath string `json:"originalRequestPath" dc:"原始请求参数保存路径"`
MediaType string `json:"mediaType" dc:"输入媒体类型快照(audio/no_video/has_video,创建任务时按请求体推导)"`
MediaType string `json:"mediaType" dc:"输入媒体类型快照(audio/video,空=无媒体引用;shop 计费词汇,创建任务时按请求体推导)"`
}
type CreateModelTaskStartRes struct {
@@ -2,7 +2,7 @@ package entity
import "gitea.redpowerfuture.com/red-future/common/beans"
type errorMemoryCol struct {
type modelErrorMemoryCol struct {
beans.SQLBaseCol
MemoryKey string
Upstream string
@@ -13,7 +13,7 @@ type errorMemoryCol struct {
AnalyzedBy string
}
var ErrorMemoryCol = errorMemoryCol{
var ModelErrorMemoryCol = modelErrorMemoryCol{
SQLBaseCol: beans.DefSQLBaseCol,
MemoryKey: "memory_key",
Upstream: "upstream",
@@ -24,8 +24,8 @@ var ErrorMemoryCol = errorMemoryCol{
AnalyzedBy: "analyzed_by",
}
// ErrorMemory 错误重试记忆(LLM 分析结论持久化,永久有效)
type ErrorMemory struct {
// ModelErrorMemory 错误重试记忆(LLM 分析结论持久化,永久有效)
type ModelErrorMemory struct {
beans.SQLBaseDO `orm:",inline"`
MemoryKey string `orm:"memory_key" json:"memoryKey" dc:"记忆键=SHA-256(upstream|code|归一化消息)"`
Upstream string `orm:"upstream" json:"upstream" dc:"失败上游BaseURL"`
+1 -1
View File
@@ -49,6 +49,6 @@ type ModelTaskStart struct {
OriginalResponseParams map[string]any `orm:"original_response_params" json:"originalResponseParams" dc:"原始响应结果"`
DurationSeconds int64 `orm:"duration_seconds" json:"durationSeconds" dc:"耗时(秒)"`
TaskId string `orm:"task_id" json:"taskId" dc:"任务ID"`
MediaType string `orm:"media_type" json:"mediaType" dc:"输入媒体类型快照(audio/no_video/has_video,创建任务时按请求体推导)"`
MediaType string `orm:"media_type" json:"mediaType" dc:"输入媒体类型快照(audio/video,空=无媒体引用;shop 计费词汇,创建任务时按请求体推导)"`
ErrorMsg string `orm:"error_msg" json:"errorMsg" dc:"错误消息"`
}
+7 -2
View File
@@ -12,6 +12,8 @@ import (
"model-gateway/model/dto"
"model-gateway/model/entity"
"github.com/gogf/gf/v2/frame/g"
)
// parseAnalysisResponse 解析分析模型输出的判定 JSON。
@@ -71,8 +73,11 @@ func truncateStr(s string, max int) string {
}
// buildAnalysisBody 构造分析请求体(OpenAI 兼容 messages 格式),纯函数便于单测。
func buildAnalysisBody(modelName, code, msg, body string) map[string]any {
func buildAnalysisBody(ctx context.Context, modelName, code, msg, body string) map[string]any {
user := fmt.Sprintf("错误码: %s\n错误消息: %s\n错误响应体: %s", code, msg, truncateStr(body, analysisMaxBody))
// 打印错误信息
g.Log().Debugf(ctx, "分析请求体: %s", user)
return map[string]any{
"model": modelName,
"messages": []map[string]string{
@@ -100,7 +105,7 @@ func resolveAnalysisModel(ctx context.Context, modelInfo *entity.ModelManage) (m
// callAnalysisLLM 调分析模型(对话模型)判定错误是否可重试。
// 独立短超时 http.Client;非 200 / 解析失败 / 超时 → 返回 err,调用方 fail-closed。
func callAnalysisLLM(ctx context.Context, model *entity.ModelManage, code, msg, body string) (retryable bool, reason string, err error) {
reqBody, err := json.Marshal(buildAnalysisBody(model.ModelName, code, msg, body))
reqBody, err := json.Marshal(buildAnalysisBody(ctx, model.ModelName, code, msg, body))
if err != nil {
return false, "", fmt.Errorf("marshal分析请求失败: %w", err)
}
-69
View File
@@ -1,69 +0,0 @@
package service
import (
"strings"
"testing"
)
func TestParseAnalysisResponse(t *testing.T) {
cases := []struct {
name string
in string
wantRetry bool
wantErr bool
}{
{"纯JSON", `{"retryable": true, "reason": "限流"}`, true, false},
{"带json代码块", "```json\n{\"retryable\": false, \"reason\": \"参数错误\"}\n```", false, false},
{"带前后文字", `分析结果: {"retryable": true, "reason": "瞬时故障"} 完毕`, true, false},
{"非法响应", `抱歉,我无法分析`, false, true},
{"空串", ``, false, true},
}
for _, c := range cases {
retry, _, err := parseAnalysisResponse(c.in)
if c.wantErr {
if err == nil {
t.Fatalf("[%s] 期望错误, got retry=%v", c.name, retry)
}
continue
}
if err != nil {
t.Fatalf("[%s] 不应报错: %v", c.name, err)
}
if retry != c.wantRetry {
t.Fatalf("[%s] retryable = %v, want %v", c.name, retry, c.wantRetry)
}
}
}
func TestParseAnalysisResponseReason(t *testing.T) {
_, reason, err := parseAnalysisResponse(`{"retryable": true, "reason": "服务过载"}`)
if err != nil {
t.Fatalf("unexpected err: %v", err)
}
if reason != "服务过载" {
t.Fatalf("reason = %q, want 服务过载", reason)
}
}
func TestBuildAnalysisBody(t *testing.T) {
body := buildAnalysisBody("doubao-lite", "429", "too many", strings.Repeat("x", 5000))
msg := body["messages"].([]map[string]string)[1]
if len(msg["content"]) >= 2000+len("429")+len("too many")+100 {
t.Fatalf("响应体应被截断到2000字符, got %d", len(msg["content"]))
}
if body["model"] != "doubao-lite" {
t.Fatalf("model 字段错误")
}
if body["temperature"] != 0 {
t.Fatalf("temperature 应为0")
}
}
func TestTruncateStr(t *testing.T) {
if truncateStr("abc", 5) != "abc" {
t.Fatalf("短串不应截断")
}
if truncateStr("abcdef", 3) != "abc" {
t.Fatalf("截断错误")
}
}
-35
View File
@@ -1,35 +0,0 @@
package service
import (
"sync"
"testing"
"time"
)
// analyzeOnce 应合并并发同键分析请求,只执行一次 fn
func TestAnalyzeOnceDedup(t *testing.T) {
var mu sync.Mutex
calls := 0
fn := func() (bool, string) {
mu.Lock()
calls++
mu.Unlock()
// 短暂阻塞确保所有并发调用在 singleflight 窗口内到达 Do,仅首次执行 fn
time.Sleep(50 * time.Millisecond)
return true, "dedup"
}
var wg sync.WaitGroup
for i := 0; i < 10; i++ {
wg.Add(1)
go func() {
defer wg.Done()
if retry, _ := analyzeOnce("k", fn); !retry {
t.Errorf("期望 retryable=true")
}
}()
}
wg.Wait()
if calls != 1 {
t.Fatalf("singleflight应只执行一次fn, got %d", calls)
}
}
-43
View File
@@ -1,43 +0,0 @@
package service
import "testing"
func TestNormalizeErrorMsg(t *testing.T) {
cases := []struct{ in, want string }{
{"rate limit exceeded for req-abc123", "rate limit exceeded for {reqid}"},
{"timeout at 2026-09-01T09:00:00Z req_88f1a2", "timeout at {time} {reqid}"},
{"uuid 0f8fad5b-d9cb-469f-a165-70867728950e remains", "uuid {uuid} remains"},
{"err code 12345678 quota exceeded", "err code {num} quota exceeded"},
{"clean message unchanged", "clean message unchanged"},
}
for _, c := range cases {
if got := normalizeErrorMsg(c.in); got != c.want {
t.Fatalf("normalizeErrorMsg(%q) = %q, want %q", c.in, got, c.want)
}
}
}
func TestBuildMemoryKey(t *testing.T) {
k1 := buildMemoryKey("https://a.com", "429", "rate limit for req-a1b2c3")
k2 := buildMemoryKey("https://a.com", "429", "rate limit for req-d4e5f6") // 同因不同请求ID
if k1 != k2 {
t.Fatalf("同因不同请求ID应同键: %s != %s", k1, k2)
}
k3 := buildMemoryKey("https://b.com", "429", "rate limit for req-a1b2c3") // 不同上游
if k1 == k3 {
t.Fatalf("不同上游应不同键")
}
k4 := buildMemoryKey("https://a.com", "500", "rate limit for req-a1b2c3") // 不同错误码
if k1 == k4 {
t.Fatalf("不同错误码应不同键")
}
if len(k1) != 64 {
t.Fatalf("memory_key 应为 SHA-256 十六进制64位, got %d", len(k1))
}
}
func TestMsgFingerprint(t *testing.T) {
if msgFingerprint("a req-a1b2c3") != msgFingerprint("a req-d4e5f6") {
t.Fatalf("同因指纹应一致")
}
}
+1 -1
View File
@@ -197,7 +197,7 @@ func (s *modelCallService) saveModelRequestParams(ctx context.Context, now time.
return 0, nil, fmt.Errorf("上传模型解析请求参数文件失败:%v", err)
}
// 3) 保存模型请求信息(快照媒体类型任务完成时算费;模型计费配置任务完成时按 modelId 现查)
// 3) 保存模型请求信息(快照媒体类型=shop 计费词汇 audio/video,空=无媒体引用,任务完成时直接用于算费;模型计费配置任务完成时按 modelId 现查)
if *modelInfo.ResponseType == *model.ResponseTypeAsync.Code() {
id, err = dao.ModelTaskStart.Insert(ctx, &dto.CreateModelTaskStartReq{
ModelId: req.ModelId,
@@ -20,7 +20,7 @@ func shouldRetryWithMemory(ctx context.Context, modelInfo *entity.ModelManage, c
return false
}
key := buildMemoryKey(modelInfo.BaseURL, code, msg)
if row, err := dao.ErrorMemory.GetByKey(ctx, key); err != nil {
if row, err := dao.ModelErrorMemory.GetByKey(ctx, key); err != nil {
g.Log().Errorf(ctx, "查询错误重试记忆失败: %v", err)
return false
} else if row != nil {
@@ -37,7 +37,7 @@ func shouldRetryWithMemory(ctx context.Context, modelInfo *entity.ModelManage, c
g.Log().Warningf(ctx, "错误分析失败,fail-closed不重试: %v", err)
return false, ""
}
row := &entity.ErrorMemory{
row := &entity.ModelErrorMemory{
MemoryKey: key,
Upstream: modelInfo.BaseURL,
ErrorCode: code,
@@ -46,7 +46,7 @@ func shouldRetryWithMemory(ctx context.Context, modelInfo *entity.ModelManage, c
Reason: reason,
AnalyzedBy: model.ModelName,
}
if err := dao.ErrorMemory.Upsert(ctx, row); err != nil {
if err := dao.ModelErrorMemory.Upsert(ctx, row); err != nil {
g.Log().Errorf(ctx, "错误重试记忆落库失败: %v", err)
}
return r, reason
@@ -68,12 +68,12 @@ func analyzeOnce(key string, fn func() (bool, string)) (bool, string) {
return vals[0].(bool), vals[1].(string)
}
var ErrorMemory = &errorMemoryService{}
var ModelErrorMemory = &modelErrorMemoryService{}
type errorMemoryService struct{}
type modelErrorMemoryService struct{}
// List 错误重试记忆列表
func (s *errorMemoryService) List(ctx context.Context, req *dto.GetErrorMemoryListReq) (res *dto.GetErrorMemoryListRes, err error) {
func (s *modelErrorMemoryService) List(ctx context.Context, req *dto.GetErrorMemoryListReq) (res *dto.GetErrorMemoryListRes, err error) {
page, size := 1, 20
if req.Page != nil && req.Page.PageNum > 0 {
page = int(req.Page.PageNum)
@@ -81,7 +81,7 @@ func (s *errorMemoryService) List(ctx context.Context, req *dto.GetErrorMemoryLi
if req.Page != nil && req.Page.PageSize > 0 {
size = int(req.Page.PageSize)
}
list, total, err := dao.ErrorMemory.List(ctx, page, size)
list, total, err := dao.ModelErrorMemory.List(ctx, page, size)
if err != nil {
return nil, err
}
@@ -102,6 +102,6 @@ func (s *errorMemoryService) List(ctx context.Context, req *dto.GetErrorMemoryLi
}
// Delete 删除错误重试记忆(手动纠错永久记忆)
func (s *errorMemoryService) Delete(ctx context.Context, req *dto.DeleteErrorMemoryReq) (err error) {
return dao.ErrorMemory.Delete(ctx, req.Id)
func (s *modelErrorMemoryService) Delete(ctx context.Context, req *dto.DeleteErrorMemoryReq) (err error) {
return dao.ModelErrorMemory.Delete(ctx, req.Id)
}
+2
View File
@@ -314,6 +314,8 @@ LOOP:
}
// 调 shop-user-trade 按用量算费(媒体类型取任务创建时的快照;subject=解析后的系统模型 id)
docMsg.ModelId = modelInfo.Id // 引用行=系统模型 id,供 per_token 结算按系统模型计价
docMsg.MediaType = item.MediaType
docMsg.Cost = calcModelCost(asyncCtx, modelInfo.Id,
buildModelUsage(docMsg.PromptTokens, docMsg.CompletionTokens, 0, item.MediaType, docMsg.Duration))
}
+4 -14
View File
@@ -36,18 +36,6 @@ func pricingURL(sub string) string {
return "shop-user-trade/pricing/controller/" + sub
}
// mgMediaTypeToShop 媒体类型词汇映射:model-gatewayaudio/no_video/has_video)→ shop-user-tradetext/audio/video/image
func mgMediaTypeToShop(mg string) string {
switch mg {
case "audio":
return "audio"
case "has_video":
return "video"
default: // no_video
return "text"
}
}
// walletURL 组装 shop-user-trade 钱包接口地址(accountController → account/controller,与 pricing 同 RouteRegister 推导规则)
func walletURL(sub string) string {
return "shop-user-trade/account/controller/" + sub
@@ -87,7 +75,9 @@ func modelBillable(ctx context.Context, modelId int64) error {
return nil
}
// buildModelUsage 组装算费用量 JSON 对象(ChargeUsage 形状)。mediaType 为 model-gateway 词汇(DetectMediaType/异步快照)。
// buildModelUsage 组装算费用量 JSON 对象(ChargeUsage 形状)。mediaType 为 shop 计费词汇
// audio/video,空=无媒体引用走默认价;DetectMediaType/异步快照已直接为该词汇,不再二次转换)。
// per_char 模型把输出字数映射到 completionTokensTokenMapping),随该字段传给 shop /calc 计价。
func buildModelUsage(prompt, completion, cached int64, mediaType string, durationSec int64) map[string]any {
if durationSec < 0 {
durationSec = 0
@@ -96,7 +86,7 @@ func buildModelUsage(prompt, completion, cached int64, mediaType string, duratio
"promptTokens": prompt,
"completionTokens": completion,
"cachedTokens": cached,
"mediaType": mgMediaTypeToShop(mediaType),
"mediaType": mediaType,
"durationSec": durationSec,
}
}
+4
View File
@@ -130,6 +130,8 @@ LOOP:
updateModelSessionReq.DurationSeconds = int64(time.Since(startTime).Seconds())
// 调 shop-user-trade 按用量算费(不本地换算;调用前门禁已保证配置存在,失败→0 不阻塞)
mediaType := modelUtils.DetectMediaType(modelInfo.RequestBusinessFieldMapping, newRequestParams)
docMsg.ModelId = modelInfo.Id // 引用行=系统模型 id,供 per_token 结算按系统模型计价
docMsg.MediaType = mediaType
docMsg.Cost = calcModelCost(ctx, modelInfo.Id,
buildModelUsage(docMsg.PromptTokens, docMsg.CompletionTokens, 0, mediaType, 0))
updateModelSessionReq.TotalCost = docMsg.Cost
@@ -241,6 +243,8 @@ func (s *modelSessionService) CreateSessionStream(ctx context.Context, w http.Re
// 流结束:调 shop-user-trade 按用量算费(不本地换算;调用前门禁已保证配置存在,失败→0 不阻塞)
mediaType := modelUtils.DetectMediaType(modelInfo.RequestBusinessFieldMapping, newRequestParams)
docMsg.ModelId = modelInfo.Id // 引用行=系统模型 id,供 per_token 结算按系统模型计价
docMsg.MediaType = mediaType
docMsg.Cost = calcModelCost(ctx, modelInfo.Id,
buildModelUsage(docMsg.PromptTokens, docMsg.CompletionTokens, 0, mediaType, 0))
+2
View File
@@ -120,6 +120,8 @@ LOOP:
updateModelSessionReq.DurationSeconds = int64(time.Since(startTime).Seconds())
// 9.5) 调 shop-user-trade 按用量算费(不本地换算;调用前门禁已保证配置存在,失败→0 不阻塞)
mediaType := modelUtils.DetectMediaType(modelInfo.RequestBusinessFieldMapping, newRequestParams)
docMsg.ModelId = modelInfo.Id // 引用行=系统模型 id,供 per_token 结算按系统模型计价
docMsg.MediaType = mediaType
docMsg.Cost = calcModelCost(ctx, modelInfo.Id,
buildModelUsage(docMsg.PromptTokens, docMsg.CompletionTokens, 0, mediaType, 0))
updateModelSessionReq.TotalCost = docMsg.Cost
+8 -6
View File
@@ -1,21 +1,23 @@
package utils
// DetectMediaType 按模型业务字段映射从请求体推导输入媒体类型(替代硬编码的 media.type 路径)
// DetectMediaType 按模型业务字段映射从请求体推导输入媒体类型(替代硬编码的 media.type 路径)
// 直接返回 shop-user-trade 计费词汇(audio/video,对齐 ChargeUsage.MediaType),
// 无需二次转换(原 mgMediaTypeToShop 已删除):
// - reference_audio 映射路径在请求体中有值 → "audio"
// - reference_video 映射路径在请求体中有值 → "has_video"
// - 否则 → "no_video"
// - reference_video 映射路径在请求体中有值 → "video"
// - 否则 → ""(无媒体引用,shop 侧 pickModelPrice 查不到 mediaPrices 键 → 落默认价)
//
// 判定完全由模型配置(RequestBusinessFieldMapping,业务字段名见 ChatFieldsReq/VideoFields)驱动,
// 无请求结构硬编码;映射路径值即 GetByPathAll 路径(如 input.media?type=audio&url=#)。
// 媒体类型仅供 shop-user-trade 算费用量(见 service/pricing_client.go mgMediaTypeToShop)。
// 媒体类型仅供 shop-user-trade 算费用量(见 service/pricing_client.go buildModelUsage)。
func DetectMediaType(reqBizMapping map[string]string, reqParams map[string]any) string {
if hasMediaValue(reqBizMapping, reqParams, "reference_audio") {
return "audio"
}
if hasMediaValue(reqBizMapping, reqParams, "reference_video") {
return "has_video"
return "video"
}
return "no_video"
return ""
}
// hasMediaValue 业务字段映射路径在请求体中是否命中值
+8 -7
View File
@@ -312,7 +312,7 @@ COMMENT ON COLUMN model_gateway_session.total_cost
ALTER TABLE model_gateway_model_task_start
ADD COLUMN IF NOT EXISTS media_type VARCHAR(32) DEFAULT NULL;
COMMENT ON COLUMN model_gateway_model_task_start.media_type
IS '输入媒体类型快照(audio/no_video/has_video,创建任务时按请求体参考媒体字段推导)';
IS '输入媒体类型快照(audio/video,空=无媒体引用;shop 计费词汇,创建任务时按请求体参考媒体字段推导)';
ALTER TABLE model_gateway_model_task_end
ADD COLUMN IF NOT EXISTS total_cost NUMERIC DEFAULT 0;
@@ -363,14 +363,15 @@ ALTER TABLE model_gateway_model_manage
COMMENT ON COLUMN model_gateway_model_manage.error_message_mapping
IS '错误消息映射:{code,message} 的 schema 树(type/attrs/value/defaultValue),解析模型错误响应,defaultValue 为成功码';
-- =========================
-- 错误重试记忆:LLM 分析上游模型错误是否可重试的持久知识库
-- memory_key = SHA-256(upstream|error_code|归一化消息),唯一;命中直接复用,永久有效
-- =========================
CREATE TABLE IF NOT EXISTS model_gateway_error_memory (
id BIGSERIAL PRIMARY KEY,
tenant_id BIGINT DEFAULT 0,
creator VARCHAR(64) DEFAULT '',
CREATE TABLE IF NOT EXISTS model_gateway_model_error_memory (
id BIGSERIAL PRIMARY KEY,
tenant_id BIGINT DEFAULT 0,
creator VARCHAR(64) DEFAULT '',
created_at TIMESTAMPTZ DEFAULT now(),
updater VARCHAR(64) DEFAULT '',
updated_at TIMESTAMPTZ DEFAULT now(),
@@ -382,7 +383,7 @@ CREATE TABLE IF NOT EXISTS model_gateway_error_memory (
retryable BOOLEAN NOT NULL DEFAULT false,
reason VARCHAR(512) NOT NULL DEFAULT '',
analyzed_by VARCHAR(128) NOT NULL DEFAULT ''
);
);
CREATE UNIQUE INDEX IF NOT EXISTS uk_error_memory_memory_key
ON model_gateway_error_memory (memory_key)
ON model_gateway_model_error_memory (memory_key)
WHERE deleted_at IS NULL;