{"_id":"@asenajs/asena-kafka","_rev":"4-c2a2f397065357779fb283a6444762ee","name":"@asenajs/asena-kafka","dist-tags":{"latest":"4.0.0"},"versions":{"1.0.0":{"name":"@asenajs/asena-kafka","version":"1.0.0","keywords":["kafka","kafkajs","asena","asenajs","bun","microservice","transport","ioc","dependency-injection","messaging","event-driven"],"author":{"name":"LibirSoft"},"license":"MIT","_id":"@asenajs/asena-kafka@1.0.0","maintainers":[{"name":"libir","email":"libirsoft@gmail.com"}],"homepage":"https://asena.sh","bugs":{"url":"https://github.com/AsenaJs/asena-kafka/issues"},"dist":{"shasum":"94c3d17e4dfbb5bfe868fe92f6ca90244deb78d5","tarball":"https://registry.npmjs.org/@asenajs/asena-kafka/-/asena-kafka-1.0.0.tgz","fileCount":42,"integrity":"sha512-r+TZkSMZu4MaCVABijfNnne3x7dN7BQk2Rjo+WUui9YRJQbVh3ISWwQbrOPMXkibpCAgeTaEg7l3JpKm/hNhmQ==","signatures":[{"sig":"MEUCICiiwR98/eRfQEduZxeSEm7P3BGmm+u/JT2OelN+fleSAiEAvhW+ysS8PeogM3VnoZTm50gIgtwEzzlS/LyfKW4SRD8=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":161991},"main":"dist/index.js","types":"dist/index.d.ts","module":"index.ts","gitHead":"7923815dbabba49fed94a5ac56095c7779d7423b","scripts":{"test":"bun test","build":"bun run clean && tsc","clean":"rm -rf dist","release":"bun run build && changeset publish","version":"changeset version","changeset":"changeset","test:watch":"bun test --watch","test:coverage":"bun test --coverage","prepublishOnly":"bun run build"},"_npmUser":{"name":"libir","email":"libirsoft@gmail.com"},"repository":{"url":"git+https://github.com/AsenaJs/asena-kafka.git","type":"git"},"_npmVersion":"12.0.1","description":"Kafka integration for AsenaJS - service client and microservice transport","directories":{},"_nodeVersion":"26.5.0","_hasShrinkwrap":false,"devDependencies":{"eslint":"^8.57.1","kafkajs":"^2.2.4","prettier":"^3.5.3","@types/bun":"latest","@asenajs/asena":"^0.8.0","@changesets/cli":"^2.29.3","eslint-plugin-n":"^16.6.2","reflect-metadata":"^0.2.2","eslint-config-alloy":"^5.1.2","eslint-plugin-alloy":"^1.2.1","eslint-plugin-import":"^2.31.0","eslint-plugin-promise":"^6.6.0","eslint-config-prettier":"^9.1.0","eslint-plugin-prettier":"^5.4.0","@typescript-eslint/eslint-plugin":"^6.21.0"},"peerDependencies":{"kafkajs":"^2.2.4","typescript":"^5.8.2","@asenajs/asena":"^0.8.0","reflect-metadata":"^0.2.2"},"_npmOperationalInternal":{"tmp":"tmp/asena-kafka_1.0.0_1785007215764_0.7101613863339098","host":"s3://npm-registry-packages-npm-production"}},"2.0.0":{"name":"@asenajs/asena-kafka","version":"2.0.0","keywords":["kafka","kafkajs","asena","asenajs","bun","microservice","transport","ioc","dependency-injection","messaging","event-driven"],"author":{"name":"LibirSoft"},"license":"MIT","_id":"@asenajs/asena-kafka@2.0.0","maintainers":[{"name":"libir","email":"libirsoft@gmail.com"}],"homepage":"https://asena.sh","bugs":{"url":"https://github.com/AsenaJs/asena-kafka/issues"},"dist":{"shasum":"29f81e32ec3846d57159aa1bacc6e953bb10ee8a","tarball":"https://registry.npmjs.org/@asenajs/asena-kafka/-/asena-kafka-2.0.0.tgz","fileCount":42,"integrity":"sha512-ygpzF5/kbfNSPdIOnf0wQxfv83rEf6LnWT+1seMt4NWxVxTHnh9nu8WtP6nCD2LXLwgemVAYI0631PRnwumHpA==","signatures":[{"sig":"MEUCIG6wwaLcZv2ze7r0F2+bqdpFK4FxtU6D14YoKekRqZK7AiEAtKotq7Zk9flmtARjltVYM3nw17pZpbglRtORu7Vv0z4=","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":183429},"main":"dist/index.js","types":"dist/index.d.ts","module":"index.ts","engines":{"bun":">=1.3.12"},"gitHead":"efc19835edc14996faeb2d3a83e150a994d5cd61","scripts":{"lint":"eslint . --ext .ts","test":"bun test","build":"bun run clean && tsc","check":"bun run lint && bun run format:check","clean":"rm -rf dist","format":"prettier --write .","release":"bun run build && changeset publish","version":"changeset version","lint:fix":"eslint . --ext .ts --fix","changeset":"changeset","check:fix":"bun run lint:fix && bun run format","typecheck":"tsc -p tsconfig.typecheck.json","test:watch":"bun test --watch","format:check":"prettier --check .","test:coverage":"bun test --coverage","prepublishOnly":"bun run build"},"_npmUser":{"name":"libir","email":"libirsoft@gmail.com"},"repository":{"url":"git+https://github.com/AsenaJs/asena-kafka.git","type":"git"},"_npmVersion":"12.0.1","description":"Kafka integration for AsenaJS - service client and microservice transport","directories":{},"_nodeVersion":"26.5.0","_hasShrinkwrap":false,"devDependencies":{"eslint":"^8.57.1","kafkajs":"^2.2.4","prettier":"^3.5.3","@types/bun":"latest","@asenajs/asena":"^0.9.0","@changesets/cli":"^2.29.3","eslint-plugin-n":"^16.6.2","reflect-metadata":"^0.2.2","eslint-config-alloy":"^5.1.2","eslint-plugin-alloy":"^1.2.1","eslint-plugin-import":"^2.31.0","eslint-plugin-promise":"^6.6.0","eslint-config-prettier":"^9.1.0","eslint-plugin-prettier":"^5.4.0","@typescript-eslint/eslint-plugin":"^6.21.0"},"peerDependencies":{"kafkajs":"^2.2.4","typescript":"^5.8.2","@asenajs/asena":"^0.9.0","reflect-metadata":"^0.2.2"},"_npmOperationalInternal":{"tmp":"tmp/asena-kafka_2.0.0_1785197146095_0.8204699335062946","host":"s3://npm-registry-packages-npm-production"}},"3.0.0":{"name":"@asenajs/asena-kafka","version":"3.0.0","keywords":["kafka","kafkajs","asena","asenajs","bun","microservice","transport","ioc","dependency-injection","messaging","event-driven"],"author":{"name":"LibirSoft"},"license":"MIT","_id":"@asenajs/asena-kafka@3.0.0","maintainers":[{"name":"libir","email":"libirsoft@gmail.com"}],"homepage":"https://asena.sh","bugs":{"url":"https://github.com/AsenaJs/asena-kafka/issues"},"dist":{"shasum":"c09d90fe56add8230b95196b1627b0d8ea5f4084","tarball":"https://registry.npmjs.org/@asenajs/asena-kafka/-/asena-kafka-3.0.0.tgz","fileCount":42,"integrity":"sha512-bQIaA411Lo3cGWiuaMq6l2ZEZCqI+mbhegcFxopz3SAPfMlgMrXwceaeo2vsRs7dAAVE96aHkbImoZ92sI+WOw==","signatures":[{"sig":"MEQCIE9Ar/TSQxHugCVGbe2J2HzPWORjxRoCbEm7/D6RKQDQAiAaNbQ6+a5CDby05AXkBXDN7AdoCfaFAUTJUZvUhxMc3g==","keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U"}],"unpackedSize":187774},"main":"dist/index.js","types":"dist/index.d.ts","module":"index.ts","engines":{"bun":">=1.3.12"},"gitHead":"a19498024d21d9fa7fa2fe32c92ff32f16bf2fb9","scripts":{"lint":"eslint . --ext .ts","test":"bun test","build":"bun run clean && tsc","check":"bun run lint && bun run format:check","clean":"rm -rf dist","format":"prettier --write .","release":"bun run build && changeset publish","version":"changeset version","lint:fix":"eslint . --ext .ts --fix","changeset":"changeset","check:fix":"bun run lint:fix && bun run format","typecheck":"tsc -p tsconfig.typecheck.json","test:watch":"bun test --watch","format:check":"prettier --check .","test:coverage":"bun test --coverage","prepublishOnly":"bun run build"},"_npmUser":{"name":"libir","email":"libirsoft@gmail.com"},"repository":{"url":"git+https://github.com/AsenaJs/asena-kafka.git","type":"git"},"_npmVersion":"12.0.1","description":"Kafka integration for AsenaJS - service client and microservice transport","directories":{},"_nodeVersion":"26.5.0","_hasShrinkwrap":false,"devDependencies":{"eslint":"^8.57.1","kafkajs":"^2.2.4","prettier":"^3.5.3","@types/bun":"latest","@asenajs/asena":"^0.10.0","@changesets/cli":"^2.29.3","eslint-plugin-n":"^16.6.2","reflect-metadata":"^0.2.2","eslint-config-alloy":"^5.1.2","eslint-plugin-alloy":"^1.2.1","eslint-plugin-import":"^2.31.0","eslint-plugin-promise":"^6.6.0","eslint-config-prettier":"^9.1.0","eslint-plugin-prettier":"^5.4.0","@typescript-eslint/eslint-plugin":"^6.21.0"},"peerDependencies":{"kafkajs":"^2.2.4","typescript":"^5.8.2","@asenajs/asena":"^0.10.0","reflect-metadata":"^0.2.2"},"_npmOperationalInternal":{"tmp":"tmp/asena-kafka_3.0.0_1785360722410_0.8116441891557289","host":"s3://npm-registry-packages-npm-production"}},"4.0.0":{"name":"@asenajs/asena-kafka","version":"4.0.0","author":{"name":"LibirSoft"},"description":"Kafka integration for AsenaJS - service client and microservice transport","main":"dist/index.js","module":"index.ts","types":"dist/index.d.ts","license":"MIT","engines":{"bun":">=1.4.0"},"keywords":["kafka","kafkajs","asena","asenajs","bun","microservice","transport","ioc","dependency-injection","messaging","event-driven"],"repository":{"type":"git","url":"git+https://github.com/AsenaJs/asena-kafka.git"},"bugs":{"url":"https://github.com/AsenaJs/asena-kafka/issues"},"homepage":"https://asena.sh","scripts":{"test":"bun test","test:watch":"bun test --watch","test:coverage":"bun test --coverage","build":"bun run clean && tsc","typecheck":"tsc -p tsconfig.typecheck.json","clean":"rm -rf dist","prepublishOnly":"bun run build","changeset":"changeset","version":"changeset version","release":"bun run build && changeset publish","lint":"eslint . --ext .ts","lint:fix":"eslint . --ext .ts --fix","format":"prettier --write .","format:check":"prettier --check .","check":"bun run lint && bun run format:check","check:fix":"bun run lint:fix && bun run format"},"devDependencies":{"@asenajs/asena":"^0.11.0","@changesets/cli":"^2.31.1","@types/bun":"latest","@typescript-eslint/eslint-plugin":"^6.21.0","eslint":"^8.57.1","eslint-config-alloy":"^5.1.2","eslint-config-prettier":"^9.1.2","eslint-plugin-alloy":"^1.2.1","eslint-plugin-import":"^2.32.0","eslint-plugin-n":"^16.6.2","eslint-plugin-prettier":"^5.5.6","eslint-plugin-promise":"^6.6.0","kafkajs":"^2.2.4","prettier":"^3.9.6","reflect-metadata":"^0.2.2"},"peerDependencies":{"@asenajs/asena":"^0.11.0","kafkajs":"^2.2.4","reflect-metadata":"^0.2.2","typescript":"^5.9.3"},"gitHead":"256e30bc0dffa2ab036dcfeac6012d9edbc0f6c0","_id":"@asenajs/asena-kafka@4.0.0","_nodeVersion":"26.7.0","_npmVersion":"12.0.2","dist":{"integrity":"sha512-7F5CGh4s2dO1nDSnMRNghsEPNTALmvDGY3+b7LHzDX6GVZLqXkN2LrjSFz+wlYuYF7DsbeKOJiQOYbBkmLbWNg==","shasum":"cb04df28888b0691b7c876727c4b06e635258560","tarball":"https://registry.npmjs.org/@asenajs/asena-kafka/-/asena-kafka-4.0.0.tgz","fileCount":42,"unpackedSize":187767,"signatures":[{"keyid":"SHA256:DhQ8wR5APBvFHLF/+Tc+AYvPOdTpcIDqOhxsBHRwC7U","sig":"MEQCIDwNYQEPOj0BUrhtcbYMfg4YkoIuYIRDOLW2kWxAJY3bAiBeo/wD26lyK1y2dMO85pDWRvsDxaXVwasDCI9CdEOkOw=="}]},"_npmUser":{"name":"libir","email":"libirsoft@gmail.com"},"directories":{},"maintainers":[{"name":"libir","email":"libirsoft@gmail.com"}],"_npmOperationalInternal":{"host":"s3://npm-registry-packages-npm-production","tmp":"tmp/asena-kafka_4.0.0_1787696185075_0.6264463493426595"},"_hasShrinkwrap":false}},"time":{"created":"2026-07-25T19:20:15.581Z","modified":"2026-08-25T22:16:25.466Z","1.0.0":"2026-07-25T19:20:15.885Z","2.0.0":"2026-07-28T00:05:46.255Z","3.0.0":"2026-07-29T21:32:02.579Z","4.0.0":"2026-08-25T22:16:25.239Z"},"bugs":{"url":"https://github.com/AsenaJs/asena-kafka/issues"},"author":{"name":"LibirSoft"},"license":"MIT","homepage":"https://asena.sh","keywords":["kafka","kafkajs","asena","asenajs","bun","microservice","transport","ioc","dependency-injection","messaging","event-driven"],"repository":{"type":"git","url":"git+https://github.com/AsenaJs/asena-kafka.git"},"description":"Kafka integration for AsenaJS - service client and microservice transport","maintainers":[{"name":"libir","email":"libirsoft@gmail.com"}],"readme":"<p width=\"%100\" align=\"center\">\n  <img src=\"https://avatars.githubusercontent.com/u/179836938?s=200&v=4\" width=\"150\" align=\"center\"/>\n</p>\n\n# @asenajs/asena-kafka\n\n[![Version](https://img.shields.io/badge/version-4.0.0-blue.svg)](https://github.com/AsenaJs/asena-kafka)\n[![License: MIT](https://img.shields.io/badge/License-MIT-green.svg)](https://opensource.org/licenses/MIT)\n[![Bun Version](https://img.shields.io/badge/Bun-1.4%2B-blueviolet)](https://bun.sh)\n\nKafka integration for AsenaJS — service client and microservice transport.\n\n`KafkaMicroserviceTransport` plugs Kafka into Asena's microservice layer (`@MessageController`, `@MessagePattern`, `@EventPattern`, RPC over `ulak`), and the `@Kafka` decorated service gives you raw produce/consume access with automatic IoC registration.\n\n## Features\n\n- **Microservice Transport** - Full `MicroserviceTransport` implementation: RPC (request/reply), event fan-out with wildcards, retries, DLQ\n- **Broker-Tracked Delivery Attempts** - `context.attempt` survives crashes and rebalances (persisted in offset-commit metadata, no side store)\n- **Decorator-Based Setup** - `@Kafka` decorator handles IoC registration and connection lifecycle\n- **Adapter Seam** - kafkajs today behind a `KafkaClientAdapter` interface (see [Client Roadmap](#client-roadmap))\n- **Deterministic Topic Management** - The transport creates its topics explicitly and waits for partition leaders; broker auto-create is never relied on\n- **External-Topic Interop** - Consume from and emit to envelope-less foreign topics (Quarkus/SmallRye, CDC, plain Kafka clients) with raw headers and traceparent continuity\n- **Zero Runtime Dependencies** - Only peer deps (asena, kafkajs, reflect-metadata)\n\n## Requirements\n\n- [Bun](https://bun.sh) v1.4 or higher\n- [@asenajs/asena](https://github.com/AsenaJs/Asena) v0.11.0 or higher\n- [kafkajs](https://kafka.js.org) v2.2.4 (peer dependency)\n- Apache Kafka **2.8 – 3.9**. Kafka 4.0 removed old protocol API versions (KIP-896) and kafkajs 2.2.4 has reported incompatibilities — pin your broker to 3.9.x. See [Client Roadmap](#client-roadmap).\n\n## Installation\n\n```bash\nbun add @asenajs/asena-kafka kafkajs\n```\n\n### Local development broker\n\nA single-node KRaft container is all you need:\n\n```bash\ndocker run -d --name asena-kafka --restart always -p 9092:9092 \\\n  -e KAFKA_NODE_ID=1 -e KAFKA_PROCESS_ROLES=broker,controller \\\n  -e KAFKA_LISTENERS=PLAINTEXT://:9092,CONTROLLER://:9093 \\\n  -e KAFKA_ADVERTISED_LISTENERS=PLAINTEXT://localhost:9092 \\\n  -e KAFKA_CONTROLLER_LISTENER_NAMES=CONTROLLER \\\n  -e KAFKA_CONTROLLER_QUORUM_VOTERS=1@localhost:9093 \\\n  -e KAFKA_LISTENER_SECURITY_PROTOCOL_MAP=CONTROLLER:PLAINTEXT,PLAINTEXT:PLAINTEXT \\\n  -e KAFKA_OFFSETS_TOPIC_REPLICATION_FACTOR=1 \\\n  -e KAFKA_TRANSACTION_STATE_LOG_REPLICATION_FACTOR=1 \\\n  -e KAFKA_TRANSACTION_STATE_LOG_MIN_ISR=1 \\\n  -e KAFKA_GROUP_INITIAL_REBALANCE_DELAY_MS=0 \\\n  -e KAFKA_GROUP_MIN_SESSION_TIMEOUT_MS=1000 \\\n  apache/kafka:3.9.1\n```\n\n`KAFKA_GROUP_MIN_SESSION_TIMEOUT_MS=1000` lets tests use low `sessionTimeout` values for fast crash detection; production brokers can keep the default.\n\n## Quick Start\n\n### Microservice Transport\n\n```typescript\nimport { Config, MessageController } from '@asenajs/asena/decorators';\nimport { MessagePattern, EventPattern } from '@asenajs/asena/microservice';\nimport { KafkaMicroserviceTransport } from '@asenajs/asena-kafka';\n\n@Config()\nexport class ServerConfig {\n  transport() {\n    return {\n      microservice: new KafkaMicroserviceTransport(\n        { brokers: ['localhost:9092'] },\n        { serviceName: 'order-service' },\n      ),\n    };\n  }\n}\n\n@MessageController('order')\nexport class OrderController {\n  @MessagePattern('create') // answers RPC 'order.create'\n  async create(data: CreateOrderDto, context: MessageContext) {\n    return { id: 42, ...data };\n  }\n\n  @EventPattern('*') // wildcard event subscription - handles 'order.*'\n  async onOrderEvent(data: any, context: MessageContext) {\n    // context.messageId is the dedup key, context.attempt > 1 marks redelivery\n  }\n\n  // Another service's vocabulary - opt out of the 'order' prefix\n  @EventPattern({ pattern: 'payment.*', prefix: false })\n  async onPayment(data: any, context: MessageContext) {}\n}\n```\n\nAny service can then call `ulak.messages('order').send('create', dto)` or `emit(...)` — see the [Microservices concepts doc](https://asena.sh/docs/concepts/microservices) for the full picture (headless mode, named transports, interceptors, OTel tracing).\n\n### Service Client\n\n```typescript\nimport { Kafka, AsenaKafkaService } from '@asenajs/asena-kafka';\n\n@Kafka({\n  config: { brokers: ['localhost:9092'], clientId: 'my-app' },\n  name: 'AppKafka',\n})\nexport class AppKafka extends AsenaKafkaService {\n\n  async publishAudit(entry: AuditEntry) {\n    await this.sendMessage('audit.log', [{ value: JSON.stringify(entry) }]);\n  }\n}\n```\n\n`AsenaKafkaService` exposes `sendMessage(topic, messages)`, `createProducer()`, `createConsumer(config)`, `createAdmin()`, `client`, `disconnect()` and `testConnection()`. The transport can borrow a decorated service too: `new KafkaMicroserviceTransport(appKafka, { serviceName })` — it still creates its own producer/consumers, your client is never touched.\n\nThe service connects on `server.start()` and disconnects itself on `server.stop()`. Objects from `createProducer()` / `createConsumer(config)` / `createAdmin()` stay yours: close them from an `@OnStop()` on the component that created them, which runs while the service is still up.\n\n## Delivery Model\n\n| Concern | Kafka shape |\n|---|---|\n| Events | Single shared topic `{prefix}.evt` (keyless round-robin over `eventPartitions`), one consumer group per service. Wildcards are matched locally; non-matching records are committed immediately. |\n| Requests (RPC) | One topic per exact pattern `{prefix}.req.{pattern}`, consumer group per responding service. RPC errors are FINAL — the caller gets the error, the offset is committed, no broker retry. |\n| Replies | One shared topic `{prefix}.reply`; every caller instance runs an ephemeral consumer group and filters by correlationId. |\n| DLQ | `{prefix}.dlq` with provenance headers (`origin_group`, `origin_stream`, `origin_offset`, `delivery_count`, `dlq_ts`). |\n\n### Attempt tracking (broker-persisted)\n\nKafka has no per-message delivery counter, so the transport persists one in **offset-commit metadata**: before dispatching the record at offset X it commits offset X with `{\"a\":attempt}`; on success it commits X+1 clean. A crash therefore redelivers exactly that record, and the successor (loaded on every group join) derives `attempt = a + 1` — `context.attempt` is trustworthy across processes with no extra infrastructure.\n\nThe cost is two synchronous commits per record per partition. Processing inside a partition is strictly sequential (that is what makes the marker sound); throughput scales by adding partitions (`maxInFlight` maps to kafkajs `partitionsConsumedConcurrently`).\n\n### Delivery guarantees\n\n- Events are **at-least-once**: a failed handler leaves the offset uncommitted; the partition is paused, seeked back and re-fetched after `retryBackoffMs`, up to `maxRetries`, then the record moves to the DLQ. Make handlers idempotent — dedup on `context.messageId` (one id per emit, identical for every group and every redelivery).\n- Requests older than the caller's own timeout (envelope header `to`) are **dropped without execution** — a restarting service never burns through a backlog of dead RPCs.\n- On its very first boot a consumer group starts at **latest**: messages sent before a service ever existed are invisible (deploy the consumer before the producer). The start position is pinned to a committed offset immediately, so later crashes can never skip records.\n\n## External Topics (interop)\n\nEverything above rides Asena's own envelope on Asena-owned topics. The `external` option is the escape hatch for **foreign systems** — a Quarkus/SmallRye service, a CDC pipeline, any plain Kafka producer/consumer — whose topics carry no envelope at all:\n\n```typescript\nconst transport = new KafkaMicroserviceTransport(\n  { brokers: ['localhost:9092'] },\n  {\n    serviceName: 'billing-service',\n    external: {\n      topics: [\n        'orders',                                    // string shorthand\n        { name: 'invoices', keyHeader: 'x-tenant' }, // outbound partition affinity\n      ],\n      fromBeginning: false, // start position on FIRST subscribe (default: latest)\n    },\n  },\n);\n```\n\nControllers stay completely transport-agnostic — external topics surface through the same decorators:\n\n```typescript\n@MessageController()\nexport class OrdersListener {\n  // Pattern = the external TOPIC NAME. Segment wildcards match it like any\n  // event: an 'upstream.*' handler pulls an 'upstream.orders' external topic\n  @EventPattern({ pattern: 'orders', prefix: false })\n  async onOrder(order: OrderPayload, context: MessageContext) {\n    // context.headers = ALL raw record headers - a Quarkus producer's\n    // traceparent / ce-* headers arrive verbatim (OTel picks traceparent up\n    // automatically, so the foreign trace continues through your handler)\n    // context.messageId = 'mid' header if present, else 'topic:partition:offset'\n  }\n}\n```\n\n**Inbound rules** (topics with a matching `@EventPattern`):\n\n- The dispatch pattern is **always the topic name** — a foreign record's incidental `p` header never steers routing.\n- The handler must resolve to exactly the topic name. The `@MessageController` prefix is joined onto `@EventPattern` too, so on a prefixed controller you **must** pass `prefix: false` — otherwise the handler registers as `billing.orders`, the topic is never consumed, and the service still boots green. `listen()` prints a hint naming the shadowed handler.\n- `context.headers` exposes **all raw record headers**; `context.messageId` honors a `mid` header, otherwise it falls back to the record position `topic:partition:offset`, which is stable across groups and redeliveries so messageId-dedup still works. Caveat: a replay *tool* that re-produces a DLQ'd record creates a new position and therefore a new id.\n- Retry, DLQ (with `origin_stream` = the foreign topic, value + headers preserved) and `context.attempt` work exactly as for own topics.\n\n**Outbound rules** (`emit('topicName')` where the name is configured external):\n\n- The value is plain JSON of your payload, headers are your `options.headers` **verbatim** — no `p`/`mid`/`h`/`ts` envelope keys. If OTel messaging instrumentation is active, its injected `traceparent` goes over as a plain header and the foreign consumer resumes the trace.\n- With `keyHeader` set, that header's value becomes the record **key** (partition affinity on the foreign topic); the header itself is still sent.\n- External topics are **event-only**: `send()` to an external name rejects immediately with `SEND_FAILED` (no request/reply without an envelope).\n\n**Boot behavior**: the transport **never creates external topics** — they are foreign property. Subscribed external topics are awaited (bounded) at `listen()` and missing ones fail the boot loudly; configured topics with no matching handler are outbound-only, logged and never an error. `fromBeginning: true` reads each external topic from the earliest retained record on the group's first subscribe (start-position pinning is skipped for those topics so the semantic actually holds).\n\n> Note for rolling deploys: changing the `external` list changes the group's subscription. Restart all replicas together for a deterministic assignment instead of mixing generations.\n\n## Configuration\n\n### KafkaConfig\n\nPassed through to the kafkajs `Kafka` constructor: `brokers` (required), `clientId`, `ssl`, `sasl`, `connectionTimeout`, `requestTimeout`, `retry`, `logLevel` — plus Asena's display `name`.\n\n### KafkaMicroserviceOptions\n\n| Option | Default | Description |\n|---|---|---|\n| `serviceName` | **required** | Consumer group identity, shared by all replicas (Kafka group id is `{topicPrefix}.{serviceName}`) |\n| `topicPrefix` | `'asena.ms'` | Topic/group namespace, `[a-zA-Z0-9._-]` only |\n| `requestTimeout` | `30000` | Default reply timeout for `send()` |\n| `maxRetries` | `3` | Event delivery attempts before DLQ (RPC never retried) |\n| `retryBackoffMs` | `5000` | Pause before the failed partition is resumed for the retry fetch |\n| `handlerTimeout` | `min(30000, sessionTimeout)` | Per-handler timeout; values above `sessionTimeout` throw at construction |\n| `maxInFlight` | `16` | Partitions processed concurrently |\n| `drainTimeout` | `10000` | Graceful drain window for `destroy()` |\n| `sessionTimeout` | `30000` | Consumer group session timeout (crash detection speed) |\n| `heartbeatInterval` | `3000` | Consumer heartbeat (keep ≤ 1/3 of sessionTimeout) |\n| `rebalanceTimeout` | `60000` | Max rebalance duration |\n| `maxWaitTimeInMs` | `1000` | Idle fetch long-poll (boot readiness / shutdown responsiveness) |\n| `eventPartitions` / `requestPartitions` / `replyPartitions` | `4` | Partition counts for transport-created topics |\n| `replicationFactor` | `-1` | Broker default |\n| `healthCheckIntervalMs` | `5000` | Interval of the active broker probe — one of the two inputs to `isConnected` (see Operational Notes) |\n| `external` | — | Foreign-topic interop: `{ topics: (string \\| { name, keyHeader? })[], fromBeginning? }` — see [External Topics](#external-topics-interop) |\n\n## Operational Notes\n\n- **Handler duration must stay below `sessionTimeout`** — kafkajs cannot heartbeat while a handler runs; a longer handler gets the member evicted and the record concurrently redelivered elsewhere.\n- **`isConnected` means \"can serve\", not \"a broker answers\"** — it requires both a passing metadata probe and a reply consumer that has rejoined and is fetching. After a broker outage the probe recovers in milliseconds while the ephemeral reply group can still be rejoining ~20 seconds later, and until it fetches nothing consumes replies, so every `send()` would time out. Expect an instance to stay 503 for the length of that rejoin, and keep the liveness probe more forgiving than the readiness probe.\n- **Rolling deploys**: kafkajs uses eager rebalancing — every membership change briefly pauses the whole group. Graceful shutdown (`destroy()`) leaves the group cleanly so the pause is short; SIGKILL costs a full `sessionTimeout`.\n- **Ordering**: with the default multi-partition event topic, cross-partition ordering is not preserved. Set `eventPartitions: 1` if you need strict publish order (at the cost of parallelism).\n- Under Bun you may see a cosmetic `TimeoutNegativeWarning` from kafkajs's request queue — Bun warns where Node silently clamps negative timers to 1ms; behavior is identical.\n\n## Client Roadmap\n\n**Why kafkajs?** It is the only full-featured Kafka client that actually runs on Bun today. It is also effectively unmaintained (2.2.4 is the final release), which is why every kafkajs call sits behind the `KafkaClientAdapter` interface.\n\n**Why not `@confluentinc/kafka-javascript`?** Confluent's official client ships a KafkaJS-compatibility mode and looks like the natural successor. It was tested directly on Bun 1.3.14 and does not work, for two independent reasons:\n\n1. **It cannot load on Bun.** The addon is built on NAN (the V8 C++ API), not N-API. Bun's ABI (`NODE_MODULE_VERSION` 137) matches a published prebuilt, and that prebuilt downloads and installs cleanly — so this is neither an install-script nor an ABI problem. Loading still dies inside `dlopen`:\n\n   ```\n   bun: symbol lookup error: .../confluent-kafka-javascript.node:\n        undefined symbol: _ZN2v816FunctionTemplate12SetClassNameENS_5LocalINS_6StringEEE\n   ```\n\n   That is `v8::FunctionTemplate::SetClassName`, which JavaScriptCore does not provide. Tracking: [bun#4290](https://github.com/oven-sh/bun/issues/4290), [confluent-kafka-javascript#264](https://github.com/confluentinc/confluent-kafka-javascript/issues/264).\n\n2. **Even on Node it would cost a feature.** Confluent's KafkaJS mode does not support sending metadata with offset commits, and the broker-tracked `context.attempt` protocol is built entirely on offset-commit metadata. `consumer.stop()` and the consumer instrumentation events used for ready-gating and crash detection are also unsupported.\n\n**What would change this:** [confluent-kafka-javascript#471](https://github.com/confluentinc/confluent-kafka-javascript/pull/471) migrates the addon from NAN to N-API. If it lands, blocker 1 disappears; blocker 2 would still need upstream work.\n\n## Contributing\n\nContributions are welcome! Please open an issue or a pull request on [GitHub](https://github.com/AsenaJs/asena-kafka).\n\n## License\n\nMIT — see [LICENSE](./LICENSE).\n\n## Support\n\n- Documentation: [asena.sh](https://asena.sh)\n- Issues: [GitHub Issues](https://github.com/AsenaJs/asena-kafka/issues)","readmeFilename":"README.md"}