Skip to content
New issue

Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.

By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.

Already on GitHub? Sign in to your account

feat(nestjs-cqrs-kafka-events): init #353

Merged
merged 5 commits into from
Mar 14, 2025
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
2 changes: 1 addition & 1 deletion .github/ISSUE_TEMPLATE/bug.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -58,4 +58,4 @@ body:
validations:
required: true

projects: ['atls/11']
projects: ['atls/11']
2 changes: 1 addition & 1 deletion .github/ISSUE_TEMPLATE/docs.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -41,4 +41,4 @@ body:
validations:
required: true

projects: ['atls/11']
projects: ['atls/11']
2 changes: 1 addition & 1 deletion .github/ISSUE_TEMPLATE/feature.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -44,4 +44,4 @@ body:
validations:
required: true

projects: ['atls/11']
projects: ['atls/11']
188 changes: 158 additions & 30 deletions .pnp.cjs

Large diffs are not rendered by default.

3 changes: 3 additions & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -13,5 +13,8 @@
"typescript": "5.4.2"
},
"packageManager": "[email protected]",
"formatterIgnorePatterns": [
"**/*/CHANGELOG.md"
],
"typecheckSkipLibCheck": true
}
49 changes: 49 additions & 0 deletions packages/nestjs-cqrs-kafka-events/package.json
Original file line number Diff line number Diff line change
@@ -0,0 +1,49 @@
{
"name": "@atls/nestjs-cqrs-kafka-events",
"version": "0.0.0",
"license": "BSD-3-Clause",
"type": "module",
"exports": {
"./package.json": "./package.json",
".": "./src/index.ts"
},
"main": "src/index.ts",
"files": [
"dist"
],
"scripts": {
"build": "yarn library build",
"prepack": "yarn run build",
"postpack": "rm -rf dist"
},
"dependencies": {
"@atls/nestjs-kafka": "workspace:0.0.2",
"telejson": "7.2.0"
},
"devDependencies": {
"@nestjs/common": "10.4.15",
"@nestjs/core": "10.4.15",
"@nestjs/cqrs": "10.2.8",
"reflect-metadata": "0.2.2",
"rxjs": "^7.8.1"
},
"peerDependencies": {
"@nestjs/common": "^10",
"@nestjs/core": "^10",
"@nestjs/cqrs": "^10",
"reflect-metadata": "^0.2",
"rxjs": "^7"
},
"publishConfig": {
"exports": {
"./package.json": "./package.json",
".": {
"import": "./dist/index.js",
"types": "./dist/index.d.ts",
"default": "./dist/index.js"
}
},
"main": "dist/index.js",
"typings": "dist/index.d.ts"
}
}
4 changes: 4 additions & 0 deletions packages/nestjs-cqrs-kafka-events/src/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,4 @@
export * from '@atls/nestjs-kafka'

export * from './messaging/index.js'
export * from './module/index.js'
2 changes: 2 additions & 0 deletions packages/nestjs-cqrs-kafka-events/src/messaging/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,2 @@
export * from './kafka.subscriber.js'
export * from './kafka.publisher.js'
40 changes: 40 additions & 0 deletions packages/nestjs-cqrs-kafka-events/src/messaging/kafka.publisher.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,40 @@
import type { Kafka } from '@atls/nestjs-kafka'
import type { Producer } from '@atls/nestjs-kafka'
import type { RecordMetadata } from '@atls/nestjs-kafka'
import type { OnModuleDestroy } from '@nestjs/common'
import type { IEventPublisher } from '@nestjs/cqrs'
import type { IEvent } from '@nestjs/cqrs'

import { stringify } from 'telejson'

export class KafkaPublisher implements IEventPublisher, OnModuleDestroy {
private readonly kafkaProducer: Producer

constructor(kafka: Kafka) {
this.kafkaProducer = kafka.producer()
}

async onModuleDestroy(): Promise<void> {
await this.kafkaProducer.disconnect()
}

async connect(): Promise<void> {
await this.kafkaProducer.connect()
}

async publish(event: IEvent): Promise<Array<RecordMetadata>> {
return this.kafkaProducer.send({
topic: event.constructor.name,
messages: [{ value: stringify(event) }],
})
}

async publishAll(events: Array<IEvent>): Promise<Array<RecordMetadata>> {
return this.kafkaProducer.sendBatch({
topicMessages: events.map((event) => ({
topic: event.constructor.name,
messages: [{ value: stringify(event) }],
})),
})
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,50 @@
import type { Consumer } from '@atls/nestjs-kafka'
import type { Kafka } from '@atls/nestjs-kafka'
import type { OnModuleDestroy } from '@nestjs/common'
import type { IEvent } from '@nestjs/cqrs'
import type { IMessageSource } from '@nestjs/cqrs'
import type { Subject } from 'rxjs'

import { parse } from 'telejson'

export class KafkaSubscriber implements IMessageSource, OnModuleDestroy {
private readonly kafkaConsumer: Consumer

// eslint-disable-next-line @typescript-eslint/no-explicit-any
private bridge!: Subject<any>

constructor(kafka: Kafka, groupId: string) {
this.kafkaConsumer = kafka.consumer({ groupId })
}

async onModuleDestroy(): Promise<void> {
await this.kafkaConsumer.disconnect()
}

async connect(events: Array<FunctionConstructor>): Promise<void> {
await this.kafkaConsumer.connect()

for await (const event of events) {
await this.kafkaConsumer.subscribe({ topic: event.name, fromBeginning: false })
}

await this.kafkaConsumer.run({
eachMessage: async ({ topic, message }) => {
if (this.bridge) {
for (const Event of events) {
if (Event.name === topic) {
const parsedJson = parse((message.value || '').toString())
const receivedEvent: IEvent = Object.assign(new Event(), parsedJson)

this.bridge.next(receivedEvent)
}
}
}
},
})
}

bridgeEventsTo<T extends IEvent>(subject: Subject<T>): void {
this.bridge = subject
}
}
Original file line number Diff line number Diff line change
@@ -0,0 +1 @@
export const CQRS_KAFKA_EVENTS_MODULE_OPTIONS = Symbol('cqrs-kafka-events-module-options')
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
import type { KafkaConfig } from '@atls/nestjs-kafka'
import type { ModuleMetadata } from '@nestjs/common/interfaces'
import type { Type } from '@nestjs/common/interfaces'

export interface CqrsKafkaEventsModuleOptions extends Partial<KafkaConfig> {
groupId?: string
}

export interface CqrsKafkaEventsOptionsFactory {
createCqrsKafkaEventsOptions: () =>
| CqrsKafkaEventsModuleOptions
| Promise<CqrsKafkaEventsModuleOptions>
}

export interface CqrsKafkaEventsModuleAsyncOptions extends Pick<ModuleMetadata, 'imports'> {
useExisting?: Type<CqrsKafkaEventsOptionsFactory>
useClass?: Type<CqrsKafkaEventsOptionsFactory>
useFactory?: (
...args: Array<any>
) => CqrsKafkaEventsModuleOptions | Promise<CqrsKafkaEventsModuleOptions>
// eslint-disable-next-line @typescript-eslint/no-explicit-any
inject?: Array<any>
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,150 @@
import type { KafkaConfig } from '@atls/nestjs-kafka'
import type { DynamicModule } from '@nestjs/common'
import type { OnModuleInit } from '@nestjs/common'
import type { Provider } from '@nestjs/common'

import type { CqrsKafkaEventsModuleOptions } from './cqrs-kafka-events.module.interfaces.js'
import type { CqrsKafkaEventsModuleAsyncOptions } from './cqrs-kafka-events.module.interfaces.js'
import type { CqrsKafkaEventsOptionsFactory } from './cqrs-kafka-events.module.interfaces.js'

import { Module } from '@nestjs/common'
import { EventBus } from '@nestjs/cqrs'
import { EVENTS_HANDLER_METADATA } from '@nestjs/cqrs/dist/decorators/constants.js'
import { ExplorerService } from '@nestjs/cqrs/dist/services/explorer.service.js'

import { KafkaModule } from '@atls/nestjs-kafka'
import { KafkaFactory } from '@atls/nestjs-kafka'
import { Kafka } from '@atls/nestjs-kafka'

import { KafkaPublisher } from '../messaging/index.js'
import { KafkaSubscriber } from '../messaging/index.js'
import { CQRS_KAFKA_EVENTS_MODULE_OPTIONS } from './cqrs-kafka-events.module.constants.js'

@Module({})
export class CqrsKafkaEventsModule implements OnModuleInit {
constructor(
private readonly eventBus: EventBus,
private readonly kafkaPublisher: KafkaPublisher,
private readonly kafkaSubscriber: KafkaSubscriber,
private readonly explorerService: ExplorerService
) {}

static register(options: CqrsKafkaEventsModuleOptions): DynamicModule {
return {
module: CqrsKafkaEventsModule,
imports: [KafkaModule.register(options)],
providers: [
{
provide: CQRS_KAFKA_EVENTS_MODULE_OPTIONS,
useValue: options,
},
{
provide: Kafka,
useFactory: (kafkaFactory: KafkaFactory, config: Partial<KafkaConfig>): Kafka =>
kafkaFactory.create(config),
inject: [KafkaFactory, CQRS_KAFKA_EVENTS_MODULE_OPTIONS],
},
{
provide: KafkaSubscriber,
useFactory: (
kafka: Kafka,
moduleOptions: CqrsKafkaEventsModuleOptions
): KafkaSubscriber =>
new KafkaSubscriber(
kafka,
moduleOptions.groupId || process.env.CQRS_KAFKA_EVENTS_GROUP_ID || 'default'
),
inject: [Kafka, CQRS_KAFKA_EVENTS_MODULE_OPTIONS],
},
{
provide: KafkaPublisher,
useFactory: (kafka: Kafka): KafkaPublisher => new KafkaPublisher(kafka),
inject: [Kafka],
},
],
}
}

static registerAsync(options: CqrsKafkaEventsModuleAsyncOptions): DynamicModule {
return {
module: CqrsKafkaEventsModule,
imports: [KafkaModule.register(), ...(options.imports || [])],
providers: [
...this.createAsyncProviders(options),
{
provide: Kafka,
useFactory: (kafkaFactory: KafkaFactory, config: Partial<KafkaConfig>): Kafka =>
kafkaFactory.create(config),
inject: [KafkaFactory, CQRS_KAFKA_EVENTS_MODULE_OPTIONS],
},
{
provide: KafkaSubscriber,
useFactory: (
kafka: Kafka,
moduleOptions: CqrsKafkaEventsModuleOptions
): KafkaSubscriber =>
new KafkaSubscriber(
kafka,
moduleOptions.groupId || process.env.CQRS_KAFKA_EVENTS_GROUP_ID || 'default'
),
inject: [Kafka, CQRS_KAFKA_EVENTS_MODULE_OPTIONS],
},
{
provide: KafkaPublisher,
useFactory: (kafka: Kafka): KafkaPublisher => new KafkaPublisher(kafka),
inject: [Kafka],
},
],
}
}

private static createAsyncProviders(options: CqrsKafkaEventsModuleAsyncOptions): Array<Provider> {
if (options.useExisting || options.useFactory) {
return [this.createAsyncOptionsProvider(options)]
}

return [
this.createAsyncOptionsProvider(options),
{
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
provide: options.useClass!,
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
useClass: options.useClass!,
},
]
}

private static createAsyncOptionsProvider(options: CqrsKafkaEventsModuleAsyncOptions): Provider {
if (options.useFactory) {
return {
provide: CQRS_KAFKA_EVENTS_MODULE_OPTIONS,
useFactory: options.useFactory,
inject: options.inject || [],
}
}

return {
provide: CQRS_KAFKA_EVENTS_MODULE_OPTIONS,
useFactory: (
optionsFactory: CqrsKafkaEventsOptionsFactory
): CqrsKafkaEventsModuleOptions | Promise<CqrsKafkaEventsModuleOptions> =>
optionsFactory.createCqrsKafkaEventsOptions(),
// eslint-disable-next-line @typescript-eslint/no-non-null-assertion
inject: [options.useExisting! || options.useClass!],
}
}

async onModuleInit(): Promise<void> {
await this.kafkaPublisher.connect()
await this.kafkaSubscriber.connect(
(this.explorerService.explore().events || [])
.map(
(handler) => Reflect.getMetadata(EVENTS_HANDLER_METADATA, handler) as FunctionConstructor
)
.flat()
)

this.eventBus.publisher = this.kafkaPublisher
this.kafkaSubscriber.bridgeEventsTo(this.eventBus.subject$)
}
}
3 changes: 3 additions & 0 deletions packages/nestjs-cqrs-kafka-events/src/module/index.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,3 @@
export * from './cqrs-kafka-events.module.constants.js'
export * from './cqrs-kafka-events.module.js'
export type * from './cqrs-kafka-events.module.interfaces.js'
3 changes: 2 additions & 1 deletion packages/nestjs-cqrs/package.json
Original file line number Diff line number Diff line change
Expand Up @@ -41,5 +41,6 @@
},
"main": "dist/index.js",
"typings": "dist/index.d.ts"
}
},
"typecheckSkipLibCheck": true
}
Loading
Loading