Lines
62.50%
20 / 32
Methods
66.66%
4 / 6
Classes
0.00%
0 / 1
| Name | Lines | Methods | CRAP | ||||
|---|---|---|---|---|---|---|---|
| __construct | 100.00% | 5 / 5 | 100.00% | 1 / 1 | 2 | ||
| push | 100.00% | 4 / 4 | 100.00% | 1 / 1 | 2 | ||
| close | 100.00% | 2 / 2 | 100.00% | 1 / 1 | 1 | ||
| stream | 35.29% | 6 / 17 | 0.00% | 0 / 1 | 30.94 | ||
| response | 100.00% | 1 / 1 | 100.00% | 1 / 1 | 1 | ||
| signal | 66.66% | 2 / 3 | 0.00% | 0 / 1 | 2.15 | ||
| 30 | final class EventStream | |
| 31 | { | |
| 32 | private SplQueue $queue; | |
| 33 | private bool $closed = false; | |
| 34 | ||
| 35 | /** @var resource|null Self-pipe write end — push()/close() signal it */ | |
| 36 | private $signalWrite = null; | |
| 37 | ||
| 38 | /** @var resource|null Self-pipe read end — stream() blocks on it */ | |
| 39 | private $signalRead = null; | |
| 40 | ||
| 41 | public function __construct() | |
| 42 | { | |
| 43 | $this->queue = new SplQueue(); | |
| 44 | ||
| 45 | // Self-pipe trick: a socket pair lets push()/close() wake a | |
| 46 | // blocked stream() immediately — no polling, no lost wakeups. | |
| 47 | ||
| 48 | // Protocol 0 = no protocol — correct for UNIX domain sockets. | |
| 49 | // (STREAM_IPPROTO_IP is for TCP/IP and only works here by accident | |
| 50 | // on Linux; 0 is portable.) | |
| 51 | ||
| 52 | $pair = stream_socket_pair(STREAM_PF_UNIX, STREAM_SOCK_STREAM, 0); | |
| 53 | if ($pair !== false) { | |
| 54 | [$this->signalRead, $this->signalWrite] = $pair; | |
| 55 | stream_set_blocking($this->signalWrite, false); | |
| 56 | } | |
| 57 | } | |
| 58 | ||
| 59 | /** | |
| 60 | * Queue an event to be streamed. Safe to call from any context | |
| 61 | * (event listener, worker, timer, signal handler). Ignored after close(). | |
| 62 | * Wakes a blocked stream() immediately (via the self-pipe). | |
| 63 | */ | |
| 64 | public function push(Event $event): void | |
| 65 | { | |
| 66 | if ($this->closed) { | |
| 67 | return; | |
| 68 | } | |
| 69 | $this->queue->enqueue($event->toSSE()); | |
| 70 | $this->signal(); | |
| 71 | } | |
| 72 | ||
| 73 | /** | |
| 74 | * End the stream: the generator returns after draining pending events. | |
| 75 | * Wakes a blocked stream() immediately (via the self-pipe). | |
| 76 | * | |
| 77 | * Useful for finite streams (e.g. send N events then finish). | |
| 78 | */ | |
| 79 | public function close(): void | |
| 80 | { | |
| 81 | $this->closed = true; | |
| 82 | $this->signal(); | |
| 83 | } | |
| 84 | ||
| 85 | /** | |
| 86 | * Pull queued events as a generator — pass this to withEventStream(). | |
| 87 | * | |
| 88 | * Blocks (via stream_select on a self-pipe) until an event is available | |
| 89 | * or the stream is closed — zero idle latency, no polling. If the socket | |
| 90 | * pair could not be created, falls back to a short usleep poll. | |
| 91 | * | |
| 92 | * The self-pipe makes push()/close() from signal handlers, threads, or | |
| 93 | * other processes wake the blocked generator immediately. A signal | |
| 94 | * interrupting stream_select (EINTR) is handled by re-checking the queue. | |
| 95 | * | |
| 96 | * @return Generator | |
| 97 | */ | |
| 98 | public function stream(): Generator | |
| 99 | { | |
| 100 | while (true) { | |
| 101 | if (! $this->queue->isEmpty()) { | |
| 102 | yield $this->queue->dequeue(); | |
| 103 | continue; | |
| 104 | } | |
| 105 | if ($this->closed) { | |
| 106 | return; | |
| 107 | } | |
| 108 | ||
| 109 | if ($this->signalRead !== null) { | |
| 110 | // Block until push()/close() writes a byte (or a signal | |
| 111 | // interrupts us — then re-check the queue). | |
| 112 | $read = [$this->signalRead]; | |
| 113 | $write = null; | |
| 114 | $except = null; | |
| 115 | if (@stream_select($read, $write, $except, null) === false) { | |
| 116 | continue; // interrupted (e.g. pcntl signal) — re-check | |
| 117 | } | |
| 118 | // Drain the wake-up byte(s) so the next stream_select blocks. | |
| 119 | ||
| 120 | while (!feof($this->signalRead)) { | |
| 121 | $chunk = fread($this->signalRead, 8192); | |
| 122 | if ($chunk === '' || $chunk === false) { | |
| 123 | break; | |
| 124 | } | |
| 125 | } | |
| 126 | } else { | |
| 127 | usleep(50_000); // fallback: pair creation failed | |
| 128 | } | |
| 129 | } | |
| 130 | } | |
| 131 | ||
| 132 | /** | |
| 133 | * Build an SSE response wired to this stream's generator. | |
| 134 | */ | |
| 135 | public function response(): Response | |
| 136 | { | |
| 137 | return (new Response())->withEventStream($this->stream()); | |
| 138 | } | |
| 139 | ||
| 140 | /** | |
| 141 | * Wake a blocked stream() by writing a byte to the self-pipe. | |
| 142 | * | |
| 143 | * Non-blocking write; failures (e.g. pipe full) are ignored — the | |
| 144 | * queue is the source of truth, the pipe is just a wake-up signal. | |
| 145 | * | |
| 146 | * @return void | |
| 147 | */ | |
| 148 | private function signal(): void | |
| 149 | { | |
| 150 | if ($this->signalWrite === null) { | |
| 151 | return; | |
| 152 | } | |
| 153 | @fwrite($this->signalWrite, "\0"); | |
| 154 | } | |
| 155 | } |