Documentation
¶
Index ¶
- Constants
- Variables
- func Accept(key []byte) string
- func IsCloseError(err error, codes ...uint16) bool
- func IsNormalClose(err error) bool
- func New(handler func(*Conn) error) fh.HandlerFunc
- func NewEventHandler(h *EventHub, wsCfg Config, metadata MetadataFunc) fh.HandlerFunc
- func NewWithConfig(cfg Config, handler func(*Conn) error) fh.HandlerFunc
- type AuthFunc
- type AuthorizeFunc
- type BroadcastOptions
- type ClientIDFunc
- type Clock
- type CloseError
- type Config
- type Conn
- func (c *Conn) Close() error
- func (c *Conn) CloseWithStatus(code uint16, reason string) error
- func (c *Conn) Closed() bool
- func (c *Conn) LocalAddr() net.Addr
- func (c *Conn) NetConn() net.Conn
- func (c *Conn) NextReader() (opcode byte, r io.ReadCloser, err error)
- func (c *Conn) NextWriter(opcode byte) (io.WriteCloser, error)
- func (c *Conn) Ping(payload []byte) error
- func (c *Conn) Pong(payload []byte) error
- func (c *Conn) ReadJSON(v any) error
- func (c *Conn) ReadMessage() (opcode byte, payload []byte, err error)
- func (c *Conn) RemoteAddr() net.Addr
- func (c *Conn) SetReadDeadline(t time.Time) error
- func (c *Conn) SetWriteDeadline(t time.Time) error
- func (c *Conn) StartHeartbeat(interval, pongTimeout time.Duration, payload []byte, done <-chan struct{})
- func (c *Conn) Subprotocol() string
- func (c *Conn) WriteBinary(b []byte) error
- func (c *Conn) WriteJSON(v any) error
- func (c *Conn) WriteMessage(opcode byte, payload []byte) error
- func (c *Conn) WriteText(s string) error
- type Context
- type Envelope
- type EventConn
- func (c *EventConn) Ack(replyTo string, payload any) error
- func (c *EventConn) Close(code uint16, reason string) error
- func (c *EventConn) Closed() bool
- func (c *EventConn) CreatedAt() time.Time
- func (c *EventConn) Emit(event string, payload any) error
- func (c *EventConn) EmitScoped(topic, channel, event string, payload any) error
- func (c *EventConn) GetMeta(key string) string
- func (c *EventConn) Join(topic, channel string) error
- func (c *EventConn) LastSeen() time.Time
- func (c *EventConn) Leave(topic, channel string) error
- func (c *EventConn) Request(event string, payload any, timeout time.Duration) (Envelope, error)
- func (c *EventConn) RequestScoped(topic, channel, event string, payload any, timeout time.Duration) (Envelope, error)
- func (c *EventConn) Send(env Envelope) error
- func (c *EventConn) SendError(replyTo, code, msg string) error
- func (c *EventConn) SendRawJSON(raw []byte) error
- func (c *EventConn) SetMeta(key, value string)
- func (c *EventConn) Subscriptions() []Subscription
- type EventHub
- func (h *EventHub) Accept(ws *Conn, metadata map[string]string, requestedID string) (*EventConn, error)
- func (h *EventHub) Add(ws *Conn, metadata map[string]string) *EventConn
- func (h *EventHub) AddWithID(ws *Conn, metadata map[string]string, clientID string) *EventConn
- func (h *EventHub) Broadcast(topic, channel string, env Envelope, opts ...BroadcastOptions) error
- func (h *EventHub) BroadcastEvent(topic, channel, event string, payload any, opts ...BroadcastOptions) error
- func (h *EventHub) Client(id string) *EventConn
- func (h *EventHub) Clients() []*EventConn
- func (h *EventHub) Close() error
- func (h *EventHub) Count() int
- func (h *EventHub) EmitTo(clientID, event string, payload any) error
- func (h *EventHub) Handler(wsCfg Config, metadata MetadataFunc) fh.HandlerFunc
- func (h *EventHub) HandlerWithContext(wsCfg Config, metadata MetadataFunc) fh.HandlerFunc
- func (h *EventHub) Notify(topic, channel string, env Envelope, opts ...BroadcastOptions) error
- func (h *EventHub) NotifyEvent(topic, channel, event string, payload any, opts ...BroadcastOptions) error
- func (h *EventHub) Off(event string)
- func (h *EventHub) On(event string, handler Handler, middleware ...Middleware)
- func (h *EventHub) OnAny(handler Handler)
- func (h *EventHub) Remove(c *EventConn, cause error)
- func (h *EventHub) RequestTo(clientID, event string, payload any, timeout time.Duration) (Envelope, error)
- func (h *EventHub) Serve(c *EventConn) error
- func (h *EventHub) Stats() EventStats
- func (h *EventHub) Subscribe(c *EventConn, topic, channel string) error
- func (h *EventHub) Subscriptions() []Subscription
- func (h *EventHub) Unsubscribe(c *EventConn, topic, channel string) error
- func (h *EventHub) Use(m Middleware)
- type EventHubConfig
- type EventStats
- type Handler
- type HandlerContext
- type Manager
- func (m *Manager) Add(c *Conn)
- func (m *Manager) Broadcast(opcode byte, payload []byte) int
- func (m *Manager) BroadcastJSON(v any) int
- func (m *Manager) BroadcastText(s string) int
- func (m *Manager) CloseAll(code uint16, reason string, timeout time.Duration)
- func (m *Manager) Count() int
- func (m *Manager) Remove(c *Conn)
- func (m *Manager) Snapshot() []*Conn
- func (m *Manager) TryAdd(c *Conn) bool
- type MessageReader
- type MessageWriter
- type MetadataFunc
- type Middleware
- type OutboundMessage
- type PresencePayload
- type Subscription
- type Writer
Constants ¶
const ( EnvelopeHello = "hello" EnvelopeEmit = "emit" EnvelopeAck = "ack" EnvelopeError = "error" EnvelopeSubscribe = "subscribe" EnvelopeUnsubscribe = "unsubscribe" EnvelopePresence = "presence" EnvelopePing = "ping" EnvelopePong = "pong" )
const ( ActionSubscribe = "subscribe" ActionPublish = "publish" ActionEmit = "emit" ActionNotify = "notify" ActionPresence = "presence" )
const ( Continuation = byte(0x0) Text = byte(0x1) Binary = byte(0x2) Close = byte(0x8) Ping = byte(0x9) Pong = byte(0xa) )
const ( CloseNormalClosure = 1000 CloseGoingAway = 1001 CloseProtocolError = 1002 CloseUnsupportedData = 1003 CloseNoStatusReceived = 1005 CloseAbnormalClosure = 1006 CloseInvalidFramePayloadData = 1007 ClosePolicyViolation = 1008 CloseMessageTooBig = 1009 CloseMandatoryExtension = 1010 CloseInternalServerErr = 1011 CloseServiceRestart = 1012 CloseTryAgainLater = 1013 CloseTLSHandshake = 1015 )
const GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"
Variables ¶
var ( ErrEventClosed = errors.New("event websocket closed") ErrEventForbidden = errors.New("event websocket forbidden") ErrEventNotFound = errors.New("event handler not found") ErrEventBadEnvelope = errors.New("bad event envelope") ErrEventPayloadTooLarge = errors.New("event payload too large") ErrEventTooManySubs = errors.New("too many websocket subscriptions") ErrEventSlowClient = errors.New("slow websocket client") ErrEventAckTimeout = errors.New("event ack timeout") ErrEventDuplicateClientID = errors.New("duplicate websocket client id") )
var ( ErrWebSocketHandshake = errors.New("invalid websocket handshake") ErrWebSocketProtocol = errors.New("websocket protocol error") ErrWebSocketTooLarge = errors.New("websocket message too large") ErrWebSocketRateLimit = errors.New("websocket message rate limit exceeded") ErrWebSocketClosed = errors.New("websocket closed") )
Functions ¶
func IsCloseError ¶
func IsNormalClose ¶
func NewEventHandler ¶
func NewEventHandler(h *EventHub, wsCfg Config, metadata MetadataFunc) fh.HandlerFunc
NewEventHandler is the production-ready combined HTTP upgrade + event hub integration. It intentionally mirrors NewWithConfig from the low-level layer so metadata and client IDs can be derived from fh.Ctx before the context is no longer safe to read.
func NewWithConfig ¶
func NewWithConfig(cfg Config, handler func(*Conn) error) fh.HandlerFunc
Types ¶
type AuthorizeFunc ¶
AuthorizeFunc controls topic/channel access and server fanout policies.
type BroadcastOptions ¶
type ClientIDFunc ¶
ClientIDFunc can derive a stable client ID from metadata. Return empty to use a generated random ID.
type CloseError ¶
func (*CloseError) Error ¶
func (e *CloseError) Error() string
type Config ¶
type Config struct {
MaxMessageSize int
MaxFrameSize int
MaxFragments int
ReadTimeout time.Duration
WriteTimeout time.Duration
PingInterval time.Duration
PongTimeout time.Duration
MaxMessagesPerSecond int
// Secure default:
// - If CheckOrigin is nil and AllowedOrigins is empty:
// browser requests with Origin are rejected, non-browser requests without Origin are allowed.
// - Set AllowAllOrigins=true only for trusted/internal/dev use.
AllowAllOrigins bool
AllowedOrigins []string
CheckOrigin func(fh.Ctx) bool
Subprotocols []string
EnableHeartbeat bool
Manager *Manager
OnOpen func(*Conn)
OnClose func(*Conn, error)
OnError func(*Conn, error)
OnMessage func(*Conn, byte, int64)
}
func DefaultConfig ¶
func DefaultConfig() Config
type Conn ¶
type Conn struct {
// contains filtered or unexported fields
}
func (*Conn) NetConn ¶
NetConn is exposed only for advanced integration. Do not read/write directly after WebSocket upgrade. Direct reads/writes corrupt WebSocket framing.
func (*Conn) NextReader ¶
func (c *Conn) NextReader() (opcode byte, r io.ReadCloser, err error)
func (*Conn) NextWriter ¶
func (c *Conn) NextWriter(opcode byte) (io.WriteCloser, error)
func (*Conn) RemoteAddr ¶
func (*Conn) StartHeartbeat ¶
func (*Conn) Subprotocol ¶
func (*Conn) WriteBinary ¶
type Context ¶
type Context = HandlerContext
Context is a backward-compatible alias for HandlerContext. Use websocket.Context or websocket.HandlerContext in application code.
type Envelope ¶
type Envelope struct {
Type string `json:"type,omitempty"`
ID string `json:"id,omitempty"`
ReplyTo string `json:"replyTo,omitempty"`
Event string `json:"event,omitempty"`
Topic string `json:"topic,omitempty"`
Channel string `json:"channel,omitempty"`
Payload json.RawMessage `json:"payload,omitempty"`
Error string `json:"error,omitempty"`
Code string `json:"code,omitempty"`
Timestamp int64 `json:"ts,omitempty"`
}
Envelope is the JSON wire format used by both Go and TS/JS clients.
type EventConn ¶
type EventConn struct {
ID string `json:"id"`
Hub *EventHub `json:"-"`
Conn *Conn `json:"-"`
Writer *Writer `json:"-"`
Meta map[string]string `json:"meta,omitempty"`
// contains filtered or unexported fields
}
func (*EventConn) EmitScoped ¶
func (*EventConn) RequestScoped ¶
func (*EventConn) SendRawJSON ¶
func (*EventConn) Subscriptions ¶
func (c *EventConn) Subscriptions() []Subscription
type EventHub ¶
type EventHub struct {
// contains filtered or unexported fields
}
func NewEventHub ¶
func NewEventHub(cfg ...EventHubConfig) *EventHub
func (*EventHub) Accept ¶
func (h *EventHub) Accept(ws *Conn, metadata map[string]string, requestedID string) (*EventConn, error)
Accept registers an already-upgraded Conn in the hub. Use this from your own NewWithConfig callback when you need exact metadata control.
func (*EventHub) Add ¶
Add registers an already-upgraded Conn with generated client ID. It is a convenience wrapper kept for application code that does:
ec := hub.Add(conn, metadata) return hub.Serve(ec)
Prefer Accept when you want to handle duplicate IDs or registration errors explicitly.
func (*EventHub) AddWithID ¶
AddWithID registers an already-upgraded Conn using a caller-supplied client ID. It returns nil on registration failure. Prefer Accept when the error must be surfaced.
func (*EventHub) Broadcast ¶
func (h *EventHub) Broadcast(topic, channel string, env Envelope, opts ...BroadcastOptions) error
func (*EventHub) BroadcastEvent ¶
func (h *EventHub) BroadcastEvent(topic, channel, event string, payload any, opts ...BroadcastOptions) error
func (*EventHub) Handler ¶
func (h *EventHub) Handler(wsCfg Config, metadata MetadataFunc) fh.HandlerFunc
Handler returns an fh-compatible WebSocket upgrade handler. It performs the low-level RFC6455 upgrade, extracts trusted metadata from the HTTP request, registers the connection in the hub, and starts the event loop.
func (*EventHub) HandlerWithContext ¶
func (h *EventHub) HandlerWithContext(wsCfg Config, metadata MetadataFunc) fh.HandlerFunc
HandlerWithContext is kept as a clearer alias for Handler.
func (*EventHub) Notify ¶
func (h *EventHub) Notify(topic, channel string, env Envelope, opts ...BroadcastOptions) error
func (*EventHub) NotifyEvent ¶
func (h *EventHub) NotifyEvent(topic, channel, event string, payload any, opts ...BroadcastOptions) error
func (*EventHub) On ¶
func (h *EventHub) On(event string, handler Handler, middleware ...Middleware)
func (*EventHub) Stats ¶
func (h *EventHub) Stats() EventStats
func (*EventHub) Subscriptions ¶
func (h *EventHub) Subscriptions() []Subscription
func (*EventHub) Unsubscribe ¶
func (*EventHub) Use ¶
func (h *EventHub) Use(m Middleware)
type EventHubConfig ¶
type EventHubConfig struct {
WriterQueueSize int
// MaxEnvelopeBytes limits the entire incoming JSON envelope.
MaxEnvelopeBytes int
// MaxPayloadBytes limits the payload JSON field after decode.
MaxPayloadBytes int
AckTimeout time.Duration
ClientIdleTimeout time.Duration
CleanupInterval time.Duration
// MaxSubscriptions is per connection.
MaxSubscriptions int
// MaxPendingRequests is per connection for server-side request/ack tracking.
MaxPendingRequests int
// If false, client messages with Event=server.broadcast/server.notify are rejected.
AllowClientBroadcast bool
AllowClientNotify bool
// CloseSlowClient closes a connection when the writer queue is full. If false,
// the message is dropped and Send returns ErrEventSlowClient.
CloseSlowClient bool
// EnablePresence emits presence messages to topic/channel subscribers.
EnablePresence bool
// SendHello sends an initial hello envelope after upgrade.
SendHello bool
// Debug adds debug fields to metrics/log callbacks only; it does not print.
Debug bool
Auth AuthFunc
Authorize AuthorizeFunc
ClientID ClientIDFunc
Now Clock
OnError func(*EventConn, error)
OnConnect func(*EventConn)
OnDisconnect func(*EventConn, error)
OnSubscribe func(*EventConn, string, string)
OnUnsubscribe func(*EventConn, string, string)
}
func DefaultEventHubConfig ¶
func DefaultEventHubConfig() EventHubConfig
type EventStats ¶
type EventStats struct {
ConnectedClients int64 `json:"connectedClients"`
TotalConnections int64 `json:"totalConnections"`
TotalDisconnects int64 `json:"totalDisconnects"`
IncomingEnvelopes int64 `json:"incomingEnvelopes"`
OutgoingEnvelopes int64 `json:"outgoingEnvelopes"`
HandledEvents int64 `json:"handledEvents"`
HandlerErrors int64 `json:"handlerErrors"`
DroppedEnvelopes int64 `json:"droppedEnvelopes"`
Subscriptions int64 `json:"subscriptions"`
PublishedEnvelopes int64 `json:"publishedEnvelopes"`
AckTimeouts int64 `json:"ackTimeouts"`
}
type Handler ¶
type Handler func(*HandlerContext) (any, error)
type HandlerContext ¶
type HandlerContext struct {
context.Context
Hub *EventHub
Conn *EventConn
Envelope Envelope
Payload json.RawMessage
}
HandlerContext is passed to event handlers.
func (*HandlerContext) Bind ¶
func (c *HandlerContext) Bind(v any) error
func (*HandlerContext) Channel ¶
func (c *HandlerContext) Channel() string
func (*HandlerContext) Topic ¶
func (c *HandlerContext) Topic() string
type Manager ¶
type Manager struct {
MaxConnections int
// contains filtered or unexported fields
}
func NewManager ¶
func NewManager() *Manager
func (*Manager) BroadcastJSON ¶
func (*Manager) BroadcastText ¶
type MessageReader ¶
type MessageReader struct {
// contains filtered or unexported fields
}
func (*MessageReader) Close ¶
func (r *MessageReader) Close() error
type MessageWriter ¶
type MessageWriter struct {
// contains filtered or unexported fields
}
func (*MessageWriter) Close ¶
func (w *MessageWriter) Close() error
type MetadataFunc ¶
MetadataFunc extracts metadata from the HTTP upgrade request. Store only safe, trusted server-side fields here. Do not trust user controlled values unless verified in this function.
type Middleware ¶
func JSONOnlyMiddleware ¶
func JSONOnlyMiddleware() Middleware
func MaxPayloadMiddleware ¶
func MaxPayloadMiddleware(max int) Middleware
func RecoverMiddleware ¶
func RecoverMiddleware(onPanic func(any)) Middleware
Production middleware.
func RequireMeta ¶
func RequireMeta(key, expected string) Middleware
func RequireTopicPrefix ¶
func RequireTopicPrefix(prefix string) Middleware
type OutboundMessage ¶
type PresencePayload ¶
type PresencePayload struct {
Action string `json:"action"`
ClientID string `json:"clientId"`
Topic string `json:"topic,omitempty"`
Channel string `json:"channel,omitempty"`
Timestamp int64 `json:"ts"`
}
PresencePayload is emitted when presence is enabled and clients join/leave.
type Subscription ¶
type Writer ¶
type Writer struct {
Conn *Conn
Queue chan OutboundMessage
Done chan struct{}
// contains filtered or unexported fields
}