diff --git a/README.md b/README.md index a220b45..28ec53a 100644 --- a/README.md +++ b/README.md @@ -109,7 +109,7 @@ begin ], fail_if: account_scope.of_type("CreditsUsed"), ) -rescue En57::AppendConditionViolated, PG::TRSerializationFailure +rescue En57::AppendConditionViolated # lost the race; another writer already consumed credits end ``` diff --git a/lib/en57/active_record_adapter.rb b/lib/en57/active_record_adapter.rb index 0ce0578..5ee5429 100644 --- a/lib/en57/active_record_adapter.rb +++ b/lib/en57/active_record_adapter.rb @@ -1,5 +1,7 @@ # frozen_string_literal: true +require "pg" + module En57 class ActiveRecordAdapter def initialize(connection_pool) @@ -18,6 +20,10 @@ module En57 yield connection.raw_connection end end + rescue ActiveRecord::StatementInvalid => e + raise e.cause if e.cause.is_a?(PG::Error) + + raise end end end diff --git a/lib/en57/repository.rb b/lib/en57/repository.rb index 553091b..fb8f022 100644 --- a/lib/en57/repository.rb +++ b/lib/en57/repository.rb @@ -4,6 +4,11 @@ require "pg" module En57 class Repository + SQL_STATES = [ + RAISE_EXCEPTION = "P0001", + SERIALIZATION_FAILURE = "40001", + ].freeze + def initialize(adapter, serializer) @adapter = adapter @serializer = serializer @@ -45,7 +50,7 @@ module En57 sqlstate = e.result&.error_field(PG::Result::PG_DIAG_SQLSTATE) || (e.sqlstate if e.respond_to?(:sqlstate)) - raise AppendConditionViolated if sqlstate == "P0001" + raise AppendConditionViolated if SQL_STATES.include?(sqlstate) raise end @@ -53,19 +58,21 @@ module En57 def read(query) criteria = query.encoded_criteria.map { |item| JSON.generate(item) } - @adapter.with_connection do |connection| - connection.exec_params( - "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", - [@array_encoder.encode(criteria)], - ) - end.map do |row| - Event.new( - id: row.fetch("id"), - type: row.fetch("type"), - data: @serializer.load(row.fetch("data"), row.fetch("meta")), - tags: @array_decoder.decode(row.fetch("tags")), - ) - end + @adapter + .with_connection do |connection| + connection.exec_params( + "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", + [@array_encoder.encode(criteria)], + ) + end + .map do |row| + Event.new( + id: row.fetch("id"), + type: row.fetch("type"), + data: @serializer.load(row.fetch("data"), row.fetch("meta")), + tags: @array_decoder.decode(row.fetch("tags")), + ) + end end end end diff --git a/test/test_active_record_adapter.rb b/test/test_active_record_adapter.rb index 9987747..9016853 100644 --- a/test/test_active_record_adapter.rb +++ b/test/test_active_record_adapter.rb @@ -79,8 +79,59 @@ module En57 assert_same error, raised end + def test_with_serializable_transaction_unwraps_pg_errors + assert_unwraps_pg_error(PG::Error.new("boom")) + end + + def test_with_serializable_transaction_unwraps_pg_error_subclasses + assert_unwraps_pg_error(PG::TRSerializationFailure.new("boom")) + end + + def test_with_serializable_transaction_reraises_non_pg_active_record_errors + ar_error = ActiveRecord::StatementInvalid.new("wrapped") + ar_error.define_singleton_method(:cause) { RuntimeError.new("boom") } + + raised = + assert_raises(ActiveRecord::StatementInvalid) do + with_mock_adapter do |pool, connection, _raw_connection, adapter| + expect_failed_transaction(pool, connection, ar_error) + + adapter.with_serializable_transaction { flunk "not yielded" } + end + end + + assert_same ar_error, raised + end + private + def assert_unwraps_pg_error(pg_error) + ar_error = ActiveRecord::StatementInvalid.new("wrapped") + ar_error.define_singleton_method(:cause) { pg_error } + + raised = + assert_raises(PG::Error) do + with_mock_adapter do |pool, connection, _raw_connection, adapter| + expect_failed_transaction(pool, connection, ar_error) + + adapter.with_serializable_transaction { flunk "not yielded" } + end + end + + assert_same pg_error, raised + end + + def expect_failed_transaction(pool, connection, error) + pool.expect(:with_connection, nil) do |&block| + block.call(connection) + true + end + connection.expect(:transaction, nil) do |options, &_block| + assert_equal({ isolation: :serializable }, options) + raise error + end + end + def with_mock_adapter pool = Minitest::Mock.new connection = Minitest::Mock.new diff --git a/test/test_repository.rb b/test/test_repository.rb index 9474c4b..206dd15 100644 --- a/test/test_repository.rb +++ b/test/test_repository.rb @@ -160,6 +160,40 @@ module En57 end end + def test_append_raises_append_condition_violated_from_serialization_failure_result_sqlstate + with_connection_to(connection_uri) do |connection| + connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) + connection.expect(:exec, nil, ["ROLLBACK"]) + connection.expect(:exec_params, nil) do + raise pg_error(result_sqlstate: "40001") + end + + assert_raises(AppendConditionViolated) do + Repository.new( + PgAdapter.new(connection_uri), + JsonSerializer.new, + ).append([], fail_if: Query.all) + end + end + end + + def test_append_raises_append_condition_violated_from_serialization_failure_sqlstate + with_connection_to(connection_uri) do |connection| + connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) + connection.expect(:exec, nil, ["ROLLBACK"]) + connection.expect(:exec_params, nil) do + raise pg_error(sqlstate: "40001") + end + + assert_raises(AppendConditionViolated) do + Repository.new( + PgAdapter.new(connection_uri), + JsonSerializer.new, + ).append([], fail_if: Query.all) + end + end + end + def test_append_reraises_pg_error_for_non_append_condition_sqlstate with_connection_to(connection_uri) do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) diff --git a/test/test_stress.rb b/test/test_stress.rb index 6650839..820b236 100644 --- a/test/test_stress.rb +++ b/test/test_stress.rb @@ -44,7 +44,7 @@ module En57 ) end end - rescue AppendConditionViolated, PG::TRSerializationFailure + rescue AppendConditionViolated end end threads.each(&:join)