diff --git a/http/index.js b/http/index.js index ebdede8..65a26d3 100644 --- a/http/index.js +++ b/http/index.js @@ -21,10 +21,10 @@ require("./utils/interceptors").enable(module.exports); const PollRequestManager = require("./utils/PollRequestManager"); const rm = new PollRequestManager(module.exports.fetch); -module.exports.poll = function (url, options, connectionTimeout, delayStart) { +module.exports.poll = function (url, options, connectionTimeout, delayStart, abortController) { connectionTimeout = connectionTimeout || 10000; rm.setConnectionTimeout(connectionTimeout); - const request = rm.createRequest(url, options, delayStart); + const request = rm.createRequest(url, options, delayStart, abortController); return request; }; diff --git a/http/utils/PollRequestManager.js b/http/utils/PollRequestManager.js index 83fc85f..a45e9cb 100644 --- a/http/utils/PollRequestManager.js +++ b/http/utils/PollRequestManager.js @@ -2,25 +2,15 @@ function PollRequestManager(fetchFunction, connectionTimeout = 10000, pollingTi const requests = new Map(); - function Request(url, options, delay = 0) { - let promiseHandlers = {}; - let currentState = undefined; + function Request(url, options, delay = 0, abortController) { + let timeout; this.url = url; - let abortController; - let previousAbortController; this.execute = function() { - if (typeof AbortController !== "undefined") { - if (typeof abortController === "undefined") { - previousAbortController = new AbortController() - } else { - previousAbortController = abortController; - } - abortController = new AbortController(); - options.signal = previousAbortController.signal; - } - if (!currentState && delay) { + let currentState; + options.signal = abortController.signal; // Always overwrite request options to prevent external overwriting + if (!timeout && delay) { currentState = new Promise((resolve, reject) => { timeout = setTimeout(() => { fetchFunction(url, options).then((response) => { @@ -36,73 +26,25 @@ function PollRequestManager(fetchFunction, connectionTimeout = 10000, pollingTi return currentState; } - this.cancelExecution = function() { + this.abort = () => { clearTimeout(timeout); timeout = undefined; - if(typeof currentState !== "undefined"){ - currentState = undefined; - } - promiseHandlers.resolve = (...args) => {console.log("(not important) Resolve called after cancel execution with the following args", ...args)}; - promiseHandlers.reject = (...args) => {console.log("(not important) Reject called after cancel execution with the following args", ...args)}; - } - - this.setExecutor = function(resolve, reject) { - if(promiseHandlers.resolve){ - return reject(new Error("Request already in progress")); - } - promiseHandlers.resolve = resolve; - promiseHandlers.reject = reject; - } - - this.resolve = function(...args) { - promiseHandlers.resolve(...args); - this.destroy(); - promiseHandlers = {}; - } - - this.reject = function(...args) { - if(promiseHandlers.reject){ - promiseHandlers.reject(...args); - } - this.destroy(); - promiseHandlers = {}; - } - - this.destroy = function(removeFromPool = true) { - this.cancelExecution(); - - if (!removeFromPool) { - return; - } - - // Find our identifier - const requestsEntries = requests.entries() - let identifier; - for (const [key, value] of requestsEntries) { - if (value === this) { - identifier = key; - break; - } - } - - if (identifier) { - requests.delete(identifier); - } - } - - this.abort = () => { - if (typeof previousAbortController !== "undefined") { - previousAbortController.abort(); - } + abortController.abort(); } } - this.createRequest = function (url, options, delayedStart = 0) { - const request = new Request(url, options, delayedStart); + this.createRequest = function (url, options, delayedStart = 0, abortController) { + if (!abortController) + abortController = new AbortController(); + + const request = new Request(url, options, delayedStart, abortController); const promise = new Promise((resolve, reject) => { - request.setExecutor(resolve, reject); - createPollingTask(request); + createPollingTask(request).then((response) => { + resolve(response); + }).catch((err) => { + reject(err); + }) }); promise.abort = () => { this.cancelRequest(promise); @@ -120,8 +62,10 @@ function PollRequestManager(fetchFunction, connectionTimeout = 10000, pollingTi const request = requests.get(promiseOfRequest); if (request) { - request.destroy(false); + request.abort(); requests.delete(promiseOfRequest); + } else { + console.warn("No active request found."); } } @@ -131,84 +75,91 @@ function PollRequestManager(fetchFunction, connectionTimeout = 10000, pollingTi /* *************************** polling zone ****************************/ function createPollingTask(request) { - let safePeriodTimeoutHandler; - let serverResponded = false; - /** - * default connection timeout in api-hub is @connectionTimeout - * we wait double the time before aborting the request - */ - function beginSafePeriod() { - safePeriodTimeoutHandler = setTimeout(() => { - if (!serverResponded) { - request.abort(); - } - serverResponded = false; - beginSafePeriod() - }, connectionTimeout * 2); - reArm(); - } + return new Promise((resolve, reject) => { + let safePeriodTimeoutHandler; + let serverResponded = false; + /** + * default connection timeout in api-hub is @connectionTimeout + * we wait double the time before aborting the request + */ + function beginSafePeriod() { + safePeriodTimeoutHandler = setTimeout(() => { + if (!serverResponded) { + request.abort(); + } + serverResponded = false; + beginSafePeriod() + }, connectionTimeout * 2); + reArm(); + } - function endSafePeriod(serverHasResponded) { - serverResponded = serverHasResponded; + function endSafePeriod(serverHasResponded) { + serverResponded = serverHasResponded; - clearTimeout(safePeriodTimeoutHandler); - } + clearTimeout(safePeriodTimeoutHandler); + } - function reArm() { - request.execute().then( (response) => { - if (!response.ok) { - endSafePeriod(true); + function reArm() { + request.execute().then( (response) => { + if (!response.ok) { + endSafePeriod(true); + + //todo check for http errors like 404 + if (response.status === 403) { + reject(Error("Token expired")); + return + } + + if (response.status === 503){ + let err = Error(response.statusText || "Service unavailable"); + err.code = 503; + throw err; + return; + } - //todo check for http errors like 404 - if (response.status === 403) { - request.reject(Error("Token expired")); - return + return beginSafePeriod(); } - if (response.status === 503){ - let err = Error(response.statusText || "Service unavailable"); - err.code = 503; - throw err; + if (response.status === 204) { + endSafePeriod(true); + beginSafePeriod(); return; } - return beginSafePeriod(); - } - - if (response.status === 204) { - endSafePeriod(true); - beginSafePeriod(); - return; - } + if (safePeriodTimeoutHandler) { + clearTimeout(safePeriodTimeoutHandler); + } - if (safePeriodTimeoutHandler) { - clearTimeout(safePeriodTimeoutHandler); - } + resolve(response); + }).catch( (err) => { + switch (err.code) { + case "ETIMEDOUT": + case "ECONNREFUSED": + endSafePeriod(true); + beginSafePeriod(); + break; + case 20: + case "ERR_NETWORK_IO_SUSPENDED": + //reproduced when user is idle on ios (chrome). + case "ERR_INTERNET_DISCONNECTED": + //indicates a general network failure. + break; + case "ABORT_ERR": + endSafePeriod(true); + reject(err); + break; + default: + console.log("abnormal error: ", err); + endSafePeriod(true); + reject(err); + } + }); - request.resolve(response); - }).catch( (err) => { - switch (err.code) { - case "ETIMEDOUT": - case "ECONNREFUSED": - endSafePeriod(true); - beginSafePeriod(); - break; - case 20: - case "ERR_NETWORK_IO_SUSPENDED": - //reproduced when user is idle on ios (chrome). - case "ERR_INTERNET_DISCONNECTED": - //indicates a general network failure. - break; - default: - console.log("abnormal error: ", err); - endSafePeriod(true); - request.reject(err); - } - }); + } - } + beginSafePeriod(); - beginSafePeriod(); + }) } } diff --git a/mq/mqClient.js b/mq/mqClient.js index 8f61a91..c1a623e 100644 --- a/mq/mqClient.js +++ b/mq/mqClient.js @@ -144,7 +144,12 @@ function MQHandler(didDocument, domain, pollingTimeout) { } } - function ensureAuth(callback) { + function ensureAuth(options, callback) { + if (typeof options === "function") { + callback = options; + options = {}; + } + getURL(queueName, "token", (err, url) => { if (err) { return callback(err); @@ -152,7 +157,7 @@ function MQHandler(didDocument, domain, pollingTimeout) { if (!token || (expiryTime && Date.now() + 2000 > expiryTime)) { callback = $$.makeSaneCallback(callback); - return http.fetch(url) + return http.fetch(url, options) .then(response => { connectionTimeout = parseInt(response.headers.get("connection-timeout")); return response.json() @@ -186,42 +191,68 @@ function MQHandler(didDocument, domain, pollingTimeout) { } - function consumeMessage(action, waitForMore, callback) { - if (typeof waitForMore === "function") { - callback = waitForMore; - waitForMore = false; + function consumeMessage(action, waitForMore, abortController, callback) { + + if (typeof abortController === "function") { + callback = abortController; + abortController = new AbortController(); + } + + const signal = abortController.signal; + + function abortCB(callback) { + let msg = "Aborted by client"; + if (self.stopReceivingMessages) { + msg = "Rejected by client. Message handler is set to not receive messages"; + } + callback(new Error(msg)); } - callback.__requestInProgress = true; + + // callback.__requestInProgress = true; ensureAuth((err, token) => { if (err) { return callback(err); } - //somebody called abort before the ensureAuth resolved - if (!callback.__requestInProgress) { - return; + + if (signal.aborted || self.stopReceivingMessages) { + return abortCB(callback); } + didDocument.sign(token, (err, signature) => { if (err) { return callback(createOpenDSUErrorWrapper(`Failed to sign token`, err)); } + if (signal.aborted || self.stopReceivingMessages) { + return abortCB(callback); + } + getURL(queueName, action, signature.toString("hex"), (err, url) => { if (err) { return callback(err); } - let originalCb = callback; + + if (signal.aborted || self.stopReceivingMessages) { + return abortCB(callback); + } + + //let originalCb = callback; //callback = $$.makeSaneCallback(callback); - let options = { headers: { Authorization: token } }; + let options = { headers: { Authorization: token }, signal: signal }; function makeRequest() { - let request = http.poll(url, options, connectionTimeout, timeout); - originalCb.__requestInProgress = request; + if (signal.aborted || self.stopReceivingMessages) { + return abortCB(callback); + } + + let request = http.poll(url, options, connectionTimeout, timeout, abortController); + // originalCb.__requestInProgress = request; request.then(response => response.json()) .then((response) => { - if(self.stopReceivingMessages){ - return callback(new Error("Message rejected by client")); + if (signal.aborted || self.stopReceivingMessages) { + return abortCB(callback); } //the return value of the listing callback helps to stop the polling mechanism in case that //we need to stop to listen for more messages @@ -239,10 +270,6 @@ function MQHandler(didDocument, domain, pollingTimeout) { }); } - //somebody called abort before we arrived here - if (!originalCb.__requestInProgress) { - return; - } makeRequest(); }) }) @@ -255,8 +282,12 @@ function MQHandler(didDocument, domain, pollingTimeout) { }, callback); } - this.previewMessage = (callback) => { - consumeMessage("get", callback); + this.previewMessage = (abortController, callback) => { + if (typeof abortController === "function") { + callback = abortController; + abortController = new AbortController(); + } + consumeMessage("get", false, abortController, callback); }; @@ -278,12 +309,20 @@ function MQHandler(didDocument, domain, pollingTimeout) { } } - this.readMessage = (callback) => { - consumeMessage("get", getSafeMessageRead(callback)); + this.readMessage = (abortController, callback) => { + if (typeof abortController === "function") { + callback = abortController; + abortController = new AbortController(); + } + consumeMessage("get", false, abortController, getSafeMessageRead(callback)); }; - this.readAndWaitForMessages = (callback) => { - consumeMessage("take", true, getSafeMessageRead(callback)); + this.readAndWaitForMessages = (abortController, callback) => { + if (typeof abortController === "function") { + callback = abortController; + abortController = new AbortController(); + } + consumeMessage("take", true, abortController, getSafeMessageRead(callback)); }; this.readAndWaitForMore = (waitForMore, callback) => { @@ -291,25 +330,18 @@ function MQHandler(didDocument, domain, pollingTimeout) { callback = waitForMore; waitForMore = undefined; } - consumeMessage("get", waitForMore ? waitForMore() : true, getSafeMessageRead(callback)); + const abortController = new AbortController(); + consumeMessage("get", waitForMore ? waitForMore() : true, abortController, getSafeMessageRead(callback)); }; this.subscribe = this.readAndWaitForMore; this.abort = (callback) => { - let request = callback.__requestInProgress; - //if we have an object it means that a http.poll request is in progress - if (typeof request === "object") { - request.abort(); - callback.__requestInProgress = undefined; - delete callback.__requestInProgress; + if (callback && callback.abort) { + callback.abort(); console.log("A request was aborted programmatically"); } else { - //if we have true value it means that an ensureAuth is in progress - if (request) { - callback.__requestInProgress = false; - console.log("A request was aborted programmatically"); - } + console.warn("The request could not be aborted."); } } diff --git a/w3cdid/W3CDID_Mixin.js b/w3cdid/W3CDID_Mixin.js index 84bd399..966c68e 100644 --- a/w3cdid/W3CDID_Mixin.js +++ b/w3cdid/W3CDID_Mixin.js @@ -113,7 +113,12 @@ function W3CDID_Mixin(target, enclave) { const mqHandler = require("opendsu") .loadAPI("mq") .getMQHandlerForDID(target); - mqHandler.previewMessage((err, encryptedMessage) => { + + if (target.stopReceivingMessages === true) + return callback("DID is set to not receive messages"); + + const abortController = new AbortController(); + mqHandler.previewMessage(abortController, (err, encryptedMessage) => { if (err) { return callback(createOpenDSUErrorWrapper(`Failed to read message`, err)); } @@ -132,6 +137,12 @@ function W3CDID_Mixin(target, enclave) { target.decryptMessage(message, callback); }) }); + + return { + abort: () => { + abortController.abort(); + } + }; }; target.subscribe = function (callback) { @@ -169,10 +180,12 @@ function W3CDID_Mixin(target, enclave) { .loadAPI("mq") .getMQHandlerForDID(target); + if (target.stopReceivingMessages === true) + return callback("DID is set to not receive messages"); + target.onCallback = (err, encryptedMessage, notificationHandler) => { if (target.stopReceivingMessages) { - console.log(`Received message for unsubscribed DID`); - return; + return callback("DID is set to not receive messages"); } if (err) {