Added the ability to set priority to various downloading helpers

This commit is contained in:
Kieran Eglin 2025-01-21 15:13:46 -08:00
parent 8c8be10083
commit cb5cf3833d
No known key found for this signature in database
GPG key ID: 193984967FCF432D
5 changed files with 68 additions and 34 deletions

View file

@ -27,13 +27,15 @@ defmodule Pinchflat.Downloading.DownloadingHelpers do
Returns :ok Returns :ok
""" """
def enqueue_pending_download_tasks(%Source{download_media: true} = source) do def enqueue_pending_download_tasks(source, job_opts \\ [])
def enqueue_pending_download_tasks(%Source{download_media: true} = source, job_opts) do
source source
|> Media.list_pending_media_items_for() |> Media.list_pending_media_items_for()
|> Enum.each(&MediaDownloadWorker.kickoff_with_task/1) |> Enum.each(&MediaDownloadWorker.kickoff_with_task(&1, %{}, job_opts))
end end
def enqueue_pending_download_tasks(%Source{download_media: false}) do def enqueue_pending_download_tasks(%Source{download_media: false}, _job_opts) do
:ok :ok
end end
@ -55,13 +57,13 @@ defmodule Pinchflat.Downloading.DownloadingHelpers do
Returns {:ok, %Task{}} | {:error, :should_not_download} | {:error, any()} Returns {:ok, %Task{}} | {:error, :should_not_download} | {:error, any()}
""" """
def kickoff_download_if_pending(%MediaItem{} = media_item) do def kickoff_download_if_pending(%MediaItem{} = media_item, job_opts \\ []) do
media_item = Repo.preload(media_item, :source) media_item = Repo.preload(media_item, :source)
if media_item.source.download_media && Media.pending_download?(media_item) do if media_item.source.download_media && Media.pending_download?(media_item) do
Logger.info("Kicking off download for media item ##{media_item.id} (#{media_item.media_id})") Logger.info("Kicking off download for media item ##{media_item.id} (#{media_item.media_id})")
MediaDownloadWorker.kickoff_with_task(media_item) MediaDownloadWorker.kickoff_with_task(media_item, %{}, job_opts)
else else
{:error, :should_not_download} {:error, :should_not_download}
end end

View file

@ -40,7 +40,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpers do
Returns [%MediaItem{}] where each item is a new media item that was created _but not necessarily Returns [%MediaItem{}] where each item is a new media item that was created _but not necessarily
downloaded_. downloaded_.
""" """
def kickoff_download_tasks_from_youtube_rss_feed(%Source{} = source) do def index_and_kickoff_downloads(%Source{} = source) do
# The media_profile is needed to determine the quality options to _then_ determine a more # The media_profile is needed to determine the quality options to _then_ determine a more
# accurate predicted filepath # accurate predicted filepath
source = Repo.preload(source, [:media_profile]) source = Repo.preload(source, [:media_profile])
@ -53,6 +53,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpers do
Enum.map(new_media_ids, fn media_id -> Enum.map(new_media_ids, fn media_id ->
case create_media_item_from_media_id(source, media_id) do case create_media_item_from_media_id(source, media_id) do
{:ok, media_item} -> {:ok, media_item} ->
DownloadingHelpers.kickoff_download_if_pending(media_item, priority: 0)
media_item media_item
err -> err ->
@ -61,7 +62,9 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpers do
end end
end) end)
DownloadingHelpers.enqueue_pending_download_tasks(source) # Pick up any stragglers. Intentionally has a lower priority than the per-media item
# kickoff above
DownloadingHelpers.enqueue_pending_download_tasks(source, priority: 1)
Enum.filter(maybe_new_media_items, & &1) Enum.filter(maybe_new_media_items, & &1)
end end

View file

@ -38,8 +38,8 @@ defmodule Pinchflat.FastIndexing.FastIndexingWorker do
Order of operations: Order of operations:
1. FastIndexingWorker (this module) periodically checks the YouTube RSS feed for new media. 1. FastIndexingWorker (this module) periodically checks the YouTube RSS feed for new media.
with `FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed` with `FastIndexingHelpers.index_and_kickoff_downloads`
2. If the above `kickoff_download_tasks_from_youtube_rss_feed` finds new media items in the RSS feed, 2. If the above `index_and_kickoff_downloads` finds new media items in the RSS feed,
it indexes them with a yt-dlp call to create the media item records then kicks off downloading it indexes them with a yt-dlp call to create the media item records then kicks off downloading
tasks (MediaDownloadWorker) for any new media items _that should be downloaded_. tasks (MediaDownloadWorker) for any new media items _that should be downloaded_.
3. Once downloads are kicked off, this worker sends a notification to the apprise server if applicable 3. Once downloads are kicked off, this worker sends a notification to the apprise server if applicable
@ -67,7 +67,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingWorker do
new_media_items = new_media_items =
source source
|> FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed() |> FastIndexingHelpers.index_and_kickoff_downloads()
|> Enum.filter(&Media.pending_download?(&1)) |> Enum.filter(&Media.pending_download?(&1))
if source.download_media do if source.download_media do

View file

@ -10,7 +10,7 @@ defmodule Pinchflat.Downloading.DownloadingHelpersTest do
alias Pinchflat.Downloading.MediaDownloadWorker alias Pinchflat.Downloading.MediaDownloadWorker
describe "enqueue_pending_download_tasks/1" do describe "enqueue_pending_download_tasks/1" do
test "it enqueues a job for each pending media item" do test "enqueues a job for each pending media item" do
source = source_fixture() source = source_fixture()
media_item = media_item_fixture(source_id: source.id, media_filepath: nil) media_item = media_item_fixture(source_id: source.id, media_filepath: nil)
@ -19,7 +19,7 @@ defmodule Pinchflat.Downloading.DownloadingHelpersTest do
assert_enqueued(worker: MediaDownloadWorker, args: %{"id" => media_item.id}) assert_enqueued(worker: MediaDownloadWorker, args: %{"id" => media_item.id})
end end
test "it does not enqueue a job for media items with a filepath" do test "does not enqueue a job for media items with a filepath" do
source = source_fixture() source = source_fixture()
_media_item = media_item_fixture(source_id: source.id, media_filepath: "some/filepath.mp4") _media_item = media_item_fixture(source_id: source.id, media_filepath: "some/filepath.mp4")
@ -28,7 +28,7 @@ defmodule Pinchflat.Downloading.DownloadingHelpersTest do
refute_enqueued(worker: MediaDownloadWorker) refute_enqueued(worker: MediaDownloadWorker)
end end
test "it attaches a task to each enqueued job" do test "attaches a task to each enqueued job" do
source = source_fixture() source = source_fixture()
media_item = media_item_fixture(source_id: source.id, media_filepath: nil) media_item = media_item_fixture(source_id: source.id, media_filepath: nil)
@ -39,7 +39,7 @@ defmodule Pinchflat.Downloading.DownloadingHelpersTest do
assert [_] = Tasks.list_tasks_for(media_item) assert [_] = Tasks.list_tasks_for(media_item)
end end
test "it does not create a job if the source is set to not download" do test "does not create a job if the source is set to not download" do
source = source_fixture(download_media: false) source = source_fixture(download_media: false)
assert :ok = DownloadingHelpers.enqueue_pending_download_tasks(source) assert :ok = DownloadingHelpers.enqueue_pending_download_tasks(source)
@ -47,17 +47,26 @@ defmodule Pinchflat.Downloading.DownloadingHelpersTest do
refute_enqueued(worker: MediaDownloadWorker) refute_enqueued(worker: MediaDownloadWorker)
end end
test "it does not attach tasks if the source is set to not download" do test "does not attach tasks if the source is set to not download" do
source = source_fixture(download_media: false) source = source_fixture(download_media: false)
media_item = media_item_fixture(source_id: source.id, media_filepath: nil) media_item = media_item_fixture(source_id: source.id, media_filepath: nil)
assert :ok = DownloadingHelpers.enqueue_pending_download_tasks(source) assert :ok = DownloadingHelpers.enqueue_pending_download_tasks(source)
assert [] = Tasks.list_tasks_for(media_item) assert [] = Tasks.list_tasks_for(media_item)
end end
test "can pass job options" do
source = source_fixture()
media_item = media_item_fixture(source_id: source.id, media_filepath: nil)
assert :ok = DownloadingHelpers.enqueue_pending_download_tasks(source, priority: 1)
assert_enqueued(worker: MediaDownloadWorker, args: %{"id" => media_item.id}, priority: 1)
end
end end
describe "dequeue_pending_download_tasks/1" do describe "dequeue_pending_download_tasks/1" do
test "it deletes all pending tasks for a source's media items" do test "deletes all pending tasks for a source's media items" do
source = source_fixture() source = source_fixture()
media_item = media_item_fixture(source_id: source.id, media_filepath: nil) media_item = media_item_fixture(source_id: source.id, media_filepath: nil)
@ -109,6 +118,14 @@ defmodule Pinchflat.Downloading.DownloadingHelpersTest do
refute_enqueued(worker: MediaDownloadWorker) refute_enqueued(worker: MediaDownloadWorker)
end end
test "can pass job options" do
media_item = media_item_fixture(media_filepath: nil)
assert {:ok, _} = DownloadingHelpers.kickoff_download_if_pending(media_item, priority: 1)
assert_enqueued(worker: MediaDownloadWorker, args: %{"id" => media_item.id}, priority: 1)
end
end end
describe "kickoff_redownload_for_existing_media/1" do describe "kickoff_redownload_for_existing_media/1" do

View file

@ -38,36 +38,48 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
end end
end end
describe "kickoff_download_tasks_from_youtube_rss_feed/1" do describe "index_and_kickoff_downloads/1" do
test "enqueues a new worker for each new media_id in the source's RSS feed", %{source: source} do test "enqueues a worker for each new media_id in the source's RSS feed", %{source: source} do
expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end) expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end)
assert [media_item] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [media_item] = FastIndexingHelpers.index_and_kickoff_downloads(source)
assert [worker] = all_enqueued(worker: MediaDownloadWorker) assert [worker] = all_enqueued(worker: MediaDownloadWorker)
assert worker.args["id"] == media_item.id assert worker.args["id"] == media_item.id
assert worker.priority == 0
end end
test "does not enqueue a new worker for the source's media IDs we already know about", %{source: source} do test "does not enqueue a new worker for the source's media IDs we already know about", %{source: source} do
expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end) expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end)
media_item_fixture(source_id: source.id, media_id: "test_1") media_item_fixture(source_id: source.id, media_id: "test_1")
assert [] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [] = FastIndexingHelpers.index_and_kickoff_downloads(source)
refute_enqueued(worker: MediaDownloadWorker) refute_enqueued(worker: MediaDownloadWorker)
end end
test "kicks off a download task for all pending media but at a lower priority", %{source: source} do
pending_item = media_item_fixture(source_id: source.id, media_filepath: nil)
expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end)
assert [%MediaItem{}] = FastIndexingHelpers.index_and_kickoff_downloads(source)
assert [worker_1, _worker_2] = all_enqueued(worker: MediaDownloadWorker)
assert worker_1.args["id"] == pending_item.id
assert worker_1.priority == 1
end
test "returns the found media items", %{source: source} do test "returns the found media items", %{source: source} do
expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end) expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end)
assert [%MediaItem{}] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [%MediaItem{}] = FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
test "does not enqueue a download job if the source does not allow it" do test "does not enqueue a download job if the source does not allow it" do
expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end) expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end)
source = source_fixture(%{download_media: false}) source = source_fixture(%{download_media: false})
assert [%MediaItem{}] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [%MediaItem{}] = FastIndexingHelpers.index_and_kickoff_downloads(source)
refute_enqueued(worker: MediaDownloadWorker) refute_enqueued(worker: MediaDownloadWorker)
end end
@ -75,7 +87,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
test "creates a download task record", %{source: source} do test "creates a download task record", %{source: source} do
expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end) expect(HTTPClientMock, :get, fn _url -> {:ok, "<yt:videoId>test_1</yt:videoId>"} end)
assert [media_item] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [media_item] = FastIndexingHelpers.index_and_kickoff_downloads(source)
assert [_] = Tasks.list_tasks_for(media_item, "MediaDownloadWorker") assert [_] = Tasks.list_tasks_for(media_item, "MediaDownloadWorker")
end end
@ -89,7 +101,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
{:ok, media_attributes_return_fixture()} {:ok, media_attributes_return_fixture()}
end) end)
FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
test "sets use_cookies if the source uses cookies" do test "sets use_cookies if the source uses cookies" do
@ -103,7 +115,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
source = source_fixture(%{use_cookies: true}) source = source_fixture(%{use_cookies: true})
assert [%MediaItem{}] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [%MediaItem{}] = FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
test "does not set use_cookies if the source does not use cookies" do test "does not set use_cookies if the source does not use cookies" do
@ -117,7 +129,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
source = source_fixture(%{use_cookies: false}) source = source_fixture(%{use_cookies: false})
assert [%MediaItem{}] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [%MediaItem{}] = FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
test "does not enqueue a download job if the media item does not match the format rules" do test "does not enqueue a download job if the media item does not match the format rules" do
@ -142,7 +154,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
{:ok, output} {:ok, output}
end) end)
assert [%MediaItem{}] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [%MediaItem{}] = FastIndexingHelpers.index_and_kickoff_downloads(source)
refute_enqueued(worker: MediaDownloadWorker) refute_enqueued(worker: MediaDownloadWorker)
end end
@ -154,7 +166,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
{:ok, "{}"} {:ok, "{}"}
end) end)
assert [] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [] = FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
test "does not blow up if a media item causes a yt-dlp error", %{source: source} do test "does not blow up if a media item causes a yt-dlp error", %{source: source} do
@ -164,11 +176,11 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
{:error, "message", 1} {:error, "message", 1}
end) end)
assert [] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [] = FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
end end
describe "kickoff_download_tasks_from_youtube_rss_feed/1 when testing backends" do describe "index_and_kickoff_downloads/1 when testing backends" do
test "uses the YouTube API if it is enabled", %{source: source} do test "uses the YouTube API if it is enabled", %{source: source} do
expect(HTTPClientMock, :get, fn url, _headers -> expect(HTTPClientMock, :get, fn url, _headers ->
assert url =~ "https://youtube.googleapis.com/youtube/v3/playlistItems" assert url =~ "https://youtube.googleapis.com/youtube/v3/playlistItems"
@ -178,7 +190,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
Settings.set(youtube_api_key: "test_key") Settings.set(youtube_api_key: "test_key")
assert [] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [] = FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
test "the YouTube API creates records as expected", %{source: source} do test "the YouTube API creates records as expected", %{source: source} do
@ -188,7 +200,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
Settings.set(youtube_api_key: "test_key") Settings.set(youtube_api_key: "test_key")
assert [%MediaItem{}] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [%MediaItem{}] = FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
test "RSS is used as a backup if the API fails", %{source: source} do test "RSS is used as a backup if the API fails", %{source: source} do
@ -197,7 +209,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
Settings.set(youtube_api_key: "test_key") Settings.set(youtube_api_key: "test_key")
assert [%MediaItem{}] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [%MediaItem{}] = FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
test "RSS is used if the API is not enabled", %{source: source} do test "RSS is used if the API is not enabled", %{source: source} do
@ -209,7 +221,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpersTest do
Settings.set(youtube_api_key: nil) Settings.set(youtube_api_key: nil)
assert [%MediaItem{}] = FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed(source) assert [%MediaItem{}] = FastIndexingHelpers.index_and_kickoff_downloads(source)
end end
end end
end end