who_need_help/priv/repo/migrations/20260827173037_add_discovery_rollups.exs

546 lines
17 KiB
Elixir

defmodule WhoNeedHelp.Repo.Migrations.AddDiscoveryRollups do
use Ecto.Migration
@maximum_rollup_level 12
def up do
create table(:help_discovery_rollups, primary_key: false) do
add(:level, :smallint, null: false)
add(:cell_x, :integer, null: false)
add(:cell_y, :integer, null: false)
add(:category_id, :binary_id, null: false)
add(:urgency, :string, null: false)
add(:visible_until_hour, :utc_datetime, null: false)
add(:point_count, :bigint, null: false)
end
create(
constraint(:help_discovery_rollups, :help_discovery_rollups_level_range,
check: "level BETWEEN 0 AND #{@maximum_rollup_level}"
)
)
create(
constraint(:help_discovery_rollups, :help_discovery_rollups_non_negative_count,
check: "point_count >= 0"
)
)
create table(:activity_discovery_rollups, primary_key: false) do
add(:level, :smallint, null: false)
add(:cell_x, :integer, null: false)
add(:cell_y, :integer, null: false)
add(:category_id, :binary_id, null: false)
add(:visible_until_hour, :utc_datetime, null: false)
add(:point_count, :bigint, null: false)
end
create(
constraint(:activity_discovery_rollups, :activity_discovery_rollups_level_range,
check: "level BETWEEN 0 AND #{@maximum_rollup_level}"
)
)
create(
constraint(
:activity_discovery_rollups,
:activity_discovery_rollups_non_negative_count,
check: "point_count >= 0"
)
)
# The initial snapshot and trigger installation must describe one database
# state. Reads remain available, while writes to either discovery source wait
# until this migration transaction commits.
execute("LOCK TABLE help_requests, activities IN SHARE ROW EXCLUSIVE MODE")
execute(rebuild_help_rollups_sql())
execute(rebuild_activity_rollups_sql())
execute(validate_initial_rollups_sql())
create(
unique_index(
:help_discovery_rollups,
[:level, :cell_x, :cell_y, :category_id, :urgency, :visible_until_hour],
name: :help_discovery_rollups_identity
)
)
create(
index(
:help_discovery_rollups,
[:level, :cell_x, :cell_y, :visible_until_hour],
include: [:category_id, :urgency, :point_count],
name: :help_discovery_rollups_viewport
)
)
create(
index(
:help_discovery_rollups,
[:level, :category_id, :urgency, :cell_x, :cell_y, :visible_until_hour],
include: [:point_count],
name: :help_discovery_rollups_filtered_viewport
)
)
create(
unique_index(
:activity_discovery_rollups,
[:level, :cell_x, :cell_y, :category_id, :visible_until_hour],
name: :activity_discovery_rollups_identity
)
)
create(
index(
:activity_discovery_rollups,
[:level, :cell_x, :cell_y, :visible_until_hour],
include: [:category_id, :point_count],
name: :activity_discovery_rollups_viewport
)
)
create(
index(
:activity_discovery_rollups,
[:level, :category_id, :cell_x, :cell_y, :visible_until_hour],
include: [:point_count],
name: :activity_discovery_rollups_filtered_viewport
)
)
execute(help_trigger_function("insert", help_change_source("new_rows", 1)))
execute(help_trigger_function("delete", help_change_source("old_rows", -1)))
execute(
help_trigger_function(
"update",
help_change_source("new_rows", 1) <>
" UNION ALL " <>
help_change_source("old_rows", -1)
)
)
execute(activity_trigger_function("insert", activity_change_source("new_rows", 1)))
execute(activity_trigger_function("delete", activity_change_source("old_rows", -1)))
execute(
activity_trigger_function(
"update",
activity_change_source("new_rows", 1) <>
" UNION ALL " <>
activity_change_source("old_rows", -1)
)
)
execute("""
CREATE TRIGGER help_discovery_rollups_after_insert
AFTER INSERT ON help_requests
REFERENCING NEW TABLE AS new_rows
FOR EACH STATEMENT EXECUTE FUNCTION wnh_help_discovery_rollups_insert()
""")
execute("""
CREATE TRIGGER help_discovery_rollups_after_delete
AFTER DELETE ON help_requests
REFERENCING OLD TABLE AS old_rows
FOR EACH STATEMENT EXECUTE FUNCTION wnh_help_discovery_rollups_delete()
""")
execute("""
CREATE TRIGGER help_discovery_rollups_after_update
AFTER UPDATE ON help_requests
REFERENCING OLD TABLE AS old_rows NEW TABLE AS new_rows
FOR EACH STATEMENT EXECUTE FUNCTION wnh_help_discovery_rollups_update()
""")
execute("""
CREATE TRIGGER activity_discovery_rollups_after_insert
AFTER INSERT ON activities
REFERENCING NEW TABLE AS new_rows
FOR EACH STATEMENT EXECUTE FUNCTION wnh_activity_discovery_rollups_insert()
""")
execute("""
CREATE TRIGGER activity_discovery_rollups_after_delete
AFTER DELETE ON activities
REFERENCING OLD TABLE AS old_rows
FOR EACH STATEMENT EXECUTE FUNCTION wnh_activity_discovery_rollups_delete()
""")
execute("""
CREATE TRIGGER activity_discovery_rollups_after_update
AFTER UPDATE ON activities
REFERENCING OLD TABLE AS old_rows NEW TABLE AS new_rows
FOR EACH STATEMENT EXECUTE FUNCTION wnh_activity_discovery_rollups_update()
""")
end
def down do
execute("DROP TRIGGER IF EXISTS activity_discovery_rollups_after_update ON activities")
execute("DROP TRIGGER IF EXISTS activity_discovery_rollups_after_delete ON activities")
execute("DROP TRIGGER IF EXISTS activity_discovery_rollups_after_insert ON activities")
execute("DROP TRIGGER IF EXISTS help_discovery_rollups_after_update ON help_requests")
execute("DROP TRIGGER IF EXISTS help_discovery_rollups_after_delete ON help_requests")
execute("DROP TRIGGER IF EXISTS help_discovery_rollups_after_insert ON help_requests")
for operation <- ~w(insert delete update) do
execute("DROP FUNCTION IF EXISTS wnh_activity_discovery_rollups_#{operation}()")
execute("DROP FUNCTION IF EXISTS wnh_help_discovery_rollups_#{operation}()")
end
drop(table(:activity_discovery_rollups))
drop(table(:help_discovery_rollups))
end
defp rebuild_help_rollups_sql do
"""
INSERT INTO help_discovery_rollups
(level, cell_x, cell_y, category_id, urgency, visible_until_hour, point_count)
SELECT
level,
#{cell_x_sql("request.public_longitude", "level")},
#{cell_y_sql("request.public_latitude", "level")},
request.category_id,
request.urgency,
date_trunc('hour', request.expires_at),
count(*)::bigint
FROM help_requests AS request
CROSS JOIN generate_series(0, #{@maximum_rollup_level}) AS level
WHERE #{eligible_help_sql("request")}
GROUP BY level, cell_x, cell_y, request.category_id, request.urgency,
date_trunc('hour', request.expires_at)
"""
end
defp rebuild_activity_rollups_sql do
"""
INSERT INTO activity_discovery_rollups
(level, cell_x, cell_y, category_id, visible_until_hour, point_count)
SELECT
level,
#{cell_x_sql("activity.public_longitude", "level")},
#{cell_y_sql("activity.public_latitude", "level")},
activity.category_id,
date_trunc('hour', least(activity.starts_at, activity.join_deadline)),
count(*)::bigint
FROM activities AS activity
CROSS JOIN generate_series(0, #{@maximum_rollup_level}) AS level
WHERE #{eligible_activity_sql("activity")}
GROUP BY level, cell_x, cell_y, activity.category_id,
date_trunc('hour', least(activity.starts_at, activity.join_deadline))
"""
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 #{eligible_help_sql("request")}
"""
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 #{eligible_activity_sql("activity")}
"""
end
defp help_trigger_function(operation, changes_sql) do
"""
CREATE FUNCTION wnh_help_discovery_rollups_#{operation}()
RETURNS trigger
LANGUAGE plpgsql
AS $function$
BEGIN
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
)
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
$function$
"""
end
defp activity_trigger_function(operation, changes_sql) do
"""
CREATE FUNCTION wnh_activity_discovery_rollups_#{operation}()
RETURNS trigger
LANGUAGE plpgsql
AS $function$
BEGIN
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
)
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
$function$
"""
end
defp eligible_help_sql(alias_name) do
"""
#{alias_name}.status = 'open'
AND #{alias_name}.hidden_at IS NULL
AND #{alias_name}.location_visibility <> 'hidden'
AND #{alias_name}.public_longitude IS NOT NULL
AND #{alias_name}.public_latitude IS NOT NULL
AND #{alias_name}.category_id IS NOT NULL
AND #{alias_name}.urgency IS NOT NULL
AND #{alias_name}.expires_at IS NOT NULL
"""
end
defp eligible_activity_sql(alias_name) do
"""
#{alias_name}.status = 'open'
AND #{alias_name}.hidden_at IS NULL
AND #{alias_name}.location_visibility = 'approximate_public'
AND #{alias_name}.public_longitude IS NOT NULL
AND #{alias_name}.public_latitude IS NOT NULL
AND #{alias_name}.category_id IS NOT NULL
AND #{alias_name}.starts_at IS NOT NULL
AND #{alias_name}.join_deadline IS NOT NULL
"""
end
defp validate_initial_rollups_sql do
"""
DO $validation$
DECLARE
expected_help bigint;
actual_help bigint;
expected_activities bigint;
actual_activities bigint;
BEGIN
SELECT count(*) INTO expected_help
FROM help_requests AS request
WHERE #{eligible_help_sql("request")};
SELECT coalesce(sum(point_count), 0) INTO actual_help
FROM help_discovery_rollups
WHERE level = 0;
IF expected_help <> actual_help THEN
RAISE EXCEPTION
'help discovery rollup validation failed: expected %, got %',
expected_help, actual_help;
END IF;
SELECT count(*) INTO expected_activities
FROM activities AS activity
WHERE #{eligible_activity_sql("activity")};
SELECT coalesce(sum(point_count), 0) INTO actual_activities
FROM activity_discovery_rollups
WHERE level = 0;
IF expected_activities <> actual_activities THEN
RAISE EXCEPTION
'activity discovery rollup validation failed: expected %, got %',
expected_activities, actual_activities;
END IF;
END
$validation$
"""
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