请求明细支持查看原始请求/响应:OT_PROXY_LOG_RAW 开关控制,仅管理员记录与可见,流式全量捕获
- 配置: ProxyConfig.LogRaw (OT_PROXY_LOG_RAW, 默认 false) - 存储: usage_logs 新增 raw_request/raw_response 文本列 (AutoMigrate) - 网关: NewGateway 接收 logRaw 参数 - handlers: 三个协议入口按 开关+管理员 条件记录原始请求体 - passthrough: 非流式 copyAndCapture 捕获响应, 流式 streamCopy 累积全部原始 SSE 行, finishUsage 统一写入 - admin API: AdminUsage 返回 raw_request/raw_response (仅管理员) - 前端: 用量页新增查看入口, 弹窗 tab 切换请求/响应 - gitignore: 修正 server/web/ 忽略规则(尾随空格导致未生效)
This commit is contained in:
@@ -31,7 +31,7 @@ func main() {
|
||||
}
|
||||
defer a.Shutdown(context.Background())
|
||||
|
||||
gw := proxy.NewGateway(a.DB, a.Enc, a.Usage, a.Limit, cfg.RateLimit.UserRPS)
|
||||
gw := proxy.NewGateway(a.DB, a.Enc, a.Usage, a.Limit, cfg.RateLimit.UserRPS, cfg.Proxy.LogRaw)
|
||||
router := api.NewRouter(a, gw)
|
||||
|
||||
srv := &http.Server{
|
||||
|
||||
@@ -95,6 +95,7 @@ func (h *Handler) AdminUsage(c *gin.Context) {
|
||||
"input_tokens": l.InputTokens, "output_tokens": l.OutputTokens,
|
||||
"cache_read_tokens": l.CacheReadTokens, "cost": l.Cost,
|
||||
"latency_ms": l.LatencyMS, "status": l.Status, "error_code": l.ErrorCode,
|
||||
"raw_request": l.RawRequest, "raw_response": l.RawResponse,
|
||||
"created_at": l.CreatedAt,
|
||||
})
|
||||
}
|
||||
|
||||
@@ -66,6 +66,7 @@ type ProxyConfig struct {
|
||||
Timeout time.Duration
|
||||
HealthInterval time.Duration // 渠道健康检查周期
|
||||
HealthFailThreshold int // 连续失败 N 次进 cooldown
|
||||
LogRaw bool // 记录管理员原始请求体+响应到 usage_logs(调试用,默认关)
|
||||
}
|
||||
|
||||
// loadDotEnv 读取 .env 并把 KEY=VALUE 注入环境变量(AutomaticEnv 自动映射 OT_ 前缀)。
|
||||
@@ -129,6 +130,7 @@ func Load() (*Config, error) {
|
||||
v.SetDefault("proxy.timeout", "120s")
|
||||
v.SetDefault("proxy.health_interval", "60s")
|
||||
v.SetDefault("proxy.health_fail_threshold", 2)
|
||||
v.SetDefault("proxy.log_raw", false)
|
||||
|
||||
v.SetDefault("ratelimit.user_rps", 20)
|
||||
|
||||
@@ -168,6 +170,7 @@ func Load() (*Config, error) {
|
||||
Timeout: v.GetDuration("proxy.timeout"),
|
||||
HealthInterval: v.GetDuration("proxy.health_interval"),
|
||||
HealthFailThreshold: v.GetInt("proxy.health_fail_threshold"),
|
||||
LogRaw: v.GetBool("proxy.log_raw"),
|
||||
},
|
||||
RateLimit: RateLimitConfig{
|
||||
UserRPS: v.GetInt("ratelimit.user_rps"),
|
||||
|
||||
@@ -35,6 +35,7 @@ type Gateway struct {
|
||||
enc *crypto.Encryptor
|
||||
lim *ratelimit.Limiter
|
||||
userRPS int
|
||||
logRaw bool
|
||||
hc *http.Client
|
||||
|
||||
policyMu sync.Mutex
|
||||
@@ -108,7 +109,7 @@ func contains(list []string, s string) bool {
|
||||
return false
|
||||
}
|
||||
|
||||
func NewGateway(db *gorm.DB, enc *crypto.Encryptor, rec *usage.Recorder, lim *ratelimit.Limiter, userRPS int) *Gateway {
|
||||
func NewGateway(db *gorm.DB, enc *crypto.Encryptor, rec *usage.Recorder, lim *ratelimit.Limiter, userRPS int, logRaw bool) *Gateway {
|
||||
return &Gateway{
|
||||
db: db,
|
||||
ch: channel.NewService(db, enc),
|
||||
@@ -116,6 +117,7 @@ func NewGateway(db *gorm.DB, enc *crypto.Encryptor, rec *usage.Recorder, lim *ra
|
||||
enc: enc,
|
||||
lim: lim,
|
||||
userRPS: userRPS,
|
||||
logRaw: logRaw,
|
||||
hc: &http.Client{Timeout: 120 * time.Second},
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import (
|
||||
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/openteam/server/internal/proxy/convert"
|
||||
"github.com/openteam/server/internal/store"
|
||||
)
|
||||
|
||||
// chatCompletions POST /v1/chat/completions
|
||||
@@ -23,6 +24,7 @@ func (g *Gateway) chatCompletions(c *gin.Context) {
|
||||
}
|
||||
c.Set("protocol", convert.ProtoChat)
|
||||
c.Set("model_name", br.Model)
|
||||
g.recordRawRequest(c, u, body)
|
||||
if !g.checkModelAllowed(u, br.Model) {
|
||||
apiError(c, http.StatusForbidden, "model_not_allowed", "模型未对你开放,请联系管理员")
|
||||
return
|
||||
@@ -55,6 +57,7 @@ func (g *Gateway) responses(c *gin.Context) {
|
||||
}
|
||||
c.Set("protocol", convert.ProtoResponses)
|
||||
c.Set("model_name", br.Model)
|
||||
g.recordRawRequest(c, u, body)
|
||||
if !g.checkModelAllowed(u, br.Model) {
|
||||
apiError(c, http.StatusForbidden, "model_not_allowed", "模型未对你开放,请联系管理员")
|
||||
return
|
||||
@@ -87,6 +90,7 @@ func (g *Gateway) messages(c *gin.Context) {
|
||||
}
|
||||
c.Set("protocol", convert.ProtoMessages)
|
||||
c.Set("model_name", br.Model)
|
||||
g.recordRawRequest(c, u, body)
|
||||
if !g.checkModelAllowed(u, br.Model) {
|
||||
apiError(c, http.StatusForbidden, "model_not_allowed", "模型未对你开放,请联系管理员")
|
||||
return
|
||||
@@ -108,6 +112,14 @@ type sinkHolder struct {
|
||||
sink *usageSink
|
||||
}
|
||||
|
||||
// recordRawRequest 记录管理员原始请求体到 context(供 finishUsage 落库)。
|
||||
// 仅当开关开启且用户为管理员时记录;响应侧以 c.Get("raw_request") 是否非空判断是否需要捕获响应。
|
||||
func (g *Gateway) recordRawRequest(c *gin.Context, u *store.User, body []byte) {
|
||||
if g.logRaw && u.Role == store.RoleAdmin {
|
||||
c.Set("raw_request", string(body))
|
||||
}
|
||||
}
|
||||
|
||||
// apiError 按客户端协议返回错误体(PLANNING §5.1.4)。
|
||||
func apiError(c *gin.Context, status int, code, message string) {
|
||||
if p, _ := c.Get("protocol"); p == convert.ProtoMessages {
|
||||
|
||||
@@ -272,6 +272,9 @@ func (g *Gateway) copyAndCapture(c *gin.Context, ch *store.Channel, r io.Reader,
|
||||
}
|
||||
}
|
||||
_, _ = c.Writer.Write(out)
|
||||
if _, ok := c.Get("raw_request"); ok {
|
||||
c.Set("raw_response", string(data)) // 上游原始响应(未转换)
|
||||
}
|
||||
g.finishUsage(c, ch, start, store.UsageStatusSuccess, "")
|
||||
}
|
||||
|
||||
@@ -283,10 +286,24 @@ func (g *Gateway) streamCopy(c *gin.Context, ch *store.Channel, r io.Reader, sta
|
||||
flusher = nopFlusher{}
|
||||
}
|
||||
|
||||
// 原始响应捕获:仅管理员且开关开启(raw_request 已 set)时累积上游原始行
|
||||
_, capture := c.Get("raw_request")
|
||||
var rawResp strings.Builder
|
||||
|
||||
// commitRaw 在记账前把已累积的原始响应写入 context
|
||||
commitRaw := func() {
|
||||
if capture {
|
||||
c.Set("raw_response", rawResp.String())
|
||||
}
|
||||
}
|
||||
|
||||
scanner := newSSEScanner(r)
|
||||
for {
|
||||
line, err := scanner.Next()
|
||||
if line != nil {
|
||||
if capture {
|
||||
rawResp.Write(line)
|
||||
}
|
||||
out := line
|
||||
if lineConv != nil {
|
||||
out = lineConv(line)
|
||||
@@ -294,6 +311,7 @@ func (g *Gateway) streamCopy(c *gin.Context, ch *store.Channel, r io.Reader, sta
|
||||
if out != nil {
|
||||
if _, werr := w.Write(out); werr != nil {
|
||||
// 客户端意外断开:按已生成部分收费(canceled)
|
||||
commitRaw()
|
||||
g.finishUsage(c, ch, start, store.UsageStatusCanceled, "client_disconnect")
|
||||
return
|
||||
}
|
||||
@@ -307,6 +325,7 @@ func (g *Gateway) streamCopy(c *gin.Context, ch *store.Channel, r io.Reader, sta
|
||||
}
|
||||
}
|
||||
if err != nil {
|
||||
commitRaw()
|
||||
if err == io.EOF {
|
||||
g.finishUsage(c, ch, start, store.UsageStatusSuccess, "")
|
||||
} else if c.Request.Context().Err() != nil {
|
||||
@@ -578,6 +597,14 @@ func (g *Gateway) finishUsage(c *gin.Context, ch *store.Channel, start time.Time
|
||||
chID = ch.ID
|
||||
}
|
||||
|
||||
var rawReq, rawResp string
|
||||
if v, ok := c.Get("raw_request"); ok {
|
||||
rawReq, _ = v.(string)
|
||||
}
|
||||
if v, ok := c.Get("raw_response"); ok {
|
||||
rawResp, _ = v.(string)
|
||||
}
|
||||
|
||||
// 密钥今日 token 用量累计(配额检查用)
|
||||
if g.lim != nil && kidVal > 0 {
|
||||
g.lim.AddTokens(kidVal, in+out)
|
||||
@@ -603,6 +630,8 @@ func (g *Gateway) finishUsage(c *gin.Context, ch *store.Channel, start time.Time
|
||||
LatencyMS: latency,
|
||||
Status: status,
|
||||
ErrorCode: errCodePtr,
|
||||
RawRequest: rawReq,
|
||||
RawResponse: rawResp,
|
||||
CreatedAt: time.Now().UTC(),
|
||||
})
|
||||
}
|
||||
|
||||
@@ -188,6 +188,8 @@ type UsageLog struct {
|
||||
LatencyMS int `json:"latency_ms"`
|
||||
Status string `gorm:"size:16;not null" json:"status"`
|
||||
ErrorCode *string `json:"error_code,omitempty"`
|
||||
RawRequest string `gorm:"type:text" json:"raw_request"` // 客户端原始请求体(未转换)
|
||||
RawResponse string `gorm:"type:text" json:"raw_response"` // 上游原始响应(未转换;流式为全部 SSE 事件)
|
||||
CreatedAt time.Time `gorm:"index" json:"created_at"`
|
||||
}
|
||||
|
||||
|
||||
Reference in New Issue
Block a user