Adds what the customer sees of expiry (docs/prd-point-coin.md F6, F12, PC-504). GET /customer/wallet/expiring lists everything that will expire, per currency and day, soonest first. GET /customer/wallet already had the nearest expiry per currency. The expiry job now also sends reminders, with the settings of note N4 as decided: once, reminder_days before (7 by default, 0 for none), per currency. A customer gets one FCM push per currency and expiry day, however many lots make it up: "150 EnakPoint akan kedaluwarsa pada 31 Okt 2026. Pakai sebelum hangus.", with type WALLET_EXPIRING, the currency, amount and expiry_date in its data. Reminders cover whatever falls within the window, so a run that was missed catches up rather than skipping a day. Migration 000097 adds wallet_expiry_reminders, one row per customer, currency and expiry day. The row is written before the push is sent, so several instances of the job or a restart never remind twice; a push that then fails is logged and not retried. Lots that expire later on the same day as an earlier reminder are not reminded of again. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
238 lines
9.5 KiB
Go
238 lines
9.5 KiB
Go
package processor
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
"os"
|
|
"sync"
|
|
"testing"
|
|
"time"
|
|
|
|
"github.com/google/uuid"
|
|
"github.com/stretchr/testify/assert"
|
|
"github.com/stretchr/testify/require"
|
|
"gorm.io/driver/postgres"
|
|
"gorm.io/gorm"
|
|
"gorm.io/gorm/logger"
|
|
|
|
"apskel-pos-be/internal/constants"
|
|
"apskel-pos-be/internal/models"
|
|
"apskel-pos-be/internal/repository"
|
|
)
|
|
|
|
// fixedOrganizationSettings serves the same organization settings to every caller.
|
|
type fixedOrganizationSettings struct {
|
|
s models.OrganizationLoyaltySettings
|
|
}
|
|
|
|
func (f fixedOrganizationSettings) Organization(context.Context, uuid.UUID) (*models.OrganizationLoyaltySettings, error) {
|
|
s := f.s
|
|
return &s, nil
|
|
}
|
|
|
|
// walletMoveDB opens TEST_DATABASE_URL and creates an organization with two customers,
|
|
// removed again when the test ends. See internal/repository/wallet_repository_test.go.
|
|
func walletMoveDB(t *testing.T) (db *gorm.DB, org, a, b uuid.UUID) {
|
|
t.Helper()
|
|
dsn := os.Getenv("TEST_DATABASE_URL")
|
|
if dsn == "" {
|
|
t.Skip("TEST_DATABASE_URL not set")
|
|
}
|
|
db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
|
require.NoError(t, err)
|
|
|
|
org, a, b = uuid.New(), uuid.New(), uuid.New()
|
|
phoneA, phoneB := "08"+a.String()[:10], "08"+b.String()[:10]
|
|
require.NoError(t, db.Exec(`INSERT INTO organizations (id, name, plan_type) VALUES (?, 'wallet move test', 'basic')`, org).Error)
|
|
require.NoError(t, db.Exec(`INSERT INTO customers (id, organization_id, name, phone_number) VALUES (?, ?, 'Anita', ?), (?, ?, 'Budi Santoso', ?)`,
|
|
a, org, phoneA, b, org, phoneB).Error)
|
|
customers := []uuid.UUID{a, b}
|
|
t.Cleanup(func() {
|
|
db.Exec(`DELETE FROM wallet_lot_allocations WHERE lot_id IN (SELECT id FROM wallet_lots WHERE customer_id IN ?)`, customers)
|
|
db.Exec(`DELETE FROM wallet_lots WHERE customer_id IN ? AND origin_lot_id IS NOT NULL`, customers)
|
|
db.Exec(`DELETE FROM wallet_lots WHERE customer_id IN ?`, customers)
|
|
db.Exec(`DELETE FROM wallet_transactions WHERE customer_id IN ?`, customers)
|
|
db.Exec(`DELETE FROM customer_wallets WHERE customer_id IN ?`, customers)
|
|
db.Exec(`DELETE FROM customers WHERE id IN ?`, customers)
|
|
db.Exec(`DELETE FROM organizations WHERE id = ?`, org)
|
|
})
|
|
return db, org, a, b
|
|
}
|
|
|
|
func TestWalletExchange_AgainstPostgres(t *testing.T) {
|
|
db, _, a, _ := walletMoveDB(t)
|
|
wallet := NewWalletProcessor(repository.NewWalletRepository(db))
|
|
txm := repository.NewTxManager(db)
|
|
settings := fixedOrganizationSettings{models.OrganizationLoyaltySettings{
|
|
Exchange: models.LoyaltyExchangeSettings{CoinAmount: 10, PointAmount: 3},
|
|
}}
|
|
p := NewWalletExchangeProcessor(repository.NewWalletMoveRepository(db), settings, repository.NewWalletQueryRepository(db),
|
|
&movePinFake{good: "482913"}, wallet, txm)
|
|
|
|
expiry := time.Now().Add(24 * time.Hour).Truncate(time.Second)
|
|
require.NoError(t, txm.WithTransaction(context.Background(), func(ctx context.Context) error {
|
|
in := earn(a, 30, &expiry)
|
|
in.Currency = constants.WalletCurrencyCoin
|
|
_, err := wallet.Credit(ctx, in)
|
|
return err
|
|
}))
|
|
|
|
res, err := p.Exchange(context.Background(), a, 20, "482913", "db-key", models.CustomerPinRequestInfo{})
|
|
require.NoError(t, err)
|
|
assert.Equal(t, int64(6), res.Points)
|
|
assert.Equal(t, int64(10), res.CoinBalance)
|
|
assert.Equal(t, int64(6), res.PointBalance)
|
|
require.Len(t, res.Lots, 1)
|
|
require.NotNil(t, res.Lots[0].ExpiresAt)
|
|
assert.True(t, res.Lots[0].ExpiresAt.Equal(expiry))
|
|
|
|
// The retry reads the frozen rate back out of JSONB and replays.
|
|
again, err := p.Exchange(context.Background(), a, 20, "482913", "db-key", models.CustomerPinRequestInfo{})
|
|
require.NoError(t, err)
|
|
assert.True(t, again.Replayed)
|
|
assert.Equal(t, int64(6), again.Points)
|
|
|
|
var rows int64
|
|
require.NoError(t, db.Raw(`SELECT COUNT(*) FROM wallet_transactions WHERE group_id = ?`, res.GroupID).Scan(&rows).Error)
|
|
assert.Equal(t, int64(2), rows)
|
|
}
|
|
|
|
// Transfers in both directions at once must not deadlock: both lock the two wallets
|
|
// in customer_id order. Every one of them lands, and the totals still reconcile.
|
|
func TestWalletTransfer_BothWaysAtOnceAgainstPostgres(t *testing.T) {
|
|
db, _, a, b := walletMoveDB(t)
|
|
wallet := NewWalletProcessor(repository.NewWalletRepository(db))
|
|
txm := repository.NewTxManager(db)
|
|
moves := repository.NewWalletMoveRepository(db)
|
|
settings := fixedOrganizationSettings{models.OrganizationLoyaltySettings{
|
|
Transfer: models.LoyaltyTransferSettings{Enabled: true, MinAmount: 1},
|
|
}}
|
|
p := NewWalletTransferProcessor(moves, settings, repository.NewWalletQueryRepository(db), &movePinFake{good: "482913"}, wallet, txm, nil)
|
|
|
|
expiry := time.Now().Add(24 * time.Hour).Truncate(time.Second)
|
|
require.NoError(t, txm.WithTransaction(context.Background(), func(ctx context.Context) error {
|
|
if _, err := wallet.Credit(ctx, earn(a, 100, &expiry)); err != nil {
|
|
return err
|
|
}
|
|
_, err := wallet.Credit(ctx, earn(b, 100, nil))
|
|
return err
|
|
}))
|
|
phone := func(id uuid.UUID) string { return "08" + id.String()[:10] }
|
|
|
|
const rounds = 10
|
|
errs := make(chan error, 2*rounds)
|
|
var wg sync.WaitGroup
|
|
for i := 0; i < rounds; i++ {
|
|
for _, pair := range [][2]uuid.UUID{{a, b}, {b, a}} {
|
|
wg.Add(1)
|
|
go func(from, to uuid.UUID, i int) {
|
|
defer wg.Done()
|
|
_, err := p.Transfer(context.Background(), from, sendPoints(1, phone(to)), "482913", fmt.Sprintf("race-%d", i), models.CustomerPinRequestInfo{})
|
|
errs <- err
|
|
}(pair[0], pair[1], i)
|
|
}
|
|
}
|
|
wg.Wait()
|
|
close(errs)
|
|
for err := range errs {
|
|
assert.NoError(t, err)
|
|
}
|
|
|
|
var balances []int64
|
|
require.NoError(t, db.Raw(`SELECT point_balance FROM customer_wallets WHERE customer_id IN ? ORDER BY point_balance`, []uuid.UUID{a, b}).Scan(&balances).Error)
|
|
assert.Equal(t, []int64{100, 100}, balances)
|
|
|
|
// B's lots that came from A keep A's expiry to the second.
|
|
var mismatched int64
|
|
require.NoError(t, db.Raw(`
|
|
SELECT COUNT(*) FROM wallet_lots l JOIN wallet_lots o ON o.id = l.origin_lot_id
|
|
WHERE l.customer_id = ? AND o.customer_id = ? AND l.expires_at IS DISTINCT FROM o.expires_at`, b, a).Scan(&mismatched).Error)
|
|
assert.Zero(t, mismatched)
|
|
|
|
require.NoError(t, txm.WithTransaction(context.Background(), func(ctx context.Context) error {
|
|
sent, err := moves.TransferredOutSince(ctx, a, constants.WalletCurrencyPoint, startOfWalletDay(time.Now()))
|
|
assert.Equal(t, int64(rounds), sent)
|
|
return err
|
|
}))
|
|
}
|
|
|
|
// The example of ยง8 against Postgres: B's payment of 30 traces back to A's #ORD-1.
|
|
func TestWalletTrace_AgainstPostgres(t *testing.T) {
|
|
db, org, a, b := walletMoveDB(t)
|
|
wallet := NewWalletProcessor(repository.NewWalletRepository(db))
|
|
txm := repository.NewTxManager(db)
|
|
settings := fixedOrganizationSettings{models.OrganizationLoyaltySettings{
|
|
Transfer: models.LoyaltyTransferSettings{Enabled: true, MinAmount: 1},
|
|
}}
|
|
transfers := NewWalletTransferProcessor(repository.NewWalletMoveRepository(db), settings, repository.NewWalletQueryRepository(db), &movePinFake{good: "482913"}, wallet, txm, nil)
|
|
|
|
dec, jan := time.Now().Add(30*24*time.Hour), time.Now().Add(60*24*time.Hour)
|
|
ord1 := earn(a, 100, &dec)
|
|
ord1.Description = "Belanja #ORD-1"
|
|
require.NoError(t, txm.WithTransaction(context.Background(), func(ctx context.Context) error {
|
|
if _, err := wallet.Credit(ctx, ord1); err != nil {
|
|
return err
|
|
}
|
|
_, err := wallet.Credit(ctx, earn(a, 50, &jan))
|
|
return err
|
|
}))
|
|
_, err := transfers.Transfer(context.Background(), a, sendPoints(120, "08"+b.String()[:10]), "482913", "trace", models.CustomerPinRequestInfo{})
|
|
require.NoError(t, err)
|
|
var payment *WalletResult
|
|
require.NoError(t, txm.WithTransaction(context.Background(), func(ctx context.Context) error {
|
|
payment, err = wallet.Debit(ctx, pay(b, 30))
|
|
return err
|
|
}))
|
|
|
|
trace, err := NewWalletTraceProcessor(repository.NewWalletTraceRepository(db)).Trace(context.Background(), org, payment.Transaction.ID)
|
|
require.NoError(t, err)
|
|
require.Len(t, trace.Lots, 1)
|
|
chain := trace.Lots[0].Chain
|
|
require.Len(t, chain, 2)
|
|
assert.Equal(t, constants.WalletTxTypeTransferIn, chain[0].Source.Type)
|
|
assert.Equal(t, "Anita", chain[1].Source.Customer.Name)
|
|
assert.Equal(t, ord1.ReferenceID, chain[1].Source.ReferenceID)
|
|
|
|
_, err = NewWalletTraceProcessor(repository.NewWalletTraceRepository(db)).Trace(context.Background(), uuid.New(), payment.Transaction.ID)
|
|
assert.ErrorIs(t, err, repository.ErrWalletTransactionNotFound)
|
|
}
|
|
|
|
// Two instances of the expiry job at once expire each lot exactly once (PC-503).
|
|
func TestWalletExpiry_TwoInstancesAgainstPostgres(t *testing.T) {
|
|
db, _, a, b := walletMoveDB(t)
|
|
wallet := NewWalletProcessor(repository.NewWalletRepository(db))
|
|
txm := repository.NewTxManager(db)
|
|
past := time.Now().Add(-time.Hour)
|
|
require.NoError(t, txm.WithTransaction(context.Background(), func(ctx context.Context) error {
|
|
for i := 0; i < 5; i++ {
|
|
for _, c := range []uuid.UUID{a, b} {
|
|
if _, err := wallet.Credit(ctx, earn(c, 10, &past)); err != nil {
|
|
return err
|
|
}
|
|
}
|
|
}
|
|
return nil
|
|
}))
|
|
|
|
var wg sync.WaitGroup
|
|
counts := make([]int, 2)
|
|
for i := range counts {
|
|
wg.Add(1)
|
|
go func(i int) {
|
|
defer wg.Done()
|
|
p := NewWalletExpiryProcessor(repository.NewWalletExpiryRepository(db), fixedOrganizationSettings{}, wallet, txm, nil)
|
|
n, err := p.ExpireDue(context.Background())
|
|
assert.NoError(t, err)
|
|
counts[i] = n
|
|
}(i)
|
|
}
|
|
wg.Wait()
|
|
assert.Equal(t, 10, counts[0]+counts[1], "every lot once, between them")
|
|
|
|
var expires, left int64
|
|
require.NoError(t, db.Raw(`SELECT COUNT(*) FROM wallet_transactions WHERE customer_id IN ? AND type = 'EXPIRE'`, []uuid.UUID{a, b}).Scan(&expires).Error)
|
|
require.NoError(t, db.Raw(`SELECT COALESCE(SUM(point_balance), 0) FROM customer_wallets WHERE customer_id IN ?`, []uuid.UUID{a, b}).Scan(&left).Error)
|
|
assert.Equal(t, int64(10), expires)
|
|
assert.Zero(t, left)
|
|
}
|