Resolved bug where non-pending media would notify when fast indexing (#187)
This commit is contained in:
@@ -7,9 +7,11 @@ defmodule Pinchflat.Downloading.DownloadingHelpers do
|
||||
|
||||
require Logger
|
||||
|
||||
alias Pinchflat.Repo
|
||||
alias Pinchflat.Media
|
||||
alias Pinchflat.Tasks
|
||||
alias Pinchflat.Sources.Source
|
||||
alias Pinchflat.Media.MediaItem
|
||||
alias Pinchflat.Downloading.MediaDownloadWorker
|
||||
|
||||
@doc """
|
||||
@@ -43,4 +45,23 @@ defmodule Pinchflat.Downloading.DownloadingHelpers do
|
||||
|> Media.list_pending_media_items_for()
|
||||
|> Enum.each(&Tasks.delete_pending_tasks_for/1)
|
||||
end
|
||||
|
||||
@doc """
|
||||
Takes a single media item and enqueues a download job if the media should be
|
||||
downloaded, based on the source's download settings and whether media is
|
||||
considered pending.
|
||||
|
||||
Returns {:ok, %Task{}} | {:error, :should_not_download} | {:error, any()}
|
||||
"""
|
||||
def kickoff_download_if_pending(%MediaItem{} = media_item) do
|
||||
media_item = Repo.preload(media_item, :source)
|
||||
|
||||
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})")
|
||||
|
||||
MediaDownloadWorker.kickoff_with_task(media_item)
|
||||
else
|
||||
{:error, :should_not_download}
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -5,64 +5,45 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpers do
|
||||
Many of these methods are made to be kickoff or be consumed by workers.
|
||||
"""
|
||||
|
||||
require Logger
|
||||
|
||||
alias Pinchflat.Repo
|
||||
alias Pinchflat.Media
|
||||
alias Pinchflat.Sources.Source
|
||||
alias Pinchflat.Media.MediaQuery
|
||||
alias Pinchflat.FastIndexing.YoutubeRss
|
||||
alias Pinchflat.Downloading.MediaDownloadWorker
|
||||
alias Pinchflat.FastIndexing.MediaIndexingWorker
|
||||
alias Pinchflat.Downloading.DownloadingHelpers
|
||||
|
||||
alias Pinchflat.YtDlp.Media, as: YtDlpMedia
|
||||
|
||||
@doc """
|
||||
Fetches new media IDs from a source's YouTube RSS feed and kicks off indexing tasks
|
||||
for any new media items. See comments in `MediaIndexingWorker` for more info on the
|
||||
Fetches new media IDs from a source's YouTube RSS feed, indexes them, and kicks off downloading
|
||||
tasks for any pending media items. See comments in `FastIndexingWorker` for more info on the
|
||||
order of operations and how this fits into the indexing process.
|
||||
|
||||
Despite the similar name to `kickoff_fast_indexing_task`, this does work differently.
|
||||
`kickoff_fast_indexing_task` starts a task that _calls_ this function whereas this
|
||||
function starts individual indexing tasks for each new media item. I think it does
|
||||
make sense grammatically, but I could see how that's confusing.
|
||||
|
||||
Returns [binary()] where each binary is the media ID of a new media item.
|
||||
Returns [%MediaItem{}] where each item is a new media item that was created _but not necessarily
|
||||
downloaded_.
|
||||
"""
|
||||
def kickoff_indexing_tasks_from_youtube_rss_feed(%Source{} = source) do
|
||||
def kickoff_download_tasks_from_youtube_rss_feed(%Source{} = source) do
|
||||
{:ok, media_ids} = YoutubeRss.get_recent_media_ids_from_rss(source)
|
||||
existing_media_items = list_media_items_by_media_id_for(source, media_ids)
|
||||
new_media_ids = media_ids -- Enum.map(existing_media_items, & &1.media_id)
|
||||
|
||||
Enum.each(new_media_ids, fn media_id ->
|
||||
url = "https://www.youtube.com/watch?v=#{media_id}"
|
||||
maybe_new_media_items =
|
||||
Enum.map(new_media_ids, fn media_id ->
|
||||
case create_media_item_from_media_id(source, media_id) do
|
||||
{:ok, media_item} ->
|
||||
media_item
|
||||
|
||||
MediaIndexingWorker.kickoff_with_task(source, url)
|
||||
end)
|
||||
|
||||
new_media_ids
|
||||
end
|
||||
|
||||
@doc """
|
||||
Indexes a single media item for a source and enqueues a download job if the
|
||||
media should be downloaded. This method creates the media item record so it's
|
||||
the one-stop-shop for adding a media item (and possibly downloading it) just
|
||||
by a URL and source.
|
||||
|
||||
Returns {:ok, media_item} | {:error, any()}
|
||||
"""
|
||||
def index_and_enqueue_download_for_media_item(%Source{} = source, url) do
|
||||
maybe_media_item = create_media_item_from_url(source, url)
|
||||
|
||||
case maybe_media_item do
|
||||
{:ok, media_item} ->
|
||||
if source.download_media && Media.pending_download?(media_item) do
|
||||
MediaDownloadWorker.kickoff_with_task(media_item)
|
||||
err ->
|
||||
Logger.error("Error creating media item '#{media_id}' from URL: #{inspect(err)}")
|
||||
nil
|
||||
end
|
||||
end)
|
||||
|
||||
{:ok, media_item}
|
||||
DownloadingHelpers.enqueue_pending_download_tasks(source)
|
||||
|
||||
err ->
|
||||
err
|
||||
end
|
||||
Enum.filter(maybe_new_media_items, & &1)
|
||||
end
|
||||
|
||||
defp list_media_items_by_media_id_for(source, media_ids) do
|
||||
@@ -72,9 +53,15 @@ defmodule Pinchflat.FastIndexing.FastIndexingHelpers do
|
||||
|> Repo.all()
|
||||
end
|
||||
|
||||
defp create_media_item_from_url(source, url) do
|
||||
{:ok, media_attrs} = YtDlpMedia.get_media_attributes(url)
|
||||
defp create_media_item_from_media_id(source, media_id) do
|
||||
url = "https://www.youtube.com/watch?v=#{media_id}"
|
||||
|
||||
Media.create_media_item_from_backend_attrs(source, media_attrs)
|
||||
case YtDlpMedia.get_media_attributes(url) do
|
||||
{:ok, media_attrs} ->
|
||||
Media.create_media_item_from_backend_attrs(source, media_attrs)
|
||||
|
||||
err ->
|
||||
err
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
@@ -10,6 +10,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingWorker do
|
||||
|
||||
alias __MODULE__
|
||||
alias Pinchflat.Tasks
|
||||
alias Pinchflat.Media
|
||||
alias Pinchflat.Sources
|
||||
alias Pinchflat.Settings
|
||||
alias Pinchflat.Sources.Source
|
||||
@@ -28,9 +29,21 @@ defmodule Pinchflat.FastIndexing.FastIndexingWorker do
|
||||
end
|
||||
|
||||
@doc """
|
||||
Kicks off the fast indexing process for a source, reschedules the job to run again
|
||||
once complete. See `MediaCollectionIndexingWorker` and `MediaIndexingWorker` comments
|
||||
for more
|
||||
Similar to `MediaCollectionIndexingWorker`, but for working with RSS feeds.
|
||||
`MediaCollectionIndexingWorker` should be preferred in general, but this is
|
||||
useful for downloading small batches of media items via fast indexing.
|
||||
|
||||
Only kicks off downloads for media that _should_ be downloaded
|
||||
(ie: the source is set to download and the media matches the profile's format preferences)
|
||||
|
||||
Order of operations:
|
||||
1. FastIndexingWorker (this module) periodically checks the YouTube RSS feed for new media.
|
||||
with `FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed`
|
||||
2. If the above `kickoff_download_tasks_from_youtube_rss_feed` 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
|
||||
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
|
||||
then reschedules itself to run again in the future.
|
||||
|
||||
Returns :ok | {:ok, :job_exists} | {:ok, %Task{}}
|
||||
"""
|
||||
@@ -39,7 +52,7 @@ defmodule Pinchflat.FastIndexing.FastIndexingWorker do
|
||||
source = Sources.get_source!(source_id)
|
||||
|
||||
if source.fast_index do
|
||||
perform_indexing_and_notification(source)
|
||||
perform_indexing_and_send_notification(source)
|
||||
reschedule_indexing(source)
|
||||
else
|
||||
:ok
|
||||
@@ -49,11 +62,17 @@ defmodule Pinchflat.FastIndexing.FastIndexingWorker do
|
||||
Ecto.StaleEntryError -> Logger.info("#{__MODULE__} discarded: source #{source_id} stale")
|
||||
end
|
||||
|
||||
defp perform_indexing_and_notification(source) do
|
||||
defp perform_indexing_and_send_notification(source) do
|
||||
apprise_server = Settings.get!(:apprise_server)
|
||||
new_media_items = FastIndexingHelpers.kickoff_indexing_tasks_from_youtube_rss_feed(source)
|
||||
|
||||
SourceNotifications.send_new_media_notification(apprise_server, source, length(new_media_items))
|
||||
new_media_items =
|
||||
source
|
||||
|> FastIndexingHelpers.kickoff_download_tasks_from_youtube_rss_feed()
|
||||
|> Enum.filter(&Media.pending_download?(&1))
|
||||
|
||||
if source.download_media do
|
||||
SourceNotifications.send_new_media_notification(apprise_server, source, length(new_media_items))
|
||||
end
|
||||
end
|
||||
|
||||
defp reschedule_indexing(source) do
|
||||
|
||||
@@ -1,69 +0,0 @@
|
||||
defmodule Pinchflat.FastIndexing.MediaIndexingWorker do
|
||||
@moduledoc false
|
||||
|
||||
use Oban.Worker,
|
||||
queue: :media_indexing,
|
||||
unique: [period: :infinity, states: [:available, :scheduled, :retryable]],
|
||||
tags: ["media_source", "media_indexing"]
|
||||
|
||||
require Logger
|
||||
|
||||
alias __MODULE__
|
||||
alias Pinchflat.Tasks
|
||||
alias Pinchflat.Sources
|
||||
alias Pinchflat.FastIndexing.FastIndexingHelpers
|
||||
|
||||
@doc """
|
||||
Starts the fast media indexing worker and creates a task for the source.
|
||||
|
||||
Returns {:ok, %Task{}} | {:error, :duplicate_job} | {:error, %Ecto.Changeset{}}
|
||||
"""
|
||||
def kickoff_with_task(source, media_url, opts \\ []) do
|
||||
%{id: source.id, media_url: media_url}
|
||||
|> MediaIndexingWorker.new(opts)
|
||||
|> Tasks.create_job_with_task(source)
|
||||
end
|
||||
|
||||
@doc """
|
||||
Similar to `MediaCollectionIndexingWorker`, but for individual media items.
|
||||
Does not reschedule or check anything to do with a source's indexing
|
||||
frequency - only collects initial metadata then kicks off a download.
|
||||
`MediaCollectionIndexingWorker` should be preferred in general, but this is
|
||||
useful for downloading one-off media items based on a URL (like for fast indexing).
|
||||
|
||||
Only downloads media that _should_ be downloaded (ie: the source is set to download
|
||||
and the media matches the profile's format preferences)
|
||||
|
||||
Order of operations:
|
||||
1. FastIndexingHelpers.kickoff_indexing_tasks_from_youtube_rss_feed/1 (which is running
|
||||
in its own worker) periodically checks the YouTube RSS feed for new media
|
||||
2. If new media is found, it enqueues a MediaIndexingWorker (this module) for each new media
|
||||
item
|
||||
3. This worker fetches the media metadata and uses that to determine if it should be
|
||||
downloaded. If so, it enqueues a MediaDownloadWorker
|
||||
|
||||
Each is a worker because they all either need to be scheduled periodically or call out to
|
||||
an external service and will be long-running. They're split into different jobs to separate
|
||||
retry logic for each step and allow us to better optimize various queues (eg: the indexing
|
||||
steps can keep running while the slow download steps are worked through).
|
||||
|
||||
Returns :ok
|
||||
"""
|
||||
@impl Oban.Worker
|
||||
def perform(%Oban.Job{args: %{"id" => source_id, "media_url" => media_url}}) do
|
||||
source = Sources.get_source!(source_id)
|
||||
|
||||
case FastIndexingHelpers.index_and_enqueue_download_for_media_item(source, media_url) do
|
||||
{:ok, media_item} ->
|
||||
Logger.debug("Indexed and enqueued download for url: #{media_url} (media item: #{media_item.id})")
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.debug("Failed to index and enqueue download for url: #{media_url} (reason: #{inspect(reason)})")
|
||||
end
|
||||
|
||||
:ok
|
||||
rescue
|
||||
Ecto.NoResultsError -> Logger.info("#{__MODULE__} discarded: source #{source_id} not found")
|
||||
Ecto.StaleEntryError -> Logger.info("#{__MODULE__} discarded: source #{source_id} stale")
|
||||
end
|
||||
end
|
||||
@@ -61,9 +61,8 @@ defmodule Pinchflat.SlowIndexing.MediaCollectionIndexingWorker do
|
||||
5. If the source uses fast indexing, that job is kicked off as well. It
|
||||
uses RSS to run a smaller, faster, and more frequent index. That job
|
||||
handles rescheduling itself but largely has a similar behaviour to this
|
||||
job in that it kicks off index and maybe download jobs. The biggest difference
|
||||
is that an index job is kicked off _for each new media item_ as opposed
|
||||
to one larger index job. Check out `MediaIndexingWorker` comments for more.
|
||||
job in that it runs and index and maybe kicks off media download jobs.
|
||||
Check out `FastIndexingWorker` comments for more.
|
||||
6. If the job reschedules, the cycle from step 3 repeats until the heat death
|
||||
of the universe. The user changing things like the index frequency can
|
||||
dequeue or reschedule jobs as well
|
||||
|
||||
@@ -16,7 +16,6 @@ defmodule Pinchflat.SlowIndexing.SlowIndexingHelpers do
|
||||
alias Pinchflat.YtDlp.MediaCollection
|
||||
alias Pinchflat.Downloading.DownloadingHelpers
|
||||
alias Pinchflat.SlowIndexing.FileFollowerServer
|
||||
alias Pinchflat.Downloading.MediaDownloadWorker
|
||||
alias Pinchflat.SlowIndexing.MediaCollectionIndexingWorker
|
||||
|
||||
alias Pinchflat.YtDlp.Media, as: YtDlpMedia
|
||||
@@ -29,7 +28,6 @@ defmodule Pinchflat.SlowIndexing.SlowIndexingHelpers do
|
||||
"""
|
||||
def kickoff_indexing_task(%Source{} = source, job_args \\ %{}, job_opts \\ []) do
|
||||
Tasks.delete_pending_tasks_for(source, "FastIndexingWorker")
|
||||
Tasks.delete_pending_tasks_for(source, "MediaIndexingWorker")
|
||||
Tasks.delete_pending_tasks_for(source, "MediaCollectionIndexingWorker")
|
||||
|
||||
MediaCollectionIndexingWorker.kickoff_with_task(source, job_args, job_opts)
|
||||
@@ -127,11 +125,7 @@ defmodule Pinchflat.SlowIndexing.SlowIndexingHelpers do
|
||||
|
||||
case Media.create_media_item_from_backend_attrs(source, media_attrs) do
|
||||
{:ok, %MediaItem{} = media_item} ->
|
||||
if source.download_media && Media.pending_download?(media_item) do
|
||||
Logger.debug("FileFollowerServer Handler: Enqueuing download task for #{inspect(media_attrs)}")
|
||||
|
||||
MediaDownloadWorker.kickoff_with_task(media_item)
|
||||
end
|
||||
DownloadingHelpers.kickoff_download_if_pending(media_item)
|
||||
|
||||
{:error, changeset} ->
|
||||
changeset
|
||||
|
||||
@@ -288,7 +288,6 @@ defmodule Pinchflat.Sources do
|
||||
|
||||
%{index_frequency_minutes: _} ->
|
||||
Tasks.delete_pending_tasks_for(source, "FastIndexingWorker")
|
||||
Tasks.delete_pending_tasks_for(source, "MediaIndexingWorker")
|
||||
Tasks.delete_pending_tasks_for(source, "MediaCollectionIndexingWorker")
|
||||
|
||||
_ ->
|
||||
|
||||
Reference in New Issue
Block a user