package generic import ( "context" "sig-pub/pkg/zlog" "time" "github.com/bytedance/sonic" "github.com/jhump/protoreflect/v2/grpcdynamic" "github.com/jhump/protoreflect/v2/grpcreflect" "google.golang.org/grpc" "google.golang.org/protobuf/encoding/protojson" "google.golang.org/protobuf/proto" "google.golang.org/protobuf/reflect/protoreflect" "google.golang.org/protobuf/types/dynamicpb" refv1 "google.golang.org/grpc/reflection/grpc_reflection_v1" ) type GrpcGenericClient struct { serviceName string conn *grpc.ClientConn serviceDesc protoreflect.ServiceDescriptor stub *grpcdynamic.Stub } func NewGpcGenericClient(serviceName string, conn *grpc.ClientConn) *GrpcGenericClient { return &GrpcGenericClient{ serviceName: serviceName, conn: conn, } } func (c *GrpcGenericClient) InitStub(ctx context.Context) (err error) { client := grpcreflect.NewClientV1(ctx, refv1.NewServerReflectionClient(c.conn)) marketServiceSymbol, err := client.FileContainingSymbol(protoreflect.FullName(c.serviceName)) if err != nil { return } c.serviceDesc = marketServiceSymbol.Services().ByName(protoreflect.Name(c.serviceName)) c.stub = grpcdynamic.NewStub(c.conn) return } func (c *GrpcGenericClient) ServiceName() string { return c.serviceName } func (c *GrpcGenericClient) InvokeUnary(ctx context.Context, method string, reqBytes []byte, opts ...grpc.CallOption) (resp proto.Message, err error) { caller, err := c.getMethodCaller(method) if err != nil { return } request := dynamicpb.NewMessage(caller.Input()) if err = proto.Unmarshal(reqBytes, request); err != nil { return } ms := time.Now().UnixMilli() resp, err = c.stub.InvokeRpc(ctx, caller, request, opts...) // resp -> *dynamicpb.Message if delay := time.Now().UnixMilli() - ms; delay > 2000 { zlog.Warningf("api %s.%s process use %dms\n", c.serviceName, method, delay) } return } func (c *GrpcGenericClient) InvokeUnaryJson(ctx context.Context, method string, jsonBody any, opts ...grpc.CallOption) (resp proto.Message, err error) { jsonBytes, err := sonic.Marshal(jsonBody) if err != nil { return } resp, err = c.InvokeUnaryJsonBytes(ctx, method, jsonBytes, opts...) return } func (c *GrpcGenericClient) InvokeUnaryJsonBytes(ctx context.Context, method string, jsonBytes []byte, opts ...grpc.CallOption) (resp proto.Message, err error) { caller, err := c.getMethodCaller(method) if err != nil { return } request := dynamicpb.NewMessage(caller.Input()) if err = protojson.Unmarshal(jsonBytes, request); err != nil { return } resp, err = c.stub.InvokeRpc(ctx, caller, request, opts...) return } // getMethodCaller load generic resource from cache func (c *GrpcGenericClient) getMethodCaller(method string) (caller protoreflect.MethodDescriptor, err error) { caller = c.serviceDesc.Methods().ByName(protoreflect.Name(method)) if caller == nil { err = ErrorMethodNotExists return } return }