feat(wallet): migrate legacy points and tokens into the wallet
Adds cmd/wallet-migrate (make wallet-migrate, args=-dry-run to only report), which moves customer_points and customer_tokens into the wallet (docs/prd-point-coin.md §10, PC-105). Each customer gets a MIGRATION ledger row and a non-expiring lot per currency, written through WalletProcessor in one transaction per customer. EnakCoin is the sum of every token type (Q6), with the legacy rows listed in the row's metadata. It credits the difference between the legacy balance and what earlier runs migrated, so running it again never doubles a balance and picks up only what the old code added since. A legacy balance that shrank after being migrated is reported and left alone, since only an admin adjustment may take balance away, and the command then exits non-zero. It ends with a legacy / migrated / wallet total per currency. Migration 000092 renames TOKENS to COINS in campaigns.type and campaign_rules.reward_type. The campaign API now validates COINS; it still accepts TOKENS, including as a list filter, and stores it as COINS so older dashboards keep working while they are updated. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 5.5
parent
6af97f5696
commit
41b75810fd
@@ -70,7 +70,7 @@ func (p *campaignProcessor) ListCampaigns(ctx context.Context, req *contract.Lis
|
||||
Page: req.Page,
|
||||
Limit: req.Limit,
|
||||
Search: req.Search,
|
||||
Type: req.Type,
|
||||
Type: string(entities.NormalizeCampaignType(req.Type)),
|
||||
IsActive: req.IsActive,
|
||||
ShowOnApp: req.ShowOnApp,
|
||||
StartDate: req.StartDate,
|
||||
@@ -178,7 +178,7 @@ func (p *campaignRuleProcessor) CreateCampaignRule(ctx context.Context, req *con
|
||||
CampaignID: req.CampaignID,
|
||||
RuleType: entities.RuleType(req.RuleType),
|
||||
ConditionValue: req.ConditionValue,
|
||||
RewardType: entities.CampaignRewardType(req.RewardType),
|
||||
RewardType: entities.NormalizeCampaignRewardType(req.RewardType),
|
||||
RewardValue: req.RewardValue,
|
||||
RewardSubtype: (*entities.RewardSubtype)(req.RewardSubtype),
|
||||
RewardRefID: req.RewardRefID,
|
||||
@@ -218,7 +218,7 @@ func (p *campaignRuleProcessor) ListCampaignRules(ctx context.Context, req *cont
|
||||
Limit: req.Limit,
|
||||
CampaignID: req.CampaignID,
|
||||
RuleType: req.RuleType,
|
||||
RewardType: req.RewardType,
|
||||
RewardType: string(entities.NormalizeCampaignRewardType(req.RewardType)),
|
||||
}
|
||||
|
||||
// Get from repository
|
||||
@@ -247,7 +247,7 @@ func (p *campaignRuleProcessor) UpdateCampaignRule(ctx context.Context, req *con
|
||||
CampaignID: req.CampaignID,
|
||||
RuleType: entities.RuleType(req.RuleType),
|
||||
ConditionValue: req.ConditionValue,
|
||||
RewardType: entities.CampaignRewardType(req.RewardType),
|
||||
RewardType: entities.NormalizeCampaignRewardType(req.RewardType),
|
||||
RewardValue: req.RewardValue,
|
||||
RewardSubtype: (*entities.RewardSubtype)(req.RewardSubtype),
|
||||
RewardRefID: req.RewardRefID,
|
||||
|
||||
@@ -0,0 +1,184 @@
|
||||
package processor
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/google/uuid"
|
||||
|
||||
"apskel-pos-be/internal/constants"
|
||||
"apskel-pos-be/internal/entities"
|
||||
"apskel-pos-be/internal/repository"
|
||||
)
|
||||
|
||||
// TxRunner runs fn inside a database transaction. repository.TxManager is one.
|
||||
type TxRunner interface {
|
||||
WithTransaction(ctx context.Context, fn func(ctx context.Context) error) error
|
||||
}
|
||||
|
||||
// WalletMigrationDiscrepancy is a customer whose legacy balance is now lower than
|
||||
// what was already migrated: the old code spent from it after the migration ran.
|
||||
// The wallet is left alone, because only an admin adjustment can take balance away.
|
||||
type WalletMigrationDiscrepancy struct {
|
||||
CustomerID uuid.UUID
|
||||
Currency string
|
||||
Legacy int64
|
||||
Migrated int64
|
||||
}
|
||||
|
||||
type WalletMigrationReport struct {
|
||||
DryRun bool
|
||||
CustomersScanned int
|
||||
// Ledger rows written (or, on a dry run, that would be written) and their sum.
|
||||
PointCredits int
|
||||
PointsCredited int64
|
||||
CoinCredits int
|
||||
CoinsCredited int64
|
||||
Discrepancies []WalletMigrationDiscrepancy
|
||||
// Taken after the run. On a dry run they show the state before it.
|
||||
Totals *repository.WalletMigrationTotals
|
||||
}
|
||||
|
||||
// Balanced reports whether everything in the legacy tables is now in the wallet.
|
||||
func (r *WalletMigrationReport) Balanced() bool {
|
||||
return len(r.Discrepancies) == 0 && r.Totals != nil &&
|
||||
r.Totals.LegacyPoints == r.Totals.MigratedPoints &&
|
||||
r.Totals.LegacyCoins == r.Totals.MigratedCoins
|
||||
}
|
||||
|
||||
// WalletMigrationProcessor moves the balances in customer_points and customer_tokens
|
||||
// into the wallet (docs/prd-point-coin.md §10, PC-105). Each customer gets a MIGRATION
|
||||
// ledger row and a non-expiring lot per currency, through WalletProcessor like any
|
||||
// other credit, so the wallet reconciles from the first row.
|
||||
//
|
||||
// It credits the difference between the legacy balance and what earlier runs already
|
||||
// migrated, so running it again never doubles a balance, and a run after the old code
|
||||
// kept writing to the legacy tables picks up only what was added since.
|
||||
type WalletMigrationProcessor struct {
|
||||
repo repository.WalletMigrationRepository
|
||||
wallet *WalletProcessor
|
||||
tx TxRunner
|
||||
}
|
||||
|
||||
func NewWalletMigrationProcessor(repo repository.WalletMigrationRepository, wallet *WalletProcessor, tx TxRunner) *WalletMigrationProcessor {
|
||||
return &WalletMigrationProcessor{repo: repo, wallet: wallet, tx: tx}
|
||||
}
|
||||
|
||||
// Run migrates every customer with a legacy balance, one transaction per customer.
|
||||
// With dryRun it only reports what it would credit.
|
||||
func (p *WalletMigrationProcessor) Run(ctx context.Context, dryRun bool, batchSize int) (*WalletMigrationReport, error) {
|
||||
if batchSize <= 0 {
|
||||
batchSize = 500
|
||||
}
|
||||
report := &WalletMigrationReport{DryRun: dryRun}
|
||||
|
||||
after := uuid.Nil
|
||||
for {
|
||||
ids, err := p.repo.ListLegacyCustomers(ctx, after, batchSize)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if len(ids) == 0 {
|
||||
break
|
||||
}
|
||||
for _, id := range ids {
|
||||
if dryRun {
|
||||
err = p.migrateCustomer(ctx, id, true, report)
|
||||
} else {
|
||||
err = p.tx.WithTransaction(ctx, func(ctx context.Context) error {
|
||||
return p.migrateCustomer(ctx, id, false, report)
|
||||
})
|
||||
}
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("customer %s: %w", id, err)
|
||||
}
|
||||
report.CustomersScanned++
|
||||
}
|
||||
after = ids[len(ids)-1]
|
||||
}
|
||||
|
||||
totals, err := p.repo.Totals(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
report.Totals = totals
|
||||
return report, nil
|
||||
}
|
||||
|
||||
func (p *WalletMigrationProcessor) migrateCustomer(ctx context.Context, customerID uuid.UUID, dryRun bool, report *WalletMigrationReport) error {
|
||||
// Lock before reading what was migrated, so two runs at once cannot both see the
|
||||
// same gap and fill it twice.
|
||||
if !dryRun {
|
||||
if err := p.wallet.LockWallet(ctx, customerID); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
legacy, err := p.repo.GetLegacyBalance(ctx, customerID)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Points come from the single customer_points row. Tokens come from several rows,
|
||||
// one per type, so the ledger row points at the customer and lists the rows.
|
||||
pointsRef := customerID
|
||||
if legacy.PointsRowID != nil {
|
||||
pointsRef = *legacy.PointsRowID
|
||||
}
|
||||
tokens := make([]map[string]any, 0, len(legacy.Tokens))
|
||||
for _, t := range legacy.Tokens {
|
||||
tokens = append(tokens, map[string]any{"id": t.ID, "token_type": string(t.TokenType), "balance": t.Balance})
|
||||
}
|
||||
|
||||
for _, c := range []struct {
|
||||
currency, refType string
|
||||
refID uuid.UUID
|
||||
legacy int64
|
||||
metadata entities.Metadata
|
||||
credits *int
|
||||
credited *int64
|
||||
}{
|
||||
{constants.WalletCurrencyPoint, constants.WalletRefTypeLegacyPoints, pointsRef, legacy.Points,
|
||||
entities.Metadata{}, &report.PointCredits, &report.PointsCredited},
|
||||
{constants.WalletCurrencyCoin, constants.WalletRefTypeLegacyTokens, customerID, legacy.Coins(),
|
||||
entities.Metadata{"legacy_tokens": tokens}, &report.CoinCredits, &report.CoinsCredited},
|
||||
} {
|
||||
migrated, err := p.repo.SumMigrated(ctx, customerID, c.currency)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
delta := c.legacy - migrated
|
||||
if delta < 0 {
|
||||
report.Discrepancies = append(report.Discrepancies, WalletMigrationDiscrepancy{
|
||||
CustomerID: customerID, Currency: c.currency, Legacy: c.legacy, Migrated: migrated,
|
||||
})
|
||||
continue
|
||||
}
|
||||
if delta == 0 {
|
||||
continue
|
||||
}
|
||||
|
||||
if !dryRun {
|
||||
c.metadata["legacy_balance"] = c.legacy
|
||||
c.metadata["previously_migrated"] = migrated
|
||||
_, err = p.wallet.Credit(ctx, WalletCreditInput{WalletEntry: WalletEntry{
|
||||
CustomerID: customerID,
|
||||
Currency: c.currency,
|
||||
Type: constants.WalletTxTypeMigration,
|
||||
Amount: delta,
|
||||
ReferenceType: c.refType,
|
||||
ReferenceID: c.refID,
|
||||
Description: "Saldo awal dari sistem lama",
|
||||
Metadata: c.metadata,
|
||||
// The legacy total in the key lets a later run top up a balance that
|
||||
// grew, while a retry of the same run is still recognised.
|
||||
IdempotencyKey: fmt.Sprintf("migration:%s:%s:%d", c.currency, customerID, c.legacy),
|
||||
}})
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
*c.credits++
|
||||
*c.credited += delta
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,167 @@
|
||||
package processor
|
||||
|
||||
import (
|
||||
"context"
|
||||
"os"
|
||||
"testing"
|
||||
|
||||
"github.com/google/uuid"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/driver/postgres"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/gorm/logger"
|
||||
|
||||
"apskel-pos-be/internal/constants"
|
||||
"apskel-pos-be/internal/repository"
|
||||
)
|
||||
|
||||
// Needs TEST_DATABASE_URL pointing at a migrated database; see
|
||||
// internal/repository/wallet_repository_test.go. Other packages' tests may use the
|
||||
// same database at the same time, so everything here is scoped to its own customers.
|
||||
func TestWalletMigrationProcessor_AgainstPostgres(t *testing.T) {
|
||||
dsn := os.Getenv("TEST_DATABASE_URL")
|
||||
if dsn == "" {
|
||||
t.Skip("TEST_DATABASE_URL not set")
|
||||
}
|
||||
db, err := gorm.Open(postgres.Open(dsn), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)})
|
||||
require.NoError(t, err)
|
||||
ctx := context.Background()
|
||||
|
||||
org := uuid.New()
|
||||
full, tokensOnly, pointsOnly, none := uuid.New(), uuid.New(), uuid.New(), uuid.New()
|
||||
customers := []uuid.UUID{full, tokensOnly, pointsOnly, none}
|
||||
exec := func(q string, args ...any) {
|
||||
t.Helper()
|
||||
require.NoError(t, db.Exec(q, args...).Error)
|
||||
}
|
||||
exec(`INSERT INTO organizations (id, name, plan_type) VALUES (?, 'migration test', 'basic')`, org)
|
||||
for _, c := range customers {
|
||||
exec(`INSERT INTO customers (id, organization_id, name) VALUES (?, ?, 'migration test')`, c, org)
|
||||
}
|
||||
t.Cleanup(func() {
|
||||
db.Exec(`DELETE FROM wallet_lot_allocations WHERE lot_id IN (SELECT id FROM wallet_lots WHERE customer_id IN ?)`, customers)
|
||||
db.Exec(`DELETE FROM wallet_lots WHERE customer_id IN ?`, customers)
|
||||
db.Exec(`DELETE FROM wallet_transactions WHERE customer_id IN ?`, customers)
|
||||
db.Exec(`DELETE FROM customer_wallets WHERE customer_id IN ?`, customers)
|
||||
db.Exec(`DELETE FROM customers WHERE id IN ?`, customers)
|
||||
db.Exec(`DELETE FROM organizations WHERE id = ?`, org)
|
||||
})
|
||||
|
||||
// The example from §10: SPIN 5 + RAFFLE 2 + MINIGAME 1 = 8 EnakCoin.
|
||||
exec(`INSERT INTO customer_points (customer_id, balance) VALUES (?, 100), (?, 0), (?, 40)`, full, tokensOnly, pointsOnly)
|
||||
exec(`INSERT INTO customer_tokens (customer_id, token_type, balance) VALUES
|
||||
(?, 'SPIN', 5), (?, 'RAFFLE', 2), (?, 'MINIGAME', 1), (?, 'SPIN', 3)`, full, full, full, tokensOnly)
|
||||
|
||||
migrator := NewWalletMigrationProcessor(
|
||||
repository.NewWalletMigrationRepository(db),
|
||||
NewWalletProcessor(repository.NewWalletRepository(db)),
|
||||
repository.NewTxManager(db),
|
||||
)
|
||||
|
||||
type balance struct{ Point, Coin int64 }
|
||||
balances := func() map[uuid.UUID]balance {
|
||||
t.Helper()
|
||||
var rows []struct {
|
||||
CustomerID uuid.UUID
|
||||
PointBalance, CoinBalance int64
|
||||
}
|
||||
require.NoError(t, db.Raw(`SELECT customer_id, point_balance, coin_balance FROM customer_wallets WHERE customer_id IN ?`, customers).Scan(&rows).Error)
|
||||
out := map[uuid.UUID]balance{}
|
||||
for _, r := range rows {
|
||||
out[r.CustomerID] = balance{r.PointBalance, r.CoinBalance}
|
||||
}
|
||||
return out
|
||||
}
|
||||
countRows := func() int64 {
|
||||
t.Helper()
|
||||
var n int64
|
||||
require.NoError(t, db.Raw(`SELECT COUNT(*) FROM wallet_transactions WHERE customer_id IN ?`, customers).Scan(&n).Error)
|
||||
return n
|
||||
}
|
||||
|
||||
// A dry run reports and writes nothing, not even the wallets.
|
||||
report, err := migrator.Run(ctx, true, 2)
|
||||
require.NoError(t, err)
|
||||
assert.GreaterOrEqual(t, report.PointsCredited, int64(140))
|
||||
assert.GreaterOrEqual(t, report.CoinsCredited, int64(11))
|
||||
assert.Empty(t, balances())
|
||||
assert.Zero(t, countRows())
|
||||
|
||||
// The real run. A batch of 2 makes it page through the customers.
|
||||
_, err = migrator.Run(ctx, false, 2)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, map[uuid.UUID]balance{
|
||||
full: {Point: 100, Coin: 8},
|
||||
tokensOnly: {Point: 0, Coin: 3},
|
||||
pointsOnly: {Point: 40, Coin: 0},
|
||||
}, balances(), "a customer without legacy rows gets no wallet")
|
||||
assert.Equal(t, int64(4), countRows(), "one row per customer per currency with a balance")
|
||||
|
||||
var coinRow struct {
|
||||
ReferenceType string
|
||||
ReferenceID uuid.UUID
|
||||
Metadata string
|
||||
}
|
||||
require.NoError(t, db.Raw(`SELECT reference_type, reference_id, metadata::text AS metadata FROM wallet_transactions
|
||||
WHERE customer_id = ? AND currency = 'COIN'`, full).Scan(&coinRow).Error)
|
||||
assert.Equal(t, constants.WalletRefTypeLegacyTokens, coinRow.ReferenceType)
|
||||
assert.Equal(t, full, coinRow.ReferenceID)
|
||||
for _, part := range []string{`"token_type": "SPIN"`, `"token_type": "RAFFLE"`, `"token_type": "MINIGAME"`, `"legacy_balance": 8`} {
|
||||
assert.Contains(t, coinRow.Metadata, part)
|
||||
}
|
||||
|
||||
var pointRef, pointsRowID string
|
||||
require.NoError(t, db.Raw(`SELECT reference_id::text FROM wallet_transactions WHERE customer_id = ? AND currency = 'POINT'`, full).Scan(&pointRef).Error)
|
||||
require.NoError(t, db.Raw(`SELECT id::text FROM customer_points WHERE customer_id = ?`, full).Scan(&pointsRowID).Error)
|
||||
assert.NotEmpty(t, pointRef)
|
||||
assert.Equal(t, pointsRowID, pointRef, "points row points at the customer_points row")
|
||||
|
||||
var expiring int64
|
||||
require.NoError(t, db.Raw(`SELECT COUNT(*) FROM wallet_lots WHERE customer_id IN ? AND expires_at IS NOT NULL`, customers).Scan(&expiring).Error)
|
||||
assert.Zero(t, expiring, "migrated lots never expire")
|
||||
|
||||
// Running again changes nothing.
|
||||
report, err = migrator.Run(ctx, false, 2)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, int64(4), countRows())
|
||||
assert.Empty(t, discrepanciesFor(report, customers))
|
||||
|
||||
// The old code kept writing: one balance grew, one shrank. Only the growth is
|
||||
// migrated; the shrink is reported and left alone.
|
||||
exec(`UPDATE customer_tokens SET balance = 9 WHERE customer_id = ? AND token_type = 'SPIN'`, full)
|
||||
exec(`UPDATE customer_points SET balance = 30 WHERE customer_id = ?`, pointsOnly)
|
||||
report, err = migrator.Run(ctx, false, 2)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, balance{Point: 100, Coin: 12}, balances()[full])
|
||||
assert.Equal(t, balance{Point: 40, Coin: 0}, balances()[pointsOnly])
|
||||
assert.Equal(t, []WalletMigrationDiscrepancy{{CustomerID: pointsOnly, Currency: constants.WalletCurrencyPoint, Legacy: 30, Migrated: 40}},
|
||||
discrepanciesFor(report, customers))
|
||||
assert.Equal(t, int64(5), countRows())
|
||||
|
||||
// §7.5 for these customers.
|
||||
var broken int64
|
||||
require.NoError(t, db.Raw(`
|
||||
SELECT COUNT(*) FROM customer_wallets w
|
||||
WHERE w.customer_id IN ? AND (
|
||||
w.point_balance <> (SELECT COALESCE(SUM(amount), 0) FROM wallet_transactions t WHERE t.customer_id = w.customer_id AND t.currency = 'POINT')
|
||||
OR w.coin_balance <> (SELECT COALESCE(SUM(amount), 0) FROM wallet_transactions t WHERE t.customer_id = w.customer_id AND t.currency = 'COIN')
|
||||
OR w.point_balance <> (SELECT COALESCE(SUM(remaining_amount), 0) FROM wallet_lots l WHERE l.customer_id = w.customer_id AND l.currency = 'POINT')
|
||||
OR w.coin_balance <> (SELECT COALESCE(SUM(remaining_amount), 0) FROM wallet_lots l WHERE l.customer_id = w.customer_id AND l.currency = 'COIN'))`,
|
||||
customers).Scan(&broken).Error)
|
||||
assert.Zero(t, broken)
|
||||
}
|
||||
|
||||
func discrepanciesFor(report *WalletMigrationReport, customers []uuid.UUID) []WalletMigrationDiscrepancy {
|
||||
mine := map[uuid.UUID]bool{}
|
||||
for _, c := range customers {
|
||||
mine[c] = true
|
||||
}
|
||||
var out []WalletMigrationDiscrepancy
|
||||
for _, d := range report.Discrepancies {
|
||||
if mine[d.CustomerID] {
|
||||
out = append(out, d)
|
||||
}
|
||||
}
|
||||
return out
|
||||
}
|
||||
Reference in New Issue
Block a user