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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,6 +102,8 @@ Every setting can be configured via PHP constant (in `wp-config.php`) or environ
| `QUEUE_WORKER_JOB_TIMEOUT` | `300` | Per-job timeout in seconds (SIGALRM) |
| `QUEUE_WORKER_BATCH_TIMEOUT` | `3600` | Subprocess timeout in seconds (safety net) |
| `QUEUE_WORKER_RESCAN_INTERVAL` | `60` | Seconds between database rescans |
| `QUEUE_WORKER_SCHEDULING_HORIZON` | `3600` | Only keep timers for jobs due within this many seconds; never shorter than the rescan interval |
| `QUEUE_WORKER_SCAN_TIMEOUT` | `300` | Full-network scanner subprocess timeout in seconds |
| `QUEUE_WORKER_MEMORY_LIMIT` | `200` | Memory limit in MB before auto-restart |
| `QUEUE_WORKER_UPTIME_LIMIT` | `3600` | Max uptime in seconds before auto-restart |
| `QUEUE_WORKER_LOG_FILE` | auto-detect | Path to log for admin viewer |
Expand Down
285 changes: 279 additions & 6 deletions bin/scan-cron.php
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@
}

use QueueWorker\Bootstrap;
use QueueWorker\Config;
use QueueWorker\Cron_Event_Filter;
use QueueWorker\Job_Payload;

Expand All @@ -90,9 +91,13 @@

require_once $wp_load;

$scheduling_horizon = max(1, (int) ($payload['scheduling_horizon'] ?? Config::scheduling_horizon()));
$scan_timeout = max(1, (int) ($payload['scan_timeout'] ?? Config::scan_timeout()));
$payloads = [];
$isolated_network_id = (int) ($payload['isolated_network_id'] ?? 0);
if ($isolated_network_id > 0) {
if (!empty($payload['full_network'])) {
$payloads = qw_scan_full_network_jobs($scheduling_horizon, $scan_timeout);
} elseif ($isolated_network_id > 0) {
if (!defined('WU_MT_LEGACY_ISOLATED_NETWORK')
|| (int) WU_MT_LEGACY_ISOLATED_NETWORK !== $isolated_network_id
) {
Expand All @@ -113,7 +118,7 @@
switch_to_blog($local_site_id);
}

$payloads = array_merge($payloads, qw_scan_current_site_jobs());
$payloads = array_merge($payloads, qw_scan_current_site_jobs($scheduling_horizon));

if ($switched) {
restore_current_blog();
Expand All @@ -131,7 +136,7 @@
switch_to_blog($site_id);
}

$payloads = qw_scan_current_site_jobs();
$payloads = qw_scan_current_site_jobs($scheduling_horizon);
}

// A plugin can leave nested output buffers open. Discard every buffer opened
Expand All @@ -158,16 +163,96 @@

fwrite(STDOUT, $encoded_payloads);

function qw_scan_current_site_jobs(): array
function qw_scan_full_network_jobs(int $scheduling_horizon, int $scan_timeout): array
{
$payloads = [];
$deadline = time() + $scan_timeout;
$sovereign_sites = qw_scan_sovereign_site_entries();
$isolated_networks = qw_scan_isolated_network_entries();
$initial_blog_id = get_current_blog_id();
$site_ids = is_multisite() ? get_sites(['number' => 0, 'fields' => 'ids']) : [$initial_blog_id];

foreach ($site_ids as $site_id) {
$site_id = (int) $site_id;
if (isset($sovereign_sites[$site_id])) {
continue;
}

$site = get_site($site_id);
$network_id = $site ? (int) $site->site_id : 0;
if (isset($isolated_networks[$network_id])) {
continue;
}

$switched = $site_id !== get_current_blog_id();
if ($switched) {
switch_to_blog($site_id);
}
foreach (qw_scan_current_site_jobs($scheduling_horizon) as $payload) {
$payloads[] = $payload;
}
if ($switched) {
restore_current_blog();
}
}

if ($initial_blog_id !== get_current_blog_id()) {
switch_to_blog($initial_blog_id);
}

foreach ($sovereign_sites as $site_id => $entry) {
$remaining = $deadline - time();
if ($remaining <= 0) {
fwrite(STDERR, sprintf("Scan budget exhausted before sovereign site %d.\n", $site_id));
return $payloads;
}

$site_payloads = qw_scan_registry_jobs([
'site_id' => (int) $site_id,
'site_url' => qw_scan_site_url_from_registry_entry($entry),
'scheduling_horizon' => $scheduling_horizon,
'scan_timeout' => $remaining,
], 'sovereign site ' . $site_id, $remaining);
foreach ($site_payloads as $payload) {
$payloads[] = $payload;
}
}

foreach ($isolated_networks as $network_id => $entry) {
$remaining = $deadline - time();
if ($remaining <= 0) {
fwrite(STDERR, sprintf("Scan budget exhausted before isolated network %d.\n", $network_id));
return $payloads;
}

$network_payloads = qw_scan_registry_jobs([
'isolated_network_id' => (int) $network_id,
'site_url' => qw_scan_site_url_from_registry_entry($entry),
'scheduling_horizon' => $scheduling_horizon,
'scan_timeout' => $remaining,
], 'isolated network ' . $network_id, $remaining);
foreach ($network_payloads as $payload) {
$payloads[] = $payload;
}
}
Comment thread
coderabbitai[bot] marked this conversation as resolved.

return $payloads;
}

function qw_scan_current_site_jobs(int $scheduling_horizon): array
{
wp_cache_delete('cron', 'options');
wp_cache_delete('alloptions', 'options');

$payloads = [];
$deadline = time() + $scheduling_horizon;
$crons = _get_cron_array();
if (is_array($crons)) {
$seen_cron_signatures = [];
foreach ($crons as $timestamp => $hooks) {
if ((int) $timestamp > $deadline) {
continue;
}
if (!is_array($hooks)) {
continue;
}
Expand Down Expand Up @@ -200,9 +285,11 @@ function qw_scan_current_site_jobs(): array
]);
foreach ($actions as $action_id => $action) {
$action_payload = Job_Payload::from_as_action($action_id);
if ($action_payload) {
$payloads[] = json_decode($action_payload->to_json(), true);
if (!$action_payload || $action_payload->timestamp > $deadline) {
continue;
}

$payloads[] = json_decode($action_payload->to_json(), true);
}
} catch (\Throwable $e) {
fwrite(STDERR, "Action Scheduler scan failed: " . $e->getMessage() . "\n");
Expand All @@ -212,6 +299,192 @@ function qw_scan_current_site_jobs(): array
return $payloads;
}

function qw_scan_sovereign_site_entries(): array
{
if (!defined('WP_CONTENT_DIR')) {
return [];
}

$data = qw_scan_registry_data(WP_CONTENT_DIR . '/site-registry.data.json');
if (empty($data['sites']) || !is_array($data['sites'])) {
return [];
}

$entries = [];
foreach ($data['sites'] as $site_id => $entry) {
if (!is_array($entry)
|| ($entry['isolation_model'] ?? '') !== 'sovereign'
|| ($entry['status'] ?? 'active') !== 'active'
|| qw_scan_site_url_from_registry_entry($entry) === ''
) {
continue;
}
$entries[(int) $site_id] = $entry;
}

return $entries;
}

function qw_scan_isolated_network_entries(): array
{
if (!defined('WP_CONTENT_DIR')) {
return [];
}

$data = qw_scan_registry_data(WP_CONTENT_DIR . '/network-registry.data.json');
if (empty($data['networks']) || !is_array($data['networks'])) {
return [];
}

$entries = [];
foreach ($data['networks'] as $registry_id => $entry) {
if (!is_array($entry)
|| ($entry['tier'] ?? '') !== 'isolated'
|| ($entry['status'] ?? 'active') !== 'active'
|| qw_scan_site_url_from_registry_entry($entry) === ''
) {
continue;
}

$network_id = (int) ($entry['network_id'] ?? $entry['id'] ?? $registry_id);
if ($network_id > 0) {
$entries[$network_id] = $entry;
}
}

return $entries;
}

function qw_scan_registry_data(string $path): array
{
if (!is_readable($path)) {
return [];
}

$data = json_decode((string) file_get_contents($path), true);
return is_array($data) ? $data : [];
}

function qw_scan_site_url_from_registry_entry(array $entry): string
{
$domains = $entry['domains'] ?? [];
if (!is_array($domains)) {
$domains = [];
}
if (!empty($entry['domain'])) {
array_unshift($domains, $entry['domain']);
}

foreach ($domains as $domain) {
$domain = trim((string) $domain);
if ($domain !== '') {
return 'https://' . $domain . '/';
}
}

return '';
}

function qw_scan_registry_jobs(array $payload, string $description, int $scan_timeout): array
{
$process = proc_open([PHP_BINARY, __FILE__, '--stdin'], [
0 => ['pipe', 'r'],
1 => ['pipe', 'w'],
2 => ['pipe', 'w'],
], $pipes);
if (!is_resource($process)) {
fwrite(STDERR, sprintf("Could not start scanner for %s.\n", $description));
return [];
}

$encoded_payload = json_encode($payload);
if ($encoded_payload === false || fwrite($pipes[0], $encoded_payload) === false) {
fclose($pipes[0]);
fclose($pipes[1]);
fclose($pipes[2]);
proc_terminate($process, 9);
proc_close($process);
fwrite(STDERR, sprintf("Could not configure scanner for %s.\n", $description));
return [];
}

fclose($pipes[0]);
stream_set_blocking($pipes[1], false);
stream_set_blocking($pipes[2], false);
$stdout = '';
$stderr = '';
$started = time();
do {
$out = stream_get_contents($pipes[1]);
if ($out !== false && $out !== '') {
$stdout .= $out;
}
$err = stream_get_contents($pipes[2]);
if ($err !== false && $err !== '') {
$stderr .= $err;
}

$status = proc_get_status($process);
if (!$status['running']) {
break;
}
if (time() - $started >= $scan_timeout) {
proc_terminate($process, 9);
fwrite(STDERR, sprintf("Scanner for %s exceeded %d seconds.\n", $description, $scan_timeout));
break;
}
usleep(10000);
} while (true);

$remaining = stream_get_contents($pipes[1]);
if ($remaining !== false && $remaining !== '') {
$stdout .= $remaining;
}
$remaining_error = stream_get_contents($pipes[2]);
if ($remaining_error !== false && $remaining_error !== '') {
$stderr .= $remaining_error;
}
fclose($pipes[1]);
fclose($pipes[2]);
$close_code = proc_close($process);
$exit_code = (int) ($status['exitcode'] ?? $close_code);
if ($exit_code === -1) {
$exit_code = $close_code;
}

if ($exit_code !== 0 || $stdout === '') {
fwrite(STDERR, sprintf(
"Scanner for %s failed (exit %d, stdout_bytes=%d, stderr_bytes=%d).\n",
$description,
$exit_code,
strlen($stdout),
strlen($stderr)
));
return [];
}

try {
$jobs = json_decode($stdout, true, 512, JSON_THROW_ON_ERROR);
} catch (\JsonException $e) {
fwrite(STDERR, sprintf("Scanner for %s returned invalid JSON.\n", $description));
return [];
}

if (!array_is_list($jobs)) {
fwrite(STDERR, sprintf("Scanner for %s returned an invalid payload shape.\n", $description));
return [];
}

foreach ($jobs as $job) {
if (!is_array($job)) {
fwrite(STDERR, sprintf("Scanner for %s returned an invalid payload shape.\n", $description));
return [];
}
}

return $jobs;
}

function qw_scan_action_scheduler_tables_exist(): bool
{
global $wpdb;
Expand Down
1 change: 1 addition & 0 deletions bin/worker.php
Original file line number Diff line number Diff line change
Expand Up @@ -91,6 +91,7 @@
$process = new Worker_Process($wp_load, $primary_domain, $execute_script, $scan_script);

$worker->onWorkerStart = [$process, 'on_worker_start'];
$worker->onWorkerStop = [$process, 'on_worker_stop'];
$worker->onMessage = [$process, 'on_message'];

Worker::runAll();
10 changes: 10 additions & 0 deletions src/class-config.php
Original file line number Diff line number Diff line change
Expand Up @@ -150,6 +150,16 @@ public static function rescan_interval(): int
return (int) self::get('QUEUE_WORKER_RESCAN_INTERVAL', 60);
}

public static function scheduling_horizon(): int
{
return max(1, self::rescan_interval(), (int) self::get('QUEUE_WORKER_SCHEDULING_HORIZON', 3600));
}

public static function scan_timeout(): int
{
return max(1, (int) self::get('QUEUE_WORKER_SCAN_TIMEOUT', 300));
}

public static function action_scheduler_rescan_interval(): int
{
return max(1, (int) self::get('QUEUE_WORKER_AS_RESCAN_INTERVAL', 5));
Expand Down
Loading
Loading