From e376ce6f2e63bbd044c01e5f230cc97f1cafcaf9 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pawe=C5=82=20Pacana?= Date: Tue, 28 Apr 2026 12:27:20 +0200 Subject: [PATCH] Harmonize with remaining adapters - Let applications pass an existing pg connection or pool instead of having PgAdapter create and configure its own pool. - Keep single-connection use serialized while delegating pooled use to the pool checkout mechanism. --- README.md | 9 ++- lib/en57/pg_adapter.rb | 16 +++-- test/test_helper.rb | 5 +- test/test_pg_adapter.rb | 119 +++++++++------------------------ test/test_repository.rb | 145 +++++++++++++++++----------------------- 5 files changed, 110 insertions(+), 184 deletions(-) diff --git a/README.md b/README.md index 08623b6..dbbbd7c 100644 --- a/README.md +++ b/README.md @@ -6,16 +6,15 @@ DCB-compatible event store library in Ruby with support for PostgreSQL. ### Connect with raw pg -Use `PgAdapter` when En57 should own its PostgreSQL connection. +Use `PgAdapter` when your app already owns a pg connection or connection pool. ```ruby +pool = ConnectionPool.new(size: 8) { PG.connect("postgres://localhost:5432/en57") } + store = En57::EventStore.new( En57::Repository.new( - En57::PgAdapter.new( - "postgres://localhost:5432/en57", - max_connections: 8, - ), + En57::PgAdapter.new(pool), En57::JsonSerializer.new, ), ) diff --git a/lib/en57/pg_adapter.rb b/lib/en57/pg_adapter.rb index 529b512..5549680 100644 --- a/lib/en57/pg_adapter.rb +++ b/lib/en57/pg_adapter.rb @@ -1,17 +1,19 @@ # frozen_string_literal: true -require "connection_pool" -require "pg" - module En57 class PgAdapter - def initialize(connection_uri, max_connections: 5) - @connection_pool = - ConnectionPool.new(size: max_connections) { PG.connect(connection_uri) } + def initialize(connection_or_pool) + @with_connection = + if connection_or_pool.respond_to?(:with) + connection_or_pool.public_method(:with) + else + mutex = Mutex.new + ->(&block) { mutex.synchronize { block.call(connection_or_pool) } } + end end def with_connection = - @connection_pool.with { |connection| yield connection } + @with_connection.call { |connection| yield connection } def with_serializable_transaction with_connection do |connection| diff --git a/test/test_helper.rb b/test/test_helper.rb index 55fae94..4507021 100644 --- a/test/test_helper.rb +++ b/test/test_helper.rb @@ -9,6 +9,7 @@ require "active_record" require "en57" require "securerandom" require "concurrent-ruby" +require "connection_pool" require "pg_ephemeral" module En57 @@ -16,6 +17,7 @@ module En57 POOL_SIZE = 8 SERVER = PgEphemeral.start CONNECTION = PG.connect(SERVER.url) + PG_POOL = ConnectionPool.new(size: POOL_SIZE) { PG.connect(SERVER.url) } SEQUEL_DB = Sequel.connect( SERVER.url, @@ -25,7 +27,7 @@ module En57 ActiveRecord::Base.establish_connection("#{SERVER.url}&pool=#{POOL_SIZE}") AR_POOL = ActiveRecord::Base.connection_pool ADAPTERS = { - pg: -> { PgAdapter.new(SERVER.url, max_connections: POOL_SIZE) }, + pg: -> { PgAdapter.new(PG_POOL) }, sequel: -> { SequelAdapter.new(SEQUEL_DB) }, active_record: -> { ActiveRecordAdapter.new(AR_POOL) }, } @@ -36,6 +38,7 @@ module En57 Minitest.after_run do ActiveRecord::Base.connection_pool.disconnect! SEQUEL_DB.disconnect + PG_POOL.shutdown(&:close) CONNECTION.close SERVER.shutdown end diff --git a/test/test_pg_adapter.rb b/test/test_pg_adapter.rb index 57485aa..b57d22e 100644 --- a/test/test_pg_adapter.rb +++ b/test/test_pg_adapter.rb @@ -6,73 +6,48 @@ module En57 class TestPgAdapter < Minitest::Test cover PgAdapter - def test_with_connection_connects_once_and_yields_connection - with_mock_adapter do |connection, adapter, connection_count| - connection.expect(:exec, :first, ["SELECT 1"]) - connection.expect(:exec, :second, ["SELECT 2"]) + def test_with_connection_yields_connection + with_mock_adapter do |connection, adapter| + connection.expect(:exec, :selected, ["SELECT 1"]) - assert_equal :first, + assert_equal :selected, adapter.with_connection { |conn| conn.exec("SELECT 1") } - assert_equal :second, - adapter.with_connection { |conn| conn.exec("SELECT 2") } - assert_equal 1, connection_count.call end end - def test_with_connection_uses_five_connections_by_default - with_connection_count( - 6, - ) do |connections, acquired, release, connection_count| - adapter = PgAdapter.new(connection_uri) - - threads = - 6.times.map do - Thread.new do - adapter.with_connection do |connection| - acquired << connection - release.pop - end - end - end - - assert_equal connections.take(5), 5.times.map { acquired.pop } - assert_equal 5, connection_count.call - assert_raises(ThreadError) { acquired.pop(true) } - - release << true - assert_includes connections.take(5), acquired.pop + def test_with_connection_uses_connection_pool + connection = Object.new + pool = Object.new + pool.define_singleton_method(:with) { |&block| block.call(connection) } + adapter = PgAdapter.new(pool) - 5.times { release << true } - threads.each(&:value) - end + assert_same connection, adapter.with_connection { |conn| conn } end - def test_with_connection_honors_max_connections - with_connection_count( - 3, - ) do |connections, acquired, release, connection_count| - adapter = PgAdapter.new(connection_uri, max_connections: 2) - - threads = - 3.times.map do - Thread.new do - adapter.with_connection do |connection| - acquired << connection - release.pop - end + def test_with_connection_synchronizes_access + connection = Object.new + adapter = PgAdapter.new(connection) + acquired = Queue.new + release = Queue.new + + threads = + 2.times.map do + Thread.new do + adapter.with_connection do |conn| + acquired << conn + release.pop end end + end - assert_equal connections.take(2), 2.times.map { acquired.pop } - assert_equal 2, connection_count.call - assert_raises(ThreadError) { acquired.pop(true) } + assert_same connection, acquired.pop + assert_raises(ThreadError) { acquired.pop(true) } - release << true - assert_includes connections.take(2), acquired.pop + release << true + assert_same connection, acquired.pop - 2.times { release << true } - threads.each(&:value) - end + release << true + threads.each(&:value) end def test_with_serializable_transaction_commits_on_success @@ -120,40 +95,12 @@ module En57 private - def connection_uri = "postgres://localhost:5432/en57_test" - - def with_connection_count(size) - connections = Array.new(size) { Object.new } - acquired = Queue.new - release = Queue.new - connection_count = 0 - - PG.stub( - :connect, - ->(actual_connection_uri) do - connection_count += 1 - assert_equal connection_uri, actual_connection_uri - connections.fetch(connection_count - 1) - end, - ) { yield connections, acquired, release, -> { connection_count } } - end - def with_mock_adapter connection = Minitest::Mock.new - connection_count = 0 - - PG.stub( - :connect, - ->(actual_connection_uri) do - connection_count += 1 - assert_equal connection_uri, actual_connection_uri - connection - end, - ) do - yield connection, PgAdapter.new(connection_uri), -> { connection_count } - ensure - connection.verify - end + + yield connection, PgAdapter.new(connection) + ensure + connection.verify end end end diff --git a/test/test_repository.rb b/test/test_repository.rb index 4233e7f..c527101 100644 --- a/test/test_repository.rb +++ b/test/test_repository.rb @@ -30,7 +30,7 @@ module En57 ), ], ) - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, @@ -42,10 +42,7 @@ module En57 ) connection.expect(:exec, nil, ["COMMIT"]) - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).append( + Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( [ Event.new( id: ids[0], @@ -74,7 +71,7 @@ module En57 array_encoder.encode( [record_encoder.encode([ids[0], "OrderPlaced", nil, nil, "{}"])], ) - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, @@ -86,10 +83,7 @@ module En57 ) connection.expect(:exec, nil, ["COMMIT"]) - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).append( + Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( [Event.new(id: ids[0], type: "OrderPlaced")], fail_if: Query.all, ) @@ -97,7 +91,7 @@ module En57 end def test_append_passes_fail_if_and_after_conditions - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, @@ -112,10 +106,7 @@ module En57 ) connection.expect(:exec, nil, ["COMMIT"]) - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).append( + Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( [], fail_if: Query.new( @@ -132,7 +123,7 @@ module En57 end def test_append_rolls_back_transaction_on_pg_failure - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect(:exec, nil, ["ROLLBACK"]) connection.expect(:exec_params, nil) do |sql, params| @@ -145,31 +136,31 @@ module En57 end assert_raises(PG::Error) do - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).append([], fail_if: Query.all) + Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( + [], + fail_if: Query.all, + ) end end end def test_append_rolls_back_transaction_on_failure - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect(:exec, nil, ["ROLLBACK"]) connection.expect(:exec_params, nil) { raise RuntimeError, "boom" } assert_raises(RuntimeError) do - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).append([], fail_if: Query.all) + Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( + [], + fail_if: Query.all, + ) end end end def test_read_events_with_tags - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect( :exec_params, [ @@ -221,16 +212,15 @@ module En57 2, ], ], - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).read(Query.all), + Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( + Query.all, + ), ) end end def test_read_events_with_metadata_restores_types - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect( :exec_params, [ @@ -262,16 +252,15 @@ module En57 1, ], ], - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).read(Query.all), + Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( + Query.all, + ), ) end end def test_read_events_with_null_data_returns_empty_hash - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect( :exec_params, [ @@ -292,10 +281,9 @@ module En57 assert_equal( [[Event.new(id: ids[0], type: "OrderPlaced", data: {}), 1]], - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).read(Query.all), + Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( + Query.all, + ), ) end end @@ -305,7 +293,7 @@ module En57 Query.new( criteria: [Query::Criteria.new(types: [], tags: ["order_id:123"])], ) - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect( :exec_params, [ @@ -338,17 +326,16 @@ module En57 1, ], ], - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).read(query), + Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( + query, + ), ) end end def test_read_events_with_wildcard_query_item query = Query.new(criteria: [Query::Criteria.new(types: [], tags: [])]) - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect( :exec_params, [ @@ -381,10 +368,9 @@ module En57 1, ], ], - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).read(query), + Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( + query, + ), ) end end @@ -397,7 +383,7 @@ module En57 Query::Criteria.new(types: [], tags: ["order_id:456"]), ], ) - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect( :exec_params, [ @@ -434,10 +420,9 @@ module En57 1, ], ], - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).read(query), + Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( + query, + ), ) end end @@ -447,7 +432,7 @@ module En57 Query.new( criteria: [Query::Criteria.new(types: [], tags: [], after: 42)], ) - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect( :exec_params, [], @@ -459,10 +444,9 @@ module En57 assert_equal( [], - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).read(query), + Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( + query, + ), ) end end @@ -472,7 +456,7 @@ module En57 Query.new( criteria: [Query::Criteria.new(types: ["OrderPlaced"], tags: [])], ) - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect( :exec_params, [ @@ -505,31 +489,30 @@ module En57 1, ], ], - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).read(query), + Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( + query, + ), ) end end def test_append_raises_append_condition_violated_from_pg_error_sqlstate - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect(:exec, nil, ["ROLLBACK"]) connection.expect(:exec_params, nil) { raise(PG::RaiseException.new) } assert_raises(AppendConditionViolated) do - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).append([], fail_if: Query.all) + Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( + [], + fail_if: Query.all, + ) end end end def test_append_raises_append_condition_violated_from_serialization_failure_result_sqlstate - with_connection_to(connection_uri) do |connection| + with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect(:exec, nil, ["ROLLBACK"]) connection.expect(:exec_params, nil) do @@ -537,10 +520,10 @@ module En57 end assert_raises(AppendConditionViolated) do - Repository.new( - PgAdapter.new(connection_uri), - JsonSerializer.new, - ).append([], fail_if: Query.all) + Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( + [], + fail_if: Query.all, + ) end end end @@ -549,18 +532,10 @@ module En57 def ids = @ids ||= Hash.new { |h, k| h[k] = SecureRandom.uuid_v7 } - def connection_uri = "postgres://localhost:5432/en57_test" - - def with_connection_to(connection_uri) + def with_connection connection = Minitest::Mock.new - PG.stub( - :connect, - ->(actual_connection_uri) do - assert_equal(connection_uri, actual_connection_uri) - connection - end, - ) { yield connection } + yield connection connection.verify end -- 2.51.2