Chromium Code Reviews| Index: runtime/bin/socket_patch.dart |
| diff --git a/runtime/bin/socket_patch.dart b/runtime/bin/socket_patch.dart |
| index 76289c505c7cf4e47e6f93479a3ed85dffa20d4f..a9d4d26a8d95ed635c539866cef6fe664be1e27c 100644 |
| --- a/runtime/bin/socket_patch.dart |
| +++ b/runtime/bin/socket_patch.dart |
| @@ -205,6 +205,129 @@ class _NetworkInterface implements NetworkInterface { |
| } |
| +class _Rate { |
| + final int buckets; |
| + final data; |
| + int lastValue = 0; |
| + int nextBucket = 0; |
| + |
| + _Rate(int buckets) : buckets = buckets, data = new List.filled(buckets, 0); |
| + |
| + void update(int value) { |
| + data[nextBucket] = value - lastValue; |
| + lastValue = value; |
| + nextBucket = (nextBucket + 1) % buckets; |
| + } |
| + |
| + int get rate { |
| + int sum = data.fold(0, (prev, element) => prev + element); |
| + return sum ~/ buckets; |
| + } |
| +} |
| + |
| +// Statics information for the observatory. |
| +class _SocketStat { |
| + _Rate readRate = new _Rate(5); |
| + _Rate writeRate = new _Rate(5); |
| + |
| + void update(_NativeSocket socket) { |
| + readRate.update(socket.totalRead); |
| + writeRate.update(socket.totalWritten); |
| + } |
| +} |
| + |
| +class _SocketsObservatory { |
| + static int socketCount = 0; |
| + static Map sockets = new Map<_NativeSocket, _SocketStat>(); |
| + static Timer timer; |
| + |
| + static add(_NativeSocket socket) { |
| + if (socketCount == 0) startTimer(); |
| + sockets[socket] = new _SocketStat(); |
| + socketCount++; |
| + } |
| + |
| + static remove(_NativeSicket socket) { |
| + sockets.remove(socket); |
| + socketCount--; |
| + if (socketCount == 0) stopTimer(); |
| + } |
| + |
| + static update(_) { |
| + sockets.forEach((socket, stat) { |
| + stat.update(socket); |
| + }); |
| + } |
| + |
| + static startTimer() { |
| + if (timer != null) return; |
| + timer = new Timer.periodic(new Duration(seconds: 1), update); |
| + } |
| + |
| + static stopTimer() { |
| + if (timer == null) return; |
| + timer.cancel(); |
| + timer = null; |
| + } |
| + |
| + static String generateResponse() { |
| + var response = new Map(); |
| + response['type'] = 'io'; |
|
Cutch
2014/02/19 15:56:02
This should probably be 'SocketList'.
Søren Gjesse
2014/02/19 16:02:41
Done.
|
| + var socketsStat = new List(); |
| + response['sockets'] = socketsStat; |
|
Cutch
2014/02/19 15:56:02
s/sockets/members
Søren Gjesse
2014/02/19 16:02:41
Done.
|
| + sockets.forEach((socket, stat) { |
| + var type = |
| + socket.isListening ? "LISTENING" : |
| + socket.isPipe ? "PIPE" : |
| + socket.isInternal ? "INTERNAL" : "NORMAL"; |
|
Cutch
2014/02/19 15:56:02
The "type" field has a special meaning in Observat
Søren Gjesse
2014/02/19 16:02:41
Done.
|
| + var protocol = |
| + socket.isTcp ? "tcp" : |
| + socket.isUdp ? "udp" : ""; |
| + var localAddress; |
| + var localPort; |
| + var remoteAddress; |
| + var remotePort; |
| + try { |
| + localAddress = socket.address.address; |
| + } catch (e) { |
| + localAddress = "n/a"; |
| + } |
| + try { |
| + localPort = socket.port; |
| + } catch (e) { |
| + localPort = "n/a"; |
| + } |
| + try { |
| + remoteAddress = socket.remoteAddress.address; |
| + } catch (e) { |
| + remoteAddress = "n/a"; |
| + } |
| + try { |
| + remotePort = socket.remotePort; |
| + } catch (e) { |
| + remotePort = "n/a"; |
| + } |
| + socketsStat.add({'type': type, 'protocol': protocol, |
| + 'localAddress': localAddress, 'localPort': localPort, |
| + 'remoteAddress': remoteAddress, 'remotePort': remotePort, |
| + 'totalRead': socket.totalRead, |
| + 'totalWritten': socket.totalWritten, |
| + 'readPerSec': stat.readRate.rate, |
| + 'writePerSec': stat.writeRate.rate}); |
| + }); |
| + return JSON.encode(response);; |
| + } |
| + |
| + static String toJSON() { |
| + try { |
| + return generateResponse(); |
| + } catch (e, s) { |
| + return '{"type":"Error","text":"$e","stacktrace":"$s"}'; |
| + } |
| + } |
| +} |
| + |
| + |
| // The _NativeSocket class encapsulates an OS socket. |
| class _NativeSocket extends NativeFieldWrapperClass1 { |
| // Bit flags used when communicating between the eventhandler and |
| @@ -238,6 +361,18 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| static const int TYPE_NORMAL_SOCKET = 0; |
| static const int TYPE_LISTENING_SOCKET = 1 << LISTENING_SOCKET; |
| static const int TYPE_PIPE = 1 << PIPE_SOCKET; |
| + static const int TYPE_TYPE_MASK = TYPE_LISTENING_SOCKET | PIPE_SOCKET; |
| + |
| + // Protocol flags. |
| + static const int TCP_SOCKET = 18; |
| + static const int UDP_SOCKET = 19; |
| + static const int INTERNAL_SOCKET = 20; |
| + static const int TYPE_TCP_SOCKET = 1 << TCP_SOCKET; |
| + static const int TYPE_UDP_SOCKET = 1 << UDP_SOCKET; |
| + static const int TYPE_INTERNAL_SOCKET = 1 << INTERNAL_SOCKET; |
| + static const int TYPE_PROTOCOL_MASK = |
| + TYPE_TCP_SOCKET | TYPE_UDP_SOCKET | TYPE_INTERNAL_SOCKET; |
| + |
| // Native port messages. |
| static const HOST_NAME_LOOKUP = 0; |
| @@ -277,6 +412,10 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| bool writeEventIssued = false; |
| bool writeAvailable = false; |
| + // Statistics. |
| + int totalRead = 0; |
| + int totalWritten = 0; |
| + |
| static Future<List<InternetAddress>> lookup( |
| String host, {InternetAddressType type: InternetAddressType.ANY}) { |
| return _IOService.dispatch(_SOCKET_LOOKUP, [host, type._value]) |
| @@ -431,19 +570,36 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| }); |
| } |
| - _NativeSocket.datagram(this.address) : typeFlags = TYPE_NORMAL_SOCKET; |
| + _NativeSocket.datagram(this.address) |
| + : typeFlags = TYPE_NORMAL_SOCKET | TYPE_UDP_SOCKET { |
| + _SocketsObservatory.add(this); |
| + } |
| - _NativeSocket.normal() : typeFlags = TYPE_NORMAL_SOCKET; |
| + _NativeSocket.normal() : typeFlags = TYPE_NORMAL_SOCKET | TYPE_TCP_SOCKET { |
| + _SocketsObservatory.add(this); |
| + } |
| - _NativeSocket.listen() : typeFlags = TYPE_LISTENING_SOCKET; |
| + _NativeSocket.listen() : typeFlags = TYPE_LISTENING_SOCKET | TYPE_TCP_SOCKET { |
| + _SocketsObservatory.add(this); |
| + } |
| - _NativeSocket.pipe() : typeFlags = TYPE_PIPE; |
| + _NativeSocket.pipe() : typeFlags = TYPE_PIPE { |
| + _SocketsObservatory.add(this); |
| + } |
| - _NativeSocket.watch(int id) : typeFlags = TYPE_NORMAL_SOCKET { |
| + _NativeSocket.watch(int id) |
| + : typeFlags = TYPE_NORMAL_SOCKET | TYPE_INTERNAL_SOCKET { |
| isClosedWrite = true; |
| nativeSetSocketId(id); |
| + _SocketsObservatory.add(this); |
| } |
| + bool get isListening => (typeFlags & TYPE_LISTENING_SOCKET) != 0; |
| + bool get isPipe => (typeFlags & TYPE_PIPE) != 0; |
| + bool get isInternal => (typeFlags & TYPE_INTERNAL_SOCKET) != 0; |
| + bool get isTcp => (typeFlags & TYPE_TCP_SOCKET) != 0; |
| + bool get isUdp => (typeFlags & TYPE_UDP_SOCKET) != 0; |
| + |
| List<int> read(int len) { |
| if (len != null && len <= 0) { |
| throw new ArgumentError("Illegal length $len"); |
| @@ -455,6 +611,7 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| return null; |
| } |
| if (result != null) available -= result.length; |
| + totalRead += result.length; |
| return result; |
| } |
| @@ -512,6 +669,7 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| } |
| // Negate the result, as stated above. |
| if (result < 0) result = -result; |
| + totalWritten += result; |
| return result; |
| } |
| @@ -535,9 +693,13 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| // Don't issue accept if we're closing. |
| if (isClosing || isClosed) return null; |
| var socket = new _NativeSocket.normal(); |
| - if (nativeAccept(socket) != true) return null; |
| + if (nativeAccept(socket) != true) { |
| + _SocketsObservatory.remove(socket); |
| + return null; |
| + } |
| socket.localPort = localPort; |
| socket.address = address; |
| + totalRead += 1; |
| return socket; |
| } |
| @@ -609,7 +771,7 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| if ((i == CLOSED_EVENT || i == READ_EVENT) && isClosedRead) continue; |
| if (isClosing && i != DESTROYED_EVENT) continue; |
| if (i == CLOSED_EVENT && |
| - typeFlags != TYPE_LISTENING_SOCKET && |
| + !isListening && |
| !isClosing && |
| !isClosed) { |
| isClosedRead = true; |
| @@ -623,8 +785,7 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| continue; |
| } |
| - if (i == READ_EVENT && |
| - typeFlags != TYPE_LISTENING_SOCKET) { |
| + if (i == READ_EVENT && !isListening) { |
| var avail = nativeAvailable(); |
| if (avail is int) { |
| available = avail; |
| @@ -644,6 +805,7 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| closeCompleter.complete(); |
| disconnectFromEventHandler(); |
| if (handler != null) handler(); |
| + _SocketsObservatory.remove(this); |
| continue; |
| } |
| @@ -674,7 +836,7 @@ class _NativeSocket extends NativeFieldWrapperClass1 { |
| if (read) issueReadEvent(); |
| if (write) issueWriteEvent(); |
| if (eventPort == null) { |
| - int flags = typeFlags; |
| + int flags = typeFlags & TYPE_TYPE_MASK; |
| if (!isClosedRead) flags |= 1 << READ_EVENT; |
| if (!isClosedWrite) flags |= 1 << WRITE_EVENT; |
| sendToEventHandler(flags); |
| @@ -1594,3 +1756,5 @@ Datagram _makeDatagram(List<int> data, |
| new _InternetAddress(address, null, in_addr), |
| port); |
| } |
| + |
| +String _socketsStats() => _SocketsObservatory.toJSON(); |