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([all_headers/1, name/1, filename/1, content_type/1, body/1, from_datastream/1, to_datastream/1, drain_body/1, parse_stream_with_limits/3, parse_stream/2, from_part/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.
-opaque 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", 47).
?DOC(" All headers as `(name, value)` pairs in input order.\n").
-spec all_headers(stream_part()) -> list({binary(), binary()}).
all_headers(Stream_part) ->
erlang:element(2, Stream_part).
-file("src/multipartkit/stream.gleam", 53).
?DOC(
" The convenience `name` field derived from `Content-Disposition` for\n"
" `form-data` parts, or `None`.\n"
).
-spec name(stream_part()) -> gleam@option:option(binary()).
name(Stream_part) ->
erlang:element(3, Stream_part).
-file("src/multipartkit/stream.gleam", 59).
?DOC(
" The convenience `filename` field derived from `Content-Disposition`\n"
" for `form-data` parts, or `None`.\n"
).
-spec filename(stream_part()) -> gleam@option:option(binary()).
filename(Stream_part) ->
erlang:element(4, Stream_part).
-file("src/multipartkit/stream.gleam", 64).
?DOC(" The `Content-Type` header value, or `None`.\n").
-spec content_type(stream_part()) -> gleam@option:option(binary()).
content_type(Stream_part) ->
erlang:element(5, Stream_part).
-file("src/multipartkit/stream.gleam", 70).
?DOC(
" The single-pass body yielder. See the type-level doc for streaming\n"
" semantics.\n"
).
-spec body(stream_part()) -> gleam@yielder:yielder({ok, bitstring()} |
{error, multipartkit@error:multipart_error()}).
body(Stream_part) ->
erlang:element(6, Stream_part).
-file("src/multipartkit/stream.gleam", 128).
?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", 133).
?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", 220).
-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", 227).
-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 > multipartkit@limit:max_body_bytes(
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", 367).
?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", 409).
?DOC(
" Internal: split a fully buffered body into a yielder of fixed-size\n"
" chunks of at most `chunk_size` bytes each. The last chunk may be\n"
" smaller. An empty body still emits one empty chunk so that callers\n"
" that fold over `Ok(_)` items see at least one body item per part.\n"
).
-spec chunk_body(bitstring(), integer()) -> gleam@yielder:yielder({ok,
bitstring()} |
{error, multipartkit@error:multipart_error()}).
chunk_body(Body, Chunk_size) ->
Total = erlang:byte_size(Body),
case Total of
0 ->
gleam@yielder:from_list([{ok, <<>>}]);
_ ->
gleam@yielder:unfold(0, fun(Offset) -> case Offset >= Total of
true ->
done;
false ->
Take = case (Offset + Chunk_size) > Total of
true ->
Total - Offset;
false ->
Chunk_size
end,
Chunk = multipartkit@internal@bytes:slice_or_empty(
Body,
Offset,
Take
),
{next, {ok, Chunk}, Offset + Take}
end end)
end.
-file("src/multipartkit/stream.gleam", 433).
-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", 399).
?DOC(
" Consume a `StreamPart`'s body yielder and return the concatenated bytes.\n"
"\n"
" Stops at the first `Error(_)` and returns it. The body yielder may\n"
" emit multiple `Ok(BitArray)` items for large parts; this helper folds\n"
" them into a single `BitArray` for callers that only need the full\n"
" body in memory.\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, <<>>).
-file("src/multipartkit/stream.gleam", 318).
-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 > multipartkit@limit:max_part_bytes(erlang:element(8, State)) of
true ->
{failed,
{part_too_large,
multipartkit@limit:max_part_bytes(erlang:element(8, State))}};
false ->
New_count = erlang:element(9, State) + 1,
case New_count > multipartkit@limit:max_parts(
erlang:element(8, State)
) of
true ->
{failed,
{too_many_parts,
multipartkit@limit:max_parts(
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),
chunk_body(Part_body, 65536)},
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", 283).
-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
{found, Body_end_excl, Kind, After_delim} ->
finalise_part(
State,
Body_start,
Body_end_excl,
Kind,
After_delim,
Header_list,
Meta
);
incomplete ->
Max = multipartkit@limit:max_part_bytes(erlang:element(8, State)),
Buf_size = erlang:byte_size(erlang:element(2, State)),
Safety = erlang:byte_size(erlang:element(7, State)) + 4,
Confirmed_body_size = (Buf_size - Body_start) - Safety,
case Confirmed_body_size > Max of
true ->
{failed, {part_too_large, Max}};
false ->
need_more
end
end.
-file("src/multipartkit/stream.gleam", 260).
-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 > multipartkit@limit:max_header_bytes(
erlang:element(8, State)
) of
true ->
{failed,
{header_too_large,
multipartkit@limit:max_header_bytes(
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", 253).
-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", 204).
-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,
multipartkit@limit:max_body_bytes(
erlang:element(8, State)
)}
)
end
end.
-file("src/multipartkit/stream.gleam", 185).
-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,
multipartkit@limit:max_body_bytes(
erlang:element(8, State)
)}
)
end
end.
-file("src/multipartkit/stream.gleam", 172).
-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", 99).
?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` and\n"
" `max_part_bytes` are both enforced incrementally as chunks arrive, so an\n"
" oversized stream — or an oversized individual part — is rejected at the\n"
" chunk that crosses the limit, not after the whole part is buffered.\n"
" Per-part memory is bounded by `max_part_bytes`. Body bytes are surfaced\n"
" to consumers in fixed-size chunks via `StreamPart.body` (see the\n"
" `StreamPart` doc).\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", 82).
?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", 383).
?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. The resulting\n"
" body yielder emits the buffered bytes in chunks of up to\n"
" `body_chunk_size`, matching the behaviour of `parse_stream`.\n"
).
-spec from_part(multipartkit@part:part()) -> stream_part().
from_part(The_part) ->
{stream_part,
multipartkit@part:all_headers(The_part),
multipartkit@part:name(The_part),
multipartkit@part:filename(The_part),
multipartkit@part:content_type(The_part),
chunk_body(multipartkit@part:body(The_part), 65536)}.