-module(dproto_tcp). -include_lib("mmath/include/mmath.hrl"). -include("dproto.hrl"). -export([ encode_metrics/1, decode_metrics/1, decode_ot/1, encode_buckets/1, decode_buckets/1, encode_bucket_info/1, decode_bucket_info/1, encode/1, decode/1, encode_get_reply/1, encode_get_stream/1, decode_get_reply/1, decode_get_stream/2, decode_stream/1, decode_batch/1 ]). -ignore_xref([ encode_metrics/1, decode_metrics/1, encode_buckets/1, decode_buckets/1, encode_bucket_info/1, decode_bucket_info/1, encode/1, decode/1, encode_get_reply/1, encode_get_stream/1, decode_get_reply/1, decode_get_stream/2, decode_stream/1, decode_batch/1 ]). -export_type([ttl/0, read_opts/0, read_repair_opt/0, read_r_opt/0, bucket_info/0, tcp_message/0, batch_message/0, stream_message/0]). -ifdef(TEST). -export([encode_aggr/1, decode_aggr/1]). -endif. -type ttl() :: pos_integer() | infinity. -type read_repair_opt() :: {rr, default} | {rr, off} | {rr, on}. -type read_r_opt() :: {r, n} | {r, default} | {r, 1..254}. -type aggr() :: {binary(), pos_integer()}. -type read_aggr_opt() :: {aggr, aggr()}. -type read_opts() :: [read_repair_opt() | read_r_opt() | read_aggr_opt()]. -type bucket_info() :: #{ resolution => pos_integer(), ppf => pos_integer(), grace => non_neg_integer(), ttl => ttl() }. -type get_stream_element() :: {more, binary()} | {done, binary()}. -type stream_message() :: flush | incomplete | {batch, Time :: non_neg_integer()} | {stream, Metric :: binary(), Time :: non_neg_integer(), Points :: binary()}. -type batch_message() :: incomplete | batch_end | {batch, Metric :: binary(), Points :: binary()}. %% Messages shorthands that can be encoded but will never be decoded. -type tcp_encode_message() :: {get, Bucket :: binary(), Metric :: binary(), Time :: pos_integer(), Count :: pos_integer()}. -type otids() :: {undefined | pos_integer(), undefined | pos_integer()} | undefined. -type tcp_message() :: {ot, otids()} | {ot, otids(), tcp_message()} | buckets | {ttl, Bucket :: binary(), TTL :: ttl()} | {list, Bucket :: binary()} | {list, Bucket :: binary(), Prefix :: binary()} | {info, Bucket :: binary()} | {delete, Bucket :: binary()} | {events, [{pos_integer(), term()}]} | events_end | {events, Bucket :: binary(), [{pos_integer(), term()}]} | {get_events, Bucket :: binary(), Start :: pos_integer(), End :: pos_integer()} | {get_events, Bucket :: binary(), Start :: pos_integer(), End :: pos_integer(), Filter :: jsxd:filter_filters()} | {get, Bucket :: binary(), Metric :: binary(), Time :: pos_integer(), Count :: pos_integer(), Opts :: read_opts()} | {stream, Bucket :: binary(), Delay :: pos_integer()} | {error, Message :: binary()}. -type encoded_metric() :: <<_:?METRICS_SS, _:_*8>>. -type encoded_bucket() :: <<_:?BUCKETS_SS, _:_*8>>. %%-------------------------------------------------------------------- %% @doc %% Encode a list of metrics to its binary form for sending it over %% the wire. %% %% @end %%-------------------------------------------------------------------- -spec encode_metrics([dproto:metric()]) -> encoded_metric(). encode_metrics(Metrics) when is_list(Metrics) -> Data = << <<(byte_size(Metric)):?METRIC_SS/?SIZE_TYPE, Metric/binary>> || Metric <- Metrics >>, <<(byte_size(Data)):?METRICS_SS/?SIZE_TYPE, Data/binary>>. %%-------------------------------------------------------------------- %% @doc %% Decodes the binary representation of a metric list to its list %% representation. %% %% Node this does not recursively decode the metrics! %% %% @end %%-------------------------------------------------------------------- -spec decode_metrics(encoded_metric()) -> [dproto:metric()]. decode_metrics(<<_Size:?METRICS_SS/?SIZE_TYPE, Metrics:_Size/binary>>) -> [ Metric || <<_S:?METRIC_SS/?SIZE_TYPE, Metric:_S/binary>> <= Metrics]. %%-------------------------------------------------------------------- %% @doc %% Encode a list of buckets to its binary form for sending it over %% the wire. %% %% @end %%-------------------------------------------------------------------- -spec encode_buckets([dproto:metric()]) -> encoded_bucket(). encode_buckets(Buckets) when is_list(Buckets) -> Data = << <<(byte_size(Bucket)):?BUCKET_SS/?SIZE_TYPE, Bucket/binary>> || Bucket <- Buckets >>, <<(byte_size(Data)):?BUCKETS_SS/?SIZE_TYPE, Data/binary>>. %%-------------------------------------------------------------------- %% @doc %% Decodes the binary representation of a bucket list to its list %% representation. %% %% @end %%-------------------------------------------------------------------- -spec decode_buckets(encoded_bucket()) -> [dproto:bucket()]. decode_buckets(<<_Size:?BUCKETS_SS/?SIZE_TYPE, Buckets:_Size/binary>>) -> [ Bucket || <<_S:?BUCKET_SS/?SIZE_TYPE, Bucket:_S/binary>> <= Buckets]. %%-------------------------------------------------------------------- %% @doc %% Encodes bucket properties such as PPF, Resolution and TTL into a binary %% form for transmission over the wire. %% %% @end %%-------------------------------------------------------------------- -spec encode_bucket_info(bucket_info()) -> <<_:192>> | <<_:256>>. encode_bucket_info(#{ resolution := Resolution, ppf := PPF, grace := Grace, ttl := infinity }) when is_integer(Resolution), Resolution > 0, is_integer(PPF), PPF > 0, is_integer(Grace), Grace >= 0 -> <>; encode_bucket_info(#{ resolution := Resolution, ppf := PPF, grace := Grace, ttl := TTL }) when is_integer(Resolution), Resolution > 0, is_integer(PPF), PPF > 0, is_integer(Grace), Grace >= 0, is_integer(TTL), TTL > 0 -> <>. %%-------------------------------------------------------------------- %% @doc %% Decodes bucket properties from the wire protocol. %% %% @end %%-------------------------------------------------------------------- -spec decode_bucket_info(<<_:192, _:_*64>>) -> bucket_info(). decode_bucket_info(<>) -> #{ resolution => Resolution, ppf => PPF, grace => Grace, ttl => infinity }; decode_bucket_info(<>) -> #{ resolution => Resolution, ppf => PPF, grace => Grace, ttl => TTL }. %%-------------------------------------------------------------------- %% @doc %% Encodes a message for the binary protocol. %% %% @end %%-------------------------------------------------------------------- -spec encode(tcp_encode_message() | tcp_message() | stream_message() | batch_message()) -> binary(). encode({ot, TIDs}) -> <> ; encode({ot, TIDs, Message}) -> <<(encode({ot, TIDs}))/binary, (encode(Message))/binary>>; encode(buckets) -> <>; %% @doc %% Encodes the TTL for a bucket. %% Note that a zero value is substituted in place of `infinity'. %% %% @end encode({ttl, Bucket, infinity}) -> <>; encode({ttl, Bucket, TTL}) when is_binary(Bucket), byte_size(Bucket) > 0, is_integer(TTL), TTL > 0 -> <>; encode({list, Bucket}) when is_binary(Bucket), byte_size(Bucket) > 0 -> <>; encode({list, Bucket, Prefix}) when is_binary(Bucket), byte_size(Bucket) > 0, is_binary(Prefix), byte_size(Prefix) > 0 -> <>; encode({info, Bucket}) when is_binary(Bucket), byte_size(Bucket) > 0 -> <>; encode({add, Bucket, Resolution, PPF, TTL}) when is_binary(Bucket), byte_size(Bucket) > 0, is_integer(Resolution), Resolution > 0, is_integer(PPF), PPF > 0, is_integer(TTL), TTL >= 0 -> <>; encode({delete, Bucket}) when is_binary(Bucket), byte_size(Bucket) > 0 -> <>; encode({get, Bucket, Metric, Time, Count}) -> encode({get, Bucket, Metric, Time, Count, []}); encode({get, Bucket, Metric, Time, Count, Opts}) when is_binary(Bucket), byte_size(Bucket) > 0, is_binary(Metric), byte_size(Metric) > 0, is_integer(Time), Time >= 0, (Time band 16#FFFFFFFFFFFFFFFF) =:= Time, %% We only want positive numbers < 32 bit is_integer(Count), Count > 0, (Count band 16#FFFFFFFF) =:= Count, is_list(Opts) -> RROpt = proplists:get_value(rr, Opts, default), ROpt = proplists:get_value(r, Opts, default), Res = <>, case proplists:get_value(aggr, Opts) of undefined -> Res; Aggr -> AggrB = encode_aggr(Aggr), <> end; encode({stream, Bucket, Delay}) when is_binary(Bucket), byte_size(Bucket) > 0, is_integer(Delay), Delay > 0, Delay < 256-> <>; encode({stream, Metric, Time, Points}) when is_binary(Metric), byte_size(Metric) > 0, is_binary(Points), byte_size(Points) rem ?DATA_SIZE == 0, is_integer(Time), Time >= 0-> <>; encode({batch, Time}) when is_integer(Time), Time >= 0 -> <>; encode({batch, Metric, Point}) when is_binary(Metric), byte_size(Metric) > 0, is_binary(Point), byte_size(Point) == ?DATA_SIZE -> <<(byte_size(Metric)):?METRIC_SS/?SIZE_TYPE, Metric/binary, Point:?DATA_SIZE/binary>>; encode({batch, Metric, Point}) when is_binary(Metric), byte_size(Metric) > 0, is_integer(Point) -> PointB = mmath_bin:from_list([Point]), <<(byte_size(Metric)):?METRIC_SS/?SIZE_TYPE, Metric/binary, PointB:?DATA_SIZE/binary>>; encode(batch_end) -> <<0:?METRIC_SS/?SIZE_TYPE>>; encode(flush) -> <>; encode({events, Bucket, Events}) -> EventsB = encode_events(Events), <>; encode({get_events, Bucket, Start, End}) -> <>; encode({get_events, Bucket, Start, End, Filter}) -> <>; encode({events, Events}) -> EventsB = encode_events(Events), <>; encode(events_end) -> <>; encode({error, Message}) -> <>. -spec encode_events([{pos_integer(), term()}]) -> binary(). encode_events(Es) -> {ok, B} = snappiest:compress(<< <<(encode_event(E))/binary>> || E <- Es >>), %% Damn you dailyzer! true = is_binary(B), B. -spec encode_event({pos_integer(), term()}) -> <<_:64, _:_*8>>. encode_event({T, E}) when T > 0, is_integer(T) -> B = term_to_binary(E), <>. %%-------------------------------------------------------------------- %% @doc %% Decodes a normal TCP message from the wire protocol. %% %% @end %%-------------------------------------------------------------------- -spec decode_ot(binary()) -> {ot, otids(), binary()}. decode_ot(<>) -> {ot, decode_traceids(TraceID, ParentID), Body}; decode_ot(Body) when is_binary(Body) -> {ot, undefined, Body}. -spec decode(binary()) -> tcp_message(). decode(<>) -> {ot, decode_traceids(TraceID, ParentID), decode(R)}; decode(<>) -> buckets; %% @doc %% Decodes the TTL for a bucket. %% Note that a zero value is interpreted to mean `infinity'. %% %% @end decode(<>) -> case TTL of 0 -> {ttl, Bucket, infinity}; _ when TTL > 0 -> {ttl, Bucket, TTL} end; decode(<>) -> {list, Bucket}; decode(<>) -> {list, Bucket, Prefix}; decode(<>) -> {info, Bucket}; decode(<>) -> {add, Bucket, Resolution, PPF, TTL}; decode(<>) -> {delete, Bucket}; decode(<>) -> Opts = [{r, default}, {rr, default}], {get, Bucket, Metric, Time, Count, Opts}; decode(<>) -> Opts = [{r, decode_r(R)}, {rr, decode_rr(RR)}], {get, Bucket, Metric, Time, Count, Opts}; decode(<>) -> Opts = [{r, decode_r(R)}, {rr, decode_rr(RR)}, {aggr, decode_aggr(AggrB)}], {get, Bucket, Metric, Time, Count, Opts}; decode(<>) -> {stream, Bucket, Delay}; decode(<>) -> {events, Bucket, decode_events(Events)}; decode(<>) -> {get_events, Bucket, Start, End}; decode(<>) -> {get_events, Bucket, Start, End, jsxd_filter:deserialize(Filter)}; decode(<>) -> {events, decode_events(Events)}; decode(<>) -> events_end; decode(<>) -> {error, Message}. decode_events(<<>>) -> []; decode_events(Compressed) -> {ok, Events} = snappiest:decompress(Compressed), [ {T, binary_to_term(E)} || <> <= Events]. %%-------------------------------------------------------------------- %% @doc %% Decodes a streaming TCP message from the wire protocol. %% %% @end %%-------------------------------------------------------------------- -spec decode_stream(binary()) -> {stream_message(), binary()}. decode_stream(<>) -> {flush, Rest}; decode_stream(<>) -> {{stream, Metric, Time, Points}, Rest}; decode_stream(<>) -> {{batch, Time}, Rest}; decode_stream(Rest) when is_binary(Rest) -> {incomplete, Rest}. %%-------------------------------------------------------------------- %% @doc %% Decodes a batched TCP message from the wire protocol. %% %% @end %%-------------------------------------------------------------------- -spec decode_batch(binary()) -> {batch_message(), binary()}. decode_batch(<<0:?METRIC_SS/?SIZE_TYPE, Rest/binary>>) -> {batch_end, Rest}; decode_batch(<<_MetricSize:?METRIC_SS/?SIZE_TYPE, Metric:_MetricSize/binary, Point:?DATA_SIZE/binary, Rest/binary>>) -> {{batch, Metric, Point}, Rest}; decode_batch(Rest) -> {incomplete, Rest}. encode_get_reply({aggr, undefined}) -> <>; encode_get_reply({aggr, Aggr}) -> AggrB = encode_aggr(Aggr), <>. encode_get_stream({data, Data}) -> {ok, Compressed} = snappiest:compress(Data), <>; encode_get_stream({data, Data, 0}) -> {ok, Compressed} = snappiest:compress(Data), <>; encode_get_stream({data, Data, Padding}) -> {ok, Compressed} = snappiest:compress(Data), <>; encode_get_stream(done) -> <>. %%-------------------------------------------------------------------- %% @doc %% Decodes streamed/compressed the initial get reply. %% %% @end %%-------------------------------------------------------------------- -spec decode_get_reply(<<_:8, _:_*8>>) -> {aggr, aggr() | undefined, get_stream_element()}. decode_get_reply(<>) -> {aggr, undefined, {more, <<>>}}; decode_get_reply(<>) -> {aggr, decode_aggr(Aggr), {more, <<>>}}; decode_get_reply(NoAggr) -> Res = decode_get_stream(NoAggr, <<>>), {aggr, undefined, Res}. %%-------------------------------------------------------------------- %% @doc %% Decodes streamed/compressed consecutive get replies. %% %% @end %%-------------------------------------------------------------------- -spec decode_get_stream(binary(), binary()) -> get_stream_element(). decode_get_stream(<>, Acc) -> {done, Acc}; decode_get_stream(<>, Acc) -> {ok, Data} = snappiest:decompress(Compressed), {more, <>}; decode_get_stream(<>, Acc) -> {ok, Data} = snappiest:decompress(Compressed), {more, <>}; %% Backwards compatibility decode_get_stream(<>, Acc) -> {ok, Data} = snappiest:decompress(Compressed), {more, <>}. %%-------------------------------------------------------------------- %% @doc %% Encodes/decodes read repair option for a read request %% %% @end %%-------------------------------------------------------------------- encode_rr(off) -> ?OPT_RR_OFF; encode_rr(on) -> ?OPT_RR_ON; encode_rr(default) -> ?OPT_RR_DEFAULT. decode_rr(?OPT_RR_OFF) -> off; decode_rr(?OPT_RR_ON) -> on; decode_rr(?OPT_RR_DEFAULT) -> default. %%-------------------------------------------------------------------- %% @doc %% Encodes/decodes read quorum(R) option for a read request %% %% @end %%-------------------------------------------------------------------- encode_r(n) -> ?OPT_R_N; encode_r(default) -> ?OPT_R_DEFAULT; encode_r(R) when is_integer(R), R >= 0, (R band 16#FF) =:= R -> R. decode_r(?OPT_R_N) -> n; decode_r(?OPT_R_DEFAULT) -> default; decode_r(R) when is_integer(R), R > 0 -> R. -type aggr_bin() :: <<_:40, _:_*8>>. -spec encode_aggr(aggr()) -> aggr_bin(). encode_aggr({Name, Count}) when is_binary(Name), byte_size(Name) =< 255, Count > 0-> NameS = byte_size(Name), <>. -spec decode_aggr(aggr_bin()) -> aggr(). decode_aggr(<>) when Count > 0-> {Name, Count}. -spec zero_to_undef(non_neg_integer()) -> undefined | pos_integer(). zero_to_undef(0) -> undefined; zero_to_undef(N) when is_integer(N), N > 0 -> N. undef_to_number(undefined) -> 0; undef_to_number(N) when N > 0 -> N. encode_traceids({TraceID, ParentID}) -> <<(undef_to_number(TraceID)):64/unsigned-integer, (undef_to_number(ParentID)):64/unsigned-integer>>; encode_traceids(undefined) -> <<(undef_to_number(undefined)):64/unsigned-integer, (undef_to_number(undefined)):64/unsigned-integer>>. decode_traceids(0, 0) -> undefined; decode_traceids(TraceID, ParentID) -> {zero_to_undef(TraceID), zero_to_undef(ParentID)}.