This repository was archived by the owner on Apr 9, 2020. It is now read-only.
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathqueue.go
More file actions
112 lines (91 loc) · 2 KB
/
Copy pathqueue.go
File metadata and controls
112 lines (91 loc) · 2 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
package serverz
import (
"context"
"net"
"sync"
)
type queueItem struct {
server Server
addr net.Addr
}
// Queue holds a list of servers and starts them at once.
type Queue struct {
servers []*queueItem
}
// NewQueue returns a new Queue.
func NewQueue() *Queue {
return new(Queue)
}
// Append appends a new server to the list of existing ones.
func (q *Queue) Append(server Server, addr net.Addr) {
q.servers = append(
q.servers,
&queueItem{
server,
addr,
},
)
}
// Prepend prepends a new server to the list of existing ones.
func (q *Queue) Prepend(server Server, addr net.Addr) {
q.servers = append(
[]*queueItem{
&queueItem{
server,
addr,
},
},
q.servers...,
)
}
// Start starts all the servers.
func (q *Queue) Start() <-chan error {
ch := make(chan error, len(q.servers))
for _, s := range q.servers {
if ls, ok := s.server.(listenServer); ok {
go func(ch chan<- error, server listenServer, addr net.Addr) {
ch <- server.ListenAndServe(addr)
}(ch, ls, s.addr)
} else {
lis, err := listen(s.addr)
if err != nil {
ch <- err
// Skip starting this server if we can't listen on the interface
// Consuming to the error channel should end up being the application terminated anyway
continue
}
go func(ch chan<- error, server Server, lis net.Listener) {
ch <- server.Serve(lis)
}(ch, s.server, lis)
}
}
return ch
}
// Shutdown tries to gracefully stop all the servers.
func (q *Queue) Shutdown(ctx context.Context) error {
wg := &sync.WaitGroup{}
merr := multiError{}
for _, s := range q.servers {
wg.Add(1)
go func(server Server) {
err := server.Shutdown(ctx)
if err != nil {
merr = append(merr, err)
}
wg.Done()
}(s.server)
}
wg.Wait()
return merr.ErrOrNil()
}
// Close immediately calls close for all servers.
func (q *Queue) Close() error {
merr := multiError{}
for _, s := range q.servers {
err := s.server.Close()
if err != nil {
merr = append(merr, err)
}
}
return merr.ErrOrNil()
}