[dart:io] Fix hanging on zero-length datagram
This is observed on Win and Linux. Here is a doc to explain the problem: https://docs.google.com/document/d/1ZzyBUMrDHLU6vNryjgSMJfUruOdpmTImtE_s2E0J8IQ/edit?usp=sharing. Bug: https://github.com/dart-lang/sdk/issues/39910 Change-Id: Ia961239f45615f14108bcd66043ac33d9a0a4abe Reviewed-on: https://dart-review.googlesource.com/c/sdk/+/137425 Commit-Queue: Zichang Guo <zichangguo@google.com> Reviewed-by: Siva Annamalai <asiva@google.com>
This commit is contained in:
committed by
commit-bot@chromium.org
parent
a980fb1849
commit
c326c587c5
@@ -208,7 +208,7 @@ void Handle::ReadComplete(OverlappedBuffer* buffer) {
|
||||
// Currently only one outstanding read at the time.
|
||||
ASSERT(pending_read_ == buffer);
|
||||
ASSERT(data_ready_ == NULL);
|
||||
if (!IsClosing() && !buffer->IsEmpty()) {
|
||||
if (!IsClosing()) {
|
||||
data_ready_ = pending_read_;
|
||||
} else {
|
||||
OverlappedBuffer::DisposeBuffer(buffer);
|
||||
@@ -611,10 +611,13 @@ intptr_t Handle::Available() {
|
||||
if (data_ready_ == NULL) {
|
||||
return 0;
|
||||
}
|
||||
ASSERT(!data_ready_->IsEmpty());
|
||||
return data_ready_->GetRemainingLength();
|
||||
}
|
||||
|
||||
bool Handle::DataReady() {
|
||||
return data_ready_ != NULL;
|
||||
}
|
||||
|
||||
intptr_t Handle::Read(void* buffer, intptr_t num_bytes) {
|
||||
MonitorLocker ml(&monitor_);
|
||||
if (data_ready_ == NULL) {
|
||||
|
||||
@@ -172,6 +172,7 @@ class Handle : public ReferenceCounted<Handle>, public DescriptorInfoBase {
|
||||
|
||||
// Socket interface exposing normal socket operations.
|
||||
intptr_t Available();
|
||||
bool DataReady();
|
||||
intptr_t Read(void* buffer, intptr_t num_bytes);
|
||||
intptr_t RecvFrom(void* buffer,
|
||||
intptr_t num_bytes,
|
||||
|
||||
@@ -132,6 +132,7 @@ namespace bin {
|
||||
V(ServerSocket_CreateBindListen, 7) \
|
||||
V(SocketBase_IsBindError, 2) \
|
||||
V(Socket_Available, 1) \
|
||||
V(Socket_AvailableDatagram, 1) \
|
||||
V(Socket_CreateBindConnect, 5) \
|
||||
V(Socket_CreateBindDatagram, 6) \
|
||||
V(Socket_CreateConnect, 4) \
|
||||
|
||||
+18
-6
@@ -399,17 +399,12 @@ void FUNCTION_NAME(Socket_RecvFrom)(Dart_NativeArguments args) {
|
||||
RawAddr addr;
|
||||
const intptr_t bytes_read = SocketBase::RecvFrom(
|
||||
socket->fd(), recv_buffer, kReceiveBufferLen, &addr, SocketBase::kAsync);
|
||||
if (bytes_read == 0) {
|
||||
Dart_SetReturnValue(args, Dart_Null());
|
||||
return;
|
||||
}
|
||||
if (bytes_read < 0) {
|
||||
ASSERT(bytes_read == -1);
|
||||
Dart_ThrowException(DartUtils::NewDartOSError());
|
||||
}
|
||||
|
||||
// Datagram data read. Copy into buffer of the exact size,
|
||||
ASSERT(bytes_read > 0);
|
||||
ASSERT(bytes_read >= 0);
|
||||
uint8_t* data_buffer = NULL;
|
||||
Dart_Handle data = IOBuffer::Allocate(bytes_read, &data_buffer);
|
||||
if (Dart_IsNull(data)) {
|
||||
@@ -421,6 +416,11 @@ void FUNCTION_NAME(Socket_RecvFrom)(Dart_NativeArguments args) {
|
||||
ASSERT(data_buffer != NULL);
|
||||
memmove(data_buffer, recv_buffer, bytes_read);
|
||||
|
||||
// Memory Sanitizer complains addr not being initialized, which is done
|
||||
// through RecvFrom().
|
||||
// Issue: https://github.com/google/sanitizers/issues/1201
|
||||
MSAN_UNPOISON(&addr, sizeof(RawAddr));
|
||||
|
||||
// Get the port and clear it in the sockaddr structure.
|
||||
int port = SocketAddress::GetAddrPort(addr);
|
||||
if (addr.addr.sa_family == AF_INET) {
|
||||
@@ -1015,6 +1015,18 @@ void FUNCTION_NAME(Socket_LeaveMulticast)(Dart_NativeArguments args) {
|
||||
}
|
||||
}
|
||||
|
||||
void FUNCTION_NAME(Socket_AvailableDatagram)(Dart_NativeArguments args) {
|
||||
const int kReceiveBufferLen = 1;
|
||||
Socket* socket =
|
||||
Socket::GetSocketIdNativeField(Dart_GetNativeArgument(args, 0));
|
||||
ASSERT(socket != NULL);
|
||||
// Ensure that a receive buffer for peeking the UDP socket exists.
|
||||
uint8_t recv_buffer[kReceiveBufferLen];
|
||||
bool available = SocketBase::AvailableDatagram(socket->fd(), recv_buffer,
|
||||
kReceiveBufferLen);
|
||||
Dart_SetBooleanReturnValue(args, available);
|
||||
}
|
||||
|
||||
static void NormalSocketFinalizer(void* isolate_data,
|
||||
Dart_WeakPersistentHandle handle,
|
||||
void* data) {
|
||||
|
||||
@@ -102,7 +102,7 @@ void SocketAddress::SetAddrPort(RawAddr* addr, intptr_t port) {
|
||||
}
|
||||
|
||||
intptr_t SocketAddress::GetAddrPort(const RawAddr& addr) {
|
||||
if (addr.ss.ss_family == AF_INET) {
|
||||
if (addr.addr.sa_family == AF_INET) {
|
||||
return ntohs(addr.in.sin_port);
|
||||
} else {
|
||||
return ntohs(addr.in6.sin6_port);
|
||||
|
||||
@@ -167,6 +167,7 @@ class SocketBase : public AllStatic {
|
||||
intptr_t num_bytes,
|
||||
RawAddr* addr,
|
||||
SocketOpKind sync);
|
||||
static bool AvailableDatagram(intptr_t fd, void* buffer, intptr_t num_bytes);
|
||||
// Returns true if the given error-number is because the system was not able
|
||||
// to bind the socket to a specific IP.
|
||||
static bool IsBindError(intptr_t error_number);
|
||||
|
||||
@@ -82,14 +82,18 @@ intptr_t SocketBase::RecvFrom(intptr_t fd,
|
||||
socklen_t addr_len = sizeof(addr->ss);
|
||||
ssize_t read_bytes = TEMP_FAILURE_RETRY(
|
||||
recvfrom(fd, buffer, num_bytes, 0, &addr->addr, &addr_len));
|
||||
if ((sync == kAsync) && (read_bytes == -1) && (errno == EWOULDBLOCK)) {
|
||||
// If the read would block we need to retry and therefore return 0
|
||||
// as the number of bytes written.
|
||||
read_bytes = 0;
|
||||
}
|
||||
return read_bytes;
|
||||
}
|
||||
|
||||
bool SocketBase::AvailableDatagram(intptr_t fd,
|
||||
void* buffer,
|
||||
intptr_t num_bytes) {
|
||||
ASSERT(fd >= 0);
|
||||
ssize_t read_bytes =
|
||||
TEMP_FAILURE_RETRY(recvfrom(fd, buffer, num_bytes, MSG_PEEK, NULL, NULL));
|
||||
return read_bytes >= 0;
|
||||
}
|
||||
|
||||
intptr_t SocketBase::Write(intptr_t fd,
|
||||
const void* buffer,
|
||||
intptr_t num_bytes,
|
||||
|
||||
@@ -82,14 +82,18 @@ intptr_t SocketBase::RecvFrom(intptr_t fd,
|
||||
socklen_t addr_len = sizeof(addr->ss);
|
||||
ssize_t read_bytes = TEMP_FAILURE_RETRY(
|
||||
recvfrom(fd, buffer, num_bytes, 0, &addr->addr, &addr_len));
|
||||
if ((sync == kAsync) && (read_bytes == -1) && (errno == EWOULDBLOCK)) {
|
||||
// If the read would block we need to retry and therefore return 0
|
||||
// as the number of bytes written.
|
||||
read_bytes = 0;
|
||||
}
|
||||
return read_bytes;
|
||||
}
|
||||
|
||||
bool SocketBase::AvailableDatagram(intptr_t fd,
|
||||
void* buffer,
|
||||
intptr_t num_bytes) {
|
||||
ASSERT(fd >= 0);
|
||||
ssize_t read_bytes =
|
||||
TEMP_FAILURE_RETRY(recvfrom(fd, buffer, num_bytes, MSG_PEEK, NULL, NULL));
|
||||
return read_bytes >= 0;
|
||||
}
|
||||
|
||||
intptr_t SocketBase::Write(intptr_t fd,
|
||||
const void* buffer,
|
||||
intptr_t num_bytes,
|
||||
|
||||
@@ -81,14 +81,18 @@ intptr_t SocketBase::RecvFrom(intptr_t fd,
|
||||
socklen_t addr_len = sizeof(addr->ss);
|
||||
ssize_t read_bytes = TEMP_FAILURE_RETRY(
|
||||
recvfrom(fd, buffer, num_bytes, 0, &addr->addr, &addr_len));
|
||||
if ((sync == kAsync) && (read_bytes == -1) && (errno == EWOULDBLOCK)) {
|
||||
// If the read would block we need to retry and therefore return 0
|
||||
// as the number of bytes written.
|
||||
read_bytes = 0;
|
||||
}
|
||||
return read_bytes;
|
||||
}
|
||||
|
||||
bool SocketBase::AvailableDatagram(intptr_t fd,
|
||||
void* buffer,
|
||||
intptr_t num_bytes) {
|
||||
ASSERT(fd >= 0);
|
||||
ssize_t read_bytes =
|
||||
TEMP_FAILURE_RETRY(recvfrom(fd, buffer, num_bytes, MSG_PEEK, NULL, NULL));
|
||||
return read_bytes >= 0;
|
||||
}
|
||||
|
||||
intptr_t SocketBase::Write(intptr_t fd,
|
||||
const void* buffer,
|
||||
intptr_t num_bytes,
|
||||
|
||||
@@ -88,6 +88,13 @@ intptr_t SocketBase::RecvFrom(intptr_t fd,
|
||||
return handle->RecvFrom(buffer, num_bytes, &addr->addr, addr_len);
|
||||
}
|
||||
|
||||
bool SocketBase::AvailableDatagram(intptr_t fd,
|
||||
void* buffer,
|
||||
intptr_t num_bytes) {
|
||||
ClientSocket* client_socket = reinterpret_cast<ClientSocket*>(fd);
|
||||
return client_socket->DataReady();
|
||||
}
|
||||
|
||||
intptr_t SocketBase::Write(intptr_t fd,
|
||||
const void* buffer,
|
||||
intptr_t num_bytes,
|
||||
|
||||
@@ -364,12 +364,18 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
// Holds the address used to connect or bind the socket.
|
||||
InternetAddress localAddress;
|
||||
|
||||
// The number of available bytes to read.
|
||||
// The size of data that is ready to be read, for TCP sockets.
|
||||
// This might be out-of-date when Read is called.
|
||||
// The number of pending connections, for Listening sockets.
|
||||
int available = 0;
|
||||
|
||||
// Only used for UDP sockets.
|
||||
bool _availableDatagram = false;
|
||||
|
||||
// The number of incoming connnections for Listening socket.
|
||||
int connections = 0;
|
||||
|
||||
// The count of received event from eventhandler.
|
||||
int tokens = 0;
|
||||
|
||||
bool sendReadEvents = false;
|
||||
@@ -759,26 +765,19 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
}
|
||||
|
||||
Datagram receive() {
|
||||
if (isClosing || isClosed) return null;
|
||||
if (isClosing || isClosed || !_availableDatagram) return null;
|
||||
try {
|
||||
Datagram result = nativeRecvFrom();
|
||||
if (result != null) {
|
||||
// Read the next available. Available is only for the next datagram, not
|
||||
// the sum of all datagrams pending, so we need to call after each
|
||||
// receive. If available becomes > 0, the _NativeSocket will continue to
|
||||
// emit read events.
|
||||
available = nativeAvailable();
|
||||
if (resourceInfo != null) {
|
||||
resourceInfo.totalRead += result.data.length;
|
||||
}
|
||||
}
|
||||
assert(result != null);
|
||||
if (resourceInfo != null) {
|
||||
resourceInfo.totalRead += result.data.length;
|
||||
resourceInfo.didRead();
|
||||
}
|
||||
if (!const bool.fromEnvironment("dart.vm.product")) {
|
||||
_SocketProfile.collectStatistic(nativeGetSocketId(),
|
||||
_SocketProfileType.readBytes, result?.data?.length);
|
||||
}
|
||||
_availableDatagram = nativeAvailableDatagram();
|
||||
return result;
|
||||
} catch (e) {
|
||||
reportError(e, StackTrace.current, "Receive failed");
|
||||
@@ -913,7 +912,7 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
readEventIssued = false;
|
||||
if (isClosing) return;
|
||||
if (!sendReadEvents) return;
|
||||
if (available == 0) {
|
||||
if (stopRead()) {
|
||||
if (isClosedRead && !closedReadEventSent) {
|
||||
if (isClosedWrite) close();
|
||||
var handler = eventHandlers[closedEvent];
|
||||
@@ -933,6 +932,14 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
scheduleMicrotask(issue);
|
||||
}
|
||||
|
||||
bool stopRead() {
|
||||
if (isUdp) {
|
||||
return !_availableDatagram;
|
||||
} else {
|
||||
return available == 0;
|
||||
}
|
||||
}
|
||||
|
||||
void issueWriteEvent({bool delayed: true}) {
|
||||
if (writeEventIssued) return;
|
||||
if (!writeAvailable) return;
|
||||
@@ -980,7 +987,11 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
if (isListening) {
|
||||
connections++;
|
||||
} else {
|
||||
available = nativeAvailable();
|
||||
if (isUdp) {
|
||||
_availableDatagram = nativeAvailableDatagram();
|
||||
} else {
|
||||
available = nativeAvailable();
|
||||
}
|
||||
issueReadEvent();
|
||||
continue;
|
||||
}
|
||||
@@ -1239,6 +1250,7 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
|
||||
void nativeSetSocketId(int id, int typeFlags) native "Socket_SetSocketId";
|
||||
int nativeAvailable() native "Socket_Available";
|
||||
bool nativeAvailableDatagram() native "Socket_AvailableDatagram";
|
||||
Uint8List nativeRead(int len) native "Socket_Read";
|
||||
Datagram nativeRecvFrom() native "Socket_RecvFrom";
|
||||
int nativeWrite(List<int> buffer, int offset, int bytes)
|
||||
|
||||
@@ -364,12 +364,18 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
// Holds the address used to connect or bind the socket.
|
||||
late InternetAddress localAddress;
|
||||
|
||||
// The number of available bytes to read.
|
||||
// The size of data that is ready to be read, for TCP sockets.
|
||||
// This might be out-of-date when Read is called.
|
||||
// The number of pending connections, for Listening sockets.
|
||||
int available = 0;
|
||||
|
||||
// Only used for UDP sockets.
|
||||
bool _availableDatagram = false;
|
||||
|
||||
// The number of incoming connnections for Listening socket.
|
||||
int connections = 0;
|
||||
|
||||
// The count of received event from eventhandler.
|
||||
int tokens = 0;
|
||||
|
||||
bool sendReadEvents = false;
|
||||
@@ -764,32 +770,21 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
}
|
||||
|
||||
Datagram? receive() {
|
||||
if (isClosing || isClosed) return null;
|
||||
if (isClosing || isClosed || !_availableDatagram) return null;
|
||||
try {
|
||||
Datagram? datagram = nativeRecvFrom();
|
||||
Datagram? result = nativeRecvFrom();
|
||||
assert(result != null);
|
||||
final resourceInformation = resourceInfo;
|
||||
assert(resourceInformation != null ||
|
||||
isPipe ||
|
||||
isInternal ||
|
||||
isInternalSignal);
|
||||
if (datagram != null) {
|
||||
// Read the next available. Available is only for the next datagram, not
|
||||
// the sum of all datagrams pending, so we need to call after each
|
||||
// receive. If available becomes > 0, the _NativeSocket will continue to
|
||||
// emit read events.
|
||||
available = nativeAvailable();
|
||||
if (resourceInformation != null) {
|
||||
resourceInformation.totalRead += datagram.data.length;
|
||||
}
|
||||
}
|
||||
if (resourceInformation != null) {
|
||||
resourceInformation.totalRead += result?.data?.length!;
|
||||
resourceInformation.didRead();
|
||||
}
|
||||
if (!const bool.fromEnvironment("dart.vm.product")) {
|
||||
_SocketProfile.collectStatistic(nativeGetSocketId(),
|
||||
_SocketProfileType.readBytes, datagram?.data.length);
|
||||
_SocketProfileType.readBytes, result?.data.length);
|
||||
}
|
||||
return datagram;
|
||||
_availableDatagram = nativeAvailableDatagram();
|
||||
return result;
|
||||
} catch (e) {
|
||||
reportError(e, StackTrace.current, "Receive failed");
|
||||
return null;
|
||||
@@ -933,7 +928,7 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
readEventIssued = false;
|
||||
if (isClosing) return;
|
||||
if (!sendReadEvents) return;
|
||||
if (available == 0) {
|
||||
if (stopRead()) {
|
||||
if (isClosedRead && !closedReadEventSent) {
|
||||
if (isClosedWrite) close();
|
||||
var handler = eventHandlers[closedEvent];
|
||||
@@ -953,6 +948,14 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
scheduleMicrotask(issue);
|
||||
}
|
||||
|
||||
bool stopRead() {
|
||||
if (isUdp) {
|
||||
return !_availableDatagram;
|
||||
} else {
|
||||
return available == 0;
|
||||
}
|
||||
}
|
||||
|
||||
void issueWriteEvent({bool delayed: true}) {
|
||||
if (writeEventIssued) return;
|
||||
if (!writeAvailable) return;
|
||||
@@ -1000,7 +1003,11 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
if (isListening) {
|
||||
connections++;
|
||||
} else {
|
||||
available = nativeAvailable();
|
||||
if (isUdp) {
|
||||
_availableDatagram = nativeAvailableDatagram();
|
||||
} else {
|
||||
available = nativeAvailable();
|
||||
}
|
||||
issueReadEvent();
|
||||
continue;
|
||||
}
|
||||
@@ -1266,6 +1273,7 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
|
||||
|
||||
void nativeSetSocketId(int id, int typeFlags) native "Socket_SetSocketId";
|
||||
int nativeAvailable() native "Socket_Available";
|
||||
bool nativeAvailableDatagram() native "Socket_AvailableDatagram";
|
||||
Uint8List? nativeRead(int len) native "Socket_Read";
|
||||
Datagram? nativeRecvFrom() native "Socket_RecvFrom";
|
||||
int nativeWrite(List<int> buffer, int offset, int bytes)
|
||||
|
||||
@@ -0,0 +1,31 @@
|
||||
// Copyright (c) 2020, the Dart project authors. Please see the AUTHORS file
|
||||
// for details. All rights reserved. Use of this source code is governed by a
|
||||
// BSD-style license that can be found in the LICENSE file.
|
||||
|
||||
import "dart:async";
|
||||
import "dart:io";
|
||||
import "dart:typed_data";
|
||||
|
||||
import "package:async_helper/async_helper.dart";
|
||||
import "package:expect/expect.dart";
|
||||
|
||||
main() async {
|
||||
asyncStart();
|
||||
var address = InternetAddress.loopbackIPv4;
|
||||
var sender = await RawDatagramSocket.bind(address, 0);
|
||||
var receiver = await RawDatagramSocket.bind(address, 0);
|
||||
|
||||
var sub;
|
||||
sub = receiver.listen((event) {
|
||||
if (event != RawSocketEvent.read) return;
|
||||
var datagram = receiver.receive();
|
||||
Expect.equals(0, datagram?.data?.length);
|
||||
Expect.isNull(receiver.receive());
|
||||
receiver.close();
|
||||
sub.cancel();
|
||||
asyncEnd();
|
||||
});
|
||||
|
||||
sender.send(Uint8List(0), address, receiver.port);
|
||||
sender.close();
|
||||
}
|
||||
@@ -0,0 +1,31 @@
|
||||
// Copyright (c) 2020, the Dart project authors. Please see the AUTHORS file
|
||||
// for details. All rights reserved. Use of this source code is governed by a
|
||||
// BSD-style license that can be found in the LICENSE file.
|
||||
|
||||
import "dart:async";
|
||||
import "dart:io";
|
||||
import "dart:typed_data";
|
||||
|
||||
import "package:async_helper/async_helper.dart";
|
||||
import "package:expect/expect.dart";
|
||||
|
||||
main() async {
|
||||
asyncStart();
|
||||
var address = InternetAddress.loopbackIPv4;
|
||||
var sender = await RawDatagramSocket.bind(address, 0);
|
||||
var receiver = await RawDatagramSocket.bind(address, 0);
|
||||
|
||||
var sub;
|
||||
sub = receiver.listen((event) {
|
||||
if (event != RawSocketEvent.read) return;
|
||||
var datagram = receiver.receive();
|
||||
Expect.equals(0, datagram.data.length);
|
||||
Expect.isNull(receiver.receive());
|
||||
receiver.close();
|
||||
sub.cancel();
|
||||
asyncEnd();
|
||||
});
|
||||
|
||||
sender.send(Uint8List(0), address, receiver.port);
|
||||
sender.close();
|
||||
}
|
||||
Reference in New Issue
Block a user