-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathbqin.go
More file actions
139 lines (117 loc) · 2.7 KB
/
Copy pathbqin.go
File metadata and controls
139 lines (117 loc) · 2.7 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
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
package bqin
import (
"context"
"errors"
"github.com/kayac/bqin/internal/logger"
)
type App struct {
*Receiver
*Resolver
*Transporter
*Loader
}
func NewApp(conf *Config) *App {
factory := &Factory{Config: conf}
return factory.NewApp()
}
func (app *App) Run(ctx context.Context, opts ...RunOption) error {
logger.Infof("Starting up bqin worker")
defer logger.Infof("Shutdown bqin worker")
settings := &RunSettings{}
for _, opt := range opts {
opt.Apply(settings)
}
if settings.QueueName != "" {
defaultQueueName := app.GetQueueName()
app.SetQueueName(settings.QueueName)
defer app.SetQueueName(defaultQueueName)
}
for {
select {
case <-ctx.Done():
if err := ctx.Err(); err != context.Canceled {
return err
}
logger.Infof("canceled")
return nil
default:
}
switch err := app.batch(context.Background()); err {
case ErrNoMessage:
if settings.ExitNoMessage {
logger.Infof("success all")
return nil
}
case nil:
//nothing todo
default:
if settings.ExitError {
return err
}
logger.Errorf("process failed. reason:%s", err)
}
}
}
func (app *App) batch(ctx context.Context) error {
urls, receiptHandle, err := app.Receive(ctx)
defer receiptHandle.Cleanup()
if err != nil {
return err
}
jobs := app.Resolve(urls)
if len(jobs) == 0 {
return errors.New("nothing to do")
}
transportHandles := make([]*TransportJobHandle, 0, len(jobs))
defer func() {
for _, h := range transportHandles {
h.Cleanup(ctx)
}
}()
for i, job := range jobs {
receiptHandle.Infof("[job %02d]%s", i, job)
transportHandle, err := app.Transport(ctx, job.TransportJob)
if err != nil {
return err
}
transportHandles = append(transportHandles, transportHandle)
if err := app.Load(ctx, job.LoadingJob); err != nil {
return err
}
receiptHandle.Infof("[job %02d]complte job", i)
}
return receiptHandle.Complete()
}
type RunOption interface {
Apply(*RunSettings)
}
type RunSettings struct {
ExitNoMessage bool
ExitError bool
QueueName string
}
func (s *RunSettings) Apply(o *RunSettings) {
o.ExitNoMessage = s.ExitNoMessage
o.ExitError = s.ExitError
}
type withExitNoMessage bool
func (opt withExitNoMessage) Apply(settings *RunSettings) {
settings.ExitNoMessage = bool(opt)
}
func WithExitNoMessage(flag bool) RunOption {
return withExitNoMessage(flag)
}
type withExitError bool
func (opt withExitError) Apply(settings *RunSettings) {
settings.ExitError = bool(opt)
}
func WithExitError(flag bool) RunOption {
return withExitError(flag)
}
type withQueueName string
func (opt withQueueName) Apply(settings *RunSettings) {
settings.QueueName = string(opt)
}
func WithQueueName(queueName string) RunOption {
return withQueueName(queueName)
}