Adds tasks for media items

This commit is contained in:
Kieran Eglin 2024-01-30 10:55:54 -08:00
parent a0d62a156a
commit cfdba60b67
No known key found for this signature in database
GPG key ID: 193984967FCF432D
14 changed files with 243 additions and 39 deletions

View file

@ -6,7 +6,9 @@ defmodule Pinchflat.Media do
import Ecto.Query, warn: false import Ecto.Query, warn: false
alias Pinchflat.Repo alias Pinchflat.Repo
alias Pinchflat.Tasks
alias Pinchflat.Media.MediaItem alias Pinchflat.Media.MediaItem
alias Pinchflat.MediaSource.Channel
@doc """ @doc """
Returns the list of media_items. Returns [%MediaItem{}, ...]. Returns the list of media_items. Returns [%MediaItem{}, ...].
@ -15,6 +17,20 @@ defmodule Pinchflat.Media do
Repo.all(MediaItem) Repo.all(MediaItem)
end 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 """ @doc """
Gets a single media_item. Gets a single media_item.
@ -41,9 +57,12 @@ defmodule Pinchflat.Media do
end end
@doc """ @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 def delete_media_item(%MediaItem{} = media_item) do
Tasks.delete_tasks_for(media_item)
Repo.delete(media_item) Repo.delete(media_item)
end end

View file

@ -6,6 +6,7 @@ defmodule Pinchflat.Media.MediaItem do
use Ecto.Schema use Ecto.Schema
import Ecto.Changeset import Ecto.Changeset
alias Pinchflat.Tasks.Task
alias Pinchflat.MediaSource.Channel alias Pinchflat.MediaSource.Channel
alias Pinchflat.Media.MediaMetadata alias Pinchflat.Media.MediaMetadata
@ -21,6 +22,8 @@ defmodule Pinchflat.Media.MediaItem do
has_one :metadata, MediaMetadata, on_replace: :update has_one :metadata, MediaMetadata, on_replace: :update
has_many :tasks, Task
timestamps(type: :utc_datetime) timestamps(type: :utc_datetime)
end end

View file

@ -2,4 +2,18 @@ defmodule Pinchflat.Repo do
use Ecto.Repo, use Ecto.Repo,
otp_app: :pinchflat, otp_app: :pinchflat,
adapter: Ecto.Adapters.Postgres 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 end

View file

@ -7,6 +7,7 @@ defmodule Pinchflat.Tasks do
alias Pinchflat.Repo alias Pinchflat.Repo
alias Pinchflat.Tasks.Task alias Pinchflat.Tasks.Task
alias Pinchflat.Media.MediaItem
alias Pinchflat.MediaSource.Channel alias Pinchflat.MediaSource.Channel
@doc """ @doc """
@ -54,7 +55,10 @@ defmodule Pinchflat.Tasks do
def get_task!(id), do: Repo.get!(Task, id) def get_task!(id), do: Repo.get!(Task, id)
@doc """ @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 def create_task(attrs) do
%Task{} %Task{}
@ -62,23 +66,30 @@ defmodule Pinchflat.Tasks do
|> Repo.insert() |> Repo.insert()
end 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` # 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{}
|> Task.changeset(%{job_id: job.id, channel_id: channel.id}) |> Task.changeset(Map.merge(%{job_id: job.id}, attached_record_attr))
|> Repo.insert() |> Repo.insert()
end end
@doc """ @doc """
Creates a job from given attrs, creating a task with an attached record 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 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) {:ok, job} -> create_task(job, task_attached_record)
{:duplicate, _} -> {:error, :duplicate_job}
err -> err err -> err
end end
end end
@ -99,8 +110,12 @@ defmodule Pinchflat.Tasks do
Returns :ok Returns :ok
""" """
def delete_tasks_for(%Channel{} = channel) do def delete_tasks_for(attached_record) do
tasks = list_tasks_for(:channel_id, channel.id) 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 -> Enum.each(tasks, fn task ->
delete_task(task) delete_task(task)
@ -112,8 +127,12 @@ defmodule Pinchflat.Tasks do
Returns :ok Returns :ok
""" """
def delete_pending_tasks_for(%Channel{} = channel) do def delete_pending_tasks_for(attached_record) do
tasks = list_pending_tasks_for(:channel_id, channel.id) 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 -> Enum.each(tasks, fn task ->
delete_task(task) delete_task(task)

View file

@ -8,7 +8,9 @@ defmodule Pinchflat.Tasks.ChannelTasks do
alias Pinchflat.Workers.MediaIndexingWorker alias Pinchflat.Workers.MediaIndexingWorker
@doc """ @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 def kickoff_indexing_task(%Channel{} = channel) do
Tasks.delete_pending_tasks_for(channel) 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 # Schedule this one immediately, but future ones will be on an interval
|> MediaIndexingWorker.new() |> MediaIndexingWorker.new()
|> Tasks.create_job_with_task(channel) |> 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 end
end end

View file

@ -6,11 +6,13 @@ defmodule Pinchflat.Tasks.Task do
use Ecto.Schema use Ecto.Schema
import Ecto.Changeset import Ecto.Changeset
alias Pinchflat.Media.MediaItem
alias Pinchflat.MediaSource.Channel alias Pinchflat.MediaSource.Channel
schema "tasks" do schema "tasks" do
belongs_to :job, Oban.Job belongs_to :job, Oban.Job
belongs_to :channel, Channel belongs_to :channel, Channel
belongs_to :media_item, MediaItem
timestamps(type: :utc_datetime) timestamps(type: :utc_datetime)
end end
@ -18,7 +20,7 @@ defmodule Pinchflat.Tasks.Task do
@doc false @doc false
def changeset(task, attrs) do def changeset(task, attrs) do
task task
|> cast(attrs, [:job_id, :channel_id]) |> cast(attrs, [:job_id, :channel_id, :media_item_id])
|> validate_required([:job_id]) |> validate_required([:job_id])
end end
end end

View file

@ -7,9 +7,9 @@ defmodule Pinchflat.Workers.MediaIndexingWorker do
tags: ["media_source", "media_indexing"] tags: ["media_source", "media_indexing"]
alias __MODULE__ alias __MODULE__
alias Pinchflat.Media
alias Pinchflat.Tasks alias Pinchflat.Tasks
alias Pinchflat.MediaSource alias Pinchflat.MediaSource
alias Pinchflat.Media.MediaItem
alias Pinchflat.Workers.VideoDownloadWorker alias Pinchflat.Workers.VideoDownloadWorker
@impl Oban.Worker @impl Oban.Worker
@ -48,24 +48,33 @@ defmodule Pinchflat.Workers.MediaIndexingWorker do
end end
defp index_media_and_reschedule(channel) do defp index_media_and_reschedule(channel) do
channel MediaSource.index_media_items(channel)
|> MediaSource.index_media_items() enqueue_video_downloads(channel)
|> 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)
channel channel
|> Map.take([:id]) |> Map.take([:id])
|> MediaIndexingWorker.new(schedule_in: channel.index_frequency_minutes * 60) |> MediaIndexingWorker.new(schedule_in: channel.index_frequency_minutes * 60)
|> Tasks.create_job_with_task(channel) |> 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
end end

View file

@ -3,8 +3,8 @@ defmodule Pinchflat.Workers.VideoDownloadWorker do
use Oban.Worker, use Oban.Worker,
queue: :media_fetching, queue: :media_fetching,
unique: [period: :infinity, states: [:available, :scheduled, :retryable]], unique: [period: :infinity, states: [:available, :scheduled, :retryable, :executing]],
tags: ["media_itwm", "media_fetching"] tags: ["media_item", "media_fetching"]
alias Pinchflat.Media alias Pinchflat.Media
alias Pinchflat.MediaClient.VideoDownloader alias Pinchflat.MediaClient.VideoDownloader
@ -19,11 +19,8 @@ defmodule Pinchflat.Workers.VideoDownloadWorker do
media_item = Media.get_media_item!(media_item_id) media_item = Media.get_media_item!(media_item_id)
case VideoDownloader.download_for_media_item(media_item) do case VideoDownloader.download_for_media_item(media_item) do
{:ok, _} -> {:ok, _} -> {:ok, media_item}
{:ok, media_item} err -> err
err ->
err
end end
end end
end end

View file

@ -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

View file

@ -1,6 +1,7 @@
defmodule Pinchflat.MediaTest do defmodule Pinchflat.MediaTest do
use Pinchflat.DataCase use Pinchflat.DataCase
import Pinchflat.TasksFixtures
import Pinchflat.MediaFixtures import Pinchflat.MediaFixtures
import Pinchflat.MediaSourceFixtures import Pinchflat.MediaSourceFixtures
@ -28,6 +29,27 @@ defmodule Pinchflat.MediaTest do
end end
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 describe "get_media_item!/1" do
test "it returns the media_item with given id" do test "it returns the media_item with given id" do
media_item = media_item_fixture() media_item = media_item_fixture()
@ -85,6 +107,14 @@ defmodule Pinchflat.MediaTest do
assert {:ok, %MediaItem{}} = Media.delete_media_item(media_item) assert {:ok, %MediaItem{}} = Media.delete_media_item(media_item)
assert_raise Ecto.NoResultsError, fn -> Media.get_media_item!(media_item.id) end assert_raise Ecto.NoResultsError, fn -> Media.get_media_item!(media_item.id) end
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 end
describe "change_media_item/1" do describe "change_media_item/1" do

View file

@ -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

View file

@ -2,6 +2,7 @@ defmodule Pinchflat.TasksTest do
use Pinchflat.DataCase use Pinchflat.DataCase
import Pinchflat.JobFixtures import Pinchflat.JobFixtures
import Pinchflat.TasksFixtures import Pinchflat.TasksFixtures
import Pinchflat.MediaFixtures
import Pinchflat.MediaSourceFixtures import Pinchflat.MediaSourceFixtures
alias Pinchflat.Tasks alias Pinchflat.Tasks
@ -92,14 +93,24 @@ defmodule Pinchflat.TasksTest do
assert task.job_id == job.id assert task.job_id == job.id
assert task.channel_id == channel.id assert task.channel_id == channel.id
end 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 end
describe "create_job_with_task/2" do describe "create_job_with_task/2" do
test "it enqueues the given job" do test "it enqueues the given job" do
channel = channel_fixture() media_item = media_item_fixture()
refute_enqueued(worker: TestJobWorker) 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) assert_enqueued(worker: TestJobWorker)
end end
@ -111,6 +122,14 @@ defmodule Pinchflat.TasksTest do
assert task.channel_id == channel.id assert task.channel_id == channel.id
end 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 test "it returns an error if the job fails to enqueue" do
channel = channel_fixture() channel = channel_fixture()
@ -143,6 +162,14 @@ defmodule Pinchflat.TasksTest do
assert :ok = Tasks.delete_tasks_for(channel) assert :ok = Tasks.delete_tasks_for(channel)
assert_raise Ecto.NoResultsError, fn -> Tasks.get_task!(task.id) end assert_raise Ecto.NoResultsError, fn -> Tasks.get_task!(task.id) end
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 end
describe "delete_pending_tasks_for/1" do describe "delete_pending_tasks_for/1" do
@ -162,6 +189,17 @@ defmodule Pinchflat.TasksTest do
assert :ok = Tasks.delete_pending_tasks_for(channel) assert :ok = Tasks.delete_pending_tasks_for(channel)
assert Tasks.get_task!(task.id) assert Tasks.get_task!(task.id)
end 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 end
describe "change_task/1" do describe "change_task/1" do

View file

@ -2,6 +2,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
use Pinchflat.DataCase use Pinchflat.DataCase
import Mox import Mox
import Pinchflat.MediaFixtures
import Pinchflat.MediaSourceFixtures import Pinchflat.MediaSourceFixtures
alias Pinchflat.Tasks alias Pinchflat.Tasks
@ -36,7 +37,7 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
perform_job(MediaIndexingWorker, %{id: channel.id}) perform_job(MediaIndexingWorker, %{id: channel.id})
end 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) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, "video1"} end)
channel = channel_fixture(index_frequency_minutes: 10) channel = channel_fixture(index_frequency_minutes: 10)
@ -45,6 +46,16 @@ defmodule Pinchflat.Workers.MediaIndexingWorkerTest do
assert [_] = all_enqueued(worker: VideoDownloadWorker) assert [_] = all_enqueued(worker: VideoDownloadWorker)
end 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 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) expect(YtDlpRunnerMock, :run, fn _url, _opts, _ot -> {:ok, "video1\nvideo1"} end)

View file

@ -38,5 +38,22 @@ defmodule Pinchflat.Workers.VideoDownloadWorkerTest do
perform_job(VideoDownloadWorker, %{id: media_item.id}) perform_job(VideoDownloadWorker, %{id: media_item.id})
assert Repo.reload(media_item).metadata != nil assert Repo.reload(media_item).metadata != nil
end 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
end end