diff --git a/db/schema.sql b/db/schema.sql index 5e8287f..5e2e727 100644 --- a/db/schema.sql +++ b/db/schema.sql @@ -32,14 +32,14 @@ DECLARE criteria jsonb[] := ARRAY ( SELECT jsonb_array_elements(COALESCE(append_condition -> 'fail_if_events_match', '[]'::jsonb))); - after_pos bigint := (append_condition ->> 'after')::bigint; + AFTER bigint := (append_condition ->> 'after')::bigint; BEGIN IF cardinality(criteria) > 0 AND EXISTS ( SELECT 1 FROM events AS e - WHERE (after_pos IS NULL OR e.position > after_pos) AND EXISTS ( + WHERE (AFTER IS NULL OR e.position > AFTER) AND EXISTS ( SELECT 1 FROM @@ -88,6 +88,7 @@ CREATE FUNCTION read_events (criteria jsonb[]) SELECT c, c -> 'tags' AS tags, + (c ->> 'after')::bigint AS after, ARRAY ( SELECT jsonb_array_elements_text(c -> 'types')) AS types @@ -110,7 +111,9 @@ filtered_events AS ( 1 FROM parsed_criteria AS pc - WHERE (pc.c -> 'types' IS NULL + WHERE (pc.after IS NULL + OR e.position > pc.after) + AND (pc.c -> 'types' IS NULL OR e.type = ANY (pc.types)) AND NOT EXISTS ( SELECT diff --git a/lib/en57/query.rb b/lib/en57/query.rb index f421235..f7f9613 100644 --- a/lib/en57/query.rb +++ b/lib/en57/query.rb @@ -3,7 +3,9 @@ module En57 class Query < Data.define(:criteria) Criteria = - Data.define(:types, :tags) do + Data.define(:types, :tags, :after) do + def initialize(types:, tags:, after: nil) = super + def self.all = new(types: [], tags: []) def with_tags(tags) @@ -14,8 +16,14 @@ module En57 with(types: [*self.types, *types].uniq) end + def with_after(position) + with(after: position) + end + def matcher - { types:, tags: }.reject { |_, value| value.empty? } + { types:, tags:, after: }.reject do |key, value| + key == :after ? value.nil? : value.empty? + end end end diff --git a/lib/en57/scope.rb b/lib/en57/scope.rb index ec3d1ec..2e365c3 100644 --- a/lib/en57/scope.rb +++ b/lib/en57/scope.rb @@ -53,6 +53,13 @@ module En57 ) end + def after(position) + self.class.new( + @repository, + @query.refine_last { |item| item.with_after(position) }, + ) + end + def or(other) MergedScope.new(repository: @repository, query: @query.or(other.to_query)) end diff --git a/test/pg_regress/expected/001_schema.out b/test/pg_regress/expected/001_schema.out index 7356096..f7025e6 100644 --- a/test/pg_regress/expected/001_schema.out +++ b/test/pg_regress/expected/001_schema.out @@ -27,14 +27,14 @@ DECLARE criteria jsonb[] := ARRAY ( SELECT jsonb_array_elements(COALESCE(append_condition -> 'fail_if_events_match', '[]'::jsonb))); - after_pos bigint := (append_condition ->> 'after')::bigint; + AFTER bigint := (append_condition ->> 'after')::bigint; BEGIN IF cardinality(criteria) > 0 AND EXISTS ( SELECT 1 FROM events AS e - WHERE (after_pos IS NULL OR e.position > after_pos) AND EXISTS ( + WHERE (AFTER IS NULL OR e.position > AFTER) AND EXISTS ( SELECT 1 FROM @@ -82,6 +82,7 @@ CREATE FUNCTION read_events (criteria jsonb[]) SELECT c, c -> 'tags' AS tags, + (c ->> 'after')::bigint AS after, ARRAY ( SELECT jsonb_array_elements_text(c -> 'types')) AS types @@ -104,7 +105,8 @@ filtered_events AS ( 1 FROM parsed_criteria AS pc - WHERE (pc.c -> 'types' IS NULL + WHERE (pc.after IS NULL OR e.position > pc.after) + AND (pc.c -> 'types' IS NULL OR e.type = ANY (pc.types)) AND NOT EXISTS ( SELECT diff --git a/test/test_integration.rb b/test/test_integration.rb index 57d07d5..db04b21 100644 --- a/test/test_integration.rb +++ b/test/test_integration.rb @@ -86,6 +86,20 @@ module En57 end end + def test_read_filters_after + with_event_store do |event_store| + events = [ + Event.new(id: ids[0], type: "OrderPlaced"), + Event.new(id: ids[1], type: "PriceChanged"), + ] + + assert_equal( + events.drop(1), + event_store.append(events).read.after(1).each.to_a, + ) + end + end + def test_read_filters_by_tags with_event_store do |event_store| events = [ diff --git a/test/test_pg_repository.rb b/test/test_pg_repository.rb index 30854b5..1e7a603 100644 --- a/test/test_pg_repository.rb +++ b/test/test_pg_repository.rb @@ -358,6 +358,28 @@ module En57 end end + def test_read_events_filtered_by_after + query = + Query.new( + criteria: [Query::Criteria.new(types: [], tags: [], after: 42)], + ) + with_connection_to(connection_uri) do |connection| + connection.expect( + :exec_params, + [], + [ + "SELECT id, type, data, meta, tags FROM read_events($1::jsonb[])", + [array_encoder.encode(['{"after":42}'])], + ], + ) + + assert_equal( + [], + PgRepository.new(connection_uri, JsonSerializer.new).read(query), + ) + end + end + def test_read_events_filtered_by_type query = Query.new( diff --git a/test/test_query.rb b/test/test_query.rb index 1a8888b..ee4b13d 100644 --- a/test/test_query.rb +++ b/test/test_query.rb @@ -69,12 +69,18 @@ module En57 criteria: [ Query::Criteria.new(types: ["OrderPlaced"], tags: []), Query::Criteria.new(types: [], tags: ["order_id:123"]), + Query::Criteria.new(types: [], tags: [], after: 42), Query::Criteria.new(types: [], tags: []), ], ) assert_equal( - [{ types: ["OrderPlaced"] }, { tags: ["order_id:123"] }, {}], + [ + { types: ["OrderPlaced"] }, + { tags: ["order_id:123"] }, + { after: 42 }, + {}, + ], query.encoded_criteria, ) end diff --git a/test/test_scope.rb b/test/test_scope.rb index f7b1966..341afca 100644 --- a/test/test_scope.rb +++ b/test/test_scope.rb @@ -176,6 +176,23 @@ module En57 repository.verify end + def test_after_refines_query + repository = Minitest::Mock.new + + repository.expect( + :read, + [], + [ + Query.new( + criteria: [Query::Criteria.new(types: [], tags: [], after: 42)], + ), + ], + ) + + assert_equal([], Scope.new(repository, Query.all).after(42).each.to_a) + repository.verify + end + def test_merged_scope_cannot_be_refined_anymore repository = Minitest::Mock.new merged = @@ -186,6 +203,7 @@ module En57 assert_raises(NoMethodError) { merged.with_tag("customer_id:456") } assert_raises(NoMethodError) { merged.of_type("OrderCancelled") } + assert_raises(NoMethodError) { merged.after(42) } end def test_pipe_aliases_or