stream

package
v0.85.0 Latest Latest
Warning

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

Go to latest
Published: Sep 8, 2026 License: MIT Imports: 20 Imported by: 0

Documentation

Overview

Package stream is part of the GoFastr framework. See https://github.com/DonaldMurillo/gofastr for documentation.

Index

Constants

View Source
const SnapshotEnvelopeType = "snapshot"

SnapshotEnvelopeType is the Type a StateChannel snapshot envelope carries. Every other Type the channel emits comes from the Publish call, so an application can reserve names besides this one.

Variables

View Source
var ErrClosed = errors.New("stream: connection closed")

ErrClosed is returned when writing to a closed connection. It is a plain errors.New sentinel so callers can compare with errors.Is.

Functions

func Encode

func Encode(e Event) string

Encode formats an Event as a W3C-compliant SSE frame.

The output uses the following fields:

  • "id:" when Event.ID is non-empty
  • "event:" for Error, Done, and Custom types
  • "data:" for the payload (multi-line data splits on \n)
  • terminated by a blank line ("\n\n")

ID and custom event names are truncated at the first CR/LF/NUL; those bytes terminate an SSE field and would let a caller-supplied value inject forged directives ("event: forged", "data: pwned"…) below it. Multi-line data is split on '\n' and each line is re-prefixed with "data: " so an injected blank line ("\n\n") still appears as a single event frame to the client. CRs and NULs in data are stripped for the same reason.

func LastEventID

func LastEventID(r *http.Request) string

LastEventID returns the Last-Event-ID from the request headers or the "last_event_id" query parameter. The value is truncated at the first CR/LF/NUL so a malicious resume token can't inject forged SSE fields when later echoed back to clients.

Types

type ChunkedWriter

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

ChunkedWriter writes raw chunks to an http.ResponseWriter with flush.

func NewChunkedWriter

func NewChunkedWriter(w http.ResponseWriter) *ChunkedWriter

NewChunkedWriter creates a ChunkedWriter wrapping w. It panics if w does not implement http.Flusher.

func (*ChunkedWriter) Close

func (c *ChunkedWriter) Close() error

Close performs a final flush. It is safe to call multiple times.

func (*ChunkedWriter) WriteChunk

func (c *ChunkedWriter) WriteChunk(data []byte) error

WriteChunk writes data to the underlying response writer and flushes.

type Event

type Event struct {
	Type EventType
	Name string // used when Type == Custom
	Data string
	ID   string // optional Last-Event-ID value
}

Event represents a single Server-Sent Event.

type EventType

type EventType int

EventType represents the kind of SSE event.

const (
	// Message is a standard data event.
	Message EventType = iota
	// Error signals an error to the client.
	Error
	// Done is the terminal sentinel event.
	Done
	// Custom is a named event with an arbitrary event type string.
	Custom
)

func (EventType) String

func (t EventType) String() string

String returns a human-readable label for the EventType.

type Hub

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

Hub manages a set of WebSocket connections and broadcasts messages to all of them. Multiple hubs can coexist (one per room, channel, etc.).

Usage:

hub := stream.NewHub()
go hub.Run()

// On connect:
hub.Register(conn)

// Broadcast:
hub.Broadcast([]byte("hello everyone"))

// On disconnect:
hub.Unregister(conn)

func NewHub

func NewHub() *Hub

NewHub creates a new Hub.

func (*Hub) Broadcast

func (h *Hub) Broadcast(msg []byte)

Broadcast sends a message to all registered connections. Non-blocking: if the hub's broadcast channel is full, the message is dropped.

func (*Hub) BroadcastWait

func (h *Hub) BroadcastWait(msg []byte)

BroadcastWait sends a message to all connections, blocking if the broadcast channel is full.

func (*Hub) Count

func (h *Hub) Count() int

Count returns the number of active connections.

func (*Hub) Register

func (h *Hub) Register(conn *WebSocketConn)

Register adds a connection to the hub. Non-blocking. Returns immediately if the hub has been stopped.

func (*Hub) Run

func (h *Hub) Run()

Run starts the hub's event loop. Block until Stop is called. Must be called in a goroutine:

go hub.Run()

Register/Unregister are now mutex-only ops on the connections map, so Run is only responsible for draining the broadcast channel and shutting down on Stop.

func (*Hub) Stop

func (h *Hub) Stop()

Stop stops the hub and closes all registered connections.

func (*Hub) Unregister

func (h *Hub) Unregister(conn *WebSocketConn)

Unregister removes a connection from the hub. Non-blocking.

type SSEBroker

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

SSEBroker fans out SSE events to multiple HTTP subscribers. Each subscriber gets a buffered channel. Default subscribers drop the oldest queued event when the buffer is full; clients that opt into ?slow=block or X-SSE-Slow: block instead backpressure Publish until buffer space is available.

Buffer size is configurable per-subscriber via query param (?buffer=128) or header (X-SSE-Buffer), with a default fallback bounded by MaxBuf.

When SSEBrokerConfig.Fanout is set, Publish also broadcasts to other replicas (topic "gofastr.sse.<Topic>") and events from other replicas are delivered locally. Delivery is lossy best-effort (the real-time lane).

func NewSSEBroker

func NewSSEBroker(cfg SSEBrokerConfig) *SSEBroker

NewSSEBroker creates a new broker for fan-out SSE delivery.

func (*SSEBroker) Close added in v0.16.0

func (b *SSEBroker) Close()

Close tears down the broker's fanout participation entirely: the receive subscription is cancelled and subsequent Publish calls stop crossing replicas. It is a no-op when no fanout is attached. The broker has no other goroutines to stop; its subscribers live and die with their HTTP request contexts. Safe to call multiple times.

func (*SSEBroker) Publish

func (b *SSEBroker) Publish(name, data string, id ...string)

Publish sends an event to all subscribers. If a default subscriber's buffer is full, the oldest event is dropped. A subscriber that opted into slow=block backpressures this call until buffer space opens or that subscriber is closed. Subscribers are snapshotted under the read lock; sends happen outside the lock to keep fan-out from holding the broker lock during slow per-channel writes.

func (*SSEBroker) Subscribe

func (b *SSEBroker) Subscribe(w http.ResponseWriter, r *http.Request)

Subscribe adds a subscriber and blocks, writing events to the response. The subscriber ID is taken from ?subscriber_id or X-Subscriber-ID header. Buffer size from ?buffer= or X-SSE-Buffer header, clamped to MaxBuf. Subscribe returns when the request context is canceled or the client disconnects.

func (*SSEBroker) SubscriberCount

func (b *SSEBroker) SubscriberCount() int

SubscriberCount returns the number of active subscribers.

type SSEBrokerConfig

type SSEBrokerConfig struct {
	Topic             string        // logical topic name (for logging/debugging)
	DefaultBuf        int           // default subscriber buffer size (0 = 64)
	MaxBuf            int           // maximum allowed subscriber buffer (0 = 1024)
	HeartbeatInterval time.Duration // 0 = 30s; emits a comment frame to keep idle connections open

	// Fanout, when set, makes Publish cross replicas (topic
	// "gofastr.sse.<Topic>") and re-delivers other replicas' events locally.
	// Remote-origin events are delivered with ALWAYS drop-oldest semantics
	// (even to slow=block subscribers): a remote replica cannot be
	// backpressured through a channel send, so a single stalled subscriber
	// must never wedge receive for the others. Local Publish keeps its
	// block-mode backpressure contract. Optional. Close cancels the
	// subscription.
	Fanout fanout.Fanout

	// AllowClientSlowMode lets a REQUEST select block mode via
	// ?slow=block / X-SSE-Slow. Off by default.
	//
	// deliver() walks subscribers sequentially on the publisher's
	// goroutine, so a block-mode subscriber that stops reading stalls
	// every other subscriber AND whatever called Publish, usually a
	// request handler. On a public endpoint that is an unauthenticated
	// denial of service, so the choice belongs to the developer who
	// knows whether the endpoint is trusted, not to the caller.
	AllowClientSlowMode bool

	// BlockTimeout bounds how long a block-mode send may stall before
	// the broker gives up on that subscriber and moves on. 0 = 5s.
	// A blocking send with no timeout is unbounded backpressure.
	BlockTimeout time.Duration

	// Principal identifies the caller behind a request, so a reconnect
	// with the same ?subscriber_id replaces its OWN entry and nobody
	// else's. Return "" for "cannot tell", which is treated as "not the
	// same caller".
	//
	// Set this to something the caller cannot choose and another caller
	// cannot guess: a session user id is the usual answer.
	//
	// Left nil, the broker never evicts. The old default keyed on
	// RemoteAddr's host, which is honest only on a direct connection:
	// behind nginx, an ALB, Cloudflare or a k8s ingress, every request's
	// TCP peer is the proxy, so all subscribers collapsed to one
	// principal and `?subscriber_id=<victim>` dropped the victim's
	// stream, repeatably. Nothing is lost by not evicting: subscriber
	// ids address nothing (deliver() broadcasts), and a dropped
	// connection already unregisters itself.
	Principal func(*http.Request) string

	// MaxSeatsPerPrincipal caps how many concurrent streams ONE
	// principal holds. 0 = 16 (the same default core/mcp's SSE stream
	// uses); a negative value lifts the cap for deployments that bound
	// seats elsewhere. The zero-value config gets the default: a single
	// low-privilege user must not be able to exhaust the process's
	// goroutines/FDs/memory for everyone else, and MaxSubscribers (the
	// all-users cap) is opt-in unlimited by design.
	//
	// When Principal is nil (or returns ""), callers cannot be told
	// apart, so every anonymous stream shares ONE seat bucket: the cap
	// then bounds all unidentifiable subscribers collectively. Raise it
	// (or set Principal) on a public stream that legitimately serves
	// more than 16 anonymous viewers.
	MaxSeatsPerPrincipal int

	// SeatOverflow selects what a principal at MaxSeatsPerPrincipal does
	// with its next stream: SeatOverflowRefuse answers 429 at connect
	// (the default), SeatOverflowEvictOldest closes that principal's
	// oldest stream and seats the new one — the multi-tab-friendly
	// policy for apps whose clients reconnect faster than a half-open
	// connection's seat is reclaimed.
	SeatOverflow SeatOverflowPolicy

	// MaxSubscribers caps concurrent subscribers; 0 = unlimited. Subscribe
	// rejects past the cap rather than evicting, and the cap is exact.
	//
	// A client whose previous connection is half-open (mobile handoff, laptop
	// sleep, an LB idle-kill the server has not noticed) still holds a seat
	// until HeartbeatInterval's next write fails and the stream unregisters
	// itself. That heartbeat is what reclaims the seat: for every client,
	// including the ones that send no subscriber_id, which is all of the
	// framework's own. A reserved slot keyed on the requested id was tried
	// and removed: nothing in the client runtime sends an id, so it could
	// never fire, while costing a scan of every subscriber under the write
	// lock and letting the cap be exceeded by one.
	MaxSubscribers int
}

SSEBrokerConfig configures the broker.

type SSEWriter

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

SSEWriter writes Server-Sent Events to an http.ResponseWriter.

It automatically sets the required headers (Content-Type, Cache-Control, Connection) on first write and flushes after every event.

func NewSSEWriter

func NewSSEWriter(w http.ResponseWriter) *SSEWriter

NewSSEWriter creates an SSEWriter wrapping w. It panics if w does not implement http.Flusher.

func (*SSEWriter) Flush

func (s *SSEWriter) Flush()

Flush sends any buffered data to the client immediately.

func (*SSEWriter) SetID

func (s *SSEWriter) SetID(id string)

SetID queues an "id:" field to be emitted before the next event.

func (*SSEWriter) SetRetry

func (s *SSEWriter) SetRetry(seconds int)

SetRetry writes the "retry:" field, telling the client how many milliseconds to wait before reconnecting. Non-positive values are dropped: `retry: 0` tells the client to reconnect with zero delay, which spins into a reconnect storm, an accidental DoS amplifier.

func (*SSEWriter) WriteComment

func (s *SSEWriter) WriteComment(comment string) error

WriteComment writes an SSE comment (keepalive):

: <comment>

followed by a blank line and a flush. The comment is truncated at the first CR/LF so a caller can't terminate the comment line and inject arbitrary SSE fields ("event: …", "data: …", …) below it.

func (*SSEWriter) WriteData

func (s *SSEWriter) WriteData(data string) error

WriteData writes an anonymous SSE event (type defaults to "message"):

data: <data>

followed by a blank line and a flush.

func (*SSEWriter) WriteDone

func (s *SSEWriter) WriteDone() error

WriteDone sends the terminal "[DONE]" sentinel.

func (*SSEWriter) WriteError

func (s *SSEWriter) WriteError(message string) error

WriteError is a convenience for writing an Error event.

func (*SSEWriter) WriteEvent

func (s *SSEWriter) WriteEvent(event, data string) error

WriteEvent writes a named SSE event:

event: <name>
data: <data>

followed by a blank line and a flush.

CR/LF characters in the event name are stripped: an event name may only occupy a single SSE field line. A caller-supplied newline would otherwise terminate the field and let following bytes appear as arbitrary SSE directives.

func (*SSEWriter) WriteMessage

func (s *SSEWriter) WriteMessage(data string) error

WriteMessage is a convenience for writing a Message event.

type SeatOverflowPolicy added in v0.85.0

type SeatOverflowPolicy uint8

SeatOverflowPolicy selects what a principal at MaxSeatsPerPrincipal does with its next stream.

const (
	// SeatOverflowRefuse answers the connection with 429 and holds no
	// seat — distinguishable from the global MaxSubscribers' 503, which
	// is about the server, not the caller.
	SeatOverflowRefuse SeatOverflowPolicy = iota
	// SeatOverflowEvictOldest closes that principal's oldest stream and
	// seats the new one. The reconnect-friendly policy: a client whose
	// previous connection is half-open (mobile handoff, laptop sleep)
	// gets its new stream instead of a 429 until the old seat's
	// heartbeat write fails.
	SeatOverflowEvictOldest
)

type SequencedEnvelope added in v0.80.0

type SequencedEnvelope[T any] struct {
	Sequence uint64    `json:"sequence"`
	Type     string    `json:"type"`
	Payload  T         `json:"payload"`
	SentAt   time.Time `json:"sentAt"`
}

SequencedEnvelope is the wire shape a StateChannel puts on every message it delivers: the hydration snapshot and every live event. The client-side companion (core-ui/runtime/src/ws.js, createSequencedReducer) applies an envelope only when its sequence is strictly greater than the last applied one, so a delayed snapshot can never resurrect state a newer event already replaced.

Sequence is a uint64 on the wire and a Number in the browser, which compares exactly up to 2^53 (nine quadrillion events on one channel). That bound is stated rather than engineered around: a string-or-BigInt protocol would cost every consumer a conversion for a limit no deployment reaches.

type SnapshotSource added in v0.80.0

type SnapshotSource[Role comparable, Snapshot any, Event any] interface {
	// SnapshotFor returns the role's view of the current state and the
	// sequence of that state, from one immutable read. Called from the
	// channel's Run loop: a slow read delays event delivery for every
	// connection on the channel, so keep it an in-memory copy.
	SnapshotFor(Role) (Snapshot, uint64)

	// FilterEvent returns the payload a role may receive for an event,
	// and whether the role may see it at all. Called once per role that
	// currently has connections, before the envelope is marshaled.
	FilterEvent(Role, Event) (any, bool)
}

SnapshotSource is the application-owned half of a StateChannel. The channel never stores business state; it asks the source for a role's view on connect and asks it to shape each event per role before anything is serialized.

Sequences share ONE space per channel, across roles. SnapshotFor must return the snapshot payload and its sequence from a single immutable read (copy-on-read), where the sequence is the channel-wide version of the state returned. FilterEvent then sees only events; the channel reconciles its own event counter to every snapshot sequence it sends, so events published after a snapshot always sort above it.

Minimization happens here, at the source: FilterEvent runs BEFORE serialization, so a field the source strips never crosses the transport. Hiding a field in the UI is not data minimization.

type StateChannel added in v0.80.0

type StateChannel[Role comparable, Snapshot any, Event any] struct {
	// contains filtered or unexported fields
}

StateChannel layers reconnect hydration and ordered live events over WebSocketConn. It is a small helper above Hub, not a state framework: persistence and business state stay application-owned, and the channel owns the envelope shape, sequencing, initial hydration, and per-role filtering.

Usage:

channel := stream.NewStateChannel(source)
go channel.Run()

// On connect (handler goroutine); returns once the snapshot is
// queued on the connection:
channel.Connect(role, conn)

// After the application has applied a mutation to its own state:
channel.Publish("cleared", event)

Ordering contract, per connection, guaranteed on the wire:

  • The hydration snapshot for a connection is sent before any event published after that connection's snapshot was queued, and every such event carries a sequence greater than the snapshot's.
  • Events reach each connection in strictly increasing sequence order, so the client reducer's reject-stale rule never discards a live event as a side effect.
  • An event whose mutation a snapshot already contains may still be delivered after it (at-least-once); sequences make the client state convergent, not exactly-once.

Delivery is best-effort like Hub: a connection whose send buffer is full has events dropped for it, and a connection that cannot accept its hydration snapshot at all is closed (it cannot catch up on its own, so it reconnects and hydrates again). Publish never blocks; when the channel's job queue is full the event is dropped, mirroring Hub.Broadcast.

func NewStateChannel added in v0.80.0

func NewStateChannel[Role comparable, Snapshot any, Event any](
	source SnapshotSource[Role, Snapshot, Event],
) *StateChannel[Role, Snapshot, Event]

NewStateChannel creates a StateChannel over the given source. Call Run in a goroutine before connecting connections.

func (*StateChannel[Role, Snapshot, Event]) CloseOnOverflow added in v0.83.0

func (c *StateChannel[Role, Snapshot, Event]) CloseOnOverflow(on bool)

CloseOnOverflow selects what happens when an event cannot be queued on a connection because its send buffer is full. Off (the default) the event is dropped for that connection alone, the Hub.Run posture, which is right for presence-shaped sources: the next snapshot or event corrects the view. On, the connection is closed instead, so the client reconnects and re-hydrates. That is right for a source whose events are not recoverable from a snapshot (addressed signaling: a dropped SDP answer leaves the far side waiting for ever, and nothing tells either side). Call before Run.

func (*StateChannel[Role, Snapshot, Event]) Connect added in v0.80.0

func (c *StateChannel[Role, Snapshot, Event]) Connect(role Role, conn *WebSocketConn)

Connect hydrates conn with the role's snapshot and keeps it live for published events. It blocks until the snapshot is queued on the connection (or the connection/channel is closed), so on return the handler knows hydration happened.

If the channel is stopped, or its job queue is full, the connection is closed: a connection that never hydrates cannot catch up.

func (*StateChannel[Role, Snapshot, Event]) Count added in v0.80.0

func (c *StateChannel[Role, Snapshot, Event]) Count() int

Count returns the number of registered connections.

func (*StateChannel[Role, Snapshot, Event]) Publish added in v0.80.0

func (c *StateChannel[Role, Snapshot, Event]) Publish(typ string, event Event)

Publish broadcasts an event to every connected role. typ becomes the envelope's Type; the payload each role receives is FilterEvent's return value for that role. Non-blocking: returns immediately; the event is dropped if the channel is stopped or its job queue is full.

func (*StateChannel[Role, Snapshot, Event]) Run added in v0.80.0

func (c *StateChannel[Role, Snapshot, Event]) Run()

Run starts the channel's event loop. Block until Stop is called:

go channel.Run()

The loop is the single owner of sequencing: it dequeues snapshot and event jobs in FIFO order, assigns event sequences at dequeue time (so the wire order per connection is the sequence order), and is the only writer of nextSeq.

func (*StateChannel[Role, Snapshot, Event]) Stop added in v0.80.0

func (c *StateChannel[Role, Snapshot, Event]) Stop()

Stop stops the channel and closes every registered connection. Safe to call more than once. Publish and Connect after Stop are no-ops (Connect closes the connection it was given).

func (*StateChannel[Role, Snapshot, Event]) Unregister added in v0.80.0

func (c *StateChannel[Role, Snapshot, Event]) Unregister(conn *WebSocketConn)

Unregister removes a connection without closing it. Connections are unregistered automatically when they close, so this is only for explicit removal.

type WSConfig

type WSConfig struct {
	// ReadLimit is the maximum message size in bytes. 0 = default 1MB.
	ReadLimit int64

	// SendBuffer is the number of messages that can be buffered before
	// Write blocks. 0 = 32.
	SendBuffer int

	// WriteTimeout bounds each frame write. 0 means default 10s. Set
	// negative to disable (not recommended): a peer with a full TCP send
	// buffer otherwise pins the writePump and keepalive goroutines forever.
	WriteTimeout time.Duration

	// CheckOrigin returns true if the Origin header is acceptable.
	// If nil, Upgrade enforces same-origin by comparing Origin host to
	// the request Host. Use a custom CheckOrigin to allow cross-origin
	// upgrades (e.g. for trusted third-party clients).
	CheckOrigin func(*http.Request) bool

	// ReadIdleTimeout bounds the longest period of read inactivity before
	// the keepalive sends a Ping. 0 means default 60s. Set negative to
	// disable keepalive entirely.
	ReadIdleTimeout time.Duration

	// PongTimeout bounds how long after a Ping we wait for the matching
	// Pong. If exceeded, the connection is closed. 0 means default 10s.
	// Set negative to disable the pong timeout check.
	PongTimeout time.Duration

	// CloseTimeout caps how long Close() waits for the peer's reciprocal
	// Close frame after sending our own. 0 means default 1s.
	CloseTimeout time.Duration

	// Subprotocols is the server's preferred list of WebSocket subprotocols
	// in priority order. During Upgrade, the first subprotocol that the
	// client offered AND we support is echoed back via
	// Sec-WebSocket-Protocol. If no match, no header is sent (RFC 6455).
	Subprotocols []string

	// ConnectionID identifies this connection in logs and reconnect
	// bookkeeping. Empty (the default) means Upgrade generates a random
	// id. Set it explicitly when the application mints its own ids
	// (e.g. one per browser session, so a client's reconnect after a
	// socket drop correlates across server-side logs). See
	// WebSocketConn.ConnectionID for the reconnect-generation contract.
	ConnectionID string

	// OnClose is called when the connection closes.
	OnClose func()
	// contains filtered or unexported fields
}

WSConfig configures the WebSocket connection.

type WebSocketConn

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

WebSocketConn wraps a hijacked HTTP connection as a simple WebSocket client. It implements a minimal WebSocket frame parser sufficient for the framework's needs: text and binary messages, close frames.

For production use with full RFC 6455 compliance, use a dedicated WebSocket library (nhooyr.io/websocket, gorilla/websocket). This implementation avoids external dependencies so the core/stream package compiles without `go get` additions.

Backpressure: writes block when the send buffer is full; reads come from a small internal buffer fed by the connection's own read pump. The read pump owns ALL socket reads and runs from the moment the connection starts, so control frames (Ping/Pong/Close) are handled even if the application never calls Read. Read() consumes complete messages from that buffer; beyond it, TCP backpressure applies. A push-only server needs no read loop; it stays alive while a healthy peer answers Pongs, and Close()'s handshake completes promptly when the peer reciprocates.

func Upgrade

func Upgrade(w http.ResponseWriter, r *http.Request, cfg WSConfig) (*WebSocketConn, error)

Upgrade upgrades an HTTP connection to a simple WebSocket. Performs the HTTP upgrade handshake and returns a managed connection.

func (*WebSocketConn) Close

func (c *WebSocketConn) Close() error

Close closes the WebSocket connection. Safe to call multiple times. Performs the RFC 6455 closing handshake: sends a Close frame, then waits up to CloseTimeout for the peer's reciprocal Close before tearing down the underlying TCP connection. This avoids the abnormal 1006 close code on the peer side.

If the peer initiated the close, the echo Close frame preserves the peer's 2-byte status code per RFC 6455 §5.5.1, sanitized per §7.4.1 so a reserved (1004/1005/1006/1015), sub-1000, or otherwise invalid code, or an illegal 1-byte body, is never echoed verbatim and is replaced by 1002. Otherwise we send an empty Close payload (status 1000 implied by absence).

func (*WebSocketConn) CloseWithStatus added in v0.83.0

func (c *WebSocketConn) CloseWithStatus(code uint16, reason string) error

CloseWithStatus is Close with a status code and reason in the close frame, for a server that accepted the handshake only to refuse the connection: a browser cannot read an HTTP status off a failed handshake (the WHATWG spec withholds it so a page cannot probe the network), but it can read a close code. The reason is kept to printable ASCII and at most 123 bytes, the frame's limit. A close the peer initiated first still echoes the peer's own code.

func (*WebSocketConn) Closed

func (c *WebSocketConn) Closed() <-chan struct{}

Closed returns a channel closed when the connection closes.

func (*WebSocketConn) ConnectionID added in v0.80.0

func (c *WebSocketConn) ConnectionID() string

ConnectionID returns the connection's stable identifier: the WSConfig.ConnectionID value, or a random id minted at Upgrade.

The id distinguishes server-side connections from each other; it does NOT by itself model reconnects. A client that reconnects gets a NEW connection and therefore a new id (or a fresh browser-side generation, see core-ui/runtime/src/ws.js). Applications that need "same user, new transport" semantics correlate the two explicitly: echo this id to the client on connect, or key the client's generation counter against it in logs.

func (*WebSocketConn) OnClose

func (c *WebSocketConn) OnClose(fn func())

OnClose registers a callback for when the connection closes.

func (*WebSocketConn) Read

func (c *WebSocketConn) Read() ([]byte, error)

Read reads a complete message from the client. It consumes from the internal read pump's buffer (readMsgs); it never reads the socket directly, so the pump remains the sole owner of the wire and control frames stay processed even between Read calls. A push-only server need not call Read at all; the pump still answers Pings, clears the keepalive's Pong watch, and completes Close()'s handshake.

Messages decoded by the pump before a terminal error are delivered before the error surfaces (drain-before-error). Concurrent Read callers are safe: the buffer channel serializes delivery.

func (*WebSocketConn) Write

func (c *WebSocketConn) Write(data []byte) error

Write sends a text message to the client. Blocks if the send buffer is full (backpressure). Returns an error if the connection is closed.

func (*WebSocketConn) WriteString

func (c *WebSocketConn) WriteString(data string) error

WriteString is a convenience for sending a text message.

Jump to

Keyboard shortcuts

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