Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
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
16 changes: 8 additions & 8 deletions src/sync/Mutex.ts
Original file line number Diff line number Diff line change
@@ -1,6 +1,6 @@
import type { Token } from "../types.js";
import { mutexIdentityData } from "./identifiers.js";
import { acquire, acquireAsync, SemaphoreInternal, type Releaser } from "./Semaphore.js";
import { acquireSync, acquire, SemaphoreInternal, type Releaser } from "./Semaphore.js";

/**
* Synchronization object that can be used to grant a single thread exclusive access to a resource or critical section.
Expand All @@ -19,7 +19,7 @@ export class Mutex extends SemaphoreInternal {
*
* @example
* ```typescript
* const releaser = await Mutex.acquireAsync(myMutexToken);
* const releaser = await Mutex.acquire(myMutexToken);
* try {
* ...
* }
Expand All @@ -30,7 +30,7 @@ export class Mutex extends SemaphoreInternal {
* @param token Mutex token to be acquired.
* @returns A releaser object that can and should be used for releasing the mutex.
*/
static acquire(token: Token): Releaser;
static acquireSync(token: Token): Releaser;
/**
* Acquires exclusivity from the specified mutex's token.
*
Expand All @@ -40,7 +40,7 @@ export class Mutex extends SemaphoreInternal {
*
* @example
* ```typescript
* const releaser = await Mutex.acquireAsync(myMutexToken);
* const releaser = await Mutex.acquire(myMutexToken);
* try {
* ...
* }
Expand All @@ -53,11 +53,11 @@ export class Mutex extends SemaphoreInternal {
* @returns A releaser object that can and should be used for releasing the mutex, or the value `'timed-out'` if
* the mutex could not be acquired before the specified timeout time elapsed.
*/
static acquire(token: Token, timeout: number): Releaser | "timed-out";
static acquireSync(token: Token, timeout: number): Releaser | "timed-out";
static acquireSync(token: Token, timeout?: number) {
return acquireSync(mutexIdentityData, token, timeout);
}
static acquire(token: Token, timeout?: number) {
return acquire(mutexIdentityData, token, timeout);
}
static acquireAsync(token: Token, timeout?: number) {
return acquireAsync(mutexIdentityData, token, timeout);
}
};
28 changes: 14 additions & 14 deletions src/sync/Semaphore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -46,7 +46,7 @@ function buildReleaser(token: Token) {
}) as Releaser;
}

export function acquire(identifierData: IdentifierData, token: Token, timeout?: number) {
export function acquireSync(identifierData: IdentifierData, token: Token, timeout?: number) {
checkToken(token, ...identifierData);
while (true) {
const available = Atomics.load(token, 0);
Expand All @@ -64,7 +64,7 @@ export function acquire(identifierData: IdentifierData, token: Token, timeout?:
}
}

export async function acquireAsync(identifierData: IdentifierData, token: Token, timeout?: number) {
export async function acquire(identifierData: IdentifierData, token: Token, timeout?: number) {
checkToken(token, ...identifierData);
while (true) {
const available = Atomics.load(token, 0);
Expand Down Expand Up @@ -101,7 +101,7 @@ export class Semaphore extends SemaphoreInternal {
*
* @example
* ```typescript
* const releaser = Semaphore.acquire(mySemaphoreToken);
* const releaser = Semaphore.acquireSync(mySemaphoreToken);
* try {
* ...
* }
Expand All @@ -112,7 +112,7 @@ export class Semaphore extends SemaphoreInternal {
* @param token Semaphore token to be acquired.
* @returns A releaser object that can and should be used for releasing the semaphore.
*/
static acquire(token: Token): Releaser;
static acquireSync(token: Token): Releaser;
/**
* Acquires a slot from the specified semaphore's token.
*
Expand All @@ -122,7 +122,7 @@ export class Semaphore extends SemaphoreInternal {
*
* @example
* ```typescript
* const releaser = Semaphore.acquire(mySemaphoreToken);
* const releaser = Semaphore.acquireSync(mySemaphoreToken);
* try {
* ...
* }
Expand All @@ -135,9 +135,9 @@ export class Semaphore extends SemaphoreInternal {
* @returns A releaser object that can and should be used for releasing the semaphore, or the value `'timed-out'`
* if the sempahore could not be acquired before the specified timeout time elapsed.
*/
static acquire(token: Token, timeout: number): 'timed-out' | Releaser;
static acquire(token: Token, timeout?: number) {
return acquire(semaphoreIdentityData, token, timeout);
static acquireSync(token: Token, timeout: number): 'timed-out' | Releaser;
static acquireSync(token: Token, timeout?: number) {
return acquireSync(semaphoreIdentityData, token, timeout);
}

/**
Expand All @@ -149,7 +149,7 @@ export class Semaphore extends SemaphoreInternal {
*
* @example
* ```typescript
* const releaser = await Semaphore.acquireAsync(mySemaphoreToken);
* const releaser = await Semaphore.acquire(mySemaphoreToken);
* try {
* ...
* }
Expand All @@ -160,7 +160,7 @@ export class Semaphore extends SemaphoreInternal {
* @param token Semaphore token to be acquired.
* @returns A releaser object that can and should be used for releasing the semaphore.
*/
static acquireAsync(token: Token): Promise<Releaser>;
static acquire(token: Token): Promise<Releaser>;
/**
* Asynchronously acquires a slot from the specified semaphore's token.
*
Expand All @@ -170,7 +170,7 @@ export class Semaphore extends SemaphoreInternal {
*
* @example
* ```typescript
* const releaser = await Semaphore.acquireAsync(mySemaphoreToken);
* const releaser = await Semaphore.acquire(mySemaphoreToken);
* try {
* ...
* }
Expand All @@ -183,8 +183,8 @@ export class Semaphore extends SemaphoreInternal {
* @returns A releaser object that can and should be used for releasing the semaphore, or the value `'timed-out'`
* if the sempahore could not be acquired before the specified timeout time elapsed.
*/
static acquireAsync(token: Token, timeout: number): Promise<"timed-out" | Releaser>;
static acquireAsync(token: Token, timeout?: number) {
return acquireAsync(semaphoreIdentityData, token, timeout);
static acquire(token: Token, timeout: number): Promise<"timed-out" | Releaser>;
static acquire(token: Token, timeout?: number) {
return acquire(semaphoreIdentityData, token, timeout);
}
}
4 changes: 2 additions & 2 deletions tests/ut/helpers/test-acquire-worker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@ self.onmessage = function (event) {
try {
self.postMessage('running');
const releaser = timeout !== undefined
? (source === "Mutex" ? Mutex.acquire(token, timeout) : Semaphore.acquire(token, timeout))
: (source === "Mutex" ? Mutex.acquire(token) : Semaphore.acquire(token));
? (source === "Mutex" ? Mutex.acquireSync(token, timeout) : Semaphore.acquireSync(token, timeout))
: (source === "Mutex" ? Mutex.acquireSync(token) : Semaphore.acquireSync(token));
self.postMessage({ success: typeof releaser === 'function' });
} catch (error) {
const errorMessage = error instanceof Error ? error.message : String(error);
Expand Down
26 changes: 13 additions & 13 deletions tests/ut/sync/Mutex.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -45,19 +45,19 @@ describe('Mutex', () => {
});
});

describe('static acquire', () => {
describe('static acquireSync', () => {
beforeEach(() => {
mutex = new Mutex();
});

it('Should acquire immediately when mutex is available.', () => {
const releaser = Mutex.acquire(mutex.token);
const releaser = Mutex.acquireSync(mutex.token);

expect(typeof releaser).toBe('function');
});

it('Should wait when mutex is not available.', async () => {
const releaser = Mutex.acquire(mutex.token);
const releaser = Mutex.acquireSync(mutex.token);
const waitAcquire = (await testMutexAcquireInWorker(mutex.token)).wait;
let mainThreadReleased = false;
setTimeout(() => {
Expand All @@ -71,25 +71,25 @@ describe('Mutex', () => {
});
});

describe('static acquireAsync', () => {
describe('static acquire', () => {
beforeEach(() => {
mutex = new Mutex();
});

it('Should acquire immediately when mutex is available.', async () => {
const releaser = await Mutex.acquireAsync(mutex.token);
const releaser = await Mutex.acquire(mutex.token);

expect(typeof releaser).toBe('function');
});

it('Should wait asynchronously when mutex is not available.', async () => {
const releaser = Mutex.acquire(mutex.token);
const releaser = Mutex.acquireSync(mutex.token);
let mainThreadReleased = false;
setTimeout(() => {
mainThreadReleased = true;
releaser();
}, 0);
await Mutex.acquireAsync(mutex.token);
await Mutex.acquire(mutex.token);

expect(mainThreadReleased).toBe(true);
});
Expand All @@ -101,7 +101,7 @@ describe('Mutex', () => {
});

it('Should release the mutex when called.', () => {
const releaser = Mutex.acquire(mutex.token);
const releaser = Mutex.acquireSync(mutex.token);

expect(typeof releaser).toBe('function');

Expand All @@ -111,7 +111,7 @@ describe('Mutex', () => {
});

it('Should throw error when called twice.', () => {
const releaser = Mutex.acquire(mutex.token);
const releaser = Mutex.acquireSync(mutex.token);

releaser(); // First release

Expand All @@ -125,22 +125,22 @@ describe('Mutex', () => {
});

it('Should ensure the mutex cannot be acquired again asynchronously.', async () => {
const releaser1 = Mutex.acquire(mutex.token);
const releaser1 = Mutex.acquireSync(mutex.token);
expect(typeof releaser1).toBe('function');
let releaser2 = await Mutex.acquireAsync(mutex.token, 0);
let releaser2 = await Mutex.acquire(mutex.token, 0);
expect(releaser2).toBe('timed-out');
let mainThreadReleased = false;
setTimeout(() => {
mainThreadReleased = true;
releaser1();
}, 0);
releaser2 = await Mutex.acquireAsync(mutex.token);
releaser2 = await Mutex.acquire(mutex.token);
expect(mainThreadReleased).toBe(true);
expect(typeof releaser2).toBe('function');
});

it("Should ensure the mutex cannot be acquired from a different thread while it's held.", async () => {
const releaser1 = Mutex.acquire(mutex.token);
const releaser1 = Mutex.acquireSync(mutex.token);
expect(typeof releaser1).toBe('function');
const waitAcquire = (await testMutexAcquireInWorker(mutex.token, 0)).wait;
const result = await waitAcquire;
Expand Down
24 changes: 12 additions & 12 deletions tests/ut/sync/Semaphore.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,7 +23,7 @@ describe('Semaphore', () => {

it('Should create disabled semaphore when requested.', () => {
semaphore = new Semaphore(3, true);
expect(Semaphore.acquire(semaphore.token, 0)).toBe('timed-out');
expect(Semaphore.acquireSync(semaphore.token, 0)).toBe('timed-out');
});
});

Expand All @@ -50,32 +50,32 @@ describe('Semaphore', () => {
});
});

describe('static acquire', () => {
describe('static acquireSync', () => {
const initialCapacity = 2;
beforeEach(() => {
semaphore = new Semaphore(initialCapacity);
});

it('Should acquire immediately when capacity is available.', () => {
const result = Semaphore.acquire(semaphore.token);
const result = Semaphore.acquireSync(semaphore.token);

expect(result).not.toBe('timed-out');
expect(typeof result).toBe('function');
});

it('Should return "timed-out" when timeout occurs.', () => {
for (let i = 0; i < initialCapacity; i++) {
Semaphore.acquire(semaphore.token);
Semaphore.acquireSync(semaphore.token);
}
const result = Semaphore.acquire(semaphore.token, 0);
const result = Semaphore.acquireSync(semaphore.token, 0);

expect(result).toBe('timed-out');
});

it('Should wait and acquire when capacity becomes available.', async () => {
let releaser: Function;
for (let i = 0; i < initialCapacity; i++) {
releaser = Semaphore.acquire(semaphore.token);
releaser = Semaphore.acquireSync(semaphore.token);
}
const waitAcquire = (await testSemaphoreAcquireInWorker(semaphore.token)).wait;
releaser!();
Expand All @@ -85,23 +85,23 @@ describe('Semaphore', () => {
});
});

describe('static acquireAsync', () => {
describe('static acquire', () => {
const initialCapacity = 2;
beforeEach(() => {
semaphore = new Semaphore(initialCapacity);
});

it('Should acquire immediately when capacity is available.', async () => {
const result = await Semaphore.acquireAsync(semaphore.token);
const result = await Semaphore.acquire(semaphore.token);

expect(typeof result).toBe('function');
});

it('Should return "timed-out" when timeout occurs.', async () => {
for (let i = 0; i < initialCapacity; i++) {
Semaphore.acquire(semaphore.token);
Semaphore.acquireSync(semaphore.token);
}
const result = await Semaphore.acquireAsync(semaphore.token, 0);
const result = await Semaphore.acquire(semaphore.token, 0);

expect(result).toBe('timed-out');
});
Expand All @@ -113,7 +113,7 @@ describe('Semaphore', () => {
});

it('Should release the semaphore when called.', () => {
const releaser = Semaphore.acquire(semaphore.token);
const releaser = Semaphore.acquireSync(semaphore.token);

expect(typeof releaser).toBe('function');
if (typeof releaser === 'function') {
Expand All @@ -123,7 +123,7 @@ describe('Semaphore', () => {
});

it('Should throw error when called twice.', () => {
const releaser = Semaphore.acquire(semaphore.token);
const releaser = Semaphore.acquireSync(semaphore.token);

if (typeof releaser === 'function') {
releaser(); // First release
Expand Down
Loading