Skip to content

Fix race in HTTP/2 buffered data writes - #396

Open
kafkiansky wants to merge 1 commit into
amphp:3.xfrom
kafkiansky:fix-http2-buffered-data-race
Open

kafkiansky wants to merge 1 commit into
amphp:3.xfrom
kafkiansky:fix-http2-buffered-data-race

Conversation

@kafkiansky

@kafkiansky kafkiansky commented Oct 1, 2026 •

Copy link
Copy Markdown
Contributor

writeBufferedData() updated the stream buffer and flow-control windows only after writeFrame(), which may suspend. A concurrent sendBufferedData() (deferred on every WINDOW_UPDATE) could then resend data that was already partly written and complete the shared deferredFuture early. Trailers with END_STREAM then went out before the body was finished. Peers received duplicated, reordered or truncated bodies, or the stream stalled.

The buffer and windows are now updated before writing, and all DATA frames of one call go out in a single write(), so frames from concurrent writers cannot interleave.

The new test testConcurrentWindowUpdatesWithSuspendedWrites makes every write suspend, sends connection and stream WINDOW_UPDATE frames together, and checks that the body arrives intact with END_STREAM on the trailers. It fails without the fix.

This overlaps with #387, which addresses the same race, but #387 fixes only part of it:

This PR updates the state before writing in both branches and sends all DATA frames of a call in a single write(). The test here makes every write suspend and sends connection and stream WINDOW_UPDATE frames together, so it covers the race end to end, not just the buffer invariant.

How to reproduce:

<?php

declare(strict_types=1);

require __DIR__ . '/vendor/autoload.php';

use Amp\Http\Client\Connection\DefaultConnectionFactory;
use Amp\Http\Client\Connection\UnlimitedConnectionPool;
use Amp\Http\Client\HttpClientBuilder;
use Amp\Http\Client\Request as ClientRequest;
use Amp\Http\Server\DefaultErrorHandler;
use Amp\Http\Server\Request;
use Amp\Http\Server\RequestHandler\ClosureRequestHandler;
use Amp\Http\Server\Response;
use Amp\Http\Server\SocketHttpServer;
use Amp\Socket\BindContext;
use Amp\Socket\Certificate;
use Amp\Socket\ClientTlsContext;
use Amp\Socket\ConnectContext;
use Amp\Socket\ServerTlsContext;
use Amp\TimeoutCancellation;
use Composer\InstalledVersions;
use Psr\Log\NullLogger;

const ADDRESS = '127.0.0.1:50443';

function body(int $size): string
{
    $body = '';

    for ($i = 0; \strlen($body) < $size; ++$i) {
        $body .= hash('sha256', (string) $i, true);
    }

    return substr($body, 0, $size);
}

$server = SocketHttpServer::createForDirectAccess(new NullLogger(), enableCompression: false);
$server->expose(ADDRESS, new BindContext()->withTlsContext(
    new ServerTlsContext()->withDefaultCertificate(new Certificate(
        __DIR__ . '/tests/Stub/tls/server.crt',
        __DIR__ . '/tests/Stub/tls/server.key',
    )),
));
$server->start(
    new ClosureRequestHandler(static function (Request $request): Response {
        parse_str($request->getUri()->getQuery(), $query);

        return new Response(body: body((int) ($query['n'] ?? 0)));
    }),
    new DefaultErrorHandler(),
);

echo 'amphp/http-server ', InstalledVersions::getPrettyVersion('amphp/http-server'), ' listening on https://', ADDRESS, "\n";

if (($argv[1] ?? '') === 'serve') {
    Amp\trapSignal([\SIGINT, \SIGTERM]);

    exit(0);
}

$client = new HttpClientBuilder()
    ->usingPool(new UnlimitedConnectionPool(new DefaultConnectionFactory(
        connectContext: new ConnectContext()->withTlsContext(new ClientTlsContext('')->withoutPeerVerification()),
    )))
    ->build();

$failed = false;

foreach ([9_000, 70_000, 300_000, 1_000_000, 5_000_000] as $size) {
    try {
        $request = new ClientRequest('https://' . ADDRESS . "/?n={$size}");
        $request->setProtocolVersions(['2']);

        $cancellation = new TimeoutCancellation(10);
        $received = $client->request($request, $cancellation)->getBody()->buffer($cancellation);
        $result = $received === body($size) ? 'ok' : 'BAD, received ' . \strlen($received) . ' bytes';
    } catch (\Throwable $e) {
        $result = 'BAD, ' . $e->getMessage();
    }

    $failed = $failed || $result !== 'ok';
    printf("%9d bytes: %s\n", $size, $result);
}

// The server is not stopped gracefully: after a corrupted response it may wait for the broken stream forever.
exit((int) $failed);

This branch has not been deployed

No deployments
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Development

Successfully merging this pull request may close these issues.

1 participant