, dispatch_config: string}> */ public const TARGETS = [ 'forum-moderation' => [ 'expected' => ['cPad\\Plugins\\Forum\\Jobs\\AnalyzeForumPostJob'], 'dispatch_config' => 'skinbase_ai_moderation.queue.dispatch_enabled', ], 'forum-security' => [ 'expected' => [ 'cPad\\Plugins\\Forum\\Jobs\\FirewallActivityMonitor', 'cPad\\Plugins\\Forum\\Jobs\\BotActivityMonitor', ], 'dispatch_config' => 'forum_security.queues.dispatch_enabled', ], 'collections' => [ 'expected' => [ RefreshCollectionHealthJob::class, RefreshCollectionRecommendationJob::class, ScanCollectionDuplicateCandidatesJob::class, ], 'dispatch_config' => 'collections.v5.queue.dispatch_enabled', ], ]; /** * @return array */ public function inspect(string $queue, int $sampleSize = 400): array { $this->assertQueueName($queue); $redis = Redis::connection(); $waitingKey = 'queues:'.$queue; $llen = (int) $redis->llen($waitingKey); $delayed = (int) $redis->zcard('queues:'.$queue.':delayed'); $reserved = (int) $redis->zcard('queues:'.$queue.':reserved'); $notify = (int) $redis->llen('queues:'.$queue.':notify'); $prefix = (string) config('database.redis.options.prefix', ''); $classes = []; $bytes = []; $sampled = 0; if ($llen > 0) { $step = max(1, (int) floor($llen / max(1, $sampleSize))); for ($i = 0; $i < $llen && $sampled < $sampleSize; $i += $step) { $raw = $redis->lindex($waitingKey, $i); if (! is_string($raw) || $raw === '') { continue; } $sampled++; $class = QueuedJobClassMatcher::className($raw) ?? 'unknown'; $classes[$class] = ($classes[$class] ?? 0) + 1; $bytes[] = strlen($raw); } } arsort($classes); $avg = $bytes !== [] ? (int) round(array_sum($bytes) / count($bytes)) : 0; return [ 'queue' => $queue, 'logical_key' => $waitingKey, 'notify_key' => 'queues:'.$queue.':notify', 'prefix' => $prefix, 'llen' => $llen, 'delayed' => $delayed, 'reserved' => $reserved, 'notify' => $notify, 'sampled' => $sampled, 'classes' => $classes, 'avg_bytes' => $avg, 'est_bytes' => $llen * $avg, ]; } /** * @return array */ public function plan(string $queue, bool $force = false): array { if (in_array($queue, self::PROTECTED_QUEUES, true)) { throw new RuntimeException("Refusing protected queue [{$queue}]."); } if (! isset(self::TARGETS[$queue])) { throw new RuntimeException("Unknown orphan target [{$queue}]."); } $meta = self::TARGETS[$queue]; $dispatchEnabled = (bool) config($meta['dispatch_config'], false); $inspect = $this->inspect($queue); $unexpected = array_values(array_diff(array_keys($inspect['classes']), $meta['expected'])); $errors = []; if ($dispatchEnabled) { $errors[] = 'producer_dispatch_enabled'; } if ($inspect['reserved'] > 0) { $errors[] = 'reserved_jobs_present'; } if ($inspect['delayed'] > 0) { $errors[] = 'delayed_jobs_present'; } if ($unexpected !== [] && ! $force) { $errors[] = 'unexpected_job_classes'; } return $inspect + [ 'expected_classes' => $meta['expected'], 'unexpected_classes' => $unexpected, 'dispatch_enabled' => $dispatchEnabled, 'errors' => $errors, 'can_execute' => $errors === [], ]; } /** * @return array */ public function execute(string $queue, bool $force = false): array { $plan = $this->plan($queue, $force); if (! $plan['can_execute']) { throw new RuntimeException('Cleanup refused: '.implode(',', $plan['errors'])); } $redis = Redis::connection(); $keys = ['queues:'.$queue, 'queues:'.$queue.':notify']; $unlinked = 0; foreach ($keys as $key) { $unlinked += (int) $redis->unlink($key); } return $plan + ['unlinked_keys' => $keys, 'unlinked' => $unlinked]; } private function assertQueueName(string $queue): void { if (! preg_match('/^[A-Za-z0-9_-]+$/', $queue)) { throw new RuntimeException('Invalid queue name.'); } } }