You can not select more than 25 topics
Topics must start with a letter or number, can include dashes ('-') and can be up to 35 characters long.
429 lines
16 KiB
429 lines
16 KiB
2 years ago
|
package msggateway
|
||
4 years ago
|
|
||
|
import (
|
||
3 years ago
|
"Open_IM/pkg/common/config"
|
||
|
"Open_IM/pkg/common/constant"
|
||
2 years ago
|
"Open_IM/pkg/common/db"
|
||
3 years ago
|
"Open_IM/pkg/common/log"
|
||
2 years ago
|
promePkg "Open_IM/pkg/common/prometheus"
|
||
2 years ago
|
pbChat "Open_IM/pkg/proto/msg"
|
||
2 years ago
|
push "Open_IM/pkg/proto/push"
|
||
3 years ago
|
pbRtc "Open_IM/pkg/proto/rtc"
|
||
2 years ago
|
sdkws "Open_IM/pkg/proto/sdkws"
|
||
3 years ago
|
"Open_IM/pkg/utils"
|
||
3 years ago
|
"bytes"
|
||
4 years ago
|
"context"
|
||
3 years ago
|
"encoding/gob"
|
||
|
"runtime"
|
||
4 years ago
|
"strings"
|
||
2 years ago
|
|
||
|
"github.com/golang/protobuf/proto"
|
||
|
"github.com/gorilla/websocket"
|
||
|
"google.golang.org/grpc"
|
||
4 years ago
|
)
|
||
|
|
||
3 years ago
|
func (ws *WServer) msgParse(conn *UserConn, binaryMsg []byte) {
|
||
|
b := bytes.NewBuffer(binaryMsg)
|
||
4 years ago
|
m := Req{}
|
||
3 years ago
|
dec := gob.NewDecoder(b)
|
||
|
err := dec.Decode(&m)
|
||
|
if err != nil {
|
||
3 years ago
|
log.NewError("", "ws Decode err", err.Error())
|
||
3 years ago
|
err = conn.Close()
|
||
|
if err != nil {
|
||
|
log.NewError("", "ws close err", err.Error())
|
||
|
}
|
||
4 years ago
|
return
|
||
|
}
|
||
|
if err := validate.Struct(m); err != nil {
|
||
3 years ago
|
log.NewError("", "ws args validate err", err.Error())
|
||
3 years ago
|
ws.sendErrMsg(conn, 201, err.Error(), m.ReqIdentifier, m.MsgIncr, m.OperationID)
|
||
4 years ago
|
return
|
||
|
}
|
||
2 years ago
|
log.NewInfo(m.OperationID, "Basic Info Authentication Success", m.SendID, m.MsgIncr, m.ReqIdentifier)
|
||
2 years ago
|
if m.SendID != conn.userID {
|
||
|
if err = conn.Close(); err != nil {
|
||
|
log.NewError(m.OperationID, "close ws conn failed", conn.userID, "send id", m.SendID, err.Error())
|
||
|
return
|
||
|
}
|
||
|
}
|
||
4 years ago
|
switch m.ReqIdentifier {
|
||
|
case constant.WSGetNewestSeq:
|
||
2 years ago
|
log.NewInfo(m.OperationID, "getSeqReq ", m.SendID, m.MsgIncr, m.ReqIdentifier)
|
||
3 years ago
|
ws.getSeqReq(conn, &m)
|
||
2 years ago
|
promePkg.PromeInc(promePkg.GetNewestSeqTotalCounter)
|
||
4 years ago
|
case constant.WSSendMsg:
|
||
2 years ago
|
log.NewInfo(m.OperationID, "sendMsgReq ", m.SendID, m.MsgIncr, m.ReqIdentifier)
|
||
3 years ago
|
ws.sendMsgReq(conn, &m)
|
||
2 years ago
|
promePkg.PromeInc(promePkg.MsgRecvTotalCounter)
|
||
3 years ago
|
case constant.WSSendSignalMsg:
|
||
2 years ago
|
log.NewInfo(m.OperationID, "sendSignalMsgReq ", m.SendID, m.MsgIncr, m.ReqIdentifier)
|
||
3 years ago
|
ws.sendSignalMsgReq(conn, &m)
|
||
3 years ago
|
case constant.WSPullMsgBySeqList:
|
||
2 years ago
|
log.NewInfo(m.OperationID, "pullMsgBySeqListReq ", m.SendID, m.MsgIncr, m.ReqIdentifier)
|
||
3 years ago
|
ws.pullMsgBySeqListReq(conn, &m)
|
||
2 years ago
|
promePkg.PromeInc(promePkg.PullMsgBySeqListTotalCounter)
|
||
2 years ago
|
case constant.WsLogoutMsg:
|
||
|
log.NewInfo(m.OperationID, "conn.Close()", m.SendID, m.MsgIncr, m.ReqIdentifier)
|
||
2 years ago
|
ws.userLogoutReq(conn, &m)
|
||
2 years ago
|
case constant.WsSetBackgroundStatus:
|
||
|
log.NewInfo(m.OperationID, "WsSetBackgroundStatus", m.SendID, m.MsgIncr, m.ReqIdentifier)
|
||
|
ws.setUserDeviceBackground(conn, &m)
|
||
4 years ago
|
default:
|
||
2 years ago
|
log.Error(m.OperationID, "ReqIdentifier failed ", m.SendID, m.MsgIncr, m.ReqIdentifier)
|
||
4 years ago
|
}
|
||
3 years ago
|
log.NewInfo(m.OperationID, "goroutine num is ", runtime.NumGoroutine())
|
||
4 years ago
|
}
|
||
2 years ago
|
|
||
3 years ago
|
func (ws *WServer) getSeqReq(conn *UserConn, m *Req) {
|
||
2 years ago
|
log.NewInfo(m.OperationID, "Ws call success to getNewSeq", m.MsgIncr, m.SendID, m.ReqIdentifier)
|
||
2 years ago
|
nReply := new(sdkws.GetMaxAndMinSeqResp)
|
||
2 years ago
|
isPass, errCode, errMsg, data := ws.argsValidate(m, constant.WSGetNewestSeq, m.OperationID)
|
||
|
log.Info(m.OperationID, "argsValidate ", isPass, errCode, errMsg)
|
||
3 years ago
|
if isPass {
|
||
2 years ago
|
rpcReq := sdkws.GetMaxAndMinSeqReq{}
|
||
|
rpcReq.GroupIDList = data.(sdkws.GetMaxAndMinSeqReq).GroupIDList
|
||
3 years ago
|
rpcReq.UserID = m.SendID
|
||
|
rpcReq.OperationID = m.OperationID
|
||
2 years ago
|
log.Debug(m.OperationID, "Ws call success to getMaxAndMinSeq", m.SendID, m.ReqIdentifier, m.MsgIncr, data.(sdkws.GetMaxAndMinSeqReq).GroupIDList)
|
||
2 years ago
|
grpcConn := rpc.GetDefaultConn(config.Config.Etcd.EtcdSchema, strings.Join(config.Config.Etcd.EtcdAddr, ","), config.Config.RpcRegisterName.OpenImMsgName, rpcReq.OperationID)
|
||
2 years ago
|
if grpcConn == nil {
|
||
2 years ago
|
errMsg := rpcReq.OperationID + "getcdv3.GetDefaultConn == nil"
|
||
2 years ago
|
nReply.ErrCode = 500
|
||
|
nReply.ErrMsg = errMsg
|
||
|
log.NewError(rpcReq.OperationID, errMsg)
|
||
|
ws.getSeqResp(conn, m, nReply)
|
||
|
return
|
||
|
}
|
||
2 years ago
|
msgClient := pbChat.NewMsgClient(grpcConn)
|
||
3 years ago
|
rpcReply, err := msgClient.GetMaxAndMinSeq(context.Background(), &rpcReq)
|
||
|
if err != nil {
|
||
|
nReply.ErrCode = 500
|
||
|
nReply.ErrMsg = err.Error()
|
||
2 years ago
|
log.Error(rpcReq.OperationID, "rpc call failed to GetMaxAndMinSeq ", nReply.String())
|
||
3 years ago
|
ws.getSeqResp(conn, m, nReply)
|
||
|
} else {
|
||
|
log.NewInfo(rpcReq.OperationID, "rpc call success to getSeqReq", rpcReply.String())
|
||
|
ws.getSeqResp(conn, m, rpcReply)
|
||
|
}
|
||
4 years ago
|
} else {
|
||
3 years ago
|
nReply.ErrCode = errCode
|
||
|
nReply.ErrMsg = errMsg
|
||
2 years ago
|
log.Error(m.OperationID, "argsValidate failed send resp: ", nReply.String())
|
||
3 years ago
|
ws.getSeqResp(conn, m, nReply)
|
||
4 years ago
|
}
|
||
3 years ago
|
}
|
||
2 years ago
|
|
||
2 years ago
|
func (ws *WServer) getSeqResp(conn *UserConn, m *Req, pb *sdkws.GetMaxAndMinSeqResp) {
|
||
2 years ago
|
|
||
3 years ago
|
b, _ := proto.Marshal(pb)
|
||
3 years ago
|
mReply := Resp{
|
||
|
ReqIdentifier: m.ReqIdentifier,
|
||
|
MsgIncr: m.MsgIncr,
|
||
|
ErrCode: pb.GetErrCode(),
|
||
|
ErrMsg: pb.GetErrMsg(),
|
||
|
OperationID: m.OperationID,
|
||
|
Data: b,
|
||
4 years ago
|
}
|
||
2 years ago
|
log.Debug(m.OperationID, "getSeqResp come here req: ", pb.String(), "send resp: ",
|
||
|
mReply.ReqIdentifier, mReply.MsgIncr, mReply.ErrCode, mReply.ErrMsg)
|
||
4 years ago
|
ws.sendMsg(conn, mReply)
|
||
|
}
|
||
|
|
||
3 years ago
|
func (ws *WServer) pullMsgBySeqListReq(conn *UserConn, m *Req) {
|
||
2 years ago
|
log.NewInfo(m.OperationID, "Ws call success to pullMsgBySeqListReq start", m.SendID, m.ReqIdentifier, m.MsgIncr, string(m.Data))
|
||
2 years ago
|
nReply := new(sdkws.PullMessageBySeqListResp)
|
||
2 years ago
|
isPass, errCode, errMsg, data := ws.argsValidate(m, constant.WSPullMsgBySeqList, m.OperationID)
|
||
3 years ago
|
if isPass {
|
||
2 years ago
|
rpcReq := sdkws.PullMessageBySeqListReq{}
|
||
|
rpcReq.SeqList = data.(sdkws.PullMessageBySeqListReq).SeqList
|
||
3 years ago
|
rpcReq.UserID = m.SendID
|
||
|
rpcReq.OperationID = m.OperationID
|
||
2 years ago
|
rpcReq.GroupSeqList = data.(sdkws.PullMessageBySeqListReq).GroupSeqList
|
||
|
log.NewInfo(m.OperationID, "Ws call success to pullMsgBySeqListReq middle", m.SendID, m.ReqIdentifier, m.MsgIncr, data.(sdkws.PullMessageBySeqListReq).SeqList)
|
||
2 years ago
|
grpcConn := rpc.GetDefaultConn(config.Config.Etcd.EtcdSchema, strings.Join(config.Config.Etcd.EtcdAddr, ","), config.Config.RpcRegisterName.OpenImMsgName, m.OperationID)
|
||
2 years ago
|
if grpcConn == nil {
|
||
2 years ago
|
errMsg := rpcReq.OperationID + "getcdv3.GetDefaultConn == nil"
|
||
2 years ago
|
nReply.ErrCode = 500
|
||
|
nReply.ErrMsg = errMsg
|
||
|
log.NewError(rpcReq.OperationID, errMsg)
|
||
|
ws.pullMsgBySeqListResp(conn, m, nReply)
|
||
|
return
|
||
|
}
|
||
2 years ago
|
msgClient := pbChat.NewMsgClient(grpcConn)
|
||
2 years ago
|
maxSizeOption := grpc.MaxCallRecvMsgSize(1024 * 1024 * 20)
|
||
|
reply, err := msgClient.PullMessageBySeqList(context.Background(), &rpcReq, maxSizeOption)
|
||
3 years ago
|
if err != nil {
|
||
3 years ago
|
log.NewError(rpcReq.OperationID, "pullMsgBySeqListReq err", err.Error())
|
||
3 years ago
|
nReply.ErrCode = 200
|
||
|
nReply.ErrMsg = err.Error()
|
||
3 years ago
|
ws.pullMsgBySeqListResp(conn, m, nReply)
|
||
3 years ago
|
} else {
|
||
3 years ago
|
log.NewInfo(rpcReq.OperationID, "rpc call success to pullMsgBySeqListReq", reply.String(), len(reply.List))
|
||
3 years ago
|
ws.pullMsgBySeqListResp(conn, m, reply)
|
||
3 years ago
|
}
|
||
|
} else {
|
||
|
nReply.ErrCode = errCode
|
||
|
nReply.ErrMsg = errMsg
|
||
3 years ago
|
ws.pullMsgBySeqListResp(conn, m, nReply)
|
||
3 years ago
|
}
|
||
|
}
|
||
2 years ago
|
func (ws *WServer) pullMsgBySeqListResp(conn *UserConn, m *Req, pb *sdkws.PullMessageBySeqListResp) {
|
||
3 years ago
|
log.NewInfo(m.OperationID, "pullMsgBySeqListResp come here ", pb.String())
|
||
|
c, _ := proto.Marshal(pb)
|
||
3 years ago
|
mReply := Resp{
|
||
|
ReqIdentifier: m.ReqIdentifier,
|
||
|
MsgIncr: m.MsgIncr,
|
||
|
ErrCode: pb.GetErrCode(),
|
||
|
ErrMsg: pb.GetErrMsg(),
|
||
|
OperationID: m.OperationID,
|
||
|
Data: c,
|
||
|
}
|
||
3 years ago
|
log.NewInfo(m.OperationID, "pullMsgBySeqListResp all data is ", mReply.ReqIdentifier, mReply.MsgIncr, mReply.ErrCode, mReply.ErrMsg,
|
||
3 years ago
|
len(mReply.Data))
|
||
|
ws.sendMsg(conn, mReply)
|
||
|
}
|
||
2 years ago
|
func (ws *WServer) userLogoutReq(conn *UserConn, m *Req) {
|
||
|
log.NewInfo(m.OperationID, "Ws call success to userLogoutReq start", m.SendID, m.ReqIdentifier, m.MsgIncr, string(m.Data))
|
||
2 years ago
|
|
||
2 years ago
|
rpcReq := push.DelUserPushTokenReq{}
|
||
|
rpcReq.UserID = m.SendID
|
||
2 years ago
|
rpcReq.PlatformID = conn.PlatformID
|
||
2 years ago
|
rpcReq.OperationID = m.OperationID
|
||
2 years ago
|
grpcConn := rpc.GetDefaultConn(config.Config.Etcd.EtcdSchema, strings.Join(config.Config.Etcd.EtcdAddr, ","), config.Config.RpcRegisterName.OpenImPushName, m.OperationID)
|
||
2 years ago
|
if grpcConn == nil {
|
||
|
errMsg := rpcReq.OperationID + "getcdv3.GetDefaultConn == nil"
|
||
|
log.NewError(rpcReq.OperationID, errMsg)
|
||
|
ws.userLogoutResp(conn, m)
|
||
|
return
|
||
|
}
|
||
|
msgClient := push.NewPushMsgServiceClient(grpcConn)
|
||
|
reply, err := msgClient.DelUserPushToken(context.Background(), &rpcReq)
|
||
|
if err != nil {
|
||
|
log.NewError(rpcReq.OperationID, "DelUserPushToken err", err.Error())
|
||
2 years ago
|
|
||
2 years ago
|
ws.userLogoutResp(conn, m)
|
||
|
} else {
|
||
|
log.NewInfo(rpcReq.OperationID, "rpc call success to DelUserPushToken", reply.String())
|
||
|
ws.userLogoutResp(conn, m)
|
||
|
}
|
||
|
ws.userLogoutResp(conn, m)
|
||
|
|
||
|
}
|
||
|
func (ws *WServer) userLogoutResp(conn *UserConn, m *Req) {
|
||
|
mReply := Resp{
|
||
|
ReqIdentifier: m.ReqIdentifier,
|
||
|
MsgIncr: m.MsgIncr,
|
||
|
OperationID: m.OperationID,
|
||
|
}
|
||
|
ws.sendMsg(conn, mReply)
|
||
|
_ = conn.Close()
|
||
|
}
|
||
3 years ago
|
func (ws *WServer) sendMsgReq(conn *UserConn, m *Req) {
|
||
3 years ago
|
sendMsgAllCountLock.Lock()
|
||
3 years ago
|
sendMsgAllCount++
|
||
3 years ago
|
sendMsgAllCountLock.Unlock()
|
||
2 years ago
|
log.NewInfo(m.OperationID, "Ws call success to sendMsgReq start", m.MsgIncr, m.ReqIdentifier, m.SendID)
|
||
3 years ago
|
|
||
3 years ago
|
nReply := new(pbChat.SendMsgResp)
|
||
2 years ago
|
isPass, errCode, errMsg, pData := ws.argsValidate(m, constant.WSSendMsg, m.OperationID)
|
||
4 years ago
|
if isPass {
|
||
2 years ago
|
data := pData.(sdkws.MsgData)
|
||
3 years ago
|
pbData := pbChat.SendMsgReq{
|
||
|
Token: m.Token,
|
||
|
OperationID: m.OperationID,
|
||
|
MsgData: &data,
|
||
4 years ago
|
}
|
||
2 years ago
|
log.NewInfo(m.OperationID, "Ws call success to sendMsgReq middle", m.ReqIdentifier, m.SendID, m.MsgIncr, data.String())
|
||
2 years ago
|
etcdConn := rpc.GetDefaultConn(config.Config.Etcd.EtcdSchema, strings.Join(config.Config.Etcd.EtcdAddr, ","), config.Config.RpcRegisterName.OpenImMsgName, m.OperationID)
|
||
2 years ago
|
if etcdConn == nil {
|
||
2 years ago
|
errMsg := m.OperationID + "getcdv3.GetDefaultConn == nil"
|
||
2 years ago
|
nReply.ErrCode = 500
|
||
|
nReply.ErrMsg = errMsg
|
||
|
log.NewError(m.OperationID, errMsg)
|
||
|
ws.sendMsgResp(conn, m, nReply)
|
||
|
return
|
||
|
}
|
||
2 years ago
|
client := pbChat.NewMsgClient(etcdConn)
|
||
3 years ago
|
reply, err := client.SendMsg(context.Background(), &pbData)
|
||
3 years ago
|
if err != nil {
|
||
|
log.NewError(pbData.OperationID, "UserSendMsg err", err.Error())
|
||
|
nReply.ErrCode = 200
|
||
|
nReply.ErrMsg = err.Error()
|
||
3 years ago
|
ws.sendMsgResp(conn, m, nReply)
|
||
3 years ago
|
} else {
|
||
|
log.NewInfo(pbData.OperationID, "rpc call success to sendMsgReq", reply.String())
|
||
3 years ago
|
ws.sendMsgResp(conn, m, reply)
|
||
3 years ago
|
}
|
||
|
|
||
4 years ago
|
} else {
|
||
3 years ago
|
nReply.ErrCode = errCode
|
||
|
nReply.ErrMsg = errMsg
|
||
3 years ago
|
ws.sendMsgResp(conn, m, nReply)
|
||
4 years ago
|
}
|
||
|
|
||
|
}
|
||
3 years ago
|
func (ws *WServer) sendMsgResp(conn *UserConn, m *Req, pb *pbChat.SendMsgResp) {
|
||
2 years ago
|
var mReplyData sdkws.UserSendMsgResp
|
||
3 years ago
|
mReplyData.ClientMsgID = pb.GetClientMsgID()
|
||
|
mReplyData.ServerMsgID = pb.GetServerMsgID()
|
||
3 years ago
|
mReplyData.SendTime = pb.GetSendTime()
|
||
3 years ago
|
b, _ := proto.Marshal(&mReplyData)
|
||
|
mReply := Resp{
|
||
|
ReqIdentifier: m.ReqIdentifier,
|
||
|
MsgIncr: m.MsgIncr,
|
||
|
ErrCode: pb.GetErrCode(),
|
||
|
ErrMsg: pb.GetErrMsg(),
|
||
|
OperationID: m.OperationID,
|
||
|
Data: b,
|
||
|
}
|
||
|
ws.sendMsg(conn, mReply)
|
||
2 years ago
|
|
||
3 years ago
|
}
|
||
4 years ago
|
|
||
3 years ago
|
func (ws *WServer) sendSignalMsgReq(conn *UserConn, m *Req) {
|
||
2 years ago
|
log.NewInfo(m.OperationID, "Ws call success to sendSignalMsgReq start", m.MsgIncr, m.ReqIdentifier, m.SendID, string(m.Data))
|
||
3 years ago
|
nReply := new(pbChat.SendMsgResp)
|
||
2 years ago
|
isPass, errCode, errMsg, pData := ws.argsValidate(m, constant.WSSendSignalMsg, m.OperationID)
|
||
3 years ago
|
if isPass {
|
||
3 years ago
|
signalResp := pbRtc.SignalResp{}
|
||
2 years ago
|
etcdConn := rpc.GetDefaultConn(config.Config.Etcd.EtcdSchema, strings.Join(config.Config.Etcd.EtcdAddr, ","), config.Config.RpcRegisterName.OpenImRealTimeCommName, m.OperationID)
|
||
2 years ago
|
if etcdConn == nil {
|
||
2 years ago
|
errMsg := m.OperationID + "getcdv3.GetDefaultConn == nil"
|
||
2 years ago
|
log.NewError(m.OperationID, errMsg)
|
||
|
ws.sendSignalMsgResp(conn, 204, errMsg, m, &signalResp)
|
||
|
return
|
||
|
}
|
||
3 years ago
|
rtcClient := pbRtc.NewRtcServiceClient(etcdConn)
|
||
3 years ago
|
req := &pbRtc.SignalMessageAssembleReq{
|
||
3 years ago
|
SignalReq: pData.(*pbRtc.SignalReq),
|
||
|
OperationID: m.OperationID,
|
||
3 years ago
|
}
|
||
|
respPb, err := rtcClient.SignalMessageAssemble(context.Background(), req)
|
||
|
if err != nil {
|
||
3 years ago
|
log.NewError(m.OperationID, utils.GetSelfFuncName(), "SignalMessageAssemble", err.Error(), config.Config.RpcRegisterName.OpenImRealTimeCommName)
|
||
3 years ago
|
ws.sendSignalMsgResp(conn, 204, "grpc SignalMessageAssemble failed: "+err.Error(), m, &signalResp)
|
||
|
return
|
||
|
}
|
||
|
signalResp.Payload = respPb.SignalResp.Payload
|
||
2 years ago
|
msgData := sdkws.MsgData{}
|
||
3 years ago
|
utils.CopyStructFields(&msgData, respPb.MsgData)
|
||
3 years ago
|
log.NewInfo(m.OperationID, utils.GetSelfFuncName(), respPb.String())
|
||
3 years ago
|
if respPb.IsPass {
|
||
3 years ago
|
pbData := pbChat.SendMsgReq{
|
||
|
Token: m.Token,
|
||
|
OperationID: m.OperationID,
|
||
3 years ago
|
MsgData: &msgData,
|
||
3 years ago
|
}
|
||
3 years ago
|
log.NewInfo(m.OperationID, utils.GetSelfFuncName(), "pbData: ", pbData)
|
||
3 years ago
|
log.NewInfo(m.OperationID, "Ws call success to sendSignalMsgReq middle", m.ReqIdentifier, m.SendID, m.MsgIncr, msgData)
|
||
2 years ago
|
etcdConn := rpc.GetDefaultConn(config.Config.Etcd.EtcdSchema, strings.Join(config.Config.Etcd.EtcdAddr, ","), config.Config.RpcRegisterName.OpenImMsgName, m.OperationID)
|
||
2 years ago
|
if etcdConn == nil {
|
||
2 years ago
|
errMsg := m.OperationID + "getcdv3.GetDefaultConn == nil"
|
||
2 years ago
|
log.NewError(m.OperationID, errMsg)
|
||
|
ws.sendSignalMsgResp(conn, 200, errMsg, m, &signalResp)
|
||
|
return
|
||
|
}
|
||
2 years ago
|
client := pbChat.NewMsgClient(etcdConn)
|
||
3 years ago
|
reply, err := client.SendMsg(context.Background(), &pbData)
|
||
|
if err != nil {
|
||
3 years ago
|
log.NewError(pbData.OperationID, utils.GetSelfFuncName(), "rpc sendMsg err", err.Error())
|
||
3 years ago
|
nReply.ErrCode = 200
|
||
|
nReply.ErrMsg = err.Error()
|
||
3 years ago
|
ws.sendSignalMsgResp(conn, 200, err.Error(), m, &signalResp)
|
||
3 years ago
|
} else {
|
||
2 years ago
|
log.NewInfo(pbData.OperationID, "rpc call success to sendMsgReq", reply.String(), signalResp.String(), m)
|
||
2 years ago
|
ws.sendSignalMsgResp(conn, 0, "", m, &signalResp)
|
||
3 years ago
|
}
|
||
3 years ago
|
} else {
|
||
3 years ago
|
log.NewError(m.OperationID, utils.GetSelfFuncName(), respPb.IsPass, respPb.CommonResp.ErrCode, respPb.CommonResp.ErrMsg)
|
||
|
ws.sendSignalMsgResp(conn, respPb.CommonResp.ErrCode, respPb.CommonResp.ErrMsg, m, &signalResp)
|
||
3 years ago
|
}
|
||
3 years ago
|
} else {
|
||
|
ws.sendSignalMsgResp(conn, errCode, errMsg, m, nil)
|
||
3 years ago
|
}
|
||
|
|
||
|
}
|
||
3 years ago
|
func (ws *WServer) sendSignalMsgResp(conn *UserConn, errCode int32, errMsg string, m *Req, pb *pbRtc.SignalResp) {
|
||
3 years ago
|
// := make(map[string]interface{})
|
||
2 years ago
|
log.Debug(m.OperationID, "sendSignalMsgResp is", pb.String())
|
||
3 years ago
|
b, _ := proto.Marshal(pb)
|
||
3 years ago
|
mReply := Resp{
|
||
|
ReqIdentifier: m.ReqIdentifier,
|
||
|
MsgIncr: m.MsgIncr,
|
||
3 years ago
|
ErrCode: errCode,
|
||
|
ErrMsg: errMsg,
|
||
3 years ago
|
OperationID: m.OperationID,
|
||
|
Data: b,
|
||
|
}
|
||
|
ws.sendMsg(conn, mReply)
|
||
|
}
|
||
3 years ago
|
func (ws *WServer) sendMsg(conn *UserConn, mReply interface{}) {
|
||
|
var b bytes.Buffer
|
||
|
enc := gob.NewEncoder(&b)
|
||
|
err := enc.Encode(mReply)
|
||
|
if err != nil {
|
||
2 years ago
|
// uid, platform := ws.getUserUid(conn)
|
||
|
log.NewError(mReply.(Resp).OperationID, mReply.(Resp).ReqIdentifier, mReply.(Resp).ErrCode, mReply.(Resp).ErrMsg, "Encode Msg error", conn.RemoteAddr().String(), err.Error())
|
||
3 years ago
|
return
|
||
|
}
|
||
|
err = ws.writeMsg(conn, websocket.BinaryMessage, b.Bytes())
|
||
4 years ago
|
if err != nil {
|
||
2 years ago
|
// uid, platform := ws.getUserUid(conn)
|
||
|
log.NewError(mReply.(Resp).OperationID, mReply.(Resp).ReqIdentifier, mReply.(Resp).ErrCode, mReply.(Resp).ErrMsg, "ws writeMsg error", conn.RemoteAddr().String(), err.Error())
|
||
2 years ago
|
} else {
|
||
|
log.Debug(mReply.(Resp).OperationID, mReply.(Resp).ReqIdentifier, mReply.(Resp).ErrCode, mReply.(Resp).ErrMsg, "ws write response success")
|
||
4 years ago
|
}
|
||
|
}
|
||
3 years ago
|
func (ws *WServer) sendErrMsg(conn *UserConn, errCode int32, errMsg string, reqIdentifier int32, msgIncr string, operationID string) {
|
||
|
mReply := Resp{
|
||
|
ReqIdentifier: reqIdentifier,
|
||
|
MsgIncr: msgIncr,
|
||
|
ErrCode: errCode,
|
||
|
ErrMsg: errMsg,
|
||
|
OperationID: operationID,
|
||
|
}
|
||
4 years ago
|
ws.sendMsg(conn, mReply)
|
||
|
}
|
||
2 years ago
|
|
||
|
func SetTokenKicked(userID string, platformID int, operationID string) {
|
||
|
m, err := db.DB.GetTokenMapByUidPid(userID, constant.PlatformIDToName(platformID))
|
||
|
if err != nil {
|
||
|
log.Error(operationID, "GetTokenMapByUidPid failed ", err.Error(), userID, constant.PlatformIDToName(platformID))
|
||
|
return
|
||
|
}
|
||
|
for k, _ := range m {
|
||
|
m[k] = constant.KickedToken
|
||
|
}
|
||
|
err = db.DB.SetTokenMapByUidPid(userID, platformID, m)
|
||
|
if err != nil {
|
||
|
log.Error(operationID, "SetTokenMapByUidPid failed ", err.Error(), userID, constant.PlatformIDToName(platformID))
|
||
|
return
|
||
|
}
|
||
|
}
|
||
2 years ago
|
|
||
|
func (ws *WServer) setUserDeviceBackground(conn *UserConn, m *Req) {
|
||
|
isPass, errCode, errMsg, pData := ws.argsValidate(m, constant.WsSetBackgroundStatus, m.OperationID)
|
||
|
if isPass {
|
||
2 years ago
|
req := pData.(*sdkws.SetAppBackgroundStatusReq)
|
||
2 years ago
|
conn.IsBackground = req.IsBackground
|
||
2 years ago
|
callbackResp := callbackUserOnline(m.OperationID, conn.userID, int(conn.PlatformID), conn.token, conn.IsBackground, conn.connID)
|
||
|
if callbackResp.ErrCode != 0 {
|
||
|
log.NewError(m.OperationID, utils.GetSelfFuncName(), "callbackUserOffline failed", callbackResp)
|
||
2 years ago
|
}
|
||
2 years ago
|
log.NewInfo(m.OperationID, "SetUserDeviceBackground", "success", *conn, req.IsBackground)
|
||
|
}
|
||
|
ws.setUserDeviceBackgroundResp(conn, m, errCode, errMsg)
|
||
|
}
|
||
|
|
||
|
func (ws *WServer) setUserDeviceBackgroundResp(conn *UserConn, m *Req, errCode int32, errMsg string) {
|
||
|
mReply := Resp{
|
||
|
ReqIdentifier: m.ReqIdentifier,
|
||
|
MsgIncr: m.MsgIncr,
|
||
|
OperationID: m.OperationID,
|
||
|
ErrCode: errCode,
|
||
|
ErrMsg: errMsg,
|
||
|
}
|
||
|
ws.sendMsg(conn, mReply)
|
||
|
}
|