% @hidden -module(kafe_protocol_join_group). -include("../include/kafe.hrl"). -export([ run/2, request/3, response/2 ]). run(GroupId, Options) -> case kafe:group_coordinator(bucs:to_binary(GroupId)) of {ok, #{coordinator_host := Host, coordinator_port := Port, error_code := none}} -> kafe_protocol:run({host_and_port, Host, Port}, {call, fun ?MODULE:request/3, [GroupId, Options], fun ?MODULE:response/2}); _ -> {error, no_broker_found} end. request(GroupId, Options, State) -> SessionTimeout = maps:get(session_timeout, Options, ?DEFAULT_JOIN_GROUP_SESSION_TIMEOUT), MemberId = maps:get(member_id, Options, ?DEFAULT_JOIN_GROUP_MEMBER_ID), ProtocolType = maps:get(protocol_type, Options, ?DEFAULT_JOIN_GROUP_PROTOCOL_TYPE), GroupProtocols = maps:get(protocols, Options, ?DEFAULT_JOIN_GROUP_PROTOCOLS), kafe_protocol:request( ?JOIN_GROUP_REQUEST, <<(kafe_protocol:encode_string(GroupId))/binary, SessionTimeout:32/signed, (kafe_protocol:encode_string(MemberId))/binary, (kafe_protocol:encode_string(ProtocolType))/binary, (kafe_protocol:encode_array(GroupProtocols))/binary>>, State, ?V0). response(<>, _ApiVersion) -> {ok, #{error_code => kafe_error:code(ErrorCode), generation_id => GenerationId, protocol_group => ProtocolGroup, leader_id => LeaderId, member_id => MemberId, members => response(MembersLength, Members, [])}}. response(0, _, Acc) -> Acc; response(N, <>, Acc) -> response(N - 1, Rest, [#{member_id => MemberId, member_metadata => MemberMetadata}|Acc]).