Chuyển đến nội dung chính

Bài 14: gRPC & Protocol Buffers

Protocol Buffers schema design, code generation. Unary, server streaming, client streaming, bidirectional streaming. gRPC interceptors, error handling, deadlines. gRPC-Gateway cho REST compatibility.

💻 Lập trình — Bài 14 Bài 14: gRPC & Protocol Buffers

Golang: Từ Cơ bản đến Nâng cao

Phần 4: Advanced Features

xdev.asia

1. Protocol Buffers

1.1. Setup

# Install protoc compiler
brew install protobuf

# Install Go plugins
go install google.golang.org/protobuf/cmd/protoc-gen-go@latest
go install google.golang.org/grpc/cmd/protoc-gen-go-grpc@latest

# Go dependencies
go get google.golang.org/grpc
go get google.golang.org/protobuf

1.2. Proto File

// proto/user/v1/user.proto
syntax = "proto3";

package user.v1;

option go_package = "myapp/gen/user/v1;userv1";

import "google/protobuf/timestamp.proto";
import "google/protobuf/empty.proto";

// Messages
message User {
    uint64 id = 1;
    string name = 2;
    string email = 3;
    Role role = 4;
    bool is_active = 5;
    google.protobuf.Timestamp created_at = 6;
}

enum Role {
    ROLE_UNSPECIFIED = 0;
    ROLE_USER = 1;
    ROLE_EDITOR = 2;
    ROLE_ADMIN = 3;
}

message CreateUserRequest {
    string name = 1;
    string email = 2;
    string password = 3;
}

message GetUserRequest {
    uint64 id = 1;
}

message ListUsersRequest {
    int32 page = 1;
    int32 page_size = 2;
    string search = 3;
}

message ListUsersResponse {
    repeated User users = 1;
    int32 total = 2;
    int32 page = 3;
}

message UpdateUserRequest {
    uint64 id = 1;
    optional string name = 2;
    optional string email = 3;
    optional Role role = 4;
}

// Service
service UserService {
    rpc CreateUser(CreateUserRequest) returns (User);
    rpc GetUser(GetUserRequest) returns (User);
    rpc ListUsers(ListUsersRequest) returns (ListUsersResponse);
    rpc UpdateUser(UpdateUserRequest) returns (User);
    rpc DeleteUser(GetUserRequest) returns (google.protobuf.Empty);
    
    // Server streaming
    rpc WatchUsers(google.protobuf.Empty) returns (stream User);
    
    // Bidirectional streaming
    rpc Chat(stream ChatMessage) returns (stream ChatMessage);
}

message ChatMessage {
    string sender = 1;
    string content = 2;
    google.protobuf.Timestamp sent_at = 3;
}

1.3. Code Generation

# Generate Go code
protoc --go_out=. --go_opt=paths=source_relative \
       --go-grpc_out=. --go-grpc_opt=paths=source_relative \
       proto/user/v1/user.proto

# Makefile
.PHONY: proto
proto:
	protoc --go_out=. --go_opt=paths=source_relative \
	       --go-grpc_out=. --go-grpc_opt=paths=source_relative \
	       proto/**/**/*.proto

2. gRPC Server

package grpcserver

import (
    "context"
    "fmt"
    
    userv1 "myapp/gen/user/v1"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/status"
    "google.golang.org/protobuf/types/known/emptypb"
    "google.golang.org/protobuf/types/known/timestamppb"
)

type UserServer struct {
    userv1.UnimplementedUserServiceServer
    repo repository.UserRepository
}

func NewUserServer(repo repository.UserRepository) *UserServer {
    return &UserServer{repo: repo}
}

func (s *UserServer) CreateUser(ctx context.Context, req *userv1.CreateUserRequest) (*userv1.User, error) {
    // Validate input
    if req.Name == "" {
        return nil, status.Error(codes.InvalidArgument, "name is required")
    }
    if req.Email == "" {
        return nil, status.Error(codes.InvalidArgument, "email is required")
    }
    
    user, err := s.repo.Create(ctx, &model.User{
        Name:  req.Name,
        Email: req.Email,
    })
    if err != nil {
        return nil, status.Errorf(codes.Internal, "create user: %v", err)
    }
    
    return toProtoUser(user), nil
}

func (s *UserServer) GetUser(ctx context.Context, req *userv1.GetUserRequest) (*userv1.User, error) {
    user, err := s.repo.GetByID(ctx, uint(req.Id))
    if err != nil {
        return nil, status.Error(codes.NotFound, "user not found")
    }
    
    return toProtoUser(user), nil
}

func (s *UserServer) ListUsers(ctx context.Context, req *userv1.ListUsersRequest) (*userv1.ListUsersResponse, error) {
    page := int(req.Page)
    if page < 1 {
        page = 1
    }
    pageSize := int(req.PageSize)
    if pageSize < 1 || pageSize > 100 {
        pageSize = 20
    }
    
    users, total, err := s.repo.List(ctx, page, pageSize, req.Search)
    if err != nil {
        return nil, status.Errorf(codes.Internal, "list users: %v", err)
    }
    
    protoUsers := make([]*userv1.User, len(users))
    for i, u := range users {
        protoUsers[i] = toProtoUser(&u)
    }
    
    return &userv1.ListUsersResponse{
        Users: protoUsers,
        Total: int32(total),
        Page:  int32(page),
    }, nil
}

func (s *UserServer) DeleteUser(ctx context.Context, req *userv1.GetUserRequest) (*emptypb.Empty, error) {
    if err := s.repo.Delete(ctx, uint(req.Id)); err != nil {
        return nil, status.Error(codes.NotFound, "user not found")
    }
    return &emptypb.Empty{}, nil
}

// Server streaming
func (s *UserServer) WatchUsers(req *emptypb.Empty, stream userv1.UserService_WatchUsersServer) error {
    // Subscribe to user changes
    ch := s.repo.Subscribe()
    defer s.repo.Unsubscribe(ch)
    
    for {
        select {
        case <-stream.Context().Done():
            return nil
        case user := <-ch:
            if err := stream.Send(toProtoUser(user)); err != nil {
                return err
            }
        }
    }
}

func toProtoUser(u *model.User) *userv1.User {
    return &userv1.User{
        Id:        uint64(u.ID),
        Name:      u.Name,
        Email:     u.Email,
        IsActive:  u.IsActive,
        CreatedAt: timestamppb.New(u.CreatedAt),
    }
}

3. Server Setup

package main

import (
    "log"
    "net"
    
    "google.golang.org/grpc"
    "google.golang.org/grpc/reflection"
    
    userv1 "myapp/gen/user/v1"
)

func main() {
    lis, err := net.Listen("tcp", ":50051")
    if err != nil {
        log.Fatalf("Failed to listen: %v", err)
    }
    
    server := grpc.NewServer(
        grpc.ChainUnaryInterceptor(
            LoggingInterceptor,
            AuthInterceptor,
            RecoveryInterceptor,
        ),
        grpc.ChainStreamInterceptor(
            StreamLoggingInterceptor,
        ),
    )
    
    // Register services
    userv1.RegisterUserServiceServer(server, NewUserServer(userRepo))
    
    // Enable reflection for debugging (grpcurl, grpcui)
    reflection.Register(server)
    
    log.Printf("gRPC server listening on :50051")
    if err := server.Serve(lis); err != nil {
        log.Fatalf("Failed to serve: %v", err)
    }
}

4. Interceptors (Middleware)

import (
    "context"
    "log"
    "time"
    
    "google.golang.org/grpc"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/metadata"
    "google.golang.org/grpc/status"
)

// Logging interceptor
func LoggingInterceptor(
    ctx context.Context,
    req interface{},
    info *grpc.UnaryServerInfo,
    handler grpc.UnaryHandler,
) (interface{}, error) {
    start := time.Now()
    
    resp, err := handler(ctx, req)
    
    duration := time.Since(start)
    statusCode := codes.OK
    if err != nil {
        statusCode = status.Code(err)
    }
    
    log.Printf("gRPC %s | %s | %v | %v",
        info.FullMethod, statusCode, duration, err)
    
    return resp, err
}

// Auth interceptor
func AuthInterceptor(
    ctx context.Context,
    req interface{},
    info *grpc.UnaryServerInfo,
    handler grpc.UnaryHandler,
) (interface{}, error) {
    // Skip auth for certain methods
    publicMethods := map[string]bool{
        "/user.v1.UserService/CreateUser": true,
    }
    if publicMethods[info.FullMethod] {
        return handler(ctx, req)
    }
    
    // Extract token from metadata
    md, ok := metadata.FromIncomingContext(ctx)
    if !ok {
        return nil, status.Error(codes.Unauthenticated, "missing metadata")
    }
    
    tokens := md.Get("authorization")
    if len(tokens) == 0 {
        return nil, status.Error(codes.Unauthenticated, "missing token")
    }
    
    // Validate token...
    claims, err := validateToken(tokens[0])
    if err != nil {
        return nil, status.Error(codes.Unauthenticated, "invalid token")
    }
    
    // Set user info in context
    ctx = context.WithValue(ctx, "user_id", claims.UserID)
    
    return handler(ctx, req)
}

// Recovery interceptor
func RecoveryInterceptor(
    ctx context.Context,
    req interface{},
    info *grpc.UnaryServerInfo,
    handler grpc.UnaryHandler,
) (resp interface{}, err error) {
    defer func() {
        if r := recover(); r != nil {
            log.Printf("Panic recovered in %s: %v", info.FullMethod, r)
            err = status.Errorf(codes.Internal, "internal error")
        }
    }()
    
    return handler(ctx, req)
}

5. gRPC Client

package main

import (
    "context"
    "log"
    "time"
    
    "google.golang.org/grpc"
    "google.golang.org/grpc/credentials/insecure"
    
    userv1 "myapp/gen/user/v1"
)

func main() {
    conn, err := grpc.NewClient("localhost:50051",
        grpc.WithTransportCredentials(insecure.NewCredentials()),
    )
    if err != nil {
        log.Fatalf("Connect failed: %v", err)
    }
    defer conn.Close()
    
    client := userv1.NewUserServiceClient(conn)
    
    // Unary call with deadline
    ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
    defer cancel()
    
    user, err := client.CreateUser(ctx, &userv1.CreateUserRequest{
        Name:  "Duy Tran",
        Email: "[email protected]",
    })
    if err != nil {
        st, ok := status.FromError(err)
        if ok {
            log.Printf("gRPC error: code=%s, msg=%s", st.Code(), st.Message())
        }
        return
    }
    
    log.Printf("Created user: %+v", user)
    
    // Server streaming
    stream, err := client.WatchUsers(ctx, &emptypb.Empty{})
    if err != nil {
        log.Fatal(err)
    }
    
    for {
        user, err := stream.Recv()
        if err != nil {
            break
        }
        log.Printf("User update: %+v", user)
    }
}

6. gRPC Error Handling

gRPC CodeHTTP CodeKhi nào dùng
OK200Thành công
InvalidArgument400Input không hợp lệ
Unauthenticated401Chưa xác thực
PermissionDenied403Không có quyền
NotFound404Resource không tồn tại
AlreadyExists409Duplicate
Internal500Lỗi server
Unavailable503Service tạm ngưng
DeadlineExceeded504Timeout
// Error with details
import "google.golang.org/genproto/googleapis/rpc/errdetails"

func detailedError() error {
    st := status.New(codes.InvalidArgument, "invalid input")
    
    st, _ = st.WithDetails(&errdetails.BadRequest{
        FieldViolations: []*errdetails.BadRequest_FieldViolation{
            {Field: "email", Description: "invalid email format"},
            {Field: "name", Description: "name is required"},
        },
    })
    
    return st.Err()
}

7. gRPC-Gateway (REST compatibility)

go install github.com/grpc-ecosystem/grpc-gateway/v2/protoc-gen-grpc-gateway@latest
go install github.com/grpc-ecosystem/grpc-gateway/v2/protoc-gen-openapiv2@latest
import "google/api/annotations.proto";

service UserService {
    rpc GetUser(GetUserRequest) returns (User) {
        option (google.api.http) = {
            get: "/api/v1/users/{id}"
        };
    }
    
    rpc CreateUser(CreateUserRequest) returns (User) {
        option (google.api.http) = {
            post: "/api/v1/users"
            body: "*"
        };
    }
    
    rpc ListUsers(ListUsersRequest) returns (ListUsersResponse) {
        option (google.api.http) = {
            get: "/api/v1/users"
        };
    }
}
// Gateway server - proxy HTTP → gRPC
func runGateway() error {
    ctx := context.Background()
    mux := runtime.NewServeMux()
    
    opts := []grpc.DialOption{grpc.WithTransportCredentials(insecure.NewCredentials())}
    
    err := userv1.RegisterUserServiceHandlerFromEndpoint(ctx, mux, "localhost:50051", opts)
    if err != nil {
        return err
    }
    
    log.Println("Gateway listening on :8080")
    return http.ListenAndServe(":8080", mux)
}

Bài tiếp theo: Message Queues & Event-Driven Architecture — RabbitMQ, Kafka, và event sourcing.