Page Menu
Home
Phorge
Search
Configure Global Search
Log In
Files
F85628683
No One
Temporary
Actions
View File
Edit File
Delete File
View Transforms
Subscribe
Award Token
Flag For Later
Size
8 KB
Referenced Files
None
Subscribers
None
View Options
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
Details
Attached
Mime Type
text/x-diff
Expires
Sat, Aug 8, 1:25 PM (1 d, 13 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1723506
Default Alt Text
(8 KB)
Attached To
Mode
rPUBE pleroma-upstream
Attached
Detach File
Event Timeline
Log In to Comment