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 84aa91700 Fix multipart attachments with serialize_worker_startup
84aa91700 is described below
commit 84aa91700c011a82e242b18e99e875d80cb57331
Author: Nick Vatamaniuc <[email protected]>
AuthorDate: Wed Jul 22 01:11:24 2026 -0400
Fix multipart attachments with serialize_worker_startup
Previously, if a multipart parser reached an eof before the declared
length, it
still let the writer wait and timeout (300s). This was especially noticed if
users toggled `serialize_worker_startup=true`. Then, the first worker would
buffer all the data (another undesirable behavior) until it reached a
premature
eof and then get stuck. Other workers wouldn't start either and the request
would eventually crash.
Here we fix both issues:
1) If we detect serialize_worker_startup=true then we switch back to the
default parallel worker startup pattern. This how the MP parser was built
to
work. Otherwise it would buffer the whole attachment into memory until the
first worker wasn't done and the others started. That defies the purpose of
incremental attachment uploads. So to go with the grain of MP parser
design we
switch back to parallel worker startup.
2) Let writers which wait on an EOF exit instead of deadlocking until a
timeout. Writer waits for more bytes when the stream already reached EOF is
because the user uploaded less than the declared number of bytes
(attachment
is too short). In that case we exit normally and let the parser monitors
fire
with `"attachment shorter than expected"` error.
Reproducer for the issue:
* start dev/run cluster with 3 nodes
* `s:multicall(fun() -> config:set("fabric",
"serialize_worker_startup","true") end)`
* curl -XPUT 'http://adm:[email protected]:15984/mptest'
Before PR:
```
curl --max-time 10 -XPUT 'http://adm:[email protected]:15984/mptest/short' -H
'Content-Type: multipart/related;boundary="abc123"' --data-binary
$'--abc123\r\nContent-Type:
application/json\r\n\r\n{"_attachments":{"ohai":{"follows":true,"content_type":"text/plain","length":
4}}}\r\n--abc123\r\n\r\noha\r\n--abc123--'
curl: (28) Operation timed out after 10005 milliseconds with 0 bytes
received
```
After PR
```
curl --max-time 10 -X PUT 'http://adm:[email protected]:15984/mptest/short'
-H 'Content-Type: multipart/related;boundary="abc123"' --data-binary
$'--abc123\r\nContent-Type:
application/json\r\n\r\n{"_attachments":{"ohai":{"follows":true,"content_type":"text/plain","length":
4}}}\r\n--abc123\r\n\r\noha\r\n--abc123--'
{"error":"bad_request","reason":"attachment shorter than expected"}
```
---
src/couch/src/couch_httpd_multipart.erl | 7 ++++-
src/fabric/src/fabric_doc_update.erl | 54 +++++++++++++++++++++++++++++----
2 files changed, 54 insertions(+), 7 deletions(-)
diff --git a/src/couch/src/couch_httpd_multipart.erl
b/src/couch/src/couch_httpd_multipart.erl
index 593c4c315..67306bcef 100644
--- a/src/couch/src/couch_httpd_multipart.erl
+++ b/src/couch/src/couch_httpd_multipart.erl
@@ -94,7 +94,12 @@ mp_parse_atts({body, Bytes}, {Ref, Chunks, Offset, Counters,
Waiting}) ->
mp_parse_atts(eof, {Ref, Chunks, Offset, Counters, Waiting}) ->
N = num_mp_writers(),
M = length(Counters),
- case (M == N) andalso Chunks == [] of
+ AllConsumed = (M == N) andalso Chunks == [],
+ % A writer waiting at eof is asking of bytes past the end of the stream, it
+ % means the stream ended before its declared length. A write then can never
+ % be satified so we should stop parsing. If we exit normal the waiting
+ % writers will fail as with "attachment shorter than expected".
+ case AllConsumed orelse Waiting =/= [] of
true ->
ok;
false ->
diff --git a/src/fabric/src/fabric_doc_update.erl
b/src/fabric/src/fabric_doc_update.erl
index 70081af73..f1199c6a7 100644
--- a/src/fabric/src/fabric_doc_update.erl
+++ b/src/fabric/src/fabric_doc_update.erl
@@ -53,7 +53,7 @@ go(DbName, AllDocs0, Opts) ->
reply = #{},
dbname = DbName,
update_options = Options,
- serialize_worker_startup = serialize_worker_startup(Options)
+ serialize_worker_startup = serialize_worker_startup(AllDocs, Options)
},
Timeout = fabric_util:request_timeout(),
Acc1 = start_workers_strategy(Acc0),
@@ -434,11 +434,25 @@ validate_atomic_update(_DbName, AllDocs, true) ->
),
throw({aborted, PreCommitFailures}).
-% replicated changes are always in parallel
-serialize_worker_startup(Options) ->
- case proplists:get_value(?REPLICATED_CHANGES, Options) of
+% 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.
+serialize_worker_startup(AllDocs, Options) ->
+ Replicated = proplists:get_value(?REPLICATED_CHANGES, Options) =:= true,
+ case Replicated orelse any_multipart_atts(AllDocs) of
true -> false;
- _ -> config:get_boolean("fabric", "serialize_worker_startup", true)
+ 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).
+
+is_multipart_att(Att) ->
+ case couch_att:fetch(data, Att) of
+ {follows, Parser, Ref} when is_pid(Parser), is_reference(Ref) -> true;
+ _ -> false
end.
start_workers_strategy(#acc{serialize_worker_startup = true} = Acc) ->
@@ -513,9 +527,10 @@ filter_conflicts(Docs, Conflicts) ->
-include_lib("eunit/include/eunit.hrl").
setup_all() ->
- meck:new([couch_log, couch_stats]),
+ meck:new([couch_log, couch_stats, config]),
meck:expect(couch_log, warning, fun(_, _) -> ok end),
meck:expect(couch_stats, increment_counter, fun(_) -> ok end),
+ meck:expect(config, get_boolean, fun(_, _, Default) -> Default end),
meck:new(rexi, [passthrough]),
meck:expect(rexi, cast_ref, fun(Ref, _Node, _Msg) -> Ref end).
@@ -546,6 +561,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_false_mode_conflict_not_final/0,
fun sws_false_mode_ok_can_outvote_conflict/0,
fun group_docs_content_and_order/0,
@@ -1157,6 +1173,32 @@ serial_filters_conflicts_at_cast() ->
{Worker, Stored} = lists:keyfind(Worker, 1, Acc1#acc.grouped_docs),
?assertEqual([Tagged2], Stored).
+% Docs with streaming attachment don't serialize workers
+sws_multipart_atts_check() ->
+ MpDoc = #doc{id = <<"m">>, atts = [mp_att()]},
+ InlineAtt = couch_att:new([
+ {name, <<"b">>}, {type, <<"text/plain">>}, {att_len, 1}, {data,
<<"x">>}
+ ]),
+ StubAtt = couch_att:new([
+ {name, <<"c">>}, {type, <<"text/plain">>}, {att_len, 1}, {data, stub}
+ ]),
+ InlineDoc = #doc{id = <<"i">>, atts = [InlineAtt, StubAtt]},
+ PlainDoc = #doc{id = <<"p">>},
+ % Without multipart atts the config default applies
+ ?assert(serialize_worker_startup([PlainDoc], [])),
+ ?assert(serialize_worker_startup([InlineDoc], [])),
+ % Any doc with a mp att forces parallel startup
+ ?assertNot(serialize_worker_startup([MpDoc], [])),
+ ?assertNot(serialize_worker_startup([PlainDoc, MpDoc], [])).
+
+mp_att() ->
+ couch_att:new([
+ {name, <<"a">>},
+ {type, <<"text/plain">>},
+ {att_len, 4},
+ {data, {follows, self(), make_ref()}}
+ ]).
+
sws_false_mode_conflict_not_final() ->
Docs =
[Doc1, Doc2] = tag_docs([