This is an automated email from the ASF dual-hosted git repository.
nickva pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/couchdb.git
The following commit(s) were added to refs/heads/main by this push:
new 3ebff7dae Fix all stream attachments for serialize_worker_startup
3ebff7dae is described below
commit 3ebff7dae9eea9c9a52d5e908629b6e799c0214a
Author: Nick Vatamaniuc <[email protected]>
AuthorDate: Wed Jul 22 16:24:49 2026 -0400
Fix all stream attachments for serialize_worker_startup
In #6071 we fixed multipart attachment uploads but we forgot to fix plain
attachment uploads with `PUT /db/doc/attachment`. We should apply the same
fix
for these too, otherwise they'll also buffer whole attachment in memory when
sws=true, defeating the whole purpose of incremental attachment streaming.
And add a few parsing tests, especially testing how workers wait on others
and
how the parser should behave if the attachment is too short
---
.../test/eunit/couch_httpd_multipart_tests.erl | 148 +++++++++++++++++++++
src/fabric/src/fabric_doc_update.erl | 51 +++++--
2 files changed, 185 insertions(+), 14 deletions(-)
diff --git a/src/couch/test/eunit/couch_httpd_multipart_tests.erl
b/src/couch/test/eunit/couch_httpd_multipart_tests.erl
new file mode 100644
index 000000000..7411f7611
--- /dev/null
+++ b/src/couch/test/eunit/couch_httpd_multipart_tests.erl
@@ -0,0 +1,148 @@
+% Licensed under the Apache License, Version 2.0 (the "License"); you may not
+% use this file except in compliance with the License. You may obtain a copy of
+% the License at
+%
+% http://www.apache.org/licenses/LICENSE-2.0
+%
+% Unless required by applicable law or agreed to in writing, software
+% distributed under the License is distributed on an "AS IS" BASIS, WITHOUT
+% WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the
+% License for the specific language governing permissions and limitations under
+% the License.
+
+-module(couch_httpd_multipart_tests).
+
+-include_lib("couch/include/couch_eunit.hrl").
+
+-define(CONTENT_TYPE, ~S'multipart/related;boundary="abc123"').
+-define(DOC_JSON,
+
~B'{"_attachments":{"ohai":{"follows":true,"content_type":"text/plain","length":4}}}'
+).
+-define(ATT_LEN, 4).
+
+couch_httpd_multipart_test_() ->
+ {
+ foreach,
+ fun setup/0,
+ fun teardown/1,
+ [
+ ?TDEF_FE(t_single_writer_short_part_fails_fast),
+ ?TDEF_FE(t_single_writer_empty_part_fails_fast),
+ ?TDEF_FE(t_single_writer_exact_part_waits_for_more_writers),
+ ?TDEF_FE(t_all_writers_short_part_fail)
+ ]
+ }.
+
+setup() ->
+ test_util:start_applications([config]).
+
+teardown(Apps) ->
+ test_util:stop_applications(Apps).
+
+% Doc with on attachment declaring length 4 but only having length 3. The
+% writer which detects doesn't have enough bytes exits.
+t_single_writer_short_part_fails_fast(_) ->
+ {Parser, ParserRef, Ref} = parse(3, body(<<"oha">>)),
+ Writer = spawn_writer(Parser, Ref, ?ATT_LEN),
+ ?assertEqual({parser_died, normal}, wait_writer(Writer, 3000)),
+ ?assertEqual(normal, wait_down(ParserRef, 3000)),
+ stop_writers([Writer]).
+
+% Empty body part. Writer won't get any chunk before it starves.
+t_single_writer_empty_part_fails_fast(_) ->
+ {Parser, ParserRef, Ref} = parse(3, body(<<>>)),
+ Writer = spawn_writer(Parser, Ref, ?ATT_LEN),
+ ?assertEqual({parser_died, normal}, wait_writer(Writer, 3000)),
+ ?assertEqual(normal, wait_down(ParserRef, 3000)),
+ stop_writers([Writer]).
+
+% First writer consuming all bytes should not exit and wait for other to also
+% connect and fetch their copies.
+t_single_writer_exact_part_waits_for_more_writers(_) ->
+ {Parser, ParserRef, Ref} = parse(3, body(<<"ohai">>)),
+ Writer1 = spawn_writer(Parser, Ref, ?ATT_LEN),
+ ?assertEqual(got_all_bytes, wait_writer(Writer1, 3000)),
+ % Check it didn't die
+ ?assertEqual(timeout, wait_down(ParserRef, 300)),
+ % Late-arriving writers
+ Writer2 = spawn_writer(Parser, Ref, ?ATT_LEN),
+ ?assertEqual(got_all_bytes, wait_writer(Writer2, 3000)),
+ % Still shouldn't die. One more left
+ ?assertEqual(timeout, wait_down(ParserRef, 300)),
+ Writer3 = spawn_writer(Parser, Ref, ?ATT_LEN),
+ ?assertEqual(got_all_bytes, wait_writer(Writer3, 3000)),
+ % Everyone got their bytes, now it can die
+ ?assertEqual(normal, wait_down(ParserRef, 3000)),
+ stop_writers([Writer1, Writer2, Writer3]).
+
+% Multuple writers get the "too short" condition. That's fine, everyone dies
+% without blocking.
+t_all_writers_short_part_fail(_) ->
+ {Parser, ParserRef, Ref} = parse(3, body(<<"oha">>)),
+ Writers = [spawn_writer(Parser, Ref, ?ATT_LEN) || _ <- lists:seq(1, 3)],
+ [?assertEqual({parser_died, normal}, wait_writer(W, 3000)) || W <-
Writers],
+ ?assertEqual(normal, wait_down(ParserRef, 3000)),
+ stop_writers(Writers).
+
+body(AttBytes) ->
+ <<
+ "--abc123\r\n"
+ "Content-Type: application/json\r\n"
+ "\r\n",
+ ?DOC_JSON/binary,
+ "\r\n"
+ "--abc123\r\n"
+ "\r\n",
+ AttBytes/binary,
+ "\r\n"
+ "--abc123--"
+ >>.
+
+% Spawn the parser like chttpd_db for MP PUT
+parse(NumWriters, Body) ->
+ Ref = make_ref(),
+ couch_httpd_multipart:num_mp_writers(NumWriters),
+ DataFun = fun() -> {Body, fun() -> throw(<<"expected more data">>) end}
end,
+ {{doc_bytes, Ref, DocBytes}, Parser, ParserRef} =
+ couch_httpd_multipart:decode_multipart_stream(?CONTENT_TYPE, DataFun,
Ref),
+ ?assertEqual(?DOC_JSON, iolist_to_binary(DocBytes)),
+ {Parser, ParserRef, Ref}.
+
+% This is like the reader loop in fabric_rpc:make_att_reader()
+spawn_writer(Parser, Ref, Need) ->
+ ReportTo = self(),
+ spawn_link(fun() ->
+ ParserRef = monitor(process, Parser),
+ Parser ! {hello_from_writer, Ref, self()},
+ ReportTo ! {self(), writer_read(Parser, ParserRef, Ref, Need)},
+ receive
+ stop -> ok
+ end
+ end).
+
+writer_read(_Parser, _ParserRef, _Ref, Need) when Need =< 0 ->
+ got_all_bytes;
+writer_read(Parser, ParserRef, Ref, Need) ->
+ Parser ! {get_bytes, Ref, self()},
+ receive
+ {bytes, Ref, Bytes} ->
+ writer_read(Parser, ParserRef, Ref, Need - iolist_size(Bytes));
+ {'DOWN', ParserRef, _, _, Reason} ->
+ {parser_died, Reason}
+ end.
+
+wait_writer(Writer, Timeout) ->
+ receive
+ {Writer, Result} -> Result
+ after Timeout -> timeout
+ end.
+
+wait_down(MonitorRef, Timeout) ->
+ receive
+ {'DOWN', MonitorRef, _, _, Reason} -> Reason
+ after Timeout -> timeout
+ end.
+
+stop_writers(Writers) ->
+ [Writer ! stop || Writer <- Writers],
+ ok.
diff --git a/src/fabric/src/fabric_doc_update.erl
b/src/fabric/src/fabric_doc_update.erl
index f1199c6a7..259bcd40b 100644
--- a/src/fabric/src/fabric_doc_update.erl
+++ b/src/fabric/src/fabric_doc_update.erl
@@ -434,24 +434,26 @@ validate_atomic_update(_DbName, AllDocs, true) ->
),
throw({aborted, PreCommitFailures}).
-% Replicated changes and multipart attachment are always in parallel. MP parser
-% is designed to distribute attachment chunks to num_mp_writers concurrently,
-% so if we serialize them we make the MP parser buffer all the attachment
-% chunks (say 1GB of data) before the subsequent workers will be started.
+% Replicated changes streamed attachments are always in parallel. Both MP
+% parser and fabric_doc_atts are designed to distribute attachment chunks
+% concurrently. If we serialize them they would always buffer the whole
+% attachment in memory (say 1GB of data) until all workers have consumed the
+% byte ranges.
serialize_worker_startup(AllDocs, Options) ->
Replicated = proplists:get_value(?REPLICATED_CHANGES, Options) =:= true,
- case Replicated orelse any_multipart_atts(AllDocs) of
+ case Replicated orelse any_streamed_atts(AllDocs) of
true -> false;
false -> config:get_boolean("fabric", "serialize_worker_startup", true)
end.
-any_multipart_atts(Docs) ->
- DocHasMpAtt = fun(#doc{atts = Atts}) -> lists:any(fun is_multipart_att/1,
Atts) end,
- lists:any(DocHasMpAtt, Docs).
+any_streamed_atts(Docs) ->
+ HasStreamedAtt = fun(#doc{atts = Atts}) -> lists:any(fun
is_streamed_att/1, Atts) end,
+ lists:any(HasStreamedAtt, Docs).
-is_multipart_att(Att) ->
+is_streamed_att(Att) ->
case couch_att:fetch(data, Att) of
{follows, Parser, Ref} when is_pid(Parser), is_reference(Ref) -> true;
+ {fabric_attachment_receiver, Middleman, _} when is_pid(Middleman) ->
true;
_ -> false
end.
@@ -561,7 +563,7 @@ doc_update_test_() ->
fun filter_conflicts_drops_seen_docs/0,
fun parallel_in_flight_after_conflict/0,
fun serial_filters_conflicts_at_cast/0,
- fun sws_multipart_atts_check/0,
+ fun sws_streamed_atts_check/0,
fun sws_false_mode_conflict_not_final/0,
fun sws_false_mode_ok_can_outvote_conflict/0,
fun group_docs_content_and_order/0,
@@ -1174,8 +1176,10 @@ serial_filters_conflicts_at_cast() ->
?assertEqual([Tagged2], Stored).
% Docs with streaming attachment don't serialize workers
-sws_multipart_atts_check() ->
+sws_streamed_atts_check() ->
MpDoc = #doc{id = <<"m">>, atts = [mp_att()]},
+ ReceiverDoc = #doc{id = <<"r">>, atts = [receiver_att()]},
+ ChunkedReceiverDoc = #doc{id = <<"cr">>, atts = [chunked_receiver_att()]},
InlineAtt = couch_att:new([
{name, <<"b">>}, {type, <<"text/plain">>}, {att_len, 1}, {data,
<<"x">>}
]),
@@ -1184,12 +1188,15 @@ sws_multipart_atts_check() ->
]),
InlineDoc = #doc{id = <<"i">>, atts = [InlineAtt, StubAtt]},
PlainDoc = #doc{id = <<"p">>},
- % Without multipart atts the config default applies
+ % Without streamed atts the config default applies
?assert(serialize_worker_startup([PlainDoc], [])),
?assert(serialize_worker_startup([InlineDoc], [])),
- % Any doc with a mp att forces parallel startup
+ % Parallel for mp parser and fabric attachment receiver
?assertNot(serialize_worker_startup([MpDoc], [])),
- ?assertNot(serialize_worker_startup([PlainDoc, MpDoc], [])).
+ ?assertNot(serialize_worker_startup([PlainDoc, MpDoc], [])),
+ ?assertNot(serialize_worker_startup([ReceiverDoc], [])),
+ ?assertNot(serialize_worker_startup([PlainDoc, ReceiverDoc], [])),
+ ?assertNot(serialize_worker_startup([ChunkedReceiverDoc], [])).
mp_att() ->
couch_att:new([
@@ -1199,6 +1206,22 @@ mp_att() ->
{data, {follows, self(), make_ref()}}
]).
+receiver_att() ->
+ couch_att:new([
+ {name, <<"a">>},
+ {type, <<"text/plain">>},
+ {att_len, 4},
+ {data, {fabric_attachment_receiver, self(), 4}}
+ ]).
+
+chunked_receiver_att() ->
+ couch_att:new([
+ {name, <<"a">>},
+ {type, <<"text/plain">>},
+ {att_len, undefined},
+ {data, {fabric_attachment_receiver, self(), chunked}}
+ ]).
+
sws_false_mode_conflict_not_final() ->
Docs =
[Doc1, Doc2] = tag_docs([