4141% where we may end up processing an unbounded number of messages.
4242-define (MAX_DISCARDED_MESSAGES , 100 ).
4343
44- should_compress_request (Body ) when is_binary (Body ) ->
45- MinSize = config :get_integer (" replicator" , " compress_min_size" , 1024 ),
46- byte_size (Body ) >= MinSize ;
47- should_compress_request (Body ) when is_list (Body ) ->
48- should_compress_request (iolist_to_binary (Body ));
49- should_compress_request (_ ) ->
50- false .
51-
52- get_compression_algorithm () ->
53- % Supported: gzip (default), deflate
54- Algorithm = config :get (" replicator" , " compression_algorithm" , " gzip" ),
55- case Algorithm of
56- " gzip" -> gzip ;
57- " deflate" -> deflate ;
58- _ ->
59- couch_log :warning (
60- " couch_replicator_httpc: Unknown compression algorithm ~p , using gzip" ,
61- [Algorithm ]
62- ),
63- gzip
64- end .
65-
66- compress_body (Body ) when is_binary (Body ) ->
67- compress_body_with_algorithm (Body , get_compression_algorithm ());
68- compress_body (Body ) when is_list (Body ) ->
69- compress_body (iolist_to_binary (Body ));
70- compress_body (Body ) ->
71- Body .
72-
73- compress_body_with_algorithm (Body , gzip ) ->
74- zlib :gzip (Body );
75- compress_body_with_algorithm (Body , deflate ) ->
76- zlib :compress (Body ).
77-
78- get_content_encoding (gzip ) -> " gzip" ;
79- get_content_encoding (deflate ) -> " deflate" .
80-
81- decompress_body (Headers , Body ) ->
82- case lists :keyfind (" Content-Encoding" , 1 , Headers ) of
83- {" Content-Encoding" , Encoding } ->
84- decompress_body_with_encoding (Encoding , Body );
85- _ ->
86- Body
87- end .
88-
89- decompress_body_with_encoding (" gzip" , Body ) ->
90- try
91- zlib :gunzip (Body )
92- catch
93- error :data_error ->
94- couch_log :warning (
95- " couch_replicator_httpc: Failed to decompress gzip response, using original" ,
96- []
97- ),
98- Body
99- end ;
100- decompress_body_with_encoding (" deflate" , Body ) ->
101- try
102- zlib :uncompress (Body )
103- catch
104- error :data_error ->
105- couch_log :warning (
106- " couch_replicator_httpc: Failed to decompress deflate response, using original" ,
107- []
108- ),
109- Body
110- end ;
111- decompress_body_with_encoding (Other , Body ) ->
112- couch_log :warning (
113- " couch_replicator_httpc: Unknown content encoding ~p , using original body" ,
114- [Other ]
115- ),
116- Body .
117-
11844setup (Db ) ->
11945 # httpdb {
12046 httpc_pool = nil ,
@@ -183,36 +109,11 @@ stop_http_worker() ->
183109
184110send_ibrowse_req (# httpdb {headers = BaseHeaders } = HttpDb0 , Params ) ->
185111 Method = get_value (method , Params , get ),
186- UserHeaders0 = get_value (headers , Params , []),
187- % Accept multiple compression algorithms
188- AcceptEncodings = config :get (" replicator" , " accept_encodings" , " gzip, deflate, zstd" ),
189- UserHeaders1 = case lists :keyfind (" Accept-Encoding" , 1 , UserHeaders0 ) of
190- false -> [{" Accept-Encoding" , AcceptEncodings } | UserHeaders0 ];
191- _ -> UserHeaders0
192- end ,
193- Body0 = get_value (body , Params , []),
194- CompressEnabled = config :get_boolean (" replicator" , " compress_requests" , true ),
195- ShouldCompress = should_compress_request (Body0 ),
196- {Body , UserHeaders2 } = case CompressEnabled andalso ShouldCompress of
197- true ->
198- Algorithm = get_compression_algorithm (),
199- CompressedBody = compress_body (Body0 ),
200- ContentEncoding = get_content_encoding (Algorithm ),
201- UpdatedHeaders = case lists :keyfind (" Content-Encoding" , 1 , UserHeaders1 ) of
202- false -> [{" Content-Encoding" , ContentEncoding } | UserHeaders1 ];
203- _ -> UserHeaders1
204- end ,
205- % Track compression algorithm usage
206- couch_stats :increment_counter ([couch_replicator , requests_compressed ]),
207- couch_stats :increment_counter ([couch_replicator , requests_compressed , Algorithm ]),
208- {CompressedBody , UpdatedHeaders };
209- false ->
210- {Body0 , UserHeaders1 }
211- end ,
212-
213- Headers1 = merge_headers (BaseHeaders , UserHeaders2 ),
112+ UserHeaders = get_value (headers , Params , []),
113+ Headers1 = merge_headers (BaseHeaders , UserHeaders ),
214114 {Headers2 , HttpDb } = couch_replicator_auth :update_headers (HttpDb0 , Headers1 ),
215115 Url0 = full_url (HttpDb , Params ),
116+ Body = get_value (body , Params , []),
216117 case get_value (path , Params ) == " _changes" of
217118 true ->
218119 Timeout = infinity ;
@@ -282,9 +183,6 @@ process_response({error, connection_closing}, Worker, HttpDb, Params, _Cb) ->
282183process_response ({ibrowse_req_id , ReqId }, Worker , HttpDb , Params , Callback ) ->
283184 process_stream_response (ReqId , Worker , HttpDb , Params , Callback );
284185process_response ({ok , Code , Headers , Body }, Worker , HttpDb , Params , Callback ) ->
285- % Decompress body if it's gzip compressed
286- DecompressedBody = decompress_body (Headers , Body ),
287-
288186 case list_to_integer (Code ) of
289187 R when R =:= 301 ; R =:= 302 ; R =:= 303 ->
290188 backoff_success (HttpDb , Params ),
@@ -298,7 +196,7 @@ process_response({ok, Code, Headers, Body}, Worker, HttpDb, Params, Callback) ->
298196 backoff_success (HttpDb , Params ),
299197 couch_stats :increment_counter ([couch_replicator , responses , success ]),
300198 EJson =
301- case DecompressedBody of
199+ case Body of
302200 <<>> ->
303201 null ;
304202 Json ->
0 commit comments