diff --git a/.changeset/hooks-add-pipe.md b/.changeset/hooks-add-pipe.md new file mode 100644 index 00000000..c268559b --- /dev/null +++ b/.changeset/hooks-add-pipe.md @@ -0,0 +1,15 @@ +--- +"@logosdx/hooks": minor +"@logosdx/fetch": patch +--- + +`HookEngine.addPipe()` registers pipe middleware (#147) + +`@logosdx/hooks`: + +- New `addPipe(name, callback, options?)` method, typed against the lifecycle's `(next, ...args, ctx)` shape. Pipe middleware — retry, dedupe, caching execution — now registers with full type inference instead of requiring an `as any` cast on `add()`. +- `addPipe` shares the same registry, `AddOptions` semantics (`priority`, `once`, `times`, `ignoreOnFail`), cleanup-function return, and `register()` strict-mode enforcement as `add()`. Runtime behavior of `add`, `pipe`, `pipeSync` is unchanged. + +`@logosdx/fetch`: + +- `retryPlugin` and `dedupePlugin` register their `execute` middleware via `addPipe` instead of `add(... as any)`. No behavior change. diff --git a/docs/packages/hooks.md b/docs/packages/hooks.md index 80218df3..d49b1391 100644 --- a/docs/packages/hooks.md +++ b/docs/packages/hooks.md @@ -127,6 +127,7 @@ new HookEngine(options?) |--------|-------------| | `register(...names)` | Enable strict mode. Returns `this` for chaining. | | `add(name, callback, options?)` | Subscribe. Returns cleanup function. | +| `addPipe(name, callback, options?)` | Subscribe pipe middleware. Returns cleanup function. | | `run(name, ...args)` | Run hook async. Returns `Promise`. | | `runSync(name, ...args)` | Run hook sync. Returns `RunResult`. | | `pipe(name, coreFn, ...args)` | Pipe hook async (onion middleware). Returns result. | @@ -195,6 +196,8 @@ const actualResult = await doWork(...args); ### AddOptions +Same options for `add()` and `addPipe()`: + ```typescript hooks.add('name', callback, { once: true, // Remove after first run (sugar for times: 1) @@ -262,9 +265,9 @@ const wrappedValidate = hooks.wrapSync( // Post: receives (result, ...args, ctx) — can transform result ``` -### pipe() / pipeSync() +### addPipe() -Onion/middleware composition where each callback wraps the next. Used for cross-cutting concerns like retry, deduplication, and caching execution. +Subscribe pipe middleware to a lifecycle hook, typed against the lifecycle's `(next, ...args, ctx)` shape — middleware registers with full inference, no `as any` cast required. ```typescript interface PipeLifecycle { @@ -274,18 +277,24 @@ interface PipeLifecycle { const hooks = new HookEngine() .register('execute'); -// Add middleware — receives (next, ...args, ctx) -hooks.add('execute', async (next, opts, ctx) => { +// addPipe infers (next, opts, ctx) from the lifecycle signature +hooks.addPipe('execute', async (next, opts, ctx) => { console.log('before core'); const result = await next(); // call next middleware or core function console.log('after core'); return result; }, { priority: -10 }); +``` -// Run the pipe — core function is the innermost call +### pipe() / pipeSync() + +Run middleware registered via `addPipe()` as an onion/middleware composition — each callback wraps the next. Used for cross-cutting concerns like retry, deduplication, and caching execution. + +```typescript +// Run the pipe — core function is the innermost call, closes over opts const result = await hooks.pipe('execute', - async (opts) => fetch(opts.url, opts), // core function - opts // spread args + () => fetch(opts.url, opts), // core function + opts // passed to each middleware call ); ``` @@ -304,7 +313,7 @@ const result = await hooks.pipe('execute', ```typescript const result = hooks.pipeSync('validate', - (data) => validate(data), + () => validate(data), data ); ``` diff --git a/docs/spec/hooks-add-pipe.md b/docs/spec/hooks-add-pipe.md new file mode 100644 index 00000000..73f75f26 --- /dev/null +++ b/docs/spec/hooks-add-pipe.md @@ -0,0 +1,97 @@ +# Spec: `HookEngine.addPipe()` — typed registration for pipe middleware + +Resolves [#147](https://github.com/logosdx/monorepo/issues/147). The issue body is the design +analysis: three options were weighed there (an `add` overload, a branded `Pipe<>` lifecycle +marker, a sibling method) and the sibling method was chosen. This spec is the implementation +contract for that choice. + +## Problem + +`HookEngine.add()` types every callback as `HookCallback` (the `run()` shape: `(...args, ctx)` +returning `void | EarlyReturnSignal`). Pipe middleware is `(next, ...args, ctx)` and must return +the response, so it cannot be registered without `as any` or `@ts-expect-error`. The exported +`PipeCallback` type describes the correct shape but is accepted nowhere. In this repo, +`packages/fetch/src/plugins/retry.ts:86` and `packages/fetch/src/plugins/dedupe.ts:131` carry +the cast today. + +## Decision + +Add one public method, `addPipe`, typed with `PipeCallback` against the lifecycle signature. +It stores into the **same registry** `add()` uses — `pipe()`/`pipeSync()` already invoke +whatever was stored, so the fix is type-level; runtime behavior of `add`, `pipe`, `pipeSync` +is unchanged. + +## API contract + +```ts +addPipe>( + name: K, + callback: PipeCallback< + Parameters>, + Awaited>>, + FailArgs + >, + options?: HookEngine.AddOptions, +): () => void; +``` + +- Same runtime validation as `add`: assert `name` is a string, assert `callback` is a function, + enforce `register()` strict mode (error message names `addPipe`). +- Same `AddOptions` semantics: `priority` (lower = outermost layer), `once`, `times`, + `ignoreOnFail` — all already honored by `pipe()`'s chain builder. +- Returns the same cleanup-function shape as `add`. +- Implementation shares the insertion logic with `add` (extract a private helper both call); + do not duplicate the priority-insert loop. No public-facing casts; the internal registry is + already `HookEntry`. +- JSDoc with a WHY-bearing example mirroring `pipe()`'s retry/dedupe example, and the + `add` JSDoc's pipe-shaped example moves to `addPipe` (that example currently does not + compile against `add` — it is the bug). + +## Checkpoints + +| # | Deliverable | Where | Done when | +|---|-------------|-------|-----------| +| 1 | `addPipe` method + shared insert helper + JSDoc | `packages/hooks/src/index.ts` | Issue #147's repro compiles cast-free using `addPipe`; `pnpm build` green | +| 2 | Runtime + type coverage | `tests/src/hooks.ts` | New `engine.addPipe()`, `engine.pipe()`, `engine.pipeSync()` describes pass (see Verification) | +| 3 | Migrate fetch plugins to `addPipe` | `packages/fetch/src/plugins/retry.ts`, `packages/fetch/src/plugins/dedupe.ts` | The outer `(... ) as any` on both registrations is gone; fetch test suite green | +| 4 | Docs + changeset | `docs/packages/hooks.md`, `skills/logosdx/references/hooks.md`, `.changeset/` | Both doc surfaces show `addPipe`; changeset: minor `@logosdx/hooks`, patch `@logosdx/fetch` | + +## Verification + +- Checkpoint 2 must cover, at minimum: + - `addPipe`-registered middleware executes via `pipe()` and `pipeSync()` (onion order + follows `priority`, lower = outermost). + - Short-circuit by not calling `next()`; result propagation from `next()`. + - `ctx.args()` replacement reaching inner layers; `ctx.fail()`; `ctx.removeHook()`. + - `once`, `times`, `ignoreOnFail` options through the pipe path. + - Cleanup function removes the middleware. + - Strict-mode `register()` enforcement for `addPipe`. + - Type-level: the issue's repro (a `PipeCallback`-typed const and an inline callback with + inferred params) compiles when passed to `addPipe`; test file compiles under the suite's + type checking with relative imports to `packages/hooks/src/index.ts`. +- **Build before test**: test-suite package imports resolve `@logosdx/*` to `dist/`. Run + `pnpm build` after touching `packages/hooks/src` or the fetch plugins, or new exports throw + "not a function" at test time. +- Full gate: `pnpm build` then `pnpm test` from repo root, all green. + +## Non-goals + +- Tightening `pipe()`/`pipeSync()` argument typing against the lifecycle (`...args: unknown[]`, + unbound `R`) — potentially breaking for existing consumers; the issue lists it as Related, + not as the fix. +- The branded `Pipe<>` lifecycle-marker design — rejected in the issue (requires every + consumer lifecycle to annotate). +- `PipeOptions.append` typing — Related-listed, out of scope. +- Any behavior change to `add`, `run`, `runSync`, `pipe`, `pipeSync`. +- De-`any`-ing the interiors of the fetch retry/dedupe plugins beyond removing the outer + registration cast. + +## Change log + +- 2026-08-08: Initial spec from issue #147 (autopilot). +- 2026-08-09: Implemented as specified, no contract deviations; shipped as a single + squashed commit (PR #148). Checkpoints 1–2: `addPipe` + `#insertEntry` extraction + + 26 tests — first direct `pipe()`/`pipeSync()` coverage in the suite. Checkpoint 3: + retry/dedupe migration, 6-line diff. Checkpoint 4: docs + changeset; also fixed + three pre-existing doc examples that passed args to the zero-arg `coreFn`. Full + gate green: build, 2434 tests, tsc clean. diff --git a/docs/wiki/fetch.md b/docs/wiki/fetch.md index e216acd5..9004a453 100644 --- a/docs/wiki/fetch.md +++ b/docs/wiki/fetch.md @@ -42,7 +42,7 @@ description: HTTP client (`FetchEngine`) with resilience policies, plugins, cook - Depends on `@logosdx/utils` for `attempt`, flow control, and validation helpers. - Depends on `@logosdx/observer` for event emission; `FetchEngine` emits typed lifecycle events. -- Depends on `@logosdx/hooks` for the entire plugin architecture: `FetchEngine.hooks` is a `HookEngine` ([`packages/fetch/src/engine/index.ts`](../../packages/fetch/src/engine/index.ts)); every request runs the three-phase pipeline `beforeRequest` (run) → `execute` (pipe, onion-wrapped: retry wraps dedupe wraps the network call) → `afterRequest` (run) in [`packages/fetch/src/engine/executor.ts`](../../packages/fetch/src/engine/executor.ts); built-in policies install hooks at negative priorities so user hooks run after them; per-call `CallConfig.hooks` appends one-request hooks; a `HookScope` carries per-request state between phases. See [`docs/wiki/hooks.md`](hooks.md). +- Depends on `@logosdx/hooks` for the entire plugin architecture: `FetchEngine.hooks` is a `HookEngine` ([`packages/fetch/src/engine/index.ts`](../../packages/fetch/src/engine/index.ts)); every request runs the three-phase pipeline `beforeRequest` (run) → `execute` (pipe, onion-wrapped: retry wraps dedupe wraps the network call) → `afterRequest` (run) in [`packages/fetch/src/engine/executor.ts`](../../packages/fetch/src/engine/executor.ts); built-in policies install hooks at negative priorities so user hooks run after them; per-call `CallConfig.hooks` appends one-request hooks; a `HookScope` carries per-request state between phases; `retryPlugin` and `dedupePlugin` register their `execute`-phase middleware via `engine.hooks.addPipe('execute', ...)` — the typed, cast-free pipe registration, replacing a prior `engine.hooks.add((... ) as any, ...)` call. See [`docs/wiki/hooks.md`](hooks.md). - `@logosdx/react` wraps `FetchEngine` via `createFetchContext` in [`packages/react/src/fetch.ts`](../../packages/react/src/fetch.ts); also exports `useQuery`, `useMutation`, `useAsync`, `createApiHooks`. A change to `FetchResponse`/`FetchError` here forces a matching change to `FetchFailure` in [`packages/react/src/types.ts`](../../packages/react/src/types.ts). - Tests in [`tests/src/fetch/`](../../tests/src/fetch) cover engine, cookies, policies, serializers, state, adapters; [`tests/src/fetch/engine/plugin-resolution.test.ts`](../../tests/src/fetch/engine/plugin-resolution.test.ts) (17 tests) covers config-key/`plugins` array composition, the construction-time conflict throws, runtime `config.set()` ownership-conflict rejection (single-key and multi-key merge), convenience-method (`cacheStats`/`clearCache`/inflight count) parity across install paths, and the falsy-key-plus-plugin warn-vs-throw behavior. [`tests/src/fetch/executor/per-call-overrides.test.ts`](../../tests/src/fetch/executor/per-call-overrides.test.ts) (4 tests) covers per-call `retry: false` overriding an engine configured with retries, per-call `retry` config overriding an engine configured with `retry: false`, `res.config.retry` reflecting the engine default when no per-call override is given, and per-call `skipCache` bypassing both cache lookup and store. - [`tests/src/fetch/engine/configuration.test.ts`](../../tests/src/fetch/engine/configuration.test.ts) (13 tests) covers serializer/config edge cases plus runtime `config.set()` reconfigure for each policy: `retry.maxAttempts` taking effect on the next request, rate-limit token buckets rebuilding with a new capacity, cache TTL/rules reconfiguring while existing cached entries survive, the cache-adapter `reconfigureGuard` throwing and leaving the store unchanged (single-key and multi-key merge), the dedupe rule cache rebuilding, and a cookie `exclude` update applying without clearing the jar. [`tests/src/fetch/executor/retry.test.ts`](../../tests/src/fetch/executor/retry.test.ts) (42 tests) covers retry resolution including `attemptTimeout` firing under `retry: false`/`maxAttempts: 0`. diff --git a/docs/wiki/hooks.md b/docs/wiki/hooks.md index 590b8e12..2504875b 100644 --- a/docs/wiki/hooks.md +++ b/docs/wiki/hooks.md @@ -10,23 +10,25 @@ type: Domain ## Artifacts -- [`skills/logosdx/references/hooks.md`](../../skills/logosdx/references/hooks.md) — skill reference (326 LOC) covering lifecycle hooks, middleware, plugins, priority chains +- [`skills/logosdx/references/hooks.md`](../../skills/logosdx/references/hooks.md) — skill reference (335 LOC) covering lifecycle hooks, middleware, plugins, priority chains; includes an `## addPipe() vs add()` section distinguishing run-shaped `(...args, ctx)` callbacks from onion-style `(next, ...args, ctx)` middleware ## CLI code -- [`packages/hooks/src/index.ts`](../../packages/hooks/src/index.ts) — all hook engine implementation in a single file (1252 LOC, 36k chars); exports `HookError`, `isHookError`, and the `HookEngine` class +- [`packages/hooks/src/index.ts`](../../packages/hooks/src/index.ts) — all hook engine implementation in a single file (1316 LOC, 39k chars); exports `HookError`, `isHookError`, `PipeCallback`, and the `HookEngine` class +- `HookEngine.addPipe(name, callback, options?)` registers onion-style middleware typed against `PipeCallback` (the `(next, ...args, ctx)` shape) without an `as any` cast; it runs the same validation as `add()` (string `name`, function `callback`, `register()` strict-mode via `#assertRegistered`, which names `addPipe` in its thrown error when called from that path) and writes into the same internal registry `add()` uses, sharing `AddOptions` semantics (`priority`, `once`, `times`, `ignoreOnFail`) and the cleanup-function return. The priority-ordered splice-and-cleanup logic formerly inline in `add()` is now the private `#insertEntry(name, callback, options)`, shared by both `add()` and `addPipe()`. - [`packages/hooks/notes.md`](../../packages/hooks/notes.md) — internal design notes (700 LOC) ## Docs -- [`docs/packages/hooks.md`](../packages/hooks.md) — combined hooks reference (523 LOC) +- [`docs/packages/hooks.md`](../packages/hooks.md) — combined hooks reference (532 LOC) +- [`docs/spec/hooks-add-pipe.md`](../spec/hooks-add-pipe.md) — implementation spec for the `addPipe()` decision (issue #147); sibling-method approach chosen over an `add()` overload or a branded `Pipe<>` lifecycle marker ## Coupling - Depends on `@logosdx/utils` for `attempt`, `attemptSync`, `assert`, `isFunction`, `isObject`, and `FunctionProps`. - No dependency on `@logosdx/observer`. -- `@logosdx/fetch` is the largest in-repo consumer: `FetchEngine` composes a `HookEngine` ([`packages/fetch/src/engine/index.ts`](../../packages/fetch/src/engine/index.ts)) and runs every request through a three-phase pipeline — `hooks.run('beforeRequest')` → `hooks.pipe('execute')` → `hooks.run('afterRequest')` ([`packages/fetch/src/engine/executor.ts`](../../packages/fetch/src/engine/executor.ts)). Every fetch policy plugin (retry, dedupe, cache, rate-limit, cookies) is a hook installation whose `install()` returns the hook cleanup; a `HookScope` carries per-request state between phases (e.g. the cache key set in `beforeRequest` and read in `afterRequest`). A behavior change to `run`/`pipe`/`HookScope` semantics is a behavior change to the entire fetch pipeline. -- Tests in [`tests/src/hooks.ts`](../../tests/src/hooks.ts) (1462 LOC, 43k chars). +- `@logosdx/fetch` is the largest in-repo consumer: `FetchEngine` composes a `HookEngine` ([`packages/fetch/src/engine/index.ts`](../../packages/fetch/src/engine/index.ts)) and runs every request through a three-phase pipeline — `hooks.run('beforeRequest')` → `hooks.pipe('execute')` → `hooks.run('afterRequest')` ([`packages/fetch/src/engine/executor.ts`](../../packages/fetch/src/engine/executor.ts)). Every fetch policy plugin (retry, dedupe, cache, rate-limit, cookies) is a hook installation whose `install()` returns the hook cleanup; a `HookScope` carries per-request state between phases (e.g. the cache key set in `beforeRequest` and read in `afterRequest`). [`packages/fetch/src/plugins/retry.ts`](../../packages/fetch/src/plugins/retry.ts) and [`packages/fetch/src/plugins/dedupe.ts`](../../packages/fetch/src/plugins/dedupe.ts) register their `execute` middleware via `engine.hooks.addPipe(...)`, replacing the prior `engine.hooks.add((...) as any, ...)` cast pattern. A behavior change to `run`/`pipe`/`HookScope` semantics is a behavior change to the entire fetch pipeline. +- Tests in [`tests/src/hooks.ts`](../../tests/src/hooks.ts) (1923 LOC, 58k chars), including `describe('engine.addPipe()', ...)` (unsubscribe behavior, invalid name/callback rejection, `register()` strict-mode error naming `addPipe`, type-level `PipeCallback` registration with no casts) and `describe('engine.pipe()', ...)` (onion execution order by priority). - [`tests/src/smoke/hooks.test.ts`](../../tests/src/smoke/hooks.test.ts) runs browser smoke tests. ## Conventions worth knowing @@ -35,4 +37,4 @@ type: Domain - `run` executes a hook chain where a hook may modify args or short-circuit with a result; `pipe` onion-wraps a core function middleware-style (in fetch: retry wraps dedupe wraps the network call). - `ctx.fail(message)` within a hook handler halts the pipeline and throws `HookError` (or a custom error type if `handleFail` is overridden). - `isHookError(err)` type guard identifies hook pipeline failures. -- The entire implementation lives in one 1252-line file rather than being split across modules. +- The entire implementation lives in one 1316-line file rather than being split across modules. diff --git a/docs/wiki/testing.md b/docs/wiki/testing.md index c07fe01c..7688918c 100644 --- a/docs/wiki/testing.md +++ b/docs/wiki/testing.md @@ -24,7 +24,7 @@ The [`tests/`](../../tests) workspace runs the full validation suite for all pac - [`tests/src/observable/`](../../tests/src/observable) — unit tests for `@logosdx/observer` (engine, queue, relay) - [`tests/src/react/`](../../tests/src/react) — unit tests for `@logosdx/react` (10 files) - [`tests/src/storage/`](../../tests/src/storage) — unit tests for `@logosdx/storage` -- [`tests/src/hooks.ts`](../../tests/src/hooks.ts) — unit tests for `@logosdx/hooks` (1462 LOC) +- [`tests/src/hooks.ts`](../../tests/src/hooks.ts) — unit tests for `@logosdx/hooks` (1923 LOC), including direct coverage for `addPipe()`/`pipe()`/`pipeSync()` (onion middleware, priority ordering, cast-free typed registration) - [`tests/src/localize.ts`](../../tests/src/localize.ts) — unit tests for `@logosdx/localize` (835 LOC) - [`tests/src/localize-extractor.ts`](../../tests/src/localize-extractor.ts) — unit tests for the localize type extractor (426 LOC) - [`tests/src/state-machine.ts`](../../tests/src/state-machine.ts) — unit tests for `@logosdx/state-machine` (1499 LOC) diff --git a/packages/fetch/src/plugins/dedupe.ts b/packages/fetch/src/plugins/dedupe.ts index 8026728b..0ded6994 100644 --- a/packages/fetch/src/plugins/dedupe.ts +++ b/packages/fetch/src/plugins/dedupe.ts @@ -128,7 +128,7 @@ export function dedupePlugin( install(engine: FetchEnginePublic): () => void { - const cleanup = engine.hooks.add('execute', (async (next: () => Promise, opts: any) => { + const cleanup = engine.hooks.addPipe('execute', async (next, opts) => { const normalizedOpts = opts as InternalReqOptions; const { method, path } = normalizedOpts; @@ -172,7 +172,7 @@ export function dedupePlugin( throw err; } - return new Promise((resolve, reject) => { + return new Promise((resolve, reject) => { let settled = false; @@ -242,7 +242,7 @@ export function dedupePlugin( inflightMap.delete(key); } - }) as any, { priority: -30 }); + }, { priority: -30 }); return cleanup; } diff --git a/packages/fetch/src/plugins/retry.ts b/packages/fetch/src/plugins/retry.ts index b40ce0e8..7d6f22b2 100644 --- a/packages/fetch/src/plugins/retry.ts +++ b/packages/fetch/src/plugins/retry.ts @@ -83,7 +83,7 @@ export function retryPlugin( install(engine: FetchEnginePublic): () => void { - const cleanup = engine.hooks.add('execute', (async (next: () => Promise, opts: any, _ctx: any) => { + const cleanup = engine.hooks.addPipe('execute', async (next, opts, _ctx) => { const normalizedOpts = opts as InternalReqOptions; const mergedRetry = resolveConfigForRequest(normalizedOpts); @@ -129,7 +129,7 @@ export function retryPlugin( (normalizedOpts as any).signal = attemptController.signal; } - const [result, err] = await attempt(async () => next()); + const [result, err] = await attempt(async (): Promise => next()); // Restore original controller/signal if (normalizedOpts.attemptTimeout !== undefined) { @@ -280,7 +280,7 @@ export function retryPlugin( throw new FetchError('Unexpected end of retry logic'); - }) as any, { priority: -20 }); + }, { priority: -20 }); return cleanup; } diff --git a/packages/hooks/src/index.ts b/packages/hooks/src/index.ts index f69cbee5..bef52ba4 100644 --- a/packages/hooks/src/index.ts +++ b/packages/hooks/src/index.ts @@ -297,7 +297,7 @@ export class HookContext< * flow by calling or not calling `next()`. * * @example - * hooks.add('execute', async (next, opts, ctx) => { + * hooks.addPipe('execute', async (next, opts, ctx) => { * // Modify opts for inner layers * ctx.args({ ...opts, headers: { ...opts.headers, Auth: token } }); * @@ -574,44 +574,57 @@ export class HookEngine = { callback, options, priority }; - - const hooks = this.#hooks.get(name as string) ?? []; - - let inserted = false; - - for (let i = 0; i < hooks.length; i++) { - - if (hooks[i]!.priority > priority) { - - hooks.splice(i, 0, entry); - inserted = true; - break; - } - } - - if (!inserted) { - - hooks.push(entry); - } - - this.#hooks.set(name as string, hooks); - - return () => { - - const arr = this.#hooks.get(name as string); + return this.#insertEntry(name as string, callback, options); + } - if (arr) { + /** + * Subscribe pipe middleware to a lifecycle hook, typed against the lifecycle's + * `(next, ...args, ctx)` shape so retry/dedupe-style wrappers register without + * `as any` or `@ts-expect-error`. + * + * Stores into the same registry `add()` uses — `pipe()`/`pipeSync()` already + * invoke whatever is stored there, so this is a type-level fix, not a new + * runtime path. + * + * @param name - Name of the lifecycle hook + * @param callback - Middleware receiving `(next, ...args, ctx)` + * @param options - Options for this subscription + * @returns Cleanup function to remove the subscription + * + * @example + * // Retry middleware: call next(), retry on failure + * hooks.addPipe('execute', async (next, opts, ctx) => { + * for (let i = 0; i < 3; i++) { + * const [result, err] = await attempt(next); + * if (!err) return result; + * await wait(1000 * i); + * } + * throw lastError; + * }, { priority: -20 }); + * + * // Dedupe middleware: short-circuit by not calling next() + * hooks.addPipe('execute', async (next, opts, ctx) => { + * const inflight = getInflight(key); + * if (inflight) return inflight; + * return next(); + * }, { priority: -30 }); + */ + addPipe>( + name: K, + callback: PipeCallback< + Parameters>, + Awaited>>, + FailArgs + >, + options: HookEngine.AddOptions = {} + ): () => void { - const idx = arr.indexOf(entry); + assert(typeof name === 'string', '"name" must be a string'); + assert(isFunction(callback), '"callback" must be a function'); - if (idx !== -1) { + this.#assertRegistered(name, 'addPipe'); - arr.splice(idx, 1); - } - } - }; + return this.#insertEntry(name as string, callback, options); } /** @@ -913,7 +926,7 @@ export class HookEngine { + * hooks.addPipe('execute', async (next, opts, ctx) => { * for (let i = 0; i < 3; i++) { * const [result, err] = await attempt(next); * if (!err) return result; @@ -923,7 +936,7 @@ export class HookEngine { + * hooks.addPipe('execute', async (next, opts, ctx) => { * const inflight = getInflight(key); * if (inflight) return inflight; * const result = await next(); @@ -1161,6 +1174,57 @@ export class HookEngine any, + options: HookEngine.AddOptions + ): () => void { + + const priority = options.priority ?? 0; + const entry: HookEntry = { callback, options, priority }; + + const hooks = this.#hooks.get(name) ?? []; + + let inserted = false; + + for (let i = 0; i < hooks.length; i++) { + + if (hooks[i]!.priority > priority) { + + hooks.splice(i, 0, entry); + inserted = true; + break; + } + } + + if (!inserted) { + + hooks.push(entry); + } + + this.#hooks.set(name, hooks); + + return () => { + + const arr = this.#hooks.get(name); + + if (arr) { + + const idx = arr.indexOf(entry); + + if (idx !== -1) { + + arr.splice(idx, 1); + } + } + }; + } + /** * Remove an entry from the hooks array. */ diff --git a/skills/logosdx/references/hooks.md b/skills/logosdx/references/hooks.md index 19d43289..46a7c63e 100644 --- a/skills/logosdx/references/hooks.md +++ b/skills/logosdx/references/hooks.md @@ -55,7 +55,8 @@ cleanup(); | Method | Description | |--------|-------------| | `register(...names)` | Register hooks for runtime validation. Returns `this`. | -| `add(name, callback, options?)` | Subscribe to hook. Returns cleanup function. | +| `add(name, callback, options?)` | Subscribe run-shaped callback `(...args, ctx)`. Returns cleanup function. | +| `addPipe(name, callback, options?)` | Subscribe pipe middleware `(next, ...args, ctx)`. Returns cleanup function. | | `run(name, ...args)` | Run hook async. Returns `Promise`. | | `runSync(name, ...args)` | Run hook sync. Returns `RunResult`. | | `pipe(name, coreFn, ...args)` | Pipe hook async (onion middleware). Returns `Promise`. | @@ -122,6 +123,22 @@ await hooks.run('beforeRequest', url, opts, { ``` +## addPipe() vs add() + +`add()` registers a run-shaped callback `(...args, ctx)`, invoked by `run()`/`runSync()`. `addPipe()` registers pipe middleware `(next, ...args, ctx)`, invoked by `pipe()`/`pipeSync()` — typed against the lifecycle so `next`/args/`ctx` infer correctly, no `as any` cast needed. + +Use `add()` for linear before/after hooks. Use `addPipe()` for onion-style middleware wrapping a core call (retry, dedupe, caching execution) that controls flow via `next()`. + +```typescript +hooks.addPipe('execute', async (next, opts, ctx) => { + const start = Date.now(); + const result = await next(); + console.log(`Took ${Date.now() - start}ms`); + return result; +}, { priority: -10 }); +``` + + ## pipe() — Middleware Composition Onion/middleware pattern where each callback wraps the next. Used for cross-cutting concerns like retry, dedupe, caching execution. @@ -129,21 +146,13 @@ Onion/middleware pattern where each callback wraps the next. Used for cross-cutt ```typescript // Core function is the innermost call const result = await hooks.pipe('execute', - async (opts) => { + async () => { const [res, err] = await attempt(() => fetch(opts.url, opts)); if (err) throw err; return res; }, opts ); - -// Callbacks: (next, ...args, ctx) — call next() to proceed -hooks.add('execute', async (next, opts, ctx) => { - const start = Date.now(); - const result = await next(); - console.log(`Took ${Date.now() - start}ms`); - return result; -}, { priority: -10 }); ``` **PipeContext** — simpler than HookContext: diff --git a/tests/src/hooks.ts b/tests/src/hooks.ts index 47eaf68f..bdb0ef73 100644 --- a/tests/src/hooks.ts +++ b/tests/src/hooks.ts @@ -10,7 +10,9 @@ import { HookError, isHookError, HookContext, - HookScope + HookScope, + PipeContext, + PipeCallback } from '../../packages/hooks/src/index.ts'; import { attempt, attemptSync } from '../../packages/utils/src/index.ts'; @@ -951,6 +953,465 @@ describe('@logosdx/hooks', () => { }); }); + describe('engine.addPipe()', () => { + + it('returns unsubscribe function that removes the middleware', async () => { + + const engine = new HookEngine(); + const middleware = vi.fn((next: any, ..._args: any[]) => next()); + + const unsub = engine.addPipe('test', middleware); + + await engine.pipe('test', async () => 'core-result'); + expect(middleware).toHaveBeenCalledOnce(); + + unsub(); + + await engine.pipe('test', async () => 'core-result'); + expect(middleware).toHaveBeenCalledOnce(); + }); + + it('rejects invalid name', () => { + + const engine = new HookEngine(); + + expect(() => engine.addPipe(null as any, (next: any, ..._args: any[]) => next())).to.throw(); + expect(() => engine.addPipe(123 as any, (next: any, ..._args: any[]) => next())).to.throw(); + }); + + it('rejects invalid callback', () => { + + const engine = new HookEngine(); + + expect(() => engine.addPipe('test', null as any)).to.throw(); + expect(() => engine.addPipe('test', 'notAFunction' as any)).to.throw(); + }); + + it('enforces register() strict mode, naming addPipe in the error', () => { + + const engine = new HookEngine(); + engine.register('validHook'); + + expect(() => engine.addPipe('invalidHook', (next: any, ..._args: any[]) => next())).to.throw( + /Hook "invalidHook" is not registered/ + ); + expect(() => engine.addPipe('invalidHook', (next: any, ..._args: any[]) => next())).to.throw( + /using addPipe\(\)/ + ); + }); + + describe('type-level: PipeCallback registers cast-free (issue #147)', () => { + + interface Lifecycle { + execute(opts: { url: string }): Promise; + } + + it('accepts a PipeCallback-typed const with no casts', async () => { + + const engine = new HookEngine(); + + const middleware: PipeCallback<[opts: { url: string }], string> = async (next, opts) => { + + expect(opts.url).to.equal('/api/users'); + return next(); + }; + + engine.addPipe('execute', middleware); + + const result = await engine.pipe('execute', async () => 'core-result', { url: '/api/users' }); + + expect(result).to.equal('core-result'); + }); + + it('accepts an inline callback with inferred params', async () => { + + const engine = new HookEngine(); + + engine.addPipe('execute', async (next, opts) => { + + expect(opts.url).to.equal('/api/users'); + return next(); + }); + + const result = await engine.pipe('execute', async () => 'core-result', { url: '/api/users' }); + + expect(result).to.equal('core-result'); + }); + }); + }); + + describe('engine.pipe()', () => { + + interface Lifecycle { + execute(): Promise; + executeWithOpts(opts: { retried: boolean }): Promise; + test(a: string, b: number): Promise; + } + + it('executes middleware in onion order by priority (lower = outermost)', async () => { + + const engine = new HookEngine(); + const order: string[] = []; + + engine.addPipe('execute', async (next) => { + + order.push('enter:outer'); + const result = await next(); + order.push('exit:outer'); + return result; + }, { priority: -20 }); + + engine.addPipe('execute', async (next) => { + + order.push('enter:inner'); + const result = await next(); + order.push('exit:inner'); + return result; + }, { priority: -10 }); + + const result = await engine.pipe('execute', async () => { + + order.push('core'); + return 'core-result'; + }); + + expect(order).to.deep.equal([ + 'enter:outer', 'enter:inner', 'core', 'exit:inner', 'exit:outer' + ]); + expect(result).to.equal('core-result'); + }); + + it('passes next fn + spread args + PipeContext as ctx', async () => { + + const engine = new HookEngine(); + let receivedArgs: unknown[] = []; + + engine.addPipe('test', async (next, a, b, ctx) => { + + receivedArgs = [a, b, ctx]; + return next(); + }); + + await engine.pipe('test', async () => 'core-result', 'hello', 42); + + const ctx = receivedArgs.pop(); + expect(receivedArgs).to.deep.equal(['hello', 42]); + expect(ctx).to.be.instanceOf(PipeContext); + }); + + it('short-circuits when middleware does not call next()', async () => { + + const engine = new HookEngine(); + const coreFn = vi.fn(async () => 'core-result'); + + engine.addPipe('execute', async () => 'cached-result'); + + const result = await engine.pipe('execute', coreFn); + + expect(result).to.equal('cached-result'); + expect(coreFn).not.toHaveBeenCalled(); + }); + + it('propagates the result from next() back through the chain', async () => { + + const engine = new HookEngine(); + + engine.addPipe('execute', async (next) => { + + const result = await next(); + return `wrapped:${result}`; + }); + + const result = await engine.pipe('execute', async () => 'core-result'); + + expect(result).to.equal('wrapped:core-result'); + }); + + it('ctx.args() replacement reaches inner layers', async () => { + + const engine = new HookEngine(); + let innerSawOpts: { retried: boolean } | undefined; + + engine.addPipe('executeWithOpts', async (next, opts, ctx) => { + + ctx.args({ ...opts, retried: true }); + return next(); + }, { priority: -10 }); + + engine.addPipe('executeWithOpts', async (next, opts) => { + + innerSawOpts = opts; + return next(); + }, { priority: 0 }); + + await engine.pipe('executeWithOpts', async () => 'core-result', { retried: false }); + + expect(innerSawOpts).to.deep.equal({ retried: true }); + }); + + it('ctx.fail() aborts the chain', async () => { + + const engine = new HookEngine(); + + engine.addPipe('execute', async (_next, ctx) => { + + return ctx.fail('Rate limit exceeded'); + }); + + const [, err] = await attempt(() => engine.pipe('execute', async () => 'core-result')); + + expect(isHookError(err)).to.be.true; + expect(err).to.have.property('message').that.includes('Rate limit exceeded'); + }); + + it('ctx.removeHook() removes the middleware from future pipes', async () => { + + const engine = new HookEngine(); + let callCount = 0; + + engine.addPipe('execute', async (next, ctx) => { + + callCount++; + + if (callCount >= 2) ctx.removeHook(); + + return next(); + }); + + await engine.pipe('execute', async () => 'a'); + await engine.pipe('execute', async () => 'b'); + await engine.pipe('execute', async () => 'c'); + + expect(callCount).to.equal(2); + }); + + it('supports once option', async () => { + + const engine = new HookEngine(); + const middleware = vi.fn>((next) => next()); + + engine.addPipe('execute', middleware, { once: true }); + + await engine.pipe('execute', async () => 'a'); + await engine.pipe('execute', async () => 'b'); + + expect(middleware).toHaveBeenCalledOnce(); + }); + + it('supports times option', async () => { + + const engine = new HookEngine(); + const middleware = vi.fn>((next) => next()); + + engine.addPipe('execute', middleware, { times: 2 }); + + await engine.pipe('execute', async () => 'a'); + await engine.pipe('execute', async () => 'b'); + await engine.pipe('execute', async () => 'c'); + + expect(middleware).toHaveBeenCalledTimes(2); + }); + + it('supports ignoreOnFail option, falling through to next()', async () => { + + const engine = new HookEngine(); + const coreFn = vi.fn(async () => 'core-result'); + + engine.addPipe('execute', async () => { throw new Error('boom'); }, { ignoreOnFail: true }); + + const result = await engine.pipe('execute', coreFn); + + expect(result).to.equal('core-result'); + expect(coreFn).toHaveBeenCalledOnce(); + }); + }); + + describe('engine.pipeSync()', () => { + + interface Lifecycle { + execute(): string; + executeWithOpts(opts: { retried: boolean }): string; + test(a: string, b: number): string; + } + + it('executes middleware in onion order by priority (lower = outermost)', () => { + + const engine = new HookEngine(); + const order: string[] = []; + + engine.addPipe('execute', (next) => { + + order.push('enter:outer'); + const result = next(); + order.push('exit:outer'); + return result; + }, { priority: -20 }); + + engine.addPipe('execute', (next) => { + + order.push('enter:inner'); + const result = next(); + order.push('exit:inner'); + return result; + }, { priority: -10 }); + + const result = engine.pipeSync('execute', () => { + + order.push('core'); + return 'core-result'; + }); + + expect(order).to.deep.equal([ + 'enter:outer', 'enter:inner', 'core', 'exit:inner', 'exit:outer' + ]); + expect(result).to.equal('core-result'); + }); + + it('passes next fn + spread args + PipeContext as ctx', () => { + + const engine = new HookEngine(); + let receivedArgs: unknown[] = []; + + engine.addPipe('test', (next, a, b, ctx) => { + + receivedArgs = [a, b, ctx]; + return next(); + }); + + engine.pipeSync('test', () => 'core-result', 'hello', 42); + + const ctx = receivedArgs.pop(); + expect(receivedArgs).to.deep.equal(['hello', 42]); + expect(ctx).to.be.instanceOf(PipeContext); + }); + + it('short-circuits when middleware does not call next()', () => { + + const engine = new HookEngine(); + const coreFn = vi.fn(() => 'core-result'); + + engine.addPipe('execute', () => 'cached-result'); + + const result = engine.pipeSync('execute', coreFn); + + expect(result).to.equal('cached-result'); + expect(coreFn).not.toHaveBeenCalled(); + }); + + it('propagates the result from next() back through the chain', () => { + + const engine = new HookEngine(); + + engine.addPipe('execute', (next) => { + + const result = next(); + return `wrapped:${result}`; + }); + + const result = engine.pipeSync('execute', () => 'core-result'); + + expect(result).to.equal('wrapped:core-result'); + }); + + it('ctx.args() replacement reaches inner layers', () => { + + const engine = new HookEngine(); + let innerSawOpts: { retried: boolean } | undefined; + + engine.addPipe('executeWithOpts', (next, opts, ctx) => { + + ctx.args({ ...opts, retried: true }); + return next(); + }, { priority: -10 }); + + engine.addPipe('executeWithOpts', (next, opts) => { + + innerSawOpts = opts; + return next(); + }, { priority: 0 }); + + engine.pipeSync('executeWithOpts', () => 'core-result', { retried: false }); + + expect(innerSawOpts).to.deep.equal({ retried: true }); + }); + + it('ctx.fail() aborts the chain', () => { + + const engine = new HookEngine(); + + engine.addPipe('execute', (_next, ctx) => { + + return ctx.fail('Rate limit exceeded'); + }); + + const [, err] = attemptSync(() => engine.pipeSync('execute', () => 'core-result')); + + expect(isHookError(err)).to.be.true; + expect(err).to.have.property('message').that.includes('Rate limit exceeded'); + }); + + it('ctx.removeHook() removes the middleware from future pipes', () => { + + const engine = new HookEngine(); + let callCount = 0; + + engine.addPipe('execute', (next, ctx) => { + + callCount++; + + if (callCount >= 2) ctx.removeHook(); + + return next(); + }); + + engine.pipeSync('execute', () => 'a'); + engine.pipeSync('execute', () => 'b'); + engine.pipeSync('execute', () => 'c'); + + expect(callCount).to.equal(2); + }); + + it('supports once option', () => { + + const engine = new HookEngine(); + const middleware = vi.fn>((next) => next()); + + engine.addPipe('execute', middleware, { once: true }); + + engine.pipeSync('execute', () => 'a'); + engine.pipeSync('execute', () => 'b'); + + expect(middleware).toHaveBeenCalledOnce(); + }); + + it('supports times option', () => { + + const engine = new HookEngine(); + const middleware = vi.fn>((next) => next()); + + engine.addPipe('execute', middleware, { times: 2 }); + + engine.pipeSync('execute', () => 'a'); + engine.pipeSync('execute', () => 'b'); + engine.pipeSync('execute', () => 'c'); + + expect(middleware).toHaveBeenCalledTimes(2); + }); + + it('supports ignoreOnFail option, falling through to next()', () => { + + const engine = new HookEngine(); + const coreFn = vi.fn(() => 'core-result'); + + engine.addPipe('execute', () => { throw new Error('boom'); }, { ignoreOnFail: true }); + + const result = engine.pipeSync('execute', coreFn); + + expect(result).to.equal('core-result'); + expect(coreFn).toHaveBeenCalledOnce(); + }); + }); + describe('engine.clear()', () => { it('removes all hooks', async () => {