diff --git a/internal/app/app.go b/internal/app/app.go index ee0c19b..934cc47 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -33,6 +33,7 @@ type App struct { omsetScheduler *service.OmsetMilestoneScheduler walletRecon *service.WalletReconciliationJob earningRetry *service.EarningBackfillJob + walletExpiry *service.WalletExpiryJob } func NewApp(db *gorm.DB, redisClient *redis.Client) *App { @@ -63,6 +64,9 @@ func (a *App) Initialize(cfg *config.Config) error { ) // Earns for paid orders whose earning failed at payment time (docs/prd-point-coin.md F3) a.earningRetry = service.NewEarningBackfillJob(processors.earningProcessor) + // Expires balances whose time is up (docs/prd-point-coin.md F12) + a.walletExpiry = service.NewWalletExpiryJob(processor.NewWalletExpiryProcessor( + repository.NewWalletExpiryRepository(a.db), processor.NewWalletProcessor(repos.walletRepo), repos.txManager, processors.customerDeviceProcessor)) services := a.initServices(processors, repos, cfg) validators := a.initValidators() @@ -178,6 +182,9 @@ func (a *App) Start(port string) error { if a.earningRetry != nil { a.earningRetry.Start(30 * time.Minute) } + if a.walletExpiry != nil { + a.walletExpiry.Start(15 * time.Minute) + } engine := a.router.Init() @@ -223,6 +230,9 @@ func (a *App) Shutdown() { if a.earningRetry != nil { a.earningRetry.Stop() } + if a.walletExpiry != nil { + a.walletExpiry.Stop() + } close(a.shutdown) } diff --git a/internal/processor/wallet_expiry_processor.go b/internal/processor/wallet_expiry_processor.go new file mode 100644 index 0000000..c6a4ba7 --- /dev/null +++ b/internal/processor/wallet_expiry_processor.go @@ -0,0 +1,118 @@ +package processor + +import ( + "context" + "fmt" + "strconv" + "time" + + "github.com/google/uuid" + + "apskel-pos-be/internal/logger" + "apskel-pos-be/internal/repository" +) + +const ( + // Lots expired per query; a run keeps going until nothing is due. + walletExpiryBatchSize = 500 + // Batches per run at most, so one run cannot run away. + walletExpiryMaxBatches = 40 +) + +// NotificationTypeWalletExpired is the data type of the push a customer gets when +// part of their balance expires. +const NotificationTypeWalletExpired = "WALLET_EXPIRED" + +// WalletExpiryProcessor takes what is left in lots whose expiry has passed +// (docs/prd-point-coin.md F12, PC-503). It is safe to run on several instances at +// once: every lot is expired under its wallet's lock with the key expire:{lot_id}. +type WalletExpiryProcessor struct { + repo repository.WalletExpiryRepository + wallet *WalletProcessor + tx TxRunner + notifier customerNotifier + now func() time.Time +} + +func NewWalletExpiryProcessor(repo repository.WalletExpiryRepository, wallet *WalletProcessor, tx TxRunner, notifier customerNotifier) *WalletExpiryProcessor { + return &WalletExpiryProcessor{repo: repo, wallet: wallet, tx: tx, notifier: notifier, now: time.Now} +} + +type walletExpiredKey struct { + customerID uuid.UUID + currency string +} + +// ExpireDue expires every lot due now and tells each customer how much of each +// currency they lost, in one push per currency. It returns how many lots it expired. +// A lot that fails is logged and left for the next run; it does not stop the others. +func (p *WalletExpiryProcessor) ExpireDue(ctx context.Context) (int, error) { + asOf := p.now() + expired := map[walletExpiredKey]int64{} + count := 0 + for batch := 0; batch < walletExpiryMaxBatches; batch++ { + due, err := p.repo.ListDueLots(ctx, asOf, walletExpiryBatchSize) + if err != nil { + p.notify(ctx, expired) + return count, err + } + progressed := false + for _, lot := range due { + var res *WalletResult + err := p.tx.WithTransaction(ctx, func(ctx context.Context) error { + var err error + res, err = p.wallet.ExpireLot(ctx, lot.ID, func(amount int64) string { + return expiryDescription(amount, lot.Currency, lot.SourceDescription) + }, asOf) + return err + }) + if err != nil { + logger.NonContext.Error(fmt.Sprintf("Could not expire wallet lot %s; it will be retried", lot.ID), err) + continue + } + if res == nil || res.Transaction == nil || res.Replayed { + // Another run got there first, or a payment used it up. + continue + } + progressed = true + count++ + expired[walletExpiredKey{lot.CustomerID, lot.Currency}] += -res.Transaction.Amount + } + // A short batch was the last; a batch that moved nothing would only come back + // the same, whether failing or taken by another instance. + if len(due) < walletExpiryBatchSize || !progressed { + break + } + } + p.notify(ctx, expired) + return count, nil +} + +// notify is best effort: the balance has already expired. +func (p *WalletExpiryProcessor) notify(ctx context.Context, expired map[walletExpiredKey]int64) { + if p.notifier == nil { + return + } + for key, amount := range expired { + name := walletCurrencyName(key.currency) + data := map[string]string{ + "type": NotificationTypeWalletExpired, + "currency": key.currency, + "amount": strconv.FormatInt(amount, 10), + } + body := fmt.Sprintf("%d %s kamu sudah kedaluwarsa.", amount, name) + if err := p.notifier.Notify(ctx, key.customerID, name+" kedaluwarsa", body, data); err != nil { + logger.NonContext.Error(fmt.Sprintf("Could not tell customer %s about expired %s", key.customerID, name), err) + } + } +} + +// expiryDescription is the EXPIRE row's frozen description (ยง8.1): +// "Kedaluwarsa: 150 EnakPoint dari Belanja #ORD-0098". +func expiryDescription(amount int64, currency, sourceDescription string) string { + description := fmt.Sprintf("Kedaluwarsa: %d %s", amount, walletCurrencyName(currency)) + if sourceDescription != "" { + description += " dari " + sourceDescription + } + return truncateRunes(description, walletDescriptionLimit) +} diff --git a/internal/processor/wallet_expiry_processor_test.go b/internal/processor/wallet_expiry_processor_test.go new file mode 100644 index 0000000..28b22dd --- /dev/null +++ b/internal/processor/wallet_expiry_processor_test.go @@ -0,0 +1,152 @@ +package processor + +import ( + "context" + "sort" + "testing" + "time" + + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "apskel-pos-be/internal/constants" + "apskel-pos-be/internal/entities" + "apskel-pos-be/internal/repository" +) + +// walletExpiryRepoFake lists the fake wallet's due lots the way the query does, plus +// any extra lots a test wants listed. +type walletExpiryRepoFake struct { + wallet *walletRepoFake + extra []repository.DueLot +} + +func (f *walletExpiryRepoFake) ListDueLots(_ context.Context, asOf time.Time, limit int) ([]repository.DueLot, error) { + descriptions := map[uuid.UUID]string{} + for _, tx := range f.wallet.transactions { + descriptions[tx.ID] = tx.Description + } + due := append([]repository.DueLot(nil), f.extra...) + for _, lot := range f.wallet.lots { + if lot.RemainingAmount > 0 && lot.ExpiresAt != nil && !lot.ExpiresAt.After(asOf) { + due = append(due, repository.DueLot{ + ID: lot.ID, CustomerID: lot.CustomerID, Currency: lot.Currency, Remaining: lot.RemainingAmount, + ExpiresAt: *lot.ExpiresAt, SourceDescription: descriptions[lot.SourceTransactionID], + }) + } + } + sort.SliceStable(due, func(i, j int) bool { return due[i].ExpiresAt.Before(due[j].ExpiresAt) }) + if len(due) > limit { + due = due[:limit] + } + return due, nil +} + +func (e *walletMoveEnv) expiry(notifier customerNotifier) (*WalletExpiryProcessor, *walletExpiryRepoFake) { + repo := &walletExpiryRepoFake{wallet: e.repo} + p := NewWalletExpiryProcessor(repo, e.p, txRunnerFake{}, notifier) + p.now = func() time.Time { return e.now } + return p, repo +} + +func TestWalletExpiry_ExpiresWhatIsDueAndTellsTheCustomer(t *testing.T) { + e := newWalletMoveEnv(t) + a := e.member("Anita", "081200005678") + ord := earn(a, 150, e.at(-time.Hour)) + ord.Description = "Belanja #ORD-0098" + due := e.credit(t, ord) + e.credit(t, earn(a, 50, e.at(-2*time.Hour))) + e.credit(t, earn(a, 70, e.at(time.Hour))) // not yet + e.credit(t, earn(a, 30, nil)) // never + e.earnCoins(t, a, 4, e.at(-time.Minute)) + // A payment already used part of the first lot; only the rest expires. + _, err := e.p.Debit(e.ctx, WalletDebitInput{WalletEntry: pay(a, 20).WalletEntry, PreferredLotIDs: []uuid.UUID{due.Lots[0].ID}}) + require.NoError(t, err) + + notifier := ¬ifierFake{} + p, _ := e.expiry(notifier) + count, err := p.ExpireDue(e.ctx) + require.NoError(t, err) + + assert.Equal(t, 3, count) + assert.Equal(t, int64(100), e.balance(t, a), "70 not due yet + 30 that never expires") + assert.Equal(t, int64(0), e.coinBalance(t, a)) + + var expire *entities.WalletTransaction + for _, tx := range e.repo.transactions { + if tx.Type == constants.WalletTxTypeExpire && tx.ReferenceID == due.Lots[0].ID { + expire = tx + } + } + require.NotNil(t, expire) + assert.Equal(t, int64(-130), expire.Amount) + assert.Equal(t, "Kedaluwarsa: 130 EnakPoint dari Belanja #ORD-0098", expire.Description) + assert.Equal(t, constants.WalletRefTypeLot, expire.ReferenceType) + assert.Equal(t, "expire:"+due.Lots[0].ID.String(), *expire.IdempotencyKey) + + // One push per currency, with the total. + pushes := notifier.pushes[a] + require.Len(t, pushes, 2) + byCurrency := map[string]pushFake{} + for _, p := range pushes { + byCurrency[p.data["currency"]] = p + } + assert.Equal(t, "EnakPoint kedaluwarsa", byCurrency["POINT"].title) + assert.Equal(t, "180 EnakPoint kamu sudah kedaluwarsa.", byCurrency["POINT"].body) + assert.Equal(t, NotificationTypeWalletExpired, byCurrency["POINT"].data["type"]) + assert.Equal(t, "4", byCurrency["COIN"].data["amount"]) +} + +func TestWalletExpiry_RunningAgainExpiresNothingMore(t *testing.T) { + e := newWalletMoveEnv(t) + a := e.member("Anita", "081200005678") + e.credit(t, earn(a, 150, e.at(-time.Hour))) + notifier := ¬ifierFake{} + + // Two instances, one after the other. + first, _ := e.expiry(notifier) + second, _ := e.expiry(notifier) + n1, err := first.ExpireDue(e.ctx) + require.NoError(t, err) + n2, err := second.ExpireDue(e.ctx) + require.NoError(t, err) + + assert.Equal(t, 1, n1) + assert.Equal(t, 0, n2) + assert.Len(t, notifier.pushes[a], 1) + var expires int + for _, tx := range e.repo.transactions { + if tx.Type == constants.WalletTxTypeExpire { + expires++ + } + } + assert.Equal(t, 1, expires) +} + +func TestWalletExpiry_OneFailingLotDoesNotStopTheOthers(t *testing.T) { + e := newWalletMoveEnv(t) + a := e.member("Anita", "081200005678") + e.credit(t, earn(a, 150, e.at(-time.Hour))) + p, repo := e.expiry(nil) + // A lot listed that ExpireLot cannot find. + repo.extra = []repository.DueLot{{ID: uuid.New(), CustomerID: a, Currency: "POINT", Remaining: 5, ExpiresAt: e.now.Add(-3 * time.Hour)}} + + count, err := p.ExpireDue(e.ctx) + require.NoError(t, err) + assert.Equal(t, 1, count) + assert.Equal(t, int64(0), e.balance(t, a)) +} + +func TestWalletExpiry_NothingDue(t *testing.T) { + e := newWalletMoveEnv(t) + a := e.member("Anita", "081200005678") + e.credit(t, earn(a, 150, e.at(time.Hour))) + notifier := ¬ifierFake{} + p, _ := e.expiry(notifier) + + count, err := p.ExpireDue(e.ctx) + require.NoError(t, err) + assert.Zero(t, count) + assert.Empty(t, notifier.pushes) +} diff --git a/internal/processor/wallet_move_db_test.go b/internal/processor/wallet_move_db_test.go index 0f654d5..d425cda 100644 --- a/internal/processor/wallet_move_db_test.go +++ b/internal/processor/wallet_move_db_test.go @@ -196,3 +196,42 @@ func TestWalletTrace_AgainstPostgres(t *testing.T) { _, 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), 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) +} diff --git a/internal/processor/wallet_processor.go b/internal/processor/wallet_processor.go index 28a22a2..17bb8de 100644 --- a/internal/processor/wallet_processor.go +++ b/internal/processor/wallet_processor.go @@ -149,6 +149,55 @@ func (p *WalletProcessor) FindTransaction(ctx context.Context, idempotencyKey st return p.repo.GetTransactionByIdempotencyKey(ctx, idempotencyKey) } +// ExpireLot takes what is left in a lot whose expiry has passed at asOf, as an EXPIRE +// row pointing at the lot (F12). It locks the wallet before reading the lot again, so +// it never races a payment for the same balance, and the key expire:{lot_id} makes a +// second run, on this instance or another, take nothing more. describe gives the +// row's description for the amount taken. It returns nil when +// there is nothing to take: the lot is empty, not due, or already expired. +func (p *WalletProcessor) ExpireLot(ctx context.Context, lotID uuid.UUID, describe func(amount int64) string, asOf time.Time) (*WalletResult, error) { + lot, err := p.getLot(ctx, lotID) + if err != nil { + return nil, err + } + if _, err := p.repo.LockWallet(ctx, lot.CustomerID); err != nil { + return nil, err + } + // Read again under the lock: a payment may have used it up meanwhile. + if lot, err = p.getLot(ctx, lotID); err != nil { + return nil, err + } + if lot.RemainingAmount == 0 || lot.ExpiresAt == nil || lot.ExpiresAt.After(asOf) { + return nil, nil + } + return p.Debit(ctx, WalletDebitInput{ + WalletEntry: WalletEntry{ + CustomerID: lot.CustomerID, + Currency: lot.Currency, + Type: constants.WalletTxTypeExpire, + Amount: lot.RemainingAmount, + ReferenceType: constants.WalletRefTypeLot, + ReferenceID: lot.ID, + Description: describe(lot.RemainingAmount), + Metadata: entities.Metadata{"expires_at": lot.ExpiresAt.UTC().Format(time.RFC3339)}, + IdempotencyKey: "expire:" + lot.ID.String(), + }, + // Exactly what the lot holds, from the lot itself, even though it has expired. + PreferredLotIDs: []uuid.UUID{lot.ID}, + }) +} + +func (p *WalletProcessor) getLot(ctx context.Context, lotID uuid.UUID) (*entities.WalletLot, error) { + lots, err := p.repo.GetLotsByIDs(ctx, []uuid.UUID{lotID}) + if err != nil { + return nil, err + } + if len(lots) == 0 { + return nil, fmt.Errorf("%w: lot %s does not exist", ErrWalletInvalidEntry, lotID) + } + return &lots[0], nil +} + // Credit adds Amount to the wallet and creates its lots. func (p *WalletProcessor) Credit(ctx context.Context, in WalletCreditInput) (*WalletResult, error) { if err := validateWalletEntry(&in.WalletEntry, true); err != nil { diff --git a/internal/repository/wallet_expiry_repository.go b/internal/repository/wallet_expiry_repository.go new file mode 100644 index 0000000..a4ed31e --- /dev/null +++ b/internal/repository/wallet_expiry_repository.go @@ -0,0 +1,54 @@ +package repository + +import ( + "context" + "fmt" + "time" + + "github.com/google/uuid" + "gorm.io/gorm" +) + +// DueLot is a lot whose expiry has passed and that still holds something. +type DueLot struct { + ID uuid.UUID + CustomerID uuid.UUID + Currency string + Remaining int64 + ExpiresAt time.Time + // The description of the row that created the lot, for the EXPIRE row's. + SourceDescription string +} + +// WalletExpiryRepository finds what the expiry job has to do (docs/prd-point-coin.md +// F12). Balances only change through WalletProcessor. +type WalletExpiryRepository interface { + // ListDueLots returns up to limit lots due at asOf, the longest overdue first. It + // takes no lock: locking a lot before its wallet would deadlock against payments, + // which lock the wallet first. WalletProcessor.ExpireLot locks and reads again. + ListDueLots(ctx context.Context, asOf time.Time, limit int) ([]DueLot, error) +} + +type walletExpiryRepository struct { + db *gorm.DB +} + +func NewWalletExpiryRepository(db *gorm.DB) WalletExpiryRepository { + return &walletExpiryRepository{db: db} +} + +func (r *walletExpiryRepository) ListDueLots(ctx context.Context, asOf time.Time, limit int) ([]DueLot, error) { + var lots []DueLot + err := DBFromContext(ctx, r.db).WithContext(ctx).Raw(` + SELECT l.id, l.customer_id, l.currency, l.remaining_amount AS remaining, l.expires_at, + t.description AS source_description + FROM wallet_lots l + JOIN wallet_transactions t ON t.id = l.source_transaction_id + WHERE l.remaining_amount > 0 AND l.expires_at <= ? + ORDER BY l.expires_at, l.id + LIMIT ?`, asOf, limit).Scan(&lots).Error + if err != nil { + return nil, fmt.Errorf("failed to list due wallet lots: %w", err) + } + return lots, nil +} diff --git a/internal/service/wallet_expiry_job.go b/internal/service/wallet_expiry_job.go new file mode 100644 index 0000000..d671d15 --- /dev/null +++ b/internal/service/wallet_expiry_job.go @@ -0,0 +1,67 @@ +package service + +import ( + "context" + "sync" + "time" + + "apskel-pos-be/internal/logger" +) + +// Lots expire at the end of their day, so running every quarter of an hour keeps any +// lot from staying past its expiry for more than about that long (PC-503). +const defaultWalletExpiryInterval = 15 * time.Minute + +type dueExpirer interface { + ExpireDue(ctx context.Context) (int, error) +} + +// WalletExpiryJob expires the balances whose time is up (docs/prd-point-coin.md F12). +// Unlike OmsetMilestoneScheduler it keeps no state in memory: several instances can +// run it at once, and a restart repeats nothing, because every lot is expired under +// its wallet's lock with an idempotency key. +type WalletExpiryJob struct { + expirer dueExpirer + stopCh chan struct{} + stopOnce sync.Once +} + +func NewWalletExpiryJob(expirer dueExpirer) *WalletExpiryJob { + return &WalletExpiryJob{expirer: expirer, stopCh: make(chan struct{})} +} + +func (j *WalletExpiryJob) Start(interval time.Duration) { + if interval <= 0 { + interval = defaultWalletExpiryInterval + } + go func() { + j.RunOnce(context.Background()) + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-ticker.C: + j.RunOnce(context.Background()) + case <-j.stopCh: + return + } + } + }() + logger.NonContext.Infof("Wallet expiry job started (interval: %s)", interval) +} + +func (j *WalletExpiryJob) Stop() { + j.stopOnce.Do(func() { close(j.stopCh) }) +} + +// RunOnce expires what is due and returns how many lots it expired. +func (j *WalletExpiryJob) RunOnce(ctx context.Context) int { + expired, err := j.expirer.ExpireDue(ctx) + if err != nil { + logger.NonContext.Error("Wallet expiry failed to run", err) + } + if expired > 0 { + logger.NonContext.Infof("Wallet expiry expired %d lots", expired) + } + return expired +}