Мини-проект: события в BackendStart (LessonCompleted → Achievements/Email/Analytics)

Мини-проект: события в BackendStart (LessonCompleted → Achievements/Email/Analytics)

В этом уроке мы соберём всё, что изучили за трек, в работающую систему. Возьмём реальный сценарий BackendStart: пользователь завершает урок, и это порождает цепочку побочных эффектов - проверку ачивок, запись аналитики, отправку уведомления. Каждый эффект реализован отдельным подписчиком. Ни один из них не знает о существовании остальных. Добавление нового подписчика - это один файл и одна строка регистрации.

Мы реализуем: доменное событие, EventBus (in-process), outbox-таблицу, три подписчика, wiring в main.go и unit-тесты.

Архитектура мини-проекта: use-case CompleteLesson пишет данные и outbox в одной транзакции; poller досылает в EventBus; три подписчика - Achievement, Analytics, Notification - читают параллельно

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

backend/
├── internal/
│   ├── domain/
│   │   └── events.go              # Доменные события
│   ├── eventbus/
│   │   ├── bus.go                  # EventBus interface + in-process реализация
│   │   └── bus_test.go             # Тесты EventBus
│   ├── outbox/
│   │   ├── repository.go           # Outbox-таблица (запись + чтение)
│   │   └── poller.go               # Outbox poller (горутина-отправитель)
│   ├── usecase/
│   │   └── complete_lesson.go      # Use-case: завершение урока
│   └── subscribers/
│       ├── achievement_checker.go   # Подписчик: проверка ачивок
│       ├── analytics_recorder.go    # Подписчик: запись аналитики
│       └── notification_sender.go   # Подписчик: отправка уведомления
└── cmd/
    └── server/
        └── main.go                  # Wiring: bus + subscribers + use-case

Domain Event: LessonCompleted

Доменное событие - это факт, который произошёл в системе. Неизменяемая структура с временной меткой:

// domain/events.go

// LessonCompleted - доменное событие: пользователь завершил урок.
type LessonCompleted struct {
    EventID       string    `json:"event_id"`
    Type          string    `json:"type"`            // всегда "lesson.completed"
    CorrelationID string    `json:"correlation_id"`
    OccurredAt    time.Time `json:"occurred_at"`
    UserID        int64     `json:"user_id"`
    LessonID      int64     `json:"lesson_id"`
    TrackSlug     string    `json:"track_slug"`
    LessonSlug    string    `json:"lesson_slug"`
}

// NewLessonCompleted создаёт событие с автоматическим event_id и timestamp.
func NewLessonCompleted(correlationID string, userID, lessonID int64, trackSlug, lessonSlug string) LessonCompleted {
    return LessonCompleted{
        EventID:       ulid.Make().String(),
        Type:          "lesson.completed",
        CorrelationID: correlationID,
        OccurredAt:    time.Now().UTC(),
        UserID:        userID,
        LessonID:      lessonID,
        TrackSlug:     trackSlug,
        LessonSlug:    lessonSlug,
    }
}
<?php
// src/Domain/Event/LessonCompleted.php
declare(strict_types=1);

namespace App\Domain\Event;

use Symfony\Component\Uid\Uuid;

// final readonly - событие immutable по определению.
final readonly class LessonCompleted implements DomainEvent
{
    public function __construct(
        public string $eventId,
        public string $correlationId,
        public \DateTimeImmutable $occurredAt,
        public int $userId,
        public int $lessonId,
        public string $trackSlug,
        public string $lessonSlug,
    ) {}

    public static function create(
        string $correlationId,
        int $userId,
        int $lessonId,
        string $trackSlug,
        string $lessonSlug,
    ): self {
        return new self(
            eventId: Uuid::v4()->toRfc4122(),
            correlationId: $correlationId,
            occurredAt: new \DateTimeImmutable('now', new \DateTimeZone('UTC')),
            userId: $userId,
            lessonId: $lessonId,
            trackSlug: $trackSlug,
            lessonSlug: $lessonSlug,
        );
    }

    public function eventType(): string
    {
        return 'lesson.completed';
    }

    public function occurredAt(): \DateTimeImmutable
    {
        return $this->occurredAt;
    }
}

Поле Type - строковая константа. Подписчики фильтруют события по этому полю. CorrelationID приходит из HTTP-middleware (см. урок 09).

EventBus: интерфейс и in-process реализация

EventBus - это порт (interface) в терминах Hexagonal Architecture. Use-case зависит от интерфейса, а не от конкретной реализации:

// eventbus/bus.go

// Event - обобщённый интерфейс события.
type Event interface {
    EventType() string
}

// Для LessonCompleted добавляем метод:
func (e LessonCompleted) EventType() string { return e.Type }

// Handler - функция-обработчик события.
type Handler func(ctx context.Context, event Event) error

// EventBus - порт для публикации и подписки на события.
type EventBus interface {
    // Publish отправляет событие всем зарегистрированным подписчикам.
    Publish(ctx context.Context, event Event) error

    // Subscribe регистрирует обработчик для данного типа события.
    Subscribe(eventType string, handler Handler)
}
<?php
// src/Application/Port/EventBus.php
declare(strict_types=1);

namespace App\Application\Port;

use App\Domain\Event\DomainEvent;

// final НЕ ставится на interface.
interface EventBus
{
    public function publish(DomainEvent $event): void;

    public function subscribe(string $eventClass, callable $handler): void;
}

In-process реализация - подходит для монолита и для начала разработки. Позже можно заменить на Kafka/NATS без изменения use-case:

// eventbus/bus.go

// InProcessBus - реализация EventBus через Go channels.
// Подписчики вызываются асинхронно в отдельных горутинах.
type InProcessBus struct {
    mu       sync.RWMutex
    handlers map[string][]Handler
}

// NewInProcessBus создаёт bus.
func NewInProcessBus() *InProcessBus {
    return &InProcessBus{
        handlers: make(map[string][]Handler),
    }
}

// Subscribe регистрирует обработчик для типа события.
func (b *InProcessBus) Subscribe(eventType string, handler Handler) {
    b.mu.Lock()
    defer b.mu.Unlock()
    b.handlers[eventType] = append(b.handlers[eventType], handler)
}

// Publish отправляет событие всем подписчикам данного типа.
// Каждый подписчик вызывается в отдельной горутине.
// Ошибки логируются, но не останавливают остальных подписчиков.
func (b *InProcessBus) Publish(ctx context.Context, event Event) error {
    b.mu.RLock()
    handlers := b.handlers[event.EventType()]
    b.mu.RUnlock()

    var wg sync.WaitGroup
    for _, h := range handlers {
        wg.Add(1)
        go func(handler Handler) {
            defer wg.Done()
            if err := handler(ctx, event); err != nil {
                slog.ErrorContext(ctx, "subscriber failed",
                    slog.String("event_type", event.EventType()),
                    slog.String("err", err.Error()),
                )
            }
        }(h)
    }
    wg.Wait()
    return nil
}
<?php
// src/Infrastructure/EventBus/InProcessBus.php
declare(strict_types=1);

namespace App\Infrastructure\EventBus;

use App\Application\Port\EventBus;
use App\Domain\Event\DomainEvent;
use Psr\Log\LoggerInterface;

// PHP-FPM однопоточный per request - подписчики выполняются последовательно,
// без горутин. Для параллельной обработки используй Symfony Messenger
// с async-транспортом (Redis Streams/AMQP) и несколькими workers.
final class InProcessBus implements EventBus
{
    /** @var array<string, list<callable>> */
    private array $handlers = [];

    public function __construct(
        private readonly LoggerInterface $logger,
    ) {}

    public function subscribe(string $eventClass, callable $handler): void
    {
        $this->handlers[$eventClass][] = $handler;
    }

    // Ошибки логируются, но не останавливают остальных подписчиков.
    public function publish(DomainEvent $event): void
    {
        foreach ($this->handlers[$event::class] ?? [] as $handler) {
            try {
                $handler($event);
            } catch (\Throwable $e) {
                $this->logger->error('subscriber failed', [
                    'event_type' => $event->eventType(),
                    'err' => $e->getMessage(),
                ]);
            }
        }
    }
}

Use-case: CompleteLesson с outbox

Use-case завершает урок, сохраняет прогресс и записывает событие в outbox - всё в одной транзакции. После коммита публикует событие в in-process bus для немедленной обработки:

// usecase/complete_lesson.go

// CompleteLessonUseCase - завершение урока с публикацией события.
type CompleteLessonUseCase struct {
    db       *sql.DB
    outbox   *outbox.Repository
    bus      eventbus.EventBus
}

// NewCompleteLessonUseCase создаёт use-case с зависимостями.
func NewCompleteLessonUseCase(db *sql.DB, outbox *outbox.Repository, bus eventbus.EventBus) *CompleteLessonUseCase {
    return &CompleteLessonUseCase{db: db, outbox: outbox, bus: bus}
}

// Execute завершает урок для пользователя.
func (uc *CompleteLessonUseCase) Execute(ctx context.Context, req CompleteLessonRequest) error {
    // Создаём доменное событие
    evt := domain.NewLessonCompleted(
        middleware.CorrelationFromCtx(ctx),
        req.UserID, req.LessonID,
        req.TrackSlug, req.LessonSlug,
    )

    // Сериализуем payload для outbox
    payload, err := json.Marshal(evt)
    if err != nil {
        return fmt.Errorf("marshal event: %w", err)
    }

    // Транзакция: прогресс + outbox
    tx, err := uc.db.BeginTx(ctx, nil)
    if err != nil {
        return fmt.Errorf("begin tx: %w", err)
    }
    defer tx.Rollback()

    // 1. Обновляем прогресс
    _, err = tx.ExecContext(ctx,
        `UPDATE lesson_progress
         SET completed = true, completed_at = now()
         WHERE user_id = $1 AND lesson_id = $2`,
        req.UserID, req.LessonID,
    )
    if err != nil {
        return fmt.Errorf("update progress: %w", err)
    }

    // 2. Записываем событие в outbox (в той же транзакции)
    _, err = tx.ExecContext(ctx,
        `INSERT INTO outbox (event_id, event_type, payload)
         VALUES ($1, $2, $3)`,
        evt.EventID, evt.Type, payload,
    )
    if err != nil {
        return fmt.Errorf("insert outbox: %w", err)
    }

    if err := tx.Commit(); err != nil {
        return fmt.Errorf("commit: %w", err)
    }

    // После коммита: публикуем в in-process bus для немедленной обработки
    // Если bus недоступен - outbox poller доставит событие позже
    if pubErr := uc.bus.Publish(ctx, evt); pubErr != nil {
        slog.WarnContext(ctx, "in-process publish failed, outbox poller will retry",
            slog.String("event_id", evt.EventID),
            slog.String("err", pubErr.Error()),
        )
    }

    return nil
}
<?php
// src/Application/UseCase/CompleteLessonUseCase.php
declare(strict_types=1);

namespace App\Application\UseCase;

use App\Application\Port\EventBus;
use App\Application\Port\OutboxRepository;
use App\Domain\Event\LessonCompleted;
use Doctrine\DBAL\Connection;
use Psr\Log\LoggerInterface;

final readonly class CompleteLessonRequest
{
    public function __construct(
        public int $userId,
        public int $lessonId,
        public string $trackSlug,
        public string $lessonSlug,
        public string $correlationId,
    ) {}
}

final class CompleteLessonUseCase
{
    public function __construct(
        private readonly Connection $connection,
        private readonly OutboxRepository $outbox,
        private readonly EventBus $bus,
        private readonly LoggerInterface $logger,
    ) {}

    public function execute(CompleteLessonRequest $req): void
    {
        $event = LessonCompleted::create(
            correlationId: $req->correlationId,
            userId: $req->userId,
            lessonId: $req->lessonId,
            trackSlug: $req->trackSlug,
            lessonSlug: $req->lessonSlug,
        );

        // Транзакция: progress + outbox
        $this->connection->transactional(function (Connection $tx) use ($req, $event): void {
            $tx->executeStatement(
                'UPDATE lesson_progress SET completed = true, completed_at = now()
                 WHERE user_id = :user_id AND lesson_id = :lesson_id',
                ['user_id' => $req->userId, 'lesson_id' => $req->lessonId],
            );

            $this->outbox->saveInTx($tx, $event);
        });

        // После коммита: пробуем сразу опубликовать в in-process bus.
        // Если упадёт - outbox poller всё равно доставит.
        try {
            $this->bus->publish($event);
        } catch (\Throwable $e) {
            $this->logger->warning('in-process publish failed, outbox poller will retry', [
                'event_id' => $event->eventId,
                'err' => $e->getMessage(),
            ]);
        }
    }
}

Обратите внимание: ошибка bus.Publish не приводит к ошибке use-case. Данные уже закоммичены. Outbox poller гарантирует доставку.

Подписчик 1: AchievementChecker

Проверяет, заработал ли пользователь бейдж после завершения урока:

// subscribers/achievement_checker.go

// AchievementChecker проверяет условия ачивок при завершении урока.
type AchievementChecker struct {
    db *sql.DB
}

func NewAchievementChecker(db *sql.DB) *AchievementChecker {
    return &AchievementChecker{db: db}
}

// Handle обрабатывает событие lesson.completed.
func (a *AchievementChecker) Handle(ctx context.Context, event eventbus.Event) error {
    evt, ok := event.(domain.LessonCompleted)
    if !ok {
        return fmt.Errorf("unexpected event type: %T", event)
    }

    slog.InfoContext(ctx, "checking achievements",
        slog.Int64("user_id", evt.UserID),
        slog.String("track_slug", evt.TrackSlug),
        slog.String("correlation_id", evt.CorrelationID),
    )

    // Подсчитываем завершённые уроки в треке
    var completedCount int
    err := a.db.QueryRowContext(ctx,
        `SELECT COUNT(*) FROM lesson_progress
         WHERE user_id = $1 AND track_slug = $2 AND completed = true`,
        evt.UserID, evt.TrackSlug,
    ).Scan(&completedCount)
    if err != nil {
        return fmt.Errorf("count completed: %w", err)
    }

    // Проверяем условия ачивок
    achievements := checkAchievementRules(evt.TrackSlug, completedCount)
    for _, achievement := range achievements {
        // UPSERT - идемпотентно: повторный вызов не создаёт дубль
        _, err := a.db.ExecContext(ctx,
            `INSERT INTO user_achievements (user_id, achievement_slug, earned_at)
             VALUES ($1, $2, now())
             ON CONFLICT (user_id, achievement_slug) DO NOTHING`,
            evt.UserID, achievement,
        )
        if err != nil {
            slog.ErrorContext(ctx, "grant achievement failed",
                slog.String("achievement", achievement),
                slog.String("err", err.Error()),
            )
        }
    }

    return nil
}

// checkAchievementRules возвращает список заслуженных ачивок.
func checkAchievementRules(trackSlug string, completedCount int) []string {
    var result []string
    if completedCount >= 1 {
        result = append(result, "first_lesson")
    }
    if completedCount >= 5 {
        result = append(result, "five_lessons")
    }
    if completedCount >= 10 {
        result = append(result, fmt.Sprintf("track_%s_master", trackSlug))
    }
    return result
}
<?php
// src/Application/Subscriber/AchievementChecker.php
declare(strict_types=1);

namespace App\Application\Subscriber;

use App\Domain\Event\LessonCompleted;
use Doctrine\DBAL\Connection;
use Psr\Log\LoggerInterface;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;

#[AsMessageHandler]
final class AchievementChecker
{
    public function __construct(
        private readonly Connection $connection,
        private readonly LoggerInterface $logger,
    ) {}

    public function __invoke(LessonCompleted $event): void
    {
        $this->logger->info('checking achievements', [
            'user_id' => $event->userId,
            'track_slug' => $event->trackSlug,
            'correlation_id' => $event->correlationId,
        ]);

        $completedCount = (int) $this->connection->fetchOne(
            'SELECT COUNT(*) FROM lesson_progress
             WHERE user_id = :uid AND track_slug = :slug AND completed = true',
            ['uid' => $event->userId, 'slug' => $event->trackSlug],
        );

        foreach ($this->checkAchievementRules($event->trackSlug, $completedCount) as $achievement) {
            // UPSERT - идемпотентно: повторный вызов не создаёт дубль
            $this->connection->executeStatement(
                'INSERT INTO user_achievements (user_id, achievement_slug, earned_at)
                 VALUES (:uid, :slug, now())
                 ON CONFLICT (user_id, achievement_slug) DO NOTHING',
                ['uid' => $event->userId, 'slug' => $achievement],
            );
        }
    }

    /** @return list<string> */
    private function checkAchievementRules(string $trackSlug, int $completedCount): array
    {
        $result = [];
        if ($completedCount >= 1) {
            $result[] = 'first_lesson';
        }
        if ($completedCount >= 5) {
            $result[] = 'five_lessons';
        }
        if ($completedCount >= 10) {
            $result[] = sprintf('track_%s_master', $trackSlug);
        }
        return $result;
    }
}

Подписчик 2: AnalyticsRecorder

Записывает факт завершения урока для статистики:

// subscribers/analytics_recorder.go

// AnalyticsRecorder записывает метрику завершения урока.
type AnalyticsRecorder struct {
    db *sql.DB
}

func NewAnalyticsRecorder(db *sql.DB) *AnalyticsRecorder {
    return &AnalyticsRecorder{db: db}
}

// Handle записывает событие в таблицу аналитики.
func (a *AnalyticsRecorder) Handle(ctx context.Context, event eventbus.Event) error {
    evt, ok := event.(domain.LessonCompleted)
    if !ok {
        return fmt.Errorf("unexpected event type: %T", event)
    }

    // UPSERT: идемпотентно - повторный вызов обновляет timestamp
    _, err := a.db.ExecContext(ctx,
        `INSERT INTO analytics_events (event_id, event_type, user_id, track_slug, lesson_slug, occurred_at)
         VALUES ($1, $2, $3, $4, $5, $6)
         ON CONFLICT (event_id) DO NOTHING`,
        evt.EventID, evt.Type, evt.UserID,
        evt.TrackSlug, evt.LessonSlug, evt.OccurredAt,
    )
    if err != nil {
        return fmt.Errorf("insert analytics: %w", err)
    }

    slog.InfoContext(ctx, "analytics event recorded",
        slog.String("event_id", evt.EventID),
        slog.Int64("user_id", evt.UserID),
        slog.String("correlation_id", evt.CorrelationID),
    )
    return nil
}
<?php
// src/Application/Subscriber/AnalyticsRecorder.php
declare(strict_types=1);

namespace App\Application\Subscriber;

use App\Domain\Event\LessonCompleted;
use Doctrine\DBAL\Connection;
use Psr\Log\LoggerInterface;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;

#[AsMessageHandler]
final class AnalyticsRecorder
{
    public function __construct(
        private readonly Connection $connection,
        private readonly LoggerInterface $logger,
    ) {}

    public function __invoke(LessonCompleted $event): void
    {
        // ON CONFLICT DO NOTHING - идемпотентно по event_id
        $this->connection->executeStatement(
            'INSERT INTO analytics_events (event_id, event_type, user_id, track_slug, lesson_slug, occurred_at)
             VALUES (:eid, :etype, :uid, :track, :lesson, :occurred)
             ON CONFLICT (event_id) DO NOTHING',
            [
                'eid' => $event->eventId,
                'etype' => $event->eventType(),
                'uid' => $event->userId,
                'track' => $event->trackSlug,
                'lesson' => $event->lessonSlug,
                'occurred' => $event->occurredAt->format('Y-m-d H:i:s'),
            ],
        );

        $this->logger->info('analytics event recorded', [
            'event_id' => $event->eventId,
            'user_id' => $event->userId,
            'correlation_id' => $event->correlationId,
        ]);
    }
}

Подписчик 3: NotificationSender

Отправляет поздравление пользователю (заглушка - в реальности это может быть email, push или in-app уведомление):

// subscribers/notification_sender.go

// NotificationSender отправляет уведомление при завершении урока.
type NotificationSender struct {
    notifier Notifier
}

// Notifier - порт для отправки уведомлений.
type Notifier interface {
    Send(ctx context.Context, userID int64, message string) error
}

func NewNotificationSender(notifier Notifier) *NotificationSender {
    return &NotificationSender{notifier: notifier}
}

// Handle отправляет поздравление.
func (n *NotificationSender) Handle(ctx context.Context, event eventbus.Event) error {
    evt, ok := event.(domain.LessonCompleted)
    if !ok {
        return fmt.Errorf("unexpected event type: %T", event)
    }

    message := fmt.Sprintf("Отлично! Вы завершили урок %q в треке %q.",
        evt.LessonSlug, evt.TrackSlug)

    if err := n.notifier.Send(ctx, evt.UserID, message); err != nil {
        return fmt.Errorf("send notification: %w", err)
    }

    slog.InfoContext(ctx, "notification sent",
        slog.Int64("user_id", evt.UserID),
        slog.String("correlation_id", evt.CorrelationID),
    )
    return nil
}
<?php
// src/Application/Subscriber/NotificationSender.php
declare(strict_types=1);

namespace App\Application\Subscriber;

use App\Application\Port\NotifierPort;
use App\Domain\Event\LessonCompleted;
use Psr\Log\LoggerInterface;
use Symfony\Component\Messenger\Attribute\AsMessageHandler;

#[AsMessageHandler]
final class NotificationSender
{
    public function __construct(
        private readonly NotifierPort $notifier,
        private readonly LoggerInterface $logger,
    ) {}

    public function __invoke(LessonCompleted $event): void
    {
        $message = sprintf(
            'Отлично! Вы завершили урок '%s' в треке '%s'.',
            $event->lessonSlug,
            $event->trackSlug,
        );

        $this->notifier->send($event->userId, $message);

        $this->logger->info('notification sent', [
            'user_id' => $event->userId,
            'correlation_id' => $event->correlationId,
        ]);
    }
}

Outbox: DDL и Repository

CREATE TABLE outbox (
    id           BIGSERIAL    PRIMARY KEY,
    event_id     VARCHAR(64)  NOT NULL UNIQUE,
    event_type   VARCHAR(128) NOT NULL,
    payload      JSONB        NOT NULL,
    created_at   TIMESTAMPTZ  NOT NULL DEFAULT now(),
    published_at TIMESTAMPTZ,
    attempts     INT          NOT NULL DEFAULT 0
);

CREATE INDEX idx_outbox_unpublished
    ON outbox (created_at)
    WHERE published_at IS NULL;
// outbox/repository.go

// Repository предоставляет доступ к outbox-таблице.
type Repository struct {
    db *sql.DB
}

func NewRepository(db *sql.DB) *Repository {
    return &Repository{db: db}
}

// FetchUnpublished возвращает пачку неотправленных событий.
func (r *Repository) FetchUnpublished(ctx context.Context, limit int) ([]OutboxEntry, error) {
    rows, err := r.db.QueryContext(ctx,
        `SELECT id, event_id, event_type, payload
         FROM outbox
         WHERE published_at IS NULL
         ORDER BY created_at
         LIMIT $1`,
        limit,
    )
    if err != nil {
        return nil, fmt.Errorf("query outbox: %w", err)
    }
    defer rows.Close()

    var entries []OutboxEntry
    for rows.Next() {
        var e OutboxEntry
        if err := rows.Scan(&e.ID, &e.EventID, &e.EventType, &e.Payload); err != nil {
            return nil, fmt.Errorf("scan: %w", err)
        }
        entries = append(entries, e)
    }
    return entries, rows.Err()
}

// MarkPublished отмечает событие как отправленное.
func (r *Repository) MarkPublished(ctx context.Context, id int64) error {
    _, err := r.db.ExecContext(ctx,
        `UPDATE outbox SET published_at = now() WHERE id = $1`, id)
    return err
}
<?php
// src/Infrastructure/Outbox/OutboxRepository.php
declare(strict_types=1);

namespace App\Infrastructure\Outbox;

use App\Application\Port\OutboxRepository as OutboxRepositoryPort;
use App\Domain\Event\LessonCompleted;
use Doctrine\DBAL\Connection;

final readonly class OutboxEntry
{
    public function __construct(
        public int $id,
        public string $eventId,
        public string $eventType,
        public string $payload,
    ) {}
}

final class OutboxRepository implements OutboxRepositoryPort
{
    public function __construct(
        private readonly Connection $connection,
    ) {}

    public function saveInTx(Connection $tx, LessonCompleted $event): void
    {
        $tx->insert('outbox', [
            'event_id' => $event->eventId,
            'event_type' => $event->eventType(),
            'payload' => json_encode([
                'event_id' => $event->eventId,
                'correlation_id' => $event->correlationId,
                'occurred_at' => $event->occurredAt->format(\DateTimeInterface::RFC3339),
                'user_id' => $event->userId,
                'lesson_id' => $event->lessonId,
                'track_slug' => $event->trackSlug,
                'lesson_slug' => $event->lessonSlug,
            ], JSON_THROW_ON_ERROR),
        ]);
    }

    /** @return list<OutboxEntry> */
    public function fetchUnpublished(int $limit): array
    {
        $rows = $this->connection->fetchAllAssociative(
            'SELECT id, event_id, event_type, payload
             FROM outbox
             WHERE published_at IS NULL
             ORDER BY created_at
             LIMIT :lim
             FOR UPDATE SKIP LOCKED',
            ['lim' => $limit],
        );

        return array_map(
            static fn(array $r) => new OutboxEntry(
                id: (int) $r['id'],
                eventId: $r['event_id'],
                eventType: $r['event_type'],
                payload: (string) $r['payload'],
            ),
            $rows,
        );
    }

    public function markPublished(int $id): void
    {
        $this->connection->executeStatement(
            'UPDATE outbox SET published_at = now() WHERE id = :id',
            ['id' => $id],
        );
    }
}

Outbox Poller

// outbox/poller.go

// Poller периодически читает outbox и публикует события в EventBus.
type Poller struct {
    repo     *Repository
    bus      eventbus.EventBus
    interval time.Duration
    batch    int
}

func NewPoller(repo *Repository, bus eventbus.EventBus, interval time.Duration, batch int) *Poller {
    return &Poller{repo: repo, bus: bus, interval: interval, batch: batch}
}

// Run запускает polling. Останавливается по отмене контекста.
func (p *Poller) Run(ctx context.Context) {
    ticker := time.NewTicker(p.interval)
    defer ticker.Stop()

    slog.Info("outbox poller started",
        slog.Duration("interval", p.interval),
        slog.Int("batch_size", p.batch),
    )

    for {
        select {
        case <-ctx.Done():
            slog.Info("outbox poller stopped")
            return
        case <-ticker.C:
            entries, err := p.repo.FetchUnpublished(ctx, p.batch)
            if err != nil {
                slog.ErrorContext(ctx, "outbox fetch failed",
                    slog.String("err", err.Error()),
                )
                continue
            }

            for _, entry := range entries {
                // Десериализуем событие по типу
                evt, err := deserializeEvent(entry.EventType, entry.Payload)
                if err != nil {
                    slog.ErrorContext(ctx, "outbox deserialize failed",
                        slog.String("event_id", entry.EventID),
                        slog.String("err", err.Error()),
                    )
                    continue
                }

                if err := p.bus.Publish(ctx, evt); err != nil {
                    slog.ErrorContext(ctx, "outbox publish failed",
                        slog.String("event_id", entry.EventID),
                        slog.String("err", err.Error()),
                    )
                    continue
                }

                if err := p.repo.MarkPublished(ctx, entry.ID); err != nil {
                    slog.ErrorContext(ctx, "outbox mark failed",
                        slog.String("event_id", entry.EventID),
                        slog.String("err", err.Error()),
                    )
                }
            }
        }
    }
}

// deserializeEvent восстанавливает событие из JSON по типу.
func deserializeEvent(eventType string, payload []byte) (eventbus.Event, error) {
    switch eventType {
    case "lesson.completed":
        var evt domain.LessonCompleted
        if err := json.Unmarshal(payload, &evt); err != nil {
            return nil, err
        }
        return evt, nil
    default:
        return nil, fmt.Errorf("unknown event type: %s", eventType)
    }
}
<?php
// src/Application/Cron/OutboxPollCommand.php
declare(strict_types=1);

namespace App\Application\Cron;

use App\Application\Port\EventBus;
use App\Domain\Event\LessonCompleted;
use App\Infrastructure\Outbox\OutboxRepository;
use Psr\Log\LoggerInterface;
use Symfony\Component\Console\Attribute\AsCommand;
use Symfony\Component\Console\Command\Command;
use Symfony\Component\Console\Input\InputInterface;
use Symfony\Component\Console\Output\OutputInterface;

// Альтернатива - Symfony Messenger worker (`messenger:consume outbox-transport`).
// Здесь - однопроходная команда из cron, проще в развёртывании.
#[AsCommand(name: 'app:outbox-poll')]
final class OutboxPollCommand extends Command
{
    private const BATCH_SIZE = 10;

    public function __construct(
        private readonly OutboxRepository $repo,
        private readonly EventBus $bus,
        private readonly LoggerInterface $logger,
    ) {
        parent::__construct();
    }

    protected function execute(InputInterface $input, OutputInterface $output): int
    {
        foreach ($this->repo->fetchUnpublished(self::BATCH_SIZE) as $entry) {
            try {
                $event = $this->deserializeEvent($entry->eventType, $entry->payload);
                $this->bus->publish($event);
                $this->repo->markPublished($entry->id);
            } catch (\Throwable $e) {
                $this->logger->error('outbox publish failed', [
                    'event_id' => $entry->eventId,
                    'err' => $e->getMessage(),
                ]);
            }
        }

        return Command::SUCCESS;
    }

    private function deserializeEvent(string $eventType, string $payload): LessonCompleted
    {
        $data = json_decode($payload, true, 512, JSON_THROW_ON_ERROR);

        return match ($eventType) {
            'lesson.completed' => new LessonCompleted(
                eventId: $data['event_id'],
                correlationId: $data['correlation_id'],
                occurredAt: new \DateTimeImmutable($data['occurred_at']),
                userId: $data['user_id'],
                lessonId: $data['lesson_id'],
                trackSlug: $data['track_slug'],
                lessonSlug: $data['lesson_slug'],
            ),
            default => throw new \LogicException(sprintf('unknown event type: %s', $eventType)),
        };
    }
}

Wiring в main.go

Всё собирается в main.go - создаём bus, регистрируем подписчиков, передаём bus в use-case:

// cmd/server/main.go (фрагмент)

func main() {
    // ... DB, router, middleware setup ...

    // 1. Создаём EventBus
    bus := eventbus.NewInProcessBus()

    // 2. Создаём подписчиков
    achievementChecker := subscribers.NewAchievementChecker(db)
    analyticsRecorder := subscribers.NewAnalyticsRecorder(db)
    notificationSender := subscribers.NewNotificationSender(stubNotifier{})

    // 3. Регистрируем подписчиков (одна строка на подписчика!)
    bus.Subscribe("lesson.completed", achievementChecker.Handle)
    bus.Subscribe("lesson.completed", analyticsRecorder.Handle)
    bus.Subscribe("lesson.completed", notificationSender.Handle)

    // 4. Создаём outbox и use-case
    outboxRepo := outbox.NewRepository(db)
    completeLessonUC := usecase.NewCompleteLessonUseCase(db, outboxRepo, bus)

    // 5. Запускаём outbox poller в фоне
    pollerCtx, pollerCancel := context.WithCancel(context.Background())
    defer pollerCancel()

    poller := outbox.NewPoller(outboxRepo, bus, 500*time.Millisecond, 10)
    go poller.Run(pollerCtx)

    // 6. Регистрируем handler
    router.Put("/progress/lessons/{id}/complete", handler.CompleteLesson(completeLessonUC))

    // ... server start, graceful shutdown ...
}
<?php
// config/services.yaml - autoconfigure через #[AsMessageHandler].
// Контроллер просто инжектит MessageBusInterface и dispatches DomainEvent.

// src/Controller/ProgressController.php
declare(strict_types=1);

namespace App\Controller;

use App\Application\UseCase\CompleteLessonRequest;
use App\Application\UseCase\CompleteLessonUseCase;
use Symfony\Bundle\FrameworkBundle\Controller\AbstractController;
use Symfony\Component\HttpFoundation\JsonResponse;
use Symfony\Component\HttpFoundation\Request;
use Symfony\Component\Routing\Attribute\Route;

final class ProgressController extends AbstractController
{
    public function __construct(
        private readonly CompleteLessonUseCase $completeLesson,
    ) {}

    #[Route('/api/progress/lessons/{id}/complete', methods: ['PUT'])]
    public function complete(int $id, Request $request): JsonResponse
    {
        $payload = json_decode($request->getContent(), true, 512, JSON_THROW_ON_ERROR);

        ($this->completeLesson)(new CompleteLessonRequest(
            userId: $this->getUser()->getId(),
            lessonId: $id,
            trackSlug: $payload['track_slug'],
            lessonSlug: $payload['lesson_slug'],
            correlationId: $request->headers->get('X-Correlation-Id', 'unknown'),
        ));

        return new JsonResponse(['status' => 'ok']);
    }
}

// Подписчики (AchievementChecker, AnalyticsRecorder, NotificationSender)
// с #[AsMessageHandler] - Symfony автоматически подвязывает их к типу события.
// Outbox poller: `bin/console app:outbox-poll` через systemd-timer или
// `bin/console messenger:consume outbox` для async-транспорта.

Добавить нового подписчика = один файл с реализацией + одна строка bus.Subscribe(...). Ни один существующий файл не меняется (Open-Closed Principle).

Тестирование: FakeEventBus

Для unit-тестов use-case не нужен реальный bus. Создаём FakeEventBus, который записывает все опубликованные события:

// eventbus/bus_test.go

// FakeEventBus записывает все опубликованные события для проверки в тестах.
type FakeEventBus struct {
    mu     sync.Mutex
    Events []Event
}

func NewFakeEventBus() *FakeEventBus {
    return &FakeEventBus{}
}

func (f *FakeEventBus) Publish(_ context.Context, event Event) error {
    f.mu.Lock()
    defer f.mu.Unlock()
    f.Events = append(f.Events, event)
    return nil
}

func (f *FakeEventBus) Subscribe(_ string, _ Handler) {
    // No-op для тестов
}

// Published возвращает все опубликованные события данного типа.
func (f *FakeEventBus) Published(eventType string) []Event {
    f.mu.Lock()
    defer f.mu.Unlock()
    var result []Event
    for _, e := range f.Events {
        if e.EventType() == eventType {
            result = append(result, e)
        }
    }
    return result
}
<?php
// tests/Fake/FakeEventBus.php
declare(strict_types=1);

namespace App\Tests\Fake;

use App\Application\Port\EventBus;
use App\Domain\Event\DomainEvent;

// Каждый t.Run в PHPUnit/Pest создаёт новый bus - класс stateful, не singleton.
final class FakeEventBus implements EventBus
{
    /** @var list<DomainEvent> */
    public array $events = [];

    public function subscribe(string $eventClass, callable $handler): void
    {
        // No-op для тестов
    }

    public function publish(DomainEvent $event): void
    {
        $this->events[] = $event;
    }

    /** @return list<DomainEvent> */
    public function published(string $eventType): array
    {
        return array_values(array_filter(
            $this->events,
            static fn(DomainEvent $e) => $e->eventType() === $eventType,
        ));
    }
}

Тест use-case:

func TestCompleteLesson_PublishesEvent(t *testing.T) {
    db := setupTestDB(t) // in-memory или testcontainers
    fakeBus := eventbus.NewFakeEventBus()
    outboxRepo := outbox.NewRepository(db)
    uc := usecase.NewCompleteLessonUseCase(db, outboxRepo, fakeBus)

    err := uc.Execute(context.Background(), usecase.CompleteLessonRequest{
        UserID:     1,
        LessonID:   42,
        TrackSlug:  "go",
        LessonSlug: "goroutines",
    })
    require.NoError(t, err)

    // Проверяем, что событие опубликовано
    events := fakeBus.Published("lesson.completed")
    require.Len(t, events, 1)

    evt := events[0].(domain.LessonCompleted)
    assert.Equal(t, int64(1), evt.UserID)
    assert.Equal(t, int64(42), evt.LessonID)
    assert.Equal(t, "go", evt.TrackSlug)
    assert.Equal(t, "lesson.completed", evt.Type)
    assert.NotEmpty(t, evt.EventID)
}

func TestCompleteLesson_WritesToOutbox(t *testing.T) {
    db := setupTestDB(t)
    fakeBus := eventbus.NewFakeEventBus()
    outboxRepo := outbox.NewRepository(db)
    uc := usecase.NewCompleteLessonUseCase(db, outboxRepo, fakeBus)

    err := uc.Execute(context.Background(), usecase.CompleteLessonRequest{
        UserID:     1,
        LessonID:   42,
        TrackSlug:  "go",
        LessonSlug: "goroutines",
    })
    require.NoError(t, err)

    // Проверяем, что запись появилась в outbox
    entries, err := outboxRepo.FetchUnpublished(context.Background(), 10)
    require.NoError(t, err)
    require.Len(t, entries, 1)
    assert.Equal(t, "lesson.completed", entries[0].EventType)
}
<?php
// tests/Application/UseCase/CompleteLessonUseCaseTest.php
declare(strict_types=1);

namespace App\Tests\Application\UseCase;

use App\Application\UseCase\CompleteLessonRequest;
use App\Application\UseCase\CompleteLessonUseCase;
use App\Domain\Event\LessonCompleted;
use App\Tests\Fake\FakeEventBus;
use PHPUnit\Framework\TestCase;

final class CompleteLessonUseCaseTest extends TestCase
{
    public function testPublishesEvent(): void
    {
        $db = $this->setUpTestDB();
        $fakeBus = new FakeEventBus();
        $outboxRepo = new OutboxRepository($db);
        $uc = new CompleteLessonUseCase($db, $outboxRepo, $fakeBus, $this->createMock(\Psr\Log\LoggerInterface::class));

        $uc(new CompleteLessonRequest(
            userId: 1,
            lessonId: 42,
            trackSlug: 'go',
            lessonSlug: 'goroutines',
            correlationId: 'test-corr',
        ));

        // Проверяем, что событие опубликовано
        $events = $fakeBus->published('lesson.completed');
        self::assertCount(1, $events);

        /** @var LessonCompleted $evt */
        $evt = $events[0];
        self::assertSame(1, $evt->userId);
        self::assertSame(42, $evt->lessonId);
        self::assertSame('go', $evt->trackSlug);
        self::assertSame('lesson.completed', $evt->eventType());
        self::assertNotEmpty($evt->eventId);
    }

    public function testWritesToOutbox(): void
    {
        $db = $this->setUpTestDB();
        $fakeBus = new FakeEventBus();
        $outboxRepo = new OutboxRepository($db);
        $uc = new CompleteLessonUseCase($db, $outboxRepo, $fakeBus, $this->createMock(\Psr\Log\LoggerInterface::class));

        $uc(new CompleteLessonRequest(
            userId: 1,
            lessonId: 42,
            trackSlug: 'go',
            lessonSlug: 'goroutines',
            correlationId: 'test-corr',
        ));

        $entries = $outboxRepo->fetchUnpublished(10);
        self::assertCount(1, $entries);
        self::assertSame('lesson.completed', $entries[0]->eventType);
    }
}

Тест подписчика:

func TestAchievementChecker_GrantsFirstLesson(t *testing.T) {
    db := setupTestDB(t)
    checker := subscribers.NewAchievementChecker(db)

    // Подготовка: один завершённый урок
    insertLessonProgress(t, db, 1, "go", true)

    evt := domain.LessonCompleted{
        EventID:   "test-event-1",
        Type:      "lesson.completed",
        UserID:    1,
        TrackSlug: "go",
    }

    err := checker.Handle(context.Background(), evt)
    require.NoError(t, err)

    // Проверяем ачивку
    var count int
    db.QueryRow(`SELECT COUNT(*) FROM user_achievements
        WHERE user_id = 1 AND achievement_slug = 'first_lesson'`).Scan(&count)
    assert.Equal(t, 1, count)

    // Повторный вызов - идемпотентно, ачивка не дублируется
    err = checker.Handle(context.Background(), evt)
    require.NoError(t, err)

    db.QueryRow(`SELECT COUNT(*) FROM user_achievements
        WHERE user_id = 1 AND achievement_slug = 'first_lesson'`).Scan(&count)
    assert.Equal(t, 1, count) // по-прежнему 1
}
<?php
// tests/Application/Subscriber/AchievementCheckerTest.php
declare(strict_types=1);

namespace App\Tests\Application\Subscriber;

use App\Application\Subscriber\AchievementChecker;
use App\Domain\Event\LessonCompleted;
use PHPUnit\Framework\TestCase;

final class AchievementCheckerTest extends TestCase
{
    public function testGrantsFirstLesson(): void
    {
        $db = $this->setUpTestDB();
        $checker = new AchievementChecker($db, $this->createMock(\Psr\Log\LoggerInterface::class));

        // Подготовка: один завершённый урок
        $this->insertLessonProgress($db, 1, 'go', true);

        $event = new LessonCompleted(
            eventId: 'test-event-1',
            correlationId: 'test-corr',
            occurredAt: new \DateTimeImmutable(),
            userId: 1,
            lessonId: 1,
            trackSlug: 'go',
            lessonSlug: 'goroutines',
        );

        $checker($event);

        // Проверяем ачивку
        $count = (int) $db->fetchOne(
            'SELECT COUNT(*) FROM user_achievements
             WHERE user_id = 1 AND achievement_slug = :slug',
            ['slug' => 'first_lesson'],
        );
        self::assertSame(1, $count);

        // Повторный вызов - идемпотентно, ачивка не дублируется
        $checker($event);

        $count = (int) $db->fetchOne(
            'SELECT COUNT(*) FROM user_achievements
             WHERE user_id = 1 AND achievement_slug = :slug',
            ['slug' => 'first_lesson'],
        );
        self::assertSame(1, $count); // по-прежнему 1
    }
}
Это главное преимущество event-driven архитектуры. Завтра понадобится отправлять данные в рекомендательную систему? Создайте `RecommendationUpdater`, реализуйте `Handle`, добавьте `bus.Subscribe("lesson.completed", recommender.Handle)` в main.go. Use-case, другие подписчики, outbox - всё остаётся нетронутым.

Чеклист проекта

Убедитесь, что реализовано:

  • Доменное событие LessonCompleted с полями event_id, type, correlation_id, occurred_at, user_id, lesson_id, track_slug, lesson_slug
  • Интерфейс EventBus с методами Publish и Subscribe
  • InProcessBus - реализация через горутины с логированием ошибок подписчиков
  • Таблица outbox с DDL (id, event_id, event_type, payload, created_at, published_at, attempts)
  • Use-case CompleteLessonUseCase, который сохраняет прогресс и outbox в одной транзакции
  • Подписчик AchievementChecker с идемпотентным UPSERT
  • Подписчик AnalyticsRecorder с дедупликацией по event_id
  • Подписчик NotificationSender с интерфейсом Notifier (заглушка)
  • OutboxPoller - горутина с интервалом, batch-чтением и mark published
  • Wiring в main.go: bus -> subscribers -> use-case -> poller
  • FakeEventBus для unit-тестов
  • Тест: CompleteLesson публикует событие и записывает в outbox
  • Тест: AchievementChecker идемпотентен при повторном вызове

Типичные ошибки в этом проекте

  • bus.Publish напрямую из use-case вместо outbox - публикация в брокер до commit транзакции. Откат БД → событие уже ушло. Use-case пишет только в outbox-таблицу в той же транзакции, что и progress; брокеру отправляет poller.
  • OutboxPoller - два инстанса без leader election - оба читают одну запись, дубли × N. Либо FOR UPDATE SKIP LOCKED, либо один инстанс с k8s Recreate strategy.
  • AchievementChecker через INSERT вместо UPSERT - повторное событие → unique violation, retry, DLQ. Идемпотентность через INSERT ... ON CONFLICT DO NOTHING или ON CONFLICT (event_id) DO NOTHING - фундамент проекта.
  • AnalyticsRecorder пишет в Redis без AOF - событие записано, Redis перезагрузился, аналитика потеряна. Для аналитики достаточно at-least-once; для критики - PostgreSQL/долговечное хранилище.
  • NotificationSender отправляет email синхронно - SMTP отвечает за 5 секунд, subscriber обрабатывает 1 событие в 5с, очередь растёт. Notifier с deadline + retry + DLQ, либо отдельный сервис отправки.
  • FakeEventBus для теста разделён между параллельными t.Run - события одного теста ловит handler другого. Создавай новый bus в каждом setup.
  • Outbox-таблица без индекса на (published_at, created_at) - poller-запрос WHERE published_at IS NULL под нагрузкой делает seqscan на миллионе строк. Индекс - обязательно с первого дня.
  • correlation_id не пробрасывается в horoutu subscriber-а - лог HTTP-запроса и лог subscriber-а не сопоставить. Subscriber извлекает correlation_id из envelope и добавляет в slog-контекст.

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

  • Реализуй полный pipeline: use-case -> outbox -> poller -> bus -> subscribers
  • Напиши unit-тест с FakeEventBus, проверяющий публикацию LessonCompleted
  • Убедись, что каждый подписчик идемпотентен (вызови Handle дважды - результат не меняется)
  • Добавь четвёртого подписчика (например, RecommendationUpdater), не меняя ни одного существующего файла
  • Запусти go test ./... и убедись, что все тесты проходят

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