From b636142e27c51c0ea6d2967ba895bb40a4adf686 Mon Sep 17 00:00:00 2001 From: zizi Date: Wed, 20 May 2026 15:56:46 +0800 Subject: [PATCH] fix: schedule affiliate rebate releases Run an hourly rebate release task at startup so frozen affiliate rebates become transferable after expiry. --- main.go | 2 + service/rebate_release_task.go | 85 ++++++++++++++++++++++++++++++++++ service/rebate_test.go | 27 +++++++++++ 3 files changed, 114 insertions(+) create mode 100644 service/rebate_release_task.go diff --git a/main.go b/main.go index 91a6f45c..f1892649 100644 --- a/main.go +++ b/main.go @@ -123,6 +123,8 @@ func main() { // Channel monitor runner and maintenance tasks service.StartChannelMonitorRunner() service.StartChannelMonitorMaintenanceTask() + // Affiliate rebate release task + service.StartRebateReleaseTask() // Wire task polling adaptor factory (breaks service -> relay import cycle) service.GetTaskAdaptorFunc = func(platform constant.TaskPlatform) service.TaskPollingAdaptor { diff --git a/service/rebate_release_task.go b/service/rebate_release_task.go new file mode 100644 index 00000000..6c2daad3 --- /dev/null +++ b/service/rebate_release_task.go @@ -0,0 +1,85 @@ +package service + +import ( + "log" + "sync" + "time" +) + +var ( + rebateReleaseOnce sync.Once + rebateReleaseStopMu sync.Mutex + rebateReleaseStopCh chan struct{} + rebateReleaseStopped bool + + rebateReleaseRunMu sync.Mutex + rebateReleaseActive bool +) + +func StartRebateReleaseTask() { + rebateReleaseOnce.Do(func() { + rebateReleaseStopMu.Lock() + defer rebateReleaseStopMu.Unlock() + + rebateReleaseStopCh = make(chan struct{}) + rebateReleaseStopped = false + go runRebateReleaseLoop(rebateReleaseStopCh) + }) +} + +func StopRebateReleaseTask() { + rebateReleaseStopMu.Lock() + defer rebateReleaseStopMu.Unlock() + + if rebateReleaseStopCh != nil && !rebateReleaseStopped { + close(rebateReleaseStopCh) + rebateReleaseStopped = true + } +} + +func RunRebateReleaseOnce() error { + if !markRebateReleaseRunning() { + return nil + } + defer unmarkRebateReleaseRunning() + + return ReleaseExpiredRebates() +} + +func runRebateReleaseLoop(stopCh <-chan struct{}) { + if err := RunRebateReleaseOnce(); err != nil { + log.Printf("rebate release: %v", err) + } + + ticker := time.NewTicker(time.Hour) + defer ticker.Stop() + + for { + select { + case <-ticker.C: + if err := RunRebateReleaseOnce(); err != nil { + log.Printf("rebate release: %v", err) + } + case <-stopCh: + return + } + } +} + +func markRebateReleaseRunning() bool { + rebateReleaseRunMu.Lock() + defer rebateReleaseRunMu.Unlock() + + if rebateReleaseActive { + return false + } + rebateReleaseActive = true + return true +} + +func unmarkRebateReleaseRunning() { + rebateReleaseRunMu.Lock() + defer rebateReleaseRunMu.Unlock() + + rebateReleaseActive = false +} diff --git a/service/rebate_test.go b/service/rebate_test.go index 1575daa3..f486e1c9 100644 --- a/service/rebate_test.go +++ b/service/rebate_test.go @@ -234,3 +234,30 @@ func TestReleaseExpiredRebatesMarksOnlyDueFrozenRecords(t *testing.T) { require.Equal(t, 50, reloadedInviter.AffQuota) require.Equal(t, 50, reloadedInviter.AffHistoryQuota) } + +func TestRunRebateReleaseOnceReleasesDueFrozenRecords(t *testing.T) { + db := setupRebateTestDB(t) + inviter, invitee := createRebateUsers(t, db) + past := time.Now().Add(-time.Hour) + require.NoError(t, db.Create(&model.RebateRecord{ + InviterId: inviter.Id, + InviteeId: invitee.Id, + OrderId: 1007, + OrderType: "topup", + OrderAmount: 500, + RebateAmount: 50, + RatePercent: 10, + Status: "frozen", + FrozenUntil: &past, + }).Error) + + require.NoError(t, RunRebateReleaseOnce()) + + var record model.RebateRecord + require.NoError(t, db.Where("order_id = ?", 1007).First(&record).Error) + require.Equal(t, "released", record.Status) + require.NotNil(t, record.ReleasedAt) + + reloadedInviter := loadRebateUser(t, db, inviter.Id) + require.Equal(t, 50, reloadedInviter.AffQuota) +}