Files
apskel-pos-backend/internal/processor/game_session_processor.go
efrilmandClaude Opus 5.5 52e8fe11c6 feat(enakgame): filter play history by game and status
GET /customer/enakgame/sessions takes optional game_id and status, so a game
reloaded mid-play finds the session it was running (status=STARTED) instead
of starting a new one and charging EnakCoin again. An invalid game_id or
status is refused.

integration-enakgame.md §4.4 now describes recovery after a reload: keep the
session_id in sessionStorage, continue a STARTED session before expires_at,
and call complete again for a COMPLETED one to get the full answer, prize
included. The mobile guide and RFC §11 mention the filters.

Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
2026-10-08 12:51:39 +07:00

477 lines
17 KiB
Go

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, of one game
// or status when asked.
func (p *GameSessionProcessor) ListSessions(ctx context.Context, customerID uuid.UUID, q models.GameSessionListQuery) (*models.PaginatedResponse[models.CustomerGameSession], error) {
page, limit := enakGamePage(q.Page, q.Limit)
filter := repository.CustomerSessionFilter{CustomerID: customerID, Offset: (page - 1) * limit, Limit: limit}
if s := strings.TrimSpace(q.GameID); s != "" {
id, err := uuid.Parse(s)
if err != nil {
return nil, gameSessionRejected("game_id must be a UUID")
}
filter.GameID = &id
}
switch status := strings.ToUpper(strings.TrimSpace(q.Status)); status {
case "", constants.GameSessionStatusStarted, constants.GameSessionStatusCompleted,
constants.GameSessionStatusRefunded, constants.GameSessionStatusExpired:
filter.Status = status
default:
return nil, gameSessionRejected("status must be STARTED, COMPLETED, REFUNDED or EXPIRED")
}
sessions, total, err := p.sessions.ListCustomerSessions(ctx, filter)
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])
}