/
mikopbx
/
ModuleNotifier
Обзор
Документация
Войти
/
mikopbx
/
ModuleNotifier
Код
Запросы
0
Задачи
Вики
Пакеты
0
Релизы
0
CI/CD
Аналитика
master
bin/ConnectorDB.php
649 строк
24 KB
Alexey Portnov
Feat: add VK as notification channel alongside Telegram
19 мар 2026, 21:28
19 мар 2026, 21:28
05b7f25
Код
Авторство
О чём код?
<?php /* * MikoPBX - free phone system for small business * Copyright © 2017-2023 Alexey Portnov and Nikolay Beketov * * This program is free software: you can redistribute it and/or modify * it under the terms of the GNU General Public License as published by * the Free Software Foundation; either version 3 of the License, or * (at your option) any later version. * * This program is distributed in the hope that it will be useful, * but WITHOUT ANY WARRANTY; without even the implied warranty of * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the * GNU General Public License for more details. * * You should have received a copy of the GNU General Public License along with this program. * If not, see <https://www.gnu.org/licenses/>. */ namespace Modules\ModuleNotifier\bin; use MikoPBX\Common\Models\Sip; use MikoPBX\Common\Providers\LanguageProvider; use MikoPBX\Common\Providers\TranslationProvider; use MikoPBX\Core\System\Util; use MikoPBX\Core\Workers\WorkerBase; use MikoPBX\Core\System\BeanstalkClient; use MikoPBX\PBXCoreREST\Lib\PBXApiResult; use Modules\ModuleNotifier\Lib\MikoPBXVersion; use Modules\ModuleNotifier\Lib\CacheManager; use Modules\ModuleNotifier\Lib\HistoryParser; use Modules\ModuleNotifier\Models\MessageData; use Modules\ModuleNotifier\Models\ModuleNotifier; use Modules\ModuleNotifier\Lib\Logger; use Modules\ModuleNotifier\Models\CallHistory; use Exception; use Modules\ModuleNotifier\Lib\Providers\CdrDbProvider; use Phalcon\Db\Enum; use DateTime; use Throwable; require_once 'Globals.php'; class ConnectorDB extends WorkerBase { private Logger $logger; public int $cdrOffset = 1; public string $referenceDate = ''; public bool $disableIvr = true; public array $providerName = []; private int $lastSyncTime = 0; /** * Старт работы листнера. * * @param $argv */ public function start($argv):void { $this->logger = new Logger('ConnectorDB', 'ModuleNotifier'); $this->logger->writeInfo('Starting...'); $this->updateSettings(); $beanstalk = new BeanstalkClient(self::class); $beanstalk->subscribe(self::class, [$this, 'onEvents']); $beanstalk->subscribe($this->makePingTubeName(self::class), [$this, 'pingCallBack']); while ($this->needRestart === false) { $beanstalk->wait(10); $this->logger->rotate(); $this->syncCdrData(); } } /** * Получение настроек модуля. * @param int $newCdrOffset * @return void */ public function updateSettings(int $newCdrOffset=0):void { $settings = ModuleNotifier::findFirst(); if(!$settings){ $settings = new ModuleNotifier(); } if($newCdrOffset > 0){ $minOffset = HistoryParser::getMinCdrId(); $settings->cdrOffset = max($newCdrOffset,$minOffset); $settings->save(); } if(empty($settings->referenceDate)){ $settings->referenceDate = date("Y-m-d H:i:s.0", strtotime("-1 days")); $settings->save(); } if(empty($settings->cdrOffset)){ $filter = [ 'limit' => 1, 'order' => 'id DESC' ]; $data = HistoryParser::getCdr($filter); $settings->cdrOffset = $data[0]['id']??1; $settings->save(); $this->logger->writeInfo('Update cdrOffset from cdr: '.$settings->cdrOffset.'...'); } $this->cdrOffset = (int)$settings->cdrOffset; $this->referenceDate = $settings->referenceDate; $providers = Sip::find("type='friend'"); foreach ($providers as $provider) { $this->providerName[$provider->uniqid] = $provider->description; } } /** * Ответ на запрос состояния сервиса. * @param BeanstalkClient $message * @return void */ public function pingCallBack(BeanstalkClient $message): void { parent::pingCallBack($message); $this->updateSettings(); $this->syncCdrData(); } /** * Получение запросов на идентификацию номера телефона. * @param $tube * @return void */ public function onEvents($tube): void { try { $data = json_decode($tube->getBody(), true, 512, JSON_THROW_ON_ERROR); }catch (Exception $e){ return; } if($data['action'] === 'invoke'){ $res_data = []; $funcName = $data['function']??''; if(method_exists($this, $funcName)){ if(count($data['args']) === 0){ $res_data = $this->$funcName(); }else{ $res_data = $this->$funcName(...$data['args']??[]); } }else{ $this->logger->writeError($data); } if(isset($data['need-ret'])){ $res_data = $this->saveResultInTmpFile($res_data); $tube->reply($res_data); } } } /** * Сериализует данные и сохраняет их во временный файл. * @param $data * @return string */ private function saveResultInTmpFile($data):string { try { $res_data = json_encode($data, JSON_THROW_ON_ERROR); }catch (\JsonException $e){ return ''; } $dirsConfig = $this->di->getShared('config'); $tmoDirName = $dirsConfig->path('core.tempDir') . '/ModuleNotifier'; Util::mwMkdir($tmoDirName, true); if (file_exists($tmoDirName)) { $tmpDir = $tmoDirName; }else{ $tmpDir = '/tmp/'; } $downloadCacheDir = $dirsConfig->path('www.downloadCacheDir'); if (!file_exists($downloadCacheDir)) { $downloadCacheDir = ''; } $fileBaseName = md5(microtime(true)); // "temp-" in the filename is necessary for the file to be automatically deleted after 5 minutes. $filename = $tmpDir . '/temp-' . $fileBaseName; file_put_contents($filename, $res_data); if (!empty($downloadCacheDir)) { $linkName = $downloadCacheDir . '/' . $fileBaseName; // For automatic file deletion. // A file with such a symlink will be deleted after 5 minutes by cron. Util::createUpdateSymlink($filename, $linkName, true); } chown($filename, 'www'); return $filename; } /** * Метод следует вызывать при работе с API из прочих процессов. * @param string $function * @param array $args * @param bool $retVal * @return array */ public static function invoke(string $function, array $args = [], bool $retVal = true):array { $req = [ 'action' => 'invoke', 'function' => $function, 'args' => $args ]; $client = new BeanstalkClient(self::class); $object = []; try { if($retVal){ $req['need-ret'] = true; $result = $client->request(json_encode($req, JSON_THROW_ON_ERROR), 60); }else{ $client->publish(json_encode($req, JSON_THROW_ON_ERROR)); return []; } if(file_exists($result)){ $object = json_decode(file_get_contents($result), true, 512, JSON_THROW_ON_ERROR); unlink($result); } } catch (\Throwable $e) { $object = []; } return $object; } /** * Возвращает усеченный слева номер телефона. * * @param $number * * @return bool|string */ public static function getPhoneIndex($number) { $number = preg_replace('/\D+/', '', $number); return substr($number, -9); } /** * Запускаем парсер истории звонков. Парсер сохраняет кэш, кто последний говорил с клиентом. * @return void */ public function syncCdrData():void { if(time() - $this->lastSyncTime < 5){ return; } $this->lastSyncTime = time(); $oldOffset = $this->cdrOffset; $cdrData = HistoryParser::getHistoryData($this->cdrOffset); $arrKeys = (new CallHistory())->toArray(); unset($arrKeys['id']); foreach ($cdrData as $key => $cdr){ $cdr['linkedid'] = $key; $this->sendEditMessage($cdr); foreach ($cdr['rows'] as $row){ /** @var CallHistory $dbData */ $dbData = CallHistory::findFirst("UNIQUEID='{$row['UNIQUEID']}'"); if(!$dbData){ $dbData = new CallHistory(); } foreach ($row as $key => $value){ if(!array_key_exists($key, $arrKeys)){ continue; } $dbData->$key = $value; } foreach ($cdr as $key => $value){ if(!array_key_exists($key, $arrKeys)){ continue; } $dbData->$key = $value; } $this->setCallType($dbData); $dbData->save(); $this->sendAudioToTelegram($row); unset($dbData); } } if($oldOffset !== $this->cdrOffset){ $this->logger->writeInfo("Update cdrOffset from $oldOffset to $this->cdrOffset "); $lastCdrData = HistoryParser::getLastCdrData(); if(!empty($lastCdrData)){ $tmpCdrData = [ 'lastId' => 1*$lastCdrData['id'], 'lastDate' => $lastCdrData['start'], 'nowId' => $this->cdrOffset ]; CacheManager::setCacheData(HistoryParser::CDR_SYNC_PROGRESS_KEY, $tmpCdrData); } $this->updateSettings($this->cdrOffset); } } public function sendEditMessage($cdr):array { $this->logger->writeInfo('sendEditMessage: '.json_encode($cdr)); $messageId = ''; $data = MessageData::findFirst(["linkedId=:linkedId:", 'bind' => ['linkedId' => $cdr['linkedid']]]); if($data){ $messageId.= $data->messageId; $this->logger->writeInfo('find message data: '.json_encode($data->toArray())); } $response = $this->sendCallMessageToTelegram($cdr, $messageId); if(!$data){ $this->logger->writeInfo('add new message data: '.json_encode($response->data)); $data = new MessageData(); $data->linkedId = $cdr['linkedid']; $data->messageId = $response->data['result']['message_id']; $data->messageText = $response->data['result']['text']; $data->save(); } return $response->data; } /** * Format phone number: 7XXXXXXXXXX, 8XXXXXXXXXX, XXXXXXXXXX (10 digits) → +7XXXXXXXXXX * @param string $number * @return string */ public static function formatPhone(string $number): string { $digits = preg_replace('/\D/', '', $number); if (strlen($digits) === 11 && ($digits[0] === '7' || $digits[0] === '8')) { return '+7' . substr($digits, 1); } if (strlen($digits) === 10) { return '+7' . $digits; } return $number; } private function sendCallMessageToTelegram(array $cdr, string $messageId = ""):PBXApiResult { $result = new PBXApiResult(); $line = $this->providerName[$cdr['line']]??$cdr['line']; $src = self::formatPhone($cdr['src'] ?? ''); $dst = self::formatPhone($cdr['dst'] ?? ''); $did = self::formatPhone($cdr['did'] ?? ''); if($cdr['typeCall'] === CallHistory::CALL_TYPE_INCOMING){ $message = self::translate('module_notifier_CALL_TYPE_INCOMING', ['src' => $src, 'dst' => $line, 'did' => $did]); }elseif($cdr['typeCall'] === CallHistory::CALL_TYPE_INNER){ return $result; }elseif($cdr['typeCall'] === CallHistory::CALL_TYPE_OUTGOING && (int)$cdr['answered'] === 0){ $message = self::translate('module_notifier_CALL_TYPE_OUTGOING_FAIL', ['src' => $src, 'dst' => $dst, 'line' => $line]); }elseif($cdr['typeCall'] === CallHistory::CALL_TYPE_OUTGOING){ $message = self::translate('module_notifier_CALL_TYPE_OUTGOING', ['src' => $src, 'dst' => $dst, 'line' => $line]); }elseif($cdr['typeCall'] === CallHistory::CALL_TYPE_MISSED){ $message = self::translate('module_notifier_CALL_TYPE_MISSED', ['src' => $src, 'dst' => $line, 'did' => $did]); }else{ return $result; } $message .= " ($cdr[linkedid])"; if(!empty($messageId)){ $res = Notifier::invoke(Notifier::ACTION_EDIT_MESSAGE, [$messageId, $message]); }else{ $res = Notifier::invoke(Notifier::ACTION_SEND_MESSAGE, [$message]); } return $res; } private function sendAudioToTelegram($dbData) { if($dbData['typeCall'] === CallHistory::CALL_TYPE_INNER){ return; } if(!file_exists($dbData['recordingfile'])){ return; } $messageId = 0; $data = MessageData::findFirst(["linkedId=:linkedId:", 'bind' => ['linkedId' => $dbData['linkedid']]]); if($data){ $messageId = $data->messageId; } $srcFormatted = self::formatPhone($dbData['src_num']); $dstFormatted = self::formatPhone($dbData['dst_num']); $title = "$srcFormatted - $dstFormatted"; $messageText = self::translate('module_notifier_CALL_AUDIO', ['src' => $srcFormatted, 'dst' => $dstFormatted]); Notifier::invoke(Notifier::ACTION_SEND_AUDIO, [$messageText, $dbData['recordingfile'], $title, $messageId], false); } /** * Translates a text string. * * @param string $text The text to translate. * @param array $params * * @return string The translated text. */ private static $moduleTranslations = null; public static function translate(string $text, array $params): string { // Load module translations directly from file if (self::$moduleTranslations === null) { $di = MikoPBXVersion::getDefaultDi(); $lang = 'ru'; if ($di !== null) { $lang = $di->getShared('config')->path('General.WebAdminLanguage') ?: 'ru'; } $moduleDir = dirname(__DIR__); $langFile = $moduleDir . '/Messages/' . $lang . '.php'; if (!file_exists($langFile)) { $langFile = $moduleDir . '/Messages/en.php'; } self::$moduleTranslations = file_exists($langFile) ? include $langFile : []; } $newText = self::$moduleTranslations[$text] ?? $text; foreach ($params as $key => $value) { $newText = str_replace('%' . $key . '%', $value, $newText); } return $newText; } /** * @param $dbData * @return void */ private function setCallType($dbData):void { $number = ''; if($dbData->typeCall === CallHistory::CALL_TYPE_OUTGOING){ if($dbData->billsec === '0'){ $dbData->stateCall = CallHistory::CALL_STATE_OUTGOING_FAIL; }else{ // Успешный исходящий. $dbData->stateCall = CallHistory::CALL_STATE_OK; $number = $dbData->dst_num; } }elseif($dbData->typeCall === CallHistory::CALL_TYPE_INCOMING && $dbData->is_app === '1'){ $dbData->stateCall = CallHistory::CALL_STATE_APPLICATION; }elseif($dbData->typeCall === CallHistory::CALL_TYPE_MISSED){ $dbData->stateCall = CallHistory::CALL_STATE_MISSED; }elseif ($dbData->typeCall === CallHistory::CALL_TYPE_INCOMING){ $dbData->stateCall = CallHistory::CALL_STATE_OK; $number = $dbData->src_num; }elseif($dbData->billsec === '0'){ // Внутренний. $dbData->stateCall = CallHistory::CALL_STATE_OUTGOING_FAIL; }else{ $dbData->stateCall = CallHistory::CALL_STATE_OK; } if($dbData->stateCall === CallHistory::CALL_STATE_OK && !empty($number)){ try { $dateTime = new DateTime($dbData->start); }catch (Exception $e){ return; } $dateTime->modify('-60 minutes'); $oldStart = $dateTime->format('Y-m-d H:i:s'); // Ищем последний вызов по этому номеру телефона. $filter = [ '(dstIndex = :number: OR srcIndex = :number:) AND start BETWEEN :dateFromPhrase1: AND :dateFromPhrase2: AND linkedid<>:linkedid:', 'columns' => 'typeCall', 'bind' => [ 'linkedid' => $dbData->linkedid, 'number' => self::getPhoneIndex($number), 'dateFromPhrase1' => $oldStart, 'dateFromPhrase2' => $dbData->start, ], 'order' => ['start desc'], 'limit' => 1 ]; $oldHistory = CallHistory::find($filter); foreach ($oldHistory as $oldCdr){ if($oldCdr->typeCall === CallHistory::CALL_TYPE_MISSED && $dbData->typeCall === CallHistory::CALL_TYPE_INCOMING){ $dbData->stateCall = CallHistory::CALL_STATE_RECALL_CLIENT; }elseif($oldCdr->typeCall === CallHistory::CALL_TYPE_MISSED && $dbData->typeCall === CallHistory::CALL_TYPE_OUTGOING){ $dbData->stateCall = CallHistory::CALL_STATE_RECALL_USER; } } $filter = [ 'billsec>0 AND start < :dateFromPhrase1: AND linkedid=:linkedid:', 'columns' => 'typeCall', 'bind' => [ 'linkedid' => $dbData->linkedid, 'dateFromPhrase1' => $dbData->start, ], 'order' => ['start desc'], 'limit' => 1 ]; $oldHistory = CallHistory::find($filter); if(!empty($oldHistory->toArray())){ $dbData->stateCall = CallHistory::CALL_STATE_TRANSFER; } } } public function getCdr(array $filter = []): array { $res_data = []; if ($this->filterNotValid($filter)) { return $res_data; } try { $res = CallHistory::find($filter); $res_data = $res->toArray(); } catch (\Throwable $e) { $res_data = []; } return $res_data; } /** * Возвращает количество записпей за период с отбором по номерам. * @param string $start * @param string $end * @param array $numbers * @param array $additionalNumbers * @param array $additionalFilter * @param int $minBilSec * @return array */ public function getCountCdr(string $start, string $end, array $numbers, array $additionalNumbers, array $additionalFilter, int $minBilSec = 0): array { $bindParams = [ ':start' => $start, ':end' => $end ]; $condition = "cdr_general.start BETWEEN :start AND :end"; if (!empty($numbers)) { foreach ($numbers as $value) { $bindParams[":Index$value"] = $value; } $placeholders = implode( ', ', array_map(static function ($value){ return ":Index$value"; }, $numbers) ); $condition .= " AND (cdr_general.dstIndex IN ($placeholders) OR cdr_general.srcIndex IN ($placeholders))"; } if (!empty($additionalNumbers)) { foreach ($additionalNumbers as $value) { $bindParams[":IndexAdd$value"] = $value; } $placeholders = implode( ', ', array_map(static function ($value){ return ":IndexAdd$value"; }, $additionalNumbers) ); $condition .= " AND (cdr_general.dstIndex IN ($placeholders) OR cdr_general.srcIndex IN ($placeholders))"; } $extFilter = $additionalFilter['bind']['filteredExtensions']??[]; if(!empty($extFilter)){ foreach ($extFilter as &$value) { $value = self::getPhoneIndex($value); $bindParams[":IndexAdd$value"] = $value; } unset($value); $placeholders = implode( ', ', array_map(static function ($value){ return ":IndexAdd$value"; }, $extFilter) ); $condition .= ' AND '. str_replace( ['{filteredExtensions:array}', 'dst_num', 'src_num', 'AND ()'], [$placeholders, 'cdr_general.dstIndex', 'cdr_general.srcIndex', ''], $additionalFilter['conditions']??'' ); } if (!$this->di->has(CdrDbProvider::SERVICE_NAME)) { $this->di->register(new CdrDbProvider()); } $billSecFilter = ''; if($minBilSec>0){ $billSecFilter = "billsec > $minBilSec AND "; } $db = $this->di->getShared(CdrDbProvider::SERVICE_NAME); $sql = " SELECT COALESCE(SUM(IIF(t.typeCall=0,1,0)),0) AS cINNER, COALESCE(SUM(IIF(t.typeCall=1,1,0)),0) AS cOUTGOING, COALESCE(SUM(IIF(t.typeCall=2,1,0)),0) AS cINCOMING, COALESCE(SUM(IIF(t.typeCall=3,1,0)),0) AS cMISSED, COUNT(t.linkedid) AS cCalls FROM ( SELECT MIN(cdr_general.id) AS id, MAX(cdr_general.typeCall) AS typeCall, cdr_general.linkedid AS linkedid FROM cdr_general WHERE {$billSecFilter} {$condition} GROUP BY cdr_general.linkedid ) AS t "; try { $result = $db->query($sql, $bindParams); $result->setFetchMode(Enum::FETCH_ASSOC); $row = $result->fetch(); }catch (Throwable $e) { $row = []; Util::sysLogMsg('ERROR- EXTENDED CDR', $sql.PHP_EOL.print_r($bindParams, true)); } return is_array($row)?$row:[]; } /** * Check if the filter has any invalid bind parameters. * * @param array $filter The filter to validate. * @return bool True if the filter has invalid bind parameters, false otherwise. */ private function filterNotValid(array $filter): bool { $haveErrors = false; $validValue = ['0', '']; if (isset($filter['bind'])) { if (is_array($filter['bind'])) { foreach ($filter['bind'] as $bindValue) { if (empty($bindValue) && !in_array($bindValue, $validValue, true)) { $haveErrors = true; } } } else { $haveErrors = true; } } return $haveErrors; } } if(isset($argv) && count($argv) !== 1 && Util::getFilePathByClassName(ConnectorDB::class) === $argv[0]){ ini_set('memory_limit', '512M'); ConnectorDB::startWorker($argv??[]); }