291 lines
10 KiB
Go
291 lines
10 KiB
Go
package broadcast
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"strconv"
|
||
"strings"
|
||
|
||
"google.golang.org/grpc/codes"
|
||
"google.golang.org/grpc/status"
|
||
"hyapp/pkg/appcode"
|
||
"hyapp/pkg/gamemq"
|
||
"hyapp/pkg/xerr"
|
||
broadcastdomain "hyapp/services/activity-service/internal/domain/broadcast"
|
||
)
|
||
|
||
const maxGameWinMinCoins = int64(1_000_000_000_000)
|
||
|
||
type cachedGameWinConfig struct {
|
||
config broadcastdomain.GameWinConfig
|
||
expiresAtMS int64
|
||
}
|
||
|
||
// HandleGameEvent 把 game-service 已提交的返奖事实投影成区域 game_win 飘屏。
|
||
// 游戏结算、阈值策略、房间锁状态和 IM 投递分别由 owner 服务持有;这里仅在消费时组装稳定展示快照。
|
||
func (s *Service) HandleGameEvent(ctx context.Context, message gamemq.GameOutboxMessage) (broadcastdomain.ConsumeGameEventResult, error) {
|
||
result := broadcastdomain.ConsumeGameEventResult{EventID: strings.TrimSpace(message.EventID), Status: broadcastdomain.StatusSkipped}
|
||
if strings.TrimSpace(message.EventType) != gamemq.EventTypeGameOrderSettled {
|
||
return result, nil
|
||
}
|
||
opType := strings.ToLower(strings.TrimSpace(message.OpType))
|
||
if opType != "credit" && opType != "payout" {
|
||
// debit/bet/refund 不是中奖到账,确认 MQ 位点但绝不能生成正向中奖飘屏。
|
||
return result, nil
|
||
}
|
||
if strings.TrimSpace(message.AppCode) == "" || strings.TrimSpace(message.EventID) == "" || strings.TrimSpace(message.OrderID) == "" || message.UserID <= 0 || message.CoinAmount < 0 {
|
||
return broadcastdomain.ConsumeGameEventResult{}, xerr.New(xerr.InvalidArgument, "game order settled fact is incomplete")
|
||
}
|
||
|
||
eventCtx := appcode.WithContext(ctx, message.AppCode)
|
||
config, err := s.cachedGameWinBroadcastConfig(eventCtx)
|
||
if err != nil {
|
||
return broadcastdomain.ConsumeGameEventResult{}, err
|
||
}
|
||
if !config.Enabled || message.CoinAmount < config.MinWinCoins {
|
||
return result, nil
|
||
}
|
||
|
||
payload, err := decodeGameWinSourcePayload(message.PayloadJSON)
|
||
if err != nil {
|
||
return broadcastdomain.ConsumeGameEventResult{}, err
|
||
}
|
||
regionID := firstGamePayloadInt64(payload, "visible_region_id", "region_id")
|
||
if regionID <= 0 {
|
||
// 区域未知时不能扩大成全局播报;来源事实仍被合法确认,避免坏数据反复进入 DLQ。
|
||
return result, nil
|
||
}
|
||
|
||
profile, err := s.broadcastUserProfile(eventCtx, message.UserID)
|
||
if err != nil {
|
||
return broadcastdomain.ConsumeGameEventResult{}, fmt.Errorf("resolve game win user profile: %w", err)
|
||
}
|
||
roomID := firstGamePayloadString(payload, "room_id")
|
||
roomSnapshot := RoomBroadcastSnapshot{}
|
||
if roomID != "" {
|
||
if s.roomSnapshotSource == nil {
|
||
return broadcastdomain.ConsumeGameEventResult{}, xerr.New(xerr.Unavailable, "room snapshot source is not configured")
|
||
}
|
||
roomSnapshot, err = s.roomSnapshotSource.GetRoomBroadcastSnapshot(eventCtx, roomID)
|
||
if err != nil {
|
||
if status.Code(err) == codes.NotFound {
|
||
// 房间已删除时仍可展示中奖,但必须移除失效入口,不能把未知 locked=false 发给客户端。
|
||
roomID = ""
|
||
roomSnapshot = RoomBroadcastSnapshot{}
|
||
} else {
|
||
return broadcastdomain.ConsumeGameEventResult{}, fmt.Errorf("resolve game win room snapshot: %w", err)
|
||
}
|
||
}
|
||
}
|
||
|
||
broadcastEventID := "game_win:" + strings.TrimSpace(message.OrderID)
|
||
payloadJSON, err := gameWinPayload(message, payload, broadcastEventID, regionID, profile, roomID, roomSnapshot, s.now().UTC().UnixMilli())
|
||
if err != nil {
|
||
return broadcastdomain.ConsumeGameEventResult{}, err
|
||
}
|
||
published, err := s.PublishRegionBroadcast(eventCtx, PublishInput{
|
||
EventID: broadcastEventID,
|
||
BroadcastType: broadcastdomain.TypeGameWin,
|
||
RegionID: regionID,
|
||
PayloadJSON: payloadJSON,
|
||
})
|
||
if err != nil {
|
||
return broadcastdomain.ConsumeGameEventResult{}, err
|
||
}
|
||
return broadcastdomain.ConsumeGameEventResult{
|
||
EventID: message.EventID,
|
||
Status: published.Status,
|
||
BroadcastEventID: published.EventID,
|
||
BroadcastCreated: published.Created,
|
||
}, nil
|
||
}
|
||
|
||
// GetGameWinBroadcastConfig 直读当前 App 配置并刷新本实例缓存,后台读取不会拿到另一个副本尚未过期的旧值。
|
||
func (s *Service) GetGameWinBroadcastConfig(ctx context.Context) (broadcastdomain.GameWinConfig, error) {
|
||
if s == nil || s.gameWinConfigRepo == nil {
|
||
return broadcastdomain.GameWinConfig{}, xerr.New(xerr.Unavailable, "game win broadcast config repository is not configured")
|
||
}
|
||
config, exists, err := s.gameWinConfigRepo.GetGameWinBroadcastConfig(ctx)
|
||
if err != nil {
|
||
return broadcastdomain.GameWinConfig{}, err
|
||
}
|
||
if !exists {
|
||
config = broadcastdomain.GameWinConfig{Enabled: true, MinWinCoins: s.cfg.GameWinMinCoins}
|
||
}
|
||
s.storeCachedGameWinConfig(appcode.FromContext(ctx), config)
|
||
return config, nil
|
||
}
|
||
|
||
// UpdateGameWinBroadcastConfig 保存完整替换策略;enabled=false 是显式停发,不能被省略值自动恢复为 true。
|
||
func (s *Service) UpdateGameWinBroadcastConfig(ctx context.Context, enabled bool, minWinCoins int64, operatorAdminID int64) (broadcastdomain.GameWinConfig, error) {
|
||
if s == nil || s.gameWinConfigRepo == nil {
|
||
return broadcastdomain.GameWinConfig{}, xerr.New(xerr.Unavailable, "game win broadcast config repository is not configured")
|
||
}
|
||
if minWinCoins <= 0 || minWinCoins > maxGameWinMinCoins {
|
||
return broadcastdomain.GameWinConfig{}, xerr.New(xerr.InvalidArgument, "min_win_coins must be between 1 and 1000000000000")
|
||
}
|
||
if operatorAdminID < 0 {
|
||
return broadcastdomain.GameWinConfig{}, xerr.New(xerr.InvalidArgument, "operator_admin_id is invalid")
|
||
}
|
||
nowMS := s.now().UTC().UnixMilli()
|
||
stored, err := s.gameWinConfigRepo.UpsertGameWinBroadcastConfig(ctx, broadcastdomain.GameWinConfig{
|
||
Enabled: enabled,
|
||
MinWinCoins: minWinCoins,
|
||
UpdatedByAdminID: operatorAdminID,
|
||
CreatedAtMS: nowMS,
|
||
UpdatedAtMS: nowMS,
|
||
})
|
||
if err != nil {
|
||
return broadcastdomain.GameWinConfig{}, err
|
||
}
|
||
s.storeCachedGameWinConfig(appcode.FromContext(ctx), stored)
|
||
return stored, nil
|
||
}
|
||
|
||
func (s *Service) cachedGameWinBroadcastConfig(ctx context.Context) (broadcastdomain.GameWinConfig, error) {
|
||
app := appcode.FromContext(ctx)
|
||
nowMS := s.now().UTC().UnixMilli()
|
||
s.gameWinConfigMu.RLock()
|
||
cached, ok := s.gameWinConfigCache[app]
|
||
s.gameWinConfigMu.RUnlock()
|
||
if ok && cached.expiresAtMS > nowMS {
|
||
return cached.config, nil
|
||
}
|
||
return s.GetGameWinBroadcastConfig(ctx)
|
||
}
|
||
|
||
func (s *Service) storeCachedGameWinConfig(app string, config broadcastdomain.GameWinConfig) {
|
||
if s == nil {
|
||
return
|
||
}
|
||
app = appcode.Normalize(app)
|
||
s.gameWinConfigMu.Lock()
|
||
s.gameWinConfigCache[app] = cachedGameWinConfig{
|
||
config: config,
|
||
expiresAtMS: s.now().UTC().Add(s.cfg.GameWinConfigCacheTTL).UnixMilli(),
|
||
}
|
||
s.gameWinConfigMu.Unlock()
|
||
}
|
||
|
||
func gameWinPayload(message gamemq.GameOutboxMessage, source map[string]any, eventID string, regionID int64, profile SenderProfile, roomID string, room RoomBroadcastSnapshot, sentAtMS int64) (string, error) {
|
||
gameID := strings.TrimSpace(message.GameID)
|
||
providerGameID := firstGamePayloadString(source, "provider_game_id")
|
||
gameName := firstGamePayloadString(source, "game_name")
|
||
gameIconURL := firstGamePayloadString(source, "game_icon_url", "game_cover_url")
|
||
roundID := firstGamePayloadString(source, "provider_round_id", "round_id")
|
||
payload := map[string]any{
|
||
"event_id": eventID,
|
||
"source_event_id": message.EventID,
|
||
"broadcast_type": broadcastdomain.TypeGameWin,
|
||
"scope": broadcastdomain.ScopeRegion,
|
||
"app_code": message.AppCode,
|
||
"region_id": regionID,
|
||
"visible_region_id": regionID,
|
||
"order_id": message.OrderID,
|
||
"platform_code": message.PlatformCode,
|
||
"game_id": gameID,
|
||
"provider_game_id": providerGameID,
|
||
"game_name": gameName,
|
||
"game_icon_url": gameIconURL,
|
||
"game_round_id": roundID,
|
||
"round_id": roundID,
|
||
"room_id": roomID,
|
||
"room_locked": room.Locked,
|
||
"room_name": room.Title,
|
||
"user_id": message.UserID,
|
||
"sender_user_id": message.UserID,
|
||
"account": profile.Account,
|
||
"actual_account": profile.Account,
|
||
"user_name": profile.Nickname,
|
||
"user_nickname": profile.Nickname,
|
||
"user_avatar_url": profile.Avatar,
|
||
"coins": message.CoinAmount,
|
||
"win_amount": message.CoinAmount,
|
||
"award_amount": message.CoinAmount,
|
||
"currency_diff": message.CoinAmount,
|
||
"wallet_transaction_id": firstGamePayloadString(source, "wallet_transaction_id"),
|
||
"occurred_at_ms": message.OccurredAtMS,
|
||
"sent_at_ms": sentAtMS,
|
||
"user": map[string]any{
|
||
"user_id": profile.UserID,
|
||
"display_user_id": profile.Account,
|
||
"nickname": profile.Nickname,
|
||
"avatar": profile.Avatar,
|
||
},
|
||
}
|
||
if roomID != "" {
|
||
payload["jump_type"] = "room_window"
|
||
payload["jump_payload"] = map[string]any{
|
||
"room_id": roomID,
|
||
"room_locked": room.Locked,
|
||
"window": "game_window",
|
||
"game_id": gameID,
|
||
}
|
||
payload["action"] = map[string]any{
|
||
"type": "room_window",
|
||
"room_id": roomID,
|
||
"room_locked": room.Locked,
|
||
"window": "game_window",
|
||
"game_id": gameID,
|
||
}
|
||
}
|
||
encoded, err := json.Marshal(payload)
|
||
return string(encoded), err
|
||
}
|
||
|
||
func decodeGameWinSourcePayload(raw string) (map[string]any, error) {
|
||
if strings.TrimSpace(raw) == "" {
|
||
return map[string]any{}, nil
|
||
}
|
||
var payload map[string]any
|
||
if err := json.Unmarshal([]byte(raw), &payload); err != nil {
|
||
return nil, fmt.Errorf("decode game order payload: %w", err)
|
||
}
|
||
if payload == nil {
|
||
payload = map[string]any{}
|
||
}
|
||
return payload, nil
|
||
}
|
||
|
||
func firstGamePayloadString(payload map[string]any, keys ...string) string {
|
||
for _, key := range keys {
|
||
value, exists := payload[key]
|
||
if !exists || value == nil {
|
||
continue
|
||
}
|
||
if normalized := strings.TrimSpace(fmt.Sprint(value)); normalized != "" && normalized != "<nil>" {
|
||
return normalized
|
||
}
|
||
}
|
||
return ""
|
||
}
|
||
|
||
func firstGamePayloadInt64(payload map[string]any, keys ...string) int64 {
|
||
for _, key := range keys {
|
||
value, exists := payload[key]
|
||
if !exists {
|
||
continue
|
||
}
|
||
switch typed := value.(type) {
|
||
case float64:
|
||
return int64(typed)
|
||
case float32:
|
||
return int64(typed)
|
||
case int64:
|
||
return typed
|
||
case int:
|
||
return int64(typed)
|
||
case json.Number:
|
||
if parsed, err := typed.Int64(); err == nil {
|
||
return parsed
|
||
}
|
||
case string:
|
||
if parsed, err := strconv.ParseInt(strings.TrimSpace(typed), 10, 64); err == nil {
|
||
return parsed
|
||
}
|
||
}
|
||
}
|
||
return 0
|
||
}
|