diff --git a/db/schema/0.1.0.sql b/db/schema/0.1.0.sql index d7d03dd..94e2ae2 100644 --- a/db/schema/0.1.0.sql +++ b/db/schema/0.1.0.sql @@ -27,10 +27,12 @@ CREATE TYPE en57.event AS ( ); CREATE FUNCTION en57.append_events (new_events en57.event[], append_condition jsonb DEFAULT '{}'::jsonb) - RETURNS TABLE ("position" bigint, conflicting_events text) + RETURNS TABLE ( + "position" bigint, + conflicting_events text) LANGUAGE plpgsql AS $$ -#variable_conflict use_column + # variable_conflict use_column DECLARE criteria jsonb[] := ARRAY ( SELECT @@ -44,10 +46,11 @@ BEGIN SELECT e.position, jsonb_build_object('id', e.id, 'type', e.type, 'data', e.data::text, 'meta', e.meta::text, 'tags', COALESCE(( - SELECT - array_agg(t.value ORDER BY t.value) FROM en57.tags AS t - WHERE - t.event_id = e.id), ARRAY[]::text[])) AS event_obj + SELECT + array_agg(t.value ORDER BY t.value) + FROM en57.tags AS t + WHERE + t.event_id = e.id), ARRAY[]::text[])) AS event_obj FROM en57.events AS e WHERE @@ -56,35 +59,44 @@ BEGIN 1 FROM unnest(criteria) AS c - WHERE ((c ->> 'after')::bigint IS NULL OR e.position > (c ->> 'after')::bigint) AND (c -> 'types' IS NULL OR e.type IN ( -SELECT - jsonb_array_elements_text(c -> 'types'))) AND NOT EXISTS ( -SELECT - 1 -FROM - jsonb_array_elements_text(COALESCE(c -> 'tags', '[]'::jsonb)) AS req (value) -WHERE - NOT EXISTS ( - SELECT - 1 - FROM - en57.tags AS t - WHERE - t.event_id = e.id AND t.value = req.value))) -) - SELECT - max(position), - jsonb_agg(event_obj ORDER BY position) INTO conflict_max_position, - conflicts - FROM - matches; - IF conflict_max_position IS NOT NULL THEN - RETURN QUERY SELECT conflict_max_position, conflicts::text; - RETURN; - END IF; + WHERE ((c ->> 'after')::bigint IS NULL + OR e.position > (c ->> 'after')::bigint) + AND (c -> 'types' IS NULL + OR e.type IN ( + SELECT + jsonb_array_elements_text(c -> 'types'))) + AND NOT EXISTS ( + SELECT + 1 + FROM + jsonb_array_elements_text(COALESCE(c -> 'tags', '[]'::jsonb)) AS req (value) + WHERE + NOT EXISTS ( + SELECT + 1 + FROM + en57.tags AS t + WHERE + t.event_id = e.id + AND t.value = req.value)))) + SELECT + max(position), + jsonb_agg(event_obj ORDER BY position) + INTO + conflict_max_position, + conflicts + FROM + matches; + IF conflict_max_position IS NOT NULL THEN + RETURN QUERY + SELECT + conflict_max_position, + conflicts::text; + RETURN; + END IF; END IF; WITH inserted AS ( - INSERT INTO en57.events (id, type, data, meta) +INSERT INTO en57.events (id, type, data, meta) SELECT e.id, e.type, @@ -92,9 +104,15 @@ WHERE e.meta FROM unnest(new_events) AS e - RETURNING position - ) - SELECT max(position) INTO last_position FROM inserted; + RETURNING + position +) + SELECT + max(position) + INTO + last_position + FROM + inserted; INSERT INTO en57.tags (event_id, value) SELECT e.id, @@ -102,7 +120,10 @@ WHERE FROM unnest(new_events) AS e CROSS JOIN LATERAL unnest(COALESCE(e.tags, ARRAY[]::text[])) AS t (value); - RETURN QUERY SELECT last_position, '[]'::text; + RETURN QUERY + SELECT + last_position, + '[]'::text; END; $$; diff --git a/lib/en57/configuration.rb b/lib/en57/configuration.rb index 1670cbf..9f79e4e 100644 --- a/lib/en57/configuration.rb +++ b/lib/en57/configuration.rb @@ -6,10 +6,11 @@ module En57 class Configuration include Singleton - attr_accessor :serializer + attr_accessor :serializer, :max_retries def initialize @serializer = JsonSerializer.new + @max_retries = 10 end end end diff --git a/lib/en57/repository.rb b/lib/en57/repository.rb index 466c907..8eb59a0 100644 --- a/lib/en57/repository.rb +++ b/lib/en57/repository.rb @@ -4,9 +4,14 @@ require "pg" module En57 class Repository - def initialize(adapter, serializer = En57.configuration.serializer) + def initialize( + adapter, + serializer = En57.configuration.serializer, + max_retries: En57.configuration.max_retries + ) @adapter = adapter @serializer = serializer + @max_retries = max_retries @record_encoder = PG::TextEncoder::Record.new @array_encoder = PG::TextEncoder::Array.new @array_decoder = PG::TextDecoder::Array.new @@ -32,39 +37,41 @@ module En57 :fail_if_events_match ] = fail_if_events_match unless fail_if_events_match.empty? - row = - @adapter - .with_serializable_transaction do |connection| - connection.exec_params( - "SELECT position, conflicting_events FROM en57.append_events($1::en57.event[], $2::jsonb)", - [ - @array_encoder.encode(event_records), - JSON.generate(append_condition), - ], - ) - end - .first - raw_position = row.fetch("position") - position = raw_position && Integer(raw_position) - conflicting_events = - JSON - .parse(row.fetch("conflicting_events")) - .map do |conflict| - Event.new( - id: conflict.fetch("id"), - type: conflict.fetch("type"), - data: - @serializer.load( - conflict.fetch("data") || "{}", - conflict.fetch("meta"), - ), - tags: conflict.fetch("tags"), - ) - end - if conflicting_events.empty? - Success.new(position:) - else - Failure.new(position:, conflicting_events:) + with_retries do + row = + @adapter + .with_serializable_transaction do |connection| + connection.exec_params( + "SELECT position, conflicting_events FROM en57.append_events($1::en57.event[], $2::jsonb)", + [ + @array_encoder.encode(event_records), + JSON.generate(append_condition), + ], + ) + end + .first + raw_position = row.fetch("position") + position = raw_position && Integer(raw_position) + conflicting_events = + JSON + .parse(row.fetch("conflicting_events")) + .map do |conflict| + Event.new( + id: conflict.fetch("id"), + type: conflict.fetch("type"), + data: + @serializer.load( + conflict.fetch("data") || "{}", + conflict.fetch("meta"), + ), + tags: conflict.fetch("tags"), + ) + end + if conflicting_events.empty? + Success.new(position:) + else + Failure.new(position:, conflicting_events:) + end end end @@ -91,5 +98,18 @@ module En57 ] end end + + private + + def with_retries + attempts = 0 + begin + yield + rescue PG::TRSerializationFailure + raise if attempts == @max_retries + attempts += 1 + retry + end + end end end diff --git a/test/test_en57.rb b/test/test_en57.rb index 3a0f030..53f4a8a 100644 --- a/test/test_en57.rb +++ b/test/test_en57.rb @@ -28,6 +28,12 @@ module En57 assert_kind_of JsonSerializer, En57.configuration.serializer end + def test_configuration_default_retry_settings + with_empty_configuration do + assert_equal 10, En57.configuration.max_retries + end + end + def test_configure with_empty_configuration do serializer = Object.new diff --git a/test/test_repository.rb b/test/test_repository.rb index 7dcd24b..8d90f24 100644 --- a/test/test_repository.rb +++ b/test/test_repository.rb @@ -595,19 +595,78 @@ module En57 end end - def test_append_lets_serialization_failure_surface + def test_append_retries_on_serialization_failure_then_succeeds with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect(:exec, nil, ["ROLLBACK"]) connection.expect(:exec_params, nil) do raise PG::TRSerializationFailure.new end + connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) + connection.expect( + :exec_params, + [{ "position" => "1", "conflicting_events" => "[]" }], + [ + "SELECT position, conflicting_events FROM en57.append_events($1::en57.event[], $2::jsonb)", + [array_encoder.encode([]), "{}"], + ], + ) + connection.expect(:exec, nil, ["COMMIT"]) + + repository = + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + max_retries: 1, + ) + + repository.stub(:sleep, ->(*) { raise "sleep should not be called" }) do + result = repository.append([], fail_if: Query.all) + assert_equal(Success.new(position: 1), result) + end + end + end + + def test_append_raises_after_exhausting_retries + with_connection do |connection| + 2.times do + connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) + connection.expect(:exec, nil, ["ROLLBACK"]) + connection.expect(:exec_params, nil) do + raise PG::TRSerializationFailure.new + end + end + + repository = + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + max_retries: 1, + ) assert_raises(PG::TRSerializationFailure) do + repository.append([], fail_if: Query.all) + end + end + end + + def test_append_raises_without_retrying_when_max_retries_is_zero + with_connection do |connection| + connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) + connection.expect(:exec, nil, ["ROLLBACK"]) + connection.expect(:exec_params, nil) do + raise PG::TRSerializationFailure.new + end + + repository = Repository.new( PgAdapter.for_connection(connection), JsonSerializer.new, - ).append([], fail_if: Query.all) + max_retries: 0, + ) + + assert_raises(PG::TRSerializationFailure) do + repository.append([], fail_if: Query.all) end end end diff --git a/test/test_stress.rb b/test/test_stress.rb index eb22720..e361afc 100644 --- a/test/test_stress.rb +++ b/test/test_stress.rb @@ -42,14 +42,12 @@ module En57 ], fail_if: account_scope.of_type("CreditsUsed"), ) - rescue PG::TRSerializationFailure => e - e end end assert_equal( (concurrency - 1), - threads.map(&:value).reject { Success === it }.size, + threads.map(&:value).count(&:failure?), ) assert_equal( 1,