Packages
Typed distributed messaging for Gleam on the BEAM.
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, sized_decoder/3, decoder/3, codec/3, encoder/1, encode_tagged/4, decode_tagged/4]).
-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(IWQ) :: {tagged_message, binary(), integer(), IWQ}.
-file("src/distribute/codec/tagged.gleam", 20).
?DOC(" Create a tagged message.\n").
-spec new(binary(), integer(), IWR) -> tagged_message(IWR).
new(Tag, Version, Payload) ->
{tagged_message, Tag, Version, Payload}.
-file("src/distribute/codec/tagged.gleam", 29).
?DOC(" Extract the payload.\n").
-spec payload(tagged_message(IWT)) -> IWT.
payload(Msg) ->
erlang:element(4, Msg).
-file("src/distribute/codec/tagged.gleam", 34).
?DOC(" Get the tag.\n").
-spec tag(tagged_message(any())) -> binary().
tag(Msg) ->
erlang:element(2, Msg).
-file("src/distribute/codec/tagged.gleam", 39).
?DOC(" Get the version.\n").
-spec version(tagged_message(any())) -> integer().
version(Msg) ->
erlang:element(3, Msg).
-file("src/distribute/codec/tagged.gleam", 134).
?DOC(
" Parse and validate the `[tag_len:32][tag][version:32]` header shared by\n"
" `decoder` and `sized_decoder`. Returns the verified tag, version, and\n"
" the payload bytes that follow the header.\n"
).
-spec parse_header(bitstring(), binary(), integer()) -> {ok,
{binary(), integer(), bitstring()}} |
{error, distribute@codec:decode_error()}.
parse_header(Binary, Expected_tag, Expected_version) ->
case Binary of
<<Tag_len:32, Rest/bitstring>> ->
gleam@result:'try'(
begin
_pipe = gleam_stdlib:bit_array_slice(Rest, 0, Tag_len),
gleam@result:replace_error(
_pipe,
{insufficient_data, <<"tag truncated"/utf8>>}
)
end,
fun(Tag_bytes) ->
gleam@result:'try'(
begin
_pipe@1 = gleam@bit_array:to_string(Tag_bytes),
gleam@result:replace_error(
_pipe@1,
{invalid_binary, <<"tag not valid UTF-8"/utf8>>}
)
end,
fun(Tag_string) -> case Tag_string =:= Expected_tag of
false ->
{error,
{tag_mismatch, Expected_tag, Tag_string}};
true ->
gleam@result:'try'(
begin
_pipe@2 = gleam_stdlib:bit_array_slice(
Rest,
Tag_len,
erlang:byte_size(Rest) - Tag_len
),
gleam@result:replace_error(
_pipe@2,
{insufficient_data,
<<"missing version"/utf8>>}
)
end,
fun(After_tag) -> case After_tag of
<<Ver:32,
Payload_rest/bitstring>> ->
case Ver =:= Expected_version of
false ->
{error,
{version_mismatch,
Expected_version,
Ver}};
true ->
{ok,
{Tag_string,
Ver,
Payload_rest}}
end;
_ ->
{error,
{insufficient_data,
<<"missing version or payload"/utf8>>}}
end end
)
end end
)
end
);
_ ->
{error, {invalid_binary, <<"missing tag length"/utf8>>}}
end.
-file("src/distribute/codec/tagged.gleam", 95).
?DOC(
" Sized decoder for tagged messages with tag and version validation.\n"
"\n"
" Returns `TagMismatch` or `VersionMismatch` errors for protocol mismatches.\n"
" Unlike `decoder`, returns remaining bytes for use in composite codecs.\n"
).
-spec sized_decoder(
binary(),
integer(),
fun((bitstring()) -> {ok, {IXD, bitstring()}} |
{error, distribute@codec:decode_error()})
) -> fun((bitstring()) -> {ok, {tagged_message(IXD), bitstring()}} |
{error, distribute@codec:decode_error()}).
sized_decoder(Expected_tag, Expected_version, Payload_sized_decoder) ->
fun(Binary) ->
gleam@result:'try'(
parse_header(Binary, Expected_tag, Expected_version),
fun(_use0) ->
{Tag_string, Ver, Payload_rest} = _use0,
gleam@result:'try'(
Payload_sized_decoder(Payload_rest),
fun(_use0@1) ->
{Payload, Remaining} = _use0@1,
{ok,
{{tagged_message, Tag_string, Ver, Payload},
Remaining}}
end
)
end
)
end.
-file("src/distribute/codec/tagged.gleam", 115).
?DOC(
" Decoder for tagged messages with tag and version validation.\n"
"\n"
" Unlike `sized_decoder`, the payload decoder receives every byte after\n"
" the header and is responsible for its own trailing-byte policy.\n"
).
-spec decoder(
binary(),
integer(),
fun((bitstring()) -> {ok, IXH} | {error, distribute@codec:decode_error()})
) -> fun((bitstring()) -> {ok, tagged_message(IXH)} |
{error, distribute@codec:decode_error()}).
decoder(Expected_tag, Expected_version, Payload_decoder) ->
fun(Binary) ->
gleam@result:'try'(
parse_header(Binary, Expected_tag, Expected_version),
fun(_use0) ->
{Tag_string, Ver, Payload_rest} = _use0,
gleam@result:'try'(
distribute@codec:decode(Payload_decoder, Payload_rest),
fun(Payload) ->
{ok, {tagged_message, Tag_string, Ver, Payload}}
end
)
end
)
end.
-file("src/distribute/codec/tagged.gleam", 181).
?DOC(
" Bundled `Codec` for tagged messages.\n"
"\n"
" ```gleam\n"
" let tagged_int = tagged.codec(\"counter\", 1, codec.int())\n"
" ```\n"
).
-spec codec(binary(), integer(), distribute@codec:codec(IXN)) -> distribute@codec:codec(tagged_message(IXN)).
codec(Expected_tag, Expected_version, Payload_codec) ->
{codec,
encoder(erlang:element(2, Payload_codec)),
decoder(
Expected_tag,
Expected_version,
erlang:element(3, Payload_codec)
),
sized_decoder(
Expected_tag,
Expected_version,
erlang:element(4, Payload_codec)
)}.
-file("src/distribute/codec/tagged.gleam", 52).
?DOC(
" Encoder for tagged messages.\n"
"\n"
" Format: `[tag_len:32][tag:utf8][version:32][payload]`\n"
"\n"
" Rejects values that would silently wrap when written to the fixed-width\n"
" 32-bit fields:\n"
" - `version` must be in `[0, 2^32 - 1]`\n"
" - `tag` byte-length must be in `[0, 2^32 - 1]` (a 4 GiB tag is absurd\n"
" in practice but the bound is enforced for defence in depth)\n"
).
-spec encoder(
fun((IWZ) -> {ok, bitstring()} | {error, distribute@codec:encode_error()})
) -> fun((tagged_message(IWZ)) -> {ok, bitstring()} |
{error, distribute@codec:encode_error()}).
encoder(Payload_encoder) ->
fun(Msg) ->
case (erlang:element(3, Msg) < 0) orelse (erlang:element(3, Msg) > 4294967295) of
true ->
{error,
{value_too_large,
<<<<"tagged version "/utf8,
(erlang:integer_to_binary(
erlang:element(3, Msg)
))/binary>>/binary,
" out of unsigned 32-bit range"/utf8>>}};
false ->
Tag_bytes = gleam_stdlib:identity(erlang:element(2, Msg)),
Tag_len = erlang:byte_size(Tag_bytes),
case Tag_len > 4294967295 of
true ->
{error,
{value_too_large,
<<<<"tagged tag length "/utf8,
(erlang:integer_to_binary(Tag_len))/binary>>/binary,
" exceeds unsigned 32-bit range"/utf8>>}};
false ->
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
end
end.
-file("src/distribute/codec/tagged.gleam", 198).
?DOC(" Convenience: encode a value as a tagged message in one step.\n").
-spec encode_tagged(
binary(),
integer(),
fun((IXR) -> {ok, bitstring()} | {error, distribute@codec:encode_error()}),
IXR
) -> {ok, bitstring()} | {error, distribute@codec:encode_error()}.
encode_tagged(Tag, Version, Payload_encoder, Value) ->
Msg = new(Tag, Version, Value),
distribute@codec:encode(encoder(Payload_encoder), Msg).
-file("src/distribute/codec/tagged.gleam", 209).
?DOC(" Convenience: decode binary as a tagged message and extract the payload.\n").
-spec decode_tagged(
binary(),
integer(),
fun((bitstring()) -> {ok, IXV} | {error, distribute@codec:decode_error()}),
bitstring()
) -> {ok, IXV} | {error, distribute@codec:decode_error()}.
decode_tagged(Expected_tag, Expected_version, Payload_decoder, Data) ->
gleam@result:'try'(
distribute@codec:decode(
decoder(Expected_tag, Expected_version, Payload_decoder),
Data
),
fun(Msg) -> {ok, payload(Msg)} end
).