From 9f8297687c4ba46045ba6153d9a3e3ded92e30e6 Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Mon, 5 Oct 2026 16:42:49 -0700 Subject: [PATCH 01/12] fix: stop the topology monitor before dropping a database in aws_postgresql PostgreSQL refuses to drop a database that other sessions are connected to. ActiveRecord disconnects its own connections before dropping, but the wrapper's topology monitor, which is shared by the cluster's connections, kept its own connections to the database it was started for. db:drop, db:reset, db:test:prepare, and db:purge therefore failed with "database is being accessed by other users" whenever an Aurora dialect was in use. AwsPostgreSQLAdapter#drop_database now stops the topology monitor through the new WrapperPgConnection#stop_topology_monitor before dropping. The next connection that needs topology starts a new monitor. MySQL does not refuse to drop a database in use, so aws_mysql2 is unchanged. --- CHANGELOG.md | 1 + .../active_record/aws_postgresql_adapter.rb | 8 ++++++++ .../postgresql.rb | 8 ++++++++ spec/unit/aws_postgresql_adapter_spec.rb | 18 ++++++++++++++++++ spec/unit/wrapper_pg_connection_spec.rb | 12 ++++++++++++ 5 files changed, 47 insertions(+) diff --git a/CHANGELOG.md b/CHANGELOG.md index 1a2c9105..8cbc9fef 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), - Fixed the `aws_postgresql` ActiveRecord adapter passing ActiveRecord-only `database.yml` keys, such as the Rails 8.1 pool settings (`max_connections`, `min_connections`, `keepalive`, `max_age`) and multi-database settings (`replica`, `migrations_paths`, `database_tasks`), to `pg`, which rejected them with `PG::Error: invalid connection option` ([PR #191](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/191)). - Fixed the `aws_mysql2` ActiveRecord adapter dropping the `flags`, `encoding`, `socket`, and `reconnect` client options, which removed ActiveRecord's `FOUND_ROWS` flag so that statements such as `update_all` reported changed rows instead of matched rows ([PR #193](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/193)). - Fixed the `aws_postgresql` and `aws_mysql2` ActiveRecord adapters not translating connection errors the way the standard adapters do, so a missing database was not reported as `ActiveRecord::NoDatabaseError` and `db:prepare` (and `bin/setup`) failed instead of creating it ([PR #193](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/193)). +- Fixed `db:drop`, `db:reset`, and `db:test:prepare` failing on Aurora PostgreSQL with `database "" is being accessed by other users`, because the topology monitor kept its own connections to the database being dropped ([PR #TBD](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/TBD)). ## [1.0.0] - 2026-10-05 diff --git a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb index c596854d..1f86c380 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb @@ -90,6 +90,14 @@ def active? super end + # Used by db:drop, db:reset, db:test:prepare and db:purge. ActiveRecord disconnects its own + # connections first, but the topology monitor shared by this cluster's connections keeps its own + # connections to the database, and PostgreSQL refuses to drop a database other sessions are using. + def drop_database(name) + raw_connection.stop_topology_monitor + super + end + def translate_exception(exception, message:, sql:, binds:) return super unless exception.is_a?(AwsAdvancedRubyDriverWrapper::Errors::AwsError) diff --git a/lib/aws_advanced_ruby_driver_wrapper/postgresql.rb b/lib/aws_advanced_ruby_driver_wrapper/postgresql.rb index 4346b76c..c84f7f69 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/postgresql.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/postgresql.rb @@ -224,6 +224,14 @@ def close alias finish close + # Stops the background topology monitor of this connection's cluster, which closes the connections it + # holds to the database it was started for. The next connection that needs topology starts a new one. + # PostgreSQL refuses to drop a database that other sessions are connected to, so this has to happen + # before dropping the database the monitor was started for. + def stop_topology_monitor + @service_container.host_service.host_list_provider&.stop_monitor + end + # Resets the connection through the pipeline and returns this wrapper, so the reset connection stays # usable through it. The driver's own reset tears down and re-establishes the underlying socket, which # clears any server-side session state, so the tracked session state is reset to match. diff --git a/spec/unit/aws_postgresql_adapter_spec.rb b/spec/unit/aws_postgresql_adapter_spec.rb index d1597a62..1ecf0de9 100644 --- a/spec/unit/aws_postgresql_adapter_spec.rb +++ b/spec/unit/aws_postgresql_adapter_spec.rb @@ -108,6 +108,24 @@ def connect_failing_with(message, client_config = config) end end + describe '#drop_database' do + # DROP DATABASE fails while any other session is connected to the database, and the topology monitor + # keeps connections to the database it was created for after ActiveRecord has disconnected. + it 'stops the topology monitor before dropping the database' do + adapter = described_class.allocate + wrapper_connection = instance_double(AwsAdvancedRubyDriverWrapper::WrapperPgConnection) + calls = [] + allow(adapter).to receive(:raw_connection).and_return(wrapper_connection) + allow(wrapper_connection).to receive(:stop_topology_monitor) { calls << :stop_topology_monitor } + allow(adapter).to receive(:quote_table_name) { |name| %("#{name}") } + allow(adapter).to receive(:execute) { |sql| calls << sql } + + adapter.drop_database('mydb') + + expect(calls).to eq([:stop_topology_monitor, 'DROP DATABASE IF EXISTS "mydb"']) + end + end + describe '#translate_exception' do let(:adapter) { described_class.allocate } let(:sql) { 'SELECT 1' } diff --git a/spec/unit/wrapper_pg_connection_spec.rb b/spec/unit/wrapper_pg_connection_spec.rb index faab6719..83defa67 100644 --- a/spec/unit/wrapper_pg_connection_spec.rb +++ b/spec/unit/wrapper_pg_connection_spec.rb @@ -525,6 +525,18 @@ def build_wrapper(container) end end + describe '#stop_topology_monitor' do + it 'stops the topology monitor of the host list provider without calling the plugins' do + provider = double('host_list_provider', stop_monitor: nil) + host_service = double('host_service', host_list_provider: provider) + wrapper = build_wrapper(double('service_container', host_service: host_service)) + + wrapper.stop_topology_monitor + + expect(provider).to have_received(:stop_monitor) + end + end + describe '#reset' do let(:session_state_service) { AwsAdvancedRubyDriverWrapper::Services::SessionStateService.new } From 4c78114e080679dac1e890d8a6627bd96a04e915 Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Mon, 5 Oct 2026 17:56:40 -0700 Subject: [PATCH 02/12] fix: stop the wrapper monitors when ActiveRecord clears all connections bin/rails test checks the schema from the test process, which starts a topology monitor there, then clears all ActiveRecord connections and runs db:test:prepare in a child process to drop and reload the test database. The monitor in the parent process kept its connection to the test database, so the child's DROP DATABASE failed with "database is being accessed by other users". Stopping the monitor in drop_database cannot help, because that runs in the child. The ActiveRecord adapters now prepend a ConnectionHandler extension that, after clear_all_connections!, stops the wrapper's monitors and clears the cached topology, so the next connection starts a new monitor straight away. StorageService#registered? lets it skip the cache when no wrapper connection has registered it. --- .../active_record/aws_connection_handler.rb | 44 +++++++++++ .../active_record/aws_mysql2_adapter.rb | 1 + .../active_record/aws_postgresql_adapter.rb | 1 + .../utils/storage/storage_service.rb | 6 ++ spec/unit/aws_connection_handler_spec.rb | 79 +++++++++++++++++++ .../utils/storage/storage_service_spec.rb | 8 ++ 6 files changed, 139 insertions(+) create mode 100644 lib/aws_advanced_ruby_driver_wrapper/active_record/aws_connection_handler.rb create mode 100644 spec/unit/aws_connection_handler_spec.rb diff --git a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_connection_handler.rb b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_connection_handler.rb new file mode 100644 index 00000000..5b47f87a --- /dev/null +++ b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_connection_handler.rb @@ -0,0 +1,44 @@ +# frozen_string_literal: true + +# Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +require 'active_record' +require_relative '../services/service_utility' +require_relative '../host/rds_host_list_provider' + +module ActiveRecord + module ConnectionAdapters + # Stops the wrapper's background monitors when ActiveRecord closes every connection pool, and forgets the + # cached topology so that the next connection starts a new monitor straight away. + # + # The monitors keep their own connections to the database, outside ActiveRecord's pools. Closing + # everything is what ActiveRecord does before another process takes over the database, for example + # when bin/rails test runs db:test:prepare in a child process to drop and reload the test database, + # and PostgreSQL refuses to drop a database that a monitor is still connected to. The next + # connection that needs a monitor starts a new one. + module AwsConnectionHandler + def clear_all_connections!(...) + super + core = AwsAdvancedRubyDriverWrapper::Services::CoreServices + core.monitor_service.stop_and_remove_all + # Without this, connections would keep reading the cached topology with no monitor refreshing it. + topology = AwsAdvancedRubyDriverWrapper::Host::RdsHostListProvider::TOPOLOGY_CACHE_NAME + core.storage_service.clear(topology) if core.storage_service.registered?(topology) + end + end + + ConnectionHandler.prepend(AwsConnectionHandler) + end +end diff --git a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_mysql2_adapter.rb b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_mysql2_adapter.rb index bd5c444a..15ea481a 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_mysql2_adapter.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_mysql2_adapter.rb @@ -17,6 +17,7 @@ require 'active_record/connection_adapters/mysql2_adapter' require_relative '../mysql' require_relative '../errors' +require_relative 'aws_connection_handler' module ActiveRecord module ConnectionAdapters diff --git a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb index 1f86c380..dd44b4f5 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb @@ -17,6 +17,7 @@ require 'active_record/connection_adapters/postgresql_adapter' require_relative '../postgresql' require_relative '../errors' +require_relative 'aws_connection_handler' module ActiveRecord module ConnectionAdapters diff --git a/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb b/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb index eadebdf9..40ed4ef5 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb @@ -107,6 +107,12 @@ def remove(name, key) fetch_cache!(name).remove(key) end + # @param name [Symbol] cache name. + # @return [Boolean] whether a cache is registered under the name + def registered?(name) + @lock.synchronize { @caches.key?(name) } + end + # Clears all items for a given name. # @param name [Symbol] registered cache name. def clear(name) diff --git a/spec/unit/aws_connection_handler_spec.rb b/spec/unit/aws_connection_handler_spec.rb new file mode 100644 index 00000000..6df1f076 --- /dev/null +++ b/spec/unit/aws_connection_handler_spec.rb @@ -0,0 +1,79 @@ +# frozen_string_literal: true + +# Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +require_relative '../spec_helper' +require 'aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter' +require 'aws_advanced_ruby_driver_wrapper/active_record/aws_mysql2_adapter' + +# Rails closes every pool and then drops the test database from a separate process (bin/rails test +# runs db:test:prepare as a child when the schema has changed). The wrapper's monitors keep their own +# connections to the database, so they have to stop when ActiveRecord closes everything. +RSpec.describe 'ActiveRecord::ConnectionAdapters::AwsConnectionHandler' do + let(:handler) { ActiveRecord::ConnectionAdapters::ConnectionHandler.new } + let(:monitor_service) { AwsAdvancedRubyDriverWrapper::Services::CoreServices.monitor_service } + let(:storage) { AwsAdvancedRubyDriverWrapper::Services::CoreServices.storage_service } + let(:topology) { AwsAdvancedRubyDriverWrapper::Host::RdsHostListProvider::TOPOLOGY_CACHE_NAME } + let(:pool) { double('pool') } + let(:calls) { [] } + + before do + allow(pool).to receive(:disconnect!) { calls << :disconnect_pool } + allow(monitor_service).to receive(:stop_and_remove_all) { calls << :stop_monitors } + end + + it 'stops the wrapper monitors after clearing all connections' do + allow(handler).to receive(:each_connection_pool).with(nil).and_return([pool]) + + handler.clear_all_connections! + + expect(calls).to eq(%i[disconnect_pool stop_monitors]) + end + + it 'clears the cached topology so the next connection starts a new monitor' do + storage.register(topology, ttl: 300) + storage.set(topology, 'my-cluster', [:host]) + allow(handler).to receive(:each_connection_pool).with(nil).and_return([pool]) + + handler.clear_all_connections! + + expect(storage.get(topology, 'my-cluster', register_access: false)).to be_nil + end + + it 'does not fail when no wrapper connection has registered the topology cache' do + allow(storage).to receive(:registered?).with(topology).and_return(false) + allow(storage).to receive(:clear) + allow(handler).to receive(:each_connection_pool).with(nil).and_return([pool]) + + expect { handler.clear_all_connections! }.not_to raise_error + expect(storage).not_to have_received(:clear) + end + + it 'passes the role through to ActiveRecord' do + allow(handler).to receive(:each_connection_pool).with(:all).and_return([pool]) + + handler.clear_all_connections!(:all) + + expect(calls).to eq(%i[disconnect_pool stop_monitors]) + end + + it 'does not stop the monitors when only idle connections are flushed' do + allow(handler).to receive(:each_connection_pool).and_return([]) + + handler.flush_idle_connections! + + expect(monitor_service).not_to have_received(:stop_and_remove_all) + end +end diff --git a/spec/unit/utils/storage/storage_service_spec.rb b/spec/unit/utils/storage/storage_service_spec.rb index 23878d64..3bc51415 100644 --- a/spec/unit/utils/storage/storage_service_spec.rb +++ b/spec/unit/utils/storage/storage_service_spec.rb @@ -100,6 +100,14 @@ end end + describe '#registered?' do + it 'is true only for a registered cache' do + service.register(:data, ttl: 60) + expect(service.registered?(:data)).to be true + expect(service.registered?(:other)).to be false + end + end + describe '#clear' do it 'removes all items for a name' do service.register(:data, ttl: 60) From f996d635ef18b06a382273467fc2f20c5c84b414 Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Mon, 5 Oct 2026 18:31:56 -0700 Subject: [PATCH 03/12] docs: update changelog entry for the db:drop fix --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 8cbc9fef..1cb76a38 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,7 +9,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), - Fixed the `aws_postgresql` ActiveRecord adapter passing ActiveRecord-only `database.yml` keys, such as the Rails 8.1 pool settings (`max_connections`, `min_connections`, `keepalive`, `max_age`) and multi-database settings (`replica`, `migrations_paths`, `database_tasks`), to `pg`, which rejected them with `PG::Error: invalid connection option` ([PR #191](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/191)). - Fixed the `aws_mysql2` ActiveRecord adapter dropping the `flags`, `encoding`, `socket`, and `reconnect` client options, which removed ActiveRecord's `FOUND_ROWS` flag so that statements such as `update_all` reported changed rows instead of matched rows ([PR #193](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/193)). - Fixed the `aws_postgresql` and `aws_mysql2` ActiveRecord adapters not translating connection errors the way the standard adapters do, so a missing database was not reported as `ActiveRecord::NoDatabaseError` and `db:prepare` (and `bin/setup`) failed instead of creating it ([PR #193](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/193)). -- Fixed `db:drop`, `db:reset`, and `db:test:prepare` failing on Aurora PostgreSQL with `database "" is being accessed by other users`, because the topology monitor kept its own connections to the database being dropped ([PR #TBD](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/TBD)). +- Fixed `db:drop`, `db:reset`, `db:test:prepare`, and `bin/rails test` after a schema change failing on Aurora PostgreSQL with `database "" is being accessed by other users`, because the topology monitor kept its own connections to the database being dropped ([PR #195](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/195)). ## [1.0.0] - 2026-10-05 From 5c67a5941b6392bd8efc52ce78ea182360093bee Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Mon, 5 Oct 2026 19:57:03 -0700 Subject: [PATCH 04/12] fix: also release the Blue/Green status providers before dropping a database The Blue/Green status providers keep their own monitoring connections to the database, outside ActiveRecord's pools and outside the monitor service, so with the bg plugin enabled they blocked DROP DATABASE the same way the topology monitor did. drop_database and the clear_all_connections! extension now also stop them through BlueGreenPlugin.clean_up_providers, and the extension clears the cached Blue/Green status along with the cached topology. The next initial connection starts a new provider. --- .../active_record/aws_connection_handler.rb | 13 ++++++----- .../active_record/aws_postgresql_adapter.rb | 6 +++-- spec/unit/aws_connection_handler_spec.rb | 22 ++++++++++++++----- spec/unit/aws_postgresql_adapter_spec.rb | 5 +++-- 4 files changed, 32 insertions(+), 14 deletions(-) diff --git a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_connection_handler.rb b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_connection_handler.rb index 5b47f87a..785f0809 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_connection_handler.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_connection_handler.rb @@ -17,11 +17,12 @@ require 'active_record' require_relative '../services/service_utility' require_relative '../host/rds_host_list_provider' +require_relative '../plugins/blue_green/blue_green_plugin' module ActiveRecord module ConnectionAdapters - # Stops the wrapper's background monitors when ActiveRecord closes every connection pool, and forgets the - # cached topology so that the next connection starts a new monitor straight away. + # Stops the wrapper's background monitors and Blue/Green status providers when ActiveRecord closes every + # connection pool, and forgets the state they cached, so that the next connection starts new ones straight away. # # The monitors keep their own connections to the database, outside ActiveRecord's pools. Closing # everything is what ActiveRecord does before another process takes over the database, for example @@ -33,9 +34,11 @@ def clear_all_connections!(...) super core = AwsAdvancedRubyDriverWrapper::Services::CoreServices core.monitor_service.stop_and_remove_all - # Without this, connections would keep reading the cached topology with no monitor refreshing it. - topology = AwsAdvancedRubyDriverWrapper::Host::RdsHostListProvider::TOPOLOGY_CACHE_NAME - core.storage_service.clear(topology) if core.storage_service.registered?(topology) + AwsAdvancedRubyDriverWrapper::Plugins::BlueGreen::BlueGreenPlugin.clean_up_providers + # Without this, connections would keep reading cached state that nothing refreshes any more. + cached = [AwsAdvancedRubyDriverWrapper::Host::RdsHostListProvider::TOPOLOGY_CACHE_NAME, + AwsAdvancedRubyDriverWrapper::Plugins::BlueGreen::BlueGreenPlugin::BLUE_GREEN_NAME] + cached.each { |name| core.storage_service.clear(name) if core.storage_service.registered?(name) } end end diff --git a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb index 6d69e040..4b5078f5 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb @@ -92,10 +92,12 @@ def active? end # Used by db:drop, db:reset, db:test:prepare and db:purge. ActiveRecord disconnects its own - # connections first, but the topology monitor shared by this cluster's connections keeps its own - # connections to the database, and PostgreSQL refuses to drop a database other sessions are using. + # connections first, but the topology monitor shared by this cluster's connections, and the Blue/Green + # status providers, keep their own connections to the database, and PostgreSQL refuses to drop a + # database other sessions are using. def drop_database(name) raw_connection.stop_topology_monitor + AwsAdvancedRubyDriverWrapper::Plugins::BlueGreen::BlueGreenPlugin.clean_up_providers super end diff --git a/spec/unit/aws_connection_handler_spec.rb b/spec/unit/aws_connection_handler_spec.rb index 6df1f076..632a0910 100644 --- a/spec/unit/aws_connection_handler_spec.rb +++ b/spec/unit/aws_connection_handler_spec.rb @@ -26,20 +26,22 @@ let(:monitor_service) { AwsAdvancedRubyDriverWrapper::Services::CoreServices.monitor_service } let(:storage) { AwsAdvancedRubyDriverWrapper::Services::CoreServices.storage_service } let(:topology) { AwsAdvancedRubyDriverWrapper::Host::RdsHostListProvider::TOPOLOGY_CACHE_NAME } + let(:blue_green) { AwsAdvancedRubyDriverWrapper::Plugins::BlueGreen::BlueGreenPlugin } let(:pool) { double('pool') } let(:calls) { [] } before do allow(pool).to receive(:disconnect!) { calls << :disconnect_pool } allow(monitor_service).to receive(:stop_and_remove_all) { calls << :stop_monitors } + allow(blue_green).to receive(:clean_up_providers) { calls << :stop_blue_green_providers } end - it 'stops the wrapper monitors after clearing all connections' do + it 'stops the wrapper monitors and Blue/Green status providers after clearing all connections' do allow(handler).to receive(:each_connection_pool).with(nil).and_return([pool]) handler.clear_all_connections! - expect(calls).to eq(%i[disconnect_pool stop_monitors]) + expect(calls).to eq(%i[disconnect_pool stop_monitors stop_blue_green_providers]) end it 'clears the cached topology so the next connection starts a new monitor' do @@ -52,8 +54,18 @@ expect(storage.get(topology, 'my-cluster', register_access: false)).to be_nil end - it 'does not fail when no wrapper connection has registered the topology cache' do - allow(storage).to receive(:registered?).with(topology).and_return(false) + it 'clears the cached Blue/Green status so the next connection starts a new provider' do + storage.register(blue_green::BLUE_GREEN_NAME, ttl: 3600) + storage.set(blue_green::BLUE_GREEN_NAME, '1', :status) + allow(handler).to receive(:each_connection_pool).with(nil).and_return([pool]) + + handler.clear_all_connections! + + expect(storage.get(blue_green::BLUE_GREEN_NAME, '1', register_access: false)).to be_nil + end + + it 'does not fail when no wrapper connection has registered the caches' do + allow(storage).to receive(:registered?).and_return(false) allow(storage).to receive(:clear) allow(handler).to receive(:each_connection_pool).with(nil).and_return([pool]) @@ -66,7 +78,7 @@ handler.clear_all_connections!(:all) - expect(calls).to eq(%i[disconnect_pool stop_monitors]) + expect(calls).to eq(%i[disconnect_pool stop_monitors stop_blue_green_providers]) end it 'does not stop the monitors when only idle connections are flushed' do diff --git a/spec/unit/aws_postgresql_adapter_spec.rb b/spec/unit/aws_postgresql_adapter_spec.rb index 6390a071..3bb41472 100644 --- a/spec/unit/aws_postgresql_adapter_spec.rb +++ b/spec/unit/aws_postgresql_adapter_spec.rb @@ -111,7 +111,7 @@ def connect_failing_with(message, client_config = config) describe '#drop_database' do # DROP DATABASE fails while any other session is connected to the database, and the topology monitor # keeps connections to the database it was created for after ActiveRecord has disconnected. - it 'stops the topology monitor before dropping the database' do + it 'stops the topology monitor and Blue/Green status providers before dropping the database' do adapter = described_class.allocate wrapper_connection = instance_double(AwsAdvancedRubyDriverWrapper::WrapperPgConnection) calls = [] @@ -119,10 +119,11 @@ def connect_failing_with(message, client_config = config) allow(wrapper_connection).to receive(:stop_topology_monitor) { calls << :stop_topology_monitor } allow(adapter).to receive(:quote_table_name) { |name| %("#{name}") } allow(adapter).to receive(:execute) { |sql| calls << sql } + allow(AwsAdvancedRubyDriverWrapper::Plugins::BlueGreen::BlueGreenPlugin).to receive(:clean_up_providers) { calls << :stop_blue_green_providers } adapter.drop_database('mydb') - expect(calls).to eq([:stop_topology_monitor, 'DROP DATABASE IF EXISTS "mydb"']) + expect(calls).to eq([:stop_topology_monitor, :stop_blue_green_providers, 'DROP DATABASE IF EXISTS "mydb"']) end end From b190a8309d909f51ce7f9fea208da1660cfd4f10 Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Mon, 5 Oct 2026 20:50:49 -0700 Subject: [PATCH 05/12] fix: restart wrapper monitors in forked processes without closing the parent's connections Only the forking thread survives a fork, so a forked child (Puma preload_app!, Unicorn, Resque) inherited monitors whose threads were dead. run_if_absent kept returning them, the topology never refreshed, and failover in the child could not find the new writer. The child's at_exit shutdown also closed the inherited monitoring connections, ending the parent's sessions. A Process._fork hook now runs in the child only: it forgets inherited monitors and Blue/Green providers, abandons their connections without closing them on the server (pg: socket to /dev/null; mysql2: automatic_close off), and restarts the monitor service, event publisher, and storage cleanup threads. --- CHANGELOG.md | 1 + lib/aws_advanced_ruby_driver_wrapper.rb | 25 ++++ .../driver_dialects/driver_dialect.rb | 5 + .../driver_dialects/mysql_driver_dialect.rb | 10 ++ .../driver_dialects/pg_driver_dialect.rb | 10 ++ .../monitoring/cluster_topology_monitor.rb | 7 + .../monitoring/monitor.rb | 14 ++ .../monitoring/monitor_connection.rb | 5 + .../plugins/blue_green/blue_green_plugin.rb | 5 + .../plugins/blue_green/status_monitor.rb | 4 + .../plugins/blue_green/status_provider.rb | 5 + .../services/monitor_service.rb | 11 ++ .../utils/events/batching_event_publisher.rb | 7 + .../utils/storage/storage_service.rb | 7 + .../mysql_driver_dialect_spec.rb | 9 ++ .../driver_dialects/pg_driver_dialect_spec.rb | 16 ++ spec/unit/fork_safety_spec.rb | 139 ++++++++++++++++++ .../cluster_topology_monitor_spec.rb | 18 +++ .../monitoring/monitor_connection_spec.rb | 10 ++ 19 files changed, 308 insertions(+) create mode 100644 spec/unit/fork_safety_spec.rb diff --git a/CHANGELOG.md b/CHANGELOG.md index 35b7660a..2f8079f6 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), - Fixed the `aws_postgresql` and `aws_mysql2` ActiveRecord adapters not translating connection errors the way the standard adapters do, so a missing database was not reported as `ActiveRecord::NoDatabaseError` and `db:prepare` (and `bin/setup`) failed instead of creating it ([PR #193](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/193)). - Fixed the `aws_postgresql` and `aws_mysql2` ActiveRecord adapters keeping ActiveRecord's prepared statement cache after a successful failover, so every query prepared before the failover failed on that connection with `Method invoked against old connection` until the process restarted ([PR #194](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/194), [documentation](https://aws.github.io/aws-advanced-wrapper-docs/ruby/enhanced-failover)). - Fixed `db:drop`, `db:reset`, `db:test:prepare`, and `bin/rails test` after a schema change failing on Aurora PostgreSQL with `database "" is being accessed by other users`, because the topology monitor kept its own connections to the database being dropped ([PR #195](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/195)). +- Fixed forked processes, such as Puma workers with `preload_app!`, Unicorn workers, and Resque jobs, inheriting the parent's background monitors after their threads had stopped, so the cluster topology never refreshed in the child and failover there could not find the new writer. A forked child now starts its own monitors and leaves the parent's monitoring connections open ([PR #196](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/196)). ## [1.0.0] - 2026-10-05 diff --git a/lib/aws_advanced_ruby_driver_wrapper.rb b/lib/aws_advanced_ruby_driver_wrapper.rb index 9ae43623..a54432f6 100644 --- a/lib/aws_advanced_ruby_driver_wrapper.rb +++ b/lib/aws_advanced_ruby_driver_wrapper.rb @@ -73,6 +73,29 @@ def self.clear_caches Services::HostService.clear_id_cache end + # Runs in a forked child (see ForkHook). Only the forking thread survives a fork, so every background + # thread is gone and the inherited monitors look alive but never refresh. The monitors and Blue/Green + # providers are forgotten without closing their connections, which the parent still uses, and the + # shared services restart their threads. Cached data such as topology stays valid and is kept. + def self.after_fork + Plugins::BlueGreen::BlueGreenPlugin.release_providers_after_fork + return unless defined?(Services::CoreServices) + + Services::CoreServices.monitor_service.restart_after_fork + Services::CoreServices.event_publisher.restart_after_fork + Services::CoreServices.storage_service.restart_after_fork + end + + # Process._fork is the single entry point for Kernel#fork, Process.fork and IO.popen('-'), so + # hooking it covers every way an application server or job runner forks a worker. + module ForkHook + def _fork + pid = super + AwsAdvancedRubyDriverWrapper.after_fork if pid.zero? + pid + end + end + def self.release_resources(grace_period_sec: 5) require_relative 'aws_advanced_ruby_driver_wrapper/services/service_utility' Services::CoreServices.monitor_service.shutdown(grace_period: grace_period_sec) @@ -88,6 +111,8 @@ def self.release_resources(grace_period_sec: 5) at_exit { AwsAdvancedRubyDriverWrapper.shutdown } +Process.singleton_class.prepend(AwsAdvancedRubyDriverWrapper::ForkHook) + # Register adapters with ActiveRecord via a lazy-load hook so the require order between this file and # ActiveRecord does not matter. The block runs the first time ActiveRecord::Base is referenced, whether # ActiveRecord is loaded before or after this file. The register call is itself lazy — each adapter file diff --git a/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/driver_dialect.rb b/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/driver_dialect.rb index bb02e6e1..450e0193 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/driver_dialect.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/driver_dialect.rb @@ -115,6 +115,11 @@ def close_connection(connection) raise NotImplementedError end + # Releases a connection inherited from the parent process after a fork, without telling the + # server: the parent still owns the session. Garbage collection or process exit in the child + # must not end it. + def abandon_connection(connection); end + def sql_state(_exception) nil end diff --git a/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/mysql_driver_dialect.rb b/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/mysql_driver_dialect.rb index dd48ced6..923c7ef3 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/mysql_driver_dialect.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/mysql_driver_dialect.rb @@ -131,6 +131,16 @@ def close_connection(connection) logger.error("Failed to close MySQL connection: #{e.message}") end + # With automatic_close off, mysql2 invalidates the socket instead of sending COM_QUIT when + # the client is collected, leaving the parent's session open. + def abandon_connection(connection) + return if connection.nil? || connection.closed? + + connection.automatic_close = false + rescue StandardError => e + logger.debug("Failed to abandon MySQL connection: #{e.message}") + end + def sql_state(exception) return nil unless exception.respond_to?(:sql_state) diff --git a/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/pg_driver_dialect.rb b/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/pg_driver_dialect.rb index 48b3266c..0d410879 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/pg_driver_dialect.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/driver_dialects/pg_driver_dialect.rb @@ -162,6 +162,16 @@ def close_connection(connection) logger.error("Failed to close PostgreSQL connection: #{e.message}") end + # Points the socket at /dev/null, so the Terminate message libpq sends when the connection + # is finished or collected never reaches the server. + def abandon_connection(connection) + return if connection.nil? || connection.finished? + + connection.socket_io.reopen(IO::NULL) + rescue StandardError => e + logger.debug("Failed to abandon PostgreSQL connection: #{e.message}") + end + def sql_state(exception) return nil unless exception.is_a?(::PG::Error) && exception.result diff --git a/lib/aws_advanced_ruby_driver_wrapper/monitoring/cluster_topology_monitor.rb b/lib/aws_advanced_ruby_driver_wrapper/monitoring/cluster_topology_monitor.rb index 7db87298..32fd0588 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/monitoring/cluster_topology_monitor.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/monitoring/cluster_topology_monitor.rb @@ -131,6 +131,13 @@ def close @instance_monitors_writer_conn.close end + def abandon_connections + @monitoring_connection.abandon + @instance_monitors_writer_conn.abandon + driver_dialect = @service_container.dialect_service.driver_dialect + @instance_monitor_connections.each_value { |conn| driver_dialect.abandon_connection(conn) } + end + private # Main monitoring loop. diff --git a/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor.rb b/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor.rb index 1f403367..149bf200 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor.rb @@ -73,6 +73,20 @@ def stopped? def close; end + # Detaches this monitor in a forked child. Its thread did not survive the fork and its + # connections still belong to the parent, so the thread reference is dropped and the + # connections are abandoned rather than closed: closing them would end the parent's sessions. + def release_after_fork + @stop_flag.make_true + @thread = nil + @state.set(MonitorState::STOPPED) + abandon_connections + end + + # Releases connections inherited from the parent process without closing them on the server. + # Subclasses that hold connections override this. + def abandon_connections; end + private def run diff --git a/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor_connection.rb b/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor_connection.rb index 5b28f417..6bc2220e 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor_connection.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor_connection.rb @@ -52,6 +52,11 @@ def compare_and_set(expected, new_conn) def close set(nil) end + + # Releases a connection inherited across a fork without closing it on the server. + def abandon + @driver_dialect.abandon_connection(get) + end end end end diff --git a/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/blue_green_plugin.rb b/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/blue_green_plugin.rb index 754b592e..8e1fa4a8 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/blue_green_plugin.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/blue_green_plugin.rb @@ -130,6 +130,11 @@ def self.clean_up_providers PROVIDERS.each_key { |k| PROVIDERS.delete(k)&.stop } end + # Forgets the providers inherited by a forked child so the child starts its own on first use. + def self.release_providers_after_fork + PROVIDERS.each_key { |k| PROVIDERS.delete(k)&.release_after_fork } + end + private def route_connect(host_info, driver_props, is_initial_connection, pipeline_callable, is_internal: false) diff --git a/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/status_monitor.rb b/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/status_monitor.rb index dd6eb5be..35d2cf07 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/status_monitor.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/status_monitor.rb @@ -338,6 +338,10 @@ def notify_changes @event.reset end + def abandon_connections + @connection.abandon + end + private def attempt_open_connection diff --git a/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/status_provider.rb b/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/status_provider.rb index c3faaeb9..e8eaaa5d 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/status_provider.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/status_provider.rb @@ -93,6 +93,11 @@ def stop @monitors.each_value { |m| m&.stop } end + # Detaches the monitors inherited by a forked child without closing the parent's connections. + def release_after_fork + @monitors.each_value { |m| m&.release_after_fork } + end + def log_switchover_final_summary return unless switchover_finalized? diff --git a/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb b/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb index 50110163..9c29a504 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb @@ -122,6 +122,17 @@ def shutdown(grace_period:) stop_and_remove_all end + # Resets this service in a forked child. Monitor threads do not survive a fork, so the inherited + # monitors are detached (keeping the parent's connections open) and forgotten, letting + # run_if_absent start live ones. The cleanup thread is restarted. + def restart_after_fork + @lock.synchronize { @caches.values }.each do |container| + container.cache.entries.each_key { |key| container.cache.remove(key)&.release_after_fork } + end + @running = true + @cleanup_thread = start_cleanup_thread + end + # Processes events from the event publisher. # @param event [Event] the event to process. def process_event(event) diff --git a/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb b/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb index 3eb6e9b4..41137857 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb @@ -27,6 +27,7 @@ module Events # unsubscribe(subscriber, event_classes) — deregister # publish(event) — deliver immediate or queue batched # release_resources — stop background thread + # restart_after_fork - restart background thread in a forked child class BatchingEventPublisher include Logging @@ -79,6 +80,12 @@ def release_resources @thread&.join(@message_interval_sec) end + # Restarts the publishing thread in a forked child, where it no longer runs. + def restart_after_fork + @lock.synchronize { @running = true } + @thread = start_publishing_thread + end + private def start_publishing_thread diff --git a/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb b/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb index 40ed4ef5..74952c3e 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb @@ -35,6 +35,7 @@ def initialize(event_publisher:, cleanup_interval: DEFAULT_CLEANUP_INTERVAL) @event_publisher = event_publisher @lock = Mutex.new @running = true + @cleanup_interval = cleanup_interval @cleanup_thread = start_cleanup_thread(cleanup_interval) end @@ -142,6 +143,12 @@ def shutdown @cleanup_thread&.join(5) end + # Restarts the cleanup thread in a forked child, where it no longer runs. + def restart_after_fork + @running = true + @cleanup_thread = start_cleanup_thread(@cleanup_interval) + end + private def fetch_cache!(name) diff --git a/spec/unit/driver_dialects/mysql_driver_dialect_spec.rb b/spec/unit/driver_dialects/mysql_driver_dialect_spec.rb index b258ec61..6b87f1cf 100644 --- a/spec/unit/driver_dialects/mysql_driver_dialect_spec.rb +++ b/spec/unit/driver_dialects/mysql_driver_dialect_spec.rb @@ -128,6 +128,15 @@ end end + describe '#abandon_connection' do + it 'turns off automatic_close instead of closing the connection' do + allow(connection).to receive(:closed?).and_return(false) + expect(connection).to receive(:automatic_close=).with(false) + expect(connection).not_to receive(:close) + dialect.abandon_connection(connection) + end + end + describe '#sql_state' do it 'extracts sql_state from exception' do exception = double(sql_state: '42S02') diff --git a/spec/unit/driver_dialects/pg_driver_dialect_spec.rb b/spec/unit/driver_dialects/pg_driver_dialect_spec.rb index 366f158d..76fdf84e 100644 --- a/spec/unit/driver_dialects/pg_driver_dialect_spec.rb +++ b/spec/unit/driver_dialects/pg_driver_dialect_spec.rb @@ -68,6 +68,22 @@ end end + describe '#abandon_connection' do + it 'points the socket at the null device instead of closing the connection' do + socket_io = instance_double(IO) + allow(connection).to receive_messages(finished?: false, socket_io: socket_io) + expect(socket_io).to receive(:reopen).with(IO::NULL) + expect(connection).not_to receive(:close) + dialect.abandon_connection(connection) + end + + it 'skips finished connections' do + allow(connection).to receive(:finished?).and_return(true) + expect(connection).not_to receive(:socket_io) + dialect.abandon_connection(connection) + end + end + describe '#sql_state' do it 'extracts SQLSTATE from PG::Error' do pg_result = double(error_field: '23505') diff --git a/spec/unit/fork_safety_spec.rb b/spec/unit/fork_safety_spec.rb new file mode 100644 index 00000000..da8c2ead --- /dev/null +++ b/spec/unit/fork_safety_spec.rb @@ -0,0 +1,139 @@ +# frozen_string_literal: true + +# Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +require 'rspec' +require 'aws_advanced_ruby_driver_wrapper' +require 'aws_advanced_ruby_driver_wrapper/monitoring/monitor' +require 'aws_advanced_ruby_driver_wrapper/services/service_utility' + +RSpec.describe AwsAdvancedRubyDriverWrapper, 'after fork' do + let(:core) { AwsAdvancedRubyDriverWrapper::Services::CoreServices } + let(:providers) { AwsAdvancedRubyDriverWrapper::Plugins::BlueGreen::BlueGreenPlugin::PROVIDERS } + let(:test_monitor_class) do + Class.new(AwsAdvancedRubyDriverWrapper::Monitoring::Monitor) do + attr_reader :closed, :abandoned + + def monitor + sleep(0.01) until stopped? + end + + def close + @closed = true + end + + def abandon_connections + @abandoned = true + end + end + end + + after do + providers.clear + core.reset! + end + + # Runs the block in a forked child and returns its (marshalled) result to the parent. + def in_forked_child + reader, writer = IO.pipe + pid = fork do + reader.close + result = begin + yield + rescue Exception => e # rubocop:disable Lint/RescueException + "child raised #{e.class}: #{e.message}" + end + writer.write(Marshal.dump(result)) + ensure + # Never return into RSpec from the child, or the rest of the suite would run there too. + exit!(0) + end + writer.close + result = reader.read + Process.wait(pid) + Marshal.load(result) # rubocop:disable Security/MarshalLoad + end + + def start_test_monitor + core.monitor_service.register_type(:fork_test_monitor, expiration_timeout_sec: 60) + monitor = core.monitor_service.run_if_absent(:fork_test_monitor, 'k1', double('container')) { |_| test_monitor_class.new } + sleep(0.05) + monitor + end + + it 'starts a fresh monitor in the child instead of returning the dead inherited one' do + inherited = start_test_monitor + + result = in_forked_child do + fresh = core.monitor_service.run_if_absent(:fork_test_monitor, 'k1', nil) { |_| test_monitor_class.new } + sleep(0.05) + { same: fresh.equal?(inherited), fresh_thread_alive: fresh.instance_variable_get(:@thread)&.alive?, + inherited_closed: inherited.closed, inherited_abandoned: inherited.abandoned } + end + + expect(result).to eq(same: false, fresh_thread_alive: true, inherited_closed: nil, inherited_abandoned: true) + end + + it 'leaves the parent monitor running and registered' do + inherited = start_test_monitor + + in_forked_child { core.monitor_service.get(:fork_test_monitor, 'k1') && nil } + + expect(core.monitor_service.get(:fork_test_monitor, 'k1')).to equal(inherited) + expect(inherited.state).to eq(:running) + expect(inherited.closed).to be_nil + end + + it "does not close the parent's monitors or providers when the child shuts down" do + inherited = start_test_monitor + provider = Struct.new(:events) do + def stop = events << :stop + def release_after_fork = events << :release_after_fork + end.new([]) + providers['bgd-1'] = provider + + # The at_exit hook runs this when a forked worker exits normally. + result = in_forked_child do + described_class.shutdown(grace_period_sec: 1) + { inherited_closed: inherited.closed, provider_events: provider.events } + end + + expect(result).to eq(inherited_closed: nil, provider_events: [:release_after_fork]) + end + + it 'restarts the core background threads in the child' do + threads = lambda do + [core.monitor_service.instance_variable_get(:@cleanup_thread), + core.event_publisher.instance_variable_get(:@thread), + core.storage_service.instance_variable_get(:@cleanup_thread)] + end + + expect(in_forked_child { threads.call.map(&:alive?) }).to eq([true, true, true]) + end + + it 'forgets inherited blue/green providers without stopping them' do + provider_class = Struct.new(:events) do + def stop = events << :stop + def release_after_fork = events << :release_after_fork + end + provider = provider_class.new([]) + providers['bgd-1'] = provider + + result = in_forked_child { { remaining: providers.keys, events: provider.events } } + + expect(result).to eq(remaining: [], events: [:release_after_fork]) + expect(providers['bgd-1']).to equal(provider) + end +end diff --git a/spec/unit/monitoring/cluster_topology_monitor_spec.rb b/spec/unit/monitoring/cluster_topology_monitor_spec.rb index 32a5a1f1..b984f85c 100644 --- a/spec/unit/monitoring/cluster_topology_monitor_spec.rb +++ b/spec/unit/monitoring/cluster_topology_monitor_spec.rb @@ -336,6 +336,24 @@ end end + describe '#release_after_fork' do + it 'abandons every connection without closing any' do + allow(driver_dialect).to receive(:abandon_connection) + conn1 = instance_double('Connection') + conn2 = instance_double('Connection') + conn3 = instance_double('Connection') + monitor.instance_variable_get(:@monitoring_connection).set(conn1, close_old: false) + monitor.instance_variable_get(:@instance_monitors_writer_conn).set(conn2, close_old: false) + monitor.instance_variable_get(:@instance_monitor_connections)['reader-1'] = conn3 + + monitor.release_after_fork + + [conn1, conn2, conn3].each { |conn| expect(driver_dialect).to have_received(:abandon_connection).with(conn) } + expect(driver_dialect).not_to have_received(:close_connection) + expect(monitor.state).to eq(:stopped) + end + end + describe 'stable reader topologies' do it 'accepts topology when all readers agree for the required duration' do # Simulate reader topologies being stored diff --git a/spec/unit/monitoring/monitor_connection_spec.rb b/spec/unit/monitoring/monitor_connection_spec.rb index 01d4118c..477cc016 100644 --- a/spec/unit/monitoring/monitor_connection_spec.rb +++ b/spec/unit/monitoring/monitor_connection_spec.rb @@ -123,4 +123,14 @@ expect { monitor_connection.close }.not_to raise_error end end + + describe '#abandon' do + it 'abandons the connection through the driver dialect without closing it' do + allow(mock_driver_dialect).to receive(:abandon_connection) + monitor_connection.set(conn1) + monitor_connection.abandon + expect(mock_driver_dialect).to have_received(:abandon_connection).with(conn1) + expect(mock_driver_dialect).not_to have_received(:close_connection) + end + end end From 5d02b0e437d14c0ebe06396e5b35d7ba65489ce1 Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Tue, 6 Oct 2026 00:17:10 -0700 Subject: [PATCH 06/12] test: wait for topology after connecting in the ActiveRecord reader failover spec The spec warmed the topology cache before establish_connection, but clear_all_connections! now stops the monitors and clears that cache, so the warm-up was lost. If the new monitor's first connection raced with the reader outage, it never learned the writer and failover timed out (seen on MySQL in CI). Wait for the connection's own monitor to cache every instance before cutting connectivity instead. --- .../container/failover_activerecord_spec.rb | 32 +++++++------------ 1 file changed, 11 insertions(+), 21 deletions(-) diff --git a/spec/integration/container/failover_activerecord_spec.rb b/spec/integration/container/failover_activerecord_spec.rb index f61daa60..e41e371f 100644 --- a/spec/integration/container/failover_activerecord_spec.rb +++ b/spec/integration/container/failover_activerecord_spec.rb @@ -26,6 +26,7 @@ require_relative 'utils/database_engine' require_relative 'utils/rds_test_utility' require_relative 'utils/retry_helper' +require_relative 'utils/topology_helper' require 'aws_advanced_ruby_driver_wrapper' require 'aws_advanced_ruby_driver_wrapper/active_record/aws_mysql2_adapter' require 'aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter' @@ -100,23 +101,12 @@ def establish_failover_connection(host:, port:, variables: nil, prepared_stateme ) end - def warm_failover_topology - config = Integration::DriverHelper.native_config( - drv, - host: proxy_info.cluster_endpoint, - port: proxy_info.cluster_endpoint_port, - user: proxy_info.username, - password: proxy_info.password, - dbname: proxy_info.default_dbname - ) - props = { - AwsAdvancedRubyDriverWrapper::PropertyDefinition::PLUGINS.name => 'failover', - AwsAdvancedRubyDriverWrapper::PropertyDefinition::CLUSTER_INSTANCE_HOST_PATTERN.name => - "?.#{proxy_info.instance_endpoint_suffix}:#{proxy_info.instance_endpoint_port}" - } - discovered = Integration::TopologyHelper.warm_topology_cache( - drv: drv, config: config, props: props, min_instances: proxy_info.instances.size - ) + # Blocks until the topology monitor started by the ActiveRecord connection has cached every instance. + # This must run after the connection is established: establish_failover_connection calls + # clear_all_connections!, which stops the monitors and clears the topology cache, so warming the cache + # beforehand would be undone. + def wait_for_failover_topology + discovered = Integration::TopologyHelper.wait_for_topology(min_instances: proxy_info.instances.size) expect(discovered).to be(true), 'Topology was not discovered before failover' end @@ -314,10 +304,6 @@ def instance_id_via(conn) features: [Integration::TestEnvironmentFeatures::NETWORK_OUTAGES_ENABLED] do enable_on_num_instances(min_instances: 2, max_instances: 2) - # Warm the topology cache first: this test connects to a bare reader instance, so without a cached - # writer host, writer failover would have only the dead reader to probe once connectivity is cut. - warm_failover_topology - probe = session_probe reader_instance = proxy_info.instances[1] establish_failover_connection( @@ -330,6 +316,10 @@ def instance_id_via(conn) # Session variables are applied via SET SESSION, so they set fine on a reader connection too. expect(adapter.select_value(probe[:read_sql])).to eq(probe[:expected]) + # This test connects to a bare reader instance, so without a cached writer host, writer failover + # would have only the dead reader to probe once connectivity is cut. + wait_for_failover_topology + Integration::ProxyHelper.disable_connectivity(reader_instance.instance_id) expect { current_instance_id }.to raise_error( From 42baeb0065abf8145e3e7477583ebd2323cd6d75 Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Tue, 6 Oct 2026 09:03:31 -0700 Subject: [PATCH 07/12] PR suggestion --- CHANGELOG.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 35b7660a..2acc222b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -10,7 +10,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), - Fixed the `aws_mysql2` ActiveRecord adapter dropping the `flags`, `encoding`, `socket`, and `reconnect` client options, which removed ActiveRecord's `FOUND_ROWS` flag so that statements such as `update_all` reported changed rows instead of matched rows ([PR #193](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/193)). - Fixed the `aws_postgresql` and `aws_mysql2` ActiveRecord adapters not translating connection errors the way the standard adapters do, so a missing database was not reported as `ActiveRecord::NoDatabaseError` and `db:prepare` (and `bin/setup`) failed instead of creating it ([PR #193](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/193)). - Fixed the `aws_postgresql` and `aws_mysql2` ActiveRecord adapters keeping ActiveRecord's prepared statement cache after a successful failover, so every query prepared before the failover failed on that connection with `Method invoked against old connection` until the process restarted ([PR #194](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/194), [documentation](https://aws.github.io/aws-advanced-wrapper-docs/ruby/enhanced-failover)). -- Fixed `db:drop`, `db:reset`, `db:test:prepare`, and `bin/rails test` after a schema change failing on Aurora PostgreSQL with `database "" is being accessed by other users`, because the topology monitor kept its own connections to the database being dropped ([PR #195](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/195)). +- Fixed `db:drop`, `db:reset`, and `db:test:prepare` failing on Aurora PostgreSQL with `database "" is being accessed by other users`. The same error occurred when running `bin/rails test` after a schema change. The cause was the wrapper's topology and Blue/Green monitors, which kept their own connections open to the database being dropped. The wrapper now stops these monitors before dropping a database and when ActiveRecord clears all its connections ([PR #195](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/195)). ## [1.0.0] - 2026-10-05 From b3c9e684d6117af7ec78daab425e35d86107596c Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Tue, 6 Oct 2026 09:42:03 -0700 Subject: [PATCH 08/12] chore: trigger CI From c50bf0bcb9a8ebc982bdb848e23367db781538e7 Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Tue, 6 Oct 2026 09:50:19 -0700 Subject: [PATCH 09/12] Address PR comments --- lib/aws_advanced_ruby_driver_wrapper.rb | 11 ++++-- .../plugins/blue_green/blue_green_plugin.rb | 7 +++- .../services/monitor_service.rb | 9 +++-- .../utils/events/batching_event_publisher.rb | 8 ++--- spec/unit/fork_safety_spec.rb | 34 +++++++++++++++++++ 5 files changed, 60 insertions(+), 9 deletions(-) diff --git a/lib/aws_advanced_ruby_driver_wrapper.rb b/lib/aws_advanced_ruby_driver_wrapper.rb index a54432f6..0b6e0f12 100644 --- a/lib/aws_advanced_ruby_driver_wrapper.rb +++ b/lib/aws_advanced_ruby_driver_wrapper.rb @@ -87,11 +87,18 @@ def self.after_fork end # Process._fork is the single entry point for Kernel#fork, Process.fork and IO.popen('-'), so - # hooking it covers every way an application server or job runner forks a worker. + # hooking it covers every way an application server or job runner forks a worker. An error in the + # child is logged rather than raised: raising here would fail the application's own fork call. module ForkHook def _fork pid = super - AwsAdvancedRubyDriverWrapper.after_fork if pid.zero? + if pid.zero? + begin + AwsAdvancedRubyDriverWrapper.after_fork + rescue StandardError => e + AwsAdvancedRubyDriverWrapper.logger.error("Failed to reset the wrapper's background services after fork: #{e.message}") + end + end pid end end diff --git a/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/blue_green_plugin.rb b/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/blue_green_plugin.rb index 8e1fa4a8..f8e8fe75 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/blue_green_plugin.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/plugins/blue_green/blue_green_plugin.rb @@ -131,8 +131,13 @@ def self.clean_up_providers end # Forgets the providers inherited by a forked child so the child starts its own on first use. + # A provider that fails to release is still forgotten, and the rest are released. def self.release_providers_after_fork - PROVIDERS.each_key { |k| PROVIDERS.delete(k)&.release_after_fork } + PROVIDERS.each_key do |key| + PROVIDERS.delete(key)&.release_after_fork + rescue StandardError => e + AwsAdvancedRubyDriverWrapper.logger.warn("Failed to release Blue/Green status provider #{key} after fork: #{e.message}") + end end private diff --git a/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb b/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb index 9c29a504..7b043cf4 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb @@ -124,10 +124,15 @@ def shutdown(grace_period:) # Resets this service in a forked child. Monitor threads do not survive a fork, so the inherited # monitors are detached (keeping the parent's connections open) and forgotten, letting - # run_if_absent start live ones. The cleanup thread is restarted. + # run_if_absent start live ones. A monitor that fails to detach is still forgotten. The cleanup + # thread is restarted. def restart_after_fork @lock.synchronize { @caches.values }.each do |container| - container.cache.entries.each_key { |key| container.cache.remove(key)&.release_after_fork } + container.cache.entries.each_key do |key| + container.cache.remove(key)&.release_after_fork + rescue StandardError => e + logger.warn("Failed to release monitor #{key} after fork: #{e.message}") + end end @running = true @cleanup_thread = start_cleanup_thread diff --git a/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb b/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb index 41137857..874dba8a 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb @@ -23,10 +23,10 @@ module Events # Batches deduplicate events via Set semantics (eql?/hash). # # Public API: - # subscribe(subscriber, event_classes) — register for event types - # unsubscribe(subscriber, event_classes) — deregister - # publish(event) — deliver immediate or queue batched - # release_resources — stop background thread + # subscribe(subscriber, event_classes) - register for event types + # unsubscribe(subscriber, event_classes) - deregister + # publish(event) - deliver immediate or queue batched + # release_resources - stop background thread # restart_after_fork - restart background thread in a forked child class BatchingEventPublisher include Logging diff --git a/spec/unit/fork_safety_spec.rb b/spec/unit/fork_safety_spec.rb index da8c2ead..95f84a9a 100644 --- a/spec/unit/fork_safety_spec.rb +++ b/spec/unit/fork_safety_spec.rb @@ -136,4 +136,38 @@ def release_after_fork = events << :release_after_fork expect(result).to eq(remaining: [], events: [:release_after_fork]) expect(providers['bgd-1']).to equal(provider) end + + it 'keeps releasing and restarting in the child when releasing one inherited object raises' do + failing = Struct.new(:name) do + def release_after_fork = raise("release failed for #{name}") + def stop; end + end + provider_class = Struct.new(:events) { def release_after_fork = events << :release_after_fork } + good_provider = provider_class.new([]) + providers['bgd-bad'] = failing.new('provider') + providers['bgd-good'] = good_provider + core.monitor_service.register_type(:fork_test_monitor, expiration_timeout_sec: 60) + core.monitor_service.run_if_absent(:fork_test_monitor, 'bad', nil) { |_| failing.new('monitor').tap { |m| def m.start; end } } + good_monitor = start_test_monitor + allow(described_class.logger).to receive(:warn) + + # A failure here would otherwise escape from Process._fork, so the child's own fork block would never run. + result = in_forked_child do + threads = [core.monitor_service.instance_variable_get(:@cleanup_thread), core.event_publisher.instance_variable_get(:@thread), + core.storage_service.instance_variable_get(:@cleanup_thread)] + monitors_left = %w[bad k1].filter_map { |key| core.monitor_service.get(:fork_test_monitor, key) } + { block_ran: true, providers_left: providers.keys, good_provider_events: good_provider.events, + good_monitor_abandoned: good_monitor.abandoned, monitors_left: monitors_left.size, threads_alive: threads.map(&:alive?) } + end + + expect(result).to eq(block_ran: true, providers_left: [], good_provider_events: [:release_after_fork], + good_monitor_abandoned: true, monitors_left: 0, threads_alive: [true, true, true]) + end + + it "logs instead of failing the application's fork when resetting the wrapper raises" do + allow(described_class).to receive(:after_fork).and_raise(ThreadError, "can't create Thread") + allow(described_class.logger).to receive(:error) + + expect(in_forked_child { :block_ran }).to eq(:block_ran) + end end From 366636fa938593899634955e5a5dfd3ece776c9f Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Tue, 6 Oct 2026 09:57:40 -0700 Subject: [PATCH 10/12] refactor: use an atomic flag for the running state of the shared background services MonitorService and StorageService set and read @running without a lock, while BatchingEventPublisher took its lock for every access. All three now use a Concurrent::AtomicBoolean, so the flag is thread-safe without a lock. The publisher's lock still guards its subscribers and event queue. --- .../services/monitor_service.rb | 9 +++++---- .../utils/events/batching_event_publisher.rb | 9 +++++---- .../utils/storage/storage_service.rb | 9 +++++---- 3 files changed, 15 insertions(+), 12 deletions(-) diff --git a/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb b/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb index 7b043cf4..f5bcef4a 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb @@ -14,6 +14,7 @@ # See the License for the specific language governing permissions and # limitations under the License. +require 'concurrent' require_relative '../monitoring/monitor_state' require_relative '../logging' require_relative '../utils/storage/sliding_expiration_cache' @@ -36,7 +37,7 @@ class MonitorService def initialize(event_publisher:) @caches = {} @lock = Mutex.new - @running = true + @running = Concurrent::AtomicBoolean.new(true) @cleanup_thread = start_cleanup_thread event_publisher.subscribe( self, @@ -112,7 +113,7 @@ def stop_and_remove_all end def shutdown(grace_period:) - @running = false + @running.make_false begin @cleanup_thread&.wakeup rescue ThreadError @@ -134,7 +135,7 @@ def restart_after_fork logger.warn("Failed to release monitor #{key} after fork: #{e.message}") end end - @running = true + @running.make_true @cleanup_thread = start_cleanup_thread end @@ -158,7 +159,7 @@ def handle_data_access_event(event) def start_cleanup_thread thread = Thread.new do - while @running + while @running.true? sleep(CLEANUP_INTERVAL_SEC) run_cleanup end diff --git a/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb b/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb index 874dba8a..38235123 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb @@ -14,6 +14,7 @@ # See the License for the specific language governing permissions and # limitations under the License. +require 'concurrent' require_relative '../../logging' module AwsAdvancedRubyDriverWrapper @@ -38,7 +39,7 @@ def initialize(message_interval_sec: DEFAULT_MESSAGE_INTERVAL_SEC) @subscribers = {} @event_queue = Set.new @lock = Mutex.new - @running = true + @running = Concurrent::AtomicBoolean.new(true) @thread = start_publishing_thread end @@ -71,7 +72,7 @@ def publish(event) end def release_resources - @lock.synchronize { @running = false } + @running.make_false begin @thread&.wakeup rescue ThreadError @@ -82,7 +83,7 @@ def release_resources # Restarts the publishing thread in a forked child, where it no longer runs. def restart_after_fork - @lock.synchronize { @running = true } + @running.make_true @thread = start_publishing_thread end @@ -90,7 +91,7 @@ def restart_after_fork def start_publishing_thread thread = Thread.new do - while @lock.synchronize { @running } + while @running.true? sleep(@message_interval_sec) send_messages end diff --git a/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb b/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb index 74952c3e..92355778 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb @@ -14,6 +14,7 @@ # See the License for the specific language governing permissions and # limitations under the License. +require 'concurrent' require_relative 'expiration_cache' require_relative '../../logging' require_relative '../events/data_access_event' @@ -34,7 +35,7 @@ def initialize(event_publisher:, cleanup_interval: DEFAULT_CLEANUP_INTERVAL) @caches = {} @event_publisher = event_publisher @lock = Mutex.new - @running = true + @running = Concurrent::AtomicBoolean.new(true) @cleanup_interval = cleanup_interval @cleanup_thread = start_cleanup_thread(cleanup_interval) end @@ -134,7 +135,7 @@ def size(name) # Stops the cleanup thread. def shutdown - @running = false + @running.make_false begin @cleanup_thread&.wakeup rescue ThreadError @@ -145,7 +146,7 @@ def shutdown # Restarts the cleanup thread in a forked child, where it no longer runs. def restart_after_fork - @running = true + @running.make_true @cleanup_thread = start_cleanup_thread(@cleanup_interval) end @@ -160,7 +161,7 @@ def fetch_cache!(name) def start_cleanup_thread(interval) thread = Thread.new do - while @running + while @running.true? sleep(interval) remove_expired_items end From 5ae8d0b31cdf04e319c5b908d332929795b51053 Mon Sep 17 00:00:00 2001 From: aaron-congo Date: Tue, 6 Oct 2026 12:00:09 -0700 Subject: [PATCH 11/12] fix: harden fork handling and report Secrets Manager fetch timeouts - Forget Secrets Manager fetches in flight at fork time, so a forked child does not wait on a future whose thread did not survive the fork. - Raise the timeout error when a Secrets Manager fetch does not finish in time. Future#value! returns nil on timeout instead of raising, so the fetch previously failed later with a generic auth error. - Drop the references to abandoned monitor connections, so a later close cannot send COM_QUIT on a socket shared with the parent. - Keep a still-running background thread in restart_after_fork, so calling after_fork again does not start duplicate threads. - Kill a hung forked child in the fork safety spec after a timeout instead of stalling the suite. --- CHANGELOG.md | 3 +- lib/aws_advanced_ruby_driver_wrapper.rb | 6 ++- .../monitoring/cluster_topology_monitor.rb | 5 ++- .../monitoring/monitor_connection.rb | 5 ++- .../plugins/secrets_manager_plugin.rb | 10 ++++- .../services/monitor_service.rb | 4 +- .../utils/events/batching_event_publisher.rb | 5 ++- .../utils/storage/storage_service.rb | 5 ++- spec/unit/fork_safety_spec.rb | 37 ++++++++++++++++++- .../cluster_topology_monitor_spec.rb | 3 ++ .../monitoring/monitor_connection_spec.rb | 11 ++++++ .../plugins/secrets_manager_plugin_spec.rb | 24 ++++++++++++ 12 files changed, 106 insertions(+), 12 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 61ec4505..079687dd 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,7 +11,8 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), - Fixed the `aws_postgresql` and `aws_mysql2` ActiveRecord adapters not translating connection errors the way the standard adapters do, so a missing database was not reported as `ActiveRecord::NoDatabaseError` and `db:prepare` (and `bin/setup`) failed instead of creating it ([PR #193](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/193)). - Fixed the `aws_postgresql` and `aws_mysql2` ActiveRecord adapters keeping ActiveRecord's prepared statement cache after a successful failover, so every query prepared before the failover failed on that connection with `Method invoked against old connection` until the process restarted ([PR #194](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/194), [documentation](https://aws.github.io/aws-advanced-wrapper-docs/ruby/enhanced-failover)). - Fixed `db:drop`, `db:reset`, and `db:test:prepare` failing on Aurora PostgreSQL with `database "" is being accessed by other users`. The same error occurred when running `bin/rails test` after a schema change. The cause was the wrapper's topology and Blue/Green monitors, which kept their own connections open to the database being dropped. The wrapper now stops these monitors before dropping a database and when ActiveRecord clears all its connections ([PR #195](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/195)). -- Fixed forked processes, such as Puma workers with `preload_app!`, Unicorn workers, and Resque jobs, inheriting the parent's background monitors after their threads had stopped, so the cluster topology never refreshed in the child and failover there could not find the new writer. A forked child now starts its own monitors and leaves the parent's monitoring connections open ([PR #196](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/196)). +- Fixed forked processes, such as Puma workers with `preload_app!`, Unicorn workers, and Resque jobs, inheriting the parent's background monitors after their threads had stopped, so the cluster topology never refreshed in the child and failover there could not find the new writer. A forked child now starts its own monitors and leaves the parent's monitoring connections open. A forked child also no longer waits on a Secrets Manager fetch that was in progress in the parent when it forked ([PR #196](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/196)). +- Fixed the AWS Secrets Manager plugin not reporting a timed-out secret fetch as a timeout. A fetch that did not finish within 60 seconds failed with a generic `SecretsManagerAuthError` (`Failed to fetch database credentials from AWS Secrets Manager`) instead of `Timed out fetching secret after 60s` ([PR #196](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/196)). ## [1.0.0] - 2026-10-05 diff --git a/lib/aws_advanced_ruby_driver_wrapper.rb b/lib/aws_advanced_ruby_driver_wrapper.rb index 0b6e0f12..d8ccdf4b 100644 --- a/lib/aws_advanced_ruby_driver_wrapper.rb +++ b/lib/aws_advanced_ruby_driver_wrapper.rb @@ -75,10 +75,12 @@ def self.clear_caches # Runs in a forked child (see ForkHook). Only the forking thread survives a fork, so every background # thread is gone and the inherited monitors look alive but never refresh. The monitors and Blue/Green - # providers are forgotten without closing their connections, which the parent still uses, and the - # shared services restart their threads. Cached data such as topology stays valid and is kept. + # providers are forgotten without closing their connections, which the parent still uses, Secrets + # Manager fetches that were in flight are forgotten, and the shared services restart their threads. + # Cached data such as topology stays valid and is kept. def self.after_fork Plugins::BlueGreen::BlueGreenPlugin.release_providers_after_fork + Plugins::SecretsManagerPlugin.release_pending_refreshes_after_fork if defined?(Plugins::SecretsManagerPlugin) return unless defined?(Services::CoreServices) Services::CoreServices.monitor_service.restart_after_fork diff --git a/lib/aws_advanced_ruby_driver_wrapper/monitoring/cluster_topology_monitor.rb b/lib/aws_advanced_ruby_driver_wrapper/monitoring/cluster_topology_monitor.rb index 32fd0588..74bdfc25 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/monitoring/cluster_topology_monitor.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/monitoring/cluster_topology_monitor.rb @@ -135,7 +135,10 @@ def abandon_connections @monitoring_connection.abandon @instance_monitors_writer_conn.abandon driver_dialect = @service_container.dialect_service.driver_dialect - @instance_monitor_connections.each_value { |conn| driver_dialect.abandon_connection(conn) } + @instance_monitor_connections.each_pair do |host, conn| + driver_dialect.abandon_connection(conn) + @instance_monitor_connections.delete(host) + end end private diff --git a/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor_connection.rb b/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor_connection.rb index 6bc2220e..b6574908 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor_connection.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor_connection.rb @@ -53,9 +53,10 @@ def close set(nil) end - # Releases a connection inherited across a fork without closing it on the server. + # Releases a connection inherited across a fork without closing it on the server. The reference + # is dropped too, so a later close or set cannot reach the parent's session. def abandon - @driver_dialect.abandon_connection(get) + @driver_dialect.abandon_connection(@connection.get_and_set(nil)) end end end diff --git a/lib/aws_advanced_ruby_driver_wrapper/plugins/secrets_manager_plugin.rb b/lib/aws_advanced_ruby_driver_wrapper/plugins/secrets_manager_plugin.rb index bd3ba478..b1698afd 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/plugins/secrets_manager_plugin.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/plugins/secrets_manager_plugin.rb @@ -79,6 +79,12 @@ def clear_cache(storage_service) storage_service.clear(SECRETS_MANAGER_CACHE_NAME) end + # Forgets fetches inherited by a forked child. A fetch running at fork time lost its thread, + # so its future never resolves and would otherwise be returned to every later fetch for its key. + def release_pending_refreshes_after_fork + @pending_refreshes.clear + end + # The secret's value depends on which secret is read, from which region and endpoint, and # with which AWS credentials, so all of them are part of the key. Connections that differ in # any of them neither share a cached secret nor wait on each other's fetch. @@ -244,7 +250,9 @@ def fetch_synchronously(cache_key, credentials) else Concurrent::Promises.future_on(:io) { fetch_and_store_secret(nil, credentials) } end - future.value!(SYNC_FETCH_TIMEOUT_SEC) + # value! returns nil instead of raising when the timeout elapses; a completed fetch always + # returns an entry or raises. + future.value!(SYNC_FETCH_TIMEOUT_SEC) || raise(Timeout::Error) rescue Concurrent::CancelledOperationError, Timeout::Error raise Errors::SecretsManagerAuthError, "Timed out fetching secret after #{SYNC_FETCH_TIMEOUT_SEC}s" diff --git a/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb b/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb index f5bcef4a..c6dbd708 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/services/monitor_service.rb @@ -126,7 +126,7 @@ def shutdown(grace_period:) # Resets this service in a forked child. Monitor threads do not survive a fork, so the inherited # monitors are detached (keeping the parent's connections open) and forgotten, letting # run_if_absent start live ones. A monitor that fails to detach is still forgotten. The cleanup - # thread is restarted. + # thread is restarted unless it is still alive, so calling this again does not start a second one. def restart_after_fork @lock.synchronize { @caches.values }.each do |container| container.cache.entries.each_key do |key| @@ -135,6 +135,8 @@ def restart_after_fork logger.warn("Failed to release monitor #{key} after fork: #{e.message}") end end + return if @cleanup_thread&.alive? + @running.make_true @cleanup_thread = start_cleanup_thread end diff --git a/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb b/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb index 38235123..a45a2341 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/utils/events/batching_event_publisher.rb @@ -81,8 +81,11 @@ def release_resources @thread&.join(@message_interval_sec) end - # Restarts the publishing thread in a forked child, where it no longer runs. + # Restarts the publishing thread in a forked child, where it no longer runs. A thread that is + # still alive is kept, so calling this again does not start a second one. def restart_after_fork + return if @thread&.alive? + @running.make_true @thread = start_publishing_thread end diff --git a/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb b/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb index 92355778..31c5ddfb 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/utils/storage/storage_service.rb @@ -144,8 +144,11 @@ def shutdown @cleanup_thread&.join(5) end - # Restarts the cleanup thread in a forked child, where it no longer runs. + # Restarts the cleanup thread in a forked child, where it no longer runs. A thread that is + # still alive is kept, so calling this again does not start a second one. def restart_after_fork + return if @cleanup_thread&.alive? + @running.make_true @cleanup_thread = start_cleanup_thread(@cleanup_interval) end diff --git a/spec/unit/fork_safety_spec.rb b/spec/unit/fork_safety_spec.rb index 95f84a9a..30d1bdc4 100644 --- a/spec/unit/fork_safety_spec.rb +++ b/spec/unit/fork_safety_spec.rb @@ -18,10 +18,12 @@ require 'aws_advanced_ruby_driver_wrapper' require 'aws_advanced_ruby_driver_wrapper/monitoring/monitor' require 'aws_advanced_ruby_driver_wrapper/services/service_utility' +require 'aws_advanced_ruby_driver_wrapper/plugins/secrets_manager_plugin' RSpec.describe AwsAdvancedRubyDriverWrapper, 'after fork' do let(:core) { AwsAdvancedRubyDriverWrapper::Services::CoreServices } let(:providers) { AwsAdvancedRubyDriverWrapper::Plugins::BlueGreen::BlueGreenPlugin::PROVIDERS } + let(:pending_secret_fetches) { AwsAdvancedRubyDriverWrapper::Plugins::SecretsManagerPlugin.pending_refreshes } let(:test_monitor_class) do Class.new(AwsAdvancedRubyDriverWrapper::Monitoring::Monitor) do attr_reader :closed, :abandoned @@ -42,11 +44,13 @@ def abandon_connections after do providers.clear + pending_secret_fetches.clear core.reset! end - # Runs the block in a forked child and returns its (marshalled) result to the parent. - def in_forked_child + # Runs the block in a forked child and returns its (marshalled) result to the parent. A child that + # hangs is killed after timeout_sec so the suite fails instead of stalling. + def in_forked_child(timeout_sec: 10) reader, writer = IO.pipe pid = fork do reader.close @@ -61,9 +65,18 @@ def in_forked_child exit!(0) end writer.close + unless reader.wait_readable(timeout_sec) + Process.kill(:KILL, pid) + Process.wait(pid) + raise "forked child did not finish within #{timeout_sec}s" + end result = reader.read Process.wait(pid) + raise 'forked child exited without returning a result' if result.empty? + Marshal.load(result) # rubocop:disable Security/MarshalLoad + ensure + reader&.close end def start_test_monitor @@ -164,6 +177,26 @@ def stop; end good_monitor_abandoned: true, monitors_left: 0, threads_alive: [true, true, true]) end + it 'forgets Secrets Manager fetches that were in flight at fork time' do + in_flight = Concurrent::Promises.resolvable_future + pending_secret_fetches['secret-key'] = in_flight + + result = in_forked_child { pending_secret_fetches.keys } + + expect(result).to eq([]) + expect(pending_secret_fetches['secret-key']).to equal(in_flight) + end + + it 'does not start extra background threads when the reset runs again in the same process' do + result = in_forked_child do + before = Thread.list.count(&:alive?) + described_class.after_fork + Thread.list.count(&:alive?) - before + end + + expect(result).to eq(0) + end + it "logs instead of failing the application's fork when resetting the wrapper raises" do allow(described_class).to receive(:after_fork).and_raise(ThreadError, "can't create Thread") allow(described_class.logger).to receive(:error) diff --git a/spec/unit/monitoring/cluster_topology_monitor_spec.rb b/spec/unit/monitoring/cluster_topology_monitor_spec.rb index b984f85c..2ce58ff3 100644 --- a/spec/unit/monitoring/cluster_topology_monitor_spec.rb +++ b/spec/unit/monitoring/cluster_topology_monitor_spec.rb @@ -351,6 +351,9 @@ [conn1, conn2, conn3].each { |conn| expect(driver_dialect).to have_received(:abandon_connection).with(conn) } expect(driver_dialect).not_to have_received(:close_connection) expect(monitor.state).to eq(:stopped) + expect(monitor.instance_variable_get(:@monitoring_connection).get).to be_nil + expect(monitor.instance_variable_get(:@instance_monitors_writer_conn).get).to be_nil + expect(monitor.instance_variable_get(:@instance_monitor_connections)).to be_empty end end diff --git a/spec/unit/monitoring/monitor_connection_spec.rb b/spec/unit/monitoring/monitor_connection_spec.rb index 477cc016..e6e037f2 100644 --- a/spec/unit/monitoring/monitor_connection_spec.rb +++ b/spec/unit/monitoring/monitor_connection_spec.rb @@ -132,5 +132,16 @@ expect(mock_driver_dialect).to have_received(:abandon_connection).with(conn1) expect(mock_driver_dialect).not_to have_received(:close_connection) end + + it 'drops the reference so a later close does not reach the abandoned connection' do + allow(mock_driver_dialect).to receive(:abandon_connection) + allow(mock_driver_dialect).to receive(:close_connection) + monitor_connection.set(conn1) + monitor_connection.abandon + monitor_connection.close + + expect(monitor_connection.get).to be_nil + expect(mock_driver_dialect).not_to have_received(:close_connection) + end end end diff --git a/spec/unit/plugins/secrets_manager_plugin_spec.rb b/spec/unit/plugins/secrets_manager_plugin_spec.rb index 97ad0b2f..ec1d904c 100644 --- a/spec/unit/plugins/secrets_manager_plugin_spec.rb +++ b/spec/unit/plugins/secrets_manager_plugin_spec.rb @@ -445,6 +445,30 @@ def keys_read_by(plugin) end end + describe 'fetch timeout' do + it 'raises an auth error instead of returning no secret when the fetch does not finish in time' do + stub_const("#{described_class}::SYNC_FETCH_TIMEOUT_SEC", 0.05) + allow(mock_sm_client).to receive(:get_secret_value) do + sleep(0.5) + secret_response + end + + expect { build_plugin.connect(host_info, Concurrent::Map.new, true, -> {}) } + .to raise_error(AwsAdvancedRubyDriverWrapper::Errors::SecretsManagerAuthError, /Timed out fetching secret/) + end + end + + describe '.release_pending_refreshes_after_fork' do + it 'forgets in-flight fetches so a later fetch for the same key starts a new one' do + stuck = Concurrent::Promises.resolvable_future + described_class.pending_refreshes['some-key'] = stuck + + described_class.release_pending_refreshes_after_fork + + expect(described_class.pending_refreshes).to be_empty + end + end + describe 'thundering herd protection' do it 'deduplicates concurrent fetches for the same key' do call_count = Concurrent::AtomicFixnum.new(0) From 8a3309231016478861046185373fb8eb4c76bcf3 Mon Sep 17 00:00:00 2001 From: Aaron <69273634+aaron-congo@users.noreply.github.com> Date: Tue, 6 Oct 2026 13:46:53 -0700 Subject: [PATCH 12/12] fix: let Action Cable's PostgreSQL subscription adapter run on aws_postgresql (#197) Action Cable's PostgreSQL subscription adapter refuses any raw connection that is not a PG::Connection, and aws_postgresql hands out a WrapperPgConnection, so every subscribe and broadcast (including Turbo Stream broadcasts) raised "The Active Record database must be PostgreSQL in order to use the PostgreSQL Action Cable storage adapter". The aws_postgresql adapter now registers an :action_cable load hook that prepends a verify! override to ActionCable::SubscriptionAdapter::PostgreSQL, accepting a WrapperPgConnection and leaving every other connection to Action Cable's own check. The calls Action Cable makes on the connection (exec, escape_identifier, escape_string, wait_for_notify) are forwarded to the wrapped PG::Connection. --- CHANGELOG.md | 1 + .../aws_action_cable_postgresql_support.rb | 40 +++++++++++++++++++ .../active_record/aws_postgresql_adapter.rb | 1 + .../subscription_adapter/postgresql.rb | 35 ++++++++++++++++ ...ws_action_cable_postgresql_support_spec.rb | 36 +++++++++++++++++ 5 files changed, 113 insertions(+) create mode 100644 lib/aws_advanced_ruby_driver_wrapper/active_record/aws_action_cable_postgresql_support.rb create mode 100644 spec/fixtures/action_cable_stub/action_cable/subscription_adapter/postgresql.rb create mode 100644 spec/unit/aws_action_cable_postgresql_support_spec.rb diff --git a/CHANGELOG.md b/CHANGELOG.md index 079687dd..8566978a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,7 @@ The format is based on [Keep a Changelog](https://keepachangelog.com/en/1.0.0/), - Fixed `db:drop`, `db:reset`, and `db:test:prepare` failing on Aurora PostgreSQL with `database "" is being accessed by other users`. The same error occurred when running `bin/rails test` after a schema change. The cause was the wrapper's topology and Blue/Green monitors, which kept their own connections open to the database being dropped. The wrapper now stops these monitors before dropping a database and when ActiveRecord clears all its connections ([PR #195](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/195)). - Fixed forked processes, such as Puma workers with `preload_app!`, Unicorn workers, and Resque jobs, inheriting the parent's background monitors after their threads had stopped, so the cluster topology never refreshed in the child and failover there could not find the new writer. A forked child now starts its own monitors and leaves the parent's monitoring connections open. A forked child also no longer waits on a Secrets Manager fetch that was in progress in the parent when it forked ([PR #196](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/196)). - Fixed the AWS Secrets Manager plugin not reporting a timed-out secret fetch as a timeout. A fetch that did not finish within 60 seconds failed with a generic `SecretsManagerAuthError` (`Failed to fetch database credentials from AWS Secrets Manager`) instead of `Timed out fetching secret after 60s` ([PR #196](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/196)). +- Fixed Action Cable's PostgreSQL subscription adapter (`adapter: postgresql` in `cable.yml`) refusing to run on the `aws_postgresql` ActiveRecord adapter with `The Active Record database must be PostgreSQL in order to use the PostgreSQL Action Cable storage adapter`, which also stopped Turbo Stream broadcasts ([PR #197](https://github.com/aws/aws-advanced-ruby-driver-wrapper/pull/197)). ## [1.0.0] - 2026-10-05 diff --git a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_action_cable_postgresql_support.rb b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_action_cable_postgresql_support.rb new file mode 100644 index 00000000..903c367d --- /dev/null +++ b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_action_cable_postgresql_support.rb @@ -0,0 +1,40 @@ +# frozen_string_literal: true + +# Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +require_relative '../postgresql' + +module ActiveRecord + module ConnectionAdapters + # Lets Action Cable's PostgreSQL subscription adapter (cable.yml `adapter: postgresql`) run on the + # aws_postgresql adapter. Action Cable refuses any raw connection that is not a PG::Connection, and + # aws_postgresql hands out a WrapperPgConnection, which forwards the calls Action Cable makes + # (exec, escape_identifier, escape_string, wait_for_notify) to the PG::Connection it wraps. + module AwsActionCablePostgreSQLSupport + private + + def verify!(pg_conn) + super unless pg_conn.is_a?(AwsAdvancedRubyDriverWrapper::WrapperPgConnection) + end + end + end +end + +# Action Cable runs this hook when its server loads, before it loads a subscription adapter, so the +# PostgreSQL one is loaded here to patch it. Applications without Action Cable never run the hook. +ActiveSupport.on_load(:action_cable) do + require 'action_cable/subscription_adapter/postgresql' + ActionCable::SubscriptionAdapter::PostgreSQL.prepend(ActiveRecord::ConnectionAdapters::AwsActionCablePostgreSQLSupport) +end diff --git a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb index 4b5078f5..5527b572 100644 --- a/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb +++ b/lib/aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter.rb @@ -18,6 +18,7 @@ require_relative '../postgresql' require_relative '../errors' require_relative 'aws_connection_handler' +require_relative 'aws_action_cable_postgresql_support' module ActiveRecord module ConnectionAdapters diff --git a/spec/fixtures/action_cable_stub/action_cable/subscription_adapter/postgresql.rb b/spec/fixtures/action_cable_stub/action_cable/subscription_adapter/postgresql.rb new file mode 100644 index 00000000..90db0084 --- /dev/null +++ b/spec/fixtures/action_cable_stub/action_cable/subscription_adapter/postgresql.rb @@ -0,0 +1,35 @@ +# frozen_string_literal: true + +# Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +# Stand-in for Action Cable's PostgreSQL subscription adapter, which is not a dependency of this gem. +# It keeps the check every supported Action Cable version (7.2 to 8.1) runs on the connection it is given. +module ActionCable + module SubscriptionAdapter + class PostgreSQL + def check(pg_conn) + verify!(pg_conn) + end + + private + + def verify!(pg_conn) + return if pg_conn.is_a?(PG::Connection) + + raise 'The Active Record database must be PostgreSQL in order to use the PostgreSQL Action Cable storage adapter' + end + end + end +end diff --git a/spec/unit/aws_action_cable_postgresql_support_spec.rb b/spec/unit/aws_action_cable_postgresql_support_spec.rb new file mode 100644 index 00000000..95f0c4dd --- /dev/null +++ b/spec/unit/aws_action_cable_postgresql_support_spec.rb @@ -0,0 +1,36 @@ +# frozen_string_literal: true + +# Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. +# +# Licensed under the Apache License, Version 2.0 (the "License"). +# You may not use this file except in compliance with the License. +# You may obtain a copy of the License at +# +# http://www.apache.org/licenses/LICENSE-2.0 +# +# Unless required by applicable law or agreed to in writing, software +# distributed under the License is distributed on an "AS IS" BASIS, +# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +# See the License for the specific language governing permissions and +# limitations under the License. + +require_relative '../spec_helper' + +# Action Cable is not a dependency of this gem, so its PostgreSQL subscription adapter is replaced by a stand-in. +$LOAD_PATH.unshift(File.expand_path('../fixtures/action_cable_stub', __dir__)) +require 'aws_advanced_ruby_driver_wrapper/active_record/aws_postgresql_adapter' + +RSpec.describe 'Action Cable PostgreSQL subscription adapter with aws_postgresql' do + # Action Cable runs its load hooks when its server loads. + before(:all) { ActiveSupport.run_load_hooks(:action_cable, Object.new) } + + let(:adapter) { ActionCable::SubscriptionAdapter::PostgreSQL.new } + + it 'accepts the wrapper connection that aws_postgresql hands out as its raw connection' do + expect { adapter.check(AwsAdvancedRubyDriverWrapper::WrapperPgConnection.allocate) }.not_to raise_error + end + + it 'still rejects connections that are not PostgreSQL' do + expect { adapter.check(Object.new) }.to raise_error(RuntimeError, /must be PostgreSQL/) + end +end