diff --git a/internal/app/app.go b/internal/app/app.go index ea9cf0c..11093f0 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -32,6 +32,7 @@ type App struct { shutdown chan os.Signal omsetScheduler *service.OmsetMilestoneScheduler walletRecon *service.WalletReconciliationJob + earningRetry *service.EarningBackfillJob } func NewApp(db *gorm.DB, redisClient *redis.Client) *App { @@ -60,6 +61,8 @@ func (a *App) Initialize(cfg *config.Config) error { repos.userRepo, processors.notificationProcessor, ) + // Earns for paid orders whose earning failed at payment time (docs/prd-point-coin.md F3) + a.earningRetry = service.NewEarningBackfillJob(processors.earningProcessor) services := a.initServices(processors, repos, cfg) validators := a.initValidators() @@ -167,6 +170,9 @@ func (a *App) Start(port string) error { if a.walletRecon != nil { a.walletRecon.Start(6 * time.Hour) } + if a.earningRetry != nil { + a.earningRetry.Start(30 * time.Minute) + } engine := a.router.Init() @@ -209,6 +215,9 @@ func (a *App) Shutdown() { if a.walletRecon != nil { a.walletRecon.Stop() } + if a.earningRetry != nil { + a.earningRetry.Stop() + } close(a.shutdown) } @@ -376,6 +385,7 @@ type processors struct { walletProcessor *processor.WalletProcessor walletAdminProcessor *processor.WalletAdminProcessor loyaltySettingsProcessor *processor.LoyaltySettingsProcessor + earningProcessor *processor.EarningProcessor } func (a *App) initProcessors(cfg *config.Config, repos *repositories) *processors { @@ -384,6 +394,12 @@ func (a *App) initProcessors(cfg *config.Config, repos *repositories) *processor otpProcessor := processor.NewOtpProcessor(fonnteClient, repos.otpRepo) inventoryMovementService := service.NewInventoryMovementService(repos.inventoryMovementRepo, repos.ingredientRepo) + orderProcessor := processor.NewOrderProcessorImpl(repos.orderRepo, repos.orderItemRepo, repos.paymentRepo, repos.paymentOrderItemRepo, repos.productRepo, repos.paymentMethodRepo, repos.inventoryRepo, repos.inventoryMovementRepo, repos.productVariantRepo, repos.outletRepo, repos.customerRepo, repos.txManager, repos.productRecipeRepo, repos.ingredientRepo, inventoryMovementService, repos.productOutletPriceRepo) + loyaltySettingsProcessor := processor.NewLoyaltySettingsProcessor(repos.loyaltySettingsRepo, repos.txManager) + // Earn EnakPoint and EnakCoin when an order becomes fully paid (docs/prd-point-coin.md F3) + earningProcessor := processor.NewEarningProcessor(repository.NewEarningRepository(a.db), loyaltySettingsProcessor, processor.NewWalletProcessor(repos.walletRepo), repos.txManager) + orderProcessor.SetOrderPaidHook(earningProcessor) + return &processors{ userProcessor: processor.NewUserProcessor(repos.userRepo, repos.organizationRepo, repos.outletRepo), organizationProcessor: processor.NewOrganizationProcessorImpl(repos.organizationRepo, repos.outletRepo, repos.userRepo), @@ -393,7 +409,7 @@ func (a *App) initProcessors(cfg *config.Config, repos *repositories) *processor productProcessor: processor.NewProductProcessorImpl(repos.productRepo, repos.categoryRepo, repos.productVariantRepo, repos.inventoryRepo, repos.outletRepo, repos.productOutletPriceRepo), productVariantProcessor: processor.NewProductVariantProcessorImpl(repos.productVariantRepo, repos.productRepo), inventoryProcessor: processor.NewInventoryProcessorImpl(repos.inventoryRepo, repos.productRepo, repos.outletRepo, repos.ingredientRepo, repos.inventoryMovementRepo), - orderProcessor: processor.NewOrderProcessorImpl(repos.orderRepo, repos.orderItemRepo, repos.paymentRepo, repos.paymentOrderItemRepo, repos.productRepo, repos.paymentMethodRepo, repos.inventoryRepo, repos.inventoryMovementRepo, repos.productVariantRepo, repos.outletRepo, repos.customerRepo, repos.txManager, repos.productRecipeRepo, repos.ingredientRepo, inventoryMovementService, repos.productOutletPriceRepo), + orderProcessor: orderProcessor, paymentMethodProcessor: processor.NewPaymentMethodProcessorImpl(repos.paymentMethodRepo), fileProcessor: processor.NewFileProcessorImpl(repos.fileRepo, fileClient), customerProcessor: processor.NewCustomerProcessor(repos.customerRepo), @@ -430,7 +446,8 @@ func (a *App) initProcessors(cfg *config.Config, repos *repositories) *processor expenseProcessor: processor.NewExpenseProcessorImpl(repos.expenseRepo, repos.purchaseCategoryRepo, repos.cashAdvanceRepo), cashAdvanceProcessor: processor.NewCashAdvanceProcessorImpl(repos.cashAdvanceRepo, repos.categoryRepo), walletProcessor: processor.NewWalletProcessor(repos.walletRepo), - loyaltySettingsProcessor: processor.NewLoyaltySettingsProcessor(repos.loyaltySettingsRepo, repos.txManager), + loyaltySettingsProcessor: loyaltySettingsProcessor, + earningProcessor: earningProcessor, walletAdminProcessor: processor.NewWalletAdminProcessor(repository.NewWalletAdminRepository(a.db), repos.walletQueryRepo, processor.NewWalletProcessor(repos.walletRepo), repos.txManager), } } diff --git a/internal/constants/payment.go b/internal/constants/payment.go index a95330e..75138db 100644 --- a/internal/constants/payment.go +++ b/internal/constants/payment.go @@ -8,6 +8,9 @@ const ( PaymentMethodTypeDigitalWallet PaymentMethodType = "digital_wallet" PaymentMethodTypeQR PaymentMethodType = "qr" PaymentMethodTypeEDC PaymentMethodType = "edc" + // Paying with EnakPoint (docs/prd-point-coin.md F9). Not accepted as a payment method + // type until that phase ships. + PaymentMethodTypePoint PaymentMethodType = "point" ) type PaymentStatus string diff --git a/internal/processor/earning_processor.go b/internal/processor/earning_processor.go new file mode 100644 index 0000000..02eb649 --- /dev/null +++ b/internal/processor/earning_processor.go @@ -0,0 +1,184 @@ +package processor + +import ( + "context" + "fmt" + "time" + + "github.com/google/uuid" + + "apskel-pos-be/internal/constants" + "apskel-pos-be/internal/entities" + "apskel-pos-be/internal/logger" + "apskel-pos-be/internal/models" + "apskel-pos-be/internal/repository" +) + +// Why an order earned nothing. +const ( + EarningSkipNotPaid = "NOT_PAID" + EarningSkipVoid = "VOID" + EarningSkipNoCustomer = "NO_CUSTOMER" + EarningSkipDefaultCustomer = "DEFAULT_CUSTOMER" + EarningSkipInactiveCustomer = "INACTIVE_CUSTOMER" + EarningSkipNothingToEarn = "NOTHING_TO_EARN" +) + +// EarningOutcome is what earning did for one order. +type EarningOutcome struct { + Points int64 + Coins int64 + // Set when the order earned nothing, to say why. + Skipped string +} + +type outletSettingsReader interface { + Outlet(ctx context.Context, outletID uuid.UUID) (*models.OutletLoyaltySettings, error) +} + +// EarningProcessor credits EnakPoint and EnakCoin for paid orders +// (docs/prd-point-coin.md F3). +type EarningProcessor struct { + orders repository.EarningRepository + settings outletSettingsReader + wallet *WalletProcessor + tx TxRunner +} + +func NewEarningProcessor(orders repository.EarningRepository, settings outletSettingsReader, wallet *WalletProcessor, tx TxRunner) *EarningProcessor { + return &EarningProcessor{orders: orders, settings: settings, wallet: wallet, tx: tx} +} + +// OnOrderPaid is called once an order has become fully paid and the payment has +// committed. It never fails the caller: a failed earning is logged and picked up later +// by EarnMissing, and the idempotency keys make that retry safe. +func (p *EarningProcessor) OnOrderPaid(ctx context.Context, orderID uuid.UUID) { + defer func() { + if r := recover(); r != nil { + logger.NonContext.Error(fmt.Sprintf("Earning for order %s panicked; it will be retried", orderID), fmt.Errorf("%v", r)) + } + }() + if _, err := p.EarnForOrder(ctx, orderID); err != nil { + logger.NonContext.Error(fmt.Sprintf("Earning for order %s failed; it will be retried", orderID), err) + } +} + +// EarnForOrder credits what a paid order earns. Calling it again for the same order +// credits nothing more. +func (p *EarningProcessor) EarnForOrder(ctx context.Context, orderID uuid.UUID) (*EarningOutcome, error) { + order, err := p.orders.GetOrderForEarning(ctx, orderID) + if err != nil { + return nil, err + } + if skip := earningSkipReason(order); skip != "" { + return &EarningOutcome{Skipped: skip}, nil + } + + settings, err := p.settings.Outlet(ctx, order.OutletID) + if err != nil { + return nil, err + } + pointPaid, err := p.orders.PointPaidAmount(ctx, orderID) + if err != nil { + return nil, err + } + result := CalculateEarning(&entities.Order{Subtotal: order.Subtotal, DiscountAmount: order.DiscountAmount}, pointPaid, *settings) + if result.Point.Amount == 0 && result.Coin.Amount == 0 { + return &EarningOutcome{Skipped: EarningSkipNothingToEarn}, nil + } + + outcome := &EarningOutcome{} + err = p.tx.WithTransaction(ctx, func(ctx context.Context) error { + for _, c := range []struct { + currency string + line EarningLine + total *int64 + }{ + {constants.WalletCurrencyPoint, result.Point, &outcome.Points}, + {constants.WalletCurrencyCoin, result.Coin, &outcome.Coins}, + } { + if c.line.Amount == 0 { + continue + } + outletID := order.OutletID + res, err := p.wallet.Credit(ctx, WalletCreditInput{WalletEntry: WalletEntry{ + CustomerID: *order.CustomerID, + Currency: c.currency, + Type: constants.WalletTxTypeEarn, + Amount: c.line.Amount, + ReferenceType: constants.WalletRefTypeOrder, + ReferenceID: order.ID, + OutletID: &outletID, + Description: earningDescription(order), + Metadata: result.Metadata(c.line), + IdempotencyKey: fmt.Sprintf("earn:%s:%s", order.ID, c.currency), + // Lots never expire until the expiry model is decided (F12, note N4). + }}) + if err != nil { + return fmt.Errorf("crediting %s: %w", c.currency, err) + } + *c.total = res.Transaction.Amount + } + return nil + }) + if err != nil { + return nil, err + } + return outcome, nil +} + +// EarnMissing is the safety net behind OnOrderPaid: it looks for orders paid since the +// given time that should have earned and have no EARN row, and earns for them. It +// returns how many orders it looked at and how many now earned. One order failing does +// not stop the others. +func (p *EarningProcessor) EarnMissing(ctx context.Context, since time.Time, maxOrders int) (checked, earned int, err error) { + const page = 200 + var after *repository.EarningCursor + for checked < maxOrders { + batch, err := p.orders.ListPaidOrdersWithoutEarning(ctx, since, after, page) + if err != nil { + return checked, earned, err + } + if len(batch) == 0 { + break + } + for _, candidate := range batch { + checked++ + outcome, err := p.EarnForOrder(ctx, candidate.ID) + if err != nil { + logger.NonContext.Error(fmt.Sprintf("Earning retry for order %s failed", candidate.ID), err) + continue + } + if outcome.Skipped == "" { + earned++ + } + } + last := batch[len(batch)-1] + after = &last + } + return checked, earned, nil +} + +func earningSkipReason(order *repository.EarningOrder) string { + switch { + case order.PaymentStatus != string(entities.PaymentStatusCompleted): + return EarningSkipNotPaid + case order.IsVoid: + return EarningSkipVoid + case order.CustomerID == nil || order.CustomerIsDefault == nil: + return EarningSkipNoCustomer + case *order.CustomerIsDefault: + return EarningSkipDefaultCustomer + case order.CustomerIsActive == nil || !*order.CustomerIsActive: + return EarningSkipInactiveCustomer + } + return "" +} + +func earningDescription(order *repository.EarningOrder) string { + description := "Belanja #" + order.OrderNumber + if order.OutletName != "" { + description += " di " + order.OutletName + } + return truncateRunes(description, walletDescriptionLimit) +} diff --git a/internal/processor/earning_processor_db_test.go b/internal/processor/earning_processor_db_test.go new file mode 100644 index 0000000..ae6565f --- /dev/null +++ b/internal/processor/earning_processor_db_test.go @@ -0,0 +1,221 @@ +package processor + +import ( + "context" + "errors" + "os" + "sync" + "testing" + "time" + + "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/models" + "apskel-pos-be/internal/repository" +) + +// failingSettings fails for the outlet settings until healed, to stand in for the +// database being unreachable right after a payment. +type failingSettings struct { + real outletSettingsReader + mu sync.Mutex + fail bool +} + +func (f *failingSettings) Outlet(ctx context.Context, outletID uuid.UUID) (*models.OutletLoyaltySettings, error) { + f.mu.Lock() + fail := f.fail + f.mu.Unlock() + if fail { + return nil, errors.New("connection refused") + } + return f.real.Outlet(ctx, outletID) +} + +// Needs TEST_DATABASE_URL pointing at a migrated database; see +// internal/repository/wallet_repository_test.go. +func TestEarningProcessor_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, user := uuid.New(), uuid.New() + earningOutlet, quietOutlet := uuid.New(), uuid.New() + regular, inactive := uuid.New(), uuid.New() + var walkIn uuid.UUID + var customers []uuid.UUID + 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 (?, 'earning test', 'basic')`, org) + exec(`INSERT INTO users (id, organization_id, name, email, password_hash, role) VALUES (?, ?, 'Kasir', ?, 'x', 'cashier')`, user, org, user.String()+"@test") + exec(`INSERT INTO outlets (id, organization_id, name) VALUES (?, ?, 'Kemang'), (?, ?, 'Tanpa Poin')`, earningOutlet, org, quietOutlet, org) + exec(`INSERT INTO customers (id, organization_id, name, is_default, is_active) VALUES + (?, ?, 'Budi', false, true), (?, ?, 'Nonaktif', false, false)`, regular, org, inactive, org) + // Creating the organization created its walk-in customer (trigger_create_default_customer). + var walkInID string + require.NoError(t, db.Raw(`SELECT id::text FROM customers WHERE organization_id = ? AND is_default`, org).Scan(&walkInID).Error) + walkIn = uuid.MustParse(walkInID) + customers = []uuid.UUID{regular, walkIn, inactive} + 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 orders WHERE organization_id = ?`, org) + db.Exec(`DELETE FROM loyalty_setting_changes WHERE organization_id = ?`, org) + db.Exec(`DELETE FROM outlet_settings WHERE outlet_id IN ?`, []uuid.UUID{earningOutlet, quietOutlet}) + db.Exec(`DELETE FROM customers WHERE id IN ?`, customers) + db.Exec(`DELETE FROM outlets WHERE id IN ?`, []uuid.UUID{earningOutlet, quietOutlet}) + db.Exec(`DELETE FROM users WHERE id = ?`, user) + db.Exec(`DELETE FROM organizations WHERE id = ?`, org) + }) + + txm := repository.NewTxManager(db) + settingsProcessor := NewLoyaltySettingsProcessor(repository.NewLoyaltySettingsRepository(db), txm) + outletSettings, err := settingsProcessor.Outlet(ctx, earningOutlet) + require.NoError(t, err) + outletSettings.Point.Enabled = true + outletSettings.Coin.Enabled = true + _, err = settingsProcessor.UpdateOutlet(ctx, org, earningOutlet, user, *outletSettings) + require.NoError(t, err) + + settings := &failingSettings{real: settingsProcessor} + earning := NewEarningProcessor(repository.NewEarningRepository(db), settings, NewWalletProcessor(repository.NewWalletRepository(db)), txm) + + orderNo := 0 + newOrder := func(outlet uuid.UUID, customer *uuid.UUID, paymentStatus string, isVoid bool, subtotal, discount float64) uuid.UUID { + t.Helper() + orderNo++ + id := uuid.New() + exec(`INSERT INTO orders (id, organization_id, outlet_id, user_id, customer_id, order_number, order_type, + subtotal, discount_amount, tax_amount, total_amount, payment_status, is_void) + VALUES (?, ?, ?, ?, ?, ?, 'dine_in', ?, ?, 0, ?, ?, ?)`, + id, org, outlet, user, customer, id.String()[:8]+"-"+string(rune('A'+orderNo)), subtotal, discount, subtotal-discount, paymentStatus, isVoid) + return id + } + earnRows := func(orderID uuid.UUID) map[string]int64 { + t.Helper() + var rows []struct { + Currency string + Amount int64 + } + require.NoError(t, db.Raw(`SELECT currency, amount FROM wallet_transactions WHERE reference_id = ? AND type = 'EARN'`, orderID).Scan(&rows).Error) + out := map[string]int64{} + for _, r := range rows { + out[r.Currency] += r.Amount + } + return out + } + + // Paid in full: the PRD example, 875 EnakPoint and 3 EnakCoin. + paid := newOrder(earningOutlet, ®ular, "completed", false, 97500, 10000) + earning.OnOrderPaid(ctx, paid) + assert.Equal(t, map[string]int64{"POINT": 875, "COIN": 3}, earnRows(paid)) + + var row struct { + Description string + Metadata string + OutletID string + } + require.NoError(t, db.Raw(`SELECT description, metadata::text AS metadata, outlet_id::text AS outlet_id FROM wallet_transactions + WHERE reference_id = ? AND currency = 'POINT'`, paid).Scan(&row).Error) + assert.Contains(t, row.Description, "di Kemang") + assert.Contains(t, row.Metadata, `"earn_per_amount": 100`, "the settings used are frozen on the row") + assert.Contains(t, row.Metadata, `"basis": 87500`) + assert.Equal(t, earningOutlet.String(), row.OutletID) + + // Called again, and five times at once, it is still one earning. + var wg sync.WaitGroup + for i := 0; i < 5; i++ { + wg.Add(1) + go func() { defer wg.Done(); earning.OnOrderPaid(ctx, paid) }() + } + wg.Wait() + outcome, err := earning.EarnForOrder(ctx, paid) + require.NoError(t, err) + assert.Equal(t, int64(875), outcome.Points, "a repeat reports the first earning") + assert.Equal(t, map[string]int64{"POINT": 875, "COIN": 3}, earnRows(paid)) + + // A self-order goes through the same payment path and the same rules. + selfOrder := newOrder(earningOutlet, ®ular, "completed", false, 25000, 0) + earning.OnOrderPaid(ctx, selfOrder) + assert.Equal(t, map[string]int64{"POINT": 250, "COIN": 1}, earnRows(selfOrder)) + + // A split bill earns once, on the payment that settles it: while partial, nothing. + split := newOrder(earningOutlet, ®ular, "partial", false, 60000, 0) + earning.OnOrderPaid(ctx, split) + assert.Empty(t, earnRows(split)) + exec(`UPDATE orders SET payment_status = 'completed' WHERE id = ?`, split) + earning.OnOrderPaid(ctx, split) + assert.Equal(t, map[string]int64{"POINT": 600, "COIN": 2}, earnRows(split)) + + // Orders that must not earn. + for name, c := range map[string]struct { + id uuid.UUID + skip string + }{ + "walk-in customer": {newOrder(earningOutlet, &walkIn, "completed", false, 50000, 0), EarningSkipDefaultCustomer}, + "inactive customer": {newOrder(earningOutlet, &inactive, "completed", false, 50000, 0), EarningSkipInactiveCustomer}, + "no customer": {newOrder(earningOutlet, nil, "completed", false, 50000, 0), EarningSkipNoCustomer}, + "void": {newOrder(earningOutlet, ®ular, "completed", true, 50000, 0), EarningSkipVoid}, + "unpaid": {newOrder(earningOutlet, ®ular, "pending", false, 50000, 0), EarningSkipNotPaid}, + "outlet not earning": {newOrder(quietOutlet, ®ular, "completed", false, 50000, 0), EarningSkipNothingToEarn}, + } { + outcome, err := earning.EarnForOrder(ctx, c.id) + require.NoError(t, err, name) + assert.Equal(t, c.skip, outcome.Skipped, name) + assert.Empty(t, earnRows(c.id), name) + } + + // An earning that fails does not surface to the payment, and the job picks it up. + settings.mu.Lock() + settings.fail = true + settings.mu.Unlock() + missed := newOrder(earningOutlet, ®ular, "completed", false, 40000, 0) + assert.NotPanics(t, func() { earning.OnOrderPaid(ctx, missed) }) + assert.Empty(t, earnRows(missed)) + + settings.mu.Lock() + settings.fail = false + settings.mu.Unlock() + since := time.Now().Add(-time.Hour) + checked, earned, err := earning.EarnMissing(ctx, since, 1000) + require.NoError(t, err) + assert.GreaterOrEqual(t, earned, 1) + assert.GreaterOrEqual(t, checked, earned) + assert.Equal(t, map[string]int64{"POINT": 400, "COIN": 1}, earnRows(missed)) + + // The job only looks at orders that could earn and have not. + candidates, err := repository.NewEarningRepository(db).ListPaidOrdersWithoutEarning(ctx, since, nil, 1000) + require.NoError(t, err) + var ours []uuid.UUID + for _, c := range candidates { + var n int64 + db.Raw(`SELECT COUNT(*) FROM orders WHERE id = ? AND organization_id = ?`, c.ID, org).Scan(&n) + if n > 0 { + ours = append(ours, c.ID) + } + } + assert.Empty(t, ours, "every eligible order of ours has earned; walk-in, inactive, void, unpaid and non-earning outlets are never candidates") + + // A second run finds nothing more to do for these orders. + _, _, err = earning.EarnMissing(ctx, since, 1000) + require.NoError(t, err) + assert.Equal(t, map[string]int64{"POINT": 400, "COIN": 1}, earnRows(missed)) + + balance := struct{ PointBalance, CoinBalance int64 }{} + require.NoError(t, db.Raw(`SELECT point_balance, coin_balance FROM customer_wallets WHERE customer_id = ?`, regular).Scan(&balance).Error) + assert.Equal(t, int64(875+250+600+400), balance.PointBalance) + assert.Equal(t, int64(3+1+2+1), balance.CoinBalance) +} diff --git a/internal/processor/order_paid_hook_test.go b/internal/processor/order_paid_hook_test.go new file mode 100644 index 0000000..b2ab4ba --- /dev/null +++ b/internal/processor/order_paid_hook_test.go @@ -0,0 +1,181 @@ +package processor + +import ( + "context" + "testing" + + "github.com/google/uuid" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + "apskel-pos-be/internal/entities" + "apskel-pos-be/internal/models" + "apskel-pos-be/internal/repository" +) + +// These fakes embed the interface they stand for and implement only what the paths +// under test call; anything else would panic, which would show a path doing more than +// expected. +type hookOrderRepo struct { + OrderRepository + order *entities.Order + statusUpdates int + hookCallsAtPay int // hook calls seen when the status was written + hook *orderPaidHookFake +} + +func (r *hookOrderRepo) GetByID(context.Context, uuid.UUID) (*entities.Order, error) { + o := *r.order + return &o, nil +} + +func (r *hookOrderRepo) GetWithRelations(context.Context, uuid.UUID) (*entities.Order, error) { + o := *r.order + return &o, nil +} + +func (r *hookOrderRepo) UpdateStatusSuccess(_ context.Context, _ uuid.UUID, status entities.OrderStatus, payment entities.PaymentStatus) error { + r.statusUpdates++ + r.hookCallsAtPay = len(r.hook.calls) + r.order.Status, r.order.PaymentStatus = status, payment + return nil +} + +type hookPaymentRepo struct { + PaymentRepository + created []*entities.Payment +} + +func (r *hookPaymentRepo) GetTotalPaidByOrderID(context.Context, uuid.UUID) (float64, error) { + return 0, nil +} + +func (r *hookPaymentRepo) Create(_ context.Context, p *entities.Payment) error { + p.ID = uuid.New() + r.created = append(r.created, p) + return nil +} + +func (r *hookPaymentRepo) GetByID(_ context.Context, id uuid.UUID) (*entities.Payment, error) { + for _, p := range r.created { + if p.ID == id { + return p, nil + } + } + return nil, nil +} + +func (r *hookPaymentRepo) GetByOrderID(context.Context, uuid.UUID) ([]*entities.Payment, error) { + return r.created, nil +} + +type hookPaymentMethodRepo struct{} + +func (hookPaymentMethodRepo) GetByID(_ context.Context, id uuid.UUID) (*entities.PaymentMethod, error) { + return &entities.PaymentMethod{ID: id}, nil +} + +type hookOrderItemRepo struct{ OrderItemRepository } + +func (hookOrderItemRepo) GetByOrderID(context.Context, uuid.UUID) ([]*entities.OrderItem, error) { + return nil, nil +} + +// splitFake settles the order on the payment that covers what is left, as the real +// split bill processor does. +type splitFake struct{ settle bool } + +func (f splitFake) split(order *entities.Order) (*models.SplitBillResponse, error) { + if f.settle { + order.PaymentStatus = entities.PaymentStatusCompleted + } else { + order.PaymentStatus = entities.PaymentStatusPartial + } + return &models.SplitBillResponse{OrderID: order.ID}, nil +} + +func (f splitFake) SplitByAmount(_ context.Context, _ *models.SplitBillRequest, order *entities.Order, _ *entities.PaymentMethod, _ *entities.Customer) (*models.SplitBillResponse, error) { + return f.split(order) +} + +func (f splitFake) SplitByItem(_ context.Context, _ *models.SplitBillRequest, order *entities.Order, _ *entities.PaymentMethod, _ *entities.Customer) (*models.SplitBillResponse, error) { + return f.split(order) +} + +type orderPaidHookFake struct { + calls []uuid.UUID + ctxs []context.Context +} + +func (h *orderPaidHookFake) OnOrderPaid(ctx context.Context, orderID uuid.UUID) { + h.calls = append(h.calls, orderID) + h.ctxs = append(h.ctxs, ctx) +} + +func newHookedOrderProcessor(split SplitBillProcessor) (*OrderProcessorImpl, *hookOrderRepo, *orderPaidHookFake) { + hook := &orderPaidHookFake{} + orders := &hookOrderRepo{ + order: &entities.Order{ID: uuid.New(), OrganizationID: uuid.New(), OutletID: uuid.New(), TotalAmount: 100000, PaymentStatus: entities.PaymentStatusPending}, + hook: hook, + } + p := &OrderProcessorImpl{ + orderRepo: orders, + orderItemRepo: hookOrderItemRepo{}, + paymentRepo: &hookPaymentRepo{}, + paymentMethodRepo: hookPaymentMethodRepo{}, + splitBillProcessor: split, + txManager: repository.NewTxManager(nil), + } + p.SetOrderPaidHook(hook) + return p, orders, hook +} + +func TestOrderPaidHook_CreatePayment(t *testing.T) { + p, orders, hook := newHookedOrderProcessor(nil) + ctx, cancel := context.WithCancel(context.Background()) + + _, err := p.CreatePayment(ctx, &models.CreatePaymentRequest{OrderID: orders.order.ID, PaymentMethodID: uuid.New(), Amount: 100000}) + require.NoError(t, err) + assert.Equal(t, []uuid.UUID{orders.order.ID}, hook.calls) + assert.Zero(t, orders.hookCallsAtPay, "the hook runs after the payment, not inside its transaction") + + // The hook's context outlives the request. + cancel() + assert.NoError(t, hook.ctxs[0].Err()) +} + +func TestOrderPaidHook_UpdateOrder(t *testing.T) { + p, orders, hook := newHookedOrderProcessor(nil) + _, err := p.UpdateOrder(context.Background(), orders.order.ID, &models.UpdateOrderRequest{}) + require.NoError(t, err) + assert.Equal(t, []uuid.UUID{orders.order.ID}, hook.calls) + assert.Equal(t, 1, orders.statusUpdates) + assert.Zero(t, orders.hookCallsAtPay) +} + +func TestOrderPaidHook_SplitBillOnlyOnTheSettlingPayment(t *testing.T) { + for _, splitType := range []string{"AMOUNT", "ITEM"} { + t.Run(splitType, func(t *testing.T) { + req := &models.SplitBillRequest{Type: splitType, PaymentMethodID: uuid.New()} + + p, orders, hook := newHookedOrderProcessor(splitFake{settle: false}) + req.OrderID = orders.order.ID + _, err := p.SplitBill(context.Background(), req) + require.NoError(t, err) + assert.Empty(t, hook.calls, "a partial split payment does not make the order paid") + + p, orders, hook = newHookedOrderProcessor(splitFake{settle: true}) + req.OrderID = orders.order.ID + _, err = p.SplitBill(context.Background(), req) + require.NoError(t, err) + assert.Equal(t, []uuid.UUID{orders.order.ID}, hook.calls) + }) + } +} + +func TestOrderPaidHook_NoHookIsFine(t *testing.T) { + p, orders, _ := newHookedOrderProcessor(nil) + p.SetOrderPaidHook(nil) + _, err := p.UpdateOrder(context.Background(), orders.order.ID, &models.UpdateOrderRequest{}) + assert.NoError(t, err) +} diff --git a/internal/processor/order_processor.go b/internal/processor/order_processor.go index 3e14e23..6cf9eb0 100644 --- a/internal/processor/order_processor.go +++ b/internal/processor/order_processor.go @@ -109,6 +109,30 @@ type OrderProcessorImpl struct { ingredientRepo IngredientRepository inventoryMovementService InventoryMovementService productOutletPriceRepo repository.ProductOutletPriceRepository + orderPaidHook OrderPaidHook +} + +// OrderPaidHook is told when an order has just become fully paid and the payment has +// committed. EarningProcessor is one (docs/prd-point-coin.md F3). +type OrderPaidHook interface { + OnOrderPaid(ctx context.Context, orderID uuid.UUID) +} + +// SetOrderPaidHook sets what runs when an order becomes fully paid. +func (p *OrderProcessorImpl) SetOrderPaidHook(hook OrderPaidHook) { + p.orderPaidHook = hook +} + +// onOrderPaid is the single place every path that completes an order's payment goes +// through: UpdateOrder, CreatePayment and both kinds of split bill. It must be called +// after the payment has committed. The hook runs detached from the caller's +// transaction and from the request being cancelled, and anything it does cannot fail +// the payment. +func (p *OrderProcessorImpl) onOrderPaid(ctx context.Context, orderID uuid.UUID) { + if p.orderPaidHook == nil { + return + } + p.orderPaidHook.OnOrderPaid(repository.DetachTransaction(context.WithoutCancel(ctx)), orderID) } func NewOrderProcessorImpl( @@ -494,6 +518,7 @@ func (p *OrderProcessorImpl) UpdateOrder(ctx context.Context, id uuid.UUID, req if err := p.orderRepo.UpdateStatusSuccess(ctx, order.ID, order.Status, order.PaymentStatus); err != nil { return nil, fmt.Errorf("failed to update order: %w", err) } + p.onOrderPaid(ctx, order.ID) orderWithRelations, err := p.orderRepo.GetWithRelations(ctx, id) if err != nil { @@ -818,6 +843,8 @@ func (p *OrderProcessorImpl) CreatePayment(ctx context.Context, req *models.Crea if err != nil { return nil, err } + // Not from updateOrderStatus: that runs inside the payment's transaction. + p.onOrderPaid(ctx, req.OrderID) paymentWithRelations, err := p.paymentRepo.GetByID(ctx, payment.ID) if err != nil { @@ -1207,6 +1234,10 @@ func (p *OrderProcessorImpl) SplitBill(ctx context.Context, req *models.SplitBil if err != nil { return nil, err } + // Both split paths mark the order paid on the payment that settles it. + if order.PaymentStatus == entities.PaymentStatusCompleted { + p.onOrderPaid(ctx, order.ID) + } return response, nil } diff --git a/internal/repository/earning_repository.go b/internal/repository/earning_repository.go new file mode 100644 index 0000000..f583c9d --- /dev/null +++ b/internal/repository/earning_repository.go @@ -0,0 +1,193 @@ +package repository + +import ( + "context" + "errors" + "fmt" + "time" + + "github.com/google/uuid" + "gorm.io/gorm" + + "apskel-pos-be/internal/constants" + "apskel-pos-be/internal/entities" +) + +// ErrEarningOrderNotFound means the order does not exist. +var ErrEarningOrderNotFound = errors.New("earning: order not found") + +// EarningOrder is what earning needs to know about an order. +type EarningOrder struct { + ID uuid.UUID + OrganizationID uuid.UUID + OutletID uuid.UUID + OrderNumber string + OutletName string + CustomerID *uuid.UUID + Subtotal float64 + DiscountAmount float64 + PaymentStatus string + IsVoid bool + // Nil when the order has no customer, or the customer row is gone. + CustomerIsDefault *bool + CustomerIsActive *bool +} + +// EarningCursor pages through orders by (updated_at, id). +type EarningCursor struct { + UpdatedAt time.Time + ID uuid.UUID +} + +// EarningRepository reads orders for loyalty earning (docs/prd-point-coin.md F3). +type EarningRepository interface { + GetOrderForEarning(ctx context.Context, orderID uuid.UUID) (*EarningOrder, error) + // PointPaidAmount is the rupiah part of the order paid with EnakPoint, which earns + // nothing (Q10). Zero until EnakPoint payment exists (phase 3). + PointPaidAmount(ctx context.Context, orderID uuid.UUID) (float64, error) + // ListPaidOrdersWithoutEarning pages, oldest first, through orders updated since + // the given time that are paid, not void, have an eligible customer, belong to an + // outlet that earns something, and have no EARN row yet. Pass the previous page's + // last cursor to continue; nil starts at the beginning. + ListPaidOrdersWithoutEarning(ctx context.Context, since time.Time, after *EarningCursor, limit int) ([]EarningCursor, error) + // ListEarnTransactions returns the EARN rows written for an order. + ListEarnTransactions(ctx context.Context, orderID uuid.UUID) ([]entities.WalletTransaction, error) +} + +type earningRepository struct { + db *gorm.DB +} + +func NewEarningRepository(db *gorm.DB) EarningRepository { + return &earningRepository{db: db} +} + +func (r *earningRepository) GetOrderForEarning(ctx context.Context, orderID uuid.UUID) (*EarningOrder, error) { + var rows []struct { + ID string + OrganizationID string + OutletID string + OrderNumber string + OutletName string + CustomerID *string + Subtotal float64 + DiscountAmount float64 + PaymentStatus string + IsVoid bool + CustomerIsDefault *bool + CustomerIsActive *bool + } + err := DBFromContext(ctx, r.db).WithContext(ctx).Raw(` + SELECT o.id::text AS id, o.organization_id::text AS organization_id, o.outlet_id::text AS outlet_id, + o.order_number, COALESCE(ou.name, '') AS outlet_name, o.customer_id::text AS customer_id, + o.subtotal, COALESCE(o.discount_amount, 0) AS discount_amount, o.payment_status, + COALESCE(o.is_void, false) AS is_void, + c.is_default AS customer_is_default, c.is_active AS customer_is_active + FROM orders o + LEFT JOIN outlets ou ON ou.id = o.outlet_id + LEFT JOIN customers c ON c.id = o.customer_id + WHERE o.id = ? + LIMIT 1`, orderID).Scan(&rows).Error + if err != nil { + return nil, fmt.Errorf("failed to get order for earning: %w", err) + } + if len(rows) == 0 { + return nil, ErrEarningOrderNotFound + } + row := rows[0] + order := &EarningOrder{ + OrderNumber: row.OrderNumber, + OutletName: row.OutletName, + Subtotal: row.Subtotal, + DiscountAmount: row.DiscountAmount, + PaymentStatus: row.PaymentStatus, + IsVoid: row.IsVoid, + CustomerIsDefault: row.CustomerIsDefault, + CustomerIsActive: row.CustomerIsActive, + } + order.ID, _ = uuid.Parse(row.ID) + order.OrganizationID, _ = uuid.Parse(row.OrganizationID) + order.OutletID, _ = uuid.Parse(row.OutletID) + if row.CustomerID != nil { + if id, err := uuid.Parse(*row.CustomerID); err == nil { + order.CustomerID = &id + } + } + return order, nil +} + +func (r *earningRepository) PointPaidAmount(ctx context.Context, orderID uuid.UUID) (float64, error) { + var total float64 + err := DBFromContext(ctx, r.db).WithContext(ctx).Raw(` + SELECT COALESCE(SUM(p.amount), 0) + FROM payments p + JOIN payment_methods pm ON pm.id = p.payment_method_id + WHERE p.order_id = ? AND pm.type = ? AND p.status = ?`, + orderID, constants.PaymentMethodTypePoint, entities.PaymentTransactionStatusCompleted). + Scan(&total).Error + if err != nil { + return 0, fmt.Errorf("failed to sum EnakPoint payments: %w", err) + } + return total, nil +} + +func (r *earningRepository) ListPaidOrdersWithoutEarning(ctx context.Context, since time.Time, after *EarningCursor, limit int) ([]EarningCursor, error) { + cursorAt, cursorID := since, uuid.Nil + if after != nil { + cursorAt, cursorID = after.UpdatedAt, after.ID + } + var rows []struct { + ID string + UpdatedAt time.Time + } + // An outlet that has neither currency switched on can never earn, so its orders are + // not candidates; otherwise every order of such an outlet would be rescanned on + // every run. + err := DBFromContext(ctx, r.db).WithContext(ctx).Raw(` + SELECT o.id::text AS id, o.updated_at + FROM orders o + JOIN customers c ON c.id = o.customer_id + WHERE o.payment_status = ? + AND COALESCE(o.is_void, false) = false + AND c.is_default = false AND c.is_active = true + AND o.updated_at >= ? + AND (o.updated_at, o.id) > (?, ?) + AND EXISTS ( + SELECT 1 FROM outlet_settings s + WHERE s.outlet_id = o.outlet_id + AND s.key IN (?, ?) + AND lower(trim(s.value)) IN ('true', 't', '1') + ) + AND NOT EXISTS ( + SELECT 1 FROM wallet_transactions t + WHERE t.reference_type = ? AND t.reference_id = o.id AND t.type = ? + ) + ORDER BY o.updated_at, o.id + LIMIT ?`, + entities.PaymentStatusCompleted, since, cursorAt, cursorID, + constants.LoyaltyPointEnabledKey, constants.LoyaltyCoinEnabledKey, + constants.WalletRefTypeOrder, constants.WalletTxTypeEarn, limit). + Scan(&rows).Error + if err != nil { + return nil, fmt.Errorf("failed to list paid orders without earning: %w", err) + } + out := make([]EarningCursor, 0, len(rows)) + for _, row := range rows { + if id, err := uuid.Parse(row.ID); err == nil { + out = append(out, EarningCursor{UpdatedAt: row.UpdatedAt, ID: id}) + } + } + return out, nil +} + +func (r *earningRepository) ListEarnTransactions(ctx context.Context, orderID uuid.UUID) ([]entities.WalletTransaction, error) { + var rows []entities.WalletTransaction + err := DBFromContext(ctx, r.db).WithContext(ctx). + Where("reference_type = ? AND reference_id = ? AND type = ?", constants.WalletRefTypeOrder, orderID, constants.WalletTxTypeEarn). + Order("currency"). + Find(&rows).Error + if err != nil { + return nil, fmt.Errorf("failed to list EARN rows: %w", err) + } + return rows, nil +} diff --git a/internal/repository/tx_manager.go b/internal/repository/tx_manager.go index ee9a3f8..33760c2 100644 --- a/internal/repository/tx_manager.go +++ b/internal/repository/tx_manager.go @@ -50,3 +50,10 @@ func (m *TxManager) WithTransactionOptions(ctx context.Context, opts *sql.TxOpti return fn(ctxTx) }, opts) } + +// DetachTransaction returns ctx without the caller's transaction, so work started from +// it (such as loyalty earning after a payment) reads committed data and commits on its +// own, whatever happens to the caller's transaction. +func DetachTransaction(ctx context.Context) context.Context { + return context.WithValue(ctx, txKey, (*gorm.DB)(nil)) +} diff --git a/internal/service/earning_backfill_job.go b/internal/service/earning_backfill_job.go new file mode 100644 index 0000000..131434c --- /dev/null +++ b/internal/service/earning_backfill_job.go @@ -0,0 +1,74 @@ +package service + +import ( + "context" + "sync" + "time" + + "apskel-pos-be/internal/logger" +) + +const ( + defaultEarningBackfillInterval = 30 * time.Minute + // How far back to look for paid orders that never earned. + earningBackfillWindow = 72 * time.Hour + // Orders looked at per run at most, so one run cannot run away. + earningBackfillMaxOrders = 5000 +) + +type missingEarner interface { + EarnMissing(ctx context.Context, since time.Time, maxOrders int) (checked, earned int, err error) +} + +// EarningBackfillJob is the safety net behind earning at payment time +// (docs/prd-point-coin.md F3, PC-203). Every run it earns for orders paid in the last +// few days that should have earned and did not, for example because the database was +// briefly unreachable right after the payment committed. +type EarningBackfillJob struct { + earner missingEarner + now func() time.Time + stopCh chan struct{} + stopOnce sync.Once +} + +func NewEarningBackfillJob(earner missingEarner) *EarningBackfillJob { + return &EarningBackfillJob{earner: earner, now: time.Now, stopCh: make(chan struct{})} +} + +func (j *EarningBackfillJob) Start(interval time.Duration) { + if interval <= 0 { + interval = defaultEarningBackfillInterval + } + go func() { + j.RunOnce(context.Background()) + ticker := time.NewTicker(interval) + defer ticker.Stop() + for { + select { + case <-ticker.C: + j.RunOnce(context.Background()) + case <-j.stopCh: + return + } + } + }() + logger.NonContext.Infof("Earning backfill job started (interval: %s)", interval) +} + +func (j *EarningBackfillJob) Stop() { + j.stopOnce.Do(func() { close(j.stopCh) }) +} + +// RunOnce earns for every missed order in the window and reports how many it fixed. +// It is quiet when nothing was missed. +func (j *EarningBackfillJob) RunOnce(ctx context.Context) int { + checked, earned, err := j.earner.EarnMissing(ctx, j.now().Add(-earningBackfillWindow), earningBackfillMaxOrders) + if err != nil { + logger.NonContext.Error("Earning backfill failed to run", err) + } + if earned > 0 { + logger.NonContext.WarnWithFields("Earning backfill credited orders that had missed their earning", + map[string]interface{}{"checked": checked, "earned": earned}, nil) + } + return earned +} diff --git a/internal/service/earning_backfill_job_test.go b/internal/service/earning_backfill_job_test.go new file mode 100644 index 0000000..31146b5 --- /dev/null +++ b/internal/service/earning_backfill_job_test.go @@ -0,0 +1,46 @@ +package service + +import ( + "context" + "errors" + "testing" + "time" + + "github.com/stretchr/testify/assert" + + "apskel-pos-be/internal/logger" +) + +type missingEarnerFake struct { + since time.Time + max int + earned int + err error + calls int +} + +func (f *missingEarnerFake) EarnMissing(_ context.Context, since time.Time, maxOrders int) (int, int, error) { + f.calls++ + f.since, f.max = since, maxOrders + return f.earned * 2, f.earned, f.err +} + +func TestEarningBackfillJob(t *testing.T) { + logger.Setup("fatal", "json") + now := time.Date(2026, 9, 30, 12, 0, 0, 0, time.UTC) + earner := &missingEarnerFake{earned: 3} + job := NewEarningBackfillJob(earner) + job.now = func() time.Time { return now } + + assert.Equal(t, 3, job.RunOnce(context.Background())) + assert.Equal(t, now.Add(-72*time.Hour), earner.since, "looks back three days") + assert.Equal(t, earningBackfillMaxOrders, earner.max) + + // A failing run is logged, not fatal. + earner.earned, earner.err = 0, errors.New("db down") + assert.Equal(t, 0, job.RunOnce(context.Background())) + + job.Start(time.Hour) + job.Stop() + job.Stop() +}