356 lines
13 KiB
Elixir
356 lines
13 KiB
Elixir
defmodule WhoNeedHelp.ProductionPlayPhysicalFixture do
|
|
import Ecto.Query
|
|
|
|
alias Oban.Job
|
|
alias WhoNeedHelp.Accounts.{Scope, User, UserToken}
|
|
alias WhoNeedHelp.Catalog.Category
|
|
alias WhoNeedHelp.Help.{Assignment, HelpRequest}
|
|
alias WhoNeedHelp.Messaging.Message
|
|
alias WhoNeedHelp.Notifications.Notification
|
|
alias WhoNeedHelp.Repo
|
|
alias WhoNeedHelp.Tracking
|
|
alias WhoNeedHelp.Tracking.{Position, TrackingSession}
|
|
alias WhoNeedHelp.Trust.{AuditEvent, RateLimitBucket, Report, Review}
|
|
|
|
@allowed_actions ~w(prepare verify-active verify-stopped cleanup)
|
|
|
|
def run(action, options) when is_binary(action) and is_map(options) do
|
|
unless action in @allowed_actions, do: raise("unsupported fixture action")
|
|
context = verified_context(action, options)
|
|
|
|
case action do
|
|
"prepare" -> prepare(context)
|
|
"verify-active" -> verify_active(context)
|
|
"verify-stopped" -> verify_stopped(context)
|
|
"cleanup" -> cleanup(context)
|
|
end
|
|
end
|
|
|
|
defp verified_context(action, options) do
|
|
run_id = required_option!(options, :run_id)
|
|
expected_database = required_option!(options, :expected_database)
|
|
helper_email = required_option!(options, :helper_email) |> String.downcase()
|
|
manifest_path = required_option!(options, :manifest_path)
|
|
|
|
unless Regex.match?(~r/^[a-z0-9-]+$/, run_id), do: raise("invalid run id")
|
|
|
|
unless String.starts_with?(manifest_path, "/tmp/wnh-play-physical-") do
|
|
raise "invalid manifest path"
|
|
end
|
|
|
|
%Postgrex.Result{rows: [[actual_database]]} =
|
|
Repo.query!("SELECT current_database()", [], log: false)
|
|
|
|
unless actual_database == expected_database, do: raise("database identity mismatch")
|
|
|
|
%{
|
|
action: action,
|
|
run_id: run_id,
|
|
database: actual_database,
|
|
helper_email: helper_email,
|
|
requester_email: "wnh-play-physical-#{run_id}-requester@example.invalid",
|
|
manifest_path: manifest_path
|
|
}
|
|
end
|
|
|
|
defp prepare(context) do
|
|
if File.exists?(context.manifest_path), do: raise("manifest already exists")
|
|
assert_requester_absent!(context)
|
|
|
|
helper = Repo.get_by(User, email: context.helper_email)
|
|
|
|
unless match?(%User{moderation_status: :active, confirmed_at: %DateTime{}}, helper) do
|
|
raise "helper is missing, unconfirmed, or inactive"
|
|
end
|
|
|
|
if Repo.exists?(from(s in TrackingSession, where: s.user_id == ^helper.id and s.active)) do
|
|
raise "helper already has an active tracking session"
|
|
end
|
|
|
|
now = DateTime.utc_now(:second)
|
|
category = Repo.get_by!(Category, slug: "medicine-pickup", active: true)
|
|
|
|
{:ok, fixture} =
|
|
Repo.transaction(fn ->
|
|
requester =
|
|
%User{}
|
|
|> User.registration_changeset(%{
|
|
"email" => context.requester_email,
|
|
"display_name" => "Play physical smoke requester",
|
|
"locale" => "en",
|
|
"terms_accepted" => true
|
|
})
|
|
|> Ecto.Changeset.put_change(:confirmed_at, now)
|
|
|> Repo.insert!()
|
|
|
|
request =
|
|
%HelpRequest{requester_id: requester.id}
|
|
|> HelpRequest.create_changeset(%{
|
|
"title" => "Play physical location smoke #{context.run_id}",
|
|
"description" => "Run-scoped synthetic request for foreground location verification.",
|
|
"structured_data" => %{"pickup_status" => "reserved"},
|
|
"location_label" => "Synthetic Play verification area",
|
|
"latitude" => "50.4501",
|
|
"longitude" => "30.5234",
|
|
"location_radius_meters" => 500,
|
|
"urgency" => "now",
|
|
"location_visibility" => "exact_for_active_match",
|
|
"expires_at" => DateTime.add(now, 2 * 60 * 60, :second),
|
|
"category_id" => category.id,
|
|
"safety_confirmed" => true
|
|
})
|
|
|> Ecto.Changeset.put_change(:status, :matched)
|
|
|> Repo.insert!()
|
|
|
|
assignment =
|
|
%Assignment{}
|
|
|> Assignment.changeset(%{
|
|
request_id: request.id,
|
|
helper_id: helper.id,
|
|
status: :accepted,
|
|
accepted_at: now,
|
|
handover_code_hash: :crypto.hash(:sha256, context.run_id)
|
|
})
|
|
|> Repo.insert!()
|
|
|
|
%{requester: requester, helper: helper, request: request, assignment: assignment}
|
|
end)
|
|
|
|
manifest = %{
|
|
"schema_version" => 1,
|
|
"run_id" => context.run_id,
|
|
"database" => context.database,
|
|
"requester" => %{"id" => fixture.requester.id, "email" => fixture.requester.email},
|
|
"helper" => %{"id" => fixture.helper.id, "email" => fixture.helper.email},
|
|
"request" => %{
|
|
"id" => fixture.request.id,
|
|
"path" => "/requests/#{fixture.request.id}"
|
|
},
|
|
"assignment" => %{"id" => fixture.assignment.id}
|
|
}
|
|
|
|
File.write!(context.manifest_path, Jason.encode_to_iodata!(manifest, pretty: true))
|
|
File.chmod!(context.manifest_path, 0o600)
|
|
IO.puts("fixture_prepared=true")
|
|
IO.puts("request_path=#{manifest["request"]["path"]}")
|
|
end
|
|
|
|
defp verify_active(context) do
|
|
fixture = context |> load_and_validate_manifest!() |> validate_fixture_links!()
|
|
sessions = fixture_sessions(fixture.assignment.id, fixture.helper.id)
|
|
positions = fixture_position_count(sessions)
|
|
|
|
unless match?([%TrackingSession{active: true, sample_count: n}] when n >= 1, sessions) and
|
|
positions == 1 do
|
|
raise "expected one active sampled session and one current raw position"
|
|
end
|
|
|
|
IO.puts("active_tracking_verified=true")
|
|
IO.puts("sample_count=#{hd(sessions).sample_count}")
|
|
IO.puts("retained_position_count=#{positions}")
|
|
end
|
|
|
|
defp verify_stopped(context) do
|
|
fixture = context |> load_and_validate_manifest!() |> validate_fixture_links!()
|
|
sessions = fixture_sessions(fixture.assignment.id, fixture.helper.id)
|
|
positions = fixture_position_count(sessions)
|
|
|
|
unless match?(
|
|
[%TrackingSession{active: false, ended_at: %DateTime{}, sample_count: n}]
|
|
when n >= 1,
|
|
sessions
|
|
) and positions == 0 do
|
|
raise "expected one stopped sampled session and zero retained raw positions"
|
|
end
|
|
|
|
IO.puts("stopped_tracking_verified=true")
|
|
IO.puts("sample_count=#{hd(sessions).sample_count}")
|
|
IO.puts("retained_position_count=#{positions}")
|
|
end
|
|
|
|
defp cleanup(context) do
|
|
fixture = context |> load_and_validate_manifest!() |> validate_fixture_links!()
|
|
|
|
if Repo.exists?(
|
|
from(s in TrackingSession,
|
|
where:
|
|
s.assignment_id == ^fixture.assignment.id and s.user_id == ^fixture.helper.id and
|
|
s.active
|
|
)
|
|
) do
|
|
case Tracking.stop_session(Scope.for_user(fixture.helper), fixture.assignment) do
|
|
{:ok, _} -> :ok
|
|
{:error, reason} -> raise "could not stop fixture tracking: #{inspect(reason)}"
|
|
end
|
|
end
|
|
|
|
sessions = fixture_sessions(fixture.assignment.id, fixture.helper.id)
|
|
session_ids = Enum.map(sessions, & &1.id)
|
|
message_ids = exact_ids(Message, :assignment_id, fixture.assignment.id)
|
|
|
|
notification_ids =
|
|
Notification
|
|
|> where(
|
|
[n],
|
|
fragment("?->>'assignment_id' = ?", n.data, ^fixture.assignment.id) or
|
|
fragment("?->>'request_id' = ?", n.data, ^fixture.request.id)
|
|
)
|
|
|> select([n], n.id)
|
|
|> Repo.all()
|
|
|
|
job_ids = fixture_job_ids(fixture, notification_ids)
|
|
target_ids = [fixture.requester.id, fixture.request.id, fixture.assignment.id] ++ message_ids
|
|
requester_scope_hash = :crypto.hash(:sha256, fixture.requester.id)
|
|
|
|
{:ok, deleted} =
|
|
Repo.transaction(fn ->
|
|
%{
|
|
jobs: delete_ids(Job, job_ids),
|
|
notifications: delete_ids(Notification, notification_ids),
|
|
audit_events: AuditEvent |> where([e], e.target_id in ^target_ids) |> delete_count(),
|
|
reports:
|
|
Report
|
|
|> where(
|
|
[r],
|
|
r.request_id == ^fixture.request.id or r.assignment_id == ^fixture.assignment.id or
|
|
r.message_id in ^message_ids
|
|
)
|
|
|> delete_count(),
|
|
reviews:
|
|
Review
|
|
|> where([r], r.assignment_id == ^fixture.assignment.id)
|
|
|> delete_count(),
|
|
positions:
|
|
Position
|
|
|> where([p], p.tracking_session_id in ^session_ids)
|
|
|> delete_count(),
|
|
sessions: delete_ids(TrackingSession, session_ids),
|
|
messages: delete_ids(Message, message_ids),
|
|
assignment: delete_ids(Assignment, [fixture.assignment.id]),
|
|
request: delete_ids(HelpRequest, [fixture.request.id]),
|
|
requester_tokens:
|
|
UserToken
|
|
|> where([t], t.user_id == ^fixture.requester.id)
|
|
|> delete_count(),
|
|
requester_rate_limits:
|
|
RateLimitBucket
|
|
|> where([b], b.scope_hash == ^requester_scope_hash)
|
|
|> delete_count(),
|
|
requester: delete_ids(User, [fixture.requester.id])
|
|
}
|
|
end)
|
|
|
|
unless deleted.assignment == 1 and deleted.request == 1 and deleted.requester == 1 do
|
|
raise "cleanup did not remove the exact fixture"
|
|
end
|
|
|
|
assert_requester_absent!(context)
|
|
File.rm!(context.manifest_path)
|
|
IO.puts("fixture_cleanup_verified=true")
|
|
IO.puts("deleted=#{inspect(deleted, limit: :infinity)}")
|
|
end
|
|
|
|
defp load_and_validate_manifest!(context) do
|
|
manifest = context.manifest_path |> File.read!() |> Jason.decode!()
|
|
|
|
valid? =
|
|
manifest["schema_version"] == 1 and manifest["run_id"] == context.run_id and
|
|
manifest["database"] == context.database and
|
|
manifest["requester"]["email"] == context.requester_email and
|
|
manifest["helper"]["email"] == context.helper_email and
|
|
uuid?(manifest["requester"]["id"]) and uuid?(manifest["helper"]["id"]) and
|
|
uuid?(manifest["request"]["id"]) and uuid?(manifest["assignment"]["id"]) and
|
|
manifest["request"]["path"] == "/requests/#{manifest["request"]["id"]}"
|
|
|
|
unless valid?, do: raise("manifest does not match this exact run")
|
|
manifest
|
|
end
|
|
|
|
defp validate_fixture_links!(manifest) do
|
|
requester = Repo.get(User, manifest["requester"]["id"])
|
|
helper = Repo.get(User, manifest["helper"]["id"])
|
|
request = Repo.get(HelpRequest, manifest["request"]["id"])
|
|
assignment = Repo.get(Assignment, manifest["assignment"]["id"])
|
|
|
|
valid? =
|
|
match?(%User{}, requester) and match?(%User{}, helper) and
|
|
match?(%HelpRequest{}, request) and match?(%Assignment{}, assignment) and
|
|
requester.email == manifest["requester"]["email"] and
|
|
helper.email == manifest["helper"]["email"] and request.requester_id == requester.id and
|
|
assignment.request_id == request.id and assignment.helper_id == helper.id and
|
|
assignment.status in [:accepted, :in_progress] and assignment.active
|
|
|
|
unless valid?, do: raise("fixture links no longer match the manifest")
|
|
|
|
if Repo.exists?(
|
|
from(r in HelpRequest,
|
|
where: r.requester_id == ^requester.id and r.id != ^request.id
|
|
)
|
|
) do
|
|
raise "synthetic requester owns unexpected requests"
|
|
end
|
|
|
|
%{requester: requester, helper: helper, request: request, assignment: assignment}
|
|
end
|
|
|
|
defp fixture_sessions(assignment_id, helper_id) do
|
|
TrackingSession
|
|
|> where([s], s.assignment_id == ^assignment_id and s.user_id == ^helper_id)
|
|
|> order_by([s], asc: s.inserted_at)
|
|
|> Repo.all()
|
|
end
|
|
|
|
defp fixture_position_count([]), do: 0
|
|
|
|
defp fixture_position_count(sessions) do
|
|
ids = Enum.map(sessions, & &1.id)
|
|
Position |> where([p], p.tracking_session_id in ^ids) |> Repo.aggregate(:count)
|
|
end
|
|
|
|
defp fixture_job_ids(fixture, notification_ids) do
|
|
notification_ids = MapSet.new(notification_ids)
|
|
|
|
Job
|
|
|> where(
|
|
[job],
|
|
job.worker in [
|
|
"WhoNeedHelp.Push.DeliveryWorker",
|
|
"WhoNeedHelp.Push.NotificationDispatchWorker",
|
|
"WhoNeedHelp.Push.DeviceDeliveryWorker",
|
|
"WhoNeedHelp.Push.NotificationEmailWorker"
|
|
]
|
|
)
|
|
|> select([job], {job.id, job.args})
|
|
|> Repo.all()
|
|
|> Enum.filter(fn {_id, args} ->
|
|
args["assignment_id"] == fixture.assignment.id or
|
|
args["request_id"] == fixture.request.id or
|
|
MapSet.member?(notification_ids, args["notification_id"])
|
|
end)
|
|
|> Enum.map(&elem(&1, 0))
|
|
end
|
|
|
|
defp exact_ids(schema, name, value) do
|
|
schema |> where([row], field(row, ^name) == ^value) |> select([row], row.id) |> Repo.all()
|
|
end
|
|
|
|
defp delete_ids(_schema, []), do: 0
|
|
defp delete_ids(schema, ids), do: schema |> where([row], row.id in ^ids) |> delete_count()
|
|
defp delete_count(query), do: query |> Repo.delete_all() |> elem(0)
|
|
|
|
defp assert_requester_absent!(context) do
|
|
if Repo.exists?(from(u in User, where: u.email == ^context.requester_email)) do
|
|
raise "run-scoped requester already exists"
|
|
end
|
|
end
|
|
|
|
defp required_option!(options, name) do
|
|
case Map.get(options, name) do
|
|
value when is_binary(value) and value != "" -> value
|
|
_ -> raise "#{name} is required"
|
|
end
|
|
end
|
|
|
|
defp uuid?(value) when is_binary(value), do: match?({:ok, _}, Ecto.UUID.cast(value))
|
|
defp uuid?(_), do: false
|
|
end |