% @hidden -module(kafe_protocol_metadata). -include("../include/kafe.hrl"). -export([ run/1, request/2, response/2 ]). run(Topics) -> kafe_protocol:run({call, fun ?MODULE:request/2, [Topics], fun ?MODULE:response/2}). % Metadata Request (Version: 0) => [topics] request(TopicNames, State) -> kafe_protocol:request( ?METADATA_REQUEST, <<(kafe_protocol:encode_array( [kafe_protocol:encode_string(bucs:to_binary(Name)) || Name <- TopicNames]))/binary>>, State, ?V0). % Metadata Response (Version: 0) => [brokers] [topic_metadata] % brokers => node_id host port % node_id => INT32 % host => STRING % port => INT32 % topic_metadata => topic_error_code topic [partition_metadata] % topic_error_code => INT16 % topic => STRING % partition_metadata => partition_error_code partition_id leader [replicas] [isr] % partition_error_code => INT16 % partition_id => INT32 % leader => INT32 response(<>, _ApiVersion) -> { Brokers, <> } = brokers(NumberOfBrokers, BrokerRemainder, []), {Topics, _} = topics(NumberOfTopics, TopicMetadataRemainder, []), {ok, #{brokers => Brokers, topics => Topics}}. % Private brokers(0, Remainder, Acc) -> {Acc, Remainder}; brokers( N, << NodeId:32/signed, HostLength:16/signed, Host:HostLength/bytes, Port:32/signed, Remainder/binary >>, Acc) -> brokers(N-1, Remainder, [#{id => NodeId, host => Host, port => Port} | Acc]). topics(0, <>, Acc) -> {Acc, Remainder}; topics( N, << ErrorCode:16/signed, TopicNameLen:16/signed, TopicName:TopicNameLen/bytes, PartitionLength:32/signed, PartitionsRemainder/binary >>, Acc) -> {Partitions, Remainder} = partitions(PartitionLength, PartitionsRemainder, []), topics(N-1, Remainder, [#{error_code => kafe_error:code(ErrorCode), name => TopicName, partitions => Partitions} | Acc]). partitions(0, <>, Acc) -> {Acc, Remainder}; partitions( N, << ErrorCode:16/signed, Id:32/signed, Leader:32/signed, NumberOfReplicas:32/signed, ReplicasRemainder/binary >>, Acc) -> { Replicas, <> } = replicas(NumberOfReplicas, ReplicasRemainder, []), {ISR, Remainder} = isrs(NumberOfISR, ISRRemainder, []), partitions(N-1, Remainder, [#{error_code => kafe_error:code(ErrorCode), id => Id, leader => Leader, replicas => Replicas, isr => ISR} | Acc]). isrs(0, Remainder, Acc) -> {Acc, Remainder}; isrs(N, <>, Acc) -> isrs(N-1, Remainder, [InSyncReplica | Acc]). replicas(0, Remainder, Acc) -> {Acc, Remainder}; replicas(N, <>, Acc) -> replicas(N-1, Remainder, [Replica | Acc]).