mirror of
https://github.com/basecamp/once-campfire.git
synced 2026-10-09 08:10:08 +09:00
Interrupt WAL checkpoint shutdown and retain exclusive lock ownership
This commit is contained in:
@@ -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
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user