diff --git a/README.md b/README.md index 9cd066c..d32c00d 100644 --- a/README.md +++ b/README.md @@ -121,8 +121,8 @@ result = event_store.append( ) case result -in En57::Success - # credits consumed +in En57::Success(position:) + # credits consumed at event position in En57::Failure # lost the race; another writer already consumed credits end @@ -159,8 +159,8 @@ result = event_store.append( ) case result -in En57::Success - # user registered +in En57::Success(position:) + # user registered at event position in En57::Failure # email already used end diff --git a/db/schema/0.1.0.sql b/db/schema/0.1.0.sql index 1e3e6e6..67c05f1 100644 --- a/db/schema/0.1.0.sql +++ b/db/schema/0.1.0.sql @@ -27,7 +27,8 @@ CREATE TYPE en57.event AS ( ); CREATE TYPE en57.append_result AS ( - status text + status text, + "position" bigint ); CREATE FUNCTION en57.append_events (new_events en57.event[], append_condition jsonb DEFAULT '{}'::jsonb) @@ -43,6 +44,7 @@ DECLARE req_types text[]; req_tags text[]; req_after bigint; + appended_position bigint; BEGIN FOREACH criterion IN ARRAY criteria LOOP req_after := (criterion ->> 'after')::bigint; @@ -71,7 +73,8 @@ BEGIN OR e.position > req_after) AND (criterion -> 'types' IS NULL OR e.type = ANY (req_types))) THEN - RETURN ROW ('append_condition_violated')::en57.append_result; + RETURN ROW ('append_condition_violated', + NULL)::en57.append_result; END IF; ELSE IF EXISTS ( @@ -83,26 +86,38 @@ BEGIN OR e.position > req_after) AND (criterion -> 'types' IS NULL OR e.type = ANY (req_types))) THEN - RETURN ROW ('append_condition_violated')::en57.append_result; + RETURN ROW ('append_condition_violated', + NULL)::en57.append_result; END IF; END IF; END LOOP; + WITH inserted_events AS ( INSERT INTO en57.events (id, type, data, meta) -SELECT - e.id, - e.type, - e.data, - e.meta -FROM - unnest(new_events) AS e; -INSERT INTO en57.tags (event_id, value) -SELECT - e.id, - t.value -FROM - unnest(new_events) AS e + SELECT + e.id, + e.type, + e.data, + e.meta + FROM + unnest(new_events) AS e + RETURNING + "position" +) + SELECT + max("position") + FROM + inserted_events + INTO + appended_position; + INSERT INTO en57.tags (event_id, value) + SELECT + e.id, + t.value + FROM + unnest(new_events) AS e CROSS JOIN LATERAL unnest(COALESCE(e.tags, ARRAY[]::text[])) AS t (value); - RETURN ROW ('success')::en57.append_result; + RETURN ROW ('success', + appended_position)::en57.append_result; END; $$; diff --git a/lib/en57.rb b/lib/en57.rb index 8ef0864..171dcc0 100644 --- a/lib/en57.rb +++ b/lib/en57.rb @@ -14,7 +14,7 @@ require_relative "en57/event_store" require_relative "en57/configuration" module En57 - Success = Data.define + Success = Data.define(:position) Failure = Data.define AppendRetriesExhausted = Class.new(StandardError) diff --git a/lib/en57/repository.rb b/lib/en57/repository.rb index 0d1b02c..2e5ab05 100644 --- a/lib/en57/repository.rb +++ b/lib/en57/repository.rb @@ -33,7 +33,7 @@ module En57 ] = fail_if_events_match unless fail_if_events_match.empty? statement = - "SELECT status FROM en57.append_events($1::en57.event[], $2::jsonb)" + "SELECT status, position FROM en57.append_events($1::en57.event[], $2::jsonb)" params = [ @array_encoder.encode(event_records), JSON.generate(append_condition), @@ -54,7 +54,7 @@ module En57 case row.first.fetch("status") when "success" - Success.new + Success.new(position: row.first.fetch("position")&.then { Integer(it) }) when "append_condition_violated" Failure.new end diff --git a/test/pg_regress/expected/001_schema.out b/test/pg_regress/expected/001_schema.out index 1e3e6e6..67c05f1 100644 --- a/test/pg_regress/expected/001_schema.out +++ b/test/pg_regress/expected/001_schema.out @@ -27,7 +27,8 @@ CREATE TYPE en57.event AS ( ); CREATE TYPE en57.append_result AS ( - status text + status text, + "position" bigint ); CREATE FUNCTION en57.append_events (new_events en57.event[], append_condition jsonb DEFAULT '{}'::jsonb) @@ -43,6 +44,7 @@ DECLARE req_types text[]; req_tags text[]; req_after bigint; + appended_position bigint; BEGIN FOREACH criterion IN ARRAY criteria LOOP req_after := (criterion ->> 'after')::bigint; @@ -71,7 +73,8 @@ BEGIN OR e.position > req_after) AND (criterion -> 'types' IS NULL OR e.type = ANY (req_types))) THEN - RETURN ROW ('append_condition_violated')::en57.append_result; + RETURN ROW ('append_condition_violated', + NULL)::en57.append_result; END IF; ELSE IF EXISTS ( @@ -83,26 +86,38 @@ BEGIN OR e.position > req_after) AND (criterion -> 'types' IS NULL OR e.type = ANY (req_types))) THEN - RETURN ROW ('append_condition_violated')::en57.append_result; + RETURN ROW ('append_condition_violated', + NULL)::en57.append_result; END IF; END IF; END LOOP; + WITH inserted_events AS ( INSERT INTO en57.events (id, type, data, meta) -SELECT - e.id, - e.type, - e.data, - e.meta -FROM - unnest(new_events) AS e; -INSERT INTO en57.tags (event_id, value) -SELECT - e.id, - t.value -FROM - unnest(new_events) AS e + SELECT + e.id, + e.type, + e.data, + e.meta + FROM + unnest(new_events) AS e + RETURNING + "position" +) + SELECT + max("position") + FROM + inserted_events + INTO + appended_position; + INSERT INTO en57.tags (event_id, value) + SELECT + e.id, + t.value + FROM + unnest(new_events) AS e CROSS JOIN LATERAL unnest(COALESCE(e.tags, ARRAY[]::text[])) AS t (value); - RETURN ROW ('success')::en57.append_result; + RETURN ROW ('success', + appended_position)::en57.append_result; END; $$; diff --git a/test/test_event_store.rb b/test/test_event_store.rb index 0f92989..a8e893c 100644 --- a/test/test_event_store.rb +++ b/test/test_event_store.rb @@ -10,7 +10,12 @@ module En57 event = Event.new(type: "CreditsToppedUp") with_repository do |repository| - repository.expect(:append, Success.new, [[event]], fail_if: Query.all) + repository.expect( + :append, + Success.new(position: 1), + [[event]], + fail_if: Query.all, + ) EventStore.new(repository).append([event]) end @@ -33,9 +38,17 @@ module En57 event = Event.new(type: "CreditsToppedUp") with_repository do |repository| - repository.expect(:append, Success.new, [[event]], fail_if: Query.all) + repository.expect( + :append, + Success.new(position: 1), + [[event]], + fail_if: Query.all, + ) - assert_equal(Success.new, EventStore.new(repository).append([event])) + assert_equal( + Success.new(position: 1), + EventStore.new(repository).append([event]), + ) end end @@ -47,7 +60,7 @@ module En57 fail_if = event_store.read.with_tag("order_id:123") repository.expect( :append, - Success.new, + Success.new(position: 1), [[event]], fail_if: fail_if.to_query, ) diff --git a/test/test_factories.rb b/test/test_factories.rb index e352a0c..cb2f92c 100644 --- a/test/test_factories.rb +++ b/test/test_factories.rb @@ -44,7 +44,7 @@ module En57 def assert_round_trip(event_store) event = Event.new(type: "FactoryTested") - assert_equal Success.new, event_store.append([event]) + assert_equal Success.new(position: 1), event_store.append([event]) assert_equal [event], event_store.read.each.to_a end end diff --git a/test/test_integration.rb b/test/test_integration.rb index 90429de..1f3a788 100644 --- a/test/test_integration.rb +++ b/test/test_integration.rb @@ -24,7 +24,7 @@ module En57 ), ] - assert_equal(Success.new, event_store.append(events)) + assert_equal(Success.new(position: 2), event_store.append(events)) assert_equal(events, event_store.read.each.to_a) end end @@ -36,7 +36,7 @@ module En57 Event.new(id: ids[1], type: "PriceChanged"), ] - assert_equal(Success.new, event_store.append(events)) + assert_equal(Success.new(position: 2), event_store.append(events)) assert_equal( events.map.with_index(1) { |event, position| [event, position] }, event_store.read.each_with_position.to_a, @@ -48,7 +48,7 @@ module En57 with_event_store(factory) do |event_store| event = Event.new(id: ids[0], type: "OrderPlaced") assert_equal( - Success.new, + Success.new(position: 1), event_store.append( [event], fail_if: event_store.read.of_type("PriceChanged"), @@ -62,7 +62,10 @@ module En57 define_method "test_#{name}_append_with_fail_if_and_matches_returns_failure" do with_event_store(factory) do |event_store| existing_event = Event.new(id: ids[0], type: "OrderPlaced") - assert_equal(Success.new, event_store.append([existing_event])) + assert_equal( + Success.new(position: 1), + event_store.append([existing_event]), + ) assert_equal( Failure.new, @@ -79,9 +82,12 @@ module En57 define_method "test_#{name}_append_with_after_ignores_matches_at_or_before_cutoff" do with_event_store(factory) do |event_store| existing_event = Event.new(id: ids[0], type: "OrderPlaced") - assert_equal(Success.new, event_store.append([existing_event])) assert_equal( - Success.new, + Success.new(position: 1), + event_store.append([existing_event]), + ) + assert_equal( + Success.new(position: 2), event_store.append( [Event.new(id: ids[1], type: "ShipmentScheduled")], fail_if: event_store.read.of_type("OrderPlaced").after(1), @@ -98,7 +104,10 @@ module En57 define_method "test_#{name}_append_with_after_returns_failure_if_match_is_after_cutoff" do with_event_store(factory) do |event_store| existing_event = Event.new(id: ids[0], type: "OrderPlaced") - assert_equal(Success.new, event_store.append([existing_event])) + assert_equal( + Success.new(position: 1), + event_store.append([existing_event]), + ) assert_equal( Failure.new, @@ -115,7 +124,10 @@ module En57 define_method "test_#{name}_append_with_duplicate_id_raises_unique_violation" do with_event_store(factory) do |event_store| existing_event = Event.new(id: ids[0], type: "OrderPlaced") - assert_equal(Success.new, event_store.append([existing_event])) + assert_equal( + Success.new(position: 1), + event_store.append([existing_event]), + ) assert_raises(PG::UniqueViolation) do event_store.append( @@ -132,7 +144,7 @@ module En57 event = Event.new(id: ids[0], type: "OrderPlaced", tags: ["order_id:123"]) - assert_equal(Success.new, event_store.append([event])) + assert_equal(Success.new(position: 1), event_store.append([event])) assert_equal([event], event_store.read.each.to_a) end end @@ -144,7 +156,7 @@ module En57 Event.new(id: ids[1], type: "PriceChanged"), ] - assert_equal(Success.new, event_store.append(events)) + assert_equal(Success.new(position: 2), event_store.append(events)) assert_equal(events.drop(1), event_store.read.after(1).each.to_a) end end @@ -164,7 +176,7 @@ module En57 ), ] - assert_equal(Success.new, event_store.append(events)) + assert_equal(Success.new(position: 2), event_store.append(events)) assert_equal( events.take(1), event_store @@ -183,7 +195,7 @@ module En57 Event.new(id: ids[1], type: "PriceChanged"), ] - assert_equal(Success.new, event_store.append(events)) + assert_equal(Success.new(position: 2), event_store.append(events)) assert_equal( events.take(1), event_store.read.of_type("OrderPlaced").each.to_a, @@ -199,7 +211,7 @@ module En57 Event.new(id: ids[2], type: "OrderCancelled"), ] - assert_equal(Success.new, event_store.append(events)) + assert_equal(Success.new(position: 3), event_store.append(events)) assert_equal( events.drop(1), event_store.read.of_type("OrderPlaced", "OrderCancelled").each.to_a, @@ -215,7 +227,7 @@ module En57 Event.new(id: ids[2], type: "PriceChanged", tags: ["order_id:123"]), ] - assert_equal(Success.new, event_store.append(events)) + assert_equal(Success.new(position: 3), event_store.append(events)) assert_equal( events.take(1), event_store @@ -241,7 +253,7 @@ module En57 tags: ["order_id:123"], ), ] - assert_equal(Success.new, event_store.append(events)) + assert_equal(Success.new(position: 5), event_store.append(events)) orders = event_store.read.of_type("OrderPlaced").with_tag("order_id:123") diff --git a/test/test_migrator.rb b/test/test_migrator.rb index f06468d..db74e45 100644 --- a/test/test_migrator.rb +++ b/test/test_migrator.rb @@ -55,7 +55,7 @@ module En57 ), ) - assert_equal Success.new, event_store.append([event]) + assert_equal Success.new(position: 1), event_store.append([event]) assert_equal [event], event_store.read.each.to_a ensure connection&.close diff --git a/test/test_repository.rb b/test/test_repository.rb index 3aefccf..f2789b5 100644 --- a/test/test_repository.rb +++ b/test/test_repository.rb @@ -36,7 +36,7 @@ module En57 :exec_params, success_result, [ - "SELECT status FROM en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT status, position FROM en57.append_events($1::en57.event[], $2::jsonb)", [expected_events, "{}"], ], ) @@ -80,7 +80,7 @@ module En57 :exec_params, success_result, [ - "SELECT status FROM en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT status, position FROM en57.append_events($1::en57.event[], $2::jsonb)", [expected_events, "{}"], ], ) @@ -101,9 +101,9 @@ module En57 connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, - success_result, + [{ "status" => "success", "position" => nil }], [ - "SELECT status FROM en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT status, position FROM en57.append_events($1::en57.event[], $2::jsonb)", [ array_encoder.encode([]), '{"fail_if_events_match":[{"types":["OrderPlaced"],"after":42}]}', @@ -112,21 +112,24 @@ module En57 ) connection.expect(:exec, nil, ["COMMIT"]) - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).append( - [], - fail_if: - Query.new( - criteria: [ - Query::Criteria.new( - types: ["OrderPlaced"], - tags: [], - after: 42, - ), - ], - ), + assert_equal( + Success.new(position: nil), + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append( + [], + fail_if: + Query.new( + criteria: [ + Query::Criteria.new( + types: ["OrderPlaced"], + tags: [], + after: 42, + ), + ], + ), + ), ) end end @@ -136,7 +139,7 @@ module En57 connection.expect(:exec, nil, ["BEGIN"]) connection.expect(:exec_params, nil) do |sql, params| assert_equal( - "SELECT status FROM en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT status, position FROM en57.append_events($1::en57.event[], $2::jsonb)", sql, ) assert_equal([array_encoder.encode([]), "{}"], params) @@ -541,7 +544,7 @@ module En57 connection.expect(:exec, nil, ["COMMIT"]) assert_equal( - Success.new, + Success.new(position: 1), Repository.new( PgAdapter.for_connection(connection), JsonSerializer.new, @@ -619,13 +622,13 @@ module En57 def record_encoder = @record_encoder ||= PG::TextEncoder::Record.new - def success_result = [{ "status" => "success" }] + def success_result = [{ "status" => "success", "position" => "1" }] def failure_result = [{ "status" => "append_condition_violated" }] def append_args [ - "SELECT status FROM en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT status, position FROM en57.append_events($1::en57.event[], $2::jsonb)", [ array_encoder.encode([]), '{"fail_if_events_match":[{"types":["OrderPlaced"]}]}',