From 7ee397f2283dfb6765b563dc6517c8a38c0e2a72 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pawe=C5=82=20Pacana?= Date: Thu, 30 Apr 2026 13:03:19 +0200 Subject: [PATCH] Clarify PgAdapter construction - Split pool and single-connection setup into explicit factories so callers do not depend on implicit duck-typing. - Wrap single PG connections in a pool-like Mono adapter to keep serialized access while simplifying PgAdapter itself. --- lib/en57/pg_adapter.rb | 37 ++++++++++----- test/test_helper.rb | 2 +- test/test_migrator.rb | 5 +- test/test_pg_adapter.rb | 10 ++-- test/test_repository.rb | 103 +++++++++++++++++++++++----------------- 5 files changed, 95 insertions(+), 62 deletions(-) diff --git a/lib/en57/pg_adapter.rb b/lib/en57/pg_adapter.rb index e51efc3..3324cd6 100644 --- a/lib/en57/pg_adapter.rb +++ b/lib/en57/pg_adapter.rb @@ -4,18 +4,15 @@ require "pg" module En57 class PgAdapter - 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 + def self.for_pool(connection_pool) = new(connection_pool) + + def self.for_connection(connection) = new(Mono.new(connection)) + + def initialize(connection_pool) + @connection_pool = connection_pool end - def with_connection = - @with_connection.call { |connection| yield connection } + def with_connection(&block) = @connection_pool.with(&block) def with_serializable_transaction with_connection do |connection| @@ -27,11 +24,27 @@ module En57 raise end end + + class Mono + def initialize(connection) + @connection = connection + @mutex = Mutex.new + end + + def with + @mutex.synchronize { yield @connection } + end + end end class EventStore def self.for_pg(connection_uri) - new(Repository.new(PgAdapter.new(PG.connect(connection_uri)), JsonSerializer.new)) + new( + Repository.new( + PgAdapter.for_connection(PG.connect(connection_uri)), + JsonSerializer.new, + ), + ) end end @@ -40,7 +53,7 @@ module En57 def self.for_pooled_pg(connection_uri, max_connections: 5) new( Repository.new( - PgAdapter.new( + PgAdapter.for_pool( ConnectionPool.new(size: max_connections) do PG.connect(connection_uri) end, diff --git a/test/test_helper.rb b/test/test_helper.rb index 749d78d..f707f91 100644 --- a/test/test_helper.rb +++ b/test/test_helper.rb @@ -39,7 +39,7 @@ module En57 end.call ADAPTERS = { - pg: -> { PgAdapter.new(PG_POOL) }, + pg: -> { PgAdapter.for_pool(PG_POOL) }, sequel: -> { SequelAdapter.new(SEQUEL_DB) }, active_record: -> { ActiveRecordAdapter.new(AR_POOL) }, } diff --git a/test/test_migrator.rb b/test/test_migrator.rb index 5816f69..952057d 100644 --- a/test/test_migrator.rb +++ b/test/test_migrator.rb @@ -49,7 +49,10 @@ module En57 connection = PG.connect(url) event_store = EventStore.new( - Repository.new(PgAdapter.new(connection), JsonSerializer.new), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ), ) assert_equal [event], event_store.append([event]).read.each.to_a diff --git a/test/test_pg_adapter.rb b/test/test_pg_adapter.rb index 944344e..495dc18 100644 --- a/test/test_pg_adapter.rb +++ b/test/test_pg_adapter.rb @@ -15,18 +15,18 @@ module En57 end end - def test_with_connection_uses_connection_pool + def test_for_pool_uses_connection_pool connection = Object.new pool = Object.new pool.define_singleton_method(:with) { |&block| block.call(connection) } - adapter = PgAdapter.new(pool) + adapter = PgAdapter.for_pool(pool) assert_same connection, adapter.with_connection { |conn| conn } end - def test_with_connection_synchronizes_access + def test_for_connection_synchronizes_access connection = Object.new - adapter = PgAdapter.new(connection) + adapter = PgAdapter.for_connection(connection) acquired = Queue.new release = Queue.new @@ -99,7 +99,7 @@ module En57 def with_mock_adapter connection = Minitest::Mock.new - yield connection, PgAdapter.new(connection) + yield connection, PgAdapter.for_connection(connection) ensure connection.verify end diff --git a/test/test_repository.rb b/test/test_repository.rb index 691de31..c50828b 100644 --- a/test/test_repository.rb +++ b/test/test_repository.rb @@ -42,7 +42,10 @@ module En57 ) connection.expect(:exec, nil, ["COMMIT"]) - Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append( [ Event.new( id: ids[0], @@ -83,7 +86,10 @@ module En57 ) connection.expect(:exec, nil, ["COMMIT"]) - Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append( [Event.new(id: ids[0], type: "OrderPlaced")], fail_if: Query.all, ) @@ -106,7 +112,10 @@ module En57 ) connection.expect(:exec, nil, ["COMMIT"]) - Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append( [], fail_if: Query.new( @@ -136,10 +145,10 @@ module En57 end assert_raises(PG::Error) do - Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( - [], - fail_if: Query.all, - ) + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append([], fail_if: Query.all) end end end @@ -151,10 +160,10 @@ module En57 connection.expect(:exec_params, nil) { raise RuntimeError, "boom" } assert_raises(RuntimeError) do - Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( - [], - fail_if: Query.all, - ) + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append([], fail_if: Query.all) end end end @@ -212,9 +221,10 @@ module En57 2, ], ], - Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( - Query.all, - ), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).read(Query.all), ) end end @@ -252,9 +262,10 @@ module En57 1, ], ], - Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( - Query.all, - ), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).read(Query.all), ) end end @@ -281,9 +292,10 @@ module En57 assert_equal( [[Event.new(id: ids[0], type: "OrderPlaced", data: {}), 1]], - Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( - Query.all, - ), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).read(Query.all), ) end end @@ -326,9 +338,10 @@ module En57 1, ], ], - Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( - query, - ), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).read(query), ) end end @@ -368,9 +381,10 @@ module En57 1, ], ], - Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( - query, - ), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).read(query), ) end end @@ -420,9 +434,10 @@ module En57 1, ], ], - Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( - query, - ), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).read(query), ) end end @@ -444,9 +459,10 @@ module En57 assert_equal( [], - Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( - query, - ), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).read(query), ) end end @@ -489,9 +505,10 @@ module En57 1, ], ], - Repository.new(PgAdapter.new(connection), JsonSerializer.new).read( - query, - ), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).read(query), ) end end @@ -503,10 +520,10 @@ module En57 connection.expect(:exec_params, nil) { raise(PG::RaiseException.new) } assert_raises(AppendConditionViolated) do - Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( - [], - fail_if: Query.all, - ) + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append([], fail_if: Query.all) end end end @@ -520,10 +537,10 @@ module En57 end assert_raises(AppendConditionViolated) do - Repository.new(PgAdapter.new(connection), JsonSerializer.new).append( - [], - fail_if: Query.all, - ) + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append([], fail_if: Query.all) end end end -- 2.51.2