Files
apskel-pos-backend/internal/processor/wallet_migration_processor.go
T
2026-09-30 15:31:11 +07:00

185 lines
5.9 KiB
Go

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
}