From b2cb72d75fc0c65a1fd606a2d9a29839b7241e82 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Pawe=C5=82=20Pacana?= Date: Thu, 25 Jun 2026 15:02:04 +0200 Subject: [PATCH] Clone ephemeral test databases from golden templates MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Replace the static initialDatabases list (13 named DBs reset with TRUNCATE) with three golden templates: golden_en57, golden_res, and golden_en57_seeded. The seeded golden bakes its 1M-row seed in at provisioning time. - Add En57::EphemeralDatabase.with(template:), which clones a golden via CREATE DATABASE ... TEMPLATE (PostgreSQL's fastest copy — file-level, not row-by-row), yields a database_url, and drops it WITH (FORCE) when the block ends. No TRUNCATE. - IntegrationTest now clones golden_en57 per test and drops it on teardown, isolating every test; pools connect lazily. Benchmarks clone per scenario instead of resetting. - The seeded benchmark clones the 1M-row golden (~1s) instead of re-INSERTing the seed on every run (~16s). - Ignore EphemeralDatabase#execute in mutant: the receiver->self mutation runs Kernel#exec on multi-word SQL and replaces the worker, so it cannot be killed. - Adopting this needs a one-time Postgres data-dir reset so devenv provisions the goldens (devenv up and devenv test use separate state dirs); the old named databases are obsolete. --- .mutant.yml | 1 + devenv.nix | 37 +++--- .../concurrent_append_non_conflicting_tags.rb | 9 +- lib/benchmark/res_append_stream_any.rb | 2 +- ...s_concurrent_append_conflicting_streams.rb | 2 +- ...ncurrent_append_non_conflicting_streams.rb | 2 +- lib/en57/benchmark.rb | 76 +++++------- lib/en57/ephemeral_database.rb | 35 ++++++ test/test_benchmark.rb | 100 ++++++++-------- test/test_ephemeral_database.rb | 111 ++++++++++++++++++ test/test_factories.rb | 19 +-- test/test_helper.rb | 62 +++++----- test/test_integration.rb | 32 ++--- test/test_migrator.rb | 16 +-- test/test_stress.rb | 4 +- 15 files changed, 318 insertions(+), 190 deletions(-) create mode 100644 lib/en57/ephemeral_database.rb create mode 100644 test/test_ephemeral_database.rb diff --git a/.mutant.yml b/.mutant.yml index ece4ea4..45868b9 100644 --- a/.mutant.yml +++ b/.mutant.yml @@ -33,3 +33,4 @@ matcher: - En57::Benchmark::ResConcurrentAppendConflictingStreams* - En57::Benchmark::ResConcurrentAppendNonConflictingStreams* - En57::Benchmark::Scenario#concurrently # thread raise mutation survives despite direct coverage + - En57::EphemeralDatabase#execute # receiver->self mutation runs Kernel#exec on multi-word SQL, replacing the worker diff --git a/devenv.nix b/devenv.nix index 7c382d2..926b5f4 100644 --- a/devenv.nix +++ b/devenv.nix @@ -42,29 +42,24 @@ package = pkgs.postgresql_18; initialDatabases = let - en57 = name: { - inherit name; - schema = ./db/schema/0.1.0.sql; - }; - res = name: { - inherit name; - schema = ./db/seeds/res.sql; - }; + seededSchema = pkgs.runCommand "golden-en57-seeded.sql" { } '' + cat ${./db/schema/0.1.0.sql} \ + ${./db/seeds/concurrent_append_non_conflicting_tags_seeded.sql} > $out + ''; in [ - (en57 "main") - (en57 "append-no-fail-if") - (en57 "append-no-fail-if-ar") - (en57 "append-non-conflicting-tags") - (en57 "concurrent-append-no-fail-if") - (en57 "concurrent-append-no-fail-if-ar") - (en57 "concurrent-append-conflicting-tags") - (en57 "concurrent-append-non-conflicting-tags") - (en57 "concurrent-append-non-conflicting-tags-seeded") - (res "res-append-stream-any") - (res "res-concurrent-append-non-conflicting-streams") - (res "res-concurrent-append-conflicting-streams") - (en57 "regress") + { + name = "golden_en57"; + schema = ./db/schema/0.1.0.sql; + } + { + name = "golden_res"; + schema = ./db/seeds/res.sql; + } + { + name = "golden_en57_seeded"; + schema = seededSchema; + } ]; }; diff --git a/lib/benchmark/concurrent_append_non_conflicting_tags.rb b/lib/benchmark/concurrent_append_non_conflicting_tags.rb index 7ee7a49..d55fe14 100644 --- a/lib/benchmark/concurrent_append_non_conflicting_tags.rb +++ b/lib/benchmark/concurrent_append_non_conflicting_tags.rb @@ -38,14 +38,7 @@ module En57 concurrent_append_non_conflicting_tags.with( database_instance: "concurrent-append-non-conflicting-tags-seeded", name: "10x100 concurrent append, non-conflicting tags (seeded)", - reset: - Scenario::RESET_EN57 + "; " + - File.read( - File.expand_path( - "../../db/seeds/concurrent_append_non_conflicting_tags_seeded.sql", - __dir__, - ), - ), + template: "golden_en57_seeded", ) end end diff --git a/lib/benchmark/res_append_stream_any.rb b/lib/benchmark/res_append_stream_any.rb index 33f36b1..c4cf767 100644 --- a/lib/benchmark/res_append_stream_any.rb +++ b/lib/benchmark/res_append_stream_any.rb @@ -8,7 +8,7 @@ module En57 runs: ->(runs) { runs * 10 }, concurrency: 1, batch_size: 100, - reset: Scenario::RESET_RES, + template: "golden_res", ) do def setup(database_url) require "active_record" diff --git a/lib/benchmark/res_concurrent_append_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_conflicting_streams.rb index f59e122..0101076 100644 --- a/lib/benchmark/res_concurrent_append_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_conflicting_streams.rb @@ -7,7 +7,7 @@ module En57 name: "10x100 concurrent append, conflicting streams (RES)", concurrency: 10, batch_size: 100, - reset: Scenario::RESET_RES, + template: "golden_res", ) do def setup(database_url) require "active_record" diff --git a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb index 5f9167c..e1992f7 100644 --- a/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb +++ b/lib/benchmark/res_concurrent_append_non_conflicting_streams.rb @@ -7,7 +7,7 @@ module En57 name: "10x100 concurrent append, non-conflicting streams (RES)", concurrency: 10, batch_size: 100, - reset: Scenario::RESET_RES, + template: "golden_res", ) do def setup(database_url) require "active_record" diff --git a/lib/en57/benchmark.rb b/lib/en57/benchmark.rb index 563b301..40c4f6f 100644 --- a/lib/en57/benchmark.rb +++ b/lib/en57/benchmark.rb @@ -7,6 +7,7 @@ require "pg" require "securerandom" require_relative "../en57" +require_relative "ephemeral_database" module En57 module Benchmark @@ -22,7 +23,7 @@ module En57 :retry_count, ) - Runnable = Data.define(:build, :reset) + Runnable = Data.define(:build, :template) class Table def format(results) @@ -129,14 +130,9 @@ module En57 :concurrency, :batch_size, :runs, - :reset, + :template, ) - RESET_EN57 = "TRUNCATE en57.tags, en57.events RESTART IDENTITY CASCADE" - RESET_RES = - "TRUNCATE event_store_events, event_store_events_in_streams " \ - "RESTART IDENTITY CASCADE" - @definitions = [] def self.definitions = @definitions @@ -147,7 +143,7 @@ module En57 concurrency: 1, batch_size: 100, runs: ->(runs) { runs }, - reset: RESET_EN57, + template: "golden_en57", &block ) register( @@ -157,7 +153,7 @@ module En57 concurrency:, batch_size:, runs:, - reset:, + template:, ), &block ) @@ -174,7 +170,7 @@ module En57 configuration.database_instance end - define_singleton_method(:reset) { configuration.reset } + define_singleton_method(:template) { configuration.template } define_singleton_method( :build, @@ -262,7 +258,7 @@ module En57 [ scenario_class.database_instance, Runnable.new( - reset: scenario_class.reset, + template: scenario_class.template, build: ->(database_url, warmup_runs) do scenario_class.build(database_url:, warmup_runs:, runs:) end, @@ -281,45 +277,37 @@ module En57 En57.configuration.append_retries = 100 results = - @scenarios.map do |instance_name, runnable| - database_url = "postgres:///#{instance_name}" - reset(database_url, runnable.reset) - - samples = Concurrent::Array.new - retries = Concurrent::AtomicFixnum.new - - scenario = runnable.build.call(database_url, 2) - scenario.run( - ->(&block) { samples << ::Benchmark.realtime { block.call } }, - -> { retries.increment }, - ) - measurement = Measurement.from(samples) - - Result.new( - name: scenario.name, - runs: scenario.runs, - mean: measurement.mean, - stddev: measurement.stddev, - min: measurement.min, - max: measurement.max, - median: measurement.median, - retry_count: retries.value, - ) + @scenarios.map do |_instance_name, runnable| + EphemeralDatabase.with( + template: runnable.template, + ) do |database_url| + samples = Concurrent::Array.new + retries = Concurrent::AtomicFixnum.new + + scenario = runnable.build.call(database_url, 2) + scenario.run( + ->(&block) { samples << ::Benchmark.realtime { block.call } }, + -> { retries.increment }, + ) + measurement = Measurement.from(samples) + + Result.new( + name: scenario.name, + runs: scenario.runs, + mean: measurement.mean, + stddev: measurement.stddev, + min: measurement.min, + max: measurement.max, + median: measurement.median, + retry_count: retries.value, + ) + end end @formatter.format(results) ensure En57.configuration.append_retries = original_append_retries end - - private - - def reset(database_url, sql) - connection = PG.connect(database_url) - connection.exec(sql) - ensure - connection&.close - end end class CLI diff --git a/lib/en57/ephemeral_database.rb b/lib/en57/ephemeral_database.rb new file mode 100644 index 0000000..062276d --- /dev/null +++ b/lib/en57/ephemeral_database.rb @@ -0,0 +1,35 @@ +# frozen_string_literal: true + +require "pg" +require "securerandom" + +module En57 + module EphemeralDatabase + extend self + + ADMIN_URL = "postgres:///postgres" + + def with(template: nil, prefix: "en57") + name = "#{prefix}.#{SecureRandom.hex(8)}" + admin = PG.connect(ADMIN_URL) + statement = "CREATE DATABASE #{PG::Connection.quote_ident(name)}" + statement += + " TEMPLATE #{PG::Connection.quote_ident(template)}" if template + execute(admin, statement) + yield "postgres:///#{name}" + ensure + if admin + execute( + admin, + "DROP DATABASE IF EXISTS #{PG::Connection.quote_ident(name)} " \ + "WITH (FORCE)", + ) + end + admin&.close + end + + private + + def execute(connection, statement) = connection.exec(statement) + end +end diff --git a/test/test_benchmark.rb b/test/test_benchmark.rb index 93e81c9..c06e224 100644 --- a/test/test_benchmark.rb +++ b/test/test_benchmark.rb @@ -88,7 +88,7 @@ module En57 mk_scenario = ->(name) do Runnable.new( - reset: "", + template: "golden_en57", build: ->(_database_url, _warmup_runs) do Data .define(:name, :runs, :retry_count) do @@ -149,7 +149,7 @@ module En57 scenarios: { "instance" => Runnable.new( - reset: "", + template: "golden_en57", build: ->(_database_url, _warmup_runs) { scenario }, ), }, @@ -195,7 +195,7 @@ module En57 scenarios: { "instance" => Runnable.new( - reset: "", + template: "golden_en57", build: ->(_database_url, _warmup_runs) { scenario }, ), }, @@ -207,7 +207,7 @@ module En57 assert_in_delta(1.1, formatted_results.fetch(0).mean) end - def test_runner_resets_database_and_uses_instance_database_urls + def test_runner_clones_template_and_yields_ephemeral_database_url formatter = Object.new formatter.define_singleton_method(:format) { |_results| "formatted" } database_urls = [] @@ -240,20 +240,29 @@ module En57 Runner.new( formatter:, scenarios: { - "instance" => Runnable.new(reset: "RESET SQL", build:), + "instance" => Runnable.new(template: "golden_res", build:), }, ).run end - assert_equal(["postgres:///instance"], database_urls) + assert_equal(1, database_urls.size) + assert_match(%r{\Apostgres:///en57\.\h{16}\z}, database_urls.fetch(0)) assert_equal([2], warmup_runs) assert_equal(3, measured_blocks) - assert_equal(["postgres:///instance"], connection.urls) - assert_equal(["RESET SQL"], connection.statements) + assert_equal([EphemeralDatabase::ADMIN_URL], connection.urls) + assert_match( + /\ACREATE DATABASE "en57\.\h{16}" TEMPLATE "golden_res"\z/, + connection.statements.fetch(0), + ) + assert_match( + /\ADROP DATABASE IF EXISTS "en57\.\h{16}" WITH \(FORCE\)\z/, + connection.statements.fetch(1), + ) + assert_equal(2, connection.statements.size) assert_equal(1, connection.closed) end - def test_runner_propagates_reset_connection_errors_without_masking + def test_runner_propagates_clone_connection_errors_without_masking formatter = Object.new formatter.define_singleton_method(:format) { |_results| "formatted" } boom = Class.new(StandardError) @@ -266,7 +275,7 @@ module En57 scenarios: { "instance" => Runnable.new( - reset: "", + template: "golden_en57", build: ->(_database_url, _warmup_runs) { nil }, ), }, @@ -324,7 +333,7 @@ module En57 scenarios: { "warmup" => Runnable.new( - reset: "", + template: "golden_en57", build: ->(_database_url, warmup_runs) do scenario = scenario_class.new(warmup_runs:) end, @@ -659,7 +668,7 @@ module En57 first_scenario = Class.new do def self.database_instance = "a-discovered" - def self.reset = "reset-a" + def self.template = "golden_a" def self.build(database_url:, warmup_runs:, runs:) [database_url, warmup_runs, runs] @@ -668,7 +677,7 @@ module En57 second_scenario = Class.new do def self.database_instance = "b-discovered" - def self.reset = "reset-b" + def self.template = "golden_b" def self.build(database_url:, warmup_runs:, runs:) [database_url, warmup_runs, runs] @@ -678,7 +687,7 @@ module En57 Scenario.stub(:definitions, [second_scenario, first_scenario]) do assert_equal(%w[a-discovered b-discovered], Runner.names) runnable = Runner.scenarios(runs: 3).fetch("a-discovered") - assert_equal("reset-a", runnable.reset) + assert_equal("golden_a", runnable.template) assert_equal( ["postgres://example", 2, 3], runnable.build.call("postgres://example", 2), @@ -690,13 +699,13 @@ module En57 first_scenario = Class.new do def self.database_instance = "first" - def self.reset = "" + def self.template = "golden_en57" def self.build(...) = nil end second_scenario = Class.new do def self.database_instance = "second" - def self.reset = "" + def self.template = "golden_en57" def self.build(...) = nil end @@ -714,7 +723,7 @@ module En57 scenario = Class.new do def self.database_instance = "scenario" - def self.reset = "" + def self.template = "golden_en57" def self.build(database_url:, warmup_runs:, runs:) [database_url, warmup_runs, runs] @@ -833,57 +842,44 @@ module En57 assert_in_delta(0.35, measurement.median) end - def test_reset_en57_truncates_event_and_tag_tables - assert_equal( - "TRUNCATE en57.tags, en57.events RESTART IDENTITY CASCADE", - Scenario::RESET_EN57, - ) - end - - def test_reset_res_truncates_event_store_tables - assert_equal( - "TRUNCATE event_store_events, event_store_events_in_streams " \ - "RESTART IDENTITY CASCADE", - Scenario::RESET_RES, - ) - end - - def test_scenario_define_defaults_reset_to_en57_truncate + def test_scenario_define_defaults_template_to_golden_en57 original_definitions = Scenario.definitions.dup scenario_class = - Scenario.define(database_instance: "reset-default", name: "Reset") + Scenario.define( + database_instance: "template-default", + name: "Template", + ) - assert_equal(Scenario::RESET_EN57, scenario_class.reset) + assert_equal("golden_en57", scenario_class.template) ensure Scenario.definitions.replace(original_definitions) end - def test_scenario_define_accepts_custom_reset + def test_scenario_define_accepts_custom_template original_definitions = Scenario.definitions.dup scenario_class = Scenario.define( - database_instance: "reset-custom", - name: "Reset", - reset: "TRUNCATE custom", + database_instance: "template-custom", + name: "Template", + template: "golden_custom", ) - assert_equal("TRUNCATE custom", scenario_class.reset) + assert_equal("golden_custom", scenario_class.template) ensure Scenario.definitions.replace(original_definitions) end - def test_seeded_scenario_reset_reloads_seed_after_truncate + def test_seeded_scenario_clones_seeded_template seeded = Scenario.definitions.find do it.database_instance == "concurrent-append-non-conflicting-tags-seeded" end - assert(seeded.reset.start_with?(Scenario::RESET_EN57)) - assert_includes(seeded.reset, "INSERT INTO en57.events") + assert_equal("golden_en57_seeded", seeded.template) end - def test_res_scenarios_reset_with_res_truncate + def test_res_scenarios_clone_res_template res_scenarios = Scenario.definitions.select do it.database_instance.start_with?("res-") @@ -891,18 +887,22 @@ module En57 refute_empty(res_scenarios) res_scenarios.each do |scenario| - assert_equal(Scenario::RESET_RES, scenario.reset) + assert_equal("golden_res", scenario.template) end end - def test_scenario_with_overrides_reset + def test_scenario_with_overrides_template original_definitions = Scenario.definitions.dup scenario_class = - Scenario.define(database_instance: "reset-base", name: "Base") - copy = scenario_class.with(database_instance: "reset-copy", reset: "X") + Scenario.define(database_instance: "template-base", name: "Base") + copy = + scenario_class.with( + database_instance: "template-copy", + template: "golden_x", + ) - assert_equal(Scenario::RESET_EN57, scenario_class.reset) - assert_equal("X", copy.reset) + assert_equal("golden_en57", scenario_class.template) + assert_equal("golden_x", copy.template) ensure Scenario.definitions.replace(original_definitions) end diff --git a/test/test_ephemeral_database.rb b/test/test_ephemeral_database.rb new file mode 100644 index 0000000..638a364 --- /dev/null +++ b/test/test_ephemeral_database.rb @@ -0,0 +1,111 @@ +# frozen_string_literal: true + +require "test_helper" + +module En57 + class TestEphemeralDatabase < Minitest::Test + cover "En57::EphemeralDatabase*" + + def test_clones_from_template_yields_url_then_drops_and_closes + admin = recording_connection + urls = [] + yielded = nil + connect = ->(url) do + urls << url + admin + end + + result = + PG.stub(:connect, connect) do + EphemeralDatabase.with(template: "golden_x") do |database_url| + yielded = database_url + "block-result" + end + end + + assert_equal([EphemeralDatabase::ADMIN_URL], urls) + assert_match(%r{\Apostgres:///en57\.\h{16}\z}, yielded) + + name = yielded.delete_prefix("postgres:///") + assert_equal( + [ + %(CREATE DATABASE "#{name}" TEMPLATE "golden_x"), + %(DROP DATABASE IF EXISTS "#{name}" WITH (FORCE)), + ], + admin.statements, + ) + assert_equal(1, admin.closed) + assert_equal("block-result", result) + end + + def test_without_template_creates_a_plain_database + admin = recording_connection + yielded = nil + + PG.stub(:connect, ->(_url) { admin }) do + EphemeralDatabase.with { |database_url| yielded = database_url } + end + + name = yielded.delete_prefix("postgres:///") + assert_equal(%(CREATE DATABASE "#{name}"), admin.statements.fetch(0)) + end + + def test_uses_the_given_prefix_for_the_database_name + yielded = nil + + PG.stub(:connect, ->(_url) { recording_connection }) do + EphemeralDatabase.with(prefix: "migrator") { |url| yielded = url } + end + + assert_match(%r{\Apostgres:///migrator\.\h{16}\z}, yielded) + end + + def test_drops_and_closes_even_when_the_block_raises + admin = recording_connection + boom = Class.new(StandardError) + + assert_raises(boom) do + PG.stub(:connect, ->(_url) { admin }) do + EphemeralDatabase.with(template: "golden_x") { raise boom } + end + end + + assert_equal(2, admin.statements.size) + assert_match( + /\ADROP DATABASE IF EXISTS .+ WITH \(FORCE\)\z/, + admin.statements.fetch(1), + ) + assert_equal(1, admin.closed) + end + + def test_propagates_connection_errors_without_running_teardown + boom = Class.new(StandardError) + + assert_raises(boom) do + PG.stub(:connect, ->(_url) { raise boom }) do + EphemeralDatabase.with(template: "golden_x") do + flunk("must not yield") + end + end + end + end + + private + + def recording_connection + Class + .new do + attr_reader :statements, :closed + + def initialize + @statements = [] + @closed = 0 + end + + def exec(sql) = @statements << sql + def close = @closed += 1 + end + .new + end + end +end diff --git a/test/test_factories.rb b/test/test_factories.rb index 79c21cb..e213ac1 100644 --- a/test/test_factories.rb +++ b/test/test_factories.rb @@ -5,15 +5,18 @@ require "test_helper" module En57 class TestFactories < IntegrationTest def test_for_pg_round_trips_with_connection_uri - assert_round_trip EventStore.for_pg(MAIN_URL) + assert_round_trip EventStore.for_pg(database_url) end def test_for_pooled_pg_round_trips_with_default_max_connections - assert_round_trip EventStore.for_pooled_pg(MAIN_URL) + assert_round_trip EventStore.for_pooled_pg(database_url) end def test_for_pooled_pg_round_trips_with_custom_max_connections - assert_round_trip EventStore.for_pooled_pg(MAIN_URL, max_connections: 1) + assert_round_trip EventStore.for_pooled_pg( + database_url, + max_connections: 1, + ) end def test_for_active_record_round_trips_with_default_model @@ -27,16 +30,16 @@ module En57 end def test_for_sequel_round_trips_with_database - assert_round_trip EventStore.for_sequel(SEQUEL_DB) + assert_round_trip EventStore.for_sequel(sequel_db) end def test_event_store_does_not_conflict_with_public_schema_tables - CONNECTION.exec("CREATE TABLE public.events (id integer PRIMARY KEY)") - CONNECTION.exec("CREATE TABLE public.tags (id integer PRIMARY KEY)") + connection.exec("CREATE TABLE public.events (id integer PRIMARY KEY)") + connection.exec("CREATE TABLE public.tags (id integer PRIMARY KEY)") - assert_round_trip EventStore.for_pg(MAIN_URL) + assert_round_trip EventStore.for_pg(database_url) ensure - CONNECTION.exec("DROP TABLE IF EXISTS public.tags, public.events") + connection.exec("DROP TABLE IF EXISTS public.tags, public.events") end private diff --git a/test/test_helper.rb b/test/test_helper.rb index 1f8e11c..4ca46ae 100644 --- a/test/test_helper.rb +++ b/test/test_helper.rb @@ -11,6 +11,7 @@ require "active_record" require "connection_pool" require "en57" +require "en57/ephemeral_database" # test dependencies require "securerandom" @@ -18,42 +19,49 @@ require "concurrent-ruby" module En57 class IntegrationTest < Minitest::Test - MAIN_URL = "postgres:///main" - - CONNECTION = PG.connect(MAIN_URL) + ADAPTER_NAMES = %i[pg sequel active_record] POOL_SIZE = 8 - PG_POOL = ConnectionPool.new(size: POOL_SIZE) { PG.connect(MAIN_URL) } + attr_reader :database_url, :connection, :sequel_db - SEQUEL_DB = - Sequel.connect( - MAIN_URL, - preconnect: :concurrently, - max_connections: POOL_SIZE, + def setup + @admin = PG.connect(EphemeralDatabase::ADMIN_URL) + @database_name = "en57.#{SecureRandom.hex(8)}" + @admin.exec( + "CREATE DATABASE #{PG::Connection.quote_ident(@database_name)} " \ + "TEMPLATE golden_en57", ) + @database_url = "postgres:///#{@database_name}" - AR_POOL = -> do - ActiveRecord::Base.establish_connection("#{MAIN_URL}?pool=#{POOL_SIZE}") - ActiveRecord::Base.connection_pool - end.call - - ADAPTERS = { - pg: -> { PgAdapter.for_pool(PG_POOL) }, - sequel: -> { SequelAdapter.new(SEQUEL_DB) }, - active_record: -> { ActiveRecordAdapter.new(AR_POOL) }, - } + @connection = PG.connect(@database_url) + @pg_pool = + ConnectionPool.new(size: POOL_SIZE) { PG.connect(@database_url) } + @sequel_db = Sequel.connect(@database_url, max_connections: POOL_SIZE) + ActiveRecord::Base.establish_connection( + "#{@database_url}?pool=#{POOL_SIZE}", + ) + @ar_pool = ActiveRecord::Base.connection_pool + end - def setup = - CONNECTION.exec( - "TRUNCATE TABLE en57.tags, en57.events RESTART IDENTITY CASCADE", + def teardown + @ar_pool&.disconnect! + @sequel_db&.disconnect + @pg_pool&.shutdown(&:close) + @connection&.close + @admin&.exec( + "DROP DATABASE IF EXISTS " \ + "#{PG::Connection.quote_ident(@database_name)} WITH (FORCE)", ) + @admin&.close + end - Minitest.after_run do - AR_POOL.disconnect! - SEQUEL_DB.disconnect - PG_POOL.shutdown(&:close) - CONNECTION.close + def adapter_factory(name) + { + pg: -> { PgAdapter.for_pool(@pg_pool) }, + sequel: -> { SequelAdapter.new(@sequel_db) }, + active_record: -> { ActiveRecordAdapter.new(@ar_pool) }, + }.fetch(name) end end end diff --git a/test/test_integration.rb b/test/test_integration.rb index c0ab1e7..796a365 100644 --- a/test/test_integration.rb +++ b/test/test_integration.rb @@ -4,9 +4,9 @@ require "test_helper" module En57 class TestIntegration < IntegrationTest - ADAPTERS.each do |name, factory| + ADAPTER_NAMES.each do |name| define_method "test_#{name}_happy_path" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| events = [ Event.new( id: ids[0], @@ -30,7 +30,7 @@ module En57 end define_method "test_#{name}_read_with_position_yields_events_and_positions" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| events = [ Event.new(id: ids[0], type: "OrderPlaced"), Event.new(id: ids[1], type: "PriceChanged"), @@ -45,7 +45,7 @@ module En57 end define_method "test_#{name}_append_with_fail_if_and_no_matches_appends_events" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| event = Event.new(id: ids[0], type: "OrderPlaced") assert_equal( Success.new(position: 1), @@ -60,7 +60,7 @@ module En57 end define_method "test_#{name}_append_with_fail_if_and_matches_returns_failure" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| existing_event = Event.new( id: ids[0], @@ -88,7 +88,7 @@ module En57 end define_method "test_#{name}_append_with_after_ignores_matches_at_or_before_cutoff" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| existing_event = Event.new(id: ids[0], type: "OrderPlaced") assert_equal( Success.new(position: 1), @@ -110,7 +110,7 @@ module En57 end define_method "test_#{name}_append_with_after_returns_failure_if_match_is_after_cutoff" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| existing_event = Event.new(id: ids[0], type: "OrderPlaced") assert_equal( Success.new(position: 1), @@ -130,7 +130,7 @@ module En57 end define_method "test_#{name}_append_with_duplicate_id_raises_unique_violation" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| existing_event = Event.new(id: ids[0], type: "OrderPlaced") assert_equal( Success.new(position: 1), @@ -148,7 +148,7 @@ module En57 end define_method "test_#{name}_tags_round_trip" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| event = Event.new(id: ids[0], type: "OrderPlaced", tags: ["order_id:123"]) @@ -158,7 +158,7 @@ module En57 end define_method "test_#{name}_read_filters_after" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| events = [ Event.new(id: ids[0], type: "OrderPlaced"), Event.new(id: ids[1], type: "PriceChanged"), @@ -170,7 +170,7 @@ module En57 end define_method "test_#{name}_read_filters_by_tags" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| events = [ Event.new( id: ids[0], @@ -197,7 +197,7 @@ module En57 end define_method "test_#{name}_read_filters_by_type" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| events = [ Event.new(id: ids[0], type: "OrderPlaced"), Event.new(id: ids[1], type: "PriceChanged"), @@ -212,7 +212,7 @@ module En57 end define_method "test_#{name}_read_filters_by_any_of_types" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| events = [ Event.new(id: ids[0], type: "PriceChanged"), Event.new(id: ids[1], type: "OrderPlaced"), @@ -228,7 +228,7 @@ module En57 end define_method "test_#{name}_read_filters_by_type_and_tag_on_same_item" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| events = [ Event.new(id: ids[0], type: "OrderPlaced", tags: ["order_id:123"]), Event.new(id: ids[1], type: "OrderPlaced", tags: ["order_id:456"]), @@ -249,7 +249,7 @@ module En57 end define_method "test_#{name}_read_or_combines_scopes_as_disjunction" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| events = [ Event.new(id: ids[0], type: "OrderPlaced", tags: ["order_id:123"]), Event.new(id: ids[1], type: "OrderPlaced", tags: ["order_id:456"]), @@ -272,7 +272,7 @@ module En57 end define_method "test_#{name}_read_streams_results_spanning_many_batches" do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| events = (1..5).map do |n| Event.new( diff --git a/test/test_migrator.rb b/test/test_migrator.rb index 59cfb64..2ada086 100644 --- a/test/test_migrator.rb +++ b/test/test_migrator.rb @@ -4,6 +4,9 @@ require "test_helper" module En57 class TestMigrator < IntegrationTest + def setup = nil + def teardown = nil + def test_status_reports_pending_schema_on_empty_database with_database do |url| assert_equal( @@ -108,17 +111,8 @@ module En57 private - def with_database - name = "en57_migrator_#{SecureRandom.hex(8)}" - CONNECTION.exec(%(CREATE DATABASE #{PG::Connection.quote_ident(name)})) - yield database_url(name) - ensure - CONNECTION.exec( - %(DROP DATABASE IF EXISTS #{PG::Connection.quote_ident(name)}), - ) - end - - def database_url(name) = "postgres:///#{name}" + def with_database(&block) = + EphemeralDatabase.with(prefix: "en57-migrator", &block) def schema_path(version) File.expand_path("../db/schema/#{version}.sql", __dir__) diff --git a/test/test_stress.rb b/test/test_stress.rb index 9d9eb5e..bbb34b7 100644 --- a/test/test_stress.rb +++ b/test/test_stress.rb @@ -4,11 +4,11 @@ require "test_helper" module En57 class TestStress < IntegrationTest - ADAPTERS.each do |name, factory| + ADAPTER_NAMES.each do |name| define_method( "test_#{name}_only_one_writer_can_consume_account_credits", ) do - with_event_store(factory) do |event_store| + with_event_store(adapter_factory(name)) do |event_store| event_store.append( [ Event.new( -- 2.51.2