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>
135 lines
4.8 KiB
Go
135 lines
4.8 KiB
Go
package repository
|
|
|
|
import (
|
|
"context"
|
|
"fmt"
|
|
|
|
"github.com/google/uuid"
|
|
"gorm.io/gorm"
|
|
|
|
"apskel-pos-be/internal/constants"
|
|
"apskel-pos-be/internal/entities"
|
|
)
|
|
|
|
// LegacyBalance is what one customer holds in customer_points and customer_tokens,
|
|
// the tables the wallet replaces (docs/prd-point-coin.md §10).
|
|
type LegacyBalance struct {
|
|
CustomerID uuid.UUID
|
|
// Nil when the customer has no customer_points row.
|
|
PointsRowID *uuid.UUID
|
|
Points int64
|
|
Tokens []entities.CustomerTokens
|
|
}
|
|
|
|
// Coins is the sum of every token type: all of them become EnakCoin (Q6).
|
|
func (b LegacyBalance) Coins() int64 {
|
|
var total int64
|
|
for _, t := range b.Tokens {
|
|
total += t.Balance
|
|
}
|
|
return total
|
|
}
|
|
|
|
// WalletMigrationTotals compares the legacy tables with what has been migrated.
|
|
type WalletMigrationTotals struct {
|
|
LegacyPoints int64
|
|
LegacyCoins int64
|
|
MigratedPoints int64
|
|
MigratedCoins int64
|
|
WalletPoints int64
|
|
WalletCoins int64
|
|
}
|
|
|
|
// WalletMigrationRepository reads the legacy balances for the one-time move into the
|
|
// wallet. The writes go through the wallet processor like any other credit.
|
|
type WalletMigrationRepository interface {
|
|
// ListLegacyCustomers returns, in id order, up to limit customers after the given
|
|
// id that have a row in customer_points or customer_tokens.
|
|
ListLegacyCustomers(ctx context.Context, after uuid.UUID, limit int) ([]uuid.UUID, error)
|
|
GetLegacyBalance(ctx context.Context, customerID uuid.UUID) (*LegacyBalance, error)
|
|
// SumMigrated returns how much has already been credited to the customer by
|
|
// MIGRATION ledger rows in the currency.
|
|
SumMigrated(ctx context.Context, customerID uuid.UUID, currency string) (int64, error)
|
|
Totals(ctx context.Context) (*WalletMigrationTotals, error)
|
|
}
|
|
|
|
type walletMigrationRepository struct {
|
|
db *gorm.DB
|
|
}
|
|
|
|
func NewWalletMigrationRepository(db *gorm.DB) WalletMigrationRepository {
|
|
return &walletMigrationRepository{db: db}
|
|
}
|
|
|
|
func (r *walletMigrationRepository) ListLegacyCustomers(ctx context.Context, after uuid.UUID, limit int) ([]uuid.UUID, error) {
|
|
var ids []uuid.UUID
|
|
err := DBFromContext(ctx, r.db).WithContext(ctx).Raw(`
|
|
SELECT customer_id FROM (
|
|
SELECT customer_id FROM customer_points
|
|
UNION
|
|
SELECT customer_id FROM customer_tokens
|
|
) legacy
|
|
WHERE customer_id > ?
|
|
ORDER BY customer_id
|
|
LIMIT ?`, after, limit).
|
|
Scan(&ids).Error
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to list legacy customers: %w", err)
|
|
}
|
|
return ids, nil
|
|
}
|
|
|
|
func (r *walletMigrationRepository) GetLegacyBalance(ctx context.Context, customerID uuid.UUID) (*LegacyBalance, error) {
|
|
db := DBFromContext(ctx, r.db).WithContext(ctx)
|
|
balance := &LegacyBalance{CustomerID: customerID}
|
|
|
|
// Find rather than First: many customers have tokens but no points row, and First
|
|
// would log each of them as a "record not found" error.
|
|
var points []entities.CustomerPoints
|
|
if err := db.Where("customer_id = ?", customerID).Limit(1).Find(&points).Error; err != nil {
|
|
return nil, fmt.Errorf("failed to get legacy points: %w", err)
|
|
}
|
|
if len(points) > 0 {
|
|
balance.PointsRowID = &points[0].ID
|
|
balance.Points = points[0].Balance
|
|
}
|
|
|
|
err := db.Where("customer_id = ?", customerID).Order("token_type").Find(&balance.Tokens).Error
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to get legacy tokens: %w", err)
|
|
}
|
|
return balance, nil
|
|
}
|
|
|
|
func (r *walletMigrationRepository) SumMigrated(ctx context.Context, customerID uuid.UUID, currency string) (int64, error) {
|
|
var total int64
|
|
err := DBFromContext(ctx, r.db).WithContext(ctx).
|
|
Model(&entities.WalletTransaction{}).
|
|
Where("customer_id = ? AND currency = ? AND type = ?", customerID, currency, constants.WalletTxTypeMigration).
|
|
Select("COALESCE(SUM(amount), 0)").
|
|
Scan(&total).Error
|
|
if err != nil {
|
|
return 0, fmt.Errorf("failed to sum migrated balance: %w", err)
|
|
}
|
|
return total, nil
|
|
}
|
|
|
|
func (r *walletMigrationRepository) Totals(ctx context.Context) (*WalletMigrationTotals, error) {
|
|
var totals WalletMigrationTotals
|
|
err := DBFromContext(ctx, r.db).WithContext(ctx).Raw(`
|
|
SELECT
|
|
(SELECT COALESCE(SUM(balance), 0) FROM customer_points) AS legacy_points,
|
|
(SELECT COALESCE(SUM(balance), 0) FROM customer_tokens) AS legacy_coins,
|
|
(SELECT COALESCE(SUM(amount), 0) FROM wallet_transactions WHERE type = ? AND currency = ?) AS migrated_points,
|
|
(SELECT COALESCE(SUM(amount), 0) FROM wallet_transactions WHERE type = ? AND currency = ?) AS migrated_coins,
|
|
(SELECT COALESCE(SUM(point_balance), 0) FROM customer_wallets) AS wallet_points,
|
|
(SELECT COALESCE(SUM(coin_balance), 0) FROM customer_wallets) AS wallet_coins`,
|
|
constants.WalletTxTypeMigration, constants.WalletCurrencyPoint,
|
|
constants.WalletTxTypeMigration, constants.WalletCurrencyCoin).
|
|
Scan(&totals).Error
|
|
if err != nil {
|
|
return nil, fmt.Errorf("failed to compute wallet migration totals: %w", err)
|
|
}
|
|
return &totals, nil
|
|
}
|