EnakGame phase 10 of docs/tasks-enakgame.md (EG-1001 to EG-1003). Spin (EG-1001) - PROBABILITY entries take an optional label (a wheel segment). The customer game list shows a PROBABILITY game's prizes (entry, label, amount, never weights), and completing returns the drawn prize, so the client can draw the wheel and stop it on the server's draw. - docs/enakgame-spin.md: the admin steps to set up spin per organization (no seeder) and the customer app flow. An HTTP test plays it end to end. Old game flow removed (EG-1002) - Routes POST /customer/spin, GET /customer/games, GET /customer/ferris-wheel, and admin /marketing/games, /marketing/game-prizes, /marketing/rewards, with their handlers, services, processors, repositories, validators, models, contracts, mappers and tests (GamePlayProcessor, SpinGameService, rewards, ...). This also closes RFC §15 findings 1 and 2 (double charge, spinning another org's game). - Tables games, game_prizes, game_plays and rewards stay for ledger history. entities.StringSlice moves to its own file; the omset tracker (unrouted) keeps game_id but no longer embeds the old game response. games.is_active dropped (EG-1003) - Migration 000115; nothing reads metadata.coin_cost any more. The EnakPoint integration docs now point at /customer/enakgame. The Postgres tests were not run: no test database here. Migration 000115 has not been run anywhere. The customer app must stop calling the removed endpoints before this is deployed. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com>
461 lines
16 KiB
Go
461 lines
16 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.
|
|
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])
|
|
}
|