wsclient

package
v1.6.7 Latest Latest
Warning

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

Go to latest
Published: Oct 6, 2026 License: Apache-2.0 Imports: 33 Imported by: 0

Documentation

Index

Constants

View Source
const (
	// WSSpecificConfPrefix is the named sub-section of the http config options that contains websocket specific config
	WSSpecificConfPrefix = "ws"
	// WSConfigKeyWriteBufferSize is the write buffer size
	WSConfigKeyWriteBufferSize = "ws.writeBufferSize"
	// WSConfigKeyReadBufferSize is the read buffer size
	WSConfigKeyReadBufferSize = "ws.readBufferSize"
	// WSConfigKeyInitialConnectAttempts sets how many times the websocket should attempt to connect on startup, before failing (after initial connection, retry is indefinite)
	WSConfigKeyInitialConnectAttempts = "ws.initialConnectAttempts"
	// WSConfigKeyBackgroundConnect is recommended instead of initialConnectAttempts for new uses of this library, and makes initial connection and reconnection identical in behavior
	WSConfigKeyBackgroundConnect = "ws.backgroundConnect"
	// WSConfigKeyPath if set will define the path to connect to - allows sharing of the same URL between HTTP and WebSocket connection info
	WSConfigKeyPath = "ws.path"
	// WSConfigURL if set will be a completely separate URL for WebSockets (must be a ws: or wss: scheme)
	WSConfigURL = "ws.url"
	// WSConfigKeyHeartbeatInterval is the frequency of ping/pong requests, and also used for the timeout to receive a response to the heartbeat
	WSConfigKeyHeartbeatInterval = "ws.heartbeatInterval"
	// WSConnectionTimeout is the amount of time to wait while attempting to establish a connection (or automatic reconnection)
	WSConfigKeyConnectionTimeout = "ws.connectionTimeout"
	// WSConfigKeyConnectionCycleInterval when non-zero enables proactive cycling of the connection on this interval - the new connection is fully established before the old one is quiesced and closed
	WSConfigKeyConnectionCycleInterval = "ws.connectionCycleInterval"
	// WSConfigKeyConnectionCycleQuiesceTime is how long the old connection continues to deliver inbound messages after a connection cycle, before it is closed
	WSConfigKeyConnectionCycleQuiesceTime = "ws.connectionCycleQuiesceTime"
	// WSConfigDelayFactor the exponential backoff factor for delay
	WSConfigDelayFactor = "retry.factor"
)

Variables

This section is empty.

Functions

func GenerateTLSCertficates

func GenerateTLSCertficates(t *testing.T) (publicKeyFile *os.File, privateKeyFile *os.File)

GenerateTLSCertificates creates a key pair for server and client auth

func InitConfig

func InitConfig(conf config.Section)

InitConfig ensures the config is initialized for HTTP too, as WS and HTTP can share the same tree of configuration (and all the HTTP options apply to the initial upgrade)

func InitConfigWrap

func InitConfigWrap(conf config.Section)

func NewTestTLSWSServer

func NewTestTLSWSServer(testReq func(req *http.Request), publicKeyFile *os.File, privateKeyFile *os.File) (toServer, fromServer chan string, url string, done func(), err error)

NewTestTLSWSServer creates a little test server for packages (including wsclient itself) to use in unit tests and secured with mTLS by passing in a key pair

func NewTestWSServer

func NewTestWSServer(testReq func(req *http.Request)) (toServer, fromServer chan string, url string, done func())

NewTestWSServer creates a little test server for packages (including wsclient itself) to use in unit tests

func NewTestWSServerMulti added in v1.6.6

func NewTestWSServerMulti(testReq func(req *http.Request)) (connections chan *TestWSConnection, url string, rejectNext func(bool), done func())

NewTestWSServerMulti creates a test server that accepts an unlimited sequence of connections, emitting each accepted connection on the connections channel - allowing tests of reconnect and connection cycling, where more than one connection can be active at the same time

Types

type TestWSConnection added in v1.6.6

type TestWSConnection struct {
	ToServer   chan string   // messages the server received on this connection
	FromServer chan string   // push messages here to send them to the client on this connection
	Done       chan struct{} // closed when this connection has closed
	CloseConn  func()        // server-side close of this connection
}

TestWSConnection is a single server-side connection accepted by NewTestWSServerMulti

type WSClient

type WSClient interface {
	Connect() error
	Receive() <-chan []byte
	ReceiveExt() <-chan *WSPayload
	URL() string
	SetURL(url string)
	SetHeader(header, value string)
	Send(ctx context.Context, message []byte) error
	Close()
}

func New

func New(ctx context.Context, config *WSConfig, beforeConnect WSPreConnectHandler, afterConnect WSPostConnectHandler) (WSClient, error)

New creates a new outbound client that can be connected to a remote server. ** Recommend using NewWithConfig directly **

func NewWithConfig added in v1.6.6

func NewWithConfig(ctx context.Context, config *WSConfig) (WSClient, error)

NewWithConfig creates a new outbound WebSocket client with configuration, including lifecycle hooks

func Wrap

func Wrap(ctx context.Context, config WSWrapConfig, wsconn *websocket.Conn, onClose func()) WSClient

Wrap an existing connection (including an inbound server connection) with heartbeating and throttling. No reconnect functions are supported when wrapping an existing connection like this, but the supplied callback will be invoked when the connection closes (allowing cleanup/tracking).

type WSConfig

type WSConfig struct {
	HTTPURL                   string             `json:"httpUrl,omitempty"`
	WebSocketURL              string             `json:"wsUrl,omitempty"`
	WSKeyPath                 string             `json:"wsKeyPath,omitempty"`
	ReadBufferSize            int                `json:"readBufferSize,omitempty"`
	WriteBufferSize           int                `json:"writeBufferSize,omitempty"`
	InitialDelay              time.Duration      `json:"initialDelay,omitempty"`
	MaximumDelay              time.Duration      `json:"maximumDelay,omitempty"`
	DelayFactor               float64            `json:"delayFactor,omitempty"`
	BackgroundConnect         bool               `json:"backgroundConnect,omitempty"`
	InitialConnectAttempts    int                `json:"initialConnectAttempts,omitempty"` // recommend backgroundConnect instead
	DisableReconnect          bool               `json:"disableReconnect"`
	AuthUsername              string             `json:"authUsername,omitempty"`
	AuthPassword              string             `json:"authPassword,omitempty"`
	ThrottleRequestsPerSecond int                `json:"requestsPerSecond,omitempty"`
	ThrottleBurst             int                `json:"burst,omitempty"`
	HTTPHeaders               fftypes.JSONObject `json:"headers,omitempty"`
	HeartbeatInterval         time.Duration      `json:"heartbeatInterval,omitempty"`
	TLSClientConfig           *tls.Config        `json:"tlsClientConfig,omitempty"`
	ConnectionTimeout         time.Duration      `json:"connectionTimeout,omitempty"`
	// NetDialer carries the custom DNS resolver and SSRF egress guard (CIDR denylist) for the
	// underlying TCP connection. Built by GenerateConfig from the net config; cannot be set in
	// JSON. Left nil for hand-built configs, in which case the default net dialer is used.
	NetDialer *net.Dialer `json:"-"`
	// ConnectionCycleInterval when non-zero enables proactive replacement of the connection
	// on this interval - the new connection is fully established (including afterConnect)
	// before sends switch over and the old connection is quiesced and closed.
	// The interval restarts from the end of each quiesce/close, and from any reconnect
	// due to a connection error - so at most two connections ever exist concurrently.
	ConnectionCycleInterval time.Duration `json:"connectionCycleInterval,omitempty"`
	// ConnectionCycleQuiesceTime is how long the old connection continues to deliver inbound
	// messages after a connection cycle switches sends to the new connection, before it is closed
	ConnectionCycleQuiesceTime time.Duration `json:"connectionCycleQuiesceTime,omitempty"`
	// The lifecycle handlers cannot be set in JSON - they must be configured on the code interface
	PreConnectHandler    WSPreConnectHandler    `json:"-"`
	PostConnectHandler   WSPostConnectHandler   `json:"-"`
	PreDisconnectHandler WSPreDisconnectHandler `json:"-"`
	// This one cannot be set in JSON - must be configured on the code interface
	ReceiveExt bool
}

func GenerateConfig

func GenerateConfig(ctx context.Context, conf config.Section) (*WSConfig, error)

type WSPayload

type WSPayload struct {
	MessageType int
	Reader      io.Reader
	// contains filtered or unexported fields
}

WSPayload allows API consumers of this package to stream data, and inspect the message type, rather than just being passed the bytes directly.

func NewWSPayload

func NewWSPayload(mt int, r io.Reader) *WSPayload

func (*WSPayload) Processed

func (wsp *WSPayload) Processed()

Must call done on each payload, before being delivered the next

type WSPostConnectHandler

type WSPostConnectHandler func(ctx context.Context, establishing WSClient) error

WSPostConnectHandler will be called after every connect/reconnect. Can send data over ws, but must not block listening for data on the ws. During auto-cycle it is passed an establishing WSConn handle that will route to the NEW connection. You should not use this connection beyond the scope of the callback.

type WSPreConnectHandler

type WSPreConnectHandler func(ctx context.Context, w WSClient) error

WSPreConnectHandler will be called before every connect/reconnect. Any error returned will prevent the websocket from connecting.

type WSPreDisconnectHandler added in v1.6.6

type WSPreDisconnectHandler func(ctx context.Context, disconnecting WSClient) error

WSPreDisconnectHandler is called before a graceful close, to allow cleanup (such as unsubscribe):

  • When closed explicitly
  • When cycling the connection (after the new connection is established, before post-connect is called)

Passed a disconnecting WSConn handle that will route to the OLD connection. You should not use this connection beyond the scope of the callback.

type WSWrapConfig

type WSWrapConfig struct {
	HeartbeatInterval         time.Duration `json:"heartbeatInterval,omitempty"`
	ThrottleRequestsPerSecond int           `json:"requestsPerSecond,omitempty"`
	ThrottleBurst             int           `json:"burst,omitempty"`
	// This one cannot be set in JSON - must be configured on the code interface
	ReceiveExt bool
}

Jump to

Keyboard shortcuts

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