<?php

namespace xmvc\libQueue;

use xmvc\DB;
use xmvc\Queue;
use xmvc\Sqler;
use Pheanstalk\Pheanstalk;
use Pheanstalk\Values\Job;
use xmvc\interface\IQueue;
use xmvc\QueueJob;
use Pheanstalk\Values\TubeName;

/**
 * Beanstalk 是一个轻量的快速的内存消耗极小的队列系统
 * - 适合处理不需要持久化的大量应用请求, 例导入,导出请求
 */
class Beanstalk implements IQueue
{
    /** pda/pheanstalk */
    protected Pheanstalk $pheanstalk;
    /** 管道 */
    protected TubeName $tube;
    /** 管道字符串名称 */
    protected string $tubeName;
    /** 是否持久化 */
    protected bool $persist;
    /** 是否支持多任务消费 */
    protected bool $multiTasking;
    /** 是否入库 */
    protected bool $isToDB;
    /** DB对象 */
    protected DB $db;
    /** beanstalkd 队列存哪个数据库 */
    protected bool $dbCfg;
    /** beanstalkd 队列存哪个表 */
    protected string $table;

    /**
     * 构造
     * - $options = [
     *     'host' => '127.0.0.1',
     *     'port' => 11300,
     *     'tubeName' => 'tube1', // 管道名称
     *     'persist' => false, // 是否持久化 (为true将入库)
     *     'multiTasking' => false, // 是否支持多任务消费 (为true也将入库)
     *     'dbCfg' => 'xmvc', // 持久化存入哪个标签指定的数据库
     *     'table' => 'beanstalkd', // 持久化存入数据库的哪个表
     * ];
     * - 默认端口 11300
     */
    public function __construct(array $options)
    {
        $this->pheanstalk = Pheanstalk::create($options['host'], $options['port']);
        $this->tubeName = $options['tubeName'];
        $this->persist = $options['persist'];
        $this->multiTasking = $options['multiTasking'];
        $this->dbCfg = $options['dbCfg'];
        $this->table = $options['table'];
        if (empty($this->tubeName)) {
            $this->error(lg('管道名称不能为空'));
        }
        $this->tube = new TubeName($this->tubeName);
        if ($this->persist || $this->multiTasking) {
            $this->db = new DB($this->dbCfg);
            $this->isToDB = true;
        } else {
            $this->isToDB = false;
        }

        if ($this->persist) {
            $this->load(); //恢复持久化数据
        }
    }

    protected function error(string $message)
    {
        error($message, null, [__FILE__]);
    }

    /**
     * 恢复队列
     *
     */
    public function load(): self
    {
        $isFindTube = false;
        $isHaveQueue = false;
        $tubes = $this->pheanstalk->listTubes();
        foreach ($tubes as $tubeName) {
            if ($this->tubeName == $tubeName) {
                $isFindTube = true;
                break;
            }
        }
        if ($isFindTube) {
            $stats = $this->pheanstalk->statsTube($this->tube);
            $isHaveQueue = $stats->currentJobsReady > 0;
        } else {
            $this->pheanstalk->useTube($this->tube);
        }
        if (!$isHaveQueue) {
            $q = new Sqler($this->db);
            $q->table($this->table)
                ->where('tubename', $this->tubeName)
                ->where('status', 'DOING')
                ->set('status', 'WAIT');
            $this->db->update($q);
            $q->table($this->table)
                ->where('tubename', $this->tubeName)
                ->where('status', 'WAIT');
            $records = $q->count();
            if ($records > 0) {
                $pageSize = 500;
                $total = ceil($records / $pageSize);
                for ($i = 1; $i <= $total; $i++) {
                    $rows = $q->limit(($i - 1) * $pageSize, $pageSize)->all();
                    foreach ($rows as $row) {
                        $job = (new QueueJob())->fromArray($row);
                        $json = json_encode($job, JSON_UNESCAPED_UNICODE);
                        $this->pheanstalk->put($json, Pheanstalk::DEFAULT_PRIORITY, 15); //延时15秒后消费
                    }
                }
            }
        }
        return $this;
    }
    /**
     * 将任务加入队列
     *
     */
    public function put(array $data): self
    {
        $job = new QueueJob();
        $job->id = 0;
        $job->tubename = $this->tubeName;
        $job->status = 'WAIT';
        $job->trylimit = 0;
        $job->persist = $this->persist;
        $job->data = $data;
        try {
            $json = json_encode($job, JSON_THROW_ON_ERROR | JSON_UNESCAPED_UNICODE);
        } catch (\JsonException $e) {
            $this->error(lg("Beanstalk put() 编码JSON出现异常: ") . $e->getMessage());
        }
        if ($this->isToDB) {
            $q = (new Sqler())->table($this->table)->set($job->toArray());
            if ($this->db->insert($q) == 0) {
                $this->error(lg("Beanstalk 队列 put() 写数据库出错: ") . $json);
            } else {
                $job->id = $this->db->insertId();
                //更新ID后再次编码
                try {
                    $json = json_encode($job->toArray(), JSON_THROW_ON_ERROR | JSON_UNESCAPED_UNICODE);
                } catch (\JsonException $e) {
                    $this->error(lg("Beanstalk put() 编码JSON出现异常: ") . $e->getMessage());
                }
            }
        }
        $this->pheanstalk->useTube($this->tube);
        $this->pheanstalk->put($json, Pheanstalk::DEFAULT_PRIORITY, 0);
        return $this;
    }

    /**
     * 消费队列任务 - 样例
     */
    public function reserveSample()
    {
        $queue = Queue::create('Beanstalk', 'queue.beanstalk.importExcel'); //取 MVC::$cfg['queue']['beanstalk']['importExcel']
        $job =  $queue->reserve();
        if ($job->status == 'WAIT') {
            try {
                $queue->touch($job);
                dump($job);
                $queue->delete($job);
            } catch (\Exception $e) {
                $queue->release($job);
            }
        }
    }

    /**
     * 接收与等待任务
     */
    public function reserve(): QueueJob
    {
        $this->pheanstalk->watch($this->tube);
        $originalJob = $this->pheanstalk->reserve();
        $json = $originalJob->getData();
        $arr = json_decode($json, true);
        /*
        if ($this->isToDB) {
            $arr = (new Sqler($this->db))->table($this->table)->where('id', $arr['id'])->row();
            $arr['data'] = json_decode($arr['data'], true);
        }
        */
        $job = (new QueueJob())->fromArray($arr);
        $job->originalJob = $originalJob;
        return $job;
    }

    /**
     * 保持活动任务
     */
    public function touch(QueueJob $job): void
    {
        $this->pheanstalk->touch($job->originalJob);
        if ($this->isToDB) {
            $q = (new Sqler())->table($this->table)
                ->set('status', 'DOING')
                ->set('atime', date('Y-m-d H:i:s'))
                ->where('id', $job->id);
            if ($this->db->update($q) < 1) {
                $this->error(lg("Beanstalk 队列 touch() 更新数据库出错: ") . json_encode($job->toArray(), JSON_UNESCAPED_UNICODE));
            }
        }
    }

    /**
     * 删除任务
     */
    public function delete(QueueJob $job): void
    {
        if ($this->isToDB) {
            $q = (new Sqler())->table($this->table)->where('id', $job->id);
            if ($this->db->delete($q) < 1) {
                $this->error(lg("Beanstalk 队列 delete() 从数据库删除出错: ") . json_encode($job->toArray(), JSON_UNESCAPED_UNICODE));
            }
        }
        $this->pheanstalk->delete($job->originalJob);
    }

    /**
     * 释放一个任务回到队列中
     */
    public function release(QueueJob $job): void
    {
        if ($this->isToDB) {
            $q = (new Sqler())->table($this->table)
                ->set('status', 'WAIT')
                ->set('trylimit', '@trylimit + 1')
                ->set('atime', date('Y-m-d H:i:s'))
                ->where('id', $job->id);
            if ($job->trylimit >= 3) {
                //超过x次尝试后,将设置为失败
                $q->set('status', 'FAILL');
            }
            if ($this->db->update($q) < 1) {
                $this->error(lg("Beanstalk 队列 release() 更新数据库出错: ") . json_encode($job->toArray(), JSON_UNESCAPED_UNICODE));
            }
        }
        //释放一个任务回到队列中，并设置优先级(越小优先级越高,默认1024)和延迟时间。
        $this->pheanstalk->release($job->originalJob, 2048, 5);
    }

    /**
     * 获取 Beanstalkd 服务器的统计信息
     */
    public function status(): array
    {
        $data = [];
        $data['tubes'] = [];
        $tubes = $this->pheanstalk->listTubes();
        foreach ($tubes as $tubeName) {
            $data['tubes'][(string)$tubeName] = (array)$this->pheanstalk->statsTube(new TubeName($tubeName));
        }
        $data['status'] = $this->pheanstalk->stats();
        return $data;
    }

    /**
     * 获取 Beanstalkd 管道的统计信息
     */
    public function statusTube(string $tubeName)
    {
        return $this->pheanstalk->statsTube(new TubeName($tubeName));
    }

    /**
     * 暂停指定的 Beanstalkd 管道, 在指定的延迟时间后恢复
     */
    public function pauseTube(string $tubeName, int $delay)
    {
        return $this->pheanstalk->pauseTube(new TubeName($tubeName), $delay);
    }
}
