Page MenuHomePhorge

No OneTemporary

Size
8 KB
Referenced Files
None
Subscribers
None
diff --git a/lib/pleroma/migrators/media_repository_migrator.ex b/lib/pleroma/migrators/media_repository_migrator.ex
index ec7402fb3..3acc22145 100644
--- a/lib/pleroma/migrators/media_repository_migrator.ex
+++ b/lib/pleroma/migrators/media_repository_migrator.ex
@@ -1,152 +1,162 @@
# Pleroma: A lightweight social networking server
# Copyright © 2017-2023 Pleroma Authors <https://pleroma.social/>
# SPDX-License-Identifier: AGPL-3.0-only
defmodule Pleroma.Migrators.MediaRepositoryMigrator do
defmodule State do
use Pleroma.Migrators.Support.BaseMigratorState
@impl Pleroma.Migrators.Support.BaseMigratorState
defdelegate data_migration(), to: Pleroma.DataMigration, as: :add_media_repository_info
end
use Pleroma.Migrators.Support.BaseMigrator
alias Pleroma.Migrators.Support.BaseMigrator
alias Pleroma.Object
alias Pleroma.UploadedFile
@doc "This migration removes objects created exclusively for contexts, containing only an `id` field."
@impl BaseMigrator
def feature_config_path, do: [:features, :add_media_repository_info]
@impl BaseMigrator
def fault_rate_allowance, do: Config.get([:add_media_repository_info, :fault_rate_allowance], 0)
@impl BaseMigrator
def perform do
data_migration_id = data_migration_id()
max_processed_id = get_stat(:max_processed_id, 0)
Logger.info("Adding media info (from oid: #{max_processed_id})...")
query()
|> where([object], object.id > ^max_processed_id)
|> Repo.chunk_stream(100, :batches, timeout: :infinity)
|> Stream.each(fn objects ->
object_ids = Enum.map(objects, & &1.id)
results = Enum.map(objects, &add_media_info(&1))
failed_ids =
results
|> Enum.filter(&(elem(&1, 0) == :error))
|> Enum.map(&elem(&1, 1))
chunk_affected_count =
results
|> Enum.filter(&(elem(&1, 0) == :ok))
|> length()
for failed_id <- failed_ids do
_ =
Repo.query(
"INSERT INTO data_migration_failed_ids(data_migration_id, record_id) " <>
"VALUES ($1, $2) ON CONFLICT DO NOTHING;",
[data_migration_id, failed_id]
)
end
_ =
Repo.query(
"DELETE FROM data_migration_failed_ids " <>
"WHERE data_migration_id = $1 AND record_id = ANY($2)",
[data_migration_id, object_ids -- failed_ids]
)
max_object_id = Enum.at(object_ids, -1)
put_stat(:max_processed_id, max_object_id)
increment_stat(:iteration_processed_count, length(object_ids))
increment_stat(:processed_count, length(object_ids))
increment_stat(:failed_count, length(failed_ids))
increment_stat(:affected_count, chunk_affected_count)
put_stat(:records_per_second, records_per_second())
persist_state()
# A quick and dirty approach to controlling the load this background migration imposes
sleep_interval = Config.get([:delete_context_objects, :sleep_interval_ms], 15000)
Process.sleep(sleep_interval)
end)
|> Stream.run()
end
@impl BaseMigrator
def query do
from(
object in Object,
where: fragment("(?)->>'type' IN ('Document', 'Image')", object.data),
left_join: uploaded_file in UploadedFile,
on: object.id == uploaded_file.object_id,
where: is_nil(uploaded_file)
)
end
defp get_url_spec_from_object(object) do
url =
object.data["url"]
|> Enum.at(0)
|> Map.get("href")
- upload_base = Pleroma.Upload.base_url()
- if String.starts_with?(url, upload_base) do
- String.replace_prefix(url, upload_base, "")
+ base_urls = [
+ Pleroma.Upload.base_url()
+ | Config.get([:add_media_repository_info, :additional_base_urls], [])
+ ]
+
+ try_get_url_spec(url, base_urls)
+ end
+
+ defp try_get_url_spec(url, [base_url | next]) do
+ if String.starts_with?(url, base_url) do
+ String.replace_prefix(url, base_url, "")
else
- nil
+ try_get_url_spec(url, next)
end
end
+ defp try_get_url_spec(_url, []), do: nil
+
@spec add_media_info(Object.t()) :: {:ok | :error, integer()}
def add_media_info(object) do
with url_spec when not is_nil(url_spec) <- get_url_spec_from_object(object),
{:ok, _uploaded_file} <- UploadedFile.create(%{object: object, path: url_spec}) do
{:ok, object.id}
else
_ -> {:error, object.id}
end
end
@impl BaseMigrator
def retry_failed do
data_migration_id = data_migration_id()
failed_objects_query()
|> Repo.chunk_stream(100, :one)
|> Stream.each(fn object ->
with {res, _} when res != :error <- add_media_info(object.id) do
_ =
Repo.query(
"DELETE FROM data_migration_failed_ids " <>
"WHERE data_migration_id = $1 AND record_id = $2",
[data_migration_id, object.id]
)
end
end)
|> Stream.run()
put_stat(:failed_count, failures_count())
persist_state()
force_continue()
end
defp failed_objects_query do
from(o in Object)
|> join(:inner, [o], dmf in fragment("SELECT * FROM data_migration_failed_ids"),
on: dmf.record_id == o.id
)
|> where([_o, dmf], dmf.data_migration_id == ^data_migration_id())
|> order_by([o], asc: o.id)
end
end
diff --git a/test/pleroma/migrators/media_repository_migrator_test.exs b/test/pleroma/migrators/media_repository_migrator_test.exs
index 2a2cd03fe..d5d0ef45f 100644
--- a/test/pleroma/migrators/media_repository_migrator_test.exs
+++ b/test/pleroma/migrators/media_repository_migrator_test.exs
@@ -1,56 +1,90 @@
# Pleroma: A lightweight social networking server
# Copyright © 2017-2022 Pleroma Authors <https://pleroma.social/>
# SPDX-License-Identifier: AGPL-3.0-only
defmodule Pleroma.Migrators.MediaRepositoryMigratorTest do
- use Pleroma.DataCase
+ use Pleroma.DataCase, async: false
alias Pleroma.Migrators.MediaRepositoryMigrator
alias Pleroma.Repo
alias Pleroma.UploadedFile
import Pleroma.Factory
@url_prefix "#{Pleroma.Web.Endpoint.url()}/media/"
defp get_url_spec(url) do
String.replace_prefix(url, @url_prefix, "")
end
describe "query/0" do
test "it returns Documents and Images" do
document = insert(:attachment)
image = insert(:attachment, data: %{"type" => "Image"})
_note = insert(:note)
result = MediaRepositoryMigrator.query() |> Repo.all()
assert [_, _] = result
assert image in result
assert document in result
end
test "it does not return objects with a corresponding UploadedFile" do
document = insert(:attachment)
{:ok, _uploaded_file} = UploadedFile.create(%{object: document, path: "some/path"})
result = MediaRepositoryMigrator.query() |> Repo.all()
assert [] == result
end
end
describe "add_media_info/1" do
test "it adds an UploadedFile" do
document = insert(:attachment)
+
+ url_spec =
+ document.data["url"]
+ |> Enum.at(0)
+ |> Map.get("href")
+ |> get_url_spec()
+
+ {:ok, _} = MediaRepositoryMigrator.add_media_info(document)
+
+ uploaded_file = UploadedFile.get_by_object(document)
+ assert uploaded_file.path == url_spec
+ end
+
+ test "if the url spec does not start with base url, it should fail" do
+ clear_config([Pleroma.Upload, :base_url], "https://some.example/media/")
+
+ document = insert(:attachment)
+
+ {:error, _} = MediaRepositoryMigrator.add_media_info(document)
+
+ uploaded_file = UploadedFile.get_by_object(document)
+ assert uploaded_file == nil
+ end
+
+ test "it adds an UploadedFile for urls starting with one of the additional_base_urls" do
+ clear_config([Pleroma.Upload, :base_url], "https://some.example/media/")
+
+ clear_config(
+ [:add_media_repository_info, :additional_base_urls],
+ ["#{Pleroma.Web.Endpoint.url()}/media/"]
+ )
+
+ document = insert(:attachment)
+
url_spec =
document.data["url"]
|> Enum.at(0)
|> Map.get("href")
|> get_url_spec()
{:ok, _} = MediaRepositoryMigrator.add_media_info(document)
uploaded_file = UploadedFile.get_by_object(document)
assert uploaded_file.path == url_spec
end
end
end

File Metadata

Mime Type
text/x-diff
Expires
Sun, Aug 9, 4:38 AM (1 d, 16 h)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1724209
Default Alt Text
(8 KB)

Event Timeline