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
| Module | Description | Stack |
|---|---|---|
| API Gateway | Point d'entrée unique | Chi, JWT, Rate Limiting |
| User Service | Gestion utilisateurs | gRPC, PostgreSQL, Redis |
| Event Service | Traitement événements | gRPC, Kafka, Protobuf |
| Notification Service | Notifications | Events, Email, Webhooks |
| CLI Admin | Administration | Cobra, Bubble Tea, Lip Gloss |
| Infrastructure | Déploiement | Docker, K8s, Terraform |
| CI/CD | Intégration continue | GitHub 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
- Architecture microservices : API Gateway + 3 services spécialisés
- Communication : gRPC interne, REST externe, Kafka pour events
- Sécurité : JWT auth, bcrypt passwords, secure headers
- CLI : Cobra + Bubble Tea pour l'administration
- Infrastructure : Docker multi-stage, K8s, GoReleaser
- Observabilité : OpenTelemetry, Prometheus, structured logging