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,
sendblokuje (lub zawiesza Fiber);trySendzgł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
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łabySHARED_MAX_CHANNEL_BYTES(domyślnie 64 MiB, ~2 mln slotów), powoduje wyrzucenieCapacityException. 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ść$capacityto osobnyTypeException.) 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_BYTESnigdy 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 wyrzucenieValueTooLargeException. Dotyczy tosend,trySend,sendTimeoutorazsendMany(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 egzekwujeShared\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:
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
// 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 (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,trySendzwracająSendResult::Closed.sendManyzwraca liczbę rzeczywiście przyjętych elementów przed zamknięciem (0, jeśli kanał był już zamknięty).recv,recvTimeout,tryRecvnadal wyczytują zbuforowane elementy jakoRecvResult::Ok($value), a po opróżnieniu zwracająRecvResult::Closed.recvManyzwraca elementy, które może wyczytać, możliwe że pustą tablicę.isClosed()zwracatrue.- Zablokowani nadawcy budzą się z
SendResult::Closed; zablokowani odbiorcy budzą się zRecvResult::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());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.
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
$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.
sendManyna już zamkniętym kanale zwraca0.recvManyna zamkniętym i pustym zwraca[].$msmusi być> 0. Nie ma przeciążenia typu „wyczytaj cokolwiek jest w buforze” — jeśli ten wzorzec ma znaczenie, wywołujtryRecv()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/summaryzawiera koszykChannelz polamicount,bytesiops.GET /__ox_shared/entrieswypisuje wpisy rejestru wraz z ich identyfikatorami (przyjmuje parametr zapytanialimit).GET /__ox_shared/entry?id=<id>zwraca stan pojedynczego kanału:capacity,count,pending(przestarzały alias dlacount),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:
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 jest zwiększany o końcówkę częściowego sendMany, która się nie zmieściła.
Typowe wzorce
Producent HTTP, asynchroniczny konsument
<?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
$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
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żywajtryRecv(). 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,
nulloraz zagnieżdżone tablice wartości współdzielnych. Przekazanie obiektu, który nie jest instancjąShared\*, powoduje wyrzucenieTypeExceptionprzy wysyłce. - Klonowanie jest zabronione.
clone $channelwyrzucaSharedException; zamiast tego przekaż kanał przezusedomknię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żyjvalueOr($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 wyrzucaOxPHP\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/recvi 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\Counterlub 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 przekazaniaChanneldo Fibera działającego w tle. - Multipleksowanie Fiberów — wyjaśnia, jak zawieszanie utrzymuje produktywność wątku workera, gdy operacje na kanale czekają.