Page MenuHomePhorge

No OneTemporary

Size
14 KB
Referenced Files
None
Subscribers
None
diff --git a/lib/pleroma/federation_failure.ex b/lib/pleroma/federation_failure.ex
index bda1dfbab..e7c7f0e24 100644
--- a/lib/pleroma/federation_failure.ex
+++ b/lib/pleroma/federation_failure.ex
@@ -1,60 +1,59 @@
# Pleroma: A lightweight social networking server
# Copyright © 2017-2019 Pleroma Authors <https://pleroma.social/>
# SPDX-License-Identifier: AGPL-3.0-only
defmodule Pleroma.FederationFailure do
use Ecto.Schema
alias Pleroma.Activity
alias Pleroma.FederationFailure
alias Pleroma.FlakeId
alias Pleroma.Repo
import Ecto.Changeset
import Ecto.Query
@type t :: %__MODULE__{}
schema "federation_failures" do
field(:recipient, :string)
field(:transport, :string)
- field(:data, :map)
field(:retries_count, :integer, default: 0)
belongs_to(:activity, Activity, type: FlakeId)
timestamps()
end
def changeset(federation_failure, params \\ %{}) do
federation_failure
- |> cast(params, [:activity_id, :recipient, :transport, :data, :retries_count])
- |> validate_required([:activity_id, :recipient, :transport, :data])
+ |> cast(params, [:activity_id, :recipient, :transport, :retries_count])
+ |> validate_required([:activity_id, :recipient, :transport])
|> foreign_key_constraint(:activity_id)
|> unique_constraint(:activity_id,
name: :federation_failures_activity_id_recipient_transport_index
)
end
def create_or_rewrite_with(%{activity_id: nil}), do: :noop
def create_or_rewrite_with(%{} = params) do
%FederationFailure{}
|> changeset(params)
|> Repo.insert(
on_conflict: :replace_all_except_primary_key,
conflict_target: [:activity_id, :recipient, :transport]
)
end
def delete_by(%{activity_id: nil}), do: :noop
def delete_by(%{activity_id: activity_id, recipient: recipient, transport: transport}) do
from(ff in FederationFailure,
where:
ff.activity_id == ^activity_id and ff.recipient == ^recipient and
ff.transport == ^transport
)
|> Repo.delete_all()
end
end
diff --git a/lib/pleroma/web/federator/publisher.ex b/lib/pleroma/web/federator/publisher.ex
index 6f0a03814..071f3bc61 100644
--- a/lib/pleroma/web/federator/publisher.ex
+++ b/lib/pleroma/web/federator/publisher.ex
@@ -1,111 +1,108 @@
# Pleroma: A lightweight social networking server
# Copyright © 2017-2019 Pleroma Authors <https://pleroma.social/>
# SPDX-License-Identifier: AGPL-3.0-only
defmodule Pleroma.Web.Federator.Publisher do
alias Pleroma.Activity
alias Pleroma.Config
alias Pleroma.FederationFailure
alias Pleroma.User
alias Pleroma.Web.Federator.RetryQueue
require Logger
@moduledoc """
Defines the contract used by federation implementations to publish messages to
their peers.
"""
@doc """
Determine whether an activity can be relayed using the federation module.
"""
@callback is_representable?(Pleroma.Activity.t()) :: boolean()
@doc """
Relays an activity to a specified peer, determined by the parameters. The
parameters used are controlled by the federation module.
"""
@callback publish_one(Map.t()) :: {:ok, Map.t()} | {:error, any()}
@doc """
Enqueue publishing a single activity.
"""
@spec enqueue_one(module(), Map.t()) :: :ok
def enqueue_one(transport, %{job: %{activity_id: _, recipient: _}} = params) do
PleromaJobQueue.enqueue(:federator_outgoing, __MODULE__, [:publish_one, transport, params])
end
@spec perform(atom(), module(), any()) :: {:ok, any()} | {:error, any()}
def perform(
:publish_one,
transport,
%{job: %{activity_id: _, recipient: _}} = params
) do
federation_failure_id_data = Map.put(params[:job], :transport, to_string(transport))
case apply(transport, :publish_one, [params]) do
{:ok, _} ->
FederationFailure.delete_by(federation_failure_id_data)
:ok
{:error, _e} ->
federation_failure_id_data
- # Note: `params` must be JSON-serializable (it's not for Salmon, see salmon_test.exs)
- # |> Map.put(:data, params)
- |> Map.put(:data, %{})
|> Map.put(:retries_count, 0)
|> FederationFailure.create_or_rewrite_with()
RetryQueue.enqueue(params, transport)
end
end
def perform(type, _, _) do
Logger.debug("Unknown task: #{type}")
{:error, "Don't know what to do with this"}
end
@doc """
Relays an activity to all specified peers.
"""
@callback publish(User.t(), Activity.t()) :: :ok | {:error, any()}
@spec publish(User.t(), Activity.t()) :: :ok
def publish(%User{} = user, %Activity{} = activity) do
Config.get([:instance, :federation_publisher_modules])
|> Enum.each(fn module ->
if module.is_representable?(activity) do
Logger.info("Publishing #{activity.data["id"]} using #{inspect(module)}")
module.publish(user, activity)
end
end)
:ok
end
@doc """
Gathers links used by an outgoing federation module for WebFinger output.
"""
@callback gather_webfinger_links(User.t()) :: list()
@spec gather_webfinger_links(User.t()) :: list()
def gather_webfinger_links(%User{} = user) do
Config.get([:instance, :federation_publisher_modules])
|> Enum.reduce([], fn module, links ->
links ++ module.gather_webfinger_links(user)
end)
end
@doc """
Gathers nodeinfo protocol names supported by the federation module.
"""
@callback gather_nodeinfo_protocol_names() :: list()
@spec gather_nodeinfo_protocol_names() :: list()
def gather_nodeinfo_protocol_names do
Config.get([:instance, :federation_publisher_modules])
|> Enum.reduce([], fn module, links ->
links ++ module.gather_nodeinfo_protocol_names()
end)
end
end
diff --git a/lib/pleroma/web/federator/retry_queue.ex b/lib/pleroma/web/federator/retry_queue.ex
index cba7f72e7..c025a5eb3 100644
--- a/lib/pleroma/web/federator/retry_queue.ex
+++ b/lib/pleroma/web/federator/retry_queue.ex
@@ -1,258 +1,255 @@
# Pleroma: A lightweight social networking server
# Copyright © 2017-2019 Pleroma Authors <https://pleroma.social/>
# SPDX-License-Identifier: AGPL-3.0-only
defmodule Pleroma.Web.Federator.RetryQueue do
use GenServer
alias Pleroma.FederationFailure
require Logger
def init(args) do
queue_table = :ets.new(:pleroma_retry_queue, [:bag, :protected])
{:ok, %{args | queue_table: queue_table, running_jobs: :sets.new()}}
end
def start_link do
enabled =
if Pleroma.Config.get(:env) == :test,
do: true,
else: Pleroma.Config.get([__MODULE__, :enabled], false)
if enabled do
Logger.info("Starting retry queue")
linkres =
GenServer.start_link(
__MODULE__,
%{delivered: 0, dropped: 0, queue_table: nil, running_jobs: nil},
name: __MODULE__
)
maybe_kickoff_timer()
linkres
else
Logger.info("Retry queue disabled")
:ignore
end
end
def enqueue(data, transport, retries \\ 0) do
if job_params = job_params(data, transport) do
job_params
- # Note: `params` must be JSON-serializable (it's not for Salmon, see salmon_test.exs)
- # |> Map.put(:data, data)
- |> Map.put(:data, %{})
|> Map.put(:retries_count, retries + 1)
|> FederationFailure.create_or_rewrite_with()
end
GenServer.cast(__MODULE__, {:maybe_enqueue, data, transport, retries + 1})
end
def get_stats do
GenServer.call(__MODULE__, :get_stats)
end
def reset_stats do
GenServer.call(__MODULE__, :reset_stats)
end
def get_retry_params(retries) do
if retries > Pleroma.Config.get([__MODULE__, :max_retries]) do
{:drop, "Max retries reached"}
else
{:retry, growth_function(retries)}
end
end
def get_retry_timer_interval do
Pleroma.Config.get([:retry_queue, :interval], 1000)
end
defp ets_count_expires(table, current_time) do
:ets.select_count(
table,
[
{
{:"$1", :"$2"},
[{:"=<", :"$1", {:const, current_time}}],
[true]
}
]
)
end
defp ets_pop_n_expired(table, current_time, desired) do
{popped, _continuation} =
:ets.select(
table,
[
{
{:"$1", :"$2"},
[{:"=<", :"$1", {:const, current_time}}],
[:"$_"]
}
],
desired
)
popped
|> Enum.each(fn e ->
:ets.delete_object(table, e)
end)
popped
end
def maybe_start_job(running_jobs, queue_table) do
# we don't want to hit the ets or the DateTime more times than we have to
# could optimize slightly further by not using the count, and instead grabbing
# up to N objects early...
current_time = DateTime.to_unix(DateTime.utc_now())
n_running_jobs = :sets.size(running_jobs)
if n_running_jobs < Pleroma.Config.get([__MODULE__, :max_jobs]) do
n_ready_jobs = ets_count_expires(queue_table, current_time)
if n_ready_jobs > 0 do
# figure out how many we could start
available_job_slots = Pleroma.Config.get([__MODULE__, :max_jobs]) - n_running_jobs
start_n_jobs(running_jobs, queue_table, current_time, available_job_slots)
else
running_jobs
end
else
running_jobs
end
end
defp start_n_jobs(running_jobs, _queue_table, _current_time, 0) do
running_jobs
end
defp start_n_jobs(running_jobs, queue_table, current_time, available_job_slots)
when available_job_slots > 0 do
candidates = ets_pop_n_expired(queue_table, current_time, available_job_slots)
candidates
|> List.foldl(running_jobs, fn {_, e}, rj ->
{:ok, pid} = Task.start(fn -> worker(e) end)
mref = Process.monitor(pid)
:sets.add_element(mref, rj)
end)
end
def worker({:send, data, transport, retries}) do
case transport.publish_one(data) do
{:ok, _} ->
if job_params = job_params(data, transport), do: FederationFailure.delete_by(job_params)
GenServer.cast(__MODULE__, :inc_delivered)
:delivered
{:error, _reason} ->
enqueue(data, transport, retries)
:retry
end
end
def handle_call(:get_stats, _from, %{delivered: delivery_count, dropped: drop_count} = state) do
{:reply, %{delivered: delivery_count, dropped: drop_count}, state}
end
def handle_call(:reset_stats, _from, %{delivered: delivery_count, dropped: drop_count} = state) do
{:reply, %{delivered: delivery_count, dropped: drop_count},
%{state | delivered: 0, dropped: 0}}
end
def handle_cast(:reset_stats, state) do
{:noreply, %{state | delivered: 0, dropped: 0}}
end
def handle_cast(
{:maybe_enqueue, data, transport, retries},
%{dropped: drop_count, queue_table: queue_table, running_jobs: running_jobs} = state
) do
case get_retry_params(retries) do
{:retry, timeout} ->
:ets.insert(queue_table, {timeout, {:send, data, transport, retries}})
running_jobs = maybe_start_job(running_jobs, queue_table)
{:noreply, %{state | running_jobs: running_jobs}}
{:drop, message} ->
Logger.debug(message)
{:noreply, %{state | dropped: drop_count + 1}}
end
end
def handle_cast(:kickoff_timer, state) do
retry_interval = get_retry_timer_interval()
Process.send_after(__MODULE__, :retry_timer_run, retry_interval)
{:noreply, state}
end
def handle_cast(:inc_delivered, %{delivered: delivery_count} = state) do
{:noreply, %{state | delivered: delivery_count + 1}}
end
def handle_cast(:inc_dropped, %{dropped: drop_count} = state) do
{:noreply, %{state | dropped: drop_count + 1}}
end
def handle_info({:send, data, transport, retries}, %{delivered: delivery_count} = state) do
case transport.publish_one(data) do
{:ok, _} ->
if job_params = job_params(data, transport), do: FederationFailure.delete_by(job_params)
{:noreply, %{state | delivered: delivery_count + 1}}
{:error, _reason} ->
enqueue(data, transport, retries)
{:noreply, state}
end
end
def handle_info(
:retry_timer_run,
%{queue_table: queue_table, running_jobs: running_jobs} = state
) do
maybe_kickoff_timer()
running_jobs = maybe_start_job(running_jobs, queue_table)
{:noreply, %{state | running_jobs: running_jobs}}
end
def handle_info({:DOWN, ref, :process, _pid, _reason}, state) do
%{running_jobs: running_jobs, queue_table: queue_table} = state
running_jobs = :sets.del_element(ref, running_jobs)
running_jobs = maybe_start_job(running_jobs, queue_table)
{:noreply, %{state | running_jobs: running_jobs}}
end
def handle_info(unknown, state) do
Logger.debug("RetryQueue: don't know what to do with #{inspect(unknown)}, ignoring")
{:noreply, state}
end
if Pleroma.Config.get(:env) == :test do
defp growth_function(_retries) do
_shutit = Pleroma.Config.get([__MODULE__, :initial_timeout])
DateTime.to_unix(DateTime.utc_now()) - 1
end
else
defp growth_function(retries) do
round(Pleroma.Config.get([__MODULE__, :initial_timeout]) * :math.pow(retries, 3)) +
DateTime.to_unix(DateTime.utc_now())
end
end
defp maybe_kickoff_timer do
GenServer.cast(__MODULE__, :kickoff_timer)
end
defp job_params(%{job: job_data}, transport) do
Map.put(job_data, :transport, to_string(transport))
end
defp job_params(_, _), do: nil
end
diff --git a/priv/repo/migrations/20190719163442_create_federation_failures.exs b/priv/repo/migrations/20190719163442_create_federation_failures.exs
index 6143dd58a..40616bc2c 100644
--- a/priv/repo/migrations/20190719163442_create_federation_failures.exs
+++ b/priv/repo/migrations/20190719163442_create_federation_failures.exs
@@ -1,17 +1,16 @@
defmodule Pleroma.Repo.Migrations.CreateFederationFailures do
use Ecto.Migration
def change do
create_if_not_exists table(:federation_failures) do
add(:activity_id, references(:activities, type: :uuid, on_delete: :delete_all))
add(:recipient, :string, null: false)
add(:transport, :string, null: false)
- add(:data, :map)
add(:retries_count, :integer, default: 0)
timestamps()
end
create_if_not_exists(unique_index(:federation_failures, [:activity_id, :recipient, :transport]))
end
end

File Metadata

Mime Type
text/x-diff
Expires
Sat, Oct 10, 5:40 PM (1 d, 9 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1785375
Default Alt Text
(14 KB)

Event Timeline