Chromium Code Reviews
chromiumcodereview-hr@appspot.gserviceaccount.com (chromiumcodereview-hr) | Please choose your nickname with Settings | Help | Chromium Project | Gerrit Changes | Sign out
(136)

Side by Side Diff: lib/isolate/dart2js/ports.dart

Issue 10828410: Unify dart:isolate. (Closed) Base URL: https://dart.googlecode.com/svn/branches/bleeding_edge/dart
Patch Set: Fix mapping to dart:isolate. Created 8 years, 4 months ago
Use n/p to move between diff chunks; N/P to move between comments. Draft comments are only viewable by you.
Jump to:
View unified diff | Download patch | Annotate | Revision Log
« no previous file with comments | « lib/isolate/dart2js/messages.dart ('k') | lib/isolate/dart2js/timer_provider.dart » ('j') | no next file with comments »
Toggle Intra-line Diffs ('i') | Expand Comments ('e') | Collapse Comments ('c') | Show Comments Hide Comments ('s')
OLDNEW
(Empty)
1 // Copyright (c) 2012, the Dart project authors. Please see the AUTHORS file
2 // for details. All rights reserved. Use of this source code is governed by a
3 // BSD-style license that can be found in the LICENSE file.
4
5 /** Common functionality to all send ports. */
6 class _BaseSendPort implements SendPort {
7 /** Id for the destination isolate. */
8 final int _isolateId;
9
10 const _BaseSendPort(this._isolateId);
11
12 void _checkReplyTo(SendPort replyTo) {
13 if (replyTo !== null
14 && replyTo is! _NativeJsSendPort
15 && replyTo is! _WorkerSendPort
16 && replyTo is! _BufferingSendPort) {
17 throw new Exception("SendPort.send: Illegal replyTo port type");
18 }
19 }
20
21 Future call(var message) {
22 final completer = new Completer();
23 final port = new _ReceivePortImpl();
24 send(message, port.toSendPort());
25 port.receive((value, ignoreReplyTo) {
26 port.close();
27 if (value is Exception) {
28 completer.completeException(value);
29 } else {
30 completer.complete(value);
31 }
32 });
33 return completer.future;
34 }
35
36 abstract void send(var message, [SendPort replyTo]);
37 abstract bool operator ==(var other);
38 abstract int hashCode();
39 }
40
41 /** A send port that delivers messages in-memory via native JavaScript calls. */
42 class _NativeJsSendPort extends _BaseSendPort implements SendPort {
43 final _ReceivePortImpl _receivePort;
44
45 const _NativeJsSendPort(this._receivePort, int isolateId) : super(isolateId);
46
47 void send(var message, [SendPort replyTo = null]) {
48 _waitForPendingPorts([message, replyTo], () {
49 _checkReplyTo(replyTo);
50 // Check that the isolate still runs and the port is still open
51 final isolate = _globalState.isolates[_isolateId];
52 if (isolate == null) return;
53 if (_receivePort._callback == null) return;
54
55 // We force serialization/deserialization as a simple way to ensure
56 // isolate communication restrictions are respected between isolates that
57 // live in the same worker. [_NativeJsSendPort] delivers both messages
58 // from the same worker and messages from other workers. In particular,
59 // messages sent from a worker via a [_WorkerSendPort] are received at
60 // [_processWorkerMessage] and forwarded to a native port. In such cases,
61 // here we'll see [_globalState.currentContext == null].
62 final shouldSerialize = _globalState.currentContext != null
63 && _globalState.currentContext.id != _isolateId;
64 var msg = message;
65 var reply = replyTo;
66 if (shouldSerialize) {
67 msg = _serializeMessage(msg);
68 reply = _serializeMessage(reply);
69 }
70 _globalState.topEventLoop.enqueue(isolate, () {
71 if (_receivePort._callback != null) {
72 if (shouldSerialize) {
73 msg = _deserializeMessage(msg);
74 reply = _deserializeMessage(reply);
75 }
76 _receivePort._callback(msg, reply);
77 }
78 }, 'receive $message');
79 });
80 }
81
82 bool operator ==(var other) => (other is _NativeJsSendPort) &&
83 (_receivePort == other._receivePort);
84
85 int hashCode() => _receivePort._id;
86 }
87
88 /** A send port that delivers messages via worker.postMessage. */
89 // TODO(eub): abstract this for iframes.
90 class _WorkerSendPort extends _BaseSendPort implements SendPort {
91 final int _workerId;
92 final int _receivePortId;
93
94 const _WorkerSendPort(this._workerId, int isolateId, this._receivePortId)
95 : super(isolateId);
96
97 void send(var message, [SendPort replyTo = null]) {
98 _waitForPendingPorts([message, replyTo], () {
99 _checkReplyTo(replyTo);
100 final workerMessage = _serializeMessage({
101 'command': 'message',
102 'port': this,
103 'msg': message,
104 'replyTo': replyTo});
105
106 if (_globalState.isWorker) {
107 // communication from one worker to another go through the main worker:
108 _globalState.mainManager.postMessage(workerMessage);
109 } else {
110 _globalState.managers[_workerId].postMessage(workerMessage);
111 }
112 });
113 }
114
115 bool operator ==(var other) {
116 return (other is _WorkerSendPort) &&
117 (_workerId == other._workerId) &&
118 (_isolateId == other._isolateId) &&
119 (_receivePortId == other._receivePortId);
120 }
121
122 int hashCode() {
123 // TODO(sigmund): use a standard hash when we get one available in corelib.
124 return (_workerId << 16) ^ (_isolateId << 8) ^ _receivePortId;
125 }
126 }
127
128 /** A port that buffers messages until an underlying port gets resolved. */
129 class _BufferingSendPort extends _BaseSendPort implements SendPort {
130 /** Internal counter to assign unique ids to each port. */
131 static int _idCount = 0;
132
133 /** For implementing equals and hashcode. */
134 final int _id;
135
136 /** Underlying port, when resolved. */
137 SendPort _port;
138
139 /**
140 * Future of the underlying port, so that we can detect when this port can be
141 * sent on messages.
142 */
143 Future<SendPort> _futurePort;
144
145 /** Pending messages (and reply ports). */
146 List pending;
147
148 _BufferingSendPort(isolateId, this._futurePort)
149 : super(isolateId), _id = _idCount, pending = [] {
150 _idCount++;
151 _futurePort.then((p) {
152 _port = p;
153 for (final item in pending) {
154 p.send(item['message'], item['replyTo']);
155 }
156 pending = null;
157 });
158 }
159
160 _BufferingSendPort.fromPort(isolateId, this._port)
161 : super(isolateId), _id = _idCount {
162 _idCount++;
163 }
164
165 void send(var message, [SendPort replyTo]) {
166 if (_port != null) {
167 _port.send(message, replyTo);
168 } else {
169 pending.add({'message': message, 'replyTo': replyTo});
170 }
171 }
172
173 bool operator ==(var other) =>
174 other is _BufferingSendPort && _id == other._id;
175 int hashCode() => _id;
176 }
177
178 /** Default factory for receive ports. */
179 class _ReceivePortFactory {
180
181 factory ReceivePort() {
182 return new _ReceivePortImpl();
183 }
184
185 }
186
187 /** Implementation of a multi-use [ReceivePort] on top of JavaScript. */
188 class _ReceivePortImpl implements ReceivePort {
189 int _id;
190 Function _callback;
191 static int _nextFreeId = 1;
192
193 _ReceivePortImpl()
194 : _id = _nextFreeId++ {
195 _globalState.currentContext.register(_id, this);
196 }
197
198 void receive(void onMessage(var message, SendPort replyTo)) {
199 _callback = onMessage;
200 }
201
202 void close() {
203 _callback = null;
204 _globalState.currentContext.unregister(_id);
205 }
206
207 SendPort toSendPort() {
208 return new _NativeJsSendPort(this, _globalState.currentContext.id);
209 }
210 }
211
212 /** Wait until all ports in a message are resolved. */
213 _waitForPendingPorts(var message, void callback()) {
214 final finder = new _PendingSendPortFinder();
215 finder.traverse(message);
216 Futures.wait(finder.ports).then((_) => callback());
217 }
218
219
220 /** Visitor that finds all unresolved [SendPort]s in a message. */
221 class _PendingSendPortFinder extends _MessageTraverser {
222 List<Future<SendPort>> ports;
223 _PendingSendPortFinder() : super(), ports = [] {
224 _visited = new _JsVisitedMap();
225 }
226
227 visitPrimitive(x) {}
228
229 visitList(List list) {
230 final seen = _visited[list];
231 if (seen !== null) return;
232 _visited[list] = true;
233 // TODO(sigmund): replace with the following: (bug #1660)
234 // list.forEach(_dispatch);
235 list.forEach((e) => _dispatch(e));
236 }
237
238 visitMap(Map map) {
239 final seen = _visited[map];
240 if (seen !== null) return;
241
242 _visited[map] = true;
243 // TODO(sigmund): replace with the following: (bug #1660)
244 // map.getValues().forEach(_dispatch);
245 map.getValues().forEach((e) => _dispatch(e));
246 }
247
248 visitSendPort(SendPort port) {
249 if (port is _BufferingSendPort && port._port == null) {
250 ports.add(port._futurePort);
251 }
252 }
253 }
OLDNEW
« no previous file with comments | « lib/isolate/dart2js/messages.dart ('k') | lib/isolate/dart2js/timer_provider.dart » ('j') | no next file with comments »

Powered by Google App Engine
This is Rietveld 408576698