Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 18 additions & 7 deletions lib/RateLimiterQueue.js
Original file line number Diff line number Diff line change
Expand Up @@ -25,7 +25,7 @@ module.exports = class RateLimiterQueue {
}
}

removeTokens(tokens, key = KEY_DEFAULT) {
removeTokens(tokens, key = KEY_DEFAULT, expiresUnixAt = 0) {
if (!this._queueLimiters[key]) {
this._queueLimiters[key] = new RateLimiterQueueInternal(
this._limiterFlexible, {
Expand All @@ -34,7 +34,7 @@ module.exports = class RateLimiterQueue {
})
}

return this._queueLimiters[key].removeTokens(tokens)
return this._queueLimiters[key].removeTokens(tokens, expiresUnixAt)
}
};

Expand All @@ -59,7 +59,7 @@ class RateLimiterQueueInternal {
})
}

removeTokens(tokens) {
removeTokens(tokens, expiresUnixAt = 0) {
const _this = this;

return new Promise((resolve, reject) => {
Expand All @@ -69,7 +69,7 @@ class RateLimiterQueueInternal {
}

if (_this._queue.length > 0) {
_this._queueRequest.call(_this, resolve, reject, tokens);
_this._queueRequest.call(_this, resolve, reject, tokens, expiresUnixAt);
Comment thread
animir marked this conversation as resolved.
} else {
_this._limiterFlexible.consume(_this._key, tokens)
.then((res) => {
Expand All @@ -79,7 +79,7 @@ class RateLimiterQueueInternal {
if (rej instanceof Error) {
reject(rej);
} else {
_this._queueRequest.call(_this, resolve, reject, tokens);
_this._queueRequest.call(_this, resolve, reject, tokens, expiresUnixAt);
if (_this._waitTimeout === null) {
_this._waitTimeout = setTimeout(_this._processFIFO.bind(_this), rej.msBeforeNext);
}
Expand All @@ -89,10 +89,10 @@ class RateLimiterQueueInternal {
})
}

_queueRequest(resolve, reject, tokens) {
_queueRequest(resolve, reject, tokens, expiresUnixAt = 0) {
const _this = this;
if (_this._queue.length < _this._maxQueueSize) {
_this._queue.push({resolve, reject, tokens});
_this._queue.push({resolve, reject, tokens, expiresUnixAt});
} else {
reject(new RateLimiterQueueError(`Number of requests reached it's maximum ${_this._maxQueueSize}`))
}
Expand All @@ -106,6 +106,17 @@ class RateLimiterQueueInternal {
_this._waitTimeout = null;
}

// Reject any queued requests that have reached their expiration deadline
// (expiresUnixAt, in Unix seconds) before they could be fulfilled.
const nowSecs = Math.floor(Date.now() / 1000);
_this._queue = _this._queue.filter((item) => {
if (item.expiresUnixAt && nowSecs >= item.expiresUnixAt) {
item.reject(new RateLimiterQueueError('The request to remove tokens expired before it could be fulfilled'));
return false;
}
return true;
});

if (_this._queue.length === 0) {
return;
}
Expand Down
33 changes: 33 additions & 0 deletions test/RateLimiterQueue.test.js
Original file line number Diff line number Diff line change
Expand Up @@ -129,6 +129,39 @@ describe('RateLimiterQueue with FIFO queue', function RateLimiterQueueTest() {
});
});

it('rejects a queued request that expires before it is fulfilled (expiresUnixAt)', (done) => {
const rlMemory = new RateLimiterMemory({ points: 1, duration: 2 });
const rlQueue = new RateLimiterQueue(rlMemory);
// Consume the only available token so the next request has to be queued.
rlQueue.removeTokens(1).then(() => {
// Allow the queued request to wait only until the current second, so it is
// still queued (and overdue) when the FIFO processor next runs.
const expiresUnixAt = Math.floor(Date.now() / 1000);
rlQueue.removeTokens(1, 'limiter', expiresUnixAt)
.then(() => {
Comment thread
animir marked this conversation as resolved.
done(new Error('queued request should have been rejected as expired'));
})
.catch((err) => {
expect(err instanceof RateLimiterQueueError).to.equal(true);
done();
});
});
});

it('does not reject a queued request whose expiresUnixAt is still in the future', (done) => {
const rlMemory = new RateLimiterMemory({ points: 1, duration: 1 });
const rlQueue = new RateLimiterQueue(rlMemory);
rlQueue.removeTokens(1).then(() => {
const expiresUnixAt = Math.floor(Date.now() / 1000) + 10;
rlQueue.removeTokens(1, 'limiter', expiresUnixAt)
.then((remainingTokens) => {
expect(remainingTokens).to.equal(0);
done();
})
.catch(done);
});
});

it('getTokensRemaining works', (done) => {
const rlMemory = new RateLimiterMemory({ points: 2, duration: 1 });
const rlQueue = new RateLimiterQueue(rlMemory);
Expand Down
11 changes: 10 additions & 1 deletion types.d.ts
Original file line number Diff line number Diff line change
Expand Up @@ -514,7 +514,16 @@ export class RateLimiterQueue {

getTokensRemaining(key?: string | number): Promise<number>;

removeTokens(tokens: number, key?: string | number): Promise<number>;
/**
* Remove tokens from the queue.
*
* @param tokens Number of tokens to remove.
* @param key Optional queue key for separate FIFO queues.
* @param expiresUnixAt Optional absolute deadline as a Unix timestamp in
* seconds. If the request is still queued when this time is reached, it is
* rejected with a `RateLimiterQueueError`. Defaults to `0` (never expires).
*/
removeTokens(tokens: number, key?: string | number, expiresUnixAt?: number): Promise<number>;
}

export class BurstyRateLimiter {
Expand Down
Loading