summaryrefslogtreecommitdiffhomepage
path: root/packages/kernel/src/bus/pure.ts
diff options
context:
space:
mode:
Diffstat (limited to 'packages/kernel/src/bus/pure.ts')
-rw-r--r--packages/kernel/src/bus/pure.ts122
1 files changed, 61 insertions, 61 deletions
diff --git a/packages/kernel/src/bus/pure.ts b/packages/kernel/src/bus/pure.ts
index 4d90fc6..a1c7a86 100644
--- a/packages/kernel/src/bus/pure.ts
+++ b/packages/kernel/src/bus/pure.ts
@@ -2,82 +2,82 @@ import type { Logger } from "../contracts/extension.js";
import type { EventHandler, FilterHandler } from "../contracts/hooks.js";
export function dispatchEventSync<T>(
- handlers: ReadonlyArray<EventHandler<T>>,
- payload: T,
- logger: Logger,
- hookId: string,
+ handlers: ReadonlyArray<EventHandler<T>>,
+ payload: T,
+ logger: Logger,
+ hookId: string,
): void {
- for (const handler of handlers) {
- try {
- const result = handler(payload);
- if (result instanceof Promise) {
- result.catch((err: unknown) => {
- logger.error(`Event hook "${hookId}" handler rejected`, { err });
- });
- }
- } catch (err) {
- logger.error(`Event hook "${hookId}" handler threw`, { err });
- }
- }
+ for (const handler of handlers) {
+ try {
+ const result = handler(payload);
+ if (result instanceof Promise) {
+ result.catch((err: unknown) => {
+ logger.error(`Event hook "${hookId}" handler rejected`, { err });
+ });
+ }
+ } catch (err) {
+ logger.error(`Event hook "${hookId}" handler threw`, { err });
+ }
+ }
}
export async function dispatchEventAsync<T>(
- handlers: ReadonlyArray<EventHandler<T>>,
- payload: T,
- logger: Logger,
- hookId: string,
- timeoutMs?: number,
+ handlers: ReadonlyArray<EventHandler<T>>,
+ payload: T,
+ logger: Logger,
+ hookId: string,
+ timeoutMs?: number,
): Promise<void> {
- const promises = handlers.map(async (handler) => {
- try {
- await handler(payload);
- } catch (err) {
- logger.error(`Event hook "${hookId}" handler threw`, { err });
- }
- });
+ const promises = handlers.map(async (handler) => {
+ try {
+ await handler(payload);
+ } catch (err) {
+ logger.error(`Event hook "${hookId}" handler threw`, { err });
+ }
+ });
- if (timeoutMs !== undefined) {
- await Promise.race([
- Promise.all(promises),
- new Promise<void>((resolve) => {
- setTimeout(resolve, timeoutMs);
- }),
- ]);
- } else {
- await Promise.all(promises);
- }
+ if (timeoutMs !== undefined) {
+ await Promise.race([
+ Promise.all(promises),
+ new Promise<void>((resolve) => {
+ setTimeout(resolve, timeoutMs);
+ }),
+ ]);
+ } else {
+ await Promise.all(promises);
+ }
}
export interface FilterEntry<T> {
- readonly fn: FilterHandler<T>;
- readonly priority: number;
- readonly order: number;
+ readonly fn: FilterHandler<T>;
+ readonly priority: number;
+ readonly order: number;
}
export function sortFilters<T>(
- entries: ReadonlyArray<FilterEntry<T>>,
+ entries: ReadonlyArray<FilterEntry<T>>,
): ReadonlyArray<FilterEntry<T>> {
- return [...entries].sort((a, b) => {
- if (a.priority !== b.priority) return a.priority - b.priority;
- return a.order - b.order;
- });
+ return [...entries].sort((a, b) => {
+ if (a.priority !== b.priority) return a.priority - b.priority;
+ return a.order - b.order;
+ });
}
export async function applyFilterChain<T>(
- filters: ReadonlyArray<FilterHandler<T>>,
- value: T,
- logger: Logger,
- hookId: string,
- failClosed: boolean,
+ filters: ReadonlyArray<FilterHandler<T>>,
+ value: T,
+ logger: Logger,
+ hookId: string,
+ failClosed: boolean,
): Promise<T> {
- let current = value;
- for (const fn of filters) {
- try {
- current = await fn(current);
- } catch (err) {
- if (failClosed) throw err;
- logger.error(`Filter "${hookId}" handler threw (fail-open, passing through)`, { err });
- }
- }
- return current;
+ let current = value;
+ for (const fn of filters) {
+ try {
+ current = await fn(current);
+ } catch (err) {
+ if (failClosed) throw err;
+ logger.error(`Filter "${hookId}" handler threw (fail-open, passing through)`, { err });
+ }
+ }
+ return current;
}