Мини-проект: CQRS-система с Command → Read Model

Мини-проект: CQRS-система на Go

Соберём полный CQRS-поток: Command → запись → Read Model. Никаких брокеров - чистое разделение команд и запросов с оптимизированными read-проекциями.

Если нужны брокеры - это трек Message Brokers: там RabbitMQ + Kafka + Outbox паттерн.

Архитектура

HTTP POST /api/lessons/:id/complete
    │
    ▼
CompleteLesson (Command Handler)
    │
    ├─ INSERT completions (write model)
    └─ UPDATE user_progress (read model - в одной транзакции)

HTTP GET /api/users/:id/progress
    │
    ▼
GetUserProgress (Query Handler)
    │
    └─ SELECT user_progress (read model - денормализован, без JOIN)

Ключевое: write model - таблица completions (иммутабельный журнал). Read model - user_progress (денормализованный, быстрый SELECT).

Структура проекта

cmd/
  api/main.go
internal/
  domain/
    lesson.go        -- LessonCompletion struct, интерфейсы
  usecase/
    complete_lesson.go
    get_progress.go
  repository/
    completion_repo.go
    progress_repo.go
  handler/
    lesson_handler.go
  config/
    wire.go          -- DI через Wire

Доменные типы

// internal/domain/lesson.go

type LessonCompletion struct {
    ID        uuid.UUID
    UserID    uuid.UUID
    LessonID  string
    TrackID   string
    CompletedAt time.Time
}

type UserProgress struct {
    UserID           uuid.UUID
    CompletedLessons int
    LastActivityAt   time.Time
    ProgressByTrack  map[string]int
}

type CompletionRepository interface {
    Save(ctx context.Context, c *LessonCompletion) error
    FindByUser(ctx context.Context, userID uuid.UUID) ([]*LessonCompletion, error)
}

type ProgressRepository interface {
    Update(ctx context.Context, userID uuid.UUID, trackID string) error
    Get(ctx context.Context, userID uuid.UUID) (*UserProgress, error)
}

Command: CompleteLesson

// internal/usecase/complete_lesson.go

type CompleteLessonUseCase struct {
    completions CompletionRepository
    progress    ProgressRepository
    db          *sql.DB
}

type CompleteLessonInput struct {
    UserID   uuid.UUID
    LessonID string
    TrackID  string
}

func (uc *CompleteLessonUseCase) Execute(ctx context.Context, in CompleteLessonInput) error {
    tx, err := uc.db.BeginTx(ctx, nil)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback()

    completion := &domain.LessonCompletion{
        ID:          uuid.New(),
        UserID:      in.UserID,
        LessonID:    in.LessonID,
        TrackID:     in.TrackID,
        CompletedAt: time.Now(),
    }

    // Write model: записать факт прохождения
    if err := uc.completions.SaveTx(ctx, tx, completion); err != nil {
        return fmt.Errorf("save completion: %w", err)
    }

    // Read model: обновить денормализованную проекцию
    if err := uc.progress.UpdateTx(ctx, tx, in.UserID, in.TrackID); err != nil {
        return fmt.Errorf("update progress: %w", err)
    }

    return tx.Commit()
}

Оба INSERT в одной транзакции - read model всегда консистентен с write model.

Query: GetUserProgress

// internal/usecase/get_progress.go

type GetProgressUseCase struct {
    progress ProgressRepository
}

func (uc *GetProgressUseCase) Execute(ctx context.Context, userID uuid.UUID) (*domain.UserProgress, error) {
    p, err := uc.progress.Get(ctx, userID)
    if err != nil {
        return nil, fmt.Errorf("get progress: %w", err)
    }
    return p, nil
}

Query не трогает completions - читает только из денормализованного user_progress. Никаких JOIN, никаких агрегаций в реальном времени.

Схема БД

-- Write model: иммутабельный журнал
CREATE TABLE completions (
    id           UUID PRIMARY KEY,
    user_id      UUID NOT NULL,
    lesson_id    TEXT NOT NULL,
    track_id     TEXT NOT NULL,
    completed_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
    UNIQUE (user_id, lesson_id)
);

-- Read model: денормализованная проекция
CREATE TABLE user_progress (
    user_id            UUID PRIMARY KEY,
    completed_lessons  INT NOT NULL DEFAULT 0,
    last_activity_at   TIMESTAMPTZ,
    progress_by_track  JSONB NOT NULL DEFAULT '{}'
);

CREATE INDEX ON completions (user_id);
CREATE INDEX ON user_progress (user_id);

UNIQUE (user_id, lesson_id) в completions - встроенная защита от повторного прохождения.

HTTP Handler

// internal/handler/lesson_handler.go

func (h *LessonHandler) CompleteLesson(w http.ResponseWriter, r *http.Request) {
    userID := middleware.UserIDFromContext(r.Context())
    lessonID := chi.URLParam(r, "lessonID")
    trackID := chi.URLParam(r, "trackID")

    err := h.completeLesson.Execute(r.Context(), usecase.CompleteLessonInput{
        UserID:   userID,
        LessonID: lessonID,
        TrackID:  trackID,
    })
    if err != nil {
        http.Error(w, "internal error", http.StatusInternalServerError)
        return
    }

    w.WriteHeader(http.StatusNoContent)
}

func (h *LessonHandler) GetProgress(w http.ResponseWriter, r *http.Request) {
    userID := middleware.UserIDFromContext(r.Context())

    progress, err := h.getProgress.Execute(r.Context(), userID)
    if err != nil {
        http.Error(w, "not found", http.StatusNotFound)
        return
    }

    json.NewEncoder(w).Encode(progress)
}

Docker Compose для разработки

services:
  postgres:
    image: postgres:16-alpine
    environment:
      POSTGRES_DB: cqrs_demo
      POSTGRES_USER: app
      POSTGRES_PASSWORD: secret
    ports:
      - "5432:5432"
    healthcheck:
      test: pg_isready -U app
      interval: 5s
      retries: 5

  api:
    build: .
    depends_on:
      postgres:
        condition: service_healthy
    environment:
      DATABASE_URL: postgres://app:secret@postgres:5432/cqrs_demo
    ports:
      - "8080:8080"

Только postgres и api - никаких брокеров. Это и есть чистый CQRS без event-driven.

Тест use case

func TestCompleteLessonUseCase(t *testing.T) {
    db := testdb.Open(t) // тестовая БД из testcontainers
    repo := repository.NewCompletionRepo(db)
    progressRepo := repository.NewProgressRepo(db)
    uc := usecase.NewCompleteLessonUseCase(repo, progressRepo, db)

    userID := uuid.New()
    err := uc.Execute(context.Background(), usecase.CompleteLessonInput{
        UserID:   userID,
        LessonID: "go-01",
        TrackID:  "go",
    })
    require.NoError(t, err)

    // Проверяем write model
    completions, _ := repo.FindByUser(context.Background(), userID)
    assert.Len(t, completions, 1)

    // Проверяем read model
    progress, _ := progressRepo.Get(context.Background(), userID)
    assert.Equal(t, 1, progress.CompletedLessons)
}

Типичная ошибка

Обновлять read model асинхронно (через фоновую горутину или брокер), когда в этом нет необходимости. Если read model обновляется в той же транзакции - он всегда консистентен, нет eventual consistency, нет задержек. Брокер нужен, когда событие должны обработать несколько систем - тогда смотри Message Brokers.

Мини-задание

  • Реализуй CompletionRepository.SaveTx и ProgressRepository.UpdateTx принимающие *sql.Tx
  • Добавь UNIQUE constraint и проверь что повторный вызов возвращает специфичную ошибку (не 500)
  • Напиши интеграционный тест с testcontainers: PostgreSQL без моков
  • Добавь endpoint GET /api/users/:id/completions возвращающий полный список пройденных уроков из write model

Зарегистрируйтесь бесплатно, чтобы пройти квиз, решить задание с автопроверкой и вести прогресс.