package app import ( "encoding/json" "strings" "hyapp/pkg/appcode" "hyapp/pkg/userleaderboard" "hyapp/pkg/walletmq" broadcastdomain "hyapp/services/activity-service/internal/domain/broadcast" cumulativerechargedomain "hyapp/services/activity-service/internal/domain/cumulativerecharge" firstrechargedomain "hyapp/services/activity-service/internal/domain/firstrecharge" inviteactivitydomain "hyapp/services/activity-service/internal/domain/inviteactivity" ) func firstRechargeEventFromWalletMessage(body []byte) (firstrechargedomain.RechargeEvent, bool, error) { message, err := walletmq.DecodeWalletOutboxMessage(body) if err != nil { return firstrechargedomain.RechargeEvent{}, false, err } if message.EventType != walletmq.EventTypeWalletRechargeRecorded { return firstrechargedomain.RechargeEvent{}, false, nil } fields := firstRechargeFactFromWalletPayload(message.PayloadJSON) // 首充奖励按 Google 商品 ID 精确匹配运营档位;币商转账、H5 或其他充值事实没有 Google 商品维度, // 这里必须在进入业务校验前确认并跳过,避免合法钱包事实被首充 consumer 当成 poison message 反复重投。 if !isCumulativeGoogleRecharge(fields.rechargeType) { return firstrechargedomain.RechargeEvent{}, false, nil } return firstrechargedomain.RechargeEvent{ AppCode: appcode.Normalize(message.AppCode), EventID: message.EventID, EventType: message.EventType, TransactionID: message.TransactionID, CommandID: message.CommandID, UserID: message.UserID, RechargeCoinAmount: message.AvailableDelta, RechargeSequence: fields.sequence, RechargeType: fields.rechargeType, RechargeUSDMinor: fields.rechargeUSDMinor, GoogleProductID: fields.googleProductID, PayloadJSON: message.PayloadJSON, OccurredAtMS: message.OccurredAtMS, }, true, nil } func cumulativeRechargeEventFromWalletMessage(body []byte) (cumulativerechargedomain.RechargeEvent, bool, error) { message, err := walletmq.DecodeWalletOutboxMessage(body) if err != nil { return cumulativerechargedomain.RechargeEvent{}, false, err } if message.EventType != walletmq.EventTypeWalletRechargeRecorded { return cumulativerechargedomain.RechargeEvent{}, false, nil } // WalletRechargeRecorded 是累充唯一事实源;payload 里没有可折算 USD 金额时直接跳过,避免把普通钱包流水计入活动。 fields := cumulativeRechargeFactFromWalletPayload(message.PayloadJSON, message.AvailableDelta) if fields.qualifyingUSDMinor <= 0 { return cumulativerechargedomain.RechargeEvent{}, false, nil } return cumulativerechargedomain.RechargeEvent{ AppCode: appcode.Normalize(message.AppCode), EventID: message.EventID, EventType: message.EventType, TransactionID: message.TransactionID, CommandID: message.CommandID, UserID: message.UserID, RechargeCoinAmount: fields.coinAmount, RechargeSequence: fields.sequence, RechargeType: fields.rechargeType, QualifyingUSDMinor: fields.qualifyingUSDMinor, PayloadJSON: message.PayloadJSON, OccurredAtMS: message.OccurredAtMS, }, true, nil } func inviteActivityRechargeEventFromWalletMessage(body []byte) (inviteactivitydomain.RechargeEvent, bool, error) { message, err := walletmq.DecodeWalletOutboxMessage(body) if err != nil { return inviteactivitydomain.RechargeEvent{}, false, err } if message.EventType != walletmq.EventTypeWalletRechargeRecorded || message.AssetType != "COIN" || message.AvailableDelta <= 0 { return inviteactivitydomain.RechargeEvent{}, false, nil } fields := cumulativeRechargeFactFromWalletPayload(message.PayloadJSON, message.AvailableDelta) return inviteactivitydomain.RechargeEvent{ AppCode: appcode.Normalize(message.AppCode), EventID: message.EventID, EventType: message.EventType, TransactionID: message.TransactionID, CommandID: message.CommandID, InvitedUserID: message.UserID, RechargeCoinAmount: fields.coinAmount, RechargeType: fields.rechargeType, PayloadJSON: message.PayloadJSON, OccurredAtMS: message.OccurredAtMS, }, true, nil } func redPacketEventFromWalletMessage(body []byte) (broadcastdomain.RedPacketWalletEvent, bool, error) { message, err := walletmq.DecodeWalletOutboxMessage(body) if err != nil { return broadcastdomain.RedPacketWalletEvent{}, false, err } if message.EventType != walletmq.EventTypeWalletRedPacketCreated && message.EventType != walletmq.EventTypeWalletRedPacketClaimed && message.EventType != walletmq.EventTypeWalletRedPacketRefunded { return broadcastdomain.RedPacketWalletEvent{}, false, nil } fields := redPacketFieldsFromWalletPayload(message.PayloadJSON, message.EventType) if fields.PacketID == "" { return broadcastdomain.RedPacketWalletEvent{}, false, nil } return broadcastdomain.RedPacketWalletEvent{ AppCode: appcode.Normalize(message.AppCode), EventID: message.EventID, EventType: message.EventType, TransactionID: message.TransactionID, CommandID: message.CommandID, PayloadJSON: message.PayloadJSON, PacketID: fields.PacketID, PacketType: fields.PacketType, RoomID: fields.RoomID, RegionID: fields.RegionID, RegionCode: fields.RegionCode, RoomLocked: fields.RoomLocked, OwnerUserID: fields.OwnerUserID, SenderUserID: fields.SenderUserID, TotalAmount: fields.TotalAmount, TotalCount: fields.TotalCount, RemainingAmount: fields.RemainingAmount, RemainingCount: fields.RemainingCount, RefundedAmount: fields.RefundedAmount, Status: fields.Status, OpenAtMS: fields.OpenAtMS, ExpiresAtMS: fields.ExpiresAtMS, CreatedAtMS: message.OccurredAtMS, }, true, nil } func userLeaderboardGiftEventFromWalletMessage(body []byte) (userleaderboard.GiftEvent, bool, error) { message, err := walletmq.DecodeWalletOutboxMessage(body) if err != nil { return userleaderboard.GiftEvent{}, false, err } if message.EventType != walletmq.EventTypeWalletGiftDebited { return userleaderboard.GiftEvent{}, false, nil } var decoded map[string]any if err := json.Unmarshal([]byte(message.PayloadJSON), &decoded); err != nil { return userleaderboard.GiftEvent{}, false, err } if boolFromDecoded(decoded, "direct_gift") { // 旧 gateway SQL 只统计 biz_type='gift_debit',direct_gift_debit 不进榜;Redis 投影必须保持这个口径。 return userleaderboard.GiftEvent{}, false, nil } senderUserID := firstNonZeroInt64(int64FromDecoded(decoded, "sender_user_id"), message.UserID) targetUserID := int64FromDecoded(decoded, "target_user_id") giftValue, _ := firstPresentInt64(decoded, "heat_value", "coin_spent", "charge_amount") if senderUserID <= 0 || targetUserID <= 0 || giftValue <= 0 { return userleaderboard.GiftEvent{}, false, nil } return userleaderboard.GiftEvent{ AppCode: appcode.Normalize(message.AppCode), EventID: message.EventID, SenderUserID: senderUserID, SenderRegionID: int64FromDecoded(decoded, "sender_region_id"), TargetUserID: targetUserID, // 主播身份区域是魅力榜的稳定归属;普通收礼用户没有主播快照时,回退送礼发生房间的可见区域。 TargetRegionID: firstNonZeroInt64(int64FromDecoded(decoded, "target_host_region_id"), int64FromDecoded(decoded, "room_region_id")), GiftValue: giftValue, GiftCount: int64FromDecoded(decoded, "gift_count"), OccurredAtMS: message.OccurredAtMS, }, true, nil } type firstRechargePayloadFields struct { sequence int64 rechargeType string rechargeUSDMinor int64 googleProductID string } func firstRechargeFactFromWalletPayload(payload string) firstRechargePayloadFields { var decoded map[string]any if err := json.Unmarshal([]byte(payload), &decoded); err != nil { decoded = map[string]any{} } rechargeType := firstNonEmptyString( stringFromDecoded(decoded, "recharge_type"), stringFromDecoded(decoded, "channel"), stringFromDecoded(decoded, "provider"), firstrechargedomain.RechargeTypeCoinSeller, ) rechargeUSDMinor := firstNonZeroInt64( int64FromDecoded(decoded, "recharge_usd_minor"), int64FromDecoded(decoded, "usd_minor_amount"), int64FromDecoded(decoded, "amount_usd_minor"), ) if rechargeUSDMinor <= 0 && isCumulativeGoogleRecharge(rechargeType) { rechargeUSDMinor = int64FromDecoded(decoded, "amount_micro") / 10000 } return firstRechargePayloadFields{ sequence: int64FromDecoded(decoded, "recharge_sequence"), rechargeType: rechargeType, rechargeUSDMinor: rechargeUSDMinor, googleProductID: firstNonEmptyString( stringFromDecoded(decoded, "google_product_id"), stringFromDecoded(decoded, "product_name"), stringFromDecoded(decoded, "product_code"), ), } } type cumulativeRechargePayloadFields struct { sequence int64 rechargeType string coinAmount int64 qualifyingUSDMinor int64 } func cumulativeRechargeFactFromWalletPayload(payload string, availableDelta int64) cumulativeRechargePayloadFields { var decoded map[string]any if err := json.Unmarshal([]byte(payload), &decoded); err != nil { decoded = map[string]any{} } // 来源字段兼容 wallet 历史 payload;没有明确来源时按币商转账处理,因为币商链路通常只带金币增量。 rechargeType := firstNonEmptyString( stringFromDecoded(decoded, "recharge_type"), stringFromDecoded(decoded, "channel"), stringFromDecoded(decoded, "provider"), cumulativerechargedomain.RechargeTypeCoinSeller, ) // coinAmount 用于审计和币商折算,优先使用 payload 明细,最后回退 wallet outbox 的 available_delta。 coinAmount := firstNonZeroInt64( int64FromDecoded(decoded, "coin_amount"), int64FromDecoded(decoded, "amount"), int64FromDecoded(decoded, "available_delta"), availableDelta, ) qualifyingUSDMinor := cumulativeRechargeUSDMinorFromPayload(decoded, rechargeType, coinAmount) return cumulativeRechargePayloadFields{ sequence: int64FromDecoded(decoded, "recharge_sequence"), rechargeType: rechargeType, coinAmount: coinAmount, qualifyingUSDMinor: qualifyingUSDMinor, } } func cumulativeRechargeUSDMinorFromPayload(decoded map[string]any, rechargeType string, coinAmount int64) int64 { if isCumulativeCoinSellerRecharge(rechargeType) { // 币商给用户转账按 90000 coins = 1 USD 折算;这里返回美分并向下取整,尾差只留在原始金币字段审计。 return coinAmount * 100 / 90000 } // Google/Mifapay 充值优先使用支付侧传入的美元美分,避免用金币包配置反推造成汇率误差。 if value := firstNonZeroInt64( int64FromDecoded(decoded, "recharge_usd_minor"), int64FromDecoded(decoded, "usd_minor_amount"), int64FromDecoded(decoded, "exchange_usd_minor_amount"), int64FromDecoded(decoded, "amount_usd_minor"), ); value > 0 { return value } if isCumulativeGoogleRecharge(rechargeType) { // Google Billing 常见 amount_micro 是 USD 微单位,除以 10000 后得到美分。 return int64FromDecoded(decoded, "amount_micro") / 10000 } return 0 } func isCumulativeCoinSellerRecharge(value string) bool { switch strings.ToLower(strings.TrimSpace(value)) { case "coin_seller_transfer", "coin_seller": return true default: return false } } func isCumulativeGoogleRecharge(value string) bool { switch strings.ToLower(strings.TrimSpace(value)) { case "google", "google_play": return true default: return false } } type redPacketPayloadFields struct { PacketID string PacketType string RoomID string RegionID int64 RegionCode string RoomLocked bool OwnerUserID int64 SenderUserID int64 TotalAmount int64 TotalCount int32 RemainingAmount int64 RemainingCount int32 RefundedAmount int64 Status string OpenAtMS int64 ExpiresAtMS int64 } func redPacketFieldsFromWalletPayload(payload string, eventType string) redPacketPayloadFields { var decoded map[string]any decoder := json.NewDecoder(strings.NewReader(payload)) // wallet payload 中的用户 ID 可能是 Snowflake int64;UseNumber 必须在首次解码时启用,先落成 float64 后无法恢复精度。 decoder.UseNumber() if err := decoder.Decode(&decoded); err != nil { return redPacketPayloadFields{} } packetID := stringFromDecoded(decoded, "packet_id") if packetID == "" { packetID = stringFromDecoded(decoded, "packetId") } packetType := stringFromDecoded(decoded, "packet_type") roomID := stringFromDecoded(decoded, "room_id") regionID := int64FromDecoded(decoded, "region_id") status := stringFromDecoded(decoded, "status") openAtMS := int64FromDecoded(decoded, "open_at_ms") expiresAtMS := int64FromDecoded(decoded, "expires_at_ms") if eventType == walletmq.EventTypeWalletRedPacketClaimed || eventType == walletmq.EventTypeWalletRedPacketRefunded { status = stringFromDecoded(decoded, "packet_status") } totalCount := int32FromDecoded(decoded, "packet_count") remainingCount := int32FromDecoded(decoded, "remaining_count") if remainingCount == 0 && eventType == walletmq.EventTypeWalletRedPacketCreated { remainingCount = totalCount } return redPacketPayloadFields{ PacketID: packetID, PacketType: packetType, RoomID: roomID, RegionID: regionID, RegionCode: stringFromDecoded(decoded, "region_code"), RoomLocked: boolFromDecoded(decoded, "room_locked"), OwnerUserID: int64FromDecoded(decoded, "owner_user_id"), SenderUserID: int64FromDecoded(decoded, "sender_user_id"), TotalAmount: int64FromDecoded(decoded, "total_amount"), TotalCount: totalCount, RemainingAmount: int64FromDecoded(decoded, "remaining_amount"), RemainingCount: remainingCount, RefundedAmount: int64FromDecoded(decoded, "refund_amount"), Status: status, OpenAtMS: openAtMS, ExpiresAtMS: expiresAtMS, } } func stringFromDecoded(decoded map[string]any, key string) string { value, _ := decoded[key].(string) return value } func firstNonEmptyString(values ...string) string { for _, value := range values { if value = strings.TrimSpace(value); value != "" { return value } } return "" } func firstNonZeroInt64(values ...int64) int64 { for _, value := range values { if value != 0 { return value } } return 0 } func int64FromDecoded(decoded map[string]any, key string) int64 { switch value := decoded[key].(type) { case float64: return int64(value) case int64: return value case json.Number: parsed, _ := value.Int64() return parsed case string: parsed, _ := json.Number(strings.TrimSpace(value)).Int64() return parsed default: return 0 } } func int32FromDecoded(decoded map[string]any, key string) int32 { return int32(int64FromDecoded(decoded, key)) } func firstPresentInt64(decoded map[string]any, keys ...string) (int64, bool) { for _, key := range keys { value, ok := int64FromDecodedIfPresent(decoded, key) if ok { return value, true } } return 0, false } func int64FromDecodedIfPresent(decoded map[string]any, key string) (int64, bool) { raw, exists := decoded[key] if !exists { return 0, false } switch value := raw.(type) { case float64: return int64(value), true case int64: return value, true case json.Number: parsed, err := value.Int64() return parsed, err == nil case string: value = strings.TrimSpace(value) if value == "" { return 0, false } parsed, err := json.Number(value).Int64() return parsed, err == nil default: return 0, false } } func boolFromDecoded(decoded map[string]any, key string) bool { switch value := decoded[key].(type) { case bool: return value case string: return strings.EqualFold(strings.TrimSpace(value), "true") || strings.TrimSpace(value) == "1" case float64: return value != 0 case json.Number: parsed, _ := value.Int64() return parsed != 0 default: return false } }