-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathtaskqueue.go
More file actions
107 lines (95 loc) · 1.68 KB
/
taskqueue.go
File metadata and controls
107 lines (95 loc) · 1.68 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
package taskqueue
import (
"errors"
"sync"
"time"
)
// Task is Task
type Task struct {
F func() error
RetryTimes int
}
// Err is tasks fail error
type Err struct {
Task Task
Err error
}
// TaskQueue is taskqueue struct
type TaskQueue struct {
interval time.Duration
sync.RWMutex
tasks []Task
breakFlag bool
closeCh chan int
Error chan Err
}
// New Create new struct
func New(interval time.Duration) *TaskQueue {
return &TaskQueue{
interval: interval,
closeCh: make(chan int),
Error: make(chan Err),
}
}
// Add is add job to task queue
func (t *TaskQueue) Add(f func() error, retryTimes int) error {
if t.breakFlag {
return errors.New("taskQueue.Stop is called")
}
t.Lock()
t.tasks = append(t.tasks, Task{
F: f,
RetryTimes: retryTimes,
})
t.Unlock()
return nil
}
func (t *TaskQueue) addNotCheckBreakFlag(f func() error, retryTimes int) {
t.Lock()
t.tasks = append(t.tasks, Task{
F: f,
RetryTimes: retryTimes,
})
t.Unlock()
}
// Run run taskqueue
func (t *TaskQueue) Run() {
L:
for {
if len(t.tasks) > 0 {
tt := t.pop()
if err := tt.F(); err != nil {
t.retry(tt, err)
}
}
if len(t.tasks) <= 0 && t.breakFlag {
t.closeCh <- 1
break L
}
time.Sleep(t.interval)
}
}
func (t *TaskQueue) retry(tt Task, err error) {
tt.RetryTimes = tt.RetryTimes - 1
if tt.RetryTimes >= 1 {
t.addNotCheckBreakFlag(tt.F, tt.RetryTimes)
}
t.Error <- Err{
Task: tt,
Err: err,
}
}
func (t *TaskQueue) pop() Task {
t.Lock()
task := t.tasks[0]
t.tasks = t.tasks[1:]
t.Unlock()
return task
}
// Stop stop taskqueue
func (t *TaskQueue) Stop() {
t.breakFlag = true
<-t.closeCh
close(t.closeCh)
close(t.Error)
}