node_modules updated, buefy add to lib folder
This commit is contained in:
+6
-3
@@ -43,9 +43,11 @@ export declare class AjaxObservable<T> extends Observable<T> {
|
||||
* url, headers, etc or a string for a URL.
|
||||
*
|
||||
* ## Example
|
||||
* ```javascript
|
||||
* source = Rx.Observable.ajax('/products');
|
||||
* source = Rx.Observable.ajax({ url: 'products', method: 'GET' });
|
||||
* ```ts
|
||||
* import { ajax } from 'rxjs/ajax';
|
||||
*
|
||||
* const source1 = ajax('/products');
|
||||
* const source2 = ajax({ url: 'products', method: 'GET' });
|
||||
* ```
|
||||
*
|
||||
* @param {string|Object} request Can be one of the following:
|
||||
@@ -87,6 +89,7 @@ export declare class AjaxSubscriber<T> extends Subscriber<Event> {
|
||||
private send;
|
||||
private serializeBody;
|
||||
private setHeaders;
|
||||
private getHeader;
|
||||
private setupEvents;
|
||||
unsubscribe(): void;
|
||||
}
|
||||
|
||||
+50
-46
@@ -14,8 +14,6 @@ var __extends = (this && this.__extends) || (function () {
|
||||
})();
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
var root_1 = require("../../util/root");
|
||||
var tryCatch_1 = require("../../util/tryCatch");
|
||||
var errorObject_1 = require("../../util/errorObject");
|
||||
var Observable_1 = require("../../Observable");
|
||||
var Subscriber_1 = require("../../Subscriber");
|
||||
var map_1 = require("../../operators/map");
|
||||
@@ -140,47 +138,39 @@ var AjaxSubscriber = (function (_super) {
|
||||
_this.request = request;
|
||||
_this.done = false;
|
||||
var headers = request.headers = request.headers || {};
|
||||
if (!request.crossDomain && !headers['X-Requested-With']) {
|
||||
if (!request.crossDomain && !_this.getHeader(headers, 'X-Requested-With')) {
|
||||
headers['X-Requested-With'] = 'XMLHttpRequest';
|
||||
}
|
||||
if (!('Content-Type' in headers) && !(root_1.root.FormData && request.body instanceof root_1.root.FormData) && typeof request.body !== 'undefined') {
|
||||
var contentTypeHeader = _this.getHeader(headers, 'Content-Type');
|
||||
if (!contentTypeHeader && !(root_1.root.FormData && request.body instanceof root_1.root.FormData) && typeof request.body !== 'undefined') {
|
||||
headers['Content-Type'] = 'application/x-www-form-urlencoded; charset=UTF-8';
|
||||
}
|
||||
request.body = _this.serializeBody(request.body, request.headers['Content-Type']);
|
||||
request.body = _this.serializeBody(request.body, _this.getHeader(request.headers, 'Content-Type'));
|
||||
_this.send();
|
||||
return _this;
|
||||
}
|
||||
AjaxSubscriber.prototype.next = function (e) {
|
||||
this.done = true;
|
||||
var _a = this, xhr = _a.xhr, request = _a.request, destination = _a.destination;
|
||||
var response = new AjaxResponse(e, xhr, request);
|
||||
if (response.response === errorObject_1.errorObject) {
|
||||
destination.error(errorObject_1.errorObject.e);
|
||||
var result;
|
||||
try {
|
||||
result = new AjaxResponse(e, xhr, request);
|
||||
}
|
||||
else {
|
||||
destination.next(response);
|
||||
catch (err) {
|
||||
return destination.error(err);
|
||||
}
|
||||
destination.next(result);
|
||||
};
|
||||
AjaxSubscriber.prototype.send = function () {
|
||||
var _a = this, request = _a.request, _b = _a.request, user = _b.user, method = _b.method, url = _b.url, async = _b.async, password = _b.password, headers = _b.headers, body = _b.body;
|
||||
var createXHR = request.createXHR;
|
||||
var xhr = tryCatch_1.tryCatch(createXHR).call(request);
|
||||
if (xhr === errorObject_1.errorObject) {
|
||||
this.error(errorObject_1.errorObject.e);
|
||||
}
|
||||
else {
|
||||
this.xhr = xhr;
|
||||
try {
|
||||
var xhr = this.xhr = request.createXHR();
|
||||
this.setupEvents(xhr, request);
|
||||
var result = void 0;
|
||||
if (user) {
|
||||
result = tryCatch_1.tryCatch(xhr.open).call(xhr, method, url, async, user, password);
|
||||
xhr.open(method, url, async, user, password);
|
||||
}
|
||||
else {
|
||||
result = tryCatch_1.tryCatch(xhr.open).call(xhr, method, url, async);
|
||||
}
|
||||
if (result === errorObject_1.errorObject) {
|
||||
this.error(errorObject_1.errorObject.e);
|
||||
return null;
|
||||
xhr.open(method, url, async);
|
||||
}
|
||||
if (async) {
|
||||
xhr.timeout = request.timeout;
|
||||
@@ -190,13 +180,16 @@ var AjaxSubscriber = (function (_super) {
|
||||
xhr.withCredentials = !!request.withCredentials;
|
||||
}
|
||||
this.setHeaders(xhr, headers);
|
||||
result = body ? tryCatch_1.tryCatch(xhr.send).call(xhr, body) : tryCatch_1.tryCatch(xhr.send).call(xhr);
|
||||
if (result === errorObject_1.errorObject) {
|
||||
this.error(errorObject_1.errorObject.e);
|
||||
return null;
|
||||
if (body) {
|
||||
xhr.send(body);
|
||||
}
|
||||
else {
|
||||
xhr.send();
|
||||
}
|
||||
}
|
||||
return xhr;
|
||||
catch (err) {
|
||||
this.error(err);
|
||||
}
|
||||
};
|
||||
AjaxSubscriber.prototype.serializeBody = function (body, contentType) {
|
||||
if (!body || typeof body === 'string') {
|
||||
@@ -227,6 +220,14 @@ var AjaxSubscriber = (function (_super) {
|
||||
}
|
||||
}
|
||||
};
|
||||
AjaxSubscriber.prototype.getHeader = function (headers, headerName) {
|
||||
for (var key in headers) {
|
||||
if (key.toLowerCase() === headerName.toLowerCase()) {
|
||||
return headers[key];
|
||||
}
|
||||
}
|
||||
return undefined;
|
||||
};
|
||||
AjaxSubscriber.prototype.setupEvents = function (xhr, request) {
|
||||
var progressSubscriber = request.progressSubscriber;
|
||||
function xhrTimeout(e) {
|
||||
@@ -234,13 +235,14 @@ var AjaxSubscriber = (function (_super) {
|
||||
if (progressSubscriber) {
|
||||
progressSubscriber.error(e);
|
||||
}
|
||||
var ajaxTimeoutError = new exports.AjaxTimeoutError(this, request);
|
||||
if (ajaxTimeoutError.response === errorObject_1.errorObject) {
|
||||
subscriber.error(errorObject_1.errorObject.e);
|
||||
var error;
|
||||
try {
|
||||
error = new exports.AjaxTimeoutError(this, request);
|
||||
}
|
||||
else {
|
||||
subscriber.error(ajaxTimeoutError);
|
||||
catch (err) {
|
||||
error = err;
|
||||
}
|
||||
subscriber.error(error);
|
||||
}
|
||||
xhr.ontimeout = xhrTimeout;
|
||||
xhrTimeout.request = request;
|
||||
@@ -267,13 +269,14 @@ var AjaxSubscriber = (function (_super) {
|
||||
if (progressSubscriber) {
|
||||
progressSubscriber.error(e);
|
||||
}
|
||||
var ajaxError = new exports.AjaxError('ajax error', this, request);
|
||||
if (ajaxError.response === errorObject_1.errorObject) {
|
||||
subscriber.error(errorObject_1.errorObject.e);
|
||||
var error;
|
||||
try {
|
||||
error = new exports.AjaxError('ajax error', this, request);
|
||||
}
|
||||
else {
|
||||
subscriber.error(ajaxError);
|
||||
catch (err) {
|
||||
error = err;
|
||||
}
|
||||
subscriber.error(error);
|
||||
};
|
||||
xhr.onerror = xhrError_1;
|
||||
xhrError_1.request = request;
|
||||
@@ -306,13 +309,14 @@ var AjaxSubscriber = (function (_super) {
|
||||
if (progressSubscriber) {
|
||||
progressSubscriber.error(e);
|
||||
}
|
||||
var ajaxError = new exports.AjaxError('ajax error ' + status_1, this, request);
|
||||
if (ajaxError.response === errorObject_1.errorObject) {
|
||||
subscriber.error(errorObject_1.errorObject.e);
|
||||
var error = void 0;
|
||||
try {
|
||||
error = new exports.AjaxError('ajax error ' + status_1, this, request);
|
||||
}
|
||||
else {
|
||||
subscriber.error(ajaxError);
|
||||
catch (err) {
|
||||
error = err;
|
||||
}
|
||||
subscriber.error(error);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -367,7 +371,7 @@ function parseJson(xhr) {
|
||||
function parseXhrResponse(responseType, xhr) {
|
||||
switch (responseType) {
|
||||
case 'json':
|
||||
return tryCatch_1.tryCatch(parseJson)(xhr);
|
||||
return parseJson(xhr);
|
||||
case 'xml':
|
||||
return xhr.responseXML;
|
||||
case 'text':
|
||||
@@ -381,4 +385,4 @@ function AjaxTimeoutErrorImpl(xhr, request) {
|
||||
return this;
|
||||
}
|
||||
exports.AjaxTimeoutError = AjaxTimeoutErrorImpl;
|
||||
//# sourceMappingURL=AjaxObservable.js.map
|
||||
//# sourceMappingURL=AjaxObservable.js.map
|
||||
+1
-1
File diff suppressed because one or more lines are too long
+90
-5
@@ -4,6 +4,96 @@ import { Observable } from '../../Observable';
|
||||
import { Subscription } from '../../Subscription';
|
||||
import { Operator } from '../../Operator';
|
||||
import { Observer, NextObserver } from '../../types';
|
||||
/**
|
||||
* WebSocketSubjectConfig is a plain Object that allows us to make our
|
||||
* webSocket configurable.
|
||||
*
|
||||
* <span class="informal">Provides flexibility to {@link webSocket}</span>
|
||||
*
|
||||
* It defines a set of properties to provide custom behavior in specific
|
||||
* moments of the socket's lifecycle. When the connection opens we can
|
||||
* use `openObserver`, when the connection is closed `closeObserver`, if we
|
||||
* are interested in listening for data comming from server: `deserializer`,
|
||||
* which allows us to customize the deserialization strategy of data before passing it
|
||||
* to the socket client. By default `deserializer` is going to apply `JSON.parse` to each message comming
|
||||
* from the Server.
|
||||
*
|
||||
* ## Example
|
||||
* **deserializer**, the default for this property is `JSON.parse` but since there are just two options
|
||||
* for incomming data, either be text or binarydata. We can apply a custom deserialization strategy
|
||||
* or just simply skip the default behaviour.
|
||||
* ```ts
|
||||
* import { webSocket } from 'rxjs/webSocket';
|
||||
*
|
||||
* const wsSubject = webSocket({
|
||||
* url: 'ws://localhost:8081',
|
||||
* //Apply any transformation of your choice.
|
||||
* deserializer: ({data}) => data
|
||||
* });
|
||||
*
|
||||
* wsSubject.subscribe(console.log);
|
||||
*
|
||||
* // Let's suppose we have this on the Server: ws.send("This is a msg from the server")
|
||||
* //output
|
||||
* //
|
||||
* // This is a msg from the server
|
||||
* ```
|
||||
*
|
||||
* **serializer** allows us tom apply custom serialization strategy but for the outgoing messages
|
||||
* ```ts
|
||||
* import { webSocket } from 'rxjs/webSocket';
|
||||
*
|
||||
* const wsSubject = webSocket({
|
||||
* url: 'ws://localhost:8081',
|
||||
* //Apply any transformation of your choice.
|
||||
* serializer: msg => JSON.stringify({channel: "webDevelopment", msg: msg})
|
||||
* });
|
||||
*
|
||||
* wsSubject.subscribe(() => subject.next("msg to the server"));
|
||||
*
|
||||
* // Let's suppose we have this on the Server: ws.send("This is a msg from the server")
|
||||
* //output
|
||||
* //
|
||||
* // {"channel":"webDevelopment","msg":"msg to the server"}
|
||||
* ```
|
||||
*
|
||||
* **closeObserver** allows us to set a custom error when an error raise up.
|
||||
* ```ts
|
||||
* import { webSocket } from 'rxjs/webSocket';
|
||||
*
|
||||
* const wsSubject = webSocket({
|
||||
* url: 'ws://localhost:8081',
|
||||
* closeObserver: {
|
||||
next(closeEvent) {
|
||||
const customError = { code: 6666, reason: "Custom evil reason" }
|
||||
console.log(`code: ${customError.code}, reason: ${customError.reason}`);
|
||||
}
|
||||
}
|
||||
* });
|
||||
*
|
||||
* //output
|
||||
* // code: 6666, reason: Custom evil reason
|
||||
* ```
|
||||
*
|
||||
* **openObserver**, Let's say we need to make some kind of init task before sending/receiving msgs to the
|
||||
* webSocket or sending notification that the connection was successful, this is when
|
||||
* openObserver is usefull for.
|
||||
* ```ts
|
||||
* import { webSocket } from 'rxjs/webSocket';
|
||||
*
|
||||
* const wsSubject = webSocket({
|
||||
* url: 'ws://localhost:8081',
|
||||
* openObserver: {
|
||||
* next: () => {
|
||||
* console.log('connetion ok');
|
||||
* }
|
||||
* },
|
||||
* });
|
||||
*
|
||||
* //output
|
||||
* // connetion ok`
|
||||
* ```
|
||||
* */
|
||||
export interface WebSocketSubjectConfig<T> {
|
||||
/** The url of the socket server to connect to */
|
||||
url: string;
|
||||
@@ -46,11 +136,6 @@ export interface WebSocketSubjectConfig<T> {
|
||||
binaryType?: 'blob' | 'arraybuffer';
|
||||
}
|
||||
export declare type WebSocketMessage = string | ArrayBuffer | Blob | ArrayBufferView;
|
||||
/**
|
||||
* We need this JSDoc comment for affecting ESDoc.
|
||||
* @extends {Ignored}
|
||||
* @hide true
|
||||
*/
|
||||
export declare class WebSocketSubject<T> extends AnonymousSubject<T> {
|
||||
private _config;
|
||||
/** @deprecated This is an internal implementation detail, do not use. */
|
||||
|
||||
+34
-35
@@ -29,8 +29,6 @@ var Subscriber_1 = require("../../Subscriber");
|
||||
var Observable_1 = require("../../Observable");
|
||||
var Subscription_1 = require("../../Subscription");
|
||||
var ReplaySubject_1 = require("../../ReplaySubject");
|
||||
var tryCatch_1 = require("../../util/tryCatch");
|
||||
var errorObject_1 = require("../../util/errorObject");
|
||||
var DEFAULT_WEBSOCKET_CONFIG = {
|
||||
url: '',
|
||||
deserializer: function (e) { return JSON.parse(e.data); },
|
||||
@@ -84,29 +82,28 @@ var WebSocketSubject = (function (_super) {
|
||||
WebSocketSubject.prototype.multiplex = function (subMsg, unsubMsg, messageFilter) {
|
||||
var self = this;
|
||||
return new Observable_1.Observable(function (observer) {
|
||||
var result = tryCatch_1.tryCatch(subMsg)();
|
||||
if (result === errorObject_1.errorObject) {
|
||||
observer.error(errorObject_1.errorObject.e);
|
||||
try {
|
||||
self.next(subMsg());
|
||||
}
|
||||
else {
|
||||
self.next(result);
|
||||
catch (err) {
|
||||
observer.error(err);
|
||||
}
|
||||
var subscription = self.subscribe(function (x) {
|
||||
var result = tryCatch_1.tryCatch(messageFilter)(x);
|
||||
if (result === errorObject_1.errorObject) {
|
||||
observer.error(errorObject_1.errorObject.e);
|
||||
try {
|
||||
if (messageFilter(x)) {
|
||||
observer.next(x);
|
||||
}
|
||||
}
|
||||
else if (result) {
|
||||
observer.next(x);
|
||||
catch (err) {
|
||||
observer.error(err);
|
||||
}
|
||||
}, function (err) { return observer.error(err); }, function () { return observer.complete(); });
|
||||
return function () {
|
||||
var result = tryCatch_1.tryCatch(unsubMsg)();
|
||||
if (result === errorObject_1.errorObject) {
|
||||
observer.error(errorObject_1.errorObject.e);
|
||||
try {
|
||||
self.next(unsubMsg());
|
||||
}
|
||||
else {
|
||||
self.next(result);
|
||||
catch (err) {
|
||||
observer.error(err);
|
||||
}
|
||||
subscription.unsubscribe();
|
||||
};
|
||||
@@ -137,6 +134,12 @@ var WebSocketSubject = (function (_super) {
|
||||
}
|
||||
});
|
||||
socket.onopen = function (e) {
|
||||
var _socket = _this._socket;
|
||||
if (!_socket) {
|
||||
socket.close();
|
||||
_this._resetState();
|
||||
return;
|
||||
}
|
||||
var openObserver = _this._config.openObserver;
|
||||
if (openObserver) {
|
||||
openObserver.next(e);
|
||||
@@ -144,13 +147,13 @@ var WebSocketSubject = (function (_super) {
|
||||
var queue = _this.destination;
|
||||
_this.destination = Subscriber_1.Subscriber.create(function (x) {
|
||||
if (socket.readyState === 1) {
|
||||
var serializer = _this._config.serializer;
|
||||
var msg = tryCatch_1.tryCatch(serializer)(x);
|
||||
if (msg === errorObject_1.errorObject) {
|
||||
_this.destination.error(errorObject_1.errorObject.e);
|
||||
return;
|
||||
try {
|
||||
var serializer = _this._config.serializer;
|
||||
socket.send(serializer(x));
|
||||
}
|
||||
catch (e) {
|
||||
_this.destination.error(e);
|
||||
}
|
||||
socket.send(msg);
|
||||
}
|
||||
}, function (e) {
|
||||
var closingObserver = _this._config.closingObserver;
|
||||
@@ -194,13 +197,12 @@ var WebSocketSubject = (function (_super) {
|
||||
}
|
||||
};
|
||||
socket.onmessage = function (e) {
|
||||
var deserializer = _this._config.deserializer;
|
||||
var result = tryCatch_1.tryCatch(deserializer)(e);
|
||||
if (result === errorObject_1.errorObject) {
|
||||
observer.error(errorObject_1.errorObject.e);
|
||||
try {
|
||||
var deserializer = _this._config.deserializer;
|
||||
observer.next(deserializer(e));
|
||||
}
|
||||
else {
|
||||
observer.next(result);
|
||||
catch (err) {
|
||||
observer.error(err);
|
||||
}
|
||||
};
|
||||
};
|
||||
@@ -226,17 +228,14 @@ var WebSocketSubject = (function (_super) {
|
||||
return subscriber;
|
||||
};
|
||||
WebSocketSubject.prototype.unsubscribe = function () {
|
||||
var _a = this, source = _a.source, _socket = _a._socket;
|
||||
var _socket = this._socket;
|
||||
if (_socket && _socket.readyState === 1) {
|
||||
_socket.close();
|
||||
this._resetState();
|
||||
}
|
||||
this._resetState();
|
||||
_super.prototype.unsubscribe.call(this);
|
||||
if (!source) {
|
||||
this.destination = new ReplaySubject_1.ReplaySubject();
|
||||
}
|
||||
};
|
||||
return WebSocketSubject;
|
||||
}(Subject_1.AnonymousSubject));
|
||||
exports.WebSocketSubject = WebSocketSubject;
|
||||
//# sourceMappingURL=WebSocketSubject.js.map
|
||||
//# sourceMappingURL=WebSocketSubject.js.map
|
||||
+1
-1
File diff suppressed because one or more lines are too long
+67
-4
@@ -5,15 +5,78 @@ import { AjaxCreationMethod } from './AjaxObservable';
|
||||
* It creates an observable for an Ajax request with either a request object with
|
||||
* url, headers, etc or a string for a URL.
|
||||
*
|
||||
* ## Using ajax.getJSON() to fetch data from API.
|
||||
* ```javascript
|
||||
*
|
||||
* ## Using ajax() to fetch the response object that is being returned from API.
|
||||
* ```ts
|
||||
* import { ajax } from 'rxjs/ajax';
|
||||
* import { map, catchError } from 'rxjs/operators';
|
||||
* import { of } from 'rxjs';
|
||||
*
|
||||
* const obs$ = ajax(`https://api.github.com/users?per_page=5`).pipe(
|
||||
* map(userResponse => console.log('users: ', userResponse)),
|
||||
* catchError(error => {
|
||||
* console.log('error: ', error);
|
||||
* return of(error);
|
||||
* })
|
||||
* );
|
||||
*
|
||||
* ```
|
||||
*
|
||||
* ## Using ajax.getJSON() to fetch data from API.
|
||||
* ```ts
|
||||
* import { ajax } from 'rxjs/ajax';
|
||||
* import { map, catchError } from 'rxjs/operators';
|
||||
* import { of } from 'rxjs';
|
||||
*
|
||||
* const obs$ = ajax.getJSON(`https://api.github.com/users?per_page=5`).pipe(
|
||||
* map(userResponse => console.log('users: ', userResponse)),
|
||||
* catchError(error => console.log('error: ', error))
|
||||
* ));
|
||||
* catchError(error => {
|
||||
* console.log('error: ', error);
|
||||
* return of(error);
|
||||
* })
|
||||
* );
|
||||
*
|
||||
* ```
|
||||
*
|
||||
* ## Using ajax() with object as argument and method POST with a two seconds delay.
|
||||
* ```ts
|
||||
* import { ajax } from 'rxjs/ajax';
|
||||
* import { of } from 'rxjs';
|
||||
*
|
||||
* const users = ajax({
|
||||
* url: 'https://httpbin.org/delay/2',
|
||||
* method: 'POST',
|
||||
* headers: {
|
||||
* 'Content-Type': 'application/json',
|
||||
* 'rxjs-custom-header': 'Rxjs'
|
||||
* },
|
||||
* body: {
|
||||
* rxjs: 'Hello World!'
|
||||
* }
|
||||
* }).pipe(
|
||||
* map(response => console.log('response: ', response)),
|
||||
* catchError(error => {
|
||||
* console.log('error: ', error);
|
||||
* return of(error);
|
||||
* })
|
||||
* );
|
||||
*
|
||||
* ```
|
||||
*
|
||||
* ## Using ajax() to fetch. An error object that is being returned from the request.
|
||||
* ```ts
|
||||
* import { ajax } from 'rxjs/ajax';
|
||||
* import { map, catchError } from 'rxjs/operators';
|
||||
* import { of } from 'rxjs';
|
||||
*
|
||||
* const obs$ = ajax(`https://api.github.com/404`).pipe(
|
||||
* map(userResponse => console.log('users: ', userResponse)),
|
||||
* catchError(error => {
|
||||
* console.log('error: ', error);
|
||||
* return of(error);
|
||||
* })
|
||||
* );
|
||||
*
|
||||
* ```
|
||||
*/
|
||||
export declare const ajax: AjaxCreationMethod;
|
||||
|
||||
+1
-1
@@ -2,4 +2,4 @@
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
var AjaxObservable_1 = require("./AjaxObservable");
|
||||
exports.ajax = AjaxObservable_1.AjaxObservable.create;
|
||||
//# sourceMappingURL=ajax.js.map
|
||||
//# sourceMappingURL=ajax.js.map
|
||||
+1
-1
@@ -1 +1 @@
|
||||
{"version":3,"file":"ajax.js","sources":["../../../src/internal/observable/dom/ajax.ts"],"names":[],"mappings":";;AAAA,mDAAwE;AAkB3D,QAAA,IAAI,GAAuB,+BAAc,CAAC,MAAM,CAAC"}
|
||||
{"version":3,"file":"ajax.js","sources":["../../../src/internal/observable/dom/ajax.ts"],"names":[],"mappings":";;AAAA,mDAAwE;AAiF3D,QAAA,IAAI,GAAuB,+BAAc,CAAC,MAAM,CAAC"}
|
||||
|
||||
+52
@@ -0,0 +1,52 @@
|
||||
import { Observable } from '../../Observable';
|
||||
/**
|
||||
* Uses [the Fetch API](https://developer.mozilla.org/en-US/docs/Web/API/Fetch_API) to
|
||||
* make an HTTP request.
|
||||
*
|
||||
* **WARNING** Parts of the fetch API are still experimental. `AbortController` is
|
||||
* required for this implementation to work and use cancellation appropriately.
|
||||
*
|
||||
* Will automatically set up an internal [AbortController](https://developer.mozilla.org/en-US/docs/Web/API/AbortController)
|
||||
* in order to teardown the internal `fetch` when the subscription tears down.
|
||||
*
|
||||
* If a `signal` is provided via the `init` argument, it will behave like it usually does with
|
||||
* `fetch`. If the provided `signal` aborts, the error that `fetch` normally rejects with
|
||||
* in that scenario will be emitted as an error from the observable.
|
||||
*
|
||||
* ### Basic Use
|
||||
*
|
||||
* ```ts
|
||||
* import { of } from 'rxjs';
|
||||
* import { fromFetch } from 'rxjs/fetch';
|
||||
* import { switchMap, catchError } from 'rxjs/operators';
|
||||
*
|
||||
* const data$ = fromFetch('https://api.github.com/users?per_page=5').pipe(
|
||||
* switchMap(response => {
|
||||
* if (response.ok) {
|
||||
* // OK return data
|
||||
* return response.json();
|
||||
* } else {
|
||||
* // Server is returning a status requiring the client to try something else.
|
||||
* return of({ error: true, message: `Error ${response.status}` });
|
||||
* }
|
||||
* }),
|
||||
* catchError(err => {
|
||||
* // Network or other error, handle appropriately
|
||||
* console.error(err);
|
||||
* return of({ error: true, message: err.message })
|
||||
* })
|
||||
* );
|
||||
*
|
||||
* data$.subscribe({
|
||||
* next: result => console.log(result),
|
||||
* complete: () => console.log('done')
|
||||
* })
|
||||
* ```
|
||||
*
|
||||
* @param input The resource you would like to fetch. Can be a url or a request object.
|
||||
* @param init A configuration object for the fetch.
|
||||
* [See MDN for more details](https://developer.mozilla.org/en-US/docs/Web/API/WindowOrWorkerGlobalScope/fetch#Parameters)
|
||||
* @returns An Observable, that when subscribed to performs an HTTP request using the native `fetch`
|
||||
* function. The {@link Subscription} is tied to an `AbortController` for the the fetch.
|
||||
*/
|
||||
export declare function fromFetch(input: string | Request, init?: RequestInit): Observable<Response>;
|
||||
+44
@@ -0,0 +1,44 @@
|
||||
"use strict";
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
var Observable_1 = require("../../Observable");
|
||||
function fromFetch(input, init) {
|
||||
return new Observable_1.Observable(function (subscriber) {
|
||||
var controller = new AbortController();
|
||||
var signal = controller.signal;
|
||||
var outerSignalHandler;
|
||||
var abortable = true;
|
||||
var unsubscribed = false;
|
||||
if (init) {
|
||||
if (init.signal) {
|
||||
outerSignalHandler = function () {
|
||||
if (!signal.aborted) {
|
||||
controller.abort();
|
||||
}
|
||||
};
|
||||
init.signal.addEventListener('abort', outerSignalHandler);
|
||||
}
|
||||
init.signal = signal;
|
||||
}
|
||||
else {
|
||||
init = { signal: signal };
|
||||
}
|
||||
fetch(input, init).then(function (response) {
|
||||
abortable = false;
|
||||
subscriber.next(response);
|
||||
subscriber.complete();
|
||||
}).catch(function (err) {
|
||||
abortable = false;
|
||||
if (!unsubscribed) {
|
||||
subscriber.error(err);
|
||||
}
|
||||
});
|
||||
return function () {
|
||||
unsubscribed = true;
|
||||
if (abortable) {
|
||||
controller.abort();
|
||||
}
|
||||
};
|
||||
});
|
||||
}
|
||||
exports.fromFetch = fromFetch;
|
||||
//# sourceMappingURL=fetch.js.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"fetch.js","sources":["../../../src/internal/observable/dom/fetch.ts"],"names":[],"mappings":";;AAAA,+CAA8C;AAoD9C,SAAgB,SAAS,CAAC,KAAuB,EAAE,IAAkB;IACnE,OAAO,IAAI,uBAAU,CAAW,UAAA,UAAU;QACxC,IAAM,UAAU,GAAG,IAAI,eAAe,EAAE,CAAC;QACzC,IAAM,MAAM,GAAG,UAAU,CAAC,MAAM,CAAC;QACjC,IAAI,kBAA8B,CAAC;QACnC,IAAI,SAAS,GAAG,IAAI,CAAC;QACrB,IAAI,YAAY,GAAG,KAAK,CAAC;QAEzB,IAAI,IAAI,EAAE;YAER,IAAI,IAAI,CAAC,MAAM,EAAE;gBACf,kBAAkB,GAAG;oBACnB,IAAI,CAAC,MAAM,CAAC,OAAO,EAAE;wBACnB,UAAU,CAAC,KAAK,EAAE,CAAC;qBACpB;gBACH,CAAC,CAAC;gBACF,IAAI,CAAC,MAAM,CAAC,gBAAgB,CAAC,OAAO,EAAE,kBAAkB,CAAC,CAAC;aAC3D;YACD,IAAI,CAAC,MAAM,GAAG,MAAM,CAAC;SACtB;aAAM;YACL,IAAI,GAAG,EAAE,MAAM,QAAA,EAAE,CAAC;SACnB;QAED,KAAK,CAAC,KAAK,EAAE,IAAI,CAAC,CAAC,IAAI,CAAC,UAAA,QAAQ;YAC9B,SAAS,GAAG,KAAK,CAAC;YAClB,UAAU,CAAC,IAAI,CAAC,QAAQ,CAAC,CAAC;YAC1B,UAAU,CAAC,QAAQ,EAAE,CAAC;QACxB,CAAC,CAAC,CAAC,KAAK,CAAC,UAAA,GAAG;YACV,SAAS,GAAG,KAAK,CAAC;YAClB,IAAI,CAAC,YAAY,EAAE;gBAEjB,UAAU,CAAC,KAAK,CAAC,GAAG,CAAC,CAAC;aACvB;QACH,CAAC,CAAC,CAAC;QAEH,OAAO;YACL,YAAY,GAAG,IAAI,CAAC;YACpB,IAAI,SAAS,EAAE;gBACb,UAAU,CAAC,KAAK,EAAE,CAAC;aACpB;QACH,CAAC,CAAC;IACJ,CAAC,CAAC,CAAC;AACL,CAAC;AA1CD,8BA0CC"}
|
||||
+37
-40
@@ -65,16 +65,17 @@ import { WebSocketSubject, WebSocketSubjectConfig } from './WebSocketSubject';
|
||||
* subscribes and unsubscribes. Server can use them to verify that some kind of messages should start or stop
|
||||
* being forwarded to the client. In case of the above example application, after getting subscription message with proper identifier,
|
||||
* gateway server can decide that it should connect to real sport news service and start forwarding messages from it.
|
||||
* Note that both messages will be sent as returned by the functions, meaning they will have to be serialized manually, just
|
||||
* Note that both messages will be sent as returned by the functions, they are by default serialized using JSON.stringify, just
|
||||
* as messages pushed via `next`. Also bear in mind that these messages will be sent on *every* subscription and
|
||||
* unsubscription. This is potentially dangerous, because one consumer of an Observable may unsubscribe and the server
|
||||
* might stop sending messages, since it got unsubscription message. This needs to be handled
|
||||
* on the server or using {@link publish} on a Observable returned from 'multiplex'.
|
||||
*
|
||||
* Last argument to `multiplex` is a `messageFilter` function which filters out messages
|
||||
* Last argument to `multiplex` is a `messageFilter` function which should return a boolean. It is used to filter out messages
|
||||
* sent by the server to only those that belong to simulated WebSocket stream. For example, server might mark these
|
||||
* messages with some kind of string identifier on a message object and `messageFilter` would return `true`
|
||||
* if there is such identifier on an object emitted by the socket.
|
||||
* if there is such identifier on an object emitted by the socket. Messages which returns `false` in `messageFilter` are simply skipped,
|
||||
* and are not passed down the stream.
|
||||
*
|
||||
* Return value of `multiplex` is an Observable with messages incoming from emulated socket connection. Note that this
|
||||
* is not a `WebSocketSubject`, so calling `next` or `multiplex` again will fail. For pushing values to the
|
||||
@@ -82,71 +83,67 @@ import { WebSocketSubject, WebSocketSubjectConfig } from './WebSocketSubject';
|
||||
*
|
||||
* ### Examples
|
||||
* #### Listening for messages from the server
|
||||
* const subject = Rx.Observable.webSocket('ws://localhost:8081');
|
||||
* ```ts
|
||||
* import { webSocket } from "rxjs/webSocket";
|
||||
* const subject = webSocket("ws://localhost:8081");
|
||||
*
|
||||
* subject.subscribe(
|
||||
* (msg) => console.log('message received: ' + msg), // Called whenever there is a message from the server.
|
||||
* (err) => console.log(err), // Called if at any point WebSocket API signals some kind of error.
|
||||
* msg => console.log('message received: ' + msg), // Called whenever there is a message from the server.
|
||||
* err => console.log(err), // Called if at any point WebSocket API signals some kind of error.
|
||||
* () => console.log('complete') // Called when connection is closed (for whatever reason).
|
||||
* );
|
||||
*
|
||||
* ```
|
||||
*
|
||||
* #### Pushing messages to the server
|
||||
* const subject = Rx.Observable.webSocket('ws://localhost:8081');
|
||||
* ```ts
|
||||
* import { webSocket } from "rxjs/webSocket";
|
||||
* const subject = webSocket('ws://localhost:8081');
|
||||
*
|
||||
* subject.subscribe(); // Note that at least one consumer has to subscribe to
|
||||
* // the created subject - otherwise "nexted" values will be just
|
||||
* // buffered and not sent, since no connection was established!
|
||||
* subject.subscribe();
|
||||
* // Note that at least one consumer has to subscribe to the created subject - otherwise "nexted" values will be just buffered and not sent,
|
||||
* // since no connection was established!
|
||||
*
|
||||
* subject.next(JSON.stringify({message: 'some message'})); // This will send a message to the server
|
||||
* // once a connection is made.
|
||||
* // Remember to serialize sent value first!
|
||||
* subject.next({message: 'some message'});
|
||||
* // This will send a message to the server once a connection is made. Remember value is serialized with JSON.stringify by default!
|
||||
*
|
||||
* subject.complete(); // Closes the connection.
|
||||
*
|
||||
*
|
||||
* subject.error({code: 4000, reason: 'I think our app just broke!'}); // Also closes the connection,
|
||||
* // but let's the server know that
|
||||
* // this closing is caused by some error.
|
||||
*
|
||||
* subject.error({code: 4000, reason: 'I think our app just broke!'});
|
||||
* // Also closes the connection, but let's the server know that this closing is caused by some error.
|
||||
* ```
|
||||
*
|
||||
* #### Multiplexing WebSocket
|
||||
* const subject = Rx.Observable.webSocket('ws://localhost:8081');
|
||||
* ```ts
|
||||
* import { webSocket } from "rxjs/webSocket";
|
||||
* const subject = webSocket('ws://localhost:8081');
|
||||
*
|
||||
* const observableA = subject.multiplex(
|
||||
* () => JSON.stringify({subscribe: 'A'}), // When server gets this message, it will start sending messages for 'A'...
|
||||
* () => JSON.stringify({unsubscribe: 'A'}), // ...and when gets this one, it will stop.
|
||||
* message => message.type === 'A' // Server will tag all messages for 'A' with type property.
|
||||
* () => ({subscribe: 'A'}), // When server gets this message, it will start sending messages for 'A'...
|
||||
* () => ({unsubscribe: 'A'}), // ...and when gets this one, it will stop.
|
||||
* message => message.type === 'A' // If the function returns `true` message is passed down the stream. Skipped if the function returns false.
|
||||
* );
|
||||
*
|
||||
* const observableB = subject.multiplex( // And the same goes for 'B'.
|
||||
* () => JSON.stringify({subscribe: 'B'}),
|
||||
* () => JSON.stringify({unsubscribe: 'B'}),
|
||||
* () => ({subscribe: 'B'}),
|
||||
* () => ({unsubscribe: 'B'}),
|
||||
* message => message.type === 'B'
|
||||
* );
|
||||
*
|
||||
* const subA = observableA.subscribe(messageForA => console.log(messageForA));
|
||||
* // At this moment WebSocket connection
|
||||
* // is established. Server gets '{"subscribe": "A"}'
|
||||
* // message and starts sending messages for 'A',
|
||||
* // At this moment WebSocket connection is established. Server gets '{"subscribe": "A"}' message and starts sending messages for 'A',
|
||||
* // which we log here.
|
||||
*
|
||||
* const subB = observableB.subscribe(messageForB => console.log(messageForB));
|
||||
* // Since we already have a connection,
|
||||
* // we just send '{"subscribe": "B"}' message
|
||||
* // to the server. It starts sending
|
||||
* // messages for 'B', which we log here.
|
||||
* // Since we already have a connection, we just send '{"subscribe": "B"}' message to the server. It starts sending messages for 'B',
|
||||
* // which we log here.
|
||||
*
|
||||
* subB.unsubscribe();
|
||||
* // Message '{"unsubscribe": "B"}' is sent to the
|
||||
* // server, which stops sending 'B' messages.
|
||||
*
|
||||
* subA.unubscribe();
|
||||
* // Message '{"unsubscribe": "A"}' makes the server
|
||||
* // stop sending messages for 'A'. Since there is
|
||||
* // no more subscribers to root Subject, socket
|
||||
* // connection closes.
|
||||
* // Message '{"unsubscribe": "B"}' is sent to the server, which stops sending 'B' messages.
|
||||
*
|
||||
* subA.unsubscribe();
|
||||
* // Message '{"unsubscribe": "A"}' makes the server stop sending messages for 'A'. Since there is no more subscribers to root Subject,
|
||||
* // socket connection closes.
|
||||
* ```
|
||||
*
|
||||
*
|
||||
* @param {string|WebSocketSubjectConfig} urlConfigOrSource The WebSocket endpoint as an url or an object with
|
||||
|
||||
+1
-1
@@ -5,4 +5,4 @@ function webSocket(urlConfigOrSource) {
|
||||
return new WebSocketSubject_1.WebSocketSubject(urlConfigOrSource);
|
||||
}
|
||||
exports.webSocket = webSocket;
|
||||
//# sourceMappingURL=webSocket.js.map
|
||||
//# sourceMappingURL=webSocket.js.map
|
||||
+1
-1
@@ -1 +1 @@
|
||||
{"version":3,"file":"webSocket.js","sources":["../../../src/internal/observable/dom/webSocket.ts"],"names":[],"mappings":";;AAAA,uDAA8E;AA4J9E,SAAgB,SAAS,CAAI,iBAAqD;IAChF,OAAO,IAAI,mCAAgB,CAAI,iBAAiB,CAAC,CAAC;AACpD,CAAC;AAFD,8BAEC"}
|
||||
{"version":3,"file":"webSocket.js","sources":["../../../src/internal/observable/dom/webSocket.ts"],"names":[],"mappings":";;AAAA,uDAA8E;AAyJ9E,SAAgB,SAAS,CAAI,iBAAqD;IAChF,OAAO,IAAI,mCAAgB,CAAI,iBAAiB,CAAC,CAAC;AACpD,CAAC;AAFD,8BAEC"}
|
||||
|
||||
Reference in New Issue
Block a user