feat: 实现完整可观测性架构与火山v3适配器重构
重构整体架构: 1. 新增telemetry包实现零依赖的Prometheus指标系统 2. 新增metrics包集中管理业务埋点指标 3. 重构火山v3适配器,拆分client/request/response等模块 4. 替换旧的service/stats统计系统为标准指标埋点 5. 新增/metrics观测端点与完整仪表盘支持 功能更新: - 实现基于IP的限流与并发限制,添加指标埋点 - 重构TTS控制器,支持多格式输出与完整错误分类 - 更新.env.example配置示例,新增多项可选参数 - 替换旧的volcano适配器实现,支持完整的v3 API特性 - 清理冗余代码,移除service/stats与旧adapter实现
This commit is contained in:
+105
-45
@@ -1,29 +1,34 @@
|
||||
package controller
|
||||
package controller
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io"
|
||||
"log"
|
||||
"net/http"
|
||||
"runtime"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/volcano-tts/tts-api/adapter/volcano"
|
||||
"github.com/volcano-tts/tts-api/common"
|
||||
"github.com/volcano-tts/tts-api/dto"
|
||||
"github.com/volcano-tts/tts-api/metrics"
|
||||
"github.com/volcano-tts/tts-api/middleware"
|
||||
"github.com/volcano-tts/tts-api/service"
|
||||
"github.com/volcano-tts/tts-api/setting"
|
||||
"github.com/volcano-tts/tts-api/telemetry"
|
||||
)
|
||||
|
||||
var volcanoClient *volcano.HTTPClient
|
||||
var (
|
||||
volcanoClient *volcano.HTTPClient
|
||||
adapterRec volcano.MetricsRecorder = metrics.AdapterRecorder{}
|
||||
)
|
||||
|
||||
func InitController() {
|
||||
volcanoClient = volcano.NewHTTPClient()
|
||||
}
|
||||
|
||||
// truncateForLog 用于在日志中安全地展示请求内容(截断避免日志爆炸、控制不可打印字符)
|
||||
func truncateForLog(b []byte, max int) string {
|
||||
if len(b) > max {
|
||||
return string(b[:max]) + fmt.Sprintf("...(truncated, total %d bytes)", len(b))
|
||||
@@ -31,15 +36,36 @@ func truncateForLog(b []byte, max int) string {
|
||||
return string(b)
|
||||
}
|
||||
|
||||
// resolveClientFormat 把 OpenAI 风格的 response_format 映射为最终输出格式;
|
||||
// 不识别或未指定时回退到 setting.TTSOptions.Format。
|
||||
func resolveClientFormat(reqFmt string) string {
|
||||
switch strings.ToLower(reqFmt) {
|
||||
case "mp3", "wav", "opus", "pcm", "aac", "flac":
|
||||
if reqFmt == "opus" {
|
||||
return "ogg_opus"
|
||||
}
|
||||
return strings.ToLower(reqFmt)
|
||||
}
|
||||
if reqFmt == "" {
|
||||
return setting.TTSOptions.Format
|
||||
}
|
||||
return setting.TTSOptions.Format
|
||||
}
|
||||
|
||||
// OpenaiTTSHandler 是 /v1/audio/speech 的入口。
|
||||
func OpenaiTTSHandler(w http.ResponseWriter, r *http.Request) {
|
||||
start := time.Now()
|
||||
|
||||
if r.Method != http.MethodPost {
|
||||
log.Printf("警告: 错误的方法 - 方法=%s 期望=POST 路径=%s 客户端=%s",
|
||||
r.Method, r.URL.Path, middleware.GetClientIP(r))
|
||||
metrics.RequestTotal.Inc(telemetry.Labels{"status": "method_not_allowed", "format": "", "speaker": "", "model": ""})
|
||||
http.Error(w, "Method not allowed", http.StatusMethodNotAllowed)
|
||||
return
|
||||
}
|
||||
|
||||
if !middleware.ValidateAPIKey(r) {
|
||||
metrics.AuthFailed.Inc(telemetry.Labels{})
|
||||
log.Printf("警告: API Key 鉴权失败 - 路径=%s 客户端=%s 远端=%s",
|
||||
r.URL.Path, middleware.GetClientIP(r), r.RemoteAddr)
|
||||
middleware.SendJSONError(w, http.StatusUnauthorized, "Invalid API key provided.", "invalid_request_error", "invalid_api_key")
|
||||
@@ -47,7 +73,7 @@ func OpenaiTTSHandler(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
if setting.TTSConfigErr != nil {
|
||||
log.Printf("警告: TTS配置未就绪,拒绝请求 - 错误=%v 路径=%s 客户端=%s",
|
||||
log.Printf("警告: TTS配置未就绪,拒绝请求 - 错误=%v 路径=%s 客户端=%s",
|
||||
setting.TTSConfigErr, r.URL.Path, middleware.GetClientIP(r))
|
||||
middleware.SendJSONError(w, http.StatusServiceUnavailable, "TTS service configuration error. Please check environment variables and restart the service.", "configuration_error", "service_unavailable")
|
||||
return
|
||||
@@ -96,7 +122,6 @@ func OpenaiTTSHandler(w http.ResponseWriter, r *http.Request) {
|
||||
http.Error(w, "Input text is required", http.StatusBadRequest)
|
||||
return
|
||||
}
|
||||
|
||||
if len(req.Input) > common.MaxTextLength {
|
||||
log.Printf("警告: input 文本过长 - 路径=%s 客户端=%s 长度=%d 限制=%d",
|
||||
r.URL.Path, middleware.GetClientIP(r), len(req.Input), common.MaxTextLength)
|
||||
@@ -115,27 +140,78 @@ func OpenaiTTSHandler(w http.ResponseWriter, r *http.Request) {
|
||||
speed = common.MaxSpeed
|
||||
}
|
||||
|
||||
ttsStart := time.Now()
|
||||
result, err := volcano.Synthesis(&setting.TTSConfig, volcanoClient, req.Input, speed)
|
||||
duration := time.Since(ttsStart)
|
||||
clientFormat := resolveClientFormat(req.ResponseFormat)
|
||||
|
||||
opts := setting.TTSOptions
|
||||
opts.Text = req.Input
|
||||
|
||||
ctx, cancel := context.WithTimeout(r.Context(), setting.TTSTimeout)
|
||||
defer cancel()
|
||||
|
||||
result, err := volcano.Synthesis(ctx, volcanoClient, opts, req.Input, clientFormat, speed, adapterRec)
|
||||
duration := time.Since(start)
|
||||
|
||||
finalLabels := telemetry.Labels{
|
||||
"format": clientFormat,
|
||||
"speaker": opts.Speaker,
|
||||
"model": opts.Model,
|
||||
}
|
||||
if err != nil {
|
||||
service.GlobalStats.AddRequest(false, duration, err.Error())
|
||||
finalLabels["status"] = classifyStatus(err)
|
||||
metrics.RequestTotal.Inc(finalLabels)
|
||||
metrics.RequestDuration.Observe(duration.Seconds(), telemetry.Labels{"status": finalLabels["status"], "format": clientFormat})
|
||||
log.Printf("警告: TTS 合成失败 - 路径=%s 客户端=%s 文本长度=%d 耗时=%v 错误=%v",
|
||||
r.URL.Path, middleware.GetClientIP(r), len(req.Input), duration, err)
|
||||
http.Error(w, "TTS synthesis failed", http.StatusInternalServerError)
|
||||
return
|
||||
}
|
||||
|
||||
service.GlobalStats.AddRequest(true, duration, "")
|
||||
finalLabels["status"] = "ok"
|
||||
metrics.RequestTotal.Inc(finalLabels)
|
||||
metrics.RequestDuration.Observe(duration.Seconds(), telemetry.Labels{"status": "ok", "format": clientFormat})
|
||||
|
||||
w.Header().Set("Content-Type", "audio/wav")
|
||||
w.Header().Set("Content-Type", contentTypeFor(result.Format))
|
||||
w.Header().Set("Content-Length", fmt.Sprintf("%d", len(result.AudioData)))
|
||||
w.Header().Set("X-Request-Id", result.ReqID)
|
||||
w.WriteHeader(http.StatusOK)
|
||||
w.Write(result.AudioData)
|
||||
}
|
||||
|
||||
func classifyStatus(err error) string {
|
||||
if ue, ok := err.(*volcano.UpstreamError); ok {
|
||||
switch ue.Stage {
|
||||
case "request":
|
||||
return "request_error"
|
||||
case "http":
|
||||
return fmt.Sprintf("http_%d", ue.Code)
|
||||
case "stream":
|
||||
return "upstream_error"
|
||||
case "wrap":
|
||||
return "wrap_error"
|
||||
}
|
||||
}
|
||||
return "internal_error"
|
||||
}
|
||||
|
||||
func contentTypeFor(format string) string {
|
||||
switch strings.ToLower(format) {
|
||||
case "wav":
|
||||
return "audio/wav"
|
||||
case "mp3":
|
||||
return "audio/mpeg"
|
||||
case "ogg_opus", "opus":
|
||||
return "audio/ogg"
|
||||
case "pcm":
|
||||
return "audio/L16"
|
||||
case "aac":
|
||||
return "audio/aac"
|
||||
case "flac":
|
||||
return "audio/flac"
|
||||
}
|
||||
return "application/octet-stream"
|
||||
}
|
||||
|
||||
// HealthHandler 暴露运行期状态;无鉴权。
|
||||
func HealthHandler(w http.ResponseWriter, r *http.Request) {
|
||||
w.Header().Set("Content-Type", "application/json")
|
||||
|
||||
@@ -145,55 +221,39 @@ func HealthHandler(w http.ResponseWriter, r *http.Request) {
|
||||
w.WriteHeader(http.StatusOK)
|
||||
}
|
||||
|
||||
totalRequests, successfulRequests, failedRequests, totalResponseTime, recentResponseTimes, lastErrors := service.GlobalStats.GetSnapshot()
|
||||
|
||||
var errorRate float64
|
||||
if totalRequests > 0 {
|
||||
errorRate = float64(failedRequests) / float64(totalRequests) * 100
|
||||
}
|
||||
|
||||
var avgResponseTime float64
|
||||
if totalRequests > 0 {
|
||||
avgResponseTime = totalResponseTime.Seconds() * 1000 / float64(totalRequests)
|
||||
}
|
||||
|
||||
envCheckStatus := setting.CheckEnvironmentVariables()
|
||||
allEnvVarsSet := envCheckStatus["all_required_vars_set"].(bool)
|
||||
env := setting.CheckEnvironmentVariables()
|
||||
allRequired := env["all_required_vars_set"].(bool)
|
||||
|
||||
status := "ok"
|
||||
if !allEnvVarsSet {
|
||||
if !allRequired {
|
||||
status = "configuration_error"
|
||||
}
|
||||
|
||||
response := dto.HealthResponse{
|
||||
resp := dto.HealthResponse{
|
||||
Status: status,
|
||||
Service: "ByteDance TTS to OpenAI API Adapter",
|
||||
Version: "2.0.0 (v3 API)",
|
||||
Uptime: fmt.Sprintf("%.0f seconds", time.Since(startTime).Seconds()),
|
||||
StartTime: startTime.Format(time.RFC3339),
|
||||
Memory: service.GetMemoryInfo(),
|
||||
APIStats: dto.APIStatsResponse{
|
||||
TotalRequests: int(totalRequests),
|
||||
SuccessfulRequests: successfulRequests,
|
||||
FailedRequests: failedRequests,
|
||||
ErrorRatePercent: fmt.Sprintf("%.2f", errorRate),
|
||||
AvgResponseTimeMs: fmt.Sprintf("%.2f", avgResponseTime),
|
||||
RecentResponseTimesMs: recentResponseTimes,
|
||||
},
|
||||
Errors: dto.ErrorResponse{
|
||||
RecentErrorsCount: len(lastErrors),
|
||||
},
|
||||
Memory: collectMemorySnapshot(),
|
||||
ConfigStatus: dto.ConfigStatusResponse{
|
||||
AllRequiredVarsSet: allEnvVarsSet,
|
||||
AllRequiredVarsSet: allRequired,
|
||||
ConfigError: setting.TTSConfigErr != nil,
|
||||
},
|
||||
}
|
||||
|
||||
json.NewEncoder(w).Encode(response)
|
||||
json.NewEncoder(w).Encode(resp)
|
||||
}
|
||||
|
||||
var startTime time.Time
|
||||
|
||||
func SetStartTime(t time.Time) {
|
||||
startTime = t
|
||||
func SetStartTime(t time.Time) { startTime = t }
|
||||
|
||||
func collectMemorySnapshot() map[string]interface{} {
|
||||
var ms runtime.MemStats
|
||||
runtime.ReadMemStats(&ms)
|
||||
return map[string]interface{}{
|
||||
"heap_alloc": ms.HeapAlloc,
|
||||
"heap_inuse": ms.HeapInuse,
|
||||
"goroutines": runtime.NumGoroutine(),
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user