ws

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

ws — send events over a WebSocket

A Sink that writes each essessey.Event as one JSON frame, and a Source that turns a read loop back into an essessey.Source.

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

Contents

Sending

This package never imports a WebSocket library. It asks for one method:

type Conn interface {
	WriteJSON(v any) error
}

gorilla's *websocket.Conn satisfies it as-is, so you pass your own connection:

func handler(w http.ResponseWriter, r *http.Request) {
	conn, err := upgrader.Upgrade(w, r, nil)
	if err != nil {
		return
	}
	defer conn.Close()

	sink := ws.NewSink(conn)
	pub := essessey.NewPublisher(r.Context(), sink)

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

	streamer := essessey.NewTextStreamer(pub, 0)
	for _, chunk := range []string{"Hello, ", "world!"} {
		if err := streamer.Write(r.Context(), chunk); err != nil {
			return
		}
	}

	_ = streamer.Close(r.Context())
	_ = pub.SendStreamEpilogue(essessey.StopReasonEndTurn, 0)
}

Receiving

Source is a queue with an essessey.Source face. Decode frames in your read loop and hand them to Deliver:

src := ws.NewSource()

go func() {
	defer src.Close()

	for {
		var ev essessey.Event
		if err := conn.ReadJSON(&ev); err != nil {
			return
		}

		src.Deliver(ev)
	}
}()

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

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

What the client actually gets

One JSON object per frame — the whole Event, envelope included:

{"id":"42","event":"content_block_delta","data":{"text":"hi"}}

This is the least work of any binding: nothing to parse out of a subject, no framing to scan. id is omitted entirely when empty rather than sent as "", which matters because in the SSE format an empty id RESETS a client's resume point — keeping the two bindings' semantics aligned means a client can treat a missing id the same way on both.

Concurrency — read this one

gorilla permits exactly one concurrent writer, and this package cannot enforce that for you.

Sink.Emit takes its own lock, so concurrent Emit calls are safe. What is not safe is your application writing to the same connection behind the sink's back — a ping, a control frame, a close message — while a stream is in flight. That is a data race gorilla will not protect you from and this package cannot see.

If anything else writes to that connection, put your own mutex around every writer including the sink, or funnel all writes through a single goroutine. The minimal Conn interface is what keeps this package free of a WebSocket dependency, and the cost of that choice is that the locking discipline for non-essessey writes stays yours.

Documentation

Overview

Package ws provides essessey.Sink and essessey.Source bindings for WebSocket connections.

WebSocket is message-oriented: every write is already a discrete frame, so this package needs no SSE-style framing — that framing exists only because an HTTP body is an undelimited byte stream. Sink sends the whole essessey.Event as one JSON message; Source hands events back exactly as they were delivered.

This package deliberately does NOT import a websocket library such as github.com/gorilla/websocket. Conn declares the one method (WriteJSON(v any) error) this package actually needs, and gorilla's *websocket.Conn already satisfies it structurally — the caller passes their own connection, and essessey adds zero transport dependency. This is a design choice, not a missing feature: pulling in a websocket library would tie essessey's module graph to one 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 Conn

type Conn interface {
	WriteJSON(v any) error
}

Conn is the minimal capability this package needs from a WebSocket connection: write a value as one JSON-encoded frame.

gorilla's *websocket.Conn (github.com/gorilla/websocket) satisfies this method as-is. Callers pass their own connection — see the package doc comment for why this package never imports a websocket library itself.

type Sink

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

Sink writes essessey events to a WebSocket connection.

WebSocket already delimits frames, so Emit sends the whole essessey.Event as one WriteJSON call — no SSE-style framing. The client receives {"event": ..., "data": ...} on the wire.

Unlike a NATS connection, gorilla's *websocket.Conn permits only one writer at a time, so Emit is mutex-guarded to stay safe for concurrent use.

func NewSink

func NewSink(c Conn) *Sink

NewSink returns a Sink that writes through c.

func (*Sink) Emit

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

Emit writes ev to the connection as a single JSON message.

type Source

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

Source turns a WebSocket read loop's per-message delivery into an essessey.Source. Wire Deliver as the callback the read loop invokes per inbound message; 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