-
Notifications
You must be signed in to change notification settings - Fork 6
/
service.go
162 lines (145 loc) Β· 3.04 KB
/
service.go
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
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
package orbital
import (
"context"
"fmt"
"io"
"os"
"sync"
"time"
"github.com/segmentio/stats"
)
// Service runs all registered TestCases on the schedule specified during
// registration.
type Service struct {
// list of tests to run
tests []TestCase
stats *stats.Engine
mu sync.Mutex
started bool
defaultTimeout time.Duration
w io.Writer
done chan struct{}
once sync.Once
wg sync.WaitGroup
}
func WithStats(s *stats.Engine) func(*Service) {
return func(svc *Service) {
svc.stats = s
}
}
func WithTimeout(d time.Duration) func(*Service) {
return func(svc *Service) {
svc.defaultTimeout = d
}
}
func New(opts ...func(*Service)) *Service {
s := &Service{
w: os.Stderr,
done: make(chan struct{}),
tests: make([]TestCase, 0),
defaultTimeout: 10 * time.Minute,
}
for _, o := range opts {
o(s)
}
return s
}
var DefaultService = New()
func Register(tc TestCase) {
DefaultService.Register(tc)
}
// Register a test case to be run.
func (s *Service) Register(tc TestCase) {
s.mu.Lock()
defer s.mu.Unlock()
s.tests = append(s.tests, tc)
}
func (s *Service) Run() {
s.mu.Lock()
defer s.mu.Unlock()
if s.stats == nil {
s.stats = stats.DefaultEngine
}
if !s.started {
for _, v := range s.tests {
s.wg.Add(1)
go s.run(v)
}
}
s.started = true
}
func (s *Service) handle(ctx context.Context, tc TestCase) {
start := time.Now()
o := &O{
w: s.w,
stats: s.stats,
}
to := s.defaultTimeout
if tc.Timeout > 10*time.Millisecond {
to = tc.Timeout
}
c, cancel := context.WithTimeout(ctx, to)
defer cancel()
tc.Func(c, o)
if c.Err() != nil && !o.failed {
o.Errorf("failed on context error: %v", c.Err())
}
dur := time.Now().Sub(start)
if o.failed {
tags := append([]stats.Tag{
stats.T("case", tc.Name),
stats.T("result", "fail"),
}, tc.Tags...)
s.stats.Observe("case", dur, tags...)
fmt.Fprintf(s.w, "--- FAIL: %s (%s)\n", tc.Name, dur)
} else {
tags := append([]stats.Tag{
stats.T("case", tc.Name),
stats.T("result", "pass"),
}, tc.Tags...)
s.stats.Observe("case", dur, tags...)
fmt.Fprintf(s.w, "--- PASS: %s (%s)\n", tc.Name, dur)
}
}
func (s *Service) run(tc TestCase) {
tick := time.NewTicker(tc.Period)
// Waitgroup for different invocations of this test case
var wg sync.WaitGroup
loop:
for {
select {
case <-tick.C:
case <-s.done:
tick.Stop()
break loop
}
ctx, cancel := context.WithCancel(context.Background())
complete := make(chan struct{})
wg.Add(1)
go func(c context.Context) {
s.handle(c, tc)
close(complete)
wg.Done()
}(ctx)
// Cancel the above goroutine on shutdown without blocking the loop
go func(comp chan struct{}, c context.CancelFunc) {
select {
case <-s.done:
// TODO: make this configurable at the service level. i.e.
// provide the option to allow the inflight tests to finish
// without failing
c()
case <-comp:
}
}(complete, cancel)
}
wg.Wait()
s.wg.Done()
}
func (s *Service) Close() error {
s.once.Do(func() {
close(s.done)
})
s.wg.Wait()
return nil
}