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()}.
-file("src/ssevents/decoder.gleam", 70).
-spec new_decoder_with_limits(ssevents@limit:limits()) -> decode_state().
new_decoder_with_limits(Limits) ->
{decode_state, <<>>, none, [], 0, none, none, 0, Limits}.
-file("src/ssevents/decoder.gleam", 66).
-spec new_decoder() -> decode_state().
new_decoder() ->
new_decoder_with_limits(ssevents@limit:default()).
-file("src/ssevents/decoder.gleam", 245).
-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)},
[]}}
end.
-file("src/ssevents/decoder.gleam", 267).
-spec apply_validated_field(
decode_state(),
{ok, DOH} | {error, ssevents@error:sse_error()},
integer(),
fun((decode_state(), DOH) -> 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)},
Value
),
[]}}
end.
-file("src/ssevents/decoder.gleam", 311).
-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)}.
-file("src/ssevents/decoder.gleam", 323).
-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", 280).
-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", 300).
-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", 362).
-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", 369).
-spec decode_comment_text(binary()) -> binary().
decode_comment_text(Value) ->
trim_optional_leading_space(Value).
-file("src/ssevents/decoder.gleam", 395).
-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", 388).
-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", 373).
-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", 206).
-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)}
end
);
<<"id"/utf8>> ->
apply_validated_field(
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)}
end
);
<<"retry"/utf8>> ->
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)}
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)},
[]}}
end
end.
-file("src/ssevents/decoder.gleam", 404).
-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", 439).
-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", 340).
-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", 336).
-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", 194).
-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", 162).
-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", 417).
-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};
<<10, Rest/binary>> ->
case Acc_rev of
<<13, Acc_rest/bitstring>> ->
{ok,
{some,
{reverse_bytes(Acc_rest, <<>>),
Rest,
Line_bytes - 1}}};
_ ->
{ok,
{some, {reverse_bytes(Acc_rev, <<>>), Rest, Line_bytes}}}
end;
<<Byte, Rest@1/binary>> ->
find_newline(Rest@1, <<Byte, Acc_rev/bitstring>>, Line_bytes + 1);
_ ->
{error, invalid_utf8}
end.
-file("src/ssevents/decoder.gleam", 411).
-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", 121).
-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)},
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", 83).
-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 = {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)},
process_lines(State@1, []).
-file("src/ssevents/decoder.gleam", 447).
-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", 92).
-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);
_ ->
case ends_with_cr(erlang:element(2, State)) of
true ->
{error, unexpected_end};
false ->
case decode_line(erlang:element(2, State)) 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)},
Line,
erlang:byte_size(erlang:element(2, State))
) 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
end.
-file("src/ssevents/decoder.gleam", 52).
-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", 41).
-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", 45).
-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", 37).
-spec decode(binary()) -> {ok, list(ssevents@event:item())} |
{error, ssevents@error:sse_error()}.
decode(Input) ->
decode_with_limits(Input, ssevents@limit:default()).