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