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 a8b399fb0 Always apply user configured compression
a8b399fb0 is described below

commit a8b399fb0ec4c5591919321e38f7c693e503dc33
Author: Nick Vatamaniuc <[email protected]>
AuthorDate: Sun Jun 7 00:46:02 2026 -0400

    Always apply user configured compression
    
    Previously, even when user picked a non-default compression algorithm in 
some
    instances, like the purge trees and emsort, snappy would still be used such.
    Make sure to always respect user's choice, and if they picked something
    non-default always use that.
    
    [1] There is a good reason to avoid snappy as doesn't properly yield when
    encoding or decoding. In a low concurrency environment it might not matter, 
but
    as the concurrency ramps up it could lead to scheduler collapse (it could
    manifest as high p90+ latencies while cpus usage remains low). A similar 
effect
    is illustrated in scheduler responsiveness test for jiffy vs other C json
    libraries in https://github.com/nickva/bench#scheduler-responsiveness-test
---
 src/couch/src/couch_bt_engine.erl                   | 12 ++++++++----
 src/couch/src/couch_bt_engine_compactor.erl         | 18 +++++++++++-------
 src/couch/src/couch_emsort.erl                      | 21 ++++++++++++++++-----
 .../test/eunit/couch_bt_engine_compactor_tests.erl  | 10 +++++++---
 src/smoosh/test/smoosh_tests.erl                    | 17 ++++++++++++-----
 5 files changed, 54 insertions(+), 24 deletions(-)

diff --git a/src/couch/src/couch_bt_engine.erl 
b/src/couch/src/couch_bt_engine.erl
index 072a75c74..4bdef732d 100644
--- a/src/couch/src/couch_bt_engine.erl
+++ b/src/couch/src/couch_bt_engine.erl
@@ -844,7 +844,7 @@ init_state(FilePath, Fd, Header0, Options) ->
 
     Header1 = couch_bt_engine_header:upgrade(Header0),
     Header2 = set_default_security_object(Fd, Header1, Compression, Options),
-    Header = upgrade_purge_info(Fd, Header2),
+    Header = upgrade_purge_info(Fd, Header2, Compression),
 
     IdTreeState = couch_bt_engine_header:id_tree_state(Header),
     {ok, IdTree} = couch_btree:open(IdTreeState, Fd, [
@@ -877,6 +877,7 @@ init_state(FilePath, Fd, Header0, Options) ->
         {split, fun ?MODULE:purge_tree_split/1},
         {join, fun ?MODULE:purge_tree_join/2},
         {reduce, fun ?MODULE:purge_tree_reduce/2},
+        {compression, Compression},
         {cache_depth, btree_cache_depth()}
     ]),
 
@@ -885,6 +886,7 @@ init_state(FilePath, Fd, Header0, Options) ->
         {split, fun ?MODULE:purge_seq_tree_split/1},
         {join, fun ?MODULE:purge_seq_tree_join/2},
         {reduce, fun ?MODULE:purge_tree_reduce/2},
+        {compression, Compression},
         {cache_depth, btree_cache_depth()}
     ]),
 
@@ -943,7 +945,7 @@ set_default_security_object(Fd, Header, Compression, 
Options) ->
 
 % This function is here, and not in couch_bt_engine_header
 % because it requires modifying file contents
-upgrade_purge_info(Fd, Header) ->
+upgrade_purge_info(Fd, Header, Compression) ->
     case couch_bt_engine_header:get(Header, purge_tree_state) of
         nil ->
             Header;
@@ -979,7 +981,8 @@ upgrade_purge_info(Fd, Header) ->
                     {ok, PurgeTree} = couch_btree:open(nil, Fd, [
                         {split, fun ?MODULE:purge_tree_split/1},
                         {join, fun ?MODULE:purge_tree_join/2},
-                        {reduce, fun ?MODULE:purge_tree_reduce/2}
+                        {reduce, fun ?MODULE:purge_tree_reduce/2},
+                        {compression, Compression}
                     ]),
                     {ok, PurgeTree2} = couch_btree:add(PurgeTree, Infos),
                     PurgeTreeSt = couch_btree:get_state(PurgeTree2),
@@ -987,7 +990,8 @@ upgrade_purge_info(Fd, Header) ->
                     {ok, PurgeSeqTree} = couch_btree:open(nil, Fd, [
                         {split, fun ?MODULE:purge_seq_tree_split/1},
                         {join, fun ?MODULE:purge_seq_tree_join/2},
-                        {reduce, fun ?MODULE:purge_tree_reduce/2}
+                        {reduce, fun ?MODULE:purge_tree_reduce/2},
+                        {compression, Compression}
                     ]),
                     {ok, PurgeSeqTree2} = couch_btree:add(PurgeSeqTree, Infos),
                     PurgeSeqTreeSt = couch_btree:get_state(PurgeSeqTree2),
diff --git a/src/couch/src/couch_bt_engine_compactor.erl 
b/src/couch/src/couch_bt_engine_compactor.erl
index 12c1ea3ac..ee349337d 100644
--- a/src/couch/src/couch_bt_engine_compactor.erl
+++ b/src/couch/src/couch_bt_engine_compactor.erl
@@ -509,7 +509,9 @@ copy_docs(St, #st{} = NewSt, MixedInfos, Retry) ->
     ),
 
     EMSortFd = couch_emsort:get_fd(NewSt#st.id_tree),
-    {ok, LocSizes} = couch_file:append_terms(EMSortFd, NewInfos),
+    {ok, LocSizes} = couch_file:append_terms(
+        EMSortFd, NewInfos, [{compression, NewSt#st.compression}]
+    ),
     EMSortEntries = lists:zipwith(
         fun(FDI, {Loc, _}) ->
             #full_doc_info{
@@ -603,7 +605,8 @@ copy_meta_data(#comp_st{new_st = St} = CompSt) ->
     {ok, IdTree0} = couch_btree:open(DstState, Fd, [
         {split, fun couch_bt_engine:id_tree_split/1},
         {join, fun couch_bt_engine:id_tree_join/2},
-        {reduce, fun couch_bt_engine:id_tree_reduce/2}
+        {reduce, fun couch_bt_engine:id_tree_reduce/2},
+        {compression, St#st.compression}
     ]),
     {ok, Iter} = couch_emsort:iter(Src),
     Acc0 = #merge_st{
@@ -689,20 +692,21 @@ commit_compaction_data(#st{header = OldHeader} = St0, Fd) 
->
     bind_emsort(St2, MetaFd, MetaState).
 
 bind_emsort(St, Fd, nil) ->
-    {ok, Ems} = couch_emsort:open(Fd),
-    St#st{id_tree = Ems};
+    bind_emsort(St, Fd, []);
 bind_emsort(St, Fd, {BB, _} = Root) when is_list(BB) ->
     % Upgrade clause when we find old compaction files
     bind_emsort(St, Fd, [{root, Root}]);
-bind_emsort(St, Fd, State) ->
-    {ok, Ems} = couch_emsort:open(Fd, State),
+bind_emsort(St, Fd, State) when is_list(State) ->
+    Options = [{compression, St#st.compression} | State],
+    {ok, Ems} = couch_emsort:open(Fd, Options),
     St#st{id_tree = Ems}.
 
 bind_id_tree(St, Fd, State) ->
     {ok, IdBtree} = couch_btree:open(State, Fd, [
         {split, fun couch_bt_engine:id_tree_split/1},
         {join, fun couch_bt_engine:id_tree_join/2},
-        {reduce, fun couch_bt_engine:id_tree_reduce/2}
+        {reduce, fun couch_bt_engine:id_tree_reduce/2},
+        {compression, St#st.compression}
     ]),
     St#st{id_tree = IdBtree}.
 
diff --git a/src/couch/src/couch_emsort.erl b/src/couch/src/couch_emsort.erl
index 4250b45e9..6a7856839 100644
--- a/src/couch/src/couch_emsort.erl
+++ b/src/couch/src/couch_emsort.erl
@@ -133,13 +133,16 @@
 -export([add/2, merge/2, iter/1, next/1]).
 -export([num_kvs/1, num_merges/1]).
 
+-include_lib("couch/include/couch_db.hrl").
+
 -record(ems, {
     fd,
     root,
     bb_chunk = 10,
     chain_chunk = 100,
     num_kvs = 0,
-    num_bb = 0
+    num_bb = 0,
+    compression = ?DEFAULT_COMPRESSION
 }).
 
 -define(REPORT_INTERVAL, 1000).
@@ -161,7 +164,9 @@ set_options(Ems, [{back_bone_chunk, Count} | Rest]) when 
is_integer(Count) ->
 set_options(Ems, [{num_kvs, NumKVs} | Rest]) when is_integer(NumKVs) ->
     set_options(Ems#ems{num_kvs = NumKVs}, Rest);
 set_options(Ems, [{num_bb, NumBB} | Rest]) when is_integer(NumBB) ->
-    set_options(Ems#ems{num_bb = NumBB}, Rest).
+    set_options(Ems#ems{num_bb = NumBB}, Rest);
+set_options(Ems, [{compression, Comp} | Rest]) ->
+    set_options(Ems#ems{compression = Comp}, Rest).
 
 get_fd(#ems{fd = Fd}) ->
     Fd.
@@ -231,7 +236,9 @@ write_kvs(Ems, KVs) ->
             {[], nil},
             lists:sort(KVs)
         ),
-    {ok, Final, _} = couch_file:append_term(Ems#ems.fd, {LastKVs, LastPos}),
+    {ok, Final, _} = couch_file:append_term(
+        Ems#ems.fd, {LastKVs, LastPos}, [{compression, Ems#ems.compression}]
+    ),
     Final.
 
 decimate(#ems{root = {_BB, nil}} = Ems, _Reporter) ->
@@ -274,7 +281,9 @@ merge_chains(Ems, Choose, BB, Reporter) ->
     merge_chains(Ems, Choose, Chains, {[], nil}, Reporter, 0).
 
 merge_chains(Ems, _Choose, [], ChainAcc, _Reporter, _Count) ->
-    {ok, CPos, _} = couch_file:append_term(Ems#ems.fd, ChainAcc),
+    {ok, CPos, _} = couch_file:append_term(
+        Ems#ems.fd, ChainAcc, [{compression, Ems#ems.compression}]
+    ),
     CPos;
 merge_chains(#ems{chain_chunk = CC} = Ems, Choose, Chains, Acc, Reporter, 
Count0) ->
     {KV, RestChains} = choose_kv(Choose, Ems, Chains),
@@ -327,7 +336,9 @@ ins_big_chain(Rest, Chain, Acc) ->
     lists:reverse(Acc, [Chain | Rest]).
 
 append_item(Ems, {List, Prev}, Pos, Size) when length(List) >= Size ->
-    {ok, PrevList, _} = couch_file:append_term(Ems#ems.fd, {List, Prev}),
+    {ok, PrevList, _} = couch_file:append_term(
+        Ems#ems.fd, {List, Prev}, [{compression, Ems#ems.compression}]
+    ),
     {[Pos], PrevList};
 append_item(_Ems, {List, Prev}, Pos, _Size) ->
     {[Pos | List], Prev}.
diff --git a/src/couch/test/eunit/couch_bt_engine_compactor_tests.erl 
b/src/couch/test/eunit/couch_bt_engine_compactor_tests.erl
index a7fec17db..63df30fb0 100644
--- a/src/couch/test/eunit/couch_bt_engine_compactor_tests.erl
+++ b/src/couch/test/eunit/couch_bt_engine_compactor_tests.erl
@@ -70,7 +70,7 @@ is_compacting_works(DbName) ->
         {_Pid, Ref} = spawn_monitor(fun() -> compact_db(DbName) end),
         receive
             {in_emsort_bind, From} ->
-                % When emsort:open(Fd) is called the files should
+                % When emsort:open(Fd, Options) is called the files should
                 % have been created already
                 ?assert(couch_db:is_compacting(DbName)),
                 From ! {please_continue, self()}
@@ -145,11 +145,15 @@ wait_db_compact_done(DbName, N) ->
     end.
 
 wait_in_emsort_bind(WaitingPid) when is_pid(WaitingPid) ->
-    meck:expect(couch_emsort, open, fun(Fd) ->
+    meck:expect(couch_emsort, open, fun(Fd, Options) ->
+        % Intercept first call to bind only the rest should pass through
+        meck:expect(couch_emsort, open, fun(F, O) ->
+            meck:passthrough([F, O])
+        end),
         WaitingPid ! {in_emsort_bind, self()},
         receive
             {please_continue, WaitingPid} ->
                 ok
         end,
-        meck:passthrough([Fd])
+        meck:passthrough([Fd, Options])
     end).
diff --git a/src/smoosh/test/smoosh_tests.erl b/src/smoosh/test/smoosh_tests.erl
index 517024875..7f540af64 100644
--- a/src/smoosh/test/smoosh_tests.erl
+++ b/src/smoosh/test/smoosh_tests.erl
@@ -473,11 +473,18 @@ wait_update_status() ->
 
 setup_db_compactor_intercept() ->
     TestPid = self(),
-    meck:expect(couch_emsort, open, fun(Fd) ->
-        TestPid ! {compactor_paused, self()},
-        receive
-            continue -> meck:passthrough([Fd]);
-            {raise, Tag, Reason} -> meck:exception(Tag, Reason)
+    meck:expect(couch_emsort, open, fun(Fd, Options) ->
+        % in the intercept is called a few times on first bind it has not
+        % persistent root so pause there and let the bind pass through.
+        case lists:keymember(root, 1, Options) of
+            false ->
+                TestPid ! {compactor_paused, self()},
+                receive
+                    continue -> meck:passthrough([Fd, Options]);
+                    {raise, Tag, Reason} -> meck:exception(Tag, Reason)
+                end;
+            true ->
+                meck:passthrough([Fd, Options])
         end
     end).
 

Reply via email to