handleIdentifyResponse method
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);
}
}