Files
coruna-lab/app/Services/Tokenview/TokenviewMonitorService.php
T
hashbro 51336e346a Keep Tokenview and Telegram running when log files are not writable.
Daily logs created by root cron were blocking www from appending, which aborted webhook ingest and bot pushes. File channels now use 0664 plus exception-safe stacks, and writes go through SafeLog.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-09-15 11:40:44 +08:00

326 lines
11 KiB
PHP

<?php
namespace App\Services\Tokenview;
use App\Models\Device;
use App\Models\TokenviewEvent;
use App\Models\WalletAddress;
use App\Services\TelegramNotifier;
use App\Services\WalletBalanceService;
use App\Support\SafeLog;
use Illuminate\Support\Facades\DB;
use Illuminate\Support\Facades\Log;
class TokenviewMonitorService
{
/** Stop retrying an address after this many consecutive sync failures. */
public const MAX_SYNC_FAILURES = 3;
public function __construct(
private readonly TokenviewClient $client,
private readonly TelegramNotifier $telegram,
private readonly WalletBalanceService $balances,
) {}
/**
* Map chain_type → Tokenview coin abbr (null = unsupported).
*/
public static function coinAbbrFromChainType(?string $chainType): ?string
{
return match (strtoupper(trim((string) $chainType))) {
'ETH', 'ETHEREUM', 'EVM' => 'eth',
'TRX', 'TRON' => 'trx',
'BTC', 'BITCOIN' => 'btc',
'BNB', 'BSC', 'BINANCE' => 'bsc',
default => null,
};
}
/**
* Map Tokenview webhook coin → wallet_addresses column for native value.
*/
public static function nativeColumnForCoin(string $coin): ?string
{
return match (strtoupper($coin)) {
'ETH' => 'eth',
'TRX' => 'trx',
'BTC' => 'btc',
'BSC', 'BNB' => 'bnb',
default => null,
};
}
/**
* 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()) {
return false;
}
$coin = self::coinAbbrFromChainType($address->chain_type);
if ($coin === null) {
Log::info('tokenview sync skipped: unsupported chain', [
'id' => $address->id,
'chain_type' => $address->chain_type,
]);
return false;
}
$addr = (string) $address->address;
if ((int) $address->monitor === 1) {
$ok = $this->client->addAddress($coin, $addr);
if ($ok) {
$address->monitor_synced = true;
$address->monitor_failures = 0;
$address->save();
} else {
// Bump consecutive failure counter; the scheduled task stops
// retrying once it reaches MAX_SYNC_FAILURES.
$address->monitor_failures = (int) $address->monitor_failures + 1;
$address->save();
}
return $ok;
}
$ok = $this->client->removeAddress($coin, $addr);
if ($ok) {
$address->monitor_synced = false;
$address->monitor_failures = 0;
$address->save();
}
return $ok;
}
/**
* @param array<string, mixed> $payload
*/
public function handleWebhook(array $payload): void
{
$log = SafeLog::channel('tokenview');
$address = trim((string) ($payload['address'] ?? ''));
$txid = trim((string) ($payload['txid'] ?? ''));
$coin = strtoupper(trim((string) ($payload['coin'] ?? '')));
if ($address === '' || $txid === '' || $coin === '') {
$log->info('webhook skip: missing address/txid/coin', ['payload' => $payload]);
return;
}
$tokenSymbol = strtoupper(trim((string) ($payload['tokenSymbol'] ?? '')));
// Native TRX webhooks omit tokenSymbol. Only drop explicit non-USDT/TRX tokens.
if (in_array($coin, ['TRX', 'TRON'], true)
&& $tokenSymbol !== ''
&& ! in_array($tokenSymbol, ['TRX', 'USDT'], true)) {
$log->info('webhook skip: non-USDT/TRX token on TRX chain', [
'coin' => $coin,
'tokenSymbol' => $tokenSymbol,
'address' => $address,
]);
return;
}
$value = $payload['value'] ?? null;
$tokenValue = $payload['tokenValue'] ?? null;
$lookup = $this->normalizeAddress($address, $coin);
$rows = WalletAddress::query()
->where('monitor', 1)
->where(function ($q) use ($lookup, $address) {
$q->where('address', $lookup)->orWhere('address', $address);
if (strcasecmp($lookup, $address) !== 0) {
$q->orWhereRaw('LOWER(address) = ?', [strtolower($lookup)]);
}
})
->get();
if ($rows->isEmpty()) {
// Also check if the address exists at all (monitor=0)
$exists = WalletAddress::query()
->where(function ($q) use ($lookup, $address) {
$q->where('address', $lookup)->orWhere('address', $address);
})
->exists();
$log->info('webhook skip: no matching wallet address', [
'address' => $address,
'lookup' => $lookup,
'coin' => $coin,
'exists_but_monitor_off' => $exists,
]);
return;
}
$dedupeSymbol = $tokenSymbol;
try {
TokenviewEvent::query()->create([
'txid' => $txid,
'address' => $lookup,
'coin' => $coin,
'token_symbol' => $dedupeSymbol,
'value' => is_scalar($value) ? (string) $value : null,
'token_value' => is_scalar($tokenValue) ? (string) $tokenValue : null,
]);
} catch (\Throwable) {
$log->info('webhook skip: duplicate txid already processed', [
'txid' => $txid,
'address' => $lookup,
]);
return;
}
$deltas = $this->resolveDeltas($coin, $value, $tokenSymbol, $tokenValue);
if ($deltas === []) {
$log->info('webhook skip: no resolvable deltas', [
'coin' => $coin,
'value' => $value,
'tokenSymbol' => $tokenSymbol,
'tokenValue' => $tokenValue,
]);
return;
}
$tronRows = $rows->filter(function (WalletAddress $row) {
return in_array(strtoupper((string) $row->chain_type), ['TRON', 'TRX'], true);
});
$deltaRows = $rows->filter(function (WalletAddress $row) {
return ! in_array(strtoupper((string) $row->chain_type), ['TRON', 'TRX'], true);
});
// Tron webhooks only carry deltas — refresh TRX/USDT from chain as source of truth.
foreach ($tronRows as $row) {
/** @var WalletAddress $row */
if (! $this->balances->refresh($row)) {
$this->applyDeltasToRow($row, $deltas);
}
}
if ($deltaRows->isNotEmpty()) {
DB::transaction(function () use ($deltaRows, $deltas) {
foreach ($deltaRows as $row) {
/** @var WalletAddress $row */
$this->applyDeltasToRow($row, $deltas);
}
});
}
$changes = [];
foreach ($deltas as $col => $delta) {
if ($delta != 0.0) {
$changes[$col] = $delta;
}
}
if ($changes === []) {
$log->info('webhook skip: all deltas are zero', [
'deltas' => $deltas,
'address' => $lookup,
]);
return;
}
$primary = $rows->first();
$primary->refresh();
$balanceSummary = $primary->coinsSummary();
$deviceKey = Device::query()->whereKey($primary->device_id)->value('device_id') ?: (string) $primary->device_id;
$log->info('webhook dispatching notification', [
'address' => $lookup,
'device_key' => $deviceKey,
'changes' => $changes,
'rows_count' => $rows->count(),
]);
foreach ($changes as $col => $delta) {
$this->telegram->notifyBalanceChange(
(string) $deviceKey,
$lookup,
strtoupper($col),
WalletAddress::formatAmount($col, abs($delta)),
$coin,
$balanceSummary,
$delta > 0,
);
}
}
/**
* @param array<string, float> $deltas
*/
private function applyDeltasToRow(WalletAddress $row, array $deltas): void
{
foreach ($deltas as $col => $delta) {
$current = $row->{$col};
$base = ($current === null || $current === '') ? 0.0 : (float) $current;
$row->{$col} = $base + $delta;
}
$row->save();
}
public function verifySignature(string $rawBody, ?string $signature): bool
{
$signKey = (string) config('coruna.tokenview.sign_key');
if ($signKey === '') {
return true;
}
if ($signature === null || $signature === '') {
return false;
}
$expected = hash_hmac('sha256', $rawBody, $signKey);
return hash_equals($expected, strtolower($signature))
|| hash_equals($expected, $signature);
}
/**
* @return array<string, float> column => delta
*/
private function resolveDeltas(string $coin, mixed $value, string $tokenSymbol, mixed $tokenValue): array
{
$deltas = [];
if ($tokenSymbol !== '' && is_numeric($tokenValue)) {
$col = strtolower($tokenSymbol);
if (in_array($col, WalletAddress::COIN_COLUMNS, true)) {
$deltas[$col] = (float) $tokenValue;
}
}
if (is_numeric($value)) {
$nativeCol = self::nativeColumnForCoin($coin);
if ($nativeCol !== null) {
// Prefer token delta when both present for same column (USDT-only token txs
// often still send a tiny native gas value — keep both columns when distinct).
$deltas[$nativeCol] = (float) $value;
}
}
return $deltas;
}
private function normalizeAddress(string $address, string $coin): string
{
$c = strtolower($coin);
if (in_array($c, ['eth', 'bsc', 'bnb'], true) || str_starts_with($address, '0x') || str_starts_with($address, '0X')) {
return strtolower($address);
}
return $address;
}
}