Page MenuHomePhorge

No OneTemporary

Size
8 KB
Referenced Files
None
Subscribers
None
diff --git a/lib/pleroma/web/fed_sockets/adapter.ex b/lib/pleroma/web/fed_sockets/adapter.ex
index 00ab2d056..edb22d1d9 100644
--- a/lib/pleroma/web/fed_sockets/adapter.ex
+++ b/lib/pleroma/web/fed_sockets/adapter.ex
@@ -1,90 +1,90 @@
defmodule Pleroma.Web.FedSockets.Adapter do
@moduledoc """
A behavior both types of sockets (server and client) should implement
and a collection of helper functions useful to both
"""
@type adapter_state :: map()
@doc "A synchronous fetch."
@callback fetch(pid(), adapter_state(), term(), timeout()) :: {:ok, term()} | {:error, term()}
@doc "An asynchronous publish."
@callback publish(pid(), adapter_state(), term()) :: :ok | {:error, term()}
alias Pleroma.Object
alias Pleroma.Object.Containment
alias Pleroma.User
alias Pleroma.Web.ActivityPub.ObjectView
alias Pleroma.Web.ActivityPub.UserView
alias Pleroma.Web.ActivityPub.Visibility
alias Pleroma.Web.FedSockets.IngesterWorker
@type origin :: String.t()
@type fetch_id :: integer()
@type waiting_fetches :: %{required(fetch_id()) => pid()}
@doc "Processes incoming messages. Returns {:reply, websocket_frame, waiting_fetches} or `{:noreply, waiting_fetches}`"
@spec process_message(binary() | map(), origin(), waiting_fetches()) ::
{:reply, term(), waiting_fetches()} | {:noreply, waiting_fetches()}
def process_message(message, origin, waiting_fetches) when is_binary(message) do
case Jason.decode(message) do
{:ok, message} -> process_message(message, origin, waiting_fetches)
# 1003 indicates that an endpoint is terminating the connection
# because it has received a type of data it cannot accept.
{:error, decode_error} -> {:reply, {:close, 1003, Exception.message(decode_error)}}
end
end
def process_message(%{"action" => "publish", "data" => data}, origin, waiting_fetches) do
if Containment.contain_origin(origin, data) do
IngesterWorker.enqueue("ingest", %{"object" => data})
end
{:noreply, waiting_fetches}
end
- def process_message(%{"action" => "fetch", "uuid" => uuid, "data" => ap_id}, _, _) do
+ def process_message(%{"action" => "fetch", "uuid" => uuid, "data" => ap_id}, _, waiting_fetches) do
data = %{
"action" => "fetch_reply",
"status" => "processed",
"uuid" => uuid,
"data" => represent_item(ap_id)
}
- {:reply, {:text, Jason.encode!(data)}}
+ {:reply, {:text, Jason.encode!(data)}, waiting_fetches}
end
def process_message(
%{"action" => "fetch_reply", "uuid" => uuid, "data" => data},
_,
waiting_fetches
) do
with {pid, waiting_fetches} when is_pid(pid) <- Map.pop(waiting_fetches, uuid) do
send(pid, {:fetch_reply, uuid, data})
{:noreply, waiting_fetches}
else
_ ->
{:noreply, waiting_fetches}
end
end
def process_message(_, _, waiting_fetches) do
{:reply, {:close, 1003, "Unknown message type."}, waiting_fetches}
end
defp represent_item(ap_id) do
case User.get_by_ap_id(ap_id) do
nil ->
object = Object.get_cached_by_ap_id(ap_id)
if Visibility.is_public?(object) do
Phoenix.View.render(ObjectView, "object.json", object: object)
else
nil
end
user ->
Phoenix.View.render(UserView, "user.json", user: user)
end
end
end
diff --git a/lib/pleroma/web/fed_sockets/adapter/cowboy.ex b/lib/pleroma/web/fed_sockets/adapter/cowboy.ex
index 880af929c..107d80849 100644
--- a/lib/pleroma/web/fed_sockets/adapter/cowboy.ex
+++ b/lib/pleroma/web/fed_sockets/adapter/cowboy.ex
@@ -1,144 +1,145 @@
# Pleroma: A lightweight social networking server
# Copyright © 2017-2020 Pleroma Authors <https://pleroma.social/>
# SPDX-License-Identifier: AGPL-3.0-only
defmodule Pleroma.Web.FedSockets.Adapter.Cowboy do
require Logger
alias Pleroma.Web.FedSockets.Adapter
alias Pleroma.Web.FedSockets.FedSocket
alias Pleroma.Web.FedSockets.Registry.Value
import HTTPSignatures, only: [validate_conn: 1, split_signature: 1]
require Logger
@behaviour :cowboy_websocket
@behaviour Adapter
@impl true
def fetch(pid, %{last_fetch_id_ref: last_fetch_id_ref}, id, timeout) do
fetch_id = :atomics.add_get(last_fetch_id_ref, 1, 1)
message = %{action: :fetch, data: id, uuid: fetch_id}
send(pid, {:send_fetch, Jason.encode!(message), fetch_id, self()})
receive do
{:fetch_reply, ^fetch_id, data} -> {:ok, data}
after
timeout -> {:error, :timeout}
end
end
@impl true
def publish(pid, _, data) do
message = %{action: :publish, data: data}
send(pid, {:send, Jason.encode!(message)})
:ok
end
@impl true
def init(req, state) do
shake = FedSocket.shake()
with {_, true} <- {:enabled, Pleroma.Config.get([:fed_sockets, :enabled])},
sec_protocol <- :cowboy_req.header("sec-websocket-protocol", req, nil),
{_, %{"(request-target)" => ^shake} = headers} <-
{:has_request_target, :cowboy_req.headers(req)},
{_, true} <- {:signature_validated, validate_conn(%{req_headers: headers})},
%{"keyId" => origin} <- split_signature(headers["signature"]) do
req =
if is_nil(sec_protocol) do
req
else
:cowboy_req.set_resp_header("sec-websocket-protocol", sec_protocol, req)
end
{:cowboy_websocket, req, origin, %{}}
else
{:has_request_target, headers} ->
Logger.debug(fn ->
"#{__MODULE__}: Wrong or no \"(request-target)\" header. Rejecting websocket switch. Headers:\n#{
inspect(headers)
}"
end)
:cowboy_req.reply(400, req)
{:ok, req, state}
{:signature_validated, false} ->
Logger.debug(fn ->
"#{__MODULE__}: Signature validation failed. Rejecting websocket switch."
end)
:cowboy_req.reply(401, req)
{:ok, req, state}
e ->
Logger.debug(fn -> "#{__MODULE__}: Websocket switch failed, #{inspect(e)}" end)
:cowboy_req.reply(500, req)
{:ok, req, state}
end
end
@registry Pleroma.Web.FedSockets.Registry
@impl true
def websocket_init(origin) do
key = Pleroma.Web.FedSockets.Registry.key_from_uri(URI.parse(origin))
# Since, unlike with gun, we don't have calls.
# We store last fetch id in an atomic counter and use casts.
last_fetch_id_ref = :atomics.new(1, [])
:ok = :atomics.put(last_fetch_id_ref, 1, 0)
case Registry.register(@registry, key, %Value{
adapter: __MODULE__,
adapter_state: %{last_fetch_id_ref: last_fetch_id_ref, waiting_fetches: %{}}
}) do
{:ok, _owner} ->
- {:ok, %{origin: origin}}
+ {:ok, %{origin: origin, waiting_fetches: %{}}}
{:error, {:already_registered, _}} ->
{:stop, origin}
end
end
@impl true
- def websocket_handle(:ping, socket_info), do: {:ok, socket_info}
+ def websocket_handle(:ping, state), do: {:ok, state}
def websocket_handle(
{:text, raw_message},
%{origin: origin, waiting_fetches: waiting_fetches} = state
) do
case Adapter.process_message(raw_message, origin, waiting_fetches) do
{:reply, frame, waiting_fetches} ->
- {:reply, frame, %{state | waiting_fetches: waiting_fetches}}
+ state = %{state | waiting_fetches: waiting_fetches}
+ {:reply, frame, state}
{:noreply, waiting_fetches} ->
{:ok, %{state | waiting_fetches: waiting_fetches}}
end
end
@impl true
def websocket_info(
{:send_fetch, message, fetch_id, pid},
%{waiting_fetches: waiting_fetches} = state
) do
waiting_fetches = Map.put(waiting_fetches, fetch_id, pid)
{:reply, {:text, message}, %{state | waiting_fetches: waiting_fetches}}
end
@impl true
def websocket_info({:send, message}, state) do
{:reply, {:text, message}, state}
end
@impl true
def websocket_info(:close, state) do
{:stop, state}
end
def websocket_info(message, state) do
Logger.debug("#{__MODULE__} unknown message #{inspect(message)}")
{:ok, state}
end
end

File Metadata

Mime Type
text/x-diff
Expires
Sat, Aug 8, 1:25 PM (1 d, 14 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1723506
Default Alt Text
(8 KB)

Event Timeline