In one line. anyq gives TypeScript and Go applications a single
IProducer/IConsumerinterface that runs unchanged across nine brokers, so switching queues becomes an import change instead of a rewrite.
Every team that uses a message queue eventually faces the same moment: the broker that made sense in year one stops making sense in year two. Maybe the cloud bill changed, or the company moved clouds, or a new service needs fan-out instead of point-to-point. Whatever the reason, the queue is suddenly wrong and the application code is married to it. Producers reference broker-specific clients; consumers call broker-specific acknowledge methods; error-handling wraps broker-specific exceptions. anyq is a dual TypeScript and Go library that decouples application logic from that choice by providing one consistent interface across nine brokers, with built-in Dead Letter Queue support and a circuit breaker already wired in.
Figure 1. Left: producer and consumer code tightly coupled to a specific broker SDK. Right: the same code running against anyq's shared interface; the broker is an import, not a rewrite.
Why this matters
- For teams facing cloud migration or multi-cloud deployments: queue lock-in is one of the quietest infrastructure taxes; anyq makes the broker a configuration decision rather than an architectural commitment, so a move from AWS SQS to Google Pub/Sub does not require rewriting message-handling logic.
- For TypeScript and Go shops building services that need to run portably: the two languages share the same interface contract, so a TS service and a Go service can consume from the same topic with code that reads identically at the abstraction layer.
- For teams that want resilience primitives without wiring them up themselves: built-in DLQ routing and a circuit breaker with exponential backoff are available across all nine adapters, not just the one broker that happened to ship those features natively.
If you only read this far: anyq lets you write producer and consumer code once, swap the broker by changing an import, and keep dead-letter and circuit-breaker behavior no matter which queue runs underneath.
Terms in 30 seconds
Domain engineers: skip ahead, this is orientation for everyone else.
- Message queue / message broker: a service that receives messages from a producer and delivers them to one or more consumers asynchronously. Examples: Redis Streams, Apache Kafka, AWS SQS.
- Producer: the application component that sends messages to a queue.
- Consumer: the application component that receives messages from a queue and processes them.
- Dead Letter Queue (DLQ): a secondary queue where messages that cannot be processed successfully are routed, so they can be inspected or retried without blocking normal flow.
- Circuit breaker: a pattern that stops attempts to call a failing service after a threshold of errors, preventing cascading failures, and automatically retries after a cooldown.
- AMQP (Advanced Message Queuing Protocol): an open wire-level protocol for message brokers; RabbitMQ implements it. (amqp.org)
- CloudEvents: a CNCF specification for a common envelope format for event data, to improve interoperability between event producers and consumers. (cloudevents.io)
The problem, for engineers who don't live in this space
Imagine you rent a car with a manual gearbox. You learn where the gear positions are, how much clutch to give it, how it responds. Now imagine next month the rental agency only has automatics. Everything you memorized is useless in a different way: not because driving changed, but because the interface did. Software message queues work exactly like that. You write code for SQS, and it knows about ReceiptHandle for acknowledgment, MaxNumberOfMessages for batch size, and VisibilityTimeout for lease extension. Move to RabbitMQ and none of those names exist; the acknowledgment is channel.ack(msg), the consumer model is pull-via-channel, and error handling routes through AMQP channel events. Move to Kafka and it changes again.
Here is the concrete version. Say you build a TypeScript service that processes order events off an AWS SQS queue. You use the AWS SDK directly: new SQSClient(), ReceiveMessageCommand, delete-on-ack. A year later the company moves to Google Cloud. Now you need @google-cloud/pubsub, a Subscription object, a message.ack() call with a completely different signature, and you need to re-test your retry and DLQ logic because Pub/Sub's dead-lettering is configured through a separate subscription, not your application code. Every consumer file changes. Every integration test changes.
Figure 2. Current state: the same application logic (process, acknowledge, handle errors) is expressed differently for each broker, making broker changes expensive rewrites rather than configuration changes.
Why this is genuinely hard. The naive fix is a thin wrapper, but brokers differ not just in API surface but in delivery semantics: SQS is at-least-once with visibility timeouts, Kafka is log-based with consumer group offsets, SNS is pub/sub fan-out with no consumer-side state. An abstraction that papers over those differences silently produces bugs; an abstraction that exposes them completely is no abstraction. The honest middle path is a shared interface for the operations every broker supports, with well-documented semantic differences called out where they exist.
What exists today (and the gap)
| Tool / Library | Single interface for produce + consume | Built-in DLQ support | Built-in circuit breaker | TypeScript support | Go support | Works across 7+ brokers |
|---|---|---|---|---|---|---|
Cloud SDKs directly (AWS SDK, @google-cloud/pubsub, etc.) | ❌ | ⚠️ broker-specific | ❌ | ✅ | ✅ | ❌ |
| BullMQ | ✅ | ✅ | ❌ | ✅ | ❌ | ❌ Redis only |
| node-rdkafka | ❌ | ❌ | ❌ | ✅ | ❌ | ❌ Kafka only |
| Spring Cloud Stream (JVM) | ✅ | ⚠️ varies | ❌ | ❌ | ❌ | ✅ |
| Watermill (Go) | ✅ | ❌ | ❌ | ❌ | ✅ | ⚠️ partial |
| anyq | ✅ | ✅ | ✅ | ✅ | ✅ | ✅ 9 adapters |
No existing library delivers one identical IProducer/IConsumer interface across seven or more brokers in both TypeScript and Go, with Dead Letter Queue routing and a circuit breaker included across all adapters. BullMQ and Watermill each solve the abstraction problem within one language and one or few brokers. Spring Cloud Stream covers multiple brokers in the JVM but not TypeScript or Go. The cloud SDKs themselves are broker-specific by design.
How it works
anyq provides a shared interface layer backed by broker-specific adapters, organized as a Bun monorepo in TypeScript and mirrored in Go. It has three layers:
- Core: the
IProducer<T>,IConsumer<T>, andIMessage<T>interfaces; abstract base classes; middleware pipeline (circuit breaker, retry with exponential backoff); JSON serialization utilities. - Adapters: one package per broker, each implementing the core interfaces with adapter-specific configuration. Nine are currently available.
- Testers: lightweight Hono web servers in
apps/testersfor exercising each adapter locally against a live broker instance.
Figure 3. anyq's layered architecture. Application code targets core interfaces (blue). Adapters (gray) handle broker-specific wiring. The middleware pipeline (DLQ, circuit breaker, retry) sits between them and applies identically regardless of which adapter is active.
Core interfaces
IProducer<T> exposes connect, disconnect, publish, publishBatch, and isConnected. IConsumer<T> exposes connect, disconnect, subscribe, unsubscribe, pause, resume, isConnected, and isPaused. Every message delivered to a subscriber implements IMessage<T> with fields id, body, metadata, plus ack(), nack(requeue?), and an optional extendDeadline(seconds). These three interfaces are the only surface application code needs to know. The purpose of keeping the interface this minimal is that every method in it has a clear analogue across all nine brokers; nothing in the core contract requires broker-specific knowledge to satisfy.
Adapters
The nine TypeScript adapters are: @anyq/memory (in-process, for dev and test), @anyq/redis-streams, @anyq/rabbitmq (AMQP), @anyq/sqs, @anyq/sns (fan-out), @anyq/google-pubsub, @anyq/kafka (Avro serialization supported in addition to JSON), @anyq/nats (JetStream), and @anyq/azure-servicebus. Each adapter takes adapter-specific connection config in its constructor and is otherwise invisible to calling code.
Middleware
The middleware pipeline wraps any consumer with circuit breaker logic and retry-with-exponential-backoff before the message reaches application code. DLQ routing is configured at the consumer level; messages that exhaust retries are published to the configured dead-letter destination using the same IProducer interface, so DLQ behavior is portable across brokers as well. The pipeline is composable: you can apply it to any adapter without modifying the adapter.
The decisions that actually mattered
One shared interface vs. a per-broker "ergonomic" wrapper
The fork. Two plausible designs: provide a thin convenience wrapper per broker (pleasant to use, exposes each broker's strengths), or enforce a single interface contract across all brokers (portable, but constrained to the intersection of capabilities).
Options. Ergonomic per-broker wrappers vs. a strict shared interface that all adapters must satisfy.
Chosen. Strict shared interface. Trade-off accepted. The interface deliberately does not expose broker-specific primitives: Kafka partition assignment, SNS subscription filters, NATS JetStream stream configuration. Those are reachable through adapter-specific config at construction time, but not through the shared interface. That constraint is the point: code written against IConsumer<T> cannot accidentally become Kafka-specific, which is what keeps it portable. The cost is that some broker-specific capabilities require dropping to the adapter layer directly.
Middleware at the library layer vs. leaving it to the application
The fork. Many queue libraries ship a minimal client and leave retry, DLQ, and circuit-breaking to application code or to a separate framework.
Options. Minimal client with no middleware vs. built-in middleware pipeline.
Chosen. Built-in middleware. Trade-off accepted. The library makes an opinionated choice about retry semantics (exponential backoff) and DLQ routing that not every application will want. Teams with strong existing opinions on retry policy may need to disable or replace the defaults. The benefit is that the resilience baseline is consistent across all adapters without any per-broker wiring: a consumer on SQS and a consumer on Kafka get the same circuit breaker behavior from the same configuration.
Why this is new, and why it's significant
What's novel. anyq is, to my knowledge, the first open-source library to expose one identical IProducer/IConsumer interface across nine message brokers in both TypeScript and Go simultaneously, with Dead Letter Queue and circuit breaker behavior that applies uniformly regardless of the adapter in use. The dual-language aspect is particularly significant: prior multi-broker abstractions are language-specific (BullMQ for JS/TS, Watermill for Go), so a polyglot team gets two different abstraction models and two different interface contracts.
Why it's significant to the field. Queue portability is a named, recurring cost in cloud-native development. The CloudEvents specification (cloudevents.io, CNCF Incubating) exists precisely because the ecosystem recognized that event schemas and delivery are too broker-specific, and it standardizes the envelope. anyq addresses the application-code layer of that same problem: not what the message looks like on the wire, but what the producer and consumer code looks like in the application. That is a distinct and complementary gap. The growth of multi-cloud and cloud-exit strategies makes broker portability more relevant over time, not less: organizations that want to avoid single-cloud lock-in at the infrastructure layer benefit from having the application layer already decoupled.
Standards and ecosystem alignment. anyq's interface design is consistent with:
- AMQP: the
@anyq/rabbitmqadapter communicates over AMQP, a wire-level open standard for message brokers, and theack()/nack()semantics inIMessage<T>mirror AMQP's explicit acknowledgment model. - CloudEvents (CNCF Incubating): anyq's
IMessage<T>envelope (id,body,metadata) is structurally compatible with CloudEvents' required attributes (id,data, metadata fields), making it straightforward to serialize messages as CloudEvents at the adapter layer. - Go standard library conventions: the Go implementation follows the interface-based composition pattern idiomatic to Go; the same
IProducer/IConsumercontract is expressed as Go interfaces rather than classes.
Honest scope. anyq normalizes the API surface, not the delivery semantics. Switching from
@anyq/sqsto@anyq/snschanges fan-out behavior; switching from@anyq/redis-streamsto@anyq/kafkachanges ordering and partition guarantees. The interface is identical, but the operational behavior of the broker underneath is not, and developers need to understand both. The Go implementation mirrors the same interface model; confirm exact Go package paths and API signatures before relying on them in production. anyq has not yet been independently adopted at scale, and no adapter has been proposed for inclusion in a standards body.
See it run
Install the core and one adapter (TypeScript, using Bun):
bun add @anyq/core @anyq/redis-streams
import { createRedisProducer, createRedisConsumer } from '@anyq/redis-streams';
const producer = createRedisProducer({ url: 'redis://localhost:6379' });
await producer.connect();
await producer.publish({ body: { orderId: 'ord-001', amount: 42 } });
const consumer = createRedisConsumer({ url: 'redis://localhost:6379' });
await consumer.connect();
await consumer.subscribe(async (message) => {
console.log(message.body); // ← { orderId: 'ord-001', amount: 42 }
await message.ack(); // ← same call regardless of broker
});
To switch to Google Pub/Sub, change the import and the connection config:
// Before: import { createRedisProducer, createRedisConsumer } from '@anyq/redis-streams';
import { createPubSubProducer, createPubSubConsumer } from '@anyq/google-pubsub';
// message-handling code below is unchanged
The subscribe handler, the message.ack() call, the publish invocation, none of those lines change. Only the import and the connection object change.
In Go, the same flow runs against the same interface contract, with adapters selected the same way; see the repository for the current Go module path and a runnable producer/consumer example.
Go deeper. The
packages/coredirectory contains the interface definitions, abstract base classes, and middleware pipeline. Thepackages/redis-streamsadapter is a good concrete example of how a broker maps onto the shared interfaces. (Pin to a specific commit SHA when linking to source.)
What's next
The immediate roadmap includes additional adapter coverage and formalizing the Go module structure. anyq sits at the application-code layer of a broader portability problem: it decouples what you write from which queue runs it, the same way an ORM decouples data access from which database runs it. As multi-cloud deployments become the default rather than the exception, that kind of broker-agnostic application layer becomes less of a convenience and more of a standard expectation.
Appendix / references
- Repository: github.com/sns45/anyq
- Standards referenced: AMQP (amqp.org) · CloudEvents CNCF Incubating (cloudevents.io)
- Discussion / feedback: github.com/sns45/anyq/issues