Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
16 commits
Select commit Hold shift + click to select a range
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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 "<name>" 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

Expand Down
34 changes: 34 additions & 0 deletions lib/aws_advanced_ruby_driver_wrapper.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
@@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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)

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down
14 changes: 14 additions & 0 deletions lib/aws_advanced_ruby_driver_wrapper/monitoring/monitor.rb
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Original file line number Diff line number Diff line change
Expand Up @@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -338,6 +338,10 @@ def notify_changes
@event.reset
end

def abandon_connections
@connection.abandon
end

private

def attempt_open_connection
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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?

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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"
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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'
Expand All @@ -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,
Expand Down Expand Up @@ -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
Expand All @@ -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)
Expand All @@ -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
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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

Expand All @@ -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

Expand Down Expand Up @@ -70,7 +72,7 @@ def publish(event)
end

def release_resources
@lock.synchronize { @running = false }
@running.make_false
begin
@thread&.wakeup
rescue ThreadError
Expand All @@ -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
Expand Down
Loading
Loading