handleIncomingDatagram method
Processes an incoming datagram from the multiplexer.
Implementation
Future<void> handleIncomingDatagram(Uint8List data, InternetAddress fromAddress, int fromPort) async {
// Any datagram is proof the path is alive; restart the idle clock (RFC 9000
// section 10.1: the timer resets on receiving and processing a packet).
_lastActivity = DateTime.now();
// Track received bytes for anti-amplification
if (!_addressValidated) {
_bytesReceivedBeforeValidation += data.length;
// Address is validated after receiving a certain amount of data
// or on successful handshake completion
if (_bytesReceivedBeforeValidation >= 1000 || _handshakeCompleted) {
_onAddressValidated();
}
}
try {
final packet = UDXPacket.fromBytes(data);
// Check if version is supported
if (!UdxVersion.isSupported(packet.version) && !_handshakeCompleted) {
// Send VERSION_NEGOTIATION packet
final versionNegPacket = VersionNegotiationPacket(
destinationCid: packet.sourceCid,
sourceCid: packet.destinationCid,
supportedVersions: UdxVersion.supportedVersions,
);
multiplexer.send(versionNegPacket.toBytes(), fromAddress, fromPort);
emit('versionNegotiation', {'clientVersion': packet.version, 'supportedVersions': UdxVersion.supportedVersions});
return;
}
// Always update the remote CID from the packet's source CID.
// This ensures that even during retransmissions or path migrations,
// we are targeting the correct peer identifier.
cids.remoteCid = packet.sourceCid;
if (!_handshakeCompleted) {
// The handshake is considered complete on the first valid packet received.
_handshakeCompleted = true;
if (!_handshakeCompleter.isCompleted) {
_handshakeCompleter.complete();
}
// Notify observer of successful handshake
if (_handshakeStartTime != null) {
final duration = DateTime.now().difference(_handshakeStartTime!);
metricsObserver?.onHandshakeComplete(cids.localCid, duration, true, null);
}
emit('connect');
}
} catch (e) {
// Ignore if it's not a valid packet.
return;
}
try {
final udxPacket = UDXPacket.fromBytes(data);
// v3 moved the STREAM frame's data length behind an eight-byte offset, so
// an older peer's frame parses without error into nonsense. Drop anything
// that is not the current version: a mismatch has to look like an
// unreachable peer, which is diagnosable, rather than corrupt data.
if (udxPacket.version != UdxVersion.current) return;
if (UdxLogging.info) {
UdxLogging.infoLog('[DIAG-UDX-RECV] seq=${udxPacket.sequence} frames=${udxPacket.frames.length} from=${fromAddress.address}:$fromPort');
}
// --- Path Migration Logic ---
final pathHasChanged = remoteAddress.address != fromAddress.address || remotePort != fromPort;
if (pathHasChanged && _pathChallengeData == null) {
_initiatePathValidation(fromAddress, fromPort);
}
// --- PMTUD: Handle ACKs for Probes & Trigger New Probes ---
for (final frame in udxPacket.frames.whereType<AckFrame>()) {
_handleAckFrameForPmtud(frame);
}
_sendMtuProbeIfNeeded();
// --- Process connection-level frames (not sequence-dependent) ---
for (final frame in udxPacket.frames) {
if (frame is ConnectionCloseFrame) {
emit('connectionClose', {
'errorCode': frame.errorCode,
'frameType': frame.frameType,
'reason': frame.reasonPhrase
});
await close();
return;
} else if (frame is MaxDataFrame) {
_handleMaxDataFrame(frame);
} else if (frame is MaxStreamsFrame) {
_handleMaxStreamsFrame(frame);
} else if (frame is PathChallengeFrame) {
_handlePathChallenge(frame, fromAddress, fromPort);
} else if (frame is PathResponseFrame) {
_handlePathResponse(frame, fromAddress, fromPort);
} else if (frame is DataBlockedFrame) {
emit('dataBlocked', {'maxData': frame.maxData});
}
}
// --- Process ACK frames at connection level ---
for (final frame in udxPacket.frames.whereType<AckFrame>()) {
_handleConnectionAckFrame(frame);
}
// --- Process RESET, STOP_SENDING, WINDOW_UPDATE, STREAM_DATA_BLOCKED immediately (not sequence-dependent) ---
for (final frame in udxPacket.frames) {
if (frame is ResetStreamFrame) {
_findStream(udxPacket.destinationStreamId, udxPacket.sourceStreamId)
?.deliverReset(frame.errorCode);
} else if (frame is StopSendingFrame) {
_findStream(udxPacket.destinationStreamId, udxPacket.sourceStreamId)
?.deliverStopSending(frame.errorCode);
} else if (frame is WindowUpdateFrame) {
_findStream(udxPacket.destinationStreamId, udxPacket.sourceStreamId)
?.deliverWindowUpdate(frame.windowSize);
} else if (frame is StreamDataBlockedFrame) {
_findStream(udxPacket.destinationStreamId, udxPacket.sourceStreamId)
?.deliverStreamDataBlocked();
}
}
// Every frame is handled on arrival. Nothing waits for an earlier packet:
// STREAM frames carry their own byte offset, so each stream places its
// own bytes and is held up only by its own gaps.
//
// This used to reorder whole packets here, on the connection's sequence
// number, before any stream saw them. That made a gap anywhere stall
// every stream on the connection — one lossy stream held up all its
// siblings — and it was the only thing keeping the byte stream contiguous,
// so it could not simply be removed. Offsets replace it.
bool containsAckElicitingFrames =
udxPacket.frames.any((f) => f is StreamFrame || f is PingFrame);
bool hasStreamData = udxPacket.frames
.any((f) => f is StreamFrame && (f.data.isNotEmpty || f.isFin || f.isSyn));
if (hasStreamData) {
_largestAckedPacketArrivalTime = DateTime.now();
}
_processPacketFrames(udxPacket, fromAddress, fromPort);
if (containsAckElicitingFrames) {
// ACKs reflect receipt, not delivery. A packet whose bytes are waiting
// on an earlier gap has still arrived, and saying so lets the peer
// retransmit only what is genuinely missing.
_receivedPacketSequences.add(udxPacket.sequence);
_sendConnectionAck();
}
} catch (e) {
// Ignore invalid packets
}
}