Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
91 changes: 77 additions & 14 deletions src/hackney_conn.erl
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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}}]};

Expand Down Expand Up @@ -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;

Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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} =
Expand Down Expand Up @@ -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} =
Expand Down Expand Up @@ -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) ->
Expand Down Expand Up @@ -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)}.

Expand All @@ -3959,16 +3985,45 @@ 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
%% made checkout_h2/h2_conn_usable keep handing it out, so every reused
%% 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
Expand All @@ -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) ->
Expand All @@ -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).


%%====================================================================
Expand Down
Loading
Loading