Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 58 additions & 0 deletions Lib/AtomicBatch.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,58 @@
<?php

namespace Modules\ModuleExtendedCDRs\Lib;

use RuntimeException;
use Throwable;

final class AtomicBatch
{
/**
* @param object $db Phalcon-compatible transaction adapter.
* @param callable $operation Required database mutations.
* @return array{ok:bool,value:mixed,error:string,rollbackOk:?bool,rollbackError:string}
*/
public static function run(object $db, callable $operation): array
{
$started = false;
try {
if ($db->begin() !== true) {
throw new RuntimeException('transaction_begin_failed');
}
$started = true;
$value = $operation();
if ($db->commit() !== true) {
throw new RuntimeException('transaction_commit_failed');
}
$started = false;
return [
'ok' => true,
'value' => $value,
'error' => '',
'rollbackOk' => null,
'rollbackError' => '',
];
} catch (Throwable $e) {
$rollbackOk = null;
$rollbackError = '';
if ($started) {
try {
$rollbackOk = $db->rollback() === true;
if (!$rollbackOk) {
$rollbackError = 'transaction_rollback_failed';
}
} catch (Throwable $rollbackException) {
$rollbackOk = false;
$rollbackError = $rollbackException->getMessage();
}
}
return [
'ok' => false,
'value' => null,
'error' => $e->getMessage(),
'rollbackOk' => $rollbackOk,
'rollbackError' => $rollbackError,
];
}
}
}
32 changes: 32 additions & 0 deletions Lib/BatchLogContext.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,32 @@
<?php

namespace Modules\ModuleExtendedCDRs\Lib;

final class BatchLogContext
{
/** @return array<string,mixed> */
public static function make(array $state): array
{
$oldOffset = (int)($state['oldOffset'] ?? 0);
$proposedOffset = (int)($state['proposedOffset'] ?? $oldOffset);
$sourceLastId = (int)($state['sourceLastId'] ?? $oldOffset);
return [
'event' => 'cdr_sync_batch',
'oldOffset' => $oldOffset,
'proposedOffset' => $proposedOffset,
'sourceLastId' => $sourceLastId,
'lagBefore' => max(0, $sourceLastId - $oldOffset),
'lagAfter' => max(0, $sourceLastId - $proposedOffset),
'minId' => (int)($state['minId'] ?? 0),
'maxId' => (int)($state['maxId'] ?? 0),
'linkedIdCount' => (int)($state['linkedIdCount'] ?? 0),
'rowCount' => (int)($state['rowCount'] ?? 0),
'inserted' => (int)($state['inserted'] ?? 0),
'updated' => (int)($state['updated'] ?? 0),
'mode' => (string)($state['mode'] ?? 'normal'),
'elapsedMs' => max(0, (int)($state['elapsedMs'] ?? 0)),
'outcome' => (string)($state['outcome'] ?? 'unknown'),
'errorCategory' => (string)($state['errorCategory'] ?? ''),
];
}
}
30 changes: 30 additions & 0 deletions Lib/BatchPersistenceResult.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
<?php

namespace Modules\ModuleExtendedCDRs\Lib;

final class BatchPersistenceResult
{
/** @return array{ok:bool,inserted:int,updated:int,errorCategory:string,message:string} */
public static function success(int $inserted, int $updated): array
{
return [
'ok' => true,
'inserted' => $inserted,
'updated' => $updated,
'errorCategory' => '',
'message' => '',
];
}

/** @return array{ok:bool,inserted:int,updated:int,errorCategory:string,message:string} */
public static function failure(string $category, string $message = ''): array
{
return [
'ok' => false,
'inserted' => 0,
'updated' => 0,
'errorCategory' => $category,
'message' => $message,
];
}
}
21 changes: 21 additions & 0 deletions Lib/CheckpointPolicy.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,21 @@
<?php

namespace Modules\ModuleExtendedCDRs\Lib;

final class CheckpointPolicy
{
/**
* @param array<string,mixed> $batch
*/
public static function nextOffset(array $batch): int
{
$oldOffset = (int)$batch['oldOffset'];
if (empty($batch['requestOk']) || empty($batch['saveOk']) || !empty($batch['newQuarantine'])) {
return $oldOffset;
}

// SelectCDR intentionally filters and groups source rows, so gaps in raw IDs
// are normal. The parser checkpoint is the authoritative safe boundary.
return max($oldOffset, (int)$batch['parsedOffset']);
}
}
52 changes: 18 additions & 34 deletions Lib/GetReport.php
Original file line number Diff line number Diff line change
Expand Up @@ -302,26 +302,15 @@ private function prepareCdrData($selectedRecords):array
];

$providers = Sip::find("type='friend'");
$providerName = [];
$providerNameByLogin = [];
$providerRows = [];
foreach ($providers as $provider) {
$providerName[$provider->uniqid] = $provider->description;
// Несколько учёток одного провайдера приходят на одну линию и различаются только по DID.
// Строим карту "логин провайдера (username) => название" для приоритетного сопоставления по DID.
$login = (string)$provider->username;
if ($login === '') {
continue;
}
if (array_key_exists($login, $providerNameByLogin)) {
// Один и тот же логин у нескольких учёток — сопоставление по DID неоднозначно,
// отключаем его для этого логина (null), чтобы не показать чужого провайдера.
if ($providerNameByLogin[$login] !== $provider->description) {
$providerNameByLogin[$login] = null;
}
} else {
$providerNameByLogin[$login] = $provider->description;
}
$providerRows[] = [
'uniqid' => $provider->uniqid,
'username' => $provider->username,
'description' => $provider->description,
];
}
$trunkResolver = new TrunkResolver($providerRows);
unset($providers);

$statsCall = [
Expand Down Expand Up @@ -374,23 +363,18 @@ private function prepareCdrData($selectedRecords):array
// Приоритет: если DID входящего звонка совпадает с логином (username) учётки провайдера —
// показываем именно её. Так различаются несколько учёток одного провайдера на одной линии
// (ограничение Asterisk: плечи приходят на один peer и отличаются только по DID).
$didName = '';
if ($record->typeCall === CallHistory::CALL_TYPE_INCOMING
&& !empty($record->did)
&& !empty($providerNameByLogin[$record->did])) {
$didName = $providerNameByLogin[$record->did];
}
if ($didName !== '') {
$linkedRecord->line = $didName;
if (!empty($record->line)) {
$linkedRecord->lineId = $record->line;
}
$resolution = $trunkResolver->resolve(
['line' => $record->line, 'did' => $record->did],
$record->typeCall
);
if ($resolution['status'] === 'resolved' && $resolution['source'] === 'did_username') {
$linkedRecord->line = $resolution['name'];
$linkedRecord->lineId = $resolution['id'];
$lineFixedByDid[$record->linkedid] = true;
} elseif (empty($lineFixedByDid[$record->linkedid])) {
$newLine = $providerName[$record->line] ?? $record->line;
if(!empty($newLine)){
$linkedRecord->line = $newLine;
$linkedRecord->lineId = $record->line;
if ($resolution['name'] !== '') {
$linkedRecord->line = $resolution['name'];
$linkedRecord->lineId = $resolution['id'];
}
}
$linkedRecord->disposition = $linkedRecord->disposition !== 'ANSWERED' ? $disposition : 'ANSWERED';
Expand Down Expand Up @@ -1329,4 +1313,4 @@ public function getDateRanges(string $srcPeriod): string
];
return $dateRanges[$srcPeriod] ?? $srcPeriod;
}
}
}
42 changes: 42 additions & 0 deletions Lib/HistoryBatchResult.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,42 @@
<?php

namespace Modules\ModuleExtendedCDRs\Lib;

final class HistoryBatchResult
{
/**
* @param array<string,array> $data
* @return array<string,mixed>
*/
public static function make(
int $oldOffset,
array $data,
bool $ok,
int $linkedIdLimit,
string $error = '',
?int $newOffset = null
): array {
$ids = [];
foreach ($data as $call) {
foreach ($call['rows'] ?? [] as $row) {
$id = (int)($row['id'] ?? 0);
if ($id > 0) {
$ids[] = $id;
}
}
}

return [
'ok' => $ok,
'data' => $data,
'oldOffset' => $oldOffset,
'newOffset' => $newOffset ?? $oldOffset,
'minId' => empty($ids) ? 0 : min($ids),
'maxId' => empty($ids) ? 0 : max($ids),
'rowCount' => count($ids),
'linkedIdCount' => count($data),
'limitReached' => count($data) >= $linkedIdLimit,
'error' => $ok ? '' : ($error ?: 'source_request_failed'),
];
}
}
56 changes: 47 additions & 9 deletions Lib/HistoryParser.php
Original file line number Diff line number Diff line change
Expand Up @@ -71,8 +71,9 @@ public static function getCdr(array $filter = [], ?bool &$requestOk = null): arr
try {
[$result, $message] = $client->sendRequest(json_encode($filter), 30);
if ($result!==false){
$requestOk = true;
$filename = json_decode($message, true, 512, JSON_THROW_ON_ERROR);
// SelectCDR may return a successful empty response without creating a file.
$requestOk = empty($filename);
}
} catch (\Throwable $e) {
$filename = '';
Expand All @@ -81,6 +82,7 @@ public static function getCdr(array $filter = [], ?bool &$requestOk = null): arr
if (is_string($filename) && file_exists($filename)) {
try {
$result_data = json_decode(file_get_contents($filename), true, 512, JSON_THROW_ON_ERROR);
$requestOk = is_array($result_data);
} catch (\Throwable $e) {
SystemMessages::sysLogMsg('HistoryParser:SELECT_CDR_TUBE', 'Error parse response.');
}
Expand All @@ -92,6 +94,9 @@ public static function getCdr(array $filter = [], ?bool &$requestOk = null): arr
shell_exec("$findPath -L $downloadCacheDir -samefile $filename -delete");
}
unlink($filename);
} elseif (!empty($filename)) {
// A non-empty response promised a result file, but it is unavailable.
$requestOk = false;
}

return $result_data;
Expand All @@ -116,7 +121,11 @@ public static function getQueues():array
* @param int $offset
* @return void
*/
public static function getHistoryData(int $offset = 1, array $excludeLinkedIds = []):array
public static function getHistoryData(
int $offset = 1,
array $excludeLinkedIds = [],
int $linkedIdLimit = self::LIMIT_CDR
):array
{
$filter = [
"type = :extType:",
Expand Down Expand Up @@ -144,7 +153,7 @@ public static function getHistoryData(int $offset = 1, array $excludeLinkedIds =
'order' => 'id ASC',
'group' => 'linkedid',
'columns' => 'linkedid',
'limit' => self::LIMIT_CDR,
'limit' => $linkedIdLimit,
'add_pack_query' => $add_query,
];

Expand All @@ -155,7 +164,17 @@ public static function getHistoryData(int $offset = 1, array $excludeLinkedIds =
$filter['bind']['exclude'] = array_values($excludeLinkedIds);
}

$cdrData = self::getCdr($filter);
$requestOk = false;
$cdrData = self::getCdr($filter, $requestOk);
if (!$requestOk) {
return HistoryBatchResult::make(
$offset,
[],
false,
$linkedIdLimit,
'source_request_failed'
);
}
$resultRows = [];
if(count($cdrData)>0){
$queues = self::getQueues();
Expand Down Expand Up @@ -268,7 +287,7 @@ public static function getHistoryData(int $offset = 1, array $excludeLinkedIds =
unset($firstQueue);
$resultRows[$cdr['linkedid']]['rows'][] = $cdr;
}
$calculatedOffset = min($offset + self::LIMIT_CDR, $newOffset);
$calculatedOffset = min($offset + $linkedIdLimit, $newOffset);
$calculatedOffset = max($calculatedOffset, $minNewOffset);
}

Expand All @@ -285,7 +304,14 @@ public static function getHistoryData(int $offset = 1, array $excludeLinkedIds =
}
}

return ['data' => $resultRows, 'newOffset' => $calculatedOffset ?? $offset];
return HistoryBatchResult::make(
$offset,
$resultRows,
true,
$linkedIdLimit,
'',
$calculatedOffset ?? $offset
);
}

/**
Expand Down Expand Up @@ -329,14 +355,26 @@ public static function getActiveLinkedIds(array $linkedIds, int $offset):?array
* @return array
*/
public static function getLastCdrData():array
{
$state = self::getLastCdrState();
return $state['data'];
}

/**
* Returns the last source CDR and an explicit request status.
*
* @return array{ok:bool,data:array}
*/
public static function getLastCdrState(): array
{
$filter = [
'columns' => 'id,start',
'order' => 'id DESC',
'limit' => 1,
];
$res = \Modules\ModuleExtendedCDRs\Lib\HistoryParser::getCdr($filter);
return $res[0]??[];
$requestOk = false;
$res = self::getCdr($filter, $requestOk);
return ['ok' => $requestOk, 'data' => $res[0] ?? []];
}

/**
Expand Down Expand Up @@ -376,4 +414,4 @@ public static function getMinCdrId():int
}
return $id;
}
}
}
Loading
Loading