diff --git a/.changeset/typed-idempotent-operations.md b/.changeset/typed-idempotent-operations.md new file mode 100644 index 00000000..9c55df32 --- /dev/null +++ b/.changeset/typed-idempotent-operations.md @@ -0,0 +1,11 @@ +--- +"@cleverbrush/server": minor +"@cleverbrush/client": minor +"@cleverbrush/server-openapi": minor +--- + +Add contract-declared mutation replay with typed async request preparation, +explicit authorization scopes, and consistent endpoint error policies. Generate +client keys and coordinate HTTP retries/batching from the contract without +changing ordinary client calls. Document replay headers and framework errors in +OpenAPI while preserving domain response schemas. diff --git a/libs/client/README.md b/libs/client/README.md index 9591c604..6484087a 100644 --- a/libs/client/README.md +++ b/libs/client/README.md @@ -243,6 +243,27 @@ const client = createClient(api, { ## Resilience Middlewares +### Contract-declared mutation retries + +For endpoints declared with `.idempotent()`, the typed client generates one +`X-Idempotency-Key` per call before middleware execution. `retry()` recognizes +the contract and reuses the key/body for every HTTP attempt; other POSTs retain +their existing retry policy. Ordinary endpoint calls require no extra options: + +```ts +await client.items.create({ body: item }); +``` + +Explicit middleware or per-call `retry.methods` takes precedence, and +`retry: { limit: 0 }` disables automatic retries. `batching()` sends these calls +directly so their individual timeout signals still reach the transport. +Each new client invocation gets a fresh key, even when its input is identical. +This does not track form submissions or user-initiated retries across calls. + +The server requires an authorization scope and replays responses within a +bounded, process-local store. This does not provide durable exactly-once +execution across replicas or restarts. + ### Retry — `@cleverbrush/client/retry` ```ts diff --git a/libs/client/src/client.ts b/libs/client/src/client.ts index bb0bd4f7..ede1b48d 100644 --- a/libs/client/src/client.ts +++ b/libs/client/src/client.ts @@ -196,6 +196,22 @@ export function createClient( ...args?.headers }; + if (meta.idempotent) { + const existing = Object.keys(reqHeaders).filter( + name => name.toLowerCase() === 'x-idempotency-key' + ); + const key = + new Headers(args?.headers).get('x-idempotency-key') ?? + new Headers(extraHeaders).get('x-idempotency-key') ?? + crypto.randomUUID(); + if (typeof key !== 'string' || key.length === 0 || key.length > 256) + throw new TypeError( + 'Idempotency key must contain 1 to 256 characters' + ); + for (const name of existing) delete reqHeaders[name]; + reqHeaders['x-idempotency-key'] = key; + } + const token = getToken?.(); if (token && meta.authRoles !== null) { reqHeaders['Authorization'] = `Bearer ${token}`; @@ -255,26 +271,14 @@ export function createClient( return { url, method, headers: reqHeaders, body }; } - // The actual fetch logic, shared by every endpoint proxy method. - async function execute( + // Every response mode carries the same contract metadata and retry options. + function attachRequestOptions( ep: any, args: any, + init: RequestInit, groupName?: string, endpointName?: string - ): Promise { - const { - url, - method, - headers: reqHeaders, - body - } = buildRequest(ep, args); - - const init: RequestInit = { - method, - headers: reqHeaders, - body - }; - + ): void { // Attach per-call middleware overrides if provided. const perCallOptions: Record = {}; if (args?.retry !== undefined) perCallOptions.retry = args.retry; @@ -330,13 +334,37 @@ export function createClient( headers: args?.headers ?? ({} as Record), operationId: meta.operationId ?? null, tags: meta.tags ?? [], - cacheTags: meta.cacheTags ?? [] + cacheTags: meta.cacheTags ?? [], + idempotent: meta.idempotent ?? false }; if (!(init as any).__endpointMeta) { (init as any).__endpointMeta = epMeta; } } + } + + // The actual fetch logic, shared by every endpoint proxy method. + async function execute( + ep: any, + args: any, + groupName?: string, + endpointName?: string + ): Promise { + const { + url, + method, + headers: reqHeaders, + body + } = buildRequest(ep, args); + + const init: RequestInit = { + method, + headers: reqHeaders, + body + }; + + attachRequestOptions(ep, args, init, groupName, endpointName); // -- beforeRequest hooks -- await runBeforeRequest(hooks, url, init); @@ -392,7 +420,12 @@ export function createClient( } // Streaming fetch — yields newline-delimited chunks (e.g. NDJSON). - async function* streamLines(ep: any, args: any): AsyncIterable { + async function* streamLines( + ep: any, + args: any, + groupName?: string, + endpointName?: string + ): AsyncIterable { const { url, method, @@ -407,6 +440,8 @@ export function createClient( signal: args?.signal }; + attachRequestOptions(ep, args, init, groupName, endpointName); + // -- beforeRequest hooks -- await runBeforeRequest(hooks, url, init); @@ -489,7 +524,8 @@ export function createClient( // Regular HTTP endpoints return a callable with .stream() and .file() const call = (args?: any) => execute(ep, args, groupName, endpointName); - call.stream = (args?: any) => streamLines(ep, args); + call.stream = (args?: any) => + streamLines(ep, args, groupName, endpointName); call.file = async (args?: any): Promise => { const { url, @@ -502,6 +538,13 @@ export function createClient( headers: reqHeaders, body }; + attachRequestOptions( + ep, + args, + init, + groupName, + endpointName + ); await runBeforeRequest(hooks, url, init); const response = await composedFetch(url, init); if (!response.ok) { diff --git a/libs/client/src/contractIdempotency.test.ts b/libs/client/src/contractIdempotency.test.ts new file mode 100644 index 00000000..167e7bc9 --- /dev/null +++ b/libs/client/src/contractIdempotency.test.ts @@ -0,0 +1,119 @@ +import { number, object, string } from '@cleverbrush/schema'; +import { defineApi, endpoint } from '@cleverbrush/server/contract'; +import { expect, it, vi } from 'vitest'; +import { batching } from './batching.js'; +import { createClient } from './client.js'; +import { retry } from './retry.js'; + +const api = defineApi({ + items: { + create: endpoint + .post('/items') + .idempotent() + .body(object({ amount: number() })), + other: endpoint.post('/other') + } +}); + +it('retries contract-declared mutations with the original key/body and bypasses batching', async () => { + const requests: { url: string; key: string | null; body: unknown }[] = []; + const fetch = vi.fn(async (url, init) => { + requests.push({ + url: String(url), + key: new Headers(init.headers).get('x-idempotency-key'), + body: init.body + }); + if (requests.length === 1) throw new TypeError('Lost response'); + return Response.json({ ok: true }); + }); + const client = createClient(api, { + baseUrl: 'https://example.test', + fetch, + middlewares: [retry({ delay: () => 0 }), batching({ windowMs: 1 })] + }); + await client.items.create({ body: { amount: 1 } }); + expect(requests).toHaveLength(2); + expect(requests[1]).toEqual(requests[0]); + expect(requests[0].key).toBeTruthy(); + await Promise.all([ + client.items.create({ body: { amount: 1 } }), + client.items.create({ body: { amount: 1 } }) + ]); + expect(requests).toHaveLength(4); + expect(requests.every(value => value.key)).toBe(true); + expect(new Set(requests.map(value => value.key)).size).toBe(3); + expect(requests.every(value => !value.url.includes('__batch'))).toBe(true); + const failing = vi.fn().mockRejectedValue(new TypeError('Lost')); + const ordinary = createClient(api, { + fetch: failing, + middlewares: [retry({ delay: () => 0 })] + }); + await expect(ordinary.items.other()).rejects.toThrow(); + expect(failing).toHaveBeenCalledOnce(); +}); + +it('preserves supplied headers and respects retry opt-out', async () => { + const fetch = vi.fn().mockRejectedValue(new TypeError('Lost')); + const client = createClient(api, { + fetch, + headers: { 'X-Idempotency-Key': 'manual' }, + middlewares: [retry({ delay: () => 0 })] + }); + await expect( + client.items.create({ body: { amount: 1 }, retry: { limit: 0 } }) + ).rejects.toThrow(); + expect(fetch).toHaveBeenCalledOnce(); + expect( + new Headers(fetch.mock.calls[0][1].headers).get('x-idempotency-key') + ).toBe('manual'); +}); + +it('preserves the contract retry policy when requesting a binary response', async () => { + const fetch = vi + .fn() + .mockRejectedValueOnce(new TypeError('Lost response')) + .mockImplementation(async () => new Response('receipt')); + const client = createClient(api, { + fetch, + middlewares: [retry({ delay: () => 0 }), batching({ windowMs: 1 })] + }); + const result = await client.items.create.file({ + body: { amount: 1 } + }); + expect(await result.text()).toBe('receipt'); + expect(fetch).toHaveBeenCalledTimes(2); + const keys = fetch.mock.calls.map(([, init]) => + new Headers(init.headers).get('x-idempotency-key') + ); + expect(keys[0]).toBeTruthy(); + expect(keys[1]).toBe(keys[0]); +}); + +it('honors explicit retry method restrictions and removes duplicate header casing', async () => { + const fetch = vi.fn().mockRejectedValue(new TypeError('Lost response')); + const customHeadersApi = defineApi({ + items: { + create: api.items.create.headers( + object({ 'x-idempotency-key': string() }) + ) + } + }); + const client = createClient(customHeadersApi, { + fetch, + headers: { + 'X-Idempotency-Key': 'default', + 'x-idempotency-key': 'duplicate' + }, + middlewares: [retry({ methods: ['GET'], delay: () => 0 })] + }); + await expect( + client.items.create({ + body: { amount: 1 }, + headers: { 'x-idempotency-key': 'explicit' } + }) + ).rejects.toThrow(); + expect(fetch).toHaveBeenCalledOnce(); + expect( + new Headers(fetch.mock.calls[0][1].headers).get('x-idempotency-key') + ).toBe('explicit'); +}); diff --git a/libs/client/src/middleware.ts b/libs/client/src/middleware.ts index d9b3871a..a2d91a97 100644 --- a/libs/client/src/middleware.ts +++ b/libs/client/src/middleware.ts @@ -130,6 +130,8 @@ export function getPerCallOptions( * Used by `throttlingCache` for cache-invalidation callbacks. */ export interface EndpointMeta { + /** Contract declares server-side mutation replay. */ + idempotent?: boolean; /** Contract group name, e.g. `"todos"`. */ group: string; /** Endpoint name within the group, e.g. `"update"`. */ @@ -185,3 +187,13 @@ export interface EndpointMeta { /** Request headers from the call, e.g. `{ 'x-request-id': 'abc' }`. */ headers: Readonly>; } + +/** @internal Only contract-declared, keyed mutations are automatically retryable. */ +export function isIdempotentRequest(init: RequestInit): boolean { + const meta = (init as RequestInit & { __endpointMeta?: EndpointMeta }) + .__endpointMeta; + return ( + meta?.idempotent === true && + !!new Headers(init.headers).get('x-idempotency-key') + ); +} diff --git a/libs/client/src/middlewares/batching.ts b/libs/client/src/middlewares/batching.ts index 2c12e971..a5a9a6e5 100644 --- a/libs/client/src/middlewares/batching.ts +++ b/libs/client/src/middlewares/batching.ts @@ -29,7 +29,11 @@ * @module */ -import type { FetchLike, Middleware } from '../middleware.js'; +import { + type FetchLike, + isIdempotentRequest, + type Middleware +} from '../middleware.js'; // --------------------------------------------------------------------------- // Types @@ -293,7 +297,7 @@ export function batching(options: BatchingOptions = {}): Middleware { } // Honour the user-provided skip predicate. - if (skip?.(url, init)) { + if (isIdempotentRequest(init) || skip?.(url, init)) { return next(url, init); } diff --git a/libs/client/src/middlewares/retry.ts b/libs/client/src/middlewares/retry.ts index 3f96563a..feaed55a 100644 --- a/libs/client/src/middlewares/retry.ts +++ b/libs/client/src/middlewares/retry.ts @@ -20,7 +20,7 @@ import { ApiError } from '../errors.js'; import type { Middleware } from '../middleware.js'; -import { getPerCallOptions } from '../middleware.js'; +import { getPerCallOptions, isIdempotentRequest } from '../middleware.js'; // --------------------------------------------------------------------------- // Types @@ -182,7 +182,11 @@ export function retry(options: RetryOptions = {}): Middleware { const method = (init.method ?? 'GET').toUpperCase(); // Non-retryable methods go straight through. - if (!methodSet.has(method)) { + const allowed = perCall?.methods + ? perCall.methods.some(value => value.toUpperCase() === method) + : methodSet.has(method) || + (options.methods === undefined && isIdempotentRequest(init)); + if (!allowed) { return next(url, init); } diff --git a/libs/server-openapi/src/generateOpenApiSpec.ts b/libs/server-openapi/src/generateOpenApiSpec.ts index eb3f6f9a..45c0f297 100644 --- a/libs/server-openapi/src/generateOpenApiSpec.ts +++ b/libs/server-openapi/src/generateOpenApiSpec.ts @@ -380,6 +380,27 @@ function buildResponses( } } + if (meta.idempotent) { + for (const [status, description] of [ + [400, 'Invalid idempotency key'], + [ + 409, + 'Previous response cannot be replayed; verify the operation outcome' + ], + [503, 'Idempotency capacity reached'] + ] as const) { + const response = (result[status] ?? { description }) as Record< + string, + any + >; + response.content = { + ...response.content, + 'application/problem+json': { schema: PROBLEM_DETAILS_SCHEMA } + }; + result[status] = response; + } + } + // Multiple content types — augment each response's content map with extra // MIME types from .produces(). producesFile already handled above (binary wins). if (meta.produces && !meta.producesFile) { @@ -564,6 +585,24 @@ function buildOperation( } } + if ( + meta.idempotent && + !parameters.some( + parameter => + parameter.in === 'header' && + String(parameter.name).toLowerCase() === 'x-idempotency-key' + ) + ) { + parameters.push( + buildParameterObject( + 'X-Idempotency-Key', + 'header', + { type: 'string', minLength: 1, maxLength: 256 }, + false, + 'Reuse for retries of the same mutation. Replay is bounded and process-local.' + ) + ); + } if (parameters.length > 0) operation['parameters'] = parameters; // Request body diff --git a/libs/server-openapi/src/idempotency.test.ts b/libs/server-openapi/src/idempotency.test.ts new file mode 100644 index 00000000..2ac814f1 --- /dev/null +++ b/libs/server-openapi/src/idempotency.test.ts @@ -0,0 +1,53 @@ +import { object, string } from '@cleverbrush/schema'; +import { endpoint } from '@cleverbrush/server'; +import { expect, it } from 'vitest'; +import { generateOpenApiSpec } from './generateOpenApiSpec.js'; + +it('documents optional replay headers and framework errors without replacing domain errors', () => { + const ep = endpoint + .post('/items') + .idempotent() + .responses({ 201: string(), 409: object({ message: string() }) }); + const spec = generateOpenApiSpec({ + info: { title: 'Items', version: '1' }, + registrations: [{ endpoint: ep.introspect(), handler: () => undefined }] + }) as any; + const operation = spec.paths['/items'].post; + expect(operation.parameters).toContainEqual( + expect.objectContaining({ + in: 'header', + name: 'X-Idempotency-Key' + }) + ); + expect(operation.parameters[0].required).not.toBe(true); + expect(operation.responses['409'].content).toHaveProperty( + 'application/json' + ); + expect(operation.responses['409'].content).toHaveProperty( + 'application/problem+json' + ); + expect(operation.responses['503'].content).toHaveProperty( + 'application/problem+json' + ); +}); + +it('does not duplicate a manually declared key header', () => { + const ep = endpoint + .post('/items') + .idempotent() + .headers( + object({ + 'x-idempotency-key': string().optional(), + 'x-request-id': string().optional() + }) + ); + const spec = generateOpenApiSpec({ + info: { title: 'Items', version: '1' }, + registrations: [{ endpoint: ep.introspect(), handler: () => undefined }] + }) as any; + expect( + spec.paths['/items'].post.parameters.filter( + (value: any) => value.name.toLowerCase() === 'x-idempotency-key' + ) + ).toHaveLength(1); +}); diff --git a/libs/server/README.md b/libs/server/README.md index 7884bc30..b13f3c80 100644 --- a/libs/server/README.md +++ b/libs/server/README.md @@ -916,6 +916,51 @@ apps/ ## Response replay and caching +Declare mutation replay in the shared contract with `.idempotent()`. Bind an +explicit scope on the server; registration fails at startup when it is missing. +The optional `prepare` callback receives validated request data and injected +services and returns the request passed to scope resolution and the handler. +Use it to authorize resource access and resolve request defaults before replay. + +```ts +const createItem = endpoint.post('/items') + .idempotent() + .body(CreateItemSchema) + .authorize(UserPrincipal) + .inject({ db: DbToken }) + .responses({ 201: ItemSchema, 403: MessageSchema }); + +server.handle(createItem, createItemHandler, { + prepare: async (request, { db }) => { + const workspace = await requireWorkspaceAccess(db, request.principal, + request.body.workspaceId); + return { ...request, body: { ...request.body, workspaceId: workspace.id } }; + }, + idempotency: { + scope: ({ principal, body }) => [principal.userId, body.workspaceId] + }, + errors: itemErrors +}); +``` + +The same options work in `mapHandlers` and `implement(api).group(...).withHandlers`. +Authentication, validation, preparation and scope resolution run before each +replay, including batch subrequests. Preparation/scope exceptions use the bound +error policy without reserving a key. Handler exceptions translated by that +policy become replayable responses. Validation, DI and serialization errors +retain their normal Framework handling; standalone `withErrors` still wraps +only its handler. + +`X-Idempotency-Key` is optional, case-insensitive, and must contain 1–256 +characters when present. It is transport metadata, so no `.headers()` schema is +needed. OpenAPI includes this header and the 400/409/503 Problem Details responses. +When enabling cross-origin browser access, explicitly include this header in +your CORS `allowedHeaders`. Each server endpoint retains its own bounded store +with the defaults below; `idempotency` options can override its limits. + +### Low-level replay and caching + + `idempotency({ scope })` coalesces concurrent mutations and replays their completed responses within one middleware instance. Install it after authentication and authorization. Derive `scope(ctx)` from verified identity/tenant context; returning diff --git a/libs/server/src/Endpoint.ts b/libs/server/src/Endpoint.ts index 04b9cd67..db07f5f5 100644 --- a/libs/server/src/Endpoint.ts +++ b/libs/server/src/Endpoint.ts @@ -1,4 +1,5 @@ // biome-ignore-all lint/suspicious/useAdjacentOverloadSignatures: each method in ScopedEndpointFactoryMethods and EndpointFactory has a single signature; they are separate methods, not overloads + import type { InferType, ObjectSchemaBuilder, @@ -20,6 +21,10 @@ import type { } from './ActionResult.js'; import type { CacheTagDefinition } from './CacheTag.js'; import { createCacheTagTree, serializeTag } from './CacheTag.js'; +import type { + EndpointOptions, + RuntimeEndpointOptions +} from './EndpointOptions.js'; import type { RequestContext } from './RequestContext.js'; import { createSubscription, @@ -270,7 +275,7 @@ type AnySubscriptionBuilder = SubscriptionBuilder< */ export type HandlerEntry = | Handler - | { handler: Handler; middlewares?: Middleware[] }; + | ({ handler: Handler } & EndpointOptions); /** * A compile-time complete mapping from an endpoint group structure to @@ -304,6 +309,10 @@ export interface HandlerMapping { endpoint: AnyEndpoint; handler: (...args: any[]) => any; middlewares?: Middleware[]; + handlerErrorsMapped?: boolean; + prepare?: RuntimeEndpointOptions['prepare']; + idempotency?: RuntimeEndpointOptions['idempotency']; + errors?: RuntimeEndpointOptions['errors']; }>; /** @internal */ readonly _subscriptions: ReadonlyArray<{ @@ -372,7 +381,19 @@ export function mapHandlers< entries.push({ endpoint: ep as AnyEndpoint, handler, - middlewares + middlewares, + handlerErrorsMapped: + typeof entry === 'function' + ? false + : entry.handlerErrorsMapped, + prepare: + typeof entry === 'function' ? undefined : entry.prepare, + idempotency: + typeof entry === 'function' + ? undefined + : entry.idempotency, + errors: + typeof entry === 'function' ? undefined : entry.errors }); } } @@ -597,6 +618,8 @@ export interface EndpointMetadata { * key computation for the client middleware. */ readonly cacheTags: readonly CacheTagDefinition[]; + /** Opt-in to bounded mutation response replay. */ + readonly idempotent?: boolean; } /** @@ -760,6 +783,7 @@ export class EndpointBuilder< readonly #callbacks: Record | null; readonly #fileUpload: UploadConfiguration | null; readonly #cacheTags: readonly CacheTagDefinition[]; + readonly #idempotent: boolean; constructor( method: string, @@ -825,7 +849,8 @@ export class EndpointBuilder< links: Record | null = null, callbacks: Record | null = null, fileUpload: UploadConfiguration | null = null, - cacheTags: readonly CacheTagDefinition[] = [] + cacheTags: readonly CacheTagDefinition[] = [], + idempotent = false ) { validateUploadConfiguration(fileUpload, bodySchema); this.#method = method; @@ -853,6 +878,60 @@ export class EndpointBuilder< this.#callbacks = callbacks; this.#fileUpload = fileUpload; this.#cacheTags = cacheTags; + this.#idempotent = idempotent; + } + + /** + * Enable optional X-Idempotency-Key replay for this mutation. The server + * registration must provide an explicit authorization scope. Retention is + * process-local; reusing a key asserts that the input is unchanged. + */ + idempotent(): EndpointBuilder< + TParams, + TBody, + TQuery, + THeaders, + TServices, + TPrincipal, + TRoles, + TResponse, + TResponses, + TUpload + > { + if ( + !['POST', 'PUT', 'PATCH', 'DELETE'].includes( + this.#method.toUpperCase() + ) + ) + throw new TypeError('Idempotency requires a mutation endpoint'); + return new EndpointBuilder( + this.#method, + this.#basePath, + this.#pathTemplate, + this.#bodySchema, + this.#querySchema, + this.#headerSchema, + this.#serviceSchemas, + this.#authRoles, + this.#summary, + this.#description, + this.#tags, + this.#operationId, + this.#deprecated, + this.#responseSchema, + this.#responsesSchemas, + this.#example, + this.#examples, + this.#producesFile, + this.#produces, + this.#responseHeaderSchema, + this.#externalDocs, + this.#links, + this.#callbacks, + this.#fileUpload, + this.#cacheTags, + true + ); } /** Define the request body schema. Validation failures return 422 Problem Details. */ @@ -895,7 +974,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -941,7 +1021,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -987,7 +1068,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1033,7 +1115,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1130,7 +1213,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1178,7 +1262,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1249,7 +1334,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1319,7 +1405,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1363,7 +1450,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1407,7 +1495,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1451,7 +1540,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1495,7 +1585,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1537,7 +1628,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1588,7 +1680,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1642,7 +1735,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1695,7 +1789,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1810,7 +1905,8 @@ export class EndpointBuilder< maxFileCount: config?.maxFileCount ?? 10, schema }, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1867,7 +1963,8 @@ export class EndpointBuilder< links: this.#links, callbacks: this.#callbacks, fileUpload: this.#fileUpload, - cacheTags: this.#cacheTags + cacheTags: this.#cacheTags, + idempotent: this.#idempotent }; } @@ -1934,7 +2031,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -1997,7 +2095,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -2049,7 +2148,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -2112,7 +2212,8 @@ export class EndpointBuilder< defs as Record, this.#callbacks, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -2177,7 +2278,8 @@ export class EndpointBuilder< this.#links, defs as Record, this.#fileUpload, - this.#cacheTags + this.#cacheTags, + this.#idempotent ); } @@ -2336,7 +2438,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - [...this.#cacheTags, { name, properties: {} }] + [...this.#cacheTags, { name, properties: {} }], + this.#idempotent ); } @@ -2385,7 +2488,8 @@ export class EndpointBuilder< this.#links, this.#callbacks, this.#fileUpload, - [...this.#cacheTags, definition] + [...this.#cacheTags, definition], + this.#idempotent ); } } diff --git a/libs/server/src/EndpointIdempotency.test-d.ts b/libs/server/src/EndpointIdempotency.test-d.ts new file mode 100644 index 00000000..c0db204e --- /dev/null +++ b/libs/server/src/EndpointIdempotency.test-d.ts @@ -0,0 +1,43 @@ +import { number, object, string } from '@cleverbrush/schema'; +import { expectTypeOf } from 'vitest'; +import { + ActionResult, + createServer, + defineApi, + type EndpointOptions, + endpoint, + implement +} from './index.js'; + +const Db = object({ tenant: string() }); +const api = defineApi({ + items: { + create: endpoint + .post('/items') + .idempotent() + .body(object({ amount: number(), tenant: string().optional() })) + .responses({ 201: string() }) + } +}); +const scope = implement(api).group('items', { inject: { db: Db } }); +const options: EndpointOptions = { + prepare: async (request, { db }) => { + expectTypeOf(request.body.amount).toEqualTypeOf(); + expectTypeOf(db.tenant).toEqualTypeOf(); + return { ...request, body: { ...request.body, tenant: db.tenant } }; + }, + idempotency: { scope: async ({ body }, { db }) => [db.tenant, body.amount] } +}; +createServer().handle( + scope.endpoints.create, + () => ActionResult.created('ok'), + options +); +scope.withHandlers({ + create: { ...options, handler: () => ActionResult.created('ok') } +}); +const invalid: EndpointOptions = { + // @ts-expect-error preparation must preserve the typed action context + prepare: () => ({ body: { amount: 'invalid' } }) +}; +void invalid; diff --git a/libs/server/src/EndpointIdempotency.test.ts b/libs/server/src/EndpointIdempotency.test.ts new file mode 100644 index 00000000..24cbd8ac --- /dev/null +++ b/libs/server/src/EndpointIdempotency.test.ts @@ -0,0 +1,189 @@ +import { number, object, string } from '@cleverbrush/schema'; +import { afterEach, expect, it, vi } from 'vitest'; +import { + ActionResult, + createServer, + defineApi, + endpoint, + errorMap, + implement, + type Server +} from './index.js'; + +const servers: Server[] = []; +afterEach(async () => { + await Promise.all(servers.splice(0).map(server => server.close())); +}); +const Body = object({ tenant: string().optional(), amount: number() }); +const Message = object({ message: string() }); +const Db = object({ tenant: string() }); +class Denied extends Error {} + +async function fixture( + limits: { + maxEntries?: number; + maxResponseBytes?: number; + ttl?: number; + } = {} +) { + let allowed = true; + const api = defineApi({ + items: { + create: endpoint + .post('/items') + .idempotent() + .body(Body) + .responses({ 201: Body, 403: Message }) + } + }); + const scope = implement(api).group('items', { inject: { db: Db } }); + const prepare = vi.fn(async (request, { db }) => { + if (!allowed) throw new Denied(); + return { + ...request, + body: { ...request.body, tenant: request.body.tenant ?? db.tenant } + }; + }); + const handler = vi.fn(async ({ body }) => { + await new Promise(resolve => setTimeout(resolve, 5)); + return ActionResult.created(body, '/items/1'); + }); + const module = scope.withHandlers({ + create: { + prepare, + handler, + idempotency: { + ...limits, + scope: async ({ body }) => ['verified-user', body.tenant!] + }, + errors: errorMap().on(Denied, () => + ActionResult.forbidden({ message: 'Denied' }) + ) + } + }); + const server = await createServer() + .services(services => + services.addSingleton(Db, () => ({ tenant: 'default' })) + ) + .useBatching() + .handleAll(implement(api).use(module).complete()) + .listen(0, '127.0.0.1'); + servers.push(server); + const url = `http://127.0.0.1:${server.address!.port}`; + const send = ( + key: string | undefined = 'same', + body: unknown = { amount: 1 } + ) => + fetch(url + '/items', { + method: 'POST', + headers: { + 'content-type': 'application/json', + ...(key === undefined ? {} : { 'x-idempotency-key': key }) + }, + body: JSON.stringify(body) + }); + return { + url, + send, + prepare, + handler, + revoke: () => { + allowed = false; + } + }; +} + +it('prepares validated input with DI on every replay, captures serialized responses, and isolates resolved scope', async () => { + const f = await fixture(); + const responses = await Promise.all( + Array.from({ length: 5 }, () => f.send()) + ); + expect(responses.map(response => response.status)).toEqual([ + 201, 201, 201, 201, 201 + ]); + expect(responses.map(response => response.headers.get('location'))).toEqual( + Array(5).fill('/items/1') + ); + expect( + await Promise.all(responses.map(response => response.json())) + ).toEqual(Array(5).fill({ tenant: 'default', amount: 1 })); + expect(f.handler).toHaveBeenCalledOnce(); + expect(f.prepare).toHaveBeenCalledTimes(5); + expect((await f.send('same', { tenant: 'other', amount: 1 })).status).toBe( + 201 + ); + expect(f.handler).toHaveBeenCalledTimes(2); + f.revoke(); + const denied = await f.send(); + expect(denied.status).toBe(403); + expect(await denied.json()).toEqual({ message: 'Denied' }); + expect(f.handler).toHaveBeenCalledTimes(2); +}); + +it('validates input and transport keys before application preparation', async () => { + const f = await fixture(); + for (const response of [ + await f.send('same', { amount: 'bad' }), + await f.send(''), + await f.send('x'.repeat(257)) + ]) + expect(response.status).toBe(400); + expect(f.prepare).not.toHaveBeenCalled(); + expect((await f.send()).status).toBe(201); +}); + +it('uses the same endpoint policy for explicit batch subrequests', async () => { + const f = await fixture(); + await f.send(); + const response = await fetch(f.url + '/__batch', { + method: 'POST', + headers: { 'content-type': 'application/json' }, + body: JSON.stringify({ + requests: [ + { + id: '1', + method: 'POST', + url: '/items', + headers: { + 'content-type': 'application/json', + 'x-idempotency-key': 'same' + }, + body: JSON.stringify({ amount: 1 }) + } + ] + }) + }); + expect(response.status).toBe(200); + const result = await response.json(); + expect(JSON.stringify(result)).toContain('default'); + expect(f.prepare).toHaveBeenCalledTimes(2); + expect(f.handler).toHaveBeenCalledOnce(); +}); + +it('enforces capacity and retains non-replayable responses', async () => { + const full = await fixture({ maxEntries: 1 }); + expect((await full.send('one')).status).toBe(201); + expect((await full.send('two')).status).toBe(503); + expect(full.handler).toHaveBeenCalledOnce(); + const small = await fixture({ maxResponseBytes: 1 }); + expect((await small.send()).status).toBe(201); + expect((await small.send()).status).toBe(409); + expect(small.handler).toHaveBeenCalledOnce(); +}); + +it('requires explicit policies at startup and keeps contract opt-in immutable', async () => { + const original = endpoint.post('/items').body(Body); + const keyed = original + .idempotent() + .summary('Create') + .clearsCacheTag('items') + .returns(Body); + expect(original.introspect().idempotent).toBe(false); + expect(keyed.introspect().idempotent).toBe(true); + await expect( + createServer() + .handle(keyed, ({ body }) => body) + .listen(0) + ).rejects.toThrow('explicit scope'); + expect(() => endpoint.get('/items').idempotent()).toThrow('mutation'); +}); diff --git a/libs/server/src/EndpointOptions.ts b/libs/server/src/EndpointOptions.ts new file mode 100644 index 00000000..e6b7ca74 --- /dev/null +++ b/libs/server/src/EndpointOptions.ts @@ -0,0 +1,41 @@ +import type { ActionContext, Handler, ResponsesOf } from './Endpoint.js'; +import type { ErrorMap, ErrorResponsesOf } from './ErrorMap.js'; +import type { Middleware } from './types.js'; + +/** Identity components are serialized without ambiguous string concatenation. */ +export type IdempotencyScope = string | readonly (string | number)[]; + +/** Per-endpoint, process-local replay limits. */ +export interface EndpointIdempotencyLimits { + ttl?: number; + maxEntries?: number; + maxResponseBytes?: number; +} + +/** Server-only callbacks, inferred from the endpoint's validated request and DI. */ +export type EndpointOptions = { + middlewares?: Middleware[]; + /** Runs on every request, including replays; return the request used by scope and handler. */ + prepare?: ( + ...args: Parameters> + ) => ActionContext | Promise>; + idempotency?: EndpointIdempotencyLimits & { + scope: ( + ...args: Parameters> + ) => IdempotencyScope | Promise; + }; + errors?: keyof ResponsesOf extends never + ? never + : ErrorMap>; +}; + +/** @internal Erased callback types retained through registration/composition. */ +export interface RuntimeEndpointOptions { + /** @internal Handler invocation already owns its error policy. */ + handlerErrorsMapped?: boolean; + prepare?: (...args: any[]) => any; + idempotency?: EndpointIdempotencyLimits & { + scope: (...args: any[]) => IdempotencyScope | Promise; + }; + errors?: ErrorMap; +} diff --git a/libs/server/src/ErrorMap.ts b/libs/server/src/ErrorMap.ts index 655ea8ee..31a647e0 100644 --- a/libs/server/src/ErrorMap.ts +++ b/libs/server/src/ErrorMap.ts @@ -150,12 +150,7 @@ export function withErrors< : ErrorMap>>, handler: Handler> ): Handler { - const responses = endpoint.introspect().responsesSchemas; - if (!responses || Object.keys(responses).length === 0) { - throw new TypeError( - 'Error policies require explicit endpoint responses' - ); - } + validateErrorMap(endpoint); return (async (...args: unknown[]) => { try { return await (handler as (...args: unknown[]) => unknown)(...args); @@ -164,3 +159,15 @@ export function withErrors< } }) as Handler; } + +/** @internal Shared registration validation. */ +export function validateErrorMap(endpoint: { + introspect(): { responsesSchemas: unknown }; +}): void { + const responses = endpoint.introspect().responsesSchemas; + if (!responses || Object.keys(responses).length === 0) { + throw new TypeError( + 'Error policies require explicit endpoint responses' + ); + } +} diff --git a/libs/server/src/Implementation.ts b/libs/server/src/Implementation.ts index b2bd9ad1..fc83d1e7 100644 --- a/libs/server/src/Implementation.ts +++ b/libs/server/src/Implementation.ts @@ -4,14 +4,13 @@ import { type EndpointBuilder, type Handler, type HandlerMapping, - mapHandlers, - type ResponsesOf + mapHandlers } from './Endpoint.js'; -import { - type ErrorMap, - type ErrorResponsesOf, - withErrors -} from './ErrorMap.js'; +import type { + EndpointOptions, + RuntimeEndpointOptions +} from './EndpointOptions.js'; +import { type ErrorMap, withErrors } from './ErrorMap.js'; import { isSubscriptionBuilder, type SubscriptionBuilder, @@ -151,15 +150,7 @@ type ConfiguredGroup = { export type ImplementationHandlerEntry = E extends AnySubscription ? SubscriptionHandlerEntry & { errors?: never } : E extends AnyEndpoint - ? - | Handler - | { - handler: Handler; - middlewares?: Middleware[]; - errors?: keyof ResponsesOf extends never - ? never - : ErrorMap>; - } + ? Handler | ({ handler: Handler } & EndpointOptions) : never; /** Complete, endpoint-specific handler bindings for one configured scope. */ @@ -175,7 +166,7 @@ type ExactOperations = O extends { readonly operations: infer Ops } ? { readonly operations: ExactKeys } : unknown; -type Entry = { +type Entry = RuntimeEndpointOptions & { readonly group: string; readonly name: string; readonly source: Definition; @@ -185,11 +176,11 @@ type Entry = { }; type RuntimeBinding = | Entry['handler'] - | { + | (RuntimeEndpointOptions & { handler: Entry['handler']; middlewares?: Middleware[]; errors?: ErrorMap; - }; + }); const moduleContract: unique symbol = Symbol('implementationContract'); const moduleEntries: unique symbol = Symbol('implementationEntries'); @@ -325,6 +316,16 @@ export class ImplementationScope< handler: errors ? withErrors(endpoint as AnyEndpoint, errors, handler) : handler, + handlerErrorsMapped: !!errors, + errors, + prepare: + typeof binding === 'function' + ? undefined + : binding?.prepare, + idempotency: + typeof binding === 'function' + ? undefined + : binding?.idempotency, middlewares: typeof binding === 'function' ? undefined @@ -519,6 +520,10 @@ export class ApiImplementation< endpoints[entry.group][entry.name] = entry.endpoint; handlers[entry.group][entry.name] = { handler: entry.handler, + handlerErrorsMapped: entry.handlerErrorsMapped, + prepare: entry.prepare, + idempotency: entry.idempotency, + errors: entry.errors, middlewares: entry.middlewares ? [...entry.middlewares] : undefined diff --git a/libs/server/src/Server.ts b/libs/server/src/Server.ts index 0eab94ef..abd5c4cb 100644 --- a/libs/server/src/Server.ts +++ b/libs/server/src/Server.ts @@ -19,8 +19,11 @@ import { ActionResult, JsonResult } from './ActionResult.js'; import { ContentNegotiator } from './ContentNegotiator.js'; import { CorsPolicy, type ServerCorsOptions } from './Cors.js'; import type { EndpointBuilder, Handler, HandlerMapping } from './Endpoint.js'; +import type { EndpointOptions } from './EndpointOptions.js'; +import { validateErrorMap } from './ErrorMap.js'; import { HttpError } from './HttpError.js'; import { MiddlewarePipeline } from './MiddlewarePipeline.js'; +import { idempotency } from './middlewares/Idempotency.js'; import { parseMultipart } from './multipart.js'; import { needsBody, resolveArgs } from './ParameterResolver.js'; import { @@ -237,15 +240,14 @@ export class ServerBuilder { */ handle< E extends EndpointBuilder - >( - endpointDef: E, - handler: Handler, - options?: { middlewares?: Middleware[] } - ): this { + >(endpointDef: E, handler: Handler, options?: EndpointOptions): this { this.#registrations.push({ endpoint: endpointDef.introspect(), handler, - middlewares: options?.middlewares + middlewares: options?.middlewares, + prepare: options?.prepare, + idempotency: options?.idempotency, + errors: options?.errors }); return this; } @@ -262,7 +264,11 @@ export class ServerBuilder { this.#registrations.push({ endpoint: entry.endpoint.introspect(), handler: entry.handler, - middlewares: entry.middlewares + middlewares: entry.middlewares, + handlerErrorsMapped: entry.handlerErrorsMapped, + prepare: entry.prepare, + idempotency: entry.idempotency, + errors: entry.errors }); } for (const entry of mapping._subscriptions) { @@ -331,6 +337,23 @@ export class ServerBuilder { const router = new Router(); for (const reg of this.#registrations) { + if ( + reg.endpoint.idempotent && + typeof reg.idempotency?.scope !== 'function' + ) + throw new TypeError( + 'Idempotent endpoints require an explicit scope policy' + ); + if (reg.idempotency && !reg.endpoint.idempotent) + throw new TypeError( + 'An idempotency policy requires an idempotent contract' + ); + if (reg.errors) + validateErrorMap({ introspect: () => reg.endpoint }); + if (reg.idempotency) { + // Validate limits before opening the listening socket. + idempotency({ ...reg.idempotency, scope: () => undefined }); + } router.addRoute(reg); } for (const reg of this.#subscriptionRegistrations) { @@ -401,6 +424,8 @@ const MAX_WS_QUEUE_SIZE = 1024; export class Server { readonly #router: Router; + readonly #replay = new Map(); + readonly #replayScopes = new WeakMap(); readonly #serviceProvider: ServiceProvider; readonly #contentNegotiator: ContentNegotiator; readonly #globalMiddlewares: Middleware[]; @@ -804,15 +829,79 @@ export class Server { } } - // Call handler - let result = registration.handler(...resolveResult.args); - if (result instanceof Promise) { - result = await result; + const args = resolveResult.args; + const send = async (result: unknown) => { + if (ctx.responded) return; + await this.#sendResult(req, res, result); + ctx.responded = true; + }; + const translate = async (error: unknown) => { + if (!registration.errors) throw error; + return registration.errors.translate(error); + }; + // Framework transport validation precedes application callbacks. + const key = ctx.headers['x-idempotency-key']; + if ( + meta.idempotent && + key !== undefined && + (key.length === 0 || key.length > 256) + ) + throw new HttpError( + 400, + 'Idempotency key must contain 1 to 256 characters' + ); + try { + if (registration.prepare) + args[0] = await registration.prepare(...args); + if (registration.idempotency && key !== undefined) { + const scope = await registration.idempotency.scope( + ...args + ); + if ( + typeof scope !== 'string' && + !( + Array.isArray(scope) && + scope.length > 0 && + scope.every( + part => + typeof part === 'string' || + (typeof part === 'number' && + Number.isFinite(part)) + ) + ) + ) + throw new TypeError( + 'Idempotency scope must be a string or identity components' + ); + this.#replayScopes.set(ctx, JSON.stringify(scope)); + } + } catch (error) { + await send(await translate(error)); + return; + } + const execute = async () => { + let result: unknown; + try { + result = await registration.handler(...args); + } catch (error) { + if (registration.handlerErrorsMapped) throw error; + result = await translate(error); + } + await send(result); + }; + if (registration.idempotency) { + let replay = this.#replay.get(registration); + if (!replay) { + replay = idempotency({ + ...registration.idempotency, + scope: request => this.#replayScopes.get(request) + }); + this.#replay.set(registration, replay); + } + await replay(ctx, execute); + } else { + await execute(); } - - if (ctx.responded) return; - await this.#sendResult(req, res, result); - ctx.responded = true; }); } catch (err) { if (res.headersSent) return; @@ -927,7 +1016,10 @@ export class Server { const virtualReq = new VirtualIncomingMessage({ method: (item.method ?? 'GET').toUpperCase(), url: item.url, - headers: item.headers ?? {}, + headers: { + host: req.headers.host ?? 'localhost', + ...item.headers + }, body: item.body }); virtualReq.__cleverbrushBatchSubrequest = { diff --git a/libs/server/src/index.ts b/libs/server/src/index.ts index a0f0dfb2..226cf89a 100644 --- a/libs/server/src/index.ts +++ b/libs/server/src/index.ts @@ -51,6 +51,11 @@ export { type ResponsesOf, type ScopedEndpointFactory } from './Endpoint.js'; +export type { + EndpointIdempotencyLimits, + EndpointOptions, + IdempotencyScope +} from './EndpointOptions.js'; export { ErrorMap, type ErrorResponse, diff --git a/libs/server/src/types.ts b/libs/server/src/types.ts index 7b8b150c..589dfe37 100644 --- a/libs/server/src/types.ts +++ b/libs/server/src/types.ts @@ -1,4 +1,5 @@ import type { EndpointMetadata } from './Endpoint.js'; +import type { RuntimeEndpointOptions } from './EndpointOptions.js'; import type { RequestContext } from './RequestContext.js'; import type { SubscriptionMetadata } from './Subscription.js'; @@ -10,7 +11,7 @@ import type { SubscriptionMetadata } from './Subscription.js'; * A registered endpoint pairing its metadata (method, path, schemas) with * the handler function and any per-endpoint middleware. */ -export interface EndpointRegistration { +export interface EndpointRegistration extends RuntimeEndpointOptions { readonly endpoint: EndpointMetadata; readonly handler: (...args: any[]) => any; readonly middlewares?: readonly Middleware[]; diff --git a/websites/docs/app/client/sections/idempotency.tsx b/websites/docs/app/client/sections/idempotency.tsx index 242995ec..88698f85 100644 --- a/websites/docs/app/client/sections/idempotency.tsx +++ b/websites/docs/app/client/sections/idempotency.tsx @@ -5,7 +5,7 @@ export default function IdempotencySection() { return ( <>
-

Idempotency Middleware

+

HTTP Idempotency

Deduplicate replays of mutating requests via idempotency keys @@ -13,23 +13,21 @@ export default function IdempotencySection() {

-

Basic Usage

+

Contract-declared retries

                     
@@ -39,22 +37,24 @@ await client.todos.create({ body: { title: 'Buy milk' } });
             

Server Integration

- The server-side idempotency() middleware reads - the header, stores the response, and replays it for - duplicate keys within an explicit scope, method and URL. - Install it after authentication and authorization. - Concurrent duplicates share one execution; this - process-local store is not durable exactly-once execution. + Endpoint preparation runs after authentication and + validation, before every replay. It receives typed request + data and injected services. Scope resolution uses the + prepared request. Concurrent duplicates share one execution + within a bounded, process-local store.

                      'public-todos', ttl: 86_400_000 })],
+                            __html: highlightTS(`server.handle(CreateTodo.authorize(UserPrincipal).inject({ db: DbToken }), createHandler, {
+    prepare: authorizeAndResolveWorkspace,
+    idempotency: {
+        scope: ({ principal, body }) => [principal.userId, body.workspaceId]
+    },
+    errors: todoErrors
 });
+
+// The same options work in implement(api).group(...).withHandlers(...).
 `)
                         }}
                     />
@@ -62,7 +62,7 @@ server.handle(CreateTodo, createHandler, {
             
-

How It Works

+

How replay works

  • On mutation: Client auto-generates a @@ -82,7 +82,7 @@ server.handle(CreateTodo, createHandler, {
-

Options (Client)

+

Low-level middleware options (Client)