From cfdba60b67748d605b4a28b2c95bdf83cbddaed0 Mon Sep 17 00:00:00 2001 From: Kieran Eglin Date: Tue, 30 Jan 2024 10:55:54 -0800 Subject: [PATCH] Adds tasks for media items --- lib/pinchflat/media.ex | 21 +++++++++- lib/pinchflat/media/media_item.ex | 3 ++ lib/pinchflat/repo.ex | 14 +++++++ lib/pinchflat/tasks.ex | 41 +++++++++++++----- lib/pinchflat/tasks/channel_tasks.ex | 9 +++- lib/pinchflat/tasks/task.ex | 4 +- .../workers/media_indexing_worker.ex | 39 ++++++++++------- .../workers/video_download_worker.ex | 11 ++--- ...20240130161316_add_media_item_to_tasks.exs | 12 ++++++ test/pinchflat/media_test.exs | 30 +++++++++++++ test/pinchflat/repo_test.exs | 26 ++++++++++++ test/pinchflat/tasks_test.exs | 42 ++++++++++++++++++- .../workers/media_indexing_worker_test.exs | 13 +++++- .../workers/video_download_worker_test.exs | 17 ++++++++ 14 files changed, 243 insertions(+), 39 deletions(-) create mode 100644 priv/repo/migrations/20240130161316_add_media_item_to_tasks.exs create mode 100644 test/pinchflat/repo_test.exs diff --git a/lib/pinchflat/media.ex b/lib/pinchflat/media.ex index 9178e72..0ec21d1 100644 --- a/lib/pinchflat/media.ex +++ b/lib/pinchflat/media.ex @@ -6,7 +6,9 @@ defmodule Pinchflat.Media do import Ecto.Query, warn: false alias Pinchflat.Repo + alias Pinchflat.Tasks alias Pinchflat.Media.MediaItem + alias Pinchflat.MediaSource.Channel @doc """ Returns the list of media_items. Returns [%MediaItem{}, ...]. @@ -15,6 +17,20 @@ defmodule Pinchflat.Media do Repo.all(MediaItem) end + @doc """ + Returns a list of pending media_items for a given channel, where + pending means the `video_filepath` is `nil`. + + Returns [%MediaItem{}, ...]. + """ + def list_pending_media_items_for(%Channel{} = channel) do + from( + m in MediaItem, + where: m.channel_id == ^channel.id and is_nil(m.video_filepath) + ) + |> Repo.all() + end + @doc """ Gets a single media_item. @@ -41,9 +57,12 @@ defmodule Pinchflat.Media do end @doc """ - Deletes a media_item. Returns {:ok, %MediaItem{}} | {:error, %Ecto.Changeset{}}. + Deletes a media_item and its associated tasks. + + Returns {:ok, %MediaItem{}} | {:error, %Ecto.Changeset{}}. """ def delete_media_item(%MediaItem{} = media_item) do + Tasks.delete_tasks_for(media_item) Repo.delete(media_item) end diff --git a/lib/pinchflat/media/media_item.ex b/lib/pinchflat/media/media_item.ex index 48dd779..a15427a 100644 --- a/lib/pinchflat/media/media_item.ex +++ b/lib/pinchflat/media/media_item.ex @@ -6,6 +6,7 @@ defmodule Pinchflat.Media.MediaItem do use Ecto.Schema import Ecto.Changeset + alias Pinchflat.Tasks.Task alias Pinchflat.MediaSource.Channel alias Pinchflat.Media.MediaMetadata @@ -21,6 +22,8 @@ defmodule Pinchflat.Media.MediaItem do has_one :metadata, MediaMetadata, on_replace: :update + has_many :tasks, Task + timestamps(type: :utc_datetime) end diff --git a/lib/pinchflat/repo.ex b/lib/pinchflat/repo.ex index 05ad6a6..9751f75 100644 --- a/lib/pinchflat/repo.ex +++ b/lib/pinchflat/repo.ex @@ -2,4 +2,18 @@ defmodule Pinchflat.Repo do use Ecto.Repo, otp_app: :pinchflat, adapter: Ecto.Adapters.Postgres + + @doc """ + It's not immediately obvious if an Oban job qualifies as unique, so this method + attempts creating a job and checks for the `conflict?` field in the returned job. + + Returns {:ok, %Oban.Job{}} | {:duplicate, %Oban.Job{}} | {:error, any()}. + """ + def insert_unique_job(job_struct) do + case Oban.insert(job_struct) do + {:ok, %Oban.Job{conflict?: false} = job} -> {:ok, job} + {:ok, %Oban.Job{conflict?: true} = job} -> {:duplicate, job} + err -> err + end + end end diff --git a/lib/pinchflat/tasks.ex b/lib/pinchflat/tasks.ex index d6c0b72..052b137 100644 --- a/lib/pinchflat/tasks.ex +++ b/lib/pinchflat/tasks.ex @@ -7,6 +7,7 @@ defmodule Pinchflat.Tasks do alias Pinchflat.Repo alias Pinchflat.Tasks.Task + alias Pinchflat.Media.MediaItem alias Pinchflat.MediaSource.Channel @doc """ @@ -54,7 +55,10 @@ defmodule Pinchflat.Tasks do def get_task!(id), do: Repo.get!(Task, id) @doc """ - Creates a task. Returns {:ok, %Task{}} | {:error, %Ecto.Changeset{}}. + Creates a task. + + Accepts map() | %Oban.Job{}, %Channel{} | %Oban.Job{}, %MediaItem{}. + Returns {:ok, %Task{}} | {:error, %Ecto.Changeset{}}. """ def create_task(attrs) do %Task{} @@ -62,23 +66,30 @@ defmodule Pinchflat.Tasks do |> Repo.insert() end - # This one's function signature is designed to help simplify + # This function's signature is designed to help simplify # usage of `create_job_with_task/2` - def create_task(%Oban.Job{} = job, %Channel{} = channel) do + def create_task(%Oban.Job{} = job, attached_record) do + attached_record_attr = + case attached_record do + %Channel{} = channel -> %{channel_id: channel.id} + %MediaItem{} = media_item -> %{media_item_id: media_item.id} + end + %Task{} - |> Task.changeset(%{job_id: job.id, channel_id: channel.id}) + |> Task.changeset(Map.merge(%{job_id: job.id}, attached_record_attr)) |> Repo.insert() end @doc """ Creates a job from given attrs, creating a task with an attached record - if successful. + if successful. Returns an error if the job already exists. - Returns {:ok, %Task{}} | {:error, %Ecto.Changeset{}}. + Returns {:ok, %Task{}} | {:error, :duplicate_job} | {:error, %Ecto.Changeset{}}. """ def create_job_with_task(job_attrs, task_attached_record) do - case Oban.insert(job_attrs) do + case Repo.insert_unique_job(job_attrs) do {:ok, job} -> create_task(job, task_attached_record) + {:duplicate, _} -> {:error, :duplicate_job} err -> err end end @@ -99,8 +110,12 @@ defmodule Pinchflat.Tasks do Returns :ok """ - def delete_tasks_for(%Channel{} = channel) do - tasks = list_tasks_for(:channel_id, channel.id) + def delete_tasks_for(attached_record) do + tasks = + case attached_record do + %Channel{} = channel -> list_tasks_for(:channel_id, channel.id) + %MediaItem{} = media_item -> list_tasks_for(:media_item_id, media_item.id) + end Enum.each(tasks, fn task -> delete_task(task) @@ -112,8 +127,12 @@ defmodule Pinchflat.Tasks do Returns :ok """ - def delete_pending_tasks_for(%Channel{} = channel) do - tasks = list_pending_tasks_for(:channel_id, channel.id) + def delete_pending_tasks_for(attached_record) do + tasks = + case attached_record do + %Channel{} = channel -> list_pending_tasks_for(:channel_id, channel.id) + %MediaItem{} = media_item -> list_pending_tasks_for(:media_item_id, media_item.id) + end Enum.each(tasks, fn task -> delete_task(task) diff --git a/lib/pinchflat/tasks/channel_tasks.ex b/lib/pinchflat/tasks/channel_tasks.ex index 095c90f..a2b1cd2 100644 --- a/lib/pinchflat/tasks/channel_tasks.ex +++ b/lib/pinchflat/tasks/channel_tasks.ex @@ -8,7 +8,9 @@ defmodule Pinchflat.Tasks.ChannelTasks do alias Pinchflat.Workers.MediaIndexingWorker @doc """ - Starts tasks for indexing a channel's media. Returns {:ok, :should_not_index} | {:ok, %Task{}}. + Starts tasks for indexing a channel's media. + + Returns {:ok, :should_not_index} | {:ok, %Task{}}. """ def kickoff_indexing_task(%Channel{} = channel) do Tasks.delete_pending_tasks_for(channel) @@ -21,6 +23,11 @@ defmodule Pinchflat.Tasks.ChannelTasks do # Schedule this one immediately, but future ones will be on an interval |> MediaIndexingWorker.new() |> Tasks.create_job_with_task(channel) + |> case do + # This should never return {:error, :duplicate_job} since we just deleted + # any pending tasks. I'm being assertive about it so it's obvious if I'm wrong + {:ok, task} -> {:ok, task} + end end end end diff --git a/lib/pinchflat/tasks/task.ex b/lib/pinchflat/tasks/task.ex index 4184ee8..9d53b60 100644 --- a/lib/pinchflat/tasks/task.ex +++ b/lib/pinchflat/tasks/task.ex @@ -6,11 +6,13 @@ defmodule Pinchflat.Tasks.Task do use Ecto.Schema import Ecto.Changeset + alias Pinchflat.Media.MediaItem alias Pinchflat.MediaSource.Channel schema "tasks" do belongs_to :job, Oban.Job belongs_to :channel, Channel + belongs_to :media_item, MediaItem timestamps(type: :utc_datetime) end @@ -18,7 +20,7 @@ defmodule Pinchflat.Tasks.Task do @doc false def changeset(task, attrs) do task - |> cast(attrs, [:job_id, :channel_id]) + |> cast(attrs, [:job_id, :channel_id, :media_item_id]) |> validate_required([:job_id]) end end diff --git a/lib/pinchflat/workers/media_indexing_worker.ex b/lib/pinchflat/workers/media_indexing_worker.ex index c377147..4869b0d 100644 --- a/lib/pinchflat/workers/media_indexing_worker.ex +++ b/lib/pinchflat/workers/media_indexing_worker.ex @@ -7,9 +7,9 @@ defmodule Pinchflat.Workers.MediaIndexingWorker do tags: ["media_source", "media_indexing"] alias __MODULE__ + alias Pinchflat.Media alias Pinchflat.Tasks alias Pinchflat.MediaSource - alias Pinchflat.Media.MediaItem alias Pinchflat.Workers.VideoDownloadWorker @impl Oban.Worker @@ -48,24 +48,33 @@ defmodule Pinchflat.Workers.MediaIndexingWorker do end defp index_media_and_reschedule(channel) do - channel - |> MediaSource.index_media_items() - |> Enum.each(fn media_item_or_changeset -> - case media_item_or_changeset do - %MediaItem{} = media_item -> - media_item - |> Map.take([:id]) - |> VideoDownloadWorker.new() - |> Oban.insert() - - _ -> - nil - end - end) + MediaSource.index_media_items(channel) + enqueue_video_downloads(channel) channel |> Map.take([:id]) |> MediaIndexingWorker.new(schedule_in: channel.index_frequency_minutes * 60) |> Tasks.create_job_with_task(channel) + |> case do + {:ok, task} -> {:ok, task} + {:error, :duplicate_job} -> {:ok, :job_exists} + end + end + + # NOTE: this starts a download for each media item that is pending, + # not just the ones that were indexed in this job run. This should ensure + # that any stragglers are caught if, for some reason, they weren't enqueued + # or somehow got de-queued. + # + # I'm not sure of a case where this would happen, but it's cheap insurance. + defp enqueue_video_downloads(channel) do + channel + |> Media.list_pending_media_items_for() + |> Enum.each(fn media_item -> + media_item + |> Map.take([:id]) + |> VideoDownloadWorker.new() + |> Tasks.create_job_with_task(media_item) + end) end end diff --git a/lib/pinchflat/workers/video_download_worker.ex b/lib/pinchflat/workers/video_download_worker.ex index 338cba3..f0047c1 100644 --- a/lib/pinchflat/workers/video_download_worker.ex +++ b/lib/pinchflat/workers/video_download_worker.ex @@ -3,8 +3,8 @@ defmodule Pinchflat.Workers.VideoDownloadWorker do use Oban.Worker, queue: :media_fetching, - unique: [period: :infinity, states: [:available, :scheduled, :retryable]], - tags: ["media_itwm", "media_fetching"] + unique: [period: :infinity, states: [:available, :scheduled, :retryable, :executing]], + tags: ["media_item", "media_fetching"] alias Pinchflat.Media alias Pinchflat.MediaClient.VideoDownloader @@ -19,11 +19,8 @@ defmodule Pinchflat.Workers.VideoDownloadWorker do media_item = Media.get_media_item!(media_item_id) case VideoDownloader.download_for_media_item(media_item) do - {:ok, _} -> - {:ok, media_item} - - err -> - err + {:ok, _} -> {:ok, media_item} + err -> err end end end diff --git a/priv/repo/migrations/20240130161316_add_media_item_to_tasks.exs b/priv/repo/migrations/20240130161316_add_media_item_to_tasks.exs new file mode 100644 index 0000000..bd8a9c9 --- /dev/null +++ b/priv/repo/migrations/20240130161316_add_media_item_to_tasks.exs @@ -0,0 +1,12 @@ +defmodule Pinchflat.Repo.Migrations.AddMediaItemToTasks do + use Ecto.Migration + + def change do + alter table(:tasks) do + # `restrict` because we need to be sure to delete pending tasks when a channel is deleted + add :media_item_id, references(:media_items, on_delete: :restrict), null: true + end + + create index(:tasks, [:media_item_id]) + end +end diff --git a/test/pinchflat/media_test.exs b/test/pinchflat/media_test.exs index 8d7ca72..d7a99cf 100644 --- a/test/pinchflat/media_test.exs +++ b/test/pinchflat/media_test.exs @@ -1,6 +1,7 @@ defmodule Pinchflat.MediaTest do use Pinchflat.DataCase + import Pinchflat.TasksFixtures import Pinchflat.MediaFixtures import Pinchflat.MediaSourceFixtures @@ -28,6 +29,27 @@ defmodule Pinchflat.MediaTest do end end + describe "list_pending_media_items_for/1" do + test "it returns pending media_items for a given channel" do + channel = channel_fixture() + media_item = media_item_fixture(%{channel_id: channel.id, video_filepath: nil}) + + assert Media.list_pending_media_items_for(channel) == [media_item] + end + + test "it does not return media_items with video_filepath" do + channel = channel_fixture() + + _media_item = + media_item_fixture(%{ + channel_id: channel.id, + video_filepath: "/video/#{Faker.File.file_name(:video)}" + }) + + assert Media.list_pending_media_items_for(channel) == [] + end + end + describe "get_media_item!/1" do test "it returns the media_item with given id" do media_item = media_item_fixture() @@ -85,6 +107,14 @@ defmodule Pinchflat.MediaTest do assert {:ok, %MediaItem{}} = Media.delete_media_item(media_item) assert_raise Ecto.NoResultsError, fn -> Media.get_media_item!(media_item.id) end end + + test "it also deletes attached tasks" do + media_item = media_item_fixture() + task = task_fixture(%{media_item_id: media_item.id}) + + assert {:ok, %MediaItem{}} = Media.delete_media_item(media_item) + assert_raise Ecto.NoResultsError, fn -> Repo.reload!(task) end + end end describe "change_media_item/1" do diff --git a/test/pinchflat/repo_test.exs b/test/pinchflat/repo_test.exs new file mode 100644 index 0000000..0d7eea2 --- /dev/null +++ b/test/pinchflat/repo_test.exs @@ -0,0 +1,26 @@ +defmodule Pinchflat.RepoTest do + use Pinchflat.DataCase + + alias Pinchflat.JobFixtures.TestJobWorker + + describe "insert_unique_job/1" do + test "returns {:ok, job} if there is no conflict" do + job = TestJobWorker.new(%{}) + + assert {:ok, %Oban.Job{}} = Pinchflat.Repo.insert_unique_job(job) + end + + test "returns {:duplicate, original_job} if there is a conflict" do + job = TestJobWorker.new(%{foo: "bar"}, unique: [period: :infinity]) + + {:ok, saved_job_1} = Pinchflat.Repo.insert_unique_job(job) + + assert {:duplicate, saved_job_2} = Pinchflat.Repo.insert_unique_job(job) + assert saved_job_1.id == saved_job_2.id + end + + test "returns the error if there is an error" do + assert {:error, _} = Pinchflat.Repo.insert_unique_job(%Ecto.Changeset{}) + end + end +end diff --git a/test/pinchflat/tasks_test.exs b/test/pinchflat/tasks_test.exs index b19bf84..fbcc670 100644 --- a/test/pinchflat/tasks_test.exs +++ b/test/pinchflat/tasks_test.exs @@ -2,6 +2,7 @@ defmodule Pinchflat.TasksTest do use Pinchflat.DataCase import Pinchflat.JobFixtures import Pinchflat.TasksFixtures + import Pinchflat.MediaFixtures import Pinchflat.MediaSourceFixtures alias Pinchflat.Tasks @@ -92,14 +93,24 @@ defmodule Pinchflat.TasksTest do assert task.job_id == job.id assert task.channel_id == channel.id end + + test "accepts a job and media item" do + job = job_fixture() + media_item = media_item_fixture() + + assert {:ok, %Task{} = task} = Tasks.create_task(job, media_item) + + assert task.job_id == job.id + assert task.media_item_id == media_item.id + end end describe "create_job_with_task/2" do test "it enqueues the given job" do - channel = channel_fixture() + media_item = media_item_fixture() refute_enqueued(worker: TestJobWorker) - assert {:ok, %Task{}} = Tasks.create_job_with_task(TestJobWorker.new(%{}), channel) + assert {:ok, %Task{}} = Tasks.create_job_with_task(TestJobWorker.new(%{}), media_item) assert_enqueued(worker: TestJobWorker) end @@ -111,6 +122,14 @@ defmodule Pinchflat.TasksTest do assert task.channel_id == channel.id end + test "it returns an error if the job already exists" do + channel = channel_fixture() + job = TestJobWorker.new(%{foo: "bar"}, unique: [period: :infinity]) + + assert {:ok, %Task{}} = Tasks.create_job_with_task(job, channel) + assert {:error, :duplicate_job} = Tasks.create_job_with_task(job, channel) + end + test "it returns an error if the job fails to enqueue" do channel = channel_fixture() @@ -143,6 +162,14 @@ defmodule Pinchflat.TasksTest do assert :ok = Tasks.delete_tasks_for(channel) assert_raise Ecto.NoResultsError, fn -> Tasks.get_task!(task.id) end end + + test "it deletes the tasks attached to a media_item" do + media_item = media_item_fixture() + task = task_fixture(media_item_id: media_item.id) + + assert :ok = Tasks.delete_tasks_for(media_item) + assert_raise Ecto.NoResultsError, fn -> Tasks.get_task!(task.id) end + end end describe "delete_pending_tasks_for/1" do @@ -162,6 +189,17 @@ defmodule Pinchflat.TasksTest do assert :ok = Tasks.delete_pending_tasks_for(channel) assert Tasks.get_task!(task.id) end + + test "it works on media_items" do + media_item = media_item_fixture() + pending_task = task_fixture(media_item_id: media_item.id) + cancelled_task = Repo.preload(task_fixture(media_item_id: media_item.id), :job) + :ok = Oban.cancel_job(cancelled_task.job) + + assert :ok = Tasks.delete_pending_tasks_for(media_item) + assert Tasks.get_task!(cancelled_task.id) + assert_raise Ecto.NoResultsError, fn -> Repo.reload!(pending_task) end + end end describe "change_task/1" do diff --git a/test/pinchflat/workers/media_indexing_worker_test.exs b/test/pinchflat/workers/media_indexing_worker_test.exs index 2683e1e..154f9ab 100644 --- a/test/pinchflat/workers/media_indexing_worker_test.exs +++ b/test/pinchflat/workers/media_indexing_worker_test.exs @@ -2,6 +2,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do use Pinchflat.DataCase import Mox + import Pinchflat.MediaFixtures import Pinchflat.MediaSourceFixtures alias Pinchflat.Tasks @@ -36,7 +37,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do perform_job(MediaIndexingWorker, %{id: channel.id}) end - test "it kicks off a download job for each new media item" do + test "it kicks off a download job for each pending media item" do expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, "video1"} end) channel = channel_fixture(index_frequency_minutes: 10) @@ -45,6 +46,16 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do assert [_] = all_enqueued(worker: VideoDownloadWorker) end + test "it starts a job for any pending media item even if it's from another run" do + expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, "video1"} end) + + channel = channel_fixture(index_frequency_minutes: 10) + media_item_fixture(%{channel_id: channel.id, video_filepath: nil}) + perform_job(MediaIndexingWorker, %{id: channel.id}) + + assert [_, _] = all_enqueued(worker: VideoDownloadWorker) + end + test "it does not kick off a job for media items that could not be saved" do expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, "video1\nvideo1"} end) diff --git a/test/pinchflat/workers/video_download_worker_test.exs b/test/pinchflat/workers/video_download_worker_test.exs index 3c09820..d4f0995 100644 --- a/test/pinchflat/workers/video_download_worker_test.exs +++ b/test/pinchflat/workers/video_download_worker_test.exs @@ -38,5 +38,22 @@ defmodule Pinchflat.Workers.VideoDownloadWorkerTest do perform_job(VideoDownloadWorker, %{id: media_item.id}) assert Repo.reload(media_item).metadata != nil end + + test "it won't double-schedule downloading jobs", %{media_item: media_item} do + Oban.insert(VideoDownloadWorker.new(%{id: media_item.id})) + Oban.insert(VideoDownloadWorker.new(%{id: media_item.id})) + + assert [_] = all_enqueued(worker: VideoDownloadWorker) + end + + test "it sets the job to retryable if the download fails", %{media_item: media_item} do + expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:error, "error"} end) + + Oban.Testing.with_testing_mode(:inline, fn -> + {:ok, job} = Oban.insert(VideoDownloadWorker.new(%{id: media_item.id})) + + assert job.state == "retryable" + end) + end end end