init
This commit is contained in:
+7
@@ -0,0 +1,7 @@
|
||||
import { Client, ClientOptions } from '.';
|
||||
export default class BroadcastClient extends Client {
|
||||
private readonly clients;
|
||||
constructor(servers: string[], options?: ClientOptions);
|
||||
private getMethodNames;
|
||||
}
|
||||
//# sourceMappingURL=BroadcastClient.d.ts.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"BroadcastClient.d.ts","sourceRoot":"","sources":["../../../src/client/BroadcastClient.ts"],"names":[],"mappings":"AAAA,OAAO,EAAE,MAAM,EAAE,aAAa,EAAE,MAAM,GAAG,CAAA;AAiBzC,MAAM,CAAC,OAAO,OAAO,eAAgB,SAAQ,MAAM;IACjD,OAAO,CAAC,QAAQ,CAAC,OAAO,CAAU;gBASf,OAAO,EAAE,MAAM,EAAE,EAAE,OAAO,GAAE,aAAkB;IAsCjE,OAAO,CAAC,cAAc;CAkBvB"}
|
||||
+48
@@ -0,0 +1,48 @@
|
||||
"use strict";
|
||||
var __awaiter = (this && this.__awaiter) || function (thisArg, _arguments, P, generator) {
|
||||
function adopt(value) { return value instanceof P ? value : new P(function (resolve) { resolve(value); }); }
|
||||
return new (P || (P = Promise))(function (resolve, reject) {
|
||||
function fulfilled(value) { try { step(generator.next(value)); } catch (e) { reject(e); } }
|
||||
function rejected(value) { try { step(generator["throw"](value)); } catch (e) { reject(e); } }
|
||||
function step(result) { result.done ? resolve(result.value) : adopt(result.value).then(fulfilled, rejected); }
|
||||
step((generator = generator.apply(thisArg, _arguments || [])).next());
|
||||
});
|
||||
};
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
const _1 = require(".");
|
||||
class BroadcastClient extends _1.Client {
|
||||
constructor(servers, options = {}) {
|
||||
super(servers[0], options);
|
||||
const clients = servers.map((server) => new _1.Client(server, options));
|
||||
this.clients = clients;
|
||||
this.getMethodNames().forEach((name) => {
|
||||
this[name] = (...args) => __awaiter(this, void 0, void 0, function* () { return Promise.race(clients.map((client) => __awaiter(this, void 0, void 0, function* () { return client[name](...args); }))); });
|
||||
});
|
||||
this.connect = () => __awaiter(this, void 0, void 0, function* () {
|
||||
yield Promise.all(clients.map((client) => __awaiter(this, void 0, void 0, function* () { return client.connect(); })));
|
||||
});
|
||||
this.disconnect = () => __awaiter(this, void 0, void 0, function* () {
|
||||
yield Promise.all(clients.map((client) => __awaiter(this, void 0, void 0, function* () { return client.disconnect(); })));
|
||||
});
|
||||
this.isConnected = () => clients.map((client) => client.isConnected()).every(Boolean);
|
||||
clients.forEach((client) => {
|
||||
client.on('error', (errorCode, errorMessage, data) => this.emit('error', errorCode, errorMessage, data));
|
||||
});
|
||||
}
|
||||
getMethodNames() {
|
||||
const methodNames = [];
|
||||
const firstClient = this.clients[0];
|
||||
const methods = Object.getOwnPropertyNames(firstClient);
|
||||
methods.push(...Object.getOwnPropertyNames(Object.getPrototypeOf(firstClient)));
|
||||
for (const name of methods) {
|
||||
if (typeof firstClient[name] === 'function' &&
|
||||
name !== 'constructor' &&
|
||||
name !== 'on') {
|
||||
methodNames.push(name);
|
||||
}
|
||||
}
|
||||
return methodNames;
|
||||
}
|
||||
}
|
||||
exports.default = BroadcastClient;
|
||||
//# sourceMappingURL=BroadcastClient.js.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"BroadcastClient.js","sourceRoot":"","sources":["../../../src/client/BroadcastClient.ts"],"names":[],"mappings":";;;;;;;;;;;AAAA,wBAAyC;AAiBzC,MAAqB,eAAgB,SAAQ,SAAM;IAUjD,YAAmB,OAAiB,EAAE,UAAyB,EAAE;QAC/D,KAAK,CAAC,OAAO,CAAC,CAAC,CAAC,EAAE,OAAO,CAAC,CAAA;QAE1B,MAAM,OAAO,GAAa,OAAO,CAAC,GAAG,CACnC,CAAC,MAAM,EAAE,EAAE,CAAC,IAAI,SAAM,CAAC,MAAM,EAAE,OAAO,CAAC,CACxC,CAAA;QAED,IAAI,CAAC,OAAO,GAAG,OAAO,CAAA;QACtB,IAAI,CAAC,cAAc,EAAE,CAAC,OAAO,CAAC,CAAC,IAAY,EAAE,EAAE;YAC7C,IAAI,CAAC,IAAI,CAAC,GAAG,CAAO,GAAG,IAAI,EAAoB,EAAE,gDAG/C,OAAA,OAAO,CAAC,IAAI,CAAC,OAAO,CAAC,GAAG,CAAC,CAAO,MAAM,EAAE,EAAE,gDAAC,OAAA,MAAM,CAAC,IAAI,CAAC,CAAC,GAAG,IAAI,CAAC,CAAA,GAAA,CAAC,CAAC,CAAA,GAAA,CAAA;QAEtE,CAAC,CAAC,CAAA;QAGF,IAAI,CAAC,OAAO,GAAG,GAAwB,EAAE;YACvC,MAAM,OAAO,CAAC,GAAG,CAAC,OAAO,CAAC,GAAG,CAAC,CAAO,MAAM,EAAE,EAAE,gDAAC,OAAA,MAAM,CAAC,OAAO,EAAE,CAAA,GAAA,CAAC,CAAC,CAAA;QACpE,CAAC,CAAA,CAAA;QACD,IAAI,CAAC,UAAU,GAAG,GAAwB,EAAE;YAC1C,MAAM,OAAO,CAAC,GAAG,CAAC,OAAO,CAAC,GAAG,CAAC,CAAO,MAAM,EAAE,EAAE,gDAAC,OAAA,MAAM,CAAC,UAAU,EAAE,CAAA,GAAA,CAAC,CAAC,CAAA;QACvE,CAAC,CAAA,CAAA;QACD,IAAI,CAAC,WAAW,GAAG,GAAY,EAAE,CAC/B,OAAO,CAAC,GAAG,CAAC,CAAC,MAAM,EAAE,EAAE,CAAC,MAAM,CAAC,WAAW,EAAE,CAAC,CAAC,KAAK,CAAC,OAAO,CAAC,CAAA;QAE9D,OAAO,CAAC,OAAO,CAAC,CAAC,MAAM,EAAE,EAAE;YACzB,MAAM,CAAC,EAAE,CAAC,OAAO,EAAE,CAAC,SAAS,EAAE,YAAY,EAAE,IAAI,EAAE,EAAE,CACnD,IAAI,CAAC,IAAI,CAAC,OAAO,EAAE,SAAS,EAAE,YAAY,EAAE,IAAI,CAAC,CAClD,CAAA;QACH,CAAC,CAAC,CAAA;IACJ,CAAC;IAOO,cAAc;QACpB,MAAM,WAAW,GAAa,EAAE,CAAA;QAChC,MAAM,WAAW,GAAG,IAAI,CAAC,OAAO,CAAC,CAAC,CAAC,CAAA;QACnC,MAAM,OAAO,GAAG,MAAM,CAAC,mBAAmB,CAAC,WAAW,CAAC,CAAA;QACvD,OAAO,CAAC,IAAI,CACV,GAAG,MAAM,CAAC,mBAAmB,CAAC,MAAM,CAAC,cAAc,CAAC,WAAW,CAAC,CAAC,CAClE,CAAA;QACD,KAAK,MAAM,IAAI,IAAI,OAAO,EAAE;YAC1B,IACE,OAAO,WAAW,CAAC,IAAI,CAAC,KAAK,UAAU;gBACvC,IAAI,KAAK,aAAa;gBACtB,IAAI,KAAK,IAAI,EACb;gBACA,WAAW,CAAC,IAAI,CAAC,IAAI,CAAC,CAAA;aACvB;SACF;QACD,OAAO,WAAW,CAAA;IACpB,CAAC;CACF;AAlED,kCAkEC"}
|
||||
+7
@@ -0,0 +1,7 @@
|
||||
export default class ConnectionManager {
|
||||
private promisesAwaitingConnection;
|
||||
resolveAllAwaiting(): void;
|
||||
rejectAllAwaiting(error: Error): void;
|
||||
awaitConnection(): Promise<void>;
|
||||
}
|
||||
//# sourceMappingURL=ConnectionManager.d.ts.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"ConnectionManager.d.ts","sourceRoot":"","sources":["../../../src/client/ConnectionManager.ts"],"names":[],"mappings":"AAKA,MAAM,CAAC,OAAO,OAAO,iBAAiB;IACpC,OAAO,CAAC,0BAA0B,CAG3B;IAKA,kBAAkB,IAAI,IAAI;IAU1B,iBAAiB,CAAC,KAAK,EAAE,KAAK,GAAG,IAAI;IAU/B,eAAe,IAAI,OAAO,CAAC,IAAI,CAAC;CAK9C"}
|
||||
+33
@@ -0,0 +1,33 @@
|
||||
"use strict";
|
||||
var __awaiter = (this && this.__awaiter) || function (thisArg, _arguments, P, generator) {
|
||||
function adopt(value) { return value instanceof P ? value : new P(function (resolve) { resolve(value); }); }
|
||||
return new (P || (P = Promise))(function (resolve, reject) {
|
||||
function fulfilled(value) { try { step(generator.next(value)); } catch (e) { reject(e); } }
|
||||
function rejected(value) { try { step(generator["throw"](value)); } catch (e) { reject(e); } }
|
||||
function step(result) { result.done ? resolve(result.value) : adopt(result.value).then(fulfilled, rejected); }
|
||||
step((generator = generator.apply(thisArg, _arguments || [])).next());
|
||||
});
|
||||
};
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
class ConnectionManager {
|
||||
constructor() {
|
||||
this.promisesAwaitingConnection = [];
|
||||
}
|
||||
resolveAllAwaiting() {
|
||||
this.promisesAwaitingConnection.map(({ resolve }) => resolve());
|
||||
this.promisesAwaitingConnection = [];
|
||||
}
|
||||
rejectAllAwaiting(error) {
|
||||
this.promisesAwaitingConnection.map(({ reject }) => reject(error));
|
||||
this.promisesAwaitingConnection = [];
|
||||
}
|
||||
awaitConnection() {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
return new Promise((resolve, reject) => {
|
||||
this.promisesAwaitingConnection.push({ resolve, reject });
|
||||
});
|
||||
});
|
||||
}
|
||||
}
|
||||
exports.default = ConnectionManager;
|
||||
//# sourceMappingURL=ConnectionManager.js.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"ConnectionManager.js","sourceRoot":"","sources":["../../../src/client/ConnectionManager.ts"],"names":[],"mappings":";;;;;;;;;;;AAKA,MAAqB,iBAAiB;IAAtC;QACU,+BAA0B,GAG7B,EAAE,CAAA;IA8BT,CAAC;IAzBQ,kBAAkB;QACvB,IAAI,CAAC,0BAA0B,CAAC,GAAG,CAAC,CAAC,EAAE,OAAO,EAAE,EAAE,EAAE,CAAC,OAAO,EAAE,CAAC,CAAA;QAC/D,IAAI,CAAC,0BAA0B,GAAG,EAAE,CAAA;IACtC,CAAC;IAOM,iBAAiB,CAAC,KAAY;QACnC,IAAI,CAAC,0BAA0B,CAAC,GAAG,CAAC,CAAC,EAAE,MAAM,EAAE,EAAE,EAAE,CAAC,MAAM,CAAC,KAAK,CAAC,CAAC,CAAA;QAClE,IAAI,CAAC,0BAA0B,GAAG,EAAE,CAAA;IACtC,CAAC;IAOY,eAAe;;YAC1B,OAAO,IAAI,OAAO,CAAC,CAAC,OAAO,EAAE,MAAM,EAAE,EAAE;gBACrC,IAAI,CAAC,0BAA0B,CAAC,IAAI,CAAC,EAAE,OAAO,EAAE,MAAM,EAAE,CAAC,CAAA;YAC3D,CAAC,CAAC,CAAA;QACJ,CAAC;KAAA;CACF;AAlCD,oCAkCC"}
|
||||
+16
@@ -0,0 +1,16 @@
|
||||
interface ExponentialBackoffOptions {
|
||||
min?: number;
|
||||
max?: number;
|
||||
}
|
||||
export default class ExponentialBackoff {
|
||||
private readonly ms;
|
||||
private readonly max;
|
||||
private readonly factor;
|
||||
private numAttempts;
|
||||
constructor(opts?: ExponentialBackoffOptions);
|
||||
get attempts(): number;
|
||||
duration(): number;
|
||||
reset(): void;
|
||||
}
|
||||
export {};
|
||||
//# sourceMappingURL=ExponentialBackoff.d.ts.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"ExponentialBackoff.d.ts","sourceRoot":"","sources":["../../../src/client/ExponentialBackoff.ts"],"names":[],"mappings":"AAcA,UAAU,yBAAyB;IAEjC,GAAG,CAAC,EAAE,MAAM,CAAA;IAEZ,GAAG,CAAC,EAAE,MAAM,CAAA;CACb;AASD,MAAM,CAAC,OAAO,OAAO,kBAAkB;IACrC,OAAO,CAAC,QAAQ,CAAC,EAAE,CAAQ;IAC3B,OAAO,CAAC,QAAQ,CAAC,GAAG,CAAQ;IAC5B,OAAO,CAAC,QAAQ,CAAC,MAAM,CAAY;IACnC,OAAO,CAAC,WAAW,CAAI;gBAOJ,IAAI,GAAE,yBAA8B;IAUvD,IAAW,QAAQ,IAAI,MAAM,CAE5B;IAOM,QAAQ,IAAI,MAAM;IASlB,KAAK,IAAI,IAAI;CAGrB"}
|
||||
+26
@@ -0,0 +1,26 @@
|
||||
"use strict";
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
const DEFAULT_MIN = 100;
|
||||
const DEFAULT_MAX = 1000;
|
||||
class ExponentialBackoff {
|
||||
constructor(opts = {}) {
|
||||
var _a, _b;
|
||||
this.factor = 2;
|
||||
this.numAttempts = 0;
|
||||
this.ms = (_a = opts.min) !== null && _a !== void 0 ? _a : DEFAULT_MIN;
|
||||
this.max = (_b = opts.max) !== null && _b !== void 0 ? _b : DEFAULT_MAX;
|
||||
}
|
||||
get attempts() {
|
||||
return this.numAttempts;
|
||||
}
|
||||
duration() {
|
||||
const ms = this.ms * Math.pow(this.factor, this.numAttempts);
|
||||
this.numAttempts += 1;
|
||||
return Math.floor(Math.min(ms, this.max));
|
||||
}
|
||||
reset() {
|
||||
this.numAttempts = 0;
|
||||
}
|
||||
}
|
||||
exports.default = ExponentialBackoff;
|
||||
//# sourceMappingURL=ExponentialBackoff.js.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"ExponentialBackoff.js","sourceRoot":"","sources":["../../../src/client/ExponentialBackoff.ts"],"names":[],"mappings":";;AAqBA,MAAM,WAAW,GAAG,GAAG,CAAA;AACvB,MAAM,WAAW,GAAG,IAAI,CAAA;AAMxB,MAAqB,kBAAkB;IAWrC,YAAmB,OAAkC,EAAE;;QARtC,WAAM,GAAW,CAAC,CAAA;QAC3B,gBAAW,GAAG,CAAC,CAAA;QAQrB,IAAI,CAAC,EAAE,GAAG,MAAA,IAAI,CAAC,GAAG,mCAAI,WAAW,CAAA;QACjC,IAAI,CAAC,GAAG,GAAG,MAAA,IAAI,CAAC,GAAG,mCAAI,WAAW,CAAA;IACpC,CAAC;IAOD,IAAW,QAAQ;QACjB,OAAO,IAAI,CAAC,WAAW,CAAA;IACzB,CAAC;IAOM,QAAQ;QACb,MAAM,EAAE,GAAG,IAAI,CAAC,EAAE,GAAG,SAAA,IAAI,CAAC,MAAM,EAAI,IAAI,CAAC,WAAW,CAAA,CAAA;QACpD,IAAI,CAAC,WAAW,IAAI,CAAC,CAAA;QACrB,OAAO,IAAI,CAAC,KAAK,CAAC,IAAI,CAAC,GAAG,CAAC,EAAE,EAAE,IAAI,CAAC,GAAG,CAAC,CAAC,CAAA;IAC3C,CAAC;IAKM,KAAK;QACV,IAAI,CAAC,WAAW,GAAG,CAAC,CAAA;IACtB,CAAC;CACF;AA1CD,qCA0CC"}
|
||||
+13
@@ -0,0 +1,13 @@
|
||||
import { Response } from '../models/methods';
|
||||
import { BaseRequest, ErrorResponse } from '../models/methods/baseMethod';
|
||||
export default class RequestManager {
|
||||
private nextId;
|
||||
private readonly promisesAwaitingResponse;
|
||||
resolve(id: string | number, response: Response): void;
|
||||
reject(id: string | number, error: Error): void;
|
||||
rejectAll(error: Error): void;
|
||||
createRequest<T extends BaseRequest>(request: T, timeout: number): [string | number, string, Promise<Response>];
|
||||
handleResponse(response: Partial<Response | ErrorResponse>): void;
|
||||
private deletePromise;
|
||||
}
|
||||
//# sourceMappingURL=RequestManager.d.ts.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"RequestManager.d.ts","sourceRoot":"","sources":["../../../src/client/RequestManager.ts"],"names":[],"mappings":"AAMA,OAAO,EAAE,QAAQ,EAAE,MAAM,mBAAmB,CAAA;AAC5C,OAAO,EAAE,WAAW,EAAE,aAAa,EAAE,MAAM,8BAA8B,CAAA;AAQzE,MAAM,CAAC,OAAO,OAAO,cAAc;IACjC,OAAO,CAAC,MAAM,CAAI;IAClB,OAAO,CAAC,QAAQ,CAAC,wBAAwB,CAOtC;IASI,OAAO,CAAC,EAAE,EAAE,MAAM,GAAG,MAAM,EAAE,QAAQ,EAAE,QAAQ,GAAG,IAAI;IAoBtD,MAAM,CAAC,EAAE,EAAE,MAAM,GAAG,MAAM,EAAE,KAAK,EAAE,KAAK,GAAG,IAAI;IAmB/C,SAAS,CAAC,KAAK,EAAE,KAAK,GAAG,IAAI;IAiB7B,aAAa,CAAC,CAAC,SAAS,WAAW,EACxC,OAAO,EAAE,CAAC,EACV,OAAO,EAAE,MAAM,GACd,CAAC,MAAM,GAAG,MAAM,EAAE,MAAM,EAAE,OAAO,CAAC,QAAQ,CAAC,CAAC;IAuDxC,cAAc,CAAC,QAAQ,EAAE,OAAO,CAAC,QAAQ,GAAG,aAAa,CAAC,GAAG,IAAI;IA2CxE,OAAO,CAAC,aAAa;CAGtB"}
|
||||
+97
@@ -0,0 +1,97 @@
|
||||
"use strict";
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
const errors_1 = require("../errors");
|
||||
class RequestManager {
|
||||
constructor() {
|
||||
this.nextId = 0;
|
||||
this.promisesAwaitingResponse = new Map();
|
||||
}
|
||||
resolve(id, response) {
|
||||
const promise = this.promisesAwaitingResponse.get(id);
|
||||
if (promise == null) {
|
||||
throw new errors_1.XrplError(`No existing promise with id ${id}`, {
|
||||
type: 'resolve',
|
||||
response,
|
||||
});
|
||||
}
|
||||
clearTimeout(promise.timer);
|
||||
promise.resolve(response);
|
||||
this.deletePromise(id);
|
||||
}
|
||||
reject(id, error) {
|
||||
const promise = this.promisesAwaitingResponse.get(id);
|
||||
if (promise == null) {
|
||||
throw new errors_1.XrplError(`No existing promise with id ${id}`, {
|
||||
type: 'reject',
|
||||
error,
|
||||
});
|
||||
}
|
||||
clearTimeout(promise.timer);
|
||||
promise.reject(error);
|
||||
this.deletePromise(id);
|
||||
}
|
||||
rejectAll(error) {
|
||||
this.promisesAwaitingResponse.forEach((_promise, id, _map) => {
|
||||
this.reject(id, error);
|
||||
this.deletePromise(id);
|
||||
});
|
||||
}
|
||||
createRequest(request, timeout) {
|
||||
let newId;
|
||||
if (request.id == null) {
|
||||
newId = this.nextId;
|
||||
this.nextId += 1;
|
||||
}
|
||||
else {
|
||||
newId = request.id;
|
||||
}
|
||||
const newRequest = JSON.stringify(Object.assign(Object.assign({}, request), { id: newId }));
|
||||
const timer = setTimeout(() => {
|
||||
this.reject(newId, new errors_1.TimeoutError(`Timeout for request: ${JSON.stringify(request)} with id ${newId}`, request));
|
||||
}, timeout);
|
||||
if (timer.unref) {
|
||||
;
|
||||
timer.unref();
|
||||
}
|
||||
if (this.promisesAwaitingResponse.has(newId)) {
|
||||
clearTimeout(timer);
|
||||
throw new errors_1.XrplError(`Response with id '${newId}' is already pending`, request);
|
||||
}
|
||||
const newPromise = new Promise((resolve, reject) => {
|
||||
this.promisesAwaitingResponse.set(newId, { resolve, reject, timer });
|
||||
});
|
||||
return [newId, newRequest, newPromise];
|
||||
}
|
||||
handleResponse(response) {
|
||||
var _a, _b;
|
||||
if (response.id == null ||
|
||||
!(typeof response.id === 'string' || typeof response.id === 'number')) {
|
||||
throw new errors_1.ResponseFormatError('valid id not found in response', response);
|
||||
}
|
||||
if (!this.promisesAwaitingResponse.has(response.id)) {
|
||||
return;
|
||||
}
|
||||
if (response.status == null) {
|
||||
const error = new errors_1.ResponseFormatError('Response has no status');
|
||||
this.reject(response.id, error);
|
||||
}
|
||||
if (response.status === 'error') {
|
||||
const errorResponse = response;
|
||||
const error = new errors_1.RippledError((_a = errorResponse.error_message) !== null && _a !== void 0 ? _a : errorResponse.error, errorResponse);
|
||||
this.reject(response.id, error);
|
||||
return;
|
||||
}
|
||||
if (response.status !== 'success') {
|
||||
const error = new errors_1.ResponseFormatError(`unrecognized response.status: ${(_b = response.status) !== null && _b !== void 0 ? _b : ''}`, response);
|
||||
this.reject(response.id, error);
|
||||
return;
|
||||
}
|
||||
delete response.status;
|
||||
this.resolve(response.id, response);
|
||||
}
|
||||
deletePromise(id) {
|
||||
this.promisesAwaitingResponse.delete(id);
|
||||
}
|
||||
}
|
||||
exports.default = RequestManager;
|
||||
//# sourceMappingURL=RequestManager.js.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"RequestManager.js","sourceRoot":"","sources":["../../../src/client/RequestManager.ts"],"names":[],"mappings":";;AAAA,sCAKkB;AAUlB,MAAqB,cAAc;IAAnC;QACU,WAAM,GAAG,CAAC,CAAA;QACD,6BAAwB,GAAG,IAAI,GAAG,EAOhD,CAAA;IAyKL,CAAC;IAhKQ,OAAO,CAAC,EAAmB,EAAE,QAAkB;QACpD,MAAM,OAAO,GAAG,IAAI,CAAC,wBAAwB,CAAC,GAAG,CAAC,EAAE,CAAC,CAAA;QACrD,IAAI,OAAO,IAAI,IAAI,EAAE;YACnB,MAAM,IAAI,kBAAS,CAAC,+BAA+B,EAAE,EAAE,EAAE;gBACvD,IAAI,EAAE,SAAS;gBACf,QAAQ;aACT,CAAC,CAAA;SACH;QACD,YAAY,CAAC,OAAO,CAAC,KAAK,CAAC,CAAA;QAC3B,OAAO,CAAC,OAAO,CAAC,QAAQ,CAAC,CAAA;QACzB,IAAI,CAAC,aAAa,CAAC,EAAE,CAAC,CAAA;IACxB,CAAC;IASM,MAAM,CAAC,EAAmB,EAAE,KAAY;QAC7C,MAAM,OAAO,GAAG,IAAI,CAAC,wBAAwB,CAAC,GAAG,CAAC,EAAE,CAAC,CAAA;QACrD,IAAI,OAAO,IAAI,IAAI,EAAE;YACnB,MAAM,IAAI,kBAAS,CAAC,+BAA+B,EAAE,EAAE,EAAE;gBACvD,IAAI,EAAE,QAAQ;gBACd,KAAK;aACN,CAAC,CAAA;SACH;QACD,YAAY,CAAC,OAAO,CAAC,KAAK,CAAC,CAAA;QAE3B,OAAO,CAAC,MAAM,CAAC,KAAK,CAAC,CAAA;QACrB,IAAI,CAAC,aAAa,CAAC,EAAE,CAAC,CAAA;IACxB,CAAC;IAOM,SAAS,CAAC,KAAY;QAC3B,IAAI,CAAC,wBAAwB,CAAC,OAAO,CAAC,CAAC,QAAQ,EAAE,EAAE,EAAE,IAAI,EAAE,EAAE;YAC3D,IAAI,CAAC,MAAM,CAAC,EAAE,EAAE,KAAK,CAAC,CAAA;YACtB,IAAI,CAAC,aAAa,CAAC,EAAE,CAAC,CAAA;QACxB,CAAC,CAAC,CAAA;IACJ,CAAC;IAYM,aAAa,CAClB,OAAU,EACV,OAAe;QAEf,IAAI,KAAsB,CAAA;QAC1B,IAAI,OAAO,CAAC,EAAE,IAAI,IAAI,EAAE;YACtB,KAAK,GAAG,IAAI,CAAC,MAAM,CAAA;YACnB,IAAI,CAAC,MAAM,IAAI,CAAC,CAAA;SACjB;aAAM;YACL,KAAK,GAAG,OAAO,CAAC,EAAE,CAAA;SACnB;QACD,MAAM,UAAU,GAAG,IAAI,CAAC,SAAS,iCAAM,OAAO,KAAE,EAAE,EAAE,KAAK,IAAG,CAAA;QAE5D,MAAM,KAAK,GAAkC,UAAU,CAAC,GAAG,EAAE;YAC3D,IAAI,CAAC,MAAM,CACT,KAAK,EACL,IAAI,qBAAY,CACd,wBAAwB,IAAI,CAAC,SAAS,CAAC,OAAO,CAAC,YAAY,KAAK,EAAE,EAClE,OAAO,CACR,CACF,CAAA;QACH,CAAC,EAAE,OAAO,CAAC,CAAA;QASX,IAAK,KAAwB,CAAC,KAAK,EAAE;YAGnC,CAAC;YAAC,KAAwB,CAAC,KAAK,EAAE,CAAA;SACnC;QACD,IAAI,IAAI,CAAC,wBAAwB,CAAC,GAAG,CAAC,KAAK,CAAC,EAAE;YAC5C,YAAY,CAAC,KAAK,CAAC,CAAA;YACnB,MAAM,IAAI,kBAAS,CACjB,qBAAqB,KAAK,sBAAsB,EAChD,OAAO,CACR,CAAA;SACF;QACD,MAAM,UAAU,GAAG,IAAI,OAAO,CAC5B,CAAC,OAA0D,EAAE,MAAM,EAAE,EAAE;YACrE,IAAI,CAAC,wBAAwB,CAAC,GAAG,CAAC,KAAK,EAAE,EAAE,OAAO,EAAE,MAAM,EAAE,KAAK,EAAE,CAAC,CAAA;QACtE,CAAC,CACF,CAAA;QAED,OAAO,CAAC,KAAK,EAAE,UAAU,EAAE,UAAU,CAAC,CAAA;IACxC,CAAC;IASM,cAAc,CAAC,QAA2C;;QAC/D,IACE,QAAQ,CAAC,EAAE,IAAI,IAAI;YACnB,CAAC,CAAC,OAAO,QAAQ,CAAC,EAAE,KAAK,QAAQ,IAAI,OAAO,QAAQ,CAAC,EAAE,KAAK,QAAQ,CAAC,EACrE;YACA,MAAM,IAAI,4BAAmB,CAAC,gCAAgC,EAAE,QAAQ,CAAC,CAAA;SAC1E;QACD,IAAI,CAAC,IAAI,CAAC,wBAAwB,CAAC,GAAG,CAAC,QAAQ,CAAC,EAAE,CAAC,EAAE;YACnD,OAAM;SACP;QACD,IAAI,QAAQ,CAAC,MAAM,IAAI,IAAI,EAAE;YAC3B,MAAM,KAAK,GAAG,IAAI,4BAAmB,CAAC,wBAAwB,CAAC,CAAA;YAC/D,IAAI,CAAC,MAAM,CAAC,QAAQ,CAAC,EAAE,EAAE,KAAK,CAAC,CAAA;SAChC;QACD,IAAI,QAAQ,CAAC,MAAM,KAAK,OAAO,EAAE;YAE/B,MAAM,aAAa,GAAG,QAAkC,CAAA;YACxD,MAAM,KAAK,GAAG,IAAI,qBAAY,CAC5B,MAAA,aAAa,CAAC,aAAa,mCAAI,aAAa,CAAC,KAAK,EAClD,aAAa,CACd,CAAA;YACD,IAAI,CAAC,MAAM,CAAC,QAAQ,CAAC,EAAE,EAAE,KAAK,CAAC,CAAA;YAC/B,OAAM;SACP;QACD,IAAI,QAAQ,CAAC,MAAM,KAAK,SAAS,EAAE;YACjC,MAAM,KAAK,GAAG,IAAI,4BAAmB,CACnC,iCAAiC,MAAA,QAAQ,CAAC,MAAM,mCAAI,EAAE,EAAE,EACxD,QAAQ,CACT,CAAA;YACD,IAAI,CAAC,MAAM,CAAC,QAAQ,CAAC,EAAE,EAAE,KAAK,CAAC,CAAA;YAC/B,OAAM;SACP;QAED,OAAO,QAAQ,CAAC,MAAM,CAAA;QAEtB,IAAI,CAAC,OAAO,CAAC,QAAQ,CAAC,EAAE,EAAE,QAA+B,CAAC,CAAA;IAC5D,CAAC;IAOO,aAAa,CAAC,EAAmB;QACvC,IAAI,CAAC,wBAAwB,CAAC,MAAM,CAAC,EAAE,CAAC,CAAA;IAC1C,CAAC;CACF;AAlLD,iCAkLC"}
|
||||
+25
@@ -0,0 +1,25 @@
|
||||
/// <reference types="node" />
|
||||
/// <reference types="node" />
|
||||
import { EventEmitter } from 'events';
|
||||
interface WSWrapperOptions {
|
||||
perMessageDeflate: boolean;
|
||||
handshakeTimeout: number;
|
||||
protocolVersion: number;
|
||||
origin: string;
|
||||
maxPayload: number;
|
||||
followRedirects: boolean;
|
||||
maxRedirects: number;
|
||||
}
|
||||
export default class WSWrapper extends EventEmitter {
|
||||
static CONNECTING: number;
|
||||
static OPEN: number;
|
||||
static CLOSING: number;
|
||||
static CLOSED: number;
|
||||
private readonly ws;
|
||||
constructor(url: string, _protocols: string | string[] | WSWrapperOptions | undefined, _websocketOptions: WSWrapperOptions);
|
||||
get readyState(): number;
|
||||
close(code?: number, reason?: Buffer): void;
|
||||
send(message: string): void;
|
||||
}
|
||||
export {};
|
||||
//# sourceMappingURL=WSWrapper.d.ts.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"WSWrapper.d.ts","sourceRoot":"","sources":["../../../src/client/WSWrapper.ts"],"names":[],"mappings":";;AACA,OAAO,EAAE,YAAY,EAAE,MAAM,QAAQ,CAAA;AAcrC,UAAU,gBAAgB;IACxB,iBAAiB,EAAE,OAAO,CAAA;IAC1B,gBAAgB,EAAE,MAAM,CAAA;IACxB,eAAe,EAAE,MAAM,CAAA;IACvB,MAAM,EAAE,MAAM,CAAA;IACd,UAAU,EAAE,MAAM,CAAA;IAClB,eAAe,EAAE,OAAO,CAAA;IACxB,YAAY,EAAE,MAAM,CAAA;CACrB;AAMD,MAAM,CAAC,OAAO,OAAO,SAAU,SAAQ,YAAY;IACjD,OAAc,UAAU,SAAI;IAC5B,OAAc,IAAI,SAAI;IACtB,OAAc,OAAO,SAAI;IAEzB,OAAc,MAAM,SAAI;IACxB,OAAO,CAAC,QAAQ,CAAC,EAAE,CAAW;gBAU5B,GAAG,EAAE,MAAM,EACX,UAAU,EAAE,MAAM,GAAG,MAAM,EAAE,GAAG,gBAAgB,GAAG,SAAS,EAC5D,iBAAiB,EAAE,gBAAgB;IAkCrC,IAAW,UAAU,IAAI,MAAM,CAE9B;IAQM,KAAK,CAAC,IAAI,CAAC,EAAE,MAAM,EAAE,MAAM,CAAC,EAAE,MAAM,GAAG,IAAI;IAW3C,IAAI,CAAC,OAAO,EAAE,MAAM,GAAG,IAAI;CAGnC"}
|
||||
+44
@@ -0,0 +1,44 @@
|
||||
"use strict";
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
const events_1 = require("events");
|
||||
class WSWrapper extends events_1.EventEmitter {
|
||||
constructor(url, _protocols, _websocketOptions) {
|
||||
super();
|
||||
this.setMaxListeners(Infinity);
|
||||
this.ws = new WebSocket(url);
|
||||
this.ws.onclose = (closeEvent) => {
|
||||
let reason;
|
||||
if (closeEvent.reason) {
|
||||
const enc = new TextEncoder();
|
||||
reason = enc.encode(closeEvent.reason);
|
||||
}
|
||||
this.emit('close', closeEvent.code, reason);
|
||||
};
|
||||
this.ws.onopen = () => {
|
||||
this.emit('open');
|
||||
};
|
||||
this.ws.onerror = (error) => {
|
||||
this.emit('error', error);
|
||||
};
|
||||
this.ws.onmessage = (message) => {
|
||||
this.emit('message', message.data);
|
||||
};
|
||||
}
|
||||
get readyState() {
|
||||
return this.ws.readyState;
|
||||
}
|
||||
close(code, reason) {
|
||||
if (this.readyState === 1) {
|
||||
this.ws.close(code, reason);
|
||||
}
|
||||
}
|
||||
send(message) {
|
||||
this.ws.send(message);
|
||||
}
|
||||
}
|
||||
exports.default = WSWrapper;
|
||||
WSWrapper.CONNECTING = 0;
|
||||
WSWrapper.OPEN = 1;
|
||||
WSWrapper.CLOSING = 2;
|
||||
WSWrapper.CLOSED = 3;
|
||||
//# sourceMappingURL=WSWrapper.js.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"WSWrapper.js","sourceRoot":"","sources":["../../../src/client/WSWrapper.ts"],"names":[],"mappings":";;AACA,mCAAqC;AA4BrC,MAAqB,SAAU,SAAQ,qBAAY;IAejD,YACE,GAAW,EACX,UAA4D,EAC5D,iBAAmC;QAEnC,KAAK,EAAE,CAAA;QACP,IAAI,CAAC,eAAe,CAAC,QAAQ,CAAC,CAAA;QAE9B,IAAI,CAAC,EAAE,GAAG,IAAI,SAAS,CAAC,GAAG,CAAC,CAAA;QAE5B,IAAI,CAAC,EAAE,CAAC,OAAO,GAAG,CAAC,UAAsB,EAAQ,EAAE;YACjD,IAAI,MAA8B,CAAA;YAClC,IAAI,UAAU,CAAC,MAAM,EAAE;gBACrB,MAAM,GAAG,GAAG,IAAI,WAAW,EAAE,CAAA;gBAC7B,MAAM,GAAG,GAAG,CAAC,MAAM,CAAC,UAAU,CAAC,MAAM,CAAC,CAAA;aACvC;YACD,IAAI,CAAC,IAAI,CAAC,OAAO,EAAE,UAAU,CAAC,IAAI,EAAE,MAAM,CAAC,CAAA;QAC7C,CAAC,CAAA;QAED,IAAI,CAAC,EAAE,CAAC,MAAM,GAAG,GAAS,EAAE;YAC1B,IAAI,CAAC,IAAI,CAAC,MAAM,CAAC,CAAA;QACnB,CAAC,CAAA;QAED,IAAI,CAAC,EAAE,CAAC,OAAO,GAAG,CAAC,KAAK,EAAQ,EAAE;YAChC,IAAI,CAAC,IAAI,CAAC,OAAO,EAAE,KAAK,CAAC,CAAA;QAC3B,CAAC,CAAA;QAED,IAAI,CAAC,EAAE,CAAC,SAAS,GAAG,CAAC,OAAqB,EAAQ,EAAE;YAClD,IAAI,CAAC,IAAI,CAAC,SAAS,EAAE,OAAO,CAAC,IAAI,CAAC,CAAA;QACpC,CAAC,CAAA;IACH,CAAC;IAOD,IAAW,UAAU;QACnB,OAAO,IAAI,CAAC,EAAE,CAAC,UAAU,CAAA;IAC3B,CAAC;IAQM,KAAK,CAAC,IAAa,EAAE,MAAe;QACzC,IAAI,IAAI,CAAC,UAAU,KAAK,CAAC,EAAE;YACzB,IAAI,CAAC,EAAE,CAAC,KAAK,CAAC,IAAI,EAAE,MAAM,CAAC,CAAA;SAC5B;IACH,CAAC;IAOM,IAAI,CAAC,OAAe;QACzB,IAAI,CAAC,EAAE,CAAC,IAAI,CAAC,OAAO,CAAC,CAAA;IACvB,CAAC;;AA3EH,4BA4EC;AA3Ee,oBAAU,GAAG,CAAC,CAAA;AACd,cAAI,GAAG,CAAC,CAAA;AACR,iBAAO,GAAG,CAAC,CAAA;AAEX,gBAAM,GAAG,CAAC,CAAA"}
|
||||
+49
@@ -0,0 +1,49 @@
|
||||
/// <reference types="node" />
|
||||
import { EventEmitter } from 'events';
|
||||
import { BaseRequest } from '../models/methods/baseMethod';
|
||||
interface ConnectionOptions {
|
||||
trace?: boolean | ((id: string, message: string) => void);
|
||||
proxy?: string;
|
||||
proxyAuthorization?: string;
|
||||
authorization?: string;
|
||||
trustedCertificates?: string[];
|
||||
key?: string;
|
||||
passphrase?: string;
|
||||
certificate?: string;
|
||||
timeout: number;
|
||||
connectionTimeout: number;
|
||||
headers?: {
|
||||
[key: string]: string;
|
||||
};
|
||||
}
|
||||
export type ConnectionUserOptions = Partial<ConnectionOptions>;
|
||||
export declare const INTENTIONAL_DISCONNECT_CODE = 4000;
|
||||
export declare class Connection extends EventEmitter {
|
||||
private readonly url;
|
||||
private ws;
|
||||
private reconnectTimeoutID;
|
||||
private heartbeatIntervalID;
|
||||
private readonly retryConnectionBackoff;
|
||||
private readonly config;
|
||||
private readonly requestManager;
|
||||
private readonly connectionManager;
|
||||
constructor(url?: string, options?: ConnectionUserOptions);
|
||||
private get state();
|
||||
private get shouldBeConnected();
|
||||
isConnected(): boolean;
|
||||
connect(): Promise<void>;
|
||||
disconnect(): Promise<number | undefined>;
|
||||
reconnect(): Promise<void>;
|
||||
request<T extends BaseRequest>(request: T, timeout?: number): Promise<unknown>;
|
||||
getUrl(): string;
|
||||
readonly trace: (id: string, message: string) => void;
|
||||
private onMessage;
|
||||
private onceOpen;
|
||||
private intentionalDisconnect;
|
||||
private clearHeartbeatInterval;
|
||||
private startHeartbeatInterval;
|
||||
private heartbeat;
|
||||
private onConnectionFailed;
|
||||
}
|
||||
export {};
|
||||
//# sourceMappingURL=connection.d.ts.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"connection.d.ts","sourceRoot":"","sources":["../../../src/client/connection.ts"],"names":[],"mappings":";AACA,OAAO,EAAE,YAAY,EAAE,MAAM,QAAQ,CAAA;AAYrC,OAAO,EAAE,WAAW,EAAE,MAAM,8BAA8B,CAAA;AAa1D,UAAU,iBAAiB;IACzB,KAAK,CAAC,EAAE,OAAO,GAAG,CAAC,CAAC,EAAE,EAAE,MAAM,EAAE,OAAO,EAAE,MAAM,KAAK,IAAI,CAAC,CAAA;IACzD,KAAK,CAAC,EAAE,MAAM,CAAA;IACd,kBAAkB,CAAC,EAAE,MAAM,CAAA;IAC3B,aAAa,CAAC,EAAE,MAAM,CAAA;IACtB,mBAAmB,CAAC,EAAE,MAAM,EAAE,CAAA;IAC9B,GAAG,CAAC,EAAE,MAAM,CAAA;IACZ,UAAU,CAAC,EAAE,MAAM,CAAA;IACnB,WAAW,CAAC,EAAE,MAAM,CAAA;IAEpB,OAAO,EAAE,MAAM,CAAA;IACf,iBAAiB,EAAE,MAAM,CAAA;IACzB,OAAO,CAAC,EAAE;QAAE,CAAC,GAAG,EAAE,MAAM,GAAG,MAAM,CAAA;KAAE,CAAA;CACpC;AAOD,MAAM,MAAM,qBAAqB,GAAG,OAAO,CAAC,iBAAiB,CAAC,CAAA;AAO9D,eAAO,MAAM,2BAA2B,OAAO,CAAA;AAyH/C,qBAAa,UAAW,SAAQ,YAAY;IAC1C,OAAO,CAAC,QAAQ,CAAC,GAAG,CAAoB;IACxC,OAAO,CAAC,EAAE,CAAyB;IAEnC,OAAO,CAAC,kBAAkB,CAA6C;IAEvE,OAAO,CAAC,mBAAmB,CAA6C;IACxE,OAAO,CAAC,QAAQ,CAAC,sBAAsB,CAGrC;IAEF,OAAO,CAAC,QAAQ,CAAC,MAAM,CAAmB;IAC1C,OAAO,CAAC,QAAQ,CAAC,cAAc,CAAuB;IACtD,OAAO,CAAC,QAAQ,CAAC,iBAAiB,CAA0B;gBAQzC,GAAG,CAAC,EAAE,MAAM,EAAE,OAAO,GAAE,qBAA0B;IAsBpE,OAAO,KAAK,KAAK,GAEhB;IAOD,OAAO,KAAK,iBAAiB,GAE5B;IAOM,WAAW,IAAI,OAAO;IAWhB,OAAO,IAAI,OAAO,CAAC,IAAI,CAAC;IA0DxB,UAAU,IAAI,OAAO,CAAC,MAAM,GAAG,SAAS,CAAC;IAkCzC,SAAS,IAAI,OAAO,CAAC,IAAI,CAAC;IAmB1B,OAAO,CAAC,CAAC,SAAS,WAAW,EACxC,OAAO,EAAE,CAAC,EACV,OAAO,CAAC,EAAE,MAAM,GACf,OAAO,CAAC,OAAO,CAAC;IAqBZ,MAAM,IAAI,MAAM;IAKvB,SAAgB,KAAK,EAAE,CAAC,EAAE,EAAE,MAAM,EAAE,OAAO,EAAE,MAAM,KAAK,IAAI,CAAW;IAOvE,OAAO,CAAC,SAAS;YA2CH,QAAQ;IA0EtB,OAAO,CAAC,qBAAqB;IAkB7B,OAAO,CAAC,sBAAsB;IAS9B,OAAO,CAAC,sBAAsB;YAahB,SAAS;IAavB,OAAO,CAAC,kBAAkB;CA4B3B"}
|
||||
+342
@@ -0,0 +1,342 @@
|
||||
"use strict";
|
||||
var __awaiter = (this && this.__awaiter) || function (thisArg, _arguments, P, generator) {
|
||||
function adopt(value) { return value instanceof P ? value : new P(function (resolve) { resolve(value); }); }
|
||||
return new (P || (P = Promise))(function (resolve, reject) {
|
||||
function fulfilled(value) { try { step(generator.next(value)); } catch (e) { reject(e); } }
|
||||
function rejected(value) { try { step(generator["throw"](value)); } catch (e) { reject(e); } }
|
||||
function step(result) { result.done ? resolve(result.value) : adopt(result.value).then(fulfilled, rejected); }
|
||||
step((generator = generator.apply(thisArg, _arguments || [])).next());
|
||||
});
|
||||
};
|
||||
var __importDefault = (this && this.__importDefault) || function (mod) {
|
||||
return (mod && mod.__esModule) ? mod : { "default": mod };
|
||||
};
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
exports.Connection = exports.INTENTIONAL_DISCONNECT_CODE = void 0;
|
||||
const events_1 = require("events");
|
||||
const omitBy_1 = __importDefault(require("lodash/omitBy"));
|
||||
const ws_1 = __importDefault(require("ws"));
|
||||
const errors_1 = require("../errors");
|
||||
const ConnectionManager_1 = __importDefault(require("./ConnectionManager"));
|
||||
const ExponentialBackoff_1 = __importDefault(require("./ExponentialBackoff"));
|
||||
const RequestManager_1 = __importDefault(require("./RequestManager"));
|
||||
const SECONDS_PER_MINUTE = 60;
|
||||
const TIMEOUT = 20;
|
||||
const CONNECTION_TIMEOUT = 5;
|
||||
exports.INTENTIONAL_DISCONNECT_CODE = 4000;
|
||||
function getAgent(url, config) {
|
||||
if (config.proxy == null) {
|
||||
return undefined;
|
||||
}
|
||||
const parsedURL = new URL(url);
|
||||
const parsedProxyURL = new URL(config.proxy);
|
||||
const proxyOptions = (0, omitBy_1.default)({
|
||||
secureEndpoint: parsedURL.protocol === 'wss:',
|
||||
secureProxy: parsedProxyURL.protocol === 'https:',
|
||||
auth: config.proxyAuthorization,
|
||||
ca: config.trustedCertificates,
|
||||
key: config.key,
|
||||
passphrase: config.passphrase,
|
||||
cert: config.certificate,
|
||||
href: parsedProxyURL.href,
|
||||
origin: parsedProxyURL.origin,
|
||||
protocol: parsedProxyURL.protocol,
|
||||
username: parsedProxyURL.username,
|
||||
password: parsedProxyURL.password,
|
||||
host: parsedProxyURL.host,
|
||||
hostname: parsedProxyURL.hostname,
|
||||
port: parsedProxyURL.port,
|
||||
pathname: parsedProxyURL.pathname,
|
||||
search: parsedProxyURL.search,
|
||||
hash: parsedProxyURL.hash,
|
||||
}, (value) => value == null);
|
||||
let HttpsProxyAgent;
|
||||
try {
|
||||
HttpsProxyAgent = require('https-proxy-agent');
|
||||
}
|
||||
catch (_error) {
|
||||
throw new Error('"proxy" option is not supported in the browser');
|
||||
}
|
||||
return new HttpsProxyAgent(proxyOptions);
|
||||
}
|
||||
function createWebSocket(url, config) {
|
||||
const options = {};
|
||||
options.agent = getAgent(url, config);
|
||||
if (config.headers) {
|
||||
options.headers = config.headers;
|
||||
}
|
||||
if (config.authorization != null) {
|
||||
const base64 = Buffer.from(config.authorization).toString('base64');
|
||||
options.headers = Object.assign(Object.assign({}, options.headers), { Authorization: `Basic ${base64}` });
|
||||
}
|
||||
const optionsOverrides = (0, omitBy_1.default)({
|
||||
ca: config.trustedCertificates,
|
||||
key: config.key,
|
||||
passphrase: config.passphrase,
|
||||
cert: config.certificate,
|
||||
}, (value) => value == null);
|
||||
const websocketOptions = Object.assign(Object.assign({}, options), optionsOverrides);
|
||||
const websocket = new ws_1.default(url, websocketOptions);
|
||||
if (typeof websocket.setMaxListeners === 'function') {
|
||||
websocket.setMaxListeners(Infinity);
|
||||
}
|
||||
return websocket;
|
||||
}
|
||||
function websocketSendAsync(ws, message) {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
return new Promise((resolve, reject) => {
|
||||
ws.send(message, (error) => {
|
||||
if (error) {
|
||||
reject(new errors_1.DisconnectedError(error.message, error));
|
||||
}
|
||||
else {
|
||||
resolve();
|
||||
}
|
||||
});
|
||||
});
|
||||
});
|
||||
}
|
||||
class Connection extends events_1.EventEmitter {
|
||||
constructor(url, options = {}) {
|
||||
super();
|
||||
this.ws = null;
|
||||
this.reconnectTimeoutID = null;
|
||||
this.heartbeatIntervalID = null;
|
||||
this.retryConnectionBackoff = new ExponentialBackoff_1.default({
|
||||
min: 100,
|
||||
max: SECONDS_PER_MINUTE * 1000,
|
||||
});
|
||||
this.requestManager = new RequestManager_1.default();
|
||||
this.connectionManager = new ConnectionManager_1.default();
|
||||
this.trace = () => { };
|
||||
this.setMaxListeners(Infinity);
|
||||
this.url = url;
|
||||
this.config = Object.assign({ timeout: TIMEOUT * 1000, connectionTimeout: CONNECTION_TIMEOUT * 1000 }, options);
|
||||
if (typeof options.trace === 'function') {
|
||||
this.trace = options.trace;
|
||||
}
|
||||
else if (options.trace) {
|
||||
this.trace = console.log;
|
||||
}
|
||||
}
|
||||
get state() {
|
||||
return this.ws ? this.ws.readyState : ws_1.default.CLOSED;
|
||||
}
|
||||
get shouldBeConnected() {
|
||||
return this.ws !== null;
|
||||
}
|
||||
isConnected() {
|
||||
return this.state === ws_1.default.OPEN;
|
||||
}
|
||||
connect() {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
if (this.isConnected()) {
|
||||
return Promise.resolve();
|
||||
}
|
||||
if (this.state === ws_1.default.CONNECTING) {
|
||||
return this.connectionManager.awaitConnection();
|
||||
}
|
||||
if (!this.url) {
|
||||
return Promise.reject(new errors_1.ConnectionError('Cannot connect because no server was specified'));
|
||||
}
|
||||
if (this.ws != null) {
|
||||
return Promise.reject(new errors_1.XrplError('Websocket connection never cleaned up.', {
|
||||
state: this.state,
|
||||
}));
|
||||
}
|
||||
const connectionTimeoutID = setTimeout(() => {
|
||||
this.onConnectionFailed(new errors_1.ConnectionError(`Error: connect() timed out after ${this.config.connectionTimeout} ms. If your internet connection is working, the ` +
|
||||
`rippled server may be blocked or inaccessible. You can also try setting the 'connectionTimeout' option in the Client constructor.`));
|
||||
}, this.config.connectionTimeout);
|
||||
this.ws = createWebSocket(this.url, this.config);
|
||||
if (this.ws == null) {
|
||||
throw new errors_1.XrplError('Connect: created null websocket');
|
||||
}
|
||||
this.ws.on('error', (error) => this.onConnectionFailed(error));
|
||||
this.ws.on('error', () => clearTimeout(connectionTimeoutID));
|
||||
this.ws.on('close', (reason) => this.onConnectionFailed(reason));
|
||||
this.ws.on('close', () => clearTimeout(connectionTimeoutID));
|
||||
this.ws.once('open', () => {
|
||||
void this.onceOpen(connectionTimeoutID);
|
||||
});
|
||||
return this.connectionManager.awaitConnection();
|
||||
});
|
||||
}
|
||||
disconnect() {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
this.clearHeartbeatInterval();
|
||||
if (this.reconnectTimeoutID !== null) {
|
||||
clearTimeout(this.reconnectTimeoutID);
|
||||
this.reconnectTimeoutID = null;
|
||||
}
|
||||
if (this.state === ws_1.default.CLOSED) {
|
||||
return Promise.resolve(undefined);
|
||||
}
|
||||
if (this.ws == null) {
|
||||
return Promise.resolve(undefined);
|
||||
}
|
||||
return new Promise((resolve) => {
|
||||
if (this.ws == null) {
|
||||
resolve(undefined);
|
||||
}
|
||||
if (this.ws != null) {
|
||||
this.ws.once('close', (code) => resolve(code));
|
||||
}
|
||||
if (this.ws != null && this.state !== ws_1.default.CLOSING) {
|
||||
this.ws.close(exports.INTENTIONAL_DISCONNECT_CODE);
|
||||
}
|
||||
});
|
||||
});
|
||||
}
|
||||
reconnect() {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
this.emit('reconnect');
|
||||
yield this.disconnect();
|
||||
yield this.connect();
|
||||
});
|
||||
}
|
||||
request(request, timeout) {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
if (!this.shouldBeConnected || this.ws == null) {
|
||||
throw new errors_1.NotConnectedError(JSON.stringify(request), request);
|
||||
}
|
||||
const [id, message, responsePromise] = this.requestManager.createRequest(request, timeout !== null && timeout !== void 0 ? timeout : this.config.timeout);
|
||||
this.trace('send', message);
|
||||
websocketSendAsync(this.ws, message).catch((error) => {
|
||||
this.requestManager.reject(id, error);
|
||||
});
|
||||
return responsePromise;
|
||||
});
|
||||
}
|
||||
getUrl() {
|
||||
var _a;
|
||||
return (_a = this.url) !== null && _a !== void 0 ? _a : '';
|
||||
}
|
||||
onMessage(message) {
|
||||
this.trace('receive', message);
|
||||
let data;
|
||||
try {
|
||||
data = JSON.parse(message);
|
||||
}
|
||||
catch (error) {
|
||||
if (error instanceof Error) {
|
||||
this.emit('error', 'badMessage', error.message, message);
|
||||
}
|
||||
return;
|
||||
}
|
||||
if (data.type == null && data.error) {
|
||||
this.emit('error', data.error, data.error_message, data);
|
||||
return;
|
||||
}
|
||||
if (data.type) {
|
||||
this.emit(data.type, data);
|
||||
}
|
||||
if (data.type === 'response') {
|
||||
try {
|
||||
this.requestManager.handleResponse(data);
|
||||
}
|
||||
catch (error) {
|
||||
if (error instanceof Error) {
|
||||
this.emit('error', 'badMessage', error.message, message);
|
||||
}
|
||||
else {
|
||||
this.emit('error', 'badMessage', error, error);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
onceOpen(connectionTimeoutID) {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
if (this.ws == null) {
|
||||
throw new errors_1.XrplError('onceOpen: ws is null');
|
||||
}
|
||||
this.ws.removeAllListeners();
|
||||
clearTimeout(connectionTimeoutID);
|
||||
this.ws.on('message', (message) => this.onMessage(message));
|
||||
this.ws.on('error', (error) => this.emit('error', 'websocket', error.message, error));
|
||||
this.ws.once('close', (code, reason) => {
|
||||
if (this.ws == null) {
|
||||
throw new errors_1.XrplError('onceClose: ws is null');
|
||||
}
|
||||
this.clearHeartbeatInterval();
|
||||
this.requestManager.rejectAll(new errors_1.DisconnectedError(`websocket was closed, ${new TextDecoder('utf-8').decode(reason)}`));
|
||||
this.ws.removeAllListeners();
|
||||
this.ws = null;
|
||||
if (code === undefined) {
|
||||
const internalErrorCode = 1011;
|
||||
this.emit('disconnected', internalErrorCode);
|
||||
}
|
||||
else {
|
||||
this.emit('disconnected', code);
|
||||
}
|
||||
if (code !== exports.INTENTIONAL_DISCONNECT_CODE && code !== undefined) {
|
||||
this.intentionalDisconnect();
|
||||
}
|
||||
});
|
||||
try {
|
||||
this.retryConnectionBackoff.reset();
|
||||
this.startHeartbeatInterval();
|
||||
this.connectionManager.resolveAllAwaiting();
|
||||
this.emit('connected');
|
||||
}
|
||||
catch (error) {
|
||||
if (error instanceof Error) {
|
||||
this.connectionManager.rejectAllAwaiting(error);
|
||||
yield this.disconnect().catch(() => { });
|
||||
}
|
||||
}
|
||||
});
|
||||
}
|
||||
intentionalDisconnect() {
|
||||
const retryTimeout = this.retryConnectionBackoff.duration();
|
||||
this.trace('reconnect', `Retrying connection in ${retryTimeout}ms.`);
|
||||
this.emit('reconnecting', this.retryConnectionBackoff.attempts);
|
||||
this.reconnectTimeoutID = setTimeout(() => {
|
||||
this.reconnect().catch((error) => {
|
||||
this.emit('error', 'reconnect', error.message, error);
|
||||
});
|
||||
}, retryTimeout);
|
||||
}
|
||||
clearHeartbeatInterval() {
|
||||
if (this.heartbeatIntervalID) {
|
||||
clearInterval(this.heartbeatIntervalID);
|
||||
}
|
||||
}
|
||||
startHeartbeatInterval() {
|
||||
this.clearHeartbeatInterval();
|
||||
this.heartbeatIntervalID = setInterval(() => {
|
||||
void this.heartbeat();
|
||||
}, this.config.timeout);
|
||||
}
|
||||
heartbeat() {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
this.request({ command: 'ping' }).catch(() => __awaiter(this, void 0, void 0, function* () {
|
||||
return this.reconnect().catch((error) => {
|
||||
this.emit('error', 'reconnect', error.message, error);
|
||||
});
|
||||
}));
|
||||
});
|
||||
}
|
||||
onConnectionFailed(errorOrCode) {
|
||||
if (this.ws) {
|
||||
this.ws.removeAllListeners();
|
||||
this.ws.on('error', () => {
|
||||
});
|
||||
this.ws.close();
|
||||
this.ws = null;
|
||||
}
|
||||
if (typeof errorOrCode === 'number') {
|
||||
this.connectionManager.rejectAllAwaiting(new errors_1.NotConnectedError(`Connection failed with code ${errorOrCode}.`, {
|
||||
code: errorOrCode,
|
||||
}));
|
||||
}
|
||||
else if (errorOrCode === null || errorOrCode === void 0 ? void 0 : errorOrCode.message) {
|
||||
this.connectionManager.rejectAllAwaiting(new errors_1.NotConnectedError(errorOrCode.message, errorOrCode));
|
||||
}
|
||||
else {
|
||||
this.connectionManager.rejectAllAwaiting(new errors_1.NotConnectedError('Connection failed.'));
|
||||
}
|
||||
}
|
||||
}
|
||||
exports.Connection = Connection;
|
||||
//# sourceMappingURL=connection.js.map
|
||||
+1
File diff suppressed because one or more lines are too long
+92
@@ -0,0 +1,92 @@
|
||||
/// <reference types="node" />
|
||||
import { EventEmitter } from 'events';
|
||||
import { AccountChannelsRequest, AccountChannelsResponse, AccountCurrenciesRequest, AccountCurrenciesResponse, AccountInfoRequest, AccountInfoResponse, AccountLinesRequest, AccountLinesResponse, AccountNFTsRequest, AccountNFTsResponse, AccountObjectsRequest, AccountObjectsResponse, AccountOffersRequest, AccountOffersResponse, AccountTxRequest, AccountTxResponse, GatewayBalancesRequest, GatewayBalancesResponse, NoRippleCheckRequest, NoRippleCheckResponse, LedgerRequest, LedgerResponse, LedgerClosedRequest, LedgerClosedResponse, LedgerCurrentRequest, LedgerCurrentResponse, LedgerDataRequest, LedgerDataResponse, LedgerEntryRequest, LedgerEntryResponse, SubmitRequest, SubmitResponse, SubmitMultisignedRequest, SubmitMultisignedResponse, TransactionEntryRequest, TransactionEntryResponse, TxRequest, TxResponse, BookOffersRequest, BookOffersResponse, DepositAuthorizedRequest, DepositAuthorizedResponse, PathFindRequest, PathFindResponse, RipplePathFindRequest, RipplePathFindResponse, ChannelVerifyRequest, ChannelVerifyResponse, FeeRequest, FeeResponse, ManifestRequest, ManifestResponse, ServerInfoRequest, ServerInfoResponse, ServerStateRequest, ServerStateResponse, PingRequest, PingResponse, RandomRequest, RandomResponse, LedgerStream, ValidationStream, TransactionStream, PathFindStream, PeerStatusStream, ConsensusStream, SubscribeRequest, SubscribeResponse, UnsubscribeRequest, UnsubscribeResponse, NFTBuyOffersRequest, NFTBuyOffersResponse, NFTSellOffersRequest, NFTSellOffersResponse } from '../models/methods';
|
||||
import { BaseRequest, BaseResponse } from '../models/methods/baseMethod';
|
||||
import { autofill, getLedgerIndex, getOrderbook, getBalances, getXrpBalance, submit, submitAndWait } from '../sugar';
|
||||
import fundWallet from '../Wallet/fundWallet';
|
||||
import { Connection, ConnectionUserOptions } from './connection';
|
||||
export interface ClientOptions extends ConnectionUserOptions {
|
||||
feeCushion?: number;
|
||||
maxFeeXRP?: string;
|
||||
proxy?: string;
|
||||
timeout?: number;
|
||||
}
|
||||
declare class Client extends EventEmitter {
|
||||
readonly connection: Connection;
|
||||
readonly feeCushion: number;
|
||||
readonly maxFeeXRP: string;
|
||||
constructor(server: string, options?: ClientOptions);
|
||||
get url(): string;
|
||||
request(r: AccountChannelsRequest): Promise<AccountChannelsResponse>;
|
||||
request(r: AccountCurrenciesRequest): Promise<AccountCurrenciesResponse>;
|
||||
request(r: AccountInfoRequest): Promise<AccountInfoResponse>;
|
||||
request(r: AccountLinesRequest): Promise<AccountLinesResponse>;
|
||||
request(r: AccountNFTsRequest): Promise<AccountNFTsResponse>;
|
||||
request(r: AccountObjectsRequest): Promise<AccountObjectsResponse>;
|
||||
request(r: AccountOffersRequest): Promise<AccountOffersResponse>;
|
||||
request(r: AccountTxRequest): Promise<AccountTxResponse>;
|
||||
request(r: BookOffersRequest): Promise<BookOffersResponse>;
|
||||
request(r: ChannelVerifyRequest): Promise<ChannelVerifyResponse>;
|
||||
request(r: DepositAuthorizedRequest): Promise<DepositAuthorizedResponse>;
|
||||
request(r: FeeRequest): Promise<FeeResponse>;
|
||||
request(r: GatewayBalancesRequest): Promise<GatewayBalancesResponse>;
|
||||
request(r: LedgerRequest): Promise<LedgerResponse>;
|
||||
request(r: LedgerClosedRequest): Promise<LedgerClosedResponse>;
|
||||
request(r: LedgerCurrentRequest): Promise<LedgerCurrentResponse>;
|
||||
request(r: LedgerDataRequest): Promise<LedgerDataResponse>;
|
||||
request(r: LedgerEntryRequest): Promise<LedgerEntryResponse>;
|
||||
request(r: ManifestRequest): Promise<ManifestResponse>;
|
||||
request(r: NFTBuyOffersRequest): Promise<NFTBuyOffersResponse>;
|
||||
request(r: NFTSellOffersRequest): Promise<NFTSellOffersResponse>;
|
||||
request(r: NoRippleCheckRequest): Promise<NoRippleCheckResponse>;
|
||||
request(r: PathFindRequest): Promise<PathFindResponse>;
|
||||
request(r: PingRequest): Promise<PingResponse>;
|
||||
request(r: RandomRequest): Promise<RandomResponse>;
|
||||
request(r: RipplePathFindRequest): Promise<RipplePathFindResponse>;
|
||||
request(r: ServerInfoRequest): Promise<ServerInfoResponse>;
|
||||
request(r: ServerStateRequest): Promise<ServerStateResponse>;
|
||||
request(r: SubmitRequest): Promise<SubmitResponse>;
|
||||
request(r: SubmitMultisignedRequest): Promise<SubmitMultisignedResponse>;
|
||||
request(r: SubscribeRequest): Promise<SubscribeResponse>;
|
||||
request(r: UnsubscribeRequest): Promise<UnsubscribeResponse>;
|
||||
request(r: TransactionEntryRequest): Promise<TransactionEntryResponse>;
|
||||
request(r: TxRequest): Promise<TxResponse>;
|
||||
request<R extends BaseRequest, T extends BaseResponse>(r: R): Promise<T>;
|
||||
requestNextPage(req: AccountChannelsRequest, resp: AccountChannelsResponse): Promise<AccountChannelsResponse>;
|
||||
requestNextPage(req: AccountLinesRequest, resp: AccountLinesResponse): Promise<AccountLinesResponse>;
|
||||
requestNextPage(req: AccountObjectsRequest, resp: AccountObjectsResponse): Promise<AccountObjectsResponse>;
|
||||
requestNextPage(req: AccountOffersRequest, resp: AccountOffersResponse): Promise<AccountOffersResponse>;
|
||||
requestNextPage(req: AccountTxRequest, resp: AccountTxResponse): Promise<AccountTxResponse>;
|
||||
requestNextPage(req: LedgerDataRequest, resp: LedgerDataResponse): Promise<LedgerDataResponse>;
|
||||
on(event: 'connected', listener: () => void): this;
|
||||
on(event: 'disconnected', listener: (code: number) => void): this;
|
||||
on(event: 'ledgerClosed', listener: (ledger: LedgerStream) => void): this;
|
||||
on(event: 'validationReceived', listener: (validation: ValidationStream) => void): this;
|
||||
on(event: 'transaction', listener: (tx: TransactionStream) => void): this;
|
||||
on(event: 'peerStatusChange', listener: (status: PeerStatusStream) => void): this;
|
||||
on(event: 'consensusPhase', listener: (phase: ConsensusStream) => void): this;
|
||||
on(event: 'manifestReceived', listener: (manifest: ManifestResponse) => void): this;
|
||||
on(event: 'path_find', listener: (path: PathFindStream) => void): this;
|
||||
on(event: 'error', listener: (...err: any[]) => void): this;
|
||||
requestAll(req: AccountChannelsRequest): Promise<AccountChannelsResponse[]>;
|
||||
requestAll(req: AccountLinesRequest): Promise<AccountLinesResponse[]>;
|
||||
requestAll(req: AccountObjectsRequest): Promise<AccountObjectsResponse[]>;
|
||||
requestAll(req: AccountOffersRequest): Promise<AccountOffersResponse[]>;
|
||||
requestAll(req: AccountTxRequest): Promise<AccountTxResponse[]>;
|
||||
requestAll(req: BookOffersRequest): Promise<BookOffersResponse[]>;
|
||||
requestAll(req: LedgerDataRequest): Promise<LedgerDataResponse[]>;
|
||||
connect(): Promise<void>;
|
||||
disconnect(): Promise<void>;
|
||||
isConnected(): boolean;
|
||||
autofill: typeof autofill;
|
||||
submit: typeof submit;
|
||||
submitAndWait: typeof submitAndWait;
|
||||
prepareTransaction: typeof autofill;
|
||||
getXrpBalance: typeof getXrpBalance;
|
||||
getBalances: typeof getBalances;
|
||||
getOrderbook: typeof getOrderbook;
|
||||
getLedgerIndex: typeof getLedgerIndex;
|
||||
fundWallet: typeof fundWallet;
|
||||
}
|
||||
export { Client };
|
||||
//# sourceMappingURL=index.d.ts.map
|
||||
+1
File diff suppressed because one or more lines are too long
+202
@@ -0,0 +1,202 @@
|
||||
"use strict";
|
||||
var __createBinding = (this && this.__createBinding) || (Object.create ? (function(o, m, k, k2) {
|
||||
if (k2 === undefined) k2 = k;
|
||||
var desc = Object.getOwnPropertyDescriptor(m, k);
|
||||
if (!desc || ("get" in desc ? !m.__esModule : desc.writable || desc.configurable)) {
|
||||
desc = { enumerable: true, get: function() { return m[k]; } };
|
||||
}
|
||||
Object.defineProperty(o, k2, desc);
|
||||
}) : (function(o, m, k, k2) {
|
||||
if (k2 === undefined) k2 = k;
|
||||
o[k2] = m[k];
|
||||
}));
|
||||
var __setModuleDefault = (this && this.__setModuleDefault) || (Object.create ? (function(o, v) {
|
||||
Object.defineProperty(o, "default", { enumerable: true, value: v });
|
||||
}) : function(o, v) {
|
||||
o["default"] = v;
|
||||
});
|
||||
var __importStar = (this && this.__importStar) || function (mod) {
|
||||
if (mod && mod.__esModule) return mod;
|
||||
var result = {};
|
||||
if (mod != null) for (var k in mod) if (k !== "default" && Object.prototype.hasOwnProperty.call(mod, k)) __createBinding(result, mod, k);
|
||||
__setModuleDefault(result, mod);
|
||||
return result;
|
||||
};
|
||||
var __awaiter = (this && this.__awaiter) || function (thisArg, _arguments, P, generator) {
|
||||
function adopt(value) { return value instanceof P ? value : new P(function (resolve) { resolve(value); }); }
|
||||
return new (P || (P = Promise))(function (resolve, reject) {
|
||||
function fulfilled(value) { try { step(generator.next(value)); } catch (e) { reject(e); } }
|
||||
function rejected(value) { try { step(generator["throw"](value)); } catch (e) { reject(e); } }
|
||||
function step(result) { result.done ? resolve(result.value) : adopt(result.value).then(fulfilled, rejected); }
|
||||
step((generator = generator.apply(thisArg, _arguments || [])).next());
|
||||
});
|
||||
};
|
||||
var __importDefault = (this && this.__importDefault) || function (mod) {
|
||||
return (mod && mod.__esModule) ? mod : { "default": mod };
|
||||
};
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
exports.Client = void 0;
|
||||
const assert = __importStar(require("assert"));
|
||||
const events_1 = require("events");
|
||||
const errors_1 = require("../errors");
|
||||
const sugar_1 = require("../sugar");
|
||||
const fundWallet_1 = __importDefault(require("../Wallet/fundWallet"));
|
||||
const connection_1 = require("./connection");
|
||||
const partialPayment_1 = require("./partialPayment");
|
||||
function getCollectKeyFromCommand(command) {
|
||||
switch (command) {
|
||||
case 'account_channels':
|
||||
return 'channels';
|
||||
case 'account_lines':
|
||||
return 'lines';
|
||||
case 'account_objects':
|
||||
return 'account_objects';
|
||||
case 'account_tx':
|
||||
return 'transactions';
|
||||
case 'account_offers':
|
||||
case 'book_offers':
|
||||
return 'offers';
|
||||
case 'ledger_data':
|
||||
return 'state';
|
||||
default:
|
||||
return null;
|
||||
}
|
||||
}
|
||||
function clamp(value, min, max) {
|
||||
assert.ok(min <= max, 'Illegal clamp bounds');
|
||||
return Math.min(Math.max(value, min), max);
|
||||
}
|
||||
const DEFAULT_FEE_CUSHION = 1.2;
|
||||
const DEFAULT_MAX_FEE_XRP = '2';
|
||||
const MIN_LIMIT = 10;
|
||||
const MAX_LIMIT = 400;
|
||||
const NORMAL_DISCONNECT_CODE = 1000;
|
||||
class Client extends events_1.EventEmitter {
|
||||
constructor(server, options = {}) {
|
||||
var _a, _b;
|
||||
super();
|
||||
this.autofill = sugar_1.autofill;
|
||||
this.submit = sugar_1.submit;
|
||||
this.submitAndWait = sugar_1.submitAndWait;
|
||||
this.prepareTransaction = sugar_1.autofill;
|
||||
this.getXrpBalance = sugar_1.getXrpBalance;
|
||||
this.getBalances = sugar_1.getBalances;
|
||||
this.getOrderbook = sugar_1.getOrderbook;
|
||||
this.getLedgerIndex = sugar_1.getLedgerIndex;
|
||||
this.fundWallet = fundWallet_1.default;
|
||||
if (typeof server !== 'string' || !/wss?(?:\+unix)?:\/\//u.exec(server)) {
|
||||
throw new errors_1.ValidationError('server URI must start with `wss://`, `ws://`, `wss+unix://`, or `ws+unix://`.');
|
||||
}
|
||||
this.feeCushion = (_a = options.feeCushion) !== null && _a !== void 0 ? _a : DEFAULT_FEE_CUSHION;
|
||||
this.maxFeeXRP = (_b = options.maxFeeXRP) !== null && _b !== void 0 ? _b : DEFAULT_MAX_FEE_XRP;
|
||||
this.connection = new connection_1.Connection(server, options);
|
||||
this.connection.on('error', (errorCode, errorMessage, data) => {
|
||||
this.emit('error', errorCode, errorMessage, data);
|
||||
});
|
||||
this.connection.on('connected', () => {
|
||||
this.emit('connected');
|
||||
});
|
||||
this.connection.on('disconnected', (code) => {
|
||||
let finalCode = code;
|
||||
if (finalCode === connection_1.INTENTIONAL_DISCONNECT_CODE) {
|
||||
finalCode = NORMAL_DISCONNECT_CODE;
|
||||
}
|
||||
this.emit('disconnected', finalCode);
|
||||
});
|
||||
this.connection.on('ledgerClosed', (ledger) => {
|
||||
this.emit('ledgerClosed', ledger);
|
||||
});
|
||||
this.connection.on('transaction', (tx) => {
|
||||
(0, partialPayment_1.handleStreamPartialPayment)(tx, this.connection.trace);
|
||||
this.emit('transaction', tx);
|
||||
});
|
||||
this.connection.on('validationReceived', (validation) => {
|
||||
this.emit('validationReceived', validation);
|
||||
});
|
||||
this.connection.on('manifestReceived', (manifest) => {
|
||||
this.emit('manifestReceived', manifest);
|
||||
});
|
||||
this.connection.on('peerStatusChange', (status) => {
|
||||
this.emit('peerStatusChange', status);
|
||||
});
|
||||
this.connection.on('consensusPhase', (consensus) => {
|
||||
this.emit('consensusPhase', consensus);
|
||||
});
|
||||
this.connection.on('path_find', (path) => {
|
||||
this.emit('path_find', path);
|
||||
});
|
||||
}
|
||||
get url() {
|
||||
return this.connection.getUrl();
|
||||
}
|
||||
request(req) {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
const response = (yield this.connection.request(Object.assign(Object.assign({}, req), { account: req.account
|
||||
?
|
||||
(0, sugar_1.ensureClassicAddress)(req.account)
|
||||
: undefined })));
|
||||
(0, partialPayment_1.handlePartialPayment)(req.command, response);
|
||||
return response;
|
||||
});
|
||||
}
|
||||
requestNextPage(req, resp) {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
if (!resp.result.marker) {
|
||||
return Promise.reject(new errors_1.NotFoundError('response does not have a next page'));
|
||||
}
|
||||
const nextPageRequest = Object.assign(Object.assign({}, req), { marker: resp.result.marker });
|
||||
return this.request(nextPageRequest);
|
||||
});
|
||||
}
|
||||
on(eventName, listener) {
|
||||
return super.on(eventName, listener);
|
||||
}
|
||||
requestAll(request, collect) {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
const collectKey = collect !== null && collect !== void 0 ? collect : getCollectKeyFromCommand(request.command);
|
||||
if (!collectKey) {
|
||||
throw new errors_1.ValidationError(`no collect key for command ${request.command}`);
|
||||
}
|
||||
const countTo = request.limit == null ? Infinity : request.limit;
|
||||
let count = 0;
|
||||
let marker = request.marker;
|
||||
let lastBatchLength;
|
||||
const results = [];
|
||||
do {
|
||||
const countRemaining = clamp(countTo - count, MIN_LIMIT, MAX_LIMIT);
|
||||
const repeatProps = Object.assign(Object.assign({}, request), { limit: countRemaining, marker });
|
||||
const singleResponse = yield this.connection.request(repeatProps);
|
||||
const singleResult = singleResponse.result;
|
||||
if (!(collectKey in singleResult)) {
|
||||
throw new errors_1.XrplError(`${collectKey} not in result`);
|
||||
}
|
||||
const collectedData = singleResult[collectKey];
|
||||
marker = singleResult.marker;
|
||||
results.push(singleResponse);
|
||||
if (Array.isArray(collectedData)) {
|
||||
count += collectedData.length;
|
||||
lastBatchLength = collectedData.length;
|
||||
}
|
||||
else {
|
||||
lastBatchLength = 0;
|
||||
}
|
||||
} while (Boolean(marker) && count < countTo && lastBatchLength !== 0);
|
||||
return results;
|
||||
});
|
||||
}
|
||||
connect() {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
return this.connection.connect();
|
||||
});
|
||||
}
|
||||
disconnect() {
|
||||
return __awaiter(this, void 0, void 0, function* () {
|
||||
yield this.connection.disconnect();
|
||||
});
|
||||
}
|
||||
isConnected() {
|
||||
return this.connection.isConnected();
|
||||
}
|
||||
}
|
||||
exports.Client = Client;
|
||||
//# sourceMappingURL=index.js.map
|
||||
+1
File diff suppressed because one or more lines are too long
+4
@@ -0,0 +1,4 @@
|
||||
import type { Response, TransactionStream } from '..';
|
||||
export declare function handlePartialPayment(command: string, response: Response): void;
|
||||
export declare function handleStreamPartialPayment(stream: TransactionStream, log: (id: string, message: string) => void): void;
|
||||
//# sourceMappingURL=partialPayment.d.ts.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"partialPayment.d.ts","sourceRoot":"","sources":["../../../src/client/partialPayment.ts"],"names":[],"mappings":"AAGA,OAAO,KAAK,EAEV,QAAQ,EAER,iBAAiB,EAElB,MAAM,IAAI,CAAA;AAmGX,wBAAgB,oBAAoB,CAClC,OAAO,EAAE,MAAM,EACf,QAAQ,EAAE,QAAQ,GACjB,IAAI,CAaN;AAQD,wBAAgB,0BAA0B,CACxC,MAAM,EAAE,iBAAiB,EACzB,GAAG,EAAE,CAAC,EAAE,EAAE,MAAM,EAAE,OAAO,EAAE,MAAM,KAAK,IAAI,GACzC,IAAI,CAgBN"}
|
||||
+100
@@ -0,0 +1,100 @@
|
||||
"use strict";
|
||||
var __importDefault = (this && this.__importDefault) || function (mod) {
|
||||
return (mod && mod.__esModule) ? mod : { "default": mod };
|
||||
};
|
||||
Object.defineProperty(exports, "__esModule", { value: true });
|
||||
exports.handleStreamPartialPayment = exports.handlePartialPayment = void 0;
|
||||
const bignumber_js_1 = __importDefault(require("bignumber.js"));
|
||||
const ripple_binary_codec_1 = require("ripple-binary-codec");
|
||||
const transactions_1 = require("../models/transactions");
|
||||
const utils_1 = require("../models/utils");
|
||||
const WARN_PARTIAL_PAYMENT_CODE = 2001;
|
||||
function amountsEqual(amt1, amt2) {
|
||||
if (typeof amt1 === 'string' && typeof amt2 === 'string') {
|
||||
return amt1 === amt2;
|
||||
}
|
||||
if (typeof amt1 === 'string' || typeof amt2 === 'string') {
|
||||
return false;
|
||||
}
|
||||
const aValue = new bignumber_js_1.default(amt1.value);
|
||||
const bValue = new bignumber_js_1.default(amt2.value);
|
||||
return (amt1.currency === amt2.currency &&
|
||||
amt1.issuer === amt2.issuer &&
|
||||
aValue.isEqualTo(bValue));
|
||||
}
|
||||
function isPartialPayment(tx, metadata) {
|
||||
var _a;
|
||||
if (tx == null || metadata == null || tx.TransactionType !== 'Payment') {
|
||||
return false;
|
||||
}
|
||||
let meta = metadata;
|
||||
if (typeof meta === 'string') {
|
||||
if (meta === 'unavailable') {
|
||||
return false;
|
||||
}
|
||||
meta = (0, ripple_binary_codec_1.decode)(meta);
|
||||
}
|
||||
const tfPartial = typeof tx.Flags === 'number'
|
||||
? (0, utils_1.isFlagEnabled)(tx.Flags, transactions_1.PaymentFlags.tfPartialPayment)
|
||||
: (_a = tx.Flags) === null || _a === void 0 ? void 0 : _a.tfPartialPayment;
|
||||
if (!tfPartial) {
|
||||
return false;
|
||||
}
|
||||
const delivered = meta.delivered_amount;
|
||||
const amount = tx.Amount;
|
||||
if (delivered === undefined) {
|
||||
return false;
|
||||
}
|
||||
return !amountsEqual(delivered, amount);
|
||||
}
|
||||
function txHasPartialPayment(response) {
|
||||
return isPartialPayment(response.result, response.result.meta);
|
||||
}
|
||||
function txEntryHasPartialPayment(response) {
|
||||
return isPartialPayment(response.result.tx_json, response.result.metadata);
|
||||
}
|
||||
function accountTxHasPartialPayment(response) {
|
||||
const { transactions } = response.result;
|
||||
const foo = transactions.some((tx) => isPartialPayment(tx.tx, tx.meta));
|
||||
return foo;
|
||||
}
|
||||
function hasPartialPayment(command, response) {
|
||||
switch (command) {
|
||||
case 'tx':
|
||||
return txHasPartialPayment(response);
|
||||
case 'transaction_entry':
|
||||
return txEntryHasPartialPayment(response);
|
||||
case 'account_tx':
|
||||
return accountTxHasPartialPayment(response);
|
||||
default:
|
||||
return false;
|
||||
}
|
||||
}
|
||||
function handlePartialPayment(command, response) {
|
||||
var _a;
|
||||
if (hasPartialPayment(command, response)) {
|
||||
const warnings = (_a = response.warnings) !== null && _a !== void 0 ? _a : [];
|
||||
const warning = {
|
||||
id: WARN_PARTIAL_PAYMENT_CODE,
|
||||
message: 'This response contains a Partial Payment',
|
||||
};
|
||||
warnings.push(warning);
|
||||
response.warnings = warnings;
|
||||
}
|
||||
}
|
||||
exports.handlePartialPayment = handlePartialPayment;
|
||||
function handleStreamPartialPayment(stream, log) {
|
||||
var _a;
|
||||
if (isPartialPayment(stream.transaction, stream.meta)) {
|
||||
const warnings = (_a = stream.warnings) !== null && _a !== void 0 ? _a : [];
|
||||
const warning = {
|
||||
id: WARN_PARTIAL_PAYMENT_CODE,
|
||||
message: 'This response contains a Partial Payment',
|
||||
};
|
||||
warnings.push(warning);
|
||||
stream.warnings = warnings;
|
||||
log('Partial payment received', JSON.stringify(stream));
|
||||
}
|
||||
}
|
||||
exports.handleStreamPartialPayment = handleStreamPartialPayment;
|
||||
//# sourceMappingURL=partialPayment.js.map
|
||||
+1
@@ -0,0 +1 @@
|
||||
{"version":3,"file":"partialPayment.js","sourceRoot":"","sources":["../../../src/client/partialPayment.ts"],"names":[],"mappings":";;;;;;AAAA,gEAAoC;AACpC,6DAA4C;AAU5C,yDAAkE;AAElE,2CAA+C;AAE/C,MAAM,yBAAyB,GAAG,IAAI,CAAA;AAEtC,SAAS,YAAY,CAAC,IAAY,EAAE,IAAY;IAC9C,IAAI,OAAO,IAAI,KAAK,QAAQ,IAAI,OAAO,IAAI,KAAK,QAAQ,EAAE;QACxD,OAAO,IAAI,KAAK,IAAI,CAAA;KACrB;IAED,IAAI,OAAO,IAAI,KAAK,QAAQ,IAAI,OAAO,IAAI,KAAK,QAAQ,EAAE;QACxD,OAAO,KAAK,CAAA;KACb;IAED,MAAM,MAAM,GAAG,IAAI,sBAAS,CAAC,IAAI,CAAC,KAAK,CAAC,CAAA;IACxC,MAAM,MAAM,GAAG,IAAI,sBAAS,CAAC,IAAI,CAAC,KAAK,CAAC,CAAA;IAExC,OAAO,CACL,IAAI,CAAC,QAAQ,KAAK,IAAI,CAAC,QAAQ;QAC/B,IAAI,CAAC,MAAM,KAAK,IAAI,CAAC,MAAM;QAC3B,MAAM,CAAC,SAAS,CAAC,MAAM,CAAC,CACzB,CAAA;AACH,CAAC;AAED,SAAS,gBAAgB,CACvB,EAAgB,EAChB,QAAuC;;IAEvC,IAAI,EAAE,IAAI,IAAI,IAAI,QAAQ,IAAI,IAAI,IAAI,EAAE,CAAC,eAAe,KAAK,SAAS,EAAE;QACtE,OAAO,KAAK,CAAA;KACb;IAED,IAAI,IAAI,GAAG,QAAQ,CAAA;IACnB,IAAI,OAAO,IAAI,KAAK,QAAQ,EAAE;QAC5B,IAAI,IAAI,KAAK,aAAa,EAAE;YAC1B,OAAO,KAAK,CAAA;SACb;QAGD,IAAI,GAAG,IAAA,4BAAM,EAAC,IAAI,CAAmC,CAAA;KACtD;IAED,MAAM,SAAS,GACb,OAAO,EAAE,CAAC,KAAK,KAAK,QAAQ;QAC1B,CAAC,CAAC,IAAA,qBAAa,EAAC,EAAE,CAAC,KAAK,EAAE,2BAAY,CAAC,gBAAgB,CAAC;QACxD,CAAC,CAAC,MAAA,EAAE,CAAC,KAAK,0CAAE,gBAAgB,CAAA;IAEhC,IAAI,CAAC,SAAS,EAAE;QACd,OAAO,KAAK,CAAA;KACb;IAED,MAAM,SAAS,GAAG,IAAI,CAAC,gBAAgB,CAAA;IACvC,MAAM,MAAM,GAAG,EAAE,CAAC,MAAM,CAAA;IAExB,IAAI,SAAS,KAAK,SAAS,EAAE;QAC3B,OAAO,KAAK,CAAA;KACb;IAED,OAAO,CAAC,YAAY,CAAC,SAAS,EAAE,MAAM,CAAC,CAAA;AACzC,CAAC;AAED,SAAS,mBAAmB,CAAC,QAAoB;IAC/C,OAAO,gBAAgB,CAAC,QAAQ,CAAC,MAAM,EAAE,QAAQ,CAAC,MAAM,CAAC,IAAI,CAAC,CAAA;AAChE,CAAC;AAED,SAAS,wBAAwB,CAAC,QAAkC;IAClE,OAAO,gBAAgB,CAAC,QAAQ,CAAC,MAAM,CAAC,OAAO,EAAE,QAAQ,CAAC,MAAM,CAAC,QAAQ,CAAC,CAAA;AAC5E,CAAC;AAED,SAAS,0BAA0B,CAAC,QAA2B;IAC7D,MAAM,EAAE,YAAY,EAAE,GAAG,QAAQ,CAAC,MAAM,CAAA;IACxC,MAAM,GAAG,GAAG,YAAY,CAAC,IAAI,CAAC,CAAC,EAAE,EAAE,EAAE,CAAC,gBAAgB,CAAC,EAAE,CAAC,EAAE,EAAE,EAAE,CAAC,IAAI,CAAC,CAAC,CAAA;IACvE,OAAO,GAAG,CAAA;AACZ,CAAC;AAED,SAAS,iBAAiB,CAAC,OAAe,EAAE,QAAkB;IAE5D,QAAQ,OAAO,EAAE;QACf,KAAK,IAAI;YACP,OAAO,mBAAmB,CAAC,QAAsB,CAAC,CAAA;QACpD,KAAK,mBAAmB;YACtB,OAAO,wBAAwB,CAAC,QAAoC,CAAC,CAAA;QACvE,KAAK,YAAY;YACf,OAAO,0BAA0B,CAAC,QAA6B,CAAC,CAAA;QAClE;YACE,OAAO,KAAK,CAAA;KACf;AAEH,CAAC;AAQD,SAAgB,oBAAoB,CAClC,OAAe,EACf,QAAkB;;IAElB,IAAI,iBAAiB,CAAC,OAAO,EAAE,QAAQ,CAAC,EAAE;QACxC,MAAM,QAAQ,GAAG,MAAA,QAAQ,CAAC,QAAQ,mCAAI,EAAE,CAAA;QAExC,MAAM,OAAO,GAAG;YACd,EAAE,EAAE,yBAAyB;YAC7B,OAAO,EAAE,0CAA0C;SACpD,CAAA;QAED,QAAQ,CAAC,IAAI,CAAC,OAAO,CAAC,CAAA;QAEtB,QAAQ,CAAC,QAAQ,GAAG,QAAQ,CAAA;KAC7B;AACH,CAAC;AAhBD,oDAgBC;AAQD,SAAgB,0BAA0B,CACxC,MAAyB,EACzB,GAA0C;;IAE1C,IAAI,gBAAgB,CAAC,MAAM,CAAC,WAAW,EAAE,MAAM,CAAC,IAAI,CAAC,EAAE;QACrD,MAAM,QAAQ,GAAG,MAAA,MAAM,CAAC,QAAQ,mCAAI,EAAE,CAAA;QAEtC,MAAM,OAAO,GAAG;YACd,EAAE,EAAE,yBAAyB;YAC7B,OAAO,EAAE,0CAA0C;SACpD,CAAA;QAED,QAAQ,CAAC,IAAI,CAAC,OAAO,CAAC,CAAA;QAGtB,MAAM,CAAC,QAAQ,GAAG,QAAQ,CAAA;QAE1B,GAAG,CAAC,0BAA0B,EAAE,IAAI,CAAC,SAAS,CAAC,MAAM,CAAC,CAAC,CAAA;KACxD;AACH,CAAC;AAnBD,gEAmBC"}
|
||||
Reference in New Issue
Block a user