You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Go版capnp异步客户端回调实现:注册回调跨方法失效问题

Cap'n Proto Go RPC 回调存储失效问题

我基于Go语言的capnp实现搭建RPC服务器,希望接收客户端回调函数并在后续触发调用,但遇到问题:在RegisterCallback方法内可正常调用回调,方法执行完毕后存储的callback实例失效,InvokeCallback无法触发回调。

示例代码

example_schema.capnp

using Go = import "/go.capnp";

$Go.package("rpc_example");
$Go.import("rpc_example");

interface ClientCallbackInterface {
    callbackMethod @0 () -> ();
}

interface ServerInterface {
    registerCallback @0 (callback :ClientCallbackInterface);
    invokeCallback @1 ();
}

example_server.go

package main

import (
    "capnproto.org/go/capnp/v3"
    "capnproto.org/go/capnp/v3/rpc"
    "context"
    "example-capnp/rpc_example"
    "fmt"
    "net"
)

type ServerImpl struct {
    callback rpc_example.ClientCallbackInterface
}

func (s *ServerImpl) RegisterCallback(ctx context.Context, call rpc_example.ServerInterface_registerCallback) error {
    fmt.Println("Server: Registering a callback.")

    params := call.Args()
    cb := params.Callback()
    s.callback = cb

    fmt.Println("Server: Callback registered.")

    // 此处调用正常
    _, _ = s.callback.CallbackMethod(ctx, nil)
    return nil
}

func (s *ServerImpl) InvokeCallback(_ context.Context, _ rpc_example.ServerInterface_invokeCallback) error {
    fmt.Println("Server: Invoking the callback.")

    // 此处调用失败
    _, _ = s.callback.CallbackMethod(context.Background(), nil)

    fmt.Println("Server: Callback invoked.")
    return nil
}

func main() {
    ctx := context.Background()

    listener, err := net.Listen("tcp", "127.0.0.1:2000")
    if err != nil {
        fmt.Printf("%s", err.Error())
    }

    server := ServerImpl{}
    client := rpc_example.ServerInterface_ServerToClient(&server)

    rwc, err := listener.Accept()
    if err != nil {
        fmt.Printf("%s", err.Error())
    }

    conn := rpc.NewConn(rpc.NewStreamTransport(rwc), &rpc.Options{
        BootstrapClient: capnp.Client(client),
    })

    // 阻塞直到连接终止
    select {
    case <-conn.Done():
        client.Release()
    case <-ctx.Done():
        _ = conn.Close()
    }
}

example_client.go

package main

import (
    "capnproto.org/go/capnp/v3/rpc"
    "context"
    "example-capnp/rpc_example"
    "fmt"
    "net"
    "time"
)

type Callback struct{}

func (c Callback) CallbackMethod(_ context.Context, call rpc_example.ClientCallbackInterface_callbackMethod) error {
    fmt.Println("Client: CallbackMethod has been invoked.")
    return nil
}

func main() {
    ctx := context.Background()

    rwc, err := net.Dial("tcp", "127.0.0.1:2000")
    if err != nil {
        panic(err)
    }

    conn := rpc.NewConn(rpc.NewStreamTransport(rwc), nil)
    defer conn.Close()

    d := rpc_example.ServerInterface(conn.Bootstrap(ctx))

    // 注册回调
    fmt.Println("Client: Registering a callback.")
    callback := Callback{}
    d.RegisterCallback(ctx, func(params rpc_example.ServerInterface_registerCallback_Params) error {
        cb := rpc_example.ClientCallbackInterface_ServerToClient(callback)
        return params.SetCallback(cb)
    })

    time.Sleep(time.Second * 1)

    fmt.Println("Client: Invoking callback.")
    _, release := d.InvokeCallback(ctx, nil)
    defer release()
    fmt.Println("Client: Invoked callback.")

    time.Sleep(time.Second * 1)

}

问题原因与修复方案

问题出在未保留capnp.Client的引用:生成的接口类型(如rpc_example.ClientCallbackInterface)是对capnp.Client的轻量包装,本身不持有引用。当RegisterCallback方法返回后,原始capnp.Client会被GC回收,导致存储的callback实例失效。

修复步骤

  1. 修改ServerImpl结构体,直接存储capnp.Client类型的回调引用:
type ServerImpl struct {
    callback capnp.Client
}
  1. 在RegisterCallback中,获取回调的capnp.Client并存储,转换为接口类型调用:
func (s *ServerImpl) RegisterCallback(ctx context.Context, call rpc_example.ServerInterface_registerCallback) error {
    fmt.Println("Server: Registering a callback.")

    params := call.Args()
    cb := params.Callback()
    // 保留capnp.Client引用
    s.callback = capnp.Client(cb)
    // 转换为接口类型调用
    callback := rpc_example.ClientCallbackInterface(s.callback)

    fmt.Println("Server: Callback registered.")

    _, _ = callback.CallbackMethod(ctx, nil)
    return nil
}
  1. 在InvokeCallback中,将存储的capnp.Client转换回接口类型调用:
func (s *ServerImpl) InvokeCallback(_ context.Context, _ rpc_example.ServerInterface_invokeCallback) error {
    fmt.Println("Server: Invoking the callback.")

    if s.callback.IsValid() {
        callback := rpc_example.ClientCallbackInterface(s.callback)
        _, _ = callback.CallbackMethod(context.Background(), nil)
    }

    fmt.Println("Server: Callback invoked.")
    return nil
}

原理说明

通过直接存储capnp.Client,可以确保回调的核心引用被保留,避免GC回收导致的实例失效。调用时再将其转换为生成的接口类型,即可正常触发客户端回调。

内容的提问来源于stack exchange,提问作者William Hicklin

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.06 03:28:31