%%%------------------------------------------------------------------- %%% @author aresei %%% @copyright (C) 2023, %%% @doc %%% %%% @end %%% Created : 06. 7月 2023 12:02 %%%------------------------------------------------------------------- -module(endpoint_kafka). -include("endpoint.hrl"). -behaviour(gen_statem). %% API -export([start_link/2]). %% gen_statem callbacks -export([callback_mode/0, init/1, terminate/3, code_change/4]). -export([disconnected/3, connected/3]). %% 消息重发间隔 -define(RETRY_INTERVAL, 5000). -record(state, { endpoint :: #endpoint{}, buffer :: endpoint_buffer:buffer(), client_id :: atom(), client_pid :: undefined | pid() }). -type kafka_state() :: disconnected | connected. %%%=================================================================== %%% API %%%=================================================================== -spec start_link(LocalName :: atom(), Endpoint :: #endpoint{}) -> {ok, pid()} | ignore | {error, term()}. start_link(LocalName, Endpoint = #endpoint{}) when is_atom(LocalName) -> gen_statem:start_link({local, LocalName}, ?MODULE, [Endpoint], []). %%%=================================================================== %%% gen_statem callbacks %%%=================================================================== -spec callback_mode() -> state_functions. callback_mode() -> state_functions. -spec init(term()) -> gen_statem:init_result(kafka_state(), #state{}). init([Endpoint = #endpoint{id = Id, matcher = Matcher}]) -> ok = iot_log:set_metadata(), erlang:process_flag(trap_exit, true), ok = endpoint_subscription:subscribe(Matcher, self()), Buffer = endpoint_buffer:new(Endpoint, 10), ClientId = list_to_atom("brod_client:" ++ integer_to_list(Id)), {ok, disconnected, #state{endpoint = Endpoint, buffer = Buffer, client_id = ClientId}, [{state_timeout, 0, connect}]}. -spec disconnected(gen_statem:event_type(), term(), #state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). disconnected({call, From}, get_stat, State = #state{buffer = Buffer}) -> reply_stat(From, Buffer, State); disconnected(cast, {forward, Metric}, State = #state{buffer = Buffer}) -> forward_metric(Metric, Buffer, State); disconnected(cast, cleanup, State = #state{buffer = Buffer}) -> cleanup_buffer(Buffer, State); disconnected(cast, {reload, NEndpoint = #endpoint{matcher = NMatcher}}, State = #state{endpoint = #endpoint{matcher = Matcher}, client_id = ClientId}) -> reload_endpoint(Matcher, NMatcher, ClientId, NEndpoint, State); disconnected(state_timeout, connect, State) -> try_connect(State); disconnected(info, {next_data, _Id, _Tuple}, State) -> {keep_state, State}; disconnected(info, {ack, Id}, State = #state{buffer = Buffer}) -> ack_buffer(Id, Buffer, State); disconnected(info, Info, State) -> unknown_info(Info, disconnected, State); disconnected(EventType, EventContent, State) -> unknown_event(EventType, EventContent, disconnected, State). -spec connected(gen_statem:event_type(), term(), #state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). connected({call, From}, get_stat, State = #state{buffer = Buffer}) -> reply_stat(From, Buffer, State); connected(cast, {forward, Metric}, State = #state{buffer = Buffer}) -> forward_metric(Metric, Buffer, State); connected(cast, cleanup, State = #state{buffer = Buffer}) -> cleanup_buffer(Buffer, State); connected(cast, {reload, NEndpoint = #endpoint{matcher = NMatcher}}, State = #state{endpoint = #endpoint{matcher = Matcher}, client_id = ClientId}) -> reload_endpoint(Matcher, NMatcher, ClientId, NEndpoint, State); connected(info, {next_data, Id, Metric}, State = #state{ buffer = Buffer, client_pid = ClientPid, client_id = ClientId, endpoint = #endpoint{config = #kafka_endpoint{topic = Topic}} }) -> ReceiverPid = self(), AckCb = fun(Partition, BaseOffset) -> logger:debug("[endpoint_kafka] ack partion: ~p, offset: ~p", [Partition, BaseOffset]), ReceiverPid ! {ack, Id} end, case catch brod:produce_cb(ClientPid, Topic, random, <<>>, Metric, AckCb) of {ok, _CallRef} -> {keep_state, State}; {ok, _CallRef, _ProducerPid} -> {keep_state, State}; {error, Reason} -> logger:warning("[endpoint_kafka] produce topic: ~p, get error: ~p", [Topic, Reason]), stop_kafka_client(ClientId), NBuffer = endpoint_buffer:recover_inflight(Buffer), {next_state, disconnected, State#state{client_pid = undefined, buffer = NBuffer}, [{state_timeout, ?RETRY_INTERVAL, connect}]}; {'EXIT', Reason} -> logger:warning("[endpoint_kafka] produce topic: ~p, exit with reason: ~p", [Topic, Reason]), stop_kafka_client(ClientId), NBuffer = endpoint_buffer:recover_inflight(Buffer), {next_state, disconnected, State#state{client_pid = undefined, buffer = NBuffer}, [{state_timeout, ?RETRY_INTERVAL, connect}]} end; connected(info, {ack, Id}, State = #state{buffer = Buffer}) -> ack_buffer(Id, Buffer, State); connected(info, {'EXIT', ClientPid, Reason}, State = #state{client_pid = ClientPid, endpoint = #endpoint{title = Title}, buffer = Buffer}) -> logger:warning("[endpoint_kafka] endpoint: ~p, conn pid exit with reason: ~p", [Title, Reason]), NBuffer = endpoint_buffer:recover_inflight(Buffer), {next_state, disconnected, State#state{client_pid = undefined, buffer = NBuffer}, [{state_timeout, ?RETRY_INTERVAL, connect}]}; connected(info, Info, State) -> unknown_info(Info, connected, State); connected(EventType, EventContent, State) -> unknown_event(EventType, EventContent, connected, State). -spec terminate(term(), kafka_state(), #state{}) -> term(). terminate(Reason, _StateName, #state{endpoint = #endpoint{title = Title}, buffer = Buffer, client_id = ClientId}) -> logger:debug("[endpoint_kafka] endpoint: ~p, terminate with reason: ~p", [Title, Reason]), stop_kafka_client(ClientId), endpoint_buffer:cleanup(Buffer), ok. -spec code_change(term() | {down, term()}, kafka_state(), #state{}, term()) -> {ok, kafka_state(), #state{}} | {error, term()}. code_change(_OldVsn, StateName, State = #state{}, _Extra) -> {ok, StateName, State}. %%%=================================================================== %%% Internal functions %%%=================================================================== -spec reply_stat(gen_statem:from(), endpoint_buffer:buffer(), #state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). reply_stat(From, Buffer, State) -> Stat = endpoint_buffer:stat(Buffer), {keep_state, State, [{reply, From, {ok, Stat}}]}. -spec forward_metric(binary(), endpoint_buffer:buffer(), #state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). forward_metric(Metric, Buffer, State) -> NBuffer = endpoint_buffer:append(Metric, Buffer), {keep_state, State#state{buffer = NBuffer}}. -spec cleanup_buffer(endpoint_buffer:buffer(), #state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). cleanup_buffer(Buffer, _State) -> NBuffer = endpoint_buffer:cleanup(Buffer), {keep_state, _State#state{buffer = NBuffer}}. -spec ack_buffer(integer(), endpoint_buffer:buffer(), #state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). ack_buffer(Id, Buffer, State) -> NBuffer = endpoint_buffer:ack(Id, Buffer), {keep_state, State#state{buffer = NBuffer}}. -spec reload_endpoint(binary(), binary(), atom(), #endpoint{}, #state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). reload_endpoint(Matcher, NMatcher, ClientId, NEndpoint, State = #state{}) -> ensure_subscription(Matcher, NMatcher), stop_kafka_client(ClientId), NBuffer = endpoint_buffer:recover_inflight(State#state.buffer), {next_state, disconnected, State#state{endpoint = NEndpoint, client_pid = undefined, buffer = NBuffer}, [{state_timeout, 0, connect}]}. -spec try_connect(#state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). try_connect(State = #state{ buffer = Buffer, client_id = ClientId, endpoint = #endpoint{ title = Title, config = #kafka_endpoint{ sasl_config = SaslConfig, bootstrap_servers = BootstrapServers, topic = Topic } } }) -> logger:debug("[endpoint_kafka] endpoint: ~p, create postman", [Title]), BaseConfig = [ {reconnect_cool_down_seconds, 5}, {socket_options, [{keepalive, true}]} ], ClientConfig = case SaslConfig of {Mechanism, Username, Password} -> [{sasl, {Mechanism, Username, Password}} | BaseConfig]; undefined -> BaseConfig end, case catch brod:start_link_client(BootstrapServers, ClientId, ClientConfig) of {ok, ClientPid} -> case brod:start_producer(ClientId, Topic, []) of ok -> NBuffer = endpoint_buffer:trigger_n(Buffer), {next_state, connected, State#state{buffer = NBuffer, client_pid = ClientPid}}; {error, Reason} -> logger:debug("[endpoint_kafka] start_producer: ~p, get error: ~p", [ClientId, Reason]), stop_kafka_client(ClientId), {keep_state, State#state{client_pid = undefined}, [{state_timeout, ?RETRY_INTERVAL, connect}]} end; Error -> logger:debug("[endpoint_kafka] start_client: ~p, get error: ~p", [ClientId, Error]), {keep_state, State#state{client_pid = undefined}, [{state_timeout, ?RETRY_INTERVAL, connect}]} end. -spec unknown_info(term(), kafka_state(), #state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). unknown_info(Info, StateName, State) -> logger:warning("[endpoint_kafka] unknown message: ~p, status: ~p", [Info, StateName]), {keep_state, State}. -spec unknown_event(gen_statem:event_type(), term(), kafka_state(), #state{}) -> gen_statem:event_handler_result(kafka_state(), #state{}). unknown_event(EventType, EventContent, StateName, State) -> logger:warning("[endpoint_kafka] unknown event: ~p, content: ~p, status: ~p", [EventType, EventContent, StateName]), {keep_state, State}. -spec ensure_subscription(binary(), binary()) -> ok. ensure_subscription(Matcher, Matcher) -> ok; ensure_subscription(Matcher, NMatcher) -> ok = endpoint_subscription:unsubscribe(Matcher, self()), endpoint_subscription:subscribe(NMatcher, self()). -spec stop_kafka_client(atom()) -> ok. stop_kafka_client(ClientId) when is_atom(ClientId) -> _ = catch brod:stop_client(ClientId), ok.