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
6 changes: 3 additions & 3 deletions .agents/rules/code-style.md
Original file line number Diff line number Diff line change
Expand Up @@ -84,14 +84,14 @@ const contract = defineContract({

```typescript
// Bad — using async handlers
processOrder: async ({ payload }) => {
processOrder: async ({ input: { payload } }) => {
await process(payload);
};

// Good — use the AsyncResult pattern from unthrown.
// fromPromise REQUIRES the error mapper as the second argument; chaining
// .mapErr afterwards is a type error since fromPromise has no `unknown` overload.
processOrder: ({ payload }) =>
processOrder: ({ input: { payload } }) =>
fromPromise(
process(payload),
(e) => new RetryableError("Failed", e),
Expand All @@ -103,7 +103,7 @@ processOrder: (message) => {
};

// Good — destructure payload
processOrder: ({ payload }) => {
processOrder: ({ input: { payload } }) => {
console.log(payload.orderId);
};

Expand Down
35 changes: 22 additions & 13 deletions .agents/rules/handlers.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,37 +4,46 @@ This project uses [unthrown](https://github.com/btravstack/unthrown) for explici

## Regular consumer handler

A consumer handler receives `({ payload, headers }, rawMessage)` and returns `AsyncResult<void, HandlerError>`:
A consumer handler receives one record — `{ input, context, errors, raw, retryable, nonRetryable }`, where `input` is the validated `{ payload, headers }` — and returns `AsyncResult<void, HandlerError>`. The message is repeated as a second positional parameter, oRPC's shape, for a caller who prefers it:

```typescript
import { fromPromise, OkAsync } from "unthrown";
import { RetryableError, NonRetryableError } from "@amqp-contract/worker";

// Sync OK case — lift a sync Result into an AsyncResult with .toAsync()
const handler = ({ payload }, rawMessage) => {
console.log(payload.orderId);
const handler = ({ raw, input: { payload } }) => {
console.log(payload.orderId, raw.fields.deliveryTag);
return OkAsync(undefined);
};

// Async case — fromPromise REQUIRES the qualify mapper as the second arg
const asyncHandler = ({ payload }) =>
const asyncHandler = ({ input: { payload } }) =>
fromPromise(processPayment(payload), (error) => new RetryableError("Payment failed", error)).map(
() => undefined,
);
```

### Parameters

1. **`message`** — `{ payload, headers }`
1. **`helpers`** — `{ input, context, errors, raw, retryable, nonRetryable }`
- `input`: the validated `{ payload, headers }`, the same value the second parameter carries — oRPC's shape and its word for it, so `({ errors, input }) => ...` and `({ errors }, message) => ...` are the same call, and the same name is destructured on all three transports
Comment thread
coderabbitai[bot] marked this conversation as resolved.
- `context`: seeded by `createContext`, accumulated by the middleware chain
- `errors`: typed constructors for the RPC's declared errors (empty for consumers)
- `raw`: the raw amqplib `ConsumeMessage` (e.g. `raw.fields.deliveryTag`, `raw.properties.messageId`)
- `retryable` / `nonRetryable`: the two modeled failures as factories — `ErrAsync(retryable("db down", cause))` is `new RetryableError(...)` without the import. They sit beside `errors` rather than inside it: `errors` is the contract's declared map, which is what the name means on the other two transports.
2. **`message`** (positional) — the same `{ payload, headers }` as `helpers.input`
- `payload`: validated against the message's payload schema
- `headers`: validated against the message's optional headers schema (otherwise `undefined`)
2. **`rawMessage`** — the raw amqplib `ConsumeMessage` (e.g. `msg.fields.deliveryTag`, `msg.properties.messageId`)

Helpers first is oRPC's parameter order, which this family converged on
(temporal-contract's activity leaf moved with it). A handler that needs none of
them still names the position: `({ input: { payload } }) => ...`.

## RPC handler

`defineRpc` creates a request-reply slot. RPC handlers return `AsyncResult<TResponse, HandlerError | WorkerInferRpcErrors<...>>` — the worker validates the response against the RPC's response schema and publishes it back to the caller's `replyTo` with the same `correlationId`.

All handlers (consumer and RPC) receive a third `helpers` argument — `{ context, errors }`. `context` is seeded by `createContext` and accumulated by the middleware chain (`TypedAmqpWorker.create({ createContext, middleware: composeMiddleware(...) })`); `errors` carries typed constructors for the RPC's declared errors (`ErrAsync(errors.CODE({ ... }))`), empty for consumers. Middleware `next({ payload })` substitutes the payload with re-validation before the handler. See `packages/worker/src/middleware.ts`, `packages/worker/src/worker.ts` (`runHandler`), and [docs/guide/middleware-and-interceptors.md](../../docs/guide/middleware-and-interceptors.md).
All handlers (consumer and RPC) receive the `helpers` record first — `{ context, errors, raw }`. `context` is seeded by `createContext` and accumulated by the middleware chain (`TypedAmqpWorker.create({ createContext, middleware: composeMiddleware(...) })`); `errors` carries typed constructors for the RPC's declared errors (`ErrAsync(errors.CODE({ ... }))`), empty for consumers. Middleware `next({ payload })` substitutes the payload with re-validation before the handler. See `packages/worker/src/middleware.ts`, `packages/worker/src/worker.ts` (`runHandler`), and [docs/guide/middleware-and-interceptors.md](../../docs/guide/middleware-and-interceptors.md).

When the RPC declares an `errors` map (`defineRpc(queue, { request, response, errors })`), the handler may also return `Err(rpcError(code, data))` for a declared code — the worker validates `data` against the declared schema, publishes an error reply (marked by the `RPC_ERROR_CODE_HEADER` header), and **acks the request**; typed business errors never enter the retry/DLQ pipeline. Undeclared codes or invalid error data are contract violations routed to the DLQ. See `packages/worker/src/worker.ts` (`publishRpcErrorReply`) and [docs/guide/error-model.md](../../docs/guide/error-model.md#typed-rpc-errors-rpcerror).

Expand All @@ -48,13 +57,13 @@ const result = await TypedAmqpWorker.create({
contract,
handlers: {
// Regular consumer — `payload` typed from the consumer's message schema
processOrder: ({ payload }) => OkAsync(undefined),
processOrder: ({ input: { payload } }) => OkAsync(undefined),

// RPC handler — must return the typed response payload
calculate: ({ payload }) => OkAsync({ sum: payload.a + payload.b }),
calculate: ({ input: { payload } }) => OkAsync({ sum: payload.a + payload.b }),

// RPC with async work
lookupUser: ({ payload }) =>
lookupUser: ({ input: { payload } }) =>
fromPromise(
db.users.findById(payload.userId),
(error) => new RetryableError("DB unavailable", error),
Expand Down Expand Up @@ -99,15 +108,15 @@ Use `declareHandler` (single) or `declareHandlers` (object) for full type infere
import { declareHandler, RetryableError, NonRetryableError } from "@amqp-contract/worker";
import { ErrAsync, fromPromise, OkAsync } from "unthrown";

const processOrderHandler = declareHandler(contract, "processOrder", ({ payload }) =>
const processOrderHandler = declareHandler(contract, "processOrder", ({ input: { payload } }) =>
fromPromise(
processPayment(payload.orderId),
(error) => new RetryableError("Payment service unavailable", error),
).map(() => undefined),
);

// Permanent failures use NonRetryableError → DLQ, never retried
const validateOrderHandler = declareHandler(contract, "validateOrder", ({ payload }) => {
const validateOrderHandler = declareHandler(contract, "validateOrder", ({ input: { payload } }) => {
if (payload.amount < 1) {
return ErrAsync(new NonRetryableError("Invalid amount"));
}
Expand Down Expand Up @@ -149,7 +158,7 @@ Helpers: `qualifyRetryable(message)` / `qualifyNonRetryable(message)` build `fro

```typescript
// Conditional error mapping inside fromPromise's qualify
({ payload }) =>
({ input: { payload } }) =>
fromPromise(process(payload), (error) => {
if (error instanceof ValidationError) return new NonRetryableError("Invalid data");
return new RetryableError("Temporary failure", error);
Expand Down
8 changes: 4 additions & 4 deletions .agents/rules/recipes.md
Original file line number Diff line number Diff line change
Expand Up @@ -8,7 +8,7 @@ End-to-end how-tos for the changes that come up most. Each recipe lists the exac
2. **Queue** — `defineQueue(...)` with a `deadLetter` and a `retry` mode (immediate-requeue or ttl-backoff). Quorum by default; classic only if you need priority/exclusive/auto-delete.
3. **Consumer entry** — `defineEventConsumer(eventPublisher, queue, { routingKey: ... })`. The queue↔exchange binding is auto-generated.
4. **Add to `defineContract`** under `consumers: { ... }`. Don't add the queue or binding yourself — they're auto-extracted.
5. **Handler** — implement with `declareHandler(contract, "yourConsumerName", ({ payload, headers }) => …)` returning `AsyncResult<void, HandlerError>`. See [handlers.md](./handlers.md).
5. **Handler** — implement with `declareHandler(contract, "yourConsumerName", ({ input: { payload, headers } }) => …)` returning `AsyncResult<void, HandlerError>`. See [handlers.md](./handlers.md).
6. **Tests** — integration test in `src/__tests__/<consumer>.spec.ts` using `it` from `@amqp-contract/testing/extension`. Mock the handler with `vi.fn().mockReturnValue(OkAsync(undefined))`.
7. **Changeset** — `pnpm changeset` with a minor bump. Public API surface grew.

Expand All @@ -18,7 +18,7 @@ End-to-end how-tos for the changes that come up most. Each recipe lists the exac
2. **Queue** — `defineQueue(...)` for the RPC. Quorum by default. **Configure a `deadLetter`** even though replies, not the queue, drive most failure modes: missing `replyTo` / `correlationId` and response-schema mismatches are surfaced as `NonRetryableError` and the worker `nack`s them without requeue, so without a DLX they're dropped silently.
3. **RPC entry** — `defineRpc(queue, { request, response })`. Typed business errors go in an optional `errors` map whose entries are `{ data: schema, message?: string }` (the raw Standard Schema, NOT `defineMessage`); the optional `message` is the default human message when the handler constructs the error without one.
4. **Add to `defineContract`** under `rpcs: { ... }`.
5. **Server-side handler** — define it with `declareHandler(contract, "yourRpcName", ({ payload }) => OkAsync({ /* response */ }))`, via `declareHandlers`, or inline in the `handlers` object passed to `TypedAmqpWorker.create({ handlers: { … } })`. All three are RPC-aware: `declareHandler` / `declareHandlers` are overloaded against `InferRpcNames` and validate the name against both `contract.consumers` and `contract.rpcs`. The worker validates the response against the response schema and publishes back automatically.
5. **Server-side handler** — define it with `declareHandler(contract, "yourRpcName", ({ input: { payload } }) => OkAsync({ /* response */ }))`, via `declareHandlers`, or inline in the `handlers` object passed to `TypedAmqpWorker.create({ handlers: { … } })`. All three are RPC-aware: `declareHandler` / `declareHandlers` are overloaded against `InferRpcNames` and validate the name against both `contract.consumers` and `contract.rpcs`. The worker validates the response against the response schema and publishes back automatically.
6. **Client call** — `client.call("yourRpcName", request, { timeoutMs: 5_000 })`. `timeoutMs` is required.
7. **Tests** — round-trip integration test (worker + client both wired up). For "no server" scenarios, just create the client without a worker; for "request validation fails", pass a deliberately wrong payload through `as unknown as ...`.
8. **Changeset** — minor bump.
Expand Down Expand Up @@ -51,15 +51,15 @@ If you're spinning up a new `@amqp-contract/*` package:
Old shape (now banned):

```typescript
processOrder: async ({ payload }) => {
processOrder: async ({ input: { payload } }) => {
await processPayment(payload);
};
```

New shape:

```typescript
processOrder: ({ payload }) =>
processOrder: ({ input: { payload } }) =>
fromPromise(processPayment(payload), (error) => new RetryableError("Payment failed", error)).map(
() => undefined,
);
Expand Down
2 changes: 1 addition & 1 deletion .agents/rules/testing.md
Original file line number Diff line number Diff line change
Expand Up @@ -68,10 +68,10 @@ const mockHandler = vi.fn().mockReturnValue(OkAsync(undefined));

// Assertion pattern
expect(mockHandler).toHaveBeenCalledWith(
expect.anything(), // helpers: { context, errors, raw }
expect.objectContaining({
payload: expect.objectContaining({ orderId: "123" }),
}),
expect.anything(), // rawMessage
);
```

Expand Down
62 changes: 62 additions & 0 deletions .changeset/handler-helpers-first.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,62 @@
---
"@amqp-contract/worker": major
---

Handlers take **one record** — `{ input, context, errors, raw, retryable,
nonRetryable }`, where `input` is the validated `{ payload, headers }` — with
that message repeated as a second positional parameter. It was
`({ payload, headers }, rawMessage, { context, errors })`: the raw amqplib
delivery moved onto the record as `raw`, and the message is reachable from
either place.

oRPC is the reference shape for this family, being the most widely used of the
three transports a `@btravstack/*` application composes: a developer arriving
here has more likely seen `({ errors, input })` than either of the others —
down to the word, since oRPC's `ProcedureHandlerOptions` carries `input` and
its handler still takes it positionally, which is exactly the pair of spellings
offered here. The mint and compose calls already agreed across the three; the leaf a
developer types by hand did not, and it is the one they relearn per transport.
`@temporal-contract`'s activity leaf moves with it.

It is also what makes the AMQP triage site the same SHAPE as the other two, and
that half is not cosmetic: `retryable` and `nonRetryable` ride the helpers
record beside `errors`, so a handler that wants "infrastructure comes back"
reaches for the constructor it was handed instead of importing `RetryableError`
and constructing it by hand.

```ts
processOrder: ({ retryable, input: { payload } }) =>
fromPromise(save(payload), (cause) => retryable("database unavailable", cause)),
```

They sit BESIDE `errors` rather than inside it: `errors` is the
contract-declared error map — `errors.ORDER_NOT_FOUND({ orderId })` — which is
what it means on the other two transports, and folding the framework's own two
into that namespace would both break the mirror and collide with a declared code
called `retryable`.

```diff
- processOrder: ({ payload }) => save(payload),
+ processOrder: ({ input: { payload } }) => save(payload),
- handleFailed: ({ payload }, rawMessage) => log(rawMessage.properties.headers),
+ handleFailed: ({ raw, input: { payload } }) => log(raw.properties.headers),
- getOrder: ({ payload }, _raw, { errors }) => lookup(payload, errors),
+ getOrder: ({ errors, input: { payload } }) => lookup(payload, errors),
```

The message is on the helpers record as well as in the second parameter, which
is oRPC's own shape — `ProcedureHandlerOptions` carries `input` and the handler
still takes it positionally — so both spellings are the same call:

```ts
getOrder: ({ errors, input }) => lookup(input.payload, errors),
getOrder: ({ errors, input: { payload } }) => lookup(payload, errors),
```

A handler that reads its payload fails to compile until it is swapped, since
the first parameter is the helpers record now; one that ignores its message
keeps compiling with a parameter whose name lies — grep the handlers object for
a leaf whose first parameter is neither `_` nor a helpers destructuring. A
handler that wants only its message is `({ input: { payload } }) => ...`, with no placeholder to spell.

Closes #670.
2 changes: 1 addition & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -94,7 +94,7 @@ import { contract } from "./contract.js";
const worker = await TypedAmqpWorker.create({
contract,
handlers: {
processOrder: ({ payload }) => {
processOrder: ({ input: { payload } }) => {
console.log(payload.orderId); // ✅ TypeScript knows!
return OkAsync();
},
Expand Down
12 changes: 6 additions & 6 deletions docs/examples/basic-order-processing.md
Original file line number Diff line number Diff line change
Expand Up @@ -447,7 +447,7 @@ const worker = await TypedAmqpWorker.create({
contract: orderContract,
handlers: declareHandlers(orderContract, {
// Event handler for NEW orders (order.created) — headers are typed
processOrder: ({ payload, headers }) => {
processOrder: ({ input: { payload, headers } }) => {
console.log(`[PROCESSING] Order ${payload.orderId}`, {
customer: payload.customerId,
total: payload.totalAmount,
Expand All @@ -460,7 +460,7 @@ const worker = await TypedAmqpWorker.create({
},

// Event handler for ALL order events (order.#) — payload is the union type
notifyOrder: ({ payload }) => {
notifyOrder: ({ input: { payload } }) => {
if ("items" in payload) {
console.log(`[NOTIFICATIONS] New order ${payload.orderId}`);
} else {
Expand All @@ -472,31 +472,31 @@ const worker = await TypedAmqpWorker.create({
},

// Event handler for SHIPPED orders (order.shipped)
shipOrder: ({ payload }) => {
shipOrder: ({ input: { payload } }) => {
console.log(`[SHIPPING] Order ${payload.orderId} - ${payload.status}`);
return fromPromise(prepareShipping(payload), qualifyRetryable("Shipping failed")).map(
() => undefined,
);
},

// Event handler for URGENT orders (order.*.urgent)
handleUrgentOrder: ({ payload }) => {
handleUrgentOrder: ({ input: { payload } }) => {
console.warn(`[URGENT] Order ${payload.orderId} - ${payload.status}`);
return fromPromise(escalate(payload), qualifyRetryable("Urgent handling failed")).map(
() => undefined,
);
},

// Command handler (task queue): reaches exactly one worker
fulfillOrder: ({ payload }) => {
fulfillOrder: ({ input: { payload } }) => {
console.log(`[FULFILLMENT] Order ${payload.orderId} → ${payload.warehouseId}`);
return fromPromise(fulfill(payload), qualifyRetryable("Fulfillment failed")).map(
() => undefined,
);
},

// Dead-letter handler: messages that failed in order-processing
handleFailedOrders: ({ payload }) => {
handleFailedOrders: ({ input: { payload } }) => {
console.error(`[DLX] Failed order ${payload.orderId}`);
return fromPromise(recordFailure(payload), qualifyRetryable("DLX handling failed")).map(
() => undefined,
Expand Down
2 changes: 1 addition & 1 deletion docs/examples/command-pattern.md
Original file line number Diff line number Diff line change
Expand Up @@ -118,7 +118,7 @@ import {
import { fromPromise } from "unthrown";
import { contract } from "@org/payment-contract";

const chargeHandler = declareHandler(contract, "chargeCustomer", ({ payload }) =>
const chargeHandler = declareHandler(contract, "chargeCustomer", ({ input: { payload } }) =>
fromPromise(
chargeProvider({
customerId: payload.customerId,
Expand Down
2 changes: 1 addition & 1 deletion docs/explanation/core-concepts.md
Original file line number Diff line number Diff line change
Expand Up @@ -113,7 +113,7 @@ const orderMessage = defineMessage(z.object({ orderId: z.string() }));
When the worker asks for `handlers.processOrder`, it looks up that key and knows the payload is `{ orderId: string }`. Which is why this fails to compile:

```typescript
processOrder: ({ payload }) => {
processOrder: ({ input: { payload } }) => {
console.log(payload.orderNumber); // Property 'orderNumber' does not exist
return OkAsync(undefined);
};
Expand Down
4 changes: 2 additions & 2 deletions docs/explanation/delivery-guarantees.md
Original file line number Diff line number Diff line change
Expand Up @@ -58,8 +58,8 @@ await client
A handler reads it from the raw message:

```typescript
processOrder: ({ payload }, rawMessage) => {
const { messageId } = rawMessage.properties;
processOrder: ({ raw, input: { payload } }) => {
const { messageId } = raw.properties;
const id = typeof messageId === "string" ? messageId : undefined;
return upsertOrder(payload, id).map(() => undefined);
},
Expand Down
2 changes: 1 addition & 1 deletion docs/explanation/the-retry-model.md
Original file line number Diff line number Diff line change
Expand Up @@ -22,7 +22,7 @@ Conflating them is the usual design, and it goes wrong in a familiar way: retry
So they are separated. The handler classifies:

```typescript
processOrder: ({ payload }) =>
processOrder: ({ input: { payload } }) =>
fromPromise(chargeCard(payload), (cause) =>
cause instanceof CardDeclined
? new NonRetryableError("card declined", cause)
Expand Down
2 changes: 1 addition & 1 deletion docs/explanation/why-amqp-contract.md
Original file line number Diff line number Diff line change
Expand Up @@ -101,7 +101,7 @@ An exception is a poor way to carry that decision. It has type `unknown`, it can
So handlers return their outcome instead:

```typescript
processOrder: ({ payload }) =>
processOrder: ({ input: { payload } }) =>
fromPromise(chargeCard(payload), (cause) =>
cause instanceof CardDeclined
? new NonRetryableError("declined", cause) // → dead letter
Expand Down
Loading
Loading