diff --git a/README.md b/README.md index 645589a..e68223d 100644 --- a/README.md +++ b/README.md @@ -57,8 +57,11 @@ event_store = En57::EventStore.for_active_record ### Append events unconditionally +`append` returns a result. On success it carries the position of the last +appended event; on failure the position is `nil`. + ```ruby -event_store.append( +result = event_store.append( [ En57::Event.new( type: "OrderPlaced", @@ -67,6 +70,20 @@ event_store.append( ), ], ) + +result.success? # => true +result.position # => 1 +``` + +Or with pattern matching: + +```ruby +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" +end ``` ### Read all events @@ -109,18 +126,18 @@ Example: consume credits only once per account. ```ruby account_scope = event_store.read.with_tag("account:x") -begin - event_store.append( - [ - En57::Event.new( - type: "CreditsUsed", - data: { amount: 100 }, - tags: ["account:x"], - ), - ], - fail_if: account_scope.of_type("CreditsUsed"), - ) -rescue En57::AppendConditionViolated +result = event_store.append( + [ + En57::Event.new( + type: "CreditsUsed", + data: { amount: 100 }, + tags: ["account:x"], + ), + ], + fail_if: account_scope.of_type("CreditsUsed"), +) + +if result.failure? # lost the race; another writer already consumed credits end ``` @@ -144,18 +161,18 @@ Example: ensure no event exists with this email tag before writing. ```ruby email_tag = "email:alice@example.com" -begin - event_store.append( - [ - En57::Event.new( - type: "UserRegistered", - data: { name: "Alice" }, - tags: [email_tag], - ), - ], - fail_if: event_store.read.with_tag(email_tag), - ) -rescue En57::AppendConditionViolated +result = event_store.append( + [ + En57::Event.new( + type: "UserRegistered", + data: { name: "Alice" }, + tags: [email_tag], + ), + ], + fail_if: event_store.read.with_tag(email_tag), +) + +if result.failure? # email already used end ``` diff --git a/db/schema/0.1.0.sql b/db/schema/0.1.0.sql index dea0716..e2de6a0 100644 --- a/db/schema/0.1.0.sql +++ b/db/schema/0.1.0.sql @@ -27,13 +27,15 @@ CREATE TYPE en57.event AS ( ); CREATE FUNCTION en57.append_events (new_events en57.event[], append_condition jsonb DEFAULT '{}'::jsonb) - RETURNS void + RETURNS TABLE ("position" bigint) LANGUAGE plpgsql AS $$ +#variable_conflict use_column DECLARE criteria jsonb[] := ARRAY ( SELECT jsonb_array_elements(COALESCE(append_condition -> 'fail_if_events_match', '[]'::jsonb))); + last_position bigint; BEGIN IF cardinality(criteria) > 0 AND EXISTS ( SELECT @@ -63,14 +65,18 @@ WHERE t.event_id = e.id AND t.value = req.value)))) THEN RAISE EXCEPTION 'append_condition_violated'; END IF; - INSERT INTO en57.events (id, type, data, meta) - SELECT - e.id, - e.type, - e.data, - e.meta - FROM - unnest(new_events) AS e; + WITH inserted AS ( + INSERT INTO en57.events (id, type, data, meta) + SELECT + e.id, + e.type, + e.data, + e.meta + FROM + unnest(new_events) AS e + RETURNING position + ) + SELECT max(position) INTO last_position FROM inserted; INSERT INTO en57.tags (event_id, value) SELECT e.id, @@ -78,6 +84,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; END; $$; diff --git a/lib/en57.rb b/lib/en57.rb index 42ff441..fe2b914 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/result" require_relative "en57/pg_adapter" require_relative "en57/sequel_adapter" if defined?(Sequel) require_relative "en57/active_record_adapter" if defined?(ActiveRecord) @@ -14,8 +15,6 @@ require_relative "en57/event_store" require_relative "en57/configuration" module En57 - AppendConditionViolated = Class.new(StandardError) - def self.configuration = Configuration.instance def self.configure diff --git a/lib/en57/event_store.rb b/lib/en57/event_store.rb index b29238e..d475f39 100644 --- a/lib/en57/event_store.rb +++ b/lib/en57/event_store.rb @@ -8,7 +8,6 @@ module En57 def append(events, fail_if: EmptyScope.new) @repository.append(events, fail_if: fail_if.to_query) - self end def read diff --git a/lib/en57/pg_adapter.rb b/lib/en57/pg_adapter.rb index 66a4cf1..07c2bf5 100644 --- a/lib/en57/pg_adapter.rb +++ b/lib/en57/pg_adapter.rb @@ -17,8 +17,9 @@ module En57 def with_serializable_transaction with_connection do |connection| connection.exec("BEGIN ISOLATION LEVEL SERIALIZABLE") - yield connection + result = yield connection connection.exec("COMMIT") + result rescue StandardError connection.exec("ROLLBACK") raise diff --git a/lib/en57/repository.rb b/lib/en57/repository.rb index dddf761..4d1c1d1 100644 --- a/lib/en57/repository.rb +++ b/lib/en57/repository.rb @@ -32,17 +32,20 @@ module En57 :fail_if_events_match ] = fail_if_events_match unless fail_if_events_match.empty? - @adapter.with_serializable_transaction do |connection| - connection.exec_params( - "SELECT en57.append_events($1::en57.event[], $2::jsonb)", - [ - @array_encoder.encode(event_records), - JSON.generate(append_condition), - ], - ) - end + 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 - raise AppendConditionViolated + Result.failure(position: nil) end def read(query) diff --git a/lib/en57/result.rb b/lib/en57/result.rb new file mode 100644 index 0000000..6b73232 --- /dev/null +++ b/lib/en57/result.rb @@ -0,0 +1,21 @@ +# frozen_string_literal: true + +module En57 + module Result + Success = + Data.define(:position) do + def success? = true + def failure? = false + end + + Failure = + Data.define(:position) do + def success? = false + def failure? = true + end + + def self.success(position:) = Success.new(position:) + + def self.failure(position:) = Failure.new(position:) + end +end diff --git a/test/test_event_store.rb b/test/test_event_store.rb index 0270029..43fdcd2 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, nil, [[event]], fail_if: Query.all) + repository.expect( + :append, + Result.success(position: 1), + [[event]], + fail_if: Query.all, + ) EventStore.new(repository).append([event]) end @@ -29,15 +34,14 @@ module En57 end end - def test_return_self_from_append + def test_append_returns_repository_result event = Event.new(type: "CreditsToppedUp") + success = Result.success(position: 1) with_repository do |repository| - repository.expect(:append, nil, [[event]], fail_if: Query.all) + repository.expect(:append, success, [[event]], fail_if: Query.all) - event_store = EventStore.new(repository) - - assert_equal(event_store, event_store.append([event])) + assert_equal(success, EventStore.new(repository).append([event])) end end @@ -47,7 +51,12 @@ module En57 with_repository do |repository| event_store = EventStore.new(repository) fail_if = event_store.read.with_tag("order_id:123") - repository.expect(:append, nil, [[event]], fail_if: fail_if.to_query) + repository.expect( + :append, + Result.success(position: 1), + [[event]], + fail_if: fail_if.to_query, + ) event_store.append([event], fail_if:) end diff --git a/test/test_factories.rb b/test/test_factories.rb index d7f0e6d..3fbda63 100644 --- a/test/test_factories.rb +++ b/test/test_factories.rb @@ -43,8 +43,9 @@ module En57 def assert_round_trip(event_store) event = Event.new(type: "FactoryTested") + event_store.append([event]) - assert_equal [event], event_store.append([event]).read.each.to_a + assert_equal [event], event_store.read.each.to_a end end end diff --git a/test/test_integration.rb b/test/test_integration.rb index 2fea4df..4f5f956 100644 --- a/test/test_integration.rb +++ b/test/test_integration.rb @@ -24,7 +24,8 @@ module En57 ), ] - assert_equal(events, event_store.append(events).read.each.to_a) + assert_equal(Result.success(position: 2), event_store.append(events)) + assert_equal(events, event_store.read.each.to_a) end end @@ -34,10 +35,11 @@ module En57 Event.new(id: ids[0], type: "OrderPlaced"), Event.new(id: ids[1], type: "PriceChanged"), ] + event_store.append(events) assert_equal( events.map.with_index(1) { |event, position| [event, position] }, - event_store.append(events).read.each_with_position.to_a, + event_store.read.each_with_position.to_a, ) end end @@ -54,18 +56,18 @@ module En57 end end - define_method "test_#{name}_append_with_fail_if_and_matches_raises_append_condition_violated" do + 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") event_store.append([existing_event]) - assert_raises(AppendConditionViolated) do + result = event_store.append( [Event.new(id: ids[1], type: "ShipmentScheduled")], fail_if: event_store.read.of_type("OrderPlaced"), ) - end + assert(result.failure?) assert_equal([existing_event], event_store.read.each.to_a) end end @@ -86,18 +88,18 @@ module En57 end end - define_method "test_#{name}_append_with_after_raises_if_match_is_after_cutoff" do + 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") event_store.append([existing_event]) - assert_raises(AppendConditionViolated) do + result = event_store.append( [Event.new(id: ids[1], type: "ShipmentScheduled")], fail_if: event_store.read.of_type("OrderPlaced").after(0), ) - end + assert(result.failure?) assert_equal([existing_event], event_store.read.each.to_a) end end @@ -122,7 +124,9 @@ module En57 event = Event.new(id: ids[0], type: "OrderPlaced", tags: ["order_id:123"]) - assert_equal([event], event_store.append([event]).read.each.to_a) + event_store.append([event]) + + assert_equal([event], event_store.read.each.to_a) end end @@ -132,11 +136,9 @@ module En57 Event.new(id: ids[0], type: "OrderPlaced"), Event.new(id: ids[1], type: "PriceChanged"), ] + event_store.append(events) - assert_equal( - events.drop(1), - event_store.append(events).read.after(1).each.to_a, - ) + assert_equal(events.drop(1), event_store.read.after(1).each.to_a) end end @@ -154,11 +156,11 @@ module En57 tags: %w[order_id:456 tenant_id:acme], ), ] + event_store.append(events) assert_equal( events.take(1), event_store - .append(events) .read .with_tag("order_id:123", "tenant_id:acme") .each @@ -173,10 +175,11 @@ module En57 Event.new(id: ids[0], type: "OrderPlaced"), Event.new(id: ids[1], type: "PriceChanged"), ] + event_store.append(events) assert_equal( events.take(1), - event_store.append(events).read.of_type("OrderPlaced").each.to_a, + event_store.read.of_type("OrderPlaced").each.to_a, ) end end @@ -188,15 +191,11 @@ module En57 Event.new(id: ids[1], type: "OrderPlaced"), Event.new(id: ids[2], type: "OrderCancelled"), ] + event_store.append(events) assert_equal( events.drop(1), - event_store - .append(events) - .read - .of_type("OrderPlaced", "OrderCancelled") - .each - .to_a, + event_store.read.of_type("OrderPlaced", "OrderCancelled").each.to_a, ) end end @@ -208,11 +207,11 @@ module En57 Event.new(id: ids[1], type: "OrderPlaced", tags: ["order_id:456"]), Event.new(id: ids[2], type: "PriceChanged", tags: ["order_id:123"]), ] + event_store.append(events) assert_equal( events.take(1), event_store - .append(events) .read .of_type("OrderPlaced") .with_tag("order_id:123") diff --git a/test/test_migrator.rb b/test/test_migrator.rb index 952057d..ca64d76 100644 --- a/test/test_migrator.rb +++ b/test/test_migrator.rb @@ -55,7 +55,9 @@ module En57 ), ) - assert_equal [event], event_store.append([event]).read.each.to_a + event_store.append([event]) + + assert_equal [event], event_store.read.each.to_a ensure connection&.close end diff --git a/test/test_pg_adapter.rb b/test/test_pg_adapter.rb index 495dc18..0538388 100644 --- a/test/test_pg_adapter.rb +++ b/test/test_pg_adapter.rb @@ -60,10 +60,12 @@ module En57 ) connection.expect(:exec, nil, ["COMMIT"]) - adapter.with_serializable_transaction do |conn| - assert_equal :written, - conn.exec_params("SELECT en57.append_events()", []) - end + assert_equal( + :written, + adapter.with_serializable_transaction do |conn| + conn.exec_params("SELECT en57.append_events()", []) + end, + ) end end diff --git a/test/test_repository.rb b/test/test_repository.rb index c50828b..f301327 100644 --- a/test/test_repository.rb +++ b/test/test_repository.rb @@ -34,38 +34,41 @@ module En57 connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, - nil, + [{ "position" => "2" }], [ - "SELECT en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT position FROM en57.append_events($1::en57.event[], $2::jsonb)", [expected_events, "{}"], ], ) connection.expect(:exec, nil, ["COMMIT"]) - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).append( - [ - Event.new( - id: ids[0], - type: "CreditsToppedUp", - data: { - amount: 100, - }, - tags: ["order_id:123"], - ), - Event.new( - id: ids[1], - type: "CreditsToppedUp", - data: { - amount: 50, - }, - tags: ["order_id:234"], - ), - ], - fail_if: Query.all, - ) + result = + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append( + [ + Event.new( + id: ids[0], + type: "CreditsToppedUp", + data: { + amount: 100, + }, + tags: ["order_id:123"], + ), + Event.new( + id: ids[1], + type: "CreditsToppedUp", + data: { + amount: 50, + }, + tags: ["order_id:234"], + ), + ], + fail_if: Query.all, + ) + + assert_equal(Result.success(position: 2), result) end end @@ -78,21 +81,24 @@ module En57 connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, - nil, + [{ "position" => "1" }], [ - "SELECT en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT position FROM en57.append_events($1::en57.event[], $2::jsonb)", [expected_events, "{}"], ], ) connection.expect(:exec, nil, ["COMMIT"]) - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).append( - [Event.new(id: ids[0], type: "OrderPlaced")], - fail_if: Query.all, - ) + result = + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append( + [Event.new(id: ids[0], type: "OrderPlaced")], + fail_if: Query.all, + ) + + assert_equal(Result.success(position: 1), result) end end @@ -101,9 +107,9 @@ module En57 connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect( :exec_params, - nil, + [{ "position" => nil }], [ - "SELECT en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT position FROM en57.append_events($1::en57.event[], $2::jsonb)", [ array_encoder.encode([]), '{"fail_if_events_match":[{"types":["OrderPlaced"],"after":42}]}', @@ -112,22 +118,25 @@ 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, - ), - ], - ), - ) + result = + Repository.new( + PgAdapter.for_connection(connection), + JsonSerializer.new, + ).append( + [], + fail_if: + Query.new( + criteria: [ + Query::Criteria.new( + types: ["OrderPlaced"], + tags: [], + after: 42, + ), + ], + ), + ) + + assert_equal(Result.success(position: nil), result) end end @@ -137,7 +146,7 @@ module En57 connection.expect(:exec, nil, ["ROLLBACK"]) connection.expect(:exec_params, nil) do |sql, params| assert_equal( - "SELECT en57.append_events($1::en57.event[], $2::jsonb)", + "SELECT position FROM en57.append_events($1::en57.event[], $2::jsonb)", sql, ) assert_equal([array_encoder.encode([]), "{}"], params) @@ -513,22 +522,23 @@ module En57 end end - def test_append_raises_append_condition_violated_from_pg_error_sqlstate + def test_append_returns_failure_on_pg_raise_exception 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) } - assert_raises(AppendConditionViolated) do + result = Repository.new( PgAdapter.for_connection(connection), JsonSerializer.new, ).append([], fail_if: Query.all) - end + + assert_equal(Result.failure(position: nil), result) end end - def test_append_raises_append_condition_violated_from_serialization_failure_result_sqlstate + def test_append_returns_failure_on_serialization_failure with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) connection.expect(:exec, nil, ["ROLLBACK"]) @@ -536,12 +546,13 @@ module En57 raise PG::TRSerializationFailure.new end - assert_raises(AppendConditionViolated) do + result = Repository.new( PgAdapter.for_connection(connection), JsonSerializer.new, ).append([], fail_if: Query.all) - end + + assert_equal(Result.failure(position: nil), result) end end diff --git a/test/test_result.rb b/test/test_result.rb new file mode 100644 index 0000000..3834e84 --- /dev/null +++ b/test/test_result.rb @@ -0,0 +1,25 @@ +# frozen_string_literal: true + +require "test_helper" + +module En57 + class TestResult < Minitest::Test + cover Result + + def test_success_carries_position + result = Result.success(position: 42) + + assert(result.success?) + refute(result.failure?) + assert_equal(42, result.position) + end + + def test_failure_carries_position + result = Result.failure(position: 7) + + assert(result.failure?) + refute(result.success?) + assert_equal(7, result.position) + end + end +end diff --git a/test/test_stress.rb b/test/test_stress.rb index ff35966..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 AppendConditionViolated => e - e end end assert_equal( (concurrency - 1), - threads.map(&:value).select { AppendConditionViolated === it }.size, + threads.map(&:value).count(&:failure?), ) assert_equal( 1,