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
30final 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}