InteractiveFrameworks

Streams and timeouts

Replies that arrive over time and replies that never arrive: the cats service streams search results one cat at a time, the gateway collects them, and a census that takes too long is cut off with a timeout and answered 504.

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

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 with EmptyError: 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.

  1. 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 equals breed (if given) and its age is at least minAge (if given).
  2. In the gateway, GET /cats/search?breed=&minAge= sends 'cats.search' with both values, minAge as a number. Both are optional. Answer with every cat the service streams back as one array, [] when nothing matches. Declare the route above GET /cats/:id.
  3. Never wait more than 500 ms for the census. When the service is slower, answer 504 with the message The cats service did not answer in time.

When it fails

  • /cats/search?breed=Tabby answers one cat instead of two: the route returned the Observable (you get the last reply) or used firstValueFrom() (you get the first).
  • /cats/search?breed=Sphynx is a 500 with EmptyError: 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, and toArray() wrapped it again. Return from(...) to send one reply per cat.
  • /cats/search answers 400 Validation failed (numeric string is expected): the request reached GET /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() with lastValueFrom() 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.