mirror of
https://github.com/basecamp/once-campfire.git
synced 2026-10-09 08:10:08 +09:00
Merge #339: checkpoint SQLite WAL outside request threads
This commit is contained in:
@@ -9,6 +9,11 @@ default: &default
|
||||
pool: <%= ENV.fetch("RAILS_MAX_THREADS") { 10 } %>
|
||||
timeout: 5000
|
||||
default_transaction_mode: immediate
|
||||
# Checkpoint on a background connection instead (SqliteWalCheckpoint).
|
||||
# Non-test processes start a contender; forking servers stop before fork and
|
||||
# start again in the child so the flock is never inherited.
|
||||
pragmas:
|
||||
wal_autocheckpoint: 0
|
||||
|
||||
development:
|
||||
primary:
|
||||
|
||||
@@ -0,0 +1,8 @@
|
||||
# Loaded from lib/rails_ext (autoload_lib ignore list) so reloads do not orphan
|
||||
# the contender thread. Non-test processes start here; Puma/Resque stop before
|
||||
# fork and start again in the child.
|
||||
require Rails.root.join("lib/rails_ext/sqlite_wal_checkpoint")
|
||||
|
||||
Rails.application.config.after_initialize do
|
||||
SqliteWalCheckpoint.start unless Rails.env.test?
|
||||
end
|
||||
@@ -50,6 +50,12 @@ plugin :tmp_restart
|
||||
# Reset all membership connections
|
||||
Membership.disconnect_all
|
||||
|
||||
# Initializer starts a contender in this process. Stop before fork so workers
|
||||
# do not inherit the flock; each worker starts its own contender. Single-process
|
||||
# mode never forks, so the initializer's contender keeps running.
|
||||
before_fork { SqliteWalCheckpoint.stop }
|
||||
before_worker_boot { SqliteWalCheckpoint.start }
|
||||
|
||||
Signal.trap :SIGPROF do
|
||||
Thread.list.each do |t|
|
||||
puts t
|
||||
|
||||
@@ -0,0 +1,173 @@
|
||||
# SQLite's default WAL auto-checkpoint (~1,000 pages) runs inside the committing
|
||||
# writer and fsyncs during the request. database.yml sets wal_autocheckpoint=0;
|
||||
# this module copies pages off the request thread with PASSIVE checkpoints.
|
||||
#
|
||||
# Every non-test process starts a contender. Callers that fork (Puma, Resque pool)
|
||||
# must stop before fork and start again in the child so the flock is never shared
|
||||
# across an inherited file descriptor.
|
||||
module SqliteWalCheckpoint
|
||||
INTERVAL = 0.25
|
||||
MAX_BACKOFF = 30.0
|
||||
LOCK_PATH = Rails.root.join("tmp/pids/sqlite_wal_checkpoint.lock")
|
||||
|
||||
class << self
|
||||
attr_writer :lock_path, :database_path_override
|
||||
|
||||
def start(interval: INTERVAL, enabled: !Rails.env.test?)
|
||||
return unless enabled
|
||||
|
||||
@mutex ||= Mutex.new
|
||||
@mutex.synchronize do
|
||||
return if @thread&.alive?
|
||||
|
||||
@stop = false
|
||||
install_exit_checkpoint
|
||||
@thread = Thread.new { run(interval) }
|
||||
@thread.report_on_exception = false
|
||||
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
|
||||
|
||||
thread.join
|
||||
@mutex&.synchronize { @thread = nil if @thread.equal?(thread) }
|
||||
end
|
||||
|
||||
def checkpoint
|
||||
with_database { |database| checkpoint_on(database) }
|
||||
end
|
||||
|
||||
# One leadership attempt for tests: acquire the lock, checkpoint once, release.
|
||||
def tick
|
||||
return :no_database unless database_path
|
||||
return :busy unless acquire_lock
|
||||
|
||||
begin
|
||||
result = checkpoint
|
||||
result.nil? ? :no_database : :checkpointed
|
||||
ensure
|
||||
release_lock
|
||||
end
|
||||
end
|
||||
|
||||
def lock_path
|
||||
@lock_path || LOCK_PATH
|
||||
end
|
||||
|
||||
def reset!
|
||||
stop
|
||||
@lock_path = nil
|
||||
@database_path_override = nil
|
||||
@stop = false
|
||||
@exit_checkpoint_installed = false
|
||||
end
|
||||
|
||||
private
|
||||
def run(interval)
|
||||
Thread.current.name = "sqlite-wal-checkpoint"
|
||||
backoff = interval
|
||||
|
||||
until @stop
|
||||
begin
|
||||
path = database_path
|
||||
unless path && File.exist?(path)
|
||||
sleep interval
|
||||
next
|
||||
end
|
||||
|
||||
unless acquire_lock
|
||||
sleep interval
|
||||
next
|
||||
end
|
||||
|
||||
begin
|
||||
ran = false
|
||||
with_database do |database|
|
||||
ran = true
|
||||
until @stop
|
||||
checkpoint_on(database)
|
||||
backoff = interval
|
||||
sleep interval
|
||||
end
|
||||
end
|
||||
ensure
|
||||
release_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
|
||||
rescue => error
|
||||
Rails.logger.warn "SQLite WAL checkpoint failed: #{error.class}: #{error.message}"
|
||||
sleep backoff
|
||||
backoff = [ backoff * 2, MAX_BACKOFF ].min
|
||||
end
|
||||
end
|
||||
ensure
|
||||
release_lock
|
||||
end
|
||||
|
||||
def checkpoint_on(database)
|
||||
database.execute("PRAGMA wal_checkpoint(PASSIVE)").first
|
||||
end
|
||||
|
||||
def with_database
|
||||
path = database_path
|
||||
return unless path && File.exist?(path)
|
||||
|
||||
result = nil
|
||||
SQLite3::Database.new(path) do |database|
|
||||
# Keep below stop's join timeout so before_fork can finish cleanly.
|
||||
database.busy_handler_timeout = 1_000
|
||||
result = yield database
|
||||
end
|
||||
result
|
||||
end
|
||||
|
||||
def database_path
|
||||
return @database_path_override if @database_path_override
|
||||
|
||||
config = ActiveRecord::Base.connection_db_config
|
||||
return unless config.adapter.to_s == "sqlite3"
|
||||
|
||||
ActiveRecord::ConnectionAdapters::SQLite3Adapter.resolve_path(config.database)
|
||||
end
|
||||
|
||||
def acquire_lock
|
||||
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
|
||||
else
|
||||
file.close
|
||||
false
|
||||
end
|
||||
end
|
||||
|
||||
def release_lock
|
||||
return unless @lock_file
|
||||
|
||||
@lock_file.flock(File::LOCK_UN)
|
||||
@lock_file.close
|
||||
ensure
|
||||
@lock_file = nil
|
||||
end
|
||||
|
||||
# Best-effort PASSIVE for short-lived console/rake writers that exit before
|
||||
# the contender acquires the flock. Does not touch lock ownership.
|
||||
def install_exit_checkpoint
|
||||
return if @exit_checkpoint_installed
|
||||
|
||||
@exit_checkpoint_installed = true
|
||||
at_exit { checkpoint rescue nil }
|
||||
end
|
||||
end
|
||||
end
|
||||
@@ -4,9 +4,13 @@ end
|
||||
|
||||
task "resque:pool:setup" do
|
||||
ActiveRecord::Base.connection.disconnect!
|
||||
# Initializer may have started a contender in the pool master; drop it before
|
||||
# workers fork so they do not inherit the flock.
|
||||
SqliteWalCheckpoint.stop
|
||||
|
||||
Resque::Pool.after_prefork do |job|
|
||||
ActiveRecord::Base.establish_connection
|
||||
Resque.redis.client.close
|
||||
SqliteWalCheckpoint.start
|
||||
end
|
||||
end
|
||||
|
||||
@@ -0,0 +1,187 @@
|
||||
require "test_helper"
|
||||
|
||||
class SqliteWalCheckpointTest < ActiveSupport::TestCase
|
||||
setup do
|
||||
@tmpdir = Dir.mktmpdir("sqlite-wal-checkpoint")
|
||||
SqliteWalCheckpoint.reset!
|
||||
SqliteWalCheckpoint.lock_path = File.join(@tmpdir, "checkpoint.lock")
|
||||
end
|
||||
|
||||
teardown do
|
||||
SqliteWalCheckpoint.reset!
|
||||
FileUtils.remove_entry(@tmpdir) if @tmpdir && File.directory?(@tmpdir)
|
||||
end
|
||||
|
||||
test "connections disable WAL auto-checkpoint so commits do not fsync on the writer" do
|
||||
assert_equal 0, ActiveRecord::Base.connection.select_value("PRAGMA wal_autocheckpoint").to_i
|
||||
end
|
||||
|
||||
test "a passive checkpoint against the primary database does not raise" do
|
||||
assert_nothing_raised { SqliteWalCheckpoint.checkpoint }
|
||||
end
|
||||
|
||||
test "does not start a background thread in test by default" do
|
||||
assert_nil SqliteWalCheckpoint.start
|
||||
assert_empty checkpoint_threads
|
||||
end
|
||||
|
||||
test "start with enabled: true runs a named contender thread that stop joins" do
|
||||
thread = SqliteWalCheckpoint.start(interval: 0.05, enabled: true)
|
||||
assert thread.alive?
|
||||
wait_until { thread.name == "sqlite-wal-checkpoint" }
|
||||
|
||||
SqliteWalCheckpoint.stop
|
||||
assert_not thread.alive?
|
||||
assert_empty checkpoint_threads
|
||||
end
|
||||
|
||||
test "stop before start again mimics a fork-safe worker boot" do
|
||||
db_path = build_wal_database(rows: 10)
|
||||
SqliteWalCheckpoint.database_path_override = db_path
|
||||
|
||||
first = SqliteWalCheckpoint.start(interval: 0.05, enabled: true)
|
||||
wait_until { first.name == "sqlite-wal-checkpoint" }
|
||||
SqliteWalCheckpoint.stop
|
||||
assert_not first.alive?
|
||||
|
||||
second = SqliteWalCheckpoint.start(interval: 0.05, enabled: true)
|
||||
wait_until { second.name == "sqlite-wal-checkpoint" }
|
||||
assert second.alive?
|
||||
|
||||
SqliteWalCheckpoint.stop
|
||||
assert_not second.alive?
|
||||
assert_equal :checkpointed, SqliteWalCheckpoint.tick
|
||||
end
|
||||
|
||||
test "tick checkpoints through the elected lock holder" do
|
||||
db_path = build_wal_database(rows: 50)
|
||||
SqliteWalCheckpoint.database_path_override = db_path
|
||||
|
||||
assert_equal :checkpointed, SqliteWalCheckpoint.tick
|
||||
end
|
||||
|
||||
test "tick reports busy while another process holds the lock and succeeds after it exits" do
|
||||
db_path = build_wal_database(rows: 20)
|
||||
SqliteWalCheckpoint.database_path_override = db_path
|
||||
lock = SqliteWalCheckpoint.lock_path
|
||||
ready = File.join(@tmpdir, "holder-ready")
|
||||
|
||||
holder = spawn_lock_holder(lock, ready_path: ready, hold_for: 0.4)
|
||||
|
||||
begin
|
||||
wait_until { File.exist?(ready) }
|
||||
assert_equal :busy, SqliteWalCheckpoint.tick
|
||||
ensure
|
||||
Process.wait(holder)
|
||||
end
|
||||
|
||||
assert_equal :checkpointed, SqliteWalCheckpoint.tick
|
||||
end
|
||||
|
||||
test "passive checkpoint moves WAL pages after writes with autocheckpoint disabled" do
|
||||
db_path = File.join(@tmpdir, "writer-#{SecureRandom.hex(4)}.sqlite3")
|
||||
|
||||
SQLite3::Database.new(db_path) do |writer|
|
||||
writer.execute("PRAGMA journal_mode=WAL")
|
||||
writer.execute("PRAGMA wal_autocheckpoint=0")
|
||||
writer.execute("CREATE TABLE items (id INTEGER PRIMARY KEY, body TEXT)")
|
||||
200.times do |i|
|
||||
writer.execute("INSERT INTO items (body) VALUES (?)", "row-#{i}-#{"x" * 200}")
|
||||
end
|
||||
|
||||
wal_path = "#{db_path}-wal"
|
||||
assert File.exist?(wal_path), "expected a WAL file after inserts"
|
||||
assert File.size(wal_path) > 0
|
||||
|
||||
SqliteWalCheckpoint.database_path_override = db_path
|
||||
assert SqliteWalCheckpoint.send(:acquire_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)
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
test "a contender thread takes over checkpointing after the lock holder stops" do
|
||||
db_path = build_wal_database(rows: 30)
|
||||
SqliteWalCheckpoint.database_path_override = db_path
|
||||
lock = SqliteWalCheckpoint.lock_path
|
||||
ready = File.join(@tmpdir, "holder-ready")
|
||||
|
||||
holder = spawn_lock_holder(lock, ready_path: ready, hold_for: 0.25)
|
||||
wait_until { File.exist?(ready) }
|
||||
|
||||
contender = SqliteWalCheckpoint.start(interval: 0.05, enabled: true)
|
||||
|
||||
begin
|
||||
assert_equal :busy, SqliteWalCheckpoint.tick
|
||||
Process.wait(holder)
|
||||
holder = nil
|
||||
|
||||
wait_until(timeout: 2) { lock_held_by_other_process?(lock) }
|
||||
ensure
|
||||
Process.wait(holder) if holder
|
||||
SqliteWalCheckpoint.stop
|
||||
assert_not contender.alive?
|
||||
end
|
||||
end
|
||||
|
||||
private
|
||||
def checkpoint_threads
|
||||
Thread.list.select { |thread| thread.name == "sqlite-wal-checkpoint" }
|
||||
end
|
||||
|
||||
def build_wal_database(rows:)
|
||||
path = File.join(@tmpdir, "writer-#{SecureRandom.hex(4)}.sqlite3")
|
||||
|
||||
SQLite3::Database.new(path) do |database|
|
||||
database.execute("PRAGMA journal_mode=WAL")
|
||||
database.execute("PRAGMA wal_autocheckpoint=0")
|
||||
database.execute("CREATE TABLE items (id INTEGER PRIMARY KEY, body TEXT)")
|
||||
rows.times do |i|
|
||||
database.execute("INSERT INTO items (body) VALUES (?)", "row-#{i}-#{"x" * 200}")
|
||||
end
|
||||
end
|
||||
|
||||
path
|
||||
end
|
||||
|
||||
def spawn_lock_holder(lock, ready_path:, hold_for:)
|
||||
Process.spawn(
|
||||
RbConfig.ruby, "-e", <<~RUBY
|
||||
require "fileutils"
|
||||
lock = #{lock.inspect}
|
||||
ready = #{ready_path.inspect}
|
||||
FileUtils.mkdir_p(File.dirname(lock))
|
||||
file = File.open(lock, File::RDWR | File::CREAT, 0644)
|
||||
abort "lock failed" unless file.flock(File::LOCK_EX | File::LOCK_NB)
|
||||
File.write(ready, "1")
|
||||
sleep #{hold_for}
|
||||
RUBY
|
||||
)
|
||||
end
|
||||
|
||||
def lock_held_by_other_process?(lock)
|
||||
return false unless File.exist?(lock)
|
||||
|
||||
probe = File.open(lock, File::RDWR | File::CREAT, 0644)
|
||||
if probe.flock(File::LOCK_EX | File::LOCK_NB)
|
||||
probe.flock(File::LOCK_UN)
|
||||
probe.close
|
||||
false
|
||||
else
|
||||
probe.close
|
||||
true
|
||||
end
|
||||
end
|
||||
|
||||
def wait_until(timeout: 1)
|
||||
deadline = Process.clock_gettime(Process::CLOCK_MONOTONIC) + timeout
|
||||
until yield
|
||||
raise "condition not met within #{timeout}s" if Process.clock_gettime(Process::CLOCK_MONOTONIC) > deadline
|
||||
sleep 0.02
|
||||
end
|
||||
end
|
||||
end
|
||||
Reference in New Issue
Block a user