package processor import ( "context" "encoding/json" "errors" "fmt" "strings" "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" ) // ErrGameSessionRejected wraps every reason a customer cannot start a game: the key, // the game, its configuration, the budget, the EnakCoin balance. The message says // which. var ErrGameSessionRejected = errors.New("game session refused") // errGameSessionMoved rolls back a refund whose session another flow finished first. var errGameSessionMoved = errors.New("game session is no longer started") const ( // gameSessionKeyLimit keeps the client's Idempotency-Key short enough to fit, with // the prefix scoping it to the customer, in wallet_transactions.idempotency_key. gameSessionKeyLimit = 50 // gameSessionJobBatch is how many sessions the job reads at a time, and // gameSessionJobBatches how many batches one run goes through at most. gameSessionJobBatch = 100 gameSessionJobBatches = 10 // customerGameListLimit caps the games the customer app gets in one list. customerGameListLimit = 100 ) func gameSessionRejected(format string, args ...any) error { return fmt.Errorf("%w: %s", ErrGameSessionRejected, fmt.Sprintf(format, args...)) } type gameCustomerReader interface { GetCustomer(ctx context.Context, customerID uuid.UUID) (*repository.WalletMoveCustomer, error) } // GameSessionProcessor starts and completes EnakGame sessions, refunds their entry // cost when the system is at fault, and ends the ones left behind // (docs/rfc-enakgame.md §7.1–§7.3). // // Every flow locks the customer's wallet before anything else (P4), and a session // leaves STARTED only through a conditional update (D4), so a refund never races a // completion or another refund of the same session. type GameSessionProcessor struct { customers gameCustomerReader games repository.EnakGameRepository sessions repository.GameSessionRepository budgets repository.GameBudgetRepository events repository.GameEventRepository counters repository.GameRewardCounterRepository settings organizationSettingsReader spendable spendableReader wallet *WalletProcessor audit *AuditLogger tx TxRunner rng RewardRNG now func() time.Time } func NewGameSessionProcessor(customers gameCustomerReader, games repository.EnakGameRepository, sessions repository.GameSessionRepository, budgets repository.GameBudgetRepository, events repository.GameEventRepository, counters repository.GameRewardCounterRepository, settings organizationSettingsReader, spendable spendableReader, wallet *WalletProcessor, audit *AuditLogger, tx TxRunner) *GameSessionProcessor { return &GameSessionProcessor{ customers: customers, games: games, sessions: sessions, budgets: budgets, events: events, counters: counters, settings: settings, spendable: spendable, wallet: wallet, audit: audit, tx: tx, rng: CryptoRewardRNG{}, now: time.Now, } } // Start takes a game's entry cost in EnakCoin and opens a session, in one transaction // (§7.1, PRD §10.1). The entry cost and the active reward configuration are frozen on // the session, so changing the game later does not change it. // // idempotencyKey is the client's Idempotency-Key: a retry with the same key returns // the session the first attempt opened and takes nothing more. func (p *GameSessionProcessor) Start(ctx context.Context, customerID, gameID uuid.UUID, idempotencyKey string) (*models.GameSessionStart, error) { key := strings.TrimSpace(idempotencyKey) if key == "" { return nil, gameSessionRejected("the Idempotency-Key header is required") } if len(key) > gameSessionKeyLimit { return nil, gameSessionRejected("the Idempotency-Key header must be at most %d characters", gameSessionKeyLimit) } customer, err := p.customers.GetCustomer(ctx, customerID) if err != nil { return nil, err } if !customer.IsActive { return nil, gameSessionRejected("the customer is not active") } walletKey := fmt.Sprintf("game-entry:%s:%s", customerID, key) var session *entities.GameSession replayed := false err = p.tx.WithTransaction(ctx, func(ctx context.Context) error { if err := p.wallet.LockWallet(ctx, customerID); err != nil { return err } previous, err := p.wallet.FindTransaction(ctx, walletKey) if err != nil { return err } if previous != nil { session, err = p.sessions.GetSessionBySpendTransaction(ctx, previous.ID) if err != nil { return err } if session.GameID != gameID { return gameSessionRejected("this Idempotency-Key was already used to start another game") } replayed = true return nil } game, err := p.games.GetGame(ctx, customer.OrganizationID, gameID) if err != nil { return err } if game.Status != constants.GameStatusActive { return gameSessionRejected("the game is not available") } config, err := p.games.GetActiveRewardConfig(ctx, customer.OrganizationID, gameID) if errors.Is(err, repository.ErrGameRewardConfigNotFound) { return gameSessionRejected("the game has no active reward configuration") } if err != nil { return err } // A reward must name the budget paying for it (D6), so no game starts while // none is set for the period. now := p.now() if _, err := p.budgets.GetGlobalBudgetOn(ctx, customer.OrganizationID, walletDay(now)); err != nil { if errors.Is(err, repository.ErrGameBudgetNotFound) { return gameSessionRejected("no EnakGame budget is set for this period") } return err } sessionID := uuid.New() spend, err := p.wallet.Debit(ctx, WalletDebitInput{WalletEntry: WalletEntry{ CustomerID: customerID, Currency: constants.WalletCurrencyCoin, Type: constants.WalletTxTypeGameSpend, Amount: game.EntryCost, ReferenceType: constants.WalletRefTypeGameSession, ReferenceID: sessionID, Description: truncateDescription("Main " + game.Name), Metadata: entities.Metadata{"game_id": game.ID.String(), "entry_cost": game.EntryCost}, IdempotencyKey: walletKey, }}) if errors.Is(err, repository.ErrWalletInsufficientBalance) { return gameSessionRejected("not enough EnakCoin") } if err != nil { return err } session = &entities.GameSession{ ID: sessionID, OrganizationID: customer.OrganizationID, CustomerID: customerID, GameID: game.ID, RewardConfigID: config.ID, EntryCost: game.EntryCost, Status: constants.GameSessionStatusStarted, StartedAt: now, ExpiresAt: now.Add(time.Duration(game.SessionTTLSeconds) * time.Second), SpendTransactionID: spend.Transaction.ID, } return p.sessions.CreateSession(ctx, session) }) if err != nil { return nil, err } balances, err := p.spendable.SpendableBalances(ctx, customerID, p.now()) if err != nil { return nil, err } return &models.GameSessionStart{ SessionID: session.ID, GameID: session.GameID, EntryCost: session.EntryCost, ExpiresAt: session.ExpiresAt, CoinBalance: balances[constants.WalletCurrencyCoin], Replayed: replayed, }, nil } // RefundSession gives back the entry cost of a STARTED session, in a transaction of // its own (§7.3, PRD §10.2). It reports false when the session had already left // STARTED, which changes nothing. func (p *GameSessionProcessor) RefundSession(ctx context.Context, organizationID, sessionID uuid.UUID, reason, source string) (bool, error) { refunded := false err := p.tx.WithTransaction(ctx, func(ctx context.Context) error { session, err := p.sessions.GetSession(ctx, organizationID, sessionID) if err != nil { return err } if err := p.wallet.LockWallet(ctx, session.CustomerID); err != nil { return err } // Read again under the lock: another flow may have finished it meanwhile. if session, err = p.sessions.GetSession(ctx, organizationID, sessionID); err != nil { return err } refunded, err = p.refundLocked(ctx, session, reason, source) return err }) if errors.Is(err, errGameSessionMoved) { return false, nil } return refunded, err } // refundLocked refunds a session inside the caller's transaction, with the // customer's wallet already locked. The EnakCoin go back to lots that keep the expiry // of those the entry cost was taken from, but last at least seven days (§6.3). // // When the session leaves STARTED between the read and the update, it returns // errGameSessionMoved so the caller rolls the credit back. func (p *GameSessionProcessor) refundLocked(ctx context.Context, session *entities.GameSession, reason, source string) (bool, error) { if session.Status != constants.GameSessionStatusStarted { return false, nil } lots, err := p.wallet.RefundLots(ctx, session.SpendTransactionID) if err != nil { return false, err } spendID := session.SpendTransactionID refund, err := p.wallet.Credit(ctx, WalletCreditInput{ WalletEntry: WalletEntry{ CustomerID: session.CustomerID, Currency: constants.WalletCurrencyCoin, Type: constants.WalletTxTypeGameSpendRefund, Amount: session.EntryCost, ReferenceType: constants.WalletRefTypeGameSession, ReferenceID: session.ID, ReversesTransactionID: &spendID, Description: "Pengembalian biaya main game", Metadata: entities.Metadata{"game_id": session.GameID.String(), "refund_reason": reason}, IdempotencyKey: "game-refund:" + session.ID.String(), }, Lots: lots, }) if err != nil { return false, err } now := p.now() moved, err := p.sessions.RefundSession(ctx, session.ID, refund.Transaction.ID, reason, now) if err != nil { return false, err } if !moved { return false, errGameSessionMoved } err = p.audit.Record(ctx, AuditEntry{ OrganizationID: session.OrganizationID, ActorType: constants.AuditActorSystem, EntityType: constants.AuditEntityGameSession, EntityID: session.ID, Action: "REFUNDED", Before: map[string]any{"status": constants.GameSessionStatusStarted}, After: map[string]any{ "status": constants.GameSessionStatusRefunded, "refund_reason": reason, "refund_transaction_id": refund.Transaction.ID, "amount": session.EntryCost, }, Reason: &reason, Source: source, }) if err != nil { return false, err } session.Status, session.RefundReason, session.RefundTransactionID, session.EndedAt = constants.GameSessionStatusRefunded, &reason, &refund.Transaction.ID, &now return true, nil } // ProcessDueSessions is the session job (§7.3): a STARTED session whose game is no // longer ACTIVE is refunded at once; an expired one is refunded when completing it // failed on a system error, and otherwise only expired, keeping the entry cost. // Each session is handled in its own transaction. It returns how many it refunded and // expired. func (p *GameSessionProcessor) ProcessDueSessions(ctx context.Context) (refunded, expired int, err error) { var failures []error for batch := 0; batch < gameSessionJobBatches; batch++ { due, err := p.sessions.ListDueSessions(ctx, p.now(), gameSessionJobBatch) if err != nil { return refunded, expired, err } changed := 0 for _, d := range due { action, err := p.processDue(ctx, d) switch { case err != nil: failures = append(failures, fmt.Errorf("session %s: %w", d.ID, err)) logger.NonContext.Error("Game session job failed on a session", err) case action == constants.GameSessionStatusRefunded: refunded++ changed++ case action == constants.GameSessionStatusExpired: expired++ changed++ } } // A short batch was the last; a batch where nothing changed would only be // read again. if len(due) < gameSessionJobBatch || changed == 0 { break } } return refunded, expired, errors.Join(failures...) } // processDue handles one session and returns the status it moved it to, or "" when it // left it alone. func (p *GameSessionProcessor) processDue(ctx context.Context, d repository.DueGameSession) (string, error) { action := "" err := p.tx.WithTransaction(ctx, func(ctx context.Context) error { if err := p.wallet.LockWallet(ctx, d.CustomerID); err != nil { return err } session, err := p.sessions.GetSession(ctx, d.OrganizationID, d.ID) if err != nil { return err } if session.Status != constants.GameSessionStatusStarted { return nil } game, err := p.games.GetGame(ctx, session.OrganizationID, session.GameID) if err != nil { return err } now := p.now() reason := "" switch { case game.Status != constants.GameStatusActive: reason = constants.GameSessionRefundGameDeactivated case now.Before(session.ExpiresAt): return nil case session.CompletionFailedAt != nil: reason = constants.GameSessionRefundSystemError default: // Left behind by the customer: no refund (PRD §10.2). moved, err := p.sessions.ExpireSession(ctx, session.ID, now) if moved { action = constants.GameSessionStatusExpired } return err } refunded, err := p.refundLocked(ctx, session, reason, constants.AuditSourceSessionJob) if refunded { action = constants.GameSessionStatusRefunded } return err }) if errors.Is(err, errGameSessionMoved) { return "", nil } return action, err } // ListGames returns the ACTIVE games of the customer's organization, each with the // events making it pay more right now. func (p *GameSessionProcessor) ListGames(ctx context.Context, customerID uuid.UUID) ([]models.CustomerEnakGame, error) { customer, err := p.customers.GetCustomer(ctx, customerID) if err != nil { return nil, err } games, _, err := p.games.ListGames(ctx, repository.EnakGameFilter{ OrganizationID: customer.OrganizationID, Statuses: []string{constants.GameStatusActive}, Limit: customerGameListLimit, }) if err != nil { return nil, err } events, err := p.events.ActiveEventsByGame(ctx, customer.OrganizationID, p.now()) if err != nil { return nil, err } configs, err := p.games.ListActiveRewardConfigs(ctx, customer.OrganizationID) if err != nil { return nil, err } prizes := make(map[uuid.UUID][]models.GamePrize, len(configs)) for _, c := range configs { if c.RewardType != constants.GameRewardTypeProbability { continue } // Rules passed Validate when the configuration was made. if list, err := probabilityPrizes(json.RawMessage(c.Rules)); err == nil { prizes[c.GameID] = list } } out := make([]models.CustomerEnakGame, 0, len(games)) for _, g := range games { item := models.CustomerEnakGame{ ID: g.ID, Name: g.Name, Description: g.Description, ThumbnailURL: g.ThumbnailURL, GameURL: g.GameURL, Version: g.Version, EntryCost: g.EntryCost, SessionTTLSeconds: g.SessionTTLSeconds, Events: []models.CustomerGameEvent{}, Prizes: prizes[g.ID], } for _, e := range events[g.ID] { item.Events = append(item.Events, models.CustomerGameEvent{ ID: e.ID, Name: e.Name, BannerURL: e.BannerURL, Multiplier: e.Multiplier, Bonus: e.Bonus, EndAt: e.EndAt, }) } if g.Slug != nil { item.Slug = *g.Slug } out = append(out, item) } return out, nil } // ListSessions returns a page of the customer's sessions, newest first. func (p *GameSessionProcessor) ListSessions(ctx context.Context, customerID uuid.UUID, page, limit int) (*models.PaginatedResponse[models.CustomerGameSession], error) { page, limit = enakGamePage(page, limit) sessions, total, err := p.sessions.ListCustomerSessions(ctx, customerID, (page-1)*limit, limit) if err != nil { return nil, err } items := make([]models.CustomerGameSession, 0, len(sessions)) for i := range sessions { items = append(items, customerGameSessionModel(&sessions[i])) } return &models.PaginatedResponse[models.CustomerGameSession]{Data: items, Pagination: enakGamePagination(page, limit, total)}, nil } // GetSession returns one of the customer's sessions; another customer's is not found. func (p *GameSessionProcessor) GetSession(ctx context.Context, customerID, sessionID uuid.UUID) (*models.CustomerGameSession, error) { session, err := p.sessions.GetCustomerSession(ctx, customerID, sessionID) if err != nil { return nil, err } m := customerGameSessionModel(session) return &m, nil } func customerGameSessionModel(s *entities.GameSession) models.CustomerGameSession { return models.CustomerGameSession{ ID: s.ID, GameID: s.GameID, Status: s.Status, EntryCost: s.EntryCost, RewardTotal: s.RewardTotal, StartedAt: s.StartedAt, ExpiresAt: s.ExpiresAt, EndedAt: s.EndedAt, RefundReason: s.RefundReason, } } // truncateDescription keeps a ledger description within its 255 characters. func truncateDescription(s string) string { runes := []rune(s) if len(runes) <= 255 { return s } return string(runes[:255]) }