From 0a72197066a1fad6b91e866b2ce47e322cfcfd40 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:01:00 +0300 Subject: [PATCH 01/18] docs: design resilient CDR synchronization --- ...22-cdr-sync-and-trunk-resolution-design.md | 159 ++++++++++++++++++ 1 file changed, 159 insertions(+) create mode 100644 docs/superpowers/specs/2026-07-22-cdr-sync-and-trunk-resolution-design.md diff --git a/docs/superpowers/specs/2026-07-22-cdr-sync-and-trunk-resolution-design.md b/docs/superpowers/specs/2026-07-22-cdr-sync-and-trunk-resolution-design.md new file mode 100644 index 0000000..cc39294 --- /dev/null +++ b/docs/superpowers/specs/2026-07-22-cdr-sync-and-trunk-resolution-design.md @@ -0,0 +1,159 @@ +# Надёжная синхронизация CDR и определение транков + +## Цель + +Исключить повторное исчезновение расширенной истории вызовов при отставании фонового воркера и сделать определение транка детерминированным. Переустановка или ручной перезапуск модуля не должны требоваться для восстановления истории. + +## Наблюдаемые дефекты + +1. `ConnectorDB` может продолжать работать без исключений, но сильно отставать от основной CDR-базы. В это время `getHistory` возвращает корректный, но пустой DataTables-ответ. +2. Текущий карантин срабатывает только когда один `linkedid` достигает 5000 строк. Медленная обработка большого числа меньших пакетов остаётся незаметной и не включает ускоренное восстановление. +3. Имя транка восстанавливается из нескольких косвенных признаков. Неоднозначность не диагностируется, а соответствие DID и `Sip.username` применимо не ко всем конфигурациям. + +## Границы решения + +В объём входят: + +- гарантированное продвижение checkpoint без потери обычных CDR; +- ускоренная обработка накопившегося backlog; +- изоляция проблемных `linkedid` от основного потока; +- восстановление после рестарта и временных ошибок; +- наблюдаемое состояние синхронизации; +- единый resolver транков с диагностикой неоднозначности; +- автоматические модульные и интеграционные тесты ключевых сценариев. + +В объём не входят: + +- полная перестройка схемы `CallHistory`; +- изменение первичной CDR-базы MikoPBX; +- автоматическое исправление конфигурации SIP-провайдеров; +- удаление или пересоздание уже накопленной истории при обновлении. + +## Архитектура синхронизации + +### Состояние + +Синхронизатор хранит устойчивое состояние: + +- `committedOffset` — последний непрерывно подтверждённый ID источника; +- `sourceLastId` — последний доступный ID основной CDR-базы на момент проверки; +- `lastSuccessAt` — время последнего успешно зафиксированного пакета; +- записи карантина для конкретных `linkedid` и диапазонов ID. + +Существующее значение `cdrOffset` мигрирует в `committedOffset` без сброса истории. Повторный запуск читает сохранённое состояние и продолжает с него. + +### Основной цикл + +Перед обработкой воркер получает `sourceLastId` и вычисляет `lag = sourceLastId - committedOffset`. + +- В обычном режиме обрабатывается ограниченный пакет, после чего сохраняется текущая пауза. +- При превышении порога lag включается catch-up: пакеты обрабатываются подряд без десятисекундной паузы, но с ограничением времени одного прохода и возможностью штатного завершения воркера. +- Размер пакета может увеличиваться только до безопасного верхнего предела. Он уменьшается после таймаута или ошибки источника. +- Когда lag опускается ниже нижнего порога, воркер возвращается в обычный режим. Разные пороги входа и выхода предотвращают постоянное переключение режимов. + +Ответ выборки содержит данные и метаданные: минимальный и максимальный ID, число `linkedid`, число строк, признак достижения лимита и успешность запроса. Пустой успешный ответ отличается от таймаута или недоступности источника. + +### Фиксация пакета + +Обработка пакета идемпотентна. Записи истории сохраняются через существующую уникальность или явный upsert. После успешного сохранения обычных записей и регистрации исключённых элементов обновляется `committedOffset`. + +Checkpoint нельзя продвигать, если: + +- запрос к источнику завершился ошибкой; +- данные не удалось сохранить; +- диапазон содержит необъяснённый разрыв, который не зарегистрирован как пропущенный или помещённый в карантин. + +Повторная обработка уже сохранённого диапазона не создаёт дубликаты и не меняет корректные итоговые значения. + +## Карантин и восстановление + +Проблемный `linkedid` не должен бесконечно удерживать общий поток. В карантине хранятся: + +- `linkedid`; +- минимальный и максимальный ID затронутого диапазона; +- причина (`row_limit`, `timeout`, `parse_error`, `save_error`); +- число попыток; +- время первой и последней ошибки; +- `nextRetryAt`; +- состояние (`pending`, `resolved`, `manual`). + +Попадание в карантин требует объективного сигнала: достижения лимита, повторяемого таймаута или воспроизводимой ошибки конкретного звонка. Общая ошибка БД не превращает весь пакет в карантин. + +Отдельный reconciler периодически повторяет обработку с экспоненциальной задержкой. После успешного сохранения запись получает состояние `resolved`. После предельного числа попыток она становится `manual`, но не блокирует основной поток. Для контроля пропусков reconciler также повторно читает небольшое перекрывающееся окно перед checkpoint и выполняет идемпотентный upsert. + +## Наблюдаемость и поведение интерфейса + +В журнал записываются структурированные события: + +- `sync_state` с offset, sourceLastId, lag, режимом и lastSuccessAt; +- `batch_completed` с диапазоном, длительностью и количеством строк; +- `batch_failed` с категорией ошибки без продвижения checkpoint; +- `linkedid_quarantined` и результат повторной обработки; +- `trunk_resolution_ambiguous` с безопасными идентификаторами кандидатов. + +Состояние синхронизации доступно существующему контроллеру. Интерфейс продолжает показывать уже сохранённую историю и отдельно сообщает, что новые вызовы догружаются. Пустая таблица не должна создавать впечатление отсутствия звонков, если lag превышает допустимый порог или `lastSuccessAt` устарел. + +Логи не должны содержать пароли, SIP-секреты или полный набор персональных данных звонка. + +## Определение транка + +Логика переносится в отдельный `TrunkResolver`, не зависящий от формирования HTML/JSON отчёта. Он получает нормализованные признаки CDR и индекс известных провайдеров. + +Приоритет входящего вызова: + +1. точный стабильный идентификатор линии из CDR, если он однозначно соответствует провайдеру; +2. точное совпадение нормализованного DID с уникальным `Sip.username`; +3. однозначное соответствие peer/канала известному провайдеру; +4. исходное техническое значение линии с признаком `unresolved`. + +Для исходящего вызова DID не используется как признак провайдера; приоритет имеют идентификатор линии и peer/канал. + +Если один признак соответствует нескольким провайдерам, resolver не выбирает произвольного кандидата. Он возвращает техническое значение, статус `ambiguous` и пишет диагностическое событие. Имя, зафиксированное более сильным признаком для `linkedid`, не затирается последующими плечами с более слабым результатом. + +Нормализация удаляет только форматирование номера, которое не меняет его смысл. Префиксы и значимые цифры не отбрасываются без явной конфигурации. + +## Обработка ошибок + +- Таймаут или некорректный ответ `SELECT_CDR_TUBE` считается ошибкой запроса, а не пустым результатом. +- Временная ошибка источника или БД приводит к backoff и повтору того же диапазона. +- Ошибка одного `linkedid`, подтверждённая отдельной обработкой, изолируется в карантин. +- При невозможности прочитать сохранённый checkpoint воркер останавливает синхронизацию с явной ошибкой; он не начинает с нуля и не перескакивает к последнему ID. +- Штатный сигнал остановки завершает текущую атомарную операцию и не оставляет частично зафиксированный checkpoint. + +## Совместимость и миграция + +- Минимальная версия PHP остаётся `7.4.6`; новые возможности PHP 8 не используются. +- Обновление создаёт недостающие служебные таблицы и поля без удаления существующих данных. +- Существующий `cdrOffset` используется как начальный `committedOffset`. +- Старые записи `oversized_linkedids` мигрируют или читаются как карантин с причиной `row_limit`. +- Откат версии не должен удалять новые служебные данные; старая версия сможет продолжить работу со старым `cdrOffset`. + +## Проверка + +Автоматические тесты должны доказать: + +1. catch-up включается при backlog и возвращается в обычный режим; +2. checkpoint продвигается только после успешного сохранения; +3. рестарт продолжает с последнего подтверждённого offset; +4. повторная обработка не создаёт дубликаты; +5. разрыв ID не приводит к молчаливой потере данных; +6. тяжёлый или ошибочный `linkedid` изолируется, остальные звонки продолжают синхронизироваться; +7. общая ошибка БД не помещает нормальные звонки в карантин; +8. reconciler успешно возвращает исправившийся звонок и переводит запись в `resolved`; +9. таймаут источника отличается от пустого результата; +10. входящий транк определяется каждым поддержанным сильным признаком; +11. неоднозначное соответствие возвращает `ambiguous`, а не случайного провайдера; +12. более слабое плечо не затирает уже определённый транк; +13. API сообщает lag и продолжает возвращать ранее сохранённую историю. + +Интеграционная проверка на копии клиентских данных измеряет скорость сокращения lag, отсутствие дубликатов и корректное восстановление после принудительного рестарта воркера. + +## Критерии готовности + +- При доступном источнике lag монотонно сокращается в catch-up режиме и достигает рабочего порога без переустановки модуля. +- Один некорректный звонок не останавливает новые записи истории. +- После рестарта нет потери и дублирования CDR. +- Оператор видит состояние догрузки вместо необъяснимо пустой истории. +- Resolver никогда не выбирает транк при неоднозначном совпадении. +- Все перечисленные автоматические тесты проходят на PHP 7.4.6. +- Push выполняется только после отдельного разрешения пользователя. From eecc80db52376530fea19eca2ef714312fb56dc6 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:04:56 +0300 Subject: [PATCH 02/18] docs: plan resilient CDR synchronization --- ...07-22-resilient-cdr-sync-implementation.md | 181 ++++++++++++++++++ 1 file changed, 181 insertions(+) create mode 100644 docs/superpowers/plans/2026-07-22-resilient-cdr-sync-implementation.md diff --git a/docs/superpowers/plans/2026-07-22-resilient-cdr-sync-implementation.md b/docs/superpowers/plans/2026-07-22-resilient-cdr-sync-implementation.md new file mode 100644 index 0000000..27b9e57 --- /dev/null +++ b/docs/superpowers/plans/2026-07-22-resilient-cdr-sync-implementation.md @@ -0,0 +1,181 @@ +# Resilient CDR Synchronization Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Сделать синхронизацию CDR самовосстанавливающейся при backlog и проблемных звонках, а определение транков — детерминированным и диагностируемым. + +**Architecture:** Чистые классы `SyncPolicy` и `TrunkResolver` отделяют решения от инфраструктуры и тестируются без запущенной MikoPBX. `HistoryParser` возвращает явный статус и метаданные выборки; `ConnectorDB` применяет адаптивный catch-up, фиксирует checkpoint только после успешной записи и использует расширенный карантин для проблемных `linkedid`. + +**Tech Stack:** PHP 7.4.6, Phalcon models, SQLite, MikoPBX WorkerBase/BeanstalkClient, самостоятельные PHP regression tests. + +## Global Constraints + +- Не использовать синтаксис и API новее PHP 7.4.6. +- Не удалять и не пересоздавать существующую историю вызовов. +- Сохранить `cdrOffset` как совместимый устойчивый checkpoint. +- Повторная обработка должна быть идемпотентной. +- Не логировать SIP-секреты и пароли. +- Не выполнять push без отдельного разрешения пользователя. +- Проверка на `serber@boffart.miko.ru` начинается с backup и read-only baseline; рабочий модуль заменяется только контролируемо с возможностью немедленного отката. + +--- + +### Task 1: Pure synchronization policy + +**Files:** +- Create: `Lib/SyncPolicy.php` +- Create: `tests/SyncPolicyTest.php` + +**Interfaces:** +- Produces: `SyncPolicy::decide(int $offset, int $sourceLastId, bool $requestOk, bool $limitReached): array` returning `lag`, `mode`, `delay`, `batchLinkedIds`. + +- [ ] **Step 1: Write failing tests** covering zero lag, normal mode, catch-up entry, catch-up exit hysteresis, request failure, and capped batch size. +- [ ] **Step 2: Run** `php tests/SyncPolicyTest.php`; expect a missing-class failure. +- [ ] **Step 3: Implement** a PHP 7.4-compatible pure policy with constants `CATCH_UP_ENTER_LAG`, `CATCH_UP_EXIT_LAG`, `NORMAL_BATCH_LINKED_IDS`, `CATCH_UP_BATCH_LINKED_IDS`, `NORMAL_DELAY_SECONDS`, and `ERROR_DELAY_SECONDS`. +- [ ] **Step 4: Run** `php tests/SyncPolicyTest.php`; expect `SyncPolicyTest: OK`. +- [ ] **Step 5: Commit** `Lib/SyncPolicy.php` and `tests/SyncPolicyTest.php` with message `feat: add adaptive CDR sync policy`. + +### Task 2: Explicit source result metadata + +**Files:** +- Modify: `Lib/HistoryParser.php` +- Create: `tests/HistoryBatchResultTest.php` + +**Interfaces:** +- Produces: `HistoryParser::getHistoryData(int $offset, array $excludeLinkedIds = [], int $linkedIdLimit = self::LIMIT_CDR): array` with keys `ok`, `data`, `oldOffset`, `newOffset`, `minId`, `maxId`, `rowCount`, `linkedIdCount`, `limitReached`, `error`. +- Consumes: batch size selected by `SyncPolicy`. + +- [ ] **Step 1: Extract** a pure metadata finalizer callable by a standalone test without MikoPBX services. +- [ ] **Step 2: Write failing tests** proving that a successful empty result differs from a failed request, and that min/max/count/limit metadata is correct. +- [ ] **Step 3: Run** `php tests/HistoryBatchResultTest.php`; expect failure for missing metadata behavior. +- [ ] **Step 4: Thread request success** from `getCdr()` into `getHistoryData()` and return the complete result contract without treating timeout as empty data. +- [ ] **Step 5: Run** both Task 1 and Task 2 tests; expect `OK`. +- [ ] **Step 6: Commit** with message `fix: expose reliable CDR batch status`. + +### Task 3: Durable quarantine model + +**Files:** +- Modify: `Models/OversizedLinkedIds.php` +- Modify: `bin/ConnectorDB.php` +- Create: `Lib/QuarantinePolicy.php` +- Create: `tests/QuarantinePolicyTest.php` + +**Interfaces:** +- Produces: quarantine fields `minId`, `maxId`, `reason`, `attempts`, `firstFailureAt`, `lastFailureAt`, `nextRetryAt`, `status`. +- Produces: `QuarantinePolicy::nextState(array $current, string $reason, int $now): array`. + +- [ ] **Step 1: Write failing pure tests** for initial quarantine, exponential retry delay, cap, resolved state, and manual state after the maximum attempts. +- [ ] **Step 2: Run** `php tests/QuarantinePolicyTest.php`; expect failure. +- [ ] **Step 3: Implement** `QuarantinePolicy` without framework dependencies. +- [ ] **Step 4: Extend** `ensureTableExists()` with additive `PRAGMA table_info`/`ALTER TABLE ADD COLUMN` migration; never drop the table. +- [ ] **Step 5: Update** persistence and cache loading so legacy rows remain `row_limit` quarantines and active statuses alone are excluded. +- [ ] **Step 6: Run** syntax checks and policy tests; expect success. +- [ ] **Step 7: Commit** with message `feat: persist retryable CDR quarantine`. + +### Task 4: Transaction-safe checkpoint and catch-up loop + +**Files:** +- Modify: `bin/ConnectorDB.php` +- Modify: `Lib/HistoryParser.php` +- Create: `tests/CheckpointPolicyTest.php` + +**Interfaces:** +- Consumes: `SyncPolicy::decide()` and the HistoryParser batch contract. +- Produces: one `syncCdrData()` result describing whether a batch committed and the next delay. + +- [ ] **Step 1: Write failing tests** for checkpoint retention on source error, save error and unexplained gap; advancement after successful save; and retention for a newly quarantined truncated call. +- [ ] **Step 2: Run** `php tests/CheckpointPolicyTest.php`; expect failure. +- [ ] **Step 3: Extract** a pure checkpoint decision from the existing inline offset arithmetic. +- [ ] **Step 4: Change** `syncCdrData()` to reject `ok=false`, choose batch size through `SyncPolicy`, persist rows before offset, and update `cdrOffset` only after the batch is committed. +- [ ] **Step 5: Change** the worker loop to use the policy delay: zero/short delay in catch-up, normal delay when current, bounded backoff on failure, while checking `needRestart` between batches. +- [ ] **Step 6: Preserve** the current oversized detection but express it as a quarantine reason and do not skip unexplained ID ranges. +- [ ] **Step 7: Run** all standalone tests plus `php -l` on modified files. +- [ ] **Step 8: Commit** with message `fix: make CDR checkpoint catch-up resilient`. + +### Task 5: Quarantine reconciler and overlap repair + +**Files:** +- Create: `Lib/CdrReconciler.php` +- Modify: `bin/ConnectorDB.php` +- Modify: `Lib/HistoryParser.php` +- Create: `tests/CdrReconcilerTest.php` + +**Interfaces:** +- Produces: `CdrReconciler::due(array $records, int $now): array` and retry-result state transitions. +- Consumes: idempotent existing CallHistory insert/update behavior. + +- [ ] **Step 1: Write failing tests** for due-record selection, successful resolution, delayed retry and transition to manual review. +- [ ] **Step 2: Run** `php tests/CdrReconcilerTest.php`; expect failure. +- [ ] **Step 3: Implement** pure scheduling behavior and a thin infrastructure method in `ConnectorDB` that processes at most one due quarantine record per normal cycle. +- [ ] **Step 4: Add** a bounded overlap query before `committedOffset`; save through the same upsert path without moving the checkpoint backward. +- [ ] **Step 5: Ensure** a global source/DB failure leaves quarantine attempts unchanged. +- [ ] **Step 6: Run** all tests and syntax checks. +- [ ] **Step 7: Commit** with message `feat: reconcile quarantined CDR calls`. + +### Task 6: Deterministic trunk resolver + +**Files:** +- Create: `Lib/TrunkResolver.php` +- Modify: `Lib/GetReport.php` +- Create: `tests/TrunkResolverTest.php` + +**Interfaces:** +- Produces: `TrunkResolver::__construct(array $providers)` and `resolve(array $record, string $callType): array` returning `name`, `id`, `status`, `source`, `candidates`. +- Consumes provider entries with `uniqid`, `username`, `description` and normalized channel/line evidence. + +- [ ] **Step 1: Write failing tests** for exact line ID, inbound unique DID/username, peer fallback, outbound behavior, duplicate username ambiguity, unresolved technical value, and evidence priority. +- [ ] **Step 2: Run** `php tests/TrunkResolverTest.php`; expect failure. +- [ ] **Step 3: Implement** resolver indexes and conservative number normalization. +- [ ] **Step 4: Replace** provider maps in `GetReport::prepareCdrData()` with resolver calls while preserving sticky stronger evidence per `linkedid`. +- [ ] **Step 5: Log** ambiguous results without SIP credentials or full call payloads. +- [ ] **Step 6: Run** resolver and full standalone test suite. +- [ ] **Step 7: Commit** with message `fix: resolve CDR trunks deterministically`. + +### Task 7: Sync health API and UI state + +**Files:** +- Modify: `bin/ConnectorDB.php` +- Modify: `App/Controllers/ModuleExtendedCDRsController.php` +- Modify: `public/assets/js/src/module-export-records-index.js` +- Regenerate: `public/assets/js/module-export-records-index.js` +- Regenerate: `public/assets/js/module-export-records-index.js.map` +- Create: `tests/SyncStateTest.php` + +**Interfaces:** +- Produces cached state keys `offset`, `sourceLastId`, `lag`, `mode`, `lastSuccessAt`, `lastError`, `quarantinePending`. + +- [ ] **Step 1: Write failing state serialization tests** for current, catching-up, stale and failed states. +- [ ] **Step 2: Run** `php tests/SyncStateTest.php`; expect failure. +- [ ] **Step 3: Publish** the state after each source check and successful batch; do not overwrite `lastSuccessAt` on failure. +- [ ] **Step 4: Return** the state from `getStateAction()` with stable defaults for upgraded installations. +- [ ] **Step 5: Display** a non-blocking “history is catching up” status while leaving saved rows visible. +- [ ] **Step 6: Rebuild** frontend assets using the repository build command. +- [ ] **Step 7: Run** state tests, JS build and syntax checks. +- [ ] **Step 8: Commit** with message `feat: expose CDR synchronization health`. + +### Task 8: Local regression and package verification + +**Files:** +- Modify only if failures expose defects in files from Tasks 1–7. + +- [ ] **Step 1: Run** every `tests/*Test.php` with the available PHP interpreter. +- [ ] **Step 2: Run** `php -l` for every changed PHP file. +- [ ] **Step 3: Run** the existing module build and verify `git diff --check`. +- [ ] **Step 4: Confirm** unrelated `.DS_Store` files remain untracked and untouched. +- [ ] **Step 5: Record** test commands and exact results in the final handoff; do not push. + +### Task 9: Controlled verification on boffart.miko.ru + +**Files:** +- Remote backup under a timestamped `/tmp/ModuleExtendedCDRs-production-fix-*` directory. +- Remote installed module: `/storage/usbdisk1/mikopbx/custom_modules/ModuleExtendedCDRs` only after baseline and backup succeed. + +- [ ] **Step 1: Read-only baseline** over SSH: MikoPBX/PHP versions, installed module revision/files, worker PID, `cdrOffset`, source last ID, current lag, quarantine schema, free disk and recent ConnectorDB errors. +- [ ] **Step 2: Create** a timestamped backup of the installed module files and `db/cdr.db`; verify the backup before any replacement. +- [ ] **Step 3: Upload** a build artifact to a timestamped `/tmp` directory and run syntax/tests there before activation. +- [ ] **Step 4: Activate** files with the existing module-safe workflow, restart only `ConnectorDB`, and verify a single worker PID. +- [ ] **Step 5: Observe** offset/sourceLastId/lag, batch durations, error count and HTTP `getHistory` responses until lag reaches the normal threshold or a real blocker is identified. +- [ ] **Step 6: Exercise** representative inbound and outbound calls supplied by existing CDR data; compare resolver evidence and displayed trunk without exposing secrets. +- [ ] **Step 7: Restart** `ConnectorDB` once and verify checkpoint continuity, no duplicate rows and resumed catch-up. +- [ ] **Step 8: Roll back** immediately on fatal errors, growing lag, DB write failures, duplicate growth or broken history API; restore files/database from the verified backup. +- [ ] **Step 9: Report** remote evidence and leave the server either on the verified build or fully restored. Do not push repository changes. From 034e30e7b86c1b4666c95e1e873a569c3d43d501 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:06:27 +0300 Subject: [PATCH 03/18] feat: add adaptive CDR sync policy --- Lib/SyncPolicy.php | 61 ++++++++++++++++++++++++++++++++++++++++ tests/SyncPolicyTest.php | 44 +++++++++++++++++++++++++++++ 2 files changed, 105 insertions(+) create mode 100644 Lib/SyncPolicy.php create mode 100644 tests/SyncPolicyTest.php diff --git a/Lib/SyncPolicy.php b/Lib/SyncPolicy.php new file mode 100644 index 0000000..4cf4496 --- /dev/null +++ b/Lib/SyncPolicy.php @@ -0,0 +1,61 @@ + $lag, + 'mode' => self::MODE_ERROR, + 'delay' => self::ERROR_DELAY_SECONDS, + 'batchLinkedIds' => self::NORMAL_BATCH_LINKED_IDS, + ]; + } + + $catchUp = $wasCatchUp + ? $lag > self::CATCH_UP_EXIT_LAG + : $lag >= self::CATCH_UP_ENTER_LAG; + + $batchLinkedIds = $catchUp + ? self::CATCH_UP_BATCH_LINKED_IDS + : self::NORMAL_BATCH_LINKED_IDS; + if ($limitReached) { + $batchLinkedIds = self::MAX_BATCH_LINKED_IDS; + } + + return [ + 'lag' => $lag, + 'mode' => $catchUp ? self::MODE_CATCH_UP : self::MODE_NORMAL, + 'delay' => $catchUp ? 0 : self::NORMAL_DELAY_SECONDS, + 'batchLinkedIds' => min(self::MAX_BATCH_LINKED_IDS, $batchLinkedIds), + ]; + } +} diff --git a/tests/SyncPolicyTest.php b/tests/SyncPolicyTest.php new file mode 100644 index 0000000..2b1890c --- /dev/null +++ b/tests/SyncPolicyTest.php @@ -0,0 +1,44 @@ + Date: Wed, 22 Jul 2026 14:07:15 +0300 Subject: [PATCH 04/18] fix: expose reliable CDR batch status --- Lib/HistoryBatchResult.php | 42 ++++++++++++++++++++++++++++++++ Lib/HistoryParser.php | 33 ++++++++++++++++++++----- tests/HistoryBatchResultTest.php | 38 +++++++++++++++++++++++++++++ 3 files changed, 107 insertions(+), 6 deletions(-) create mode 100644 Lib/HistoryBatchResult.php create mode 100644 tests/HistoryBatchResultTest.php diff --git a/Lib/HistoryBatchResult.php b/Lib/HistoryBatchResult.php new file mode 100644 index 0000000..9b56fd4 --- /dev/null +++ b/Lib/HistoryBatchResult.php @@ -0,0 +1,42 @@ + $data + * @return array + */ + 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'), + ]; + } +} diff --git a/Lib/HistoryParser.php b/Lib/HistoryParser.php index 5248e37..b0400b5 100644 --- a/Lib/HistoryParser.php +++ b/Lib/HistoryParser.php @@ -116,7 +116,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:", @@ -144,7 +148,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, ]; @@ -155,7 +159,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(); @@ -268,7 +282,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); } @@ -285,7 +299,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 + ); } /** @@ -376,4 +397,4 @@ public static function getMinCdrId():int } return $id; } -} \ No newline at end of file +} diff --git a/tests/HistoryBatchResultTest.php b/tests/HistoryBatchResultTest.php new file mode 100644 index 0000000..8b008a4 --- /dev/null +++ b/tests/HistoryBatchResultTest.php @@ -0,0 +1,38 @@ + ['rows' => [['id' => 44], ['id' => 47]]], + 'call-b' => ['rows' => [['id' => 51]]], +]; +$batch = HistoryBatchResult::make(42, $data, true, 2); +expectSame(44, $batch['minId'], 'minimum row id'); +expectSame(51, $batch['maxId'], 'maximum row id'); +expectSame(3, $batch['rowCount'], 'total row count'); +expectSame(2, $batch['linkedIdCount'], 'linked id count'); +expectSame(true, $batch['limitReached'], 'linked id limit signal'); + +echo "HistoryBatchResultTest: OK\n"; From d50a33abb4bcd8543febc2ca868238b94747b575 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:10:25 +0300 Subject: [PATCH 05/18] fix: make CDR checkpoint catch-up resilient --- Lib/CheckpointPolicy.php | 32 +++++++++ Lib/HistoryParser.php | 18 ++++- Lib/QuarantinePolicy.php | 34 ++++++++++ Models/OversizedLinkedIds.php | 44 ++++++++++++- bin/ConnectorDB.php | 117 ++++++++++++++++++++++----------- tests/CheckpointPolicyTest.php | 26 ++++++++ tests/QuarantinePolicyTest.php | 36 ++++++++++ 7 files changed, 264 insertions(+), 43 deletions(-) create mode 100644 Lib/CheckpointPolicy.php create mode 100644 Lib/QuarantinePolicy.php create mode 100644 tests/CheckpointPolicyTest.php create mode 100644 tests/QuarantinePolicyTest.php diff --git a/Lib/CheckpointPolicy.php b/Lib/CheckpointPolicy.php new file mode 100644 index 0000000..829d1bd --- /dev/null +++ b/Lib/CheckpointPolicy.php @@ -0,0 +1,32 @@ + $batch + */ + public static function nextOffset(array $batch): int + { + $oldOffset = (int)$batch['oldOffset']; + if (empty($batch['requestOk']) || empty($batch['saveOk']) || !empty($batch['newQuarantine'])) { + return $oldOffset; + } + + $parsedOffset = max($oldOffset, (int)$batch['parsedOffset']); + $ids = array_values(array_unique(array_map('intval', $batch['rowIds'] ?? []))); + sort($ids); + if (empty($ids) || $ids[0] > $oldOffset + 1) { + return $parsedOffset; + } + + $set = array_flip($ids); + $contiguous = $oldOffset; + while (isset($set[$contiguous + 1])) { + $contiguous++; + } + + return min($parsedOffset, $contiguous); + } +} diff --git a/Lib/HistoryParser.php b/Lib/HistoryParser.php index b0400b5..c53f925 100644 --- a/Lib/HistoryParser.php +++ b/Lib/HistoryParser.php @@ -71,7 +71,6 @@ 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); } } catch (\Throwable $e) { @@ -81,6 +80,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.'); } @@ -350,14 +350,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] ?? []]; } /** diff --git a/Lib/QuarantinePolicy.php b/Lib/QuarantinePolicy.php new file mode 100644 index 0000000..7af296a --- /dev/null +++ b/Lib/QuarantinePolicy.php @@ -0,0 +1,34 @@ + */ + public static function failed(array $current, string $reason, int $now): array + { + $attempts = (int)($current['attempts'] ?? 0) + 1; + $delay = min(self::MAX_RETRY_SECONDS, self::BASE_RETRY_SECONDS * (2 ** ($attempts - 1))); + return [ + 'reason' => $reason, + 'attempts' => $attempts, + 'firstFailureAt' => (int)($current['firstFailureAt'] ?? $now), + 'lastFailureAt' => $now, + 'nextRetryAt' => $now + $delay, + 'status' => $attempts >= self::MAX_ATTEMPTS ? 'manual' : 'pending', + ]; + } + + /** @return array */ + public static function resolved(array $current, int $now): array + { + $current['status'] = 'resolved'; + $current['nextRetryAt'] = 0; + $current['lastFailureAt'] = $now; + return $current; + } +} diff --git a/Models/OversizedLinkedIds.php b/Models/OversizedLinkedIds.php index 82cbc85..f3d4794 100644 --- a/Models/OversizedLinkedIds.php +++ b/Models/OversizedLinkedIds.php @@ -61,6 +61,23 @@ class OversizedLinkedIds extends ModelsBase */ public ?string $detectedAt = ''; + /** @Column(type="integer", nullable=true) */ + public ?int $minId = 0; + /** @Column(type="integer", nullable=true) */ + public ?int $maxRangeId = 0; + /** @Column(type="string", nullable=true) */ + public ?string $reason = 'row_limit'; + /** @Column(type="integer", nullable=true) */ + public ?int $attempts = 0; + /** @Column(type="string", nullable=true) */ + public ?string $firstFailureAt = ''; + /** @Column(type="string", nullable=true) */ + public ?string $lastFailureAt = ''; + /** @Column(type="string", nullable=true) */ + public ?string $nextRetryAt = ''; + /** @Column(type="string", nullable=true) */ + public ?string $status = 'pending'; + /** * Создаёт служебную таблицу oversized_linkedids, если её ещё нет. * Единый источник схемы: вызывается воркерами ConnectorDB и SyncRecords при старте. @@ -81,8 +98,33 @@ public static function ensureTableExists(): void linkedid TEXT NOT NULL UNIQUE, rowCount INTEGER DEFAULT 0, maxId INTEGER DEFAULT 0, - detectedAt TEXT DEFAULT '' + detectedAt TEXT DEFAULT '', + minId INTEGER DEFAULT 0, + maxRangeId INTEGER DEFAULT 0, + reason TEXT DEFAULT 'row_limit', + attempts INTEGER DEFAULT 0, + firstFailureAt TEXT DEFAULT '', + lastFailureAt TEXT DEFAULT '', + nextRetryAt TEXT DEFAULT '', + status TEXT DEFAULT 'pending' )"); + $columns = $db->fetchAll('PRAGMA table_info(oversized_linkedids)'); + $existing = array_column($columns, 'name'); + $additions = [ + 'minId' => 'INTEGER DEFAULT 0', + 'maxRangeId' => 'INTEGER DEFAULT 0', + 'reason' => "TEXT DEFAULT 'row_limit'", + 'attempts' => 'INTEGER DEFAULT 0', + 'firstFailureAt' => "TEXT DEFAULT ''", + 'lastFailureAt' => "TEXT DEFAULT ''", + 'nextRetryAt' => "TEXT DEFAULT ''", + 'status' => "TEXT DEFAULT 'pending'", + ]; + foreach ($additions as $name => $definition) { + if (!in_array($name, $existing, true)) { + $db->execute("ALTER TABLE oversized_linkedids ADD COLUMN $name $definition"); + } + } } /** diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index c236ea3..fc22579 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -28,6 +28,8 @@ use Modules\ModuleExtendedCDRs\Lib\Logger; use Modules\ModuleExtendedCDRs\Lib\Mp3TagService; use Modules\ModuleExtendedCDRs\Lib\CdrQueryBuilder; +use Modules\ModuleExtendedCDRs\Lib\CheckpointPolicy; +use Modules\ModuleExtendedCDRs\Lib\SyncPolicy; use Exception; use Modules\ModuleExtendedCDRs\Lib\MikoPBXVersion; use Modules\ModuleExtendedCDRs\Lib\Providers\CdrDbProvider; @@ -51,6 +53,8 @@ class ConnectorDB extends WorkerBase public string $referenceDate = ''; private int $lastSyncTime = 0; + private int $nextSyncDelay = SyncPolicy::NORMAL_DELAY_SECONDS; + private bool $catchUpMode = false; private Mp3TagService $mp3TagService; /** @var string[] Кэш списка "раздутых" linkedid, исключённых из синхронизации. */ @@ -107,12 +111,13 @@ public function start($argv):void $beanstalk->subscribe($this->makePingTubeName(self::class), [$this, 'pingCallBack']); while ($this->needRestart === false) { try { - $this->syncCdrData(); + $this->syncCdrData(true); $this->pruneOversizedLinkedIds(); }catch (Throwable $exception){ $this->logger->writeError("Throwable:".$exception->getMessage(). ' Line: '.$exception->getLine()); + $this->nextSyncDelay = SyncPolicy::ERROR_DELAY_SECONDS; } - $beanstalk->wait(10); + $beanstalk->wait(max(1, $this->nextSyncDelay)); $this->logger->rotate(); } } @@ -306,7 +311,28 @@ public function syncCdrData(bool $force = false):void $oldOffset = $this->cdrOffset; $this->logger->writeInfo('...Start sync with offset...'. $oldOffset); - $historyResult = HistoryParser::getHistoryData($this->cdrOffset, $this->loadOversizedLinkedIds()); + $sourceState = HistoryParser::getLastCdrState(); + $sourceLastId = (int)($sourceState['data']['id'] ?? $oldOffset); + $policy = SyncPolicy::decide($oldOffset, $sourceLastId, $sourceState['ok'], false, $this->catchUpMode); + $this->nextSyncDelay = $policy['delay']; + $this->catchUpMode = $policy['mode'] === SyncPolicy::MODE_CATCH_UP; + if (!$sourceState['ok']) { + $this->publishSyncState($oldOffset, $sourceLastId, $policy, 'source_last_id_failed'); + $this->logger->writeError('batch_failed: source_last_id_failed'); + return; + } + + $historyResult = HistoryParser::getHistoryData( + $this->cdrOffset, + $this->loadOversizedLinkedIds(), + $policy['batchLinkedIds'] + ); + if (!$historyResult['ok']) { + $this->nextSyncDelay = SyncPolicy::ERROR_DELAY_SECONDS; + $this->publishSyncState($oldOffset, $sourceLastId, $policy, $historyResult['error']); + $this->logger->writeError('batch_failed: ' . $historyResult['error']); + return; + } $cdrData = $historyResult['data']; $parsedOffset = $historyResult['newOffset']; $totalRows = array_sum(array_map(fn($cdr) => count($cdr['rows'] ?? []), $cdrData)); @@ -501,39 +527,14 @@ public function syncCdrData(bool $force = false):void return; } - // Применяем offset от getHistoryData после успешного сохранения - $this->cdrOffset = $parsedOffset; - - // Уточняем offset: ищем максимальный последовательный id начиная от oldOffset - if (!empty($allRowIds)) { - $allRowIds = array_unique($allRowIds); - sort($allRowIds); - $allRowIds = array_values($allRowIds); - - // Создаём set для быстрой проверки O(1) - $idSet = array_flip($allRowIds); - - // Ищем максимальный последовательный id начиная от oldOffset+1 - $sequentialMaxId = $oldOffset; - $nextExpected = $oldOffset + 1; - - while (isset($idSet[$nextExpected])) { - $sequentialMaxId = $nextExpected; - $nextExpected++; - } - - if ($sequentialMaxId > $oldOffset) { - // Берём максимум между sequential и parsed offset - $newOffset = max($sequentialMaxId, $this->cdrOffset); - $this->logger->writeInfo("Adjusting offset from {$this->cdrOffset} to $newOffset (sequential from $oldOffset to $sequentialMaxId)"); - $this->cdrOffset = $newOffset; - } else { - // Нет последовательных id от текущего offset — разрыв - $minId = min($allRowIds); - $maxId = max($allRowIds); - $this->logger->writeInfo("Gap detected: offset={$this->cdrOffset}, minId=$minId, maxId=$maxId, count=" . count($allRowIds)); - } - } + $this->cdrOffset = CheckpointPolicy::nextOffset([ + 'oldOffset' => $oldOffset, + 'parsedOffset' => $parsedOffset, + 'requestOk' => true, + 'saveOk' => true, + 'newQuarantine' => false, + 'rowIds' => $allRowIds, + ]); $this->logger->writeInfo([ 'CallHistoryFindTime' => round($CallHistoryFindTime, 4), @@ -546,7 +547,7 @@ public function syncCdrData(bool $force = false):void "Timing"); if($oldOffset !== $this->cdrOffset){ $this->logger->writeInfo("Update progress, offset $oldOffset to new value $this->cdrOffset "); - $lastCdrData = HistoryParser::getLastCdrData(); + $lastCdrData = $sourceState['data']; if(!empty($lastCdrData)){ $tmpCdrData = [ 'lastId' => intval($lastCdrData['id']), @@ -557,10 +558,34 @@ public function syncCdrData(bool $force = false):void } $this->updateSettings($this->cdrOffset); } + $policy = SyncPolicy::decide( + $this->cdrOffset, + $sourceLastId, + true, + $historyResult['limitReached'], + $this->catchUpMode + ); + $this->nextSyncDelay = $policy['delay']; + $this->catchUpMode = $policy['mode'] === SyncPolicy::MODE_CATCH_UP; + $this->publishSyncState($this->cdrOffset, $sourceLastId, $policy, ''); $offsetDelta = $this->cdrOffset - $oldOffset; $this->logger->writeInfo("End sync with offset {$this->cdrOffset} (+$offsetDelta)"); } + private function publishSyncState(int $offset, int $sourceLastId, array $policy, string $error): void + { + CacheManager::setCacheData(HistoryParser::CDR_SYNC_PROGRESS_KEY, [ + 'lastId' => $sourceLastId, + 'nowId' => $offset, + 'offset' => $offset, + 'sourceLastId' => $sourceLastId, + 'lag' => max(0, $sourceLastId - $offset), + 'mode' => $policy['mode'], + 'lastSuccessAt' => $error === '' ? date('c') : '', + 'lastError' => $error, + ]); + } + /** * Возвращает путь к файлу записи по ID. * @param string $id @@ -996,7 +1021,10 @@ private function loadOversizedLinkedIds(): array { if (time() - $this->oversizedCacheTime > 60) { try { - $rows = OversizedLinkedIds::find(['columns' => 'linkedid']); + $rows = OversizedLinkedIds::find([ + "status IS NULL OR status <> 'resolved'", + 'columns' => 'linkedid' + ]); $dbList = array_column($rows->toArray(), 'linkedid'); // Кэш = записи из БД ∪ session-only (не потерянные при неуспешном save()). $this->oversizedCache = array_values(array_unique(array_merge($dbList, $this->oversizedPending))); @@ -1023,9 +1051,12 @@ private function persistOversizedLinkedIds(array $linkedIds, array $cdrData): vo } $rows = $cdrData[$linkedId]['rows'] ?? []; $rowCount = count($rows); + $minId = 0; $maxId = 0; foreach ($rows as $row) { - $maxId = max($maxId, (int)($row['id'] ?? 0)); + $rowId = (int)($row['id'] ?? 0); + $maxId = max($maxId, $rowId); + $minId = $minId === 0 ? $rowId : min($minId, $rowId); } $record = new OversizedLinkedIds(); @@ -1033,6 +1064,14 @@ private function persistOversizedLinkedIds(array $linkedIds, array $cdrData): vo $record->rowCount = $rowCount; $record->maxId = $maxId; $record->detectedAt = date('Y-m-d H:i:s'); + $record->minId = $minId; + $record->maxRangeId = $maxId; + $record->reason = 'row_limit'; + $record->attempts = 0; + $record->firstFailureAt = $record->detectedAt; + $record->lastFailureAt = $record->detectedAt; + $record->nextRetryAt = date('Y-m-d H:i:s', time() + 60); + $record->status = 'pending'; $saved = $record->save(); // В любом случае исключаем linkedid в пределах текущей сессии воркера, diff --git a/tests/CheckpointPolicyTest.php b/tests/CheckpointPolicyTest.php new file mode 100644 index 0000000..57f4f17 --- /dev/null +++ b/tests/CheckpointPolicyTest.php @@ -0,0 +1,26 @@ + 100, 'parsedOffset' => 105, 'requestOk' => true, + 'saveOk' => true, 'newQuarantine' => false, 'rowIds' => [101, 102, 103, 104, 105]]; +checkPointExpected(105, $base, 'successful contiguous batch advances'); +checkPointExpected(100, array_merge($base, ['requestOk' => false]), 'source error retains offset'); +checkPointExpected(100, array_merge($base, ['saveOk' => false]), 'save error retains offset'); +checkPointExpected(100, array_merge($base, ['newQuarantine' => true]), 'new quarantine retains offset for replay'); +checkPointExpected(102, array_merge($base, ['rowIds' => [101, 102, 104, 105]]), 'gap advances only contiguous prefix'); +checkPointExpected(105, array_merge($base, ['rowIds' => [102, 105]]), 'parser offset is authoritative when source groups omit unrelated IDs'); + +echo "CheckpointPolicyTest: OK\n"; diff --git a/tests/QuarantinePolicyTest.php b/tests/QuarantinePolicyTest.php new file mode 100644 index 0000000..5e8f068 --- /dev/null +++ b/tests/QuarantinePolicyTest.php @@ -0,0 +1,36 @@ + $first['nextRetryAt'], 'backoff grows'); + +$state = $second; +for ($attempt = 2; $attempt < QuarantinePolicy::MAX_ATTEMPTS; $attempt++) { + $state = QuarantinePolicy::failed($state, 'row_limit', 1100 + $attempt); +} +quarantineAssert('manual', $state['status'], 'maximum attempts require manual review'); + +$resolved = QuarantinePolicy::resolved($second, 2000); +quarantineAssert('resolved', $resolved['status'], 'successful retry resolves'); +quarantineAssert(2000, $resolved['lastFailureAt'], 'resolution timestamp retained for audit'); + +echo "QuarantinePolicyTest: OK\n"; From 94c3f2d81373947bf0cc2209f25499d88897eca7 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:11:51 +0300 Subject: [PATCH 06/18] fix: resolve CDR trunks deterministically --- Lib/GetReport.php | 52 +++++++++--------------- Lib/TrunkResolver.php | 79 +++++++++++++++++++++++++++++++++++++ tests/TrunkResolverTest.php | 44 +++++++++++++++++++++ 3 files changed, 141 insertions(+), 34 deletions(-) create mode 100644 Lib/TrunkResolver.php create mode 100644 tests/TrunkResolverTest.php diff --git a/Lib/GetReport.php b/Lib/GetReport.php index 80046b4..ead95e7 100644 --- a/Lib/GetReport.php +++ b/Lib/GetReport.php @@ -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 = [ @@ -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'; @@ -1329,4 +1313,4 @@ public function getDateRanges(string $srcPeriod): string ]; return $dateRanges[$srcPeriod] ?? $srcPeriod; } -} \ No newline at end of file +} diff --git a/Lib/TrunkResolver.php b/Lib/TrunkResolver.php new file mode 100644 index 0000000..311ac49 --- /dev/null +++ b/Lib/TrunkResolver.php @@ -0,0 +1,79 @@ + */ + private array $byId = []; + /** @var array> */ + private array $byUsername = []; + + public function __construct(iterable $providers) + { + foreach ($providers as $provider) { + $provider = (array)$provider; + $id = (string)($provider['uniqid'] ?? ''); + $name = (string)($provider['description'] ?? ''); + if ($id === '' || $name === '') { + continue; + } + $candidate = ['name' => $name, 'id' => $id]; + $this->byId[$id] = $candidate; + $username = self::normalizeNumber((string)($provider['username'] ?? '')); + if ($username !== '') { + $this->byUsername[$username][] = $candidate; + } + } + } + + /** @return array{name:string,id:string,status:string,source:string,candidates:array} */ + public function resolve(array $record, string $callType): array + { + $technical = (string)($record['line'] ?? ''); + if (isset($this->byId[$technical])) { + return $this->resolved($this->byId[$technical], 'line_id'); + } + + if ($callType === 'incoming' || $callType === '2') { + $did = self::normalizeNumber((string)($record['did'] ?? '')); + $candidates = $did === '' ? [] : ($this->byUsername[$did] ?? []); + if (count($candidates) === 1) { + return $this->resolved($candidates[0], 'did_username'); + } + if (count($candidates) > 1) { + return [ + 'name' => $technical, + 'id' => $technical, + 'status' => 'ambiguous', + 'source' => 'did_username', + 'candidates' => array_column($candidates, 'id'), + ]; + } + } + + return [ + 'name' => $technical, + 'id' => $technical, + 'status' => 'unresolved', + 'source' => 'technical', + 'candidates' => [], + ]; + } + + private function resolved(array $candidate, string $source): array + { + return [ + 'name' => $candidate['name'], + 'id' => $candidate['id'], + 'status' => 'resolved', + 'source' => $source, + 'candidates' => [], + ]; + } + + private static function normalizeNumber(string $number): string + { + return preg_replace('/\D+/', '', $number) ?: ''; + } +} diff --git a/tests/TrunkResolverTest.php b/tests/TrunkResolverTest.php new file mode 100644 index 0000000..48bea0c --- /dev/null +++ b/tests/TrunkResolverTest.php @@ -0,0 +1,44 @@ + 'SIP-A', 'username' => '+7 (495) 111-22-33', 'description' => 'Provider A'], + ['uniqid' => 'SIP-B', 'username' => '74952223344', 'description' => 'Provider B'], +]); + +$line = $resolver->resolve(['line' => 'SIP-A', 'did' => '74952223344'], 'incoming'); +trunkAssert('Provider A', $line['name'], 'stable line id has highest priority'); +trunkAssert('line_id', $line['source'], 'line id evidence source'); + +$did = $resolver->resolve(['line' => 'unknown-peer', 'did' => '7 495 222-33-44'], 'incoming'); +trunkAssert('Provider B', $did['name'], 'incoming unique DID resolves provider'); +trunkAssert('did_username', $did['source'], 'DID evidence source'); +$storedType = $resolver->resolve(['line' => 'unknown-peer', 'did' => '74952223344'], '2'); +trunkAssert('Provider B', $storedType['name'], 'stored incoming call type resolves DID'); + +$outgoing = $resolver->resolve(['line' => 'unknown-peer', 'did' => '74952223344'], 'outgoing'); +trunkAssert('unresolved', $outgoing['status'], 'outgoing call does not use DID'); + +$ambiguous = new TrunkResolver([ + ['uniqid' => 'SIP-A', 'username' => '100500', 'description' => 'Provider A'], + ['uniqid' => 'SIP-B', 'username' => '100500', 'description' => 'Provider B'], +]); +$result = $ambiguous->resolve(['line' => 'shared-peer', 'did' => '100500'], 'incoming'); +trunkAssert('ambiguous', $result['status'], 'duplicate usernames are ambiguous'); +trunkAssert('shared-peer', $result['name'], 'ambiguous result preserves technical value'); +trunkAssert(2, count($result['candidates']), 'ambiguous candidates exposed'); + +echo "TrunkResolverTest: OK\n"; From fed8bfee93cc59bed56f62b1d8e6b4b72dc0a263 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:13:11 +0300 Subject: [PATCH 07/18] feat: expose CDR synchronization health --- bin/ConnectorDB.php | 5 ++++- 1 file changed, 4 insertions(+), 1 deletion(-) diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index fc22579..9ae6d97 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -574,6 +574,8 @@ public function syncCdrData(bool $force = false):void private function publishSyncState(int $offset, int $sourceLastId, array $policy, string $error): void { + $previous = CacheManager::getCacheData(HistoryParser::CDR_SYNC_PROGRESS_KEY); + $previous = is_array($previous) ? $previous : []; CacheManager::setCacheData(HistoryParser::CDR_SYNC_PROGRESS_KEY, [ 'lastId' => $sourceLastId, 'nowId' => $offset, @@ -581,7 +583,8 @@ private function publishSyncState(int $offset, int $sourceLastId, array $policy, 'sourceLastId' => $sourceLastId, 'lag' => max(0, $sourceLastId - $offset), 'mode' => $policy['mode'], - 'lastSuccessAt' => $error === '' ? date('c') : '', + 'lastDate' => $previous['lastDate'] ?? '', + 'lastSuccessAt' => $error === '' ? date('c') : ($previous['lastSuccessAt'] ?? ''), 'lastError' => $error, ]); } From b17c30ede3b25097f12771d663dc0641b3efbb5e Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:18:17 +0300 Subject: [PATCH 08/18] fix: accept successful empty CDR responses --- Lib/HistoryParser.php | 5 +++++ 1 file changed, 5 insertions(+) diff --git a/Lib/HistoryParser.php b/Lib/HistoryParser.php index c53f925..4d78247 100644 --- a/Lib/HistoryParser.php +++ b/Lib/HistoryParser.php @@ -72,6 +72,8 @@ public static function getCdr(array $filter = [], ?bool &$requestOk = null): arr [$result, $message] = $client->sendRequest(json_encode($filter), 30); if ($result!==false){ $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 = ''; @@ -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; From bc6056af3f5aab3320c2d936ffdb592a21d86171 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:22:37 +0300 Subject: [PATCH 09/18] fix: advance catch-up across filtered CDR gaps --- Lib/CheckpointPolicy.php | 17 +++-------------- bin/ConnectorDB.php | 9 --------- tests/CheckpointPolicyTest.php | 2 +- 3 files changed, 4 insertions(+), 24 deletions(-) diff --git a/Lib/CheckpointPolicy.php b/Lib/CheckpointPolicy.php index 829d1bd..47833a8 100644 --- a/Lib/CheckpointPolicy.php +++ b/Lib/CheckpointPolicy.php @@ -14,19 +14,8 @@ public static function nextOffset(array $batch): int return $oldOffset; } - $parsedOffset = max($oldOffset, (int)$batch['parsedOffset']); - $ids = array_values(array_unique(array_map('intval', $batch['rowIds'] ?? []))); - sort($ids); - if (empty($ids) || $ids[0] > $oldOffset + 1) { - return $parsedOffset; - } - - $set = array_flip($ids); - $contiguous = $oldOffset; - while (isset($set[$contiguous + 1])) { - $contiguous++; - } - - return min($parsedOffset, $contiguous); + // 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']); } } diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index 9ae6d97..0312db4 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -547,15 +547,6 @@ public function syncCdrData(bool $force = false):void "Timing"); if($oldOffset !== $this->cdrOffset){ $this->logger->writeInfo("Update progress, offset $oldOffset to new value $this->cdrOffset "); - $lastCdrData = $sourceState['data']; - if(!empty($lastCdrData)){ - $tmpCdrData = [ - 'lastId' => intval($lastCdrData['id']), - 'lastDate' => $lastCdrData['start'], - 'nowId' => $this->cdrOffset - ]; - CacheManager::setCacheData(HistoryParser::CDR_SYNC_PROGRESS_KEY, $tmpCdrData); - } $this->updateSettings($this->cdrOffset); } $policy = SyncPolicy::decide( diff --git a/tests/CheckpointPolicyTest.php b/tests/CheckpointPolicyTest.php index 57f4f17..c3e21fc 100644 --- a/tests/CheckpointPolicyTest.php +++ b/tests/CheckpointPolicyTest.php @@ -20,7 +20,7 @@ function checkPointExpected(int $expected, array $input, string $message): void checkPointExpected(100, array_merge($base, ['requestOk' => false]), 'source error retains offset'); checkPointExpected(100, array_merge($base, ['saveOk' => false]), 'save error retains offset'); checkPointExpected(100, array_merge($base, ['newQuarantine' => true]), 'new quarantine retains offset for replay'); -checkPointExpected(102, array_merge($base, ['rowIds' => [101, 102, 104, 105]]), 'gap advances only contiguous prefix'); +checkPointExpected(105, array_merge($base, ['rowIds' => [101, 102, 104, 105]]), 'filtered ID gaps do not block parser checkpoint'); checkPointExpected(105, array_merge($base, ['rowIds' => [102, 105]]), 'parser offset is authoritative when source groups omit unrelated IDs'); echo "CheckpointPolicyTest: OK\n"; From cdf46fba535f5159109483d2096e0575ff014100 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:24:32 +0300 Subject: [PATCH 10/18] fix: preserve CDR batch metadata during catch-up --- bin/ConnectorDB.php | 25 +++++++++++++++---------- 1 file changed, 15 insertions(+), 10 deletions(-) diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index 0312db4..2e012d3 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -114,7 +114,12 @@ public function start($argv):void $this->syncCdrData(true); $this->pruneOversizedLinkedIds(); }catch (Throwable $exception){ - $this->logger->writeError("Throwable:".$exception->getMessage(). ' Line: '.$exception->getLine()); + $this->logger->writeError( + "Throwable:" . $exception->getMessage() + . ' File:' . $exception->getFile() + . ' Line:' . $exception->getLine() + . ' Trace:' . $exception->getTraceAsString() + ); $this->nextSyncDelay = SyncPolicy::ERROR_DELAY_SECONDS; } $beanstalk->wait(max(1, $this->nextSyncDelay)); @@ -322,19 +327,19 @@ public function syncCdrData(bool $force = false):void return; } - $historyResult = HistoryParser::getHistoryData( + $batchResult = HistoryParser::getHistoryData( $this->cdrOffset, $this->loadOversizedLinkedIds(), $policy['batchLinkedIds'] ); - if (!$historyResult['ok']) { + if (!$batchResult['ok']) { $this->nextSyncDelay = SyncPolicy::ERROR_DELAY_SECONDS; - $this->publishSyncState($oldOffset, $sourceLastId, $policy, $historyResult['error']); - $this->logger->writeError('batch_failed: ' . $historyResult['error']); + $this->publishSyncState($oldOffset, $sourceLastId, $policy, $batchResult['error']); + $this->logger->writeError('batch_failed: ' . $batchResult['error']); return; } - $cdrData = $historyResult['data']; - $parsedOffset = $historyResult['newOffset']; + $cdrData = $batchResult['data']; + $parsedOffset = $batchResult['newOffset']; $totalRows = array_sum(array_map(fn($cdr) => count($cdr['rows'] ?? []), $cdrData)); $this->logger->writeInfo("Parsed offset $parsedOffset. linkedIds:" . count($cdrData) . ", totalRows:$totalRows"); @@ -408,11 +413,11 @@ public function syncCdrData(bool $force = false):void $start = microtime(true); $existingHistory = []; if (!empty($normalLinkedIds)) { - $historyResult = CallHistory::find([ + $historyRecords = CallHistory::find([ 'linkedid IN ({ids:array})', 'bind' => ['ids' => $normalLinkedIds] ]); - foreach ($historyResult as $h) { + foreach ($historyRecords as $h) { $existingHistory[$h->UNIQUEID] = $h; } } @@ -553,7 +558,7 @@ public function syncCdrData(bool $force = false):void $this->cdrOffset, $sourceLastId, true, - $historyResult['limitReached'], + $batchResult['limitReached'], $this->catchUpMode ); $this->nextSyncDelay = $policy['delay']; From 53dbd0adb68c73b4ad4993bb2257401b8409b921 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:44:02 +0300 Subject: [PATCH 11/18] docs: design atomic CDR batch persistence --- .../2026-07-22-atomic-cdr-batch-design.md | 86 +++++++++++++++++++ 1 file changed, 86 insertions(+) create mode 100644 docs/superpowers/specs/2026-07-22-atomic-cdr-batch-design.md diff --git a/docs/superpowers/specs/2026-07-22-atomic-cdr-batch-design.md b/docs/superpowers/specs/2026-07-22-atomic-cdr-batch-design.md new file mode 100644 index 0000000..23c1ce6 --- /dev/null +++ b/docs/superpowers/specs/2026-07-22-atomic-cdr-batch-design.md @@ -0,0 +1,86 @@ +# Atomic CDR Batch Persistence Design + +## Goal + +Prevent permanent CDR loss when any database write in a synchronization batch fails. A batch may advance the persisted CDR offset only after every required write has committed successfully. + +## Scope + +This change covers writes performed by `ConnectorDB::syncCdrData()`: + +- `cdr_general` inserts and updates; +- queue-history writes required by the same source batch; +- recall/transfer state changes derived from the saved rows; +- creation of an oversized-linkedid quarantine record; +- persistence of the new synchronization offset after the CDR transaction commits. + +A background reconciler for oversized calls is outside this change. Quarantine records must remain available for manual analysis and must not be automatically deleted. + +## Transaction Boundary + +All module CDR tables involved in a batch use the module CDR SQLite connection. `ConnectorDB` opens one transaction on that connection before the first mutation. + +Within the transaction it: + +1. saves queue-history changes; +2. inserts or updates call-history rows; +3. updates recall/transfer state; +4. saves quarantine metadata for newly oversized linked IDs. + +Every model save must be checked for a `false` result as well as exceptions. Batch INSERT helpers must propagate failures instead of logging and continuing. + +If any required write fails, the transaction rolls back, the in-memory and persisted offsets remain unchanged, and the worker retries from the same source boundary after the error delay. + +After a successful database commit, `ConnectorDB` computes and persists the new offset. Offset persistence failure is treated as a synchronization failure: the CDR rows may already be committed, but replay is safe because `UNIQUEID` processing is idempotent. The next run starts from the old persisted offset and updates existing rows rather than creating duplicates. + +## Oversized Calls + +An oversized linked ID is not added to the active exclusion cache before its available CDR rows and quarantine record commit successfully. + +On a successful commit: + +- its available rows are durable; +- its quarantine record is durable; +- it is added to the exclusion cache; +- the offset remains at the prior value for one cycle so the next query can exclude it and process ordinary calls behind it. + +On failure, neither the transaction nor the exclusion becomes active. The same linked ID is retried from the unchanged offset. + +Quarantine audit records are retained. Automatic pruning/deletion is disabled until a separate reconciler can mark records resolved with an auditable outcome. + +## Failure State and Logging + +Each batch produces a stable outcome event containing only operational metadata: + +- old offset and proposed offset; +- source last ID and lag; +- parsed minimum/maximum ID; +- linked-ID and row counts; +- insert/update counts; +- mode and elapsed time; +- outcome (`committed`, `rolled_back`, `offset_persist_failed`, `quarantined`); +- normalized error category. + +Raw phone numbers and raw linked IDs are not included in normal batch events. Exception details remain in the error log, while the health state receives the normalized category and retains the previous successful-progress timestamp. + +## Testing + +Tests must demonstrate failure before implementation and cover: + +1. a later INSERT chunk fails after an earlier chunk succeeded: transaction rolls back and offset remains unchanged; +2. an UPDATE save returns `false`: transaction rolls back and offset remains unchanged; +3. queue-history or state save returns `false`: transaction rolls back; +4. quarantine persistence fails: no exclusion is activated and offset remains unchanged; +5. successful oversized batch: rows and quarantine commit, exclusion activates, offset is held for the next cycle; +6. restart/replay after a failed offset persistence produces no duplicate `UNIQUEID` rows; +7. successful ordinary batch commits and advances the offset. + +Local policy tests, PHP syntax checks, and `git diff --check` remain mandatory. Server verification uses a database backup and a controlled offset rewind. A destructive database fault is simulated only against the test server and must be restored immediately. + +## Release Criteria + +- No code path advances the persisted offset after a failed required write. +- No oversized linked ID is excluded before its rows and quarantine metadata are durable. +- Failed batches are distinguishable from successful empty batches in both log and health state. +- Replaying a committed batch is idempotent. +- Existing history and trunk-resolution tests continue to pass. From 628d17b7fa308d88731bda1e598a830d1d4aa8bf Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 15:05:46 +0300 Subject: [PATCH 12/18] docs: plan atomic CDR batch implementation --- ...6-07-22-atomic-cdr-batch-implementation.md | 223 ++++++++++++++++++ 1 file changed, 223 insertions(+) create mode 100644 docs/superpowers/plans/2026-07-22-atomic-cdr-batch-implementation.md diff --git a/docs/superpowers/plans/2026-07-22-atomic-cdr-batch-implementation.md b/docs/superpowers/plans/2026-07-22-atomic-cdr-batch-implementation.md new file mode 100644 index 0000000..0e47fb7 --- /dev/null +++ b/docs/superpowers/plans/2026-07-22-atomic-cdr-batch-implementation.md @@ -0,0 +1,223 @@ +# Atomic CDR Batch Persistence Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Guarantee that CDR synchronization never advances its checkpoint after a failed required database write and never activates oversized-call exclusion before durable commit. + +**Architecture:** Add a small transaction coordinator with a testable result contract, then route all CDR, queue, state, and quarantine mutations through one shared `CdrDbProvider` transaction. Write failures throw into the coordinator, which rolls back and returns a normalized failure; cache and offset mutations occur only after commit. + +**Tech Stack:** PHP 7.4, Phalcon models/database adapter, SQLite, standalone PHP regression tests. + +## Global Constraints + +- Work in the current `develop` branch and do not push. +- Keep module settings/offset persistence outside the CDR transaction because it uses `module.db`; replay after offset-save failure must remain idempotent. +- Do not delete quarantine audit records automatically. +- Do not log raw phone numbers or raw linked IDs in normal batch events. +- Every production behavior change requires a failing test first. + +--- + +### Task 1: Transaction coordinator + +**Files:** +- Create: `Lib/AtomicBatch.php` +- Test: `tests/AtomicBatchTest.php` + +**Interfaces:** +- Consumes: a Phalcon-compatible adapter exposing `begin()`, `commit()`, and `rollback()`. +- Produces: `AtomicBatch::run(object $db, callable $operation): array{ok:bool,value:mixed,error:string}`. + +- [ ] **Step 1: Write the failing test** + +Create a fake adapter that records transaction calls. Assert success calls `begin,commit`; an exception calls `begin,rollback`; false `begin` and false `commit` return `ok=false`; and a failed commit attempts rollback. + +- [ ] **Step 2: Run test to verify RED** + +Run: `php tests/AtomicBatchTest.php` +Expected: FAIL because `Lib/AtomicBatch.php` does not exist. + +- [ ] **Step 3: Implement the coordinator** + +Implement `AtomicBatch::run()` so it validates each adapter result, executes the callback once, returns its value after commit, catches `Throwable`, rolls back only when a transaction began, and reports a normalized error without logging payload data. + +- [ ] **Step 4: Run test to verify GREEN** + +Run: `php tests/AtomicBatchTest.php` +Expected: `AtomicBatchTest: OK`. + +- [ ] **Step 5: Commit** + +```bash +git add Lib/AtomicBatch.php tests/AtomicBatchTest.php +git commit -m "feat: add atomic CDR batch coordinator" +``` + +### Task 2: Explicit persistence result and failure propagation + +**Files:** +- Create: `Lib/BatchPersistenceResult.php` +- Test: `tests/BatchPersistenceResultTest.php` +- Modify: `bin/ConnectorDB.php:474-568, 788-858` + +**Interfaces:** +- Produces: `BatchPersistenceResult::success(int $inserted, int $updated): array` and `BatchPersistenceResult::failure(string $category, string $message=''): array`. +- `batchSaveCallHistory()` returns the explicit result and throws on failed adapter execution or model `save() === false` while inside an atomic operation. + +- [ ] **Step 1: Write the failing result-contract test** + +Assert success contains `ok=true`, counts, and an empty error; failure contains `ok=false`, a stable category, zero counts, and a diagnostic message. + +- [ ] **Step 2: Run test to verify RED** + +Run: `php tests/BatchPersistenceResultTest.php` +Expected: FAIL because the result class does not exist. + +- [ ] **Step 3: Implement the result class and propagate failures** + +Remove catch-and-continue behavior from INSERT/UPDATE loops. Treat adapter `execute() !== true`, queue save false, state save false, and model validation errors as exceptions categorized as `insert_failed`, `update_failed`, `queue_save_failed`, or `state_save_failed`. + +- [ ] **Step 4: Run focused and existing tests** + +Run: `php tests/BatchPersistenceResultTest.php && for test_file in tests/*Test.php; do php "$test_file" || exit 1; done` +Expected: all tests print `OK`. + +- [ ] **Step 5: Commit** + +```bash +git add Lib/BatchPersistenceResult.php tests/BatchPersistenceResultTest.php bin/ConnectorDB.php +git commit -m "fix: propagate CDR persistence failures" +``` + +### Task 3: Make synchronization atomic + +**Files:** +- Modify: `bin/ConnectorDB.php:310-585` +- Test: `tests/AtomicBatchTest.php` +- Test: `tests/CheckpointPolicyTest.php` + +**Interfaces:** +- Consumes: `AtomicBatch::run()` and the explicit persistence result. +- Produces: no offset or success-health update unless the CDR transaction commits. + +- [ ] **Step 1: Add failing transaction/checkpoint scenarios** + +Extend tests to assert that an operation failing after an earlier mutation rolls back and that `CheckpointPolicy::nextOffset()` holds the old offset for `saveOk=false`. + +- [ ] **Step 2: Run tests to verify RED** + +Run: `php tests/AtomicBatchTest.php && php tests/CheckpointPolicyTest.php` +Expected: the new late-failure assertion fails before integration. + +- [ ] **Step 3: Wrap all required writes** + +Acquire the shared `CdrDbProvider` adapter, run queue/history/state/quarantine mutations inside `AtomicBatch::run()`, and on failure publish `batch_write_failed`, log a structured rollback event, set error delay, retain the old in-memory offset, and return. Compute and persist the next offset only after commit. + +- [ ] **Step 4: Verify focused tests and syntax** + +Run: `php tests/AtomicBatchTest.php && php tests/CheckpointPolicyTest.php && php -l bin/ConnectorDB.php` +Expected: tests pass and syntax is valid. + +- [ ] **Step 5: Commit** + +```bash +git add bin/ConnectorDB.php tests/AtomicBatchTest.php tests/CheckpointPolicyTest.php +git commit -m "fix: commit CDR batches atomically" +``` + +### Task 4: Commit-gated quarantine and retained audit + +**Files:** +- Create: `Lib/QuarantineActivation.php` +- Test: `tests/QuarantineActivationTest.php` +- Modify: `bin/ConnectorDB.php:1015-1143` + +**Interfaces:** +- Produces: `QuarantineActivation::afterCommit(array $current, array $committed): array` for deterministic cache activation. +- `persistOversizedLinkedIds()` persists metadata only and returns IDs eligible for activation after commit. + +- [ ] **Step 1: Write failing cache-activation tests** + +Assert failed/uncommitted IDs never enter the exclusion list; committed IDs enter once; existing IDs remain; raw IDs are not required in the result log context. + +- [ ] **Step 2: Run test to verify RED** + +Run: `php tests/QuarantineActivationTest.php` +Expected: FAIL because the activation class does not exist. + +- [ ] **Step 3: Implement commit-gated activation** + +Move all `$oversizedCache` and `$oversizedPending` mutations after a successful transaction. Make quarantine save false throw. Disable `pruneOversizedLinkedIds()` deletion; retain records until a separate reconciler marks an auditable resolution. + +- [ ] **Step 4: Run focused tests** + +Run: `php tests/QuarantineActivationTest.php && php tests/QuarantinePolicyTest.php` +Expected: both pass. + +- [ ] **Step 5: Commit** + +```bash +git add Lib/QuarantineActivation.php tests/QuarantineActivationTest.php bin/ConnectorDB.php +git commit -m "fix: activate CDR quarantine only after commit" +``` + +### Task 5: Production-safe batch diagnostics + +**Files:** +- Create: `Lib/BatchLogContext.php` +- Test: `tests/BatchLogContextTest.php` +- Modify: `bin/ConnectorDB.php:310-585` + +**Interfaces:** +- Produces: `BatchLogContext::make(array $state): array` with stable operational keys and no raw call identifiers. + +- [ ] **Step 1: Write the failing log-context test** + +Assert the context includes event, offsets, sourceLastId, lag, parsed range, counts, mode, elapsedMs, outcome, and errorCategory; assert phone, linkedid, UNIQUEID, and recording fields are absent. + +- [ ] **Step 2: Run test to verify RED** + +Run: `php tests/BatchLogContextTest.php` +Expected: FAIL because the context class does not exist. + +- [ ] **Step 3: Implement and emit one outcome event per batch** + +Emit `cdr_sync_batch` for committed, rolled-back, quarantined, and offset-persist-failed outcomes. Do not refresh `lastSuccessAt` on rollback. Include exception stack separately only for unexpected top-level failures. + +- [ ] **Step 4: Run focused tests** + +Run: `php tests/BatchLogContextTest.php && php -l bin/ConnectorDB.php` +Expected: pass. + +- [ ] **Step 5: Commit** + +```bash +git add Lib/BatchLogContext.php tests/BatchLogContextTest.php bin/ConnectorDB.php +git commit -m "feat: add safe CDR batch diagnostics" +``` + +### Task 6: Full verification and test-server fault exercise + +**Files:** +- Modify only if a test exposes a defect in files already listed above. + +- [ ] **Step 1: Run the complete local suite** + +Run every `tests/*Test.php`, lint all changed PHP files, and run `git diff --check`. +Expected: all tests pass, no syntax errors, no whitespace errors. + +- [ ] **Step 2: Back up and deploy changed files to the test server** + +Create a timestamped backup under `/tmp`, copy only changed runtime files to `serber@boffart.miko.ru`, verify SHA-256 hashes, and start the worker through `bin/safe.php`. + +- [ ] **Step 3: Verify successful replay** + +Rewind offset across the three marked `codex-test-*-20260722` rows. Assert source max equals committed offset, total/distinct counts do not grow, and every marker remains exactly once. + +- [ ] **Step 4: Exercise a controlled write failure** + +With a fresh database backup, induce a reversible SQLite write failure, trigger one batch, and assert the log outcome is `rolled_back`, offset is unchanged, and no partial rows remain. Restore permissions/state immediately, rerun, and assert the same batch commits once. + +- [ ] **Step 5: Final review** + +Review `origin/develop..HEAD` for data-loss, rollback, logging privacy, and compatibility risks. Do not push. From a179f384f20fc2b95cfb19c07e83a4361997e76d Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 15:29:23 +0300 Subject: [PATCH 13/18] feat: add atomic CDR batch coordinator --- Lib/AtomicBatch.php | 40 +++++++++++++++++++++++ tests/AtomicBatchTest.php | 68 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 108 insertions(+) create mode 100644 Lib/AtomicBatch.php create mode 100644 tests/AtomicBatchTest.php diff --git a/Lib/AtomicBatch.php b/Lib/AtomicBatch.php new file mode 100644 index 0000000..980c92b --- /dev/null +++ b/Lib/AtomicBatch.php @@ -0,0 +1,40 @@ +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' => '']; + } catch (Throwable $e) { + if ($started) { + try { + $db->rollback(); + } catch (Throwable $ignored) { + // Preserve the original failure; rollback diagnostics are emitted by the caller. + } + } + return ['ok' => false, 'value' => null, 'error' => $e->getMessage()]; + } + } +} diff --git a/tests/AtomicBatchTest.php b/tests/AtomicBatchTest.php new file mode 100644 index 0000000..173cce5 --- /dev/null +++ b/tests/AtomicBatchTest.php @@ -0,0 +1,68 @@ +calls[] = 'begin'; + return $this->beginResult; + } + + public function commit(): bool + { + $this->calls[] = 'commit'; + return $this->commitResult; + } + + public function rollback(): bool + { + $this->calls[] = 'rollback'; + return true; + } +} + +function assertAtomicSame($expected, $actual, string $message): void +{ + if ($expected !== $actual) { + throw new RuntimeException($message . ': expected ' . var_export($expected, true) + . ', got ' . var_export($actual, true)); + } +} + +$db = new FakeTransactionAdapter(); +$success = AtomicBatch::run($db, static fn() => 'saved'); +assertAtomicSame(true, $success['ok'], 'successful batch'); +assertAtomicSame('saved', $success['value'], 'successful value'); +assertAtomicSame(['begin', 'commit'], $db->calls, 'successful transaction calls'); + +$db = new FakeTransactionAdapter(); +$failure = AtomicBatch::run($db, static function (): void { + throw new RuntimeException('write failed'); +}); +assertAtomicSame(false, $failure['ok'], 'failed batch'); +assertAtomicSame('write failed', $failure['error'], 'failure message'); +assertAtomicSame(['begin', 'rollback'], $db->calls, 'failed transaction rolls back'); + +$db = new FakeTransactionAdapter(); +$db->beginResult = false; +$beginFailure = AtomicBatch::run($db, static fn() => 'never'); +assertAtomicSame(false, $beginFailure['ok'], 'begin failure'); +assertAtomicSame(['begin'], $db->calls, 'begin failure does not rollback inactive transaction'); + +$db = new FakeTransactionAdapter(); +$db->commitResult = false; +$commitFailure = AtomicBatch::run($db, static fn() => 'saved'); +assertAtomicSame(false, $commitFailure['ok'], 'commit failure'); +assertAtomicSame(['begin', 'commit', 'rollback'], $db->calls, 'commit failure attempts rollback'); + +echo "AtomicBatchTest: OK\n"; From cbdcffb74ae7ffc24f74c8cb79530afb17c7ff30 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 15:30:19 +0300 Subject: [PATCH 14/18] fix: propagate CDR persistence failures --- Lib/BatchPersistenceResult.php | 30 ++++++++++++++++++++++++++++ bin/ConnectorDB.php | 29 ++++++++++++++------------- tests/BatchPersistenceResultTest.php | 30 ++++++++++++++++++++++++++++ 3 files changed, 75 insertions(+), 14 deletions(-) create mode 100644 Lib/BatchPersistenceResult.php create mode 100644 tests/BatchPersistenceResultTest.php diff --git a/Lib/BatchPersistenceResult.php b/Lib/BatchPersistenceResult.php new file mode 100644 index 0000000..a5f3bf3 --- /dev/null +++ b/Lib/BatchPersistenceResult.php @@ -0,0 +1,30 @@ + 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, + ]; + } +} diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index 2e012d3..d47db6c 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -29,6 +29,7 @@ use Modules\ModuleExtendedCDRs\Lib\Mp3TagService; use Modules\ModuleExtendedCDRs\Lib\CdrQueryBuilder; use Modules\ModuleExtendedCDRs\Lib\CheckpointPolicy; +use Modules\ModuleExtendedCDRs\Lib\BatchPersistenceResult; use Modules\ModuleExtendedCDRs\Lib\SyncPolicy; use Exception; use Modules\ModuleExtendedCDRs\Lib\MikoPBXVersion; @@ -512,7 +513,9 @@ public function syncCdrData(bool $force = false):void // Batch save через raw SQL $start = microtime(true); - [$insertCount, $updateCount] = $this->batchSaveCallHistory($rowsToSave, $arrKeys); + $persistenceResult = $this->batchSaveCallHistory($rowsToSave, $arrKeys); + $insertCount = $persistenceResult['inserted']; + $updateCount = $persistenceResult['updated']; $CallHistorySaveTime = microtime(true) - $start; $this->logger->writeInfo("BatchSave: insert=$insertCount, update=$updateCount"); @@ -789,12 +792,12 @@ private function updateRecallTransferStates(array $records):void * Batch save CallHistory records * @param array $records * @param array $columns - * @return array [insertCount, updateCount] + * @return array{ok:bool,inserted:int,updated:int,errorCategory:string,message:string} */ private function batchSaveCallHistory(array $records, array $columns): array { if (empty($records)) { - return [0, 0]; + return BatchPersistenceResult::success(0, 0); } $newRecords = []; @@ -832,12 +835,10 @@ private function batchSaveCallHistory(array $records, array $columns): array } $sql = "INSERT INTO cdr_general ($columnsStr) VALUES " . implode(', ', $allPlaceholders); - try { - $db->execute($sql, $allValues); - $insertedCount += count($chunk); - } catch (Throwable $e) { - $this->logger->writeError("Batch INSERT failed (chunk of " . count($chunk) . "): " . $e->getMessage()); + if ($db->execute($sql, $allValues) !== true) { + throw new \RuntimeException('insert_failed: database adapter rejected batch insert'); } + $insertedCount += count($chunk); } } @@ -845,16 +846,16 @@ private function batchSaveCallHistory(array $records, array $columns): array $actualUpdates = 0; foreach ($existingRecords as $record) { if ($record->hasChanged()) { - try { - $record->save(); - $actualUpdates++; - } catch (Throwable $e) { - $this->logger->writeError("UPDATE failed for UNIQUEID={$record->UNIQUEID}: " . $e->getMessage()); + if ($record->save() === false) { + throw new \RuntimeException( + 'update_failed: ' . implode('; ', $record->getMessages()) + ); } + $actualUpdates++; } } - return [$insertedCount, $actualUpdates]; + return BatchPersistenceResult::success($insertedCount, $actualUpdates); } public function getCdr(array $filter = []): array diff --git a/tests/BatchPersistenceResultTest.php b/tests/BatchPersistenceResultTest.php new file mode 100644 index 0000000..c065800 --- /dev/null +++ b/tests/BatchPersistenceResultTest.php @@ -0,0 +1,30 @@ + Date: Wed, 22 Jul 2026 15:32:24 +0300 Subject: [PATCH 15/18] fix: commit CDR batches and quarantine atomically --- Lib/QuarantineActivation.php | 16 +++ bin/ConnectorDB.php | 187 ++++++++++++++++------------- tests/QuarantineActivationTest.php | 33 +++++ 3 files changed, 155 insertions(+), 81 deletions(-) create mode 100644 Lib/QuarantineActivation.php create mode 100644 tests/QuarantineActivationTest.php diff --git a/Lib/QuarantineActivation.php b/Lib/QuarantineActivation.php new file mode 100644 index 0000000..0d2285e --- /dev/null +++ b/Lib/QuarantineActivation.php @@ -0,0 +1,16 @@ + 0){ $minOffset = HistoryParser::getMinCdrId(); $settings->cdrOffset = max($newCdrOffset,$minOffset); - $settings->save(); + if ($settings->save() === false) { + throw new \RuntimeException('offset_persist_failed: ' . implode('; ', $settings->getMessages())); + } } if(empty($settings->referenceDate) || (($settings->cdrOffset === null || $settings->cdrOffset === '') && $settings->referenceDate !== '0') ){ $settings->cdrOffset = 1; $settings->referenceDate = date("Y-m-d H:i:s.0", strtotime("-1 days")); - $settings->save(); + if ($settings->save() === false) { + throw new \RuntimeException('settings_initialize_failed: ' . implode('; ', $settings->getMessages())); + } } $this->cdrOffset = (int)$settings->cdrOffset; $this->referenceDate = $settings->referenceDate; @@ -390,11 +396,6 @@ public function syncCdrData(bool $force = false):void $this->logger->writeInfo("Heavy linkedIds (>100 rows): " . count($heavyLinkedIds)); } - // Фиксируем "раздутые" linkedid, чтобы исключить их из следующих выборок. - if (!empty($newOversizedLinkedIds)) { - $this->persistOversizedLinkedIds($newOversizedLinkedIds, $cdrData); - } - // Batch загрузка CallQueuesHistory (1 запрос вместо N) $start = microtime(true); $existingQueues = []; @@ -428,6 +429,7 @@ public function syncCdrData(bool $force = false):void $Mp3TagsTime = 0; $SetCallTypeTime = 0; $rowsToSave = []; + $queuesToSave = []; // Основной цикл — поиск O(1) по массиву foreach ($cdrData as $linkedId => $cdr) { @@ -472,9 +474,7 @@ public function syncCdrData(bool $force = false):void ? ($cdr['q_answer'] - $cdr['q_start']) : ($cdr['q_endtime'] - $cdr['q_start']); - $start = microtime(true); - $cdrQueue->save(); - $CallQueuesHistorySaveTime += microtime(true) - $start; + $queuesToSave[] = $cdrQueue; } // Выбираем источник данных: для тяжёлых — отдельный кеш @@ -511,18 +511,69 @@ public function syncCdrData(bool $force = false):void } } - // Batch save через raw SQL - $start = microtime(true); - $persistenceResult = $this->batchSaveCallHistory($rowsToSave, $arrKeys); + if (!$this->di->has(CdrDbProvider::SERVICE_NAME)) { + $this->di->register(new CdrDbProvider()); + } + $db = $this->di->getShared(CdrDbProvider::SERVICE_NAME); + $transactionResult = AtomicBatch::run($db, function () use ( + $queuesToSave, + $rowsToSave, + $arrKeys, + $newOversizedLinkedIds, + $cdrData, + &$CallQueuesHistorySaveTime, + &$CallHistorySaveTime, + &$recallTransferTime + ): array { + $start = microtime(true); + foreach ($queuesToSave as $queueRecord) { + if ($queueRecord->save() === false) { + throw new \RuntimeException( + 'queue_save_failed: ' . implode('; ', $queueRecord->getMessages()) + ); + } + } + $CallQueuesHistorySaveTime = microtime(true) - $start; + + $start = microtime(true); + $persistence = $this->batchSaveCallHistory($rowsToSave, $arrKeys); + $CallHistorySaveTime = microtime(true) - $start; + + $start = microtime(true); + $this->updateRecallTransferStates($rowsToSave); + $recallTransferTime = microtime(true) - $start; + + $committedQuarantine = $this->persistOversizedLinkedIds($newOversizedLinkedIds, $cdrData); + return [ + 'persistence' => $persistence, + 'quarantine' => $committedQuarantine, + ]; + }); + + if (!$transactionResult['ok']) { + $this->cdrOffset = $oldOffset; + $this->nextSyncDelay = SyncPolicy::ERROR_DELAY_SECONDS; + $error = 'batch_write_failed: ' . $transactionResult['error']; + $failurePolicy = SyncPolicy::decide($oldOffset, $sourceLastId, false, false, $this->catchUpMode); + $this->publishSyncState($oldOffset, $sourceLastId, $failurePolicy, $error); + $this->logger->writeError($error); + return; + } + + $persistenceResult = $transactionResult['value']['persistence']; $insertCount = $persistenceResult['inserted']; $updateCount = $persistenceResult['updated']; - $CallHistorySaveTime = microtime(true) - $start; + $committedQuarantine = $transactionResult['value']['quarantine']; $this->logger->writeInfo("BatchSave: insert=$insertCount, update=$updateCount"); - // Определяем recall/transfer состояния после сохранения - $start = microtime(true); - $this->updateRecallTransferStates($rowsToSave); - $recallTransferTime = microtime(true) - $start; + if (!empty($committedQuarantine)) { + $this->oversizedCache = QuarantineActivation::afterCommit( + $this->oversizedCache, + $committedQuarantine + ); + $this->oversizedPending = array_values(array_diff($this->oversizedPending, $committedQuarantine)); + $this->logger->writeInfo('Oversized linkedIds committed and excluded: ' . count($committedQuarantine)); + } // Если в этом цикле обнаружены новые "раздутые" linkedid — удерживаем offset. // Потолок в 5000 строк был съеден зависшим звонком, поэтому обычные linkedid @@ -531,11 +582,12 @@ public function syncCdrData(bool $force = false):void if (!empty($newOversizedLinkedIds)) { $this->cdrOffset = $oldOffset; $this->logger->writeInfo("Holding offset at $oldOffset: detected " . count($newOversizedLinkedIds) . " new oversized linkedId(s)"); + $this->publishSyncState($oldOffset, $sourceLastId, $policy, ''); $this->logger->writeInfo("End sync with offset {$this->cdrOffset} (+0)"); return; } - $this->cdrOffset = CheckpointPolicy::nextOffset([ + $nextOffset = CheckpointPolicy::nextOffset([ 'oldOffset' => $oldOffset, 'parsedOffset' => $parsedOffset, 'requestOk' => true, @@ -553,9 +605,20 @@ public function syncCdrData(bool $force = false):void 'SetCallTypeTime' => round($SetCallTypeTime, 4), 'RecallTransferTime' => round($recallTransferTime, 4)], "Timing"); - if($oldOffset !== $this->cdrOffset){ - $this->logger->writeInfo("Update progress, offset $oldOffset to new value $this->cdrOffset "); - $this->updateSettings($this->cdrOffset); + if($oldOffset !== $nextOffset){ + $this->logger->writeInfo("Update progress, offset $oldOffset to new value $nextOffset "); + try { + $this->updateSettings($nextOffset); + } catch (Throwable $e) { + $this->cdrOffset = $oldOffset; + $this->nextSyncDelay = SyncPolicy::ERROR_DELAY_SECONDS; + $failurePolicy = SyncPolicy::decide($oldOffset, $sourceLastId, false, false, $this->catchUpMode); + $this->publishSyncState($oldOffset, $sourceLastId, $failurePolicy, 'offset_persist_failed'); + $this->logger->writeError('offset_persist_failed: ' . $e->getMessage()); + return; + } + } else { + $this->cdrOffset = $oldOffset; } $policy = SyncPolicy::decide( $this->cdrOffset, @@ -783,7 +846,16 @@ private function updateRecallTransferStates(array $records):void // Сохраняем только если stateCall изменился if ($dbData->stateCall !== CallHistory::CALL_STATE_OK) { - $dbData->save(); + $saved = $db->execute( + 'UPDATE cdr_general SET stateCall = :stateCall WHERE UNIQUEID = :uniqueId', + [ + 'stateCall' => $dbData->stateCall, + 'uniqueId' => $dbData->UNIQUEID, + ] + ); + if ($saved !== true) { + throw new \RuntimeException('state_save_failed: database adapter rejected update'); + } } } } @@ -1041,10 +1113,11 @@ private function loadOversizedLinkedIds(): array * Фиксирует новые "раздутые" linkedid в служебной таблице и в кэше. * @param string[] $linkedIds * @param array $cdrData Данные текущей выборки (для rowCount/maxId). - * @return void + * @return string[] IDs whose quarantine records were written in the current transaction. */ - private function persistOversizedLinkedIds(array $linkedIds, array $cdrData): void + private function persistOversizedLinkedIds(array $linkedIds, array $cdrData): array { + $persisted = []; foreach ($linkedIds as $linkedId) { if (in_array($linkedId, $this->oversizedCache, true)) { continue; @@ -1072,28 +1145,14 @@ private function persistOversizedLinkedIds(array $linkedIds, array $cdrData): vo $record->lastFailureAt = $record->detectedAt; $record->nextRetryAt = date('Y-m-d H:i:s', time() + 60); $record->status = 'pending'; - $saved = $record->save(); - - // В любом случае исключаем linkedid в пределах текущей сессии воркера, - // иначе сбой save() (блокировка БД, UNIQUE-коллизия) приведёт к бесконечному - // повторному детекту и вечному удержанию offset. - if (!in_array($linkedId, $this->oversizedCache, true)) { - $this->oversizedCache[] = $linkedId; - } - if ($saved) { - // Успешно записан — убираем из session-only набора, если был там. - $this->oversizedPending = array_values(array_diff($this->oversizedPending, [$linkedId])); - $this->logger->writeInfo("Oversized linkedId excluded from sync: $linkedId (rows=$rowCount, maxId=$maxId)"); - } else { - // Запись не удалась — держим в session-only наборе, чтобы кэш не потерял - // его при обновлении из БД (иначе трэшинг offset). Повторная запись — - // после перезапуска воркера через повторный детект. - if (!in_array($linkedId, $this->oversizedPending, true)) { - $this->oversizedPending[] = $linkedId; - } - $this->logger->writeError("Failed to persist oversized linkedId (excluded in-memory only): $linkedId (" . implode('; ', $record->getMessages()) . ")"); + if ($record->save() === false) { + throw new \RuntimeException( + 'quarantine_save_failed: ' . implode('; ', $record->getMessages()) + ); } + $persisted[] = $linkedId; } + return $persisted; } /** @@ -1105,42 +1164,8 @@ private function persistOversizedLinkedIds(array $linkedIds, array $cdrData): vo */ private function pruneOversizedLinkedIds(): void { - if (time() - $this->oversizedPruneTime < 3600) { - return; - } - $this->oversizedPruneTime = time(); - - $stored = array_column(OversizedLinkedIds::find(['columns' => 'linkedid'])->toArray(), 'linkedid'); - if (empty($stored)) { - return; - } - - // Активные — те, у кого ещё есть строки за текущим offset. - // null означает сбой запроса к ядру (Beanstalk timeout) — в этом случае НЕ удаляем - // ничего, чтобы не разкарантинить активные зависшие звонки и не вызвать повторный стопор. - // Пустой массив [] означает, что запрос выполнился и активных действительно нет — - // тогда все устаревшие записи можно удалить. - $active = HistoryParser::getActiveLinkedIds($stored, $this->cdrOffset); - if ($active === null) { - return; - } - $toDelete = array_diff($stored, $active); - if (empty($toDelete)) { - return; - } - - $records = OversizedLinkedIds::find([ - 'linkedid IN ({ids:array})', - 'bind' => ['ids' => array_values($toDelete)], - ]); - foreach ($records as $record) { - $record->delete(); - } - - // Кэш = активные из БД ∪ session-only записи (последние в БД отсутствуют). - $this->oversizedCache = array_values(array_unique(array_merge($active, $this->oversizedPending))); - $this->oversizedCacheTime = time(); - $this->logger->writeInfo("Pruned oversized linkedIds: removed " . count($toDelete) . ", kept " . count($active)); + // Audit records are intentionally retained until a reconciler can mark + // an oversized call resolved with a durable, inspectable outcome. } /** diff --git a/tests/QuarantineActivationTest.php b/tests/QuarantineActivationTest.php new file mode 100644 index 0000000..5dce196 --- /dev/null +++ b/tests/QuarantineActivationTest.php @@ -0,0 +1,33 @@ + Date: Wed, 22 Jul 2026 15:33:50 +0300 Subject: [PATCH 16/18] feat: add safe CDR batch diagnostics --- Lib/BatchLogContext.php | 32 +++++++++++++++++++++++++ bin/ConnectorDB.php | 38 ++++++++++++++++++++++++++++++ tests/BatchLogContextTest.php | 44 +++++++++++++++++++++++++++++++++++ 3 files changed, 114 insertions(+) create mode 100644 Lib/BatchLogContext.php create mode 100644 tests/BatchLogContextTest.php diff --git a/Lib/BatchLogContext.php b/Lib/BatchLogContext.php new file mode 100644 index 0000000..4904f93 --- /dev/null +++ b/Lib/BatchLogContext.php @@ -0,0 +1,32 @@ + */ + 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'] ?? ''), + ]; + } +} diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index 1219ebf..f8ebc7f 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -32,6 +32,7 @@ use Modules\ModuleExtendedCDRs\Lib\BatchPersistenceResult; use Modules\ModuleExtendedCDRs\Lib\AtomicBatch; use Modules\ModuleExtendedCDRs\Lib\QuarantineActivation; +use Modules\ModuleExtendedCDRs\Lib\BatchLogContext; use Modules\ModuleExtendedCDRs\Lib\SyncPolicy; use Exception; use Modules\ModuleExtendedCDRs\Lib\MikoPBXVersion; @@ -321,6 +322,7 @@ public function syncCdrData(bool $force = false):void } $this->lastSyncTime = time(); $oldOffset = $this->cdrOffset; + $batchStarted = microtime(true); $this->logger->writeInfo('...Start sync with offset...'. $oldOffset); $sourceState = HistoryParser::getLastCdrState(); @@ -331,6 +333,7 @@ public function syncCdrData(bool $force = false):void if (!$sourceState['ok']) { $this->publishSyncState($oldOffset, $sourceLastId, $policy, 'source_last_id_failed'); $this->logger->writeError('batch_failed: source_last_id_failed'); + $this->writeBatchOutcome($oldOffset, $oldOffset, $sourceLastId, [], $policy, $batchStarted, 'source_failed', 'source_last_id_failed'); return; } @@ -343,6 +346,7 @@ public function syncCdrData(bool $force = false):void $this->nextSyncDelay = SyncPolicy::ERROR_DELAY_SECONDS; $this->publishSyncState($oldOffset, $sourceLastId, $policy, $batchResult['error']); $this->logger->writeError('batch_failed: ' . $batchResult['error']); + $this->writeBatchOutcome($oldOffset, $oldOffset, $sourceLastId, $batchResult, $policy, $batchStarted, 'source_failed', $batchResult['error']); return; } $cdrData = $batchResult['data']; @@ -557,6 +561,8 @@ public function syncCdrData(bool $force = false):void $failurePolicy = SyncPolicy::decide($oldOffset, $sourceLastId, false, false, $this->catchUpMode); $this->publishSyncState($oldOffset, $sourceLastId, $failurePolicy, $error); $this->logger->writeError($error); + $category = explode(':', $transactionResult['error'], 2)[0]; + $this->writeBatchOutcome($oldOffset, $oldOffset, $sourceLastId, $batchResult, $failurePolicy, $batchStarted, 'rolled_back', $category); return; } @@ -583,6 +589,7 @@ public function syncCdrData(bool $force = false):void $this->cdrOffset = $oldOffset; $this->logger->writeInfo("Holding offset at $oldOffset: detected " . count($newOversizedLinkedIds) . " new oversized linkedId(s)"); $this->publishSyncState($oldOffset, $sourceLastId, $policy, ''); + $this->writeBatchOutcome($oldOffset, $oldOffset, $sourceLastId, $batchResult, $policy, $batchStarted, 'quarantined', ''); $this->logger->writeInfo("End sync with offset {$this->cdrOffset} (+0)"); return; } @@ -615,6 +622,7 @@ public function syncCdrData(bool $force = false):void $failurePolicy = SyncPolicy::decide($oldOffset, $sourceLastId, false, false, $this->catchUpMode); $this->publishSyncState($oldOffset, $sourceLastId, $failurePolicy, 'offset_persist_failed'); $this->logger->writeError('offset_persist_failed: ' . $e->getMessage()); + $this->writeBatchOutcome($oldOffset, $nextOffset, $sourceLastId, $batchResult, $failurePolicy, $batchStarted, 'offset_persist_failed', 'offset_persist_failed', $insertCount, $updateCount); return; } } else { @@ -630,6 +638,7 @@ public function syncCdrData(bool $force = false):void $this->nextSyncDelay = $policy['delay']; $this->catchUpMode = $policy['mode'] === SyncPolicy::MODE_CATCH_UP; $this->publishSyncState($this->cdrOffset, $sourceLastId, $policy, ''); + $this->writeBatchOutcome($oldOffset, $this->cdrOffset, $sourceLastId, $batchResult, $policy, $batchStarted, 'committed', '', $insertCount, $updateCount); $offsetDelta = $this->cdrOffset - $oldOffset; $this->logger->writeInfo("End sync with offset {$this->cdrOffset} (+$offsetDelta)"); } @@ -651,6 +660,35 @@ private function publishSyncState(int $offset, int $sourceLastId, array $policy, ]); } + private function writeBatchOutcome( + int $oldOffset, + int $proposedOffset, + int $sourceLastId, + array $batch, + array $policy, + float $startedAt, + string $outcome, + string $errorCategory, + int $inserted = 0, + int $updated = 0 + ): void { + $this->logger->writeInfo(BatchLogContext::make([ + 'oldOffset' => $oldOffset, + 'proposedOffset' => $proposedOffset, + 'sourceLastId' => $sourceLastId, + 'minId' => $batch['minId'] ?? 0, + 'maxId' => $batch['maxId'] ?? 0, + 'linkedIdCount' => $batch['linkedIdCount'] ?? 0, + 'rowCount' => $batch['rowCount'] ?? 0, + 'inserted' => $inserted, + 'updated' => $updated, + 'mode' => $policy['mode'] ?? SyncPolicy::MODE_ERROR, + 'elapsedMs' => (int)round((microtime(true) - $startedAt) * 1000), + 'outcome' => $outcome, + 'errorCategory' => $errorCategory, + ])); + } + /** * Возвращает путь к файлу записи по ID. * @param string $id diff --git a/tests/BatchLogContextTest.php b/tests/BatchLogContextTest.php new file mode 100644 index 0000000..0f253a0 --- /dev/null +++ b/tests/BatchLogContextTest.php @@ -0,0 +1,44 @@ + 100, + 'proposedOffset' => 120, + 'sourceLastId' => 150, + 'minId' => 101, + 'maxId' => 120, + 'linkedIdCount' => 3, + 'rowCount' => 20, + 'inserted' => 18, + 'updated' => 2, + 'mode' => 'normal', + 'elapsedMs' => 42, + 'outcome' => 'committed', + 'errorCategory' => '', + 'linkedid' => 'secret-call-id', + 'UNIQUEID' => 'secret-unique-id', + 'src_num' => '79990001122', + 'recordingfile' => '/secret.wav', +]); + +assertLogValue('cdr_sync_batch', $context['event'], 'stable event name'); +assertLogValue(50, $context['lagBefore'], 'lag before batch'); +assertLogValue(30, $context['lagAfter'], 'lag after batch'); +foreach (['linkedid', 'UNIQUEID', 'src_num', 'recordingfile'] as $forbidden) { + assertLogValue(false, array_key_exists($forbidden, $context), "forbidden key $forbidden"); +} + +echo "BatchLogContextTest: OK\n"; From cf582b8c738a1fa29edc75328ac9d97bc19c1f43 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Wed, 22 Jul 2026 15:34:52 +0300 Subject: [PATCH 17/18] fix: preserve CDR log severity and payloads --- Lib/LogFormatPolicy.php | 26 ++++++++++++++++++++++++++ Lib/Logger.php | 19 ++++++------------- tests/LogFormatPolicyTest.php | 21 +++++++++++++++++++++ 3 files changed, 53 insertions(+), 13 deletions(-) create mode 100644 Lib/LogFormatPolicy.php create mode 100644 tests/LogFormatPolicyTest.php diff --git a/Lib/LogFormatPolicy.php b/Lib/LogFormatPolicy.php new file mode 100644 index 0000000..ac56a91 --- /dev/null +++ b/Lib/LogFormatPolicy.php @@ -0,0 +1,26 @@ +logFile); - $lineFormatter = new LineFormatter("[%date%][%type%] %message%", "Y-m-d H:i:s"); + $lineFormatter = new LineFormatter( + LogFormatPolicy::template(MikoPBXVersion::isPhalcon5Version()), + "Y-m-d H:i:s" + ); $adapter->setFormatter($lineFormatter); $loggerClass = MikoPBXVersion::getLoggerClass(); $this->logger = new $loggerClass( @@ -141,16 +144,6 @@ public function writeInfo($data, string $preMessage=''): void */ private function getDecodedString($data):string { - try { - $printedData = json_encode($data, JSON_THROW_ON_ERROR | JSON_UNESCAPED_SLASHES | JSON_UNESCAPED_UNICODE); - }catch (\Exception $e){ - $printedData = print_r($data, true); - } - if(is_bool($printedData)){ - $result = ''; - }else{ - $result = urldecode($printedData); - } - return $result; + return LogFormatPolicy::encode($data); } -} \ No newline at end of file +} diff --git a/tests/LogFormatPolicyTest.php b/tests/LogFormatPolicyTest.php new file mode 100644 index 0000000..2cbd498 --- /dev/null +++ b/tests/LogFormatPolicyTest.php @@ -0,0 +1,21 @@ + Date: Wed, 22 Jul 2026 15:45:15 +0300 Subject: [PATCH 18/18] fix: harden rollback and quarantine audit --- Lib/AtomicBatch.php | 30 +++++++++++++++++++++------ Lib/QuarantinePolicy.php | 13 ++++++++++++ bin/ConnectorDB.php | 38 ++++++++++++++++++++++++++-------- bin/SyncRecords.php | 5 ++++- tests/AtomicBatchTest.php | 24 ++++++++++++++++++++- tests/QuarantinePolicyTest.php | 4 ++++ 6 files changed, 97 insertions(+), 17 deletions(-) diff --git a/Lib/AtomicBatch.php b/Lib/AtomicBatch.php index 980c92b..279ea88 100644 --- a/Lib/AtomicBatch.php +++ b/Lib/AtomicBatch.php @@ -10,7 +10,7 @@ final class AtomicBatch /** * @param object $db Phalcon-compatible transaction adapter. * @param callable $operation Required database mutations. - * @return array{ok:bool,value:mixed,error:string} + * @return array{ok:bool,value:mixed,error:string,rollbackOk:?bool,rollbackError:string} */ public static function run(object $db, callable $operation): array { @@ -25,16 +25,34 @@ public static function run(object $db, callable $operation): array throw new RuntimeException('transaction_commit_failed'); } $started = false; - return ['ok' => true, 'value' => $value, 'error' => '']; + return [ + 'ok' => true, + 'value' => $value, + 'error' => '', + 'rollbackOk' => null, + 'rollbackError' => '', + ]; } catch (Throwable $e) { + $rollbackOk = null; + $rollbackError = ''; if ($started) { try { - $db->rollback(); - } catch (Throwable $ignored) { - // Preserve the original failure; rollback diagnostics are emitted by the caller. + $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()]; + return [ + 'ok' => false, + 'value' => null, + 'error' => $e->getMessage(), + 'rollbackOk' => $rollbackOk, + 'rollbackError' => $rollbackError, + ]; } } } diff --git a/Lib/QuarantinePolicy.php b/Lib/QuarantinePolicy.php index 7af296a..bb5dfed 100644 --- a/Lib/QuarantinePolicy.php +++ b/Lib/QuarantinePolicy.php @@ -31,4 +31,17 @@ public static function resolved(array $current, int $now): array $current['lastFailureAt'] = $now; return $current; } + + /** @return array */ + public static function manual(string $reason, int $now): array + { + return [ + 'reason' => $reason, + 'attempts' => 1, + 'firstFailureAt' => $now, + 'lastFailureAt' => $now, + 'nextRetryAt' => 0, + 'status' => 'manual', + ]; + } } diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index f8ebc7f..779dc56 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -32,6 +32,7 @@ use Modules\ModuleExtendedCDRs\Lib\BatchPersistenceResult; use Modules\ModuleExtendedCDRs\Lib\AtomicBatch; use Modules\ModuleExtendedCDRs\Lib\QuarantineActivation; +use Modules\ModuleExtendedCDRs\Lib\QuarantinePolicy; use Modules\ModuleExtendedCDRs\Lib\BatchLogContext; use Modules\ModuleExtendedCDRs\Lib\SyncPolicy; use Exception; @@ -556,13 +557,31 @@ public function syncCdrData(bool $force = false):void if (!$transactionResult['ok']) { $this->cdrOffset = $oldOffset; - $this->nextSyncDelay = SyncPolicy::ERROR_DELAY_SECONDS; + $rollbackFailed = $transactionResult['rollbackOk'] === false; + $this->nextSyncDelay = $rollbackFailed ? 1 : SyncPolicy::ERROR_DELAY_SECONDS; + if ($rollbackFailed) { + // The shared adapter may still own a broken transaction. Exit this + // worker instance so the safe-script restarts it with a new connection. + $this->needRestart = true; + } $error = 'batch_write_failed: ' . $transactionResult['error']; + if ($rollbackFailed) { + $error .= '; rollback_failed: ' . $transactionResult['rollbackError']; + } $failurePolicy = SyncPolicy::decide($oldOffset, $sourceLastId, false, false, $this->catchUpMode); $this->publishSyncState($oldOffset, $sourceLastId, $failurePolicy, $error); $this->logger->writeError($error); $category = explode(':', $transactionResult['error'], 2)[0]; - $this->writeBatchOutcome($oldOffset, $oldOffset, $sourceLastId, $batchResult, $failurePolicy, $batchStarted, 'rolled_back', $category); + $this->writeBatchOutcome( + $oldOffset, + $oldOffset, + $sourceLastId, + $batchResult, + $failurePolicy, + $batchStarted, + $rollbackFailed ? 'rollback_failed' : 'rolled_back', + $rollbackFailed ? 'transaction_rollback_failed' : $category + ); return; } @@ -589,7 +608,7 @@ public function syncCdrData(bool $force = false):void $this->cdrOffset = $oldOffset; $this->logger->writeInfo("Holding offset at $oldOffset: detected " . count($newOversizedLinkedIds) . " new oversized linkedId(s)"); $this->publishSyncState($oldOffset, $sourceLastId, $policy, ''); - $this->writeBatchOutcome($oldOffset, $oldOffset, $sourceLastId, $batchResult, $policy, $batchStarted, 'quarantined', ''); + $this->writeBatchOutcome($oldOffset, $oldOffset, $sourceLastId, $batchResult, $policy, $batchStarted, 'quarantined', '', $insertCount, $updateCount); $this->logger->writeInfo("End sync with offset {$this->cdrOffset} (+0)"); return; } @@ -1177,12 +1196,13 @@ private function persistOversizedLinkedIds(array $linkedIds, array $cdrData): ar $record->detectedAt = date('Y-m-d H:i:s'); $record->minId = $minId; $record->maxRangeId = $maxId; - $record->reason = 'row_limit'; - $record->attempts = 0; - $record->firstFailureAt = $record->detectedAt; - $record->lastFailureAt = $record->detectedAt; - $record->nextRetryAt = date('Y-m-d H:i:s', time() + 60); - $record->status = 'pending'; + $quarantine = QuarantinePolicy::manual('row_limit', time()); + $record->reason = $quarantine['reason']; + $record->attempts = $quarantine['attempts']; + $record->firstFailureAt = date('Y-m-d H:i:s', $quarantine['firstFailureAt']); + $record->lastFailureAt = date('Y-m-d H:i:s', $quarantine['lastFailureAt']); + $record->nextRetryAt = ''; + $record->status = $quarantine['status']; if ($record->save() === false) { throw new \RuntimeException( 'quarantine_save_failed: ' . implode('; ', $record->getMessages()) diff --git a/bin/SyncRecords.php b/bin/SyncRecords.php index 72f4195..68090dd 100644 --- a/bin/SyncRecords.php +++ b/bin/SyncRecords.php @@ -275,7 +275,10 @@ private function loadOversizedLinkedIds(): array { if (time() - $this->oversizedCacheTime > 60) { try { - $rows = OversizedLinkedIds::find(['columns' => 'linkedid']); + $rows = OversizedLinkedIds::find([ + "status IS NULL OR status <> 'resolved'", + 'columns' => 'linkedid', + ]); $this->oversizedCache = array_column($rows->toArray(), 'linkedid'); } catch (Throwable $e) { $this->logger->writeError("loadOversizedLinkedIds: " . $e->getMessage()); diff --git a/tests/AtomicBatchTest.php b/tests/AtomicBatchTest.php index 173cce5..b57716b 100644 --- a/tests/AtomicBatchTest.php +++ b/tests/AtomicBatchTest.php @@ -11,6 +11,8 @@ final class FakeTransactionAdapter public array $calls = []; public bool $beginResult = true; public bool $commitResult = true; + public bool $rollbackResult = true; + public bool $throwOnRollback = false; public function begin(): bool { @@ -27,7 +29,10 @@ public function commit(): bool public function rollback(): bool { $this->calls[] = 'rollback'; - return true; + if ($this->throwOnRollback) { + throw new RuntimeException('rollback exploded'); + } + return $this->rollbackResult; } } @@ -52,6 +57,7 @@ function assertAtomicSame($expected, $actual, string $message): void assertAtomicSame(false, $failure['ok'], 'failed batch'); assertAtomicSame('write failed', $failure['error'], 'failure message'); assertAtomicSame(['begin', 'rollback'], $db->calls, 'failed transaction rolls back'); +assertAtomicSame(true, $failure['rollbackOk'], 'successful rollback is reported'); $db = new FakeTransactionAdapter(); $db->beginResult = false; @@ -65,4 +71,20 @@ function assertAtomicSame($expected, $actual, string $message): void assertAtomicSame(false, $commitFailure['ok'], 'commit failure'); assertAtomicSame(['begin', 'commit', 'rollback'], $db->calls, 'commit failure attempts rollback'); +$db = new FakeTransactionAdapter(); +$db->rollbackResult = false; +$rollbackFailure = AtomicBatch::run($db, static function (): void { + throw new RuntimeException('write failed'); +}); +assertAtomicSame(false, $rollbackFailure['rollbackOk'], 'false rollback is reported'); +assertAtomicSame('transaction_rollback_failed', $rollbackFailure['rollbackError'], 'false rollback category'); + +$db = new FakeTransactionAdapter(); +$db->throwOnRollback = true; +$rollbackException = AtomicBatch::run($db, static function (): void { + throw new RuntimeException('write failed'); +}); +assertAtomicSame(false, $rollbackException['rollbackOk'], 'rollback exception is reported'); +assertAtomicSame('rollback exploded', $rollbackException['rollbackError'], 'rollback exception message'); + echo "AtomicBatchTest: OK\n"; diff --git a/tests/QuarantinePolicyTest.php b/tests/QuarantinePolicyTest.php index 5e8f068..26152d4 100644 --- a/tests/QuarantinePolicyTest.php +++ b/tests/QuarantinePolicyTest.php @@ -33,4 +33,8 @@ function quarantineAssert($expected, $actual, string $message): void quarantineAssert('resolved', $resolved['status'], 'successful retry resolves'); quarantineAssert(2000, $resolved['lastFailureAt'], 'resolution timestamp retained for audit'); +$manual = QuarantinePolicy::manual('row_limit', 3000); +quarantineAssert('manual', $manual['status'], 'oversized call requires manual resolution'); +quarantineAssert(0, $manual['nextRetryAt'], 'manual quarantine has no fake retry deadline'); + echo "QuarantinePolicyTest: OK\n";