API reference / @node-ts/bus-core / Receiver
Interface: Receiver<TReceivedMessage, TTransportMessage, TReceiveResult>
Defined in: packages/bus-core/src/receiver/receiver.ts:10
When used, the app will be responsible for receiving its own messages rather than subscribing directly to the transport. This can be useful for local testing or in serverless environments, where the cloud will manage receiving a message from the transport and passing it directly to the serverless container.
Type Parameters
| Type Parameter | Default type |
|---|---|
TReceivedMessage | unknown |
TTransportMessage extends TransportMessage<unknown> | TransportMessage<unknown> |
TReceiveResult | unknown |
Methods
receive()
receive(receivedMessage, messageSerializer): Promise<TTransportMessage | TTransportMessage[]>;Defined in: packages/bus-core/src/receiver/receiver.ts:23
Invoked when a message is received by the application and needs to be converted into a transport message so that it can be passed to the dispatcher and send to handlers.
Parameters
| Parameter | Type | Description |
|---|---|---|
receivedMessage | TReceivedMessage | The message received by the app |
messageSerializer | MessageSerializer | The configured serializer, which can be used to deserialize the incoming message |
Returns
Promise<TTransportMessage | TTransportMessage[]>
toReceiveResult()?
optional toReceiveResult(failures): TReceiveResult | Promise<TReceiveResult>;Defined in: packages/bus-core/src/receiver/receiver.ts:69
Optional. When implemented, each received message is handled independently: a failing message no longer makes bus.receive() reject straight away. Instead, once every message has been handled, this is called with the ones that failed and its return value is returned by bus.receive(). Throw from it to fail the whole batch. Messages without a handler are discarded and are not reported as failures. Messages that a handler returned with bus.returnMessage() are reported as failures with a ReceivedMessageReturnedToQueue error.
When not implemented, bus.receive() rejects as soon as any message fails.
Parameters
| Parameter | Type | Description |
|---|---|---|
failures | ReceivedMessageFailure<TTransportMessage>[] | The messages that failed to be handled, and why. Empty when all succeeded. |
Returns
TReceiveResult | Promise<TReceiveResult>
The value that bus.receive() resolves with, e.g. a partial batch response for the host.
Example
// Without this hook, one bad message in a batch of 10 makes bus.receive() reject, and the host retries all 10,
// re-running the 9 handlers that already succeeded. With it, the receiver tells the host which ones to retry.
// Here, a Lambda with `ReportBatchItemFailures` enabled retries only the listed SQS records and deletes the rest.
class SqsBatchReceiver implements Receiver<SQSEvent, TransportMessage<SQSRecord>, SQSBatchResponse> {
async receive(event: SQSEvent, serializer: MessageSerializer) {
return event.Records.map(record => ({
id: record.messageId,
domainMessage: serializer.deserialize(record.body),
attributes: { attributes: {}, stickyAttributes: {} },
raw: record
}))
}
toReceiveResult(failures: ReceivedMessageFailure<TransportMessage<SQSRecord>>[]): SQSBatchResponse {
// Every failure is visible in one place, e.g. to count retries separately from real errors
failures.forEach(({ message, error }) =>
metrics.increment(error instanceof ReceivedMessageReturnedToQueue ? 'returned' : 'failed', {
messageName: message.domainMessage.$name
})
)
return {
batchItemFailures: failures.map(({ message }) => ({ itemIdentifier: message.raw.messageId }))
}
}
}
const bus = Bus.configure().withReceiver(new SqsBatchReceiver()).build()
export const handler: SQSHandler = event => bus.receive<SQSBatchResponse>(event)