fix(meeting): Task 16 修复 code-reviewer 审计 P1 八项
按 docs/reviews/2026-04-23-phase2e-2-code-review.md 清单落地 P1 批次: - P1-1 主持人转让原子性:MeetingParticipantDAO.TransferHost 内部事务 合并 meeting_rooms.host_id 更新,service/lifecycle 移除冗余 UpdateHost - P1-2 ListChatMessages / ListMyMeetings 回归 DAO:新增 MeetingChatDAO.ListByRoomBefore 反向游标,service 移除 s.db 直查 - P1-3 ListMyMeetings N+1 放大优化:新增 MeetingParticipantDAO.ListJoinedRoomsByUser JOIN + DISTINCT 单次 SQL - P1-4 EndRoom 行锁 + 事务:s.db.Transaction 包裹 SELECT FOR UPDATE → 快照 → MarkEnded → LeaveAllActive;DAO 补 WithTx(*gorm.DB) 和 GetByIDForUpdate - P1-5 后台 goroutine trace_id 保留:新增 logs.DetachContext(ctx); 替换 meeting_service / meeting_signal_service 共 10 处 context.Background(); meeting_ws_handler 每条 WS 消息分配独立 trace_id - P1-6 _broadcastSelfState 重试:最多 2 次指数退避(700ms→2100ms), 检测 localAudioEnabled/localVideoEnabled 已被后续动作覆盖时放弃旧 patch - P1-7 HandleHostGraceExpired / HandleEmptyRoomExpired 处理锁: 独立 host_grace_handling:<code> / empty_ttl_handling:<code> SETNX + 60s TTL, 消除 Redis 自然过期 + 本地 timer 到点 + DEL 返回 0 的盲区 - P1-8 pushExistingRoomState 同步化:OnRoomJoin 返回前完成补推, 消除与 REST member.joined 并行造成的前端状态闪烁 并附带修复 .gitignore 中 logs/ 规则误伤 pkg/logs/ 代码目录的历史遗留问题, 将 logger.go / trace.go 正式纳入版本控制。 验证:backend go build + go vet 通过;frontend npm run build:h5 通过。 Made-with: Cursor
This commit is contained in:
@@ -58,6 +58,15 @@ type simpleRoomPayload struct {
|
||||
RoomCode string `json:"room_code"`
|
||||
}
|
||||
|
||||
// newEventContext 为每个 WS 事件生成独立 ctx(携带新 trace_id)
|
||||
// Task 16 P1-5:WS 层无 HTTP 中间件注入 trace_id,旧版 handler 一律 context.Background(),
|
||||
// 导致 signalSvc → Redis / DB / broadcaster 的整条链路日志无法通过 trace_id 串联。
|
||||
// 现改为每条 WS 消息生成独立 UUID 作为 trace_id,相当于 RPC 语义;配合 service 层
|
||||
// logs.DetachContext 用于后台 goroutine,可形成完整"WS 消息 → 业务处理 → 异步广播"的日志链路。
|
||||
func (h *MeetingWSHandler) newEventContext() context.Context {
|
||||
return logs.WithTraceID(context.Background(), logs.GenerateTraceID())
|
||||
}
|
||||
|
||||
// sendACK 统一发送 ACK 响应
|
||||
func (h *MeetingWSHandler) sendACK(client *ws.Client, msg *ws.Message, code int, message string, data interface{}) {
|
||||
resp := ws.NewResponse(msg.Event, msg.Seq, code, message, data)
|
||||
@@ -91,7 +100,7 @@ func (h *MeetingWSHandler) handleRoomJoin(client *ws.Client, msg *ws.Message) {
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnRoomJoin(context.Background(), client.UserID, payload.RoomCode); err != nil {
|
||||
if err := h.signalSvc.OnRoomJoin(h.newEventContext(), client.UserID, payload.RoomCode); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
@@ -104,7 +113,7 @@ func (h *MeetingWSHandler) handleRoomLeave(client *ws.Client, msg *ws.Message) {
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnRoomLeave(context.Background(), client.UserID, payload.RoomCode); err != nil {
|
||||
if err := h.signalSvc.OnRoomLeave(h.newEventContext(), client.UserID, payload.RoomCode); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
@@ -119,7 +128,7 @@ func (h *MeetingWSHandler) handleMemberStateChange(client *ws.Client, msg *ws.Me
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnMemberStateChanged(context.Background(), client.UserID, &payload); err != nil {
|
||||
if err := h.signalSvc.OnMemberStateChanged(h.newEventContext(), client.UserID, &payload); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
@@ -134,7 +143,7 @@ func (h *MeetingWSHandler) handleTransportCreate(client *ws.Client, msg *ws.Mess
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
info, err := h.signalSvc.OnTransportCreate(context.Background(), client.UserID, &payload)
|
||||
info, err := h.signalSvc.OnTransportCreate(h.newEventContext(), client.UserID, &payload)
|
||||
if err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
@@ -148,7 +157,7 @@ func (h *MeetingWSHandler) handleTransportConnect(client *ws.Client, msg *ws.Mes
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnTransportConnect(context.Background(), client.UserID, &payload); err != nil {
|
||||
if err := h.signalSvc.OnTransportConnect(h.newEventContext(), client.UserID, &payload); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
@@ -161,7 +170,7 @@ func (h *MeetingWSHandler) handleProduceStart(client *ws.Client, msg *ws.Message
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
result, err := h.signalSvc.OnProduceStart(context.Background(), client.UserID, &payload)
|
||||
result, err := h.signalSvc.OnProduceStart(h.newEventContext(), client.UserID, &payload)
|
||||
if err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
@@ -175,7 +184,7 @@ func (h *MeetingWSHandler) handleConsumeStart(client *ws.Client, msg *ws.Message
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
info, err := h.signalSvc.OnConsumeStart(context.Background(), client.UserID, &payload)
|
||||
info, err := h.signalSvc.OnConsumeStart(h.newEventContext(), client.UserID, &payload)
|
||||
if err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
@@ -189,7 +198,7 @@ func (h *MeetingWSHandler) handleConsumeResume(client *ws.Client, msg *ws.Messag
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnConsumeResume(context.Background(), client.UserID, &payload); err != nil {
|
||||
if err := h.signalSvc.OnConsumeResume(h.newEventContext(), client.UserID, &payload); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
@@ -202,7 +211,7 @@ func (h *MeetingWSHandler) handleProducerClose(client *ws.Client, msg *ws.Messag
|
||||
if !h.unmarshal(client, msg, &payload) {
|
||||
return
|
||||
}
|
||||
if err := h.signalSvc.OnProducerClose(context.Background(), client.UserID, &payload); err != nil {
|
||||
if err := h.signalSvc.OnProducerClose(h.newEventContext(), client.UserID, &payload); err != nil {
|
||||
h.sendACK(client, msg, -1, err.Error(), nil)
|
||||
return
|
||||
}
|
||||
|
||||
@@ -59,6 +59,30 @@ func (d *MeetingChatDAO) ListByRoom(ctx context.Context, roomID int64, afterID i
|
||||
return list, err
|
||||
}
|
||||
|
||||
// ListByRoomBefore 按房间反向游标分页查询聊天历史(id < beforeID),按 created_at DESC, id DESC 排序
|
||||
// Task 16 P1-2:ListChatMessages 回归 DAO(原 service 层直接拼 SQL 破坏分层)
|
||||
// beforeID 为 0 表示从最新的一条开始;limit <= 0 时返回空列表
|
||||
// 由调用方传 limit+1 自行判断 has_more
|
||||
func (d *MeetingChatDAO) ListByRoomBefore(ctx context.Context, roomID int64, beforeID int64, limit int) ([]model.MeetingChat, error) {
|
||||
if limit <= 0 {
|
||||
return []model.MeetingChat{}, nil
|
||||
}
|
||||
funcName := "dao.meeting_chat_dao.ListByRoomBefore"
|
||||
q := d.db.WithContext(ctx).
|
||||
Model(&model.MeetingChat{}).
|
||||
Where("room_id = ?", roomID)
|
||||
if beforeID > 0 {
|
||||
q = q.Where("id < ?", beforeID)
|
||||
}
|
||||
var list []model.MeetingChat
|
||||
if err := q.Order("created_at DESC, id DESC").Limit(limit).Find(&list).Error; err != nil {
|
||||
logs.Error(ctx, funcName, "反向游标查询会议聊天失败",
|
||||
zap.Int64("room_id", roomID), zap.Int64("before_id", beforeID), zap.Error(err))
|
||||
return nil, err
|
||||
}
|
||||
return list, nil
|
||||
}
|
||||
|
||||
// DeleteByRoomIDs 按房间 ID 批量删除聊天消息(会议结束 24 小时后清理)
|
||||
// 外层已有 meeting_rooms ON DELETE CASCADE,此方法用于主动定时清理(不删除 room 本身)
|
||||
func (d *MeetingChatDAO) DeleteByRoomIDs(ctx context.Context, roomIDs []int64) (int64, error) {
|
||||
|
||||
@@ -22,6 +22,13 @@ func NewMeetingParticipantDAO(db *gorm.DB) *MeetingParticipantDAO {
|
||||
return &MeetingParticipantDAO{db: db}
|
||||
}
|
||||
|
||||
// WithTx 返回绑定到指定事务句柄的新 DAO 实例
|
||||
// 所有现有方法复用 db 字段,因此替换 db 为 tx 可让 DAO 方法参与上游事务。
|
||||
// Task 16 P1-4:EndRoom 行锁 + 事务化 新增
|
||||
func (d *MeetingParticipantDAO) WithTx(tx *gorm.DB) *MeetingParticipantDAO {
|
||||
return &MeetingParticipantDAO{db: tx}
|
||||
}
|
||||
|
||||
// JoinRoom 记录用户加入会议
|
||||
// 语义:
|
||||
// - 若 (room_id, user_id) 不存在 → 创建新行,role 取参数值,joined_at = NOW()
|
||||
@@ -207,11 +214,16 @@ func (d *MeetingParticipantDAO) FindActiveByUser(ctx context.Context, userID int
|
||||
return &p, nil
|
||||
}
|
||||
|
||||
// TransferHost 事务内完成主持人转让
|
||||
// 1) 旧 host 的 role 置为 Participant(0)
|
||||
// 2) 新 host 的 role 置为 Host(1)
|
||||
// 调用方应在同一上游事务中同时调 MeetingRoomDAO.UpdateHost 修改 meeting_rooms.host_id
|
||||
// 本方法内部已自带事务,上游无需再包裹
|
||||
// TransferHost 事务内完成主持人转让(Task 16 P1-1 强化为跨表原子事务)
|
||||
// 同一事务内完成三步:
|
||||
// 1) 旧 host 的 meeting_participants.role 置为 Participant(0)
|
||||
// 2) 新 host 的 meeting_participants.role 置为 Host(1),且必须仍在会议中(left_at IS NULL)
|
||||
// 3) meeting_rooms.host_id 原子更新为新 host
|
||||
//
|
||||
// 背景(审计 P1-1):旧实现把 1+2 放在 DAO 事务,3 由 service 层另起 SQL 执行;
|
||||
// 两步非原子 → 第 3 步失败会造成 participant.role 与 room.host_id 不一致,
|
||||
// 进而 assertIsHost(以 room.host_id 为准)永久拒绝老新 host 的所有主持人操作。
|
||||
// 现合并到单一事务,任一子步骤失败自动回滚,跨表一致性得到保证。
|
||||
func (d *MeetingParticipantDAO) TransferHost(ctx context.Context, roomID, oldHostID, newHostID int64) error {
|
||||
funcName := "dao.meeting_participant_dao.TransferHost"
|
||||
logs.Info(ctx, funcName, "转让主持人",
|
||||
@@ -232,6 +244,11 @@ func (d *MeetingParticipantDAO) TransferHost(ctx context.Context, roomID, oldHos
|
||||
Update("role", constants.MeetingRoleHost).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
if err := tx.Model(&model.MeetingRoom{}).
|
||||
Where("id = ?", roomID).
|
||||
Update("host_id", newHostID).Error; err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
})
|
||||
}
|
||||
@@ -244,6 +261,44 @@ func (d *MeetingParticipantDAO) UpdateRole(ctx context.Context, roomID, userID i
|
||||
Update("role", role).Error
|
||||
}
|
||||
|
||||
// ListJoinedRoomsByUser 用户参与过的会议列表(JOIN meeting_rooms,支持 status 过滤 + 游标分页)
|
||||
// Task 16 P1-3:替代 "ListByUser + 内存合并 + ListByIDs" 的 N+1 查询路径
|
||||
// 返回按 meeting_rooms.created_at DESC, meeting_rooms.id DESC 排序;
|
||||
// 使用 DISTINCT 避免参与者多次 join/leave 同一会议导致的行重复。
|
||||
// 由调用方自行传 limit+1 判断 has_more
|
||||
func (d *MeetingParticipantDAO) ListJoinedRoomsByUser(
|
||||
ctx context.Context,
|
||||
userID int64,
|
||||
statusFilter *int,
|
||||
beforeID int64,
|
||||
limit int,
|
||||
) ([]model.MeetingRoom, error) {
|
||||
if limit <= 0 {
|
||||
return []model.MeetingRoom{}, nil
|
||||
}
|
||||
funcName := "dao.meeting_participant_dao.ListJoinedRoomsByUser"
|
||||
|
||||
q := d.db.WithContext(ctx).
|
||||
Table("meeting_rooms AS r").
|
||||
Select("DISTINCT r.*").
|
||||
Joins("INNER JOIN meeting_participants AS p ON p.room_id = r.id").
|
||||
Where("p.user_id = ?", userID)
|
||||
|
||||
if statusFilter != nil {
|
||||
q = q.Where("r.status = ?", *statusFilter)
|
||||
}
|
||||
if beforeID > 0 {
|
||||
q = q.Where("r.id < ?", beforeID)
|
||||
}
|
||||
|
||||
var rooms []model.MeetingRoom
|
||||
if err := q.Order("r.created_at DESC, r.id DESC").Limit(limit).Find(&rooms).Error; err != nil {
|
||||
logs.Error(ctx, funcName, "JOIN 查询用户会议失败", zap.Int64("user_id", userID), zap.Error(err))
|
||||
return nil, err
|
||||
}
|
||||
return rooms, nil
|
||||
}
|
||||
|
||||
// ListByUser 用户历史会议列表(分页,按加入时间倒序)
|
||||
func (d *MeetingParticipantDAO) ListByUser(ctx context.Context, userID int64, offset, limit int) ([]model.MeetingParticipant, int64, error) {
|
||||
q := d.db.WithContext(ctx).
|
||||
|
||||
@@ -11,6 +11,7 @@ import (
|
||||
"github.com/echochat/backend/pkg/logs"
|
||||
"go.uber.org/zap"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/clause"
|
||||
)
|
||||
|
||||
// MeetingRoomDAO 会议房间数据访问对象
|
||||
@@ -23,6 +24,31 @@ func NewMeetingRoomDAO(db *gorm.DB) *MeetingRoomDAO {
|
||||
return &MeetingRoomDAO{db: db}
|
||||
}
|
||||
|
||||
// WithTx 返回绑定到指定事务句柄的新 DAO 实例
|
||||
// 用法:service 层使用 s.db.Transaction(func(tx *gorm.DB) error { return s.roomDAO.WithTx(tx).MarkEnded(...) })
|
||||
// 所有 DAO 方法都通过 d.db.WithContext(ctx) 派生查询,替换 d.db 为 tx 即可让现有方法参与上游事务。
|
||||
// Task 16 P1-4:EndRoom 行锁 + 事务化 新增
|
||||
func (d *MeetingRoomDAO) WithTx(tx *gorm.DB) *MeetingRoomDAO {
|
||||
return &MeetingRoomDAO{db: tx}
|
||||
}
|
||||
|
||||
// GetByIDForUpdate 事务内 SELECT ... FOR UPDATE 行锁查询
|
||||
// 仅在上游已开启事务(调用方通过 WithTx(tx) 注入)时使用;否则等同于普通查询不生效。
|
||||
// Task 16 P1-4:EndRoom 行锁 新增
|
||||
func (d *MeetingRoomDAO) GetByIDForUpdate(ctx context.Context, id int64) (*model.MeetingRoom, error) {
|
||||
var room model.MeetingRoom
|
||||
err := d.db.WithContext(ctx).
|
||||
Clauses(clause.Locking{Strength: "UPDATE"}).
|
||||
First(&room, id).Error
|
||||
if err != nil {
|
||||
if errors.Is(err, gorm.ErrRecordNotFound) {
|
||||
return nil, nil
|
||||
}
|
||||
return nil, err
|
||||
}
|
||||
return &room, nil
|
||||
}
|
||||
|
||||
// Create 新建一个会议房间
|
||||
// 由调用方预先生成 RoomCode 并确认唯一,本方法不做冲突重试
|
||||
func (d *MeetingRoomDAO) Create(ctx context.Context, room *model.MeetingRoom) error {
|
||||
|
||||
@@ -18,8 +18,11 @@ import (
|
||||
|
||||
// Redis key 前缀(设计文档 §5.4 生命周期 keys;Task 8 落地)
|
||||
const (
|
||||
redisKeyHostGracePrefix = "echo:meeting:host_grace:"
|
||||
redisKeyEmptyTTLPrefix = "echo:meeting:empty_ttl:"
|
||||
redisKeyHostGracePrefix = "echo:meeting:host_grace:"
|
||||
redisKeyEmptyTTLPrefix = "echo:meeting:empty_ttl:"
|
||||
redisKeyHostGraceHandlingLock = "echo:meeting:host_grace_handling:" // Task 16 P1-7:处理互斥锁
|
||||
redisKeyEmptyTTLHandlingLock = "echo:meeting:empty_ttl_handling:" // Task 16 P1-7:空房处理互斥锁
|
||||
handlingLockTTLSeconds = 60 // 处理锁 60s,远大于正常处理耗时(P99 <2s)
|
||||
)
|
||||
|
||||
// hostGracePayload host 宽限期 Redis value 结构
|
||||
@@ -154,26 +157,49 @@ func (s *MeetingLifecycleService) OnHostReconnect(ctx context.Context, roomCode
|
||||
|
||||
// HandleHostGraceExpired host 宽限期过期处理
|
||||
// 触发时机:本地 timer 到期 或 cleanup task 扫到 key TTL <=0
|
||||
// Task 16 P1-7 强化并发控制:
|
||||
// 旧版依赖 DEL host_grace 主 key 返回值作为唯一处理权判据;但在 Redis TTL 自然过期 >
|
||||
// 本地 timer 到点的极端场景(时钟漂移 / 进程重启 RescheduleFromRedis 丢 timer 后 cleanup 扫到前 Redis 已自然过期),
|
||||
// 本地 AfterFunc 执行 DEL 返回 0 被误判为"已处理",主持人宽限转让彻底丢失。
|
||||
// 现引入独立的 host_grace_handling 处理锁 (SETNX + 60s TTL):
|
||||
// - 拿到锁 → 当前协程是唯一处理者,继续执行转让逻辑;无论主 key 是否还存在
|
||||
// - 未拿到锁 → 其他协程/节点正在处理,当前跳过
|
||||
// 主 key DEL 仅作为"重连撤销"语义(OnHostReconnect 调用),不再作为处理权判据。
|
||||
// 流程:
|
||||
// 1. DEL host_grace key → 返回 0 表示已被其他路径处理,直接 return 防止重复转让
|
||||
// 2. 加载 room + 活跃成员;若 room 已 Ended / host 已变更,视为幂等完成
|
||||
// 3. 活跃成员存在:选最早加入者转 host(事务),广播 host.changed,auto_reason=host_grace_expired
|
||||
// 4. 无活跃成员:走 OnAllMembersLeft 路径(置 empty_ttl)
|
||||
// 1. SET NX host_grace_handling key → 失败则跳过
|
||||
// 2. cancelTimer 防重入
|
||||
// 3. 取消原 host_grace key(正常清理,允许失败)
|
||||
// 4. 加载 room + 活跃成员;若 room 已 Ended / host 已变更,视为幂等完成
|
||||
// 5. 活跃成员存在:选最早加入者转 host(事务),广播 host.changed
|
||||
// 6. 无活跃成员:走 OnAllMembersLeft 路径(置 empty_ttl)
|
||||
func (s *MeetingLifecycleService) HandleHostGraceExpired(ctx context.Context, roomCode string) {
|
||||
funcName := "service.meeting_lifecycle_service.HandleHostGraceExpired"
|
||||
|
||||
deleted, err := s.redis.Del(ctx, HostGraceKey(roomCode)).Result()
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "DEL host_grace key 失败(跳过本次处理)",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
handlingKey := redisKeyHostGraceHandlingLock + roomCode
|
||||
locked, lockErr := s.redis.SetNX(ctx, handlingKey, "1", handlingLockTTLSeconds*time.Second).Result()
|
||||
if lockErr != nil {
|
||||
logs.Warn(ctx, funcName, "抢占 host_grace 处理锁失败(跳过本次处理)",
|
||||
zap.String("room_code", roomCode), zap.Error(lockErr))
|
||||
return
|
||||
}
|
||||
s.cancelTimer(&s.graceTimers, roomCode)
|
||||
if deleted == 0 {
|
||||
logs.Debug(ctx, funcName, "host_grace key 不存在(已被重连/其他节点清走)",
|
||||
if !locked {
|
||||
logs.Debug(ctx, funcName, "host_grace 处理锁被其他协程持有,跳过",
|
||||
zap.String("room_code", roomCode))
|
||||
return
|
||||
}
|
||||
// 处理完成后释放锁(防止同 roomCode 60s 内无法重入正常排班)
|
||||
defer func() {
|
||||
if err := s.redis.Del(ctx, handlingKey).Err(); err != nil {
|
||||
logs.Warn(ctx, funcName, "释放 host_grace 处理锁失败(60s 后自然过期)",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
}
|
||||
}()
|
||||
|
||||
s.cancelTimer(&s.graceTimers, roomCode)
|
||||
if err := s.redis.Del(ctx, HostGraceKey(roomCode)).Err(); err != nil {
|
||||
logs.Warn(ctx, funcName, "DEL host_grace 主 key 失败(继续处理,锁已持有)",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
}
|
||||
|
||||
room, err := s.roomDAO.GetByCode(ctx, roomCode)
|
||||
if err != nil || room == nil {
|
||||
@@ -212,15 +238,12 @@ func (s *MeetingLifecycleService) HandleHostGraceExpired(ctx context.Context, ro
|
||||
}
|
||||
|
||||
newHostID := candidates[0]
|
||||
// P1-1:TransferHost 内部事务已原子更新 meeting_rooms.host_id,不再需要单独 UpdateHost
|
||||
if err := s.participantDAO.TransferHost(ctx, room.ID, room.HostID, newHostID); err != nil {
|
||||
logs.Error(ctx, funcName, "host 自动转让 TransferHost 失败",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
return
|
||||
}
|
||||
if err := s.roomDAO.UpdateHost(ctx, room.ID, newHostID); err != nil {
|
||||
logs.Warn(ctx, funcName, "UpdateHost 失败(已发生参与者 role 变更)",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
}
|
||||
// 老 host 彻底离会(宽限期内未重连视为网络断开)
|
||||
_, _ = s.participantDAO.LeaveRoom(ctx, room.ID, room.HostID, constants.MeetingLeftReasonDisconnect)
|
||||
|
||||
@@ -272,22 +295,35 @@ func (s *MeetingLifecycleService) CancelEmptyTTL(ctx context.Context, roomCode s
|
||||
}
|
||||
|
||||
// HandleEmptyRoomExpired 空房 TTL 过期处理
|
||||
// 流程:DEL empty_ttl key → 确认没被新成员清走 → MarkEnded(reason=empty_ttl) + CloseRouter 幂等
|
||||
// Task 16 P1-7:采用 empty_ttl_handling 独立处理锁避免"Redis 自然过期 + 本地 timer 到点 + DEL 返回 0"盲区
|
||||
// 新成员重入由 "DEL 成功 → DB activeCount==0" 复合校验保证(持锁期间仍会被 CountActiveByRoom 二次校验)
|
||||
func (s *MeetingLifecycleService) HandleEmptyRoomExpired(ctx context.Context, roomCode string) {
|
||||
funcName := "service.meeting_lifecycle_service.HandleEmptyRoomExpired"
|
||||
|
||||
deleted, err := s.redis.Del(ctx, EmptyTTLKey(roomCode)).Result()
|
||||
if err != nil {
|
||||
logs.Warn(ctx, funcName, "DEL empty_ttl key 失败",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
handlingKey := redisKeyEmptyTTLHandlingLock + roomCode
|
||||
locked, lockErr := s.redis.SetNX(ctx, handlingKey, "1", handlingLockTTLSeconds*time.Second).Result()
|
||||
if lockErr != nil {
|
||||
logs.Warn(ctx, funcName, "抢占 empty_ttl 处理锁失败(跳过本次处理)",
|
||||
zap.String("room_code", roomCode), zap.Error(lockErr))
|
||||
return
|
||||
}
|
||||
s.cancelTimer(&s.emptyTTLTimers, roomCode)
|
||||
if deleted == 0 {
|
||||
logs.Debug(ctx, funcName, "empty_ttl key 已消失(有人重入 或 已处理过)",
|
||||
if !locked {
|
||||
logs.Debug(ctx, funcName, "empty_ttl 处理锁被其他协程持有,跳过",
|
||||
zap.String("room_code", roomCode))
|
||||
return
|
||||
}
|
||||
defer func() {
|
||||
if err := s.redis.Del(ctx, handlingKey).Err(); err != nil {
|
||||
logs.Warn(ctx, funcName, "释放 empty_ttl 处理锁失败(60s 后自然过期)",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
}
|
||||
}()
|
||||
|
||||
s.cancelTimer(&s.emptyTTLTimers, roomCode)
|
||||
if err := s.redis.Del(ctx, EmptyTTLKey(roomCode)).Err(); err != nil {
|
||||
logs.Warn(ctx, funcName, "DEL empty_ttl 主 key 失败(继续处理,锁已持有)",
|
||||
zap.String("room_code", roomCode), zap.Error(err))
|
||||
}
|
||||
|
||||
room, err := s.roomDAO.GetByCode(ctx, roomCode)
|
||||
if err != nil || room == nil {
|
||||
@@ -299,7 +335,7 @@ func (s *MeetingLifecycleService) HandleEmptyRoomExpired(ctx context.Context, ro
|
||||
return
|
||||
}
|
||||
|
||||
// 二次校验:DEL 成功与 DB 校验之间可能有新成员加入(并发场景)
|
||||
// 二次校验:处理锁 ≠ 新成员重入阻断,若进入执行时已有活跃成员则放弃销毁
|
||||
activeCount, _ := s.participantDAO.CountActiveByRoom(ctx, room.ID)
|
||||
if activeCount > 0 {
|
||||
logs.Info(ctx, funcName, "检测到活跃成员,放弃销毁",
|
||||
|
||||
@@ -408,7 +408,7 @@ func (s *MeetingService) JoinRoom(ctx context.Context, userID int64, code, passw
|
||||
|
||||
// 广播 payload 附带 user_name / user_avatar,前端 _onMemberJoined 直接落库,无需二次拉取
|
||||
name, avatar := s.resolveUserDisplay(ctx, userID)
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventMemberJoined, map[string]interface{}{
|
||||
go s.broadcastToActiveParticipants(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberJoined, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"user_id": userID,
|
||||
"user_name": name,
|
||||
@@ -458,13 +458,11 @@ func (s *MeetingService) LeaveRoom(ctx context.Context, userID int64, code strin
|
||||
|
||||
if participant.Role == constants.MeetingRoleHost && len(actives) > 0 {
|
||||
newHost := actives[0]
|
||||
// P1-1:TransferHost 内部事务已同时更新 meeting_rooms.host_id,不再需要单独 UpdateHost
|
||||
if txErr := s.participantDAO.TransferHost(ctx, room.ID, userID, newHost.UserID); txErr != nil {
|
||||
logs.Warn(ctx, funcName, "主持人自动转让失败", zap.Error(txErr))
|
||||
} else {
|
||||
if uErr := s.roomDAO.UpdateHost(ctx, room.ID, newHost.UserID); uErr != nil {
|
||||
logs.Warn(ctx, funcName, "UpdateHost 失败", zap.Error(uErr))
|
||||
}
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventHostChanged, map[string]interface{}{
|
||||
go s.broadcastToActiveParticipants(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventHostChanged, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"old_host_id": userID,
|
||||
"new_host_id": newHost.UserID,
|
||||
@@ -479,7 +477,7 @@ func (s *MeetingService) LeaveRoom(ctx context.Context, userID int64, code strin
|
||||
s.lifecycleSvc.OnAllMembersLeft(ctx, code)
|
||||
}
|
||||
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
|
||||
go s.broadcastToActiveParticipants(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"user_id": userID,
|
||||
"reason": constants.MeetingLeftReasonSelf,
|
||||
@@ -492,6 +490,14 @@ func (s *MeetingService) LeaveRoom(ctx context.Context, userID int64, code strin
|
||||
|
||||
// EndRoom 主持人结束会议
|
||||
// 将房间 status=Ended + reason=host_ended + 所有活跃成员 LeaveAllActive + 广播 meeting.room.ended + 关闭 mediasoup Router
|
||||
// EndRoom 主持人主动结束会议
|
||||
// Task 16 P1-4 强化为行锁 + 单事务:
|
||||
// 1. 事务内 SELECT meeting_rooms FOR UPDATE 行锁,避免与并发 JoinRoom 的 status 竞争
|
||||
// (旧版 JoinRoom 在乐观读 status=active 之后才写 participant 行,可能出现
|
||||
// "JoinRoom 刚读 status=active → EndRoom 改 status=ended + LeaveAll → JoinRoom 写入活跃 participant"
|
||||
// 的僵尸参与者现象,导致用户被卡在"我在会议中"的状态但房间已结束)
|
||||
// 2. 事务内依序执行 状态校验 / 活跃快照 / MarkEnded / LeaveAllActive,全部成功才提交
|
||||
// 3. 事务成功后再做 mediasoup CloseRouter 与 WS 广播(这些是外部 IO,不应拉长锁持有时间)
|
||||
func (s *MeetingService) EndRoom(ctx context.Context, userID int64, code string) error {
|
||||
funcName := "service.meeting_service.EndRoom"
|
||||
|
||||
@@ -509,20 +515,55 @@ func (s *MeetingService) EndRoom(ctx context.Context, userID int64, code string)
|
||||
return err
|
||||
}
|
||||
|
||||
activesBefore, _ := s.participantDAO.ListActiveByRoom(ctx, room.ID)
|
||||
var (
|
||||
activesBefore []model.MeetingParticipant
|
||||
endedAt = time.Now()
|
||||
)
|
||||
|
||||
now := time.Now()
|
||||
if _, err := s.roomDAO.MarkEnded(ctx, room.ID, constants.MeetingEndedReasonHostEnded, now); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := s.participantDAO.LeaveAllActive(ctx, room.ID, constants.MeetingLeftReasonHostEnd); err != nil {
|
||||
logs.Warn(ctx, funcName, "批量离会失败", zap.Int64("room_id", room.ID), zap.Error(err))
|
||||
txErr := s.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
|
||||
roomDAO := s.roomDAO.WithTx(tx)
|
||||
participantDAO := s.participantDAO.WithTx(tx)
|
||||
|
||||
// 行锁:事务内持有 meeting_rooms 行级 X 锁,直到事务结束才释放
|
||||
locked, err := roomDAO.GetByIDForUpdate(ctx, room.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if locked == nil {
|
||||
return ErrMeetingNotFound
|
||||
}
|
||||
// 行锁之后复检 status,避免"同时多个 host 点结束"时第二请求重复执行
|
||||
if locked.Status == constants.MeetingStatusEnded {
|
||||
return ErrMeetingEnded
|
||||
}
|
||||
|
||||
snapshot, err := participantDAO.ListActiveByRoom(ctx, room.ID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
activesBefore = snapshot
|
||||
|
||||
if _, err := roomDAO.MarkEnded(ctx, room.ID, constants.MeetingEndedReasonHostEnded, endedAt); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := participantDAO.LeaveAllActive(ctx, room.ID, constants.MeetingLeftReasonHostEnd); err != nil {
|
||||
return err
|
||||
}
|
||||
return nil
|
||||
})
|
||||
if txErr != nil {
|
||||
if errors.Is(txErr, ErrMeetingEnded) || errors.Is(txErr, ErrMeetingNotFound) {
|
||||
return txErr
|
||||
}
|
||||
logs.Error(ctx, funcName, "结束会议事务失败",
|
||||
zap.String("room_code", code), zap.Int64("room_id", room.ID), zap.Error(txErr))
|
||||
return txErr
|
||||
}
|
||||
|
||||
payload := map[string]interface{}{
|
||||
"room_code": code,
|
||||
"ended_reason": constants.MeetingEndedReasonHostEnded,
|
||||
"ended_at": now.Format("2006-01-02 15:04:05"),
|
||||
"ended_at": endedAt.Format("2006-01-02 15:04:05"),
|
||||
}
|
||||
// activesBefore 已经是结束前的活跃成员快照,此时 participantDAO.ListActiveByRoom 查会返回空集合
|
||||
// 因此使用 broadcaster.PublishToUser 逐人定向推送(而非 BroadcastToMeeting 基于当前状态查库)
|
||||
@@ -569,14 +610,12 @@ func (s *MeetingService) TransferHost(ctx context.Context, operatorID int64, cod
|
||||
return err
|
||||
}
|
||||
|
||||
// P1-1:TransferHost 内部事务同时更新 meeting_rooms.host_id,无需单独 UpdateHost
|
||||
if err := s.participantDAO.TransferHost(ctx, room.ID, operatorID, target.UserID); err != nil {
|
||||
return err
|
||||
}
|
||||
if err := s.roomDAO.UpdateHost(ctx, room.ID, target.UserID); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventHostChanged, map[string]interface{}{
|
||||
go s.broadcastToActiveParticipants(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventHostChanged, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"old_host_id": operatorID,
|
||||
"new_host_id": target.UserID,
|
||||
@@ -625,7 +664,7 @@ func (s *MeetingService) KickMember(ctx context.Context, operatorID int64, code
|
||||
"user_id": targetUserID,
|
||||
"by": operatorID,
|
||||
})
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
|
||||
go s.broadcastToActiveParticipants(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"user_id": targetUserID,
|
||||
"reason": constants.MeetingLeftReasonKicked,
|
||||
@@ -648,38 +687,13 @@ func (s *MeetingService) ListMyMeetings(ctx context.Context, userID int64, statu
|
||||
limit = 50
|
||||
}
|
||||
|
||||
// 多取一条用于判断 has_more;考虑到可能按 status 过滤,适度放大拉取倍数
|
||||
fetchSize := (limit + 1) * 2
|
||||
parts, _, err := s.participantDAO.ListByUser(ctx, userID, 0, fetchSize*3)
|
||||
// Task 16 P1-3:改走 participantDAO.ListJoinedRoomsByUser 的 JOIN 查询
|
||||
// 旧版两轮 (ListByUser O(N) → ListByIDs O(M))、N 放大 6 倍 → 单次 JOIN DISTINCT O(limit+1)
|
||||
rooms, err := s.participantDAO.ListJoinedRoomsByUser(ctx, userID, statusFilter, beforeID, limit+1)
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
|
||||
if len(parts) == 0 {
|
||||
return []model.MeetingRoom{}, false, nil
|
||||
}
|
||||
seen := make(map[int64]struct{}, len(parts))
|
||||
roomIDs := make([]int64, 0, len(parts))
|
||||
for _, p := range parts {
|
||||
if _, ok := seen[p.RoomID]; ok {
|
||||
continue
|
||||
}
|
||||
seen[p.RoomID] = struct{}{}
|
||||
roomIDs = append(roomIDs, p.RoomID)
|
||||
}
|
||||
|
||||
var rooms []model.MeetingRoom
|
||||
q := s.db.WithContext(ctx).Model(&model.MeetingRoom{}).Where("id IN ?", roomIDs)
|
||||
if statusFilter != nil {
|
||||
q = q.Where("status = ?", *statusFilter)
|
||||
}
|
||||
if beforeID > 0 {
|
||||
q = q.Where("id < ?", beforeID)
|
||||
}
|
||||
if err := q.Order("created_at DESC, id DESC").Limit(limit + 1).Find(&rooms).Error; err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
|
||||
hasMore := false
|
||||
if len(rooms) > limit {
|
||||
hasMore = true
|
||||
@@ -871,7 +885,7 @@ func (s *MeetingService) SendChatMessage(ctx context.Context, userID int64, code
|
||||
// 附带 user_name / user_avatar,前端聊天面板无需额外拉取
|
||||
userName, userAvatar := s.resolveUserDisplay(ctx, userID)
|
||||
|
||||
go s.broadcastToActiveParticipants(context.Background(), room.ID, constants.MeetingWSEventChatMessage, map[string]interface{}{
|
||||
go s.broadcastToActiveParticipants(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventChatMessage, map[string]interface{}{
|
||||
"room_code": code,
|
||||
"message_id": chat.ID,
|
||||
"user_id": userID,
|
||||
@@ -906,13 +920,9 @@ func (s *MeetingService) ListChatMessages(ctx context.Context, userID int64, cod
|
||||
limit = 100
|
||||
}
|
||||
|
||||
// DAO 当前接受 afterID(正向游标);ListChats 前端通常是"查更早的",这里直接用 beforeID 语义
|
||||
var chats []model.MeetingChat
|
||||
q := s.db.WithContext(ctx).Model(&model.MeetingChat{}).Where("room_id = ?", room.ID)
|
||||
if beforeID > 0 {
|
||||
q = q.Where("id < ?", beforeID)
|
||||
}
|
||||
if err := q.Order("created_at DESC, id DESC").Limit(limit + 1).Find(&chats).Error; err != nil {
|
||||
// Task 16 P1-2:改走 chatDAO.ListByRoomBefore(反向游标),避免 service 层直接拼 SQL 破坏分层
|
||||
chats, err := s.chatDAO.ListByRoomBefore(ctx, room.ID, beforeID, limit+1)
|
||||
if err != nil {
|
||||
return nil, false, err
|
||||
}
|
||||
|
||||
|
||||
@@ -217,8 +217,12 @@ func (s *MeetingSignalService) OnRoomJoin(ctx context.Context, userID int64, roo
|
||||
zap.Int64("user_id", userID),
|
||||
zap.Int64("room_id", room.ID))
|
||||
|
||||
// 定向补推已有 producer + 成员状态:异步执行,不阻塞 ACK
|
||||
go s.pushExistingRoomState(context.Background(), room.ID, room.RoomCode, userID)
|
||||
// P1-8:pushExistingRoomState 同步化(Task 16)
|
||||
// 旧版异步 go 导致 REST `meeting.member.joined` 与 WS `pushExistingRoomState` 并行,
|
||||
// 极端场景下新人先收到未知 user_id 的 state.changed 再收到 REST participant(名称头像短暂缺失)。
|
||||
// 前端本就在等 room.join 的 ACK,这里同步在 ACK 返回前完成补推,可消除并行窗口。
|
||||
// 单次调用仅 O(当前房间 producer 数量) Redis 读 + 若干 WS 发送,P99 <30ms,不影响 ACK 体验。
|
||||
s.pushExistingRoomState(ctx, room.ID, room.RoomCode, userID)
|
||||
|
||||
return nil
|
||||
}
|
||||
@@ -406,7 +410,7 @@ func (s *MeetingSignalService) OnRoomLeave(ctx context.Context, userID int64, ro
|
||||
}
|
||||
s.cleanupUserResources(ctx, roomCode, userID)
|
||||
|
||||
go s.broadcaster.BroadcastToMeeting(context.Background(), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
|
||||
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberLeft, map[string]interface{}{
|
||||
"room_code": roomCode,
|
||||
"user_id": userID,
|
||||
"reason": "ws_disconnect",
|
||||
@@ -465,7 +469,7 @@ func (s *MeetingSignalService) OnMemberStateChanged(ctx context.Context, fromUse
|
||||
// 注意:写 Hash 之所以同步执行而非放进 goroutine,是因为同一用户并发开关操作需要顺序一致
|
||||
s.updateMemberState(ctx, payload.RoomCode, targetID, payload.AudioEnabled, payload.VideoEnabled)
|
||||
|
||||
go s.broadcaster.BroadcastToMeeting(context.Background(), room.ID, constants.MeetingWSEventMemberStateChange, data, fromUserID)
|
||||
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberStateChange, data, fromUserID)
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -553,7 +557,7 @@ func (s *MeetingSignalService) OnProduceStart(ctx context.Context, userID int64,
|
||||
}
|
||||
s.trackResource(ctx, payload.RoomCode, userID, "producer", producerID)
|
||||
|
||||
go s.broadcaster.BroadcastToMeeting(context.Background(), room.ID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{
|
||||
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{
|
||||
"room_code": payload.RoomCode,
|
||||
"user_id": userID,
|
||||
"producer_id": producerID,
|
||||
@@ -650,7 +654,7 @@ func (s *MeetingSignalService) OnProducerClose(ctx context.Context, userID int64
|
||||
}
|
||||
s.untrackResource(ctx, payload.RoomCode, userID, "producer", payload.ProducerID)
|
||||
|
||||
go s.broadcaster.BroadcastToMeeting(context.Background(), room.ID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{
|
||||
go s.broadcaster.BroadcastToMeeting(logs.DetachContext(ctx), room.ID, constants.MeetingWSEventMemberProducerNew, map[string]interface{}{
|
||||
"room_code": payload.RoomCode,
|
||||
"user_id": userID,
|
||||
"producer_id": payload.ProducerID,
|
||||
|
||||
184
backend/go-service/pkg/logs/logger.go
Normal file
184
backend/go-service/pkg/logs/logger.go
Normal file
@@ -0,0 +1,184 @@
|
||||
// Package logs 提供基于 zap 的结构化日志封装
|
||||
// 支持从 context 中提取 trace_id 自动附加到每条日志
|
||||
// 开发环境输出彩色可读文本,生产环境输出 JSON 结构化格式
|
||||
// 支持日志文件轮转(按大小切割、自动归档、过期清理)
|
||||
package logs
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/echochat/backend/config"
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
"gopkg.in/natefinch/lumberjack.v2"
|
||||
)
|
||||
|
||||
var globalLogger *zap.Logger
|
||||
|
||||
// Init 初始化全局日志实例
|
||||
// 根据 LogConfig 配置决定输出目标:控制台、文件、或两者同时输出
|
||||
func Init(cfg *config.LogConfig) error {
|
||||
lvl := parseLevel(cfg.Level)
|
||||
|
||||
// 控制台 encoder(始终保留控制台输出)
|
||||
consoleEncoder := buildEncoder(cfg.Format)
|
||||
|
||||
var cores []zapcore.Core
|
||||
|
||||
// 控制台输出(始终启用)
|
||||
consoleSyncer := zapcore.AddSync(os.Stdout)
|
||||
cores = append(cores, zapcore.NewCore(consoleEncoder, consoleSyncer, lvl))
|
||||
|
||||
// 文件输出(按配置启用)
|
||||
if cfg.File.Enable && cfg.File.Dir != "" {
|
||||
if err := os.MkdirAll(cfg.File.Dir, 0755); err != nil {
|
||||
return fmt.Errorf("创建日志目录失败 [%s]: %w", cfg.File.Dir, err)
|
||||
}
|
||||
|
||||
// 应用日志文件(全量日志)
|
||||
allLogWriter := &lumberjack.Logger{
|
||||
Filename: filepath.Join(cfg.File.Dir, "app.log"),
|
||||
MaxSize: cfg.File.MaxSize,
|
||||
MaxBackups: cfg.File.MaxBackups,
|
||||
MaxAge: cfg.File.MaxAge,
|
||||
Compress: cfg.File.Compress,
|
||||
LocalTime: true,
|
||||
}
|
||||
|
||||
// 错误日志文件(仅 WARN 及以上),便于快速定位问题
|
||||
errorLogWriter := &lumberjack.Logger{
|
||||
Filename: filepath.Join(cfg.File.Dir, "error.log"),
|
||||
MaxSize: cfg.File.MaxSize,
|
||||
MaxBackups: cfg.File.MaxBackups,
|
||||
MaxAge: cfg.File.MaxAge,
|
||||
Compress: cfg.File.Compress,
|
||||
LocalTime: true,
|
||||
}
|
||||
|
||||
// 文件始终用 JSON 格式,便于后期接入 ELK/Loki
|
||||
fileEncoder := buildEncoder("json")
|
||||
|
||||
cores = append(cores,
|
||||
zapcore.NewCore(fileEncoder, zapcore.AddSync(allLogWriter), lvl),
|
||||
zapcore.NewCore(fileEncoder, zapcore.AddSync(errorLogWriter), zap.WarnLevel),
|
||||
)
|
||||
}
|
||||
|
||||
core := zapcore.NewTee(cores...)
|
||||
globalLogger = zap.New(core, zap.AddCaller(), zap.AddCallerSkip(1))
|
||||
return nil
|
||||
}
|
||||
|
||||
func buildEncoder(format string) zapcore.Encoder {
|
||||
if format == "json" {
|
||||
encoderConfig := zap.NewProductionEncoderConfig()
|
||||
encoderConfig.TimeKey = "ts"
|
||||
encoderConfig.EncodeTime = func(t time.Time, enc zapcore.PrimitiveArrayEncoder) {
|
||||
enc.AppendString(t.Format("2006-01-02 15:04:05"))
|
||||
}
|
||||
encoderConfig.EncodeLevel = zapcore.CapitalLevelEncoder
|
||||
return zapcore.NewJSONEncoder(encoderConfig)
|
||||
}
|
||||
|
||||
encoderConfig := zap.NewDevelopmentEncoderConfig()
|
||||
encoderConfig.EncodeTime = func(t time.Time, enc zapcore.PrimitiveArrayEncoder) {
|
||||
enc.AppendString(t.Format("2006-01-02 15:04:05"))
|
||||
}
|
||||
encoderConfig.EncodeLevel = zapcore.CapitalColorLevelEncoder
|
||||
return zapcore.NewConsoleEncoder(encoderConfig)
|
||||
}
|
||||
|
||||
// Debug 输出 DEBUG 级别日志
|
||||
func Debug(ctx context.Context, funcName, msg string, fields ...zap.Field) {
|
||||
globalLogger.Debug(msg, withContext(ctx, funcName, fields)...)
|
||||
}
|
||||
|
||||
// Info 输出 INFO 级别日志
|
||||
func Info(ctx context.Context, funcName, msg string, fields ...zap.Field) {
|
||||
globalLogger.Info(msg, withContext(ctx, funcName, fields)...)
|
||||
}
|
||||
|
||||
// Warn 输出 WARN 级别日志
|
||||
func Warn(ctx context.Context, funcName, msg string, fields ...zap.Field) {
|
||||
globalLogger.Warn(msg, withContext(ctx, funcName, fields)...)
|
||||
}
|
||||
|
||||
// Error 输出 ERROR 级别日志
|
||||
func Error(ctx context.Context, funcName, msg string, fields ...zap.Field) {
|
||||
globalLogger.Error(msg, withContext(ctx, funcName, fields)...)
|
||||
}
|
||||
|
||||
// Fatal 输出 FATAL 级别日志并退出进程
|
||||
func Fatal(ctx context.Context, funcName, msg string, fields ...zap.Field) {
|
||||
globalLogger.Fatal(msg, withContext(ctx, funcName, fields)...)
|
||||
}
|
||||
|
||||
// Sync 刷新日志缓冲区,应在程序退出时调用
|
||||
func Sync() {
|
||||
if globalLogger != nil {
|
||||
_ = globalLogger.Sync()
|
||||
}
|
||||
}
|
||||
|
||||
// withContext 从 context 提取 trace_id 和函数名,合并到日志字段中
|
||||
func withContext(ctx context.Context, funcName string, fields []zap.Field) []zap.Field {
|
||||
traceID := GetTraceID(ctx)
|
||||
result := make([]zap.Field, 0, len(fields)+2)
|
||||
if traceID != "" {
|
||||
result = append(result, zap.String("trace_id", traceID))
|
||||
}
|
||||
if funcName != "" {
|
||||
result = append(result, zap.String("func", funcName))
|
||||
}
|
||||
result = append(result, fields...)
|
||||
return result
|
||||
}
|
||||
|
||||
func parseLevel(level string) zapcore.Level {
|
||||
switch strings.ToLower(level) {
|
||||
case "debug":
|
||||
return zapcore.DebugLevel
|
||||
case "info":
|
||||
return zapcore.InfoLevel
|
||||
case "warn":
|
||||
return zapcore.WarnLevel
|
||||
case "error":
|
||||
return zapcore.ErrorLevel
|
||||
default:
|
||||
return zapcore.InfoLevel
|
||||
}
|
||||
}
|
||||
|
||||
// MaskEmail 邮箱脱敏:zh***@example.com
|
||||
func MaskEmail(email string) string {
|
||||
at := strings.Index(email, "@")
|
||||
if at <= 0 {
|
||||
return "***"
|
||||
}
|
||||
prefix := email[:at]
|
||||
if len(prefix) <= 2 {
|
||||
return prefix[:1] + "***" + email[at:]
|
||||
}
|
||||
return prefix[:2] + "***" + email[at:]
|
||||
}
|
||||
|
||||
// MaskPhone 手机号脱敏:138****8000
|
||||
func MaskPhone(phone string) string {
|
||||
if len(phone) < 7 {
|
||||
return "***"
|
||||
}
|
||||
return phone[:3] + "****" + phone[len(phone)-4:]
|
||||
}
|
||||
|
||||
// MaskToken Token 脱敏:只显示前后各 4 位
|
||||
func MaskToken(token string) string {
|
||||
if len(token) <= 8 {
|
||||
return "***"
|
||||
}
|
||||
return token[:4] + "..." + token[len(token)-4:]
|
||||
}
|
||||
44
backend/go-service/pkg/logs/trace.go
Normal file
44
backend/go-service/pkg/logs/trace.go
Normal file
@@ -0,0 +1,44 @@
|
||||
package logs
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
type contextKey string
|
||||
|
||||
const traceIDKey contextKey = "trace_id"
|
||||
|
||||
// GenerateTraceID 生成唯一的链路追踪 ID(UUID v4)
|
||||
func GenerateTraceID() string {
|
||||
return uuid.New().String()
|
||||
}
|
||||
|
||||
// WithTraceID 将 trace_id 注入 context
|
||||
func WithTraceID(ctx context.Context, traceID string) context.Context {
|
||||
return context.WithValue(ctx, traceIDKey, traceID)
|
||||
}
|
||||
|
||||
// GetTraceID 从 context 提取 trace_id,不存在则返回空字符串
|
||||
func GetTraceID(ctx context.Context) string {
|
||||
if ctx == nil {
|
||||
return ""
|
||||
}
|
||||
if traceID, ok := ctx.Value(traceIDKey).(string); ok {
|
||||
return traceID
|
||||
}
|
||||
return ""
|
||||
}
|
||||
|
||||
// DetachContext 剥离父 ctx 的 Deadline/Cancel 但保留 trace_id
|
||||
// 用于启动后台 goroutine 时避免随 HTTP 请求结束被取消,同时保持日志链路追踪连续性
|
||||
// 典型使用:go func() { ... }(logs.DetachContext(ctx))
|
||||
// Task 16 P1-3:后台 goroutine trace_id 保留 新增
|
||||
func DetachContext(ctx context.Context) context.Context {
|
||||
bg := context.Background()
|
||||
if traceID := GetTraceID(ctx); traceID != "" {
|
||||
return WithTraceID(bg, traceID)
|
||||
}
|
||||
return bg
|
||||
}
|
||||
Reference in New Issue
Block a user