Page Menu
Home
Phorge
Search
Configure Global Search
Log In
Files
F85650411
No One
Temporary
Actions
View File
Edit File
Delete File
View Transforms
Subscribe
Award Token
Flag For Later
Size
30 KB
Referenced Files
None
Subscribers
None
View Options
diff --git a/lib/pleroma/object/fetcher.ex b/lib/pleroma/object/fetcher.ex
index 169298b34..eb9a4c478 100644
--- a/lib/pleroma/object/fetcher.ex
+++ b/lib/pleroma/object/fetcher.ex
@@ -1,251 +1,250 @@
# Pleroma: A lightweight social networking server
# Copyright © 2017-2020 Pleroma Authors <https://pleroma.social/>
# SPDX-License-Identifier: AGPL-3.0-only
defmodule Pleroma.Object.Fetcher do
alias Pleroma.HTTP
alias Pleroma.Object
alias Pleroma.Object.Containment
alias Pleroma.Repo
alias Pleroma.Signature
alias Pleroma.Web.ActivityPub.InternalFetchActor
alias Pleroma.Web.ActivityPub.ObjectValidator
alias Pleroma.Web.ActivityPub.Transmogrifier
alias Pleroma.Web.Federator
alias Pleroma.Web.FedSockets
require Logger
require Pleroma.Constants
defp touch_changeset(changeset) do
updated_at =
NaiveDateTime.utc_now()
|> NaiveDateTime.truncate(:second)
Ecto.Changeset.put_change(changeset, :updated_at, updated_at)
end
defp maybe_reinject_internal_fields(%{data: %{} = old_data}, new_data) do
internal_fields = Map.take(old_data, Pleroma.Constants.object_internal_fields())
Map.merge(new_data, internal_fields)
end
defp maybe_reinject_internal_fields(_, new_data), do: new_data
@spec reinject_object(struct(), map()) :: {:ok, Object.t()} | {:error, any()}
defp reinject_object(%Object{data: %{"type" => "Question"}} = object, new_data) do
Logger.debug("Reinjecting object #{new_data["id"]}")
with data <- maybe_reinject_internal_fields(object, new_data),
{:ok, data, _} <- ObjectValidator.validate(data, %{}),
changeset <- Object.change(object, %{data: data}),
changeset <- touch_changeset(changeset),
{:ok, object} <- Repo.insert_or_update(changeset),
{:ok, object} <- Object.set_cache(object) do
{:ok, object}
else
e ->
Logger.error("Error while processing object: #{inspect(e)}")
{:error, e}
end
end
defp reinject_object(%Object{} = object, new_data) do
Logger.debug("Reinjecting object #{new_data["id"]}")
with new_data <- Transmogrifier.fix_object(new_data),
data <- maybe_reinject_internal_fields(object, new_data),
changeset <- Object.change(object, %{data: data}),
changeset <- touch_changeset(changeset),
{:ok, object} <- Repo.insert_or_update(changeset),
{:ok, object} <- Object.set_cache(object) do
{:ok, object}
else
e ->
Logger.error("Error while processing object: #{inspect(e)}")
{:error, e}
end
end
def refetch_object(%Object{data: %{"id" => id}} = object) do
with {:local, false} <- {:local, Object.local?(object)},
{:ok, new_data} <- fetch_and_contain_remote_object_from_id(id),
{:ok, object} <- reinject_object(object, new_data) do
{:ok, object}
else
{:local, true} -> {:ok, object}
e -> {:error, e}
end
end
# Note: will create a Create activity, which we need internally at the moment.
def fetch_object_from_id(id, options \\ []) do
with {_, nil} <- {:fetch_object, Object.get_cached_by_ap_id(id)},
{_, true} <- {:allowed_depth, Federator.allowed_thread_distance?(options[:depth])},
{_, {:ok, data}} <- {:fetch, fetch_and_contain_remote_object_from_id(id)},
{_, nil} <- {:normalize, Object.normalize(data, false)},
params <- prepare_activity_params(data),
{_, :ok} <- {:containment, Containment.contain_origin(id, params)},
{_, {:ok, activity}} <-
{:transmogrifier, Transmogrifier.handle_incoming(params, options)},
{_, _data, %Object{} = object} <-
{:object, data, Object.normalize(activity, false)} do
{:ok, object}
else
{:allowed_depth, false} ->
{:error, "Max thread distance exceeded."}
{:containment, _} ->
{:error, "Object containment failed."}
{:transmogrifier, {:error, {:reject, e}}} ->
{:reject, e}
{:transmogrifier, _} = e ->
{:error, e}
{:object, data, nil} ->
reinject_object(%Object{}, data)
{:normalize, object = %Object{}} ->
{:ok, object}
{:fetch_object, %Object{} = object} ->
{:ok, object}
{:fetch, {:error, error}} ->
{:error, error}
e ->
e
end
end
defp prepare_activity_params(data) do
%{
"type" => "Create",
"to" => data["to"] || [],
"cc" => data["cc"] || [],
# Should we seriously keep this attributedTo thing?
"actor" => data["actor"] || data["attributedTo"],
"object" => data
}
end
def fetch_object_from_id!(id, options \\ []) do
with {:ok, object} <- fetch_object_from_id(id, options) do
object
else
{:error, %Tesla.Mock.Error{}} ->
nil
{:error, "Object has been deleted"} ->
nil
{:reject, reason} ->
Logger.info("Rejected #{id} while fetching: #{inspect(reason)}")
nil
e ->
Logger.error("Error while fetching #{id}: #{inspect(e)}")
nil
end
end
defp make_signature(id, date) do
uri = URI.parse(id)
signature =
InternalFetchActor.get_actor()
|> Signature.sign(%{
"(request-target)": "get #{uri.path}",
host: uri.host,
date: date
})
{"signature", signature}
end
defp sign_fetch(headers, id, date) do
if Pleroma.Config.get([:activitypub, :sign_object_fetches]) do
[make_signature(id, date) | headers]
else
headers
end
end
defp maybe_date_fetch(headers, date) do
if Pleroma.Config.get([:activitypub, :sign_object_fetches]) do
[{"date", date} | headers]
else
headers
end
end
def fetch_and_contain_remote_object_from_id(prm, opts \\ [])
def fetch_and_contain_remote_object_from_id(%{"id" => id}, opts),
do: fetch_and_contain_remote_object_from_id(id, opts)
def fetch_and_contain_remote_object_from_id(id, opts) when is_binary(id) do
Logger.debug("Fetching object #{id} via AP")
with {:scheme, true} <- {:scheme, String.starts_with?(id, "http")},
{:ok, body} <- get_object(id, opts),
{:ok, data} <- safe_json_decode(body),
:ok <- Containment.contain_origin_from_id(id, data) do
{:ok, data}
else
{:scheme, _} ->
{:error, "Unsupported URI scheme"}
{:error, e} ->
{:error, e}
e ->
{:error, e}
end
end
def fetch_and_contain_remote_object_from_id(_id, _opts),
do: {:error, "id must be a string"}
defp get_object(id, opts) do
- with false <- Keyword.get(opts, :force_http, false),
- {:ok, fedsocket} <- FedSockets.get_or_create_fed_socket(id) do
+ with false <- Keyword.get(opts, :force_http, false) do
Logger.debug("fetching via fedsocket - #{inspect(id)}")
- FedSockets.fetch(fedsocket, id)
+ FedSockets.Registry.fetch(id)
else
_other ->
Logger.debug("fetching via http - #{inspect(id)}")
get_object_http(id)
end
end
defp get_object_http(id) do
date = Pleroma.Signature.signed_date()
headers =
[{"accept", "application/activity+json"}]
|> maybe_date_fetch(date)
|> sign_fetch(id, date)
case HTTP.get(id, headers) do
{:ok, %{body: body, status: code}} when code in 200..299 ->
{:ok, body}
{:ok, %{status: code}} when code in [404, 410] ->
{:error, "Object has been deleted"}
{:error, e} ->
{:error, e}
e ->
{:error, e}
end
end
defp safe_json_decode(nil), do: {:ok, nil}
defp safe_json_decode(json), do: Jason.decode(json)
end
diff --git a/lib/pleroma/web/fed_sockets/adapter.ex b/lib/pleroma/web/fed_sockets/adapter.ex
new file mode 100644
index 000000000..5bd9caec9
--- /dev/null
+++ b/lib/pleroma/web/fed_sockets/adapter.ex
@@ -0,0 +1,69 @@
+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
+
+ @typedoc """
+ Should be "fetch" or "publish"
+ """
+ @type common_action :: String.t()
+ @type origin :: String.t()
+ @doc "Process non adapter-specific messages."
+ @spec process_message(map(), origin()) ::
+ {:reply, term()} | :noreply | {:error, :unknown_action}
+ def process_message(%{"action" => "publish", "data" => data}, origin) do
+ if Containment.contain_origin(origin, data) do
+ IngesterWorker.enqueue("ingest", %{"object" => data})
+ end
+
+ :noreply
+ end
+
+ def process_message(%{"action" => "fetch", "uuid" => uuid, "data" => ap_id}, _) do
+ data = %{
+ "action" => "fetch_reply",
+ "status" => "processed",
+ "uuid" => uuid,
+ "data" => represent_item(ap_id)
+ }
+
+ {:reply, data}
+ end
+
+ def process_message(_, _) do
+ {:error, :unknown_action}
+ 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
new file mode 100644
index 000000000..b4f785c57
--- /dev/null
+++ b/lib/pleroma/web/fed_sockets/adapter/cowboy.ex
@@ -0,0 +1,150 @@
+# 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(_socket, _data), do: {:error, :not_implemented}
+
+ @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}}
+
+ {:error, {:already_registered, _}} ->
+ {:stop, origin}
+ end
+ end
+
+ @impl true
+ def websocket_handle(:ping, socket_info), do: {:ok, socket_info}
+
+ def websocket_handle({:text, raw_message}, %{origin: origin} = state) do
+ case Jason.decode(raw_message) do
+ {:ok, message} ->
+ case message do
+ %{"action" => "fetch_reply", "uuid" => uuid, "data" => data} ->
+ with {pid, waiting_fetches} when is_pid(pid) <- Map.pop(state.waiting_fetches, uuid) do
+ send(pid, {:fetch_reply, uuid, data})
+ {:ok, %{state | waiting_fetches: waiting_fetches}}
+ else
+ _ ->
+ {:ok, state}
+ end
+
+ message ->
+ case Adapter.process_message(message, origin) do
+ :noreply -> {:ok, state}
+ {:reply, data} -> {:reply, Jason.encode!(data), state}
+ end
+ end
+
+ {:error, decode_error} ->
+ exit({:malformed_message, decode_error})
+ 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
diff --git a/lib/pleroma/web/fed_sockets/socket/client.ex b/lib/pleroma/web/fed_sockets/adapter/gun.ex
similarity index 63%
rename from lib/pleroma/web/fed_sockets/socket/client.ex
rename to lib/pleroma/web/fed_sockets/adapter/gun.ex
index 480cc2274..9a93d3b83 100644
--- a/lib/pleroma/web/fed_sockets/socket/client.ex
+++ b/lib/pleroma/web/fed_sockets/adapter/gun.ex
@@ -1,214 +1,240 @@
# 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.Socket.Client do
- use GenServer
+defmodule Pleroma.Web.FedSockets.Adapter.Gun do
+ use GenServer, restart: :temporary
require Logger
alias Pleroma.Web.ActivityPub.InternalFetchActor
- alias Pleroma.Web.FedSockets
- alias Pleroma.Web.FedSockets.FedRegistry
+ alias Pleroma.Web.FedSockets.Registry.Value
alias Pleroma.Web.FedSockets.FedSocket
- alias Pleroma.Web.FedSockets.Socket
- alias Pleroma.Web.FedSockets.SocketInfo
+ alias Pleroma.Web.FedSockets.Adapter
- @behaviour Socket
+ @behaviour Adapter
+
+ @registry Pleroma.Web.FedSockets.Registry
+
+ @impl true
+ def fetch(pid, state, id, timeout) do
+ # TODO: refactor with atomics and encoding on the client
+ with {:ok, _} <- await_connected(pid, state) do
+ GenServer.call(pid, {:fetch, id}, timeout)
+ end
+ end
+
+ @impl true
+ def publish(pid, state, data) do
+ with {:ok, conn_pid} <- await_connected(pid, state) do
+ send_json(conn_pid, %{action: :publish, data: data})
+ end
+ end
+
+ defp await_connected(_pid, %{conn_pid: conn_pid}), do: {:ok, conn_pid}
+
+ defp await_connected(pid, _) do
+ monitor = Process.monitor(pid)
+ GenServer.cast(pid, {:await_connected, self()})
+
+ receive do
+ {:DOWN, ^monitor, _, _, {:shutdown, reason}} -> reason
+ {:DOWN, ^monitor, _, _, reason} -> {:error, reason}
+ {:await_connected, ^pid, conn_pid} -> {:ok, conn_pid}
+ end
+ end
+
+ def start_link([key | _] = opts) do
+ GenServer.start_link(__MODULE__, opts,
+ name: {:via, Registry, {@registry, key, %Value{adapter: __MODULE__, adapter_state: %{}}}}
+ )
+ end
@impl true
- def fetch(%{pid: handler_pid}, data, timeout) do
- GenServer.call(handler_pid, {:fetch, data}, timeout)
+ def init(opts) do
+ {:ok, nil, {:continue, {:connect, opts}}}
end
@impl true
- def publish(socket_info, data) do
- send_json(socket_info, %{action: :publish, data: data})
+ def handle_continue({:connect, [key, uri] = opts}, _) do
+ case initiate_connection(uri) do
+ {:ok, conn_pid} ->
+ Registry.update_value(@registry, key, fn value ->
+ %{value | adapter_state: %{conn_pid: conn_pid}}
+ end)
+
+ {:noreply,
+ %{conn_pid: conn_pid, waiting_fetches: %{}, last_fetch_id: 0, origin: uri, key: key}}
+
+ {:error, reason} = e ->
+ Logger.debug("Outgoing connection failed - #{inspect(reason)}")
+ {:stop, {:shutdown, e}, opts}
+ end
end
- def start_link(uri) do
- GenServer.start_link(__MODULE__, %{uri: uri})
+ @impl true
+ def handle_cast({:await_connected, pid}, %{conn_pid: conn_pid} = state) do
+ send(pid, {:await_connected, self(), conn_pid})
+ {:noreply, state}
end
@impl true
def handle_call(
{:fetch, data},
from,
%{
last_fetch_id: last_fetch_id,
- socket_info: socket_info,
+ conn_pid: conn_pid,
waiting_fetches: waiting_fetches
} = state
) do
last_fetch_id = last_fetch_id + 1
request = %{action: :fetch, data: data, uuid: last_fetch_id}
- socket_info = send_json(socket_info, request)
+ :ok = send_json(conn_pid, request)
waiting_fetches = Map.put(waiting_fetches, last_fetch_id, from)
{:noreply,
%{
state
- | socket_info: socket_info,
- waiting_fetches: waiting_fetches,
+ | waiting_fetches: waiting_fetches,
last_fetch_id: last_fetch_id
}}
end
- defp send_json(%{conn_pid: conn_pid} = socket_info, data) do
- socket_info = SocketInfo.touch(socket_info)
+ defp send_json(conn_pid, data) do
:gun.ws_send(conn_pid, {:text, Jason.encode!(data)})
- socket_info
- end
-
- @impl true
- def init(%{uri: uri}) do
- case initiate_connection(uri) do
- {:ok, ws_origin, conn_pid} ->
- {:ok, socket_info} = FedRegistry.add_fed_socket(ws_origin, conn_pid)
- {:ok, %{socket_info: socket_info, waiting_fetches: %{}, last_fetch_id: 0}}
-
- {:error, reason} ->
- Logger.debug("Outgoing connection failed - #{inspect(reason)}")
- :ignore
- end
end
@impl true
def handle_info(
{:gun_ws, _conn_pid, _ref, {:text, raw_message}},
- %{socket_info: socket_info} = state
+ %{conn_pid: conn_pid, origin: origin} = state
) do
- socket_info = SocketInfo.touch(socket_info)
-
state =
case Jason.decode(raw_message) do
{:ok, message} ->
case message do
%{"action" => "fetch_reply", "uuid" => uuid, "data" => data} ->
with {{_, _} = client, waiting_fetches} <- Map.pop(state.waiting_fetches, uuid) do
GenServer.reply(client, {:ok, data})
%{state | waiting_fetches: waiting_fetches}
else
_ ->
state
end
message ->
- case Socket.process_message(socket_info, message) do
+ case Socket.process_message(message, origin) do
:noreply -> :noop
- {:reply, data} -> send_json(socket_info, data)
+ {:reply, data} -> send_json(conn_pid, data)
end
state
end
- {:error, _e} ->
- state
+ {:error, decode_error} ->
+ exit({:malformed_message, decode_error})
end
- {:noreply, %{state | socket_info: socket_info}}
+ {:noreply, state}
end
@impl true
def handle_info(:close, state) do
Logger.debug("Sending close frame !!!!!!!")
{:close, state}
end
@impl true
def handle_info({:gun_down, _pid, _prot, :closed, _}, state) do
{:stop, :normal, state}
end
@impl true
def handle_info({:gun_ws, _, _, :pong}, state) do
{:noreply, state, :hibernate}
end
@impl true
def handle_info(msg, state) do
Logger.debug("#{__MODULE__} unhandled event #{inspect(msg)}")
{:noreply, state}
end
@impl true
def terminate(reason, state) do
Logger.debug(
"#{__MODULE__} terminating outgoing connection for #{inspect(state)} for #{inspect(reason)}"
)
{:ok, state}
end
+ @path '/api/fedsocket/v1'
def initiate_connection(uri) do
- ws_uri =
- uri
- |> SocketInfo.origin()
- |> FedSockets.uri_for_origin()
-
- %{host: host, port: port, path: path} = URI.parse(ws_uri)
+ %{host: host, port: port} = URI.parse(uri)
with {:ok, conn_pid} <- :gun.open(to_charlist(host), port, %{protocols: [:http]}),
{:ok, _} <- :gun.await_up(conn_pid),
# TODO: nodeinfo-based support detection
# reference <- :gun.get(conn_pid, to_charlist(path)),
# {:response, :fin, 204, _} <- :gun.await(conn_pid, reference) |> IO.inspect(),
# :ok <- :gun.flush(conn_pid),
headers <- build_headers(uri),
- ref <- :gun.ws_upgrade(conn_pid, to_charlist(path), headers, %{silence_pings: false}) do
+ ref <- :gun.ws_upgrade(conn_pid, @path, headers, %{silence_pings: false}) do
receive do
{:gun_upgrade, ^conn_pid, ^ref, [<<"websocket">>], _} ->
- {:ok, ws_uri, conn_pid}
+ {:ok, conn_pid}
# mes ->
# IO.inspect(mes)
after
15_000 ->
Logger.debug("Fedsocket timeout connecting to #{inspect(uri)}")
{:error, :timeout}
end
else
{:response, :nofin, 404, _} ->
{:error, :fedsockets_not_supported}
e ->
Logger.debug("Fedsocket error connecting to #{inspect(uri)}")
{:error, e}
end
end
defp build_headers(uri) do
host_for_sig = uri |> URI.parse() |> host_signature()
shake = FedSocket.shake()
digest = "SHA-256=" <> (:crypto.hash(:sha256, shake) |> Base.encode64())
date = Pleroma.Signature.signed_date()
shake_size = byte_size(shake)
signature_opts = %{
"(request-target)": shake,
"content-length": to_charlist("#{shake_size}"),
date: date,
digest: digest,
host: host_for_sig
}
signature = Pleroma.Signature.sign(InternalFetchActor.get_actor(), signature_opts)
[
{"signature", signature},
{"date", date},
{"digest", digest},
{"content-length", to_string(shake_size)},
{"(request-target)", shake}
]
end
defp host_signature(%{host: host, scheme: scheme, port: port}) do
if port == URI.default_port(scheme) do
host
else
"#{host}:#{port}"
end
end
end
diff --git a/lib/pleroma/web/fed_sockets/registry.ex b/lib/pleroma/web/fed_sockets/registry.ex
new file mode 100644
index 000000000..b8e254fd6
--- /dev/null
+++ b/lib/pleroma/web/fed_sockets/registry.ex
@@ -0,0 +1,72 @@
+defmodule Pleroma.Web.FedSockets.Registry.Value do
+ defstruct [:adapter, :adapter_state]
+end
+
+defmodule Pleroma.Web.FedSockets.Registry do
+ alias Pleroma.Web.FedSockets.Registry.Value
+
+ @registry __MODULE__
+
+ @spec fetch(binary()) :: {:ok, term()} | {:error, term()}
+ def fetch(object_id) do
+ case get_socket(object_id) do
+ {:ok, pid, %Value{adapter: adapter, adapter_state: adapter_state}} ->
+ apply(adapter, :fetch, [pid, adapter_state, object_id, 5_000])
+
+ e ->
+ e
+ end
+ end
+
+ @spec publish(binary(), term()) :: :ok | {:error, term()}
+ def publish(inbox, data) do
+ case get_socket(inbox) do
+ {:ok, pid, %Value{adapter: adapter, adapter_state: adapter_state}} ->
+ apply(adapter, :publish, [pid, adapter_state, data])
+
+ e ->
+ e
+ end
+ end
+
+ @doc "Get a registry key from a URI. For internal use by adapters."
+ @spec key_from_uri(URI.t()) :: String.t()
+ def key_from_uri(%URI{scheme: scheme, host: host, port: port}), do: "#{scheme}:#{host}:#{port}"
+
+ @client_adapter Pleroma.Web.FedSockets.Adapter.Gun
+ defp get_socket(%URI{} = uri) do
+ key = key_from_uri(uri)
+
+ case Registry.lookup(@registry, key) do
+ [] ->
+ case DynamicSupervisor.start_child(
+ Pleroma.Web.FedSockets.ClientSupervisor,
+ {@client_adapter, [key, uri]}
+ ) do
+ {:ok, pid} ->
+ {:ok, pid, %Value{adapter: @client_adapter, adapter_state: %{}}}
+
+ {:error, {:already_started, pid}} ->
+ {:ok, pid, %Value{adapter: @client_adapter, adapter_state: %{}}}
+
+ {:error, _} = e ->
+ e
+
+ e ->
+ {:error, e}
+ end
+
+ [{pid, value}] ->
+ {:ok, pid, value}
+ end
+ end
+
+ defp get_socket(uri) when is_binary(uri) do
+ case URI.parse(uri) do
+ %{host: nil} = uri -> {:error, {:invalid_uri, uri}}
+ %{port: nil} = uri -> {:error, {:invalid_uri, uri}}
+ %{scheme: nil} = uri -> {:error, {:invalid_uri, uri}}
+ uri -> get_socket(uri)
+ end
+ end
+end
diff --git a/lib/pleroma/web/fed_sockets/socket/server.ex b/lib/pleroma/web/fed_sockets/socket/server.ex
deleted file mode 100644
index 71b45ede0..000000000
--- a/lib/pleroma/web/fed_sockets/socket/server.ex
+++ /dev/null
@@ -1,97 +0,0 @@
-# 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.Socket.Server do
- require Logger
-
- alias Pleroma.Web.FedSockets.Socket
- alias Pleroma.Web.FedSockets.FedRegistry
- alias Pleroma.Web.FedSockets.FedSocket
- alias Pleroma.Web.FedSockets.SocketInfo
-
- import HTTPSignatures, only: [validate_conn: 1, split_signature: 1]
-
- require Logger
-
- @behaviour :cowboy_websocket
- @behaviour Socket
-
- @impl true
- def fetch(_socket, _data, _timeout), do: {:error, :not_implemented}
-
- @impl true
- def publish(_socket, _data), do: {:error, :not_implemented}
-
- @impl true
- def init(req, state) do
- shake = FedSocket.shake()
-
- with true <- Pleroma.Config.get([:fed_sockets, :enabled]),
- sec_protocol <- :cowboy_req.header("sec-websocket-protocol", req, nil),
- headers = %{"(request-target)" => ^shake} <- :cowboy_req.headers(req),
- true <- 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: origin}, %{}}
- else
- e ->
- Logger.debug(fn -> "#{__MODULE__}: Websocket switch failed, #{inspect(e)}" end)
- {:ok, req, state}
- end
- end
-
- @impl true
- def websocket_init(%{origin: origin}) do
- case FedRegistry.add_fed_socket(origin) do
- {:ok, socket_info} ->
- {:ok, socket_info}
-
- e ->
- Logger.error("FedSocket websocket_init failed - #{inspect(e)}")
- {:error, inspect(e)}
- end
- end
-
- @impl true
- def websocket_handle(:ping, socket_info), do: {:ok, socket_info}
-
- def websocket_handle({:text, raw_message}, socket_info) do
- socket_info = SocketInfo.touch(socket_info)
-
- case Jason.decode(raw_message) do
- {:ok, message} ->
- case message do
- message ->
- case Socket.process_message(socket_info, message) do
- :noreply -> {:ok, socket_info}
- {:reply, data} -> {:reply, Jason.encode!(data), socket_info}
- end
- end
-
- {:error, decode_error} ->
- exit({:malformed_message, decode_error})
- end
- end
-
- def websocket_info({:send, message}, socket_info) do
- socket_info = SocketInfo.touch(socket_info)
-
- {:reply, {:text, message}, socket_info}
- end
-
- 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
diff --git a/lib/pleroma/web/fed_sockets/supervisor.ex b/lib/pleroma/web/fed_sockets/supervisor.ex
index 721ac39c5..300bc1289 100644
--- a/lib/pleroma/web/fed_sockets/supervisor.ex
+++ b/lib/pleroma/web/fed_sockets/supervisor.ex
@@ -1,57 +1,58 @@
# 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.Supervisor do
use Supervisor
import Cachex.Spec
def start_link(opts) do
Supervisor.start_link(__MODULE__, opts, name: __MODULE__)
end
def init(args) do
children = [
build_cache(:fed_socket_rejections, args),
- {Registry, keys: :unique, name: FedSockets.Registry}
+ {Registry, keys: :unique, name: Pleroma.Web.FedSockets.Registry},
+ {DynamicSupervisor, name: Pleroma.Web.FedSockets.ClientSupervisor, strategy: :one_for_one}
]
opts = [strategy: :one_for_all, name: Pleroma.Web.Streamer.Supervisor]
Supervisor.init(children, opts)
end
defp build_cache(name, args) do
opts = get_opts(name, args)
%{
id: String.to_atom("#{name}_cache"),
start: {Cachex, :start_link, [name, opts]},
type: :worker
}
end
defp get_opts(:fed_socket_rejections = cache_name, args) do
default = get_opts_or_config(args, cache_name, :default, 15_000)
interval = get_opts_or_config(args, cache_name, :interval, 3_000)
lazy = get_opts_or_config(args, cache_name, :lazy, false)
[expiration: expiration(default: default, interval: interval, lazy: lazy)]
end
defp get_opts(name, args) do
Keyword.get(args, name, [])
end
defp get_opts_or_config(args, name, key, default) do
args
|> Keyword.get(name, [])
|> Keyword.get(key)
|> case do
nil ->
Pleroma.Config.get([:fed_sockets, name, key], default)
value ->
value
end
end
end
File Metadata
Details
Attached
Mime Type
text/x-diff
Expires
Sun, Aug 30, 2:04 AM (1 d, 14 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1737695
Default Alt Text
(30 KB)
Attached To
Mode
rPUBE pleroma-upstream
Attached
Detach File
Event Timeline
Log In to Comment