Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
50 changes: 50 additions & 0 deletions graph/typescript/engine/config-resolution/knobs.dl
Original file line number Diff line number Diff line change
Expand Up @@ -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 `<client>.emit(key, data)` / `<client>.send(key, data)`
// kafka_producer `<producer>.send({ topic, messages })`
// kafka_consumer `<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").
143 changes: 139 additions & 4 deletions graph/typescript/engine/framework-behavior/destinations.dl
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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).
27 changes: 27 additions & 0 deletions tests/cases/typescript/cross-process-message/case.json
Original file line number Diff line number Diff line change
@@ -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"]}
]}
Original file line number Diff line number Diff line change
@@ -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 });
}
Original file line number Diff line number Diff line change
@@ -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<unknown>; }
export interface Consumer { subscribe(s: { topic?: string; topics?: string[] }): Promise<void>; run(c: unknown): Promise<void>; }
}
Original file line number Diff line number Diff line change
@@ -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();
});
Original file line number Diff line number Diff line change
@@ -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, {}); }
}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
{ "name": "@acme/shared", "version": "0.0.0", "main": "./dist/topics.js", "types": "./dist/topics.d.ts" }
Loading
Loading