diff --git a/src/hackney_conn.erl b/src/hackney_conn.erl index 95593b91..150d954e 100644 --- a/src/hackney_conn.erl +++ b/src/hackney_conn.erl @@ -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}}]}; @@ -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, @@ -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) -> @@ -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) -> @@ -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, []). @@ -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 @@ -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; diff --git a/test/hackney_http2_goaway_streaming_tests.erl b/test/hackney_http2_goaway_streaming_tests.erl index 25ab73c6..2d7c258b 100644 --- a/test/hackney_http2_goaway_streaming_tests.erl +++ b/test/hackney_http2_goaway_streaming_tests.erl @@ -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), @@ -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).