Adds backfill job to startup tasks
This commit is contained in:
parent
bc24ac3808
commit
bc25bf67e4
4 changed files with 66 additions and 4 deletions
|
|
@ -10,9 +10,9 @@ defmodule Pinchflat.Application do
|
||||||
children = [
|
children = [
|
||||||
PinchflatWeb.Telemetry,
|
PinchflatWeb.Telemetry,
|
||||||
Pinchflat.Repo,
|
Pinchflat.Repo,
|
||||||
# {Task, &run_startup_tasks/0},
|
# Must be before startup tasks
|
||||||
Pinchflat.StartupTasks,
|
|
||||||
{Oban, Application.fetch_env!(:pinchflat, Oban)},
|
{Oban, Application.fetch_env!(:pinchflat, Oban)},
|
||||||
|
Pinchflat.StartupTasks,
|
||||||
{DNSCluster, query: Application.get_env(:pinchflat, :dns_cluster_query) || :ignore},
|
{DNSCluster, query: Application.get_env(:pinchflat, :dns_cluster_query) || :ignore},
|
||||||
{Phoenix.PubSub, name: Pinchflat.PubSub},
|
{Phoenix.PubSub, name: Pinchflat.PubSub},
|
||||||
# Start the Finch HTTP client for sending emails
|
# Start the Finch HTTP client for sending emails
|
||||||
|
|
|
||||||
|
|
@ -8,8 +8,11 @@ defmodule Pinchflat.StartupTasks do
|
||||||
|
|
||||||
# restart: :temporary means that this process will never be restarted (ie: will run once and then die)
|
# restart: :temporary means that this process will never be restarted (ie: will run once and then die)
|
||||||
use GenServer, restart: :temporary
|
use GenServer, restart: :temporary
|
||||||
|
import Ecto.Query, warn: false
|
||||||
|
|
||||||
|
alias Pinchflat.Repo
|
||||||
alias Pinchflat.Settings
|
alias Pinchflat.Settings
|
||||||
|
alias Pinchflat.Workers.DataBackfillWorker
|
||||||
|
|
||||||
def start_link(opts \\ []) do
|
def start_link(opts \\ []) do
|
||||||
GenServer.start_link(__MODULE__, %{}, opts)
|
GenServer.start_link(__MODULE__, %{}, opts)
|
||||||
|
|
@ -21,11 +24,13 @@ defmodule Pinchflat.StartupTasks do
|
||||||
Any code defined here will run every time the application starts. You must
|
Any code defined here will run every time the application starts. You must
|
||||||
make sure that the code is idempotent and safe to run multiple times.
|
make sure that the code is idempotent and safe to run multiple times.
|
||||||
|
|
||||||
This is a good place to set up default settings, create initial records, stuff like that
|
This is a good place to set up default settings, create initial records, stuff like that.
|
||||||
|
Should be fast - anything with the potential to be slow should be kicked off as a job instead.
|
||||||
"""
|
"""
|
||||||
@impl true
|
@impl true
|
||||||
def init(state) do
|
def init(state) do
|
||||||
apply_default_settings()
|
apply_default_settings()
|
||||||
|
enqueue_backfill_worker()
|
||||||
|
|
||||||
{:ok, state}
|
{:ok, state}
|
||||||
end
|
end
|
||||||
|
|
@ -34,4 +39,12 @@ defmodule Pinchflat.StartupTasks do
|
||||||
Settings.fetch!(:onboarding, true)
|
Settings.fetch!(:onboarding, true)
|
||||||
Settings.fetch!(:pro_enabled, false)
|
Settings.fetch!(:pro_enabled, false)
|
||||||
end
|
end
|
||||||
|
|
||||||
|
defp enqueue_backfill_worker do
|
||||||
|
DataBackfillWorker.cancel_pending_backfill_jobs()
|
||||||
|
|
||||||
|
%{}
|
||||||
|
|> DataBackfillWorker.new()
|
||||||
|
|> Repo.insert_unique_job()
|
||||||
|
end
|
||||||
end
|
end
|
||||||
|
|
|
||||||
|
|
@ -14,11 +14,24 @@ defmodule Pinchflat.Workers.DataBackfillWorker do
|
||||||
# I'm just trying out that pattern and seeing if I like it better
|
# I'm just trying out that pattern and seeing if I like it better
|
||||||
# so this may change.
|
# so this may change.
|
||||||
import Ecto.Query, warn: false
|
import Ecto.Query, warn: false
|
||||||
|
require Logger
|
||||||
|
|
||||||
alias __MODULE__
|
alias __MODULE__
|
||||||
alias Pinchflat.Repo
|
alias Pinchflat.Repo
|
||||||
alias Pinchflat.Media.MediaItem
|
alias Pinchflat.Media.MediaItem
|
||||||
|
|
||||||
|
@doc """
|
||||||
|
Cancels all pending backfill jobs. Useful for ensuring worker runs immediately
|
||||||
|
on app boot.
|
||||||
|
|
||||||
|
Returns {:ok, integer()}
|
||||||
|
"""
|
||||||
|
def cancel_pending_backfill_jobs do
|
||||||
|
Oban.Job
|
||||||
|
|> where(worker: "Pinchflat.Workers.DataBackfillWorker")
|
||||||
|
|> Oban.cancel_all_jobs()
|
||||||
|
end
|
||||||
|
|
||||||
@impl Oban.Worker
|
@impl Oban.Worker
|
||||||
@doc """
|
@doc """
|
||||||
Performs one-off tasks to get data in the right shape.
|
Performs one-off tasks to get data in the right shape.
|
||||||
|
|
@ -30,6 +43,7 @@ defmodule Pinchflat.Workers.DataBackfillWorker do
|
||||||
Returns :ok
|
Returns :ok
|
||||||
"""
|
"""
|
||||||
def perform(%Oban.Job{}) do
|
def perform(%Oban.Job{}) do
|
||||||
|
Logger.info("Running data backfill worker")
|
||||||
backfill_shorts_data()
|
backfill_shorts_data()
|
||||||
|
|
||||||
reschedule_backfill()
|
reschedule_backfill()
|
||||||
|
|
@ -45,7 +59,9 @@ defmodule Pinchflat.Workers.DataBackfillWorker do
|
||||||
where: m.short_form_content == false
|
where: m.short_form_content == false
|
||||||
)
|
)
|
||||||
|
|
||||||
Repo.update_all(query, set: [short_form_content: true])
|
{count, _} = Repo.update_all(query, set: [short_form_content: true])
|
||||||
|
|
||||||
|
Logger.info("Backfill worker set short_form_content to true for #{count} media items.")
|
||||||
end
|
end
|
||||||
|
|
||||||
defp reschedule_backfill do
|
defp reschedule_backfill do
|
||||||
|
|
|
||||||
|
|
@ -4,8 +4,41 @@ defmodule Pinchflat.Workers.DataBackfillWorkerTest do
|
||||||
import Pinchflat.MediaFixtures
|
import Pinchflat.MediaFixtures
|
||||||
|
|
||||||
alias Pinchflat.Workers.DataBackfillWorker
|
alias Pinchflat.Workers.DataBackfillWorker
|
||||||
|
alias Pinchflat.Workers.FilesystemDataWorker
|
||||||
|
|
||||||
|
describe "cancel_pending_backfill_jobs/0" do
|
||||||
|
test "cancels all pending backfill jobs" do
|
||||||
|
%{}
|
||||||
|
|> DataBackfillWorker.new()
|
||||||
|
|> Repo.insert_unique_job()
|
||||||
|
|
||||||
|
assert_enqueued(worker: DataBackfillWorker)
|
||||||
|
|
||||||
|
DataBackfillWorker.cancel_pending_backfill_jobs()
|
||||||
|
|
||||||
|
refute_enqueued(worker: DataBackfillWorker)
|
||||||
|
end
|
||||||
|
|
||||||
|
test "does not cancel jobs for other workers" do
|
||||||
|
%{id: 0}
|
||||||
|
|> FilesystemDataWorker.new()
|
||||||
|
|> Repo.insert_unique_job()
|
||||||
|
|
||||||
|
assert_enqueued(worker: FilesystemDataWorker)
|
||||||
|
|
||||||
|
DataBackfillWorker.cancel_pending_backfill_jobs()
|
||||||
|
|
||||||
|
assert_enqueued(worker: FilesystemDataWorker)
|
||||||
|
end
|
||||||
|
end
|
||||||
|
|
||||||
describe "perform/1" do
|
describe "perform/1" do
|
||||||
|
setup do
|
||||||
|
DataBackfillWorker.cancel_pending_backfill_jobs()
|
||||||
|
|
||||||
|
:ok
|
||||||
|
end
|
||||||
|
|
||||||
test "reschedules itself once complete" do
|
test "reschedules itself once complete" do
|
||||||
perform_job(DataBackfillWorker, %{})
|
perform_job(DataBackfillWorker, %{})
|
||||||
|
|
||||||
|
|
|
||||||
Loading…
Reference in a new issue