2026-07-22 13:55:55 +08:00

215 lines
11 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

package mysql
import (
"context"
"database/sql"
"encoding/json"
"errors"
"fmt"
"strings"
"time"
"hyapp/pkg/appcode"
"hyapp/pkg/xerr"
"hyapp/services/wallet-service/internal/domain/ledger"
)
type pointToCoinExchangeMetadata struct {
AppCode string `json:"app_code"`
UserID int64 `json:"user_id"`
PointAmount int64 `json:"point_amount"`
SourceAssetType string `json:"source_asset_type"`
CoinAmount int64 `json:"coin_amount"`
RatioPointAmount int64 `json:"ratio_point_amount"`
RatioCoinAmount int64 `json:"ratio_coin_amount"`
PointBalanceAfter int64 `json:"point_balance_after"`
CoinBalanceAfter int64 `json:"coin_balance_after"`
CreatedAtMS int64 `json:"created_at_ms"`
PolicyID uint64 `json:"policy_id"`
PolicyVersion uint64 `json:"policy_version"`
}
// ExchangePointToCoin 在同一事务中解析 active POINT 政策并锁定用户 POINT/COIN 账户;
// 客户端不能提交 coin_amount避免篡改比例或在后台改价后重放得到不同结果。
func (r *Repository) ExchangePointToCoin(ctx context.Context, command ledger.PointToCoinCommand) (ledger.PointToCoinReceipt, error) {
if r == nil || r.db == nil {
return ledger.PointToCoinReceipt{}, xerr.New(xerr.Unavailable, "mysql repository is not configured")
}
ctx = contextWithCommandApp(ctx, command.AppCode)
command.AppCode = appcode.FromContext(ctx)
tx, err := r.db.BeginTx(ctx, nil)
if err != nil {
return ledger.PointToCoinReceipt{}, err
}
defer func() { _ = tx.Rollback() }()
sourceAssetType := strings.ToUpper(strings.TrimSpace(command.SourceAssetType))
if sourceAssetType == "" {
sourceAssetType = ledger.AssetPoint
}
if sourceAssetType != ledger.AssetPoint && sourceAssetType != ledger.AssetPointDiamond {
return ledger.PointToCoinReceipt{}, xerr.New(xerr.InvalidArgument, "point source asset type is invalid")
}
// region 是服务端从当前用户资料派生的政策定位,不属于客户端稳定意图。同一 command_id
// 在用户区域变化后仍应返回首次回执;用户、资产和兑换数量变化仍构成幂等冲突。
requestHash := stableHash(fmt.Sprintf("point_exchange_to_coin|%s|%d|%s|%d", command.AppCode, command.UserID, sourceAssetType, command.PointAmount))
if txRow, exists, lookupErr := r.lookupTransactionWithConflictCode(ctx, tx, command.CommandID, requestHash, bizTypePointExchangeToCoin, xerr.IdempotencyConflict); lookupErr != nil || exists {
if lookupErr != nil || !exists {
return ledger.PointToCoinReceipt{}, lookupErr
}
return r.receiptForPointToCoinExchange(ctx, tx, txRow.TransactionID)
}
// POINT 分母与 USDT 提现读取同一份已发布 wallet 政策COIN 分子复用工资兑金币的统一产品比例。
// 在账变事务内重新解析,避免页面预览后 Admin 改政策导致按旧分母扣账。
var ratioPointAmount, ratioCoinAmount int64
var policyID, policyVersion uint64
if sourceAssetType == ledger.AssetPointDiamond {
policy, err := r.resolvePointDiamondHostPolicy(ctx, tx, command.RegionID, command.NowMS)
if err != nil {
return ledger.PointToCoinReceipt{}, err
}
ratioPointAmount, ratioCoinAmount = policy.PointDiamondsPerUSD, policy.CoinsPerUSD
policyID, policyVersion = policy.PolicyID, policy.PolicyVersion
} else {
policy, found, err := r.resolveActiveWalletPolicy(ctx, tx, command.AppCode, command.RegionID, command.NowMS)
if err != nil {
return ledger.PointToCoinReceipt{}, err
}
if !found {
return ledger.PointToCoinReceipt{}, xerr.New(xerr.PermissionDenied, "point wallet policy is not configured")
}
ratioPointAmount, ratioCoinAmount = policy.PointsPerUSD, ledger.SalaryExchangeCoinPerUSD
}
numerator, err := checkedMul(command.PointAmount, ratioCoinAmount)
if err != nil {
return ledger.PointToCoinReceipt{}, err
}
coinAmount := numerator / ratioPointAmount
if coinAmount <= 0 {
return ledger.PointToCoinReceipt{}, xerr.New(xerr.InvalidArgument, "point amount is below exchange precision")
}
nowMS := time.Now().UTC().UnixMilli()
states, err := r.lockAccountsOrdered(ctx, tx, []walletAccountLockRequest{
{UserID: command.UserID, AssetType: sourceAssetType, CreateIfMissing: false},
{UserID: command.UserID, AssetType: ledger.AssetCoin, CreateIfMissing: true},
}, nowMS)
if err != nil {
return ledger.PointToCoinReceipt{}, err
}
pointAccount := states[walletAccountLockKey(command.UserID, sourceAssetType)].account
if pointAccount.AvailableAmount < command.PointAmount {
return ledger.PointToCoinReceipt{}, xerr.New(xerr.InsufficientBalance, "insufficient point balance")
}
coinAccount := states[walletAccountLockKey(command.UserID, ledger.AssetCoin)].account
pointAfter := pointAccount.AvailableAmount - command.PointAmount
coinAfter, err := checkedAdd(coinAccount.AvailableAmount, coinAmount)
if err != nil {
return ledger.PointToCoinReceipt{}, err
}
transactionID := transactionID(command.AppCode, command.CommandID)
metadata := pointToCoinExchangeMetadata{
AppCode: command.AppCode, UserID: command.UserID, PointAmount: command.PointAmount, SourceAssetType: sourceAssetType, CoinAmount: coinAmount,
RatioPointAmount: ratioPointAmount, RatioCoinAmount: ratioCoinAmount,
PointBalanceAfter: pointAfter, CoinBalanceAfter: coinAfter, CreatedAtMS: nowMS,
PolicyID: policyID, PolicyVersion: policyVersion,
}
if err := r.insertTransaction(ctx, tx, transactionID, command.CommandID, bizTypePointExchangeToCoin, requestHash, fmt.Sprintf("point_exchange:%d", command.UserID), metadata, nowMS); err != nil {
return ledger.PointToCoinReceipt{}, err
}
if err := r.applyAccountDelta(ctx, tx, pointAccount, -command.PointAmount, 0, nowMS); err != nil {
return ledger.PointToCoinReceipt{}, err
}
if err := r.insertEntry(ctx, tx, walletEntry{TransactionID: transactionID, UserID: command.UserID, AssetType: sourceAssetType, AvailableDelta: -command.PointAmount, AvailableAfter: pointAfter, FrozenAfter: pointAccount.FrozenAmount, CreatedAtMS: nowMS}); err != nil {
return ledger.PointToCoinReceipt{}, err
}
if err := r.applyAccountDelta(ctx, tx, coinAccount, coinAmount, 0, nowMS); err != nil {
return ledger.PointToCoinReceipt{}, err
}
if err := r.insertEntry(ctx, tx, walletEntry{TransactionID: transactionID, UserID: command.UserID, AssetType: ledger.AssetCoin, AvailableDelta: coinAmount, AvailableAfter: coinAfter, FrozenAfter: coinAccount.FrozenAmount, CreatedAtMS: nowMS}); err != nil {
return ledger.PointToCoinReceipt{}, err
}
if err := r.insertWalletOutbox(ctx, tx, []walletOutboxEvent{
balanceChangedEvent(transactionID, command.CommandID, command.UserID, sourceAssetType, -command.PointAmount, 0, pointAfter, pointAccount.FrozenAmount, pointAccount.Version+1, metadata, nowMS, bizTypePointExchangeToCoin),
balanceChangedEvent(transactionID, command.CommandID, command.UserID, ledger.AssetCoin, coinAmount, 0, coinAfter, coinAccount.FrozenAmount, coinAccount.Version+1, metadata, nowMS, bizTypePointExchangeToCoin),
}); err != nil {
return ledger.PointToCoinReceipt{}, err
}
if err := tx.Commit(); err != nil {
return ledger.PointToCoinReceipt{}, err
}
return receiptFromPointToCoinExchangeMetadata(transactionID, metadata), nil
}
func (r *Repository) receiptForPointToCoinExchange(ctx context.Context, tx *sql.Tx, transactionID string) (ledger.PointToCoinReceipt, error) {
var metadataJSON string
if err := tx.QueryRowContext(ctx, `SELECT COALESCE(CAST(metadata_json AS CHAR), '{}') FROM wallet_transactions WHERE app_code = ? AND transaction_id = ?`, appcode.FromContext(ctx), transactionID).Scan(&metadataJSON); err != nil {
return ledger.PointToCoinReceipt{}, err
}
var metadata pointToCoinExchangeMetadata
if err := json.Unmarshal([]byte(metadataJSON), &metadata); err != nil {
return ledger.PointToCoinReceipt{}, err
}
return receiptFromPointToCoinExchangeMetadata(transactionID, metadata), nil
}
func receiptFromPointToCoinExchangeMetadata(transactionID string, metadata pointToCoinExchangeMetadata) ledger.PointToCoinReceipt {
return ledger.PointToCoinReceipt{
TransactionID: transactionID, UserID: metadata.UserID, PointAmount: metadata.PointAmount, CoinAmount: metadata.CoinAmount,
PointBalanceAfter: metadata.PointBalanceAfter, CoinBalanceAfter: metadata.CoinBalanceAfter,
RatioPointAmount: metadata.RatioPointAmount, RatioCoinAmount: metadata.RatioCoinAmount, CreatedAtMS: metadata.CreatedAtMS,
SourceAssetType: metadata.SourceAssetType, PolicyID: metadata.PolicyID, PolicyVersion: metadata.PolicyVersion,
}
}
// resolvePointDiamondHostPolicy 读取当前区域不晚于本月的最近一份 POINT_DIAMOND 发布快照。
// 永久积分不会随周期清空,因此当前政策切回工资型后,旧余额仍必须按最后一份积分政策兑换和退出。
func (r *Repository) resolvePointDiamondHostPolicy(ctx context.Context, tx *sql.Tx, regionID int64, nowMS int64) (ledger.HostSalaryPolicy, error) {
if nowMS <= 0 {
nowMS = time.Now().UTC().UnixMilli()
}
cycleKey := time.UnixMilli(nowMS).UTC().Format("2006-01")
var policy ledger.HostSalaryPolicy
// 已发布政策快照不可修改,资金动作只需要在自身事务的一致性视图中读取并把版本写入回执;
// 对这行加 FOR UPDATE 不增加正确性,反而会让同一区域的 overview 和兑换/提现互相串行。
err := tx.QueryRowContext(ctx, `
SELECT p.policy_id, p.policy_version, binding.cycle_key, binding.region_id,
p.point_diamonds_per_usd, p.coins_per_usd, p.minimum_withdraw_usd_minor,
p.withdraw_fee_bps, p.agency_point_share_bps,
p.coin_seller_withdrawal_limit_period, p.coin_seller_withdrawal_limit_count,
p.platform_withdrawal_limit_period, p.platform_withdrawal_limit_count,
p.platform_withdrawal_allowed_days
FROM host_salary_policy_cycle_bindings binding
JOIN host_agency_salary_policies p
ON p.app_code = binding.app_code
AND p.policy_id = binding.policy_id
AND p.policy_version = binding.policy_version
WHERE binding.app_code = ?
AND binding.region_id = ?
AND binding.cycle_key <= ?
AND p.status = 'active'
AND p.policy_type = ?
ORDER BY binding.cycle_key DESC, p.policy_version DESC
LIMIT 1`, appcode.FromContext(ctx), regionID, cycleKey, ledger.HostPolicyTypePointDiamond).Scan(
&policy.PolicyID, &policy.PolicyVersion, &policy.CycleKey, &policy.RegionID,
&policy.PointDiamondsPerUSD, &policy.CoinsPerUSD, &policy.MinimumWithdrawUSDMinor,
&policy.WithdrawFeeBPS, &policy.AgencyPointShareBPS,
&policy.CoinSellerWithdrawalLimitPeriod, &policy.CoinSellerWithdrawalLimitCount,
&policy.PlatformWithdrawalLimitPeriod, &policy.PlatformWithdrawalLimitCount,
&policy.PlatformWithdrawalAllowedDays,
)
if errors.Is(err, sql.ErrNoRows) {
return ledger.HostSalaryPolicy{}, xerr.New(xerr.PermissionDenied, "point diamond policy is not configured")
}
if err != nil {
return ledger.HostSalaryPolicy{}, err
}
policy.PolicyType = ledger.HostPolicyTypePointDiamond
if policy.PointDiamondsPerUSD <= 0 || policy.CoinsPerUSD <= 0 || policy.MinimumWithdrawUSDMinor <= 0 || policy.WithdrawFeeBPS < 0 || policy.WithdrawFeeBPS > 10_000 {
return ledger.HostSalaryPolicy{}, xerr.New(xerr.Internal, "point diamond policy is invalid")
}
return policy, nil
}