Shared\Channel

OxPHP\Shared\Channel 是一个有界的多生产者多消费者通道,它存活于共享注册表中,进程内每一个 PHP 工作进程都能看到它。当一个请求处理器与一个后台工作进程、或者两个工作进程之间需要按 FIFO 顺序交换工作项时,就可以使用它。在纤程内部,阻塞调用会协作式挂起,从而让底层的工作进程线程腾出来处理其他请求。

概览

  • 有界。 容量在构造时固定。一旦满了,send 就会阻塞(或挂起纤程);trySend 则不等待,直接报告结果。
  • MPMC。 跨线程的发送者和接收者数量不限。投递遵循 FIFO。
  • 纤程感知。 在工作进程模式配合异步池时,阻塞变体会挂起纤程而非 OS 线程。在传统模式下,它们会阻塞 OS 线程。
  • 结果类型化。 每次 send/recv 都返回一个 Channel\SendResult / Channel\RecvResult——对于扇出分发器而言,closed / full / timeout 都是正常的结果,因此它们表现为结果变体而非异常。

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

大小限制

有两道上限用于防范来自不可信输入的分配炸弹;两者都可以 通过环境变量配置,并且抛出可捕获的异常, 而不会中止进程:

  • 容量 —— new Channel($capacity) 会预分配一个槽位数组。如果某个 容量对应的槽位数组会超过 SHARED_MAX_CHANNEL_BYTES (默认 64 MiB,约 2M 个槽位),则抛出 CapacityException。通道 通常只有几十到几千的深度,因此这只在病态 数值下才会触发。(非正的 $capacity 是另一个 TypeException。) 有效预算会被向上钳制到至少一个槽位,因此一个为零或 过小的 SHARED_MAX_CHANNEL_BYTES 永远不会拒绝一个最小的 容量为 1 的通道——它只会对大容量设上限。
  • 每个值 —— 每个发送的值都会被序列化;如果某个值的序列化后 大小超过 SHARED_MAX_VALUE_SIZE(默认 1 MiB),则抛出 ValueTooLargeException。这适用于 sendtrySendsendTimeoutsendMany(若任一元素超过上限,它会拒绝整个批次, 什么都不发送)。这与 Shared\Map 强制执行的每值上限相同。

等待策略

每个方向都有三种方法变体;后缀编码了等待策略:

后缀 行为
try* 非阻塞。如果调用无法继续,则报告 Empty / Full
(裸名) 永远阻塞,或直到请求纤程被取消。
*Timeout 有界等待。$ms 必须 > 0——若要永远等待或非阻塞,请使用其他形式。

对于 recvTimeout / sendTimeout / sendMany / recvMany 来说,$ms 始终是一个以毫秒为单位的正整数。零、负数、非整数以及缺省值都会在桥接层抛出 OxPHP\Shared\TypeException——这一约束已从方法体中移出,从而让这三分法自我说明。

每个方法可达的结果变体

SendResult::FullRecvResult::Empty 只由非阻塞的 try* 调用发出——阻塞变体要么获取到槽位/项,要么耗尽预算,此时结果是 Timeout,而不是 Full / Empty。完整的可达性矩阵:

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

对来自阻塞调用的结果做 isFull() / isEmpty() 检查是死代码。下面的 match 示例把这些分支留作 unreachable 注释,以便让这种不对称性对读者一目了然。

结果分发

RecvResultSendResult 携带一个 status() 判别值,外加(仅 RecvResult)一个负载。两种等价写法:

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

如果在非 Ok 变体上调用 RecvResult::value(),它会抛出 OxPHP\Shared\SharedException——请先用 isOk() / valueOr() / status()。这是有意为之:它把对缺失负载采取行动这一 bug 暴露为一次响亮的失败,而不是悄无声息的 null

纤程行为与阻塞行为

同样的方法调用,其行为会根据 PHP 当前是否运行在纤程内部而有所不同:

  • 在纤程内部(工作进程模式 + oxphp_async(...)):阻塞变体(sendrecvsendTimeoutrecvTimeout)会分配一个合成 promise,向通道注册一个唤醒器,并挂起纤程。工作进程线程回到调度器,处理其他纤程,直到通道通知该唤醒器。
  • 在纤程外部(传统模式,或非异步的调用路径):阻塞变体会通过 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() 是幂等的;再次调用它是一次空操作。关闭之后:

  • sendsendTimeouttrySend 返回 SendResult::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 个以上的项时,优先使用它们:每个批次只是一次 FFI 往返而非 N 次,这能在吞吐量受限的循环中明显削减每项开销。

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

值得注意的语义:

  • 两个批量方法在超时或批处理中途关闭时都返回部分结果——绝不抛出异常。检查整数 / 数组长度以检测部分进展。
  • 对一个已经关闭的通道调用 sendMany 返回 0
  • 对已关闭且为空的通道调用 recvMany 返回 []
  • $ms 必须 > 0。没有“排空一切已缓冲内容”的重载——如果需要那种模式,请在紧凑循环中调用 tryRecv();等到有真实调用点提出需求时,我们会加上一个专门的批量变体。

异常

异常 由谁抛出
CapacityException new Channel($capacity),当槽位数组会超过 SHARED_MAX_CHANNEL_BYTES(默认 64 MiB,约 2M 个槽位)时。
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,而非某个 Result 变体。

可观测性

内部服务器(默认 INTERNAL_ADDR=127.0.0.1:9090)在通用的共享注册表端点中暴露通道:

  • GET /__ox_shared/summary 包含一个 Channel 分桶,带有 countbytesops
  • 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 个异步消费者;注册表确保它们中恰好有一个收到每个项。

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 $channel 会抛出 SharedException;请改为通过闭包的 use 传递通道——oxphp_async(function () use ($ch) { ... })——这样双方看到的是同一个注册表条目。
  • 在非 Ok 上调用 value() 会抛异常。 这是特性,不是 bug——它把“我忘了检查 isOk”的失误转化为一次响亮的失败。如果一个默认值确实可以接受,就用 valueOr($default)
  • 纤程取消以异常形式传播,而非某个 Result 变体。一个被请求取消打断的 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

相关

  • 工作进程模式 —— 纤程挂起式阻塞变体的前提。
  • 异步 Promise —— oxphp_async() 闭包是把 Channel 交给后台纤程的常规方式。
  • 纤程多路复用 —— 解释挂起如何在通道操作等待期间保持工作进程线程的生产力。