diff --git a/lib/benchmark/append_no_fail_if.rb b/lib/benchmark/append_no_fail_if.rb index ed2c9e2..1804bc6 100644 --- a/lib/benchmark/append_no_fail_if.rb +++ b/lib/benchmark/append_no_fail_if.rb @@ -11,13 +11,14 @@ module En57 private - def call + def call(measure) type = "event_benchmarked" tags = %W[writer:#{SecureRandom.hex(4)}] events = Array.new(@batch_size) { En57::Event.new(type: type, tags: tags) } - @measure.call { @event_store.append(events) } + measure.call { @event_store.append(events) } + verify end def verify = @event_store.read.each.to_a.size == total_runs * @batch_size diff --git a/lib/benchmark/append_non_conflicting_tags.rb b/lib/benchmark/append_non_conflicting_tags.rb index 6876f52..1e5a434 100644 --- a/lib/benchmark/append_non_conflicting_tags.rb +++ b/lib/benchmark/append_non_conflicting_tags.rb @@ -11,14 +11,14 @@ module En57 private - def call + def call(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) } - @measure.call do + measure.call do begin @event_store.append(events, fail_if: scope.after(position = 0)) rescue AppendConditionViolated @@ -26,6 +26,7 @@ module En57 retry end end + verify end def verify = @event_store.read.each.to_a.size == total_runs * @batch_size diff --git a/lib/benchmark/concurrent_append_conflicting_tags.rb b/lib/benchmark/concurrent_append_conflicting_tags.rb index e5dadcd..9660e76 100644 --- a/lib/benchmark/concurrent_append_conflicting_tags.rb +++ b/lib/benchmark/concurrent_append_conflicting_tags.rb @@ -11,7 +11,7 @@ module En57 private - def call + def call(measure) type = "event_benchmarked" tags = %W[writer:#{SecureRandom.hex(4)}] barrier = Concurrent::CyclicBarrier.new(@concurrency) @@ -24,7 +24,7 @@ module En57 barrier.wait - @measure.call do + measure.call do begin @event_store.append(events, fail_if: scope.after(position)) rescue AppendConditionViolated @@ -36,6 +36,7 @@ module En57 end end end + verify end def verify = diff --git a/lib/benchmark/concurrent_append_no_fail_if.rb b/lib/benchmark/concurrent_append_no_fail_if.rb index 6477c2b..cc98bce 100644 --- a/lib/benchmark/concurrent_append_no_fail_if.rb +++ b/lib/benchmark/concurrent_append_no_fail_if.rb @@ -11,7 +11,7 @@ module En57 private - def call + def call(measure) type = "event_benchmarked" barrier = Concurrent::CyclicBarrier.new(@concurrency) @@ -22,8 +22,9 @@ module En57 barrier.wait - @measure.call { @event_store.append(events) } + measure.call { @event_store.append(events) } end + verify end def verify = diff --git a/lib/benchmark/concurrent_append_non_conflicting_tags.rb b/lib/benchmark/concurrent_append_non_conflicting_tags.rb index d8ea284..e9794e1 100644 --- a/lib/benchmark/concurrent_append_non_conflicting_tags.rb +++ b/lib/benchmark/concurrent_append_non_conflicting_tags.rb @@ -11,7 +11,7 @@ module En57 private - def call + def call(measure) type = "event_benchmarked" barrier = Concurrent::CyclicBarrier.new(@concurrency) @@ -23,7 +23,7 @@ module En57 barrier.wait - @measure.call do + measure.call do begin @event_store.append(events, fail_if: scope.after(position = 0)) rescue AppendConditionViolated @@ -32,6 +32,7 @@ module En57 end end end + verify end def verify = diff --git a/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb b/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb index 7aba247..2dec96a 100644 --- a/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb +++ b/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb @@ -13,7 +13,7 @@ module En57 private - def call + def call(measure) type = "event_benchmarked" barrier = Concurrent::CyclicBarrier.new(@concurrency) @@ -25,7 +25,7 @@ module En57 barrier.wait - @measure.call do + measure.call do begin @event_store.append(events, fail_if: scope.after(position = 0)) rescue AppendConditionViolated @@ -34,6 +34,7 @@ module En57 end end end + verify end def verify = diff --git a/lib/benchmark/res_append_stream_any.rb b/lib/benchmark/res_append_stream_any.rb index d22751d..520a4e8 100644 --- a/lib/benchmark/res_append_stream_any.rb +++ b/lib/benchmark/res_append_stream_any.rb @@ -14,7 +14,7 @@ module En57 private - def call + def call(measure) type = "event_benchmarked" tag = "writer:#{SecureRandom.hex(4)}" events = @@ -22,9 +22,10 @@ module En57 RubyEventStore::Event.new(metadata: { event_type: type }) end - @measure.call do + 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 diff --git a/lib/benchmark/res_concurrent_append_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_conflicting_streams.rb index c0de45c..6e511e8 100644 --- a/lib/benchmark/res_concurrent_append_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_conflicting_streams.rb @@ -14,7 +14,7 @@ module En57 private - def call + def call(measure) type = "event_benchmarked" stream_name = "writer:#{SecureRandom.hex(4)}" barrier = Concurrent::CyclicBarrier.new(@concurrency) @@ -28,7 +28,7 @@ module En57 barrier.wait - @measure.call do + measure.call do begin @event_store.append( events, @@ -42,6 +42,7 @@ module En57 end end end + verify end def verify = diff --git a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb index 38f58ed..23e1440 100644 --- a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb @@ -14,7 +14,7 @@ module En57 private - def call + def call(measure) type = "event_benchmarked" barrier = Concurrent::CyclicBarrier.new(@concurrency) @@ -28,7 +28,7 @@ module En57 barrier.wait - @measure.call do + measure.call do @event_store.append( events, stream_name: tag, @@ -36,6 +36,7 @@ module En57 ) end end + verify end def verify = diff --git a/lib/en57/benchmark.rb b/lib/en57/benchmark.rb index e877494..36a8c3e 100644 --- a/lib/en57/benchmark.rb +++ b/lib/en57/benchmark.rb @@ -124,7 +124,6 @@ module En57 def initialize( name:, database_url:, - measure:, runs:, warmup_runs:, concurrency:, @@ -134,7 +133,6 @@ module En57 @batch_size = batch_size @concurrency = concurrency @database_url = database_url - @measure = measure @runs = runs @retry_count = Concurrent::AtomicFixnum.new @warmup_runs = warmup_runs @@ -144,22 +142,23 @@ module En57 def retry_count = @retry_count.value - def run + NOOP_MEASURE = ->(&block) { block.call } + + def run(measure) warmup reset_retry_count - @runs.times { call } - verify + verified = true + @runs.times { verified = call(measure) } + verified end private def total_runs = @runs + @warmup_runs - def call - end + def call(_measure) = true def record_retry = @retry_count.increment def reset_retry_count = @retry_count.value = 0 - def verify = true - def warmup = @warmup_runs.times { call } + def warmup = @warmup_runs.times { call(NOOP_MEASURE) } def concurrently(concurrency) Array.new(concurrency) { Thread.new { yield } }.each(&:value) @@ -188,41 +187,30 @@ module En57 def self.scenarios(runs:) { - "append-no-fail-if" => ->(database_url, warmup_runs, measure) do + "append-no-fail-if" => ->(database_url, warmup_runs) do AppendNoFailIf.new( name: "1x100 append, no fail_if", database_url:, - measure:, warmup_runs:, runs: runs * 10, concurrency: 1, batch_size: 100, ) end, - "append-non-conflicting-tags" => ->( - database_url, - warmup_runs, - measure - ) do + "append-non-conflicting-tags" => ->(database_url, warmup_runs) do AppendNonConflictingTags.new( name: "1x100 append, non-conflicting tags", database_url:, - measure:, warmup_runs:, runs: runs * 10, concurrency: 1, batch_size: 100, ) end, - "concurrent-append-no-fail-if" => ->( - database_url, - warmup_runs, - measure - ) do + "concurrent-append-no-fail-if" => ->(database_url, warmup_runs) do ConcurrentAppendNoFailIf.new( name: "10x100 concurrent append, no fail_if", database_url:, - measure:, warmup_runs:, runs:, concurrency: 10, @@ -231,13 +219,11 @@ module En57 end, "concurrent-append-non-conflicting-tags" => ->( database_url, - warmup_runs, - measure + warmup_runs ) do ConcurrentAppendNonConflictingTags.new( name: "10x100 concurrent append, non-conflicting tags", database_url:, - measure:, warmup_runs:, runs:, concurrency: 10, @@ -246,13 +232,11 @@ module En57 end, "concurrent-append-non-conflicting-tags-seeded" => ->( database_url, - warmup_runs, - measure + warmup_runs ) do ConcurrentAppendNonConflictingTagsSeeded.new( name: "10x100 concurrent append, non-conflicting tags (seeded)", database_url:, - measure:, warmup_runs:, runs:, concurrency: 10, @@ -261,24 +245,21 @@ module En57 end, "concurrent-append-conflicting-tags" => ->( database_url, - warmup_runs, - measure + warmup_runs ) do ConcurrentAppendConflictingTags.new( name: "10x100 concurrent append, conflicting tags", database_url:, - measure:, warmup_runs:, runs:, concurrency: 10, batch_size: 100, ) end, - "res-append-stream-any" => ->(database_url, warmup_runs, measure) do + "res-append-stream-any" => ->(database_url, warmup_runs) do ResAppendStreamAny.new( name: "1x100 append, expected_version :any (RES)", database_url:, - measure:, warmup_runs:, runs: runs * 10, concurrency: 1, @@ -287,13 +268,11 @@ module En57 end, "res-concurrent-append-non-conflicting-streams" => ->( database_url, - warmup_runs, - measure + warmup_runs ) do ResConcurrentAppendNonConflictingStreams.new( name: "10x100 concurrent append, non-conflicting streams (RES)", database_url:, - measure:, warmup_runs:, runs:, concurrency: 10, @@ -302,13 +281,11 @@ module En57 end, "res-concurrent-append-conflicting-streams" => ->( database_url, - warmup_runs, - measure + warmup_runs ) do ResConcurrentAppendConflictingStreams.new( name: "10x100 concurrent append, conflicting streams (RES)", database_url:, - measure:, warmup_runs:, runs:, concurrency: 10, @@ -328,14 +305,12 @@ module En57 @scenarios.map do |instance_name, mk_scenario| PgEphemeral.with_server(instance_name:) do |server| samples = [] - scenario = - mk_scenario.call( - server.url, - warmup_runs = 2, + scenario = mk_scenario.call(server.url, 2) + verified = + scenario.run( ->(&block) { samples << ::Benchmark.realtime { block.call } }, ) - verified = scenario.run - measurement = Measurement.from(samples.drop(warmup_runs)) + measurement = Measurement.from(samples) Result.new( name: scenario.name, diff --git a/test/test_benchmark.rb b/test/test_benchmark.rb index 8b230a1..f5f8959 100644 --- a/test/test_benchmark.rb +++ b/test/test_benchmark.rb @@ -108,15 +108,15 @@ module En57 server = Data.define(:url).new("postgres://example") mk_scenario = ->(name, verified) do - ->(_database_url, _warmup_runs, measure) do + ->(_database_url, _warmup_runs) do Data - .define(:name, :runs, :measure, :verified, :retry_count) do - def run + .define(:name, :runs, :verified, :retry_count) do + def run(measure) 3.times { measure.call { nil } } verified end end - .new(name, 1, measure, verified, 3) + .new(name, 1, verified, 3) end end @@ -146,19 +146,17 @@ module En57 server = Data.define(:url).new("postgres://example") instance_names = [] database_urls = [] + warmup_runs = [] measured_blocks = 0 - mk_scenario = ->(database_url, _warmup_runs, measure) do + mk_scenario = ->(database_url, warmup_run_count) do database_urls << database_url + warmup_runs << warmup_run_count Class .new do - define_method(:initialize) { @measure = measure } - - attr_reader :measure - def name = "scenario" def retry_count = 0 def runs = 1 - def run + def run(measure) 3.times { measure.call { @measured_blocks.call } } true end @@ -184,10 +182,11 @@ module En57 assert_equal(["instance"], instance_names) assert_equal(["postgres://example"], database_urls) + assert_equal([2], warmup_runs) assert_equal(3, measured_blocks) end - def test_runner_discards_two_warmup_measurements + def test_scenario_uses_noop_measure_for_warmup formatter = Object.new formatted_results = nil @@ -198,11 +197,13 @@ module En57 scenario_class = Class.new(Scenario) do - def initialize(measure:, warmup_runs:) + attr_reader :call_count + + def initialize(warmup_runs:) + @call_count = 0 super( name: "warmup", database_url: "postgres://example", - measure:, runs: 2, warmup_runs:, concurrency: 1, @@ -210,10 +211,15 @@ module En57 ) end - def call = @measure.call { nil } + def call(measure) + @call_count += 1 + measure.call { nil } + true + end end + scenario = nil server = Data.define(:url).new("postgres://example") - durations = [0.1, 0.2, 0.3, 0.5] + durations = [0.1, 0.2] PgEphemeral.stub( :with_server, @@ -229,19 +235,20 @@ module En57 Runner.new( formatter:, scenarios: { - "warmup" => ->(_database_url, warmup_runs, measure) do - scenario_class.new(measure:, warmup_runs:) + "warmup" => ->(_database_url, warmup_runs) do + scenario = scenario_class.new(warmup_runs:) end, }, ).run end end - assert_equal(0.4, formatted_results.fetch(0).mean) - assert_equal(0.1, formatted_results.fetch(0).stddev) - assert_equal(0.3, formatted_results.fetch(0).min) - assert_equal(0.5, formatted_results.fetch(0).max) - assert_equal(0.4, formatted_results.fetch(0).median) + assert_equal(4, scenario.call_count) + assert_in_delta(0.15, formatted_results.fetch(0).mean) + assert_in_delta(0.05, formatted_results.fetch(0).stddev) + assert_equal(0.1, formatted_results.fetch(0).min) + assert_equal(0.2, formatted_results.fetch(0).max) + assert_in_delta(0.15, formatted_results.fetch(0).median) end def test_scenario_calculates_total_runs @@ -252,7 +259,6 @@ module En57 super( name: "total", database_url: "postgres://example", - measure: ->(&block) { block.call }, runs: 2, warmup_runs: 3, concurrency: 1, @@ -272,7 +278,6 @@ module En57 Scenario.new( name: "noop", database_url: "postgres://example", - measure: ->(&block) { block.call }, runs: 1, warmup_runs: 1, concurrency: 1, @@ -280,7 +285,21 @@ module En57 ) assert_equal(0, scenario.retry_count) - assert_equal(true, scenario.run) + assert_equal(true, scenario.run(->(&block) { block.call })) + end + + def test_scenario_verifies_when_no_measured_runs + scenario = + Scenario.new( + name: "empty", + database_url: "postgres://example", + runs: 0, + warmup_runs: 0, + concurrency: 1, + batch_size: 1, + ) + + assert_equal(true, scenario.run(->(&block) { block.call })) end def test_scenario_counts_retries_after_warmup @@ -291,7 +310,6 @@ module En57 super( name: "retrying", database_url: "postgres://example", - measure: ->(&block) { block.call }, runs: 2, warmup_runs: 1, concurrency: 1, @@ -299,11 +317,14 @@ module En57 ) end - def call = record_retry + def call(_measure) + record_retry + true + end end .new - assert_equal(true, scenario.run) + assert_equal(true, scenario.run(->(&block) { block.call })) assert_equal(2, scenario.retry_count) end @@ -317,7 +338,6 @@ module En57 super( name: "concurrent", database_url: "postgres://example", - measure: ->(&block) { block.call }, runs: 1, warmup_runs: 0, concurrency: 1, @@ -325,11 +345,14 @@ module En57 ) end - def call = concurrently(2) { @calls.increment } + def call(_measure) + concurrently(2) { @calls.increment } + true + end end .new(calls) - assert_equal(true, scenario.run) + assert_equal(true, scenario.run(->(&block) { block.call })) assert_equal(2, calls.value) end @@ -351,7 +374,6 @@ module En57 end def test_classic_runner_builds_scenarios - measure = ->(&block) { block.call } scenarios = Runner.classic.instance_variable_get(:@scenarios) [ @@ -398,7 +420,7 @@ module En57 10, ], ].each do |key, scenario_class, name, runs, concurrency| - scenario = scenarios.fetch(key).call("postgres://example", 2, measure) + scenario = scenarios.fetch(key).call("postgres://example", 2) assert_instance_of(scenario_class, scenario) assert_equal(name, scenario.name) @@ -406,7 +428,6 @@ module En57 "postgres://example", scenario.instance_variable_get(:@database_url), ) - assert_same(measure, scenario.instance_variable_get(:@measure)) assert_equal(2, scenario.instance_variable_get(:@warmup_runs)) assert_equal( concurrency,