Repair messages_count triggers in one SQLite write transaction

ensure! now rechecks trigger presence under BEGIN IMMEDIATE, backfills
drifted rooms.messages_count, and reinstalls all three triggers before
commit so concurrent writers and concurrent boot repairs cannot observe
a partial install or keep a wrong total. install! uses the same
immediate write lock for drop/recreate. Lifecycle regressions cover
missing/partial trigger drift, concurrent ensure!, and writes racing
repair.

Co-authored-by: Thomas Klemm <github@tklemm.eu>
This commit is contained in:
Cursor Agent
2026-10-07 20:58:32 +00:00
parent 32748e7fd6
commit b4ab2df20f
2 changed files with 198 additions and 32 deletions
+54 -23
View File
@@ -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
@@ -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