Current section

Files

Jump to
multipartkit src multipartkit@stream.erl
Raw

src/multipartkit@stream.erl

-module(multipartkit@stream).
-compile([no_auto_import, nowarn_unused_vars, nowarn_unused_function, nowarn_nomatch, inline]).
-define(FILEPATH, "src/multipartkit/stream.gleam").
-export([from_datastream/1, to_datastream/1, parse_stream_with_limits/3, parse_stream/2, from_part/1, drain_body/1]).
-export_type([stream_part/0, stream_state/0, pull_outcome/0, part_outcome/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.
-type stream_part() :: {stream_part,
list({binary(), binary()}),
gleam@option:option(binary()),
gleam@option:option(binary()),
gleam@option:option(binary()),
gleam@yielder:yielder({ok, bitstring()} |
{error, multipartkit@error:multipart_error()})}.
-type stream_state() :: {stream_state,
bitstring(),
integer(),
gleam@yielder:yielder(bitstring()),
boolean(),
integer(),
bitstring(),
multipartkit@limit:limits(),
integer(),
boolean(),
boolean()}.
-type pull_outcome() :: {pulled, stream_state()} | exhausted | over_body.
-type part_outcome() :: {part_ready, stream_part(), stream_state()} |
need_more |
{failed, multipartkit@error:multipart_error()}.
-file("src/multipartkit/stream.gleam", 90).
?DOC(
" Adapter that returns the source unchanged. Exists so callers can write\n"
" `source |> from_datastream` in a pipeline.\n"
).
-spec from_datastream(gleam@yielder:yielder(bitstring())) -> gleam@yielder:yielder(bitstring()).
from_datastream(Source) ->
Source.
-file("src/multipartkit/stream.gleam", 95).
?DOC(" Adapter that returns the source unchanged. Mirror of `from_datastream`.\n").
-spec to_datastream(gleam@yielder:yielder(bitstring())) -> gleam@yielder:yielder(bitstring()).
to_datastream(Source) ->
Source.
-file("src/multipartkit/stream.gleam", 172).
-spec halt_with(stream_state(), multipartkit@error:multipart_error()) -> gleam@yielder:step({ok,
stream_part()} |
{error, multipartkit@error:multipart_error()}, stream_state()).
halt_with(State, Err) ->
{next,
{error, Err},
{stream_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),
erlang:element(10, State),
true}}.
-file("src/multipartkit/stream.gleam", 179).
-spec pull_chunk(stream_state()) -> pull_outcome().
pull_chunk(State) ->
case erlang:element(5, State) of
true ->
exhausted;
false ->
case gleam@yielder:step(erlang:element(4, State)) of
done ->
exhausted;
{next, Chunk, Rest} ->
Size = erlang:byte_size(Chunk),
New_pulled = erlang:element(6, State) + Size,
case New_pulled > erlang:element(
2,
erlang:element(8, State)
) of
true ->
over_body;
false ->
{pulled,
{stream_state,
gleam@bit_array:append(
erlang:element(2, State),
Chunk
),
erlang:element(3, State),
Rest,
erlang:element(5, State),
New_pulled,
erlang:element(7, State),
erlang:element(8, State),
erlang:element(9, State),
erlang:element(10, State),
erlang:element(11, State)}}
end
end
end.
-file("src/multipartkit/stream.gleam", 305).
?DOC(
" Internal: drop already-consumed bytes from the front of `buf` so memory\n"
" is bounded by the current part rather than the entire input.\n"
).
-spec compact_buffer(stream_state()) -> stream_state().
compact_buffer(State) ->
case erlang:element(3, State) of
0 ->
State;
_ ->
Remaining = multipartkit@internal@bytes:drop(
erlang:element(2, State),
erlang:element(3, State)
),
{stream_state,
Remaining,
0,
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),
erlang:element(11, State)}
end.
-file("src/multipartkit/stream.gleam", 256).
-spec finalise_part(
stream_state(),
integer(),
integer(),
multipartkit@internal@scan:delim_kind(),
integer(),
list({binary(), binary()}),
multipartkit@internal@headers:derived_meta()
) -> part_outcome().
finalise_part(
State,
Body_start,
Body_end_excl,
Kind,
After_delim,
Header_list,
Meta
) ->
Body_size = Body_end_excl - Body_start,
case Body_size > erlang:element(3, erlang:element(8, State)) of
true ->
{failed,
{part_too_large, erlang:element(3, erlang:element(8, State))}};
false ->
New_count = erlang:element(9, State) + 1,
case New_count > erlang:element(4, erlang:element(8, State)) of
true ->
{failed,
{too_many_parts,
erlang:element(4, erlang:element(8, State))}};
false ->
Part_body = multipartkit@internal@bytes:slice_or_empty(
erlang:element(2, State),
Body_start,
Body_size
),
Stream_part = {stream_part,
Header_list,
erlang:element(2, Meta),
erlang:element(3, Meta),
erlang:element(4, Meta),
gleam@yielder:from_list([{ok, Part_body}])},
Next_done = case Kind of
closing ->
true;
delimiter ->
false
end,
{part_ready,
Stream_part,
compact_buffer(
{stream_state,
erlang:element(2, State),
After_delim,
erlang:element(4, State),
erlang:element(5, State),
erlang:element(6, State),
erlang:element(7, State),
erlang:element(8, State),
New_count,
erlang:element(10, State),
Next_done}
)}
end
end.
-file("src/multipartkit/stream.gleam", 235).
-spec finalise_body(
stream_state(),
integer(),
list({binary(), binary()}),
multipartkit@internal@headers:derived_meta()
) -> part_outcome().
finalise_body(State, Body_start, Header_list, Meta) ->
case multipartkit@internal@scan:find_delimiter(
erlang:element(2, State),
erlang:element(7, State),
Body_start
) of
incomplete ->
need_more;
{found, Body_end_excl, Kind, After_delim} ->
finalise_part(
State,
Body_start,
Body_end_excl,
Kind,
After_delim,
Header_list,
Meta
)
end.
-file("src/multipartkit/stream.gleam", 212).
-spec finalise_headers(stream_state(), integer(), integer()) -> part_outcome().
finalise_headers(State, Blank_at, Body_start) ->
Header_block_size = Body_start - erlang:element(3, State),
case Header_block_size > erlang:element(5, erlang:element(8, State)) of
true ->
{failed,
{header_too_large, erlang:element(5, erlang:element(8, State))}};
false ->
Header_block = multipartkit@internal@bytes:slice_or_empty(
erlang:element(2, State),
erlang:element(3, State),
Blank_at - erlang:element(3, State)
),
case multipartkit@internal@headers:parse_block(Header_block) of
{error, Err} ->
{failed, Err};
{ok, Header_list} ->
case multipartkit@internal@headers:derive_meta(Header_list) of
{error, Err@1} ->
{failed, Err@1};
{ok, Meta} ->
finalise_body(State, Body_start, Header_list, Meta)
end
end
end.
-file("src/multipartkit/stream.gleam", 205).
-spec parse_one_part(stream_state()) -> part_outcome().
parse_one_part(State) ->
case multipartkit@internal@bytes:find_blank_line(
erlang:element(2, State),
erlang:element(3, State)
) of
{error, nil} ->
need_more;
{ok, {Blank_at, Body_start}} ->
finalise_headers(State, Blank_at, Body_start)
end.
-file("src/multipartkit/stream.gleam", 157).
-spec step_produce_next(stream_state()) -> gleam@yielder:step({ok,
stream_part()} |
{error, multipartkit@error:multipart_error()}, stream_state()).
step_produce_next(State) ->
case parse_one_part(State) of
{part_ready, Stream_part, Next} ->
{next, {ok, Stream_part}, Next};
{failed, Err} ->
halt_with(State, Err);
need_more ->
case pull_chunk(State) of
{pulled, Next@1} ->
step_produce_next(Next@1);
exhausted ->
halt_with(State, unexpected_end_of_input);
over_body ->
halt_with(
State,
{body_too_large,
erlang:element(2, erlang:element(8, State))}
)
end
end.
-file("src/multipartkit/stream.gleam", 139).
-spec step_find_first(stream_state()) -> gleam@yielder:step({ok, stream_part()} |
{error, multipartkit@error:multipart_error()}, stream_state()).
step_find_first(State) ->
case multipartkit@internal@scan:find_delimiter(
erlang:element(2, State),
erlang:element(7, State),
0
) of
{found, _, closing, _} ->
done;
{found, _, delimiter, After_first} ->
step_produce_next(
{stream_state,
erlang:element(2, State),
After_first,
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,
erlang:element(11, State)}
);
incomplete ->
case pull_chunk(State) of
{pulled, Next} ->
step_find_first(Next);
exhausted ->
halt_with(State, unexpected_end_of_input);
over_body ->
halt_with(
State,
{body_too_large,
erlang:element(2, erlang:element(8, State))}
)
end
end.
-file("src/multipartkit/stream.gleam", 126).
-spec step(stream_state()) -> gleam@yielder:step({ok, stream_part()} |
{error, multipartkit@error:multipart_error()}, stream_state()).
step(State) ->
case erlang:element(11, State) of
true ->
done;
false ->
case erlang:element(10, State) of
false ->
step_find_first(State);
true ->
step_produce_next(State)
end
end.
-file("src/multipartkit/stream.gleam", 61).
?DOC(
" Parse a stream of input chunks with caller-supplied limits.\n"
"\n"
" Chunks are pulled lazily: bytes are consumed from `chunks` only as needed\n"
" to deliver the next part's headers and body. `max_body_bytes` is enforced\n"
" incrementally as chunks arrive, so an oversized stream is rejected before\n"
" it is fully buffered. Per-part memory is bounded by `max_part_bytes`.\n"
).
-spec parse_stream_with_limits(
gleam@yielder:yielder(bitstring()),
binary(),
multipartkit@limit:limits()
) -> {ok,
gleam@yielder:yielder({ok, stream_part()} |
{error, multipartkit@error:multipart_error()})} |
{error, multipartkit@error:multipart_error()}.
parse_stream_with_limits(Chunks, Content_type, Limits) ->
case multipartkit@header:boundary(Content_type) of
{error, Err} ->
{error, Err};
{ok, Boundary_value} ->
Pattern = multipartkit@internal@scan:dash_pattern(Boundary_value),
Initial = {stream_state,
<<>>,
0,
Chunks,
false,
0,
Pattern,
Limits,
0,
false,
false},
{ok, gleam@yielder:unfold(Initial, fun step/1)}
end.
-file("src/multipartkit/stream.gleam", 48).
?DOC(
" Parse a stream of input chunks using `default_limits()`.\n"
"\n"
" The outer `Result` reports errors decidable from `content_type` alone.\n"
" Each yielded item is a `Result` because once body consumption begins,\n"
" failures can arise later from message structure, limits, or truncation.\n"
" After the first `Error(_)` is yielded the iterator is exhausted.\n"
).
-spec parse_stream(gleam@yielder:yielder(bitstring()), binary()) -> {ok,
gleam@yielder:yielder({ok, stream_part()} |
{error, multipartkit@error:multipart_error()})} |
{error, multipartkit@error:multipart_error()}.
parse_stream(Chunks, Content_type) ->
parse_stream_with_limits(
Chunks,
Content_type,
multipartkit@limit:default_limits()
).
-file("src/multipartkit/stream.gleam", 319).
?DOC(
" Build a `StreamPart` from a fully buffered `Part`.\n"
"\n"
" Useful when feeding parts into `encode_stream` or when adapting a\n"
" buffered parse result into the streaming API surface.\n"
).
-spec from_part(multipartkit@part:part()) -> stream_part().
from_part(The_part) ->
{stream_part,
erlang:element(2, The_part),
erlang:element(3, The_part),
erlang:element(4, The_part),
erlang:element(5, The_part),
gleam@yielder:from_list([{ok, erlang:element(6, The_part)}])}.
-file("src/multipartkit/stream.gleam", 341).
-spec drain_loop(
gleam@yielder:yielder({ok, bitstring()} |
{error, multipartkit@error:multipart_error()}),
bitstring()
) -> {ok, bitstring()} | {error, multipartkit@error:multipart_error()}.
drain_loop(Source, Acc) ->
case gleam@yielder:step(Source) of
done ->
{ok, Acc};
{next, {ok, Chunk}, Rest} ->
drain_loop(Rest, gleam@bit_array:append(Acc, Chunk));
{next, {error, Err}, _} ->
{error, Err}
end.
-file("src/multipartkit/stream.gleam", 335).
?DOC(
" Consume a `StreamPart`'s body yielder and return the concatenated bytes.\n"
"\n"
" Stops at the first `Error(_)` and returns it. Because `StreamPart.body`\n"
" in v0.1.0 always emits a single buffered chunk, this is a constant-time\n"
" pull over a one-element yielder; future releases that switch to true\n"
" chunked body streaming will still let this helper drain the whole body.\n"
).
-spec drain_body(
gleam@yielder:yielder({ok, bitstring()} |
{error, multipartkit@error:multipart_error()})
) -> {ok, bitstring()} | {error, multipartkit@error:multipart_error()}.
drain_body(Source) ->
drain_loop(Source, <<>>).