...
|
...
|
@@ -24,13 +24,11 @@ import ( |
|
|
"sync"
|
|
|
"sync/atomic"
|
|
|
"time"
|
|
|
|
|
|
"github.com/garyburd/redigo/internal"
|
|
|
)
|
|
|
|
|
|
var (
|
|
|
_ ConnWithTimeout = (*pooledConnection)(nil)
|
|
|
_ ConnWithTimeout = (*errorConnection)(nil)
|
|
|
_ ConnWithTimeout = (*activeConn)(nil)
|
|
|
_ ConnWithTimeout = (*errorConn)(nil)
|
|
|
)
|
|
|
|
|
|
var nowFunc = time.Now // for testing
|
...
|
...
|
@@ -150,6 +148,10 @@ type Pool struct { |
|
|
// for a connection to be returned to the pool before returning.
|
|
|
Wait bool
|
|
|
|
|
|
// Close connections older than this duration. If the value is zero, then
|
|
|
// the pool does not close connections based on age.
|
|
|
MaxConnLifetime time.Duration
|
|
|
|
|
|
chInitialized uint32 // set to 1 when field ch is initialized
|
|
|
|
|
|
mu sync.Mutex // mu protects the following fields
|
...
|
...
|
@@ -172,11 +174,11 @@ func NewPool(newFn func() (Conn, error), maxIdle int) *Pool { |
|
|
// getting an underlying connection, then the connection Err, Do, Send, Flush
|
|
|
// and Receive methods return that error.
|
|
|
func (p *Pool) Get() Conn {
|
|
|
c, err := p.get(nil)
|
|
|
pc, err := p.get(nil)
|
|
|
if err != nil {
|
|
|
return errorConnection{err}
|
|
|
return errorConn{err}
|
|
|
}
|
|
|
return &pooledConnection{p: p, c: c}
|
|
|
return &activeConn{p: p, pc: pc}
|
|
|
}
|
|
|
|
|
|
// PoolStats contains pool statistics.
|
...
|
...
|
@@ -226,15 +228,15 @@ func (p *Pool) Close() error { |
|
|
}
|
|
|
p.closed = true
|
|
|
p.active -= p.idle.count
|
|
|
ic := p.idle.front
|
|
|
pc := p.idle.front
|
|
|
p.idle.count = 0
|
|
|
p.idle.front, p.idle.back = nil, nil
|
|
|
if p.ch != nil {
|
|
|
close(p.ch)
|
|
|
}
|
|
|
p.mu.Unlock()
|
|
|
for ; ic != nil; ic = ic.next {
|
|
|
ic.c.Close()
|
|
|
for ; pc != nil; pc = pc.next {
|
|
|
pc.c.Close()
|
|
|
}
|
|
|
return nil
|
|
|
}
|
...
|
...
|
@@ -265,7 +267,7 @@ func (p *Pool) lazyInit() { |
|
|
func (p *Pool) get(ctx interface {
|
|
|
Done() <-chan struct{}
|
|
|
Err() error
|
|
|
}) (Conn, error) {
|
|
|
}) (*poolConn, error) {
|
|
|
|
|
|
// Handle limit for p.Wait == true.
|
|
|
if p.Wait && p.MaxActive > 0 {
|
...
|
...
|
@@ -287,10 +289,10 @@ func (p *Pool) get(ctx interface { |
|
|
if p.IdleTimeout > 0 {
|
|
|
n := p.idle.count
|
|
|
for i := 0; i < n && p.idle.back != nil && p.idle.back.t.Add(p.IdleTimeout).Before(nowFunc()); i++ {
|
|
|
c := p.idle.back.c
|
|
|
pc := p.idle.back
|
|
|
p.idle.popBack()
|
|
|
p.mu.Unlock()
|
|
|
c.Close()
|
|
|
pc.c.Close()
|
|
|
p.mu.Lock()
|
|
|
p.active--
|
|
|
}
|
...
|
...
|
@@ -298,13 +300,14 @@ func (p *Pool) get(ctx interface { |
|
|
|
|
|
// Get idle connection from the front of idle list.
|
|
|
for p.idle.front != nil {
|
|
|
ic := p.idle.front
|
|
|
pc := p.idle.front
|
|
|
p.idle.popFront()
|
|
|
p.mu.Unlock()
|
|
|
if p.TestOnBorrow == nil || p.TestOnBorrow(ic.c, ic.t) == nil {
|
|
|
return ic.c, nil
|
|
|
if (p.TestOnBorrow == nil || p.TestOnBorrow(pc.c, pc.t) == nil) &&
|
|
|
(p.MaxConnLifetime == 0 || nowFunc().Sub(pc.created) < p.MaxConnLifetime) {
|
|
|
return pc, nil
|
|
|
}
|
|
|
ic.c.Close()
|
|
|
pc.c.Close()
|
|
|
p.mu.Lock()
|
|
|
p.active--
|
|
|
}
|
...
|
...
|
@@ -333,24 +336,25 @@ func (p *Pool) get(ctx interface { |
|
|
}
|
|
|
p.mu.Unlock()
|
|
|
}
|
|
|
return c, err
|
|
|
return &poolConn{c: c, created: nowFunc()}, err
|
|
|
}
|
|
|
|
|
|
func (p *Pool) put(c Conn, forceClose bool) error {
|
|
|
func (p *Pool) put(pc *poolConn, forceClose bool) error {
|
|
|
p.mu.Lock()
|
|
|
if !p.closed && !forceClose {
|
|
|
p.idle.pushFront(&idleConn{t: nowFunc(), c: c})
|
|
|
pc.t = nowFunc()
|
|
|
p.idle.pushFront(pc)
|
|
|
if p.idle.count > p.MaxIdle {
|
|
|
c = p.idle.back.c
|
|
|
pc = p.idle.back
|
|
|
p.idle.popBack()
|
|
|
} else {
|
|
|
c = nil
|
|
|
pc = nil
|
|
|
}
|
|
|
}
|
|
|
|
|
|
if c != nil {
|
|
|
if pc != nil {
|
|
|
p.mu.Unlock()
|
|
|
c.Close()
|
|
|
pc.c.Close()
|
|
|
p.mu.Lock()
|
|
|
p.active--
|
|
|
}
|
...
|
...
|
@@ -362,9 +366,9 @@ func (p *Pool) put(c Conn, forceClose bool) error { |
|
|
return nil
|
|
|
}
|
|
|
|
|
|
type pooledConnection struct {
|
|
|
type activeConn struct {
|
|
|
p *Pool
|
|
|
c Conn
|
|
|
pc *poolConn
|
|
|
state int
|
|
|
}
|
|
|
|
...
|
...
|
@@ -385,79 +389,107 @@ func initSentinel() { |
|
|
}
|
|
|
}
|
|
|
|
|
|
func (pc *pooledConnection) Close() error {
|
|
|
c := pc.c
|
|
|
if _, ok := c.(errorConnection); ok {
|
|
|
func (ac *activeConn) Close() error {
|
|
|
pc := ac.pc
|
|
|
if pc == nil {
|
|
|
return nil
|
|
|
}
|
|
|
pc.c = errorConnection{errConnClosed}
|
|
|
|
|
|
if pc.state&internal.MultiState != 0 {
|
|
|
c.Send("DISCARD")
|
|
|
pc.state &^= (internal.MultiState | internal.WatchState)
|
|
|
} else if pc.state&internal.WatchState != 0 {
|
|
|
c.Send("UNWATCH")
|
|
|
pc.state &^= internal.WatchState
|
|
|
ac.pc = nil
|
|
|
|
|
|
if ac.state&connectionMultiState != 0 {
|
|
|
pc.c.Send("DISCARD")
|
|
|
ac.state &^= (connectionMultiState | connectionWatchState)
|
|
|
} else if ac.state&connectionWatchState != 0 {
|
|
|
pc.c.Send("UNWATCH")
|
|
|
ac.state &^= connectionWatchState
|
|
|
}
|
|
|
if pc.state&internal.SubscribeState != 0 {
|
|
|
c.Send("UNSUBSCRIBE")
|
|
|
c.Send("PUNSUBSCRIBE")
|
|
|
if ac.state&connectionSubscribeState != 0 {
|
|
|
pc.c.Send("UNSUBSCRIBE")
|
|
|
pc.c.Send("PUNSUBSCRIBE")
|
|
|
// To detect the end of the message stream, ask the server to echo
|
|
|
// a sentinel value and read until we see that value.
|
|
|
sentinelOnce.Do(initSentinel)
|
|
|
c.Send("ECHO", sentinel)
|
|
|
c.Flush()
|
|
|
pc.c.Send("ECHO", sentinel)
|
|
|
pc.c.Flush()
|
|
|
for {
|
|
|
p, err := c.Receive()
|
|
|
p, err := pc.c.Receive()
|
|
|
if err != nil {
|
|
|
break
|
|
|
}
|
|
|
if p, ok := p.([]byte); ok && bytes.Equal(p, sentinel) {
|
|
|
pc.state &^= internal.SubscribeState
|
|
|
ac.state &^= connectionSubscribeState
|
|
|
break
|
|
|
}
|
|
|
}
|
|
|
}
|
|
|
c.Do("")
|
|
|
pc.p.put(c, pc.state != 0 || c.Err() != nil)
|
|
|
pc.c.Do("")
|
|
|
ac.p.put(pc, ac.state != 0 || pc.c.Err() != nil)
|
|
|
return nil
|
|
|
}
|
|
|
|
|
|
func (pc *pooledConnection) Err() error {
|
|
|
func (ac *activeConn) Err() error {
|
|
|
pc := ac.pc
|
|
|
if pc == nil {
|
|
|
return errConnClosed
|
|
|
}
|
|
|
return pc.c.Err()
|
|
|
}
|
|
|
|
|
|
func (pc *pooledConnection) Do(commandName string, args ...interface{}) (reply interface{}, err error) {
|
|
|
ci := internal.LookupCommandInfo(commandName)
|
|
|
pc.state = (pc.state | ci.Set) &^ ci.Clear
|
|
|
func (ac *activeConn) Do(commandName string, args ...interface{}) (reply interface{}, err error) {
|
|
|
pc := ac.pc
|
|
|
if pc == nil {
|
|
|
return nil, errConnClosed
|
|
|
}
|
|
|
ci := lookupCommandInfo(commandName)
|
|
|
ac.state = (ac.state | ci.Set) &^ ci.Clear
|
|
|
return pc.c.Do(commandName, args...)
|
|
|
}
|
|
|
|
|
|
func (pc *pooledConnection) DoWithTimeout(timeout time.Duration, commandName string, args ...interface{}) (reply interface{}, err error) {
|
|
|
func (ac *activeConn) DoWithTimeout(timeout time.Duration, commandName string, args ...interface{}) (reply interface{}, err error) {
|
|
|
pc := ac.pc
|
|
|
if pc == nil {
|
|
|
return nil, errConnClosed
|
|
|
}
|
|
|
cwt, ok := pc.c.(ConnWithTimeout)
|
|
|
if !ok {
|
|
|
return nil, errTimeoutNotSupported
|
|
|
}
|
|
|
ci := internal.LookupCommandInfo(commandName)
|
|
|
pc.state = (pc.state | ci.Set) &^ ci.Clear
|
|
|
ci := lookupCommandInfo(commandName)
|
|
|
ac.state = (ac.state | ci.Set) &^ ci.Clear
|
|
|
return cwt.DoWithTimeout(timeout, commandName, args...)
|
|
|
}
|
|
|
|
|
|
func (pc *pooledConnection) Send(commandName string, args ...interface{}) error {
|
|
|
ci := internal.LookupCommandInfo(commandName)
|
|
|
pc.state = (pc.state | ci.Set) &^ ci.Clear
|
|
|
func (ac *activeConn) Send(commandName string, args ...interface{}) error {
|
|
|
pc := ac.pc
|
|
|
if pc == nil {
|
|
|
return errConnClosed
|
|
|
}
|
|
|
ci := lookupCommandInfo(commandName)
|
|
|
ac.state = (ac.state | ci.Set) &^ ci.Clear
|
|
|
return pc.c.Send(commandName, args...)
|
|
|
}
|
|
|
|
|
|
func (pc *pooledConnection) Flush() error {
|
|
|
func (ac *activeConn) Flush() error {
|
|
|
pc := ac.pc
|
|
|
if pc == nil {
|
|
|
return errConnClosed
|
|
|
}
|
|
|
return pc.c.Flush()
|
|
|
}
|
|
|
|
|
|
func (pc *pooledConnection) Receive() (reply interface{}, err error) {
|
|
|
func (ac *activeConn) Receive() (reply interface{}, err error) {
|
|
|
pc := ac.pc
|
|
|
if pc == nil {
|
|
|
return nil, errConnClosed
|
|
|
}
|
|
|
return pc.c.Receive()
|
|
|
}
|
|
|
|
|
|
func (pc *pooledConnection) ReceiveWithTimeout(timeout time.Duration) (reply interface{}, err error) {
|
|
|
func (ac *activeConn) ReceiveWithTimeout(timeout time.Duration) (reply interface{}, err error) {
|
|
|
pc := ac.pc
|
|
|
if pc == nil {
|
|
|
return nil, errConnClosed
|
|
|
}
|
|
|
cwt, ok := pc.c.(ConnWithTimeout)
|
|
|
if !ok {
|
|
|
return nil, errTimeoutNotSupported
|
...
|
...
|
@@ -465,63 +497,64 @@ func (pc *pooledConnection) ReceiveWithTimeout(timeout time.Duration) (reply int |
|
|
return cwt.ReceiveWithTimeout(timeout)
|
|
|
}
|
|
|
|
|
|
type errorConnection struct{ err error }
|
|
|
type errorConn struct{ err error }
|
|
|
|
|
|
func (ec errorConnection) Do(string, ...interface{}) (interface{}, error) { return nil, ec.err }
|
|
|
func (ec errorConnection) DoWithTimeout(time.Duration, string, ...interface{}) (interface{}, error) {
|
|
|
func (ec errorConn) Do(string, ...interface{}) (interface{}, error) { return nil, ec.err }
|
|
|
func (ec errorConn) DoWithTimeout(time.Duration, string, ...interface{}) (interface{}, error) {
|
|
|
return nil, ec.err
|
|
|
}
|
|
|
func (ec errorConnection) Send(string, ...interface{}) error { return ec.err }
|
|
|
func (ec errorConnection) Err() error { return ec.err }
|
|
|
func (ec errorConnection) Close() error { return nil }
|
|
|
func (ec errorConnection) Flush() error { return ec.err }
|
|
|
func (ec errorConnection) Receive() (interface{}, error) { return nil, ec.err }
|
|
|
func (ec errorConnection) ReceiveWithTimeout(time.Duration) (interface{}, error) { return nil, ec.err }
|
|
|
func (ec errorConn) Send(string, ...interface{}) error { return ec.err }
|
|
|
func (ec errorConn) Err() error { return ec.err }
|
|
|
func (ec errorConn) Close() error { return nil }
|
|
|
func (ec errorConn) Flush() error { return ec.err }
|
|
|
func (ec errorConn) Receive() (interface{}, error) { return nil, ec.err }
|
|
|
func (ec errorConn) ReceiveWithTimeout(time.Duration) (interface{}, error) { return nil, ec.err }
|
|
|
|
|
|
type idleList struct {
|
|
|
count int
|
|
|
front, back *idleConn
|
|
|
front, back *poolConn
|
|
|
}
|
|
|
|
|
|
type idleConn struct {
|
|
|
type poolConn struct {
|
|
|
c Conn
|
|
|
t time.Time
|
|
|
next, prev *idleConn
|
|
|
created time.Time
|
|
|
next, prev *poolConn
|
|
|
}
|
|
|
|
|
|
func (l *idleList) pushFront(ic *idleConn) {
|
|
|
ic.next = l.front
|
|
|
ic.prev = nil
|
|
|
func (l *idleList) pushFront(pc *poolConn) {
|
|
|
pc.next = l.front
|
|
|
pc.prev = nil
|
|
|
if l.count == 0 {
|
|
|
l.back = ic
|
|
|
l.back = pc
|
|
|
} else {
|
|
|
l.front.prev = ic
|
|
|
l.front.prev = pc
|
|
|
}
|
|
|
l.front = ic
|
|
|
l.front = pc
|
|
|
l.count++
|
|
|
return
|
|
|
}
|
|
|
|
|
|
func (l *idleList) popFront() {
|
|
|
ic := l.front
|
|
|
pc := l.front
|
|
|
l.count--
|
|
|
if l.count == 0 {
|
|
|
l.front, l.back = nil, nil
|
|
|
} else {
|
|
|
ic.next.prev = nil
|
|
|
l.front = ic.next
|
|
|
pc.next.prev = nil
|
|
|
l.front = pc.next
|
|
|
}
|
|
|
ic.next, ic.prev = nil, nil
|
|
|
pc.next, pc.prev = nil, nil
|
|
|
}
|
|
|
|
|
|
func (l *idleList) popBack() {
|
|
|
ic := l.back
|
|
|
pc := l.back
|
|
|
l.count--
|
|
|
if l.count == 0 {
|
|
|
l.front, l.back = nil, nil
|
|
|
} else {
|
|
|
ic.prev.next = nil
|
|
|
l.back = ic.prev
|
|
|
pc.prev.next = nil
|
|
|
l.back = pc.prev
|
|
|
}
|
|
|
ic.next, ic.prev = nil, nil
|
|
|
pc.next, pc.prev = nil, nil
|
|
|
} |
...
|
...
|
|