Files
open-im-server/internal/push/push_rpc_server.go
T

94 lines
2.8 KiB
Go
Raw Normal View History

2023-02-20 10:13:29 +08:00
package push
2021-05-26 19:22:11 +08:00
import (
2023-02-23 19:15:30 +08:00
"OpenIM/pkg/common/config"
"OpenIM/pkg/common/constant"
"OpenIM/pkg/common/db/cache"
"OpenIM/pkg/common/db/controller"
"OpenIM/pkg/common/log"
"OpenIM/pkg/common/prome"
pbPush "OpenIM/pkg/proto/push"
"OpenIM/pkg/utils"
2021-05-26 19:22:11 +08:00
"context"
"net"
2022-05-07 17:05:05 +08:00
"strconv"
2021-05-26 19:22:11 +08:00
"strings"
2022-09-12 19:32:24 +08:00
2022-09-15 16:27:36 +08:00
grpcPrometheus "github.com/grpc-ecosystem/go-grpc-prometheus"
2022-09-12 19:32:24 +08:00
"google.golang.org/grpc"
2021-05-26 19:22:11 +08:00
)
type RPCServer struct {
rpcPort int
rpcRegisterName string
2023-02-23 17:28:57 +08:00
pushInterface controller.PushInterface
pusher Pusher
2021-05-26 19:22:11 +08:00
}
2023-02-23 18:17:17 +08:00
func (r *RPCServer) Init(rpcPort int, cache cache.MsgCache) {
2021-05-26 19:22:11 +08:00
r.rpcPort = rpcPort
r.rpcRegisterName = config.Config.RpcRegisterName.OpenImPushName
}
2023-02-22 19:51:14 +08:00
2021-05-26 19:22:11 +08:00
func (r *RPCServer) run() {
2022-05-07 17:05:05 +08:00
listenIP := ""
if config.Config.ListenIP == "" {
listenIP = "0.0.0.0"
} else {
listenIP = config.Config.ListenIP
}
address := listenIP + ":" + strconv.Itoa(r.rpcPort)
listener, err := net.Listen("tcp", address)
2021-05-26 19:22:11 +08:00
if err != nil {
2022-05-10 09:09:37 +08:00
panic("listening err:" + err.Error() + r.rpcRegisterName)
2021-05-26 19:22:11 +08:00
}
defer listener.Close()
2022-09-15 01:22:20 +08:00
var grpcOpts []grpc.ServerOption
if config.Config.Prometheus.Enable {
2023-02-15 15:52:32 +08:00
prome.NewGrpcRequestCounter()
prome.NewGrpcRequestFailedCounter()
prome.NewGrpcRequestSuccessCounter()
2022-09-15 16:27:36 +08:00
grpcOpts = append(grpcOpts, []grpc.ServerOption{
2023-02-15 15:52:32 +08:00
// grpc.UnaryInterceptor(prome.UnaryServerInterceptorProme),
2022-09-15 16:27:36 +08:00
grpc.StreamInterceptor(grpcPrometheus.StreamServerInterceptor),
grpc.UnaryInterceptor(grpcPrometheus.UnaryServerInterceptor),
}...)
2022-09-15 01:22:20 +08:00
}
srv := grpc.NewServer(grpcOpts...)
2021-05-26 19:22:11 +08:00
defer srv.GracefulStop()
pbPush.RegisterPushMsgServiceServer(srv, r)
2022-06-23 09:24:05 +08:00
rpcRegisterIP := config.Config.RpcRegisterIP
2022-05-07 17:05:05 +08:00
if config.Config.RpcRegisterIP == "" {
rpcRegisterIP, err = utils.GetLocalIP()
if err != nil {
log.Error("", "GetLocalIP failed ", err.Error())
}
}
2023-02-07 20:24:20 +08:00
err = rpc.RegisterEtcd(r.etcdSchema, strings.Join(r.etcdAddr, ","), rpcRegisterIP, r.rpcPort, r.rpcRegisterName, 10)
2021-05-26 19:22:11 +08:00
if err != nil {
2022-05-07 17:05:05 +08:00
log.Error("", "register push module rpc to etcd err", err.Error(), r.etcdSchema, strings.Join(r.etcdAddr, ","), rpcRegisterIP, r.rpcPort, r.rpcRegisterName)
2022-08-26 17:41:58 +08:00
panic(utils.Wrap(err, "register push module rpc to etcd err"))
2021-05-26 19:22:11 +08:00
}
err = srv.Serve(listener)
if err != nil {
2022-05-07 17:05:05 +08:00
log.Error("", "push module rpc start err", err.Error())
2021-05-26 19:22:11 +08:00
return
}
}
2023-02-22 19:51:14 +08:00
func (r *RPCServer) PushMsg(ctx context.Context, pbData *pbPush.PushMsgReq) (resp *pbPush.PushMsgResp, err error) {
2022-05-30 17:46:06 +08:00
switch pbData.MsgData.SessionType {
case constant.SuperGroupChatType:
2023-02-23 17:28:57 +08:00
err = r.pusher.MsgToSuperGroupUser(ctx, pbData.SourceID, pbData.MsgData)
2022-05-30 17:46:06 +08:00
default:
2023-02-23 17:28:57 +08:00
err = r.pusher.MsgToUser(ctx, pbData.SourceID, pbData.MsgData)
2022-05-30 17:46:06 +08:00
}
2023-02-23 17:28:57 +08:00
return &pbPush.PushMsgResp{}, err
2021-05-26 19:22:11 +08:00
}
2022-10-27 17:15:06 +08:00
2023-02-22 19:51:14 +08:00
func (r *RPCServer) DelUserPushToken(ctx context.Context, req *pbPush.DelUserPushTokenReq) (resp *pbPush.DelUserPushTokenResp, err error) {
2023-02-23 17:28:57 +08:00
return &pbPush.DelUserPushTokenResp{}, r.pushInterface.DelFcmToken(ctx, req.UserID, int(req.PlatformID))
2022-10-27 17:15:06 +08:00
}