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) requester_password = temporary_password() {: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!() |> User.password_changeset(%{"password" => requester_password}) |> Repo.update!() 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" => 2, "run_id" => context.run_id, "database" => context.database, "requester" => %{ "id" => fixture.requester.id, "email" => fixture.requester.email, "password" => requester_password }, "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"]}") IO.puts("requester_email=#{fixture.requester.email}") IO.puts("requester_password=#{requester_password}") IO.puts("credentials_are_temporary=true") 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"] == 2 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 temporary_password?(manifest["requester"]["password"]) 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 User.valid_password?(requester, manifest["requester"]["password"]) 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 temporary_password do 32 |> :crypto.strong_rand_bytes() |> Base.url_encode64(padding: false) end defp temporary_password?(password) when is_binary(password), do: byte_size(password) >= 32 and byte_size(password) <= 72 defp temporary_password?(_password), do: false defp uuid?(value) when is_binary(value), do: match?({:ok, _}, Ecto.UUID.cast(value)) defp uuid?(_), do: false end