diff --git a/AGENTS.md b/AGENTS.md index 6a354efb5..bb130e378 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -71,7 +71,7 @@ Phoenix 1.8 server-rendered MVC — **no LiveView**. Views/templates use `phoeni Four Elixir app namespaces under `lib/`: -- **`Philomena`** — domain contexts (images, tags, forums, comments, users, filters, galleries, notifications, ...). Each context is a `.ex` module plus a `/` directory of Ecto schemas and helpers. Background jobs live in `lib/philomena/workers/` and run via Exq (Redis/Valkey-backed); contexts enqueue e.g. `IndexWorker`, `ThumbnailWorker`. +- **`Philomena`** — domain contexts (images, tags, forums, comments, users, filters, galleries, notifications, ...). Each context is a `.ex` module plus a `/` directory of Ecto schemas and helpers. Background jobs live in `lib/philomena/workers/` and run via Oban (Postgres-backed); contexts enqueue e.g. `IndexWorker`, `ThumbnailWorker`. - **`PhilomenaWeb`** — controllers, plugs, views, templates. Routing is aggressively RESTful: instead of custom actions there are many small nested singleton controllers (e.g. `Image.VoteController`, `Topic.SubscriptionController`) with only `create`/`delete`. The public JSON API is `lib/philomena_web/controllers/api/json/` and is documented by `openapi.yaml` at the repo root — keep the two in sync. Authorization uses Canada/Canary (`can?` protocols + plugs). - **`PhilomenaQuery`** — the search layer. `parse/` is a nimble_parsec-based parser for the user-facing search query language; `search.ex` + `search/` is the OpenSearch client. Each searchable domain implements the `PhilomenaQuery.Search.Index` behaviour (e.g. `Philomena.Images.SearchIndex`) defining the index mapping and document serialization. Data flow: writes go to Postgres, then documents are (re)indexed into OpenSearch via `PhilomenaQuery.Search.reindex`/`IndexWorker`. - **`PhilomenaMedia`** — media intake pipeline: `analyzers/` (mime/dimension/duration detection), `processors/` (per-format thumbnailing/optimization), intensities for duplicate detection, and `objects.ex` for S3 storage (ex_aws; s3proxy in dev). diff --git a/config/config.exs b/config/config.exs index 006c1fc41..9e1247a8b 100644 --- a/config/config.exs +++ b/config/config.exs @@ -24,10 +24,10 @@ config :philomena, search_target_poll_interval_ms: 5_000, search_migration_settle_ms: 15_000 -config :exq, - max_retries: 5, - scheduler_enable: true, - start_on_application: false +config :philomena, Oban, + repo: Philomena.Repo, + plugins: [Oban.Pruner], + queues: [videos: 2, images: 4, indexing: 12, notifications: 2] # Configures the endpoint config :philomena, PhilomenaWeb.Endpoint, diff --git a/config/runtime.exs b/config/runtime.exs index d84d979dc..216fbf188 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -58,18 +58,9 @@ json_config = config :philomena, config: json_config -config :exq, - host: System.get_env("REDIS_HOST", "localhost"), - queues: [ - {"videos", 2}, - {"images", 4}, - {"indexing", 12}, - {"notifications", 2} - ] - if is_nil(System.get_env("START_WORKER")) do # Make queueing available but don't process any jobs - config :exq, queues: [] + config :philomena, Oban, queues: false end # S3/Object store config diff --git a/config/test.exs b/config/test.exs index 60509a63b..216994a4b 100644 --- a/config/test.exs +++ b/config/test.exs @@ -21,11 +21,11 @@ config :philomena, pwned_passwords: false, captcha: false -# Keep test enqueues in memory. The application still exercises the same -# enqueue calls, but the test suite cannot fill the shared development Valkey -# instance with jobs that reference the test database. -config :exq, - queue_adapter: Exq.Adapters.Queue.Mock +# Keep test enqueues in the sandbox database without starting queue consumers. +config :philomena, Oban, + testing: :manual, + queues: false, + plugins: false # Namespace OpenSearch indexes so test runs cannot touch dev data on the # shared cluster. Search-backed tests recreate their index in setup; see diff --git a/docker-compose.yml b/docker-compose.yml index cb343ff0d..c4ba9a54a 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -22,7 +22,7 @@ services: # event's module and function. The last entry in the list of filters should # be a bare `level` which will be used as a catch-all for all other log # events that do not match any of the previous filters. - - PHILOMENA_LOG=${PHILOMENA_LOG-Ecto=debug,Exq=none,PhilomenaMedia.Objects=info,debug} + - PHILOMENA_LOG=${PHILOMENA_LOG-Ecto=debug,Oban=none,PhilomenaMedia.Objects=info,debug} - MIX_ENV=dev - PGPASSWORD=postgres - ANONYMOUS_NAME_SALT=2fmJRo0OgMFe65kyAJBxPT0QtkVes/jnKDdtP21fexsRqiw8TlSY7yO+uFyMZycp diff --git a/lib/philomena/application.ex b/lib/philomena/application.ex index 6182e25f8..3d808f4f7 100644 --- a/lib/philomena/application.ex +++ b/lib/philomena/application.ex @@ -15,11 +15,11 @@ defmodule Philomena.Application do # Search write-target tracking, so document writes fan out to both # indices during a search index migration. Must start before anything - # which writes documents (the endpoint and the Exq workers). + # which writes documents (the endpoint and the Oban workers). {PhilomenaQuery.Search.WriteTargets, []}, # Background queueing system - Philomena.ExqSupervisor, + {Oban, Application.fetch_env!(:philomena, Oban)}, # Mailer {Task.Supervisor, name: Philomena.AsyncEmailSupervisor}, diff --git a/lib/philomena/bans.ex b/lib/philomena/bans.ex index 3244a966f..1badb3bb0 100644 --- a/lib/philomena/bans.ex +++ b/lib/philomena/bans.ex @@ -29,6 +29,7 @@ defmodule Philomena.Bans do alias Philomena.Bans.UserQueryForm alias Philomena.ModerationLogs alias Philomena.Multi + alias Philomena.Workers.IndexJob alias Philomena.UserIps alias Philomena.Users @@ -553,9 +554,7 @@ defmodule Philomena.Bans do |> repo.insert() end end) - |> Multi.on_commit(fn _changes -> - Users.reindex_user(%Users.User{id: target.id}) - end) + |> IndexJob.put_enqueue("Users", :id, [target.id]) end @doc """ diff --git a/lib/philomena/comments.ex b/lib/philomena/comments.ex index a7e3fe27b..ffc4cfcb0 100644 --- a/lib/philomena/comments.ex +++ b/lib/philomena/comments.ex @@ -19,7 +19,7 @@ defmodule Philomena.Comments do alias Philomena.Filters.Filter alias Philomena.Images alias Philomena.Images.Image - alias Philomena.IndexWorker + alias Philomena.Workers.IndexJob alias Philomena.IntegerId alias Philomena.Loader alias Philomena.ModerationLogs @@ -83,7 +83,7 @@ defmodule Philomena.Comments do end defp put_reindex_comment(%Multi{} = multi, step \\ :comment) do - Multi.on_commit(multi, fn %{^step => comment} -> reindex_comment(comment) end) + IndexJob.put_enqueue(multi, "Comments", :id, fn %{^step => comment} -> [comment.id] end) end defp put_approval_report(%Multi{} = multi) do @@ -354,8 +354,9 @@ defmodule Philomena.Comments do Write access, image commenting permission, the Images-owned forced-filter prerequisite, and the 15-second creation limit are checked before insertion. The transaction updates the image count, notification, and subscription state. - Indexing, statistics/reporting, rate tracking, and the firehose broadcast run - after commit. The image is returned for the caller to reuse. + Indexing, statistics/reporting, and rate tracking are committed with the + write; the firehose broadcast runs after commit. The image is returned for + the caller to reuse. ## Examples @@ -477,7 +478,8 @@ defmodule Philomena.Comments do Write access is checked before image authorization, forced-filter enforcement, and comment authorization. A successful transaction records the prior version - Reporting, indexing, and the firehose broadcast run after commit. Validation + Reporting and indexing are committed with the write; the firehose broadcast + runs after commit. Validation returns the changeset preserving the loaded comment and image. On success, the image is returned for the caller to reuse. @@ -828,51 +830,6 @@ defmodule Philomena.Comments do Search.update_by_query(Comment, data.query, data.set_replacements, data.replacements) end - @doc """ - Queues one comment for search indexing and returns it unchanged. - - ## Examples - - iex> reindex_comment(comment) - %Comment{} - - """ - @spec reindex_comment(Comment.t()) :: Comment.t() - def reindex_comment(%Comment{} = comment) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Comments", "id", [comment.id]]) - comment - end - - @doc """ - Queues every comment on `image` for indexing and returns the image unchanged. - - ## Examples - - iex> reindex_comments_on_image(image) - %Image{} - - """ - @spec reindex_comments_on_image(Image.t()) :: Image.t() - def reindex_comments_on_image(%Image{} = image) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Comments", "image_id", [image.id]]) - image - end - - @doc """ - Queues comments on the given image IDs for reindexing and returns the list unchanged. - - ## Examples - - iex> reindex_comments_on_images([1, 2, 3]) - [1, 2, 3] - - """ - @spec reindex_comments_on_images([integer()]) :: [integer()] - def reindex_comments_on_images(image_ids) when is_list(image_ids) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Comments", "image_id", image_ids]) - image_ids - end - @doc """ Returns the association queries required to serialize comment search records. diff --git a/lib/philomena/exq_supervisor.ex b/lib/philomena/exq_supervisor.ex deleted file mode 100644 index cdbce815b..000000000 --- a/lib/philomena/exq_supervisor.ex +++ /dev/null @@ -1,17 +0,0 @@ -defmodule Philomena.ExqSupervisor do - use Supervisor - - def start_link(init_arg) do - Supervisor.start_link(__MODULE__, init_arg, name: __MODULE__) - end - - @impl true - def init(_init_arg) do - child_spec = %{ - id: Exq, - start: {Exq, :start_link, []} - } - - Supervisor.init([child_spec], strategy: :one_for_one) - end -end diff --git a/lib/philomena/filters.ex b/lib/philomena/filters.ex index d90c9d42e..ab9c86904 100644 --- a/lib/philomena/filters.ex +++ b/lib/philomena/filters.ex @@ -24,7 +24,7 @@ defmodule Philomena.Filters do alias Philomena.Users alias Philomena.Users.User alias PhilomenaQuery.Search - alias Philomena.IndexWorker + alias Philomena.Workers.IndexJob defp ensure_current_filter(%User{current_filter: current_filter} = user) do if current_filter do @@ -59,14 +59,7 @@ defmodule Philomena.Filters do end defp put_reindex_filter(multi, step) do - Multi.on_commit(multi, fn %{^step => filter} -> reindex_filter(filter) end) - end - - defp reindex_filter_ids([]), do: [] - - defp reindex_filter_ids(filter_ids) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Filters", "id", filter_ids]) - filter_ids + IndexJob.put_enqueue(multi, "Filters", :id, fn %{^step => filter} -> [filter.id] end) end @doc """ @@ -921,8 +914,9 @@ defmodule Philomena.Filters do |> Multi.all(spoilered_ids_step, select(exclude(spoilered_filters, :update), [f], f.id)) |> Multi.update_all(hidden_step, hidden_filters, []) |> Multi.update_all(spoilered_step, spoilered_filters, []) - |> Multi.on_commit(fn %{^hidden_ids_step => hidden_ids, ^spoilered_ids_step => spoilered_ids} -> - reindex_filter_ids(Enum.uniq(hidden_ids ++ spoilered_ids)) + |> IndexJob.put_enqueue("Filters", :id, fn + %{^hidden_ids_step => hidden_ids, ^spoilered_ids_step => spoilered_ids} -> + Enum.uniq(hidden_ids ++ spoilered_ids) end) end @@ -942,23 +936,6 @@ defmodule Philomena.Filters do Search.update_by_query(Filter, data.query, data.set_replacements, data.replacements) end - @doc """ - Queues a single filter for search index updates. - Returns the filter struct unchanged, for use in a pipeline. - - ## Examples - - iex> reindex_filter(filter) - %Filter{} - - """ - @spec reindex_filter(Filter.t()) :: Filter.t() - def reindex_filter(%Filter{} = filter) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Filters", "id", [filter.id]]) - - filter - end - @doc """ Returns a list of associations to preload when indexing filters. diff --git a/lib/philomena/galleries.ex b/lib/philomena/galleries.ex index e1383c21b..f0c8c6459 100644 --- a/lib/philomena/galleries.ex +++ b/lib/philomena/galleries.ex @@ -36,7 +36,7 @@ defmodule Philomena.Galleries do alias Philomena.Galleries.QueryForm alias Philomena.Galleries.ReorderForm alias Philomena.Galleries - alias Philomena.IndexWorker + alias Philomena.Workers.IndexJob alias Philomena.Interactions alias Philomena.Notifications alias Philomena.Images @@ -68,7 +68,7 @@ defmodule Philomena.Galleries do end defp put_reindex_gallery(%Multi{} = multi, step \\ :gallery) do - Multi.on_commit(multi, fn %{^step => gallery} -> reindex_gallery(gallery) end) + IndexJob.put_enqueue(multi, "Galleries", :id, fn %{^step => gallery} -> [gallery.id] end) end defp cleanup_gallery(%Gallery{} = gallery) do @@ -88,8 +88,8 @@ defmodule Philomena.Galleries do end, [] ) - |> Multi.on_commit(fn %{interactions: {_count, image_ids}} -> - Images.reindex_images(image_ids) + |> IndexJob.put_enqueue("Images", :id, fn %{interactions: {_count, image_ids}} -> + image_ids end) |> Multi.transact() end) @@ -115,7 +115,7 @@ defmodule Philomena.Galleries do |> Reports.put_close_reports(:reports, closing_user, gallery_id: gallery.id) |> Multi.delete(:gallery, fn %{locked_gallery: gallery} -> gallery end) |> Multi.on_commit(fn %{gallery: gallery} -> unindex_gallery(gallery) end) - |> Multi.on_commit(fn %{interactions: {_, image_ids}} -> Images.reindex_images(image_ids) end) + |> IndexJob.put_enqueue("Images", :id, fn %{interactions: {_, image_ids}} -> image_ids end) |> Multi.transact() |> case do {:ok, %{gallery: %Gallery{} = gallery}} -> @@ -750,8 +750,8 @@ defmodule Philomena.Galleries do |> Multi.run(:reorder, fn _repo, %{locked_gallery: gallery, reorder_form: reorder_form} -> persist_reorder_positions(gallery, reorder_form.image_ids) end) - |> Multi.on_commit(fn %{reorder_form: reorder_form} -> - Images.reindex_images(reorder_form.image_ids) + |> IndexJob.put_enqueue("Images", :id, fn %{reorder_form: reorder_form} -> + reorder_form.image_ids end) |> Multi.transact() |> case do @@ -843,7 +843,7 @@ defmodule Philomena.Galleries do counters and interaction rows therefore change atomically in the caller's transaction. - Affected galleries are reindexed after the transaction commits. + A job to reindex affected galleries is enqueued in the same transaction. """ @spec put_remove_image_interactions(Multi.t(), Image.t()) :: Multi.t() def put_remove_image_interactions(%Multi{} = multi, %Image{} = image) do @@ -856,7 +856,9 @@ defmodule Philomena.Galleries do multi |> Multi.update_all(:galleries, galleries, []) |> Multi.delete_all(:gallery_interactions, where(Interaction, image_id: ^image.id), []) - |> Multi.on_commit(fn %{galleries: {_, gallery_ids}} -> reindex_galleries(gallery_ids) end) + |> IndexJob.put_enqueue("Galleries", :id, fn %{galleries: {_, gallery_ids}} -> + gallery_ids + end) end @doc """ @@ -866,7 +868,7 @@ defmodule Philomena.Galleries do affected gallery counters are adjusted. Image merge workflows must hold both image locks before composing this operation. - Affected galleries are reindexed after the transaction commits. + A jobs to reindex affected galleries is enqueued in the same transaction. """ @spec put_migrate_image_interactions(Multi.t(), Image.t(), Image.t()) :: Multi.t() def put_migrate_image_interactions(%Multi{} = multi, %Image{} = source, %Image{} = target) do @@ -898,11 +900,12 @@ defmodule Philomena.Galleries do {:ok, {count, gallery_ids}} end) - |> Multi.on_commit(fn %{ - migrated_gallery_interactions: {_, migrated_gallery_ids}, - galleries: {_, removed_gallery_ids} - } -> - reindex_galleries(Enum.uniq(migrated_gallery_ids ++ removed_gallery_ids)) + |> IndexJob.put_enqueue("Galleries", :id, fn + %{ + migrated_gallery_interactions: {_, migrated_gallery_ids}, + galleries: {_, removed_gallery_ids} + } -> + Enum.uniq(migrated_gallery_ids ++ removed_gallery_ids) end) end @@ -929,42 +932,6 @@ defmodule Philomena.Galleries do Search.update_by_query(Gallery, data.query, data.set_replacements, data.replacements) end - @doc """ - Queues a gallery for reindexing. - - Adds the gallery to the indexing queue to update its search index. - - ## Examples - - iex> reindex_gallery(gallery) - %Gallery{} - - """ - @spec reindex_gallery(Gallery.t()) :: Gallery.t() - def reindex_gallery(%Gallery{} = gallery) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Galleries", "id", [gallery.id]]) - - gallery - end - - @doc """ - Queues multiple galleries for reindexing by their ids. - - ## Examples - - iex> reindex_galleries([1, 2, 3]) - [1, 2, 3] - - """ - @spec reindex_galleries([integer()]) :: [integer()] - def reindex_galleries([]), do: [] - - def reindex_galleries(gallery_ids) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Galleries", "id", gallery_ids]) - - gallery_ids - end - @doc """ Removes a gallery from the search index. diff --git a/lib/philomena/images.ex b/lib/philomena/images.ex index 41bef896c..ba4170a3a 100644 --- a/lib/philomena/images.ex +++ b/lib/philomena/images.ex @@ -18,8 +18,8 @@ defmodule Philomena.Images do alias Philomena.Repo alias PhilomenaQuery.Search - alias Philomena.ThumbnailWorker - alias Philomena.ImagePurgeWorker + alias Philomena.Workers.ThumbnailJob + alias Philomena.Workers.ImagePurgeJob alias Philomena.DuplicateReports alias Philomena.DnpEntries alias Philomena.Images.Image @@ -37,7 +37,7 @@ defmodule Philomena.Images do alias Philomena.Images.Subscription alias Philomena.Images alias Philomena.IntegerId - alias Philomena.IndexWorker + alias Philomena.Workers.IndexJob alias Philomena.Loader alias Philomena.RateLimiter alias Philomena.Attribution.Actor @@ -210,11 +210,13 @@ defmodule Philomena.Images do |> Multi.on_commit(fn %{image: image} -> spawn(fn -> Thumbnailer.hide_thumbnails(image, image.hidden_image_key) - purge_files(image, image.hidden_image_key) end) - - Comments.reindex_comments_on_image(image) - reindex_image(image) + end) + |> IndexJob.put_enqueue("Images", :id, fn %{image: image} -> [image.id] end) + |> IndexJob.put_enqueue("Comments", :image_id, fn %{image: image} -> [image.id] end) + |> ImagePurgeJob.put_enqueue(fn %{image: image} -> + Thumbnailer.thumbnail_urls(image, image.hidden_image_key) ++ + Thumbnailer.thumbnail_urls(image, nil) end) end @@ -336,9 +338,9 @@ defmodule Philomena.Images do |> Tags.put_image_tag_count_changes() |> UserStatistics.put_increment(actor.user, :metadata_updates_count) |> put_reindex_image(:image) + |> IndexJob.put_enqueue("Comments", :image_id, fn %{image: image} -> [image.id] end) |> Multi.on_commit(fn %{image: %{added_tags: added, removed_tags: removed} = image} -> image = Repo.preload(image, [:user, :sources, tags: :aliases]) - Comments.reindex_comments_on_image(image) broadcast_tag_update(image, added, removed) end) |> Multi.transact_with_automatic_retry() @@ -519,18 +521,17 @@ defmodule Philomena.Images do |> where(id: ^image.id) |> Repo.update_all(set: [thumbnails_generated: false, processed: false]) - enqueue_image_repair(image) + queue_image_repair(image) end - defp enqueue_image_repair(image) do - Exq.enqueue(Exq, queue(image.image_mime_type), ThumbnailWorker, [image.id]) + defp queue_image_repair(image) do + Multi.new() + |> ThumbnailJob.put_enqueue(image.id, image.image_mime_type) + |> Multi.transact() image end - defp queue("video/webm"), do: "videos" - defp queue(_mime_type), do: "images" - defp purge_files(image, hidden_key) do files = if is_nil(hidden_key) do @@ -540,7 +541,9 @@ defmodule Philomena.Images do Thumbnailer.thumbnail_urls(image, nil) end - Exq.enqueue(Exq, "indexing", ImagePurgeWorker, [files]) + Multi.new() + |> ImagePurgeJob.put_enqueue(fn _ -> files end) + |> Multi.transact() end ## Bulk operations @@ -628,9 +631,9 @@ defmodule Philomena.Images do end) |> TagChanges.put_batch_tag_changes(:inserted_taggings, :deleted_taggings, attributes) |> Tags.put_batch_image_count_changes(:inserted_taggings, :deleted_taggings, :visible_images) - |> Multi.on_commit(fn %{locked_image_ids: image_ids} -> - reindex_images(image_ids) - Comments.reindex_comments_on_images(image_ids) + |> IndexJob.put_enqueue("Images", :id, fn %{locked_image_ids: image_ids} -> image_ids end) + |> IndexJob.put_enqueue("Comments", :image_id, fn %{locked_image_ids: image_ids} -> + image_ids end) end @@ -1250,11 +1253,9 @@ defmodule Philomena.Images do |> Multi.run(:notification, fn _repo, _changes -> Notifications.broadcast_image_merge(image, duplicate_of_image) end) - |> Multi.on_commit(fn result -> - reindex_image(duplicate_of_image) - Comments.reindex_comments_on_image(duplicate_of_image) - broadcast_image_merge(result.image, duplicate_of_image) - end) + |> IndexJob.put_enqueue("Images", :id, [duplicate_of_image.id]) + |> IndexJob.put_enqueue("Comments", :image_id, [duplicate_of_image.id]) + |> Multi.on_commit(fn result -> broadcast_image_merge(result.image, duplicate_of_image) end) end @doc group: "Cross-context transaction helpers" @@ -1262,9 +1263,7 @@ defmodule Philomena.Images do Adds the inverse of a source change to `multi` without recording a new source-change row. - The image is locked and reindexed after the transaction commits. The - corresponding source-change row can be deleted by composing this operation - with `Philomena.SourceChanges.put_erase_source_change/2`. + The image is locked and a reindex job is enqueued in the same transaction. ## Examples @@ -1329,7 +1328,7 @@ defmodule Philomena.Images do @doc """ Adds multiple denormalized image counter adjustments to `multi`. - The owner attaches one image reindex after commit for the complete update. + The owner attaches one image reindex in the same transaction for the complete update. """ @spec put_image_counter_deltas( Multi.t(), @@ -1345,7 +1344,7 @@ defmodule Philomena.Images do {:ok, count} end) - |> Multi.on_commit(fn _changes -> reindex_images([image_id]) end) + |> IndexJob.put_enqueue("Images", :id, [image_id]) end @doc group: "Cross-context transaction helpers" @@ -1362,7 +1361,7 @@ defmodule Philomena.Images do multi |> Multi.all(image_ids_step, image_ids_query) |> Multi.delete_all(step, query) - |> Multi.on_commit(fn %{^image_ids_step => image_ids} -> reindex_images(image_ids) end) + |> IndexJob.put_enqueue("Images", :id, fn %{^image_ids_step => image_ids} -> image_ids end) end @doc group: "Cross-context transaction helpers" @@ -1378,10 +1377,8 @@ defmodule Philomena.Images do on_conflict: :nothing, returning: [:image_id, :tag_id] ) - |> Multi.on_commit(fn %{^step => {_count, taggings}} -> - taggings - |> Enum.map(& &1.image_id) - |> reindex_images() + |> IndexJob.put_enqueue("Images", :id, fn %{^step => {_count, taggings}} -> + Enum.map(taggings, & &1.image_id) end) end @@ -1395,7 +1392,7 @@ defmodule Philomena.Images do on_conflict: :nothing, returning: [:image_id, :tag_id] ) - |> Multi.on_commit(fn %{^image_ids_step => image_ids} -> reindex_images(image_ids) end) + |> IndexJob.put_enqueue("Images", :id, fn %{^image_ids_step => image_ids} -> image_ids end) end @doc group: "Cross-context transaction helpers" @@ -1427,7 +1424,7 @@ defmodule Philomena.Images do |> Multi.run(:copied_tag_ids, fn _repo, %{target_taggings: {_count, taggings}} -> {:ok, Enum.map(taggings, & &1.tag_id)} end) - |> Multi.on_commit(fn _changes -> reindex_images([target.id]) end) + |> IndexJob.put_enqueue("Images", :id, [target.id]) end @doc group: "Forms and uploads" @@ -2001,11 +1998,14 @@ defmodule Philomena.Images do Paths.image_path(image), "Repaired image #{image.id}" ) + |> ThumbnailJob.put_enqueue(image.id, image.image_mime_type) + |> ImagePurgeJob.put_enqueue(fn _changes -> + Thumbnailer.thumbnail_urls(image, image.hidden_image_key) ++ + Thumbnailer.thumbnail_urls(image, nil) + end) |> Multi.transact() |> case do {:ok, _changes} -> - enqueue_image_repair(image) - purge_files(image, image.hidden_image_key) {:ok, image} error -> @@ -2109,10 +2109,7 @@ defmodule Philomena.Images do 1 ) |> put_reindex_image(:image) - |> Multi.on_commit(fn %{image: image, tags: tags} -> - Comments.reindex_comments_on_image(image) - {image, tags} - end) + |> IndexJob.put_enqueue("Comments", :image_id, fn %{image: image} -> [image.id] end) |> ModerationLogs.put_log(:moderation_log, actor, fn %{image: image} -> {"Image.Delete:delete", Paths.image_path(image), "Restored image #{image.id}"} end) @@ -2173,7 +2170,7 @@ defmodule Philomena.Images do "Deleted #{vote_type} by #{user.name} on image #{image.id}" } end) - |> Multi.on_commit(fn _changes -> reindex_image(image) end) + |> IndexJob.put_enqueue("Images", :id, [image.id]) |> Multi.transact() |> case do {:ok, _changes} -> {:ok, image} @@ -2978,15 +2975,22 @@ defmodule Philomena.Images do """ @spec update_thumbnail_metadata!(Image.t(), map(), :thumbnail | :process) :: Image.t() def update_thumbnail_metadata!(%Image{} = image, attrs, :thumbnail) do - image - |> Image.thumbnail_changeset(attrs) - |> Repo.update!() + update_thumbnail_metadata!(image, Image.thumbnail_changeset(image, attrs)) end def update_thumbnail_metadata!(%Image{} = image, attrs, :process) do - image - |> Image.process_changeset(attrs) - |> Repo.update!() + update_thumbnail_metadata!(image, Image.process_changeset(image, attrs)) + end + + defp update_thumbnail_metadata!(%Image{}, changeset) do + Multi.new() + |> Multi.update(:image, changeset) + |> IndexJob.put_enqueue("Images", :id, fn %{image: image} -> [image.id] end) + |> Multi.transact() + |> case do + {:ok, %{image: %Image{} = image}} -> + image + end end @doc group: "Background jobs" @@ -3384,14 +3388,14 @@ defmodule Philomena.Images do @doc group: "Search indexing" @doc """ - Adds an after-commit image reindex step to a transaction workflow. + Adds an image reindex job to a transaction workflow. The referenced step must resolve to an image. The indexing job is enqueued - only after the database transaction commits. + in the same database transaction. """ @spec put_reindex_image(Multi.t(), Ecto.Multi.name()) :: Multi.t() def put_reindex_image(%Multi{} = multi, step) do - Multi.on_commit(multi, fn %{^step => image} -> reindex_image(image) end) + IndexJob.put_enqueue(multi, "Images", :id, fn %{^step => image} -> [image.id] end) end @doc group: "Search indexing" @@ -3412,42 +3416,6 @@ defmodule Philomena.Images do Repo.one!(Image |> where(id: ^image_id) |> preload(:tags)) end - @doc group: "Search indexing" - @doc """ - Queues a single image for search index updates. - Returns the image struct unchanged, for use in a pipeline. - - ## Examples - - iex> reindex_image(image) - %Image{} - - """ - @spec reindex_image(Image.t()) :: Image.t() - def reindex_image(%Image{} = image) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Images", "id", [image.id]]) - - image - end - - @doc group: "Search indexing" - @doc """ - Queues all listed image IDs for search index updates. - Returns the list unchanged, for use in a pipeline. - - ## Examples - - iex> reindex_images([1, 2, 3]) - [1, 2, 3] - - """ - @spec reindex_images([integer()]) :: [integer()] - def reindex_images(image_ids) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Images", "id", image_ids]) - - image_ids - end - @doc group: "Search indexing" @doc """ Returns the preload configuration for image indexing. diff --git a/lib/philomena/images/thumbnailer.ex b/lib/philomena/images/thumbnailer.ex index 1cde081d8..bb91b2abf 100644 --- a/lib/philomena/images/thumbnailer.ex +++ b/lib/philomena/images/thumbnailer.ex @@ -12,7 +12,7 @@ defmodule Philomena.Images.Thumbnailer do alias Philomena.DuplicateReports alias Philomena.ImageIntensities alias Philomena.Images - alias Philomena.ImagePurgeWorker + alias Philomena.Workers.ImagePurgeJob alias Philomena.Images.Image alias Philomena.Repo @@ -115,9 +115,7 @@ defmodule Philomena.Images.Thumbnailer do full = "full.#{image.image_format}" upload_file(image, new_file, full) - Exq.enqueue(Exq, "indexing", ImagePurgeWorker, [ - Path.join(image_url_base(image, nil), full) - ]) + ImagePurgeJob.enqueue([Path.join(image_url_base(image, nil), full)]) end defp apply_change(image, {:thumbnails, thumbnails}), diff --git a/lib/philomena/job_queue.ex b/lib/philomena/job_queue.ex new file mode 100644 index 000000000..1f627f125 --- /dev/null +++ b/lib/philomena/job_queue.ex @@ -0,0 +1,37 @@ +defmodule Philomena.JobQueue do + @moduledoc """ + The application boundary for enqueueing Oban jobs. + + Job arguments are JSON-safe maps so they can be consumed directly by the + worker modules and remain self-documenting when inspected in Oban. + """ + + alias Philomena.Multi + + @doc "Adds a worker insert to a Philomena.Multi." + @spec put_enqueue( + Multi.t(), + module(), + map() | (Ecto.Multi.changes() -> map()), + keyword() | (Ecto.Multi.changes() -> keyword()) + ) :: Multi.t() + def put_enqueue(%Multi{multi: ecto_multi} = multi, worker, args, opts \\ []) do + changeset = fn changes -> + args = if is_function(args, 1), do: args.(changes), else: args + opts = if is_function(opts, 1), do: opts.(changes), else: opts + worker.new(args, opts) + end + + %{multi | multi: Oban.insert(ecto_multi, {:oban, make_ref()}, changeset)} + end + + @doc "Inserts a worker immediately for standalone operations." + @spec insert(module(), map(), keyword()) :: :ok + def insert(worker, args, opts \\ []) do + args + |> worker.new(opts) + |> Oban.insert!() + + :ok + end +end diff --git a/lib/philomena/multi.ex b/lib/philomena/multi.ex index da02f91e2..aaba1da28 100644 --- a/lib/philomena/multi.ex +++ b/lib/philomena/multi.ex @@ -764,14 +764,16 @@ defmodule Philomena.Multi do The callback receives the transaction changes and runs only after a successful transaction. There is no ordering guarantee of post-commit - callback execution. Use this for side effects that must occur after - transaction completion, like object storage or indexing. + callback execution. Use this only for side effects that deliberately occur + after transaction completion. + + The return value of the callback is ignored. ## Example Multi.new() |> Multi.run(:user, fn _repo, _changes -> {:ok, user} end) - |> Multi.on_commit(fn %{user: user} -> Users.reindex_user(user) end) + |> Multi.on_commit(fn _changes -> :ok end) |> Multi.transact() """ @@ -785,6 +787,16 @@ defmodule Philomena.Multi do The callback receives the changes from the failed transaction. This is useful for compensating external reservations made by a `Multi.run/3` step. + + The return value of the callback is ignored. + + ## Example + + Multi.new() + |> Multi.run(:user, fn _repo, _changes -> {:error, changeset} end) + |> Multi.on_rollback(fn _changes -> :ok end) + |> Multi.transact() + """ @spec on_rollback(t(), (Ecto.Multi.changes() -> any())) :: t() def on_rollback(%__MODULE__{} = multi, callback) when is_function(callback, 1) do diff --git a/lib/philomena/posts.ex b/lib/philomena/posts.ex index 4aefbad78..77a893d7f 100644 --- a/lib/philomena/posts.ex +++ b/lib/philomena/posts.ex @@ -23,7 +23,7 @@ defmodule Philomena.Posts do alias Philomena.Users.User alias Philomena.Posts.{Post, PostVersion} alias Philomena.Posts - alias Philomena.IndexWorker + alias Philomena.Workers.IndexJob alias Philomena.Forums alias Philomena.Forums.Forum alias Philomena.Forums.Visibility @@ -63,14 +63,14 @@ defmodule Philomena.Posts do defp broadcast_post_creation(result), do: result defp put_reindex_post(%Multi{} = multi, step \\ :post) do - Multi.on_commit(multi, fn %{^step => post} -> reindex_post(post) end) + IndexJob.put_enqueue(multi, "Posts", :id, fn %{^step => post} -> [post.id] end) end @doc """ - Adds an `on_commit` step that reindexes all posts in a topic. + Adds a transactional indexing job for all posts in a topic. `step` names the transaction result containing the topic and defaults to - `:topic`. Indexing runs only after the transaction commits. + `:topic`. The indexing job is inserted with the transaction. ## Examples @@ -80,7 +80,7 @@ defmodule Philomena.Posts do """ @spec put_reindex_posts_in_topic(Multi.t(), Multi.name()) :: Multi.t() def put_reindex_posts_in_topic(%Multi{} = multi, step \\ :topic) do - Multi.on_commit(multi, fn %{^step => topic} -> reindex_posts_in_topic(topic) end) + IndexJob.put_enqueue(multi, "Posts", :topic_id, fn %{^step => topic} -> [topic.id] end) end @doc """ @@ -829,43 +829,6 @@ defmodule Philomena.Posts do Search.update_by_query(Post, data.query, data.set_replacements, data.replacements) end - @doc """ - Queues a single post for search index updates. - Returns the post struct unchanged, for use in a pipeline. - - ## Examples - - iex> reindex_post(post) - %Post{} - - """ - @spec reindex_post(Post.t()) :: Post.t() - def reindex_post(%Post{} = post) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Posts", "id", [post.id]]) - - post - end - - @doc """ - Queues every post in the given topic for search index updates. - - Used when a topic's visibility changes: a post's indexed `hidden_from_users` - reflects its topic's hidden state (see `Philomena.Posts.SearchIndex`), so - hiding or unhiding a topic must refresh all of its posts. - - ## Examples - - iex> reindex_posts_in_topic(topic_id) - :ok - - """ - @spec reindex_posts_in_topic(Topic.t()) :: :ok - def reindex_posts_in_topic(%Topic{} = topic) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Posts", "topic_id", [topic.id]]) - - :ok - end - @doc """ Provides preload queries for post indexing operations. diff --git a/lib/philomena/reports.ex b/lib/philomena/reports.ex index 85d814a74..a5da6cce5 100644 --- a/lib/philomena/reports.ex +++ b/lib/philomena/reports.ex @@ -18,7 +18,7 @@ defmodule Philomena.Reports do alias Philomena.Galleries.Gallery alias Philomena.Images alias Philomena.Images.Image - alias Philomena.IndexWorker + alias Philomena.Workers.IndexJob alias Philomena.Loader alias Philomena.ModerationLogs alias Philomena.ModerationLogs.Paths @@ -147,14 +147,8 @@ defmodule Philomena.Reports do ] end - defp reindex_closed_reports(report_ids) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Reports", "id", report_ids]) - end - defp put_reindex_report(%Multi{} = multi, report_step \\ :report) do - Multi.on_commit(multi, fn %{^report_step => report} -> - Exq.enqueue(Exq, "indexing", IndexWorker, ["Reports", "id", [report.id]]) - end) + IndexJob.put_enqueue(multi, "Reports", :id, fn %{^report_step => report} -> [report.id] end) end @doc """ @@ -594,8 +588,8 @@ defmodule Philomena.Reports do def put_close_reports(%Multi{} = multi, step, closing_user, target) do multi |> Multi.update_all(step, fn _ -> close_report_query(closing_user, target) end, []) - |> Multi.on_commit(fn %{^step => {_count, report_ids}} -> - reindex_closed_reports(report_ids) + |> IndexJob.put_enqueue("Reports", :id, fn %{^step => {_count, report_ids}} -> + report_ids end) end diff --git a/lib/philomena/tag_changes.ex b/lib/philomena/tag_changes.ex index dc1925405..1797dba00 100644 --- a/lib/philomena/tag_changes.ex +++ b/lib/philomena/tag_changes.ex @@ -9,14 +9,14 @@ defmodule Philomena.TagChanges do alias Philomena.Attribution.Actor alias Philomena.Images alias Philomena.Images.Image - alias Philomena.IndexWorker + alias Philomena.Workers.IndexJob alias Philomena.IntegerId alias Philomena.Loader alias Philomena.ModerationLogs alias Philomena.ModerationLogs.Paths alias Philomena.Multi alias Philomena.Repo - alias Philomena.TagChangeRevertWorker + alias Philomena.Workers.TagChangeRevertJob alias Philomena.TagChanges.QueryBuilder alias Philomena.TagChanges.QueryForm alias Philomena.TagChanges.RevertForm @@ -139,11 +139,7 @@ defmodule Philomena.TagChanges do batch_size: 100 } - Multi.on_commit(multi, fn _changes -> - Exq.enqueue(Exq, "indexing", TagChangeRevertWorker, [ - Map.put(target, :attributes, attributes) - ]) - end) + TagChangeRevertJob.put_enqueue(multi, target, attributes) end @doc """ @@ -456,8 +452,8 @@ defmodule Philomena.TagChanges do {:ok, {added_count, removed_count}} end) - |> Multi.on_commit(fn %{tag_change: tag_change} -> - Exq.enqueue(Exq, "indexing", IndexWorker, ["TagChanges", "id", [tag_change.id]]) + |> IndexJob.put_enqueue("TagChanges", :id, fn %{tag_change: tag_change} -> + [tag_change.id] end) end @@ -699,7 +695,7 @@ defmodule Philomena.TagChanges do @doc """ Adds deletion of tag change tag join table rows represented by `query` to - `multi` and queues the affected tag changes after commit. + `multi` and queues the affected tag changes transactionally. """ @spec put_delete_tag_change_tags(Multi.t(), Multi.name(), Ecto.Query.t()) :: Multi.t() def put_delete_tag_change_tags(%Multi{} = multi, step, %Ecto.Query{} = query) do @@ -711,8 +707,8 @@ defmodule Philomena.TagChanges do select(query, [tag_change_tag], tag_change_tag.tag_change_id) ) |> Multi.delete_all(step, query) - |> Multi.on_commit(fn %{^tag_change_ids_step => tag_change_ids} -> - reindex_tag_changes(tag_change_ids) + |> IndexJob.put_enqueue("TagChanges", :id, fn %{^tag_change_ids_step => tag_change_ids} -> + tag_change_ids end) end @@ -720,7 +716,7 @@ defmodule Philomena.TagChanges do Records tag changes from inserted and deleted image tagging rows in `multi`. Images supplies the two prior Multi steps. This function generates the tag - changes and queues the affected rows after commit. + changes and queues the affected rows transactionally. """ @spec put_batch_tag_changes(Multi.t(), Multi.name(), Multi.name(), map()) :: Multi.t() def put_batch_tag_changes(%Multi{} = multi, inserted_step, deleted_step, attributes) do @@ -778,8 +774,8 @@ defmodule Philomena.TagChanges do {:ok, tag_change_ids} end) - |> Multi.on_commit(fn %{batch_tag_changes: tag_change_ids} -> - reindex_tag_changes(tag_change_ids) + |> IndexJob.put_enqueue("TagChanges", :id, fn %{batch_tag_changes: tag_change_ids} -> + tag_change_ids end) end @@ -798,21 +794,6 @@ defmodule Philomena.TagChanges do Search.update_by_query(TagChange, data.query, data.set_replacements, data.replacements) end - @doc """ - Queues every tag change for worker reindexing. - - ## Examples - - iex> reindex_tag_changes([12, 13]) - [12, 13] - - """ - @spec reindex_tag_changes([integer()]) :: [integer()] - def reindex_tag_changes(ids) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["TagChanges", "id", ids]) - ids - end - @doc """ Returns the association projection required to serialize tag-change search documents. diff --git a/lib/philomena/tags.ex b/lib/philomena/tags.ex index 724feb005..e3d872554 100644 --- a/lib/philomena/tags.ex +++ b/lib/philomena/tags.ex @@ -24,16 +24,16 @@ defmodule Philomena.Tags do alias Philomena.Images.Search, as: ImageSearch alias Philomena.Images.Search.Scope alias Philomena.Images.Tagging - alias Philomena.IndexWorker + alias Philomena.Workers.IndexJob alias Philomena.Interactions alias Philomena.Loader alias Philomena.ModerationLogs alias Philomena.ModerationLogs.Paths alias Philomena.Multi alias Philomena.Repo - alias Philomena.TagAliasWorker - alias Philomena.TagDeleteWorker - alias Philomena.TagReindexWorker + alias Philomena.Workers.TagAliasJob + alias Philomena.Workers.TagDeleteJob + alias Philomena.Workers.TagReindexJob alias Philomena.TagChanges alias Philomena.TagChanges.TagChange alias Philomena.TagChanges.TagChangeTag @@ -83,18 +83,6 @@ defmodule Philomena.Tags do defp load_tag_for_action(_actor, _action, _slug, _preloads), do: {:error, :not_found} - defp reindex_tag_images(%Tag{} = tag) do - Exq.enqueue(Exq, "indexing", TagReindexWorker, [tag.id]) - tag - end - - defp reindex_tag_ids([]), do: [] - - defp reindex_tag_ids(tag_ids) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Tags", "id", tag_ids]) - tag_ids - end - # Computes the search query that lists the tag's images. A tag whose name # compiles back to itself is used verbatim. Anything else is escaped so the # search parser does not reinterpret it. @@ -214,10 +202,8 @@ defmodule Philomena.Tags do multi |> Multi.insert_all(:new_tags, Tag, insert_rows, insert_options) - |> Multi.on_commit(fn %{new_tags: {_count, new_tags}} -> - if Enum.any?(new_tags) do - reindex_tags(new_tags) - end + |> IndexJob.put_enqueue("Tags", :id, fn %{new_tags: {_count, new_tags}} -> + Enum.map(new_tags, & &1.id) end) end @@ -861,7 +847,8 @@ defmodule Philomena.Tags do Updates the tag named by `slug` on behalf of `actor`. Write access, `:update` authorization, the tag update, and its audit log share - one workflow. Search and affected-image reindexing run only after commit. + one workflow. Search and affected-image reindexing are enqueued in the same + transaction. ## Examples @@ -906,13 +893,13 @@ defmodule Philomena.Tags do Paths.tag_path(tag), "Updated details on tag '#{tag.name}'" ) - |> Multi.on_commit(fn %{tag: updated_tag} -> - # credo:disable-for-next-line + |> IndexJob.put_enqueue("Tags", :id, fn %{tag: updated_tag} -> [updated_tag.id] end) + |> Multi.merge(fn %{tag: updated_tag} -> if updated_tag.category != tag.category do - reindex_tag_images(updated_tag) + TagReindexJob.put_enqueue(Multi.new(), updated_tag.id) + else + Multi.new() end - - reindex_tags([updated_tag]) end) |> Multi.transact() |> case do @@ -1040,9 +1027,7 @@ defmodule Philomena.Tags do Paths.tag_path(tag), "Deleted tag '#{tag.name}'" ) - |> Multi.on_commit(fn _changes -> - Exq.enqueue(Exq, "indexing", TagDeleteWorker, [tag.id]) - end) + |> TagDeleteJob.put_enqueue(tag.id) |> Multi.transact() |> case do {:ok, _changes} -> @@ -1056,7 +1041,7 @@ defmodule Philomena.Tags do Write access, `:alias` authorization, the association migration, and its audit log share one transaction. Tagging migration and alias finalization - are queued after commit. + are queued transactionally. ## Examples @@ -1126,9 +1111,10 @@ defmodule Philomena.Tags do "Aliased tag '#{source_tag.name}' into '#{target_tag.name}'" } end) - |> Multi.on_commit(fn %{tags: {source_tag, target_tag}} -> - Exq.enqueue(Exq, "indexing", TagAliasWorker, [source_tag.id, target_tag.id]) - end) + |> TagAliasJob.put_enqueue( + fn %{tags: {source_tag, _target_tag}} -> source_tag.id end, + fn %{tags: {_source_tag, target_tag}} -> target_tag.id end + ) |> Multi.transact() |> case do {:ok, %{tag: %Tag{} = source_tag}} -> @@ -1161,18 +1147,23 @@ defmodule Philomena.Tags do def create_tag_reindex(%Actor{} = actor, slug) do with :ok <- verify_write_access(actor), {:ok, tag} <- load_tag_for_action(actor, :reindex, slug, @alias_preloads) do - reindex_tag_images(tag) - reindex_tags([tag]) - - {:ok, tag} + Multi.new() + |> TagReindexJob.put_enqueue(tag.id) + |> IndexJob.put_enqueue("Tags", :id, [tag.id]) + |> Multi.transact() + |> case do + {:ok, _changes} -> {:ok, tag} + error -> error + end end end @doc """ Removes the alias on the tag named by `slug`, on behalf of `actor`. - Write access, `:unalias` authorization, the relationship update, and the - audit log share one transaction. Image and tag reindexing runs after commit. + Write access, `:unalias` authorization, the relationship update, the + audit log, and enqueues for image and tag reindexing share one + transaction. ## Examples @@ -1203,9 +1194,12 @@ defmodule Philomena.Tags do Paths.tag_path(tag), "Dealiased tag '#{tag.name}'" ) - |> Multi.on_commit(fn %{locked_tag: %{aliased_tag: former_alias}, tag: tag} -> - reindex_tag_images(former_alias) - reindex_tags([tag, former_alias]) + |> Multi.merge(fn %{locked_tag: %{aliased_tag: former_alias}} -> + TagReindexJob.put_enqueue(Multi.new(), former_alias.id) + end) + |> IndexJob.put_enqueue("Tags", :id, fn + %{tag: tag, locked_tag: %{aliased_tag: former_alias}} -> + [tag.id, former_alias.id] end) |> Multi.transact() |> case do @@ -1253,8 +1247,8 @@ defmodule Philomena.Tags do `image_step` must resolve to an `%Image{}` with added and removed tag lists. Tags owns its image count update rule: hidden images do not contribute to tag - image counts. Counter updates are performed in ascending tag ID order and - affected tags are reindexed after the transaction commits. + image counts. Counter updates are performed in ascending tag ID order. + Affected tags have reindexing jobs enqueued in the same transaction. ## Examples @@ -1268,9 +1262,7 @@ defmodule Philomena.Tags do |> Multi.run(:image_tag_counts_tag_ids, fn repo, %{^image_step => image} -> {:ok, update_image_count_changes(repo, image)} end) - |> Multi.on_commit(fn %{image_tag_counts_tag_ids: tag_ids} -> - reindex_tag_ids(tag_ids) - end) + |> IndexJob.put_enqueue("Tags", :id, fn %{image_tag_counts_tag_ids: tag_ids} -> tag_ids end) end @doc """ @@ -1290,9 +1282,7 @@ defmodule Philomena.Tags do Multi.run(multi, step, fn repo, changes -> {:ok, update_image_counts(repo, diff, tag_ids_callback.(changes))} end) - |> Multi.on_commit(fn changes -> - reindex_tag_ids(tag_ids_callback.(changes)) - end) + |> IndexJob.put_enqueue("Tags", :id, tag_ids_callback) end @doc """ @@ -1300,7 +1290,7 @@ defmodule Philomena.Tags do Images supplies inserted and deleted tagging steps. This function updates image counts for visible images and queues every affected tag for - indexing after commit. `visible_image_step` must resolve to the images + indexing transactionally. `visible_image_step` must resolve to the images matching the `hidden_from_users == false` precondition. """ @spec put_batch_image_count_changes(Multi.t(), Multi.name(), Multi.name(), Multi.name()) :: @@ -1362,7 +1352,7 @@ defmodule Philomena.Tags do {:ok, Enum.map(rows, & &1.id)} end) - |> Multi.on_commit(fn %{batch_tag_counts: tag_ids} -> reindex_tag_ids(tag_ids) end) + |> IndexJob.put_enqueue("Tags", :id, fn %{batch_tag_counts: tag_ids} -> tag_ids end) end @doc """ @@ -1520,10 +1510,8 @@ defmodule Philomena.Tags do end, [] ) - |> Multi.on_commit(fn _changes -> - reindex_tag_images(target_tag) - reindex_tags([tag, target_tag]) - end) + |> TagReindexJob.put_enqueue(target_tag.id) + |> IndexJob.put_enqueue("Tags", :id, [tag.id, target_tag.id]) |> Multi.transact() |> case do {:ok, _changes} -> @@ -1569,7 +1557,7 @@ defmodule Philomena.Tags do end, [] ) - |> Multi.on_commit(fn _changes -> reindex_tags([tag]) end) + |> IndexJob.put_enqueue("Tags", :id, [tag.id]) |> Multi.transact_with_automatic_retry(isolation: :serializable) # Then, reindex. @@ -1645,22 +1633,6 @@ defmodule Philomena.Tags do end end - @doc """ - Queues a list of tags for search index updates. - Returns the list of tags unchanged, for use in a pipeline. - - ## Examples - - iex> reindex_tags([%Tag{}, %Tag{}, ...]) - [%Tag{}, %Tag{}, ...] - - """ - @spec reindex_tags([Tag.t()]) :: [Tag.t()] - def reindex_tags(tags) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Tags", "id", Enum.map(tags, & &1.id)]) - tags - end - @doc """ Returns the list of associations to preload for tag indexing. diff --git a/lib/philomena/user_name_changes.ex b/lib/philomena/user_name_changes.ex index 2765b84c6..57c0e7451 100644 --- a/lib/philomena/user_name_changes.ex +++ b/lib/philomena/user_name_changes.ex @@ -32,15 +32,15 @@ defmodule Philomena.UserNameChanges do ## Examples - iex> record_rename(Multi.new(), :name_change, user) + iex> record_rename(Multi.new(), :name_change, :locked_user) %Multi{} """ - @spec record_rename(Multi.t(), Multi.name(), User.t()) :: Multi.t() - def record_rename(%Multi{} = multi, step, %User{} = user) do - changeset = UserNameChange.changeset(%UserNameChange{user_id: user.id}, user.name) - - Multi.insert(multi, step, changeset) + @spec put_record_rename(Multi.t(), Multi.name(), Multi.name()) :: Multi.t() + def put_record_rename(%Multi{} = multi, step, locked_user_step) do + Multi.insert(multi, step, fn %{^locked_user_step => user} -> + UserNameChange.changeset(%UserNameChange{user_id: user.id}, user.name) + end) end @doc """ diff --git a/lib/philomena/user_statistics.ex b/lib/philomena/user_statistics.ex index b5219460d..75f15cbbb 100644 --- a/lib/philomena/user_statistics.ex +++ b/lib/philomena/user_statistics.ex @@ -7,7 +7,6 @@ defmodule Philomena.UserStatistics do """ alias Philomena.Multi - alias Philomena.Repo alias Philomena.Users alias Philomena.Users.User alias Philomena.UserStatistics.UserStatistic @@ -32,58 +31,12 @@ defmodule Philomena.UserStatistics do | :posts_count | :topics_count - defp persist_increment(user_id, statistic, amount) do - day = Date.utc_today() - - Repo.transact(fn -> - case Users.increment_counter(Repo, user_id, statistic, amount) do - {1, nil} -> - Repo.insert( - Map.put(%UserStatistic{day: day, user_id: user_id}, statistic, amount), - on_conflict: [inc: [{statistic, amount}]], - conflict_target: [:day, :user_id] - ) - - {0, nil} -> - {:error, :not_found} - end - end) - end - - defp reindex_result({:ok, %UserStatistic{}}, user_id) do - Users.reindex_user(%User{id: user_id}) - {:ok, nil} - end - - defp reindex_result(error, _user_id), do: error - - defp persist_bulk_increment(repo, user_ids, statistic, amount) do - case Users.increment_counters(repo, user_ids, statistic, amount) do - {count, nil} when count == length(user_ids) -> - entries = - Enum.map(user_ids, fn user_id -> - %{day: Date.utc_today(), user_id: user_id} - |> Map.put(statistic, amount) - end) - - repo.insert_all(UserStatistic, entries, - on_conflict: [inc: [{statistic, amount}]], - conflict_target: [:day, :user_id] - ) - - {:ok, nil} - - _ -> - {:error, :not_found} - end - end - @doc """ Adds an atomic statistic increment to `multi`. The Multi updates both the user's lifetime counter and current UTC-daily counter. Passing `nil` leaves the Multi unchanged, which supports anonymous - activity. After the transaction commits, it reindexes the user. + activity. The user reindex job is inserted with the transaction. ## Example @@ -118,12 +71,20 @@ defmodule Philomena.UserStatistics do def put_increment(multi, user_id, statistic, amount) when is_integer(user_id) and statistic in @permitted_actions and is_integer(amount) do + counter_step = {:put_increment_counter, make_ref()} + multi - |> Multi.run({:put_increment, make_ref()}, fn _repo, _changes -> - persist_increment(user_id, statistic, amount) - end) - |> Multi.on_commit(fn _changes -> - Users.reindex_user(%User{id: user_id}) + |> Users.put_increment_counter(counter_step, user_id, statistic, amount) + |> Multi.run({:put_increment, make_ref()}, fn repo, %{^counter_step => {count, nil}} -> + if count == 1 do + repo.insert( + Map.put(%UserStatistic{day: Date.utc_today(), user_id: user_id}, statistic, amount), + on_conflict: [inc: [{statistic, amount}]], + conflict_target: [:day, :user_id] + ) + else + {:error, :not_found} + end end) end @@ -163,53 +124,27 @@ defmodule Philomena.UserStatistics do end) |> Enum.uniq() - multi - |> Multi.run({:put_bulk_increment, make_ref()}, fn repo, _changes -> - persist_bulk_increment(repo, user_ids, statistic, amount) - end) - |> Multi.on_commit(fn _changes -> Users.reindex_user_ids(user_ids) end) - end - - @doc """ - Atomically increments one lifetime and UTC-daily statistic for `user_or_id`. - - A `nil` user is an intentional no-op for anonymous activity. A missing user - ID is `{:error, :not_found}`. Negative amounts decrement both counters. - Unknown statistic keys and non-integer amounts do not match this API. - - The database increments join an ambient transaction when called from an - `Ecto.Multi` callback, so an owning action rollback also rolls them back. A - successful call enqueues a user reindex; that queue side effect is best-effort - and is not part of the database transaction. + counter_step = {:put_bulk_increment_counter, make_ref()} - ## Examples - - iex> increment(user, :images_count) - {:ok, nil} - - iex> increment(user.id, :images_count, -1) - {:ok, nil} - - iex> increment(nil, :comments_count) - {:ok, nil} - - """ - @spec increment(User.t() | integer() | nil, statistic(), integer()) :: - {:ok, nil} | {:error, :not_found | Ecto.Changeset.t()} - def increment(user_or_id, statistic, amount \\ 1) - - def increment(nil, statistic, amount) - when statistic in @permitted_actions and is_integer(amount), - do: {:ok, nil} + multi + |> Users.put_increment_counters(counter_step, user_ids, statistic, amount) + |> Multi.run({:put_bulk_increment, make_ref()}, fn repo, %{^counter_step => {count, nil}} -> + if count == length(user_ids) do + entries = + Enum.map(user_ids, fn user_id -> + %{day: Date.utc_today(), user_id: user_id} + |> Map.put(statistic, amount) + end) - def increment(%User{} = user, statistic, amount) - when statistic in @permitted_actions and is_integer(amount), - do: increment(user.id, statistic, amount) + repo.insert_all(UserStatistic, entries, + on_conflict: [inc: [{statistic, amount}]], + conflict_target: [:day, :user_id] + ) - def increment(user_id, statistic, amount) - when is_integer(user_id) and statistic in @permitted_actions and is_integer(amount) do - user_id - |> persist_increment(statistic, amount) - |> reindex_result(user_id) + {:ok, nil} + else + {:error, :not_found} + end + end) end end diff --git a/lib/philomena/users.ex b/lib/philomena/users.ex index 21f9da660..167551fb0 100644 --- a/lib/philomena/users.ex +++ b/lib/philomena/users.ex @@ -52,12 +52,12 @@ defmodule Philomena.Users do alias Philomena.Filters.Filter alias Philomena.ModerationLogs alias Philomena.ModerationLogs.Paths - alias Philomena.IndexWorker + alias Philomena.Workers.IndexJob alias Philomena.Loader - alias Philomena.UserEraseWorker - alias Philomena.UserRenameWorker - alias Philomena.UserUnvoteWorker - alias Philomena.UserWipeWorker + alias Philomena.Workers.UserEraseJob + alias Philomena.Workers.UserRenameJob + alias Philomena.Workers.UserUnvoteJob + alias Philomena.Workers.UserWipeJob ## Shared locators @@ -230,34 +230,29 @@ defmodule Philomena.Users do end) end - ## Post-commit hooks + ## Transactional jobs defp put_reindex_user(multi) do - Multi.on_commit(multi, fn %{user: user} -> reindex_user(user) end) + IndexJob.put_enqueue(multi, "Users", :id, fn %{user: user} -> [user.id] end) end defp put_wipe_user_votes_job(multi, [{:upvotes_and_faves?, upvotes_and_faves?}]) do - Multi.on_commit(multi, fn %{user: user} -> - Exq.enqueue(Exq, "indexing", UserUnvoteWorker, [user.id, upvotes_and_faves?]) - end) + UserUnvoteJob.put_enqueue(multi, fn %{user: user} -> user.id end, upvotes_and_faves?) end defp put_wipe_user_job(multi) do - Multi.on_commit(multi, fn %{user: user} -> - Exq.enqueue(Exq, "indexing", UserWipeWorker, [user.id]) - end) + UserWipeJob.put_enqueue(multi, fn %{user: user} -> user.id end) end - defp put_rename_user_job(multi, [{:old_name, old_name}]) do - Multi.on_commit(multi, fn %{user: user} -> - Exq.enqueue(Exq, "indexing", UserRenameWorker, [old_name, user.name]) + defp put_rename_user_job(multi) do + UserRenameJob.put_enqueue(multi, fn + %{locked_user: %{name: old_name}, user: %{name: new_name}} -> + {old_name, new_name} end) end - defp put_erase_user_job(multi, %Actor{} = actor) do - Multi.on_commit(multi, fn %{user: user} -> - Exq.enqueue(Exq, "indexing", UserEraseWorker, [user.id, actor.user.id]) - end) + defp put_erase_user_job(multi, %User{} = target, %Actor{} = actor) do + UserEraseJob.put_enqueue(multi, target.id, actor.user.id) end ## Public reads @@ -1744,12 +1739,9 @@ defmodule Philomena.Users do {:ok, User.t()} | {:error, :ban | :unauthorized | Ecto.Changeset.t()} def update_name(%Actor{user: user} = actor, user_params) do with :ok <- verify_write_access(actor) do - old_name = user.name - Multi.new() |> Multi.lock_one(:locked_user, user_lock_query(user)) |> Multi.run(:authorize, fn _repo, %{locked_user: user} -> - # credo:disable-for-next-line with :ok <- authorize(user, :change_username, user) do {:ok, nil} end @@ -1757,9 +1749,9 @@ defmodule Philomena.Users do |> Multi.update(:user, fn %{locked_user: user} -> User.name_changeset(user, user_params) end) - |> UserNameChanges.record_rename(:name_change, user) + |> UserNameChanges.put_record_rename(:name_change, :locked_user) |> put_reindex_user() - |> put_rename_user_job(old_name: old_name) + |> put_rename_user_job() |> Multi.transact() |> case do {:ok, %{user: %User{} = user}} -> @@ -2192,7 +2184,7 @@ defmodule Philomena.Users do {"Admin.User.Erase:create", Paths.profile_path(user), "Erased #{original_name}"} end) |> put_reindex_user() - |> put_erase_user_job(actor) + |> put_erase_user_job(user, actor) |> Multi.transact() |> case do {:ok, %{user: %User{} = user}} -> @@ -2699,27 +2691,31 @@ defmodule Philomena.Users do @doc group: "Cross-context helpers" @doc """ - Increments one lifetime counter on a user through the supplied repository. - - The return value is the normal `update_all/3` row count. + Increments one lifetime counter on a user within `multi` and queues the + resulting user index update in the same transaction. """ - @spec increment_counter(module(), integer(), atom(), integer()) :: {non_neg_integer(), nil} - def increment_counter(repo, user_id, field, amount) + @spec put_increment_counter(Multi.t(), Multi.name(), integer(), atom(), integer()) :: Multi.t() + def put_increment_counter(%Multi{} = multi, step, user_id, field, amount) when is_integer(user_id) and is_atom(field) and is_integer(amount) do - repo.update_all(where(User, id: ^user_id), inc: [{field, amount}]) + multi + |> Multi.update_all(step, where(User, id: ^user_id), inc: [{field, amount}]) + |> IndexJob.put_enqueue("Users", :id, [user_id]) end @doc group: "Cross-context helpers" @doc """ - Increments one lifetime counter for each supplied user through the supplied - repository. - - The return value is the normal `update_all/3` row count. + Increments one lifetime counter for each supplied user within `multi` and + queues the resulting user index updates in the same transaction. """ - @spec increment_counters(module(), [integer()], atom(), integer()) :: {non_neg_integer(), nil} - def increment_counters(repo, user_ids, field, amount) + @spec put_increment_counters(Multi.t(), Multi.name(), [integer()], atom(), integer()) :: + Multi.t() + def put_increment_counters(%Multi{} = multi, _step, [], _field, _amount), do: multi + + def put_increment_counters(%Multi{} = multi, step, user_ids, field, amount) when is_list(user_ids) and is_atom(field) and is_integer(amount) do - repo.update_all(where(User, [user], user.id in ^user_ids), inc: [{field, amount}]) + multi + |> Multi.update_all(step, where(User, [user], user.id in ^user_ids), inc: [{field, amount}]) + |> IndexJob.put_enqueue("Users", :id, user_ids) end @doc group: "Cross-context helpers" @@ -2743,7 +2739,7 @@ defmodule Philomena.Users do :ok """ - @spec perform_rename(String.t(), String.t()) :: term() + @spec perform_rename(String.t(), String.t()) :: :ok def perform_rename(old_name, new_name) do Images.user_name_reindex(old_name, new_name) Comments.user_name_reindex(old_name, new_name) @@ -2753,42 +2749,8 @@ defmodule Philomena.Users do Filters.user_name_reindex(old_name, new_name) TagChanges.user_name_reindex(old_name, new_name) Users.user_name_reindex(old_name, new_name) - end - - @doc group: "Background jobs" - @doc """ - Queues a single user for search index updates. - Returns the user struct unchanged, for use in a pipeline. - - ## Examples - - iex> reindex_user(user) - %User{} - - """ - @spec reindex_user(User.t()) :: User.t() - def reindex_user(%User{} = user) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Users", "id", [user.id]]) - - user - end - @doc group: "Background jobs" - @doc """ - Queues a list of user IDs for search index updates. - Returns the list unchanged, for use in a pipeline. - - ## Examples - - iex> reindex_user_ids([1, 2, 3]) - [1, 2, 3] - - """ - @spec reindex_user_ids(list(integer())) :: list(integer()) - def reindex_user_ids(user_ids) do - Exq.enqueue(Exq, "indexing", IndexWorker, ["Users", "id", user_ids]) - - user_ids + :ok end @doc group: "Background jobs" diff --git a/lib/philomena/users/user_downvote_wipe.ex b/lib/philomena/users/user_downvote_wipe.ex index 955a89eab..f6b8f2f65 100644 --- a/lib/philomena/users/user_downvote_wipe.ex +++ b/lib/philomena/users/user_downvote_wipe.ex @@ -3,13 +3,14 @@ defmodule Philomena.Users.UserDownvoteWipe do Performs the asynchronous vote/favorite cleanup owned by the Users context. The public entry point accepts only a trusted persisted user ID and is called - by `Philomena.UserUnvoteWorker` after an authorized Users service enqueues it. + by `Philomena.Workers.UserUnvoteJob` after an authorized Users service enqueues it. """ import Ecto.Query alias PhilomenaQuery.Search alias Philomena.Users + alias Philomena.Users.User alias Philomena.Images.Image alias Philomena.Images alias Philomena.ImageVotes @@ -47,18 +48,18 @@ defmodule Philomena.Users.UserDownvoteWipe do {count, image_ids} = ImageVotes.delete_user_votes!(user.id, false) Images.decrement_vote_counters!(image_ids, false) - Users.increment_counter(Repo, user.id, :image_votes_count, -count) + Repo.update_all(where(User, id: ^user.id), inc: [image_votes_count: -count]) reindex(image_ids) if upvotes_and_faves_too do {count, image_ids} = ImageVotes.delete_user_votes!(user.id, true) Images.decrement_vote_counters!(image_ids, true) - Users.increment_counter(Repo, user.id, :image_votes_count, -count) + Repo.update_all(where(User, id: ^user.id), inc: [image_votes_count: -count]) reindex(image_ids) {count, image_ids} = ImageFaves.delete_user_faves!(user.id) Images.decrement_fave_counters!(image_ids) - Users.increment_counter(Repo, user.id, :image_faves_count, -count) + Repo.update_all(where(User, id: ^user.id), inc: [image_faves_count: -count]) reindex(image_ids) end diff --git a/lib/philomena/users/user_wipe.ex b/lib/philomena/users/user_wipe.ex index 941cc88bd..872aed68f 100644 --- a/lib/philomena/users/user_wipe.ex +++ b/lib/philomena/users/user_wipe.ex @@ -4,7 +4,7 @@ defmodule Philomena.Users.UserWipe do the Users context. The public entry point accepts only a trusted persisted user ID and is called - by `Philomena.UserWipeWorker` after an authorized Users service enqueues it. + by `Philomena.Workers.UserWipeJob` after an authorized Users service enqueues it. """ alias Philomena.Comments @@ -16,7 +16,6 @@ defmodule Philomena.Users.UserWipe do alias Philomena.UserIps alias Philomena.UserFingerprints alias Philomena.Users - alias Philomena.Users.User @wipe_ip %Postgrex.INET{address: {127, 0, 1, 1}, netmask: 32} @wipe_fp "ffff" @@ -30,9 +29,10 @@ defmodule Philomena.Users.UserWipe do ## Examples iex> UserWipe.perform(user.id) - %User{} + :ok + """ - @spec perform(integer()) :: User.t() + @spec perform(integer()) :: :ok def perform(user_id) do user = Users.fetch_user_for_worker!(user_id) @@ -48,6 +48,8 @@ defmodule Philomena.Users.UserWipe do UserFingerprints.delete_for_user!(user.id) Users.replace_email_for_wipe!(user.id, "deactivated#{random_hex}@example.com") - Users.reindex_user(user) + Users.perform_reindex(:id, [user.id]) + + :ok end end diff --git a/lib/philomena/workers/image_purge_job.ex b/lib/philomena/workers/image_purge_job.ex new file mode 100644 index 000000000..59651d0b6 --- /dev/null +++ b/lib/philomena/workers/image_purge_job.ex @@ -0,0 +1,24 @@ +defmodule Philomena.Workers.ImagePurgeJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + + alias Philomena.Images + alias Philomena.Multi + + @spec put_enqueue(Multi.t(), (Multi.changes() -> [String.t()])) :: Multi.t() + def put_enqueue(%Multi{} = multi, files) when is_function(files, 1) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, fn changes -> + %{files: files.(changes)} + end) + end + + @spec enqueue([String.t()]) :: any() + def enqueue(files) do + # Legacy variation for thumbnailer. + Philomena.JobQueue.insert(__MODULE__, %{files: files}) + end + + @impl Oban.Worker + def perform(%Oban.Job{args: %{"files" => files}}) do + Images.perform_purge(files) + end +end diff --git a/lib/philomena/workers/image_purge_worker.ex b/lib/philomena/workers/image_purge_worker.ex deleted file mode 100644 index 2ab2f17c6..000000000 --- a/lib/philomena/workers/image_purge_worker.ex +++ /dev/null @@ -1,7 +0,0 @@ -defmodule Philomena.ImagePurgeWorker do - alias Philomena.Images - - def perform(files) do - Images.perform_purge(files) - end -end diff --git a/lib/philomena/workers/index_job.ex b/lib/philomena/workers/index_job.ex new file mode 100644 index 000000000..556cd6e3e --- /dev/null +++ b/lib/philomena/workers/index_job.ex @@ -0,0 +1,47 @@ +defmodule Philomena.Workers.IndexJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + + alias Philomena.Multi + + @modules %{ + "Comments" => Philomena.Comments, + "Galleries" => Philomena.Galleries, + "Images" => Philomena.Images, + "Posts" => Philomena.Posts, + "Reports" => Philomena.Reports, + "Tags" => Philomena.Tags, + "Filters" => Philomena.Filters, + "TagChanges" => Philomena.TagChanges, + "Users" => Philomena.Users + } + + @spec put_enqueue(Multi.t(), String.t(), atom(), [integer()] | (Multi.changes() -> [integer()])) :: + Multi.t() + def put_enqueue(%Multi{} = multi, module, column, condition) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, fn changes -> + condition = if is_function(condition, 1), do: condition.(changes), else: condition + + %{ + module: module, + column: to_string(column), + condition: condition + } + end) + end + + # Perform the queued index. Context function looks like the following: + # + # def perform_reindex(column, condition) do + # Image + # |> preload(^indexing_preloads()) + # |> where([i], field(i, ^column) in ^condition) + # |> Search.reindex(Image) + # end + # + @impl Oban.Worker + def perform(%Oban.Job{ + args: %{"module" => module, "column" => column, "condition" => condition} + }) do + @modules[module].perform_reindex(String.to_existing_atom(column), condition) + end +end diff --git a/lib/philomena/workers/index_worker.ex b/lib/philomena/workers/index_worker.ex deleted file mode 100644 index cce9ccc65..000000000 --- a/lib/philomena/workers/index_worker.ex +++ /dev/null @@ -1,26 +0,0 @@ -defmodule Philomena.IndexWorker do - @modules %{ - "Comments" => Philomena.Comments, - "Galleries" => Philomena.Galleries, - "Images" => Philomena.Images, - "Posts" => Philomena.Posts, - "Reports" => Philomena.Reports, - "Tags" => Philomena.Tags, - "Filters" => Philomena.Filters, - "TagChanges" => Philomena.TagChanges, - "Users" => Philomena.Users - } - - # Perform the queued index. Context function looks like the following: - # - # def perform_reindex(column, condition) do - # Image - # |> preload(^indexing_preloads()) - # |> where([i], field(i, ^column) in ^condition) - # |> Search.reindex(Image) - # end - # - def perform(module, column, condition) do - @modules[module].perform_reindex(String.to_existing_atom(column), condition) - end -end diff --git a/lib/philomena/workers/tag_alias_job.ex b/lib/philomena/workers/tag_alias_job.ex new file mode 100644 index 000000000..cf6502bef --- /dev/null +++ b/lib/philomena/workers/tag_alias_job.ex @@ -0,0 +1,28 @@ +defmodule Philomena.Workers.TagAliasJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + + alias Philomena.Tags + alias Philomena.Multi + + @spec put_enqueue(Multi.t(), function(), function()) :: Multi.t() + def put_enqueue(%Multi{} = multi, tag_id, target_tag_id) + when is_function(tag_id, 1) and is_function(target_tag_id, 1) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, fn changes -> + %{ + tag_id: tag_id.(changes), + target_tag_id: target_tag_id.(changes) + } + end) + end + + @impl Oban.Worker + def perform(%Oban.Job{args: %{"tag_id" => tag_id, "target_tag_id" => target_tag_id}}) do + case Tags.perform_alias(tag_id, target_tag_id) do + {:error, :stale_target} -> + {:cancel, :stale_target} + + result -> + result + end + end +end diff --git a/lib/philomena/workers/tag_alias_worker.ex b/lib/philomena/workers/tag_alias_worker.ex deleted file mode 100644 index 48639946b..000000000 --- a/lib/philomena/workers/tag_alias_worker.ex +++ /dev/null @@ -1,7 +0,0 @@ -defmodule Philomena.TagAliasWorker do - alias Philomena.Tags - - def perform(tag_id, target_tag_id) do - Tags.perform_alias(tag_id, target_tag_id) - end -end diff --git a/lib/philomena/workers/tag_change_revert_worker.ex b/lib/philomena/workers/tag_change_revert_job.ex similarity index 51% rename from lib/philomena/workers/tag_change_revert_worker.ex rename to lib/philomena/workers/tag_change_revert_job.ex index f4dbc8303..2f40aa74c 100644 --- a/lib/philomena/workers/tag_change_revert_worker.ex +++ b/lib/philomena/workers/tag_change_revert_job.ex @@ -1,4 +1,6 @@ -defmodule Philomena.TagChangeRevertWorker do +defmodule Philomena.Workers.TagChangeRevertJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + @moduledoc """ Reverts every tag change made by a user, IP, or fingerprint, batching by image so each image's tag history is reverted in a single operation. @@ -6,21 +8,33 @@ defmodule Philomena.TagChangeRevertWorker do alias Philomena.TagChanges alias Philomena.TagChanges.TagChange + alias Philomena.Multi import Ecto.Query - def perform(%{"user_id" => user_id, "attributes" => attributes}) do + @spec put_enqueue(Multi.t(), map(), map()) :: Multi.t() + def put_enqueue(%Multi{} = multi, target, attributes) + when is_map(target) and is_map(attributes) do + args = Map.put(target, :attributes, attributes) + + Philomena.JobQueue.put_enqueue(multi, __MODULE__, args) + end + + @impl Oban.Worker + def perform(job) + + def perform(%Oban.Job{args: %{"user_id" => user_id, "attributes" => attributes}}) do TagChange |> where(user_id: ^user_id) |> revert_all(attributes) end - def perform(%{"ip" => ip, "attributes" => attributes}) do + def perform(%Oban.Job{args: %{"ip" => ip, "attributes" => attributes}}) do TagChange |> where(ip: ^ip) |> revert_all(attributes) end - def perform(%{"fingerprint" => fp, "attributes" => attributes}) do + def perform(%Oban.Job{args: %{"fingerprint" => fp, "attributes" => attributes}}) do TagChange |> where(fingerprint: ^fp) |> revert_all(attributes) @@ -29,10 +43,7 @@ defmodule Philomena.TagChangeRevertWorker do defp revert_all(queryable, attributes) do attributes = cast_ip(atomify_keys(attributes)) - case TagChanges.revert_all_for_worker(queryable, attributes) do - :ok -> :ok - {:error, reason} -> raise "tag change batch revert failed: #{inspect(reason)}" - end + TagChanges.revert_all_for_worker(queryable, attributes) end defp atomify_keys(map) do diff --git a/lib/philomena/workers/tag_delete_job.ex b/lib/philomena/workers/tag_delete_job.ex new file mode 100644 index 000000000..6a9749e90 --- /dev/null +++ b/lib/philomena/workers/tag_delete_job.ex @@ -0,0 +1,16 @@ +defmodule Philomena.Workers.TagDeleteJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + + alias Philomena.Tags + alias Philomena.Multi + + @spec put_enqueue(Multi.t(), integer()) :: Multi.t() + def put_enqueue(%Multi{} = multi, tag_id) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, %{tag_id: tag_id}) + end + + @impl Oban.Worker + def perform(%Oban.Job{args: %{"tag_id" => tag_id}}) do + Tags.perform_delete(tag_id) + end +end diff --git a/lib/philomena/workers/tag_delete_worker.ex b/lib/philomena/workers/tag_delete_worker.ex deleted file mode 100644 index c20c80b0f..000000000 --- a/lib/philomena/workers/tag_delete_worker.ex +++ /dev/null @@ -1,7 +0,0 @@ -defmodule Philomena.TagDeleteWorker do - alias Philomena.Tags - - def perform(tag_id) do - Tags.perform_delete(tag_id) - end -end diff --git a/lib/philomena/workers/tag_reindex_job.ex b/lib/philomena/workers/tag_reindex_job.ex new file mode 100644 index 000000000..2c461edfb --- /dev/null +++ b/lib/philomena/workers/tag_reindex_job.ex @@ -0,0 +1,16 @@ +defmodule Philomena.Workers.TagReindexJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + + alias Philomena.Tags + alias Philomena.Multi + + @spec put_enqueue(Multi.t(), integer()) :: Multi.t() + def put_enqueue(%Multi{} = multi, tag_id) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, %{tag_id: tag_id}) + end + + @impl Oban.Worker + def perform(%Oban.Job{args: %{"tag_id" => tag_id}}) do + Tags.perform_reindex_images(tag_id) + end +end diff --git a/lib/philomena/workers/tag_reindex_worker.ex b/lib/philomena/workers/tag_reindex_worker.ex deleted file mode 100644 index d050bb469..000000000 --- a/lib/philomena/workers/tag_reindex_worker.ex +++ /dev/null @@ -1,7 +0,0 @@ -defmodule Philomena.TagReindexWorker do - alias Philomena.Tags - - def perform(tag_id) do - Tags.perform_reindex_images(tag_id) - end -end diff --git a/lib/philomena/workers/thumbnail_job.ex b/lib/philomena/workers/thumbnail_job.ex new file mode 100644 index 000000000..a9a6790c0 --- /dev/null +++ b/lib/philomena/workers/thumbnail_job.ex @@ -0,0 +1,29 @@ +defmodule Philomena.Workers.ThumbnailJob do + use Oban.Worker, queue: :images, max_attempts: 5 + + alias Philomena.Images.Thumbnailer + alias Philomena.Multi + + @spec put_enqueue(Multi.t(), integer(), String.t()) :: Multi.t() + def put_enqueue(%Multi{} = multi, image_id, image_mime_type) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, %{image_id: image_id}, + queue: queue(image_mime_type) + ) + end + + defp queue("video/webm"), do: :videos + defp queue(_mime_type), do: :images + + @impl Oban.Worker + def perform(%Oban.Job{args: %{"image_id" => image_id}}) do + Thumbnailer.generate_thumbnails(image_id) + + PhilomenaWeb.Endpoint.broadcast!( + "firehose", + "image:process", + %{image_id: image_id} + ) + + :ok + end +end diff --git a/lib/philomena/workers/thumbnail_worker.ex b/lib/philomena/workers/thumbnail_worker.ex deleted file mode 100644 index 225023664..000000000 --- a/lib/philomena/workers/thumbnail_worker.ex +++ /dev/null @@ -1,18 +0,0 @@ -defmodule Philomena.ThumbnailWorker do - alias Philomena.Images.Thumbnailer - alias Philomena.Images - - def perform(image_id) do - Thumbnailer.generate_thumbnails(image_id) - - PhilomenaWeb.Endpoint.broadcast!( - "firehose", - "image:process", - %{image_id: image_id} - ) - - image_id - |> Images.load_image_for_reindex!() - |> Images.reindex_image() - end -end diff --git a/lib/philomena/workers/user_erase_job.ex b/lib/philomena/workers/user_erase_job.ex new file mode 100644 index 000000000..9e6e2bf6d --- /dev/null +++ b/lib/philomena/workers/user_erase_job.ex @@ -0,0 +1,23 @@ +defmodule Philomena.Workers.UserEraseJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + + alias Philomena.Users.Eraser + alias Philomena.Users + alias Philomena.Multi + + @spec put_enqueue(Multi.t(), integer(), integer()) :: Multi.t() + def put_enqueue(%Multi{} = multi, user_id, moderator_id) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, %{ + user_id: user_id, + moderator_id: moderator_id + }) + end + + @impl Oban.Worker + def perform(%Oban.Job{args: %{"user_id" => user_id, "moderator_id" => moderator_id}}) do + moderator = Users.fetch_user_for_erase!(moderator_id) + user = Users.fetch_user_for_worker!(user_id) + + Eraser.erase_permanently!(user, moderator) + end +end diff --git a/lib/philomena/workers/user_erase_worker.ex b/lib/philomena/workers/user_erase_worker.ex deleted file mode 100644 index 895058e93..000000000 --- a/lib/philomena/workers/user_erase_worker.ex +++ /dev/null @@ -1,11 +0,0 @@ -defmodule Philomena.UserEraseWorker do - alias Philomena.Users.Eraser - alias Philomena.Users - - def perform(user_id, moderator_id) do - moderator = Users.fetch_user_for_erase!(moderator_id) - user = Users.fetch_user_for_worker!(user_id) - - Eraser.erase_permanently!(user, moderator) - end -end diff --git a/lib/philomena/workers/user_rename_job.ex b/lib/philomena/workers/user_rename_job.ex new file mode 100644 index 000000000..af0a2da5f --- /dev/null +++ b/lib/philomena/workers/user_rename_job.ex @@ -0,0 +1,23 @@ +defmodule Philomena.Workers.UserRenameJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + + alias Philomena.Users + alias Philomena.Multi + + @spec put_enqueue(Multi.t(), (Multi.changes() -> {String.t(), String.t()})) :: Multi.t() + def put_enqueue(%Multi{} = multi, old_and_new_names) when is_function(old_and_new_names, 1) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, fn changes -> + {old_name, new_name} = old_and_new_names.(changes) + + %{ + old_name: old_name, + new_name: new_name + } + end) + end + + @impl Oban.Worker + def perform(%Oban.Job{args: %{"old_name" => old_name, "new_name" => new_name}}) do + Users.perform_rename(old_name, new_name) + end +end diff --git a/lib/philomena/workers/user_rename_worker.ex b/lib/philomena/workers/user_rename_worker.ex deleted file mode 100644 index f79657a76..000000000 --- a/lib/philomena/workers/user_rename_worker.ex +++ /dev/null @@ -1,7 +0,0 @@ -defmodule Philomena.UserRenameWorker do - alias Philomena.Users - - def perform(old_name, new_name) do - Users.perform_rename(old_name, new_name) - end -end diff --git a/lib/philomena/workers/user_unvote_job.ex b/lib/philomena/workers/user_unvote_job.ex new file mode 100644 index 000000000..227e73b4f --- /dev/null +++ b/lib/philomena/workers/user_unvote_job.ex @@ -0,0 +1,24 @@ +defmodule Philomena.Workers.UserUnvoteJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + + alias Philomena.Users.UserDownvoteWipe + alias Philomena.Multi + + @spec put_enqueue(Multi.t(), function(), boolean()) :: Multi.t() + def put_enqueue(%Multi{} = multi, user_id, votes_and_faves_too?) + when is_function(user_id, 1) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, fn changes -> + %{ + user_id: user_id.(changes), + votes_and_faves_too?: votes_and_faves_too? + } + end) + end + + @impl Oban.Worker + def perform(%Oban.Job{ + args: %{"user_id" => user_id, "votes_and_faves_too?" => votes_and_faves_too?} + }) do + UserDownvoteWipe.perform(user_id, votes_and_faves_too?) + end +end diff --git a/lib/philomena/workers/user_unvote_worker.ex b/lib/philomena/workers/user_unvote_worker.ex deleted file mode 100644 index 8c8f8cb23..000000000 --- a/lib/philomena/workers/user_unvote_worker.ex +++ /dev/null @@ -1,7 +0,0 @@ -defmodule Philomena.UserUnvoteWorker do - alias Philomena.Users.UserDownvoteWipe - - def perform(user_id, votes_and_faves_too?) do - UserDownvoteWipe.perform(user_id, votes_and_faves_too?) - end -end diff --git a/lib/philomena/workers/user_wipe_job.ex b/lib/philomena/workers/user_wipe_job.ex new file mode 100644 index 000000000..067056ba4 --- /dev/null +++ b/lib/philomena/workers/user_wipe_job.ex @@ -0,0 +1,18 @@ +defmodule Philomena.Workers.UserWipeJob do + use Oban.Worker, queue: :indexing, max_attempts: 5 + + alias Philomena.Users.UserWipe + alias Philomena.Multi + + @spec put_enqueue(Multi.t(), function()) :: Multi.t() + def put_enqueue(%Multi{} = multi, user_id) when is_function(user_id, 1) do + Philomena.JobQueue.put_enqueue(multi, __MODULE__, fn changes -> + %{user_id: user_id.(changes)} + end) + end + + @impl Oban.Worker + def perform(%Oban.Job{args: %{"user_id" => user_id}}) do + UserWipe.perform(user_id) + end +end diff --git a/lib/philomena/workers/user_wipe_worker.ex b/lib/philomena/workers/user_wipe_worker.ex deleted file mode 100644 index c4f2475a1..000000000 --- a/lib/philomena/workers/user_wipe_worker.ex +++ /dev/null @@ -1,7 +0,0 @@ -defmodule Philomena.UserWipeWorker do - alias Philomena.Users.UserWipe - - def perform(user_id) do - UserWipe.perform(user_id) - end -end diff --git a/mix.exs b/mix.exs index 0bc7fee02..54be168cd 100644 --- a/mix.exs +++ b/mix.exs @@ -69,7 +69,8 @@ defmodule Philomena.MixProject do {:remote_ip, "~> 1.2"}, {:briefly, "~> 0.5"}, {:req, "~> 0.7.4"}, - {:exq, "~> 0.21"}, + {:elixir_uuid, "~> 1.2"}, + {:oban, "~> 2.18"}, {:ex_aws, "~> 2.6"}, {:ex_aws_s3, "~> 2.5"}, {:sweet_xml, "~> 0.7"}, diff --git a/mix.lock b/mix.lock index 4665977c5..8be29c778 100644 --- a/mix.lock +++ b/mix.lock @@ -23,7 +23,6 @@ "ex_aws_s3": {:hex, :ex_aws_s3, "2.5.9", "862b7792f2e60d7010e2920d79964e3fab289bc0fd951b0ba8457a3f7f9d1199", [:mix], [{:ex_aws, "~> 2.0", [hex: :ex_aws, repo: "hexpm", optional: false]}, {:sweet_xml, ">= 0.0.0", [hex: :sweet_xml, repo: "hexpm", optional: true]}], "hexpm", "a480d2bb2da64610014021629800e1e9457ca5e4a62f6775bffd963360c2bf90"}, "ex_doc": {:hex, :ex_doc, "0.40.4", "66f2e42bf588594d5a8aab31cad87f2ddad09d0da1b1a2f379340ec2c2e497cb", [:mix], [{:earmark_parser, "~> 1.4.46", [hex: :earmark_parser, repo: "hexpm", optional: false]}, {:makeup_c, ">= 0.1.0", [hex: :makeup_c, repo: "hexpm", optional: true]}, {:makeup_elixir, "~> 0.14 or ~> 1.0", [hex: :makeup_elixir, repo: "hexpm", optional: false]}, {:makeup_erlang, "~> 0.1 or ~> 1.0", [hex: :makeup_erlang, repo: "hexpm", optional: false]}, {:makeup_html, ">= 0.1.0", [hex: :makeup_html, repo: "hexpm", optional: true]}], "hexpm", "6222b9e423d76584ee34df2c82a5ed72c2d53dc153f7f483ad28b378694186cc"}, "expo": {:hex, :expo, "1.1.1", "4202e1d2ca6e2b3b63e02f69cfe0a404f77702b041d02b58597c00992b601db5", [:mix], [], "hexpm", "5fb308b9cb359ae200b7e23d37c76978673aa1b06e2b3075d814ce12c5811640"}, - "exq": {:hex, :exq, "0.24.0", "326a66708144496bbc4e2745c1171eeb3f9ac5140969d9724dd6db35386f4b15", [:mix], [{:elixir_uuid, ">= 1.2.0", [hex: :elixir_uuid, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:poison, ">= 1.2.0 and < 7.0.0", [hex: :poison, repo: "hexpm", optional: true]}, {:redix, ">= 0.9.0", [hex: :redix, repo: "hexpm", optional: false]}], "hexpm", "7a6849408c78ba889f540a5930a80eb7ff2ca000b1ea7ac6581a837d08591a07"}, "file_system": {:hex, :file_system, "1.1.1", "31864f4685b0148f25bd3fbef2b1228457c0c89024ad67f7a81a3ffbc0bbad3a", [:mix], [], "hexpm", "7a15ff97dfe526aeefb090a7a9d3d03aa907e100e262a0f8f7746b78f8f87a5d"}, "finch": {:hex, :finch, "0.23.0", "e3f9287ac25a8832f848b144c2b57346aac65b205e2e0629a52adfe6507fd837", [:mix], [{:mime, "~> 1.0 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:mint, "~> 1.8", [hex: :mint, repo: "hexpm", optional: false]}, {:nimble_options, "~> 0.4 or ~> 1.0", [hex: :nimble_options, repo: "hexpm", optional: false]}, {:nimble_pool, "~> 1.1", [hex: :nimble_pool, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "80e58d3f936f57e3fdf404f83a3642897ae6d9fb642934e46da4d8fe761b99d5"}, "gettext": {:hex, :gettext, "1.0.2", "5457e1fd3f4abe47b0e13ff85086aabae760497a3497909b8473e0acee57673b", [:mix], [{:expo, "~> 0.5.1 or ~> 1.0", [hex: :expo, repo: "hexpm", optional: false]}], "hexpm", "eab805501886802071ad290714515c8c4a17196ea76e5afc9d06ca85fb1bfeb3"}, @@ -43,6 +42,7 @@ "nimble_options": {:hex, :nimble_options, "1.1.1", "e3a492d54d85fc3fd7c5baf411d9d2852922f66e69476317787a7b2bb000a61b", [:mix], [], "hexpm", "821b2470ca9442c4b6984882fe9bb0389371b8ddec4d45a9504f00a66f650b44"}, "nimble_parsec": {:hex, :nimble_parsec, "1.4.2", "8efba0122db06df95bfaa78f791344a89352ba04baedd3849593bfce4d0dc1c6", [:mix], [], "hexpm", "4b21398942dda052b403bbe1da991ccd03a053668d147d53fb8c4e0efe09c973"}, "nimble_pool": {:hex, :nimble_pool, "1.1.0", "bf9c29fbdcba3564a8b800d1eeb5a3c58f36e1e11d7b7fb2e084a643f645f06b", [:mix], [], "hexpm", "af2e4e6b34197db81f7aad230c1118eac993acc0dae6bc83bac0126d4ae0813a"}, + "oban": {:hex, :oban, "2.24.1", "2a609c54697ad2c44ba339df30491df2a40eda0758c95b5b879426e4e478bd1f", [:mix], [{:ecto_sql, "~> 3.10", [hex: :ecto_sql, repo: "hexpm", optional: false]}, {:ecto_sqlite3, "~> 0.9", [hex: :ecto_sqlite3, repo: "hexpm", optional: true]}, {:igniter, "~> 0.5", [hex: :igniter, repo: "hexpm", optional: true]}, {:jason, "~> 1.1", [hex: :jason, repo: "hexpm", optional: true]}, {:myxql, "~> 0.7", [hex: :myxql, repo: "hexpm", optional: true]}, {:postgrex, "~> 0.20", [hex: :postgrex, repo: "hexpm", optional: true]}, {:telemetry, "~> 1.3", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "ef8482472cf198554400b7f8e36a0ffee75c3a64de425d5c7ee625d271925ac7"}, "patch": {:hex, :patch, "0.16.0", "dfa2e1c381959d821ab1ddf4acd55f42d1c798003f0fef270154c0c5bd139b08", [:mix], [], "hexpm", "50e06ef77a9b4987edfc717a9bf047429f3d69af397a42d9ff311f27ccb23408"}, "pbkdf2": {:git, "https://github.com/basho/erlang-pbkdf2.git", "7e9bd5fcd3cc3062159e4c9214bb628aa6feb5ca", [ref: "7e9bd5fcd3cc3062159e4c9214bb628aa6feb5ca"]}, "phoenix": {:hex, :phoenix, "1.8.13", "e33192826d9bed4022bdb3f5a7b36c04362049d9390fd1b581d5ad6261779268", [:mix], [{:bandit, "~> 1.0", [hex: :bandit, repo: "hexpm", optional: true]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:phoenix_pubsub, "~> 2.1", [hex: :phoenix_pubsub, repo: "hexpm", optional: false]}, {:phoenix_template, "~> 1.0", [hex: :phoenix_template, repo: "hexpm", optional: false]}, {:phoenix_view, "~> 2.0", [hex: :phoenix_view, repo: "hexpm", optional: true]}, {:plug, "~> 1.14", [hex: :plug, repo: "hexpm", optional: false]}, {:plug_cowboy, "~> 2.7", [hex: :plug_cowboy, repo: "hexpm", optional: true]}, {:plug_crypto, "~> 2.2", [hex: :plug_crypto, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}, {:websock_adapter, "~> 0.5", [hex: :websock_adapter, repo: "hexpm", optional: false]}], "hexpm", "ad14e24d10e5a52d5f80429053bbe3a5d124311a2868fceb0a01a2e859c44539"}, diff --git a/priv/repo/migrations/20260905193823_add_oban.exs b/priv/repo/migrations/20260905193823_add_oban.exs new file mode 100644 index 000000000..48ac23c06 --- /dev/null +++ b/priv/repo/migrations/20260905193823_add_oban.exs @@ -0,0 +1,11 @@ +defmodule Philomena.Repo.Migrations.AddOban do + use Ecto.Migration + + def up do + Oban.Migration.up() + end + + def down do + Oban.Migration.down(version: 1) + end +end diff --git a/priv/repo/seeds.exs b/priv/repo/seeds.exs index 7d3bd6b9c..4c4d76d4b 100644 --- a/priv/repo/seeds.exs +++ b/priv/repo/seeds.exs @@ -23,9 +23,9 @@ alias Philomena.{ } alias PhilomenaQuery.Search +alias Philomena.Filters alias Philomena.Users alias Philomena.Tags -alias Philomena.Filters import Ecto.Query IO.puts("---- Creating search indices") @@ -62,13 +62,6 @@ for filter_def <- resources["system_filters"] do } ) |> Repo.insert(on_conflict: :nothing) - |> case do - {:ok, filter} -> - Filters.reindex_filter(filter) - - {:error, changeset} -> - IO.inspect(changeset.errors) - end end IO.puts("---- Generating forums") @@ -149,6 +142,7 @@ for rule_def <- resources["rules"] do end IO.puts("---- Indexing content") +Search.reindex(Filter |> preload(^Filters.indexing_preloads()), Filter) Search.reindex(Tag |> preload(^Tags.indexing_preloads()), Tag) IO.puts("---- Done.") diff --git a/priv/repo/structure.sql b/priv/repo/structure.sql index 8dd8befad..453fe3a44 100644 --- a/priv/repo/structure.sql +++ b/priv/repo/structure.sql @@ -2,7 +2,7 @@ -- PostgreSQL database dump -- -\restrict O6OYrjIYDwvGhgBIGh2SR1j95V8aVnMGs9WcuKijctXbT5ypaF3nbIkiah3iFpk +\restrict DiOFiU8d5t1AmxDoKBF0A1E487owRGbnVYjRKkg1IFCSL1or1ZGLi8wPipZ954H -- Dumped from database version 18.4 -- Dumped by pg_dump version 18.4 @@ -33,6 +33,22 @@ CREATE EXTENSION IF NOT EXISTS citext WITH SCHEMA public; COMMENT ON EXTENSION citext IS 'data type for case-insensitive character strings'; +-- +-- Name: oban_job_state; Type: TYPE; Schema: public; Owner: - +-- + +CREATE TYPE public.oban_job_state AS ENUM ( + 'available', + 'suspended', + 'scheduled', + 'executing', + 'retryable', + 'completed', + 'discarded', + 'cancelled' +); + + SET default_tablespace = ''; SET default_table_access_method = heap; @@ -1188,6 +1204,74 @@ CREATE SEQUENCE public.notifications_id_seq ALTER SEQUENCE public.notifications_id_seq OWNED BY public.notifications.id; +-- +-- Name: oban_jobs; Type: TABLE; Schema: public; Owner: - +-- + +CREATE TABLE public.oban_jobs ( + id bigint NOT NULL, + state public.oban_job_state DEFAULT 'available'::public.oban_job_state NOT NULL, + queue text DEFAULT 'default'::text NOT NULL, + worker text NOT NULL, + args jsonb DEFAULT '{}'::jsonb NOT NULL, + errors jsonb[] DEFAULT ARRAY[]::jsonb[] NOT NULL, + attempt integer DEFAULT 0 NOT NULL, + max_attempts integer DEFAULT 20 NOT NULL, + inserted_at timestamp without time zone DEFAULT timezone('UTC'::text, now()) NOT NULL, + scheduled_at timestamp without time zone DEFAULT timezone('UTC'::text, now()) NOT NULL, + attempted_at timestamp without time zone, + completed_at timestamp without time zone, + attempted_by text[], + discarded_at timestamp without time zone, + priority integer DEFAULT 0 NOT NULL, + tags text[] DEFAULT ARRAY[]::text[], + meta jsonb DEFAULT '{}'::jsonb, + cancelled_at timestamp without time zone, + CONSTRAINT attempt_range CHECK (((attempt >= 0) AND (attempt <= max_attempts))), + CONSTRAINT positive_max_attempts CHECK ((max_attempts > 0)), + CONSTRAINT queue_length CHECK (((char_length(queue) > 0) AND (char_length(queue) < 128))), + CONSTRAINT worker_length CHECK (((char_length(worker) > 0) AND (char_length(worker) < 128))) +); + + +-- +-- Name: TABLE oban_jobs; Type: COMMENT; Schema: public; Owner: - +-- + +COMMENT ON TABLE public.oban_jobs IS '14'; + + +-- +-- Name: oban_jobs_id_seq; Type: SEQUENCE; Schema: public; Owner: - +-- + +CREATE SEQUENCE public.oban_jobs_id_seq + START WITH 1 + INCREMENT BY 1 + NO MINVALUE + NO MAXVALUE + CACHE 1; + + +-- +-- Name: oban_jobs_id_seq; Type: SEQUENCE OWNED BY; Schema: public; Owner: - +-- + +ALTER SEQUENCE public.oban_jobs_id_seq OWNED BY public.oban_jobs.id; + + +-- +-- Name: oban_peers; Type: TABLE; Schema: public; Owner: - +-- + +CREATE UNLOGGED TABLE public.oban_peers ( + name text NOT NULL, + node text NOT NULL, + started_at timestamp without time zone NOT NULL, + expires_at timestamp without time zone NOT NULL +); + + -- -- Name: old_source_changes; Type: TABLE; Schema: public; Owner: - -- @@ -2550,6 +2634,13 @@ ALTER TABLE ONLY public.moderation_logs ALTER COLUMN id SET DEFAULT nextval('pub ALTER TABLE ONLY public.notifications ALTER COLUMN id SET DEFAULT nextval('public.notifications_id_seq'::regclass); +-- +-- Name: oban_jobs id; Type: DEFAULT; Schema: public; Owner: - +-- + +ALTER TABLE ONLY public.oban_jobs ALTER COLUMN id SET DEFAULT nextval('public.oban_jobs_id_seq'::regclass); + + -- -- Name: old_source_changes id; Type: DEFAULT; Schema: public; Owner: - -- @@ -2938,6 +3029,14 @@ ALTER TABLE ONLY public.moderation_logs ADD CONSTRAINT moderation_logs_pkey PRIMARY KEY (id); +-- +-- Name: oban_jobs non_negative_priority; Type: CHECK CONSTRAINT; Schema: public; Owner: - +-- + +ALTER TABLE public.oban_jobs + ADD CONSTRAINT non_negative_priority CHECK ((priority >= 0)) NOT VALID; + + -- -- Name: notifications notifications_pkey; Type: CONSTRAINT; Schema: public; Owner: - -- @@ -2946,6 +3045,22 @@ ALTER TABLE ONLY public.notifications ADD CONSTRAINT notifications_pkey PRIMARY KEY (id); +-- +-- Name: oban_jobs oban_jobs_pkey; Type: CONSTRAINT; Schema: public; Owner: - +-- + +ALTER TABLE ONLY public.oban_jobs + ADD CONSTRAINT oban_jobs_pkey PRIMARY KEY (id); + + +-- +-- Name: oban_peers oban_peers_pkey; Type: CONSTRAINT; Schema: public; Owner: - +-- + +ALTER TABLE ONLY public.oban_peers + ADD CONSTRAINT oban_peers_pkey PRIMARY KEY (name); + + -- -- Name: old_source_changes old_source_changes_pkey; Type: CONSTRAINT; Schema: public; Owner: - -- @@ -4629,6 +4744,41 @@ CREATE INDEX moderation_logs_user_id_created_at_index ON public.moderation_logs CREATE INDEX moderation_logs_user_id_index ON public.moderation_logs USING btree (user_id); +-- +-- Name: oban_jobs_args_index; Type: INDEX; Schema: public; Owner: - +-- + +CREATE INDEX oban_jobs_args_index ON public.oban_jobs USING gin (args); + + +-- +-- Name: oban_jobs_meta_index; Type: INDEX; Schema: public; Owner: - +-- + +CREATE INDEX oban_jobs_meta_index ON public.oban_jobs USING gin (meta); + + +-- +-- Name: oban_jobs_state_cancelled_at_index; Type: INDEX; Schema: public; Owner: - +-- + +CREATE INDEX oban_jobs_state_cancelled_at_index ON public.oban_jobs USING btree (state, cancelled_at); + + +-- +-- Name: oban_jobs_state_discarded_at_index; Type: INDEX; Schema: public; Owner: - +-- + +CREATE INDEX oban_jobs_state_discarded_at_index ON public.oban_jobs USING btree (state, discarded_at); + + +-- +-- Name: oban_jobs_state_queue_priority_scheduled_at_id_index; Type: INDEX; Schema: public; Owner: - +-- + +CREATE INDEX oban_jobs_state_queue_priority_scheduled_at_id_index ON public.oban_jobs USING btree (state, queue, priority, scheduled_at, id); + + -- -- Name: post_versions_post_id_created_at_index; Type: INDEX; Schema: public; Owner: - -- @@ -5952,7 +6102,7 @@ ALTER TABLE ONLY public.users -- PostgreSQL database dump complete -- -\unrestrict O6OYrjIYDwvGhgBIGh2SR1j95V8aVnMGs9WcuKijctXbT5ypaF3nbIkiah3iFpk +\unrestrict DiOFiU8d5t1AmxDoKBF0A1E487owRGbnVYjRKkg1IFCSL1or1ZGLi8wPipZ954H INSERT INTO public."schema_migrations" (version) VALUES (20200503002523); INSERT INTO public."schema_migrations" (version) VALUES (20200607000511); @@ -5998,3 +6148,4 @@ INSERT INTO public."schema_migrations" (version) VALUES (20260719123611); INSERT INTO public."schema_migrations" (version) VALUES (20260806180557); INSERT INTO public."schema_migrations" (version) VALUES (20260810212302); INSERT INTO public."schema_migrations" (version) VALUES (20260831235832); +INSERT INTO public."schema_migrations" (version) VALUES (20260905193823); diff --git a/test/CONVENTIONS.md b/test/CONVENTIONS.md index 563a9bae4..7cb130d56 100644 --- a/test/CONVENTIONS.md +++ b/test/CONVENTIONS.md @@ -59,8 +59,8 @@ One test per auth level that can reach the action: - Attribution-taking contexts (topics, posts, comments, reports) use `Philomena.AttributionFixtures.attribution/1`; pass a user or `nil` for anonymous. Attrs for these fixtures are string-keyed, controller-style. -- Context functions that enqueue Exq jobs are safe: test config uses Exq's - in-memory fake queue, so jobs neither reach Valkey nor OpenSearch. +- Context functions that enqueue Oban jobs are safe: test config stores jobs in + the sandbox database without starting queue consumers. ## Singleton toggle controllers (phase 3) diff --git a/test/philomena/background_jobs_test.exs b/test/philomena/background_jobs_test.exs index 46a7b99f7..311e5f1fd 100644 --- a/test/philomena/background_jobs_test.exs +++ b/test/philomena/background_jobs_test.exs @@ -1,104 +1,59 @@ defmodule Philomena.BackgroundJobsTest do use Philomena.DataCase, async: false - use Patch + import Ecto.Query import Philomena.AttributionFixtures import Philomena.ImagesFixtures import Philomena.TagsFixtures import Philomena.UsersFixtures - alias Philomena.Comments - alias Philomena.Comments.Comment - alias Philomena.Filters - alias Philomena.Filters.Filter - alias Philomena.Galleries - alias Philomena.Galleries.Gallery alias Philomena.Images - alias Philomena.Images.Image alias Philomena.Images.Thumbnailer alias Philomena.Multi - alias Philomena.Posts - alias Philomena.Posts.Post alias Philomena.Reports alias Philomena.Reports.Report alias Philomena.TagChanges alias Philomena.TagChanges.TagChange alias Philomena.Tags alias Philomena.Tags.Tag - alias Philomena.Topics.Topic alias Philomena.Users alias Philomena.Users.User - defp assert_enqueued(queue, worker, arguments, count \\ 1) do - call = {:enqueue, [Exq, queue, worker, arguments]} - assert Enum.count(history(Exq), &(&1 == call)) == count + setup do + baseline_id = Repo.one(from job in Oban.Job, select: max(job.id)) || 0 + Process.put(:oban_test_baseline_id, baseline_id) + :ok end - describe "search indexing jobs" do - test "comments enqueue the selected id column and values" do - comment = %Comment{id: 11} - image = %Image{id: 12} - spy(Exq) - - assert Comments.reindex_comment(comment) == comment - assert Comments.reindex_comments_on_image(image) == image - assert Comments.reindex_comments_on_images([12, 13]) == [12, 13] - - assert_enqueued("indexing", Philomena.IndexWorker, ["Comments", "id", [11]]) - assert_enqueued("indexing", Philomena.IndexWorker, ["Comments", "image_id", [12]]) - assert_enqueued("indexing", Philomena.IndexWorker, ["Comments", "image_id", [12, 13]]) - end - - test "filters enqueue their id" do - filter = %Filter{id: 21} - spy(Exq) - - assert Filters.reindex_filter(filter) == filter - - assert_enqueued("indexing", Philomena.IndexWorker, ["Filters", "id", [21]]) - end - - test "galleries enqueue one or many ids and skip an empty batch" do - gallery = %Gallery{id: 31} - spy(Exq) - - assert Galleries.reindex_gallery(gallery) == gallery - assert Galleries.reindex_galleries([31, 32]) == [31, 32] - assert Galleries.reindex_galleries([]) == [] - - assert_enqueued("indexing", Philomena.IndexWorker, ["Galleries", "id", [31]]) - assert_enqueued("indexing", Philomena.IndexWorker, ["Galleries", "id", [31, 32]]) - end - - test "images enqueue one or many ids" do - image = %Image{id: 41} - spy(Exq) - - assert Images.reindex_image(image) == image - assert Images.reindex_images([41, 42]) == [41, 42] - - assert_enqueued("indexing", Philomena.IndexWorker, ["Images", "id", [41]]) - assert_enqueued("indexing", Philomena.IndexWorker, ["Images", "id", [41, 42]]) - end + defp assert_enqueued(queue, worker, arguments, count \\ 1) do + baseline_id = Process.get(:oban_test_baseline_id) + + jobs = + from(job in Oban.Job, + where: + job.queue == ^queue and job.worker == ^inspect(worker) and + job.id > ^baseline_id, + select: job.args + ) + + expected = stringify_keys(arguments) + assert Enum.count(Repo.all(jobs), &(&1 == expected)) >= count + end - test "posts enqueue a post id or a topic id" do - post = %Post{id: 51} - topic = %Topic{id: 52} - spy(Exq) + defp stringify_keys(value) when is_list(value), do: Enum.map(value, &stringify_keys/1) - assert Posts.reindex_post(post) == post - assert Posts.reindex_posts_in_topic(topic) == :ok + defp stringify_keys(value) when is_map(value) do + Map.new(value, fn {key, value} -> {to_string(key), stringify_keys(value)} end) + end - assert_enqueued("indexing", Philomena.IndexWorker, ["Posts", "id", [51]]) - assert_enqueued("indexing", Philomena.IndexWorker, ["Posts", "topic_id", [52]]) - end + defp stringify_keys(value), do: value + describe "search indexing jobs" do test "persisted tag changes enqueue their generated id after commit" do user = confirmed_user_fixture() tags = "safe, background base one, background base two" image = image_fixture(tags: tags) reset_tag_change_limits(attribution(user)) - spy(Exq) assert {:ok, _image} = Images.update_image_tags(actor(user), image.id, %{ @@ -108,30 +63,11 @@ defmodule Philomena.BackgroundJobsTest do tag_change = Repo.get_by!(TagChange, image_id: image.id) - assert_enqueued("indexing", Philomena.IndexWorker, [ - "TagChanges", - "id", - [tag_change.id] - ]) - end - - test "tags enqueue a batch of ids" do - tag = %Tag{id: 71} - other_tag = %Tag{id: 72} - spy(Exq) - - assert Tags.reindex_tags([tag, other_tag]) == [tag, other_tag] - - assert_enqueued("indexing", Philomena.IndexWorker, ["Tags", "id", [71, 72]]) - end - - test "users enqueue their id" do - user = %User{id: 81} - spy(Exq) - - assert Users.reindex_user(user) == user - - assert_enqueued("indexing", Philomena.IndexWorker, ["Users", "id", [81]]) + assert_enqueued("indexing", Philomena.Workers.IndexJob, %{ + module: "TagChanges", + column: "id", + condition: [tag_change.id] + }) end end @@ -142,15 +78,13 @@ defmodule Philomena.BackgroundJobsTest do video = image_fixture(image_mime_type: "video/webm", image_format: "webm") image_id = image.id video_id = video.id - spy(Exq) - assert {:ok, repaired_image} = Images.create_image_repair(actor(moderator), image.id) assert {:ok, repaired_video} = Images.create_image_repair(actor(moderator), video.id) assert repaired_image.id == image.id assert repaired_video.id == video.id - assert_enqueued("images", Philomena.ThumbnailWorker, [image_id]) - assert_enqueued("videos", Philomena.ThumbnailWorker, [video_id]) + assert_enqueued("images", Philomena.Workers.ThumbnailJob, %{image_id: image_id}) + assert_enqueued("videos", Philomena.Workers.ThumbnailJob, %{image_id: video_id}) end test "purges every visible and hidden thumbnail path" do @@ -161,11 +95,9 @@ defmodule Philomena.BackgroundJobsTest do expected_files = Thumbnailer.thumbnail_urls(image, hidden_key) ++ Thumbnailer.thumbnail_urls(image, nil) - spy(Exq) - assert {:ok, _image} = Images.create_image_repair(actor(moderator), image.id) - assert_enqueued("indexing", Philomena.ImagePurgeWorker, [expected_files]) + assert_enqueued("indexing", Philomena.Workers.ImagePurgeJob, %{files: expected_files}) end end @@ -178,8 +110,6 @@ defmodule Philomena.BackgroundJobsTest do target_id = target.id delete_tag_id = delete_tag.id alias_tag_id = alias_tag.id - spy(Exq) - assert {:ok, %Tag{id: ^delete_tag_id}} = Tags.delete_tag(actor(admin), delete_tag.slug) assert {:ok, aliased} = @@ -192,10 +122,20 @@ defmodule Philomena.BackgroundJobsTest do assert {:ok, %Tag{id: ^target_id}} = Tags.create_tag_reindex(actor(admin), target.slug) - assert_enqueued("indexing", Philomena.TagDeleteWorker, [delete_tag_id]) - assert_enqueued("indexing", Philomena.TagAliasWorker, [alias_tag_id, target_id]) - assert_enqueued("indexing", Philomena.TagReindexWorker, [target_id], 2) - assert_enqueued("indexing", Philomena.IndexWorker, ["Tags", "id", [target_id]]) + assert_enqueued("indexing", Philomena.Workers.TagDeleteJob, %{tag_id: delete_tag_id}) + + assert_enqueued("indexing", Philomena.Workers.TagAliasJob, %{ + tag_id: alias_tag_id, + target_tag_id: target_id + }) + + assert_enqueued("indexing", Philomena.Workers.TagReindexJob, %{tag_id: target_id}, 2) + + assert_enqueued("indexing", Philomena.Workers.IndexJob, %{ + module: "Tags", + column: "id", + condition: [target_id] + }) end end @@ -204,28 +144,45 @@ defmodule Philomena.BackgroundJobsTest do admin = admin_user_fixture() target = user_fixture(name: "background target") target_id = target.id - spy(Exq) - assert {:ok, %User{id: ^target_id}} = Users.delete_user_downvotes(actor(admin), target.slug) assert {:ok, %User{id: ^target_id}} = Users.delete_user_votes(actor(admin), target.slug) assert {:ok, %User{id: ^target_id}} = Users.create_user_wipe(actor(admin), target.slug) - assert_enqueued("indexing", Philomena.UserUnvoteWorker, [target_id, false]) - assert_enqueued("indexing", Philomena.UserUnvoteWorker, [target_id, true]) - assert_enqueued("indexing", Philomena.UserWipeWorker, [target_id]) + assert_enqueued("indexing", Philomena.Workers.UserUnvoteJob, %{ + user_id: target_id, + votes_and_faves_too?: false + }) + + assert_enqueued("indexing", Philomena.Workers.UserUnvoteJob, %{ + user_id: target_id, + votes_and_faves_too?: true + }) + + assert_enqueued("indexing", Philomena.Workers.UserWipeJob, %{user_id: target_id}) old_name = admin.name assert {:ok, renamed_admin} = Users.update_name(actor(admin), %{"name" => "renamed admin"}) new_name = renamed_admin.name admin_id = renamed_admin.id - assert_enqueued("indexing", Philomena.UserRenameWorker, [old_name, new_name]) + assert_enqueued("indexing", Philomena.Workers.UserRenameJob, %{ + old_name: old_name, + new_name: new_name + }) assert {:ok, erased} = Users.create_user_erase(actor(renamed_admin), target.slug) erased_id = erased.id - assert_enqueued("indexing", Philomena.UserEraseWorker, [target_id, admin_id]) - assert_enqueued("indexing", Philomena.IndexWorker, ["Users", "id", [erased_id]]) + assert_enqueued("indexing", Philomena.Workers.UserEraseJob, %{ + user_id: target_id, + moderator_id: admin_id + }) + + assert_enqueued("indexing", Philomena.Workers.IndexJob, %{ + module: "Users", + column: "id", + condition: [erased_id] + }) end end @@ -243,8 +200,6 @@ defmodule Philomena.BackgroundJobsTest do batch_size: 100 } - spy(Exq) - assert {:ok, %User{id: ^target_id}} = TagChanges.create_user_tag_change_revert(moderator_actor, target.slug) @@ -254,17 +209,20 @@ defmodule Philomena.BackgroundJobsTest do assert {:ok, "c1774"} = TagChanges.create_fingerprint_tag_change_revert(moderator_actor, "c1774") - assert_enqueued("indexing", Philomena.TagChangeRevertWorker, [ - %{user_id: target_id, attributes: attributes} - ]) + assert_enqueued("indexing", Philomena.Workers.TagChangeRevertJob, %{ + user_id: target_id, + attributes: attributes + }) - assert_enqueued("indexing", Philomena.TagChangeRevertWorker, [ - %{ip: "203.0.113.9", attributes: attributes} - ]) + assert_enqueued("indexing", Philomena.Workers.TagChangeRevertJob, %{ + ip: "203.0.113.9", + attributes: attributes + }) - assert_enqueued("indexing", Philomena.TagChangeRevertWorker, [ - %{fingerprint: "c1774", attributes: attributes} - ]) + assert_enqueued("indexing", Philomena.Workers.TagChangeRevertJob, %{ + fingerprint: "c1774", + attributes: attributes + }) end end @@ -275,12 +233,15 @@ defmodule Philomena.BackgroundJobsTest do moderator = moderator_user_fixture() report = Philomena.ReportsFixtures.report_fixture(reporter, %{}, image_id: image.id) report_id = report.id - spy(Exq) assert {:ok, %Report{id: ^report_id}} = Reports.create_report_close(actor(moderator), report_id) - assert_enqueued("indexing", Philomena.IndexWorker, ["Reports", "id", [report_id]]) + assert_enqueued("indexing", Philomena.Workers.IndexJob, %{ + module: "Reports", + column: "id", + condition: [report_id] + }) end test "bulk report closure enqueues the returned id list" do @@ -288,7 +249,6 @@ defmodule Philomena.BackgroundJobsTest do reporter = confirmed_user_fixture() moderator = moderator_user_fixture() report = Philomena.ReportsFixtures.report_fixture(reporter, %{}, image_id: image.id) - spy(Exq) assert {:ok, %{reports: {1, [report_id]}}} = Multi.new() @@ -296,7 +256,12 @@ defmodule Philomena.BackgroundJobsTest do |> Multi.transact() assert report_id == report.id - assert_enqueued("indexing", Philomena.IndexWorker, ["Reports", "id", [report_id]]) + + assert_enqueued("indexing", Philomena.Workers.IndexJob, %{ + module: "Reports", + column: "id", + condition: [report_id] + }) end end end diff --git a/test/philomena/user_name_changes_test.exs b/test/philomena/user_name_changes_test.exs index b985c8665..b933fd0f9 100644 --- a/test/philomena/user_name_changes_test.exs +++ b/test/philomena/user_name_changes_test.exs @@ -15,7 +15,8 @@ defmodule Philomena.UserNameChangesTest do defp record_name!(user) do {:ok, %{name_change: change}} = Multi.new() - |> UserNameChanges.record_rename(:name_change, user) + |> Multi.put(:locked_user, user) + |> UserNameChanges.put_record_rename(:name_change, :locked_user) |> Multi.transact() change @@ -36,7 +37,8 @@ defmodule Philomena.UserNameChangesTest do assert {:error, :later_step, :forced, %{}} = Multi.new() - |> UserNameChanges.record_rename(:name_change, user) + |> Multi.put(:locked_user, user) + |> UserNameChanges.put_record_rename(:name_change, :locked_user) |> Multi.run(:later_step, fn _repo, _changes -> {:error, :forced} end) |> Multi.transact() diff --git a/test/philomena/user_statistics_concurrency_test.exs b/test/philomena/user_statistics_concurrency_test.exs index d50d26d17..50b8c6242 100644 --- a/test/philomena/user_statistics_concurrency_test.exs +++ b/test/philomena/user_statistics_concurrency_test.exs @@ -3,6 +3,7 @@ defmodule Philomena.UserStatisticsConcurrencyTest do import Philomena.UsersFixtures + alias Philomena.Multi alias Philomena.Repo alias Philomena.UserStatistics alias Philomena.UserStatistics.UserStatistic @@ -14,11 +15,15 @@ defmodule Philomena.UserStatisticsConcurrencyTest do results = concurrently( for _ <- 1..8 do - fn -> UserStatistics.increment(user.id, :image_votes_count) end + fn -> + Multi.new() + |> UserStatistics.put_increment(user.id, :image_votes_count) + |> Multi.transact() + end end ) - assert results == List.duplicate({:ok, nil}, 8) + assert Enum.all?(results, &match?({:ok, _}, &1)) assert Repo.get!(User, user.id).image_votes_count == 8 assert Repo.get_by!(UserStatistic, user_id: user.id).image_votes_count == 8 end diff --git a/test/philomena/user_statistics_test.exs b/test/philomena/user_statistics_test.exs index d6a9f1b5c..229c511cc 100644 --- a/test/philomena/user_statistics_test.exs +++ b/test/philomena/user_statistics_test.exs @@ -9,10 +9,16 @@ defmodule Philomena.UserStatisticsTest do alias Philomena.UserStatistics alias Philomena.UserStatistics.UserStatistic + defp transact_increment(user_or_id, statistic, amount \\ 1) do + Multi.new() + |> UserStatistics.put_increment(user_or_id, statistic, amount) + |> Multi.transact() + end + test "increments a loaded user's lifetime and current UTC-day counters" do user = confirmed_user_fixture() - assert UserStatistics.increment(user, :images_count) == {:ok, nil} + assert {:ok, _changes} = transact_increment(user, :images_count) assert Repo.get!(User, user.id).images_count == 1 @@ -25,44 +31,47 @@ defmodule Philomena.UserStatisticsTest do test "accepts an ID and negative amounts" do user = confirmed_user_fixture() - assert UserStatistics.increment(user.id, :comments_count, 4) == {:ok, nil} - assert UserStatistics.increment(user.id, :comments_count, -2) == {:ok, nil} + assert {:ok, _changes} = transact_increment(user.id, :comments_count, 4) + assert {:ok, _changes} = transact_increment(user.id, :comments_count, -2) assert Repo.get!(User, user.id).comments_count == 2 assert Repo.get_by!(UserStatistic, user_id: user.id).comments_count == 2 end test "nil users are a no-op for anonymous activity" do - assert UserStatistics.increment(nil, :comments_count) == {:ok, nil} + assert {:ok, _changes} = transact_increment(nil, :comments_count) assert Repo.aggregate(UserStatistic, :count) == 0 end test "a missing user ID is not-found and creates no daily row" do - assert UserStatistics.increment(2_000_000_000, :posts_count) == {:error, :not_found} + assert {:error, _step, :not_found, _changes} = + transact_increment(2_000_000_000, :posts_count) + assert Repo.aggregate(UserStatistic, :count) == 0 end - test "unknown keys and non-integer amounts do not match the service API" do + test "unknown keys and non-integer amounts do not match the transactional API" do user = confirmed_user_fixture() assert_raise FunctionClauseError, fn -> # credo:disable-for-next-line Credo.Check.Refactor.Apply - apply(UserStatistics, :increment, [user, :email, 1]) + apply(UserStatistics, :put_increment, [Multi.new(), user, :email, 1]) end assert_raise FunctionClauseError, fn -> # credo:disable-for-next-line Credo.Check.Refactor.Apply - apply(UserStatistics, :increment, [user, :images_count, 1.5]) + apply(UserStatistics, :put_increment, [Multi.new(), user, :images_count, 1.5]) end end test "an owning transaction rollback restores both counters" do user = confirmed_user_fixture() - assert Repo.transact(fn -> - assert UserStatistics.increment(user, :topics_count) == {:ok, nil} - {:error, :forced_rollback} - end) == {:error, :forced_rollback} + assert {:error, :rollback, :forced_rollback, _changes} = + Multi.new() + |> UserStatistics.put_increment(user, :topics_count) + |> Multi.run(:rollback, fn _repo, _changes -> {:error, :forced_rollback} end) + |> Multi.transact() assert Repo.get!(User, user.id).topics_count == 0 refute Repo.get_by(UserStatistic, user_id: user.id) @@ -98,11 +107,12 @@ defmodule Philomena.UserStatisticsTest do test "daily rows cascade on user deletion and deleted IDs are not-found" do user = confirmed_user_fixture() - assert UserStatistics.increment(user, :posts_count) == {:ok, nil} + assert {:ok, _changes} = transact_increment(user, :posts_count) Repo.delete!(user) refute Repo.get_by(UserStatistic, user_id: user.id) - assert UserStatistics.increment(user.id, :posts_count) == {:error, :not_found} + + assert {:error, _step, :not_found, _changes} = transact_increment(user.id, :posts_count) end end diff --git a/test/philomena/workers/tag_change_revert_worker_test.exs b/test/philomena/workers/tag_change_revert_worker_test.exs index dfb2ea940..70e6c26b1 100644 --- a/test/philomena/workers/tag_change_revert_worker_test.exs +++ b/test/philomena/workers/tag_change_revert_worker_test.exs @@ -1,15 +1,15 @@ -defmodule Philomena.TagChangeRevertWorkerTest do +defmodule Philomena.Workers.TagChangeRevertJobTest do use Philomena.DataCase, async: true # The worker runs synchronously here; only its reindex side effects are - # dead Exq enqueues, so the tests stay Postgres-only. + # dead Oban enqueues, so the tests stay Postgres-only. import Philomena.AttributionFixtures import Philomena.ImagesFixtures import Philomena.UsersFixtures alias Philomena.Images - alias Philomena.TagChangeRevertWorker + alias Philomena.Workers.TagChangeRevertJob alias Philomena.TagChanges.TagChange # Images validate a 3-tag minimum, so every input keeps these on top of @@ -38,13 +38,15 @@ defmodule Philomena.TagChangeRevertWorkerTest do end defp full_revert!(user, batch_size) do - TagChangeRevertWorker.perform(%{ - "user_id" => user.id, - "attributes" => %{ - "ip" => "203.0.113.99", - "fingerprint" => "c1774e9294a", - "user_id" => moderator_user_fixture().id, - "batch_size" => batch_size + TagChangeRevertJob.perform(%Oban.Job{ + args: %{ + "user_id" => user.id, + "attributes" => %{ + "ip" => "203.0.113.99", + "fingerprint" => "c1774e9294a", + "user_id" => moderator_user_fixture().id, + "batch_size" => batch_size + } } }) end diff --git a/test/philomena/workers/tag_workers_test.exs b/test/philomena/workers/tag_workers_test.exs index a2763f687..feee38866 100644 --- a/test/philomena/workers/tag_workers_test.exs +++ b/test/philomena/workers/tag_workers_test.exs @@ -10,8 +10,8 @@ defmodule Philomena.TagWorkersTest do alias Philomena.Repo alias Philomena.ArtistLinks.ArtistLink - alias Philomena.TagAliasWorker - alias Philomena.TagDeleteWorker + alias Philomena.Workers.TagAliasJob + alias Philomena.Workers.TagDeleteJob alias Philomena.TagChanges.TagChange alias Philomena.TagChanges.TagChangeTag alias Philomena.Tags.Tag @@ -34,7 +34,10 @@ defmodule Philomena.TagWorkersTest do |> Ecto.Changeset.change(aliased_tag_id: target.id) |> Repo.update!() - assert :ok = TagAliasWorker.perform(source.id, target.id) + assert :ok = + TagAliasJob.perform(%Oban.Job{ + args: %{"tag_id" => source.id, "target_tag_id" => target.id} + }) ids = tag_ids(image) assert target.id in ids @@ -54,7 +57,10 @@ defmodule Philomena.TagWorkersTest do |> Ecto.Changeset.change(aliased_tag_id: target.id) |> Repo.update!() - assert :ok = TagAliasWorker.perform(source.id, target.id) + assert :ok = + TagAliasJob.perform(%Oban.Job{ + args: %{"tag_id" => source.id, "target_tag_id" => target.id} + }) assert Repo.reload!(source).images_count == 0 assert Repo.reload!(target).images_count == 1 @@ -72,7 +78,10 @@ defmodule Philomena.TagWorkersTest do |> Ecto.Changeset.change(aliased_tag_id: target.id) |> Repo.update!() - assert :ok = TagAliasWorker.perform(source.id, target.id) + assert :ok = + TagAliasJob.perform(%Oban.Job{ + args: %{"tag_id" => source.id, "target_tag_id" => target.id} + }) assert Repo.reload!(source).images_count == 0 assert Repo.reload!(target).images_count == 0 @@ -107,7 +116,7 @@ defmodule Philomena.TagWorkersTest do tag = tag_fixture(name: "worker delete tag") image = image_fixture(tags: "safe, #{tag.name}") - assert :ok = TagDeleteWorker.perform(tag.id) + assert :ok = TagDeleteJob.perform(%Oban.Job{args: %{"tag_id" => tag.id}}) assert Repo.get(Tag, tag.id) == nil refute tag.id in tag_ids(image) @@ -135,7 +144,7 @@ defmodule Philomena.TagWorkersTest do assert Repo.get_by(TagChangeTag, tag_change_id: tag_change.id, tag_id: tag.id) - assert :ok = TagDeleteWorker.perform(tag.id) + assert :ok = TagDeleteJob.perform(%Oban.Job{args: %{"tag_id" => tag.id}}) refute Repo.get_by(TagChangeTag, tag_change_id: tag_change.id, tag_id: tag.id) refute Repo.get(TagChange, tag_change.id) diff --git a/test/philomena/workers/user_workers_test.exs b/test/philomena/workers/user_workers_test.exs index 59ecfb5f1..b3709f642 100644 --- a/test/philomena/workers/user_workers_test.exs +++ b/test/philomena/workers/user_workers_test.exs @@ -15,11 +15,11 @@ defmodule Philomena.UserWorkersTest do alias Philomena.ImageVotes.ImageVote alias Philomena.Repo alias Philomena.SourceChanges.SourceChange - alias Philomena.UserEraseWorker + alias Philomena.Workers.UserEraseJob alias Philomena.UserFingerprints.UserFingerprint alias Philomena.UserIps.UserIp - alias Philomena.UserUnvoteWorker - alias Philomena.UserWipeWorker + alias Philomena.Workers.UserUnvoteJob + alias Philomena.Workers.UserWipeJob alias Philomena.Users.User alias Philomena.Users.UserDownvoteWipe @@ -33,7 +33,10 @@ defmodule Philomena.UserWorkersTest do assert {:ok, _image} = Images.create_image_vote(actor(user), downvoted.id, %{up: false}) assert {:ok, _image} = Images.create_image_fave(actor(user), faved.id) - assert :ok = UserUnvoteWorker.perform(user.id, true) + assert :ok = + UserUnvoteJob.perform(%Oban.Job{ + args: %{"user_id" => user.id, "votes_and_faves_too?" => true} + }) refute Repo.get_by(ImageVote, user_id: user.id, image_id: downvoted.id) refute Repo.get_by(ImageVote, user_id: user.id, image_id: faved.id) @@ -56,8 +59,7 @@ defmodule Philomena.UserWorkersTest do fingerprint: "worker-fingerprint" ) - assert %User{id: user_id} = UserWipeWorker.perform(user.id) - assert user_id == user.id + assert :ok = UserWipeJob.perform(%Oban.Job{args: %{"user_id" => user.id}}) wiped_user = Repo.reload!(user) assert wiped_user.email =~ ~r/^deactivated[0-9a-f]{32}@example\.com$/ @@ -82,7 +84,10 @@ defmodule Philomena.UserWorkersTest do }) |> Repo.update!() - assert :ok = UserEraseWorker.perform(user.id, moderator.id) + assert :ok = + UserEraseJob.perform(%Oban.Job{ + args: %{"user_id" => user.id, "moderator_id" => moderator.id} + }) erased = Repo.reload!(user) assert erased.description in [nil, ""] @@ -118,7 +123,10 @@ defmodule Philomena.UserWorkersTest do added: false ) - assert :ok = UserEraseWorker.perform(user.id, moderator.id) + assert :ok = + UserEraseJob.perform(%Oban.Job{ + args: %{"user_id" => user.id, "moderator_id" => moderator.id} + }) assert Repo.reload!(added_image) |> Repo.preload(:sources) |> Map.fetch!(:sources) == [] diff --git a/test/philomena/workers/worker_test.exs b/test/philomena/workers/worker_test.exs index 8816c5a65..a0941e1ab 100644 --- a/test/philomena/workers/worker_test.exs +++ b/test/philomena/workers/worker_test.exs @@ -6,13 +6,12 @@ defmodule Philomena.WorkerTest do import Philomena.ImagesFixtures - alias Philomena.ImagePurgeWorker - alias Philomena.Images + alias Philomena.Workers.ImagePurgeJob alias Philomena.Images.Image alias Philomena.Images.Thumbnailer - alias Philomena.IndexWorker - alias Philomena.ThumbnailWorker - alias Philomena.UserRenameWorker + alias Philomena.Workers.IndexJob + alias Philomena.Workers.ThumbnailJob + alias Philomena.Workers.UserRenameJob alias PhilomenaQuery.Search @index_contexts [ @@ -51,7 +50,11 @@ defmodule Philomena.WorkerTest do test "the index worker indexes matching records" do image = image_fixture() - assert :ok = IndexWorker.perform("Images", "id", [image.id]) + assert :ok = + IndexJob.perform(%Oban.Job{ + args: %{"module" => "Images", "column" => "id", "condition" => [image.id]} + }) + :ok = Search.refresh_index!(Image) hits = Search.search(Image, %{query: %{match_all: %{}}})["hits"]["hits"] @@ -64,7 +67,10 @@ defmodule Philomena.WorkerTest do end) Enum.each(@index_contexts, fn {name, _context} -> - assert :ok = IndexWorker.perform(name, "id", [123]) + assert :ok = + IndexJob.perform(%Oban.Job{ + args: %{"module" => name, "column" => "id", "condition" => [123]} + }) end) Enum.each(@index_contexts, fn {_name, context} -> @@ -72,24 +78,19 @@ defmodule Philomena.WorkerTest do end) end - test "the thumbnail worker generates media, broadcasts completion, and reindexes" do - image = %Image{id: 321} + test "the thumbnail worker generates media and broadcasts completion" do patch(Thumbnailer, :generate_thumbnails, :ok) - patch(Images, :load_image_for_reindex!, image) - patch(Images, :reindex_image, image) - assert image == ThumbnailWorker.perform(image.id) + assert :ok == ThumbnailJob.perform(%Oban.Job{args: %{"image_id" => 321}}) - assert_exact_call(Thumbnailer, :generate_thumbnails, [image.id]) - assert_exact_call(Images, :load_image_for_reindex!, [image.id]) - assert_exact_call(Images, :reindex_image, [image]) + assert_exact_call(Thumbnailer, :generate_thumbnails, [321]) end test "the purge worker passes the complete file list to the purge operation" do files = ["/img/1/full.png", "/img/1/thumb.png"] patch(System, :cmd, {"", 0}) - assert :ok = ImagePurgeWorker.perform(files) + assert :ok = ImagePurgeJob.perform(%Oban.Job{args: %{"files" => files}}) assert_exact_call(System, :cmd, [ "purge-cache", @@ -100,7 +101,10 @@ defmodule Philomena.WorkerTest do test "the rename worker updates every index containing a user name" do Enum.each(@rename_contexts, &patch(&1, :user_name_reindex, :ok)) - assert :ok = UserRenameWorker.perform("Old Name", "New Name") + assert :ok = + UserRenameJob.perform(%Oban.Job{ + args: %{"old_name" => "Old Name", "new_name" => "New Name"} + }) Enum.each(@rename_contexts, fn context -> assert_exact_call(context, :user_name_reindex, ["Old Name", "New Name"]) diff --git a/test/philomena_web/controllers/admin/user/avatar_controller_test.exs b/test/philomena_web/controllers/admin/user/avatar_controller_test.exs index e5c6587c9..a8bc404f5 100644 --- a/test/philomena_web/controllers/admin/user/avatar_controller_test.exs +++ b/test/philomena_web/controllers/admin/user/avatar_controller_test.exs @@ -2,7 +2,7 @@ defmodule PhilomenaWeb.Admin.User.AvatarControllerTest do use PhilomenaWeb.ConnCase, async: true # Postgres-only: the S3 delete in Users.delete_avatar/1 goes through the - # stubbed ex_aws client and the reindex is a dead Exq enqueue; + # stubbed ex_aws client and the reindex is a dead Oban enqueue; # moderation_log/2 is a synchronous insert. import Philomena.UsersFixtures diff --git a/test/philomena_web/controllers/admin/user/downvote_controller_test.exs b/test/philomena_web/controllers/admin/user/downvote_controller_test.exs index cdbeaf140..2045ee35c 100644 --- a/test/philomena_web/controllers/admin/user/downvote_controller_test.exs +++ b/test/philomena_web/controllers/admin/user/downvote_controller_test.exs @@ -1,8 +1,8 @@ defmodule PhilomenaWeb.Admin.User.DownvoteControllerTest do use PhilomenaWeb.ConnCase, async: true - # Postgres-only. The actual downvote wipe is performed by UserUnvoteWorker, - # which is only enqueued (a dead Exq enqueue in test), so only the + # Postgres-only. The actual downvote wipe is performed by UserUnvoteJob, + # which is only enqueued (a dead Oban enqueue in test), so only the # flash/redirect and the synchronous moderation_log insert are observable # here. @@ -50,7 +50,7 @@ defmodule PhilomenaWeb.Admin.User.DownvoteControllerTest do describe "DELETE /admin/users/:user_id/downvotes as a plain moderator" do setup [:register_and_log_in_moderator] - # NOTE: the wipe is performed by an (unconsumed) UserUnvoteWorker enqueue, so + # NOTE: the wipe is performed by an (unconsumed) UserUnvoteJob enqueue, so # there is no observable side effect to assert absent - the denial redirect + # flash is the pin. test "is denied to a plain moderator", %{conn: conn} do diff --git a/test/philomena_web/controllers/admin/user/erase_controller_test.exs b/test/philomena_web/controllers/admin/user/erase_controller_test.exs index 7855068d7..a813c4882 100644 --- a/test/philomena_web/controllers/admin/user/erase_controller_test.exs +++ b/test/philomena_web/controllers/admin/user/erase_controller_test.exs @@ -2,8 +2,8 @@ defmodule PhilomenaWeb.Admin.User.EraseControllerTest do use PhilomenaWeb.ConnCase, async: true # Postgres-only. Users.erase_user/2 synchronously deactivates and renames - # the account, then enqueues UserEraseWorker for the rest (a dead Exq - # enqueue in test); the rename also enqueues (dead) UserRenameWorker. So + # the account, then enqueues UserEraseJob for the rest (a dead Oban + # enqueue in test); the rename also enqueues (dead) UserRenameJob. So # the deactivation and rename are observable, but the deeper deletion is # not. diff --git a/test/philomena_web/controllers/admin/user/verification_controller_test.exs b/test/philomena_web/controllers/admin/user/verification_controller_test.exs index cba5293d9..35ff4f19b 100644 --- a/test/philomena_web/controllers/admin/user/verification_controller_test.exs +++ b/test/philomena_web/controllers/admin/user/verification_controller_test.exs @@ -1,7 +1,7 @@ defmodule PhilomenaWeb.Admin.User.VerificationControllerTest do use PhilomenaWeb.ConnCase, async: true - # Postgres-only: reindex is a dead Exq enqueue, moderation_log/2 is a + # Postgres-only: reindex is a dead Oban enqueue, moderation_log/2 is a # synchronous insert. import Philomena.UsersFixtures diff --git a/test/philomena_web/controllers/admin/user/vote_controller_test.exs b/test/philomena_web/controllers/admin/user/vote_controller_test.exs index 06e9a3a32..24e107177 100644 --- a/test/philomena_web/controllers/admin/user/vote_controller_test.exs +++ b/test/philomena_web/controllers/admin/user/vote_controller_test.exs @@ -1,8 +1,8 @@ defmodule PhilomenaWeb.Admin.User.VoteControllerTest do use PhilomenaWeb.ConnCase, async: true - # Postgres-only. The actual vote/fave wipe is performed by UserUnvoteWorker, - # which is only enqueued (a dead Exq enqueue in test), so only the + # Postgres-only. The actual vote/fave wipe is performed by UserUnvoteJob, + # which is only enqueued (a dead Oban enqueue in test), so only the # flash/redirect and the synchronous moderation_log insert are observable # here. @@ -50,7 +50,7 @@ defmodule PhilomenaWeb.Admin.User.VoteControllerTest do describe "DELETE /admin/users/:user_id/votes as a plain moderator" do setup [:register_and_log_in_moderator] - # NOTE: the wipe is performed by an (unconsumed) UserUnvoteWorker enqueue, so + # NOTE: the wipe is performed by an (unconsumed) UserUnvoteJob enqueue, so # there is no observable side effect to assert absent - the denial redirect + # flash is the pin. test "is denied to a plain moderator", %{conn: conn} do diff --git a/test/philomena_web/controllers/admin/user/wipe_controller_test.exs b/test/philomena_web/controllers/admin/user/wipe_controller_test.exs index 0204d099e..409a6c323 100644 --- a/test/philomena_web/controllers/admin/user/wipe_controller_test.exs +++ b/test/philomena_web/controllers/admin/user/wipe_controller_test.exs @@ -1,8 +1,8 @@ defmodule PhilomenaWeb.Admin.User.WipeControllerTest do use PhilomenaWeb.ConnCase, async: true - # Postgres-only. The actual PII wipe is performed by UserWipeWorker, which - # is only enqueued (a dead Exq enqueue in test), so only the flash/redirect + # Postgres-only. The actual PII wipe is performed by UserWipeJob, which + # is only enqueued (a dead Oban enqueue in test), so only the flash/redirect # and the synchronous moderation_log insert are observable here. import Philomena.UsersFixtures @@ -52,7 +52,7 @@ defmodule PhilomenaWeb.Admin.User.WipeControllerTest do describe "POST /admin/users/:user_id/wipe as a plain moderator" do setup [:register_and_log_in_moderator] - # NOTE: the wipe is performed by an (unconsumed) UserWipeWorker enqueue, so + # NOTE: the wipe is performed by an (unconsumed) UserWipeJob enqueue, so # there is no observable side effect to assert absent - the denial redirect + # flash is the pin. test "is denied to a plain moderator", %{conn: conn} do diff --git a/test/philomena_web/controllers/duplicate_report/accept_controller_test.exs b/test/philomena_web/controllers/duplicate_report/accept_controller_test.exs index df43440bb..7c8d4d592 100644 --- a/test/philomena_web/controllers/duplicate_report/accept_controller_test.exs +++ b/test/philomena_web/controllers/duplicate_report/accept_controller_test.exs @@ -3,7 +3,7 @@ defmodule PhilomenaWeb.DuplicateReport.AcceptControllerTest do # Accepting a report merges the source image into the target: the merge # runs synchronously (S3 ops go through the ex_aws stub, reindexing is a - # dead Exq enqueue). + # dead Oban enqueue). import Philomena.ImagesFixtures import Philomena.DuplicateReportsFixtures diff --git a/test/philomena_web/controllers/image/file_controller_test.exs b/test/philomena_web/controllers/image/file_controller_test.exs index 8fee6fbbb..95b13f5e2 100644 --- a/test/philomena_web/controllers/image/file_controller_test.exs +++ b/test/philomena_web/controllers/image/file_controller_test.exs @@ -3,7 +3,7 @@ defmodule PhilomenaWeb.Image.FileControllerTest do # `Images.update_image_file/3` drives the media pipeline synchronously # (analyze, persist to the stubbed S3, enqueue the dead - # ThumbnailWorker/reindex jobs), with no spawned upload process, so this + # ThumbnailJob/reindex jobs), with no spawned upload process, so this # file stays `async: true`. import Philomena.ImagesFixtures diff --git a/test/philomena_web/controllers/tag/alias_controller_test.exs b/test/philomena_web/controllers/tag/alias_controller_test.exs index eb87402ef..d9e39329e 100644 --- a/test/philomena_web/controllers/tag/alias_controller_test.exs +++ b/test/philomena_web/controllers/tag/alias_controller_test.exs @@ -4,7 +4,7 @@ defmodule PhilomenaWeb.Tag.AliasControllerTest do # All three actions authorize :alias on the tag, which a *plain* # moderator lacks (they only have :edit), so only an admin or a Tag-admin # role_map moderator can reach them. The actual alias/unalias work is a - # dead Exq enqueue; only the synchronous aliased_tag_id write on :update + # dead Oban enqueue; only the synchronous aliased_tag_id write on :update # is observable. Tags are slug-keyed. import Philomena.TagsFixtures diff --git a/test/philomena_web/controllers/tag_change/full_revert_controller_test.exs b/test/philomena_web/controllers/tag_change/full_revert_controller_test.exs index 51cc7bce6..ac6431cef 100644 --- a/test/philomena_web/controllers/tag_change/full_revert_controller_test.exs +++ b/test/philomena_web/controllers/tag_change/full_revert_controller_test.exs @@ -1,7 +1,7 @@ defmodule PhilomenaWeb.Profile.TagChange.RevertControllerTest do use PhilomenaWeb.ConnCase, async: true - # full_revert only enqueues a (dead) TagChangeRevertWorker, so there is + # full_revert only enqueues a (dead) TagChangeRevertJob, so there is # nothing to observe beyond the flash and redirect. import Philomena.UsersFixtures diff --git a/test/philomena_web/controllers/tag_controller_test.exs b/test/philomena_web/controllers/tag_controller_test.exs index fa1e9b23f..20fdd6639 100644 --- a/test/philomena_web/controllers/tag_controller_test.exs +++ b/test/philomena_web/controllers/tag_controller_test.exs @@ -224,7 +224,7 @@ defmodule PhilomenaWeb.TagControllerTest do assert redirected_to(conn) == "/" assert Phoenix.Flash.get(conn.assigns.flash, :info) =~ "Tag queued for deletion" - # NOTE: delete_tag only enqueues a (dead) TagDeleteWorker, so the row is + # NOTE: delete_tag only enqueues a (dead) TagDeleteJob, so the row is # still present synchronously. assert Repo.get(Tag, tag.id) end diff --git a/test/support/context_boundary_check.ex b/test/support/context_boundary_check.ex index a4d88cfe0..59a723ecf 100644 --- a/test/support/context_boundary_check.ex +++ b/test/support/context_boundary_check.ex @@ -12,7 +12,6 @@ defmodule Philomena.ContextBoundaryCheck do Philomena.Application, Philomena.Attribution, Philomena.Config, - Philomena.ExqSupervisor, Philomena.IntegerId, Philomena.Mailer, Philomena.Maintenance, diff --git a/test/support/data_case.ex b/test/support/data_case.ex index 85a61d494..6ac30b084 100644 --- a/test/support/data_case.ex +++ b/test/support/data_case.ex @@ -16,6 +16,8 @@ defmodule Philomena.DataCase do using do quote do + use Oban.Testing, repo: Philomena.Repo + alias Philomena.Repo import Ecto diff --git a/test/support/search_helpers.ex b/test/support/search_helpers.ex index ea27318a5..8f860bd4e 100644 --- a/test/support/search_helpers.ex +++ b/test/support/search_helpers.ex @@ -16,7 +16,7 @@ defmodule PhilomenaQuery.SearchHelpers do * clear the indexes they read in `setup` with `PhilomenaQuery.Search.clear_index!/1`, * index their fixtures explicitly (fixture inserts only enqueue dead - Exq jobs) with `reindex_all!/1` before querying. + Oban jobs) with `reindex_all!/1` before querying. ## Example diff --git a/test/test_helper.exs b/test/test_helper.exs index 4daf993dd..e4dd1cc78 100644 --- a/test/test_helper.exs +++ b/test/test_helper.exs @@ -11,10 +11,6 @@ Supervisor.terminate_child(Philomena.Supervisor, Philomena.Adverts.Server) Supervisor.terminate_child(Philomena.Supervisor, Philomena.UserIps.Server) Supervisor.terminate_child(Philomena.Supervisor, Philomena.UserFingerprints.Server) -# Use Exq's in-memory fake queue in tests. This keeps enqueue side effects -# observable to tests without writing jobs to the shared Valkey instance. -{:ok, _exq_mock} = Exq.Mock.start_link(mode: :fake) - # Create every searchable index once, with the current mappings. Tests get # per-test isolation from PhilomenaQuery.Search.clear_index!/1, which only # deletes documents - dropping and recreating an index costs ~95 ms and used to