Packages
Typed distributed messaging for Gleam on the BEAM.
Retired package: Deprecated - The project needs to be redesigned around a much smaller and clearer core.
Current section
Files
Jump to
Current section
Files
src/distribute@codec@tagged.erl
-module(distribute@codec@tagged).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/distribute/codec/tagged.gleam").
-export([new/3, payload/1, tag/1, version/1, encoder/1, decoder/3]).
-export_type([tagged_message/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.
-opaque tagged_message(IMD) :: {tagged_message, binary(), integer(), IMD}.
-file("src/distribute/codec/tagged.gleam", 16).
?DOC(" Create a tagged message.\n").
-spec new(binary(), integer(), IME) -> tagged_message(IME).
new(Tag, Version, Payload) ->
{tagged_message, Tag, Version, Payload}.
-file("src/distribute/codec/tagged.gleam", 25).
?DOC(" Extract the payload from a tagged message.\n").
-spec payload(tagged_message(IMG)) -> IMG.
payload(Msg) ->
erlang:element(4, Msg).
-file("src/distribute/codec/tagged.gleam", 30).
?DOC(" Get the tag from a tagged message.\n").
-spec tag(tagged_message(any())) -> binary().
tag(Msg) ->
erlang:element(2, Msg).
-file("src/distribute/codec/tagged.gleam", 35).
?DOC(" Get the version from a tagged message.\n").
-spec version(tagged_message(any())) -> integer().
version(Msg) ->
erlang:element(3, Msg).
-file("src/distribute/codec/tagged.gleam", 41).
?DOC(
" Create an encoder for tagged messages.\n"
" This embeds the tag and version in the binary format.\n"
).
-spec encoder(
fun((IMM) -> {ok, bitstring()} | {error, distribute@codec:encode_error()})
) -> fun((tagged_message(IMM)) -> {ok, bitstring()} |
{error, distribute@codec:encode_error()}).
encoder(Payload_encoder) ->
fun(Msg) ->
Tag_bytes = gleam_stdlib:identity(erlang:element(2, Msg)),
Tag_len = erlang:byte_size(Tag_bytes),
gleam@result:'try'(
distribute@codec:encode(Payload_encoder, erlang:element(4, Msg)),
fun(Payload_bytes) ->
{ok,
<<Tag_len:32,
Tag_bytes/bitstring,
(erlang:element(3, Msg)):32,
Payload_bytes/bitstring>>}
end
)
end.
-file("src/distribute/codec/tagged.gleam", 59).
?DOC(
" Create a decoder for tagged messages with validation.\n"
" This ensures received messages match the expected tag and version.\n"
).
-spec decoder(
binary(),
integer(),
fun((bitstring()) -> {ok, IMQ} | {error, distribute@codec:decode_error()})
) -> fun((bitstring()) -> {ok, tagged_message(IMQ)} |
{error, distribute@codec:decode_error()}).
decoder(Expected_tag, Expected_version, Payload_decoder) ->
fun(Binary) -> case Binary of
<<Tag_len:32, Rest/bitstring>> ->
case gleam_stdlib:bit_array_slice(Rest, 0, Tag_len) of
{ok, Tag_bytes} ->
Tag = gleam@bit_array:to_string(Tag_bytes),
case Tag of
{ok, Tag_string} ->
case Tag_string =:= Expected_tag of
false ->
{error,
{type_mismatch,
<<<<<<"Tag mismatch: expected "/utf8,
Expected_tag/binary>>/binary,
", got "/utf8>>/binary,
Tag_string/binary>>}};
true ->
Remaining = gleam_stdlib:bit_array_slice(
Rest,
Tag_len,
erlang:byte_size(Rest) - Tag_len
),
case Remaining of
{ok,
<<Version:32,
Payload_bytes/bitstring>>} ->
case Version =:= Expected_version of
false ->
{error,
{type_mismatch,
<<<<<<"Version mismatch: expected "/utf8,
(erlang:integer_to_binary(
Expected_version
))/binary>>/binary,
", got "/utf8>>/binary,
(erlang:integer_to_binary(
Version
))/binary>>}};
true ->
gleam@result:'try'(
distribute@codec:decode(
Payload_decoder,
Payload_bytes
),
fun(Payload) ->
{ok,
{tagged_message,
Tag_string,
Version,
Payload}}
end
)
end;
_ ->
{error,
{insufficient_data,
<<"Missing version or payload"/utf8>>}}
end
end;
{error, _} ->
{error,
{invalid_binary,
<<"Tag is not valid UTF-8"/utf8>>}}
end;
{error, _} ->
{error, {insufficient_data, <<"Tag truncated"/utf8>>}}
end;
_ ->
{error, {invalid_binary, <<"Missing tag length"/utf8>>}}
end end.