diff --git a/lib/en57.rb b/lib/en57.rb index 8dee0c3..f21b649 100644 --- a/lib/en57.rb +++ b/lib/en57.rb @@ -5,6 +5,7 @@ require_relative "en57/event" require_relative "en57/json_serializer" require_relative "en57/query" require_relative "en57/scope" +require_relative "en57/pg_adapter" require_relative "en57/pg_repository" require_relative "en57/event_store" diff --git a/lib/en57/pg_adapter.rb b/lib/en57/pg_adapter.rb new file mode 100644 index 0000000..22f1cab --- /dev/null +++ b/lib/en57/pg_adapter.rb @@ -0,0 +1,27 @@ +# frozen_string_literal: true + +require "pg" + +module En57 + class PgAdapter + def initialize(connection_uri) + @connection_uri = connection_uri + end + + def with_connection + @connection ||= PG.connect(@connection_uri) + yield @connection + end + + def with_serializable_transaction + with_connection do |connection| + connection.exec("BEGIN ISOLATION LEVEL SERIALIZABLE") + yield connection + connection.exec("COMMIT") + rescue StandardError + connection.exec("ROLLBACK") + raise + end + end + end +end diff --git a/lib/en57/pg_repository.rb b/lib/en57/pg_repository.rb index 9f1c65c..d69f46d 100644 --- a/lib/en57/pg_repository.rb +++ b/lib/en57/pg_repository.rb @@ -5,7 +5,7 @@ require "pg" module En57 class PgRepository def initialize(connection_uri, serializer) - @connection_uri = connection_uri + @adapter = PgAdapter.new(connection_uri) @serializer = serializer @record_encoder = PG::TextEncoder::Record.new @array_encoder = PG::TextEncoder::Array.new @@ -32,7 +32,7 @@ module En57 :fail_if_events_match ] = fail_if_events_match unless fail_if_events_match.empty? - with_serializable_transaction do |connection| + @adapter.with_serializable_transaction do |connection| connection.exec_params( "SELECT append_events($1::event_with_tags[], $2::jsonb)", [ @@ -53,7 +53,7 @@ module En57 def read(query) criteria = query.encoded_criteria.map { |item| JSON.generate(item) } - with_connection do |connection| + @adapter.with_connection do |connection| connection.exec_params( "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", [@array_encoder.encode(criteria)], @@ -67,23 +67,5 @@ module En57 ) end end - - private - - def with_serializable_transaction - with_connection do |connection| - connection.exec("BEGIN ISOLATION LEVEL SERIALIZABLE") - yield connection - connection.exec("COMMIT") - rescue StandardError - connection.exec("ROLLBACK") - raise - end - end - - def with_connection - @connection ||= PG.connect(@connection_uri) - yield @connection - end end end