dmutex
v0.3.0
Published
A small TypeScript distributed mutex and semaphore library for MongoDB and Redis.
Maintainers
Readme
dmutex
A small TypeScript distributed mutex and semaphore library that can use MongoDB or Redis as its backend.
DMutex.acquire() allows only one caller to hold a lock for a given key at a time. Each lock stores an ownership token, so a stale worker cannot release a lock that was later acquired by another worker. DMutex is implemented as a single-permit semaphore, and applications can use the same interface while choosing either MongoDB or Redis as the implementation.
DSemaphore.acquire() allows up to maxPermits callers to hold permits for a given key at the same time. Each permit also carries an ownership token and TTL.
Installation
Bun:
bun add dmutexNode.js with npm:
npm install dmutexNode.js with pnpm:
pnpm add dmutexNode.js with Yarn:
yarn add dmutexdmutex does not force a specific MongoDB or Redis client package as a runtime dependency or peer dependency. Pass in the database client your application already uses.
If you use the official mongodb, redis, or ioredis packages, install the versions you want in your application.
Bun:
bun add mongodb
bun add redisNode.js with npm:
npm install mongodb redisNode.js with pnpm:
pnpm add mongodb redisNode.js with Yarn:
yarn add mongodb redisRedis compatibility is currently pinned with real package tests for:
redis/@redis/client: uses thesendCommand(args)pathioredis: uses theset(...args)/eval(...args)path
Usage
MongoDB
import { MongoClient } from "mongodb";
import { DMutex } from "dmutex";
const mongoClient = new MongoClient("mongodb://localhost:27017");
await mongoClient.connect();
const dmutex = new DMutex("my-service", mongoClient);
await dmutex.ready();
const result = await dmutex.run("job:daily-report", async () => {
// Run protected work.
return "done";
}, 60);
if (result === null) {
// Another process already holds this lock.
process.exit(0);
}
await mongoClient.close();Redis
import { createClient } from "redis";
import { DMutex } from "dmutex";
const redisClient = createClient({ url: "redis://localhost:6379" });
await redisClient.connect();
const dmutex = new DMutex("my-service", redisClient);
const result = await dmutex.run("job:daily-report", async () => {
// Run protected work.
return "done";
}, 60);
if (result === null) {
// Another process already holds this lock.
process.exit(0);
}
await redisClient.close();Semaphore
import { createClient } from "redis";
import { DSemaphore } from "dmutex";
const redisClient = createClient({ url: "redis://localhost:6379" });
await redisClient.connect();
const semaphore = new DSemaphore("my-service", redisClient, {
maxPermits: 3,
});
const result = await semaphore.run("api:partner", async (permit) => {
await permit.extend(120);
return "done";
}, 120);
if (result === null) {
// All permits are currently held.
}
await redisClient.close();API
new DMutex(serviceName, client, options?)
Creates a mutex instance for a service.
serviceName: service identifier used for backend-specific namespacingclient: MongoDB or Redis client.dmutexdetects the backend from the injected client shape.options.defaultTtlSeconds: default lock TTL. Defaults to 300 seconds.options.backend: optional explicit backend override, eithermongodborredis. Use this when a wrapped client matches more than one backend contract.
MongoDB options:
options.dbName: database name. Defaults todmutex.options.collectionName: collection name. If set, this takes precedence overcollectionPrefixandserviceName.options.collectionPrefix: collection prefix. Defaults to_dmutex_.
Redis options:
options.keyPrefix: Redis key prefix. Defaults to_dmutex_${serviceName}:.
MongoDB uses the _dmutex_${serviceName} collection in the dmutex database by default. Redis uses keys under the _dmutex_${serviceName}: prefix by default. Backend keys include internal permit-slot names.
ready()
await dmutex.ready();Waits for backend initialization. For MongoDB, this waits for the TTL index to be created. For Redis, this is a no-op. acquire(), lock(), unlock(), and extend() also wait for any required initialization internally, but calling ready() during application startup surfaces MongoDB initialization failures earlier.
run(key, callback, ttl?)
const result = await dmutex.run("some-key", async (lock) => {
await lock.extend(300);
return "done";
}, 300);
if (result === null) {
// Another process already holds this lock.
}Attempts to acquire a lock, runs the callback while the lock is held, and releases the lock in a finally block.
key: lock identifiercallback: function to run while holding the lock. It receives the acquiredDMutexLock.ttl: lock TTL in seconds. Defaults to 300 seconds.- returns: the callback result when the lock is acquired, or
nullwhen another holder already owns the key
If the callback throws, run() releases the lock and rethrows the callback error. run() does not automatically renew long-running locks; use the callback's lock.extend(ttl) when the protected work may run longer than the TTL.
runWithRetry(key, callback, options?)
const result = await dmutex.runWithRetry("some-key", async (lock) => {
return "done";
}, {
ttl: 300,
timeoutMs: 10_000,
retryDelayMs: 100,
});
if (result === null) {
// The lock was not acquired before timeoutMs elapsed.
}Attempts to acquire a lock until it succeeds or timeoutMs elapses, then runs the callback and releases the lock in a finally block.
options.ttl: lock TTL in seconds. Defaults to 300 seconds.options.timeoutMs: maximum time to wait. Defaults to 30,000 milliseconds.options.retryDelayMs: delay between attempts. Defaults to 100 milliseconds.- returns: the callback result when the lock is acquired, or
nullwhen the timeout elapses
acquire(key, ttl?)
const lock = await dmutex.acquire("some-key", 300);
if (lock) {
try {
// protected work
} finally {
await lock.release();
}
}Attempts to acquire a lock for the given key.
key: lock identifierttl: lock TTL in seconds. Defaults to 300 seconds.- returns:
DMutexLockwhen the lock is acquired, ornullwhen another holder already owns the key
acquireWithRetry(key, options?)
const lock = await dmutex.acquireWithRetry("some-key", {
ttl: 300,
timeoutMs: 10_000,
retryDelayMs: 100,
});
if (!lock) {
// The lock was not acquired before timeoutMs elapsed.
}Attempts to acquire a lock until it succeeds or timeoutMs elapses.
key: lock identifieroptions.ttl: lock TTL in seconds. Defaults to 300 seconds.options.timeoutMs: maximum time to wait. Defaults to 30,000 milliseconds.options.retryDelayMs: delay between attempts. Defaults to 100 milliseconds.- returns:
DMutexLockwhen the lock is acquired, ornullwhen the timeout elapses
DMutexLock contains:
key: lock keytoken: ownership tokenexpiredAt: current lock expiration time, updated after a successfullock.extend()release(): releases only the lock with the matching ownership tokenextend(ttl?): extends only the active lock with the matching ownership token
lock(key, ttl?)
Deprecated: use acquire() instead. acquire() returns a lock handle that carries its ownership token and is safer across async boundaries.
const acquired = await dmutex.lock("some-key", 300);Attempts to acquire a lock for the given key.
key: lock identifierttl: lock TTL in seconds. Defaults to 300 seconds.- returns:
truewhen the lock is acquired, orfalsewhen another holder already owns the key
This is the legacy boolean-style API. New code should prefer acquire(), which exposes ownership explicitly.
unlock(key, token?)
Deprecated: prefer lock.release() from the lock handle returned by acquire(). unlock(key) depends on token state stored in the same DMutex instance unless a token is provided.
await dmutex.unlock("some-key");Deletes the lock for the given key. If token is provided, only a lock with the matching token is released. Locks acquired through lock() can be released with unlock(key) from the same DMutex instance because the instance keeps the internal token.
extend(key, token, ttl?)
await dmutex.extend("some-key", lock.token, 300);Extends the TTL for an active lock with the matching token. Returns true on success, or false when the token does not match or the lock is already expired.
Semaphore API
new DSemaphore(serviceName, client, options)
Creates a semaphore instance for a service.
serviceName: service identifier used for backend-specific namespacingclient: MongoDB or Redis client. Backend detection is the same asDMutex.options.maxPermits: maximum concurrent permits per key. Must be a positive integer.options.defaultTtlSeconds,options.backend, MongoDB options, and Redis options are the same asDMutex.
MongoDB uses the _dsemaphore_${serviceName} collection by default. Redis uses _dsemaphore_${serviceName}: as the default key prefix. Backend keys include internal permit-slot names. Explicit collectionName, collectionPrefix, and keyPrefix options override these defaults.
semaphore.acquire(key, ttl?)
const permit = await semaphore.acquire("some-key", 300);
if (permit) {
try {
// limited-concurrency work
} finally {
await permit.release();
}
}Attempts to acquire one permit for the given key.
key: semaphore identifierttl: permit TTL in seconds. Defaults to 300 seconds.- returns:
DSemaphorePermitwhen a permit is acquired, ornullwhen all permits are held
DSemaphorePermit contains:
key: original semaphore keytoken: ownership tokenslot: internal permit slot numberexpiredAt: current permit expiration time, updated after a successfulpermit.extend()release(): releases only this permitextend(ttl?): extends only this active permit
semaphore.run(key, callback, ttl?)
Attempts to acquire one permit, runs the callback while the permit is held, and releases the permit in a finally block. It returns the callback result, or null when all permits are held.
semaphore.acquireWithRetry(key, options?)
Attempts to acquire one permit until it succeeds or timeoutMs elapses. The options are the same as DMutex.acquireWithRetry().
semaphore.runWithRetry(key, callback, options?)
Attempts to acquire one permit with retry, runs the callback, and releases the permit in a finally block.
semaphore.release(key, token)
Releases the active permit with the matching token for the given key. Returns true on success, or false when no active permit matches.
semaphore.extend(key, token, ttl?)
Extends the TTL for an active permit with the matching token. Returns true on success, or false when no active permit matches.
Backend Behavior and Caveats
MongoDB:
- Permit-slot acquisition uses MongoDB
insertOne(). Only one document can exist for a given internal slot_id, so only one concurrent caller can win each slot. - TTL is handled with the
expiredAtfield and a MongoDB TTL index. MongoDB's TTL monitor runs periodically, so expired locks are not deleted exactly at their expiration time. - If an expired lock document has not yet been removed by the TTL monitor, a new acquisition attempt checks expiration and atomically attempts takeover.
- Duplicate key conflicts are treated as normal lock contention. Connection, authorization, and other MongoDB errors are thrown to the caller.
Redis:
- Permit-slot acquisition uses
SET key token PX ttl NX. - Release and extension use Lua
EVALscripts to verify the token and runDELorPEXPIREatomically. DMutexLock.expiredAtandDSemaphorePermit.expiredAtare calculated with the client clock. Actual expiration is enforced by Redis TTL.
Common:
- Release and extension verify the ownership token. The safest usage is to call
release()andextend()on the lock handle returned byacquire(). DMutexis backed byDSemaphorewithmaxPermits: 1.DSemaphoreis implemented as a fixed set of token-protected internal permit slots. Use the samemaxPermitsfor all callers coordinating on the same semaphore key.- The package does not import
mongodb,redis, orioredisat runtime. Client implementations are injected by the application. - Backend detection is structural. MongoDB clients must expose
db(). Redis clients must expose eithersendCommand(args)or bothset(...args)andeval(...args). - A client that matches multiple backend contracts is rejected because backend selection would be ambiguous. Pass
options.backendto choose explicitly. - MongoDB clients must provide
db,collection,createIndex,insertOne,updateOne, anddeleteOne. - Redis clients must provide either
sendCommand(args)orset(...args)pluseval(...args).
Development
Install dependencies:
bun installProject layout:
src/ Runtime library source
tests/unit/ Fast tests that do not require external services
tests/integration/ MongoDB and Redis integration tests
docs/ Project planning and maintenance notesBuild:
bun run buildRun the default test suite:
bun run testThis runs the fast unit suite only and does not require MongoDB or Redis.
Running bun test directly is also safe: integration tests are skipped unless DMUTEX_INTEGRATION=1 is set.
Run all unit tests explicitly:
bun run test:unitRun only Redis adapter unit tests:
bun run test:redis:unitRun only real redis and ioredis client integration tests:
bun run test:redis:integrationRun only MongoDB integration tests:
bun run test:mongodbRun all integration tests:
bun run test:integrationThe integration suite requires MongoDB and Redis. The default MongoDB URL is mongodb://localhost:27017, and the default Redis URL is redis://localhost:6379. Set MONGODB_URL and REDIS_URL to use different endpoints.
MONGODB_URL=mongodb://localhost:27017 REDIS_URL=redis://localhost:6379 bun run test:integrationIntegration Tests with Docker Compose
Run the integration suite with Docker Compose-managed MongoDB and Redis:
bun run test:integration:dockerThis starts the services, waits for their healthchecks, runs the integration tests, and stops the services when the test command exits.
To manage services manually, start MongoDB and Redis:
docker compose up -dRun the integration suite against those services:
MONGODB_URL=mongodb://localhost:27017 REDIS_URL=redis://localhost:6379 bun run test:integrationStop the services:
docker compose down