diff --git a/config/database.yml b/config/database.yml index ec220a4..845eb70 100644 --- a/config/database.yml +++ b/config/database.yml @@ -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: diff --git a/config/initializers/sqlite_wal_checkpoint.rb b/config/initializers/sqlite_wal_checkpoint.rb new file mode 100644 index 0000000..aa14b80 --- /dev/null +++ b/config/initializers/sqlite_wal_checkpoint.rb @@ -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 diff --git a/config/puma.rb b/config/puma.rb index 960ab6a..06220b8 100644 --- a/config/puma.rb +++ b/config/puma.rb @@ -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 diff --git a/lib/rails_ext/sqlite_wal_checkpoint.rb b/lib/rails_ext/sqlite_wal_checkpoint.rb new file mode 100644 index 0000000..1ccda79 --- /dev/null +++ b/lib/rails_ext/sqlite_wal_checkpoint.rb @@ -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 diff --git a/lib/tasks/resque.rake b/lib/tasks/resque.rake index 63ecc9d..6e04c74 100644 --- a/lib/tasks/resque.rake +++ b/lib/tasks/resque.rake @@ -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 diff --git a/test/lib/sqlite_wal_checkpoint_test.rb b/test/lib/sqlite_wal_checkpoint_test.rb new file mode 100644 index 0000000..508fddd --- /dev/null +++ b/test/lib/sqlite_wal_checkpoint_test.rb @@ -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