From 1de74a16f40fdfb27834c7a4b322206e0fcda9c8 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:04:03 +0300 Subject: [PATCH 01/10] docs: design resilient module workers --- ...6-08-03-module-worker-resilience-design.md | 84 +++++++++++++++++++ 1 file changed, 84 insertions(+) create mode 100644 docs/superpowers/specs/2026-08-03-module-worker-resilience-design.md diff --git a/docs/superpowers/specs/2026-08-03-module-worker-resilience-design.md b/docs/superpowers/specs/2026-08-03-module-worker-resilience-design.md new file mode 100644 index 0000000..4015439 --- /dev/null +++ b/docs/superpowers/specs/2026-08-03-module-worker-resilience-design.md @@ -0,0 +1,84 @@ +# ModuleExtendedCDRs Worker Resilience Design + +## Goal + +Prevent overlapping module watchdog runs and make `ConnectorDB` recover cleanly from Beanstalk failures without changing MikoPBX Core, Redis, nginx, or Monit. + +## Scope + +Only files shipped by `ModuleExtendedCDRs` may change. The work covers the module cron entry, `bin/safe.php`, `bin/ConnectorDB.php`, focused pure-PHP support code, tests, and module-owned logging. + +The production-like verification target is the test server `serber@boffart.miko.ru`. Deployment must use the existing module installation/update workflow rather than copying an incomplete source tree over the installed module. + +## Watchdog Singleton + +`bin/safe.php` will acquire a non-blocking exclusive file lock before it calls MikoPBX process or Beanstalk APIs. The lock file will live in the configured runtime/temp area when available, with a module-specific directory under `/tmp` as a fallback. + +If another watchdog invocation owns the lock, the new invocation will log one rate-limited skip message and exit successfully without inspecting or starting workers. The lock contents will include the current PID and start timestamp for diagnostics, but lock ownership will be determined by the operating-system lock rather than trusting PID-file contents. + +The file handle will remain open for the whole watchdog run. Normal return, exceptions, and PHP shutdown will release the lock automatically. Stale file contents do not block startup when no process owns the OS lock. + +## Bounded Watchdog Execution + +The watchdog will set a module-owned execution deadline. On platforms with `pcntl`, an alarm will terminate an overlong watchdog run after logging its phase and elapsed time. On platforms without `pcntl`, the singleton still prevents overlap and elapsed-time diagnostics identify the blocking phase; no unsafe asynchronous signal emulation will be added. + +Worker discovery, duplicate cleanup, and worker startup will be logged as named phases with PID, elapsed milliseconds, and outcome. Existing MikoPBX process helpers remain the source of truth for worker detection and startup. + +The cron schedule remains once per minute. The cron command will also use the available BusyBox `timeout` command as an outer process boundary so a PHP call blocked below userland cannot accumulate indefinitely. The timeout must exceed the internal deadline and must not kill a normally completing run. + +## ConnectorDB Beanstalk Recovery + +`ConnectorDB` will treat Beanstalk construction, subscription, wait, publish, and request failures as recoverable worker-process failures. It will log a structured event containing operation, exception class, message, PID, worker uptime, memory usage, open-file-descriptor count when `/proc` is available, and elapsed milliseconds. + +Failures in the main listener connection will cause the worker process to leave its loop and exit. The next singleton-protected watchdog run will start a fresh process and therefore a fresh Beanstalk connection. The worker will not run an unbounded reconnect loop inside a process with potentially corrupted connection state. + +Synchronous `invoke()` keeps its existing external contract: failures return an empty array. It will clean request/response temporary files in all paths and emit a bounded diagnostic without exposing request data or recording paths. + +## Diagnostics + +Module-owned logs will provide structured events for: + +- watchdog acquisition, overlap skip, timeout, completion, and phase duration; +- detected ConnectorDB PIDs and duplicate cleanup; +- ConnectorDB start, shutdown, Beanstalk failure, and periodic health; +- process uptime, RSS/memory, open file descriptors, and TCP socket count where supported. + +Health diagnostics will be rate-limited to avoid turning a degraded system into a logging storm. Missing `/proc` data will be represented as unavailable and will never terminate a worker. + +## Error Handling + +- Lock contention is an expected successful no-op. +- A watchdog exception is logged and returns a non-zero exit status after releasing the lock. +- An internal or outer deadline terminates only that watchdog invocation; it does not send signals to a healthy ConnectorDB. +- Existing duplicate ConnectorDB cleanup remains, but it executes only from the lock owner. +- A listener-side Beanstalk error shuts down ConnectorDB cleanly so the watchdog can restart it. +- Diagnostic failures are swallowed after a compact warning and cannot become a worker failure. + +## Tests + +Pure-PHP tests will cover lock contention, stale lock contents, lock release, diagnostic collection without `/proc`, rate limiting, Beanstalk failure classification, and temporary-file cleanup policy. Script-level tests will verify that cron configuration contains both the timeout boundary and the module watchdog command. + +All existing repository PHP tests and syntax checks for changed PHP files must pass. + +## Test-Server Verification + +Before deployment, capture a read-only baseline on `serber@boffart.miko.ru`: installed module version/files, MikoPBX and PHP versions, cron entry, ConnectorDB and `safe.php` processes, Beanstalk status, recent module/system errors, and web/API health. + +After deploying through the module-safe workflow: + +1. Trigger several concurrent `safe.php` invocations and verify only one owns the lock. +2. Verify exactly one ConnectorDB worker remains after watchdog convergence. +3. Verify normal CDR synchronization continues and the worker answers its existing API calls. +4. Exercise a controlled Beanstalk interruption only if it can be isolated safely on the test server; otherwise validate recovery by restarting only the ConnectorDB worker and observing fresh subscription/health events. +5. Verify the web interface and authenticated module endpoints remain available. +6. Observe at least two cron intervals for duplicate processes, stale locks, error loops, and excessive log volume. + +Any failed activation or regression requires restoration of the timestamped installed-module backup and restart of only the affected module worker. + +## Non-Goals + +- Changing Redis or nginx behavior. +- Changing Monit policies. +- Modifying MikoPBX Core classes such as `WorkerSafeScriptsCore`, `Processes`, or `BeanstalkClient`. +- Adding a general-purpose process supervisor to the module. +- Destructive Beanstalk or Redis fault injection on a server carrying non-test traffic. From 944f5844b4020d35a718311c61d8f2404576fa0b Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:06:16 +0300 Subject: [PATCH 02/10] docs: plan resilient module workers --- .../2026-08-03-module-worker-resilience.md | 241 ++++++++++++++++++ 1 file changed, 241 insertions(+) create mode 100644 docs/superpowers/plans/2026-08-03-module-worker-resilience.md diff --git a/docs/superpowers/plans/2026-08-03-module-worker-resilience.md b/docs/superpowers/plans/2026-08-03-module-worker-resilience.md new file mode 100644 index 0000000..967fa89 --- /dev/null +++ b/docs/superpowers/plans/2026-08-03-module-worker-resilience.md @@ -0,0 +1,241 @@ +# ModuleExtendedCDRs Worker Resilience 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:** Prevent overlapping ModuleExtendedCDRs watchdog processes, bound their runtime, and make ConnectorDB exit cleanly with useful diagnostics after Beanstalk failures. + +**Architecture:** Small pure-PHP policy/resource classes own locking, deadlines, process metrics, and diagnostic payload construction. `safe.php` and `ConnectorDB.php` remain thin infrastructure adapters around MikoPBX Core APIs. The cron entry supplies an outer BusyBox timeout, while the ConnectorDB worker exits after a broken listener connection so the next singleton watchdog run establishes a fresh connection. + +**Tech Stack:** PHP 7.4.6, `flock`, optional `pcntl_alarm`, Linux `/proc`, MikoPBX WorkerBase/Processes/BeanstalkClient, standalone PHP regression tests. + +## Global Constraints + +- Change only files shipped by `ModuleExtendedCDRs`. +- Do not change MikoPBX Core, Redis, nginx, or Monit. +- Preserve the public contract of `ConnectorDB::invoke()`: failures return an empty array. +- Diagnostics must not contain request payloads, phone numbers, linked IDs, or recording paths. +- Missing `pcntl` or `/proc` support must degrade safely. +- Deploy to `serber@boffart.miko.ru` through the existing module-safe installation/update workflow and retain a timestamped rollback backup. + +--- + +### Task 1: Exclusive Module Watchdog Lease + +**Files:** +- Create: `Lib/WorkerWatchdogLease.php` +- Create: `tests/WorkerWatchdogLeaseTest.php` + +**Interfaces:** +- Produces: `WorkerWatchdogLease::tryAcquire(string $path, int $pid, int $startedAt): ?WorkerWatchdogLease` +- Produces: `WorkerWatchdogLease::release(): void` +- Produces: automatic release from `__destruct()` + +- [ ] **Step 1: Write the failing lease test** + +Create two real lease attempts against one temporary file. Assert the first succeeds, the second returns `null`, stale pre-existing contents do not prevent acquisition, and acquisition succeeds after release. + +```php +$first = WorkerWatchdogLease::tryAcquire($path, 101, 1700000000); +assertLease(true, $first instanceof WorkerWatchdogLease, 'first owner acquires'); +assertLease(null, WorkerWatchdogLease::tryAcquire($path, 202, 1700000001), 'contender skips'); +$first->release(); +assertLease(true, WorkerWatchdogLease::tryAcquire($path, 303, 1700000002) instanceof WorkerWatchdogLease, 'released lease reacquires'); +``` + +- [ ] **Step 2: Run the test and verify RED** + +Run: `php tests/WorkerWatchdogLeaseTest.php` +Expected: FAIL because `Lib/WorkerWatchdogLease.php` does not exist. + +- [ ] **Step 3: Implement the minimal lease** + +Open with `fopen($path, 'c+')`, acquire `LOCK_EX | LOCK_NB`, truncate and write a JSON diagnostic record only after acquisition, keep the resource in the object, and use an idempotent `release()`. + +- [ ] **Step 4: Run the test and verify GREEN** + +Run: `php tests/WorkerWatchdogLeaseTest.php` +Expected: `WorkerWatchdogLeaseTest: OK`. + +- [ ] **Step 5: Commit the lease** + +```bash +git add Lib/WorkerWatchdogLease.php tests/WorkerWatchdogLeaseTest.php +git commit -m "feat: serialize module watchdog runs" +``` + +### Task 2: Watchdog Runtime Policy and Process Diagnostics + +**Files:** +- Create: `Lib/WorkerRuntimePolicy.php` +- Create: `Lib/WorkerProcessMetrics.php` +- Create: `tests/WorkerRuntimePolicyTest.php` +- Create: `tests/WorkerProcessMetricsTest.php` + +**Interfaces:** +- Produces: `WorkerRuntimePolicy::watchdogDeadlineSeconds(): int` +- Produces: `WorkerRuntimePolicy::outerTimeoutSeconds(): int` +- Produces: `WorkerRuntimePolicy::shouldLogHealth(int $now, int $lastLoggedAt): bool` +- Produces: `WorkerProcessMetrics::collect(int $pid, int $startedAt, string $procRoot = '/proc'): array` + +- [ ] **Step 1: Write failing policy and metrics tests** + +Assert literal boundaries: a 40-second internal deadline, a 50-second outer timeout, health logging no more often than every 300 seconds, non-negative uptime/memory, counted entries in a synthetic `fd` directory, and unavailable socket/FD fields for a missing proc tree. + +- [ ] **Step 2: Run both tests and verify RED** + +Run: `php tests/WorkerRuntimePolicyTest.php && php tests/WorkerProcessMetricsTest.php` +Expected: FAIL because both production classes are missing. + +- [ ] **Step 3: Implement the pure policies** + +Use constants for 40/50/300-second boundaries. Collect `pid`, `uptimeSeconds`, `memoryBytes`, `peakMemoryBytes`, `openFdCount`, and `tcpSocketCount`; count socket symlink targets beginning with `socket:[` and return `null` when proc data cannot be read. + +- [ ] **Step 4: Run both tests and verify GREEN** + +Run: `php tests/WorkerRuntimePolicyTest.php && php tests/WorkerProcessMetricsTest.php` +Expected: both tests print `OK`. + +- [ ] **Step 5: Commit policies and diagnostics** + +```bash +git add Lib/WorkerRuntimePolicy.php Lib/WorkerProcessMetrics.php tests/WorkerRuntimePolicyTest.php tests/WorkerProcessMetricsTest.php +git commit -m "feat: add bounded worker diagnostics" +``` + +### Task 3: Harden safe.php and Its Cron Boundary + +**Files:** +- Create: `Lib/WorkerWatchdogRunner.php` +- Create: `Lib/ModuleWatchdogCommand.php` +- Modify: `bin/safe.php` +- Modify: `Lib/ExtendedCDRsConf.php` +- Create: `tests/WorkerWatchdogRunnerTest.php` +- Create: `tests/ModuleCronPolicyTest.php` + +**Interfaces:** +- Consumes: `WorkerWatchdogLease`, `WorkerRuntimePolicy`, `WorkerProcessMetrics` +- Produces: `WorkerWatchdogRunner::run(array $workers, callable $findPids, callable $startWorker, callable $signalDuplicates, callable $log): int` +- Produces: `ModuleWatchdogCommand::build(string $busybox, string $php, string $moduleDir, int $timeoutSeconds): string` +- Produces: cron command `busybox timeout 50 php -f /bin/safe.php` + +- [ ] **Step 1: Write a failing runner behavior test** + +Use real runner logic with local callbacks and assert: no PID starts one worker; one PID starts none; multiple PIDs signal every PID except the final canonical PID; each phase produces a sanitized diagnostic; a thrown callback returns a non-zero status. + +- [ ] **Step 2: Write a failing cron behavior test** + +Build a command with `ModuleWatchdogCommand::build(..., 1)` around a temporary PHP fixture that blocks for longer than one second. Execute the real command and assert timeout exit status `124` and elapsed time below three seconds. Also assert paths containing spaces are shell-escaped by executing a fixture from such a path. + +- [ ] **Step 3: Run both tests and verify RED** + +Run: `php tests/WorkerWatchdogRunnerTest.php && php tests/ModuleCronPolicyTest.php` +Expected: FAIL because runner and command builder are missing. + +- [ ] **Step 4: Implement the runner and safe.php adapter** + +Acquire the lease before any MikoPBX helper call. Install an optional `pcntl_async_signals(true)`/`pcntl_alarm(40)` handler, track the active phase, delegate worker decisions to the runner, catch `Throwable`, log structured events through `SystemMessages::sysLogMsg`, cancel the alarm, release the lease in `finally`, and exit with the runner status. + +- [ ] **Step 5: Add the outer cron timeout** + +Build the once-per-minute command through `ModuleWatchdogCommand` with the discovered BusyBox and PHP paths and `WorkerRuntimePolicy::outerTimeoutSeconds()`. Preserve output redirection and existing cleanup/report cron tasks. + +- [ ] **Step 6: Run tests and syntax checks; verify GREEN** + +Run: `php tests/WorkerWatchdogRunnerTest.php && php tests/ModuleCronPolicyTest.php && php -l bin/safe.php && php -l Lib/ExtendedCDRsConf.php` +Expected: tests print `OK`; syntax checks report no errors. + +- [ ] **Step 7: Commit watchdog integration** + +```bash +git add Lib/WorkerWatchdogRunner.php Lib/ModuleWatchdogCommand.php bin/safe.php Lib/ExtendedCDRsConf.php tests/WorkerWatchdogRunnerTest.php tests/ModuleCronPolicyTest.php +git commit -m "fix: bound module watchdog execution" +``` + +### Task 4: ConnectorDB Beanstalk Failure Boundary + +**Files:** +- Create: `Lib/WorkerFailureContext.php` +- Create: `Lib/TemporaryFileGuard.php` +- Modify: `bin/ConnectorDB.php` +- Create: `tests/WorkerFailureContextTest.php` +- Create: `tests/TemporaryFileGuardTest.php` + +**Interfaces:** +- Consumes: `WorkerProcessMetrics`, `WorkerRuntimePolicy` +- Produces: `WorkerFailureContext::make(string $operation, Throwable $error, array $metrics, int $elapsedMs): array` +- Produces: `TemporaryFileGuard::track(string $path): void`, `forget(string $path): void`, and idempotent `cleanup(): void` + +- [ ] **Step 1: Write the failing diagnostic-context test** + +Assert stable event name `worker_dependency_failure`, operation, exception class, elapsed time, and metrics. Pass secrets in the exception message and assert the context stores only a bounded generic category/message without paths, linked IDs, phone numbers, or request content. + +- [ ] **Step 2: Write the failing temporary-file test** + +Track two real temporary files, forget one after simulating ownership transfer, call cleanup twice, and assert only the still-owned file is removed. + +- [ ] **Step 3: Run both tests and verify RED** + +Run: `php tests/WorkerFailureContextTest.php && php tests/TemporaryFileGuardTest.php` +Expected: FAIL because both production classes are missing. + +- [ ] **Step 4: Implement context and file guard; verify GREEN** + +Run: `php tests/WorkerFailureContextTest.php && php tests/TemporaryFileGuardTest.php` +Expected: both tests print `OK`. + +- [ ] **Step 5: Integrate listener recovery** + +Wrap Beanstalk construction/subscription/wait in operation timers. On listener-side `Throwable`, log `WorkerFailureContext`, set restart intent, leave the loop, and return from `start()` so WorkerBase shutdown can complete. Emit rate-limited health events from the normal loop. + +- [ ] **Step 6: Integrate invoke cleanup without changing its contract** + +Track request and response files; forget each only after another component owns or the method has removed it; always call cleanup in `finally`; continue returning `[]` on any failure and log only the sanitized failure context. + +- [ ] **Step 7: Run focused tests and syntax checks** + +Run: `php tests/WorkerFailureContextTest.php && php tests/TemporaryFileGuardTest.php && php -l bin/ConnectorDB.php` +Expected: tests print `OK`; syntax check reports no errors. + +- [ ] **Step 8: Commit ConnectorDB recovery** + +```bash +git add Lib/WorkerFailureContext.php Lib/TemporaryFileGuard.php bin/ConnectorDB.php tests/WorkerFailureContextTest.php tests/TemporaryFileGuardTest.php +git commit -m "fix: recover connector after beanstalk failures" +``` + +### Task 5: Full Verification and Test-Server Deployment + +**Files:** +- Modify only if verification exposes a module-owned defect in files listed above. + +**Interfaces:** +- Consumes: completed module archive/install workflow and SSH access to `serber@boffart.miko.ru` +- Produces: verified installed module state with rollback evidence. + +- [ ] **Step 1: Run the complete local PHP suite** + +Run every `tests/*.php` file in sorted order and stop at the first failure. Then run `php -l` for every changed PHP file and `git diff --check`. + +- [ ] **Step 2: Capture read-only server baseline** + +Over SSH record MikoPBX/PHP/module versions, installed module checksum/status, cron line, PIDs and elapsed times for ConnectorDB/safe.php, Beanstalk availability, recent module/system errors, and an HTTP health response. Do not restart services during baseline. + +- [ ] **Step 3: Build and inspect the module package** + +Use the repository's existing build/package workflow. Inspect the archive to ensure new runtime classes are included and `.git`, tests, docs, `.DS_Store`, credentials, and local artifacts are excluded. + +- [ ] **Step 4: Back up and deploy on the test server** + +Create a timestamped backup of the installed module under `/tmp`, install/update using the module-safe workflow, and verify installed file hashes against the package. + +- [ ] **Step 5: Exercise singleton and worker recovery** + +Start concurrent `safe.php` invocations, verify one owner and successful overlap skips, observe at least two cron intervals, verify exactly one ConnectorDB, invoke an existing harmless ConnectorDB API, and confirm fresh health/phase diagnostics without secrets or log flooding. + +- [ ] **Step 6: Verify web health and rollback readiness** + +Confirm authenticated web/module endpoints respond, CDR synchronization advances or remains caught up, no stale lock/process appears, and the timestamped backup remains available. Restore the backup and restart only ConnectorDB if any module regression appears. + +- [ ] **Step 7: Record final evidence** + +Capture local test counts, server PIDs, cron command, relevant structured log events, HTTP status, deployed checksums, and `git status --short` for the final report. From 2f7dfba84c9e0122c3ebab1c90e258e1ad558319 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:13:34 +0300 Subject: [PATCH 03/10] chore: ignore local worktrees --- .gitignore | 1 + 1 file changed, 1 insertion(+) create mode 100644 .gitignore diff --git a/.gitignore b/.gitignore new file mode 100644 index 0000000..e458ed5 --- /dev/null +++ b/.gitignore @@ -0,0 +1 @@ +.worktrees/ From 4275da54d3580208623395309f46eace7d3a877e Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:14:29 +0300 Subject: [PATCH 04/10] feat: serialize module watchdog runs --- Lib/WorkerWatchdogLease.php | 61 +++++++++++++++++++++++++++++++ tests/WorkerWatchdogLeaseTest.php | 45 +++++++++++++++++++++++ 2 files changed, 106 insertions(+) create mode 100644 Lib/WorkerWatchdogLease.php create mode 100644 tests/WorkerWatchdogLeaseTest.php diff --git a/Lib/WorkerWatchdogLease.php b/Lib/WorkerWatchdogLease.php new file mode 100644 index 0000000..5767023 --- /dev/null +++ b/Lib/WorkerWatchdogLease.php @@ -0,0 +1,61 @@ +handle = $handle; + } + + public static function tryAcquire(string $path, int $pid, int $startedAt): ?self + { + $directory = dirname($path); + if (!is_dir($directory) && !mkdir($directory, 0775, true) && !is_dir($directory)) { + throw new \RuntimeException("Unable to create watchdog lock directory: {$directory}"); + } + + $handle = fopen($path, 'c+'); + if ($handle === false) { + throw new \RuntimeException("Unable to open watchdog lock: {$path}"); + } + if (!flock($handle, LOCK_EX | LOCK_NB)) { + fclose($handle); + return null; + } + + $payload = json_encode(['pid' => $pid, 'startedAt' => $startedAt], JSON_UNESCAPED_SLASHES); + if ($payload === false || !ftruncate($handle, 0) || fseek($handle, 0) !== 0 || fwrite($handle, $payload) === false) { + flock($handle, LOCK_UN); + fclose($handle); + throw new \RuntimeException("Unable to write watchdog lock diagnostics: {$path}"); + } + fflush($handle); + + return new self($handle); + } + + public function release(): void + { + if (!is_resource($this->handle)) { + return; + } + flock($this->handle, LOCK_UN); + fclose($this->handle); + $this->handle = null; + } + + public function __destruct() + { + $this->release(); + } +} diff --git a/tests/WorkerWatchdogLeaseTest.php b/tests/WorkerWatchdogLeaseTest.php new file mode 100644 index 0000000..d839b9f --- /dev/null +++ b/tests/WorkerWatchdogLeaseTest.php @@ -0,0 +1,45 @@ + 101, 'startedAt' => 1700000000], + json_decode((string) file_get_contents($path), true), + 'owner diagnostics are recorded' + ); + + $second = WorkerWatchdogLease::tryAcquire($path, 202, 1700000001); + assertLease(null, $second, 'contender skips while first owner holds lock'); + + $first->release(); + $third = WorkerWatchdogLease::tryAcquire($path, 303, 1700000002); + assertLease(true, $third instanceof WorkerWatchdogLease, 'released lease can be reacquired'); + $third->release(); + $third->release(); +} finally { + @unlink($path); +} + +echo "WorkerWatchdogLeaseTest: OK\n"; From 897202f3ee92658727663eae009af8c2375b2d9e Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:15:20 +0300 Subject: [PATCH 05/10] feat: add bounded worker diagnostics --- Lib/WorkerProcessMetrics.php | 41 +++++++++++++++++++++++++++ Lib/WorkerRuntimePolicy.php | 27 ++++++++++++++++++ tests/WorkerProcessMetricsTest.php | 45 ++++++++++++++++++++++++++++++ tests/WorkerRuntimePolicyTest.php | 23 +++++++++++++++ 4 files changed, 136 insertions(+) create mode 100644 Lib/WorkerProcessMetrics.php create mode 100644 Lib/WorkerRuntimePolicy.php create mode 100644 tests/WorkerProcessMetricsTest.php create mode 100644 tests/WorkerRuntimePolicyTest.php diff --git a/Lib/WorkerProcessMetrics.php b/Lib/WorkerProcessMetrics.php new file mode 100644 index 0000000..6eb4a9c --- /dev/null +++ b/Lib/WorkerProcessMetrics.php @@ -0,0 +1,41 @@ + + */ + public static function collect(int $pid, int $startedAt, string $procRoot = '/proc'): array + { + $fdDirectory = rtrim($procRoot, '/') . '/' . $pid . '/fd'; + $openFdCount = null; + $socketCount = null; + + if (is_dir($fdDirectory)) { + $entries = glob($fdDirectory . '/*'); + if (is_array($entries)) { + $openFdCount = count($entries); + $socketCount = 0; + foreach ($entries as $entry) { + $target = @readlink($entry); + if (is_string($target) && strpos($target, 'socket:[') === 0) { + ++$socketCount; + } + } + } + } + + return [ + 'pid' => $pid, + 'uptimeSeconds' => max(0, time() - $startedAt), + 'memoryBytes' => memory_get_usage(true), + 'peakMemoryBytes' => memory_get_peak_usage(true), + 'openFdCount' => $openFdCount, + 'tcpSocketCount' => $socketCount, + ]; + } +} diff --git a/Lib/WorkerRuntimePolicy.php b/Lib/WorkerRuntimePolicy.php new file mode 100644 index 0000000..3e7cada --- /dev/null +++ b/Lib/WorkerRuntimePolicy.php @@ -0,0 +1,27 @@ += self::HEALTH_LOG_INTERVAL_SECONDS; + } +} diff --git a/tests/WorkerProcessMetricsTest.php b/tests/WorkerProcessMetricsTest.php new file mode 100644 index 0000000..2f9be62 --- /dev/null +++ b/tests/WorkerProcessMetricsTest.php @@ -0,0 +1,45 @@ += 12, 'uptime is non-negative'); + assertProcessMetric(true, $metrics['memoryBytes'] > 0, 'memory is available'); + assertProcessMetric(true, $metrics['peakMemoryBytes'] > 0, 'peak memory is available'); + assertProcessMetric(3, $metrics['openFdCount'], 'open descriptor count'); + assertProcessMetric(2, $metrics['tcpSocketCount'], 'socket descriptor count'); + + $missing = WorkerProcessMetrics::collect($pid, time(), $root . '/missing'); + assertProcessMetric(null, $missing['openFdCount'], 'missing proc descriptor count'); + assertProcessMetric(null, $missing['tcpSocketCount'], 'missing proc socket count'); +} finally { + foreach (glob($root . '/' . $pid . '/fd/*') ?: [] as $entry) { + unlink($entry); + } + rmdir($root . '/' . $pid . '/fd'); + rmdir($root . '/' . $pid); + rmdir($root); +} + +echo "WorkerProcessMetricsTest: OK\n"; diff --git a/tests/WorkerRuntimePolicyTest.php b/tests/WorkerRuntimePolicyTest.php new file mode 100644 index 0000000..8c45c96 --- /dev/null +++ b/tests/WorkerRuntimePolicyTest.php @@ -0,0 +1,23 @@ + Date: Mon, 3 Aug 2026 14:17:27 +0300 Subject: [PATCH 06/10] fix: bound module watchdog execution --- Lib/ExtendedCDRsConf.php | 10 +++- Lib/ModuleWatchdogCommand.php | 24 ++++++++ Lib/WorkerWatchdogRunner.php | 62 +++++++++++++++++++ bin/safe.php | 96 +++++++++++++++++++++++++----- tests/ModuleCronPolicyTest.php | 49 +++++++++++++++ tests/WorkerWatchdogRunnerTest.php | 67 +++++++++++++++++++++ 6 files changed, 292 insertions(+), 16 deletions(-) create mode 100644 Lib/ModuleWatchdogCommand.php create mode 100644 Lib/WorkerWatchdogRunner.php create mode 100644 tests/ModuleCronPolicyTest.php create mode 100644 tests/WorkerWatchdogRunnerTest.php diff --git a/Lib/ExtendedCDRsConf.php b/Lib/ExtendedCDRsConf.php index 84092bf..7359ca2 100644 --- a/Lib/ExtendedCDRsConf.php +++ b/Lib/ExtendedCDRsConf.php @@ -15,7 +15,9 @@ use MikoPBX\Modules\Config\ConfigClass; use MikoPBX\PBXCoreREST\Lib\PBXApiResult; use Modules\ModuleExtendedCDRs\bin\ConnectorDB; +use Modules\ModuleExtendedCDRs\Lib\ModuleWatchdogCommand; use Modules\ModuleExtendedCDRs\Lib\RestAPI\Controllers\ApiController; +use Modules\ModuleExtendedCDRs\Lib\WorkerRuntimePolicy; use Modules\ModuleExtendedCDRs\Models\ReportSettings; class ExtendedCDRsConf extends ConfigClass @@ -131,7 +133,13 @@ public function createCronTasks(array &$tasks): void $busyboxPath= Util::which('busybox'); $tasks[] = "*/1 * * * * $busyboxPath find /storage/usbdisk*/mikopbx/tmp/ModuleExtendedCDRs/ -mmin +5 -type f -delete> /dev/null 2>&1".PHP_EOL; $phpPath = Util::which('php'); - $tasks[] = "*/1 * * * * $phpPath -f {$this->moduleDir}/bin/safe.php > /dev/null 2>&1".PHP_EOL; + $watchdogCommand = ModuleWatchdogCommand::build( + $busyboxPath, + $phpPath, + $this->moduleDir, + WorkerRuntimePolicy::outerTimeoutSeconds() + ); + $tasks[] = "*/1 * * * * $watchdogCommand > /dev/null 2>&1".PHP_EOL; $reportsData = ReportSettings::find('sendingScheduledReport=1'); foreach ($reportsData as $settings) { diff --git a/Lib/ModuleWatchdogCommand.php b/Lib/ModuleWatchdogCommand.php new file mode 100644 index 0000000..2712223 --- /dev/null +++ b/Lib/ModuleWatchdogCommand.php @@ -0,0 +1,24 @@ + 1) { + $duplicates = array_slice($pids, 0, -1); + $signalDuplicates($duplicates); + $outcome = 'duplicates_signalled'; + } else { + $outcome = 'running'; + } + + $log(self::event($worker, $pids, $outcome, $startedAt)); + } catch (Throwable $error) { + $log(self::event($worker, [], 'failed', $startedAt) + [ + 'errorClass' => get_class($error), + ]); + return 1; + } + } + + return 0; + } + + /** + * @return array + */ + private static function event(string $worker, array $pids, string $outcome, float $startedAt): array + { + return [ + 'event' => 'worker_watchdog_phase', + 'worker' => $worker, + 'pids' => $pids, + 'outcome' => $outcome, + 'elapsedMs' => (int) round((microtime(true) - $startedAt) * 1000), + ]; + } +} diff --git a/bin/safe.php b/bin/safe.php index e8c46d4..6a19bbe 100644 --- a/bin/safe.php +++ b/bin/safe.php @@ -22,27 +22,93 @@ use MikoPBX\Core\System\SystemMessages; use MikoPBX\Modules\PbxExtensionUtils; use Modules\ModuleExtendedCDRs\Lib\ExtendedCDRsConf; +use Modules\ModuleExtendedCDRs\Lib\WorkerRuntimePolicy; +use Modules\ModuleExtendedCDRs\Lib\WorkerWatchdogLease; +use Modules\ModuleExtendedCDRs\Lib\WorkerWatchdogRunner; require_once 'Globals.php'; +require_once dirname(__DIR__) . '/vendor/autoload.php'; $moduleEnable = PbxExtensionUtils::isEnabled('ModuleExtendedCDRs'); if(!$moduleEnable){ exit(1); } -$conf = new ExtendedCDRsConf(); -$workers = $conf->getModuleWorkers(); -foreach ($workers as $workerData) { - $WorkerPID = Processes::getPidOfProcess($workerData['worker']); - print_r($WorkerPID.PHP_EOL); - if (empty($WorkerPID)) { - Processes::processPHPWorker($workerData['worker']); - SystemMessages::sysLogMsg('ModuleExtendedCDRs_SAFE', "Service {$workerData['worker']} started.", LOG_NOTICE); - }else{ - // Проверка дубликата процесса. - $allButLast = array_slice(explode(' ', $WorkerPID), 0, -1); - if(!empty($allButLast)){ - // Завершаем дубликаты процессов. - $bbPath = Util::which('busybox'); - shell_exec("$bbPath kill -SIGUSR2 ". implode(" ", $allButLast)); + +$startedAt = time(); +$lockPath = '/tmp/ModuleExtendedCDRs/worker-watchdog.lock'; +$lease = null; +$activePhase = 'acquire_lock'; + +try { + $lease = WorkerWatchdogLease::tryAcquire($lockPath, getmypid(), $startedAt); + if ($lease === null) { + SystemMessages::sysLogMsg( + 'ModuleExtendedCDRs_SAFE', + json_encode(['event' => 'worker_watchdog_skipped', 'reason' => 'lock_busy']), + LOG_NOTICE + ); + exit(0); + } + + if (function_exists('pcntl_async_signals') && function_exists('pcntl_alarm')) { + pcntl_async_signals(true); + pcntl_signal(SIGALRM, static function () use (&$activePhase, $startedAt): void { + SystemMessages::sysLogMsg( + 'ModuleExtendedCDRs_SAFE', + json_encode([ + 'event' => 'worker_watchdog_timeout', + 'phase' => $activePhase, + 'elapsedMs' => (time() - $startedAt) * 1000, + ]), + LOG_ERR + ); + exit(124); + }); + pcntl_alarm(WorkerRuntimePolicy::watchdogDeadlineSeconds()); + } + + $activePhase = 'load_workers'; + $conf = new ExtendedCDRsConf(); + $workers = array_column($conf->getModuleWorkers(), 'worker'); + $busyboxPath = Util::which('busybox'); + + $status = WorkerWatchdogRunner::run( + $workers, + static function (string $worker) use (&$activePhase): string { + $activePhase = 'find_worker'; + return (string) Processes::getPidOfProcess($worker); + }, + static function (string $worker) use (&$activePhase): void { + $activePhase = 'start_worker'; + Processes::processPHPWorker($worker); + }, + static function (array $duplicates) use (&$activePhase, $busyboxPath): void { + $activePhase = 'signal_duplicates'; + $pids = implode(' ', array_map('intval', $duplicates)); + shell_exec(escapeshellarg($busyboxPath) . ' kill -SIGUSR2 ' . $pids); + }, + static function (array $event): void { + SystemMessages::sysLogMsg('ModuleExtendedCDRs_SAFE', json_encode($event), LOG_NOTICE); } + ); +} catch (Throwable $error) { + SystemMessages::sysLogMsg( + 'ModuleExtendedCDRs_SAFE', + json_encode([ + 'event' => 'worker_watchdog_failed', + 'phase' => $activePhase, + 'errorClass' => get_class($error), + 'elapsedMs' => (time() - $startedAt) * 1000, + ]), + LOG_ERR + ); + $status = 1; +} finally { + if (function_exists('pcntl_alarm')) { + pcntl_alarm(0); + } + if ($lease instanceof WorkerWatchdogLease) { + $lease->release(); } } + +exit($status ?? 1); diff --git a/tests/ModuleCronPolicyTest.php b/tests/ModuleCronPolicyTest.php new file mode 100644 index 0000000..60a7fd2 --- /dev/null +++ b/tests/ModuleCronPolicyTest.php @@ -0,0 +1,49 @@ + " . escapeshellarg($capture) . "\n"); +file_put_contents($fakePhp, "#!/bin/sh\nexit 0\n"); +file_put_contents($moduleDir . '/bin/safe.php', " '', + 'WorkerB' => '42', + 'WorkerC' => '11 22 33', +]; + +$status = WorkerWatchdogRunner::run( + array_keys($pids), + static function (string $worker) use ($pids): string { + return $pids[$worker]; + }, + static function (string $worker) use (&$started): void { + $started[] = $worker; + }, + static function (array $duplicates) use (&$signalled): void { + $signalled[] = $duplicates; + }, + static function (array $event) use (&$events): void { + $events[] = $event; + } +); + +assertWatchdogRunner(0, $status, 'successful run status'); +assertWatchdogRunner(['WorkerA'], $started, 'missing worker is started once'); +assertWatchdogRunner([[11, 22]], $signalled, 'all duplicate pids except canonical final pid are signalled'); +assertWatchdogRunner('worker_watchdog_phase', $events[0]['event'], 'stable diagnostic event'); +assertWatchdogRunner('WorkerA', $events[0]['worker'], 'diagnostic identifies worker'); +assertWatchdogRunner(true, isset($events[0]['elapsedMs']), 'diagnostic contains duration'); + +$failureEvents = []; +$failureStatus = WorkerWatchdogRunner::run( + ['BrokenWorker'], + static function (): string { + throw new RuntimeException('secret /recordings/call.wav 79990001122'); + }, + static function (): void { + }, + static function (): void { + }, + static function (array $event) use (&$failureEvents): void { + $failureEvents[] = $event; + } +); +assertWatchdogRunner(1, $failureStatus, 'dependency exception returns failure'); +assertWatchdogRunner('failed', $failureEvents[0]['outcome'], 'failure outcome is logged'); +assertWatchdogRunner(false, isset($failureEvents[0]['message']), 'dependency message is not logged'); + +echo "WorkerWatchdogRunnerTest: OK\n"; From 4f8c17749130f3071500210fc711c107fef1d595 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:19:12 +0300 Subject: [PATCH 07/10] fix: recover connector after beanstalk failures --- Lib/TemporaryFileGuard.php | 38 +++++++++++++++ Lib/WorkerFailureContext.php | 41 ++++++++++++++++ bin/ConnectorDB.php | 77 ++++++++++++++++++++++++++---- tests/TemporaryFileGuardTest.php | 34 +++++++++++++ tests/WorkerFailureContextTest.php | 40 ++++++++++++++++ 5 files changed, 221 insertions(+), 9 deletions(-) create mode 100644 Lib/TemporaryFileGuard.php create mode 100644 Lib/WorkerFailureContext.php create mode 100644 tests/TemporaryFileGuardTest.php create mode 100644 tests/WorkerFailureContextTest.php diff --git a/Lib/TemporaryFileGuard.php b/Lib/TemporaryFileGuard.php new file mode 100644 index 0000000..97aaca0 --- /dev/null +++ b/Lib/TemporaryFileGuard.php @@ -0,0 +1,38 @@ + */ + private array $paths = []; + + public function track(string $path): void + { + if ($path !== '') { + $this->paths[$path] = true; + } + } + + public function forget(string $path): void + { + unset($this->paths[$path]); + } + + public function cleanup(): void + { + foreach (array_keys($this->paths) as $path) { + if (is_file($path) || is_link($path)) { + @unlink($path); + } + unset($this->paths[$path]); + } + } + + public function __destruct() + { + $this->cleanup(); + } +} diff --git a/Lib/WorkerFailureContext.php b/Lib/WorkerFailureContext.php new file mode 100644 index 0000000..57127a0 --- /dev/null +++ b/Lib/WorkerFailureContext.php @@ -0,0 +1,41 @@ + + */ + public static function make(string $operation, Throwable $error, array $metrics, int $elapsedMs): array + { + $context = [ + 'event' => 'worker_dependency_failure', + 'operation' => $operation, + 'errorClass' => get_class($error), + 'errorCategory' => 'dependency_failure', + 'elapsedMs' => max(0, $elapsedMs), + ]; + + foreach (self::METRIC_KEYS as $key) { + if (array_key_exists($key, $metrics)) { + $context[$key] = $metrics[$key]; + } + } + + return $context; + } +} diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index 779dc56..8d0b75a 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -35,6 +35,10 @@ use Modules\ModuleExtendedCDRs\Lib\QuarantinePolicy; use Modules\ModuleExtendedCDRs\Lib\BatchLogContext; use Modules\ModuleExtendedCDRs\Lib\SyncPolicy; +use Modules\ModuleExtendedCDRs\Lib\TemporaryFileGuard; +use Modules\ModuleExtendedCDRs\Lib\WorkerFailureContext; +use Modules\ModuleExtendedCDRs\Lib\WorkerProcessMetrics; +use Modules\ModuleExtendedCDRs\Lib\WorkerRuntimePolicy; use Exception; use Modules\ModuleExtendedCDRs\Lib\MikoPBXVersion; use Modules\ModuleExtendedCDRs\Lib\Providers\CdrDbProvider; @@ -68,6 +72,8 @@ class ConnectorDB extends WorkerBase private array $oversizedPending = []; private int $oversizedCacheTime = 0; private int $oversizedPruneTime = 0; + private int $workerStartedAt = 0; + private int $lastHealthLogAt = 0; /** * Белый список методов, разрешённых для вызова через onEvents/invoke. @@ -105,15 +111,22 @@ public function signalHandler(int $signal): void */ public function start($argv):void { + $this->workerStartedAt = time(); $this->logger = new Logger('ConnectorDB', 'ModuleExtendedCDRs'); $this->mp3TagService = new Mp3TagService(dirname(__DIR__)); $this->logger->writeInfo('Starting...'); $this->ensureDailyStatsTableExists(); $this->ensureOversizedTableExists(); $this->updateSettings(); - $beanstalk = new BeanstalkClient(self::class); - $beanstalk->subscribe(self::class, [$this, 'onEvents']); - $beanstalk->subscribe($this->makePingTubeName(self::class), [$this, 'pingCallBack']); + $operationStartedAt = microtime(true); + try { + $beanstalk = new BeanstalkClient(self::class); + $beanstalk->subscribe(self::class, [$this, 'onEvents']); + $beanstalk->subscribe($this->makePingTubeName(self::class), [$this, 'pingCallBack']); + } catch (Throwable $exception) { + $this->logDependencyFailure('beanstalk_subscribe', $exception, $operationStartedAt); + return; + } while ($this->needRestart === false) { try { $this->syncCdrData(true); @@ -127,11 +140,37 @@ public function start($argv):void ); $this->nextSyncDelay = SyncPolicy::ERROR_DELAY_SECONDS; } - $beanstalk->wait(max(1, $this->nextSyncDelay)); + $operationStartedAt = microtime(true); + try { + $beanstalk->wait(max(1, $this->nextSyncDelay)); + } catch (Throwable $exception) { + $this->logDependencyFailure('beanstalk_wait', $exception, $operationStartedAt); + return; + } + $this->logWorkerHealthIfDue(); $this->logger->rotate(); } } + private function logDependencyFailure(string $operation, Throwable $exception, float $startedAt): void + { + $metrics = WorkerProcessMetrics::collect(getmypid(), $this->workerStartedAt ?: time()); + $elapsedMs = (int) round((microtime(true) - $startedAt) * 1000); + $this->logger->writeError(WorkerFailureContext::make($operation, $exception, $metrics, $elapsedMs)); + } + + private function logWorkerHealthIfDue(): void + { + $now = time(); + if (!WorkerRuntimePolicy::shouldLogHealth($now, $this->lastHealthLogAt)) { + return; + } + $this->lastHealthLogAt = $now; + $this->logger->writeInfo([ + 'event' => 'worker_health', + ] + WorkerProcessMetrics::collect(getmypid(), $this->workerStartedAt)); + } + /** * Получение настроек модуля. @@ -277,24 +316,44 @@ public static function invoke(string $function, array $args = [], bool $retVal = 'function' => $function, 'args' => $args ]; - $client = new BeanstalkClient(self::class); + if ($retVal) { + $req['need-ret'] = true; + } $object = []; + $guard = new TemporaryFileGuard(); + $operationStartedAt = microtime(true); try { + $client = new BeanstalkClient(self::class); + $pathToData = self::saveInTmpFile($req); + $guard->track($pathToData); if($retVal){ - $req['need-ret'] = true; - $pathToData = self::saveInTmpFile($req); $result = $client->request($pathToData, 20); + $guard->track((string) $result); }else{ - $pathToData = self::saveInTmpFile($req); $client->publish($pathToData); + $guard->forget($pathToData); return []; } if(file_exists($result)){ $object = json_decode(file_get_contents($result), true, 512, JSON_THROW_ON_ERROR); unlink($result); + $guard->forget($result); } - } catch (\Throwable $e) { + } catch (Throwable $e) { $object = []; + try { + $logger = new Logger('ConnectorDB', 'ModuleExtendedCDRs'); + $logger->writeError(WorkerFailureContext::make( + 'beanstalk_request', + $e, + WorkerProcessMetrics::collect(getmypid(), time()), + (int) round((microtime(true) - $operationStartedAt) * 1000) + )); + } catch (Throwable $loggingError) { + // Diagnostics must never change the public invoke() contract. + } + } finally { + $guard->cleanup(); } return $object; } diff --git a/tests/TemporaryFileGuardTest.php b/tests/TemporaryFileGuardTest.php new file mode 100644 index 0000000..23a6a00 --- /dev/null +++ b/tests/TemporaryFileGuardTest.php @@ -0,0 +1,34 @@ +track($owned); +$guard->track($transferred); +$guard->forget($transferred); +$guard->cleanup(); +$guard->cleanup(); + +assertTemporaryGuard(false, file_exists($owned), 'owned file is removed'); +assertTemporaryGuard(true, file_exists($transferred), 'transferred file is preserved'); +unlink($transferred); + +echo "TemporaryFileGuardTest: OK\n"; diff --git a/tests/WorkerFailureContextTest.php b/tests/WorkerFailureContextTest.php new file mode 100644 index 0000000..8ce9b34 --- /dev/null +++ b/tests/WorkerFailureContextTest.php @@ -0,0 +1,40 @@ + 321, + 'uptimeSeconds' => 600, + 'memoryBytes' => 1048576, + 'peakMemoryBytes' => 2097152, + 'openFdCount' => 17, + 'tcpSocketCount' => 4, + 'recordingfile' => '/secret/from-metrics.wav', +]; +$error = new RuntimeException('failed /storage/recordings/call.wav linked-secret 79990001122'); +$context = WorkerFailureContext::make('beanstalk_wait', $error, $metrics, 1250); + +assertFailureContext('worker_dependency_failure', $context['event'], 'stable event'); +assertFailureContext('beanstalk_wait', $context['operation'], 'operation'); +assertFailureContext(RuntimeException::class, $context['errorClass'], 'exception class'); +assertFailureContext('dependency_failure', $context['errorCategory'], 'generic error category'); +assertFailureContext(1250, $context['elapsedMs'], 'elapsed milliseconds'); +assertFailureContext(321, $context['pid'], 'allowed metric'); +assertFailureContext(false, isset($context['message']), 'exception message is excluded'); +assertFailureContext(false, isset($context['recordingfile']), 'unknown metric is excluded'); +assertFailureContext(false, strpos(json_encode($context), '79990001122') !== false, 'phone number is absent'); +assertFailureContext(false, strpos(json_encode($context), 'recordings') !== false, 'path is absent'); + +echo "WorkerFailureContextTest: OK\n"; From 1d427645b2c94baf180d058dc792884d29e127fa Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:29:10 +0300 Subject: [PATCH 08/10] fix: secure watchdog and worker diagnostics --- Lib/WorkerDependencyException.php | 24 +++++++++++++++ Lib/WorkerEventContext.php | 27 +++++++++++++++++ Lib/WorkerLogRateLimiter.php | 34 +++++++++++++++++++++ Lib/WorkerWatchdogLease.php | 12 +++++++- bin/ConnectorDB.php | 40 ++++++++++++++++++++----- bin/safe.php | 16 ++++++---- tests/WorkerDependencyExceptionTest.php | 18 +++++++++++ tests/WorkerEventContextTest.php | 36 ++++++++++++++++++++++ tests/WorkerLogRateLimiterTest.php | 26 ++++++++++++++++ tests/WorkerWatchdogLeaseTest.php | 21 ++++++++++--- 10 files changed, 236 insertions(+), 18 deletions(-) create mode 100644 Lib/WorkerDependencyException.php create mode 100644 Lib/WorkerEventContext.php create mode 100644 Lib/WorkerLogRateLimiter.php create mode 100644 tests/WorkerDependencyExceptionTest.php create mode 100644 tests/WorkerEventContextTest.php create mode 100644 tests/WorkerLogRateLimiterTest.php diff --git a/Lib/WorkerDependencyException.php b/Lib/WorkerDependencyException.php new file mode 100644 index 0000000..ba9d714 --- /dev/null +++ b/Lib/WorkerDependencyException.php @@ -0,0 +1,24 @@ +operation = $operation; + } + + public function operation(): string + { + return $this->operation; + } +} diff --git a/Lib/WorkerEventContext.php b/Lib/WorkerEventContext.php new file mode 100644 index 0000000..e6e5b9f --- /dev/null +++ b/Lib/WorkerEventContext.php @@ -0,0 +1,27 @@ + + */ + public static function make(array $request, string $outcome): array + { + return [ + 'event' => 'worker_event', + 'action' => self::identifier((string) ($request['action'] ?? '')), + 'function' => self::identifier((string) ($request['function'] ?? '')), + 'needsReply' => isset($request['need-ret']), + 'outcome' => self::identifier($outcome), + ]; + } + + private static function identifier(string $value): string + { + return substr((string) preg_replace('/[^A-Za-z0-9_-]/', '', $value), 0, 64); + } +} diff --git a/Lib/WorkerLogRateLimiter.php b/Lib/WorkerLogRateLimiter.php new file mode 100644 index 0000000..7034982 --- /dev/null +++ b/Lib/WorkerLogRateLimiter.php @@ -0,0 +1,34 @@ += $intervalSeconds; + if ($allowed) { + ftruncate($handle, 0); + fseek($handle, 0); + fwrite($handle, (string) $now); + fflush($handle); + } + + flock($handle, LOCK_UN); + fclose($handle); + return $allowed; + } +} diff --git a/Lib/WorkerWatchdogLease.php b/Lib/WorkerWatchdogLease.php index 5767023..349aa26 100644 --- a/Lib/WorkerWatchdogLease.php +++ b/Lib/WorkerWatchdogLease.php @@ -20,14 +20,24 @@ private function __construct($handle) public static function tryAcquire(string $path, int $pid, int $startedAt): ?self { $directory = dirname($path); - if (!is_dir($directory) && !mkdir($directory, 0775, true) && !is_dir($directory)) { + if (!is_dir($directory) && !mkdir($directory, 0700, true) && !is_dir($directory)) { throw new \RuntimeException("Unable to create watchdog lock directory: {$directory}"); } + if (is_link($path)) { + throw new \RuntimeException("Refusing symlink watchdog lock: {$path}"); + } $handle = fopen($path, 'c+'); if ($handle === false) { throw new \RuntimeException("Unable to open watchdog lock: {$path}"); } + $pathStat = lstat($path); + $handleStat = fstat($handle); + if ($pathStat === false || $handleStat === false + || $pathStat['dev'] !== $handleStat['dev'] || $pathStat['ino'] !== $handleStat['ino']) { + fclose($handle); + throw new \RuntimeException("Watchdog lock path changed while opening: {$path}"); + } if (!flock($handle, LOCK_EX | LOCK_NB)) { fclose($handle); return null; diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index 8d0b75a..1cbc67f 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -36,6 +36,8 @@ use Modules\ModuleExtendedCDRs\Lib\BatchLogContext; use Modules\ModuleExtendedCDRs\Lib\SyncPolicy; use Modules\ModuleExtendedCDRs\Lib\TemporaryFileGuard; +use Modules\ModuleExtendedCDRs\Lib\WorkerDependencyException; +use Modules\ModuleExtendedCDRs\Lib\WorkerEventContext; use Modules\ModuleExtendedCDRs\Lib\WorkerFailureContext; use Modules\ModuleExtendedCDRs\Lib\WorkerProcessMetrics; use Modules\ModuleExtendedCDRs\Lib\WorkerRuntimePolicy; @@ -121,6 +123,12 @@ public function start($argv):void $operationStartedAt = microtime(true); try { $beanstalk = new BeanstalkClient(self::class); + } catch (Throwable $exception) { + $this->logDependencyFailure('beanstalk_connect', $exception, $operationStartedAt); + return; + } + $operationStartedAt = microtime(true); + try { $beanstalk->subscribe(self::class, [$this, 'onEvents']); $beanstalk->subscribe($this->makePingTubeName(self::class), [$this, 'pingCallBack']); } catch (Throwable $exception) { @@ -144,7 +152,13 @@ public function start($argv):void try { $beanstalk->wait(max(1, $this->nextSyncDelay)); } catch (Throwable $exception) { - $this->logDependencyFailure('beanstalk_wait', $exception, $operationStartedAt); + $operation = $exception instanceof WorkerDependencyException + ? $exception->operation() + : 'beanstalk_wait'; + $cause = $exception->getPrevious() instanceof Throwable + ? $exception->getPrevious() + : $exception; + $this->logDependencyFailure($operation, $cause, $operationStartedAt); return; } $this->logWorkerHealthIfDue(); @@ -227,9 +241,10 @@ public function onEvents($tube): void $data = json_decode(file_get_contents($pathToData), true, 512, JSON_THROW_ON_ERROR); unlink($pathToData); } - $this->logger->writeInfo(['data'=> $data, 'pathToData' => $pathToData], 'onEvents'); }catch (Throwable $exception){ - $this->logger->writeError("Throwable:".$exception->getMessage(). ' Line: '.$exception->getLine()); + $this->logger->writeError(WorkerEventContext::make($data, 'invalid_request') + [ + 'errorClass' => get_class($exception), + ]); return; } $action = $data['action']??''; @@ -245,16 +260,27 @@ public function onEvents($tube): void $res_data = $this->$funcName(...$data['args']??[]); } }else{ - $this->logger->writeError($data); + $this->logger->writeError(WorkerEventContext::make($data, 'rejected')); } if(isset($data['need-ret'])){ $res_data = self::saveInTmpFile($res_data); - $tube->reply($res_data); + $replyGuard = new TemporaryFileGuard(); + $replyGuard->track($res_data); + try { + $tube->reply($res_data); + $replyGuard->forget($res_data); + } catch (Throwable $exception) { + throw new WorkerDependencyException('beanstalk_reply', $exception); + } } - $this->logger->writeInfo(['data'=> $data, 'result' => $res_data], 'invoke'); + $this->logger->writeInfo(WorkerEventContext::make($data, 'completed')); } + }catch (WorkerDependencyException $exception){ + throw $exception; }catch (Throwable $exception){ - $this->logger->writeError($data, "Throwable:".$exception->getMessage(). ' Line: '.$exception->getLine()); + $this->logger->writeError(WorkerEventContext::make($data, 'failed') + [ + 'errorClass' => get_class($exception), + ]); return; } } diff --git a/bin/safe.php b/bin/safe.php index 6a19bbe..eaca917 100644 --- a/bin/safe.php +++ b/bin/safe.php @@ -23,6 +23,7 @@ use MikoPBX\Modules\PbxExtensionUtils; use Modules\ModuleExtendedCDRs\Lib\ExtendedCDRsConf; use Modules\ModuleExtendedCDRs\Lib\WorkerRuntimePolicy; +use Modules\ModuleExtendedCDRs\Lib\WorkerLogRateLimiter; use Modules\ModuleExtendedCDRs\Lib\WorkerWatchdogLease; use Modules\ModuleExtendedCDRs\Lib\WorkerWatchdogRunner; require_once 'Globals.php'; @@ -34,18 +35,21 @@ } $startedAt = time(); -$lockPath = '/tmp/ModuleExtendedCDRs/worker-watchdog.lock'; +$lockPath = '/var/run/php-workers/ModuleExtendedCDRs-watchdog.lock'; +$skipLogPath = '/var/run/php-workers/ModuleExtendedCDRs-watchdog-skip.log'; $lease = null; $activePhase = 'acquire_lock'; try { $lease = WorkerWatchdogLease::tryAcquire($lockPath, getmypid(), $startedAt); if ($lease === null) { - SystemMessages::sysLogMsg( - 'ModuleExtendedCDRs_SAFE', - json_encode(['event' => 'worker_watchdog_skipped', 'reason' => 'lock_busy']), - LOG_NOTICE - ); + if (WorkerLogRateLimiter::shouldLog($skipLogPath, time(), 300)) { + SystemMessages::sysLogMsg( + 'ModuleExtendedCDRs_SAFE', + json_encode(['event' => 'worker_watchdog_skipped', 'reason' => 'lock_busy']), + LOG_NOTICE + ); + } exit(0); } diff --git a/tests/WorkerDependencyExceptionTest.php b/tests/WorkerDependencyExceptionTest.php new file mode 100644 index 0000000..a989625 --- /dev/null +++ b/tests/WorkerDependencyExceptionTest.php @@ -0,0 +1,18 @@ +operation() !== 'beanstalk_reply' || $error->getPrevious() !== $previous) { + throw new RuntimeException('Dependency exception must retain operation and cause'); +} +if (strpos($error->getMessage(), 'secret') !== false) { + throw new RuntimeException('Dependency exception message must be sanitized'); +} + +echo "WorkerDependencyExceptionTest: OK\n"; diff --git a/tests/WorkerEventContextTest.php b/tests/WorkerEventContextTest.php new file mode 100644 index 0000000..65f57e4 --- /dev/null +++ b/tests/WorkerEventContextTest.php @@ -0,0 +1,36 @@ + 'invoke', + 'function' => 'getRecordingPathByID', + 'args' => ['79990001122', '/storage/recordings/secret.wav'], + 'need-ret' => true, + 'linkedid' => 'secret-linked-id', +]; +$context = WorkerEventContext::make($request, 'completed'); + +assertWorkerEvent('worker_event', $context['event'], 'stable event'); +assertWorkerEvent('invoke', $context['action'], 'action'); +assertWorkerEvent('getRecordingPathByID', $context['function'], 'function'); +assertWorkerEvent(true, $context['needsReply'], 'reply flag'); +assertWorkerEvent('completed', $context['outcome'], 'outcome'); +$encoded = json_encode($context); +assertWorkerEvent(false, strpos($encoded, '79990001122') !== false, 'arguments are excluded'); +assertWorkerEvent(false, strpos($encoded, 'recordings') !== false, 'paths are excluded'); +assertWorkerEvent(false, strpos($encoded, 'linked-id') !== false, 'linked IDs are excluded'); + +echo "WorkerEventContextTest: OK\n"; diff --git a/tests/WorkerLogRateLimiterTest.php b/tests/WorkerLogRateLimiterTest.php new file mode 100644 index 0000000..fd939bf --- /dev/null +++ b/tests/WorkerLogRateLimiterTest.php @@ -0,0 +1,26 @@ +release(); $third->release(); + + $target = $directory . '/target'; + $link = $directory . '/symlink.lock'; + file_put_contents($target, 'must-not-change'); + symlink($target, $link); + try { + WorkerWatchdogLease::tryAcquire($link, 404, 1700000003); + throw new RuntimeException('symlink lock path must be rejected'); + } catch (RuntimeException $error) { + assertLease('must-not-change', file_get_contents($target), 'symlink target remains unchanged'); + } } finally { @unlink($path); + @unlink($link ?? ''); + @unlink($target ?? ''); + @rmdir($directory); } echo "WorkerWatchdogLeaseTest: OK\n"; From af6bb57075f380db27adb39e3c11a4896e43bb56 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:30:20 +0300 Subject: [PATCH 09/10] fix: keep watchdog lock in root runtime --- Lib/WorkerWatchdogLease.php | 4 ++++ bin/safe.php | 4 ++-- tests/WorkerWatchdogLeaseTest.php | 12 ++++++++++++ 3 files changed, 18 insertions(+), 2 deletions(-) diff --git a/Lib/WorkerWatchdogLease.php b/Lib/WorkerWatchdogLease.php index 349aa26..c8bbd6b 100644 --- a/Lib/WorkerWatchdogLease.php +++ b/Lib/WorkerWatchdogLease.php @@ -23,6 +23,10 @@ public static function tryAcquire(string $path, int $pid, int $startedAt): ?self if (!is_dir($directory) && !mkdir($directory, 0700, true) && !is_dir($directory)) { throw new \RuntimeException("Unable to create watchdog lock directory: {$directory}"); } + $directoryMode = fileperms($directory); + if (is_link($directory) || $directoryMode === false || (($directoryMode & 0022) !== 0)) { + throw new \RuntimeException("Refusing insecure watchdog lock directory: {$directory}"); + } if (is_link($path)) { throw new \RuntimeException("Refusing symlink watchdog lock: {$path}"); } diff --git a/bin/safe.php b/bin/safe.php index eaca917..4a5aa40 100644 --- a/bin/safe.php +++ b/bin/safe.php @@ -35,8 +35,8 @@ } $startedAt = time(); -$lockPath = '/var/run/php-workers/ModuleExtendedCDRs-watchdog.lock'; -$skipLogPath = '/var/run/php-workers/ModuleExtendedCDRs-watchdog-skip.log'; +$lockPath = '/var/run/ModuleExtendedCDRs/watchdog.lock'; +$skipLogPath = '/var/run/ModuleExtendedCDRs/watchdog-skip.log'; $lease = null; $activePhase = 'acquire_lock'; diff --git a/tests/WorkerWatchdogLeaseTest.php b/tests/WorkerWatchdogLeaseTest.php index b953ee1..4d0573e 100644 --- a/tests/WorkerWatchdogLeaseTest.php +++ b/tests/WorkerWatchdogLeaseTest.php @@ -48,11 +48,23 @@ function assertLease($expected, $actual, string $message): void } catch (RuntimeException $error) { assertLease('must-not-change', file_get_contents($target), 'symlink target remains unchanged'); } + + $unsafeDirectory = sys_get_temp_dir() . '/extended-cdr-unsafe-' . bin2hex(random_bytes(6)); + mkdir($unsafeDirectory, 0777); + chmod($unsafeDirectory, 0777); + try { + WorkerWatchdogLease::tryAcquire($unsafeDirectory . '/watchdog.lock', 505, 1700000004); + throw new RuntimeException('group/world-writable lock directory must be rejected'); + } catch (RuntimeException $error) { + assertLease(false, file_exists($unsafeDirectory . '/watchdog.lock'), 'unsafe directory gets no lock file'); + } } finally { @unlink($path); @unlink($link ?? ''); @unlink($target ?? ''); @rmdir($directory); + @unlink(($unsafeDirectory ?? '') . '/watchdog.lock'); + @rmdir($unsafeDirectory ?? ''); } echo "WorkerWatchdogLeaseTest: OK\n"; From 02b23d3318328013a3db171d312533e5710b0b40 Mon Sep 17 00:00:00 2001 From: boffart <5922739+boffart@users.noreply.github.com> Date: Mon, 3 Aug 2026 14:59:24 +0300 Subject: [PATCH 10/10] fix: classify connector dependency failures --- Lib/WorkerFailureContext.php | 8 ++++++++ bin/ConnectorDB.php | 4 +++- tests/WorkerFailureContextTest.php | 3 +++ 3 files changed, 14 insertions(+), 1 deletion(-) diff --git a/Lib/WorkerFailureContext.php b/Lib/WorkerFailureContext.php index 57127a0..96b1101 100644 --- a/Lib/WorkerFailureContext.php +++ b/Lib/WorkerFailureContext.php @@ -17,6 +17,14 @@ final class WorkerFailureContext 'tcpSocketCount', ]; + public static function invokeOperation(bool $clientCreated, bool $expectsReply): string + { + if (!$clientCreated) { + return 'beanstalk_connect'; + } + return $expectsReply ? 'beanstalk_request' : 'beanstalk_publish'; + } + /** * @return array */ diff --git a/bin/ConnectorDB.php b/bin/ConnectorDB.php index 1cbc67f..0112ff5 100644 --- a/bin/ConnectorDB.php +++ b/bin/ConnectorDB.php @@ -348,8 +348,10 @@ public static function invoke(string $function, array $args = [], bool $retVal = $object = []; $guard = new TemporaryFileGuard(); $operationStartedAt = microtime(true); + $clientCreated = false; try { $client = new BeanstalkClient(self::class); + $clientCreated = true; $pathToData = self::saveInTmpFile($req); $guard->track($pathToData); if($retVal){ @@ -370,7 +372,7 @@ public static function invoke(string $function, array $args = [], bool $retVal = try { $logger = new Logger('ConnectorDB', 'ModuleExtendedCDRs'); $logger->writeError(WorkerFailureContext::make( - 'beanstalk_request', + WorkerFailureContext::invokeOperation($clientCreated, $retVal), $e, WorkerProcessMetrics::collect(getmypid(), time()), (int) round((microtime(true) - $operationStartedAt) * 1000) diff --git a/tests/WorkerFailureContextTest.php b/tests/WorkerFailureContextTest.php index 8ce9b34..8d0bf24 100644 --- a/tests/WorkerFailureContextTest.php +++ b/tests/WorkerFailureContextTest.php @@ -36,5 +36,8 @@ function assertFailureContext($expected, $actual, string $message): void assertFailureContext(false, isset($context['recordingfile']), 'unknown metric is excluded'); assertFailureContext(false, strpos(json_encode($context), '79990001122') !== false, 'phone number is absent'); assertFailureContext(false, strpos(json_encode($context), 'recordings') !== false, 'path is absent'); +assertFailureContext('beanstalk_connect', WorkerFailureContext::invokeOperation(false, true), 'construction operation'); +assertFailureContext('beanstalk_request', WorkerFailureContext::invokeOperation(true, true), 'request operation'); +assertFailureContext('beanstalk_publish', WorkerFailureContext::invokeOperation(true, false), 'publish operation'); echo "WorkerFailureContextTest: OK\n";