Hexagonal Architecture (гексагональная архитектура) предложена Алистером Кокбёрном (Alistair Cockburn) в 2005 году. Другое название -- Ports and Adapters.
Главная идея: приложение не должно зависеть от способа доставки данных (HTTP, CLI, очередь) или способа хранения (PostgreSQL, файлы, API).
Порт -- это интерфейс, который определяет контракт взаимодействия с приложением. Порты бывают двух видов:
Тип порта
Направление
Назначение
Пример
Primary (Input)
Внешний мир → Приложение
Определяет что приложение умеет делать
CreateOrderUseCase
Secondary (Output)
Приложение → Внешний мир
Определяет что приложению нужно извне
OrderRepositoryInterface
Адаптеры (Adapters)
Адаптер -- это реализация порта для конкретной технологии.
Тип адаптера
Сторона
Примеры
Primary (Driving)
Входная
REST Controller, CLI Command, GraphQL Resolver
Secondary (Driven)
Выходная
PostgreSQL Repository, Redis Cache, SMTP Mailer
Полный пример на PHP
Output Port (интерфейс)
<?php
declare(strict_types=1);
namespace App\Domain\Port;
use App\Domain\Entity\Order;
use App\Domain\ValueObject\OrderId;
// Output Port: определяет ЧТО нужно, но не КАК
interface OrderRepositoryInterface
{
public function save(Order $order): void;
public function findById(OrderId $id): ?Order;
/** @return Order[] */
public function findByCustomerId(string $customerId): array;
public function nextIdentity(): OrderId;
}
package port
import "context"
// OrderID represents a unique order identifier.
type OrderID string
// Order represents an order domain entity (defined in domain package).
type Order struct {
ID OrderID
CustomerID string
TotalAmount int
Status string
}
// OrderRepository is the output port: defines WHAT is needed, not HOW.
type OrderRepository interface {
Save(ctx context.Context, order *Order) error
FindByID(ctx context.Context, id OrderID) (*Order, error)
FindByCustomerID(ctx context.Context, customerID string) ([]*Order, error)
NextIdentity(ctx context.Context) (OrderID, error)
}
namespace App.Domain.Ports;
/// Strongly typed order identifier (value object).
public readonly record struct OrderId(Guid Value)
{
public static OrderId New() => new(Guid.CreateVersion7());
public static OrderId Parse(string raw) => new(Guid.Parse(raw));
public override string ToString() => Value.ToString();
}
// Output Port: defines WHAT is needed, not HOW.
// It lives in the domain, so the core never references Npgsql or EF Core.
public interface IOrderRepository
{
Task SaveAsync(Order order, CancellationToken ct = default);
Task<Order?> FindByIdAsync(OrderId id, CancellationToken ct = default);
Task<IReadOnlyList<Order>> FindByCustomerIdAsync(string customerId, CancellationToken ct = default);
OrderId NextIdentity();
}
import uuid
from dataclasses import dataclass
from typing import Protocol, Sequence
@dataclass(frozen=True, slots=True)
class OrderId:
"""Strongly typed order identifier (value object)."""
value: uuid.UUID
@classmethod
def generate(cls) -> "OrderId":
return cls(uuid.uuid4())
@classmethod
def parse(cls, raw: str) -> "OrderId":
return cls(uuid.UUID(raw))
def __str__(self) -> str:
return str(self.value)
class OrderRepository(Protocol):
# Output Port: defines WHAT is needed, not HOW.
# Python has no interfaces; typing.Protocol gives structural typing,
# so adapters satisfy the port without importing the domain.
async def save(self, order: Order) -> None: ...
async def find_by_id(self, order_id: OrderId) -> Order | None: ...
async def find_by_customer_id(self, customer_id: str) -> Sequence[Order]: ...
def next_identity(self) -> OrderId: ...
### Domain Entity
<?php
declare(strict_types=1);
namespace App\Domain\Entity;
use App\Domain\ValueObject\OrderId;
use App\Domain\ValueObject\Money;
use App\Domain\Event\OrderCreatedEvent;
// Domain Entity: чистая бизнес-логика, никаких зависимостей от инфраструктуры
final class Order
{
private string $status = 'pending';
/** @var OrderCreatedEvent[] */
private array $domainEvents = [];
public function __construct(
private readonly OrderId $id,
private readonly string $customerId,
private Money $totalAmount,
private readonly \DateTimeImmutable $createdAt,
) {
$this->domainEvents[] = new OrderCreatedEvent($this->id, $this->customerId);
}
public function confirm(): void
{
if ($this->status !== 'pending') {
throw new \DomainException(
sprintf('Cannot confirm order in status "%s"', $this->status)
);
}
$this->status = 'confirmed';
}
public function cancel(): void
{
if ($this->status === 'shipped') {
throw new \DomainException('Cannot cancel shipped order');
}
$this->status = 'cancelled';
}
public function applyDiscount(int $percentage): void
{
if ($percentage < 0 || $percentage > 50) {
throw new \DomainException('Discount must be between 0 and 50%');
}
$this->totalAmount = $this->totalAmount->multiply(1 - $percentage / 100);
}
public function getId(): OrderId { return $this->id; }
public function getStatus(): string { return $this->status; }
public function getTotalAmount(): Money { return $this->totalAmount; }
/** @return OrderCreatedEvent[] */
public function pullDomainEvents(): array
{
$events = $this->domainEvents;
$this->domainEvents = [];
return $events;
}
}
package domain
import (
"errors"
"fmt"
"time"
)
// OrderCreatedEvent is emitted when a new order is created.
type OrderCreatedEvent struct {
OrderID string
CustomerID string
}
// Order is a domain entity: pure business logic, no infrastructure dependencies.
type Order struct {
id string
customerID string
totalAmount Money
status string
createdAt time.Time
domainEvents []OrderCreatedEvent
}
// NewOrder creates a new order and records an OrderCreatedEvent.
func NewOrder(id, customerID string, totalAmount Money, createdAt time.Time) *Order {
o := &Order{
id: id,
customerID: customerID,
totalAmount: totalAmount,
status: "pending",
createdAt: createdAt,
}
o.domainEvents = append(o.domainEvents, OrderCreatedEvent{
OrderID: id,
CustomerID: customerID,
})
return o
}
// Confirm transitions the order to confirmed status.
func (o *Order) Confirm() error {
if o.status != "pending" {
return fmt.Errorf("cannot confirm order in status %q", o.status)
}
o.status = "confirmed"
return nil
}
// Cancel transitions the order to cancelled status.
func (o *Order) Cancel() error {
if o.status == "shipped" {
return errors.New("cannot cancel shipped order")
}
o.status = "cancelled"
return nil
}
// ApplyDiscount applies a percentage discount (0-50%).
func (o *Order) ApplyDiscount(percentage int) error {
if percentage < 0 || percentage > 50 {
return errors.New("discount must be between 0 and 50%")
}
o.totalAmount = o.totalAmount.Multiply(1 - float64(percentage)/100)
return nil
}
func (o *Order) ID() string { return o.id }
func (o *Order) Status() string { return o.status }
func (o *Order) TotalAmount() Money { return o.totalAmount }
// PullDomainEvents returns and clears recorded domain events.
func (o *Order) PullDomainEvents() []OrderCreatedEvent {
events := o.domainEvents
o.domainEvents = nil
return events
}
namespace App.Domain.Entities;
/// Emitted when a new order is created.
public sealed record OrderCreatedEvent(OrderId OrderId, string CustomerId);
public enum OrderStatus
{
Pending,
Confirmed,
Shipped,
Cancelled,
}
/// Raised when a domain invariant is violated.
public sealed class OrderDomainException : Exception
{
public OrderDomainException(string message) : base(message) { }
}
// Domain Entity: pure business logic, no infrastructure dependencies.
// No EF Core attributes here — mapping lives in the secondary adapter.
public sealed class Order
{
private const int MaxDiscountPercent = 50;
private readonly List<object> _domainEvents = [];
public Order(OrderId id, string customerId, Money totalAmount, DateTimeOffset createdAt)
{
Id = id;
CustomerId = customerId;
TotalAmount = totalAmount;
CreatedAt = createdAt;
Status = OrderStatus.Pending;
_domainEvents.Add(new OrderCreatedEvent(id, customerId));
}
public OrderId Id { get; }
public string CustomerId { get; }
public Money TotalAmount { get; private set; }
public OrderStatus Status { get; private set; }
public DateTimeOffset CreatedAt { get; }
public void Confirm()
{
if (Status is not OrderStatus.Pending)
{
throw new OrderDomainException($"Cannot confirm order in status \"{Status}\"");
}
Status = OrderStatus.Confirmed;
}
public void Cancel()
{
if (Status is OrderStatus.Shipped)
{
throw new OrderDomainException("Cannot cancel shipped order");
}
Status = OrderStatus.Cancelled;
}
public void ApplyDiscount(int percentage)
{
if (percentage is < 0 or > MaxDiscountPercent)
{
throw new OrderDomainException($"Discount must be between 0 and {MaxDiscountPercent}%");
}
TotalAmount = TotalAmount.Multiply(1 - percentage / 100m);
}
// Return and clear recorded domain events.
public IReadOnlyList<object> PullDomainEvents()
{
var events = _domainEvents.ToArray();
_domainEvents.Clear();
return events;
}
}
from dataclasses import dataclass, field
from datetime import datetime
from decimal import Decimal
from enum import Enum
MAX_DISCOUNT_PERCENT = 50
class OrderStatus(str, Enum):
PENDING = "pending"
CONFIRMED = "confirmed"
SHIPPED = "shipped"
CANCELLED = "cancelled"
@dataclass(frozen=True, slots=True)
class OrderCreatedEvent:
"""Emitted when a new order is created."""
order_id: OrderId
customer_id: str
class OrderDomainError(Exception):
"""Raised when a domain invariant is violated."""
@dataclass(slots=True)
class Order:
"""Domain entity: pure business logic, no infrastructure dependencies.
No ORM mapping or table metadata here — that belongs to the
secondary adapter, which keeps the core independent of storage.
"""
id: OrderId
customer_id: str
total_amount: Money
created_at: datetime
status: OrderStatus = OrderStatus.PENDING
_domain_events: list[object] = field(default_factory=list, repr=False)
def __post_init__(self) -> None:
self._domain_events.append(OrderCreatedEvent(self.id, self.customer_id))
def confirm(self) -> None:
if self.status is not OrderStatus.PENDING:
raise OrderDomainError(f'Cannot confirm order in status "{self.status.value}"')
self.status = OrderStatus.CONFIRMED
def cancel(self) -> None:
if self.status is OrderStatus.SHIPPED:
raise OrderDomainError("Cannot cancel shipped order")
self.status = OrderStatus.CANCELLED
def apply_discount(self, percentage: int) -> None:
if not 0 <= percentage <= MAX_DISCOUNT_PERCENT:
raise OrderDomainError(
f"Discount must be between 0 and {MAX_DISCOUNT_PERCENT}%"
)
self.total_amount = self.total_amount.multiply(
1 - Decimal(percentage) / Decimal(100)
)
def pull_domain_events(self) -> list[object]:
"""Return and clear recorded domain events."""
events = list(self._domain_events)
self._domain_events.clear()
return events
### Input Port (Use Case Interface)
<?php
declare(strict_types=1);
namespace App\Application\Port;
use App\Application\DTO\CreateOrderCommand;
use App\Application\DTO\OrderResult;
// Input Port: определяет ЧТО приложение умеет делать
interface CreateOrderUseCaseInterface
{
public function execute(CreateOrderCommand $command): OrderResult;
}
package port
import "context"
// CreateOrderCommand describes the intent to create an order.
type CreateOrderCommand struct {
CustomerID string
Amount int
Currency string
DiscountPercent int
}
// OrderResult is the response DTO for order creation.
type OrderResult struct {
ID string
Status string
Amount int
}
// CreateOrderUseCase is the input port: defines WHAT the application can do.
type CreateOrderUseCase interface {
Execute(ctx context.Context, cmd CreateOrderCommand) (OrderResult, error)
}
namespace App.Application.Ports;
/// Describes the intent to create an order.
public sealed record CreateOrderCommand(
string CustomerId,
decimal Amount,
string Currency,
int DiscountPercent = 0);
/// Response DTO returned by the use case — never a domain entity,
/// so primary adapters cannot reach into the domain model.
public sealed record OrderResult(string Id, string Status, decimal Amount);
// Input Port: defines WHAT the application can do.
// Primary adapters depend on this interface, never on the implementation.
public interface ICreateOrderUseCase
{
Task<OrderResult> ExecuteAsync(CreateOrderCommand command, CancellationToken ct = default);
}
from dataclasses import dataclass
from decimal import Decimal
from typing import Protocol
@dataclass(frozen=True, slots=True)
class CreateOrderCommand:
"""Describes the intent to create an order."""
customer_id: str
amount: Decimal
currency: str
discount_percent: int = 0
@dataclass(frozen=True, slots=True)
class OrderResult:
"""Response DTO returned by the use case — never a domain entity,
so primary adapters cannot reach into the domain model.
"""
id: str
status: str
amount: Decimal
class CreateOrderUseCase(Protocol):
# Input Port: defines WHAT the application can do.
# Python has no interfaces; typing.Protocol keeps primary adapters
# bound to the contract rather than to a concrete service class.
async def execute(self, command: CreateOrderCommand) -> OrderResult: ...
### Application Service (реализация Input Port)
<?php
declare(strict_types=1);
namespace App\Application\Service;
use App\Application\Port\CreateOrderUseCaseInterface;
use App\Application\DTO\CreateOrderCommand;
use App\Application\DTO\OrderResult;
use App\Domain\Entity\Order;
use App\Domain\Port\OrderRepositoryInterface;
use App\Domain\Port\EventDispatcherInterface;
use App\Domain\ValueObject\Money;
// Application Service: оркестрирует бизнес-логику
final readonly class CreateOrderService implements CreateOrderUseCaseInterface
{
public function __construct(
private OrderRepositoryInterface $orderRepository,
private EventDispatcherInterface $eventDispatcher,
) {}
public function execute(CreateOrderCommand $command): OrderResult
{
$orderId = $this->orderRepository->nextIdentity();
$order = new Order(
id: $orderId,
customerId: $command->customerId,
totalAmount: Money::fromAmount($command->amount, $command->currency),
createdAt: new \DateTimeImmutable(),
);
if ($command->discountPercent > 0) {
$order->applyDiscount($command->discountPercent);
}
$this->orderRepository->save($order);
// Publish domain events
foreach ($order->pullDomainEvents() as $event) {
$this->eventDispatcher->dispatch($event);
}
return new OrderResult(
id: (string) $orderId,
status: $order->getStatus(),
amount: $order->getTotalAmount()->getAmount(),
);
}
}
package service
import (
"context"
"fmt"
"time"
"app/domain"
"app/port"
)
// EventDispatcher publishes domain events.
type EventDispatcher interface {
Dispatch(ctx context.Context, event any) error
}
// CreateOrderService is the application service that orchestrates business logic.
type CreateOrderService struct {
repo port.OrderRepository
dispatcher EventDispatcher
}
// NewCreateOrderService creates a new service with injected dependencies.
func NewCreateOrderService(repo port.OrderRepository, dispatcher EventDispatcher) *CreateOrderService {
return &CreateOrderService{repo: repo, dispatcher: dispatcher}
}
// Execute implements the CreateOrderUseCase input port.
func (s *CreateOrderService) Execute(ctx context.Context, cmd port.CreateOrderCommand) (port.OrderResult, error) {
orderID, err := s.repo.NextIdentity(ctx)
if err != nil {
return port.OrderResult{}, fmt.Errorf("generate order id: %w", err)
}
totalAmount := domain.MoneyFromAmount(cmd.Amount, cmd.Currency)
order := domain.NewOrder(string(orderID), cmd.CustomerID, totalAmount, time.Now())
if cmd.DiscountPercent > 0 {
if err := order.ApplyDiscount(cmd.DiscountPercent); err != nil {
return port.OrderResult{}, fmt.Errorf("apply discount: %w", err)
}
}
if err := s.repo.Save(ctx, &port.Order{
ID: orderID,
CustomerID: cmd.CustomerID,
TotalAmount: order.TotalAmount().Amount(),
Status: order.Status(),
}); err != nil {
return port.OrderResult{}, fmt.Errorf("save order: %w", err)
}
// Publish domain events.
for _, event := range order.PullDomainEvents() {
if err := s.dispatcher.Dispatch(ctx, event); err != nil {
return port.OrderResult{}, fmt.Errorf("dispatch event: %w", err)
}
}
return port.OrderResult{
ID: string(orderID),
Status: order.Status(),
Amount: order.TotalAmount().Amount(),
}, nil
}
namespace App.Application.Services;
/// Output port for publishing domain events.
public interface IEventDispatcher
{
Task DispatchAsync(object domainEvent, CancellationToken ct = default);
}
// Application Service: orchestrates business logic and implements the input port.
// Both dependencies are output ports — no infrastructure type is referenced here.
public sealed class CreateOrderService : ICreateOrderUseCase
{
private readonly IOrderRepository _orderRepository;
private readonly IEventDispatcher _eventDispatcher;
private readonly TimeProvider _clock;
// Constructor injection: the DI container wires concrete adapters at startup.
public CreateOrderService(
IOrderRepository orderRepository,
IEventDispatcher eventDispatcher,
TimeProvider clock)
{
_orderRepository = orderRepository;
_eventDispatcher = eventDispatcher;
_clock = clock;
}
public async Task<OrderResult> ExecuteAsync(
CreateOrderCommand command,
CancellationToken ct = default)
{
var orderId = _orderRepository.NextIdentity();
var order = new Order(
orderId,
command.CustomerId,
Money.FromAmount(command.Amount, command.Currency),
_clock.GetUtcNow());
if (command.DiscountPercent > 0)
{
order.ApplyDiscount(command.DiscountPercent);
}
await _orderRepository.SaveAsync(order, ct);
// Publish domain events
foreach (var domainEvent in order.PullDomainEvents())
{
await _eventDispatcher.DispatchAsync(domainEvent, ct);
}
return new OrderResult(
orderId.ToString(),
order.Status.ToString(),
order.TotalAmount.Amount);
}
}
from datetime import datetime, timezone
from typing import Protocol
class EventDispatcher(Protocol):
"""Output port for publishing domain events."""
async def dispatch(self, domain_event: object) -> None: ...
class CreateOrderService:
"""Application service: orchestrates business logic and implements
the CreateOrderUseCase input port.
Both dependencies are output ports, so no infrastructure module is
imported here — that keeps the core independent of storage and messaging.
"""
# Python has no compile-time DI container; the composition root
# constructs the adapters and passes them in here.
def __init__(
self,
order_repository: OrderRepository,
event_dispatcher: EventDispatcher,
) -> None:
self._order_repository = order_repository
self._event_dispatcher = event_dispatcher
async def execute(self, command: CreateOrderCommand) -> OrderResult:
order_id = self._order_repository.next_identity()
order = Order(
id=order_id,
customer_id=command.customer_id,
total_amount=Money.from_amount(command.amount, command.currency),
created_at=datetime.now(timezone.utc),
)
if command.discount_percent > 0:
order.apply_discount(command.discount_percent)
await self._order_repository.save(order)
# Publish domain events
for domain_event in order.pull_domain_events():
await self._event_dispatcher.dispatch(domain_event)
return OrderResult(
id=str(order_id),
status=order.status.value,
amount=order.total_amount.amount,
)
### Primary Adapter (REST Controller)
<?php
declare(strict_types=1);
namespace App\Infrastructure\Adapter\Primary;
use App\Application\Port\CreateOrderUseCaseInterface;
use App\Application\DTO\CreateOrderCommand;
use Symfony\Component\HttpFoundation\JsonResponse;
use Symfony\Component\HttpFoundation\Response;
use Symfony\Component\Routing\Attribute\Route;
use Symfony\Component\HttpKernel\Attribute\MapRequestPayload;
// Primary Adapter: преобразует HTTP-запрос в команду для Use Case
final class OrderHttpAdapter
{
public function __construct(
private readonly CreateOrderUseCaseInterface $createOrder,
) {}
#[Route('/api/orders', methods: ['POST'])]
public function create(
#[MapRequestPayload] CreateOrderCommand $command,
): JsonResponse {
$result = $this->createOrder->execute($command);
return new JsonResponse($result, Response::HTTP_CREATED);
}
}
package primary
import (
"encoding/json"
"net/http"
"app/port"
)
// OrderHTTPAdapter is a primary adapter that transforms HTTP requests into use case commands.
type OrderHTTPAdapter struct {
createOrder port.CreateOrderUseCase
}
// NewOrderHTTPAdapter creates a new HTTP adapter for orders.
func NewOrderHTTPAdapter(createOrder port.CreateOrderUseCase) *OrderHTTPAdapter {
return &OrderHTTPAdapter{createOrder: createOrder}
}
// CreateOrder handles POST /api/orders.
func (a *OrderHTTPAdapter) CreateOrder(w http.ResponseWriter, r *http.Request) {
var cmd port.CreateOrderCommand
if err := json.NewDecoder(r.Body).Decode(&cmd); err != nil {
http.Error(w, "invalid request body", http.StatusBadRequest)
return
}
result, err := a.createOrder.Execute(r.Context(), cmd)
if err != nil {
http.Error(w, err.Error(), http.StatusUnprocessableEntity)
return
}
w.Header().Set("Content-Type", "application/json")
w.WriteHeader(http.StatusCreated)
json.NewEncoder(w).Encode(result)
}
using Microsoft.AspNetCore.Mvc;
namespace App.Infrastructure.Adapters.Primary;
// Primary (driving) Adapter: transforms an HTTP request into a use case command.
// It depends only on the input port, so the same use case is reachable
// from a CLI or a queue consumer through different adapters.
[ApiController]
[Route("api/orders")]
public sealed class OrderHttpAdapter : ControllerBase
{
private readonly ICreateOrderUseCase _createOrder;
public OrderHttpAdapter(ICreateOrderUseCase createOrder) => _createOrder = createOrder;
[HttpPost]
public async Task<IActionResult> Create(
[FromBody] CreateOrderRequest request,
CancellationToken ct)
{
// Map the transport DTO to the application command —
// HTTP concerns never leak into the core.
var command = new CreateOrderCommand(
request.CustomerId,
request.Amount,
request.Currency,
request.DiscountPercent);
var result = await _createOrder.ExecuteAsync(command, ct);
return StatusCode(StatusCodes.Status201Created, result);
}
/// Transport-level DTO owned by the adapter, not by the core.
public sealed record CreateOrderRequest(
string CustomerId,
decimal Amount,
string Currency,
int DiscountPercent = 0);
}
from dataclasses import dataclass
from decimal import Decimal
from typing import Annotated
from fastapi import APIRouter, Depends, status
from pydantic import BaseModel
router = APIRouter(prefix="/api/orders")
class CreateOrderRequest(BaseModel):
"""Transport-level DTO owned by the adapter, not by the core."""
customer_id: str
amount: Decimal
currency: str
discount_percent: int = 0
# Primary (driving) Adapter: transforms an HTTP request into a use case command.
# It depends only on the input port, so the same use case is reachable
# from a CLI or a queue consumer through different adapters.
@router.post("", status_code=status.HTTP_201_CREATED)
async def create_order(
request: CreateOrderRequest,
create_order_use_case: Annotated[CreateOrderUseCase, Depends(get_create_order_use_case)],
) -> OrderResult:
# Map the transport DTO to the application command —
# HTTP concerns never leak into the core.
command = CreateOrderCommand(
customer_id=request.customer_id,
amount=request.amount,
currency=request.currency,
discount_percent=request.discount_percent,
)
return await create_order_use_case.execute(command)
### Secondary Adapter (PostgreSQL Repository)
<?php
declare(strict_types=1);
namespace App\Infrastructure\Adapter\Secondary;
use App\Domain\Entity\Order;
use App\Domain\Port\OrderRepositoryInterface;
use App\Domain\ValueObject\OrderId;
use Doctrine\DBAL\Connection;
// Secondary Adapter: реализует Output Port через PostgreSQL
final readonly class PostgresOrderRepository implements OrderRepositoryInterface
{
public function __construct(
private Connection $connection,
) {}
public function save(Order $order): void
{
$this->connection->executeStatement(
'INSERT INTO orders (id, customer_id, total_amount, status, created_at)
VALUES (:id, :customer_id, :amount, :status, :created_at)
ON CONFLICT (id) DO UPDATE SET
status = :status,
total_amount = :amount',
[
'id' => (string) $order->getId(),
'customer_id' => $order->getCustomerId(),
'amount' => $order->getTotalAmount()->getAmount(),
'status' => $order->getStatus(),
'created_at' => $order->getCreatedAt()->format('c'),
]
);
}
public function findById(OrderId $id): ?Order
{
$row = $this->connection->fetchAssociative(
'SELECT * FROM orders WHERE id = :id',
['id' => (string) $id]
);
return $row ? $this->hydrate($row) : null;
}
public function findByCustomerId(string $customerId): array
{
$rows = $this->connection->fetchAllAssociative(
'SELECT * FROM orders WHERE customer_id = :cid ORDER BY created_at DESC',
['cid' => $customerId]
);
return array_map($this->hydrate(...), $rows);
}
public function nextIdentity(): OrderId
{
return OrderId::generate();
}
private function hydrate(array $row): Order
{
// Reconstruct domain entity from database row
return new Order(
id: OrderId::fromString($row['id']),
customerId: $row['customer_id'],
totalAmount: Money::fromAmount((int) $row['total_amount'], 'RUB'),
createdAt: new \DateTimeImmutable($row['created_at']),
);
}
}
package secondary
import (
"context"
"database/sql"
"fmt"
"github.com/google/uuid"
"app/port"
)
// PostgresOrderRepository is a secondary adapter implementing the output port via PostgreSQL.
type PostgresOrderRepository struct {
db *sql.DB
}
// NewPostgresOrderRepository creates a new repository with a database connection.
func NewPostgresOrderRepository(db *sql.DB) *PostgresOrderRepository {
return &PostgresOrderRepository{db: db}
}
// Save persists an order using upsert.
func (r *PostgresOrderRepository) Save(ctx context.Context, order *port.Order) error {
_, err := r.db.ExecContext(ctx,
`INSERT INTO orders (id, customer_id, total_amount, status, created_at)
VALUES ($1, $2, $3, $4, NOW())
ON CONFLICT (id) DO UPDATE SET status = $4, total_amount = $3`,
order.ID, order.CustomerID, order.TotalAmount, order.Status,
)
if err != nil {
return fmt.Errorf("save order: %w", err)
}
return nil
}
// FindByID retrieves an order by its identifier.
func (r *PostgresOrderRepository) FindByID(ctx context.Context, id port.OrderID) (*port.Order, error) {
row := r.db.QueryRowContext(ctx,
"SELECT id, customer_id, total_amount, status FROM orders WHERE id = $1", id,
)
var o port.Order
if err := row.Scan(&o.ID, &o.CustomerID, &o.TotalAmount, &o.Status); err != nil {
if err == sql.ErrNoRows {
return nil, nil
}
return nil, fmt.Errorf("find order by id: %w", err)
}
return &o, nil
}
// FindByCustomerID retrieves all orders for a customer.
func (r *PostgresOrderRepository) FindByCustomerID(ctx context.Context, customerID string) ([]*port.Order, error) {
rows, err := r.db.QueryContext(ctx,
"SELECT id, customer_id, total_amount, status FROM orders WHERE customer_id = $1 ORDER BY created_at DESC",
customerID,
)
if err != nil {
return nil, fmt.Errorf("find orders by customer: %w", err)
}
defer rows.Close()
var orders []*port.Order
for rows.Next() {
var o port.Order
if err := rows.Scan(&o.ID, &o.CustomerID, &o.TotalAmount, &o.Status); err != nil {
return nil, fmt.Errorf("scan order row: %w", err)
}
orders = append(orders, &o)
}
return orders, rows.Err()
}
// NextIdentity generates a new unique order identifier.
func (r *PostgresOrderRepository) NextIdentity(_ context.Context) (port.OrderID, error) {
return port.OrderID(uuid.New().String()), nil
}
using Dapper;
using Npgsql;
namespace App.Infrastructure.Adapters.Secondary;
// Secondary (driven) Adapter: implements the output port via PostgreSQL.
// Swapping this for MongoDB or an in-memory fake touches nothing in the core.
public sealed class PostgresOrderRepository : IOrderRepository
{
private readonly NpgsqlDataSource _db;
public PostgresOrderRepository(NpgsqlDataSource db) => _db = db;
public async Task SaveAsync(Order order, CancellationToken ct = default)
{
const string sql = """
INSERT INTO orders (id, customer_id, total_amount, currency, status, created_at)
VALUES (@Id, @CustomerId, @Amount, @Currency, @Status, @CreatedAt)
ON CONFLICT (id) DO UPDATE SET
status = EXCLUDED.status,
total_amount = EXCLUDED.total_amount
""";
await using var connection = await _db.OpenConnectionAsync(ct);
await connection.ExecuteAsync(new CommandDefinition(
sql,
new
{
Id = order.Id.Value,
order.CustomerId,
Amount = order.TotalAmount.Amount,
Currency = order.TotalAmount.Currency,
Status = order.Status.ToString(),
order.CreatedAt,
},
cancellationToken: ct));
}
public async Task<Order?> FindByIdAsync(OrderId id, CancellationToken ct = default)
{
const string sql = """
SELECT id, customer_id, total_amount, currency, status, created_at
FROM orders
WHERE id = @id
""";
await using var connection = await _db.OpenConnectionAsync(ct);
var row = await connection.QuerySingleOrDefaultAsync<OrderRow>(
new CommandDefinition(sql, new { id = id.Value }, cancellationToken: ct));
return row is null ? null : Hydrate(row);
}
public async Task<IReadOnlyList<Order>> FindByCustomerIdAsync(
string customerId,
CancellationToken ct = default)
{
const string sql = """
SELECT id, customer_id, total_amount, currency, status, created_at
FROM orders
WHERE customer_id = @customerId
ORDER BY created_at DESC
""";
await using var connection = await _db.OpenConnectionAsync(ct);
var rows = await connection.QueryAsync<OrderRow>(
new CommandDefinition(sql, new { customerId }, cancellationToken: ct));
return rows.Select(Hydrate).ToArray();
}
public OrderId NextIdentity() => OrderId.New();
// Reconstruct the domain entity from a database row.
// Mapping lives here so the entity stays free of persistence concerns.
private static Order Hydrate(OrderRow row)
{
var order = new Order(
new OrderId(row.Id),
row.CustomerId,
Money.FromAmount(row.TotalAmount, row.Currency),
row.CreatedAt);
// Drop the creation event: this instance is restored, not newly created.
order.PullDomainEvents();
return order;
}
private sealed record OrderRow(
Guid Id,
string CustomerId,
decimal TotalAmount,
string Currency,
string Status,
DateTimeOffset CreatedAt);
}
from datetime import datetime
from decimal import Decimal
import asyncpg
class PostgresOrderRepository:
"""Secondary (driven) Adapter: implements the OrderRepository output port
via PostgreSQL.
It satisfies the port structurally — no base class to inherit — so
swapping in an in-memory fake for tests touches nothing in the core.
"""
def __init__(self, pool: asyncpg.Pool) -> None:
self._pool = pool
async def save(self, order: Order) -> None:
await self._pool.execute(
"""
INSERT INTO orders (id, customer_id, total_amount, currency, status, created_at)
VALUES ($1, $2, $3, $4, $5, $6)
ON CONFLICT (id) DO UPDATE SET
status = EXCLUDED.status,
total_amount = EXCLUDED.total_amount
""",
order.id.value,
order.customer_id,
order.total_amount.amount,
order.total_amount.currency,
order.status.value,
order.created_at,
)
async def find_by_id(self, order_id: OrderId) -> Order | None:
row = await self._pool.fetchrow(
"""
SELECT id, customer_id, total_amount, currency, status, created_at
FROM orders
WHERE id = $1
""",
order_id.value,
)
return self._hydrate(row) if row else None
async def find_by_customer_id(self, customer_id: str) -> list[Order]:
rows = await self._pool.fetch(
"""
SELECT id, customer_id, total_amount, currency, status, created_at
FROM orders
WHERE customer_id = $1
ORDER BY created_at DESC
""",
customer_id,
)
return [self._hydrate(row) for row in rows]
def next_identity(self) -> OrderId:
return OrderId.generate()
@staticmethod
def _hydrate(row: asyncpg.Record) -> Order:
"""Reconstruct the domain entity from a database row.
Mapping lives here so the entity stays free of persistence concerns.
"""
order = Order(
id=OrderId(row["id"]),
customer_id=row["customer_id"],
total_amount=Money.from_amount(row["total_amount"], row["currency"]),
created_at=row["created_at"],
status=OrderStatus(row["status"]),
)
# Drop the creation event: this instance is restored, not newly created.
order.pull_domain_events()
return order