InteractiveFrameworks

A transport of your own

Carry the cats service's messages over something Nest has no transporter for: the shelter's message bus. Write the server strategy that finds and runs the handlers, and the client proxy that publishes requests, matches their replies and dispatches events, then plug both in without touching a handler.

What you'll learn

  • Write a CustomTransportStrategy: a Server subclass whose listen() receives messages and hands them to the handlers Nest registered, found with getHandlerByPattern()
  • Reply through Server.send() and transformToObservable(), so values, promises, streams and errors all become WritePackets
  • Write a ClientProxy subclass: publish() with a callback for each reply and a teardown, dispatchEvent() for events, connect() and close()
  • Plug both in with createMicroservice's strategy and ClientsModule's customClass, and know that handlers and callers do not change

The shelter's partners already run a message bus, a hosted service in the style of Google Pub/Sub: named topics, and every message published to a topic delivered to whoever subscribes to it. They want the cats service on it, not on a TCP port they would have to open. Nest ships transporters for TCP, Redis, NATS, MQTT, RabbitMQ, Kafka and gRPC, and none for this bus. But a transporter is only two classes: a server strategy that receives messages and hands them to your handlers, and a client proxy that sends them. Write those two, and every @MessagePattern, every send() and every emit() keeps working unchanged.

The server strategy

A custom strategy extends Server and implements CustomTransportStrategy. Nest calls listen() when the app starts and close() when it stops:

export class QueueServer extends Server implements CustomTransportStrategy {
  listen(callback: () => void) {
    this.queue.onMessage((msg) => this.handle(msg));
    callback();
  }
  close() {
    this.queue.stop();
  }
  on() { throw new Error('Not supported'); }
  unwrap<T>(): T { return this.queue as T; }
}

By the time listen() runs, Nest has registered every handler it found. this.messageHandlers is a Map from pattern to handler. Object patterns are keyed by their JSON, {"cmd":"find-all"}, and getHandlerByPattern(pattern) finds one. A handler is Nest's wrapper around your method, with the pipes, guards, interceptors and filters already applied. Call it with the payload and a context (new BaseRpcContext([]) when the transport has nothing to add). It resolves to a value or to an Observable, depending on what your method and the interceptors returned, and a thrown exception comes back as an erroring Observable.

The base class turns all of that into replies. transformToObservable() makes whatever the handler resolved to into an Observable. send(stream$, respond) subscribes to it and calls respond with Nest's WritePackets:

  • { response } for each value, the last one marked isDisposed: true;
  • { err, isDisposed: true } when it errors;
  • { isDisposed: true } alone when the stream completes empty.

Your transport's only job is to get each packet back to the caller that asked. For events there is handleEvent(pattern, packet, context), which runs every event handler for the pattern and logs a missing one.

The client proxy

A client extends ClientProxy. send() and emit() are already written, in terms of four methods you provide:

  • connect(), awaited before every message, so it must be cheap after the first time;
  • publish(packet, callback) for a request: send it, and call callback with each reply's packet. It returns a teardown that Nest calls when the caller stops listening, after the last reply, on an error, or on a timeout();
  • dispatchEvent(packet) for an event: send it, expect nothing;
  • close(), when the app shuts down.

A bus has no idea of replies, so the client makes them: it gives every request an id and a topic to reply on, and keeps the callbacks by id until their teardown. this.normalizePattern(pattern) gives the string form Nest keys its handlers by. Send that, not the object.

Plugging them in

The server replaces the transport options in createMicroservice(): { strategy: new QueueServer(queue) }. The client is registered by class: ClientsModule.register([{ name: 'ORDERS', customClass: QueueClient, options: {...} }]). Nest creates it with new QueueClient(options), so anything it needs, the bus included, goes in options.

Your task

bus/shelter-bus.ts is the bus: subscribe(topic, listener) returns an unsubscribe function, publish(topic, message) delivers to that topic's subscribers a moment later. It also counts what each topic carried, which the gateway serves at GET /bus. The service and the gateway still talk over TCP. Their handlers and routes do not change in this lesson.

  1. Write ShelterBusServer. Subscribe to requests and to events. For a request ({ id, pattern, data, replyTo }), run its handler and publish every reply packet to replyTo, each carrying the request's id. A pattern nobody handles gets one reply whose err is There is no matching message handler defined in the remote service.. Events go to Nest's own event handling. close() unsubscribes.
  2. Write ShelterBusClient. connect() subscribes once to the replyTo topic from its options. publish() gives the request an id, publishes it to requests with the normalized pattern, and hands each reply with that id to the callback. dispatchEvent() publishes to events. close() unsubscribes.
  3. Switch both sides over: the service listens with strategy: new ShelterBusServer(shelterBus), and CATS_SERVICE becomes a ShelterBusClient with the options { bus: shelterBus, replyTo: 'replies.gateway' }.

When it fails

  • A request never answers and the gateway waits until the request times out: the reply was published without the request's id, to another topic, or not at all. The no-handler case needs a reply too.
  • A search answers [{"source":{...}}] and a missing cat answers 200 {}: the server published the handler's result as it was, an Observable. Reply through transformToObservable() and send().
  • GET /bus still shows nothing: the app still uses TCP on one side. Both the strategy and the custom class must change.
  • Property 'dispatchEvent' in type 'ShelterBusClient' is not assignable to the same property in base type: it must return Promise<any>, not Promise<void>.

Remember

  • A transporter is a Server strategy and a ClientProxy; handlers and callers never change.
  • getHandlerByPattern() finds the handler, transformToObservable() and send() turn anything it gives back into packets.
  • publish() delivers every reply to its callback and returns a teardown; dispatchEvent() is fire and forget.
  • Plug in with strategy on the server and customClass on the client.
Stuck? Show a hint

Server: in listen(), bus.subscribe('requests', ...) and bus.subscribe('events', ...); for a request, this.getHandlerByPattern(request.pattern), then this.send(this.transformToObservable(await handler(request.data, new BaseRpcContext([]))), packet => bus.publish(request.replyTo, { ...packet, id: request.id })); for an event, this.handleEvent(event.pattern, event, new BaseRpcContext([])). Client: subscribe to replyTo in connect(), keep a Map from id to callback, publish { id, pattern: this.normalizePattern(packet.pattern), data, replyTo }, and return a function that deletes the id. main.ts: { strategy: new ShelterBusServer(shelterBus) }; the gateway: { name, customClass: ShelterBusClient, options: { bus: shelterBus, replyTo: 'replies.gateway' } }.