MFormations
Modern Go Engineering

Chapitre 15

15 - Cloud avec Go

15 - Cloud avec Go

Cours 15 : Cloud avec Go

1. Docker Multi-stage pour Go

1.1 Dockerfile optimisé

# Stage 1: Build
FROM golang:1.22-alpine AS builder
RUN apk add --no-cache git ca-certificates
WORKDIR /app
COPY go.mod go.sum ./
RUN go mod download
COPY . .
RUN CGO_ENABLED=0 GOOS=linux GOARCH=amd64 \
    go build -ldflags="-s -w" -o /app/service ./cmd/server

# Stage 2: Runtime (scratch)
FROM scratch
COPY --from=builder /etc/ssl/certs/ca-certificates.crt /etc/ssl/certs/
COPY --from=builder /app/service /service
EXPOSE 8080
ENTRYPOINT ["/service"]

# Alternative : distroless
FROM gcr.io/distroless/static-debian12:nonroot
COPY --from=builder /app/service /service
USER nonroot:nonroot
ENTRYPOINT ["/service"]

# Alternative : alpine (avec shell)
FROM alpine:3.19
RUN apk add --no-cache ca-certificates tzdata
COPY --from=builder /app/service /service
EXPOSE 8080
ENTRYPOINT ["/service"]

1.2 Docker compose

version: '3.8'
services:
  app:
    build:
      context: .
      dockerfile: Dockerfile
    ports:
      - "8080:8080"
    environment:
      - DB_HOST=postgres
      - REDIS_HOST=redis
      - KAFKA_BROKERS=kafka:9092
    depends_on:
      postgres:
        condition: service_healthy
      redis:
        condition: service_started
    healthcheck:
      test: ["CMD", "/service", "health"]
      interval: 30s
      timeout: 10s
      retries: 3
    deploy:
      resources:
        limits:
          cpus: '0.5'
          memory: 256M

  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_DB: myapp
      POSTGRES_USER: app
      POSTGRES_PASSWORD: secret
    volumes:
      - pgdata:/var/lib/postgresql/data
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U app"]
      interval: 5s

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

volumes:
  pgdata:
  redisdata:

2. AWS SDK v2

2.1 S3 (Stockage)

package s3

import (
    "context"
    "fmt"
    "io"

    "github.com/aws/aws-sdk-go-v2/aws"
    "github.com/aws/aws-sdk-go-v2/config"
    "github.com/aws/aws-sdk-go-v2/service/s3"
)

type S3Client struct {
    client *s3.Client
    bucket string
}

func NewS3Client(ctx context.Context, bucket, region string) (*S3Client, error) {
    cfg, err := config.LoadDefaultConfig(ctx, config.WithRegion(region))
    if err != nil {
        return nil, fmt.Errorf("load config: %w", err)
    }

    return &S3Client{
        client: s3.NewFromConfig(cfg),
        bucket: bucket,
    }, nil
}

func (c *S3Client) Upload(ctx context.Context, key string, body io.Reader) error {
    _, err := c.client.PutObject(ctx, &s3.PutObjectInput{
        Bucket: aws.String(c.bucket),
        Key:    aws.String(key),
        Body:   body,
    })
    return err
}

func (c *S3Client) Download(ctx context.Context, key string) (io.ReadCloser, error) {
    output, err := c.client.GetObject(ctx, &s3.GetObjectInput{
        Bucket: aws.String(c.bucket),
        Key:    aws.String(key),
    })
    if err != nil {
        return nil, err
    }
    return output.Body, nil
}

func (c *S3Client) Delete(ctx context.Context, key string) error {
    _, err := c.client.DeleteObject(ctx, &s3.DeleteObjectInput{
        Bucket: aws.String(c.bucket),
        Key:    aws.String(key),
    })
    return err
}

func (c *S3Client) List(ctx context.Context, prefix string) ([]string, error) {
    var keys []string
    paginator := s3.NewListObjectsV2Paginator(c.client, &s3.ListObjectsV2Input{
        Bucket: aws.String(c.bucket),
        Prefix: aws.String(prefix),
    })

    for paginator.HasMorePages() {
        page, err := paginator.NextPage(ctx)
        if err != nil {
            return nil, err
        }
        for _, obj := range page.Contents {
            keys = append(keys, *obj.Key)
        }
    }
    return keys, nil
}

2.2 DynamoDB (NoSQL)

package dynamodb

import (
    "context"
    "fmt"

    "github.com/aws/aws-sdk-go-v2/aws"
    "github.com/aws/aws-sdk-go-v2/config"
    "github.com/aws/aws-sdk-go-v2/feature/dynamodb/attributevalue"
    "github.com/aws/aws-sdk-go-v2/feature/dynamodb/expression"
    "github.com/aws/aws-sdk-go-v2/service/dynamodb"
)

type User struct {
    ID    string   `dynamodbav:"id"`
    Name  string   `dynamodbav:"name"`
    Email string   `dynamodbav:"email"`
    Roles []string `dynamodbav:"roles"`
}

type DynamoDBClient struct {
    client *dynamodb.Client
    table  string
}

func New(ctx context.Context, table, region string) (*DynamoDBClient, error) {
    cfg, err := config.LoadDefaultConfig(ctx, config.WithRegion(region))
    if err != nil {
        return nil, err
    }
    return &DynamoDBClient{
        client: dynamodb.NewFromConfig(cfg),
        table:  table,
    }, nil
}

func (d *DynamoDBClient) Put(ctx context.Context, user User) error {
    item, err := attributevalue.MarshalMap(user)
    if err != nil {
        return err
    }
    _, err = d.client.PutItem(ctx, &dynamodb.PutItemInput{
        TableName: aws.String(d.table),
        Item:      item,
    })
    return err
}

func (d *DynamoDBClient) Get(ctx context.Context, id string) (*User, error) {
    result, err := d.client.GetItem(ctx, &dynamodb.GetItemInput{
        TableName: aws.String(d.table),
        Key: map[string]types.AttributeValue{
            "id": &types.AttributeValueMemberS{Value: id},
        },
    })
    if err != nil {
        return nil, err
    }
    if result.Item == nil {
        return nil, fmt.Errorf("not found")
    }

    var user User
    if err := attributevalue.UnmarshalMap(result.Item, &user); err != nil {
        return nil, err
    }
    return &user, nil
}

func (d *DynamoDBClient) Query(ctx context.Context, email string) ([]User, error) {
    expr, err := expression.NewBuilder().
        WithKeyCondition(expression.KeyEqual(expression.Key("email"), expression.Value(email))).
        Build()
    if err != nil {
        return nil, err
    }

    result, err := d.client.Query(ctx, &dynamodb.QueryInput{
        TableName:                 aws.String(d.table),
        KeyConditionExpression:    expr.KeyCondition(),
        ExpressionAttributeNames:  expr.Names(),
        ExpressionAttributeValues: expr.Values(),
    })
    if err != nil {
        return nil, err
    }

    var users []User
    for _, item := range result.Items {
        var user User
        attributevalue.UnmarshalMap(item, &user)
        users = append(users, user)
    }
    return users, nil
}

2.3 SQS (Queues)

package sqs

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

    "github.com/aws/aws-sdk-go-v2/aws"
    "github.com/aws/aws-sdk-go-v2/config"
    "github.com/aws/aws-sdk-go-v2/service/sqs"
    "github.com/aws/aws-sdk-go-v2/service/sqs/types"
)

type SQSClient struct {
    client *sqs.Client
    queueURL string
}

func New(ctx context.Context, queueURL, region string) (*SQSClient, error) {
    cfg, err := config.LoadDefaultConfig(ctx, config.WithRegion(region))
    if err != nil {
        return nil, err
    }
    return &SQSClient{
        client:   sqs.NewFromConfig(cfg),
        queueURL: queueURL,
    }, nil
}

func (c *SQSClient) SendMessage(ctx context.Context, msg any) error {
    body, err := json.Marshal(msg)
    if err != nil {
        return err
    }
    _, err = c.client.SendMessage(ctx, &sqs.SendMessageInput{
        QueueUrl:    aws.String(c.queueURL),
        MessageBody: aws.String(string(body)),
    })
    return err
}

func (c *SQSClient) ReceiveMessages(ctx context.Context, maxMessages int32) ([]types.Message, error) {
    result, err := c.client.ReceiveMessage(ctx, &sqs.ReceiveMessageInput{
        QueueUrl:            aws.String(c.queueURL),
        MaxNumberOfMessages: maxMessages,
        WaitTimeSeconds:     20, // Long polling
    })
    if err != nil {
        return nil, err
    }
    return result.Messages, nil
}

func (c *SQSClient) DeleteMessage(ctx context.Context, receiptHandle string) error {
    _, err := c.client.DeleteMessage(ctx, &sqs.DeleteMessageInput{
        QueueUrl:      aws.String(c.queueURL),
        ReceiptHandle: aws.String(receiptHandle),
    })
    return err
}

2.4 Lambda

package lambda

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

    "github.com/aws/aws-lambda-go/events"
    "github.com/aws/aws-lambda-go/lambda"
    "github.com/aws/aws-sdk-go-v2/service/dynamodb"
)

type Event struct {
    UserID string `json:"user_id"`
    Action string `json:"action"`
}

type Handler struct {
    db *dynamodb.Client
}

func NewHandler() *Handler {
    return &Handler{}
}

func (h *Handler) HandleRequest(ctx context.Context, event json.RawMessage) (string, error) {
    var e Event
    if err := json.Unmarshal(event, &e); err != nil {
        return "", err
    }

    log.Printf("Processing event: user=%s action=%s", e.UserID, e.Action)
    return "OK", nil
}

// API Gateway event
func (h *Handler) HandleAPIGateway(ctx context.Context, event events.APIGatewayProxyRequest) (events.APIGatewayProxyResponse, error) {
    return events.APIGatewayProxyResponse{
        StatusCode: 200,
        Headers:    map[string]string{"Content-Type": "application/json"},
        Body:       `{"message": "Hello from Lambda"}`,
    }, nil
}

// SQS event
func (h *Handler) HandleSQS(ctx context.Context, event events.SQSEvent) error {
    for _, msg := range event.Records {
        log.Printf("SQS message: %s", msg.Body)
    }
    return nil
}

func main() {
    handler := NewHandler()
    lambda.Start(handler.HandleRequest)
}

3. Kubernetes avec client-go

3.1 Informers et Watchers

package k8s

import (
    "context"
    "fmt"
    "time"

    appsv1 "k8s.io/api/apps/v1"
    corev1 "k8s.io/api/core/v1"
    metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
    "k8s.io/client-go/informers"
    "k8s.io/client-go/kubernetes"
    "k8s.io/client-go/rest"
    "k8s.io/client-go/tools/cache"
    "k8s.io/client-go/tools/clientcmd"
)

type K8sClient struct {
    clientset *kubernetes.Clientset
}

func NewInClusterClient() (*K8sClient, error) {
    config, err := rest.InClusterConfig()
    if err != nil {
        return nil, fmt.Errorf("in-cluster config: %w", err)
    }
    clientset, err := kubernetes.NewForConfig(config)
    if err != nil {
        return nil, err
    }
    return &K8sClient{clientset: clientset}, nil
}

func NewOutOfClusterClient(kubeconfig string) (*K8sClient, error) {
    config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
    if err != nil {
        return nil, err
    }
    clientset, err := kubernetes.NewForConfig(config)
    if err != nil {
        return nil, err
    }
    return &K8sClient{clientset: clientset}, nil
}

// Deployment operations
func (c *K8sClient) CreateDeployment(ctx context.Context, deploy *appsv1.Deployment) error {
    _, err := c.clientset.AppsV1().Deployments(deploy.Namespace).Create(ctx, deploy, metav1.CreateOptions{})
    return err
}

func (c *K8sClient) GetDeployment(ctx context.Context, namespace, name string) (*appsv1.Deployment, error) {
    return c.clientset.AppsV1().Deployments(namespace).Get(ctx, name, metav1.GetOptions{})
}

func (c *K8sClient) UpdateDeployment(ctx context.Context, deploy *appsv1.Deployment) error {
    _, err := c.clientset.AppsV1().Deployments(deploy.Namespace).Update(ctx, deploy, metav1.UpdateOptions{})
    return err
}

func (c *K8sClient) DeleteDeployment(ctx context.Context, namespace, name string) error {
    return c.clientset.AppsV1().Deployments(namespace).Delete(ctx, name, metav1.DeleteOptions{})
}

// Pod operations
func (c *K8sClient) ListPods(ctx context.Context, namespace string, labelSelector string) (*corev1.PodList, error) {
    return c.clientset.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{
        LabelSelector: labelSelector,
    })
}

func (c *K8sClient) GetPodLogs(ctx context.Context, namespace, name string, tailLines int64) (string, error) {
    req := c.clientset.CoreV1().Pods(namespace).GetLogs(name, &corev1.PodLogOptions{
        TailLines: &tailLines,
    })
    data, err := req.DoRaw(ctx)
    if err != nil {
        return "", err
    }
    return string(data), nil
}

// Informer pattern (watch)
func (c *K8sClient) WatchPods(ctx context.Context, namespace string, handlers cache.ResourceEventHandler) {
    factory := informers.NewSharedInformerFactoryWithOptions(
        c.clientset,
        10*time.Minute,
        informers.WithNamespace(namespace),
    )
    
    informer := factory.Core().V1().Pods().Informer()
    informer.AddEventHandler(handlers)
    
    stop := make(chan struct{})
    defer close(stop)
    informer.Run(stop)
}

3.2 Controller-runtime (Operators)

package controllers

import (
    "context"
    "time"

    corev1 "k8s.io/api/core/v1"
    "k8s.io/apimachinery/pkg/runtime"
    ctrl "sigs.k8s.io/controller-runtime"
    "sigs.k8s.io/controller-runtime/pkg/client"
    "sigs.k8s.io/controller-runtime/pkg/log"
    "sigs.k8s.io/controller-runtime/pkg/reconcile"
)

type PodReconciler struct {
    client.Client
    Scheme *runtime.Scheme
}

func (r *PodReconciler) Reconcile(ctx context.Context, req ctrl.Request) (ctrl.Result, error) {
    log := log.FromContext(ctx)

    var pod corev1.Pod
    if err := r.Get(ctx, req.NamespacedName, &pod); err != nil {
        return ctrl.Result{}, client.IgnoreNotFound(err)
    }

    log.Info("Reconciling Pod", "name", pod.Name, "phase", pod.Status.Phase)

    // Handle pod phases
    switch pod.Status.Phase {
    case corev1.PodPending:
        log.Info("Pod is pending")
    case corev1.PodRunning:
        log.Info("Pod is running")
    case corev1.PodFailed:
        log.Info("Pod failed", "reason", pod.Status.Reason)
    }

    return ctrl.Result{RequeueAfter: 30 * time.Second}, nil
}

func (r *PodReconciler) SetupWithManager(mgr ctrl.Manager) error {
    return ctrl.NewControllerManagedBy(mgr).
        For(&corev1.Pod{}).
        Complete(r)
}

4. CloudEvents

package cloudevents

import (
    "context"
    "time"

    ce "github.com/cloudevents/sdk-go/v2"
    "github.com/cloudevents/sdk-go/v2/client"
    "github.com/cloudevents/sdk-go/v2/protocol/http"
)

type EventPublisher struct {
    client client.Client
}

func NewEventPublisher() (*EventPublisher, error) {
    c, err := ce.NewClientHTTP()
    if err != nil {
        return nil, err
    }
    return &EventPublisher{client: c}, nil
}

func (p *EventPublisher) Publish(ctx context.Context, eventType, source string, data any) error {
    event := ce.NewEvent()
    event.SetID(GenerateID())
    event.SetType(eventType)
    event.SetSource(source)
    event.SetTime(time.Now())
    event.SetDataContentType("application/json")

    if err := event.SetData("application/json", data); err != nil {
        return err
    }

    ctx = ce.ContextWithTarget(ctx, "http://event-broker:8080")
    return p.client.Send(ctx, event)
}

// Consumer
type EventConsumer struct {
    client client.Client
}

func NewEventConsumer() (*EventConsumer, error) {
    c, err := ce.NewClientHTTP(ce.WithPort(8080))
    if err != nil {
        return nil, err
    }
    return &EventConsumer{client: c}, nil
}

func (c *EventConsumer) Start(ctx context.Context, handler func(ce.Event)) {
    c.client.StartReceiver(ctx, func(event ce.Event) {
        handler(event)
    })
}

5. Diagrammes

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

Points Clés

  1. Docker multi-stage : scratch pour taille minimale, distroless pour sécurité
  2. AWS SDK v2 : modular, context-aware, paginators
  3. client-go : informers pour watch efficace, controller-runtime pour operators
  4. Lambda : handler simple pour API Gateway, SQS, S3 events
  5. CloudEvents : standard CNCF pour events cloud
  6. K8s operators : controller-runtime pattern