diff --git a/lib/dalli/options.rb b/lib/dalli/options.rb index a48cd4ba..b30b7630 100644 --- a/lib/dalli/options.rb +++ b/lib/dalli/options.rb @@ -31,6 +31,12 @@ def close end end + def pipeline_get_setup + @lock.synchronize do + super + end + end + def pipeline_response_setup @lock.synchronize do super diff --git a/lib/dalli/pipelined_getter.rb b/lib/dalli/pipelined_getter.rb index cc34bc5c..26261770 100644 --- a/lib/dalli/pipelined_getter.rb +++ b/lib/dalli/pipelined_getter.rb @@ -45,6 +45,7 @@ def setup_requests(keys) ## def make_getkq_requests(groups) groups.each do |server, keys_for_server| + server.pipeline_get_setup server.request(:pipelined_get, keys_for_server) rescue DalliError, NetworkError => e Dalli.logger.debug { e.inspect } diff --git a/lib/dalli/protocol/base.rb b/lib/dalli/protocol/base.rb index 89c4fd25..8f8b7d65 100644 --- a/lib/dalli/protocol/base.rb +++ b/lib/dalli/protocol/base.rb @@ -18,7 +18,7 @@ class Base def_delegators :@value_marshaller, :serializer, :compressor, :compression_min_size, :compress_by_default? def_delegators :@connection_manager, :name, :sock, :hostname, :port, :close, :connected?, :socket_timeout, - :socket_type, :up!, :down!, :write, :reconnect_down_server?, :raise_down_error + :socket_type, :up!, :down!, :write, :write_nonblock, :reconnect_down_server?, :raise_down_error def initialize(attribs, client_options = {}) hostname, port, socket_type, @weight, user_creds = ServerConfigParser.parse(attribs) @@ -59,16 +59,19 @@ def lock!; end def unlock!; end + # verify and start request before sending GETKQ commands, + # otherwise the socket could get corrupted if we fail to reach pipline_response stage + def pipeline_get_setup + verify_state(:getkq) + end + # Start reading key/value pairs from this connection. This is usually called # after a series of GETKQ commands. A NOOP is sent, and the server begins # flushing responses for kv pairs that were found. # # Returns nothing. def pipeline_response_setup - verify_state(:getkq) - write_noop response_buffer.reset - @connection_manager.start_request! end # Attempt to receive and parse as many key/value pairs as possible @@ -204,8 +207,9 @@ def pipelined_get(keys) keys.each do |key| req << quiet_get_request(key) end - # Could send noop here instead of in pipeline_response_setup - write(req) + @connection_manager.start_request! + write_nonblock(req) + write_noop(nonblock: true) end def response_buffer diff --git a/lib/dalli/protocol/binary.rb b/lib/dalli/protocol/binary.rb index 66f71516..4d549d20 100644 --- a/lib/dalli/protocol/binary.rb +++ b/lib/dalli/protocol/binary.rb @@ -158,9 +158,13 @@ def version response_processor.version end - def write_noop + def write_noop(nonblock: false) req = RequestFormatter.standard_request(opkey: :noop) - write(req) + if nonblock + write_nonblock(req) + else + write(req) + end end require_relative 'binary/request_formatter' diff --git a/lib/dalli/protocol/connection_manager.rb b/lib/dalli/protocol/connection_manager.rb index 5a0b6ff0..0ec18239 100644 --- a/lib/dalli/protocol/connection_manager.rb +++ b/lib/dalli/protocol/connection_manager.rb @@ -171,6 +171,16 @@ def read_nonblock @sock.read_available end + # Non-blocking write. Should only be used in the context + # of a caller who has called start_request!, but not yet + # called finish_request!. Here to support the operation + # of the get_multi operation. + def write_nonblock(bytes) + @sock.write(bytes) + rescue SystemCallError, Timeout::Error => e + error_on_request!(e) + end + def max_allowed_failures @max_allowed_failures ||= @options[:socket_max_failures] || 2 end diff --git a/lib/dalli/protocol/meta.rb b/lib/dalli/protocol/meta.rb index b2e66c37..8dd293d0 100644 --- a/lib/dalli/protocol/meta.rb +++ b/lib/dalli/protocol/meta.rb @@ -162,8 +162,12 @@ def version response_processor.version end - def write_noop - write(RequestFormatter.meta_noop) + def write_noop(nonblock: false) + if nonblock + write_nonblock(RequestFormatter.meta_noop) + else + write(RequestFormatter.meta_noop) + end end def authenticate_connection