Skip to content

Day 3 技术设计文档 — 客户端

Client 生命周期

Client 有四个状态,由 closingshutdown 两个布尔字段控制:

  1. ConnectedDial() 成功后,TCP 连接已建立、Option 握手完成、receive 协程已启动。此时 closing = false, shutdown = false,客户端可用。
  2. Active:正常工作状态。应用代码通过 Go()(异步)或 Call()(同步)发起 RPC 调用,每次调用会注册到 pending map 中,调用完成后由 receive 协程通知。客户端在此状态中可以持续发起多次调用。
  3. Closing:用户主动调用 client.Close(),设置 closing = true,底层连接关闭。此后新的调用会被拒绝(registerCall 返回 ErrShutdown),但已经在 pending 中的调用会等待 receive 协程处理完毕后通过 terminateCalls 统一通知错误。
  4. Shutdown:连接异常断开(如服务端崩溃、网络中断),receive 协程读取失败后设置 shutdown = true,并调用 terminateCalls 将错误通知给所有未完成的调用。

closingshutdown 的区别在于:closing 是用户主动的优雅关闭,shutdown 是被动的异常关闭。两者都会导致 IsAvailable() 返回 false,后续调用被拒绝。

send / receive 流程

发送流程(send)

当应用代码调用 client.Call(ctx, "Greeter.SayHello", args, &reply) 时:

  1. Call → GoCall 内部调用 Go() 创建一个 Call 结构体,包含方法名、参数、返回值指针和一个 Done channel,然后将其交给 send
  2. registerCall:在 sending 锁的保护下,先获取 mu 锁将 Call 注册到 pending map 中,分配一个递增的序列号 seq。此时 pending[seq] = call,表示"这个 seq 对应的调用正在等待响应"。
  3. cc.Write:将 header(含 ServiceMethodSeq)与 args 通过 Codec 编码后写入 TCP 连接。sending 锁保证多个 goroutine 并发调用时,请求的 Header + Body 不会交叉写入。
  4. 错误处理:如果写入失败,从 pending 中移除该调用,设置错误并通过 call.done() 通知调用方。

接收流程(receive)

receive 是在 newClientCodec 中启动的一个独立后台 goroutine,持续从连接中读取服务端响应:

  1. ReadHeader:从连接中读取响应头,获取 Seq(序列号)和 Error(服务端错误信息)。
  2. removeCall:通过 Seqpending map 中找到并移除对应的 Call。移除是为了防止重复处理。
  3. 分三种情况处理
    • call == nil:该调用已被移除(比如客户端因 context 超时已经取消了这个调用),此时调用 ReadBody(nil) 丢弃 Body 数据,保持流的正确位置。
    • h.Error != "":服务端返回了错误,将错误信息填入 call.Error,调用 ReadBody(nil) 丢弃 Body(错误时 Body 是无效占位值),然后 call.done() 通知调用方。
    • 正常响应:调用 ReadBody(call.Reply) 将响应数据反序列化到调用方提供的 reply 指针中,然后 call.done() 通知调用方。
  4. 循环退出:当 ReadHeader 返回错误(连接关闭/EOF),退出循环,调用 terminateCalls(err) 将错误通知给所有仍在 pending 中的调用。

send 和 receive 的并发关系

应用 goroutine 1 ── send(call1) ──┐
应用 goroutine 2 ── send(call2) ──┤ sending 锁串行写入
应用 goroutine 3 ── send(call3) ──┘        │
                                           ▼ TCP 连接(全双工)
receive goroutine ◄── 持续读取响应 ─────────┘

        ├─ h.Seq=2 → pending[2] → call2.done()
        ├─ h.Seq=1 → pending[1] → call1.done()
        └─ h.Seq=3 → pending[3] → call3.done()

关键点:发送和接收是完全独立的。TCP 是全双工的,多个 goroutine 可以同时发送请求,receive goroutine 同时在读取响应。响应不一定按请求顺序到达(服务端并行处理,先完成的先返回),Seq 序列号就是用来将乱序到达的响应匹配回正确的调用的。

pending map 并发安全

go
type Client struct {
    mu      sync.Mutex       // 保护以下字段
    seq     uint64           // 递增序列号
    pending map[uint64]*Call // seq -> *Call
    closing bool
    shutdown bool

    sending sync.Mutex // 保护 cc.Write 的串行化
}
  • mu 保护 pending map 和状态字段的读写
  • sending 保护编码和发送操作的串行化
  • 两把锁分离,避免 send 和 receive 之间不必要的竞争

Done channel 设计

Done 使用 buffered channel(容量 10):

  • 防止 receive goroutine 在 call.done() 时阻塞
  • 允许应用代码延迟消费结果
  • Call() 方法创建容量为 1 的 channel,恰好满足单次使用

Context 超时机制

go
func (client *Client) Call(ctx context.Context, ...) error {
    call := client.Go(...)
    select {
    case <-ctx.Done():
        client.removeCall(call.Seq) // 清理 pending
        return ctx.Err()
    case c := <-call.Done:
        return c.Error
    }
}

通过 select 同时监听 context 取消和调用完成,先到者优先。