nats

package
v0.7.5 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Sep 9, 2026 License: MIT Imports: 5 Imported by: 0

README

nats — publish events to NATS subjects

A Sink that publishes essessey.Event to NATS, and a Source that turns a subscription callback back into an essessey.Source.

NATS already delivers discrete messages, so nothing here frames anything — that is the sse package's job and only its job.

Contents

Publishing

This package never imports a NATS client. It asks for one method:

type Publisher interface {
	Publish(subject string, data []byte) error
}

*nats.Conn satisfies it as-is, so you pass your own connection and this package stays dependency-free:

conn, err := nats.Connect(nats.DefaultURL)
if err != nil {
	return err
}

sink := essnats.NewSink(conn, "chat.conv_1")
pub := essessey.NewPublisher(ctx, sink)

_ = pub.SendStreamPreamble("msg_1", "conv_1", "some-model")

Emit publishes ev.Data raw — no framing, no envelope. A *nats.Conn is safe for concurrent publishes, so the sink needs no lock of its own.

The subject scheme

The subject is subjectPrefix + . + the event type:

chat.conv_1.message_start
chat.conv_1.content_block_delta
chat.conv_1.message_stop

An empty prefix publishes to the bare event type. Putting the type in the subject is what lets a subscriber filter with a wildcard instead of decoding every payload to find out whether it cares:

conn.Subscribe("chat.stream_1.*", handler)            // one stream
conn.Subscribe("chat.*.content_block_delta", handler) // deltas everywhere

Subscribing

Source is a queue with an essessey.Source face. Wire Deliver into your subscription callback and pull events off with Next:

src := essnats.NewSource()

sub, err := conn.Subscribe("chat.conv_1.*", func(msg *nats.Msg) {
	src.Deliver(essessey.Event{
		Event: eventTypeFromSubject(msg.Subject),
		Data:  msg.Data,
	})
})

defer func() {
	_ = sub.Unsubscribe()
	src.Close()
}()

parsed := essessey.Reassemble(ctx, src)
fmt.Println(parsed.Text)

Close stops the source; Next drains whatever is already buffered and then returns essessey.ErrNoMoreEvents. Delivering after Close is a no-op that warns rather than panicking, because a subscription callback can still be in flight when you tear down.

The envelope is not identical to other bindings

This is the part that surprises people, so it is stated plainly rather than implied away.

binding how the client receives the envelope
ws the whole Event as one JSON object — id, event, data inline
sse id:, event: and data: as wire fields
nats payload raw; the event type is in the SUBJECT; the id is not carried at all

The Data payload is byte-identical everywhere. The envelope is not, and a NATS subscriber has the most reconstruction to do — you rebuild the event type from the subject you matched, as eventTypeFromSubject does above.

Event.ID has nowhere to go here. Core NATS messages carry the payload and the subject; there is no envelope field for it. If you need ids on this binding, carry them in a message header yourself, or use a binding that has room for them. Resume-from-Last-Event-ID is an SSE concept anyway — NATS has no equivalent reconnect semantics.

Ordering and delivery caveats

Worth knowing before streaming tokens over this:

  • The subject-per-event-type scheme fans one ordered stream across N subjects. Core NATS preserves order per publisher-subscriber pair on a connection, which usually holds — but it stops being a guarantee the moment there is a queue group, a cluster hop, or JetStream with multiple subjects. If strict order matters, publish to a single subject and put the type in the payload instead.
  • Core NATS is at-most-once with no redelivery. A slow consumer is dropped silently. For a token stream that is a corrupted render with no error anywhere.

Neither is a bug in this package — they are properties of the transport you chose — but they decide whether this binding suits your stream.

Documentation

Overview

Package nats provides essessey.Sink and essessey.Source bindings for NATS.

NATS is message-oriented: every publish is already a discrete message, so this package needs no SSE-style framing — that framing exists only because an HTTP body is an undelimited byte stream. Sink publishes essessey.Event.Data raw; Source hands events back exactly as they were delivered.

This package deliberately does NOT import github.com/nats-io/nats.go. Publisher declares the one method (Publish(subject string, data []byte) error) this package actually needs, and *nats.Conn already satisfies it structurally — the caller passes their own client, and essessey adds zero transport dependency. This is a design choice, not a missing feature: pulling in the real client would tie essessey's module graph to one NATS driver version for the sake of a single method.

Index

Constants

This section is empty.

Variables

This section is empty.

Functions

This section is empty.

Types

type Publisher

type Publisher interface {
	Publish(subject string, data []byte) error
}

Publisher is the minimal capability this package needs from a NATS client: publish raw bytes to a subject.

*nats.Conn (github.com/nats-io/nats.go) satisfies this method as-is. Callers pass their own connection — see the package doc comment for why this package never imports the nats client library itself.

type Sink

type Sink struct {
	// contains filtered or unexported fields
}

Sink publishes essessey events to NATS subjects.

NATS already delimits messages, so Emit publishes ev.Data RAW — no SSE-style framing. The subject is the configured prefix plus the event's type, so a subscriber can filter with a NATS wildcard subject (subjectPrefix.*) without inspecting payloads.

A real *nats.Conn is safe for concurrent Publish calls, so Emit needs no mutex of its own — it holds no other mutable state.

func NewSink

func NewSink(p Publisher, subjectPrefix string) *Sink

NewSink returns a Sink that publishes through p, with subjects prefixed by subjectPrefix. An empty subjectPrefix publishes to the bare event type.

func (*Sink) Emit

func (s *Sink) Emit(ctx context.Context, ev essessey.Event) error

Emit publishes ev.Data, unframed, to the subject derived from the configured prefix and ev.Event.

type Source

type Source struct {
	// contains filtered or unexported fields
}

Source turns a NATS subscription's push-callback delivery into an essessey.Source. Wire Deliver as the subscription callback; Next pulls events off in the order they arrived.

func NewSource

func NewSource() *Source

NewSource returns a ready-to-use Source.

func (*Source) Close

func (s *Source) Close()

Close stops the source. Next drains any already-buffered events first, then returns essessey.ErrNoMoreEvents. Close is idempotent.

func (*Source) Deliver

func (s *Source) Deliver(ev essessey.Event)

Deliver pushes ev onto the queue for Next to pick up. Deliver after Close is a no-op — the source has already stopped accepting events.

func (*Source) Next

func (s *Source) Next(ctx context.Context) (essessey.Event, error)

Next blocks until an event is delivered, the source is closed and drained, or ctx is done.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL