From 26b98087128c39ddebe2799fb9f536547596f25a Mon Sep 17 00:00:00 2001 From: Lennart Schoch Date: Wed, 7 Oct 2026 14:25:58 +0100 Subject: [PATCH] Wait for a stream slot at the peer's limit instead of failing A request that finds the shared HTTP/2 connection at the peer's max_concurrent_streams came back as {error, max_streams_exceeded}. It now waits for a stream to end and goes out then, for at most connect_timeout, the bound a pool checkout has, and runs out with the same checkout_timeout. A slot freed while a request body is being sent is handed out once the connection takes requests again. A GOAWAY, a close or the idle timer fail the waiting requests with that reason: nothing of theirs was sent. --- src/hackney_conn.erl | 132 ++++++++++++++++++---- test/hackney_h2_goaway_server.erl | 29 ++++- test/hackney_http2_stream_limit_tests.erl | 125 ++++++++++++++++++++ 3 files changed, 257 insertions(+), 29 deletions(-) create mode 100644 test/hackney_http2_stream_limit_tests.erl diff --git a/src/hackney_conn.erl b/src/hackney_conn.erl index 034445a8..2e5fc5f8 100644 --- a/src/hackney_conn.erl +++ b/src/hackney_conn.erl @@ -247,6 +247,9 @@ %% new requests are refused and the connection closes when the last %% stream ends. h2_goaway :: atom() | undefined, + %% Requests the h2 layer refused for want of a stream slot, waiting for + %% one to end: [{Ref, From, CallEvent}], oldest first. + h2_waiters = [] :: [{reference(), gen_statem:from(), term()}], %% Per-stream caller monitors of a shared connection: StreamId => ref h2_stream_monitors = #{} :: #{pos_integer() => reference()}, %% Current HTTP/2 stream ID for streaming body mode (body = stream) @@ -934,11 +937,15 @@ connected(enter, OldState, #conn_data{transport = Transport, socket = Socket, %% Transfer ownership back to pool and notify it auto_release_to_pool(Data) end, + %% A slot may have freed while a request body was being sent: the + %% requests waiting for one get their turn now that the connection is + %% back. A state enter call may not queue an event, so a zero timeout + %% stands in for it. case {Timeout, Data2#conn_data.protocol} of - {infinity, _} -> {keep_state, Data2}; + {infinity, _} -> {keep_state, Data2, h2_slot_enter_actions(Data2)}; %% HTTP/2 multiplexes: the connection idles when its last stream ends, %% not when it enters `connected' (#836). - {_, http2} -> {keep_state, Data2, h2_idle_actions(Data2)}; + {_, http2} -> {keep_state, Data2, h2_idle_actions(Data2) ++ h2_slot_enter_actions(Data2)}; _ -> {keep_state, Data2, [{state_timeout, Timeout, idle_timeout}]} end; @@ -1045,6 +1052,18 @@ connected({call, From}, get_state, #conn_data{h2_goaway = ErrorCode}) connected({call, From}, get_state, _Data) -> {keep_state_and_data, [{reply, From, {ok, connected}}]}; +%% A stream has ended, or the connection is back from sending a request body +%% (h2_stream_result/2 and connected(enter) queue this): the request that has +%% waited longest for a slot is made again, as the call it came in as. +connected(internal, h2_stream_slot, #conn_data{h2_waiters = [{Ref, From, Event} | Rest]} = Data) -> + {keep_state, Data#conn_data{h2_waiters = Rest}, + [{{timeout, {h2_stream_wait, Ref}}, cancel}, + {next_event, {call, From}, Event}]}; +connected(internal, h2_stream_slot, _Data) -> + keep_state_and_data; +connected({timeout, h2_stream_slot}, h2_stream_slot, _Data) -> + {keep_state_and_data, [{next_event, internal, h2_stream_slot}]}; + connected({call, From}, verify_socket, #conn_data{socket = undefined} = Data) -> %% Socket not connected {next_state, closed, Data, [{reply, From, {error, closed}}]}; @@ -1159,7 +1178,8 @@ connected({call, From}, checkin_info, Data) -> Map = (checkin_info_map(Data))#{ready => socket_ready(Data)}, {keep_state_and_data, [{reply, From, Map}]}; -connected({call, From}, {request, Method, Path, Headers, Body, ReqOpts}, #conn_data{protocol = http2} = Data) -> +connected({call, From}, {request, Method, Path, Headers, Body, ReqOpts} = Event, + #conn_data{protocol = http2} = Data) -> %% HTTP/2 request - use h2_machine (1xx not applicable for HTTP/2) %% Allow recv_timeout to be overridden per-request (fix for issue #832) RecvTimeout = proplists:get_value(recv_timeout, ReqOpts, Data#conn_data.recv_timeout), @@ -1167,7 +1187,8 @@ connected({call, From}, {request, Method, Path, Headers, Body, ReqOpts}, #conn_d %% not leak into later requests on a pooled connection. SendTimeout = proplists:get_value(send_timeout, ReqOpts, Data#conn_data.send_timeout), NewData = Data#conn_data{recv_timeout = RecvTimeout}, - do_h2_request(From, Method, Path, Headers, Body, SendTimeout, NewData); + h2_or_wait(do_h2_request(From, Method, Path, Headers, Body, SendTimeout, NewData), + From, Event, Data); connected({call, From}, {request, Method, Path, Headers, Body, ReqOpts}, #conn_data{protocol = http3} = Data) -> %% HTTP/3 request - use hackney_h3 (1xx not applicable for HTTP/3) @@ -1176,15 +1197,19 @@ connected({call, From}, {request, Method, Path, Headers, Body, ReqOpts}, #conn_d NewData = Data#conn_data{recv_timeout = RecvTimeout}, do_h3_request(From, Method, Path, Headers, Body, NewData); -connected({call, From}, {request_async, Method, Path, Headers, Body, AsyncMode, StreamTo}, #conn_data{protocol = http2} = Data) -> +connected({call, From}, {request_async, Method, Path, Headers, Body, AsyncMode, StreamTo} = Event, + #conn_data{protocol = http2} = Data) -> %% HTTP/2 async request - do_h2_request_async(From, Method, Path, Headers, Body, AsyncMode, StreamTo, false, - Data#conn_data.send_timeout, Data); + h2_or_wait(do_h2_request_async(From, Method, Path, Headers, Body, AsyncMode, StreamTo, false, + Data#conn_data.send_timeout, Data), + From, Event, Data); -connected({call, From}, {request_async, Method, Path, Headers, Body, AsyncMode, StreamTo, FollowRedirect}, #conn_data{protocol = http2} = Data) -> +connected({call, From}, {request_async, Method, Path, Headers, Body, AsyncMode, StreamTo, FollowRedirect} = Event, + #conn_data{protocol = http2} = Data) -> %% HTTP/2 async request with redirect option - do_h2_request_async(From, Method, Path, Headers, Body, AsyncMode, StreamTo, FollowRedirect, - Data#conn_data.send_timeout, Data); + h2_or_wait(do_h2_request_async(From, Method, Path, Headers, Body, AsyncMode, StreamTo, FollowRedirect, + Data#conn_data.send_timeout, Data), + From, Event, Data); connected({call, From}, {request_async, Method, Path, Headers, Body, AsyncMode, StreamTo}, #conn_data{protocol = http3} = Data) -> %% HTTP/3 async request @@ -1198,11 +1223,13 @@ connected({call, From}, {request_streaming, Method, Path, Headers, Body}, #conn_ %% HTTP/3 request with streaming body reads (returns headers, then stream_body for chunks) do_h3_request_streaming(From, Method, Path, Headers, Body, Data); -connected({call, From}, {request_streaming, Method, Path, Headers, Body}, #conn_data{protocol = http2} = Data) -> +connected({call, From}, {request_streaming, Method, Path, Headers, Body} = Event, + #conn_data{protocol = http2} = Data) -> %% HTTP/2: reply with status and headers, leave the body for %% stream_body/1 or body/1, like the HTTP/3 clause above. - do_h2_send(From, Method, Path, Headers, Body, {stream, waiting_headers, From}, - sync, Data#conn_data.send_timeout, Data); + h2_or_wait(do_h2_send(From, Method, Path, Headers, Body, {stream, waiting_headers, From}, + sync, Data#conn_data.send_timeout, Data), + From, Event, Data); connected({call, From}, {request_streaming, Method, Path, Headers, Body}, _Data) -> %% HTTP/1.1 leaves the body on the connection anyway, so this is the @@ -1246,13 +1273,16 @@ connected({call, From}, {request_async, Method, Path, Headers, Body, AsyncMode, %% Start a new async request with redirect option (HTTP/1.1) do_request_async(From, Method, Path, Headers, Body, AsyncMode, StreamTo, FollowRedirect, Data); -connected({call, From}, {request_async, Method, Path, Headers, Body, AsyncMode, StreamTo, FollowRedirect, ReqOpts}, #conn_data{protocol = http2} = Data) -> +connected({call, From}, {request_async, Method, Path, Headers, Body, AsyncMode, StreamTo, FollowRedirect, ReqOpts} = Event, + #conn_data{protocol = http2} = Data) -> %% HTTP/2 async request with ReqOpts (fix for issue #832) RecvTimeout = proplists:get_value(recv_timeout, ReqOpts, Data#conn_data.recv_timeout), %% send_timeout passed along, never stored (no leak across pooled requests) SendTimeout = proplists:get_value(send_timeout, ReqOpts, Data#conn_data.send_timeout), NewData = Data#conn_data{recv_timeout = RecvTimeout}, - do_h2_request_async(From, Method, Path, Headers, Body, AsyncMode, StreamTo, FollowRedirect, SendTimeout, NewData); + h2_or_wait(do_h2_request_async(From, Method, Path, Headers, Body, AsyncMode, StreamTo, FollowRedirect, + SendTimeout, NewData), + From, Event, Data); connected({call, From}, {request_async, Method, Path, Headers, Body, AsyncMode, StreamTo, _FollowRedirect, ReqOpts}, #conn_data{protocol = http3} = Data) -> %% HTTP/3 async request with ReqOpts (fix for issue #832, redirect not yet implemented for H3) @@ -1270,10 +1300,12 @@ connected({call, From}, {send_headers, Method, Path, Headers, _ReqOpts}, #conn_d %% HTTP/3 streaming body - send headers only via QUIC do_h3_send_headers(From, Method, Path, Headers, Data); -connected({call, From}, {send_headers, Method, Path, Headers, ReqOpts}, #conn_data{protocol = http2} = Data) -> +connected({call, From}, {send_headers, Method, Path, Headers, ReqOpts} = Event, + #conn_data{protocol = http2} = Data) -> %% HTTP/2 streaming body - send headers without END_STREAM, then accept body %% chunks via send_body_chunk/finish_send_body. Mirrors do_h3_send_headers/5. - do_h2_send_headers(From, Method, Path, Headers, ReqOpts, Data); + h2_or_wait(do_h2_send_headers(From, Method, Path, Headers, ReqOpts, Data), + From, Event, Data); connected({call, From}, {open_h2_stream, Method, Path, Headers, HandlerPid, Opts}, #conn_data{protocol = http2, h2_conn = H2Conn} = Data) -> @@ -2103,6 +2135,17 @@ handle_common(cast, stop, _State, Data) -> handle_common(info, {'DOWN', Ref, process, _Pid, _Reason}, State, Data) -> handle_h2_stream_owner_down(Ref, State, Data); +%% A request waited for a stream slot as long as a pool checkout may wait. +handle_common({timeout, {h2_stream_wait, Ref}}, Ref, _State, + #conn_data{h2_waiters = Waiters} = Data) -> + case lists:keytake(Ref, 1, Waiters) of + {value, {Ref, From, _Event}, Rest} -> + {keep_state, Data#conn_data{h2_waiters = Rest}, + [{reply, From, {error, checkout_timeout}}]}; + false -> + keep_state_and_data + end; + handle_common({timeout, h2_idle}, h2_idle, connected, #conn_data{h2_shared = true, h2_streams = Streams, h2_conn = H2Conn, h2_mon = H2Mon} = Data) @@ -2114,8 +2157,9 @@ handle_common({timeout, h2_idle}, h2_idle, connected, _ -> erlang:demonitor(H2Mon, [flush]) end, close_h2(H2Conn), - {next_state, closed, Data#conn_data{h2_conn = undefined, h2_mon = undefined, - socket = undefined}}; + {Replies, Data1} = fail_h2_waiters(closed, Data), + {next_state, closed, Data1#conn_data{h2_conn = undefined, h2_mon = undefined, + socket = undefined}, Replies}; handle_common({timeout, h2_idle}, h2_idle, _State, _Data) -> %% A stream opened since the timer was armed; the last one to finish @@ -2270,6 +2314,14 @@ handle_common(info, {'EXIT', _Pid, _Reason}, _State, _Data) -> keep_state_and_data; handle_common(info, _Msg, _State, _Data) -> + keep_state_and_data; + +%% A stream ended while a request body was being sent, or after the +%% connection closed. The waiting requests, if any are left, get their turn +%% when the connection is back in connected (connected(enter)). +handle_common(internal, h2_stream_slot, _State, _Data) -> + keep_state_and_data; +handle_common({timeout, h2_stream_slot}, h2_stream_slot, _State, _Data) -> keep_state_and_data. %% @private Reply to caller and stop @@ -3528,6 +3580,8 @@ do_h2_send(From, Method, Path, Headers, Body, StreamState, Mode, SendTimeout, Da sync -> {keep_state, NewData}; {async, Ref1, _, _} -> {keep_state, NewData, [{reply, From, {ok, Ref1}}]} end; + {error, max_streams_exceeded} -> + wait_for_stream; {error, Reason} -> {keep_state_and_data, [{reply, From, {error, Reason}}]} end. @@ -3566,6 +3620,8 @@ do_h2_send_headers(From, Method, Path, Headers, ReqOpts, Data) -> request_from = undefined }, {next_state, streaming_body, NewData, [{reply, From, ok}]}; + {error, max_streams_exceeded} -> + wait_for_stream; {error, Reason} -> {keep_state_and_data, [{reply, From, {error, Reason}}]} end. @@ -4020,7 +4076,34 @@ h2_stream_result(#conn_data{h2_goaway = ErrorCode, h2_streams = Streams} = Data, {next_state, closed, Data2, CloseReplies} = h2_close_on_goaway(ErrorCode, Data), {next_state, closed, Data2, Replies ++ CloseReplies}; h2_stream_result(Data, Replies) -> - {keep_state, Data, Replies ++ h2_idle_actions(Data)}. + {keep_state, Data, Replies ++ h2_idle_actions(Data) ++ h2_slot_actions(Data)}. + +%% @private A send that found the connection at the peer's +%% max_concurrent_streams answers `wait_for_stream' instead of replying: the +%% request is parked with its call event until a stream ends, for at most the +%% connection's connect_timeout, the bound a pool checkout has, with the same +%% answer when it runs out. Any other result stands. +h2_or_wait(wait_for_stream, From, Event, + #conn_data{h2_waiters = Waiters, connect_timeout = Timeout} = Data) -> + Ref = make_ref(), + {keep_state, Data#conn_data{h2_waiters = Waiters ++ [{Ref, From, Event}]}, + [{{timeout, {h2_stream_wait, Ref}}, Timeout, Ref}]}; +h2_or_wait(Result, _From, _Event, _Data) -> + Result. + +h2_slot_actions(#conn_data{h2_waiters = []}) -> []; +h2_slot_actions(_Data) -> [{next_event, internal, h2_stream_slot}]. + +h2_slot_enter_actions(#conn_data{h2_waiters = []}) -> []; +h2_slot_enter_actions(_Data) -> [{{timeout, h2_stream_slot}, 0, h2_stream_slot}]. + +%% @private Answer every waiting request with Err. Nothing of theirs was sent. +fail_h2_waiters(Err, #conn_data{h2_waiters = Waiters} = Data) -> + Actions = lists:flatmap(fun({Ref, From, _Event}) -> + [{{timeout, {h2_stream_wait, Ref}}, cancel}, + {reply, From, {error, Err}}] + end, Waiters), + {Actions, Data#conn_data{h2_waiters = []}}. %% @private Arm the idle timer of a shared HTTP/2 connection with no stream %% open. Re-arming replaces the previous timer, so it always counts from the @@ -4050,7 +4133,9 @@ h2_on_goaway(LastStreamId, ErrorCode, #conn_data{h2_streams = Streams} = Data) - {next_state, closed, Data2, Replies ++ CloseReplies}; false -> ok = leave_h2_pool(Data1), - {keep_state, Data1#conn_data{h2_goaway = ErrorCode}, Replies} + %% A waiting request would only ever be refused here now. + {WaiterReplies, Data2} = fail_h2_waiters({goaway, ErrorCode}, Data1), + {keep_state, Data2#conn_data{h2_goaway = ErrorCode}, Replies ++ WaiterReplies} end. %% @private Leave streaming_body once the upload stream is gone: back to @@ -4111,9 +4196,10 @@ h2_on_closed(Reason, Data) -> collect_h2_aborts(Err, #conn_data{h2_streams = Streams} = Data) -> Replies = h2_abort_replies(Err, Streams), + {WaiterReplies, Data0} = fail_h2_waiters(Err, Data), Data1 = clear_h2_stream_monitors( - Data#conn_data{h2_streams = #{}, request_from = undefined}), - {Replies, Data1}. + Data0#conn_data{h2_streams = #{}, request_from = undefined}), + {Replies ++ WaiterReplies, Data1}. %% @private Fail the refused streams and keep the rest of the connection %% running. A request body still being sent has no caller parked on it, so diff --git a/test/hackney_h2_goaway_server.erl b/test/hackney_h2_goaway_server.erl index e5c0de33..f1f98d0c 100644 --- a/test/hackney_h2_goaway_server.erl +++ b/test/hackney_h2_goaway_server.erl @@ -11,11 +11,17 @@ %%% Pick(FirstStreamId, SecondStreamId) as last_stream_id %%% first_data when a stream's first body chunk arrives, with that %%% stream as last_stream_id +%%% {first_stream, Ms} Ms after the first stream opens, accepting it %%% and shaped by options: %%% goaway => direct | two_step two_step first sends GOAWAY(2^31-1), the %%% graceful shutdown of RFC 9113 6.8 (default direct) -%%% answer => boolean() false answers nothing after the GOAWAY, not -%%% even the accepted streams (default true) +%%% answer => boolean() false makes the first connection answer +%%% nothing, not even streams a GOAWAY accepted +%%% (default true) +%%% answer_after => Ms delay before a complete request is answered +%%% (default 0) +%%% max_streams => N SETTINGS_MAX_CONCURRENT_STREAMS the server +%%% advertises (default unlimited) %%% notify => pid() gets {goaway_server, rst_stream, StreamId} for %%% each RST_STREAM on the first connection, then %%% {goaway_server, done} when that connection ends @@ -37,7 +43,8 @@ start(Trigger, Opts) -> {versions, ['tlsv1.2', 'tlsv1.3']}, {active, false}, {mode, binary}, {reuseaddr, true}]), {ok, {_, Port}} = ssl:sockname(LSock), - Options = maps:merge(#{goaway => direct, answer => true}, Opts), + Options = maps:merge(#{goaway => direct, answer => true, answer_after => 0, + max_streams => unlimited}, Opts), Pid = spawn(fun() -> accept_loop(LSock, Trigger, Options) end), Url = iolist_to_binary([<<"https://localhost:">>, integer_to_list(Port), <<"/">>]), {Pid, Url}. @@ -49,7 +56,7 @@ accept_loop(LSock, Trigger, Opts) -> case ssl:transport_accept(LSock, 2000) of {ok, TSock} -> spawn(fun() -> serve(TSock, Trigger, Opts) end), - accept_loop(LSock, none, maps:remove(notify, Opts)); + accept_loop(LSock, none, maps:remove(notify, Opts#{answer => true})); {error, timeout} -> accept_loop(LSock, Trigger, Opts); {error, closed} -> ok end. @@ -78,7 +85,7 @@ serve_conn(TSock, Trigger, Opts) -> {ok, Sock} -> case recv_preface(Sock, <<>>) of {ok, Rest} -> - send(Sock, h2_frame:settings([])), + send(Sock, h2_frame:settings(settings(Opts))), %% seen: stream ids in order of arrival. complete: streams %% whose request has fully arrived and is not answered yet. %% accepted: after the GOAWAY, the streams it let through. @@ -141,6 +148,9 @@ goaway_for(headers, _Sid, _EndStream, #{trigger := {second_stream, Pick}, seen : Pick(First, Second); goaway_for(data, Sid, false, #{trigger := first_data}) -> Sid; +goaway_for(headers, Sid, _EndStream, #{trigger := {first_stream, Ms}, seen := [Sid]}) -> + timer:sleep(Ms), + Sid; goaway_for(_Kind, _Sid, _EndStream, _St) -> none. @@ -166,12 +176,19 @@ send_goaway(Sock, LastStreamId, #{seen := Seen, goaway := Shape, answer := Answe %% flight when it does. answer_ready(_Sock, #{trigger := Trigger} = St) when Trigger =/= none -> St; -answer_ready(Sock, #{complete := Complete, accepted := Accepted} = St) -> +answer_ready(_Sock, #{answer := false} = St) -> + St; +answer_ready(Sock, #{complete := Complete, accepted := Accepted, + answer_after := Delay} = St) -> Ready = [S || S <- lists:reverse(Complete), Accepted =:= all orelse lists:member(S, Accepted)], + _ = [timer:sleep(Delay) || Ready =/= [], Delay > 0], lists:foldl(fun(S, Acc) -> respond(Sock, S, Acc) end, St#{complete := Complete -- Ready}, Ready). +settings(#{max_streams := unlimited}) -> []; +settings(#{max_streams := N}) -> [{max_concurrent_streams, N}]. + respond(Sock, Sid, #{enc := Enc} = St) -> {HBlock, Enc2} = h2_hpack:encode([{<<":status">>, <<"200">>}], Enc), send(Sock, h2_frame:headers(Sid, HBlock, false)), diff --git a/test/hackney_http2_stream_limit_tests.erl b/test/hackney_http2_stream_limit_tests.erl new file mode 100644 index 00000000..55ea5cfa --- /dev/null +++ b/test/hackney_http2_stream_limit_tests.erl @@ -0,0 +1,125 @@ +%%% A request that finds the shared HTTP/2 connection at the peer's +%%% SETTINGS_MAX_CONCURRENT_STREAMS waits for a stream to end instead of +%%% failing with max_streams_exceeded. +%%% +%%% The wait is bounded by connect_timeout and answers checkout_timeout when it +%%% runs out, the way a pool with every connection in use does. A GOAWAY while +%%% a request waits fails it with the goaway reason: nothing of it was sent. +%%% +%%% The server (hackney_h2_goaway_server) advertises max_streams and answers +%%% each request with its stream id as the body. +-module(hackney_http2_stream_limit_tests). + +-include_lib("eunit/include/eunit.hrl"). + +-define(POOL, stream_limit_test_pool). + +stream_limit_test_() -> + {foreach, fun setup/0, fun cleanup/1, + [{timeout, 30, fun waits_for_a_free_stream/0}, + {timeout, 30, fun gives_up_like_a_full_pool/0}, + {timeout, 30, fun goaway_fails_the_waiting_request/0}, + {timeout, 30, fun slot_freed_during_an_upload_is_used_after_it/0}]}. + +setup() -> + _ = application:ensure_all_started(hackney), + _ = application:ensure_all_started(h2), + stop_pool(), + ok = hackney_pool:start_pool(?POOL, [{max_connections, 5}]), + ok. + +cleanup(_) -> + stop_pool(). + +stop_pool() -> + try hackney_pool:stop_pool(?POOL) catch _:_ -> ok end, + ok. + +%% One stream at a time, each answered after 600ms. The second request has +%% to wait for the first to end, then goes out as stream 3. +waits_for_a_free_stream() -> + {Server, Url} = hackney_h2_goaway_server:start(none, #{max_streams => 1, answer_after => 600}), + try + [R1, R2] = concurrent_fetches(Url, [0, 300], []), + ?assertEqual({ok, 200, <<"1">>}, R1), + ?assertEqual({ok, 200, <<"3">>}, R2) + after + hackney_h2_goaway_server:stop(Server) + end. + +%% The first stream is never answered, so the slot never frees. The waiting +%% request gives up after connect_timeout with the answer a full pool gives. +gives_up_like_a_full_pool() -> + {Server, Url} = hackney_h2_goaway_server:start(none, #{max_streams => 1, answer => false}), + try + [R1, R2] = concurrent_fetches(Url, [0, 300], [{connect_timeout, 300}, {recv_timeout, 1000}]), + ?assertEqual({error, timeout}, R1), + ?assertEqual({error, checkout_timeout}, R2) + after + hackney_h2_goaway_server:stop(Server) + end. + +%% The server sends GOAWAY covering stream 1 while a second request waits for +%% its slot. The waiting request is refused at once with the goaway reason, +%% since nothing of it was sent, while stream 1 finishes. +goaway_fails_the_waiting_request() -> + {Server, Url} = hackney_h2_goaway_server:start({first_stream, 600}, #{max_streams => 1}), + try + [R1, R2] = concurrent_fetches(Url, [0, 300], []), + ?assertEqual({ok, 200, <<"1">>}, R1), + ?assertEqual({error, {goaway, no_error}}, R2) + after + hackney_h2_goaway_server:stop(Server) + end. + +%% Two streams at a time. With both busy, a request with a streamed body and +%% then a plain request queue up. The first slot goes to the body: while it is +%% being sent, the second stream ends, and that slot goes to the plain request +%% once the body is done and the connection takes requests again. +slot_freed_during_an_upload_is_used_after_it() -> + {Server, Url} = hackney_h2_goaway_server:start(none, #{max_streams => 2, answer_after => 600}), + try + Self = self(), + %% The first request registers the shared connection before the second + %% checks one out, so both streams share it. + Plain = [begin + timer:sleep(Delay), + spawn_link(fun() -> Self ! {self(), fetch(Url, [])} end) + end || Delay <- [0, 300]], + timer:sleep(100), + %% Queues behind the streamed request below, while that one still waits. + Waiter = spawn_link(fun() -> timer:sleep(50), Self ! {self(), fetch(Url, [])} end), + {ok, Conn} = hackney:request(post, Url, [], stream, opts([])), + %% The second plain stream ends while the body is still being sent. + timer:sleep(1000), + ok = hackney:finish_send_body(Conn), + {ok, 200, _Headers, Conn} = hackney:start_response(Conn), + ?assertEqual({ok, <<"5">>}, hackney:body(Conn)), + ?assertEqual([{ok, 200, <<"1">>}, {ok, 200, <<"3">>}], [await(P) || P <- Plain]), + ?assertEqual({ok, 200, <<"7">>}, await(Waiter)) + after + hackney_h2_goaway_server:stop(Server) + end. + +%% The first request registers the shared connection before the next checks +%% one out, so they share it: 300ms covers the TLS handshake on localhost. +concurrent_fetches(Url, Delays, Extra) -> + Self = self(), + Pids = [begin + timer:sleep(Delay), + spawn_link(fun() -> Self ! {self(), fetch(Url, Extra)} end) + end || Delay <- Delays], + [await(P) || P <- Pids]. + +await(Pid) -> + receive {Pid, R} -> R after 10000 -> {error, test_timeout} end. + +fetch(Url, Extra) -> + case hackney:request(get, Url, [], <<>>, opts(Extra)) of + {ok, S, _H, B} when is_binary(B) -> {ok, S, B}; + {error, E} -> {error, E} + end. + +opts(Extra) -> + Extra ++ [{pool, ?POOL}, {protocols, [http2]}, {recv_timeout, 5000}, + {ssl_options, [{insecure, true}, {verify, verify_none}]}].