-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathpool.go
More file actions
191 lines (165 loc) · 3.67 KB
/
Copy pathpool.go
File metadata and controls
191 lines (165 loc) · 3.67 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
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
package easy_pool
import (
"context"
"errors"
"fmt"
"github.com/google/uuid"
"io"
"sync"
"sync/atomic"
"time"
)
type Pool interface {
Take() (Conn, error)
Put(conn Conn) error
Capacity() int
io.Closer
fmt.Stringer
}
type poolRoutines struct {
ch chan struct{}
cap int
once sync.Once
}
func (r *poolRoutines) Take() (Conn, error) {
<-r.ch
return emptyConn, nil
}
func (r *poolRoutines) Put(conn Conn) error {
r.ch <- struct{}{}
return nil
}
func (r *poolRoutines) Capacity() int {
return r.cap
}
func (r *poolRoutines) Close() error {
r.once.Do(func() {
close(r.ch)
})
return nil
}
func (r *poolRoutines) String() string {
return "pool.routines"
}
// NewRoutines is create a goroutines pool
// it's use to limit goroutines
func NewRoutines(capacity int) *poolRoutines {
c := poolRoutines{cap: capacity, ch: make(chan struct{}, capacity)}
for i := 0; i < c.cap; i++ {
c.ch <- struct{}{}
}
return &c
}
type ConnectionCreator interface {
Create() (Conn, error)
}
type ConnectionChecker interface {
Check(conn Conn) bool
}
var ErrorTakeTimeout = errors.New("get connection from pool timeout")
var ErrorPoolClosed = errors.New("pool is closed")
var ErrorConNil = errors.New("connection is nil")
var emptyConn = Conn{}
type poolConn struct {
ch chan Conn
cap int
factory ConnectionCreator
checker ConnectionChecker
max int64
len int64
closed int32
}
// NewPoolConn will create a connection pool
// capacity is pool capacity
// limit is max connection, it must bigger or equal capacity , default limit = 2 x capacity
// if pool take empty and not put yet,
// factory will build new connection (not bigger then limit),
// when connection put to pool, if pool full ,the connection will close.
func NewPoolConn(capacity int, limit int, factory ConnectionCreator, checker ConnectionChecker) *poolConn {
if limit < capacity {
limit = capacity * 2
}
p := poolConn{cap: capacity, factory: factory, checker: checker, ch: make(chan Conn, capacity), max: int64(limit), len: 0, closed: 0}
return &p
}
func (r *poolConn) Len() int64 {
return atomic.LoadInt64(&r.len)
}
func (p *poolConn) Take() (Conn, error) {
if atomic.LoadInt32(&p.closed) == 1 {
return emptyConn, ErrorPoolClosed
}
ctx, cancel := context.WithTimeout(context.Background(), time.Second*2)
defer cancel()
loop:
for {
select {
case conn := <-p.ch:
if !p.checker.Check(conn) {
return createConnWithFactory(p.factory)
}
return conn, nil
case <-ctx.Done():
break loop
default:
length := atomic.LoadInt64(&p.len)
if length < p.max {
atomic.AddInt64(&p.len, 1)
return createConnWithFactory(p.factory)
}
time.Sleep(time.Microsecond * 10)
}
}
return emptyConn, ErrorTakeTimeout
}
func (p *poolConn) Put(conn Conn) error {
if atomic.LoadInt32(&p.closed) == 1 {
defer conn.Close()
return ErrorPoolClosed
}
if conn == emptyConn {
return ErrorConNil
}
select {
case p.ch <- conn:
return nil
default:
atomic.AddInt64(&p.len, -1)
return conn.Close()
}
}
func (p *poolConn) Capacity() int {
return p.cap
}
func (p *poolConn) Close() error {
if atomic.CompareAndSwapInt32(&p.closed, 0, 1) {
var err error
for i := 0; i < p.cap; i++ {
select {
case conn := <-p.ch:
err = conn.Close()
default:
}
}
close(p.ch)
return err
}
return nil
}
func (r *poolConn) String() string {
return "pool.connection_pool"
}
func createConnWithFactory(factory ConnectionCreator) (Conn, error) {
conn, err := factory.Create()
if err != nil {
return emptyConn, err
}
n := time.Now()
conn.sec = n.Unix()
conn.nsec = uint32(n.Nanosecond())
conn.UUID = uuid.New().String()
return conn, nil
}
func EmptyConn() Conn {
return emptyConn
}