Files
What you'll learn
- Return an Observable from a message handler to send several replies, and know that each one crosses the wire as it is emitted
- Read parts of a payload with @Payload('key') and the transport's context with @Ctx(), and know which arguments an undecorated handler receives
- Collect a stream on the client with toArray() and lastValueFrom(), and know why firstValueFrom() and returning the Observable from a route both answer with a single value
- Bound a call with the timeout() operator and turn its TimeoutError into a 504 at the gateway
Streams and timeouts
A search over thousands of cats should not wait until the last row is read before the first one goes back. And a gateway should never wait forever for a service that is overloaded, restarting or gone: a visitor's request that hangs for a minute is worse than a clear error after half a second. This lesson covers both sides of time in a reply: replies that arrive one by one, and replies that do not arrive at all.
Many replies to one message
A message handler may return an Observable. Nest subscribes to it and sends every value it emits as a separate reply, then a final packet saying the stream is done:
@MessagePattern('orders.lines')
lines(@Payload('orderId') orderId: number): Observable<OrderLine> {
return from(this.orders.linesOf(orderId));
}
On the calling side send() returns an Observable that emits each reply as it arrives and completes after the last. What you do with it decides what you get:
lastValueFrom(send(...).pipe(toArray()))collects every reply into one array, and gives[]when the stream was empty.firstValueFrom(send(...))takes the first reply and ignores the rest. On an empty stream it fails withEmptyError: no elements in sequence.- Returning the Observable from a route handler makes Nest answer with the last value, and it fails the same way when there is none.
Returning an array from the handler is not a stream: it is one reply that happens to be an array. A client that collects with toArray() then gets an array inside an array.
Parts of a payload
@Payload('orderId') hands the handler one property of the payload instead of the whole object, like @Body('name') over HTTP. @Ctx() gives the transport's context, a TcpContext here, whose getPattern() names the pattern that was matched (an object pattern as its JSON, {"cmd":"find-one"}).
Without any decorator a handler receives the payload as its first argument, and nothing else: a second parameter is undefined. Once one parameter is decorated, only decorated parameters receive anything, so a lone @Ctx() leaves the payload undefined. That is why the docs pair them: (@Payload() data, @Ctx() context).
Only JSON crosses the wire, so an optional value that is undefined does not travel at all. The key is simply missing from the payload, and @Payload('breed') reads undefined on the other side.
Waiting, but not forever
The client has no timeout of its own: a reply that never comes is waited for as long as the connection lives. RxJS's timeout() operator bounds it. If no value arrives within the time, the Observable errors with a TimeoutError and the subscription is dropped:
return this.inventory.send('stock.count', { sku }).pipe(
timeout(2000),
catchError((err) =>
throwError(() => (err instanceof TimeoutError ? new ServiceUnavailableException('Stock is not answering') : err)),
),
);
Uncaught, the TimeoutError reaches Nest's HTTP exception layer as an unknown error, a 500. Turning it into an HTTP exception tells the caller what actually happened. For a gateway that is 504 Gateway Timeout, the status for "the server behind me did not answer in time". Other errors pass through unchanged.
Dropping the subscription does not stop the service: it keeps working and sends its reply, which the client now ignores. A timeout protects the caller, not the callee.
Your task
The cats service knows five cats now. It answers 'cats.census' with a count across partner shelters, and each shelter takes 100 ms, so the census gets slower with every shelter.
- In the service's
CatsController, answer'cats.search'. Its payload is{ breed?, minAge? }. Reply with a stream: one reply per matching cat, in the service's order. A cat matches when its breed equalsbreed(if given) and its age is at leastminAge(if given). - In the gateway,
GET /cats/search?breed=&minAge=sends'cats.search'with both values,minAgeas a number. Both are optional. Answer with every cat the service streams back as one array,[]when nothing matches. Declare the route aboveGET /cats/:id. - Never wait more than 500 ms for the census. When the service is slower, answer
504with the messageThe cats service did not answer in time.
When it fails
/cats/search?breed=Tabbyanswers one cat instead of two: the route returned the Observable (you get the last reply) or usedfirstValueFrom()(you get the first)./cats/search?breed=Sphynxis a 500 withEmptyError: no elements in sequence: the stream had no replies, and something expected at least one.toArray()turns nothing into[].- The search answers
[[...]]: the handler returned an array, one reply, andtoArray()wrapped it again. Returnfrom(...)to send one reply per cat. /cats/searchanswers 400Validation failed (numeric string is expected): the request reachedGET /cats/:id, declared above the search route.- The slow census is a 500 and the console says
TimeoutError: Timeout has occurred: the timeout works, but nothing turned its error into a 504.
Remember
- A handler that returns an Observable replies once per value;
send()emits them as they arrive. toArray()withlastValueFrom()collects a stream;firstValueFrom()and returned Observables keep one value.@Payload('key')reads part of a payload; once one parameter is decorated, decorate them all.- Bound every call with
timeout(), and answer 504 when the service behind the gateway is too slow.
Stuck? Show a hint
Service: @MessagePattern('cats.search') with @Payload('breed') and @Payload('minAge') parameters, returning from(cats).pipe(filter(...)). Gateway: @Get('search') above @Get(':id'), @Query('minAge', new ParseIntPipe({ optional: true })), lastValueFrom(send(...).pipe(toArray())); the census pipes timeout(500) and catchError, rethrowing a GatewayTimeoutException for a TimeoutError.
Press Run tests to start the app. Its log appears here.Graded endpoints
The service emits one reply per match and the gateway collects them into one array, in the order they were sent
minAge travels as a number and the service keeps cats at least that old
Both keys of the payload are read: Siamese cats aged 3 or more
The stream completes without a single reply; toArray() makes that an empty array rather than an error
Without a breed or an age every cat matches, in the service's order
Two partner shelters answer in 200 ms, inside the gateway's 500 ms
Nine shelters take 900 ms; the gateway stops waiting at 500 ms and answers 504 instead of hanging
GET /cats/:id from the first lesson is unchanged, and /cats/search did not take its place
{ cmd: 'find-all' } still answers with a single reply that is an array, unlike the search's stream
The adoption event from the previous lesson is published as before
Both event handlers ran: Luna's adopter is recorded, and the stream carries the current records
NoticesController announced the adoption