From 01ccdaa9c46fed14375ea4ee9ca76373fbf35cbd Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pawe=C5=82=20Pacana?= Date: Wed, 20 May 2026 13:20:14 +0200 Subject: [PATCH] Add conflicting tags benchmark --- database.toml | 4 ++++ lib/en57/benchmark.rb | 48 ++++++++++++++++++++++++++++++++++++++++++- 2 files changed, 51 insertions(+), 1 deletion(-) diff --git a/database.toml b/database.toml index c528047..5d9a070 100644 --- a/database.toml +++ b/database.toml @@ -11,3 +11,7 @@ path = "db/schema/0.1.0.sql" [instances.concurrent-append-no-fail-if.seeds.schema] type = "sql-file" path = "db/schema/0.1.0.sql" + +[instances.concurrent-append-conflicting-tags.seeds.schema] +type = "sql-file" +path = "db/schema/0.1.0.sql" diff --git a/lib/en57/benchmark.rb b/lib/en57/benchmark.rb index 7d8f4f8..ea0d5ca 100644 --- a/lib/en57/benchmark.rb +++ b/lib/en57/benchmark.rb @@ -144,6 +144,16 @@ module En57 batch_size: 100, ) end, + "concurrent-append-conflicting-tags" => ->(database_url, measure) do + ConcurrentAppendConflictingTags.new( + name: "Concurrent append, conflicting tags", + database_url:, + measure:, + runs: ENV.fetch("BENCHMARK_RUNS", 1), + concurrency: 10, + batch_size: 100, + ) + end, }, ) end @@ -194,7 +204,6 @@ module En57 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) } @@ -248,5 +257,42 @@ module En57 def verify = @event_store.read.each.to_a.size == @runs * @concurrency * @batch_size end + + class ConcurrentAppendConflictingTags < Scenario + def initialize(...) + super + @event_store = + EventStore.for_pooled_pg(@database_url, max_connections: @concurrency) + end + + private + + def call + 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) } + position = 0 + + barrier.wait + + @measure.call do + begin + @event_store.append(events, fail_if: scope.after(position)) + rescue AppendConditionViolated + scope.each_with_position { |_event, event_position| position = event_position } + retry + end + end + end + end + + def verify = + @event_store.read.each.to_a.size == @runs * @concurrency * @batch_size + end end end -- 2.51.2