feat(loyalty): remind customers before balances expire
Adds what the customer sees of expiry (docs/prd-point-coin.md F6, F12, PC-504). GET /customer/wallet/expiring lists everything that will expire, per currency and day, soonest first. GET /customer/wallet already had the nearest expiry per currency. The expiry job now also sends reminders, with the settings of note N4 as decided: once, reminder_days before (7 by default, 0 for none), per currency. A customer gets one FCM push per currency and expiry day, however many lots make it up: "150 EnakPoint akan kedaluwarsa pada 31 Okt 2026. Pakai sebelum hangus.", with type WALLET_EXPIRING, the currency, amount and expiry_date in its data. Reminders cover whatever falls within the window, so a run that was missed catches up rather than skipping a day. Migration 000097 adds wallet_expiry_reminders, one row per customer, currency and expiry day. The row is written before the push is sent, so several instances of the job or a restart never remind twice; a push that then fails is logged and not retried. Lots that expire later on the same day as an earlier reminder are not reminded of again. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
d293786cde
commit
550122f29c
+2
-2
@@ -64,9 +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)
|
// Expires balances whose time is up and reminds customers before (docs/prd-point-coin.md F12)
|
||||||
a.walletExpiry = service.NewWalletExpiryJob(processor.NewWalletExpiryProcessor(
|
a.walletExpiry = service.NewWalletExpiryJob(processor.NewWalletExpiryProcessor(
|
||||||
repository.NewWalletExpiryRepository(a.db), processor.NewWalletProcessor(repos.walletRepo), repos.txManager, processors.customerDeviceProcessor))
|
repository.NewWalletExpiryRepository(a.db), processors.loyaltySettingsProcessor, 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()
|
||||||
|
|||||||
@@ -204,3 +204,26 @@ func walletErrorCode(err error) string {
|
|||||||
return constants.InternalServerErrorCode
|
return constants.InternalServerErrorCode
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// GetCustomerWalletExpiring is GET /customer/wallet/expiring: what will expire, per
|
||||||
|
// currency and day (docs/prd-point-coin.md F6).
|
||||||
|
func (h *CustomerPointsHandler) GetCustomerWalletExpiring(c *gin.Context) {
|
||||||
|
ctx := c.Request.Context()
|
||||||
|
customerID, ok := c.Get("customer_id")
|
||||||
|
customerIDStr, isString := customerID.(string)
|
||||||
|
if !ok || !isString {
|
||||||
|
util.HandleResponse(c.Writer, c.Request, contract.BuildErrorResponse([]*contract.ResponseError{
|
||||||
|
contract.NewResponseError(constants.ValidationErrorCode, constants.AuthHandlerEntity, "Customer ID not found"),
|
||||||
|
}), "CustomerPointsHandler::GetCustomerWalletExpiring")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
response, err := h.customerPointsService.GetCustomerWalletExpiring(ctx, customerIDStr)
|
||||||
|
if err != nil {
|
||||||
|
logger.FromContext(ctx).WithError(err).Error("CustomerPointsHandler::GetCustomerWalletExpiring -> service call failed")
|
||||||
|
util.HandleResponse(c.Writer, c.Request, contract.BuildErrorResponse([]*contract.ResponseError{
|
||||||
|
contract.NewResponseError(walletErrorCode(err), constants.RequestEntity, err.Error()),
|
||||||
|
}), "CustomerPointsHandler::GetCustomerWalletExpiring")
|
||||||
|
return
|
||||||
|
}
|
||||||
|
util.HandleResponse(c.Writer, c.Request, contract.BuildSuccessResponse(response), "CustomerPointsHandler::GetCustomerWalletExpiring")
|
||||||
|
}
|
||||||
|
|||||||
@@ -157,3 +157,10 @@ type PointPaymentPreview struct {
|
|||||||
// Rupiah covered by MaxPoints.
|
// Rupiah covered by MaxPoints.
|
||||||
MaxAmount int64 `json:"max_amount"`
|
MaxAmount int64 `json:"max_amount"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// CustomerWalletExpiringList is GET /customer/wallet/expiring (docs/prd-point-coin.md
|
||||||
|
// F6): everything that will expire, per currency and day, soonest first.
|
||||||
|
type CustomerWalletExpiringList struct {
|
||||||
|
Point []CustomerWalletExpiring `json:"point"`
|
||||||
|
Coin []CustomerWalletExpiring `json:"coin"`
|
||||||
|
}
|
||||||
|
|||||||
@@ -224,3 +224,11 @@ func (p *CustomerPointsProcessor) GetFerrisWheelGameAPI(ctx context.Context) (*m
|
|||||||
},
|
},
|
||||||
}, nil
|
}, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (p *CustomerPointsProcessor) GetCustomerWalletExpiringAPI(ctx context.Context, customerID string) (*models.CustomerWalletExpiringList, error) {
|
||||||
|
id, err := parseWalletCustomerID(customerID)
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
return p.walletQuery.Expiring(ctx, id)
|
||||||
|
}
|
||||||
|
|||||||
@@ -8,6 +8,7 @@ import (
|
|||||||
|
|
||||||
"github.com/google/uuid"
|
"github.com/google/uuid"
|
||||||
|
|
||||||
|
"apskel-pos-be/internal/constants"
|
||||||
"apskel-pos-be/internal/logger"
|
"apskel-pos-be/internal/logger"
|
||||||
"apskel-pos-be/internal/repository"
|
"apskel-pos-be/internal/repository"
|
||||||
)
|
)
|
||||||
@@ -23,19 +24,24 @@ const (
|
|||||||
// part of their balance expires.
|
// part of their balance expires.
|
||||||
const NotificationTypeWalletExpired = "WALLET_EXPIRED"
|
const NotificationTypeWalletExpired = "WALLET_EXPIRED"
|
||||||
|
|
||||||
|
// NotificationTypeWalletExpiring is the data type of the reminder a customer gets
|
||||||
|
// before part of their balance expires.
|
||||||
|
const NotificationTypeWalletExpiring = "WALLET_EXPIRING"
|
||||||
|
|
||||||
// WalletExpiryProcessor takes what is left in lots whose expiry has passed
|
// 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
|
// (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}.
|
// once: every lot is expired under its wallet's lock with the key expire:{lot_id}.
|
||||||
type WalletExpiryProcessor struct {
|
type WalletExpiryProcessor struct {
|
||||||
repo repository.WalletExpiryRepository
|
repo repository.WalletExpiryRepository
|
||||||
|
settings organizationSettingsReader
|
||||||
wallet *WalletProcessor
|
wallet *WalletProcessor
|
||||||
tx TxRunner
|
tx TxRunner
|
||||||
notifier customerNotifier
|
notifier customerNotifier
|
||||||
now func() time.Time
|
now func() time.Time
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewWalletExpiryProcessor(repo repository.WalletExpiryRepository, wallet *WalletProcessor, tx TxRunner, notifier customerNotifier) *WalletExpiryProcessor {
|
func NewWalletExpiryProcessor(repo repository.WalletExpiryRepository, settings organizationSettingsReader, wallet *WalletProcessor, tx TxRunner, notifier customerNotifier) *WalletExpiryProcessor {
|
||||||
return &WalletExpiryProcessor{repo: repo, wallet: wallet, tx: tx, notifier: notifier, now: time.Now}
|
return &WalletExpiryProcessor{repo: repo, settings: settings, wallet: wallet, tx: tx, notifier: notifier, now: time.Now}
|
||||||
}
|
}
|
||||||
|
|
||||||
type walletExpiredKey struct {
|
type walletExpiredKey struct {
|
||||||
@@ -116,3 +122,75 @@ func expiryDescription(amount int64, currency, sourceDescription string) string
|
|||||||
}
|
}
|
||||||
return truncateRunes(description, walletDescriptionLimit)
|
return truncateRunes(description, walletDescriptionLimit)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// SendReminders tells customers, reminder_days before, how much of their balance
|
||||||
|
// expires on a day (F12): one push per customer, currency and expiry day, however many
|
||||||
|
// lots make it up. A reminder is recorded before it is sent, so another instance or a
|
||||||
|
// later run never sends it again; a push that then fails is logged and not retried.
|
||||||
|
// It returns how many reminders it sent.
|
||||||
|
func (p *WalletExpiryProcessor) SendReminders(ctx context.Context) (int, error) {
|
||||||
|
now := p.now()
|
||||||
|
organizations, err := p.repo.OrganizationsWithUpcomingExpiry(ctx, now)
|
||||||
|
if err != nil {
|
||||||
|
return 0, err
|
||||||
|
}
|
||||||
|
sent := 0
|
||||||
|
for _, organizationID := range organizations {
|
||||||
|
settings, err := p.settings.Organization(ctx, organizationID)
|
||||||
|
if err != nil {
|
||||||
|
logger.NonContext.Error(fmt.Sprintf("Could not read the expiry settings of organization %s; its reminders wait for the next run", organizationID), err)
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
for _, currency := range []string{constants.WalletCurrencyPoint, constants.WalletCurrencyCoin} {
|
||||||
|
days := ExpirySettings(settings, currency).ReminderDays
|
||||||
|
if days <= 0 {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
until := endOfWalletDay(walletDay(now).AddDate(0, 0, int(days)))
|
||||||
|
upcoming, err := p.repo.UpcomingUnreminded(ctx, organizationID, currency, now, *until)
|
||||||
|
if err != nil {
|
||||||
|
return sent, err
|
||||||
|
}
|
||||||
|
for _, u := range upcoming {
|
||||||
|
first, err := p.repo.MarkReminded(ctx, u, currency)
|
||||||
|
if err != nil {
|
||||||
|
return sent, err
|
||||||
|
}
|
||||||
|
if !first {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
sent++
|
||||||
|
p.remind(ctx, u, currency)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return sent, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (p *WalletExpiryProcessor) remind(ctx context.Context, u repository.UpcomingExpiry, currency string) {
|
||||||
|
if p.notifier == nil {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
name := walletCurrencyName(currency)
|
||||||
|
body := fmt.Sprintf("%d %s akan kedaluwarsa pada %s. Pakai sebelum hangus.", u.Amount, name, formatWalletDate(u.Date))
|
||||||
|
data := map[string]string{
|
||||||
|
"type": NotificationTypeWalletExpiring,
|
||||||
|
"currency": currency,
|
||||||
|
"amount": strconv.FormatInt(u.Amount, 10),
|
||||||
|
"expiry_date": u.Date,
|
||||||
|
}
|
||||||
|
if err := p.notifier.Notify(ctx, u.CustomerID, name+" akan kedaluwarsa", body, data); err != nil {
|
||||||
|
logger.NonContext.Error(fmt.Sprintf("Could not remind customer %s of expiring %s", u.CustomerID, name), err)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
var walletMonthNames = [...]string{"Jan", "Feb", "Mar", "Apr", "Mei", "Jun", "Jul", "Agu", "Sep", "Okt", "Nov", "Des"}
|
||||||
|
|
||||||
|
// formatWalletDate writes a YYYY-MM-DD date the way the apps do: "31 Okt 2026".
|
||||||
|
func formatWalletDate(date string) string {
|
||||||
|
d, err := time.Parse("2006-01-02", date)
|
||||||
|
if err != nil {
|
||||||
|
return date
|
||||||
|
}
|
||||||
|
return fmt.Sprintf("%d %s %d", d.Day(), walletMonthNames[d.Month()-1], d.Year())
|
||||||
|
}
|
||||||
|
|||||||
@@ -20,6 +20,8 @@ import (
|
|||||||
type walletExpiryRepoFake struct {
|
type walletExpiryRepoFake struct {
|
||||||
wallet *walletRepoFake
|
wallet *walletRepoFake
|
||||||
extra []repository.DueLot
|
extra []repository.DueLot
|
||||||
|
// customer/currency/date of the reminders recorded.
|
||||||
|
reminded map[string]bool
|
||||||
}
|
}
|
||||||
|
|
||||||
func (f *walletExpiryRepoFake) ListDueLots(_ context.Context, asOf time.Time, limit int) ([]repository.DueLot, error) {
|
func (f *walletExpiryRepoFake) ListDueLots(_ context.Context, asOf time.Time, limit int) ([]repository.DueLot, error) {
|
||||||
@@ -45,7 +47,7 @@ func (f *walletExpiryRepoFake) ListDueLots(_ context.Context, asOf time.Time, li
|
|||||||
|
|
||||||
func (e *walletMoveEnv) expiry(notifier customerNotifier) (*WalletExpiryProcessor, *walletExpiryRepoFake) {
|
func (e *walletMoveEnv) expiry(notifier customerNotifier) (*WalletExpiryProcessor, *walletExpiryRepoFake) {
|
||||||
repo := &walletExpiryRepoFake{wallet: e.repo}
|
repo := &walletExpiryRepoFake{wallet: e.repo}
|
||||||
p := NewWalletExpiryProcessor(repo, e.p, txRunnerFake{}, notifier)
|
p := NewWalletExpiryProcessor(repo, e, e.p, txRunnerFake{}, notifier)
|
||||||
p.now = func() time.Time { return e.now }
|
p.now = func() time.Time { return e.now }
|
||||||
return p, repo
|
return p, repo
|
||||||
}
|
}
|
||||||
@@ -150,3 +152,93 @@ func TestWalletExpiry_NothingDue(t *testing.T) {
|
|||||||
assert.Zero(t, count)
|
assert.Zero(t, count)
|
||||||
assert.Empty(t, notifier.pushes)
|
assert.Empty(t, notifier.pushes)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (f *walletExpiryRepoFake) OrganizationsWithUpcomingExpiry(_ context.Context, asOf time.Time) ([]uuid.UUID, error) {
|
||||||
|
seen := map[uuid.UUID]bool{}
|
||||||
|
var out []uuid.UUID
|
||||||
|
for _, lot := range f.wallet.lots {
|
||||||
|
if lot.RemainingAmount > 0 && lot.ExpiresAt != nil && lot.ExpiresAt.After(asOf) && !seen[lot.OrganizationID] {
|
||||||
|
seen[lot.OrganizationID] = true
|
||||||
|
out = append(out, lot.OrganizationID)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *walletExpiryRepoFake) UpcomingUnreminded(_ context.Context, organizationID uuid.UUID, currency string, asOf, until time.Time) ([]repository.UpcomingExpiry, error) {
|
||||||
|
sums := map[[2]string]int64{}
|
||||||
|
var order [][2]string
|
||||||
|
for _, lot := range f.wallet.lots {
|
||||||
|
if lot.OrganizationID != organizationID || lot.Currency != currency || lot.RemainingAmount == 0 ||
|
||||||
|
lot.ExpiresAt == nil || !lot.ExpiresAt.After(asOf) || lot.ExpiresAt.After(until) {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
key := [2]string{lot.CustomerID.String(), lot.ExpiresAt.In(walletDisplayLocation).Format("2006-01-02")}
|
||||||
|
if f.reminded[key[0]+"/"+currency+"/"+key[1]] {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
if _, ok := sums[key]; !ok {
|
||||||
|
order = append(order, key)
|
||||||
|
}
|
||||||
|
sums[key] += lot.RemainingAmount
|
||||||
|
}
|
||||||
|
var out []repository.UpcomingExpiry
|
||||||
|
for _, key := range order {
|
||||||
|
out = append(out, repository.UpcomingExpiry{CustomerID: uuid.MustParse(key[0]), Date: key[1], Amount: sums[key]})
|
||||||
|
}
|
||||||
|
return out, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (f *walletExpiryRepoFake) MarkReminded(_ context.Context, u repository.UpcomingExpiry, currency string) (bool, error) {
|
||||||
|
if f.reminded == nil {
|
||||||
|
f.reminded = map[string]bool{}
|
||||||
|
}
|
||||||
|
key := u.CustomerID.String() + "/" + currency + "/" + u.Date
|
||||||
|
if f.reminded[key] {
|
||||||
|
return false, nil
|
||||||
|
}
|
||||||
|
f.reminded[key] = true
|
||||||
|
return true, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestWalletExpiry_RemindsOncePerDayBeforeExpiry(t *testing.T) {
|
||||||
|
e := newWalletMoveEnv(t)
|
||||||
|
e.now = wib(2026, 10, 25, 9, 0)
|
||||||
|
e.settings.PointExpiry.ReminderDays = 7
|
||||||
|
e.settings.CoinExpiry.ReminderDays = 0 // no reminders for EnakCoin
|
||||||
|
a := e.member("Anita", "081200005678")
|
||||||
|
oct31 := wib(2026, 10, 31, 23, 59)
|
||||||
|
nov30 := wib(2026, 11, 30, 23, 59)
|
||||||
|
e.credit(t, earn(a, 100, &oct31))
|
||||||
|
e.credit(t, earn(a, 50, &oct31))
|
||||||
|
e.credit(t, earn(a, 70, &nov30)) // too far off yet
|
||||||
|
e.earnCoins(t, a, 5, &oct31)
|
||||||
|
notifier := ¬ifierFake{}
|
||||||
|
p, _ := e.expiry(notifier)
|
||||||
|
|
||||||
|
sent, err := p.SendReminders(e.ctx)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, 1, sent)
|
||||||
|
require.Len(t, notifier.pushes[a], 1)
|
||||||
|
push := notifier.pushes[a][0]
|
||||||
|
assert.Equal(t, "EnakPoint akan kedaluwarsa", push.title)
|
||||||
|
assert.Equal(t, "150 EnakPoint akan kedaluwarsa pada 31 Okt 2026. Pakai sebelum hangus.", push.body)
|
||||||
|
assert.Equal(t, map[string]string{"type": NotificationTypeWalletExpiring, "currency": "POINT", "amount": "150", "expiry_date": "2026-10-31"}, push.data)
|
||||||
|
|
||||||
|
// The next run, on this instance or another, sends nothing again.
|
||||||
|
again, err := p.SendReminders(e.ctx)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Zero(t, again)
|
||||||
|
|
||||||
|
// Once 30 Nov comes within seven days, it gets its own reminder.
|
||||||
|
e.now = wib(2026, 11, 23, 9, 0)
|
||||||
|
sent, err = p.SendReminders(e.ctx)
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, 1, sent)
|
||||||
|
assert.Equal(t, "70", notifier.pushes[a][1].data["amount"])
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestFormatWalletDate(t *testing.T) {
|
||||||
|
assert.Equal(t, "31 Okt 2026", formatWalletDate("2026-10-31"))
|
||||||
|
assert.Equal(t, "1 Mei 2027", formatWalletDate("2027-05-01"))
|
||||||
|
}
|
||||||
|
|||||||
@@ -220,7 +220,7 @@ func TestWalletExpiry_TwoInstancesAgainstPostgres(t *testing.T) {
|
|||||||
wg.Add(1)
|
wg.Add(1)
|
||||||
go func(i int) {
|
go func(i int) {
|
||||||
defer wg.Done()
|
defer wg.Done()
|
||||||
p := NewWalletExpiryProcessor(repository.NewWalletExpiryRepository(db), wallet, txm, nil)
|
p := NewWalletExpiryProcessor(repository.NewWalletExpiryRepository(db), fixedOrganizationSettings{}, wallet, txm, nil)
|
||||||
n, err := p.ExpireDue(context.Background())
|
n, err := p.ExpireDue(context.Background())
|
||||||
assert.NoError(t, err)
|
assert.NoError(t, err)
|
||||||
counts[i] = n
|
counts[i] = n
|
||||||
|
|||||||
@@ -315,3 +315,22 @@ func walletTransactionFilter(customerID uuid.UUID, q models.ListCustomerWalletTr
|
|||||||
}
|
}
|
||||||
return filter, page, nil
|
return filter, page, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Expiring is GET /customer/wallet/expiring: what will expire, grouped by day (F6).
|
||||||
|
func (p *WalletQueryProcessor) Expiring(ctx context.Context, customerID uuid.UUID) (*models.CustomerWalletExpiringList, error) {
|
||||||
|
rows, err := p.repo.ExpiringByDay(ctx, customerID, p.now())
|
||||||
|
if err != nil {
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
list := &models.CustomerWalletExpiringList{Point: []models.CustomerWalletExpiring{}, Coin: []models.CustomerWalletExpiring{}}
|
||||||
|
for _, row := range rows {
|
||||||
|
item := models.CustomerWalletExpiring{Amount: row.Amount, Date: row.Date}
|
||||||
|
switch row.Currency {
|
||||||
|
case constants.WalletCurrencyPoint:
|
||||||
|
list.Point = append(list.Point, item)
|
||||||
|
case constants.WalletCurrencyCoin:
|
||||||
|
list.Coin = append(list.Coin, item)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return list, nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -243,3 +243,24 @@ func TestWalletQueryProcessor_RejectsBadQueries(t *testing.T) {
|
|||||||
func (f *walletQueryRepoFake) OrganizationOutstanding(context.Context, uuid.UUID) (int64, int64, error) {
|
func (f *walletQueryRepoFake) OrganizationOutstanding(context.Context, uuid.UUID) (int64, int64, error) {
|
||||||
return 0, 0, nil
|
return 0, 0, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (f *walletQueryRepoFake) ExpiringByDay(context.Context, uuid.UUID, time.Time) ([]repository.WalletExpiringAmount, error) {
|
||||||
|
return f.expiring, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestWalletQueryProcessor_ExpiringGroupsByCurrencyAndDay(t *testing.T) {
|
||||||
|
repo := &walletQueryRepoFake{org: uuid.New(), expiring: []repository.WalletExpiringAmount{
|
||||||
|
{Currency: "POINT", Date: "2026-10-31", Amount: 150},
|
||||||
|
{Currency: "COIN", Date: "2026-10-31", Amount: 4},
|
||||||
|
{Currency: "POINT", Date: "2026-12-31", Amount: 200},
|
||||||
|
}}
|
||||||
|
got, err := newWalletQueryTest(repo, nil).Expiring(context.Background(), uuid.New())
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.Equal(t, []models.CustomerWalletExpiring{{Amount: 150, Date: "2026-10-31"}, {Amount: 200, Date: "2026-12-31"}}, got.Point)
|
||||||
|
assert.Equal(t, []models.CustomerWalletExpiring{{Amount: 4, Date: "2026-10-31"}}, got.Coin)
|
||||||
|
|
||||||
|
empty, err := newWalletQueryTest(&walletQueryRepoFake{org: uuid.New()}, nil).Expiring(context.Background(), uuid.New())
|
||||||
|
require.NoError(t, err)
|
||||||
|
assert.NotNil(t, empty.Point, "an empty list, not null")
|
||||||
|
assert.NotNil(t, empty.Coin)
|
||||||
|
}
|
||||||
|
|||||||
@@ -20,6 +20,14 @@ type DueLot struct {
|
|||||||
SourceDescription string
|
SourceDescription string
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// UpcomingExpiry is how much of a customer's balance expires on one day.
|
||||||
|
type UpcomingExpiry struct {
|
||||||
|
CustomerID uuid.UUID
|
||||||
|
// A calendar date in walletDisplayTimeZone, formatted YYYY-MM-DD.
|
||||||
|
Date string
|
||||||
|
Amount int64
|
||||||
|
}
|
||||||
|
|
||||||
// WalletExpiryRepository finds what the expiry job has to do (docs/prd-point-coin.md
|
// WalletExpiryRepository finds what the expiry job has to do (docs/prd-point-coin.md
|
||||||
// F12). Balances only change through WalletProcessor.
|
// F12). Balances only change through WalletProcessor.
|
||||||
type WalletExpiryRepository interface {
|
type WalletExpiryRepository interface {
|
||||||
@@ -27,6 +35,17 @@ type WalletExpiryRepository interface {
|
|||||||
// takes no lock: locking a lot before its wallet would deadlock against payments,
|
// takes no lock: locking a lot before its wallet would deadlock against payments,
|
||||||
// which lock the wallet first. WalletProcessor.ExpireLot locks and reads again.
|
// which lock the wallet first. WalletProcessor.ExpireLot locks and reads again.
|
||||||
ListDueLots(ctx context.Context, asOf time.Time, limit int) ([]DueLot, error)
|
ListDueLots(ctx context.Context, asOf time.Time, limit int) ([]DueLot, error)
|
||||||
|
|
||||||
|
// OrganizationsWithUpcomingExpiry lists the organizations that have balance
|
||||||
|
// expiring after asOf.
|
||||||
|
OrganizationsWithUpcomingExpiry(ctx context.Context, asOf time.Time) ([]uuid.UUID, error)
|
||||||
|
// UpcomingUnreminded sums, per customer and expiry day, the balance of one currency
|
||||||
|
// of an organization expiring after asOf and up to until, leaving out the days the
|
||||||
|
// customer has already been reminded of.
|
||||||
|
UpcomingUnreminded(ctx context.Context, organizationID uuid.UUID, currency string, asOf, until time.Time) ([]UpcomingExpiry, error)
|
||||||
|
// MarkReminded records a reminder, and reports false when it was already recorded,
|
||||||
|
// by this run or another.
|
||||||
|
MarkReminded(ctx context.Context, reminder UpcomingExpiry, currency string) (bool, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
type walletExpiryRepository struct {
|
type walletExpiryRepository struct {
|
||||||
@@ -52,3 +71,48 @@ func (r *walletExpiryRepository) ListDueLots(ctx context.Context, asOf time.Time
|
|||||||
}
|
}
|
||||||
return lots, nil
|
return lots, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (r *walletExpiryRepository) OrganizationsWithUpcomingExpiry(ctx context.Context, asOf time.Time) ([]uuid.UUID, error) {
|
||||||
|
var ids []uuid.UUID
|
||||||
|
err := DBFromContext(ctx, r.db).WithContext(ctx).Raw(`
|
||||||
|
SELECT DISTINCT organization_id FROM wallet_lots
|
||||||
|
WHERE remaining_amount > 0 AND expires_at > ?`, asOf).Scan(&ids).Error
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("failed to list organizations with expiring balances: %w", err)
|
||||||
|
}
|
||||||
|
return ids, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *walletExpiryRepository) UpcomingUnreminded(ctx context.Context, organizationID uuid.UUID, currency string, asOf, until time.Time) ([]UpcomingExpiry, error) {
|
||||||
|
var rows []UpcomingExpiry
|
||||||
|
err := DBFromContext(ctx, r.db).WithContext(ctx).Raw(`
|
||||||
|
WITH by_day AS (
|
||||||
|
SELECT customer_id, (expires_at AT TIME ZONE ?)::date AS day, SUM(remaining_amount) AS amount
|
||||||
|
FROM wallet_lots
|
||||||
|
WHERE organization_id = ? AND currency = ? AND remaining_amount > 0
|
||||||
|
AND expires_at > ? AND expires_at <= ?
|
||||||
|
GROUP BY customer_id, day
|
||||||
|
)
|
||||||
|
SELECT d.customer_id, to_char(d.day, 'YYYY-MM-DD') AS date, d.amount
|
||||||
|
FROM by_day d
|
||||||
|
LEFT JOIN wallet_expiry_reminders w
|
||||||
|
ON w.customer_id = d.customer_id AND w.currency = ? AND w.expiry_date = d.day
|
||||||
|
WHERE w.customer_id IS NULL
|
||||||
|
ORDER BY d.day, d.customer_id`,
|
||||||
|
walletDisplayTimeZone, organizationID, currency, asOf, until, currency).Scan(&rows).Error
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("failed to list upcoming expiry: %w", err)
|
||||||
|
}
|
||||||
|
return rows, nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func (r *walletExpiryRepository) MarkReminded(ctx context.Context, reminder UpcomingExpiry, currency string) (bool, error) {
|
||||||
|
res := DBFromContext(ctx, r.db).WithContext(ctx).Exec(`
|
||||||
|
INSERT INTO wallet_expiry_reminders (customer_id, currency, expiry_date, amount)
|
||||||
|
VALUES (?, ?, ?::date, ?)
|
||||||
|
ON CONFLICT DO NOTHING`, reminder.CustomerID, currency, reminder.Date, reminder.Amount)
|
||||||
|
if res.Error != nil {
|
||||||
|
return false, fmt.Errorf("failed to record expiry reminder: %w", res.Error)
|
||||||
|
}
|
||||||
|
return res.RowsAffected == 1, nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -47,6 +47,9 @@ type WalletQueryRepository interface {
|
|||||||
// NearestExpiring returns, per currency, the earliest day after asOf on which some
|
// NearestExpiring returns, per currency, the earliest day after asOf on which some
|
||||||
// balance expires, and how much expires that day.
|
// balance expires, and how much expires that day.
|
||||||
NearestExpiring(ctx context.Context, customerID uuid.UUID, asOf time.Time) ([]WalletExpiringAmount, error)
|
NearestExpiring(ctx context.Context, customerID uuid.UUID, asOf time.Time) ([]WalletExpiringAmount, error)
|
||||||
|
// ExpiringByDay returns, per currency and day, everything that expires after asOf,
|
||||||
|
// soonest first.
|
||||||
|
ExpiringByDay(ctx context.Context, customerID uuid.UUID, asOf time.Time) ([]WalletExpiringAmount, error)
|
||||||
// ListTransactions returns a page of the ledger, newest first, and the total count.
|
// ListTransactions returns a page of the ledger, newest first, and the total count.
|
||||||
ListTransactions(ctx context.Context, filter WalletTransactionFilter) ([]entities.WalletTransaction, int64, error)
|
ListTransactions(ctx context.Context, filter WalletTransactionFilter) ([]entities.WalletTransaction, int64, error)
|
||||||
// OrganizationOutstanding sums every wallet balance of an organization.
|
// OrganizationOutstanding sums every wallet balance of an organization.
|
||||||
@@ -181,3 +184,19 @@ func (r *walletQueryRepository) OrganizationOutstanding(ctx context.Context, org
|
|||||||
}
|
}
|
||||||
return totals.Points, totals.Coins, nil
|
return totals.Points, totals.Coins, nil
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (r *walletQueryRepository) ExpiringByDay(ctx context.Context, customerID uuid.UUID, asOf time.Time) ([]WalletExpiringAmount, error) {
|
||||||
|
var rows []WalletExpiringAmount
|
||||||
|
err := DBFromContext(ctx, r.db).WithContext(ctx).Raw(`
|
||||||
|
SELECT currency, to_char((expires_at AT TIME ZONE ?)::date, 'YYYY-MM-DD') AS date,
|
||||||
|
SUM(remaining_amount) AS amount
|
||||||
|
FROM wallet_lots
|
||||||
|
WHERE customer_id = ? AND remaining_amount > 0 AND expires_at > ?
|
||||||
|
GROUP BY currency, date
|
||||||
|
ORDER BY date, currency`, walletDisplayTimeZone, customerID, asOf).
|
||||||
|
Scan(&rows).Error
|
||||||
|
if err != nil {
|
||||||
|
return nil, fmt.Errorf("failed to list expiring wallet balance: %w", err)
|
||||||
|
}
|
||||||
|
return rows, nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -172,6 +172,7 @@ func (r *Router) addAppRoutes(rg *gin.Engine) {
|
|||||||
customer.GET("/tokens", r.customerPointsHandler.GetCustomerTokens)
|
customer.GET("/tokens", r.customerPointsHandler.GetCustomerTokens)
|
||||||
customer.GET("/wallet", r.customerPointsHandler.GetCustomerWallet)
|
customer.GET("/wallet", r.customerPointsHandler.GetCustomerWallet)
|
||||||
customer.GET("/wallet/transactions", r.customerPointsHandler.GetCustomerWalletTransactions)
|
customer.GET("/wallet/transactions", r.customerPointsHandler.GetCustomerWalletTransactions)
|
||||||
|
customer.GET("/wallet/expiring", r.customerPointsHandler.GetCustomerWalletExpiring)
|
||||||
customer.POST("/wallet/payment-code", r.customerPinHandler.IssuePaymentCode)
|
customer.POST("/wallet/payment-code", r.customerPinHandler.IssuePaymentCode)
|
||||||
customer.GET("/wallet/exchange/preview", r.customerWalletHandler.PreviewExchange)
|
customer.GET("/wallet/exchange/preview", r.customerWalletHandler.PreviewExchange)
|
||||||
customer.POST("/wallet/exchange", r.customerWalletHandler.Exchange)
|
customer.POST("/wallet/exchange", r.customerWalletHandler.Exchange)
|
||||||
|
|||||||
@@ -28,6 +28,7 @@ func TestAllRoutesRegister(t *testing.T) {
|
|||||||
for _, want := range []string{
|
for _, want := range []string{
|
||||||
"GET /api/v1/customer/wallet",
|
"GET /api/v1/customer/wallet",
|
||||||
"GET /api/v1/customer/wallet/transactions",
|
"GET /api/v1/customer/wallet/transactions",
|
||||||
|
"GET /api/v1/customer/wallet/expiring",
|
||||||
"GET /api/v1/marketing/customers/:id/wallet",
|
"GET /api/v1/marketing/customers/:id/wallet",
|
||||||
"POST /api/v1/marketing/customers/:id/wallet/adjust",
|
"POST /api/v1/marketing/customers/:id/wallet/adjust",
|
||||||
"GET /api/v1/marketing/wallet-transactions/:id/trace",
|
"GET /api/v1/marketing/wallet-transactions/:id/trace",
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ type CustomerPointsService interface {
|
|||||||
GetCustomerTokens(ctx context.Context, customerID string) (*models.GetCustomerTokensResponse, error)
|
GetCustomerTokens(ctx context.Context, customerID string) (*models.GetCustomerTokensResponse, error)
|
||||||
GetCustomerWallet(ctx context.Context, customerID string) (*models.GetCustomerWalletResponse, error)
|
GetCustomerWallet(ctx context.Context, customerID string) (*models.GetCustomerWalletResponse, error)
|
||||||
GetCustomerWalletTransactions(ctx context.Context, customerID string, query models.ListCustomerWalletTransactionsQuery) (*models.PaginatedResponse[models.CustomerWalletTransaction], error)
|
GetCustomerWalletTransactions(ctx context.Context, customerID string, query models.ListCustomerWalletTransactionsQuery) (*models.PaginatedResponse[models.CustomerWalletTransaction], error)
|
||||||
|
GetCustomerWalletExpiring(ctx context.Context, customerID string) (*models.CustomerWalletExpiringList, error)
|
||||||
GetCustomerGames(ctx context.Context) (*models.GetCustomerGamesResponse, error)
|
GetCustomerGames(ctx context.Context) (*models.GetCustomerGamesResponse, error)
|
||||||
GetFerrisWheelGame(ctx context.Context) (*models.GetFerrisWheelGameResponse, error)
|
GetFerrisWheelGame(ctx context.Context) (*models.GetFerrisWheelGameResponse, error)
|
||||||
}
|
}
|
||||||
@@ -90,3 +91,10 @@ func (s *customerPointsService) GetCustomerWalletTransactions(ctx context.Contex
|
|||||||
}
|
}
|
||||||
return s.customerPointsProcessor.GetCustomerWalletTransactionsAPI(ctx, customerID, query)
|
return s.customerPointsProcessor.GetCustomerWalletTransactionsAPI(ctx, customerID, query)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (s *customerPointsService) GetCustomerWalletExpiring(ctx context.Context, customerID string) (*models.CustomerWalletExpiringList, error) {
|
||||||
|
if customerID == "" {
|
||||||
|
return nil, fmt.Errorf("customer ID is required")
|
||||||
|
}
|
||||||
|
return s.customerPointsProcessor.GetCustomerWalletExpiringAPI(ctx, customerID)
|
||||||
|
}
|
||||||
|
|||||||
@@ -12,21 +12,23 @@ import (
|
|||||||
// lot from staying past its expiry for more than about that long (PC-503).
|
// lot from staying past its expiry for more than about that long (PC-503).
|
||||||
const defaultWalletExpiryInterval = 15 * time.Minute
|
const defaultWalletExpiryInterval = 15 * time.Minute
|
||||||
|
|
||||||
type dueExpirer interface {
|
type walletExpiryWork interface {
|
||||||
ExpireDue(ctx context.Context) (int, error)
|
ExpireDue(ctx context.Context) (int, error)
|
||||||
|
SendReminders(ctx context.Context) (int, error)
|
||||||
}
|
}
|
||||||
|
|
||||||
// WalletExpiryJob expires the balances whose time is up (docs/prd-point-coin.md F12).
|
// WalletExpiryJob expires the balances whose time is up and reminds customers of what
|
||||||
|
// is about to (docs/prd-point-coin.md F12, PC-503, PC-504).
|
||||||
// Unlike OmsetMilestoneScheduler it keeps no state in memory: several instances can
|
// 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
|
// run it at once, and a restart repeats nothing, because every lot is expired under
|
||||||
// its wallet's lock with an idempotency key.
|
// its wallet's lock with an idempotency key.
|
||||||
type WalletExpiryJob struct {
|
type WalletExpiryJob struct {
|
||||||
expirer dueExpirer
|
expirer walletExpiryWork
|
||||||
stopCh chan struct{}
|
stopCh chan struct{}
|
||||||
stopOnce sync.Once
|
stopOnce sync.Once
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewWalletExpiryJob(expirer dueExpirer) *WalletExpiryJob {
|
func NewWalletExpiryJob(expirer walletExpiryWork) *WalletExpiryJob {
|
||||||
return &WalletExpiryJob{expirer: expirer, stopCh: make(chan struct{})}
|
return &WalletExpiryJob{expirer: expirer, stopCh: make(chan struct{})}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -54,7 +56,8 @@ func (j *WalletExpiryJob) Stop() {
|
|||||||
j.stopOnce.Do(func() { close(j.stopCh) })
|
j.stopOnce.Do(func() { close(j.stopCh) })
|
||||||
}
|
}
|
||||||
|
|
||||||
// RunOnce expires what is due and returns how many lots it expired.
|
// RunOnce expires what is due, sends the reminders that are due, and returns how many
|
||||||
|
// lots it expired.
|
||||||
func (j *WalletExpiryJob) RunOnce(ctx context.Context) int {
|
func (j *WalletExpiryJob) RunOnce(ctx context.Context) int {
|
||||||
expired, err := j.expirer.ExpireDue(ctx)
|
expired, err := j.expirer.ExpireDue(ctx)
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -63,5 +66,12 @@ func (j *WalletExpiryJob) RunOnce(ctx context.Context) int {
|
|||||||
if expired > 0 {
|
if expired > 0 {
|
||||||
logger.NonContext.Infof("Wallet expiry expired %d lots", expired)
|
logger.NonContext.Infof("Wallet expiry expired %d lots", expired)
|
||||||
}
|
}
|
||||||
|
reminded, err := j.expirer.SendReminders(ctx)
|
||||||
|
if err != nil {
|
||||||
|
logger.NonContext.Error("Wallet expiry reminders failed to run", err)
|
||||||
|
}
|
||||||
|
if reminded > 0 {
|
||||||
|
logger.NonContext.Infof("Wallet expiry sent %d reminders", reminded)
|
||||||
|
}
|
||||||
return expired
|
return expired
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1 @@
|
|||||||
|
DROP TABLE IF EXISTS wallet_expiry_reminders;
|
||||||
@@ -0,0 +1,11 @@
|
|||||||
|
-- Which expiry reminders have gone out (docs/prd-point-coin.md F12): one per
|
||||||
|
-- customer, currency and expiry day. The row is written before the push is sent, so
|
||||||
|
-- several instances of the job, or a restart, never remind twice.
|
||||||
|
CREATE TABLE wallet_expiry_reminders (
|
||||||
|
customer_id UUID NOT NULL REFERENCES customers(id) ON DELETE CASCADE,
|
||||||
|
currency VARCHAR(10) NOT NULL CHECK (currency IN ('POINT','COIN')),
|
||||||
|
expiry_date DATE NOT NULL,
|
||||||
|
amount BIGINT NOT NULL,
|
||||||
|
created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
|
||||||
|
PRIMARY KEY (customer_id, currency, expiry_date)
|
||||||
|
);
|
||||||
Reference in New Issue
Block a user