package logic import ( "context" "fmt" "google.golang.org/grpc" "google.golang.org/grpc/reflection" "google.golang.org/protobuf/types/known/emptypb" "net" "sonet/api/gen/chat" "sonet/pkg/config" "sonet/pkg/grpc/discovery" "sonet/pkg/grpc/interceptor" "sonet/pkg/protocol/deliver" "sonet/pkg/protocol/session" "sonet/pkg/utils/collect" "sonet/pkg/utils/logger" "sonet/pkg/utils/shutdown" "strconv" "sync/atomic" "time" ) type ChatServer struct { chat.UnimplementedChatServer deliver *deliver.Deliver groupDeliver *deliver.GroupDeliver roomId int64 rooms *collect.ConcurrentMap[string, *chat.Room] } func NewChatServer(deliver *deliver.Deliver, groupDeliver *deliver.GroupDeliver) *ChatServer { rooms := collect.NewConcurrentMap[string, *chat.Room](32, func(k string) string { return k }) return &ChatServer{ deliver: deliver, groupDeliver: groupDeliver, roomId: 90000, rooms: rooms, } } func (s *ChatServer) Run(conf config.GrpcConfig, registry discovery.Registry) (err error) { server := grpc.NewServer( config.GetGrpcOptions( conf, grpc.UnaryInterceptor(interceptor.RecoverInterceptor), )..., ) if !conf.NoReflection { // 注册反射服务 reflection.Register(server) } chat.RegisterChatServer(server, s) listen, err := net.Listen("tcp", conf.Address) if err != nil { return } // registry discovery register := conf.Register if register.Name == "" { register.Name = chat.Chat_ServiceDesc.ServiceName } if register.Addr == "" { register.Addr = conf.Address } ctx, cancel := context.WithCancel(context.Background()) err = registry.Registry(ctx, register) if err != nil { panic(err) } shutdown.AddHook(cancel) // run serve logger.Infof("%s grpc server running %s\n", register.Name, listen.Addr().String()) err = server.Serve(listen) return } func (s *ChatServer) Send(ctx context.Context, req *chat.ReqSend) (*chat.ResSend, error) { subject, err := session.GetSubject(ctx) if err != nil { return nil, err } // TODO 消息队列消峰发送,ack 队列重发(可异步发送),asc超时第二次发送时同步判断是否不在线? // 投递消息 message := &chat.ChatMessage{Sender: subject.Uid, Content: req.Content} _, err = s.deliver.Deliver(ctx, message, req.Receiver) if err != nil { return nil, err } return &chat.ResSend{}, nil } func (s *ChatServer) RoomSend(ctx context.Context, req *chat.ReqRoomSend) (res *chat.ResSend, err error) { subject, err := session.GetSubject(ctx) if err != nil { return } gMessage := &chat.ChatMessage{ Type: 1, Sender: subject.Uid, Gid: req.Rid, Content: req.Message, } err = s.groupDeliver.DeliverGroup(ctx, req.Rid, gMessage) res = &chat.ResSend{} return } func (s *ChatServer) RoomCreate(ctx context.Context, req *chat.ReqRoomCreate) (res *chat.ResRoomCreate, err error) { subject, err := session.GetSubject(ctx) if err != nil { return } logger.Infof("create room: %s", req.Rname) room := &chat.Room{ Rid: strconv.Itoa(int(atomic.AddInt64(&s.roomId, 1))), MasterUid: subject.Uid, UserLimit: 10000, UserNum: 1, Gname: req.Rname, CreateUid: subject.Uid, CreateAt: time.Now().UnixMilli(), } // join postal room err = s.groupDeliver.GroupJoin(ctx, subject.Uid, []string{room.Rid}) if err != nil { return } s.rooms.Store(room.Rid, room) res = &chat.ResRoomCreate{Room: room} return } func (s *ChatServer) RoomJoin(ctx context.Context, req *chat.ReqRoomJoin) (emp *emptypb.Empty, err error) { emp = &emptypb.Empty{} subject, err := session.GetSubject(ctx) if err != nil { return } room, ok := s.rooms.Load(req.Rid) if !ok { err = fmt.Errorf("room %s not exists", req.Rid) return } atomic.AddInt32(&room.UserNum, 1) err = s.groupDeliver.GroupJoin(ctx, subject.Uid, []string{req.Rid}) return } func (s *ChatServer) RoomLeave(ctx context.Context, req *chat.ReqRoomLeave) (emp *emptypb.Empty, err error) { emp = &emptypb.Empty{} subject, err := session.GetSubject(ctx) if err != nil { return } room, ok := s.rooms.Load(req.Rid) if !ok { err = fmt.Errorf("room %s not exists", req.Rid) return } atomic.AddInt32(&room.UserNum, -1) err = s.groupDeliver.GroupLeave(ctx, subject.Uid, []string{req.Rid}) return } func (s *ChatServer) RoomKickOut(ctx context.Context, req *chat.ReqRoomKickOut) (emp *emptypb.Empty, err error) { emp = &emptypb.Empty{} return } func (s *ChatServer) RoomInfo(ctx context.Context, req *chat.ReqRoomInfo) (res *chat.ResRoomInfo, err error) { room, ok := s.rooms.Load(req.Rid) if !ok { err = fmt.Errorf("room %s not exists", req.Rid) return } res = &chat.ResRoomInfo{Room: room} return } // RoomDissolve 解散群组 func (s *ChatServer) RoomDissolve(ctx context.Context, req *chat.ReqRoomDissolve) (emp *emptypb.Empty, err error) { emp = &emptypb.Empty{} err = s.groupDeliver.GroupDissolve(ctx, req.Rid) if err != nil { return } s.rooms.Delete(req.Rid) return } func (s *ChatServer) RoomList(ctx context.Context, req *chat.ReqRoomList) (res *chat.ResRoomList, err error) { rooms := make([]*chat.Room, 16) s.rooms.Range(func(_ string, room *chat.Room) bool { rooms = append(rooms, room) return true }) res = &chat.ResRoomList{Rooms: rooms} return }