first commit

This commit is contained in:
Ichitux
2026-04-05 03:08:53 +02:00
commit 1082d36c12
28015 changed files with 3767672 additions and 0 deletions
+21
View File
@@ -0,0 +1,21 @@
The MIT License (MIT)
Copyright (c) 2016 Christian Speckner <cnspeckn@googlemail.com>
Permission is hereby granted, free of charge, to any person obtaining a copy
of this software and associated documentation files (the "Software"), to deal
in the Software without restriction, including without limitation the rights
to use, copy, modify, merge, publish, distribute, sublicense, and/or sell
copies of the Software, and to permit persons to whom the Software is
furnished to do so, subject to the following conditions:
The above copyright notice and this permission notice shall be included in
all copies or substantial portions of the Software.
THE SOFTWARE IS PROVIDED "AS IS", WITHOUT WARRANTY OF ANY KIND, EXPRESS OR
IMPLIED, INCLUDING BUT NOT LIMITED TO THE WARRANTIES OF MERCHANTABILITY,
FITNESS FOR A PARTICULAR PURPOSE AND NONINFRINGEMENT. IN NO EVENT SHALL THE
AUTHORS OR COPYRIGHT HOLDERS BE LIABLE FOR ANY CLAIM, DAMAGES OR OTHER
LIABILITY, WHETHER IN AN ACTION OF CONTRACT, TORT OR OTHERWISE, ARISING FROM,
OUT OF OR IN CONNECTION WITH THE SOFTWARE OR THE USE OR OTHER DEALINGS IN
THE SOFTWARE.
+540
View File
@@ -0,0 +1,540 @@
[![Build status](https://github.com/DirtyHairy/async-mutex/workflows/Build%20and%20Tests/badge.svg)](https://github.com/DirtyHairy/async-mutex/actions?query=workflow%3A%22Build+and+Tests%22)
[![NPM version](https://badge.fury.io/js/async-mutex.svg)](https://badge.fury.io/js/async-mutex)
[![Coverage Status](https://coveralls.io/repos/github/DirtyHairy/async-mutex/badge.svg?branch=master)](https://coveralls.io/github/DirtyHairy/async-mutex?branch=master)
# What is it?
This package implements primitives for synchronizing asynchronous operations in
Javascript.
## Mutex
The term "mutex" usually refers to a data structure used to synchronize
concurrent processes running on different threads. For example, before accessing
a non-threadsafe resource, a thread will lock the mutex. This is guaranteed
to block the thread until no other thread holds a lock on the mutex and thus
enforces exclusive access to the resource. Once the operation is complete, the
thread releases the lock, allowing other threads to acquire a lock and access the
resource.
While Javascript is strictly single-threaded, the asynchronous nature of its
execution model allows for race conditions that require similar synchronization
primitives. Consider for example a library communicating with a web worker that
needs to exchange several subsequent messages with the worker in order to achieve
a task. As these messages are exchanged in an asynchronous manner, it is perfectly
possible that the library is called again during this process. Depending on the
way state is handled during the async process, this will lead to race conditions
that are hard to fix and even harder to track down.
This library solves the problem by applying the concept of mutexes to Javascript.
Locking the mutex will return a promise that resolves once the mutex becomes
available. Once the async process is complete (usually taking multiple
spins of the event loop), a callback supplied to the caller should be called in order
to release the mutex, allowing the next scheduled worker to execute.
# Semaphore
Imagine a situation where you need to control access to several instances of
a shared resource. For example, you might want to distribute images between several
worker processes that perform transformations, or you might want to create a web
crawler that performs a defined number of requests in parallel.
A semaphore is a data structure that is initialized with an arbitrary integer value and that
can be locked multiple times.
As long as the semaphore value is positive, locking it will return the current value
and the locking process will continue execution immediately; the semaphore will
be decremented upon locking. Releasing the lock will increment the semaphore again.
Once the semaphore has reached zero, the next process that attempts to acquire a lock
will be suspended until another process releases its lock and this increments the semaphore
again.
This library provides a semaphore implementation for Javascript that is similar to the
mutex implementation described above.
# How to use it?
## Installation
You can install the library into your project via npm
npm install async-mutex
The library is written in TypeScript and will work in any environment that
supports ES5, ES6 promises and `Array.isArray`. On ancient browsers,
a shim can be used (e.g. [core-js](https://github.com/zloirock/core-js)).
No external typings are required for using this library with
TypeScript (version >= 2).
Starting with Node 12.16 and 13.7, native ES6 style imports are supported.
**WARNING:** Node 13 versions < 13.2.0 fail to import this package correctly.
Node 12 and earlier are fine, as are newer versions of Node 13.
## Importing
**CommonJS:**
```javascript
var Mutex = require('async-mutex').Mutex;
var Semaphore = require('async-mutex').Semaphore;
var withTimeout = require('async-mutex').withTimeout;
```
**ES6:**
```javascript
import {Mutex, Semaphore, withTimeout} from 'async-mutex';
```
**TypeScript:**
```typescript
import {Mutex, MutexInterface, Semaphore, SemaphoreInterface, withTimeout} from 'async-mutex';
```
With the latest version of Node, native ES6 style imports are supported.
## Mutex API
### Creating
```typescript
const mutex = new Mutex();
```
Create a new mutex.
### Synchronized code execution
Promise style:
```typescript
mutex
.runExclusive(() => {
// ...
})
.then((result) => {
// ...
});
```
async/await:
```typescript
await mutex.runExclusive(async () => {
// ...
});
```
`runExclusive` schedules the supplied callback to be run once the mutex is unlocked.
The function may return a promise. Once the promise is resolved or rejected (or immediately after
execution if an immediate value was returned),
the mutex is released. `runExclusive` returns a promise that adopts the state of the function result.
The mutex is released and the result rejected if an exception occurs during execution
of the callback.
### Manual locking / releasing
Promise style:
```typescript
mutex
.acquire()
.then(function(release) {
// ...
release();
});
```
async/await:
```typescript
const release = await mutex.acquire();
try {
// ...
} finally {
release();
}
```
`acquire` returns an (ES6) promise that will resolve as soon as the mutex is
available. The promise resolves with a function `release` that
must be called once the mutex should be released again. The `release` callback
is idempotent.
**IMPORTANT:** Failure to call `release` will hold the mutex locked and will
likely deadlock the application. Make sure to call `release` under all circumstances
and handle exceptions accordingly.
### Unscoped release
As an alternative to calling the `release` callback returned by `acquire`, the mutex
can be released by calling `release` directly on it:
```typescript
mutex.release();
```
### Checking whether the mutex is locked
```typescript
mutex.isLocked();
```
### Cancelling pending locks
Pending locks can be cancelled by calling `cancel()` on the mutex. This will reject
all pending locks with `E_CANCELED`:
Promise style:
```typescript
import {E_CANCELED} from 'async-mutex';
mutex
.runExclusive(() => {
// ...
})
.then(() => {
// ...
})
.catch(e => {
if (e === E_CANCELED) {
// ...
}
});
```
async/await:
```typescript
import {E_CANCELED} from 'async-mutex';
try {
await mutex.runExclusive(() => {
// ...
});
} catch (e) {
if (e === E_CANCELED) {
// ...
}
}
```
This works with `acquire`, too:
if `acquire` is used for locking, the resulting promise will reject with `E_CANCELED`.
The error that is thrown can be customized by passing a different error to the `Mutex`
constructor:
```typescript
const mutex = new Mutex(new Error('fancy custom error'));
```
Note that while all pending locks are cancelled, a currently held lock will not be
revoked. In consequence, the mutex may not be available even after `cancel()` has been called.
### Waiting until the mutex is available
You can wait until the mutex is available without locking it by calling `waitForUnlock()`.
This will return a promise that resolve once the mutex can be acquired again. This operation
will not lock the mutex, and there is no guarantee that the mutex will still be available
once an async barrier has been encountered.
Promise style:
```typescript
mutex
.waitForUnlock()
.then(() => {
// ...
});
```
Async/await:
```typescript
await mutex.waitForUnlock();
// ...
```
## Semaphore API
### Creating
```typescript
const semaphore = new Semaphore(initialValue);
```
Creates a new semaphore. `initialValue` is an arbitrary integer that defines the
initial value of the semaphore.
### Synchronized code execution
Promise style:
```typescript
semaphore
.runExclusive(function(value) {
// ...
})
.then(function(result) {
// ...
});
```
async/await:
```typescript
await semaphore.runExclusive(async (value) => {
// ...
});
```
`runExclusive` schedules the supplied callback to be run once the semaphore is available.
The callback will receive the current value of the semaphore as its argument.
The function may return a promise. Once the promise is resolved or rejected (or immediately after
execution if an immediate value was returned),
the semaphore is released. `runExclusive` returns a promise that adopts the state of the function result.
The semaphore is released and the result rejected if an exception occurs during execution
of the callback.
`runExclusive` accepts a first optional argument `weight`. Specifying a `weight` will decrement the
semaphore by the specified value, and the callback will only be invoked once the semaphore's
value greater or equal to `weight`.
`runExclusive` accepts a second optional argument `priority`. Specifying a greater value for `priority`
tells the scheduler to run this task before other tasks. `priority` can be any real number. The default
is zero.
### Manual locking / releasing
Promise style:
```typescript
semaphore
.acquire()
.then(function([value, release]) {
// ...
release();
});
```
async/await:
```typescript
const [value, release] = await semaphore.acquire();
try {
// ...
} finally {
release();
}
```
`acquire` returns an (ES6) promise that will resolve as soon as the semaphore is
available. The promise resolves to an array with the
first entry being the current value of the semaphore, and the second value a
function that must be called to release the semaphore once the critical operation
has completed. The `release` callback is idempotent.
**IMPORTANT:** Failure to call `release` will hold the semaphore locked and will
likely deadlock the application. Make sure to call `release` under all circumstances
and handle exceptions accordingly.
`acquire` accepts a first optional argument `weight`. Specifying a `weight` will decrement the
semaphore by the specified value, and the semaphore will only be acquired once its
value is greater or equal to `weight`.
`acquire` accepts a second optional argument `priority`. Specifying a greater value for `priority`
tells the scheduler to release the semaphore to the caller before other callers. `priority` can be
any real number. The default is zero.
### Unscoped release
As an alternative to calling the `release` callback returned by `acquire`, the semaphore
can be released by calling `release` directly on it:
```typescript
semaphore.release();
```
`release` accepts an optional argument `weight` and increments the semaphore accordingly.
**IMPORTANT:** Releasing a previously acquired semaphore with the releaser that was
returned by acquire will automatically increment the semaphore by the correct weight. If
you release by calling the unscoped `release` you have to supply the correct weight
yourself!
### Getting the semaphore value
```typescript
semaphore.getValue()
```
### Checking whether the semaphore is locked
```typescript
semaphore.isLocked();
```
The semaphore is considered to be locked if its value is either zero or negative.
### Setting the semaphore value
The value of a semaphore can be set directly to a desired value. A positive value will
cause the semaphore to schedule any pending waiters accordingly.
```typescript
semaphore.setValue();
```
### Cancelling pending locks
Pending locks can be cancelled by calling `cancel()` on the semaphore. This will reject
all pending locks with `E_CANCELED`:
Promise style:
```typescript
import {E_CANCELED} from 'async-mutex';
semaphore
.runExclusive(() => {
// ...
})
.then(() => {
// ...
})
.catch(e => {
if (e === E_CANCELED) {
// ...
}
});
```
async/await:
```typescript
import {E_CANCELED} from 'async-mutex';
try {
await semaphore.runExclusive(() => {
// ...
});
} catch (e) {
if (e === E_CANCELED) {
// ...
}
}
```
This works with `acquire`, too:
if `acquire` is used for locking, the resulting promise will reject with `E_CANCELED`.
The error that is thrown can be customized by passing a different error to the `Semaphore`
constructor:
```typescript
const semaphore = new Semaphore(2, new Error('fancy custom error'));
```
Note that while all pending locks are cancelled, any currently held locks will not be
revoked. In consequence, the semaphore may not be available even after `cancel()` has been called.
### Waiting until the semaphore is available
You can wait until the semaphore is available without locking it by calling `waitForUnlock()`.
This will return a promise that resolve once the semaphore can be acquired again. This operation
will not lock the semaphore, and there is no guarantee that the semaphore will still be available
once an async barrier has been encountered.
Promise style:
```typescript
semaphore
.waitForUnlock()
.then(() => {
// ...
});
```
Async/await:
```typescript
await semaphore.waitForUnlock();
// ...
```
`waitForUnlock` accepts optional arguments `weight` and `priority`. The promise will resolve as soon
as it is possible to `acquire` the semaphore with the given weight and priority. Scheduled tasks with
the greatest `priority` values execute first.
## Limiting the time waiting for a mutex or semaphore to become available
Sometimes it is desirable to limit the time a program waits for a mutex or
semaphore to become available. The `withTimeout` decorator can be applied
to both semaphores and mutexes and changes the behavior of `acquire` and
`runExclusive` accordingly.
```typescript
import {withTimeout, E_TIMEOUT} from 'async-mutex';
const mutexWithTimeout = withTimeout(new Mutex(), 100);
const semaphoreWithTimeout = withTimeout(new Semaphore(5), 100);
```
The API of the decorated mutex or semaphore is unchanged.
The second argument of `withTimeout` is the timeout in milliseconds. After the
timeout is exceeded, the promise returned by `acquire` and `runExclusive` will
reject with `E_TIMEOUT`. The latter will not run the provided callback in case
of an timeout.
The third argument of `withTimeout` is optional and can be used to
customize the error with which the promise is rejected.
```typescript
const mutexWithTimeout = withTimeout(new Mutex(), 100, new Error('new fancy error'));
const semaphoreWithTimeout = withTimeout(new Semaphore(5), 100, new Error('new fancy error'));
```
### Failing early if the mutex or semaphore is not available
A shortcut exists for the case where you do not want to wait for a lock to
be available at all. The `tryAcquire` decorator can be applied to both mutexes
and semaphores and changes the behavior of `acquire` and `runExclusive` to
immediately throw `E_ALREADY_LOCKED` if the mutex is not available.
Promise style:
```typescript
import {tryAcquire, E_ALREADY_LOCKED} from 'async-mutex';
tryAcquire(semaphoreOrMutex)
.runExclusive(() => {
// ...
})
.then(() => {
// ...
})
.catch(e => {
if (e === E_ALREADY_LOCKED) {
// ...
}
});
```
async/await:
```typescript
import {tryAcquire, E_ALREADY_LOCKED} from 'async-mutex';
try {
await tryAcquire(semaphoreOrMutex).runExclusive(() => {
// ...
});
} catch (e) {
if (e === E_ALREADY_LOCKED) {
// ...
}
}
```
Again, the error can be customized by providing a custom error as second argument to
`tryAcquire`.
```typescript
tryAcquire(semaphoreOrMutex, new Error('new fancy error'))
.runExclusive(() => {
// ...
});
```
# License
Feel free to use this library under the conditions of the MIT license.
+41
View File
@@ -0,0 +1,41 @@
import { __awaiter, __generator } from "tslib";
import Semaphore from './Semaphore';
var Mutex = /** @class */ (function () {
function Mutex(cancelError) {
this._semaphore = new Semaphore(1, cancelError);
}
Mutex.prototype.acquire = function () {
return __awaiter(this, arguments, void 0, function (priority) {
var _a, releaser;
if (priority === void 0) { priority = 0; }
return __generator(this, function (_b) {
switch (_b.label) {
case 0: return [4 /*yield*/, this._semaphore.acquire(1, priority)];
case 1:
_a = _b.sent(), releaser = _a[1];
return [2 /*return*/, releaser];
}
});
});
};
Mutex.prototype.runExclusive = function (callback, priority) {
if (priority === void 0) { priority = 0; }
return this._semaphore.runExclusive(function () { return callback(); }, 1, priority);
};
Mutex.prototype.isLocked = function () {
return this._semaphore.isLocked();
};
Mutex.prototype.waitForUnlock = function (priority) {
if (priority === void 0) { priority = 0; }
return this._semaphore.waitForUnlock(1, priority);
};
Mutex.prototype.release = function () {
if (this._semaphore.isLocked())
this._semaphore.release();
};
Mutex.prototype.cancel = function () {
return this._semaphore.cancel();
};
return Mutex;
}());
export default Mutex;
+1
View File
@@ -0,0 +1 @@
export {};
+153
View File
@@ -0,0 +1,153 @@
import { __awaiter, __generator } from "tslib";
import { E_CANCELED } from './errors';
var Semaphore = /** @class */ (function () {
function Semaphore(_value, _cancelError) {
if (_cancelError === void 0) { _cancelError = E_CANCELED; }
this._value = _value;
this._cancelError = _cancelError;
this._queue = [];
this._weightedWaiters = [];
}
Semaphore.prototype.acquire = function (weight, priority) {
var _this = this;
if (weight === void 0) { weight = 1; }
if (priority === void 0) { priority = 0; }
if (weight <= 0)
throw new Error("invalid weight ".concat(weight, ": must be positive"));
return new Promise(function (resolve, reject) {
var task = { resolve: resolve, reject: reject, weight: weight, priority: priority };
var i = findIndexFromEnd(_this._queue, function (other) { return priority <= other.priority; });
if (i === -1 && weight <= _this._value) {
// Needs immediate dispatch, skip the queue
_this._dispatchItem(task);
}
else {
_this._queue.splice(i + 1, 0, task);
}
});
};
Semaphore.prototype.runExclusive = function (callback_1) {
return __awaiter(this, arguments, void 0, function (callback, weight, priority) {
var _a, value, release;
if (weight === void 0) { weight = 1; }
if (priority === void 0) { priority = 0; }
return __generator(this, function (_b) {
switch (_b.label) {
case 0: return [4 /*yield*/, this.acquire(weight, priority)];
case 1:
_a = _b.sent(), value = _a[0], release = _a[1];
_b.label = 2;
case 2:
_b.trys.push([2, , 4, 5]);
return [4 /*yield*/, callback(value)];
case 3: return [2 /*return*/, _b.sent()];
case 4:
release();
return [7 /*endfinally*/];
case 5: return [2 /*return*/];
}
});
});
};
Semaphore.prototype.waitForUnlock = function (weight, priority) {
var _this = this;
if (weight === void 0) { weight = 1; }
if (priority === void 0) { priority = 0; }
if (weight <= 0)
throw new Error("invalid weight ".concat(weight, ": must be positive"));
if (this._couldLockImmediately(weight, priority)) {
return Promise.resolve();
}
else {
return new Promise(function (resolve) {
if (!_this._weightedWaiters[weight - 1])
_this._weightedWaiters[weight - 1] = [];
insertSorted(_this._weightedWaiters[weight - 1], { resolve: resolve, priority: priority });
});
}
};
Semaphore.prototype.isLocked = function () {
return this._value <= 0;
};
Semaphore.prototype.getValue = function () {
return this._value;
};
Semaphore.prototype.setValue = function (value) {
this._value = value;
this._dispatchQueue();
};
Semaphore.prototype.release = function (weight) {
if (weight === void 0) { weight = 1; }
if (weight <= 0)
throw new Error("invalid weight ".concat(weight, ": must be positive"));
this._value += weight;
this._dispatchQueue();
};
Semaphore.prototype.cancel = function () {
var _this = this;
this._queue.forEach(function (entry) { return entry.reject(_this._cancelError); });
this._queue = [];
};
Semaphore.prototype._dispatchQueue = function () {
this._drainUnlockWaiters();
while (this._queue.length > 0 && this._queue[0].weight <= this._value) {
this._dispatchItem(this._queue.shift());
this._drainUnlockWaiters();
}
};
Semaphore.prototype._dispatchItem = function (item) {
var previousValue = this._value;
this._value -= item.weight;
item.resolve([previousValue, this._newReleaser(item.weight)]);
};
Semaphore.prototype._newReleaser = function (weight) {
var _this = this;
var called = false;
return function () {
if (called)
return;
called = true;
_this.release(weight);
};
};
Semaphore.prototype._drainUnlockWaiters = function () {
if (this._queue.length === 0) {
for (var weight = this._value; weight > 0; weight--) {
var waiters = this._weightedWaiters[weight - 1];
if (!waiters)
continue;
waiters.forEach(function (waiter) { return waiter.resolve(); });
this._weightedWaiters[weight - 1] = [];
}
}
else {
var queuedPriority_1 = this._queue[0].priority;
for (var weight = this._value; weight > 0; weight--) {
var waiters = this._weightedWaiters[weight - 1];
if (!waiters)
continue;
var i = waiters.findIndex(function (waiter) { return waiter.priority <= queuedPriority_1; });
(i === -1 ? waiters : waiters.splice(0, i))
.forEach((function (waiter) { return waiter.resolve(); }));
}
}
};
Semaphore.prototype._couldLockImmediately = function (weight, priority) {
return (this._queue.length === 0 || this._queue[0].priority < priority) &&
weight <= this._value;
};
return Semaphore;
}());
function insertSorted(a, v) {
var i = findIndexFromEnd(a, function (other) { return v.priority <= other.priority; });
a.splice(i + 1, 0, v);
}
function findIndexFromEnd(a, predicate) {
for (var i = a.length - 1; i >= 0; i--) {
if (predicate(a[i])) {
return i;
}
}
return -1;
}
export default Semaphore;
+1
View File
@@ -0,0 +1 @@
export {};
+3
View File
@@ -0,0 +1,3 @@
export var E_TIMEOUT = new Error('timeout while waiting for mutex to become available');
export var E_ALREADY_LOCKED = new Error('mutex already locked');
export var E_CANCELED = new Error('request for lock canceled');
+5
View File
@@ -0,0 +1,5 @@
export { default as Mutex } from './Mutex';
export { default as Semaphore } from './Semaphore';
export { withTimeout } from './withTimeout';
export { tryAcquire } from './tryAcquire';
export * from './errors';
+8
View File
@@ -0,0 +1,8 @@
import { E_ALREADY_LOCKED } from './errors';
import { withTimeout } from './withTimeout';
// eslint-disable-next-lisne @typescript-eslint/explicit-module-boundary-types
export function tryAcquire(sync, alreadyAcquiredError) {
if (alreadyAcquiredError === void 0) { alreadyAcquiredError = E_ALREADY_LOCKED; }
// eslint-disable-next-line @typescript-eslint/no-explicit-any
return withTimeout(sync, 0, alreadyAcquiredError);
}
+124
View File
@@ -0,0 +1,124 @@
import { __awaiter, __generator } from "tslib";
/* eslint-disable @typescript-eslint/no-explicit-any */
import { E_TIMEOUT } from './errors';
export function withTimeout(sync, timeout, timeoutError) {
var _this = this;
if (timeoutError === void 0) { timeoutError = E_TIMEOUT; }
return {
acquire: function (weightOrPriority, priority) {
var weight;
if (isSemaphore(sync)) {
weight = weightOrPriority;
}
else {
weight = undefined;
priority = weightOrPriority;
}
if (weight !== undefined && weight <= 0) {
throw new Error("invalid weight ".concat(weight, ": must be positive"));
}
return new Promise(function (resolve, reject) { return __awaiter(_this, void 0, void 0, function () {
var isTimeout, handle, ticket, release, e_1;
return __generator(this, function (_a) {
switch (_a.label) {
case 0:
isTimeout = false;
handle = setTimeout(function () {
isTimeout = true;
reject(timeoutError);
}, timeout);
_a.label = 1;
case 1:
_a.trys.push([1, 3, , 4]);
return [4 /*yield*/, (isSemaphore(sync)
? sync.acquire(weight, priority)
: sync.acquire(priority))];
case 2:
ticket = _a.sent();
if (isTimeout) {
release = Array.isArray(ticket) ? ticket[1] : ticket;
release();
}
else {
clearTimeout(handle);
resolve(ticket);
}
return [3 /*break*/, 4];
case 3:
e_1 = _a.sent();
if (!isTimeout) {
clearTimeout(handle);
reject(e_1);
}
return [3 /*break*/, 4];
case 4: return [2 /*return*/];
}
});
}); });
},
runExclusive: function (callback, weight, priority) {
return __awaiter(this, void 0, void 0, function () {
var release, ticket;
return __generator(this, function (_a) {
switch (_a.label) {
case 0:
release = function () { return undefined; };
_a.label = 1;
case 1:
_a.trys.push([1, , 7, 8]);
return [4 /*yield*/, this.acquire(weight, priority)];
case 2:
ticket = _a.sent();
if (!Array.isArray(ticket)) return [3 /*break*/, 4];
release = ticket[1];
return [4 /*yield*/, callback(ticket[0])];
case 3: return [2 /*return*/, _a.sent()];
case 4:
release = ticket;
return [4 /*yield*/, callback()];
case 5: return [2 /*return*/, _a.sent()];
case 6: return [3 /*break*/, 8];
case 7:
release();
return [7 /*endfinally*/];
case 8: return [2 /*return*/];
}
});
});
},
release: function (weight) {
sync.release(weight);
},
cancel: function () {
return sync.cancel();
},
waitForUnlock: function (weightOrPriority, priority) {
var weight;
if (isSemaphore(sync)) {
weight = weightOrPriority;
}
else {
weight = undefined;
priority = weightOrPriority;
}
if (weight !== undefined && weight <= 0) {
throw new Error("invalid weight ".concat(weight, ": must be positive"));
}
return new Promise(function (resolve, reject) {
var handle = setTimeout(function () { return reject(timeoutError); }, timeout);
(isSemaphore(sync)
? sync.waitForUnlock(weight, priority)
: sync.waitForUnlock(priority)).then(function () {
clearTimeout(handle);
resolve();
});
});
},
isLocked: function () { return sync.isLocked(); },
getValue: function () { return sync.getValue(); },
setValue: function (value) { return sync.setValue(value); },
};
}
function isSemaphore(sync) {
return sync.getValue !== undefined;
}
+291
View File
@@ -0,0 +1,291 @@
const E_TIMEOUT = new Error('timeout while waiting for mutex to become available');
const E_ALREADY_LOCKED = new Error('mutex already locked');
const E_CANCELED = new Error('request for lock canceled');
var __awaiter$2 = (undefined && undefined.__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());
});
};
class Semaphore {
constructor(_value, _cancelError = E_CANCELED) {
this._value = _value;
this._cancelError = _cancelError;
this._queue = [];
this._weightedWaiters = [];
}
acquire(weight = 1, priority = 0) {
if (weight <= 0)
throw new Error(`invalid weight ${weight}: must be positive`);
return new Promise((resolve, reject) => {
const task = { resolve, reject, weight, priority };
const i = findIndexFromEnd(this._queue, (other) => priority <= other.priority);
if (i === -1 && weight <= this._value) {
// Needs immediate dispatch, skip the queue
this._dispatchItem(task);
}
else {
this._queue.splice(i + 1, 0, task);
}
});
}
runExclusive(callback_1) {
return __awaiter$2(this, arguments, void 0, function* (callback, weight = 1, priority = 0) {
const [value, release] = yield this.acquire(weight, priority);
try {
return yield callback(value);
}
finally {
release();
}
});
}
waitForUnlock(weight = 1, priority = 0) {
if (weight <= 0)
throw new Error(`invalid weight ${weight}: must be positive`);
if (this._couldLockImmediately(weight, priority)) {
return Promise.resolve();
}
else {
return new Promise((resolve) => {
if (!this._weightedWaiters[weight - 1])
this._weightedWaiters[weight - 1] = [];
insertSorted(this._weightedWaiters[weight - 1], { resolve, priority });
});
}
}
isLocked() {
return this._value <= 0;
}
getValue() {
return this._value;
}
setValue(value) {
this._value = value;
this._dispatchQueue();
}
release(weight = 1) {
if (weight <= 0)
throw new Error(`invalid weight ${weight}: must be positive`);
this._value += weight;
this._dispatchQueue();
}
cancel() {
this._queue.forEach((entry) => entry.reject(this._cancelError));
this._queue = [];
}
_dispatchQueue() {
this._drainUnlockWaiters();
while (this._queue.length > 0 && this._queue[0].weight <= this._value) {
this._dispatchItem(this._queue.shift());
this._drainUnlockWaiters();
}
}
_dispatchItem(item) {
const previousValue = this._value;
this._value -= item.weight;
item.resolve([previousValue, this._newReleaser(item.weight)]);
}
_newReleaser(weight) {
let called = false;
return () => {
if (called)
return;
called = true;
this.release(weight);
};
}
_drainUnlockWaiters() {
if (this._queue.length === 0) {
for (let weight = this._value; weight > 0; weight--) {
const waiters = this._weightedWaiters[weight - 1];
if (!waiters)
continue;
waiters.forEach((waiter) => waiter.resolve());
this._weightedWaiters[weight - 1] = [];
}
}
else {
const queuedPriority = this._queue[0].priority;
for (let weight = this._value; weight > 0; weight--) {
const waiters = this._weightedWaiters[weight - 1];
if (!waiters)
continue;
const i = waiters.findIndex((waiter) => waiter.priority <= queuedPriority);
(i === -1 ? waiters : waiters.splice(0, i))
.forEach((waiter => waiter.resolve()));
}
}
}
_couldLockImmediately(weight, priority) {
return (this._queue.length === 0 || this._queue[0].priority < priority) &&
weight <= this._value;
}
}
function insertSorted(a, v) {
const i = findIndexFromEnd(a, (other) => v.priority <= other.priority);
a.splice(i + 1, 0, v);
}
function findIndexFromEnd(a, predicate) {
for (let i = a.length - 1; i >= 0; i--) {
if (predicate(a[i])) {
return i;
}
}
return -1;
}
var __awaiter$1 = (undefined && undefined.__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());
});
};
class Mutex {
constructor(cancelError) {
this._semaphore = new Semaphore(1, cancelError);
}
acquire() {
return __awaiter$1(this, arguments, void 0, function* (priority = 0) {
const [, releaser] = yield this._semaphore.acquire(1, priority);
return releaser;
});
}
runExclusive(callback, priority = 0) {
return this._semaphore.runExclusive(() => callback(), 1, priority);
}
isLocked() {
return this._semaphore.isLocked();
}
waitForUnlock(priority = 0) {
return this._semaphore.waitForUnlock(1, priority);
}
release() {
if (this._semaphore.isLocked())
this._semaphore.release();
}
cancel() {
return this._semaphore.cancel();
}
}
var __awaiter = (undefined && undefined.__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());
});
};
function withTimeout(sync, timeout, timeoutError = E_TIMEOUT) {
return {
acquire: (weightOrPriority, priority) => {
let weight;
if (isSemaphore(sync)) {
weight = weightOrPriority;
}
else {
weight = undefined;
priority = weightOrPriority;
}
if (weight !== undefined && weight <= 0) {
throw new Error(`invalid weight ${weight}: must be positive`);
}
return new Promise((resolve, reject) => __awaiter(this, void 0, void 0, function* () {
let isTimeout = false;
const handle = setTimeout(() => {
isTimeout = true;
reject(timeoutError);
}, timeout);
try {
const ticket = yield (isSemaphore(sync)
? sync.acquire(weight, priority)
: sync.acquire(priority));
if (isTimeout) {
const release = Array.isArray(ticket) ? ticket[1] : ticket;
release();
}
else {
clearTimeout(handle);
resolve(ticket);
}
}
catch (e) {
if (!isTimeout) {
clearTimeout(handle);
reject(e);
}
}
}));
},
runExclusive(callback, weight, priority) {
return __awaiter(this, void 0, void 0, function* () {
let release = () => undefined;
try {
const ticket = yield this.acquire(weight, priority);
if (Array.isArray(ticket)) {
release = ticket[1];
return yield callback(ticket[0]);
}
else {
release = ticket;
return yield callback();
}
}
finally {
release();
}
});
},
release(weight) {
sync.release(weight);
},
cancel() {
return sync.cancel();
},
waitForUnlock: (weightOrPriority, priority) => {
let weight;
if (isSemaphore(sync)) {
weight = weightOrPriority;
}
else {
weight = undefined;
priority = weightOrPriority;
}
if (weight !== undefined && weight <= 0) {
throw new Error(`invalid weight ${weight}: must be positive`);
}
return new Promise((resolve, reject) => {
const handle = setTimeout(() => reject(timeoutError), timeout);
(isSemaphore(sync)
? sync.waitForUnlock(weight, priority)
: sync.waitForUnlock(priority)).then(() => {
clearTimeout(handle);
resolve();
});
});
},
isLocked: () => sync.isLocked(),
getValue: () => sync.getValue(),
setValue: (value) => sync.setValue(value),
};
}
function isSemaphore(sync) {
return sync.getValue !== undefined;
}
// eslint-disable-next-lisne @typescript-eslint/explicit-module-boundary-types
function tryAcquire(sync, alreadyAcquiredError = E_ALREADY_LOCKED) {
// eslint-disable-next-line @typescript-eslint/no-explicit-any
return withTimeout(sync, 0, alreadyAcquiredError);
}
export { E_ALREADY_LOCKED, E_CANCELED, E_TIMEOUT, Mutex, Semaphore, tryAcquire, withTimeout };
+12
View File
@@ -0,0 +1,12 @@
import MutexInterface from './MutexInterface';
declare class Mutex implements MutexInterface {
constructor(cancelError?: Error);
acquire(priority?: number): Promise<MutexInterface.Releaser>;
runExclusive<T>(callback: MutexInterface.Worker<T>, priority?: number): Promise<T>;
isLocked(): boolean;
waitForUnlock(priority?: number): Promise<void>;
release(): void;
cancel(): void;
private _semaphore;
}
export default Mutex;
+43
View File
@@ -0,0 +1,43 @@
"use strict";
Object.defineProperty(exports, "__esModule", { value: true });
var tslib_1 = require("tslib");
var Semaphore_1 = require("./Semaphore");
var Mutex = /** @class */ (function () {
function Mutex(cancelError) {
this._semaphore = new Semaphore_1.default(1, cancelError);
}
Mutex.prototype.acquire = function () {
return tslib_1.__awaiter(this, arguments, void 0, function (priority) {
var _a, releaser;
if (priority === void 0) { priority = 0; }
return tslib_1.__generator(this, function (_b) {
switch (_b.label) {
case 0: return [4 /*yield*/, this._semaphore.acquire(1, priority)];
case 1:
_a = _b.sent(), releaser = _a[1];
return [2 /*return*/, releaser];
}
});
});
};
Mutex.prototype.runExclusive = function (callback, priority) {
if (priority === void 0) { priority = 0; }
return this._semaphore.runExclusive(function () { return callback(); }, 1, priority);
};
Mutex.prototype.isLocked = function () {
return this._semaphore.isLocked();
};
Mutex.prototype.waitForUnlock = function (priority) {
if (priority === void 0) { priority = 0; }
return this._semaphore.waitForUnlock(1, priority);
};
Mutex.prototype.release = function () {
if (this._semaphore.isLocked())
this._semaphore.release();
};
Mutex.prototype.cancel = function () {
return this._semaphore.cancel();
};
return Mutex;
}());
exports.default = Mutex;
+17
View File
@@ -0,0 +1,17 @@
interface MutexInterface {
acquire(priority?: number): Promise<MutexInterface.Releaser>;
runExclusive<T>(callback: MutexInterface.Worker<T>, priority?: number): Promise<T>;
waitForUnlock(priority?: number): Promise<void>;
isLocked(): boolean;
release(): void;
cancel(): void;
}
declare namespace MutexInterface {
interface Releaser {
(): void;
}
interface Worker<T> {
(): Promise<T> | T;
}
}
export default MutexInterface;
+2
View File
@@ -0,0 +1,2 @@
"use strict";
Object.defineProperty(exports, "__esModule", { value: true });
+22
View File
@@ -0,0 +1,22 @@
import SemaphoreInterface from './SemaphoreInterface';
declare class Semaphore implements SemaphoreInterface {
private _value;
private _cancelError;
constructor(_value: number, _cancelError?: Error);
acquire(weight?: number, priority?: number): Promise<[number, SemaphoreInterface.Releaser]>;
runExclusive<T>(callback: SemaphoreInterface.Worker<T>, weight?: number, priority?: number): Promise<T>;
waitForUnlock(weight?: number, priority?: number): Promise<void>;
isLocked(): boolean;
getValue(): number;
setValue(value: number): void;
release(weight?: number): void;
cancel(): void;
private _dispatchQueue;
private _dispatchItem;
private _newReleaser;
private _drainUnlockWaiters;
private _couldLockImmediately;
private _queue;
private _weightedWaiters;
}
export default Semaphore;
+155
View File
@@ -0,0 +1,155 @@
"use strict";
Object.defineProperty(exports, "__esModule", { value: true });
var tslib_1 = require("tslib");
var errors_1 = require("./errors");
var Semaphore = /** @class */ (function () {
function Semaphore(_value, _cancelError) {
if (_cancelError === void 0) { _cancelError = errors_1.E_CANCELED; }
this._value = _value;
this._cancelError = _cancelError;
this._queue = [];
this._weightedWaiters = [];
}
Semaphore.prototype.acquire = function (weight, priority) {
var _this = this;
if (weight === void 0) { weight = 1; }
if (priority === void 0) { priority = 0; }
if (weight <= 0)
throw new Error("invalid weight ".concat(weight, ": must be positive"));
return new Promise(function (resolve, reject) {
var task = { resolve: resolve, reject: reject, weight: weight, priority: priority };
var i = findIndexFromEnd(_this._queue, function (other) { return priority <= other.priority; });
if (i === -1 && weight <= _this._value) {
// Needs immediate dispatch, skip the queue
_this._dispatchItem(task);
}
else {
_this._queue.splice(i + 1, 0, task);
}
});
};
Semaphore.prototype.runExclusive = function (callback_1) {
return tslib_1.__awaiter(this, arguments, void 0, function (callback, weight, priority) {
var _a, value, release;
if (weight === void 0) { weight = 1; }
if (priority === void 0) { priority = 0; }
return tslib_1.__generator(this, function (_b) {
switch (_b.label) {
case 0: return [4 /*yield*/, this.acquire(weight, priority)];
case 1:
_a = _b.sent(), value = _a[0], release = _a[1];
_b.label = 2;
case 2:
_b.trys.push([2, , 4, 5]);
return [4 /*yield*/, callback(value)];
case 3: return [2 /*return*/, _b.sent()];
case 4:
release();
return [7 /*endfinally*/];
case 5: return [2 /*return*/];
}
});
});
};
Semaphore.prototype.waitForUnlock = function (weight, priority) {
var _this = this;
if (weight === void 0) { weight = 1; }
if (priority === void 0) { priority = 0; }
if (weight <= 0)
throw new Error("invalid weight ".concat(weight, ": must be positive"));
if (this._couldLockImmediately(weight, priority)) {
return Promise.resolve();
}
else {
return new Promise(function (resolve) {
if (!_this._weightedWaiters[weight - 1])
_this._weightedWaiters[weight - 1] = [];
insertSorted(_this._weightedWaiters[weight - 1], { resolve: resolve, priority: priority });
});
}
};
Semaphore.prototype.isLocked = function () {
return this._value <= 0;
};
Semaphore.prototype.getValue = function () {
return this._value;
};
Semaphore.prototype.setValue = function (value) {
this._value = value;
this._dispatchQueue();
};
Semaphore.prototype.release = function (weight) {
if (weight === void 0) { weight = 1; }
if (weight <= 0)
throw new Error("invalid weight ".concat(weight, ": must be positive"));
this._value += weight;
this._dispatchQueue();
};
Semaphore.prototype.cancel = function () {
var _this = this;
this._queue.forEach(function (entry) { return entry.reject(_this._cancelError); });
this._queue = [];
};
Semaphore.prototype._dispatchQueue = function () {
this._drainUnlockWaiters();
while (this._queue.length > 0 && this._queue[0].weight <= this._value) {
this._dispatchItem(this._queue.shift());
this._drainUnlockWaiters();
}
};
Semaphore.prototype._dispatchItem = function (item) {
var previousValue = this._value;
this._value -= item.weight;
item.resolve([previousValue, this._newReleaser(item.weight)]);
};
Semaphore.prototype._newReleaser = function (weight) {
var _this = this;
var called = false;
return function () {
if (called)
return;
called = true;
_this.release(weight);
};
};
Semaphore.prototype._drainUnlockWaiters = function () {
if (this._queue.length === 0) {
for (var weight = this._value; weight > 0; weight--) {
var waiters = this._weightedWaiters[weight - 1];
if (!waiters)
continue;
waiters.forEach(function (waiter) { return waiter.resolve(); });
this._weightedWaiters[weight - 1] = [];
}
}
else {
var queuedPriority_1 = this._queue[0].priority;
for (var weight = this._value; weight > 0; weight--) {
var waiters = this._weightedWaiters[weight - 1];
if (!waiters)
continue;
var i = waiters.findIndex(function (waiter) { return waiter.priority <= queuedPriority_1; });
(i === -1 ? waiters : waiters.splice(0, i))
.forEach((function (waiter) { return waiter.resolve(); }));
}
}
};
Semaphore.prototype._couldLockImmediately = function (weight, priority) {
return (this._queue.length === 0 || this._queue[0].priority < priority) &&
weight <= this._value;
};
return Semaphore;
}());
function insertSorted(a, v) {
var i = findIndexFromEnd(a, function (other) { return v.priority <= other.priority; });
a.splice(i + 1, 0, v);
}
function findIndexFromEnd(a, predicate) {
for (var i = a.length - 1; i >= 0; i--) {
if (predicate(a[i])) {
return i;
}
}
return -1;
}
exports.default = Semaphore;
+19
View File
@@ -0,0 +1,19 @@
interface SemaphoreInterface {
acquire(weight?: number, priority?: number): Promise<[number, SemaphoreInterface.Releaser]>;
runExclusive<T>(callback: SemaphoreInterface.Worker<T>, weight?: number, priority?: number): Promise<T>;
waitForUnlock(weight?: number, priority?: number): Promise<void>;
isLocked(): boolean;
getValue(): number;
setValue(value: number): void;
release(weight?: number): void;
cancel(): void;
}
declare namespace SemaphoreInterface {
interface Releaser {
(): void;
}
interface Worker<T> {
(value: number): Promise<T> | T;
}
}
export default SemaphoreInterface;
+2
View File
@@ -0,0 +1,2 @@
"use strict";
Object.defineProperty(exports, "__esModule", { value: true });
+3
View File
@@ -0,0 +1,3 @@
export declare const E_TIMEOUT: Error;
export declare const E_ALREADY_LOCKED: Error;
export declare const E_CANCELED: Error;
+6
View File
@@ -0,0 +1,6 @@
"use strict";
Object.defineProperty(exports, "__esModule", { value: true });
exports.E_CANCELED = exports.E_ALREADY_LOCKED = exports.E_TIMEOUT = void 0;
exports.E_TIMEOUT = new Error('timeout while waiting for mutex to become available');
exports.E_ALREADY_LOCKED = new Error('mutex already locked');
exports.E_CANCELED = new Error('request for lock canceled');
+7
View File
@@ -0,0 +1,7 @@
export { default as Mutex } from './Mutex';
export { default as MutexInterface } from './MutexInterface';
export { default as Semaphore } from './Semaphore';
export { default as SemaphoreInterface } from './SemaphoreInterface';
export { withTimeout } from './withTimeout';
export { tryAcquire } from './tryAcquire';
export * from './errors';
+13
View File
@@ -0,0 +1,13 @@
"use strict";
Object.defineProperty(exports, "__esModule", { value: true });
exports.tryAcquire = exports.withTimeout = exports.Semaphore = exports.Mutex = void 0;
var tslib_1 = require("tslib");
var Mutex_1 = require("./Mutex");
Object.defineProperty(exports, "Mutex", { enumerable: true, get: function () { return Mutex_1.default; } });
var Semaphore_1 = require("./Semaphore");
Object.defineProperty(exports, "Semaphore", { enumerable: true, get: function () { return Semaphore_1.default; } });
var withTimeout_1 = require("./withTimeout");
Object.defineProperty(exports, "withTimeout", { enumerable: true, get: function () { return withTimeout_1.withTimeout; } });
var tryAcquire_1 = require("./tryAcquire");
Object.defineProperty(exports, "tryAcquire", { enumerable: true, get: function () { return tryAcquire_1.tryAcquire; } });
tslib_1.__exportStar(require("./errors"), exports);
+4
View File
@@ -0,0 +1,4 @@
import MutexInterface from './MutexInterface';
import SemaphoreInterface from './SemaphoreInterface';
export declare function tryAcquire(mutex: MutexInterface, alreadyAcquiredError?: Error): MutexInterface;
export declare function tryAcquire(semaphore: SemaphoreInterface, alreadyAcquiredError?: Error): SemaphoreInterface;
+12
View File
@@ -0,0 +1,12 @@
"use strict";
Object.defineProperty(exports, "__esModule", { value: true });
exports.tryAcquire = void 0;
var errors_1 = require("./errors");
var withTimeout_1 = require("./withTimeout");
// eslint-disable-next-lisne @typescript-eslint/explicit-module-boundary-types
function tryAcquire(sync, alreadyAcquiredError) {
if (alreadyAcquiredError === void 0) { alreadyAcquiredError = errors_1.E_ALREADY_LOCKED; }
// eslint-disable-next-line @typescript-eslint/no-explicit-any
return (0, withTimeout_1.withTimeout)(sync, 0, alreadyAcquiredError);
}
exports.tryAcquire = tryAcquire;
+4
View File
@@ -0,0 +1,4 @@
import MutexInterface from './MutexInterface';
import SemaphoreInterface from './SemaphoreInterface';
export declare function withTimeout(mutex: MutexInterface, timeout: number, timeoutError?: Error): MutexInterface;
export declare function withTimeout(semaphore: SemaphoreInterface, timeout: number, timeoutError?: Error): SemaphoreInterface;
+128
View File
@@ -0,0 +1,128 @@
"use strict";
Object.defineProperty(exports, "__esModule", { value: true });
exports.withTimeout = void 0;
var tslib_1 = require("tslib");
/* eslint-disable @typescript-eslint/no-explicit-any */
var errors_1 = require("./errors");
function withTimeout(sync, timeout, timeoutError) {
var _this = this;
if (timeoutError === void 0) { timeoutError = errors_1.E_TIMEOUT; }
return {
acquire: function (weightOrPriority, priority) {
var weight;
if (isSemaphore(sync)) {
weight = weightOrPriority;
}
else {
weight = undefined;
priority = weightOrPriority;
}
if (weight !== undefined && weight <= 0) {
throw new Error("invalid weight ".concat(weight, ": must be positive"));
}
return new Promise(function (resolve, reject) { return tslib_1.__awaiter(_this, void 0, void 0, function () {
var isTimeout, handle, ticket, release, e_1;
return tslib_1.__generator(this, function (_a) {
switch (_a.label) {
case 0:
isTimeout = false;
handle = setTimeout(function () {
isTimeout = true;
reject(timeoutError);
}, timeout);
_a.label = 1;
case 1:
_a.trys.push([1, 3, , 4]);
return [4 /*yield*/, (isSemaphore(sync)
? sync.acquire(weight, priority)
: sync.acquire(priority))];
case 2:
ticket = _a.sent();
if (isTimeout) {
release = Array.isArray(ticket) ? ticket[1] : ticket;
release();
}
else {
clearTimeout(handle);
resolve(ticket);
}
return [3 /*break*/, 4];
case 3:
e_1 = _a.sent();
if (!isTimeout) {
clearTimeout(handle);
reject(e_1);
}
return [3 /*break*/, 4];
case 4: return [2 /*return*/];
}
});
}); });
},
runExclusive: function (callback, weight, priority) {
return tslib_1.__awaiter(this, void 0, void 0, function () {
var release, ticket;
return tslib_1.__generator(this, function (_a) {
switch (_a.label) {
case 0:
release = function () { return undefined; };
_a.label = 1;
case 1:
_a.trys.push([1, , 7, 8]);
return [4 /*yield*/, this.acquire(weight, priority)];
case 2:
ticket = _a.sent();
if (!Array.isArray(ticket)) return [3 /*break*/, 4];
release = ticket[1];
return [4 /*yield*/, callback(ticket[0])];
case 3: return [2 /*return*/, _a.sent()];
case 4:
release = ticket;
return [4 /*yield*/, callback()];
case 5: return [2 /*return*/, _a.sent()];
case 6: return [3 /*break*/, 8];
case 7:
release();
return [7 /*endfinally*/];
case 8: return [2 /*return*/];
}
});
});
},
release: function (weight) {
sync.release(weight);
},
cancel: function () {
return sync.cancel();
},
waitForUnlock: function (weightOrPriority, priority) {
var weight;
if (isSemaphore(sync)) {
weight = weightOrPriority;
}
else {
weight = undefined;
priority = weightOrPriority;
}
if (weight !== undefined && weight <= 0) {
throw new Error("invalid weight ".concat(weight, ": must be positive"));
}
return new Promise(function (resolve, reject) {
var handle = setTimeout(function () { return reject(timeoutError); }, timeout);
(isSemaphore(sync)
? sync.waitForUnlock(weight, priority)
: sync.waitForUnlock(priority)).then(function () {
clearTimeout(handle);
resolve();
});
});
},
isLocked: function () { return sync.isLocked(); },
getValue: function () { return sync.getValue(); },
setValue: function (value) { return sync.setValue(value); },
};
}
exports.withTimeout = withTimeout;
function isSemaphore(sync) {
return sync.getValue !== undefined;
}
+88
View File
@@ -0,0 +1,88 @@
{
"name": "async-mutex",
"version": "0.5.0",
"description": "A mutex for guarding async workflows",
"scripts": {
"lint": "eslint src/**/*.ts test/**/*.ts",
"build": "tsc && tsc -p tsconfig.es6.json && tsc -p tsconfig.mjs.json && rollup -o index.mjs mjs/index.js",
"prepublishOnly": "yarn test && yarn build",
"test": "yarn lint && nyc --reporter=text --reporter=html --reporter=lcov mocha test/*.ts",
"coveralls": "cat ./coverage/lcov.info | coveralls"
},
"author": "Christian Speckner <cnspeckn@googlemail.com> (https://github.com/DirtyHairy/)",
"license": "MIT",
"repository": {
"type": "git",
"url": "https://github.com/DirtyHairy/async-mutex"
},
"prettier": {
"printWidth": 120,
"tabWidth": 4,
"singleQuote": true,
"parser": "typescript"
},
"importSort": {
".js, .jsx, .ts, .tsx": {
"style": "eslint",
"parser": "typescript"
}
},
"eslintConfig": {
"root": true,
"parser": "@typescript-eslint/parser",
"plugins": [
"@typescript-eslint"
],
"extends": [
"eslint:recommended",
"plugin:@typescript-eslint/eslint-recommended",
"plugin:@typescript-eslint/recommended"
],
"rules": {
"eqeqeq": "error",
"@typescript-eslint/no-namespace": "off",
"no-async-promise-executor": "off"
}
},
"keywords": [
"mutex",
"async"
],
"files": [
"lib",
"es6",
"index.mjs"
],
"devDependencies": {
"@sinonjs/fake-timers": "^11.2.2",
"@types/mocha": "^10.0.6",
"@types/node": "^20.11.25",
"@types/sinonjs__fake-timers": "^8.1.2",
"@typescript-eslint/eslint-plugin": "^7.2.0",
"@typescript-eslint/parser": "^7.2.0",
"coveralls": "^3.1.1",
"eslint": "^8.57.0",
"import-sort-style-eslint": "^6.0.0",
"mocha": "^10.3.0",
"nyc": "^15.1.0",
"prettier": "^3.2.5",
"prettier-plugin-import-sort": "^0.0.7",
"rollup": "^4.12.1",
"ts-node": "^10.9.1",
"typescript": "^5.4.2"
},
"main": "lib/index.js",
"module": "es6/index.js",
"types": "lib/index.d.ts",
"exports": {
".": {
"import": "./index.mjs",
"require": "./lib/index.js",
"default": "./lib/index.js"
},
"./package.json": "./package.json"
},
"dependencies": {
"tslib": "^2.4.0"
}
}