From 4bcc745f7abbc1a2ab5c4914258c6d06da955f0f Mon Sep 17 00:00:00 2001 From: GPT on behalf of DHH <2741+dhh@users.noreply.github.com> Date: Thu, 8 Oct 2026 15:28:42 +0200 Subject: [PATCH] Interrupt WAL checkpoint shutdown and retain exclusive lock ownership --- lib/rails_ext/sqlite_wal_checkpoint.rb | 90 +++++++++-------- test/lib/sqlite_wal_checkpoint_test.rb | 130 ++++++++++++++++++++++++- 2 files changed, 180 insertions(+), 40 deletions(-) diff --git a/lib/rails_ext/sqlite_wal_checkpoint.rb b/lib/rails_ext/sqlite_wal_checkpoint.rb index 1ccda79..7b1e35b 100644 --- a/lib/rails_ext/sqlite_wal_checkpoint.rb +++ b/lib/rails_ext/sqlite_wal_checkpoint.rb @@ -8,7 +8,12 @@ module SqliteWalCheckpoint INTERVAL = 0.25 MAX_BACKOFF = 30.0 + STOP_TIMEOUT = 5.0 LOCK_PATH = Rails.root.join("tmp/pids/sqlite_wal_checkpoint.lock") + LIFECYCLE_MUTEX = Mutex.new + private_constant :LIFECYCLE_MUTEX + + class StopTimeout < StandardError; end class << self attr_writer :lock_path, :database_path_override @@ -16,28 +21,34 @@ module SqliteWalCheckpoint def start(interval: INTERVAL, enabled: !Rails.env.test?) return unless enabled - @mutex ||= Mutex.new - @mutex.synchronize do - return if @thread&.alive? + LIFECYCLE_MUTEX.synchronize do + if @thread&.alive? + raise StopTimeout, "SQLite WAL checkpointer is still stopping" if @wakeup.closed? - @stop = false + return @thread + end + + @wakeup = Thread::Queue.new install_exit_checkpoint - @thread = Thread.new { run(interval) } + @thread = Thread.new(@wakeup) { |wakeup| run(interval, wakeup) } @thread.report_on_exception = false + @thread end - - @thread end # Signal the contender to exit and wait until it has released the flock and # closed its SQLite connection. before_fork must not return while those are open. - def stop - @stop = true - thread = @mutex&.synchronize { @thread } - return unless thread + def stop(timeout: STOP_TIMEOUT) + LIFECYCLE_MUTEX.synchronize do + return unless @thread - thread.join - @mutex&.synchronize { @thread = nil if @thread.equal?(thread) } + @wakeup.close + unless @thread.join(timeout) + # Never let a caller fork with live SQLite or flock ownership. + raise StopTimeout, "SQLite WAL checkpointer did not stop within #{timeout}s" + end + @thread = @wakeup = nil + end end def checkpoint @@ -47,13 +58,14 @@ module SqliteWalCheckpoint # One leadership attempt for tests: acquire the lock, checkpoint once, release. def tick return :no_database unless database_path - return :busy unless acquire_lock + lock = acquire_lock + return :busy unless lock begin result = checkpoint result.nil? ? :no_database : :checkpointed ensure - release_lock + release_lock(lock) end end @@ -65,25 +77,24 @@ module SqliteWalCheckpoint stop @lock_path = nil @database_path_override = nil - @stop = false - @exit_checkpoint_installed = false end private - def run(interval) + def run(interval, wakeup) Thread.current.name = "sqlite-wal-checkpoint" backoff = interval - until @stop + until wakeup.closed? begin path = database_path unless path && File.exist?(path) - sleep interval + wait(wakeup, interval) next end - unless acquire_lock - sleep interval + lock = acquire_lock + unless lock + wait(wakeup, interval) next end @@ -91,27 +102,29 @@ module SqliteWalCheckpoint ran = false with_database do |database| ran = true - until @stop + until wakeup.closed? checkpoint_on(database) backoff = interval - sleep interval + wait(wakeup, interval) end end ensure - release_lock + release_lock(lock) end # with_database no-ops if the file vanished between the exist? check # and open; sleep so we do not spin on the lock file. - sleep interval unless ran || @stop + wait(wakeup, interval) unless ran || wakeup.closed? rescue => error Rails.logger.warn "SQLite WAL checkpoint failed: #{error.class}: #{error.message}" - sleep backoff + wait(wakeup, backoff) backoff = [ backoff * 2, MAX_BACKOFF ].min end end - ensure - release_lock + end + + def wait(wakeup, timeout) + wakeup.pop(timeout: timeout) end def checkpoint_on(database) @@ -124,7 +137,7 @@ module SqliteWalCheckpoint result = nil SQLite3::Database.new(path) do |database| - # Keep below stop's join timeout so before_fork can finish cleanly. + # Bounds busy-handler retries, not PASSIVE checkpoint I/O. database.busy_handler_timeout = 1_000 result = yield database end @@ -144,21 +157,22 @@ module SqliteWalCheckpoint FileUtils.mkdir_p(File.dirname(lock_path)) file = File.open(lock_path, File::RDWR | File::CREAT, 0644) if file.flock(File::LOCK_EX | File::LOCK_NB) - @lock_file = file - true + file else file.close - false + nil end + rescue + file&.close + raise end - def release_lock - return unless @lock_file + def release_lock(file) + return unless file - @lock_file.flock(File::LOCK_UN) - @lock_file.close + file.flock(File::LOCK_UN) ensure - @lock_file = nil + file&.close end # Best-effort PASSIVE for short-lived console/rake writers that exit before diff --git a/test/lib/sqlite_wal_checkpoint_test.rb b/test/lib/sqlite_wal_checkpoint_test.rb index 508fddd..88bb8f8 100644 --- a/test/lib/sqlite_wal_checkpoint_test.rb +++ b/test/lib/sqlite_wal_checkpoint_test.rb @@ -53,6 +53,101 @@ class SqliteWalCheckpointTest < ActiveSupport::TestCase assert_equal :checkpointed, SqliteWalCheckpoint.tick end + test "stop interrupts a long checkpoint interval" do + waits = observe_waits + contender = SqliteWalCheckpoint.start(interval: 30, enabled: true) + assert_equal 30, receive(waits) + + SqliteWalCheckpoint.stop(timeout: 0.5) + assert_not contender.alive? + assert_empty checkpoint_threads + end + + test "stop interrupts checkpoint error backoff" do + SqliteWalCheckpoint.database_path_override = build_wal_database(rows: 10) + SqliteWalCheckpoint.stubs(:checkpoint_on).raises(SQLite3::Exception, "checkpoint failed") + waits = observe_waits + contender = SqliteWalCheckpoint.start(interval: 30, enabled: true) + assert_equal 30, receive(waits) + + SqliteWalCheckpoint.stop(timeout: 0.5) + assert_not contender.alive? + assert_not lock_held_by_other_process?(SqliteWalCheckpoint.lock_path) + end + + test "timed out stop prevents replacement until the checkpoint releases its resources" do + entered, release = block_checkpoint + contender = SqliteWalCheckpoint.start(enabled: true) + database = receive(entered) + assert lock_held_by_other_process?(SqliteWalCheckpoint.lock_path) + + assert_raises(SqliteWalCheckpoint::StopTimeout) { SqliteWalCheckpoint.stop(timeout: 0.01) } + assert contender.alive? + assert lock_held_by_other_process?(SqliteWalCheckpoint.lock_path) + assert_raises(SqliteWalCheckpoint::StopTimeout) { SqliteWalCheckpoint.start(enabled: true) } + assert_equal [ contender ], checkpoint_threads + + release << true + SqliteWalCheckpoint.stop + assert_not contender.alive? + assert database.closed? + assert_not lock_held_by_other_process?(SqliteWalCheckpoint.lock_path) + ensure + release&.push(true) + SqliteWalCheckpoint.stop + end + + test "start waits for a concurrent stop to release the previous contender" do + entered, release = block_checkpoint + contender = SqliteWalCheckpoint.start(enabled: true) + receive(entered) + stopping = Thread.new { SqliteWalCheckpoint.stop } + wait_until { stopping.status == "sleep" } + starting = Thread.new { SqliteWalCheckpoint.start(enabled: true) } + wait_until { starting.status == "sleep" } + assert contender.alive? + + release << true + stopping.value + replacement = starting.value + assert_not contender.alive? + assert replacement.alive? + assert_not_equal contender, replacement + ensure + release&.push(true) + stopping&.join + starting&.join + SqliteWalCheckpoint.stop + end + + test "a contender cannot release a concurrent tick's lock" do + entered, release = block_checkpoint + ticking = Thread.new { SqliteWalCheckpoint.tick } + receive(entered) + waits = observe_waits + contender = SqliteWalCheckpoint.start(interval: 30, enabled: true) + assert_equal 30, receive(waits) + + SqliteWalCheckpoint.stop(timeout: 0.5) + assert_not contender.alive? + assert lock_held_by_other_process?(SqliteWalCheckpoint.lock_path) + release << true + assert_equal :checkpointed, ticking.value + assert_not lock_held_by_other_process?(SqliteWalCheckpoint.lock_path) + ensure + release&.push(true) + ticking&.join + SqliteWalCheckpoint.stop + end + + test "reset does not reinstall the process exit checkpoint" do + SqliteWalCheckpoint.start(enabled: true) + SqliteWalCheckpoint.reset! + SqliteWalCheckpoint.expects(:at_exit).never + SqliteWalCheckpoint.start(enabled: true) + SqliteWalCheckpoint.stop + end + test "tick checkpoints through the elected lock holder" do db_path = build_wal_database(rows: 50) SqliteWalCheckpoint.database_path_override = db_path @@ -94,12 +189,13 @@ class SqliteWalCheckpointTest < ActiveSupport::TestCase assert File.size(wal_path) > 0 SqliteWalCheckpoint.database_path_override = db_path - assert SqliteWalCheckpoint.send(:acquire_lock) + lock = SqliteWalCheckpoint.send(:acquire_lock) + assert lock begin _busy, _log, checkpointed = SqliteWalCheckpoint.checkpoint assert checkpointed.to_i > 0, "expected PASSIVE checkpoint to copy WAL pages, got #{checkpointed.inspect}" ensure - SqliteWalCheckpoint.send(:release_lock) + SqliteWalCheckpoint.send(:release_lock, lock) end end end @@ -129,6 +225,36 @@ class SqliteWalCheckpointTest < ActiveSupport::TestCase end private + def observe_waits + waits = Thread::Queue.new + SqliteWalCheckpoint.stubs(:wait).with do |wakeup, timeout| + waits << timeout + wakeup.pop(timeout: timeout) + true + end + waits + end + + def block_checkpoint + SqliteWalCheckpoint.database_path_override = build_wal_database(rows: 10) + entered = Thread::Queue.new + release = Thread::Queue.new + blocked = false + SqliteWalCheckpoint.stubs(:checkpoint_on).returns([ 0, 0, 0 ]).with do |database| + unless blocked + blocked = true + entered << database + release.pop + end + true + end + [ entered, release ] + end + + def receive(queue) + queue.pop(timeout: 2).tap { |value| assert_not_nil value, "expected checkpoint barrier" } + end + def checkpoint_threads Thread.list.select { |thread| thread.name == "sqlite-wal-checkpoint" } end