166 lines
5.4 KiB
Go
166 lines
5.4 KiB
Go
package service
|
|||
|
|
|
||
|
|
import (
|
||
|
|
"context"
|
||
|
|
"sync"
|
||
|
|
"time"
|
||
|
|
|
||
|
|
"apskel-pos-be/internal/logger"
|
||
|
|
)
|
||
|
|
|
||
|
|
const (
|
||
|
|
defaultGameSessionJobInterval = time.Minute
|
||
|
|
defaultGameBudgetPeriodJobInterval = 24 * time.Hour
|
||
|
|
)
|
||
|
|
|
||
|
|
// tickerJob runs a function now and then on every tick until stopped. It keeps no
|
||
|
|
// state in memory: the work it runs is safe on several instances at once.
|
||
|
|
type tickerJob struct {
|
||
|
|
name string
|
||
|
|
run func(ctx context.Context)
|
||
|
|
stopCh chan struct{}
|
||
|
|
stopOnce sync.Once
|
||
|
|
}
|
||
|
|
|
||
|
|
func (j *tickerJob) start(interval, fallback time.Duration) {
|
||
|
|
if interval <= 0 {
|
||
|
|
interval = fallback
|
||
|
|
}
|
||
|
|
go func() {
|
||
|
|
j.run(context.Background())
|
||
|
|
ticker := time.NewTicker(interval)
|
||
|
|
defer ticker.Stop()
|
||
|
|
for {
|
||
|
|
select {
|
||
|
|
case <-ticker.C:
|
||
|
|
j.run(context.Background())
|
||
|
|
case <-j.stopCh:
|
||
|
|
return
|
||
|
|
}
|
||
|
|
}
|
||
|
|
}()
|
||
|
|
logger.NonContext.Infof("%s started (interval: %s)", j.name, interval)
|
||
|
|
}
|
||
|
|
|
||
|
|
func (j *tickerJob) stop() {
|
||
|
|
j.stopOnce.Do(func() { close(j.stopCh) })
|
||
|
|
}
|
||
|
|
|
||
|
|
type gameSessionWork interface {
|
||
|
|
ProcessDueSessions(ctx context.Context) (refunded, expired int, err error)
|
||
|
|
}
|
||
|
|
|
||
|
|
// GameSessionJob refunds or expires the EnakGame sessions left STARTED
|
||
|
|
// (docs/rfc-enakgame.md §7.3). Every session is handled under its customer's wallet
|
||
|
|
// lock with a conditional update, and the refund carries an idempotency key, so
|
||
|
|
// several instances and repeated runs refund nothing twice.
|
||
|
|
type GameSessionJob struct{ job tickerJob }
|
||
|
|
|
||
|
|
func NewGameSessionJob(work gameSessionWork) *GameSessionJob {
|
||
|
|
j := &GameSessionJob{}
|
||
|
|
j.job = tickerJob{name: "Game session job", stopCh: make(chan struct{}), run: func(ctx context.Context) {
|
||
|
|
refunded, expired, err := work.ProcessDueSessions(ctx)
|
||
|
|
if err != nil {
|
||
|
|
logger.NonContext.Error("Game session job failed", err)
|
||
|
|
}
|
||
|
|
if refunded > 0 || expired > 0 {
|
||
|
|
logger.NonContext.Infof("Game session job refunded %d and expired %d sessions", refunded, expired)
|
||
|
|
}
|
||
|
|
}}
|
||
|
|
return j
|
||
|
|
}
|
||
|
|
|
||
|
|
func (j *GameSessionJob) Start(interval time.Duration) {
|
||
|
|
j.job.start(interval, defaultGameSessionJobInterval)
|
||
|
|
}
|
||
|
|
func (j *GameSessionJob) Stop() { j.job.stop() }
|
||
|
|
|
||
|
|
type gameBudgetPeriodWork interface {
|
||
|
|
CreateNextPeriods(ctx context.Context) (int, error)
|
||
|
|
}
|
||
|
|
|
||
|
|
// GameBudgetPeriodJob gives each running global EnakGame budget its successor for the
|
||
|
|
// following month, when there is none yet (§12). The unique index on the period start
|
||
|
|
// keeps repeated runs and several instances from creating it twice.
|
||
|
|
type GameBudgetPeriodJob struct{ job tickerJob }
|
||
|
|
|
||
|
|
func NewGameBudgetPeriodJob(work gameBudgetPeriodWork) *GameBudgetPeriodJob {
|
||
|
|
j := &GameBudgetPeriodJob{}
|
||
|
|
j.job = tickerJob{name: "Game budget period job", stopCh: make(chan struct{}), run: func(ctx context.Context) {
|
||
|
|
created, err := work.CreateNextPeriods(ctx)
|
||
|
|
if err != nil {
|
||
|
|
logger.NonContext.Error("Game budget period job failed", err)
|
||
|
|
}
|
||
|
|
if created > 0 {
|
||
|
|
logger.NonContext.Infof("Game budget period job created %d budgets", created)
|
||
|
|
}
|
||
|
|
}}
|
||
|
|
return j
|
||
|
|
}
|
||
|
|
|
||
|
|
func (j *GameBudgetPeriodJob) Start(interval time.Duration) {
|
||
|
|
j.job.start(interval, defaultGameBudgetPeriodJobInterval)
|
||
|
|
}
|
||
|
|
func (j *GameBudgetPeriodJob) Stop() { j.job.stop() }
|
||
|
|
|
||
|
|
const defaultVoucherCodeExpiryJobInterval = time.Hour
|
||
|
|
|
||
|
|
type voucherCodeExpiryWork interface {
|
||
|
|
ExpireCodes(ctx context.Context, now time.Time) (int64, error)
|
||
|
|
}
|
||
|
|
|
||
|
|
// VoucherCodeExpiryJob moves the voucher codes past their expiry from AVAILABLE to
|
||
|
|
// EXPIRED (docs/rfc-enakgame.md §12). Redemption already skips them; this keeps the
|
||
|
|
// counts honest. Rows are taken with SKIP LOCKED, so instances share the work.
|
||
|
|
type VoucherCodeExpiryJob struct{ job tickerJob }
|
||
|
|
|
||
|
|
func NewVoucherCodeExpiryJob(work voucherCodeExpiryWork) *VoucherCodeExpiryJob {
|
||
|
|
j := &VoucherCodeExpiryJob{}
|
||
|
|
j.job = tickerJob{name: "Voucher code expiry job", stopCh: make(chan struct{}), run: func(ctx context.Context) {
|
||
|
|
expired, err := work.ExpireCodes(ctx, time.Now())
|
||
|
|
if err != nil {
|
||
|
|
logger.NonContext.Error("Voucher code expiry job failed", err)
|
||
|
|
}
|
||
|
|
if expired > 0 {
|
||
|
|
logger.NonContext.Infof("Voucher code expiry job expired %d codes", expired)
|
||
|
|
}
|
||
|
|
}}
|
||
|
|
return j
|
||
|
|
}
|
||
|
|
|
||
|
|
func (j *VoucherCodeExpiryJob) Start(interval time.Duration) {
|
||
|
|
j.job.start(interval, defaultVoucherCodeExpiryJobInterval)
|
||
|
|
}
|
||
|
|
func (j *VoucherCodeExpiryJob) Stop() { j.job.stop() }
|
||
|
|
|
||
|
|
const defaultVoucherRecoveryJobInterval = time.Minute
|
||
|
|
|
||
|
|
type voucherRecoveryWork interface {
|
||
|
|
RecoverPending(ctx context.Context) (completed, failed int, err error)
|
||
|
|
}
|
||
|
|
|
||
|
|
// VoucherRedemptionRecoveryJob settles the EXTERNAL voucher redemptions their provider
|
||
|
|
// left without an answer (docs/rfc-enakgame.md §7.5 step 4). Redemptions are claimed
|
||
|
|
// with SKIP LOCKED and settled with conditional updates, so several instances and
|
||
|
|
// repeated runs never refund one twice.
|
||
|
|
type VoucherRedemptionRecoveryJob struct{ job tickerJob }
|
||
|
|
|
||
|
|
func NewVoucherRedemptionRecoveryJob(work voucherRecoveryWork) *VoucherRedemptionRecoveryJob {
|
||
|
|
j := &VoucherRedemptionRecoveryJob{}
|
||
|
|
j.job = tickerJob{name: "Voucher redemption recovery job", stopCh: make(chan struct{}), run: func(ctx context.Context) {
|
||
|
|
completed, failed, err := work.RecoverPending(ctx)
|
||
|
|
if err != nil {
|
||
|
|
logger.NonContext.Error("Voucher redemption recovery job failed", err)
|
||
|
|
}
|
||
|
|
if completed > 0 || failed > 0 {
|
||
|
|
logger.NonContext.Infof("Voucher redemption recovery job completed %d and failed %d redemptions", completed, failed)
|
||
|
|
}
|
||
|
|
}}
|
||
|
|
return j
|
||
|
|
}
|
||
|
|
|
||
|
|
func (j *VoucherRedemptionRecoveryJob) Start(interval time.Duration) {
|
||
|
|
j.job.start(interval, defaultVoucherRecoveryJobInterval)
|
||
|
|
}
|
||
|
|
func (j *VoucherRedemptionRecoveryJob) Stop() { j.job.stop() }
|