1/* 2 * Copyright (c) 2008-2018 the original author or authors. 3 * 4 * Licensed under the Apache License, Version 2.0 (the "License"); 5 * you may not use this file except in compliance with the License. 6 * You may obtain a copy of the License at 7 * 8 * http://www.apache.org/licenses/LICENSE-2.0 9 * 10 * Unless required by applicable law or agreed to in writing, software 11 * distributed under the License is distributed on an "AS IS" BASIS, 12 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. 13 * See the License for the specific language governing permissions and 14 * limitations under the License. 15 */ 16 17/* CometD Version ${project.version} */ 18 19(function(root, factory){ 20 if (typeof exports === 'object') { 21 // CommonJS. 22 module.exports = factory(); 23 } else if (typeof define === 'function' && define.amd) { 24 // AMD. 25 define([], factory); 26 } else { 27 // Globals. 28 root.org = root.org || {}; 29 root.org.cometd = factory(); 30 } 31}(this, function() { 32 /** 33 * Utility functions. 34 */ 35 var Utils = { 36 isString: function(value) { 37 if (value === undefined || value === null) { 38 return false; 39 } 40 return typeof value === 'string' || value instanceof String; 41 }, 42 isArray: function(value) { 43 if (value === undefined || value === null) { 44 return false; 45 } 46 return value instanceof Array; 47 }, 48 /** 49 * Returns whether the given element is contained into the given array. 50 * @param element the element to check presence for 51 * @param array the array to check for the element presence 52 * @return the index of the element, if present, or a negative index if the element is not present 53 */ 54 inArray: function(element, array) { 55 for (var i = 0; i < array.length; ++i) { 56 if (element === array[i]) { 57 return i; 58 } 59 } 60 return -1; 61 }, 62 setTimeout: function(cometd, funktion, delay) { 63 return window.setTimeout(function() { 64 try { 65 cometd._debug('Invoking timed function', funktion); 66 funktion(); 67 } catch (x) { 68 cometd._debug('Exception invoking timed function', funktion, x); 69 } 70 }, delay); 71 }, 72 clearTimeout: function(timeoutHandle) { 73 window.clearTimeout(timeoutHandle); 74 } 75 }; 76 77 78 /** 79 * A registry for transports used by the CometD object. 80 */ 81 var TransportRegistry = function() { 82 var _types = []; 83 var _transports = {}; 84 85 this.getTransportTypes = function() { 86 return _types.slice(0); 87 }; 88 89 this.findTransportTypes = function(version, crossDomain, url) { 90 var result = []; 91 for (var i = 0; i < _types.length; ++i) { 92 var type = _types[i]; 93 if (_transports[type].accept(version, crossDomain, url) === true) { 94 result.push(type); 95 } 96 } 97 return result; 98 }; 99 100 this.negotiateTransport = function(types, version, crossDomain, url) { 101 for (var i = 0; i < _types.length; ++i) { 102 var type = _types[i]; 103 for (var j = 0; j < types.length; ++j) { 104 if (type === types[j]) { 105 var transport = _transports[type]; 106 if (transport.accept(version, crossDomain, url) === true) { 107 return transport; 108 } 109 } 110 } 111 } 112 return null; 113 }; 114 115 this.add = function(type, transport, index) { 116 var existing = false; 117 for (var i = 0; i < _types.length; ++i) { 118 if (_types[i] === type) { 119 existing = true; 120 break; 121 } 122 } 123 124 if (!existing) { 125 if (typeof index !== 'number') { 126 _types.push(type); 127 } else { 128 _types.splice(index, 0, type); 129 } 130 _transports[type] = transport; 131 } 132 133 return !existing; 134 }; 135 136 this.find = function(type) { 137 for (var i = 0; i < _types.length; ++i) { 138 if (_types[i] === type) { 139 return _transports[type]; 140 } 141 } 142 return null; 143 }; 144 145 this.remove = function(type) { 146 for (var i = 0; i < _types.length; ++i) { 147 if (_types[i] === type) { 148 _types.splice(i, 1); 149 var transport = _transports[type]; 150 delete _transports[type]; 151 return transport; 152 } 153 } 154 return null; 155 }; 156 157 this.clear = function() { 158 _types = []; 159 _transports = {}; 160 }; 161 162 this.reset = function(init) { 163 for (var i = 0; i < _types.length; ++i) { 164 _transports[_types[i]].reset(init); 165 } 166 }; 167 }; 168 169 170 /** 171 * Base object with the common functionality for transports. 172 */ 173 var Transport = function() { 174 var _type; 175 var _cometd; 176 var _url; 177 178 /** 179 * Function invoked just after a transport has been successfully registered. 180 * @param type the type of transport (for example 'long-polling') 181 * @param cometd the cometd object this transport has been registered to 182 * @see #unregistered() 183 */ 184 this.registered = function(type, cometd) { 185 _type = type; 186 _cometd = cometd; 187 }; 188 189 /** 190 * Function invoked just after a transport has been successfully unregistered. 191 * @see #registered(type, cometd) 192 */ 193 this.unregistered = function() { 194 _type = null; 195 _cometd = null; 196 }; 197 198 this._debug = function() { 199 _cometd._debug.apply(_cometd, arguments); 200 }; 201 202 this._mixin = function() { 203 return _cometd._mixin.apply(_cometd, arguments); 204 }; 205 206 this.getConfiguration = function() { 207 return _cometd.getConfiguration(); 208 }; 209 210 this.getAdvice = function() { 211 return _cometd.getAdvice(); 212 }; 213 214 this.setTimeout = function(funktion, delay) { 215 return Utils.setTimeout(_cometd, funktion, delay); 216 }; 217 218 this.clearTimeout = function(handle) { 219 Utils.clearTimeout(handle); 220 }; 221 222 /** 223 * Converts the given response into an array of bayeux messages 224 * @param response the response to convert 225 * @return an array of bayeux messages obtained by convert
225ing the response 226 */ 227 this.convertToMessages = function(response) { 228 if (Utils.isString(response)) { 229 try { 230 return JSON.parse(response); 231 } catch (x) { 232 this._debug('Could not convert to JSON the following string', '"' + response + '"'); 233 throw x; 234 } 235 } 236 if (Utils.isArray(response)) { 237 return response; 238 } 239 if (response === undefined || response === null) { 240 return []; 241 } 242 if (response instanceof Object) { 243 return [response]; 244 } 245 throw 'Conversion Error ' + response + ', typeof ' + (typeof response); 246 }; 247 248 /** 249 * Returns whether this transport can work for the given version and cross domain communication case. 250 * @param version a string indicating the transport version 251 * @param crossDomain a boolean indicating whether the communication is cross domain 252 * @param url the URL to connect to 253 * @return true if this transport can work for the given version and cross domain communication case, 254 * false otherwise 255 */ 256 this.accept = function(version, crossDomain, url) { 257 throw 'Abstract'; 258 }; 259 260 /** 261 * Returns the type of this transport. 262 * @see #registered(type, cometd) 263 */ 264 this.getType = function() { 265 return _type; 266 }; 267 268 this.getURL = function() { 269 return _url; 270 }; 271 272 this.setURL = function(url) { 273 _url = url; 274 }; 275 276 this.send = function(envelope, metaConnect) { 277 throw 'Abstract'; 278 }; 279 280 this.reset = function(init) { 281 this._debug('Transport', _type, 'reset', init ? 'initial' : 'retry'); 282 }; 283 284 this.abort = function() { 285 this._debug('Transport', _type, 'aborted'); 286 }; 287 288 this.toString = function() { 289 return this.getType(); 290 }; 291 }; 292 293 Transport.derive = function(baseObject) { 294 function F() { 295 } 296 297 F.prototype = baseObject; 298 return new F(); 299 }; 300 301 302 /** 303 * Base object with the common functionality for transports based on requests. 304 * The key responsibility is to allow at most 2 outstanding requests to the server, 305 * to avoid that requests are sent behind a long poll. 306 * To achieve this, we have one reserved request for the long poll, and all other 307 * requests are serialized one after the other. 308 */ 309 var RequestTransport = function() { 310 var _super = new Transport(); 311 var _self = Transport.derive(_super); 312 var _requestIds = 0; 313 var _metaConnectRequest = null; 314 var _requests = []; 315 var _envelopes = []; 316 317 function _coalesceEnvelopes(envelope) { 318 while (_envelopes.length > 0) { 319 var envelopeAndRequest = _envelopes[0]; 320 var newEnvelope = envelopeAndRequest[0]; 321 var newRequest = envelopeAndRequest[1]; 322 if (newEnvelope.url === envelope.url && 323 newEnvelope.sync === envelope.sync) { 324 _envelopes.shift(); 325 envelope.messages = envelope.messages.concat(newEnvelope.messages); 326 this._debug('Coalesced', newEnvelope.messages.length, 'messages from request', newRequest.id); 327 continue; 328 } 329 break; 330 } 331 } 332 333 function _transportSend(envelope, request) { 334 this.transportSend(envelope, request); 335 request.expired = false; 336 337 if (!envelope.sync) { 338 var maxDelay = this.getConfiguration().maxNetworkDelay; 339 var delay = maxDelay; 340 if (request.metaConnect === true) { 341 delay += this.getAdvice().timeout; 342 } 343 344 this._debug('Transport', this.getType(), 'waiting at most', delay, 'ms for the response, maxNetworkDelay', maxDelay); 345 346 var self = this; 347 request.timeout = this.setTimeout(function() { 348 request.expired = true; 349 var errorMessage = 'Request ' + request.id + ' of transport ' + self.getType() + ' exceeded ' + delay + ' ms max network delay'; 350 var failure = { 351 reason: errorMessage 352 }; 353 var xhr = request.xhr; 354 failure.httpCode = self.xhrStatus(xhr); 355 self.abortXHR(xhr); 356 self._debug(errorMessage); 357 self.complete(request, false, request.metaConnect); 358 envelope.onFailure(xhr, envelope.messages, failure); 359 }, delay); 360 } 361 } 362 363 function _queueSend(envelope) { 364 var requestId = ++_requestIds; 365 var request = { 366 id: requestId, 367 metaConnect: false, 368 envelope: envelope 369 }; 370 371 // Consider the metaConnect requests which should always be present 372 if (_requests.length < this.getConfiguration().maxConnections - 1) { 373 _requests.push(request); 374 _transportSend.call(this, envelope, request); 375 } else { 376 this._debug('Transport', this.getType(), 'queueing request', requestId, 'envelope', envelope); 377 _envelopes.push([envelope, request]);
378 } 379 } 380 381 function _metaConnectComplete(request) { 382 var requestId = request.id; 383 this._debug('Transport', this.getType(), 'metaConnect complete, request', requestId); 384 if (_metaConnectRequest !== null && _metaConnectRequest.id !== requestId) { 385 throw 'Longpoll request mismatch, completing request ' + requestId; 386 } 387 388 // Reset metaConnect request 389 _metaConnectRequest = null; 390 } 391 392 function _complete(request, success) { 393 var index = Utils.inArray(request, _requests); 394 // The index can be negative if the request has been aborted 395 if (index >= 0) { 396 _requests.splice(index, 1); 397 } 398 399 if (_envelopes.length > 0) { 400 var envelopeAndRequest = _envelopes.shift(); 401 var nextEnvelope = envelopeAndRequest[0]; 402 var nextRequest = envelopeAndRequest[1]; 403 this._debug('Transport dequeued request', nextRequest.id); 404 if (success) { 405 if (this.getConfiguration().autoBatch) { 406 _coalesceEnvelopes.call(this, nextEnvelope); 407 } 408 _queueSend.call(this, nextEnvelope); 409 this._debug('Transport completed request', request.id, nextEnvelope); 410 } else { 411 // Keep the semantic of calling response callbacks asynchronously after the request 412 var self = this; 413 this.setTimeout(function() { 414 self.complete(nextRequest, false, nextRequest.metaConnect); 415 var failure = { 416 reason: 'Previous request failed' 417 }; 418 var xhr = nextRequest.xhr; 419 failure.httpCode = self.xhrStatus(xhr); 420 nextEnvelope.onFailure(xhr, nextEnvelope.messages, failure); 421 }, 0); 422 } 423 } 424 } 425 426 _self.complete = function(request, success, metaConnect) { 427 if (metaConnect) { 428 _metaConnectComplete.call(this, request); 429 } else { 430 _complete.call(this, request, success); 431 } 432 }; 433 434 /** 435 * Performs the actual send depending on the transport type details. 436 * @param envelope the envelope to send 437 * @param request the request information 438 */ 439 _self.transportSend = function(envelope, request) { 440 throw 'Abstract'; 441 }; 442 443 _self.transportSuccess = function(envelope, request, responses) { 444 if (!request.expired) { 445 this.clearTimeout(request.timeout); 446 this.complete(request, true, request.metaConnect); 447 if (responses && responses.length > 0) { 448 envelope.onSuccess(responses); 449 } else { 450 envelope.onFailure(request.xhr, envelope.messages, { 451 httpCode: 204 452 }); 453 } 454 } 455 }; 456 457 _self.transportFailure = function(envelope, request, failure) { 458 if (!request.expired) { 459 this.clearTimeout(request.timeout); 460 this.complete(request, false, request.metaConnect); 461 envelope.onFailure(request.xhr, envelope.messages, failure); 462 } 463 }; 464 465 function _metaConnectSend(envelope) { 466 if (_metaConnectRequest !== null) { 467 throw 'Concurrent metaConnect requests not allowed, request id=' + _metaConnectRequest.id + ' not yet completed'; 468 } 469 470 var requestId = ++_requestIds; 471 this._debug('Transport', this.getType(), 'metaConnect send, request', requestId, 'envelope', envelope); 472 var request = { 473 id: requestId, 474 metaConnect: true, 475 envelope: envelope 476 }; 477 _transportSend.call(this, envelope, request); 478 _metaConnectRequest = request; 479 } 480 481 _self.send = function(envelope, metaConnect) { 482 if (metaConnect) { 483 _metaConnectSend.call(this, envelope); 484 } else {
485 _queueSend.call(this, envelope); 486 } 487 }; 488 489 _self.abort = function() { 490 _super.abort(); 491 for (var i = 0; i < _requests.length; ++i) { 492 var request = _requests[i]; 493 if (request) { 494 this._debug('Aborting request', request); 495 if (!this.abortXHR(request.xhr)) { 496 this.transportFailure(request.envelope, request, {reason: 'abort'}); 497 } 498 } 499 } 500 var metaConnectRequest = _metaConnectRequest; 501 if (metaConnectRequest) { 502 this._debug('Aborting metaConnect request', metaConnectRequest); 503 if (!this.abortXHR(metaConnectRequest.xhr)) { 504 this.transportFailure(metaConnectRequest.envelope, metaConnectRequest, {reason: 'abort'}); 505 } 506 } 507 this.reset(true); 508 }; 509 510 _self.reset = function(init) { 511 _super.reset(init); 512 _metaConnectRequest = null; 513 _requests = []; 514 _envelopes = []; 515 }; 516 517 _self.abortXHR = function(xhr) { 518 if (xhr) { 519 try { 520 var state = xhr.readyState; 521 xhr.abort(); 522 return state !== window.XMLHttpRequest.UNSENT; 523 } catch (x) { 524 this._debug(x); 525 } 526 } 527 return false; 528 }; 529 530 _self.xhrStatus = function(xhr) { 531 if (xhr) { 532 try { 533 return xhr.status; 534 } catch (x) { 535 this._debug(x); 536 } 537 } 538 return -1; 539 }; 540 541 return _self; 542 }; 543 544 545 var LongPollingTransport = function() { 546 var _super = new RequestTransport(); 547 var _self = Transport.derive(_super); 548 // By default, support cross domain 549 var _supportsCrossDomain = true; 550 551 _self.accept = function(version, crossDomain, url) { 552 return _supportsCrossDomain || !crossDomain; 553 }; 554 555 _self.newXMLHttpRequest = function() { 556 return new window.XMLHttpRequest(); 557 }; 558 559 _self.xhrSend = function(packet) { 560 var xhr = _self.newXMLHttpRequest(); 561 // Copy external context, to be used in other environments. 562 xhr.context = _self.context; 563 xhr.withCredentials = true; 564 xhr.open('POST', packet.url, packet.sync !== true); 565 var headers = packet.headers; 566 if (headers) { 567 for (var headerName in headers) { 568 if (headers.hasOwnProperty(headerName)) { 569 xhr.setRequestHeader(headerName, headers[headerName]); 570 } 571 } 572 } 573 xhr.setRequestHeader('Content-Type', 'application/json;charset=UTF-8'); 574 xhr.onload = function() { 575 if (xhr.status === 200) { 576 packet.onSuccess(xhr.responseText); 577 } else { 578 packet.onError(xhr.statusText); 579 } 580 }; 581 xhr.onerror = function() { 582 packet.onError(xhr.statusText); 583 }; 584 xhr.send(packet.body); 585 return xhr; 586 }; 587 588 _self.transportSend = function(envelope, request) { 589 this._debug('Transport', this.getType(), 'sending request', request.id, 'envelope', envelope); 590 591 var self = this; 592 try { 593 var sameStack = true; 594 request.xhr = this.xhrSend({ 595 transport: this, 596 url: envelope.url, 597 sync: envelope.sync, 598 headers: this.getConfiguration().requestHeaders, 599 body: JSON.stringify(envelope.messages), 600 onSuccess: function(response) { 601 self._debug('Transport', self.getType(), 'received response', response); 602 var success = false; 603 try { 604 var received = self.convertToMessages(response); 605 if (received.length === 0) { 606 _supportsCrossDomain = false; 607 self.transportFailure(envelope, request, { 608 httpCode: 204 609 }); 610 } else { 611 success = true; 612 self.transportSuccess(envelope, request, received); 613 } 614 } catch (x) { 615 self._debug(x); 616 if (!success) { 617 _supportsCrossDomain = false; 618 var failure = { 619 exception: x 620 }; 621 failure.httpCode = self.xhrStatus(request.xhr); 622 self.transportFailure(envelope, request, failure); 623 } 624 } 625 }, 626 onError: function(reason, exception) { 627 self._debug('Transport', self.getType(), 'received error', reason, exception); 628 _supportsCrossDomain = false; 629 var failure = { 630 reason: reason, 631 exception: exception 632 }; 633 failure.httpCode = self.xhrStatus(request.xhr); 634 if (sameStack) { 635 // Keep the semantic of calling response callbacks asynchronously after the request 636 self.setTimeout(function() { 637 self.transportFailure(envelope, request, failure); 638 }, 0); 639 } else { 640 self.transportFailure(envelope, request, failure); 641 } 642 } 643 }); 644 sameStack = false; 645 } catch (x) { 646 _supportsCrossDomain = false; 647 // Keep the semantic of calling response callbacks asynchronously after the request 648 this.setTimeout(function() { 649 self.transportFailure(envelope, request, { 650 exception: x 651 }); 652 }, 0); 653 } 654 }; 655 656 _self.reset = function(init) { 657 _super.reset(init); 658 _supportsCrossDomain = true; 659 }; 660 661 return _self; 662 }; 663 664 665 var CallbackPollingTransport = function() { 666 var _super = new RequestTransport(); 667 var _self = Transport.derive(_super); 668 var jsonp = 0; 669 670 _self.accept = function(version, crossDomain, url) { 671 return true; 672 }; 673 674 _self.jsonpSend = function(packet) { 675 var head = document.getElementsByTagName('head')[0]; 676 var script = document.createElement('script'); 677 678 var callbackName = '_cometd_jsonp_' + jsonp++; 679 window[callbackName] = function(responseText) { 680 head.removeChild(script); 681 delete window[callbackName]; 682 packet.onSuccess(responseText); 683 }; 684 685 var url = packet.url; 686 url += url.indexOf('?') < 0 ? '?' : '&'; 687 url += 'jsonp=' + callbackName; 688 url += '&message=' + encodeURIComponent(packet.body); 689 script.src = url; 690 script.async = packet.sync !== true; 691 script.type = 'application/javascript'; 692 script.onerror = function(e) { 693 packet.onError('jsonp ' + e.type); 694 }; 695 head.appendChild(script); 696 }; 697 698 function _failTransportFn(envelope, request, x) { 699 var self = this; 700 return function() { 701 self.transportFailure(envelope, request, 'error', x); 702 }; 703 } 704 705 _self.transportSend = function(envelope, request) { 706 var self = this; 707 708 // Microsoft Internet Explorer has a 2083 URL max length 709 // We must ensure that we stay within that length 710 var start = 0; 711 var length = envelope.messages.length; 712 var lengths = []; 713 while (length > 0) { 714 // Encode the messages because all brackets, quotes, commas, colons, etc 715 // present in the JSON will be URL encoded, taking many more characters 716 var json = JSON.stringify(envelope.messages.slice(start, start + length)); 717 var urlLength = envelope.url.length + encodeURI(json).length; 718 719 var maxLength = this.getConfiguration().maxURILength; 720 if (urlLength > maxLength) { 721 if (length === 1) { 722 var x = 'Bayeux message too big (' + urlLength + ' bytes, max is ' + maxLength + ') ' + 723 'for transport ' + this.getType(); 724 // Keep the semantic of calling response callbacks asynchronously after the request 725 this.setTimeout(_failTransportFn.call(this, envelope, request, x), 0); 726 return; 727 } 728 729 --length; 730 continue; 731 } 732 733 lengths.push(length); 734 start += length; 735 length = envelope.messages.length - start; 736 } 737 738 // Here we are sure that the messages can be sent within the URL limit 739 740 var envelopeToSend = envelope; 741 if (lengths.length > 1) { 742 var begin = 0; 743 var end = lengths[0]; 744 this._debug('Transport', this.getType(), 'split', envelope.messages.length, 'messages into', lengths.join(' + ')); 745 envelopeToSend = this._mixin(false, {}, envelope); 746 envelopeToSend.messages = envelope.messages.slice(begin, end); 747 envelopeToSend.onSuccess = envelope.onSuccess; 748 envelopeToSend.onFailure = envelope.onFailure; 749 750 for (var i = 1; i < lengths.length; ++i) { 751 var nextEnvelope = this._mixin(false, {}, envelope); 752 begin = end; 753 end += lengths[i]; 754 nextEnvelope.messages = envelope.messages.slice(begin, end); 755 nextEnvelope.onSuccess = envelope.onSuccess; 756 nextEnvelope.onFailure = envelope.onFailure; 757 this.send(nextEnvelope, request.metaConnect); 758 } 759 } 760 761 this._debug('Transport', this.getType(), 'sending request', request.id, 'envelope', envelopeToSend); 762 763 try { 764 var sameStack = true; 765 this.jsonpSend({ 766 transport: this, 767 url: envelopeToSend.url, 768 sync: envelopeToSend.sync, 769 headers: this.getConfiguration().requestHeaders, 770 body: JSON.stringify(envelopeToSend.messages), 771 onSuccess: function(responses) { 772 var success = false; 773 try { 774 var received = self.convertToMessages(responses); 775 if (received.length === 0) { 776 self.transportFailure(envelopeToSend, request, { 777 httpCode: 204 778 }); 779 } else { 780 success = true; 781 self.transportSuccess(envelopeToSend, request, received); 782 } 783 } catch (x) { 784 self._debug(x); 785 if (!success) { 786 self.transportFailure(envelopeToSend, request, { 787 exception: x 788 }); 789 } 790 } 791 }, 792 onError: function(reason, exception) { 793 var failure = { 794 reason: reason, 795 exception: exception 796 }; 797 if (sameStack) { 798 // Keep the semantic of calling response callbacks asynchronously after the request 799 self.setTimeout(function() { 800 self.transportFailure(envelopeToSend, request, failure); 801 }, 0); 802 } else { 803 self.transportFailure(envelopeToSend, request, failure); 804 } 805 } 806 }); 807 sameStack = false; 808 } catch (xx) { 809 // Keep the semantic of calling response callbacks asynchronously after the request 810 this.setTimeout(function() { 811 self.transportFailure(envelopeToSend, request, { 812 exception: xx 813 }); 814 }, 0); 815 } 816 }; 817 818 return _self; 819 }; 820 821 822 var WebSocketTransport = function() { 823 var _super = new Transport(); 824 var _self = Transport.derive(_super); 825 var _cometd; 826 // By default WebSocket is supported 827 var _webSocketSupported = true; 828 // Whether we were able to establish a WebSocket connection 829 var _webSocketConnected = false; 830 var _stickyReconnect = true; 831 // The context contains the envelopes that have been sent 832 // and the timeouts for the messages that have been sent. 833 var _context = null; 834 var _connecting = null; 835 var _connected = false; 836 var _successCallback = null; 837 838 _self.reset = function(init) { 839 _super.reset(init); 840 _webSocketSupported = true; 841 if (init) { 842 _webSocketConnected = false; 843 } 844 _stickyReconnect = true; 845 _context = null; 846 _connecting = null; 847 _connected = false; 848 }; 849 850 function _forceClose(context, event) { 851 if (context) { 852 this.webSocketClose(context, event.code, event.reason); 853 // Force immediate failure of pending messages to trigger reconnect. 854 // This is needed because the server may not reply to our close() 855 // and therefore the onclose function is never called. 856 this.onClose(context, event); 857 } 858 } 859 860 function _sameContext(context) { 861 return context === _connecting || context === _context; 862 } 863 864 function _storeEnvelope(context, envelope, metaConnect) { 865 var messageIds = []; 866 for (var i = 0; i < envelope.messages.length; ++i) { 867 var message = envelope.messages[i]; 868 if (message.id) { 869 messageIds.push(message.id); 870 } 871 } 872 context.envelopes[messageIds.join(',')] = [envelope, metaConnect]; 873 this._debug('Transport', this.getType(), 'stored envelope, envelopes', context.envelopes
873); 874 } 875 876 function _websocketConnect(context) { 877 // We may have multiple attempts to open a WebSocket 878 // connection, for example a /meta/connect request that 879 // may take time, along with a user-triggered publish. 880 // Early return if we are already connecting. 881 if (_connecting) { 882 return; 883 } 884 885 // Mangle the URL, changing the scheme from 'http' to 'ws'. 886 var url = _cometd.getURL().replace(/^http/, 'ws'); 887 this._debug('Transport', this.getType(), 'connecting to URL', url); 888 889 try { 890 var protocol = _cometd.getConfiguration().protocol; 891 context.webSocket = protocol ? new window.WebSocket(url, protocol) : new window.WebSocket(url); 892 _connecting = context; 893 } catch (x) { 894 _webSocketSupported = false; 895 this._debug('Exception while creating WebSocket object', x); 896 throw x; 897 } 898 899 // By default use sticky reconnects. 900 _stickyReconnect = _cometd.getConfiguration().stickyReconnect !== false; 901 902 var self = this; 903 var connectTimeout = _cometd.getConfiguration().connectTimeout; 904 if (connectTimeout > 0) { 905 context.connectTimer = this.setTimeout(function() { 906 _cometd._debug('Transport', self.getType(), 'timed out while connecting to URL', url, ':', connectTimeout, 'ms'); 907 // The connection was not opened, close anyway. 908 _forceClose.call(self, context, {code: 1000, reason: 'Connect Timeout'}); 909 }, connectTimeout); 910 } 911 912 var onopen = function() { 913 _cometd._debug('WebSocket onopen', context); 914 if (context.connectTimer) { 915 self.clearTimeout(context.connectTimer); 916 } 917
918 if (_sameContext(context)) { 919 _connecting = null; 920 _context = context; 921 _webSocketConnected = true; 922 self.onOpen(context); 923 } else { 924 // We have a valid connection already, close this one. 925 _cometd._warn('Closing extra WebSocket connection', this, 'active connection', _context); 926 _forceClose.call(self, context, {code: 1000, reason: 'Extra Connection'}); 927 } 928 }; 929 930 // This callback is invoked when the server sends the close frame. 931 // The close frame for a connection may arrive *after* another 932 // connection has been opened, so we must make sure that actions 933 // are performed only if it's the same connection. 934 var onclose = function(event) { 935 event = event || {code: 1000}; 936 _cometd._debug('WebSocket onclose', context, event, 'connecting', _connecting, 'current', _context); 937 938 if (context.connectTimer) { 939 self.clearTimeout(context.connectTimer); 940 } 941 942 self.onClose(context, event); 943 }; 944 945 var onmessage = function(wsMessage) { 946 _cometd._debug('WebSocket onmessage', wsMessage, context); 947 self.onMessage(context, wsMessage); 948 }; 949 950 context.webSocket.onopen = onopen; 951 context.webSocket.onclose = onclose; 952 context.webSocket.onerror = function() { 953 // Clients should call onclose(), but if they do not we do it here for safety. 954 onclose({code: 1000, reason: 'Error'}); 955 }; 956 context.webSocket.onmessage = onmessage; 957 958 this._debug('Transport', this.getType(), 'configured callbacks on', context); 959 } 960 961 function _webSocketSend(context, envelope, metaConnect) { 962 var json = JSON.stringify(envelope.messages); 963 context.webSocket.send(json); 964 this._debug('Transport', this.getType(), 'sent', envelope, 'metaConnect =', metaConnect); 965 966 // Manage the timeout waiting for the response. 967 var maxDelay = this.getConfiguration().maxNetworkDelay; 968 var delay = maxDelay; 969 if (metaConnect) { 970 delay += this.getAdvice().timeout; 971 _connected = true; 972 } 973 974 var self = this; 975 var messageIds = []; 976 for (var i = 0; i < envelope.messages.length; ++i) { 977 (function() { 978 var message = envelope.messages[i]; 979 if (message.id) { 980 messageIds.push(message.id); 981 context.timeouts[message.id] = self.setTimeout(function() { 982 _cometd._debug('Transport', self.getType(), 'timing out message', message.id, 'after', delay, 'on', context); 983 _forceClose.call(self, context, {code: 1000, reason: 'Message Timeout'}); 984 }, delay); 985 } 986 })(); 987 } 988 989 this._debug('Transport', this.getType(), 'waiting at most', delay, 'ms for messages', messageIds, 'maxNetworkDelay', maxDelay, ', timeouts:', context.timeouts); 990 } 991 992 _self._notifySuccess = function(fn, messages) { 993 fn.call(this, messages); 994 }; 995 996 _self._notifyFailure = function(fn, context, messages, failure) { 997 fn.call(this, context, messages, failure); 998 }; 999 1000 function _send(context, envelope, metaConnect) { 1001 try { 1002 if (context === null) { 1003 context = _connecting || { 1004 envelopes: {}, 1005 timeouts: {} 1006 }; 1007 _storeEnvelope.call(this, context, envelope, metaConnect); 1008 _websocketConnect.call(this, context); 1009 } else { 1010 _storeEnvelope.call(this, context, envelope, metaConnect); 1011 _webSocketSend.call(this, context, envelope, metaConnect); 1012 } 1013 } catch (x) { 1014 // Keep the semantic of calling response callbacks asynchronously after the request. 1015 var self = this; 1016 this.setTimeout(function() { 1017 _forceClose.call(self, context, { 1018 code: 1000, 1019 reason: 'Exception', 1020 exception: x 1021 }); 1022 }, 0); 1023 } 1024 } 1025 1026 _self.onOpen = function(context) { 1027 var envelopes = context.envelopes; 1028 this._debug('Transport', this.getType(), 'opened', context, 'pending messages', envelopes
1028); 1029 for (var key in envelopes) { 1030 if (envelopes.hasOwnProperty(key)) { 1031 var element = envelopes[key]; 1032 var envelope = element[0]; 1033 var metaConnect = element[1]; 1034 // Store the success callback, which is independent from the envelope, 1035 // so that it can be used to notify arrival of messages. 1036 _successCallback = envelope.onSuccess; 1037 _webSocketSend.call(this, context, envelope, metaConnect); 1038 } 1039 } 1040 }; 1041 1042 _self.onMessage = function(context, wsMessage) { 1043 this._debug('Transport', this.getType(), 'received websocket message', wsMessage, context); 1044 1045 var close = false; 1046 var messages = this.convertToMessages(wsMessage.data); 1047 var messageIds = []; 1048 for (var i = 0; i < messages.length; ++i) { 1049 var message = messages[i]; 1050 1051 // Detect if the message is a response to a request we made. 1052 // If it's a meta message, for sure it's a response; otherwise it's 1053 // a publish message and publish responses don't have the data field. 1054 if (/^\/meta\//.test(message.channel) || message.data === undefined) { 1055 if (message.id) { 1056 messageIds.push(message.id); 1057 1058 var timeout = context.timeouts[message.id]; 1059 if (timeout) { 1060 this.clearTimeout(timeout); 1061 delete context.timeouts[message.id]; 1062 this._debug('Transport', this.getType(), 'removed timeout for message', message.id, ', timeouts', context.timeouts); 1063 } 1064 } 1065 } 1066 1067 if ('/meta/connect' === message.channel) { 1068 _connected = false; 1069 } 1070 if ('/meta/disconnect' === message.channel && !_connected) { 1071 close = true; 1072 } 1073 } 1074 1075 // Remove the envelope corresponding to the messages. 1076 var removed = false; 1077 var envelopes = context.envelopes; 1078 for (var j = 0; j < messageIds.length; ++j) { 1079 var id = messageIds[j]; 1080 for (var key in envelopes) { 1081 if (envelopes.hasOwnProperty(key)) { 1082 var ids = key.split(','); 1083 var index = Utils.inArray(id, ids); 1084 if (index >= 0) { 1085 removed = true; 1086 ids.splice(index, 1); 1087 var envelope = envelopes[key][0]; 1088 var metaConnect = envelopes[key][1]; 1089 delete envelopes[key]; 1090 if (ids.length > 0) { 1091 envelopes[ids.join(',')] = [envelope, metaConnect]; 1092 } 1093 break; 1094 } 1095 } 1096 } 1097 } 1098 if (removed) { 1099 this._debug('Transport', this.getType(), 'removed envelope, envelopes', envelopes); 1100 } 1101 1102 this._notifySuccess(_successCallback, messages); 1103 1104 if (close) { 1105 this.webSocketClose(context, 1000, 'Disconnect'); 1106 } 1107 }; 1108 1109 _self.onClose = function(context, event) { 1110 this._debug('Transport', this.getType(), 'closed', context, event); 1111
1112 if (_sameContext(context)) { 1113 // Remember if we were able to connect. 1114 // This close event could be due to server shutdown, 1115 // and if it restarts we want to try websocket again. 1116 _webSocketSupported = _stickyReconnect && _webSocketConnected; 1117 _connecting = null; 1118 _context = null; 1119 } 1120 1121 var timeouts = context.timeouts; 1122 context.timeouts = {}; 1123 for (var id in timeouts) { 1124 if (timeouts.hasOwnProperty(id)) { 1125 this.clearTimeout(timeouts[id]); 1126 } 1127 } 1128 1129 var envelopes = context.envelopes; 1130 context.envelopes = {}; 1131 for (var key in envelopes) { 1132 if (envelopes.hasOwnProperty(key)) { 1133 var envelope = envelopes[key][0]; 1134 var metaConnect = envelopes[key][1]; 1135 if (metaConnect) { 1136 _connected = false; 1137 } 1138 var failure = { 1139 websocketCode: event.code, 1140 reason: event.reason 1141 }; 1142 if (event.exception) { 1143 failure.exception = event.exception; 1144 } 1145 this._notifyFailure(envelope.onFailure, context, envelope.messages, failure); 1146 } 1147 } 1148 }; 1149 1150 _self.registered = function(type, cometd) { 1151 _super.registered(type, cometd); 1152 _cometd = cometd; 1153 }; 1154 1155 _self.accept = function(version, crossDomain, url) { 1156 this._debug('Transport', this.getType(), 'accept, supported:', _webSocketSupported); 1157 // Using !! to return a boolean (and not the WebSocket object). 1158 return _webSocketSupported && !!window.WebSocket && _cometd.websocketEnabled !== false; 1159 }; 1160 1161 _self.send = function(envelope, metaConnect) { 1162 this._debug('Transport', this.getType(), 'sending', envelope, 'metaConnect =', metaConnect); 1163 _send.call(this, _context, envelope, metaConnect); 1164 }; 1165 1166 _self.webSocketClose = function(context, code, reason) { 1167 try { 1168 if (context.webSocket) { 1169 context.webSocket.close(code, reason); 1170 } 1171 } catch (x) { 1172 this._debug(x); 1173 } 1174 }; 1175 1176 _self.abort = function() { 1177 _super.abort(); 1178 _forceClose.call(this, _context, {code: 1000, reason: 'Abort'}); 1179 this.reset(true); 1180 }; 1181 1182 return _self; 1183 }; 1184 1185 1186 /** 1187 * The constructor for a CometD object, identified by an optional name. 1188 * The default name is the string 'default'. 1189 * @param name the optional name of this cometd object 1190 */ 1191 var CometD = function (name) { 1192 var _cometd = this; 1193 var _name = name || 'default'; 1194 var _crossDomain = false; 1195 var _transports = new TransportRegistry(); 1196 var _transport; 1197 var _status = 'disconnected'; 1198 var _messageId = 0; 1199 var _clientId = null; 1200 var _batch = 0; 1201 var _messageQueue = []; 1202 var _internalBatch = false; 1203 var _listenerId = 0; 1204 var _listeners = {}; 1205 var _backoff = 0; 1206 var _scheduledSend = null; 1207 var _extensions = []; 1208 var _advice = {}; 1209 var _handshakeProps; 1210 var _handshakeCallback; 1211 var _callbacks = {}; 1212 var _remoteCalls = {}; 1213 var _reestablish = false; 1214 var _connected = false; 1215 var _unconnectTime = 0; 1216 var _handshakeMessages = 0; 1217 var _config = { 1218 protocol: null, 1219 stickyReconnect: true, 1220 connectTimeout: 0, 1221 maxConnections: 2, 1222 backoffIncrement: 1000, 1223 maxBackoff: 60000, 1224 logLevel: 'info', 1225 maxNetworkDelay: 10000, 1226 requestHeaders: {}, 1227 appendMessageTypeToURL: true, 1228 autoBatch: false, 1229 urls: {}, 1230 maxURILength: 2000, 1231 advice: { 1232 timeout: 60000, 1233 interval: 0, 1234 reconnect: undefined, 1235 maxInterval: 0 1236 } 1237 }; 1238 1239 function _fieldValue(object, name) { 1240 try { 1241 return object[name]; 1242 } catch (x) { 1243 return undefined; 1244 } 1245 } 1246 1247 /** 1248 * Mixes in the given objects into the target object by copying the properties. 1249 * @param deep if the copy must be deep 1250 * @param target the target object 1251 * @param objects the objects whose properties are copied into the target 1252 */ 1253 this._mixin = function(deep, target, objects) { 1254 var result = target || {}; 1255 1256 // Skip first 2 parameters (deep and target), and loop over the others 1257 for (var i = 2; i < arguments.length; ++i) { 1258 var object = arguments[i]; 1259 1260 if (object === undefined || object === null) { 1261 continue; 1262 } 1263 1264 for (var propName in object) { 1265 if (object.hasOwnProperty(propName)) { 1266 var prop = _fieldValue(object, propName); 1267 var targ = _fieldValue(result, propName); 1268 1269 // Avoid infinite loops 1270 if (prop === target) { 1271 continue; 1272 }
1273 // Do not mixin undefined values 1274 if (prop === undefined) { 1275 continue; 1276 } 1277 1278 if (deep && typeof prop === 'object' && prop !== null) { 1279 if (prop instanceof Array) { 1280 result[propName] = this._mixin(deep, targ instanceof Array ? targ : [], prop); 1281 } else { 1282 var source = typeof targ === 'object' && !(targ instanceof Array) ? targ : {}; 1283 result[propName] = this._mixin(deep, source, prop); 1284 } 1285 } else { 1286 result[propName] = prop; 1287 } 1288 } 1289 } 1290 } 1291 1292 return result; 1293 }; 1294 1295 function _isString(value) { 1296 return Utils.isString(value); 1297 } 1298 1299 function _isFunction(value) { 1300 if (value === undefined || value === null) { 1301 return false; 1302 } 1303 return typeof value === 'function'; 1304 } 1305 1306 function _zeroPad(value, length) { 1307 var result = ''; 1308 while (--length > 0) { 1309 if (value >= Math.pow(10, length)) { 1310 break; 1311 } 1312 result += '0'; 1313 } 1314 result += value; 1315 return result; 1316 } 1317 1318 function _log(level, args) { 1319 if (window.console) { 1320 var logger = window.console[level]; 1321 if (_isFunction(logger)) { 1322 var now = new Date(); 1323 [].splice.call(args, 0, 0, _zeroPad(now.getHours(), 2) + ':' + _zeroPad(now.getMinutes(), 2) + ':' + 1324 _zeroPad(now.getSeconds(), 2) + '.' + _zeroPad(now.getMilliseconds(), 3)); 1325 logger.apply(window.console, args); 1326 } 1327 } 1328 } 1329 1330 this._warn = function() { 1331 _log('warn', arguments); 1332 }; 1333 1334 this._info = function() { 1335 if (_config.logLevel !== 'warn') { 1336 _log('info', arguments); 1337 } 1338 }; 1339 1340 this._debug = function() { 1341 if (_config.logLevel === 'debug') { 1342 _log('debug', arguments); 1343 } 1344 }; 1345 1346 function _splitURL(url) { 1347 // [1] = protocol://, 1348 // [2] = host:port, 1349 // [3] = host, 1350 // [4] = IPv6_host, 1351 // [5] = IPv4_host, 1352 // [6] = :port, 1353 // [7] = port, 1354 // [8] = uri, 1355 // [9] = rest (query / fragment) 1356 return /(^https?:\/\/)?(((\[[^\]]+\])|([^:\/\?#]+))(:(\d+))?)?([^\?#]*)(.*)?/.exec(url); 1357 } 1358 1359 /** 1360 * Returns whether the given hostAndPort is cross domain. 1361 * The default implementation checks against window.location.host 1362 * but this function can be overridden to make it work in non-browser 1363 * environments. 1364 * 1365 * @param hostAndPort the host and port in format host:port 1366 * @return whether the given hostAndPort is cross domain 1367 */ 1368 this._isCrossDomain = function(hostAndPort) { 1369 if (window.location && window.location.host) { 1370 if (hostAndPort) { 1371 return hostAndPort !== window.location.host; 1372 } 1373 } 1374 return false; 1375 }; 1376 1377 function _configure(configuration) { 1378 _cometd._debug('Configuring cometd object with', configuration); 1379 // Support old style param, where only the Bayeux server URL was passed 1380 if (_isString(configuration)) { 1381 configuration = { url: configuration }; 1382 } 1383 if (!configuration) { 1384 configuration = {}; 1385 } 1386 1387 _config = _cometd._mixin(false, _config, configuration); 1388 1389 var url = _cometd.getURL(); 1390 if (!url) { 1391 throw 'Missing required configuration parameter \'url\' specifying the Bayeux server URL'; 1392 } 1393 1394 // Check if we're cross domain. 1395 var urlParts = _splitURL(url); 1396 var hostAndPort = urlParts[2]; 1397 var uri = urlParts[8]; 1398 var afterURI = urlParts[9];
1399 _crossDomain = _cometd._isCrossDomain(hostAndPort); 1400 1401 // Check if appending extra path is supported 1402 if (_config.appendMessageTypeToURL) { 1403 if (afterURI !== undefined && afterURI.length > 0) { 1404 _cometd._info('Appending message type to URI ' + uri + afterURI + ' is not supported, disabling \'appendMessageTypeToURL\' configuration'); 1405 _config.appendMessageTypeToURL = false; 1406 } else { 1407 var uriSegments = uri.split('/'); 1408 var lastSegmentIndex = uriSegments.length - 1; 1409 if (uri.match(/\/$/)) { 1410 lastSegmentIndex -= 1; 1411 } 1412 if (uriSegments[lastSegmentIndex].indexOf('.') >= 0) { 1413 // Very likely the CometD servlet's URL pattern is mapped to an extension, such as *.cometd 1414 // It will be difficult to add the extra path in this case 1415 _cometd._info('Appending message type to URI ' + uri + ' is not supported, disabling \'appendMessageTypeToURL\' configuration'); 1416 _config.appendMessageTypeToURL = false; 1417 } 1418 } 1419 } 1420 } 1421 1422 function _removeListener(subscription) { 1423 if (subscription) { 1424 var subscriptions = _listeners[subscription.channel]; 1425 if (subscriptions && subscriptions[subscription.id]) { 1426 delete subscriptions[subscription.id]; 1427 _cometd._debug('Removed', subscription.listener ? 'listener' : 'subscription', subscription); 1428 } 1429 } 1430 } 1431 1432 function _removeSubscription(subscription) { 1433 if (subscription && !subscription.listener) { 1434 _removeListener(subscription); 1435 } 1436 } 1437 1438 function _clearSubscriptions() { 1439 for (var channel in _listeners) { 1440 if (_listeners.hasOwnProperty(channel)) { 1441 var subscriptions = _listeners[channel]; 1442 if (subscriptions) { 1443 for (var id in subscriptions) { 1444 if (subscriptions.hasOwnProperty(id)) { 1445 _removeSubscription(subscriptions[id]); 1446 } 1447 } 1448 } 1449 } 1450 } 1451 } 1452 1453 function _setStatus(newStatus) { 1454 if (_status !== newStatus) { 1455 _cometd._debug('Status', _status, '->', newStatus); 1456 _status = newStatus; 1457 } 1458 } 1459 1460 function _isDisconnected() { 1461 return _status === 'disconnecting' || _status === 'disconnected'; 1462 } 1463 1464 function _nextMessageId() { 1465 var result = ++_messageId; 1466 return '' + result; 1467 } 1468 1469 function _applyExtension(scope, callback, name, message, outgoing) { 1470 try { 1471 return callback.call(scope, message); 1472 } catch (x) { 1473 var handler = _cometd.onExtensionException; 1474 if (_isFunction(handler)) { 1475 _cometd._debug('Invoking extension exception handler', name, x); 1476 try { 1477 handler.call(_cometd, x, name, outgoing, message); 1478 } catch (xx) { 1479 _cometd._info('Exception during execution of extension exception handler', name, xx); 1480 } 1481 } else { 1482 _cometd._info('Exception during execution of extension', name, x); 1483 } 1484 return message; 1485 } 1486 } 1487 1488 function _applyIncomingExtensions(message) { 1489 for (var i = 0; i < _extensions.length; ++i) { 1490 if (message === undefined || message === null) { 1491 break; 1492 } 1493 1494 var extension = _extensions[i]; 1495 var callback = extension.extension.incoming; 1496 if (_isFunction(callback)) { 1497 var result = _applyExtension(extension.extension, callback, extension.name, message, false); 1498 message = result === undefined ? message : result; 1499 } 1500 } 1501 return message; 1502 } 1503 1504 function _applyOutgoingExtensions(message) { 1505 for (var i = _extensions.length - 1; i >= 0 ; --i) { 1506 if (message === undefined || message === null) { 1507 break; 1508 } 1509 1510 var extension = _extensions[i]; 1511 var callback = extension.extension.outgoing; 1512 if (_isFunction(callback)) { 1513 var result = _applyExtension(extension.extension, callback, extension.name, message, true); 1514 message = result === undefined ? message : result; 1515 } 1516 } 1517 return message; 1518 } 1519 1520 function _notify(channel, message) { 1521 var subscriptions = _listeners[channel]; 1522 if (subscriptions) { 1523 for (var id in subscriptions) { 1524 if (subscriptions.hasOwnProperty(id)) { 1525 var subscription = subscriptions[id]; 1526 // Subscriptions may come and go, so the array may have 'holes' 1527 if (subscription) { 1528 try { 1529 subscription.callback.call(subscription.scope, message); 1530 } catch (x) { 1531 var handler = _cometd.onListenerException; 1532 if (_isFunction(handler)) { 1533 _cometd._debug('Invoking listener exception handler', subscription, x); 1534 try { 1535 handler.call(_cometd, x, subscription, subscription.listener, message); 1536 } catch (xx) { 1537 _cometd._info('Exception during execution of listener exception handler', subscription, xx); 1538 } 1539 } else { 1540 _cometd._info('Exception during execution of listener', subscription, message, x); 1541 } 1542 } 1543 } 1544 } 1545 } 1546 } 1547 } 1548 1549 function _notifyListeners(channel, message) { 1550 // Notify direct listeners 1551 _notify(channel, message); 1552 1553 // Notify the globbing listeners 1554 var channelParts = channel.split('/'); 1555 var last = channelParts.length - 1; 1556 for (var i = last; i > 0; --i) { 1557 var channelPart = channelParts.slice(0, i).join('/') + '/*'; 1558 // We don't want to notify /foo/* if the channel is /foo/bar/baz, 1559 // so we stop at the first non recursive globbing 1560 if (i === last) { 1561 _notify(channelPart, message); 1562 } 1563 // Add the recursive globber and notify 1564 channelPart += '*'; 1565 _notify(channelPart, message); 1566 } 1567 } 1568 1569 function _cancelDelayedSend() { 1570 if (_scheduledSend !== null) { 1571 Utils.clearTimeout(_scheduledSend); 1572 } 1573 _scheduledSend = null; 1574 } 1575 1576 function _delayedSend(operation, delay) { 1577 _cancelDelayedSend(); 1578 var time = _advice.interval + delay; 1579 _cometd._debug('Function scheduled in', time, 'ms, interval =', _advice.interval, 'backoff =', _backoff, operation); 1580 _scheduledSend = Utils.setTimeout(_cometd, operation, time); 1581 } 1582 1583 // Needed to break cyclic dependencies between function definitions 1584 var _handleMessages; 1585 var _handleFailure; 1586 1587 /** 1588 * Delivers the messages to the CometD server 1589 * @param sync whether the send is synchronous 1590 * @param messages the array of messages to send 1591 * @param metaConnect true if this send is on /meta/connect 1592 * @param extraPath an extra path to append to the Bayeux server URL 1593 */ 1594 function _send(sync, messages, metaConnect, extraPath) { 1595 // We must be sure that the messages have a clientId. 1596 // This is not guaranteed since the handshake may take time to return 1597 // (and hence the clientId is not known yet) and the application 1598 // may create other messages. 1599 for (var i = 0; i < messages.length; ++i) { 1600 var message = messages[i]; 1601 var messageId = message.id; 1602 1603 if (_clientId) { 1604 message.clientId = _clientId; 1605 } 1606 1607 message = _applyOutgoingExtensions(message); 1608 if (message !== undefined && message !== null) { 1609 // Extensions may have modified the message id, but we need to own it. 1610 message.id = messageId; 1611 messages[i] = message; 1612 } else { 1613 delete _callbacks[messageId]; 1614 messages.splice(i--, 1); 1615 } 1616 } 1617 1618 if (messages.length === 0) { 1619 return; 1620 } 1621 1622 var url = _cometd.getURL(); 1623 if (_config.appendMessageTypeToURL) { 1624 // If url does not end with '/', then append it 1625 if (!url.match(/\/$/)) { 1626 url = url + '/'; 1627 } 1628 if (extraPath) { 1629 url = url + extraPath; 1630 } 1631 } 1632 1633 var envelope = { 1634 url: url, 1635 sync: sync, 1636 messages: messages, 1637 onSuccess: function(rcvdMessages) { 1638 try { 1639 _handleMessages.call(_cometd, rcvdMessages); 1640 } catch (x) { 1641 _cometd._info('Exception during handling of messages', x); 1642 } 1643 }, 1644 onFailure: function(conduit, messages, failure) { 1645 try { 1646 var transport = _cometd.getTransport(); 1647 failure.connectionType = transport ? transport.getType() : "unknown"; 1648 _handleFailure.call(_cometd, conduit, messages, failure); 1649 } catch (x) { 1650 _cometd._info('Exception during handling of failure', x); 1651 } 1652 } 1653 }; 1654 _cometd._debug('Send', envelope); 1655 _transport.send(envelope, metaConnect); 1656 } 1657 1658 function _queueSend(message) { 1659 if (_batch > 0 || _internalBatch === true) { 1660 _messageQueue.push(message); 1661 } else { 1662 _send(false, [message], false); 1663 } 1664 } 1665 1666 /** 1667 * Sends a complete bayeux message. 1668 * This method is exposed as a public so that extensions may use it 1669 * to send bayeux message directly, for example in case of re-sending 1670 * messages that have already been sent but that for some reason must 1671 * be resent. 1672 */ 1673 this.send = _queueSend; 1674 1675 function _resetBackoff() { 1676 _backoff = 0; 1677 } 1678 1679 function _increaseBackoff() { 1680 if (_backoff < _config.maxBackoff) { 1681 _backoff += _config.backoffIncrement; 1682 } 1683 return _backoff; 1684 } 1685 1686 /** 1687 * Starts a the batch of messages to be sent in a single request. 1688 * @see #_endBatch(sendMessages) 1689 */ 1690 function _startBatch() { 1691 ++_batch; 1692 _cometd._debug('Starting batch, depth', _batch); 1693 } 1694 1695 function _flushBatch() { 1696 var messages = _messageQueue; 1697 _messageQueue = []; 1698 if (messages.length > 0) { 1699 _send(false, messages, false); 1700 } 1701 } 1702 1703 /** 1704 * Ends the batch of messages to be sent in a single request, 1705 * optionally sending messages present in the message queue depending 1706 * on the given argument. 1707 * @see #_startBatch() 1708 */ 1709 function _endBatch() { 1710 --_batch; 1711 _cometd._debug('Ending batch, depth', _batch); 1712 if (_batch < 0) { 1713 throw 'Calls to startBatch() and endBatch() are not paired'; 1714 } 1715 1716 if (_batch === 0 && !_isDisconnected() && !_internalBatch) { 1717 _flushBatch(); 1718 } 1719 } 1720 1721 /** 1722 * Sends the connect message 1723 */ 1724 function _connect() { 1725 if (!_isDisconnected()) { 1726 var bayeuxMessage = { 1727 id: _nextMessageId(), 1728 channel: '/meta/connect', 1729 connectionType: _transport.getType() 1730 }; 1731 1732 // In case of reload or temporary loss of connection 1733 // we want the next successful connect to return immediately 1734 // instead of being held by the server, so that connect listeners 1735 // can be notified that the connection has been re-established 1736 if (!_connected) { 1737 bayeuxMessage.advice = { timeout: 0 }; 1738 } 1739 1740 _setStatus('connecting'); 1741 _cometd._debug('Connect sent', bayeuxMessage); 1742 _send(false, [bayeuxMessage], true, 'connect'); 1743 _setStatus('connected'); 1744 } 1745 } 1746 1747 function _delayedConnect(delay) { 1748 _setStatus('connecting'); 1749 _delayedSend(function() { 1750 _connect(); 1751 }, delay); 1752 } 1753
1754 function _updateAdvice(newAdvice) { 1755 if (newAdvice) { 1756 _advice = _cometd._mixin(false, {}, _config.advice, newAdvice); 1757 _cometd._debug('New advice', _advice); 1758 } 1759 } 1760 1761 function _disconnect(abort) { 1762 _cancelDelayedSend(); 1763 if (abort && _transport) { 1764 _transport.abort(); 1765 } 1766 _clientId = null; 1767 _setStatus('disconnected'); 1768 _batch = 0; 1769 _resetBackoff(); 1770 _transport = null; 1771 _reestablish = false; 1772 _connected = false; 1773 1774 // Fail any existing queued message 1775 if (_messageQueue.length > 0) { 1776 var messages = _messageQueue; 1777 _messageQueue = []; 1778 _handleFailure.call(_cometd, undefined, messages, { 1779 reason: 'Disconnected' 1780 }); 1781 } 1782 } 1783 1784 function _notifyTransportException(oldTransport, newTransport, failure) { 1785 var handler = _cometd.onTransportException; 1786 if (_isFunction(handler)) { 1787 _cometd._debug('Invoking transport exception handler', oldTransport, newTransport, failure); 1788 try { 1789 handler.call(_cometd, failure, oldTransport, newTransport); 1790 } catch (x) { 1791 _cometd._info('Exception during execution of transport exception handler', x); 1792 } 1793 } 1794 } 1795 1796 /** 1797 * Sends the initial handshake message 1798 */ 1799 function _handshake(handshakeProps, handshakeCallback) { 1800 if (_isFunction(handshakeProps)) { 1801 handshakeCallback = handshakeProps; 1802 handshakeProps = undefined; 1803 } 1804 1805 _clientId = null; 1806 1807 _clearSubscriptions(); 1808 1809 // Reset the transports if we're not retrying the handshake 1810 if (_isDisconnected()) { 1811 _transports.reset(true); 1812 } 1813 1814 // Reset the advice. 1815 _updateAdvice({}); 1816 1817 _batch = 0; 1818 1819 // Mark the start of an internal batch. 1820 // This is needed because handshake and connect are async. 1821 // It may happen that the application calls init() then subscribe() 1822 // and the subscribe message is sent before the connect message, if 1823 // the subscribe message is not held until the connect message is sent. 1824 // So here we start a batch to hold temporarily any message until 1825 // the connection is fully established. 1826 _internalBatch = true; 1827 1828 // Save the properties provided by the user, so that 1829 // we can reuse them during automatic re-handshake 1830 _handshakeProps = handshakeProps; 1831 _handshakeCallback = handshakeCallback; 1832 1833 var version = '1.0'; 1834 1835 // Figure out the transports to send to the server 1836 var url = _cometd.getURL(); 1837 var transportTypes = _transports.findTransportTypes(version, _crossDomain, url); 1838 1839 var bayeuxMessage = { 1840 id: _nextMessageId(), 1841 version: version, 1842 minimumVersion: version, 1843 channel: '/meta/handshake', 1844 supportedConnectionTypes: transportTypes, 1845 advice: { 1846 timeout: _advice.timeout, 1847 interval: _advice.interval 1848 } 1849 }; 1850 // Do not allow the user to override important fields. 1851 var message = _cometd._mixin(false, {}, _handshakeProps, bayeuxMessage); 1852 1853 // Save the callback. 1854 _cometd._putCallback(message.id, handshakeCallback); 1855 1856 // Pick up the first available transport as initial transport 1857 // since we don't know if the server supports it 1858 if (!_transport) { 1859 _transport = _transports.negotiateTransport(transportTypes, version, _crossDomain, url); 1860 if (!_transport) { 1861 var failure = 'Could not find initial transport among: ' + _transports.getTransportTypes(); 1862 _cometd._warn(failure); 1863 throw failure; 1864 } 1865 } 1866 1867 _cometd._debug('Initial transport is', _transport.getType()); 1868 1869 // We started a batch to hold the application messages, 1870 // so here we must bypass it and send immediately. 1871 _setStatus('handshaking'); 1872 _cometd._debug('Handshake sent', message); 1873 _send(false, [message], false, 'handshake'); 1874 } 1875 1876 function _delayedHandshake(delay) { 1877 _setStatus('handshaking'); 1878 1879 // We will call _handshake() which will reset _clientId, but we want to avoid 1880 // that between the end of this method and the call to _handshake() someone may 1881 // call publish() (or other methods that call _queueSend()). 1882 _internalBatch = true; 1883 1884 _delayedSend(function() { 1885 _handshake(_handshakeProps, _handshakeCallback); 1886 }, delay); 1887 } 1888 1889 function _notifyCallback(callback, message) { 1890 try { 1891 callback.call(_cometd, message); 1892 } catch (x) { 1893 var handler = _cometd.onCallbackException; 1894 if (_isFunction(handler)) { 1895 _cometd._debug('Invoking callback exception handler', x); 1896 try { 1897 handler.call(_cometd, x, message); 1898 } catch (xx) { 1899 _cometd._info('Exception during execution of callback exception handler', xx); 1900 } 1901 } else { 1902 _cometd._info('Exception during execution of message callback', x); 1903 } 1904 } 1905 } 1906 1907 this._getCallback = function(messageId) { 1908 return _callbacks[messageId]; 1909 }; 1910 1911 this._putCallback = function(messageId, callback) { 1912 var result = this._getCallback(messageId); 1913 if (_isFunction(callback)) { 1914 _callbacks[messageId] = callback; 1915 } 1916 return result; 1917 }; 1918 1919 function _handleCallback(message) { 1920 var callback = _cometd._getCallback([message.id]); 1921 if (_isFunction(callback)) { 1922 delete _callbacks[message.id]; 1923 _notifyCallback(callback, message); 1924 } 1925 } 1926 1927 function _handleRemoteCall(message) { 1928 var context = _remoteCalls[message.id]; 1929 delete _remoteCalls[message.id]; 1930 if (context) { 1931 _cometd._debug('Handling remote call response for', message, 'with context', context); 1932 1933 // Clear the timeout, if present. 1934 var timeout = context.timeout; 1935 if (timeout) { 1936 Utils.clearTimeout(timeout); 1937 } 1938 1939 var callback = context.callback; 1940 if (_isFunction(callback)) { 1941 _notifyCallback(callback, message); 1942 return true; 1943 } 1944 } 1945 return false; 1946 } 1947 1948 this.onTransportFailure = function(message, failureInfo, failureHandler) { 1949 this._debug('Transport failure', failureInfo, 'for', message); 1950 1951 var transports = this.getTransportRegistry(); 1952 var url = this.getURL(); 1953 var crossDomain = this._isCrossDomain(_splitURL(url)[2]); 1954 var version = '1.0'; 1955 var transportTypes = transports.findTransportTypes(version, crossDomain, url); 1956 1957 if (failureInfo.action === 'none') { 1958 if (message.channel === '/meta/handshake') { 1959 if (!failureInfo.transport) { 1960 var failure = 'Could not negotiate transport, client=[' + transportTypes + '], server=[' + message.supportedConnectionTypes + ']'; 1961 this._warn(failure); 1962 _notifyTransportException(_transport.getType(), null, { 1963 reason: failure, 1964 connectionType: _transport.getType(), 1965 transport: _transport 1966 }); 1967 } 1968 } 1969 } else { 1970 failureInfo.delay = this.getBackoffPeriod(); 1971 // Different logic depending on whether we are handshaking or connecting. 1972 if (message.channel === '/meta/handshake') { 1973 if (!failureInfo.transport) {
1974 // The transport is invalid, try to negotiate again. 1975 var newTransport = transports.negotiateTransport(transportTypes, version, crossDomain, url); 1976 if (!newTransport) { 1977 this._warn('Could not negotiate transport, client=[' + transportTypes + ']'); 1978 _notifyTransportException(_transport.getType(), null, message.failure); 1979 failureInfo.action = 'none'; 1980 } else { 1981 this._debug('Transport', _transport.getType(), '->', newTransport.getType()); 1982 _notifyTransportException(_transport.getType(), newTransport.getType(), message.failure); 1983 failureInfo.action = 'handshake'; 1984 failureInfo.transport = newTransport; 1985 } 1986 } 1987 1988 if (failureInfo.action !== 'none') { 1989 this.increaseBackoffPeriod(); 1990 } 1991 } else { 1992 var now = new Date().getTime(); 1993 1994 if (_unconnectTime === 0) { 1995 _unconnectTime = now; 1996 } 1997 1998 if (failureInfo.action === 'retry') { 1999 failureInfo.delay = this.increaseBackoffPeriod(); 2000 // Check whether we may switch to handshaking. 2001 var maxInterval = _advice.maxInterval; 2002 if (maxInterval > 0) { 2003 var expiration = _advice.timeout + _advice.interval + maxInterval; 2004 var unconnected = now - _unconnectTime; 2005 if (unconnected + _backoff > expiration) { 2006 failureInfo.action = 'handshake'; 2007 } 2008 } 2009 } 2010 2011 if (failureInfo.action === 'handshake') { 2012 failureInfo.delay = 0; 2013 transports.reset(false); 2014 this.resetBackoffPeriod(); 2015 } 2016 } 2017 } 2018 2019 failureHandler.call(_cometd, failureInfo); 2020 }; 2021 2022 function _handleTransportFailure(failureInfo) { 2023 _cometd._debug('Transport failure handling', failureInfo); 2024 2025 if (failureInfo.transport) { 2026 _transport = failureInfo.transport; 2027 } 2028 2029 if (failureInfo.url) { 2030 _transport.setURL(failureInfo.url); 2031 } 2032 2033 var action = failureInfo.action; 2034 var delay = failureInfo.delay || 0; 2035 switch (action) { 2036 case 'handshake': 2037 _delayedHandshake(delay); 2038 break; 2039 case 'retry': 2040 _delayedConnect(delay); 2041 break; 2042 case 'none': 2043 _disconnect(true); 2044 break; 2045 default: 2046 throw 'Unknown action ' + action; 2047 } 2048 } 2049 2050 function _failHandshake(message, failureInfo) { 2051 _handleCallback(message); 2052 _notifyListeners('/meta/handshake', message); 2053 _notifyListeners('/meta/unsuccessful', message); 2054 2055 // The listeners may have disconnected. 2056 if (_isDisconnected()) { 2057 failureInfo.action = 'none'; 2058 } 2059 2060 _cometd.onTransportFailure.call(_cometd, message, failureInfo, _handleTransportFailure); 2061 } 2062 2063 function _handshakeResponse(message) { 2064 var url = _cometd.getURL(); 2065 if (message.successful) { 2066 var crossDomain = _cometd._isCrossDomain(_splitURL(url)[2]); 2067 var newTransport = _transports.negotiateTransport(message.supportedConnectionTypes, message.version, crossDomain, url); 2068 if (newTransport === null) { 2069 message.successful = false; 2070 _failHandshake(message, { 2071 cause: 'negotiation', 2072 action: 'none', 2073 transport: null 2074 }); 2075 return; 2076 } else if (_transport !== newTransport) { 2077 _cometd._debug('Transport', _transport.getType(), '->', newTransport.getType()); 2078 _transport = newTransport; 2079 } 2080 2081 _clientId = message.clientId; 2082 2083 // End the internal batch and allow held messages from the application 2084 // to go to the server (see _handshake() where we start the internal batch). 2085 _internalBatch = false; 2086 _flushBatch(); 2087 2088 // Here the new transport is in place, as well as the clientId, so 2089 // the listeners can perform a publish() if they want. 2090 // Notify the listeners before the connect below. 2091 message.reestablish = _reestablish; 2092 _reestablish = true; 2093 2094 _handleCallback(message); 2095 _notifyListeners('/meta/handshake', message); 2096 2097 _handshakeMessages = message['x-messages'] || 0; 2098 2099 var action = _isDisconnected() ? 'none' : _advice.reconnect || 'retry'; 2100 switch (action) { 2101 case 'retry': 2102 _resetBackoff(); 2103 if (_handshakeMessages === 0) { 2104 _delayedConnect(0); 2105 } else { 2106 _cometd._debug('Processing', _handshakeMessages, 'handshake-delivered messages'); 2107 } 2108 break; 2109 case 'none': 2110 _disconnect(true); 2111 break; 2112 default: 2113 throw 'Unrecognized advice action ' + action; 2114 } 2115 } else { 2116 _failHandshake(message, { 2117 cause: 'unsuccessful', 2118 action: _advice.reconnect || 'handshake', 2119 transport: _transport 2120 }); 2121 } 2122 } 2123 2124 function _handshakeFailure(message) { 2125 _failHandshake(message, { 2126 cause: 'failure', 2127 action: 'handshake', 2128 transport: null 2129 }); 2130 } 2131 2132 function _failConnect(message, failureInfo) { 2133 // Notify the listeners after the status change but before the next action. 2134 _notifyListeners('/meta/connect', message); 2135 _notifyListeners('/meta/unsuccessful', message); 2136 2137 // The listeners may have disconnected. 2138 if (_isDisconnected()) {
2139 failureInfo.action = 'none'; 2140 } 2141 2142 _cometd.onTransportFailure.call(_cometd, message, failureInfo, _handleTransportFailure); 2143 } 2144 2145 function _connectResponse(message) { 2146 _connected = message.successful; 2147 2148 if (_connected) { 2149 _notifyListeners('/meta/connect', message); 2150 2151 // Normally, the advice will say "reconnect: 'retry', interval: 0" 2152 // and the server will hold the request, so when a response returns 2153 // we immediately call the server again (long polling). 2154 // Listeners can call disconnect(), so check the state after they run. 2155 var action = _isDisconnected() ? 'none' : _advice.reconnect || 'retry'; 2156 switch (action) { 2157 case 'retry': 2158 _resetBackoff(); 2159 _delayedConnect(_backoff); 2160 break; 2161 case 'none': 2162 _disconnect(false); 2163 break; 2164 default: 2165 throw 'Unrecognized advice action ' + action; 2166 } 2167 } else { 2168 _failConnect(message, { 2169 cause: 'unsuccessful', 2170 action: _advice.reconnect || 'retry', 2171 transport: _transport 2172 }); 2173 } 2174 } 2175 2176 function _connectFailure(message) { 2177 _connected = false; 2178 2179 _failConnect(message, { 2180 cause: 'failure', 2181 action: 'retry', 2182 transport: null 2183 }); 2184 } 2185 2186 function _failDisconnect(message) { 2187 _disconnect(true); 2188 _handleCallback(message); 2189 _notifyListeners('/meta/disconnect', message); 2190 _notifyListeners('/meta/unsuccessful', message); 2191 } 2192 2193 function _disconnectResponse(message) { 2194 if (message.successful) { 2195 // Wait for the /meta/connect to arrive. 2196 _disconnect(false); 2197 _handleCallback(message); 2198 _notifyListeners('/meta/disconnect', message); 2199 } else { 2200 _failDisconnect(message); 2201 } 2202 } 2203 2204 function _disconnectFailure(message) { 2205 _failDisconnect(message); 2206 } 2207 2208 function _failSubscribe(message) { 2209 var subscriptions = _listeners[message.subscription]; 2210 if (subscriptions) { 2211 for (var id in subscriptions) { 2212 if (subscriptions.hasOwnProperty(id)) { 2213 var subscription = subscriptions[id]; 2214 if (subscription && !subscription.listener) { 2215 delete subscriptions[id]; 2216 _cometd._debug('Removed failed subscription', subscription); 2217 } 2218 } 2219 } 2220 } 2221 _handleCallback(message); 2222 _notifyListeners('/meta/subscribe', message); 2223 _notifyListeners('/meta/unsuccessful', message); 2224 } 2225 2226 function _subscribeResponse(message) { 2227 if (message.successful) { 2228 _handleCallback(message); 2229 _notifyListeners('/meta/subscribe', message); 2230 } else { 2231 _failSubscribe(message); 2232 } 2233 } 2234 2235 function _subscribeFailure(message) { 2236 _failSubscribe(message); 2237 } 2238 2239 function _failUnsubscribe(message) { 2240 _handleCallback(message); 2241 _notifyListeners('/meta/unsubscribe', message); 2242 _notifyListeners('/meta/unsuccessful', message); 2243 } 2244 2245 function _unsubscribeResponse(message) { 2246 if (message.successful) { 2247 _handleCallback(message); 2248 _notifyListeners('/meta/unsubscribe', message); 2249 } else { 2250 _failUnsubscribe(message); 2251 } 2252 } 2253 2254 function _unsubscribeFailure(message) { 2255 _failUnsubscribe(message); 2256 } 2257 2258 function _failMessage(message) { 2259 if (!_handleRemoteCall(message)) { 2260 _handleCallback(message); 2261 _notifyListeners('/meta/publish', message); 2262 _notifyListeners('/meta/unsuccessful', message); 2263 } 2264 } 2265 2266 function _messageResponse(message) { 2267 if (message.data !== undefined) { 2268 if (!_handleRemoteCall(message)) { 2269 _notifyListeners(message.channel, message); 2270 if (_handshakeMessages > 0) { 2271 --_handshakeMessages; 2272 if (_handshakeMessages === 0) { 2273 _cometd._debug('Processed last handshake-delivered message'); 2274 _delayedConnect(0); 2275 } 2276 } 2277 } 2278 } else { 2279 if (message.successful === undefined) { 2280 _cometd._warn('Unknown Bayeux Message', message); 2281 } else { 2282 if (message.successful) { 2283 _handleCallback(message); 2284 _notifyListeners('/meta/publish', message); 2285 } else { 2286 _failMessage(message); 2287 } 2288 } 2289 } 2290 } 2291 2292 function _messageFailure(failure) { 2293 _failMessage(failure); 2294 } 2295 2296 function _receive(message) { 2297 _unconnectTime = 0; 2298 2299 message = _applyIncomingExtensions(message); 2300 if (message === undefined || message === null) { 2301 return; 2302 } 2303
2304 _updateAdvice(message.advice); 2305 2306 var channel = message.channel; 2307 switch (channel) { 2308 case '/meta/handshake': 2309 _handshakeResponse(message); 2310 break; 2311 case '/meta/connect': 2312 _connectResponse(message); 2313 break; 2314 case '/meta/disconnect': 2315 _disconnectResponse(message); 2316 break; 2317 case '/meta/subscribe': 2318 _subscribeResponse(message); 2319 break; 2320 case '/meta/unsubscribe': 2321 _unsubscribeResponse(message); 2322 break; 2323 default: 2324 _messageResponse(message); 2325 break; 2326 } 2327 } 2328 2329 /** 2330 * Receives a message. 2331 * This method is exposed as a public so that extensions may inject 2332 * messages simulating that they had been received. 2333 */ 2334 this.receive = _receive; 2335 2336 _handleMessages = function(rcvdMessages) { 2337 _cometd._debug('Received', rcvdMessages); 2338 2339 for (var i = 0; i < rcvdMessages.length; ++i) { 2340 var message = rcvdMessages[i]; 2341 _receive(message); 2342 } 2343 }; 2344 2345 _handleFailure = function(conduit, messages, failure) { 2346 _cometd._debug('handleFailure', conduit, messages, failure); 2347 2348 failure.transport = conduit; 2349 for (var i = 0; i < messages.length; ++i) { 2350 var message = messages[i]; 2351 var failureMessage = { 2352 id: message.id, 2353 successful: false, 2354 channel: message.channel, 2355 failure: failure 2356 }; 2357 failure.message = message; 2358 switch (message.channel) { 2359 case '/meta/handshake': 2360 _handshakeFailure(failureMessage); 2361 break; 2362 case '/meta/connect': 2363 _connectFailure(failureMessage); 2364 break; 2365 case '/meta/disconnect': 2366 _disconnectFailure(failureMessage); 2367 break; 2368 case '/meta/subscribe': 2369 failureMessage.subscription = message.subscription; 2370 _subscribeFailure(failureMessage); 2371 break; 2372 case '/meta/unsubscribe': 2373 failureMessage.subscription = message.subscription; 2374 _unsubscribeFailure(failureMessage); 2375 break; 2376 default: 2377 _messageFailure(failureMessage); 2378 break; 2379 } 2380 } 2381 }; 2382 2383 function _hasSubscriptions(channel) { 2384 var subscriptions = _listeners[channel]; 2385 if (subscriptions) { 2386 for (var id in subscriptions) { 2387 if (subscriptions.hasOwnProperty(id)) { 2388 if (subscriptions[id]) { 2389 return true; 2390 } 2391 } 2392 } 2393 } 2394 return false; 2395 } 2396 2397 function _resolveScopedCallback(scope, callback) { 2398 var delegate = { 2399 scope: scope, 2400 method: callback 2401 }; 2402 if (_isFunction(scope)) { 2403 delegate.scope = undefined; 2404 delegate.method = scope; 2405 } else { 2406 if (_isString(callback)) { 2407 if (!scope) { 2408 throw 'Invalid scope ' + scope; 2409 } 2410 delegate.method = scope[callback]; 2411 if (!_isFunction(delegate.method)) { 2412 throw 'Invalid callback ' + callback + ' for scope ' + scope; 2413 } 2414 }
2414 else if (!_isFunction(callback)) { 2415 throw 'Invalid callback ' + callback; 2416 } 2417 } 2418 return delegate; 2419 } 2420 2421 function _addListener(channel, scope, callback, isListener) { 2422 // The data structure is a map<channel, subscription[]>, where each subscription 2423 // holds the callback to be called and its scope. 2424 2425 var delegate = _resolveScopedCallback(scope, callback); 2426 _cometd._debug('Adding', isListener ? 'listener' : 'subscription', 'on', channel, 'with scope', delegate.scope, 'and callback', delegate.method); 2427 2428 var id = ++_listenerId; 2429 var subscription = { 2430 id: id, 2431 channel: channel, 2432 scope: delegate.scope, 2433 callback: delegate.method, 2434 listener: isListener 2435 }; 2436 2437 var subscriptions = _listeners[channel]; 2438 if (!subscriptions) { 2439 subscriptions = {}; 2440 _listeners[channel] = subscriptions; 2441 } 2442 2443 subscriptions[id] = subscription; 2444 2445 _cometd._debug('Added', isListener ? 'listener' : 'subscription', subscription); 2446 2447 return subscription; 2448 } 2449 2450 // 2451 // PUBLIC API 2452 // 2453 2454 /** 2455 * Registers the given transport under the given transport type. 2456 * The optional index parameter specifies the "priority" at which the 2457 * transport is registered (where 0 is the max priority). 2458 * If a transport with the same type is already registered, this function 2459 * does nothing and returns false. 2460 * @param type the transport type 2461 * @param transport the transport object 2462 * @param index the index at which this transport is to be registered 2463 * @return true if the transport has been registered, false otherwise 2464 * @see #unregisterTransport(type) 2465 */ 2466 this.registerTransport = function(type, transport, index) { 2467 var result = _transports.add(type, transport, index); 2468 if (result) { 2469 this._debug('Registered transport', type); 2470 2471 if (_isFunction(transport.registered)) { 2472 transport.registered(type, this); 2473 } 2474 } 2475 return result; 2476 }; 2477 2478 /** 2479 * Unregisters the transport with the given transport type. 2480 * @param type the transport type to unregister 2481 * @return the transport that has been unregistered, 2482 * or null if no transport was previously registered under the given transport type 2483 */ 2484 this.unregisterTransport = function(type) { 2485 var transport = _transports.remove(type); 2486 if (transport !== null) { 2487 this._debug('Unregistered transport', type); 2488 2489 if (_isFunction(transport.unregistered)) { 2490 transport.unregistered(); 2491 } 2492 } 2493 return transport; 2494 }; 2495 2496 this.unregisterTransports = function() { 2497 _transports.clear(); 2498 }; 2499 2500 /** 2501 * @return an array of all registered transport types 2502 */ 2503 this.getTransportTypes = function() { 2504 return _transports.getTransportTypes(); 2505 }; 2506 2507 this.findTransport = function(name) { 2508 return _transports.find(name); 2509 }; 2510 2511 /** 2512 * @returns the TransportRegistry object 2513 */ 2514 this.getTransportRegistry = function() { 2515 return _transports; 2516 }; 2517 2518 /** 2519 * Configures the initial Bayeux communication with the Bayeux server. 2520 * Configuration is passed via an object that must contain a mandatory field <code>url</code> 2521 * of type string containing the URL of the Bayeux server. 2522 * @param configuration the configuration object 2523 */ 2524 this.configure = function(configuration) { 2525 _configure.call(this, configuration); 2526 }; 2527 2528 /** 2529 * Configures and establishes the Bayeux communication with the Bayeux server 2530 * via a handshake and a subsequent connect. 2531 * @param configuration the configuration object 2532 * @param handshakeProps an object to be merged with the handshake message 2533 * @see #configure(configuration) 2534 * @see #handshake(handshakeProps) 2535 */ 2536 this.init = function(configuration, handshakeProps) { 2537 this.configure(configuration); 2538 this.handshake(handshakeProps); 2539 }; 2540 2541 /** 2542 * Establishes the Bayeux communication with the Bayeux server 2543 * via a handshake and a subsequent connect. 2544 * @param handshakeProps an object to be merged with the handshake message 2545 * @param handshakeCallback a function to be invoked when the handshake is acknowledged 2546 */ 2547 this.handshake = function(handshakeProps, handshakeCallback) { 2548 if (_status !== 'disconnected') { 2549 throw 'Illegal state: handshaken'; 2550 } 2551 _handshake(handshakeProps, handshakeCallback); 2552 }; 2553 2554 /** 2555 * Disconnects from the Bayeux server.
2556 * It is possible to suggest to attempt a synchronous disconnect, but this feature 2557 * may only be available in certain transports (for example, long-polling may support 2558 * it, callback-polling certainly does not). 2559 * @param sync whether attempt to perform a synchronous disconnect 2560 * @param disconnectProps an object to be merged with the disconnect message 2561 * @param disconnectCallback a function to be invoked when the disconnect is acknowledged 2562 */ 2563 this.disconnect = function(sync, disconnectProps, disconnectCallback) { 2564 if (_isDisconnected()) { 2565 return; 2566 } 2567 2568 if (typeof sync !== 'boolean') { 2569 disconnectCallback = disconnectProps; 2570 disconnectProps = sync; 2571 sync = false; 2572 } 2573 if (_isFunction(disconnectProps)) { 2574 disconnectCallback = disconnectProps; 2575 disconnectProps = undefined; 2576 } 2577 2578 var bayeuxMessage = { 2579 id: _nextMessageId(), 2580 channel: '/meta/disconnect' 2581 }; 2582 // Do not allow the user to override important fields. 2583 var message = this._mixin(false, {}, disconnectProps, bayeuxMessage); 2584 2585 // Save the callback. 2586 _cometd._putCallback(message.id, disconnectCallback); 2587 2588 _setStatus('disconnecting'); 2589 _send(sync === true, [message], false, 'disconnect'); 2590 }; 2591 2592 /** 2593 * Marks the start of a batch of application messages to be sent to the server 2594 * in a single request, obtaining a single response containing (possibly) many 2595 * application reply messages. 2596 * Messages are held in a queue and not sent until {@link #endBatch()} is called. 2597 * If startBatch() is called multiple times, then an equal number of endBatch() 2598 * calls must be made to close and send the batch of messages. 2599 * @see #endBatch() 2600 */ 2601 this.startBatch = function() { 2602 _startBatch(); 2603 }; 2604 2605 /** 2606 * Marks the end of a batch of application messages to be sent to the server 2607 * in a single request. 2608 * @see #startBatch() 2609 */ 2610 this.endBatch = function() { 2611 _endBatch(); 2612 }; 2613 2614 /** 2615 * Executes the given callback in the given scope, surrounded by a {@link #startBatch()} 2616 * and {@link #endBatch()} calls. 2617 * @param scope the scope of the callback, may be omitted 2618 * @param callback the callback to be executed within {@link #startBatch()} and {@link #endBatch()} calls 2619 */ 2620 this.batch = function(scope, callback) { 2621 var delegate = _resolveScopedCallback(scope, callback); 2622 this.startBatch(); 2623 try { 2624 delegate.method.call(delegate.scope); 2625 this.endBatch(); 2626 } catch (x) { 2627 this._info('Exception during execution of batch', x); 2628 this.endBatch(); 2629 throw x; 2630 } 2631 }; 2632 2633 /** 2634 * Adds a listener for bayeux messages, performing the given callback in the given scope 2635 * when a message for the given channel arrives. 2636 * @param channel the channel the listener is interested to 2637 * @param scope the scope of the callback, may be omitted 2638 * @param callback the callback to call when a message is sent to the channel 2639 * @returns the subscription handle to be passed to {@link #removeListener(object)} 2640 * @see #removeListener(subscription) 2641 */ 2642 this.addListener = function(channel, scope, callback) { 2643 if (arguments.length < 2) { 2644 throw 'Illegal arguments number: required 2, got ' + arguments.length; 2645 } 2646 if (!_isString(channel)) { 2647 throw 'Illegal argument type: channel must be a string'; 2648 } 2649 2650 return _addListener(channel, scope, callback, true); 2651 }; 2652 2653 /** 2654 * Removes the subscription obtained with a call to {@link #addListener(string, object, function)}. 2655 * @param subscription the subscription to unsubscribe. 2656 * @see #addListener(channel, scope, callback) 2657 */ 2658 this.removeListener = function(subscription) { 2659 // Beware of subscription.id == 0, which is falsy => cannot use !subscription.id 2660 if (!subscription || !subscription.channel || !("id" in subscription)) { 2661 throw 'Invalid argument: expected subscription, not ' + subscription; 2662 } 2663 2664 _removeListener(subscription); 2665 }; 2666 2667 /** 2668 * Removes all listeners registered with {@link #addListener(channel, scope, callback)} or 2669 * {@link #subscribe(channel, scope, callback)}. 2670 */ 2671 this.clearListeners = function() { 2672 _listeners = {}; 2673 }; 2674 2675 /** 2676 * Subscribes to the given channel, performing the given callback in the given scope 2677 * when a message for the channel arrives. 2678 * @param channel the channel to subscribe to 2679 * @param scope the scope of the callback, may be omitted 2680 * @param callback the callback to call when a message is sent to the channel 2681 * @param subscribeProps an object to be merged with the subscribe message 2682 * @param subscribeCallback a function to be invoked when the subscription is acknowledged 2683 * @return the subscription handle to be passed to {@link #unsubscribe(object)} 2684 */
2685 this.subscribe = function(channel, scope, callback, subscribeProps, subscribeCallback) { 2686 if (arguments.length < 2) { 2687 throw 'Illegal arguments number: required 2, got ' + arguments.length; 2688 } 2689 if (!_isString(channel)) { 2690 throw 'Illegal argument type: channel must be a string'; 2691 } 2692 if (_isDisconnected()) { 2693 throw 'Illegal state: disconnected'; 2694 } 2695 2696 // Normalize arguments 2697 if (_isFunction(scope)) { 2698 subscribeCallback = subscribeProps; 2699 subscribeProps = callback; 2700 callback = scope; 2701 scope = undefined; 2702 } 2703 if (_isFunction(subscribeProps)) { 2704 subscribeCallback = subscribeProps; 2705 subscribeProps = undefined; 2706 } 2707 2708 // Only send the message to the server if this client has not yet subscribed to the channel 2709 var send = !_hasSubscriptions(channel); 2710 2711 var subscription = _addListener(channel, scope, callback, false); 2712 2713 if (send) { 2714 // Send the subscription message after the subscription registration to avoid 2715 // races where the server would send a message to the subscribers, but here 2716 // on the client the subscription has not been added yet to the data structures 2717 var bayeuxMessage = { 2718 id: _nextMessageId(), 2719 channel: '/meta/subscribe', 2720 subscription: channel 2721 }; 2722 // Do not allow the user to override important fields. 2723 var message = this._mixin(false, {}, subscribeProps, bayeuxMessage); 2724 2725 // Save the callback. 2726 _cometd._putCallback(message.id, subscribeCallback); 2727 2728 _queueSend(message); 2729 } 2730 2731 return subscription; 2732 }; 2733 2734 /** 2735 * Unsubscribes the subscription obtained with a call to {@link #subscribe(string, object, function)}. 2736 * @param subscription the subscription to unsubscribe. 2737 * @param unsubscribeProps an object to be merged with the unsubscribe message 2738 * @param unsubscribeCallback a function to be invoked when the unsubscription is acknowledged 2739 */ 2740 this.unsubscribe = function(subscription, unsubscribeProps, unsubscribeCallback) { 2741 if (arguments.length < 1) { 2742 throw 'Illegal arguments number: required 1, got ' + arguments.length; 2743 } 2744 if (_isDisconnected()) { 2745 throw 'Illegal state: disconnected'; 2746 } 2747 2748 if (_isFunction(unsubscribeProps)) { 2749 unsubscribeCallback = unsubscribeProps; 2750 unsubscribeProps = undefined; 2751 } 2752 2753 // Remove the local listener before sending the message 2754 // This ensures that if the server fails, this client does not get notifications 2755 this.removeListener(subscription); 2756 2757 var channel = subscription.channel; 2758 // Only send the message to the server if this client unsubscribes the last subscription 2759 if (!_hasSubscriptions(channel)) { 2760 var bayeuxMessage = { 2761 id: _nextMessageId(), 2762 channel: '/meta/unsubscribe', 2763 subscription: channel 2764 }; 2765 // Do not allow the user to override important fields. 2766 var message = this._mixin(false, {}, unsubscribeProps, bayeuxMessage); 2767 2768 // Save the callback. 2769 _cometd._putCallback(message.id, unsubscribeCallback); 2770 2771 _queueSend(message); 2772 } 2773 }; 2774 2775 this.resubscribe = function(subscription, subscribeProps) { 2776 _removeSubscription(subscription); 2777 if (subscription) { 2778 return this.subscribe(subscription.channel, subscription.scope, subscription.callback, subscribeProps); 2779 } 2780 return undefined; 2781 }; 2782 2783 /** 2784 * Removes all subscriptions added via {@link #subscribe(channel, scope, callback, subscribeProps)}, 2785 * but does not remove the listeners added via {@link addListener(channel, scope, callback)}. 2786 */ 2787 this.clearSubscriptions = function() { 2788 _clearSubscriptions(); 2789 }; 2790 2791 /** 2792 * Publishes a message on the given channel, containing the given content. 2793 * @param channel the channel to publish the message to 2794 * @param content the content of the message 2795 * @param publishProps an object to be merged with the publish message 2796 * @param publishCallback a function to be invoked when the publish is acknowledged by the server 2797 */
2798 this.publish = function(channel, content, publishProps, publishCallback) { 2799 if (arguments.length < 1) { 2800 throw 'Illegal arguments number: required 1, got ' + arguments.length; 2801 } 2802 if (!_isString(channel)) { 2803 throw 'Illegal argument type: channel must be a string'; 2804 } 2805 if (/^\/meta\//.test(channel)) { 2806 throw 'Illegal argument: cannot publish to meta channels'; 2807 } 2808 if (_isDisconnected()) { 2809 throw 'Illegal state: disconnected'; 2810 } 2811 2812 if (_isFunction(content)) { 2813 publishCallback = content; 2814 content = {}; 2815 publishProps = undefined; 2816 } else if (_isFunction(publishProps)) { 2817 publishCallback = publishProps; 2818 publishProps = undefined; 2819 } 2820 2821 var bayeuxMessage = { 2822 id: _nextMessageId(), 2823 channel: channel, 2824 data: content 2825 }; 2826 // Do not allow the user to override important fields. 2827 var message = this._mixin(false, {}, publishProps, bayeuxMessage); 2828 2829 // Save the callback. 2830 _cometd._putCallback(message.id, publishCallback); 2831 2832 _queueSend(message); 2833 }; 2834 2835 /** 2836 * Publishes a message with binary data on the given channel. 2837 * The binary data chunk may be an ArrayBuffer, a DataView, a TypedArray 2838 * (such as Uint8Array) or a plain integer array. 2839 * The meta data object may contain additional application data such as 2840 * a file name, a mime type, etc. 2841 * @param channel the channel to publish the message to 2842 * @param data the binary data to publish 2843 * @param last whether the binary data chunk is the last
2844 * @param meta an object containing meta data associated to the binary chunk 2845 * @param callback a function to be invoked when the publish is acknowledged by the server 2846 */ 2847 this.publishBinary = function(channel, data, last, meta, callback) { 2848 if (_isFunction(data)) { 2849 callback = data; 2850 data = new ArrayBuffer(0); 2851 last = true; 2852 meta = undefined; 2853 } else if (_isFunction(last)) { 2854 callback = last; 2855 last = true; 2856 meta = undefined; 2857 } else if (_isFunction(meta)) { 2858 callback = meta; 2859 meta = undefined; 2860 } 2861 var content = { 2862 meta: meta, 2863 data: data, 2864 last: last 2865 }; 2866 var ext = { 2867 ext: { 2868 binary: { 2869 } 2870 } 2871 }; 2872 this.publish(channel, content, ext, callback); 2873 }; 2874 2875 this.remoteCall = function(target, content, timeout, callProps, callback) { 2876 if (arguments.length < 1) { 2877 throw 'Illegal arguments number: required 1, got ' + arguments.length; 2878 } 2879 if (!_isString(target)) { 2880 throw 'Illegal argument type: target must be a string'; 2881 } 2882 if (_isDisconnected()) { 2883 throw 'Illegal state: disconnected'; 2884 } 2885 2886 if (_isFunction(content)) { 2887 callback = content; 2888 content = {}; 2889 timeout = _config.maxNetworkDelay; 2890 callProps = undefined; 2891 } else if (_isFunction(timeout)) { 2892 callback = timeout; 2893 timeout = _config.maxNetworkDelay; 2894 callProps = undefined; 2895 } else if (_isFunction(callProps)) { 2896 callback = callProps; 2897 callProps = undefined; 2898 } 2899 2900 if (typeof timeout !== 'number') { 2901 throw 'Illegal argument type: timeout must be a number'; 2902 } 2903 2904 if (!target.match(/^\//)) { 2905 target = '/' + target; 2906 } 2907 var channel = '/service' + target; 2908 2909 var bayeuxMessage = { 2910 id: _nextMessageId(), 2911 channel: channel, 2912 data: content 2913 }; 2914 var message = this._mixin(false, {}, callProps, bayeuxMessage); 2915 2916 var context = { 2917 callback: callback 2918 }; 2919 if (timeout > 0) { 2920 context.timeout = Utils.setTimeout(_cometd, function() { 2921 _cometd._debug('Timing out remote call', message, 'after', timeout, 'ms'); 2922 _failMessage({ 2923 id: message.id, 2924 error: '406::timeout', 2925 successful: false, 2926 failure: { 2927 message : message, 2928 reason: 'Remote Call Timeout' 2929 } 2930 }); 2931 }, timeout); 2932 _cometd._debug('Scheduled remote call timeout', message, 'in', timeout, 'ms'); 2933 } 2934 _remoteCalls[message.id] = context; 2935 2936 _queueSend(message); 2937 }; 2938 2939 this.remoteCallBinary = function(target, data, last, meta, timeout, callback) { 2940 if (_isFunction(data)) { 2941 callback = data; 2942 data = new ArrayBuffer(0); 2943 last = true; 2944 meta = undefined; 2945 timeout = _config.maxNetworkDelay; 2946 } else if (_isFunction(last)) { 2947 callback = last; 2948 last = true; 2949 meta = undefined; 2950 timeout = _config.maxNetworkDelay; 2951 } else if (_isFunction(meta)) { 2952 callback = meta; 2953 meta = undefined; 2954 timeout = _config.maxNetworkDelay; 2955 } else if (_isFunction(timeout)) { 2956 callback = timeout; 2957 timeout = _config.maxNetworkDelay; 2958 } 2959 2960 var content = { 2961 meta: meta, 2962 data: data, 2963 last: last 2964 }; 2965 var ext = { 2966 ext: { 2967 binary: { 2968 } 2969 } 2970 }; 2971 2972 this.remoteCall(target, content, timeout, ext, callback); 2973 }; 2974 2975 /** 2976 * Returns a string representing the status of the bayeux communication with the Bayeux server.
2977 */ 2978 this.getStatus = function() { 2979 return _status; 2980 }; 2981 2982 /** 2983 * Returns whether this instance has been disconnected. 2984 */ 2985 this.isDisconnected = _isDisconnected; 2986 2987 /** 2988 * Sets the backoff period used to increase the backoff time when retrying an unsuccessful or failed message. 2989 * Default value is 1 second, which means if there is a persistent failure the retries will happen 2990 * after 1 second, then after 2 seconds, then after 3 seconds, etc. So for example with 15 seconds of 2991 * elapsed time, there will be 5 retries (at 1, 3, 6, 10 and 15 seconds elapsed). 2992 * @param period the backoff period to set 2993 * @see #getBackoffIncrement() 2994 */ 2995 this.setBackoffIncrement = function(period) { 2996 _config.backoffIncrement = period; 2997 }; 2998 2999 /** 3000 * Returns the backoff period used to increase the backoff time when retrying an unsuccessful or failed message. 3001 * @see #setBackoffIncrement(period) 3002 */ 3003 this.getBackoffIncrement = function() { 3004 return _config.backoffIncrement; 3005 }; 3006 3007 /** 3008 * Returns the backoff period to wait before retrying an unsuccessful or failed message. 3009 */ 3010 this.getBackoffPeriod = function() { 3011 return _backoff; 3012 }; 3013 3014 /** 3015 * Increases the backoff period up to the maximum value configured. 3016 * @returns the backoff period after increment 3017 * @see getBackoffIncrement 3018 */ 3019 this.increaseBackoffPeriod = function() { 3020 return _increaseBackoff(); 3021 }; 3022 3023 /** 3024 * Resets the backoff period to zero. 3025 */ 3026 this.resetBackoffPeriod = function() { 3027 _resetBackoff(); 3028 }; 3029 3030 /** 3031 * Sets the log level for console logging. 3032 * Valid values are the strings 'error', 'warn', 'info' and 'debug', from 3033 * less verbose to more verbose. 3034 * @param level the log level string 3035 */ 3036 this.setLogLevel = function(level) { 3037 _config.logLevel = level; 3038 }; 3039 3040 /** 3041 * Registers an extension whose callbacks are called for every incoming message 3042 * (that comes from the server to this client implementation) and for every 3043 * outgoing message (that originates from this client implementation for the 3044 * server). 3045 * The format of the extension object is the following: 3046 * <pre> 3047 * { 3048 * incoming: function(message) { ... }, 3049 * outgoing: function(message) { ... } 3050 * } 3051 * </pre> 3052 * Both properties are optional, but if they are present they will be called 3053 * respectively for each incoming message and for each outgoing message. 3054 * @param name the name of the extension 3055 * @param extension the extension to register 3056 * @return true if the extension was registered, false otherwise 3057 * @see #unregisterExtension(name) 3058 */ 3059 this.registerExtension = function(name, extension) { 3060 if (arguments.length < 2) { 3061 throw 'Illegal arguments number: required 2, got ' + arguments.length; 3062 } 3063 if (!_isString(name)) { 3064 throw 'Illegal argument type: extension name must be a string'; 3065 } 3066 3067 var existing = false; 3068 for (var i = 0; i < _extensions.length; ++i) { 3069 var existingExtension = _extensions[i]; 3070 if (existingExtension.name === name) { 3071 existing = true; 3072 break; 3073 } 3074 } 3075 if (!existing) { 3076 _extensions.push({ 3077 name: name, 3078 extension: extension 3079 }); 3080 this._debug('Registered extension', name); 3081 3082 // Callback for extensions 3083 if (_isFunction(extension.registered)) { 3084 extension.registered(name, this); 3085 } 3086 3087 return true; 3088 } else { 3089 this._info('Could not register extension with name', name, 'since another extension with the same name already exists'); 3090 return false; 3091 } 3092 }; 3093 3094 /** 3095 * Unregister an extension previously registered with 3096 * {@link #registerExtension(name, extension)}. 3097 * @param name the name of the extension to unregister. 3098 * @return true if the extension was unregistered, false otherwise 3099 */ 3100 this.unregisterExtension = function(name) { 3101 if (!_isString(name)) { 3102 throw 'Illegal argument type: extension name must be a string'; 3103 } 3104 3105 var unregistered = false; 3106 for (var i = 0; i < _extensions.length; ++i) { 3107 var extension = _extensions[i]; 3108 if (extension.name === name) { 3109 _extensions.splice(i, 1); 3110 unregistered = true; 3111 this._debug('Unregistered extension', name); 3112 3113 // Callback for extensions 3114 var ext = extension.extension; 3115 if (_isFunction(ext.unregistered)) { 3116 ext.unregistered(); 3117 } 3118 3119 break; 3120 } 3121 } 3122 return unregistered; 3123 }; 3124 3125 /** 3126 * Find the extension registered with the given name. 3127 * @param name the name of the extension to find 3128 * @return the extension found or null if no extension with the given name has been registered 3129 */ 3130 this.getExtension = function(name) { 3131 for (var i = 0; i < _extensions.length; ++i) { 3132 var extension = _extensions[i]; 3133 if (extension.name === name) { 3134 return extension.extension; 3135 } 3136 } 3137 return null; 3138 }; 3139 3140 /** 3141 * Returns the name assigned to this CometD object, or the string 'default'
3142 * if no name has been explicitly passed as parameter to the constructor. 3143 */ 3144 this.getName = function() { 3145 return _name; 3146 }; 3147 3148 /** 3149 * Returns the clientId assigned by the Bayeux server during handshake. 3150 */ 3151 this.getClientId = function() { 3152 return _clientId; 3153 }; 3154 3155 /** 3156 * Returns the URL of the Bayeux server. 3157 */ 3158 this.getURL = function() { 3159 if (_transport) { 3160 var url = _transport.getURL(); 3161 if (url) { 3162 return url; 3163 } 3164 url = _config.urls[_transport.getType()]; 3165 if (url) { 3166 return url; 3167 } 3168 } 3169 return _config.url; 3170 }; 3171 3172 this.getTransport = function() { 3173 return _transport; 3174 }; 3175 3176 this.getConfiguration = function() { 3177 return this._mixin(true, {}, _config); 3178 }; 3179 3180 this.getAdvice = function() { 3181 return this._mixin(true, {}, _advice); 3182 }; 3183 3184 // Initialize transports. 3185 if (window.WebSocket) { 3186 this.registerTransport('websocket', new WebSocketTransport()); 3187 } 3188 this.registerTransport('long-polling', new LongPollingTransport()); 3189 this.registerTransport('callback-polling', new CallbackPollingTransport()); 3190 }; 3191 3192 var _z85EncodeTable = [ 3193 '0', '1', '2', '3', '4', '5', '6', '7', '8', '9', 3194 'a', 'b', 'c', 'd', 'e', 'f', 'g', 'h', 'i', 'j', 3195 'k', 'l', 'm', 'n', 'o', 'p', 'q', 'r', 's', 't', 3196 'u', 'v', 'w', 'x', 'y', 'z', 'A', 'B', 'C', 'D', 3197 'E', 'F', 'G', 'H', 'I', 'J', 'K', 'L', 'M', 'N', 3198 'O', 'P', 'Q', 'R', 'S', 'T', 'U', 'V', 'W', 'X', 3199 'Y', 'Z', '.', '-', ':', '+', '=', '^', '!', '/', 3200 '*', '?', '&', '<', '>', '(', ')', '[', ']', '{', 3201 '}', '@', '%', '$', '#' 3202 ]; 3203 var _z85DecodeTable = [ 3204 0x00, 0x44, 0x00, 0x54, 0x53, 0x52, 0x48, 0x00, 3205 0x4B, 0x4C, 0x46, 0x41, 0x00, 0x3F, 0x3E, 0x45, 3206 0x00, 0x01, 0x02, 0x03, 0x04, 0x05, 0x06, 0x07, 3207 0x08, 0x09, 0x40, 0x00, 0x49, 0x42, 0x4A, 0x47, 3208 0x51, 0x24, 0x25, 0x26, 0x27, 0x28, 0x29, 0x2A, 3209 0x2B, 0x2C, 0x2D, 0x2E, 0x2F, 0x30, 0x31, 0x32,
3210 0x33, 0x34, 0x35, 0x36, 0x37, 0x38, 0x39, 0x3A, 3211 0x3B, 0x3C, 0x3D, 0x4D, 0x00, 0x4E, 0x43, 0x00, 3212 0x00, 0x0A, 0x0B, 0x0C, 0x0D, 0x0E, 0x0F, 0x10, 3213 0x11, 0x12, 0x13, 0x14, 0x15, 0x16, 0x17, 0x18, 3214 0x19, 0x1A, 0x1B, 0x1C, 0x1D, 0x1E, 0x1F, 0x20, 3215 0x21, 0x22, 0x23, 0x4F, 0x00, 0x50, 0x00, 0x00 3216 ]; 3217 var Z85 = { 3218 encode: function(bytes) { 3219 var buffer = null; 3220 if (bytes instanceof ArrayBuffer) { 3221 buffer = bytes; 3222 } else if (bytes.buffer instanceof ArrayBuffer) { 3223 buffer = bytes.buffer; 3224 } else if (Array.isArray(bytes)) { 3225 buffer = new Uint8Array(bytes).buffer; 3226 } 3227 if (buffer == null) { 3228 throw 'Cannot Z85 encode ' + bytes; 3229 } 3230 3231 var length = buffer.byteLength; 3232 var remainder = length % 4; 3233 var padding = 4 - (remainder === 0 ? 4 : remainder); 3234 var view = new DataView(buffer); 3235 var result = ''; 3236 var value = 0; 3237 for (var i = 0; i < length + padding; ++i) { 3238 var isPadding = i >= length; 3239 value = value * 256 + (isPadding ? 0 : view.getUint8(i)); 3240 if ((i + 1) % 4 === 0) { 3241 var divisor = 85 * 85 * 85 * 85; 3242 for (var j = 5; j > 0; --j) { 3243 if (!isPadding || j > padding) { 3244 var code = Math.floor(value / divisor) % 85; 3245 result += _z85EncodeTable[code]; 3246 } 3247 divisor /= 85; 3248 } 3249 value = 0; 3250 } 3251 } 3252 3253 return result; 3254 }, 3255 decode: function(string) { 3256 var remainder = string.length % 5; 3257 var padding = 5 - (remainder === 0 ? 5 : remainder); 3258 for (var p = 0; p < padding; ++p) { 3259 string += _z85EncodeTable[_z85EncodeTable.length - 1]; 3260 } 3261 var length = string.length; 3262 3263 var buffer = new ArrayBuffer((length * 4 / 5) - padding); 3264 var view = new DataView(buffer); 3265 var value = 0; 3266 var charIdx = 0; 3267 var byteIdx = 0; 3268 for (var i = 0; i < length; ++i) { 3269 var code = string.charCodeAt(charIdx++) - 32; 3270 value = value * 85 + _z85DecodeTable[code]; 3271 if (charIdx % 5 === 0) { 3272 var divisor = 256 * 256 * 256; 3273 while (divisor >= 1) { 3274 if (byteIdx < view.byteLength) { 3275 view.setUint8(byteIdx++, Math.floor(value / divisor) % 256); 3276 } 3277 divisor /= 256; 3278 } 3279 value = 0; 3280 } 3281 } 3282 3283 return buffer; 3284 } 3285 }; 3286 3287 return { 3288 CometD: CometD, 3289 Transport: Transport, 3290 RequestTransport: RequestTransport,
3291 LongPollingTransport: LongPollingTransport, 3292 CallbackPollingTransport: CallbackPollingTransport, 3293 WebSocketTransport: WebSocketTransport, 3294 Utils: Utils, 3295 Z85: Z85 3296 }; 3297}));
Line numbers count LF bytes from the start of the resource, as the search results do. Vendor segments are library code the classifier recognised; they are stored but not indexed. Bytes are shown as Latin1 characters, one per byte.