290 lines
10 KiB
Go
290 lines
10 KiB
Go
package processor
|
||
|
||
import (
|
||
"context"
|
||
"errors"
|
||
"fmt"
|
||
"strings"
|
||
"time"
|
||
|
||
"github.com/google/uuid"
|
||
|
||
"apskel-pos-be/internal/constants"
|
||
"apskel-pos-be/internal/entities"
|
||
"apskel-pos-be/internal/models"
|
||
"apskel-pos-be/internal/repository"
|
||
)
|
||
|
||
// ErrWalletMoveRejected wraps every reason an exchange or a transfer is refused on
|
||
// the customer's side: the amount, the limits, the recipient or the balance. The
|
||
// message says which.
|
||
var ErrWalletMoveRejected = errors.New("wallet move refused")
|
||
|
||
// walletMoveKeyLimit keeps a client's Idempotency-Key short enough to fit, with the
|
||
// prefix that scopes it to the customer, in wallet_transactions.idempotency_key.
|
||
const walletMoveKeyLimit = 50
|
||
|
||
type organizationSettingsReader interface {
|
||
Organization(ctx context.Context, organizationID uuid.UUID) (*models.OrganizationLoyaltySettings, error)
|
||
}
|
||
|
||
// WalletExchangeProcessor exchanges EnakCoin into EnakPoint (docs/prd-point-coin.md
|
||
// F4, K3). It is one way only; nothing turns EnakPoint back into EnakCoin.
|
||
type WalletExchangeProcessor struct {
|
||
customers repository.WalletMoveRepository
|
||
settings organizationSettingsReader
|
||
spendable spendableReader
|
||
pins pinVerifier
|
||
wallet *WalletProcessor
|
||
tx TxRunner
|
||
now func() time.Time
|
||
}
|
||
|
||
func NewWalletExchangeProcessor(customers repository.WalletMoveRepository, settings organizationSettingsReader, spendable spendableReader, pins pinVerifier, wallet *WalletProcessor, tx TxRunner) *WalletExchangeProcessor {
|
||
return &WalletExchangeProcessor{customers: customers, settings: settings, spendable: spendable, pins: pins, wallet: wallet, tx: tx, now: time.Now}
|
||
}
|
||
|
||
// Preview is GET /customer/wallet/exchange/preview: the organization's rate and what
|
||
// exchanging coins would give, so the app can show it before asking for the PIN.
|
||
func (p *WalletExchangeProcessor) Preview(ctx context.Context, customerID uuid.UUID, coins int64) (*models.WalletExchangePreview, error) {
|
||
customer, err := p.customers.GetCustomer(ctx, customerID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
settings, err := p.settings.Organization(ctx, customer.OrganizationID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
balances, err := p.spendable.SpendableBalances(ctx, customerID, p.now())
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
rate := settings.Exchange
|
||
preview := &models.WalletExchangePreview{
|
||
CoinAmount: rate.CoinAmount,
|
||
PointAmount: rate.PointAmount,
|
||
CoinBalance: balances[constants.WalletCurrencyCoin],
|
||
Coins: coins,
|
||
}
|
||
if reason := exchangeProblem(customer, rate, coins); reason != "" {
|
||
preview.Reason = reason
|
||
return preview, nil
|
||
}
|
||
preview.Points = exchangePoints(coins, rate)
|
||
if coins > preview.CoinBalance {
|
||
preview.Reason = "not enough EnakCoin"
|
||
return preview, nil
|
||
}
|
||
preview.Valid = true
|
||
return preview, nil
|
||
}
|
||
|
||
// Exchange takes coins EnakCoin and gives the EnakPoint they are worth at the
|
||
// organization's rate, in one transaction. The two ledger rows share a group and
|
||
// point at each other, and both freeze the rate. Each EnakPoint lot keeps the expiry
|
||
// of the EnakCoin lot it came from, so exchanging cannot extend a balance's life
|
||
// (K9). The PIN approves it (K8).
|
||
//
|
||
// idempotencyKey is the client's Idempotency-Key: a retry with the same key returns
|
||
// the first exchange, at the rate it was made, without moving anything again.
|
||
func (p *WalletExchangeProcessor) Exchange(ctx context.Context, customerID uuid.UUID, coins int64, pin, idempotencyKey string, info models.CustomerPinRequestInfo) (*models.WalletExchangeResult, error) {
|
||
key, err := walletMoveKey(idempotencyKey)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
customer, err := p.customers.GetCustomer(ctx, customerID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
settings, err := p.settings.Organization(ctx, customer.OrganizationID)
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
// Refuse a malformed request before the PIN is checked, so a typo in the amount
|
||
// costs the customer no PIN attempt.
|
||
if reason := exchangeProblem(customer, settings.Exchange, coins); reason != "" {
|
||
return nil, fmt.Errorf("%w: %s", ErrWalletMoveRejected, reason)
|
||
}
|
||
if err := p.pins.VerifyPin(ctx, customerID, pin, PinActionExchange, info); err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
outKey := fmt.Sprintf("exchange:%s:%s:out", customerID, key)
|
||
inKey := fmt.Sprintf("exchange:%s:%s:in", customerID, key)
|
||
result := &models.WalletExchangeResult{Coins: coins}
|
||
err = p.tx.WithTransaction(ctx, func(ctx context.Context) error {
|
||
if err := p.wallet.LockWallet(ctx, customerID); err != nil {
|
||
return err
|
||
}
|
||
rate := settings.Exchange
|
||
groupID, outID, inID := uuid.New(), uuid.New(), uuid.New()
|
||
previous, err := p.wallet.FindTransaction(ctx, outKey)
|
||
if err != nil {
|
||
return err
|
||
}
|
||
if previous != nil && previous.GroupID != nil {
|
||
// A retry: repeat it with the ids and the rate the first attempt froze, so
|
||
// both rows replay even if the rate has changed since.
|
||
outID, inID, groupID = previous.ID, previous.ReferenceID, *previous.GroupID
|
||
rate = frozenExchangeRate(previous.Metadata, rate)
|
||
}
|
||
points := exchangePoints(coins, rate)
|
||
metadata := entities.Metadata{
|
||
"coins": coins,
|
||
"points": points,
|
||
"coin_amount": rate.CoinAmount,
|
||
"point_amount": rate.PointAmount,
|
||
}
|
||
|
||
out, err := p.wallet.Debit(ctx, WalletDebitInput{WalletEntry: WalletEntry{
|
||
TransactionID: outID,
|
||
CustomerID: customerID,
|
||
Currency: constants.WalletCurrencyCoin,
|
||
Type: constants.WalletTxTypeExchangeOut,
|
||
Amount: coins,
|
||
ReferenceType: constants.WalletRefTypeWalletTx,
|
||
ReferenceID: inID,
|
||
GroupID: &groupID,
|
||
Description: fmt.Sprintf("Tukar %d EnakCoin ke EnakPoint", coins),
|
||
Metadata: metadata,
|
||
IdempotencyKey: outKey,
|
||
}})
|
||
if errors.Is(err, repository.ErrWalletInsufficientBalance) {
|
||
return fmt.Errorf("%w: not enough EnakCoin", ErrWalletMoveRejected)
|
||
}
|
||
if err != nil {
|
||
return err
|
||
}
|
||
in, err := p.wallet.Credit(ctx, WalletCreditInput{
|
||
WalletEntry: WalletEntry{
|
||
TransactionID: inID,
|
||
CustomerID: customerID,
|
||
Currency: constants.WalletCurrencyPoint,
|
||
Type: constants.WalletTxTypeExchangeIn,
|
||
Amount: points,
|
||
ReferenceType: constants.WalletRefTypeWalletTx,
|
||
ReferenceID: outID,
|
||
GroupID: &groupID,
|
||
Description: fmt.Sprintf("Dari tukar %d EnakCoin", coins),
|
||
Metadata: metadata,
|
||
IdempotencyKey: inKey,
|
||
},
|
||
Lots: exchangeLots(out.Allocations, rate, ComputeExpiry(settings.PointExpiry, p.now())),
|
||
})
|
||
if err != nil {
|
||
return err
|
||
}
|
||
|
||
result.GroupID = groupID
|
||
result.Points = points
|
||
result.CoinAmount = rate.CoinAmount
|
||
result.PointAmount = rate.PointAmount
|
||
result.Lots = movedLots(in.Lots)
|
||
result.Replayed = out.Replayed
|
||
return nil
|
||
})
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
|
||
balances, err := p.spendable.SpendableBalances(ctx, customerID, p.now())
|
||
if err != nil {
|
||
return nil, err
|
||
}
|
||
result.CoinBalance = balances[constants.WalletCurrencyCoin]
|
||
result.PointBalance = balances[constants.WalletCurrencyPoint]
|
||
return result, nil
|
||
}
|
||
|
||
// exchangeProblem says why coins cannot be exchanged, or "" when they can as far as
|
||
// the request goes. The balance is checked under the wallet lock.
|
||
func exchangeProblem(customer *repository.WalletMoveCustomer, rate models.LoyaltyExchangeSettings, coins int64) string {
|
||
switch {
|
||
case !customer.IsActive:
|
||
return "the customer is not active"
|
||
case coins <= 0:
|
||
return "the number of EnakCoin must be positive"
|
||
case rate.CoinAmount <= 0 || rate.PointAmount <= 0:
|
||
return "exchange is not available"
|
||
case coins%rate.CoinAmount != 0:
|
||
// Otherwise part of the EnakCoin would be lost to rounding (F4).
|
||
return fmt.Sprintf("EnakCoin are exchanged in multiples of %d", rate.CoinAmount)
|
||
}
|
||
return ""
|
||
}
|
||
|
||
// exchangePoints is (coins / coin_amount) × point_amount, for coins that are a
|
||
// multiple of coin_amount.
|
||
func exchangePoints(coins int64, rate models.LoyaltyExchangeSettings) int64 {
|
||
return coins / rate.CoinAmount * rate.PointAmount
|
||
}
|
||
|
||
// exchangeLots splits the EnakPoint of an exchange over the EnakCoin lots it took,
|
||
// so each part keeps the expiry of its lot and points back at it (K9). The share of
|
||
// a lot is the difference of floor(coins so far × point_amount / coin_amount) before
|
||
// and after it, which adds up exactly because the total is a multiple of
|
||
// coin_amount. A lot too small to earn a whole EnakPoint on its own gives none.
|
||
//
|
||
// Each part expires at the sooner of its EnakCoin lot's expiry and pointExpiry, when an
|
||
// EnakPoint received now would expire (F4); nil means never.
|
||
func exchangeLots(allocations []WalletAllocation, rate models.LoyaltyExchangeSettings, pointExpiry *time.Time) []WalletLotInput {
|
||
var lots []WalletLotInput
|
||
var coinsSoFar int64
|
||
for _, a := range allocations {
|
||
before := coinsSoFar * rate.PointAmount / rate.CoinAmount
|
||
coinsSoFar += a.Amount
|
||
points := coinsSoFar*rate.PointAmount/rate.CoinAmount - before
|
||
if points == 0 {
|
||
continue
|
||
}
|
||
lotID := a.LotID
|
||
lots = append(lots, WalletLotInput{Amount: points, ExpiresAt: EarlierExpiry(a.ExpiresAt, pointExpiry), OriginLotID: &lotID})
|
||
}
|
||
return lots
|
||
}
|
||
|
||
// frozenExchangeRate reads the rate an exchange was made at from its ledger row.
|
||
func frozenExchangeRate(metadata entities.Metadata, fallback models.LoyaltyExchangeSettings) models.LoyaltyExchangeSettings {
|
||
coinAmount, ok1 := metadataInt(metadata, "coin_amount")
|
||
pointAmount, ok2 := metadataInt(metadata, "point_amount")
|
||
if !ok1 || !ok2 || coinAmount <= 0 || pointAmount <= 0 {
|
||
return fallback
|
||
}
|
||
return models.LoyaltyExchangeSettings{CoinAmount: coinAmount, PointAmount: pointAmount}
|
||
}
|
||
|
||
// metadataInt reads a whole number from metadata that may have been through JSONB,
|
||
// which gives numbers back as float64.
|
||
func metadataInt(metadata entities.Metadata, key string) (int64, bool) {
|
||
switch v := metadata[key].(type) {
|
||
case float64:
|
||
return int64(v), true
|
||
case int64:
|
||
return v, true
|
||
case int:
|
||
return int64(v), true
|
||
}
|
||
return 0, false
|
||
}
|
||
|
||
func movedLots(lots []entities.WalletLot) []models.WalletMovedLot {
|
||
out := make([]models.WalletMovedLot, 0, len(lots))
|
||
for _, lot := range lots {
|
||
out = append(out, models.WalletMovedLot{Amount: lot.OriginalAmount, ExpiresAt: lot.ExpiresAt})
|
||
}
|
||
return out
|
||
}
|
||
|
||
// walletMoveKey checks the client's Idempotency-Key, which exchanges and transfers
|
||
// require (F4, F5).
|
||
func walletMoveKey(key string) (string, error) {
|
||
key = strings.TrimSpace(key)
|
||
if key == "" {
|
||
return "", fmt.Errorf("%w: the Idempotency-Key header is required", ErrWalletMoveRejected)
|
||
}
|
||
if len(key) > walletMoveKeyLimit {
|
||
return "", fmt.Errorf("%w: the Idempotency-Key header must be at most %d characters", ErrWalletMoveRejected, walletMoveKeyLimit)
|
||
}
|
||
return key, nil
|
||
}
|