feat(loyalty): expire balances whose time is up

Adds the expiry job (docs/prd-point-coin.md F12, PC-503).

Every 15 minutes it lists the lots whose expiry has passed and that still
hold something, the longest overdue first, 500 at a time, and expires each
in its own transaction through WalletProcessor.ExpireLot: lock the wallet,
read the lot again, and take what is left with an EXPIRE row pointing at
the lot, keyed expire:{lot_id}. The description is frozen as
"Kedaluwarsa: 130 EnakPoint dari Belanja #ORD-0098", using the amount read
under the lock. Lots expire at the end of their day, so none stays past it
for more than about a quarter of an hour.

It is safe on several instances and across restarts, keeping no state in
memory as OmsetMilestoneScheduler does. Selecting the lots FOR UPDATE SKIP
LOCKED, as PC-503 suggested, would lock a lot before its wallet and
deadlock against payments, which lock the wallet first; instead the
listing takes no lock, and the wallet lock plus the idempotency key make a
second instance find the lot empty or the key used and take nothing.

A lot that fails is logged and retried on the next run without stopping
the others. Each customer gets one FCM push per currency with the total
that expired ("180 EnakPoint kamu sudah kedaluwarsa.", type
WALLET_EXPIRED).

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
efrilm
2026-09-30 14:31:38 +07:00
co-authored by Claude Opus 5.5
parent 4d63673a25
commit d293786cde
7 changed files with 489 additions and 0 deletions
+10
View File
@@ -33,6 +33,7 @@ type App struct {
omsetScheduler *service.OmsetMilestoneScheduler omsetScheduler *service.OmsetMilestoneScheduler
walletRecon *service.WalletReconciliationJob walletRecon *service.WalletReconciliationJob
earningRetry *service.EarningBackfillJob earningRetry *service.EarningBackfillJob
walletExpiry *service.WalletExpiryJob
} }
func NewApp(db *gorm.DB, redisClient *redis.Client) *App { 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) // Earns for paid orders whose earning failed at payment time (docs/prd-point-coin.md F3)
a.earningRetry = service.NewEarningBackfillJob(processors.earningProcessor) 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) services := a.initServices(processors, repos, cfg)
validators := a.initValidators() validators := a.initValidators()
@@ -178,6 +182,9 @@ func (a *App) Start(port string) error {
if a.earningRetry != nil { if a.earningRetry != nil {
a.earningRetry.Start(30 * time.Minute) a.earningRetry.Start(30 * time.Minute)
} }
if a.walletExpiry != nil {
a.walletExpiry.Start(15 * time.Minute)
}
engine := a.router.Init() engine := a.router.Init()
@@ -223,6 +230,9 @@ func (a *App) Shutdown() {
if a.earningRetry != nil { if a.earningRetry != nil {
a.earningRetry.Stop() a.earningRetry.Stop()
} }
if a.walletExpiry != nil {
a.walletExpiry.Stop()
}
close(a.shutdown) close(a.shutdown)
} }
@@ -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)
}
@@ -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 := &notifierFake{}
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 := &notifierFake{}
// 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 := &notifierFake{}
p, _ := e.expiry(notifier)
count, err := p.ExpireDue(e.ctx)
require.NoError(t, err)
assert.Zero(t, count)
assert.Empty(t, notifier.pushes)
}
+39
View File
@@ -196,3 +196,42 @@ func TestWalletTrace_AgainstPostgres(t *testing.T) {
_, err = NewWalletTraceProcessor(repository.NewWalletTraceRepository(db)).Trace(context.Background(), uuid.New(), payment.Transaction.ID) _, err = NewWalletTraceProcessor(repository.NewWalletTraceRepository(db)).Trace(context.Background(), uuid.New(), payment.Transaction.ID)
assert.ErrorIs(t, err, repository.ErrWalletTransactionNotFound) 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)
}
+49
View File
@@ -149,6 +149,55 @@ func (p *WalletProcessor) FindTransaction(ctx context.Context, idempotencyKey st
return p.repo.GetTransactionByIdempotencyKey(ctx, idempotencyKey) 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. // Credit adds Amount to the wallet and creates its lots.
func (p *WalletProcessor) Credit(ctx context.Context, in WalletCreditInput) (*WalletResult, error) { func (p *WalletProcessor) Credit(ctx context.Context, in WalletCreditInput) (*WalletResult, error) {
if err := validateWalletEntry(&in.WalletEntry, true); err != nil { if err := validateWalletEntry(&in.WalletEntry, true); err != nil {
@@ -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
}
+67
View File
@@ -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
}