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
32 changes: 32 additions & 0 deletions src/sync/AutoResetEvent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,38 @@ export class AutoResetEvent extends Event {
constructor() {
super(autoResetEventIdentityData[0], 1);
}
/**
* Checks whether or not this auto-resettable event is in its signaled state.
*
* Auto-reset events automatically reset when a thread is freed by it, so expect this function to only return
* `true` if the token has been signaled and there were no threads blocked by it.
* @returns `true` if the token is signaled, or `false` otherwise.
*/
isSignaled() {
return AutoResetEvent.isSignaled(this.token);
}
/**
* Waits on this auto-resettable event to be signaled.
*
* Use this method to block a (worker's) thread for whatever reason.
* @param timeout Maximum time to wait. Don't specify a value to wait indefinitely.
* @returns `'ok'` when the waiting is over because the event signaled while waiting on it, `'timed-out'` when the
* specified timeout elapsed and the event did not signal, or `'not-equal'` if no wait took place.
*/
waitSync(timeout?: number) {
return AutoResetEvent.waitSync(this.token, timeout);
}
/**
* Asynchronously waits on this auto-resettable event to be signaled.
*
* Use this method to stop the current work and release the worker thread (to pick up on new messages, perhaps).
* @param timeout Maximum time to wait. Don't specify a value to wait indefinitely.
* @returns `'ok'` when the waiting is over because the event signaled while waiting on it, `'timed-out'` when the
* specified timeout elapsed and the event did not signal, or `'not-equal'` if no wait took place.
*/
wait(timeout?: number) {
return AutoResetEvent.wait(this.token, timeout);
}
/**
* Checks whether or not an auto-resettable event's token is in its signaled state.
*
Expand Down
27 changes: 27 additions & 0 deletions src/sync/ManualResetEvent.ts
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,33 @@ export class ManualResetEvent extends Event {
reset() {
Atomics.store(super.token, 0, 0);
}
/**
* Checks whether or not this manually-resettable event is in its signaled state.
* @returns `true` if the event is signaled, or `false` otherwise.
*/
isSignaled() {
return ManualResetEvent.isSignaled(this.token);
}
/**
* Waits on the specified manually-resettable event to be signaled.
* @param timeout Maximum time to wait. Don't specify a value to wait indefinitely.
* @returns `'ok'` when the waiting is over because the token signaled while waiting on it, `'timed-out'` when the
* specified timeout elapsed and the token did not signal, or `'not-equal'` if no wait took place.
*/
waitSync(timeout?: number) {
return ManualResetEvent.waitSync(this.token, timeout);
}
/**
* Asynchronously waits on the specified manually-resettable event to be signaled.
*
* Use this method to stop the current work and release the thread.
* @param timeout Maximum time to wait. Don't specify a value to wait indefinitely.
* @returns `'ok'` when the waiting is over because the token signaled while waiting on it, `'timed-out'` when the
* specified timeout elapsed and the token did not signal, or `'not-equal'` if no wait took place.
*/
wait(timeout?: number) {
return ManualResetEvent.wait(this.token, timeout);
}
/**
* Checks whether or not a manually-resettable event's token is in its signaled state.
*
Expand Down
69 changes: 58 additions & 11 deletions src/sync/Mutex.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,19 +7,18 @@ import { acquireSync, acquire, SemaphoreInternal, type Releaser } from "./Semaph
*/
export class Mutex extends SemaphoreInternal {
constructor(createDisabled: boolean = false) {
super(mutexIdentityData[0], 1, createDisabled);
super(mutexIdentityData, 1, createDisabled);
}

/**
* Acquires exclusivity from the specified mutex's token.
* Acquires exclusivity over the specified mutex's token.
*
* **IMPORTANT**: Acquiring a mutex halts other work. Releasing the mutex by invoking the releaser function that
* this method returns is *imperative*. Use a `try..finally` construct to make sure that the releaser function is
* always executed, even in the event of an error.
*
* @example
* ```typescript
* const releaser = await Mutex.acquire(myMutexToken);
* const releaser = Mutex.acquireSync(myMutexToken);
* try {
* ...
* }
Expand All @@ -28,35 +27,83 @@ 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.
* @returns A releaser function that can and should be used for releasing the mutex.
*/
static acquireSync(token: Token): Releaser;
/**
* Acquires exclusivity from the specified mutex's token.
* Acquires exclusivity over the specified mutex's token.
*
* **IMPORTANT**: Acquiring a mutex halts other work. Releasing the mutex by invoking the releaser function that
* this method returns is *imperative*. Use a `try..finally` construct to make sure that the releaser function is
* always executed, even in the event of an error.
*
* @example
* ```typescript
* const releaser = await Mutex.acquire(myMutexToken);
* try {
* const releaser = Mutex.acquireSync(myMutexToken);
* if (releaser !== "timed-out") {
* try {
* ...
* }
* finally {
* }
* finally {
* releaser();
* }
* }
* ```
* @param token Mutex token to be acquired.
* @param timeout Timeout value in milliseconds to wait for acquisition.
* @returns A releaser object that can and should be used for releasing the mutex, or the value `'timed-out'` if
* @returns A releaser function 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 acquireSync(token: Token, timeout: number): Releaser | "timed-out";
static acquireSync(token: Token, timeout?: number) {
return acquireSync(mutexIdentityData, token, timeout);
}
/**
* Asynchronously acquires exclusivity over the specified mutex's token.
*
* **IMPORTANT**: Acquiring a mutex halts other work. Releasing the mutex by invoking the releaser function that
* this method returns is *imperative*. Use a `try..finally` construct to make sure that the releaser function is
* always executed, even in the event of an error.
*
* @example
* ```typescript
* const releaser = await Mutex.acquire(myMutexToken);
* try {
* ...
* }
* finally {
* releaser();
* }
* ```
* @param token Mutex token to be acquired.
* @returns A releaser function that can and should be used for releasing the mutex.
*/
static acquire(token: Token): Promise<Releaser>;
/**
* Asynchronously acquires exclusivity over the specified mutex's token.
*
* **IMPORTANT**: Acquiring a mutex halts other work. Releasing the mutex by invoking the releaser function that
* this method returns is *imperative*. Use a `try..finally` construct to make sure that the releaser function is
* always executed, even in the event of an error.
*
* @example
* ```typescript
* const releaser = await myMutex.acquire(500);
* if (releaser !== "timed-out") {
* try {
* ...
* }
* finally {
* releaser();
* }
* }
* ```
* @param token Mutex token to be acquired.
* @param timeout Timeout value in milliseconds to wait for acquisition.
* @returns A releaser function 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): Promise<Releaser | "timed-out">;
static acquire(token: Token, timeout?: number) {
return acquire(mutexIdentityData, token, timeout);
}
Expand Down
130 changes: 115 additions & 15 deletions src/sync/Semaphore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -10,13 +10,15 @@ export type Releaser = () => void;
export class SemaphoreInternal extends SyncObject {
#capacity;
#disabled;
constructor(identifier: number, capacity: number, createDisabled: boolean) {
#identifierData;
constructor(identifierData: IdentifierData, capacity: number, createDisabled: boolean) {
if (capacity <= 0 || !Number.isInteger(capacity)) {
throw new Error("A semaphore's capacity can only be a positive integer.");
}
super(identifier, createDisabled ? 0 : capacity);
super(identifierData[0], createDisabled ? 0 : capacity);
this.#capacity = capacity;
this.#disabled = createDisabled;
this.#identifierData = identifierData;
}

/**
Expand All @@ -32,6 +34,100 @@ export class SemaphoreInternal extends SyncObject {
Atomics.notify(this.token, 0);
return true;
}
/**
* Acquires exclusivity over the mutex or a slot in the semaphore.
*
* **IMPORTANT**: Acquiring a mutex or a semaphore slot halts other work. Releasing the mutex or semaphore slot
* by invoking the releaser function that this method returns is *imperative*. Use a `try..finally` construct to
* make sure that the releaser function is always executed, even in the event of an error.
*
* @example
* ```typescript
* const releaser = myMutexOrSemaphore.acquireSync();
* try {
* ...
* }
* finally {
* releaser();
* }
* ```
* @returns A releaser function that can and should be used for releasing the mutex or semaphore slot.
*/
acquireSync(): Releaser;
/**
* Acquires exclusivity over the mutex or a slot in the semaphore.
*
* **IMPORTANT**: Acquiring a mutex or a semaphore slot halts other work. Releasing the mutex or semaphore slot
* by invoking the releaser function that this method returns is *imperative*. Use a `try..finally` construct to
* make sure that the releaser function is always executed, even in the event of an error.
*
* @example
* ```typescript
* const releaser = myMutexOrSemaphore.acquireSync(500);
* if (releaser !== "timed-out") {
* try {
* ...
* }
* finally {
* releaser();
* }
* }
* ```
* @param timeout Timeout value in milliseconds to wait for acquisition.
* @returns A releaser function that can and should be used for releasing the mutex or semaphore slot, or the value
* `'timed-out'` if the mutex or semaphore slot could not be acquired before the specified timeout time elapsed.
*/
acquireSync(timeout: number): Releaser | "timed-out";
acquireSync(timeout?: number) {
return acquireSync(this.#identifierData, this.token, timeout);
}
/**
* Asynchronously acquires exclusivity over the mutex or a slot in the semaphore.
*
* **IMPORTANT**: Acquiring a mutex or a semaphore slot halts other work. Releasing the mutex or semaphore slot
* by invoking the releaser function that this method returns is *imperative*. Use a `try..finally` construct to
* make sure that the releaser function is always executed, even in the event of an error.
*
* @example
* ```typescript
* const releaser = await myMutexOrSemaphore.acquire();
* try {
* ...
* }
* finally {
* releaser();
* }
* ```
* @returns A releaser function that can and should be used for releasing the mutex or semaphore slot.
*/
acquire(): Promise<Releaser>;
/**
* Asynchronously acquires exclusivity over the mutex or a slot in the semaphore.
*
* **IMPORTANT**: Acquiring a mutex or a semaphore slot halts other work. Releasing the mutex or semaphore slot
* by invoking the releaser function that this method returns is *imperative*. Use a `try..finally` construct to
* make sure that the releaser function is always executed, even in the event of an error.
*
* @example
* ```typescript
* const releaser = await myMutexOrSemaphore.acquire(500);
* if (releaser !== "timed-out") {
* try {
* ...
* }
* finally {
* releaser();
* }
* }
* ```
* @param timeout Timeout value in milliseconds to wait for acquisition.
* A releaser function that can and should be used for releasing the mutex or semaphore slot, or the value
* `'timed-out'` if the mutex or semaphore slot could not be acquired before the specified timeout time elapsed.
*/
acquire(timeout: number): Promise<"timed-out" | Releaser>;
acquire(timeout?: number) {
return acquire(this.#identifierData, this.token, timeout);
}
};

function buildReleaser(token: Token) {
Expand Down Expand Up @@ -90,7 +186,7 @@ export async function acquire(identifierData: IdentifierData, token: Token, time
*/
export class Semaphore extends SemaphoreInternal {
constructor(capacity: number, createDisabled: boolean = false) {
super(semaphoreIdentityData[0], capacity, createDisabled);
super(semaphoreIdentityData, capacity, createDisabled);
}
/**
* Acquires a slot from the specified semaphore's token.
Expand All @@ -103,10 +199,10 @@ export class Semaphore extends SemaphoreInternal {
* ```typescript
* const releaser = Semaphore.acquireSync(mySemaphoreToken);
* try {
* ...
* ...
* }
* finally {
* releaser();
* releaser();
* }
* ```
* @param token Semaphore token to be acquired.
Expand All @@ -123,17 +219,19 @@ export class Semaphore extends SemaphoreInternal {
* @example
* ```typescript
* const releaser = Semaphore.acquireSync(mySemaphoreToken);
* try {
* if (releaser !== "timed-out") {
* try {
* ...
* }
* finally {
* }
* finally {
* releaser();
* }
* }
* ```
* @param token Semaphore token to be acquired.
* @param timeout Optional timeout value in milliseconds to wait for acquisition.
* @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.
* if the semaphore could not be acquired before the specified timeout time elapsed.
*/
static acquireSync(token: Token, timeout: number): 'timed-out' | Releaser;
static acquireSync(token: Token, timeout?: number) {
Expand All @@ -151,10 +249,10 @@ export class Semaphore extends SemaphoreInternal {
* ```typescript
* const releaser = await Semaphore.acquire(mySemaphoreToken);
* try {
* ...
* ...
* }
* finally {
* releaser();
* releaser();
* }
* ```
* @param token Semaphore token to be acquired.
Expand All @@ -171,17 +269,19 @@ export class Semaphore extends SemaphoreInternal {
* @example
* ```typescript
* const releaser = await Semaphore.acquire(mySemaphoreToken);
* try {
* if (releaser !== "timed-out") {
* try {
* ...
* }
* finally {
* }
* finally {
* releaser();
* }
* }
* ```
* @param token Semaphore token to be acquired.
* @param timeout Timeout value in milliseconds to wait for acquisition.
* @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.
* if the semaphore could not be acquired before the specified timeout time elapsed.
*/
static acquire(token: Token, timeout: number): Promise<"timed-out" | Releaser>;
static acquire(token: Token, timeout?: number) {
Expand Down
Loading