Current section
Files
Jump to
Current section
Files
src/gemqtt@subscriber.erl
-module(gemqtt@subscriber).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch]).
-export([puback_from_dynamic/1, selecting_mqtt_pubacks/2, new/1, set_local_echo/2, set_qos/2, set_retain_as_published/2, set_retain_handling/2, set_property/3, add/2, remove/2, message_from_dynamic/1, selecting_mqtt_messages/2]).
-export_type([message/0, pub_ack/0, retain_handling/0, subscriber/0, subscribe_option/0]).
-type message() :: {message,
gemqtt:client(),
boolean(),
gleam@option:option(integer()),
bitstring(),
gemqtt:qos(),
boolean(),
binary()}.
-type pub_ack() :: {pub_ack, integer(), integer()}.
-type retain_handling() :: sent_always | sent_on_new_subscription | sent_never.
-opaque subscriber() :: {subscriber,
gemqtt:client(),
boolean(),
gemqtt:qos(),
boolean(),
retain_handling(),
gemqtt:properties()}.
-type subscribe_option() :: {nl, boolean()} |
{qos, gemqtt:qos()} |
{rap, boolean()} |
{rh, integer()}.
-spec puback_from_dynamic(gleam@dynamic:dynamic_()) -> {ok, pub_ack()} |
{error, list(gleam@dynamic:decode_error())}.
puback_from_dynamic(Input) ->
gleam@result:'try'(
(gleam@dynamic:field(packet_id, fun gleam@dynamic:int/1))(Input),
fun(Packet_id) ->
gleam@result:'try'(
(gleam@dynamic:field(reason_code, fun gleam@dynamic:int/1))(
Input
),
fun(Reason_code) -> {ok, {pub_ack, Packet_id, Reason_code}} end
)
end
).
-spec map_dynamic_puback(fun((pub_ack()) -> GAA)) -> fun((gleam@dynamic:dynamic_()) -> GAA).
map_dynamic_puback(Mapper) ->
fun(Input) ->
_assert_subject = puback_from_dynamic(Input),
{ok, Puback} = case _assert_subject of
{ok, _} -> _assert_subject;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Assertion pattern match failed"/utf8>>,
value => _assert_fail,
module => <<"gemqtt/subscriber"/utf8>>,
function => <<"map_dynamic_puback"/utf8>>,
line => 109})
end,
Mapper(Puback)
end.
-spec selecting_mqtt_pubacks(
gleam@erlang@process:selector(FZU),
fun((pub_ack()) -> FZU)
) -> gleam@erlang@process:selector(FZU).
selecting_mqtt_pubacks(Selector, Mapper) ->
Puback = erlang:binary_to_atom(<<"puback"/utf8>>),
_pipe = Selector,
gleam@erlang@process:selecting_record2(
_pipe,
Puback,
map_dynamic_puback(Mapper)
).
-spec new(gemqtt:client()) -> subscriber().
new(Client) ->
{subscriber,
Client,
false,
at_most_once,
false,
sent_always,
{properties, gleam@dict:new()}}.
-spec set_local_echo(subscriber(), boolean()) -> subscriber().
set_local_echo(Subscriber, Value) ->
erlang:setelement(3, Subscriber, not Value).
-spec set_qos(subscriber(), gemqtt:qos()) -> subscriber().
set_qos(Subscriber, Value) ->
erlang:setelement(4, Subscriber, Value).
-spec set_retain_as_published(subscriber(), boolean()) -> subscriber().
set_retain_as_published(Subscriber, Value) ->
erlang:setelement(5, Subscriber, Value).
-spec set_retain_handling(subscriber(), retain_handling()) -> subscriber().
set_retain_handling(Subscriber, Value) ->
erlang:setelement(6, Subscriber, Value).
-spec set_property(subscriber(), binary(), any()) -> subscriber().
set_property(Sub, Name, Value) ->
{properties, Props} = erlang:element(7, Sub),
erlang:setelement(
7,
Sub,
{properties,
gleam@dict:insert(
Props,
erlang:binary_to_atom(Name),
gleam@dynamic:from(Value)
)}
).
-spec add(subscriber(), list(binary())) -> {ok,
{gleam@option:option(gemqtt:properties()), list(integer())}} |
{error, nil}.
add(Sub, Topics) ->
Rh = case erlang:element(6, Sub) of
sent_always ->
0;
sent_on_new_subscription ->
1;
sent_never ->
2
end,
Opts = [{nl, erlang:element(3, Sub)},
{qos, erlang:element(4, Sub)},
{rap, erlang:element(5, Sub)},
{rh, Rh}],
{properties, Props} = erlang:element(7, Sub),
emqtt_ffi:subscribe(erlang:element(2, Sub), Opts, Props, Topics).
-spec remove(gemqtt:client(), list(binary())) -> {ok,
{gleam@option:option(gemqtt:properties()), list(integer())}} |
{error, nil}.
remove(Client, Topics) ->
emqtt_ffi:unsubscribe(Client, Topics).
-spec decode_qos(gleam@dynamic:dynamic_()) -> {ok, gemqtt:qos()} |
{error, list(gleam@dynamic:decode_error())}.
decode_qos(Data) ->
_pipe = Data,
_pipe@1 = gleam@dynamic:int(_pipe),
gleam@result:'try'(_pipe@1, fun(Qos) -> case Qos of
0 ->
{ok, at_most_once};
1 ->
{ok, at_least_once};
2 ->
{ok, exactly_once};
_ ->
{error,
[{decode_error,
<<"0,1,2"/utf8>>,
gleam@int:to_string(Qos),
[]}]}
end end).
-spec message_from_dynamic(gleam@dynamic:dynamic_()) -> {ok, message()} |
{error, list(gleam@dynamic:decode_error())}.
message_from_dynamic(Input) ->
gleam@result:'try'(
(gleam@dynamic:field(client_pid, fun emqtt_ffi:decode_client/1))(Input),
fun(Client) ->
gleam@result:'try'(
(gleam@dynamic:field(dup, fun gleam@dynamic:bool/1))(Input),
fun(Duplicate) ->
gleam@result:'try'(
(gleam@dynamic:optional_field(
packet_id,
fun gleam@dynamic:int/1
))(Input),
fun(Packet_id) ->
gleam@result:'try'(
(gleam@dynamic:field(
payload,
fun gleam@dynamic:bit_array/1
))(Input),
fun(Payload) ->
gleam@result:'try'(
(gleam@dynamic:field(
qos,
fun decode_qos/1
))(Input),
fun(Qos) ->
gleam@result:'try'(
(gleam@dynamic:field(
retain,
fun gleam@dynamic:bool/1
))(Input),
fun(Retain) ->
gleam@result:'try'(
(gleam@dynamic:field(
topic,
fun gleam@dynamic:string/1
))(Input),
fun(Topic) ->
{ok,
{message,
Client,
Duplicate,
Packet_id,
Payload,
Qos,
Retain,
Topic}}
end
)
end
)
end
)
end
)
end
)
end
)
end
).
-spec map_dynamic_message(fun((message()) -> FZT)) -> fun((gleam@dynamic:dynamic_()) -> FZT).
map_dynamic_message(Mapper) ->
fun(Input) ->
_assert_subject = message_from_dynamic(Input),
{ok, Message} = case _assert_subject of
{ok, _} -> _assert_subject;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Assertion pattern match failed"/utf8>>,
value => _assert_fail,
module => <<"gemqtt/subscriber"/utf8>>,
function => <<"map_dynamic_message"/utf8>>,
line => 70})
end,
Mapper(Message)
end.
-spec selecting_mqtt_messages(
gleam@erlang@process:selector(FZN),
fun((message()) -> FZN)
) -> gleam@erlang@process:selector(FZN).
selecting_mqtt_messages(Selector, Mapper) ->
Publish = erlang:binary_to_atom(<<"publish"/utf8>>),
_pipe = Selector,
gleam@erlang@process:selecting_record2(
_pipe,
Publish,
map_dynamic_message(Mapper)
).