This is an automated email from the ASF dual-hosted git repository. nickva pushed a commit to branch fix-conflicts-sws in repository https://gitbox.apache.org/repos/asf/couchdb.git
commit 620c1f02c6c6cd12c5cb14bed7bb39afa303fd57 Author: Nick Vatamaniuc <[email protected]> AuthorDate: Tue Aug 11 18:32:05 2026 -0400 Fix secondary worker conflicts case for sws=true In sws=true case, secondary replies (from non-owner) should generate accepted responses instead of conflicts. With sws=true the owner already replied with `ok` means there at least one good update of that document. Secondaries returning conflicts means they are simply lagging. That's the same calculation logic as previously (sws=false) when we got one ok and two errors. Otherwise we get into a situation when clients cannot make progress: they get a rev, do the update (on the owner) but then keep getting conflicts until the internal replicator catches up. Retrying will fix the problem eventually but we don't want to emit too many 409 and that's not optimal. --- src/fabric/src/fabric_doc_update.erl | 92 ++++++++++++++++++++++++++++++------ 1 file changed, 77 insertions(+), 15 deletions(-) diff --git a/src/fabric/src/fabric_doc_update.erl b/src/fabric/src/fabric_doc_update.erl index 259bcd40b..e37e9fc8b 100644 --- a/src/fabric/src/fabric_doc_update.erl +++ b/src/fabric/src/fabric_doc_update.erl @@ -116,7 +116,7 @@ handle_message({ok, Replies}, Worker, #acc{} = Acc0) -> serialize_worker_startup = SWS } = Acc0, {value, {_, Docs}, NewGrpDocs} = lists:keytake(Worker, 1, GroupedDocs), - Conflicts = collect_conflicts(Docs, Replies, Conflicts0, SWS), + Conflicts = collect_conflicts(Docs, Replies, DocReplyDict0, Conflicts0, SWS), DocReplyDict = append_update_replies(Docs, Replies, DocReplyDict0), Acc1 = Acc0#acc{conflicts = Conflicts}, case {WaitingCount, map_size(DocReplyDict)} of @@ -321,23 +321,34 @@ check_forbidden_msg(Replies) -> end. update_quorum_met(W, Replies, SWS) -> + % With sws=true conflict after the owner returns ok should not turn into a conflict + % it should be become an accepted (202) instead. But if we don't have any ok responses, + % say owner returned an error then we should return a conflict from a secondary node. + ConflictsFinal = SWS andalso not any_ok(Replies), Counters = lists:foldl( fun(R, D) -> orddict:update_counter(R, 1, D) end, orddict:new(), Replies ), - GoodReplies = lists:filter(fun(C) -> good_reply(C, SWS) end, Counters), - case lists:dropwhile(quorum_pred(W, SWS), GoodReplies) of + GoodReplies = lists:filter(fun(C) -> good_reply(C, ConflictsFinal) end, Counters), + case lists:dropwhile(quorum_pred(W, ConflictsFinal), GoodReplies) of [] -> false; [{FinalReply, _} | _] -> {true, FinalReply} end. -% With a conflict in sws we stop the quorum early -quorum_pred(W, SWS) when is_boolean(SWS) -> +any_ok(Replies) -> + Fun = fun + ({ok, _}) -> true; + (_) -> false + end, + lists:any(Fun, Replies). + +% Final confluct => stop the quorum early +quorum_pred(W, ConflictsFinal) when is_boolean(ConflictsFinal) -> fun - ({conflict, _}) when SWS -> false; + ({conflict, _}) when ConflictsFinal -> false; ({_, Count}) -> Count < W end. @@ -345,8 +356,8 @@ good_reply({{ok, _}, _}, _) -> true; good_reply({noreply, _}, _) -> true; -good_reply({conflict, _}, SWS) -> - SWS; +good_reply({conflict, _}, ConflictsFinal) -> + ConflictsFinal; good_reply(_, _) -> false. @@ -511,14 +522,23 @@ start_worker(#shard{}, _Docs, #acc{} = Acc) -> % With sws=false conflicts act like normal values and conflicts list stays % empty (filter_conflicts/2 is a no-op) -collect_conflicts(_Docs, _Replies, Conflicts, false) -> +collect_conflicts(_Docs, _Replies, _DocReplyDict, Conflicts, false) -> Conflicts; -collect_conflicts(Docs, Replies, Conflicts, true) when length(Docs) == length(Replies) -> - [untag_doc(D) || {D, conflict} <- lists:zip(Docs, Replies)] ++ Conflicts; -collect_conflicts(_Docs, _Replies, Conflicts, true) -> +collect_conflicts(Docs, Replies, DocReplyDict, Conflicts, true) when + length(Docs) == length(Replies) +-> + % In case of a conflict from a lagging copy don't throw a final conflict, + % consider it an "accepted" case. The copy hasn't caught up yet. + Zipped = lists:zip(Docs, Replies), + [untag_doc(D) || {D, conflict} <- Zipped, not doc_ok_reply(D, DocReplyDict)] ++ Conflicts; +collect_conflicts(_Docs, _Replies, _DocReplyDict, Conflicts, true) -> % Replicated changes return no replies by default Conflicts. +doc_ok_reply(#doc{} = Doc, DocReplyDict) -> + {_, Replies} = maps:get(doc_tag(Doc), DocReplyDict, {Doc, []}), + any_ok(Replies). + filter_conflicts(Docs, []) -> Docs; filter_conflicts(Docs, Conflicts) -> @@ -566,6 +586,8 @@ doc_update_test_() -> fun sws_streamed_atts_check/0, fun sws_false_mode_conflict_not_final/0, fun sws_false_mode_ok_can_outvote_conflict/0, + fun sws_secondary_conflict_is_accepted/0, + fun sws_secondary_conflict_and_ok_is_quorum/0, fun group_docs_content_and_order/0, fun rotate_ranges_rotates_each_range/0, fun rotate_ranges_varies_by_db/0, @@ -1062,7 +1084,7 @@ one_success_two_forbid() -> ). worker_before_doc_update_forbidden() -> - Docs = [Doc1] = tag_docs([#doc{revs = {1, [<<"foo">>]}}]), + Docs = [_] = tag_docs([#doc{revs = {1, [<<"foo">>]}}]), Shards = mem3_util:create_partition_map("foo", 3, 1, ["node1", "node2", "node3"]), GroupedDocs = group_docs_by_shard_hack(<<"foo">>, Shards, Docs), @@ -1076,7 +1098,7 @@ worker_before_doc_update_forbidden() -> ?assertThrow({forbidden, <<"msg">>}, handle_message({forbidden, <<"msg">>}, hd(Shards), Acc)). handle_bad_request() -> - Docs = [Doc1] = tag_docs([#doc{revs = {1, [<<"foo">>]}}]), + Docs = [_] = tag_docs([#doc{revs = {1, [<<"foo">>]}}]), Shards = mem3_util:create_partition_map("foo", 3, 1, ["node1", "node2", "node3"]), GroupedDocs = group_docs_by_shard_hack(<<"foo">>, Shards, Docs), @@ -1224,7 +1246,7 @@ chunked_receiver_att() -> sws_false_mode_conflict_not_final() -> Docs = - [Doc1, Doc2] = tag_docs([ + [_Doc1, Doc2] = tag_docs([ #doc{revs = {1, [<<"foo">>]}}, #doc{revs = {1, [<<"bar">>]}} ]), @@ -1279,6 +1301,46 @@ sws_false_mode_ok_can_outvote_conflict() -> lists:sort(Reply) ). +% When sws=true if owner returns ok and secondaries (possibly stale) +% return conflict => return accepted +sws_secondary_conflict_is_accepted() -> + Docs = [Doc1] = tag_docs([#doc{revs = {1, [<<"foo">>]}}]), + [S1, S2, S3] = + Shards = + mem3_util:create_partition_map("foo", 3, 1, ["node1", "node2", "node3"]), + GroupedDocs = group_docs_by_shard_hack(<<"foo">>, Shards, Docs), + Acc0 = #acc{ + waiting_count = length(Shards), + doc_count = length(Docs), + w = 2, + grouped_docs = GroupedDocs, + reply = reply_map(Docs) + }, + {ok, #acc{} = Acc1} = handle_message({ok, [{ok, Doc1}]}, S1, Acc0), + {ok, #acc{conflicts = []} = Acc2} = handle_message({ok, [conflict]}, S2, Acc1), + {stop, Res} = handle_message({ok, [conflict]}, S3, Acc2), + ?assertEqual({accepted, [{Doc1, {accepted, Doc1}}]}, Res). + +% When sws=true and get owner=ok, conflict from secondary but then another ok +% from the 3rd node. Since we got 2 OKs we can reply ok (201) for w=2 quorum +sws_secondary_conflict_and_ok_is_quorum() -> + Docs = [Doc1] = tag_docs([#doc{revs = {1, [<<"foo">>]}}]), + [S1, S2, S3] = + Shards = + mem3_util:create_partition_map("foo", 3, 1, ["node1", "node2", "node3"]), + GroupedDocs = group_docs_by_shard_hack(<<"foo">>, Shards, Docs), + Acc0 = #acc{ + waiting_count = length(Shards), + doc_count = length(Docs), + w = 2, + grouped_docs = GroupedDocs, + reply = reply_map(Docs) + }, + {ok, Acc1} = handle_message({ok, [{ok, Doc1}]}, S1, Acc0), + {ok, #acc{conflicts = []} = Acc2} = handle_message({ok, [conflict]}, S2, Acc1), + {stop, Res} = handle_message({ok, [{ok, Doc1}]}, S3, Acc2), + ?assertEqual({ok, [{Doc1, {ok, Doc1}}]}, Res). + owner_order_test() -> S1 = mk_shard(<<"r1">>, n1, [0, 10]), S2 = mk_shard(<<"r1">>, n2, [0, 10]),
