Reland "[dart:io] Fix hanging on zero-length datagram"

This is a reland of c326c587c5

The fix is to remove "_availableDatagram" check at the start of receive(). Since receive() can be called without receiver being listening, this check will block the receive().

Original change's description:
> [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>

Bug: https://github.com/dart-lang/sdk/issues/39910
Change-Id: Iefc0af96ed4e1b433396bb5fe11276e8a04c00d5
Reviewed-on: https://dart-review.googlesource.com/c/sdk/+/140164
Reviewed-by: Siva Annamalai <asiva@google.com>
Commit-Queue: Zichang Guo <zichangguo@google.com>
This commit is contained in:
Zichang Guo
2020-03-27 15:28:23 +00:00
committed by commit-bot@chromium.org
parent 4f0adf328a
commit fdf4762c99
13 changed files with 173 additions and 29 deletions
+5 -2
View File
@@ -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) {
+1
View File
@@ -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,
+1
View File
@@ -134,6 +134,7 @@ namespace bin {
V(ServerSocket_CreateUnixDomainBindListen, 5) \
V(SocketBase_IsBindError, 2) \
V(Socket_Available, 1) \
V(Socket_AvailableDatagram, 1) \
V(Socket_CreateBindConnect, 5) \
V(Socket_CreateUnixDomainBindConnect, 4) \
V(Socket_CreateBindDatagram, 6) \
+18 -1
View File
@@ -557,7 +557,7 @@ void FUNCTION_NAME(Socket_RecvFrom)(Dart_NativeArguments args) {
}
// 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)) {
@@ -569,6 +569,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);
// TODO(21403): Add checks for AF_UNIX, if unix domain sockets
@@ -1200,6 +1205,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) {
+1
View File
@@ -185,6 +185,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);
+9
View File
@@ -98,6 +98,15 @@ intptr_t SocketBase::RecvFrom(intptr_t fd,
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,
+9
View File
@@ -98,6 +98,15 @@ intptr_t SocketBase::RecvFrom(intptr_t fd,
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,
+9
View File
@@ -97,6 +97,15 @@ intptr_t SocketBase::RecvFrom(intptr_t fd,
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,
+7
View File
@@ -91,6 +91,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,
+23 -8
View File
@@ -425,12 +425,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;
@@ -845,11 +851,6 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
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;
}
@@ -861,6 +862,7 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
_SocketProfile.collectStatistic(nativeGetSocketId(),
_SocketProfileType.readBytes, result?.data?.length);
}
_availableDatagram = nativeAvailableDatagram();
return result;
} catch (e) {
reportError(e, StackTrace.current, "Receive failed");
@@ -1001,7 +1003,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];
@@ -1021,6 +1023,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;
@@ -1068,7 +1078,11 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
if (isListening) {
connections++;
} else {
available = nativeAvailable();
if (isUdp) {
_availableDatagram = nativeAvailableDatagram();
} else {
available = nativeAvailable();
}
issueReadEvent();
continue;
}
@@ -1327,6 +1341,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)
+30 -18
View File
@@ -427,12 +427,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;
@@ -850,30 +856,23 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
Datagram? receive() {
if (isClosing || isClosed) return null;
try {
Datagram? datagram = nativeRecvFrom();
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();
Datagram? result = nativeRecvFrom();
if (result != null) {
final resourceInformation = resourceInfo;
if (resourceInformation != null) {
resourceInformation.totalRead += datagram.data.length;
resourceInformation.totalRead += result.data.length;
}
}
final resourceInformation = resourceInfo;
if (resourceInformation != null) {
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;
@@ -1023,7 +1022,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];
@@ -1043,6 +1042,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;
@@ -1090,7 +1097,11 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
if (isListening) {
connections++;
} else {
available = nativeAvailable();
if (isUdp) {
_availableDatagram = nativeAvailableDatagram();
} else {
available = nativeAvailable();
}
issueReadEvent();
continue;
}
@@ -1356,6 +1367,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,30 @@
// 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) {
Expect.isNull(receiver.receive());
}
receiver.close();
sub.cancel();
asyncEnd();
});
sender.send(Uint8List(0), address, receiver.port);
sender.close();
}
@@ -0,0 +1,30 @@
// 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) {
Expect.isNull(receiver.receive());
}
receiver.close();
sub.cancel();
asyncEnd();
});
sender.send(Uint8List(0), address, receiver.port);
sender.close();
}