{"_id":"@armaghanzahid/kafka-module","_rev":"6-e7a950dfb447e411bbb47e38e05c63f1","name":"@armaghanzahid/kafka-module","dist-tags":{"latest":"1.0.5"},"versions":{"1.0.0":{"name":"@armaghanzahid/kafka-module","version":"1.0.0","keywords":["nestjs","kafka","microservices","typescript"],"author":{"name":"HeartBug"},"license":"ISC","_id":"@armaghanzahid/kafka-module@1.0.0","maintainers":[{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"}],"homepage":"https://github.com/your-org/kafka-module#readme","bugs":{"url":"https://github.com/your-org/kafka-module/issues"},"dist":{"shasum":"35c4d90dbbe20be659217fc6ac6e76ac32c1fbcd","tarball":"https://registry.npmjs.org/@armaghanzahid/kafka-module/-/kafka-module-1.0.0.tgz","fileCount":22,"integrity":"sha512-qeBVkBCc2rfaJzapfizvU/V4np10elQXCeMW++dqK3EtTkdQPyX2JhXSPuv1xrbU1Nac+927tKjPZKdbfuAR4w==","signatures":[{"sig":"MEQCIGACq5+dyqYGYaDnsj4ltW36eeRaI371tCEDZ1D5kM3OAiAnTYidR/v41RO3orI/AcKIvHFn03b0ETuRedRoJE1NIQ==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":120000},"main":"dist/index.js","types":"dist/index.d.ts","engines":{"node":">=18.0.0"},"gitHead":"ff655eef276270bb92f58d7b5abeab3e6e6cfac8","scripts":{"lint":"eslint \"{src,apps,libs,test}/**/*.ts\" --fix","test":"jest","build":"nest build","format":"prettier --write \"src/**/*.ts\"","test:cov":"jest --coverage","test:e2e":"jest --config ./test/jest-e2e.json","test:debug":"node --inspect-brk -r tsconfig-paths/register -r ts-node/register node_modules/.bin/jest --runInBand","test:watch":"jest --watch","prepublishOnly":"npm run build"},"_npmUser":{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"},"repository":{"url":"git+https://github.com/your-org/kafka-module.git","type":"git"},"_npmVersion":"10.7.0","description":"A robust NestJS module for Kafka integration in HeartBug microservices","directories":{},"_nodeVersion":"20.14.0","dependencies":{"rxjs":"^7.8.1","kafkajs":"^2.2.4","@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0","reflect-metadata":"^0.1.13"},"publishConfig":{"access":"public","registry":"https://registry.npmjs.org/"},"_hasShrinkwrap":false,"devDependencies":{"jest":"^29.5.0","eslint":"^8.42.0","ts-jest":"^29.1.0","ts-node":"^10.9.1","prettier":"^3.0.0","supertest":"^6.3.3","ts-loader":"^9.4.3","typescript":"^5.1.3","@nestjs/cli":"^10.0.0","@types/jest":"^29.5.2","@types/node":"^20.3.1","@types/express":"^4.17.17","tsconfig-paths":"^4.2.0","@nestjs/testing":"^10.0.0","@types/supertest":"^2.0.12","@nestjs/schematics":"^10.0.0","source-map-support":"^0.5.21","eslint-config-prettier":"^9.0.0","eslint-plugin-prettier":"^5.0.0","@typescript-eslint/parser":"^6.0.0","@typescript-eslint/eslint-plugin":"^6.0.0"},"peerDependencies":{"@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0"},"_npmOperationalInternal":{"tmp":"tmp/kafka-module_1.0.0_1745460250954_0.2128517904209548","host":"s3://npm-registry-packages-npm-production"}},"1.0.1":{"name":"@armaghanzahid/kafka-module","version":"1.0.1","keywords":["nestjs","kafka","microservices","typescript"],"author":{"name":"HeartBug"},"license":"ISC","_id":"@armaghanzahid/kafka-module@1.0.1","maintainers":[{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"}],"homepage":"https://github.com/your-org/kafka-module#readme","bugs":{"url":"https://github.com/your-org/kafka-module/issues"},"dist":{"shasum":"f4fe6ee14a15f7d1ed9e0a635a50ce8ec5da19dc","tarball":"https://registry.npmjs.org/@armaghanzahid/kafka-module/-/kafka-module-1.0.1.tgz","fileCount":22,"integrity":"sha512-i/E/K3So+/qNQVLiPDmyWgJQI0eTsf33ei+ftzKJpDlkLK5obx6RRa16U5TvMzxvsYDzfIACCaTdyXsf9jvbGA==","signatures":[{"sig":"MEUCIQDcbCPPDZLQWqzKZvNQ0XlLaZPldSO7NEZlXV3eBQkLhQIgexQtNTJmwLW6PZ52xNOG3dHPz8BdWDRS3a7BTMhC2SM=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":122712},"main":"dist/index.js","types":"dist/index.d.ts","engines":{"node":">=18.0.0"},"gitHead":"ff655eef276270bb92f58d7b5abeab3e6e6cfac8","scripts":{"lint":"eslint \"{src,apps,libs,test}/**/*.ts\" --fix","test":"jest","build":"nest build","format":"prettier --write \"src/**/*.ts\"","test:cov":"jest --coverage","test:e2e":"jest --config ./test/jest-e2e.json","test:debug":"node --inspect-brk -r tsconfig-paths/register -r ts-node/register node_modules/.bin/jest --runInBand","test:watch":"jest --watch","prepublishOnly":"npm run build"},"_npmUser":{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"},"repository":{"url":"git+https://github.com/your-org/kafka-module.git","type":"git"},"_npmVersion":"10.7.0","description":"A robust NestJS module for Kafka integration in HeartBug microservices","directories":{},"_nodeVersion":"20.14.0","dependencies":{"rxjs":"^7.8.1","kafkajs":"^2.2.4","@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0","reflect-metadata":"^0.1.13"},"publishConfig":{"access":"public","registry":"https://registry.npmjs.org/"},"_hasShrinkwrap":false,"devDependencies":{"jest":"^29.5.0","eslint":"^8.42.0","ts-jest":"^29.1.0","ts-node":"^10.9.1","prettier":"^3.0.0","supertest":"^6.3.3","ts-loader":"^9.4.3","typescript":"^5.1.3","@nestjs/cli":"^10.0.0","@types/jest":"^29.5.2","@types/node":"^20.3.1","@types/express":"^4.17.17","tsconfig-paths":"^4.2.0","@nestjs/testing":"^10.0.0","@types/supertest":"^2.0.12","@nestjs/schematics":"^10.0.0","source-map-support":"^0.5.21","eslint-config-prettier":"^9.0.0","eslint-plugin-prettier":"^5.0.0","@typescript-eslint/parser":"^6.0.0","@typescript-eslint/eslint-plugin":"^6.0.0"},"peerDependencies":{"@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0"},"_npmOperationalInternal":{"tmp":"tmp/kafka-module_1.0.1_1745460728848_0.35580944069999676","host":"s3://npm-registry-packages-npm-production"}},"1.0.2":{"name":"@armaghanzahid/kafka-module","version":"1.0.2","keywords":["nestjs","kafka","microservices","typescript"],"author":{"name":"Armaghan Zahid"},"license":"ISC","_id":"@armaghanzahid/kafka-module@1.0.2","maintainers":[{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"}],"homepage":"https://github.com/your-org/kafka-module#readme","bugs":{"url":"https://github.com/your-org/kafka-module/issues"},"dist":{"shasum":"78396080cfb47096dda2d4d1e4126a64f4bf10b8","tarball":"https://registry.npmjs.org/@armaghanzahid/kafka-module/-/kafka-module-1.0.2.tgz","fileCount":22,"integrity":"sha512-IoATa0YpqQ/cGsyFJVdW1NyVrUhabF7cyLUPVMG1nbHeG2VfjNLWu1RnJcZNjXXhUnzSvAiBsZTUfmFZeZAabQ==","signatures":[{"sig":"MEYCIQD+dFMCX2XbyHGCTfpSCwYeYNQym8Krv8orLZn/eVfaHgIhAOEqwEe+y4EoIO0tfDRmLJSoo0i+9PM5auGZN4eIh/Z5","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":122709},"main":"dist/index.js","types":"dist/index.d.ts","engines":{"node":">=18.0.0"},"gitHead":"ff655eef276270bb92f58d7b5abeab3e6e6cfac8","scripts":{"lint":"eslint \"{src,apps,libs,test}/**/*.ts\" --fix","test":"jest","build":"nest build","format":"prettier --write \"src/**/*.ts\"","test:cov":"jest --coverage","test:e2e":"jest --config ./test/jest-e2e.json","test:debug":"node --inspect-brk -r tsconfig-paths/register -r ts-node/register node_modules/.bin/jest --runInBand","test:watch":"jest --watch","prepublishOnly":"npm run build"},"_npmUser":{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"},"repository":{"url":"git+https://github.com/your-org/kafka-module.git","type":"git"},"_npmVersion":"10.7.0","description":"A robust NestJS module for Kafka integration in microservices","directories":{},"_nodeVersion":"20.14.0","dependencies":{"rxjs":"^7.8.1","kafkajs":"^2.2.4","@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0","reflect-metadata":"^0.1.13"},"publishConfig":{"access":"public","registry":"https://registry.npmjs.org/"},"_hasShrinkwrap":false,"devDependencies":{"jest":"^29.5.0","eslint":"^8.42.0","ts-jest":"^29.1.0","ts-node":"^10.9.1","prettier":"^3.0.0","supertest":"^6.3.3","ts-loader":"^9.4.3","typescript":"^5.1.3","@nestjs/cli":"^10.0.0","@types/jest":"^29.5.2","@types/node":"^20.3.1","@types/express":"^4.17.17","tsconfig-paths":"^4.2.0","@nestjs/testing":"^10.0.0","@types/supertest":"^2.0.12","@nestjs/schematics":"^10.0.0","source-map-support":"^0.5.21","eslint-config-prettier":"^9.0.0","eslint-plugin-prettier":"^5.0.0","@typescript-eslint/parser":"^6.0.0","@typescript-eslint/eslint-plugin":"^6.0.0"},"peerDependencies":{"@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0"},"_npmOperationalInternal":{"tmp":"tmp/kafka-module_1.0.2_1745461474369_0.2890158422432816","host":"s3://npm-registry-packages-npm-production"}},"1.0.3":{"name":"@armaghanzahid/kafka-module","version":"1.0.3","keywords":["nestjs","kafka","microservices","typescript"],"author":{"name":"Armaghan Zahid"},"license":"ISC","_id":"@armaghanzahid/kafka-module@1.0.3","maintainers":[{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"}],"homepage":"https://github.com/your-org/kafka-module#readme","bugs":{"url":"https://github.com/your-org/kafka-module/issues"},"dist":{"shasum":"0bd87321b1e5dbe3036e4692e07902eae816b9e1","tarball":"https://registry.npmjs.org/@armaghanzahid/kafka-module/-/kafka-module-1.0.3.tgz","fileCount":22,"integrity":"sha512-8nvSueEsulpR0euQ1tb1Q0PWrgNlP1QdWHVBr8PwUn67BzQaiP7+k33f83N4L0WbLTD7562IAgQyght6vdwvtA==","signatures":[{"sig":"MEQCIH3S3aHCzZ4dUSGDZJh19AjvKDbnLCqSCFN8rbZnOlGvAiAu0d7AewDt8oWeDg07GXrhJ3Mbif0Fii8vyQJZl/6G6A==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":122709},"main":"dist/index.js","types":"dist/index.d.ts","engines":{"node":">=18.0.0"},"gitHead":"ff655eef276270bb92f58d7b5abeab3e6e6cfac8","scripts":{"lint":"eslint \"{src,apps,libs,test}/**/*.ts\" --fix","test":"jest","build":"nest build","format":"prettier --write \"src/**/*.ts\"","test:cov":"jest --coverage","test:e2e":"jest --config ./test/jest-e2e.json","test:debug":"node --inspect-brk -r tsconfig-paths/register -r ts-node/register node_modules/.bin/jest --runInBand","test:watch":"jest --watch","prepublishOnly":"npm run build"},"_npmUser":{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"},"repository":{"url":"git+https://github.com/your-org/kafka-module.git","type":"git"},"_npmVersion":"10.7.0","description":"A robust NestJS module for Kafka integration in microservices","directories":{},"_nodeVersion":"20.14.0","dependencies":{"rxjs":"^7.8.1","kafkajs":"^2.2.4","@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0","reflect-metadata":"^0.1.13"},"publishConfig":{"access":"public","registry":"https://registry.npmjs.org/"},"_hasShrinkwrap":false,"devDependencies":{"jest":"^29.5.0","eslint":"^8.42.0","ts-jest":"^29.1.0","ts-node":"^10.9.1","prettier":"^3.0.0","supertest":"^6.3.3","ts-loader":"^9.4.3","typescript":"^5.1.3","@nestjs/cli":"^10.0.0","@types/jest":"^29.5.2","@types/node":"^20.3.1","@types/express":"^4.17.17","tsconfig-paths":"^4.2.0","@nestjs/testing":"^10.0.0","@types/supertest":"^2.0.12","@nestjs/schematics":"^10.0.0","source-map-support":"^0.5.21","eslint-config-prettier":"^9.0.0","eslint-plugin-prettier":"^5.0.0","@typescript-eslint/parser":"^6.0.0","@typescript-eslint/eslint-plugin":"^6.0.0"},"peerDependencies":{"@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0"},"_npmOperationalInternal":{"tmp":"tmp/kafka-module_1.0.3_1745461577900_0.2880694943301685","host":"s3://npm-registry-packages-npm-production"}},"1.0.4":{"name":"@armaghanzahid/kafka-module","version":"1.0.4","author":{"name":"Armaghan Zahid"},"license":"ISC","_id":"@armaghanzahid/kafka-module@1.0.4","maintainers":[{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"}],"dist":{"shasum":"a1b4c0d110f4de326ed199518a06aafe87fcd83a","tarball":"https://registry.npmjs.org/@armaghanzahid/kafka-module/-/kafka-module-1.0.4.tgz","fileCount":28,"integrity":"sha512-nXgZWGBSLlskC9Yfk8sWmDNwOVaRCRSVqa3W2ehPnb8VuMGUw82zTOgxLrhBl2XWcQ9sc/4H358UWZ0BBn0+UA==","signatures":[{"sig":"MEUCIQC0oyHYS0hNA5drAJwUNuPlg3FfZHwCqDpCihJywp6CzwIgMN9ZIfBRZgzZ4wyoaXJTiQIxpNX/ddqdSKvLlc8Yto4=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":191050},"main":"dist/index.js","types":"dist/index.d.ts","engines":{"node":">=18.0.0"},"gitHead":"1fa35b613adb23fee285723289e0f69b3590f47a","scripts":{"lint":"eslint \"{src,apps,libs,test}/**/*.ts\" --fix","test":"jest","build":"nest build","format":"prettier --write \"src/**/*.ts\"","test:cov":"jest --coverage","test:e2e":"jest --config ./test/jest-e2e.json","test:debug":"node --inspect-brk -r tsconfig-paths/register -r ts-node/register node_modules/.bin/jest --runInBand","test:watch":"jest --watch","prepublishOnly":"npm run build"},"_npmUser":{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"},"_npmVersion":"10.7.0","description":"A robust NestJS module for Kafka integration in microservices","directories":{},"_nodeVersion":"20.14.0","dependencies":{"joi":"^17.13.3","zod":"^3.24.3","rxjs":"^7.8.1","kafkajs":"^2.2.4","@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0","@azure/core-amqp":"^4.3.6","reflect-metadata":"^0.1.13"},"publishConfig":{"access":"public","registry":"https://registry.npmjs.org/"},"_hasShrinkwrap":false,"devDependencies":{"jest":"^29.5.0","eslint":"^8.42.0","ts-jest":"^29.1.0","ts-node":"^10.9.1","prettier":"^3.0.0","supertest":"^6.3.3","ts-loader":"^9.4.3","typescript":"^5.1.3","@nestjs/cli":"^10.0.0","@types/jest":"^29.5.2","@types/node":"^20.3.1","@types/express":"^4.17.17","tsconfig-paths":"^4.2.0","@nestjs/testing":"^10.0.0","@types/supertest":"^2.0.12","@nestjs/schematics":"^10.0.0","source-map-support":"^0.5.21","eslint-config-prettier":"^9.0.0","eslint-plugin-prettier":"^5.0.0","@typescript-eslint/parser":"^6.0.0","@typescript-eslint/eslint-plugin":"^6.0.0"},"peerDependencies":{"@nestjs/core":"^10.0.0","@nestjs/common":"^10.0.0"},"_npmOperationalInternal":{"tmp":"tmp/kafka-module_1.0.4_1746168189315_0.29479329061504855","host":"s3://npm-registry-packages-npm-production"}},"1.0.5":{"name":"@armaghanzahid/kafka-module","version":"1.0.5","description":"A robust NestJS module for Kafka integration in microservices","main":"dist/index.js","types":"dist/index.d.ts","scripts":{"build":"nest build","format":"prettier --write \"src/**/*.ts\"","lint":"eslint \"{src,apps,libs,test}/**/*.ts\" --fix","test":"jest","test:watch":"jest --watch","test:cov":"jest --coverage","test:debug":"node --inspect-brk -r tsconfig-paths/register -r ts-node/register node_modules/.bin/jest --runInBand","test:e2e":"jest --config ./test/jest-e2e.json","prepublishOnly":"npm run build"},"author":{"name":"Armaghan Zahid"},"license":"ISC","dependencies":{"@azure/core-amqp":"^4.3.6","@nestjs/common":"^10.0.0","@nestjs/core":"^10.0.0","kafkajs":"^2.2.4","reflect-metadata":"^0.1.13","rxjs":"^7.8.1","zod":"^3.24.3"},"devDependencies":{"@nestjs/cli":"^10.0.0","@nestjs/schematics":"^10.0.0","@nestjs/testing":"^10.0.0","@types/express":"^4.17.17","@types/jest":"^29.5.2","@types/node":"^20.3.1","@types/supertest":"^2.0.12","@typescript-eslint/eslint-plugin":"^6.0.0","@typescript-eslint/parser":"^6.0.0","eslint":"^8.42.0","eslint-config-prettier":"^9.0.0","eslint-plugin-prettier":"^5.0.0","jest":"^29.5.0","prettier":"^3.0.0","source-map-support":"^0.5.21","supertest":"^6.3.3","ts-jest":"^29.1.0","ts-loader":"^9.4.3","ts-node":"^10.9.1","tsconfig-paths":"^4.2.0","typescript":"^5.1.3"},"peerDependencies":{"@nestjs/common":"^10.0.0","@nestjs/core":"^10.0.0"},"engines":{"node":">=18.0.0"},"publishConfig":{"registry":"https://registry.npmjs.org/","access":"public"},"_id":"@armaghanzahid/kafka-module@1.0.5","gitHead":"cba4df5e18ee22bc44d3f6f9ea71731f3413e9a7","_nodeVersion":"20.14.0","_npmVersion":"10.7.0","dist":{"integrity":"sha512-RezuozRvtNcYwjlVfJkLe/E0MoNdFFjscN1+/92qHoCyE7NjvNkN3rZ52j3/ixapRLkZ4cPjcveWX5HXaic84Q==","shasum":"c4eadc8bce9021616756beaa4cb3cf9dc13ec6e4","tarball":"https://registry.npmjs.org/@armaghanzahid/kafka-module/-/kafka-module-1.0.5.tgz","fileCount":28,"unpackedSize":196005,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEUCIGva/PiR5Q9QAzf/B6OWBF2Yx7OoGJpka43y915iCDIGAiEAwKkUS21Q6zFcDA94eUE7CQRZ1mmo+FNd2tzMLLZf67c="}]},"_npmUser":{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"},"directories":{},"maintainers":[{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/kafka-module_1.0.5_1749019168615_0.49230414962954483"},"_hasShrinkwrap":false}},"time":{"created":"2025-04-24T02:04:10.849Z","modified":"2025-06-04T06:39:28.987Z","1.0.0":"2025-04-24T02:04:11.150Z","1.0.1":"2025-04-24T02:12:09.103Z","1.0.2":"2025-04-24T02:24:34.570Z","1.0.3":"2025-04-24T02:26:18.134Z","1.0.4":"2025-05-02T06:43:09.506Z","1.0.5":"2025-06-04T06:39:28.822Z"},"author":{"name":"Armaghan Zahid"},"license":"ISC","description":"A robust NestJS module for Kafka integration in microservices","maintainers":[{"name":"armaghanzahid","email":"armaghan@heartbug.com.au"}],"readme":"# Kafka Module for NestJS\n\nA robust Kafka module for NestJS applications that provides type-safe message handling, automatic JSON parsing, and cloud provider support.\n\n## Features\n\n- Type-safe message handling with automatic JSON parsing\n- Support for AWS MSK, Azure Event Hubs, and local Kafka\n- Automatic SSL/SASL configuration based on cloud provider\n- Decorator-based topic subscription\n- Configurable error handling\n- Environment-based configuration\n- Zod schema validation for all configurations\n\n## Installation\n\n```bash\nnpm install @your-org/kafka-module\n# or\nyarn add @your-org/kafka-module\n```\n\n## Module Usage\n\n### Basic Module Registration\n\n```typescript\nimport { KafkaModule } from \"@your-org/kafka-module\";\n\n@Module({\n  imports: [\n    KafkaModule.register({\n      client: {\n        clientId: \"my-app\",\n        brokers: [\"localhost:9092\"],\n      },\n      consumer: {\n        groupId: \"my-group\",\n      },\n    }),\n  ],\n})\nexport class AppModule {}\n```\n\n### Async Module Registration\n\n```typescript\nimport { KafkaModule } from \"@your-org/kafka-module\";\nimport { ConfigService } from \"@nestjs/config\";\n\n@Module({\n  imports: [\n    KafkaModule.registerAsync({\n      imports: [ConfigModule],\n      useFactory: (configService: ConfigService) => ({\n        client: {\n          clientId: configService.get(\"KAFKA_CLIENT_ID\"),\n          brokers: configService.get(\"KAFKA_BROKERS\").split(\",\"),\n        },\n        consumer: {\n          groupId: configService.get(\"KAFKA_GROUP_ID\"),\n        },\n      }),\n      inject: [ConfigService],\n    }),\n  ],\n})\nexport class AppModule {}\n```\n\n## Message Handling\n\n### How the KafkaSubscribe Decorator Works\n\nThe `@KafkaSubscribe` decorator uses TypeScript's metadata reflection to register message handlers. Here's how it works:\n\n1. **Metadata Registration**\n\n   ```typescript\n   @KafkaSubscribe<MyMessage>('my-topic')\n   async handleMessage(payload: MyMessage, topic: string, partition: number) {\n     // Handler implementation\n   }\n   ```\n\n   - The decorator stores the topic and method name in the class metadata\n   - The generic type parameter `<MyMessage>` defines the expected payload type\n   - The method signature is validated at runtime\n\n2. **Handler Discovery**\n\n   - During module initialization, the service scans for classes with `@KafkaSubscribe` decorators\n   - Each decorated method is registered as a handler for its specified topic\n   - The service maintains a map of topics to their handlers\n\n3. **Message Processing**\n\n   - When a message arrives, the service:\n     1. Parses the message value as JSON\n     2. Validates the payload against the handler's type\n     3. Calls the handler with the parsed payload\n     4. Handles any errors that occur\n\n4. **Type Safety**\n\n   ```typescript\n   // The decorator ensures type safety through generics\n   @KafkaSubscribe<MyMessage>('my-topic')\n   async handleMessage(payload: MyMessage, topic: string, partition: number) {\n     // TypeScript knows the shape of 'payload'\n     console.log(payload.id); // OK\n     console.log(payload.unknown); // TypeScript error\n   }\n   ```\n\n5. **Error Handling**\n   - The decorator validates the method signature at runtime\n   - Throws errors if:\n     - The topic is empty or invalid\n     - The method is not async\n     - The decorator is used on a non-method property\n\n### Using the KafkaSubscribe Decorator\n\nThe `@KafkaSubscribe` decorator marks a method as a Kafka message handler. The decorated method must:\n\n- Be async\n- Have a generic type parameter for the message payload\n\n```typescript\nimport { KafkaSubscribe } from \"@your-org/kafka-module\";\nimport { Injectable } from \"@nestjs/common\";\n\n// Define your message type\ninterface MyMessage {\n  id: string;\n  data: string;\n}\n\n@Injectable()\nexport class MyService {\n  @KafkaSubscribe<MyMessage>(\"my-topic\")\n  async handleMessage(payload: MyMessage, topic: string, partition: number) {\n    console.log(`Received message:`, payload);\n  }\n}\n```\n\n### Real-World Example\n\nHere's a practical example showing type-safe message handling with DTOs:\n\n```typescript\nimport { KafkaSubscribe } from \"@your-org/kafka-module\";\nimport { Injectable } from \"@nestjs/common\";\n\n// Define your topics enum\nenum MedicalRecordTopics {\n  RECEPTION_UPLOAD = \"medical.records.reception.upload\",\n  PATIENT_UPDATE = \"medical.records.patient.update\",\n}\n\n// Define your DTO\ninterface UploadReceptionDto {\n  patientId: string;\n  recordType: string;\n  fileUrl: string;\n  uploadedBy: string;\n  timestamp: Date;\n}\n\n@Injectable()\nexport class MedicalRecordsService {\n  @KafkaSubscribe<UploadReceptionDto>(MedicalRecordTopics.RECEPTION_UPLOAD)\n  async handleReceptionUpload(\n    reception: UploadReceptionDto, // <-- automatically parsed payload\n    topic: string,\n    partition: number,\n  ) {\n    // No JSON.parse needed - service handles parsing\n    if (this.consumer) {\n      await this.consumer.onReceptionUpload(reception);\n    }\n  }\n}\n```\n\n### Error Handling in Handlers\n\n```typescript\n@KafkaSubscribe<MyMessage>('my-topic')\nasync handleMessage(payload: MyMessage, topic: string, partition: number) {\n  try {\n    // Process message\n    await this.processMessage(payload);\n  } catch (error) {\n    // Handle error\n    this.logger.error(`Failed to process message on ${topic}[${partition}]`, error);\n    throw error; // Will be caught by the module's error handler\n  }\n}\n```\n\n## KafkaService\n\nThe `KafkaService` provides methods for interacting with Kafka.\n\n### Publishing Messages\n\n```typescript\nimport { KafkaService } from \"@your-org/kafka-module\";\nimport { Injectable } from \"@nestjs/common\";\n\n@Injectable()\nexport class MyService {\n  constructor(private readonly kafkaService: KafkaService) {}\n\n  async publishMessage() {\n    // Publish with payload only\n    await this.kafkaService.publish(\"my-topic\", { id: \"1\", data: \"test\" });\n\n    // Publish with key and headers\n    await this.kafkaService.publish(\n      \"my-topic\",\n      { id: \"1\", data: \"test\" },\n      \"message-key\",\n      { \"custom-header\": \"value\" },\n    );\n  }\n}\n```\n\n### Accessing Kafka Clients\n\n```typescript\n@Injectable()\nexport class MyService {\n  constructor(private readonly kafkaService: KafkaService) {}\n\n  async getKafkaInfo() {\n    const kafka = this.kafkaService.getKafkaClient();\n    const producer = this.kafkaService.getProducer();\n    const consumer = this.kafkaService.getConsumer();\n  }\n}\n```\n\n## Configuration\n\n### Required Configuration\n\n```typescript\ninterface KafkaModuleOptions {\n  client: {\n    clientId: string; // Required: Unique identifier for the client\n    brokers: string[]; // Required: At least one broker address\n  };\n  consumer?: {\n    groupId: string; // Required if using consumer\n  };\n}\n```\n\n### Optional Configuration\n\n```typescript\ninterface KafkaModuleOptions {\n  client: {\n    ssl?: boolean; // Default: false\n    logLevel?: number; // Default: INFO (4)\n    sasl?: SASLOptions; // Optional: SASL configuration\n    connectionTimeout?: number; // Optional: Connection timeout in ms\n    requestTimeout?: number; // Optional: Request timeout in ms\n  };\n  consumer?: {\n    fromBeginning?: boolean; // Default: false\n  };\n  producer?: {\n    allowAutoTopicCreation?: boolean; // Default: true\n    transactionTimeout?: number; // Optional: Transaction timeout in ms\n    idempotent?: boolean; // Optional: Enable idempotent producer\n  };\n  cloudProvider?: \"aws\" | \"azure\" | \"local\"; // Default: 'local'\n  onMessageError?: (\n    error: Error,\n    context: { topic: string; partition: number; message: any },\n  ) => Promise<void>;\n}\n```\n\n### Cloud Provider Configuration\n\n#### AWS MSK\n\n```typescript\nKafkaModule.register({\n  client: {\n    clientId: \"my-app\",\n    brokers: [\"your-msk-endpoint:9092\"],\n  },\n  cloudProvider: \"aws\",\n});\n```\n\nRequired environment variables:\n\n- `AWS_ACCESS_KEY_ID`\n- `AWS_SECRET_ACCESS_KEY`\n- `AWS_REGION` (optional)\n- `AWS_SESSION_TOKEN` (optional)\n- `AWS_ROLE_ARN` (optional)\n\n#### Azure Event Hubs\n\n```typescript\nKafkaModule.register({\n  client: {\n    clientId: \"my-app\",\n    brokers: [\"your-eventhub-endpoint:9093\"],\n  },\n  cloudProvider: \"azure\",\n});\n```\n\nRequired environment variable:\n\n- `AZURE_EVENT_HUB_CONNECTION_STRING`\n\n## Environment Variables\n\nThe module can be configured using environment variables:\n\n```env\n# Required\nKAFKA_CLIENT_ID=my-app\nKAFKA_BROKERS=localhost:9092\n\n# Optional\nKAFKA_CLOUD_PROVIDER=local|aws|azure\nKAFKA_GROUP_ID=my-group\nKAFKA_CONNECTION_TIMEOUT_MS=3000\nKAFKA_REQUEST_TIMEOUT_MS=30000\n\n# AWS specific\nAWS_ACCESS_KEY_ID=your-access-key\nAWS_SECRET_ACCESS_KEY=your-secret-key\nAWS_REGION=your-region\nAWS_SESSION_TOKEN=your-session-token\nAWS_ROLE_ARN=your-role-arn\n\n# Azure specific\nAZURE_EVENT_HUB_CONNECTION_STRING=your-connection-string\n```\n\n## Best Practices\n\n1. **Type Safety**\n\n   - Always define message types for type-safe handling\n   - Use interfaces or types for message payloads\n   - Leverage TypeScript's type inference\n\n2. **Configuration**\n\n   - Use environment variables for configuration\n   - Validate configuration using the provided Zod schema\n   - Set appropriate timeouts for your use case\n\n3. **Error Handling**\n\n   - Implement proper error handling in message handlers\n   - Use the `onMessageError` callback for global error handling\n   - Log errors with appropriate context\n\n4. **Cloud Integration**\n\n   - Use appropriate cloud provider settings\n   - Configure SSL/SASL based on your environment\n   - Follow cloud provider best practices\n\n5. **Performance**\n   - Configure appropriate batch sizes\n   - Set reasonable timeouts\n   - Monitor consumer lag\n\n## Contributing\n\nContributions are welcome! Please read our contributing guidelines.\n\n## License\n\nMIT\n","readmeFilename":"README.md"}