fix(WebWorker): Add zone support to MessageBus

Closes #4053
This commit is contained in:
Jason Teplitz
2015-09-08 10:52:06 -07:00
parent 3b9c08676a
commit f3da37c92f
35 changed files with 628 additions and 365 deletions
@@ -1,25 +1,29 @@
library angular2.src.web_workers.debug_tools.multi_client_server_message_bus;
import "package:angular2/src/web_workers/shared/message_bus.dart"
show MessageBus, MessageBusSink, MessageBusSource;
import 'dart:io';
import 'dart:convert' show JSON;
import 'dart:async';
import 'package:angular2/src/core/facade/async.dart' show EventEmitter;
import 'package:angular2/src/web_workers/shared/messaging_api.dart';
import 'package:angular2/src/web_workers/shared/generic_message_bus.dart';
// TODO(jteplitz602): Remove hard coded result type and
// clear messageHistory once app is done with it #3859
class MultiClientServerMessageBus implements MessageBus {
final MultiClientServerMessageBusSink sink;
MultiClientServerMessageBusSource source;
class MultiClientServerMessageBus extends GenericMessageBus {
bool hasPrimary = false;
MultiClientServerMessageBus(this.sink, this.source);
@override
MultiClientServerMessageBusSink get sink => super.sink;
@override
MultiClientServerMessageBusSource get source => super.source;
MultiClientServerMessageBus(MultiClientServerMessageBusSink sink,
MultiClientServerMessageBusSource source)
: super(sink, source);
MultiClientServerMessageBus.fromHttpServer(HttpServer server)
: sink = new MultiClientServerMessageBusSink() {
source = new MultiClientServerMessageBusSource(resultReceived);
: super(new MultiClientServerMessageBusSink(),
new MultiClientServerMessageBusSource()) {
source.onResult.listen(_resultReceived);
server.listen((HttpRequest request) {
if (request.uri.path == "/ws") {
WebSocketTransformer.upgrade(request).then((WebSocket socket) {
@@ -38,18 +42,10 @@ class MultiClientServerMessageBus implements MessageBus {
});
}
void resultReceived() {
void _resultReceived(_) {
sink.resultReceived();
}
EventEmitter from(String channel) {
return source.from(channel);
}
EventEmitter to(String channel) {
return sink.to(channel);
}
Function _handleDisconnect(WebSocketWrapper wrapper) {
return () {
sink.removeConnection(wrapper);
@@ -72,12 +68,15 @@ class WebSocketWrapper {
WebSocketWrapper(this._messageHistory, this._resultMarkers, this.socket) {
stream = socket.asBroadcastStream();
stream.listen((encodedMessage) {
var message = JSON.decode(encodedMessage)['message'];
if (message is Map && message.containsKey("type")) {
if (message['type'] == 'result') {
resultReceived();
var messages = JSON.decode(encodedMessage);
messages.forEach((data) {
var message = data['message'];
if (message is Map && message.containsKey("type")) {
if (message['type'] == 'result') {
resultReceived();
}
}
}
});
});
}
@@ -121,10 +120,9 @@ class WebSocketWrapper {
}
}
class MultiClientServerMessageBusSink implements MessageBusSink {
class MultiClientServerMessageBusSink extends GenericMessageBusSink {
final List<String> messageHistory = new List<String>();
final Set<WebSocketWrapper> openConnections = new Set<WebSocketWrapper>();
final Map<String, EventEmitter> _channels = new Map<String, EventEmitter>();
final List<int> resultMarkers = new List<int>();
void resultReceived() {
@@ -141,76 +139,77 @@ class MultiClientServerMessageBusSink implements MessageBusSink {
openConnections.remove(webSocket);
}
EventEmitter to(String channel) {
if (_channels.containsKey(channel)) {
return _channels[channel];
} else {
var emitter = new EventEmitter();
emitter.listen((message) {
_send({'channel': channel, 'message': message});
});
return emitter;
}
}
void _send(dynamic message) {
String encodedMessage = JSON.encode(message);
@override
void sendMessages(List<dynamic> messages) {
String encodedMessages = JSON.encode(messages);
openConnections.forEach((WebSocketWrapper webSocket) {
if (webSocket.caughtUp) {
webSocket.socket.add(encodedMessage);
webSocket.socket.add(encodedMessages);
}
});
messageHistory.add(encodedMessage);
messageHistory.add(encodedMessages);
}
}
class MultiClientServerMessageBusSource implements MessageBusSource {
final Map<String, EventEmitter> _channels = new Map<String, EventEmitter>();
class MultiClientServerMessageBusSource extends GenericMessageBusSource {
Function onResultReceived;
final StreamController mainController;
final StreamController resultController = new StreamController();
MultiClientServerMessageBusSource(this.onResultReceived);
MultiClientServerMessageBusSource._(controller)
: mainController = controller,
super(controller.stream);
EventEmitter from(String channel) {
if (_channels.containsKey(channel)) {
return _channels[channel];
} else {
var emitter = new EventEmitter();
_channels[channel] = emitter;
return emitter;
}
factory MultiClientServerMessageBusSource() {
return new MultiClientServerMessageBusSource._(
new StreamController.broadcast());
}
Stream get onResult => resultController.stream;
void addConnection(WebSocketWrapper webSocket) {
if (webSocket.isPrimary) {
webSocket.stream.listen((encodedMessage) {
var decodedMessage = decodeMessage(encodedMessage);
var channel = decodedMessage['channel'];
var message = decodedMessage['message'];
if (message is Map && message.containsKey("type")) {
if (message['type'] == 'result') {
// tell the bus that a result was received on the primary
onResultReceived();
webSocket.stream.listen((encodedMessages) {
var decodedMessages = _decodeMessages(encodedMessages);
decodedMessages.forEach((decodedMessage) {
var message = decodedMessage['message'];
if (message is Map && message.containsKey("type")) {
if (message['type'] == 'result') {
// tell the bus that a result was received on the primary
resultController.add(message);
}
}
}
});
if (_channels.containsKey(channel)) {
_channels[channel].add(message);
}
mainController.add(decodedMessages);
});
} else {
webSocket.stream.listen((encodedMessage) {
// handle events from non-primary browser
var decodedMessage = decodeMessage(encodedMessage);
var channel = decodedMessage['channel'];
var message = decodedMessage['message'];
if (_channels.containsKey(EVENT_CHANNEL) && channel == EVENT_CHANNEL) {
_channels[channel].add(message);
webSocket.stream.listen((encodedMessages) {
// handle events from non-primary connection.
var decodedMessages = _decodeMessages(encodedMessages);
var eventMessages = new List<Map<String, dynamic>>();
decodedMessages.forEach((decodedMessage) {
var channel = decodedMessage['channel'];
if (channel == EVENT_CHANNEL) {
eventMessages.add(decodedMessage);
}
});
if (eventMessages.length > 0) {
mainController.add(eventMessages);
}
});
}
}
Map<String, dynamic> decodeMessage(dynamic message) {
return JSON.decode(message);
List<dynamic> _decodeMessages(dynamic messages) {
return JSON.decode(messages);
}
// This is a noop for the MultiClientBus because it has to decode the JSON messages before
// the generic bus receives them in order to check for results and forward events
// from the non-primary connection.
@override
List<dynamic> decodeMessages(dynamic messages) {
return messages;
}
}
@@ -1,22 +1,24 @@
library angular2.src.web_workers.debug_tools.single_client_server_message_bus;
import "package:angular2/src/web_workers/shared/message_bus.dart"
show MessageBus, MessageBusSink, MessageBusSource;
import 'dart:io';
import 'dart:convert' show JSON;
import 'dart:async';
import "package:angular2/src/core/facade/async.dart" show EventEmitter;
import 'package:angular2/src/web_workers/shared/generic_message_bus.dart';
class SingleClientServerMessageBus implements MessageBus {
final SingleClientServerMessageBusSink sink;
SingleClientServerMessageBusSource source;
class SingleClientServerMessageBus extends GenericMessageBus {
bool connected = false;
SingleClientServerMessageBus(this.sink, this.source);
@override
SingleClientServerMessageBusSink get sink => super.sink;
@override
SingleClientServerMessageBusSource get source => super.source;
SingleClientServerMessageBus(SingleClientServerMessageBusSink sink,
SingleClientServerMessageBusSource source)
: super(sink, source);
SingleClientServerMessageBus.fromHttpServer(HttpServer server)
: sink = new SingleClientServerMessageBusSink() {
source = new SingleClientServerMessageBusSource();
: super(new SingleClientServerMessageBusSink(),
new SingleClientServerMessageBusSource()) {
server.listen((HttpRequest request) {
if (request.uri.path == "/ws") {
if (!connected) {
@@ -24,7 +26,7 @@ class SingleClientServerMessageBus implements MessageBus {
sink.setConnection(socket);
var stream = socket.asBroadcastStream();
source.setConnectionFromStream(stream);
source.attachTo(stream);
stream.listen(null, onDone: _handleDisconnect);
}).catchError((error) {
throw error;
@@ -43,51 +45,30 @@ class SingleClientServerMessageBus implements MessageBus {
void _handleDisconnect() {
sink.removeConnection();
source.removeConnection();
connected = false;
}
EventEmitter from(String channel) {
return source.from(channel);
}
EventEmitter to(String channel) {
return sink.to(channel);
}
}
class SingleClientServerMessageBusSink implements MessageBusSink {
class SingleClientServerMessageBusSink extends GenericMessageBusSink {
final List<String> _messageBuffer = new List<String>();
WebSocket _socket = null;
final Map<String, EventEmitter> _channels = new Map<String, EventEmitter>();
void setConnection(WebSocket webSocket) {
_socket = webSocket;
_sendBufferedMessages();
}
EventEmitter to(String channel) {
if (_channels.containsKey(channel)) {
return _channels[channel];
} else {
var emitter = new EventEmitter();
emitter.listen((message) {
_send({'channel': channel, 'message': message});
});
return emitter;
}
}
void removeConnection() {
_socket = null;
}
void _send(dynamic message) {
String encodedMessage = JSON.encode(message);
@override
void sendMessages(List<dynamic> message) {
String encodedMessages = JSON.encode(message);
if (_socket != null) {
_socket.add(encodedMessage);
_socket.add(encodedMessages);
} else {
_messageBuffer.add(encodedMessage);
_messageBuffer.add(encodedMessages);
}
}
@@ -97,44 +78,11 @@ class SingleClientServerMessageBusSink implements MessageBusSink {
}
}
class SingleClientServerMessageBusSource implements MessageBusSource {
final Map<String, EventEmitter> _channels = new Map<String, EventEmitter>();
Stream _stream;
class SingleClientServerMessageBusSource extends GenericMessageBusSource {
SingleClientServerMessageBusSource() : super(null);
SingleClientServerMessageBusSource();
EventEmitter from(String channel) {
if (_channels.containsKey(channel)) {
return _channels[channel];
} else {
var emitter = new EventEmitter();
_channels[channel] = emitter;
return emitter;
}
}
void setConnectionFromWebSocket(WebSocket socket) {
setConnectionFromStream(socket.asBroadcastStream());
}
void setConnectionFromStream(Stream stream) {
_stream = stream;
_stream.listen((encodedMessage) {
var decodedMessage = decodeMessage(encodedMessage);
var channel = decodedMessage['channel'];
var message = decodedMessage['message'];
if (_channels.containsKey(channel)) {
_channels[channel].add(message);
}
});
}
void removeConnection() {
_stream = null;
}
Map<String, dynamic> decodeMessage(dynamic message) {
return JSON.decode(message);
@override
List<dynamic> decodeMessages(dynamic messages) {
return JSON.decode(messages);
}
}
@@ -2,77 +2,33 @@ library angular2.src.web_workers.worker.web_socket_message_bus;
import 'dart:html';
import 'dart:convert' show JSON;
import "package:angular2/src/web_workers/shared/message_bus.dart"
show MessageBus, MessageBusSink, MessageBusSource;
import 'package:angular2/src/core/facade/async.dart' show EventEmitter;
import 'package:angular2/src/web_workers/shared/generic_message_bus.dart';
class WebSocketMessageBus implements MessageBus {
final WebSocketMessageBusSink sink;
final WebSocketMessageBusSource source;
WebSocketMessageBus(this.sink, this.source);
class WebSocketMessageBus extends GenericMessageBus {
WebSocketMessageBus(
WebSocketMessageBusSink sink, WebSocketMessageBusSource source)
: super(sink, source);
WebSocketMessageBus.fromWebSocket(WebSocket webSocket)
: sink = new WebSocketMessageBusSink(webSocket),
source = new WebSocketMessageBusSource(webSocket);
EventEmitter from(String channel) {
return source.from(channel);
}
EventEmitter to(String channel) {
return sink.to(channel);
}
: super(new WebSocketMessageBusSink(webSocket),
new WebSocketMessageBusSource(webSocket));
}
class WebSocketMessageBusSink implements MessageBusSink {
class WebSocketMessageBusSink extends GenericMessageBusSink {
final WebSocket _webSocket;
final Map<String, EventEmitter> _channels = new Map<String, EventEmitter>();
WebSocketMessageBusSink(this._webSocket);
EventEmitter to(String channel) {
if (_channels.containsKey(channel)) {
return _channels[channel];
} else {
var emitter = new EventEmitter();
emitter.listen((message) {
_send({'channel': channel, 'message': message});
});
_channels[channel] = emitter;
return emitter;
}
}
void _send(message) {
_webSocket.send(JSON.encode(message));
void sendMessages(List<dynamic> messages) {
_webSocket.send(JSON.encode(messages));
}
}
class WebSocketMessageBusSource implements MessageBusSource {
final Map<String, EventEmitter> _channels = new Map<String, EventEmitter>();
class WebSocketMessageBusSource extends GenericMessageBusSource {
WebSocketMessageBusSource(WebSocket webSocket) : super(webSocket.onMessage);
WebSocketMessageBusSource(WebSocket webSocket) {
webSocket.onMessage.listen((MessageEvent encodedMessage) {
var message = decodeMessage(encodedMessage.data);
var channel = message['channel'];
if (_channels.containsKey(channel)) {
_channels[channel].add(message['message']);
}
});
}
EventEmitter from(String channel) {
if (_channels.containsKey(channel)) {
return _channels[channel];
} else {
var emitter = new EventEmitter();
_channels[channel] = emitter;
return emitter;
}
}
Map<String, dynamic> decodeMessage(dynamic message) {
return JSON.decode(message);
List<dynamic> decodeMessages(MessageEvent event) {
var messages = event.data;
return JSON.decode(messages);
}
}