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
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) полезную нагрузку. Две эквивалентные идиомы:
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
// 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
$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
$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:
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>"} counteritems_dropped_total увеличивается на хвост частичного sendMany, который не поместился.
Распространённые паттерны
HTTP-производитель, асинхронный потребитель
<?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
$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
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фоновому файберу. - Мультиплексирование файберов — объясняет, как приостановка поддерживает продуктивность потока воркера, пока операции с каналом ожидают.