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 ¶
- Variables
- type ContentBlock
- type ContentBlockDeltaData
- type ContentBlockDeltaToolInputData
- type ContentBlockDeltaToolResultData
- type ContentBlockStartData
- type ContentBlockStartToolResultData
- type ContentBlockStartToolUseData
- type ContentBlockStopData
- type ContentBlockType
- type Event
- type EventStore
- type EventType
- type InMemoryEventStore
- func (s *InMemoryEventStore) Append(_ context.Context, streamID string, ev Event) error
- func (s *InMemoryEventStore) Clear(_ context.Context, streamID string) error
- func (s *InMemoryEventStore) Since(_ context.Context, streamID string, lastEventID string) ([]Event, bool, error)
- func (s *InMemoryEventStore) SinkFor(streamID string) Sink
- type InMemorySink
- type InputJSONDelta
- type LineStreamer
- type LineTransform
- type MessageDeltaData
- type MessageDeltaInfo
- type MessageMeta
- type MessageStartData
- type MessageStopData
- type MessageType
- type MultiSink
- type ParsedStream
- type PingData
- type Publisher
- func (p *Publisher) Publish(eventType EventType, data any) error
- func (p *Publisher) SendContentBlockDeltaText(index int, text string) error
- func (p *Publisher) SendContentBlockDeltaThinking(index int, text string) error
- func (p *Publisher) SendContentBlockStartText(index int) error
- func (p *Publisher) SendContentBlockStartThinking(index int) error
- func (p *Publisher) SendContentBlockStop(index int) error
- func (p *Publisher) SendMessageDelta(stopReason StopReason, outputTokens int) error
- func (p *Publisher) SendMessageStart(msgID, streamID, model string) error
- func (p *Publisher) SendMessageStop() error
- func (p *Publisher) SendPing() error
- func (p *Publisher) SendStreamEpilogue(stopReason StopReason, outputTokens int) error
- func (p *Publisher) SendStreamPreamble(msgID, streamID, model string) error
- func (p *Publisher) SendToolInputDelta(index int, inputJSON string) error
- func (p *Publisher) SendToolResultBlock(index int, toolUseID, resultText string, isError bool) error
- func (p *Publisher) SendToolResultDelta(index int, text string) error
- func (p *Publisher) SendToolResultStart(index int, toolUseID string, isError bool) error
- func (p *Publisher) SendToolUseBlock(index int, toolUseID, name, inputJSON string) error
- func (p *Publisher) SendToolUseStart(index int, toolUseID, name string) error
- type Role
- type Sink
- type SliceSource
- type Source
- type StopReason
- type TextDelta
- type TextStreamer
- type TimelineItem
- type TimelineItemKind
- type ToolCall
- type ToolExecution
- type ToolResultBlock
- type ToolResultDelta
- type ToolUseBlock
- type UsageEnd
- type UsageStart
Constants ¶
This section is empty.
Variables ¶
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.
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 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 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
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.
type LineTransform ¶
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
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
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 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 ¶
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 ¶
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 ¶
SendContentBlockDeltaText appends a text delta to the block at index.
func (*Publisher) SendContentBlockDeltaThinking ¶
SendContentBlockDeltaThinking appends a reasoning delta at index.
func (*Publisher) SendContentBlockStartText ¶
SendContentBlockStartText opens a text content block at index.
func (*Publisher) SendContentBlockStartThinking ¶
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 ¶
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 ¶
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 ¶
SendMessageStop emits the terminal message_stop 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 ¶
SendStreamPreamble emits message_start + ping to open a stream.
func (*Publisher) SendToolInputDelta ¶
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 ¶
SendToolResultDelta streams the tool result payload at index.
func (*Publisher) SendToolResultStart ¶
SendToolResultStart opens a tool_result content block for toolUseID.
func (*Publisher) SendToolUseBlock ¶
SendToolUseBlock emits a full tool_use block (start + input delta + stop).
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 ¶
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.
type Source ¶
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) 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 UsageStart ¶
Source Files
¶
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. |