websocket

package
v0.0.26 Latest Latest
Warning

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

Go to latest
Published: Sep 24, 2026 License: MIT Imports: 18 Imported by: 0

Documentation

Index

Constants

View Source
const (
	EnvelopeHello       = "hello"
	EnvelopeEmit        = "emit"
	EnvelopeAck         = "ack"
	EnvelopeError       = "error"
	EnvelopeSubscribe   = "subscribe"
	EnvelopeUnsubscribe = "unsubscribe"
	EnvelopePresence    = "presence"
	EnvelopePing        = "ping"
	EnvelopePong        = "pong"
)
View Source
const (
	ActionSubscribe = "subscribe"
	ActionPublish   = "publish"
	ActionEmit      = "emit"
	ActionNotify    = "notify"
	ActionPresence  = "presence"
)
View Source
const (
	Continuation = byte(0x0)
	Text         = byte(0x1)
	Binary       = byte(0x2)
	Close        = byte(0x8)
	Ping         = byte(0x9)
	Pong         = byte(0xa)
)
View Source
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
)
View Source
const GUID = "258EAFA5-E914-47DA-95CA-C5AB0DC85B11"

Variables

View Source
var (
	ErrEventClosed            = errors.New("event websocket closed")
	ErrEventUnauthorized      = errors.New("event websocket unauthorized")
	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")
)
View Source
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 Accept

func Accept(key []byte) string

func IsCloseError

func IsCloseError(err error, codes ...uint16) bool

func IsNormalClose

func IsNormalClose(err error) bool

func New

func New(handler func(*Conn) error) fh.HandlerFunc

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 AuthFunc

type AuthFunc func(*EventConn, Envelope) error

AuthFunc runs for every non-ack envelope. Use it for session/JWT checks.

type AuthorizeFunc

type AuthorizeFunc func(*EventConn, string, string, string) error

AuthorizeFunc controls topic/channel access and server fanout policies.

type BroadcastOptions

type BroadcastOptions struct {
	SkipClientID      string
	OnlyClientID      string
	RequireSubscribed bool
}

type ClientIDFunc

type ClientIDFunc func(fh.Ctx, map[string]string) string

ClientIDFunc can derive a stable client ID from metadata. Return empty to use a generated random ID.

type Clock

type Clock func() time.Time

Clock exists to simplify deterministic testing.

type CloseError

type CloseError struct {
	Code   uint16
	Reason string
}

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) Close

func (c *Conn) Close() error

func (*Conn) CloseWithStatus

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

func (*Conn) Closed

func (c *Conn) Closed() bool

func (*Conn) LocalAddr

func (c *Conn) LocalAddr() net.Addr

func (*Conn) NetConn

func (c *Conn) NetConn() net.Conn

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) Ping

func (c *Conn) Ping(payload []byte) error

func (*Conn) Pong

func (c *Conn) Pong(payload []byte) error

func (*Conn) ReadJSON

func (c *Conn) ReadJSON(v any) error

func (*Conn) ReadMessage

func (c *Conn) ReadMessage() (opcode byte, payload []byte, err error)

func (*Conn) RemoteAddr

func (c *Conn) RemoteAddr() net.Addr

func (*Conn) SetReadDeadline

func (c *Conn) SetReadDeadline(t time.Time) error

func (*Conn) SetWriteDeadline

func (c *Conn) SetWriteDeadline(t time.Time) error

func (*Conn) StartHeartbeat

func (c *Conn) StartHeartbeat(interval, pongTimeout time.Duration, payload []byte, done <-chan struct{})

func (*Conn) Subprotocol

func (c *Conn) Subprotocol() string

func (*Conn) WriteBinary

func (c *Conn) WriteBinary(b []byte) error

func (*Conn) WriteJSON

func (c *Conn) WriteJSON(v any) error

func (*Conn) WriteMessage

func (c *Conn) WriteMessage(opcode byte, payload []byte) error

func (*Conn) WriteText

func (c *Conn) WriteText(s string) error

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) Ack

func (c *EventConn) Ack(replyTo string, payload any) error

func (*EventConn) Close

func (c *EventConn) Close(code uint16, reason string) error

func (*EventConn) Closed

func (c *EventConn) Closed() bool

func (*EventConn) CreatedAt

func (c *EventConn) CreatedAt() time.Time

func (*EventConn) Emit

func (c *EventConn) Emit(event string, payload any) error

func (*EventConn) EmitScoped

func (c *EventConn) EmitScoped(topic, channel, event string, payload any) error

func (*EventConn) GetMeta

func (c *EventConn) GetMeta(key string) string

func (*EventConn) Join

func (c *EventConn) Join(topic, channel string) error

func (*EventConn) LastSeen

func (c *EventConn) LastSeen() time.Time

func (*EventConn) Leave

func (c *EventConn) Leave(topic, channel string) error

func (*EventConn) Request

func (c *EventConn) Request(event string, payload any, timeout time.Duration) (Envelope, error)

func (*EventConn) RequestScoped

func (c *EventConn) RequestScoped(topic, channel, event string, payload any, timeout time.Duration) (Envelope, error)

func (*EventConn) Send

func (c *EventConn) Send(env Envelope) error

func (*EventConn) SendError

func (c *EventConn) SendError(replyTo, code, msg string) error

func (*EventConn) SendRawJSON

func (c *EventConn) SendRawJSON(raw []byte) error

func (*EventConn) SetMeta

func (c *EventConn) SetMeta(key, value string)

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

func (h *EventHub) Add(ws *Conn, metadata map[string]string) *EventConn

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

func (h *EventHub) AddWithID(ws *Conn, metadata map[string]string, clientID string) *EventConn

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) Client

func (h *EventHub) Client(id string) *EventConn

func (*EventHub) Clients

func (h *EventHub) Clients() []*EventConn

func (*EventHub) Close

func (h *EventHub) Close() error

func (*EventHub) Count

func (h *EventHub) Count() int

func (*EventHub) EmitTo

func (h *EventHub) EmitTo(clientID, event string, payload any) 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) Off

func (h *EventHub) Off(event string)

func (*EventHub) On

func (h *EventHub) On(event string, handler Handler, middleware ...Middleware)

func (*EventHub) OnAny

func (h *EventHub) OnAny(handler Handler)

OnAny receives unmatched events after normal lookup fails.

func (*EventHub) Remove

func (h *EventHub) Remove(c *EventConn, cause error)

func (*EventHub) RequestTo

func (h *EventHub) RequestTo(clientID, event string, payload any, timeout time.Duration) (Envelope, error)

func (*EventHub) Serve

func (h *EventHub) Serve(c *EventConn) error

func (*EventHub) Stats

func (h *EventHub) Stats() EventStats

func (*EventHub) Subscribe

func (h *EventHub) Subscribe(c *EventConn, topic, channel string) error

func (*EventHub) Subscriptions

func (h *EventHub) Subscriptions() []Subscription

func (*EventHub) Unsubscribe

func (h *EventHub) Unsubscribe(c *EventConn, topic, channel string) error

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) Add

func (m *Manager) Add(c *Conn)

func (*Manager) Broadcast

func (m *Manager) Broadcast(opcode byte, payload []byte) int

func (*Manager) BroadcastJSON

func (m *Manager) BroadcastJSON(v any) int

func (*Manager) BroadcastText

func (m *Manager) BroadcastText(s string) int

func (*Manager) CloseAll

func (m *Manager) CloseAll(code uint16, reason string, timeout time.Duration)

func (*Manager) Count

func (m *Manager) Count() int

func (*Manager) Remove

func (m *Manager) Remove(c *Conn)

func (*Manager) Snapshot

func (m *Manager) Snapshot() []*Conn

func (*Manager) TryAdd

func (m *Manager) TryAdd(c *Conn) bool

TryAdd atomically enforces MaxConnections and registers c. It returns false when the manager is full or either argument is nil.

type MessageReader

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

func (*MessageReader) Close

func (r *MessageReader) Close() error

func (*MessageReader) Read

func (r *MessageReader) Read(p []byte) (int, error)

type MessageWriter

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

func (*MessageWriter) Close

func (w *MessageWriter) Close() error

func (*MessageWriter) Write

func (w *MessageWriter) Write(p []byte) (int, error)

type MetadataFunc

type MetadataFunc func(fh.Ctx) map[string]string

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

type Middleware func(Handler) Handler

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 OutboundMessage struct {
	Opcode  byte
	Payload []byte
}

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 Subscription struct {
	Topic     string    `json:"topic"`
	Channel   string    `json:"channel"`
	CreatedAt time.Time `json:"createdAt"`
}

type Writer

type Writer struct {
	Conn  *Conn
	Queue chan OutboundMessage
	Done  chan struct{}
	// contains filtered or unexported fields
}

func NewWriter

func NewWriter(conn *Conn, queueSize int) *Writer

func (*Writer) CloseWithStatus

func (w *Writer) CloseWithStatus(code uint16, reason string) error

func (*Writer) Send

func (w *Writer) Send(opcode byte, payload []byte) bool

func (*Writer) SendJSON

func (w *Writer) SendJSON(v any) bool

func (*Writer) SendText

func (w *Writer) SendText(s string) bool

Jump to

Keyboard shortcuts

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