Files
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
A transport of your own
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 markedisDisposed: 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 callcallbackwith 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 atimeout();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.
- Write
ShelterBusServer. Subscribe torequestsand toevents. For a request ({ id, pattern, data, replyTo }), run its handler and publish every reply packet toreplyTo, each carrying the request'sid. A pattern nobody handles gets one reply whoseerrisThere is no matching message handler defined in the remote service.. Events go to Nest's own event handling.close()unsubscribes. - Write
ShelterBusClient.connect()subscribes once to thereplyTotopic from its options.publish()gives the request an id, publishes it torequestswith the normalized pattern, and hands each reply with that id to the callback.dispatchEvent()publishes toevents.close()unsubscribes. - Switch both sides over: the service listens with
strategy: new ShelterBusServer(shelterBus), andCATS_SERVICEbecomes aShelterBusClientwith 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 throughtransformToObservable()andsend(). GET /busstill 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 returnPromise<any>, notPromise<void>.
Remember
- A transporter is a
Serverstrategy and aClientProxy; handlers and callers never change. getHandlerByPattern()finds the handler,transformToObservable()andsend()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
strategyon the server andcustomClasson 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' } }.
Press Run tests to start the app. Its log appears here.Graded endpoints
The client published { cmd: 'find-all' } to 'requests', the server found the handler and published its reply to 'replies.gateway' with the request's id
The data travels with the pattern, and the reply goes to the call that asked
The handler's RpcException became an err packet through Server.send(), and the gateway's operator made it a 404
The handler returns an Observable: two replies with the same id, the last marked isDisposed, and the client passes each one on
The server replies with the error itself, so the caller gets a 502 instead of waiting forever
emit() reaches the client's dispatchEvent(), which publishes to 'events'
The server passed the event to Nest's handleEvent(), and the handler recorded the adopter
The bus carried six requests and one event: none of this went over TCP