essessey

package module
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: 7 Imported by: 0

README

essessey

Go Reference CI coverage version license imported by

Say the letters out loud. That's the name. S-S-E.

Getting a model's answer to whoever's waiting for it — token by token, block by block, in the order it actually fucking happened.

Here's the hill this package is prepared to die on: SSE is a format, not a transport. Everybody lists it next to WebSocket and NATS like they're three flavors of the same damn thing. They're not. event: / data: / blank line exists for exactly one reason — an HTTP response body is a pipe with no seams, so something has to mark where one event stops and the next starts. Hand those same events to NATS or a WebSocket and that framing is dead weight: those already deliver discrete messages. So here SSE is one binding that adds framing, the message-oriented ones don't, and all of them carry the same Event. A browser EventSource and a NATS subscriber decode the identical JSON PAYLOAD — the envelope around it is what differs per binding, and that difference is spelled out below rather than glossed over. That's the whole fucking point.

It doesn't talk to a model — that's elelem's job, and elelemstream is the seam between them. It doesn't own your HTTP handler, doesn't pick your broker, and doesn't drag half of npm's Go equivalent into your build just to prove it supports one: the NATS and WebSocket bindings are written against the smallest interface each actually needs, so you hand them the connection you already have and this package stays at zero transport dependencies. Zero. Check the go.mod yourself.

What it DOES own is the boring shit nobody wants to write twice — which content block index a tool result belongs to, when a thinking block has to close before the answer starts, and gluing a stream back together at the other end.

sink := essessey.NewInMemorySink() // or sse.NewWriterSink(w), nats.NewSink(conn, "turn"), ws.NewSink(conn)
pub := essessey.NewPublisher(ctx, sink)

if err := pub.SendStreamPreamble(msgID, streamID, model); err != nil {
	return err
}

streamer := essessey.NewTextStreamer(pub, 0)
for chunk := range delta {
	if err := streamer.Write(ctx, chunk); err != nil {
		return err
	}
}

if err := streamer.Close(ctx); err != nil {
	return err
}

return pub.SendStreamEpilogue(essessey.StopReasonEndTurn, outputTokens)

Contents

Quick start

go get github.com/psyb0t/essessey

A Publisher writes to any Sink. InMemorySink needs nothing at all to try — it just collects whatever got emitted:

ctx := context.Background()

sink := essessey.NewInMemorySink()
pub := essessey.NewPublisher(ctx, sink)

if err := pub.SendStreamPreamble("msg_1", "conv_1", "some-model"); err != nil {
	panic(err)
}

streamer := essessey.NewTextStreamer(pub, 0)

if err := streamer.Write(ctx, "Hello, "); err != nil {
	panic(err)
}

if err := streamer.Write(ctx, "world!"); err != nil {
	panic(err)
}

if err := streamer.Close(ctx); err != nil {
	panic(err)
}

if err := pub.SendStreamEpilogue(essessey.StopReasonEndTurn, 0); err != nil {
	panic(err)
}

fmt.Println(sink.Len(), "events emitted")

Swap InMemorySink for sse.NewWriterSink(w) (or sse.NewHTTPSink(w) behind a flushing http.ResponseWriter), nats.NewSink(conn, subjectPrefix), or ws.NewSink(conn) and every line above the sink construction stays exactly the fucking same — the Publisher, the streamer, and the event sequence have no idea which delivery is on the other end, and no reason to.

Reading a stream back

The other half. Reassemble drains any Source and hands back the finished turn — you do not walk events yourself unless you want to:

src := sse.NewSource(resp.Body) // or nats/ws Source, or SliceSource in a test

parsed := essessey.Reassemble(ctx, src)

fmt.Println(parsed.Text)       // every text delta, concatenated
fmt.Println(parsed.Thinking)   // every reasoning delta, concatenated
fmt.Println(parsed.ToolNames)  // tools the model called
fmt.Println(parsed.Timeline)   // reasoning, text, and tool activity in order

ParsedStream also carries Tools (each call matched to its result by content block index — the bookkeeping this package exists to do for you), Executions, StreamID, and Error if the stream carried one. Separate thinking blocks remain separate timeline entries, while Thinking holds their concatenated content.

If you do want the raw events, a Source is just an iterator:

for {
	ev, err := src.Next(ctx)
	if errors.Is(err, essessey.ErrNoMoreEvents) {
		break
	}
	if err != nil {
		return err
	}

	fmt.Println(ev.Event, string(ev.Data))
}

ErrNoMoreEvents ends the stream cleanly — it is the terminator, not a failure, so a caller can range over a Source without special-casing it.

Why one Event, many bindings

type Event struct {
	ID    string          `json:"id,omitempty"`
	Event EventType       `json:"event"`
	Data  json.RawMessage `json:"data"`
}

That's the whole wire model: an optional id, a name, and a JSON payload. Sink delivers it (Emit(ctx, Event) error); Source reads it back (Next(ctx) (Event, error), ending the stream with ErrNoMoreEvents). Neither interface knows what "framing" even means — that's a property of the binding underneath, not of the event.

ID is what makes a dropped connection recoverable. A browser EventSource remembers the last id it saw and sends it back as Last-Event-ID on reconnect, so a server can resume instead of restarting the stream — and a subscriber that tracks ids can tell it MISSED one rather than silently rendering a gap. It is omitempty deliberately: an empty id: field on the wire does not mean "no id", it RESETS the receiver's resume point, so emitting one for an event that simply has no id would throw away the position every earlier event established.

binding needs framing? why
SSE (io.Writer, http.ResponseWriter) yes an HTTP response body is an undelimited byte stream; something has to mark where one event ends
NATS no every publish is already a discrete message
WebSocket no every write is already a discrete frame

So the SSE binding owns a codec the other two never need, and it implements the format as specified rather than the subset one consumer happens to use: one data: field per line of the payload, id: and retry:, comment lines, CRLF/LF/lone-CR terminators, the optional space after a field's colon, a stripped byte order mark, and the rule that an event with no data field is discarded rather than delivered empty. FrameComment writes the keep-alive that stops an intermediary dropping an idle connection; FrameRetry tells the client how long to wait before reconnecting.

The per-line data: split is not a detail. The format has no escaping and no length prefix, so a newline inside a payload ends the FIELD and a blank line ends the EVENT — emitting a multi-line payload as one data: line puts different bytes on the wire than the caller passed, and for a payload containing a blank line it forges an extra event out of the remainder.

What travels as Data is the same json.RawMessage either way, but be precise about what "identical" means per binding: WebSocket writes the whole Event as one JSON object, so a client gets id/event/data inline. NATS publishes the payload RAW with the event type in the SUBJECT, so a subscriber reconstructs the envelope from the subject it matched and does not see the id at all. SSE carries all three as wire fields. The PAYLOAD is identical everywhere; the envelope takes a different route on each binding, and a NATS subscriber has the most work to do.

What each package does

Package Responsibility
Core (this package) Event, the Sink/Source interfaces, Publisher (one Send* method per protocol event, plus SendStreamPreamble/SendStreamEpilogue for the open/close pair), TextStreamer/LineStreamer for turning a chunk-at-a-time answer into correctly-indexed content blocks, and Reassemble, which drains a Source back into a ParsedStream with accumulated reasoning and text, tool calls matched to their results by content-block index, and an ordered timeline.
sse The SSE format itself: FrameLines renders the wire bytes, WriterSink/HTTPSink write framed events to an io.Writer or a flushing http.ResponseWriter, and Source scans them back off an io.Reader — a malformed frame gets warn-logged and skipped instead of nuking the whole stream.
nats A Sink that publishes Event.Data unframed to subjectPrefix.<eventType>, and a Source whose Deliver method you wire in as a subscription callback.
ws A Sink that writes the whole Event as one WriteJSON call, and a Source whose Deliver method you wire into a read loop.
elelemstream Bridges elelem's callbacks to this protocol — see below.
Retention (store.go, multisink.go) MultiSink fans one Emit out to several sinks, and EventStore retains recent events per stream so a reconnecting client can be resumed — with InMemoryEventStore as a bounded, per-stream default. See Resuming a dropped stream.
Test doubles (memory.go) InMemorySink collects events instead of delivering them (not test-only — it's also what you want when a turn has to be fully produced before any of it gets released), and SliceSource replays a fixed slice, so feeding one InMemorySink's Events() into a SliceSource round-trips a whole stream with no transport involved whatsoever.

Resuming a dropped stream

Event.ID gives a client a resume point; EventStore is what lets a server honour it. The protocol gives you the mechanism, not the retention — you can only resume to an event you still have.

Capture composes rather than wrapping. MultiSink fans one Emit to several sinks, and store.SinkFor(id) is the store's write half, so retention is just another destination:

store, err := essessey.NewInMemoryEventStore(256) // per stream, oldest evicted
// A capacity of zero or less returns ErrInvalidCapacity rather than a store
// that accepts every append and resumes nothing.
live := essessey.NewMultiSink(httpSink, store.SinkFor(chatID))

Replay needs nothing new. Ask what came after the client's last id, hand it to a SliceSource, and pump it through the ordinary sink — the same code path a live stream uses, so the two cannot drift and the client cannot tell them apart:

events, known, err := store.Since(ctx, chatID, lastEventID)
if !known {
	// Not retained: evicted, never seen, or from a previous process.
	// Restart the stream — do NOT guess.
}

That known flag is the whole design. When a resume point is missing, replaying from the start duplicates everything the client already rendered and replaying from now silently drops the events in between — the exact gap ids exist to prevent. Only the caller knows which is acceptable, so the store refuses to choose.

Four things worth knowing before wiring it up:

  • Replay through the wire sink alone, never the MultiSink — otherwise each reconnect re-appends what it is replaying and the store grows without bound.
  • streamID is a security boundary. Replay keyed only by event id would let a client presenting an id receive someone else's events. On reconnect the streamID arrives from the client, so authorize the resume exactly as you authorize opening the stream — Since only knows map keys, not owners.
  • Events without an id are still replayed, they just cannot be resumed TO. Dropping them would silently skip real content.
  • An id-less frame still moves a client's resume point. The format sets the last-event-id before it discards an event with no data, so a bare id: frame advances the client without delivering anything.

InMemoryEventStore is per-process and bounded — a reasonable default and a reference implementation. Anything that must survive a restart or span replicas wants its own EventStore against shared storage.

Zero transport dependencies

sse needs nothing beyond the standard library. nats and ws each declare the one single method they actually need from a client:

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

// ws.Conn
type Conn interface {
	WriteJSON(v any) error
}

*nats.Conn (github.com/nats-io/nats.go) and a gorilla *websocket.Conn already satisfy these as-is — you pass your own client in, and essessey never imports either SDK. Which keeps your module graph clear of the go mod vendor avalanche a real transport client drags along behind it, all for the sake of one method per binding. One method. That's not worth a dependency.

elelemstream

elelemstream is the one subpackage that imports elelem: it translates elelem's callback stream (text deltas, reasoning deltas, tool-call starts and tool results) into this package's block protocol, so an elelem-backed handler gets the same message_start → content blocks → message_stop sequence without you hand-rolling the translation. Everything else in this module — the core package, sse, nats, ws — stays free of an elelem import; a caller who isn't using elelem never pulls it in.

The block-index arithmetic it owns is the part worth reading before you go poking at it — which index a tool result lands on, why parallel calls wreck the naive implementation everyone writes first, and which invariants the tests pin down. That lives next to the code, in elelemstream/README.md. Read it before you "simplify" anything in there.

Layout

types.go, event.go             the wire types, EventType/Role/etc. constants, Event, Sink, Source
publisher.go                   Publisher and one Send* method per protocol event
streamer.go                    TextStreamer, LineStreamer
reassemble.go                  Source -> ParsedStream reconstruction
memory.go                      InMemorySink, SliceSource
multisink.go                   MultiSink — fan one Emit out to several Sinks
store.go                       EventStore, InMemoryEventStore — retention for resume
sse/                           the SSE format: codec, WriterSink, HTTPSink, Source
nats/                          NATS binding over a minimal Publisher interface
ws/                            WebSocket binding over a minimal Conn interface
elelemstream/                  elelem callbacks -> block protocol (imports elelem)

Development

Every target runs inside the Dockerfile.dev image, so a clean host needs Docker rather than a curated Go toolchain.

make dep           # tidy the module and re-vendor
make lint          # go fix, then golangci-lint at full strictness
make lint-fix      # the same, applying whatever it can fix itself
make test          # the suite, always with -race
make test-coverage # the suite plus the coverage floor CI enforces
make sec           # govulncheck + semgrep, merged into sec.sarif
make help          # the rest of it

License

MIT. See LICENSE.

See CHANGELOG.md for release notes.

Documentation

Overview

Package essessey streams an LLM turn to a client over whichever delivery the caller already has.

The whole wire model is one Event: a name plus a JSON payload. A Sink delivers it, a Source reads it back, and neither knows what "framing" means — that belongs to the binding underneath.

SSE is a format, not a transport

SSE gets listed alongside WebSocket and NATS as though the three were interchangeable. They are not. The event:/data:/blank-line framing exists because an HTTP response body is an undelimited byte stream, so something has to mark where one event ends. NATS and WebSocket already deliver discrete messages, so that framing would be dead weight there.

Hence the split: the sse subpackage owns a codec, the nats and ws subpackages own none, and all three move the same Event. A browser EventSource and a NATS subscriber decode identical JSON.

What lives here

Publisher turns protocol events into Sink emissions, with one Send method per event plus SendStreamPreamble/SendStreamEpilogue for the open/close pair. TextStreamer and LineStreamer accumulate a chunk-at-a-time answer into correctly-indexed content blocks. Reassemble drains a Source back into a ParsedStream with reasoning, text, tool calls matched to their results, and an ordered timeline.

Blocks open LAZILY, on first content, and the index advances only when a block that actually opened is closed. A round producing no text of a given kind must therefore emit nothing and burn no index — get that wrong and the client renders a blank card, or every later block shifts by one.

What lives elsewhere

Nothing here talks to a model; that is elelem's job, and the elelemstream subpackage is the seam. It is the only subpackage importing elelem, so this package and the transport bindings stay free of it.

Index

Constants

This section is empty.

Variables

View Source
var ErrInvalidCapacity = errors.New("capacity must be positive")

ErrInvalidCapacity rejects a non-positive retention size at construction.

A zero-capacity store would accept every Append and answer every Since with "not retained" — a resume buffer that silently never resumes. Failing at construction turns that into a startup error rather than a mystery later.

View Source
var ErrNoMoreEvents = errors.New("no more events")

ErrNoMoreEvents ends a Source's stream. Reassembly treats it as a clean end of input, not a failure, so a caller can range over a Source without special-casing the terminator.

Functions

This section is empty.

Types

type ContentBlock

type ContentBlock struct {
	Type ContentBlockType `json:"type"`
	Text string           `json:"text,omitempty"`
}

type ContentBlockDeltaData

type ContentBlockDeltaData struct {
	Type  EventType `json:"type"`
	Index int       `json:"index"`
	Delta TextDelta `json:"delta"`
}

type ContentBlockDeltaToolInputData

type ContentBlockDeltaToolInputData struct {
	Type  EventType      `json:"type"`
	Index int            `json:"index"`
	Delta InputJSONDelta `json:"delta"`
}

type ContentBlockDeltaToolResultData

type ContentBlockDeltaToolResultData struct {
	Type  EventType       `json:"type"`
	Index int             `json:"index"`
	Delta ToolResultDelta `json:"delta"`
}

type ContentBlockStartData

type ContentBlockStartData struct {
	Type         EventType    `json:"type"`
	Index        int          `json:"index"`
	ContentBlock ContentBlock `json:"content_block"`
}

type ContentBlockStartToolResultData

type ContentBlockStartToolResultData struct {
	Type         EventType       `json:"type"`
	Index        int             `json:"index"`
	ContentBlock ToolResultBlock `json:"content_block"`
}

type ContentBlockStartToolUseData

type ContentBlockStartToolUseData struct {
	Type         EventType    `json:"type"`
	Index        int          `json:"index"`
	ContentBlock ToolUseBlock `json:"content_block"`
}

type ContentBlockStopData

type ContentBlockStopData struct {
	Type  EventType `json:"type"`
	Index int       `json:"index"`
}

type ContentBlockType

type ContentBlockType = string

ContentBlockType is the type tag inside a content block or delta.

Deliberately a string ALIAS, not a defined type: a producer may emit block types this package has never heard of, and a consumer must be able to name them without patching essessey. The constants below are the built-in set, not the permitted set.

const (
	ContentBlockTypeText          ContentBlockType = "text"
	ContentBlockTypeTextDelta     ContentBlockType = "text_delta"
	ContentBlockTypeThinking      ContentBlockType = "thinking"
	ContentBlockTypeThinkingDelta ContentBlockType = "thinking_delta"
	ContentBlockTypeToolUse       ContentBlockType = "tool_use"
	ContentBlockTypeToolResult    ContentBlockType = "tool_result"
	ContentBlockTypeInputJSON     ContentBlockType = "input_json_delta"
	ContentBlockTypeJSONPartial   ContentBlockType = "json_partial"
)

type Event

type Event struct {
	// ID is the event's identifier. Optional, and omitted from the wire when
	// empty.
	//
	// This is what makes a dropped connection recoverable. An SSE client
	// remembers the last ID it saw and sends it back as the Last-Event-ID
	// header when it reconnects, so a server can resume from that point
	// instead of restarting the stream. It is also the only way a subscriber
	// can notice it MISSED an event rather than silently rendering a gap.
	//
	// Empty is meaningful, which is why this is omitempty rather than always
	// emitted: per the SSE specification an EMPTY id field RESETS the client's
	// last-event-ID to the empty string, so writing `id:` for an event that
	// simply has no ID would destroy the resume point of the events before it.
	ID string `json:"id,omitempty"`

	Event EventType       `json:"event"`
	Data  json.RawMessage `json:"data"`
}

Event is one protocol event on its way to a client: a name plus a JSON payload.

This is the ONLY thing a binding has to carry, and it is deliberately free of any delivery detail. Data is json.RawMessage rather than string because it has always held JSON — the type now says so, and every non-byte-stream binding is spared a []byte conversion per event.

type EventStore added in v0.6.0

type EventStore interface {
	// Append records ev as the most recent event of streamID.
	Append(ctx context.Context, streamID string, ev Event) error

	// Since returns the events recorded AFTER lastEventID, oldest first.
	//
	// known=false means the id is not in retention — evicted, never seen, or
	// from a previous process — and the caller must decide what to do rather
	// than receive a silently wrong slice. known=true with an empty slice means
	// the client is already up to date.
	Since(
		ctx context.Context,
		streamID string,
		lastEventID string,
	) ([]Event, bool, error)

	// Clear drops everything retained for streamID.
	Clear(ctx context.Context, streamID string) error
}

EventStore retains recent events per stream so a client that reconnects can resume from where it left off.

streamID is the same identifier Publisher.SendMessageStart puts on the wire as MessageMeta.StreamID — whatever the caller's domain calls one stream (a conversation, a job, a document build). essessey never mints it; it is an opaque key. Nothing forces the two to match, and keying retention at a finer grain (per turn, per connection) is a legitimate choice that trades a smaller buffer for losing replay of earlier turns. Using the id the client already learned from message_start is the common case.

It is also untrusted on the way back in: a reconnecting client supplies the streamID, and Since resolves it as a map key without any notion of who owns it. Authorize the resume exactly as you authorize opening the stream.

The read side deliberately reports whether the resume point is KNOWN instead of guessing. If an id has aged out, or never existed, neither available answer is safe: replaying from the start duplicates everything the client already rendered, and replaying from now silently drops the events in between — the exact gap event ids exist to prevent. Only the caller knows whether to restart the stream or fail the request, so Since hands that decision back.

type EventType

type EventType = string

EventType names one protocol event.

On an SSE byte stream this is the `event:` line; over NATS it is the subject suffix; over WebSocket it is the envelope's type field. The NAME is the same everywhere — only the delivery differs.

const (
	EventTypeMessageStart      EventType = "message_start"
	EventTypeContentBlockStart EventType = "content_block_start"
	EventTypePing              EventType = "ping"
	EventTypeContentBlockDelta EventType = "content_block_delta"
	EventTypeContentBlockStop  EventType = "content_block_stop"
	EventTypeMessageDelta      EventType = "message_delta"
	EventTypeMessageStop       EventType = "message_stop"
)

type InMemoryEventStore added in v0.6.0

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

InMemoryEventStore is a bounded, per-stream, in-process EventStore.

Per stream it keeps a ring buffer for order and O(1) eviction, plus an index from event id to sequence number for O(1) resume lookup. The index is what a ring alone cannot provide: event ids are opaque strings with no ordering, so without it every resume would be a linear scan.

Two behaviours are worth knowing before relying on it:

  • Events with NO id are still retained and still replayed. They cannot be resumed TO — nothing can name them — but they must come back if they fall after the resume point, or the resume would silently skip them.
  • Ids are assumed unique per stream. Reusing one moves the resume point to the later occurrence, so resuming from it skips everything in between.

Safe for concurrent use. Compose it for capture via SinkFor.

func NewInMemoryEventStore added in v0.6.0

func NewInMemoryEventStore(capacity int) (*InMemoryEventStore, error)

NewInMemoryEventStore returns a store retaining up to capacity events per stream, evicting the oldest first.

func (*InMemoryEventStore) Append added in v0.6.0

func (s *InMemoryEventStore) Append(
	_ context.Context,
	streamID string,
	ev Event,
) error

Append records ev as the most recent event of streamID, evicting the oldest once the stream is at capacity.

func (*InMemoryEventStore) Clear added in v0.6.0

func (s *InMemoryEventStore) Clear(_ context.Context, streamID string) error

Clear drops everything retained for streamID.

func (*InMemoryEventStore) Since added in v0.6.0

func (s *InMemoryEventStore) Since(
	_ context.Context,
	streamID string,
	lastEventID string,
) ([]Event, bool, error)

Since returns the events of streamID recorded after lastEventID.

func (*InMemoryEventStore) SinkFor added in v0.6.0

func (s *InMemoryEventStore) SinkFor(streamID string) Sink

SinkFor returns a Sink that appends everything it receives to streamID.

This is the capture half, and it exists so retention composes with MultiSink rather than needing a wrapper type:

live := essessey.NewMultiSink(httpSink, store.SinkFor(chatID))

Do NOT replay INTO that MultiSink — replay through the wire sink alone, or every reconnect re-appends what it is replaying and the store grows without bound.

It returns the Sink interface deliberately: this value exists to be composed alongside other Sinks, and a concrete unexported type would force callers to name something they cannot reference.

type InMemorySink

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

InMemorySink collects emitted events instead of delivering them.

Not test-only: it is also the buffer you want when a turn must be fully produced before any of it is released (a durable replay, a moderation pass). It lives in the core package rather than a test helper package for that reason, and because every binding's tests need it.

func NewInMemorySink

func NewInMemorySink() *InMemorySink

NewInMemorySink builds an empty InMemorySink.

func (*InMemorySink) Emit

func (s *InMemorySink) Emit(_ context.Context, ev Event) error

Emit appends the event.

func (*InMemorySink) Events

func (s *InMemorySink) Events() []Event

Events returns a COPY of what was emitted, so a caller ranging over the result cannot race a concurrent Emit or mutate the sink's own slice.

func (*InMemorySink) Len

func (s *InMemorySink) Len() int

Len reports how many events have been emitted.

type InputJSONDelta

type InputJSONDelta struct {
	Type        ContentBlockType `json:"type"`
	PartialJSON string           `json:"partial_json"`
}

type LineStreamer

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

LineStreamer is the newline-buffered sibling of TextStreamer: it splits input at newlines and passes each complete line through a LineTransform before emitting. Used for streaming line-oriented payloads (e.g. UI-component specs) where each line must be validated / rewritten before it reaches the client.

func NewLineStreamer

func NewLineStreamer(
	publisher *Publisher,
	startIndex int,
	transform LineTransform,
) *LineStreamer

NewLineStreamer builds a LineStreamer emitting from startIndex. transform may be nil, in which case lines pass through unchanged.

func (*LineStreamer) BlockIndex

func (s *LineStreamer) BlockIndex() int

BlockIndex returns the next available content-block index.

func (*LineStreamer) BlockStarted

func (s *LineStreamer) BlockStarted() bool

BlockStarted reports whether a block is currently open.

func (*LineStreamer) Close

func (s *LineStreamer) Close(ctx context.Context) error

Close flushes any trailing partial line and emits content_block_stop.

func (*LineStreamer) Text

func (s *LineStreamer) Text() string

Text returns the transformed output accumulated so far.

func (*LineStreamer) Write

func (s *LineStreamer) Write(ctx context.Context, chunk string) error

Write buffers a chunk, emitting each complete (newline-terminated) line through the transform; a trailing partial line is held for the next chunk.

type LineTransform

type LineTransform func(line string) string

LineTransform rewrites one buffered line before it is emitted. Returning the line unchanged is the identity transform.

type MessageDeltaData

type MessageDeltaData struct {
	Type  EventType        `json:"type"`
	Delta MessageDeltaInfo `json:"delta"`
	Usage UsageEnd         `json:"usage"`
}

type MessageDeltaInfo

type MessageDeltaInfo struct {
	StopReason   StopReason `json:"stop_reason"`
	StopSequence *string    `json:"stop_sequence"`
}

type MessageMeta

type MessageMeta struct {
	ID           string      `json:"id"`
	StreamID     string      `json:"stream_id"`
	Type         MessageType `json:"type"`
	Role         Role        `json:"role"`
	Content      []any       `json:"content"`
	Model        string      `json:"model"`
	StopReason   *StopReason `json:"stop_reason"`
	StopSequence *string     `json:"stop_sequence"`
	Usage        UsageStart  `json:"usage"`
}

type MessageStartData

type MessageStartData struct {
	Type    EventType   `json:"type"`
	Message MessageMeta `json:"message"`
}

type MessageStopData

type MessageStopData struct {
	Type EventType `json:"type"`
}

type MessageType

type MessageType = string

MessageType tags a message envelope.

const MessageTypeMessage MessageType = "message"

type MultiSink added in v0.6.0

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

MultiSink fans one Emit out to several Sinks.

It exists because every other Sink is terminal — there was no way to send the same event to two places. The motivating case is retention: a stream that writes to the client AND to an EventStore, so a reconnecting client can be resumed. It covers the other obvious ones too (a wire sink plus an audit log, or SSE plus NATS) without any of them needing a bespoke wrapper.

Emit is safe for concurrent use if the wrapped Sinks are; MultiSink holds no mutable state of its own.

func NewMultiSink added in v0.6.0

func NewMultiSink(sinks ...Sink) *MultiSink

NewMultiSink returns a Sink that forwards to every sink given, in order.

Zero sinks is allowed and discards everything. That is deliberate: it lets a caller build the list conditionally without special-casing empty.

func (*MultiSink) Emit added in v0.6.0

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

Emit forwards ev to every sink and returns the joined error of any that failed.

It does NOT stop at the first failure. A full store or a broken audit sink must not abort delivery to the client — the user's stream is the thing that matters and the rest are copies. Equally the failure is not swallowed: every error comes back joined, so a caller that cares can inspect it with errors.Is and one that does not still sees a non-nil error rather than silence.

The consequence worth knowing: a non-nil error does NOT mean the event went nowhere. It means at least one destination missed it.

type ParsedStream

type ParsedStream struct {
	StreamID   string
	Thinking   string
	Text       string
	ToolNames  []string
	Tools      []ToolCall
	Executions []ToolExecution
	Timeline   []TimelineItem
	Error      string
}

ParsedStream is the structured reconstruction of a full streamed turn.

func Reassemble

func Reassemble(ctx context.Context, src Source) ParsedStream

Reassemble drains src and reconstructs a full streamed turn: accumulated reasoning, text, tool_use blocks matched against their tool_result blocks by content-block index, and an ordered timeline of all three.

A malformed individual event is warn-logged and dropped rather than aborting reassembly — a single corrupted event degrades the result, not the caller. src.Next returning ErrNoMoreEvents ends the stream cleanly; the accumulated result is returned, never an error. Any other error from Next is recorded on the result's Error field and reassembly stops there, still returning whatever was accumulated so far.

type PingData

type PingData struct {
	Type EventType `json:"type"`
}

type Publisher

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

Publisher emits protocol events to a Sink. It is constructed per stream with the request context so the streamer helpers can call the Send* methods without threading ctx through every call.

func NewPublisher

func NewPublisher(ctx context.Context, sink Sink) *Publisher

NewPublisher builds a Publisher that writes to sink for the lifetime of ctx.

sink decides the delivery — framed bytes for a stream, a native message for NATS or WebSocket. Nothing below this line knows or cares which.

func (*Publisher) Publish

func (p *Publisher) Publish(eventType EventType, data any) error

Publish marshals data and emits it as an event of eventType. Every event gets a DEBUG log line (event kind + block index/type where applicable) so the wire-level flow of a turn is reconstructable from logs alone. For text/thinking/tool-result deltas ONLY the content LENGTH is logged, never the content itself — this is the single choke point all SendXxx helpers funnel through, so instrumenting it here covers every event without a duplicated log call in each helper, and keeps model output / tool-result payloads (which may carry sensitive data) out of logs at high frequency.

func (*Publisher) SendContentBlockDeltaText

func (p *Publisher) SendContentBlockDeltaText(index int, text string) error

SendContentBlockDeltaText appends a text delta to the block at index.

func (*Publisher) SendContentBlockDeltaThinking

func (p *Publisher) SendContentBlockDeltaThinking(
	index int,
	text string,
) error

SendContentBlockDeltaThinking appends a reasoning delta at index.

func (*Publisher) SendContentBlockStartText

func (p *Publisher) SendContentBlockStartText(index int) error

SendContentBlockStartText opens a text content block at index.

func (*Publisher) SendContentBlockStartThinking

func (p *Publisher) SendContentBlockStartThinking(index int) error

SendContentBlockStartThinking opens a thinking (reasoning) content block at index. Reasoning streams as its own block type so the client renders it apart from the final answer text.

func (*Publisher) SendContentBlockStop

func (p *Publisher) SendContentBlockStop(index int) error

SendContentBlockStop closes the content block at index.

func (*Publisher) SendMessageDelta

func (p *Publisher) SendMessageDelta(
	stopReason StopReason,
	outputTokens int,
) error

SendMessageDelta emits the trailing message_delta with stop reason + usage.

func (*Publisher) SendMessageStart

func (p *Publisher) SendMessageStart(
	msgID, streamID, model string,
) error

SendMessageStart emits message_start, carrying the stream id so the client learns a newly-created stream's id from the stream itself.

streamID is the caller's own identifier for whatever this stream represents — a conversation, a job, a document build. essessey never mints one and only echoes it back; it is also the natural key to retain the stream under in an EventStore, so the same value passed here goes to EventStore.SinkFor.

func (*Publisher) SendMessageStop

func (p *Publisher) SendMessageStop() error

SendMessageStop emits the terminal message_stop event.

func (*Publisher) SendPing

func (p *Publisher) SendPing() error

SendPing emits a ping keep-alive event.

func (*Publisher) SendStreamEpilogue

func (p *Publisher) SendStreamEpilogue(
	stopReason StopReason,
	outputTokens int,
) error

SendStreamEpilogue emits message_delta + message_stop to close a stream, carrying stopReason so the client knows WHY the stream ended (a clean answer, a tool request, or a truncation like a token cap) instead of always being told end_turn regardless of what actually happened.

func (*Publisher) SendStreamPreamble

func (p *Publisher) SendStreamPreamble(
	msgID, streamID, model string,
) error

SendStreamPreamble emits message_start + ping to open a stream.

func (*Publisher) SendToolInputDelta

func (p *Publisher) SendToolInputDelta(index int, inputJSON string) error

SendToolInputDelta streams the tool's partial input JSON at index.

func (*Publisher) SendToolResultBlock

func (p *Publisher) SendToolResultBlock(
	index int,
	toolUseID, resultText string,
	isError bool,
) error

SendToolResultBlock emits a full tool_result block (start + delta + stop).

func (*Publisher) SendToolResultDelta

func (p *Publisher) SendToolResultDelta(index int, text string) error

SendToolResultDelta streams the tool result payload at index.

func (*Publisher) SendToolResultStart

func (p *Publisher) SendToolResultStart(
	index int,
	toolUseID string,
	isError bool,
) error

SendToolResultStart opens a tool_result content block for toolUseID.

func (*Publisher) SendToolUseBlock

func (p *Publisher) SendToolUseBlock(
	index int,
	toolUseID, name, inputJSON string,
) error

SendToolUseBlock emits a full tool_use block (start + input delta + stop).

func (*Publisher) SendToolUseStart

func (p *Publisher) SendToolUseStart(
	index int,
	toolUseID, name string,
) error

SendToolUseStart opens a tool_use content block (name known, args pending).

type Role

type Role = string

Role is the author of a message.

Declared here rather than imported so a wire-protocol package stays self-describing. An engine's Role (elelem's, say) is a different concept that happens to share values today, and the two must be free to diverge without breaking each other.

const RoleAssistant Role = "assistant"

type Sink

type Sink interface {
	Emit(ctx context.Context, ev Event) error
}

Sink delivers events to a client.

SSE, NATS and WebSocket are NOT peers: SSE is a FORMAT whose framing exists because an HTTP body is an undelimited byte stream, while NATS and WebSocket are message-oriented and already have boundaries. So a byte-stream sink frames each event; a message-oriented sink carries Data as-is. Both satisfy this one method, and the payload a client deserializes is identical either way.

Implementations live in subpackages — essessey/sse, essessey/nats, essessey/ws — and callers add their own by satisfying Emit.

Emit must be safe for concurrent use: the tool loop can produce blocks from several goroutines.

type SliceSource

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

SliceSource replays a fixed slice of events, the Source counterpart to InMemorySink. Feeding an InMemorySink's Events into one round-trips a stream with no transport involved — which is how every binding proves it preserves the event sequence.

func NewSliceSource

func NewSliceSource(events []Event) *SliceSource

NewSliceSource builds a SliceSource over a copy of events, so a later mutation by the caller cannot rewrite a replay already in progress.

func (*SliceSource) Next

func (s *SliceSource) Next(_ context.Context) (Event, error)

Next returns the next event, or ErrNoMoreEvents once the slice is exhausted.

type Source

type Source interface {
	Next(ctx context.Context) (Event, error)
}

Source reads events back from a delivery, the mirror of Sink.

This exists because the read side cannot assume a byte stream. An SSE source scans framed text out of an io.Reader, but NATS and WebSocket hand over DISCRETE events with no stream to scan — there is no io.Reader to pass. Both shapes satisfy Next, so reassembly (see Reassemble) works against any delivery rather than against SSE alone.

Next returns ErrNoMoreEvents when the stream ends normally. Any other error is a real failure.

type StopReason

type StopReason = string

StopReason is why a message stopped.

const (
	StopReasonEndTurn StopReason = "end_turn"
	StopReasonToolUse StopReason = "tool_use"
	// StopReasonMaxTokens signals the model's response was cut off by a token
	// cap (upstream finish_reason "length") rather than finishing cleanly —
	// the client must NOT treat the turn as a complete answer.
	StopReasonMaxTokens StopReason = "max_tokens"
)

type TextDelta

type TextDelta struct {
	Type ContentBlockType `json:"type"`
	Text string           `json:"text"`
}

type TextStreamer

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

TextStreamer accumulates streamed text and, when a Publisher is set, emits it as content blocks — text blocks by default, or thinking blocks when built via NewThinkingStreamer. Publisher may be nil for non-streaming use — the text still accumulates and is readable via Text.

It also owns the block index, which is why block accounting stays consistent across a turn: Close advances the index exactly once per closed block, so the next streamer starts where this one stopped.

func NewTextStreamer

func NewTextStreamer(publisher *Publisher, startIndex int) *TextStreamer

NewTextStreamer builds a TextStreamer emitting text blocks from startIndex.

func NewThinkingStreamer

func NewThinkingStreamer(publisher *Publisher, startIndex int) *TextStreamer

NewThinkingStreamer builds a TextStreamer emitting reasoning (thinking) blocks from startIndex. Same accumulation + block-index accounting as the text streamer; only the emitted block/delta type differs.

func (*TextStreamer) BlockIndex

func (s *TextStreamer) BlockIndex() int

BlockIndex returns the next available content-block index.

func (*TextStreamer) BlockStarted

func (s *TextStreamer) BlockStarted() bool

BlockStarted reports whether a block is currently open.

func (*TextStreamer) Close

func (s *TextStreamer) Close(_ context.Context) error

Close emits content_block_stop when a block is open, then advances the block index. A no-op when no block is open — which is what keeps the index honest for a streamer that never received content.

func (*TextStreamer) Text

func (s *TextStreamer) Text() string

Text returns the accumulated text.

func (*TextStreamer) Write

func (s *TextStreamer) Write(_ context.Context, chunk string) error

Write appends a chunk, opening the block on the first non-empty chunk and emitting a delta when a Publisher is set.

Opening LAZILY is deliberate: a round that produces no text of its kind must not emit an empty block, or the client renders a stray blank card and the index advances for nothing.

type TimelineItem

type TimelineItem struct {
	Kind      TimelineItemKind `json:"kind"`
	Text      string           `json:"text,omitempty"`
	Execution *ToolExecution   `json:"execution,omitempty"`
}

TimelineItem is one ordered entry produced by reassembly: a reasoning or text segment, or a completed tool execution (call + result pair).

type TimelineItemKind

type TimelineItemKind = string

TimelineItemKind marks whether a timeline entry is reasoning, text, or a tool execution.

const (
	TimelineKindThinking TimelineItemKind = "thinking"
	TimelineKindText     TimelineItemKind = "text"
	TimelineKindTool     TimelineItemKind = "tool"
)

type ToolCall

type ToolCall struct {
	Name      string          `json:"name"`
	Params    json.RawMessage `json:"params"`
	ToolUseID string          `json:"tool_use_id"`
}

ToolCall is one reassembled tool_use block.

type ToolExecution

type ToolExecution struct {
	Name      string          `json:"name"`
	Params    json.RawMessage `json:"params"`
	Result    string          `json:"result"`
	ToolUseID string          `json:"tool_use_id"`
}

ToolExecution is a matched tool_use + tool_result pair.

type ToolResultBlock

type ToolResultBlock struct {
	Type      ContentBlockType `json:"type"`
	ToolUseID string           `json:"tool_use_id"`
	Content   string           `json:"content,omitempty"`
	IsError   bool             `json:"is_error,omitempty"`
}

type ToolResultDelta

type ToolResultDelta struct {
	Type ContentBlockType `json:"type"`
	Text string           `json:"text"`
}

type ToolUseBlock

type ToolUseBlock struct {
	Type  ContentBlockType `json:"type"`
	ID    string           `json:"id"`
	Name  string           `json:"name"`
	Input any              `json:"input"`
}

type UsageEnd

type UsageEnd struct {
	OutputTokens int `json:"output_tokens"`
}

type UsageStart

type UsageStart struct {
	InputTokens  int `json:"input_tokens"`
	OutputTokens int `json:"output_tokens"`
}

Directories

Path Synopsis
Package elelemstream bridges elelem's streaming callbacks to essessey's content-block protocol.
Package elelemstream bridges elelem's streaming callbacks to essessey's content-block protocol.
Package nats provides essessey.Sink and essessey.Source bindings for NATS.
Package nats provides essessey.Sink and essessey.Source bindings for NATS.
Package sse is the SSE FORMAT binding for essessey: the wire codec plus the byte-stream Sinks (WriterSink, HTTPSink) and the byte-stream Source that frame/parse it.
Package sse is the SSE FORMAT binding for essessey: the wire codec plus the byte-stream Sinks (WriterSink, HTTPSink) and the byte-stream Source that frame/parse it.
Package ws provides essessey.Sink and essessey.Source bindings for WebSocket connections.
Package ws provides essessey.Sink and essessey.Source bindings for WebSocket connections.

Jump to

Keyboard shortcuts

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