Packages
macula
0.32.1
7.0.0
6.0.0
5.2.2
5.2.1
5.2.0
5.1.0
5.0.0
4.8.0
4.7.1
4.7.0
4.6.0
4.5.0
4.4.10
4.4.9
4.4.8
4.4.7
4.4.6
4.4.5
4.4.4
4.4.3
4.4.2
4.4.1
4.4.0
4.3.1
4.3.0
4.2.9
4.2.8
4.2.7
4.2.6
4.2.5
4.2.4
4.2.3
4.2.2
4.2.1
4.2.0
4.1.1
4.1.0
4.0.0
3.16.0
3.15.3
3.15.2
3.15.1
3.14.0
3.13.0
3.12.1
3.12.0
3.11.1
3.11.0
3.10.3
3.10.2
3.10.1
3.9.0
3.8.0
3.7.0
3.5.0
3.4.0
3.3.0
3.2.0
3.1.0
3.0.0
2.1.1
2.1.0
2.0.0
1.5.2
1.5.1
1.4.30
1.4.29
1.4.28
1.4.27
1.4.26
1.4.25
1.4.24
1.4.23
1.4.22
1.4.21
1.4.20
1.4.19
1.4.18
1.4.17
1.4.16
1.4.15
1.4.14
1.4.13
1.4.11
1.4.10
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.1
1.3.0
1.2.0
1.1.0
1.0.10
1.0.9
1.0.8
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.48.6
0.48.5
0.48.4
0.48.3
0.48.2
0.48.1
0.48.0
0.47.1
0.47.0
0.46.3
0.46.1
0.46.0
0.45.3
0.45.2
0.45.1
0.45.0
0.44.2
0.44.1
0.44.0
0.43.3
0.43.2
0.43.1
0.43.0
0.42.9
0.42.8
0.42.7
0.42.6
0.42.5
0.42.4
0.42.3
0.42.2
0.42.1
0.42.0
0.41.1
0.41.0
0.40.1
0.40.0
0.39.9
0.39.8
0.39.7
0.39.6
0.39.5
0.39.4
0.39.3
0.39.2
0.39.1
0.39.0
0.38.8
0.38.7
0.38.6
0.38.5
0.38.4
0.38.3
0.38.2
0.38.1
0.38.0
0.37.7
0.37.6
0.37.5
0.37.4
0.37.3
0.37.2
0.37.1
0.37.0
0.36.6
0.36.5
0.36.4
0.36.3
0.36.2
0.36.1
0.36.0
0.35.4
0.35.3
0.35.2
0.35.1
0.35.0
0.34.1
0.34.0
0.33.1
0.33.0
0.32.5
0.32.4
0.32.3
0.32.2
0.32.1
0.32.0
0.31.9
0.31.8
0.31.7
0.31.6
0.31.5
0.31.4
0.31.3
0.31.2
0.31.1
0.31.0
0.30.10
0.30.9
0.30.8
0.30.7
0.30.6
0.30.5
0.30.4
0.30.3
0.30.2
0.30.1
0.30.0
0.29.0
0.28.3
0.28.2
0.28.1
0.28.0
0.27.1
0.27.0
0.26.1
0.26.0
0.25.6
0.25.5
0.25.4
0.25.3
0.25.2
0.25.1
0.25.0
0.24.6
0.24.5
0.24.4
0.24.3
0.24.2
0.24.1
0.24.0
0.23.3
0.23.2
0.23.1
0.23.0
0.22.12
0.22.11
0.22.10
0.22.9
0.22.8
0.22.7
0.22.6
0.22.5
0.22.4
0.22.3
0.22.2
0.22.1
0.22.0
0.21.7
0.21.6
0.21.5
0.21.4
0.21.2
0.21.1
0.21.0
0.20.25
0.20.24
0.20.23
0.20.22
0.20.21
0.20.20
0.20.19
0.20.18
0.20.17
0.20.16
0.20.15
0.20.14
0.20.13
0.20.12
0.20.11
0.20.10
0.20.9
0.20.8
0.20.7
0.20.6
0.20.5
0.20.3
0.20.2
0.20.1
0.20.0
0.19.2
0.19.1
0.19.0
0.18.1
0.18.0
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.6
0.16.5
0.16.4
0.16.3
0.16.2
0.16.1
0.16.0
0.15.1
0.15.0
0.14.3
0.14.2
0.14.1
0.14.0
0.12.6
0.12.5
0.12.3
0.11.3
0.10.2
0.10.1
0.10.0
0.9.2
0.9.1
0.9.0
0.8.25
0.8.24
0.8.23
0.8.22
0.8.21
0.8.20
0.8.19
0.8.18
0.8.17
0.8.16
0.8.15
0.8.14
0.8.13
0.8.12
0.8.11
0.8.10
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.30
0.7.29
0.7.28
0.7.27
0.7.26
0.7.25
0.7.24
0.7.23
0.7.22
0.7.21
0.7.20
0.7.19
0.7.18
0.7.17
0.7.16
0.7.15
0.7.14
0.7.13
0.7.12
0.7.11
0.7.10
0.7.9
0.7.8
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.7
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.4
0.3.3
0.3.2
0.3.1
Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise
Current section
Files
Jump to
Current section
Files
src/macula_relay_client.erl
%%%-------------------------------------------------------------------
%%% @doc Macula Relay Client — connects to relay servers via QUIC.
%%%
%%% Single gen_server managing ONE persistent QUIC connection with
%%% automatic failover across multiple relays. On disconnect, cycles
%%% to the next relay in the list with exponential backoff + jitter.
%%% Replays all subscriptions and procedure registrations on reconnect.
%%%
%%% Usage:
%%%
%%% ```
%%% Opts = #{relays => [Relay0, Relay1, Relay2]},
%%% {ok, Client} = macula_relay_client:start_link(Opts).
%%% '''
%%%
%%% Backward compatible with single relay via url key.
%%% @end
%%%-------------------------------------------------------------------
-module(macula_relay_client).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
-export([start_link/1, stop/1]).
-export([subscribe/3, unsubscribe/2, publish/3]).
-export([advertise/3, unadvertise/2, call/4]).
-export([init/1, handle_call/3, handle_cast/2, handle_info/2, terminate/2]).
%% Test exports
-ifdef(TEST).
-export([backoff_ms/1, parse_url/1]).
-endif.
-define(RECONNECT_BASE_MS, 1000).
-define(RECONNECT_MAX_MS, 30000).
-define(RECONNECT_JITTER, 0.3).
-define(CALL_TIMEOUT, 5000).
-record(state, {
relays :: [binary()], %% all configured relay URLs
relay_index :: non_neg_integer(), %% current relay (0-based)
url :: binary(), %% current relay URL
host :: string(),
port :: integer(),
realm :: binary(),
identity :: binary(),
client_type :: binary(), %% <<"node">> or <<"relay">> (for CONNECT handshake)
conn :: reference() | undefined,
stream :: reference() | undefined,
status :: connecting | connected | disconnected,
recv_buffer :: binary(),
reconnect_attempt :: non_neg_integer(), %% for exponential backoff
%% Application state (survives reconnects)
subscriptions :: #{reference() => {binary(), fun()}}, %% ref => {topic, callback}
procedures :: #{binary() => fun()}, %% procedure => handler
pending_calls :: #{binary() => {pid(), reference()}} %% call_id => {from, timer_ref}
}).
%%====================================================================
%% API
%%====================================================================
start_link(Opts) ->
gen_server:start_link(?MODULE, Opts, []).
stop(Pid) ->
gen_server:stop(Pid).
-spec subscribe(pid(), binary(), fun((map()) -> ok)) -> {ok, reference()}.
subscribe(Pid, Topic, Callback) ->
gen_server:call(Pid, {subscribe, Topic, Callback}).
-spec unsubscribe(pid(), reference()) -> ok.
unsubscribe(Pid, Ref) ->
gen_server:call(Pid, {unsubscribe, Ref}).
-spec publish(pid(), binary(), binary() | map()) -> ok.
publish(Pid, Topic, Payload) ->
gen_server:cast(Pid, {publish, Topic, Payload}).
-spec advertise(pid(), binary(), fun((map()) -> {ok, term()} | {error, term()})) -> {ok, reference()}.
advertise(Pid, Procedure, Handler) ->
gen_server:call(Pid, {advertise, Procedure, Handler}).
-spec unadvertise(pid(), binary()) -> ok.
unadvertise(Pid, Procedure) ->
gen_server:call(Pid, {unadvertise, Procedure}).
-spec call(pid(), binary(), map(), timeout()) -> {ok, term()} | {error, term()}.
call(Pid, Procedure, Args, Timeout) ->
gen_server:call(Pid, {rpc_call, Procedure, Args, Timeout}, Timeout + 1000).
%%====================================================================
%% gen_server callbacks
%%====================================================================
init(Opts) ->
Relays = case maps:find(relays, Opts) of
{ok, List} when is_list(List), length(List) > 0 -> List;
_ -> [maps:get(url, Opts, <<"https://localhost:4433">>)]
end,
%% Randomize initial relay to distribute load across relays
Index = rand:uniform(length(Relays)) - 1,
Url = lists:nth(Index + 1, Relays),
{Host, Port} = parse_url(Url),
Realm = maps:get(realm, Opts, <<"io.macula">>),
Identity = maps:get(identity, Opts, <<"anonymous">>),
ClientType = maps:get(type, Opts, <<"node">>),
State = #state{
relays = Relays,
relay_index = Index,
url = Url,
host = Host,
port = Port,
realm = Realm,
identity = Identity,
client_type = ClientType,
conn = undefined,
stream = undefined,
status = connecting,
recv_buffer = <<>>,
reconnect_attempt = 0,
subscriptions = #{},
procedures = #{},
pending_calls = #{}
},
self() ! connect,
{ok, State}.
%%====================================================================
%% Subscribe / Unsubscribe
%%====================================================================
handle_call({subscribe, Topic, Callback}, _From, State) ->
Ref = make_ref(),
Subs = maps:put(Ref, {Topic, Callback}, State#state.subscriptions),
%% Send SUBSCRIBE if connected
maybe_send(subscribe, #{<<"topics">> => [Topic], <<"qos">> => 0}, State),
{reply, {ok, Ref}, State#state{subscriptions = Subs}};
handle_call({unsubscribe, Ref}, _From, State) ->
case maps:get(Ref, State#state.subscriptions, undefined) of
undefined ->
{reply, {error, not_found}, State};
{Topic, _Callback} ->
Subs = maps:remove(Ref, State#state.subscriptions),
maybe_send(unsubscribe, #{<<"topics">> => [Topic]}, State),
{reply, ok, State#state{subscriptions = Subs}}
end;
%%====================================================================
%% Advertise / Unadvertise
%%====================================================================
handle_call({advertise, Procedure, Handler}, _From, State) ->
Procs = maps:put(Procedure, Handler, State#state.procedures),
maybe_send(register_procedure, #{<<"procedure">> => Procedure}, State),
{reply, {ok, make_ref()}, State#state{procedures = Procs}};
handle_call({unadvertise, Procedure}, _From, State) ->
Procs = maps:remove(Procedure, State#state.procedures),
%% No protocol message for unadvertise yet — procedure removed on disconnect
{reply, ok, State#state{procedures = Procs}};
%%====================================================================
%% RPC Call
%%====================================================================
handle_call({rpc_call, Procedure, Args, Timeout}, From, State) ->
CallId = base64:encode(crypto:strong_rand_bytes(12)),
TimerRef = erlang:send_after(Timeout, self(), {call_timeout, CallId}),
PendingCalls = maps:put(CallId, {From, TimerRef}, State#state.pending_calls),
maybe_send(call, #{
<<"call_id">> => CallId,
<<"procedure">> => Procedure,
<<"args">> => Args
}, State),
{noreply, State#state{pending_calls = PendingCalls}};
handle_call(_Request, _From, State) ->
{reply, {error, unknown}, State}.
%%====================================================================
%% Publish (async)
%%====================================================================
handle_cast({publish, Topic, Payload}, State) ->
BinPayload = to_binary_payload(Payload),
maybe_send(publish, #{
<<"topic">> => Topic,
<<"payload">> => BinPayload,
<<"qos">> => 0,
<<"retain">> => false,
<<"message_id">> => crypto:strong_rand_bytes(16)
}, State),
{noreply, State};
handle_cast(_Msg, State) ->
{noreply, State}.
%%====================================================================
%% Connection lifecycle
%%====================================================================
handle_info(connect, State) ->
QuicOpts = [{alpn, ["macula"]} | macula_tls:quic_client_opts()],
case macula_quic:connect(State#state.host, State#state.port, QuicOpts, 10000) of
{ok, Conn} ->
case macula_quic:open_stream(Conn) of
{ok, Stream} ->
quicer:setopt(Stream, active, true),
?LOG_INFO("[relay_client] Connected to ~s", [State#state.url]),
State2 = State#state{conn = Conn, stream = Stream, status = connected,
recv_buffer = <<>>, reconnect_attempt = 0},
send_connect(State2),
replay_state(State2),
{noreply, State2};
{error, StreamErr} ->
?LOG_WARNING("[relay_client] Stream open failed on ~s: ~p", [State#state.url, StreamErr]),
catch macula_quic:close(Conn),
{noreply, schedule_failover(State)}
end;
{error, Reason} ->
?LOG_WARNING("[relay_client] Connect to ~s failed: ~p", [State#state.url, Reason]),
{noreply, schedule_failover(State)};
{error, Type, Detail} ->
?LOG_WARNING("[relay_client] Connect to ~s failed: ~p ~p", [State#state.url, Type, Detail]),
{noreply, schedule_failover(State)}
end;
handle_info({call_timeout, CallId}, State) ->
case maps:get(CallId, State#state.pending_calls, undefined) of
undefined -> {noreply, State};
{From, _TimerRef} ->
gen_server:reply(From, {error, timeout}),
{noreply, State#state{pending_calls = maps:remove(CallId, State#state.pending_calls)}}
end;
%%====================================================================
%% QUIC data received
%%====================================================================
handle_info({quic, Data, _Stream, _Flags}, State) when is_binary(Data) ->
Buffer = <<(State#state.recv_buffer)/binary, Data/binary>>,
{NewBuffer, State2} = process_buffer(Buffer, State),
{noreply, State2#state{recv_buffer = NewBuffer}};
%%====================================================================
%% QUIC lifecycle — reconnect on any closure
%%====================================================================
handle_info({quic, peer_send_shutdown, _Stream, _}, State) ->
handle_disconnect(State);
handle_info({quic, peer_send_aborted, _Stream, _}, State) ->
handle_disconnect(State);
handle_info({quic, send_shutdown_complete, _Stream, _}, State) ->
handle_disconnect(State);
handle_info({quic, shutdown, _Ref, _Reason}, State) ->
handle_disconnect(State);
handle_info({quic, closed, _Ref, _Reason}, State) ->
handle_disconnect(State);
handle_info({quic, transport_shutdown, _Ref, _Reason}, State) ->
handle_disconnect(State);
handle_info({error, transport_down, _Detail}, State) ->
handle_disconnect(State);
handle_info({quic, streams_available, _Conn, _Info}, State) ->
{noreply, State};
handle_info({quic, peer_needs_streams, _Conn, _Info}, State) ->
{noreply, State};
handle_info(Info, State) ->
?LOG_DEBUG("[relay_client] Unhandled: ~p", [Info]),
{noreply, State}.
terminate(_Reason, #state{conn = Conn, stream = Stream}) ->
catch macula_quic:close(Stream),
catch macula_quic:close(Conn),
ok.
%%====================================================================
%% Internal: message processing
%%====================================================================
process_buffer(Buffer, State) when byte_size(Buffer) < 8 ->
{Buffer, State};
process_buffer(<<_Version:8, _TypeId:8, _Flags:8, _Reserved:8,
PayloadLen:32/big-unsigned, Rest/binary>> = Buffer, State) ->
case byte_size(Rest) >= PayloadLen of
true ->
MsgBytes = binary:part(Buffer, 0, 8 + PayloadLen),
Remaining = binary:part(Buffer, 8 + PayloadLen, byte_size(Buffer) - 8 - PayloadLen),
State2 = handle_message(macula_protocol_decoder:decode(MsgBytes), State),
process_buffer(Remaining, State2);
false ->
{Buffer, State}
end.
%% Incoming PUBLISH from relay
handle_message({ok, {publish, Msg}}, State) ->
Topic = maps:get(<<"topic">>, Msg, <<>>),
Payload = maps:get(<<"payload">>, Msg, <<>>),
invoke_matching_callbacks(Topic, Payload, State#state.subscriptions),
State;
%% Incoming CALL from relay (we're the provider)
handle_message({ok, {call, Msg}}, State) ->
Procedure = maps:get(<<"procedure">>, Msg, <<>>),
CallId = maps:get(<<"call_id">>, Msg, <<>>),
Args = maps:get(<<"args">>, Msg, #{}),
case maps:get(Procedure, State#state.procedures, undefined) of
undefined ->
maybe_send(reply, #{<<"call_id">> => CallId,
<<"error">> => #{<<"code">> => <<"not_found">>}}, State);
Handler ->
%% Execute handler in separate process to avoid blocking
Stream = State#state.stream,
spawn(fun() ->
Result = try Handler(Args)
catch _:Err -> {error, Err}
end,
ReplyMsg = case Result of
{ok, R} -> #{<<"call_id">> => CallId, <<"result">> => R};
{error, R} -> #{<<"call_id">> => CallId,
<<"error">> => #{<<"code">> => <<"handler_error">>,
<<"message">> => iolist_to_binary(
io_lib:format("~p", [R]))}}
end,
Binary = macula_protocol_encoder:encode(reply, ReplyMsg),
macula_quic:async_send(Stream, Binary)
end)
end,
State;
%% Incoming REPLY from relay (we're the caller)
handle_message({ok, {reply, Msg}}, State) ->
CallId = maps:get(<<"call_id">>, Msg, <<>>),
case maps:get(CallId, State#state.pending_calls, undefined) of
undefined ->
?LOG_DEBUG("[relay_client] Reply for unknown call_id"),
State;
{From, TimerRef} ->
erlang:cancel_timer(TimerRef),
Result = case maps:get(<<"result">>, Msg, undefined) of
undefined ->
Error = maps:get(<<"error">>, Msg, #{}),
{error, Error};
R ->
{ok, R}
end,
gen_server:reply(From, Result),
State#state{pending_calls = maps:remove(CallId, State#state.pending_calls)}
end;
%% PONG — connection alive
handle_message({ok, {pong, _Msg}}, State) ->
State;
handle_message({ok, {Type, _Msg}}, State) ->
?LOG_DEBUG("[relay_client] Ignoring message type: ~p", [Type]),
State;
handle_message({error, Reason}, State) ->
?LOG_WARNING("[relay_client] Decode error: ~p", [Reason]),
State.
%%====================================================================
%% Internal helpers
%%====================================================================
send_connect(State) ->
NodeId = crypto:strong_rand_bytes(32),
NodeName = case net_kernel:nodename() of
nonode@nohost -> State#state.identity;
Name -> atom_to_binary(Name)
end,
Base = #{
version => <<"1.0">>,
node_id => NodeId,
node_name => NodeName,
realm_id => State#state.realm,
capabilities => [<<"pubsub">>, <<"rpc">>],
endpoint => State#state.url,
type => State#state.client_type
},
Msg = case macula_connection:build_node_identity(#{}) of
Identity when map_size(Identity) > 0 -> Base#{identity => Identity};
_ -> Base
end,
maybe_send(connect, Msg, State).
%% Replay all subscriptions and procedure registrations after reconnect
replay_state(#state{subscriptions = Subs, procedures = Procs} = State) ->
%% Batch all topics in one SUBSCRIBE
Topics = lists:usort([T || {_Ref, {T, _Cb}} <- maps:to_list(Subs)]),
case Topics of
[] -> ok;
_ ->
?LOG_INFO("[relay_client] Replaying ~p subscription(s)", [length(Topics)]),
maybe_send(subscribe, #{<<"topics">> => Topics, <<"qos">> => 0}, State)
end,
%% Register all procedures
maps:foreach(fun(Proc, _Handler) ->
?LOG_INFO("[relay_client] Replaying procedure: ~s", [Proc]),
maybe_send(register_procedure, #{<<"procedure">> => Proc}, State)
end, Procs).
maybe_send(_Type, _Msg, #state{stream = undefined}) -> ok;
maybe_send(Type, Msg, #state{stream = Stream}) ->
Binary = macula_protocol_encoder:encode(Type, Msg),
macula_quic:async_send(Stream, Binary).
%% Ensure payload is a flat binary before passing to msgpack.
%% Callers may pass: map, iolist (from json:encode), or binary.
to_binary_payload(Payload) when is_binary(Payload) -> Payload;
to_binary_payload(Payload) when is_map(Payload) -> iolist_to_binary(json:encode(Payload));
to_binary_payload(Payload) when is_list(Payload) -> iolist_to_binary(Payload);
to_binary_payload(Payload) -> iolist_to_binary(json:encode(Payload)).
handle_disconnect(#state{status = disconnected} = State) ->
{noreply, State};
handle_disconnect(State) ->
?LOG_WARNING("[relay_client] Disconnected from ~s", [State#state.url]),
catch macula_quic:close(State#state.stream),
catch macula_quic:close(State#state.conn),
State2 = State#state{conn = undefined, stream = undefined, status = disconnected,
recv_buffer = <<>>},
{noreply, schedule_failover(State2)}.
%% Pick next relay (round-robin) and schedule reconnect with exponential backoff + jitter.
schedule_failover(#state{relays = Relays, relay_index = CurrentIdx,
reconnect_attempt = Attempt} = State) ->
NumRelays = length(Relays),
%% Try a different relay if available
NextIdx = case NumRelays > 1 of
true -> (CurrentIdx + 1) rem NumRelays;
false -> CurrentIdx
end,
NextUrl = lists:nth(NextIdx + 1, Relays),
{NextHost, NextPort} = parse_url(NextUrl),
%% Exponential backoff with jitter (±30%)
BackoffMs = backoff_ms(Attempt),
erlang:send_after(BackoffMs, self(), connect),
?LOG_INFO("[relay_client] Failover to ~s in ~bms (attempt ~b)",
[NextUrl, BackoffMs, Attempt + 1]),
State#state{relay_index = NextIdx, url = NextUrl, host = NextHost,
port = NextPort, status = disconnected,
reconnect_attempt = Attempt + 1}.
backoff_ms(Attempt) ->
Base = min(?RECONNECT_BASE_MS * (1 bsl min(Attempt, 10)), ?RECONNECT_MAX_MS),
Jitter = round(Base * ?RECONNECT_JITTER),
Base + rand:uniform(2 * Jitter + 1) - Jitter - 1.
parse_url(Url) when is_binary(Url) ->
parse_url(binary_to_list(Url));
parse_url("https://" ++ Rest) ->
parse_host_port(Rest);
parse_url(Rest) ->
parse_host_port(Rest).
parse_host_port(HostPort) ->
case string:split(HostPort, ":", trailing) of
[Host, PortStr] -> {Host, list_to_integer(PortStr)};
[Host] -> {Host, 4433}
end.
invoke_matching_callbacks(Topic, Payload, Subscriptions) ->
DecodedPayload = safe_decode_json(Payload),
maps:foreach(fun(_Ref, {SubTopic, Callback}) ->
case topic_matches(SubTopic, Topic) of
true ->
spawn(fun() ->
try Callback(#{topic => Topic, payload => DecodedPayload})
catch C:E -> ?LOG_WARNING("[relay_client] Callback error: ~p:~p", [C, E])
end
end);
false -> ok
end
end, Subscriptions).
topic_matches(Pattern, Topic) when Pattern =:= Topic -> true;
topic_matches(<<"**">>, _Topic) -> true;
topic_matches(Pattern, Topic) ->
PatternParts = binary:split(Pattern, <<".">>, [global]),
TopicParts = binary:split(Topic, <<".">>, [global]),
parts_match(PatternParts, TopicParts).
parts_match([], []) -> true;
parts_match([<<"**">>], _) -> true;
parts_match([<<"*">> | PR], [_ | TR]) -> parts_match(PR, TR);
parts_match([P | PR], [T | TR]) when P =:= T -> parts_match(PR, TR);
parts_match(_, _) -> false.
safe_decode_json(Payload) when is_binary(Payload) ->
try json:decode(Payload) catch _:_ -> Payload end;
safe_decode_json(Payload) -> Payload.