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
45 changes: 41 additions & 4 deletions lib/ex_aws/bedrock/event_stream.ex
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,11 @@ defmodule ExAws.Bedrock.EventStream do
{"user-agent", @user_agent},
{"x-amzn-bedrock-accept", "*/*"}
]
@hackney_options [{:async, :once}]
# :protocols forces HTTP/1.1: hackney >= 4 negotiates HTTP/2 via ALPN by
# default, but its h2 path does not deliver async body messages, so the
# event stream would hang waiting for chunks that never arrive. HTTP/1.1
# also guarantees the chunked transfer-encoding this module verifies.
@hackney_options [{:async, :once}, {:protocols, [:http1]}]

@doc """
Stream of chunks from the response stream.
Expand All @@ -57,8 +61,10 @@ defmodule ExAws.Bedrock.EventStream do
encoded_data
)

hackney_opts = hackney_options(config)

request_fun = fn [] ->
{:ok, ref} = :hackney.post(url, full_headers, encoded_data, @hackney_options)
{:ok, ref} = :hackney.post(url, full_headers, encoded_data, hackney_opts)

receive do
{:hackney_response, ^ref, {:status, 200, _reason}} ->
Expand All @@ -69,6 +75,9 @@ defmodule ExAws.Bedrock.EventStream do

{:hackney_response, ^ref, {:error, {:closed, :timeout}}} ->
:closed

{:hackney_response, ^ref, {:error, reason}} ->
raise ExAws.Error, "Bedrock stream request failed: #{inspect(reason)}"
end
end

Expand All @@ -82,7 +91,9 @@ defmodule ExAws.Bedrock.EventStream do
{:error, status, reason} ->
raise ExAws.Error, "#{to_string(status)}: #{to_string(reason)}"

ref when is_reference(ref) ->
# hackney < 4 identifies async responses by reference,
# hackney >= 4 by the connection pid.
ref when is_reference(ref) or is_pid(ref) ->
:ok = :hackney.stream_next(ref)

receive do
Expand All @@ -94,8 +105,15 @@ defmodule ExAws.Bedrock.EventStream do
{:hackney_response, ^ref, :done} ->
{:halt, []}

{:hackney_response, ^ref, data} ->
{:hackney_response, ^ref, {:error, reason}} ->
raise ExAws.Error, "Bedrock stream failed mid-stream: #{inspect(reason)}"

{:hackney_response, ^ref, data} when is_binary(data) ->
{[data], ref}

{:hackney_response, ^ref, other} ->
raise ExAws.Error,
"Bedrock stream received unexpected message: #{inspect(other)}"
end
end,
&Function.identity/1
Expand All @@ -104,6 +122,25 @@ defmodule ExAws.Bedrock.EventStream do
Stream.flat_map(stream, &decode_chunk/1)
end

@doc """
Builds the hackney options for the streaming request.

Merges caller-provided options from the ExAws config `:http_opts` (e.g.
`recv_timeout`, `connect_timeout`, `pool`) on top of the async-streaming
defaults. Without this, the stream would always use hackney's built-in
`recv_timeout` (5s) and drop slow responses regardless of the timeout the
caller configured.

The streaming defaults win on conflicting keys, so the async-streaming mode
(`async: :once`) and the forced HTTP/1.1 protocol can't be accidentally
disabled by caller options.
"""
def hackney_options(config) do
config
|> Map.get(:http_opts, [])
|> Keyword.merge(@hackney_options)
end

defp verify_event_stream!(headers) do
verify_header!(headers, "Content-Type", @content_type)
end
Expand Down
29 changes: 29 additions & 0 deletions test/ex_aws/bedrock/event_stream_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,35 @@ defmodule ExAws.Bedrock.EventStreamTest do
end
end

describe "hackney_options/1" do
test "defaults to async streaming over HTTP/1.1 when no http_opts are given" do
assert EventStream.hackney_options(%{}) == [async: :once, protocols: [:http1]]
end

test "honors caller recv_timeout / connect_timeout / pool from :http_opts" do
opts =
EventStream.hackney_options(%{
http_opts: [recv_timeout: 600_000, connect_timeout: 10_000, pool: :ex_aws]
})

assert Keyword.get(opts, :async) == :once
assert Keyword.get(opts, :recv_timeout) == 600_000
assert Keyword.get(opts, :connect_timeout) == 10_000
assert Keyword.get(opts, :pool) == :ex_aws
end

test "streaming defaults win so async: :once and HTTP/1.1 cannot be disabled" do
opts =
EventStream.hackney_options(%{
http_opts: [async: false, protocols: [:http2], recv_timeout: 1_000]
})

assert Keyword.get(opts, :async) == :once
assert Keyword.get(opts, :protocols) == [:http1]
assert Keyword.get(opts, :recv_timeout) == 1_000
end
end

setup_all do
# Claude 3.5 Sonnet Multi-Chunk response
multipart_chunk =
Expand Down
Loading