numberSync = new SplitTicketNumberSyncService(); $this->ruleService = new SplitTicketRuleService(); $this->lockService = new SplitTicketSyncLockService(); } /** * 同步单条工单 * * @param bool $dueValidated 为 true 时跳过 shouldSkip(由 syncDueTickets 预筛后传入) * @param string $syncMode auto=定时自动 manual=手动触发 * @return array{success:bool,message:string,skipped?:bool} */ public function syncOne(int $ticketId, bool $force = false, bool $dueValidated = false, string $syncMode = 'manual'): array { $ticket = Ticket::get($ticketId); if (!$ticket) { SplitTicketSyncLogger::log('sync', 'ticket not found', ['ticketId' => $ticketId]); return ['success' => false, 'message' => '工单不存在']; } SplitTicketSyncLogger::setTicketContext($ticketId, (string) $ticket['ticket_type']); SplitTicketSyncLogger::log('sync', 'syncOne start', [ 'force' => $force, 'syncMode' => $syncMode, 'status' => (string) $ticket['status'], 'syncFailCount' => (int) ($ticket['sync_fail_count'] ?? 0), 'syncTime' => (int) ($ticket['sync_time'] ?? 0), 'pageUrl' => (string) $ticket['ticket_url'], 'nodeHost' => SplitSyncConfigService::getNodeHost(), ]); if (!$force && !$dueValidated) { $skip = $this->shouldSkip($ticket); if ($skip !== null) { SplitTicketSyncLogger::log('sync', 'skipped', ['reason' => $skip]); SplitTicketSyncLogger::clearTicketContext(); return ['success' => false, 'message' => $skip, 'skipped' => true]; } } if (!$this->lockService->acquire($ticketId)) { SplitTicketSyncLogger::log('sync', 'lock busy', ['ticketId' => $ticketId]); SplitTicketSyncLogger::clearTicketContext(); return ['success' => false, 'message' => '工单正在同步中', 'skipped' => true]; } try { $result = $this->doSync($ticket, $syncMode); SplitTicketSyncLogger::log('sync', 'syncOne end', $result); return $result; } finally { $this->lockService->release($ticketId); SplitTicketSyncLogger::clearTicketContext(); } } /** * 扫描到期工单并同步(限流 + Node 队列感知) */ public function syncDueTickets(): int { $maxPerRun = SplitSyncConfigService::getMaxTicketsPerCronRun(); $autoSyncTypes = $this->resolveAutoSyncTicketTypes(); $failThreshold = SplitSyncConfigService::getFailPauseThreshold(); if ($autoSyncTypes === []) { SplitTicketSyncLogger::logCron('scan skip', ['reason' => 'no auto-sync ticket types']); return 0; } if (!SplitNodeHealthService::canAcceptWork()) { SplitTicketSyncLogger::logCron('scan skip', [ 'reason' => 'node queue busy', 'queue' => SplitNodeHealthService::getQueueSnapshot(), ]); return 0; } $query = Ticket::where('status', 'normal')->whereIn('ticket_type', $autoSyncTypes); if ($failThreshold > 0) { $query->where('sync_fail_count', '<', $failThreshold); } $list = $query->order('sync_time', 'asc')->order('id', 'asc') ->limit($maxPerRun * 5) ->select(); $list = $this->sortTicketsBySessionKey($list); $candidateCount = count($list); SplitTicketSyncLogger::logCron('scan start', [ 'candidateCount' => $candidateCount, 'maxPerRun' => $maxPerRun, 'autoSyncTypes' => $autoSyncTypes, 'failThreshold' => $failThreshold, 'nodeQueue' => SplitNodeHealthService::getQueueSnapshot(), 'groupedBy' => 'sessionKey', ]); SplitTicketSyncLogger::log('cron', 'scan start', [ 'candidateCount' => $candidateCount, 'maxPerRun' => $maxPerRun, 'autoSyncTypes' => $autoSyncTypes, 'nodeQueue' => SplitNodeHealthService::getQueueSnapshot(), ]); /** @var list> $processed */ $processed = []; /** @var list> $skipped */ $skipped = []; $count = 0; $stoppedByNodeBusy = false; foreach ($list as $index => $ticket) { $ticketId = (int) $ticket['id']; $ticketUrl = (string) $ticket['ticket_url']; $ticketType = (string) $ticket['ticket_type']; $sessionKey = AntiBotConfigBuilder::resolveSessionKey($ticketType, $ticketUrl); if ($count >= $maxPerRun) { $reason = '本轮处理上限已满'; $skipped[] = $this->buildCronTicketEntry($ticketId, $ticketType, $ticketUrl, $reason); $this->logCronCandidateSkipped($ticketId, $ticketType, $ticketUrl, $reason); continue; } $skip = $this->shouldSkip($ticket); if ($skip !== null) { $skipped[] = $this->buildCronTicketEntry($ticketId, $ticketType, $ticketUrl, $skip); $this->logCronCandidateSkipped($ticketId, $ticketType, $ticketUrl, $skip); continue; } if (!SplitNodeHealthService::canAcceptWork()) { $reason = 'Node 队列繁忙,停止本轮后续同步'; $skipped[] = $this->buildCronTicketEntry($ticketId, $ticketType, $ticketUrl, $reason); $this->logCronCandidateSkipped($ticketId, $ticketType, $ticketUrl, $reason); $stoppedByNodeBusy = true; for ($i = $index + 1; $i < $candidateCount; $i++) { $remain = $list[$i]; $remainReason = 'Node 队列繁忙,本轮未执行'; $skipped[] = $this->buildCronTicketEntry( (int) $remain['id'], (string) $remain['ticket_type'], (string) $remain['ticket_url'], $remainReason ); } SplitTicketSyncLogger::log('cron', 'node busy, stop batch', [ 'processed' => $count, 'queue' => SplitNodeHealthService::getQueueSnapshot(), ]); break; } $result = $this->syncOne($ticketId, false, true, 'auto'); if (!empty($result['skipped'])) { $skipReason = (string) ($result['message'] ?? '已跳过'); $skipped[] = $this->buildCronTicketEntry($ticketId, $ticketType, $ticketUrl, $skipReason); $this->logCronCandidateSkipped($ticketId, $ticketType, $ticketUrl, $skipReason); continue; } $processed[] = [ 'ticketId' => $ticketId, 'ticketType' => $ticketType, 'ticketUrl' => $ticketUrl, 'sessionKey' => $sessionKey, 'success' => !empty($result['success']), 'message' => (string) ($result['message'] ?? ''), ]; $count++; } $candidateIds = array_map(static function ($row) { return (int) (is_array($row) ? ($row['id'] ?? 0) : ($row['id'] ?? 0)); }, $list); $openNotSynced = $this->buildOpenNotSyncedAudit( $autoSyncTypes, $failThreshold, $maxPerRun, $candidateIds, $processed, $skipped, $stoppedByNodeBusy ); $successCount = count(array_filter($processed, static function (array $row): bool { return !empty($row['success']); })); $failedCount = count($processed) - $successCount; SplitTicketSyncLogger::logCron('scan end', [ 'processedCount' => $count, 'successCount' => $successCount, 'failedCount' => $failedCount, 'skippedCount' => count($skipped), 'processed' => $processed, 'skipped' => $skipped, 'openNotSynced' => $openNotSynced, ]); SplitTicketSyncLogger::log('cron', 'scan end', [ 'processedCount' => $count, 'successCount' => $successCount, 'failedCount' => $failedCount, ]); return $count; } /** * 已配置自动同步周期且已实现蜘蛛的工单类型 * * @return list */ private function resolveAutoSyncTicketTypes(): array { $types = []; foreach (SplitScrmSpiderFactory::listSupportedTypes() as $ticketType) { if (SplitSyncConfigService::getIntervalMinutes($ticketType) > 0) { $types[] = $ticketType; } } return $types; } /** * @return array{success:bool,message:string} */ private function doSync(Ticket $ticket, string $syncMode): array { $ticketType = (string) $ticket['ticket_type']; $pageUrl = trim((string) $ticket['ticket_url']); if ($pageUrl === '') { SplitTicketSyncLogger::log('sync', 'empty pageUrl'); $this->markFailure($ticket, '工单链接为空'); return ['success' => false, 'message' => '工单链接为空']; } if (!SplitScrmSpiderFactory::isSupported($ticketType)) { SplitTicketSyncLogger::log('sync', 'spider not supported', ['ticketType' => $ticketType]); $this->markFailure($ticket, '工单类型尚未实现蜘蛛'); return ['success' => false, 'message' => '工单类型尚未实现蜘蛛']; } SplitTicketSyncLogger::log('sync', 'create spider', [ 'ticketType' => $ticketType, 'hasAccount' => trim((string) ($ticket['account'] ?? '')) !== '', ]); $spider = SplitScrmSpiderFactory::create( $ticketType, $pageUrl, (string) ($ticket['account'] ?? ''), (string) ($ticket['password'] ?? '') ); if ($spider === null) { $this->markFailure($ticket, '无法创建蜘蛛实例'); return ['success' => false, 'message' => '无法创建蜘蛛实例']; } Db::startTrans(); try { SplitTicketSyncLogger::log('sync', 'spider run begin'); $finalData = $spider->run(); if (!$finalData instanceof UnifiedScrmData) { throw new Exception('蜘蛛返回数据无效'); } $this->numberSync->syncFromUnifiedData($ticket, $finalData); $completeCount = max(0, $finalData->todayNewCount); $this->ruleService->applyTicketStatusRules($ticket, $completeCount); $freshTicket = Ticket::get((int) $ticket['id']) ?: $ticket; if ((string) $freshTicket['status'] === 'hidden') { $this->ruleService->cascadeTicketClosedToNumbers($freshTicket); } $ticket = $freshTicket; $inboundCount = $this->numberSync->sumInboundForTicket($ticket); $speed = $this->calcSpeedPerHour($ticket, $completeCount); $payload = [ 'complete_count' => $completeCount, 'inbound_count' => $inboundCount, 'speed_per_hour' => $speed['speed'], 'number_count' => max(0, $finalData->total), 'number_offline_count' => max(0, $finalData->totalOffline), 'number_banned_count' => 0, 'online_count' => max(0, $finalData->totalOnline), 'sync_fail_count' => 0, 'speed_snapshot_count' => $speed['snapshot_count'], 'speed_snapshot_time' => $speed['snapshot_time'], ]; $this->applySyncResult($ticket, $payload, true, '', $syncMode); Db::commit(); SplitTicketSyncLogger::log('sync', 'db commit ok', $payload); return ['success' => true, 'message' => '同步成功']; } catch (\Throwable $e) { Db::rollback(); $msg = $e->getMessage(); SplitTicketSyncLogger::log('sync', 'exception', [ 'type' => get_class($e), 'message' => $msg, 'file' => $e->getFile(), 'line' => $e->getLine(), ]); $this->markFailure($ticket, $msg); return ['success' => false, 'message' => $msg]; } } private function shouldSkip(Ticket $ticket): ?string { if ((string) $ticket['status'] === 'hidden') { return '工单已关闭'; } $failThreshold = SplitSyncConfigService::getFailPauseThreshold(); if ($failThreshold > 0 && (int) ($ticket['sync_fail_count'] ?? 0) >= $failThreshold) { return sprintf('连续同步失败超过%d次已暂停', $failThreshold); } if (!SplitScrmSpiderFactory::isSupported((string) $ticket['ticket_type'])) { return '工单类型尚未实现'; } $interval = SplitSyncConfigService::getIntervalMinutes((string) $ticket['ticket_type']); if ($interval <= 0) { return '该类型未配置自动同步周期'; } $lastSync = (int) ($ticket['sync_time'] ?? 0); $elapsed = $lastSync > 0 ? (time() - $lastSync) : null; if ($lastSync > 0 && $elapsed !== null && $elapsed < ($interval * 60)) { SplitTicketSyncLogger::log('sync', 'interval not reached', [ 'intervalMinutes' => $interval, 'elapsedSeconds' => $elapsed, 'needSeconds' => $interval * 60, ]); return '未到同步周期'; } return null; } /** * @param array $payload */ public function applySyncResult(Ticket $ticket, array $payload, bool $success, string $message = '', string $syncMode = 'manual'): void { $now = time(); $data = [ 'complete_count' => max(0, (int) ($payload['complete_count'] ?? 0)), 'inbound_count' => max(0, (int) ($payload['inbound_count'] ?? 0)), 'speed_per_hour' => max(0, (float) ($payload['speed_per_hour'] ?? 0)), 'number_count' => max(0, (int) ($payload['number_count'] ?? 0)), 'number_offline_count' => max(0, (int) ($payload['number_offline_count'] ?? 0)), 'number_banned_count' => max(0, (int) ($payload['number_banned_count'] ?? 0)), 'online_count' => max(0, (int) ($payload['online_count'] ?? 0)), 'sync_status' => $success ? 'success' : 'error', 'sync_time' => $now, 'sync_message' => $success ? '' : $message, 'sync_fail_count' => $success ? 0 : ((int) ($ticket['sync_fail_count'] ?? 0) + 1), 'speed_snapshot_count' => (int) ($payload['speed_snapshot_count'] ?? $ticket['speed_snapshot_count'] ?? 0), 'speed_snapshot_time' => (int) ($payload['speed_snapshot_time'] ?? $ticket['speed_snapshot_time'] ?? 0), ]; if ($success) { $data['sync_success_time'] = $now; $data['sync_success_mode'] = in_array($syncMode, ['auto', 'manual'], true) ? $syncMode : 'manual'; } if (!$ticket->allowField(array_keys($data))->save($data)) { throw new Exception('工单同步结果保存失败'); } } private function markFailure(Ticket $ticket, string $message): void { $failCount = (int) ($ticket['sync_fail_count'] ?? 0) + 1; $failThreshold = SplitSyncConfigService::getFailPauseThreshold(); $previousSyncStatus = (string) ($ticket['sync_status'] ?? 'pending'); $neverSyncedSuccessfully = $previousSyncStatus === 'pending' && (int) ($ticket['sync_time'] ?? 0) <= 0; $update = [ 'sync_status' => 'error', 'sync_time' => time(), 'sync_message' => $message, 'sync_fail_count' => $failCount, ]; if ($neverSyncedSuccessfully || ($failThreshold > 0 && $failCount >= $failThreshold)) { $update['status'] = 'hidden'; } $ticket->save($update); if (isset($update['status']) && $update['status'] === 'hidden') { $fresh = Ticket::get((int) $ticket['id']); if ($fresh) { $this->ruleService->cascadeTicketClosedToNumbers($fresh); } } } /** * @return array{speed:float,snapshot_count:int,snapshot_time:int} */ private function calcSpeedPerHour(Ticket $ticket, int $currentComplete): array { $now = time(); $snapshotTime = (int) ($ticket['speed_snapshot_time'] ?? 0); $snapshotCount = (int) ($ticket['speed_snapshot_count'] ?? 0); if ($snapshotTime <= 0) { return [ 'speed' => 0.0, 'snapshot_count' => $currentComplete, 'snapshot_time' => $now, ]; } $elapsed = $now - $snapshotTime; if ($elapsed >= 3600) { return [ 'speed' => 0.0, 'snapshot_count' => $currentComplete, 'snapshot_time' => $now, ]; } $hours = $elapsed > 0 ? ($elapsed / 3600) : 0; $delta = $currentComplete - $snapshotCount; $speed = ($delta < 0 || $hours <= 0) ? 0.0 : round($delta / $hours, 2); return [ 'speed' => $speed, 'snapshot_count' => $snapshotCount, 'snapshot_time' => $snapshotTime, ]; } /** * @return array{ticketId:int,ticketType:string,ticketUrl:string,reason:string} */ private function buildCronTicketEntry(int $ticketId, string $ticketType, string $ticketUrl, string $reason): array { return [ 'ticketId' => $ticketId, 'ticketType' => $ticketType, 'ticketUrl' => $ticketUrl, 'reason' => $reason, ]; } private function logCronCandidateSkipped(int $ticketId, string $ticketType, string $ticketUrl, string $reason): void { $ctx = [ 'ticketId' => $ticketId, 'ticketType' => $ticketType, 'ticketUrl' => $ticketUrl, 'reason' => $reason, ]; SplitTicketSyncLogger::logCron('candidate skipped', $ctx); SplitTicketSyncLogger::log('cron', 'candidate skipped', $ctx); } /** * 开启但未在本轮同步的工单及原因(供 cron 汇总) * * @param list $candidateIds * @param list> $processed * @param list> $skipped * @return list> */ private function buildOpenNotSyncedAudit( array $autoSyncTypes, int $failThreshold, int $maxPerRun, array $candidateIds, array $processed, array $skipped, bool $stoppedByNodeBusy ): array { $handledIds = []; foreach ($processed as $row) { $handledIds[(int) ($row['ticketId'] ?? 0)] = true; } foreach ($skipped as $row) { $handledIds[(int) ($row['ticketId'] ?? 0)] = true; } $candidateIdMap = array_fill_keys($candidateIds, true); $supportedTypes = SplitScrmSpiderFactory::listSupportedTypes(); if ($supportedTypes === []) { return []; } $openList = Ticket::where('status', 'normal') ->whereIn('ticket_type', $supportedTypes) ->field('id,ticket_type,ticket_url,sync_fail_count,sync_time,sync_status') ->select(); $result = []; foreach ($openList as $ticket) { $ticketId = (int) $ticket['id']; if (isset($handledIds[$ticketId])) { continue; } $ticketType = (string) $ticket['ticket_type']; $ticketUrl = (string) $ticket['ticket_url']; $reason = $this->resolveOpenNotSyncedReason( $ticket, $autoSyncTypes, $failThreshold, $maxPerRun, $candidateIdMap, $stoppedByNodeBusy ); if ($reason === null) { continue; } $result[] = [ 'ticketId' => $ticketId, 'ticketType' => $ticketType, 'ticketUrl' => $ticketUrl, 'reason' => $reason, ]; } return $result; } /** * @param array $ticket * @param array $candidateIdMap */ private function resolveOpenNotSyncedReason( array $ticket, array $autoSyncTypes, int $failThreshold, int $maxPerRun, array $candidateIdMap, bool $stoppedByNodeBusy ): ?string { $ticketType = (string) ($ticket['ticket_type'] ?? ''); $failCount = (int) ($ticket['sync_fail_count'] ?? 0); if ($failThreshold > 0 && $failCount >= $failThreshold) { return sprintf('连续同步失败超过%d次已暂停自动同步', $failThreshold); } if (!in_array($ticketType, $autoSyncTypes, true)) { return '该类型未配置自动同步周期'; } if (!isset($candidateIdMap[(int) ($ticket['id'] ?? 0)])) { return sprintf('未进入本轮候选队列(仅扫描最久未同步的前%d条)', $maxPerRun * 5); } if ($stoppedByNodeBusy) { return 'Node 队列繁忙,本轮未执行'; } $skip = $this->shouldSkip(Ticket::get((int) $ticket['id']) ?: new Ticket($ticket)); return $skip ?? '未知原因,未进入执行队列'; } /** * 按 sessionKey(ticketType:host)分组排序,使同域名工单连续执行以命中 Browser 温池 * * @param \think\Collection|array> $list * @return list> */ private function sortTicketsBySessionKey($list): array { $rows = $list instanceof \think\Collection ? $list->all() : (array) $list; usort($rows, static function ($a, $b): int { $typeA = (string) (is_array($a) ? ($a['ticket_type'] ?? '') : ($a['ticket_type'] ?? '')); $urlA = (string) (is_array($a) ? ($a['ticket_url'] ?? '') : ($a['ticket_url'] ?? '')); $typeB = (string) (is_array($b) ? ($b['ticket_type'] ?? '') : ($b['ticket_type'] ?? '')); $urlB = (string) (is_array($b) ? ($b['ticket_url'] ?? '') : ($b['ticket_url'] ?? '')); $keyA = AntiBotConfigBuilder::resolveSessionKey($typeA, $urlA); $keyB = AntiBotConfigBuilder::resolveSessionKey($typeB, $urlB); if ($keyA === $keyB) { $timeA = (int) (is_array($a) ? ($a['sync_time'] ?? 0) : ($a['sync_time'] ?? 0)); $timeB = (int) (is_array($b) ? ($b['sync_time'] ?? 0) : ($b['sync_time'] ?? 0)); if ($timeA === $timeB) { $idA = (int) (is_array($a) ? ($a['id'] ?? 0) : ($a['id'] ?? 0)); $idB = (int) (is_array($b) ? ($b['id'] ?? 0) : ($b['id'] ?? 0)); return $idA <=> $idB; } return $timeA <=> $timeB; } return strcmp($keyA, $keyB); }); return $rows; } }