Current section

113 Versions

Jump to

Compare versions

5 files changed
+10 additions
-7 deletions
  @@ -1,6 +1,6 @@
1 1 PROJECT = brod
2 2 PROJECT_DESCRIPTION = Kafka client library in Erlang
3 - PROJECT_VERSION = 2.2.14
3 + PROJECT_VERSION = 2.2.15
4 4
5 5 DEPS = supervisor3 kafka_protocol
  @@ -1,5 +1,5 @@
1 1 {<<"name">>,<<"brod">>}.
2 - {<<"version">>,<<"2.2.14">>}.
2 + {<<"version">>,<<"2.2.15">>}.
3 3 {<<"requirements">>,
4 4 #{<<"kafka_protocol">> => #{<<"app">> => <<"kafka_protocol">>,
5 5 <<"optional">> => false,
  @@ -1,6 +1,6 @@
1 1 {application,brod,
2 2 [{description,"Apache Kafka Erlang client library"},
3 - {vsn,"2.2.14"},
3 + {vsn,"2.2.15"},
4 4 {registered,[]},
5 5 {applications,[kernel,stdlib,ssl,kafka_protocol,supervisor3]},
6 6 {env,[]},
  @@ -392,7 +392,7 @@ handle_add_offset(#pending_acks{ offsets_queue = Queue
392 392 _ ->
393 393 handle_add_offset(PendingAcks, Offsets)
394 394 end;
395 - handle_add_offset(PendingAcks, []) ->
395 + handle_add_offset(#pending_acks{} = PendingAcks, []) ->
396 396 PendingAcks.
397 397
398 398 maybe_shrink_max_bytes(#state{ prefetch_count = PrefetchCount
  @@ -48,7 +48,7 @@
48 48 ]).
49 49
50 50 -define(DEFAULT_CONNECT_TIMEOUT, timer:seconds(5)).
51 - -define(DEFAULT_REQUEST_TIMEOUT, timer:seconds(120)).
51 + -define(DEFAULT_REQUEST_TIMEOUT, timer:minutes(4)).
52 52
53 53 %%%_* Includes =================================================================
54 54 -include("brod_int.hrl").
  @@ -303,7 +303,7 @@ handle_msg({From, stop}, #state{mod = Mod, sock = Sock}, _Debug) ->
303 303 Mod:close(Sock),
304 304 _ = reply(From, ok),
305 305 ok;
306 - handle_msg(Msg, State, Debug) ->
306 + handle_msg(Msg, #state{} = State, Debug) ->
307 307 error_logger:warning_msg("[~p] ~p got unrecognized message: ~p",
308 308 [?MODULE, self(), Msg]),
309 309 ?MODULE:loop(State, Debug).
  @@ -387,7 +387,10 @@ assert_max_req_age(Requests, Timeout) ->
387 387 %% @end
388 388 -spec send_assert_max_req_age(pid(), timeout()) -> ok.
389 389 send_assert_max_req_age(Pid, Timeout) when Timeout >= 1000 ->
390 - _ = erlang:send_after(Timeout div 2, Pid, assert_max_req_age),
390 + %% Check every 1 minute
391 + %% or every half of the timeout value if it's less than 2 minute
392 + SendAfter = erlang:min(Timeout div 2, timer:minutes(1)),
393 + _ = erlang:send_after(SendAfter, Pid, assert_max_req_age),
391 394 ok.
392 395
393 396 %%%_* Emacs ====================================================================