rpc

package
v0.0.0-...-da72ffe Latest Latest
Warning

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

Go to latest
Published: May 16, 2025 License: MIT Imports: 30 Imported by: 0

README

RPC Package

The rpc package implements a JSON-RPC server that supports both HTTP and WebSocket protocols. It provides an efficient way to handle RPC requests, including subscription mechanisms for real-time notifications.

Architecture Overview

graph TD
    Client -->|HTTP/WebSocket| TCP_Server
    TCP_Server -->|OnTraffic| MultiplexingHandler
    MultiplexingHandler -->|Determine Protocol| HTTP_Handler
    MultiplexingHandler -->|Determine Protocol| WebSocket_Handler
    HTTP_Handler -->|Process Request| RPC_Server
    WebSocket_Handler -->|Process Message| RPC_Server
    RPC_Server -->|Handle Request| Registered_Methods
  • Client: Makes RPC requests over HTTP or WebSocket.
  • TCP_Server: Listens for incoming TCP connections.
  • MultiplexingHandler: Determines if the connection is HTTP or WebSocket and forwards it accordingly.
  • HTTP_Handler: Handles HTTP requests, supports keep-alive connections.
  • WebSocket_Handler: Handles WebSocket messages and maintains client connections.
  • RPC_Server: Manages registered methods and dispatches requests.
  • Registered_Methods: User-defined methods registered with the RPC server.
Features
  • Supports JSON-RPC over HTTP and WebSocket.
  • Efficient buffer and connection management.
  • HTTP keep-alive support for improved performance.
  • Subscription mechanism over WebSocket for real-time notifications.
  • Reusable buffer pools to reduce memory allocations.
  • Asynchronous writes in WebSocket handler to improve performance.

Getting Started

1.) Register RPC Methods:
rpcServer := rpc.NewServer(logger, observability)
rpcServer.RegisterMethod("echo", EchoHandler)
2.) Start the RPC Server:
if err := rpcServer.Start(context.Background()); err != nil {
    log.Fatal("Failed to start RPC server:", err)
}
3.) Implement Handlers:
func EchoHandler(ctx context.Context, params json.RawMessage) (interface{}, *rpc.Error) {
    var message struct {
        Message string `json:"message"`
    }
    if err := json.Unmarshal(params, &message); err != nil {
        return nil, &rpc.Error{
            Code:    rpc.InvalidParams,
            Message: "Invalid parameters",
        }
    }
    return message.Message, nil
}

Benchmark Results

The following are the benchmark results obtained from running the benchmarks:

System Specifications
  • OS: Linux (amd64)
  • CPU: AMD Ryzen Threadripper 3960X 24-Core Processor
Benchmark Results Overview
BenchmarkRPCOverHTTP

Description: Measures the performance of the RPC server handling requests over HTTP in a single-threaded manner.

Results:

  • Iterations: 5,467
  • Time per Operation: 190,065 ns/op (approximately 190 microseconds)
  • Memory Usage:
    • Bytes per Operation: 10,081 B/op
    • Allocations per Operation: 119 allocs/op

Throughput:

  • Operations per Second: ~5,262 ops/sec
BenchmarkRPCOverWebSocket

Description: Measures the performance of the RPC server handling requests over WebSocket in a single-threaded manner.

Results:

  • Iterations: 10,513
  • Time per Operation: 95,381 ns/op (approximately 95 microseconds)
  • Memory Usage:
    • Bytes per Operation: 1,820 B/op
    • Allocations per Operation: 32 allocs/op

Throughput:

  • Operations per Second: ~10,489 ops/sec
BenchmarkRPCOverHTTPParallel

Description: Simulates concurrent clients making RPC requests over HTTP.

Results:

Iterations: 71,864
Time per Operation: 14,393 ns/op (approximately 14 microseconds)
Memory Usage:
    Bytes per Operation: 9,521 B/op
    Allocations per Operation: 117 allocs/op

Throughput:

Operations per Second: ~69,486 ops/sec
BenchmarkRPCOverWebSocketConcurrent

Description: Measures the performance of the RPC server handling concurrent WebSocket connections.

Results:

Iterations: 153,610
Time per Operation: 6,969 ns/op (approximately 7 microseconds)
Memory Usage:
    Bytes per Operation: 1,978 B/op
    Allocations per Operation: 32 allocs/op

Throughput:

Operations per Second: ~143,460 ops/sec
BenchmarkSubscriptionBroadcast

Description: Measures the performance of broadcasting events to subscribers (1,000 subscribers).

Results:

Iterations: 626
Time per Operation: 1,820,635 ns/op (approximately 1.82 milliseconds)
Memory Usage:
    Bytes per Operation: 200,184 B/op
    Allocations per Operation: 5,004 allocs/op

Throughput:

Operations per Second: ~549 ops/sec
Interpretation of Results
  • RPC over WebSocket is faster than RPC over HTTP, especially in concurrent scenarios.
  • Memory allocations are significantly lower for WebSocket compared to HTTP, indicating more efficient memory usage.

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

View Source
const (
	RpcStateType state.StateType = "rpc"

	ParseError     = -32700
	InvalidRequest = -32600
	MethodNotFound = -32601
	InvalidParams  = -32602
	InternalError  = -32603

	// Custom error codes
	Unauthorized = -32604
)

Error codes as defined by the JSON-RPC 2.0 specification.

Variables

This section is empty.

Functions

func MarshalRequest

func MarshalRequest(req Request) ([]byte, error)

MarshalRequest marshals a Request into JSON.

func RegisterTestHandlers

func RegisterTestHandlers(server *Server)

RegisterTestHandlers registers the default RPC test methods.

func UnmarshalResponse

func UnmarshalResponse(data []byte, resp *Response) error

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 NewError

func NewError(code int, message string) *Error

NewError creates a new Error instance with the given code and message.

func SubscribeHandler

func SubscribeHandler(ctx context.Context, params json.RawMessage) (any, *Error)

SubscribeHandler handles subscription requests.

func UnsubscribeHandler

func UnsubscribeHandler(ctx context.Context, params json.RawMessage) (any, *Error)

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

type HandlerFunc func(ctx context.Context, params json.RawMessage) (any, *Error)

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 NotificationParams struct {
	Subscription string `json:"subscription"`
	Result       any    `json:"result"`
}

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 SetupRPCServerForTest(t testing.TB) (rpcInstance *RPC, addr string, cleanup func())

func (*RPC) Call

func (rpc *RPC) Call(ctx context.Context, req Request) (Response, error)

Call sends an RPC request using the connection pool.

func (*RPC) GetConnectionPool

func (rpc *RPC) GetConnectionPool() *ConnectionPool

GetConnectionPool returns the connection pool for external usage if needed.

func (*RPC) Server

func (rpc *RPC) Server() *Server

Server returns the underlying RPC server.

func (*RPC) Start

func (rpc *RPC) Start(ctx context.Context) error

Start launches the RPC server and transport.

func (*RPC) Stop

func (rpc *RPC) Stop() error

Stop gracefully stops the RPC server and transport.

func (*RPC) Transport

func (rpc *RPC) Transport() *tcp.Server

Transport returns the underlying transport server.

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.

func (*RPCClient) Call

func (c *RPCClient) Call(ctx context.Context, req Request) (Response, error)

Call sends an RPC request and returns the response.

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

func (s *Server) BroadcastEvent(eventType string, eventData interface{})

BroadcastEvent broadcasts an event to all relevant subscribers.

func (*Server) HandleRawRequest

func (s *Server) HandleRawRequest(ctx context.Context, data []byte) ([]byte, error)

HandleRawRequest handles raw JSON-RPC request bytes and returns the response bytes.

func (*Server) HandleRequest

func (s *Server) HandleRequest(ctx context.Context, req Request) Response

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.

Jump to

Keyboard shortcuts

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