Shared\Channel
OxPHP\Shared\Channel est un canal borné multi-producteur multi-consommateur qui réside dans le registre partagé et qui est visible par chaque worker PHP du processus. Utilisez-le lorsqu'un gestionnaire de requête et un worker en arrière-plan, ou deux workers, doivent échanger des éléments de travail dans l'ordre FIFO. À l'intérieur d'une Fiber, les appels bloquants se suspendent de manière coopérative afin que le thread worker sous-jacent reste disponible pour d'autres requêtes.
Vue d'ensemble
- Borné. La capacité est fixée à la construction. Une fois plein,
sendbloque (ou suspend une Fiber) ;trySendrenvoie le résultat sans attendre. - MPMC. N'importe quel nombre d'émetteurs et de récepteurs répartis sur les threads. La livraison est FIFO.
- Compatible Fiber. En mode worker avec le pool asynchrone, les variantes bloquantes suspendent la Fiber au lieu du thread de l'OS. En mode traditionnel, elles bloquent le thread de l'OS.
- Typé par résultat. Chaque send/recv renvoie un
Channel\SendResult/Channel\RecvResult— closed / full / timeout sont des issues normales pour les répartiteurs en fan-out, elles apparaissent donc comme des variantes de résultat plutôt que comme des exceptions.
Référence de l'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;
}Limites de taille
Deux plafonds protègent contre les bombes d'allocation issues d'entrées non fiables ; les deux sont configurables via des variables d'environnement et lèvent une exception rattrapable plutôt que d'interrompre le processus :
- Capacité —
new Channel($capacity)préalloue un tableau d'emplacements. Une capacité dont le tableau d'emplacements dépasseraitSHARED_MAX_CHANNEL_BYTES(64 MiB par défaut, ~2M emplacements) lèveCapacityException. Les canaux ont normalement une profondeur allant de quelques dizaines à quelques milliers, ce déclenchement ne survient donc que pour des valeurs pathologiques. (Une$capacitynon positive constitue uneTypeExceptiondistincte.) Le budget effectif est ramené vers le haut à au moins un emplacement, si bien qu'une valeurSHARED_MAX_CHANNEL_BYTESnulle ou trop faible ne peut jamais rejeter un canal minimal de capacité 1 — elle ne plafonne que les grandes capacités. - Par valeur — chaque valeur envoyée est sérialisée ; une valeur dont la
taille sérialisée dépasse
SHARED_MAX_VALUE_SIZE(1 MiB par défaut) lèveValueTooLargeException. Cela s'applique àsend,trySend,sendTimeoutetsendMany(qui rejette le lot entier, sans rien envoyer, si un seul élément dépasse le plafond). C'est le même plafond par valeur que celui appliqué parShared\Map.
Politiques d'attente
Chaque direction possède trois variantes de méthode ; le suffixe encode la politique d'attente :
| Suffixe | Comportement |
|---|---|
try* |
Non bloquant. Renvoie Empty / Full si l'appel ne peut aboutir. |
| (nom nu) | Bloque indéfiniment, ou jusqu'à l'annulation de la Fiber de requête. |
*Timeout |
Attente bornée. $ms doit être > 0 — utilisez l'une des autres formes pour attendre indéfiniment ou ne pas bloquer. |
$ms est toujours un entier positif de millisecondes pour recvTimeout / sendTimeout / sendMany / recvMany. Les valeurs nulles, négatives, non entières et absentes lèvent OxPHP\Shared\TypeException au niveau du pont — cette contrainte a été sortie du corps de la méthode afin que la trichotomie se documente d'elle-même.
Variantes de résultat atteignables par méthode
SendResult::Full et RecvResult::Empty ne sont émis que par les appels non bloquants try* — une variante bloquante acquiert soit l'emplacement/l'élément, soit épuise son budget, auquel cas le résultat est Timeout, et non Full / Empty. La matrice complète d'atteignabilité :
| Méthode | Ok |
Full / Empty |
Timeout |
Closed |
|---|---|---|---|---|
trySend |
✓ | Full ✓ |
— | ✓ |
send |
✓ | — | — | ✓ |
sendTimeout |
✓ | — | ✓ | ✓ |
tryRecv |
✓ | Empty ✓ |
— | ✓ |
recv |
✓ | — | — | ✓ |
recvTimeout |
✓ | — | ✓ | ✓ |
Les vérifications isFull() / isEmpty() sur un résultat issu d'un appel bloquant sont du code mort. Les exemples match ci-dessous laissent ces branches sous forme de commentaires unreachable afin de rendre l'asymétrie évidente pour le lecteur.
Répartition des résultats
RecvResult et SendResult portent un discriminant status() ainsi qu'une charge utile (pour RecvResult uniquement). Deux idiomes équivalents :
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() lève OxPHP\Shared\SharedException s'il est appelé sur une variante non-Ok — utilisez d'abord isOk() / valueOr() / status(). C'est intentionnel : cela fait ressortir sous forme d'échec bruyant le bug consistant à agir sur une charge utile absente, plutôt que sous celle d'un null silencieux.
Comportement Fiber vs bloquant
Les mêmes appels de méthode se comportent différemment selon que PHP s'exécute actuellement ou non à l'intérieur d'une Fiber :
- À l'intérieur d'une Fiber (mode worker +
oxphp_async(...)) : les variantes bloquantes (send,recv,sendTimeout,recvTimeout) allouent une promesse synthétique, enregistrent un waker auprès du canal et suspendent la Fiber. Le thread worker retourne à l'ordonnanceur et traite d'autres Fibers jusqu'à ce que le canal notifie le waker. - En dehors d'une Fiber (mode traditionnel, ou un chemin d'appel non asynchrone) : les variantes bloquantes bloquent le thread worker de l'OS via
crossbeam_channel. Aucun autre travail ne s'exécute sur ce thread jusqu'à ce que l'appel retourne.
Le mode traditionnel bénéficie tout de même de la sémantique du canal ; il la paie simplement d'un thread bloqué. Le mode worker est le déploiement recommandé pour tout pipeline qui repose sur l'attente.
<?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);
});L'annulation d'une Fiber (abandon de requête, échéance atteinte au niveau du SAPI) se manifeste toujours sous la forme d'une OxPHP\Async\AsyncException. Elle n'est pas traduite en RecvResult::Closed. Le canal ne s'est pas fermé ; c'est le runtime qui vous a annulé. Relancez ou rattrapez selon le cas.
Sémantique de fermeture
close() est idempotent ; l'appeler une seconde fois est sans effet. Après la fermeture :
send,sendTimeout,trySendrenvoientSendResult::Closed.sendManyrenvoie le nombre effectivement accepté avant la fermeture (0si déjà fermé).recv,recvTimeout,tryRecvcontinuent de vider les éléments en tampon sous forme deRecvResult::Ok($value), puis renvoientRecvResult::Closedune fois vide.recvManyrenvoie les éléments qu'il peut vider, éventuellement un tableau vide.isClosed()renvoietrue.- Les émetteurs bloqués se réveillent avec
SendResult::Closed; les récepteurs bloqués se réveillent avecRecvResult::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());Le schéma d'un arrêt gracieux de pipeline est le suivant : les producteurs s'arrêtent, le côté producteur appelle close(), les consommateurs vident dans une boucle for (;;) { $r = $ch->recv(); if ($r->isClosed()) break; … } et se terminent naturellement. Aucun consommateur n'observe jamais par erreur une charge utile null ; c'est la variante qui porte ce signal.
Vidage à l'arrêt
Lorsque le processus OxPHP s'arrête, le registre OxPHP\Shared appelle close() sur chaque entrée, y compris les canaux. Du point de vue de PHP, cela est identique à un close() explicite :
- Les appels
recv*bloqués renvoientRecvResult::Closed. - Les appels
send*bloqués renvoientSendResult::Closed.
Un appelant qui suppose que recv()->value() est toujours sûr lèvera SharedException à l'arrêt ou dès qu'un autre détenteur ferme le canal.
Opérations par lots
sendMany et recvMany existent pour les pipelines qui déplacent les éléments par groupes. Préférez-les lorsque vous manipulez régulièrement plus de 10 éléments à la fois : chaque lot représente un seul aller-retour FFI au lieu de N, ce qui réduit sensiblement le surcoût par élément dans les boucles limitées par le débit.
<?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);Points de sémantique à noter :
- Les deux méthodes par lots renvoient des résultats partiels en cas de timeout ou de fermeture en cours de lot — jamais une exception. Vérifiez l'entier / la longueur du tableau pour détecter une progression partielle.
sendManysur un canal déjà fermé renvoie0.recvManysur un canal fermé et vide renvoie[].$msdoit être> 0. Il n'existe pas de surcharge « vide tout ce qui est en tampon » — appeleztryRecv()dans une boucle serrée si ce schéma vous importe ; nous ajouterons une variante par lots dédiée lorsqu'un vrai site d'appel le demandera.
Exceptions
| Exception | Levée par |
|---|---|
CapacityException |
new Channel($capacity) lorsque le tableau d'emplacements dépasserait SHARED_MAX_CHANNEL_BYTES (64 MiB par défaut, ~2M emplacements). |
ValueTooLargeException |
send / trySend / sendTimeout / sendMany lorsque la taille sérialisée d'une valeur dépasse SHARED_MAX_VALUE_SIZE (1 MiB par défaut). |
TypeException |
$capacity non positive ; un $ms nul, négatif, non entier ou absent ; ou l'envoi d'une valeur qui n'est ni un scalaire, ni null, ni un tableau imbriqué d'éléments partageables. |
SharedException |
RecvResult::value() sur une variante non-Ok, ou clone $channel. |
AsyncException |
Un recv / send bloquant interrompu par une annulation de Fiber ; se manifeste sous la forme d'une OxPHP\Async\AsyncException, et non d'une variante de résultat. |
Observabilité
Le serveur interne (INTERNAL_ADDR=127.0.0.1:9090 par défaut) expose les canaux dans les endpoints génériques du registre partagé :
GET /__ox_shared/summaryinclut un compartimentChannelaveccount,bytesetops.GET /__ox_shared/entriesliste les entrées du registre avec leurs identifiants (accepte un paramètre de requêtelimit).GET /__ox_shared/entry?id=<id>renvoie l'état par canal :capacity,count,pending(alias déprécié decount),closed,senders_blocked,receivers_blocked.
Le compteur ops (et la métrique oxphp_shared_operations_total{type="Channel"} à l'échelle du registre) comptabilise chaque tentative de recv quelle qu'en soit l'issue — un succès, un canal vide/fermé et un recvTimeout expiré l'incrémentent tous. La métrique suit les accès au canal, et pas seulement les valeurs transférées avec succès.
Exposition Prometheus sur /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 s'incrémente pour la fin d'un sendMany partiel qui n'a pas pu tenir.
Schémas courants
Producteur HTTP, consommateur asynchrone
<?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 sur plusieurs consommateurs
Lancez N consommateurs asynchrones sur le même canal ; le registre garantit qu'exactement l'un d'eux reçoit chaque élément.
<?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; }
}
});
}Pipeline borné avec contre-pression
trySend associé à un compteur d'abandons permet à un producteur de délester la charge plutôt que de bloquer en cas de surcharge :
<?php
if ($ch->trySend($event)->isFull()) {
increment_dropped_metric();
}Pièges
recvTimeout(0)est une TypeException, pas un appel « non bloquant ». UtiliseztryRecv()pour un appel non bloquant. Les noms de méthode dédiés ont supprimé l'ambiguïté de surcharge du timeout que présentait l'ancien?float $timeout.- Les valeurs doivent être partageables. Les scalaires,
nullet les tableaux imbriqués d'éléments partageables sont autorisés. Passer un objet qui n'est pas une instanceShared\*lèveTypeExceptionà l'envoi. - Le clonage est interdit.
clone $channellèveSharedException; transférez plutôt le canal via unusede closure —oxphp_async(function () use ($ch) { ... })— afin que les deux côtés voient la même entrée de registre. value()sur une variante non-Ok lève une exception. C'est une fonctionnalité, pas un bug — cela transforme l'erreur « j'ai oublié de vérifier isOk » en échec bruyant. UtilisezvalueOr($default)si une valeur par défaut convient réellement.- L'annulation d'une Fiber se propage sous forme d'exception, et non de variante de résultat. Un
recv()interrompu par l'annulation d'une requête lèveOxPHP\Async\AsyncException. Les canaux qui se ferment d'eux-mêmes produisent toujoursRecvResult::Closed. - Attentes annulées avec charges utiles en transit. Si de nombreuses Fibers attendent sur
send/recvet sont annulées alors que leur charge utile était sur le point de passer, celle-ci peut rester référencée jusqu'au prochain réveil. Gardez le nombre d'attentes borné (par exemple, plafonnez la concurrence avec unShared\Counterou un sémaphore basé sur la capacité du canal).
Migration depuis l'ancienne API
| Avant | Maintenant |
|---|---|
$ch->tryRecv() → null si vide, lève ClosedException |
$ch->tryRecv() → RecvResult::Empty / Closed |
$ch->recv($secs) → null en cas de timeout/fermeture |
$ch->recv() (indéfiniment) / $ch->recvTimeout($ms) → RecvResult |
$ch->trySend($v): bool |
$ch->trySend($v): SendResult |
$ch->send($v, $secs): bool (lève TimeoutException/ClosedException) |
$ch->send($v) / $ch->sendTimeout($v, $ms) → SendResult |
$ch->sendMany($vs, $secs) (lève TimeoutException en cas de partiel) |
$ch->sendMany($vs, $ms): int (nombre partiel, sans exception) |
?float $timeout (secondes, NaN/INF tolérés) |
int $ms (millisecondes, doit être > 0) |
Voir aussi
- Mode worker — prérequis pour les variantes bloquantes qui suspendent les Fibers.
- Promesses asynchrones — la closure
oxphp_async()est la manière normale de confier unChannelà une Fiber en arrière-plan. - Multiplexage de Fibers — explique comment la suspension maintient le thread worker productif pendant que les opérations sur le canal attendent.