who_need_help/lib/mix/tasks/wnh.scale_fixtures.ex

695 lines
23 KiB
Elixir

defmodule Mix.Tasks.Wnh.ScaleFixtures do
use Mix.Task
alias WhoNeedHelp.Accounts.{Scope, User}
alias WhoNeedHelp.{Activities, Catalog, Help, Repo}
alias WhoNeedHelp.Help.DiscoveryViewport
@shortdoc "Prepares, appends, or verifies a persistent local discovery-scale data set"
@isolated_confirmation "persistent-local-scale"
@additive_confirmation "persistent-dev-scale"
@default_rows 1_000_000
@synthetic_users 1_000
@viewer_email "scale-viewer@example.invalid"
@impl Mix.Task
def run(arguments) do
Mix.Task.run("app.start")
{options, rest, invalid} =
OptionParser.parse(arguments,
strict: [requests: :integer, activities: :integer, samples: :integer, output: :string]
)
if invalid != [], do: Mix.raise("invalid options: #{inspect(invalid)}")
action =
case rest do
[value] when value in ["prepare", "append", "verify"] -> value
_other -> Mix.raise(usage())
end
request_count = positive_count!(options, :requests)
activity_count = positive_count!(options, :activities)
samples = positive_count!(options, :samples, 3)
output = Keyword.get(options, :output) || Mix.raise("--output PATH is required")
context = verified_context!(output, action)
case action do
"prepare" -> prepare!(context, request_count, activity_count)
"append" -> append!(context, request_count, activity_count)
"verify" -> verify!(context, request_count, activity_count, samples)
end
end
defp usage do
"usage: mix wnh.scale_fixtures prepare|append|verify " <>
"[--requests COUNT] [--activities COUNT] [--samples COUNT] --output /output/FILE.json"
end
defp positive_count!(options, key, default \\ @default_rows) do
value = Keyword.get(options, key, default)
if is_integer(value) and value > 0, do: value, else: Mix.raise("--#{key} must be positive")
end
defp verified_context!(output, action) do
expected_database = required_env!("WNH_SCALE_EXPECTED_DATABASE")
output = Path.expand(output)
expected_confirmation =
if action == "append", do: @additive_confirmation, else: @isolated_confirmation
unless System.get_env("WNH_SCALE_FIXTURE_CONFIRM") == expected_confirmation do
Mix.raise("WNH_SCALE_FIXTURE_CONFIRM must equal #{expected_confirmation}")
end
unless String.starts_with?(output, "/output/") do
Mix.raise("--output must resolve below /output")
end
%Postgrex.Result{rows: [[actual_database]]} =
Repo.query!("SELECT current_database()", [], log: false)
unless actual_database == expected_database do
Mix.raise(
"refusing scale mutation: expected database #{inspect(expected_database)}, " <>
"observed #{inspect(actual_database)}"
)
end
%{database: actual_database, output: output}
end
defp append!(context, request_count, activity_count) do
before_totals = discovery_counts()
before_synthetic = synthetic_discovery_counts()
if before_synthetic.help_requests > request_count or
before_synthetic.activities > activity_count do
Mix.raise(
"the development database already contains more synthetic scale rows than requested: " <>
"#{before_synthetic.help_requests} requests and #{before_synthetic.activities} activities"
)
end
categories = seed_categories!()
ensure_viewer!()
Repo.transaction(
fn ->
seed_users!()
seed_requests!(request_count, categories.help)
seed_activities!(activity_count, categories.activity)
end,
timeout: :infinity
)
Repo.query!(
"ANALYZE users, categories, help_requests, activities",
[],
timeout: :infinity,
log: false
)
assert_synthetic_counts!(request_count, activity_count)
write_summary!(context, "appended", request_count, activity_count, %{
expected_kind: "synthetic",
totals_before: before_totals,
synthetic_before: before_synthetic,
synthetic_after: synthetic_discovery_counts()
})
totals = discovery_counts()
Mix.shell().info(
"persistent synthetic scale data is ready alongside existing data in #{context.database}: " <>
"#{totals.help_requests} total requests and #{totals.activities} total activities"
)
end
defp prepare!(context, request_count, activity_count) do
counts = discovery_counts()
cond do
counts.help_requests == request_count and counts.activities == activity_count ->
ensure_viewer!()
write_summary!(context, "prepared", request_count, activity_count, %{})
counts.help_requests != 0 or counts.activities != 0 ->
Mix.raise(
"the isolated scale database already contains #{counts.help_requests} help requests " <>
"and #{counts.activities} activities; refusing to mix data sets"
)
true ->
categories = seed_categories!()
ensure_viewer!()
Repo.transaction(
fn ->
seed_users!()
seed_requests!(request_count, categories.help)
seed_activities!(activity_count, categories.activity)
end,
timeout: :infinity
)
Repo.query!(
"ANALYZE users, categories, help_requests, activities",
[],
timeout: :infinity,
log: false
)
assert_counts!(request_count, activity_count)
write_summary!(context, "prepared", request_count, activity_count, %{})
end
Mix.shell().info(
"persistent scale data is ready in #{context.database}: " <>
"#{request_count} requests and #{activity_count} activities"
)
end
defp verify!(context, request_count, activity_count, samples) do
assert_counts!(request_count, activity_count)
viewer = Repo.get_by!(User, email: @viewer_email)
scope = Scope.for_user(viewer)
viewports = %{
kyiv: viewport!(30.30, 50.25, 30.80, 50.65, 10),
europe: viewport!(-12.0, 34.0, 35.0, 61.0, 4),
world: viewport!(-179.99, -80.0, 179.99, 80.0, 1)
}
measurements =
Map.new(viewports, fn {name, viewport} ->
request_page =
timed_samples(
fn -> Help.paginate_open_requests(scope, %{}, viewport: viewport) end,
samples
)
request_map =
timed_samples(fn -> Help.map_discovery_items(scope, %{}, viewport) end, samples)
activity_page =
timed_samples(
fn -> Activities.paginate_open_activities(scope, %{}, viewport: viewport) end,
samples
)
activity_map =
timed_samples(fn -> Activities.map_discovery_items(scope, %{}, viewport) end, samples)
{name,
%{
request_page: page_measurement(request_page),
request_map: map_measurement(request_map),
activity_page: page_measurement(activity_page),
activity_map: map_measurement(activity_map)
}}
end)
write_summary!(context, "verified", request_count, activity_count, %{
measurement_samples: samples,
viewport_measurements: measurements
})
Mix.shell().info("scale verification written to #{context.output}")
end
defp seed_categories! do
Catalog.seed_defaults()
categories = Catalog.list_all_categories()
%{
help: category_ids!(categories, :help),
activity: category_ids!(categories, :activity)
}
end
defp category_ids!(categories, mode) do
leaves =
Enum.filter(categories, fn category ->
category.mode == mode and not is_nil(category.parent_id)
end)
ids =
case leaves do
[] -> categories |> Enum.filter(&(&1.mode == mode)) |> Enum.map(& &1.id)
values -> Enum.map(values, & &1.id)
end
if ids == [], do: Mix.raise("no #{mode} categories were seeded"), else: ids
end
defp ensure_viewer! do
password = required_env!("SCALE_VIEWER_PASSWORD")
unless byte_size(password) in 12..72 do
Mix.raise("SCALE_VIEWER_PASSWORD must contain between 12 and 72 bytes")
end
case Repo.get_by(User, email: @viewer_email) do
nil ->
now = DateTime.utc_now(:second)
%User{
email: @viewer_email,
display_name: "Scale Viewer",
hashed_password: Bcrypt.hash_pwd_salt(password),
confirmed_at: now,
accepted_terms_at: now,
locale: "en"
}
|> Repo.insert!()
%User{} ->
:ok
end
end
defp seed_users! do
Repo.query!(
"""
INSERT INTO users (
id, email, display_name, confirmed_at, accepted_terms_at, locale,
inserted_at, updated_at
)
SELECT
md5('wnh-scale-user-' || value)::uuid,
'scale-user-' || value || '@example.invalid',
'Scale user ' || value,
date_trunc('second', now()),
date_trunc('second', now()),
CASE value % 3 WHEN 0 THEN 'uk' WHEN 1 THEN 'en' ELSE 'ru' END,
date_trunc('second', now()) - value * interval '1 second',
date_trunc('second', now()) - value * interval '1 second'
FROM generate_series(1, $1) AS value
ON CONFLICT DO NOTHING
""",
[@synthetic_users],
timeout: :infinity,
log: false
)
end
defp seed_requests!(rows, category_ids) do
Repo.query!(
"""
INSERT INTO help_requests (
id, title, description, pickup_instructions, structured_data,
location_label, location, status, urgency, location_visibility,
location_radius_meters, expires_at, requester_id, category_id,
inserted_at, updated_at
)
SELECT
md5('wnh-scale-request-v1-' || value)::uuid,
'Scale help request ' || value,
'Persistent synthetic local scale request used to test viewport discovery and UI.',
'Synthetic data only; no real-world action is requested.',
'{"scale_fixture":true,"scale_fixture_version":3}'::jsonb,
city.name || ' synthetic zone ' ||
((variation.category_key + variation.time_key) % 100 + 1),
CASE WHEN variation.visibility_key % 97 = 0 THEN NULL ELSE ST_SetSRID(
ST_MakePoint(
city.longitude +
((variation.longitude_key % 4001) - 2000) / 1000.0,
city.latitude +
((variation.latitude_key % 4003) - 2001) / 1000.0
),
4326
) END,
CASE WHEN variation.status_key % 10 < 8 THEN 'open'
WHEN variation.status_key % 10 = 8 THEN 'completed'
ELSE 'cancelled' END,
CASE (variation.status_key / 10) % 3
WHEN 0 THEN 'now'
WHEN 1 THEN 'today'
ELSE 'scheduled'
END,
CASE WHEN variation.visibility_key % 97 = 0 THEN 'hidden'
WHEN (variation.visibility_key / 97) % 4 = 0 THEN 'exact_public'
WHEN (variation.visibility_key / 97) % 4 = 1
THEN 'exact_for_active_match'
ELSE 'approximate_public' END,
CASE WHEN variation.visibility_key % 97 = 0
OR (variation.visibility_key / 97) % 4 = 0 THEN NULL
WHEN (variation.visibility_key / 388) % 3 = 0 THEN 500
WHEN (variation.visibility_key / 388) % 3 = 1 THEN 1000
ELSE 2000 END,
date_trunc('second', now()) +
((variation.time_key % 30) + 1) * interval '1 day',
md5(
'wnh-scale-user-' ||
((variation.longitude_key * 37 + variation.latitude_key) % $2 + 1)
)::uuid,
(($3::text[])[
(variation.category_key % cardinality($3::text[])) + 1
])::uuid,
date_trunc('second', now()) - value * interval '1 second',
date_trunc('second', now()) - value * interval '1 second'
FROM generate_series(1, $1) AS value
CROSS JOIN LATERAL (
SELECT
((value - 1) % 20)::integer AS city_index,
md5('wnh-scale-request-v3-' || value) AS fixture_hash
) AS fixture
CROSS JOIN LATERAL (
SELECT
('x' || substr(fixture.fixture_hash, 1, 8))::bit(32)::bigint
AS longitude_key,
('x' || substr(fixture.fixture_hash, 9, 8))::bit(32)::bigint
AS latitude_key,
('x' || substr(fixture.fixture_hash, 17, 4))::bit(16)::bigint
AS status_key,
('x' || substr(fixture.fixture_hash, 21, 4))::bit(16)::bigint
AS visibility_key,
('x' || substr(fixture.fixture_hash, 25, 4))::bit(16)::bigint
AS time_key,
('x' || substr(fixture.fixture_hash, 29, 4))::bit(16)::bigint
AS category_key
) AS variation
CROSS JOIN LATERAL (
SELECT *
FROM (
VALUES
(0, 'Kyiv', 30.5234::double precision, 50.4501::double precision),
(1, 'Frankfurt', 8.6821, 50.1109),
(2, 'Berlin', 13.4050, 52.5200),
(3, 'Warsaw', 21.0122, 52.2297),
(4, 'Prague', 14.4378, 50.0755),
(5, 'London', -0.1276, 51.5072),
(6, 'Paris', 2.3522, 48.8566),
(7, 'Madrid', -3.7038, 40.4168),
(8, 'Rome', 12.4964, 41.9028),
(9, 'New York', -74.0060, 40.7128),
(10, 'Toronto', -79.3832, 43.6532),
(11, 'Mexico City', -99.1332, 19.4326),
(12, 'São Paulo', -46.6333, -23.5505),
(13, 'Cape Town', 18.4241, -33.9249),
(14, 'Nairobi', 36.8219, -1.2921),
(15, 'Delhi', 77.1025, 28.7041),
(16, 'Bangkok', 100.5018, 13.7563),
(17, 'Tokyo', 139.6917, 35.6895),
(18, 'Sydney', 151.2093, -33.8688),
(19, 'Auckland', 174.7633, -36.8485)
) AS cities(index, name, longitude, latitude)
WHERE cities.index = fixture.city_index
) AS city
ON CONFLICT (id) DO UPDATE SET
structured_data = EXCLUDED.structured_data,
location_label = EXCLUDED.location_label,
location = EXCLUDED.location,
status = EXCLUDED.status,
urgency = EXCLUDED.urgency,
location_visibility = EXCLUDED.location_visibility,
location_radius_meters = EXCLUDED.location_radius_meters,
expires_at = EXCLUDED.expires_at,
requester_id = EXCLUDED.requester_id,
category_id = EXCLUDED.category_id,
updated_at = EXCLUDED.updated_at
WHERE help_requests.structured_data @> '{"scale_fixture":true}'::jsonb
AND help_requests.structured_data ->> 'scale_fixture_version' IS DISTINCT FROM '3'
""",
[rows, @synthetic_users, category_ids],
timeout: :infinity,
log: false
)
end
defp seed_activities!(rows, category_ids) do
Repo.query!(
"""
INSERT INTO activities (
id, title, description, structured_data, location_label, location,
location_visibility, status, starts_at, join_deadline, capacity,
creator_id, category_id, inserted_at, updated_at
)
SELECT
md5('wnh-scale-activity-v1-' || value)::uuid,
'Scale community activity ' || value,
'Persistent synthetic local scale activity used to test viewport discovery and UI.',
'{"scale_fixture":true,"scale_fixture_version":3}'::jsonb,
city.name || ' synthetic zone ' ||
((variation.category_key + variation.time_key) % 100 + 1),
ST_SetSRID(
ST_MakePoint(
city.longitude +
((variation.longitude_key % 4001) - 2000) / 1000.0,
city.latitude +
((variation.latitude_key % 4003) - 2001) / 1000.0
),
4326
),
CASE WHEN variation.visibility_key % 97 = 0
THEN 'hidden'
ELSE 'approximate_public'
END,
CASE WHEN variation.status_key % 10 < 8 THEN 'open'
WHEN variation.status_key % 10 = 8 THEN 'completed'
ELSE 'cancelled' END,
date_trunc('second', now()) +
((variation.time_key % 60) + 2) * interval '1 day',
date_trunc('second', now()) +
((variation.time_key % 60) + 1) * interval '1 day',
2 + ((variation.status_key + variation.time_key) % 19),
md5(
'wnh-scale-user-' ||
((variation.longitude_key * 37 + variation.latitude_key) % $2 + 1)
)::uuid,
(($3::text[])[
(variation.category_key % cardinality($3::text[])) + 1
])::uuid,
date_trunc('second', now()) - value * interval '1 second',
date_trunc('second', now()) - value * interval '1 second'
FROM generate_series(1, $1) AS value
CROSS JOIN LATERAL (
SELECT
((value - 1) % 20)::integer AS city_index,
md5('wnh-scale-activity-v3-' || value) AS fixture_hash
) AS fixture
CROSS JOIN LATERAL (
SELECT
('x' || substr(fixture.fixture_hash, 1, 8))::bit(32)::bigint
AS longitude_key,
('x' || substr(fixture.fixture_hash, 9, 8))::bit(32)::bigint
AS latitude_key,
('x' || substr(fixture.fixture_hash, 17, 4))::bit(16)::bigint
AS status_key,
('x' || substr(fixture.fixture_hash, 21, 4))::bit(16)::bigint
AS visibility_key,
('x' || substr(fixture.fixture_hash, 25, 4))::bit(16)::bigint
AS time_key,
('x' || substr(fixture.fixture_hash, 29, 4))::bit(16)::bigint
AS category_key
) AS variation
CROSS JOIN LATERAL (
SELECT *
FROM (
VALUES
(0, 'Kyiv', 30.5234::double precision, 50.4501::double precision),
(1, 'Frankfurt', 8.6821, 50.1109),
(2, 'Berlin', 13.4050, 52.5200),
(3, 'Warsaw', 21.0122, 52.2297),
(4, 'Prague', 14.4378, 50.0755),
(5, 'London', -0.1276, 51.5072),
(6, 'Paris', 2.3522, 48.8566),
(7, 'Madrid', -3.7038, 40.4168),
(8, 'Rome', 12.4964, 41.9028),
(9, 'New York', -74.0060, 40.7128),
(10, 'Toronto', -79.3832, 43.6532),
(11, 'Mexico City', -99.1332, 19.4326),
(12, 'São Paulo', -46.6333, -23.5505),
(13, 'Cape Town', 18.4241, -33.9249),
(14, 'Nairobi', 36.8219, -1.2921),
(15, 'Delhi', 77.1025, 28.7041),
(16, 'Bangkok', 100.5018, 13.7563),
(17, 'Tokyo', 139.6917, 35.6895),
(18, 'Sydney', 151.2093, -33.8688),
(19, 'Auckland', 174.7633, -36.8485)
) AS cities(index, name, longitude, latitude)
WHERE cities.index = fixture.city_index
) AS city
ON CONFLICT (id) DO UPDATE SET
structured_data = EXCLUDED.structured_data,
location_label = EXCLUDED.location_label,
location = EXCLUDED.location,
location_visibility = EXCLUDED.location_visibility,
status = EXCLUDED.status,
starts_at = EXCLUDED.starts_at,
join_deadline = EXCLUDED.join_deadline,
capacity = EXCLUDED.capacity,
creator_id = EXCLUDED.creator_id,
category_id = EXCLUDED.category_id,
updated_at = EXCLUDED.updated_at
WHERE activities.structured_data @> '{"scale_fixture":true}'::jsonb
AND activities.structured_data ->> 'scale_fixture_version' IS DISTINCT FROM '3'
""",
[rows, @synthetic_users, category_ids],
timeout: :infinity,
log: false
)
end
defp viewport!(west, south, east, north, zoom) do
{:ok, viewport} =
DiscoveryViewport.cast(%{
"west" => west,
"south" => south,
"east" => east,
"north" => north,
"zoom" => zoom,
"width" => 1_280,
"height" => 720
})
viewport
end
defp timed_samples(callback, samples) do
_warm_result = callback.()
measurements =
Enum.map(1..samples, fn _sample ->
{microseconds, result} = :timer.tc(callback)
%{elapsed_ms: Float.round(microseconds / 1_000, 3), result: result}
end)
elapsed = measurements |> Enum.map(& &1.elapsed_ms) |> Enum.sort()
%{
elapsed_ms: median(elapsed),
min_elapsed_ms: hd(elapsed),
max_elapsed_ms: List.last(elapsed),
samples_ms: elapsed,
result: measurements |> List.last() |> Map.fetch!(:result)
}
end
defp median(values) do
middle = div(length(values), 2)
if rem(length(values), 2) == 1 do
Enum.at(values, middle)
else
Float.round((Enum.at(values, middle - 1) + Enum.at(values, middle)) / 2, 3)
end
end
defp page_measurement(%{result: page} = measurement) do
measurement
|> timing_measurement()
|> Map.merge(%{
returned_rows: length(page.entries),
has_next_page: not is_nil(page.next_cursor)
})
end
defp map_measurement(%{result: items} = measurement) do
measurement
|> timing_measurement()
|> Map.merge(%{
returned_rows: length(items),
represented_records: Enum.reduce(items, 0, &(Map.get(&1, :count, 1) + &2)),
json_bytes: items |> Jason.encode!() |> byte_size()
})
end
defp timing_measurement(measurement) do
Map.take(measurement, [:elapsed_ms, :min_elapsed_ms, :max_elapsed_ms, :samples_ms])
end
defp assert_counts!(request_count, activity_count) do
actual = discovery_counts()
unless actual.help_requests == request_count and actual.activities == activity_count do
Mix.raise(
"expected #{request_count} requests and #{activity_count} activities, " <>
"found #{actual.help_requests} and #{actual.activities}"
)
end
end
defp assert_synthetic_counts!(request_count, activity_count) do
actual = synthetic_discovery_counts()
unless actual.help_requests == request_count and actual.activities == activity_count do
Mix.raise(
"expected #{request_count} synthetic requests and #{activity_count} synthetic activities, " <>
"found #{actual.help_requests} and #{actual.activities}"
)
end
end
defp discovery_counts do
%{
help_requests: scalar!("SELECT count(*) FROM help_requests"),
activities: scalar!("SELECT count(*) FROM activities")
}
end
defp synthetic_discovery_counts do
%{
help_requests:
scalar!(
"SELECT count(*) FROM help_requests " <>
"WHERE structured_data @> '{\"scale_fixture\": true}'::jsonb"
),
activities:
scalar!(
"SELECT count(*) FROM activities " <>
"WHERE structured_data @> '{\"scale_fixture\": true}'::jsonb"
)
}
end
defp write_summary!(context, status, request_count, activity_count, extra) do
File.mkdir_p!(Path.dirname(context.output))
summary =
Map.merge(
%{
schema_version: 1,
status: status,
database: context.database,
expected: %{help_requests: request_count, activities: activity_count},
counts: discovery_counts(),
database_bytes: scalar!("SELECT pg_database_size(current_database())"),
relation_bytes: relation_sizes(),
captured_at: DateTime.utc_now() |> DateTime.truncate(:second) |> DateTime.to_iso8601()
},
extra
)
File.write!(context.output, Jason.encode_to_iodata!(summary, pretty: true))
end
defp relation_sizes do
Repo.query!(
"""
SELECT relation, pg_total_relation_size(relation::regclass)
FROM unnest(ARRAY['users', 'categories', 'help_requests', 'activities']) AS relation
ORDER BY relation
""",
[],
log: false
).rows
|> Map.new(fn [name, bytes] -> {name, bytes} end)
end
defp scalar!(sql) do
Repo.query!(sql, [], timeout: :infinity, log: false).rows |> hd() |> hd()
end
defp required_env!(name) do
case System.get_env(name) do
value when is_binary(value) and value != "" -> value
_missing -> Mix.raise("#{name} is required")
end
end
end