Shared\Channel

OxPHP\Shared\Channel は、共有レジストリ内に存在し、プロセス内のすべての PHP ワーカーから見える、有界のマルチプロデューサー・マルチコンシューマーチャネルです。リクエストハンドラーとバックグラウンドワーカー、あるいは 2 つのワーカーが、FIFO 順で作業アイテムをやり取りする必要がある場合に使用します。ファイバー内では、ブロッキング呼び出しが協調的にサスペンドするため、基盤となるワーカースレッドは他のリクエストのために解放されたままになります。

概要

  • 有界。 容量は構築時に固定されます。いっぱいになると send はブロックします(またはファイバーをサスペンドします)。trySend は待機せずに結果を報告します。
  • MPMC。 スレッドをまたいで任意の数の送信側と受信側を持てます。配信は FIFO です。
  • ファイバー対応。 非同期プールを備えたワーカーモードでは、ブロッキング系のメソッドは OS スレッドではなくファイバーをサスペンドします。従来モードでは OS スレッドをブロックします。
  • 結果型付き。 すべての send/recv は Channel\SendResult / Channel\RecvResult を返します。クローズ / フル / タイムアウトはファンアウトのディスパッチャーにとって通常の結果なので、例外ではなく結果のバリアントとして現れます。

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

サイズ制限

信頼できない入力によるアロケーション爆弾から保護するために、2 つの上限が設けられています。どちらも環境変数で設定でき、プロセスを中断させるのではなく、キャッチ可能な例外を発生させます。

  • 容量new Channel($capacity) はスロット配列を事前に確保します。スロット 配列が SHARED_MAX_CHANNEL_BYTES(デフォルト 64 MiB、約 200 万スロット)を 超えるような容量を指定すると CapacityException がスローされます。チャネルの 深さは通常、数十から数千程度なので、これは病的な値のときにのみ発生します。 ($capacity が正でない場合は別の TypeException になります。)実効的な予算は 最低でも 1 スロットまで引き上げてクランプされるため、ゼロや小さすぎる SHARED_MAX_CHANNEL_BYTES であっても、最小の容量 1 のチャネルを拒否することは 決してありません。あくまで大きな容量に上限を設けるだけです。
  • 値ごと — 送信される各値はシリアライズされます。シリアライズ後のサイズが SHARED_MAX_VALUE_SIZE(デフォルト 1 MiB)を超える値は ValueTooLargeException をスローします。これは sendtrySendsendTimeout、および sendMany に適用されます(sendMany では、いずれかの 要素が上限を超えている場合、バッチ全体を拒否し、何も送信しません)。これは Shared\Map が適用するのと同じ、値ごとの上限です。

待機ポリシー

どちらの方向にも 3 つのメソッドのバリアントがあり、サフィックスが待機ポリシーを表します。

サフィックス 動作
try* 非ブロッキング。呼び出しを続行できない場合は Empty / Full を報告します。
(サフィックスなし) 永久に、またはリクエストファイバーがキャンセルされるまでブロックします。
*Timeout 制限付き待機。$ms> 0 でなければなりません。永久待機や非ブロッキングには他の形式のいずれかを使用してください。

recvTimeout / sendTimeout / sendMany / recvMany において、$ms常にミリ秒単位の正の整数です。ゼロ、負の値、整数以外、および値の欠落は、ブリッジで OxPHP\Shared\TypeException を発生させます。この制約はメソッド本体から外に移されたため、3 つの形式の区別が自己文書化されるようになりました。

メソッドごとに到達可能な結果バリアント

SendResult::FullRecvResult::Empty は、非ブロッキングの try* 呼び出しでのみ生成されます。ブロッキング系のメソッドは、スロット / アイテムを獲得するか予算を使い切るかのいずれかであり、後者の場合の結果は Full / Empty ではなく Timeout になります。到達可能性の完全なマトリクスは次のとおりです。

メソッド Ok Full / Empty Timeout Closed
trySend Full
send
sendTimeout
tryRecv Empty
recv
recvTimeout

ブロッキング呼び出しの結果に対する isFull() / isEmpty() のチェックはデッドコードです。以下の match の例では、この非対称性を読者に明示するために、それらのアームを unreachable コメントとして残しています。

結果のディスパッチ

RecvResultSendResult は、status() という判別子に加えて(RecvResult のみ)ペイロードを保持します。等価な 2 つのイディオムを示します。

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() は、Ok 以外のバリアントに対して呼び出されると OxPHP\Shared\SharedException をスローします。先に isOk() / valueOr() / status() を使用してください。これは意図的なものです。存在しないペイロードを操作するというバグを、静かな null ではなく、大きな失敗として表面化させます。

ファイバーとブロッキングの動作

同じメソッド呼び出しでも、PHP が現在ファイバー内で実行されているかどうかによって動作が異なります。

  • ファイバー内(ワーカーモード + oxphp_async(...)): ブロッキング系のメソッド(sendrecvsendTimeoutrecvTimeout)は、合成の Promise を確保し、チャネルに waker を登録して、ファイバーをサスペンドします。ワーカースレッドはスケジューラに戻り、チャネルが waker に通知するまで他のファイバーを処理します。
  • ファイバー外(従来モード、または非同期でない呼び出しパス): ブロッキング系のメソッドは crossbeam_channel を介して OS ワーカースレッドをブロックします。呼び出しが返るまで、そのスレッドでは他の作業は実行されません。

従来モードでもチャネルのセマンティクスは得られます。ただし、スレッドがブロックされるという代償を払うだけです。待機に依存するパイプラインには、ワーカーモードでのデプロイが推奨されます。

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() は冪等であり、2 回目の呼び出しは何もしません。クローズ後は次のようになります。

  • sendsendTimeouttrySendSendResult::Closed を返します。
  • sendMany は、クローズ前に実際に受け入れられた数を返します(すでにクローズされている場合は 0)。
  • recvrecvTimeouttryRecv は、バッファされたアイテムを 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 をスローします。

バッチ操作

sendManyrecvMany は、アイテムをグループ単位で移動させるパイプラインのために存在します。日常的に一度に 10 個以上のアイテムを扱う場合は、これらを優先してください。各バッチは N 回ではなく 1 回の FFI ラウンドトリップになるため、スループット律速のループにおけるアイテムごとのオーバーヘッドを大幅に削減できます。

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

注目すべきセマンティクスは次のとおりです。

  • どちらのバッチメソッドも、タイムアウトやバッチ途中でのクローズ時には部分的な結果を返し、例外を投げることはありません。部分的な進捗を検出するには、整数 / 配列の長さを確認してください。
  • すでにクローズされたチャネルに対する sendMany0 を返します。
  • クローズ済みかつ空のチャネルに対する recvMany[] を返します。
  • $ms> 0 でなければなりません。「バッファされているものをすべて取り出す」オーバーロードは存在しません。そのパターンが必要なら tryRecv() をタイトなループで呼び出してください。実際の呼び出し箇所から要望があれば、専用のバッチバリアントを追加します。

例外

例外 発生元
CapacityException new Channel($capacity) で、スロット配列が SHARED_MAX_CHANNEL_BYTES(デフォルト 64 MiB、約 200 万スロット)を超える場合。
ValueTooLargeException send / trySend / sendTimeout / sendMany で、値のシリアライズ後のサイズが SHARED_MAX_VALUE_SIZE(デフォルト 1 MiB)を超える場合。
TypeException 正でない $capacity、ゼロ・負・整数以外・欠落した $ms、あるいはスカラー・null・共有可能な値のネストした配列のいずれでもない値の送信。
SharedException Ok 以外のバリアントに対する RecvResult::value()、または clone $channel
AsyncException ファイバーのキャンセルによって中断されたブロッキングの recv / send。結果のバリアントではなく OxPHP\Async\AsyncException として表面化します。

可観測性

内部サーバー(デフォルト INTERNAL_ADDR=127.0.0.1:9090)は、汎用の共有レジストリエンドポイントでチャネルを公開します。

  • GET /__ox_shared/summary には、countbytesops を持つ Channel バケットが含まれます。
  • GET /__ox_shared/entries は、レジストリのエントリーを ID 付きで一覧します(limit クエリパラメータを受け付けます)。
  • GET /__ox_shared/entry?id=<id> は、チャネルごとの状態を返します: capacitycountpending (count の非推奨エイリアス)closedsenders_blockedreceivers_blocked

ops カウンター(およびレジストリ全体の oxphp_shared_operations_total{type="Channel"})は、結果にかかわらずすべての recv 試行をカウントします。ヒット、空 / クローズ済みのチャネル、タイムアウトした recvTimeout のいずれもこれをインクリメントします。このメトリクスは、正常に転送された値だけでなく、チャネルへのアクセスを追跡します。

/metrics での Prometheus エクスポジション:

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

複数のコンシューマーへのファンアウト

同じチャネル上に N 個の非同期コンシューマーを生成します。レジストリは、各アイテムをそのうちのちょうど 1 つが受信することを保証します。

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

バックプレッシャーを備えた有界パイプライン

trySend とドロップカウンターを組み合わせると、プロデューサーは過負荷時にブロックする代わりに負荷を切り捨てられます。

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

落とし穴

  • recvTimeout(0) は「非ブロッキング」ではなく TypeException です。 非ブロッキングには tryRecv() を使用してください。専用のメソッド名により、以前の ?float $timeout が持っていたタイムアウトオーバーロードの曖昧さが解消されました。
  • 値は共有可能でなければなりません。 スカラー、null、および共有可能な値のネストした配列が許可されます。Shared\* のインスタンスでないオブジェクトを渡すと、送信時に TypeException が発生します。
  • クローンは禁止されています。 clone $channelSharedException をスローします。代わりにクロージャの use を通じてチャネルを渡してください(oxphp_async(function () use ($ch) { ... }))。そうすれば、両側が同じレジストリエントリーを参照します。
  • Ok 以外に対する value() はスローします。 これはバグではなく機能です。「isOk のチェックを忘れた」というミスを、大きな失敗に変えてくれます。デフォルト値が本当に問題ないなら valueOr($default) を使用してください。
  • ファイバーのキャンセルは、結果のバリアントではなく例外として伝播します。 リクエストのキャンセルによって中断された recv()OxPHP\Async\AsyncException を発生させます。自らクローズするチャネルは、それでも RecvResult::Closed を生成します。
  • 転送中のペイロードを持つ、キャンセルされた待機者。 多数のファイバーが send / recv を待機していて、ペイロードがまさに受け渡される直前にキャンセルされた場合、そのペイロードは次の目覚めまで参照されたままになることがあります。待機者の数は制限しておいてください(例: Shared\Counter やチャネル容量のセマフォで並行処理に上限を設ける)。

以前の API からの移行

以前 現在
$ch->tryRecv() → 空のとき nullClosedException をスロー $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 でなければならない)

関連

  • ワーカーモード — ファイバーをサスペンドするブロッキング系のメソッドの前提条件です。
  • Async Promiseoxphp_async() クロージャは、Channel をバックグラウンドファイバーに渡す通常の方法です。
  • ファイバー多重化 — チャネル操作が待機している間、サスペンドによってワーカースレッドがどのように生産的に保たれるかを説明します。