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, "email_enabled" => subscription.email_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 email_allowed?( %Preference{} = preference, %Notification{kind: :nearby_request} = notification ) do preference.email_enabled and preference.nearby_email_enabled and Map.get(notification.data, "email_enabled", false) end def email_allowed?(%Preference{}, %Notification{}), do: false 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