Туториали

Reliably Integrating PHP Services with Transactional Outbox and Inbox Patterns

Сигурна интеграција на PHP услуги со обрасците за трансакциско сандаче за испраќање и сандаче за примање

Опасниот момент во интеграцијата на услуги не е HTTP-барањето. Тоа е потврдувањето непосредно пред или по него. Ако нарачката е потврдена, но нејзиниот настан никогаш не е испратен, надолните системи трајно остануваат несвесни за неа. Ако настанот е испратен пред потврдувањето и трансакцијата подоцна се поништи, потрошувачите постапуваат по нарачка што не постои.

Шаблонот за трансакциско сандаче за испраќање го затвора тој јаз со складирање на деловната промена и настан во една трансакција во базата на податоци. Посебен пренесувач ги доставува потврдените настани. Шаблонот за приемно сандаче го заокружува дизајнот така што ефектот во базата на податоци на потрошувачот го прави идемпотентен.

Ова упатство гради две PHP 8.3 услуги поддржани од PostgreSQL: услуга Orders со сандаче за испраќање и услуга Billing со приемно сандаче. Испораката намерно е опишана како најмалку еднаш. Повторните обиди можат да создадат дупликат-барања, но не можат да создадат дупликат-фактури.

Предуслови и архитектура

Потребни ви се PHP 8.3 со PDO PostgreSQL, cURL и JSON екстензии; PostgreSQL 14 или понов; и пристап од командна линија до двете бази на податоци. Продукцискиот HTTP-сообраќај треба да користи TLS, иако локалните команди за проверка користат loopback HTTP.

Текот е:

  1. Orders API внесува нарачка и настан OrderCreated во една трансакција.
  2. Пренесувач накратко презема ред од сандачето за испраќање, го потврдува преземањето, а потоа го извршува мрежното барање.
  3. Billing API го внесува ID-то на настанот во своето приемно сандаче и ја создава фактурата во една трансакција.
  4. Пренесувачот го означува настанот како објавен само откако Billing ќе врати успех.

Пренесувачот никогаш не држи отворена трансакција во базата на податоци за време на HTTP. Неговиот закуп овозможува опоравување по пад, додека FOR UPDATE SKIP LOCKED им овозможува на повеќе процеси на пренесувачот да преземат различни редови.

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

reliable-integration/
  config.php
  orders.php
  inbox.php
  relay.php
  orders.sql
  billing.sql

Создадете ги базите на податоци

Во продукција користете одделни бази на податоци и акредитиви, така што ниту една услуга не може да ги менува табелите на другата услуга. Следните шеми припаѓаат во orders.sql и billing.sql, соодветно.

-- orders.sql
CREATE TABLE orders (
    id uuid PRIMARY KEY,
    amount_cents bigint NOT NULL CHECK (amount_cents > 0),
    currency char(3) NOT NULL,
    created_at timestamptz NOT NULL DEFAULT now()
);

CREATE TABLE outbox (
    id uuid PRIMARY KEY,
    aggregate_id uuid NOT NULL REFERENCES orders(id),
    event_type text NOT NULL,
    payload jsonb NOT NULL,
    occurred_at timestamptz NOT NULL DEFAULT now(),
    available_at timestamptz NOT NULL DEFAULT now(),
    attempts integer NOT NULL DEFAULT 0,
    claimed_by text,
    lease_until timestamptz,
    published_at timestamptz
);

CREATE INDEX outbox_ready_idx
    ON outbox (available_at, occurred_at)
    WHERE published_at IS NULL;

-- billing.sql
CREATE TABLE processed_events (
    event_id uuid PRIMARY KEY,
    processed_at timestamptz NOT NULL DEFAULT now()
);

CREATE TABLE invoices (
    id uuid PRIMARY KEY,
    order_id uuid NOT NULL UNIQUE,
    amount_cents bigint NOT NULL CHECK (amount_cents > 0),
    currency char(3) NOT NULL,
    created_at timestamptz NOT NULL DEFAULT now()
);

Создадете две PostgreSQL бази на податоци со вашиот вообичаен административен процес, а потоа применете ја секоја датотека со сметка со соодветни привилегии:

psql 'postgresql://[email protected]/orders' -f orders.sql
psql 'postgresql://[email protected]/billing' -f billing.sql

php -m | grep -E 'curl|json|pdo_pgsql'

На апликациските сметки им се потребни само права за поврзување, користење на шема, користење на секвенци каде што е применливо, и SELECT, INSERT и UPDATE врз сопствените табели. Тие не треба да ги поседуваат базите на податоци ниту да добијат привилегии за создавање шема.

Заедничка конфигурација и ограничени повици кон базата на податоци

Ставете ги создавањето врска и генерирањето UUID во config.php. Истекот на време за поврзување не ги ограничува барањата, па кодот исто така ги конфигурира истеците на време за PostgreSQL изјави и заклучувања.

<?php
declare(strict_types=1);

function requiredEnv(string $name): string
{
    $value = getenv($name);
    if ($value === false || $value === '') {
        throw new RuntimeException("Missing environment variable: {$name}");
    }
    return $value;
}

function database(string $dsnName): PDO
{
    $pdo = new PDO(
        requiredEnv($dsnName),
        requiredEnv($dsnName . '_USER'),
        requiredEnv($dsnName . '_PASSWORD'),
        [
            PDO::ATTR_ERRMODE => PDO::ERRMODE_EXCEPTION,
            PDO::ATTR_DEFAULT_FETCH_MODE => PDO::FETCH_ASSOC,
            PDO::ATTR_EMULATE_PREPARES => false,
            PDO::ATTR_PERSISTENT => false,
        ]
    );

    $pdo->exec("SET statement_timeout = '3000ms'");
    $pdo->exec("SET lock_timeout = '1000ms'");
    return $pdo;
}

function uuidV4(): string
{
    $bytes = random_bytes(16);
    $bytes[6] = chr((ord($bytes[6]) & 0x0f) | 0x40);
    $bytes[8] = chr((ord($bytes[8]) & 0x3f) | 0x80);
    return vsprintf('%s%s-%s-%s-%s-%s%s%s', str_split(bin2hex($bytes), 4));
}

function jsonResponse(int $status, array $body = []): never
{
    http_response_code($status);
    header('Content-Type: application/json');
    echo json_encode($body, JSON_THROW_ON_ERROR);
    exit;
}

Вклучете connect_timeout=3 во секој PostgreSQL DSN, на пример pgsql:host=127.0.0.1;port=5432;dbname=orders;connect_timeout=3. Истекот на време за изјава од три секунди е независен и експлицитен.

Запишете ги нарачката и настанот атомски

Крајната точка Orders валидира мала JSON-команда и ги запишува двата записи во една трансакција. Зачувајте ја како orders.php.

<?php
declare(strict_types=1);
require __DIR__ . '/config.php';

if ($_SERVER['REQUEST_METHOD'] !== 'POST') {
    header('Allow: POST');
    jsonResponse(405, ['error' => 'method_not_allowed']);
}

try {
    $input = json_decode(file_get_contents('php://input'), true, 32,
        JSON_THROW_ON_ERROR);
    $amount = filter_var($input['amount_cents'] ?? null, FILTER_VALIDATE_INT);
    $currency = strtoupper((string) ($input['currency'] ?? ''));

    if ($amount === false || $amount < 1 ||
        preg_match('/^[A-Z]{3}$/', $currency) !== 1) {
        jsonResponse(422, ['error' => 'invalid_order']);
    }

    $db = database('ORDERS_DSN');
    $orderId = uuidV4();
    $eventId = uuidV4();
    $payload = json_encode([
        'order_id' => $orderId,
        'amount_cents' => $amount,
        'currency' => $currency,
    ], JSON_THROW_ON_ERROR);

    $db->beginTransaction();

    $stmt = $db->prepare(
        'INSERT INTO orders (id, amount_cents, currency) VALUES (?, ?, ?)'
    );
    $stmt->execute([$orderId, $amount, $currency]);

    $stmt = $db->prepare(
        'INSERT INTO outbox
         (id, aggregate_id, event_type, payload)
         VALUES (?, ?, ?, CAST(? AS jsonb))'
    );
    $stmt->execute([$eventId, $orderId, 'OrderCreated', $payload]);

    $db->commit();
    jsonResponse(201, ['order_id' => $orderId, 'event_id' => $eventId]);
} catch (JsonException) {
    jsonResponse(400, ['error' => 'invalid_json']);
} catch (Throwable $error) {
    if (isset($db) && $db->inTransaction()) {
        $db->rollBack();
    }
    error_log($error->getMessage());
    jsonResponse(500, ['error' => 'internal_error']);
}

Прекин во базата на податоци го неуспешува барањето наместо да прифати нарачка без настан. Обратно, прекин на пренесувачот не го блокира создавањето нарачки; потврдените настани остануваат достапни за подоцнежна испорака.

Направете го Billing идемпотентен со приемно сандаче

Пренесувачот го автентицира точниот текст на барањето со HMAC-SHA256. Billing го запишува ID-то на настанот пред да го примени неговиот ефект, во истата трансакција. Зачувајте го ова како inbox.php.

<?php
declare(strict_types=1);
require __DIR__ . '/config.php';

if ($_SERVER['REQUEST_METHOD'] !== 'POST') {
    header('Allow: POST');
    jsonResponse(405, ['error' => 'method_not_allowed']);
}

$body = file_get_contents('php://input');
$provided = $_SERVER['HTTP_X_EVENT_SIGNATURE'] ?? '';
$expected = hash_hmac('sha256', $body, requiredEnv('EVENT_SECRET'));

if (!hash_equals($expected, $provided)) {
    jsonResponse(401, ['error' => 'invalid_signature']);
}

try {
    $event = json_decode($body, true, 32, JSON_THROW_ON_ERROR);
    $id = (string) ($event['id'] ?? '');
    $type = (string) ($event['type'] ?? '');
    $data = $event['data'] ?? [];

    if ($type !== 'OrderCreated' ||
        preg_match('/^[0-9a-f-]{36}$/D', $id) !== 1 ||
        preg_match('/^[0-9a-f-]{36}$/D', $data['order_id'] ?? '') !== 1 ||
        !is_int($data['amount_cents'] ?? null) ||
        ($data['amount_cents'] ?? 0) < 1 ||
        preg_match('/^[A-Z]{3}$/D', $data['currency'] ?? '') !== 1) {
        jsonResponse(422, ['error' => 'invalid_event']);
    }

    $db = database('BILLING_DSN');
    $db->beginTransaction();

    $insert = $db->prepare(
        'INSERT INTO processed_events (event_id)
         VALUES (?) ON CONFLICT DO NOTHING RETURNING event_id'
    );
    $insert->execute([$id]);

    if ($insert->fetchColumn() !== false) {
        $invoice = $db->prepare(
            'INSERT INTO invoices
             (id, order_id, amount_cents, currency)
             VALUES (?, ?, ?, ?)'
        );
        $invoice->execute([
            uuidV4(),
            $data['order_id'],
            $data['amount_cents'],
            $data['currency'],
        ]);
    }

    $db->commit();
    http_response_code(204);
} catch (JsonException) {
    jsonResponse(400, ['error' => 'invalid_json']);
} catch (Throwable $error) {
    if (isset($db) && $db->inTransaction()) {
        $db->rollBack();
    }
    error_log($error->getMessage());
    jsonResponse(500, ['error' => 'internal_error']);
}

Ако создавањето фактура не успее, и внесувањето во приемното сандаче се поништува. Ако Billing потврди, но неговиот одговор се изгуби, пренесувачот повторно го испраќа настанот; конфликтот во приемното сандаче го претвора тој повторен обид во успешно дејство без операција. Единственото ограничување на invoices.order_id обезбедува дополнителна инваријанта, а не замена за приемното сандаче.

Изградете пренесувач со кратки трансакции

Зачувајте го работникот како relay.php. Неговиот закуп од 20 секунди е удобно подолг од истекот на време за поврзување од една секунда и вкупниот HTTP-истек од пет секунди.

<?php
declare(strict_types=1);
require __DIR__ . '/config.php';

pcntl_async_signals(true);
$stopping = false;
pcntl_signal(SIGTERM, function () use (&$stopping): void {
    $stopping = true;
});
pcntl_signal(SIGINT, function () use (&$stopping): void {
    $stopping = true;
});

$db = database('ORDERS_DSN');
$worker = gethostname() . ':' . getmypid();
$url = requiredEnv('BILLING_URL');
$secret = requiredEnv('EVENT_SECRET');

while (!$stopping) {
    $db->beginTransaction();
    $claim = $db->prepare(
        "WITH picked AS (
            SELECT id FROM outbox
            WHERE published_at IS NULL
              AND available_at <= now()
              AND (lease_until IS NULL OR lease_until < now())
            ORDER BY occurred_at
            FOR UPDATE SKIP LOCKED
            LIMIT 1
         )
         UPDATE outbox AS o
         SET claimed_by = ?, lease_until = now() + interval '20 seconds',
             attempts = attempts + 1
         FROM picked
         WHERE o.id = picked.id
         RETURNING o.id, o.event_type, o.payload::text, o.attempts"
    );
    $claim->execute([$worker]);
    $event = $claim->fetch();
    $db->commit();

    if ($event === false) {
        usleep(250000);
        continue;
    }

    if ($stopping) {
        $release = $db->prepare(
            'UPDATE outbox SET claimed_by = NULL, lease_until = NULL
             WHERE id = ? AND claimed_by = ? AND published_at IS NULL'
        );
        $release->execute([$event['id'], $worker]);
        break;
    }

    $body = json_encode([
        'id' => $event['id'],
        'type' => $event['event_type'],
        'data' => json_decode($event['payload'], true, 32,
            JSON_THROW_ON_ERROR),
    ], JSON_THROW_ON_ERROR);

    $curl = curl_init($url);
    curl_setopt_array($curl, [
        CURLOPT_POST => true,
        CURLOPT_POSTFIELDS => $body,
        CURLOPT_HTTPHEADER => [
            'Content-Type: application/json',
            'X-Event-Signature: ' . hash_hmac('sha256', $body, $secret),
        ],
        CURLOPT_RETURNTRANSFER => true,
        CURLOPT_CONNECTTIMEOUT => 1,
        CURLOPT_TIMEOUT => 5,
    ]);
    curl_exec($curl);
    $status = curl_getinfo($curl, CURLINFO_RESPONSE_CODE);
    $ok = curl_errno($curl) === 0 && $status >= 200 && $status < 300;
    curl_close($curl);

    if ($ok) {
        $done = $db->prepare(
            'UPDATE outbox
             SET published_at = now(), claimed_by = NULL, lease_until = NULL
             WHERE id = ? AND claimed_by = ? AND lease_until > now()'
        );
        $done->execute([$event['id'], $worker]);
    } else {
        $delay = min(60, 2 ** min((int) $event['attempts'], 6));
        $retry = $db->prepare(
            "UPDATE outbox
             SET claimed_by = NULL, lease_until = NULL,
                 available_at = now() + CAST(? AS integer) * interval '1 second'
             WHERE id = ? AND claimed_by = ?"
        );
        $retry->execute([$delay, $event['id'], $worker]);
    }
}

Условите за сопственост се важни. Ако барањето го надживее својот закуп, друг работник може повторно да го преземе настанот. Првобитниот работник не смее да го означи преземањето на поновиот работник како завршено. Дупликат-испораката останува безбедна бидејќи Billing ја поседува границата на идемпотентност.

Извршете и тестирајте ја целосната патека

Користете одделни терминали за овие команди. Променливите на околината ги држат акредитивите надвор од контролата на изворниот код, иако продукциски управувач со тајни или заштитена датотека со околина е попожелна од историјата на школката.

export ORDERS_DSN='pgsql:host=127.0.0.1;port=5432;dbname=orders;connect_timeout=3'
export ORDERS_DSN_USER='orders_app'
export ORDERS_DSN_PASSWORD='replace-locally'
export BILLING_DSN='pgsql:host=127.0.0.1;port=5432;dbname=billing;connect_timeout=3'
export BILLING_DSN_USER='billing_app'
export BILLING_DSN_PASSWORD='replace-locally'
export EVENT_SECRET='replace-with-a-long-random-secret'
export BILLING_URL='http://127.0.0.1:8081/inbox.php'

php -S 127.0.0.1:8080
php -S 127.0.0.1:8081
php relay.php

curl --fail-with-body \
  -H 'Content-Type: application/json' \
  --data '{"amount_cents":2599,"currency":"EUR"}' \
  http://127.0.0.1:8080/orders.php

Испитајте ги двете бази на податоци и потврдете еден објавен ред во сандачето за испраќање, еден ред во приемното сандаче и една фактура. За да тестирате опоравување, запрете го Billing, создадете друга нарачка и набљудувајте ги повторните обиди со растечко available_at. Рестартирајте го Billing и потврдете евентуална обработка. За да тестирате дедупликација, привремено поставете го published_at на настанот на null во тест-база за еднократна употреба, повторно извршете го пренесувачот и потврдете дека бројот на фактури не се зголемува.

Продукциска безбедност, набљудливост и перформанси

Заменете го развојниот сервер на PHP со поддржан веб-сервер и PHP-FPM. Не изложувајте ги јавно ниту PostgreSQL ниту крајната точка Billing. Дозволете ја само потребната патека услуга-до-услуга во заштитните ѕидови на хостот и облакот, користете TLS, ротирајте ја HMAC-тајната и разгледајте взаемен TLS таму каде што тоа го оправдува идентитетот на услугата. Применете ограничувања на големината на барањата пред PHP и никогаш не евидентирајте тајни или целосни чувствителни содржини.

Извршувајте го пренесувачот под посветена непривилегирана сметка. Конфигурирајте го управувачот со услуги со истек на запирање поголем од HTTP-буџетот од пет секунди плус времето за чистење на базата на податоци. Повеќе инстанци се безбедни, но зголемувајте ги само откако ќе ги измерите капацитетот надолу по текот и спорноста во базата на податоци.

Корисните метрики вклучуваат број на необјавени редови, старост на најстариот подготвен настан, обиди по тип на настан, истекувања на закуп, латентност на испорака, статус на одговор, конфликти во приемното сандаче и неуспеси при обработка. Алармите треба да ја нагласуваат староста на најстариот настан наместо само длабочината на редицата: мала редица што содржи еден трајно заглавен настан е оперативно значајна.

За поголем проток, преземете мала серија во една кратка трансакција, а потоа обработете ја надвор од трансакцијата со ограничена конкурентност. Одржувајте ги закупите подолги од најлошото дозволено времетраење на барањето, додајте варијација на доцнењата при повторни обиди и дефинирајте политика за отровни настани. Не обидувајте се бесконечно повторно со трајни неуспеси на валидацијата 4xx; ставете ги во карантин со нивните метаподатоци за грешка за контролирана истрага. Продолжете да се обидувате повторно со минливи мрежни неуспеси и соодветни одговори 5xx.

Вообичаени режими на неуспех

  • Објавување во трансакцијата на нарачката: мрежната латентност ги продолжува заклучувањата, а поништувањето може да остави веќе испорачан настан.
  • Држење преземања при повикување Billing: бавниот HTTP го претвора сандачето за испраќање во систем со спорност за заклучувања.
  • Означување успех без проверка на сопственоста: работник со истечен закуп може да ја презапише состојбата на поновиот носител на закупот.
  • Дедуплицирање надвор од трансакцијата на потрошувачот: пад помеѓу внесувањето во приемното сандаче и деловниот ефект може да потисне незавршена работа.
  • Претпоставка дека повторните обиди подразбираат ефекти точно еднаш: надворешните ефекти како повици за е-пошта или плаќање бараат сопствени клучеви за идемпотентност и трајни машини на состојби.
  • Користење едно истекување на време насекаде: истеците на време за поврзување, барање, заклучување и HTTP штитат различни граници и мора да се вклопат под буџетите за закуп и исклучување.

Конечна контролна листа за проверка

  • Нарачката и настанот во сандачето за испраќање се потврдуваат или поништуваат заедно.
  • Преземањата се потврдуваат пред да започне каква било мрежна или деловна работа.
  • Истечените закупи може повторно да се преземат, а ажурирањата ја проверуваат сопственоста.
  • Сигнал примен за време на резервирањето предизвикува итно ослободување без нова работа за испорака.
  • Операциите за поврзување, заклучување, барање и HTTP имаат одделни ограничени истеци на време.
  • Записот во приемното сандаче на Billing и фактурата делат една трансакција.
  • Повторното пуштање на настан враќа успех без дуплирање на неговиот ефект.
  • Акредитивите се со најмали привилегии, сообраќајот е ограничен, а продукцискиот HTTP користи TLS.
  • Контролните табли ја прикажуваат староста на заостатокот, повторните обиди, истекот на закупот и неуспесите на потрошувачот.
  • Буџетите за исклучување при распоредување ја надминуваат ограничената операција на работникот во тек.

Сигурната интеграција не доаѓа од преправањето дека неуспесите можат да се отстранат. Таа доаѓа од одлучувањето каде точно живее трајната вистина, од тоа секоја неизвесна граница да биде повторлива и од обезбедувањето секој повторен обид да е безопасен. Со сандаче за испраќање, пренесувач со закуп и трансакциско приемно сандаче, падовите престануваат да бидат мистериозни гранични случаи и стануваат вообичаени премини на состојби што вашиот систем е дизајниран да ги преживее.

Портрет на автор на блогот

Mihajlo

Јас сум Михајло - развивач поттикнат од љубопитност, дисциплина и постојаната желба да создадам нешто значајно. Споделувам увиди, упатства и бесплатни услуги за да им помогнам на другите да ја поедностават својата работа и да растат во постојано развивачкиот свет на софтверот и вештачката интелигенција.