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
88 changes: 41 additions & 47 deletions src/hackney_conn.erl
Original file line number Diff line number Diff line change
Expand Up @@ -1034,17 +1034,6 @@ connected({call, From}, get_state, #conn_data{h2_goaway = ErrorCode})
connected({call, From}, get_state, _Data) ->
{keep_state_and_data, [{reply, From, {ok, connected}}]};

%% The streaming-body calls only make sense in streaming_body, except for a
%% request a GOAWAY refused while its caller was between them: the connection
%% came back here then (h2_on_goaway/3), and the caller's next call is the
%% first chance to tell it.
connected({call, From}, {send_body_chunk, _}, #conn_data{protocol = http2} = Data) ->
h2_answer_refused_request(From, Data);
connected({call, From}, finish_send_body, #conn_data{protocol = http2} = Data) ->
h2_answer_refused_request(From, Data);
connected({call, From}, start_response, #conn_data{protocol = http2} = Data) ->
h2_answer_refused_request(From, Data);

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 @@ -1533,6 +1522,19 @@ streaming_body({call, From}, {send_body_chunk, BodyData}, #conn_data{protocol =
{next_state, closed, Data, [{reply, From, {error, Reason}}]}
end;

%% A GOAWAY refused this upload between two of its caller's calls (see
%% h2_on_goaway/3): its next call is the first chance to tell it.
streaming_body({call, From}, Req, #conn_data{protocol = http2, h2_stream_id = StreamId,
h2_streams = Streams} = Data)
when (Req =:= finish_send_body orelse Req =:= start_response orelse
element(1, Req) =:= send_body_chunk),
is_map_key(StreamId, Streams),
element(1, element(2, map_get(StreamId, Streams))) =:= stream,
element(2, element(2, map_get(StreamId, Streams))) =:= refused ->
{_, {stream, refused, ErrorCode}} = maps:get(StreamId, Streams),
h2_upload_done(drop_h2_stream(StreamId, Data),
[{reply, From, {error, {goaway, ErrorCode}}}]);

streaming_body({call, From}, {send_body_chunk, BodyData}, #conn_data{protocol = http2} = Data) ->
%% HTTP/2 - send a DATA frame without END_STREAM.
#conn_data{h2_conn = H2Conn, h2_stream_id = StreamId,
Expand All @@ -1550,7 +1552,7 @@ streaming_body({call, From}, {send_body_chunk, BodyData}, #conn_data{protocol =
ok ->
{keep_state_and_data, [{reply, From, ok}]};
{error, Reason} ->
{next_state, closed, Data, [{reply, From, {error, Reason}}]}
h2_upload_failed(From, Reason, Data)
end;

streaming_body({call, From}, {send_body_chunk, BodyData}, Data) ->
Expand Down Expand Up @@ -1589,7 +1591,7 @@ streaming_body({call, From}, finish_send_body, #conn_data{protocol = http2} = Da
ok ->
{keep_state, Data, [{reply, From, ok}]};
{error, Reason} ->
{next_state, closed, Data, [{reply, From, {error, Reason}}]}
h2_upload_failed(From, Reason, Data)
end;

streaming_body({call, From}, finish_send_body, Data) ->
Expand Down Expand Up @@ -3323,14 +3325,7 @@ h2_stream_owner_down_result(streaming_body, StreamId,
%% The caller streaming a request body died: its stream is gone, so the
%% connection can take requests again, unless that stream was the last
%% one of a GOAWAY drain.
case h2_stream_result(Data#conn_data{h2_stream_id = undefined,
request_from = undefined}, []) of
{keep_state, Data1, Actions} ->
{next_state, connected, Data1,
[{state_timeout, infinity, idle_timeout} | Actions]};
Closed ->
Closed
end;
h2_upload_done(Data, []);
h2_stream_owner_down_result(_State, _StreamId, Data) ->
h2_stream_result(Data, []).

Expand Down Expand Up @@ -4011,48 +4006,44 @@ h2_idle_actions(_Data) ->
%% requests the server went on to complete as failed.
%%
%% A streaming-body request the GOAWAY refuses has no caller parked on it, so
%% it is kept as a `refused' entry and the connection leaves streaming_body:
%% the caller learns of it on its next call (h2_answer_refused_request/2).
%% it is kept as a `refused' entry and the connection stays in streaming_body
%% until the caller's next call learns of it.
h2_on_goaway(LastStreamId, ErrorCode, #conn_data{h2_streams = Streams} = Data) ->
Refused = [SId || SId <- maps:keys(Streams), SId > LastStreamId],
{Replies, Data1} = abort_h2_streams(Refused, ErrorCode, Data),
#conn_data{h2_streams = Streams1, h2_stream_id = StreamId} = Data1,
case map_size(Streams1) =:= 0 of
case map_size(Data1#conn_data.h2_streams) =:= 0 of
true ->
{next_state, closed, Data2, CloseReplies} = h2_close_on_goaway(ErrorCode, Data1),
{next_state, closed, Data2, Replies ++ CloseReplies};
false ->
ok = leave_h2_pool(Data1),
Data3 = Data1#conn_data{h2_goaway = ErrorCode},
case maps:get(StreamId, Streams1, undefined) of
{_, {stream, refused, _}} ->
{next_state, connected, Data3,
[{state_timeout, infinity, idle_timeout} | Replies]};
_ ->
{keep_state, Data3, Replies}
end
{keep_state, Data1#conn_data{h2_goaway = ErrorCode}, Replies}
end.

%% @private Answer a streaming-body call made after the request's stream was
%% refused by a GOAWAY (see h2_on_goaway/3). Any other such call in
%% `connected' is a misuse, as before.
h2_answer_refused_request({Caller, _} = From, #conn_data{h2_streams = Streams} = Data) ->
case [{Id, Code} || {Id, {Owner, {stream, refused, Code}}} <- maps:to_list(Streams),
h2_stream_owner_pid(Owner) =:= Caller] of
[{StreamId, ErrorCode}] ->
Data1 = Data#conn_data{h2_stream_id = undefined, request_from = undefined},
h2_stream_result(drop_h2_stream(StreamId, Data1),
[{reply, From, {error, {goaway, ErrorCode}}}]);
[] ->
{keep_state_and_data, [{reply, From, {error, invalid_state}}]}
%% @private Leave streaming_body once the upload stream is gone: back to
%% connected, or closed when it was the last stream of a GOAWAY drain.
h2_upload_done(Data, Replies) ->
case h2_stream_result(Data#conn_data{h2_stream_id = undefined,
request_from = undefined}, Replies) of
{keep_state, Data1, Actions} ->
{next_state, connected, Data1,
[{state_timeout, infinity, idle_timeout} | Actions]};
Closed ->
Closed
end.

%% @private A body send failed and the connection is unusable: fail the other
%% streams on it too instead of leaving their callers waiting.
h2_upload_failed(From, Reason, Data) ->
{next_state, closed, Data1, Replies} = h2_teardown({closed, Reason}, Data),
{next_state, closed, Data1, [{reply, From, {error, Reason}} | Replies]}.

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) ->
h2_close_on_goaway(ErrorCode, 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
Expand All @@ -4061,7 +4052,10 @@ h2_close_on_goaway(ErrorCode, #conn_data{h2_conn = H2Conn, h2_mon = H2Mon} = Dat
%% 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.
{Replies, Data1} = collect_h2_aborts({goaway, ErrorCode}, Data),
h2_teardown({goaway, ErrorCode}, Data).

h2_teardown(Err, #conn_data{h2_conn = H2Conn, h2_mon = H2Mon} = Data) ->
{Replies, Data1} = collect_h2_aborts(Err, Data),
Data2 = cancel_all_h2_timers(Data1),
_ = case H2Mon of
undefined -> ok;
Expand Down
36 changes: 35 additions & 1 deletion test/hackney_http2_goaway_streaming_tests.erl
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,8 @@ goaway_streaming_test_() ->
{timeout, 30, fun() -> refused_upload_fails_on(finish_send_body) end},
{timeout, 30, fun() -> refused_upload_fails_on(start_response) end},
{timeout, 30, fun accepted_streamed_response_completes/0},
{timeout, 30, fun refused_streamed_response_fails_fast/0}]}.
{timeout, 30, fun refused_streamed_response_fails_fast/0},
{timeout, 30, fun failed_upload_send_fails_other_streams/0}]}.

setup() ->
_ = application:ensure_all_started(hackney),
Expand Down Expand Up @@ -112,6 +113,39 @@ refused_streamed_response_fails_fast() ->
hackney_h2_goaway_server:stop(Server)
end.

%% A body send that fails takes the connection down. The other streams on it
%% must fail with it rather than wait for a response that can no longer come.
%% The server never answers and never opens its flow-control window, so a
%% chunk larger than the initial window times out.
failed_upload_send_fails_other_streams() ->
{Server, Url} = hackney_h2_goaway_server:start(never),
#{port := Port} = uri_string:parse(Url),
{ok, Conn} = hackney_conn_sup:start_conn(#{
host => "localhost",
port => Port,
transport => hackney_ssl,
connect_options => [{protocols, [http2]}],
ssl_options => [{insecure, true}, {verify, verify_none}]
}),
try
ok = hackney_conn:connect(Conn),
{ok, Ref} = hackney_conn:request_async(Conn, <<"GET">>, <<"/">>, [], <<>>,
false),
ok = hackney_conn:send_request_headers(Conn, <<"POST">>, <<"/">>, [],
[{send_timeout, 200}]),
{error, Reason} = hackney_conn:send_body_chunk(Conn, binary:copy(<<"x">>, 70000)),
receive
{hackney_response, Ref, Msg} ->
?assertEqual({error, {closed, Reason}}, Msg)
after 5000 ->
?assert(false)
end,
?assertEqual({ok, closed}, hackney_conn:get_state(Conn))
after
try hackney_conn:stop(Conn) catch _:_ -> ok end,
hackney_h2_goaway_server:stop(Server)
end.

spawn_fetch(Url) ->
Self = self(),
spawn_link(fun() -> Self ! {self(), fetch(Url)} end).
Expand Down
Loading