58 lines
1.0 KiB
Go
58 lines
1.0 KiB
Go
package grpc
|
|
|
|
import (
|
|
"context"
|
|
"log"
|
|
"net"
|
|
|
|
pb "spektr/internal/grpc/pb"
|
|
|
|
"google.golang.org/grpc"
|
|
)
|
|
|
|
type server struct {
|
|
pb.UnimplementedReceiverServiceServer
|
|
handler Handler
|
|
}
|
|
|
|
func (s *server) Connect(stream pb.ReceiverService_ConnectServer) error {
|
|
log.Println("client connected")
|
|
for {
|
|
msg, err := stream.Recv()
|
|
if err != nil {
|
|
log.Println("client disconnected:", err)
|
|
return err
|
|
}
|
|
|
|
log.Printf("recv request_id=%s", msg.GetRequestId())
|
|
resp := s.handler.Handle(stream.Context(), msg)
|
|
if resp == nil {
|
|
continue
|
|
}
|
|
|
|
if err := stream.Send(resp); err != nil {
|
|
log.Println("send error:", err)
|
|
return err
|
|
}
|
|
log.Printf("sent request_id=%s", resp.GetRequestId())
|
|
}
|
|
}
|
|
|
|
func Serve(ctx context.Context, addr string, handler Handler) error {
|
|
lis, err := net.Listen("tcp", addr)
|
|
if err != nil {
|
|
return err
|
|
}
|
|
|
|
srv := grpc.NewServer()
|
|
pb.RegisterReceiverServiceServer(srv, &server{handler: handler})
|
|
|
|
go func() {
|
|
<-ctx.Done()
|
|
srv.GracefulStop()
|
|
}()
|
|
|
|
log.Println("gRPC listening on", addr)
|
|
return srv.Serve(lis)
|
|
}
|