grpc_pool/pool.go

197 lines
3.3 KiB
Go
Raw Normal View History

2017-11-09 06:47:07 +00:00
package grpc_pool
2017-09-05 05:42:54 +00:00
import (
"errors"
2017-11-09 05:46:04 +00:00
"fmt"
"sync"
2017-09-05 05:42:54 +00:00
"time"
"google.golang.org/grpc"
)
var (
// ErrFullPool is the error when the pool is already full
ErrPoolFull = errors.New("grpc pool: Pool is reached maximum capacity.")
)
type pooledInstance struct {
clientConn *grpc.ClientConn
idleTime time.Time
instance interface{}
inUse bool
}
type Pool interface {
2017-11-09 05:46:04 +00:00
Start() error
Stop()
2017-09-05 05:42:54 +00:00
Get() (interface{}, error)
Put(i interface{})
2017-11-09 05:46:04 +00:00
2017-09-05 05:42:54 +00:00
Capacity() int
Available() int
}
type pool struct {
2017-11-09 05:46:04 +00:00
ph PoolHandler
instances map[interface{}]*pooledInstance
instanceQueueChan chan *pooledInstance
stopChan chan struct{}
stopWg sync.WaitGroup
2017-09-05 05:42:54 +00:00
}
2017-11-09 07:43:43 +00:00
func New(ph PoolHandler) Pool {
2017-09-05 05:42:54 +00:00
p := &pool{
2017-11-09 05:46:04 +00:00
ph: ph,
}
2017-11-09 07:50:32 +00:00
return p, nil
2017-11-09 05:46:04 +00:00
}
func (p *pool) Start() error {
if nil == p.ph {
panic("GRPC Pool: pool handler must be specified.")
2017-09-05 05:42:54 +00:00
}
2017-11-09 05:46:04 +00:00
p.ph.Validate()
2017-09-05 05:42:54 +00:00
2017-11-09 05:46:04 +00:00
if p.stopChan != nil {
return fmt.Errorf("GRPC Pool: pool is already running. Stop it before starting it again")
}
p.stopChan = make(chan struct{})
p.instances = make(map[interface{}]*pooledInstance, p.ph.GetMaxCapacity())
p.instanceQueueChan = make(chan *pooledInstance, p.ph.GetMaxCapacity())
if 0 < p.ph.GetMaxIdle() {
for i := 0; i < p.ph.GetMaxIdle(); i++ {
if err := p.createInstance(); nil != err {
p.Stop()
return err
2017-09-05 05:42:54 +00:00
}
}
}
2017-11-09 05:46:04 +00:00
return nil
}
func (p *pool) Stop() {
if p.stopChan == nil {
panic("GRPC Pool: pool must be started before stopping it")
}
close(p.instanceQueueChan)
for k, _ := range p.instances {
p.destroyInstance(k)
}
p.stopWg.Wait()
p.stopChan = nil
2017-09-05 05:42:54 +00:00
}
func (p *pool) createInstance() error {
2017-11-09 05:46:04 +00:00
maxCapacity := p.ph.GetMaxCapacity()
2017-09-05 05:42:54 +00:00
currentInstanceCount := len(p.instances)
if currentInstanceCount >= maxCapacity {
return ErrPoolFull
}
2017-11-09 06:47:07 +00:00
conn, i, err := p.ph.Dial()
2017-09-05 05:42:54 +00:00
if err != nil {
return err
}
pi := &pooledInstance{
clientConn: conn,
idleTime: time.Now(),
instance: i,
inUse: false,
}
p.instances[i] = pi
2017-11-09 05:46:04 +00:00
p.instanceQueueChan <- pi
p.stopWg.Add(1)
2017-09-05 05:42:54 +00:00
return nil
}
func (p *pool) destroyInstance(i interface{}) error {
var err error
pi, ok := p.instances[i]
if !ok {
return nil
}
err = pi.clientConn.Close()
if nil != err {
return err
}
delete(p.instances, i)
2017-11-09 05:46:04 +00:00
p.stopWg.Done()
2017-09-05 05:42:54 +00:00
return nil
}
func (p *pool) get() (*pooledInstance, error) {
var pi *pooledInstance
var err error
2017-11-09 05:46:04 +00:00
avail := len(p.instanceQueueChan)
2017-09-05 05:42:54 +00:00
if 0 == avail {
if err = p.createInstance(); nil != err {
return nil, err
}
}
for {
select {
2017-11-09 05:46:04 +00:00
case pi = <-p.instanceQueueChan:
2017-09-05 05:42:54 +00:00
// All good
default:
}
2017-11-09 05:46:04 +00:00
idleTimeout := p.ph.GetIdleTimeout()
if 1 < len(p.instanceQueueChan) && 0 < idleTimeout && pi.idleTime.Add(idleTimeout).Before(time.Now()) {
2017-09-05 05:42:54 +00:00
go p.destroyInstance(pi.instance)
continue
}
break
}
return pi, nil
}
func (p *pool) Get() (interface{}, error) {
var err error
var pi *pooledInstance
if pi, err = p.get(); nil != err {
return nil, err
}
pi.inUse = true
return pi.instance, nil
}
func (p *pool) Put(i interface{}) {
pi, ok := p.instances[i]
if !ok {
return
}
pi.idleTime = time.Now()
pi.inUse = false
select {
2017-11-09 05:46:04 +00:00
case p.instanceQueueChan <- pi:
2017-09-05 05:42:54 +00:00
// All good
default:
// channel is full
return
}
}
func (p *pool) Capacity() int {
2017-11-09 05:46:04 +00:00
return cap(p.instanceQueueChan)
2017-09-05 05:42:54 +00:00
}
func (p *pool) Available() int {
2017-11-09 05:46:04 +00:00
return len(p.instanceQueueChan)
2017-09-05 05:42:54 +00:00
}