Current section

113 Versions

Jump to

Compare versions

5 files changed
+87 additions
-7 deletions
  @@ -1,6 +1,6 @@
1 1 PROJECT = brod
2 2 PROJECT_DESCRIPTION = Kafka client library in Erlang
3 - PROJECT_VERSION = 2.3.5
3 + PROJECT_VERSION = 2.3.6
4 4
5 5 DEPS = supervisor3 kafka_protocol
  @@ -1,5 +1,5 @@
1 1 {<<"name">>,<<"brod">>}.
2 - {<<"version">>,<<"2.3.5">>}.
2 + {<<"version">>,<<"2.3.6">>}.
3 3 {<<"requirements">>,
4 4 #{<<"kafka_protocol">> => #{<<"app">> => <<"kafka_protocol">>,
5 5 <<"optional">> => false,
unknownscripts/brod
File is too large to be displayed (100 KB limit).
  @@ -1,6 +1,6 @@
1 1 {application,brod,
2 2 [{description,"Apache Kafka Erlang client library"},
3 - {vsn,"2.3.5"},
3 + {vsn,"2.3.6"},
4 4 {registered,[]},
5 5 {applications,[kernel,stdlib,ssl,kafka_protocol,supervisor3]},
6 6 {env,[]},
  @@ -49,17 +49,26 @@
49 49
50 50 -define(DEFAULT_CONNECT_TIMEOUT, timer:seconds(5)).
51 51 -define(DEFAULT_REQUEST_TIMEOUT, timer:minutes(4)).
52 + -define(SIZE_HEAD_BYTES, 4).
52 53
53 54 %%%_* Includes =================================================================
54 55 -include("brod_int.hrl").
55 56
56 57 -type options() :: proplists:proplist().
57 58 -type requests() :: brod_kafka_requests:requests().
59 + -type byte_count() :: non_neg_integer().
60 +
61 + -record(acc, { expected_size = error(bad_init) :: byte_count()
62 + , acc_size = 0 :: byte_count()
63 + , acc_buffer = [] :: [binary()] %% received bytes in reversed order
64 + }).
65 +
66 + -type acc() :: binary() | #acc{}.
58 67
59 68 -record(state, { client_id :: binary()
60 69 , parent :: pid()
61 70 , sock :: port()
62 - , tail = <<>> :: binary() %% leftover of last data stream
71 + , acc = <<>> :: acc()
63 72 , requests :: requests()
64 73 , mod :: gen_tcp | ssl
65 74 , req_timeout :: timeout()
  @@ -277,7 +286,7 @@ decode_msg(Msg, State, Debug0) ->
277 286 handle_msg(Msg, State, Debug).
278 287
279 288 handle_msg({_, Sock, Bin}, #state{ sock = Sock
280 - , tail = Tail0
289 + , acc = Acc0
281 290 , requests = Requests
282 291 , mod = Mod
283 292 } = State, Debug) when is_binary(Bin) ->
  @@ -285,7 +294,8 @@ handle_msg({_, Sock, Bin}, #state{ sock = Sock
285 294 gen_tcp -> ok = inet:setopts(Sock, [{active, once}]);
286 295 ssl -> ok = ssl:setopts(Sock, [{active, once}])
287 296 end,
288 - {Responses, Tail} = kpro:decode_response(<<Tail0/binary, Bin/binary>>),
297 + Acc1 = acc_recv_bytes(Acc0, Bin),
298 + {Responses, Acc} = decode_response(Acc1),
289 299 NewRequests =
290 300 lists:foldl(
291 301 fun(#kpro_Response{ correlationId = CorrId
  @@ -295,7 +305,7 @@ handle_msg({_, Sock, Bin}, #state{ sock = Sock
295 305 cast(Caller, {msg, self(), CorrId, Response}),
296 306 brod_kafka_requests:del(Reqs, CorrId)
297 307 end, Requests, Responses),
298 - ?MODULE:loop(State#state{tail = Tail, requests = NewRequests}, Debug);
308 + ?MODULE:loop(State#state{acc = Acc, requests = NewRequests}, Debug);
299 309 handle_msg(assert_max_req_age, #state{ requests = Requests
300 310 , req_timeout = ReqTimeout
301 311 } = State, Debug) ->
  @@ -434,6 +444,76 @@ send_assert_max_req_age(Pid, Timeout) when Timeout >= 1000 ->
434 444 _ = erlang:send_after(SendAfter, Pid, assert_max_req_age),
435 445 ok.
436 446
447 + %% @private Accumulate newly received bytes.
448 + -spec acc_recv_bytes(acc(), binary()) -> acc().
449 + acc_recv_bytes(Acc, NewBytes) when is_binary(Acc) ->
450 + case <<Acc/binary, NewBytes/binary>> of
451 + <<Size:32/signed-integer, _/binary>> = AccBytes ->
452 + do_acc(#acc{expected_size = Size + ?SIZE_HEAD_BYTES}, AccBytes);
453 + AccBytes ->
454 + AccBytes
455 + end;
456 + acc_recv_bytes(#acc{} = Acc, NewBytes) ->
457 + do_acc(Acc, NewBytes).
458 +
459 + %% @private Add newly received bytes to buffer.
460 + -spec do_acc(acc(), binary()) -> acc().
461 + do_acc(#acc{acc_size = AccSize, acc_buffer = AccBuffer} = Acc, NewBytes) ->
462 + Acc#acc{acc_size = AccSize + size(NewBytes),
463 + acc_buffer = [NewBytes | AccBuffer]
464 + }.
465 +
466 + %% @private Decode response when accumulated enough bytes.
467 + -spec decode_response(acc()) -> {[kpro_Response()], acc()}.
468 + decode_response(#acc{expected_size = ExpectedSize,
469 + acc_size = AccSize,
470 + acc_buffer = AccBuffer}) when AccSize >= ExpectedSize ->
471 + %% iolist_to_binary here to simplify kafka_protocol implementation
472 + %% maybe make it smarter in the next version
473 + kpro:decode_response(iolist_to_binary(lists:reverse(AccBuffer)));
474 + decode_response(Acc) ->
475 + {[], Acc}.
476 +
477 + %%%_* Eunit ====================================================================
478 +
479 + -ifdef(TEST).
480 +
481 + -include_lib("eunit/include/eunit.hrl").
482 +
483 + acc_test_() ->
484 + [{"clean start flow",
485 + fun() ->
486 + Acc0 = acc_recv_bytes(<<>>, <<0, 0>>),
487 + ?assertEqual(Acc0, <<0, 0>>),
488 + Acc1 = acc_recv_bytes(Acc0, <<0, 1, 0, 0>>),
489 + ?assertEqual(#acc{expected_size = 5,
490 + acc_size = 6,
491 + acc_buffer = [<<0, 0, 0, 1, 0, 0>>]
492 + }, Acc1)
493 + end},
494 + {"old tail leftover",
495 + fun() ->
496 + Acc0 = acc_recv_bytes(<<0, 0>>, <<0, 4>>),
497 + ?assertEqual(#acc{expected_size = 8,
498 + acc_size = 4,
499 + acc_buffer = [<<0, 0, 0, 4>>]
500 + }, Acc0),
501 + Acc1 = acc_recv_bytes(Acc0, <<0, 0>>),
502 + ?assertEqual(#acc{expected_size = 8,
503 + acc_size = 6,
504 + acc_buffer = [<<0, 0>>, <<0, 0, 0, 4>>]
505 + }, Acc1),
506 + Acc2 = acc_recv_bytes(Acc1, <<1, 1>>),
507 + ?assertEqual(#acc{expected_size = 8,
508 + acc_size = 8,
509 + acc_buffer = [<<1, 1>>, <<0, 0>>, <<0, 0, 0, 4>>]
510 + }, Acc2)
511 + end
512 + }
513 + ].
514 +
515 + -endif.
516 +
437 517 %%%_* Emacs ====================================================================
438 518 %%% Local Variables:
439 519 %%% allout-layout: t