67 lines
3.5 KiB
Go
67 lines
3.5 KiB
Go
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
|
||
}
|