168 lines
6.3 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package roomnotice
import (
"context"
"encoding/json"
"fmt"
"log/slog"
"time"
roomeventsv1 "hyapp.local/api/proto/events/room/v1"
"hyapp/pkg/logx"
"hyapp/pkg/tencentim"
)
// Repository 是 room notice 模块需要的持久化边界。
type Repository interface {
ClaimRoomKickEvents(ctx context.Context, workerID string, limit int, lockTTL time.Duration) ([]RoomKickEvent, error)
ClaimRoomKickEnvelope(ctx context.Context, workerID string, envelope *roomeventsv1.EventEnvelope, lockTTL time.Duration) (RoomKickEvent, bool, error)
MarkRoomKickDelivered(ctx context.Context, event RoomKickEvent, deliveredPayload []byte, nowMs int64) error
MarkRoomKickRetryable(ctx context.Context, event RoomKickEvent, retryCount int, nextRetryAtMS int64, lastErr string, nowMs int64) error
MarkRoomKickFailed(ctx context.Context, event RoomKickEvent, retryCount int, lastErr string, nowMs int64) error
}
// UserMessagePublisher 把房间私有 notice 投递到用户个人实时通道。
type UserMessagePublisher interface {
PublishUserCustomMessage(ctx context.Context, message tencentim.CustomUserMessage) error
}
// Config 保存 room notice 的进程级配置。
type Config struct {
NodeID string
}
// Service 消费 room_outbox 中需要私有通知的房间事实,不拥有房间状态。
type Service struct {
cfg Config
repository Repository
publisher UserMessagePublisher
}
// New 创建 room notice 服务。
func New(cfg Config, repository Repository, publisher UserMessagePublisher) *Service {
return &Service{cfg: cfg, repository: repository, publisher: publisher}
}
// ProcessRoomKickNotices 处理一批“被踢用户本人”的私有通知。
func (s *Service) ProcessRoomKickNotices(ctx context.Context, options RoomNoticeWorkerOptions) (int, error) {
options = normalizeRoomNoticeWorkerOptions(options, s.cfg.NodeID)
if s.repository == nil {
return 0, fmt.Errorf("notice repository is not configured")
}
if s.publisher == nil {
return 0, fmt.Errorf("notice publisher is not configured")
}
events, err := s.repository.ClaimRoomKickEvents(ctx, options.WorkerID, options.BatchSize, options.LockTTL)
if err != nil {
return 0, err
}
processed := 0
for _, event := range events {
if err := s.publishRoomKickEvent(ctx, event, options); err != nil {
return processed, err
}
processed++
}
return processed, nil
}
// ProcessRoomOutboxEnvelope handles one room_outbox MQ message. Unknown event
// types are acknowledged by returning handled=false; poison messages for known
// types still return an error so MQ can retry or dead-letter them.
func (s *Service) ProcessRoomOutboxEnvelope(ctx context.Context, envelope *roomeventsv1.EventEnvelope, options RoomNoticeWorkerOptions) (bool, error) {
if envelope == nil || envelope.GetEventType() != eventRoomUserKicked {
return false, nil
}
options = normalizeRoomNoticeWorkerOptions(options, s.cfg.NodeID)
if s.repository == nil {
return true, fmt.Errorf("notice repository is not configured")
}
if s.publisher == nil {
return true, fmt.Errorf("notice publisher is not configured")
}
event, claimed, err := s.repository.ClaimRoomKickEnvelope(ctx, options.WorkerID, envelope, options.LockTTL)
if err != nil || !claimed {
return true, err
}
return true, s.publishRoomKickEvent(ctx, event, options)
}
func (s *Service) publishRoomKickEvent(ctx context.Context, event RoomKickEvent, options RoomNoticeWorkerOptions) error {
payload, err := roomKickNoticePayload(event)
if err != nil {
return s.markFailed(ctx, event, options, err)
}
publishCtx, cancel := context.WithTimeout(ctx, options.PublishTimeout)
defer cancel()
err = s.publisher.PublishUserCustomMessage(publishCtx, tencentim.CustomUserMessage{
ToAccount: tencentim.FormatUserID(event.TargetUserID),
EventID: event.EventID,
Desc: "room_user_kicked",
// 被踢用户可能已先被移出腾讯 IM 房间群,因此必须用 C2C 兜底Ext 沿用
// Flutter 已消费的房间系统消息协议,保证当前线上客户端不会把私有通知过滤掉。
Ext: "room_system_message",
PayloadJSON: payload,
})
nowMs := time.Now().UnixMilli()
if err == nil {
if markErr := s.repository.MarkRoomKickDelivered(ctx, event, payload, nowMs); markErr != nil {
return markErr
}
logx.Info(ctx, "notice_room_kick_delivered",
slog.String("event_id", event.EventID),
slog.String("app_code", event.AppCode),
slog.String("room_id", event.RoomID),
slog.Int64("target_user_id", event.TargetUserID),
)
return nil
}
return s.markFailed(ctx, event, options, err)
}
func (s *Service) markFailed(ctx context.Context, event RoomKickEvent, options RoomNoticeWorkerOptions, cause error) error {
nowMs := time.Now().UnixMilli()
nextRetryCount := event.RetryCount + 1
if nextRetryCount >= options.MaxRetryCount {
if err := s.repository.MarkRoomKickFailed(ctx, event, nextRetryCount, cause.Error(), nowMs); err != nil {
return err
}
logx.Error(ctx, "notice_room_kick_dead_letter", cause,
slog.String("event_id", event.EventID),
slog.String("app_code", event.AppCode),
slog.String("room_id", event.RoomID),
slog.Int64("target_user_id", event.TargetUserID),
slog.Int("retry_count", nextRetryCount),
)
return nil
}
nextRetryAtMS := nowMs + roomNoticeBackoff(nextRetryCount, options).Milliseconds()
if err := s.repository.MarkRoomKickRetryable(ctx, event, nextRetryCount, nextRetryAtMS, cause.Error(), nowMs); err != nil {
return err
}
logx.Error(ctx, "notice_room_kick_retryable", cause,
slog.String("event_id", event.EventID),
slog.String("app_code", event.AppCode),
slog.String("room_id", event.RoomID),
slog.Int64("target_user_id", event.TargetUserID),
slog.Int("retry_count", nextRetryCount),
slog.Int64("next_retry_at_ms", nextRetryAtMS),
)
return nil
}
func roomKickNoticePayload(event RoomKickEvent) (json.RawMessage, error) {
payload := map[string]any{
"event_id": event.EventID,
"event_type": "room_user_kicked",
"source_event_type": event.EventType,
"app_code": event.AppCode,
"room_id": event.RoomID,
"room_version": event.RoomVersion,
"actor_user_id": event.ActorUserID,
"target_user_id": event.TargetUserID,
"occurred_at_ms": event.OccurredAtMS,
"source_created_at_ms": event.CreatedAtMS,
}
return json.Marshal(payload)
}