Shared\Channel

OxPHP\Shared\Channel — это ограниченный по ёмкости канал с несколькими производителями и несколькими потребителями (multi-producer multi-consumer), который живёт в разделяемом реестре и виден каждому PHP-воркеру в процессе. Используйте его, когда обработчику запроса и фоновому воркеру или двум воркерам нужно обмениваться рабочими элементами в порядке FIFO. Внутри файбера блокирующие вызовы приостанавливаются кооперативно, поэтому нижележащий поток воркера остаётся свободным для других запросов.

Обзор

  • Ограниченный. Ёмкость фиксируется при создании. Когда канал заполнен, send блокируется (или приостанавливает файбер); trySend сообщает результат без ожидания.
  • MPMC. Любое количество отправителей и получателей в разных потоках. Доставка происходит в порядке FIFO.
  • С поддержкой файберов. В режиме воркеров с асинхронным пулом блокирующие варианты приостанавливают файбер, а не поток ОС. В традиционном режиме они блокируют поток ОС.
  • Типизированные результаты. Каждый send/recv возвращает Channel\SendResult / Channel\RecvResult — closed / full / timeout являются нормальными исходами для диспетчеров разветвления, поэтому они возникают как варианты результата, а не как исключения.

Справочник API

php
namespace OxPHP\Shared; final class Channel implements Shareable, \Countable { public function __construct(int $capacity); // ── Receive ── public function tryRecv(): Channel\RecvResult; // non-blocking public function recv(): Channel\RecvResult; // block forever / fiber-cancel public function recvTimeout(int $ms): Channel\RecvResult; // bounded ($ms > 0) // ── Send ── public function trySend(mixed $value): Channel\SendResult; public function send(mixed $value): Channel\SendResult; public function sendTimeout(mixed $value, int $ms): Channel\SendResult; // ── Batch ── public function sendMany(array $values, int $ms): int; // partial count, no throw public function recvMany(int $max, int $ms): array; // partial array, no throw // ── Lifecycle ── public function close(): bool; public function isClosed(): bool; public function count(): int; public function id(): int; } namespace OxPHP\Shared\Channel; enum RecvStatus { case Ok; case Empty; case Timeout; case Closed; } enum SendStatus { case Ok; case Full; case Timeout; case Closed; } final class RecvResult { public function isOk(): bool; public function isEmpty(): bool; public function isTimeout(): bool; public function isClosed(): bool; public function value(): mixed; // throws SharedException if not Ok public function valueOr(mixed $default): mixed; public function status(): RecvStatus; } final class SendResult { public function isOk(): bool; public function isFull(): bool; public function isTimeout(): bool; public function isClosed(): bool; public function status(): SendStatus; }

Ограничения размера

Два потолка защищают от «бомб выделения памяти» из недоверенного ввода; оба настраиваются через переменные окружения и вызывают перехватываемое исключение, а не прерывают процесс:

  • Ёмкостьnew Channel($capacity) заранее выделяет массив слотов. Ёмкость, при которой массив слотов превысил бы SHARED_MAX_CHANNEL_BYTES (по умолчанию 64 MiB, ~2M слотов), приводит к выбросу CapacityException. Каналы обычно имеют глубину от десятков до нескольких тысяч, поэтому это срабатывает только на патологических значениях. (Неположительное $capacity — это отдельный TypeException.) Эффективный бюджет зажимается снизу минимум до одного слота, поэтому нулевое или слишком малое SHARED_MAX_CHANNEL_BYTES никогда не сможет отклонить минимальный канал ёмкостью 1 — оно лишь ограничивает большие значения ёмкости.
  • На значение — каждое отправляемое значение сериализуется; значение, чей сериализованный размер превышает SHARED_MAX_VALUE_SIZE (по умолчанию 1 MiB), приводит к выбросу ValueTooLargeException. Это касается send, trySend, sendTimeout и sendMany (которое отклоняет весь пакет, не отправив ничего, если хотя бы один элемент превышает лимит). Это тот же лимит на значение, который применяется в Shared\Map.

Политики ожидания

У каждого направления есть три варианта метода; суффикс кодирует политику ожидания:

Суффикс Поведение
try* Неблокирующий. Сообщает Empty / Full, если вызов не может продолжиться.
(без суффикса) Блокируется навсегда или до отмены файбера запроса.
*Timeout Ограниченное ожидание. $ms должно быть > 0 — для бесконечного или неблокирующего ожидания используйте одну из других форм.

$ms — это всегда положительное целое число миллисекунд для recvTimeout / sendTimeout / sendMany / recvMany. Нулевые, отрицательные, нецелые и отсутствующие значения вызывают OxPHP\Shared\TypeException на уровне моста — это ограничение вынесено из тела метода, чтобы трихотомия была самодокументируемой.

Достижимые варианты результата по методам

SendResult::Full и RecvResult::Empty возвращаются только неблокирующими вызовами try* — блокирующий вариант либо получает слот/элемент, либо исчерпывает бюджет, и в этом случае результатом будет Timeout, а не Full / Empty. Полная матрица достижимости:

Метод Ok Full / Empty Timeout Closed
trySend Full
send
sendTimeout
tryRecv Empty
recv
recvTimeout

Проверки isFull() / isEmpty() для результата блокирующего вызова — это мёртвый код. В примерах с match ниже эти ветви оставлены в виде комментариев unreachable, чтобы асимметрия была очевидна читателю.

Диспетчеризация результата

RecvResult и SendResult несут дискриминант status() плюс (только для RecvResult) полезную нагрузку. Две эквивалентные идиомы:

php
use OxPHP\Shared\Channel; use OxPHP\Shared\Channel\RecvStatus; $ch = new Channel(capacity: 64); // Boolean accessors — terse when only one outcome matters. $result = $ch->tryRecv(); if ($result->isOk()) { process($result->value()); } elseif ($result->isEmpty()) { backoff(); } else { break; // closed } // Exhaustive match — exhaustiveness checking catches new variants. $r = $ch->recvTimeout(1500); match ($r->status()) { RecvStatus::Ok => process($r->value()), RecvStatus::Timeout => $logger->debug('idle'), RecvStatus::Closed => break, RecvStatus::Empty => /* unreachable: only tryRecv returns Empty */ , }; // Symmetric send-side dispatch — only Ok / Timeout / Closed are reachable. use OxPHP\Shared\Channel\SendStatus; $s = $ch->sendTimeout($value, 1500); match ($s->status()) { SendStatus::Ok => /* delivered */ , SendStatus::Timeout => $logger->debug('backpressure'), SendStatus::Closed => break, SendStatus::Full => /* unreachable: only trySend returns Full */ , }; // Safe accessor — never throws. $value = $ch->tryRecv()->valueOr('fallback');

RecvResult::value() выбрасывает OxPHP\Shared\SharedException, если вызван на варианте, отличном от Ok — сначала используйте isOk() / valueOr() / status(). Это сделано намеренно: так ошибка обращения к отсутствующей полезной нагрузке проявляется как громкий сбой, а не как тихий null.

Поведение в файбере и при блокировке

Одни и те же вызовы методов ведут себя по-разному в зависимости от того, выполняется ли PHP в данный момент внутри файбера:

  • Внутри файбера (режим воркеров + oxphp_async(...)): блокирующие варианты (send, recv, sendTimeout, recvTimeout) выделяют синтетический промис, регистрируют waker в канале и приостанавливают файбер. Поток воркера возвращается к планировщику и обрабатывает другие файберы, пока канал не уведомит waker.
  • Вне файбера (традиционный режим или неасинхронный путь вызова): блокирующие варианты блокируют поток воркера ОС через crossbeam_channel. Никакая другая работа на этом потоке не выполняется, пока вызов не вернётся.

Традиционный режим по-прежнему получает семантику канала; он просто расплачивается заблокированным потоком. Режим воркеров — рекомендуемый вариант развёртывания для любого конвейера, который полагается на ожидание.

php
<?php // Traditional mode: this recvTimeout blocks the worker thread for up to 2 seconds. $ch = new OxPHP\Shared\Channel(16); $r = $ch->recvTimeout(2000); // Worker mode: wrap in oxphp_async and recv suspends cooperatively. oxphp_worker(function () use ($ch) { $consumer = oxphp_async(function () use ($ch) { for (;;) { $r = $ch->recvTimeout(5000); if ($r->isOk()) { process($r->value()); continue; } if ($r->isClosed()) { break; } // Timeout — keep waiting. } }); oxphp_async_await($consumer); });
Отмена файбера

Отмена файбера (прерывание запроса, истечение дедлайна на уровне SAPI) по-прежнему проявляется как OxPHP\Async\AsyncException. Она не преобразуется в RecvResult::Closed. Канал не закрывался; вас отменила среда выполнения. Пробросьте исключение повторно или перехватите его по обстоятельствам.

Семантика закрытия

close() идемпотентен; повторный вызов ничего не делает. После закрытия:

  • send, sendTimeout, trySend возвращают SendResult::Closed.
  • sendMany возвращает количество, фактически принятое до закрытия (0, если канал уже закрыт).
  • recv, recvTimeout, tryRecv продолжают вычитывать буферизованные элементы как RecvResult::Ok($value), а затем возвращают RecvResult::Closed, когда канал опустеет.
  • recvMany возвращает элементы, которые может вычитать, возможно, пустой массив.
  • isClosed() возвращает true.
  • Заблокированные отправители пробуждаются с SendResult::Closed; заблокированные получатели пробуждаются с RecvResult::Closed.
php
<?php $ch = new OxPHP\Shared\Channel(4); $ch->send('one'); $ch->send('two'); $ch->close(); // Drain leftovers. for (;;) { $r = $ch->recv(); if ($r->isClosed()) { break; } echo $r->value(), "\n"; // one, two } // Further sends report Closed without throwing. $result = $ch->send('three'); assert($result->isClosed());

Паттерн для корректного завершения работы конвейера такой: производители останавливаются, сторона производителя вызывает close(), потребители вычитывают остаток в цикле for (;;) { $r = $ch->recv(); if ($r->isClosed()) break; … } и завершаются естественным образом. Ни один потребитель при этом по ошибке не увидит null в качестве полезной нагрузки; этот сигнал несёт сам вариант результата.

Вычитывание при завершении работы

Когда процесс OxPHP завершает работу, реестр OxPHP\Shared вызывает close() для каждой записи, включая каналы. С точки зрения PHP это выглядит идентично явному close():

  • Заблокированные вызовы recv* возвращают RecvResult::Closed.
  • Заблокированные вызовы send* возвращают SendResult::Closed.
Всегда проверяйте вариант результата

Вызывающий код, который считает, что recv()->value() всегда безопасен, выбросит SharedException при завершении работы или всякий раз, когда другой держатель закроет канал.

Пакетные операции

sendMany и recvMany существуют для конвейеров, которые перемещают элементы группами. Предпочитайте их, когда регулярно обрабатываете 10+ элементов за раз: каждый пакет — это один обход через FFI вместо N, что заметно снижает накладные расходы на элемент в циклах, ограниченных пропускной способностью.

php
<?php $ch = new OxPHP\Shared\Channel(1024); // Send an array; returns the count actually accepted. $sent = $ch->sendMany([1, 2, 3, 4, 5], 100); // 5 within 100ms budget // Drain up to 10 items with a 100ms deadline. $batch = $ch->recvMany(10, 100);

Стоит отметить следующую семантику:

  • Оба пакетных метода возвращают частичные результаты при таймауте или закрытии в середине пакета — но никогда не исключение. Проверяйте целое число / длину массива, чтобы обнаружить частичный прогресс.
  • sendMany на уже закрытом канале возвращает 0.
  • recvMany на закрытом и пустом канале возвращает [].
  • $ms должно быть > 0. Перегрузки «вычитать всё, что буферизовано» нет — вызывайте tryRecv() в плотном цикле, если этот паттерн важен; мы добавим отдельный пакетный вариант, когда об этом попросит реальное место вызова.

Исключения

Исключение Кем вызывается
CapacityException new Channel($capacity), когда массив слотов превысил бы SHARED_MAX_CHANNEL_BYTES (по умолчанию 64 MiB, ~2M слотов).
ValueTooLargeException send / trySend / sendTimeout / sendMany, когда сериализованный размер значения превышает SHARED_MAX_VALUE_SIZE (по умолчанию 1 MiB).
TypeException Неположительное $capacity; $ms, равное нулю, отрицательное, нецелое или отсутствующее; или отправка значения, которое не является скаляром, null или вложенным массивом разделяемых значений.
SharedException RecvResult::value() на варианте, отличном от Ok, или clone $channel.
AsyncException Блокирующий recv / send, прерванный отменой файбера; проявляется как OxPHP\Async\AsyncException, а не как вариант Result.

Наблюдаемость

Внутренний сервер (по умолчанию INTERNAL_ADDR=127.0.0.1:9090) предоставляет каналы через общие эндпоинты разделяемого реестра:

  • GET /__ox_shared/summary включает корзину Channel с count, bytes и ops.
  • GET /__ox_shared/entries перечисляет записи реестра с их идентификаторами (принимает query-параметр limit).
  • GET /__ox_shared/entry?id=<id> возвращает состояние отдельного канала: capacity, count, pending (устаревший псевдоним для count), closed, senders_blocked, receivers_blocked.

Счётчик ops (и общерегистровый oxphp_shared_operations_total{type="Channel"}) считает каждую попытку recv независимо от результата — попадание, пустой/закрытый канал и recvTimeout, завершившийся по таймауту, все увеличивают его. Метрика отслеживает обращения к каналу, а не только успешно переданные значения.

Экспозиция Prometheus на /metrics:

text
oxphp_shared_channel_count{channel_id="<id>"} gauge oxphp_shared_channel_pending{channel_id="<id>"} gauge (deprecated, alias of _count) oxphp_shared_channel_senders_blocked{channel_id="<id>"} gauge oxphp_shared_channel_receivers_blocked{channel_id="<id>"} gauge oxphp_shared_channel_items_sent_total{channel_id="<id>"} counter oxphp_shared_channel_items_dropped_total{channel_id="<id>"} counter

items_dropped_total увеличивается на хвост частичного sendMany, который не поместился.

Распространённые паттерны

HTTP-производитель, асинхронный потребитель

php
<?php $work = new OxPHP\Shared\Channel(256); $consumer = oxphp_async(function () use ($work) { for (;;) { $r = $work->recvTimeout(30000); if ($r->isOk()) { process_job($r->value()); continue; } if ($r->isClosed()) { break; } // Timeout — health-check, then keep waiting. } }); oxphp_worker(function () use ($work) { $work->send(['url' => $_POST['url'], 'tries' => 3]); echo "queued"; });

Разветвление (fan-out) на несколько потребителей

Породите N асинхронных потребителей на одном канале; реестр гарантирует, что каждый элемент получит ровно один из них.

php
<?php $ch = new OxPHP\Shared\Channel(1024); for ($i = 0; $i < 4; $i++) { oxphp_async(function () use ($ch, $i) { for (;;) { $r = $ch->recvTimeout(60000); if ($r->isOk()) { handle($i, $r->value()); continue; } if ($r->isClosed()) { break; } } }); }

Ограниченный конвейер с противодавлением (backpressure)

trySend плюс счётчик отброшенных элементов позволяет производителю сбрасывать нагрузку, а не блокироваться при перегрузке:

php
<?php if ($ch->trySend($event)->isFull()) { increment_dropped_metric(); }

Подводные камни

  • recvTimeout(0) — это TypeException, а не «неблокирующий вызов». Для неблокирующего вызова используйте tryRecv(). Выделенные имена методов устранили неоднозначность перегрузки по таймауту, которая была у старого ?float $timeout.
  • Значения должны быть разделяемыми. Разрешены скаляры, null и вложенные массивы разделяемых значений. Передача объекта, не являющегося экземпляром Shared\*, вызывает TypeException при отправке.
  • Клонирование запрещено. clone $channel выбрасывает SharedException; вместо этого передавайте канал через use замыкания — oxphp_async(function () use ($ch) { ... }) — чтобы обе стороны видели одну и ту же запись реестра.
  • value() на варианте, отличном от Ok, выбрасывает исключение. Это фича, а не баг — она превращает ошибку «я забыл проверить isOk» в громкий сбой. Используйте valueOr($default), если значение по умолчанию действительно приемлемо.
  • Отмена файбера распространяется как исключение, а не как вариант Result. recv(), прерванный отменой запроса, вызывает OxPHP\Async\AsyncException. Каналы, которые закрываются сами, по-прежнему выдают RecvResult::Closed.
  • Отменённые ожидающие с полезной нагрузкой «в полёте». Если множество файберов ожидает на send / recv и отменяется в тот момент, когда их полезная нагрузка вот-вот должна была пересечь границу, на неё может сохраняться ссылка до следующего пробуждения. Держите число ожидающих ограниченным (например, ограничьте конкурентность с помощью Shared\Counter или семафора на ёмкость канала).

Миграция с предыдущего API

Было Стало
$ch->tryRecv()null при пустом канале, выбрасывает ClosedException $ch->tryRecv()RecvResult::Empty / Closed
$ch->recv($secs)null при таймауте/закрытии $ch->recv() (навсегда) / $ch->recvTimeout($ms)RecvResult
$ch->trySend($v): bool $ch->trySend($v): SendResult
$ch->send($v, $secs): bool (выбрасывает TimeoutException/ClosedException) $ch->send($v) / $ch->sendTimeout($v, $ms)SendResult
$ch->sendMany($vs, $secs) (выбрасывает TimeoutException при частичной отправке) $ch->sendMany($vs, $ms): int (частичное количество, без выброса)
?float $timeout (секунды, NaN/INF допускались) int $ms (миллисекунды, должно быть > 0)

Связанные материалы

  • Режим воркеров — необходимое условие для блокирующих вариантов, приостанавливающих файбер.
  • Асинхронные промисы — замыкание oxphp_async() — это стандартный способ передать Channel фоновому файберу.
  • Мультиплексирование файберов — объясняет, как приостановка поддерживает продуктивность потока воркера, пока операции с каналом ожидают.