Fast indexing (#58)

* Made method to getting singular media details; Renamed other related method

* Takes a fun and flirty digression to remove abstractions around yt-dlp since I'm 100% committed to using it exclusively

* Removed commented test code

* Lays the groundwork for fast indexing

* Added module for working with youtube RSS feed

* Added methods to kick off indexing workers from RSS response

* Improve short detection (#59)

* Made media attribute-related yt-dlp calls return a struct

* Added shorts attribute to media items

* Added ability to discern a short from yt-dlp response

* Updated search to use new shorts attribute

* Fast index UI (#63)

* Added fast_index field and adds it to source form

* Added fast indexing to source changeset operations

* Added fast indexing worker and updated other modules to start using it

* Handled fast index worker on source update

* Add support modals (#65)

* Added fast indexing upgrade modal

* Improved modal on smaller screens

* Updated links to work again

* Added donation modal

* Reverted source fast index to 15 minutes

* Removed unneeded HTML attributes from old alpine approach
This commit is contained in:
Kieran
2024-03-10 14:36:34 -07:00
committed by GitHub
parent 1caddd86c7
commit dc0313d875
72 changed files with 1698 additions and 626 deletions
+51
View File
@@ -0,0 +1,51 @@
defmodule Pinchflat.Api.YoutubeRss do
@moduledoc """
Methods for interacting with YouTube RSS feeds
"""
require Logger
alias Pinchflat.Sources.Source
@doc """
Fetches the recent media IDs from a YouTube RSS feed for a given source.
Returns {:ok, [binary()]} | {:error, binary()}
"""
def get_recent_media_ids_from_rss(%Source{} = source) do
Logger.debug("Fetching recent media IDs from YouTube RSS feed for source: #{source.collection_id}")
case http_client().get(rss_url_for_source(source)) do
{:ok, response} ->
response = to_string(response)
media_id_regex = ~r/<yt:videoId>(.*?)<\/yt:videoId>/
# Don't get on me about using regex to search XML.
# The content is known, well-formed, and simple.
media_ids =
media_id_regex
|> Regex.scan(response)
|> Enum.map(fn [_, id] -> String.trim(id) end)
|> Enum.filter(&(String.length(&1) > 0))
|> Enum.uniq()
Logger.debug("Media ids fetched from RSS: #{inspect(media_ids)}")
{:ok, media_ids}
{:error, _reason} ->
{:error, "Failed to fetch RSS feed"}
end
end
defp rss_url_for_source(source) do
case source.collection_type do
:channel -> "https://www.youtube.com/feeds/videos.xml?channel_id=#{source.collection_id}"
:playlist -> "https://www.youtube.com/feeds/videos.xml?playlist_id=#{source.collection_id}"
end
end
defp http_client do
Application.get_env(:pinchflat, :http_client, Pinchflat.HTTP.HTTPClient)
end
end
+2
View File
@@ -4,5 +4,7 @@ defmodule Pinchflat.HTTP.HTTPBehaviour do
so I can use Mox to create an HTTP mock
"""
@callback get(String.t()) :: {:ok, String.t()} | {:error, String.t()}
@callback get(String.t(), Keyword.t()) :: {:ok, String.t()} | {:error, String.t()}
@callback get(String.t(), Keyword.t(), Keyword.t()) :: {:ok, String.t()} | {:error, String.t()}
end
+38 -9
View File
@@ -31,6 +31,22 @@ defmodule Pinchflat.Media do
|> Repo.all()
end
@doc """
Fetches all media items belonging to a given source that have a media_id in the given list.
Useful for determining the what media items we DON'T already have for fast indexing.
NOTE: These queries are getting a little tedious. When I have the time, I should see about
implementing a query pattern and having these compose queries from a common base. This would
also let me compose simple queries in the module using them for one-off methods
Returns [%MediaItem{}, ...].
"""
def list_media_items_by_media_id_for(%Source{} = source, media_ids) do
MediaItem
|> where([mi], mi.source_id == ^source.id and mi.media_id in ^media_ids)
|> Repo.all()
end
@doc """
Returns a list of pending media_items for a given source, where
pending means the `media_filepath` is `nil` AND the media_item
@@ -156,7 +172,9 @@ defmodule Pinchflat.Media do
end
@doc """
Creates a media_item. Returns {:ok, %MediaItem{}} | {:error, %Ecto.Changeset{}}.
Creates a media_item.
Returns {:ok, %MediaItem{}} | {:error, %Ecto.Changeset{}}
"""
def create_media_item(attrs) do
%MediaItem{}
@@ -165,7 +183,21 @@ defmodule Pinchflat.Media do
end
@doc """
Updates a media_item. Returns {:ok, %MediaItem{}} | {:error, %Ecto.Changeset{}}.
Creates a media item from the attributes returned by the video backend
(read: yt-dlp)
Returns {:ok, %MediaItem{}} | {:error, %Ecto.Changeset{}}
"""
def create_media_item_from_backend_attrs(source, media_attrs_struct) do
%{source_id: source.id}
|> Map.merge(Map.from_struct(media_attrs_struct))
|> create_media_item()
end
@doc """
Updates a media_item.
Returns {:ok, %MediaItem{}} | {:error, %Ecto.Changeset{}}
"""
def update_media_item(%MediaItem{} = media_item, attrs) do
media_item
@@ -177,7 +209,7 @@ defmodule Pinchflat.Media do
Deletes a media_item and its associated tasks.
Can optionally delete the media_item's files.
Returns {:ok, %MediaItem{}} | {:error, %Ecto.Changeset{}}.
Returns {:ok, %MediaItem{}} | {:error, %Ecto.Changeset{}}
"""
def delete_media_item(%MediaItem{} = media_item, opts \\ []) do
delete_files = Keyword.get(opts, :delete_files, false)
@@ -225,7 +257,7 @@ defmodule Pinchflat.Media do
{{:shorts_behaviour, :only}, %{livestream_behaviour: :only}} ->
dynamic(
[mi],
^dynamic and (mi.livestream == true or fragment("LOWER(?) LIKE LOWER(?)", mi.original_url, "%/shorts/%"))
^dynamic and (mi.livestream == true or mi.short_form_content == true)
)
# Technically redundant, but makes the other clauses easier to parse
@@ -234,16 +266,13 @@ defmodule Pinchflat.Media do
dynamic
{{:shorts_behaviour, :only}, _} ->
# return records with /shorts/ in the original_url
dynamic([mi], ^dynamic and fragment("LOWER(?) LIKE LOWER(?)", mi.original_url, "%/shorts/%"))
dynamic([mi], ^dynamic and mi.short_form_content == true)
{{:livestream_behaviour, :only}, _} ->
# return records with livestream: true
dynamic([mi], ^dynamic and mi.livestream == true)
{{:shorts_behaviour, :exclude}, %{livestream_behaviour: lb}} when lb != :only ->
# return records without /shorts/ in the original_url
dynamic([mi], ^dynamic and fragment("LOWER(?) NOT LIKE LOWER(?)", mi.original_url, "%/shorts/%"))
dynamic([mi], ^dynamic and mi.short_form_content == false)
{{:livestream_behaviour, :exclude}, %{shorts_behaviour: sb}} when sb != :only ->
# return records with livestream: false
+5 -3
View File
@@ -12,14 +12,15 @@ defmodule Pinchflat.Media.MediaItem do
alias Pinchflat.Media.MediaItemSearchIndex
@allowed_fields [
# these fields are captured on indexing
# these fields are captured on indexing (and again on download)
:title,
:media_id,
:description,
:original_url,
:livestream,
:source_id,
# these fields are captured on download
:short_form_content,
# these fields are captured only on download
:media_downloaded_at,
:media_filepath,
:media_size_bytes,
@@ -27,7 +28,7 @@ defmodule Pinchflat.Media.MediaItem do
:thumbnail_filepath,
:metadata_filepath
]
@required_fields ~w(title original_url livestream media_id source_id)a
@required_fields ~w(title original_url livestream media_id source_id short_form_content)a
schema "media_items" do
field :title, :string
@@ -35,6 +36,7 @@ defmodule Pinchflat.Media.MediaItem do
field :description, :string
field :original_url, :string
field :livestream, :boolean, default: false
field :short_form_content, :boolean, default: false
field :media_downloaded_at, :utc_datetime
field :media_filepath, :string
@@ -1,26 +0,0 @@
defmodule Pinchflat.MediaClient.Backends.YtDlp.Media do
@moduledoc """
Contains utilities for working with singular pieces of media
"""
@doc """
Downloads a single piece of media (and possibly its metadata) directly to its
final destination. Returns the parsed JSON output from yt-dlp.
Returns {:ok, map()} | {:error, any, ...}.
"""
def download(url, command_opts \\ []) do
opts = [:no_simulate] ++ command_opts
with {:ok, output} <- backend_runner().run(url, opts, "after_move:%()j"),
{:ok, parsed_json} <- Phoenix.json_library().decode(output) do
{:ok, parsed_json}
else
err -> err
end
end
defp backend_runner do
Application.get_env(:pinchflat, :yt_dlp_runner)
end
end
@@ -1,47 +0,0 @@
defmodule Pinchflat.MediaClient.SourceDetails do
@moduledoc """
This is the integration layer for actually working with sources.
Technically hardcodes the yt-dlp backend for now, but should leave
it open-ish for future expansion (just in case).
"""
alias Pinchflat.Sources.Source
alias Pinchflat.MediaClient.Backends.YtDlp.MediaCollection, as: YtDlpSource
@doc """
Gets a source's ID and name from its URL using the given backend.
Returns {:ok, map()} | {:error, any, ...}.
"""
def get_source_details(source_url, backend \\ :yt_dlp) do
source_module(backend).get_source_details(source_url)
end
@doc """
Returns a list of basic media data maps for the given source URL OR
source record using the given backend.
Options:
- :file_listener_handler - a function that will be called with the path to the
file that will be written to by yt-dlp. This is useful for
setting up a file watcher to read the file as it gets written to.
Returns {:ok, [map()]} | {:error, any, ...}.
"""
def get_media_attributes(sourceable, opts \\ [], backend \\ :yt_dlp)
def get_media_attributes(%Source{} = source, opts, backend) do
get_media_attributes(source.collection_id, opts, backend)
end
def get_media_attributes(source_url, opts, backend) when is_binary(source_url) do
source_module(backend).get_media_attributes(source_url, opts)
end
defp source_module(backend) do
case backend do
:yt_dlp -> YtDlpSource
end
end
end
@@ -1,4 +1,4 @@
defmodule Pinchflat.MediaClient.Backends.YtDlp.MetadataFileHelpers do
defmodule Pinchflat.Metadata.MetadataFileHelpers do
@moduledoc """
Provides methods for creating/downloading/storing related metadata
out-of-band of the normal yt-dlp backend process.
@@ -1,4 +1,4 @@
defmodule Pinchflat.MediaClient.Backends.YtDlp.MetadataParser do
defmodule Pinchflat.Metadata.MetadataParser do
@moduledoc """
yt-dlp offers a LOT of metadata in its JSON response, some of which
needs to be extracted and included in various models.
@@ -25,9 +25,12 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.MetadataParser do
defp parse_media_metadata(metadata) do
%{
media_id: metadata["id"],
title: metadata["title"],
original_url: metadata["original_url"],
description: metadata["description"],
media_filepath: metadata["filepath"]
media_filepath: metadata["filepath"],
livestream: metadata["was_live"]
}
end
@@ -1,8 +1,6 @@
defmodule Pinchflat.Profiles.Options.YtDlp.OutputPathBuilder do
defmodule Pinchflat.Profiles.OutputPathBuilder do
@moduledoc """
Builds yt-dlp-friendly output paths for downloaded media
IDEA: consider making this a behaviour so I can add other backends later
"""
alias Pinchflat.RenderedString.Parser, as: TemplateParser
+48 -7
View File
@@ -11,7 +11,7 @@ defmodule Pinchflat.Sources do
alias Pinchflat.Sources.Source
alias Pinchflat.Tasks.SourceTasks
alias Pinchflat.Profiles.MediaProfile
alias Pinchflat.MediaClient.SourceDetails
alias Pinchflat.YtDlp.Backend.MediaCollection
@doc """
Returns the list of sources. Returns [%Source{}, ...]
@@ -46,6 +46,7 @@ defmodule Pinchflat.Sources do
def create_source(attrs) do
%Source{}
|> change_source_from_url(attrs)
|> maybe_change_indexing_frequency()
|> commit_and_handle_tasks()
end
@@ -62,6 +63,7 @@ defmodule Pinchflat.Sources do
def update_source(%Source{} = source, attrs) do
source
|> change_source_from_url(attrs)
|> maybe_change_indexing_frequency()
|> commit_and_handle_tasks()
end
@@ -116,7 +118,7 @@ defmodule Pinchflat.Sources do
defp add_source_details_to_changeset(source, changeset) do
%Ecto.Changeset{changes: changes} = changeset
case SourceDetails.get_source_details(changes.original_url) do
case MediaCollection.get_source_details(changes.original_url) do
{:ok, source_details} ->
add_source_details_by_collection_type(source, changeset, source_details)
@@ -151,6 +153,20 @@ defmodule Pinchflat.Sources do
change_source(source, Map.merge(changes, collection_changes))
end
defp maybe_change_indexing_frequency(changeset) do
fast_index = Ecto.Changeset.get_field(changeset, :fast_index)
if fast_index do
Ecto.Changeset.put_change(
changeset,
:index_frequency_minutes,
Source.index_frequency_when_fast_indexing()
)
else
changeset
end
end
defp commit_and_handle_tasks(changeset) do
case Repo.insert_or_update(changeset) do
{:ok, %Source{} = source} ->
@@ -188,13 +204,38 @@ defmodule Pinchflat.Sources do
# If the record has been persisted, only run indexing if the
# indexing frequency has been changed and is now greater than 0
%{__meta__: %{state: :loaded}} ->
case changeset.changes do
%{index_frequency_minutes: mins} when mins > 0 -> SourceTasks.kickoff_indexing_task(source)
%{index_frequency_minutes: _} -> Tasks.delete_pending_tasks_for(source, "MediaIndexingWorker")
_ -> :ok
end
maybe_update_slow_indexing_task(changeset, source)
maybe_update_fast_indexing_task(changeset, source)
end
{:ok, source}
end
defp maybe_update_slow_indexing_task(changeset, source) do
case changeset.changes do
%{index_frequency_minutes: mins} when mins > 0 ->
SourceTasks.kickoff_indexing_task(source)
%{index_frequency_minutes: _} ->
Tasks.delete_pending_tasks_for(source, "FastIndexingWorker")
Tasks.delete_pending_tasks_for(source, "MediaIndexingWorker")
Tasks.delete_pending_tasks_for(source, "MediaCollectionIndexingWorker")
_ ->
:ok
end
end
defp maybe_update_fast_indexing_task(changeset, source) do
case changeset.changes do
%{fast_index: true} ->
SourceTasks.kickoff_fast_indexing_task(source)
%{fast_index: false} ->
Tasks.delete_pending_tasks_for(source, "FastIndexingWorker")
_ ->
:ok
end
end
end
@@ -17,6 +17,7 @@ defmodule Pinchflat.Sources.Source do
collection_type
custom_name
index_frequency_minutes
fast_index
download_media
last_indexed_at
original_url
@@ -29,6 +30,7 @@ defmodule Pinchflat.Sources.Source do
collection_type
custom_name
index_frequency_minutes
fast_index
download_media
original_url
media_profile_id
@@ -40,6 +42,7 @@ defmodule Pinchflat.Sources.Source do
field :collection_id, :string
field :collection_type, Ecto.Enum, values: [:channel, :playlist]
field :index_frequency_minutes, :integer, default: 60 * 24
field :fast_index, :boolean, default: false
field :download_media, :boolean, default: true
field :last_indexed_at, :utc_datetime
# This should only be used for user reference going forward
@@ -62,4 +65,16 @@ defmodule Pinchflat.Sources.Source do
|> validate_required(@required_fields)
|> unique_constraint([:collection_id, :media_profile_id])
end
@doc false
def index_frequency_when_fast_indexing do
# 30 days in minutes
60 * 24 * 30
end
@doc false
def fast_index_frequency do
# minutes
15
end
end
+1
View File
@@ -32,5 +32,6 @@ defmodule Pinchflat.StartupTasks do
defp apply_default_settings do
Settings.fetch!(:onboarding, true)
Settings.fetch!(:pro_enabled, false)
end
end
+37
View File
@@ -6,6 +6,11 @@ defmodule Pinchflat.Tasks.MediaItemTasks do
do is also defined here. Essentially, a one-stop-shop for media-related tasks/workers.
"""
alias Pinchflat.Media
alias Pinchflat.Tasks
alias Pinchflat.Sources.Source
alias Pinchflat.Workers.MediaDownloadWorker
alias Pinchflat.YtDlp.Backend.Media, as: YtDlpMedia
@doc """
Fetches the file size of a media item and saves it to the database.
@@ -21,4 +26,36 @@ defmodule Pinchflat.Tasks.MediaItemTasks do
err
end
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
%{id: media_item.id}
|> MediaDownloadWorker.new()
|> Tasks.create_job_with_task(media_item)
end
{:ok, media_item}
err ->
err
end
end
defp create_media_item_from_url(source, url) do
{:ok, media_attrs} = YtDlpMedia.get_media_attributes(url)
Media.create_media_item_from_backend_attrs(source, media_attrs)
end
end
+62 -31
View File
@@ -12,31 +12,72 @@ defmodule Pinchflat.Tasks.SourceTasks do
alias Pinchflat.Tasks
alias Pinchflat.Sources
alias Pinchflat.Sources.Source
alias Pinchflat.Api.YoutubeRss
alias Pinchflat.Media.MediaItem
alias Pinchflat.MediaClient.SourceDetails
alias Pinchflat.Workers.MediaIndexingWorker
alias Pinchflat.Workers.FastIndexingWorker
alias Pinchflat.Workers.MediaDownloadWorker
alias Pinchflat.Workers.MediaIndexingWorker
alias Pinchflat.YtDlp.Backend.MediaCollection
alias Pinchflat.Workers.MediaCollectionIndexingWorker
alias Pinchflat.Utils.FilesystemUtils.FileFollowerServer
alias Pinchflat.YtDlp.Backend.Media, as: YtDlpMedia
@doc """
Starts tasks for indexing a source's media regardless of the source's indexing
frequency. It's assumed the caller will check for that.
frequency. It's assumed the caller will check for indexing frequency.
Returns {:ok, %Task{}}.
"""
def kickoff_indexing_task(%Source{} = source) do
Tasks.delete_pending_tasks_for(source, "FastIndexingWorker")
Tasks.delete_pending_tasks_for(source, "MediaIndexingWorker")
Tasks.delete_pending_tasks_for(source, "MediaCollectionIndexingWorker")
source
|> Map.take([:id])
%{id: source.id}
# Schedule this one immediately, but future ones will be on an interval
|> MediaIndexingWorker.new()
|> MediaCollectionIndexingWorker.new()
|> Tasks.create_job_with_task(source)
|> 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
@doc """
Starts tasks for running a fast indexing task for a source's media
regardless of the source's fast_index state. It's assumed the
caller will check for fast_index.
This is used for running fast index tasks on update. On creation, the
fast index is enqueued after the slow index is complete.
Returns {:ok, %Task{}}.
"""
def kickoff_fast_indexing_task(%Source{} = source) do
Tasks.delete_pending_tasks_for(source, "FastIndexingWorker")
%{id: source.id}
# Schedule this one immediately, but future ones will be on an interval
|> FastIndexingWorker.new()
|> Tasks.create_job_with_task(source)
end
@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
order of operations and how this fits into the indexing process.
Returns :ok
"""
def kickoff_indexing_tasks_from_youtube_rss_feed(%Source{} = source) do
{:ok, media_ids} = YoutubeRss.get_recent_media_ids_from_rss(source)
existing_media_items = Media.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}"
%{id: source.id, media_url: url}
|> MediaIndexingWorker.new()
|> Tasks.create_job_with_task(source)
end)
end
@doc """
@@ -65,9 +106,9 @@ defmodule Pinchflat.Tasks.SourceTasks do
def index_and_enqueue_download_for_media_items(%Source{} = source) do
# See the method definition below for more info on how file watchers work
# (important reading if you're not familiar with it)
{:ok, media_attributes} = get_media_attributes_and_setup_file_watcher(source)
{:ok, media_attributes} = get_media_attributes_for_collection_and_setup_file_watcher(source)
result = Enum.map(media_attributes, fn media_attrs -> create_media_item_from_attributes(source, media_attrs) end)
Sources.update_source(source, %{last_indexed_at: DateTime.utc_now()})
enqueue_pending_media_tasks(source)
@@ -89,8 +130,7 @@ defmodule Pinchflat.Tasks.SourceTasks do
source
|> Media.list_pending_media_items_for()
|> Enum.each(fn media_item ->
media_item
|> Map.take([:id])
%{id: media_item.id}
|> MediaDownloadWorker.new()
|> Tasks.create_job_with_task(media_item)
end)
@@ -116,7 +156,7 @@ defmodule Pinchflat.Tasks.SourceTasks do
# lines (ie: you should gracefully fail if you can't parse a line).
#
# This works in-tandem with the normal (blocking) media indexing behaviour. When
# the `get_media_attributes` method completes it'll return the FULL result to
# the `get_media_attributes_for_collection` method completes it'll return the FULL result to
# the caller for parsing. Ideally, every item in the list will have already
# been processed by the file follower, but if not, the caller handles creation
# of any media items that were missed/initially failed.
@@ -124,11 +164,11 @@ defmodule Pinchflat.Tasks.SourceTasks do
# It attempts a graceful shutdown of the file follower after the indexing is done,
# but the FileFollowerServer will also stop itself if it doesn't see any activity
# for a sufficiently long time.
defp get_media_attributes_and_setup_file_watcher(source) do
defp get_media_attributes_for_collection_and_setup_file_watcher(source) do
{:ok, pid} = FileFollowerServer.start_link()
handler = fn filepath -> setup_file_follower_watcher(pid, filepath, source) end
result = SourceDetails.get_media_attributes(source.original_url, file_listener_handler: handler)
result = MediaCollection.get_media_attributes_for_collection(source.original_url, file_listener_handler: handler)
FileFollowerServer.stop(pid)
@@ -141,7 +181,8 @@ defmodule Pinchflat.Tasks.SourceTasks do
{:ok, media_attrs} ->
Logger.debug("FileFollowerServer Handler: Got media attributes: #{inspect(media_attrs)}")
create_media_item_and_enqueue_download(source, media_attrs)
media_struct = YtDlpMedia.response_to_struct(media_attrs)
create_media_item_and_enqueue_download(source, media_struct)
err ->
Logger.debug("FileFollowerServer Handler: Error decoding JSON: #{inspect(err)}")
@@ -159,8 +200,7 @@ defmodule Pinchflat.Tasks.SourceTasks do
if source.download_media && Media.pending_download?(media_item) do
Logger.debug("FileFollowerServer Handler: Enqueuing download task for #{inspect(media_attrs)}")
media_item
|> Map.take([:id])
%{id: media_item.id}
|> MediaDownloadWorker.new()
|> Tasks.create_job_with_task(media_item)
end
@@ -171,16 +211,7 @@ defmodule Pinchflat.Tasks.SourceTasks do
end
defp create_media_item_from_attributes(source, media_attrs) do
attrs = %{
source_id: source.id,
title: media_attrs["title"],
media_id: media_attrs["id"],
original_url: media_attrs["original_url"],
livestream: media_attrs["was_live"],
description: media_attrs["description"]
}
case Media.create_media_item(attrs) do
case Media.create_media_item_from_backend_attrs(source, media_attrs) do
{:ok, media_item} -> media_item
{:error, changeset} -> changeset
end
@@ -0,0 +1,42 @@
defmodule Pinchflat.Workers.FastIndexingWorker do
@moduledoc false
use Oban.Worker,
queue: :fast_indexing,
unique: [period: :infinity, states: [:available, :scheduled, :retryable]],
tags: ["media_source", "fast_indexing"]
alias __MODULE__
alias Pinchflat.Tasks
alias Pinchflat.Sources
alias Pinchflat.Sources.Source
alias Pinchflat.Tasks.SourceTasks
@impl Oban.Worker
@doc """
TODO
"""
def perform(%Oban.Job{args: %{"id" => source_id}}) do
source = Sources.get_source!(source_id)
if source.fast_index do
SourceTasks.kickoff_indexing_tasks_from_youtube_rss_feed(source)
reschedule_indexing(source)
else
:ok
end
end
defp reschedule_indexing(source) do
next_run_in = Source.fast_index_frequency() * 60
%{id: source.id}
|> FastIndexingWorker.new(schedule_in: next_run_in)
|> Tasks.create_job_with_task(source)
|> case do
{:ok, task} -> {:ok, task}
{:error, :duplicate_job} -> {:ok, :job_exists}
end
end
end
@@ -0,0 +1,109 @@
defmodule Pinchflat.Workers.MediaCollectionIndexingWorker do
@moduledoc false
use Oban.Worker,
queue: :media_collection_indexing,
unique: [period: :infinity, states: [:available, :scheduled, :retryable]],
tags: ["media_source", "media_collection_indexing"]
alias __MODULE__
alias Pinchflat.Tasks
alias Pinchflat.Sources
alias Pinchflat.Sources.Source
alias Pinchflat.Tasks.SourceTasks
alias Pinchflat.Workers.FastIndexingWorker
@impl Oban.Worker
@doc """
The ID is that of a source _record_, not a YouTube channel/playlist ID. Indexes
the provided source, kicks off downloads for each new MediaItem, and
reschedules the job to run again in the future. It will ALWAYS index a source
if it's never been indexed before, but rescheduling is determined by the
`index_frequency_minutes` field.
README: Re-scheduling here works a little different than you may expect.
The reschedule time is relative to the time the job has actually _completed_.
This has some benefits but also side effects to be aware of:
- Benefit: No chance for jobs to overlap if a job takes longer than the
scheduled interval. Less likely to hit API rate limits.
- Side effect: Intervals are "soft" and _always_ walk forward. This may cause
user confusion since a 30-minute job scheduled for every hour will
actually run every 1 hour and 30 minutes. The tradeoff of not inundating
the API with requests and also not overlapping jobs is worth it, IMO.
Order of operations:
1. The user saves a source
2. This job is automatically scheduled immediately. This happens in all cases.
3. This job indexes all content for the given source. A download job is
enqueued for each media item that should be downloaded. This can be impacted
by the `download_media` field on the source as well as the profile's
shorts/livestream behaviour. At this step we also attach a file reader
to the `yt-dlp` output file so we can create media items as they come in
for a little speedup (see SourceTasks comments for more)
4. If this job is meant to reschedule (ie: has an index frequency > 0),
it reschedules itself. If not, it runs once and does not reschedule
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.
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
NOTE: Since indexing can take a LONG time, I should check what happens if an
application restart occurs while a job is running. Will the job be lost?
IDEA: Should I use paging and do indexing in chunks? Is that even faster?
Returns :ok | {:ok, %Task{}}
"""
def perform(%Oban.Job{args: %{"id" => source_id}}) do
source = Sources.get_source!(source_id)
case {source.index_frequency_minutes, source.last_indexed_at} do
{index_freq, _} when index_freq > 0 ->
# If the indexing is on a schedule simply run indexing and reschedule
SourceTasks.index_and_enqueue_download_for_media_items(source)
maybe_enqueue_fast_indexing_task(source)
reschedule_indexing(source)
{_, nil} ->
# If the source has never been indexed, index it once
# even if it's not meant to reschedule
SourceTasks.index_and_enqueue_download_for_media_items(source)
:ok
_ ->
# If the source HAS been indexed and is not meant to reschedule,
# perform a no-op
:ok
end
end
defp reschedule_indexing(source) do
next_run_in = source.index_frequency_minutes * 60
%{id: source.id}
|> MediaCollectionIndexingWorker.new(schedule_in: next_run_in)
|> Tasks.create_job_with_task(source)
|> case do
{:ok, task} -> {:ok, task}
{:error, :duplicate_job} -> {:ok, :job_exists}
end
end
defp maybe_enqueue_fast_indexing_task(source) do
if source.fast_index do
Tasks.delete_pending_tasks_for(source, "FastIndexingWorker")
next_run_in = Source.fast_index_frequency() * 60
%{id: source.id}
|> FastIndexingWorker.new(schedule_in: next_run_in)
|> Tasks.create_job_with_task(source)
end
end
end
@@ -45,8 +45,7 @@ defmodule Pinchflat.Workers.MediaDownloadWorker do
end
defp schedule_filesystem_data_worker(media_item) do
media_item
|> Map.take([:id])
%{id: media_item.id}
|> FilesystemDataWorker.new()
|> Tasks.create_job_with_task(media_item)
|> case do
+29 -48
View File
@@ -6,67 +6,48 @@ defmodule Pinchflat.Workers.MediaIndexingWorker do
unique: [period: :infinity, states: [:available, :scheduled, :retryable]],
tags: ["media_source", "media_indexing"]
alias __MODULE__
alias Pinchflat.Tasks
require Logger
alias Pinchflat.Sources
alias Pinchflat.Tasks.SourceTasks
alias Pinchflat.Tasks.MediaItemTasks
@impl Oban.Worker
@doc """
The ID is that of a source _record_, not a YouTube channel/playlist ID. Indexes
the provided source, kicks off downloads for each new MediaItem, and
reschedules the job to run again in the future. It will ALWAYS index a source
if it's never been indexed before, but rescheduling is determined by the
`index_frequency_minutes` field.
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).
README: Re-scheduling here works a little different than you may expect.
The reschedule time is relative to the time the job has actually _completed_.
This has some benefits but also side effects to be aware of:
Only downloads media that _should_ be downloaded (ie: the source is set to download
and the media matches the profile's format preferences)
- Benefit: No chance for jobs to overlap if a job takes longer than the
scheduled interval. Less likely to hit API rate limits.
- Side effect: Intervals are "soft" and _always_ walk forward. This may cause
user confusion since a 30-minute job scheduled for every hour will
actually run every 1 hour and 30 minutes. The tradeoff of not inundating
the API with requests and also not overlapping jobs is worth it, IMO.
Order of operations:
1. SourceTasks.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
NOTE: Since indexing can take a LONG time, I should check what happens if an
application restart occurs while a job is running. Will the job be lost?
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).
IDEA: Should I use paging and do indexing in chunks? Is that even faster?
Returns :ok | {:ok, %Task{}}
Returns :ok
"""
def perform(%Oban.Job{args: %{"id" => source_id}}) do
def perform(%Oban.Job{args: %{"id" => source_id, "media_url" => media_url}}) do
source = Sources.get_source!(source_id)
case {source.index_frequency_minutes, source.last_indexed_at} do
{index_freq, _} when index_freq > 0 ->
# If the indexing is on a schedule simply run indexing and reschedule
SourceTasks.index_and_enqueue_download_for_media_items(source)
reschedule_indexing(source)
case MediaItemTasks.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})")
{_, nil} ->
# If the source has never been indexed, index it once
# even if it's not meant to reschedule
SourceTasks.index_and_enqueue_download_for_media_items(source)
:ok
_ ->
# If the source HAS been indexed and is not meant to reschedule,
# perform a no-op
:ok
{:error, reason} ->
Logger.debug("Failed to index and enqueue download for url: #{media_url} (reason: #{inspect(reason)})")
end
end
defp reschedule_indexing(source) do
source
|> Map.take([:id])
|> MediaIndexingWorker.new(schedule_in: source.index_frequency_minutes * 60)
|> Tasks.create_job_with_task(source)
|> case do
{:ok, task} -> {:ok, task}
{:error, :duplicate_job} -> {:ok, :job_exists}
end
:ok
end
end
@@ -1,6 +1,9 @@
defmodule Pinchflat.MediaClient.Backends.BackendCommandRunner do
defmodule Pinchflat.YtDlp.Backend.BackendCommandRunner do
@moduledoc """
A behaviour for running CLI commands against a downloader backend
A behaviour for running CLI commands against a downloader backend (yt-dlp).
Used so we can implement Mox for testing without actually running the
yt-dlp command.
"""
@callback run(binary(), keyword(), binary()) :: {:ok, binary()} | {:error, binary(), integer()}
@@ -1,4 +1,4 @@
defmodule Pinchflat.MediaClient.Backends.YtDlp.CommandRunner do
defmodule Pinchflat.YtDlp.Backend.CommandRunner do
@moduledoc """
Runs yt-dlp commands using the `System.cmd/3` function
"""
@@ -7,7 +7,7 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.CommandRunner do
alias Pinchflat.Utils.StringUtils
alias Pinchflat.Utils.FilesystemUtils, as: FSUtils
alias Pinchflat.MediaClient.Backends.BackendCommandRunner
alias Pinchflat.YtDlp.Backend.BackendCommandRunner
@behaviour BackendCommandRunner
@@ -25,6 +25,7 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.CommandRunner do
"""
@impl BackendCommandRunner
def run(url, command_opts, output_template, addl_opts \\ []) do
# This approach lets us mock the command for testing
command = backend_executable()
# These must stay in exactly this order, hence why I'm giving it its own variable.
# Also, can't use RAM file since yt-dlp needs a concrete filepath.
+107
View File
@@ -0,0 +1,107 @@
defmodule Pinchflat.YtDlp.Backend.Media do
@moduledoc """
Contains utilities for working with singular pieces of media
"""
@enforce_keys [
:media_id,
:title,
:description,
:original_url,
:livestream,
:short_form_content
]
defstruct [
:media_id,
:title,
:description,
:original_url,
:livestream,
:short_form_content
]
alias __MODULE__
alias Pinchflat.Utils.FunctionUtils
@doc """
Downloads a single piece of media (and possibly its metadata) directly to its
final destination. Returns the parsed JSON output from yt-dlp.
Returns {:ok, map()} | {:error, any, ...}.
"""
def download(url, command_opts \\ []) do
opts = [:no_simulate] ++ command_opts
with {:ok, output} <- backend_runner().run(url, opts, "after_move:%()j"),
{:ok, parsed_json} <- Phoenix.json_library().decode(output) do
{:ok, parsed_json}
else
err -> err
end
end
@doc """
Returns a map representing the media at the given URL.
Returns {:ok, [map()]} | {:error, any, ...}.
"""
def get_media_attributes(url) do
runner = Application.get_env(:pinchflat, :yt_dlp_runner)
command_opts = [:simulate, :skip_download]
output_template = indexing_output_template()
case runner.run(url, command_opts, output_template) do
{:ok, output} ->
output
|> Phoenix.json_library().decode!()
|> response_to_struct()
|> FunctionUtils.wrap_ok()
res ->
res
end
end
@doc """
Returns the output template for yt-dlp's indexing command.
"""
def indexing_output_template do
"%(.{id,title,was_live,webpage_url,description,aspect_ratio,duration})j"
end
@doc """
Transforms a response from yt-dlp into a struct. Interprets the response to
determine if the media is short-form content.
Returns %Media{}.
"""
def response_to_struct(response) do
%Media{
media_id: response["id"],
title: response["title"],
description: response["description"],
original_url: response["webpage_url"],
livestream: response["was_live"],
short_form_content: short_form_content?(response)
}
end
defp short_form_content?(response) do
if String.contains?(response["webpage_url"], "/shorts/") do
true
else
# Sometimes shorts are returned without /shorts/ in the URL,
# so we need to do our best to determine if it's a short. This
# WILL returns false positives, but it's a best-effort approach
# that should work for most cases. The aspect_ratio check is
# based on a gut feeling and may need to be tweaked.
response["duration"] <= 60 && response["aspect_ratio"] < 0.8
end
end
defp backend_runner do
# This approach lets us mock the command for testing
Application.get_env(:pinchflat, :yt_dlp_runner)
end
end
@@ -1,4 +1,4 @@
defmodule Pinchflat.MediaClient.Backends.YtDlp.MediaCollection do
defmodule Pinchflat.YtDlp.Backend.MediaCollection do
@moduledoc """
Contains utilities for working with collections of
media (aka: a source [ie: channels, playlists]).
@@ -8,6 +8,7 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.MediaCollection do
alias Pinchflat.Utils.FunctionUtils
alias Pinchflat.Utils.FilesystemUtils
alias Pinchflat.YtDlp.Backend.Media, as: YtDlpMedia
@doc """
Returns a list of maps representing the media in the collection.
@@ -19,10 +20,10 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.MediaCollection do
Returns {:ok, [map()]} | {:error, any, ...}.
"""
def get_media_attributes(url, addl_opts \\ []) do
def get_media_attributes_for_collection(url, addl_opts \\ []) do
runner = Application.get_env(:pinchflat, :yt_dlp_runner)
command_opts = [:simulate, :skip_download]
output_template = "%(.{id,title,was_live,original_url,description})j"
output_template = YtDlpMedia.indexing_output_template()
output_filepath = FilesystemUtils.generate_metadata_tmpfile(:json)
file_listener_handler = Keyword.get(addl_opts, :file_listener_handler, false)
@@ -35,6 +36,7 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.MediaCollection do
output
|> String.split("\n", trim: true)
|> Enum.map(&Phoenix.json_library().decode!/1)
|> Enum.map(&YtDlpMedia.response_to_struct/1)
|> FunctionUtils.wrap_ok()
res ->
@@ -73,6 +75,7 @@ defmodule Pinchflat.MediaClient.Backends.YtDlp.MediaCollection do
end
defp backend_runner do
# This approach lets us mock the command for testing
Application.get_env(:pinchflat, :yt_dlp_runner)
end
end
@@ -1,12 +1,10 @@
defmodule Pinchflat.Profiles.Options.YtDlp.DownloadOptionBuilder do
defmodule Pinchflat.YtDlp.DownloadOptionBuilder do
@moduledoc """
Builds the options for yt-dlp to download media based on the given media profile.
IDEA: consider making this a behaviour so I can add other backends later
"""
alias Pinchflat.Media.MediaItem
alias Pinchflat.Profiles.Options.YtDlp.OutputPathBuilder
alias Pinchflat.Profiles.OutputPathBuilder
@doc """
Builds the options for yt-dlp to download media based on the given media's profile.
@@ -1,25 +1,22 @@
defmodule Pinchflat.MediaClient.MediaDownloader do
@moduledoc """
This is the integration layer for actually downloading medias.
This is the integration layer for actually downloading media.
It takes into account the media profile's settings in order
to download the media with the desired options.
Technically hardcodes the yt-dlp backend for now, but should leave
it open-ish for future expansion (just in case).
"""
alias Pinchflat.Repo
alias Pinchflat.Media
alias Pinchflat.Media.MediaItem
alias Pinchflat.MediaClient.Backends.YtDlp.Media, as: YtDlpMedia
alias Pinchflat.Profiles.Options.YtDlp.DownloadOptionBuilder, as: YtDlpDownloadOptionBuilder
alias Pinchflat.MediaClient.Backends.YtDlp.MetadataParser, as: YtDlpMetadataParser
alias Pinchflat.MediaClient.Backends.YtDlp.MetadataFileHelpers, as: YtDlpMetadataHelpers
alias Pinchflat.YtDlp.Backend.Media, as: YtDlpMedia
alias Pinchflat.YtDlp.DownloadOptionBuilder, as: YtDlpDownloadOptionBuilder
alias Pinchflat.Metadata.MetadataParser, as: YtDlpMetadataParser
alias Pinchflat.Metadata.MetadataFileHelpers, as: YtDlpMetadataHelpers
@doc """
Downloads media for a media item, updating the media item based on the metadata
returned by the backend. Also saves the entire metadata response to the associated
returned by yt-dlp. Also saves the entire metadata response to the associated
media_metadata record.
NOTE: related methods (like the download worker) won't download if the media item's source
@@ -28,12 +25,12 @@ defmodule Pinchflat.MediaClient.MediaDownloader do
Returns {:ok, %MediaItem{}} | {:error, any, ...any}
"""
def download_for_media_item(%MediaItem{} = media_item, backend \\ :yt_dlp) do
def download_for_media_item(%MediaItem{} = media_item) do
item_with_preloads = Repo.preload(media_item, [:metadata, source: :media_profile])
case download_with_options(media_item.original_url, item_with_preloads, backend) do
case download_with_options(media_item.original_url, item_with_preloads) do
{:ok, parsed_json} ->
{parser, helpers} = metadata_parsers(backend)
{parser, helpers} = {YtDlpMetadataParser, YtDlpMetadataHelpers}
parsed_attrs =
parsed_json
@@ -55,29 +52,14 @@ defmodule Pinchflat.MediaClient.MediaDownloader do
end
end
defp download_with_options(url, item_with_preloads, backend) do
option_builder = option_builder(backend)
media_backend = media_backend(backend)
{:ok, options} = option_builder.build(item_with_preloads)
# def download_for_source(source, url) do
# # Create MI from source and URL
# media_item = nil
# end
media_backend.download(url, options)
end
defp download_with_options(url, item_with_preloads) do
{:ok, options} = YtDlpDownloadOptionBuilder.build(item_with_preloads)
defp option_builder(backend) do
case backend do
:yt_dlp -> YtDlpDownloadOptionBuilder
end
end
defp media_backend(backend) do
case backend do
:yt_dlp -> YtDlpMedia
end
end
defp metadata_parsers(backend) do
case backend do
:yt_dlp -> {YtDlpMetadataParser, YtDlpMetadataHelpers}
end
YtDlpMedia.download(url, options)
end
end