Skip to main content

openframe-adapters-queue-nats

NATS JetStream queue adapter for the OpenFrame Microservice Development Suite.

Part of the openframe-adapters monorepo.


What it provides

Symbol Purpose
NatsSettings Pydantic-settings subclass — reads all config from env vars
NatsProducer[T] Generic async message producer — BaseProducer[T]
NatsConsumer[T] Generic async message consumer — BaseConsumer[T]
NatsPlugin BasePort (Identity + Lifecycle) — structured lifecycle via PluginRegistry/ApplicationBootstrap

Why JetStream, not core NATS

Core NATS pub/sub is fire-and-forget (at-most-once): there is no persistence, no publish acknowledgement, and no concept of redelivery — a message a subscriber isn't currently listening for is simply gone. That means there is nothing meaningful for BaseConsumer.ack()/nack() to do on top of it.

JetStream adds a persistent stream plus a durable, acknowledgeable consumer, giving ack()/nack() real at-least-once semantics: a handler failure (nack() → msg.nak()) triggers redelivery according to the durable consumer's ack_wait/max_deliver config — the same guarantee KafkaConsumer provides via manual offset commit. This adapter targets JetStream exclusively for that reason; there is no core-NATS mode.

Installation

# Via meta-package (recommended)
pip install "openframe-adapters[nats]"

# Or directly
pip install openframe-adapters-queue-nats

Quick start

from openframe.adapters.queue.nats import NatsSettings, NatsProducer, NatsConsumer

settings = NatsSettings(nats_servers="nats://localhost:4222")

# Produce
producer = NatsProducer(settings)
await producer.start()
await producer.publish({"event": "item.created", "id": "abc"})
await producer.publish_batch([{"event": "x"}, {"event": "y"}])
await producer.close()

# Consume
consumer = NatsConsumer(settings)

async def handle(event: dict) -> None:
    print(f"Received: {event}")

await consumer.subscribe(handle)   # runs until consumer.close() called

Configuration

Env var Default Description
NATS_SERVERS required nats://host:port,nats://host:port
NATS_SUBJECT "openframe" Default subject for producer and consumer
NATS_STREAM "OPENFRAME" JetStream stream name (auto-created, idempotent)
NATS_DURABLE_NAME "openframe-consumer" Durable JetStream consumer name
NATS_DELIVER_POLICY "all" "all" or "new"
NATS_ACK_WAIT_SECONDS 30.0 Seconds before an un-acked message is redelivered
NATS_MAX_DELIVER 5 Max redelivery attempts before JetStream gives up
NATS_MAX_RECONNECT_ATTEMPTS 60 NATS client reconnect attempts
NATS_USER / NATS_PASSWORD "" Basic auth, if the server requires it
NATS_TOKEN "" Token auth, if the server requires it

Typed domain objects

from openframe.adapters.queue.nats import NatsProducer, NatsConsumer, NatsSettings
from dataclasses import dataclass, asdict

@dataclass
class OrderEvent:
    order_id: str
    event_type: str

class OrderProducer(NatsProducer[OrderEvent]):
    def _serialise(self, message: OrderEvent) -> bytes:
        import json
        return json.dumps(asdict(message)).encode("utf-8")

class OrderConsumer(NatsConsumer[OrderEvent]):
    def _deserialise(self, raw: bytes) -> OrderEvent:
        import json
        return OrderEvent(**json.loads(raw.decode("utf-8")))

Plugin lifecycle (optional)

NatsPlugin is a BasePort — wire it up with ApplicationBootstrap.compose() (requires openframe-core>=3.3), the recommended zero-subclass entry point:

from openframe.core.runtime import ApplicationBootstrap
from openframe.core.ports import Capability
from openframe.adapters.queue.nats import NatsPlugin, NatsSettings

nats_plugin = NatsPlugin(NatsSettings())

async with ApplicationBootstrap.compose(nats_plugin) as app:
    plugin = app.get(Capability.QUEUE)
    producer = plugin.get_producer()
    await producer.publish({"event": "item.created"})

    consumer = plugin.make_consumer()
    await consumer.subscribe(handler)

NatsPlugin doubles as both the producer and consumer port for the QUEUE capability (get_producer() / make_consumer()), so a single instance is usually enough. If a service registers a separate producer-only and consumer-only NatsPlugin (e.g. different subjects/settings for each), pass both to compose(): ApplicationBootstrap.compose(producer_plugin, consumer_plugin). Reach for a subclassed ApplicationBootstrap with configure() only when you need per-port config=/init_timeout= or conditional registration order, and use app.registry as an escape hatch for anything neither tier covers.

Consumer acknowledgement semantics

Outcome Behaviour
Handler returns msg.ack() called → JetStream marks the message delivered
Handler raises msg.nak() called → JetStream redelivers per ack_wait/max_deliver
consumer.close() subscription unsubscribed, connection closed cleanly

Note: unlike KafkaConsumer's ack()/nack() (which commit/skip a global offset and ignore their message argument), NatsConsumer.ack()/nack() take the raw NATS message (not the deserialised payload) — JetStream acks are per-message, so the exact delivered message object is required to target the right ack. See consumer.py's module docstring for the full rationale.

Resilience (optional — requires openframe-core>=3.4)

openframe.core.resilience ships a CircuitBreakerProxy that composes around any BasePort, including this adapter's producer:

from openframe.core.resilience import CircuitBreakerProxy
from openframe.core.tracing import TracingProxy

producer = CircuitBreakerProxy(
    TracingProxy(plugin.get_producer(), prefix="nats.producer"),
    failure_threshold=5,
    reset_timeout=30.0,
)

Compose the circuit breaker around the traced producer, not the reverse — a short-circuited call should never produce a misleading adapter span for a call that never reached NATS. No adapter code changes are needed to support this; it wraps from the outside exactly like TracingProxy.

License

MIT — © Furious Meteors Engineering

Metadata

Release files for openframe-adapters-queue-nats 0.1.0

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for openframe-adapters-queue-nats 0.1.0
File Size Uploaded
openframe_adapters_queue_nats-0.1.0.tar.gz 22.9 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for openframe-adapters-queue-nats 0.1.0
File Interpreter ABI Platform
openframe_adapters_queue_nats-0.1.0-py3-none-any.whl Python 3 none any Details

Total release size: 42.0 kB

Release files / openframe_adapters_queue_nats-0.1.0.tar.gz

Download URL openframe_adapters_queue_nats-0.1.0.tar.gz
Size 22.9 kB
Tags Source
SHA-256 checksum
How to use checksums
11dae9facf1cbd0ef2a8c5c03b4440d8d3f47140cc10640d9d8e06f15ee9fe29
BLAKE2b-256 checksum
How to use checksums
6277f634380e9477769c6f0c6d4a317e2647c35825128a1b2505f7407a6e32cd
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.14

Release files / openframe_adapters_queue_nats-0.1.0-py3-none-any.whl

Download URL openframe_adapters_queue_nats-0.1.0-py3-none-any.whl
Size 19.1 kB
Tags Python 3
SHA-256 checksum
How to use checksums
8a716434580e9b6f9725012f5a87029d238083ad359126c91de69684eb1743f6
BLAKE2b-256 checksum
How to use checksums
74c5a016485021c59a47a4d6d94838b3fe0f35417e10687f98b14131b09f14ef
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/7.0.0 CPython/3.13.14

Release history Release notifications | RSS feed

This release

0.1.0 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page