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) }