handleIdentifyResponse method

Future<void> handleIdentifyResponse(
  1. P2PStream stream,
  2. bool isPush
)

Implementation

Future<void> handleIdentifyResponse(P2PStream stream, bool isPush) async {
  final peer = stream.conn.remotePeer;
  final String side = stream.conn.localPeer == host.id ? "CLIENT" : "SERVER";
  final handleStart = DateTime.now();


  try {

    await stream.scope().setService(serviceName);
    final serviceScopeTime = DateTime.now().difference(handleStart);

  } catch (e) {
    _log.log(isPush ? Level.FINE : _quietLevel(stream.conn, Level.WARNING), ' [HANDLE-IDENTIFY-RESPONSE-PHASE-1-ERROR] ($side) Error attaching stream to identify service for peer=$peer: $e. Resetting stream.');
    await stream.reset().catchError((_) {});
    rethrow;
  }

  try {

    await stream.scope().reserveMemory(signedIDSize, ReservationPriority.always);
    final memoryReserveTime = DateTime.now().difference(handleStart);

  } catch (e) {
    _log.log(isPush ? Level.FINE : _quietLevel(stream.conn, Level.WARNING), ' [HANDLE-IDENTIFY-RESPONSE-PHASE-2-ERROR] ($side) Error reserving memory for identify stream for peer=$peer: $e. Resetting stream.');
    await stream.reset().catchError((_) {});
    rethrow;
  }

  try {
    final conn = stream.conn;

    final received = await _readAllIDMessages(stream, isPush);
    if (received == null && isPush) {
      // The peer closed the stream before it sent a message, for example
      // because it shuts down. An empty message would clear its protocols
      // and addresses, so nothing is consumed.
      _log.fine('IdentifyService.handleIdentifyResponse ($side): Push from $peer ended before a message; ignored.');
      await stream.closeWrite().catchError((_) {});
      return;
    }
    final mes = received ?? Identify();

    await _consumeMessage(mes, conn, isPush);
    final consumeMessageTime = DateTime.now().difference(handleStart);


    if (metricsTracer != null) {
      metricsTracer!.identifyReceived(isPush, mes.protocols.length, mes.listenAddrs.length);
    }


    await _connsMutex.synchronized( () async {
      final e = _conns[conn];
      if (e == null) {
        _log.fine('IdentifyService.handleIdentifyResponse ($side): Connection entry for $peer already removed (disconnected). Cannot update push support.');
        return;
      }
      _log.finer('IdentifyService.handleIdentifyResponse ($side): Checking push support for $peer in peerstore.');
      final sup = await host.peerStore.protoBook.supportsProtocols(conn.remotePeer, [idPush]);
      if (sup.isNotEmpty) {
        e.pushSupport = IdentifyPushSupport.supported;
        _log.fine('IdentifyService.handleIdentifyResponse ($side): Peer $peer supports push.');
      } else {
        e.pushSupport = IdentifyPushSupport.unsupported;
        _log.fine('IdentifyService.handleIdentifyResponse ($side): Peer $peer does not support push.');
      }
      if (metricsTracer != null) {
        metricsTracer!.connPushSupport(e.pushSupport);
      }
    });
    final pushSupportTime = DateTime.now().difference(handleStart);



    await stream.closeWrite();
    final totalTime = DateTime.now().difference(handleStart);

  } catch (e, st) {
    final errorTime = DateTime.now().difference(handleStart);
    _log.log(isPush ? Level.FINE : _quietLevel(stream.conn, Level.SEVERE), ' [HANDLE-IDENTIFY-RESPONSE-ERROR] ($side) Error reading or processing identify message from peer=$peer, duration=${errorTime.inMilliseconds}ms, error=$e\n$st');
    await stream.reset().catchError((_) {});
    rethrow;
  } finally {

    stream.scope().releaseMemory(signedIDSize);
  }
}