diff --git a/README.md b/README.md index 28ec53a..eddc607 100644 --- a/README.md +++ b/README.md @@ -70,6 +70,12 @@ store.append( events = store.read.each.to_a ``` +### Read events with positions + +```ruby +event, position = store.read.each_with_position.first +``` + ### Read events filtered by tags ```ruby diff --git a/db/schema.sql b/db/schema.sql index 1a7aa7d..08db2a5 100644 --- a/db/schema.sql +++ b/db/schema.sql @@ -81,7 +81,13 @@ END; $$; CREATE FUNCTION read_events (criteria jsonb[]) - RETURNS SETOF event_with_tags + RETURNS TABLE ( + "position" bigint, + id uuid, + type text, + data jsonb, + meta jsonb, + tags text[]) LANGUAGE SQL AS $$ WITH parsed_criteria AS ( @@ -130,6 +136,7 @@ filtered_events AS ( t.event_id = e.id AND t.value = req.value)))) SELECT + e.position, e.id, e.type, e.data, diff --git a/lib/en57/repository.rb b/lib/en57/repository.rb index fb8f022..aa6ce22 100644 --- a/lib/en57/repository.rb +++ b/lib/en57/repository.rb @@ -61,17 +61,20 @@ module En57 @adapter .with_connection do |connection| connection.exec_params( - "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", + "SELECT position, id, type, data, meta, tags FROM read_events($1::jsonb[])", [@array_encoder.encode(criteria)], ) end .map do |row| - Event.new( - id: row.fetch("id"), - type: row.fetch("type"), - data: @serializer.load(row.fetch("data"), row.fetch("meta")), - tags: @array_decoder.decode(row.fetch("tags")), - ) + [ + Event.new( + id: row.fetch("id"), + type: row.fetch("type"), + data: @serializer.load(row.fetch("data"), row.fetch("meta")), + tags: @array_decoder.decode(row.fetch("tags")), + ), + Integer(row.fetch("position")), + ] end end end diff --git a/lib/en57/scope.rb b/lib/en57/scope.rb index 2e365c3..7abe981 100644 --- a/lib/en57/scope.rb +++ b/lib/en57/scope.rb @@ -14,7 +14,15 @@ module En57 def each(&block) return enum_for unless block - @repository.read(@query).each(&block) + @repository.read(@query).each { |event, _position| yield event } + end + + def each_with_position(&block) + return enum_for(__method__) unless block + + @repository.read(@query).each do |event, position| + yield event, position + end end def to_query = @query @@ -34,7 +42,15 @@ module En57 def each(&block) return enum_for unless block - @repository.read(@query).each(&block) + @repository.read(@query).each { |event, _position| yield event } + end + + def each_with_position(&block) + return enum_for(__method__) unless block + + @repository.read(@query).each do |event, position| + yield event, position + end end def to_query = @query diff --git a/test/pg_regress/expected/001_schema.out b/test/pg_regress/expected/001_schema.out index 50f4172..7f80a17 100644 --- a/test/pg_regress/expected/001_schema.out +++ b/test/pg_regress/expected/001_schema.out @@ -74,7 +74,14 @@ WHERE END; $$; CREATE FUNCTION read_events (criteria jsonb[]) - RETURNS SETOF event_with_tags + RETURNS TABLE ( + "position" bigint, + id uuid, + type text, + data jsonb, + meta jsonb, + tags text[] + ) LANGUAGE SQL AS $$ WITH parsed_criteria AS ( @@ -122,6 +129,7 @@ filtered_events AS ( t.event_id = e.id AND t.value = req.value)))) SELECT + e.position, e.id, e.type, e.data, diff --git a/test/test_event_store.rb b/test/test_event_store.rb index 43315bd..0270029 100644 --- a/test/test_event_store.rb +++ b/test/test_event_store.rb @@ -20,7 +20,7 @@ module En57 event = Event.new(type: "CreditsToppedUp") with_repository do |repository| - repository.expect(:read, [event], [Query.all]) + repository.expect(:read, [[event, 1]], [Query.all]) result = EventStore.new(repository).read diff --git a/test/test_integration.rb b/test/test_integration.rb index 621b665..2fea4df 100644 --- a/test/test_integration.rb +++ b/test/test_integration.rb @@ -28,6 +28,20 @@ module En57 end end + define_method "test_#{name}_read_with_position_yields_events_and_positions" do + with_event_store(factory) do |event_store| + events = [ + Event.new(id: ids[0], type: "OrderPlaced"), + Event.new(id: ids[1], type: "PriceChanged"), + ] + + assert_equal( + events.map.with_index(1) { |event, position| [event, position] }, + event_store.append(events).read.each_with_position.to_a, + ) + end + end + define_method "test_#{name}_append_with_fail_if_and_no_matches_appends_events" do with_event_store(factory) do |event_store| event = Event.new(id: ids[0], type: "OrderPlaced") diff --git a/test/test_repository.rb b/test/test_repository.rb index 206dd15..2a4bf22 100644 --- a/test/test_repository.rb +++ b/test/test_repository.rb @@ -232,6 +232,7 @@ module En57 :exec_params, [ { + "position" => "1", "id" => ids[0], "type" => "CreditsToppedUp", "data" => '{"amount":100}', @@ -239,6 +240,7 @@ module En57 "tags" => "{order_id:123}", }, { + "position" => "2", "id" => ids[1], "type" => "CreditsToppedUp", "data" => '{"amount":50}', @@ -247,29 +249,35 @@ module En57 }, ], [ - "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", + "SELECT position, id, type, data, meta, tags FROM read_events($1::jsonb[])", [array_encoder.encode([])], ], ) assert_equal( [ - 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"], - ), + [ + Event.new( + id: ids[0], + type: "CreditsToppedUp", + data: { + "amount" => 100, + }, + tags: ["order_id:123"], + ), + 1, + ], + [ + Event.new( + id: ids[1], + type: "CreditsToppedUp", + data: { + "amount" => 50, + }, + tags: ["order_id:234"], + ), + 2, + ], ], Repository.new( PgAdapter.new(connection_uri), @@ -289,6 +297,7 @@ module En57 :exec_params, [ { + "position" => "1", "id" => ids[0], "type" => "CreditsToppedUp", "data" => '{"amount":100}', @@ -297,21 +306,24 @@ module En57 }, ], [ - "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", + "SELECT position, id, type, data, meta, tags FROM read_events($1::jsonb[])", [array_encoder.encode(['{"tags":["order_id:123"]}'])], ], ) assert_equal( [ - Event.new( - id: ids[0], - type: "CreditsToppedUp", - data: { - "amount" => 100, - }, - tags: ["order_id:123"], - ), + [ + Event.new( + id: ids[0], + type: "CreditsToppedUp", + data: { + "amount" => 100, + }, + tags: ["order_id:123"], + ), + 1, + ], ], Repository.new( PgAdapter.new(connection_uri), @@ -328,6 +340,7 @@ module En57 :exec_params, [ { + "position" => "1", "id" => ids[0], "type" => "CreditsToppedUp", "data" => '{"amount":100}', @@ -336,21 +349,24 @@ module En57 }, ], [ - "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", + "SELECT position, id, type, data, meta, tags FROM read_events($1::jsonb[])", [array_encoder.encode(["{}"])], ], ) assert_equal( [ - Event.new( - id: ids[0], - type: "CreditsToppedUp", - data: { - "amount" => 100, - }, - tags: ["order_id:123"], - ), + [ + Event.new( + id: ids[0], + type: "CreditsToppedUp", + data: { + "amount" => 100, + }, + tags: ["order_id:123"], + ), + 1, + ], ], Repository.new( PgAdapter.new(connection_uri), @@ -373,6 +389,7 @@ module En57 :exec_params, [ { + "position" => "1", "id" => ids[0], "type" => "CreditsToppedUp", "data" => '{"amount":100}', @@ -381,7 +398,7 @@ module En57 }, ], [ - "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", + "SELECT position, id, type, data, meta, tags FROM read_events($1::jsonb[])", [ array_encoder.encode( %w[{"tags":["order_id:123"]} {"tags":["order_id:456"]}], @@ -392,14 +409,17 @@ module En57 assert_equal( [ - Event.new( - id: ids[0], - type: "CreditsToppedUp", - data: { - "amount" => 100, - }, - tags: ["order_id:123"], - ), + [ + Event.new( + id: ids[0], + type: "CreditsToppedUp", + data: { + "amount" => 100, + }, + tags: ["order_id:123"], + ), + 1, + ], ], Repository.new( PgAdapter.new(connection_uri), @@ -419,7 +439,7 @@ module En57 :exec_params, [], [ - "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", + "SELECT position, id, type, data, meta, tags FROM read_events($1::jsonb[])", [array_encoder.encode(['{"after":42}'])], ], ) @@ -444,6 +464,7 @@ module En57 :exec_params, [ { + "position" => "1", "id" => ids[0], "type" => "OrderPlaced", "data" => '{"amount":100}', @@ -452,21 +473,24 @@ module En57 }, ], [ - "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", + "SELECT position, id, type, data, meta, tags FROM read_events($1::jsonb[])", [array_encoder.encode(['{"types":["OrderPlaced"]}'])], ], ) assert_equal( [ - Event.new( - id: ids[0], - type: "OrderPlaced", - data: { - "amount" => 100, - }, - tags: [], - ), + [ + Event.new( + id: ids[0], + type: "OrderPlaced", + data: { + "amount" => 100, + }, + tags: [], + ), + 1, + ], ], Repository.new( PgAdapter.new(connection_uri), diff --git a/test/test_scope.rb b/test/test_scope.rb index ab8ab9e..cbefc8b 100644 --- a/test/test_scope.rb +++ b/test/test_scope.rb @@ -16,15 +16,39 @@ module En57 end end + def test_each_with_position_without_block_returns_enumerator + with_repository do |repository| + repository.expect(:read, [[:event_1, 1], [:event_2, 2]], [Query.all]) + + assert_equal( + [[:event_1, 1], [:event_2, 2]], + Scope.new(repository, Query.all).each_with_position.to_a, + ) + end + end + def test_each_with_block_yields_events with_repository do |repository| scope = Scope.new(repository, Query.all) yielded = [] - repository.expect(:read, [1, 2], [Query.all]) + repository.expect(:read, [[:event_1, 1], [:event_2, 2]], [Query.all]) scope.each { |event| yielded << event } - assert_equal([1, 2], yielded) + assert_equal(%i[event_1 event_2], yielded) + end + end + + def test_each_with_position_yields_events_and_positions + with_repository do |repository| + yielded = [] + repository.expect(:read, [[:event_1, 1], [:event_2, 2]], [Query.all]) + + Scope + .new(repository, Query.all) + .each_with_position { |event, position| yielded << [event, position] } + + assert_equal([[:event_1, 1], [:event_2, 2]], yielded) end end @@ -110,13 +134,29 @@ module En57 def test_merged_scope_each_without_block_returns_enumerator with_repository do |repository| - merged = - Scope - .new(repository, Query.all) - .of_type("OrderPlaced") - .or(Scope.new(repository, Query.all).with_tag("order_id:123")) + assert_instance_of(Enumerator, merged_scope(repository).each) + end + end - assert_instance_of(Enumerator, merged.each) + def test_merged_scope_each_with_position_without_block_returns_enumerator + with_repository do |repository| + repository.expect( + :read, + [[:event_1, 1], [:event_2, 2]], + [ + Query.new( + criteria: [ + Query::Criteria.new(types: ["OrderPlaced"], tags: []), + Query::Criteria.new(types: [], tags: ["order_id:123"]), + ], + ), + ], + ) + + assert_equal( + [[:event_1, 1], [:event_2, 2]], + merged_scope(repository).each_with_position.to_a, + ) end end @@ -131,7 +171,7 @@ module En57 repository.expect( :read, - [1, 2], + [[:event_1, 1], [:event_2, 2]], [ Query.new( criteria: [ @@ -144,7 +184,32 @@ module En57 merged.each { |event| yielded << event } - assert_equal([1, 2], yielded) + assert_equal(%i[event_1 event_2], yielded) + end + end + + def test_merged_scope_each_with_position_yields_events_and_positions + with_repository do |repository| + yielded = [] + + repository.expect( + :read, + [[:event_1, 1], [:event_2, 2]], + [ + Query.new( + criteria: [ + Query::Criteria.new(types: ["OrderPlaced"], tags: []), + Query::Criteria.new(types: [], tags: ["order_id:123"]), + ], + ), + ], + ) + + merged_scope(repository).each_with_position do |event, position| + yielded << [event, position] + end + + assert_equal([[:event_1, 1], [:event_2, 2]], yielded) end end @@ -252,6 +317,13 @@ module En57 private + def merged_scope(repository) + Scope + .new(repository, Query.all) + .of_type("OrderPlaced") + .or(Scope.new(repository, Query.all).with_tag("order_id:123")) + end + def with_repository repository = Minitest::Mock.new yield repository