who_need_help/lib/who_need_help/help.ex

1330 lines
42 KiB
Elixir

defmodule WhoNeedHelp.Help do
@moduledoc "Authorization-aware request matching and completion workflow."
import Ecto.Query
alias Ecto.Multi
alias WhoNeedHelp.Accounts
alias WhoNeedHelp.Accounts.Scope
alias WhoNeedHelp.Catalog
alias WhoNeedHelp.Catalog.Category
alias WhoNeedHelp.Discovery.Rollup
alias WhoNeedHelp.Help.{Assignment, DiscoveryCluster, DiscoveryViewport, HelpRequest}
alias WhoNeedHelp.Pagination
alias WhoNeedHelp.ProductAnalytics
alias WhoNeedHelp.Push
alias WhoNeedHelp.Push.NearbyMatchWorker
alias WhoNeedHelp.Repo
alias WhoNeedHelp.Trust
alias WhoNeedHelp.Trust.Block
@topic "help:requests"
@hidden_discovery_topic "help:discovery:hidden"
def subscribe, do: Phoenix.PubSub.subscribe(WhoNeedHelp.PubSub, @topic)
def subscribe_user(user_id),
do: Phoenix.PubSub.subscribe(WhoNeedHelp.PubSub, "help:user:#{user_id}")
def subscribe_discovery(%DiscoveryViewport{} = viewport, include_hidden \\ false) do
topics = discovery_topics(viewport, include_hidden)
Enum.each(topics, &Phoenix.PubSub.subscribe(WhoNeedHelp.PubSub, &1))
topics
end
def unsubscribe_discovery(topics) when is_list(topics) do
Enum.each(topics, &Phoenix.PubSub.unsubscribe(WhoNeedHelp.PubSub, &1))
:ok
end
def subscribe_request(id),
do: Phoenix.PubSub.subscribe(WhoNeedHelp.PubSub, "help:request:#{id}")
def notify_request_updated(request_id) do
request = get_request!(request_id)
broadcast({:request_updated, request})
:ok
end
def list_open_requests(%Scope{user: user}, filters \\ %{}) do
paginate_open_requests(%Scope{user: user}, filters).entries
end
def paginate_open_requests(%Scope{user: user}, filters \\ %{}, options \\ []) do
limit = Pagination.limit(options)
cursor = Pagination.cursor(options)
viewport = Keyword.get(options, :viewport)
user
|> open_requests_query(filters, viewport)
|> after_open_request(cursor)
|> order_by([request], asc: request.expires_at, asc: request.id)
|> limit(^(limit + 1))
|> Repo.all()
|> preload_request_relations()
|> Pagination.page(limit, &{&1.expires_at, &1.id})
end
def count_open_requests(%Scope{user: user}, filters \\ %{}, options \\ []) do
viewport = Keyword.get(options, :viewport)
user
|> open_requests_query(filters, viewport)
|> Repo.aggregate(:count, :id)
end
def map_discovery_items(%Scope{user: user}, filters, %DiscoveryViewport{} = viewport) do
level = DiscoveryCluster.level_for_zoom(viewport.zoom)
if DiscoveryCluster.rollup_level?(level) and blank_area_filter?(filters) do
Rollup.help_items(user, filters, viewport, level)
else
raw_map_discovery_items(user, filters, viewport, level)
end
end
def map_cluster_expansion(%Scope{user: user}, filters, cluster_id) do
Rollup.help_expansion(user, filters, cluster_id)
end
defp raw_map_discovery_items(user, filters, viewport, level) do
cell_count = Integer.pow(2, level) * 1.0
points =
user
|> open_requests_query(filters, nil)
|> where(
[request],
not is_nil(request.location) and request.location_visibility != :hidden
)
|> filter_request_map_cluster_viewport(viewport, level)
|> select([request], %{
id: request.id,
urgency: request.urgency,
public_longitude: request.public_longitude,
public_latitude: request.public_latitude
})
points
|> subquery()
|> group_by([_point], [selected_as(:cell_x), selected_as(:cell_y)])
|> select([point], %{
count: count(point.id),
urgent_count: fragment("count(*) FILTER (WHERE ? = 'now')::bigint", point.urgency),
request_id: fragment("min(?::text)", point.id),
longitude: min(point.public_longitude),
latitude: min(point.public_latitude),
west: min(point.public_longitude),
south: min(point.public_latitude),
east: max(point.public_longitude),
north: max(point.public_latitude),
cell_x:
selected_as(
fragment(
"floor(((? + 180.0) / 360.0) * ?)::bigint",
point.public_longitude,
constant(^cell_count)
),
:cell_x
),
cell_y:
selected_as(
fragment(
"floor(((ln(tan(pi() / 4.0 + radians(least(85.05112878, greatest(-85.05112878, ?))) / 2.0)) / pi() + 1.0) / 2.0) * ?)::bigint",
point.public_latitude,
constant(^cell_count)
),
:cell_y
)
})
|> order_by([_point], asc: selected_as(:cell_x), asc: selected_as(:cell_y))
|> Repo.all()
|> hydrate_request_map_singletons()
|> Enum.map(&map_discovery_item(&1, level))
end
def count_map_discovery_items(%Scope{user: user}, filters) do
user
|> open_requests_query(filters, nil)
|> where(
[request],
not is_nil(request.location) and request.location_visibility != :hidden
)
|> Repo.aggregate(:count, :id)
end
# Lists and map markers use the same displayed public center. A privacy
# radius describes uncertainty around that center; it must not make a card
# appear in a viewport where its public marker was not loaded.
defp filter_request_map_cluster_viewport(
query,
%DiscoveryViewport{} = viewport,
level
) do
envelopes = DiscoveryCluster.viewport_envelopes(viewport, level)
where(query, ^request_map_envelope_condition(envelopes))
end
defp request_map_viewport_condition(%DiscoveryViewport{} = viewport) do
request_map_envelope_condition(DiscoveryViewport.envelopes(viewport))
end
defp request_map_envelope_condition(envelopes) when is_list(envelopes) do
spatial_condition =
Enum.reduce(envelopes, dynamic(false), fn
{west, south, east, north}, condition ->
dynamic(
[request],
^condition or
fragment(
"""
(
? && ST_MakeEnvelope(?, ?, ?, ?, 4326)
AND
? >= ? AND ? <= ? AND ? >= ? AND ? <= ?
)
""",
request.location,
^(west - 0.005),
^(south - 0.005),
^(east + 0.005),
^(north + 0.005),
request.public_longitude,
^west,
request.public_longitude,
^east,
request.public_latitude,
^south,
request.public_latitude,
^north
)
)
end)
spatial_condition
end
def list_open_requests, do: raise(ArgumentError, "an authenticated scope is required")
def visible_open_request?(%Scope{user: user}, %HelpRequest{} = request, filters \\ %{}) do
request.status == :open and is_nil(request.hidden_at) and
DateTime.after?(request.expires_at, DateTime.utc_now(:second)) and
request_filter_matches?(request, filters) and
not Trust.blocked_between?(user.id, request.requester_id)
end
def list_my_requests(%Scope{user: user}) do
paginate_my_requests(%Scope{user: user}).entries
end
def paginate_my_requests(%Scope{user: user}, options \\ []) do
limit = Pagination.limit(options)
cursor = Pagination.cursor(options)
HelpRequest
|> where(
[request],
request.requester_id == ^user.id or
request.id in subquery(
from(assignment in Assignment,
where: assignment.helper_id == ^user.id,
select: assignment.request_id
)
)
)
|> before_my_request(cursor)
|> order_by([request], desc: request.inserted_at, desc: request.id)
|> limit(^(limit + 1))
|> Repo.all()
|> preload_request_relations()
|> Pagination.page(limit, &{&1.inserted_at, &1.id})
end
defp after_open_request(query, nil), do: query
defp after_open_request(query, {expires_at, id}) do
where(
query,
[request],
request.expires_at > ^expires_at or
(request.expires_at == ^expires_at and request.id > ^id)
)
end
defp before_my_request(query, nil), do: query
defp before_my_request(query, {inserted_at, id}) do
where(
query,
[request],
request.inserted_at < ^inserted_at or
(request.inserted_at == ^inserted_at and request.id < ^id)
)
end
def get_request!(id) do
HelpRequest
|> Repo.get!(id)
|> then(&preload_request_relations([&1], social_identities: true))
|> hd()
end
def get_request(%Scope{user: user} = scope, id) do
with {:ok, id} <- Ecto.UUID.cast(id),
%HelpRequest{} = request <- Repo.get(HelpRequest, id) do
request =
[request]
|> preload_request_relations(social_identities: true)
|> hd()
cond do
request.requester_id == user.id ->
{:ok, request}
Accounts.authorized?(user, :moderation_view) ->
{:ok, request}
not is_nil(request.hidden_at) ->
{:error, :not_found}
Trust.blocked_between?(user.id, request.requester_id) ->
{:error, :not_found}
request.assignment && participant?(scope, request.assignment) ->
{:ok, request}
request.status == :open and DateTime.after?(request.expires_at, DateTime.utc_now(:second)) ->
{:ok, request}
true ->
{:error, :not_found}
end
else
_invalid_or_missing -> {:error, :not_found}
end
end
def get_assignment_for_participant(%Scope{} = scope, id) do
with {:ok, id} <- Ecto.UUID.cast(id),
%Assignment{} = assignment <-
Assignment
|> Repo.get(id)
|> Repo.preload(:request),
true <- participant?(scope, assignment) do
{:ok, assignment}
else
false -> {:error, :forbidden}
nil -> {:error, :not_found}
:error -> {:error, :not_found}
end
end
def change_request(%HelpRequest{} = request, attrs \\ %{}) do
request
|> HelpRequest.create_changeset(attrs)
|> validate_structured_data()
end
def create_request(%Scope{user: user}, attrs) do
changeset =
%HelpRequest{requester_id: user.id}
|> HelpRequest.create_changeset(attrs)
|> validate_structured_data()
result =
if changeset.valid? do
with {:ok, _limit} <- Trust.authorize_action(Scope.for_user(user), :create_request) do
Repo.transact(fn ->
with {:ok, request} <- Repo.insert(changeset),
{:ok, _audit} <-
Trust.audit(user.id, "request.created", "request", request.id, %{
"urgency" => to_string(request.urgency),
"category_id" => request.category_id
}),
{:ok, _nearby_job} <- NearbyMatchWorker.enqueue(request.id) do
{:ok, request}
end
end)
end
else
{:error, changeset}
end
with {:ok, request} <- result do
_ =
ProductAnalytics.increment_for_user(
user,
"request.created",
to_string(request.urgency)
)
broadcast({:request_created, get_request!(request.id)})
{:ok, request}
end
end
def accept_request(%Scope{user: helper}, request_id) do
result =
with {:ok, request_id} <- cast_id(request_id),
{:ok, _limit} <- Trust.authorize_action(Scope.for_user(helper), :accept_request) do
now = DateTime.utc_now(:second)
code = handover_code(request_id)
Repo.transact(fn ->
request =
HelpRequest
|> where([r], r.id == ^request_id)
|> lock("FOR UPDATE")
|> Repo.one()
if request do
:ok = Trust.lock_user_pair(request.requester_id, helper.id)
cond do
request.requester_id == helper.id ->
{:error, :own_request}
not is_nil(request.hidden_at) ->
{:error, :not_open}
Trust.blocked_between?(request.requester_id, helper.id) ->
{:error, :blocked}
request.status != :open ->
{:error, :not_open}
DateTime.compare(request.expires_at, now) != :gt ->
{:error, :expired}
true ->
with {:ok, assignment} <-
%Assignment{}
|> Assignment.changeset(%{
request_id: request.id,
helper_id: helper.id,
active: true,
accepted_at: now,
handover_code_hash: code_hash(code)
})
|> Repo.insert(),
{:ok, _request} <-
request |> Ecto.Changeset.change(status: :matched) |> Repo.update(),
{:ok, _audit} <-
Trust.audit(helper.id, "request.accepted", "assignment", assignment.id, %{
"request_id" => request.id
}),
{:ok, _push_job} <-
Push.enqueue_request_accepted(
assignment.id,
request.id,
request.requester_id
) do
{:ok, assignment}
end
end
else
{:error, :not_found}
end
end)
else
:error -> {:error, :not_found}
end
case result do
{:ok, assignment} ->
_ = ProductAnalytics.increment_for_user(helper, "request.accepted")
request = get_request!(assignment.request_id)
broadcast({:request_updated, request})
{:ok,
Repo.preload(assignment,
helper: Accounts.public_user_query(),
request: []
)}
other ->
other
end
end
def start_assignment(%Scope{user: user}, assignment_id) do
transition_assignment(user, assignment_id, :start)
end
def arrive_assignment(%Scope{user: user}, assignment_id) do
transition_assignment(user, assignment_id, :arrive)
end
def confirm_completion(%Scope{user: user}, assignment_id) do
transition_assignment(user, assignment_id, :confirm)
end
def verify_handover(%Scope{user: user}, assignment_id, code) do
code = normalize_handover_code(code)
with {:ok, assignment_id} <- cast_id(assignment_id),
true <- valid_handover_code?(code),
{:ok, _limit} <- Trust.authorize_action(Scope.for_user(user), :verify_handover) do
Repo.transact(fn ->
with %Assignment{} = assignment <- locked_assignment(assignment_id) do
request = Repo.get!(HelpRequest, assignment.request_id)
:ok = Trust.lock_user_pair(assignment.helper_id, request.requester_id)
cond do
user.id != assignment.helper_id ->
{:error, :forbidden}
Trust.blocked_between?(assignment.helper_id, request.requester_id) ->
{:error, :blocked}
assignment.status not in [:accepted, :in_progress] or
not is_nil(assignment.handover_verified_at) ->
{:error, :invalid_transition}
not Plug.Crypto.secure_compare(code_hash(code), assignment.handover_code_hash) ->
{:error, :invalid_code}
true ->
now = DateTime.utc_now(:second)
assignment
|> Assignment.changeset(%{handover_verified_at: now})
|> maybe_complete(request)
|> audit_assignment_transition(user.id, "handover.verified", request.id)
|> notify_handover_verified(request)
end
else
nil -> {:error, :not_found}
end
end)
else
:error -> {:error, :not_found}
false -> {:error, :invalid_code}
end
|> after_transition()
end
def cancel_request(%Scope{user: user}, request_id, attrs \\ %{}) do
with {:ok, request_id} <- cast_id(request_id),
{:ok, _limit} <- Trust.authorize_action(Scope.for_user(user), :cancel_request) do
Repo.transact(fn ->
request =
HelpRequest
|> where([r], r.id == ^request_id)
|> lock("FOR UPDATE")
|> Repo.one()
cond do
is_nil(request) ->
{:error, :not_found}
request.requester_id == user.id and request.status in [:open, :matched, :in_progress] ->
now = DateTime.utc_now(:second)
assignment =
Assignment
|> where([assignment], assignment.request_id == ^request.id and assignment.active)
|> lock("FOR UPDATE")
|> Repo.one()
with {:ok, request} <-
request
|> HelpRequest.cancellation_changeset(cancellation_attrs(attrs, now))
|> Repo.update(),
{:ok, _assignment} <- cancel_assignment(assignment),
{:ok, _audit} <-
Trust.audit(user.id, "request.cancelled", "request", request.id, %{
"reason" => to_string(request.cancellation_reason)
}),
{:ok, _notification} <- notify_request_cancelled(request, assignment) do
{:ok, request}
end
true ->
{:error, :forbidden}
end
end)
else
:error -> {:error, :not_found}
end
|> case do
{:ok, request} ->
_ =
ProductAnalytics.increment_for_user(
user,
"request.cancelled",
to_string(request.cancellation_reason)
)
WhoNeedHelp.Tracking.cleanup_finished_sessions()
request = get_request!(request.id)
broadcast({:request_updated, request})
{:ok, request}
other ->
other
end
end
def withdraw_assignment(%Scope{user: user}, assignment_id, attrs \\ %{}) do
with {:ok, assignment_id} <- cast_id(assignment_id),
{:ok, _limit} <- Trust.authorize_action(Scope.for_user(user), :withdraw_assignment) do
Repo.transact(fn ->
with %Assignment{} = assignment <- locked_assignment(assignment_id) do
request = Repo.get!(HelpRequest, assignment.request_id)
if assignment.helper_id == user.id and
assignment.status in [:accepted, :in_progress] do
now = DateTime.utc_now(:second)
reopen? = DateTime.after?(request.expires_at, now)
with {:ok, assignment} <-
assignment
|> Assignment.withdrawal_changeset(withdrawal_attrs(attrs))
|> Repo.update(),
{:ok, _request} <-
request
|> Ecto.Changeset.change(
status: if(reopen?, do: :open, else: :expired),
cancelled_at: nil
)
|> Repo.update(),
{:ok, _audit} <-
Trust.audit(user.id, "assignment.withdrawn", "assignment", assignment.id, %{
"request_id" => request.id,
"reason" => to_string(assignment.withdrawal_reason),
"request_reopened" => reopen?
}),
{:ok, _notification} <- notify_assignment_withdrawn(request, assignment, reopen?),
{:ok, _nearby_job} <-
maybe_enqueue_reopened_request(request.id, assignment.id, reopen?) do
{:ok, assignment}
end
else
{:error, :invalid_transition}
end
else
nil -> {:error, :not_found}
end
end)
else
:error -> {:error, :not_found}
end
|> after_transition()
|> case do
{:ok, _assignment} = result ->
_ = ProductAnalytics.increment_for_user(user, "assignment.withdrawn")
WhoNeedHelp.Tracking.cleanup_finished_sessions()
result
other ->
other
end
end
def participant?(%Scope{user: user}, %Assignment{} = assignment) do
request =
case assignment.request do
%HelpRequest{} = request -> request
_ -> Repo.get!(HelpRequest, assignment.request_id)
end
user.id in [assignment.helper_id, request.requester_id]
end
defp preload_request_relations(requests, options \\ []) do
requests = Repo.preload(requests, [:category, :assignment])
user_ids =
requests
|> Enum.flat_map(fn request ->
[request.requester_id, request.assignment && request.assignment.helper_id]
end)
|> Enum.reject(&is_nil/1)
|> Enum.uniq()
users =
Accounts.public_user_query()
|> where([user], user.id in ^user_ids)
|> Repo.all()
|> maybe_preload_social_identities(options)
|> Map.new(&{&1.id, &1})
Enum.map(requests, fn request ->
assignment =
case request.assignment do
%Assignment{} = assignment ->
%{assignment | helper: Map.fetch!(users, assignment.helper_id)}
nil ->
nil
end
request
|> Map.put(:requester, Map.fetch!(users, request.requester_id))
|> Map.put(:assignment, assignment)
|> attach_assignment_request()
end)
end
defp maybe_preload_social_identities(users, options) do
if Keyword.get(options, :social_identities, false) do
Repo.preload(users, social_identities: Accounts.public_social_identity_query())
else
users
end
end
defp attach_assignment_request(%HelpRequest{assignment: %Assignment{} = assignment} = request) do
%{request | assignment: %{assignment | request: request}}
end
defp attach_assignment_request(%HelpRequest{} = request), do: request
def requester?(%Scope{user: user}, %HelpRequest{requester_id: id}), do: user.id == id
def helper?(%Scope{user: user}, %Assignment{helper_id: id}), do: user.id == id
def request_coordinates(%Scope{} = scope, %HelpRequest{} = request) do
owner = request.requester_id == scope.user.id
matched_participant =
not owner and not is_nil(request.assignment) and participant?(scope, request.assignment) and
request.assignment.status in [:accepted, :in_progress] and
not Trust.blocked_between?(scope.user.id, request.requester_id)
cond do
owner and request.location_visibility == :approximate_public ->
HelpRequest.public_coordinates(request)
owner ->
exact_request_coordinates(request)
matched_participant and request.location_visibility == :exact_for_active_match ->
exact_request_coordinates(request)
true ->
HelpRequest.public_coordinates(request)
end
end
defp exact_request_coordinates(%HelpRequest{
location: %Geo.Point{coordinates: {lng, lat}}
}) do
%{latitude: lat, longitude: lng, exact: true}
end
defp exact_request_coordinates(%HelpRequest{location: nil}), do: nil
def handover_code(request_id) do
secret = Application.fetch_env!(:who_need_help, :handover_secret)
digest = :crypto.mac(:hmac, :sha256, secret, request_id)
digest
|> :binary.decode_unsigned()
|> rem(1_000_000)
|> Integer.to_string()
|> String.pad_leading(6, "0")
end
defp transition_assignment(user, assignment_id, action) do
with {:ok, assignment_id} <- cast_id(assignment_id),
{:ok, _limit} <- Trust.authorize_action(Scope.for_user(user), action) do
Repo.transact(fn ->
with %Assignment{} = assignment <- locked_assignment(assignment_id) do
request = Repo.get!(HelpRequest, assignment.request_id)
:ok = Trust.lock_user_pair(assignment.helper_id, request.requester_id)
now = DateTime.utc_now(:second)
if Trust.blocked_between?(assignment.helper_id, request.requester_id) do
{:error, :blocked}
else
case {action, assignment.status, user.id} do
{:start, :accepted, helper_id} when helper_id == assignment.helper_id ->
with {:ok, assignment} <-
assignment
|> Assignment.changeset(%{status: :in_progress, started_at: now})
|> Repo.update(),
{:ok, _} <-
request |> Ecto.Changeset.change(status: :in_progress) |> Repo.update(),
{:ok, _audit} <-
Trust.audit(
user.id,
"assignment.started",
"assignment",
assignment.id,
%{"request_id" => request.id}
),
{:ok, _notification} <-
Push.enqueue_lifecycle(
:assignment_started,
assignment.id,
request.id,
request.requester_id,
"Your helper started",
"Open Who Need Help to follow the request status."
) do
{:ok, assignment}
end
{:arrive, :in_progress, helper_id}
when helper_id == assignment.helper_id and is_nil(assignment.arrived_at) ->
with {:ok, assignment} <-
assignment
|> Assignment.changeset(%{arrived_at: now})
|> Repo.update(),
{:ok, _audit} <-
Trust.audit(
user.id,
"assignment.arrived",
"assignment",
assignment.id,
%{"request_id" => request.id}
),
{:ok, _notification} <-
Push.enqueue_lifecycle(
:helper_arrived,
assignment.id,
request.id,
request.requester_id,
"Your helper arrived",
"Open Who Need Help to coordinate the handover safely."
) do
{:ok, assignment}
end
{:confirm, status, user_id} when status in [:accepted, :in_progress] ->
attrs =
cond do
user_id == assignment.helper_id -> %{helper_confirmed_at: now}
user_id == request.requester_id -> %{requester_confirmed_at: now}
true -> nil
end
if attrs do
assignment
|> Assignment.changeset(attrs)
|> maybe_complete(request)
|> audit_assignment_transition(user.id, "assignment.confirmed", request.id)
else
{:error, :forbidden}
end
_ ->
{:error, :invalid_transition}
end
end
else
nil -> {:error, :not_found}
end
end)
else
:error -> {:error, :not_found}
end
|> after_transition()
|> record_assignment_metric(user, action)
end
defp maybe_complete(changeset, request) do
assignment = Ecto.Changeset.apply_changes(changeset)
if assignment.handover_verified_at && assignment.requester_confirmed_at &&
assignment.helper_confirmed_at do
now = DateTime.utc_now(:second)
Multi.new()
|> Multi.update(
:assignment,
Ecto.Changeset.change(changeset, status: :completed, completed_at: now)
)
|> Multi.update(
:request,
Ecto.Changeset.change(request, status: :completed, completed_at: now)
)
|> Repo.transaction()
|> case do
{:ok, %{assignment: assignment}} -> {:ok, assignment}
{:error, _step, reason, _changes} -> {:error, reason}
end
else
Repo.update(changeset)
end
end
defp locked_assignment(id) do
Assignment
|> where([a], a.id == ^id and a.active)
|> lock("FOR UPDATE")
|> Repo.one()
end
defp valid_handover_code?(code) when is_binary(code),
do: Regex.match?(~r/^\d{6}$/, code)
defp valid_handover_code?(_code), do: false
defp normalize_handover_code(code) when is_binary(code) do
code
|> String.slice(0, 64)
|> String.replace(~r/[^0-9]/u, "")
end
defp normalize_handover_code(code), do: code
defp code_hash(code) when is_binary(code), do: :crypto.hash(:sha256, code)
defp after_transition({:ok, assignment}) do
request = get_request!(assignment.request_id)
broadcast({:request_updated, request})
if assignment.status == :completed, do: Trust.record_completion_signals(assignment)
{:ok,
Repo.preload(
assignment,
[helper: Accounts.public_user_query(), request: []],
force: true
)}
end
defp after_transition(other), do: other
defp audit_assignment_transition({:ok, assignment}, actor_id, action, request_id) do
with {:ok, _audit} <-
Trust.audit(actor_id, action, "assignment", assignment.id, %{
"request_id" => request_id,
"status" => to_string(assignment.status)
}),
{:ok, _completion_audit} <-
maybe_audit_completion(actor_id, assignment, request_id),
{:ok, _notifications} <-
maybe_notify_completion(assignment, request_id) do
{:ok, assignment}
end
end
defp audit_assignment_transition(other, _actor_id, _action, _request_id), do: other
defp maybe_audit_completion(actor_id, %Assignment{status: :completed} = assignment, request_id) do
Trust.audit(actor_id, "request.completed", "request", request_id, %{
"assignment_id" => assignment.id
})
end
defp maybe_audit_completion(_actor_id, _assignment, _request_id), do: {:ok, :not_completed}
defp cancel_assignment(nil), do: {:ok, :no_assignment}
defp cancel_assignment(assignment) do
assignment
|> Assignment.changeset(%{status: :cancelled})
|> Repo.update()
end
defp cancellation_attrs(attrs, now) do
attrs
|> normalize_attrs()
|> Map.put_new("cancellation_reason", "no_longer_needed")
|> Map.put("status", "cancelled")
|> Map.put("cancelled_at", now)
end
defp withdrawal_attrs(attrs) do
attrs
|> normalize_attrs()
|> Map.put_new("withdrawal_reason", "cannot_complete")
|> Map.put("status", "cancelled")
|> Map.put("active", false)
end
defp normalize_attrs(attrs) when is_map(attrs) do
Map.new(attrs, fn {key, value} -> {to_string(key), value} end)
end
defp notify_request_cancelled(_request, nil), do: {:ok, :no_helper}
defp notify_request_cancelled(request, assignment) do
Push.enqueue_lifecycle(
:request_cancelled,
request.id,
request.id,
assignment.helper_id,
"Request cancelled",
"The requester cancelled this request. Open Who Need Help for details.",
%{"event_variant" => "requester_cancelled"}
)
end
defp notify_assignment_withdrawn(request, _assignment, true) do
Push.enqueue_lifecycle(
:request_reopened,
request.id,
request.id,
request.requester_id,
"Your request needs a new helper",
"The previous helper withdrew, so the request is open again."
)
end
defp notify_assignment_withdrawn(request, _assignment, false) do
Push.enqueue_lifecycle(
:request_cancelled,
request.id,
request.id,
request.requester_id,
"Helper withdrew after expiry",
"The helper withdrew and the request is no longer open because it expired.",
%{"event_variant" => "helper_withdrew_after_expiry"}
)
end
defp maybe_enqueue_reopened_request(request_id, assignment_id, true) do
NearbyMatchWorker.enqueue(request_id, "reopened:#{assignment_id}")
end
defp maybe_enqueue_reopened_request(_request_id, _assignment_id, false),
do: {:ok, :not_reopened}
defp notify_handover_verified({:ok, assignment} = result, request) do
case Push.enqueue_lifecycle(
:handover_verified,
assignment.id,
request.id,
request.requester_id,
"Handover code verified",
"The helper verified the one-time handover code."
) do
{:ok, _notification} -> result
{:error, reason} -> {:error, reason}
end
end
defp notify_handover_verified(other, _request), do: other
defp maybe_notify_completion(%Assignment{status: :completed} = assignment, request_id) do
request = Repo.get!(HelpRequest, request_id)
with {:ok, _requester_notification} <-
Push.enqueue_lifecycle(
:request_completed,
assignment.id,
request_id,
request.requester_id,
"Help completed",
"Both participants confirmed completion and the handover was verified."
),
{:ok, _helper_notification} <-
Push.enqueue_lifecycle(
:request_completed,
assignment.id,
request_id,
assignment.helper_id,
"Help completed",
"Both participants confirmed completion and the handover was verified."
) do
{:ok, :notified}
end
end
defp maybe_notify_completion(_assignment, _request_id), do: {:ok, :not_completed}
defp record_assignment_metric({:ok, _assignment} = result, user, :start) do
_ = ProductAnalytics.increment_for_user(user, "assignment.started")
result
end
defp record_assignment_metric({:ok, _assignment} = result, user, :arrive) do
_ = ProductAnalytics.increment_for_user(user, "assignment.arrived")
result
end
defp record_assignment_metric(
{:ok, %Assignment{status: :completed}} = result,
user,
:confirm
) do
_ = ProductAnalytics.increment_for_user(user, "request.completed")
result
end
defp record_assignment_metric(result, _user, _action), do: result
defp broadcast(event) do
Phoenix.PubSub.broadcast(WhoNeedHelp.PubSub, @topic, event)
request =
case event do
{_, %HelpRequest{} = request} -> request
end
Phoenix.PubSub.broadcast(WhoNeedHelp.PubSub, "help:request:#{request.id}", event)
request
|> discovery_broadcast_topics()
|> Enum.each(&Phoenix.PubSub.broadcast(WhoNeedHelp.PubSub, &1, event))
[request.requester_id, request.assignment && request.assignment.helper_id]
|> Enum.reject(&is_nil/1)
|> Enum.uniq()
|> Enum.each(fn user_id ->
Phoenix.PubSub.broadcast(WhoNeedHelp.PubSub, "help:user:#{user_id}", event)
end)
end
defp discovery_topics(viewport, include_hidden) do
topics =
viewport
|> DiscoveryViewport.subscription_tiles()
|> Enum.map(&tile_topic/1)
if include_hidden, do: [@hidden_discovery_topic | topics], else: topics
end
defp discovery_broadcast_topics(%HelpRequest{location: nil}), do: [@hidden_discovery_topic]
defp discovery_broadcast_topics(%HelpRequest{
location: %Geo.Point{coordinates: {longitude, latitude}},
location_radius_meters: radius_meters
}) do
latitude
|> DiscoveryViewport.point_tiles(longitude, radius_meters || 0)
|> Enum.map(&tile_topic/1)
end
defp discovery_broadcast_topics(_request), do: []
defp tile_topic({zoom, x, y}), do: "help:discovery:#{zoom}:#{x}:#{y}"
defp filter_open_requests(query, filters) do
query
|> maybe_filter(:category_id, filters["category_id"] || filters[:category_id])
|> maybe_filter(:urgency, filters["urgency"] || filters[:urgency])
|> maybe_filter(:area, filters["area"] || filters[:area])
end
defp request_filter_matches?(request, filters) do
category_id = filters["category_id"] || filters[:category_id]
urgency = filters["urgency"] || filters[:urgency]
area = filters["area"] || filters[:area]
(category_id in [nil, ""] or to_string(request.category_id) == to_string(category_id)) and
(urgency in [nil, ""] or to_string(request.urgency) == to_string(urgency)) and
(area in [nil, ""] or
String.contains?(String.downcase(request.location_label), String.downcase(area)))
end
defp maybe_filter(query, _field, value) when value in [nil, ""], do: query
defp maybe_filter(query, :category_id, value) do
case Ecto.UUID.cast(value) do
{:ok, category_id} -> where(query, [request], request.category_id == ^category_id)
:error -> where(query, [request], false)
end
end
defp maybe_filter(query, :urgency, value) do
case normalize_urgency(value) do
{:ok, urgency} -> where(query, [request], request.urgency == ^urgency)
:error -> where(query, [request], false)
end
end
defp maybe_filter(query, :area, value) when is_binary(value) do
value = value |> String.trim() |> String.slice(0, 100)
if value == "" do
query
else
pattern = "%#{escape_like(value)}%"
where(
query,
[request],
fragment("? ILIKE ? ESCAPE E'\\\\'", request.location_label, ^pattern)
)
end
end
defp maybe_filter(query, :area, _value), do: where(query, [request], false)
defp escape_like(value) do
value
|> String.replace("\\", "\\\\")
|> String.replace("%", "\\%")
|> String.replace("_", "\\_")
end
defp blank_area_filter?(filters) do
case filters["area"] || filters[:area] do
nil -> true
area when is_binary(area) -> String.trim(area) == ""
_invalid -> false
end
end
defp open_requests_query(user, filters, viewport) do
now = DateTime.utc_now(:second)
HelpRequest
|> where(
[request],
request.status == :open and request.expires_at > ^now and is_nil(request.hidden_at)
)
|> where(
[request],
request.requester_id not in subquery(
from(block in Block, where: block.blocker_id == ^user.id, select: block.blocked_id)
)
)
|> where(
[request],
request.requester_id not in subquery(
from(block in Block, where: block.blocked_id == ^user.id, select: block.blocker_id)
)
)
|> filter_open_requests(filters)
|> filter_discovery_viewport(viewport, filters)
end
defp filter_discovery_viewport(query, nil, _filters), do: query
defp filter_discovery_viewport(query, %DiscoveryViewport{} = viewport, filters) do
spatial_condition = request_map_viewport_condition(viewport)
area = filters["area"] || filters[:area]
condition =
if is_binary(area) and String.trim(area) != "" do
dynamic(
[request],
(not is_nil(request.location) and request.location_visibility != :hidden and
^spatial_condition) or request.location_visibility == :hidden or
is_nil(request.location)
)
else
dynamic(
[request],
not is_nil(request.location) and request.location_visibility != :hidden and
^spatial_condition
)
end
where(query, ^condition)
end
defp map_discovery_item(%{count: 1} = row, level) do
%{
type: "request",
id: row.request_id,
parent_cluster_id: DiscoveryCluster.parent_id(:request, level, row.cell_x, row.cell_y),
hierarchy_level: level,
title: row.title,
location: row.location_label,
latitude: row.latitude,
longitude: row.longitude,
exact: row.location_visibility in ["exact_public", :exact_public],
radius_meters: row.radius_meters
}
end
defp map_discovery_item(row, level) do
cluster_id = DiscoveryCluster.cluster_id(:request, level, row.cell_x, row.cell_y)
{longitude, latitude} = DiscoveryCluster.cell_center(level, row.cell_x, row.cell_y)
%{
type: "cluster",
id: cluster_id,
cluster_id: cluster_id,
parent_cluster_id: DiscoveryCluster.parent_id(:request, level, row.cell_x, row.cell_y),
hierarchy_level: level,
point_count: row.count,
expansion_zoom:
DiscoveryCluster.expansion_zoom_for_bounds(
level,
row.west,
row.south,
row.east,
row.north
),
count: row.count,
urgent_count: row.urgent_count,
latitude: latitude,
longitude: longitude,
bounds: %{
west: row.west,
south: row.south,
east: row.east,
north: row.north
}
}
end
defp hydrate_request_map_singletons(rows) do
singleton_ids =
rows
|> Enum.filter(&(&1.count == 1))
|> Enum.map(& &1.request_id)
details = request_map_singleton_details(singleton_ids)
Enum.map(rows, fn
%{count: 1, request_id: request_id} = row ->
Map.merge(row, Map.fetch!(details, request_id))
row ->
row
end)
end
defp request_map_singleton_details([]), do: %{}
defp request_map_singleton_details(singleton_ids) do
HelpRequest
|> where([request], request.id in ^singleton_ids)
|> select([request], {
request.id,
%{
title: request.title,
location_label: request.location_label,
location_visibility: request.location_visibility,
radius_meters: request.location_radius_meters
}
})
|> Repo.all()
|> Map.new(fn {id, detail} -> {to_string(id), detail} end)
end
defp normalize_urgency(value) when value in [:now, "now"], do: {:ok, :now}
defp normalize_urgency(value) when value in [:today, "today"], do: {:ok, :today}
defp normalize_urgency(value) when value in [:scheduled, "scheduled"], do: {:ok, :scheduled}
defp normalize_urgency(_value), do: :error
defp cast_id(id) do
case Ecto.UUID.cast(id) do
{:ok, id} -> {:ok, id}
:error -> :error
end
end
defp validate_structured_data(changeset) do
category_id = Ecto.Changeset.get_field(changeset, :category_id)
case category_id && Repo.get(Category, category_id) do
%Category{mode: :help, active: true} = category ->
data = Ecto.Changeset.get_field(changeset, :structured_data)
case Catalog.validate_structured_data(category, data) do
{:ok, clean} ->
Ecto.Changeset.put_change(changeset, :structured_data, clean)
{:error, errors} ->
Ecto.Changeset.add_error(
changeset,
:structured_data,
Enum.join(errors, "; ")
)
end
_ ->
Ecto.Changeset.add_error(changeset, :category_id, "select an active help category")
end
end
end