first
This commit is contained in:
@@ -0,0 +1,48 @@
|
||||
package grpc
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
pb "spektr/internal/grpc/pb"
|
||||
|
||||
"google.golang.org/grpc"
|
||||
"google.golang.org/grpc/credentials/insecure"
|
||||
)
|
||||
|
||||
type Client struct {
|
||||
conn *grpc.ClientConn
|
||||
cancel context.CancelFunc
|
||||
stream pb.ReceiverService_ConnectClient
|
||||
}
|
||||
|
||||
func Dial(ctx context.Context, addr string) (*Client, error) {
|
||||
conn, err := grpc.DialContext(ctx, addr, grpc.WithBlock(), grpc.WithTransportCredentials(insecure.NewCredentials()))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
client := pb.NewReceiverServiceClient(conn)
|
||||
streamCtx, cancel := context.WithCancel(context.Background())
|
||||
stream, err := client.Connect(streamCtx)
|
||||
if err != nil {
|
||||
cancel()
|
||||
_ = conn.Close()
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Client{conn: conn, cancel: cancel, stream: stream}, nil
|
||||
}
|
||||
|
||||
func (c *Client) Send(msg *pb.ClientMessage) (*pb.ServerMessage, error) {
|
||||
if err := c.stream.Send(msg); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return c.stream.Recv()
|
||||
}
|
||||
|
||||
func (c *Client) Close() error {
|
||||
if c.cancel != nil {
|
||||
c.cancel()
|
||||
}
|
||||
return c.conn.Close()
|
||||
}
|
||||
@@ -0,0 +1,57 @@
|
||||
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)
|
||||
}
|
||||
@@ -0,0 +1,11 @@
|
||||
package grpc
|
||||
|
||||
import (
|
||||
"context"
|
||||
|
||||
pb "spektr/internal/grpc/pb"
|
||||
)
|
||||
|
||||
type Handler interface {
|
||||
Handle(context.Context, *pb.ClientMessage) *pb.ServerMessage
|
||||
}
|
||||
Reference in New Issue
Block a user