Packages

kafka_ex_tc

0.12.1-21
0.13.0 0.12.1 0.12.1-51-1 0.12.1-50-3 0.12.1-50-2 0.12.1-50-1 0.12.1-49-2 0.12.1-49-1 0.12.1-49-0 0.12.1-48-2 0.12.1-48-1 0.12.1-47-1 0.12.1-46-2 0.12.1-45-2 0.12.1-45-1 0.12.1-44-1 0.12.1-43-1 0.12.1-42-8 0.12.1-42-6 0.12.1-42-5 0.12.1-42-4 0.12.1-42-1 0.12.1-41-2 0.12.1-41-1 0.12.1-40-9 0.12.1-40-8 0.12.1-40-7 0.12.1-40-6 0.12.1-40-5 0.12.1-40-4 0.12.1-40-3 0.12.1-40-2 0.12.1-40-10 0.12.1-40-1 0.12.1-39-2 0.12.1-39-1 0.12.1-36-3 0.12.1-36-2 0.12.1-36-1 0.12.1-34-b 0.12.1-34-a 0.12.1-33-1 0.12.1-32-a1 0.12.1-29-s-3 0.12.1-29-s-2 0.12.1-29-s-1 0.12.1-28-rumble-1 0.12.1-27-2 0.12.1-27-1 0.12.1-26-rumble-2 0.12.1-26-rumble-1 0.12.1-26-dev-1 0.12.1-26-dev.1 0.12.1-26-3 0.12.1-26-2 0.12.1-26-1 0.12.1-25-s-8 0.12.1-25-s-7 0.12.1-25-s-6 0.12.1-25-s-5 0.12.1-25-s-4 0.12.1-25-s-3 0.12.1-25-s-2 0.12.1-25-s-1 0.12.1-25-dev-22 0.12.1-25-dev 0.12.1-24-2 0.12.1-24-1 0.12.1-22-poc-b-7 0.12.1-20-poc-b-6 0.12.1-20-poc-b-5 0.12.1-20-poc-b-4 0.12.1-20-poc-b-3 0.12.1-19-poc-b-3 0.12.1-19-poc-b-2 0.12.1-19-poc-b-1 0.12.1-19-poc-7 0.12.1-19-poc-6 0.12.1-19-poc-4 0.12.1-19-poc-3 0.12.1-19-poc-2 0.12.1-19-poc-1 0.12.1-19-poc 0.12.1-11-debug-1 0.12.1-11-debug 0.12.1-39 0.12.1-38.a 0.12.1-38 0.12.1-37.c 0.12.1-37.b 0.12.1-37.a 0.12.1-37.1 0.12.1-37 0.12.1-35.b 0.12.1-35.a 0.12.1-34 0.12.1-33 0.12.1-32 0.12.1-31 0.12.1-30 0.12.1-29 0.12.1-28 0.12.1-27 0.12.1-26 0.12.1-25 0.12.1-24 0.12.1-23 0.12.1-22.1 0.12.1-22 0.12.1-21 0.12.1-20 0.12.1-19 0.12.1-18 0.12.1-17 0.12.1-16 0.12.1-15 0.12.1-14 0.12.1-13 0.12.1-12 0.12.1-11 0.12.1-10 0.12.1-9 0.12.1-8 0.12.1-7 0.12.1-6 0.12.1-5 0.12.1-4 0.12.1-3 0.12.1-2 0.12.1-1

Kafka client for Elixir/Erlang.

Current section

Files

Jump to
kafka_ex_tc lib kafka_ex protocol common.ex
Raw

lib/kafka_ex/protocol/common.ex

defmodule KafkaEx.Protocol.Common do
@moduledoc """
A collection of common request generation and response parsing functions for the
Kafka wire protocol.
"""
@doc """
Generate the wire representation for a list of topics.
"""
def topic_data([]), do: <<>>
def topic_data([topic | topics]) do
<<byte_size(topic)::16-signed, topic::binary>> <> topic_data(topics)
end
def parse_topics(0, _, _), do: []
def parse_topics(
topics_size,
<<topic_size::16-signed, topic::size(topic_size)-binary,
partitions_size::32-signed, rest::binary>>,
mod
) do
struct_module = Module.concat(mod, Response)
{partitions, topics_data} =
mod.parse_partitions(partitions_size, rest, [], topic)
[
%{
__struct__: struct_module,
topic: topic,
partitions: partitions
}
| parse_topics(topics_size - 1, topics_data, mod)
]
end
def read_array(0, data_after_array, _read_one) do
{[], data_after_array}
end
def read_array(num_items, data, read_one) do
{item, rest} = read_one.(data)
{items, data_after_array} = read_array(num_items - 1, rest, read_one)
{[item | items], data_after_array}
end
@spec encode_nullable_string(String.t()) :: binary
def encode_nullable_string(text) do
case text do
nil -> <<-1::16-signed>>
_ -> encode_string(text)
end
end
@spec encode_string(String.t()) :: binary
def encode_string(text) do
<<byte_size(text)::16-signed, text::binary>>
end
def map_encode(elems, function) do
if nil == elems or [] == elems do
<<0::32-signed>>
else
<<length(elems)::32-signed>> <>
(elems
|> Enum.map(function)
|> Enum.reduce(&(&1 <> &2)))
end
end
end