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

This reverts commit c326c587c5.

Reason for revert: Analyzer and Co19 test is broken.

https://logs.chromium.org/logs/dart/buildbucket/cr-buildbucket.appspot.com/8885365974080167440/+/steps/test_results/0/logs/new_test_failures__logs_/0

https://logs.chromium.org/logs/dart/buildbucket/cr-buildbucket.appspot.com/8885364415750725248/+/steps/test_results/0/logs/new_test_failures__logs_/0

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>

TBR=sortie@google.com,lrn@google.com,zra@google.com,asiva@google.com,zichangguo@google.com

Change-Id: Id84a7cafbc14d6353fa6af8e294e8b0c751a969f
No-Presubmit: true
No-Tree-Checks: true
No-Try: true
Bug: https://github.com/dart-lang/sdk/issues/39910
Reviewed-on: https://dart-review.googlesource.com/c/sdk/+/140009
Reviewed-by: Zichang Guo <zichangguo@google.com>
Commit-Queue: Zichang Guo <zichangguo@google.com>
This commit is contained in:
Zichang Guo
2020-03-19 23:05:59 +00:00
committed by commit-bot@chromium.org
parent 532f43dcb2
commit e91e09d1e4
14 changed files with 59 additions and 178 deletions
+2 -5
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()) {
if (!IsClosing() && !buffer->IsEmpty()) {
data_ready_ = pending_read_;
} else {
OverlappedBuffer::DisposeBuffer(buffer);
@@ -611,13 +611,10 @@ 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,7 +172,6 @@ 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
@@ -132,7 +132,6 @@ 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) \
+6 -18
View File
@@ -399,12 +399,17 @@ 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)) {
@@ -416,11 +421,6 @@ 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,18 +1015,6 @@ 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 -1
View File
@@ -102,7 +102,7 @@ void SocketAddress::SetAddrPort(RawAddr* addr, intptr_t port) {
}
intptr_t SocketAddress::GetAddrPort(const RawAddr& addr) {
if (addr.addr.sa_family == AF_INET) {
if (addr.ss.ss_family == AF_INET) {
return ntohs(addr.in.sin_port);
} else {
return ntohs(addr.in6.sin6_port);
-1
View File
@@ -167,7 +167,6 @@ 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);
+5 -9
View File
@@ -82,18 +82,14 @@ 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,
+5 -9
View File
@@ -82,18 +82,14 @@ 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,
+5 -9
View File
@@ -81,18 +81,14 @@ 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,
-7
View File
@@ -88,13 +88,6 @@ 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,
+14 -26
View File
@@ -364,18 +364,12 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
// Holds the address used to connect or bind the socket.
InternetAddress localAddress;
// 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.
// The number of available bytes to read.
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;
@@ -765,19 +759,26 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
}
Datagram receive() {
if (isClosing || isClosed || !_availableDatagram) return null;
if (isClosing || isClosed) return null;
try {
Datagram result = nativeRecvFrom();
assert(result != null);
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;
}
}
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");
@@ -912,7 +913,7 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
readEventIssued = false;
if (isClosing) return;
if (!sendReadEvents) return;
if (stopRead()) {
if (available == 0) {
if (isClosedRead && !closedReadEventSent) {
if (isClosedWrite) close();
var handler = eventHandlers[closedEvent];
@@ -932,14 +933,6 @@ 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;
@@ -987,11 +980,7 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
if (isListening) {
connections++;
} else {
if (isUdp) {
_availableDatagram = nativeAvailableDatagram();
} else {
available = nativeAvailable();
}
available = nativeAvailable();
issueReadEvent();
continue;
}
@@ -1250,7 +1239,6 @@ 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)
+21 -29
View File
@@ -364,18 +364,12 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
// Holds the address used to connect or bind the socket.
late InternetAddress localAddress;
// 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.
// The number of available bytes to read.
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;
@@ -770,21 +764,32 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
}
Datagram? receive() {
if (isClosing || isClosed || !_availableDatagram) return null;
if (isClosing || isClosed) return null;
try {
Datagram? result = nativeRecvFrom();
assert(result != null);
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();
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, result?.data.length);
_SocketProfileType.readBytes, datagram?.data.length);
}
_availableDatagram = nativeAvailableDatagram();
return result;
return datagram;
} catch (e) {
reportError(e, StackTrace.current, "Receive failed");
return null;
@@ -928,7 +933,7 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
readEventIssued = false;
if (isClosing) return;
if (!sendReadEvents) return;
if (stopRead()) {
if (available == 0) {
if (isClosedRead && !closedReadEventSent) {
if (isClosedWrite) close();
var handler = eventHandlers[closedEvent];
@@ -948,14 +953,6 @@ 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;
@@ -1003,11 +1000,7 @@ class _NativeSocket extends _NativeSocketNativeWrapper with _ServiceObject {
if (isListening) {
connections++;
} else {
if (isUdp) {
_availableDatagram = nativeAvailableDatagram();
} else {
available = nativeAvailable();
}
available = nativeAvailable();
issueReadEvent();
continue;
}
@@ -1273,7 +1266,6 @@ 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)
@@ -1,31 +0,0 @@
// 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();
}
@@ -1,31 +0,0 @@
// 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();
}