package processor import ( "context" "encoding/binary" "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 := walletTestPhone(a), walletTestPhone(b) 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 := walletTestPhone 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, redeem(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) } // walletTestPhone is a phone number in its stored form, 628…, unique to the customer. func walletTestPhone(id uuid.UUID) string { return fmt.Sprintf("628%09d", binary.BigEndian.Uint64(id[:8])%1_000_000_000) }