Shared\Channel
OxPHP\Shared\Channel は、共有レジストリ内に存在し、プロセス内のすべての PHP ワーカーから見える、有界のマルチプロデューサー・マルチコンシューマーチャネルです。リクエストハンドラーとバックグラウンドワーカー、あるいは 2 つのワーカーが、FIFO 順で作業アイテムをやり取りする必要がある場合に使用します。ファイバー内では、ブロッキング呼び出しが協調的にサスペンドするため、基盤となるワーカースレッドは他のリクエストのために解放されたままになります。
概要
- 有界。 容量は構築時に固定されます。いっぱいになると
sendはブロックします(またはファイバーをサスペンドします)。trySendは待機せずに結果を報告します。 - MPMC。 スレッドをまたいで任意の数の送信側と受信側を持てます。配信は FIFO です。
- ファイバー対応。 非同期プールを備えたワーカーモードでは、ブロッキング系のメソッドは OS スレッドではなくファイバーをサスペンドします。従来モードでは OS スレッドをブロックします。
- 結果型付き。 すべての send/recv は
Channel\SendResult/Channel\RecvResultを返します。クローズ / フル / タイムアウトはファンアウトのディスパッチャーにとって通常の結果なので、例外ではなく結果のバリアントとして現れます。
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;
}サイズ制限
信頼できない入力によるアロケーション爆弾から保護するために、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をスローします。これはsend、trySend、sendTimeout、およびsendManyに適用されます(sendManyでは、いずれかの 要素が上限を超えている場合、バッチ全体を拒否し、何も送信しません)。これはShared\Mapが適用するのと同じ、値ごとの上限です。
待機ポリシー
どちらの方向にも 3 つのメソッドのバリアントがあり、サフィックスが待機ポリシーを表します。
| サフィックス | 動作 |
|---|---|
try* |
非ブロッキング。呼び出しを続行できない場合は Empty / Full を報告します。 |
| (サフィックスなし) | 永久に、またはリクエストファイバーがキャンセルされるまでブロックします。 |
*Timeout |
制限付き待機。$ms は > 0 でなければなりません。永久待機や非ブロッキングには他の形式のいずれかを使用してください。 |
recvTimeout / sendTimeout / sendMany / recvMany において、$ms は常にミリ秒単位の正の整数です。ゼロ、負の値、整数以外、および値の欠落は、ブリッジで OxPHP\Shared\TypeException を発生させます。この制約はメソッド本体から外に移されたため、3 つの形式の区別が自己文書化されるようになりました。
メソッドごとに到達可能な結果バリアント
SendResult::Full と RecvResult::Empty は、非ブロッキングの try* 呼び出しでのみ生成されます。ブロッキング系のメソッドは、スロット / アイテムを獲得するか予算を使い切るかのいずれかであり、後者の場合の結果は Full / Empty ではなく Timeout になります。到達可能性の完全なマトリクスは次のとおりです。
| メソッド | Ok |
Full / Empty |
Timeout |
Closed |
|---|---|---|---|---|
trySend |
✓ | Full ✓ |
— | ✓ |
send |
✓ | — | — | ✓ |
sendTimeout |
✓ | — | ✓ | ✓ |
tryRecv |
✓ | Empty ✓ |
— | ✓ |
recv |
✓ | — | — | ✓ |
recvTimeout |
✓ | — | ✓ | ✓ |
ブロッキング呼び出しの結果に対する isFull() / isEmpty() のチェックはデッドコードです。以下の match の例では、この非対称性を読者に明示するために、それらのアームを unreachable コメントとして残しています。
結果のディスパッチ
RecvResult と SendResult は、status() という判別子に加えて(RecvResult のみ)ペイロードを保持します。等価な 2 つのイディオムを示します。
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(...)): ブロッキング系のメソッド(send、recv、sendTimeout、recvTimeout)は、合成の Promise を確保し、チャネルに waker を登録して、ファイバーをサスペンドします。ワーカースレッドはスケジューラに戻り、チャネルが waker に通知するまで他のファイバーを処理します。 - ファイバー外(従来モード、または非同期でない呼び出しパス): ブロッキング系のメソッドは
crossbeam_channelを介して OS ワーカースレッドをブロックします。呼び出しが返るまで、そのスレッドでは他の作業は実行されません。
従来モードでもチャネルのセマンティクスは得られます。ただし、スレッドがブロックされるという代償を払うだけです。待機に依存するパイプラインには、ワーカーモードでのデプロイが推奨されます。
<?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 回目の呼び出しは何もしません。クローズ後は次のようになります。
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 個以上のアイテムを扱う場合は、これらを優先してください。各バッチは N 回ではなく 1 回の FFI ラウンドトリップになるため、スループット律速のループにおけるアイテムごとのオーバーヘッドを大幅に削減できます。
<?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、約 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には、count、bytes、opsを持つChannelバケットが含まれます。GET /__ox_shared/entriesは、レジストリのエントリーを ID 付きで一覧します(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 のいずれもこれをインクリメントします。このメトリクスは、正常に転送された値だけでなく、チャネルへのアクセスを追跡します。
/metrics での Prometheus エクスポジション:
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";
});複数のコンシューマーへのファンアウト
同じチャネル上に N 個の非同期コンシューマーを生成します。レジストリは、各アイテムをそのうちのちょうど 1 つが受信することを保証します。
<?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
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) { ... }))。そうすれば、両側が同じレジストリエントリーを参照します。 - Ok 以外に対する
value()はスローします。 これはバグではなく機能です。「isOk のチェックを忘れた」というミスを、大きな失敗に変えてくれます。デフォルト値が本当に問題ないならvalueOr($default)を使用してください。 - ファイバーのキャンセルは、結果のバリアントではなく例外として伝播します。 リクエストのキャンセルによって中断された
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 でなければならない) |
関連
- ワーカーモード — ファイバーをサスペンドするブロッキング系のメソッドの前提条件です。
- Async Promise —
oxphp_async()クロージャは、Channelをバックグラウンドファイバーに渡す通常の方法です。 - ファイバー多重化 — チャネル操作が待機している間、サスペンドによってワーカースレッドがどのように生産的に保たれるかを説明します。