diff --git a/app/models/room/messages_count.rb b/app/models/room/messages_count.rb index 53254ac..948724e 100644 --- a/app/models/room/messages_count.rb +++ b/app/models/room/messages_count.rb @@ -9,27 +9,9 @@ class Room::MessagesCount class << self def install!(connection = ActiveRecord::Base.connection) - uninstall!(connection) - connection.execute <<~SQL - CREATE TRIGGER #{INSERT_TRIGGER} AFTER INSERT ON messages - BEGIN - UPDATE rooms SET messages_count = messages_count + 1 WHERE id = NEW.room_id; - END - SQL - connection.execute <<~SQL - CREATE TRIGGER #{DELETE_TRIGGER} AFTER DELETE ON messages - BEGIN - UPDATE rooms SET messages_count = messages_count - 1 WHERE id = OLD.room_id; - END - SQL - connection.execute <<~SQL - CREATE TRIGGER #{UPDATE_TRIGGER} AFTER UPDATE OF room_id ON messages - WHEN OLD.room_id IS NOT NEW.room_id - BEGIN - UPDATE rooms SET messages_count = messages_count - 1 WHERE id = OLD.room_id; - UPDATE rooms SET messages_count = messages_count + 1 WHERE id = NEW.room_id; - END - SQL + with_immediate_write(connection) do + replace_triggers!(connection) + end end def uninstall!(connection = ActiveRecord::Base.connection) @@ -47,13 +29,21 @@ class Room::MessagesCount end # schema.rb does not dump SQLite triggers; reinstall after schema:load. + # Repair rechecks, rebuilds counts, and installs triggers in one write txn + # so concurrent writers and concurrent boot repairs cannot observe a gap. def ensure!(connection = ActiveRecord::Base.connection) return unless connection.adapter_name.match?(/sqlite/i) return unless connection.data_source_exists?(:rooms) return unless connection.column_exists?(:rooms, :messages_count) - return if TRIGGERS.all? { |name| trigger_installed?(connection, name) } + return if triggers_installed?(connection) - install!(connection) + with_immediate_write(connection) do + next if triggers_installed?(connection) + + uninstall!(connection) + backfill!(connection) + create_triggers!(connection) + end end def trigger_installed?(connection, name) @@ -61,5 +51,46 @@ class Room::MessagesCount "SELECT 1 FROM sqlite_master WHERE type = 'trigger' AND name = #{connection.quote(name)}" ).present? end + + def triggers_installed?(connection = ActiveRecord::Base.connection) + TRIGGERS.all? { |name| trigger_installed?(connection, name) } + end + + private + def with_immediate_write(connection) + if connection.transaction_open? + yield + else + connection.raw_connection.transaction(:immediate) { yield } + end + end + + def replace_triggers!(connection) + uninstall!(connection) + create_triggers!(connection) + end + + def create_triggers!(connection) + connection.execute <<~SQL + CREATE TRIGGER #{INSERT_TRIGGER} AFTER INSERT ON messages + BEGIN + UPDATE rooms SET messages_count = messages_count + 1 WHERE id = NEW.room_id; + END + SQL + connection.execute <<~SQL + CREATE TRIGGER #{DELETE_TRIGGER} AFTER DELETE ON messages + BEGIN + UPDATE rooms SET messages_count = messages_count - 1 WHERE id = OLD.room_id; + END + SQL + connection.execute <<~SQL + CREATE TRIGGER #{UPDATE_TRIGGER} AFTER UPDATE OF room_id ON messages + WHEN OLD.room_id IS NOT NEW.room_id + BEGIN + UPDATE rooms SET messages_count = messages_count - 1 WHERE id = OLD.room_id; + UPDATE rooms SET messages_count = messages_count + 1 WHERE id = NEW.room_id; + END + SQL + end end end diff --git a/test/models/room/messages_count_lifecycle_test.rb b/test/models/room/messages_count_lifecycle_test.rb index 295e178..fff8a3f 100644 --- a/test/models/room/messages_count_lifecycle_test.rb +++ b/test/models/room/messages_count_lifecycle_test.rb @@ -21,25 +21,24 @@ class Room::MessagesCountLifecycleTest < ActiveSupport::TestCase Room::MessagesCount.backfill! end - test "ensure! installs missing triggers without changing existing counts" do + test "ensure! installs missing triggers and keeps accurate counts" do before = @room.reload.messages_count Room::MessagesCount.uninstall! - assert_not Room::MessagesCount.trigger_installed?(ActiveRecord::Base.connection, Room::MessagesCount::INSERT_TRIGGER) + assert_not Room::MessagesCount.triggers_installed? Room::MessagesCount.ensure! - assert Room::MessagesCount::TRIGGERS.all? { |name| - Room::MessagesCount.trigger_installed?(ActiveRecord::Base.connection, name) - } + assert Room::MessagesCount.triggers_installed? assert_equal before, @room.reload.messages_count + assert_equal @room.messages.count, @room.messages_count assert_difference -> { @room.reload.messages_count }, +1 do @room.messages.create!(creator: users(:jason), body: "After ensure", client_message_id: "count-ensure") end end - test "ensure! repairs a partial trigger install" do + test "ensure! repairs a partial trigger install and backfills drifted counts" do Room::MessagesCount.uninstall! ActiveRecord::Base.connection.execute <<~SQL CREATE TRIGGER #{Room::MessagesCount::INSERT_TRIGGER} AFTER INSERT ON messages @@ -51,11 +50,147 @@ class Room::MessagesCountLifecycleTest < ActiveSupport::TestCase assert Room::MessagesCount.trigger_installed?(ActiveRecord::Base.connection, Room::MessagesCount::INSERT_TRIGGER) assert_not Room::MessagesCount.trigger_installed?(ActiveRecord::Base.connection, Room::MessagesCount::DELETE_TRIGGER) + # Delete fires with no delete trigger → counter drifts high. + drifted = @room.messages.create!(creator: users(:jason), body: "Drift", client_message_id: "count-partial-drift") + ActiveRecord::Base.connection.execute("DELETE FROM messages WHERE id = #{drifted.id}") + assert_operator @room.reload.messages_count, :>, @room.messages.count + Room::MessagesCount.ensure! - assert Room::MessagesCount::TRIGGERS.all? { |name| - Room::MessagesCount.trigger_installed?(ActiveRecord::Base.connection, name) - } + assert Room::MessagesCount.triggers_installed? + assert_equal @room.messages.count, @room.reload.messages_count + end + + test "ensure! backfills writes that landed while the insert trigger was missing" do + Room::MessagesCount.uninstall! + before_count = @room.messages.count + before_cached = @room.reload.messages_count + assert_equal before_count, before_cached + + path = File.expand_path(ActiveRecord::Base.connection_db_config.database) + now = Time.current.utc.strftime("%Y-%m-%d %H:%M:%S.%6N") + ActiveRecord::Base.connection_pool.release_connection + + SQLite3::Database.new(path) do |db| + db.busy_timeout = 5_000 + db.execute( + "INSERT INTO messages (room_id, creator_id, client_message_id, created_at, updated_at) VALUES (?, ?, ?, ?, ?)", + [ @room.id, users(:david).id, "count-missing-insert", now, now ] + ) + end + + assert_equal before_count + 1, @room.messages.count + assert_equal before_cached, @room.reload.messages_count + + Room::MessagesCount.ensure! + + assert Room::MessagesCount.triggers_installed? + assert_equal @room.messages.count, @room.reload.messages_count + assert_equal before_cached + 1, @room.messages_count + end + + test "concurrent ensure! repairs leave triggers installed and counts accurate" do + Room::MessagesCount.uninstall! + path = File.expand_path(ActiveRecord::Base.connection_db_config.database) + now = Time.current.utc.strftime("%Y-%m-%d %H:%M:%S.%6N") + ActiveRecord::Base.connection_pool.release_connection + + SQLite3::Database.new(path) do |db| + db.busy_timeout = 5_000 + db.execute( + "INSERT INTO messages (room_id, creator_id, client_message_id, created_at, updated_at) VALUES (?, ?, ?, ?, ?)", + [ @room.id, users(:david).id, "count-concurrent-drift", now, now ] + ) + end + + assert_not_equal @room.messages.count, @room.reload.messages_count + + errors = [] + errors_mutex = Mutex.new + + workers = 2.times.map do + Thread.new do + ActiveRecord::Base.connection_pool.with_connection do |connection| + Room::MessagesCount.ensure!(connection) + end + rescue => error + errors_mutex.synchronize { errors << error } + end + end + + workers.each(&:join) + + assert_empty errors, -> { errors.map(&:full_message).join("\n") } + assert Room::MessagesCount.triggers_installed? + assert_equal @room.messages.count, @room.reload.messages_count + end + + test "writes concurrent with ensure! repair leave accurate counts" do + Room::MessagesCount.uninstall! + path = File.expand_path(ActiveRecord::Base.connection_db_config.database) + now = Time.current.utc.strftime("%Y-%m-%d %H:%M:%S.%6N") + ActiveRecord::Base.connection_pool.release_connection + + errors = [] + errors_mutex = Mutex.new + + ensure_thread = Thread.new do + ActiveRecord::Base.connection_pool.with_connection do |connection| + Room::MessagesCount.ensure!(connection) + end + rescue => error + errors_mutex.synchronize { errors << error } + end + + insert_thread = Thread.new do + SQLite3::Database.new(path) do |db| + db.busy_timeout = 10_000 + db.execute( + "INSERT INTO messages (room_id, creator_id, client_message_id, created_at, updated_at) VALUES (?, ?, ?, ?, ?)", + [ @room.id, users(:david).id, "count-during-repair", now, now ] + ) + end + rescue => error + errors_mutex.synchronize { errors << error } + end + + [ ensure_thread, insert_thread ].each(&:join) + + assert_empty errors, -> { errors.map(&:full_message).join("\n") } + assert Room::MessagesCount.triggers_installed? + assert_equal @room.messages.count, @room.reload.messages_count + assert Message.exists?(client_message_id: "count-during-repair") + end + + test "install! never exposes a partial trigger set to other connections" do + Room::MessagesCount.uninstall! + path = File.expand_path(ActiveRecord::Base.connection_db_config.database) + ActiveRecord::Base.connection_pool.release_connection + + stop = false + partial_snapshots = [] + snapshots_mutex = Mutex.new + + watcher = Thread.new do + SQLite3::Database.new(path) do |db| + db.busy_timeout = 5_000 + until stop + names = db.execute( + "SELECT name FROM sqlite_master WHERE type = 'trigger' AND name IN (#{Room::MessagesCount::TRIGGERS.map { "'#{it}'" }.join(", ")})" + ).flatten + if names.any? && names.sort != Room::MessagesCount::TRIGGERS.sort + snapshots_mutex.synchronize { partial_snapshots << names.sort } + end + end + end + end + + 25.times { Room::MessagesCount.install! } + stop = true + watcher.join + + assert_empty partial_snapshots + assert Room::MessagesCount.triggers_installed? end test "foreign SQLite connections keep the counter in step" do