who_need_help/lib/who_need_help/notifications.ex

496 lines
16 KiB
Elixir

defmodule WhoNeedHelp.Notifications do
@moduledoc """
User-owned notification inbox, nearby-help subscriptions, and device registrations.
Exact subscription centers and device delivery credentials are private. Public
request discovery never reads or exposes them.
"""
import Ecto.Query
alias Ecto.Changeset
alias WhoNeedHelp.Accounts.{Scope, User}
alias WhoNeedHelp.Help.HelpRequest
alias WhoNeedHelp.Notifications.{NearbySubscription, Notification, Preference, PushDevice}
alias WhoNeedHelp.Pagination
alias WhoNeedHelp.Push.NotificationDispatchWorker
alias WhoNeedHelp.Repo
alias WhoNeedHelp.Trust.Block
@topic_prefix "notifications:user:"
def subscribe(%Scope{user: %User{id: user_id}}), do: subscribe_user(user_id)
def subscribe_user(user_id) when is_binary(user_id) do
Phoenix.PubSub.subscribe(WhoNeedHelp.PubSub, @topic_prefix <> user_id)
end
def paginate_notifications(%Scope{user: %User{id: user_id}}, options \\ []) do
limit = Pagination.limit(options)
cursor = Pagination.cursor(options)
Notification
|> where([notification], notification.user_id == ^user_id)
|> before_notification(cursor)
|> order_by([notification], desc: notification.inserted_at, desc: notification.id)
|> limit(^(limit + 1))
|> Repo.all()
|> Pagination.page(limit, &{&1.inserted_at, &1.id})
end
def unread_count(%Scope{user: %User{id: user_id}}) do
Notification
|> where([notification], notification.user_id == ^user_id)
|> where([notification], is_nil(notification.read_at))
|> Repo.aggregate(:count)
end
def mark_read(%Scope{user: %User{id: user_id}}, notification_id) do
with {:ok, notification_id} <- Ecto.UUID.cast(notification_id),
%Notification{} = notification <-
Repo.get_by(Notification, id: notification_id, user_id: user_id) do
now = DateTime.utc_now(:second)
case notification
|> Changeset.change(read_at: notification.read_at || now)
|> Repo.update() do
{:ok, notification} ->
broadcast(user_id, {:notification_read, notification.id})
{:ok, notification}
other ->
other
end
else
_missing_or_invalid -> {:error, :not_found}
end
end
def mark_all_read(%Scope{user: %User{id: user_id}}) do
now = DateTime.utc_now(:second)
{count, _} =
Notification
|> where([notification], notification.user_id == ^user_id)
|> where([notification], is_nil(notification.read_at))
|> Repo.update_all(set: [read_at: now, updated_at: now])
broadcast(user_id, {:notifications_read, count})
{:ok, count}
end
def notify_user(user_id, attrs) when is_binary(user_id) and is_map(attrs) do
attrs = Map.put(attrs, :user_id, user_id)
changeset = Notification.changeset(%Notification{}, attrs)
idempotency_key = Map.get(attrs, :idempotency_key) || Map.get(attrs, "idempotency_key")
Repo.transact(fn ->
case Repo.insert(changeset,
on_conflict: :nothing,
conflict_target: [:idempotency_key]
) do
{:ok, _inserted_or_conflicted} ->
notification = Repo.get_by!(Notification, idempotency_key: idempotency_key)
with {:ok, _job} <- enqueue_dispatch(notification), do: {:ok, notification}
{:error, %Changeset{} = changeset} ->
{:error, changeset}
end
end)
end
def broadcast_created(%Notification{} = notification) do
broadcast(notification.user_id, {:notification_created, notification})
:ok
end
def notification_opened(%Scope{user: user} = scope, notification_id) do
case mark_read(scope, notification_id) do
{:ok, notification} = result ->
_ =
WhoNeedHelp.ProductAnalytics.increment_for_user(
user,
"notification.opened",
to_string(notification.kind)
)
result
other ->
other
end
end
def get_preference(%Scope{user: %User{id: user_id}}) do
Repo.get_by(Preference, user_id: user_id) || %Preference{user_id: user_id}
end
def change_preference(%Preference{} = preference, attrs \\ %{}) do
Preference.changeset(preference, attrs)
end
def update_preference(%Scope{user: %User{id: user_id}}, attrs) do
preference = Repo.get_by(Preference, user_id: user_id) || %Preference{user_id: user_id}
preference
|> Preference.changeset(put_attr(attrs, :user_id, user_id))
|> Repo.insert_or_update()
end
def list_nearby_subscriptions(%Scope{user: %User{id: user_id}}) do
NearbySubscription
|> where([subscription], subscription.user_id == ^user_id)
|> order_by([subscription], desc: subscription.active, asc: subscription.name)
|> Repo.all()
|> Repo.preload(:user)
end
def change_nearby_subscription(%NearbySubscription{} = subscription, attrs \\ %{}) do
NearbySubscription.changeset(subscription, attrs)
end
def create_nearby_subscription(%Scope{user: %User{id: user_id}}, attrs) do
%NearbySubscription{user_id: user_id}
|> NearbySubscription.changeset(put_attr(attrs, :user_id, user_id))
|> Repo.insert()
end
def update_nearby_subscription(
%Scope{user: %User{id: user_id}},
subscription_id,
attrs
) do
with {:ok, subscription_id} <- Ecto.UUID.cast(subscription_id),
%NearbySubscription{} = subscription <-
Repo.get_by(NearbySubscription, id: subscription_id, user_id: user_id) do
subscription
|> NearbySubscription.changeset(put_attr(attrs, :user_id, user_id))
|> Repo.update()
else
_missing_or_invalid -> {:error, :not_found}
end
end
def delete_nearby_subscription(%Scope{user: %User{id: user_id}}, subscription_id) do
with {:ok, subscription_id} <- Ecto.UUID.cast(subscription_id),
%NearbySubscription{} = subscription <-
Repo.get_by(NearbySubscription, id: subscription_id, user_id: user_id) do
Repo.delete(subscription)
else
_missing_or_invalid -> {:error, :not_found}
end
end
def register_device(%Scope{user: %User{id: user_id}}, attrs) do
now = DateTime.utc_now(:second)
attrs =
attrs
|> put_attr(:user_id, user_id)
|> put_attr(:last_seen_at, now)
|> put_attr(:disabled_at, nil)
changeset =
%PushDevice{user_id: user_id}
|> PushDevice.changeset(attrs)
with true <- changeset.valid?,
digest when is_binary(digest) <- Changeset.get_field(changeset, :token_digest),
platform when platform in [:web, :android] <- Changeset.get_field(changeset, :platform),
installation_id when is_binary(installation_id) <-
Changeset.get_field(changeset, :installation_id) do
Repo.transact(fn ->
installation_device = locked_device_by_installation(platform, installation_id)
token_device = locked_device_by_token(digest)
cond do
installation_device && token_device && installation_device.id != token_device.id ->
{:error, :already_registered}
installation_device ->
installation_device
|> PushDevice.changeset(attrs)
|> Repo.update()
token_device && token_device.user_id == user_id ->
token_device
|> PushDevice.changeset(attrs)
|> Repo.update()
token_device ->
{:error, :already_registered}
true ->
Repo.insert(changeset)
end
end)
else
false -> {:error, changeset}
nil -> {:error, changeset}
_invalid -> {:error, changeset}
end
end
def list_devices(%Scope{user: %User{id: user_id}}) do
PushDevice
|> where([device], device.user_id == ^user_id and is_nil(device.disabled_at))
|> order_by([device], desc: device.last_seen_at)
|> Repo.all()
end
def disable_device(%Scope{user: %User{id: user_id}}, device_id) do
with {:ok, device_id} <- Ecto.UUID.cast(device_id),
%PushDevice{} = device <- Repo.get_by(PushDevice, id: device_id, user_id: user_id) do
device
|> Changeset.change(disabled_at: DateTime.utc_now(:second))
|> Repo.update()
else
_missing_or_invalid -> {:error, :not_found}
end
end
def active_devices(user_id) do
PushDevice
|> where([device], device.user_id == ^user_id and is_nil(device.disabled_at))
|> order_by([device], desc: device.last_seen_at)
|> Repo.all()
end
def disable_invalid_device(%PushDevice{} = device) do
device
|> Changeset.change(disabled_at: DateTime.utc_now(:second))
|> Repo.update()
end
defp locked_device_by_installation(platform, installation_id) do
PushDevice
|> where(
[device],
device.platform == ^platform and device.installation_id == ^installation_id
)
|> lock("FOR UPDATE")
|> Repo.one()
end
defp locked_device_by_token(digest) do
PushDevice
|> where([device], device.token_digest == ^digest)
|> lock("FOR UPDATE")
|> Repo.one()
end
def matching_nearby_subscriptions(%HelpRequest{location: nil}), do: []
def matching_nearby_subscriptions(%HelpRequest{} = request) do
now = DateTime.utc_now(:second)
category_id = Ecto.UUID.dump!(request.category_id)
NearbySubscription
|> join(:inner, [subscription], user in User, on: user.id == subscription.user_id)
|> join(:left, [subscription, _user], preference in Preference,
on: preference.user_id == subscription.user_id
)
|> join(:left, [subscription, _user, _preference], outgoing_block in Block,
on:
outgoing_block.blocker_id == subscription.user_id and
outgoing_block.blocked_id == ^request.requester_id
)
|> join(:left, [subscription, _user, _preference, _outgoing_block], incoming_block in Block,
on:
incoming_block.blocker_id == ^request.requester_id and
incoming_block.blocked_id == subscription.user_id
)
|> where(
[subscription, user, _preference, _outgoing_block, _incoming_block],
subscription.active and user.id != ^request.requester_id and
user.moderation_status == :active and not is_nil(user.confirmed_at) and
not is_nil(user.accepted_terms_at)
)
|> where(
[subscription, _user, _preference, _outgoing_block, _incoming_block],
fragment(
"ST_DWithin(?::geography, ?::geography, ?)",
subscription.center,
^request.location,
subscription.radius_meters
)
)
|> where(
[subscription, _user, _preference, _outgoing_block, _incoming_block],
fragment(
"cardinality(?) = 0 OR ? = ANY(?)",
subscription.category_ids,
^category_id,
subscription.category_ids
)
)
|> where(
[subscription, _user, _preference, _outgoing_block, _incoming_block],
fragment("? = ANY(?)", ^to_string(request.urgency), subscription.urgencies)
)
|> where(
[_subscription, _user, _preference, outgoing_block, incoming_block],
is_nil(outgoing_block.id) and is_nil(incoming_block.id)
)
|> select([subscription, _user, preference, _outgoing_block, _incoming_block], {
subscription,
preference
})
|> Repo.all()
|> Enum.filter(fn {subscription, preference} ->
available_now?(subscription, preference || %Preference{}, now)
end)
|> Enum.map(&elem(&1, 0))
|> merge_user_subscriptions()
end
def nearby_notification_attrs(
%HelpRequest{} = request,
%NearbySubscription{} = subscription,
event_key \\ "created"
) do
%{
kind: :nearby_request,
title: "New help request nearby",
body: "A request matching one of your nearby-help alerts is available.",
path: "/requests/#{request.id}",
data: %{
"request_id" => request.id,
"category_id" => request.category_id,
"urgency" => to_string(request.urgency),
"push_enabled" => subscription.push_enabled
},
idempotency_key:
"nearby-request:#{request.id}:#{subscription.user_id}:#{safe_event_key(event_key)}"
}
end
def push_allowed?(%Preference{} = preference, %Notification{kind: kind} = notification) do
preference.push_enabled and
case kind do
:nearby_request ->
preference.nearby_push_enabled and
Map.get(notification.data, "push_enabled", true)
:message_created ->
preference.message_push_enabled
_other ->
preference.lifecycle_push_enabled
end
end
def quiet_now?(%Preference{quiet_hours_enabled: false}, _now), do: false
def quiet_now?(
%Preference{quiet_start: nil, quiet_end: nil},
_now
),
do: false
def quiet_now?(%Preference{} = preference, %DateTime{} = now) do
local_time =
now
|> DateTime.add(preference.utc_offset_minutes * 60, :second)
|> DateTime.to_time()
|> Time.truncate(:second)
time_between?(local_time, preference.quiet_start, preference.quiet_end)
end
def next_quiet_end(%Preference{} = preference, %DateTime{} = now) do
local_now = DateTime.add(now, preference.utc_offset_minutes * 60, :second)
local_date = DateTime.to_date(local_now)
local_time = local_now |> DateTime.to_time() |> Time.truncate(:second)
end_date =
if Time.compare(preference.quiet_start, preference.quiet_end) == :lt or
Time.compare(local_time, preference.quiet_end) == :lt do
local_date
else
Date.add(local_date, 1)
end
{:ok, naive} = NaiveDateTime.new(end_date, preference.quiet_end)
naive
|> DateTime.from_naive!("Etc/UTC")
|> DateTime.add(-preference.utc_offset_minutes * 60, :second)
end
defp available_now?(
%NearbySubscription{} = subscription,
%Preference{} = preference,
%DateTime{} = now
) do
local = DateTime.add(now, preference.utc_offset_minutes * 60, :second)
day = local |> DateTime.to_date() |> Date.day_of_week()
time = local |> DateTime.to_time() |> Time.truncate(:second)
day in subscription.available_days and
case {subscription.available_from, subscription.available_until} do
{nil, nil} -> true
{from, until} -> time_between?(time, from, until)
end
end
defp time_between?(time, from, until) do
case Time.compare(from, until) do
:lt -> Time.compare(time, from) != :lt and Time.compare(time, until) == :lt
:gt -> Time.compare(time, from) != :lt or Time.compare(time, until) == :lt
:eq -> true
end
end
defp enqueue_dispatch(notification) do
notification.id
|> then(&NotificationDispatchWorker.new(%{notification_id: &1}))
|> Oban.insert()
end
defp safe_event_key(event_key) do
event_key
|> to_string()
|> String.replace(~r/[^a-zA-Z0-9:_-]/u, "-")
|> String.slice(0, 100)
end
defp merge_user_subscriptions(subscriptions) do
subscriptions
|> Enum.group_by(& &1.user_id)
|> Enum.map(fn {_user_id, matches} ->
base = Enum.min_by(matches, & &1.id)
%{
base
| push_enabled: Enum.any?(matches, & &1.push_enabled),
email_enabled: Enum.any?(matches, & &1.email_enabled)
}
end)
end
defp before_notification(query, nil), do: query
defp before_notification(query, {inserted_at, id}) do
where(
query,
[notification],
notification.inserted_at < ^inserted_at or
(notification.inserted_at == ^inserted_at and notification.id < ^id)
)
end
defp broadcast(user_id, event) do
Phoenix.PubSub.broadcast(WhoNeedHelp.PubSub, @topic_prefix <> user_id, event)
end
defp put_attr(attrs, key, value) when is_map(attrs) and is_atom(key) do
if Enum.any?(Map.keys(attrs), &is_binary/1) do
Map.put(attrs, Atom.to_string(key), value)
else
Map.put(attrs, key, value)
end
end
end