Documentation
¶
Overview ¶
Example (ErrorHandling) ¶
Example_errorHandling 演示各种错误场景的处理。
package main
import (
"fmt"
"net"
"time"
"github.com/23jdd/mrpc"
)
// Calculator 是一个示例 RPC 服务实现。
type Calculator struct{}
// MultiplyReq 乘法请求参数。
// 字段首字母必须大写(msgpack 序列化要求)。
type MultiplyReq struct {
A int
B int
}
// MultiplyReply 乘法响应结果。
type MultiplyReply struct {
Product int
}
// Multiply 乘法 RPC 方法。
//
// 签名必须满足 func(req T, reply *U) error。
// - req: 请求参数(可以是值类型或指针类型,msgpack 可反序列化即可)
// - reply: 响应参数,必须为指针(以便 server 将结果写入)
// - error: 方法执行错误,为 nil 时表示成功
//
// 边界条件:
// - req 类型错误时 msgpack 解码会失败
// - reply 为 nil 时会导致 panic
// - 方法中若 panic 未被 recover,会导致连接关闭
func (c *Calculator) Multiply(req MultiplyReq, reply *MultiplyReply) error {
reply.Product = req.A * req.B
return nil
}
func main() {
lis, _ := net.Listen("tcp", "127.0.0.1:0")
defer lis.Close()
server := mrpc.NewServer(lis)
server.Register("Calculator", &Calculator{})
go server.Run()
_, port, _ := net.SplitHostPort(lis.Addr().String())
var portNum int
fmt.Sscanf(port, "%d", &portNum)
client := mrpc.NewClient("127.0.0.1", portNum)
defer client.Close()
time.Sleep(50 * time.Millisecond)
// 场景1:调用不存在的方法 → 服务端关闭连接,客户端收到网络错误
var reply MultiplyReply
err := client.Call("Calculator.NoSuchMethod", &MultiplyReq{}, &reply)
if err != nil {
fmt.Println("unknown method:", err != nil)
}
// 场景2:调用不存在的服务 → 同上
err = client.Call("NoService.Method", &MultiplyReq{}, &reply)
if err != nil {
fmt.Println("unknown service:", err != nil)
}
// 场景3:传入 nil reply → msgpack 无法解码,通常不会有问题(nil reply 跳过解码)
err = client.Call("Calculator.Multiply", &MultiplyReq{A: 1, B: 2}, nil)
fmt.Println("nil reply, err:", err)
}
Output: unknown method: true unknown service: true nil reply, err: <nil>
Example (HealthCheck) ¶
Example_healthCheck 演示健康检查功能。
package main
import (
"fmt"
"log"
"net"
"time"
"github.com/23jdd/mrpc"
)
// Calculator 是一个示例 RPC 服务实现。
type Calculator struct{}
// MultiplyReq 乘法请求参数。
// 字段首字母必须大写(msgpack 序列化要求)。
type MultiplyReq struct {
A int
B int
}
// MultiplyReply 乘法响应结果。
type MultiplyReply struct {
Product int
}
// Multiply 乘法 RPC 方法。
//
// 签名必须满足 func(req T, reply *U) error。
// - req: 请求参数(可以是值类型或指针类型,msgpack 可反序列化即可)
// - reply: 响应参数,必须为指针(以便 server 将结果写入)
// - error: 方法执行错误,为 nil 时表示成功
//
// 边界条件:
// - req 类型错误时 msgpack 解码会失败
// - reply 为 nil 时会导致 panic
// - 方法中若 panic 未被 recover,会导致连接关闭
func (c *Calculator) Multiply(req MultiplyReq, reply *MultiplyReply) error {
reply.Product = req.A * req.B
return nil
}
func main() {
lis, _ := net.Listen("tcp", "127.0.0.1:0")
defer lis.Close()
server := mrpc.NewServer(lis)
server.Register("Calculator", &Calculator{})
// 注册健康检查服务
server.Register("Health", mrpc.RegisterHealth())
go server.Run()
_, port, _ := net.SplitHostPort(lis.Addr().String())
var portNum int
fmt.Sscanf(port, "%d", &portNum)
client := mrpc.NewClient("127.0.0.1", portNum)
defer client.Close()
time.Sleep(50 * time.Millisecond)
// 调用健康检查
var healthReply mrpc.HealthReply
err := client.Call("Health.Check", &mrpc.HealthRequest{}, &healthReply)
if err != nil {
log.Fatal(err)
}
fmt.Println("server healthy:", healthReply.Ok)
// 启动周期性健康检查(后台 goroutine)
checker := mrpc.NewHealthChecker(5*time.Second, 2*time.Second)
checker.Start(client, 3, func(err error) {
fmt.Println("health check failed:", err)
})
defer checker.Stop()
}
Output: server healthy: true
Example (ServerAndClient) ¶
Example_serverAndClient 演示从服务注册到客户端调用的完整流程。
流程:
- 创建 net.Listener 监听 TCP 端口
- 创建 Server 并注册服务实现
- 在后台 goroutine 启动服务端
- 创建 Client 发起 RPC 调用
- 解析响应并打印结果
注意事项:
- 单向非流式:每次调用发一个请求、收一个响应,连接随之关闭(server 端)
- Client 默认使用 msgpack 编解码,连接在多次 Call 之间复用
- Call 的 argv 和 reply 参数都必须为指针,否则编解码失败
package main
import (
"fmt"
"log"
"net"
"time"
"github.com/23jdd/mrpc"
)
// Calculator 是一个示例 RPC 服务实现。
type Calculator struct{}
// MultiplyReq 乘法请求参数。
// 字段首字母必须大写(msgpack 序列化要求)。
type MultiplyReq struct {
A int
B int
}
// MultiplyReply 乘法响应结果。
type MultiplyReply struct {
Product int
}
// Multiply 乘法 RPC 方法。
//
// 签名必须满足 func(req T, reply *U) error。
// - req: 请求参数(可以是值类型或指针类型,msgpack 可反序列化即可)
// - reply: 响应参数,必须为指针(以便 server 将结果写入)
// - error: 方法执行错误,为 nil 时表示成功
//
// 边界条件:
// - req 类型错误时 msgpack 解码会失败
// - reply 为 nil 时会导致 panic
// - 方法中若 panic 未被 recover,会导致连接关闭
func (c *Calculator) Multiply(req MultiplyReq, reply *MultiplyReply) error {
reply.Product = req.A * req.B
return nil
}
func main() {
// 1. 创建 TCP listener
lis, err := net.Listen("tcp", "127.0.0.1:0")
if err != nil {
log.Fatal(err)
}
defer lis.Close()
// 2. 创建 Server 并注册服务
server := mrpc.NewServer(lis)
if err := server.Register("Calculator", &Calculator{}); err != nil {
log.Fatal(err)
}
// 3. 后台启动服务端
go server.Run()
// 4. 创建客户端(连接服务端监听的端口)
_, port, _ := net.SplitHostPort(lis.Addr().String())
var portNum int
fmt.Sscanf(port, "%d", &portNum)
client := mrpc.NewClient("127.0.0.1", portNum)
defer client.Close()
// 给服务端一点时间
time.Sleep(50 * time.Millisecond)
// 5. 发起 RPC 调用
req := &MultiplyReq{A: 7, B: 8}
var reply MultiplyReply
if err := client.Call("Calculator.Multiply", req, &reply); err != nil {
log.Fatal(err)
}
fmt.Printf("%d * %d = %d\n", req.A, req.B, reply.Product)
}
Output: 7 * 8 = 56
Index ¶
- Constants
- Variables
- func PutFrame(payload []byte)
- func ReadFrame(r io.Reader) ([]byte, error)
- func RegisterHealth() *healthService
- func SendRequest(w io.Writer, req *Request) error
- func SendResponse(w io.Writer, resp *Response) error
- func WriteFrame(w io.Writer, payload []byte) error
- type Client
- type Codec
- type HealthChecker
- type HealthReply
- type HealthRequest
- type MsgCodec
- type RPCMethod
- type Request
- type Response
- type Server
- type TieredPool
Examples ¶
Constants ¶
const MaxPayloadSize = 10 << 20 // 10 MB
Variables ¶
var ( // ErrClosed 表示连接已关闭。 ErrClosed = errors.New("mrpc: connection has been closed") // ErrShutdown 表示客户端已手动关闭。 ErrShutdown = errors.New("mrpc: client is shut down") )
var DefaultPool = NewTieredPool( 128, 512, 2048, 8192, 32768, 65536, 262144, 524288, 1048576, MaxPayloadSize, )
DefaultPool 是协议层内部使用的默认分级缓冲池。 覆盖从 128B 到 10MB 的常见 RPC 负载大小,减少 ReadFrame 和 readString 的内存分配。
调用方也可以直接使用 DefaultPool 管理自己的缓冲区:
buf := mrpc.DefaultPool.Get(4096) defer mrpc.DefaultPool.Put(buf)
var ErrMaxPayload = errors.New("mrpc: payload exceeds maximum size")
Functions ¶
func PutFrame ¶
func PutFrame(payload []byte)
PutFrame 将从 ReadFrame 获取的帧 buffer 归还给 DefaultPool。 调用方在完成帧数据的解码后应调用此函数。 与 DefaultPool.Put 等价,方便配对使用:ReadFrame / PutFrame。
func RegisterHealth ¶
func RegisterHealth() *healthService
RegisterHealth 在服务端注册心跳检测方法。
用法:
server.Register("Health", mrpc.RegisterHealth())
这会注册一个 Health.Check 方法,客户端可周期性调用以检测连通性。
返回值是指向内部实现的指针,需作为 Register 的 target 参数。
Types ¶
type Client ¶
type Client struct {
// contains filtered or unexported fields
}
Client 是 RPC 客户端,支持单向非流式调用。
典型用法:
client := mrpc.NewClient("localhost", 8080)
defer client.Close()
var reply ReplyType
err := client.Call("Service.Method", &Request{...}, &reply)
Call 会自动维护 TCP 长连接并在首次调用或断线时重连。
func NewClient ¶
NewClient 创建一个 RPC 客户端。
参数:
- address: 服务端 IP 或主机名(支持 IPv4/IPv6)
- port: 服务端端口号
边界条件:
- address 为空时 net.Dial 会返回错误(在 Call 时暴露)
- port <= 0 或 >65535 时 net.Dial 会返回错误
func (*Client) Call ¶
Call 执行一次单向非流式 RPC 调用。
流程:编码 argv → 发送请求 → 接收响应 → 解码到 reply。 若尚未建立连接,会自动调用 Dial 建立 TCP 连接。 连接会在多次 Call 之间复用(长连接)。
参数:
- method: 方法名,格式 "ServiceName.MethodName"
- argv: 请求参数指针(msgpack 需要指针才能编码)
- reply: 响应参数指针,结果将解码到此
返回值:
- 成功时返回 nil
- 网络错误、编解码错误、服务端返回的错误均通过 error 返回
边界条件:
- argv 必须是指针或可序列化类型,否则 Encode 失败
- reply 必须是指针,否则 Decode 无法写入
- 并发调用 Call 不安全(共用同一连接),需要外部加锁
type Codec ¶
Codec 定义编解码器接口,负责请求/响应体的序列化与反序列化。
实现者需保证 Encode/Decode 的线程安全性(本库在服务端每个连接 使用独立 Codec 实例,客户端单连接复用同一实例)。
type HealthChecker ¶
type HealthChecker struct {
// contains filtered or unexported fields
}
HealthChecker 提供客户端到服务端的周期性健康检查能力。
通过注册 MRPC 方法 "Health.Check",客户端可以周期性调用此方法 来检测服务端是否可达。
func NewHealthChecker ¶
func NewHealthChecker(interval, timeout time.Duration) *HealthChecker
NewHealthChecker 创建一个健康检查器。
参数:
- interval: 健康检查间隔(必须 >0,建议 5s~60s)
- timeout: 每次检查的超时时间(必须 >0 且 ≤ interval)
type MsgCodec ¶
type MsgCodec struct{}
MsgCodec 是基于 msgpack 的 Codec 实现。
type RPCMethod ¶
RPCMethod 描述找到的符合 func(Req, *Reply) error 签名的方法。 ReqType 为请求参数类型,ReplyType 为响应参数类型(已保证为指针)。
type Request ¶
type Server ¶
type Server struct {
// contains filtered or unexported fields
}
Server 是一个基于 TCP 的反射型 RPC 服务端。 通过 Register 注册服务实现,Run 启动监听循环,对每个连接 在独立 goroutine 中处理单次单向非流式调用。
func (*Server) Register ¶
Register 从结构体指针中找出所有签名类似 func(request, *reply) error 的方法并注册。
参数:
- name: 服务名,调用时使用 "name.MethodName" 格式
- target: 服务实现,必须为指向 struct 的指针
边界条件:
- target 必须是 *struct 指针,否则返回错误
- 只注册满足 func(req T, reply *U) error 签名的方法
- 同名服务重复注册会覆盖之前的方法
- 空结构体或无符合方法时不会报错(直接返回 nil)
type TieredPool ¶
type TieredPool struct {
// contains filtered or unexported fields
}
TieredPool 分级缓冲池:按不同容量分桶复用 []byte,减少内存浪费与分配。 TieredPool is a collection of sync.Pools of different capacities, designed to reuse []byte slices of varying sizes while minimizing memory waste.
func NewTieredPool ¶
func NewTieredPool(capacities ...int) *TieredPool
NewTieredPool 按给定(升序)容量列表创建分级缓冲池。 capacities 必须非空、严格递增且全部大于 0,否则函数会 panic;构造函数会复制参数, 因此调用方之后修改原切片不会影响池配置。 NewTieredPool New creates a new TieredPool with the given capacities. Each capacity defines a pool of buffers with that exact capacity. The capacities slice must be sorted in ascending order.
func (*TieredPool) Get ¶
func (tp *TieredPool) Get(size int) []byte
Get 取出一个长度为 size、容量不小于 size 的缓冲(从能容纳的最小桶取)。 size=0 合法并使用最小桶;size<0 会 panic;size 超过最大桶时直接分配且不会被 Put 复用。 Get returns a []byte of length size with capacity at least size. The buffer is taken from the smallest pool whose capacity >= size. If no pool is large enough, a new buffer is allocated without pooling.
func (*TieredPool) Put ¶
func (tp *TieredPool) Put(buf []byte)
Put 归还缓冲:只有容量与某个桶完全匹配时才复用,其他缓冲直接丢弃。 精确匹配可以保证桶内缓冲始终满足该桶的容量约束。 buf 可以是 nil;归还后调用方不得再访问它。池不会清除底层字节,敏感数据应由调用方先覆盖。