Page MenuHomePhorge

No OneTemporary

Size
16 KB
Referenced Files
None
Subscribers
None
diff --git a/lib/mix/migrator.ex b/lib/mix/migrator.ex
new file mode 100644
index 000000000..373d46ad9
--- /dev/null
+++ b/lib/mix/migrator.ex
@@ -0,0 +1,96 @@
+defmodule Mix.Migrator do
+ import Mix.Pleroma
+ alias Pleroma.Activity
+ alias Pleroma.Object
+ alias Pleroma.Repo
+ alias Pleroma.Web.ActivityPub.Transmogrifier
+
+ @doc "Common functions to be reused in migrator tasks"
+ def keys_to_atoms(map) do
+ Map.new(map, fn {k, v} -> {String.to_atom(k), v} end)
+ end
+
+ def loop_fields(map, fields, fun) do
+ Enum.reduce(fields, map, fn key, acc ->
+ if Map.has_key?(acc, key) do
+ Map.put(acc, key, fun.(acc[key]))
+ else
+ acc
+ end
+ end)
+ end
+
+ def parse_timestamp_usec(timestamp) do
+ case NaiveDateTime.from_iso8601(timestamp) do
+ {:ok, dt} -> dt
+ {:error, reason} -> IO.puts("Error: #{reason}")
+ end
+ end
+
+ def parse_timestamp(timestamp) do
+ parse_timestamp_usec(timestamp)
+ |> NaiveDateTime.truncate(:second)
+ end
+
+ def parse_id_list(id_list) do
+ Enum.map(id_list, fn id ->
+ id
+ |> FlakeId.from_integer()
+ |> FlakeId.to_string()
+ end)
+ end
+
+ def truncate(str, max_length) do
+ String.slice(str, 0, max_length)
+ end
+
+ def try_create_activity(params) do
+ activity_params =
+ params
+ |> Map.delete(:object_data)
+
+ try do
+ {:ok, _activity} = Repo.insert(struct(Activity, activity_params))
+ shell_info("Activity created")
+
+ if params[:object_data] do
+ try_create_object(params)
+ end
+ rescue
+ Ecto.ConstraintError ->
+ shell_info("Activity already in database, skipping")
+ end
+ end
+
+ defp try_create_object(params) do
+ object_data =
+ params[:object_data]
+ # |> Transmogrifier.strip_internal_fields # We need internal fields for `likes` and `like_count`, etc
+ # |> Transmogrifier.fix_actor # Makes network requests
+ |> Transmogrifier.fix_url()
+ |> Transmogrifier.fix_attachments()
+ |> Transmogrifier.fix_context()
+ # |> Transmogrifier.fix_in_reply_to # Makes network requests
+ |> Transmogrifier.fix_emoji()
+ |> Transmogrifier.fix_tag()
+ |> Transmogrifier.fix_content_map()
+ # |> Transmogrifier.fix_addressing # Makes network requests
+ |> Transmogrifier.fix_summary()
+
+ # |> Transmogrifier.fix_type
+
+ object_params = %{
+ data: object_data,
+ inserted_at: params[:inserted_at],
+ updated_at: params[:updated_at]
+ }
+
+ try do
+ {:ok, _object} = Repo.insert(struct(Object, object_params))
+ shell_info("Object created")
+ rescue
+ Ecto.ConstraintError ->
+ shell_info("Object already in database, skipping")
+ end
+ end
+end
diff --git a/lib/mix/tasks/migrator/import.ex b/lib/mix/tasks/migrator/import.ex
new file mode 100644
index 000000000..f28ae0e8f
--- /dev/null
+++ b/lib/mix/tasks/migrator/import.ex
@@ -0,0 +1,18 @@
+defmodule Mix.Tasks.Migrator.Import do
+ use Mix.Task
+ alias Mix.Tasks.Migrator
+
+ @shortdoc "Import all dumps."
+ def run(_) do
+ Migrator.Import.Users.run(nil)
+ Migrator.Import.Statuses.run(nil)
+ Migrator.Import.Likes.run(nil)
+ Migrator.Import.Follows.run(nil)
+ Migrator.Import.Votes.run(nil)
+ Migrator.Import.Blocks.run(nil)
+ Migrator.Import.Mutes.run(nil)
+ Migrator.Import.ThreadMutes.run(nil)
+ Migrator.Import.Lists.run(nil)
+ Migrator.Import.Filters.run(nil)
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/blocks.ex b/lib/mix/tasks/migrator/import/blocks.ex
new file mode 100644
index 000000000..1f1f3c5aa
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/blocks.ex
@@ -0,0 +1,37 @@
+defmodule Mix.Tasks.Migrator.Import.Blocks do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+ alias Pleroma.User
+ alias Pleroma.UserRelationship
+
+ @shortdoc "Import blocks."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/blocks.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+ |> loop_fields([:inserted_at, :updated_at], &parse_timestamp/1)
+
+ try_create_activity(params)
+ try_create_block(params)
+ end
+
+ defp try_create_block(%{data: %{"actor" => actor, "object" => object}} = _params) do
+ try do
+ source = User.get_by_ap_id(actor)
+ target = User.get_by_ap_id(object)
+ UserRelationship.create_block(source, target)
+ shell_info("Block created")
+ rescue
+ MatchError -> shell_info("Could not create block")
+ FunctionClauseError -> shell_info("Could not create block")
+ end
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/filters.ex b/lib/mix/tasks/migrator/import/filters.ex
new file mode 100644
index 000000000..1956f3459
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/filters.ex
@@ -0,0 +1,46 @@
+defmodule Mix.Tasks.Migrator.Import.Filters do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+ alias Pleroma.User
+ alias Pleroma.Filter
+ alias Pleroma.Repo
+
+ @shortdoc "Import filters."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/filters.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+ |> loop_fields([:inserted_at, :updated_at], &parse_timestamp/1)
+ |> fix_user_id
+
+ try_create_filter(params)
+ end
+
+ defp fix_user_id(params) do
+ creator = User.get_by_ap_id(params[:user_ap_id])
+
+ Map.put(params, :user_id, creator.id)
+ |> Map.delete(:user_ap_id)
+ end
+
+ defp try_create_filter(params) do
+ changeset = struct(Filter, params)
+
+ try do
+ {:ok, _filter} = Repo.insert(changeset)
+ shell_info("Filter created")
+ rescue
+ Ecto.ConstraintError -> shell_info("Filter already exists, skipping")
+ MatchError -> shell_info("Could not create filter")
+ FunctionClauseError -> shell_info("Could not create filter")
+ end
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/follows.ex b/lib/mix/tasks/migrator/import/follows.ex
new file mode 100644
index 000000000..85af07fb1
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/follows.ex
@@ -0,0 +1,38 @@
+defmodule Mix.Tasks.Migrator.Import.Follows do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+ alias Pleroma.User
+ alias Pleroma.FollowingRelationship
+
+ @shortdoc "Import follows."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/follows.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+ |> loop_fields([:inserted_at, :updated_at], &parse_timestamp/1)
+
+ shell_info("Importing follow...")
+ try_create_activity(params)
+ create_follow(params)
+ end
+
+ defp create_follow(%{data: %{"actor" => actor, "object" => object}} = _params) do
+ try do
+ follower = User.get_by_ap_id(actor)
+ following = User.get_by_ap_id(object)
+ FollowingRelationship.follow(follower, following)
+ shell_info("Follow relationship created")
+ rescue
+ MatchError -> shell_info("Could not create follow")
+ FunctionClauseError -> shell_info("Could not create follow")
+ end
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/likes.ex b/lib/mix/tasks/migrator/import/likes.ex
new file mode 100644
index 000000000..088be0a5f
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/likes.ex
@@ -0,0 +1,23 @@
+defmodule Mix.Tasks.Migrator.Import.Likes do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+
+ @shortdoc "Import likes."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/likes.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+ |> loop_fields([:inserted_at, :updated_at], &parse_timestamp/1)
+
+ shell_info("Importing like...")
+ try_create_activity(params)
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/lists.ex b/lib/mix/tasks/migrator/import/lists.ex
new file mode 100644
index 000000000..bc2c450c1
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/lists.ex
@@ -0,0 +1,46 @@
+defmodule Mix.Tasks.Migrator.Import.Lists do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+ alias Pleroma.User
+ alias Pleroma.List
+ alias Pleroma.Repo
+
+ @shortdoc "Import lists."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/lists.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+ |> loop_fields([:inserted_at, :updated_at], &parse_timestamp/1)
+ |> fix_user_id
+
+ try_create_list(params)
+ end
+
+ defp fix_user_id(params) do
+ creator = User.get_by_ap_id(params[:user_ap_id])
+
+ Map.put(params, :user_id, creator.id)
+ |> Map.delete(:user_ap_id)
+ end
+
+ defp try_create_list(params) do
+ changeset = struct(List, params)
+
+ try do
+ {:ok, _list} = Repo.insert(changeset)
+ shell_info("List created")
+ rescue
+ Ecto.ConstraintError -> shell_info("List already exists, skipping")
+ MatchError -> shell_info("Could not create list")
+ FunctionClauseError -> shell_info("Could not create list")
+ end
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/mutes.ex b/lib/mix/tasks/migrator/import/mutes.ex
new file mode 100644
index 000000000..e16e5a27e
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/mutes.ex
@@ -0,0 +1,36 @@
+defmodule Mix.Tasks.Migrator.Import.Mutes do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+ alias Pleroma.User
+ alias Pleroma.UserRelationship
+
+ @shortdoc "Import mutes."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/mutes.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+ |> loop_fields([:inserted_at, :updated_at], &parse_timestamp/1)
+
+ try_create_mute(params)
+ end
+
+ defp try_create_mute(%{source_ap_id: source_ap_id, target_ap_id: target_ap_id} = _params) do
+ try do
+ source = User.get_by_ap_id(source_ap_id)
+ target = User.get_by_ap_id(target_ap_id)
+ UserRelationship.create_mute(source, target)
+ shell_info("Mute created")
+ rescue
+ MatchError -> shell_info("Could not create mute")
+ FunctionClauseError -> shell_info("Could not create mute")
+ end
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/statuses.ex b/lib/mix/tasks/migrator/import/statuses.ex
new file mode 100644
index 000000000..07af18af8
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/statuses.ex
@@ -0,0 +1,23 @@
+defmodule Mix.Tasks.Migrator.Import.Statuses do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+
+ @shortdoc "Import statuses."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/statuses.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+ |> loop_fields([:inserted_at, :updated_at], &parse_timestamp/1)
+
+ shell_info("Importing status #{params.id}...")
+ try_create_activity(params)
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/thread_mutes.ex b/lib/mix/tasks/migrator/import/thread_mutes.ex
new file mode 100644
index 000000000..d56439055
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/thread_mutes.ex
@@ -0,0 +1,34 @@
+defmodule Mix.Tasks.Migrator.Import.ThreadMutes do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+ alias Pleroma.User
+ alias Pleroma.ThreadMute
+
+ @shortdoc "Import thread mutes."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/thread_mutes.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+
+ try_create_thread_mute(params)
+ end
+
+ defp try_create_thread_mute(%{ap_id: ap_id, context: context} = _params) do
+ try do
+ %User{id: user_id} = User.get_by_ap_id(ap_id)
+ ThreadMute.add_mute(user_id, context)
+ shell_info("Thread mute created")
+ rescue
+ MatchError -> shell_info("Could not create thread mute")
+ FunctionClauseError -> shell_info("Could not create thread mute")
+ end
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/users.ex b/lib/mix/tasks/migrator/import/users.ex
new file mode 100644
index 000000000..657fa62f5
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/users.ex
@@ -0,0 +1,47 @@
+defmodule Mix.Tasks.Migrator.Import.Users do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+ alias Pleroma.User
+ alias Pleroma.Repo
+
+ @shortdoc "Import users."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/users.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ bio_limit = Pleroma.Config.get([:instance, :user_bio_length], 5000)
+ name_limit = Pleroma.Config.get([:instance, :user_name_length], 100)
+
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+ |> loop_fields(
+ [
+ :inserted_at,
+ :updated_at,
+ :last_digest_emailed_at
+ ],
+ &parse_timestamp/1
+ )
+ |> loop_fields([:last_refreshed_at], &parse_timestamp_usec/1)
+ |> loop_fields([:notification_settings], &keys_to_atoms/1)
+ |> loop_fields([:name], &truncate(&1, name_limit))
+ |> loop_fields([:bio], &truncate(&1, bio_limit))
+
+ changeset = struct(User, params)
+
+ try do
+ {:ok, user} = Repo.insert(changeset)
+ User.set_cache(user)
+ shell_info("User #{params.nickname} created")
+ rescue
+ Ecto.ConstraintError ->
+ shell_info("User #{params.nickname} already in database, skipping")
+ end
+ end
+end
diff --git a/lib/mix/tasks/migrator/import/votes.ex b/lib/mix/tasks/migrator/import/votes.ex
new file mode 100644
index 000000000..8ca42e551
--- /dev/null
+++ b/lib/mix/tasks/migrator/import/votes.ex
@@ -0,0 +1,23 @@
+defmodule Mix.Tasks.Migrator.Import.Votes do
+ use Mix.Task
+ import Mix.Pleroma
+ import Mix.Migrator
+
+ @shortdoc "Import poll votes."
+ def run(_) do
+ start_pleroma()
+
+ File.stream!("migrator/votes.txt")
+ |> Enum.each(&handle_line/1)
+ end
+
+ defp handle_line(line) do
+ params =
+ Jason.decode!(line)
+ |> keys_to_atoms
+ |> loop_fields([:inserted_at, :updated_at], &parse_timestamp/1)
+
+ shell_info("Importing poll vote...")
+ try_create_activity(params)
+ end
+end
diff --git a/lib/mix/tasks/migrator/rebuild.ex b/lib/mix/tasks/migrator/rebuild.ex
new file mode 100644
index 000000000..4e395e15a
--- /dev/null
+++ b/lib/mix/tasks/migrator/rebuild.ex
@@ -0,0 +1,58 @@
+defmodule Mix.Tasks.Migrator.Rebuild do
+ use Mix.Task
+ alias Pleroma.Object
+ alias Pleroma.User
+ alias Pleroma.Web.ActivityPub.ActivityPub
+ alias Pleroma.Web.ActivityPub.Transmogrifier
+
+ require Logger
+ import Ecto.Query
+ import Mix.Pleroma
+
+ @shortdoc "Rebuild remote data."
+
+ def run(["users"]) do
+ start_pleroma()
+
+ User
+ |> where([u], u.local == false)
+ |> where([u], is_nil(u.last_refreshed_at))
+ |> Pleroma.RepoStreamer.chunk_stream(500)
+ |> Stream.each(fn users ->
+ users
+ |> Enum.each(fn user ->
+ try do
+ ActivityPub.make_user_from_ap_id(user.ap_id)
+ shell_info("Updating @#{user.nickname}")
+ rescue
+ _ ->
+ shell_info("Couldn't update user. Skipping.")
+ end
+ end)
+ end)
+ |> Stream.run()
+ end
+
+ def run(["activities"]) do
+ start_pleroma()
+
+ Object
+ |> where([o], fragment("?->'pleroma_internal'->>'migrator'='true'", o.data))
+ |> Pleroma.RepoStreamer.chunk_stream(500)
+ |> Stream.each(fn objects ->
+ objects
+ |> Enum.each(fn object ->
+ shell_info("Transmogrifying #{object.data["id"]}")
+ # This doesn't write anything back to the database, it just
+ # fetches anything that's missing.
+ try do
+ Transmogrifier.fix_object(object.data)
+ rescue
+ _ ->
+ shell_info("Couldn't transmogrify. Skipping.")
+ end
+ end)
+ end)
+ |> Stream.run()
+ end
+end

File Metadata

Mime Type
text/x-diff
Expires
Fri, Aug 28, 4:23 PM (8 h, 10 m)
Storage Engine
blob
Storage Format
Raw Data
Storage Handle
1736515
Default Alt Text
(16 KB)

Event Timeline