Мини-проект: события в BackendStart (LessonCompleted → Achievements/Email/Analytics)
Мини-проект: события в BackendStart (LessonCompleted → Achievements/Email/Analytics)
В этом уроке мы соберём всё, что изучили за трек, в работающую систему. Возьмём реальный сценарий BackendStart: пользователь завершает урок, и это порождает цепочку побочных эффектов - проверку ачивок, запись аналитики, отправку уведомления. Каждый эффект реализован отдельным подписчиком. Ни один из них не знает о существовании остальных. Добавление нового подписчика - это один файл и одна строка регистрации.
Мы реализуем: доменное событие, EventBus (in-process), outbox-таблицу, три подписчика, wiring в main.go и unit-тесты.
Структура проекта
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
}
}
Чеклист проекта
Убедитесь, что реализовано:
- Доменное событие
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, либо один инстанс с k8sRecreatestrategy. 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 ./...и убедись, что все тесты проходят