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 }