Current section

Files

Jump to
ssevents src ssevents@decoder.erl
Raw

src/ssevents@decoder.erl

-module(ssevents@decoder).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/ssevents/decoder.gleam").
-export([new_decoder_with_limits/1, new_decoder/0, push/2, finish/1, decode_bytes_with_limits/2, decode_bytes/1, decode_with_limits/2, decode/1]).
-export_type([decode_state/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(
" Full-body and incremental SSE decoding.\n"
"\n"
" Semantics chosen for the initial release:\n"
" - accepted line endings: LF and CRLF\n"
" - unknown fields are ignored\n"
" - EOF dispatches the final unterminated event or trailing comment\n"
" - the first decode error fails the whole operation\n"
" - retry values must be ASCII digits and must not exceed\n"
" `Limits.max_retry_value`\n"
).
-opaque decode_state() :: {decode_state,
bitstring(),
gleam@option:option(binary()),
list(binary()),
integer(),
gleam@option:option(binary()),
gleam@option:option(integer()),
integer(),
ssevents@limit:limits(),
boolean()}.
-file("src/ssevents/decoder.gleam", 76).
-spec new_decoder_with_limits(ssevents@limit:limits()) -> decode_state().
new_decoder_with_limits(Limits) ->
{decode_state, <<>>, none, [], 0, none, none, 0, Limits, false}.
-file("src/ssevents/decoder.gleam", 72).
-spec new_decoder() -> decode_state().
new_decoder() ->
new_decoder_with_limits(ssevents@limit:default()).
-file("src/ssevents/decoder.gleam", 107).
?DOC(
" WHATWG SSE §9.2.5: discard a leading U+FEFF BYTE ORDER MARK at the\n"
" very start of the stream.\n"
"\n"
" `state.bom_handled` is set once we've either stripped the BOM or\n"
" established the stream doesn't start with one. While the buffer\n"
" holds only a 1- or 2-byte prefix of the BOM (`EF` or `EF BB`), we\n"
" hold off on deciding so the check stays correct even when the BOM\n"
" is split across two `push` calls.\n"
).
-spec maybe_strip_bom(decode_state()) -> decode_state().
maybe_strip_bom(State) ->
case {erlang:element(10, State), erlang:element(2, State)} of
{true, _} ->
State;
{false, <<16#EF, 16#BB, 16#BF, Rest/binary>>} ->
{decode_state,
Rest,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
erlang:element(9, State),
true};
{false, <<16#EF, 16#BB>>} ->
State;
{false, <<16#EF>>} ->
State;
{false, <<>>} ->
State;
{false, _} ->
{decode_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
erlang:element(9, State),
true}
end.
-file("src/ssevents/decoder.gleam", 308).
-spec apply_data_field(decode_state(), binary(), integer()) -> {ok,
{decode_state(), list(ssevents@event:item())}} |
{error, ssevents@error:sse_error()}.
apply_data_field(State, Value, Next_event_bytes) ->
case (erlang:element(5, State) + 1) > ssevents@limit:max_data_lines(
erlang:element(9, State)
) of
true ->
{error,
{too_many_data_lines,
ssevents@limit:max_data_lines(erlang:element(9, State))}};
false ->
{ok,
{{decode_state,
erlang:element(2, State),
erlang:element(3, State),
[Value | erlang:element(4, State)],
erlang:element(5, State) + 1,
erlang:element(6, State),
erlang:element(7, State),
Next_event_bytes,
erlang:element(9, State),
erlang:element(10, State)},
[]}}
end.
-file("src/ssevents/decoder.gleam", 330).
-spec apply_validated_field(
decode_state(),
{ok, DOO} | {error, ssevents@error:sse_error()},
integer(),
fun((decode_state(), DOO) -> decode_state())
) -> {ok, {decode_state(), list(ssevents@event:item())}} |
{error, ssevents@error:sse_error()}.
apply_validated_field(State, Validated, Next_event_bytes, Set) ->
case Validated of
{error, Error} ->
{error, Error};
{ok, Value} ->
{ok,
{Set(
{decode_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Next_event_bytes,
erlang:element(9, State),
erlang:element(10, State)},
Value
),
[]}}
end.
-file("src/ssevents/decoder.gleam", 348).
?DOC(
" Like `apply_validated_field` but treats validation failure as a\n"
" silent drop of the offending field (the line is still counted\n"
" against `event_bytes`, the buffered event continues to be\n"
" assembled). Used by fields whose WHATWG SSE §9.2.6 contract is\n"
" \"ignore the field\" on invalid input.\n"
).
-spec apply_validated_field_or_ignore(
decode_state(),
{ok, DOU} | {error, ssevents@error:sse_error()},
integer(),
fun((decode_state(), DOU) -> decode_state())
) -> {ok, {decode_state(), list(ssevents@event:item())}} |
{error, ssevents@error:sse_error()}.
apply_validated_field_or_ignore(State, Validated, Next_event_bytes, Set) ->
case Validated of
{error, _} ->
{ok,
{{decode_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Next_event_bytes,
erlang:element(9, State),
erlang:element(10, State)},
[]}};
{ok, Value} ->
{ok,
{Set(
{decode_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Next_event_bytes,
erlang:element(9, State),
erlang:element(10, State)},
Value
),
[]}}
end.
-file("src/ssevents/decoder.gleam", 392).
-spec reset_event_state(decode_state()) -> decode_state().
reset_event_state(State) ->
{decode_state,
erlang:element(2, State),
none,
[],
0,
none,
none,
0,
erlang:element(9, State),
erlang:element(10, State)}.
-file("src/ssevents/decoder.gleam", 404).
-spec has_meaningful_event_content(decode_state()) -> boolean().
has_meaningful_event_content(State) ->
case erlang:element(5, State) > 0 of
true ->
true;
false ->
case {erlang:element(3, State),
erlang:element(6, State),
erlang:element(7, State)} of
{{some, _}, _, _} ->
true;
{_, {some, _}, _} ->
true;
{_, _, {some, _}} ->
true;
{none, none, none} ->
false
end
end.
-file("src/ssevents/decoder.gleam", 361).
-spec dispatch_event(decode_state()) -> {ok,
{decode_state(), gleam@option:option(ssevents@event:item())}} |
{error, ssevents@error:sse_error()}.
dispatch_event(State) ->
Emitted = case has_meaningful_event_content(State) of
true ->
{some,
begin
_pipe@2 = ssevents@event:from_parts(
erlang:element(3, State),
begin
_pipe = erlang:element(4, State),
_pipe@1 = lists:reverse(_pipe),
gleam@string:join(_pipe@1, <<"\n"/utf8>>)
end,
erlang:element(6, State),
erlang:element(7, State)
),
ssevents@event:event_item(_pipe@2)
end};
false ->
none
end,
{ok, {reset_event_state(State), Emitted}}.
-file("src/ssevents/decoder.gleam", 381).
-spec finish_event(decode_state()) -> {ok, list(ssevents@event:item())} |
{error, ssevents@error:sse_error()}.
finish_event(State) ->
case dispatch_event(State) of
{error, Error} ->
{error, Error};
{ok, {_, Maybe_item}} ->
case Maybe_item of
{some, Item} ->
{ok, [Item]};
none ->
{ok, []}
end
end.
-file("src/ssevents/decoder.gleam", 443).
-spec trim_optional_leading_space(binary()) -> binary().
trim_optional_leading_space(Value) ->
case gleam_stdlib:string_starts_with(Value, <<" "/utf8>>) of
true ->
gleam@string:drop_start(Value, 1);
false ->
Value
end.
-file("src/ssevents/decoder.gleam", 450).
-spec decode_comment_text(binary()) -> binary().
decode_comment_text(Value) ->
trim_optional_leading_space(Value).
-file("src/ssevents/decoder.gleam", 476).
-spec ascii_digits_only(bitstring()) -> boolean().
ascii_digits_only(Bits) ->
case Bits of
<<>> ->
true;
<<Digit, Rest/binary>> when (Digit >= 48) andalso (Digit =< 57) ->
ascii_digits_only(Rest);
_ ->
false
end.
-file("src/ssevents/decoder.gleam", 469).
-spec is_ascii_digit_string(binary()) -> boolean().
is_ascii_digit_string(Value) ->
case Value of
<<""/utf8>> ->
false;
_ ->
ascii_digits_only(gleam_stdlib:identity(Value))
end.
-file("src/ssevents/decoder.gleam", 454).
-spec parse_retry(binary(), ssevents@limit:limits()) -> {ok, integer()} |
{error, ssevents@error:sse_error()}.
parse_retry(Value, Limits) ->
case is_ascii_digit_string(Value) of
false ->
{error, {invalid_retry, Value}};
true ->
case gleam_stdlib:parse_int(Value) of
{error, _} ->
{error, {invalid_retry, Value}};
{ok, Parsed} ->
case Parsed > ssevents@limit:max_retry_value(Limits) of
true ->
{error, {invalid_retry, Value}};
false ->
{ok, Parsed}
end
end
end.
-file("src/ssevents/decoder.gleam", 250).
-spec apply_field_parts(decode_state(), binary(), binary(), integer()) -> {ok,
{decode_state(), list(ssevents@event:item())}} |
{error, ssevents@error:sse_error()}.
apply_field_parts(State, Field, Value, Line_byte_size) ->
Next_event_bytes = erlang:element(8, State) + Line_byte_size,
case Next_event_bytes > ssevents@limit:max_event_bytes(
erlang:element(9, State)
) of
true ->
{error,
{event_too_large,
ssevents@limit:max_event_bytes(erlang:element(9, State))}};
false ->
case Field of
<<"data"/utf8>> ->
apply_data_field(State, Value, Next_event_bytes);
<<"event"/utf8>> ->
apply_validated_field(
State,
ssevents@validate:validate_event_name(Value),
Next_event_bytes,
fun(S, V) ->
{decode_state,
erlang:element(2, S),
{some, V},
erlang:element(4, S),
erlang:element(5, S),
erlang:element(6, S),
erlang:element(7, S),
erlang:element(8, S),
erlang:element(9, S),
erlang:element(10, S)}
end
);
<<"id"/utf8>> ->
apply_validated_field_or_ignore(
State,
ssevents@validate:validate_id(Value),
Next_event_bytes,
fun(S@1, V@1) ->
{decode_state,
erlang:element(2, S@1),
erlang:element(3, S@1),
erlang:element(4, S@1),
erlang:element(5, S@1),
{some, V@1},
erlang:element(7, S@1),
erlang:element(8, S@1),
erlang:element(9, S@1),
erlang:element(10, S@1)}
end
);
<<"retry"/utf8>> ->
case is_ascii_digit_string(Value) of
false ->
{ok,
{{decode_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Next_event_bytes,
erlang:element(9, State),
erlang:element(10, State)},
[]}};
true ->
apply_validated_field(
State,
parse_retry(Value, erlang:element(9, State)),
Next_event_bytes,
fun(S@2, V@2) ->
{decode_state,
erlang:element(2, S@2),
erlang:element(3, S@2),
erlang:element(4, S@2),
erlang:element(5, S@2),
erlang:element(6, S@2),
{some, V@2},
erlang:element(8, S@2),
erlang:element(9, S@2),
erlang:element(10, S@2)}
end
)
end;
_ ->
{ok,
{{decode_state,
erlang:element(2, State),
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
Next_event_bytes,
erlang:element(9, State),
erlang:element(10, State)},
[]}}
end
end.
-file("src/ssevents/decoder.gleam", 485).
-spec decode_line(bitstring()) -> {ok, binary()} |
{error, ssevents@error:sse_error()}.
decode_line(Line) ->
case gleam@bit_array:to_string(Line) of
{ok, Text} ->
{ok, Text};
{error, _} ->
{error, invalid_utf8}
end.
-file("src/ssevents/decoder.gleam", 533).
-spec reverse_bytes(bitstring(), bitstring()) -> bitstring().
reverse_bytes(Input, Acc) ->
case Input of
<<>> ->
Acc;
<<Byte, Rest/binary>> ->
reverse_bytes(Rest, <<Byte, Acc/bitstring>>);
_ ->
Acc
end.
-file("src/ssevents/decoder.gleam", 421).
-spec split_field_bytes(bitstring(), bitstring(), binary()) -> {ok,
{binary(), binary()}} |
{error, ssevents@error:sse_error()}.
split_field_bytes(Remaining, Field_rev, Fallback) ->
case Remaining of
<<>> ->
{ok, {Fallback, <<""/utf8>>}};
<<58, Rest/binary>> ->
case gleam@bit_array:to_string(reverse_bytes(Field_rev, <<>>)) of
{error, _} ->
{error, invalid_utf8};
{ok, Field} ->
case gleam@bit_array:to_string(Rest) of
{error, _} ->
{error, invalid_utf8};
{ok, Value} ->
{ok, {Field, trim_optional_leading_space(Value)}}
end
end;
<<Byte, Rest@1/binary>> ->
split_field_bytes(Rest@1, <<Byte, Field_rev/bitstring>>, Fallback);
_ ->
{error, invalid_utf8}
end.
-file("src/ssevents/decoder.gleam", 417).
-spec split_field(binary()) -> {ok, {binary(), binary()}} |
{error, ssevents@error:sse_error()}.
split_field(Line) ->
split_field_bytes(gleam_stdlib:identity(Line), <<>>, Line).
-file("src/ssevents/decoder.gleam", 238).
-spec apply_field(decode_state(), binary(), integer()) -> {ok,
{decode_state(), list(ssevents@event:item())}} |
{error, ssevents@error:sse_error()}.
apply_field(State, Line, Line_byte_size) ->
case split_field(Line) of
{error, Error} ->
{error, Error};
{ok, {Field, Value}} ->
apply_field_parts(State, Field, Value, Line_byte_size)
end.
-file("src/ssevents/decoder.gleam", 206).
-spec process_line(decode_state(), binary(), integer()) -> {ok,
{decode_state(), list(ssevents@event:item())}} |
{error, ssevents@error:sse_error()}.
process_line(State, Line, Line_byte_size) ->
case Line of
<<""/utf8>> ->
case dispatch_event(State) of
{error, Error} ->
{error, Error};
{ok, {Next_state, Maybe_item}} ->
case Maybe_item of
{some, Item} ->
{ok, {Next_state, [Item]}};
none ->
{ok, {Next_state, []}}
end
end;
_ ->
case gleam_stdlib:string_starts_with(Line, <<":"/utf8>>) of
true ->
{ok,
{State,
[{comment,
decode_comment_text(
gleam@string:drop_start(Line, 1)
)}]}};
false ->
apply_field(State, Line, Line_byte_size)
end
end.
-file("src/ssevents/decoder.gleam", 498).
-spec find_newline(bitstring(), bitstring(), integer()) -> {ok,
gleam@option:option({bitstring(), bitstring(), integer()})} |
{error, ssevents@error:sse_error()}.
find_newline(Remaining, Acc_rev, Line_bytes) ->
case Remaining of
<<>> ->
{ok, none};
<<13, 10, Rest/binary>> ->
{ok, {some, {reverse_bytes(Acc_rev, <<>>), Rest, Line_bytes}}};
<<13>> ->
{ok, none};
<<13, After_cr/binary>> ->
{ok, {some, {reverse_bytes(Acc_rev, <<>>), After_cr, Line_bytes}}};
<<10, Rest@1/binary>> ->
{ok, {some, {reverse_bytes(Acc_rev, <<>>), Rest@1, Line_bytes}}};
<<Byte, Rest@2/binary>> ->
find_newline(Rest@2, <<Byte, Acc_rev/bitstring>>, Line_bytes + 1);
_ ->
{error, invalid_utf8}
end.
-file("src/ssevents/decoder.gleam", 492).
-spec next_complete_line(bitstring()) -> {ok,
gleam@option:option({bitstring(), bitstring(), integer()})} |
{error, ssevents@error:sse_error()}.
next_complete_line(Buffer) ->
find_newline(Buffer, <<>>, 0).
-file("src/ssevents/decoder.gleam", 165).
-spec process_lines(decode_state(), list(ssevents@event:item())) -> {ok,
{decode_state(), list(ssevents@event:item())}} |
{error, ssevents@error:sse_error()}.
process_lines(State, Emitted_rev) ->
case next_complete_line(erlang:element(2, State)) of
{ok, {some, {Line_bytes, Rest, Line_byte_size}}} ->
case Line_byte_size > ssevents@limit:max_line_bytes(
erlang:element(9, State)
) of
true ->
{error,
{line_too_long,
ssevents@limit:max_line_bytes(
erlang:element(9, State)
)}};
false ->
case decode_line(Line_bytes) of
{error, Error} ->
{error, Error};
{ok, Line} ->
case process_line(
{decode_state,
Rest,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
erlang:element(9, State),
erlang:element(10, State)},
Line,
Line_byte_size
) of
{error, Error@1} ->
{error, Error@1};
{ok, {Next_state, Emitted}} ->
process_lines(
Next_state,
begin
_pipe = lists:reverse(Emitted),
lists:append(_pipe, Emitted_rev)
end
)
end
end
end;
{ok, none} ->
case erlang:byte_size(erlang:element(2, State)) > ssevents@limit:max_line_bytes(
erlang:element(9, State)
) of
true ->
{error,
{line_too_long,
ssevents@limit:max_line_bytes(
erlang:element(9, State)
)}};
false ->
{ok, {State, lists:reverse(Emitted_rev)}}
end;
{error, Error@2} ->
{error, Error@2}
end.
-file("src/ssevents/decoder.gleam", 90).
-spec push(decode_state(), bitstring()) -> {ok,
{decode_state(), list(ssevents@event:item())}} |
{error, ssevents@error:sse_error()}.
push(State, Chunk) ->
Combined = gleam@bit_array:append(erlang:element(2, State), Chunk),
State@1 = begin
_pipe = {decode_state,
Combined,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
erlang:element(9, State),
erlang:element(10, State)},
maybe_strip_bom(_pipe)
end,
process_lines(State@1, []).
-file("src/ssevents/decoder.gleam", 541).
-spec ends_with_cr(bitstring()) -> boolean().
ends_with_cr(Bits) ->
Size = erlang:byte_size(Bits),
case Size of
0 ->
false;
_ ->
gleam_stdlib:bit_array_slice(Bits, Size - 1, 1) =:= {ok, <<13>>}
end.
-file("src/ssevents/decoder.gleam", 152).
-spec strip_trailing_cr(bitstring()) -> bitstring().
strip_trailing_cr(Bits) ->
case ends_with_cr(Bits) of
false ->
Bits;
true ->
Size = erlang:byte_size(Bits),
case gleam_stdlib:bit_array_slice(Bits, 0, Size - 1) of
{ok, Prefix} ->
Prefix;
{error, _} ->
Bits
end
end.
-file("src/ssevents/decoder.gleam", 121).
-spec finish(decode_state()) -> {ok, list(ssevents@event:item())} |
{error, ssevents@error:sse_error()}.
finish(State) ->
case erlang:element(2, State) of
<<>> ->
finish_event(State);
_ ->
Trimmed = strip_trailing_cr(erlang:element(2, State)),
case decode_line(Trimmed) of
{error, Error} ->
{error, Error};
{ok, Line} ->
case process_line(
{decode_state,
<<>>,
erlang:element(3, State),
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
erlang:element(9, State),
erlang:element(10, State)},
Line,
erlang:byte_size(Trimmed)
) of
{error, Error@1} ->
{error, Error@1};
{ok, {State_after_line, Emitted}} ->
case finish_event(State_after_line) of
{error, Error@2} ->
{error, Error@2};
{ok, Trailing} ->
{ok, lists:append(Emitted, Trailing)}
end
end
end
end.
-file("src/ssevents/decoder.gleam", 58).
-spec decode_bytes_with_limits(bitstring(), ssevents@limit:limits()) -> {ok,
list(ssevents@event:item())} |
{error, ssevents@error:sse_error()}.
decode_bytes_with_limits(Input, Limits) ->
case push(new_decoder_with_limits(Limits), Input) of
{error, Error} ->
{error, Error};
{ok, {State, Items}} ->
case finish(State) of
{error, Error@1} ->
{error, Error@1};
{ok, Trailing} ->
{ok, lists:append(Items, Trailing)}
end
end.
-file("src/ssevents/decoder.gleam", 47).
-spec decode_bytes(bitstring()) -> {ok, list(ssevents@event:item())} |
{error, ssevents@error:sse_error()}.
decode_bytes(Input) ->
decode_bytes_with_limits(Input, ssevents@limit:default()).
-file("src/ssevents/decoder.gleam", 51).
-spec decode_with_limits(binary(), ssevents@limit:limits()) -> {ok,
list(ssevents@event:item())} |
{error, ssevents@error:sse_error()}.
decode_with_limits(Input, Limits) ->
decode_bytes_with_limits(gleam_stdlib:identity(Input), Limits).
-file("src/ssevents/decoder.gleam", 43).
-spec decode(binary()) -> {ok, list(ssevents@event:item())} |
{error, ssevents@error:sse_error()}.
decode(Input) ->
decode_with_limits(Input, ssevents@limit:default()).