Packages

A TCP transport channel for the spoke MQTT client.

Current section

Files

Jump to
spoke_tcp src spoke@tcp.erl
Raw

src/spoke@tcp.erl

-module(spoke@tcp).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/spoke/tcp.gleam").
-export([connector/3, connector_with_defaults/1]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
-file("src/spoke/tcp.gleam", 58).
-spec map_tcp_message(mug:tcp_message()) -> spoke@core:transport_event().
map_tcp_message(Msg) ->
case Msg of
{packet, Socket, Data} ->
mug:receive_next_packet_as_message(Socket),
{received_data, Data};
{socket_closed, _} ->
transport_closed;
{tcp_error, _, E} ->
{transport_failed,
<<"TCP error: "/utf8, (gleam@string:inspect(E))/binary>>}
end.
-file("src/spoke/tcp.gleam", 70).
-spec map_mug_error({ok, LRN} | {error, any()}, binary()) -> {ok, LRN} |
{error, binary()}.
map_mug_error(R, Reason) ->
gleam@result:map_error(
R,
fun(E) ->
<<<<Reason/binary, ": "/utf8>>/binary,
(gleam@string:inspect(E))/binary>>
end
).
-file("src/spoke/tcp.gleam", 53).
-spec send(mug:socket(), gleam@bytes_tree:bytes_tree()) -> {ok, nil} |
{error, binary()}.
send(Socket, Data) ->
_pipe = mug_ffi:send(Socket, Data),
map_mug_error(_pipe, <<"Send error"/utf8>>).
-file("src/spoke/tcp.gleam", 24).
-spec connect(binary(), integer(), integer()) -> {ok,
spoke@mqtt_actor:transport_channel()} |
{error, binary()}.
connect(Host, Port, Connect_timeout) ->
Options = {connection_options, Host, Port, Connect_timeout, ipv6_preferred},
gleam@result:'try'(
begin
_pipe = mug:connect(Options),
map_mug_error(_pipe, <<"Connect error"/utf8>>)
end,
fun(Socket) ->
Subject = gleam@erlang@process:new_subject(),
Selector = begin
_pipe@1 = gleam_erlang_ffi:new_selector(),
_pipe@2 = mug:select_tcp_messages(
_pipe@1,
fun map_tcp_message/1
),
gleam@erlang@process:select(_pipe@2, Subject)
end,
gleam@erlang@process:send(Subject, transport_established),
mug:receive_next_packet_as_message(Socket),
{ok,
{transport_channel,
Selector,
fun(_capture) -> send(Socket, _capture) end,
fun() -> _pipe@3 = mug_ffi:shutdown(Socket),
map_mug_error(_pipe@3, <<"Shutdown error"/utf8>>) end}}
end
).
-file("src/spoke/tcp.gleam", 16).
?DOC(" Constructs an (unencrypted) TCP connector.\n").
-spec connector(binary(), integer(), integer()) -> fun(() -> {ok,
spoke@mqtt_actor:transport_channel()} |
{error, binary()}).
connector(Host, Port, Connect_timeout) ->
fun() -> connect(Host, Port, Connect_timeout) end.
-file("src/spoke/tcp.gleam", 11).
?DOC(
" Constructs an (unencrypted) TCP connector using the defaults of\n"
" connecting to port 1883 on the given host.\n"
).
-spec connector_with_defaults(binary()) -> fun(() -> {ok,
spoke@mqtt_actor:transport_channel()} |
{error, binary()}).
connector_with_defaults(Host) ->
connector(Host, 1883, 5000).