Ship production optimization M1-M12.5A: queues, metrics, HTTP observability, and vector search reliability.
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.
This commit is contained in:
@@ -18,6 +18,12 @@ class DispatchCollectionMaintenanceCommand extends Command
|
||||
|
||||
public function handle(CollectionBackgroundJobService $jobs): int
|
||||
{
|
||||
if (! (bool) config('collections.v5.queue.dispatch_enabled', false)) {
|
||||
$this->warn('Collection maintenance dispatch is disabled (COLLECTIONS_V5_DISPATCH_ENABLED). No jobs queued.');
|
||||
|
||||
return self::SUCCESS;
|
||||
}
|
||||
|
||||
$runHealth = (bool) $this->option('health');
|
||||
$runRecommendations = (bool) $this->option('recommendations');
|
||||
$runDuplicates = (bool) $this->option('duplicates');
|
||||
|
||||
@@ -48,13 +48,32 @@ final class GenerateSitemapsCommand extends Command
|
||||
$this->newLine();
|
||||
|
||||
// ── Root sitemap index ────────────────────────────────────────────
|
||||
// Write several paths so nginx `location = /sitemap.xml` and child
|
||||
// listings stay in sync. `sitemaps/index.xml` is the canonical generated
|
||||
// index: it is a new filename, so a scheduler user can create it even
|
||||
// when a stale `sitemaps/sitemap.xml` is owned by another account.
|
||||
$t = microtime(true);
|
||||
$index = $build->buildIndex(force: true, persist: false, families: $families);
|
||||
$disk->put('sitemaps/sitemap.xml', $index['content']);
|
||||
$written++;
|
||||
$indexPaths = ['sitemaps/index.xml', 'sitemaps/sitemap.xml', 'sitemap.xml'];
|
||||
$indexWritten = 0;
|
||||
|
||||
foreach ($indexPaths as $path) {
|
||||
if ($this->writeXml($disk, $path, $index['content'])) {
|
||||
$written++;
|
||||
$indexWritten++;
|
||||
}
|
||||
}
|
||||
|
||||
if ($indexWritten === 0) {
|
||||
$this->error('Failed to write any root sitemap index file. Check ownership of public/sitemap.xml and public/sitemaps/.');
|
||||
|
||||
return self::FAILURE;
|
||||
}
|
||||
|
||||
$this->line(sprintf(
|
||||
' <info>✔</info> sitemaps/sitemap.xml %d entries <comment>%.3fs</comment>',
|
||||
' <info>✔</info> sitemap index %d entries %d path(s) <comment>%.3fs</comment>',
|
||||
$index['url_count'],
|
||||
$indexWritten,
|
||||
microtime(true) - $t,
|
||||
));
|
||||
|
||||
@@ -78,7 +97,12 @@ final class GenerateSitemapsCommand extends Command
|
||||
}
|
||||
|
||||
$path = 'sitemaps/' . $documentName . '.xml';
|
||||
$disk->put($path, $built['content']);
|
||||
if (! $this->writeXml($disk, $path, $built['content'])) {
|
||||
$this->line(sprintf(' <fg=red>✖</> %s write failed', $documentName . '.xml'));
|
||||
$failed++;
|
||||
|
||||
continue;
|
||||
}
|
||||
$written++;
|
||||
|
||||
$this->line(sprintf(
|
||||
@@ -114,7 +138,12 @@ final class GenerateSitemapsCommand extends Command
|
||||
}
|
||||
|
||||
$path = 'sitemaps/' . $groupName . '.xml';
|
||||
$disk->put($path, $built['content']);
|
||||
if (! $this->writeXml($disk, $path, $built['content'])) {
|
||||
$this->line(sprintf(' <fg=red>✖</> %s.xml write failed', $groupName));
|
||||
$failed++;
|
||||
|
||||
continue;
|
||||
}
|
||||
$written++;
|
||||
|
||||
$this->line(sprintf(
|
||||
@@ -159,4 +188,17 @@ final class GenerateSitemapsCommand extends Command
|
||||
|
||||
return array_values(array_filter($enabled, fn (string $f): bool => in_array($f, $only, true)));
|
||||
}
|
||||
|
||||
private function writeXml(\Illuminate\Contracts\Filesystem\Filesystem $disk, string $path, string $content): bool
|
||||
{
|
||||
$ok = $disk->put($path, $content);
|
||||
|
||||
if ($ok !== true) {
|
||||
$this->warn(' Could not write '.$path);
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -7,34 +7,137 @@ use Illuminate\Support\Facades\DB;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
|
||||
/**
|
||||
* Prune old hourly metric snapshots to prevent unbounded table growth.
|
||||
* Prune old hourly metric snapshots in bounded batches.
|
||||
*
|
||||
* Usage: php artisan nova:prune-metric-snapshots
|
||||
* php artisan nova:prune-metric-snapshots --keep-days=7
|
||||
* Usage:
|
||||
* php artisan nova:prune-metric-snapshots
|
||||
* php artisan nova:prune-metric-snapshots --keep-days=30 --chunk=5000 --dry-run
|
||||
*/
|
||||
class PruneMetricSnapshotsCommand extends Command
|
||||
{
|
||||
protected $signature = 'nova:prune-metric-snapshots
|
||||
{--keep-days=7 : Keep snapshots for this many days}';
|
||||
{--keep-days= : Keep snapshots for this many days (default: config metrics.hourly_snapshot_retention_days)}
|
||||
{--chunk= : Rows to delete per batch (default: config metrics.hourly_snapshot_prune_chunk)}
|
||||
{--sleep-ms= : Pause between batches in milliseconds}
|
||||
{--max-batches=0 : Stop after N batches (0 = until done)}
|
||||
{--dry-run : Count rows that would be deleted without deleting}';
|
||||
|
||||
protected $description = 'Delete old hourly metric snapshots beyond the retention window';
|
||||
protected $description = 'Delete old hourly metric snapshots beyond the retention window (batched)';
|
||||
|
||||
public function handle(): int
|
||||
{
|
||||
$keepDays = (int) $this->option('keep-days');
|
||||
$cutoff = now()->subDays($keepDays);
|
||||
$keepDays = $this->resolveKeepDays();
|
||||
$chunk = $this->resolvePositiveInt('chunk', (int) config('metrics.hourly_snapshot_prune_chunk', 5000), 1);
|
||||
$sleepMs = $this->resolveNonNegativeInt('sleep-ms', (int) config('metrics.hourly_snapshot_prune_sleep_ms', 50));
|
||||
$maxBatches = max(0, (int) $this->option('max-batches'));
|
||||
$dryRun = (bool) $this->option('dry-run');
|
||||
$cutoff = now()->subDays($keepDays);
|
||||
|
||||
$deleted = DB::table('artwork_metric_snapshots_hourly')
|
||||
if ($keepDays < 1) {
|
||||
$this->error('keep-days must be >= 1.');
|
||||
|
||||
return self::FAILURE;
|
||||
}
|
||||
|
||||
$eligible = (int) DB::table('artwork_metric_snapshots_hourly')
|
||||
->where('bucket_hour', '<', $cutoff)
|
||||
->delete();
|
||||
->count();
|
||||
|
||||
$this->info("Pruned {$deleted} snapshot rows older than {$keepDays} days.");
|
||||
$this->info(sprintf(
|
||||
'[nova:prune-metric-snapshots] cutoff=%s keep_days=%d eligible=%d chunk=%d sleep_ms=%d max_batches=%s%s',
|
||||
$cutoff->toDateTimeString(),
|
||||
$keepDays,
|
||||
$eligible,
|
||||
$chunk,
|
||||
$sleepMs,
|
||||
$maxBatches === 0 ? 'unlimited' : (string) $maxBatches,
|
||||
$dryRun ? ' (dry-run)' : ''
|
||||
));
|
||||
|
||||
if ($dryRun) {
|
||||
Log::info('[nova:prune-metric-snapshots] dry-run', [
|
||||
'eligible' => $eligible,
|
||||
'keep_days' => $keepDays,
|
||||
'cutoff' => $cutoff->toDateTimeString(),
|
||||
]);
|
||||
|
||||
return self::SUCCESS;
|
||||
}
|
||||
|
||||
$deleted = 0;
|
||||
$batches = 0;
|
||||
|
||||
while (true) {
|
||||
if ($maxBatches > 0 && $batches >= $maxBatches) {
|
||||
$this->warn("Stopped after max-batches={$maxBatches}.");
|
||||
break;
|
||||
}
|
||||
|
||||
$ids = DB::table('artwork_metric_snapshots_hourly')
|
||||
->where('bucket_hour', '<', $cutoff)
|
||||
->orderBy('id')
|
||||
->limit($chunk)
|
||||
->pluck('id');
|
||||
|
||||
if ($ids->isEmpty()) {
|
||||
break;
|
||||
}
|
||||
|
||||
$batchDeleted = DB::table('artwork_metric_snapshots_hourly')
|
||||
->whereIn('id', $ids->all())
|
||||
->delete();
|
||||
|
||||
$deleted += $batchDeleted;
|
||||
$batches++;
|
||||
|
||||
Log::info('[nova:prune-metric-snapshots] batch', [
|
||||
'batch' => $batches,
|
||||
'deleted' => $batchDeleted,
|
||||
'deleted_total' => $deleted,
|
||||
]);
|
||||
|
||||
if ($sleepMs > 0) {
|
||||
usleep($sleepMs * 1000);
|
||||
}
|
||||
}
|
||||
|
||||
$this->info("Pruned {$deleted} snapshot rows older than {$keepDays} days in {$batches} batch(es).");
|
||||
|
||||
Log::info('[nova:prune-metric-snapshots] completed', [
|
||||
'deleted' => $deleted,
|
||||
'deleted' => $deleted,
|
||||
'batches' => $batches,
|
||||
'keep_days' => $keepDays,
|
||||
'cutoff' => $cutoff->toDateTimeString(),
|
||||
'eligible_at_start' => $eligible,
|
||||
]);
|
||||
|
||||
return self::SUCCESS;
|
||||
}
|
||||
|
||||
private function resolveKeepDays(): int
|
||||
{
|
||||
$option = $this->option('keep-days');
|
||||
|
||||
if ($option === null || $option === '') {
|
||||
return (int) config('metrics.hourly_snapshot_retention_days', 30);
|
||||
}
|
||||
|
||||
return (int) $option;
|
||||
}
|
||||
|
||||
private function resolvePositiveInt(string $option, int $default, int $minimum): int
|
||||
{
|
||||
$raw = $this->option($option);
|
||||
$value = ($raw === null || $raw === '') ? $default : (int) $raw;
|
||||
|
||||
return max($minimum, $value);
|
||||
}
|
||||
|
||||
private function resolveNonNegativeInt(string $option, int $default): int
|
||||
{
|
||||
$raw = $this->option($option);
|
||||
$value = ($raw === null || $raw === '') ? $default : (int) $raw;
|
||||
|
||||
return max(0, $value);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,88 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Console\Commands;
|
||||
|
||||
use App\Services\Traffic\OnlineVisitorRepository;
|
||||
use App\Services\Traffic\PresenceIndexPruner;
|
||||
use Illuminate\Console\Command;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
use Illuminate\Support\Facades\Redis;
|
||||
|
||||
class PruneOnlineVisitorIndexCommand extends Command
|
||||
{
|
||||
protected $signature = 'skinbase:prune-online-visitor-index
|
||||
{--dry-run : Count stale members without SREM}
|
||||
{--execute : Remove stale members}
|
||||
{--scan-count= : SSCAN COUNT}
|
||||
{--max-batches= : Batches this run}
|
||||
{--max-members= : Members to inspect this run}
|
||||
{--sleep-ms= : Pause between batches}
|
||||
{--time-limit= : Stop after N seconds (0 = none)}';
|
||||
|
||||
protected $description = 'Remove stale members from the online-visitor Redis set (SSCAN, never SMEMBERS)';
|
||||
|
||||
public function handle(PresenceIndexPruner $pruner): int
|
||||
{
|
||||
$dryRun = (bool) $this->option('dry-run');
|
||||
$execute = (bool) $this->option('execute');
|
||||
if ($dryRun === $execute) {
|
||||
$this->error('Pass exactly one of --dry-run or --execute.');
|
||||
|
||||
return self::FAILURE;
|
||||
}
|
||||
|
||||
$scanCount = $this->intOption('scan-count', (int) config('traffic.online_visitors.prune_scan_count', 500));
|
||||
$maxBatches = $this->intOption('max-batches', (int) config('traffic.online_visitors.prune_max_batches', 200));
|
||||
$maxMembers = $this->intOption('max-members', (int) config('traffic.online_visitors.prune_max_members', 100000));
|
||||
$sleepMs = $this->intOption('sleep-ms', (int) config('traffic.online_visitors.prune_sleep_ms', 25));
|
||||
$timeLimit = $this->intOption('time-limit', 0);
|
||||
|
||||
$scard = (int) Redis::scard(OnlineVisitorRepository::INDEX_KEY);
|
||||
|
||||
$this->info(sprintf(
|
||||
'[prune-online-visitor-index] key=%s prefix=%s scard=%d scan_count=%d max_batches=%d max_members=%d sleep_ms=%d %s',
|
||||
OnlineVisitorRepository::INDEX_KEY,
|
||||
(string) config('database.redis.options.prefix'),
|
||||
$scard,
|
||||
$scanCount,
|
||||
$maxBatches,
|
||||
$maxMembers,
|
||||
$sleepMs,
|
||||
$dryRun ? 'dry-run' : 'execute'
|
||||
));
|
||||
|
||||
$result = $pruner->prune(
|
||||
$scanCount,
|
||||
$maxBatches,
|
||||
$maxMembers,
|
||||
$sleepMs,
|
||||
$dryRun,
|
||||
$timeLimit > 0 ? $timeLimit : null,
|
||||
);
|
||||
|
||||
$this->info(sprintf(
|
||||
'scanned=%d stale=%d live=%d removed=%d batches=%d',
|
||||
$result['scanned'],
|
||||
$result['stale'],
|
||||
$result['live'],
|
||||
$result['removed'],
|
||||
$result['batches']
|
||||
));
|
||||
|
||||
Log::info('[prune-online-visitor-index] completed', $result + ['scard_before' => $scard]);
|
||||
|
||||
return self::SUCCESS;
|
||||
}
|
||||
|
||||
private function intOption(string $name, int $default): int
|
||||
{
|
||||
$raw = $this->option($name);
|
||||
if ($raw === null || $raw === '') {
|
||||
return $default;
|
||||
}
|
||||
|
||||
return max(0, (int) $raw);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,94 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Console\Commands;
|
||||
|
||||
use App\Jobs\RecComputeSimilarByTagsJob;
|
||||
use App\Support\Queues\QueuedJobClassMatcher;
|
||||
use Illuminate\Console\Command;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
use Illuminate\Support\Facades\Redis;
|
||||
|
||||
/**
|
||||
* Deployment cleanup: remove waiting RecComputeSimilarByTagsJob payloads only.
|
||||
*
|
||||
* Does not run from the scheduler. Operator must pass --dry-run or --execute.
|
||||
*
|
||||
* Usage:
|
||||
* php artisan skinbase:purge-queued-rec-tags --dry-run
|
||||
* php artisan skinbase:purge-queued-rec-tags --execute
|
||||
*/
|
||||
class PurgeQueuedRecComputeTagsCommand extends Command
|
||||
{
|
||||
protected $signature = 'skinbase:purge-queued-rec-tags
|
||||
{--queue=default : Redis queue name (waiting list only)}
|
||||
{--dry-run : Count matches without removing}
|
||||
{--execute : Remove matched waiting payloads}';
|
||||
|
||||
protected $description = 'Remove waiting RecComputeSimilarByTagsJob payloads from a Redis queue (opt-in)';
|
||||
|
||||
public function handle(): int
|
||||
{
|
||||
$dryRun = (bool) $this->option('dry-run');
|
||||
$execute = (bool) $this->option('execute');
|
||||
|
||||
if ($dryRun === $execute) {
|
||||
$this->error('Pass exactly one of --dry-run or --execute.');
|
||||
|
||||
return self::FAILURE;
|
||||
}
|
||||
|
||||
$queue = (string) $this->option('queue');
|
||||
if (! preg_match('/^[A-Za-z0-9_-]+$/', $queue)) {
|
||||
$this->error('Invalid queue name.');
|
||||
|
||||
return self::FAILURE;
|
||||
}
|
||||
|
||||
$key = 'queues:'.$queue;
|
||||
$target = RecComputeSimilarByTagsJob::class;
|
||||
$redis = Redis::connection();
|
||||
$len = (int) $redis->llen($key);
|
||||
$matchedPayloads = [];
|
||||
$chunk = 200;
|
||||
|
||||
$this->info(sprintf(
|
||||
'[purge-queued-rec-tags] queue=%s waiting=%d target=%s %s',
|
||||
$queue,
|
||||
$len,
|
||||
$target,
|
||||
$dryRun ? 'dry-run' : 'execute'
|
||||
));
|
||||
|
||||
for ($start = 0; $start < $len; $start += $chunk) {
|
||||
$rows = $redis->lrange($key, $start, min($start + $chunk - 1, $len - 1));
|
||||
foreach ($rows as $raw) {
|
||||
if (QueuedJobClassMatcher::isClass((string) $raw, $target)) {
|
||||
$matchedPayloads[] = (string) $raw;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
$matched = count($matchedPayloads);
|
||||
$removed = 0;
|
||||
|
||||
if (! $dryRun) {
|
||||
foreach ($matchedPayloads as $raw) {
|
||||
$removed += (int) $redis->lrem($key, 1, $raw);
|
||||
}
|
||||
}
|
||||
|
||||
$this->info(sprintf('matched=%d removed=%d preserved_other=%s', $matched, $removed, $dryRun ? 'yes' : 'yes'));
|
||||
|
||||
Log::info('[purge-queued-rec-tags] completed', [
|
||||
'queue' => $queue,
|
||||
'waiting' => $len,
|
||||
'matched' => $matched,
|
||||
'removed' => $removed,
|
||||
'dry_run' => $dryRun,
|
||||
]);
|
||||
|
||||
return self::SUCCESS;
|
||||
}
|
||||
}
|
||||
@@ -57,7 +57,7 @@ class RecalculateHeatCommand extends Command
|
||||
$updatedCount = 0;
|
||||
$skippedCount = 0;
|
||||
|
||||
// Process in chunks using artwork IDs that have at least one snapshot in the smoothing window
|
||||
// Distinct IDs only — do not hydrate the full 24h snapshot window in one query.
|
||||
$artworkIds = DB::table('artwork_metric_snapshots_hourly')
|
||||
->whereBetween('bucket_hour', [$lookbackStart, $currentHour])
|
||||
->distinct()
|
||||
@@ -68,27 +68,33 @@ class RecalculateHeatCommand extends Command
|
||||
return self::SUCCESS;
|
||||
}
|
||||
|
||||
// Load all snapshots for the lookback window in bulk
|
||||
$snapshots = DB::table('artwork_metric_snapshots_hourly')
|
||||
->whereBetween('bucket_hour', [$lookbackStart, $currentHour])
|
||||
->whereIn('artwork_id', $artworkIds)
|
||||
->orderBy('bucket_hour')
|
||||
->get()
|
||||
->groupBy('artwork_id');
|
||||
|
||||
// Load artwork published_at dates for age factor (use published_at, fall back to created_at)
|
||||
$artworkDates = DB::table('artworks')
|
||||
->whereIn('id', $artworkIds)
|
||||
->whereNull('deleted_at')
|
||||
->where('is_approved', true)
|
||||
->select('id', 'published_at', 'created_at')
|
||||
->get()
|
||||
->mapWithKeys(fn ($row) => [
|
||||
$row->id => \Carbon\Carbon::parse($row->published_at ?? $row->created_at),
|
||||
]);
|
||||
|
||||
// Process in chunks
|
||||
foreach ($artworkIds->chunk($chunk) as $chunkIds) {
|
||||
$snapshots = DB::table('artwork_metric_snapshots_hourly')
|
||||
->select([
|
||||
'artwork_id',
|
||||
'bucket_hour',
|
||||
'views_count',
|
||||
'downloads_count',
|
||||
'favourites_count',
|
||||
'comments_count',
|
||||
'shares_count',
|
||||
])
|
||||
->whereBetween('bucket_hour', [$lookbackStart, $currentHour])
|
||||
->whereIn('artwork_id', $chunkIds)
|
||||
->orderBy('bucket_hour')
|
||||
->get()
|
||||
->groupBy('artwork_id');
|
||||
|
||||
$artworkDates = DB::table('artworks')
|
||||
->whereIn('id', $chunkIds)
|
||||
->whereNull('deleted_at')
|
||||
->where('is_approved', true)
|
||||
->select('id', 'published_at', 'created_at')
|
||||
->get()
|
||||
->mapWithKeys(fn ($row) => [
|
||||
$row->id => \Carbon\Carbon::parse($row->published_at ?? $row->created_at),
|
||||
]);
|
||||
|
||||
$upsertRows = [];
|
||||
|
||||
foreach ($chunkIds as $artworkId) {
|
||||
|
||||
@@ -0,0 +1,139 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Console\Commands;
|
||||
|
||||
use Illuminate\Console\Command;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
use Predis\Client as PredisClient;
|
||||
|
||||
class RedisCleanupLegacyPrefixCommand extends Command
|
||||
{
|
||||
protected $signature = 'skinbase:redis-cleanup-legacy-prefix
|
||||
{--dry-run : SCAN and report only}
|
||||
{--execute : UNLINK matching obsolete keys}
|
||||
{--max-keys=500 : Cap keys this run}';
|
||||
|
||||
protected $description = 'UNLINK obsolete skinbasenova-database-* and skinbasenova_horizon:* keys only';
|
||||
|
||||
private const LEGACY_PREFIXES = [
|
||||
'skinbasenova-database-',
|
||||
'skinbasenova_horizon:',
|
||||
];
|
||||
|
||||
private const PROTECTED_PREFIXES = [
|
||||
'skinbase-database-',
|
||||
'skinbase_horizon:',
|
||||
];
|
||||
|
||||
public function handle(): int
|
||||
{
|
||||
$dryRun = (bool) $this->option('dry-run');
|
||||
$execute = (bool) $this->option('execute');
|
||||
if ($dryRun === $execute) {
|
||||
$this->error('Pass exactly one of --dry-run or --execute.');
|
||||
|
||||
return self::FAILURE;
|
||||
}
|
||||
|
||||
$maxKeys = max(1, (int) $this->option('max-keys'));
|
||||
$currentPrefix = (string) config('database.redis.options.prefix');
|
||||
$horizonPrefix = (string) config('horizon.prefix');
|
||||
|
||||
$this->info(sprintf(
|
||||
'[redis-cleanup-legacy-prefix] current_prefix=%s horizon_prefix=%s max_keys=%d %s',
|
||||
$currentPrefix,
|
||||
$horizonPrefix,
|
||||
$maxKeys,
|
||||
$dryRun ? 'dry-run' : 'execute'
|
||||
));
|
||||
|
||||
if (in_array($currentPrefix, self::LEGACY_PREFIXES, true) || in_array($horizonPrefix, self::LEGACY_PREFIXES, true)) {
|
||||
$this->error('Current process still uses a legacy prefix. Aborting.');
|
||||
|
||||
return self::FAILURE;
|
||||
}
|
||||
|
||||
$client = $this->unprefixedClient();
|
||||
$found = [];
|
||||
$cursor = '0';
|
||||
do {
|
||||
[$cursor, $keys] = $client->scan($cursor, ['COUNT' => 200, 'MATCH' => '*']);
|
||||
foreach ($keys as $key) {
|
||||
$key = (string) $key;
|
||||
if ($this->isProtected($key)) {
|
||||
continue;
|
||||
}
|
||||
if (! $this->isLegacy($key)) {
|
||||
continue;
|
||||
}
|
||||
$found[] = $key;
|
||||
if (count($found) >= $maxKeys) {
|
||||
$cursor = '0';
|
||||
break;
|
||||
}
|
||||
}
|
||||
} while ($cursor !== '0' && $cursor !== 0);
|
||||
|
||||
$this->info('matched='.count($found));
|
||||
foreach (array_slice($found, 0, 20) as $key) {
|
||||
$this->line('key='.$key);
|
||||
}
|
||||
if (count($found) > 20) {
|
||||
$this->line('... truncated listing');
|
||||
}
|
||||
|
||||
$unlinked = 0;
|
||||
if (! $dryRun && $found !== []) {
|
||||
foreach (array_chunk($found, 50) as $chunk) {
|
||||
$unlinked += (int) $client->unlink(...$chunk);
|
||||
}
|
||||
}
|
||||
|
||||
$this->info('unlinked='.$unlinked);
|
||||
|
||||
Log::info('[redis-cleanup-legacy-prefix] completed', [
|
||||
'matched' => count($found),
|
||||
'unlinked' => $unlinked,
|
||||
'dry_run' => $dryRun,
|
||||
]);
|
||||
|
||||
return self::SUCCESS;
|
||||
}
|
||||
|
||||
private function isLegacy(string $key): bool
|
||||
{
|
||||
foreach (self::LEGACY_PREFIXES as $prefix) {
|
||||
if (str_starts_with($key, $prefix)) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
private function isProtected(string $key): bool
|
||||
{
|
||||
foreach (self::PROTECTED_PREFIXES as $prefix) {
|
||||
if (str_starts_with($key, $prefix)) {
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
return false;
|
||||
}
|
||||
|
||||
private function unprefixedClient(): PredisClient
|
||||
{
|
||||
$redis = config('database.redis.default', []);
|
||||
|
||||
return new PredisClient([
|
||||
'scheme' => 'tcp',
|
||||
'host' => $redis['host'] ?? '127.0.0.1',
|
||||
'port' => (int) ($redis['port'] ?? 6379),
|
||||
'password' => $redis['password'] ?? null,
|
||||
'database' => (int) ($redis['database'] ?? 0),
|
||||
]);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,79 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Console\Commands;
|
||||
|
||||
use App\Support\Redis\OrphanQueueCleanup;
|
||||
use Illuminate\Console\Command;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
use RuntimeException;
|
||||
|
||||
class RedisCleanupOrphansCommand extends Command
|
||||
{
|
||||
protected $signature = 'skinbase:redis-cleanup-orphans
|
||||
{target : forum-moderation, forum-security, or collections}
|
||||
{--dry-run : Inspect and report without UNLINK}
|
||||
{--execute : UNLINK the waiting queue and notify list}
|
||||
{--force : Allow unexpected sampled job classes}';
|
||||
|
||||
protected $description = 'UNLINK an orphan Redis queue after producer and class checks (operator only)';
|
||||
|
||||
public function handle(OrphanQueueCleanup $cleanup): int
|
||||
{
|
||||
$dryRun = (bool) $this->option('dry-run');
|
||||
$execute = (bool) $this->option('execute');
|
||||
if ($dryRun === $execute) {
|
||||
$this->error('Pass exactly one of --dry-run or --execute.');
|
||||
|
||||
return self::FAILURE;
|
||||
}
|
||||
|
||||
$target = (string) $this->argument('target');
|
||||
$force = (bool) $this->option('force');
|
||||
|
||||
$this->info(sprintf(
|
||||
'[redis-cleanup-orphans] prefix=%s connection=default target=%s force=%s mode=%s',
|
||||
(string) config('database.redis.options.prefix'),
|
||||
$target,
|
||||
$force ? 'yes' : 'no',
|
||||
$dryRun ? 'dry-run' : 'execute'
|
||||
));
|
||||
|
||||
try {
|
||||
$plan = $execute
|
||||
? $cleanup->execute($target, $force)
|
||||
: $cleanup->plan($target, $force);
|
||||
} catch (RuntimeException $e) {
|
||||
$this->error($e->getMessage());
|
||||
Log::warning('[redis-cleanup-orphans] refused', ['target' => $target, 'error' => $e->getMessage()]);
|
||||
|
||||
return self::FAILURE;
|
||||
}
|
||||
|
||||
$this->line('logical_key='.$plan['logical_key']);
|
||||
$this->line('notify_key='.$plan['notify_key']);
|
||||
$this->line('llen='.$plan['llen'].' reserved='.$plan['reserved'].' delayed='.$plan['delayed']);
|
||||
$this->line('sampled='.$plan['sampled'].' est_bytes='.$plan['est_bytes']);
|
||||
$this->line('classes='.json_encode($plan['classes']));
|
||||
if ($plan['unexpected_classes'] !== []) {
|
||||
$this->warn('unexpected_classes='.json_encode($plan['unexpected_classes']));
|
||||
}
|
||||
if ($plan['errors'] !== []) {
|
||||
$this->error('errors='.implode(',', $plan['errors']));
|
||||
}
|
||||
if ($execute) {
|
||||
$this->info('unlinked='.($plan['unlinked'] ?? 0));
|
||||
}
|
||||
|
||||
Log::info('[redis-cleanup-orphans] completed', [
|
||||
'target' => $target,
|
||||
'dry_run' => $dryRun,
|
||||
'llen' => $plan['llen'],
|
||||
'errors' => $plan['errors'],
|
||||
'unlinked' => $plan['unlinked'] ?? 0,
|
||||
]);
|
||||
|
||||
return ($plan['errors'] === [] || $dryRun) ? self::SUCCESS : self::FAILURE;
|
||||
}
|
||||
}
|
||||
@@ -6,10 +6,12 @@ namespace App\Http\Controllers\Api;
|
||||
|
||||
use App\Http\Controllers\Controller;
|
||||
use App\Models\Artwork;
|
||||
use App\Services\Vision\VectorGatewayException;
|
||||
use App\Services\Vision\VectorService;
|
||||
use Illuminate\Http\JsonResponse;
|
||||
use Illuminate\Http\Request;
|
||||
use RuntimeException;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
use Throwable;
|
||||
|
||||
final class SimilarAiArtworksController extends Controller
|
||||
{
|
||||
@@ -35,11 +37,12 @@ final class SimilarAiArtworksController extends Controller
|
||||
|
||||
try {
|
||||
$items = $this->vectors->similarToArtwork($artwork, $limit);
|
||||
} catch (RuntimeException $e) {
|
||||
} catch (Throwable $e) {
|
||||
$this->logSimilarityFailure($artwork->id, $e);
|
||||
|
||||
return response()->json([
|
||||
'data' => [],
|
||||
'reason' => 'vector_gateway_error',
|
||||
'message' => $e->getMessage(),
|
||||
], 502);
|
||||
}
|
||||
|
||||
@@ -52,4 +55,19 @@ final class SimilarAiArtworksController extends Controller
|
||||
],
|
||||
]);
|
||||
}
|
||||
|
||||
private function logSimilarityFailure(int $artworkId, Throwable $e): void
|
||||
{
|
||||
$gateway = $e instanceof VectorGatewayException ? $e : null;
|
||||
|
||||
Log::warning('Vector similarity search failed', [
|
||||
'artwork_id' => $artworkId,
|
||||
'route' => 'api.art.similar-ai',
|
||||
'failure_stage' => $gateway?->operation ?? 'unknown',
|
||||
'exception_class' => $e::class,
|
||||
'http_status' => $gateway?->httpStatus,
|
||||
'circuit_worthy' => $gateway?->circuitWorthy ?? false,
|
||||
'circuit_open' => app(\App\Services\Vision\VectorGatewayClient::class)->circuitOpen(),
|
||||
]);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -58,6 +58,16 @@ final class ArtworkDownloadController extends Controller
|
||||
abort(404);
|
||||
}
|
||||
|
||||
if (! File::isFile($filePath)) {
|
||||
Log::warning('Artwork original file missing for download.', [
|
||||
'artwork_id' => $artwork->id,
|
||||
'ext' => $ext,
|
||||
'resolved_path' => $filePath,
|
||||
]);
|
||||
|
||||
abort(404);
|
||||
}
|
||||
|
||||
$this->recordDownload($request, $artwork->id);
|
||||
$this->incrementDownloadCountIfAvailable($artwork->id);
|
||||
|
||||
@@ -70,16 +80,6 @@ final class ArtworkDownloadController extends Controller
|
||||
]);
|
||||
}
|
||||
|
||||
if (! File::isFile($filePath)) {
|
||||
Log::warning('Artwork original file missing for download.', [
|
||||
'artwork_id' => $artwork->id,
|
||||
'ext' => $ext,
|
||||
'resolved_path' => $filePath,
|
||||
]);
|
||||
|
||||
abort(404);
|
||||
}
|
||||
|
||||
$downloadName = $this->buildDownloadFilename((string) $artwork->file_name, $ext);
|
||||
|
||||
// X-Accel-Redirect is safe only when nginx is explicitly configured to
|
||||
|
||||
@@ -22,6 +22,18 @@ final class SitemapController extends Controller
|
||||
|
||||
public function index(): Response|BinaryFileResponse
|
||||
{
|
||||
if ((bool) config('sitemaps.pre_generated.enabled', true)) {
|
||||
foreach ([
|
||||
public_path('sitemaps/index.xml'),
|
||||
public_path('sitemaps/sitemap.xml'),
|
||||
public_path('sitemap.xml'),
|
||||
] as $path) {
|
||||
if (is_file($path) && is_readable($path)) {
|
||||
return $this->xmlFileResponse($path);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// 1. Published release (release management pipeline fallback).
|
||||
$published = $this->published->resolveIndex();
|
||||
if ($published !== null) {
|
||||
|
||||
@@ -0,0 +1,43 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Http\Middleware;
|
||||
|
||||
use Closure;
|
||||
use Illuminate\Http\Request;
|
||||
use Symfony\Component\HttpFoundation\Response;
|
||||
|
||||
/**
|
||||
* Negligible-overhead request duration header for production HTTP profiling.
|
||||
* Does not log queries or payloads.
|
||||
*/
|
||||
final class AddServerTiming
|
||||
{
|
||||
public function handle(Request $request, Closure $next): Response
|
||||
{
|
||||
$started = hrtime(true);
|
||||
$response = $next($request);
|
||||
$durationMs = (hrtime(true) - $started) / 1_000_000;
|
||||
|
||||
if (headers_sent()) {
|
||||
return $response;
|
||||
}
|
||||
|
||||
$metrics = ['app;desc="Laravel";dur='.number_format($durationMs, 1, '.', '')];
|
||||
$ssrMs = $request->attributes->get('http.ssr_ms');
|
||||
if (is_numeric($ssrMs) && (float) $ssrMs >= 0) {
|
||||
$metrics[] = 'ssr;desc="Inertia SSR";dur='.number_format((float) $ssrMs, 1, '.', '');
|
||||
}
|
||||
|
||||
$metric = implode(', ', $metrics);
|
||||
$existing = $response->headers->get('Server-Timing');
|
||||
|
||||
$response->headers->set(
|
||||
'Server-Timing',
|
||||
$existing ? $existing.', '.$metric : $metric,
|
||||
);
|
||||
|
||||
return $response;
|
||||
}
|
||||
}
|
||||
@@ -98,6 +98,15 @@ final class HandleInertiaRequests extends Middleware
|
||||
? $request->session()->get($key)
|
||||
: null;
|
||||
|
||||
if ($user !== null && ! $user->relationLoaded('profile')) {
|
||||
$user->load('profile');
|
||||
}
|
||||
|
||||
$studioGroups = [];
|
||||
if ($user !== null && str_starts_with($request->path(), 'studio')) {
|
||||
$studioGroups = app(GroupService::class)->studioOptionsForUser($user);
|
||||
}
|
||||
|
||||
return array_merge(parent::share($request), [
|
||||
'auth' => [
|
||||
'user' => $user ? [
|
||||
@@ -134,9 +143,7 @@ final class HandleInertiaRequests extends Middleware
|
||||
'group_assets' => (bool) config('features.group_assets', true),
|
||||
'group_activity_feed' => (bool) config('features.group_activity_feed', true),
|
||||
],
|
||||
'studio_groups' => $user
|
||||
? app(GroupService::class)->studioOptionsForUser($user)
|
||||
: [],
|
||||
'studio_groups' => $studioGroups,
|
||||
]);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,74 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Http\Middleware;
|
||||
|
||||
use App\Support\Http\HttpUriNormalizer;
|
||||
use Closure;
|
||||
use Illuminate\Http\Request;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
use Inertia\Support\Header as InertiaHeader;
|
||||
use Symfony\Component\HttpFoundation\Response;
|
||||
|
||||
final class LogSlowHttpRequest
|
||||
{
|
||||
public function __construct(private readonly HttpUriNormalizer $uris)
|
||||
{
|
||||
}
|
||||
|
||||
public function handle(Request $request, Closure $next): Response
|
||||
{
|
||||
if (! (bool) config('http_observability.slow_request.enabled', false)) {
|
||||
return $next($request);
|
||||
}
|
||||
|
||||
$started = hrtime(true);
|
||||
$response = $next($request);
|
||||
$durationMs = (hrtime(true) - $started) / 1_000_000;
|
||||
|
||||
if (app()->environment('testing')) {
|
||||
$override = config('http_observability.slow_request.test_duration_ms');
|
||||
if (is_numeric($override)) {
|
||||
$durationMs = (float) $override;
|
||||
}
|
||||
}
|
||||
|
||||
$threshold = max(1, (int) config('http_observability.slow_request.threshold_ms', 750));
|
||||
if ($durationMs < $threshold) {
|
||||
return $response;
|
||||
}
|
||||
|
||||
$route = $request->route();
|
||||
$routeUri = is_object($route) && method_exists($route, 'uri')
|
||||
? (string) $route->uri()
|
||||
: $this->uris->normalize($request->path());
|
||||
|
||||
$payload = [
|
||||
'time' => now()->toIso8601String(),
|
||||
'method' => $request->getMethod(),
|
||||
'route_name' => is_object($route) ? $route->getName() : null,
|
||||
'route_uri' => $routeUri,
|
||||
'status' => $response->getStatusCode(),
|
||||
'duration_ms' => round($durationMs, 1),
|
||||
'peak_memory_mb' => round(memory_get_peak_usage(true) / 1048576, 2),
|
||||
'authenticated' => $request->user() !== null,
|
||||
'inertia' => $this->isInertia($request, $response),
|
||||
];
|
||||
|
||||
Log::channel('slow-http')->info(json_encode($payload, JSON_UNESCAPED_SLASHES));
|
||||
|
||||
return $response;
|
||||
}
|
||||
|
||||
private function isInertia(Request $request, Response $response): bool
|
||||
{
|
||||
if ($request->headers->has(InertiaHeader::INERTIA)) {
|
||||
return true;
|
||||
}
|
||||
|
||||
return $response->headers->has(InertiaHeader::INERTIA)
|
||||
|| $response->headers->get('Vary') === 'X-Inertia'
|
||||
|| str_contains((string) $response->headers->get('Vary'), 'X-Inertia');
|
||||
}
|
||||
}
|
||||
@@ -69,6 +69,11 @@ final class TrackOnlineVisitor
|
||||
'email/*',
|
||||
'logout',
|
||||
'up',
|
||||
'rss/*',
|
||||
'rss-feeds',
|
||||
'robots.txt',
|
||||
'sitemap.xml',
|
||||
'sitemaps/*',
|
||||
])) {
|
||||
return false;
|
||||
}
|
||||
|
||||
@@ -0,0 +1,27 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Jobs\Concerns;
|
||||
|
||||
/**
|
||||
* M1 added $afterArtworkId to RecCompute* jobs. Payloads serialized before that
|
||||
* property existed leave a typed property uninitialized on unserialize, which
|
||||
* fatals on first access. Constructor defaults do not apply to unserialize.
|
||||
*/
|
||||
trait RestoresAfterArtworkIdCursor
|
||||
{
|
||||
protected function restoreAfterArtworkIdCursor(): void
|
||||
{
|
||||
if (! isset($this->afterArtworkId)) {
|
||||
$this->afterArtworkId = null;
|
||||
}
|
||||
}
|
||||
|
||||
protected function afterArtworkId(): ?int
|
||||
{
|
||||
$this->restoreAfterArtworkIdCursor();
|
||||
|
||||
return $this->afterArtworkId;
|
||||
}
|
||||
}
|
||||
@@ -6,6 +6,7 @@ namespace App\Jobs;
|
||||
|
||||
use App\Models\Artwork;
|
||||
use Illuminate\Bus\Queueable;
|
||||
use Illuminate\Contracts\Queue\ShouldBeUniqueUntilProcessing;
|
||||
use Illuminate\Contracts\Queue\ShouldQueue;
|
||||
use Illuminate\Foundation\Bus\Dispatchable;
|
||||
use Illuminate\Queue\InteractsWithQueue;
|
||||
@@ -21,16 +22,29 @@ use Meilisearch\Client as MeilisearchClient;
|
||||
* after_commit double-dispatch problem and ensures the document lands
|
||||
* in the index within this job's execution, with no extra queue hop.
|
||||
*/
|
||||
class IndexArtworkJob implements ShouldQueue
|
||||
class IndexArtworkJob implements ShouldQueue, ShouldBeUniqueUntilProcessing
|
||||
{
|
||||
use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
|
||||
|
||||
public int $tries = 3;
|
||||
public int $timeout = 60;
|
||||
|
||||
/**
|
||||
* Unique while queued so Meilisearch isn't flooded. handle() loads the
|
||||
* artwork from DB, so a dropped duplicate still indexes current state.
|
||||
* UniqueUntilProcessing: a change after the worker starts can enqueue again.
|
||||
*/
|
||||
public int $uniqueFor = 120;
|
||||
|
||||
public function __construct(public readonly int $artworkId)
|
||||
{
|
||||
$this->afterCommit = true;
|
||||
$this->onQueue((string) config('scout.queue.queue', 'search'));
|
||||
}
|
||||
|
||||
public function uniqueId(): string
|
||||
{
|
||||
return (string) $this->artworkId;
|
||||
}
|
||||
|
||||
public function handle(MeilisearchClient $client): void
|
||||
|
||||
@@ -6,6 +6,7 @@ namespace App\Jobs;
|
||||
|
||||
use App\Models\User;
|
||||
use Illuminate\Bus\Queueable;
|
||||
use Illuminate\Contracts\Queue\ShouldBeUniqueUntilProcessing;
|
||||
use Illuminate\Contracts\Queue\ShouldQueue;
|
||||
use Illuminate\Foundation\Bus\Dispatchable;
|
||||
use Illuminate\Queue\InteractsWithQueue;
|
||||
@@ -15,14 +16,28 @@ use Illuminate\Queue\SerializesModels;
|
||||
* Queued job: index (or re-index) a single User in Meilisearch.
|
||||
* Dispatched by UserStatsService whenever stats change.
|
||||
*/
|
||||
class IndexUserJob implements ShouldQueue
|
||||
class IndexUserJob implements ShouldQueue, ShouldBeUniqueUntilProcessing
|
||||
{
|
||||
use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
|
||||
|
||||
public int $tries = 3;
|
||||
public int $timeout = 30;
|
||||
|
||||
public function __construct(public readonly int $userId) {}
|
||||
/**
|
||||
* Unique while queued. handle() loads the user from DB so the queued job
|
||||
* still reflects later stats changes. uniqueFor covers a dead worker lock.
|
||||
*/
|
||||
public int $uniqueFor = 60;
|
||||
|
||||
public function __construct(public readonly int $userId)
|
||||
{
|
||||
$this->onQueue((string) config('scout.queue.queue', 'search'));
|
||||
}
|
||||
|
||||
public function uniqueId(): string
|
||||
{
|
||||
return (string) $this->userId;
|
||||
}
|
||||
|
||||
public function handle(): void
|
||||
{
|
||||
|
||||
@@ -6,6 +6,7 @@ namespace App\Jobs;
|
||||
|
||||
use App\Services\RankingService;
|
||||
use Illuminate\Bus\Queueable;
|
||||
use Illuminate\Contracts\Queue\ShouldBeUnique;
|
||||
use Illuminate\Contracts\Queue\ShouldQueue;
|
||||
use Illuminate\Foundation\Bus\Dispatchable;
|
||||
use Illuminate\Queue\InteractsWithQueue;
|
||||
@@ -13,19 +14,32 @@ use Illuminate\Queue\SerializesModels;
|
||||
use Illuminate\Support\Facades\DB;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
|
||||
class RankBuildScopeListsJob implements ShouldQueue
|
||||
class RankBuildScopeListsJob implements ShouldQueue, ShouldBeUnique
|
||||
{
|
||||
use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
|
||||
|
||||
public int $timeout = 300;
|
||||
public int $tries = 2;
|
||||
|
||||
/**
|
||||
* Unique until the job finishes or uniqueFor elapses.
|
||||
* Must outlast the queue wait (hours) plus timeout (300s), not just 360s.
|
||||
*/
|
||||
public int $uniqueFor;
|
||||
|
||||
private const LIST_TYPES = ['trending', 'new_hot', 'best'];
|
||||
|
||||
public function __construct(
|
||||
public readonly string $scopeType,
|
||||
public readonly int $scopeId,
|
||||
) {}
|
||||
) {
|
||||
$this->uniqueFor = max(3600, (int) config('ranking.scope_job_unique_for', 21600));
|
||||
}
|
||||
|
||||
public function uniqueId(): string
|
||||
{
|
||||
return $this->scopeType . ':' . $this->scopeId;
|
||||
}
|
||||
|
||||
public function handle(RankingService $ranking): void
|
||||
{
|
||||
|
||||
@@ -10,8 +10,10 @@ use Illuminate\Bus\Queueable;
|
||||
use Illuminate\Contracts\Queue\ShouldQueue;
|
||||
use Illuminate\Foundation\Bus\Dispatchable;
|
||||
use Illuminate\Queue\InteractsWithQueue;
|
||||
use App\Jobs\Concerns\RestoresAfterArtworkIdCursor;
|
||||
use Illuminate\Queue\SerializesModels;
|
||||
use Illuminate\Support\Facades\DB;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
|
||||
/**
|
||||
* Compute behavior-based (co-like) similarity from precomputed item pairs.
|
||||
@@ -21,38 +23,129 @@ use Illuminate\Support\Facades\DB;
|
||||
*/
|
||||
final class RecComputeSimilarByBehaviorJob implements ShouldQueue
|
||||
{
|
||||
use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
|
||||
use Dispatchable, InteractsWithQueue, Queueable, RestoresAfterArtworkIdCursor;
|
||||
use SerializesModels {
|
||||
__unserialize as private unserializeQueuedModels;
|
||||
}
|
||||
|
||||
public int $tries = 2;
|
||||
public int $timeout = 600;
|
||||
public int $timeout = 120;
|
||||
|
||||
private ?int $afterArtworkId = null;
|
||||
|
||||
public function __construct(
|
||||
private readonly ?int $artworkId = null,
|
||||
private readonly int $batchSize = 200,
|
||||
?int $afterArtworkId = null,
|
||||
) {
|
||||
$this->afterArtworkId = $afterArtworkId;
|
||||
$queue = (string) config('recommendations.queue', 'default');
|
||||
if ($queue !== '') {
|
||||
$this->onQueue($queue);
|
||||
}
|
||||
}
|
||||
|
||||
public function __unserialize(array $values): void
|
||||
{
|
||||
$this->unserializeQueuedModels($values);
|
||||
$this->restoreAfterArtworkIdCursor();
|
||||
}
|
||||
|
||||
public function cursorAfterArtworkId(): ?int
|
||||
{
|
||||
return $this->afterArtworkId();
|
||||
}
|
||||
|
||||
public function handle(): void
|
||||
{
|
||||
$startedAt = microtime(true);
|
||||
$modelVersion = (string) config('recommendations.similarity.model_version', 'sim_v1');
|
||||
$resultLimit = (int) config('recommendations.similarity.result_limit', 30);
|
||||
$maxPerAuthor = (int) config('recommendations.similarity.max_per_author', 2);
|
||||
|
||||
$query = Artwork::query()->public()->published()->select('id', 'user_id');
|
||||
|
||||
if ($this->artworkId !== null) {
|
||||
$query->where('id', $this->artworkId);
|
||||
$artwork = Artwork::query()->public()->published()->select('id', 'user_id')->find($this->artworkId);
|
||||
|
||||
if ($artwork instanceof Artwork) {
|
||||
$this->processArtworkSafely($artwork, $modelVersion, $resultLimit, $maxPerAuthor);
|
||||
$this->logBatchComplete($startedAt, 1, false);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
$this->logBatchComplete($startedAt, 0, false);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
$query->chunkById($this->batchSize, function ($artworks) use ($modelVersion, $resultLimit, $maxPerAuthor) {
|
||||
foreach ($artworks as $artwork) {
|
||||
$this->processArtwork($artwork, $modelVersion, $resultLimit, $maxPerAuthor);
|
||||
}
|
||||
});
|
||||
$artworks = Artwork::query()
|
||||
->public()
|
||||
->published()
|
||||
->select('id', 'user_id')
|
||||
->when($this->afterArtworkId() !== null, fn ($query) => $query->where('id', '>', $this->afterArtworkId()))
|
||||
->orderBy('id')
|
||||
->limit($this->batchSize)
|
||||
->get();
|
||||
|
||||
if ($artworks->isEmpty()) {
|
||||
$this->logBatchComplete($startedAt, 0, false);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
foreach ($artworks as $artwork) {
|
||||
$this->processArtworkSafely($artwork, $modelVersion, $resultLimit, $maxPerAuthor);
|
||||
}
|
||||
|
||||
$hasMore = $artworks->count() === $this->batchSize;
|
||||
if ($hasMore) {
|
||||
static::dispatch(null, $this->batchSize, (int) $artworks->last()->id);
|
||||
}
|
||||
|
||||
$this->logBatchComplete($startedAt, $artworks->count(), $hasMore);
|
||||
}
|
||||
|
||||
public function failed(\Throwable $exception): void
|
||||
{
|
||||
Log::error('[RecComputeSimilarByBehavior] Job failed permanently.', [
|
||||
'artwork_id' => $this->artworkId,
|
||||
'batch_size' => $this->batchSize,
|
||||
'after_artwork_id' => $this->afterArtworkId(),
|
||||
'attempts' => $this->attempts(),
|
||||
'exception_class' => $exception::class,
|
||||
'exception_message' => $exception->getMessage(),
|
||||
]);
|
||||
}
|
||||
|
||||
private function processArtworkSafely(
|
||||
Artwork $artwork,
|
||||
string $modelVersion,
|
||||
int $resultLimit,
|
||||
int $maxPerAuthor,
|
||||
): void {
|
||||
try {
|
||||
$this->processArtwork($artwork, $modelVersion, $resultLimit, $maxPerAuthor);
|
||||
} catch (\Throwable $exception) {
|
||||
Log::warning("[RecComputeSimilarByBehavior] Failed for artwork {$artwork->id}: {$exception->getMessage()}", [
|
||||
'artwork_id' => $artwork->id,
|
||||
'exception_class' => $exception::class,
|
||||
]);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* @param positive-int|0 $processed
|
||||
*/
|
||||
private function logBatchComplete(float $startedAt, int $processed, bool $hasMore): void
|
||||
{
|
||||
Log::info('[RecComputeSimilarByBehavior] Batch complete.', [
|
||||
'artwork_id' => $this->artworkId,
|
||||
'after_artwork_id' => $this->afterArtworkId(),
|
||||
'processed' => $processed,
|
||||
'has_more' => $hasMore,
|
||||
'duration_ms' => (int) round((microtime(true) - $startedAt) * 1000),
|
||||
'memory_mb' => round(memory_get_peak_usage(true) / 1048576, 1),
|
||||
]);
|
||||
}
|
||||
|
||||
private function processArtwork(
|
||||
|
||||
@@ -11,7 +11,9 @@ use Illuminate\Contracts\Queue\ShouldQueue;
|
||||
use Illuminate\Foundation\Bus\Dispatchable;
|
||||
use Illuminate\Queue\Middleware\WithoutOverlapping;
|
||||
use Illuminate\Queue\InteractsWithQueue;
|
||||
use App\Jobs\Concerns\RestoresAfterArtworkIdCursor;
|
||||
use Illuminate\Queue\SerializesModels;
|
||||
use Illuminate\Support\Facades\Cache;
|
||||
use Illuminate\Support\Facades\DB;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
|
||||
@@ -24,22 +26,44 @@ use Illuminate\Support\Facades\Log;
|
||||
*/
|
||||
final class RecComputeSimilarByTagsJob implements ShouldQueue
|
||||
{
|
||||
use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
|
||||
use Dispatchable, InteractsWithQueue, Queueable, RestoresAfterArtworkIdCursor;
|
||||
use SerializesModels {
|
||||
__unserialize as private unserializeQueuedModels;
|
||||
}
|
||||
|
||||
public int $tries = 2;
|
||||
public int $timeout = 600;
|
||||
|
||||
/**
|
||||
* Declared with a default so missing serialized payloads do not leave
|
||||
* an uninitialized typed property. Constructor defaults are not applied
|
||||
* on unserialize.
|
||||
*/
|
||||
private ?int $afterArtworkId = null;
|
||||
|
||||
public function __construct(
|
||||
private readonly ?int $artworkId = null,
|
||||
private readonly int $batchSize = 200,
|
||||
private readonly ?int $afterArtworkId = null,
|
||||
?int $afterArtworkId = null,
|
||||
) {
|
||||
$this->afterArtworkId = $afterArtworkId;
|
||||
$queue = (string) config('recommendations.queue', 'default');
|
||||
if ($queue !== '') {
|
||||
$this->onQueue($queue);
|
||||
}
|
||||
}
|
||||
|
||||
public function __unserialize(array $values): void
|
||||
{
|
||||
$this->unserializeQueuedModels($values);
|
||||
$this->restoreAfterArtworkIdCursor();
|
||||
}
|
||||
|
||||
public function cursorAfterArtworkId(): ?int
|
||||
{
|
||||
return $this->afterArtworkId();
|
||||
}
|
||||
|
||||
/**
|
||||
* @return array<int, object>
|
||||
*/
|
||||
@@ -63,12 +87,7 @@ final class RecComputeSimilarByTagsJob implements ShouldQueue
|
||||
$maxPerAuthor = (int) config('recommendations.similarity.max_per_author', 2);
|
||||
$resultLimit = (int) config('recommendations.similarity.result_limit', 30);
|
||||
|
||||
// ── Tag IDF weights (global) ───────────────────────────────────────────
|
||||
$tagFreqs = DB::table('artwork_tag')
|
||||
->select('tag_id', DB::raw('COUNT(*) as cnt'))
|
||||
->groupBy('tag_id')
|
||||
->pluck('cnt', 'tag_id')
|
||||
->all();
|
||||
$tagFreqs = $this->tagFrequencies();
|
||||
|
||||
if ($this->artworkId !== null) {
|
||||
$artwork = Artwork::query()->public()->published()->select('id', 'user_id')->find($this->artworkId);
|
||||
@@ -86,7 +105,7 @@ final class RecComputeSimilarByTagsJob implements ShouldQueue
|
||||
->public()
|
||||
->published()
|
||||
->select('id', 'user_id')
|
||||
->when($this->afterArtworkId !== null, fn ($query) => $query->where('id', '>', $this->afterArtworkId))
|
||||
->when($this->afterArtworkId() !== null, fn ($query) => $query->where('id', '>', $this->afterArtworkId()))
|
||||
->orderBy('id')
|
||||
->limit($this->batchSize)
|
||||
->get();
|
||||
@@ -109,13 +128,33 @@ final class RecComputeSimilarByTagsJob implements ShouldQueue
|
||||
Log::error('[RecComputeSimilarByTags] Job failed permanently.', [
|
||||
'artwork_id' => $this->artworkId,
|
||||
'batch_size' => $this->batchSize,
|
||||
'after_artwork_id' => $this->afterArtworkId,
|
||||
'after_artwork_id' => $this->afterArtworkId(),
|
||||
'attempts' => $this->attempts(),
|
||||
'exception_class' => $exception::class,
|
||||
'exception_message' => $exception->getMessage(),
|
||||
]);
|
||||
}
|
||||
|
||||
/**
|
||||
* Global tag frequencies change slowly relative to batch jobs.
|
||||
* Cache so each RecComputeSimilarByTagsJob batch does not re-scan artwork_tag.
|
||||
*
|
||||
* @return array<int|string, mixed>
|
||||
*/
|
||||
private function tagFrequencies(): array
|
||||
{
|
||||
$modelVersion = (string) config('recommendations.similarity.model_version', 'sim_v1');
|
||||
$ttl = (int) config('recommendations.similarity.tag_idf_cache_ttl', 3600);
|
||||
|
||||
return Cache::remember('rec:tag-idf:'.$modelVersion, max(60, $ttl), function () {
|
||||
return DB::table('artwork_tag')
|
||||
->select('tag_id', DB::raw('COUNT(*) as cnt'))
|
||||
->groupBy('tag_id')
|
||||
->pluck('cnt', 'tag_id')
|
||||
->all();
|
||||
});
|
||||
}
|
||||
|
||||
private function processArtworkSafely(
|
||||
Artwork $artwork,
|
||||
array $tagFreqs,
|
||||
|
||||
@@ -11,6 +11,7 @@ use Illuminate\Contracts\Queue\ShouldQueue;
|
||||
use Illuminate\Foundation\Bus\Dispatchable;
|
||||
use Illuminate\Queue\Middleware\WithoutOverlapping;
|
||||
use Illuminate\Queue\InteractsWithQueue;
|
||||
use App\Jobs\Concerns\RestoresAfterArtworkIdCursor;
|
||||
use Illuminate\Queue\SerializesModels;
|
||||
use Illuminate\Support\Facades\DB;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
@@ -24,24 +25,42 @@ use Illuminate\Support\Facades\Log;
|
||||
*/
|
||||
final class RecComputeSimilarHybridJob implements ShouldQueue
|
||||
{
|
||||
use Dispatchable, InteractsWithQueue, Queueable, SerializesModels;
|
||||
use Dispatchable, InteractsWithQueue, Queueable, RestoresAfterArtworkIdCursor;
|
||||
use SerializesModels {
|
||||
__unserialize as private unserializeQueuedModels;
|
||||
}
|
||||
|
||||
// This recompute is idempotent and already guards per-artwork execution.
|
||||
// Keep retries to a minimum so transient failures do not turn into
|
||||
// Horizon's max-attempt exception noise.
|
||||
public int $tries = 1;
|
||||
public int $timeout = 900;
|
||||
public int $timeout = 180;
|
||||
|
||||
private ?int $afterArtworkId = null;
|
||||
|
||||
public function __construct(
|
||||
private readonly ?int $artworkId = null,
|
||||
private readonly int $batchSize = 200,
|
||||
?int $afterArtworkId = null,
|
||||
) {
|
||||
$this->afterArtworkId = $afterArtworkId;
|
||||
$queue = (string) config('recommendations.queue', 'default');
|
||||
if ($queue !== '') {
|
||||
$this->onQueue($queue);
|
||||
}
|
||||
}
|
||||
|
||||
public function __unserialize(array $values): void
|
||||
{
|
||||
$this->unserializeQueuedModels($values);
|
||||
$this->restoreAfterArtworkIdCursor();
|
||||
}
|
||||
|
||||
public function cursorAfterArtworkId(): ?int
|
||||
{
|
||||
return $this->afterArtworkId();
|
||||
}
|
||||
|
||||
/**
|
||||
* @return array<int, object>
|
||||
*/
|
||||
@@ -72,6 +91,8 @@ final class RecComputeSimilarHybridJob implements ShouldQueue
|
||||
? (array) config('recommendations.similarity.weights_with_vector')
|
||||
: (array) config('recommendations.similarity.weights_without_vector');
|
||||
|
||||
$startedAt = microtime(true);
|
||||
|
||||
if ($this->artworkId !== null) {
|
||||
$artwork = Artwork::query()
|
||||
->public()
|
||||
@@ -80,6 +101,8 @@ final class RecComputeSimilarHybridJob implements ShouldQueue
|
||||
->find($this->artworkId);
|
||||
|
||||
if (! $artwork instanceof Artwork) {
|
||||
$this->logBatchComplete($startedAt, 0, false);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
@@ -93,26 +116,42 @@ final class RecComputeSimilarHybridJob implements ShouldQueue
|
||||
$weights,
|
||||
);
|
||||
|
||||
$this->logBatchComplete($startedAt, 1, false);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
Artwork::query()
|
||||
$artworks = Artwork::query()
|
||||
->public()
|
||||
->published()
|
||||
->select('id', 'user_id')
|
||||
->chunkById($this->batchSize, function ($artworks) use (
|
||||
$modelVersion, $vectorEnabled, $resultLimit, $maxPerAuthor, $minCatsTop12, $weights
|
||||
) {
|
||||
$this->processArtworkSafely(
|
||||
$artworks,
|
||||
$modelVersion,
|
||||
$vectorEnabled,
|
||||
$resultLimit,
|
||||
$maxPerAuthor,
|
||||
$minCatsTop12,
|
||||
$weights,
|
||||
);
|
||||
});
|
||||
->when($this->afterArtworkId() !== null, fn ($query) => $query->where('id', '>', $this->afterArtworkId()))
|
||||
->orderBy('id')
|
||||
->limit($this->batchSize)
|
||||
->get();
|
||||
|
||||
if ($artworks->isEmpty()) {
|
||||
$this->logBatchComplete($startedAt, 0, false);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
$this->processArtworkSafely(
|
||||
$artworks,
|
||||
$modelVersion,
|
||||
$vectorEnabled,
|
||||
$resultLimit,
|
||||
$maxPerAuthor,
|
||||
$minCatsTop12,
|
||||
$weights,
|
||||
);
|
||||
|
||||
$hasMore = $artworks->count() === $this->batchSize;
|
||||
if ($hasMore) {
|
||||
static::dispatch(null, $this->batchSize, (int) $artworks->last()->id);
|
||||
}
|
||||
|
||||
$this->logBatchComplete($startedAt, $artworks->count(), $hasMore);
|
||||
}
|
||||
|
||||
public function failed(\Throwable $exception): void
|
||||
@@ -120,12 +159,28 @@ final class RecComputeSimilarHybridJob implements ShouldQueue
|
||||
Log::error('[RecComputeSimilarHybrid] Job failed permanently.', [
|
||||
'artwork_id' => $this->artworkId,
|
||||
'batch_size' => $this->batchSize,
|
||||
'after_artwork_id' => $this->afterArtworkId(),
|
||||
'attempts' => $this->attempts(),
|
||||
'exception_class' => $exception::class,
|
||||
'exception_message' => $exception->getMessage(),
|
||||
]);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param positive-int|0 $processed
|
||||
*/
|
||||
private function logBatchComplete(float $startedAt, int $processed, bool $hasMore): void
|
||||
{
|
||||
Log::info('[RecComputeSimilarHybrid] Batch complete.', [
|
||||
'artwork_id' => $this->artworkId,
|
||||
'after_artwork_id' => $this->afterArtworkId(),
|
||||
'processed' => $processed,
|
||||
'has_more' => $hasMore,
|
||||
'duration_ms' => (int) round((microtime(true) - $startedAt) * 1000),
|
||||
'memory_mb' => round(memory_get_peak_usage(true) / 1048576, 1),
|
||||
]);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param iterable<Artwork> $artworks
|
||||
*/
|
||||
|
||||
@@ -58,6 +58,11 @@ class AppServiceProvider extends ServiceProvider
|
||||
*/
|
||||
public function register(): void
|
||||
{
|
||||
$this->app->bind(
|
||||
\Inertia\Ssr\Gateway::class,
|
||||
\App\Support\Http\TimedInertiaSsrGateway::class,
|
||||
);
|
||||
|
||||
$this->app->singleton(
|
||||
\App\Services\Countries\CountryRemoteProviderInterface::class,
|
||||
\App\Services\Countries\CountryRemoteProvider::class,
|
||||
|
||||
@@ -13,6 +13,7 @@ use App\Models\AcademyLessonBlock;
|
||||
use App\Models\AcademyPromptPack;
|
||||
use App\Models\AcademyPromptTemplate;
|
||||
use App\Models\User;
|
||||
use Illuminate\Support\Facades\Cache;
|
||||
use Illuminate\Support\Facades\Storage;
|
||||
use Illuminate\Support\Str;
|
||||
use Laravel\Cashier\Subscription;
|
||||
@@ -280,7 +281,9 @@ final class AcademyAccessService
|
||||
public function coursePayload(AcademyCourse $course, ?User $viewer, array $options = []): array
|
||||
{
|
||||
$progress = is_array($options['progress'] ?? null) ? $options['progress'] : null;
|
||||
$lessonCount = (int) ($course->lessons_count_cache ?: $course->courseLessons()->count());
|
||||
$lessonCount = $course->relationLoaded('courseLessons')
|
||||
? $course->courseLessons->count()
|
||||
: (int) $course->lessons_count_cache;
|
||||
|
||||
return [
|
||||
'id' => (int) $course->id,
|
||||
@@ -1089,7 +1092,13 @@ final class AcademyAccessService
|
||||
}
|
||||
|
||||
try {
|
||||
$exists = Storage::disk((string) config('uploads.object_storage.disk', 's3'))->exists($normalizedPath);
|
||||
$exists = (bool) Cache::remember(
|
||||
'academy:asset-exists:'.$cacheKey,
|
||||
3600,
|
||||
function () use ($normalizedPath): bool {
|
||||
return Storage::disk((string) config('uploads.object_storage.disk', 's3'))->exists($normalizedPath);
|
||||
},
|
||||
);
|
||||
} catch (\Throwable) {
|
||||
$exists = false;
|
||||
}
|
||||
|
||||
@@ -16,6 +16,10 @@ class CollectionBackgroundJobService
|
||||
{
|
||||
public function dispatchQualityRefresh(Collection $collection, ?User $actor = null): array
|
||||
{
|
||||
if (! $this->dispatchEnabled()) {
|
||||
return $this->disabledPayload('quality_refresh', $collection ? collect([$collection]) : collect());
|
||||
}
|
||||
|
||||
RefreshCollectionQualityJob::dispatch((int) $collection->id, $actor?->id)->afterCommit();
|
||||
|
||||
return [
|
||||
@@ -30,6 +34,10 @@ class CollectionBackgroundJobService
|
||||
|
||||
public function dispatchHealthRefresh(?Collection $collection = null, ?User $actor = null): array
|
||||
{
|
||||
if (! $this->dispatchEnabled()) {
|
||||
return $this->disabledPayload('health_refresh', $collection ? collect([$collection]) : collect());
|
||||
}
|
||||
|
||||
$targets = $collection ? collect([$collection]) : $this->healthTargets();
|
||||
|
||||
$targets->each(fn (Collection $item) => RefreshCollectionHealthJob::dispatch((int) $item->id, $actor?->id, 'programming-eligibility')->afterCommit());
|
||||
@@ -39,6 +47,10 @@ class CollectionBackgroundJobService
|
||||
|
||||
public function dispatchRecommendationRefresh(?Collection $collection = null, ?User $actor = null, string $context = 'default'): array
|
||||
{
|
||||
if (! $this->dispatchEnabled()) {
|
||||
return $this->disabledPayload('recommendation_refresh', $collection ? collect([$collection]) : collect());
|
||||
}
|
||||
|
||||
$targets = $collection ? collect([$collection]) : $this->recommendationTargets();
|
||||
|
||||
$targets->each(fn (Collection $item) => RefreshCollectionRecommendationJob::dispatch((int) $item->id, $actor?->id, $context)->afterCommit());
|
||||
@@ -48,6 +60,10 @@ class CollectionBackgroundJobService
|
||||
|
||||
public function dispatchDuplicateScan(?Collection $collection = null, ?User $actor = null): array
|
||||
{
|
||||
if (! $this->dispatchEnabled()) {
|
||||
return $this->disabledPayload('duplicate_scan', $collection ? collect([$collection]) : collect());
|
||||
}
|
||||
|
||||
$targets = $collection ? collect([$collection]) : $this->duplicateTargets();
|
||||
|
||||
$targets->each(fn (Collection $item) => ScanCollectionDuplicateCandidatesJob::dispatch((int) $item->id, $actor?->id)->afterCommit());
|
||||
@@ -128,6 +144,29 @@ class CollectionBackgroundJobService
|
||||
->get(['id']);
|
||||
}
|
||||
|
||||
private function dispatchEnabled(): bool
|
||||
{
|
||||
return (bool) config('collections.v5.queue.dispatch_enabled', false);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param SupportCollection<int, Collection> $targets
|
||||
*/
|
||||
private function disabledPayload(string $job, SupportCollection $targets): array
|
||||
{
|
||||
$ids = $targets->pluck('id')->map(static fn ($id): int => (int) $id)->values()->all();
|
||||
|
||||
return [
|
||||
'status' => 'disabled',
|
||||
'job' => $job,
|
||||
'scope' => count($ids) === 1 ? 'single' : 'batch',
|
||||
'count' => 0,
|
||||
'collection_ids' => $ids,
|
||||
'items' => [],
|
||||
'message' => 'Collection queue dispatch is disabled until a Horizon consumer exists.',
|
||||
];
|
||||
}
|
||||
|
||||
/**
|
||||
* @param SupportCollection<int, Collection> $targets
|
||||
*/
|
||||
|
||||
@@ -5,8 +5,8 @@ declare(strict_types=1);
|
||||
namespace App\Services;
|
||||
|
||||
use App\Models\Artwork;
|
||||
use App\Models\ArtworkMetricSnapshotHourly;
|
||||
use App\Models\Group;
|
||||
use App\Services\Metrics\ArtworkHourlySnapshotWindow;
|
||||
use App\Models\Leaderboard;
|
||||
use App\Models\Story;
|
||||
use App\Models\StoryLike;
|
||||
@@ -28,6 +28,7 @@ class LeaderboardService
|
||||
|
||||
public function __construct(
|
||||
private readonly GroupReputationService $groupReputation,
|
||||
private readonly ArtworkHourlySnapshotWindow $snapshotWindow,
|
||||
) {
|
||||
}
|
||||
|
||||
@@ -835,29 +836,7 @@ class LeaderboardService
|
||||
|
||||
private function artworkSnapshotDeltas(CarbonImmutable $start): \Illuminate\Database\Query\Builder
|
||||
{
|
||||
return ArtworkMetricSnapshotHourly::query()
|
||||
->from('artwork_metric_snapshots_hourly as snapshots')
|
||||
->where('snapshots.bucket_hour', '>=', $start)
|
||||
->select([
|
||||
'snapshots.artwork_id',
|
||||
DB::raw($this->nonNegativeSnapshotDelta('views_count', 'views_delta')),
|
||||
DB::raw($this->nonNegativeSnapshotDelta('downloads_count', 'downloads_delta')),
|
||||
DB::raw($this->nonNegativeSnapshotDelta('favourites_count', 'favourites_delta')),
|
||||
DB::raw($this->nonNegativeSnapshotDelta('comments_count', 'comments_delta')),
|
||||
])
|
||||
->groupBy('snapshots.artwork_id')
|
||||
->toBase();
|
||||
}
|
||||
|
||||
private function nonNegativeSnapshotDelta(string $column, string $alias): string
|
||||
{
|
||||
$delta = sprintf('MAX(snapshots.%1$s) - MIN(snapshots.%1$s)', $column);
|
||||
|
||||
if (DB::connection()->getDriverName() === 'sqlite') {
|
||||
return sprintf('CASE WHEN %1$s > 0 THEN %1$s ELSE 0 END as %2$s', $delta, $alias);
|
||||
}
|
||||
|
||||
return sprintf('GREATEST(%s, 0) as %s', $delta, $alias);
|
||||
return $this->snapshotWindow->artworkPeriodDeltas($start);
|
||||
}
|
||||
|
||||
private function creatorEntities(array $ids): array
|
||||
|
||||
@@ -0,0 +1,260 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Services\Metrics;
|
||||
|
||||
use Carbon\CarbonInterface;
|
||||
use Illuminate\Database\Query\Builder;
|
||||
use Illuminate\Support\Facades\DB;
|
||||
use Illuminate\Support\Facades\Schema;
|
||||
|
||||
/**
|
||||
* Windowed deltas from cumulative hourly artwork metric snapshots.
|
||||
*
|
||||
* Each snapshot stores running totals. Period growth is:
|
||||
* latest in window
|
||||
* - baseline at or immediately before the window start
|
||||
*
|
||||
* If the artwork was created/published during the window, baseline is 0.
|
||||
* If an older artwork has no pre-window snapshot (warm-up / missing hours),
|
||||
* fall back to MIN in the window — never SUM cumulatives, never lifetime totals.
|
||||
*/
|
||||
final class ArtworkHourlySnapshotWindow
|
||||
{
|
||||
/**
|
||||
* How far before the window start to look for a baseline snapshot.
|
||||
* Keeps the baseline scan bounded at 30-day table scale.
|
||||
*/
|
||||
public const BASELINE_LOOKBACK_HOURS = 36;
|
||||
|
||||
/**
|
||||
* @return array{
|
||||
* views: int,
|
||||
* downloads: int,
|
||||
* favourites: int,
|
||||
* comments: int,
|
||||
* shares: int,
|
||||
* requested_days: int,
|
||||
* table_coverage_days: float,
|
||||
* user_coverage_days: float,
|
||||
* window_complete: bool,
|
||||
* oldest_bucket: ?string,
|
||||
* newest_bucket: ?string,
|
||||
* user_oldest_bucket: ?string,
|
||||
* user_newest_bucket: ?string,
|
||||
* }
|
||||
*/
|
||||
public function userWindow(int $userId, int $days = 30): array
|
||||
{
|
||||
$days = max(1, $days);
|
||||
$coverage = $this->tableCoverage($days);
|
||||
$empty = $this->emptyResult($days, $coverage);
|
||||
|
||||
if ($userId <= 0 || ! Schema::hasTable('artwork_metric_snapshots_hourly')) {
|
||||
return $empty;
|
||||
}
|
||||
|
||||
$since = now()->subDays($days);
|
||||
|
||||
$row = DB::query()
|
||||
->fromSub($this->artworkPeriodDeltas($since, $userId), 'deltas')
|
||||
->selectRaw('COALESCE(SUM(views_delta), 0) as views')
|
||||
->selectRaw('COALESCE(SUM(downloads_delta), 0) as downloads')
|
||||
->selectRaw('COALESCE(SUM(favourites_delta), 0) as favourites')
|
||||
->selectRaw('COALESCE(SUM(comments_delta), 0) as comments')
|
||||
->selectRaw('COALESCE(SUM(shares_delta), 0) as shares')
|
||||
->first();
|
||||
|
||||
$userBounds = DB::table('artwork_metric_snapshots_hourly as snapshots')
|
||||
->join('artworks', 'artworks.id', '=', 'snapshots.artwork_id')
|
||||
->where('artworks.user_id', $userId)
|
||||
->whereNull('artworks.deleted_at')
|
||||
->where('snapshots.bucket_hour', '>=', $since)
|
||||
->selectRaw('MIN(snapshots.bucket_hour) as oldest_bucket')
|
||||
->selectRaw('MAX(snapshots.bucket_hour) as newest_bucket')
|
||||
->first();
|
||||
|
||||
$userCoverageDays = $this->coverageDays(
|
||||
$userBounds->oldest_bucket ?? null,
|
||||
$userBounds->newest_bucket ?? null,
|
||||
);
|
||||
|
||||
return [
|
||||
'views' => (int) ($row->views ?? 0),
|
||||
'downloads' => (int) ($row->downloads ?? 0),
|
||||
'favourites' => (int) ($row->favourites ?? 0),
|
||||
'comments' => (int) ($row->comments ?? 0),
|
||||
'shares' => (int) ($row->shares ?? 0),
|
||||
'requested_days' => $days,
|
||||
'table_coverage_days' => $coverage['coverage_days'],
|
||||
'user_coverage_days' => $userCoverageDays,
|
||||
'window_complete' => $coverage['window_complete'],
|
||||
'oldest_bucket' => $coverage['oldest_bucket'],
|
||||
'newest_bucket' => $coverage['newest_bucket'],
|
||||
'user_oldest_bucket' => $userBounds->oldest_bucket ?? null,
|
||||
'user_newest_bucket' => $userBounds->newest_bucket ?? null,
|
||||
];
|
||||
}
|
||||
|
||||
/**
|
||||
* Per-artwork non-negative period deltas for cumulative snapshot columns.
|
||||
*/
|
||||
public function artworkPeriodDeltas(CarbonInterface|\DateTimeInterface $start, ?int $userId = null): Builder
|
||||
{
|
||||
$inWindow = DB::table('artwork_metric_snapshots_hourly as snapshots');
|
||||
$this->constrainOwner($inWindow, $userId);
|
||||
$inWindow
|
||||
->where('snapshots.bucket_hour', '>=', $start)
|
||||
->groupBy('snapshots.artwork_id')
|
||||
->select('snapshots.artwork_id')
|
||||
->selectRaw('MAX(snapshots.views_count) as views_latest')
|
||||
->selectRaw('MIN(snapshots.views_count) as views_min')
|
||||
->selectRaw('MAX(snapshots.downloads_count) as downloads_latest')
|
||||
->selectRaw('MIN(snapshots.downloads_count) as downloads_min')
|
||||
->selectRaw('MAX(snapshots.favourites_count) as favourites_latest')
|
||||
->selectRaw('MIN(snapshots.favourites_count) as favourites_min')
|
||||
->selectRaw('MAX(snapshots.comments_count) as comments_latest')
|
||||
->selectRaw('MIN(snapshots.comments_count) as comments_min')
|
||||
->selectRaw('MAX(snapshots.shares_count) as shares_latest')
|
||||
->selectRaw('MIN(snapshots.shares_count) as shares_min');
|
||||
|
||||
$baselineFrom = \Carbon\Carbon::parse($start)->subHours(self::BASELINE_LOOKBACK_HOURS);
|
||||
$baseline = DB::table('artwork_metric_snapshots_hourly as snapshots');
|
||||
$this->constrainOwner($baseline, $userId);
|
||||
$baseline
|
||||
->where('snapshots.bucket_hour', '<', $start)
|
||||
->where('snapshots.bucket_hour', '>=', $baselineFrom)
|
||||
->groupBy('snapshots.artwork_id')
|
||||
->select('snapshots.artwork_id')
|
||||
->selectRaw('MAX(snapshots.views_count) as views_baseline')
|
||||
->selectRaw('MAX(snapshots.downloads_count) as downloads_baseline')
|
||||
->selectRaw('MAX(snapshots.favourites_count) as favourites_baseline')
|
||||
->selectRaw('MAX(snapshots.comments_count) as comments_baseline')
|
||||
->selectRaw('MAX(snapshots.shares_count) as shares_baseline');
|
||||
|
||||
$bornSql = $this->quotedTimestamp($start);
|
||||
|
||||
return DB::query()
|
||||
->fromSub($inWindow, 'win')
|
||||
->leftJoinSub($baseline, 'base', 'base.artwork_id', '=', 'win.artwork_id')
|
||||
->join('artworks', 'artworks.id', '=', 'win.artwork_id')
|
||||
->select('win.artwork_id')
|
||||
->selectRaw($this->deltaExpression('views', $bornSql) . ' as views_delta')
|
||||
->selectRaw($this->deltaExpression('downloads', $bornSql) . ' as downloads_delta')
|
||||
->selectRaw($this->deltaExpression('favourites', $bornSql) . ' as favourites_delta')
|
||||
->selectRaw($this->deltaExpression('comments', $bornSql) . ' as comments_delta')
|
||||
->selectRaw($this->deltaExpression('shares', $bornSql) . ' as shares_delta');
|
||||
}
|
||||
|
||||
/**
|
||||
* @return array{oldest_bucket: ?string, newest_bucket: ?string, coverage_days: float, window_complete: bool, requested_days: int}
|
||||
*/
|
||||
public function tableCoverage(int $days = 30): array
|
||||
{
|
||||
$days = max(1, $days);
|
||||
|
||||
if (! Schema::hasTable('artwork_metric_snapshots_hourly')) {
|
||||
return [
|
||||
'oldest_bucket' => null,
|
||||
'newest_bucket' => null,
|
||||
'coverage_days' => 0.0,
|
||||
'window_complete' => false,
|
||||
'requested_days' => $days,
|
||||
];
|
||||
}
|
||||
|
||||
$bounds = DB::table('artwork_metric_snapshots_hourly')
|
||||
->selectRaw('MIN(bucket_hour) as oldest_bucket')
|
||||
->selectRaw('MAX(bucket_hour) as newest_bucket')
|
||||
->first();
|
||||
|
||||
$coverageDays = $this->coverageDays(
|
||||
$bounds->oldest_bucket ?? null,
|
||||
$bounds->newest_bucket ?? null,
|
||||
);
|
||||
|
||||
return [
|
||||
'oldest_bucket' => $bounds->oldest_bucket ?? null,
|
||||
'newest_bucket' => $bounds->newest_bucket ?? null,
|
||||
'coverage_days' => $coverageDays,
|
||||
'window_complete' => $coverageDays + (1 / 24) >= $days,
|
||||
'requested_days' => $days,
|
||||
];
|
||||
}
|
||||
|
||||
private function constrainOwner(Builder $query, ?int $userId): void
|
||||
{
|
||||
if ($userId === null) {
|
||||
return;
|
||||
}
|
||||
|
||||
$query->join('artworks as owner_artworks', 'owner_artworks.id', '=', 'snapshots.artwork_id')
|
||||
->where('owner_artworks.user_id', $userId)
|
||||
->whereNull('owner_artworks.deleted_at');
|
||||
}
|
||||
|
||||
private function quotedTimestamp(CarbonInterface|\DateTimeInterface $start): string
|
||||
{
|
||||
return DB::getPdo()->quote(\Carbon\Carbon::parse($start)->toDateTimeString());
|
||||
}
|
||||
|
||||
private function deltaExpression(string $metric, string $bornSql): string
|
||||
{
|
||||
$latest = "COALESCE(win.{$metric}_latest, 0)";
|
||||
$min = "COALESCE(win.{$metric}_min, 0)";
|
||||
$base = "base.{$metric}_baseline";
|
||||
$bornInWindow = "COALESCE(artworks.published_at, artworks.created_at) >= {$bornSql}";
|
||||
$observed = $this->clampNonNegative("{$latest} - {$min}");
|
||||
$fromBaseline = $this->clampNonNegative("{$latest} - {$base}");
|
||||
|
||||
return "CASE WHEN {$bornInWindow} THEN {$latest} WHEN {$base} IS NOT NULL THEN {$fromBaseline} ELSE {$observed} END";
|
||||
}
|
||||
|
||||
private function clampNonNegative(string $expression): string
|
||||
{
|
||||
if (DB::connection()->getDriverName() === 'sqlite') {
|
||||
return "CASE WHEN ({$expression}) > 0 THEN ({$expression}) ELSE 0 END";
|
||||
}
|
||||
|
||||
return "GREATEST({$expression}, 0)";
|
||||
}
|
||||
|
||||
/**
|
||||
* @return array<string, mixed>
|
||||
*/
|
||||
private function emptyResult(int $days, array $coverage): array
|
||||
{
|
||||
return [
|
||||
'views' => 0,
|
||||
'downloads' => 0,
|
||||
'favourites' => 0,
|
||||
'comments' => 0,
|
||||
'shares' => 0,
|
||||
'requested_days' => $days,
|
||||
'table_coverage_days' => $coverage['coverage_days'],
|
||||
'user_coverage_days' => 0.0,
|
||||
'window_complete' => $coverage['window_complete'],
|
||||
'oldest_bucket' => $coverage['oldest_bucket'],
|
||||
'newest_bucket' => $coverage['newest_bucket'],
|
||||
'user_oldest_bucket' => null,
|
||||
'user_newest_bucket' => null,
|
||||
];
|
||||
}
|
||||
|
||||
private function coverageDays(?string $oldest, ?string $newest): float
|
||||
{
|
||||
if ($oldest === null || $newest === null) {
|
||||
return 0.0;
|
||||
}
|
||||
|
||||
$rangeStart = \Carbon\Carbon::parse($oldest);
|
||||
$rangeEnd = \Carbon\Carbon::parse($newest);
|
||||
|
||||
if ($rangeEnd->lessThan($rangeStart)) {
|
||||
return 0.0;
|
||||
}
|
||||
|
||||
return round(abs($rangeStart->diffInMinutes($rangeEnd)) / 1440, 2);
|
||||
}
|
||||
}
|
||||
@@ -204,7 +204,7 @@ final class CreatorJourneyService
|
||||
[
|
||||
'title' => 'Biggest download spike',
|
||||
'headline' => (string) $bestSpike['artwork']->title,
|
||||
'summary' => 'Captured the strongest one-hour download burst recorded for a public artwork.',
|
||||
'summary' => 'Captured the strongest one-hour download burst in the retained hourly snapshot window.',
|
||||
'value' => (int) $bestSpike['downloads_in_hour'] . ' downloads in 1 hour',
|
||||
'artwork' => $this->artworkSnapshot($bestSpike['artwork']),
|
||||
'metrics' => [
|
||||
@@ -425,8 +425,12 @@ final class CreatorJourneyService
|
||||
$publicArtworkIds = $artworks->pluck('id')->map(fn ($id): int => (int) $id)->all();
|
||||
|
||||
if ($publicArtworkIds !== [] && DB::getSchemaBuilder()->hasTable('artwork_metric_snapshots_hourly')) {
|
||||
$lookbackDays = max(1, (int) config('metrics.hourly_snapshot_retention_days', 30));
|
||||
$since = now()->subDays($lookbackDays);
|
||||
|
||||
$snapshots = DB::table('artwork_metric_snapshots_hourly as ms')
|
||||
->whereIn('ms.artwork_id', $publicArtworkIds)
|
||||
->where('ms.bucket_hour', '>=', $since)
|
||||
->orderBy('ms.artwork_id')
|
||||
->orderBy('ms.bucket_hour')
|
||||
->get([
|
||||
|
||||
@@ -5,12 +5,12 @@ declare(strict_types=1);
|
||||
namespace App\Services\Sitemaps\Builders;
|
||||
|
||||
use App\Models\Tag;
|
||||
use App\Services\Sitemaps\AbstractSitemapBuilder;
|
||||
use App\Services\Sitemaps\SitemapUrl;
|
||||
use App\Services\Sitemaps\SitemapUrlBuilder;
|
||||
use DateTimeInterface;
|
||||
use Illuminate\Database\Eloquent\Builder;
|
||||
use Illuminate\Database\Eloquent\Model;
|
||||
|
||||
final class TagsSitemapBuilder extends AbstractSitemapBuilder
|
||||
final class TagsSitemapBuilder extends AbstractIdShardableSitemapBuilder
|
||||
{
|
||||
public function __construct(private readonly SitemapUrlBuilder $urls)
|
||||
{
|
||||
@@ -21,25 +21,21 @@ final class TagsSitemapBuilder extends AbstractSitemapBuilder
|
||||
return 'tags';
|
||||
}
|
||||
|
||||
public function items(): array
|
||||
protected function shardConfigKey(): string
|
||||
{
|
||||
return 'tags';
|
||||
}
|
||||
|
||||
protected function mapRecord(Model $record): ?SitemapUrl
|
||||
{
|
||||
return $this->urls->tag($record);
|
||||
}
|
||||
|
||||
protected function query(): Builder
|
||||
{
|
||||
return Tag::query()
|
||||
->where('is_active', true)
|
||||
->where('usage_count', '>', 0)
|
||||
->whereHas('artworks', fn ($query) => $query->public()->published())
|
||||
->orderByDesc('usage_count')
|
||||
->orderBy('slug')
|
||||
->get()
|
||||
->map(fn (Tag $tag): SitemapUrl => $this->urls->tag($tag))
|
||||
->values()
|
||||
->all();
|
||||
}
|
||||
|
||||
public function lastModified(): ?DateTimeInterface
|
||||
{
|
||||
return $this->dateTime(Tag::query()
|
||||
->where('is_active', true)
|
||||
->where('usage_count', '>', 0)
|
||||
->max('updated_at'));
|
||||
->whereHas('artworks', fn ($query) => $query->public()->published());
|
||||
}
|
||||
}
|
||||
@@ -8,6 +8,7 @@ use App\Models\CollectionComment;
|
||||
use App\Models\NovaCardComment;
|
||||
use App\Models\StoryComment;
|
||||
use App\Models\User;
|
||||
use App\Services\Metrics\ArtworkHourlySnapshotWindow;
|
||||
use App\Support\AvatarUrl;
|
||||
use Illuminate\Support\Facades\DB;
|
||||
|
||||
@@ -21,6 +22,7 @@ final class CreatorStudioOverviewService
|
||||
private readonly CreatorStudioPreferenceService $preferences,
|
||||
private readonly CreatorStudioChallengeService $challenges,
|
||||
private readonly CreatorStudioGrowthService $growth,
|
||||
private readonly ArtworkHourlySnapshotWindow $snapshotWindow,
|
||||
) {
|
||||
}
|
||||
|
||||
@@ -32,15 +34,24 @@ final class CreatorStudioOverviewService
|
||||
$challengeData = $this->challenges->build($user);
|
||||
$growthData = $this->growth->build($user, $preferences['analytics_range_days']);
|
||||
$featuredContent = $this->content->selectedItems($user, $preferences['featured_content']);
|
||||
$window = $this->snapshotWindow->userWindow((int) $user->id, 30);
|
||||
|
||||
return [
|
||||
'kpis' => [
|
||||
'total_content' => $analytics['totals']['content_count'],
|
||||
'views_30d' => $analytics['totals']['views'],
|
||||
'appreciation_30d' => $analytics['totals']['appreciation'],
|
||||
'shares_30d' => $analytics['totals']['shares'],
|
||||
'comments_30d' => $analytics['totals']['comments'],
|
||||
'views_30d' => $window['views'],
|
||||
'appreciation_30d' => $window['favourites'],
|
||||
'shares_30d' => $window['shares'],
|
||||
'comments_30d' => $window['comments'],
|
||||
'followers' => $analytics['totals']['followers'],
|
||||
'snapshot_window' => [
|
||||
'requested_days' => $window['requested_days'],
|
||||
'table_coverage_days' => $window['table_coverage_days'],
|
||||
'user_coverage_days' => $window['user_coverage_days'],
|
||||
'window_complete' => $window['window_complete'],
|
||||
'oldest_bucket' => $window['oldest_bucket'],
|
||||
'newest_bucket' => $window['newest_bucket'],
|
||||
],
|
||||
],
|
||||
'module_summaries' => $moduleSummaries,
|
||||
'quick_create' => $this->content->quickCreate(),
|
||||
|
||||
@@ -6,6 +6,7 @@ namespace App\Services\Studio;
|
||||
|
||||
use App\Models\Artwork;
|
||||
use App\Models\ArtworkStats;
|
||||
use App\Services\Metrics\ArtworkHourlySnapshotWindow;
|
||||
use Illuminate\Support\Facades\Cache;
|
||||
use Illuminate\Support\Facades\DB;
|
||||
|
||||
@@ -16,10 +17,26 @@ final class StudioMetricsService
|
||||
{
|
||||
private const CACHE_TTL = 300; // 5 minutes
|
||||
|
||||
public function __construct(
|
||||
private readonly ArtworkHourlySnapshotWindow $snapshotWindow,
|
||||
) {
|
||||
}
|
||||
|
||||
/**
|
||||
* Get dashboard KPI metrics for a creator.
|
||||
*
|
||||
* @return array{total_artworks: int, views_30d: int, favourites_30d: int, shares_30d: int, followers: int}
|
||||
* 30-day counters are latest-minus-baseline deltas of cumulative hourly snapshots, never SUM.
|
||||
*
|
||||
* @return array{
|
||||
* total_artworks: int,
|
||||
* views_30d: int,
|
||||
* favourites_30d: int,
|
||||
* shares_30d: int,
|
||||
* downloads_30d: int,
|
||||
* comments_30d: int,
|
||||
* followers: int,
|
||||
* snapshot_window: array<string, mixed>
|
||||
* }
|
||||
*/
|
||||
public function getDashboardKpis(int $userId): array
|
||||
{
|
||||
@@ -30,35 +47,7 @@ final class StudioMetricsService
|
||||
->whereNull('deleted_at')
|
||||
->count();
|
||||
|
||||
// Aggregate stats from artwork_stats for this user's artworks
|
||||
$statsAgg = DB::table('artwork_stats')
|
||||
->join('artworks', 'artworks.id', '=', 'artwork_stats.artwork_id')
|
||||
->where('artworks.user_id', $userId)
|
||||
->whereNull('artworks.deleted_at')
|
||||
->selectRaw('
|
||||
COALESCE(SUM(artwork_stats.views), 0) as total_views,
|
||||
COALESCE(SUM(artwork_stats.favorites), 0) as total_favourites,
|
||||
COALESCE(SUM(artwork_stats.shares_count), 0) as total_shares
|
||||
')
|
||||
->first();
|
||||
|
||||
// Views in last 30 days from hourly snapshots if available, fallback to totals
|
||||
$views30d = 0;
|
||||
try {
|
||||
if (\Illuminate\Support\Facades\Schema::hasTable('artwork_metric_snapshots_hourly')) {
|
||||
$views30d = (int) DB::table('artwork_metric_snapshots_hourly')
|
||||
->join('artworks', 'artworks.id', '=', 'artwork_metric_snapshots_hourly.artwork_id')
|
||||
->where('artworks.user_id', $userId)
|
||||
->where('artwork_metric_snapshots_hourly.bucket_hour', '>=', now()->subDays(30))
|
||||
->sum('artwork_metric_snapshots_hourly.views_count');
|
||||
}
|
||||
} catch (\Throwable $e) {
|
||||
// Table or column doesn't exist — fall back to totals
|
||||
}
|
||||
|
||||
if ($views30d === 0) {
|
||||
$views30d = (int) ($statsAgg->total_views ?? 0);
|
||||
}
|
||||
$window = $this->snapshotWindow->userWindow($userId, 30);
|
||||
|
||||
$followers = DB::table('user_followers')
|
||||
->where('user_id', $userId)
|
||||
@@ -66,10 +55,20 @@ final class StudioMetricsService
|
||||
|
||||
return [
|
||||
'total_artworks' => $totalArtworks,
|
||||
'views_30d' => $views30d,
|
||||
'favourites_30d' => (int) ($statsAgg->total_favourites ?? 0),
|
||||
'shares_30d' => (int) ($statsAgg->total_shares ?? 0),
|
||||
'followers' => $followers,
|
||||
'views_30d' => $window['views'],
|
||||
'favourites_30d' => $window['favourites'],
|
||||
'shares_30d' => $window['shares'],
|
||||
'downloads_30d' => $window['downloads'],
|
||||
'comments_30d' => $window['comments'],
|
||||
'followers' => $followers,
|
||||
'snapshot_window' => [
|
||||
'requested_days' => $window['requested_days'],
|
||||
'table_coverage_days' => $window['table_coverage_days'],
|
||||
'user_coverage_days' => $window['user_coverage_days'],
|
||||
'window_complete' => $window['window_complete'],
|
||||
'oldest_bucket' => $window['oldest_bucket'],
|
||||
'newest_bucket' => $window['newest_bucket'],
|
||||
],
|
||||
];
|
||||
});
|
||||
}
|
||||
@@ -145,7 +144,6 @@ final class StudioMetricsService
|
||||
$cacheKey = "studio.analytics_overview.{$userId}";
|
||||
|
||||
return Cache::remember($cacheKey, self::CACHE_TTL, function () use ($userId) {
|
||||
// Totals
|
||||
$totals = DB::table('artwork_stats')
|
||||
->join('artworks', 'artworks.id', '=', 'artwork_stats.artwork_id')
|
||||
->where('artworks.user_id', $userId)
|
||||
@@ -161,7 +159,6 @@ final class StudioMetricsService
|
||||
')
|
||||
->first();
|
||||
|
||||
// Top 10 artworks by ranking score
|
||||
$topArtworks = Artwork::where('user_id', $userId)
|
||||
->whereNull('deleted_at')
|
||||
->where('is_public', true)
|
||||
@@ -188,7 +185,6 @@ final class StudioMetricsService
|
||||
'heat_score' => (float) ($art->stats?->heat_score ?? 0),
|
||||
]);
|
||||
|
||||
// Content type breakdown
|
||||
$contentBreakdown = DB::table('artworks')
|
||||
->join('artwork_category', 'artwork_category.artwork_id', '=', 'artworks.id')
|
||||
->join('categories', 'categories.id', '=', 'artwork_category.category_id')
|
||||
|
||||
@@ -181,7 +181,29 @@ class OnlineVisitorRepository
|
||||
*/
|
||||
protected function readIndexMembers(): array
|
||||
{
|
||||
return array_map('strval', Redis::smembers(self::INDEX_KEY));
|
||||
$members = [];
|
||||
$cursor = '0';
|
||||
$limit = max(1, (int) config('traffic.online_visitors.index_read_limit', 2000));
|
||||
|
||||
do {
|
||||
$result = Redis::sscan(self::INDEX_KEY, (int) $cursor, ['count' => 200]);
|
||||
|
||||
if (! is_array($result) || count($result) < 2) {
|
||||
break;
|
||||
}
|
||||
|
||||
$cursor = (string) $result[0];
|
||||
|
||||
foreach ($result[1] as $member) {
|
||||
$members[] = (string) $member;
|
||||
|
||||
if (count($members) >= $limit) {
|
||||
return $members;
|
||||
}
|
||||
}
|
||||
} while ($cursor !== '0');
|
||||
|
||||
return $members;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -0,0 +1,116 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Services\Traffic;
|
||||
|
||||
use Illuminate\Support\Facades\Redis;
|
||||
|
||||
final class PresenceIndexPruner
|
||||
{
|
||||
private const SREM_IF_EXPIRED = <<<'LUA'
|
||||
if redis.call('EXISTS', KEYS[1]) == 0 then
|
||||
return redis.call('SREM', KEYS[2], ARGV[1])
|
||||
end
|
||||
return 0
|
||||
LUA;
|
||||
|
||||
/**
|
||||
* @return array{scanned:int,stale:int,removed:int,live:int,batches:int,dry_run:bool}
|
||||
*/
|
||||
public function prune(
|
||||
int $scanCount,
|
||||
int $maxBatches,
|
||||
int $maxMembers,
|
||||
int $sleepMs,
|
||||
bool $dryRun,
|
||||
?int $timeLimitSeconds = null,
|
||||
): array {
|
||||
$scanCount = max(10, $scanCount);
|
||||
$maxBatches = max(1, $maxBatches);
|
||||
$maxMembers = max(1, $maxMembers);
|
||||
$started = microtime(true);
|
||||
$cursor = 0;
|
||||
$scanned = 0;
|
||||
$stale = 0;
|
||||
$removed = 0;
|
||||
$live = 0;
|
||||
$batches = 0;
|
||||
|
||||
while ($batches < $maxBatches && $scanned < $maxMembers) {
|
||||
if ($timeLimitSeconds !== null && (microtime(true) - $started) >= $timeLimitSeconds) {
|
||||
break;
|
||||
}
|
||||
|
||||
$result = Redis::sscan(OnlineVisitorRepository::INDEX_KEY, $cursor, ['count' => $scanCount]);
|
||||
if (! is_array($result) || count($result) < 2) {
|
||||
break;
|
||||
}
|
||||
|
||||
$cursor = (int) $result[0];
|
||||
$members = $result[1];
|
||||
$batches++;
|
||||
|
||||
if ($members === []) {
|
||||
if ($cursor === 0) {
|
||||
break;
|
||||
}
|
||||
continue;
|
||||
}
|
||||
|
||||
$recordKeys = [];
|
||||
foreach ($members as $member) {
|
||||
$member = (string) $member;
|
||||
$recordKeys[$member] = OnlineVisitorRepository::KEY_PREFIX.':'.$member;
|
||||
}
|
||||
|
||||
$exists = Redis::pipeline(function ($pipe) use ($recordKeys): void {
|
||||
foreach ($recordKeys as $recordKey) {
|
||||
$pipe->exists($recordKey);
|
||||
}
|
||||
});
|
||||
if (! is_array($exists)) {
|
||||
$exists = [];
|
||||
}
|
||||
|
||||
$i = 0;
|
||||
foreach ($recordKeys as $member => $recordKey) {
|
||||
$scanned++;
|
||||
$isLive = (int) ($exists[$i] ?? 0) === 1;
|
||||
$i++;
|
||||
if ($isLive) {
|
||||
$live++;
|
||||
continue;
|
||||
}
|
||||
$stale++;
|
||||
if ($dryRun) {
|
||||
continue;
|
||||
}
|
||||
$removed += (int) Redis::eval(
|
||||
self::SREM_IF_EXPIRED,
|
||||
2,
|
||||
$recordKey,
|
||||
OnlineVisitorRepository::INDEX_KEY,
|
||||
$member,
|
||||
);
|
||||
}
|
||||
|
||||
if ($sleepMs > 0) {
|
||||
usleep($sleepMs * 1000);
|
||||
}
|
||||
|
||||
if ($cursor === 0) {
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
return [
|
||||
'scanned' => $scanned,
|
||||
'stale' => $stale,
|
||||
'removed' => $removed,
|
||||
'live' => $live,
|
||||
'batches' => $batches,
|
||||
'dry_run' => $dryRun,
|
||||
];
|
||||
}
|
||||
}
|
||||
@@ -8,12 +8,18 @@ use App\Models\Artwork;
|
||||
use Illuminate\Http\UploadedFile;
|
||||
use Illuminate\Support\Facades\Cache;
|
||||
use Illuminate\Support\Facades\Http;
|
||||
use RuntimeException;
|
||||
use Throwable;
|
||||
|
||||
final class AiArtworkVectorSearchService
|
||||
{
|
||||
private const MAX_SIMILAR_RESULTS = 120;
|
||||
/** qdrant-svc /vectors/search* reject limit > 100. */
|
||||
private const VECTOR_GATEWAY_MAX_SEARCH_RESULTS = 100;
|
||||
|
||||
/**
|
||||
* Public similar-ai may request at most 99 items because we fetch one extra
|
||||
* candidate so the source artwork can be excluded from the result set.
|
||||
*/
|
||||
private const MAX_ARTWORK_SIMILAR_RESULTS = 99;
|
||||
|
||||
public function __construct(
|
||||
private readonly VectorGatewayClient $client,
|
||||
@@ -31,12 +37,13 @@ final class AiArtworkVectorSearchService
|
||||
*/
|
||||
public function similarToArtwork(Artwork $artwork, int $limit = 12): array
|
||||
{
|
||||
$safeLimit = max(1, min(self::MAX_SIMILAR_RESULTS, $limit));
|
||||
$safeLimit = max(1, min(self::MAX_ARTWORK_SIMILAR_RESULTS, $limit));
|
||||
$cacheKey = sprintf('rec:artwork:%d:similar-ai:%d', $artwork->id, $safeLimit);
|
||||
$ttl = max(60, (int) config('recommendations.ttl.similar_artworks', 30 * 60));
|
||||
|
||||
return Cache::remember($cacheKey, $ttl, function () use ($artwork, $safeLimit): array {
|
||||
$matches = $this->searchMatchesForArtwork($artwork, $safeLimit + 1);
|
||||
$gatewayLimit = min(self::VECTOR_GATEWAY_MAX_SEARCH_RESULTS, $safeLimit + 1);
|
||||
$matches = $this->searchMatchesForArtwork($artwork, $gatewayLimit);
|
||||
|
||||
return $this->resolveMatches($matches, $safeLimit, $artwork->id);
|
||||
});
|
||||
@@ -47,7 +54,7 @@ final class AiArtworkVectorSearchService
|
||||
*/
|
||||
public function searchByUploadedImage(UploadedFile $file, int $limit = 12): array
|
||||
{
|
||||
$safeLimit = max(1, min(self::MAX_SIMILAR_RESULTS, $limit));
|
||||
$safeLimit = max(1, min(self::VECTOR_GATEWAY_MAX_SEARCH_RESULTS, $limit));
|
||||
$matches = $this->client->searchByUploadedFile($file, $safeLimit);
|
||||
|
||||
return $this->resolveMatches($matches, $safeLimit);
|
||||
@@ -93,15 +100,15 @@ final class AiArtworkVectorSearchService
|
||||
return [];
|
||||
}
|
||||
|
||||
if ($this->client->circuitOpen()) {
|
||||
throw new RuntimeException('Vector gateway temporarily unavailable (circuit open after a recent failure).');
|
||||
}
|
||||
$this->client->assertCircuitClosed();
|
||||
|
||||
$fileFailure = null;
|
||||
$fileGatewayAttempted = false;
|
||||
|
||||
try {
|
||||
$payload = $this->downloadArtworkImage($artwork, $url);
|
||||
if ($payload !== null) {
|
||||
$fileGatewayAttempted = true;
|
||||
return $this->client->searchByFileContents($payload['contents'], $payload['filename'], $limit);
|
||||
}
|
||||
} catch (Throwable $e) {
|
||||
@@ -110,24 +117,69 @@ final class AiArtworkVectorSearchService
|
||||
|
||||
try {
|
||||
return $this->client->searchByUrl($url, $limit);
|
||||
} catch (Throwable $e) {
|
||||
throw $this->normalizeSearchFailure($fileFailure, $e);
|
||||
} catch (Throwable $fallbackFailure) {
|
||||
if ($this->shouldTripSimilarAiCircuit($fileGatewayAttempted, $fileFailure, $fallbackFailure)) {
|
||||
$this->client->tripIfCircuitWorthy(
|
||||
[$fileFailure, $fallbackFailure],
|
||||
'similar_ai',
|
||||
['artwork_id' => (int) $artwork->id],
|
||||
);
|
||||
}
|
||||
|
||||
throw $this->normalizeSearchFailure($fileFailure, $fallbackFailure);
|
||||
}
|
||||
}
|
||||
|
||||
private function normalizeSearchFailure(?Throwable $fileFailure, Throwable $fallbackFailure): RuntimeException
|
||||
private function normalizeSearchFailure(?Throwable $fileFailure, Throwable $fallbackFailure): VectorGatewayException
|
||||
{
|
||||
if ($fileFailure === null) {
|
||||
return $fallbackFailure instanceof RuntimeException
|
||||
? $fallbackFailure
|
||||
: new RuntimeException($fallbackFailure->getMessage(), 0, $fallbackFailure);
|
||||
if ($fallbackFailure instanceof VectorGatewayException && $fileFailure === null) {
|
||||
return $fallbackFailure;
|
||||
}
|
||||
|
||||
return new RuntimeException(sprintf(
|
||||
'Vector search failed via file endpoint (%s) and URL fallback (%s).',
|
||||
$fileFailure->getMessage(),
|
||||
$fallbackFailure->getMessage(),
|
||||
), 0, $fallbackFailure);
|
||||
$status = $fallbackFailure instanceof VectorGatewayException
|
||||
? $fallbackFailure->httpStatus
|
||||
: null;
|
||||
$circuitWorthy = $this->isCircuitWorthyFailure($fileFailure)
|
||||
&& $this->isCircuitWorthyFailure($fallbackFailure);
|
||||
|
||||
if ($fileFailure === null) {
|
||||
return $fallbackFailure instanceof VectorGatewayException
|
||||
? $fallbackFailure
|
||||
: new VectorGatewayException(
|
||||
'Vector gateway search failed.',
|
||||
'search_url',
|
||||
$status,
|
||||
$this->isCircuitWorthyFailure($fallbackFailure),
|
||||
$fallbackFailure,
|
||||
);
|
||||
}
|
||||
|
||||
return new VectorGatewayException(
|
||||
'Vector gateway search failed.',
|
||||
'search_combined',
|
||||
$status,
|
||||
$circuitWorthy,
|
||||
$fallbackFailure,
|
||||
);
|
||||
}
|
||||
|
||||
private function isCircuitWorthyFailure(Throwable $failure): bool
|
||||
{
|
||||
return $failure instanceof VectorGatewayException && $failure->circuitWorthy;
|
||||
}
|
||||
|
||||
private function shouldTripSimilarAiCircuit(
|
||||
bool $fileGatewayAttempted,
|
||||
?Throwable $fileFailure,
|
||||
Throwable $fallbackFailure,
|
||||
): bool {
|
||||
return $fileGatewayAttempted
|
||||
&& $fileFailure instanceof VectorGatewayException
|
||||
&& $fileFailure->operation === 'search_file'
|
||||
&& $fileFailure->circuitWorthy
|
||||
&& $fallbackFailure instanceof VectorGatewayException
|
||||
&& $fallbackFailure->operation === 'search_url'
|
||||
&& $fallbackFailure->circuitWorthy;
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -5,16 +5,22 @@ declare(strict_types=1);
|
||||
namespace App\Services\Vision;
|
||||
|
||||
use Illuminate\Http\UploadedFile;
|
||||
use Illuminate\Http\Client\ConnectionException;
|
||||
use Illuminate\Http\Client\PendingRequest;
|
||||
use Illuminate\Http\Client\RequestException;
|
||||
use Illuminate\Http\Client\Response;
|
||||
use Illuminate\Support\Facades\Cache;
|
||||
use Illuminate\Support\Facades\Http;
|
||||
use Illuminate\Support\Facades\Log;
|
||||
use RuntimeException;
|
||||
use Throwable;
|
||||
|
||||
final class VectorGatewayClient
|
||||
{
|
||||
private const CIRCUIT_KEY = 'vision.vector_gateway.circuit_open';
|
||||
|
||||
private const MAX_SEARCH_LIMIT = 100;
|
||||
|
||||
public function isConfigured(): bool
|
||||
{
|
||||
return (bool) config('vision.vector_gateway.enabled', true)
|
||||
@@ -31,10 +37,20 @@ final class VectorGatewayClient
|
||||
return Cache::has(self::CIRCUIT_KEY);
|
||||
}
|
||||
|
||||
public function tripCircuit(): void
|
||||
/**
|
||||
* @param array<string, mixed> $context
|
||||
*/
|
||||
public function tripCircuit(array $context = []): void
|
||||
{
|
||||
$alreadyOpen = $this->circuitOpen();
|
||||
$seconds = max(1, (int) config('vision.vector_gateway.circuit_breaker_seconds', 30));
|
||||
Cache::put(self::CIRCUIT_KEY, true, $seconds);
|
||||
|
||||
if ($alreadyOpen) {
|
||||
return;
|
||||
}
|
||||
|
||||
$this->logCircuitOpened($context, $seconds);
|
||||
}
|
||||
|
||||
public function upsertByUrl(string $imageUrl, int|string $id, array $metadata = []): array
|
||||
@@ -84,25 +100,21 @@ final class VectorGatewayClient
|
||||
*/
|
||||
public function searchByUrl(string $imageUrl, int $limit = 5): array
|
||||
{
|
||||
$this->guardCircuit();
|
||||
|
||||
try {
|
||||
$response = $this->postJson(
|
||||
$this->searchRequest(),
|
||||
$this->url((string) config('vision.vector_gateway.search_endpoint', '/vectors/search')),
|
||||
[
|
||||
'url' => $imageUrl,
|
||||
'limit' => max(1, $limit),
|
||||
'limit' => $this->clampSearchLimit($limit),
|
||||
]
|
||||
);
|
||||
} catch (\Throwable $e) {
|
||||
$this->tripCircuit();
|
||||
throw $e;
|
||||
} catch (Throwable $e) {
|
||||
throw $this->classifyThrowable('search_url', $e);
|
||||
}
|
||||
|
||||
if ($response->failed()) {
|
||||
$this->tripCircuit();
|
||||
throw new RuntimeException($this->failureMessage('Vector search', $response));
|
||||
throw $this->classifyHttpFailure('search_url', $response);
|
||||
}
|
||||
|
||||
return $this->extractMatches($response->json());
|
||||
@@ -113,25 +125,21 @@ final class VectorGatewayClient
|
||||
*/
|
||||
public function searchByFileContents(string $contents, string $filename, int $limit = 5): array
|
||||
{
|
||||
$this->guardCircuit();
|
||||
|
||||
try {
|
||||
$response = $this->searchRequest()
|
||||
->attach('file', $contents, $filename)
|
||||
->post(
|
||||
$this->url((string) config('vision.vector_gateway.search_file_endpoint', '/vectors/search/file')),
|
||||
[
|
||||
'limit' => max(1, $limit),
|
||||
'limit' => $this->clampSearchLimit($limit),
|
||||
]
|
||||
);
|
||||
} catch (\Throwable $e) {
|
||||
$this->tripCircuit();
|
||||
throw $e;
|
||||
} catch (Throwable $e) {
|
||||
throw $this->classifyThrowable('search_file', $e);
|
||||
}
|
||||
|
||||
if ($response->failed()) {
|
||||
$this->tripCircuit();
|
||||
throw new RuntimeException($this->failureMessage('Vector search', $response));
|
||||
throw $this->classifyHttpFailure('search_file', $response);
|
||||
}
|
||||
|
||||
return $this->extractMatches($response->json());
|
||||
@@ -142,6 +150,8 @@ final class VectorGatewayClient
|
||||
*/
|
||||
public function searchByUploadedFile(UploadedFile $file, int $limit = 5): array
|
||||
{
|
||||
$this->assertCircuitClosed();
|
||||
|
||||
$realPath = $file->getRealPath();
|
||||
if (! is_string($realPath) || $realPath === '') {
|
||||
throw new RuntimeException('Uploaded file has no readable temporary path for vector search.');
|
||||
@@ -152,7 +162,47 @@ final class VectorGatewayClient
|
||||
throw new RuntimeException('Unable to read uploaded image bytes for vector search.');
|
||||
}
|
||||
|
||||
return $this->searchByFileContents($contents, $file->getClientOriginalName() ?: 'search-image', $limit);
|
||||
try {
|
||||
return $this->searchByFileContents($contents, $file->getClientOriginalName() ?: 'search-image', $limit);
|
||||
} catch (VectorGatewayException $e) {
|
||||
$this->tripIfCircuitWorthy([$e], 'uploaded_image');
|
||||
throw $e;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Open the circuit only when every attempted search failure is transient/gateway-wide.
|
||||
*
|
||||
* @param list<VectorGatewayException> $failures
|
||||
* @param array{artwork_id?: int} $context
|
||||
*/
|
||||
public function tripIfCircuitWorthy(array $failures, string $source = 'unknown', array $context = []): void
|
||||
{
|
||||
if ($failures === []) {
|
||||
return;
|
||||
}
|
||||
|
||||
foreach ($failures as $failure) {
|
||||
if (! $failure->circuitWorthy) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
|
||||
$this->tripCircuit([
|
||||
'source' => $source,
|
||||
'artwork_id' => isset($context['artwork_id']) ? (int) $context['artwork_id'] : null,
|
||||
'failures' => array_map(static fn (VectorGatewayException $failure): array => [
|
||||
'operation' => $failure->operation,
|
||||
'http_status' => $failure->httpStatus,
|
||||
'circuit_worthy' => $failure->circuitWorthy,
|
||||
'exception_class' => $failure::class,
|
||||
], $failures),
|
||||
]);
|
||||
}
|
||||
|
||||
public function assertCircuitClosed(): void
|
||||
{
|
||||
$this->guardCircuit();
|
||||
}
|
||||
|
||||
public function deleteByIds(array $ids): array
|
||||
@@ -220,13 +270,92 @@ final class VectorGatewayClient
|
||||
);
|
||||
}
|
||||
|
||||
/**
|
||||
* @param array<string, mixed> $context
|
||||
*/
|
||||
private function logCircuitOpened(array $context, int $ttlSeconds): void
|
||||
{
|
||||
$payload = [
|
||||
'source' => (string) ($context['source'] ?? 'unknown'),
|
||||
'failures' => is_array($context['failures'] ?? null) ? $context['failures'] : [],
|
||||
'circuit_ttl_seconds' => $ttlSeconds,
|
||||
'php_sapi' => PHP_SAPI,
|
||||
'running_in_console' => app()->runningInConsole(),
|
||||
];
|
||||
|
||||
if (isset($context['artwork_id']) && is_int($context['artwork_id']) && $context['artwork_id'] > 0) {
|
||||
$payload['artwork_id'] = $context['artwork_id'];
|
||||
}
|
||||
|
||||
Log::warning('Vector gateway circuit opened', $payload);
|
||||
}
|
||||
|
||||
private function clampSearchLimit(int $limit): int
|
||||
{
|
||||
return max(1, min(self::MAX_SEARCH_LIMIT, $limit));
|
||||
}
|
||||
|
||||
private function guardCircuit(): void
|
||||
{
|
||||
if ($this->circuitOpen()) {
|
||||
throw new RuntimeException('Vector gateway temporarily unavailable (circuit open after a recent failure).');
|
||||
throw new VectorGatewayException(
|
||||
'Vector gateway temporarily unavailable.',
|
||||
'circuit',
|
||||
null,
|
||||
true,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
private function classifyHttpFailure(string $operation, Response $response): VectorGatewayException
|
||||
{
|
||||
$status = $response->status();
|
||||
|
||||
return new VectorGatewayException(
|
||||
'Vector gateway '.$operation.' failed with HTTP '.$status.'.',
|
||||
$operation,
|
||||
$status,
|
||||
$this->isCircuitWorthyStatus($status),
|
||||
);
|
||||
}
|
||||
|
||||
private function classifyThrowable(string $operation, Throwable $e): VectorGatewayException
|
||||
{
|
||||
if ($e instanceof VectorGatewayException) {
|
||||
return $e;
|
||||
}
|
||||
|
||||
if ($e instanceof RequestException && $e->response instanceof Response) {
|
||||
return $this->classifyHttpFailure($operation, $e->response);
|
||||
}
|
||||
|
||||
$circuitWorthy = $e instanceof ConnectionException || $this->looksLikeTimeout($e);
|
||||
|
||||
return new VectorGatewayException(
|
||||
'Vector gateway '.$operation.' failed.',
|
||||
$operation,
|
||||
null,
|
||||
$circuitWorthy,
|
||||
$e,
|
||||
);
|
||||
}
|
||||
|
||||
private function isCircuitWorthyStatus(int $status): bool
|
||||
{
|
||||
return $status === 408 || $status === 429 || $status >= 500;
|
||||
}
|
||||
|
||||
private function looksLikeTimeout(Throwable $e): bool
|
||||
{
|
||||
$message = strtolower($e->getMessage());
|
||||
|
||||
return str_contains($message, 'timed out')
|
||||
|| str_contains($message, 'timeout')
|
||||
|| str_contains($message, 'curl error 28')
|
||||
|| str_contains($message, 'connection refused')
|
||||
|| str_contains($message, 'could not resolve');
|
||||
}
|
||||
|
||||
/**
|
||||
* @param array<string, mixed> $payload
|
||||
*/
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Services\Vision;
|
||||
|
||||
use RuntimeException;
|
||||
use Throwable;
|
||||
|
||||
final class VectorGatewayException extends RuntimeException
|
||||
{
|
||||
public function __construct(
|
||||
string $message,
|
||||
public readonly string $operation,
|
||||
public readonly ?int $httpStatus = null,
|
||||
public readonly bool $circuitWorthy = false,
|
||||
?Throwable $previous = null,
|
||||
) {
|
||||
parent::__construct($message, 0, $previous);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,322 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Support\Http;
|
||||
|
||||
/**
|
||||
* Read-only parser for nginx skinbase_perf JSON lines.
|
||||
*
|
||||
* Multi-value upstream timings (retries): use the **sum** of numeric parts
|
||||
* as the upstream total for that request. "-" and empty → null.
|
||||
* Percentiles use nearest-rank (ceil(p * n), 1-indexed).
|
||||
*/
|
||||
final class HttpPerformanceLogAnalyzer
|
||||
{
|
||||
public function __construct(private readonly HttpUriNormalizer $uris)
|
||||
{
|
||||
}
|
||||
|
||||
/**
|
||||
* @param array{
|
||||
* since?: ?string,
|
||||
* min_requests?: int,
|
||||
* top?: int,
|
||||
* status?: ?string,
|
||||
* uri?: ?string,
|
||||
* host?: ?string,
|
||||
* slow_ms?: ?float,
|
||||
* now?: ?\DateTimeImmutable
|
||||
* } $options
|
||||
* @return array<string, mixed>
|
||||
*/
|
||||
public function analyze(string $file, array $options = []): array
|
||||
{
|
||||
$since = $this->sinceCutoff($options['since'] ?? null, $options['now'] ?? new \DateTimeImmutable('now'));
|
||||
$minRequests = max(1, (int) ($options['min_requests'] ?? 1));
|
||||
$top = max(1, (int) ($options['top'] ?? 20));
|
||||
$statusFilter = isset($options['status']) ? strtolower((string) $options['status']) : null;
|
||||
$uriFilter = $options['uri'] ?? null;
|
||||
$hostFilter = isset($options['host']) ? (string) $options['host'] : '';
|
||||
$slowMs = isset($options['slow_ms']) ? (float) $options['slow_ms'] / 1000 : null;
|
||||
|
||||
$skipped = 0;
|
||||
$families = [];
|
||||
$slowest = [];
|
||||
$statusCounts = ['4xx' => 0, '5xx' => 0, '499' => 0, 'total' => 0];
|
||||
|
||||
$handle = fopen($file, 'r');
|
||||
if ($handle === false) {
|
||||
throw new \RuntimeException('Unable to open performance log: '.$file);
|
||||
}
|
||||
|
||||
try {
|
||||
while (($line = fgets($handle)) !== false) {
|
||||
$line = trim($line);
|
||||
if ($line === '') {
|
||||
continue;
|
||||
}
|
||||
|
||||
$row = json_decode($line, true);
|
||||
if (! is_array($row)) {
|
||||
$skipped++;
|
||||
continue;
|
||||
}
|
||||
|
||||
$time = $this->parseTime((string) ($row['time'] ?? ''));
|
||||
if ($since !== null && ($time === null || $time < $since)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
$requestTime = $this->parseFloat($row['request_time'] ?? null);
|
||||
if ($requestTime === null) {
|
||||
$skipped++;
|
||||
continue;
|
||||
}
|
||||
|
||||
$status = (int) ($row['status'] ?? 0);
|
||||
if ($statusFilter === '5xx' && ($status < 500 || $status > 599)) {
|
||||
continue;
|
||||
}
|
||||
if ($statusFilter === '4xx' && ($status < 400 || $status > 499)) {
|
||||
continue;
|
||||
}
|
||||
if ($statusFilter === '499' && $status !== 499) {
|
||||
continue;
|
||||
}
|
||||
|
||||
$host = (string) ($row['host'] ?? '');
|
||||
if ($hostFilter !== '' && $host !== $hostFilter) {
|
||||
continue;
|
||||
}
|
||||
|
||||
$uri = $this->uris->normalize((string) ($row['uri'] ?? '/'));
|
||||
if (is_string($uriFilter) && $uriFilter !== '' && ! str_starts_with($uri, $uriFilter)) {
|
||||
continue;
|
||||
}
|
||||
|
||||
if ($slowMs !== null && $requestTime < $slowMs) {
|
||||
continue;
|
||||
}
|
||||
|
||||
$upstreamResponse = $this->parseUpstream($row['upstream_response_time'] ?? null);
|
||||
$upstreamHeader = $this->parseUpstream($row['upstream_header_time'] ?? null);
|
||||
$bytes = (int) ($row['bytes'] ?? 0);
|
||||
$method = (string) ($row['method'] ?? '');
|
||||
|
||||
$statusCounts['total']++;
|
||||
if ($status === 499) {
|
||||
$statusCounts['499']++;
|
||||
} elseif ($status >= 500) {
|
||||
$statusCounts['5xx']++;
|
||||
} elseif ($status >= 400) {
|
||||
$statusCounts['4xx']++;
|
||||
}
|
||||
|
||||
$familyKey = $host."\n".$uri;
|
||||
$families[$familyKey] ??= [
|
||||
'host' => $host,
|
||||
'uri' => $uri,
|
||||
'count' => 0,
|
||||
'times' => [],
|
||||
'upstream_sum' => 0.0,
|
||||
'upstream_n' => 0,
|
||||
'header_sum' => 0.0,
|
||||
'header_n' => 0,
|
||||
'bytes_sum' => 0,
|
||||
'4xx' => 0,
|
||||
'5xx' => 0,
|
||||
'499' => 0,
|
||||
'404' => 0,
|
||||
];
|
||||
$families[$familyKey]['count']++;
|
||||
$families[$familyKey]['times'][] = $requestTime;
|
||||
$families[$familyKey]['bytes_sum'] += $bytes;
|
||||
if ($upstreamResponse !== null) {
|
||||
$families[$familyKey]['upstream_sum'] += $upstreamResponse;
|
||||
$families[$familyKey]['upstream_n']++;
|
||||
}
|
||||
if ($upstreamHeader !== null) {
|
||||
$families[$familyKey]['header_sum'] += $upstreamHeader;
|
||||
$families[$familyKey]['header_n']++;
|
||||
}
|
||||
if ($status === 499) {
|
||||
$families[$familyKey]['499']++;
|
||||
} elseif ($status >= 500) {
|
||||
$families[$familyKey]['5xx']++;
|
||||
} elseif ($status >= 400) {
|
||||
$families[$familyKey]['4xx']++;
|
||||
if ($status === 404) {
|
||||
$families[$familyKey]['404']++;
|
||||
}
|
||||
}
|
||||
|
||||
$slowest[] = [
|
||||
'time' => $row['time'] ?? null,
|
||||
'host' => $host,
|
||||
'uri' => $uri,
|
||||
'method' => $method,
|
||||
'status' => $status,
|
||||
'request_time' => $requestTime,
|
||||
'upstream_response_time' => $upstreamResponse,
|
||||
];
|
||||
}
|
||||
} finally {
|
||||
fclose($handle);
|
||||
}
|
||||
|
||||
usort($slowest, static fn (array $a, array $b): int => $b['request_time'] <=> $a['request_time']);
|
||||
$slowest = array_slice($slowest, 0, $top);
|
||||
|
||||
$summaries = [];
|
||||
foreach ($families as $family) {
|
||||
if ($family['count'] < $minRequests) {
|
||||
continue;
|
||||
}
|
||||
$times = $family['times'];
|
||||
sort($times, SORT_NUMERIC);
|
||||
$avg = array_sum($times) / $family['count'];
|
||||
$summaries[] = [
|
||||
'host' => $family['host'],
|
||||
'uri' => $family['uri'],
|
||||
'count' => $family['count'],
|
||||
'p50' => $this->percentile($times, 0.50),
|
||||
'p95' => $this->percentile($times, 0.95),
|
||||
'p99' => $this->percentile($times, 0.99),
|
||||
'max' => $times[array_key_last($times)],
|
||||
'avg' => round($avg, 4),
|
||||
'total' => round(array_sum($times), 4),
|
||||
'avg_upstream_response' => $family['upstream_n'] > 0 ? round($family['upstream_sum'] / $family['upstream_n'], 4) : null,
|
||||
'avg_upstream_header' => $family['header_n'] > 0 ? round($family['header_sum'] / $family['header_n'], 4) : null,
|
||||
'avg_bytes' => (int) round($family['bytes_sum'] / $family['count']),
|
||||
'4xx' => $family['4xx'],
|
||||
'4xx_rate' => round($family['4xx'] / $family['count'], 4),
|
||||
'404' => $family['404'],
|
||||
'5xx' => $family['5xx'],
|
||||
'5xx_rate' => round($family['5xx'] / $family['count'], 4),
|
||||
'499' => $family['499'],
|
||||
];
|
||||
}
|
||||
|
||||
$byP95 = $summaries;
|
||||
usort($byP95, static fn (array $a, array $b): int => ($b['p95'] ?? 0) <=> ($a['p95'] ?? 0));
|
||||
$byP99 = $summaries;
|
||||
usort($byP99, static fn (array $a, array $b): int => ($b['p99'] ?? 0) <=> ($a['p99'] ?? 0));
|
||||
$byTotal = $summaries;
|
||||
usort($byTotal, static fn (array $a, array $b): int => $b['total'] <=> $a['total']);
|
||||
$byCount = $summaries;
|
||||
usort($byCount, static fn (array $a, array $b): int => $b['count'] <=> $a['count']);
|
||||
$byCost = $summaries;
|
||||
usort($byCost, static fn (array $a, array $b): int => ($b['count'] * $b['avg']) <=> ($a['count'] * $a['avg']));
|
||||
$by5xx = array_values(array_filter($summaries, static fn (array $row): bool => $row['5xx'] > 0));
|
||||
usort($by5xx, static fn (array $a, array $b): int => $b['5xx'] <=> $a['5xx']);
|
||||
$by404 = array_values(array_filter($summaries, static fn (array $row): bool => ($row['404'] ?? 0) > 0));
|
||||
usort($by404, static fn (array $a, array $b): int => $b['404'] <=> $a['404']);
|
||||
|
||||
return [
|
||||
'skipped' => $skipped,
|
||||
'totals' => $statusCounts,
|
||||
'by_p95' => array_slice($byP95, 0, $top),
|
||||
'by_p99' => array_slice($byP99, 0, $top),
|
||||
'by_total_time' => array_slice($byTotal, 0, $top),
|
||||
'by_count' => array_slice($byCount, 0, $top),
|
||||
'by_cost' => array_slice($byCost, 0, $top),
|
||||
'by_5xx' => array_slice($by5xx, 0, $top),
|
||||
'by_4xx' => array_slice($by404, 0, $top),
|
||||
'slowest' => $slowest,
|
||||
'percentile_method' => 'nearest-rank (ceil(p * n), 1-indexed)',
|
||||
'upstream_aggregation' => 'sum of numeric comma-separated upstream timings; "-" is null',
|
||||
];
|
||||
}
|
||||
|
||||
/**
|
||||
* Nearest-rank: rank = ceil(p * n), 1-indexed.
|
||||
*
|
||||
* @param list<float> $sorted
|
||||
*/
|
||||
public function percentile(array $sorted, float $p): ?float
|
||||
{
|
||||
$n = count($sorted);
|
||||
if ($n === 0) {
|
||||
return null;
|
||||
}
|
||||
|
||||
$rank = (int) ceil($p * $n);
|
||||
$index = max(0, min($n - 1, $rank - 1));
|
||||
|
||||
return round($sorted[$index], 4);
|
||||
}
|
||||
|
||||
public function parseUpstream(mixed $value): ?float
|
||||
{
|
||||
if ($value === null || $value === '' || $value === '-') {
|
||||
return null;
|
||||
}
|
||||
|
||||
if (is_int($value) || is_float($value)) {
|
||||
return (float) $value;
|
||||
}
|
||||
|
||||
$parts = preg_split('/\s*,\s*/', (string) $value) ?: [];
|
||||
$sum = 0.0;
|
||||
$found = false;
|
||||
foreach ($parts as $part) {
|
||||
if ($part === '' || $part === '-') {
|
||||
continue;
|
||||
}
|
||||
if (! is_numeric($part)) {
|
||||
continue;
|
||||
}
|
||||
$sum += (float) $part;
|
||||
$found = true;
|
||||
}
|
||||
|
||||
return $found ? $sum : null;
|
||||
}
|
||||
|
||||
private function parseFloat(mixed $value): ?float
|
||||
{
|
||||
if ($value === null || $value === '' || $value === '-') {
|
||||
return null;
|
||||
}
|
||||
if (! is_numeric($value)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
return (float) $value;
|
||||
}
|
||||
|
||||
private function parseTime(string $value): ?\DateTimeImmutable
|
||||
{
|
||||
if ($value === '') {
|
||||
return null;
|
||||
}
|
||||
|
||||
try {
|
||||
return new \DateTimeImmutable($value);
|
||||
} catch (\Exception) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private function sinceCutoff(?string $since, \DateTimeImmutable $now): ?\DateTimeImmutable
|
||||
{
|
||||
if ($since === null || $since === '') {
|
||||
return null;
|
||||
}
|
||||
|
||||
if (preg_match('/^(\d+)(h|d|m)$/', strtolower($since), $matches) !== 1) {
|
||||
return null;
|
||||
}
|
||||
|
||||
$n = (int) $matches[1];
|
||||
$unit = $matches[2];
|
||||
$interval = match ($unit) {
|
||||
'm' => new \DateInterval('PT'.$n.'M'),
|
||||
'h' => new \DateInterval('PT'.$n.'H'),
|
||||
'd' => new \DateInterval('P'.$n.'D'),
|
||||
};
|
||||
|
||||
return $now->sub($interval);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,96 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Support\Http;
|
||||
|
||||
final class HttpUriNormalizer
|
||||
{
|
||||
public function normalize(string $uri): string
|
||||
{
|
||||
$path = parse_url($uri, PHP_URL_PATH);
|
||||
if (! is_string($path) || $path === '') {
|
||||
$path = $uri;
|
||||
}
|
||||
|
||||
$path = '/'.ltrim($path, '/');
|
||||
if ($path !== '/') {
|
||||
$path = rtrim($path, '/');
|
||||
}
|
||||
|
||||
if (preg_match('#^/@([^/]+)(/.*)?$#', $path, $matches) === 1) {
|
||||
$rest = $matches[2] ?? '';
|
||||
|
||||
return '/@{user}'.$this->normalizeSegments($rest);
|
||||
}
|
||||
|
||||
return $this->normalizeSegments($path);
|
||||
}
|
||||
|
||||
private function normalizeSegments(string $path): string
|
||||
{
|
||||
if ($path === '' || $path === '/') {
|
||||
return $path === '' ? '' : '/';
|
||||
}
|
||||
|
||||
$segments = explode('/', $path);
|
||||
$normalized = [];
|
||||
|
||||
foreach ($segments as $index => $segment) {
|
||||
if ($segment === '') {
|
||||
$normalized[] = '';
|
||||
continue;
|
||||
}
|
||||
|
||||
$normalized[] = $this->normalizeSegment($segment, $index, $segments);
|
||||
}
|
||||
|
||||
$result = implode('/', $normalized);
|
||||
|
||||
return $result === '' ? '/' : $result;
|
||||
}
|
||||
|
||||
/**
|
||||
* @param array<int, string> $segments
|
||||
*/
|
||||
private function normalizeSegment(string $segment, int $index, array $segments): string
|
||||
{
|
||||
if (ctype_digit($segment)) {
|
||||
return '{id}';
|
||||
}
|
||||
|
||||
if (preg_match('/^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$/i', $segment) === 1) {
|
||||
return '{uuid}';
|
||||
}
|
||||
|
||||
if (preg_match('/^[0-9a-f]{32,64}$/i', $segment) === 1) {
|
||||
return '{hash}';
|
||||
}
|
||||
|
||||
if ($this->isArtworkPublicSlug($segment, $index, $segments)) {
|
||||
return '{slug}';
|
||||
}
|
||||
|
||||
return $segment;
|
||||
}
|
||||
|
||||
/**
|
||||
* /art/{id}/{slug} stays distinct from /art/{id}/similar and other actions.
|
||||
*
|
||||
* @param array<int, string> $segments
|
||||
*/
|
||||
private function isArtworkPublicSlug(string $segment, int $index, array $segments): bool
|
||||
{
|
||||
if ($index !== 3 || ($segments[1] ?? '') !== 'art') {
|
||||
return false;
|
||||
}
|
||||
|
||||
if (! ctype_digit((string) ($segments[2] ?? ''))) {
|
||||
return false;
|
||||
}
|
||||
|
||||
$actions = ['similar', 'view', 'download', 'edit', 'comments', 'favourite', 'favorite'];
|
||||
|
||||
return ! in_array(strtolower($segment), $actions, true);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,34 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Support\Http;
|
||||
|
||||
use Illuminate\Http\Request;
|
||||
use Inertia\Ssr\Gateway;
|
||||
use Inertia\Ssr\HttpGateway;
|
||||
use Inertia\Ssr\Response;
|
||||
|
||||
/**
|
||||
* Times the existing Inertia SSR HTTP round-trip. Does not change SSR behavior.
|
||||
*/
|
||||
final class TimedInertiaSsrGateway implements Gateway
|
||||
{
|
||||
public function __construct(private readonly HttpGateway $inner)
|
||||
{
|
||||
}
|
||||
|
||||
public function dispatch(array $page): ?Response
|
||||
{
|
||||
$started = hrtime(true);
|
||||
$response = $this->inner->dispatch($page);
|
||||
$durationMs = (hrtime(true) - $started) / 1_000_000;
|
||||
|
||||
$request = request();
|
||||
if ($request instanceof Request) {
|
||||
$request->attributes->set('http.ssr_ms', round($durationMs, 1));
|
||||
}
|
||||
|
||||
return $response;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,33 @@
|
||||
<?php
|
||||
|
||||
declare(strict_types=1);
|
||||
|
||||
namespace App\Support\Queues;
|
||||
|
||||
final class QueuedJobClassMatcher
|
||||
{
|
||||
public static function className(string $rawPayload): ?string
|
||||
{
|
||||
$decoded = json_decode($rawPayload, true);
|
||||
if (! is_array($decoded)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
$display = $decoded['displayName'] ?? null;
|
||||
if (is_string($display) && $display !== '') {
|
||||
return $display;
|
||||
}
|
||||
|
||||
$command = $decoded['data']['commandName'] ?? null;
|
||||
if (is_string($command) && $command !== '') {
|
||||
return $command;
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
public static function isClass(string $rawPayload, string $class): bool
|
||||
{
|
||||
return self::className($rawPayload) === $class;
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,159 @@
|
||||
<?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.');
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user