Documentation
¶
Overview ¶
Package sse holds the server-sent-events plumbing every clicky server needs, so handlers stop hand-writing it:
- Writer frames events onto an http.ResponseWriter (lazy headers, multi-line data, comment pings) and is what rpc's streaming operations write through.
- Snapshot streams a value that is re-loaded on a tick or a wake, sending it only when it changed (optionally ignoring volatile fields such as a rate limit or a sample counter) and pinging otherwise.
- Notifier is a fan-out wake-up: one Notify wakes every subscribed stream, coalescing bursts, unlike a shared channel where streams steal each other's signals.
- Hub multiplexes every stream of a browser tab over one connection, since a browser allows only six HTTP/1.1 connections per host across all its tabs, and Guard refuses the direct per-topic EventSource connections a stale page would otherwise hold open.
A minimal server wiring all four:
mux := http.NewServeMux()
var changed sse.Notifier // call changed.Notify() whenever the data changes
mux.HandleFunc("GET /api/items/stream", func(w http.ResponseWriter, r *http.Request) {
wake, cancel := changed.Subscribe()
defer cancel()
err := sse.ServeSnapshot(w, r, sse.SnapshotOptions{
Load: func(ctx context.Context) (any, error) { return loadItems(ctx) },
Interval: 5 * time.Second,
Wake: wake,
Exclude: []string{"fetchedAt"},
})
if err != nil {
logger.Debugf("items stream: %v", err)
}
})
hub := sse.NewHub(sse.HubOptions{Build: buildID})
root := hub.Guard(mux)
hub.Register(mux, root)
http.ListenAndServe(":8080", root)
A browser then opens GET /api/events once and subscribes to /api/items/stream over it (clicky-ui's createEventHub does this).
Index ¶
Constants ¶
const DefaultPrefix = "/api/events"
DefaultPrefix is where a Hub mounts its routes unless HubOptions.Prefix says otherwise.
Variables ¶
This section is empty.
Functions ¶
func ServeSnapshot ¶
func ServeSnapshot(w http.ResponseWriter, r *http.Request, opts SnapshotOptions) error
ServeSnapshot streams opts as the response to r through a Writer. It returns nil when the client went away and Snapshot's error otherwise (a load error has already been sent to the client as an "error" event by then).
func Snapshot ¶
func Snapshot(ctx context.Context, send entity.StreamSend, ping func() error, opts SnapshotOptions) error
Snapshot streams opts.Load's value: it loads immediately and then on every tick or wake, sends the value as JSON when it differs from the last one sent, and calls ping otherwise so idle connections stay alive and a departed client is noticed.
send and ping are separate because entity.StreamSend has no comment frame: a plain handler passes a Writer's Send and Ping (or uses ServeSnapshot), while an rpc StreamFunc passes its send and a ping of its own choosing — a func returning nil opts out of keepalives explicitly.
A Load (or encoding) failure is sent as an "error" event {"error": msg} and returned, ending the stream: a reader learns why it stopped. Snapshot returns ctx.Err() once ctx ends and a send or ping error as soon as one occurs. It panics on options it cannot run.
Types ¶
type Hub ¶
type Hub struct {
// contains filtered or unexported fields
}
Hub tracks the open events connections of one handler tree. mux resolves a sub's path to a registered route; root is the full top-level handler (middlewares included) the sub is served through.
func NewHub ¶
func NewHub(opts HubOptions) *Hub
NewHub returns a Hub; it panics on options it cannot serve.
func (*Hub) Guard ¶
Guard answers 410 Gone to a browser opening a native EventSource on a stream route the hub can serve instead. A page loaded from a bundle that predates the hub holds one such connection per topic and, under the browser's six-connections-per-host cap, starves every ordinary request of the tab and its siblings. EventSource treats a non-200 response as fatal and does not reconnect, so the stale page releases the connection.
Only browsers are refused — they alone send Sec-Fetch-Mode, so CLI and Go clients keep direct streams — and only GETs, so POST launch streams are left alone. Paths the hub cannot subscribe to (outside /api/, e.g. a page's embedded bundle serving its own streams) have no replacement and stay served.
func (*Hub) Register ¶
Register mounts the hub routes on mux. root must be the handler the server actually serves (mux wrapped in its middlewares, Guard included) so subs see exactly what a direct request to the same path would.
Every stream handler a sub can reach must return once its request context is cancelled. A handler that does not is abandoned hubSubStopTimeout after its sub is stopped and keeps its goroutine and resources until it returns.
type HubOptions ¶
type HubOptions struct {
// Build identifies the UI bundle the server ships. Every hello frame carries
// it so a page loaded from an older bundle notices it is stale and reloads.
// Required.
Build string
// Prefix is the path the hub's routes live under. Default DefaultPrefix.
Prefix string
}
HubOptions configures a Hub.
type Notifier ¶
type Notifier struct {
// contains filtered or unexported fields
}
Notifier wakes every subscriber at once. A single shared channel would hand each signal to whichever reader got there first, so concurrent streams (one per open tab) would steal each other's wake-ups. Each subscriber instead owns a channel buffered to one: a Notify never blocks, and notifications a subscriber has not consumed yet coalesce into a single pending wake.
The zero value is ready to use. A Notifier must not be copied after first use.
func (*Notifier) Notify ¶
func (n *Notifier) Notify()
Notify wakes every current subscriber without blocking.
func (*Notifier) Subscribe ¶
func (n *Notifier) Subscribe() (<-chan struct{}, func())
Subscribe returns a channel that receives a value after each Notify (or burst of them), and a cancel that unsubscribes. cancel is idempotent and does not close the channel, so a select over it never sees a spurious wake.
type SnapshotOptions ¶
type SnapshotOptions struct {
// Load produces the current value. It is called once immediately, then on
// every Interval tick and every Wake. Required.
Load func(ctx context.Context) (any, error)
// Interval is how often Load is polled. Required (> 0).
Interval time.Duration
// Wake, when non-nil, triggers an immediate Load (e.g. a Notifier
// subscription). Closing it is an error: Snapshot cannot tell a closed wake
// from a spinning one.
Wake <-chan struct{}
// Exclude lists dotted JSON paths left out when deciding whether the value
// changed — fields that churn without meaning anything to a reader, such as
// a fetch timestamp or a sampled counter. A "*" segment matches every map
// key or array element: "rateLimit", "*.processes.*.openFiles". They are
// excluded only from the comparison: a frame sent because something else
// changed carries their current values.
Exclude []string
// Event names the frames. Empty sends unnamed ("message") events.
Event string
}
SnapshotOptions describes a value Snapshot streams.
type Writer ¶
type Writer struct {
// contains filtered or unexported fields
}
Writer writes the event-stream framing onto one response. Nothing is written — not even the headers — until the first Send or Ping, so a handler that fails before producing anything can still answer with a status.
A Writer is not safe for concurrent use: one goroutine owns a response.
func (*Writer) Ping ¶
Ping writes a comment frame. Clients ignore it; it keeps idle proxies from closing the connection and lets the server notice a client that has gone.