91 lines
3.2 KiB
Go
91 lines
3.2 KiB
Go
package service
|
||
|
||
import (
|
||
"context"
|
||
"encoding/json"
|
||
"fmt"
|
||
"strings"
|
||
|
||
"hyapp/pkg/appcode"
|
||
"hyapp/pkg/walletmq"
|
||
"hyapp/pkg/xerr"
|
||
)
|
||
|
||
// HandleWalletRedPacketProjectionEvent 消费 wallet realtime outbox,异步刷新房间列表红包聚合。
|
||
// 该链路不进入 ListRooms 请求,也不读取 wallet 数据库;列表延迟只增加同一行两个标量字段的反序列化成本。
|
||
func (s *Service) HandleWalletRedPacketProjectionEvent(ctx context.Context, message walletmq.WalletOutboxMessage) error {
|
||
event, handled, err := roomRedPacketProjectionEvent(message, s.clock.Now().UnixMilli())
|
||
if err != nil || !handled {
|
||
return err
|
||
}
|
||
store, ok := s.repository.(RoomRedPacketProjectionStore)
|
||
if !ok {
|
||
return xerr.New(xerr.Unavailable, "room red packet projection store is unavailable")
|
||
}
|
||
ctx = appcode.WithContext(ctx, message.AppCode)
|
||
entry, exists, err := store.ApplyRoomRedPacketProjectionEvent(ctx, event)
|
||
if err != nil || !exists {
|
||
return err
|
||
}
|
||
// MySQL 事务是权威投影;Redis 失败只影响短期命中率,MQ 重投和定时重建都会继续收敛。
|
||
s.projectRoomListCacheBestEffort(ctx, entry)
|
||
return nil
|
||
}
|
||
|
||
func roomRedPacketProjectionEvent(message walletmq.WalletOutboxMessage, projectedAtMS int64) (RoomRedPacketProjectionEvent, bool, error) {
|
||
eventType := strings.TrimSpace(message.EventType)
|
||
eventRank := int32(0)
|
||
switch eventType {
|
||
case walletmq.EventTypeWalletRedPacketCreated:
|
||
eventRank = 1
|
||
case walletmq.EventTypeWalletRedPacketClaimed:
|
||
eventRank = 2
|
||
case walletmq.EventTypeWalletRedPacketRefunded:
|
||
eventRank = 3
|
||
default:
|
||
return RoomRedPacketProjectionEvent{}, false, nil
|
||
}
|
||
|
||
var payload struct {
|
||
PacketID string `json:"packet_id"`
|
||
RoomID string `json:"room_id"`
|
||
Status string `json:"status"`
|
||
PacketStatus string `json:"packet_status"`
|
||
RemainingAmount int64 `json:"remaining_amount"`
|
||
RemainingCount int32 `json:"remaining_count"`
|
||
ExpiresAtMS int64 `json:"expires_at_ms"`
|
||
}
|
||
if err := json.Unmarshal([]byte(message.PayloadJSON), &payload); err != nil {
|
||
return RoomRedPacketProjectionEvent{}, true, err
|
||
}
|
||
status := strings.ToLower(strings.TrimSpace(payload.Status))
|
||
if payload.PacketStatus != "" {
|
||
status = strings.ToLower(strings.TrimSpace(payload.PacketStatus))
|
||
}
|
||
if eventType == walletmq.EventTypeWalletRedPacketRefunded {
|
||
status = "refunded"
|
||
payload.RemainingAmount = 0
|
||
payload.RemainingCount = 0
|
||
}
|
||
if strings.TrimSpace(message.EventID) == "" || strings.TrimSpace(payload.PacketID) == "" ||
|
||
strings.TrimSpace(payload.RoomID) == "" || status == "" || message.OccurredAtMS <= 0 {
|
||
return RoomRedPacketProjectionEvent{}, true, fmt.Errorf("%s projection payload is incomplete", eventType)
|
||
}
|
||
if projectedAtMS <= 0 {
|
||
projectedAtMS = message.OccurredAtMS
|
||
}
|
||
return RoomRedPacketProjectionEvent{
|
||
EventID: strings.TrimSpace(message.EventID),
|
||
EventType: eventType,
|
||
PacketID: strings.TrimSpace(payload.PacketID),
|
||
RoomID: strings.TrimSpace(payload.RoomID),
|
||
Status: status,
|
||
RemainingAmount: payload.RemainingAmount,
|
||
RemainingCount: payload.RemainingCount,
|
||
ExpiresAtMS: payload.ExpiresAtMS,
|
||
OccurredAtMS: message.OccurredAtMS,
|
||
ProjectedAtMS: projectedAtMS,
|
||
EventRank: eventRank,
|
||
}, true, nil
|
||
}
|