diff --git a/README.md b/README.md index e68223d..ce09572 100644 --- a/README.md +++ b/README.md @@ -81,8 +81,8 @@ Or with pattern matching: case event_store.append([En57::Event.new(type: "OrderPlaced")]) in En57::Result::Success(position:) puts "appended up to #{position}" -in En57::Result::Failure - puts "append condition violated" +in En57::Result::Failure(position:, conflicting_events:) + puts "blocked by #{conflicting_events.size} event(s) up to position #{position}" end ``` @@ -138,7 +138,7 @@ result = event_store.append( ) if result.failure? - # lost the race; another writer already consumed credits + # lost the race; result.conflicting_events shows what got there first end ``` @@ -173,6 +173,6 @@ result = event_store.append( ) if result.failure? - # email already used + # email already used; result.conflicting_events lists the prior registration(s) end ``` diff --git a/db/schema/0.1.0.sql b/db/schema/0.1.0.sql index e2de6a0..d7d03dd 100644 --- a/db/schema/0.1.0.sql +++ b/db/schema/0.1.0.sql @@ -27,7 +27,7 @@ CREATE TYPE en57.event AS ( ); CREATE FUNCTION en57.append_events (new_events en57.event[], append_condition jsonb DEFAULT '{}'::jsonb) - RETURNS TABLE ("position" bigint) + RETURNS TABLE ("position" bigint, conflicting_events text) LANGUAGE plpgsql AS $$ #variable_conflict use_column @@ -35,20 +35,28 @@ DECLARE criteria jsonb[] := ARRAY ( SELECT jsonb_array_elements(COALESCE(append_condition -> 'fail_if_events_match', '[]'::jsonb))); + conflicts jsonb; + conflict_max_position bigint; last_position bigint; BEGIN - IF cardinality(criteria) > 0 AND EXISTS ( - SELECT - 1 - FROM - en57.events AS e - WHERE - EXISTS ( + IF cardinality(criteria) > 0 THEN + WITH matches AS ( SELECT - 1 + 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 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 ( + en57.events AS e + WHERE + EXISTS ( + SELECT + 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 @@ -62,8 +70,18 @@ WHERE FROM en57.tags AS t WHERE - t.event_id = e.id AND t.value = req.value)))) THEN - RAISE EXCEPTION 'append_condition_violated'; + 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) @@ -84,7 +102,7 @@ WHERE FROM unnest(new_events) AS e CROSS JOIN LATERAL unnest(COALESCE(e.tags, ARRAY[]::text[])) AS t (value); - RETURN QUERY SELECT last_position; + RETURN QUERY SELECT last_position, '[]'::text; END; $$; diff --git a/lib/en57/repository.rb b/lib/en57/repository.rb index 4d1c1d1..09ee4a5 100644 --- a/lib/en57/repository.rb +++ b/lib/en57/repository.rb @@ -32,20 +32,42 @@ module En57 :fail_if_events_match ] = fail_if_events_match unless fail_if_events_match.empty? - result = - @adapter.with_serializable_transaction do |connection| - connection.exec_params( - "SELECT position FROM en57.append_events($1::en57.event[], $2::jsonb)", - [ - @array_encoder.encode(event_records), - JSON.generate(append_condition), - ], - ) - end - position = result.first.fetch("position") - Result.success(position: position && Integer(position)) - rescue PG::RaiseException, PG::TRSerializationFailure - Result.failure(position: nil) + 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? + Result.success(position:) + else + Result.failure(position:, conflicting_events:) + end + rescue PG::TRSerializationFailure + Result.failure(position: nil, conflicting_events: []) end def read(query) diff --git a/lib/en57/result.rb b/lib/en57/result.rb index 6b73232..e3a7662 100644 --- a/lib/en57/result.rb +++ b/lib/en57/result.rb @@ -9,13 +9,14 @@ module En57 end Failure = - Data.define(:position) do + Data.define(:position, :conflicting_events) do def success? = false def failure? = true end def self.success(position:) = Success.new(position:) - def self.failure(position:) = Failure.new(position:) + def self.failure(position:, conflicting_events:) = + Failure.new(position:, conflicting_events:) end end diff --git a/test/test_integration.rb b/test/test_integration.rb index 4f5f956..06cd71d 100644 --- a/test/test_integration.rb +++ b/test/test_integration.rb @@ -67,7 +67,9 @@ module En57 fail_if: event_store.read.of_type("OrderPlaced"), ) - assert(result.failure?) + assert(Result::Failure === result) + assert_equal(1, result.position) + assert_equal([existing_event], result.conflicting_events) assert_equal([existing_event], event_store.read.each.to_a) end end diff --git a/test/test_repository.rb b/test/test_repository.rb index f301327..6c42ff2 100644 --- a/test/test_repository.rb +++ b/test/test_repository.rb @@ -34,9 +34,9 @@ module En57 connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, - [{ "position" => "2" }], + [{ "position" => "2", "conflicting_events" => "[]" }], [ - "SELECT position FROM en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT position, conflicting_events FROM en57.append_events($1::en57.event[], $2::jsonb)", [expected_events, "{}"], ], ) @@ -81,9 +81,9 @@ module En57 connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, - [{ "position" => "1" }], + [{ "position" => "1", "conflicting_events" => "[]" }], [ - "SELECT position FROM en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT position, conflicting_events FROM en57.append_events($1::en57.event[], $2::jsonb)", [expected_events, "{}"], ], ) @@ -107,9 +107,9 @@ module En57 connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, - [{ "position" => nil }], + [{ "position" => nil, "conflicting_events" => "[]" }], [ - "SELECT position FROM en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT position, conflicting_events FROM en57.append_events($1::en57.event[], $2::jsonb)", [ array_encoder.encode([]), '{"fail_if_events_match":[{"types":["OrderPlaced"],"after":42}]}', @@ -146,7 +146,7 @@ module En57 connection.expect(:exec, nil, ["ROLLBACK"]) connection.expect(:exec_params, nil) do |sql, params| assert_equal( - "SELECT position FROM en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT position, conflicting_events FROM en57.append_events($1::en57.event[], $2::jsonb)", sql, ) assert_equal([array_encoder.encode([]), "{}"], params) @@ -522,19 +522,76 @@ module En57 end end - def test_append_returns_failure_on_pg_raise_exception + def test_append_returns_failure_when_condition_matches with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) - connection.expect(:exec, nil, ["ROLLBACK"]) - connection.expect(:exec_params, nil) { raise(PG::RaiseException.new) } + connection.expect( + :exec_params, + [ + { + "position" => "5", + "conflicting_events" => + JSON.generate( + [ + { + id: ids[2], + type: "OrderPlaced", + data: '{"amount":100}', + meta: '{"amount":{"k":"Symbol"}}', + tags: ["order_id:123"], + }, + { + id: ids[3], + type: "OrderPlaced", + data: nil, + meta: nil, + tags: [], + }, + ], + ), + }, + ], + [ + "SELECT position, conflicting_events FROM en57.append_events($1::en57.event[], $2::jsonb)", + [ + array_encoder.encode([]), + '{"fail_if_events_match":[{"types":["OrderPlaced"]}]}', + ], + ], + ) + connection.expect(:exec, nil, ["COMMIT"]) result = Repository.new( PgAdapter.for_connection(connection), JsonSerializer.new, - ).append([], fail_if: Query.all) + ).append( + [], + fail_if: + Query.new( + criteria: [ + Query::Criteria.new(types: ["OrderPlaced"], tags: []), + ], + ), + ) - assert_equal(Result.failure(position: nil), result) + assert_equal( + Result.failure( + position: 5, + conflicting_events: [ + Event.new( + id: ids[2], + type: "OrderPlaced", + data: { + amount: 100, + }, + tags: ["order_id:123"], + ), + Event.new(id: ids[3], type: "OrderPlaced", data: {}, tags: []), + ], + ), + result, + ) end end @@ -552,7 +609,10 @@ module En57 JsonSerializer.new, ).append([], fail_if: Query.all) - assert_equal(Result.failure(position: nil), result) + assert_equal( + Result.failure(position: nil, conflicting_events: []), + result, + ) end end diff --git a/test/test_result.rb b/test/test_result.rb index 3834e84..6935c68 100644 --- a/test/test_result.rb +++ b/test/test_result.rb @@ -14,12 +14,14 @@ module En57 assert_equal(42, result.position) end - def test_failure_carries_position - result = Result.failure(position: 7) + def test_failure_carries_position_and_conflicting_events + conflict = Event.new(type: "OrderPlaced") + result = Result.failure(position: 7, conflicting_events: [conflict]) assert(result.failure?) refute(result.success?) assert_equal(7, result.position) + assert_equal([conflict], result.conflicting_events) end end end