Skip to content

Commit eb5c00b

Browse files
committed
fix: handle infinity timeout for metadata request
1 parent b471385 commit eb5c00b

2 files changed

Lines changed: 10 additions & 24 deletions

File tree

src/wolff_client.erl

Lines changed: 9 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -401,7 +401,7 @@ ensure_leader_connections2(#{metadata_conn := Pid,
401401
conn_config := ConnConfig,
402402
config := Config
403403
} = St, Group, Topic, MaxPartitions) when is_pid(Pid) ->
404-
Timeout = maps:get(request_timeout, ConnConfig, ?DEFAULT_METADATA_TIMEOUT),
404+
Timeout = metadata_request_timeout(ConnConfig),
405405
IsAutoCreateAllowed = maps:get(allow_auto_topic_creation, Config, false),
406406
case do_get_metadata(Pid, Topic, Timeout, IsAutoCreateAllowed) of
407407
{ok, {Brokers, PartitionMetaList}} ->
@@ -596,7 +596,7 @@ get_metadata([], _ConnConfig, _Topic, _IsAutoCreateAllowed, Errors) ->
596596
get_metadata([Host | Rest], ConnConfig, Topic, IsAutoCreateAllowed, Errors) ->
597597
case do_connect(Host, ConnConfig) of
598598
{ok, Pid} ->
599-
Timeout = maps:get(request_timeout, ConnConfig, ?DEFAULT_METADATA_TIMEOUT),
599+
Timeout = metadata_request_timeout(ConnConfig),
600600
case do_get_metadata(Pid, Topic, Timeout, IsAutoCreateAllowed) of
601601
{ok, Result} ->
602602
{ok, {Pid, Result}};
@@ -770,3 +770,10 @@ bin(X) ->
770770
{error, _} -> bin(io_lib:format("~0p", [X]));
771771
Addr -> bin(Addr)
772772
end.
773+
774+
metadata_request_timeout(#{request_timeout := infinity}) ->
775+
?DEFAULT_METADATA_TIMEOUT * 3;
776+
metadata_request_timeout(#{request_timeout := Timeout}) ->
777+
Timeout;
778+
metadata_request_timeout(_) ->
779+
?DEFAULT_METADATA_TIMEOUT.

test/wolff_client_SUITE.erl

Lines changed: 1 addition & 22 deletions
Original file line numberDiff line numberDiff line change
@@ -80,7 +80,7 @@ t_recv_leader_connection_with_auto_create_retry_success(Config) ->
8080
wolff_client:recv_leader_connection(Client, ?NO_GROUP, Topic, 0, self(), ?all_partitions),
8181
receive
8282
{leader_connection, Pid} when is_pid(Pid) ->
83-
?assert(is_pid(Pid));
83+
true;
8484
{leader_connection, {down, Reason}} ->
8585
%% May get error if topic doesn't exist on real Kafka
8686
?assertNotEqual(undefined, Reason)
@@ -141,24 +141,3 @@ client_config_with_auto_create() ->
141141
#{allow_auto_topic_creation => true,
142142
request_timeout => 2000
143143
}.
144-
145-
client_config_without_auto_create() ->
146-
#{allow_auto_topic_creation => false,
147-
request_timeout => 1000
148-
}.
149-
150-
%% Helper function for meck mocking (similar to pattern used in other test files)
151-
with_meck(Mod, FnName, MockedFn, Action) ->
152-
ok = meck:new(Mod, [non_strict, no_history, no_link, passthrough]),
153-
ok = meck:expect(Mod, FnName, MockedFn),
154-
try
155-
Action()
156-
after
157-
meck:unload(Mod)
158-
end.
159-
160-
%% Helper function to create a mock connection for testing
161-
mock_connection() ->
162-
%% Create a mock connection that can be used in tests
163-
%% This is a simplified mock - in real tests you might want to use a more sophisticated mock
164-
{mock_connection, self()}.

0 commit comments

Comments
 (0)