Feat: Batching up poll shipment jobs
This commit is contained in:
@@ -3,40 +3,43 @@
|
||||
namespace Modules\Core\Shipping\Jobs;
|
||||
|
||||
use Illuminate\Bus\Queueable;
|
||||
use Illuminate\Contracts\Queue\ShouldBeUnique;
|
||||
use Illuminate\Contracts\Queue\ShouldQueue;
|
||||
use Illuminate\Foundation\Bus\Dispatchable;
|
||||
use Illuminate\Queue\InteractsWithQueue;
|
||||
use Illuminate\Queue\SerializesModels;
|
||||
use Illuminate\Support\Facades\Bus;
|
||||
use Lunar\Shipping\Facades\Shipping;
|
||||
use Modules\Core\Shipping\Contracts\CarrierFulfillmentInterface;
|
||||
use Modules\Core\Shipping\Contracts\SupportsTracking;
|
||||
use Modules\Core\Shipping\Enums\TrackingStatus;
|
||||
use Modules\Core\Shipping\Models\Shipment;
|
||||
use Modules\Core\Shipping\Services\ShipmentTrackingRecorder;
|
||||
use Throwable;
|
||||
|
||||
/**
|
||||
* Carrier-agnostic: polls every Shipment not yet in a terminal state,
|
||||
* skipping carriers whose fulfillment service doesn't implement
|
||||
* SupportsTracking. New checkpoints are recorded in shipment_info and
|
||||
* dispatch ShipmentStatusUpdatedByCarrier — one event per new checkpoint.
|
||||
* Carrier-agnostic: finds every Shipment not yet in a terminal state on a
|
||||
* carrier whose fulfillment service implements SupportsTracking, and
|
||||
* dispatches one RefreshShipmentTrackingJob per shipment as a batch. The
|
||||
* carrier calls happen in those jobs — each with its own timeout, and a
|
||||
* failure in one (a carrier 500, a malformed parcel response) doesn't
|
||||
* touch the others (allowFailures).
|
||||
*
|
||||
* Each shipment's trackShipment() call is individually try/caught in
|
||||
* pollCarrierShipments() — one shipment's tracking lookup failing (a
|
||||
* carrier 500, a malformed parcel response) must not stop the rest of that
|
||||
* carrier's shipments in the same batch from being polled. The failure is
|
||||
* reported and the loop continues.
|
||||
* Unique: a scheduled run is dropped while the previous one is still
|
||||
* queued or building its batch. The schedule's withoutOverlapping() alone
|
||||
* can't do this — for a queued job it only covers the dispatch itself.
|
||||
*/
|
||||
class PollShipmentTrackingJob implements ShouldQueue
|
||||
class PollShipmentTrackingJob implements ShouldBeUnique, ShouldQueue
|
||||
{
|
||||
use Dispatchable;
|
||||
use InteractsWithQueue;
|
||||
use Queueable;
|
||||
use SerializesModels;
|
||||
|
||||
public int $tries = 3;
|
||||
/** Only queries and dispatches — no carrier calls. */
|
||||
public int $timeout = 60;
|
||||
|
||||
public int $backoff = 60;
|
||||
public bool $failOnTimeout = true;
|
||||
|
||||
public int $tries = 1;
|
||||
|
||||
public int $uniqueFor = 300;
|
||||
|
||||
public function handle(): void
|
||||
{
|
||||
@@ -48,7 +51,7 @@ class PollShipmentTrackingJob implements ShouldQueue
|
||||
return;
|
||||
}
|
||||
|
||||
Shipment::query()
|
||||
$jobs = Shipment::query()
|
||||
->whereIn('carrier', $trackableCarriers)
|
||||
->whereNull('cancelled_at')
|
||||
// Pending vouchers (issued on print) have no number to track yet.
|
||||
@@ -60,26 +63,17 @@ class PollShipmentTrackingJob implements ShouldQueue
|
||||
TrackingStatus::Cancelled->value,
|
||||
]);
|
||||
})
|
||||
->chunkById(50, function ($shipments) {
|
||||
$shipments->groupBy('carrier')->each(
|
||||
fn ($group, $carrier) => $this->pollCarrierShipments($carrier, $group)
|
||||
);
|
||||
});
|
||||
}
|
||||
->pluck('id')
|
||||
->map(fn (int $id) => new RefreshShipmentTrackingJob($id));
|
||||
|
||||
private function pollCarrierShipments(string $carrier, $shipments): void
|
||||
{
|
||||
if (! $this->fulfillmentService($carrier) instanceof SupportsTracking) {
|
||||
if ($jobs->isEmpty()) {
|
||||
return;
|
||||
}
|
||||
|
||||
foreach ($shipments as $shipment) {
|
||||
try {
|
||||
app(ShipmentTrackingRecorder::class)->refresh($shipment);
|
||||
} catch (Throwable $e) {
|
||||
report($e);
|
||||
}
|
||||
}
|
||||
Bus::batch($jobs)
|
||||
->name('Poll shipment tracking')
|
||||
->allowFailures()
|
||||
->dispatch();
|
||||
}
|
||||
|
||||
private function fulfillmentService(string $carrier): ?CarrierFulfillmentInterface
|
||||
|
||||
@@ -0,0 +1,63 @@
|
||||
<?php
|
||||
|
||||
namespace Modules\Core\Shipping\Jobs;
|
||||
|
||||
use Illuminate\Bus\Batchable;
|
||||
use Illuminate\Bus\Queueable;
|
||||
use Illuminate\Contracts\Queue\ShouldBeUnique;
|
||||
use Illuminate\Contracts\Queue\ShouldQueue;
|
||||
use Illuminate\Foundation\Bus\Dispatchable;
|
||||
use Illuminate\Queue\InteractsWithQueue;
|
||||
use Modules\Core\Shipping\Models\Shipment;
|
||||
use Modules\Core\Shipping\Services\ShipmentTrackingRecorder;
|
||||
|
||||
/**
|
||||
* Refreshes one shipment's tracking from its carrier — one job per
|
||||
* shipment in PollShipmentTrackingJob's batch, so a slow or failing
|
||||
* carrier call only costs its own job, not the rest of the run.
|
||||
*
|
||||
* Unique per shipment: while one is queued or running, another for the
|
||||
* same shipment isn't queued, so two refreshes can't both record the same
|
||||
* checkpoint (the recorder's dedup reads, then writes).
|
||||
*/
|
||||
class RefreshShipmentTrackingJob implements ShouldBeUnique, ShouldQueue
|
||||
{
|
||||
use Batchable;
|
||||
use Dispatchable;
|
||||
use InteractsWithQueue;
|
||||
use Queueable;
|
||||
|
||||
/** A few carrier calls at a 10–15s HTTP timeout each; under the queue's 90s retry_after. */
|
||||
public int $timeout = 45;
|
||||
|
||||
public bool $failOnTimeout = true;
|
||||
|
||||
/** The next poll, five minutes later, is the retry. */
|
||||
public int $tries = 1;
|
||||
|
||||
/** Releases the lock if a worker dies mid-job and never releases it. */
|
||||
public int $uniqueFor = 300;
|
||||
|
||||
public function __construct(public readonly int $shipmentId) {}
|
||||
|
||||
public function uniqueId(): string
|
||||
{
|
||||
return (string) $this->shipmentId;
|
||||
}
|
||||
|
||||
public function handle(ShipmentTrackingRecorder $recorder): void
|
||||
{
|
||||
if ($this->batch()?->cancelled()) {
|
||||
return;
|
||||
}
|
||||
|
||||
$shipment = Shipment::find($this->shipmentId);
|
||||
|
||||
// Cancelled (or deleted) since the batch was built.
|
||||
if (! $shipment || $shipment->isCancelled()) {
|
||||
return;
|
||||
}
|
||||
|
||||
$recorder->refresh($shipment);
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user