diff --git a/internal/app/app.go b/internal/app/app.go index 3e7f1be..5fba67d 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -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) } diff --git a/internal/repository/wallet_reconciliation_repository.go b/internal/repository/wallet_reconciliation_repository.go new file mode 100644 index 0000000..af41535 --- /dev/null +++ b/internal/repository/wallet_reconciliation_repository.go @@ -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 +} diff --git a/internal/repository/wallet_reconciliation_repository_test.go b/internal/repository/wallet_reconciliation_repository_test.go new file mode 100644 index 0000000..4b530f4 --- /dev/null +++ b/internal/repository/wallet_reconciliation_repository_test.go @@ -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) + } +} diff --git a/internal/service/wallet_reconciliation_job.go b/internal/service/wallet_reconciliation_job.go new file mode 100644 index 0000000..2b9aebd --- /dev/null +++ b/internal/service/wallet_reconciliation_job.go @@ -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)) +} diff --git a/internal/service/wallet_reconciliation_job_test.go b/internal/service/wallet_reconciliation_job_test.go new file mode 100644 index 0000000..c12c8e6 --- /dev/null +++ b/internal/service/wallet_reconciliation_job_test.go @@ -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 +}