diff --git a/.gitignore b/.gitignore index 786a4ac..eeca105 100644 --- a/.gitignore +++ b/.gitignore @@ -2,3 +2,4 @@ /.mutant/ /pkg/ /tmp/ +/test/pg_regress/results/ diff --git a/db/schema/0.1.0.sql b/db/schema/0.1.0.sql index f1faae5..4e2ad57 100644 --- a/db/schema/0.1.0.sql +++ b/db/schema/0.1.0.sql @@ -150,7 +150,7 @@ INSERT INTO en57.events (id, type, data, meta) END; $$; -CREATE FUNCTION en57.read_events (criteria jsonb[]) +CREATE FUNCTION en57.read_events (criteria jsonb[], batch_size int DEFAULT NULL, after_position bigint DEFAULT NULL) RETURNS TABLE ( "position" bigint, id uuid, @@ -180,8 +180,9 @@ filtered_events AS ( e.meta FROM en57.events AS e - WHERE - cardinality(criteria) = 0 + WHERE (after_position IS NULL + OR e.position > after_position) + AND (cardinality(criteria) = 0 OR EXISTS ( SELECT 1 @@ -204,14 +205,14 @@ filtered_events AS ( en57.tags AS t WHERE t.event_id = e.id - AND t.value = req.value)))) - SELECT - e.position, - e.id, - e.type, - e.data, - e.meta, - COALESCE(t.tags, ARRAY[]::text[]) AS tags + AND t.value = req.value))))) +SELECT + e.position, + e.id, + e.type, + e.data, + e.meta, + COALESCE(t.tags, ARRAY[]::text[]) AS tags FROM filtered_events AS e LEFT JOIN LATERAL ( @@ -222,6 +223,7 @@ FROM WHERE t.event_id = e.id) AS t ON TRUE ORDER BY - e.position; + e.position +LIMIT batch_size; $$; diff --git a/lib/en57/configuration.rb b/lib/en57/configuration.rb index 00aa196..3fc4ddc 100644 --- a/lib/en57/configuration.rb +++ b/lib/en57/configuration.rb @@ -6,10 +6,11 @@ module En57 class Configuration include Singleton - attr_accessor :append_retries, :serializer + attr_accessor :append_retries, :read_batch_size, :serializer def initialize @append_retries = 9 + @read_batch_size = 1000 @serializer = JsonSerializer.new end end diff --git a/lib/en57/repository.rb b/lib/en57/repository.rb index 2bc88a5..8cfc428 100644 --- a/lib/en57/repository.rb +++ b/lib/en57/repository.rb @@ -79,24 +79,41 @@ module En57 def read(query) criteria = query.encoded_criteria.map { |item| JSON.generate(item) } + batch_size = En57.configuration.read_batch_size - @adapter - .with_connection do |connection| - connection.exec_params( - "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[])", - [@array_encoder.encode(criteria)], - ) - end - .map do |row| - [ - deserialize_event(row), - Integer(row.fetch("position")), - ] + Enumerator.new do |yielder| + cursor = nil + loop do + has_more = false + read_batch( + criteria, + (batch_size + 1 if batch_size), + cursor, + ).each_with_index do |row, index| + if index == batch_size + has_more = true + break + end + position = Integer(row.fetch("position")) + yielder << [deserialize_event(row), position] + cursor = position + end + break unless has_more end + end end private + def read_batch(criteria, limit, after_position) + @adapter.with_connection do |connection| + connection.exec_params( + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)", + [@array_encoder.encode(criteria), limit, after_position], + ) + end + end + def json_string(value) value.instance_of?(String) ? value : JSON.generate(value) end @@ -111,9 +128,11 @@ module En57 row.fetch("meta").then { json_string(it) if it }, ), tags: - row.fetch("tags").then do |tags| - tags.instance_of?(Array) ? tags : @array_decoder.decode(tags) - end, + row + .fetch("tags") + .then do |tags| + tags.instance_of?(Array) ? tags : @array_decoder.decode(tags) + end, ) end end diff --git a/test/pg_regress/expected/003_read_events_batching.out b/test/pg_regress/expected/003_read_events_batching.out new file mode 100644 index 0000000..5747c93 --- /dev/null +++ b/test/pg_regress/expected/003_read_events_batching.out @@ -0,0 +1,48 @@ +TRUNCATE TABLE en57.tags, en57.events RESTART IDENTITY CASCADE; +INSERT INTO en57.events (id, type, data, meta) +VALUES + ('00000000-0000-0000-0000-000000000001', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb), + ('00000000-0000-0000-0000-000000000002', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb), + ('00000000-0000-0000-0000-000000000003', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb), + ('00000000-0000-0000-0000-000000000004', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb), + ('00000000-0000-0000-0000-000000000005', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb); +-- batch_size caps the number of rows returned (first keyset page) +SELECT + position, + id::text +FROM + en57.read_events (ARRAY[]::jsonb[], 2) +ORDER BY + position; + position | id +----------+-------------------------------------- + 1 | 00000000-0000-0000-0000-000000000001 + 2 | 00000000-0000-0000-0000-000000000002 +(2 rows) + +-- after_position is a keyset cursor; combined with batch_size it yields the next page +SELECT + position, + id::text +FROM + en57.read_events (ARRAY[]::jsonb[], 2, 2) +ORDER BY + position; + position | id +----------+-------------------------------------- + 3 | 00000000-0000-0000-0000-000000000003 + 4 | 00000000-0000-0000-0000-000000000004 +(2 rows) + +-- a cursor at or past the last position returns no further rows +SELECT + position, + id::text +FROM + en57.read_events (ARRAY[]::jsonb[], 2, 5) +ORDER BY + position; + position | id +----------+---- +(0 rows) + diff --git a/test/pg_regress/schedule b/test/pg_regress/schedule index b8594a1..5a7582c 100644 --- a/test/pg_regress/schedule +++ b/test/pg_regress/schedule @@ -1,2 +1,3 @@ test: 001_schema test: 002_read_events_disjunction +test: 003_read_events_batching diff --git a/test/pg_regress/schedule_existing b/test/pg_regress/schedule_existing index d3175e3..b6e797c 100644 --- a/test/pg_regress/schedule_existing +++ b/test/pg_regress/schedule_existing @@ -1 +1,2 @@ test: 002_read_events_disjunction +test: 003_read_events_batching diff --git a/test/pg_regress/sql/003_read_events_batching.sql b/test/pg_regress/sql/003_read_events_batching.sql new file mode 100644 index 0000000..1ffb45f --- /dev/null +++ b/test/pg_regress/sql/003_read_events_batching.sql @@ -0,0 +1,36 @@ +TRUNCATE TABLE en57.tags, en57.events RESTART IDENTITY CASCADE; + +INSERT INTO en57.events (id, type, data, meta) +VALUES + ('00000000-0000-0000-0000-000000000001', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb), + ('00000000-0000-0000-0000-000000000002', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb), + ('00000000-0000-0000-0000-000000000003', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb), + ('00000000-0000-0000-0000-000000000004', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb), + ('00000000-0000-0000-0000-000000000005', 'OrderPlaced', '{}'::jsonb, '{}'::jsonb); + +-- batch_size caps the number of rows returned (first keyset page) +SELECT + position, + id::text +FROM + en57.read_events (ARRAY[]::jsonb[], 2) +ORDER BY + position; + +-- after_position is a keyset cursor; combined with batch_size it yields the next page +SELECT + position, + id::text +FROM + en57.read_events (ARRAY[]::jsonb[], 2, 2) +ORDER BY + position; + +-- a cursor at or past the last position returns no further rows +SELECT + position, + id::text +FROM + en57.read_events (ARRAY[]::jsonb[], 2, 5) +ORDER BY + position; diff --git a/test/test_integration.rb b/test/test_integration.rb index 56e6495..c0ab1e7 100644 --- a/test/test_integration.rb +++ b/test/test_integration.rb @@ -270,6 +270,33 @@ module En57 assert_equal(events.fetch_values(0, 2), (orders | prices).each.to_a) end end + + define_method "test_#{name}_read_streams_results_spanning_many_batches" do + with_event_store(factory) do |event_store| + events = + (1..5).map do |n| + Event.new( + id: ids[n], + type: "OrderPlaced", + tags: ["order_id:#{n}"], + ) + end + assert_equal(Success.new(position: 5), event_store.append(events)) + + En57 + .configuration + .stub(:read_batch_size, 2) do + assert_equal(events, event_store.read.each.to_a) + assert_equal( + events + .map + .with_index(1) { |event, position| [event, position] }, + event_store.read.each_with_position.to_a, + ) + assert_equal(events.drop(2), event_store.read.after(2).each.to_a) + end + end + end end private diff --git a/test/test_repository.rb b/test/test_repository.rb index 2e57035..b87bfb1 100644 --- a/test/test_repository.rb +++ b/test/test_repository.rb @@ -214,8 +214,8 @@ module En57 }, ], [ - "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[])", - [array_encoder.encode([])], + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)", + [array_encoder.encode([]), 1001, nil], ], ) @@ -244,10 +244,10 @@ module En57 2, ], ], - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).read(Query.all), + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(Query.all) + .to_a, ) end end @@ -267,8 +267,8 @@ module En57 }, ], [ - "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[])", - [array_encoder.encode([])], + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)", + [array_encoder.encode([]), 1001, nil], ], ) @@ -285,10 +285,10 @@ module En57 1, ], ], - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).read(Query.all), + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(Query.all) + .to_a, ) end end @@ -308,17 +308,17 @@ module En57 }, ], [ - "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[])", - [array_encoder.encode([])], + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)", + [array_encoder.encode([]), 1001, nil], ], ) assert_equal( [[Event.new(id: ids[0], type: "OrderPlaced", data: {}), 1]], - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).read(Query.all), + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(Query.all) + .to_a, ) end end @@ -342,8 +342,8 @@ module En57 }, ], [ - "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[])", - [array_encoder.encode(['{"tags":["order_id:123"]}'])], + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)", + [array_encoder.encode(['{"tags":["order_id:123"]}']), 1001, nil], ], ) @@ -361,10 +361,10 @@ module En57 1, ], ], - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).read(query), + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(query) + .to_a, ) end end @@ -385,8 +385,8 @@ module En57 }, ], [ - "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[])", - [array_encoder.encode(["{}"])], + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)", + [array_encoder.encode(["{}"]), 1001, nil], ], ) @@ -404,10 +404,10 @@ module En57 1, ], ], - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).read(query), + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(query) + .to_a, ) end end @@ -434,11 +434,13 @@ module En57 }, ], [ - "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[])", + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)", [ array_encoder.encode( %w[{"tags":["order_id:123"]} {"tags":["order_id:456"]}], ), + 1001, + nil, ], ], ) @@ -457,10 +459,10 @@ module En57 1, ], ], - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).read(query), + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(query) + .to_a, ) end end @@ -475,17 +477,17 @@ module En57 :exec_params, [], [ - "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[])", - [array_encoder.encode(['{"after":42}'])], + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)", + [array_encoder.encode(['{"after":42}']), 1001, nil], ], ) assert_equal( [], - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).read(query), + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(query) + .to_a, ) end end @@ -509,8 +511,8 @@ module En57 }, ], [ - "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[])", - [array_encoder.encode(['{"types":["OrderPlaced"]}'])], + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)", + [array_encoder.encode(['{"types":["OrderPlaced"]}']), 1001, nil], ], ) @@ -528,14 +530,130 @@ module En57 1, ], ], - Repository.new( - PgAdapter.for_connection(connection), - JsonSerializer.new, - ).read(query), + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(query) + .to_a, ) end end + def test_read_paginates_with_keyset_cursor_across_batches + En57 + .configuration + .stub(:read_batch_size, 2) do + with_connection do |connection| + connection.expect( + :exec_params, + [ + stored_row(1, ids[0]), + stored_row(2, ids[1]), + stored_row(3, ids[2]), + ], + [read_statement, [array_encoder.encode([]), 3, nil]], + ) + connection.expect( + :exec_params, + [ + stored_row(3, ids[2]), + stored_row(4, ids[3]), + stored_row(5, ids[4]), + ], + [read_statement, [array_encoder.encode([]), 3, 2]], + ) + connection.expect( + :exec_params, + [stored_row(5, ids[4])], + [read_statement, [array_encoder.encode([]), 3, 4]], + ) + + assert_equal( + [ + [stored_event(ids[0]), 1], + [stored_event(ids[1]), 2], + [stored_event(ids[2]), 3], + [stored_event(ids[3]), 4], + [stored_event(ids[4]), 5], + ], + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(Query.all) + .to_a, + ) + end + end + end + + def test_read_makes_single_query_when_result_fills_one_batch_exactly + En57 + .configuration + .stub(:read_batch_size, 2) do + with_connection do |connection| + connection.expect( + :exec_params, + [stored_row(1, ids[0]), stored_row(2, ids[1])], + [read_statement, [array_encoder.encode([]), 3, nil]], + ) + + assert_equal( + [[stored_event(ids[0]), 1], [stored_event(ids[1]), 2]], + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(Query.all) + .to_a, + ) + end + end + end + + def test_read_fetches_only_the_first_batch_when_consumer_takes_one + En57 + .configuration + .stub(:read_batch_size, 2) do + with_connection do |connection| + connection.expect( + :exec_params, + [ + stored_row(1, ids[0]), + stored_row(2, ids[1]), + stored_row(3, ids[2]), + ], + [read_statement, [array_encoder.encode([]), 3, nil]], + ) + + assert_equal( + [stored_event(ids[0]), 1], + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(Query.all) + .first, + ) + end + end + end + + def test_read_fetches_everything_in_one_query_when_batching_disabled + En57 + .configuration + .stub(:read_batch_size, nil) do + with_connection do |connection| + connection.expect( + :exec_params, + [stored_row(1, ids[0]), stored_row(2, ids[1])], + [read_statement, [array_encoder.encode([]), nil, nil]], + ) + + assert_equal( + [[stored_event(ids[0]), 1], [stored_event(ids[1]), 2]], + Repository + .new(PgAdapter.for_connection(connection), JsonSerializer.new) + .read(Query.all) + .to_a, + ) + end + end + end + def test_append_returns_failure_when_sql_returns_conflicting_events with_connection do |connection| connection.expect(:exec, nil, ["BEGIN ISOLATION LEVEL SERIALIZABLE"]) @@ -643,6 +761,23 @@ module En57 def ids = @ids ||= Hash.new { |h, k| h[k] = SecureRandom.uuid_v7 } + def read_statement = + "SELECT position, id, type, data, meta, tags FROM en57.read_events($1::jsonb[], $2, $3)" + + def stored_row(position, id) + { + "position" => position.to_s, + "id" => id, + "type" => "OrderPlaced", + "data" => nil, + "meta" => nil, + "tags" => "{}", + } + end + + def stored_event(id) = + Event.new(id: id, type: "OrderPlaced", data: {}, tags: []) + def with_connection connection = Minitest::Mock.new