From e8b3dc6ee0f23cf652bfd6d3b56492e22012f9d7 Mon Sep 17 00:00:00 2001 From: "A. B. M. Mahmudul Hasan" Date: Sun, 4 Oct 2026 17:41:57 +0600 Subject: [PATCH 1/2] :sparkles: refactor(http): finish accepted requests and headers before drain - Refactor native HTTP workers and protocol attachments to finish pending request admissions before initiating application drain :recycle: - Update HTTP/1.1, HTTP/2, and HTTP/3 connection handlers to track uncompleted initial exchanges and header blocks :art: - Add tests covering request recycling and incomplete headers during draining across HTTP versions :white_check_mark: - Update workflow configurations to support release tags and condition checks :construction_worker: - Document graceful sequence and protocol draining behavior during application shutdown :memo: --- .github/workflows/security-standards.yml | 11 +++ docs/deployment.md | 7 +- src/Http/Http1/Http1Connection.php | 10 ++- src/Http/Http2/Http2Connection.php | 10 ++- .../Http2/Internal/RequestStreamProcessor.php | 12 ++- src/Http/Http3/Http3Session.php | 6 ++ .../Http3/Quic/PhpQuicHttp3Connection.php | 6 ++ src/Http/Http3/Quic/PhpQuicHttp3Worker.php | 6 ++ src/Http/NativeHttpConnection.php | 6 ++ .../Internal/NativeHttp3Attachment.php | 15 +++- src/Runtime/Internal/NativeHttp3Worker.php | 19 ++++- src/Runtime/Internal/NativeHttpWorker.php | 53 +++++++++--- tests/Http2ConnectionTest.php | 48 +++++++++++ tests/NativeWorkerLifecycleTest.php | 80 +++++++++++++++++++ tests/PhpQuicHttp3ConnectionTest.php | 64 +++++++++++++++ 15 files changed, 334 insertions(+), 19 deletions(-) diff --git a/.github/workflows/security-standards.yml b/.github/workflows/security-standards.yml index ad8444f9..1533da3b 100644 --- a/.github/workflows/security-standards.yml +++ b/.github/workflows/security-standards.yml @@ -5,6 +5,7 @@ on: - cron: "0 0 * * 0" push: branches: [ "main", "master" ] + tags: [ "v*", "[0-9]*" ] pull_request: branches: [ "main", "master", "develop", "development" ] @@ -14,6 +15,7 @@ concurrency: jobs: phpforge: + if: github.event_name != 'push' || !startsWith(github.ref, 'refs/tags/') uses: infocyph/phpforge/.github/workflows/security-standards.yml@main permissions: security-events: write @@ -29,7 +31,16 @@ jobs: integration_services: '[]' service_topologies: '{}' + release: + if: github.event_name == 'push' && startsWith(github.ref, 'refs/tags/') + uses: infocyph/phpforge/.github/workflows/release.yml@main + permissions: + contents: write + secrets: + COPILOT_GITHUB_TOKEN: ${{ secrets.COPILOT_GITHUB_TOKEN }} + http3-quic: + if: github.event_name != 'push' || !startsWith(github.ref, 'refs/tags/') name: "HTTP/3 + QUIC - PHP ${{ matrix.php-version }}" runs-on: ubuntu-latest container: "php:${{ matrix.php-version }}-cli" diff --git a/docs/deployment.md b/docs/deployment.md index f61cea4d..7713bfe2 100644 --- a/docs/deployment.md +++ b/docs/deployment.md @@ -250,14 +250,19 @@ Graceful sequence: ```text stop accepting -→ application drain → protocol drain / GOAWAY +→ finish admission of previously accepted requests +→ application drain → finish admitted work → timeout → force remaining work → close resources ``` +Native HTTP/1.1 connections accepted before stopping may finish their first exchange, including incomplete headers. Protocol draining prevents another keep-alive exchange. Application draining begins once these first requests have entered the application or their connections have closed; header/body deadlines and the recycling grace period remain enforced. Idle connections cannot extend recycling beyond that grace period. + +HTTP/2 includes an accepted but incomplete HEADERS/CONTINUATION block in the GOAWAY boundary and lets it finish; higher stream IDs remain refused. HTTP/3 likewise lets previously accepted request streams finish decoding their headers before application draining begins. Both protocols retain their existing resource limits and drain deadlines. Native HTTP/3 continues advancing application tasks while draining. + Portable mode drains all attachments on its shared `SelectLoop` until complete or deadline expiry. ## 11. Application lifecycle hooks diff --git a/src/Http/Http1/Http1Connection.php b/src/Http/Http1/Http1Connection.php index 7eac60fc..39fcfb8f 100644 --- a/src/Http/Http1/Http1Connection.php +++ b/src/Http/Http1/Http1Connection.php @@ -141,7 +141,7 @@ function (): void { } /** - * Stop accepting further keep-alive work and drain the active exchange. + * Stop further keep-alive work and finish the accepted initial exchange. */ public function drain(): void { @@ -158,11 +158,17 @@ public function drain(): void $this->responseCloseAfter = true; $this->writer?->forceCloseAfterResponse(); - if ($this->body === null && $this->writer === null) { + if ($this->writer === null && $this->requestCount > 0) { $this->connection->closeGracefully(); } } + /** @internal Reports an accepted connection awaiting its first request. */ + public function hasPendingRequestAdmission(): bool + { + return $this->requestCount === 0; + } + /** * Return the number of requests processed on this connection. */ diff --git a/src/Http/Http2/Http2Connection.php b/src/Http/Http2/Http2Connection.php index dfd4397f..abc5d506 100644 --- a/src/Http/Http2/Http2Connection.php +++ b/src/Http/Http2/Http2Connection.php @@ -135,7 +135,7 @@ public function drain(): void $this->draining = true; $this->requests->setDraining(true); - $this->output->sendControl(FrameWriter::goAway($this->requests->lastClientStreamId(), ErrorCode::NO_ERROR)); + $this->output->sendControl(FrameWriter::goAway($this->requests->lastAcceptedStreamId(), ErrorCode::NO_ERROR)); $this->finishDrainIfReady(); if ($this->closed) { return; @@ -150,6 +150,12 @@ public function drain(): void }); } + /** @internal Reports an accepted header block awaiting request dispatch. */ + public function hasPendingRequestAdmission(): bool + { + return $this->requests->hasOpenHeaderBlock(); + } + /** * Return the greatest client-initiated stream ID observed. */ @@ -269,7 +275,7 @@ private function failConnection(ErrorCode $code, string $message): void private function finishDrainIfReady(): void { - if (!$this->draining || $this->closed || $this->requests->count() !== 0 || !$this->output->wireIdle()) { + if (!$this->draining || $this->closed || $this->requests->count() !== 0 || $this->requests->hasOpenHeaderBlock() || !$this->output->wireIdle()) { return; } diff --git a/src/Http/Http2/Internal/RequestStreamProcessor.php b/src/Http/Http2/Internal/RequestStreamProcessor.php index d9fd5b1e..a2d67aef 100644 --- a/src/Http/Http2/Internal/RequestStreamProcessor.php +++ b/src/Http/Http2/Internal/RequestStreamProcessor.php @@ -42,7 +42,7 @@ final class RequestStreamProcessor private readonly RequestHeaderValidator $validator; - private bool $draining = false; + private ?int $drainBoundary = null; private ?int $headerBlockTimer = null; @@ -203,6 +203,12 @@ public function hasOpenHeaderBlock(): bool return $this->pendingHeaders !== null; } + /** Return the boundary including an accepted unfinished header block. */ + public function lastAcceptedStreamId(): int + { + return max($this->lastClientStreamId, $this->pendingHeaders->streamId ?? 0); + } + /** * Return the greatest client-initiated stream ID observed. */ @@ -226,7 +232,7 @@ public function reset(Http2Stream $stream): void */ public function setDraining(bool $draining): void { - $this->draining = $draining; + $this->drainBoundary = $draining ? $this->lastAcceptedStreamId() : null; } /** @@ -295,7 +301,7 @@ private function acceptData(Http2Stream $stream, Frame $frame): void private function admit(PendingHeaderBlock $pending, ?int $contentLength): bool { - if ($this->draining) { + if ($this->drainBoundary !== null && $pending->streamId > $this->drainBoundary) { $this->output->sendControl(FrameWriter::rstStream($pending->streamId, ErrorCode::REFUSED_STREAM)); return false; diff --git a/src/Http/Http3/Http3Session.php b/src/Http/Http3/Http3Session.php index 13a8a05b..e97126d0 100644 --- a/src/Http/Http3/Http3Session.php +++ b/src/Http/Http3/Http3Session.php @@ -72,6 +72,12 @@ public function finishRequestStream(int $streamId): void $this->dispatchReadyRequests(); } + /** @internal Reports whether an accepted stream still needs application admission. */ + public function hasPendingRequestAdmission(int $streamId): bool + { + return !isset($this->dispatchedRequestStreams[$streamId]); + } + /** * Feed bytes from a peer unidirectional stream into connection state. */ diff --git a/src/Http/Http3/Quic/PhpQuicHttp3Connection.php b/src/Http/Http3/Quic/PhpQuicHttp3Connection.php index 70f0327d..cfa61566 100644 --- a/src/Http/Http3/Quic/PhpQuicHttp3Connection.php +++ b/src/Http/Http3/Quic/PhpQuicHttp3Connection.php @@ -215,6 +215,12 @@ public function handleReady(array $ready, PhpQuicEventMasks $events): void } } + /** @internal Reports accepted streams awaiting headers or QPACK decoding. */ + public function hasPendingRequestAdmission(): bool + { + return array_any(array_keys($this->requestStreams), $this->session->hasPendingRequestAdmission(...)); + } + /** * Return a tracked peer stream by QUIC stream ID. */ diff --git a/src/Http/Http3/Quic/PhpQuicHttp3Worker.php b/src/Http/Http3/Quic/PhpQuicHttp3Worker.php index 141abd08..9438c3c7 100644 --- a/src/Http/Http3/Quic/PhpQuicHttp3Worker.php +++ b/src/Http/Http3/Quic/PhpQuicHttp3Worker.php @@ -122,6 +122,12 @@ public function forceClose(): void } } + /** @internal Reports accepted request streams awaiting application admission. */ + public function hasPendingRequestAdmission(): bool + { + return array_any($this->connections, fn($connection) => $connection->hasPendingRequestAdmission()); + } + /** * Stop accepting new connections and begin graceful draining. */ diff --git a/src/Http/NativeHttpConnection.php b/src/Http/NativeHttpConnection.php index d1eafb4c..1d56d867 100644 --- a/src/Http/NativeHttpConnection.php +++ b/src/Http/NativeHttpConnection.php @@ -66,4 +66,10 @@ public function drain(): void { $this->protocol->drain(); } + + /** @internal Reports accepted input whose request has not reached the application. */ + public function hasPendingRequestAdmission(): bool + { + return $this->protocol->hasPendingRequestAdmission(); + } } diff --git a/src/Runtime/Internal/NativeHttp3Attachment.php b/src/Runtime/Internal/NativeHttp3Attachment.php index 0df81240..b6ff7b63 100644 --- a/src/Runtime/Internal/NativeHttp3Attachment.php +++ b/src/Runtime/Internal/NativeHttp3Attachment.php @@ -21,6 +21,8 @@ final class NativeHttp3Attachment { private readonly WorkerStopState $state; + private bool $applicationDraining = false; + private ?int $drainTimer = null; private ?int $pollTimer = null; @@ -76,8 +78,8 @@ private function beginDrain(): void } $this->state->stop(); - $this->application->drain($this->context->shutdownReason()); $this->worker->stopAccepting(); + $this->drainApplication(); $this->drainTimer = $this->loop->delay( $this->context->recyclePolicy->gracefulTimeoutSeconds, $this->expireDrain(...), @@ -121,6 +123,16 @@ private function close(): void $this->shutdown(); } + private function drainApplication(): void + { + if (!$this->state->isStopping() || $this->applicationDraining || $this->worker->hasPendingRequestAdmission()) { + return; + } + + $this->applicationDraining = true; + $this->application->drain($this->context->shutdownReason()); + } + private function drained(): bool { return $this->state->isStopping() && $this->worker->drainComplete(); @@ -136,6 +148,7 @@ private function expireDrain(): void private function finishDrain(): void { + $this->drainApplication(); if (!$this->drained()) { return; } diff --git a/src/Runtime/Internal/NativeHttp3Worker.php b/src/Runtime/Internal/NativeHttp3Worker.php index 1e36c72c..7fef9fb5 100644 --- a/src/Runtime/Internal/NativeHttp3Worker.php +++ b/src/Runtime/Internal/NativeHttp3Worker.php @@ -145,8 +145,8 @@ public static function run( } $context->consumeStopWake(); - $application->drain($context->shutdownReason()); $worker->stopAccepting(); + $applicationDraining = false; $deadline = $context->recycling() ? MonotonicTime::deadlineAfterSeconds( MonotonicTime::nowNanoseconds(), @@ -154,15 +154,18 @@ public static function run( ) : null; while (!$worker->drainComplete()) { + self::drainApplication($application, $context, $worker, $applicationDraining); if ($deadline !== null && MonotonicTime::nowNanoseconds() >= $deadline) { $worker->forceClose(); break; } $worker->tick($options->pollTimeoutSeconds); + $taskLoop->tick(); self::observeTransport($runtimeContext, $worker); $sampler->sample(); } + self::drainApplication($application, $context, $worker, $applicationDraining); self::observeTransport($runtimeContext, $worker); $sampler->sample(true); } catch (Throwable $error) { @@ -193,6 +196,20 @@ public static function run( ApplicationShutdown::resolve($failure, $shutdownFailures); } + private static function drainApplication( + \Infocyph\Runwire\Runtime\RuntimeApplicationInterface $application, + WorkerContext $context, + PhpQuicHttp3Worker $worker, + bool &$applicationDraining, + ): void { + if ($applicationDraining || $worker->hasPendingRequestAdmission()) { + return; + } + + $applicationDraining = true; + $application->drain($context->shutdownReason()); + } + /** @return array{0: string, 1: int} */ private static function endpoint(string $address): array { diff --git a/src/Runtime/Internal/NativeHttpWorker.php b/src/Runtime/Internal/NativeHttpWorker.php index 603ee3d1..b8bc75c7 100644 --- a/src/Runtime/Internal/NativeHttpWorker.php +++ b/src/Runtime/Internal/NativeHttpWorker.php @@ -66,9 +66,13 @@ public static function attach( $loop, ); $sampler = new WorkerDiagnosticsSampler($context, $runtimeContext->metrics, $diagnostics, $loop); - $handler = self::requestHandler($application, $context, $sampler); $stopWatcher = null; $shutdown = false; + $applicationDraining = false; + $drainApplication = static function () use ($application, $context, $state, &$sessions, &$applicationDraining): void { + self::drainApplication($application, $context, $state, $sessions, $applicationDraining); + }; + $handler = self::requestHandler($application, $context, $sampler, $drainApplication); try { $application->start(); @@ -79,6 +83,7 @@ static function (Connection $connection) use ( $loop, $bound, $handler, + $drainApplication, &$sessions, &$connections, $http1Adaptive, @@ -109,12 +114,13 @@ static function (Connection $connection) use ( $runtimeContext->metrics, $sampler, ); + $connection->onClose($drainApplication); }, ); $stopWatcher = $loop->onReadable( $context->stopStream(), - static function () use ($application, $context, $bound, $loop, &$sessions, &$connections, $state, $sampler): void { - self::beginDrain($application, $context, $bound, $loop, $sessions, $connections, $state); + static function () use ($drainApplication, $context, $bound, $loop, &$sessions, &$connections, $state, $sampler): void { + self::beginDrain($drainApplication, $context, $bound, $loop, $sessions, $connections, $state); $sampler->sample(true); }, ); @@ -143,7 +149,7 @@ static function () use ($application, $context, $bound, $loop, &$sessions, &$con $context->requestStop(); }, forceStop: static function () use ( - $application, + $drainApplication, $context, $bound, $loop, @@ -153,7 +159,7 @@ static function () use ($application, $context, $bound, $loop, &$sessions, &$con $sampler, ): void { $context->requestStop(); - self::beginDrain($application, $context, $bound, $loop, $sessions, $connections, $state); + self::beginDrain($drainApplication, $context, $bound, $loop, $sessions, $connections, $state); self::forceClose($connections, $loop, $state); $sampler->sample(true); }, @@ -307,11 +313,12 @@ private static function attachConnection( } /** + * @param Closure(): void $drainApplication * @param array $sessions * @param array $connections */ private static function beginDrain( - RuntimeApplicationInterface $application, + Closure $drainApplication, WorkerContext $context, BoundServer $bound, LoopInterface $loop, @@ -324,9 +331,9 @@ private static function beginDrain( return; } - $application->drain($context->shutdownReason()); $state->stop(); $bound->listener->close(); + $drainApplication(); foreach ($sessions as $session) { $session->drain(); } @@ -342,6 +349,28 @@ static function () use (&$connections, $loop, $state): void { } } + /** @param array $sessions */ + private static function drainApplication( + RuntimeApplicationInterface $application, + WorkerContext $context, + WorkerStopState $state, + array $sessions, + bool &$applicationDraining, + ): void { + if (!$state->isStopping() || $applicationDraining) { + return; + } + + foreach ($sessions as $session) { + if ($session->hasPendingRequestAdmission()) { + return; + } + } + + $applicationDraining = true; + $application->drain($context->shutdownReason()); + } + /** @param array $connections */ private static function forceClose(array $connections, LoopInterface $loop, WorkerStopState $state): void { @@ -351,13 +380,17 @@ private static function forceClose(array $connections, LoopInterface $loop, Work $state->stopLoopIfStopping($loop); } - /** @return Closure(HttpRequest, ResponseWriterInterface): void */ + /** + * @param Closure(): void|null $drainApplication + * @return Closure(HttpRequest, ResponseWriterInterface): void + */ private static function requestHandler( RuntimeApplicationInterface $application, WorkerContext $context, WorkerDiagnosticsSampler $sampler, + ?Closure $drainApplication = null, ): Closure { - return static function (HttpRequest $request, ResponseWriterInterface $writer) use ($application, $context, $sampler): void { + return static function (HttpRequest $request, ResponseWriterInterface $writer) use ($application, $context, $sampler, $drainApplication): void { $context->recordRequestStarted(); $completed = false; $complete = static function () use ( @@ -392,6 +425,8 @@ private static function requestHandler( } throw $error; + } finally { + $drainApplication?->__invoke(); } }; } diff --git a/tests/Http2ConnectionTest.php b/tests/Http2ConnectionTest.php index 398585d5..1a2ae6a9 100644 --- a/tests/Http2ConnectionTest.php +++ b/tests/Http2ConnectionTest.php @@ -292,3 +292,51 @@ static function (HttpRequest $request, ResponseWriterInterface $writer) use (&$w ->and($http2->activeStreams())->toBe(32) ->and($ended)->toBe(3); }); + + +it('finishes accepted fragmented HTTP2 headers across drain and refuses later streams', function (): void { + [$server, $client] = stream_socket_pair(STREAM_PF_UNIX, STREAM_SOCK_STREAM, STREAM_IPPROTO_IP); + $loop = new SelectLoop(); + $connection = new Connection($loop, $server); + $requests = []; + $http2 = new Http2Connection($loop, $connection, new Http2Limits(), static function (HttpRequest $request, ResponseWriterInterface $writer) use ($loop, &$requests): void { + $requests[] = $request->target; + $loop->delay(0.02, static fn() => $writer->end('accepted-h2')); + }); + $block = (new Encoder())->encode([ + [':method', 'GET'], [':scheme', 'https'], [':authority', 'example.test'], [':path', '/first'], + ]); + + try { + fwrite($client, runwireH2ClientPrelude() . FrameWriter::encode(new Frame(FrameType::HEADERS->value, 0x1, 1, substr($block, 0, 1)))); + $loop->delay(0.01, static fn() => $http2->drain()); + $loop->delay(0.02, static function () use ($client, $block): void { + fwrite($client, FrameWriter::encode(new Frame(FrameType::CONTINUATION->value, 0x4, 1, substr($block, 1))) . runwireH2Headers(3, '/later')); + }); + $loop->delay(0.08, static fn() => $loop->stop()); + $loop->run(); + stream_set_blocking($client, false); + $frames = runwireH2Frames(stream_get_contents($client)); + $data = ''; + $boundary = null; + $refused = null; + foreach ($frames as $frame) { + if ($frame->knownType() === FrameType::DATA) { + $data .= $frame->payload; + } + if ($frame->knownType() === FrameType::GOAWAY) { + $boundary = unpack('N', substr($frame->payload, 0, 4))[1]; + } + if ($frame->knownType() === FrameType::RST_STREAM && $frame->streamId === 3) { + $refused = unpack('N', $frame->payload)[1]; + } + } + + expect($data)->toBe('accepted-h2') + ->and($requests)->toBe(['/first']) + ->and($boundary)->toBe(1) + ->and($refused)->toBe(ErrorCode::REFUSED_STREAM->value); + } finally { + fclose($client); + } +}); diff --git a/tests/NativeWorkerLifecycleTest.php b/tests/NativeWorkerLifecycleTest.php index d0f5bc26..d9c46fb8 100644 --- a/tests/NativeWorkerLifecycleTest.php +++ b/tests/NativeWorkerLifecycleTest.php @@ -9,6 +9,7 @@ use Infocyph\Runwire\Http\HttpRequest; use Infocyph\Runwire\Http\Internal\BufferedRequestBody; use Infocyph\Runwire\Http\Internal\CallbackResponseWriter; +use Infocyph\Runwire\Http\ResponseWriterInterface; use Infocyph\Runwire\Loop\SelectLoop; use Infocyph\Runwire\Metrics\DiagnosticsPolicy; use Infocyph\Runwire\RequestContext; @@ -227,3 +228,82 @@ public function reset(RequestContext $context): void $worker->close(); } }); + + +it('finishes the first accepted HTTP request when recycling precedes header completion', function (string $prefix): void { + $loop = new SelectLoop(); + $worker = new WorkerContext('test', 0, 1, getmypid(), 0, recyclePolicy: new WorkerRecyclePolicy(gracefulTimeoutSeconds: 0.08)); + $events = []; + $server = \Infocyph\Runwire\Server::http('127.0.0.1:0', static function (HttpRequest $request, ResponseWriterInterface $writer) use (&$events): void { + $events[] = $request->target; + $writer->end('accepted-first'); + }); + $listener = \Infocyph\Runwire\Network\TcpListener::bind($server->address, $server->listener, $server->connection); + $handle = NativeHttpWorker::attach( + $loop, $worker, new \Infocyph\Runwire\Runtime\Internal\BoundServer($server, $listener), + RuntimeContext::standalone(), new \Infocyph\Runwire\Runtime\RequestExecutionPolicy(), + new ApplicationLifecycleHooks(drain: static function () use (&$events): void { $events[] = 'drain'; }), + ); + $client = stream_socket_client('tcp://' . $listener->address()); + if (!is_resource($client)) { + throw new RuntimeException('Unable to connect recycling race test client.'); + } + + try { + if ($prefix !== '') { + fwrite($client, $prefix); + } + $loop->delay(0.01, static fn() => $worker->requestRecycle()); + $loop->delay(0.02, static function () use ($client, $prefix): void { + $request = "GET /first HTTP/1.1\r\nHost: localhost\r\n\r\n"; + fwrite($client, substr($request, strlen($prefix)) . "GET /second HTTP/1.1\r\nHost: localhost\r\n\r\n"); + }); + $loop->delay(0.12, static fn() => $loop->stop()); + $loop->run(); + stream_set_blocking($client, false); + $response = stream_get_contents($client); + + expect($response)->toContain('accepted-first') + ->and(substr_count($response, 'HTTP/1.1 200'))->toBe(1) + ->and($events)->toBe(['/first', 'drain']) + ->and($worker->requestsTotal())->toBe(1) + ->and($handle->drained())->toBeTrue(); + } finally { + fclose($client); + $handle->close(); + $worker->close(); + } +})->with(['before first bytes' => '', 'partial headers' => "GET /first HTTP/1.1\r\nHost:"]); + +it('bounds accepted idle connections by the recycling grace period', function (): void { + $loop = new SelectLoop(); + $worker = new WorkerContext('test', 0, 1, getmypid(), 0, recyclePolicy: new WorkerRecyclePolicy(gracefulTimeoutSeconds: 0.02)); + $drains = 0; + $server = \Infocyph\Runwire\Server::http('127.0.0.1:0', static function (): void { + throw new RuntimeException('An idle connection must not dispatch a request.'); + }); + $listener = \Infocyph\Runwire\Network\TcpListener::bind($server->address, $server->listener, $server->connection); + $handle = NativeHttpWorker::attach( + $loop, $worker, new \Infocyph\Runwire\Runtime\Internal\BoundServer($server, $listener), + RuntimeContext::standalone(), new \Infocyph\Runwire\Runtime\RequestExecutionPolicy(), + new ApplicationLifecycleHooks(drain: static function () use (&$drains): void { ++$drains; }), + ); + $client = stream_socket_client('tcp://' . $listener->address()); + if (!is_resource($client)) { + throw new RuntimeException('Unable to connect idle recycling test client.'); + } + + try { + $loop->delay(0.01, static fn() => $worker->requestRecycle()); + $loop->delay(0.06, static function () use ($handle, &$drains): void { + expect($handle->drained())->toBeTrue()->and($drains)->toBe(1); + }); + $loop->delay(0.08, static fn() => $loop->stop()); + $loop->run(); + expect($worker->requestsTotal())->toBe(0); + } finally { + fclose($client); + $handle->close(); + $worker->close(); + } +}); diff --git a/tests/PhpQuicHttp3ConnectionTest.php b/tests/PhpQuicHttp3ConnectionTest.php index f7b7cc0d..b87b4f30 100644 --- a/tests/PhpQuicHttp3ConnectionTest.php +++ b/tests/PhpQuicHttp3ConnectionTest.php @@ -539,3 +539,67 @@ static function (): void {}, expect($requestRaw->readBytes)->toBeGreaterThan(0) ->and($connection->closed())->toBeFalse(); }); + + +it('keeps accepted HTTP3 headers admissible through application draining', function (): void { + $block = (new Encoder(0, 0))->encode([ + [':method', 'GET'], [':scheme', 'https'], [':authority', 'example.com'], [':path', '/accepted'], + ], 0)->block; + $frame = new Frame(FrameType::HEADERS->value, $block)->encode(); + $stream = fakeHttp3ConnectionStream(0, true, [substr($frame, 0, 1), '', substr($frame, 1), null]); + $raw = fakeHttp3ConnectionRaw([ + fakeHttp3ConnectionStream(3, false), fakeHttp3ConnectionStream(7, false), fakeHttp3ConnectionStream(11, false), + ], [$stream]); + $events = []; + $runtime = \Infocyph\Runwire\RuntimeContext::standalone(); + $application = new \Infocyph\Runwire\Runtime\Host\RuntimeApplication( + static function (HttpRequest $request, ResponseWriterInterface $writer) use (&$events): void { + $events[] = $request->target; + $writer->end('accepted-h3'); + }, + runtimeContext: $runtime, + lifecycle: new \Infocyph\Runwire\Runtime\ApplicationLifecycleHooks(drain: static function () use (&$events): void { $events[] = 'drain'; }), + ); + $connection = new PhpQuicHttp3Connection(new PhpQuicConnection($raw), $application->handle(...)); + $connection->pump(); + $listener = new \Infocyph\Runwire\Http\Http3\Quic\PhpQuicListener(new class { + public function accept(): ?object { return null; } + public function close(): void {} + public function setBlocking(bool $blocking): void { unset($blocking); } + }); + $masks = new PhpQuicEventMasks(1, 2, 4, 8, 16); + $poller = new \Infocyph\Runwire\Http\Http3\Quic\PhpQuicHttp3Poller($masks, static function (array $items) use ($masks): array { + $ready = []; + foreach ($items as $id => $item) { + $ready[$id] = $item[1] & ~$masks->error; + } + return $ready; + }); + $worker = new \Infocyph\Runwire\Http\Http3\Quic\PhpQuicHttp3Worker($listener, $application->handle(...), new Http3Limits(), 4, $poller); + (new ReflectionProperty($worker, 'connections'))->setValue($worker, [spl_object_id($raw) => $connection]); + $loop = new \Infocyph\Runwire\Loop\SelectLoop(); + $taskLoop = new \Infocyph\Runwire\Loop\SelectLoop(); + $context = new \Infocyph\Runwire\Supervisor\WorkerContext('h3', 0, 1, getmypid(), 0); + $context->attachLoop($taskLoop); + $sampler = new \Infocyph\Runwire\Runtime\Internal\WorkerDiagnosticsSampler($context, $runtime->metrics, new \Infocyph\Runwire\Metrics\DiagnosticsPolicy(), $taskLoop); + $attachment = new \Infocyph\Runwire\Runtime\Internal\NativeHttp3Attachment($loop, $taskLoop, $context, $application, $worker, $sampler, static function (): void {}); + $handle = $attachment->start(0.01); + + try { + expect($connection->activeRequestStreams())->toBe(1); + $handle->stop(); + $loop->tick(); + $rejected = fakeHttp3ConnectionStream(4, true, ['']); + $raw->acceptedStreams[] = $rejected; + $loop->delay(0.05, static fn() => $loop->stop()); + $loop->run(); + + expect($stream->written)->toContain('accepted-h3') + ->and($stream->ended)->toBeTrue() + ->and($events)->toBe(['/accepted', 'drain']) + ->and($rejected->peerResetCode)->toBe(ErrorCode::REQUEST_REJECTED->value); + } finally { + $handle->close(); + $context->close(); + } +}); From bcaec48aba19e41861002ed5a84edd243f1426ca Mon Sep 17 00:00:00 2001 From: "A. B. M. Mahmudul Hasan" Date: Sun, 4 Oct 2026 18:23:46 +0600 Subject: [PATCH 2/2] Update ControlServerTest.php --- tests/ControlServerTest.php | 18 ++++++++++++++++-- 1 file changed, 16 insertions(+), 2 deletions(-) diff --git a/tests/ControlServerTest.php b/tests/ControlServerTest.php index 5e91627f..89f26e5f 100644 --- a/tests/ControlServerTest.php +++ b/tests/ControlServerTest.php @@ -3,11 +3,13 @@ declare(strict_types=1); use Infocyph\Runwire\Control\ControlOptions; +use Infocyph\Runwire\Internal\MonotonicTime; +use Infocyph\Runwire\Supervisor\Internal\ChildReaper; use Infocyph\Runwire\Supervisor\Supervisor; use Infocyph\Runwire\Supervisor\WorkerContext; use Infocyph\Runwire\Supervisor\WorkerGroup; -it('serves bounded runtime control over a protected unix socket', function (): void { +it('serves bounded runtime control over a protected unix socket', function (int $publicationDelayMicroseconds): void { $socket = sys_get_temp_dir() . '/runwire-control-test-' . getmypid() . '-' . bin2hex(random_bytes(4)) . '.sock'; $resultFile = tempnam(sys_get_temp_dir(), 'runwire-control-result-'); if ($resultFile === false) { @@ -98,6 +100,9 @@ 'force' => false, ]); + if ($publicationDelayMicroseconds > 0) { + usleep($publicationDelayMicroseconds); + } file_put_contents($resultFile, json_encode( compact('mode', 'status', 'stale', 'stop'), JSON_THROW_ON_ERROR, @@ -132,6 +137,9 @@ )); $supervisor->run(); + if (!ChildReaper::waitForUntil($clientPid, MonotonicTime::deadlineAfterSeconds(MonotonicTime::nowNanoseconds(), 3.0))) { + throw new RuntimeException('Control client did not exit before the result deadline.'); + } $raw = file_get_contents($resultFile); if (!is_string($raw) || $raw === '') { throw new RuntimeException('Control client did not write a result.'); @@ -151,6 +159,12 @@ ->and(file_exists($socket))->toBeFalse() ->and($supervisor->status()->workers)->toBe([]); } finally { + if (!ChildReaper::waitForUntil($clientPid, MonotonicTime::nowNanoseconds())) { + posix_kill($clientPid, SIGKILL); + if (!ChildReaper::waitForUntil($clientPid, MonotonicTime::deadlineAfterSeconds(MonotonicTime::nowNanoseconds(), 1.0))) { + throw new RuntimeException('Control client did not exit after cleanup.'); + } + } if (file_exists($socket)) { unlink($socket); } @@ -158,7 +172,7 @@ unlink($resultFile); } } -}); +})->with(['immediate publication' => 0, 'delayed publication' => 50_000]); it('rejects unsafe control options', function (): void { expect(fn () => new ControlOptions('relative.sock'))->toThrow(InvalidArgumentException::class)