fix(server): distinguish local mutex correctly (#9444)
This commit is contained in:
@@ -11,11 +11,12 @@ import { Lock } from './lock';
|
|||||||
const lockScript = `local key = KEYS[1]
|
const lockScript = `local key = KEYS[1]
|
||||||
local owner = ARGV[1]
|
local owner = ARGV[1]
|
||||||
|
|
||||||
-- if lock is not exists or lock is owned by the owner
|
-- if lock is not exists then set lock to the owner and return 1, otherwise return 0
|
||||||
-- then set lock to the owner and return 1, otherwise return 0
|
|
||||||
-- if the lock is not released correctly due to unexpected reasons
|
-- if the lock is not released correctly due to unexpected reasons
|
||||||
-- lock will be released after 60 seconds
|
-- lock will be released after 60 seconds
|
||||||
if redis.call("get", key) == owner or redis.call("set", key, owner, "NX", "EX", 60) then
|
if redis.call("get", key) == owner then
|
||||||
|
return 0
|
||||||
|
elseif redis.call("set", key, owner, "NX", "EX", 60) then
|
||||||
return 1
|
return 1
|
||||||
else
|
else
|
||||||
return 0
|
return 0
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ import { randomUUID } from 'node:crypto';
|
|||||||
import { Inject, Injectable, Logger, Scope } from '@nestjs/common';
|
import { Inject, Injectable, Logger, Scope } from '@nestjs/common';
|
||||||
import { ModuleRef, REQUEST } from '@nestjs/core';
|
import { ModuleRef, REQUEST } from '@nestjs/core';
|
||||||
import type { Request } from 'express';
|
import type { Request } from 'express';
|
||||||
|
import { nanoid } from 'nanoid';
|
||||||
|
|
||||||
import { GraphqlContext } from '../graphql';
|
import { GraphqlContext } from '../graphql';
|
||||||
import { retryable } from '../utils/promise';
|
import { retryable } from '../utils/promise';
|
||||||
@@ -14,7 +15,7 @@ export const MUTEX_WAIT = 100;
|
|||||||
@Injectable()
|
@Injectable()
|
||||||
export class Mutex {
|
export class Mutex {
|
||||||
protected logger = new Logger(Mutex.name);
|
protected logger = new Logger(Mutex.name);
|
||||||
private readonly clusterIdentifier = `cluster:${randomUUID()}`;
|
private readonly clusterIdentifier = `cluster:${nanoid()}`;
|
||||||
|
|
||||||
constructor(protected readonly locker: Locker) {}
|
constructor(protected readonly locker: Locker) {}
|
||||||
|
|
||||||
@@ -39,7 +40,10 @@ export class Mutex {
|
|||||||
* @param key resource key
|
* @param key resource key
|
||||||
* @returns LockGuard
|
* @returns LockGuard
|
||||||
*/
|
*/
|
||||||
async acquire(key: string, owner: string = this.clusterIdentifier) {
|
async acquire(
|
||||||
|
key: string,
|
||||||
|
owner: string = `${this.clusterIdentifier}:${nanoid()}`
|
||||||
|
) {
|
||||||
try {
|
try {
|
||||||
return await retryable(
|
return await retryable(
|
||||||
() => this.locker.lock(owner, key),
|
() => this.locker.lock(owner, key),
|
||||||
|
|||||||
81
packages/backend/server/tests/mutex.spec.ts
Normal file
81
packages/backend/server/tests/mutex.spec.ts
Normal file
@@ -0,0 +1,81 @@
|
|||||||
|
import { randomUUID } from 'node:crypto';
|
||||||
|
|
||||||
|
import { TestingModule } from '@nestjs/testing';
|
||||||
|
import ava, { TestFn } from 'ava';
|
||||||
|
import Sinon from 'sinon';
|
||||||
|
|
||||||
|
import { Locker, Mutex } from '../src/base/mutex';
|
||||||
|
import { SessionRedis } from '../src/base/redis';
|
||||||
|
import { createTestingModule, sleep } from './utils';
|
||||||
|
|
||||||
|
const test = ava as TestFn<{
|
||||||
|
module: TestingModule;
|
||||||
|
mutex: Mutex;
|
||||||
|
locker: Locker;
|
||||||
|
session: SessionRedis;
|
||||||
|
}>;
|
||||||
|
|
||||||
|
test.beforeEach(async t => {
|
||||||
|
const module = await createTestingModule();
|
||||||
|
|
||||||
|
t.context.module = module;
|
||||||
|
t.context.mutex = module.get(Mutex);
|
||||||
|
t.context.locker = module.get(Locker);
|
||||||
|
t.context.session = module.get(SessionRedis);
|
||||||
|
});
|
||||||
|
|
||||||
|
test.afterEach(async t => {
|
||||||
|
await t.context.module.close();
|
||||||
|
});
|
||||||
|
|
||||||
|
const lockerPrefix = randomUUID();
|
||||||
|
test('should be able to acquire lock', async t => {
|
||||||
|
const { mutex } = t.context;
|
||||||
|
|
||||||
|
{
|
||||||
|
t.truthy(
|
||||||
|
await mutex.acquire(`${lockerPrefix}1`),
|
||||||
|
'should be able to acquire lock'
|
||||||
|
);
|
||||||
|
t.falsy(
|
||||||
|
await mutex.acquire(`${lockerPrefix}1`),
|
||||||
|
'should not be able to acquire lock again'
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
{
|
||||||
|
const lock1 = await mutex.acquire(`${lockerPrefix}2`);
|
||||||
|
t.truthy(lock1);
|
||||||
|
await lock1?.release();
|
||||||
|
const lock2 = await mutex.acquire(`${lockerPrefix}2`);
|
||||||
|
t.truthy(lock2);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
test('should be able to acquire lock parallel', async t => {
|
||||||
|
const { mutex, locker } = t.context;
|
||||||
|
const spyedLocker = Sinon.spy(locker, 'lock');
|
||||||
|
const requestLock = async (key: string) => {
|
||||||
|
const lock = mutex.acquire(key);
|
||||||
|
await using _lock = await lock;
|
||||||
|
const lastCall = spyedLocker.lastCall.returnValue;
|
||||||
|
try {
|
||||||
|
// in rare cases, the lock can be acquired
|
||||||
|
// in which case skip the error message check
|
||||||
|
await lastCall;
|
||||||
|
} catch {
|
||||||
|
await t.throwsAsync(lastCall, {
|
||||||
|
message: `Failed to acquire lock for resource [${key}]`,
|
||||||
|
});
|
||||||
|
}
|
||||||
|
|
||||||
|
await sleep(100);
|
||||||
|
};
|
||||||
|
|
||||||
|
await t.notThrowsAsync(
|
||||||
|
Promise.all(
|
||||||
|
Array.from({ length: 10 }, _ => requestLock(`${lockerPrefix}3`))
|
||||||
|
),
|
||||||
|
'should be able to acquire lock parallel'
|
||||||
|
);
|
||||||
|
});
|
||||||
Reference in New Issue
Block a user