MFormations
Modern Go Engineering

Chapitre 9

09 - Microservices avec Go

09 - Microservices avec Go

Cours 09 : Microservices avec Go

1. Fondamentaux des microservices

1.1 Définition et principes

L'architecture microservices est un style d'architecture où une application est composée de petits services indépendants, chacun exécutant un processus unique et communiquant via des mécanismes légers (HTTP, gRPC, messages).

Principes clés :

  • Single Responsibility : Chaque service a une responsabilité unique et bien définie
  • Déploiement indépendant : Chaque service peut être déployé sans impact sur les autres
  • Isolation des pannes : Une défaillance dans un service ne doit pas impacter les autres
  • Autonomie des équipes : Chaque service peut être développé par une équipe différente
  • Polyglotte : Chaque service peut utiliser la techno la plus adaptée

1.2 Monolithe vs Microservices

// Monolithe : tout dans un seul package
package monolith

type OrderService struct {
    db     *sql.DB
    cache  *redis.Client
    email  *EmailClient
}

func (s *OrderService) CreateOrder(ctx context.Context, o Order) error {
    // Validation
    if err := s.validateOrder(o); err != nil {
        return err
    }
    // Base de données
    if err := s.db.SaveOrder(ctx, o); err != nil {
        return err
    }
    // Cache
    s.cache.InvalidateOrders(ctx, o.UserID)
    // Email
    s.email.SendConfirmation(o.UserEmail, o.ID)
    return nil
}
// Microservices : découplage en services spécialisés
package orderservice

type OrderService struct {
    orderRepo  *OrderRepository
    eventBus   EventBus
    userClient UserServiceClient
}

func (s *OrderService) CreateOrder(ctx context.Context, o Order) error {
    if err := s.orderRepo.Save(ctx, o); err != nil {
        return fmt.Errorf("save order: %w", err)
    }
    return s.eventBus.Publish(ctx, OrderCreatedEvent{Order: o})
}

1.3 Size des services

La taille d'un service se mesure en :

  • Responsabilité fonctionnelle : Un bounded context DDD
  • Complexité : Pas plus de quelques centaines de lignes par fichier
  • Équipe : 2 pizzas rule (Amazon) — une équipe doit pouvoir tenir autour de 2 pizzas

2. API Gateway

2.1 Rôle de l'API Gateway

L'API Gateway est le point d'entrée unique pour tous les clients. Elle gère :

  • Routage vers les services internes
  • Authentification et autorisation
  • Rate limiting
  • Aggregation de réponses
  • Transformation de protocoles
package gateway

import (
    "net/http"
    "net/http/httputil"
    "net/url"
)

type Gateway struct {
    auth      AuthMiddleware
    router    *http.ServeMux
    rateLimit *RateLimiter
    targets   map[string]*url.URL
}

func NewGateway(auth AuthMiddleware) *Gateway {
    g := &Gateway{
        auth:      auth,
        router:    http.NewServeMux(),
        rateLimit: NewRateLimiter(100, time.Minute),
        targets:   make(map[string]*url.URL),
    }
    g.registerRoutes()
    return g
}

func (g *Gateway) registerRoutes() {
    g.router.HandleFunc("/api/v1/users", g.auth.Wrap(g.proxy("user-service")))
    g.router.HandleFunc("/api/v1/orders", g.auth.Wrap(g.proxy("order-service")))
    g.router.HandleFunc("/api/v1/events", g.auth.Wrap(g.proxy("event-service")))
}

func (g *Gateway) proxy(serviceName string) http.HandlerFunc {
    target := g.targets[serviceName]
    proxy := httputil.NewSingleHostReverseProxy(target)
    return func(w http.ResponseWriter, r *http.Request) {
        if !g.rateLimit.Allow(r.RemoteAddr) {
            http.Error(w, "rate limit exceeded", http.StatusTooManyRequests)
            return
        }
        proxy.ServeHTTP(w, r)
    }
}

2.2 Patterns API Gateway

// Backend for Frontend (BFF)
type BFFGateway struct {
    mobileGateway *Gateway   // Optimisé mobile : payload réduit
    webGateway    *Gateway   // Optimisé web : HTML + JSON
    iotGateway   *Gateway   // Optimisé IoT : protobuf binaire
}

// Response aggregation
func (g *Gateway) aggregateOrderDetails(ctx context.Context, orderID string) (*OrderDetail, error) {
    // Fan-out requests
    type result struct {
        data any
        err  error
    }

    ch := make(chan result, 3)
    go func() {
        order, err := g.orderClient.GetOrder(ctx, orderID)
        ch <- result{order, err}
    }()
    go func() {
        user, err := g.userClient.GetUser(ctx, orderID)
        ch <- result{user, err}
    }()
    go func() {
        payment, err := g.paymentClient.GetPayment(ctx, orderID)
        ch <- result{payment, err}
    }()

    // Aggregate responses
    detail := &OrderDetail{}
    for i := 0; i < 3; i++ {
        r := <-ch
        if r.err != nil {
            return nil, r.err
        }
        switch v := r.data.(type) {
        case *Order:
            detail.Order = v
        case *User:
            detail.User = v
        case *Payment:
            detail.Payment = v
        }
    }
    return detail, nil
}

3. Service Registry et Discovery

3.1 Service Registry

package registry

import (
    "context"
    "sync"
    "time"
)

type ServiceInstance struct {
    ID        string
    Name      string
    Address   string
    Port      int
    Metadata  map[string]string
    Healthy   bool
    LastCheck time.Time
}

type Registry interface {
    Register(ctx context.Context, instance ServiceInstance) error
    Deregister(ctx context.Context, instanceID string) error
    Discover(ctx context.Context, serviceName string) ([]ServiceInstance, error)
    Watch(ctx context.Context, serviceName string) (<-chan []ServiceInstance, error)
    HealthCheck(ctx context.Context, instanceID string) error
}

type InMemoryRegistry struct {
    mu        sync.RWMutex
    instances map[string]map[string]*ServiceInstance // serviceName -> instanceID -> instance
    watchers  map[string][]chan []ServiceInstance
}

func NewInMemoryRegistry() *InMemoryRegistry {
    r := &InMemoryRegistry{
        instances: make(map[string]map[string]*ServiceInstance),
        watchers: make(map[string][]chan []ServiceInstance),
    }
    go r.healthCheckLoop()
    return r
}

func (r *InMemoryRegistry) Register(ctx context.Context, inst ServiceInstance) error {
    r.mu.Lock()
    defer r.mu.Unlock()
    inst.Healthy = true
    inst.LastCheck = time.Now()
    if r.instances[inst.Name] == nil {
        r.instances[inst.Name] = make(map[string]*ServiceInstance)
    }
    r.instances[inst.Name][inst.ID] = &inst
    r.notifyWatchers(inst.Name)
    return nil
}

func (r *InMemoryRegistry) Discover(ctx context.Context, serviceName string) ([]ServiceInstance, error) {
    r.mu.RLock()
    defer r.mu.RUnlock()
    instances := r.instances[serviceName]
    result := make([]ServiceInstance, 0, len(instances))
    for _, inst := range instances {
        if inst.Healthy {
            result = append(result, *inst)
        }
    }
    return result, nil
}

func (r *InMemoryRegistry) healthCheckLoop() {
    ticker := time.NewTicker(30 * time.Second)
    for range ticker.C {
        r.mu.Lock()
        for name, instances := range r.instances {
            for id, inst := range instances {
                if time.Since(inst.LastCheck) > 60*time.Second {
                    inst.Healthy = false
                    r.notifyWatchers(name)
                    _ = r.Deregister(context.Background(), id)
                }
            }
        }
        r.mu.Unlock()
    }
}

3.2 Client-side Discovery

package discovery

import (
    "math/rand"
    "sync"
)

type LoadBalancer interface {
    Pick(instances []ServiceInstance) (ServiceInstance, error)
}

type RoundRobinLB struct {
    mu    sync.Mutex
    count map[string]int
}

func (lb *RoundRobinLB) Pick(instances []ServiceInstance) (ServiceInstance, error) {
    if len(instances) == 0 {
        return ServiceInstance{}, ErrNoInstances
    }
    lb.mu.Lock()
    defer lb.mu.Unlock()
    if lb.count == nil {
        lb.count = make(map[string]int)
    }
    serviceName := instances[0].Name
    lb.count[serviceName]++
    idx := lb.count[serviceName] % len(instances)
    return instances[idx], nil
}

type RandomLB struct{}

func (lb *RandomLB) Pick(instances []ServiceInstance) (ServiceInstance, error) {
    if len(instances) == 0 {
        return ServiceInstance{}, ErrNoInstances
    }
    return instances[rand.Intn(len(instances))], nil
}

4. Résilience

4.1 Circuit Breaker

package circuitbreaker

import (
    "errors"
    "sync"
    "time"
)

type State int

const (
    StateClosed   State = iota // Normal operation
    StateOpen                  // Failing, reject requests
    StateHalfOpen              // Testing if service recovered
)

type CircuitBreaker struct {
    mu                sync.RWMutex
    state             State
    failureCount      int
    successCount      int
    failureThreshold  int
    successThreshold  int
    timeout           time.Duration
    lastStateChange   time.Time
}

func New(failureThreshold, successThreshold int, timeout time.Duration) *CircuitBreaker {
    return &CircuitBreaker{
        state:            StateClosed,
        failureThreshold: failureThreshold,
        successThreshold: successThreshold,
        timeout:          timeout,
    }
}

func (cb *CircuitBreaker) Execute(fn func() error) error {
    if !cb.allowRequest() {
        return ErrCircuitOpen
    }

    err := fn()
    cb.recordResult(err)
    return err
}

func (cb *CircuitBreaker) allowRequest() bool {
    cb.mu.RLock()
    state := cb.state
    cb.mu.RUnlock()

    switch state {
    case StateClosed:
        return true
    case StateOpen:
        if time.Since(cb.lastStateChange) > cb.timeout {
            cb.mu.Lock()
            cb.state = StateHalfOpen
            cb.successCount = 0
            cb.mu.Unlock()
            return true
        }
        return false
    case StateHalfOpen:
        return true
    default:
        return true
    }
}

func (cb *CircuitBreaker) recordResult(err error) {
    cb.mu.Lock()
    defer cb.mu.Unlock()

    if err != nil {
        cb.failureCount++
        cb.successCount = 0
        if cb.failureCount >= cb.failureThreshold {
            cb.state = StateOpen
            cb.lastStateChange = time.Now()
        }
    } else {
        cb.successCount++
        cb.failureCount = 0
        if cb.state == StateHalfOpen && cb.successCount >= cb.successThreshold {
            cb.state = StateClosed
        }
    }
}

4.2 Retry avec backoff

package retry

import (
    "context"
    "math"
    "math/rand"
    "time"
)

type Strategy int

const (
    StrategyConstant    Strategy = iota
    StrategyLinear
    StrategyExponential
)

type Config struct {
    MaxAttempts int
    Delay       time.Duration
    MaxDelay    time.Duration
    Strategy    Strategy
    Jitter      float64 // 0.0 - 1.0
}

func Do(ctx context.Context, cfg Config, fn func(context.Context) error) error {
    var lastErr error
    for attempt := 0; attempt < cfg.MaxAttempts; attempt++ {
        if attempt > 0 {
            delay := calculateDelay(cfg, attempt)
            select {
            case <-ctx.Done():
                return ctx.Err()
            case <-time.After(delay):
            }
        }
        if err := fn(ctx); err != nil {
            lastErr = err
            continue
        }
        return nil
    }
    return lastErr
}

func calculateDelay(cfg Config, attempt int) time.Duration {
    var delay time.Duration
    switch cfg.Strategy {
    case StrategyConstant:
        delay = cfg.Delay
    case StrategyLinear:
        delay = cfg.Delay * time.Duration(attempt)
    case StrategyExponential:
        delay = cfg.Delay * time.Duration(math.Pow(2, float64(attempt)))
    }
    if cfg.MaxDelay > 0 && delay > cfg.MaxDelay {
        delay = cfg.MaxDelay
    }
    if cfg.Jitter > 0 {
        jitter := time.Duration(float64(delay) * cfg.Jitter * rand.Float64())
        delay += jitter
    }
    return delay
}

4.3 Timeout avec context

package timeout

import (
    "context"
    "net/http"
    "time"
)

func TimeoutMiddleware(timeout time.Duration) func(http.Handler) http.Handler {
    return func(next http.Handler) http.Handler {
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
            ctx, cancel := context.WithTimeout(r.Context(), timeout)
            defer cancel()
            r = r.WithContext(ctx)
            done := make(chan bool)

            go func() {
                next.ServeHTTP(w, r)
                done <- true
            }()

            select {
            case <-done:
                return
            case <-ctx.Done():
                switch ctx.Err() {
                case context.DeadlineExceeded:
                    http.Error(w, "request timeout", http.StatusGatewayTimeout)
                case context.Canceled:
                    http.Error(w, "request cancelled", http.StatusRequestTimeout)
                }
            }
        })
    }
}

5. Communication synchrone

5.1 Communication HTTP

package httpclient

import (
    "bytes"
    "context"
    "encoding/json"
    "fmt"
    "io"
    "net/http"
    "time"
)

type HTTPClient struct {
    client  *http.Client
    baseURL string
    headers map[string]string
}

func NewHTTPClient(baseURL string, timeout time.Duration) *HTTPClient {
    return &HTTPClient{
        client: &http.Client{
            Timeout: timeout,
            Transport: &http.Transport{
                MaxIdleConns:        100,
                MaxIdleConnsPerHost: 100,
                IdleConnTimeout:     90 * time.Second,
            },
        },
        baseURL: baseURL,
        headers: make(map[string]string),
    }
}

func (c *HTTPClient) Get(ctx context.Context, path string, result any) error {
    return c.doRequest(ctx, http.MethodGet, path, nil, result)
}

func (c *HTTPClient) Post(ctx context.Context, path string, body, result any) error {
    return c.doRequest(ctx, http.MethodPost, path, body, result)
}

func (c *HTTPClient) doRequest(ctx context.Context, method, path string, body, result any) error {
    var reqBody io.Reader
    if body != nil {
        data, err := json.Marshal(body)
        if err != nil {
            return fmt.Errorf("marshal request: %w", err)
        }
        reqBody = bytes.NewReader(data)
    }

    req, err := http.NewRequestWithContext(ctx, method, c.baseURL+path, reqBody)
    if err != nil {
        return fmt.Errorf("create request: %w", err)
    }

    req.Header.Set("Content-Type", "application/json")
    for k, v := range c.headers {
        req.Header.Set(k, v)
    }

    resp, err := c.client.Do(req)
    if err != nil {
        return fmt.Errorf("do request: %w", err)
    }
    defer resp.Body.Close()

    if resp.StatusCode >= 400 {
        respBody, _ := io.ReadAll(resp.Body)
        return fmt.Errorf("unexpected status %d: %s", resp.StatusCode, string(respBody))
    }

    if result != nil {
        if err := json.NewDecoder(resp.Body).Decode(result); err != nil {
            return fmt.Errorf("decode response: %w", err)
        }
    }
    return nil
}

5.2 Communication gRPC

package grpcclient

import (
    "context"
    "crypto/tls"
    "crypto/x509"
    "fmt"
    "time"

    "google.golang.org/grpc"
    "google.golang.org/grpc/credentials"
    "google.golang.org/grpc/keepalive"
)

type GRPCClient struct {
    conn *grpc.ClientConn
}

func NewGRPCClient(address string, opts ...grpc.DialOption) (*GRPCClient, error) {
    defaultOpts := []grpc.DialOption{
        grpc.WithInsecure(),
        grpc.WithDefaultCallOptions(
            grpc.MaxCallRecvMsgSize(10 * 1024 * 1024),
            grpc.MaxCallSendMsgSize(10 * 1024 * 1024),
        ),
        grpc.WithKeepaliveParams(keepalive.ClientParameters{
            Time:                30 * time.Second,
            Timeout:             10 * time.Second,
            PermitWithoutStream: true,
        }),
    }

    opts = append(defaultOpts, opts...)

    conn, err := grpc.Dial(address, opts...)
    if err != nil {
        return nil, fmt.Errorf("dial %s: %w", address, err)
    }
    return &GRPCClient{conn: conn}, nil
}

func (c *GRPCClient) Close() error {
    return c.conn.Close()
}

func WithTLS(caCert, clientCert, clientKey []byte) grpc.DialOption {
    certPool := x509.NewCertPool()
    certPool.AppendCertsFromPEM(caCert)

    cert, err := tls.X509KeyPair(clientCert, clientKey)
    if err != nil {
        return grpc.WithTransportCredentials(nil) // will fail
    }

    creds := credentials.NewTLS(&tls.Config{
        Certificates: []tls.Certificate{cert},
        RootCAs:      certPool,
    })
    return grpc.WithTransportCredentials(creds)
}

6. Communication asynchrone (Events)

6.1 Event Bus

package events

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

type Event struct {
    ID        string    `json:"id"`
    Type      string    `json:"type"`
    Source    string    `json:"source"`
    Time      time.Time `json:"time"`
    Data      any       `json:"data"`
    SchemaURL string    `json:"schema_url,omitempty"`
}

type EventBus interface {
    Publish(ctx context.Context, event Event) error
    Subscribe(ctx context.Context, eventType string, handler EventHandler) error
    Close() error
}

type EventHandler func(context.Context, Event) error

type InMemoryEventBus struct {
    subscribers map[string][]EventHandler
    buffer      chan Event
}

func NewInMemoryEventBus(bufferSize int) *InMemoryEventBus {
    bus := &InMemoryEventBus{
        subscribers: make(map[string][]EventHandler),
        buffer:      make(chan Event, bufferSize),
    }
    go bus.processEvents()
    return bus
}

func (bus *InMemoryEventBus) Publish(ctx context.Context, event Event) error {
    select {
    case bus.buffer <- event:
        return nil
    case <-ctx.Done():
        return ctx.Err()
    }
}

func (bus *InMemoryEventBus) processEvents() {
    for event := range bus.buffer {
        handlers := bus.subscribers[event.Type]
        for _, handler := range handlers {
            go func(h EventHandler, e Event) {
                ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second)
                defer cancel()
                if err := h(ctx, e); err != nil {
                    fmt.Printf("handler failed for event %s: %v\n", e.Type, err)
                }
            }(handler, event)
        }
    }
}

func (bus *InMemoryEventBus) Subscribe(ctx context.Context, eventType string, handler EventHandler) error {
    bus.subscribers[eventType] = append(bus.subscribers[eventType], handler)
    return nil
}

6.2 Kafka Producer/Consumer

package kafka

import (
    "context"
    "fmt"

    "github.com/segmentio/kafka-go"
)

type Producer struct {
    writer *kafka.Writer
}

func NewProducer(brokers []string, topic string) *Producer {
    w := &kafka.Writer{
        Addr:     kafka.TCP(brokers...),
        Topic:    topic,
        Balancer: &kafka.LeastBytes{},
        BatchTimeout: 10,
        Async: true,
    }
    return &Producer{writer: w}
}

func (p *Producer) Publish(ctx context.Context, key string, value []byte) error {
    msg := kafka.Message{
        Key:   []byte(key),
        Value: value,
    }
    if err := p.writer.WriteMessages(ctx, msg); err != nil {
        return fmt.Errorf("kafka write: %w", err)
    }
    return nil
}

func (p *Producer) Close() error {
    return p.writer.Close()
}

type Consumer struct {
    reader *kafka.Reader
}

func NewConsumer(brokers []string, topic, groupID string) *Consumer {
    r := kafka.NewReader(kafka.ReaderConfig{
        Brokers:     brokers,
        Topic:       topic,
        GroupID:     groupID,
        MinBytes:    10e3,
        MaxBytes:    10e6,
        MaxWait:     1,
        StartOffset: kafka.LastOffset,
    })
    return &Consumer{reader: r}
}

func (c *Consumer) Consume(ctx context.Context, handler func(key string, value []byte) error) error {
    for {
        msg, err := c.reader.ReadMessage(ctx)
        if err != nil {
            return fmt.Errorf("read message: %w", err)
        }
        if err := handler(string(msg.Key), msg.Value); err != nil {
            return fmt.Errorf("handler: %w", err)
        }
    }
}

func (c *Consumer) Close() error {
    return c.reader.Close()
}

7. Event Sourcing

7.1 Aggregate et Events

package eventsourcing

import (
    "time"
)

type DomainEvent struct {
    ID        string    `json:"id"`
    AggregateID string  `json:"aggregate_id"`
    Type      string    `json:"type"`
    Data      []byte    `json:"data"`
    Version   int       `json:"version"`
    Timestamp time.Time `json:"timestamp"`
}

type Aggregate interface {
    AggregateID() string
    ApplyEvent(event DomainEvent)
    UncommittedEvents() []DomainEvent
    ClearEvents()
    Version() int
}

type OrderAggregate struct {
    id      string
    status  string
    items   []OrderItem
    total   float64
    version int
    events  []DomainEvent
}

func (a *OrderAggregate) CreateOrder(id string, items []OrderItem) {
    event := DomainEvent{
        AggregateID: id,
        Type:        "OrderCreated",
        Data:        mustMarshal(OrderCreatedData{Items: items}),
        Version:     1,
        Timestamp:   time.Now(),
    }
    a.ApplyEvent(event)
}

func (a *OrderAggregate) ApplyEvent(event DomainEvent) {
    switch event.Type {
    case "OrderCreated":
        var data OrderCreatedData
        mustUnmarshal(event.Data, &data)
        a.status = "created"
        a.items = data.Items
        a.version = event.Version
    case "OrderPaid":
        a.status = "paid"
        a.version = event.Version
    case "OrderShipped":
        a.status = "shipped"
        a.version = event.Version
    }
    a.events = append(a.events, event)
}

type EventStore interface {
    SaveEvents(aggregateID string, events []DomainEvent, expectedVersion int) error
    GetEvents(aggregateID string) ([]DomainEvent, error)
}

type PostgresEventStore struct {
    db *sql.DB
}

func (s *PostgresEventStore) SaveEvents(aggregateID string, events []DomainEvent, expectedVersion int) error {
    tx, err := s.db.Begin()
    if err != nil {
        return err
    }
    defer tx.Rollback()

    for _, event := range events {
        _, err := tx.Exec(
            `INSERT INTO events (aggregate_id, type, data, version, timestamp) 
             VALUES ($1, $2, $3, $4, $5)`,
            aggregateID, event.Type, event.Data, event.Version, event.Timestamp,
        )
        if err != nil {
            return err
        }
    }

    return tx.Commit()
}

8. Saga Pattern

8.1 Choreography Saga

package saga

import (
    "context"
    "fmt"
)

type SagaStep struct {
    Name     string
    Execute  func(context.Context) error
    Compensate func(context.Context) error
}

type Saga struct {
    steps []SagaStep
}

func (s *Saga) Execute(ctx context.Context) error {
    executed := make([]int, 0, len(s.steps))

    for i, step := range s.steps {
        if err := step.Execute(ctx); err != nil {
            // Compensate all executed steps in reverse order
            return s.compensate(ctx, executed)
        }
        executed = append(executed, i)
    }
    return nil
}

func (s *Saga) compensate(ctx context.Context, executed []int) error {
    var errs []error
    for i := len(executed) - 1; i >= 0; i-- {
        step := s.steps[executed[i]]
        if err := step.Compensate(ctx); err != nil {
            errs = append(errs, fmt.Errorf("compensate step %s: %w", step.Name, err))
        }
    }
    if len(errs) > 0 {
        return fmt.Errorf("saga compensation errors: %v", errs)
    }
    return nil
}

// OrderSaga example
func NewOrderSaga(orderService, paymentService, inventoryService ServiceClient) *Saga {
    return &Saga{
        steps: []SagaStep{
            {
                Name: "create-order",
                Execute: func(ctx context.Context) error {
                    return orderService.CreateOrder(ctx)
                },
                Compensate: func(ctx context.Context) error {
                    return orderService.CancelOrder(ctx)
                },
            },
            {
                Name: "reserve-inventory",
                Execute: func(ctx context.Context) error {
                    return inventoryService.Reserve(ctx)
                },
                Compensate: func(ctx context.Context) error {
                    return inventoryService.Release(ctx)
                },
            },
            {
                Name: "process-payment",
                Execute: func(ctx context.Context) error {
                    return paymentService.Charge(ctx)
                },
                Compensate: func(ctx context.Context) error {
                    return paymentService.Refund(ctx)
                },
            },
        },
    }
}

8.2 Orchestrator Pattern

package sagaorchestrator

import (
    "context"
    "fmt"
)

type SagaState int

const (
    SagaPending SagaState = iota
    SagaCompleted
    SagaCompensating
    SagaFailed
)

type Orchestrator struct {
    sagaID     string
    state      SagaState
    steps      []SagaStep
    currentStep int
}

func (o *Orchestrator) ProcessMessage(ctx context.Context, msg Message) error {
    switch msg.Type {
    case "StepCompleted":
        return o.onStepCompleted(ctx, msg)
    case "StepFailed":
        return o.onStepFailed(ctx, msg)
    case "CompensationCompleted":
        return o.onCompensationCompleted(ctx, msg)
    }
    return nil
}

func (o *Orchestrator) onStepCompleted(ctx context.Context, msg Message) error {
    o.currentStep++
    if o.currentStep >= len(o.steps) {
        o.state = SagaCompleted
        return nil
    }
    return o.executeStep(ctx, o.currentStep)
}

func (o *Orchestrator) onStepFailed(ctx context.Context, msg Message) error {
    o.state = SagaCompensating
    return o.compensate(ctx, o.currentStep-1)
}

func (o *Orchestrator) compensate(ctx context.Context, from int) error {
    for i := from; i >= 0; i-- {
        step := o.steps[i]
        if err := step.Compensate(ctx); err != nil {
            return fmt.Errorf("compensation failed at step %d: %w", i, err)
        }
    }
    o.state = SagaFailed
    return nil
}

9. Outbox Pattern

package outbox

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

type OutboxMessage struct {
    ID        string    `json:"id"`
    AggregateType string `json:"aggregate_type"`
    AggregateID   string `json:"aggregate_id"`
    EventType     string `json:"event_type"`
    Payload       []byte `json:"payload"`
    Status        string `json:"status"` // pending, published, failed
    CreatedAt     time.Time `json:"created_at"`
    PublishedAt   *time.Time `json:"published_at,omitempty"`
}

type OutboxRepository interface {
    Insert(ctx context.Context, msg OutboxMessage) error
    GetPending(ctx context.Context, limit int) ([]OutboxMessage, error)
    MarkPublished(ctx context.Context, id string) error
    MarkFailed(ctx context.Context, id string) error
}

type OutboxRelay struct {
    repo    OutboxRepository
    publisher EventPublisher
    ticker  *time.Ticker
}

func NewOutboxRelay(repo OutboxRepository, publisher EventPublisher, interval time.Duration) *OutboxRelay {
    return &OutboxRelay{
        repo:      repo,
        publisher: publisher,
        ticker:    time.NewTicker(interval),
    }
}

func (r *OutboxRelay) Start(ctx context.Context) {
    for {
        select {
        case <-r.ticker.C:
            r.processPending(ctx)
        case <-ctx.Done():
            return
        }
    }
}

func (r *OutboxRelay) processPending(ctx context.Context) {
    messages, err := r.repo.GetPending(ctx, 100)
    if err != nil {
        fmt.Printf("get pending messages: %v\n", err)
        return
    }

    for _, msg := range messages {
        if err := r.publisher.Publish(ctx, msg.EventType, msg.Payload); err != nil {
            fmt.Printf("publish message %s: %v\n", msg.ID, err)
            r.repo.MarkFailed(ctx, msg.ID)
            continue
        }
        r.repo.MarkPublished(ctx, msg.ID)
    }
}

// Usage in service
func (s *OrderService) CreateOrder(ctx context.Context, o Order) error {
    tx, err := s.db.Begin()
    if err != nil {
        return err
    }
    defer tx.Rollback()

    // 1. Save order
    if err := s.orderRepo.Save(ctx, o); err != nil {
        return err
    }

    // 2. Save outbox message in same transaction
    payload, _ := json.Marshal(OrderCreatedEvent{OrderID: o.ID})
    msg := OutboxMessage{
        AggregateID: o.ID,
        EventType:   "order.created",
        Payload:     payload,
        Status:      "pending",
        CreatedAt:   time.Now(),
    }
    if err := s.outboxRepo.Insert(ctx, msg); err != nil {
        return err
    }

    return tx.Commit() // Atomic: both saved or none
}

10. Service Mesh

10.1 Linkerd

# linkerd-profile.yaml
apiVersion: linkerd.io/v1alpha2
kind: ServiceProfile
metadata:
  name: order-service.default.svc.cluster.local
spec:
  routes:
    - name: GET /api/orders
      condition:
        method: GET
        pathRegex: /api/orders.*
      isRetryable: true
      timeout: 500ms
    - name: POST /api/orders
      condition:
        method: POST
        pathRegex: /api/orders
      timeout: 2s
  retryBudget:
    retryRatio: 0.2
    minRetriesPerSecond: 10
    ttl: 10s

10.2 Istio

# istio-virtualservice.yaml
apiVersion: networking.istio.io/v1beta1
kind: VirtualService
metadata:
  name: order-service
spec:
  hosts:
  - order-service
  http:
  - match:
    - uri:
        prefix: /api/orders
    route:
    - destination:
        host: order-service
        subset: v1
      weight: 90
    - destination:
        host: order-service
        subset: v2
      weight: 10
    retries:
      attempts: 3
      perTryTimeout: 500ms
    timeout: 2s
    fault:
      delay:
        percentage:
          value: 0.1
        fixedDelay: 5s
---
apiVersion: networking.istio.io/v1beta1
kind: DestinationRule
metadata:
  name: order-service
spec:
  host: order-service
  trafficPolicy:
    connectionPool:
      tcp:
        maxConnections: 100
      http:
        http1MaxPendingRequests: 10
        http2MaxRequests: 1000
    loadBalancer:
      simple: LEAST_CONN
    outlierDetection:
      consecutive5xxErrors: 5
      interval: 30s
      baseEjectionTime: 30s
  subsets:
  - name: v1
    labels:
      version: v1
  - name: v2
    labels:
      version: v2

11. Observabilité (OpenTelemetry)

11.1 Tracing

package telemetry

import (
    "context"
    "fmt"

    "go.opentelemetry.io/otel"
    "go.opentelemetry.io/otel/attribute"
    "go.opentelemetry.io/otel/exporters/otlp/otlptrace"
    "go.opentelemetry.io/otel/sdk/resource"
    sdktrace "go.opentelemetry.io/otel/sdk/trace"
    semconv "go.opentelemetry.io/otel/semconv/v1.17.0"
    "go.opentelemetry.io/otel/trace"
)

var tracer trace.Tracer

func InitTracing(serviceName, endpoint string) (*sdktrace.TracerProvider, error) {
    exporter, err := otlptrace.New(context.Background(),
        otlptrace.WithEndpoint(endpoint),
        otlptrace.WithInsecure(),
    )
    if err != nil {
        return nil, fmt.Errorf("create exporter: %w", err)
    }

    tp := sdktrace.NewTracerProvider(
        sdktrace.WithBatcher(exporter),
        sdktrace.WithResource(resource.NewWithAttributes(
            semconv.SchemaURL,
            semconv.ServiceName(serviceName),
            attribute.String("environment", "production"),
        )),
    )
    otel.SetTracerProvider(tp)
    tracer = tp.Tracer(serviceName)
    return tp, nil
}

// Span helpers
func StartSpan(ctx context.Context, name string, opts ...trace.SpanStartOption) (context.Context, trace.Span) {
    return tracer.Start(ctx, name, opts...)
}

// HTTP middleware
func TracingMiddleware(serviceName string) func(http.Handler) http.Handler {
    return func(next http.Handler) http.Handler {
        return http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
            ctx, span := StartSpan(r.Context(), fmt.Sprintf("%s %s", r.Method, r.URL.Path))
            defer span.End()

            span.SetAttributes(
                attribute.String("http.method", r.Method),
                attribute.String("http.url", r.URL.String()),
                attribute.String("http.host", r.Host),
            )
            next.ServeHTTP(w, r.WithContext(ctx))
        })
    }
}

11.2 Metrics

package metrics

import (
    "net/http"
    "time"

    "github.com/prometheus/client_golang/prometheus"
    "github.com/prometheus/client_golang/prometheus/promhttp"
)

type Metrics struct {
    requestDuration *prometheus.HistogramVec
    requestTotal    *prometheus.CounterVec
    activeRequests  *prometheus.GaugeVec
    errorsTotal     *prometheus.CounterVec
}

func NewMetrics(namespace, subsystem string) *Metrics {
    m := &Metrics{
        requestDuration: prometheus.NewHistogramVec(
            prometheus.HistogramOpts{
                Namespace: namespace,
                Subsystem: subsystem,
                Name:      "request_duration_seconds",
                Help:      "Request duration in seconds",
                Buckets:   prometheus.DefBuckets,
            },
            []string{"method", "path", "status"},
        ),
        requestTotal: prometheus.NewCounterVec(
            prometheus.CounterOpts{
                Namespace: namespace,
                Subsystem: subsystem,
                Name:      "requests_total",
                Help:      "Total number of requests",
            },
            []string{"method", "path", "status"},
        ),
        activeRequests: prometheus.NewGaugeVec(
            prometheus.GaugeOpts{
                Namespace: namespace,
                Subsystem: subsystem,
                Name:      "active_requests",
                Help:      "Active requests",
            },
            []string{"method"},
        ),
        errorsTotal: prometheus.NewCounterVec(
            prometheus.CounterOpts{
                Namespace: namespace,
                Subsystem: subsystem,
                Name:      "errors_total",
                Help:      "Total number of errors",
            },
            []string{"type"},
        ),
    }

    prometheus.MustRegister(
        m.requestDuration,
        m.requestTotal,
        m.activeRequests,
        m.errorsTotal,
    )
    return m
}

func (m *Metrics) Handler() http.Handler {
    return promhttp.Handler()
}

func (m *Metrics) ObserveRequest(method, path string, status int, duration time.Duration) {
    m.requestDuration.WithLabelValues(method, path, fmt.Sprintf("%d", status)).Observe(duration.Seconds())
    m.requestTotal.WithLabelValues(method, path, fmt.Sprintf("%d", status)).Inc()
}

12. Health Checks

12.1 Health Check Handler

package health

import (
    "context"
    "encoding/json"
    "net/http"
    "time"
)

type Checker interface {
    Name() string
    Check(ctx context.Context) error
}

type HealthHandler struct {
    checkers []Checker
}

func NewHealthHandler(checkers ...Checker) *HealthHandler {
    return &HealthHandler{checkers: checkers}
}

func (h *HealthHandler) Liveness(w http.ResponseWriter, r *http.Request) {
    w.WriteHeader(http.StatusOK)
    json.NewEncoder(w).Encode(map[string]string{"status": "alive"})
}

func (h *HealthHandler) Readiness(w http.ResponseWriter, r *http.Request) {
    ctx, cancel := context.WithTimeout(r.Context(), 5*time.Second)
    defer cancel()

    status := http.StatusOK
    result := make(map[string]string)

    for _, checker := range h.checkers {
        if err := checker.Check(ctx); err != nil {
            result[checker.Name()] = fmt.Sprintf("unhealthy: %v", err)
            status = http.StatusServiceUnavailable
        } else {
            result[checker.Name()] = "healthy"
        }
    }

    w.WriteHeader(status)
    json.NewEncoder(w).Encode(result)
}

// Database checker
type DBChecker struct {
    db *sql.DB
}

func (c *DBChecker) Name() string { return "database" }

func (c *DBChecker) Check(ctx context.Context) error {
    return c.db.PingContext(ctx)
}

// Redis checker
type RedisChecker struct {
    client *redis.Client
}

func (c *RedisChecker) Name() string { return "redis" }

func (c *RedisChecker) Check(ctx context.Context) error {
    return c.client.Ping(ctx).Err()
}

// Kafka checker
type KafkaChecker struct {
    brokers []string
}

func (c *KafkaChecker) Name() string { return "kafka" }

func (c *KafkaChecker) Check(ctx context.Context) error {
    conn, err := kafka.Dial("tcp", c.brokers[0])
    if err != nil {
        return err
    }
    return conn.Close()
}

12.2 Graceful Degradation

package degradation

import (
    "context"
    "net/http"
)

type DegradationState struct {
    CacheAvailable  bool
    DatabaseAvailable bool
    KafkaAvailable  bool
    mode            string // "full", "cache_only", "read_only"
}

type DegradableHandler struct {
    state    *DegradationState
    fallback http.Handler
}

func (h *DegradableHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
    switch h.state.mode {
    case "full":
        // Normal operation
        h.fallback.ServeHTTP(w, r)
    case "cache_only":
        // Serve stale cache data
        w.Header().Set("X-Cache", "HIT")
        w.Header().Set("X-Degradation-Mode", "cache_only")
        h.serveFromCache(w, r)
    case "read_only":
        // Block writes, allow reads
        if r.Method != http.MethodGet && r.Method != http.MethodHead {
            http.Error(w, "service in read-only mode", http.StatusServiceUnavailable)
            return
        }
        h.fallback.ServeHTTP(w, r)
    }
}

// Graceful shutdown
func GracefulShutdown(srv *http.Server, timeout time.Duration) {
    ctx, cancel := context.WithTimeout(context.Background(), timeout)
    defer cancel()

    if err := srv.Shutdown(ctx); err != nil {
        log.Printf("server shutdown: %v", err)
    }
}

Diagrammes Mermaid

Architecture microservices

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

Saga Pattern

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

Outbox Pattern

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

Circuit Breaker States

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

Points Clés à Retenir

  1. Microservices ≠ petits monolithes — Chaque service doit avoir une responsabilité unique
  2. Communication — Préférer async (events) pour le couplage faible, sync (gRPC) pour le temps réel
  3. Résilience — Toujours implémenter circuit breaker, retry, timeout
  4. Observabilité — Metrics, tracing, logging sont obligatoires, pas optionnels
  5. Saga — Pour les transactions distribuées, toujours prévoir la compensation
  6. Outbox — Garantit la delivery des événements sans 2PC
  7. Health Checks — Liveness (process) ≠ Readiness (dépendances)
  8. Graceful Degradation — Un service doit continuer à fonctionner partiellement

Références

  • "Building Microservices" by Sam Newman
  • "Microservices Patterns" by Chris Richardson
  • Go kit microservices toolkit
  • OpenTelemetry Go SDK documentation