Packages
brod
2.3.6
4.5.7
4.5.6
4.5.5
4.5.4
4.5.3
4.5.2
4.5.1
4.5.0
4.4.7
4.4.6
4.4.5
4.4.4
4.4.3
4.4.2
4.4.1
4.4.0
4.3.3
4.3.2
4.3.1
4.3.0
4.2.0
4.1.1
4.1.0
4.0.0
3.19.1
3.19.0
3.18.0
3.17.1
3.17.0
3.16.5
3.16.4
3.16.3
3.16.2
3.16.1
3.16.0
3.15.6
3.15.5
3.15.4
3.15.3
3.15.1
3.15.0
3.14.0
3.13.0
3.12.0
3.11.0
3.10.0
3.9.5
3.9.3
3.9.2
3.9.1
3.9.0
3.8.1
3.8.0
3.7.11
3.7.10
3.7.9
3.7.8
3.7.7
3.7.6
3.7.5
3.7.4
3.7.3
3.7.2
3.7.1
3.7.0
3.6.2
3.6.1
3.6.0
3.5.2
3.5.1
3.5.0
3.4.0
3.3.5
3.3.4
3.3.3
3.3.2
3.3.1
3.3.0
3.2.0
3.0.0
2.5.0
2.4.1
2.4.0
2.3.7
2.3.6
2.3.5
2.3.4
2.3.3
2.3.1
2.2.16
2.2.15
2.2.14
2.2.12
2.2.11
2.2.10
2.2.9
2.2.8
2.2.7
2.2.6
2.2.5
2.2.4
2.2.3
2.2.2
2.2.1
2.2.0
2.1.12
2.1.11
2.1.10
2.1.8
2.1.7
2.1.4
2.1.2
2.0.0
Apache Kafka Erlang client library
Current section
113 Versions
Jump to
Current section
113 Versions
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 |