Keep similar-ai from tripping the global circuit on a lone URL 502, clamp Qdrant search to 100, and add Server-Timing plus slow-request logging. Studio shared props, Academy S3 exists caching, heat chunking, and Redis/scheduler hygiene stay in this rollout.
160 lines
5.3 KiB
PHP
160 lines
5.3 KiB
PHP
<?php
|
|
|
|
declare(strict_types=1);
|
|
|
|
namespace App\Support\Redis;
|
|
|
|
use App\Jobs\RefreshCollectionHealthJob;
|
|
use App\Jobs\RefreshCollectionRecommendationJob;
|
|
use App\Jobs\ScanCollectionDuplicateCandidatesJob;
|
|
use App\Support\Queues\QueuedJobClassMatcher;
|
|
use Illuminate\Support\Facades\Redis;
|
|
use RuntimeException;
|
|
|
|
final class OrphanQueueCleanup
|
|
{
|
|
public const PROTECTED_QUEUES = ['default', 'search', 'mail', 'broadcasts', 'notifications'];
|
|
|
|
/**
|
|
* @var array<string, array{expected: list<string>, 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<string, mixed>
|
|
*/
|
|
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<string, mixed>
|
|
*/
|
|
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<string, mixed>
|
|
*/
|
|
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.');
|
|
}
|
|
}
|
|
}
|