MFormations
Modern Go Engineering

Chapitre 16

16 - Projet Fil Rouge : EventHub

> Plateforme SaaS de traitement d'événements en temps réel

Cours 16 : Projet Fil Rouge - EventHub

1. Présentation du projet

EventHub est une plateforme SaaS de traitement d'événements en temps réel. Le projet couvre l'ensemble des concepts abordés dans la formation : microservices, génériques, performance, sécurité, CLI, gRPC, cloud.

1.1 Architecture globale

┌─────────────────────────────────────────────────────────────────┐
│                        API Gateway (Chi)                        │
│                    JWT Auth + Rate Limiting                     │
├──────────────────────┬──────────────────┬───────────────────────┤
│    User Service      │   Event Service  │  Notification Service │
│    (REST + gRPC)     │  (gRPC + Kafka)  │    (Events + Email)   │
├──────────┬───────────┼──────────────────┼───────────────────────┤
│PostgreSQL│   Redis   │     Kafka        │     SendGrid/SES      │
└──────────┴───────────┴──────────────────┴───────────────────────┘
│                        CLI Admin (Cobra)                        │
│                    Bubble Tea Dashboard                         │
└─────────────────────────────────────────────────────────────────┘

1.2 Modules

ModuleDescriptionStack
API GatewayPoint d'entrée uniqueChi, JWT, Rate Limiting
User ServiceGestion utilisateursgRPC, PostgreSQL, Redis
Event ServiceTraitement événementsgRPC, Kafka, Protobuf
Notification ServiceNotificationsEvents, Email, Webhooks
CLI AdminAdministrationCobra, Bubble Tea, Lip Gloss
InfrastructureDéploiementDocker, K8s, Terraform
CI/CDIntégration continueGitHub Actions, GoReleaser

2. API Gateway

2.1 Router et middleware

package gateway

import (
    "net/http"
    "time"

    "github.com/go-chi/chi/v5"
    "github.com/go-chi/chi/v5/middleware"
    "github.com/go-chi/cors"
)

type Gateway struct {
    router   *chi.Mux
    userSvc  *UserServiceClient
    eventSvc *EventServiceClient
    auth     *JWTAuth
}

func New(cfg *Config) *Gateway {
    g := &Gateway{
        router: chi.NewRouter(),
        auth:   NewJWTAuth(cfg.JWTSecret),
    }

    g.router.Use(middleware.RequestID)
    g.router.Use(middleware.RealIP)
    g.router.Use(middleware.Logger)
    g.router.Use(middleware.Recoverer)
    g.router.Use(middleware.Timeout(30 * time.Second))
    g.router.Use(cors.Handler(cors.Options{
        AllowedOrigins:   cfg.CORSOrigins,
        AllowedMethods:   []string{"GET", "POST", "PUT", "DELETE", "OPTIONS"},
        AllowedHeaders:   []string{"Accept", "Authorization", "Content-Type"},
        AllowCredentials: true,
        MaxAge:           300,
    }))

    g.registerRoutes()
    return g
}

func (g *Gateway) registerRoutes() {
    g.router.Route("/api/v1", func(r chi.Router) {
        // Public routes
        r.Post("/auth/login", g.handleLogin)
        r.Post("/auth/register", g.handleRegister)

        // Protected routes
        r.Group(func(r chi.Router) {
            r.Use(g.auth.Middleware)

            r.Get("/users/{id}", g.handleGetUser)
            r.Put("/users/{id}", g.handleUpdateUser)
            r.Get("/events", g.handleListEvents)
            r.Post("/events", g.handleCreateEvent)
            r.Get("/events/{id}", g.handleGetEvent)
            r.Delete("/events/{id}", g.handleDeleteEvent)
        })

        // Admin routes
        r.Group(func(r chi.Router) {
            r.Use(g.auth.Middleware)
            r.Use(RequireRole("admin"))
            r.Get("/admin/users", g.handleAdminListUsers)
            r.Get("/admin/stats", g.handleAdminStats)
        })
    })

    // Health
    g.router.Get("/health", g.handleHealth)
    g.router.Get("/ready", g.handleReadiness)
}

func (g *Gateway) ServeHTTP(w http.ResponseWriter, r *http.Request) {
    g.router.ServeHTTP(w, r)
}

2.2 Rate Limiting

package middleware

import (
    "net/http"
    "sync"
    "time"

    "golang.org/x/time/rate"
)

type RateLimiter struct {
    visitors map[string]*visitor
    mu       sync.Mutex
    rate     rate.Limit
    burst    int
}

type visitor struct {
    limiter  *rate.Limiter
    lastSeen time.Time
}

func NewRateLimiter(r rate.Limit, burst int) *RateLimiter {
    rl := &RateLimiter{
        visitors: make(map[string]*visitor),
        rate:     r,
        burst:    burst,
    }
    go rl.cleanup()
    return rl
}

func (rl *RateLimiter) Middleware(next http.Handler) http.Handler {
    return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
        ip := r.RemoteAddr
        rl.mu.Lock()
        v, exists := rl.visitors[ip]
        if !exists {
            v = &visitor{
                limiter:  rate.NewLimiter(rl.rate, rl.burst),
                lastSeen: time.Now(),
            }
            rl.visitors[ip] = v
        }
        v.lastSeen = time.Now()
        rl.mu.Unlock()

        if !v.limiter.Allow() {
            http.Error(w, "rate limit exceeded", http.StatusTooManyRequests)
            return
        }
        next.ServeHTTP(w, r)
    })
}

func (rl *RateLimiter) cleanup() {
    for {
        time.Sleep(time.Minute)
        rl.mu.Lock()
        for ip, v := range rl.visitors {
            if time.Since(v.lastSeen) > 3*time.Minute {
                delete(rl.visitors, ip)
            }
        }
        rl.mu.Unlock()
    }
}

3. User Service

3.1 Proto definition

syntax = "proto3";
package user.v1;
option go_package = "eventhub/gen/go/user/v1;userv1";

import "google/protobuf/timestamp.proto";

message User {
  string id = 1;
  string email = 2;
  string name = 3;
  repeated string roles = 4;
  bool active = 5;
  google.protobuf.Timestamp created_at = 6;
  google.protobuf.Timestamp updated_at = 7;
}

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

message CreateUserResponse {
  User user = 1;
  string token = 2;
}

message GetUserRequest { string id = 1; }
message GetUserByEmailRequest { string email = 1; }
message UpdateUserRequest {
  string id = 1;
  string name = 2;
  repeated string roles = 3;
}
message DeleteUserRequest { string id = 1; }
message ListUsersRequest {
  int32 page_size = 1;
  string page_token = 2;
}
message ListUsersResponse {
  repeated User users = 1;
  string next_page_token = 2;
  int32 total_count = 3;
}

message LoginRequest {
  string email = 1;
  string password = 2;
}

message LoginResponse {
  User user = 1;
  string token = 2;
}

service UserService {
  rpc CreateUser(CreateUserRequest) returns (CreateUserResponse);
  rpc GetUser(GetUserRequest) returns (User);
  rpc GetUserByEmail(GetUserByEmailRequest) returns (User);
  rpc UpdateUser(UpdateUserRequest) returns (User);
  rpc DeleteUser(DeleteUserRequest) returns (google.protobuf.Empty);
  rpc ListUsers(ListUsersRequest) returns (ListUsersResponse);
  rpc Login(LoginRequest) returns (LoginResponse);
  rpc ValidateToken(ValidateTokenRequest) returns (ValidateTokenResponse);
}

message ValidateTokenRequest { string token = 1; }
message ValidateTokenResponse {
  string user_id = 1;
  repeated string roles = 2;
}

3.2 Service implementation

package userservice

import (
    "context"
    "database/sql"
    "time"

    "github.com/golang-jwt/jwt/v5"
    "github.com/google/uuid"
    "github.com/jackc/pgx/v5/pgxpool"
    "golang.org/x/crypto/bcrypt"
    "google.golang.org/grpc/codes"
    "google.golang.org/grpc/status"

    userv1 "eventhub/gen/go/user/v1"
)

type UserService struct {
    userv1.UnimplementedUserServiceServer
    db     *pgxpool.Pool
    secret []byte
}

func New(db *pgxpool.Pool, jwtSecret string) *UserService {
    return &UserService{db: db, secret: []byte(jwtSecret)}
}

func (s *UserService) CreateUser(ctx context.Context, req *userv1.CreateUserRequest) (*userv1.CreateUserResponse, error) {
    if err := validateCreateUser(req); err != nil {
        return nil, status.Error(codes.InvalidArgument, err.Error())
    }

    hash, err := bcrypt.GenerateFromPassword([]byte(req.Password), bcrypt.DefaultCost)
    if err != nil {
        return nil, status.Errorf(codes.Internal, "hash password: %v", err)
    }

    id := uuid.New().String()
    now := time.Now()

    _, err = s.db.Exec(ctx,
        `INSERT INTO users (id, email, name, password_hash, roles, active, created_at, updated_at)
         VALUES ($1, $2, $3, $4, $5, true, $6, $6)`,
        id, req.Email, req.Name, string(hash), []string{"user"}, now,
    )
    if err != nil {
        return nil, status.Errorf(codes.Internal, "create user: %v", err)
    }

    token, err := s.generateToken(id, []string{"user"})
    if err != nil {
        return nil, err
    }

    return &userv1.CreateUserResponse{
        User: &userv1.User{
            Id:        id,
            Email:     req.Email,
            Name:      req.Name,
            Roles:     []string{"user"},
            Active:    true,
            CreatedAt: timestamppb.New(now),
            UpdatedAt: timestamppb.New(now),
        },
        Token: token,
    }, nil
}

func (s *UserService) Login(ctx context.Context, req *userv1.LoginRequest) (*userv1.LoginResponse, error) {
    var user struct {
        ID           string
        Email, Name  string
        PasswordHash string
        Roles        []string
        Active       bool
    }

    err := s.db.QueryRow(ctx,
        `SELECT id, email, name, password_hash, roles, active FROM users WHERE email = $1`,
        req.Email,
    ).Scan(&user.ID, &user.Email, &user.Name, &user.PasswordHash, &user.Roles, &user.Active)
    if err != nil {
        return nil, status.Error(codes.Unauthenticated, "invalid credentials")
    }

    if !user.Active {
        return nil, status.Error(codes.PermissionDenied, "account deactivated")
    }

    if err := bcrypt.CompareHashAndPassword([]byte(user.PasswordHash), []byte(req.Password)); err != nil {
        return nil, status.Error(codes.Unauthenticated, "invalid credentials")
    }

    token, err := s.generateToken(user.ID, user.Roles)
    if err != nil {
        return nil, err
    }

    return &userv1.LoginResponse{
        User: &userv1.User{
            Id:    user.ID,
            Email: user.Email,
            Name:  user.Name,
            Roles: user.Roles,
            Active: user.Active,
        },
        Token: token,
    }, nil
}

func (s *UserService) generateToken(userID string, roles []string) (string, error) {
    claims := jwt.MapClaims{
        "sub":   userID,
        "roles": roles,
        "iat":   time.Now().Unix(),
        "exp":   time.Now().Add(24 * time.Hour).Unix(),
        "iss":   "eventhub",
    }
    token := jwt.NewWithClaims(jwt.SigningMethodHS256, claims)
    return token.SignedString(s.secret)
}

4. Event Service (gRPC + Kafka)

4.1 Proto definition

syntax = "proto3";
package event.v1;
option go_package = "eventhub/gen/go/event/v1;eventv1";

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

message Event {
  string id = 1;
  string type = 2;
  string source = 3;
  string subject = 4;
  google.protobuf.Struct data = 5;
  map<string, string> metadata = 6;
  string user_id = 7;
  google.protobuf.Timestamp time = 8;
  string specversion = 9;
  string dataschema = 10;
}

message CreateEventRequest {
  string type = 1;
  string source = 2;
  string subject = 3;
  google.protobuf.Struct data = 4;
  map<string, string> metadata = 5;
}

message CreateEventResponse {
  Event event = 1;
}

message GetEventRequest { string id = 1; }
message ListEventsRequest {
  string user_id = 1;
  string type = 2;
  int32 page_size = 3;
  string page_token = 4;
  google.protobuf.Timestamp start_time = 5;
  google.protobuf.Timestamp end_time = 6;
}

message ListEventsResponse {
  repeated Event events = 1;
  string next_page_token = 2;
  int32 total_count = 3;
}

message SubscribeRequest {
  string user_id = 1;
  string event_type = 2;
  string webhook_url = 3;
}

service EventService {
  rpc CreateEvent(CreateEventRequest) returns (CreateEventResponse);
  rpc GetEvent(GetEventRequest) returns (Event);
  rpc ListEvents(ListEventsRequest) returns (ListEventsResponse);
  rpc StreamEvents(ListEventsRequest) returns (stream Event);
  rpc Subscribe(SubscribeRequest) returns (Subscription);
  rpc Unsubscribe(UnsubscribeRequest) returns (google.protobuf.Empty);
}

4.2 Kafka Producer

package eventservice

import (
    "context"
    "encoding/json"
    "fmt"
    "time"

    "github.com/google/uuid"
    "github.com/segmentio/kafka-go"
    "google.golang.org/protobuf/types/known/structpb"
    "google.golang.org/protobuf/types/known/timestamppb"

    eventv1 "eventhub/gen/go/event/v1"
)

type EventService struct {
    eventv1.UnimplementedEventServiceServer
    producer *kafka.Writer
    consumer *kafka.Reader
    repo     EventRepository
}

func New(brokers []string, repo EventRepository) *EventService {
    return &EventService{
        producer: &kafka.Writer{
            Addr:     kafka.TCP(brokers...),
            Balancer: &kafka.LeastBytes{},
        },
        repo: repo,
    }
}

func (s *EventService) CreateEvent(ctx context.Context, req *eventv1.CreateEventRequest) (*eventv1.CreateEventResponse, error) {
    event := &eventv1.Event{
        Id:          uuid.New().String(),
        Type:        req.Type,
        Source:      req.Source,
        Subject:     req.Subject,
        Data:        req.Data,
        Metadata:    req.Metadata,
        Time:        timestamppb.Now(),
        Specversion: "1.0",
    }

    // Persist to database
    if err := s.repo.Save(ctx, event); err != nil {
        return nil, fmt.Errorf("save event: %w", err)
    }

    // Publish to Kafka
    data, _ := json.Marshal(event)
    msg := kafka.Message{
        Key:   []byte(event.Id),
        Value: data,
        Headers: []kafka.Header{
            {Key: "event-type", Value: []byte(event.Type)},
            {Key: "source", Value: []byte(event.Source)},
        },
    }

    if err := s.producer.WriteMessages(ctx, msg); err != nil {
        return nil, fmt.Errorf("publish event: %w", err)
    }

    return &eventv1.CreateEventResponse{Event: event}, nil
}

func (s *EventService) StreamEvents(req *eventv1.ListEventsRequest, stream eventv1.EventService_StreamEventsServer) error {
    // Stream from Kafka consumer
    for {
        msg, err := s.consumer.ReadMessage(stream.Context())
        if err != nil {
            return err
        }

        var event eventv1.Event
        if err := json.Unmarshal(msg.Value, &event); err != nil {
            continue
        }

        if err := stream.Send(&event); err != nil {
            return err
        }
    }
}

5. Notification Service

package notificationservice

import (
    "context"
    "encoding/json"
    "log"

    "github.com/segmentio/kafka-go"
    "github.com/sendgrid/sendgrid-go"
    "github.com/sendgrid/sendgrid-go/helpers/mail"
)

type NotificationService struct {
    consumer   *kafka.Reader
    sendGrid   *sendgrid.Client
    webhookSvc *WebhookService
}

func New(brokers []string, sendGridAPIKey string) *NotificationService {
    return &NotificationService{
        consumer: kafka.NewReader(kafka.ReaderConfig{
            Brokers:   brokers,
            Topic:     "events",
            GroupID:   "notification-service",
            MinBytes:  10e3,
            MaxBytes:  10e6,
        }),
        sendGrid: sendgrid.NewSendClient(sendGridAPIKey),
    }
}

func (s *NotificationService) Start(ctx context.Context) {
    for {
        msg, err := s.consumer.ReadMessage(ctx)
        if err != nil {
            log.Printf("read message: %v", err)
            continue
        }

        var event Event
        if err := json.Unmarshal(msg.Value, &event); err != nil {
            log.Printf("unmarshal: %v", err)
            continue
        }

        s.processEvent(ctx, event)
    }
}

func (s *NotificationService) processEvent(ctx context.Context, event Event) {
    switch event.Type {
    case "user.created":
        s.sendWelcomeEmail(ctx, event)
    case "event.created":
        s.notifySubscribers(ctx, event)
    case "alert.high":
        s.sendAlert(ctx, event)
    default:
        s.forwardToWebhooks(ctx, event)
    }
}

func (s *NotificationService) sendWelcomeEmail(ctx context.Context, event Event) error {
    from := mail.NewEmail("EventHub", "noreply@eventhub.io")
    to := mail.NewEmail(event.Subject, event.Metadata["email"])
    subject := "Welcome to EventHub!"
    content := mail.NewContent("text/plain", "Thank you for joining EventHub!")

    message := mail.NewSingleEmail(from, subject, to, content)
    _, err := s.sendGrid.Send(message)
    return err
}

6. CLI Admin (Cobra + Bubble Tea)

package cmd

import (
    "eventhub/internal/api"

    "github.com/spf13/cobra"
)

var rootCmd = &cobra.Command{
    Use:   "eventhub",
    Short: "EventHub CLI - Manage your event platform",
    Long: `EventHub CLI for managing events, users, and monitoring.
    
EventHub is a real-time event processing platform.`,
}

var userCmd = &cobra.Command{
    Use:   "user",
    Short: "Manage users",
}

var userCreateCmd = &cobra.Command{
    Use:   "create <email> <name>",
    Short: "Create a new user",
    Args:  cobra.ExactArgs(2),
    RunE: func(cmd *cobra.Command, args []string) error {
        client, err := api.NewClient(cfg)
        if err != nil {
            return err
        }
        password, _ := cmd.Flags().GetString("password")
        user, err := client.CreateUser(cmd.Context(), args[0], args[1], password)
        if err != nil {
            return err
        }
        cmd.Printf("User created: %s (%s)\n", user.ID, user.Email)
        return nil
    },
}

var userListCmd = &cobra.Command{
    Use:   "list",
    Short: "List users",
    RunE: func(cmd *cobra.Command, args []string) error {
        client, err := api.NewClient(cfg)
        if err != nil {
            return err
        }
        users, err := client.ListUsers(cmd.Context())
        if err != nil {
            return err
        }

        w := tabwriter.NewWriter(cmd.OutOrStdout(), 0, 0, 3, ' ', 0)
        fmt.Fprintln(w, "ID\tEMAIL\tNAME\tROLES\tACTIVE")
        for _, u := range users {
            fmt.Fprintf(w, "%s\t%s\t%s\t%v\t%v\n",
                u.ID[:8], u.Email, u.Name, u.Roles, u.Active)
        }
        w.Flush()
        return nil
    },
}

var eventCmd = &cobra.Command{
    Use:   "event",
    Short: "Manage events",
}

var eventListCmd = &cobra.Command{
    Use:   "list",
    Short: "List events",
    RunE: func(cmd *cobra.Command, args []string) error {
        client, err := api.NewClient(cfg)
        if err != nil {
            return err
        }
        events, err := client.ListEvents(cmd.Context())
        if err != nil {
            return err
        }

        for _, e := range events {
            fmt.Fprintf(cmd.OutOrStdout(), "%s [%s] %s\n", e.ID[:8], e.Type, e.Subject)
        }
        return nil
    },
}

var dashboardCmd = &cobra.Command{
    Use:   "dashboard",
    Short: "Open TUI dashboard",
    RunE: func(cmd *cobra.Command, args []string) error {
        return startDashboard(cfg)
    },
}

func init() {
    userCreateCmd.Flags().StringP("password", "p", "", "User password (auto-generated if empty)")
    userCreateCmd.MarkFlagRequired("password")

    userCmd.AddCommand(userCreateCmd)
    userCmd.AddCommand(userListCmd)
    eventCmd.AddCommand(eventListCmd)

    rootCmd.AddCommand(userCmd)
    rootCmd.AddCommand(eventCmd)
    rootCmd.AddCommand(dashboardCmd)
    rootCmd.AddCommand(versionCmd)
}

7. Infrastructure (Docker + K8s)

7.1 Docker Compose

version: '3.8'
services:
  gateway:
    build: ./services/gateway
    ports: ["8080:8080"]
    depends_on: [user-service, event-service]

  user-service:
    build: ./services/user
    ports: ["50051:50051"]
    environment:
      DB_HOST: postgres
      REDIS_HOST: redis
    depends_on: [postgres, redis]

  event-service:
    build: ./services/event
    ports: ["50052:50052"]
    environment:
      KAFKA_BROKERS: kafka:9092
      DB_HOST: postgres
    depends_on: [postgres, kafka]

  notification-service:
    build: ./services/notification
    environment:
      KAFKA_BROKERS: kafka:9092
      SENDGRID_API_KEY: ${SENDGRID_API_KEY}
    depends_on: [kafka]

  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_DB: eventhub
      POSTGRES_USER: eventhub
      POSTGRES_PASSWORD: ${DB_PASSWORD}
    volumes: [pgdata:/var/lib/postgresql/data]

  redis:
    image: redis:7-alpine
    volumes: [redisdata:/data]

  kafka:
    image: confluentinc/cp-kafka:latest
    environment:
      KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://kafka:9092
    depends_on: [zookeeper]

  zookeeper:
    image: confluentinc/cp-zookeeper:latest

8. CI/CD

# .github/workflows/ci.yml
name: CI/CD
on:
  push:
    branches: [main]
  pull_request:

jobs:
  test:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
      - uses: actions/setup-go@v5
        with:
          go-version: '1.22'
      - run: go test -race -coverprofile=coverage.txt ./...
      - run: go vet ./...

  lint:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
      - uses: golangci/golangci-lint-action@v4

  security:
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
      - uses: actions/setup-go@v5
      - run: go install golang.org/x/vuln/cmd/govulncheck@latest
      - run: govulncheck ./...

  release:
    needs: [test, lint, security]
    if: startsWith(github.ref, 'refs/tags/v')
    runs-on: ubuntu-latest
    steps:
      - uses: actions/checkout@v4
      - uses: actions/setup-go@v5
      - uses: goreleaser/goreleaser-action@v5
        with:
          args: release --clean

9. Diagrammes

Diagramme en cours de génération...
Diagramme en cours de génération...

Points Clés

  1. Architecture microservices : API Gateway + 3 services spécialisés
  2. Communication : gRPC interne, REST externe, Kafka pour events
  3. Sécurité : JWT auth, bcrypt passwords, secure headers
  4. CLI : Cobra + Bubble Tea pour l'administration
  5. Infrastructure : Docker multi-stage, K8s, GoReleaser
  6. Observabilité : OpenTelemetry, Prometheus, structured logging