From 29449d6403ec257c867e13e1f6936a85781a7405 Mon Sep 17 00:00:00 2001 From: swapnil <78632212+swapnilpaliwal-sd@users.noreply.github.com> Date: Wed, 30 Sep 2026 02:38:00 -0700 Subject: [PATCH] typescript: join a message key's sender to its handler (remote_edge) A producer that emits on a key (a Nest microservices client's emit/send, a kafkajs producer's send({ topic })) was not joined to the handler registered on the same key (@EventPattern / @MessagePattern, a kafkajs consumer's subscribe), so path stopped at the send and impact of a handler named no sender. destinations.dl gains a messaging section. Both ends are recognised by the receiver's DECLARED type and the package it is imported from (knobs.dl, ts_msg_client_type), never by the method name, so emit on a Node event emitter is not a send. The key is read the way a route is: a literal, a const, a const-object member (now also through `as const`), local or imported. One wrapper hop is bound: a private send(pattern, x) called with the key joins from its caller. A decorator's Transport argument names the broker; an end that names none joins any. One-sided keys are reported as remote_unserved / remote_unsent, an unreadable key as remote_undetermined. remote_unsent is judged by destination, so a handler reached by a send that names no broker is not also reported unsent. Checked: new CLI case cross-process-message (emit on an imported const-object member, send on a const, a key through a wrapper, kafkajs send/subscribe, controls: an EventEmitter emit and a handler nobody sends to) fails 5/7 before and passes 7/7; typescript case suite 212/212, engine suite 99/0. On a monorepo with three services: probes 44 -> 47 of 74, no regression, call edges unchanged, 5 amqp + 2 kafka remote edges added. Co-authored-by: axiomcode-bot[bot] <334110751+axiomcode-bot[bot]@users.noreply.github.com> --- .../engine/config-resolution/knobs.dl | 50 ++++++ .../engine/framework-behavior/destinations.dl | 143 +++++++++++++++++- .../cross-process-message/case.json | 27 ++++ .../src/consumer/items.controller.ts | 23 +++ .../src/producer/microservices.d.ts | 13 ++ .../src/producer/producer.test.ts | 7 + .../src/producer/producer.ts | 39 +++++ .../src/shared/package.json | 1 + .../src/shared/topics.ts | 4 + 9 files changed, 303 insertions(+), 4 deletions(-) create mode 100644 tests/cases/typescript/cross-process-message/case.json create mode 100644 tests/cases/typescript/cross-process-message/src/consumer/items.controller.ts create mode 100644 tests/cases/typescript/cross-process-message/src/producer/microservices.d.ts create mode 100644 tests/cases/typescript/cross-process-message/src/producer/producer.test.ts create mode 100644 tests/cases/typescript/cross-process-message/src/producer/producer.ts create mode 100644 tests/cases/typescript/cross-process-message/src/shared/package.json create mode 100644 tests/cases/typescript/cross-process-message/src/shared/topics.ts diff --git a/graph/typescript/engine/config-resolution/knobs.dl b/graph/typescript/engine/config-resolution/knobs.dl index 3e093b7a..c9064f72 100644 --- a/graph/typescript/engine/config-resolution/knobs.dl +++ b/graph/typescript/engine/config-resolution/knobs.dl @@ -202,3 +202,53 @@ ts_http_client_module("axios"). ts_http_client_factory("create"). .decl ts_http_base_url_key(c0:symbol) ts_http_base_url_key("baseURL"). + +// ── MESSAGING: a key both ends spell (a topic, a queue, a message pattern) ── +// ts_msg_client_type(Specifier, TypeName, Role): a receiver DECLARED as TypeName, +// imported from Specifier, sends or subscribes. Matched on the declared type and the +// package it comes from, never on the method name alone: `emit` on a Node event +// emitter, `send` on an HTTP reply or a socket, `subscribe` on an observable are none +// of these. +// nest_client `.emit(key, data)` / `.send(key, data)` +// kafka_producer `.send({ topic, messages })` +// kafka_consumer `.subscribe({ topic })` / `({ topics: [...] })` +.decl ts_msg_client_type(c0:symbol, c1:symbol, c2:symbol) +ts_msg_client_type("@nestjs/microservices", "ClientProxy", "nest_client"). +ts_msg_client_type("@nestjs/microservices", "ClientKafka", "nest_client"). +ts_msg_client_type("@nestjs/microservices", "ClientRMQ", "nest_client"). +ts_msg_client_type("@nestjs/microservices", "ClientNats", "nest_client"). +ts_msg_client_type("@nestjs/microservices", "ClientRedis", "nest_client"). +ts_msg_client_type("@nestjs/microservices", "ClientMqtt", "nest_client"). +ts_msg_client_type("@nestjs/microservices", "ClientTCP", "nest_client"). +ts_msg_client_type("kafkajs", "Producer", "kafka_producer"). +ts_msg_client_type("kafkajs", "Consumer", "kafka_consumer"). +// ts_msg_send(Role, Method, Transport): the call that sends; the key is argument 0, +// or its `topic` property when the argument is a record (ts_msg_record_key) +.decl ts_msg_send(c0:symbol, c1:symbol, c2:symbol) +ts_msg_send("nest_client", "emit", "message"). +ts_msg_send("nest_client", "send", "message"). +ts_msg_send("kafka_producer", "send", "kafka"). +.decl ts_msg_subscribe(c0:symbol, c1:symbol, c2:symbol) +ts_msg_subscribe("kafka_consumer", "subscribe", "kafka"). +// the property of a record argument that names the key, and the one naming several +.decl ts_msg_record_key(c0:symbol) +ts_msg_record_key("topic"). +.decl ts_msg_record_keys(c0:symbol) +ts_msg_record_keys("topics"). +// ts_msg_handler_decorator(Name): a method decorator whose argument 0 is the key the +// method handles (`@EventPattern(TOPICS.x)`, `@MessagePattern('orders.place')`) +.decl ts_msg_handler_decorator(c0:symbol) +ts_msg_handler_decorator("EventPattern"). +ts_msg_handler_decorator("MessagePattern"). +// ts_msg_transport_member(Written, Transport): the decorator's optional argument 1 +.decl ts_msg_transport_member(c0:symbol, c1:symbol) +ts_msg_transport_member("Transport.KAFKA", "kafka"). +ts_msg_transport_member("Transport.RMQ", "amqp"). +ts_msg_transport_member("Transport.NATS", "nats"). +ts_msg_transport_member("Transport.REDIS", "redis"). +ts_msg_transport_member("Transport.MQTT", "mqtt"). +ts_msg_transport_member("Transport.TCP", "tcp"). +// the transport of an end that does not say which broker it uses (a Nest client is +// bound to one in module configuration): it joins an end of any transport +.decl ts_remote_transport_message(c0:symbol) +ts_remote_transport_message("message"). diff --git a/graph/typescript/engine/framework-behavior/destinations.dl b/graph/typescript/engine/framework-behavior/destinations.dl index 7370082c..35df1c87 100644 --- a/graph/typescript/engine/framework-behavior/destinations.dl +++ b/graph/typescript/engine/framework-behavior/destinations.dl @@ -29,13 +29,18 @@ // literal, a template, a `+` concatenation, a constant (in this module or imported), // and a property of a const object literal (a config module). // +// COVERED: MESSAGING (section 6). A key (a topic, a queue, a message pattern) sent by +// a Nest microservices client (`emit` / `send`) or a kafkajs producer, joined to a +// handler registered on the same key (`@EventPattern` / `@MessagePattern`, a kafkajs +// consumer's `subscribe`), the key read as a route is; one wrapper hop is bound. +// // NOT COVERED YET, stated rather than left to be discovered: // - a router mounted under a prefix (`app.use('/api', router)`): the routes of the // router are compared without it; // - an object-form registration (`fastify.route({ method, url, handler })`); // - a URL built by a wrapper the client calls with the path (the C# rules bind one // hop of that; here the wrapper's own send is reported undetermined); -// - messaging (a topic, a queue) and gRPC. +// - gRPC. // ============================================================================ // remote_edge, remote_unserved, remote_unsent and remote_undetermined are declared in @@ -81,11 +86,17 @@ rd_demand(val) :- rd_prop_value(_, val). rd_demand(c) :- rd_demand(e), expr_kind("client", k, _, e), expr_kind_is_transparent(k), expr_child("client", e, _, _, c). +// the object literal a const is initialised with, through `as const` / `satisfies T` +.decl rd_init_obj(c0:symbol, c1:symbol) +rd_init_obj(v, o) :- var_initializer("client", _, o, v), expr_kind("client", "OBJECT_LITERAL", _, o). +rd_init_obj(v, o) :- var_initializer("client", _, i, v), expr_kind("client", k, _, i), expr_kind_is_transparent(k), + expr_child("client", i, _, _, o), expr_kind("client", "OBJECT_LITERAL", _, o). + // `config.reportsUrl`, where `config` is a const initialised with an object literal rd_prop_value(e, val) :- rd_demand(e), expr_kind("client", "PROPERTY_ACCESS", _, e), expr_child("client", e, "RECEIVER", _, r), expr_child("client", e, "PROPERTY_NAME", _, pn), expr_literal_value("client", key, pn), - rd_denotes_var(r, v), var_initializer("client", _, o, v), rd_obj_prop(o, key, val). + rd_denotes_var(r, v), rd_init_obj(v, o), rd_obj_prop(o, key, val). .decl rd_is_str(c0:symbol) .decl rd_is_tmpl(c0:symbol) @@ -411,5 +422,129 @@ remote_undetermined(from, tr, "unresolved_destination") :- ts_remote_transport_h remote_edge(from, to, tr, d, conf) :- remote_edge_via(_, from, to, tr, d, conf). // judged per SEND SITE: a site that linked is not also reported under another spelling remote_unserved(m, tr, d) :- remote_send_at(e, m, tr, d), !remote_edge_via(e, _, _, _, _, _). -// a handler nothing here sends to: its client is in another repository -remote_unsent(m, tr, d) :- remote_serves(m, tr, d), !remote_edge(_, m, tr, d, _). +// a handler nothing here sends to: its client is in another repository. Judged by +// destination alone: a message sent without naming its broker (transport "message") +// is reported under the handler's broker, and still serves it. +.decl remote_reached_at(c0:symbol, c1:symbol) +remote_reached_at(m, d) :- remote_edge(_, m, _, d, _). +remote_unsent(m, tr, d) :- remote_serves(m, tr, d), !remote_reached_at(m, d). + +// ───────────────────────────────────────────────────────────────────────────── +// 6. MESSAGING: a key both ends spell +// ───────────────────────────────────────────────────────────────────────────── +// A producer sends on a KEY (a topic, a queue, a message pattern) and a handler is +// registered on one; neither calls the other. The key is read exactly as a route is +// (section 1): a literal, a const, a const-object member, imported or local. Both +// ends are recognised by what the RECEIVER is declared as, a type imported from the +// messaging package (config-resolution/knobs.dl, ts_msg_client_type), never by the +// method name: `emit` on a Node event emitter is not a broker send. +// +// NOT COVERED YET: a pattern object (`@MessagePattern({ cmd: 'x' })`), a key held in a +// field set at run time (`this.topic`), a receiver whose type is only inferred from a +// factory call in a package outside the graph (`kafka.producer()` with no annotation). + +// the receivers of a call that could be a send or a subscribe +.decl rm_candidate_recv(c0:symbol) +rm_candidate_recv(q) :- call_site("client", "METHOD_CALL", cn, _, q, _, _), ts_msg_send(_, cn, _). +rm_candidate_recv(q) :- call_site("client", "METHOD_CALL", cn, _, q, _, _), ts_msg_subscribe(_, cn, _). +// `this.client`: the field it names, a parameter property included +.decl rm_this_field(c0:symbol, c1:symbol) +rm_this_field(q, f) :- rm_candidate_recv(q), property_access_name(q, n), + expr_child("client", q, "RECEIVER", _, t), expr_type(t, _, tt), field_in_scope(tt, n, "false", f). +// the type reference the receiver is DECLARED with +.decl rm_recv_ref(c0:symbol, c1:symbol) +rm_recv_ref(q, r) :- rm_candidate_recv(q), expr_referenced("client", "PARAMETER", p, q), param_type_ref("client", r, p). +rm_recv_ref(q, r) :- rm_candidate_recv(q), expr_referenced("client", "VARIABLE", v, q), var_type_ref("client", r, v). +rm_recv_ref(q, r) :- rm_this_field(q, f), field_type_ref("client", r, f). +rm_recv_ref(q, r) :- rm_this_field(q, f), param_declares_field("client", f, p), param_type_ref("client", r, p). +// the package a type name is imported from, under the name it exports +// (`import type { ClientProxy as Client } from '@nestjs/microservices'`) +.decl rm_ref_import(c0:symbol, c1:symbol, c2:symbol) +rm_ref_import(r, spec, orig) :- rm_recv_ref(_, r), type_ref("client", _, _, tn, _, _, _, r), type_ref_module("client", m, r), + import_binding("client", _, tn, orig, m, ih), orig != "", rd_import_path(ih, spec). +rm_ref_import(r, spec, tn) :- rm_recv_ref(_, r), type_ref("client", _, _, tn, _, _, _, r), type_ref_module("client", m, r), + import_binding("client", _, tn, "", m, ih), rd_import_path(ih, spec). +.decl rm_recv_role(c0:symbol, c1:symbol) +rm_recv_role(q, role) :- rm_recv_ref(q, r), rm_ref_import(r, spec, name), ts_msg_client_type(spec, name, role). + +// (call, role, transport) +.decl rm_send_call(c0:symbol, c1:symbol, c2:symbol) +rm_send_call(ce, role, tr) :- call_site("client", "METHOD_CALL", cn, _, q, ce, _), ts_msg_send(role, cn, tr), + rm_recv_role(q, role). +.decl rm_sub_call(c0:symbol, c1:symbol, c2:symbol) +rm_sub_call(ce, role, tr) :- call_site("client", "METHOD_CALL", cn, _, q, ce, _), ts_msg_subscribe(role, cn, tr), + rm_recv_role(q, role). + +// the expression holding the key: argument 0, or the `topic` of a record argument, or +// each element of its `topics` +.decl rm_key_arg(c0:symbol, c1:symbol) +rm_key_arg(ce, a) :- rm_send_call(ce, "nest_client", _), expr_child("client", ce, "ARGUMENT", "0", a). +rm_key_arg(ce, v) :- rm_send_call(ce, role, _), role != "nest_client", expr_child("client", ce, "ARGUMENT", "0", o), + rd_obj_prop(o, k, v), ts_msg_record_key(k). +rm_key_arg(ce, v) :- rm_sub_call(ce, _, _), expr_child("client", ce, "ARGUMENT", "0", o), + rd_obj_prop(o, k, v), ts_msg_record_key(k). +rm_key_arg(ce, el) :- rm_sub_call(ce, _, _), expr_child("client", ce, "ARGUMENT", "0", o), + rd_obj_prop(o, k, arr), ts_msg_record_keys(k), expr_kind("client", "ARRAY_LITERAL", _, arr), + expr_child("client", arr, _, _, el). +rd_demand(a) :- rm_key_arg(_, a). +// a key is a value with no computed part (a key expression of either end) +.decl rm_key_expr(c0:symbol) +rm_key_expr(a) :- rm_key_arg(_, a). +rm_key_expr(a) :- rm_handler_key_arg(_, _, a). +.decl rm_key(c0:symbol, c1:symbol) +rm_key(e, k) :- rm_key_expr(e), rd_val(e, k), k != "", !contains("{}", k). + +// ONE WRAPPER HOP: the key is a parameter of the method that sends, and each call of +// that method names it: `place() { return this.send(CMD.place, x); }` over +// `private send(pattern, x) { return this.proxy.send(pattern, x); }`. The send is then +// judged at the CALL of the wrapper, from the method that makes it. +.decl rm_key_param(c0:symbol, c1:symbol, c2:symbol) +rm_key_param(ce, w, pos) :- rm_key_arg(ce, a), expr_referenced("client", "PARAMETER", p, a), + param_decl("client", _, pos, _, w, p), expr_enclosing_method(ce, w). +.decl rm_wrap_arg(c0:symbol, c1:symbol, c2:symbol) +rm_wrap_arg(ce, oc, a) :- rm_key_param(ce, w, pos), call_chain_edge(oc, _, "-", w, "client", _, _), + expr_child("client", oc, "ARGUMENT", pos, a). +rd_demand(a) :- rm_wrap_arg(_, _, a). + +// (site, sending method, transport, key) +.decl rm_send(c0:symbol, c1:symbol, c2:symbol, c3:symbol) +rm_send(ce, from, tr, k) :- rm_send_call(ce, _, tr), rm_key_arg(ce, a), rm_key(a, k), expr_enclosing_method(ce, from). +rm_send(oc, from, tr, k) :- rm_send_call(ce, _, tr), rm_wrap_arg(ce, oc, a), rd_val(a, k), k != "", !contains("{}", k), + expr_enclosing_method(oc, from). + +// ── the handler end ───────────────────────────────────────────────────────── +// a decorated method: `@EventPattern(key, Transport.KAFKA)`; the argument's expression +// is read like any other +.decl rm_handler_key_arg(c0:symbol, c1:symbol, c2:symbol) +rm_handler_key_arg(m, d, a) :- annotation_on("client", n, _, "METHOD_DECLARATION", m, d), ts_msg_handler_decorator(n), + ts_decorator_argument(_, _, _, "0", d, _, _, _, _, _, a, _, _), a != "". +rd_demand(a) :- rm_handler_key_arg(_, _, a). +.decl rm_handler_tr(c0:symbol, c1:symbol) +rm_handler_tr(d, tr) :- rm_handler_key_arg(_, d, _), annotation_arg("client", _, v, _, "1", _, d, _), + ts_msg_transport_member(v, tr). +.decl rm_handler_has_tr(c0:symbol) +rm_handler_has_tr(d) :- rm_handler_tr(d, _). +// (handler, transport, key) +.decl rm_serve(c0:symbol, c1:symbol, c2:symbol) +rm_serve(m, tr, k) :- rm_handler_key_arg(m, d, a), rm_key(a, k), rm_handler_tr(d, tr). +rm_serve(m, tr, k) :- rm_handler_key_arg(m, d, a), rm_key(a, k), !rm_handler_has_tr(d), ts_remote_transport_message(tr). +// a subscription: the function that subscribes is where the messages arrive +rm_serve(m, tr, k) :- rm_sub_call(ce, _, tr), rm_key_arg(ce, a), rm_key(a, k), expr_enclosing_method(ce, m). + +// ── the join ──────────────────────────────────────────────────────────────── +// Same key; the transports agree, or one end does not say. The edge is reported under +// the end that names its broker. +.decl rm_edge_tr(c0:symbol, c1:symbol, c2:symbol) +rm_edge_tr(a, a, a) :- rm_send(_, _, a, _). +rm_edge_tr(a, b, b) :- rm_send(_, _, a, _), rm_serve(_, b, _), ts_remote_transport_message(a), a != b. +rm_edge_tr(a, b, a) :- rm_send(_, _, a, _), rm_serve(_, b, _), ts_remote_transport_message(b), a != b. +remote_edge_via(e, from, to, tr, k, conf) :- rm_send(e, from, st, k), rm_serve(to, ht, k), rm_edge_tr(st, ht, tr), + from != to, ts_remote_confidence_exact(conf). +remote_serves(m, tr, k) :- rm_serve(m, tr, k). +remote_send_at(e, from, tr, k) :- rm_send(e, from, tr, k). +// a send whose key cannot be read, and that no caller names either +.decl rm_send_keyed(c0:symbol) +rm_send_keyed(ce) :- rm_send(ce, _, _, _). +rm_send_keyed(ce) :- rm_wrap_arg(ce, oc, _), rm_send(oc, _, _, _). +remote_undetermined(from, tr, "unresolved_destination") :- rm_send_call(ce, _, tr), !rm_send_keyed(ce), + expr_enclosing_method(ce, from). diff --git a/tests/cases/typescript/cross-process-message/case.json b/tests/cases/typescript/cross-process-message/case.json new file mode 100644 index 00000000..f9b9ad41 --- /dev/null +++ b/tests/cases/typescript/cross-process-message/case.json @@ -0,0 +1,27 @@ +{"lang": "typescript", "src": "src", + "checks": [ + {"why": "a handler registered on a message key is depended on by the method that emits on the same key, read through an imported const-object member: a remote hop naming the transport and the key", + "run": ["impact", "ItemsController.onCreated"], + "want": ["[remote] ItemPublisher.announce", "across a process boundary (kafka) at item.created.v1"], + "avoid": ["LocalBus.fire", "ItemPublisher.shout"]}, + {"why": "path joins a request/reply send on a plain const to the handler of the same pattern", + "run": ["path", "ItemPublisher.place", "ItemsController.placeItem"], + "want": ["connected across a process: ItemPublisher.place → ItemsController.placeItem", "[amqp] at item.place (exact)"], + "expect_error": true}, + {"why": "a key passed into a private wrapper that sends it: the edge starts at the caller that names the key", + "run": ["path", "CommandClient.cancel", "ItemsController.cancelItem"], + "want": ["connected across a process: CommandClient.cancel → ItemsController.cancelItem", "at item.cancel (exact)"], + "expect_error": true}, + {"why": "kafkajs: the topic of a producer record joins the function that subscribes to it", + "run": ["impact", "listenRemovals"], + "want": ["[remote] RawProducer.removed", "at item.removed.v1"]}, + {"why": "CONTROL: a Node event emitter's emit is not a broker send, even on the same key", + "run": ["impact", "LocalBus.fire"], + "avoid": ["[remote]", "onCreated"]}, + {"why": "CONTROL: a handler on a key nothing here sends to has no sender", + "run": ["impact", "ItemsController.unused"], + "avoid": ["[remote]"]}, + {"why": "test-impact crosses the key: the test that drives the producer is selected on the remote rung", + "run": ["impact", "ItemsController.onCreated", "--tests-only"], + "want": ["producer.test.ts"]} + ]} diff --git a/tests/cases/typescript/cross-process-message/src/consumer/items.controller.ts b/tests/cases/typescript/cross-process-message/src/consumer/items.controller.ts new file mode 100644 index 00000000..d601de60 --- /dev/null +++ b/tests/cases/typescript/cross-process-message/src/consumer/items.controller.ts @@ -0,0 +1,23 @@ +import { EventPattern, MessagePattern, Transport } from '@nestjs/microservices'; +import type { Consumer } from 'kafkajs'; +import { CMD, COMMANDS, TOPICS } from '@acme/shared'; + +export class ItemsController { + @EventPattern(TOPICS.created, Transport.KAFKA) + onCreated(event: unknown) { return event; } + + @MessagePattern(CMD, Transport.RMQ) + placeItem(command: unknown) { return command; } + + @MessagePattern(COMMANDS.cancel) + cancelItem(command: unknown) { return command; } + + // no sender here: an unsent handler + @EventPattern('other.topic') + unused(event: unknown) { return event; } +} + +export async function listenRemovals(consumer: Consumer) { + await consumer.subscribe({ topic: TOPICS.removed }); + await consumer.run({ eachMessage: async () => undefined }); +} diff --git a/tests/cases/typescript/cross-process-message/src/producer/microservices.d.ts b/tests/cases/typescript/cross-process-message/src/producer/microservices.d.ts new file mode 100644 index 00000000..adb789e9 --- /dev/null +++ b/tests/cases/typescript/cross-process-message/src/producer/microservices.d.ts @@ -0,0 +1,13 @@ +declare module '@nestjs/microservices' { + export class ClientProxy { + emit(pattern: unknown, data: unknown): unknown; + send(pattern: unknown, data: unknown): unknown; + } + export enum Transport { KAFKA, RMQ } + export function EventPattern(pattern: unknown, transport?: unknown): MethodDecorator; + export function MessagePattern(pattern: unknown, transport?: unknown): MethodDecorator; +} +declare module 'kafkajs' { + export interface Producer { send(record: { topic: string; messages: unknown[] }): Promise; } + export interface Consumer { subscribe(s: { topic?: string; topics?: string[] }): Promise; run(c: unknown): Promise; } +} diff --git a/tests/cases/typescript/cross-process-message/src/producer/producer.test.ts b/tests/cases/typescript/cross-process-message/src/producer/producer.test.ts new file mode 100644 index 00000000..0c25016a --- /dev/null +++ b/tests/cases/typescript/cross-process-message/src/producer/producer.test.ts @@ -0,0 +1,7 @@ +import { it, expect } from 'vitest'; +import { ItemPublisher } from './producer'; + +it('announces an item', () => { + const publisher = new ItemPublisher({ emit: () => 1, send: () => 1 } as never); + expect(publisher.announce('a')).toBeDefined(); +}); diff --git a/tests/cases/typescript/cross-process-message/src/producer/producer.ts b/tests/cases/typescript/cross-process-message/src/producer/producer.ts new file mode 100644 index 00000000..60c09c98 --- /dev/null +++ b/tests/cases/typescript/cross-process-message/src/producer/producer.ts @@ -0,0 +1,39 @@ +import { EventEmitter } from 'events'; +import type { ClientProxy } from '@nestjs/microservices'; +import type { Producer } from 'kafkajs'; +import { CMD, COMMANDS, TOPICS } from '../shared/topics'; + +export class ItemPublisher { + constructor(private readonly client: ClientProxy) {} + + // emit on a const-object member: joins the handler that spells the same member + announce(id: string) { return this.client.emit(TOPICS.created, { id }); } + + // send on a plain const: joins the request/reply handler + place(id: string) { return this.client.send(CMD, { id }); } + + // a key nothing here handles: an unserved send + shout() { return this.client.emit('nobody.listens', {}); } +} + +export class CommandClient { + constructor(private readonly proxy: ClientProxy) {} + + // the key is the wrapper's argument: the send joins from here + cancel(id: string) { return this.dispatch(COMMANDS.cancel, { id }); } + + private dispatch(pattern: string, body: unknown) { return this.proxy.send(pattern, body); } +} + +export class RawProducer { + constructor(private readonly producer: Producer) {} + + // kafkajs: the topic is a property of the record + removed(id: string) { return this.producer.send({ topic: TOPICS.removed, messages: [{ value: id }] }); } +} + +// CONTROL: a Node event emitter's emit is not a broker send, however its key reads +export class LocalBus { + private readonly emitter = new EventEmitter(); + fire() { this.emitter.emit(TOPICS.created, {}); } +} diff --git a/tests/cases/typescript/cross-process-message/src/shared/package.json b/tests/cases/typescript/cross-process-message/src/shared/package.json new file mode 100644 index 00000000..6a3a4605 --- /dev/null +++ b/tests/cases/typescript/cross-process-message/src/shared/package.json @@ -0,0 +1 @@ +{ "name": "@acme/shared", "version": "0.0.0", "main": "./dist/topics.js", "types": "./dist/topics.d.ts" } diff --git a/tests/cases/typescript/cross-process-message/src/shared/topics.ts b/tests/cases/typescript/cross-process-message/src/shared/topics.ts new file mode 100644 index 00000000..dd3684d2 --- /dev/null +++ b/tests/cases/typescript/cross-process-message/src/shared/topics.ts @@ -0,0 +1,4 @@ +// The keys both ends spell: a const-object member, and a plain const. +export const TOPICS = { created: 'item.created.v1', removed: 'item.removed.v1' } as const; +export const CMD = 'item.place'; +export const COMMANDS = { cancel: 'item.cancel' } as const;