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
3 changes: 3 additions & 0 deletions src/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -7,9 +7,12 @@ export { AutoResetEvent } from "./sync/AutoResetEvent.js";
export { ManualResetEvent } from "./sync/ManualResetEvent.js";
export { Mutex } from "./sync/Mutex.js";
export { Semaphore } from "./sync/Semaphore.js";
export type * from "./sync/Semaphore.js";
// workers
export * from './workers/AsyncWorker.js';
export type * from './workers/AsyncWorker.js';
export * from "./workers/workerListener.js";
export type * from "./workers/workerListener.js";
export * from './workers/WorkerTerminatedMessage.js';
export * from "./workers/WorkItem.js";
export type * from "./types.js";
2 changes: 1 addition & 1 deletion src/sync/Semaphore.ts
Original file line number Diff line number Diff line change
Expand Up @@ -134,7 +134,7 @@ function buildReleaser(token: Token) {
let released = false;
return (() => {
if (released) {
throw new Error('The semaphore has already been released and cannot be released again.');
throw new Error('The mutex or semaphore has already been released and cannot be released again.');
}
Atomics.add(token, 0, 1);
Atomics.notify(token, 0, 1);
Expand Down
75 changes: 70 additions & 5 deletions src/types.ts
Original file line number Diff line number Diff line change
Expand Up @@ -35,34 +35,99 @@ export type AsyncMessageUntyped = {
cancelToken?: Token | undefined;
payload?: any;
};

/**
* Message sent from the worker back to the main thread, containing the result of a previously queued work item or an
* intermediate result or status update.
*/
export type AsyncResponse = {
/**
* The unique identifier of the work item to which this response corresponds.
*/
workItemId: number;
/**
* The payload containing the result, intermediate data, or status update from the worker.
*/
payload?: any;
};

/**
* Function type for processing messages from the worker.
* @param payload The payload received from the worker.
* @returns A boolean indicating whether the message was successfully processed.
*/
export type ProcessMessageFn = (payload: any) => boolean;

/**
* Represents the data associated with a work item in the async worker queue.
*/
export type WorkItemData<TResult> = {
/**
* The unique identifier for the work item.
*/
id: number;
/**
* The name of the task associated with the work item.
*/
task: string;
/**
* The promise associated with the work item, which will be resolved or rejected based on the work item's
* completion.
*/
promise: Promise<TResult>;
/**
* Function to resolve the promise associated with the work item.
*/
resolve: (result: any) => void;
/**
* Function to reject the promise associated with the work item.
*/
reject: (reason: any) => void;
/**
* The payload associated with the work item.
*/
payload: any,
};

/**
* Message indicating that a previously queued work item has been cancelled.
*/
export type TaskCancelledMessage = {
/**
* The unique identifier of the work item that has been cancelled.
*/
workItemId: number;
/**
* The payload indicating that the work item has been cancelled.
*/
payload: {
_$cancelled: true;
}
};

/**
* Options for queueing work items in the async worker.
*/
export type QueueingOptions = {
/**
* Indicates whether the queued work item can be cancelled. This signals the library to provide cancellation
* capabilities for the work item.
* @default false
*/
cancellable?: boolean;
/**
* Optional callback function to process messages from the worker.
* @default undefined
*/
processMessage?: ProcessMessageFn;
/**
* Optional array of transferable objects to be sent along with the message to the worker.
* @default []
*/
transferables?: Transferable[];
/**
* Indicates whether the work items can be processed out of order.
*
* **IMPORTANT**: When `outOfOrder` is set to `true`, the work item is sent immediately to the worker, skipping
* the current queue order. This may produce unexpected behavior because the work item status is not tracked and
* their completion doesn't trigger dequeueing of subsequent work items.
* @default false
*/
outOfOrder?: boolean;
}

Expand Down
12 changes: 11 additions & 1 deletion src/workers/AsyncWorker.ts
Original file line number Diff line number Diff line change
Expand Up @@ -23,9 +23,19 @@ function isSharedWorker(worker: any): worker is SharedWorker {
typeof worker.port.removeEventListener === 'function';
}

/**
* Represents the function used to enqueue a work item for a specific task.
* @template Fn The type of the function representing the task.
* @param payload The payload to be passed to the task function.
* @param options Optional queueing options.
* @returns A work item representing the enqueued task.
*/
export type EnqueueFn<Fn extends ((...args: any[]) => any) = (() => any)> =
(payload: Fn extends () => any ? void : Parameters<Fn>[0], options?: QueueingOptions) => WorkItem<ReturnType<Fn>>;

/**
* Represents the object used to enqueue work items for all tasks of a worker.
* @template T The type of the tasks object, mapping task names to their respective functions.
*/
export type Enqueue<T extends Record<string, (...args: any[]) => any>> = {
[K in keyof T]: EnqueueFn<T[K]>;
};
Expand Down
2 changes: 2 additions & 0 deletions src/workers/workerListener.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,8 @@ import type { AsyncMessageUntyped, AsyncResponse } from "../types.js";

/**
* Defines the function provided to worker tasks so workers can communicate back to the calling thread.
* @param payload The payload to be sent to the calling thread.
* @param options Optional post message options.
*/
export type PostFn = (payload: any, options?: WindowPostMessageOptions) => void;

Expand Down
Loading