2026-09-30 12:09:53 +07:00
|
|
|
|
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,
|
|
|
|
|
|
},
|
2026-09-30 13:29:25 +07:00
|
|
|
|
Lots: exchangeLots(out.Allocations, rate, ComputeExpiry(settings.PointExpiry, p.now())),
|
2026-09-30 12:09:53 +07:00
|
|
|
|
})
|
|
|
|
|
|
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.
|
|
|
|
|
|
//
|
2026-09-30 13:29:25 +07:00
|
|
|
|
// 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 {
|
2026-09-30 12:09:53 +07:00
|
|
|
|
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
|
2026-09-30 13:29:25 +07:00
|
|
|
|
lots = append(lots, WalletLotInput{Amount: points, ExpiresAt: EarlierExpiry(a.ExpiresAt, pointExpiry), OriginLotID: &lotID})
|
2026-09-30 12:09:53 +07:00
|
|
|
|
}
|
|
|
|
|
|
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
|
|
|
|
|
|
}
|