122 lines
3.8 KiB
Elixir
122 lines
3.8 KiB
Elixir
defmodule WhoNeedHelp.Messaging do
|
|
@moduledoc "Durable match chat with PubSub fan-out after commit."
|
|
|
|
import Ecto.Query
|
|
alias WhoNeedHelp.Accounts
|
|
alias WhoNeedHelp.Accounts.Scope
|
|
alias WhoNeedHelp.Help
|
|
alias WhoNeedHelp.Help.Assignment
|
|
alias WhoNeedHelp.Messaging.Message
|
|
alias WhoNeedHelp.Pagination
|
|
alias WhoNeedHelp.Push
|
|
alias WhoNeedHelp.Repo
|
|
alias WhoNeedHelp.Trust
|
|
|
|
def subscribe(assignment_id) do
|
|
Phoenix.PubSub.subscribe(WhoNeedHelp.PubSub, "messages:#{assignment_id}")
|
|
end
|
|
|
|
def list_messages(%Scope{} = scope, %Assignment{} = assignment) do
|
|
paginate_messages(scope, assignment).entries
|
|
end
|
|
|
|
def paginate_messages(%Scope{} = scope, %Assignment{} = assignment, options \\ []) do
|
|
assignment = Repo.preload(assignment, :request)
|
|
|
|
if Trust.eligible?(scope) and Help.participant?(scope, assignment) and
|
|
not blocked_assignment?(scope, assignment) do
|
|
limit = Pagination.limit(options, 50)
|
|
cursor = Pagination.cursor(options)
|
|
public_user = Accounts.public_user_query()
|
|
|
|
page =
|
|
Message
|
|
|> where([message], message.assignment_id == ^assignment.id)
|
|
|> before_message(cursor)
|
|
|> order_by([message], desc: message.inserted_at, desc: message.id)
|
|
|> limit(^(limit + 1))
|
|
|> preload([message], sender: ^public_user)
|
|
|> Repo.all()
|
|
|> Pagination.page(limit, &{&1.inserted_at, &1.id})
|
|
|
|
%{page | entries: Enum.reverse(page.entries)}
|
|
else
|
|
%Pagination.Page{}
|
|
end
|
|
end
|
|
|
|
def send_message(%Scope{user: user} = scope, %Assignment{} = assignment, attrs) do
|
|
with {:ok, _limit} <- Trust.authorize_action(scope, :send_message) do
|
|
result =
|
|
Repo.transact(fn ->
|
|
current =
|
|
Assignment
|
|
|> where([current], current.id == ^assignment.id)
|
|
|> lock("FOR UPDATE")
|
|
|> Repo.one()
|
|
|
|
if current do
|
|
current = Repo.preload(current, :request)
|
|
request = current.request
|
|
recipient_id = counterpart_id(user.id, current, request)
|
|
|
|
with true <- Help.participant?(scope, current),
|
|
:ok <- Trust.lock_user_pair(user.id, recipient_id),
|
|
false <- Trust.blocked_between?(user.id, recipient_id),
|
|
{:ok, message} <-
|
|
%Message{assignment_id: current.id, sender_id: user.id}
|
|
|> Message.changeset(attrs)
|
|
|> Repo.insert(),
|
|
{:ok, _push_job} <-
|
|
Push.enqueue_message_created(
|
|
message.id,
|
|
current.id,
|
|
request.id,
|
|
recipient_id
|
|
) do
|
|
{:ok, message}
|
|
else
|
|
true -> {:error, :blocked}
|
|
false -> {:error, :forbidden}
|
|
other -> other
|
|
end
|
|
else
|
|
{:error, :not_found}
|
|
end
|
|
end)
|
|
|
|
with {:ok, message} <- result do
|
|
message = Repo.preload(message, sender: Accounts.public_user_query())
|
|
|
|
Phoenix.PubSub.broadcast(
|
|
WhoNeedHelp.PubSub,
|
|
"messages:#{message.assignment_id}",
|
|
{:new_message, message}
|
|
)
|
|
|
|
{:ok, message}
|
|
end
|
|
end
|
|
end
|
|
|
|
defp blocked_assignment?(%Scope{user: user}, assignment) do
|
|
request = Repo.preload(assignment, :request).request
|
|
Trust.blocked_between?(user.id, counterpart_id(user.id, assignment, request))
|
|
end
|
|
|
|
defp counterpart_id(user_id, assignment, request) do
|
|
if user_id == assignment.helper_id, do: request.requester_id, else: assignment.helper_id
|
|
end
|
|
|
|
defp before_message(query, nil), do: query
|
|
|
|
defp before_message(query, {inserted_at, id}) do
|
|
where(
|
|
query,
|
|
[message],
|
|
message.inserted_at < ^inserted_at or
|
|
(message.inserted_at == ^inserted_at and message.id < ^id)
|
|
)
|
|
end
|
|
end
|