Make discovery rollups concurrency safe

This commit is contained in:
SimpleTest 2026-08-28 02:26:12 +03:00
parent 4197cb8626
commit 5576b815df
4 changed files with 602 additions and 32 deletions

View File

@ -13,3 +13,4 @@
20260813194249 application_safe additive nullable inbound-provider identity columns and a partial unique index are ignored by the previous application 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 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 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

1 # migration_version application_rollback_policy reason
13 20260813194249 application_safe additive nullable inbound-provider identity columns and a partial unique index are ignored by the previous application
14 20260827173037 application_safe additive discovery rollup tables and maintenance triggers are ignored by the previous application while continuing to follow its writes
15 20260827180241 forward_only pending lifecycle uniqueness can reject duplicate account access or deletion requests that the previous application may still attempt to insert
16 20260827225430 application_safe replacing rollup trigger functions preserves the existing schema and application contract while serializing concurrent counter inserts

View File

@ -293,23 +293,63 @@ defmodule WhoNeedHelp.Repo.Migrations.AddDiscoveryRollups do
GROUP BY level, cell_x, cell_y, category_id, urgency, visible_until_hour GROUP BY level, cell_x, cell_y, category_id, urgency, visible_until_hour
HAVING sum(delta) <> 0 HAVING sum(delta) <> 0
) )
MERGE INTO help_discovery_rollups AS target INSERT INTO help_discovery_rollups
USING deltas AS source (level, cell_x, cell_y, category_id, urgency, visible_until_hour, point_count)
ON target.level = source.level SELECT
AND target.cell_x = source.cell_x level, cell_x, cell_y, category_id, urgency, visible_until_hour, point_count
AND target.cell_y = source.cell_y FROM deltas
AND target.category_id = source.category_id WHERE point_count > 0
AND target.urgency = source.urgency ON CONFLICT
AND target.visible_until_hour = source.visible_until_hour (level, cell_x, cell_y, category_id, urgency, visible_until_hour)
WHEN MATCHED AND target.point_count + source.point_count > 0 THEN DO UPDATE SET
UPDATE SET point_count = target.point_count + source.point_count point_count = help_discovery_rollups.point_count + EXCLUDED.point_count;
WHEN MATCHED THEN DELETE
WHEN NOT MATCHED AND source.point_count > 0 THEN WITH changes AS (
INSERT #{changes_sql}
(level, cell_x, cell_y, category_id, urgency, visible_until_hour, point_count) ), expanded AS (
VALUES SELECT
(source.level, source.cell_x, source.cell_y, source.category_id, changes.delta,
source.urgency, source.visible_until_hour, source.point_count); 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; RETURN NULL;
END END
@ -348,21 +388,57 @@ defmodule WhoNeedHelp.Repo.Migrations.AddDiscoveryRollups do
GROUP BY level, cell_x, cell_y, category_id, visible_until_hour GROUP BY level, cell_x, cell_y, category_id, visible_until_hour
HAVING sum(delta) <> 0 HAVING sum(delta) <> 0
) )
MERGE INTO activity_discovery_rollups AS target INSERT INTO activity_discovery_rollups
USING deltas AS source (level, cell_x, cell_y, category_id, visible_until_hour, point_count)
ON target.level = source.level SELECT level, cell_x, cell_y, category_id, visible_until_hour, point_count
AND target.cell_x = source.cell_x FROM deltas
AND target.cell_y = source.cell_y WHERE point_count > 0
AND target.category_id = source.category_id ON CONFLICT (level, cell_x, cell_y, category_id, visible_until_hour)
AND target.visible_until_hour = source.visible_until_hour DO UPDATE SET
WHEN MATCHED AND target.point_count + source.point_count > 0 THEN point_count = activity_discovery_rollups.point_count + EXCLUDED.point_count;
UPDATE SET point_count = target.point_count + source.point_count
WHEN MATCHED THEN DELETE WITH changes AS (
WHEN NOT MATCHED AND source.point_count > 0 THEN #{changes_sql}
INSERT (level, cell_x, cell_y, category_id, visible_until_hour, point_count) ), expanded AS (
VALUES SELECT
(source.level, source.cell_x, source.cell_y, source.category_id, changes.delta,
source.visible_until_hour, source.point_count); 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; RETURN NULL;
END END

View File

@ -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

View File

@ -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