From cc2946fa0b5461cbcb93b58680c43dce44c9ce73 Mon Sep 17 00:00:00 2001 From: Lennart Schoch Date: Thu, 1 Oct 2026 18:14:47 +0100 Subject: [PATCH] Let streams accepted before a GOAWAY finish RFC 9113 6.8 only refuses the streams above last_stream_id; the peer may still process the ones up to it and keeps the connection open to finish them. h2_on_goaway failed every in-flight stream regardless, so a request the server went on to complete came back as {error, {goaway, _}}. Fail only the refused streams, leave the pool so no new request lands here, refuse one that was already checked out, and close once the accepted streams end. A streamed request or response body in flight still closes at once. --- src/hackney_conn.erl | 91 +++++++-- test/hackney_http2_goaway_drain_tests.erl | 228 ++++++++++++++++++++++ 2 files changed, 305 insertions(+), 14 deletions(-) create mode 100644 test/hackney_http2_goaway_drain_tests.erl diff --git a/src/hackney_conn.erl b/src/hackney_conn.erl index b35fb562..58611889 100644 --- a/src/hackney_conn.erl +++ b/src/hackney_conn.erl @@ -238,6 +238,10 @@ %% Shared through the pool (share_h2/1): no owner, each stream monitors %% its caller, and the connection closes itself once idle. h2_shared = false :: boolean(), + %% Error code of a GOAWAY that left accepted streams to finish. While set, + %% new requests are refused and the connection closes when the last + %% stream ends. + h2_goaway :: atom() | undefined, %% 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) @@ -1021,6 +1025,11 @@ connected(state_timeout, idle_timeout, Data) -> %% Idle timeout - close connection {next_state, closed, Data}; +connected({call, From}, get_state, #conn_data{h2_goaway = ErrorCode}) + when ErrorCode =/= undefined -> + %% Draining after a GOAWAY: still serving its accepted streams, but the + %% pool and request/5 must not pick it, as they would a closed one. + {keep_state_and_data, [{reply, From, {ok, draining}}]}; connected({call, From}, get_state, _Data) -> {keep_state_and_data, [{reply, From, {ok, connected}}]}; @@ -1973,11 +1982,13 @@ closed(enter, _OldState, #conn_data{socket = Socket, transport = Transport, pool %% surfaces as `exit:{normal, _}` in the caller (issue #836). Stay %% alive briefly so those late calls get a proper `{error, closed}` %% reply via handle_common's closed-state fallback, then stop. + %% A GOAWAY drain ends here, whichever way the connection closed. + Data2 = Data#conn_data{socket = undefined, h2_goaway = undefined}, case PoolPid of undefined -> - {keep_state, Data#conn_data{socket = undefined}}; + {keep_state, Data2}; _ -> - {keep_state, Data#conn_data{socket = undefined}, + {keep_state, Data2, [{state_timeout, ?CLOSED_GRACE_MS, closed_grace_expired}]} end; @@ -3205,7 +3216,8 @@ start_h2_connection(Socket, Data, From, Origin) -> h2_conn = H2Conn, h2_mon = Mon, h2_streams = #{}, - h2_stream_monitors = #{} + h2_stream_monitors = #{}, + h2_goaway = undefined }, %% Cancel any pending idle_timeout armed by the %% TCP-first connected(enter): HTTP/2 connections @@ -3387,6 +3399,11 @@ do_h2_request_async(From, Method, Path, Headers, Body, AsyncMode, StreamTo, do_h2_send(From, Method, Path, Headers, Body, StreamState, {async, Ref, StreamTo, AsyncMode}, SendTimeout, Data). +do_h2_send(From, _Method, _Path, _Headers, _Body, _StreamState, _Mode, _SendTimeout, + #conn_data{h2_goaway = ErrorCode}) when ErrorCode =/= undefined -> + %% Checked out before the pool dropped this draining connection. The peer + %% would ignore a new stream, so refuse it here; nothing was sent. + {keep_state_and_data, [{reply, From, {error, {goaway, ErrorCode}}}]}; do_h2_send(From, Method, Path, Headers, Body, StreamState, Mode, SendTimeout, Data) -> #conn_data{h2_conn = H2Conn} = Data, {MethodBin, PathBin, H2Headers} = @@ -3471,6 +3488,9 @@ do_h2_send(From, Method, Path, Headers, Body, StreamState, Mode, SendTimeout, Da %% @private Begin an HTTP/2 streaming-body request: send HEADERS without %% END_STREAM and transition to streaming_body so the caller can push body %% chunks via send_body_chunk/finish_send_body. Mirrors do_h3_send_headers/5. +do_h2_send_headers(From, _Method, _Path, _Headers, _ReqOpts, + #conn_data{h2_goaway = ErrorCode}) when ErrorCode =/= undefined -> + {keep_state_and_data, [{reply, From, {error, {goaway, ErrorCode}}}]}; do_h2_send_headers(From, Method, Path, Headers, ReqOpts, Data) -> #conn_data{h2_conn = H2Conn} = Data, {MethodBin, PathBin, H2Headers} = @@ -3694,8 +3714,8 @@ handle_h2_event({trailers, StreamId, _Headers}, Data) -> h2_on_data(StreamId, <<>>, true, Data); handle_h2_event({stream_reset, StreamId, ErrorCode}, Data) -> h2_on_stream_reset(StreamId, ErrorCode, Data); -handle_h2_event({goaway, _LastStreamId, ErrorCode}, Data) -> - h2_on_goaway(ErrorCode, Data); +handle_h2_event({goaway, LastStreamId, ErrorCode}, Data) -> + h2_on_goaway(LastStreamId, ErrorCode, Data); handle_h2_event({closed, Reason}, Data) -> h2_on_closed(Reason, Data); handle_h2_event(_Other, Data) -> @@ -3945,7 +3965,13 @@ h2_stream_parked_from({stream, body_full, _, _, _, From}) -> From; h2_stream_parked_from(_) -> undefined. %% @private Result for a handler that just removed an HTTP/2 stream: a shared -%% connection left with no stream arms its idle timer. +%% connection left with no stream arms its idle timer, and one draining after +%% a GOAWAY closes once its last accepted stream has ended. +h2_stream_result(#conn_data{h2_goaway = ErrorCode, h2_streams = Streams} = Data, + Replies) + when ErrorCode =/= undefined, map_size(Streams) =:= 0 -> + {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)}. @@ -3959,7 +3985,37 @@ h2_idle_actions(#conn_data{h2_shared = true, h2_streams = Streams, h2_idle_actions(_Data) -> []. -h2_on_goaway(ErrorCode, #conn_data{h2_conn = H2Conn, h2_mon = H2Mon} = Data) -> +%% RFC 9113 6.8: a GOAWAY refuses the streams above LastStreamId, but the peer +%% may still process the ones up to it and keeps the connection open to finish +%% them. Fail only the refused streams, leave the pool so no new request lands +%% here, and close once the accepted streams end. Aborting those too reported +%% requests the server went on to complete as failed. +%% +%% A streamed request or response body (a `stream' entry) drives the +%% connection through states of its own, so with one in flight the connection +%% still closes at once. +h2_on_goaway(LastStreamId, ErrorCode, #conn_data{h2_streams = Streams} = Data) -> + {Accepted, Refused} = lists:partition(fun(SId) -> SId =< LastStreamId end, + maps:keys(Streams)), + case Accepted =/= [] andalso not h2_streaming_in_flight(Streams) of + true -> + {Replies, Data1} = abort_h2_streams(Refused, {goaway, ErrorCode}, Data), + ok = leave_h2_pool(Data1), + {keep_state, Data1#conn_data{h2_goaway = ErrorCode}, Replies}; + false -> + h2_close_on_goaway(ErrorCode, Data) + end. + +h2_streaming_in_flight(Streams) -> + lists:any(fun({_Owner, Inner}) -> element(1, Inner) =:= stream end, + maps:values(Streams)). + +leave_h2_pool(#conn_data{pool_pid = PoolPid}) when is_pid(PoolPid) -> + gen_server:cast(PoolPid, {unregister_h2, self()}); +leave_h2_pool(_Data) -> + ok. + +h2_close_on_goaway(ErrorCode, #conn_data{h2_conn = H2Conn, h2_mon = H2Mon} = Data) -> %% A GOAWAY means the peer will not service new streams on this connection. %% AWS ALBs recycle connections this way, sending GOAWAY but keeping the %% socket open for a drain window. Leaving the conn `connected` and pooled @@ -3967,8 +4023,7 @@ h2_on_goaway(ErrorCode, #conn_data{h2_conn = H2Conn, h2_mon = H2Mon} = Data) -> %% request opened a stream past last_stream_id that the peer ignored and hung %% to recv_timeout. Tear the connection down and transition to `closed` (like %% h2_on_closed/2): the pool then stops reusing it (h2_conn_usable requires - %% `connected`) and new requests dial a fresh connection. in-flight streams - %% are aborted with the goaway error as before. + %% `connected`) and new requests dial a fresh connection. {Replies, Data1} = collect_h2_aborts({goaway, ErrorCode}, Data), Data2 = cancel_all_h2_timers(Data1), _ = case H2Mon of @@ -3991,7 +4046,18 @@ h2_on_closed(Reason, Data) -> {next_state, closed, Stripped, Replies}. collect_h2_aborts(Err, #conn_data{h2_streams = Streams} = Data) -> - Replies = maps:fold(fun + Replies = h2_abort_replies(Err, Streams), + Data1 = clear_h2_stream_monitors( + Data#conn_data{h2_streams = #{}, request_from = undefined}), + {Replies, Data1}. + +%% @private Fail some streams and keep the rest of the connection running. +abort_h2_streams(StreamIds, Err, #conn_data{h2_streams = Streams} = Data) -> + Replies = h2_abort_replies(Err, maps:with(StreamIds, Streams)), + {Replies, lists:foldl(fun drop_h2_stream/2, Data, StreamIds)}. + +h2_abort_replies(Err, Streams) -> + maps:fold(fun (_SId, {From, {sync, _}}, Acc) -> [{reply, From, {error, Err}} | Acc]; (_SId, {From, {sync, body, _, _, _}}, Acc) -> @@ -4011,10 +4077,7 @@ collect_h2_aborts(Err, #conn_data{h2_streams = Streams} = Data) -> From -> [{reply, From, {error, Err}} | Acc] end; (_, _, Acc) -> Acc - end, [], Streams), - Data1 = clear_h2_stream_monitors( - Data#conn_data{h2_streams = #{}, request_from = undefined}), - {Replies, Data1}. + end, [], Streams). %%==================================================================== diff --git a/test/hackney_http2_goaway_drain_tests.erl b/test/hackney_http2_goaway_drain_tests.erl new file mode 100644 index 00000000..34a4b3d4 --- /dev/null +++ b/test/hackney_http2_goaway_drain_tests.erl @@ -0,0 +1,228 @@ +%%% GOAWAY must only fail the streams the peer did not accept. +%%% +%%% RFC 9113 6.8: streams up to and including the GOAWAY's last_stream_id may +%%% still be processed, and the peer keeps the connection open to finish them. +%%% hackney used to abort every in-flight stream with {error, {goaway, _}}, so a +%%% request the server went on to complete came back as an error: for a payment +%%% API that is a charge that succeeded and was reported as failed. +%%% +%%% The server below holds the first two streams of its first connection, sends +%%% GOAWAY with a chosen last_stream_id, then answers only the streams at or +%%% below it. Later connections answer immediately. +-module(hackney_http2_goaway_drain_tests). + +-include_lib("eunit/include/eunit.hrl"). + +-define(PREFACE, <<"PRI * HTTP/2.0\r\n\r\nSM\r\n\r\n">>). +-define(POOL, goaway_drain_test_pool). + +goaway_drain_test_() -> + {foreach, fun setup/0, fun cleanup/1, + [{timeout, 30, fun accepted_streams_complete/0}, + {timeout, 30, fun unaccepted_stream_fails_fast/0}, + {timeout, 30, fun two_step_shutdown/0}, + {timeout, 30, fun stalled_drain_ends_with_the_stream/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. + +%% GOAWAY(last_stream_id = 3) after streams 1 and 3: both were accepted, so both +%% complete. A request made while they drain must not land on the draining +%% connection, where the peer would ignore it, but dial a fresh one. +accepted_streams_complete() -> + {Server, Url} = start_server(fun(_First, Second) -> Second end), + try + [R1, R2, R3] = concurrent_requests(Url, 3), + ?assertEqual({ok, 200, <<"1">>}, R1), + ?assertEqual({ok, 200, <<"3">>}, R2), + ?assertEqual({ok, 200, <<"1">>}, R3) + after + stop_server(Server) + end. + +%% GOAWAY(last_stream_id = 1) after streams 1 and 3: stream 3 was not accepted +%% and fails at once with the goaway reason, stream 1 still completes. +unaccepted_stream_fails_fast() -> + {Server, Url} = start_server(fun(First, _Second) -> First end), + try + [R1, R2] = concurrent_requests(Url, 2), + ?assertEqual({ok, 200, <<"1">>}, R1), + ?assertEqual({error, {goaway, no_error}}, R2) + after + stop_server(Server) + end. + +%% The graceful shutdown of RFC 9113 6.8: GOAWAY(2^31-1) stops new streams +%% without refusing any, then a second GOAWAY gives the real last_stream_id. +%% The first frame alone must not fail anything. +two_step_shutdown() -> + {Server, Url} = start_server(fun(First, _Second) -> {two_step, First} end), + try + [R1, R2] = concurrent_requests(Url, 2), + ?assertEqual({ok, 200, <<"1">>}, R1), + ?assertEqual({error, {goaway, no_error}}, R2) + after + stop_server(Server) + end. + +%% A server that accepts a stream and then never answers it must not keep the +%% draining connection around: the stream's own recv_timeout ends it, and with +%% it the connection, so the next request gets a fresh one. +stalled_drain_ends_with_the_stream() -> + {Server, Url} = start_server(fun(First, _Second) -> {never, First} end), + try + [R1, R2] = concurrent_requests(Url, 2, [{recv_timeout, 1000}]), + ?assertEqual({error, timeout}, R1), + ?assertEqual({error, {goaway, no_error}}, R2), + ?assertEqual({ok, 200, <<"1">>}, fetch(Url, [])) + after + stop_server(Server) + end. + +%% The first request registers the shared connection before the second checks +%% one out, so both streams share it. The server sends GOAWAY on the second and +%% waits 200ms before answering, so a third lands inside that drain. +concurrent_requests(Url, N) -> + concurrent_requests(Url, N, []). + +concurrent_requests(Url, N, Extra) -> + Self = self(), + Pids = [begin + timer:sleep(Delay), + spawn_link(fun() -> Self ! {self(), fetch(Url, Extra)} end) + end || Delay <- lists:sublist([0, 300, 100], N)], + [receive {P, R} -> R after 10000 -> {error, test_timeout} end || P <- Pids]. + +fetch(Url) -> + fetch(Url, []). + +fetch(Url, Extra) -> + Opts = Extra ++ [{pool, ?POOL}, {protocols, [http2]}, {recv_timeout, 5000}, + {ssl_options, [{insecure, true}, {verify, verify_none}]}], + case hackney:request(get, Url, [], <<>>, Opts) of + {ok, S, _H, B} when is_binary(B) -> {ok, S, B}; + {error, E} -> {error, E} + end. + +%%==================================================================== +%% Frame-level h2 server. The body of each response is its stream id. +%%==================================================================== + +start_server(PickLastStreamId) -> + Certs = cert_dir(), + {ok, LSock} = ssl:listen(0, + [{certfile, filename:join(Certs, "server.pem")}, + {keyfile, filename:join(Certs, "server.key")}, + {alpn_preferred_protocols, [<<"h2">>]}, + {versions, ['tlsv1.2', 'tlsv1.3']}, + {active, false}, {mode, binary}, {reuseaddr, true}]), + {ok, {_, Port}} = ssl:sockname(LSock), + Pid = spawn(fun() -> accept_loop(LSock, {hold, PickLastStreamId}) end), + Url = iolist_to_binary([<<"https://localhost:">>, integer_to_list(Port), <<"/">>]), + {Pid, Url}. + +stop_server(Pid) -> + exit(Pid, kill). + +accept_loop(LSock, Mode) -> + case ssl:transport_accept(LSock, 2000) of + {ok, TSock} -> + spawn(fun() -> serve(TSock, Mode) end), + accept_loop(LSock, immediate); + {error, timeout} -> accept_loop(LSock, Mode); + {error, closed} -> ok + end. + +serve(TSock, Mode) -> + case ssl:handshake(TSock, 5000) of + {ok, Sock} -> + case recv_preface(Sock, <<>>) of + {ok, Rest} -> + send(Sock, h2_frame:settings([])), + loop(Sock, Rest, #{enc => h2_hpack:new_context(), mode => Mode, + held => []}); + _ -> ok + end; + _ -> ok + end. + +recv_preface(_Sock, Acc) when byte_size(Acc) >= 24 -> + <> = Acc, + case Pre of ?PREFACE -> {ok, Rest}; _ -> {error, bad_preface} end; +recv_preface(Sock, Acc) -> + case ssl:recv(Sock, 0, 5000) of + {ok, Data} -> recv_preface(Sock, <>); + {error, R} -> {error, R} + end. + +loop(Sock, Buf, St) -> + case h2_frame:decode(Buf) of + {ok, Frame, Rest} -> + case handle(Sock, Frame, St) of + {continue, St2} -> loop(Sock, Rest, St2); + stop -> ok + end; + {more, _} -> + case ssl:recv(Sock, 0, 30000) of + {ok, Data} -> loop(Sock, <>, St); + {error, _} -> ok + end; + {error, _, Rest} -> loop(Sock, Rest, St); + {error, _} -> ok + end. + +handle(Sock, {settings, _}, St) -> send(Sock, h2_frame:settings_ack()), {continue, St}; +handle(Sock, {ping, D}, St) -> send(Sock, h2_frame:ping_ack(D)), {continue, St}; +handle(_Sock, {goaway, _, _, _}, _St) -> stop; +handle(Sock, {headers, Sid, _B, _E, _H}, #{mode := immediate} = St) -> + {continue, respond(Sock, Sid, St)}; +handle(Sock, {headers, Sid, _B, _E, _H}, #{mode := {hold, Pick}, held := Held} = St) -> + case Held ++ [Sid] of + [First, Second] -> + LastStreamId = send_goaway(Sock, Pick(First, Second)), + %% The drain: the peer is told, then the accepted streams finish. + timer:sleep(200), + St2 = lists:foldl(fun(S, Acc) -> respond(Sock, S, Acc) end, + St#{held := []}, + [S || S <- [First, Second], S =< LastStreamId]), + {continue, St2#{mode := draining}}; + Held2 -> + {continue, St#{held := Held2}} + end; +handle(_Sock, _Other, St) -> {continue, St}. + +send_goaway(Sock, {never, LastStreamId}) -> + send(Sock, h2_frame:goaway(LastStreamId, no_error, <<>>)), + %% Nothing at or below LastStreamId is answered. + 0; +send_goaway(Sock, {two_step, LastStreamId}) -> + send(Sock, h2_frame:goaway(16#7fffffff, no_error, <<>>)), + timer:sleep(50), + send_goaway(Sock, LastStreamId); +send_goaway(Sock, LastStreamId) -> + send(Sock, h2_frame:goaway(LastStreamId, no_error, <<>>)), + LastStreamId. + +respond(Sock, Sid, #{enc := Enc} = St) -> + {HBlock, Enc2} = h2_hpack:encode([{<<":status">>, <<"200">>}], Enc), + send(Sock, h2_frame:headers(Sid, HBlock, false)), + send(Sock, h2_frame:data(Sid, integer_to_binary(Sid), true)), + St#{enc := Enc2}. + +send(Sock, FrameData) -> ssl:send(Sock, h2_frame:encode(FrameData)). + +cert_dir() -> + BeamDir = filename:dirname(code:which(?MODULE)), + Root = filename:join([BeamDir, "..", "..", "..", "..", ".."]), + filename:join([filename:absname(Root), "test", "certs"]).