diff --git a/priv/repo/migration_application_compatibility.tsv b/priv/repo/migration_application_compatibility.tsv index c4d0ed0..6a2af8b 100644 --- a/priv/repo/migration_application_compatibility.tsv +++ b/priv/repo/migration_application_compatibility.tsv @@ -13,3 +13,4 @@ 20260813194249 application_safe additive nullable inbound-provider identity columns and a partial unique index are ignored by the previous application 20260827173037 application_safe additive discovery rollup tables and maintenance triggers are ignored by the previous application while continuing to follow its writes 20260827180241 forward_only pending lifecycle uniqueness can reject duplicate account access or deletion requests that the previous application may still attempt to insert +20260827225430 application_safe replacing rollup trigger functions preserves the existing schema and application contract while serializing concurrent counter inserts diff --git a/priv/repo/migrations/20260827173037_add_discovery_rollups.exs b/priv/repo/migrations/20260827173037_add_discovery_rollups.exs index 80c7942..f48b69d 100644 --- a/priv/repo/migrations/20260827173037_add_discovery_rollups.exs +++ b/priv/repo/migrations/20260827173037_add_discovery_rollups.exs @@ -293,23 +293,63 @@ defmodule WhoNeedHelp.Repo.Migrations.AddDiscoveryRollups do GROUP BY level, cell_x, cell_y, category_id, urgency, visible_until_hour HAVING sum(delta) <> 0 ) - MERGE INTO help_discovery_rollups AS target - USING deltas AS source - ON target.level = source.level - AND target.cell_x = source.cell_x - AND target.cell_y = source.cell_y - AND target.category_id = source.category_id - AND target.urgency = source.urgency - AND target.visible_until_hour = source.visible_until_hour - WHEN MATCHED AND target.point_count + source.point_count > 0 THEN - UPDATE SET point_count = target.point_count + source.point_count - WHEN MATCHED THEN DELETE - WHEN NOT MATCHED AND source.point_count > 0 THEN - INSERT - (level, cell_x, cell_y, category_id, urgency, visible_until_hour, point_count) - VALUES - (source.level, source.cell_x, source.cell_y, source.category_id, - source.urgency, source.visible_until_hour, source.point_count); + INSERT INTO help_discovery_rollups + (level, cell_x, cell_y, category_id, urgency, visible_until_hour, point_count) + SELECT + level, cell_x, cell_y, category_id, urgency, visible_until_hour, point_count + FROM deltas + WHERE point_count > 0 + ON CONFLICT + (level, cell_x, cell_y, category_id, urgency, visible_until_hour) + DO UPDATE SET + point_count = help_discovery_rollups.point_count + EXCLUDED.point_count; + + WITH changes AS ( + #{changes_sql} + ), expanded AS ( + SELECT + changes.delta, + level, + #{cell_x_sql("changes.public_longitude", "level")}, + #{cell_y_sql("changes.public_latitude", "level")}, + changes.category_id, + changes.urgency, + changes.visible_until_hour + FROM changes + CROSS JOIN generate_series(0, #{@maximum_rollup_level}) AS level + ), deltas AS ( + SELECT + level, + cell_x, + cell_y, + category_id, + urgency, + visible_until_hour, + sum(delta)::bigint AS point_count + FROM expanded + GROUP BY level, cell_x, cell_y, category_id, urgency, visible_until_hour + HAVING sum(delta) < 0 + ), decremented AS ( + UPDATE help_discovery_rollups AS target + SET point_count = target.point_count + source.point_count + FROM deltas AS source + WHERE target.level = source.level + AND target.cell_x = source.cell_x + AND target.cell_y = source.cell_y + AND target.category_id = source.category_id + AND target.urgency = source.urgency + AND target.visible_until_hour = source.visible_until_hour + RETURNING target.* + ) + DELETE FROM help_discovery_rollups AS target + USING decremented + WHERE target.level = decremented.level + AND target.cell_x = decremented.cell_x + AND target.cell_y = decremented.cell_y + AND target.category_id = decremented.category_id + AND target.urgency = decremented.urgency + AND target.visible_until_hour = decremented.visible_until_hour + AND decremented.point_count = 0; RETURN NULL; END @@ -348,21 +388,57 @@ defmodule WhoNeedHelp.Repo.Migrations.AddDiscoveryRollups do GROUP BY level, cell_x, cell_y, category_id, visible_until_hour HAVING sum(delta) <> 0 ) - MERGE INTO activity_discovery_rollups AS target - USING deltas AS source - ON target.level = source.level - AND target.cell_x = source.cell_x - AND target.cell_y = source.cell_y - AND target.category_id = source.category_id - AND target.visible_until_hour = source.visible_until_hour - WHEN MATCHED AND target.point_count + source.point_count > 0 THEN - UPDATE SET point_count = target.point_count + source.point_count - WHEN MATCHED THEN DELETE - WHEN NOT MATCHED AND source.point_count > 0 THEN - INSERT (level, cell_x, cell_y, category_id, visible_until_hour, point_count) - VALUES - (source.level, source.cell_x, source.cell_y, source.category_id, - source.visible_until_hour, source.point_count); + INSERT INTO activity_discovery_rollups + (level, cell_x, cell_y, category_id, visible_until_hour, point_count) + SELECT level, cell_x, cell_y, category_id, visible_until_hour, point_count + FROM deltas + WHERE point_count > 0 + ON CONFLICT (level, cell_x, cell_y, category_id, visible_until_hour) + DO UPDATE SET + point_count = activity_discovery_rollups.point_count + EXCLUDED.point_count; + + WITH changes AS ( + #{changes_sql} + ), expanded AS ( + SELECT + changes.delta, + level, + #{cell_x_sql("changes.public_longitude", "level")}, + #{cell_y_sql("changes.public_latitude", "level")}, + changes.category_id, + changes.visible_until_hour + FROM changes + CROSS JOIN generate_series(0, #{@maximum_rollup_level}) AS level + ), deltas AS ( + SELECT + level, + cell_x, + cell_y, + category_id, + visible_until_hour, + sum(delta)::bigint AS point_count + FROM expanded + GROUP BY level, cell_x, cell_y, category_id, visible_until_hour + HAVING sum(delta) < 0 + ), decremented AS ( + UPDATE activity_discovery_rollups AS target + SET point_count = target.point_count + source.point_count + FROM deltas AS source + WHERE target.level = source.level + AND target.cell_x = source.cell_x + AND target.cell_y = source.cell_y + AND target.category_id = source.category_id + AND target.visible_until_hour = source.visible_until_hour + RETURNING target.* + ) + DELETE FROM activity_discovery_rollups AS target + USING decremented + WHERE target.level = decremented.level + AND target.cell_x = decremented.cell_x + AND target.cell_y = decremented.cell_y + AND target.category_id = decremented.category_id + AND target.visible_until_hour = decremented.visible_until_hour + AND decremented.point_count = 0; RETURN NULL; END diff --git a/priv/repo/migrations/20260827225430_make_discovery_rollups_concurrency_safe.exs b/priv/repo/migrations/20260827225430_make_discovery_rollups_concurrency_safe.exs new file mode 100644 index 0000000..3c533ed --- /dev/null +++ b/priv/repo/migrations/20260827225430_make_discovery_rollups_concurrency_safe.exs @@ -0,0 +1,277 @@ +defmodule WhoNeedHelp.Repo.Migrations.MakeDiscoveryRollupsConcurrencySafe do + use Ecto.Migration + + @maximum_rollup_level 12 + + def change do + for {operation, changes_sql} <- help_operations() do + sql = help_trigger_function(operation, changes_sql) + execute(sql, sql) + end + + for {operation, changes_sql} <- activity_operations() do + sql = activity_trigger_function(operation, changes_sql) + execute(sql, sql) + end + end + + defp help_operations do + [ + {"insert", help_change_source("new_rows", 1)}, + {"delete", help_change_source("old_rows", -1)}, + {"update", + help_change_source("new_rows", 1) <> + " UNION ALL " <> help_change_source("old_rows", -1)} + ] + end + + defp activity_operations do + [ + {"insert", activity_change_source("new_rows", 1)}, + {"delete", activity_change_source("old_rows", -1)}, + {"update", + activity_change_source("new_rows", 1) <> + " UNION ALL " <> activity_change_source("old_rows", -1)} + ] + end + + defp help_change_source(relation, delta) do + """ + SELECT + #{delta}::bigint AS delta, + request.public_longitude, + request.public_latitude, + request.category_id, + request.urgency, + date_trunc('hour', request.expires_at) AS visible_until_hour + FROM #{relation} AS request + WHERE request.status = 'open' + AND request.hidden_at IS NULL + AND request.location_visibility <> 'hidden' + AND request.public_longitude IS NOT NULL + AND request.public_latitude IS NOT NULL + AND request.category_id IS NOT NULL + AND request.urgency IS NOT NULL + AND request.expires_at IS NOT NULL + """ + end + + defp activity_change_source(relation, delta) do + """ + SELECT + #{delta}::bigint AS delta, + activity.public_longitude, + activity.public_latitude, + activity.category_id, + date_trunc('hour', least(activity.starts_at, activity.join_deadline)) + AS visible_until_hour + FROM #{relation} AS activity + WHERE activity.status = 'open' + AND activity.hidden_at IS NULL + AND activity.location_visibility = 'approximate_public' + AND activity.public_longitude IS NOT NULL + AND activity.public_latitude IS NOT NULL + AND activity.category_id IS NOT NULL + AND activity.starts_at IS NOT NULL + AND activity.join_deadline IS NOT NULL + """ + end + + defp help_trigger_function(operation, changes_sql) do + """ + CREATE OR REPLACE FUNCTION wnh_help_discovery_rollups_#{operation}() + RETURNS trigger + LANGUAGE plpgsql + AS $function$ + BEGIN + #{help_positive_deltas_sql(changes_sql)} + #{help_negative_deltas_sql(changes_sql)} + RETURN NULL; + END + $function$ + """ + end + + defp activity_trigger_function(operation, changes_sql) do + """ + CREATE OR REPLACE FUNCTION wnh_activity_discovery_rollups_#{operation}() + RETURNS trigger + LANGUAGE plpgsql + AS $function$ + BEGIN + #{activity_positive_deltas_sql(changes_sql)} + #{activity_negative_deltas_sql(changes_sql)} + RETURN NULL; + END + $function$ + """ + end + + defp help_positive_deltas_sql(changes_sql) do + """ + #{help_deltas_cte(changes_sql, "sum(delta) <> 0")} + INSERT INTO help_discovery_rollups + (level, cell_x, cell_y, category_id, urgency, visible_until_hour, point_count) + SELECT level, cell_x, cell_y, category_id, urgency, visible_until_hour, point_count + FROM deltas + WHERE point_count > 0 + ON CONFLICT (level, cell_x, cell_y, category_id, urgency, visible_until_hour) + DO UPDATE SET + point_count = help_discovery_rollups.point_count + EXCLUDED.point_count; + """ + end + + defp help_negative_deltas_sql(changes_sql) do + """ + #{help_deltas_cte(changes_sql, "sum(delta) < 0")}, decremented AS ( + UPDATE help_discovery_rollups AS target + SET point_count = target.point_count + source.point_count + FROM deltas AS source + WHERE target.level = source.level + AND target.cell_x = source.cell_x + AND target.cell_y = source.cell_y + AND target.category_id = source.category_id + AND target.urgency = source.urgency + AND target.visible_until_hour = source.visible_until_hour + RETURNING target.* + ) + DELETE FROM help_discovery_rollups AS target + USING decremented + WHERE target.level = decremented.level + AND target.cell_x = decremented.cell_x + AND target.cell_y = decremented.cell_y + AND target.category_id = decremented.category_id + AND target.urgency = decremented.urgency + AND target.visible_until_hour = decremented.visible_until_hour + AND decremented.point_count = 0; + """ + end + + defp activity_positive_deltas_sql(changes_sql) do + """ + #{activity_deltas_cte(changes_sql, "sum(delta) <> 0")} + INSERT INTO activity_discovery_rollups + (level, cell_x, cell_y, category_id, visible_until_hour, point_count) + SELECT level, cell_x, cell_y, category_id, visible_until_hour, point_count + FROM deltas + WHERE point_count > 0 + ON CONFLICT (level, cell_x, cell_y, category_id, visible_until_hour) + DO UPDATE SET + point_count = activity_discovery_rollups.point_count + EXCLUDED.point_count; + """ + end + + defp activity_negative_deltas_sql(changes_sql) do + """ + #{activity_deltas_cte(changes_sql, "sum(delta) < 0")}, decremented AS ( + UPDATE activity_discovery_rollups AS target + SET point_count = target.point_count + source.point_count + FROM deltas AS source + WHERE target.level = source.level + AND target.cell_x = source.cell_x + AND target.cell_y = source.cell_y + AND target.category_id = source.category_id + AND target.visible_until_hour = source.visible_until_hour + RETURNING target.* + ) + DELETE FROM activity_discovery_rollups AS target + USING decremented + WHERE target.level = decremented.level + AND target.cell_x = decremented.cell_x + AND target.cell_y = decremented.cell_y + AND target.category_id = decremented.category_id + AND target.visible_until_hour = decremented.visible_until_hour + AND decremented.point_count = 0; + """ + end + + defp help_deltas_cte(changes_sql, having) do + """ + WITH changes AS ( + #{changes_sql} + ), expanded AS ( + SELECT + changes.delta, + level, + #{cell_x_sql("changes.public_longitude", "level")}, + #{cell_y_sql("changes.public_latitude", "level")}, + changes.category_id, + changes.urgency, + changes.visible_until_hour + FROM changes + CROSS JOIN generate_series(0, #{@maximum_rollup_level}) AS level + ), deltas AS ( + SELECT + level, + cell_x, + cell_y, + category_id, + urgency, + visible_until_hour, + sum(delta)::bigint AS point_count + FROM expanded + GROUP BY level, cell_x, cell_y, category_id, urgency, visible_until_hour + HAVING #{having} + ) + """ + end + + defp activity_deltas_cte(changes_sql, having) do + """ + WITH changes AS ( + #{changes_sql} + ), expanded AS ( + SELECT + changes.delta, + level, + #{cell_x_sql("changes.public_longitude", "level")}, + #{cell_y_sql("changes.public_latitude", "level")}, + changes.category_id, + changes.visible_until_hour + FROM changes + CROSS JOIN generate_series(0, #{@maximum_rollup_level}) AS level + ), deltas AS ( + SELECT + level, + cell_x, + cell_y, + category_id, + visible_until_hour, + sum(delta)::bigint AS point_count + FROM expanded + GROUP BY level, cell_x, cell_y, category_id, visible_until_hour + HAVING #{having} + ) + """ + end + + defp cell_x_sql(longitude, level) do + """ + least( + power(2, #{level})::bigint - 1, + greatest(0, floor(((#{longitude} + 180.0) / 360.0) * power(2, #{level})))::bigint + ) AS cell_x + """ + end + + defp cell_y_sql(latitude, level) do + """ + least( + power(2, #{level})::bigint - 1, + greatest( + 0, + floor( + ( + ln( + tan( + pi() / 4.0 + + radians(least(85.05112878, greatest(-85.05112878, #{latitude}))) / 2.0 + ) + ) / pi() + 1.0 + ) / 2.0 * power(2, #{level}) + )::bigint + ) + ) AS cell_y + """ + end +end diff --git a/test/who_need_help/discovery_rollup_concurrency_test.exs b/test/who_need_help/discovery_rollup_concurrency_test.exs new file mode 100644 index 0000000..027238c --- /dev/null +++ b/test/who_need_help/discovery_rollup_concurrency_test.exs @@ -0,0 +1,216 @@ +defmodule WhoNeedHelp.DiscoveryRollupConcurrencyTest do + use ExUnit.Case, async: false + + import Ecto.Query + + alias Ecto.Adapters.SQL.Sandbox + alias WhoNeedHelp.Accounts + alias WhoNeedHelp.Activities.Activity + alias WhoNeedHelp.Catalog.Category + alias WhoNeedHelp.Help.HelpRequest + alias WhoNeedHelp.Repo + + @transaction_timeout 10_000 + + test "concurrent help-request inserts serialize the first rollup bucket" do + %{category: category, user: user} = setup_entities(:help) + + on_exit(fn -> cleanup_entities(category, user) end) + + expires_at = future_hour(48) + + first = fn -> + insert_help_request(user.id, category.id, expires_at, "Concurrent request one") + end + + second = fn -> + insert_help_request(user.id, category.id, expires_at, "Concurrent request two") + end + + assert_concurrent_first_bucket(first, second) + + count = + Sandbox.unboxed_run(Repo, fn -> + rollup_count( + "help_discovery_rollups", + category.id, + expires_at, + "urgency = 'now'" + ) + end) + + assert count == 2 + end + + test "concurrent activity inserts serialize the first rollup bucket" do + %{category: category, user: user} = setup_entities(:activity) + + on_exit(fn -> cleanup_entities(category, user) end) + + join_deadline = future_hour(48) + starts_at = DateTime.add(join_deadline, 1, :hour) + + first = fn -> + insert_activity(user.id, category.id, starts_at, join_deadline, "Concurrent activity one") + end + + second = fn -> + insert_activity(user.id, category.id, starts_at, join_deadline, "Concurrent activity two") + end + + assert_concurrent_first_bucket(first, second) + + count = + Sandbox.unboxed_run(Repo, fn -> + rollup_count( + "activity_discovery_rollups", + category.id, + join_deadline, + "TRUE" + ) + end) + + assert count == 2 + end + + defp assert_concurrent_first_bucket(first_insert, second_insert) do + parent = self() + + first_task = + Task.async(fn -> + with_unboxed_connection(fn -> + Repo.transaction( + fn -> + first_insert.() + send(parent, {:first_inserted, self()}) + + receive do + :commit_first -> :ok + after + @transaction_timeout -> raise "timed out waiting to commit first insert" + end + end, + timeout: @transaction_timeout + ) + end) + end) + + assert_receive {:first_inserted, first_pid}, @transaction_timeout + + second_task = + Task.async(fn -> + with_unboxed_connection(fn -> + Repo.transaction(second_insert, timeout: @transaction_timeout) + end) + end) + + assert Task.yield(second_task, 100) == nil + send(first_pid, :commit_first) + + assert {:ok, :ok} = Task.await(first_task, @transaction_timeout) + assert {:ok, _record} = Task.await(second_task, @transaction_timeout) + end + + defp with_unboxed_connection(fun) do + :ok = Sandbox.checkout(Repo, sandbox: false) + + try do + fun.() + after + Sandbox.checkin(Repo) + end + end + + defp setup_entities(mode) do + Sandbox.unboxed_run(Repo, fn -> + suffix = System.unique_integer([:positive, :monotonic]) + + {:ok, user} = + Accounts.register_user(%{ + email: "rollup-concurrency-#{suffix}@example.invalid", + display_name: "Rollup concurrency test", + terms_accepted: true + }) + + category = + %Category{} + |> Category.changeset(%{ + slug: "rollup-concurrency-#{mode}-#{suffix}", + names: %{"en" => "Rollup concurrency test"}, + mode: mode + }) + |> Repo.insert!() + + %{category: category, user: user} + end) + end + + defp cleanup_entities(category, user) do + Sandbox.unboxed_run(Repo, fn -> + Repo.delete_all(from request in HelpRequest, where: request.category_id == ^category.id) + Repo.delete_all(from activity in Activity, where: activity.category_id == ^category.id) + Repo.delete!(category) + Repo.delete!(user) + end) + end + + defp insert_help_request(user_id, category_id, expires_at, title) do + %HelpRequest{} + |> Ecto.Changeset.change(%{ + title: title, + description: "A concurrent request used to verify rollup serialization.", + location_label: "Concurrency test point", + location: %Geo.Point{coordinates: {30.5234, 50.4501}, srid: 4326}, + location_visibility: :exact_public, + location_radius_meters: nil, + urgency: :now, + status: :open, + expires_at: expires_at, + requester_id: user_id, + category_id: category_id + }) + |> Repo.insert!() + end + + defp insert_activity(user_id, category_id, starts_at, join_deadline, title) do + %Activity{} + |> Ecto.Changeset.change(%{ + title: title, + description: "A concurrent activity used to verify rollup serialization.", + location_label: "Concurrency test point", + location: %Geo.Point{coordinates: {30.5234, 50.4501}, srid: 4326}, + location_visibility: :approximate_public, + status: :open, + starts_at: starts_at, + join_deadline: join_deadline, + capacity: 4, + creator_id: user_id, + category_id: category_id + }) + |> Repo.insert!() + end + + defp rollup_count(table, category_id, visible_until, additional_predicate) do + query = """ + SELECT point_count + FROM #{table} + WHERE level = 0 + AND category_id = $1 + AND visible_until_hour = $2 + AND #{additional_predicate} + """ + + category_id = Ecto.UUID.dump!(category_id) + %{rows: [[count]]} = Repo.query!(query, [category_id, visible_until]) + count + end + + defp future_hour(hours) do + now = DateTime.utc_now(:second) + seconds_to_next_hour = 3_600 - (now.minute * 60 + now.second) + + now + |> DateTime.add(seconds_to_next_hour + hours * 3_600, :second) + |> DateTime.truncate(:second) + end +end