185 lines
5.9 KiB
Go
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
|
||
|
|
}
|