Current section
Files
Jump to
Current section
Files
src/platybelodon_frame.erl
-module(platybelodon_frame).
-include("stomp_frame.hrl").
-export([format/1, parse/1, parse_all/1, headers_to_map/1]).
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
% Frame construction (format)
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
format([]) ->
<<"\n">>;
format(Frames) when is_list(Frames) ->
iolist_to_binary([format(Frame) || Frame <- Frames]);
format(#stomp_frame{command = Command, headers = Headers, body = Body}) ->
AdditionalHeaders = size_header_for(Body),
[command_atom_to_bin(Command), <<"\n">>, format_headers(AdditionalHeaders ++ Headers), <<"\n">>, Body, <<0>>].
command_atom_to_bin(Command) when is_atom(Command) ->
case Command of
stomp -> << "STOMP" >>;
connect -> << "CONNECT" >>;
disconnect -> << "DISCONNECT" >>;
ack -> << "ACK" >>;
nack -> << "NACK" >>;
'begin' -> << "BEGIN" >>;
commit -> << "COMMIT" >>;
abort -> << "ABORT" >>;
send -> << "SEND" >>;
subscribe -> <<"SUBSCRIBE">>;
unsubscribe -> <<"UNSUBSCRIBE">>;
connected -> << "CONNECTED" >>;
error -> << "ERROR" >>;
message -> << "MESSAGE" >>;
receipt -> << "RECEIPT" >>
end.
size_header_for(<<>>) ->
[];
size_header_for(Body) ->
[{<<"content-length">>, integer_to_binary(byte_size(Body))}].
format_headers(Headers) ->
[format_header({value_encode(K), value_encode(V)}) || {K,V} <- Headers].
format_header({Key, Value}) ->
<<Key/binary, ":", Value/binary, "\n">>.
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
% frame parsing (parse)
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
parse({command, undefined, undefined, <<"\n", Rest/binary>>}) ->
parse({command, undefined, undefined, Rest});
parse({command, undefined, undefined, <<"\r\n", Rest/binary>>}) ->
parse({command, undefined, undefined, Rest});
parse({command, undefined, undefined, Str}) ->
case binary:split(Str, [<<"\r\n">>, <<"\n">>]) of
% End of command phase, now in headers phase
[Command, Rest] -> parse({headers, #stomp_frame{command=command_bin_to_atom(Command)}, undefined, Rest});
% Too long now, something's wrong
[Rest] when byte_size(Rest) >= 14 ->
{error, command, <<"Frame too long without command match">>};
% Line unterminated, abandoned in command phase
[Rest] -> {command, undefined, undefined, Rest}
end;
parse({headers, Frame = #stomp_frame{headers=Headers}, Length, Str}) ->
case {Length, binary:split(Str, [<<"\r\n">>, <<"\n">>])} of
% Empty line means headers are over
{_, [<<>>, Rest]} ->
parse({body, Frame#stomp_frame{headers=lists:reverse(Headers)}, Length, Rest});
% Non-empty line (with content length)
{undefined, [HeaderLine = <<"content-length:", _Rest/binary>>, Rest]} ->
ContentLengthHeader = {<<"content-length">>, ContentLengthStr} = parse_header(HeaderLine),
NewLength = binary_to_integer(ContentLengthStr),
parse({headers, Frame#stomp_frame{headers=[ContentLengthHeader | Headers]}, NewLength, Rest});
% Non-empty line (no content length)
{_, [HeaderLine, Rest]} ->
parse({headers, Frame#stomp_frame{headers=[parse_header(HeaderLine) | Headers]}, Length, Rest});
% Non-line - abandoned in headers phase
{_, [Rest]} -> {headers, Frame, Length, Rest}
end;
parse({body, Frame = #stomp_frame{body=Body}, undefined, Str}) ->
case binary:split(Str, <<"\0">>) of
% There might be more that can be done, not our problem now!
[NewBody, Rest] -> {done, Frame#stomp_frame{body = <<Body/binary, NewBody/binary>>}, undefined, Rest};
[NewBody] -> {body, Frame#stomp_frame{body = <<Body/binary, NewBody/binary>>}, undefined, <<>>}
end;
parse({body, Frame = #stomp_frame{}, Length, Str}) when byte_size(Str) >= Length + 1 ->
<<NewBody:Length/binary, 0, Rest/binary>> = Str,
{done, Frame#stomp_frame{body = NewBody}, undefined, Rest};
parse({body, Frame = #stomp_frame{}, Length, Str}) ->
{body, Frame, Length, Str}.
parse_header(HeaderLine) ->
[Key, Value] = binary:split(HeaderLine, <<":">>),
{value_decode(Key), value_decode(Value)}.
command_bin_to_atom(Command) ->
case Command of
<< "STOMP" >> -> stomp;
<< "CONNECT" >> -> connect;
<< "DISCONNECT" >> -> disconnect;
<< "ACK" >> -> ack;
<< "NACK" >> -> nack;
<< "BEGIN" >> -> 'begin';
<< "COMMIT" >> -> commit;
<< "ABORT" >> -> abort;
<< "SEND" >> -> send;
<<"SUBSCRIBE">> -> subscribe;
<<"UNSUBSCRIBE">> -> unsubscribe;
<< "CONNECTED" >> -> connected;
<< "ERROR" >> -> error;
<< "MESSAGE" >> -> message;
<< "RECEIPT" >> -> receipt
end.
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
% Parse all
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
parse_all(ParseState = {_, _, _, _}) ->
parse_all(parse(ParseState), []).
parse_all(ParseState = {error, _Stage, _Message}, Accumulator) ->
{ParseState, lists:reverse(Accumulator)};
parse_all(ParseState = {Command, _, _, _}, Accumulator) when Command =/= done ->
{ParseState, lists:reverse(Accumulator)};
parse_all({done, Frame, _, Remainder}, Accumulator) ->
parse_all(parse({command, undefined, undefined, Remainder}), [Frame | Accumulator]).
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
% utils
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
headers_to_map(Headers) ->
% Reversed because in Erlang, last takes precedence, but in STOMP, first does.
maps:from_list(lists:reverse(Headers)).
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
% value encode/decode
% http://stomp.github.io/stomp-specification-1.2.html#Value_Encoding
%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%%
value_encode(Str) ->
value_encode(Str, []).
value_encode(<<":", Rest/binary>>, Acc) ->
value_encode(Rest, ["\\c" | Acc]);
value_encode(<<"\n", Rest/binary>>, Acc) ->
value_encode(Rest, ["\\n" | Acc]);
value_encode(<<"\r", Rest/binary>>, Acc) ->
value_encode(Rest, ["\\r" | Acc]);
value_encode(<<"\\", Rest/binary>>, Acc) ->
value_encode(Rest, ["\\\\" | Acc]);
value_encode(<<Ch, Rest/binary>>, Acc) ->
value_encode(Rest, [Ch | Acc]);
value_encode(<<>>, Acc) ->
list_to_binary(lists:reverse(Acc)).
value_decode(Str) ->
value_decode(Str, []).
value_decode(<<"\\c", Rest/binary>>, Acc) ->
value_decode(Rest, [":" | Acc]);
value_decode(<<"\\n", Rest/binary>>, Acc) ->
value_decode(Rest, ["\n" | Acc]);
value_decode(<<"\\r", Rest/binary>>, Acc) ->
value_decode(Rest, ["\r" | Acc]);
value_decode(<<"\\\\", Rest/binary>>, Acc) ->
value_decode(Rest, ["\\" | Acc]);
value_decode(<<Ch, Rest/binary>>, Acc) ->
value_decode(Rest, [Ch | Acc]);
value_decode(<<>>, Acc) ->
list_to_binary(lists:reverse(Acc)).