Репозитории и интерфейсы: где им жить

Репозитории и интерфейсы: где им жить

Repository - это порт для сохранения и получения данных. В терминах гексагональной архитектуры репозиторий является вторичным портом: бизнес-логика говорит «мне нужно сохранить задачу», а конкретная реализация (PostgreSQL, MongoDB, in-memory map) - это адаптер, который этот порт реализует. См. также Repository в DDD.

Главный вопрос: где объявить интерфейс, а где положить реализацию?

Правило: интерфейс рядом с потребителем

Где жить интерфейсу: рядом с реализацией плохо, рядом с потребителем хорошо

В Go принято объявлять интерфейс там, где он используется, а не там, где реализуется. В контексте Hex это означает: интерфейс репозитория живёт в слое domain или application (use-case), а реализация - в adapters.

internal/
  domain/
    task.go              # entity Task
    errors.go            # доменные ошибки
  usecase/
    ports/
      task_repo.go       # интерфейс TaskRepository
      notifier.go        # интерфейс TaskNotifier
    create_task.go       # use-case
    complete_task.go
  adapters/
    postgres/
      task_repo.go       # PostgresTaskRepository
    inmemory/
      task_repo.go       # InMemoryTaskRepository (для тестов)
    http/
      task_handler.go    # HTTP handler
Если интерфейс лежит рядом с реализацией (в `adapters/postgres/`), то use-case вынужден импортировать пакет адаптера. Это разрушает всю идею: домен начинает зависеть от инфраструктуры. Интерфейс в `usecase/ports/` - и зависимость направлена внутрь.

Интерфейс TaskRepository

Сигнатуры принимают и возвращают доменные типы. Никаких *sql.Row, pgx.Rows или bson.M в контракте порта.

// internal/usecase/ports/task_repo.go
package ports

import (
    "context"

    "myapp/internal/domain"
)

type TaskRepository interface {
    Create(ctx context.Context, task domain.Task) error
    GetByID(ctx context.Context, id string) (domain.Task, error)
    ListByUser(ctx context.Context, userID string) ([]domain.Task, error)
    Update(ctx context.Context, task domain.Task) error
    Delete(ctx context.Context, id string) error
}
<?php
// src/Application/Port/TaskRepositoryPort.php
declare(strict_types=1);

namespace App\Application\Port;

use App\Domain\Task;

interface TaskRepositoryPort
{
    public function create(Task $task): void;

    public function getById(string $id): Task;

    /** @return list<Task> */
    public function listByUser(string $userId): array;

    public function update(Task $task): void;

    public function delete(string $id): void;
}

Обрати внимание:

  • Каждый метод принимает context.Context первым аргументом - это стандарт Go.
  • Возвращаются domain.Task, а не *TaskRow или map[string]interface{}.
  • Ошибки - обычный error. Доменные ошибки (например, domain.ErrTaskNotFound) определяются в пакете domain.

Реализация: PostgresTaskRepository / DoctrineTaskRepository

Адаптер реализует интерфейс и работает с конкретной базой. Внутри - SQL, pgx, Doctrine, что угодно. Наружу - только доменные типы.

// internal/adapters/postgres/task_repo.go
package postgres

import (
    "context"
    "errors"

    "github.com/jackc/pgx/v5"
    "github.com/jackc/pgx/v5/pgxpool"

    "myapp/internal/domain"
)

type TaskRepo struct {
    pool *pgxpool.Pool
}

func NewTaskRepo(pool *pgxpool.Pool) *TaskRepo {
    return &TaskRepo{pool: pool}
}

func (r *TaskRepo) Create(ctx context.Context, task domain.Task) error {
    _, err := r.pool.Exec(ctx,
        `INSERT INTO tasks (id, user_id, title, completed)
         VALUES ($1, $2, $3, $4)`,
        task.ID, task.UserID, task.Title, task.Completed,
    )
    return err
}

func (r *TaskRepo) GetByID(ctx context.Context, id string) (domain.Task, error) {
    var t domain.Task
    err := r.pool.QueryRow(ctx,
        `SELECT id, user_id, title, completed, created_at
         FROM tasks WHERE id = $1`, id,
    ).Scan(&t.ID, &t.UserID, &t.Title, &t.Completed, &t.CreatedAt)

    if errors.Is(err, pgx.ErrNoRows) {
        return domain.Task{}, domain.ErrTaskNotFound
    }
    return t, err
}

func (r *TaskRepo) ListByUser(ctx context.Context, userID string) ([]domain.Task, error) {
    rows, err := r.pool.Query(ctx,
        `SELECT id, user_id, title, completed, created_at
         FROM tasks WHERE user_id = $1 ORDER BY created_at DESC`, userID,
    )
    if err != nil {
        return nil, err
    }
    defer rows.Close()

    var tasks []domain.Task
    for rows.Next() {
        var t domain.Task
        if err := rows.Scan(&t.ID, &t.UserID, &t.Title, &t.Completed, &t.CreatedAt); err != nil {
            return nil, err
        }
        tasks = append(tasks, t)
    }
    return tasks, rows.Err()
}
<?php
// src/Infrastructure/Persistence/Doctrine/DoctrineTaskRepository.php
declare(strict_types=1);

namespace App\Infrastructure\Persistence\Doctrine;

use App\Application\Port\TaskRepositoryPort;
use App\Domain\Task;
use App\Domain\TaskError;
use Doctrine\DBAL\Connection;

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

    public function create(Task $task): void
    {
        $this->connection->insert('tasks', [
            'id' => $task->id(),
            'user_id' => $task->userId(),
            'title' => $task->title(),
            'completed' => $task->isCompleted() ? 1 : 0,
        ]);
    }

    public function getById(string $id): Task
    {
        $row = $this->connection->fetchAssociative(
            'SELECT id, user_id, title, completed, created_at FROM tasks WHERE id = :id',
            ['id' => $id],
        );
        if ($row === false) {
            throw TaskError::notFound();
        }
        return $this->hydrate($row);
    }

    /** @return list<Task> */
    public function listByUser(string $userId): array
    {
        $rows = $this->connection->fetchAllAssociative(
            'SELECT id, user_id, title, completed, created_at FROM tasks WHERE user_id = :uid ORDER BY created_at DESC',
            ['uid' => $userId],
        );
        return array_map(fn (array $row) => $this->hydrate($row), $rows);
    }

    public function update(Task $task): void
    {
        $this->connection->update(
            'tasks',
            ['title' => $task->title(), 'completed' => $task->isCompleted() ? 1 : 0],
            ['id' => $task->id()],
        );
    }

    public function delete(string $id): void
    {
        $this->connection->delete('tasks', ['id' => $id]);
    }

    /** @param array<string,mixed> $row */
    private function hydrate(array $row): Task
    {
        // Task::fromStorage гидрирует объект без повторной валидации
        return Task::fromStorage(
            id: (string) $row['id'],
            userId: (string) $row['user_id'],
            title: (string) $row['title'],
            completed: (bool) $row['completed'],
        );
    }
}

Ключевой момент в GetByID: ошибку «нет строки» мы превращаем в доменную (domain.ErrTaskNotFound / TaskError::notFound()). Потребитель (use-case) не знает ничего про pgx или Doctrine.

Реализация: InMemoryTaskRepository

Для unit-тестов не нужна база данных. Достаточно map.

// internal/adapters/inmemory/task_repo.go
package inmemory

import (
    "context"
    "sync"

    "myapp/internal/domain"
)

type TaskRepo struct {
    mu    sync.RWMutex
    store map[string]domain.Task
}

func NewTaskRepo() *TaskRepo {
    return &TaskRepo{store: make(map[string]domain.Task)}
}

func (r *TaskRepo) Create(_ context.Context, task domain.Task) error {
    r.mu.Lock()
    defer r.mu.Unlock()
    r.store[task.ID] = task
    return nil
}

func (r *TaskRepo) GetByID(_ context.Context, id string) (domain.Task, error) {
    r.mu.RLock()
    defer r.mu.RUnlock()
    t, ok := r.store[id]
    if !ok {
        return domain.Task{}, domain.ErrTaskNotFound
    }
    return t, nil
}

func (r *TaskRepo) ListByUser(_ context.Context, userID string) ([]domain.Task, error) {
    r.mu.RLock()
    defer r.mu.RUnlock()
    var result []domain.Task
    for _, t := range r.store {
        if t.UserID == userID {
            result = append(result, t)
        }
    }
    return result, nil
}
<?php
// src/Infrastructure/Persistence/InMemory/InMemoryTaskRepository.php
declare(strict_types=1);

namespace App\Infrastructure\Persistence\InMemory;

use App\Application\Port\TaskRepositoryPort;
use App\Domain\Task;
use App\Domain\TaskError;

final class InMemoryTaskRepository implements TaskRepositoryPort
{
    /** @var array<string, Task> */
    private array $store = [];

    public function create(Task $task): void
    {
        $this->store[$task->id()] = $task;
    }

    public function getById(string $id): Task
    {
        if (!isset($this->store[$id])) {
            throw TaskError::notFound();
        }
        return $this->store[$id];
    }

    /** @return list<Task> */
    public function listByUser(string $userId): array
    {
        $result = [];
        foreach ($this->store as $task) {
            if ($task->userId() === $userId) {
                $result[] = $task;
            }
        }
        return $result;
    }

    public function update(Task $task): void
    {
        $this->store[$task->id()] = $task;
    }

    public function delete(string $id): void
    {
        unset($this->store[$id]);
    }
}

PHP-FPM однопоточный per request. Мьютекс на in-memory store для тестов не нужен - PHPUnit запускает тесты последовательно в одном процессе. Если используешь ParaTest (параллельный запуск), у каждого процесса своя память - изоляция тестов сохраняется.

Даже в тестовом репозитории ставь мьютекс. Без него `go test -race` поймает data race, если тесты запускаются параллельно. Привычка писать thread-safe код - бесплатная страховка.

Транзакции в репозиториях

Когда use-case требует атомарности (создать задачу + отправить событие), есть два подхода.

Подход 1: Unit of Work через интерфейс.

type UnitOfWork interface {
    Do(ctx context.Context, fn func(repos Repositories) error) error
}

type Repositories struct {
    Tasks  TaskRepository
    Events EventRepository
}
<?php
// src/Application/Port/UnitOfWorkPort.php
declare(strict_types=1);

namespace App\Application\Port;

interface UnitOfWorkPort
{
    /**
     * @template T
     * @param callable(Repositories): T $fn
     * @return T
     */
    public function transactional(callable $fn): mixed;
}

final readonly class Repositories
{
    public function __construct(
        public TaskRepositoryPort $tasks,
        public EventRepositoryPort $events,
    ) {}
}

// src/Infrastructure/Persistence/Doctrine/DoctrineUnitOfWork.php
final class DoctrineUnitOfWork implements UnitOfWorkPort
{
    public function __construct(
        private readonly \Doctrine\DBAL\Connection $connection,
        private readonly Repositories $repos,
    ) {}

    public function transactional(callable $fn): mixed
    {
        return $this->connection->transactional(fn () => $fn($this->repos));
    }
}

Use-case вызывает uow.Do(ctx, func(repos) { ... }) / $uow->transactional(fn ($repos) => ...), а адаптер оборачивает вызов в BEGIN / COMMIT.

Подход 2: передать транзакцию через context. Этот способ проще, но связывает слои через скрытое состояние в контексте. Первый подход предпочтительнее в большинстве случаев.

Если интерфейс возвращает `*sql.Rows`, потребитель обязан вызвать `rows.Close()`. Это утечка абстракции: use-case знает, что работает с SQL. Репозиторий должен вернуть `[]domain.Task` или ошибку - ничего между.

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

  • Объяви интерфейс UserRepository в пакете usecase/ports/ с методами Create, GetByID, GetByEmail
  • Напиши InMemoryUserRepo, который реализует этот интерфейс через map[string]domain.User
  • Убедись, что InMemoryUserRepo возвращает доменную ошибку domain.ErrUserNotFound, а не generic «not found»
  • Проверь, что use-case импортирует только пакеты domain и usecase/ports - никаких database/sql или pgx

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