diff --git a/CHANGELOG.md b/CHANGELOG.md index 2acc222b..8566978a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -11,6 +11,9 @@ 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. 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.rb b/lib/aws_advanced_ruby_driver_wrapper.rb index 9ae43623..d8ccdf4b 100644 --- a/lib/aws_advanced_ruby_driver_wrapper.rb +++ b/lib/aws_advanced_ruby_driver_wrapper.rb @@ -73,6 +73,38 @@ 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, 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 + 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. 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 + 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 + 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 +120,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/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/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..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 @@ -131,6 +131,16 @@ 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_pair do |host, conn| + driver_dialect.abandon_connection(conn) + @instance_monitor_connections.delete(host) + end + 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..b6574908 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,12 @@ 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. The reference + # is dropped too, so a later close or set cannot reach the parent's session. + def abandon + @driver_dialect.abandon_connection(@connection.get_and_set(nil)) + 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..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 @@ -130,6 +130,16 @@ 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. + # A provider that fails to release is still forgotten, and the rest are released. + def self.release_providers_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 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/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 50110163..c6dbd708 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 @@ -122,6 +123,24 @@ 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. A monitor that fails to detach is still forgotten. The cleanup + # 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| + container.cache.remove(key)&.release_after_fork + rescue StandardError => e + 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 + # Processes events from the event publisher. # @param event [Event] the event to process. def process_event(event) @@ -142,7 +161,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 3eb6e9b4..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 @@ -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 @@ -23,10 +24,11 @@ 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 @@ -37,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 @@ -70,7 +72,7 @@ def publish(event) end def release_resources - @lock.synchronize { @running = false } + @running.make_false begin @thread&.wakeup rescue ThreadError @@ -79,11 +81,20 @@ def release_resources @thread&.join(@message_interval_sec) end + # 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 + private 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 40ed4ef5..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 @@ -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,8 @@ 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 @@ -133,7 +135,7 @@ def size(name) # Stops the cleanup thread. def shutdown - @running = false + @running.make_false begin @cleanup_thread&.wakeup rescue ThreadError @@ -142,6 +144,15 @@ def shutdown @cleanup_thread&.join(5) end + # 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 + private def fetch_cache!(name) @@ -153,7 +164,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 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 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..30d1bdc4 --- /dev/null +++ b/spec/unit/fork_safety_spec.rb @@ -0,0 +1,206 @@ +# 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' +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 + + 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 + pending_secret_fetches.clear + core.reset! + end + + # 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 + 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 + 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 + 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 + + 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 '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) + + expect(in_forked_child { :block_ran }).to eq(:block_ran) + end +end diff --git a/spec/unit/monitoring/cluster_topology_monitor_spec.rb b/spec/unit/monitoring/cluster_topology_monitor_spec.rb index 32a5a1f1..2ce58ff3 100644 --- a/spec/unit/monitoring/cluster_topology_monitor_spec.rb +++ b/spec/unit/monitoring/cluster_topology_monitor_spec.rb @@ -336,6 +336,27 @@ 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) + 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 + 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..e6e037f2 100644 --- a/spec/unit/monitoring/monitor_connection_spec.rb +++ b/spec/unit/monitoring/monitor_connection_spec.rb @@ -123,4 +123,25 @@ 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 + + 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)