Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions lib/dalli/options.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
1 change: 1 addition & 0 deletions lib/dalli/pipelined_getter.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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 }
Expand Down
16 changes: 10 additions & 6 deletions lib/dalli/protocol/base.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
8 changes: 6 additions & 2 deletions lib/dalli/protocol/binary.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand Down
10 changes: 10 additions & 0 deletions lib/dalli/protocol/connection_manager.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
8 changes: 6 additions & 2 deletions lib/dalli/protocol/meta.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down