Documentation
¶
Overview ¶
pkg/protocols/rpc/client.go
Package rpc provides an implementation of a JSON-RPC server over HTTP and WebSocket. It supports HTTP keep-alive connections and efficient message handling.
The RPC server allows clients to make JSON-RPC calls over HTTP or WebSocket. It also supports subscription mechanisms over WebSocket, enabling real-time notifications.
The package includes handlers for HTTP and WebSocket connections, a connection pool, and a server that manages registered RPC methods and dispatching requests to appropriate handlers.
Example usage:
rpcServer := rpc.NewServer(logger, observability)
rpcServer.RegisterMethod("echo", EchoHandler)
// Start the RPC server
err := rpcServer.Start(context.Background())
if err != nil {
log.Fatal("Failed to start RPC server:", err)
}
pkg/protocols/rpc/handlers.go
pkg/protocols/rpc/http_handler.go
pkg/protocols/rpc/marshaller.go
pkg/protocols/rpc/middleware.go
pkg/protocols/rpc/pool.go
pkg/protocols/rpc/rpc.go
pkg/protocols/rpc/server.go
pkg/protocols/rpc/subscription.go
pkg/protocols/rpc/handlers_test.go
pkg/protocols/rpc/websocket_handler.go
Index ¶
- Constants
- func MarshalRequest(req Request) ([]byte, error)
- func RegisterTestHandlers(server *Server)
- func UnmarshalResponse(data []byte, resp *Response) error
- type ClientConnection
- type ConnectionPool
- type Error
- func AddHandler(ctx context.Context, params json.RawMessage) (interface{}, *Error)
- func EchoHandler(ctx context.Context, params json.RawMessage) (interface{}, *Error)
- func NewError(code int, message string) *Error
- func SubscribeHandler(ctx context.Context, params json.RawMessage) (any, *Error)
- func UnsubscribeHandler(ctx context.Context, params json.RawMessage) (any, *Error)
- type HTTPHandler
- type HandlerFunc
- type HandlerMethodName
- type HandlerType
- type Middleware
- type Notification
- type NotificationParams
- type RPC
- type RPCClient
- type Request
- type Response
- type Server
- func (s *Server) BroadcastEvent(eventType string, eventData interface{})
- func (s *Server) HandleRawRequest(ctx context.Context, data []byte) ([]byte, error)
- func (s *Server) HandleRequest(ctx context.Context, req Request) Response
- func (s *Server) RegisterMethod(method HandlerMethodName, handler HandlerFunc) error
- func (s *Server) SetWebSocketHandler(handler *WebSocketHandler)
- type Subscription
- type WebSocketHandler
Constants ¶
const ( RpcStateType state.StateType = "rpc" ParseError = -32700 InvalidRequest = -32600 MethodNotFound = -32601 InvalidParams = -32602 InternalError = -32603 Unauthorized = -32604 )
Error codes as defined by the JSON-RPC 2.0 specification.
Variables ¶
This section is empty.
Functions ¶
func MarshalRequest ¶
MarshalRequest marshals a Request into JSON.
func RegisterTestHandlers ¶
func RegisterTestHandlers(server *Server)
RegisterTestHandlers registers the default RPC test methods.
func UnmarshalResponse ¶
UnmarshalResponse unmarshals JSON data into a Response.
Types ¶
type ClientConnection ¶
type ClientConnection struct {
Conn gnet.Conn
Subscriptions map[string]*Subscription
Mu deadlock.Mutex // Protects Subscriptions
Ctx context.Context
Cancel context.CancelFunc
}
ClientConnection represents a client's connection and subscriptions.
type ConnectionPool ¶
type ConnectionPool struct {
// contains filtered or unexported fields
}
ConnectionPool manages a pool of reusable TCP connections.
func NewConnectionPool ¶
func NewConnectionPool(addr string, maxSize int) *ConnectionPool
NewConnectionPool creates a new ConnectionPool.
func (*ConnectionPool) Close ¶
func (cp *ConnectionPool) Close()
Close closes all connections in the pool.
func (*ConnectionPool) Get ¶
func (cp *ConnectionPool) Get() (net.Conn, error)
Get retrieves a connection from the pool or creates a new one.
func (*ConnectionPool) Put ¶
func (cp *ConnectionPool) Put(conn net.Conn)
Put returns a connection to the pool.
type Error ¶
type Error struct {
Code int `json:"code"`
Message string `json:"message"`
Data any `json:"data,omitempty"` // Optional additional information about the error
}
Error represents a JSON-RPC 2.0 error object.
func AddHandler ¶
func AddHandler(ctx context.Context, params json.RawMessage) (interface{}, *Error)
AddHandler returns the sum of two integers as float64.
func EchoHandler ¶
func EchoHandler(ctx context.Context, params json.RawMessage) (interface{}, *Error)
EchoHandler returns the received message.
func SubscribeHandler ¶
SubscribeHandler handles subscription requests.
func UnsubscribeHandler ¶
UnsubscribeHandler handles unsubscription requests.
type HTTPHandler ¶
type HTTPHandler struct {
// contains filtered or unexported fields
}
HTTPHandler handles RPC over HTTP requests.
func NewHTTPHandler ¶
func NewHTTPHandler(server *Server, logger logger.Logger) *HTTPHandler
NewHTTPHandler creates a new HTTPHandler.
func (*HTTPHandler) Handle ¶
func (h *HTTPHandler) Handle(ctx *tcp.ConnectionContext, conn gnet.Conn) gnet.Action
Handle processes incoming data and sends responses.
func (*HTTPHandler) OnClose ¶
func (h *HTTPHandler) OnClose(ctx *tcp.ConnectionContext, conn gnet.Conn)
OnClose is called when the connection is closed.
type HandlerFunc ¶
HandlerFunc defines the function signature for handling RPC methods.
func ChainMiddleware ¶
func ChainMiddleware(handler HandlerFunc, middlewares ...Middleware) HandlerFunc
ChainMiddleware applies middlewares to a handler.
type HandlerMethodName ¶
type HandlerMethodName string
func (HandlerMethodName) String ¶
func (h HandlerMethodName) String() string
type HandlerType ¶
type HandlerType string
func (HandlerType) String ¶
func (h HandlerType) String() string
type Middleware ¶
type Middleware func(HandlerFunc) HandlerFunc
Middleware defines a function to process middleware.
func BasicAuthMiddleware ¶
func BasicAuthMiddleware(validUsers map[string]string) Middleware
BasicAuthMiddleware handles HTTP Basic Authentication.
func CORSMiddleware ¶
func CORSMiddleware(allowedOrigins []string, allowedMethods []string, allowedHeaders []string) Middleware
CORSMiddleware handles CORS for HTTP transport.
func LoggingMiddleware ¶
func LoggingMiddleware(logger logger.Logger) Middleware
LoggingMiddleware logs the start and end of a request.
type Notification ¶
type Notification struct {
JSONRPC string `json:"jsonrpc"`
Method string `json:"method"`
Params NotificationParams `json:"params"`
}
type NotificationParams ¶
type RPC ¶
type RPC struct {
// contains filtered or unexported fields
}
RPC encapsulates the RPC server, transport, and connection pool.
func NewRPC ¶
func NewRPC(ctx context.Context, cfg config.Rpc, gLog logger.Logger, obs *observability.Observability, stateMgr *state.StateManager) (*RPC, error)
NewRPC initializes and returns a new RPC instance.
func SetupRPCServerForTest ¶
func (*RPC) GetConnectionPool ¶
func (rpc *RPC) GetConnectionPool() *ConnectionPool
GetConnectionPool returns the connection pool for external usage if needed.
type RPCClient ¶
type RPCClient struct {
// contains filtered or unexported fields
}
RPCClient represents a client for sending RPC requests.
func NewRPCClient ¶
func NewRPCClient(pool *ConnectionPool) *RPCClient
NewRPCClient creates a new RPCClient with the given connection pool.
type Request ¶
type Request struct {
JSONRPC string `json:"jsonrpc"` // Must be "2.0"
ID any `json:"id"` // Can be string, number, or null
Method string `json:"method"`
Params json.RawMessage `json:"params,omitempty"`
}
Request represents a JSON-RPC 2.0 request.
type Response ¶
type Response struct {
JSONRPC string `json:"jsonrpc"` // Must be "2.0"
ID any `json:"id"` // Should match the request ID
Result any `json:"result,omitempty"` // Result is mutually exclusive with Error
Error *Error `json:"error,omitempty"` // Error is mutually exclusive with Result
}
Response represents a JSON-RPC 2.0 response.
type Server ¶
type Server struct {
// contains filtered or unexported fields
}
Server manages RPC handlers and dispatching.
func NewServer ¶
func NewServer(logger logger.Logger, obs *observability.Observability) *Server
NewServer creates a new Server instance.
func (*Server) BroadcastEvent ¶
BroadcastEvent broadcasts an event to all relevant subscribers.
func (*Server) HandleRawRequest ¶
HandleRawRequest handles raw JSON-RPC request bytes and returns the response bytes.
func (*Server) HandleRequest ¶
HandleRequest processes an RPC request and returns a response.
func (*Server) RegisterMethod ¶
func (s *Server) RegisterMethod(method HandlerMethodName, handler HandlerFunc) error
RegisterMethod registers a new RPC method handler.
func (*Server) SetWebSocketHandler ¶
func (s *Server) SetWebSocketHandler(handler *WebSocketHandler)
type Subscription ¶
type Subscription struct {
ID string
Method string
Params interface{}
CreatedAt time.Time
Connection *ClientConnection
}
Subscription represents a client's subscription.
type WebSocketHandler ¶
type WebSocketHandler struct {
// contains filtered or unexported fields
}
func NewWebSocketHandler ¶
func NewWebSocketHandler(server *Server, logger logger.Logger) *WebSocketHandler
func (*WebSocketHandler) Handle ¶
func (h *WebSocketHandler) Handle(ctx *tcp.ConnectionContext, conn gnet.Conn) gnet.Action
func (*WebSocketHandler) OnClose ¶
func (h *WebSocketHandler) OnClose(ctx *tcp.ConnectionContext, conn gnet.Conn)
OnClose cleans up the client connection when the connection is closed.