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>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
a6d5a8b056
commit
040780cd2d
@@ -31,6 +31,7 @@ type App struct {
|
||||
router *router.Router
|
||||
shutdown chan os.Signal
|
||||
omsetScheduler *service.OmsetMilestoneScheduler
|
||||
walletRecon *service.WalletReconciliationJob
|
||||
}
|
||||
|
||||
func NewApp(db *gorm.DB, redisClient *redis.Client) *App {
|
||||
@@ -53,6 +54,13 @@ func (a *App) Initialize(cfg *config.Config) error {
|
||||
processors.notificationProcessor,
|
||||
)
|
||||
|
||||
// Checks that wallet balances, ledger and lots agree (docs/prd-point-coin.md §7.5)
|
||||
a.walletRecon = service.NewWalletReconciliationJob(
|
||||
repository.NewWalletReconciliationRepository(a.db),
|
||||
repos.userRepo,
|
||||
processors.notificationProcessor,
|
||||
)
|
||||
|
||||
services := a.initServices(processors, repos, cfg)
|
||||
validators := a.initValidators()
|
||||
middleware := a.initMiddleware(services, cfg)
|
||||
@@ -155,6 +163,9 @@ func (a *App) Start(port string) error {
|
||||
if a.omsetScheduler != nil {
|
||||
a.omsetScheduler.Start(5 * time.Minute)
|
||||
}
|
||||
if a.walletRecon != nil {
|
||||
a.walletRecon.Start(6 * time.Hour)
|
||||
}
|
||||
|
||||
engine := a.router.Init()
|
||||
|
||||
@@ -194,6 +205,9 @@ func (a *App) Shutdown() {
|
||||
if a.omsetScheduler != nil {
|
||||
a.omsetScheduler.Stop()
|
||||
}
|
||||
if a.walletRecon != nil {
|
||||
a.walletRecon.Stop()
|
||||
}
|
||||
close(a.shutdown)
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,160 @@
|
||||
package repository
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// The reconciliation checks of docs/prd-point-coin.md §7.5.
|
||||
const (
|
||||
// Wallet balance = SUM(amount) of the customer's ledger rows, per currency.
|
||||
WalletCheckBalanceVsLedger = "BALANCE_VS_LEDGER"
|
||||
// Wallet balance = SUM(remaining_amount) of the customer's lots, per currency.
|
||||
WalletCheckBalanceVsLots = "BALANCE_VS_LOTS"
|
||||
// For every lot: original_amount - SUM(allocations) = remaining_amount.
|
||||
WalletCheckLotVsAllocations = "LOT_VS_ALLOCATIONS"
|
||||
// For every deduction: SUM(allocations) = |amount|.
|
||||
WalletCheckDebitVsAllocations = "DEBIT_VS_ALLOCATIONS"
|
||||
// For every addition: SUM(original_amount) of the lots it created = amount. Not
|
||||
// listed in §7.5, but the engine keeps it and the other checks rely on it.
|
||||
WalletCheckCreditVsLots = "CREDIT_VS_LOTS"
|
||||
)
|
||||
|
||||
// WalletDiscrepancy is one place where the wallet tables disagree with each other.
|
||||
type WalletDiscrepancy struct {
|
||||
Check string
|
||||
OrganizationID uuid.UUID
|
||||
CustomerID uuid.UUID
|
||||
Currency string
|
||||
// The lot or ledger row the check is about. Nil for the per-wallet checks.
|
||||
ObjectID *uuid.UUID
|
||||
Expected int64
|
||||
Actual int64
|
||||
}
|
||||
|
||||
// WalletReconciliationRepository runs the §7.5 checks across every wallet.
|
||||
type WalletReconciliationRepository interface {
|
||||
// FindDiscrepancies returns every discrepancy, at most limit per check, so one
|
||||
// systematic bug cannot produce an unbounded report.
|
||||
FindDiscrepancies(ctx context.Context, limit int) ([]WalletDiscrepancy, error)
|
||||
}
|
||||
|
||||
type walletReconciliationRepository struct {
|
||||
db *gorm.DB
|
||||
}
|
||||
|
||||
func NewWalletReconciliationRepository(db *gorm.DB) WalletReconciliationRepository {
|
||||
return &walletReconciliationRepository{db: db}
|
||||
}
|
||||
|
||||
// Each query returns check, organization_id, customer_id, currency, object_id,
|
||||
// expected and actual. Aggregates are joined rather than correlated, so each check is
|
||||
// a handful of scans however many customers there are.
|
||||
var walletReconciliationQueries = []struct {
|
||||
check string
|
||||
sql string
|
||||
}{
|
||||
{WalletCheckBalanceVsLedger, `
|
||||
WITH ledger AS (
|
||||
SELECT customer_id, currency, MAX(organization_id::text) AS organization_id, SUM(amount) AS total
|
||||
FROM wallet_transactions GROUP BY customer_id, currency
|
||||
), balances AS (
|
||||
SELECT customer_id, organization_id::text AS organization_id, 'POINT' AS currency, point_balance AS balance FROM customer_wallets
|
||||
UNION ALL
|
||||
SELECT customer_id, organization_id::text, 'COIN', coin_balance FROM customer_wallets
|
||||
)
|
||||
SELECT COALESCE(b.organization_id, l.organization_id) AS organization_id,
|
||||
COALESCE(b.customer_id, l.customer_id)::text AS customer_id,
|
||||
COALESCE(b.currency, l.currency) AS currency,
|
||||
NULL AS object_id,
|
||||
COALESCE(l.total, 0) AS expected,
|
||||
COALESCE(b.balance, 0) AS actual
|
||||
FROM balances b
|
||||
FULL JOIN ledger l ON l.customer_id = b.customer_id AND l.currency = b.currency
|
||||
WHERE COALESCE(b.balance, 0) <> COALESCE(l.total, 0)
|
||||
LIMIT ?`},
|
||||
{WalletCheckBalanceVsLots, `
|
||||
WITH lots AS (
|
||||
SELECT customer_id, currency, MAX(organization_id::text) AS organization_id, SUM(remaining_amount) AS total
|
||||
FROM wallet_lots GROUP BY customer_id, currency
|
||||
), balances AS (
|
||||
SELECT customer_id, organization_id::text AS organization_id, 'POINT' AS currency, point_balance AS balance FROM customer_wallets
|
||||
UNION ALL
|
||||
SELECT customer_id, organization_id::text, 'COIN', coin_balance FROM customer_wallets
|
||||
)
|
||||
SELECT COALESCE(b.organization_id, l.organization_id) AS organization_id,
|
||||
COALESCE(b.customer_id, l.customer_id)::text AS customer_id,
|
||||
COALESCE(b.currency, l.currency) AS currency,
|
||||
NULL AS object_id,
|
||||
COALESCE(l.total, 0) AS expected,
|
||||
COALESCE(b.balance, 0) AS actual
|
||||
FROM balances b
|
||||
FULL JOIN lots l ON l.customer_id = b.customer_id AND l.currency = b.currency
|
||||
WHERE COALESCE(b.balance, 0) <> COALESCE(l.total, 0)
|
||||
LIMIT ?`},
|
||||
{WalletCheckLotVsAllocations, `
|
||||
SELECT l.organization_id::text AS organization_id, l.customer_id::text AS customer_id, l.currency,
|
||||
l.id::text AS object_id,
|
||||
l.original_amount - COALESCE(a.total, 0) AS expected,
|
||||
l.remaining_amount AS actual
|
||||
FROM wallet_lots l
|
||||
LEFT JOIN (SELECT lot_id, SUM(amount) AS total FROM wallet_lot_allocations GROUP BY lot_id) a ON a.lot_id = l.id
|
||||
WHERE l.original_amount - COALESCE(a.total, 0) <> l.remaining_amount
|
||||
LIMIT ?`},
|
||||
{WalletCheckDebitVsAllocations, `
|
||||
SELECT t.organization_id::text AS organization_id, t.customer_id::text AS customer_id, t.currency,
|
||||
t.id::text AS object_id,
|
||||
-t.amount AS expected,
|
||||
COALESCE(a.total, 0) AS actual
|
||||
FROM wallet_transactions t
|
||||
LEFT JOIN (SELECT transaction_id, SUM(amount) AS total FROM wallet_lot_allocations GROUP BY transaction_id) a ON a.transaction_id = t.id
|
||||
WHERE t.amount < 0 AND -t.amount <> COALESCE(a.total, 0)
|
||||
LIMIT ?`},
|
||||
{WalletCheckCreditVsLots, `
|
||||
SELECT t.organization_id::text AS organization_id, t.customer_id::text AS customer_id, t.currency,
|
||||
t.id::text AS object_id,
|
||||
t.amount AS expected,
|
||||
COALESCE(l.total, 0) AS actual
|
||||
FROM wallet_transactions t
|
||||
LEFT JOIN (SELECT source_transaction_id, SUM(original_amount) AS total FROM wallet_lots GROUP BY source_transaction_id) l ON l.source_transaction_id = t.id
|
||||
WHERE t.amount > 0 AND t.amount <> COALESCE(l.total, 0)
|
||||
LIMIT ?`},
|
||||
}
|
||||
|
||||
func (r *walletReconciliationRepository) FindDiscrepancies(ctx context.Context, limit int) ([]WalletDiscrepancy, error) {
|
||||
db := DBFromContext(ctx, r.db).WithContext(ctx)
|
||||
var found []WalletDiscrepancy
|
||||
for _, q := range walletReconciliationQueries {
|
||||
var rows []struct {
|
||||
OrganizationID string
|
||||
CustomerID string
|
||||
Currency string
|
||||
ObjectID *string
|
||||
Expected int64
|
||||
Actual int64
|
||||
}
|
||||
if err := db.Raw(q.sql, limit).Scan(&rows).Error; err != nil {
|
||||
return nil, fmt.Errorf("wallet reconciliation check %s failed: %w", q.check, err)
|
||||
}
|
||||
for _, row := range rows {
|
||||
d := WalletDiscrepancy{
|
||||
Check: q.check,
|
||||
Currency: row.Currency,
|
||||
Expected: row.Expected,
|
||||
Actual: row.Actual,
|
||||
}
|
||||
d.OrganizationID, _ = uuid.Parse(row.OrganizationID)
|
||||
d.CustomerID, _ = uuid.Parse(row.CustomerID)
|
||||
if row.ObjectID != nil {
|
||||
if id, err := uuid.Parse(*row.ObjectID); err == nil {
|
||||
d.ObjectID = &id
|
||||
}
|
||||
}
|
||||
found = append(found, d)
|
||||
}
|
||||
}
|
||||
return found, nil
|
||||
}
|
||||
@@ -0,0 +1,141 @@
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,212 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/sha256"
|
||||
"encoding/hex"
|
||||
"fmt"
|
||||
"sort"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"apskel-pos-be/internal/entities"
|
||||
"apskel-pos-be/internal/logger"
|
||||
"apskel-pos-be/internal/models"
|
||||
"apskel-pos-be/internal/repository"
|
||||
)
|
||||
|
||||
const (
|
||||
defaultWalletReconciliationInterval = 6 * time.Hour
|
||||
// Per check, so one systematic bug cannot flood the log or the notification.
|
||||
walletReconciliationLimit = 50
|
||||
)
|
||||
|
||||
type walletDiscrepancyFinder interface {
|
||||
FindDiscrepancies(ctx context.Context, limit int) ([]repository.WalletDiscrepancy, error)
|
||||
}
|
||||
|
||||
type organizationUserLister interface {
|
||||
GetByOrganizationID(ctx context.Context, organizationID uuid.UUID) ([]*entities.User, error)
|
||||
}
|
||||
|
||||
type notificationSender interface {
|
||||
Send(ctx context.Context, req *models.SendNotificationRequest) (*models.NotificationResponse, error)
|
||||
}
|
||||
|
||||
// WalletReconciliationJob periodically runs the §7.5 checks of
|
||||
// docs/prd-point-coin.md over every wallet (PC-108). It is silent while the data is
|
||||
// consistent. When it finds a discrepancy it logs each one and notifies the admins,
|
||||
// owners and managers of the organization concerned.
|
||||
//
|
||||
// An organization is notified again only when its set of discrepancies changes, so an
|
||||
// unfixed problem does not page the same people every run. That memory is in-process:
|
||||
// a restart notifies once more, and each running instance keeps its own.
|
||||
type WalletReconciliationJob struct {
|
||||
finder walletDiscrepancyFinder
|
||||
users organizationUserLister
|
||||
notifier notificationSender
|
||||
|
||||
mu sync.Mutex
|
||||
notified map[uuid.UUID]string // organization -> fingerprint last notified
|
||||
stopCh chan struct{}
|
||||
stopOnce sync.Once
|
||||
}
|
||||
|
||||
func NewWalletReconciliationJob(finder walletDiscrepancyFinder, users organizationUserLister, notifier notificationSender) *WalletReconciliationJob {
|
||||
return &WalletReconciliationJob{
|
||||
finder: finder,
|
||||
users: users,
|
||||
notifier: notifier,
|
||||
notified: make(map[uuid.UUID]string),
|
||||
stopCh: make(chan struct{}),
|
||||
}
|
||||
}
|
||||
|
||||
// Start runs the checks once now and then every interval, in the background.
|
||||
func (j *WalletReconciliationJob) Start(interval time.Duration) {
|
||||
if interval <= 0 {
|
||||
interval = defaultWalletReconciliationInterval
|
||||
}
|
||||
go func() {
|
||||
j.runLogged()
|
||||
ticker := time.NewTicker(interval)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
j.runLogged()
|
||||
case <-j.stopCh:
|
||||
return
|
||||
}
|
||||
}
|
||||
}()
|
||||
logger.NonContext.Infof("Wallet reconciliation job started (interval: %s)", interval)
|
||||
}
|
||||
|
||||
func (j *WalletReconciliationJob) Stop() {
|
||||
j.stopOnce.Do(func() { close(j.stopCh) })
|
||||
}
|
||||
|
||||
func (j *WalletReconciliationJob) runLogged() {
|
||||
if _, err := j.RunOnce(context.Background()); err != nil {
|
||||
logger.NonContext.Error("Wallet reconciliation failed to run", err)
|
||||
}
|
||||
}
|
||||
|
||||
// RunOnce runs every check, reports what it finds, and returns it.
|
||||
func (j *WalletReconciliationJob) RunOnce(ctx context.Context) ([]repository.WalletDiscrepancy, error) {
|
||||
found, err := j.finder.FindDiscrepancies(ctx, walletReconciliationLimit)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
byOrg := make(map[uuid.UUID][]repository.WalletDiscrepancy)
|
||||
for _, d := range found {
|
||||
fields := map[string]interface{}{
|
||||
"check": d.Check,
|
||||
"organization_id": d.OrganizationID.String(),
|
||||
"customer_id": d.CustomerID.String(),
|
||||
"currency": d.Currency,
|
||||
"expected": d.Expected,
|
||||
"actual": d.Actual,
|
||||
}
|
||||
if d.ObjectID != nil {
|
||||
fields["object_id"] = d.ObjectID.String()
|
||||
}
|
||||
logger.NonContext.WarnWithFields("Wallet reconciliation found a discrepancy", fields, nil)
|
||||
byOrg[d.OrganizationID] = append(byOrg[d.OrganizationID], d)
|
||||
}
|
||||
|
||||
j.mu.Lock()
|
||||
defer j.mu.Unlock()
|
||||
// Organizations that are clean again are forgotten, so a later problem notifies.
|
||||
for org := range j.notified {
|
||||
if _, still := byOrg[org]; !still {
|
||||
delete(j.notified, org)
|
||||
}
|
||||
}
|
||||
for org, discrepancies := range byOrg {
|
||||
fingerprint := walletDiscrepancyFingerprint(discrepancies)
|
||||
if j.notified[org] == fingerprint {
|
||||
continue
|
||||
}
|
||||
if err := j.notify(ctx, org, discrepancies); err != nil {
|
||||
logger.NonContext.Error(fmt.Sprintf("Wallet reconciliation could not notify organization %s", org), err)
|
||||
continue
|
||||
}
|
||||
j.notified[org] = fingerprint
|
||||
}
|
||||
return found, nil
|
||||
}
|
||||
|
||||
func (j *WalletReconciliationJob) notify(ctx context.Context, organizationID uuid.UUID, discrepancies []repository.WalletDiscrepancy) error {
|
||||
if organizationID == uuid.Nil {
|
||||
return fmt.Errorf("discrepancy without an organization")
|
||||
}
|
||||
users, err := j.users.GetByOrganizationID(ctx, organizationID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
var receivers []uuid.UUID
|
||||
for _, u := range users {
|
||||
switch u.Role {
|
||||
case entities.RoleAdmin, entities.RoleOwner, entities.RoleManager:
|
||||
receivers = append(receivers, u.ID)
|
||||
}
|
||||
}
|
||||
if len(receivers) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
perCheck := map[string]int{}
|
||||
customers := map[string]bool{}
|
||||
for _, d := range discrepancies {
|
||||
perCheck[d.Check]++
|
||||
customers[d.CustomerID.String()] = true
|
||||
}
|
||||
customerIDs := make([]string, 0, len(customers))
|
||||
for id := range customers {
|
||||
customerIDs = append(customerIDs, id)
|
||||
}
|
||||
sort.Strings(customerIDs)
|
||||
|
||||
_, err = j.notifier.Send(ctx, &models.SendNotificationRequest{
|
||||
Title: "Selisih saldo EnakPoint/EnakCoin terdeteksi",
|
||||
Body: fmt.Sprintf("Pemeriksaan rutin menemukan %d selisih pada saldo %d customer. Saldo belum dikoreksi otomatis; tim teknis perlu memeriksanya.",
|
||||
len(discrepancies), len(customerIDs)),
|
||||
Type: "system",
|
||||
Category: "wallet_reconciliation",
|
||||
Priority: entities.NotificationPriorityHigh,
|
||||
NotifiableType: "organization",
|
||||
NotifiableID: &organizationID,
|
||||
ReceiverIDs: receivers,
|
||||
Data: map[string]interface{}{
|
||||
"organization_id": organizationID.String(),
|
||||
"discrepancies": len(discrepancies),
|
||||
"per_check": perCheck,
|
||||
"customer_ids": customerIDs,
|
||||
},
|
||||
})
|
||||
return err
|
||||
}
|
||||
|
||||
// walletDiscrepancyFingerprint identifies a set of discrepancies regardless of order.
|
||||
func walletDiscrepancyFingerprint(discrepancies []repository.WalletDiscrepancy) string {
|
||||
keys := make([]string, 0, len(discrepancies))
|
||||
for _, d := range discrepancies {
|
||||
object := ""
|
||||
if d.ObjectID != nil {
|
||||
object = d.ObjectID.String()
|
||||
}
|
||||
keys = append(keys, fmt.Sprintf("%s|%s|%s|%s|%d|%d", d.Check, d.CustomerID, d.Currency, object, d.Expected, d.Actual))
|
||||
}
|
||||
sort.Strings(keys)
|
||||
h := sha256.New()
|
||||
for _, k := range keys {
|
||||
h.Write([]byte(k))
|
||||
h.Write([]byte{'\n'})
|
||||
}
|
||||
return hex.EncodeToString(h.Sum(nil))
|
||||
}
|
||||
@@ -0,0 +1,110 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
|
||||
"apskel-pos-be/internal/entities"
|
||||
"apskel-pos-be/internal/logger"
|
||||
"apskel-pos-be/internal/models"
|
||||
"apskel-pos-be/internal/repository"
|
||||
)
|
||||
|
||||
type discrepancyFinderFake struct {
|
||||
found []repository.WalletDiscrepancy
|
||||
}
|
||||
|
||||
func (f *discrepancyFinderFake) FindDiscrepancies(context.Context, int) ([]repository.WalletDiscrepancy, error) {
|
||||
return f.found, nil
|
||||
}
|
||||
|
||||
type orgUsersFake map[uuid.UUID][]*entities.User
|
||||
|
||||
func (f orgUsersFake) GetByOrganizationID(_ context.Context, org uuid.UUID) ([]*entities.User, error) {
|
||||
return f[org], nil
|
||||
}
|
||||
|
||||
type notifierFake struct {
|
||||
sent []*models.SendNotificationRequest
|
||||
}
|
||||
|
||||
func (f *notifierFake) Send(_ context.Context, req *models.SendNotificationRequest) (*models.NotificationResponse, error) {
|
||||
f.sent = append(f.sent, req)
|
||||
return &models.NotificationResponse{}, nil
|
||||
}
|
||||
|
||||
func TestWalletReconciliationJob(t *testing.T) {
|
||||
logger.Setup("fatal", "json")
|
||||
org := uuid.New()
|
||||
admin, owner, manager, cashier := uuid.New(), uuid.New(), uuid.New(), uuid.New()
|
||||
users := orgUsersFake{org: {
|
||||
{ID: admin, Role: entities.RoleAdmin},
|
||||
{ID: owner, Role: entities.RoleOwner},
|
||||
{ID: manager, Role: entities.RoleManager},
|
||||
{ID: cashier, Role: entities.RoleCashier},
|
||||
}}
|
||||
finder := &discrepancyFinderFake{}
|
||||
notifier := ¬ifierFake{}
|
||||
job := NewWalletReconciliationJob(finder, users, notifier)
|
||||
ctx := context.Background()
|
||||
|
||||
// Consistent data: nothing reported.
|
||||
found, err := job.RunOnce(ctx)
|
||||
require.NoError(t, err)
|
||||
assert.Empty(t, found)
|
||||
assert.Empty(t, notifier.sent)
|
||||
|
||||
// A discrepancy notifies the organization's admins, owners and managers.
|
||||
customer := uuid.New()
|
||||
lot := uuid.New()
|
||||
finder.found = []repository.WalletDiscrepancy{
|
||||
{Check: repository.WalletCheckBalanceVsLots, OrganizationID: org, CustomerID: customer, Currency: "POINT", Expected: 100, Actual: 105},
|
||||
{Check: repository.WalletCheckLotVsAllocations, OrganizationID: org, CustomerID: customer, Currency: "POINT", ObjectID: &lot, Expected: 50, Actual: 55},
|
||||
}
|
||||
found, err = job.RunOnce(ctx)
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, found, 2)
|
||||
require.Len(t, notifier.sent, 1)
|
||||
sent := notifier.sent[0]
|
||||
assert.ElementsMatch(t, []uuid.UUID{admin, owner, manager}, sent.ReceiverIDs, "cashiers are not told")
|
||||
assert.Equal(t, &org, sent.NotifiableID)
|
||||
assert.Equal(t, 2, sent.Data["discrepancies"])
|
||||
assert.Equal(t, []string{customer.String()}, sent.Data["customer_ids"])
|
||||
assert.Equal(t, entities.NotificationPriorityHigh, sent.Priority)
|
||||
|
||||
// The same problem, still unfixed and in a different order, does not notify again.
|
||||
finder.found = []repository.WalletDiscrepancy{finder.found[1], finder.found[0]}
|
||||
_, err = job.RunOnce(ctx)
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, notifier.sent, 1)
|
||||
|
||||
// A changed problem does.
|
||||
finder.found = finder.found[:1]
|
||||
_, err = job.RunOnce(ctx)
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, notifier.sent, 2)
|
||||
|
||||
// Once clean the organization is forgotten, so the same problem coming back
|
||||
// notifies again.
|
||||
previous := finder.found
|
||||
finder.found = nil
|
||||
_, err = job.RunOnce(ctx)
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, notifier.sent, 2)
|
||||
finder.found = previous
|
||||
_, err = job.RunOnce(ctx)
|
||||
require.NoError(t, err)
|
||||
assert.Len(t, notifier.sent, 3)
|
||||
}
|
||||
|
||||
func TestWalletReconciliationJobStartStop(t *testing.T) {
|
||||
logger.Setup("fatal", "json")
|
||||
job := NewWalletReconciliationJob(&discrepancyFinderFake{}, orgUsersFake{}, ¬ifierFake{})
|
||||
job.Start(0)
|
||||
job.Stop()
|
||||
job.Stop() // stopping twice is harmless
|
||||
}
|
||||
Reference in New Issue
Block a user