Skip to content
Merged
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
27 changes: 11 additions & 16 deletions lib/librevox/listener/base.rb
Original file line number Diff line number Diff line change
Expand Up @@ -38,24 +38,19 @@ def api
def send_message(msg)
promise = Async::Promise.new

# Serialize only the enqueue+send so the FIFO order of @reply_promises
# always matches the byte order on the wire, even when concurrent event
# hooks send commands. Must not cover promise.wait — holding the lock
# while awaiting a reply would deadlock every other sender.
# Send first, then record the reply promise — both under the write lock
# so @reply_promises order always matches the byte order on the wire,
# even when concurrent event hooks send commands. Recording only *after*
# send_data returns is what keeps the queue honest: a sender cancelled
# (Async::Stop) or failed mid-send never enqueues a promise, so it can
# never leave an orphaned slot for a later reply to be misrouted onto.
# There is no yield between the send and the push, so no reply for this
# command can be observed before its promise is queued. Must not cover
# promise.wait — holding the lock while awaiting a reply would deadlock
# every other sender.
@write_lock.acquire do
@connection.send_data(msg)
@reply_promises << promise
begin
@connection.send_data(msg)
rescue Exception => error
# A promise may only stay queued if its command fully reached the
# wire — otherwise no reply will ever arrive for it, and every
# later reply would resolve the wrong sender's promise. Must catch
# Exception: a sender cancelled mid-send raises Async::Stop, which
# is not a StandardError.
@reply_promises.delete(promise)
promise.reject(error)
raise
end
end

reply = promise.wait
Expand Down
29 changes: 15 additions & 14 deletions test/functional/librevox/listener/cancellation_poisoning_test.rb
Original file line number Diff line number Diff line change
Expand Up @@ -60,13 +60,14 @@ def close; end
def close_write; end
end

# Regression tests for cancellation-safety of the send path. Base#send_message
# pushes the reply promise onto @reply_promises *before* @connection.send_data
# returns; because replies are matched purely positionally
# (@reply_promises.shift in receive_message), a promise left queued by a sender
# that was cancelled mid-send would permanently shift reply correlation by one
# for every later sender on the same shared connection. send_message must
# therefore remove (and reject) the promise when send_data is interrupted.
# Regression tests for cancellation-safety of the send path. Replies are
# matched to senders purely positionally (@reply_promises.shift in
# receive_message), so a promise sitting in the queue with no command behind it
# would permanently shift reply correlation by one for every later sender on
# the shared connection. Base#send_message avoids that by enqueuing the reply
# promise only *after* @connection.send_data returns: a sender cancelled
# (Async::Stop) or failed mid-send never enqueues, so it can never orphan a
# slot for a later reply to be misrouted onto.
class TestCancellationPoisoning < Minitest::Test
prepend Librevox::Test::AsyncTest

Expand Down Expand Up @@ -119,10 +120,9 @@ def test_stopping_a_sender_mid_send_data_does_not_eat_a_later_senders_reply
resolved = bounded_yield { c_result }
assert resolved,
"C's send_message never returned — its reply promise was starved. " \
"B's cancelled-mid-send promise was not removed from @reply_promises, " \
"so the reply meant for C (the 2nd real reply on the wire) was shifted " \
"onto B's orphaned promise instead, leaving C's promise permanently " \
"unresolved."
"B was enqueued with no command behind it, so the reply meant for C " \
"(the 2nd real reply on the wire) was shifted onto B's orphaned promise " \
"instead, leaving C's promise permanently unresolved."

assert_equal "+OK reply-for-C", c_result.headers[:reply_text],
"C's send_message resolved with the wrong reply — its own reply was " \
Expand Down Expand Up @@ -191,9 +191,10 @@ def test_reply_routing_stays_correct_for_all_senders_after_a_cancellation
d&.stop
end

# A send that fails with an exception (rather than being cancelled) must
# take the same cleanup path: the error propagates to the sender, and the
# orphaned promise must not stay behind to eat a later sender's reply.
# A send that fails with an exception (rather than being cancelled) is the
# same story: send_data raises before the promise is enqueued, so the error
# propagates to the sender and no orphaned slot is ever created to eat a
# later sender's reply.
def test_a_failing_send_does_not_leave_its_promise_in_the_reply_queue
a_result = nil
c_result = nil
Expand Down
Loading