Packages
brod
3.3.2
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_apis.erl
%%%
%%% Copyright (c) 2017 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.
%%%
%% Version ranges are cached per host and per brod_sock pid in ets
-module(brod_kafka_apis).
-export([ default_version/1
, maybe_add_sock_pid/2
, pick_version/2
, start_link/0
, stop/0
, versions_received/4
]).
-export([ code_change/3
, handle_call/3
, handle_cast/2
, handle_info/2
, init/1
, terminate/2
]).
-export_type([ api/0
, vsn/0
]).
-define(SERVER, ?MODULE).
-define(ETS, ?MODULE).
-record(state, {}).
-type vsn() :: non_neg_integer().
-type range() :: {vsn(), vsn()}.
-type api() :: kpro:req_tag().
-type client_id() :: binary(). %% not brod:client_id()
-type host() :: brod:hostname().
-type versions() :: [{api(), range()}].
%% @doc Start process.
-spec start_link() -> {ok, pid()}.
start_link() ->
gen_server:start_link({local, ?SERVER}, ?MODULE, [], []).
-spec stop() -> ok.
stop() ->
gen_server:call(?SERVER, stop, infinity).
%% @doc Report API version ranges for a given `brod_sock' pid.
-spec versions_received(client_id(), pid(), versions(), host()) -> ok.
versions_received(ClientId, SockPid, Versions, Host) ->
Vsns = resolve_version_ranges(ClientId, Versions, []),
gen_server:call(?SERVER, {versions_received, SockPid, Vsns, Host}, infinity).
%% @doc Get default supported version for the given API.
-spec default_version(api()) -> vsn().
default_version(API) ->
{Min, _Max} = supported_versions(API),
Min.
%% @doc Try add pid with existing version ranges.
%% Return `{error, unknow_host}' if the host is not cached already.
%% @end
-spec maybe_add_sock_pid(host(), pid()) -> ok | {error, unknown_host}.
maybe_add_sock_pid(Host, SockPid) ->
case ets:lookup(?ETS, Host) of
[] ->
{error, unknown_host};
[{_, ResolvedVersions}] ->
gen_server:call(?SERVER, {add_sock_pid, SockPid, ResolvedVersions})
end.
%% @doc Pick API version for the given API.
-spec pick_version(pid(), api()) -> vsn().
pick_version(SockPid, API) ->
do_pick_version(SockPid, API, supported_versions(API)).
init([]) ->
?ETS = ets:new(?ETS, [named_table, protected]),
{ok, #state{}}.
handle_info({'DOWN', _Mref, process, Pid, _Reason}, State) ->
ets:delete(?ETS, Pid),
{noreply, State};
handle_info(Info, State) ->
error_logger:error_msg("unknown info ~p", [Info]),
{noreply, State}.
handle_cast(Cast, State) ->
error_logger:error_msg("unknown cast ~p", [Cast]),
{noreply, State}.
handle_call(stop, From, State) ->
gen_server:reply(From, ok),
{stop, normal, State};
handle_call({add_sock_pid, SockPid, ResolvedVersions}, _From, State) ->
_ = erlang:monitor(process, SockPid),
ets:insert(?ETS, {SockPid, ResolvedVersions}),
{reply, ok, State};
handle_call({versions_received, SockPid, Versions, Host}, _From, State) ->
_ = erlang:monitor(process, SockPid),
ets:insert(?ETS, {Host, Versions}),
ets:insert(?ETS, {SockPid, Versions}),
{reply, ok, State};
handle_call(Call, _From, State) ->
{reply, {error, {unknown_call, Call}}, State}.
code_change(_OldVsn, State, _Extra) ->
{ok, State}.
terminate(_Reason, _State) ->
ok.
%% @private
-spec do_pick_version(pid(), api(), range()) -> vsn().
do_pick_version(_SockPid, _API, {Vsn, Vsn}) ->
%% only one version supported, no need to lookup
Vsn;
do_pick_version(SockPid, API, {Min, _Max}) ->
%% query the highest supported version
case lookup_version(SockPid, API) of
none -> Min; %% no version received from kafka, use min
Vsn -> Vsn %% use max supported version
end.
%% @private Lookup API from cache, return default if not found.
-spec lookup_version(pid(), api()) -> vsn() | none.
lookup_version(SockPid, API) ->
case ets:lookup(?ETS, SockPid) of
[] -> none;
[{SockPid, Versions}] ->
case lists:keyfind(API, 1, Versions) of
{API, Vsn} -> Vsn;
false -> none
end
end.
%% @private
-spec resolve_version_ranges(client_id(), [{api(), range()}], Acc) -> Acc
when Acc :: [{api(), vsn()}].
resolve_version_ranges(_ClientId, [], Acc) -> lists:reverse(Acc);
resolve_version_ranges(ClientId, [{API, {MinKafka, MaxKafka}} | Rest], Acc) ->
case resolve_version_range(ClientId, API, MinKafka, MaxKafka,
supported_versions(API)) of
none -> resolve_version_ranges(ClientId, Rest, Acc);
Max -> resolve_version_ranges(ClientId, Rest, [{API, Max} | Acc])
end.
%% @private
-spec resolve_version_range(client_id(), api(), vsn(), vsn(),
range() | none()) -> vsn() | none.
resolve_version_range(_ClientId, _API, _MinKafka, _MaxKafka, none) ->
%% API not implemented by brod
none;
resolve_version_range(ClientId, API, MinKafka, MaxKafka, {MinBrod, MaxBrod}) ->
Min = max(MinBrod, MinKafka),
Max = min(MaxBrod, MaxKafka),
case Min =< Max of
true when MinBrod =:= MaxBrod ->
%% if brod supports only one version
%% no need to store the range in ETS
none;
true ->
Max;
false ->
log_unsupported_api(ClientId, API,
{MinBrod, MaxBrod}, {MinKafka, MaxKafka}),
none
end.
%% @private
-spec log_unsupported_api(client_id(), api(), range(), range()) -> ok.
log_unsupported_api(ClientId, API, BrodRange, KafkaRange) ->
error_logger:error_msg("Can not support API ~p for client ~p, "
"brod versions: ~p, kafka versions: ~p",
[API, ClientId, BrodRange, KafkaRange]),
ok.
%% @private Do not change range without verification.
%%% Fixed (hardcoded) version APIs
%% sasl_handshake_request: 0
%% api_versions_request: 0
%%% Missing features
%% {create_topics_request, 0, 0}
%% {delete_topics_request, 0, 0}
%%% Will not support
%% leader_and_isr_request
%% stop_replica_request
%% update_metadata_request
%% controlled_shutdown_request
%% @end
supported_versions(API) ->
case API of
produce_request -> {0, 2};
fetch_request -> {0, 3};
offsets_request -> {0, 1};
metadata_request -> {0, 2};
offset_commit_request -> {2, 2};
offset_fetch_request -> {1, 2};
group_coordinator_request -> {0, 0};
join_group_request -> {0, 0};
heartbeat_request -> {0, 0};
leave_group_request -> {0, 0};
sync_group_request -> {0, 0};
describe_groups_request -> {0, 0};
list_groups_request -> {0, 0};
_ -> none
end.
%%%_* Emacs ====================================================================
%%% Local Variables:
%%% allout-layout: t
%%% erlang-indent-level: 2
%%% End: