-module(lake_messages). -export([parse/1, message_to_correlation_id/1]). -export([peer_properties/2, sasl_handshake/1, sasl_authenticate/4, open/2, close/3, close_response/2, tune/2, declare_publisher/4, publish/2, query_publisher_sequence/3, delete_publisher/2, credit/2, create/3, delete/2, subscribe/6, store_offset/3, query_offset/3, unsubscribe/2, metadata/2, heartbeat/0, route/3, partitions/2, stream_stats/2, command_versions/0, exchange_command_versions/1, consumer_update_response/3]). -export([chunk_to_messages/1]). -include("response_codes.hrl"). -define(REQUEST, 0). -define(RESPONSE, 1). -define(VERSION_1, 1). -define(VERSION_2, 2). -define(MIN_COMMAND_INDEX, 1). -define(DECLARE_PUBLISHER, 1). -define(PUBLISH, 2). -define(PUBLISH_CONFIRM, 3). -define(PUBLISH_ERROR, 4). -define(QUERY_PUBLISHER_SEQUENCE, 5). -define(DELETE_PUBLISHER, 6). -define(SUBSCRIBE, 7). -define(DELIVER, 8). -define(CREDIT, 9). -define(STORE_OFFSET, 10). -define(QUERY_OFFSET, 11). -define(UNSUBSCRIBE, 12). -define(CREATE, 13). -define(DELETE, 14). -define(METADATA, 15). -define(METADATA_UPDATE, 16). -define(PEER_PROPERTIES, 17). -define(SASL_HANDSHAKE, 18). -define(SASL_AUTHENTICATE, 19). -define(TUNE, 20). -define(OPEN, 21). -define(CLOSE, 22). -define(HEARTBEAT, 23). -define(ROUTE, 24). -define(PARTITIONS, 25). -define(CONSUMER_UPDATE, 26). -define(EXCHANGE_COMMAND_VERSIONS, 27). -define(STREAM_STATS, 28). -define(MAX_COMMAND_INDEX, 28). % Increase this when adding a new command -define(OFFSET_TYPE_NONE, 0). -define(OFFSET_TYPE_FIRST, 1). -define(OFFSET_TYPE_LAST, 2). -define(OFFSET_TYPE_NEXT, 3). -define(OFFSET_TYPE_OFFSET, 4). -define(OFFSET_TYPE_TIMESTAMP, 5). parse(<>) -> {declare_publisher_response, Corr, ResponseCode}; parse(<>) -> {publish_confirm, PublisherId, PublishingIdCount, parse_list_of_longs(PublishingIds, [])}; parse(<>) -> ErrorById = [{Id, Code} || <> <= Details], {publish_error, PublisherId, PublishingIdCount, ErrorById}; parse(<>) -> {query_publisher_sequence_response, Corr, ResponseCode, Sequence}; parse(<>) -> {delete_publisher_response, Corr, ResponseCode}; parse(<>) -> {credit_response, SubscriptionId, ResponseCode}; parse(<>) -> {subscribe_response, Corr, ResponseCode}; parse(<>) -> {deliver, SubscriptionId, OsirisChunk}; parse(<>) -> {deliver_v2, SubscriptionId, CommittedChunkId, OsirisChunk}; parse(<>) -> {query_offset_response, Corr, ResponseCode, Offset}; parse(<>) -> {unsubscribe_response, Corr, ResponseCode}; parse(<>) -> {create_response, Corr, ResponseCode}; parse(<>) -> {delete_response, Corr, ResponseCode}; parse(<>) -> <> = Metadata0, {Endpoints, Rest} = parse_endpoints(NumEndpoints, EndpointsAndMetadata), StreamsMetadata = parse_streams_metadata(Rest), {metadata_response, Corr, Endpoints, StreamsMetadata}; parse(<>) -> {metadata_update, ResponseCode, Stream}; parse(<>) -> {peer_properties_response, Corr, ResponseCode, parse_map_binary_values(PeerProperties)}; parse(<>) -> {sasl_handshake_response, Corr, ResponseCode, parse_list_of_strings(Mechanisms)}; parse(<>) -> {sasl_authenticate_response, Corr, ResponseCode, SaslOpaque}; parse(<>) -> {tune, FrameMax, Heartbeat}; parse(<>) -> {open_response, Corr, ?RESPONSE_OK, parse_map_binary_values(ConnectionProperties)}; parse(<>) -> {close_response, Corr, ResponseCode}; parse(<>) -> {close, Corr, ResponseCode, Reason}; parse(<>) -> {open_response, Corr, ResponseCode}; parse(<>) -> {heartbeat}; parse(<>) -> {route_response, Corr, ?RESPONSE_OK, parse_list_of_strings(Streams)}; parse(<>) -> {route_response, Corr, ResponseCode, []}; parse(<>) -> {partitions_response, Corr, ?RESPONSE_OK, parse_list_of_strings(Streams)}; parse(<>) -> {partitions_response, Corr, ResponseCode, []}; parse(<>) -> {consumer_update, Corr, SubscriptionId, Active == <<"1">>}; parse(<>) -> {exchange_command_versions_response, Corr, ?RESPONSE_OK, parse_command_versions(CommandVersions)}; parse(<>) -> {exchange_command_versions_response, Corr, ResponseCode, []}; parse(<>) -> {stream_stats_response, Corr, ?RESPONSE_OK, parse_stream_stats(StreamStats)}; parse(<>) -> {stream_stats_response, Corr, ResponseCode, []}; parse(Unknown) -> {error, {unknown, Unknown}}. -define(OSIRIS_MAGIC, 5). -define(OSIRIS_VERSION_1, 0). -define(OSIRIS_CHUNK_TYPE_USER, 0). chunk_to_messages(<>) -> case erlang:crc32(Data) of DataCRC -> %% FIXME why isn't Trailer of length TrailerLength? Messages = parse_data(Data, []), Info = #{chunk_id => ChunkId, number_of_entries => NumberOfEntries, number_of_records => NumberOfRecords, timestamp => Timestamp, epoch => Epoch}, {ok, {Messages, Info}}; MismatchingCRC -> {error, {crc_mismatch, [{expected, DataCRC}, {received, MismatchingCRC}]}} end; chunk_to_messages(Other) -> {error, {invalid_osiris_chunk, Other}}. parse_data(<<>>, Acc) -> lists:reverse(Acc); parse_data(<<0:1, Size:31, Data:Size/binary, Rest/binary>>, Acc) -> parse_data(Rest, [Data | Acc]). peer_properties(CorrelationId, Properties) -> PropertiesCount = length(Properties), EncodedProperties = encode_keywords(Properties), <>. sasl_handshake(CorrelationId) -> <>. tune(FrameMax, Heartbeat) -> <>. sasl_authenticate(CorrelationId, Mechanism = <<"PLAIN">>, User, Password) -> MechanismSize = byte_size(Mechanism), Fragment = <<0:8, User/binary, 0:8, Password/binary>>, FragmentSize = byte_size(Fragment), <>. open(CorrelationId, Vhost) -> VhostSize = byte_size(Vhost), <>. close(CorrelationId, ResponseCode, Reason) -> ReasonSize = byte_size(Reason), <>. close_response(CorrelationId, ResponseCode) -> <>. declare_publisher(CorrelationId, Stream, PublisherId, PublisherReference) -> StreamSize = byte_size(Stream), PublisherReferenceSize = byte_size(PublisherReference), <>. publish(PublisherId, Messages) -> MessageCount = length(Messages), EncodedMessages = encode_messages(Messages), <>. query_publisher_sequence(CorrelationId, PublisherReference, Stream) -> PublisherReferenceSize = byte_size(PublisherReference), StreamSize = byte_size(Stream), <>. delete_publisher(CorrelationId, PublisherId) -> <>. credit(SubscriptionId, Credit) -> <>. encode_messages(Messages) -> encode_messages(Messages, <<>>). encode_messages([], Acc) -> Acc; encode_messages([{Id, Message} | Rest], Acc) when is_integer(Id), is_binary(Message) -> Size = byte_size(Message), encode_messages(Rest, <>). to_offset_bin(OffsetSpecification) -> case OffsetSpecification of %% FIXME why no OFFSET_TYPE_NONE? first -> <>; last -> <>; next -> <>; {offset, Offset} -> <>; {timestamp, Offset} -> <> end. subscribe(CorrelationId, Stream, SubscriptionId, OffsetSpecification, Credit, Properties) -> StreamSize = byte_size(Stream), OffsetBin = to_offset_bin(OffsetSpecification), EncodedProperties = encode_keywords(Properties), PropertiesCount = length(Properties), <>. store_offset(PublisherReference, Stream, Offset) -> PublisherReferenceSize = byte_size(PublisherReference), StreamSize = byte_size(Stream), <>. query_offset(CorrelationId, PublisherReference, Stream) -> PublisherReferenceSize = byte_size(PublisherReference), StreamSize = byte_size(Stream), <>. unsubscribe(CorrelationId, SubscriptionId) -> <>. create(CorrelationId, Stream, Arguments0) -> StreamSize = byte_size(Stream), ArgumentsCount = length(Arguments0), Arguments = encode_keywords(Arguments0), <>. delete(CorrelationId, Stream) -> StreamSize = byte_size(Stream), <>. metadata(CorrelationId, Streams) -> StreamsCount = length(Streams), EncodedStreams = encode_list_of_strings(Streams), <>. heartbeat() -> <>. route(CorrelationId, RoutingKey, SuperStream) -> RoutingKeySize = byte_size(RoutingKey), SuperStreamSize = byte_size(SuperStream), <>. partitions(CorrelationId, SuperStream) -> SuperStreamSize = byte_size(SuperStream), <>. stream_stats(CorrelationId, Stream) -> StreamSize = byte_size(Stream), <>. command_versions() -> maps:from_list([{Command, version_range_by_id(Command)} || Command <- lists:seq(?MIN_COMMAND_INDEX, ?MAX_COMMAND_INDEX)]). version_range_by_id(?DELIVER) -> {?VERSION_1, ?VERSION_2}; version_range_by_id(_) -> {?VERSION_1, ?VERSION_1}. exchange_command_versions(CorrelationId) -> CommandsBinary = << <> || {Command, {MinVersion, MaxVersion}} <- maps:to_list(command_versions()) >>, CommandsCount = maps:size(command_versions()), <>. consumer_update_response(CorrelationId, ResponseCode, OffsetSpecification) -> OffsetBin = to_offset_bin(OffsetSpecification), <>. parse_map_binary_values(Bin) when is_binary(Bin) -> parse_map_binary_values(Bin, #{}). parse_map_binary_values(<<>>, Acc) -> Acc; parse_map_binary_values(<>, Acc) -> parse_map_binary_values(Rest, Acc#{Key => Value}). parse_command_versions(Bin) when is_binary(Bin) -> parse_command_versions(Bin, #{}). parse_command_versions(<<>>, Acc) -> Acc; parse_command_versions(<>, Acc) -> parse_command_versions(Rest, Acc#{Key => {MinVersion, MaxVersion}}). parse_stream_stats(Bin) when is_binary(Bin) -> parse_stream_stats(Bin, #{}). parse_stream_stats(<<>>, Acc) -> Acc; parse_stream_stats(<>, Acc) -> parse_stream_stats(Rest, Acc#{Key => Value}). parse_list_of_strings(Bin) -> parse_list_of_strings(Bin, []). parse_list_of_strings(<<>>, Acc) -> lists:reverse(Acc); parse_list_of_strings(<>, Acc) -> parse_list_of_strings(Rest, [String | Acc]). parse_list_of_longs(<<>>, Acc) -> lists:reverse(Acc); parse_list_of_longs(<>, Acc) -> parse_list_of_longs(Rest, [Id | Acc]). encode_keywords(Keywords) -> encode_keywords(Keywords, <<>>). encode_keywords([], Acc) -> Acc; encode_keywords([{Key, Value} | Rest], Acc) -> SizeKey = byte_size(Key), SizeValue = byte_size(Value), encode_keywords(Rest, <>). encode_list_of_strings(List) -> encode_list_of_strings(List, <<>>). encode_list_of_strings([], Acc) -> Acc; encode_list_of_strings([Bin | Rest], Acc) -> Size = byte_size(Bin), encode_list_of_strings(Rest, <>). parse_endpoints(NumEndpoints, Bin) -> parse_endpoints(NumEndpoints, Bin, []). parse_endpoints(0, Rest, Endpoints) -> {Endpoints, Rest}; parse_endpoints(Cnt, <>, Acc) -> EndpointMetadata = #{index => Index, host => Host, port => Port}, parse_endpoints(Cnt - 1, Rest, [EndpointMetadata | Acc]). parse_streams_metadata(<>) -> parse_streams_metadata(NumStreams, Meta, []). parse_streams_metadata(0, <<>>, Acc) -> Acc; %% Error case parse_streams_metadata(NumStreams, <>, Acc) -> StreamMetadata = #{stream => Stream, code => Code}, parse_streams_metadata(NumStreams - 1, Rest, [StreamMetadata | Acc]); %% Success case parse_streams_metadata(NumStreams, <>, Acc) -> ReplicasSize = ReplicasCount * 2, <> = ReplicasAndRest, Replicas = [EndpointIndex || <> <= ReplicasBin], StreamMetadata = #{stream => Stream, leader_index => LeaderIndex, replicas_count => ReplicasCount, replicas => Replicas}, parse_streams_metadata(NumStreams - 1, Rest, [StreamMetadata | Acc]). message_to_correlation_id({declare_publisher_response, Corr, _}) -> {ok, Corr}; message_to_correlation_id({query_publisher_sequence, Corr, _, _}) -> {ok, Corr}; message_to_correlation_id({delete_publisher_response, Corr, _}) -> {ok, Corr}; message_to_correlation_id({subscribe_response, Corr, _}) -> {ok, Corr}; message_to_correlation_id({query_offset_response, Corr, _, _}) -> {ok, Corr}; message_to_correlation_id({unsubscribe_response, Corr, _}) -> {ok, Corr}; message_to_correlation_id({create_response, Corr, _}) -> {ok, Corr}; message_to_correlation_id({delete_response, Corr, _}) -> {ok, Corr}; message_to_correlation_id({metadata_response, Corr, _, _}) -> {ok, Corr}; message_to_correlation_id({route_response, Corr, _, _}) -> {ok, Corr}; message_to_correlation_id(Other) -> {error, {no_correlation_id, Other}}.