<?php
namespace lib;
/**
* mqRabbitMQ 消息队列
* http://docs.php.net/manual/da/book.amqp.php
* @version 1.0.0

//net.rabbitMQ配置
'rabbitMQ'=>[
    'main' => [
        'server' => '127.0.0.1',
        'port' => 5672,
        'username' => 'root',
        'password' => '123456',
        'vhost'=>'demo'  //虚拟主机
    ],
],

//net.mq消息队列配置
'mq'=>[
    //邮件发送队列
    'sendmail2' => [
        'type' => 'rabbitMQ', //类型
        'server' => 'net.rabbitMQ.main', //服务端配置
        'exchangeType'=> \AMQP_EX_TYPE_DIRECT,//交换机类型
        'isDurable'=> true,//是否为持久队列
        'isWait'=> true,//是否一直等待
    ],
],
//发送调用
public function rabbitMQ_send()
{
    $mq = object(\lib\mq::class,'sendmail2');
    $mq->send('mail1');
    $mq->close();
}

//接收调用
public function rabbitMQ_receive()
{
    $mq = object(\lib\mq::class,'sendmail2');
    $mq->receive(function($e,$queue){
        sleep(1);
        echo $e->getBody()."\r\n";
        $queue->ack($e->getDeliveryTag());
        return true;
    });
}
*/
class mqRabbitMQ
{
	protected $guid = '';         //唯一ID
	protected $exchangeName = ''; //交换机
    protected $queueName = '';    //队列名
    protected $routeKey = '';     //路由键
	protected $isDurable = true;  //是否为持久消息
	protected $isWait = true;     //是否一直等待
    //交换机类型
    //AMQP_EX_TYPE_DIRECT:直连交换机
    //AMQP_EX_TYPE_FANOUT:扇形交换机
    //AMQP_EX_TYPE_HEADERS:头交换机
    //AMQP_EX_TYPE_TOPIC:主题交换机
    protected $exchangeType = ''; //交换机类型
 
    protected $timeout = 60;	    //超时秒
    protected $connection = null;	//连接
	protected $channel = null;      //信道
	protected $exchange = null;     //交换机
	protected $queue = null;        //队列

    /**
    * 创建或绑定一个队列
    * @param string $queueName 队列名
    * @param array $cfg  配置数组
    * @return void
    */
	public function __construct(string $queueName, array $cfg)
    {
        $this->guid = md5(php_uname('n') . getmypid() . $queueName. uniqid(true) . microtime(true));
        $this->queueName = $queueName;
        $this->routeKey = 'route_'.$queueName;
        $this->exchangeName = 'exchange_'.$queueName;
        $this->exchangeType = $cfg['exchangeType'];
        $this->isDurable = $cfg['isDurable'];
        $this->isWait = $cfg['isWait'];
        
        $this->connection = new \AMQPConnection($cfg['server']);
        $this->connection->setReadTimeout($this->timeout);
        $this->createConnection();
        $this->createChannel();
        $this->createExchange();
    }
    
    //创建连接
    private function createConnection(){
        if(!$this->connection->connect()){
            trigger_error(lg('RabbitMQ创建连接失败'), E_USER_ERROR);
        }
    }
    
    //创建信道
    private function createChannel(){
        $this->channel = new \AMQPChannel($this->connection);
    }
    
    //创建交换机
    private function createExchange(){
		$this->exchange = new \AMQPExchange($this->channel);    
		$this->exchange->setName($this->exchangeName);
        $this->exchange->setType(\AMQP_EX_TYPE_DIRECT);  
        if($this->isDurable){
		    $this->exchange->setFlags(\AMQP_DURABLE);
        }
        $this->exchange->declareExchange();
    }
    
    //创建队列
    private function createQueue(){
        $this->queue = new \AMQPQueue($this->channel);
        $this->queue->setName($this->queueName);
        $this->queue->setFlags(\AMQP_DURABLE);
        $this->queue->declareQueue();
        $this->queue->bind($this->exchangeName, $this->routeKey);
    }

    public function send(string $message):string{
        return $this->exchange->publish($message, $this->routeKey);
    }

    public function receive($callback):bool{
        if(is_null($this->queue)){
            //创建队列
            $this->createQueue();
        }
        if($this->isWait){
            while(true){
                sleep(2);
                $this->consume($callback);
            }
        }else{
            $this->consume($callback);
        }
        return true;
    }

    public function consume($callback){
        try{
            $this->queue->consume($callback, \AMQP_NOPARAM, $this->guid);
            $this->queue->cancel($this->guid);
        }catch(\AMQPQueueException $e){
           $this->queue->cancel($this->guid);
        }catch(\AMQPException $e){
            echo $e->getCode()."\r\n";
            echo $e->getMessage()."\r\n";
            $this->createConnection();
            $this->createChannel();
            $this->createExchange();
            $this->createQueue();
        }
    }

    /**
    * 应答确认，队列收到消费者显式确认后，删除该消息
    * @param $deliveryTag //取 $envelope->getDeliveryTag()
    * @return bool
    */
    public function ack($deliveryTag):bool{
        $this->queue->ack($deliveryTag);
        return true;
    }

    public function close():bool{
        return $this->connection->disconnect();
    }
}