feat(loyalty): notify transfer recipients through FCM

The recipient of a transfer now gets a push through FCM instead of a
WhatsApp message (docs/prd-point-coin.md F5).

Customers had nowhere to keep FCM tokens: user_devices only holds staff
devices. Migration 000096 adds customer_devices, and the customer app
registers with PUT /customer/devices { device_id, fcm_token, platform,
app_version } after login and whenever FCM refreshes the token, and
unregisters with DELETE /customer/devices/:device_id on logout. A token
belongs to one customer only: registering it takes it away from whoever
had it on that phone before, so they stop getting this customer's
notifications.

The push goes to every device of the recipient after the commit, titled
"EnakPoint masuk" or "EnakCoin masuk", with type WALLET_TRANSFER_IN, the
TRANSFER_IN transaction id, the group id, the currency and the amount in
its data so the app can open it. A retried transfer sends nothing again. It
stays best effort: no device, FCM not configured or FCM failing is logged
and never undoes the transfer.

The app builds one FCM client and shares it between staff notifications
and customer pushes.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
efrilm
2026-09-30 12:30:14 +07:00
co-authored by Claude Opus 5.5
parent 8bf1d5c1a8
commit bf9651e221
14 changed files with 526 additions and 46 deletions
@@ -0,0 +1,77 @@
package processor
import (
"context"
"errors"
"fmt"
"strings"
"time"
"github.com/google/uuid"
"apskel-pos-be/internal/logger"
"apskel-pos-be/internal/repository"
)
// ErrInvalidCustomerDevice wraps every rejection of a device registration.
var ErrInvalidCustomerDevice = errors.New("invalid customer device")
type customerPushSender interface {
SendMulticastNotification(ctx context.Context, tokens []string, title string, body string, data map[string]string) error
}
// CustomerDeviceProcessor keeps the customer app's FCM tokens and sends push
// notifications to a customer's devices.
type CustomerDeviceProcessor struct {
repo repository.CustomerDeviceRepository
// Nil when FCM is not configured; notifications are then skipped.
fcm customerPushSender
now func() time.Time
}
func NewCustomerDeviceProcessor(repo repository.CustomerDeviceRepository, fcm customerPushSender) *CustomerDeviceProcessor {
return &CustomerDeviceProcessor{repo: repo, fcm: fcm, now: time.Now}
}
// Register saves the FCM token the app got for this device. The app calls it after
// login and whenever FCM gives it a new token.
func (p *CustomerDeviceProcessor) Register(ctx context.Context, device repository.CustomerDevice) error {
device.DeviceID = strings.TrimSpace(device.DeviceID)
device.FCMToken = strings.TrimSpace(device.FCMToken)
switch {
case device.DeviceID == "" || len(device.DeviceID) > 255:
return fmt.Errorf("%w: device_id is required, at most 255 characters", ErrInvalidCustomerDevice)
case device.FCMToken == "" || len(device.FCMToken) > 512:
return fmt.Errorf("%w: fcm_token is required, at most 512 characters", ErrInvalidCustomerDevice)
}
if device.Platform != nil {
platform := strings.ToLower(strings.TrimSpace(*device.Platform))
if platform != "android" && platform != "ios" && platform != "web" {
return fmt.Errorf("%w: platform must be android, ios or web", ErrInvalidCustomerDevice)
}
device.Platform = &platform
}
return p.repo.Register(ctx, device, p.now())
}
// Unregister forgets a device, so it stops getting the customer's notifications.
func (p *CustomerDeviceProcessor) Unregister(ctx context.Context, customerID uuid.UUID, deviceID string) error {
return p.repo.Unregister(ctx, customerID, strings.TrimSpace(deviceID))
}
// Notify pushes a notification to every device of the customer through FCM. A
// customer without a registered device gets nothing, which is not an error.
func (p *CustomerDeviceProcessor) Notify(ctx context.Context, customerID uuid.UUID, title, body string, data map[string]string) error {
if p.fcm == nil {
logger.NonContext.Info(fmt.Sprintf("FCM is not configured; not notifying customer %s", customerID))
return nil
}
tokens, err := p.repo.ListTokens(ctx, customerID)
if err != nil {
return err
}
if len(tokens) == 0 {
return nil
}
return p.fcm.SendMulticastNotification(ctx, tokens, title, body, data)
}
@@ -0,0 +1,136 @@
package processor
import (
"context"
"errors"
"testing"
"time"
"github.com/google/uuid"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"apskel-pos-be/internal/repository"
)
// customerDeviceRepoFake keeps devices the way the table does: one row per customer
// and device, and a token belongs to one row only.
type customerDeviceRepoFake struct{ devices []repository.CustomerDevice }
func (f *customerDeviceRepoFake) Register(_ context.Context, d repository.CustomerDevice, _ time.Time) error {
kept := f.devices[:0]
for _, existing := range f.devices {
sameRow := existing.CustomerID == d.CustomerID && existing.DeviceID == d.DeviceID
if !sameRow && existing.FCMToken != d.FCMToken {
kept = append(kept, existing)
}
}
f.devices = append(kept, d)
return nil
}
func (f *customerDeviceRepoFake) Unregister(_ context.Context, customerID uuid.UUID, deviceID string) error {
kept := f.devices[:0]
for _, d := range f.devices {
if d.CustomerID != customerID || d.DeviceID != deviceID {
kept = append(kept, d)
}
}
f.devices = kept
return nil
}
func (f *customerDeviceRepoFake) ListTokens(_ context.Context, customerID uuid.UUID) ([]string, error) {
var tokens []string
for _, d := range f.devices {
if d.CustomerID == customerID {
tokens = append(tokens, d.FCMToken)
}
}
return tokens, nil
}
type fcmFake struct {
tokens [][]string
title string
body string
data map[string]string
err error
}
func (f *fcmFake) SendMulticastNotification(_ context.Context, tokens []string, title, body string, data map[string]string) error {
f.tokens = append(f.tokens, tokens)
f.title, f.body, f.data = title, body, data
return f.err
}
func TestCustomerDevice_NotifiesEveryDeviceOfTheCustomer(t *testing.T) {
repo, fcm := &customerDeviceRepoFake{}, &fcmFake{}
p := NewCustomerDeviceProcessor(repo, fcm)
ctx := context.Background()
budi, anita := uuid.New(), uuid.New()
require.NoError(t, p.Register(ctx, repository.CustomerDevice{CustomerID: budi, DeviceID: "phone", FCMToken: "t1", Platform: ptr("Android")}))
require.NoError(t, p.Register(ctx, repository.CustomerDevice{CustomerID: budi, DeviceID: "tablet", FCMToken: "t2"}))
// A refreshed token replaces the old one of the same device.
require.NoError(t, p.Register(ctx, repository.CustomerDevice{CustomerID: budi, DeviceID: "phone", FCMToken: "t1b"}))
require.NoError(t, p.Register(ctx, repository.CustomerDevice{CustomerID: anita, DeviceID: "phone", FCMToken: "t3"}))
require.NoError(t, p.Notify(ctx, budi, "EnakPoint masuk", "Kamu menerima 10 EnakPoint", map[string]string{"type": "X"}))
assert.Equal(t, [][]string{{"t2", "t1b"}}, fcm.tokens)
assert.Equal(t, "EnakPoint masuk", fcm.title)
assert.Equal(t, map[string]string{"type": "X"}, fcm.data)
}
func TestCustomerDevice_TokenMovesToWhoeverLogsInOnThePhone(t *testing.T) {
repo, fcm := &customerDeviceRepoFake{}, &fcmFake{}
p := NewCustomerDeviceProcessor(repo, fcm)
ctx := context.Background()
budi, anita := uuid.New(), uuid.New()
require.NoError(t, p.Register(ctx, repository.CustomerDevice{CustomerID: budi, DeviceID: "phone", FCMToken: "shared"}))
require.NoError(t, p.Register(ctx, repository.CustomerDevice{CustomerID: anita, DeviceID: "phone", FCMToken: "shared"}))
// Budi's notifications no longer reach the phone Anita is now logged in on.
require.NoError(t, p.Notify(ctx, budi, "t", "b", nil))
require.NoError(t, p.Notify(ctx, anita, "t", "b", nil))
assert.Equal(t, [][]string{{"shared"}}, fcm.tokens)
}
func TestCustomerDevice_NothingToSend(t *testing.T) {
repo, fcm := &customerDeviceRepoFake{}, &fcmFake{}
ctx := context.Background()
customer := uuid.New()
// No device registered, or logged out: nothing is sent and nothing fails.
require.NoError(t, NewCustomerDeviceProcessor(repo, fcm).Notify(ctx, customer, "t", "b", nil))
require.NoError(t, NewCustomerDeviceProcessor(repo, fcm).Register(ctx, repository.CustomerDevice{CustomerID: customer, DeviceID: "phone", FCMToken: "t1"}))
require.NoError(t, NewCustomerDeviceProcessor(repo, fcm).Unregister(ctx, customer, "phone"))
require.NoError(t, NewCustomerDeviceProcessor(repo, fcm).Notify(ctx, customer, "t", "b", nil))
assert.Empty(t, fcm.tokens)
// FCM not configured.
require.NoError(t, NewCustomerDeviceProcessor(repo, nil).Notify(ctx, customer, "t", "b", nil))
}
func TestCustomerDevice_FCMFailureIsReported(t *testing.T) {
repo, fcm := &customerDeviceRepoFake{}, &fcmFake{err: errors.New("unavailable")}
p := NewCustomerDeviceProcessor(repo, fcm)
ctx := context.Background()
customer := uuid.New()
require.NoError(t, p.Register(ctx, repository.CustomerDevice{CustomerID: customer, DeviceID: "phone", FCMToken: "t1"}))
assert.Error(t, p.Notify(ctx, customer, "t", "b", nil))
}
func TestCustomerDevice_RejectsIncompleteRegistrations(t *testing.T) {
p := NewCustomerDeviceProcessor(&customerDeviceRepoFake{}, nil)
ctx := context.Background()
for name, d := range map[string]repository.CustomerDevice{
"no device": {FCMToken: "t"},
"no token": {DeviceID: "phone", FCMToken: " "},
"bad platform": {DeviceID: "phone", FCMToken: "t", Platform: ptr("symbian")},
} {
assert.ErrorIs(t, p.Register(ctx, d), ErrInvalidCustomerDevice, name)
}
}
@@ -57,8 +57,8 @@ func (e *walletMoveEnv) exchanges() *WalletExchangeProcessor {
return p
}
func (e *walletMoveEnv) transfers(messenger walletMessenger) *WalletTransferProcessor {
p := NewWalletTransferProcessor(e.customers, e, e, e.pins, e.p, txRunnerFake{}, messenger)
func (e *walletMoveEnv) transfers(notifier customerNotifier) *WalletTransferProcessor {
p := NewWalletTransferProcessor(e.customers, e, e, e.pins, e.p, txRunnerFake{}, notifier)
p.now = func() time.Time { return e.now }
return p
}
+31 -16
View File
@@ -4,6 +4,7 @@ import (
"context"
"errors"
"fmt"
"strconv"
"strings"
"time"
"unicode/utf8"
@@ -21,12 +22,16 @@ import (
// check does not reveal who uses the app elsewhere.
var ErrWalletRecipientNotFound = errors.New("no customer of this organization has that phone number")
// walletMessenger tells a customer something happened to their wallet. There is no
// push channel to customers yet, so the app sends it by WhatsApp.
type walletMessenger interface {
SendWhatsAppMessage(phoneNumber, message string) error
// customerNotifier pushes a notification to a customer's app through FCM.
// CustomerDeviceProcessor is one.
type customerNotifier interface {
Notify(ctx context.Context, customerID uuid.UUID, title, body string, data map[string]string) error
}
// NotificationTypeWalletTransferIn is the data type of the push a transfer recipient
// gets, so the app can open the transaction.
const NotificationTypeWalletTransferIn = "WALLET_TRANSFER_IN"
// WalletTransferProcessor sends EnakPoint or EnakCoin from one customer to another in
// the same organization (docs/prd-point-coin.md F5).
type WalletTransferProcessor struct {
@@ -36,12 +41,12 @@ type WalletTransferProcessor struct {
pins pinVerifier
wallet *WalletProcessor
tx TxRunner
messenger walletMessenger
notifier customerNotifier
now func() time.Time
}
func NewWalletTransferProcessor(customers repository.WalletMoveRepository, settings organizationSettingsReader, spendable spendableReader, pins pinVerifier, wallet *WalletProcessor, tx TxRunner, messenger walletMessenger) *WalletTransferProcessor {
return &WalletTransferProcessor{customers: customers, settings: settings, spendable: spendable, pins: pins, wallet: wallet, tx: tx, messenger: messenger, now: time.Now}
func NewWalletTransferProcessor(customers repository.WalletMoveRepository, settings organizationSettingsReader, spendable spendableReader, pins pinVerifier, wallet *WalletProcessor, tx TxRunner, notifier customerNotifier) *WalletTransferProcessor {
return &WalletTransferProcessor{customers: customers, settings: settings, spendable: spendable, pins: pins, wallet: wallet, tx: tx, notifier: notifier, now: time.Now}
}
// Recipient is GET /customer/wallet/transfer/recipient: the masked name and number
@@ -118,6 +123,7 @@ func (p *WalletTransferProcessor) Transfer(ctx context.Context, senderID uuid.UU
outKey := fmt.Sprintf("transfer:%s:%s:out", senderID, key)
inKey := fmt.Sprintf("transfer:%s:%s:in", senderID, key)
result := &models.WalletTransferResult{Currency: currency, Amount: in.Amount, Recipient: *to}
var receivedID uuid.UUID
err = p.tx.WithTransaction(ctx, func(ctx context.Context) error {
if err := p.wallet.LockWallets(ctx, senderID, recipient.ID); err != nil {
return err
@@ -185,6 +191,7 @@ func (p *WalletTransferProcessor) Transfer(ctx context.Context, senderID uuid.UU
result.GroupID = groupID
result.Lots = movedLots(received.Lots)
result.Replayed = out.Replayed
receivedID = received.Transaction.ID
return nil
})
if err != nil {
@@ -192,7 +199,7 @@ func (p *WalletTransferProcessor) Transfer(ctx context.Context, senderID uuid.UU
}
if !result.Replayed {
p.tellRecipient(recipient, from, currency, in.Amount)
p.tellRecipient(ctx, recipient.ID, from, currency, in.Amount, receivedID, result.GroupID)
}
balances, err := p.spendable.SpendableBalances(ctx, senderID, p.now())
if err != nil {
@@ -228,16 +235,24 @@ func (p *WalletTransferProcessor) recipient(ctx context.Context, sender *reposit
return recipient, nil
}
// tellRecipient is best effort: the transfer has already happened, so a failure to
// send the message is only logged.
func (p *WalletTransferProcessor) tellRecipient(recipient *repository.WalletMoveCustomer, sender *models.WalletTransferRecipient, currency string, amount int64) {
if p.messenger == nil || recipient.PhoneNumber == nil {
// tellRecipient pushes the transfer to the recipient's app (F5). It is best effort:
// the transfer has already happened, so a failure to send is only logged.
func (p *WalletTransferProcessor) tellRecipient(ctx context.Context, recipientID uuid.UUID, sender *models.WalletTransferRecipient, currency string, amount int64, transactionID, groupID uuid.UUID) {
if p.notifier == nil {
return
}
message := fmt.Sprintf("Kamu menerima %d %s dari %s (%s). Cek riwayatnya di aplikasi.",
amount, walletCurrencyName(currency), sender.Name, sender.PhoneNumber)
if err := p.messenger.SendWhatsAppMessage(*recipient.PhoneNumber, message); err != nil {
logger.NonContext.Error(fmt.Sprintf("Could not tell customer %s about a transfer", recipient.ID), err)
name := walletCurrencyName(currency)
title := name + " masuk"
body := fmt.Sprintf("Kamu menerima %d %s dari %s (%s).", amount, name, sender.Name, sender.PhoneNumber)
data := map[string]string{
"type": NotificationTypeWalletTransferIn,
"transaction_id": transactionID.String(),
"group_id": groupID.String(),
"currency": currency,
"amount": strconv.FormatInt(amount, 10),
}
if err := p.notifier.Notify(ctx, recipientID, title, body, data); err != nil {
logger.NonContext.Error(fmt.Sprintf("Could not tell customer %s about a transfer", recipientID), err)
}
}
@@ -1,6 +1,7 @@
package processor
import (
"context"
"errors"
"testing"
"time"
@@ -14,13 +15,19 @@ import (
"apskel-pos-be/internal/repository"
)
type messengerFake struct{ sent map[string][]string }
type pushFake struct {
title, body string
data map[string]string
}
func (f *messengerFake) SendWhatsAppMessage(phone, message string) error {
if f.sent == nil {
f.sent = map[string][]string{}
// notifierFake records the pushes each customer would get.
type notifierFake struct{ pushes map[uuid.UUID][]pushFake }
func (f *notifierFake) Notify(_ context.Context, customerID uuid.UUID, title, body string, data map[string]string) error {
if f.pushes == nil {
f.pushes = map[uuid.UUID][]pushFake{}
}
f.sent[phone] = append(f.sent[phone], message)
f.pushes[customerID] = append(f.pushes[customerID], pushFake{title: title, body: body, data: data})
return nil
}
@@ -35,10 +42,10 @@ func TestWalletTransfer_MovesBalanceWithItsExpiry(t *testing.T) {
dec, jan := e.at(30*24*time.Hour), e.at(60*24*time.Hour)
first := e.credit(t, earn(a, 100, dec))
second := e.credit(t, earn(a, 50, jan))
messenger := &messengerFake{}
notifier := &notifierFake{}
// The example in §8: A sends 120, 100 from the lot expiring first and 20 from the next.
res, err := e.transfers(messenger).Transfer(e.ctx, a, sendPoints(120, "081234561234"), "482913", "key-1", models.CustomerPinRequestInfo{})
res, err := e.transfers(notifier).Transfer(e.ctx, a, sendPoints(120, "081234561234"), "482913", "key-1", models.CustomerPinRequestInfo{})
require.NoError(t, err)
assert.Equal(t, int64(30), res.Balance)
@@ -74,7 +81,18 @@ func TestWalletTransfer_MovesBalanceWithItsExpiry(t *testing.T) {
assert.GreaterOrEqual(t, e.repo.locks[a], 1)
assert.GreaterOrEqual(t, e.repo.locks[b], 1)
assert.Equal(t, []string{"Kamu menerima 120 EnakPoint dari An*** (08**-****-5678). Cek riwayatnya di aplikasi."}, messenger.sent["081234561234"])
assert.Equal(t, []pushFake{{
title: "EnakPoint masuk",
body: "Kamu menerima 120 EnakPoint dari An*** (08**-****-5678).",
data: map[string]string{
"type": NotificationTypeWalletTransferIn,
"transaction_id": in.ID.String(),
"group_id": in.GroupID.String(),
"currency": constants.WalletCurrencyPoint,
"amount": "120",
},
}}, notifier.pushes[b])
assert.Empty(t, notifier.pushes[a], "the sender gets no push")
}
func TestWalletTransfer_Coins(t *testing.T) {
@@ -191,20 +209,20 @@ func TestWalletTransfer_RetryMovesNothingAndTellsNobodyAgain(t *testing.T) {
b := e.member("Budi", "081234561234")
e.member("Citra", "081255550000")
e.credit(t, earn(a, 100, nil))
messenger := &messengerFake{}
notifier := &notifierFake{}
first, err := e.transfers(messenger).Transfer(e.ctx, a, sendPoints(40, "081234561234"), "482913", "key-1", models.CustomerPinRequestInfo{})
first, err := e.transfers(notifier).Transfer(e.ctx, a, sendPoints(40, "081234561234"), "482913", "key-1", models.CustomerPinRequestInfo{})
require.NoError(t, err)
again, err := e.transfers(messenger).Transfer(e.ctx, a, sendPoints(40, "081234561234"), "482913", "key-1", models.CustomerPinRequestInfo{})
again, err := e.transfers(notifier).Transfer(e.ctx, a, sendPoints(40, "081234561234"), "482913", "key-1", models.CustomerPinRequestInfo{})
require.NoError(t, err)
assert.True(t, again.Replayed)
assert.Equal(t, first.GroupID, again.GroupID)
assert.Equal(t, int64(40), e.balance(t, b))
assert.Len(t, messenger.sent["081234561234"], 1)
assert.Len(t, notifier.pushes[b], 1)
// The same key to someone else is not a retry.
_, err = e.transfers(messenger).Transfer(e.ctx, a, sendPoints(40, "081255550000"), "482913", "key-1", models.CustomerPinRequestInfo{})
_, err = e.transfers(notifier).Transfer(e.ctx, a, sendPoints(40, "081255550000"), "482913", "key-1", models.CustomerPinRequestInfo{})
assert.ErrorIs(t, err, ErrWalletIdempotencyConflict)
}