Shared\Channel

OxPHP\Shared\Channel to ograniczony kanał typu multi-producer multi-consumer, który znajduje się we współdzielonym rejestrze i jest widoczny dla każdego workera PHP w procesie. Używaj go, gdy handler żądania i worker działający w tle — albo dwa workery — muszą przekazywać sobie zadania w kolejności FIFO. Wewnątrz Fibera wywołania blokujące zawieszają się kooperacyjnie, dzięki czemu bazowy wątek workera pozostaje wolny dla innych żądań.

Przegląd

  • Ograniczony. Pojemność jest ustalana przy tworzeniu. Gdy kanał się zapełni, send blokuje (lub zawiesza Fiber); trySend zgłasza wynik bez czekania.
  • MPMC. Dowolna liczba nadawców i odbiorców w wielu wątkach. Dostarczanie odbywa się w kolejności FIFO.
  • Obsługa Fiberów. W trybie worker z pulą asynchroniczną warianty blokujące zawieszają Fiber zamiast wątku systemu operacyjnego. W trybie tradycyjnym blokują wątek systemu operacyjnego.
  • Typowany wynikiem. Każde send/recv zwraca Channel\SendResult / Channel\RecvResult — closed / full / timeout to normalne wyniki dla dyspozytorów rozgałęziających (fan-out), więc pojawiają się jako warianty wyniku, a nie jako wyjątki.

Referencja 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; }

Ograniczenia rozmiaru

Dwa górne limity chronią przed bombami alokacji pochodzącymi z niezaufanych danych wejściowych; oba są konfigurowalne za pomocą zmiennych środowiskowych i zgłaszają przechwytywalny wyjątek, zamiast przerywać proces:

  • Pojemnośćnew Channel($capacity) wstępnie alokuje tablicę slotów. Pojemność, której tablica slotów przekroczyłaby SHARED_MAX_CHANNEL_BYTES (domyślnie 64 MiB, ~2 mln slotów), powoduje wyrzucenie CapacityException. Kanały mają zwykle głębokość od kilkudziesięciu do kilku tysięcy elementów, więc zadziała to tylko przy patologicznych wartościach. (Niedodatnia wartość $capacity to osobny TypeException.) Efektywny budżet jest zaokrąglany w górę do co najmniej jednego slotu, więc zerowa lub zbyt mała wartość SHARED_MAX_CHANNEL_BYTES nigdy nie może odrzucić minimalnego kanału o pojemności 1 — ogranicza wyłącznie duże pojemności.
  • Na wartość — każda wysłana wartość jest serializowana; wartość, której rozmiar po serializacji przekracza SHARED_MAX_VALUE_SIZE (domyślnie 1 MiB), powoduje wyrzucenie ValueTooLargeException. Dotyczy to send, trySend, sendTimeout oraz sendMany (które odrzuca całą paczkę, nie wysyłając niczego, jeśli którykolwiek element przekracza limit). To ten sam limit na wartość, który egzekwuje Shared\Map.

Polityki oczekiwania

Każdy kierunek ma trzy warianty metody; sufiks koduje politykę oczekiwania:

Sufiks Zachowanie
try* Nieblokujący. Zgłasza Empty / Full, jeśli wywołanie nie może być kontynuowane.
(goła nazwa) Blokuje w nieskończoność lub do momentu anulowania Fibera żądania.
*Timeout Ograniczone oczekiwanie. $ms musi być > 0 — dla oczekiwania w nieskończoność lub nieblokującego użyj jednej z pozostałych form.

$ms jest zawsze dodatnią liczbą całkowitą milisekund dla recvTimeout / sendTimeout / sendMany / recvMany. Wartości zerowe, ujemne, niecałkowite oraz brak wartości powodują wyrzucenie OxPHP\Shared\TypeException na poziomie mostka — to ograniczenie zostało wyniesione poza ciało metody, dzięki czemu trójpodział jest samodokumentujący.

Osiągalne warianty wyniku dla poszczególnych metod

SendResult::Full i RecvResult::Empty są zwracane wyłącznie przez nieblokujące wywołania try* — wariant blokujący albo pozyskuje slot/element, albo wyczerpuje budżet, w którym to przypadku wynikiem jest Timeout, a nie Full / Empty. Pełna macierz osiągalności:

Metoda Ok Full / Empty Timeout Closed
trySend Full
send
sendTimeout
tryRecv Empty
recv
recvTimeout

Sprawdzenia isFull() / isEmpty() na wyniku z wywołania blokującego to martwy kod. Przykłady z match poniżej pozostawiają te gałęzie jako komentarze unreachable, aby uczynić tę asymetrię oczywistą dla czytelników.

Rozdzielanie wyniku

RecvResult i SendResult niosą dyskryminator status() oraz (tylko w przypadku RecvResult) ładunek. Dwa równoważne idiomy:

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() wyrzuca OxPHP\Shared\SharedException, jeśli zostanie wywołane na wariancie innym niż Ok — najpierw użyj isOk() / valueOr() / status(). Jest to zamierzone: ujawnia błąd polegający na działaniu na brakującym ładunku jako głośną awarię, a nie jako ciche null.

Zachowanie w Fiberze a blokowanie

Te same wywołania metod zachowują się różnie w zależności od tego, czy PHP działa aktualnie wewnątrz Fibera:

  • Wewnątrz Fibera (tryb worker + oxphp_async(...)): warianty blokujące (send, recv, sendTimeout, recvTimeout) alokują syntetyczny promise, rejestrują waker w kanale i zawieszają Fiber. Wątek workera wraca do schedulera i przetwarza inne Fibery, dopóki kanał nie powiadomi wakera.
  • Poza Fiberem (tryb tradycyjny lub nieasynchroniczna ścieżka wywołania): warianty blokujące blokują systemowy wątek workera za pośrednictwem crossbeam_channel. Żadna inna praca nie działa na tym wątku, dopóki wywołanie nie zwróci sterowania.

Tryb tradycyjny nadal otrzymuje semantykę kanału; płaci za to po prostu zablokowanym wątkiem. Tryb worker to zalecane wdrożenie dla każdego potoku, który polega na oczekiwaniu.

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); });
Anulowanie Fibera

Anulowanie Fibera (przerwanie żądania, upłynięcie deadline'u na poziomie SAPI) nadal ujawnia się jako OxPHP\Async\AsyncException. Nie jest tłumaczone na RecvResult::Closed. Kanał się nie zamknął; to środowisko uruchomieniowe anulowało wykonanie. W zależności od sytuacji wyrzuć ponownie lub przechwyć.

Semantyka zamykania

close() jest idempotentne; wywołanie go po raz drugi nie ma żadnego efektu. Po zamknięciu:

  • send, sendTimeout, trySend zwracają SendResult::Closed.
  • sendMany zwraca liczbę rzeczywiście przyjętych elementów przed zamknięciem (0, jeśli kanał był już zamknięty).
  • recv, recvTimeout, tryRecv nadal wyczytują zbuforowane elementy jako RecvResult::Ok($value), a po opróżnieniu zwracają RecvResult::Closed.
  • recvMany zwraca elementy, które może wyczytać, możliwe że pustą tablicę.
  • isClosed() zwraca true.
  • Zablokowani nadawcy budzą się z SendResult::Closed; zablokowani odbiorcy budzą się z 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());

Wzorzec łagodnego zamknięcia potoku jest następujący: producenci przestają działać, strona producenta wywołuje close(), a konsumenci wyczytują dane w pętli for (;;) { $r = $ch->recv(); if ($r->isClosed()) break; … } i kończą działanie w naturalny sposób. Dzięki temu żaden konsument nigdy przez pomyłkę nie zaobserwuje ładunku null; ten sygnał niesie sam wariant.

Wyczytywanie przy zamykaniu procesu

Gdy proces OxPHP jest zamykany, rejestr OxPHP\Shared wywołuje close() na każdym wpisie, w tym na kanałach. Z perspektywy PHP wygląda to identycznie jak jawne close():

  • Zablokowane wywołania recv* zwracają RecvResult::Closed.
  • Zablokowane wywołania send* zwracają SendResult::Closed.
Zawsze sprawdzaj wariant wyniku

Kod wywołujący, który zakłada, że recv()->value() jest zawsze bezpieczne, wyrzuci SharedException przy zamknięciu lub za każdym razem, gdy inny posiadacz zamknie kanał.

Operacje pakietowe

sendMany i recvMany istnieją z myślą o potokach, które przenoszą elementy w grupach. Preferuj je, gdy rutynowo obsługujesz 10+ elementów naraz: każda paczka to jedna wymiana FFI (round trip) zamiast N, co znacząco zmniejsza narzut na element w pętlach ograniczonych przepustowością.

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);

Warto zwrócić uwagę na następującą semantykę:

  • Obie metody pakietowe zwracają wyniki częściowe przy przekroczeniu limitu czasu lub zamknięciu w trakcie paczki — nigdy wyjątek. Sprawdź liczbę całkowitą / długość tablicy, aby wykryć częściowy postęp.
  • sendMany na już zamkniętym kanale zwraca 0.
  • recvMany na zamkniętym i pustym zwraca [].
  • $ms musi być > 0. Nie ma przeciążenia typu „wyczytaj cokolwiek jest w buforze” — jeśli ten wzorzec ma znaczenie, wywołuj tryRecv() w ciasnej pętli; dedykowany wariant pakietowy dodamy, gdy poprosi o niego realne miejsce wywołania.

Wyjątki

Wyjątek Zgłaszany przez
CapacityException new Channel($capacity), gdy tablica slotów przekroczyłaby SHARED_MAX_CHANNEL_BYTES (domyślnie 64 MiB, ~2 mln slotów).
ValueTooLargeException send / trySend / sendTimeout / sendMany, gdy rozmiar wartości po serializacji przekracza SHARED_MAX_VALUE_SIZE (domyślnie 1 MiB).
TypeException Niedodatnia wartość $capacity; wartość $ms, która jest zerowa, ujemna, niecałkowita lub nieobecna; albo wysłanie wartości, która nie jest skalarem, null ani zagnieżdżoną tablicą wartości współdzielnych (shareable).
SharedException RecvResult::value() na wariancie innym niż Ok albo clone $channel.
AsyncException Blokujące recv / send przerwane przez anulowanie Fibera; ujawnia się jako OxPHP\Async\AsyncException, a nie jako wariant Result.

Obserwowalność

Serwer wewnętrzny (domyślnie INTERNAL_ADDR=127.0.0.1:9090) udostępnia kanały w ogólnych endpointach współdzielonego rejestru:

  • GET /__ox_shared/summary zawiera koszyk Channel z polami count, bytes i ops.
  • GET /__ox_shared/entries wypisuje wpisy rejestru wraz z ich identyfikatorami (przyjmuje parametr zapytania limit).
  • GET /__ox_shared/entry?id=<id> zwraca stan pojedynczego kanału: capacity, count, pending (przestarzały alias dla count), closed, senders_blocked, receivers_blocked.

Licznik ops (oraz obejmujący cały rejestr oxphp_shared_operations_total{type="Channel"}) zlicza każdą próbę recv niezależnie od jej wyniku — trafienie, pusty/zamknięty kanał oraz recvTimeout zakończony przekroczeniem limitu czasu, każde z nich go zwiększa. Metryka śledzi dostępy do kanału, a nie tylko pomyślnie przekazane wartości.

Ekspozycja Prometheus na /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 jest zwiększany o końcówkę częściowego sendMany, która się nie zmieściła.

Typowe wzorce

Producent HTTP, asynchroniczny konsument

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"; });

Rozgałęzianie (fan-out) na wielu konsumentów

Uruchom N asynchronicznych konsumentów na tym samym kanale; rejestr zapewnia, że dokładnie jeden z nich otrzyma każdy element.

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; } } }); }

Ograniczony potok z przeciwciśnieniem

trySend wraz z licznikiem odrzuceń pozwala producentowi odrzucać obciążenie zamiast blokować się przy przeciążeniu:

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

Pułapki

  • recvTimeout(0) to TypeException, a nie „tryb nieblokujący”. Dla wywołań nieblokujących używaj tryRecv(). Dedykowane nazwy metod usunęły niejednoznaczność przeciążenia limitu czasu, którą miał dawny ?float $timeout.
  • Wartości muszą być współdzielne (shareable). Dozwolone są skalary, null oraz zagnieżdżone tablice wartości współdzielnych. Przekazanie obiektu, który nie jest instancją Shared\*, powoduje wyrzucenie TypeException przy wysyłce.
  • Klonowanie jest zabronione. clone $channel wyrzuca SharedException; zamiast tego przekaż kanał przez use domknięcia — oxphp_async(function () use ($ch) { ... }) — tak aby obie strony widziały ten sam wpis rejestru.
  • value() na wariancie innym niż Ok wyrzuca wyjątek. To cecha, a nie błąd — zamienia pomyłkę „zapomniałem sprawdzić isOk” w głośną awarię. Użyj valueOr($default), jeśli wartość domyślna naprawdę jest w porządku.
  • Anulowanie Fibera propaguje się jako wyjątek, a nie jako wariant Result. recv() przerwane przez anulowanie żądania wyrzuca OxPHP\Async\AsyncException. Kanały, które zamykają się same, nadal zwracają RecvResult::Closed.
  • Anulowane oczekujące z ładunkami w locie. Jeśli wiele Fiberów oczekuje na send / recv i zostaje anulowanych w chwili, gdy ich ładunek miał właśnie przejść, ładunek może pozostać referencjonowany aż do następnego wybudzenia. Utrzymuj liczbę oczekujących w ograniczonych granicach (np. ogranicz współbieżność za pomocą Shared\Counter lub semafora o pojemności kanału).

Migracja z poprzedniego API

Było Jest
$ch->tryRecv()null przy pustym, wyrzuca ClosedException $ch->tryRecv()RecvResult::Empty / Closed
$ch->recv($secs)null przy przekroczeniu limitu czasu/zamknięciu $ch->recv() (w nieskończoność) / $ch->recvTimeout($ms)RecvResult
$ch->trySend($v): bool $ch->trySend($v): SendResult
$ch->send($v, $secs): bool (wyrzuca TimeoutException/ClosedException) $ch->send($v) / $ch->sendTimeout($v, $ms)SendResult
$ch->sendMany($vs, $secs) (wyrzuca TimeoutException przy częściowym) $ch->sendMany($vs, $ms): int (liczba częściowa, bez wyrzucania)
?float $timeout (sekundy, tolerowane NaN/INF) int $ms (milisekundy, musi być > 0)

Powiązane

  • Tryb worker — warunek wstępny dla wariantów blokujących zawieszających Fiber.
  • Asynchroniczne promisy — domknięcie oxphp_async() to standardowy sposób przekazania Channel do Fibera działającego w tle.
  • Multipleksowanie Fiberów — wyjaśnia, jak zawieszanie utrzymuje produktywność wątku workera, gdy operacje na kanale czekają.