From 550122f29c44989f19e5bf50c4440d0ec420c9a1 Mon Sep 17 00:00:00 2001 From: efrilm Date: Wed, 30 Sep 2026 14:35:12 +0700 Subject: [PATCH] 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 --- internal/app/app.go | 4 +- internal/handler/customer_points_handler.go | 23 +++++ internal/models/wallet.go | 7 ++ .../processor/customer_points_processor.go | 8 ++ internal/processor/wallet_expiry_processor.go | 82 +++++++++++++++- .../processor/wallet_expiry_processor_test.go | 94 ++++++++++++++++++- internal/processor/wallet_move_db_test.go | 2 +- internal/processor/wallet_query_processor.go | 19 ++++ .../processor/wallet_query_processor_test.go | 21 +++++ .../repository/wallet_expiry_repository.go | 64 +++++++++++++ .../repository/wallet_query_repository.go | 19 ++++ internal/router/router.go | 1 + internal/router/router_test.go | 1 + internal/service/customer_points_service.go | 8 ++ internal/service/wallet_expiry_job.go | 20 +++- ...97_create_wallet_expiry_reminders.down.sql | 1 + ...0097_create_wallet_expiry_reminders.up.sql | 11 +++ 17 files changed, 374 insertions(+), 11 deletions(-) create mode 100644 migrations/000097_create_wallet_expiry_reminders.down.sql create mode 100644 migrations/000097_create_wallet_expiry_reminders.up.sql diff --git a/internal/app/app.go b/internal/app/app.go index 934cc47..df1bf1b 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -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) 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( - 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) validators := a.initValidators() diff --git a/internal/handler/customer_points_handler.go b/internal/handler/customer_points_handler.go index 818d914..0df8ff2 100644 --- a/internal/handler/customer_points_handler.go +++ b/internal/handler/customer_points_handler.go @@ -204,3 +204,26 @@ func walletErrorCode(err error) string { 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") +} diff --git a/internal/models/wallet.go b/internal/models/wallet.go index da38e30..9412444 100644 --- a/internal/models/wallet.go +++ b/internal/models/wallet.go @@ -157,3 +157,10 @@ type PointPaymentPreview struct { // Rupiah covered by MaxPoints. 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"` +} diff --git a/internal/processor/customer_points_processor.go b/internal/processor/customer_points_processor.go index 020c705..bd87c33 100644 --- a/internal/processor/customer_points_processor.go +++ b/internal/processor/customer_points_processor.go @@ -224,3 +224,11 @@ func (p *CustomerPointsProcessor) GetFerrisWheelGameAPI(ctx context.Context) (*m }, }, 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) +} diff --git a/internal/processor/wallet_expiry_processor.go b/internal/processor/wallet_expiry_processor.go index c6a4ba7..b322362 100644 --- a/internal/processor/wallet_expiry_processor.go +++ b/internal/processor/wallet_expiry_processor.go @@ -8,6 +8,7 @@ import ( "github.com/google/uuid" + "apskel-pos-be/internal/constants" "apskel-pos-be/internal/logger" "apskel-pos-be/internal/repository" ) @@ -23,19 +24,24 @@ const ( // part of their balance expires. 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 // (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 + settings organizationSettingsReader 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} +func NewWalletExpiryProcessor(repo repository.WalletExpiryRepository, settings organizationSettingsReader, wallet *WalletProcessor, tx TxRunner, notifier customerNotifier) *WalletExpiryProcessor { + return &WalletExpiryProcessor{repo: repo, settings: settings, wallet: wallet, tx: tx, notifier: notifier, now: time.Now} } type walletExpiredKey struct { @@ -116,3 +122,75 @@ func expiryDescription(amount int64, currency, sourceDescription string) string } 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()) +} diff --git a/internal/processor/wallet_expiry_processor_test.go b/internal/processor/wallet_expiry_processor_test.go index 28b22dd..d5bbbe8 100644 --- a/internal/processor/wallet_expiry_processor_test.go +++ b/internal/processor/wallet_expiry_processor_test.go @@ -20,6 +20,8 @@ import ( type walletExpiryRepoFake struct { wallet *walletRepoFake 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) { @@ -45,7 +47,7 @@ func (f *walletExpiryRepoFake) ListDueLots(_ context.Context, asOf time.Time, li func (e *walletMoveEnv) expiry(notifier customerNotifier) (*WalletExpiryProcessor, *walletExpiryRepoFake) { 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 } return p, repo } @@ -150,3 +152,93 @@ func TestWalletExpiry_NothingDue(t *testing.T) { assert.Zero(t, count) 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")) +} diff --git a/internal/processor/wallet_move_db_test.go b/internal/processor/wallet_move_db_test.go index d425cda..0808d2f 100644 --- a/internal/processor/wallet_move_db_test.go +++ b/internal/processor/wallet_move_db_test.go @@ -220,7 +220,7 @@ func TestWalletExpiry_TwoInstancesAgainstPostgres(t *testing.T) { wg.Add(1) go func(i int) { 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()) assert.NoError(t, err) counts[i] = n diff --git a/internal/processor/wallet_query_processor.go b/internal/processor/wallet_query_processor.go index 02ae40d..3a83fb6 100644 --- a/internal/processor/wallet_query_processor.go +++ b/internal/processor/wallet_query_processor.go @@ -315,3 +315,22 @@ func walletTransactionFilter(customerID uuid.UUID, q models.ListCustomerWalletTr } 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 +} diff --git a/internal/processor/wallet_query_processor_test.go b/internal/processor/wallet_query_processor_test.go index d990a2b..bf6def7 100644 --- a/internal/processor/wallet_query_processor_test.go +++ b/internal/processor/wallet_query_processor_test.go @@ -243,3 +243,24 @@ func TestWalletQueryProcessor_RejectsBadQueries(t *testing.T) { func (f *walletQueryRepoFake) OrganizationOutstanding(context.Context, uuid.UUID) (int64, int64, error) { 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) +} diff --git a/internal/repository/wallet_expiry_repository.go b/internal/repository/wallet_expiry_repository.go index a4ed31e..c7e0b99 100644 --- a/internal/repository/wallet_expiry_repository.go +++ b/internal/repository/wallet_expiry_repository.go @@ -20,6 +20,14 @@ type DueLot struct { 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 // F12). Balances only change through WalletProcessor. type WalletExpiryRepository interface { @@ -27,6 +35,17 @@ type WalletExpiryRepository interface { // 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) + + // 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 { @@ -52,3 +71,48 @@ func (r *walletExpiryRepository) ListDueLots(ctx context.Context, asOf time.Time } 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 +} diff --git a/internal/repository/wallet_query_repository.go b/internal/repository/wallet_query_repository.go index 0591693..ecd446f 100644 --- a/internal/repository/wallet_query_repository.go +++ b/internal/repository/wallet_query_repository.go @@ -47,6 +47,9 @@ type WalletQueryRepository interface { // NearestExpiring returns, per currency, the earliest day after asOf on which some // balance expires, and how much expires that day. 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(ctx context.Context, filter WalletTransactionFilter) ([]entities.WalletTransaction, int64, error) // 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 } + +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 +} diff --git a/internal/router/router.go b/internal/router/router.go index d5770b0..d3a4b6f 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -172,6 +172,7 @@ func (r *Router) addAppRoutes(rg *gin.Engine) { customer.GET("/tokens", r.customerPointsHandler.GetCustomerTokens) customer.GET("/wallet", r.customerPointsHandler.GetCustomerWallet) customer.GET("/wallet/transactions", r.customerPointsHandler.GetCustomerWalletTransactions) + customer.GET("/wallet/expiring", r.customerPointsHandler.GetCustomerWalletExpiring) customer.POST("/wallet/payment-code", r.customerPinHandler.IssuePaymentCode) customer.GET("/wallet/exchange/preview", r.customerWalletHandler.PreviewExchange) customer.POST("/wallet/exchange", r.customerWalletHandler.Exchange) diff --git a/internal/router/router_test.go b/internal/router/router_test.go index 0d377c9..40db753 100644 --- a/internal/router/router_test.go +++ b/internal/router/router_test.go @@ -28,6 +28,7 @@ func TestAllRoutesRegister(t *testing.T) { for _, want := range []string{ "GET /api/v1/customer/wallet", "GET /api/v1/customer/wallet/transactions", + "GET /api/v1/customer/wallet/expiring", "GET /api/v1/marketing/customers/:id/wallet", "POST /api/v1/marketing/customers/:id/wallet/adjust", "GET /api/v1/marketing/wallet-transactions/:id/trace", diff --git a/internal/service/customer_points_service.go b/internal/service/customer_points_service.go index a247845..6e85854 100644 --- a/internal/service/customer_points_service.go +++ b/internal/service/customer_points_service.go @@ -13,6 +13,7 @@ type CustomerPointsService interface { GetCustomerTokens(ctx context.Context, customerID string) (*models.GetCustomerTokensResponse, error) GetCustomerWallet(ctx context.Context, customerID string) (*models.GetCustomerWalletResponse, 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) 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) } + +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) +} diff --git a/internal/service/wallet_expiry_job.go b/internal/service/wallet_expiry_job.go index d671d15..f2a83c0 100644 --- a/internal/service/wallet_expiry_job.go +++ b/internal/service/wallet_expiry_job.go @@ -12,21 +12,23 @@ import ( // lot from staying past its expiry for more than about that long (PC-503). const defaultWalletExpiryInterval = 15 * time.Minute -type dueExpirer interface { +type walletExpiryWork interface { 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 // 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 + expirer walletExpiryWork stopCh chan struct{} stopOnce sync.Once } -func NewWalletExpiryJob(expirer dueExpirer) *WalletExpiryJob { +func NewWalletExpiryJob(expirer walletExpiryWork) *WalletExpiryJob { return &WalletExpiryJob{expirer: expirer, stopCh: make(chan struct{})} } @@ -54,7 +56,8 @@ func (j *WalletExpiryJob) Stop() { 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 { expired, err := j.expirer.ExpireDue(ctx) if err != nil { @@ -63,5 +66,12 @@ func (j *WalletExpiryJob) RunOnce(ctx context.Context) int { if expired > 0 { 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 } diff --git a/migrations/000097_create_wallet_expiry_reminders.down.sql b/migrations/000097_create_wallet_expiry_reminders.down.sql new file mode 100644 index 0000000..a154829 --- /dev/null +++ b/migrations/000097_create_wallet_expiry_reminders.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS wallet_expiry_reminders; diff --git a/migrations/000097_create_wallet_expiry_reminders.up.sql b/migrations/000097_create_wallet_expiry_reminders.up.sql new file mode 100644 index 0000000..74e500e --- /dev/null +++ b/migrations/000097_create_wallet_expiry_reminders.up.sql @@ -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) +);