forked from JointCloud/JCC-CSScheduler
60 lines
931 B
Go
60 lines
931 B
Go
package executor
|
|
|
|
import (
|
|
"gitlink.org.cn/cloudream/common/pkgs/mq"
|
|
mymq "gitlink.org.cn/cloudream/scheduler/common/pkgs/mq"
|
|
)
|
|
|
|
type Client struct {
|
|
rabbitCli *mq.RabbitMQClient
|
|
}
|
|
|
|
func NewClient(cfg *mymq.Config) (*Client, error) {
|
|
rabbitCli, err := mq.NewRabbitMQClient(cfg.MakeConnectingURL(), ServerQueueName, "")
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return &Client{
|
|
rabbitCli: rabbitCli,
|
|
}, nil
|
|
}
|
|
|
|
func (c *Client) Close() {
|
|
c.rabbitCli.Close()
|
|
}
|
|
|
|
type PoolClient struct {
|
|
*Client
|
|
owner *Pool
|
|
}
|
|
|
|
func (c *PoolClient) Close() {
|
|
c.owner.Release(c)
|
|
}
|
|
|
|
type Pool struct {
|
|
mqcfg *mymq.Config
|
|
}
|
|
|
|
func NewPool(mqcfg *mymq.Config) *Pool {
|
|
return &Pool{
|
|
mqcfg: mqcfg,
|
|
}
|
|
}
|
|
func (p *Pool) Acquire() (*PoolClient, error) {
|
|
cli, err := NewClient(p.mqcfg)
|
|
if err != nil {
|
|
return nil, err
|
|
}
|
|
|
|
return &PoolClient{
|
|
Client: cli,
|
|
owner: p,
|
|
}, nil
|
|
}
|
|
|
|
func (p *Pool) Release(cli *PoolClient) {
|
|
cli.Client.Close()
|
|
}
|