Files
hospital-applets-api/app/Amqp/Consumer/SendSubMessageConsumer.php
T

157 lines
5.4 KiB
PHP

<?php
declare(strict_types=1);
namespace App\Amqp\Consumer;
use App\Constants\HttpEnumCode;
use App\Exception\BusinessException;
use App\Model\LogMessagePush;
use App\Model\SubTemplate;
use App\Model\User;
use App\Model\UserDoctor;
use App\Model\UserPatient;
use App\Services\UserService;
use App\Utils\Log;
use Extend\Wechat\Wechat;
use Hyperf\Amqp\Result;
use Hyperf\Amqp\Annotation\Consumer;
use Hyperf\Amqp\Message\ConsumerMessage;
use PhpAmqpLib\Message\AMQPMessage;
use Psr\Container\ContainerExceptionInterface;
use Psr\Container\NotFoundExceptionInterface;
use Symfony\Contracts\HttpClient\Exception\ClientExceptionInterface;
use Symfony\Contracts\HttpClient\Exception\DecodingExceptionInterface;
use Symfony\Contracts\HttpClient\Exception\RedirectionExceptionInterface;
use Symfony\Contracts\HttpClient\Exception\ServerExceptionInterface;
use Symfony\Contracts\HttpClient\Exception\TransportExceptionInterface;
/**
* 发送订阅消息
*/
#[Consumer(exchange: 'amqp.direct', routingKey: 'SendSubMessage', queue: 'send.sub.message.queue', nums: 1)]
class SendSubMessageConsumer extends ConsumerMessage
{
/**
* @param $data
* [
* "push_user_id" // 用户id(被推送者)
* "template_title" // 推送的模版名称
* "params" => [ // 推送所需的参数
* "page" // 跳转页面
* "data" => [
* "thing1" => [
* "value" => [
* "参数1"
* ]
* ]
* ]
*
* ]
* ]
* @param AMQPMessage $message
* @return string
* @throws ContainerExceptionInterface
* @throws NotFoundExceptionInterface
* @throws ClientExceptionInterface
* @throws DecodingExceptionInterface
* @throws RedirectionExceptionInterface
* @throws ServerExceptionInterface
* @throws TransportExceptionInterface
*/
public function consumeMessage($data, AMQPMessage $message): string
{
try {
// 获取被推送用户信息
$params = array();
$params['user_id'] = $data['push_user_id'];
$user = User::getOne($params);
if (empty($user)){
$this->addMessagePushLog($data,"未查询到被推送用户信息");
return Result::ACK;
}
// 获取open_id
$UserService = new UserService();
$open_id = $UserService->getOpenIdWithUserId($user['user_id'],$user['user_type']);
if (empty($open_id)){
$this->addMessagePushLog($data,"未获取到被推送用户open_id");
return Result::ACK;
}
// 获取订阅消息模版数据
$params = array();
$params['client_type'] = $user['user_type']; // 客户端类型(1:患者端 2:医师端 3:药师端)
$params['template_title'] = $data['template_title'];
$sub_template = SubTemplate::getOne($params);
if (empty($sub_template)){
$this->addMessagePushLog($data,"未查询到推送模版id");
return Result::ACK;
}
$sub_template = $sub_template->toArray();
// 处理发送环境
$miniprogram_state = "developer";
if (env("APP_ENV") == "prod"){
$miniprogram_state = "formal";
}
$options = [
"template_id" => $sub_template['wx_template_id'],
"page" => $data['params']['page'],
"touser" => $open_id,
"data" => $data['params']['data'],
"miniprogram_state" => $miniprogram_state,
"lang" => "zh_CN",
];
// 发起推送
$Wechat = new Wechat($user['user_type']);
$result = $Wechat->sendSubscribeMessage($options);
Log::getInstance()->info("订阅消息推送成功:" . $result);
// 记录推送记录
$this->addMessagePushLog($data,"",$sub_template['wx_template_id'],1,json_encode($options,JSON_UNESCAPED_UNICODE));
} catch (\Exception $e) {
$this->addMessagePushLog($data,$e->getMessage());
return Result::ACK;
}
return Result::ACK;
}
/**
* 增加日志
* @param array $params 推送参数
* @param string $wx_template_id 微信推送模版id
* @param int $status 推送状态(1:成功 2:失败)
* @param string $fail_reason 推送失败原因
* @param string $content 推送内容
*/
public function addMessagePushLog(array $params,string $fail_reason = '',string $wx_template_id = "",int $status = 2,string $content = '')
{
$data = array();
$data['push_user_id'] = $params['push_user_id'];
$data['template_title'] = $params['template_title'];
if (!empty($wx_template_id)){
$data['wx_template_id'] = $wx_template_id;
}
$data['params'] = json_encode($params,JSON_UNESCAPED_UNICODE);
$data['status'] = $status;
if (!empty($fail_reason)){
$data['fail_reason'] = $fail_reason;
}
if (!empty($content)){
$data['content'] = $content;
}
$log_message_push = LogMessagePush::addLogMessagePush($data);
if (empty($log_message_push)){
Log::getInstance()->error("增加推送日志失败:" . json_encode($data,JSON_UNESCAPED_UNICODE));
}
}
}