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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions resources/views/dashboard.blade.php
Original file line number Diff line number Diff line change
Expand Up @@ -38,4 +38,6 @@

<livewire:pulse.queues cols="6" />

<livewire:queue-monitor lazy />

</x-pulse>
28 changes: 28 additions & 0 deletions src/Livewire/QueueMonitor.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
<?php

namespace Laravel\Pulse\Livewire;

use Laravel\Pulse\Livewire\Concerns\HasPeriod;
use Laravel\Pulse\Livewire\Concerns\RemembersQueries;
use Laravel\Pulse\Livewire\Concerns\ShouldNotReportUsage;
use Livewire\Component;

class QueueMonitor extends Component
{
use HasPeriod, RemembersQueries, ShouldNotReportUsage;

/**
* The queue to monitor.
*/
public string $queue = 'default';

/**
* The connection.
*/
public ?string $connection = null;

public function render(callable $query)
{
[$readings, $time, $runAt] = dd($this->remember(fn ($interval) => $query($interval, $this->queue, $this->connection)));
}
}
2 changes: 2 additions & 0 deletions src/PulseServiceProvider.php
Original file line number Diff line number Diff line change
Expand Up @@ -68,6 +68,7 @@ public function register(): void
Livewire\Exceptions::class => [Queries\Exceptions::class],
Livewire\SlowRoutes::class => [Queries\SlowRoutes::class],
Livewire\SlowQueries::class => [Queries\SlowQueries::class],
Livewire\QueueMonitor::class => [Queries\QueueMonitor::class],
Livewire\SlowOutgoingRequests::class => [Queries\SlowOutgoingRequests::class],
Livewire\Cache::class => [Queries\CacheInteractions::class, Queries\MonitoredCacheInteractions::class],
] as $card => $queries) {
Expand Down Expand Up @@ -224,6 +225,7 @@ protected function registerComponents(): void
$livewire->component('pulse.queues', Livewire\Queues::class);
$livewire->component('pulse.servers', Livewire\Servers::class);
$livewire->component('pulse.slow-jobs', Livewire\SlowJobs::class);
$livewire->component('queue-monitor', Livewire\QueueMonitor::class);
$livewire->component('pulse.exceptions', Livewire\Exceptions::class);
$livewire->component('pulse.slow-routes', Livewire\SlowRoutes::class);
$livewire->component('pulse.slow-queries', Livewire\SlowQueries::class);
Expand Down
103 changes: 103 additions & 0 deletions src/Queries/QueueMonitor.php
Original file line number Diff line number Diff line change
@@ -0,0 +1,103 @@
<?php

namespace Laravel\Pulse\Queries;

use Carbon\CarbonImmutable;
use Carbon\CarbonInterval as Interval;
use Illuminate\Config\Repository;
use Illuminate\Database\DatabaseManager;

class QueueMonitor
{
use Concerns\InteractsWithConnection;

/**
* Create a new query instance.
*/
public function __construct(
protected Repository $config,
protected DatabaseManager $db,
) {
//
}

/**
* Run the query.
*/
public function __invoke(Interval $interval, string $queue, string $connection = null)
{
$now = new CarbonImmutable();

$connection ??= $this->config->get('queue.default');

$maxDataPoints = 60;

$currentBucket = CarbonImmutable::createFromTimestamp(
floor($now->getTimestamp() / ($interval->totalSeconds / $maxDataPoints)) * ($interval->totalSeconds / $maxDataPoints)
);

$secondsPerPeriod = (int) ($interval->totalSeconds / $maxDataPoints);

$padding = collect([])
->pad(60, null)
->map(fn ($value, $i) => (object) [
'date' => $currentBucket->subSeconds($i * $secondsPerPeriod)->format('Y-m-d H:i'),
'size' => null,
'failed' => null,
])
->reverse()
->keyBy('date');

return $this->connection()
->query()
->selectRaw('bucket, MAX(`size`) AS `size`, MAX(`failed`) AS `failed`')
->fromSub(
fn ($query) => $query
->from('pulse_queue_sizes')
->selectRaw('size, failed, date, FLOOR(UNIX_TIMESTAMP(CONVERT_TZ(`date`, ?, @@session.time_zone)) / ?) AS `bucket`', [$now->format('P'), $secondsPerPeriod])
->where('date', '>=', $now->ceilSeconds($interval->totalSeconds / $maxDataPoints)->subSeconds((int) $interval->totalSeconds))
->where('queue', $queue)
->when($connection, fn ($query) => $query->where('connection', $connection)),
'grouped'
)
->groupBy('bucket')
->orderByDesc('bucket')
->get()
->reverse()
->pipe(function ($readings) use ($secondsPerPeriod, $padding) {
$readings = $readings->keyBy(fn ($reading) => CarbonImmutable::createFromTimestamp($reading->bucket * $secondsPerPeriod)->format('Y-m-d H:i'));

$readings = $padding->merge($readings)->values();

$previousFailed = 0;

return $readings->each(function ($reading) use (&$previousFailed) {
if ($reading->failed === null) {
$previousFailed = 0;

return;
}

if ($previousFailed === 0 || $previousFailed > $reading->failed) {
$previousFailed = $reading->failed;

$reading->failed = null;

return;
}

with($reading->failed, function ($newPreviousFailed) use ($reading, &$previousFailed) {
$reading->failed = $reading->failed - $previousFailed;

$previousFailed = $newPreviousFailed;
});
});
})
->pipe(fn ($readings) => (object) [
'readings' => $readings->map(fn ($reading) => (object) [
'size' => $reading->size,
'failed' => $reading->failed,
])->all(),
]);
}
}