# queuekit (@mohamedhabibwork/queuekit) > Runtime-neutral TypeScript queues, messaging, pub/sub, and streams for Node.js >= 20, Bun, and Deno — Kafka, RabbitMQ, BullMQ, Redis, NATS, AWS SQS, Cloudflare Queues, Google Cloud Pub/Sub, and Azure Service Bus behind one typed contract, without losing provider-native types. The package installs with zero dependencies; provider SDKs are optional peers loaded dynamically only when that provider is created (the Cloudflare driver uses `fetch` and needs no SDK). A full in-memory provider covers development and tests with nothing installed. Dual ESM/CJS with full TypeScript declarations. ## Install ```bash npm install @mohamedhabibwork/queuekit # Add only the provider SDK you use (optional peers, loaded on creation): npm install kafkajs # or: bullmq, ioredis, amqplib, nats, redis, @aws-sdk/client-sqs, # @google-cloud/pubsub, @azure/service-bus (Cloudflare: no SDK) ``` **Entrypoints**: root (factory, manager, registry, config, errors, codec, middleware), `.../kafka`, `.../rabbitmq`, `.../bullmq`, `.../redis`, `.../nats`, `.../sqs`, `.../cloudflare`, `.../gcpubsub`, `.../azureservicebus`, `.../custom`, `.../testing` (in-memory queue + fake). A missing optional SDK throws `QueueConfigError` whose message contains the exact `npm install ` command — never a module-not-found crash. ## Quick start ```ts import { createQueue } from "@mohamedhabibwork/queuekit"; const queue = await createQueue({ type: "memory" }); // swap for bullmq/kafka in production await queue.publish("emails", { type: "welcome", payload: { email: "person@example.com" } }); const consumer = await queue.consume("emails", async ({ message }) => { await sendWelcomeEmail(message.payload.email); // return => ack; throw => retry policy }); ``` ## Providers | Provider | Config `type` | SDK (optional peer) | Notes | | --- | --- | --- | --- | | Memory | `memory` | none | dev/tests; latency injection, pending/acknowledged inspection | | BullMQ | `bullmq` | `bullmq` | Redis-backed jobs; flows, delayed jobs | | Redis | `redis` | `redis` or `ioredis` | pub/sub + streams modes; Valkey-compatible | | RabbitMQ | `rabbitmq` | `amqplib` | amqp topic/direct exchanges | | Kafka | `kafka` | `kafkajs` | consumer groups, partitions | | NATS | `nats` | `nats` | core + queue groups | | SQS | `sqs` | `@aws-sdk/client-sqs` | standard + FIFO queues | | Cloudflare Queues | `cloudflare` | none (fetch) | REST push/pull/ack; Workers binding publish; delay_seconds | | Google Cloud Pub/Sub | `gcpubsub` | `@google-cloud/pubsub` | subscriptions, ordering keys, server-side DLQ policy | | Azure Service Bus | `azureservicebus` | `@azure/service-bus` | peekLock settlement, scheduled messages, sessions | | Custom | `custom` / `defineQueueProvider` | — | register against the same contract | ## Core API ```ts queue.publish(destination, message: QueueMessage, options?); // QueueMessage: { type?, payload, headers?, correlationId?, causationId?, traceId?, metadata? } // PublishOptions: { retry?: { attempts, backoff: { type: 'fixed'|'exponential', delay, maxDelay? } }, // timeout?, signal?, metadata?, native? } queue.publishMany?(destination, messages, options?); queue.consume(destination, handler, options?); // handler context: { message, ack, signal, native? } // returns QueueConsumer: { status, pause(), resume(), close() } queue.health?(); // { ok, provider, latencyMs? } queue.native(); // the provider's native client/connection await queue.close(); ``` ## Publish with retries ```ts await queue.publish( "emails", { type: "welcome", payload: { email }, correlationId, traceId }, { retry: { attempts: 3, backoff: { type: "exponential", delay: 1_000, maxDelay: 30_000 } } }, ); ``` ## Consumers - Returning normally acknowledges; throwing triggers the provider retry policy; `ack.complete()` / `ack.retry()` / `ack.reject()` give manual control where the provider supports acknowledgement. - `consumer.pause()` / `resume()` / `close()` for graceful shutdown. - `signal` aborts in-flight handler work during close. ## Manager and typed queues ```ts import { createQueueManager, createTypedQueue } from "@mohamedhabibwork/queuekit"; const manager = createQueueManager({ providers: { emails: { type: "memory" }, events: { type: "memory" } }, default: "emails", }); const emails = await manager.provider("emails"); // typed per config; also default(), health(), close() // Typed payloads: destination + payload checked against the registry. const typed = createTypedQueue<{ welcome: { email: string } }>(emails); await typed.publish("welcome", { email: "person@example.com" }); await typed.consume("welcome", async ({ message }) => sendWelcomeEmail(message.payload.email)); ``` ## Testing The in-memory provider (`type: 'memory'`) and the fake from `.../testing` need no SDK and inspect internal state (`pending`, `acknowledged`) for assertions. Write providers against the contract once; the same tests pass against the real broker in CI. ## Runtime support - Node.js >= 20 (primary), Bun, Deno. Dual ESM/CJS with full declarations. - Root entry bundles no provider SDK; each driver loads its peer on creation. ## Links - npm: https://www.npmjs.com/package/@mohamedhabibwork/queuekit - Docs site: https://mohamedhabibwork.github.io/queuekit/ - README with full examples: ./README.md (also in the published tarball) - Guides: ./docs/frameworks.md (Express, Fastify, NestJS, Hono, Next.js, Elysia), ./docs/examples.md (end-to-end producer/worker apps), ./docs/ARCHITECTURE.md (layer map, optional-peer loading) ## Consumer patterns (provider-agnostic) ```ts import { withRetry, withTimeout, withIdempotency, scheduleRecurring, gracefulShutdown } from "@mohamedhabibwork/queuekit"; const consumer = await queue.consume("jobs", withRetry(withTimeout(handler, 10_000), { attempts: 5, backoff: { type: "exponential", delay: 1_000 } })); await queue.consume("orders", withIdempotency(handler, { store, ttl: 86_400_000 })); // store: { has, set } const job = scheduleRecurring(queue, "metrics", { payload: {} }, { every: 30_000 }); job.stop(); gracefulShutdown({ consumers: [consumer], queues: [queue], timeout: 30_000 }); ``` `scheduleAt` is converted to a delay for BullMQ, SQS, Cloudflare Queues (delay_seconds), and Azure Service Bus (scheduledEnqueueTimeUtc); Google Cloud Pub/Sub has no per-message delay. Full recipes: docs/use-cases.md.