diff --git a/lib/benchmark/append_no_fail_if.rb b/lib/benchmark/append_no_fail_if.rb index 4aae647..2c39a35 100644 --- a/lib/benchmark/append_no_fail_if.rb +++ b/lib/benchmark/append_no_fail_if.rb @@ -2,39 +2,27 @@ module En57 module Benchmark - class AppendNoFailIf < Scenario - def self.key = "append-no-fail-if" + Scenario.define do + database_instance "append-no-fail-if" + name "1x100 append, no fail_if" + runs { it * 10 } + concurrency 1 + batch_size 100 - def self.build(database_url:, warmup_runs:, runs:) - new( - name: "1x100 append, no fail_if", - database_url:, - warmup_runs:, - runs: runs * 10, - concurrency: 1, - batch_size: 100, - ) - end - - def initialize(...) - super + setup do @event_store = EventStore.for_pooled_pg(@database_url, max_connections: @concurrency) end - private - - def call(measure) + call do |measure| type = "event_benchmarked" tags = %W[writer:#{SecureRandom.hex(4)}] - events = - Array.new(@batch_size) { En57::Event.new(type: type, tags: tags) } + events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } measure.call { @event_store.append(events) } - verify end - def verify = @event_store.read.each.to_a.size == total_runs * @batch_size + verify { @event_store.read.each.to_a.size == total_runs * @batch_size } end end end diff --git a/lib/benchmark/append_non_conflicting_tags.rb b/lib/benchmark/append_non_conflicting_tags.rb index e683559..f58eb69 100644 --- a/lib/benchmark/append_non_conflicting_tags.rb +++ b/lib/benchmark/append_non_conflicting_tags.rb @@ -2,34 +2,23 @@ module En57 module Benchmark - class AppendNonConflictingTags < Scenario - def self.key = "append-non-conflicting-tags" + Scenario.define do + database_instance "append-non-conflicting-tags" + name "1x100 append, non-conflicting tags" + runs { it * 10 } + concurrency 1 + batch_size 100 - def self.build(database_url:, warmup_runs:, runs:) - new( - name: "1x100 append, non-conflicting tags", - database_url:, - warmup_runs:, - runs: runs * 10, - concurrency: 1, - batch_size: 100, - ) - end - - def initialize(...) - super + setup do @event_store = EventStore.for_pooled_pg(@database_url, max_connections: @concurrency) end - private - - def call(measure) + call do |measure| type = "event_benchmarked" tags = %W[writer:#{SecureRandom.hex(4)}] scope = @event_store.read.of_type(type).with_tag(tags) - events = - Array.new(@batch_size) { En57::Event.new(type: type, tags: tags) } + events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } measure.call do begin @@ -39,10 +28,9 @@ module En57 retry end end - verify end - def verify = @event_store.read.each.to_a.size == total_runs * @batch_size + verify { @event_store.read.each.to_a.size == total_runs * @batch_size } end end end diff --git a/lib/benchmark/concurrent_append_conflicting_tags.rb b/lib/benchmark/concurrent_append_conflicting_tags.rb index deaff5a..2f0a575 100644 --- a/lib/benchmark/concurrent_append_conflicting_tags.rb +++ b/lib/benchmark/concurrent_append_conflicting_tags.rb @@ -2,37 +2,25 @@ module En57 module Benchmark - class ConcurrentAppendConflictingTags < Scenario - def self.key = "concurrent-append-conflicting-tags" + Scenario.define do + database_instance "concurrent-append-conflicting-tags" + name "10x100 concurrent append, conflicting tags" + concurrency 10 + batch_size 100 - def self.build(database_url:, warmup_runs:, runs:) - new( - name: "10x100 concurrent append, conflicting tags", - database_url:, - warmup_runs:, - runs:, - concurrency: 10, - batch_size: 100, - ) - end - - def initialize(...) - super + setup do @event_store = EventStore.for_pooled_pg(@database_url, max_connections: @concurrency) end - private - - def call(measure) + call do |measure| type = "event_benchmarked" tags = %W[writer:#{SecureRandom.hex(4)}] barrier = Concurrent::CyclicBarrier.new(@concurrency) concurrently(@concurrency) do scope = @event_store.read.of_type(type).with_tag(tags) - events = - Array.new(@batch_size) { En57::Event.new(type: type, tags: tags) } + events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } position = 0 barrier.wait @@ -49,12 +37,12 @@ module En57 end end end - verify end - def verify = + verify do @event_store.read.each.to_a.size == total_runs * @concurrency * @batch_size + end end end end diff --git a/lib/benchmark/concurrent_append_no_fail_if.rb b/lib/benchmark/concurrent_append_no_fail_if.rb index 010296e..83043ea 100644 --- a/lib/benchmark/concurrent_append_no_fail_if.rb +++ b/lib/benchmark/concurrent_append_no_fail_if.rb @@ -2,47 +2,35 @@ module En57 module Benchmark - class ConcurrentAppendNoFailIf < Scenario - def self.key = "concurrent-append-no-fail-if" + Scenario.define do + database_instance "concurrent-append-no-fail-if" + name "10x100 concurrent append, no fail_if" + concurrency 10 + batch_size 100 - def self.build(database_url:, warmup_runs:, runs:) - new( - name: "10x100 concurrent append, no fail_if", - database_url:, - warmup_runs:, - runs:, - concurrency: 10, - batch_size: 100, - ) - end - - def initialize(...) - super + setup do @event_store = EventStore.for_pooled_pg(@database_url, max_connections: @concurrency) end - private - - def call(measure) + call do |measure| type = "event_benchmarked" barrier = Concurrent::CyclicBarrier.new(@concurrency) concurrently(@concurrency) do tags = %W[writer:#{SecureRandom.hex(4)}] - events = - Array.new(@batch_size) { En57::Event.new(type: type, tags: tags) } + events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } barrier.wait measure.call { @event_store.append(events) } end - verify end - def verify = + verify do @event_store.read.each.to_a.size == total_runs * @concurrency * @batch_size + end end end end diff --git a/lib/benchmark/concurrent_append_non_conflicting_tags.rb b/lib/benchmark/concurrent_append_non_conflicting_tags.rb index fe5afcd..63544eb 100644 --- a/lib/benchmark/concurrent_append_non_conflicting_tags.rb +++ b/lib/benchmark/concurrent_append_non_conflicting_tags.rb @@ -2,37 +2,25 @@ module En57 module Benchmark - class ConcurrentAppendNonConflictingTags < Scenario - def self.key = "concurrent-append-non-conflicting-tags" + Scenario.define do + database_instance "concurrent-append-non-conflicting-tags" + name "10x100 concurrent append, non-conflicting tags" + concurrency 10 + batch_size 100 - def self.build(database_url:, warmup_runs:, runs:) - new( - name: "10x100 concurrent append, non-conflicting tags", - database_url:, - warmup_runs:, - runs:, - concurrency: 10, - batch_size: 100, - ) - end - - def initialize(...) - super + setup do @event_store = EventStore.for_pooled_pg(@database_url, max_connections: @concurrency) end - private - - def call(measure) + call do |measure| type = "event_benchmarked" barrier = Concurrent::CyclicBarrier.new(@concurrency) concurrently(@concurrency) do tags = %W[writer:#{SecureRandom.hex(4)}] scope = @event_store.read.of_type(type).with_tag(tags) - events = - Array.new(@batch_size) { En57::Event.new(type: type, tags: tags) } + events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } barrier.wait @@ -45,12 +33,12 @@ module En57 end end end - verify end - def verify = + verify do @event_store.read.each.to_a.size == total_runs * @concurrency * @batch_size + end end end end diff --git a/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb b/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb index 01ffaa6..ac3c022 100644 --- a/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb +++ b/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb @@ -2,39 +2,25 @@ module En57 module Benchmark - class ConcurrentAppendNonConflictingTagsSeeded < Scenario - SEEDED_EVENTS = 1_000_000 + Scenario.define do + database_instance "concurrent-append-non-conflicting-tags-seeded" + name "10x100 concurrent append, non-conflicting tags (seeded)" + concurrency 10 + batch_size 100 - def self.key = "concurrent-append-non-conflicting-tags-seeded" - - def self.build(database_url:, warmup_runs:, runs:) - new( - name: "10x100 concurrent append, non-conflicting tags (seeded)", - database_url:, - warmup_runs:, - runs:, - concurrency: 10, - batch_size: 100, - ) - end - - def initialize(...) - super + setup do @event_store = EventStore.for_pooled_pg(@database_url, max_connections: @concurrency) end - private - - def call(measure) + call do |measure| type = "event_benchmarked" barrier = Concurrent::CyclicBarrier.new(@concurrency) concurrently(@concurrency) do tags = %W[writer:#{SecureRandom.hex(4)}] scope = @event_store.read.of_type(type).with_tag(tags) - events = - Array.new(@batch_size) { En57::Event.new(type: type, tags: tags) } + events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } barrier.wait @@ -47,12 +33,12 @@ module En57 end end end - verify end - def verify = + verify do @event_store.read.each.to_a.size == - SEEDED_EVENTS + total_runs * @concurrency * @batch_size + 1_000_000 + total_runs * @concurrency * @batch_size + end end end end diff --git a/lib/benchmark/res_append_stream_any.rb b/lib/benchmark/res_append_stream_any.rb index af32894..52e22e3 100644 --- a/lib/benchmark/res_append_stream_any.rb +++ b/lib/benchmark/res_append_stream_any.rb @@ -5,29 +5,19 @@ require "rails_event_store" module En57 module Benchmark - class ResAppendStreamAny < Scenario - def self.key = "res-append-stream-any" - - def self.build(database_url:, warmup_runs:, runs:) - new( - name: "1x100 append, expected_version :any (RES)", - database_url:, - warmup_runs:, - runs: runs * 10, - concurrency: 1, - batch_size: 100, - ) - end - - def initialize(...) - super + Scenario.define do + database_instance "res-append-stream-any" + name "1x100 append, expected_version :any (RES)" + runs { it * 10 } + concurrency 1 + batch_size 100 + + setup do ActiveRecord::Base.establish_connection(@database_url) @event_store = RailsEventStore::JSONClient.new end - private - - def call(measure) + call do |measure| type = "event_benchmarked" tag = "writer:#{SecureRandom.hex(4)}" events = @@ -38,10 +28,9 @@ module En57 measure.call do @event_store.append(events, stream_name: tag, expected_version: :any) end - verify end - def verify = @event_store.read.each.to_a.size == total_runs * @batch_size + verify { @event_store.read.each.to_a.size == total_runs * @batch_size } end end end diff --git a/lib/benchmark/res_concurrent_append_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_conflicting_streams.rb index 3f2eaea..bf3c05a 100644 --- a/lib/benchmark/res_concurrent_append_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_conflicting_streams.rb @@ -5,29 +5,18 @@ require "rails_event_store" module En57 module Benchmark - class ResConcurrentAppendConflictingStreams < Scenario - def self.key = "res-concurrent-append-conflicting-streams" + Scenario.define do + database_instance "res-concurrent-append-conflicting-streams" + name "10x100 concurrent append, conflicting streams (RES)" + concurrency 10 + batch_size 100 - def self.build(database_url:, warmup_runs:, runs:) - new( - name: "10x100 concurrent append, conflicting streams (RES)", - database_url:, - warmup_runs:, - runs:, - concurrency: 10, - batch_size: 100, - ) - end - - def initialize(...) - super + setup do ActiveRecord::Base.establish_connection(@database_url) @event_store = RailsEventStore::JSONClient.new end - private - - def call(measure) + call do |measure| type = "event_benchmarked" stream_name = "writer:#{SecureRandom.hex(4)}" barrier = Concurrent::CyclicBarrier.new(@concurrency) @@ -55,12 +44,12 @@ module En57 end end end - verify end - def verify = + verify do @event_store.read.each.to_a.size == total_runs * @concurrency * @batch_size + end end end end diff --git a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb index 8429e79..8b036d7 100644 --- a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb @@ -5,29 +5,18 @@ require "rails_event_store" module En57 module Benchmark - class ResConcurrentAppendNonConflictingStreams < Scenario - def self.key = "res-concurrent-append-non-conflicting-streams" - - def self.build(database_url:, warmup_runs:, runs:) - new( - name: "10x100 concurrent append, non-conflicting streams (RES)", - database_url:, - warmup_runs:, - runs:, - concurrency: 10, - batch_size: 100, - ) - end + Scenario.define do + database_instance "res-concurrent-append-non-conflicting-streams" + name "10x100 concurrent append, non-conflicting streams (RES)" + concurrency 10 + batch_size 100 - def initialize(...) - super + setup do ActiveRecord::Base.establish_connection(@database_url) @event_store = RailsEventStore::JSONClient.new end - private - - def call(measure) + call do |measure| type = "event_benchmarked" barrier = Concurrent::CyclicBarrier.new(@concurrency) @@ -49,12 +38,12 @@ module En57 ) end end - verify end - def verify = + verify do @event_store.read.each.to_a.size == total_runs * @concurrency * @batch_size + end end end end diff --git a/lib/en57/benchmark.rb b/lib/en57/benchmark.rb index 9d87c1d..75ab391 100644 --- a/lib/en57/benchmark.rb +++ b/lib/en57/benchmark.rb @@ -120,7 +120,93 @@ module En57 end end + ScenarioDefinition = + Data.define( + :database_instance, + :name, + :runs, + :concurrency, + :batch_size, + :setup, + :call_block, + :verify, + ) + + class ScenarioDSL + def initialize + @runs = ->(runs) { runs } + @concurrency = 1 + @batch_size = 100 + @setup = -> {} + @call_block = ->(_measure) {} + @verify = -> { true } + end + + def database_instance(value) = @database_instance = value + def name(value) = @name = value + + def runs(&block) = @runs = block + def concurrency(value) = @concurrency = value + def batch_size(value) = @batch_size = value + def setup(&block) = @setup = block + def call(&block) = @call_block = block + def verify(&block) = @verify = block + + def definition + ScenarioDefinition.new( + database_instance: @database_instance, + name: @name, + runs: @runs, + concurrency: @concurrency, + batch_size: @batch_size, + setup: @setup, + call_block: @call_block, + verify: @verify, + ) + end + end + class Scenario + @definitions = [] + + def self.definitions = @definitions + + def self.define(&block) + definition = ScenarioDSL.new.tap { it.instance_eval(&block) }.definition + Class + .new(self) do + define_singleton_method(:database_instance) do + definition.database_instance + end + + define_singleton_method( + :build, + ) do |database_url:, warmup_runs:, runs:| + new( + name: definition.name, + database_url:, + runs: definition.runs.call(runs), + warmup_runs:, + concurrency: definition.concurrency, + batch_size: definition.batch_size, + ) + end + + define_method(:initialize) do |**kwargs| + super(**kwargs) + instance_exec(&definition.setup) + end + + define_method(:call) do |measure| + instance_exec(measure, &definition.call_block) + verify + end + + define_method(:verify) { instance_exec(&definition.verify) } + end + .tap { definitions << it } + end + def initialize( name:, database_url:, @@ -180,25 +266,19 @@ module En57 def self.names = scenarios(runs: nil).keys def self.scenarios(runs:) - scenario_classes.to_h do |scenario_class| - [ - scenario_class.key, - ->(database_url, warmup_runs) do - scenario_class.build(database_url:, warmup_runs:, runs:) - end, - ] - end - end - - def self.scenario_classes Scenario - .subclasses - .select { it.respond_to?(:key) && it.respond_to?(:build) } - .sort_by(&:key) + .definitions + .sort_by(&:database_instance) + .to_h do |scenario_class| + [ + scenario_class.database_instance, + ->(database_url, warmup_runs) do + scenario_class.build(database_url:, warmup_runs:, runs:) + end, + ] + end end - private_class_method :scenario_classes - def initialize(scenarios:, formatter:) @formatter = formatter @scenarios = scenarios diff --git a/test/test_benchmark.rb b/test/test_benchmark.rb index d4df560..8171d43 100644 --- a/test/test_benchmark.rb +++ b/test/test_benchmark.rb @@ -251,6 +251,88 @@ module En57 assert_in_delta(0.15, formatted_results.fetch(0).median) end + def test_scenario_define_registers_configured_scenario + original_definitions = Scenario.definitions.dup + scenario_class = + Scenario.define do + database_instance "defined" + name "Defined scenario" + runs { it * 2 } + concurrency 2 + batch_size 3 + + setup { @setup_called = true } + call do |measure| + measure.call { @call_measured = true } + nil + end + verify { @setup_called && @call_measured } + end + scenario = + scenario_class.build( + database_url: "postgres://example", + warmup_runs: 0, + runs: 4, + ) + + assert_includes(Scenario.definitions, scenario_class) + assert_equal("defined", scenario_class.database_instance) + assert_equal("Defined scenario", scenario.name) + assert_equal( + "postgres://example", + scenario.instance_variable_get(:@database_url), + ) + assert_equal(8, scenario.runs) + assert_equal(2, scenario.instance_variable_get(:@concurrency)) + assert_equal(3, scenario.instance_variable_get(:@batch_size)) + assert_equal(true, scenario.run(->(&block) { block.call })) + assert_equal(true, scenario.instance_variable_get(:@call_measured)) + ensure + Scenario.definitions.replace(original_definitions) + end + + def test_scenario_define_uses_verify_block + original_definitions = Scenario.definitions.dup + scenario_class = + Scenario.define do + database_instance "unverified" + name "Unverified scenario" + verify { false } + end + scenario = + scenario_class.build( + database_url: "postgres://example", + warmup_runs: 0, + runs: 1, + ) + + assert_equal(false, scenario.run(->(&block) { block.call })) + ensure + Scenario.definitions.replace(original_definitions) + end + + def test_scenario_define_defaults + original_definitions = Scenario.definitions.dup + scenario_class = + Scenario.define do + database_instance "defaulted" + name "Defaulted scenario" + end + scenario = + scenario_class.build( + database_url: "postgres://example", + warmup_runs: 0, + runs: 7, + ) + + assert_equal(7, scenario.runs) + assert_equal(1, scenario.instance_variable_get(:@concurrency)) + assert_equal(100, scenario.instance_variable_get(:@batch_size)) + assert_equal(true, scenario.run(->(&block) { block.call })) + ensure + Scenario.definitions.replace(original_definitions) + end + def test_scenario_stores_configuration scenario = Scenario.new( @@ -376,35 +458,25 @@ module En57 assert_equal(2, calls.value) end - def test_runner_discovers_scenario_classes - first_scenario_class = + def test_runner_discovers_scenarios_by_database_instance + first_scenario = Class.new do - def self.key = "a-discovered" + def self.database_instance = "a-discovered" def self.build(database_url:, warmup_runs:, runs:) [database_url, warmup_runs, runs] end end - second_scenario_class = + second_scenario = Class.new do - def self.key = "b-discovered" + def self.database_instance = "b-discovered" def self.build(database_url:, warmup_runs:, runs:) [database_url, warmup_runs, runs] end end - incomplete_scenario_class = Class.new { def self.key = "ignored" } - anonymous_scenario_class = Class.new { def self.build(...) = nil } - - Scenario.stub( - :subclasses, - [ - second_scenario_class, - incomplete_scenario_class, - anonymous_scenario_class, - first_scenario_class, - ], - ) do + + Scenario.stub(:definitions, [second_scenario, first_scenario]) do assert_equal(%w[a-discovered b-discovered], Runner.names) assert_equal( ["postgres://example", 2, 3], @@ -417,21 +489,18 @@ module En57 end def test_classic_runner_selects_named_scenarios - first_scenario_class = + first_scenario = Class.new do - def self.key = "first" + def self.database_instance = "first" def self.build(...) = nil end - second_scenario_class = + second_scenario = Class.new do - def self.key = "second" + def self.database_instance = "second" def self.build(...) = nil end - Scenario.stub( - :subclasses, - [first_scenario_class, second_scenario_class], - ) do + Scenario.stub(:definitions, [first_scenario, second_scenario]) do runner = Runner.classic(names: ["first"]) assert_equal( @@ -442,16 +511,16 @@ module En57 end def test_classic_runner_defaults_to_fifty_runs - scenario_class = + scenario = Class.new do - def self.key = "scenario" + def self.database_instance = "scenario" def self.build(database_url:, warmup_runs:, runs:) [database_url, warmup_runs, runs] end end - Scenario.stub(:subclasses, [scenario_class]) do + Scenario.stub(:definitions, [scenario]) do scenario = Runner .classic