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
3 changes: 3 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -81,6 +81,9 @@ The options are:
- `connects_to` - set the Active Record database configuration for the Solid Cable models. All options available in Active Record can be used here.
- `polling_interval` - sets the frequency of the polling interval. (Defaults to
0.1.seconds)
- `late_commit_window` - how long after reading a message the listener keeps
looking for lower ids that commit later. A message that commits more than
this after a higher id was read is not delivered. (Defaults to 1.second)
- `message_retention` - sets the retention time for messages kept in the database. Used as the cut-off when trimming is performed. (Defaults to 1.day)
- `autotrim` - sets wether you want Solid Cable to handle autotrimming messages. (Defaults to true)
- `silence_polling` - whether to silence Active Record logs emitted when polling (Defaults to true)
Expand Down
48 changes: 36 additions & 12 deletions lib/action_cable/subscription_adapter/solid_cable.rb
Original file line number Diff line number Diff line change
Expand Up @@ -154,24 +154,48 @@ def channels
@channels ||= Concurrent::Map.new
end

# Ids are assigned when an INSERT runs but become visible when it
# commits, so a row can appear after a higher id was already read.
# last_id therefore only moves past rows that have been read for
# longer than late_commit_window; newer rows are excluded by id
# instead, so a row that commits late is still read.
def broadcast_messages
::SolidCable::Message.
broadcastable(channels.keys, last_id).each do |message|
should_broadcast_message = false
channels.compute_if_present(message.channel_hash) do |channel_last_id|
break if channel_last_id >= message.id
messages = ::SolidCable::Message.broadcastable(channels.keys, last_id)
messages = messages.where.not(id: recent_ids.keys) if recent_ids.any?

should_broadcast_message = true
message.id
end

broadcast(message) if should_broadcast_message
self.last_id = message.id
end
read_at = monotonic_time
messages.each do |message|
recent_ids[message.id] = read_at
broadcast(message) if subscribed_before?(message)
end

advance_last_id(read_at)
self.reconnect_attempt = 0
end

# A channel only receives rows above the last id when it subscribed.
def subscribed_before?(message)
subscribed_at_id = channels[message.channel_hash]
subscribed_at_id && subscribed_at_id < message.id
end

def advance_last_id(now)
window = ::SolidCable.late_commit_window
settled_id = recent_ids.filter_map { |id, read_at| id if now - read_at >= window }.max
return unless settled_id

self.last_id = settled_id
recent_ids.delete_if { |id, _| id <= settled_id }
end

def recent_ids
@recent_ids ||= {}
end

def monotonic_time
Process.clock_gettime(Process::CLOCK_MONOTONIC)
end

def broadcast(message)
super(message.channel, message.payload)
end
Expand Down
2 changes: 1 addition & 1 deletion lib/solid_cable.rb
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ class << self
delegate :connects_to, :silence_polling?, :polling_interval,
:message_retention, :autotrim?, :trim_batch_size, :use_skip_locked,
:reconnect_attempts, :writer_batch_size, :writer_batch_delay,
:encrypt?, :encryption_context_properties,
:encrypt?, :encryption_context_properties, :late_commit_window,
to: :configuration

def configuration
Expand Down
6 changes: 5 additions & 1 deletion lib/solid_cable/configuration.rb
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ def initialize(**options)
attr_writer :connects_to, :silence_polling, :polling_interval,
:message_retention, :autotrim, :trim_batch_size, :use_skip_locked,
:reconnect_attempts, :writer_batch_size, :writer_batch_delay,
:encrypt, :encryption_context_properties
:encrypt, :encryption_context_properties, :late_commit_window

def connects_to
@connects_to ||= options.connects_to.to_h.deep_transform_values(&:to_sym)
Expand All @@ -24,6 +24,10 @@ def polling_interval
parse_duration(options.polling_interval, default: 0.1.seconds)
end

def late_commit_window
@late_commit_window ||= parse_duration(options.late_commit_window, default: 1.second)
end

def message_retention
@message_retention ||= parse_duration(options.message_retention, default: 1.day)
end
Expand Down
Original file line number Diff line number Diff line change
@@ -0,0 +1,125 @@
# frozen_string_literal: true

require "test_helper"
require "concurrent"

require "active_support/core_ext/hash/indifferent_access"

# Ids are assigned when an INSERT runs but become visible when it commits, so
# two writers can commit out of id order. These tests hold one INSERT's
# transaction open while a higher id commits, which needs two real database
# sessions: transactional tests would share one connection across threads.
class ActionCable::SubscriptionAdapter::SolidCableCommitOrderTest < ActionCable::TestCase
self.use_transactional_tests = false

WAIT_WHEN_EXPECTING_EVENT = 1
WAIT_WHEN_NOT_EXPECTING_EVENT = 0.2

setup do
skip "SQLite serializes writers, so ids always commit in order" if sqlite?

server = ActionCable::Server::Base.new
server.config.cable = { adapter: "solid_cable", polling_interval: "0.01.seconds" }.with_indifferent_access
server.config.logger = Logger.new(StringIO.new).tap { |l| l.level = Logger::UNKNOWN }

@adapter = server.config.pubsub_adapter.new(server)
end

teardown do
@adapter&.shutdown
SolidCable::Message.delete_all
end

test "delivers a message that commits after a higher id" do
subscribe_as_queue("channel") do |queue|
with_open_insert("channel", "first") do
SolidCable::Message.broadcast("channel", "second")

# The listener has now read "second", an id above "first".
assert_equal "second", next_message_in_queue(queue)
end

assert_equal "first", next_message_in_queue(queue)
end
end

test "delivers a late message once" do
subscribe_as_queue("channel") do |queue|
with_open_insert("channel", "first") do
SolidCable::Message.broadcast("channel", "second")
assert_equal "second", next_message_in_queue(queue)
end

assert_equal "first", next_message_in_queue(queue)

SolidCable::Message.broadcast("channel", "third")
assert_equal "third", next_message_in_queue(queue)
end
end

test "does not deliver a late message to a channel subscribed after its id" do
subscribe_as_queue("channel") do |queue|
with_open_insert("other channel", "early") do |commit|
SolidCable::Message.broadcast("channel", "second")
assert_equal "second", next_message_in_queue(queue)

subscribe_as_queue("other channel") do |other_queue|
commit.call
sleep WAIT_WHEN_NOT_EXPECTING_EVENT
assert_empty other_queue
end
end
end
end

private
def sqlite?
SolidCable::Record.connection_db_config.adapter == "sqlite3"
end

# Inserts on its own connection and keeps the transaction open until the
# block returns or calls the yielded commit, so an id drawn in between is
# higher and commits first.
def with_open_insert(channel, payload)
inserted = Concurrent::Event.new
release = Concurrent::Event.new

writer = Thread.new do
SolidCable::Record.connection_pool.with_connection do
SolidCable::Record.transaction do
SolidCable::Message.broadcast(channel, payload)
inserted.set
release.wait(5)
end
end
end

commit = -> { release.set; writer.join }

assert inserted.wait(WAIT_WHEN_EXPECTING_EVENT), "the slow writer did not insert"
yield commit
ensure
commit&.call
end

def subscribe_as_queue(channel)
queue = Queue.new

callback = ->(data) { queue << data }
subscribed = Concurrent::Event.new
@adapter.subscribe(channel, callback, proc { subscribed.set })
subscribed.wait(WAIT_WHEN_EXPECTING_EVENT)
assert_predicate subscribed, :set?

yield queue

sleep WAIT_WHEN_NOT_EXPECTING_EVENT
assert_empty queue
ensure
@adapter.unsubscribe(channel, callback) if subscribed&.set?
end

def next_message_in_queue(queue)
Timeout.timeout(5, nil, "Failed to get next item in queue") { queue.pop }
end
end
Loading