138 lines
6.8 KiB
Go
138 lines
6.8 KiB
Go
package mysql
|
||
|
||
import (
|
||
"context"
|
||
"database/sql"
|
||
"encoding/json"
|
||
"fmt"
|
||
"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"`
|
||
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"`
|
||
}
|
||
|
||
// ExchangePointToCoin 在同一事务中锁定 active 比例和用户 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() }()
|
||
|
||
requestHash := stableHash(fmt.Sprintf("point_exchange_to_coin|%s|%d|%d|%d", command.AppCode, command.UserID, command.PointAmount, command.RegionID))
|
||
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 提现读取同一份已发布政策;COIN 分子复用工资兑金币的统一产品比例。
|
||
// 在账变事务内重新解析,避免页面预览后 Admin 改政策导致按旧分母扣账。
|
||
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()
|
||
pointAccount, err := r.lockAccount(ctx, tx, command.UserID, ledger.AssetPoint, false, nowMS)
|
||
if err != nil {
|
||
return ledger.PointToCoinReceipt{}, err
|
||
}
|
||
if pointAccount.AvailableAmount < command.PointAmount {
|
||
return ledger.PointToCoinReceipt{}, xerr.New(xerr.InsufficientBalance, "insufficient POINT balance")
|
||
}
|
||
coinAccount, err := r.lockAccount(ctx, tx, command.UserID, ledger.AssetCoin, true, nowMS)
|
||
if err != nil {
|
||
return ledger.PointToCoinReceipt{}, err
|
||
}
|
||
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, CoinAmount: coinAmount,
|
||
RatioPointAmount: ratioPointAmount, RatioCoinAmount: ratioCoinAmount,
|
||
PointBalanceAfter: pointAfter, CoinBalanceAfter: coinAfter, CreatedAtMS: nowMS,
|
||
}
|
||
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: ledger.AssetPoint, 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, ledger.AssetPoint, -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,
|
||
}
|
||
}
|