From e2949288aa615d8f6ee0cefe4018901cf34dfea3 Mon Sep 17 00:00:00 2001 From: klockky Date: Tue, 14 Jul 2026 17:32:58 +0300 Subject: [PATCH] =?UTF-8?q?fix:=20=D0=BB=D0=BE=D0=B3=D0=B8=D1=80=D0=BE?= =?UTF-8?q?=D0=B2=D0=B0=D1=82=D1=8C=20SESSION=5FINIT=20=D0=B8=20=D1=80?= =?UTF-8?q?=D0=B5=D0=B0=D0=BB=D1=8C=D0=BD=D1=8B=D0=B9=20wire-seq=20=D0=B2?= =?UTF-8?q?=20=D0=B4=D0=B0=D0=BC=D0=BF=D0=B5=20=D1=82=D1=80=D0=B0=D1=84?= =?UTF-8?q?=D0=B8=D0=BA=D0=B0?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- lib/backend/api.dart | 93 ++++++++++++++++++++++++++++++-------------- 1 file changed, 63 insertions(+), 30 deletions(-) diff --git a/lib/backend/api.dart b/lib/backend/api.dart index 1d2fb6b..fa39b16 100644 --- a/lib/backend/api.dart +++ b/lib/backend/api.dart @@ -1,4 +1,5 @@ import 'dart:async'; +import 'dart:convert'; import 'dart:io'; import 'package:device_info_plus/device_info_plus.dart'; @@ -36,6 +37,7 @@ class Api { /// ради registerHandler/pushStream). final PacketDispatcher _dispatcher = PacketDispatcher(); StreamSubscription<(int, Map)>? _pushSub; + StreamSubscription? _wireLogSub; SessionState _sessionState = SessionState.disconnected; final _stateController = StreamController.broadcast(); @@ -80,7 +82,6 @@ class Api { int _reconnectAttempts = 0; bool _autoReconnect = false; int _sessionEpoch = 0; - int _seq = 0; bool? _lastInteractive; static const Duration _connectWatchdogTimeout = Duration(seconds: 75); @@ -90,8 +91,6 @@ class Api { int get sessionEpoch => _sessionEpoch; - int _nextSeq() => _seq = (_seq + 1) & 0xFFFF; - // Публичное API /// Подключается к серверу, шлёт хэндшейк, запускает пинг. @@ -130,7 +129,7 @@ class Api { } if (gen != _connectGen) return; - final options = await _buildSessionOptions(endpoint); + final (session, wireLog) = await _buildSessionOptions(endpoint); if (gen != _connectGen) return; logger.i( @@ -145,8 +144,10 @@ class Api { } } - final session = options; _session = session; + // Подписываемся на wire-лог ядра ДО connect(), чтобы поймать пакеты + // SESSION_INIT-хендшейка (иначе они уходят до listen и теряются). + _wireLogSub = wireLog.listen(_onWireLog); _pushSub = session.pushesMap().listen(_onPush); TrafficMonitor.instance.recordEvent( 'connect', @@ -270,32 +271,21 @@ class Api { if (session == null) { throw StateError('Нет соединения (${Opcode.name(opcode)})'); } - final seq = _nextSeq(); - DebugSessionLog.instance.recordRequest(opcode, seq, payload); - TrafficMonitor.instance.recordOutgoing(opcode, payload, seq, 0); - - final KolibriResponse resp; - try { - resp = await session - .requestMapFull(opcode, Map.from(payload)) - .timeout( - ServerConfig.requestTimeout, - onTimeout: () => - throw TimeoutException('${Opcode.name(opcode)} таймаут'), - ); - } catch (e, st) { - DebugSessionLog.instance.recordError(seq, e); - Error.throwWithStackTrace(e, st); - } + // Лог запроса/ответа ведётся из wire-лога ядра (_onWireLog) по настоящему + // проводному seq, поэтому здесь ничего не пишем. + final KolibriResponse resp = await session + .requestMapFull(opcode, Map.from(payload)) + .timeout( + ServerConfig.requestTimeout, + onTimeout: () => + throw TimeoutException('${Opcode.name(opcode)} таймаут'), + ); final packet = Packet( cmd: resp.cmd, - seq: seq, opcode: resp.opcode, payload: resp.payload, ); - DebugSessionLog.instance.recordResponse(seq, packet.cmd, packet.payload); - TrafficMonitor.instance.recordIncoming(packet, 0); if (packet.isError) { if (isSessionExpiredPayload(packet.payload)) { @@ -303,7 +293,6 @@ class Api { messageFromErrorPayload(packet.payload), ); _sessionExpiredController.add(ex); - DebugSessionLog.instance.recordError(seq, ex); throw ex; } final text = _serverErrorText(packet.payload); @@ -312,7 +301,6 @@ class Api { messageFromErrorPayload(packet.payload), errorKey: resp.errorKey, ); - DebugSessionLog.instance.recordError(seq, err); throw err; } return packet; @@ -369,7 +357,7 @@ class Api { /// Строит устройство-поля и создаёт сессию ядра. Заодно заполняет /// [_userAgent] и [_deviceId] для геттеров. - Future _buildSessionOptions( + Future<(KolibriSession, Stream)> _buildSessionOptions( ({String host, int port}) endpoint, ) async { final deviceInfo = DeviceInfoPlugin(); @@ -477,7 +465,7 @@ class Api { final insecureTls = await TlsConfig.isInsecureAllowed(); final proxy = await _buildProxyUrl(); - return openSession( + return openSessionWithWireLog( host: endpoint.host, port: endpoint.port, deviceId: deviceId, @@ -519,10 +507,53 @@ class Api { opcode: event.$1, payload: event.$2, ); - TrafficMonitor.instance.recordIncoming(packet, 0); + // Учёт трафика пушей ведётся из wire-лога ядра (_onWireLog). _dispatcher.dispatch(packet); } + /// Единый источник лога трафика: ядро отдаёт сюда каждый пакет обеих сторон — + /// включая SESSION_INIT-хендшейк и пинги — с настоящим проводным seq. Раньше + /// лог вёлся вручную из [sendRequest] по локальному счётчику, из-за чего + /// хендшейк/пинги в дамп не попадали, а seq был смещён относительно провода. + void _onWireLog(WireLogEvent e) { + final payload = _decodeWireJson(e.json); + final cmd = _wireCmdCode(e.cmd); + if (e.direction == 'out') { + DebugSessionLog.instance.recordRequest(e.opcode, e.seq, payload); + TrafficMonitor.instance.recordOutgoing(e.opcode, payload, e.seq, 0); + return; + } + // Входящие: ответы матчатся по seq, пуши идут только в монитор трафика. + if (e.cmd != 'push') { + DebugSessionLog.instance.recordResponse(e.seq, cmd, payload); + } + TrafficMonitor.instance.recordIncoming( + Packet(cmd: cmd, seq: e.seq, opcode: e.opcode, payload: payload), + 0, + ); + } + + static int _wireCmdCode(String cmd) { + switch (cmd) { + case 'ok': + return CmdType.ok; + case 'not_found': + return CmdType.notFound; + case 'error': + return CmdType.error; + default: + return CmdType.request; // 'request' и 'push' + } + } + + static dynamic _decodeWireJson(String json) { + try { + return jsonDecode(json); + } catch (_) { + return json; + } + } + void _setSessionState(SessionState state) { if (_sessionState == state) return; _sessionState = state; @@ -573,6 +604,8 @@ class Api { _livenessTimer = null; _pushSub?.cancel(); _pushSub = null; + _wireLogSub?.cancel(); + _wireLogSub = null; _lastInteractive = null; final session = _session; _session = null;