diff --git a/backend/go-service/app/dto/transcribe_dto.go b/backend/go-service/app/dto/transcribe_dto.go new file mode 100644 index 0000000..44a17d5 --- /dev/null +++ b/backend/go-service/app/dto/transcribe_dto.go @@ -0,0 +1,53 @@ +// Package dto 包含 transcribe 模块的请求/响应数据结构 +package dto + +import "time" + +// AdminTranscribeRequest POST /admin/recordings/:id/transcribe 请求体 +type AdminTranscribeRequest struct { + Force bool `json:"force"` // true=已 ready 时强制重新转写 + Language string `json:"language"` // 可选 ISO-639 语言提示,"" 表示自动检测 +} + +// AdminTranscriptDTO 转写结果出参(去掉/格式化数据库字段,便于前端消费) +type AdminTranscriptDTO struct { + ID int64 `json:"id"` + RecordingID int64 `json:"recording_id"` + RoomID int64 `json:"room_id"` + Status string `json:"status"` // pending/running/ready/failed + StatusLabel string `json:"status_label"` // 中文标签 + Text string `json:"text"` + Segments []AdminTranscriptSeg `json:"segments"` // 解析后的时间轴 + Language string `json:"language"` + DurationSec int `json:"duration_sec"` + ProviderCode string `json:"provider_code"` + ModelCode string `json:"model_code"` + ErrorMsg string `json:"error_msg"` + StartedAt string `json:"started_at"` // YYYY-MM-DD HH:MM:SS,空则 "" + FinishedAt string `json:"finished_at"` + CreatedAt string `json:"created_at"` + UpdatedAt string `json:"updated_at"` +} + +// AdminTranscriptSeg 转写的单段时间轴 +type AdminTranscriptSeg struct { + Start float64 `json:"start"` + End float64 `json:"end"` + Text string `json:"text"` +} + +// FormatTime 工具:time.Time 安全格式化(零值返回 "") +func FormatTranscriptTime(t time.Time) string { + if t.IsZero() { + return "" + } + return t.Format("2006-01-02 15:04:05") +} + +// FormatTranscriptTimePtr 工具:*time.Time 安全格式化 +func FormatTranscriptTimePtr(t *time.Time) string { + if t == nil || t.IsZero() { + return "" + } + return t.Format("2006-01-02 15:04:05") +} diff --git a/backend/go-service/app/provider/provider.go b/backend/go-service/app/provider/provider.go index d39629c..b7945a0 100644 --- a/backend/go-service/app/provider/provider.go +++ b/backend/go-service/app/provider/provider.go @@ -17,6 +17,7 @@ import ( notifyController "github.com/echochat/backend/app/notify/controller" notifyService "github.com/echochat/backend/app/notify/service" notifyTask "github.com/echochat/backend/app/notify/task" + transcribeController "github.com/echochat/backend/app/transcribe/controller" wsApp "github.com/echochat/backend/app/ws" "github.com/echochat/backend/config" "github.com/echochat/backend/pkg/db" @@ -64,6 +65,8 @@ type App struct { MeetingRecordingService *meetingService.MeetingRecordingService // 会议录制业务服务(Phase B 新增) MeetingWSHandler *meetingController.MeetingWSHandler // 会议 WS 事件 Handler(构造时自动注册路由到 Hub) MeetingCleanupTask *meetingTask.MeetingCleanupTask // 会议生命周期兜底定时任务(Task 8 落地) + TranscribeController *transcribeController.TranscribeController // 语音转写控制器(Phase B 新增) + LLMSourceDB *db.LLMSourceDB // 外部 LLM 配置 MySQL 连接句柄(Phase B 新增) } // NewApp 创建应用实例 @@ -102,6 +105,8 @@ func NewApp( meetingRecordingSvc *meetingService.MeetingRecordingService, meetingWSHandler *meetingController.MeetingWSHandler, meetingCleanup *meetingTask.MeetingCleanupTask, + transcribeCtrl *transcribeController.TranscribeController, + llmSourceDB *db.LLMSourceDB, ) *App { wsHandler.SetOfflinePusher(offlinePusher) wsHandler.SetNotifyConnectHook(notifySvc) @@ -144,6 +149,8 @@ func NewApp( MeetingRecordingService: meetingRecordingSvc, MeetingWSHandler: meetingWSHandler, MeetingCleanupTask: meetingCleanup, + TranscribeController: transcribeCtrl, + LLMSourceDB: llmSourceDB, } } @@ -173,6 +180,11 @@ func provideServerConfig(cfg *config.Config) *config.ServerConfig { return &cfg.Server } +// provideLLMSourceConfig 从全局 Config 中提取 LLMSourceConfig(Phase B 语音转写) +func provideLLMSourceConfig(cfg *config.Config) *config.LLMSourceConfig { + return &cfg.LLMSource +} + // InfraSet 基础设施层 Provider Set var InfraSet = wire.NewSet( provideDBConfig, @@ -180,8 +192,10 @@ var InfraSet = wire.NewSet( provideJWTConfig, provideMinioConfig, provideServerConfig, + provideLLMSourceConfig, db.NewPostgres, db.NewRedis, + db.NewLLMSourceDB, storage.NewMinioClient, NewApp, ) diff --git a/backend/go-service/app/provider/wire.go b/backend/go-service/app/provider/wire.go index cbe3673..a1f2804 100644 --- a/backend/go-service/app/provider/wire.go +++ b/backend/go-service/app/provider/wire.go @@ -21,6 +21,7 @@ import ( meetingService "github.com/echochat/backend/app/meeting/service" notifyApp "github.com/echochat/backend/app/notify" notifyService "github.com/echochat/backend/app/notify/service" + transcribeApp "github.com/echochat/backend/app/transcribe" wsApp "github.com/echochat/backend/app/ws" "github.com/echochat/backend/config" "github.com/google/wire" @@ -39,6 +40,7 @@ func InitializeApp(cfg *config.Config) (*App, error) { groupApp.GroupSet, notifyApp.NotifySet, meetingApp.MeetingSet, + transcribeApp.TranscribeSet, wire.Bind(new(wsApp.FriendIDsGetter), new(*contactDAO.FriendshipDAO)), wire.Bind(new(groupService.UserInfoProvider), new(*contactDAO.FriendshipDAO)), wire.Bind(new(imService.GroupInfoGetter), new(*groupDAO.GroupDAO)), diff --git a/backend/go-service/app/provider/wire_gen.go b/backend/go-service/app/provider/wire_gen.go index 2be9e33..f3cf652 100644 --- a/backend/go-service/app/provider/wire_gen.go +++ b/backend/go-service/app/provider/wire_gen.go @@ -30,6 +30,9 @@ import ( dao5 "github.com/echochat/backend/app/notify/dao" service3 "github.com/echochat/backend/app/notify/service" "github.com/echochat/backend/app/notify/task" + controller8 "github.com/echochat/backend/app/transcribe/controller" + dao7 "github.com/echochat/backend/app/transcribe/dao" + service8 "github.com/echochat/backend/app/transcribe/service" "github.com/echochat/backend/app/ws" "github.com/echochat/backend/config" "github.com/echochat/backend/pkg/db" @@ -117,6 +120,16 @@ func InitializeApp(cfg *config.Config) (*App, error) { meetingRecordingController := controller7.NewMeetingRecordingController(meetingRecordingService) meetingWSHandler := controller7.NewMeetingWSHandler(meetingSignalService, hub) meetingCleanupTask := task2.NewMeetingCleanupTask(meetingLifecycleService, meetingRoomDAO, meetingChatDAO) - app := NewApp(cfg, gormDB, client, minioClient, authService, authController, adminAuthController, userManageController, onlineController, contactManageController, groupManageController, messageManageController, meetingManageController, handler, hub, pubSub, onlineService, contactController, imController, eventHandler, offlinePusher, fileController, groupController, notifyService, notificationController, cleanupTask, meetingService, meetingSignalService, meetingLifecycleService, meetingController, meetingRecordingController, meetingRecordingService, meetingWSHandler, meetingCleanupTask) + llmSourceConfig := provideLLMSourceConfig(cfg) + llmSourceDB, err := db.NewLLMSourceDB(llmSourceConfig) + if err != nil { + return nil, err + } + transcriptDAO := dao7.NewTranscriptDAO(gormDB) + llmConfigDAO := dao7.NewLLMConfigDAO(llmSourceDB) + sttClient := service8.NewOpenAICompatibleClient() + transcribeService := service8.NewTranscribeService(gormDB, transcriptDAO, llmConfigDAO, sttClient) + transcribeController := controller8.NewTranscribeController(transcribeService) + app := NewApp(cfg, gormDB, client, minioClient, authService, authController, adminAuthController, userManageController, onlineController, contactManageController, groupManageController, messageManageController, meetingManageController, handler, hub, pubSub, onlineService, contactController, imController, eventHandler, offlinePusher, fileController, groupController, notifyService, notificationController, cleanupTask, meetingService, meetingSignalService, meetingLifecycleService, meetingController, meetingRecordingController, meetingRecordingService, meetingWSHandler, meetingCleanupTask, transcribeController, llmSourceDB) return app, nil } diff --git a/backend/go-service/app/transcribe/controller/transcribe_controller.go b/backend/go-service/app/transcribe/controller/transcribe_controller.go new file mode 100644 index 0000000..f108754 --- /dev/null +++ b/backend/go-service/app/transcribe/controller/transcribe_controller.go @@ -0,0 +1,142 @@ +// Package controller 提供 transcribe 模块的 admin HTTP 接口 +package controller + +import ( + "encoding/json" + "errors" + "net/http" + "strconv" + + "github.com/echochat/backend/app/dto" + "github.com/echochat/backend/app/transcribe/dao" + "github.com/echochat/backend/app/transcribe/model" + "github.com/echochat/backend/app/transcribe/service" + "github.com/echochat/backend/pkg/utils" + "github.com/gin-gonic/gin" +) + +// TranscribeController admin 端转写控制器 +type TranscribeController struct { + svc *service.TranscribeService +} + +// NewTranscribeController 创建实例 +func NewTranscribeController(svc *service.TranscribeService) *TranscribeController { + return &TranscribeController{svc: svc} +} + +// GetTranscript GET /admin/recordings/:id/transcript +// +// 返回值: +// - 200 + transcript:存在时返回最新状态 +// - 200 + null:尚未发起转写(前端据此显示"开始转写"按钮) +func (ctl *TranscribeController) GetTranscript(c *gin.Context) { + id, err := strconv.ParseInt(c.Param("id"), 10, 64) + if err != nil || id <= 0 { + utils.ResponseBadRequest(c, "无效的录制 ID") + return + } + + t, err := ctl.svc.GetByRecording(c.Request.Context(), id) + if err != nil { + utils.ResponseError(c, "查询转写失败") + return + } + utils.ResponseOK(c, toDTO(t)) +} + +// SubmitTranscribe POST /admin/recordings/:id/transcribe +// +// body: { force: bool, language: string } +// +// 行为:异步启动 STT,返回 running 状态快照 +func (ctl *TranscribeController) SubmitTranscribe(c *gin.Context) { + id, err := strconv.ParseInt(c.Param("id"), 10, 64) + if err != nil || id <= 0 { + utils.ResponseBadRequest(c, "无效的录制 ID") + return + } + + var req dto.AdminTranscribeRequest + // body 可选;POST 不带 body 也允许(默认 force=false language="") + _ = c.ShouldBindJSON(&req) + + t, err := ctl.svc.Submit(c.Request.Context(), id, req.Force, req.Language) + if err != nil { + ctl.handleError(c, err, t) + return + } + utils.ResponseOK(c, toDTO(t)) +} + +// handleError 统一业务错误映射 +func (ctl *TranscribeController) handleError(c *gin.Context, err error, t *model.Transcript) { + switch { + case errors.Is(err, service.ErrRecordingNotFound): + utils.ResponseNotFound(c, err.Error()) + case errors.Is(err, service.ErrRecordingNotReady): + utils.ResponseBadRequest(c, err.Error()) + case errors.Is(err, service.ErrSTTNotConfigured), errors.Is(err, dao.ErrLLMSourceDisabled): + // 503 表示服务暂不可用,前端可据此提示"请联系管理员配置 STT" + c.JSON(http.StatusServiceUnavailable, gin.H{ + "code": 503, + "message": "语音转写服务未配置或不可用", + "data": nil, + }) + case errors.Is(err, service.ErrTranscribeRunning): + // 已在跑:返回当前进行中行,前端可继续轮询 + c.JSON(http.StatusAccepted, gin.H{ + "code": 0, + "message": "已有转写任务进行中", + "data": toDTO(t), + }) + case errors.Is(err, dao.ErrNoActiveSTT): + utils.ResponseBadRequest(c, "未找到可用的 STT 模型或 API Key,请检查 LLM 配置中心") + default: + utils.ResponseError(c, "发起转写失败: "+err.Error()) + } +} + +// toDTO 模型转 DTO,segments JSON 解码 + 状态中文标签 +func toDTO(t *model.Transcript) *dto.AdminTranscriptDTO { + if t == nil { + return nil + } + segs := []dto.AdminTranscriptSeg{} + if t.Segments != "" { + _ = json.Unmarshal([]byte(t.Segments), &segs) + } + return &dto.AdminTranscriptDTO{ + ID: t.ID, + RecordingID: t.RecordingID, + RoomID: t.RoomID, + Status: t.Status, + StatusLabel: statusLabel(t.Status), + Text: t.Text, + Segments: segs, + Language: t.Language, + DurationSec: t.DurationSec, + ProviderCode: t.ProviderCode, + ModelCode: t.ModelCode, + ErrorMsg: t.ErrorMsg, + StartedAt: dto.FormatTranscriptTimePtr(t.StartedAt), + FinishedAt: dto.FormatTranscriptTimePtr(t.FinishedAt), + CreatedAt: dto.FormatTranscriptTime(t.CreatedAt), + UpdatedAt: dto.FormatTranscriptTime(t.UpdatedAt), + } +} + +func statusLabel(s string) string { + switch s { + case model.TranscriptStatusPending: + return "待处理" + case model.TranscriptStatusRunning: + return "转写中" + case model.TranscriptStatusReady: + return "已就绪" + case model.TranscriptStatusFailed: + return "失败" + default: + return s + } +} diff --git a/backend/go-service/app/transcribe/dao/llm_config_dao.go b/backend/go-service/app/transcribe/dao/llm_config_dao.go new file mode 100644 index 0000000..be42064 --- /dev/null +++ b/backend/go-service/app/transcribe/dao/llm_config_dao.go @@ -0,0 +1,106 @@ +// Package dao 提供 transcribe 模块对外部 LLM 配置中心的只读访问 +package dao + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/echochat/backend/app/transcribe/model/llm" + "github.com/echochat/backend/pkg/db" +) + +// ErrLLMSourceDisabled LLM 配置库未启用 +var ErrLLMSourceDisabled = errors.New("LLM 配置中心未启用 (llm_source.enabled=false)") + +// ErrNoActiveSTT 没有可用的 STT 模型/Key +var ErrNoActiveSTT = errors.New("未找到可用的 STT 模型或 API Key") + +// LLMConfigDAO 外部 LLM 配置只读 DAO(MySQL) +// +// 选择策略说明(PickActiveSTT): +// 1. 在 t_llm_model 中找 model_type=4 (ASR) 且 status=1 deleted=0 ORDER BY sort ASC 的第一条 +// 2. 用其 provider_id 在 t_llm_provider 中找对应 provider,要求 status=1 deleted=0 +// 3. 在 t_llm_key 中按 provider_id 选 status=1 deleted=0 且 (expire_time IS NULL OR expire_time > NOW()) +// 且 (daily_limit=0 OR today_count < daily_limit) ORDER BY weight DESC, today_count ASC LIMIT 1 +// +// 任何一步落空都返回 ErrNoActiveSTT,避免上游用半成品配置发起调用。 +type LLMConfigDAO struct { + source *db.LLMSourceDB +} + +// NewLLMConfigDAO 创建实例 +func NewLLMConfigDAO(source *db.LLMSourceDB) *LLMConfigDAO { + return &LLMConfigDAO{source: source} +} + +// IsEnabled 暴露给上层做"是否能转写"的快速判断,避免每次都 catch error +func (d *LLMConfigDAO) IsEnabled() bool { + return d != nil && d.source != nil && d.source.IsEnabled() +} + +// STTConfig 三表合一的运行时配置快照,用于本次转写调用 +type STTConfig struct { + Provider llm.Provider + Model llm.Model + Key llm.Key +} + +// PickActiveSTT 按选择策略挑出一组可用的 (provider, model, key) +// +// 注意:不做事务,因为这是只读 + LLM 配置中心通常更新频率极低, +// 偶发的"挑出后 key 立刻被禁用"由调用层捕获 401/403 后重试解决。 +func (d *LLMConfigDAO) PickActiveSTT(ctx context.Context) (*STTConfig, error) { + if !d.IsEnabled() { + return nil, ErrLLMSourceDisabled + } + gdb := d.source.DB() + + // 1) 选 ASR 模型 + var model llm.Model + err := gdb.WithContext(ctx). + Where("deleted = 0 AND status = ? AND model_type = ?", llm.StatusEnabled, llm.ModelTypeASR). + Order("sort ASC, id ASC"). + First(&model).Error + if err != nil { + return nil, fmt.Errorf("%w: %v", ErrNoActiveSTT, err) + } + + // 2) 模型对应的 provider 必须也启用 + var provider llm.Provider + err = gdb.WithContext(ctx). + Where("id = ? AND deleted = 0 AND status = ?", model.ProviderID, llm.StatusEnabled). + First(&provider).Error + if err != nil { + return nil, fmt.Errorf("%w: provider not active for model %s", ErrNoActiveSTT, model.ModelCode) + } + + // 3) 选 key + now := time.Now() + var key llm.Key + err = gdb.WithContext(ctx). + Where("deleted = 0 AND status = ? AND provider_id = ?", llm.StatusEnabled, provider.ID). + Where("expire_time IS NULL OR expire_time > ?", now). + Where("daily_limit = 0 OR today_count < daily_limit"). + Order("weight DESC, today_count ASC, id ASC"). + First(&key).Error + if err != nil { + return nil, fmt.Errorf("%w: no usable key for provider %s", ErrNoActiveSTT, provider.ProviderCode) + } + + return &STTConfig{Provider: provider, Model: model, Key: key}, nil +} + +// ListASRModels 列出全部启用的 ASR 模型(admin 可视化用,可不接 UI,先备好接口) +func (d *LLMConfigDAO) ListASRModels(ctx context.Context) ([]llm.Model, error) { + if !d.IsEnabled() { + return nil, ErrLLMSourceDisabled + } + var list []llm.Model + err := d.source.DB().WithContext(ctx). + Where("deleted = 0 AND status = ? AND model_type = ?", llm.StatusEnabled, llm.ModelTypeASR). + Order("sort ASC, id ASC"). + Find(&list).Error + return list, err +} diff --git a/backend/go-service/app/transcribe/dao/transcript_dao.go b/backend/go-service/app/transcribe/dao/transcript_dao.go new file mode 100644 index 0000000..8a4688d --- /dev/null +++ b/backend/go-service/app/transcribe/dao/transcript_dao.go @@ -0,0 +1,112 @@ +// Package dao 提供 transcribe 模块的数据库访问操作 +package dao + +import ( + "context" + "errors" + "time" + + "github.com/echochat/backend/app/transcribe/model" + "gorm.io/gorm" +) + +// TranscriptDAO 转写记录数据访问对象(PostgreSQL 主库) +type TranscriptDAO struct { + db *gorm.DB +} + +// NewTranscriptDAO 创建实例 +func NewTranscriptDAO(db *gorm.DB) *TranscriptDAO { + return &TranscriptDAO{db: db} +} + +// GetByRecordingID 按 recording_id 获取转写记录(不存在返回 nil, nil) +// +// 之所以 nil 不视为错误:调用方常见模式是"取不到 → 创建新行",避免 ErrRecordNotFound 包装 +func (d *TranscriptDAO) GetByRecordingID(ctx context.Context, recordingID int64) (*model.Transcript, error) { + var t model.Transcript + err := d.db.WithContext(ctx). + Where("recording_id = ?", recordingID). + First(&t).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, nil + } + return nil, err + } + return &t, nil +} + +// ListByRoomID 拉取一场会议下所有转写(按 recording_id 升序) +// 用于会议详情页一次性带回多段录制的转写状态 +func (d *TranscriptDAO) ListByRoomID(ctx context.Context, roomID int64) ([]model.Transcript, error) { + var list []model.Transcript + err := d.db.WithContext(ctx). + Where("room_id = ?", roomID). + Order("recording_id ASC"). + Find(&list).Error + return list, err +} + +// Upsert 插入或更新(按 recording_id unique 索引) +// 用于"重新转写"路径:保留同一行 ID,状态/文本就地刷新 +func (d *TranscriptDAO) Upsert(ctx context.Context, t *model.Transcript) error { + existing, err := d.GetByRecordingID(ctx, t.RecordingID) + if err != nil { + return err + } + if existing == nil { + return d.db.WithContext(ctx).Create(t).Error + } + t.ID = existing.ID + return d.db.WithContext(ctx).Save(t).Error +} + +// MarkRunning 把指定 transcript 置为 running 状态,并填 started_at = now +// +// 用于异步任务起跑时的状态翻转,保证后端重启时 running 行不会卡死 +// (配合 RescueStuckRunning 兜底)。 +func (d *TranscriptDAO) MarkRunning(ctx context.Context, id int64) error { + now := time.Now() + return d.db.WithContext(ctx). + Model(&model.Transcript{}). + Where("id = ?", id). + Updates(map[string]any{ + "status": model.TranscriptStatusRunning, + "started_at": now, + "error_msg": "", + }).Error +} + +// MarkReady 转写成功 +func (d *TranscriptDAO) MarkReady(ctx context.Context, id int64, text, segments, language string, durationSec int) error { + now := time.Now() + return d.db.WithContext(ctx). + Model(&model.Transcript{}). + Where("id = ?", id). + Updates(map[string]any{ + "status": model.TranscriptStatusReady, + "text": text, + "segments": segments, + "language": language, + "duration_sec": durationSec, + "finished_at": now, + "error_msg": "", + }).Error +} + +// MarkFailed 转写失败 +func (d *TranscriptDAO) MarkFailed(ctx context.Context, id int64, errorMsg string) error { + now := time.Now() + if len(errorMsg) > 500 { + errorMsg = errorMsg[:500] + } + return d.db.WithContext(ctx). + Model(&model.Transcript{}). + Where("id = ?", id). + Updates(map[string]any{ + "status": model.TranscriptStatusFailed, + "error_msg": errorMsg, + "finished_at": now, + }).Error +} diff --git a/backend/go-service/app/transcribe/model/llm/llm.go b/backend/go-service/app/transcribe/model/llm/llm.go new file mode 100644 index 0000000..636b783 --- /dev/null +++ b/backend/go-service/app/transcribe/model/llm/llm.go @@ -0,0 +1,103 @@ +// Package llm 映射外部 LLM 配置中心 MySQL 的三张表 +// +// 严格只读:本服务永远不应该 INSERT / UPDATE / DELETE 这些表, +// DAO 中只用 Find / First / Where。所有写操作由配置中心自身负责。 +package llm + +import "time" + +// 模型类型枚举(对应 t_llm_model.model_type) +const ( + ModelTypeChat = 1 // 文字对话 + ModelTypeVision = 2 // 视觉(图文) + ModelTypeMultimodal = 3 // 多模态 + ModelTypeASR = 4 // 语音识别(本服务转写功能用此类型) +) + +// 通用状态枚举 +const ( + StatusDisabled = 0 // 停用 + StatusEnabled = 1 // 启用 + // t_llm_key 额外定义 + KeyStatusAutoDisabled = 2 // 连续失败自动停用 +) + +// API 协议 +const ( + APIProtocolOpenAICompatible = "openai_compatible" + APIProtocolCustom = "custom" +) + +// Provider 大模型提供商表 t_llm_provider(只读) +// +// 选择 provider 时使用:deleted=0 AND status=1 ORDER BY sort ASC +type Provider struct { + ID int64 `gorm:"column:id;primaryKey"` + ProviderName string `gorm:"column:provider_name"` + ProviderCode string `gorm:"column:provider_code"` + BaseURL string `gorm:"column:base_url"` + APIProtocol string `gorm:"column:api_protocol"` // openai_compatible / custom + Status int `gorm:"column:status"` + Sort int `gorm:"column:sort"` + Remark string `gorm:"column:remark"` + CreateTime time.Time `gorm:"column:create_time"` + UpdateTime time.Time `gorm:"column:update_time"` + Deleted int `gorm:"column:deleted"` +} + +// TableName 指定 MySQL 表名 +func (Provider) TableName() string { return "t_llm_provider" } + +// Model 大模型表 t_llm_model(只读) +// +// 选择 STT 模型时使用:deleted=0 AND status=1 AND model_type=4 ORDER BY sort ASC +type Model struct { + ID int64 `gorm:"column:id;primaryKey"` + ProviderID int64 `gorm:"column:provider_id"` + ModelCode string `gorm:"column:model_code"` + ModelName string `gorm:"column:model_name"` + ModelType int `gorm:"column:model_type"` + MaxTokens int `gorm:"column:max_tokens"` + InputPrice float64 `gorm:"column:input_price"` + OutputPrice float64 `gorm:"column:output_price"` + Status int `gorm:"column:status"` + Sort int `gorm:"column:sort"` + Remark string `gorm:"column:remark"` + CreateTime time.Time `gorm:"column:create_time"` + UpdateTime time.Time `gorm:"column:update_time"` + Deleted int `gorm:"column:deleted"` +} + +func (Model) TableName() string { return "t_llm_model" } + +// Key 大模型 Key 池表 t_llm_key(只读) +// +// 选择策略(DAO.PickKey): +// - WHERE deleted=0 AND status=1 AND provider_id=? +// - AND (expire_time IS NULL OR expire_time > NOW()) +// - AND (daily_limit = 0 OR today_count < daily_limit) +// - ORDER BY weight DESC, today_count ASC +// - LIMIT 1 +// +// 注意:本服务不写回 today_count / total_count / fail_count, +// 因为这些字段属于 LLM 网关自身的统计职责,由网关在调用时自增。 +// 我们这边仅消费它们做选择权重。 +type Key struct { + ID int64 `gorm:"column:id;primaryKey"` + ProviderID int64 `gorm:"column:provider_id"` + APIKey string `gorm:"column:api_key"` + KeyAlias string `gorm:"column:key_alias"` + Weight int `gorm:"column:weight"` + DailyLimit int `gorm:"column:daily_limit"` + TodayCount int `gorm:"column:today_count"` + TotalCount int64 `gorm:"column:total_count"` + FailCount int `gorm:"column:fail_count"` + Status int `gorm:"column:status"` // 0/1/2 + LastUsedTime *time.Time `gorm:"column:last_used_time"` + ExpireTime *time.Time `gorm:"column:expire_time"` + CreateTime time.Time `gorm:"column:create_time"` + UpdateTime time.Time `gorm:"column:update_time"` + Deleted int `gorm:"column:deleted"` +} + +func (Key) TableName() string { return "t_llm_key" } diff --git a/backend/go-service/app/transcribe/model/transcript.go b/backend/go-service/app/transcribe/model/transcript.go new file mode 100644 index 0000000..fba8d11 --- /dev/null +++ b/backend/go-service/app/transcribe/model/transcript.go @@ -0,0 +1,48 @@ +// Package model 提供 transcribe 模块的数据库模型 +package model + +import "time" + +// Transcript 会议录制语音转写结果,对应 meeting_transcripts 表 +// +// 一对一关联 meeting_recordings.id(unique 索引),同一段录制只保留最新一份转写结果, +// 重新触发转写会就地更新(保持状态机:pending → running → ready/failed)。 +// +// 字段设计: +// - Text:拼接后的全文,便于列表/搜索 +// - Segments:JSONB 数组 [{start,end,text}],便于按时间轴渲染字幕 +// - Provider/ModelCode:本次实际使用的 LLM 配置快照,调试/审计用 +// - ErrorMsg:失败原因(HTTP 错误 / Key 不可用 / 文件下载失败等) +// +// PostgreSQL JSONB 字段在 GORM 里用 string 承载,由 service 层 json.Marshal/Unmarshal。 +type Transcript struct { + ID int64 `json:"id" gorm:"primaryKey;autoIncrement"` + RecordingID int64 `json:"recording_id" gorm:"not null;uniqueIndex:uk_transcripts_recording"` // 唯一关联录制 + RoomID int64 `json:"room_id" gorm:"not null;index:idx_transcripts_room"` // 冗余 room_id 便于按会议批量查 + Status string `json:"status" gorm:"size:16;not null;default:pending"` // pending / running / ready / failed + Text string `json:"text" gorm:"type:text;not null;default:''"` + Segments string `json:"segments" gorm:"type:jsonb;not null;default:'[]'"` // JSON 数组字符串 + Language string `json:"language" gorm:"size:16;not null;default:''"` + DurationSec int `json:"duration_sec" gorm:"not null;default:0"` + ProviderCode string `json:"provider_code" gorm:"size:32;not null;default:''"` // 来自 t_llm_provider.provider_code + ModelCode string `json:"model_code" gorm:"size:64;not null;default:''"` // 来自 t_llm_model.model_code + KeyID int64 `json:"key_id" gorm:"not null;default:0"` // 来自 t_llm_key.id,便于追踪 key 用量 + ErrorMsg string `json:"error_msg" gorm:"size:512;not null;default:''"` + StartedAt *time.Time `json:"started_at" gorm:"type:timestamp(0)"` + FinishedAt *time.Time `json:"finished_at" gorm:"type:timestamp(0)"` + CreatedAt time.Time `json:"created_at" gorm:"not null;autoCreateTime;type:timestamp(0)"` + UpdatedAt time.Time `json:"updated_at" gorm:"not null;autoUpdateTime;type:timestamp(0)"` +} + +// TableName 指定数据库表名 +func (Transcript) TableName() string { + return "meeting_transcripts" +} + +// 状态常量 +const ( + TranscriptStatusPending = "pending" + TranscriptStatusRunning = "running" + TranscriptStatusReady = "ready" + TranscriptStatusFailed = "failed" +) diff --git a/backend/go-service/app/transcribe/provider.go b/backend/go-service/app/transcribe/provider.go new file mode 100644 index 0000000..1a05119 --- /dev/null +++ b/backend/go-service/app/transcribe/provider.go @@ -0,0 +1,18 @@ +// Package transcribe 提供 transcribe 模块的 Wire Provider 集合 +package transcribe + +import ( + "github.com/echochat/backend/app/transcribe/controller" + "github.com/echochat/backend/app/transcribe/dao" + "github.com/echochat/backend/app/transcribe/service" + "github.com/google/wire" +) + +// TranscribeSet 转写模块 Wire Provider 集合 +var TranscribeSet = wire.NewSet( + dao.NewTranscriptDAO, + dao.NewLLMConfigDAO, + service.NewOpenAICompatibleClient, + service.NewTranscribeService, + controller.NewTranscribeController, +) diff --git a/backend/go-service/app/transcribe/router.go b/backend/go-service/app/transcribe/router.go new file mode 100644 index 0000000..943c3e2 --- /dev/null +++ b/backend/go-service/app/transcribe/router.go @@ -0,0 +1,25 @@ +// Package transcribe 提供 transcribe 模块的路由注册 +package transcribe + +import ( + "github.com/echochat/backend/app/constants" + "github.com/echochat/backend/app/transcribe/controller" + "github.com/echochat/backend/pkg/middleware" + "github.com/gin-gonic/gin" +) + +// RegisterRoutes 注册 admin 端的转写路由 +// +// 全部走 JWT + admin 角色双重中间件,与 admin 模块同等权限。 +// +// 路由: +// - POST /api/v1/admin/recordings/:id/transcribe 提交转写任务 +// - GET /api/v1/admin/recordings/:id/transcript 查询转写结果 +func RegisterRoutes(r *gin.Engine, ctrl *controller.TranscribeController, jwtAuth gin.HandlerFunc) { + g := r.Group("/api/v1/admin") + g.Use(jwtAuth, middleware.RequireRole(constants.RoleAdmin, constants.RoleSuperAdmin)) + { + g.POST("/recordings/:id/transcribe", ctrl.SubmitTranscribe) + g.GET("/recordings/:id/transcript", ctrl.GetTranscript) + } +} diff --git a/backend/go-service/app/transcribe/service/openai_compat_client.go b/backend/go-service/app/transcribe/service/openai_compat_client.go new file mode 100644 index 0000000..eb11601 --- /dev/null +++ b/backend/go-service/app/transcribe/service/openai_compat_client.go @@ -0,0 +1,181 @@ +// Package service 提供 transcribe 模块的业务服务 +package service + +import ( + "bytes" + "context" + "encoding/json" + "fmt" + "io" + "mime/multipart" + "net/http" + "path" + "strings" + "time" + + "github.com/echochat/backend/app/transcribe/dao" +) + +// OpenAICompatibleClient 实现 OpenAI 兼容协议的 STT 调用 +// +// 标准端点:POST {base_url}/audio/transcriptions +// 形式:multipart/form-data +// - file: 音频文件(mp3/wav/m4a/mp4/...) +// - model: 模型 code(如 whisper-1 / qwen-audio-asr-flash) +// - response_format: verbose_json(拿 segments)/ json(仅 text) +// - language: 可选,提示语言以加速识别 +// +// 阿里云 DashScope / DeepSeek / OpenAI / 智谱(部分)/ 通义千问 都遵循该协议。 +// +// 调用超时:默认 5 分钟(覆盖大部分 ≤30min 会议)。超时由调用方通过 ctx 控制更精细。 +type OpenAICompatibleClient struct { + httpClient *http.Client +} + +// NewOpenAICompatibleClient 创建实例 +func NewOpenAICompatibleClient() *OpenAICompatibleClient { + return &OpenAICompatibleClient{ + httpClient: &http.Client{Timeout: 5 * time.Minute}, + } +} + +// TranscribeResult openai 兼容响应(verbose_json 模式) +// +// 字段非全列:仅取本服务关心的部分。多余字段被丢弃。 +type TranscribeResult struct { + Text string `json:"text"` + Language string `json:"language"` + Duration float64 `json:"duration"` + Segments []TranscribeSegment `json:"segments"` +} + +// TranscribeSegment 单段时间轴文字(verbose_json) +type TranscribeSegment struct { + ID int `json:"id"` + Start float64 `json:"start"` + End float64 `json:"end"` + Text string `json:"text"` +} + +// Transcribe 调用 STT +// +// 流程: +// 1. 从 fileURL 下载音频字节(HTTP GET,复用 httpClient timeout) +// 2. 构造 multipart 请求体 +// 3. POST {base_url}/audio/transcriptions,附 Authorization: Bearer {api_key} +// 4. 解析 JSON 响应 +// +// 参数: +// - fileURL:录制文件公网/内网可达 URL(meeting_recordings.file_url) +// - language:可选,传 "" 让模型自检 +func (c *OpenAICompatibleClient) Transcribe(ctx context.Context, cfg *dao.STTConfig, fileURL, language string) (*TranscribeResult, error) { + if cfg == nil { + return nil, fmt.Errorf("STT config is nil") + } + if fileURL == "" { + return nil, fmt.Errorf("file_url is empty") + } + + // 1) 下载音频 + audioBytes, filename, err := c.downloadAudio(ctx, fileURL) + if err != nil { + return nil, fmt.Errorf("下载录制文件失败: %w", err) + } + + // 2) 拼装 multipart + body := &bytes.Buffer{} + writer := multipart.NewWriter(body) + + filePart, err := writer.CreateFormFile("file", filename) + if err != nil { + return nil, fmt.Errorf("创建 multipart file 字段失败: %w", err) + } + if _, err := filePart.Write(audioBytes); err != nil { + return nil, fmt.Errorf("写入 multipart file 内容失败: %w", err) + } + if err := writer.WriteField("model", cfg.Model.ModelCode); err != nil { + return nil, fmt.Errorf("写入 model 字段失败: %w", err) + } + if err := writer.WriteField("response_format", "verbose_json"); err != nil { + return nil, fmt.Errorf("写入 response_format 字段失败: %w", err) + } + if language != "" { + _ = writer.WriteField("language", language) + } + if err := writer.Close(); err != nil { + return nil, fmt.Errorf("关闭 multipart writer 失败: %w", err) + } + + // 3) 构造请求 + endpoint := strings.TrimRight(cfg.Provider.BaseURL, "/") + "/audio/transcriptions" + req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, body) + if err != nil { + return nil, fmt.Errorf("构造 STT 请求失败: %w", err) + } + req.Header.Set("Authorization", "Bearer "+cfg.Key.APIKey) + req.Header.Set("Content-Type", writer.FormDataContentType()) + + resp, err := c.httpClient.Do(req) + if err != nil { + return nil, fmt.Errorf("调用 STT 接口失败: %w", err) + } + defer resp.Body.Close() + + respBytes, err := io.ReadAll(resp.Body) + if err != nil { + return nil, fmt.Errorf("读取 STT 响应失败: %w", err) + } + + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + // 截断响应体,避免日志爆炸 + preview := string(respBytes) + if len(preview) > 300 { + preview = preview[:300] + "..." + } + return nil, fmt.Errorf("STT 接口返回 %d: %s", resp.StatusCode, preview) + } + + // 4) 解析 + var result TranscribeResult + if err := json.Unmarshal(respBytes, &result); err != nil { + preview := string(respBytes) + if len(preview) > 300 { + preview = preview[:300] + "..." + } + return nil, fmt.Errorf("解析 STT 响应失败: %v, body=%s", err, preview) + } + if result.Text == "" && len(result.Segments) == 0 { + // 部分供应商在 response_format=json 模式只返回 text;这里没拿到任何结果视为异常 + return nil, fmt.Errorf("STT 响应内容为空") + } + return &result, nil +} + +// downloadAudio 从给定 URL 拉取音频内容,返回字节流 + 推断的文件名 +// +// 文件名仅作为 multipart 的 filename 参数,主要决定 Content-Type 推断; +// 取 URL path 末段,无后缀则默认 recording.mp4。 +func (c *OpenAICompatibleClient) downloadAudio(ctx context.Context, fileURL string) ([]byte, string, error) { + req, err := http.NewRequestWithContext(ctx, http.MethodGet, fileURL, nil) + if err != nil { + return nil, "", err + } + resp, err := c.httpClient.Do(req) + if err != nil { + return nil, "", err + } + defer resp.Body.Close() + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + return nil, "", fmt.Errorf("下载失败 status=%d", resp.StatusCode) + } + data, err := io.ReadAll(resp.Body) + if err != nil { + return nil, "", err + } + + filename := path.Base(req.URL.Path) + if filename == "" || filename == "/" || !strings.Contains(filename, ".") { + filename = "recording.mp4" + } + return data, filename, nil +} diff --git a/backend/go-service/app/transcribe/service/transcribe_service.go b/backend/go-service/app/transcribe/service/transcribe_service.go new file mode 100644 index 0000000..42683d1 --- /dev/null +++ b/backend/go-service/app/transcribe/service/transcribe_service.go @@ -0,0 +1,210 @@ +// Package service 提供 transcribe 模块的业务服务 +package service + +import ( + "context" + "encoding/json" + "errors" + "fmt" + "time" + + meetingModel "github.com/echochat/backend/app/meeting/model" + "github.com/echochat/backend/app/transcribe/dao" + "github.com/echochat/backend/app/transcribe/model" + "github.com/echochat/backend/pkg/logs" + "go.uber.org/zap" + "gorm.io/gorm" +) + +// 业务错误 +var ( + ErrRecordingNotFound = errors.New("录制不存在") + ErrRecordingNotReady = errors.New("录制尚未就绪,无法转写") + ErrTranscribeRunning = errors.New("已有转写任务进行中,请稍后再试") + ErrSTTNotConfigured = errors.New("语音转写服务未配置") +) + +// TranscribeService 转写编排服务 +// +// 设计要点: +// - Submit 立即返回 (transcript, 已 running),真正的 STT 调用走 goroutine +// - 同一 recording 只允许一个 running 任务,重复提交直接返回当前进行中行 +// - 转写超时 10 分钟;超时后任务自标记 failed +// - 服务重启会留下 running 状态的孤儿行,留给后续兜底任务处理(暂不在本期实现) +type TranscribeService struct { + db *gorm.DB + transcripts *dao.TranscriptDAO + llmConfig *dao.LLMConfigDAO + sttClient *OpenAICompatibleClient +} + +// NewTranscribeService 创建实例 +func NewTranscribeService( + db *gorm.DB, + transcripts *dao.TranscriptDAO, + llmConfig *dao.LLMConfigDAO, + sttClient *OpenAICompatibleClient, +) *TranscribeService { + return &TranscribeService{ + db: db, + transcripts: transcripts, + llmConfig: llmConfig, + sttClient: sttClient, + } +} + +// IsAvailable 转写功能是否可用(LLM 配置库连通) +func (s *TranscribeService) IsAvailable() bool { + return s.llmConfig != nil && s.llmConfig.IsEnabled() +} + +// Submit 提交转写任务 +// +// 行为: +// - 录制必须存在且 status=ready +// - 如果已有 running 行,返回 ErrTranscribeRunning +// - 如果已有 ready 行且 force=false,直接返回已有结果 +// - 否则:插入或重置 transcript 行 → 启动 goroutine 跑真实调用 → 同步返回 running 状态 +// +// 参数: +// - force:true 表示强制重跑(覆盖 ready 行) +// - language:可选,"" 表示自动检测 +// +// 返回值为 Submit 时刻的 transcript 快照,调用方可继续轮询 Get 拿最新状态。 +func (s *TranscribeService) Submit(ctx context.Context, recordingID int64, force bool, language string) (*model.Transcript, error) { + if !s.IsAvailable() { + return nil, ErrSTTNotConfigured + } + + // 1) 校验录制存在且已就绪 + var rec meetingModel.MeetingRecording + err := s.db.WithContext(ctx).Where("id = ?", recordingID).First(&rec).Error + if err != nil { + if errors.Is(err, gorm.ErrRecordNotFound) { + return nil, ErrRecordingNotFound + } + return nil, err + } + if rec.Status != meetingModel.MeetingRecordingStatusReady || rec.FileURL == "" { + return nil, ErrRecordingNotReady + } + + // 2) 检查现有 transcript + existing, err := s.transcripts.GetByRecordingID(ctx, recordingID) + if err != nil { + return nil, err + } + if existing != nil { + switch existing.Status { + case model.TranscriptStatusRunning, model.TranscriptStatusPending: + return existing, ErrTranscribeRunning + case model.TranscriptStatusReady: + if !force { + return existing, nil + } + } + } + + // 3) 创建/更新行为 pending(即将进入 running) + now := time.Now() + t := &model.Transcript{ + RecordingID: recordingID, + RoomID: rec.RoomID, + Status: model.TranscriptStatusPending, + Segments: "[]", + StartedAt: &now, + } + if existing != nil { + t.ID = existing.ID + } + if err := s.transcripts.Upsert(ctx, t); err != nil { + return nil, fmt.Errorf("写入 transcript 失败: %w", err) + } + + // 4) 取最新 ID(Upsert 后 t.ID 已填) + if err := s.transcripts.MarkRunning(ctx, t.ID); err != nil { + return nil, fmt.Errorf("切换 running 状态失败: %w", err) + } + + // 5) 异步执行真实调用 + // 使用全新 context,超时 10min;不沿用入参 ctx,避免 HTTP 请求结束后 ctx 被 cancel 中断后台任务 + go s.runJob(t.ID, rec, language) + + // 返回 running 状态快照 + updated, _ := s.transcripts.GetByRecordingID(ctx, recordingID) + if updated != nil { + return updated, nil + } + return t, nil +} + +// runJob goroutine 内执行真实 STT 调用 +// +// 不返回错误:所有失败都写回数据库 status=failed + error_msg +func (s *TranscribeService) runJob(transcriptID int64, rec meetingModel.MeetingRecording, language string) { + const funcName = "TranscribeService.runJob" + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Minute) + defer cancel() + + logs.Info(ctx, funcName, "STT 任务开始", + zap.Int64("transcript_id", transcriptID), + zap.Int64("recording_id", rec.ID), + ) + + cfg, err := s.llmConfig.PickActiveSTT(ctx) + if err != nil { + s.fail(ctx, transcriptID, "选择 STT 配置失败: "+err.Error()) + return + } + + // 把 provider/model/key 元信息写回 transcript,便于审计 + _ = s.db.WithContext(ctx).Model(&model.Transcript{}). + Where("id = ?", transcriptID). + Updates(map[string]any{ + "provider_code": cfg.Provider.ProviderCode, + "model_code": cfg.Model.ModelCode, + "key_id": cfg.Key.ID, + }).Error + + result, err := s.sttClient.Transcribe(ctx, cfg, rec.FileURL, language) + if err != nil { + s.fail(ctx, transcriptID, "STT 调用失败: "+err.Error()) + return + } + + segmentsJSON, err := json.Marshal(result.Segments) + if err != nil { + segmentsJSON = []byte("[]") + } + + if err := s.transcripts.MarkReady(ctx, transcriptID, + result.Text, string(segmentsJSON), + result.Language, int(result.Duration), + ); err != nil { + logs.Error(ctx, funcName, "写回 ready 状态失败", zap.Error(err)) + return + } + logs.Info(ctx, funcName, "STT 任务完成", + zap.Int64("transcript_id", transcriptID), + zap.Int("text_len", len(result.Text)), + zap.Int("segments", len(result.Segments)), + ) +} + +// fail 统一失败收尾:写日志 + 落库 +func (s *TranscribeService) fail(ctx context.Context, transcriptID int64, msg string) { + logs.Warn(ctx, "TranscribeService.fail", msg, zap.Int64("transcript_id", transcriptID)) + if err := s.transcripts.MarkFailed(ctx, transcriptID, msg); err != nil { + logs.Error(ctx, "TranscribeService.fail", "写回 failed 状态失败", zap.Error(err)) + } +} + +// GetByRecording 拉取某段录制的转写记录(不存在返回 nil, nil) +func (s *TranscribeService) GetByRecording(ctx context.Context, recordingID int64) (*model.Transcript, error) { + return s.transcripts.GetByRecordingID(ctx, recordingID) +} + +// ListByRoom 拉取一场会议下所有转写 +func (s *TranscribeService) ListByRoom(ctx context.Context, roomID int64) ([]model.Transcript, error) { + return s.transcripts.ListByRoomID(ctx, roomID) +} diff --git a/backend/go-service/cmd/server/main.go b/backend/go-service/cmd/server/main.go index a24218d..b60947f 100644 --- a/backend/go-service/cmd/server/main.go +++ b/backend/go-service/cmd/server/main.go @@ -12,6 +12,7 @@ import ( "github.com/echochat/backend/app/im/model" meetingModel "github.com/echochat/backend/app/meeting/model" + transcribeModel "github.com/echochat/backend/app/transcribe/model" "github.com/echochat/backend/app/provider" "github.com/echochat/backend/config" "github.com/echochat/backend/pkg/logs" @@ -59,6 +60,8 @@ func main() { &model.Message{}, // Phase B 新增:会议录制元数据表 &meetingModel.MeetingRecording{}, + // Phase B 新增:会议录制语音转写结果表 + &transcribeModel.Transcript{}, ); err != nil { logs.Fatal(ctx, "main", "IM 表迁移失败", zap.Error(err)) } diff --git a/backend/go-service/config/config.dev.yaml b/backend/go-service/config/config.dev.yaml index 6f93922..4bdae0a 100644 --- a/backend/go-service/config/config.dev.yaml +++ b/backend/go-service/config/config.dev.yaml @@ -65,6 +65,22 @@ media_server: close_retry: 2 # 关闭类接口失败重试次数(指数退避 200/500ms) create_router_retry: 1 # Task 16 Nit:CreateRouter 在 5xx/网络错误时的额外重试次数(默认 1 次,退避 300ms) +# 外部 LLM/STT 配置中心(Phase B 语音转文字) +# 仅做只读访问 t_llm_provider/t_llm_model/t_llm_key +# enabled=false 时跳过 MySQL 连接,转写接口将返回 503 +llm_source: + enabled: false # 默认关闭,本地开发不连外部 MySQL;切换到 true 后必须填齐下方字段 + host: localhost + port: 3306 + user: root + password: "" # 强烈建议通过 ECHOCHAT_LLM_SOURCE_PASSWORD 环境变量覆盖 + dbname: ai_gateway + charset: utf8mb4 + parse_time: true + loc: Local + max_idle_conns: 5 + max_open_conns: 20 + # 会议模块生命周期参数(Phase 2e-2 Task 8 会议状态机) # E2E 测试可通过环境变量 ECHOCHAT_MEETING_HOST_GRACE_SECONDS 等覆盖为更短值以加速触发 meeting: diff --git a/backend/go-service/config/config.docker.yaml b/backend/go-service/config/config.docker.yaml index a6ca6fc..78f7b59 100644 --- a/backend/go-service/config/config.docker.yaml +++ b/backend/go-service/config/config.docker.yaml @@ -57,6 +57,22 @@ media_server: close_retry: 2 create_router_retry: 1 # Task 16 Nit:CreateRouter 在 5xx/网络错误时的额外重试次数 +# 外部 LLM/STT 配置中心(Phase B 语音转文字) +# 仅做只读 t_llm_provider/t_llm_model/t_llm_key +# 通过 ECHOCHAT_LLM_SOURCE_HOST / _USER / _PASSWORD / _DBNAME 在容器中覆盖 +llm_source: + enabled: false + host: "" + port: 3306 + user: "" + password: "" + dbname: "" + charset: utf8mb4 + parse_time: true + loc: Local + max_idle_conns: 5 + max_open_conns: 20 + # 会议模块生命周期参数(Phase 2e-2 Task 8) meeting: host_grace_seconds: 120 diff --git a/backend/go-service/config/config.go b/backend/go-service/config/config.go index 7ffcd96..c110d51 100644 --- a/backend/go-service/config/config.go +++ b/backend/go-service/config/config.go @@ -19,6 +19,55 @@ type Config struct { Minio MinioConfig `mapstructure:"minio"` MediaServer MediaServerConfig `mapstructure:"media_server"` Meeting MeetingConfig `mapstructure:"meeting"` + LLMSource LLMSourceConfig `mapstructure:"llm_source"` // Phase B:外部 LLM 配置库(MySQL,t_llm_provider/model/key) +} + +// LLMSourceConfig 外部 LLM/STT 配置中心 MySQL 连接 +// +// 该库与本服务的 PostgreSQL 主库相互独立,仅做只读: +// - t_llm_provider 提供商列表(base_url / api_protocol) +// - t_llm_model 模型列表(model_type=4 表示语音识别) +// - t_llm_key Key 池(按 weight DESC + today_count ASC 选 key) +// +// Enabled=false 时跳过初始化,转写功能会以 "未配置 STT 服务" 返回 503, +// 方便在没有外部 MySQL 的开发环境只跑会议管理只读视图。 +type LLMSourceConfig struct { + Enabled bool `mapstructure:"enabled"` // false=跳过连接,转写不可用 + Host string `mapstructure:"host"` // MySQL host + Port int `mapstructure:"port"` // MySQL port,默认 3306 + User string `mapstructure:"user"` + Password string `mapstructure:"password"` + DBName string `mapstructure:"dbname"` + Charset string `mapstructure:"charset"` // 默认 utf8mb4 + ParseTime bool `mapstructure:"parse_time"` // 默认 true + Loc string `mapstructure:"loc"` // 默认 Local + MaxIdleConns int `mapstructure:"max_idle_conns"` + MaxOpenConns int `mapstructure:"max_open_conns"` +} + +// DSN 生成 MySQL DSN +// charset/parse_time/loc 缺省时填充常见安全默认值 +func (l *LLMSourceConfig) DSN() string { + charset := l.Charset + if charset == "" { + charset = "utf8mb4" + } + loc := l.Loc + if loc == "" { + loc = "Local" + } + parseTime := "True" + if !l.ParseTime { + parseTime = "False" + } + port := l.Port + if port == 0 { + port = 3306 + } + return fmt.Sprintf( + "%s:%s@tcp(%s:%d)/%s?charset=%s&parseTime=%s&loc=%s", + l.User, l.Password, l.Host, port, l.DBName, charset, parseTime, loc, + ) } // MeetingConfig 会议模块生命周期参数(Phase 2e-2 Task 8) diff --git a/backend/go-service/go.mod b/backend/go-service/go.mod index 882b513..f26c248 100644 --- a/backend/go-service/go.mod +++ b/backend/go-service/go.mod @@ -14,6 +14,7 @@ require ( go.uber.org/zap v1.27.1 golang.org/x/crypto v0.40.0 gopkg.in/natefinch/lumberjack.v2 v2.2.1 + gorm.io/driver/mysql v1.5.7 gorm.io/driver/postgres v1.6.0 gorm.io/gorm v1.31.1 ) diff --git a/backend/go-service/pkg/db/mysql.go b/backend/go-service/pkg/db/mysql.go new file mode 100644 index 0000000..27fe301 --- /dev/null +++ b/backend/go-service/pkg/db/mysql.go @@ -0,0 +1,96 @@ +// Package db 提供数据库连接管理(MySQL 二连接实例,用于外部 LLM 配置中心) +package db + +import ( + "context" + "fmt" + "time" + + "github.com/echochat/backend/config" + "github.com/echochat/backend/pkg/logs" + "go.uber.org/zap" + "gorm.io/driver/mysql" + "gorm.io/gorm" +) + +// LLMSourceDB 包装外部 LLM 配置库连接 +// +// 之所以用 wrapper 而不是直接返回 *gorm.DB: +// - 主库 (*gorm.DB) 在 Wire 中已被全局注入到所有 DAO,不能让 LLM 库被误用为业务库 +// - Enabled=false 时仍要构造一个空 wrapper,避免下游 nil 引用 +// - 下游 DAO 通过 IsEnabled() 决定是否走数据库路径或返回未配置错误 +type LLMSourceDB struct { + db *gorm.DB + enabled bool +} + +// IsEnabled 是否真正连接了 MySQL;false 时禁止访问 DB() +func (l *LLMSourceDB) IsEnabled() bool { + return l != nil && l.enabled && l.db != nil +} + +// DB 获取底层 GORM 实例(仅 IsEnabled() == true 时返回非 nil) +// 调用方应先判 IsEnabled,再使用返回值。 +func (l *LLMSourceDB) DB() *gorm.DB { + if !l.IsEnabled() { + return nil + } + return l.db +} + +// NewLLMSourceDB 初始化外部 LLM 配置 MySQL 连接 +// +// 行为约定: +// - cfg.Enabled == false:跳过连接,返回 enabled=false 的空 wrapper(不报错),保证 Wire 流程不阻塞 +// - cfg.Enabled == true 但连接失败:返回 error,由 main 决定是否硬失败 +// +// 之所以连接失败时硬报错而不是 graceful degrade: +// - 用户已显式 enabled=true,预期就要用,沉默降级反而难排查 +// - 转写功能未启用时只需保持 enabled=false 即可 +func NewLLMSourceDB(cfg *config.LLMSourceConfig) (*LLMSourceDB, error) { + funcName := "db.NewLLMSourceDB" + ctx := context.Background() + + if cfg == nil || !cfg.Enabled { + logs.Info(ctx, funcName, "外部 LLM 配置库未启用,跳过 MySQL 连接(语音转写功能将不可用)") + return &LLMSourceDB{enabled: false}, nil + } + + logs.Info(ctx, funcName, "正在连接外部 LLM 配置 MySQL", + zap.String("host", cfg.Host), + zap.Int("port", cfg.Port), + zap.String("dbname", cfg.DBName), + ) + + gdb, err := gorm.Open(mysql.Open(cfg.DSN()), &gorm.Config{ + Logger: &zapGormLogger{}, + }) + if err != nil { + logs.Error(ctx, funcName, "MySQL 连接失败", zap.Error(err)) + return nil, fmt.Errorf("连接 LLM 配置 MySQL 失败: %w", err) + } + + sqlDB, err := gdb.DB() + if err != nil { + return nil, fmt.Errorf("获取 MySQL 底层 sql.DB 失败: %w", err) + } + maxIdle := cfg.MaxIdleConns + if maxIdle <= 0 { + maxIdle = 5 + } + maxOpen := cfg.MaxOpenConns + if maxOpen <= 0 { + maxOpen = 20 + } + sqlDB.SetMaxIdleConns(maxIdle) + sqlDB.SetMaxOpenConns(maxOpen) + sqlDB.SetConnMaxLifetime(time.Hour) + + if err := sqlDB.Ping(); err != nil { + logs.Error(ctx, funcName, "MySQL Ping 失败", zap.Error(err)) + return nil, fmt.Errorf("ping LLM 配置 MySQL 失败: %w", err) + } + + logs.Info(ctx, funcName, "外部 LLM 配置 MySQL 连接成功") + return &LLMSourceDB{db: gdb, enabled: true}, nil +} diff --git a/backend/go-service/router/router.go b/backend/go-service/router/router.go index f2931f7..7a79789 100644 --- a/backend/go-service/router/router.go +++ b/backend/go-service/router/router.go @@ -15,6 +15,7 @@ import ( meetingApp "github.com/echochat/backend/app/meeting" notifyApp "github.com/echochat/backend/app/notify" "github.com/echochat/backend/app/provider" + transcribeApp "github.com/echochat/backend/app/transcribe" wsApp "github.com/echochat/backend/app/ws" "github.com/echochat/backend/pkg/middleware" "github.com/echochat/backend/pkg/utils" @@ -46,4 +47,5 @@ func Setup(engine *gin.Engine, app *provider.App) { groupApp.RegisterRoutes(engine, app.GroupController, jwtAuth) notifyApp.RegisterRoutes(engine, app.NotifyController, jwtAuth) meetingApp.RegisterRoutes(engine, app.MeetingController, app.MeetingRecordingController, jwtAuth) + transcribeApp.RegisterRoutes(engine, app.TranscribeController, jwtAuth) }