diff --git a/.mutant.yml b/.mutant.yml index 7449efd..f4c2821 100644 --- a/.mutant.yml +++ b/.mutant.yml @@ -18,3 +18,4 @@ matcher: - En57::EventStore.for_pooled_pg - En57::EventStore.for_sequel - En57::JsonSerializer#optional_big_decimal_type + - En57::Migrator* diff --git a/Rakefile b/Rakefile index 3f6abce..153bf89 100644 --- a/Rakefile +++ b/Rakefile @@ -4,6 +4,7 @@ require "bundler/gem_tasks" require "fileutils" require "minitest/test_task" require "uri" +require_relative "lib/en57/tasks" Minitest::TestTask.create diff --git a/en57.gemspec b/en57.gemspec index 874e93d..bfba555 100644 --- a/en57.gemspec +++ b/en57.gemspec @@ -17,7 +17,7 @@ Gem::Specification.new do |spec| "changelog_uri" ] = "https://github.com/mostlyobvious/en57/blob/main/CHANGELOG.md" - spec.files = Dir["lib/**/*"] + spec.files = Dir["lib/**/*", "db/schema/**/*.sql"] spec.require_paths = ["lib"] spec.extra_rdoc_files = %w[README.md] diff --git a/lib/en57.rb b/lib/en57.rb index 4c90f18..6a2bc44 100644 --- a/lib/en57.rb +++ b/lib/en57.rb @@ -9,6 +9,7 @@ require_relative "en57/pg_adapter" require_relative "en57/sequel_adapter" if defined?(Sequel) require_relative "en57/active_record_adapter" if defined?(ActiveRecord) require_relative "en57/repository" +require_relative "en57/migrator" require_relative "en57/event_store" module En57 diff --git a/lib/en57/migrator.rb b/lib/en57/migrator.rb new file mode 100644 index 0000000..4f6157a --- /dev/null +++ b/lib/en57/migrator.rb @@ -0,0 +1,146 @@ +# frozen_string_literal: true + +require "pg" +require_relative "version" + +module En57 + class Migrator + MigrationError = Class.new(StandardError) + + Status = Data.define(:current, :target, :state, :method, :pending, :warning) + + def initialize(connection_string) + @connection_string = connection_string + end + + def status + current = current_version + + Status.new( + current:, + target: SCHEMA_VERSION, + state: state_for(current), + method: (current == :fresh ? :none : :version_table), + pending: pending_for(current), + warning: warning_for(current), + ) + end + + def migrate! + current = current_version + raise partial_migration_error if current == :partial + raise MigrationError, "Cannot infer En57 schema version" if current == :unversioned + return if current == SCHEMA_VERSION + + if current == :fresh + install_fresh_schema + else + raise MigrationError, "No En57 schema diff from #{current} to #{SCHEMA_VERSION}" + end + end + + alias migrate migrate! + + private + + def state_for(current) + case current + when SCHEMA_VERSION then :up_to_date + when :partial then :partial + else :pending + end + end + + def pending_for(current) + current == SCHEMA_VERSION || current == :partial ? [] : [schema_path(SCHEMA_VERSION)] + end + + def warning_for(current) + return unless current == :partial + + "Previous En57 migration did not complete cleanly. " \ + "Inspect the database and resolve manually, then run: " \ + "UPDATE public.en57_schema_info SET in_progress = false WHERE id = 1;" + end + + def current_version + with_connection do |connection| + next :fresh unless schema_info_table?(connection) + + result = connection.exec("SELECT schema_version, in_progress FROM public.en57_schema_info WHERE id = 1") + next :unversioned if result.ntuples.zero? + + row = result.first + next :partial if row.fetch("in_progress") == "t" + + row.fetch("schema_version") + end + end + + def install_fresh_schema + with_connection do |connection| + ensure_schema_info_table(connection) + mark_in_progress(connection) + + connection.transaction do |transaction| + transaction.exec(File.read(schema_path(SCHEMA_VERSION))) + record_version(transaction) + end + end + end + + def ensure_schema_info_table(connection) + connection.exec(<<~SQL) + CREATE TABLE IF NOT EXISTS public.en57_schema_info ( + id integer PRIMARY KEY, + schema_version varchar(20) NOT NULL, + in_progress boolean NOT NULL DEFAULT false, + started_at timestamp, + applied_at timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP + ) + SQL + end + + def mark_in_progress(connection) + connection.exec_params(<<~SQL, [SCHEMA_VERSION]) + INSERT INTO public.en57_schema_info (id, schema_version, in_progress, started_at) + VALUES (1, $1, true, CURRENT_TIMESTAMP) + ON CONFLICT (id) DO UPDATE + SET in_progress = true, + started_at = CURRENT_TIMESTAMP + SQL + end + + def record_version(connection) + connection.exec_params(<<~SQL, [SCHEMA_VERSION]) + UPDATE public.en57_schema_info + SET schema_version = $1, + in_progress = false, + applied_at = CURRENT_TIMESTAMP + WHERE id = 1 + SQL + end + + def schema_info_table?(connection) + connection + .exec_params("SELECT to_regclass($1)::text", ["public.en57_schema_info"]) + .first + .fetch("to_regclass") + end + + def schema_path(version) + File.expand_path("../../db/schema/#{version}.sql", __dir__) + end + + def partial_migration_error + MigrationError.new(warning_for(:partial)) + end + + def with_connection + connection = PG.connect(@connection_string) + yield connection + ensure + connection&.close + end + end +end diff --git a/lib/en57/tasks.rb b/lib/en57/tasks.rb new file mode 100644 index 0000000..567936f --- /dev/null +++ b/lib/en57/tasks.rb @@ -0,0 +1,26 @@ +# frozen_string_literal: true + +require "rake" +require_relative "migrator" + +namespace :en57 do + desc "Apply pending En57 schema migrations" + task :migrate do + En57::Migrator.new(ENV.fetch("DATABASE_URL")).migrate! + end + + desc "List En57 schema migration status" + task :status do + status = En57::Migrator.new(ENV.fetch("DATABASE_URL")).status + + puts "Current : #{status.current}" + puts "Target : #{status.target}" + puts "Status : #{status.state}" + puts "Method : #{status.method}" + unless status.pending.empty? + puts "Pending :" + status.pending.each { |item| puts " #{item}" } + end + puts "WARNING : #{status.warning}" if status.warning + end +end diff --git a/lib/en57/version.rb b/lib/en57/version.rb index ac83b29..1af98b8 100644 --- a/lib/en57/version.rb +++ b/lib/en57/version.rb @@ -2,4 +2,5 @@ module En57 VERSION = "0.1.0" + SCHEMA_VERSION = "0.1.0" end diff --git a/test/test_migrator.rb b/test/test_migrator.rb new file mode 100644 index 0000000..5816f69 --- /dev/null +++ b/test/test_migrator.rb @@ -0,0 +1,128 @@ +# frozen_string_literal: true + +require "test_helper" +require "uri" + +module En57 + class TestMigrator < IntegrationTest + def test_status_reports_pending_schema_on_empty_database + with_database do |url| + assert_equal( + Migrator::Status.new( + current: :fresh, + target: "0.1.0", + state: :pending, + method: :none, + pending: [schema_path("0.1.0")], + warning: nil, + ), + Migrator.new(url).status, + ) + end + end + + def test_migrate_applies_pending_schema_and_records_version + with_database do |url| + migrator = Migrator.new(url) + + migrator.migrate! + migrator.migrate! + + assert_equal( + Migrator::Status.new( + current: "0.1.0", + target: "0.1.0", + state: :up_to_date, + method: :version_table, + pending: [], + warning: nil, + ), + migrator.status, + ) + end + end + + def test_migrate_installs_schema_used_by_event_store + with_database do |url| + Migrator.new(url).migrate! + event = Event.new(type: "Migrated") + connection = PG.connect(url) + event_store = + EventStore.new( + Repository.new(PgAdapter.new(connection), JsonSerializer.new), + ) + + assert_equal [event], event_store.append([event]).read.each.to_a + ensure + connection&.close + end + end + + def test_migrate_leaves_partial_status_after_failure + with_database do |url| + connection = PG.connect(url) + connection.exec("CREATE SCHEMA en57") + connection.close + + assert_raises(PG::DuplicateSchema) { Migrator.new(url).migrate! } + + status = Migrator.new(url).status + assert_equal :partial, status.current + assert_equal :partial, status.state + assert_match( + /Previous En57 migration did not complete cleanly/, + status.warning, + ) + end + end + + def test_migrate_rejects_partial_status + with_database do |url| + connection = PG.connect(url) + connection.exec(<<~SQL) + CREATE TABLE public.en57_schema_info ( + id integer PRIMARY KEY, + schema_version varchar(20) NOT NULL, + in_progress boolean NOT NULL DEFAULT false, + started_at timestamp, + applied_at timestamp NOT NULL DEFAULT CURRENT_TIMESTAMP + ) + SQL + connection.exec(<<~SQL) + INSERT INTO public.en57_schema_info (id, schema_version, in_progress) + VALUES (1, '0.1.0', true) + SQL + connection.close + + error = + assert_raises(Migrator::MigrationError) { Migrator.new(url).migrate! } + assert_match( + /Previous En57 migration did not complete cleanly/, + error.message, + ) + end + end + + 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) + uri = URI(SERVER.url) + uri.path = "/#{name}" + uri.to_s + end + + def schema_path(version) + File.expand_path("../db/schema/#{version}.sql", __dir__) + end + end +end