diff --git a/Gemfile b/Gemfile index 51d5372..c7ae7fc 100644 --- a/Gemfile +++ b/Gemfile @@ -4,6 +4,7 @@ source "https://rubygems.org" gemspec +gem "activerecord" gem "bigdecimal" gem "irb" gem "rake" diff --git a/Gemfile.lock b/Gemfile.lock index a8f1d64..92136fa 100644 --- a/Gemfile.lock +++ b/Gemfile.lock @@ -7,19 +7,44 @@ PATH GEM remote: https://rubygems.org/ specs: + activemodel (8.1.3) + activesupport (= 8.1.3) + activerecord (8.1.3) + activemodel (= 8.1.3) + activesupport (= 8.1.3) + timeout (>= 0.4.0) + activesupport (8.1.3) + base64 + bigdecimal + concurrent-ruby (~> 1.0, >= 1.3.1) + connection_pool (>= 2.2.5) + drb + i18n (>= 1.6, < 2) + json + logger (>= 1.4.2) + minitest (>= 5.1) + securerandom (>= 0.3) + tzinfo (~> 2.0, >= 2.0.5) + uri (>= 0.13.1) ast (2.4.3) + base64 (0.3.0) bigdecimal (4.1.2) concurrent-ruby (1.3.6) + connection_pool (3.0.2) date (3.5.1) diff-lcs (2.0.0) drb (2.2.3) erb (6.0.4) + i18n (1.14.8) + concurrent-ruby (~> 1.0) io-console (0.8.2) irb (1.18.0) pp (>= 0.6.0) prism (>= 1.3.0) rdoc (>= 4.0.0) reline (>= 0.4.2) + json (2.19.4) + logger (1.7.0) m (1.7.0) rake minitest (6.0.5) @@ -76,11 +101,15 @@ GEM stringio (3.2.0) syntax_tree (6.3.0) prettier_print (>= 1.2.0) + timeout (0.6.1) tsort (0.2.0) + tzinfo (2.0.6) + concurrent-ruby (~> 1.0) unparser (0.9.0) diff-lcs (>= 1.6, < 3) parser (>= 3.3.0) prism (>= 1.5.1) + uri (1.1.1) PLATFORMS aarch64-linux @@ -88,6 +117,7 @@ PLATFORMS x86_64-linux DEPENDENCIES + activerecord bigdecimal concurrent-ruby en57! @@ -102,16 +132,24 @@ DEPENDENCIES syntax_tree CHECKSUMS + activemodel (8.1.3) sha256=90c05cbe4cef3649b8f79f13016191ea94c4525ce4a5c0fb7ef909c4b91c8219 + activerecord (8.1.3) sha256=8003be7b2466ba0a2a670e603eeb0a61dd66058fccecfc49901e775260ac70ab + activesupport (8.1.3) sha256=21a5e0dfbd4c3ddd9e1317ec6a4d782fa226e7867dc70b0743acda81a1dca20e ast (2.4.3) sha256=954615157c1d6a382bc27d690d973195e79db7f55e9765ac7c481c60bdb4d383 + base64 (0.3.0) sha256=27337aeabad6ffae05c265c450490628ef3ebd4b67be58257393227588f5a97b bigdecimal (4.1.2) sha256=53d217666027eab4280346fba98e7d5b66baaae1b9c3c1c0ffe89d48188a3fbd concurrent-ruby (1.3.6) sha256=6b56837e1e7e5292f9864f34b69c5a2cbc75c0cf5338f1ce9903d10fa762d5ab + connection_pool (3.0.2) sha256=33fff5ba71a12d2aa26cb72b1db8bba2a1a01823559fb01d29eb74c286e62e0a date (3.5.1) sha256=750d06384d7b9c15d562c76291407d89e368dda4d4fff957eb94962d325a0dc0 diff-lcs (2.0.0) sha256=708a5d52ec2945b50f8f53a181174aa1ef2c496edf81c05957fe956dabb363d5 drb (2.2.3) sha256=0b00d6fdb50995fe4a45dea13663493c841112e4068656854646f418fda13373 en57 (0.1.0) erb (6.0.4) sha256=38e3803694be357fe2bfe312487c74beaf9fb4e5beb3e22498952fe1645b95d9 + i18n (1.14.8) sha256=285778639134865c5e0f6269e0b818256017e8cde89993fdfcbfb64d088824a5 io-console (0.8.2) sha256=d6e3ae7a7cc7574f4b8893b4fca2162e57a825b223a177b7afa236c5ef9814cc irb (1.18.0) sha256=de9454a0703a54704b9811a5ef31a60c86949fbf4013fcf244fabc7c775248e3 + json (2.19.4) sha256=670a7d333fb3b18ca5b29cb255eb7bef099e40d88c02c80bd42a3f30fe5239ac + logger (1.7.0) sha256=196edec7cc44b66cfb40f9755ce11b392f21f7967696af15d274dde7edff0203 m (1.7.0) sha256=058f793da8150c51353cc59366ffae7774683e868bede3c78f81d920fb9d633a minitest (6.0.5) sha256=f007d7246bf4feea549502842cd7c6aba8851cdc9c90ba06de9c476c0d01155c minitest-mock (5.27.0) sha256=7040ed7185417a966920987eaa6eaf1be4ea1fc5b25bb03ff4703f98564a55b0 @@ -141,8 +179,11 @@ CHECKSUMS sorbet-runtime (0.6.13164) sha256=4561933ac3be0b2477c146540be8e8a3a8b23971ef46cfc89f5d45272830a651 stringio (3.2.0) sha256=c37cb2e58b4ffbd33fe5cd948c05934af997b36e0b6ca6fdf43afa234cf222e1 syntax_tree (6.3.0) sha256=56e25a9692c798ec94c5442fe94c5e94af76bef91edc8bb02052cbdecf35f13d + timeout (0.6.1) sha256=78f57368a7e7bbadec56971f78a3f5ecbcfb59b7fcbb0a3ed6ddc08a5094accb tsort (0.2.0) sha256=9650a793f6859a43b6641671278f79cfead60ac714148aabe4e3f0060480089f + tzinfo (2.0.6) sha256=8daf828cc77bcf7d63b0e3bdb6caa47e2272dcfaf4fbfe46f8c3a9df087a829b unparser (0.9.0) sha256=4331f174a73a23b69250b13d47da3794ed1449711ee0f9ed8947dc020ba76067 + uri (1.1.1) sha256=379fa58d27ffb1387eaada68c749d1426738bd0f654d812fcc07e7568f5c57c6 BUNDLED WITH 4.0.10 diff --git a/lib/en57.rb b/lib/en57.rb index e5dc747..4113a62 100644 --- a/lib/en57.rb +++ b/lib/en57.rb @@ -7,6 +7,7 @@ require_relative "en57/query" require_relative "en57/scope" require_relative "en57/pg_adapter" require_relative "en57/sequel_adapter" +require_relative "en57/active_record_adapter" require_relative "en57/repository" require_relative "en57/event_store" diff --git a/lib/en57/active_record_adapter.rb b/lib/en57/active_record_adapter.rb new file mode 100644 index 0000000..0ce0578 --- /dev/null +++ b/lib/en57/active_record_adapter.rb @@ -0,0 +1,23 @@ +# frozen_string_literal: true + +module En57 + class ActiveRecordAdapter + def initialize(connection_pool) + @connection_pool = connection_pool + end + + def with_connection + @connection_pool.with_connection do |connection| + yield connection.raw_connection + end + end + + def with_serializable_transaction + @connection_pool.with_connection do |connection| + connection.transaction(isolation: :serializable) do + yield connection.raw_connection + end + end + end + end +end diff --git a/test/test_active_record_adapter.rb b/test/test_active_record_adapter.rb new file mode 100644 index 0000000..9987747 --- /dev/null +++ b/test/test_active_record_adapter.rb @@ -0,0 +1,96 @@ +# frozen_string_literal: true + +require "test_helper" + +module En57 + class TestActiveRecordAdapter < Minitest::Test + cover ActiveRecordAdapter + + def test_with_connection_checks_out_connection_and_yields_raw_connection + with_mock_adapter do |pool, connection, raw_connection, adapter| + pool.expect(:with_connection, :selected) do |&block| + block.call(connection) + true + end + connection.expect(:raw_connection, raw_connection) + raw_connection.expect(:exec, :selected, ["SELECT 1"]) + + assert_equal :selected, + adapter.with_connection { |conn| conn.exec("SELECT 1") } + end + end + + def test_with_serializable_transaction_wraps_block_in_transaction + with_mock_adapter do |pool, connection, raw_connection, adapter| + pool.expect(:with_connection, :committed) do |&block| + block.call(connection) + true + end + connection.expect(:transaction, :committed) do |options, &block| + assert_equal({ isolation: :serializable }, options) + block.call + true + end + connection.expect(:raw_connection, raw_connection) + raw_connection.expect( + :exec_params, + :written, + ["SELECT append_events()", []], + ) + + assert_equal( + :committed, + adapter.with_serializable_transaction do |conn| + assert_equal :written, + conn.exec_params("SELECT append_events()", []) + end, + ) + end + end + + def test_with_serializable_transaction_reraises_block_errors + error = RuntimeError.new("boom") + + raised = + assert_raises(RuntimeError) do + with_mock_adapter do |pool, connection, raw_connection, adapter| + pool.expect(:with_connection, nil) do |&block| + block.call(connection) + true + end + connection.expect(:transaction, nil) do |options, &block| + assert_equal({ isolation: :serializable }, options) + block.call + true + end + connection.expect(:raw_connection, raw_connection) + raw_connection.expect(:exec_params, nil) do |sql, params| + assert_equal "SELECT append_events()", sql + assert_equal [], params + raise error + end + + adapter.with_serializable_transaction do |conn| + conn.exec_params("SELECT append_events()", []) + end + end + end + + assert_same error, raised + end + + private + + def with_mock_adapter + pool = Minitest::Mock.new + connection = Minitest::Mock.new + raw_connection = Minitest::Mock.new + + yield pool, connection, raw_connection, ActiveRecordAdapter.new(pool) + ensure + pool.verify + connection.verify + raw_connection.verify + end + end +end diff --git a/test/test_helper.rb b/test/test_helper.rb index 1e6e60c..5893b6a 100644 --- a/test/test_helper.rb +++ b/test/test_helper.rb @@ -9,21 +9,26 @@ require "securerandom" require "concurrent-ruby" require "pg_ephemeral" require "sequel" +require "active_record" module En57 class IntegrationTest < Minitest::Test SERVER = PgEphemeral.start CONNECTION = PG.connect(SERVER.url) SEQUEL_DB = Sequel.connect(SERVER.url) + ActiveRecord::Base.establish_connection(SERVER.url) + AR_POOL = ActiveRecord::Base.connection_pool ADAPTERS = { pg: -> { PgAdapter.new(SERVER.url) }, sequel: -> { SequelAdapter.new(SEQUEL_DB) }, + active_record: -> { ActiveRecordAdapter.new(AR_POOL) }, } def setup = CONNECTION.exec("TRUNCATE TABLE tags, events RESTART IDENTITY CASCADE") Minitest.after_run do + ActiveRecord::Base.connection_pool.disconnect! SEQUEL_DB.disconnect CONNECTION.close SERVER.shutdown