Skip to content
Open
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
132 changes: 109 additions & 23 deletions src/hackney_conn.erl
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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;

Expand Down Expand Up @@ -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}}]};
Expand Down Expand Up @@ -1159,15 +1178,17 @@ 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),
%% send_timeout is passed along, never stored: a per-request override must
%% 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)
Expand All @@ -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
Expand All @@ -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
Expand Down Expand Up @@ -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)
Expand All @@ -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) ->
Expand Down Expand Up @@ -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)
Expand All @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
29 changes: 23 additions & 6 deletions test/hackney_h2_goaway_server.erl
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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}.
Expand All @@ -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.
Expand Down Expand Up @@ -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.
Expand Down Expand Up @@ -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.

Expand All @@ -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)),
Expand Down
Loading
Loading