Packages
brod
2.2.3
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
Files
Jump to
Current section
Files
src/brod_kafka_requests.erl
%%%
%%% Copyright (c) 2014, 2015, Klarna AB
%%%
%%% Licensed under the Apache License, Version 2.0 (the "License");
%%% you may not use this file except in compliance with the License.
%%% You may obtain a copy of the License at
%%%
%%% http://www.apache.org/licenses/LICENSE-2.0
%%%
%%% Unless required by applicable law or agreed to in writing, software
%%% distributed under the License is distributed on an "AS IS" BASIS,
%%% WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
%%% See the License for the specific language governing permissions and
%%% limitations under the License.
%%%
%%%=============================================================================
%%% @doc
%%% This module manages an opaque of sent-request collection.
%%%
%%% @copyright 2014, 2015 Klarna AB
%%% @end
%%% ============================================================================
%%%_* Module declaration =======================================================
%% @private
-module(brod_kafka_requests).
%%%_* Exports ==================================================================
%% API
-export([ new/0
, add/2
, del/2
, get_caller/2
, get_corr_id/1
, increment_corr_id/1
]).
-export_type([requests/0]).
-record(requests,
{ corr_id = 0
, sent = gb_trees:empty() :: gb_trees:tree()
}).
-opaque requests() :: #requests{}.
-define(MAX_CORR_ID_WINDOW_SIZE, (?MAX_CORR_ID div 2)).
%%%_* Includes =================================================================
-include("brod_int.hrl").
%%%_* APIs =====================================================================
-spec new() -> requests().
new() -> #requests{}.
%% @doc Add a new request to sent collection.
%% Return the last corrlation ID and the new opaque.
%% @end
-spec add(requests(), pid()) -> {corr_id(), requests()}.
add(#requests{ corr_id = CorrId
, sent = Sent
} = Requests, Caller) ->
ok = assert_corr_id(Sent),
NewSent = gb_trees:insert(CorrId, Caller, Sent),
NewRequests = Requests#requests{ corr_id = kpro:next_corr_id(CorrId)
, sent = NewSent
},
{CorrId, NewRequests}.
%% @doc Delete a request from the opaque collection.
%% Crash if correlation ID is not found.
%% @end
-spec del(requests(), corr_id()) -> requests().
del(#requests{sent = Sent} = Requests, CorrId) ->
Requests#requests{sent = gb_trees:delete(CorrId, Sent)}.
%% @doc Get caller of a request having the given correlation ID.
%% Crash if the request is not found.
%% @end
-spec get_caller(requests(), corr_id()) -> pid().
get_caller(#requests{sent = Sent}, CorrId) ->
gb_trees:get(CorrId, Sent).
%% @doc Get the correction to be sent for the next request.
-spec get_corr_id(requests()) -> corr_id().
get_corr_id(#requests{ corr_id = CorrId }) ->
CorrId.
%% @doc Fetch and increment the correlation ID
%% This is used if we don't want a response from the broker
%% @end
-spec increment_corr_id(requests()) -> {corr_id(), requests()}.
increment_corr_id(#requests{corr_id = CorrId} = Requests) ->
{CorrId, Requests#requests{ corr_id = kpro:next_corr_id(CorrId) }}.
%%%_* Internal function ========================================================
%% Assert that the in-buffer oldest and newest correlation ids are within a
%% resonable window size. Otherwise it probably means:
%% 1. kafka failed to send response to certain request(s)
%% 2. too many requests sent on to a congested tcp connection.
%% TODO: make it configurable?
assert_corr_id(Sent) ->
Size = corr_id_window_size(Sent),
case Size =< ?MAX_CORR_ID_WINDOW_SIZE of
true -> ok;
false -> erlang:error(corr_id_window_size)
end.
corr_id_window_size(Sent) ->
case gb_trees:is_empty(Sent) of
true -> 0;
false ->
{Min, _} = gb_trees:smallest(Sent),
{Max, _} = gb_trees:largest(Sent),
erlang:min(Max - Min, Min + ?MAX_CORR_ID - Max) + 1
end.
%%%_* Emacs ====================================================================
%%% Local Variables:
%%% allout-layout: t
%%% erlang-indent-level: 2
%%% End: