diff --git a/README.md b/README.md index 8a2032c..44df204 100644 --- a/README.md +++ b/README.md @@ -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) diff --git a/lib/action_cable/subscription_adapter/solid_cable.rb b/lib/action_cable/subscription_adapter/solid_cable.rb index 9692af1..5051830 100644 --- a/lib/action_cable/subscription_adapter/solid_cable.rb +++ b/lib/action_cable/subscription_adapter/solid_cable.rb @@ -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 diff --git a/lib/solid_cable.rb b/lib/solid_cable.rb index 0d62bc8..0824f8d 100644 --- a/lib/solid_cable.rb +++ b/lib/solid_cable.rb @@ -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 diff --git a/lib/solid_cable/configuration.rb b/lib/solid_cable/configuration.rb index 23d00b7..aad7dfe 100644 --- a/lib/solid_cable/configuration.rb +++ b/lib/solid_cable/configuration.rb @@ -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) @@ -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 diff --git a/test/lib/action_cable/subscription_adapter/solid_cable_commit_order_test.rb b/test/lib/action_cable/subscription_adapter/solid_cable_commit_order_test.rb new file mode 100644 index 0000000..e095fd8 --- /dev/null +++ b/test/lib/action_cable/subscription_adapter/solid_cable_commit_order_test.rb @@ -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