Current section
Files
Jump to
Current section
Files
src/spoke@core@internal@session.erl
-module(spoke@core@internal@session).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/spoke/core/internal/session.gleam").
-export([new/1, from_state/1, connect/2, pending_publishes/1, packets_to_send_after_connect/1, reserve_packet_id/1, start_qos1_publish/2, start_qos2_publish/2, handle_pubrec/2, handle_pubcomp/2, start_qos2_receive/2, handle_puback/2, handle_pubrel/2]).
-export_type([session/0, pub_ack_result/0]).
-if(?OTP_RELEASE >= 27).
-define(MODULEDOC(Str), -moduledoc(Str)).
-define(DOC(Str), -doc(Str)).
-else.
-define(MODULEDOC(Str), -compile([])).
-define(DOC(Str), -compile([])).
-endif.
?MODULEDOC(false).
-opaque session() :: {session,
boolean(),
integer(),
gleam@dict:dict(integer(), spoke@packet:message_data()),
gleam@dict:dict(integer(), spoke@packet:message_data()),
gleam@set:set(integer()),
gleam@set:set(integer())}.
-type pub_ack_result() :: {publish_finished,
session(),
list(spoke@core@session_state:storage_update())} |
invalid_pub_ack_id.
-file("src/spoke/core/internal/session.gleam", 45).
?DOC(false).
-spec new(boolean()) -> session().
new(Clean_session) ->
{session,
Clean_session,
1,
maps:new(),
maps:new(),
gleam@set:new(),
gleam@set:new()}.
-file("src/spoke/core/internal/session.gleam", 56).
?DOC(false).
-spec from_state(spoke@core@session_state:session_state()) -> session().
from_state(State) ->
Session = new(false),
Session@1 = {session,
erlang:element(2, Session),
erlang:element(2, State),
erlang:element(4, Session),
erlang:element(5, Session),
erlang:element(6, Session),
erlang:element(7, Session)},
Session@3 = begin
gleam@list:fold(
erlang:element(3, State),
Session@1,
fun(Session@2, Packet_state) -> case Packet_state of
{Id, {unacked_qo_s1, Message}} ->
{session,
erlang:element(2, Session@2),
erlang:element(3, Session@2),
gleam@dict:insert(
erlang:element(4, Session@2),
Id,
Message
),
erlang:element(5, Session@2),
erlang:element(6, Session@2),
erlang:element(7, Session@2)};
{Id@1, {unreceived_qo_s2, Message@1}} ->
{session,
erlang:element(2, Session@2),
erlang:element(3, Session@2),
erlang:element(4, Session@2),
gleam@dict:insert(
erlang:element(5, Session@2),
Id@1,
Message@1
),
erlang:element(6, Session@2),
erlang:element(7, Session@2)};
{Id@2, received_qo_s2} ->
{session,
erlang:element(2, Session@2),
erlang:element(3, Session@2),
erlang:element(4, Session@2),
erlang:element(5, Session@2),
gleam@set:insert(erlang:element(6, Session@2), Id@2),
erlang:element(7, Session@2)};
{Id@3, unreleased_qo_s2} ->
{session,
erlang:element(2, Session@2),
erlang:element(3, Session@2),
erlang:element(4, Session@2),
erlang:element(5, Session@2),
erlang:element(6, Session@2),
gleam@set:insert(erlang:element(7, Session@2), Id@3)}
end end
)
end,
Session@3.
-file("src/spoke/core/internal/session.gleam", 94).
?DOC(false).
-spec connect(session(), boolean()) -> {session(),
list(spoke@core@session_state:storage_update())}.
connect(Session, Clean_session) ->
gleam@bool:guard(
erlang:element(2, Session),
{new(Clean_session), []},
fun() -> case Clean_session of
false ->
{Session, []};
true ->
{new(true), [clear_session]}
end end
).
-file("src/spoke/core/internal/session.gleam", 113).
?DOC(false).
-spec pending_publishes(session()) -> integer().
pending_publishes(Session) ->
(maps:size(erlang:element(4, Session)) + maps:size(
erlang:element(5, Session)
))
+ gleam@set:size(erlang:element(6, Session)).
-file("src/spoke/core/internal/session.gleam", 237).
?DOC(false).
-spec packets_to_send_after_connect(session()) -> list(spoke@packet@client@outgoing:packet()).
packets_to_send_after_connect(Session) ->
Packets = gleam@list:new(),
Packets@2 = gleam@list:fold(
maps:keys(erlang:element(4, Session)),
Packets,
fun(Packets@1, Id) ->
Message@1 = case gleam_stdlib:map_get(
erlang:element(4, Session),
Id
) of
{ok, Message} -> Message;
_assert_fail ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"spoke/core/internal/session"/utf8>>,
function => <<"packets_to_send_after_connect"/utf8>>,
line => 243,
value => _assert_fail,
start => 7441,
'end' => 7500,
pattern_start => 7452,
pattern_end => 7463})
end,
Packet = {publish, {publish_data_qo_s1, Message@1, true, Id}},
[Packet | Packets@1]
end
),
Packets@4 = gleam@list:fold(
maps:keys(erlang:element(5, Session)),
Packets@2,
fun(Packets@3, Id@1) ->
Message@3 = case gleam_stdlib:map_get(
erlang:element(5, Session),
Id@1
) of
{ok, Message@2} -> Message@2;
_assert_fail@1 ->
erlang:error(#{gleam_error => let_assert,
message => <<"Pattern match failed, no pattern matched the value."/utf8>>,
file => <<?FILEPATH/utf8>>,
module => <<"spoke/core/internal/session"/utf8>>,
function => <<"packets_to_send_after_connect"/utf8>>,
line => 256,
value => _assert_fail@1,
start => 7812,
'end' => 7874,
pattern_start => 7823,
pattern_end => 7834})
end,
Packet@1 = {publish, {publish_data_qo_s2, Message@3, true, Id@1}},
[Packet@1 | Packets@3]
end
),
Packets@6 = gleam@list:fold(
gleam@set:to_list(erlang:element(6, Session)),
Packets@4,
fun(Packets@5, Id@2) -> [{pub_rel, Id@2} | Packets@5] end
),
Packets@6.
-file("src/spoke/core/internal/session.gleam", 275).
?DOC(false).
-spec store_if_persistent(
session(),
list(spoke@core@session_state:storage_update()),
fun((list(spoke@core@session_state:storage_update())) -> PUA)
) -> PUA.
store_if_persistent(Session, Updates, Operation) ->
case erlang:element(2, Session) of
false ->
Operation(Updates);
true ->
Operation([])
end.
-file("src/spoke/core/internal/session.gleam", 119).
?DOC(false).
-spec reserve_packet_id(session()) -> {session(),
integer(),
list(spoke@core@session_state:storage_update())}.
reserve_packet_id(Session) ->
Id = erlang:element(3, Session),
Next_id = case Id of
65535 ->
1;
Id@1 ->
Id@1 + 1
end,
store_if_persistent(
Session,
[{store_next_packet_id, Next_id}],
fun(Storage_updates) ->
{{session,
erlang:element(2, Session),
Next_id,
erlang:element(4, Session),
erlang:element(5, Session),
erlang:element(6, Session),
erlang:element(7, Session)},
Id,
Storage_updates}
end
).
-file("src/spoke/core/internal/session.gleam", 136).
?DOC(false).
-spec start_qos1_publish(session(), spoke@packet:message_data()) -> {session(),
spoke@packet@client@outgoing:packet(),
list(spoke@core@session_state:storage_update())}.
start_qos1_publish(Session, Message) ->
{Session@1, Id, Id_updates} = reserve_packet_id(Session),
Unacked_qos1 = gleam@dict:insert(erlang:element(4, Session@1), Id, Message),
Packet = {publish, {publish_data_qo_s1, Message, false, Id}},
store_if_persistent(
Session@1,
[{update_packet_state, Id, {unacked_qo_s1, Message}} | Id_updates],
fun(Storage_updates) ->
{{session,
erlang:element(2, Session@1),
erlang:element(3, Session@1),
Unacked_qos1,
erlang:element(5, Session@1),
erlang:element(6, Session@1),
erlang:element(7, Session@1)},
Packet,
Storage_updates}
end
).
-file("src/spoke/core/internal/session.gleam", 154).
?DOC(false).
-spec start_qos2_publish(session(), spoke@packet:message_data()) -> {session(),
spoke@packet@client@outgoing:packet(),
list(spoke@core@session_state:storage_update())}.
start_qos2_publish(Session, Message) ->
{Session@1, Id, Id_updates} = reserve_packet_id(Session),
Unreceived_qos2 = gleam@dict:insert(
erlang:element(5, Session@1),
Id,
Message
),
Packet = {publish, {publish_data_qo_s2, Message, false, Id}},
store_if_persistent(
Session@1,
[{update_packet_state, Id, {unreceived_qo_s2, Message}} | Id_updates],
fun(Storage_updates) ->
{{session,
erlang:element(2, Session@1),
erlang:element(3, Session@1),
erlang:element(4, Session@1),
Unreceived_qos2,
erlang:element(6, Session@1),
erlang:element(7, Session@1)},
Packet,
Storage_updates}
end
).
-file("src/spoke/core/internal/session.gleam", 172).
?DOC(false).
-spec handle_pubrec(session(), integer()) -> {session(),
list(spoke@core@session_state:storage_update())}.
handle_pubrec(Session, Id) ->
Unreceived_qos2 = gleam@dict:delete(erlang:element(5, Session), Id),
Unreleased_qos2 = gleam@set:insert(erlang:element(6, Session), Id),
store_if_persistent(
Session,
[{update_packet_state, Id, received_qo_s2}],
fun(Storage_updates) ->
{{session,
erlang:element(2, Session),
erlang:element(3, Session),
erlang:element(4, Session),
Unreceived_qos2,
Unreleased_qos2,
erlang:element(7, Session)},
Storage_updates}
end
).
-file("src/spoke/core/internal/session.gleam", 185).
?DOC(false).
-spec handle_pubcomp(session(), integer()) -> {session(),
list(spoke@core@session_state:storage_update())}.
handle_pubcomp(Session, Id) ->
Unreleased_qos2 = gleam@set:delete(erlang:element(6, Session), Id),
store_if_persistent(
Session,
[{clear_packet_state, Id}],
fun(Storage_updates) ->
{{session,
erlang:element(2, Session),
erlang:element(3, Session),
erlang:element(4, Session),
erlang:element(5, Session),
Unreleased_qos2,
erlang:element(7, Session)},
Storage_updates}
end
).
-file("src/spoke/core/internal/session.gleam", 197).
?DOC(false).
-spec start_qos2_receive(session(), integer()) -> {session(),
boolean(),
list(spoke@core@session_state:storage_update())}.
start_qos2_receive(Session, Packet_id) ->
case gleam@set:contains(erlang:element(7, Session), Packet_id) of
true ->
{Session, false, []};
false ->
Incomplete_qos2_in = gleam@set:insert(
erlang:element(7, Session),
Packet_id
),
store_if_persistent(
Session,
[{update_packet_state, Packet_id, unreleased_qo_s2}],
fun(Storage_updates) ->
{{session,
erlang:element(2, Session),
erlang:element(3, Session),
erlang:element(4, Session),
erlang:element(5, Session),
erlang:element(6, Session),
Incomplete_qos2_in},
true,
Storage_updates}
end
)
end.
-file("src/spoke/core/internal/session.gleam", 213).
?DOC(false).
-spec handle_puback(session(), integer()) -> pub_ack_result().
handle_puback(Session, Packet_id) ->
case gleam@dict:has_key(erlang:element(4, Session), Packet_id) of
true ->
Unacked_qos1 = gleam@dict:delete(
erlang:element(4, Session),
Packet_id
),
store_if_persistent(
Session,
[{clear_packet_state, Packet_id}],
fun(Storage_updates) ->
{publish_finished,
{session,
erlang:element(2, Session),
erlang:element(3, Session),
Unacked_qos1,
erlang:element(5, Session),
erlang:element(6, Session),
erlang:element(7, Session)},
Storage_updates}
end
);
false ->
invalid_pub_ack_id
end.
-file("src/spoke/core/internal/session.gleam", 226).
?DOC(false).
-spec handle_pubrel(session(), integer()) -> {session(),
list(spoke@core@session_state:storage_update())}.
handle_pubrel(Session, Packet_id) ->
Incomplete_qos2_in = gleam@set:delete(erlang:element(7, Session), Packet_id),
store_if_persistent(
Session,
[{clear_packet_state, Packet_id}],
fun(Storage_updates) ->
{{session,
erlang:element(2, Session),
erlang:element(3, Session),
erlang:element(4, Session),
erlang:element(5, Session),
erlang:element(6, Session),
Incomplete_qos2_in},
Storage_updates}
end
).