Skip to content

Commit 21f6fb6

Browse files
committed
sm: adjust task
1 parent babfa8b commit 21f6fb6

2 files changed

Lines changed: 19 additions & 9 deletions

File tree

taskworker/main.go

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -140,7 +140,7 @@ func initTasks(ctx context.Context) {
140140
work.ScheduleProcessPendingPayouts(clientSession, tx)
141141
// work.SchedulePopulateAccountWallets(clientSession, tx)
142142
for i := range work.DefaultCloseExpiredContractsBlockSize {
143-
work.ScheduleCloseExpiredContracts(clientSession, tx, i)
143+
work.ScheduleCloseExpiredContracts(clientSession, tx, i, false)
144144
}
145145
work.ScheduleCloseExpiredNetworkClientHandlers(clientSession, tx)
146146
work.ScheduleRemoveDisconnectedNetworkClients(clientSession, tx)

taskworker/work/subscription_work.go

Lines changed: 18 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -11,17 +11,18 @@ import (
1111
"github.com/urnetwork/server/task"
1212
)
1313

14-
const DefaultCloseExpiredContractsBlockSize = 32
14+
const DefaultCloseExpiredContractsBlockSize = 24
1515

1616
type CloseExpiredContractsArgs struct {
1717
BlockSize int `json:"block_size"`
1818
BlockIndex int `json:"block_index"`
1919
}
2020

2121
type CloseExpiredContractsResult struct {
22+
Full bool `json:"full"`
2223
}
2324

24-
func ScheduleCloseExpiredContracts(clientSession *session.ClientSession, tx server.PgTx, blockIndex int) {
25+
func ScheduleCloseExpiredContracts(clientSession *session.ClientSession, tx server.PgTx, blockIndex int, delay bool) {
2526
// runAt := func() time.Time {
2627
// now := server.NowUtc()
2728
// year, month, day := now.Date()
@@ -32,6 +33,11 @@ func ScheduleCloseExpiredContracts(clientSession *session.ClientSession, tx serv
3233
blockSize := DefaultCloseExpiredContractsBlockSize
3334
blockIndex = blockIndex % blockSize
3435

36+
runAt := server.NowUtc()
37+
if delay {
38+
runAt = runAt.Add(time.Minute)
39+
}
40+
3541
task.ScheduleTaskInTx(
3642
tx,
3743
CloseExpiredContracts,
@@ -42,7 +48,7 @@ func ScheduleCloseExpiredContracts(clientSession *session.ClientSession, tx serv
4248
clientSession,
4349
// legacy key
4450
task.RunOnce(fmt.Sprintf("close_expired_contracts_%d_%d", blockSize, blockIndex)),
45-
task.RunAt(server.NowUtc().Add(time.Minute)),
51+
task.RunAt(runAt),
4652
task.MaxTime(30*time.Minute),
4753
task.Priority(task.TaskPriorityFastest),
4854
)
@@ -54,15 +60,19 @@ func CloseExpiredContracts(
5460
) (*CloseExpiredContractsResult, error) {
5561
if closeExpiredContracts.BlockSize == DefaultCloseExpiredContractsBlockSize {
5662
minTime := server.NowUtc().Add(-5 * time.Minute)
57-
_, err := model.ForceCloseOpenContractIds(
63+
n := 100000
64+
c, err := model.ForceCloseOpenContractIds(
5865
clientSession.Ctx,
5966
minTime,
60-
1000000,
61-
48,
67+
n,
68+
92,
6269
closeExpiredContracts.BlockSize,
6370
closeExpiredContracts.BlockIndex,
6471
)
65-
return &CloseExpiredContractsResult{}, err
72+
full := int64(n/(4*DefaultCloseExpiredContractsBlockSize)) <= c
73+
return &CloseExpiredContractsResult{
74+
Full: full,
75+
}, err
6676
}
6777
// else ignore lingering tasks with older block size
6878
return &CloseExpiredContractsResult{}, nil
@@ -74,7 +84,7 @@ func CloseExpiredContractsPost(
7484
clientSession *session.ClientSession,
7585
tx server.PgTx,
7686
) error {
77-
ScheduleCloseExpiredContracts(clientSession, tx, closeExpiredContracts.BlockIndex)
87+
ScheduleCloseExpiredContracts(clientSession, tx, closeExpiredContracts.BlockIndex, !closeExpiredContractsResult.Full)
7888
return nil
7989
}
8090

0 commit comments

Comments
 (0)