Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 11 additions & 0 deletions .changeset/typed-idempotent-operations.md
Original file line number Diff line number Diff line change
@@ -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.
21 changes: 21 additions & 0 deletions libs/client/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
81 changes: 62 additions & 19 deletions libs/client/src/client.ts
Original file line number Diff line number Diff line change
Expand Up @@ -196,6 +196,22 @@ export function createClient<T extends ApiContract>(
...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}`;
Expand Down Expand Up @@ -255,26 +271,14 @@ export function createClient<T extends ApiContract>(
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<any> {
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<string, unknown> = {};
if (args?.retry !== undefined) perCallOptions.retry = args.retry;
Expand Down Expand Up @@ -330,13 +334,37 @@ export function createClient<T extends ApiContract>(
headers: args?.headers ?? ({} as Record<string, string>),
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<any> {
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);
Expand Down Expand Up @@ -392,7 +420,12 @@ export function createClient<T extends ApiContract>(
}

// Streaming fetch — yields newline-delimited chunks (e.g. NDJSON).
async function* streamLines(ep: any, args: any): AsyncIterable<string> {
async function* streamLines(
ep: any,
args: any,
groupName?: string,
endpointName?: string
): AsyncIterable<string> {
const {
url,
method,
Expand All @@ -407,6 +440,8 @@ export function createClient<T extends ApiContract>(
signal: args?.signal
};

attachRequestOptions(ep, args, init, groupName, endpointName);

// -- beforeRequest hooks --
await runBeforeRequest(hooks, url, init);

Expand Down Expand Up @@ -489,7 +524,8 @@ export function createClient<T extends ApiContract>(
// 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<Blob> => {
const {
url,
Expand All @@ -502,6 +538,13 @@ export function createClient<T extends ApiContract>(
headers: reqHeaders,
body
};
attachRequestOptions(
ep,
args,
init,
groupName,
endpointName
);
await runBeforeRequest(hooks, url, init);
const response = await composedFetch(url, init);
if (!response.ok) {
Expand Down
119 changes: 119 additions & 0 deletions libs/client/src/contractIdempotency.test.ts
Original file line number Diff line number Diff line change
@@ -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');
});
12 changes: 12 additions & 0 deletions libs/client/src/middleware.ts
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,8 @@ export function getPerCallOptions<T>(
* 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"`. */
Expand Down Expand Up @@ -185,3 +187,13 @@ export interface EndpointMeta {
/** Request headers from the call, e.g. `{ 'x-request-id': 'abc' }`. */
headers: Readonly<Record<string, string>>;
}

/** @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')
);
}
8 changes: 6 additions & 2 deletions libs/client/src/middlewares/batching.ts
Original file line number Diff line number Diff line change
Expand Up @@ -29,7 +29,11 @@
* @module
*/

import type { FetchLike, Middleware } from '../middleware.js';
import {
type FetchLike,
isIdempotentRequest,
type Middleware
} from '../middleware.js';

// ---------------------------------------------------------------------------
// Types
Expand Down Expand Up @@ -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);
}

Expand Down
8 changes: 6 additions & 2 deletions libs/client/src/middlewares/retry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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);
}

Expand Down
Loading
Loading