Current section

Files

Jump to
ssevents src ssevents@stream.erl
Raw

src/ssevents@stream.erl

-module(ssevents@stream).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/ssevents/stream.gleam").
-export([empty/0, single/1, from_list/1, next/1, map/2, append/2, encode_stream/1, to_list/1, decode_stream_with_limits/2, decode_stream/1]).
-export_type([iterator/1, step/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.
?MODULEDOC(
" Lightweight iterator utilities for chunk-based encode/decode flows.\n"
"\n"
" The current `gleam_stdlib` version used by this repository does\n"
" not ship a general-purpose iterator module, so `ssevents` provides\n"
" a tiny local one for its own stream adapters.\n"
).
-opaque iterator(EGP) :: {iterator, fun(() -> step(EGP))}.
-type step(EGQ) :: {next, EGQ, iterator(EGQ)} | done.
-file("src/ssevents/stream.gleam", 23).
-spec empty() -> iterator(any()).
empty() ->
{iterator, fun() -> done end}.
-file("src/ssevents/stream.gleam", 27).
-spec single(EGT) -> iterator(EGT).
single(Item) ->
{iterator, fun() -> {next, Item, empty()} end}.
-file("src/ssevents/stream.gleam", 31).
-spec from_list(list(EGV)) -> iterator(EGV).
from_list(Items) ->
{iterator, fun() -> case Items of
[] ->
done;
[Item | Rest] ->
{next, Item, from_list(Rest)}
end end}.
-file("src/ssevents/stream.gleam", 40).
-spec next(iterator(EGY)) -> step(EGY).
next(Iterator) ->
(erlang:element(2, Iterator))().
-file("src/ssevents/stream.gleam", 48).
-spec map(iterator(EHE), fun((EHE) -> EHG)) -> iterator(EHG).
map(Iterator, F) ->
{iterator, fun() -> case next(Iterator) of
done ->
done;
{next, Item, Rest} ->
{next, F(Item), map(Rest, F)}
end end}.
-file("src/ssevents/stream.gleam", 57).
-spec append(iterator(EHI), iterator(EHI)) -> iterator(EHI).
append(Left, Right) ->
{iterator, fun() -> case next(Left) of
done ->
next(Right);
{next, Item, Rest} ->
{next, Item, append(Rest, Right)}
end end}.
-file("src/ssevents/stream.gleam", 66).
-spec encode_stream(iterator(ssevents@event:item())) -> iterator(bitstring()).
encode_stream(Items) ->
map(Items, fun ssevents@encoder:encode_item_bytes/1).
-file("src/ssevents/stream.gleam", 88).
-spec to_list_loop(iterator(EHW), list(EHW)) -> list(EHW).
to_list_loop(Iterator, Acc_rev) ->
case next(Iterator) of
done ->
lists:reverse(Acc_rev);
{next, Item, Rest} ->
to_list_loop(Rest, [Item | Acc_rev])
end.
-file("src/ssevents/stream.gleam", 44).
-spec to_list(iterator(EHB)) -> list(EHB).
to_list(Iterator) ->
to_list_loop(Iterator, []).
-file("src/ssevents/stream.gleam", 104).
-spec step_decoded(
ssevents@decoder:decode_state(),
iterator(bitstring()),
list({ok, ssevents@event:item()} | {error, ssevents@error:sse_error()}),
boolean()
) -> step({ok, ssevents@event:item()} | {error, ssevents@error:sse_error()}).
step_decoded(Decoder_state, Chunks, Pending_rev, Finished) ->
case Pending_rev of
[Item | Rest] ->
{next, Item, pull_decoded(Decoder_state, Chunks, Rest, Finished)};
[] ->
case Finished of
true ->
done;
false ->
advance_decoder(Decoder_state, Chunks)
end
end.
-file("src/ssevents/stream.gleam", 129).
-spec advance_decoder(ssevents@decoder:decode_state(), iterator(bitstring())) -> step({ok,
ssevents@event:item()} |
{error, ssevents@error:sse_error()}).
advance_decoder(Decoder_state, Chunks) ->
case next(Chunks) of
{next, Chunk, Rest_chunks} ->
handle_push(
ssevents@decoder:push(Decoder_state, Chunk),
Decoder_state,
Rest_chunks
);
done ->
handle_finish(ssevents@decoder:finish(Decoder_state), Decoder_state)
end.
-file("src/ssevents/stream.gleam", 161).
-spec handle_finish(
{ok, list(ssevents@event:item())} | {error, ssevents@error:sse_error()},
ssevents@decoder:decode_state()
) -> step({ok, ssevents@event:item()} | {error, ssevents@error:sse_error()}).
handle_finish(Result, Decoder_state) ->
case Result of
{error, Error} ->
terminal_error(Error, Decoder_state);
{ok, Emitted} ->
step_decoded(
Decoder_state,
empty(),
begin
_pipe = Emitted,
_pipe@1 = gleam@list:map(
_pipe,
fun(Field@0) -> {ok, Field@0} end
),
lists:reverse(_pipe@1)
end,
true
)
end.
-file("src/ssevents/stream.gleam", 177).
-spec terminal_error(
ssevents@error:sse_error(),
ssevents@decoder:decode_state()
) -> step({ok, ssevents@event:item()} | {error, ssevents@error:sse_error()}).
terminal_error(Error, Decoder_state) ->
{next, {error, Error}, pull_decoded(Decoder_state, empty(), [], true)}.
-file("src/ssevents/stream.gleam", 95).
-spec pull_decoded(
ssevents@decoder:decode_state(),
iterator(bitstring()),
list({ok, ssevents@event:item()} | {error, ssevents@error:sse_error()}),
boolean()
) -> iterator({ok, ssevents@event:item()} | {error, ssevents@error:sse_error()}).
pull_decoded(Decoder_state, Chunks, Pending_rev, Finished) ->
{iterator,
fun() -> step_decoded(Decoder_state, Chunks, Pending_rev, Finished) end}.
-file("src/ssevents/stream.gleam", 76).
-spec decode_stream_with_limits(iterator(bitstring()), ssevents@limit:limits()) -> iterator({ok,
ssevents@event:item()} |
{error, ssevents@error:sse_error()}).
decode_stream_with_limits(Chunks, Limits) ->
pull_decoded(
ssevents@decoder:new_decoder_with_limits(Limits),
Chunks,
[],
false
).
-file("src/ssevents/stream.gleam", 70).
-spec decode_stream(iterator(bitstring())) -> iterator({ok,
ssevents@event:item()} |
{error, ssevents@error:sse_error()}).
decode_stream(Chunks) ->
decode_stream_with_limits(Chunks, ssevents@limit:default()).
-file("src/ssevents/stream.gleam", 144).
-spec handle_push(
{ok, {ssevents@decoder:decode_state(), list(ssevents@event:item())}} |
{error, ssevents@error:sse_error()},
ssevents@decoder:decode_state(),
iterator(bitstring())
) -> step({ok, ssevents@event:item()} | {error, ssevents@error:sse_error()}).
handle_push(Result, Prev_decoder, Rest_chunks) ->
case Result of
{error, Error} ->
terminal_error(Error, Prev_decoder);
{ok, {Next_decoder, Emitted}} ->
step_decoded(
Next_decoder,
Rest_chunks,
begin
_pipe = Emitted,
_pipe@1 = gleam@list:map(
_pipe,
fun(Field@0) -> {ok, Field@0} end
),
lists:reverse(_pipe@1)
end,
false
)
end.