% @hidden -module(kafe_protocol_sync_group). -include("../include/kafe.hrl"). -export([ run/4, request/5, response/2 ]). run(GroupId, GenerationId, MemberId, Assignments) -> 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/5, [GroupId, GenerationId, MemberId, Assignments], fun ?MODULE:response/2}); _ -> {error, no_broker_found} end. % SyncGroup Request (Version: 0) => group_id generation_id member_id [group_assignment] % group_id => STRING % generation_id => INT32 % member_id => STRING % group_assignment => member_id member_assignment % member_id => STRING % member_assignment => BYTES % % MemberAssignment => Version PartitionAssignment % Version => int16 % PartitionAssignment => [Topic [Partition]] UserData % Topic => string % Partition => int32 % UserData => bytes request(GroupId, GenerationId, MemberId, Assignments, State) -> kafe_protocol:request( ?SYNC_GROUP_REQUEST, <<(kafe_protocol:encode_string(GroupId))/binary, GenerationId:32/signed, (kafe_protocol:encode_string(MemberId))/binary, (group_assignment(Assignments, []))/binary>>, State, ?V0). group_assignment([], Acc) -> kafe_protocol:encode_array(lists:reverse(Acc)); group_assignment([#{member_id := MemberId, member_assignment := MemberAssignment}|Rest], Acc) -> Version = maps:get(version, MemberAssignment, ?DEFAULT_GROUP_PROTOCOL_VERSION), Partitions = maps:get(partition_assignment, MemberAssignment, ?DEFAULT_GROUP_PARTITION_ASSIGNMENT), UserData = maps:get(user_data, MemberAssignment, ?DEFAULT_GROUP_USER_DATA), group_assignment(Rest, [<<(kafe_protocol:encode_string(MemberId))/binary, (kafe_protocol:encode_bytes( <>))/binary>>|Acc]). partition_assignment([], Acc) -> kafe_protocol:encode_array(lists:reverse(Acc)); partition_assignment([#{topic := Topic, partitions := Partitions}|Rest], Acc) -> Partitions1 = kafe_protocol:encode_array(lists:map(fun(E) -> <> end, Partitions)), partition_assignment(Rest, [<<(kafe_protocol:encode_string(Topic))/binary, Partitions1/binary>>|Acc]). % SyncGroupResponse => ErrorCode MemberAssignment % ErrorCode => int16 % MemberAssignment => bytes response(<>, _ApiVersion) -> case MemberAssignment of <> -> {PartitionAssignment, UserData} = partition_assignment(PartitionAssignmentSize, Remainder, []), {ok, #{error_code => kafe_error:code(ErrorCode), version => Version, partition_assignment => PartitionAssignment, user_data => UserData}}; _ -> {ok, #{error_code => kafe_error:code(ErrorCode), version => -1, partition_assignment => [], user_data => <<>>}} end. partition_assignment(0, <>, Acc) -> {Acc, UserData}; partition_assignment(N, <>, Acc) -> {Partitions, Remainder1} = partitions(NbPartitions, Remainder, []), partition_assignment(N - 1, Remainder1, [#{topic => Topic, partitions => Partitions}|Acc]). partitions(0, Remainder, Acc) -> {Acc, Remainder}; partitions(N, <>, Acc) -> partitions(N - 1, Remainder, [Partition|Acc]).