package mysql import ( "context" "time" "hyapp/pkg/appcode" "hyapp/pkg/xerr" "hyapp/services/wallet-service/internal/domain/ledger" ) // GetHostRevenueStats 合并两种 Host 政策各自的送礼事实:SALARY_DIAMOND 保持读取周期钻石流水, // POINT_DIAMOND 只读取事务内角色拆分投影。它不读取钱包余额或通用分录,因此任务奖励、Agency // 分成、兑换、提现和退回都不会被误算成主播收礼收益。 func (r *Repository) GetHostRevenueStats(ctx context.Context, query ledger.HostRevenueStatsQuery) (ledger.HostRevenueStats, error) { if r == nil || r.db == nil { return ledger.HostRevenueStats{}, xerr.New(xerr.Unavailable, "mysql repository is not configured") } var stats ledger.HostRevenueStats startCycle := time.UnixMilli(query.StartAtMS).UTC().Format("2006-01") endCycle := time.UnixMilli(query.EndAtMS - 1).UTC().Format("2006-01") // 两个 UNION 分支分别命中 (app,user,cycle,time) 与 (app,host,time);外层只对已经按 // 用户/时间收敛的行去重 sender。这样政策在查询区间内切换时,同一送礼人也只计一次。 if err := r.db.QueryRowContext(ctx, ` SELECT COALESCE(SUM(CASE WHEN income_type = 'SALARY_DIAMOND' THEN host_base_amount ELSE 0 END), 0), COALESCE(SUM(CASE WHEN income_type = 'POINT_DIAMOND' THEN host_base_amount ELSE 0 END), 0), COUNT(DISTINCT CASE WHEN sender_user_id > 0 THEN sender_user_id END) FROM ( SELECT diamond_delta AS host_base_amount, sender_user_id, 'SALARY_DIAMOND' AS income_type FROM host_period_diamond_entries FORCE INDEX (idx_host_period_diamond_entries_user_cycle) WHERE app_code = ? AND user_id = ? AND cycle_key BETWEEN ? AND ? AND created_at_ms >= ? AND created_at_ms < ? UNION ALL SELECT host_base_amount, sender_user_id, 'POINT_DIAMOND' AS income_type FROM point_diamond_gift_income_entries FORCE INDEX (idx_point_diamond_gift_host_time) WHERE app_code = ? AND host_user_id = ? AND created_at_ms >= ? AND created_at_ms < ? ) AS gift_income`, appcode.FromContext(ctx), query.HostUserID, startCycle, endCycle, query.StartAtMS, query.EndAtMS, appcode.FromContext(ctx), query.HostUserID, query.StartAtMS, query.EndAtMS, ).Scan(&stats.DiamondEarnings, &stats.PointDiamondEarnings, &stats.GiftSenders); err != nil { return ledger.HostRevenueStats{}, err } // 旧 Lalu/POINT 与永久钻石积分分别聚合;旧字段绝不能混入任务 POINT 以外的新资产, // POINT_DIAMOND 在政策切回工资后仍能通过独立字段保留历史转出。 if err := r.db.QueryRowContext(ctx, ` SELECT COALESCE(SUM(CASE WHEN e.asset_type = ? THEN -e.available_delta ELSE 0 END), 0), COALESCE(SUM(CASE WHEN e.asset_type = ? THEN -e.available_delta ELSE 0 END), 0) FROM wallet_entries e FORCE INDEX (idx_wallet_entries_asset_user_time) INNER JOIN wallet_transactions t ON t.app_code = e.app_code AND t.transaction_id = e.transaction_id WHERE e.app_code = ? AND e.user_id = ? AND e.asset_type IN (?, ?) AND e.created_at_ms >= ? AND e.created_at_ms < ? AND e.available_delta < 0 AND t.biz_type IN (?, ?, ?)`, ledger.AssetPoint, ledger.AssetPointDiamond, appcode.FromContext(ctx), query.HostUserID, ledger.AssetPoint, ledger.AssetPointDiamond, query.StartAtMS, query.EndAtMS, bizTypePointExchangeToCoin, bizTypePointTransferToCoinSeller, bizTypeSalaryWithdrawalFreeze, ).Scan(&stats.DiamondExchanged, &stats.PointDiamondExchanged); err != nil { return ledger.HostRevenueStats{}, err } return stats, nil }