<?php

namespace xmvc\daemon\lib;

use xmvc\Queue;
use xmvc\QueueJob;
use xmvc\daemon\lib\BaseDaemon;
use xmvc\interface\IQueue;

/**
 * 队列服务守护进程
 */
class QueueDaemon extends BaseDaemon
{
    /** 队列对象 */
    protected IQueue $queue;
    /** 队列类型 */
    protected string $type;
    /** 队列配置键 */
    protected string $dns;
    /**
     * 构造
     *
     * @param string $type 队列类型 beanstalk 或 rabbitMQ 或 redis
     * @param string $dns 配置键
     */
    public function __construct(string $type, string $dns)
    {
        $this->type = $type;
        $this->dns = $dns;
    }

    /**
     * 清理资源
     */
    protected function cleanup()
    {
        if (isset($this->queue)) {
            $this->queue->close();
            unset($this->queue);
            static::log('队列连接已关闭');
        }
        parent::cleanup();
    }

    /**
     * 请求停止子进程进程 - 多进程模式下使用
     */
    public static function stopProcess()
    {
        static::$stopRequested = true;
        if (file_exists(static::$forceFilename)) {
            static::log('强制关闭子进程 PID:' . getmypid());
            exit(0);
        } else {
            static::log('优雅关闭子进程 PID:' . getmypid());
        }
    }

    /**
     * 消费队列
     */
    protected function run()
    {
        try {
            $pid = getmypid();
            self::log("PID:{$pid} 创建队列 {$this->type}, dsn: {$this->dns}");
            try {
                $this->queue = Queue::create($this->type, $this->dns);
            } catch (\Throwable $th) {
                $errormsg = "PID:{$pid} 创建队列失败 " . $th->getMessage();
                self::error($errormsg, __FILE__, __LINE__);
                exit(0);
            }
            while (!self::$stopRequested) {
                $queue = $this->queue;
                $job =  $queue->reserve();
                if ($job->id) {
                    self::log("PID: {$pid}, queue->reserve(), Job ID: " . $job->id);
                }
                if ($job->status == 'WAIT') {
                    try {
                        $queue->touch($job);
                        if ($this->handle($job)) {
                            $queue->delete($job);
                        } else {
                            $queue->release($job);
                        }
                    } catch (\Throwable $th) {
                        self::error('run() 出错: ' . $th->getMessage(), __FILE__, __LINE__);
                        $queue->release($job);
                    }
                } elseif ($job->status == 'NONE') {
                    $queue->delete($job);
                }
                // 1000 毫秒（1 秒）到 2000 毫秒（2 秒）之间的随机毫秒数, 多进程时, 间隔不同的时间
                usleep(mt_rand(1000, 2000) * 1000);
                pcntl_signal_dispatch();  // 处理信号
            }
        } catch (\Throwable $th) {
            pcntl_signal_dispatch();  // 处理信号
            if (strpos($th->getMessage(), $this->queue->errorMessagePartBySIGTERM()) !== false) {
                self::$stopRequested = true; // 信号 SIGTERM 导致的中断，立即退出
                self::log("PID:{$pid} run() 被 SIGTERM 信号中断");
            } else {
                $errormsg = 'run() 出错: ' . $th->getMessage();
                self::error($errormsg, __FILE__, __LINE__);
            }
        } finally {
            $this->cleanup();
        }
    }

    /**
     * 处理任务
     */
    protected function handle(QueueJob $job): bool
    {
        self::error('此方法需要在子类中实现', __FILE__, __LINE__);
        exit();
    }
}
