Files
apskel-pos-backend/internal/repository/wallet_reconciliation_repository_test.go
T
efrilmandClaude Opus 5.5 040780cd2d feat(wallet): reconcile balances, ledger and lots on a schedule
Adds the reconciliation of docs/prd-point-coin.md §7.5 (PC-108). One
aggregate query per check, across every wallet:

- wallet balance = SUM(ledger), per currency, including customers with
  ledger rows but no wallet row
- wallet balance = SUM(lot remaining)
- lot original - SUM(allocations) = remaining
- SUM(allocations) = |amount| for every deduction
- lots created = amount for every addition, which the engine keeps and the
  other checks rely on

The check on payments.points_used waits for that column (PC-305).

WalletReconciliationJob runs the checks at startup and every six hours,
alongside the omset scheduler. It is silent while the data is consistent.
Each discrepancy is logged with its check, customer, object and the
expected and actual values, and the organization's admins, owners and
managers get a high-priority notification. An organization is notified
again only when its set of discrepancies changes. Nothing is corrected
automatically. At most 50 discrepancies per check are reported.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-09-30 10:13:46 +07:00

142 lines
6.1 KiB
Go

package repository_test
import (
"context"
"os"
"sort"
"testing"
"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/processor"
"apskel-pos-be/internal/repository"
)
// Builds consistent wallets through the engine, checks the reconciliation is silent,
// then breaks each §7.5 invariant for a different customer and checks each break is
// found by the right checks and nothing else is. Needs TEST_DATABASE_URL pointing at
// a migrated database; see wallet_repository_test.go.
func TestWalletReconciliation_AgainstPostgres(t *testing.T) {
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)
ctx := context.Background()
org := uuid.New()
names := []string{"clean", "balance", "lot", "allocation", "credit"}
customers := map[string]uuid.UUID{}
var ids []uuid.UUID
require.NoError(t, db.Exec(`INSERT INTO organizations (id, name, plan_type) VALUES (?, 'recon test', 'basic')`, org).Error)
for _, n := range names {
customers[n] = uuid.New()
ids = append(ids, customers[n])
require.NoError(t, db.Exec(`INSERT INTO customers (id, organization_id, name) VALUES (?, ?, ?)`, customers[n], org, n).Error)
}
t.Cleanup(func() {
db.Exec(`DELETE FROM wallet_lot_allocations WHERE lot_id IN (SELECT id FROM wallet_lots WHERE customer_id IN ?)`, ids)
db.Exec(`DELETE FROM wallet_lot_allocations WHERE transaction_id IN (SELECT id FROM wallet_transactions WHERE customer_id IN ?)`, ids)
db.Exec(`DELETE FROM wallet_lots WHERE customer_id IN ?`, ids)
db.Exec(`DELETE FROM wallet_transactions WHERE customer_id IN ?`, ids)
db.Exec(`DELETE FROM customer_wallets WHERE customer_id IN ?`, ids)
db.Exec(`DELETE FROM customers WHERE id IN ?`, ids)
db.Exec(`DELETE FROM organizations WHERE id = ?`, org)
})
// Every customer: +100 in two lots, +30 coins, -70 points across both lots.
wallet := processor.NewWalletProcessor(repository.NewWalletRepository(db))
txm := repository.NewTxManager(db)
for _, id := range ids {
require.NoError(t, txm.WithTransaction(ctx, func(ctx context.Context) error {
if _, err := wallet.Credit(ctx, processor.WalletCreditInput{
WalletEntry: processor.WalletEntry{CustomerID: id, Currency: constants.WalletCurrencyPoint,
Type: constants.WalletTxTypeMigration, Amount: 100, ReferenceType: constants.WalletRefTypeLegacyPoints,
ReferenceID: uuid.New(), Description: "Saldo awal"},
Lots: []processor.WalletLotInput{{Amount: 60}, {Amount: 40}},
}); err != nil {
return err
}
if _, err := wallet.Credit(ctx, processor.WalletCreditInput{WalletEntry: processor.WalletEntry{
CustomerID: id, Currency: constants.WalletCurrencyCoin, Type: constants.WalletTxTypeMigration,
Amount: 30, ReferenceType: constants.WalletRefTypeLegacyTokens, ReferenceID: id, Description: "Saldo awal"}}); err != nil {
return err
}
outlet := uuid.New()
_, err := wallet.Debit(ctx, processor.WalletDebitInput{WalletEntry: processor.WalletEntry{
CustomerID: id, Currency: constants.WalletCurrencyPoint, Type: constants.WalletTxTypePayment,
Amount: 70, ReferenceType: constants.WalletRefTypePayment, ReferenceID: uuid.New(), OutletID: &outlet,
Description: "Bayar"}})
return err
}))
}
recon := repository.NewWalletReconciliationRepository(db)
// Other packages' tests may share the database, so only these customers count.
checksByCustomer := func() map[string][]string {
t.Helper()
found, err := recon.FindDiscrepancies(ctx, 1000)
require.NoError(t, err)
byName := map[uuid.UUID]string{}
for n, id := range customers {
byName[id] = n
}
out := map[string][]string{}
for _, d := range found {
if n, ok := byName[d.CustomerID]; ok {
assert.Equal(t, org, d.OrganizationID)
out[n] = append(out[n], d.Check)
}
}
for n := range out {
sort.Strings(out[n])
}
return out
}
assert.Empty(t, checksByCustomer(), "consistent data reports nothing")
exec := func(q string, args ...any) {
t.Helper()
require.NoError(t, db.Exec(q, args...).Error)
}
// A balance moved without a ledger row or a lot.
exec(`UPDATE customer_wallets SET point_balance = point_balance + 5 WHERE customer_id = ?`, customers["balance"])
// A lot's remainder changed without an allocation.
exec(`UPDATE wallet_lots SET remaining_amount = remaining_amount - 1
WHERE id = (SELECT id FROM wallet_lots WHERE customer_id = ? AND remaining_amount > 0 LIMIT 1)`, customers["lot"])
// An allocation lost.
exec(`DELETE FROM wallet_lot_allocations WHERE (transaction_id, lot_id) IN (
SELECT a.transaction_id, a.lot_id FROM wallet_lot_allocations a
JOIN wallet_transactions t ON t.id = a.transaction_id WHERE t.customer_id = ? LIMIT 1)`, customers["allocation"])
// A credit whose lot was never written, with the balance moved to match the ledger.
exec(`INSERT INTO wallet_transactions (organization_id, customer_id, currency, type, amount, balance_after, reference_type, reference_id, description)
VALUES (?, ?, 'COIN', 'MIGRATION', 10, 40, 'LEGACY_TOKENS', ?, 'x')`, org, customers["credit"], customers["credit"])
exec(`UPDATE customer_wallets SET coin_balance = coin_balance + 10 WHERE customer_id = ?`, customers["credit"])
assert.Equal(t, map[string][]string{
"balance": {repository.WalletCheckBalanceVsLedger, repository.WalletCheckBalanceVsLots},
"lot": {repository.WalletCheckBalanceVsLots, repository.WalletCheckLotVsAllocations},
"allocation": {repository.WalletCheckDebitVsAllocations, repository.WalletCheckLotVsAllocations},
"credit": {repository.WalletCheckBalanceVsLots, repository.WalletCheckCreditVsLots},
}, checksByCustomer(), "each break is found by exactly the checks it violates, and the clean customer by none")
// The per-check limit caps the report.
found, err := recon.FindDiscrepancies(ctx, 1)
require.NoError(t, err)
perCheck := map[string]int{}
for _, d := range found {
perCheck[d.Check]++
}
for check, n := range perCheck {
assert.LessOrEqual(t, n, 1, check)
}
}