diff --git a/lib/benchmark/concurrent_append_conflicting_tags.rb b/lib/benchmark/concurrent_append_conflicting_tags.rb index 8bf9a51..15fa21a 100644 --- a/lib/benchmark/concurrent_append_conflicting_tags.rb +++ b/lib/benchmark/concurrent_append_conflicting_tags.rb @@ -16,9 +16,7 @@ module En57 call do |measure| type = "event_benchmarked" tags = %W[writer:#{SecureRandom.hex(4)}] - barrier = Concurrent::CyclicBarrier.new(@concurrency) - - concurrently(@concurrency) do + concurrently do |barrier| scope = @event_store.read.of_type(type).with_tag(tags) events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } position = 0 diff --git a/lib/benchmark/concurrent_append_no_fail_if.rb b/lib/benchmark/concurrent_append_no_fail_if.rb index 34a88c3..fc3ac6e 100644 --- a/lib/benchmark/concurrent_append_no_fail_if.rb +++ b/lib/benchmark/concurrent_append_no_fail_if.rb @@ -15,9 +15,7 @@ module En57 call do |measure| type = "event_benchmarked" - barrier = Concurrent::CyclicBarrier.new(@concurrency) - - concurrently(@concurrency) do + concurrently do |barrier| tags = %W[writer:#{SecureRandom.hex(4)}] events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } diff --git a/lib/benchmark/concurrent_append_non_conflicting_tags.rb b/lib/benchmark/concurrent_append_non_conflicting_tags.rb index 4d5a247..a8465dd 100644 --- a/lib/benchmark/concurrent_append_non_conflicting_tags.rb +++ b/lib/benchmark/concurrent_append_non_conflicting_tags.rb @@ -15,9 +15,7 @@ module En57 call do |measure| type = "event_benchmarked" - barrier = Concurrent::CyclicBarrier.new(@concurrency) - - concurrently(@concurrency) do + concurrently do |barrier| tags = %W[writer:#{SecureRandom.hex(4)}] scope = @event_store.read.of_type(type).with_tag(tags) events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } diff --git a/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb b/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb index 8d40b26..9dd3cfc 100644 --- a/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb +++ b/lib/benchmark/concurrent_append_non_conflicting_tags_seeded.rb @@ -15,9 +15,7 @@ module En57 call do |measure| type = "event_benchmarked" - barrier = Concurrent::CyclicBarrier.new(@concurrency) - - concurrently(@concurrency) do + concurrently do |barrier| tags = %W[writer:#{SecureRandom.hex(4)}] scope = @event_store.read.of_type(type).with_tag(tags) events = Array.new(@batch_size) { Event.new(type: type, tags: tags) } diff --git a/lib/benchmark/res_concurrent_append_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_conflicting_streams.rb index af86bbc..236dd5e 100644 --- a/lib/benchmark/res_concurrent_append_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_conflicting_streams.rb @@ -19,9 +19,7 @@ module En57 call do |measure| type = "event_benchmarked" stream_name = "writer:#{SecureRandom.hex(4)}" - barrier = Concurrent::CyclicBarrier.new(@concurrency) - - concurrently(@concurrency) do + concurrently do |barrier| events = Array.new(@batch_size) do RubyEventStore::Event.new(metadata: { event_type: type }) diff --git a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb index f9d7e14..fb34842 100644 --- a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb @@ -18,9 +18,7 @@ module En57 call do |measure| type = "event_benchmarked" - barrier = Concurrent::CyclicBarrier.new(@concurrency) - - concurrently(@concurrency) do + concurrently do |barrier| tag = "writer:#{SecureRandom.hex(4)}" @event_store.read.stream(tag).of_type(type) events = diff --git a/lib/en57/benchmark.rb b/lib/en57/benchmark.rb index 20bde1f..b0780ec 100644 --- a/lib/en57/benchmark.rb +++ b/lib/en57/benchmark.rb @@ -236,8 +236,9 @@ module En57 def reset_retry_count = @retry_count.value = 0 def warmup = @warmup_runs.times { call(NOOP_MEASURE) } - def concurrently(concurrency) - Array.new(concurrency) { Thread.new { yield } }.each(&:value) + def concurrently + barrier = Concurrent::CyclicBarrier.new(@concurrency) + Array.new(@concurrency) { Thread.new { yield barrier } }.each(&:value) end end diff --git a/test/test_benchmark.rb b/test/test_benchmark.rb index 23d4cd5..c6b6f78 100644 --- a/test/test_benchmark.rb +++ b/test/test_benchmark.rb @@ -411,31 +411,37 @@ module En57 def test_scenario_runs_blocks_concurrently calls = Concurrent::AtomicFixnum.new(0) + barriers = Queue.new scenario = Class .new(Scenario) do - def initialize(calls) + def initialize(calls, barriers) + @barriers = barriers @calls = calls super( name: "concurrent", database_url: "postgres://example", runs: 1, warmup_runs: 0, - concurrency: 1, + concurrency: 2, batch_size: 1, ) end def call(_measure) - concurrently(2) { @calls.increment } - true + concurrently do |barrier| + @barriers << barrier + barrier.wait + @calls.increment + end end end - .new(calls) + .new(calls, barriers) scenario.run(->(&block) { block.call }) assert_equal(2, calls.value) + assert_equal(1, 2.times.map { barriers.pop }.uniq.size) end def test_runner_discovers_scenarios_by_database_instance